Skip to content

feat: make admission and shutdown safe (Phase 2) - #29

Merged
heynemann merged 13 commits into
phase-1-performance-baselinefrom
phase-2-lifecycle-safety
Aug 8, 2026
Merged

heynemann merged 13 commits into
phase-1-performance-baselinefrom
phase-2-lifecycle-safety

Conversation

@heynemann

@heynemann heynemann commented Aug 6, 2026 •

Copy link
Copy Markdown
Contributor

What

Phase 2 of the batcher performance plan: make admission and shutdown safe. All
four milestones are complete.

Commit Milestone
test: characterize batching engine before replacing rill 2.1
refactor: replace rill with an owned aggregation pipeline 2.1
test: enforce the per-batcher goroutine budget 2.1
feat: own the intake queue and add an admission gate 2.2
feat: add resumable Shutdown and forward the fx stop context 2.3
feat: bound diagnostics and recover processor panics per batch 2.4
docs: record Phase 2 measured outcomes docs
docs: document the Phase 2 lifecycle and admission APIs docs

The defects this fixes

Each was reproducible on main and is now verified fixed by direct measurement,
not just by the tests added here:

Defect Before After
Partial batch, interval > close grace 50 accepted, 0 processed, Len() = 50 50 accepted, 50 processed, Len() = 0
Add after Close panics the process counted rejection
Sustained overload unbounded, ~4.3GB in 2s queue capped, rejections reported
isClosed data race under -race atomic lifecycle state
Start called twice double-close panic idempotent
Goroutines per batcher 6 2, none leaked

Design decisions worth reviewing

The intake queue is never closed. chann's relay goroutine can only exit by
closing its own ingress channel — the channel publishers send to — so "never close
the input" and "leak no goroutines" were mutually exclusive. Owning the queue
resolves that, and removes the entire class of send-on-closed-channel panic.

Reserve before publishing, roll back before leaving the gate. Reserving after
publication would let the drain conclude while a publisher was mid-publish; rolling
back after leaving the gate would leave a phantom obligation the drain waits on
forever.

Two lost-wakeup hazards, both with sabotage-verified tests. Every gate
decrement routes through leave(), including enter()'s own rejection path, and
the coordinator performs its own post-seal gate check. Removing either mechanism
makes its test hang — I checked, because a concurrency test that cannot fail is
worthless.

Shutdown reports, never discards. Shutdown(ctx) is resumable: a later call
waits on the same drain rather than restarting it. A caller's deadline is per-caller
and never stored as the terminal result, so one impatient caller cannot poison
another's answer.

Panics are recovered per batch, not per consumer. Batcher owns goroutines the
caller cannot reach, so an unrecovered panic crashed the process and destroyed every
other queued item. Recovery is scoped to the batch so one poison batch cannot stall
everything behind it.

Performance

Removing the relay layers made the hot path substantially faster, measured with
benchstat over -count=6 against the stored baseline:

  • -54% geomean sec/op on enqueue (Add 229.8ns → 63.4ns at small batch sizes)
  • -49% to -90% bytes/op
  • Goroutines per batcher 6 → 2

This is the one place in the plan where a safety change also made things faster,
because each relay cost a channel hop and a goroutine handoff per item.

Two plan predictions were wrong

Both caught by measurement, both corrected in the plan with reasoning:

  1. Merging aggregation into the processing loop changes latency. It inverted
    the documented baseline (5ms window becoming faster than 100ms with a 50ms
    processor). I hit this twice — once in 2.1 and again in 2.2 — and reverted both
    times. Aggregation and processing stay separate until Phase 3 decides the worker
    model deliberately.
  2. Goroutines could not reach the predicted counts. Predicted 3, then 4;
    measured 5 after 2.1, then 2 after 2.2. Reaching fewer required behaviour
    changes that Phase 2 is not allowed to make.

Bug found while testing

Start racing Shutdown deadlocked: Start declined to launch the consumer when
it saw a sealed gate, but shared startOnce with the drain path, so whichever call
won left the drain waiting on a consumer that would never exist. Start now always
launches the consumer, because sealed does not mean empty. Covered across 200
trials.

Validation

  • go test -race ./... clean; go vet ./... clean.
  • All 8 characterization tests pass unchanged from before the rewrite.
  • The Phase 1 inversion baseline still holds, confirming Phase 2 did not alter the
    window/processor coupling.
  • Sabotage-verified: coordinator gate check, leave() routing, and panic recovery
    each fail their test when removed.
  • chann and rill both removed from go.mod.

Dependency context

Stacked on #28 (Phase 1); targets phase-1-performance-baseline. Review #28 first.

Phase 3 is next: wiring WithConcurrency to real workers, which is what finally
removes the effective window = max(window, processor duration) coupling and makes
5-10ms windows worthwhile.

Stack created with GitHub Stacks CLI • Give Feedback 💬

Summary by CodeRabbit

  • New Features
    • Added bounded queues with context-aware Enqueue and configurable diagnostic error buffers.
    • Added resumable shutdown with configurable grace periods, shutdown status checks, and detailed incomplete-drain errors.
    • Added Stats() snapshots for tracking pending, queued, completed, failed, rejected, and dropped work.
    • Added processor panic reporting and continued processing where possible.
  • Bug Fixes
    • Prevented work loss, deadlocks, empty batches, and late-admission panics during shutdown.
  • Documentation
    • Expanded shutdown, queue capacity, error handling, and statistics guidance.
  • Chores
    • Updated the required Go version to 1.25.0.

@coderabbitai

coderabbitai Bot commented Aug 6, 2026 •

Copy link
Copy Markdown

Review Change Stack

Important

Review skipped

Auto reviews are disabled on base/target branches other than the default branch.

Please check the settings in the CodeRabbit UI or the .coderabbit.yaml file in this repository. To trigger a single review, invoke the @coderabbitai review command.

⚙️ Run configuration

Configuration used: Organization UI

Review profile: CHILL

Plan: Pro

Run ID: ce1ad663-7132-474a-8884-2646c9b331ab

You can disable this status message by setting the reviews.review_status to false in the CodeRabbit configuration file.

Use the checkbox below for a quick retry:

  • 🔍 Trigger review
📝 Walkthrough

Walkthrough

The batcher replaces its external pipeline with internal bounded admission, queueing, aggregation, processing, diagnostics, statistics, and resumable shutdown. Tests cover concurrency, draining, backpressure, panic recovery, goroutine budgets, and lifecycle behavior. Documentation and tooling now use Go 1.25.0.

Changes

Batcher lifecycle and processing

Layer / File(s) Summary
Queue, admission, and observability primitives
pkg/batcher/queue.go, pkg/batcher/gate.go, pkg/batcher/stats.go, pkg/batcher/options.go, pkg/batcher/constants.go, pkg/batcher/shutdown.go
Adds bounded queueing, admission sealing, statistics snapshots, configuration options, defaults, and typed shutdown and panic errors.
Processing and resumable shutdown
pkg/batcher/batcher.go
Adds internal aggregation and processing, context-aware admission, panic recovery, nonblocking diagnostics, lifecycle states, and resumable draining.
Batching, admission, and shutdown validation
pkg/batcher/*_test.go
Adds characterization and concurrency tests for batching, admission races, queue bounds, shutdown draining, idempotency, errors, and statistics.
Diagnostics, panic, and goroutine validation
pkg/batcher/diagnostics_test.go, pkg/batcher/goroutine_*_test.go, pkg/batcher/batcher_test.go, pkg/batcher/fx_test.go
Validates diagnostic buffering, panic recovery, allocation behavior, goroutine budgets, leak handling, and standard-library atomic counters.
Lifecycle integration and project documentation
pkg/batcher/fx.go, README.md, docs/improvements/*, .github/workflows/bench.yml, go.mod
Wires Fx shutdown to Shutdown(ctx), documents the new APIs, updates performance measurements, selects the module Go version in benchmarks, and raises the minimum Go version to 1.25.0.

Estimated code review effort: 5 (Critical) | ~120 minutes

Sequence Diagram(s)

sequenceDiagram
  participant Producer
  participant Batcher
  participant Queue
  participant Processor
  participant ShutdownCaller
  Producer->>Batcher: Enqueue(ctx, item)
  Batcher->>Queue: publish accepted item
  Queue-->>Batcher: item-ready signal
  Batcher->>Processor: flush and process batch
  Processor-->>Batcher: result or diagnostic
  ShutdownCaller->>Batcher: Shutdown(ctx)
  Batcher->>Queue: drain accepted work
  Batcher->>Processor: process final partial batch
  Batcher-->>ShutdownCaller: completion or ShutdownIncompleteError
Loading

Possibly related PRs

  • NSXBet/batcher#28: Shares performance-planning, threshold, and benchmark workflow changes.
🚥 Pre-merge checks | ✅ 5
✅ Passed checks (5 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The title clearly summarizes the pull request's main changes to admission and shutdown safety in Phase 2.
Docstring Coverage ✅ Passed Docstring coverage is 84.75% which is sufficient. The required threshold is 80.00%.
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.
✨ Finishing Touches
📝 Generate docstrings
  • Create stacked PR
  • Commit on current branch
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch phase-2-lifecycle-safety

Comment @coderabbitai help to get the list of available commands.

@heynemann heynemann changed the title phase 2 lifecycle safety refactor: replace rill with an owned aggregation pipeline (Phase 2.1) Aug 6, 2026
@heynemann heynemann changed the title refactor: replace rill with an owned aggregation pipeline (Phase 2.1) feat: make admission and shutdown safe (Phase 2) Aug 6, 2026
@heynemann
heynemann force-pushed the phase-2-lifecycle-safety branch from 88820cc to 8f1dfeb Compare August 6, 2026 13:00

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 10

🧹 Nitpick comments (9)
pkg/batcher/characterization_test.go (2)

387-404: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value

Replace the hand-rolled itoa with strconv.Itoa.

itoa duplicates standard-library behaviour, allocates on every digit through the prepend pattern, and returns wrong output for negative input. Test helpers gain nothing from this.

♻️ Proposed refactor
 func keyFor(i int) string {
-	return "item-" + itoa(i)
-}
-
-func itoa(i int) string {
-	if i == 0 {
-		return "0"
-	}
-
-	var digits []byte
-
-	for i > 0 {
-		digits = append([]byte{byte('0' + i%10)}, digits...)
-		i /= 10
-	}
-
-	return string(digits)
+	return "item-" + strconv.Itoa(i)
 }

Add "strconv" to the import block.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@pkg/batcher/characterization_test.go` around lines 387 - 404, Remove the
hand-rolled itoa helper and update keyFor to use strconv.Itoa directly for
integer conversion. Add the strconv import and preserve keyFor’s existing
“item-” prefix and output behavior.

105-117: 🎯 Functional Correctness | 🔵 Trivial | ⚡ Quick win

The assertion under-pins the documented behaviour.

The doc comment claims the interval is measured from the first item, not from an earlier tick. elapsed >= interval/2 only proves the flush was not immediate. A periodic-ticker implementation aligned to interval boundaries would still pass roughly half the time. Assert a tighter lower bound, for example elapsed >= interval minus a small tolerance, so a ticker-based reimplementation fails.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@pkg/batcher/characterization_test.go` around lines 105 - 117, The timing
assertion in the first-item flush test should enforce that the batch does not
flush before approximately the full interval measured from Add of the first
item. Update the require.GreaterOrEqual check around start, b.Add, and elapsed
to use interval minus a small, explicit tolerance, while preserving the existing
batch-count and item-count assertions.
pkg/batcher/admission_test.go (2)

226-259: 🩺 Stability & Availability | 🔵 Trivial | ⚡ Quick win

Close the batcher when the test ends.

This test never closes b. The aggregator and processing goroutines stay alive for the rest of the package run, and the queue still holds accepted items. Other tests in this package assert goroutine budgets, so a leaked lifecycle can make them flaky.

🧹 Proposed fix
 	release := make(chan struct{})
 	defer close(release)
 
 	b := batcher.New(
 		batcher.WithBatchSize[int](1),
 		batcher.WithBatchInterval[int](time.Hour),
 		batcher.WithMaxQueueSize[int](2),
 		batcher.WithProcessor(func([]int) error {
 			<-release
 
 			return nil
 		}),
 	)
+
+	t.Cleanup(func() { _ = b.Close() })

Note: register the cleanup after defer close(release) so the processor unblocks before Close waits on the drain. Same applies to TestBlockedEnqueueIsReleasedByShutdown, which starts Close in a goroutine and never waits for it.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@pkg/batcher/admission_test.go` around lines 226 - 259, Register test cleanup
for the batcher instance after defer close(release) so the processor is
unblocked before b.Close waits for draining; apply the same lifecycle cleanup
and ensure the Close goroutine is awaited in
TestBlockedEnqueueIsReleasedByShutdown.

31-35: 🚀 Performance & Scalability | 🔵 Trivial | 💤 Low value

Unconditional high trial counts across the race tests. Five tests repeat full batcher lifecycles hundreds of times with millisecond intervals. The shared root cause is a fixed repetition count with no testing.Short() path, which makes every local go test ./... pay the full cost and multiplies under -race.

  • pkg/batcher/admission_test.go#L31-L35: make trials a variable and lower it to a small value when testing.Short() is set; apply the same to the 300-trial loops at Lines 110 and 145 in this file.
  • pkg/batcher/shutdown_test.go#L190-L190: apply the same testing.Short() reduction to this 300-trial loop.
  • pkg/batcher/shutdown_test.go#L284-L284: apply the same testing.Short() reduction to this 200-trial loop.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@pkg/batcher/admission_test.go` around lines 31 - 35, Reduce repeated
race-test cost by making the shared trials value in
pkg/batcher/admission_test.go lines 31-35 short-test aware, using a small count
when testing.Short() is enabled while preserving the current default. Apply the
same testing.Short() reduction to the 300-trial loop at
pkg/batcher/admission_test.go lines 110 and 145, the 300-trial loop at
pkg/batcher/shutdown_test.go line 190, and the 200-trial loop at
pkg/batcher/shutdown_test.go line 284; leave normal trial counts unchanged.
pkg/batcher/batcher.go (2)

166-171: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value

Rename the nonBlockingWhenFull parameter; it never applies when the queue can be full.

The flag only selects the fast path when MaxQueueSize <= 0, where the queue cannot be full. In bounded mode Add passes true and still blocks in push. The name states the opposite of the behaviour. A name such as unboundedFastPath matches the condition at Line 166.

♻️ Proposed rename
-func (b *Batcher[T]) publish(ctx context.Context, item T, nonBlockingWhenFull bool) error {
+func (b *Batcher[T]) publish(ctx context.Context, item T, allowUnboundedFastPath bool) error {
-	if nonBlockingWhenFull && b.config.MaxQueueSize <= 0 {
+	if allowUnboundedFastPath && b.config.MaxQueueSize <= 0 {
 		// Unbounded fast path: publication cannot fail.
 		_ = b.input.tryPush(item)
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@pkg/batcher/batcher.go` around lines 166 - 171, Rename the
nonBlockingWhenFull parameter and all references in the surrounding Add flow to
unboundedFastPath, preserving the existing condition and behavior so the fast
path is selected only when MaxQueueSize <= 0 while bounded queues continue using
push.

553-563: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win

The finish fallback is unreachable, and its failure mode is silent.

coordinateDrain runs once under shutdownOnce, and it CASes stateSealing→stateDraining at Line 511 before calling finish. beginSealing never lowers the state below stateSealing. So the state is always stateDraining here and the first CAS always succeeds. The comment "Another caller may already be finishing" does not match the code.

If the second branch ever did run, errorsChan would stay open and every consumer of Errors() would block forever. Make that path loud instead of silent, or drop it.

♻️ Proposed simplification
 func (b *Batcher[T]) finish() {
-	if b.gate.state.CompareAndSwap(stateDraining, stateClosed) {
-		close(b.errorsChan)
-
-		return
-	}
-
-	// Another caller may already be finishing; make sure the terminal state is
-	// visible before returning.
-	b.gate.state.CompareAndSwap(stateSealing, stateClosed)
+	// coordinateDrain runs exactly once and has already moved the state to
+	// stateDraining, so this transition cannot lose a race.
+	b.gate.state.Store(stateClosed)
+
+	close(b.errorsChan)
 }
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@pkg/batcher/batcher.go` around lines 553 - 563, Update Batcher.finish to
remove the unreachable stateSealing fallback and its misleading comment, or make
any unexpected state transition fail loudly instead of returning silently with
errorsChan open. Preserve the successful stateDraining-to-stateClosed transition
and errorsChan closure.
pkg/batcher/gate.go (1)

124-127: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win

Remove the unused isSealed method, or use it.

golangci-lint reports func (*admissionGate).isSealed is unused. That fails the lint stage. The doc comment also names the field, not the method.

♻️ Proposed removal
-// sealed reports whether admission is closed.
-func (g *admissionGate) isSealed() bool {
-	return g.sealed.Load()
-}
-
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@pkg/batcher/gate.go` around lines 124 - 127, Remove the unused
admissionGate.isSealed method from pkg/batcher/gate.go, including its inaccurate
doc comment, rather than leaving an unreferenced helper that fails lint.

Source: Linters/SAST tools

pkg/batcher/queue.go (1)

60-92: 🚀 Performance & Scalability | 🔵 Trivial | 💤 Low value

Consider re-arming notFull after a successful push while space remains.

notFull holds one pending signal. One pop therefore releases exactly one parked publisher, even when the bounded queue has many free slots. The remaining publishers stay parked until the next pop. Progress still happens, because the aggregator pops on every wakeup, but admission latency under a full bounded queue is higher than the free space justifies.

Signal notFull again after a successful append when len(q.items)-q.head < q.capacity still holds.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@pkg/batcher/queue.go` around lines 60 - 92, Update queue.push so that after a
successful append to a bounded queue, it re-signals notFull when remaining
queued items are still below capacity. Preserve the existing unbounded-queue
behavior and only re-arm the signal after the item is appended and the mutex is
released.
pkg/batcher/shutdown.go (1)

51-55: 🎯 Functional Correctness | 🔵 Trivial | 💤 Low value

Use a shallow sentinel check in Is.

errors.Is(target, ErrTimeout) checks whether target wraps ErrTimeout, so *ShutdownIncompleteError also matches wrapped timeout errors. Use target == ErrTimeout, then remove the errors import if unused.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@pkg/batcher/shutdown.go` around lines 51 - 55, Update
ShutdownIncompleteError.Is to use a direct target == ErrTimeout comparison,
limiting equivalence to the exact sentinel; remove the errors import if it
becomes unused.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Inline comments:
In `@docs/improvements/plan-perf.md`:
- Around line 1044-1046: Update the Go version in thresholds.md’s keyed
reference environment baseline from 1.22.4 to 1.25.0 to match go.mod, and ensure
any newly introduced CI tooling remains gated for Go 1.25.0 or newer.

In `@docs/improvements/thresholds.md`:
- Around line 30-38: Make the Phase 3 goroutine requirements in the thresholds
document consistent with the table’s “exactly 2 goroutines per running batcher”
gate. Update the Phase 3 table heading and rows around the existing goroutine
requirements, or explicitly scope them as Phase 3-only without contradicting the
global gate; remove the conflicting 1 and 1+n requirements.

In `@pkg/batcher/batcher.go`:
- Around line 477-486: Update the shutdown wait logic in Shutdown to check
b.shutdownDone before entering the select, then retain the existing select for
cases where completion has not yet occurred. Return nil whenever the drain is
already complete, including when ctx.Done() is simultaneously ready, while
preserving the existing ShutdownIncompleteError behavior for an incomplete
drain.

In `@pkg/batcher/goroutine_leak_test.go`:
- Around line 347-350: Replace the Stats().Pending polling in the in-flight
timeout test with a processor-owned entry signal, and wait on that signal before
calling Close(). Update the test’s processor setup to signal immediately when
processor execution begins, while preserving the existing timeout and
assertions.

In `@pkg/batcher/options.go`:
- Around line 76-77: Update the godoc comment for WithErrorBufferSize to clearly
state the diagnostic buffer capacity and separately explain that, once full, the
oldest errors are retained while newer errors are dropped.

In `@pkg/batcher/queue.go`:
- Around line 113-148: Update queue.pop to compact q.items when q.head becomes
sufficiently large: move the unconsumed suffix to the front, reslice to the
remaining count, and reset q.head to zero while preserving item order and
capacity. Keep the existing fully-drained reset and signaling behavior, and
ensure push/tryPush continue enforcing the queue bound using the compacted
state.

In `@pkg/batcher/shutdown_test.go`:
- Around line 237-241: Update the shutdown completion select around shutdownDone
to receive and assert the returned error rather than discarding it. Fail the
trial when Shutdown reports a ShutdownIncompleteError, while preserving the
existing timeout failure for a shutdown that does not complete.

In `@pkg/batcher/stats.go`:
- Around line 5-20: Correct the Stats documentation to remove the claim that
snapshots never take a lock or are entirely allocation-free, since Queued is
populated via queue.length(). State that reading Queued acquires the queue mutex
and may contend with publishers and the aggregator, while preserving the
existing O(1), eventually consistent snapshot description.

In `@README.md`:
- Around line 174-187: Update the README shutdown example around batcher.Close
so it handles the returned incomplete-drain error instead of discarding it
through defer. Show explicit Close error handling or propagate the deferred
result, while preserving the documented behavior that remaining work continues
in the background after the grace period.
- Around line 220-223: Update the batcher.New example to pass the declared
processor’s Process method to batcher.WithProcessor instead of passing the
processor struct, matching the processor function contract and the earlier
usage.

---

Nitpick comments:
In `@pkg/batcher/admission_test.go`:
- Around line 226-259: Register test cleanup for the batcher instance after
defer close(release) so the processor is unblocked before b.Close waits for
draining; apply the same lifecycle cleanup and ensure the Close goroutine is
awaited in TestBlockedEnqueueIsReleasedByShutdown.
- Around line 31-35: Reduce repeated race-test cost by making the shared trials
value in pkg/batcher/admission_test.go lines 31-35 short-test aware, using a
small count when testing.Short() is enabled while preserving the current
default. Apply the same testing.Short() reduction to the 300-trial loop at
pkg/batcher/admission_test.go lines 110 and 145, the 300-trial loop at
pkg/batcher/shutdown_test.go line 190, and the 200-trial loop at
pkg/batcher/shutdown_test.go line 284; leave normal trial counts unchanged.

In `@pkg/batcher/batcher.go`:
- Around line 166-171: Rename the nonBlockingWhenFull parameter and all
references in the surrounding Add flow to unboundedFastPath, preserving the
existing condition and behavior so the fast path is selected only when
MaxQueueSize <= 0 while bounded queues continue using push.
- Around line 553-563: Update Batcher.finish to remove the unreachable
stateSealing fallback and its misleading comment, or make any unexpected state
transition fail loudly instead of returning silently with errorsChan open.
Preserve the successful stateDraining-to-stateClosed transition and errorsChan
closure.

In `@pkg/batcher/characterization_test.go`:
- Around line 387-404: Remove the hand-rolled itoa helper and update keyFor to
use strconv.Itoa directly for integer conversion. Add the strconv import and
preserve keyFor’s existing “item-” prefix and output behavior.
- Around line 105-117: The timing assertion in the first-item flush test should
enforce that the batch does not flush before approximately the full interval
measured from Add of the first item. Update the require.GreaterOrEqual check
around start, b.Add, and elapsed to use interval minus a small, explicit
tolerance, while preserving the existing batch-count and item-count assertions.

In `@pkg/batcher/gate.go`:
- Around line 124-127: Remove the unused admissionGate.isSealed method from
pkg/batcher/gate.go, including its inaccurate doc comment, rather than leaving
an unreferenced helper that fails lint.

In `@pkg/batcher/queue.go`:
- Around line 60-92: Update queue.push so that after a successful append to a
bounded queue, it re-signals notFull when remaining queued items are still below
capacity. Preserve the existing unbounded-queue behavior and only re-arm the
signal after the item is appended and the mutex is released.

In `@pkg/batcher/shutdown.go`:
- Around line 51-55: Update ShutdownIncompleteError.Is to use a direct target ==
ErrTimeout comparison, limiting equivalence to the exact sentinel; remove the
errors import if it becomes unused.
🪄 Autofix

Fix all unresolved CodeRabbit comments on this PR:

  • Push a commit to this branch (recommended)
  • Create a new PR with the fixes

ℹ️ Review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: CHILL

Plan: Pro Plus

Run ID: 2e7f2a56-4e28-44fe-accd-fae3e81e47ed

📥 Commits

Reviewing files that changed from the base of the PR and between 32da402 and 8f1dfeb.

⛔ Files ignored due to path filters (1)
  • go.sum is excluded by !**/*.sum
📒 Files selected for processing (21)
  • README.md
  • docs/improvements/plan-perf.md
  • docs/improvements/thresholds.md
  • go.mod
  • pkg/batcher/admission_test.go
  • pkg/batcher/batcher.go
  • pkg/batcher/batcher_test.go
  • pkg/batcher/characterization_test.go
  • pkg/batcher/constants.go
  • pkg/batcher/diagnostics_test.go
  • pkg/batcher/fx.go
  • pkg/batcher/fx_test.go
  • pkg/batcher/gate.go
  • pkg/batcher/gate_internal_test.go
  • pkg/batcher/goroutine_budget_test.go
  • pkg/batcher/goroutine_leak_test.go
  • pkg/batcher/options.go
  • pkg/batcher/queue.go
  • pkg/batcher/shutdown.go
  • pkg/batcher/shutdown_test.go
  • pkg/batcher/stats.go

Comment thread docs/improvements/plan-perf.md
Comment thread docs/improvements/thresholds.md
Comment thread pkg/batcher/batcher.go
Comment thread pkg/batcher/goroutine_leak_test.go Outdated
Comment thread pkg/batcher/options.go Outdated
Comment thread pkg/batcher/queue.go
Comment thread pkg/batcher/shutdown_test.go
Comment thread pkg/batcher/stats.go
Comment thread README.md
Comment thread README.md
heynemann added a commit that referenced this pull request Aug 6, 2026
Fixes all nine unresolved review threads on PR #29.

Correctness and memory:

- Shutdown now prefers a completed drain over an already-expired context. A select
  can choose either when both are ready, so the old code could report
  ShutdownIncompleteError with Pending=0 for a batcher already closed.
- The unbounded queue now compacts its consumed prefix when it dominates the live
  suffix. This is not cosmetic: a queue kept at depth 1 grew to cap=219,136 after
  200k push/pop cycles because append grew with total throughput while the
  full-drain reset never fired. It now stays cap=2. Tests pin the memory bound and
  FIFO order across compaction.
- The close-timeout test now waits on a processor-owned signal rather than
  Stats.Pending, so it genuinely covers a processor already in flight.
- The parked-publisher shutdown test now asserts Shutdown returns nil rather than
  merely asserting it returns; a ShutdownIncompleteError no longer passes as a
  false success.

Documentation and examples:

- Stats no longer claims it never locks: Queued takes the queue mutex for one
  length check. It remains O(1) and allocation-free, but scraping in a tight loop
  can contend with push/pop.
- Fixed the garbled WithErrorBufferSize godoc.
- README examples now pass processor.Process, not the Processor struct, and both
  Close examples handle the returned incomplete-drain error instead of discarding
  it with a bare defer.
- Reconciled contradictory goroutine gates: Phase 2 enforces 2 (aggregator +
  serial processor); the n>1 1+n row is deferred until the Phase 3 code exists.
@heynemann
heynemann force-pushed the phase-2-lifecycle-safety branch from 6505fef to c58d01c Compare August 6, 2026 16:22
@coderabbitai

coderabbitai Bot commented Aug 6, 2026

Copy link
Copy Markdown

Note

GitHub couldn't provide a complete incremental comparison for this pull request, so CodeRabbit is performing a full review instead. This review may take a little longer.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 3

🧹 Nitpick comments (1)
docs/improvements/thresholds.md (1)

54-59: 🚀 Performance & Scalability | 🔵 Trivial | ⚡ Quick win

Document the queue capacity contract.

admissionGate tracks publishers; it does not limit queued items. State that MaxQueueSize <= 0 selects the unbounded path used by this allocation gate, while positive values bound the queue. In bounded mode, Add blocks when full and Enqueue returns ctx.Err() or ErrClosing when it cannot accept an item.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@docs/improvements/thresholds.md` around lines 54 - 59, Update the queue
capacity documentation near admissionGate to state that it tracks publishers
rather than limiting queued items. Document that MaxQueueSize <= 0 selects the
unbounded allocation-gate path, while positive values bound the queue; specify
that bounded Add blocks when full and Enqueue returns ctx.Err() or ErrClosing
when it cannot accept an item.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Inline comments:
In @.github/workflows/go.yml:
- Around line 110-112: The comment above go-version-file must accurately
describe the change: replace the claim about a hardcoded floor or pinned Go
version with wording that explains removing version drift and aligning CI with
the Go version declared in go.mod. Keep the go-version-file setting unchanged.
- Around line 110-112: Update the setup-go configuration in the build, test, and
coverage jobs to use go-version-file: go.mod, ensuring each job derives its Go
version from the module directive instead of the runner’s pre-installed version.

In `@pkg/batcher/admission_test.go`:
- Around line 257-258: The assertion for a cancelled Enqueue in the relevant
admission test incorrectly expects Rejected to increase. Update it to match the
accounting contract: a full-queue deadline rolls back the reservation without
incrementing Rejected, while preserving the surrounding deadline behavior
assertions.

---

Nitpick comments:
In `@docs/improvements/thresholds.md`:
- Around line 54-59: Update the queue capacity documentation near admissionGate
to state that it tracks publishers rather than limiting queued items. Document
that MaxQueueSize <= 0 selects the unbounded allocation-gate path, while
positive values bound the queue; specify that bounded Add blocks when full and
Enqueue returns ctx.Err() or ErrClosing when it cannot accept an item.
🪄 Autofix

Fix all unresolved CodeRabbit comments on this PR:

  • Push a commit to this branch (recommended)
  • Create a new PR with the fixes

ℹ️ Review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: CHILL

Plan: Pro

Run ID: 7e3e9fb2-d1dd-410f-847a-02c48669d4cf

📥 Commits

Reviewing files that changed from the base of the PR and between 3f48a40 and c58d01c.

⛔ Files ignored due to path filters (1)
  • go.sum is excluded by !**/*.sum
📒 Files selected for processing (24)
  • .github/workflows/bench.yml
  • .github/workflows/go.yml
  • README.md
  • docs/improvements/plan-perf.md
  • docs/improvements/thresholds.md
  • go.mod
  • pkg/batcher/admission_test.go
  • pkg/batcher/batcher.go
  • pkg/batcher/batcher_test.go
  • pkg/batcher/characterization_test.go
  • pkg/batcher/constants.go
  • pkg/batcher/diagnostics_test.go
  • pkg/batcher/fx.go
  • pkg/batcher/fx_test.go
  • pkg/batcher/gate.go
  • pkg/batcher/gate_internal_test.go
  • pkg/batcher/goroutine_budget_test.go
  • pkg/batcher/goroutine_leak_test.go
  • pkg/batcher/options.go
  • pkg/batcher/queue.go
  • pkg/batcher/queue_test.go
  • pkg/batcher/shutdown.go
  • pkg/batcher/shutdown_test.go
  • pkg/batcher/stats.go
🚧 Files skipped from review as they are similar to previous changes (14)
  • pkg/batcher/batcher_test.go
  • pkg/batcher/goroutine_leak_test.go
  • pkg/batcher/options.go
  • pkg/batcher/shutdown.go
  • pkg/batcher/fx_test.go
  • pkg/batcher/gate_internal_test.go
  • pkg/batcher/stats.go
  • docs/improvements/plan-perf.md
  • pkg/batcher/fx.go
  • pkg/batcher/constants.go
  • go.mod
  • pkg/batcher/shutdown_test.go
  • pkg/batcher/characterization_test.go
  • pkg/batcher/batcher.go

Comment thread .github/workflows/go.yml Outdated
Comment thread pkg/batcher/admission_test.go
@coderabbitai

coderabbitai Bot commented Aug 7, 2026

Copy link
Copy Markdown

Note

GitHub couldn't provide a complete incremental comparison for this pull request, so CodeRabbit is performing a full review instead. This review may take a little longer.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 2

🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Inline comments:
In `@pkg/batcher/gate.go`:
- Around line 124-127: Resolve the unused admissionGate.isSealed method by
removing it if no current code requires sealed-state checks; otherwise, add its
intended caller and ensure the returned state is used. Do not leave the method
unused.

In `@README.md`:
- Around line 205-218: Rename the batcher instance used in the README example
from batcher to a distinct identifier such as b, updating both Shutdown calls
while retaining batcher.ShutdownIncompleteError for the imported package type.
🪄 Autofix

Fix all unresolved CodeRabbit comments on this PR:

  • Push a commit to this branch (recommended)
  • Create a new PR with the fixes

ℹ️ Review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: CHILL

Plan: Pro

Run ID: 63d9123a-fc66-40cb-9c60-871cc13810c3

📥 Commits

Reviewing files that changed from the base of the PR and between 8b6d33c and e0353a5.

⛔ Files ignored due to path filters (1)
  • go.sum is excluded by !**/*.sum
📒 Files selected for processing (23)
  • .github/workflows/bench.yml
  • README.md
  • docs/improvements/plan-perf.md
  • docs/improvements/thresholds.md
  • go.mod
  • pkg/batcher/admission_test.go
  • pkg/batcher/batcher.go
  • pkg/batcher/batcher_test.go
  • pkg/batcher/characterization_test.go
  • pkg/batcher/constants.go
  • pkg/batcher/diagnostics_test.go
  • pkg/batcher/fx.go
  • pkg/batcher/fx_test.go
  • pkg/batcher/gate.go
  • pkg/batcher/gate_internal_test.go
  • pkg/batcher/goroutine_budget_test.go
  • pkg/batcher/goroutine_leak_test.go
  • pkg/batcher/options.go
  • pkg/batcher/queue.go
  • pkg/batcher/queue_test.go
  • pkg/batcher/shutdown.go
  • pkg/batcher/shutdown_test.go
  • pkg/batcher/stats.go
🚧 Files skipped from review as they are similar to previous changes (21)
  • pkg/batcher/fx_test.go
  • pkg/batcher/goroutine_leak_test.go
  • pkg/batcher/queue.go
  • .github/workflows/bench.yml
  • pkg/batcher/shutdown.go
  • pkg/batcher/constants.go
  • docs/improvements/thresholds.md
  • pkg/batcher/stats.go
  • pkg/batcher/diagnostics_test.go
  • pkg/batcher/goroutine_budget_test.go
  • go.mod
  • pkg/batcher/options.go
  • pkg/batcher/fx.go
  • docs/improvements/plan-perf.md
  • pkg/batcher/queue_test.go
  • pkg/batcher/shutdown_test.go
  • pkg/batcher/gate_internal_test.go
  • pkg/batcher/characterization_test.go
  • pkg/batcher/admission_test.go
  • pkg/batcher/batcher.go
  • pkg/batcher/batcher_test.go

Comment thread pkg/batcher/gate.go Outdated
Comment thread README.md
@heynemann
heynemann force-pushed the phase-2-lifecycle-safety branch from c89b39d to c578f61 Compare August 7, 2026 13:53

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 1

Caution

Some comments are outside the diff and can’t be posted inline due to platform limitations.

⚠️ Outside diff range comments (1)
docs/improvements/plan-perf.md (1)

508-510: 📐 Maintainability & Code Quality | 🟠 Major | ⚡ Quick win

Align the Phase 2.2 goroutine budget.

The plan defines incompatible budgets. Line 509 and Line 706 require one goroutine for a running n=1 batcher. Lines 1163-1165 report two goroutines, which also matches the PR objective. Update the Phase 2.2 threshold and acceptance criteria to two, and clarify whether the single-goroutine target applies only after Phase 3.

This conflict can make lifecycle tests reject the intended implementation or accept the wrong budget.

Also applies to: 589-593, 705-707, 1163-1167

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@docs/improvements/plan-perf.md` around lines 508 - 510, Update the Phase 2.2
goroutine budget references in the plan so a running n=1 batcher consistently
requires two goroutines, including the threshold and acceptance criteria at the
referenced sections. Explicitly state that the one-goroutine target applies only
after Phase 3, while preserving the zero-idle and post-closed baseline
requirements.
🧹 Nitpick comments (1)
docs/improvements/plan-perf.md (1)

164-164: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win

Add language identifiers to fenced blocks.

markdownlint-cli2 reports MD040 for these fences. Add text to pseudocode and ASCII blocks, and go where the block contains Go code.

Also applies to: 208-208, 229-229, 284-284, 321-321, 375-375, 1137-1137

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@docs/improvements/plan-perf.md` at line 164, Add explicit language
identifiers to every fenced code block called out in plan-perf.md: use text for
pseudocode and ASCII diagrams, and go for blocks containing Go code. Update all
listed fence locations so markdownlint MD040 passes without changing the block
contents.

Source: Linters/SAST tools

🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Inline comments:
In `@README.md`:
- Around line 238-241: Update the error-handling documentation around b.Enqueue
to describe both possible context errors: context.Canceled and
context.DeadlineExceeded, while retaining batcher.ErrClosing as the shutdown
case.

---

Outside diff comments:
In `@docs/improvements/plan-perf.md`:
- Around line 508-510: Update the Phase 2.2 goroutine budget references in the
plan so a running n=1 batcher consistently requires two goroutines, including
the threshold and acceptance criteria at the referenced sections. Explicitly
state that the one-goroutine target applies only after Phase 3, while preserving
the zero-idle and post-closed baseline requirements.

---

Nitpick comments:
In `@docs/improvements/plan-perf.md`:
- Line 164: Add explicit language identifiers to every fenced code block called
out in plan-perf.md: use text for pseudocode and ASCII diagrams, and go for
blocks containing Go code. Update all listed fence locations so markdownlint
MD040 passes without changing the block contents.
🪄 Autofix

Fix all unresolved CodeRabbit comments on this PR:

  • Push a commit to this branch (recommended)
  • Create a new PR with the fixes

ℹ️ Review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: CHILL

Plan: Pro

Run ID: 3716646d-393b-4d68-96fe-b336c7568872

📥 Commits

Reviewing files that changed from the base of the PR and between e0353a5 and c578f61.

📒 Files selected for processing (5)
  • README.md
  • docs/improvements/plan-perf.md
  • pkg/batcher/admission_test.go
  • pkg/batcher/gate.go
  • pkg/batcher/options.go
🚧 Files skipped from review as they are similar to previous changes (3)
  • pkg/batcher/gate.go
  • pkg/batcher/options.go
  • pkg/batcher/admission_test.go

Comment thread README.md
First commit of Milestone 2.1. Pins observable behaviour of the rill-backed
engine so the upcoming in-repo aggregation loop cannot silently change semantics.
These tests must keep passing unchanged after the swap.

Locks down eight behaviours that no existing test asserted:

- The batch interval is not a periodic tick. The timer starts when the first item
  enters an empty batch, so a sparse producer waits the full interval per item.
  This is the property the entire performance plan is built around; an
  implementation using a periodic ticker would change latency for every user.
- An idle batcher never invokes the processor, and the processor never receives an
  empty batch.
- A full batch flushes on size without waiting for the interval.
- Batch boundaries are exactly BatchSize until the remainder, and a single
  producer's items are processed in arrival order. Concurrent producers have no
  defined relative order and the test deliberately does not claim one.
- Shutdown flushes an under-full batch rather than discarding it.
- A processor error is published on Errors() without stopping the pipeline, and
  the failed batch still decrements pending work so Join can complete.
- Errors() is closed on shutdown, so a ranging consumer terminates.
- Len counts accepted work that has not finished processing, including work handed
  to a blocked processor.

The recording processor copies each batch, because the engine owns the slice it
passes to the processor and a later implementation may reuse its backing array.
Second commit of Milestone 2.1. Removes github.com/destel/rill and brings batch
aggregation in-repo, with no observable behaviour change: all eight
characterization tests from the previous commit pass unchanged.

The pipeline is now two owned stages plus a forwarder:

- forward() drains the unbounded input queue onto an unbuffered channel.
- aggregate() groups items, arming the interval timer only when the first item
  enters an empty batch, flushing on size, on timer, or on end of input, and
  never emitting an empty batch.
- startProcessing() consumes batches and invokes the processor, preserving the
  existing teardown order: close input so aggregation flushes, drain batches so
  aggregation can exit, close errors last.

Getting here required correcting two wrong assumptions in the plan, both caught
by measurement rather than review:

- Merging aggregation into the processing loop changes latency. With a 50ms
  processor at 10k items/s it inverted the documented baseline (5ms window p50
  28ms versus 100ms window p50 75ms), because the next batch only began filling
  after the processor returned. That is a Phase 3 decision, so the stages stay
  separate.
- Merging input draining into aggregation regressed sequential Add by 39-50%
  against the stored baseline, breaching the predeclared 10% gate. The unbounded
  queue has a small bounded ingress that must be drained promptly, and the
  aggregator interleaves reads with timer and slice work. A dedicated forwarding
  stage restored parity. This is what rill's ToChans/Batch/FromChans pipeline was
  buying.

Measured result: goroutines per running batcher 6 to 5 with none leaked, no
sec/op regression against the baseline (one case 6.7% faster), and 47% fewer
bytes per enqueued item because items are no longer wrapped in rill.Try[T].
Batch counts and mean batch sizes match rill's within noise, and the slow
processor inversion is preserved.

Also modernises the dependency stack: testify 1.9.0 to 1.11.1, zap 1.27.0 to
1.28.0, multierr and x/sys refreshed, and go.uber.org/atomic dropped in favour of
sync/atomic typed values, which the standard library has provided since Go 1.19.
Updating x/sys to v0.46.0 raises the module's minimum Go version to 1.25.0; that
is recorded in the Phase 5 compatibility inventory.
Third commit of Milestone 2.1. The milestone's headline claim is that removing
rill reduces owned goroutines per running batcher from 6 to 5 without changing
behaviour, and until now nothing enforced it.

Asserts the budget directly, naming each of the five goroutines: the input queue
relay, the input forwarder, the aggregator, the processing loop, and the error
queue relay. It also asserts the count returns to the pre-construction baseline
after Close, since a long-lived process that creates batchers would otherwise leak
one set per batcher.

Verified the test detects drift in both directions by temporarily setting the
constant to 4, which failed with a clear message rather than passing silently.
Goroutine counts are read after waiting for the count to stabilise, so a
scheduler that has not yet reaped finished goroutines is not mistaken for a leak.

Phase 2.2 will lower this to 2 when the queue relays and the forwarder are
removed; the constant and the plan's goroutine table are meant to change together.
Milestone 2.2. Replaces golang.design/x/chann with owned queues, introduces the
publisher admission gate, and makes shutdown non-destructive. This is the
milestone that removes the crash and data-loss paths.

The central change is that the intake queue is NEVER closed. chann's relay
goroutine can only exit by closing its own ingress channel, which is the very
channel publishers send to, so "never close the input" and "leak no goroutines"
were mutually exclusive. Owning the queue resolves that: shutdown is signalled by
sealing admission and by intake accounting instead.

What this fixes, each previously reproducible:

- Add after Close panicked with "send on closed channel", crashing the caller's
  goroutine during a graceful shutdown. It is now a counted rejection.
- Close abandoned the drain when its timeout expired, silently discarding accepted
  work. The drain now runs in an independent coordinator, so a caller that gives
  up reports an incomplete drain rather than causing one.
- Start was not idempotent; a second call raced two loops whose deferred cleanup
  closed the same channels.

Design notes worth carrying forward:

- Publishers reserve their drain obligation BEFORE publishing and roll it back
  before leaving the gate. Reserving after would let the drain conclude while a
  publisher was mid-publish; rolling back after leaving would leave a phantom
  obligation the drain waits on forever.
- Every gate decrement routes through leave(), including enter()'s own rejection
  path, and the coordinator performs its own post-seal gate check. These cover two
  distinct lost-wakeup interleavings, and both now have tests proven to hang when
  the mechanism is sabotaged.
- Aggregation and processing remain separate goroutines. Merging them inverted the
  documented latency baseline again during this work, exactly as in 2.1, and was
  reverted.

Adds Enqueue(ctx, item) error for callers who need to observe rejection or bound
their wait, WithMaxQueueSize for back-pressure (default stays unbounded),
WithCloseGrace with a flat 30s default replacing the old queue-size formula, and a
minimal Stats() snapshot that is O(1) and allocation-free.

Measured: goroutines per running batcher 6 -> 2, enqueue -54% geomean sec/op
(229.8ns -> 63.4ns at small batch sizes) and -49% to -90% bytes/op, because
removing the relays removed a channel hop and a goroutine handoff per item. All
characterization tests pass unchanged, and the Phase 1 inversion baseline still
holds.

TestBatcher_CloseTimeoutBehavior was updated to the new contract: it now asserts
Close is bounded by its grace period, reports ErrTimeout, and leaves the batcher
draining rather than falsely claiming to be closed.
Milestone 2.3. Exposes the drain coordinator built in 2.2 as a public, resumable
API and fixes a deadlock found while testing it.

Shutdown(ctx) seals admission and waits for accepted work to drain. If ctx expires
it returns *ShutdownIncompleteError and the drain continues; calling Shutdown again
with a fresh context waits on that same drain rather than starting a new one. There
is no continuation token because none is needed: the Batcher is the continuation,
sealed against new work from the first call onward. A canceled context cannot be
revived, so a new context is the honest way to express "wait longer".

Per-caller waits are independent. A caller with a short deadline gets its own
timeout error, never stored as the terminal result, so one impatient caller cannot
poison another's answer. Processor errors and recovered panics stay on Errors(),
because a shutdown result describes the drain, not the work.

ShutdownIncompleteError documents that Pending is a conservative drain obligation
rather than a queue depth: it includes items inside a processor call and publishers
that have reserved but may yet abort, and is exact only once PublishersInGate is
zero. It unwraps to the context error and matches ErrTimeout, so existing callers
written against Close keep working.

Fixed while testing: Start racing Shutdown could deadlock. Start declined to
launch the consumer when it observed a sealed gate, but it shared startOnce with
the drain path, so whichever call won the race left the drain waiting on a consumer
that would never exist. Start now always launches the consumer, because a sealed
batcher may still hold queued work. TestStartRacingShutdownYieldsOneLifecycle
covers it across 200 trials.

fx.go now forwards the stop-hook context to Shutdown instead of discarding it, so
the application's shutdown deadline governs the drain. Previously an app with a
longer grace period could not use it and one with a shorter deadline could not
bound it.

Also pins the original data-loss case directly: a partial batch with a 30s interval
and a shorter grace period is flushed rather than discarded. That configuration
previously accepted 50 items, processed 0, and still reported 50 pending.
Milestone 2.4, completing Phase 2.

Processor panics are now recovered and scoped to the batch. Batcher owns goroutines
the caller cannot reach, so an unrecovered panic there crashed the whole process and
destroyed every other queued item, and the caller had no way to install its own
recover. The panic value and stack are reported as a ProcessorPanicError so the bug
stays debuggable rather than being silently swallowed. Recovery is scoped to the
batch rather than the consumer loop, so one poison batch cannot take the consumer
down and stall every later batch.

The defer ordering is fixed and load-bearing: whatever happens, exactly one terminal
category is counted and the drain obligation is released exactly once. Leaving it
unreleased would inflate Pending permanently and hang shutdown forever. Verified by
sabotage: removing the recover makes TestProcessorPanicIsReportedAndPipelineSurvives
crash with "panic: poison batch".

Diagnostics are bounded, and the trade is deliberate because both alternatives are
worse. A blocking send would deadlock the pipeline whenever a processor fails and
nobody reads Errors() — a self-inflicted outage caused by a diagnostic. An unbounded
channel would grow without limit during an outage, exactly when memory is scarcest.
Dropping the newest error keeps the earliest errors of a storm, which usually carry
the root cause, and avoids racing the consumer. Stats().DroppedErrors makes the loss
visible and is itself worth alerting on, since a non-zero value means Errors() is
not being drained.

Terminal outcomes stay mutually exclusive: a panic counts as Panicked and never as
Failed, and a recovered panic does not make the shutdown result fail.

Adds WithErrorBufferSize (default 1024). Confirmed the non-panic path still
allocates nothing per Add despite the guard, so the safety feature does not tax
every batch.
Updates the plan and thresholds with what Phase 2 actually delivered, verified
against the original defects rather than restated from the plan:

- partial batch with a 30s interval: 50 accepted / 50 processed / 0 pending, where
  it previously reported 50 accepted / 0 processed / 50 phantom pending;
- Add after Close: counted rejection instead of a process panic;
- bounded mode: queue capped and rejections reported, instead of unbounded heap
  growth;
- goroutines per running batcher: 2, none leaked.

Also corrects the threshold table to name the tests that actually enforce each gate,
and drops the placeholder Stats() allocation gate: Stats returns a value type built
from a fixed set of atomic loads, so the Add gate is the one that matters.
The README still told users to call StopProcessing, which does not exist, and said
nothing about the APIs Phase 2 added. Left alone it would actively mislead: the
lifecycle contract changed, and the queue is unbounded by default.

Documents Close's new non-destructive semantics, resumable Shutdown with a worked
ShutdownIncompleteError example, bounded queues with Enqueue for back-pressure, and
Stats for observability. States plainly that shrinking the batch interval does not
bound queued work, since that is the assumption the plan was written to correct.

The full API inventory and migration guide remain Phase 5 work; this fixes the
parts that are wrong today.
Fixes all nine unresolved review threads on PR #29.

Correctness and memory:

- Shutdown now prefers a completed drain over an already-expired context. A select
  can choose either when both are ready, so the old code could report
  ShutdownIncompleteError with Pending=0 for a batcher already closed.
- The unbounded queue now compacts its consumed prefix when it dominates the live
  suffix. This is not cosmetic: a queue kept at depth 1 grew to cap=219,136 after
  200k push/pop cycles because append grew with total throughput while the
  full-drain reset never fired. It now stays cap=2. Tests pin the memory bound and
  FIFO order across compaction.
- The close-timeout test now waits on a processor-owned signal rather than
  Stats.Pending, so it genuinely covers a processor already in flight.
- The parked-publisher shutdown test now asserts Shutdown returns nil rather than
  merely asserting it returns; a ShutdownIncompleteError no longer passes as a
  false success.

Documentation and examples:

- Stats no longer claims it never locks: Queued takes the queue mutex for one
  length check. It remains O(1) and allocation-free, but scraping in a tight loop
  can contend with push/pop.
- Fixed the garbled WithErrorBufferSize godoc.
- README examples now pass processor.Process, not the Processor struct, and both
  Close examples handle the returned incomplete-drain error instead of discarding
  it with a bare defer.
- Reconciled contradictory goroutine gates: Phase 2 enforces 2 (aggregator +
  serial processor); the n>1 1+n row is deferred until the Phase 3 code exists.
Addresses the review finding that thresholds.md still declared Go 1.22.4 as its
reference environment while go.mod moved to 1.25.0 in Phase 2.

That row is not cosmetic: the document states thresholds are "keyed to the
environment", so a stale Go version silently invalidates every stored baseline it
governs. It now reads 1.25.0 and records that CI resolves the toolchain from
go.mod, so the row and the toolchain cannot drift apart again.

Also gated the remaining CI tooling on the module version rather than a hardcoded
string, which is the second half of the finding:

- go.yml's dependency-submission job requested ">=1.22.0", which is below the
  module floor and could select a toolchain unable to build the module.
- bench.yml requested "stable", which silently changes benchmark history whenever
  a new Go release appears, with no repository change to attribute it to.

Both now use go-version-file: go.mod, matching performance.yml. All four setup-go
steps in the repository are now derived from the module.

Fixed a flaky assertion surfaced while validating this change.
TestShutdownWithParkedPublisherOnFullBoundedQueue failed on trial 33 of 300
because it required ErrClosing from a publisher that had entered the admission
gate BEFORE sealing. Once shutdown starts the consumer and a queue slot frees,
both sealCh and notFull are ready and the select may choose either, so that
publisher can legitimately succeed. The protocol treats a pre-seal publisher that
slips through as benign precisely because the drain still accounts for it. The
test now asserts the real contract: such a publisher must not deadlock, and may
return nil or ErrClosing. Enqueue calls that start after sealing are still
required to return ErrClosing by TestEnqueueAfterSealReportsClosing.
Two review findings on Phase 2.

admissionGate.isSealed was dead code. The gate's own enter/leave path reads
g.sealed directly and the coordinator uses checkEmpty, so nothing called it;
golangci-lint flagged it as unused. Removed rather than given a synthetic caller.

The README Shutdown example could not compile. It used `batcher` as both the
instance receiver and the package qualifier in the same scope, so
`*batcher.ShutdownIncompleteError` referred to a variable rather than the package.
Verified with a minimal reproduction: "batcher.Err is not a type". The instance is
now `b`, matching every other example in the file, and `batcher` stays the package.
The Phase 2 review found the capacity bound understated by a full batch, and
measurement confirms it. With MaxQueueSize=4, BatchSize=3 and a blocked processor,
Pending reaches 10 while the doc claimed N + BatchSize = 7.

The cause is that two batches exist outside the queue at once: the aggregator holds
a partial batch that has already left it, and another batch is inside the processor.
The real bound is N + 2*BatchSize + publishers-in-gate at Concurrency 1, or
N + (1 + concurrency) * BatchSize + gate in general. Writing it as
BatchSize + concurrency * BatchSize drops the held batch, which at n=1 is a
factor-of-two error on that term -- and this is the number callers use to size
memory, so understating it is the wrong direction to be wrong in.

Corrected WithMaxQueueSize's godoc and the plan's capacity contract, and added
TestAcceptedWorkBoundIncludesHeldAndInFlightBatches. It asserts both halves: work
stays within the corrected bound, and it genuinely exceeds N + BatchSize -- so if
someone reverts the doc, the test explains why the smaller figure describes a
batcher that cannot hold a batch and process one simultaneously.

The test releases its blocked processor before Close in cleanup; without that
ordering it spent the full 30s close grace waiting on a processor it was about to
abandon.
@heynemann
heynemann force-pushed the phase-2-lifecycle-safety branch from c578f61 to 875e0d1 Compare August 8, 2026 00:50
@heynemann
heynemann merged commit 7dfe335 into main Aug 8, 2026
16 checks passed
@heynemann
heynemann deleted the phase-2-lifecycle-safety branch August 8, 2026 02:20
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant