diff --git a/projection/lifecycle/worker.go b/projection/lifecycle/worker.go index d398e8d..df95c33 100644 --- a/projection/lifecycle/worker.go +++ b/projection/lifecycle/worker.go @@ -27,7 +27,9 @@ const defaultWorkerPollInterval = time.Second // name order, signals readiness, and then tails the sequence strictly after // the mark, folding and delivering each newly recorded cutover. Flips // superseded before the worker started are never delivered: the worker is a -// convergence mechanism, not a per-event feed. +// convergence mechanism, not a per-event feed. No setter runs while the +// worker holds a read iterator open, so a setter may share a bounded resource, +// such as a connection pool, with the reader. // // The worker keeps no durable progress — no checkpoint, no cursor, no state // shared with any other worker — so any number of workers may run @@ -147,23 +149,24 @@ func (w *Worker) Run(ctx context.Context) error { // supervision of Run — and it never closes if initialization fails. func (w *Worker) Ready() <-chan struct{} { return w.ready } -// drain reads the global sequence strictly after the given position, -// folding every cutover event through the per-name continuity folds and -// handing each accepted cutover to deliver when it is non-nil. It returns -// the last observed global position. Validation precedes delivery: a -// cutover that extends no legal history stops the drain undelivered. A -// failed iterator close is a failed drain — the iterator cannot vouch for -// the completeness of what it yielded. Cancellation is checked before the -// read begins, before each event is requested, and against every result -// before it is classified, even when the reader does not observe contexts: -// a canceled drain issues no read and requests no further events, and a -// result arriving alongside cancellation — an event, the end of the -// stream, or a cancellation-shaped failure — is dropped unprocessed, while -// an independent read failure racing the cancellation is joined with it -// rather than discarded. +// drain reads the global sequence strictly after the given position and +// folds every cutover event through the per-name continuity folds. When +// deliver is non-nil, accepted cutovers are buffered and delivered in order +// only after the iterator closes successfully. An accepted prefix is still +// delivered before a later read or validation failure is returned, and a +// delivery failure takes precedence over that later failure. It returns the +// last observed global position. A failed iterator close is a failed drain +// and prevents all buffered delivery. Cancellation is checked before the +// read begins, before each event is requested, against every result before +// it is classified, and before post-close delivery, even when the reader +// does not observe contexts: a canceled drain issues no read and requests no +// further events, and a result arriving alongside cancellation — an event, +// the end of the stream, or a cancellation-shaped failure — is dropped +// unprocessed, while an independent read failure racing the cancellation is +// joined with it rather than discarded. func (w *Worker) drain(ctx context.Context, live map[string]cutoverFold, after int64, deliver func(context.Context, Cutover) error, -) (position int64, err error) { +) (int64, error) { if err := ctx.Err(); err != nil { return after, err } @@ -174,7 +177,7 @@ func (w *Worker) drain(ctx context.Context, live map[string]cutoverFold, after i // but the cancellation is the context's own story, and an independent // failure is joined with it rather than discarded. A successfully opened // iterator proceeds — the loop's entry check reports the cancellation - // after the deferred close releases it. + // before the iterator is closed. if err != nil { if ctxErr := ctx.Err(); ctxErr != nil { if leavesMatch(err, ctxErr) { @@ -187,72 +190,103 @@ func (w *Worker) drain(ctx context.Context, live map[string]cutoverFold, after i return after, fmt.Errorf("reading events: %w", err) } - defer func() { - closeCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), iteratorCloseTimeout) - defer cancel() + position := after + var accepted []Cutover - if closeErr := iter.Close(closeCtx); closeErr != nil { - err = errors.Join(err, fmt.Errorf("closing event iterator: %w", closeErr)) - } - }() + readErr := func() error { + for { + if err := ctx.Err(); err != nil { + return err + } - position = after + event, err := iter.Next(ctx) - for { - if err := ctx.Err(); err != nil { - return position, err - } + // Cancellation dominates whatever the read returned — a result + // arriving alongside it is not acted on — but an independent read + // failure is joined rather than discarded: only a result whose + // every leaf is the end of the stream or the cancellation itself + // folds into the bare cancellation. + if ctxErr := ctx.Err(); ctxErr != nil { + if err == nil || leavesMatch(err, eventstore.ErrEndOfEventStream, ctxErr) { + return ctxErr + } - event, err := iter.Next(ctx) + return errors.Join(ctxErr, fmt.Errorf("reading event: %w", err)) + } - // Cancellation dominates whatever the read returned — a result - // arriving alongside it is not acted on — but an independent read - // failure is joined rather than discarded: only a result whose - // every leaf is the end of the stream or the cancellation itself - // folds into the bare cancellation. - if ctxErr := ctx.Err(); ctxErr != nil { - if err == nil || leavesMatch(err, eventstore.ErrEndOfEventStream, ctxErr) { - return position, ctxErr + // The end of the stream is clean only when it is the read's whole + // story: a failure joined with it is a failed read, not a finished + // one. + if leavesMatch(err, eventstore.ErrEndOfEventStream) { + return nil + } else if err != nil { + return fmt.Errorf("reading event: %w", err) } - return position, errors.Join(ctxErr, fmt.Errorf("reading event: %w", err)) - } + if event.GlobalPosition != nil { + position = *event.GlobalPosition + } - // The end of the stream is clean only when it is the read's whole - // story: a failure joined with it is a failed read, not a finished - // one. - if leavesMatch(err, eventstore.ErrEndOfEventStream) { - return position, nil - } else if err != nil { - return position, fmt.Errorf("reading event: %w", err) - } + raw, ok, err := decodeCutover(event) + if err != nil { + return err + } else if !ok { + continue + } - if event.GlobalPosition != nil { - position = *event.GlobalPosition - } + next, err := live[raw.cutover.Live.Name].apply(raw) + if err != nil { + return err + } - raw, ok, err := decodeCutover(event) - if err != nil { - return position, err - } else if !ok { - continue + live[raw.cutover.Live.Name] = next + + if deliver != nil { + accepted = append(accepted, raw.cutover) + } } + }() - next, err := live[raw.cutover.Live.Name].apply(raw) - if err != nil { - return position, err + closeCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), iteratorCloseTimeout) + closeErr := iter.Close(closeCtx) + cancel() + + ctxErr := ctx.Err() + if closeErr != nil { + if ctxErr != nil { + switch { + case readErr == nil || leavesMatch(readErr, ctxErr): + readErr = ctxErr + case !errors.Is(readErr, ctxErr): + readErr = errors.Join(ctxErr, readErr) + } } - live[raw.cutover.Live.Name] = next + return position, errors.Join(readErr, fmt.Errorf("closing event iterator: %w", closeErr)) + } - if deliver == nil { - continue + if ctxErr != nil { + switch { + case readErr == nil || leavesMatch(readErr, ctxErr): + return position, ctxErr + case errors.Is(readErr, ctxErr): + return position, readErr + default: + return position, errors.Join(ctxErr, readErr) } + } - if err := deliver(ctx, raw.cutover); err != nil { + for _, cutover := range accepted { + if err := deliver(ctx, cutover); err != nil { + return position, err + } + + if err := ctx.Err(); err != nil { return position, err } } + + return position, readErr } // deliver applies one cutover through every registered setter, in diff --git a/projection/lifecycle/worker_internal_test.go b/projection/lifecycle/worker_internal_test.go index 94add26..dc9cb20 100644 --- a/projection/lifecycle/worker_internal_test.go +++ b/projection/lifecycle/worker_internal_test.go @@ -5,6 +5,7 @@ import ( "encoding/json" "errors" "fmt" + "strings" "testing" "github.com/go-estoria/estoria/eventstore" @@ -47,6 +48,59 @@ func (r stubReader) ReadAll(context.Context, eventstore.ReadAllOptions) (eventst return r.iter, nil } +// scriptedIterator yields its events, then one terminal result, and records +// whether Close has completed so setters can assert the iterator lifecycle. +type scriptedIterator struct { + events []*eventstore.Event + terminalErr error + closeErr error + onClose func() + + next int + open bool +} + +func newScriptedIterator(events ...*eventstore.Event) *scriptedIterator { + return &scriptedIterator{events: events, open: true} +} + +func (i *scriptedIterator) Next(context.Context) (*eventstore.Event, error) { + if i.next < len(i.events) { + event := i.events[i.next] + i.next++ + + return event, nil + } + + if i.terminalErr != nil { + return nil, i.terminalErr + } + + return nil, eventstore.ErrEndOfEventStream +} + +func (i *scriptedIterator) Close(context.Context) error { + i.open = false + + if i.onClose != nil { + i.onClose() + } + + return i.closeErr +} + +type callbackSetter struct { + apply func(context.Context, Cutover) error +} + +func (s callbackSetter) ApplyCutover(ctx context.Context, cutover Cutover) error { + return s.apply(ctx, cutover) +} + +func (callbackSetter) AppliedCutover(context.Context, string) (Cutover, error) { + return Cutover{}, ErrNoLiveVersion +} + // inHandIterator yields one event, firing a trigger — a cancellation, or a // wait that outlives a deadline — from inside the yielding Next call. type inHandIterator struct { @@ -92,7 +146,13 @@ func (c *deadlineTrippedCtx) trip() { c.tripped = true } func promotedEvent(t *testing.T, position int64) *eventstore.Event { t.Helper() - data, err := json.Marshal(Promoted{Next: projection.ID{Name: "orders", Version: 1}, Revision: 1}) + return promotedEventWith(t, position, Promoted{Next: projection.ID{Name: "orders", Version: 1}, Revision: 1}) +} + +func promotedEventWith(t *testing.T, position int64, promoted Promoted) *eventstore.Event { + t.Helper() + + data, err := json.Marshal(promoted) if err != nil { t.Fatalf("marshaling promoted event: %v", err) } @@ -134,6 +194,338 @@ func TestDrainCancellationPrecedesTheRead(t *testing.T) { } } +func TestDrainTailDeliveryWaitsForIteratorClose(t *testing.T) { + t.Parallel() + + ordersV1 := projection.ID{Name: "orders", Version: 1} + ordersV2 := projection.ID{Name: "orders", Version: 2} + iter := newScriptedIterator( + promotedEvent(t, 9), + promotedEventWith(t, 10, Promoted{Previous: ordersV1, Next: ordersV2, Revision: 2}), + ) + var applied []Cutover + + setter := callbackSetter{apply: func(_ context.Context, cutover Cutover) error { + if iter.open { + t.Error("setter ran while the iterator was open") + } + + applied = append(applied, cutover) + + return nil + }} + + worker, err := NewWorker(stubReader{iter: iter}, WithCutoverSetter(setter)) + if err != nil { + t.Fatalf("creating worker: %v", err) + } + + live := map[string]cutoverFold{} + position, err := worker.drain(t.Context(), live, 3, worker.deliver) + if err != nil { + t.Fatalf("draining tail: %v", err) + } + + if iter.open { + t.Error("want the iterator closed before drain returned") + } + + if position != 10 { + t.Errorf("want position 10, got %d", position) + } + + want := []Cutover{ + {Live: ordersV1, Revision: 1}, + {Live: ordersV2, Revision: 2}, + } + if len(applied) != len(want) || applied[0] != want[0] || applied[1] != want[1] { + t.Errorf("want deliveries %v in order after close, got %v", want, applied) + } + + if got := live["orders"].current; got != want[1] { + t.Errorf("want the final accepted cutover folded as %+v, got %+v", want[1], got) + } +} + +func TestDrainTailCloseFailurePreventsBufferedDelivery(t *testing.T) { + t.Parallel() + + errClose := errors.New("close failed") + iter := newScriptedIterator(promotedEvent(t, 9)) + iter.closeErr = errClose + + setter := &untouchableSetter{} + worker, err := NewWorker(stubReader{iter: iter}, WithCutoverSetter(setter)) + if err != nil { + t.Fatalf("creating worker: %v", err) + } + + live := map[string]cutoverFold{} + position, err := worker.drain(t.Context(), live, 3, worker.deliver) + if !errors.Is(err, errClose) { + t.Fatalf("want the close failure, got %v", err) + } + + if setter.touched { + t.Error("want no delivery from an iterator that failed to close") + } + + if position != 9 { + t.Errorf("want position 9, got %d", position) + } + + if got := live["orders"].current.Revision; got != 1 { + t.Errorf("want the accepted cutover retained in fold state, got revision %d", got) + } +} + +func TestDrainTailReadAndCloseFailuresPreventBufferedDelivery(t *testing.T) { + t.Parallel() + + errRead := errors.New("read failed") + errClose := errors.New("close failed") + iter := newScriptedIterator(promotedEvent(t, 9)) + iter.terminalErr = errRead + iter.closeErr = errClose + + setter := &untouchableSetter{} + worker, err := NewWorker(stubReader{iter: iter}, WithCutoverSetter(setter)) + if err != nil { + t.Fatalf("creating worker: %v", err) + } + + _, err = worker.drain(t.Context(), map[string]cutoverFold{}, 3, worker.deliver) + if !errors.Is(err, errRead) || !errors.Is(err, errClose) { + t.Fatalf("want the read and close failures joined, got %v", err) + } + + if setter.touched { + t.Error("want no buffered delivery after the close failure") + } +} + +func TestDrainTailValidationAndCloseFailuresPreventBufferedDelivery(t *testing.T) { + t.Parallel() + + errClose := errors.New("close failed") + ordersV1 := projection.ID{Name: "orders", Version: 1} + ordersV2 := projection.ID{Name: "orders", Version: 2} + iter := newScriptedIterator( + promotedEvent(t, 9), + promotedEventWith(t, 10, Promoted{Previous: ordersV1, Next: ordersV2, Revision: 5}), + ) + iter.closeErr = errClose + + setter := &untouchableSetter{} + worker, err := NewWorker(stubReader{iter: iter}, WithCutoverSetter(setter)) + if err != nil { + t.Fatalf("creating worker: %v", err) + } + + _, err = worker.drain(t.Context(), map[string]cutoverFold{}, 3, worker.deliver) + if !errors.Is(err, errClose) || !strings.Contains(err.Error(), "records revision 5 after revision 1") { + t.Fatalf("want the validation and close failures joined, got %v", err) + } + + if setter.touched { + t.Error("want no buffered delivery after the close failure") + } +} + +func TestDrainTailDeliversAcceptedPrefixBeforeLaterFailure(t *testing.T) { + t.Parallel() + + errRead := errors.New("read failed") + + for _, tt := range []struct { + name string + iterator func(*testing.T) *scriptedIterator + wantErr error + wantContains string + wantPosition int64 + }{ + { + name: "read failure", + iterator: func(t *testing.T) *scriptedIterator { + t.Helper() + + iter := newScriptedIterator(promotedEvent(t, 9)) + iter.terminalErr = errRead + + return iter + }, + wantErr: errRead, + wantPosition: 9, + }, + { + name: "validation failure", + iterator: func(t *testing.T) *scriptedIterator { + t.Helper() + + return newScriptedIterator(promotedEvent(t, 9), promotedEvent(t, 10)) + }, + wantContains: "records revision 1 after revision 1", + wantPosition: 10, + }, + } { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + + iter := tt.iterator(t) + var applied []Cutover + + setter := callbackSetter{apply: func(_ context.Context, cutover Cutover) error { + if iter.open { + t.Error("setter ran while the iterator was open") + } + + applied = append(applied, cutover) + + return nil + }} + + worker, err := NewWorker(stubReader{iter: iter}, WithCutoverSetter(setter)) + if err != nil { + t.Fatalf("creating worker: %v", err) + } + + position, err := worker.drain(t.Context(), map[string]cutoverFold{}, 3, worker.deliver) + if tt.wantErr != nil && !errors.Is(err, tt.wantErr) { + t.Fatalf("want failure %v after prefix delivery, got %v", tt.wantErr, err) + } + + if tt.wantContains != "" && (err == nil || !strings.Contains(err.Error(), tt.wantContains)) { + t.Fatalf("want failure containing %q after prefix delivery, got %v", tt.wantContains, err) + } + + if position != tt.wantPosition { + t.Errorf("want position %d, got %d", tt.wantPosition, position) + } + + want := Cutover{Live: projection.ID{Name: "orders", Version: 1}, Revision: 1} + if len(applied) != 1 || applied[0] != want { + t.Errorf("want accepted prefix %v delivered after close, got %v", want, applied) + } + }) + } +} + +func TestDrainTailDeliveryFailurePrecedesLaterReadFailure(t *testing.T) { + t.Parallel() + + errRead := errors.New("read failed") + errDelivery := errors.New("delivery failed") + iter := newScriptedIterator(promotedEvent(t, 9)) + iter.terminalErr = errRead + + setter := callbackSetter{apply: func(context.Context, Cutover) error { + if iter.open { + t.Error("setter ran while the iterator was open") + } + + return errDelivery + }} + + worker, err := NewWorker(stubReader{iter: iter}, WithCutoverSetter(setter)) + if err != nil { + t.Fatalf("creating worker: %v", err) + } + + _, err = worker.drain(t.Context(), map[string]cutoverFold{}, 3, worker.deliver) + if !errors.Is(err, errDelivery) { + t.Fatalf("want the delivery failure, got %v", err) + } + + if errors.Is(err, errRead) { + t.Fatalf("want the later read failure suppressed by the delivery failure, got %v", err) + } +} + +func TestDrainTailCancellationBeforeDeliveryDropsBufferedCutovers(t *testing.T) { + t.Parallel() + + errRead := errors.New("later read failed") + ordersV1 := projection.ID{Name: "orders", Version: 1} + ordersV2 := projection.ID{Name: "orders", Version: 2} + + for _, tt := range []struct { + name string + iterator func(*testing.T) *scriptedIterator + wantErr error + wantContains string + wantPosition int64 + exactContext bool + }{ + { + name: "without an independent failure", + iterator: func(t *testing.T) *scriptedIterator { + t.Helper() + return newScriptedIterator(promotedEvent(t, 9)) + }, + wantPosition: 9, + exactContext: true, + }, + { + name: "with an independent read failure", + iterator: func(t *testing.T) *scriptedIterator { + t.Helper() + iter := newScriptedIterator(promotedEvent(t, 9)) + iter.terminalErr = errRead + return iter + }, + wantErr: errRead, + wantPosition: 9, + }, + { + name: "with an independent validation failure", + iterator: func(t *testing.T) *scriptedIterator { + t.Helper() + return newScriptedIterator( + promotedEvent(t, 9), + promotedEventWith(t, 10, Promoted{Previous: ordersV1, Next: ordersV2, Revision: 5}), + ) + }, + wantContains: "records revision 5 after revision 1", + wantPosition: 10, + }, + } { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + + ctx, cancel := context.WithCancel(t.Context()) + t.Cleanup(cancel) + + iter := tt.iterator(t) + iter.onClose = cancel + setter := &untouchableSetter{} + worker, err := NewWorker(stubReader{iter: iter}, WithCutoverSetter(setter)) + if err != nil { + t.Fatalf("creating worker: %v", err) + } + + position, err := worker.drain(ctx, map[string]cutoverFold{}, 3, worker.deliver) + if !errors.Is(err, context.Canceled) { + t.Fatalf("want context cancellation, got %v", err) + } + if tt.exactContext && err != context.Canceled { //nolint:errorlint // Exact identity is the contract without an independent failure. + t.Fatalf("want exactly the context error, got %v", err) + } + if tt.wantErr != nil && !errors.Is(err, tt.wantErr) { + t.Fatalf("want independent failure %v joined with cancellation, got %v", tt.wantErr, err) + } + if tt.wantContains != "" && !strings.Contains(err.Error(), tt.wantContains) { + t.Fatalf("want independent failure containing %q joined with cancellation, got %v", tt.wantContains, err) + } + if setter.touched { + t.Error("want no buffered delivery after cancellation") + } + if position != tt.wantPosition { + t.Errorf("want position %d, got %d", tt.wantPosition, position) + } + }) + } +} + // TestDrainDropsTheEventInHand pins that a cutover read alongside a context // ending is wholly unprocessed — position unmoved, fold untouched — and // that the result is exactly the context's own error: not a hard-coded diff --git a/projection/lifecycle/worker_test.go b/projection/lifecycle/worker_test.go index 74e9925..c735203 100644 --- a/projection/lifecycle/worker_test.go +++ b/projection/lifecycle/worker_test.go @@ -153,6 +153,83 @@ func awaitExit(t *testing.T, runErr <-chan error) error { } } +// lease is a one-slot resource that refuses a second holder rather than +// blocking, so holding it across a call that needs it surfaces as an error. +type lease struct { + mu sync.Mutex + holder string +} + +func (l *lease) acquire(holder string) error { + l.mu.Lock() + defer l.mu.Unlock() + + if l.holder != "" { + return errors.New(holder + " needs the resource while " + l.holder + " holds it") + } + + l.holder = holder + return nil +} + +func (l *lease) release() { + l.mu.Lock() + defer l.mu.Unlock() + + l.holder = "" +} + +// leasedReader holds the lease from ReadAll until the iterator is closed, as +// a pooled connection is pinned by an open cursor. +type leasedReader struct { + inner eventstore.GlobalReader + lease *lease +} + +func (r *leasedReader) ReadAll(ctx context.Context, opts eventstore.ReadAllOptions) (eventstore.StreamIterator, error) { + if err := r.lease.acquire("an open iterator"); err != nil { + return nil, err + } + + iter, err := r.inner.ReadAll(ctx, opts) + if err != nil { + r.lease.release() + return nil, err + } + + return &leasedIterator{StreamIterator: iter, lease: r.lease}, nil +} + +type leasedIterator struct { + eventstore.StreamIterator + lease *lease + closed bool +} + +func (i *leasedIterator) Close(ctx context.Context) error { + if !i.closed { + i.closed = true + i.lease.release() + } + + return i.StreamIterator.Close(ctx) +} + +// leasedSetter records cutovers while holding the same lease as the reader. +type leasedSetter struct { + recordingSetter + lease *lease +} + +func (s *leasedSetter) ApplyCutover(ctx context.Context, cutover lifecycle.Cutover) error { + if err := s.lease.acquire("the setter"); err != nil { + return err + } + defer s.lease.release() + + return s.recordingSetter.ApplyCutover(ctx, cutover) +} + // countingReader wraps a GlobalReader, recording the AfterPosition of every // read issued through it once the read has been handed out. type countingReader struct { @@ -1263,6 +1340,58 @@ func TestWorker_TailCloseFailureStopsTheWorker(t *testing.T) { assertCutovers(t, recorder.seen(), []lifecycle.Cutover{{Live: ordersV1, Revision: 1}}) } +// TestWorker_TailDeliveryFollowsIteratorClose pins the resource order end to +// end: the reader and setter share a one-slot resource, so delivery under an +// open iterator would stop the worker instead of deadlocking the test. +func TestWorker_TailDeliveryFollowsIteratorClose(t *testing.T) { + t.Parallel() + + events := newEventStore(t) + projections, err := lifecycle.NewStore(events) + if err != nil { + t.Fatalf("creating lifecycle store: %v", err) + } + + ordersV1 := projection.ID{Name: "orders", Version: 1} + ordersV2 := projection.ID{Name: "orders", Version: 2} + recordCutover(t, projections, ordersV1, projection.ID{}, false) + + slot := &lease{} + setter := &leasedSetter{lease: slot} + worker, err := lifecycle.NewWorker(&leasedReader{inner: events, lease: slot}, + lifecycle.WithCutoverSetter(setter), + lifecycle.WithPollInterval(2*time.Millisecond), + ) + if err != nil { + t.Fatalf("creating worker: %v", err) + } + + runErr, cancel := runWorker(t, worker) + waitReady(t, worker) + recordCutover(t, projections, ordersV2, ordersV1, false) + + deadline := time.After(waitTimeout) + for len(setter.seen()) < 2 { + select { + case err := <-runErr: + t.Fatalf("want the tail cutover delivered on a one-slot resource, got %v", err) + case <-deadline: + t.Fatal("timed out waiting for the tail delivery") + case <-time.After(time.Millisecond): + } + } + + assertCutovers(t, setter.seen(), []lifecycle.Cutover{ + {Live: ordersV1, Revision: 1}, + {Live: ordersV2, Revision: 2}, + }) + + cancel() + if err := awaitExit(t, runErr); !errors.Is(err, context.Canceled) { + t.Fatalf("stopping worker: %v", err) + } +} + // TestWorker_CancellationStopsInitialization pins the entry check: a worker // started with a canceled context touches nothing — no read issued, no // setter acting, readiness withheld — and Run returns the context's error.