Repository navigation
feat: unlock small batch windows with explicit concurrency (Phase 3) - #31
Conversation
|
Important Review skippedAuto reviews are disabled on base/target branches other than the default branch. Please check the settings in the CodeRabbit UI or the ⚙️ Run configurationConfiguration used: Organization UI Review profile: CHILL Plan: Pro Run ID: You can disable this status message by setting the Use the checkbox below for a quick retry:
📝 WalkthroughWalkthroughThe batcher now defaults to serial processing and supports acknowledged concurrent workers. It validates concurrency settings, tracks in-flight items, drains workers during shutdown, preserves per-batch ordering, and adds concurrency tests, scenario measurements, and updated performance thresholds. ChangesConcurrent batch processing
Estimated code review effort: 4 (Complex) | ~45 minutes Sequence Diagram(s)sequenceDiagram
participant Producer
participant Aggregator
participant WorkerPool
participant Processor
participant Stats
Producer->>Aggregator: publish items
Aggregator->>WorkerPool: dispatch a batch
WorkerPool->>Processor: process batch items
Processor-->>WorkerPool: return result or panic
WorkerPool->>Stats: update completion and in-flight counts
WorkerPool-->>Aggregator: finish worker processing
Possibly related PRs
🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
✨ Finishing Touches 💡 1📝 Generate docstrings 💡
🧪 Generate unit tests (beta)
Comment |
58ba047 to
77153e1
Compare
37d900b to
abca783
Compare
abca783 to
a80efb5
Compare
There was a problem hiding this comment.
Actionable comments posted: 6
🤖 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 911-913: Update the goroutine budget in the Phase 2 discussion of
the improvement plan so the n=1 case is stated as two goroutines: one aggregator
and one serial processor. Align the surrounding text with the corresponding
budget in thresholds.md and preserve the explanation that the separate serial
processor remains for latency semantics.
In `@docs/improvements/thresholds.md`:
- Around line 71-80: Update the “Goroutines after terminal closed” threshold in
the table to “at most the pre-construction baseline” (or equivalent <= wording),
matching both tests’ settled <= baseline assertions; retain the existing test
references and avoid claiming exact equality.
In `@pkg/batcher/worker_test.go`:
- Around line 238-240: Update pkg/batcher/worker_test.go lines 238-240 and
294-301 to register processor-release teardown with t.Cleanup. In both affected
tests, guard close(release) with sync.Once so explicit test-body release and
cleanup cannot double-close; additionally register b.Close() with t.Cleanup at
lines 238-240, while the lines 294-301 site requires only the guarded release
cleanup.
- Around line 372-376: Update the sequential tests in batcher_test.go that call
batcher.New to close every created batcher before returning, including cleanup
for assertion or setup failures where applicable. Ensure all auto-started
batchers are stopped before TestWorkerGoroutineBudget captures its baseline via
settledGoroutines, preserving the exact goroutine-delta assertion.
In `@test/scenario/concurrency_test.go`:
- Around line 52-61: Update the listed overload scenarios to set
Config.Producers to runtime.NumCPU(): test/scenario/concurrency_test.go:52-61,
121-129, 181-206, and 226-236, plus test/scenario/matrix_test.go:95-106. In the
concurrency comparison at 181-206, assert both results have TimedOut == false
before other assertions; in the result-collection sites at 226-236 and
matrix_test.go:95-106, reject timed-out results before appending or reporting
them. Ensure runtime is imported where needed.
In `@test/scenario/run.go`:
- Around line 236-238: Protect shared procRNG access when configuring concurrent
processing through WithConcurrency: synchronize JitteredProcessor and
SlowOutlierProcessor callbacks, or provide each call an isolated RNG, and use
deferred unlocking around the callback execution.
🪄 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: 3a0eca42-3438-4665-9725-ed0fc72505d6
📒 Files selected for processing (12)
README.mddocs/improvements/plan-perf.mddocs/improvements/thresholds.mdpkg/batcher/batcher.gopkg/batcher/concurrency_test.gopkg/batcher/constants.gopkg/batcher/options.gopkg/batcher/stats.gopkg/batcher/worker_test.gotest/scenario/concurrency_test.gotest/scenario/matrix_test.gotest/scenario/run.go
eb1ca92 to
2a8be75
Compare
2a8be75 to
bf19bfa
Compare
There was a problem hiding this comment.
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (2)
test/scenario/scenario_test.go (1)
141-167: 🎯 Functional Correctness | 🟡 Minor | ⚡ Quick winAssert
TimedOutis false before comparing p50.Both runs use a 50ms
FixedProcessorunder serial processing. If either run times out, itsEndToEnddistribution describes an incomplete run, and therequire.Greatercomparison at Line 164 is meaningless.TestConcurrencyRemovesSlowProcessorCouplingintest/scenario/concurrency_test.goalready guards its comparison this way.Based on learnings: for latency comparisons across scenario runs, assert
scenario.Result.TimedOut == falsefor both results before comparing latency.💚 Proposed fix
t.Logf("100ms window: p50=%s p99=%s mean_batch=%.0f downstream/s=%.0f", large.EndToEnd.P50, large.EndToEnd.P99, large.MeanBatchSize, large.DownstreamPerSec) + require.False(t, small.TimedOut, "5ms window run must complete") + require.False(t, large.TimedOut, "100ms window run must complete") + require.Greater(t, small.EndToEnd.P50, large.EndToEnd.P50,🤖 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 `@test/scenario/scenario_test.go` around lines 141 - 167, Before comparing latencies in the window comparison test, assert that both the small and large scenario results have TimedOut == false, using the existing scenario.Result values and a clear failure message. Add these guards before the require.Greater call so p50 is only compared for complete runs.Source: Learnings
docs/improvements/plan-perf.md (1)
500-504: 🗄️ Data Integrity & Integration | 🟠 Major | ⚡ Quick winAlign the
Stats()allocation gate across both documents.The plan lists
Stats()allocations as an exact threshold, but the threshold file says that noStats()allocation gate exists. This leaves CI enforcement ambiguous. Choose one policy and update both files.
docs/improvements/plan-perf.md#L500-L504: remove theStats()gate or document its enforcing test and blocking CI behavior.docs/improvements/thresholds.md#L45-L46: align the explanation with the selected policy.🤖 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 500 - 504, Align the Stats() allocation policy across docs/improvements/plan-perf.md lines 500-504 and docs/improvements/thresholds.md lines 45-46: either remove the Stats() gate from the plan and state that no gate exists, or document the enforcing test and blocking CI behavior in both files. Apply the same selected policy and wording consistently at both sites.
🤖 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.
Outside diff comments:
In `@docs/improvements/plan-perf.md`:
- Around line 500-504: Align the Stats() allocation policy across
docs/improvements/plan-perf.md lines 500-504 and docs/improvements/thresholds.md
lines 45-46: either remove the Stats() gate from the plan and state that no gate
exists, or document the enforcing test and blocking CI behavior in both files.
Apply the same selected policy and wording consistently at both sites.
In `@test/scenario/scenario_test.go`:
- Around line 141-167: Before comparing latencies in the window comparison test,
assert that both the small and large scenario results have TimedOut == false,
using the existing scenario.Result values and a clear failure message. Add these
guards before the require.Greater call so p50 is only compared for complete
runs.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: CHILL
Plan: Pro
Run ID: 2d1b869c-d9ae-41bc-bdba-bf92039e637e
📒 Files selected for processing (11)
README.mddocs/improvements/plan-perf.mddocs/improvements/thresholds.mdpkg/batcher/batcher.gopkg/batcher/batcher_test.gopkg/batcher/options.gopkg/batcher/worker_test.gotest/scenario/concurrency_test.gotest/scenario/matrix_test.gotest/scenario/run.gotest/scenario/scenario_test.go
🚧 Files skipped from review as they are similar to previous changes (5)
- test/scenario/matrix_test.go
- pkg/batcher/options.go
- test/scenario/concurrency_test.go
- pkg/batcher/batcher.go
- pkg/batcher/worker_test.go
Milestone 3.1. Makes the concurrency contract real before the worker pool is implemented in 3.2. DefaultConcurrency changes 3 -> 1. The previous 3 was advertised but read nowhere, so observed parallelism was already 1; changing the default to match reality preserves every existing caller's behaviour. Shipping 3 as a live value would silently parallelise processors holding unsynchronised state. Adds WithConcurrency(n) and WithoutOrderedProcessing(). The latter is an acknowledgement-only gate: WithConcurrency(n>1) without it panics at construction, because concurrent processing gives up cross-batch ordering and processor mutual exclusion. This is a programming error detectable at construction, so a panic is more honest and compatible than changing New to return an error or silently falling back to n=1. The n=1 contract is now tested, not inferred: - Default is serial. - One producer's items retain publication order across size, timer, and shutdown flushes. - The processor mutual-exclusion detector proves max concurrent invocation is 1. - Concurrent producers preserve each producer's own subsequence; the test deliberately does not assert cross-producer order, because Go does not order concurrent sends. - The panic gate and mutual-exclusion detector are sabotage-verified: removing either mechanism makes its test fail. Also corrects the plan's stale Phase 3 wording. Phase 2 already has one unbuffered aggregation->processor handoff for behaviour preservation, so saying n=1 is "inline" was false. The actual Phase 3 constraint is no ADDITIONAL worker-pool dispatch at n=1; the processor remains serial until an explicitly acknowledged n>1 worker pool lands in 3.2.
Milestone 3.2. WithConcurrency is now real. At n=1 the existing serial aggregation->processor handoff stays unchanged: the processor is mutually exclusive and batches are processed in publication order. At n>1, explicitly acknowledged by WithoutOrderedProcessing, n workers pull batches from the same UNBUFFERED dispatch channel. The unbuffered part is a correctness property: a buffered dispatch queue would hold batches not counted by MaxQueueSize, silently voiding the capacity guarantee. The bound remains MaxQueueSize + BatchSize + (n*BatchSize) + PublishersInGate. Adds Stats().InFlight, updated by workers on batch boundaries rather than on the enqueue path, so observability does not add another contended producer atomic. Shutdown closes dispatch and waits for every worker; it cannot report success while a batch is inside a processor call. Measured at 10k items/s with a 50ms processor and a 5ms window: n=1 p50=70ms mean batch=476 n=2 p50=25ms mean batch=244 n=4 p50=15ms mean batch=123 n=8 p50=4ms mean batch=62 That is the Phase 3 result: acknowledged capacity breaks the slow-processor coupling that made a small window ineffective. The control also passes: with a 100ms window and a 20ms processor (timer-bound), n=1 and n=4 both measure ~70ms, so concurrency does not change the batching rule where there is no processor bottleneck to remove. Tests pin parallel invocation, within-batch order, capacity bound, active-worker shutdown, panic recovery without worker loss, exact goroutine budgets (n=1=2, n=2=3, n=4=5, n=8=9), and worker completion/close/error/enqueue races under -race.
Documents the Phase 3 concurrency contract in the README: serial processing is the default, WithConcurrency(n>1) requires WithoutOrderedProcessing, batches lose cross-batch ordering, and processors must be goroutine-safe. Explains that item order remains intact within each batch and that unbuffered dispatch preserves the queue bound instead of creating a hidden backlog.
Six review findings on Phase 3.
The scenario harness had a real data race. JitteredProcessor and
SlowOutlierProcessor read a single shared *rand.Rand from inside the processor, and
BatcherOptions can enable worker concurrency, so several workers called it at once;
matrix_test.go pairs exactly those processors with WithConcurrency. Reproduced under
-race:
WARNING: DATA RACE
math/rand.(*rngSource).Uint64()
Guarded with a mutex around the RNG read only, not around the sleep, so concurrent
workers still overlap their simulated downstream latency. The same probe now passes
clean. Beyond the race report, undefined service times would have silently
invalidated every measurement taken with those processors.
Two worker tests parked the processor on a release channel and closed it at the end
of the test body. Any earlier assertion failure calls t.FailNow, which exits the test
goroutine and leaves the channel open, so a clean failure became a package-level
timeout with workers and publishers parked forever. Both now register a
sync.Once-guarded release with t.Cleanup, and the bounded-capacity test also
registers b.Close(), which it never called on any path.
Every 10k items/s scenario now sets Producers: runtime.NumCPU(). One goroutine
cannot sustain a 100µs inter-arrival gap against real sleep granularity, so the
offered rate collapses toward the service rate and the comparison stops measuring
what it claims. The report sweeps also reject timed-out runs instead of appending
them as evidence, and the not-processor-bound control now asserts neither run timed
out.
Closed every auto-started batcher in batcher_test.go: 12 were constructed and 5
closed, so seven leaked an aggregator and processor goroutine that could start after
TestWorkerGoroutineBudget captured its baseline and make the exact delta assertion
flaky. Verified none remain unclosed, and the budget tests pass repeatedly.
Also fixed a flake this validation exposed, introduced by my own earlier review fix.
TestHarnessReportsLatenessAndInvalidatesBadRuns derived a budget from one probe run
and asserted a *later* run stayed within it, which fails under parallel CPU
contention for scheduling reasons rather than because the guard is wrong. It now
asserts both directions against budgets that cannot be marginal — an hour is
unreachable by any scheduling delay, a nanosecond is unmeetable by any scheduler — so
it tests the guard instead of the host's timer precision. Sabotage-verified: forcing
LatenessValid to true still fails the test.
Two documentation corrections. plan-perf.md stated the n=1 budget as 1 goroutine in
the same paragraph that explained why that prediction was invalidated; it now says 2
(aggregator + serial processor), matching thresholds.md and the enforcing test. The
post-close threshold said "equal to pre-construction baseline" while both tests
assert <=; it now says "at most", with the reason stated — an unrelated parallel test
can retire a goroutine during the settle window, so equality would be flaky without
indicating a leak, and only a count above the baseline is a leak signal.
Three findings from the Phase 3 adversarial review, all documentation precision on correct code. The latency formula max(BatchInterval, processor duration) is right for the steady-state flush interval but understates per-item worst case, which is additive. While the aggregator is blocked handing a batch to the busy worker, no interval timer is running: it is armed only when the next batch takes its first item. An item arriving in that window waits the remaining processor time plus a full interval. Stated in both the WithConcurrency godoc and the README, since a reader sizing a latency budget from the maximum would be optimistic by up to one interval. The "70ms -> 4ms (17x)" figures are now labelled with the host that produced them (darwin/arm64, M4 Pro, GOMAXPROCS=12, Go 1.26.5), and the doc records that a re-run on the same machine under different load measured 120ms -> 54ms: same direction, about 2x. The tests only ever asserted relative ordering, which is why they still pass; the risk was a single multiplier being quoted as a property of the library rather than of one run. Rewrote the capacity bound as N + (1 + n) * BatchSize + gate in batcher.go, worker_test.go and the plan. Arithmetically identical to BatchSize + n * BatchSize, but it names the leading term as the aggregator's held batch instead of leaving it to be inferred -- which is how the Phase 2 doc lost that batch in the first place.
bf19bfa to
6548457
Compare
What
Phase 3 of the batcher performance plan: explicit concurrency. This is the
phase that makes a 5-10ms batch window actually worth configuring.
feat: make concurrent processing an explicit opt-infeat: process acknowledged batches with a worker pooldocs: document explicit concurrent processing opt-inWhy this phase exists
Config.Concurrencydefaulted to 3 but was read nowhere, so observedparallelism was always 1. With one worker, a batch must finish before the next can
start, which makes the effective flush interval
max(BatchInterval, processor duration). That is why lowering the window did nothing — or made things worse —in Phases 1 and 2.
Measured result
10k items/s, 50ms processor, 5ms window:
The 5ms window is finally honoured at n=8 — a 17x p50 improvement. Concurrency
adds capacity, so more batches leave per second and each is smaller.
The control matters as much as the result. With a 100ms window and a 20ms
processor — timer-bound, not processor-bound — n=1 and n=4 both measure ~70ms p50.
Concurrency does not change the batching rule; it removes a bottleneck where one
exists. A test asserts this, because an "improvement" there would have meant the
measurement was wrong.
Design decisions
DefaultConcurrency3 → 1. The old value was inert, so 1 is what every calleralready had. Shipping 3 as a live value would have silently parallelised
processors holding unsynchronised state.
WithoutOrderedProcessing()is a required acknowledgement.WithConcurrency(n>1)without it panics at construction. Concurrency gives upcross-batch ordering and processor mutual exclusion; that trade belongs in the code
that makes it, not discovered in production. A panic is more honest than changing
Newto return an error (breaks every call site) or silently falling back to n=1(gives you a throughput setting that doesn't apply).
Dispatch is unbuffered. This is correctness, not style: a buffered dispatch
queue would hold batches
MaxQueueSizenever counted, silently voiding thecapacity guarantee. The bound stays
MaxQueueSize + BatchSize + (n × BatchSize) + PublishersInGate, asserted with eachterm recorded separately.
Recovery stays scoped to the batch. A panic that killed its worker would
silently and permanently degrade throughput. A test panics on the first
nbatches— enough to kill every worker — and asserts later batches still complete.
Two plan corrections
Both from measuring rather than assuming:
separate aggregation→processor handoff to preserve latency semantics, so this was
stale. The real Phase 3 constraint is no additional worker-pool dispatch at
n=1, with serial invocation.
serial processor), and
1 + nabove that: n=2→3, n=4→5, n=8→9, all with zeroleaked. Corrected in the plan and thresholds table with the reasoning.
Validation
go test -race ./...andgo vet ./...clean.making the processing loop concurrent fails the n=1 mutual-exclusion detector.
100ms → 100ms), which is the correct outcome: the coupling is real when serial,
and concurrency is the opt-in cure rather than a silent default change.
n=1enqueue regression versus the Phase 2 baseline (geomean 3.8% fastervia
benchstat,-count=6).BatcherOptionsfield so Phase 3 evidence reusesthe same open-loop, lateness-guarded framework rather than a second generator.
Dependency context
Stacked on #29 (Phase 2), which targets #28 (Phase 1). Review bottom-up.
Phase 4 is next: completing the
Stats()snapshot and the benchmark-gated adaptivebatch capacity. Phase 5 then uses this concurrency data to decide whether
DefaultBatchIntervalshould change.Stack created with GitHub Stacks CLI • Give Feedback 💬
Summary by CodeRabbit