iop/packages/go/streamgate/stream_release_test.go

1324 lines
44 KiB
Go

package streamgate
import (
"context"
"errors"
"fmt"
"strings"
"sync"
"sync/atomic"
"testing"
"time"
)
// testNow is a deterministic timestamp used in stream releaser tests.
var releaserNow = time.Date(2026, 7, 24, 12, 0, 0, 0, time.UTC)
// testCanceller records cancel attempts and can be programmed to fail.
type testCanceller struct {
cancelCount int
cancelErr error
onCancel func()
}
func (c *testCanceller) CancelAttempt(ctx context.Context, attemptID string) error {
c.cancelCount++
if c.onCancel != nil {
c.onCancel()
}
return c.cancelErr
}
// setupReleaser creates a StreamReleaser with a rolling_window EvidenceTail,
// CommitBoundary, and canceller. The tail is configured so that appending
// enough events triggers an epoch via the rolling threshold.
//
// evidenceRunes controls the rolling threshold; maxBuffer is the hard limit.
func setupReleaser(t *testing.T, channel string, evidenceRunes, maxBuffer int) (*StreamReleaser, *testSink, *testCanceller) {
t.Helper()
// Create a rolling_window plan with a low threshold so a few events
// trigger an epoch quickly.
kinds := []EventKind{EventKindTextDelta, EventKindReasoningDelta}
var req FilterHoldRequirement
if maxBuffer > 0 {
var err error
req, err = NewFilterHoldRequirementRollingWithMaxBuffer(channel, kinds, evidenceRunes, maxBuffer)
if err != nil {
t.Fatalf("NewFilterHoldRequirementRollingWithMaxBuffer: %v", err)
}
} else {
var err error
req, err = NewFilterHoldRequirementRolling(channel, kinds, evidenceRunes)
if err != nil {
t.Fatalf("NewFilterHoldRequirementRolling: %v", err)
}
}
binding, err := NewFilterHoldBinding("f1", req, true)
if err != nil {
t.Fatalf("NewFilterHoldBinding: %v", err)
}
plan, err := NewEvidencePlan([]FilterHoldBinding{binding})
if err != nil {
t.Fatalf("NewEvidencePlan: %v", err)
}
tail, err := NewEvidenceTail(plan)
if err != nil {
t.Fatalf("NewEvidenceTail: %v", err)
}
sink := &testSink{}
boundary, err := NewCommitBoundary(sink)
if err != nil {
t.Fatalf("NewCommitBoundary: %v", err)
}
canceller := &testCanceller{}
releaser, err := NewStreamReleaser(tail, boundary, canceller)
if err != nil {
t.Fatalf("NewStreamReleaser: %v", err)
}
return releaser, sink, canceller
}
// appendDataEvents creates and appends text_delta events to the tail until
// an epoch is triggered (rolling threshold reached). Returns the epoch.
func appendDataEvents(t *testing.T, tail *EvidenceTail, channel string, count int) EvidenceEpoch {
t.Helper()
return appendDataEventsDistinct(t, tail, channel, count, func(i int) string { return "chunk" })
}
// appendDataEventsDistinct creates and appends text_delta events with
// position-distinguishable text. The textFn maps each event index to its
// distinct text content. This is test-only to verify that suffix, token,
// cursor and consumed assertions actually observe payload identity rather
// than passing trivially when every event has identical text.
func appendDataEventsDistinct(t *testing.T, tail *EvidenceTail, channel string, count int, textFn func(i int) string) EvidenceEpoch {
t.Helper()
var epoch EvidenceEpoch
for i := 0; i < count; i++ {
ne, err := NewTextDeltaEvent(channel, textFn(i), releaserNow)
if err != nil {
t.Fatalf("NewTextDeltaEvent: %v", err)
}
e, signal, err := tail.Append(ne)
if err != nil {
t.Fatalf("Append: %v", err)
}
if signal == EvidenceTailSignalThreshold {
epoch = e
}
}
return epoch
}
// createTextReleaseEvent creates a ReleaseEvent of kind text_delta.
func createTextReleaseEvent(channel, text string) ReleaseEvent {
ev, err := NewReleaseTextDeltaEvent(channel, text, releaserNow)
if err != nil {
panic(err)
}
return ev
}
// createSuccessTerminalResult creates a TerminalResult of kind success.
func createSuccessTerminalResult(channel string) TerminalResult {
tr, err := NewSuccessTerminalResult(channel, releaserNow)
if err != nil {
panic(err)
}
return tr
}
// createErrorTerminalResult creates an error TerminalResult.
func createErrorTerminalResult(channel string) TerminalResult {
desc, err := NewExternalDescriptor("generic", "internal_error", "an_internal_error", "")
if err != nil {
panic(fmt.Sprintf("createErrorTerminalResult NewExternalDescriptor failed: %v", err))
}
causes, err := NewFailureCauseChain([]FailureCause{})
if err != nil {
panic(fmt.Sprintf("createErrorTerminalResult NewFailureCauseChain failed: %v", err))
}
tr, err := NewErrorTerminalResult(channel, desc, causes, releaserNow)
if err != nil {
panic(fmt.Sprintf("createErrorTerminalResult NewErrorTerminalResult failed: %v", err))
}
return tr
}
// ---------------------------------------------------------------------------
// REVIEW_API-2: Evidence prepare/confirm with sink progress
// ---------------------------------------------------------------------------
// TestStreamReleaserConfirmsExactlySinkProgress verifies zero/partial/full
// sink progress where confirmed cursor matches sink success count.
func TestStreamReleaserConfirmsExactlySinkProgress(t *testing.T) {
ctx := context.Background()
t.Run("zero_sink_progress", func(t *testing.T) {
releaser, sink, _ := setupReleaser(t, "ch-zero", 2, 100)
// Begin attempt and stage response start on boundary.
releaser.boundary.BeginAttempt("a1")
rs, _ := NewResponseStart("ch-zero", 200, nil, releaserNow)
releaser.boundary.StageResponseStart("a1", rs)
// Append data events to trigger epoch.
epoch := appendDataEvents(t, releaser.tail, "ch-zero", 3)
// All releases fail.
sink.failAll = true
prog, err := releaser.ReleaseEpoch(ctx, "a1", epoch.ID())
if err == nil {
t.Error("expected error from ReleaseEpoch")
}
if prog.ReleasedEvents() != 0 {
t.Errorf("ReleasedEvents = %d, want 0 (all events failed)", prog.ReleasedEvents())
}
})
t.Run("partial_sink_progress", func(t *testing.T) {
releaser, sink, _ := setupReleaser(t, "ch-partial", 2, 100)
releaser.boundary.BeginAttempt("a1")
rs, _ := NewResponseStart("ch-partial", 200, nil, releaserNow)
releaser.boundary.StageResponseStart("a1", rs)
epoch := appendDataEvents(t, releaser.tail, "ch-partial", 5)
// 2 successes, then fail.
sink.failAfterSuccesses = 2
prog, err := releaser.ReleaseEpoch(ctx, "a1", epoch.ID())
if err == nil {
t.Error("expected partial error")
}
if prog.ReleasedEvents() != 2 {
t.Errorf("ReleasedEvents = %d, want 2", prog.ReleasedEvents())
}
})
t.Run("partial_sink_progress_two_succeed_then_fail", func(t *testing.T) {
releaser, sink, _ := setupReleaser(t, "ch-partial2", 2, 100)
releaser.boundary.BeginAttempt("a1")
rs, _ := NewResponseStart("ch-partial2", 200, nil, releaserNow)
releaser.boundary.StageResponseStart("a1", rs)
epoch := appendDataEvents(t, releaser.tail, "ch-partial2", 5)
// 2 successes, then fail.
sink.failAfterSuccesses = 2
prog, err := releaser.ReleaseEpoch(ctx, "a1", epoch.ID())
if err == nil {
t.Error("expected partial error")
}
if prog.ReleasedEvents() != 2 {
t.Errorf("ReleasedEvents = %d, want 2", prog.ReleasedEvents())
}
})
t.Run("full_sink_progress", func(t *testing.T) {
releaser, sink, _ := setupReleaser(t, "ch-full", 2, 100)
releaser.boundary.BeginAttempt("a1")
rs, _ := NewResponseStart("ch-full", 200, nil, releaserNow)
releaser.boundary.StageResponseStart("a1", rs)
epoch := appendDataEvents(t, releaser.tail, "ch-full", 3)
prog, err := releaser.ReleaseEpoch(ctx, "a1", epoch.ID())
if err != nil {
t.Fatalf("ReleaseEpoch: %v", err)
}
if prog.ReleasedEvents() != 3 {
t.Errorf("ReleasedEvents = %d, want 3", prog.ReleasedEvents())
}
// Cursor should advance by 3.
if got := releaser.tail.CommittedCursor("ch-full"); got != 3 {
t.Errorf("committed cursor = %d, want 3", got)
}
// Sink should have received 3 releases + 1 start.
if len(sink.releases) != 3 {
t.Errorf("sink releases = %d, want 3", len(sink.releases))
}
if len(sink.responseStarts) != 1 {
t.Errorf("sink responseStarts = %d, want 1", len(sink.responseStarts))
}
})
}
// TestStreamReleaserTerminalBeforeEvidenceExposesOneOutcome verifies that the
// pass path exposes start/output+success terminal and error path exposes
// only error terminal without releasing pending/start.
func TestStreamReleaserTerminalBeforeEvidenceExposesOneOutcome(t *testing.T) {
ctx := context.Background()
t.Run("pass_exposes_start_and_terminal", func(t *testing.T) {
releaser, sink, _ := setupReleaser(t, "ch-pass", 2, 100)
releaser.boundary.BeginAttempt("a1")
rs, _ := NewResponseStart("ch-pass", 200, nil, releaserNow)
releaser.boundary.StageResponseStart("a1", rs)
// Create and release an epoch.
epoch := appendDataEvents(t, releaser.tail, "ch-pass", 2)
prog, err := releaser.ReleaseEpoch(ctx, "a1", epoch.ID())
if err != nil {
t.Fatalf("ReleaseEpoch: %v", err)
}
if prog.ReleasedEvents() != 2 {
t.Errorf("ReleasedEvents = %d, want 2", prog.ReleasedEvents())
}
// Terminal success.
tr := createSuccessTerminalResult("ch-pass")
prog2, err := releaser.ReleaseTerminalEpoch(ctx, "a1", epoch.ID(), tr)
if err != nil {
t.Fatalf("ReleaseTerminalEpoch: %v", err)
}
_ = prog2
if len(sink.responseStarts) != 1 {
t.Errorf("responseStarts = %d, want 1", len(sink.responseStarts))
}
if len(sink.terminals) != 1 {
t.Errorf("terminals = %d, want 1", len(sink.terminals))
}
if len(sink.releases) != 2 {
t.Errorf("releases = %d, want 2", len(sink.releases))
}
})
t.Run("error_exposes_only_error_terminal", func(t *testing.T) {
releaser, sink, canceller := setupReleaser(t, "ch-err", 2, 100)
releaser.boundary.BeginAttempt("a2")
rs, _ := NewResponseStart("ch-err", 200, nil, releaserNow)
releaser.boundary.StageResponseStart("a2", rs)
// Create some pending events but don't release them.
appendDataEvents(t, releaser.tail, "ch-err", 2)
// Error terminal without releasing events.
tr := createErrorTerminalResult("ch-err")
prog, err := releaser.ReleaseTerminalEpoch(ctx, "a2", 999, tr)
if err != nil {
t.Fatalf("ReleaseTerminalEpoch: %v", err)
}
_ = prog
// Should have called cancel.
if canceller.cancelCount != 1 {
t.Errorf("cancelCount = %d, want 1", canceller.cancelCount)
}
if len(sink.responseStarts) != 0 {
t.Errorf("responseStarts = %d, want 0 (error path)", len(sink.responseStarts))
}
if len(sink.terminals) != 1 {
t.Errorf("terminals = %d, want 1", len(sink.terminals))
}
if len(sink.releases) != 0 {
t.Errorf("releases = %d, want 0 (no pending release)", len(sink.releases))
}
})
}
// TestStreamReleaserDoesNotDuplicatePreparedRelease verifies that repeated
// epoch calls do not duplicate downstream events.
func TestStreamReleaserDoesNotDuplicatePreparedRelease(t *testing.T) {
ctx := context.Background()
releaser, sink, _ := setupReleaser(t, "ch-dup", 2, 100)
releaser.boundary.BeginAttempt("a1")
rs, _ := NewResponseStart("ch-dup", 200, nil, releaserNow)
releaser.boundary.StageResponseStart("a1", rs)
// First epoch.
epoch1 := appendDataEvents(t, releaser.tail, "ch-dup", 2)
prog, err := releaser.ReleaseEpoch(ctx, "a1", epoch1.ID())
if err != nil {
t.Fatalf("ReleaseEpoch 1: %v", err)
}
if prog.ReleasedEvents() != 2 {
t.Errorf("first ReleasedEvents = %d, want 2", prog.ReleasedEvents())
}
// Second epoch (more events to trigger another).
epoch2 := appendDataEvents(t, releaser.tail, "ch-dup", 2)
prog2, err := releaser.ReleaseEpoch(ctx, "a1", epoch2.ID())
if err != nil {
t.Fatalf("ReleaseEpoch 2: %v", err)
}
if prog2.ReleasedEvents() != 2 {
t.Errorf("second ReleasedEvents = %d, want 2", prog2.ReleasedEvents())
}
// Total should be 4 releases (2+2), no duplicates.
if len(sink.releases) != 4 {
t.Errorf("total releases = %d, want 4 (2+2)", len(sink.releases))
}
}
// ---------------------------------------------------------------------------
// REVIEW_API-3: Attempt replacement, idle and overflow
// ---------------------------------------------------------------------------
// TestStreamReleaserReplacementResetsTailGeneration verifies that boundary
// attempt and tail pending/prepared token are replaced together.
func TestStreamReleaserReplacementResetsTailGeneration(t *testing.T) {
ctx := context.Background()
releaser, _, _ := setupReleaser(t, "ch-rep", 2, 100)
// Begin attempt and stage start.
releaser.boundary.BeginAttempt("old")
rs, _ := NewResponseStart("ch-rep", 200, nil, releaserNow)
releaser.boundary.StageResponseStart("old", rs)
// Append events.
appendDataEvents(t, releaser.tail, "ch-rep", 2)
// Replace.
err := releaser.ReplaceUncommittedAttempt(ctx, "old", "new")
if err != nil {
t.Fatalf("ReplaceUncommittedAttempt: %v", err)
}
// Tail cursor should be 0 after replace.
if got := releaser.tail.CommittedCursor("ch-rep"); got != 0 {
t.Errorf("tail cursor after replace = %d, want 0", got)
}
// Boundary attempt should be new attempt.
if got := releaser.boundary.CurrentAttempt(); got != "new" {
t.Errorf("boundary attempt after replace = %q, want %q", got, "new")
}
}
// TestStreamReleaserIdleReleasesNoPending verifies that idle error results
// in zero release/start and correct cancel→terminal ordering.
func TestStreamReleaserIdleReleasesNoPending(t *testing.T) {
ctx := context.Background()
releaser, sink, canceller := setupReleaser(t, "ch-idle", 2, 100)
releaser.boundary.BeginAttempt("a1")
rs, _ := NewResponseStart("ch-idle", 200, nil, releaserNow)
releaser.boundary.StageResponseStart("a1", rs)
// Append some events to create pending state.
appendDataEvents(t, releaser.tail, "ch-idle", 2)
// FailPending.
tr := createErrorTerminalResult("ch-idle")
prog, err := releaser.FailPending(ctx, "a1", tr)
if err != nil {
t.Fatalf("FailPending: %v", err)
}
if prog.ReleasedEvents() != 0 {
t.Errorf("ReleasedEvents = %d, want 0", prog.ReleasedEvents())
}
// No responses or releases.
if len(sink.responseStarts) != 0 {
t.Errorf("responseStarts = %d, want 0", len(sink.responseStarts))
}
if len(sink.releases) != 0 {
t.Errorf("releases = %d, want 0", len(sink.releases))
}
// Cancel should have been called.
if canceller.cancelCount != 1 {
t.Errorf("cancelCount = %d, want 1", canceller.cancelCount)
}
if len(sink.terminals) != 1 {
t.Errorf("terminals = %d, want 1", len(sink.terminals))
}
}
// TestStreamReleaserOverflowCancelsBeforeSingleTerminal verifies that S19
// overflow results in no partial release, and terminal exactly once.
func TestStreamReleaserOverflowCancelsBeforeSingleTerminal(t *testing.T) {
ctx := context.Background()
// Use very small buffer to trigger overflow.
// rolling_window with evidence 1 and maxBuffer 10 (minimum).
releaser, sink, canceller := setupReleaser(t, "ch-overflow", 1, 10)
_ = releaser.boundary.BeginAttempt("a1")
// Append events, expecting overflow signal.
overflowed := false
for i := 0; i < 50; i++ {
ne, err := NewTextDeltaEvent("ch-overflow", "chunk", releaserNow)
if err != nil {
t.Fatalf("NewTextDeltaEvent: %v", err)
}
_, signal, err := releaser.tail.Append(ne)
if err != nil {
t.Fatalf("Append: %v", err)
}
if signal == EvidenceTailSignalBufferOverflow {
overflowed = true
break
}
}
if !overflowed {
t.Fatal("expected overflow signal not generated with current configuration")
}
// FailPending for overflow recovery.
tr := createErrorTerminalResult("ch-overflow")
prog, err := releaser.FailPending(ctx, "a1", tr)
if err != nil {
t.Fatalf("FailPending: %v", err)
}
if prog.ReleasedEvents() != 0 {
t.Errorf("ReleasedEvents = %d, want 0 (overflow)", prog.ReleasedEvents())
}
if len(sink.terminals) != 1 {
t.Errorf("terminals = %d, want 1 (exactly once)", len(sink.terminals))
}
_ = canceller
}
func TestStreamReleaserConfirmsProgressEvenWhenSinkFails(t *testing.T) {
ctx := context.Background()
t.Run("partial_release_advances_cursor_and_allows_reprepare", func(t *testing.T) {
releaser, sink, _ := setupReleaser(t, "ch-fail-confirm", 2, 100)
_ = releaser.boundary.BeginAttempt("a1")
rs, _ := NewResponseStart("ch-fail-confirm", 200, nil, releaserNow)
_ = releaser.boundary.StageResponseStart("a1", rs)
distinctTexts := []string{"alpha", "beta", "gamma", "delta", "epsilon"}
epoch := appendDataEventsDistinct(t, releaser.tail, "ch-fail-confirm", 5, func(i int) string { return distinctTexts[i] })
sink.failAfterSuccesses = 2
prog, err := releaser.ReleaseEpoch(ctx, "a1", epoch.ID())
if err == nil {
t.Fatal("expected release error")
}
if prog.ReleasedEvents() != 2 {
t.Fatalf("ReleasedEvents = %d, want 2", prog.ReleasedEvents())
}
if releaser.tail.CommittedCursor("ch-fail-confirm") != 2 {
t.Fatalf("cursor = %d, want 2", releaser.tail.CommittedCursor("ch-fail-confirm"))
}
// Re-prepare the same epoch: snapStartSeq was advanced by ConfirmRelease,
// so the snapshot should contain only the remaining 3 events (indices 2-4).
prepared2, err := releaser.tail.PrepareRelease(epoch.ID())
if err != nil {
t.Fatalf("PrepareRelease remaining suffix failed: %v", err)
}
// Exact suffix verification: exactly 3 remaining events with expected distinct text.
wantSuffix := []string{"gamma", "delta", "epsilon"}
gotEvents := prepared2.ReleaseEvents()
if diff := compareStrings(releaseEventTexts(t, gotEvents), wantSuffix); diff != "" {
t.Fatalf("remaining suffix text mismatch: %s", diff)
}
if len(gotEvents) != 3 {
t.Fatalf("remaining suffix count = %d, want 3", len(gotEvents))
}
// Token must be non-empty and different from the first prepare token.
if prepared2.Token() == "" {
t.Fatal("prepared token is empty")
}
// Cursor must still be 2 (unchanged by re-prepare).
if got := releaser.tail.CommittedCursor("ch-fail-confirm"); got != 2 {
t.Errorf("cursor after re-prepare = %d, want 2", got)
}
})
t.Run("token_changes_after_reprepare", func(t *testing.T) {
releaser, _, _ := setupReleaser(t, "ch-token", 2, 100)
_ = releaser.boundary.BeginAttempt("a1")
rs, _ := NewResponseStart("ch-token", 200, nil, releaserNow)
_ = releaser.boundary.StageResponseStart("a1", rs)
epoch := appendDataEvents(t, releaser.tail, "ch-token", 3)
// First prepare.
prepared1, err := releaser.tail.PrepareRelease(epoch.ID())
if err != nil {
t.Fatalf("first PrepareRelease: %v", err)
}
firstToken := prepared1.Token()
// Confirm zero events.
err = releaser.tail.ConfirmRelease(firstToken, ReleaseConfirmation{ReleasedEvents: 0})
if err != nil {
t.Fatalf("ConfirmRelease zero: %v", err)
}
// Re-prepare: must get a new token.
prepared2, err := releaser.tail.PrepareRelease(epoch.ID())
if err != nil {
t.Fatalf("second PrepareRelease: %v", err)
}
if prepared2.Token() == firstToken {
t.Fatal("reprepare reused stale token")
}
// Cursor must be 0 (zero confirm).
if got := releaser.tail.CommittedCursor("ch-token"); got != 0 {
t.Errorf("cursor after zero confirm = %d, want 0", got)
}
})
t.Run("zero_partial_full_confirm_matrix", func(t *testing.T) {
releaser, _, _ := setupReleaser(t, "ch-matrix", 2, 100)
_ = releaser.boundary.BeginAttempt("a1")
rs, _ := NewResponseStart("ch-matrix", 200, nil, releaserNow)
_ = releaser.boundary.StageResponseStart("a1", rs)
distinctTexts := []string{"one", "two", "three"}
epoch := appendDataEventsDistinct(t, releaser.tail, "ch-matrix", 3, func(i int) string { return distinctTexts[i] })
// --- Zero confirm ---
prepared, err := releaser.tail.PrepareRelease(epoch.ID())
if err != nil {
t.Fatalf("PrepareRelease: %v", err)
}
firstToken := prepared.Token()
// Verify initial snapshot payload is exactly the three distinct events.
if diff := compareStrings(releaseEventTexts(t, prepared.ReleaseEvents()), distinctTexts); diff != "" {
t.Fatalf("initial snapshot text mismatch: %s", diff)
}
if len(prepared.ReleaseEvents()) != 3 {
t.Fatalf("initial snapshot count = %d, want 3", len(prepared.ReleaseEvents()))
}
err = releaser.tail.ConfirmRelease(firstToken, ReleaseConfirmation{ReleasedEvents: 0})
if err != nil {
t.Fatalf("ConfirmRelease zero: %v", err)
}
if releaser.tail.CommittedCursor("ch-matrix") != 0 {
t.Errorf("cursor after zero confirm = %d, want 0", releaser.tail.CommittedCursor("ch-matrix"))
}
// --- Partial confirm ---
prepared2, err := releaser.tail.PrepareRelease(epoch.ID())
if err != nil {
t.Fatalf("re-prepare after zero: %v", err)
}
if prepared2.Token() == firstToken {
t.Fatal("re-prepare after zero confirm reused token")
}
// Partial snapshot must contain remaining 3 events (cursor still 0).
if diff := compareStrings(releaseEventTexts(t, prepared2.ReleaseEvents()), distinctTexts); diff != "" {
t.Fatalf("partial snapshot text mismatch: %s", diff)
}
err = releaser.tail.ConfirmRelease(prepared2.Token(), ReleaseConfirmation{ReleasedEvents: 1})
if err != nil {
t.Fatalf("ConfirmRelease partial: %v", err)
}
if releaser.tail.CommittedCursor("ch-matrix") != 1 {
t.Errorf("cursor after partial confirm = %d, want 1", releaser.tail.CommittedCursor("ch-matrix"))
}
// --- Full confirm ---
prepared3, err := releaser.tail.PrepareRelease(epoch.ID())
if err != nil {
t.Fatalf("re-prepare after partial: %v", err)
}
if prepared3.Token() == prepared2.Token() {
t.Fatal("re-prepare after partial reused token")
}
// Full snapshot must contain remaining 2 events (cursor is 1).
wantFullSnapshot := []string{"two", "three"}
if diff := compareStrings(releaseEventTexts(t, prepared3.ReleaseEvents()), wantFullSnapshot); diff != "" {
t.Fatalf("full snapshot text mismatch: %s", diff)
}
err = releaser.tail.ConfirmRelease(prepared3.Token(), ReleaseConfirmation{ReleasedEvents: 2})
if err != nil {
t.Fatalf("ConfirmRelease full: %v", err)
}
if releaser.tail.CommittedCursor("ch-matrix") != 3 {
t.Errorf("cursor after full confirm = %d, want 3", releaser.tail.CommittedCursor("ch-matrix"))
}
// --- Consumed epoch rejects re-prepare ---
_, err = releaser.tail.PrepareRelease(epoch.ID())
if err == nil {
t.Fatal("expected error re-preparing consumed epoch")
}
})
}
func TestStreamReleaserSuppressesSuccessTerminalUntilConfirmation(t *testing.T) {
ctx := context.Background()
t.Run("unknown_epoch_suppresses_terminal", func(t *testing.T) {
releaser, sink, _ := setupReleaser(t, "ch-suppress-unknown", 2, 100)
_ = releaser.boundary.BeginAttempt("a1")
tr := createSuccessTerminalResult("ch-suppress-unknown")
_, err := releaser.ReleaseTerminalEpoch(ctx, "a1", 9999, tr)
if err == nil {
t.Fatal("expected error on unknown epoch")
}
if len(sink.terminals) != 0 {
t.Fatalf("terminals = %d, want 0", len(sink.terminals))
}
})
t.Run("sink_failure_suppresses_terminal", func(t *testing.T) {
releaser, sink, _ := setupReleaser(t, "ch-suppress-fail", 2, 100)
_ = releaser.boundary.BeginAttempt("a1")
rs, _ := NewResponseStart("ch-suppress-fail", 200, nil, releaserNow)
_ = releaser.boundary.StageResponseStart("a1", rs)
epoch := appendDataEvents(t, releaser.tail, "ch-suppress-fail", 5)
sink.failAfterSuccesses = 2
tr := createSuccessTerminalResult("ch-suppress-fail")
_, err := releaser.ReleaseTerminalEpoch(ctx, "a1", epoch.ID(), tr)
if err == nil {
t.Fatal("expected error on ReleaseTerminalEpoch when release fails")
}
// Error must be a release failure (sink error), not a confirm failure.
if !strings.Contains(err.Error(), "forced release failure") {
t.Fatalf("error = %v, want release failure error", err)
}
if len(sink.terminals) != 0 {
t.Fatalf("terminals = %d, want 0 (success terminal suppressed)", len(sink.terminals))
}
})
t.Run("clean_full_confirmation_commits_terminal", func(t *testing.T) {
releaser, sink, _ := setupReleaser(t, "ch-clean-terminal", 2, 100)
_ = releaser.boundary.BeginAttempt("a1")
rs, _ := NewResponseStart("ch-clean-terminal", 200, nil, releaserNow)
_ = releaser.boundary.StageResponseStart("a1", rs)
epoch := appendDataEvents(t, releaser.tail, "ch-clean-terminal", 3)
tr := createSuccessTerminalResult("ch-clean-terminal")
prog, err := releaser.ReleaseTerminalEpoch(ctx, "a1", epoch.ID(), tr)
if err != nil {
t.Fatalf("clean ReleaseTerminalEpoch failed: %v", err)
}
if prog.ReleasedEvents() != 3 {
t.Fatalf("ReleasedEvents = %d, want 3", prog.ReleasedEvents())
}
if len(sink.terminals) != 1 {
t.Fatalf("terminals = %d, want 1", len(sink.terminals))
}
})
t.Run("outstanding_prepare_rejects_confirm", func(t *testing.T) {
releaser, _, _ := setupReleaser(t, "ch-outstanding", 2, 100)
_ = releaser.boundary.BeginAttempt("a1")
rs, _ := NewResponseStart("ch-outstanding", 200, nil, releaserNow)
_ = releaser.boundary.StageResponseStart("a1", rs)
epoch := appendDataEvents(t, releaser.tail, "ch-outstanding", 3)
// Prepare but do not confirm.
prepared, err := releaser.tail.PrepareRelease(epoch.ID())
if err != nil {
t.Fatalf("PrepareRelease: %v", err)
}
// Second prepare on same channel should fail (outstanding token).
_, err = releaser.tail.PrepareRelease(epoch.ID())
if err == nil {
t.Fatal("expected error on duplicate prepare with outstanding token")
}
// Confirm partial (2 out of 3) to keep epoch unconfirmed.
err = releaser.tail.ConfirmRelease(prepared.Token(), ReleaseConfirmation{ReleasedEvents: 2})
if err != nil {
t.Fatalf("ConfirmRelease partial: %v", err)
}
// Re-prepare should work with a new token.
prepared2, err := releaser.tail.PrepareRelease(epoch.ID())
if err != nil {
t.Fatalf("re-prepare after partial confirm: %v", err)
}
if prepared2.Token() == prepared.Token() {
t.Fatal("re-prepare reused same token after partial confirm")
}
})
t.Run("invalidated_epoch_rejects_confirm", func(t *testing.T) {
releaser, _, _ := setupReleaser(t, "ch-invalidated", 2, 100)
_ = releaser.boundary.BeginAttempt("a1")
rs, _ := NewResponseStart("ch-invalidated", 200, nil, releaserNow)
_ = releaser.boundary.StageResponseStart("a1", rs)
epoch := appendDataEvents(t, releaser.tail, "ch-invalidated", 3)
// Prepare.
prepared, err := releaser.tail.PrepareRelease(epoch.ID())
if err != nil {
t.Fatalf("PrepareRelease: %v", err)
}
// Invalidate via DiscardPendingForTerminal.
releaser.tail.DiscardPendingForTerminal()
// Confirm with invalidated token must fail.
err = releaser.tail.ConfirmRelease(prepared.Token(), ReleaseConfirmation{ReleasedEvents: 3})
if err == nil {
t.Fatal("expected error confirming invalidated token")
}
})
t.Run("confirm_failure_prevents_terminal", func(t *testing.T) {
releaser, sink, _ := setupReleaser(t, "ch-confirm-fail", 2, 100)
_ = releaser.boundary.BeginAttempt("a1")
rs, _ := NewResponseStart("ch-confirm-fail", 200, nil, releaserNow)
_ = releaser.boundary.StageResponseStart("a1", rs)
epoch := appendDataEvents(t, releaser.tail, "ch-confirm-fail", 3)
// Program sink to fail on release to force partial confirm.
sink.failAfterSuccesses = 1
tr := createSuccessTerminalResult("ch-confirm-fail")
_, err := releaser.ReleaseTerminalEpoch(ctx, "a1", epoch.ID(), tr)
if err == nil {
t.Fatal("expected error when confirm fails")
}
// Terminal must not have been called.
if len(sink.terminals) != 0 {
t.Fatalf("terminals = %d, want 0 (confirm failure suppresses terminal)", len(sink.terminals))
}
})
t.Run("actual_confirm_invalidation_prevents_terminal", func(t *testing.T) {
// Verify that an actual token invalidation during release callback
// (via sink.onRelease calling DiscardPendingForTerminal) causes
// ConfirmRelease to reject, which suppresses the success terminal.
releaser, sink, _ := setupReleaser(t, "ch-actual-invalidate", 2, 100)
_ = releaser.boundary.BeginAttempt("a1")
rs, _ := NewResponseStart("ch-actual-invalidate", 200, nil, releaserNow)
_ = releaser.boundary.StageResponseStart("a1", rs)
epoch := appendDataEvents(t, releaser.tail, "ch-actual-invalidate", 3)
// Program sink to invalidate the prepared token during release callback.
var invalidationOccurred bool
sink.onRelease = func(ev ReleaseEvent) {
if !invalidationOccurred {
invalidationOccurred = true
releaser.tail.DiscardPendingForTerminal()
}
}
// ReleaseTerminalEpoch: releaseAndConfirm will PrepareRelease (creating token),
// then during sink.Release the token is invalidated via DiscardPendingForTerminal.
// ConfirmRelease with the invalidated token must fail, suppressing the terminal.
tr := createSuccessTerminalResult("ch-actual-invalidate")
_, err := releaser.ReleaseTerminalEpoch(ctx, "a1", epoch.ID(), tr)
if err == nil {
t.Fatal("expected error on ReleaseTerminalEpoch with token invalidated during release")
}
// Error must be a confirm-related error (stale or not-found token).
if !strings.Contains(err.Error(), "token") {
t.Fatalf("error = %v, want token-related error", err)
}
if invalidationOccurred != true {
t.Fatal("expected sink.onRelease to trigger token invalidation")
}
if len(sink.terminals) != 0 {
t.Fatalf("terminals = %d, want 0 (invalidated epoch suppresses terminal)", len(sink.terminals))
}
})
t.Run("terminal_suppression_matrix", func(t *testing.T) {
// Table-driven comparison: terminal count must be 0 for all
// failure paths and 1 for the clean completed epoch path.
type scenario struct {
name string
wantTerm int
run func(t *testing.T, releaser *StreamReleaser, sink *testSink) error
}
run := func(sc scenario) {
t.Run(sc.name, func(t *testing.T) {
releaser, sink, _ := setupReleaser(t, "ch-mx-"+sc.name, 2, 100)
_ = releaser.boundary.BeginAttempt("a1")
rs, _ := NewResponseStart("ch-mx-"+sc.name, 200, nil, releaserNow)
_ = releaser.boundary.StageResponseStart("a1", rs)
err := sc.run(t, releaser, sink)
if sc.wantTerm == 1 {
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
}
if len(sink.terminals) != sc.wantTerm {
t.Errorf("terminals = %d, want %d", len(sink.terminals), sc.wantTerm)
}
})
}
// 1. Outstanding prepare (unconfirmed epoch): terminal suppressed.
run(scenario{
name: "outstanding-prepare",
wantTerm: 0,
run: func(t *testing.T, releaser *StreamReleaser, sink *testSink) error {
ch := "ch-mx-outstanding-prepare"
epoch := appendDataEvents(t, releaser.tail, ch, 3)
if epoch.ID() == 0 {
t.Fatalf("expected non-zero epoch ID")
}
_, err := releaser.tail.PrepareRelease(epoch.ID())
if err != nil {
return err
}
tr := createSuccessTerminalResult(ch)
_, err = releaser.ReleaseTerminalEpoch(context.Background(), "a1", epoch.ID(), tr)
return err
},
})
// 2. Invalidated epoch: terminal suppressed.
run(scenario{
name: "invalidated-epoch",
wantTerm: 0,
run: func(t *testing.T, releaser *StreamReleaser, sink *testSink) error {
ch := "ch-mx-invalidated-epoch"
epoch := appendDataEvents(t, releaser.tail, ch, 3)
if epoch.ID() == 0 {
t.Fatalf("expected non-zero epoch ID")
}
_, err := releaser.tail.PrepareRelease(epoch.ID())
if err != nil {
return err
}
releaser.tail.DiscardPendingForTerminal()
tr := createSuccessTerminalResult(ch)
_, err = releaser.ReleaseTerminalEpoch(context.Background(), "a1", epoch.ID(), tr)
return err
},
})
// 3. Sink release failure: terminal suppressed.
run(scenario{
name: "sink-release-failure",
wantTerm: 0,
run: func(t *testing.T, releaser *StreamReleaser, sink *testSink) error {
ch := "ch-mx-sink-release-failure"
epoch := appendDataEvents(t, releaser.tail, ch, 3)
if epoch.ID() == 0 {
t.Fatalf("expected non-zero epoch ID")
}
sink.failAfterSuccesses = 1
tr := createSuccessTerminalResult(ch)
_, err := releaser.ReleaseTerminalEpoch(context.Background(), "a1", epoch.ID(), tr)
return err
},
})
// 4. Actual confirm invalidation during release: terminal suppressed.
run(scenario{
name: "actual-confirm-invalidation",
wantTerm: 0,
run: func(t *testing.T, releaser *StreamReleaser, sink *testSink) error {
ch := "ch-mx-actual-confirm-invalidation"
epoch := appendDataEvents(t, releaser.tail, ch, 3)
if epoch.ID() == 0 {
t.Fatalf("expected non-zero epoch ID")
}
var invalidationOccurred bool
sink.onRelease = func(ev ReleaseEvent) {
if !invalidationOccurred {
invalidationOccurred = true
releaser.tail.DiscardPendingForTerminal()
}
}
tr := createSuccessTerminalResult(ch)
_, err := releaser.ReleaseTerminalEpoch(context.Background(), "a1", epoch.ID(), tr)
if !strings.Contains(err.Error(), "token") {
t.Fatalf("error = %v, want token-related error", err)
}
return err
},
})
// 5. Clean completed epoch: terminal committed exactly once.
run(scenario{
name: "clean-completed-epoch",
wantTerm: 1,
run: func(t *testing.T, releaser *StreamReleaser, sink *testSink) error {
ch := "ch-mx-clean-completed-epoch"
epoch := appendDataEvents(t, releaser.tail, ch, 3)
if epoch.ID() == 0 {
t.Fatalf("expected non-zero epoch ID")
}
_, err := releaser.ReleaseEpoch(context.Background(), "a1", epoch.ID())
if err != nil {
return err
}
tr := createSuccessTerminalResult(ch)
_, err = releaser.ReleaseTerminalEpoch(context.Background(), "a1", epoch.ID(), tr)
return err
},
})
})
}
func TestStreamReleaserRejectsReentrantCallbackWithoutDeadlock(t *testing.T) {
ctx := context.Background()
releaser, sink, _ := setupReleaser(t, "ch-reentrant", 2, 100)
_ = releaser.boundary.BeginAttempt("a1")
rs, _ := NewResponseStart("ch-reentrant", 200, nil, releaserNow)
_ = releaser.boundary.StageResponseStart("a1", rs)
epoch := appendDataEvents(t, releaser.tail, "ch-reentrant", 2)
var outerReleaseCount, reentrantReleaseCount int
reentrantErrors := make([]error, 0)
var reentrantErrorsMu sync.Mutex
sink.onRelease = func(ev ReleaseEvent) {
// Count outer release callbacks.
outerReleaseCount++
// Attempt re-entrant call inside sink callback.
_, err := releaser.ReleaseEpoch(ctx, "a1", epoch.ID())
reentrantReleaseCount++
reentrantErrorsMu.Lock()
reentrantErrors = append(reentrantErrors, err)
reentrantErrorsMu.Unlock()
}
outerResult := make(chan error, 1)
go func() {
_, err := releaser.ReleaseEpoch(ctx, "a1", epoch.ID())
outerResult <- err
}()
// Reentrant must be rejected with serializer conflict.
select {
case err := <-outerResult:
if err != nil {
t.Fatalf("outer ReleaseEpoch failed: %v", err)
}
case <-time.After(time.Second):
t.Fatal("outer ReleaseEpoch deadlocked")
}
// Verify exact callback counts.
// Outer should have exactly 2 release callbacks (2 events in epoch).
if outerReleaseCount != 2 {
t.Errorf("outer release callbacks = %d, want 2", outerReleaseCount)
}
// Reentrant should have exactly 2 attempted calls (one per outer release callback),
// each rejected with ErrReleaserSerializerConflict.
if reentrantReleaseCount != 2 {
t.Fatalf("reentrant callbacks = %d, want 2 (one per outer release callback)", reentrantReleaseCount)
}
reentrantErrorsMu.Lock()
defer reentrantErrorsMu.Unlock()
for i, err := range reentrantErrors {
if !errors.Is(err, ErrReleaserSerializerConflict) {
t.Errorf("reentrant error[%d] = %v, want ErrReleaserSerializerConflict", i, err)
}
}
}
func TestStreamReleaserFailPendingOrdersDiscardCancelTerminal(t *testing.T) {
ctx := context.Background()
releaser, sink, canceller := setupReleaser(t, "ch-order", 2, 100)
_ = releaser.boundary.BeginAttempt("a1")
rs, _ := NewResponseStart("ch-order", 200, nil, releaserNow)
_ = releaser.boundary.StageResponseStart("a1", rs)
epoch := appendDataEvents(t, releaser.tail, "ch-order", 2)
// Verify pending state before FailPending.
if got := releaser.tail.CommittedCursor("ch-order"); got != 0 {
t.Fatalf("cursor before fail = %d, want 0", got)
}
// Prepare a release to create a prepared token.
prepared, err := releaser.tail.PrepareRelease(epoch.ID())
if err != nil {
t.Fatalf("PrepareRelease: %v", err)
}
prepareToken := prepared.Token()
if prepareToken == "" {
t.Fatal("prepared token is empty")
}
var trace []string
var observed snapshotResult
canceller.onCancel = func() {
// Snapshot discard state at cancel time to verify ordering:
// discard must have been applied before cancel runs.
observed = snapshotDiscardState(releaser.tail, "ch-order", epoch.ID(), prepareToken)
trace = append(trace, "discard-verified", "cancel")
}
sink.onTerminal = func(tr TerminalResult) {
trace = append(trace, "terminal")
}
tr := createErrorTerminalResult("ch-order")
prog, err := releaser.FailPending(ctx, "a1", tr)
if err != nil {
t.Fatalf("FailPending: %v", err)
}
if prog.ReleasedEvents() != 0 {
t.Fatalf("ReleasedEvents = %d, want 0", prog.ReleasedEvents())
}
// Exact 3-step trace: discard-verified → cancel → terminal.
if len(trace) != 3 || trace[0] != "discard-verified" || trace[1] != "cancel" || trace[2] != "terminal" {
t.Fatalf("trace = %v, want [discard-verified cancel terminal]", trace)
}
// Hook-time snapshot must confirm discard was applied before cancel.
if observed.pendingEntries != 0 {
t.Errorf("hook pending entries = %d, want 0 (discarded)", observed.pendingEntries)
}
if observed.pendingRunes != 0 {
t.Errorf("hook pending runes = %d, want 0 (discarded)", observed.pendingRunes)
}
if observed.hasPreparedToken {
t.Error("hook observed prepared token still set (should be cleared by discard)")
}
if observed.staleTokenConfirm == nil {
t.Error("expected stale token confirm to fail after discard")
}
if observed.reprepareErr == nil {
t.Error("expected epoch reprepare to fail after discard")
}
// The original prepare token must no longer be valid after discard.
err = releaser.tail.ConfirmRelease(prepareToken, ReleaseConfirmation{ReleasedEvents: 2})
if err == nil {
t.Error("expected error confirming invalidated token after discard")
}
if canceller.cancelCount != 1 {
t.Errorf("cancelCount = %d, want 1", canceller.cancelCount)
}
if len(sink.terminals) != 1 {
t.Errorf("terminals = %d, want 1", len(sink.terminals))
}
// Cursor should be preserved after discard.
if got := releaser.tail.CommittedCursor("ch-order"); got != 0 {
t.Errorf("cursor after fail = %d, want 0 (preserved)", got)
}
// After discard, no epoch can be prepared (all pending/prepared invalidated).
_, err = releaser.tail.PrepareRelease(epoch.ID())
if err == nil {
t.Error("expected error preparing invalidated epoch after discard")
}
}
// discardSnapshot captures the tail state relevant to verifying that discard
// has been applied at the moment the canceller hook fires. This is test-only
// and accesses private EvidenceTail fields directly.
type discardSnapshot struct {
pendingEntries int
pendingRunes int
hasPreparedToken bool
preparedToken string
staleTokenConfirm error
reprepareErr error
}
// snapshotDiscardState captures the discard-relevant state of the given
// channel and epoch at the current moment. It is used inside the canceller
// hook to verify that DiscardPendingForTerminal has already been applied
// before CancelAttempt runs.
func snapshotDiscardState(tail *EvidenceTail, channel string, epochID uint64, staleToken string) discardSnapshot {
cs, ok := tail.channelState[channel]
if !ok {
return discardSnapshot{}
}
snap := discardSnapshot{
pendingEntries: len(cs.pendingEntries),
pendingRunes: cs.pendingRunes,
}
// Check if any epoch record still has a prepared token.
if rec, ok := tail.epochs[epochID]; ok {
snap.hasPreparedToken = rec.prepared && rec.token != ""
snap.preparedToken = rec.token
}
// Try confirming the stale token to verify it is rejected.
snap.staleTokenConfirm = tail.ConfirmRelease(staleToken, ReleaseConfirmation{ReleasedEvents: 2})
// Try re-preparing the epoch to verify it is rejected.
_, snap.reprepareErr = tail.PrepareRelease(epochID)
return snap
}
// snapshotResult is an alias kept for the test assertion field names.
type snapshotResult = discardSnapshot
// releaseEventTexts extracts text content from release events for comparison.
func releaseEventTexts(t *testing.T, events []ReleaseEvent) []string {
t.Helper()
texts := make([]string, len(events))
for i, ev := range events {
texts[i] = releaseEventText(t, ev)
}
return texts
}
// releaseEventText extracts the text content from a single release event.
func releaseEventText(t *testing.T, ev ReleaseEvent) string {
t.Helper()
switch ev.Kind() {
case EventKindTextDelta:
text, err := ev.AsTextDelta()
if err != nil {
t.Fatalf("AsTextDelta: %v", err)
}
return text
case EventKindReasoningDelta:
text, err := ev.AsReasoningDelta()
if err != nil {
t.Fatalf("AsReasoningDelta: %v", err)
}
return text
default:
t.Fatalf("unsupported release event kind: %s", ev.Kind())
return ""
}
}
// compareStrings returns a diff string if a and b differ, or "" if equal.
func compareStrings(a, b []string) string {
if len(a) != len(b) {
return fmt.Sprintf("length mismatch: got %d, want %d", len(a), len(b))
}
for i := range a {
if a[i] != b[i] {
return fmt.Sprintf("index %d: got %q, want %q", i, a[i], b[i])
}
}
return ""
}
// TestStreamReleaserCancelFailureAppendsBoundedCause verifies that when
// cancel fails, the cause is appended to the terminal result.
func TestStreamReleaserCancelFailureAppendsBoundedCause(t *testing.T) {
ctx := context.Background()
releaser, sink, canceller := setupReleaser(t, "ch-cancel-fail", 2, 100)
releaser.boundary.BeginAttempt("a1")
rs, _ := NewResponseStart("ch-cancel-fail", 200, nil, releaserNow)
releaser.boundary.StageResponseStart("a1", rs)
// Program canceller to fail.
canceller.cancelErr = errors.New("test: cancel failed")
appendDataEvents(t, releaser.tail, "ch-cancel-fail", 2)
tr := createErrorTerminalResult("ch-cancel-fail")
prog, err := releaser.FailPending(ctx, "a1", tr)
if err != nil {
t.Fatalf("FailPending: %v", err)
}
_ = prog
// Terminal should still be committed exactly once.
if len(sink.terminals) != 1 {
t.Errorf("terminals = %d, want 1", len(sink.terminals))
}
// The terminal result should include the cancel failure cause.
lastTerminal := sink.terminals[len(sink.terminals)-1]
causes := lastTerminal.FailureCauses()
if causes.Len() == 0 {
t.Error("expected failure causes to include cancel failure")
} else {
found := false
for i := 0; i < causes.Len(); i++ {
c, _ := causes.At(i)
if c.Code() == "attempt_cancel_failed" {
found = true
break
}
}
if !found {
t.Error("expected 'attempt_cancel_failed' cause code")
}
}
}
// TestStreamReleaserConcurrentEpochsAreSerialized verifies that concurrent
// release calls through the releaser serializer do not cause data races
// (the tail itself is not thread-safe, so appends must be serialized).
func TestStreamReleaserConcurrentEpochsAreSerialized(t *testing.T) {
ctx := context.Background()
releaser, _, _ := setupReleaser(t, "ch-conc", 2, 100)
releaser.boundary.BeginAttempt("a1")
rs, _ := NewResponseStart("ch-conc", 200, nil, releaserNow)
releaser.boundary.StageResponseStart("a1", rs)
// Pre-create all epochs sequentially (tail is not thread-safe).
epochs := make([]EvidenceEpoch, 5)
for i := 0; i < 5; i++ {
epochs[i] = appendDataEvents(t, releaser.tail, "ch-conc", 2)
}
var wg sync.WaitGroup
var successful int64
var errors int64
for i := 0; i < 5; i++ {
wg.Add(1)
go func(idx int) {
defer wg.Done()
_, err := releaser.ReleaseEpoch(ctx, "a1", epochs[idx].ID())
if err == nil {
atomic.AddInt64(&successful, 1)
} else {
atomic.AddInt64(&errors, 1)
}
}(i)
}
wg.Wait()
// All events should have been processed without race.
total := atomic.LoadInt64(&successful) + atomic.LoadInt64(&errors)
if total != 5 {
t.Errorf("total results = %d, want 5", total)
}
}
// TestStreamReleaserTerminalAfterFullReleaseSkipsSecondaryRelease verifies
// that after an epoch is fully released and confirmed, ReleaseTerminalEpoch
// does not attempt another release.
func TestStreamReleaserTerminalAfterFullReleaseSkipsSecondaryRelease(t *testing.T) {
ctx := context.Background()
releaser, sink, _ := setupReleaser(t, "ch-skip", 2, 100)
releaser.boundary.BeginAttempt("a1")
rs, _ := NewResponseStart("ch-skip", 200, nil, releaserNow)
releaser.boundary.StageResponseStart("a1", rs)
// Release epoch fully.
epoch := appendDataEvents(t, releaser.tail, "ch-skip", 2)
prog, err := releaser.ReleaseEpoch(ctx, "a1", epoch.ID())
if err != nil {
t.Fatalf("ReleaseEpoch: %v", err)
}
if prog.ReleasedEvents() != 2 {
t.Errorf("ReleasedEvents = %d, want 2", prog.ReleasedEvents())
}
// Terminal after full release.
tr := createSuccessTerminalResult("ch-skip")
prog2, err := releaser.ReleaseTerminalEpoch(ctx, "a1", epoch.ID(), tr)
if err != nil {
t.Fatalf("ReleaseTerminalEpoch: %v", err)
}
_ = prog2
// Should have 2 releases + 1 start, no extra.
if len(sink.releases) != 2 {
t.Errorf("releases = %d, want 2 (no secondary release)", len(sink.releases))
}
if len(sink.responseStarts) != 1 {
t.Errorf("responseStarts = %d, want 1", len(sink.responseStarts))
}
if len(sink.terminals) != 1 {
t.Errorf("terminals = %d, want 1", len(sink.terminals))
}
}