Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
160 changes: 97 additions & 63 deletions projection/lifecycle/worker.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
}
Expand All @@ -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) {
Expand All @@ -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
Expand Down
Loading
Loading