From aeb22a0e90872d6bdd3d5bab5eeebe65124de9d7 Mon Sep 17 00:00:00 2001 From: itseffi <15998472+itseffi@users.noreply.github.com> Date: Tue, 4 Aug 2026 22:58:55 +0200 Subject: [PATCH] v3 Co-Authored-By: Claude Opus 5 (1M context) --- pkg/productize/runs/replay.go | 4 +- pkg/productize/runs/tail_test.go | 66 ++++++++++++++++++++++++++++++++ 2 files changed, 69 insertions(+), 1 deletion(-) diff --git a/pkg/productize/runs/replay.go b/pkg/productize/runs/replay.go index 11a55116..218a0f21 100644 --- a/pkg/productize/runs/replay.go +++ b/pkg/productize/runs/replay.go @@ -22,6 +22,8 @@ func (r *Run) Replay(fromSeq uint64) iter.Seq2[events.Event, error] { return } + // errStopReplay means yield already returned false; the range-over-func + // contract forbids calling yield again, so it must not be reported. if err := replayRemoteEvents( context.Background(), r.client, @@ -29,7 +31,7 @@ func (r *Run) Replay(fromSeq uint64) iter.Seq2[events.Event, error] { fromSeq, RemoteCursor{}, yield, - ); err != nil { + ); err != nil && !errors.Is(err, errStopReplay) { yield(events.Event{}, err) } } diff --git a/pkg/productize/runs/tail_test.go b/pkg/productize/runs/tail_test.go index 17d8ef03..e3c1583d 100644 --- a/pkg/productize/runs/tail_test.go +++ b/pkg/productize/runs/tail_test.go @@ -43,6 +43,72 @@ func TestReplayPagesEventsInOrder(t *testing.T) { } } +func TestReplayStopsCleanlyWhenConsumerBreaksEarly(t *testing.T) { + t.Parallel() + + tests := []struct { + name string + pages []remoteRunEventPage + }{ + { + name: "single page", + pages: []remoteRunEventPage{{ + Events: []events.Event{ + testEvent("run-replay-break", 1, events.EventKindRunStarted), + testEvent("run-replay-break", 2, events.EventKindJobStarted), + }, + NextCursor: remoteCursor(time.Unix(2, 0).UTC(), 2), + }}, + }, + { + name: "paged reader", + pages: []remoteRunEventPage{ + { + Events: []events.Event{ + testEvent("run-replay-break", 1, events.EventKindRunStarted), + }, + NextCursor: remoteCursor(time.Unix(1, 0).UTC(), 1), + HasMore: true, + }, + { + Events: []events.Event{ + testEvent("run-replay-break", 2, events.EventKindJobStarted), + }, + NextCursor: remoteCursor(time.Unix(2, 0).UTC(), 2), + }, + }, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + + reader := &stubDaemonRunReader{eventPages: tt.pages} + run := &Run{ + summary: RunSummary{RunID: "run-replay-break"}, + client: reader, + } + + var seen []events.Event + for ev, err := range run.Replay(0) { + if err != nil { + t.Fatalf("Replay() unexpected error: %v", err) + } + seen = append(seen, ev) + break + } + + if seqs := collectedSeqs(seen); !slices.Equal(seqs, []uint64{1}) { + t.Fatalf("Replay() seqs = %v, want [1]", seqs) + } + if len(reader.eventCalls) != 1 { + t.Fatalf("ListRunEvents calls = %d, want 1", len(reader.eventCalls)) + } + }) + } +} + func TestReplayReportsIncompatibleSchemaVersion(t *testing.T) { run := &Run{ summary: RunSummary{RunID: "run-replay-schema"},