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
38 changes: 38 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -220,6 +220,44 @@ if err := b.Shutdown(ctx); err != nil {
`Shutdown` is resumable: a later call waits on the same drain rather than starting
a new one, so an expired deadline never costs you accepted work.

### Concurrent processing and ordering

By default, Batcher processes **one batch at a time**. This guarantees that a
processor is never invoked concurrently and that a single producer's batches are
processed in publication order. Use this default when the processor holds
unsynchronised state or when cross-batch ordering matters.

A slow processor can make a small batch window ineffective at this setting: a 5ms
window behind a 50ms processor is effectively bounded by the processor.

The per-item worst case is additive, not a maximum. The interval timer starts when a
batch takes its first item, so while the aggregator is blocked handing the previous
batch to the busy worker no timer is running. An item arriving in that window waits
the remaining processor time *plus* a full interval.

When the processor is goroutine-safe and cross-batch ordering does not matter,
explicitly opt into worker concurrency:

```go
b := batcher.New(
batcher.WithProcessor(processor),
batcher.WithConcurrency[Item](4),
batcher.WithoutOrderedProcessing[Item](), // required acknowledgement
)
```

`WithoutOrderedProcessing` is deliberately required. `WithConcurrency(n > 1)`
without it panics at construction, because concurrent processing gives up two
properties callers may rely on:

- batches may start, interleave, and complete in any order;
- `processor` may be invoked concurrently, so it must be goroutine-safe.

Items retain publication order **within** each batch at every concurrency. The
worker dispatch is unbuffered, so `WithMaxQueueSize` continues to bound queued
work; concurrency trades larger batches for more downstream calls, not an
unbounded hidden dispatch queue.

### Back-pressure and rejection

By default the queue is unbounded, which absorbs bursts well but turns a sustained
Expand Down
62 changes: 42 additions & 20 deletions docs/improvements/plan-perf.md
Original file line number Diff line number Diff line change
Expand Up @@ -512,7 +512,7 @@ Scenario matrix: windows `500µs, 1, 2, 5, 10, 20, 50, 100ms` × arrivals
| Sparse-window allocated bytes/flush, 4.2 trigger | > 2 KB/flush to justify work |
| 4.2 allocation regression ceiling, any scenario | ≤ +2% allocated bytes |
| Goroutines per idle batcher after 2.2 | exactly 0 |
| Goroutines per running n=1 batcher after 2.2 | exactly 1 (aggregator) |
| Goroutines per running n=1 batcher after 2.2 | exactly 2 (aggregator + processor) |
| Goroutines after `closed` | equal to pre-construction baseline |

Latency percentiles are deliberately absent: they are reported, never gated.
Expand Down Expand Up @@ -708,8 +708,10 @@ close-to-terminate model is replaced here, atomically, in one PR.
baseline. Removing the `chann` relay is expected to make `Add` *faster* than
today's 335 ns/op (prototype: 136-168 ns/op); a regression here indicates the
owned queue is mis-implemented.
- **Goroutine assertions:** 0 per idle batcher, 1 per running `n=1` batcher, and
exact return to the pre-construction baseline after `closed`. No queue owns a
- **Goroutine assertions:** 0 per unstarted batcher, 2 per running `n=1` batcher
(aggregator + serial processor), `1 + n` per running `n>1` batcher (aggregator
+ workers), and exact return to the pre-construction baseline after `closed`.
Measured: n=1=2, n=2=3, n=4=5, n=8=9, all with zero leaked. No queue owns a
goroutine.

**Dependencies**: Milestones 2.1 and 1.3.
Expand Down Expand Up @@ -839,9 +841,10 @@ Prototype evidence, 5ms window and 50ms processor at 10,000 items/s:
| decoupled, n=4 | 124 | **9.6ms** |

Naive decoupling at `n=1` is **2.5x worse**, because it adds a stage without
adding capacity. Inline aggregation already self-tunes: items pool while the
processor runs, so the next batch is larger. Therefore the `n=1` path is left
inline and unchanged; only `n>1` introduces a worker pool.
adding capacity. Phase 2 already has one unbuffered aggregation→processing handoff
for behaviour preservation; the relevant constraint is that `n=1` adds **no
additional worker-pool dispatch** and invokes the processor serially. Only `n>1`
introduces worker-pool dispatch and real concurrent processing.

### Milestone 3.1 — Concurrency configuration and ordering contract

Expand All @@ -850,11 +853,10 @@ inline and unchanged; only `n>1` introduces a worker pool.
- Add `WithConcurrency(n)`; set `DefaultConcurrency = 1`.
- Add `WithoutOrderedProcessing()` as an acknowledgement-only gate.
- Panic at construction when `n > 1` without the acknowledgement.
- Keep the `n = 1` aggregation path **without a worker-handoff stage**. It is
not byte-for-byte identical after Phase 2 — it necessarily uses the new gate,
accounting, drain, recovery, and stats machinery — but it preserves the
important performance shape: aggregation invokes the processor inline and no
additional dispatch queue exists.
- Keep the `n = 1` path without an **additional worker-pool dispatch**. Phase 2
already has one unbuffered aggregation→processing handoff; it is retained because
removing it changes observable latency. At `n=1` the processor remains serial
and no worker pool or extra dispatch queue is introduced.
- Document FIFO-by-publication-order and processor mutual-exclusion guarantees
for `n = 1`.

Expand All @@ -880,7 +882,7 @@ inline and unchanged; only `n>1` introduces a worker pool.
*within* each batch.
- **Bounded dispatch:** the aggregator-to-worker handoff is unbuffered, so
backpressure remains at admission and total accepted work stays bounded by
`N + BatchSize + (n × BatchSize)`. An unbounded dispatch queue would silently
`N + (1 + n) × BatchSize`. An unbounded dispatch queue would silently
void the `MaxQueueSize` contract.
- Add `Stats().InFlight`, tracked per worker, so the snapshot separates queued,
batch-held, and in-flight work, and so the drain waits for active workers.
Expand All @@ -893,11 +895,26 @@ inline and unchanged; only `n>1` introduces a worker pool.
- **Semantic-difference matrix:** scenarios compare `n=1` and `n>1` timer
cadence, batch-size distribution, admission blocking, pending depth, and
flush timing under the same slow processor. Documentation labels the observed
differences contractual: n=1 aggregates inline; n>1 aggregates independently
until the unbuffered worker dispatch is unavailable; cross-batch ordering is
not guaranteed at n>1.
differences contractual: n=1 uses the existing serial aggregation→processor
handoff; n>1 aggregates independently until the **unbuffered** worker dispatch
is unavailable; cross-batch ordering is not guaranteed at n>1.

The implementation was measured at 10k items/s with a 50ms processor and a 5ms
window on **darwin/arm64, Apple M4 Pro, GOMAXPROCS=12, Go 1.26.5**: n=1 p50 70ms
/ mean batch 476; n=2 p50 25ms / 244; n=4 p50 15ms / 123; n=8 p50 4ms / 62.

Treat the ratio between rows as the finding and the absolute values as this
host's. A re-run on the same machine under different load measured n=1 p50 120ms
and n=8 p50 54ms — the same direction, roughly 2x rather than 17x. Anything
quoting a single multiplier out of this table will be wrong somewhere else, which
is why the tests assert only the relative ordering (`n>1` p50 below `n=1` p50)
and never a millisecond threshold.

With a 100ms window and a 20ms processor (timer-bound rather than
processor-bound), n=1 and n=4 both measured ~70ms p50, which is the control
showing concurrency does not alter the batching rule.
- Total accepted-but-not-terminal items never exceed
`N + BatchSize + (n × BatchSize) + PublishersInGate` under a blocked processor;
`N + (1 + n) × BatchSize + PublishersInGate` under a blocked processor;
the test records `PublishersInGate` separately and demonstrates that its maximum
equals the number of intentionally parked publisher goroutines.
- **Worker-mode drain test:** shutdown while a full batch is in flight on a busy
Expand All @@ -910,13 +927,18 @@ inline and unchanged; only `n>1` introduces a worker pool.
- Goroutine count returns to the pre-construction baseline after `closed`, for
`n=1` and `n>1`, using this leak-check protocol: sample after `Shutdown`
returns, retry up to 50 times at 20ms intervals to allow scheduler settle,
then require exact equality with the pre-construction count. No unexplained
allowance is permitted; any runtime-owned goroutine must be named in the test.
then require the count to be at most the pre-construction baseline. It is `<=`
rather than `==` because an unrelated parallel test can retire a goroutine during
the settle window, which would make an equality assertion flaky without
indicating a leak. A count *above* the baseline is the leak signal, and any
runtime-owned goroutine must still be named in the test.
- **Per-batcher goroutine budget is asserted by a test table** covering `new`,
`running` at `n=1`, `running` at `n>1`, `draining`, and `closed`. Because 2.2
replaces both `chann` queues with goroutine-free owned queues, the achievable
budget is: `new` = 0; `running` at `n=1` = 1 (aggregator); `running` at `n>1`
= 1 + n; `closed` = 0. Any deviation must be explained in the test rather than
budget is: `new` = 0; `running` at `n=1` = 2 (aggregator + serial processor);
`running` at `n>1` = 1 + n; `closed` = 0. (An earlier prediction of 1 for `n=1`
was invalidated by Phase 2, which deliberately retains the separate serial
processor goroutine to preserve latency semantics.) Any deviation must be explained in the test rather than
absorbed into an allowance. The drain coordinator runs on the caller's
`Shutdown` goroutine and must not add a persistent goroutine; if an
implementation needs one, the table and this budget must be updated with it
Expand Down
31 changes: 13 additions & 18 deletions docs/improvements/thresholds.md
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ the blocking signal rather than timings.
| `Add` allocations, unbounded path | exactly 0 | `TestAddAllocatesNothingPerCall` |
| Recovery wrapper allocations, non-panic | exactly 0 | `TestRecoveredPanicAddsNoSteadyStateAllocations` |
| Scenario recorder allocations per item | ≤ 1 total, and must not grow with run length | `TestHarnessRecorderDoesNotAllocatePerItem` (`test/scenario`) |
| Goroutines per running batcher | exactly 2 | `TestGoroutineBudgetPerRunningBatcher` |
| Goroutines per running batcher, `n=1` | exactly 2 | `TestGoroutineBudgetPerRunningBatcher` |

`Stats()` is a fixed set of atomic loads returning a value type, so it has no
allocation gate of its own; the `Add` gate covers the hot path that matters.
Expand Down Expand Up @@ -68,23 +68,16 @@ Compare with `benchstat` over `-count=10`. A single run is not evidence. This is

## Goroutine gates (blocking)

These are the gates in force today. They must agree with the allocation table
above, which enforces the same running-batcher count.
These are the gates in force today. The n=1 row agrees with the allocation table
above; `n>1` is now a real, acknowledged configuration and has its own worker-pool
gate. Every count is enforced by a test, not merely documented.

| Gate | Threshold |
| ---------------------------------- | ----------------------------------------- |
| Goroutines per unstarted batcher | exactly 0 |
| Goroutines per running batcher | exactly 2 (aggregator + serial processor) |
| Goroutines after terminal `closed` | equal to pre-construction baseline |

Enforced by `TestGoroutineBudgetPerRunningBatcher`.

Phase 3 introduces explicit worker concurrency and will replace the single
running-batcher row with `1 + n` (aggregator plus workers). That row is
deliberately absent here rather than stated as a gate, because a threshold for a
configuration this code cannot express is not enforceable — and two rows claiming
different counts for the same batcher is worse than one row that is merely
incomplete.
| Gate | Threshold | Enforced by |
| ---------------------------------- | ----------------------------------------- | ----------- |
| Goroutines per unstarted batcher | exactly 0 | `TestWorkerGoroutineBudget` |
| Goroutines per running `n=1` | exactly 2 (aggregator + serial processor) | `TestGoroutineBudgetPerRunningBatcher`, `TestWorkerGoroutineBudget` |
| Goroutines per running `n>1` | exactly 1 + n (aggregator + workers) | `TestWorkerGoroutineBudget` |
| Goroutines after terminal `closed` | at most the pre-construction baseline | `TestGoroutineBudgetPerRunningBatcher`, `TestWorkerGoroutineBudget` |

Current `main` owns 6 goroutines per batcher. Phase 2.1 removed `rill`
(**measured 6 → 5**) and Phase 2.2 removed both `chann` relays and the input
Expand All @@ -105,7 +98,9 @@ draining into aggregation regressed sequential `Add` by 39-50% because the

Both earlier estimates assumed aggregation and processing could share one
goroutine. They cannot without changing observable latency, which is why the
enforced count is 2 rather than 1.
enforced count is 2 at `n=1`. Phase 3 owns the worker model on top of that base:
`n=1` is 2 goroutines and `n>1` is `1 + n`, measured as n=2→3, n=4→5, n=8→9 with
zero leaked.

## Conditional gates (Milestone 4.2 only)

Expand Down
84 changes: 70 additions & 14 deletions pkg/batcher/batcher.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package batcher
import (
"context"
"errors"
"fmt"
"runtime/debug"
"sync"
"time"
Expand Down Expand Up @@ -31,6 +32,10 @@ type Config[T any] struct {
CloseGrace time.Duration
ErrorBufferSize int
ProcessorFunc Processor[T]

// UnorderedProcessingAcknowledged records that the caller accepted the
// ordering trade required for Concurrency > 1. See WithoutOrderedProcessing.
UnorderedProcessingAcknowledged bool
}

type Batcher[T any] struct {
Expand Down Expand Up @@ -86,6 +91,8 @@ func New[T any](options ...Option[T]) *Batcher[T] {
option(b)
}

b.validateConfig()

b.input = newQueue[T](b.config.MaxQueueSize)
b.errorsChan = make(chan error, b.config.ErrorBufferSize)

Expand All @@ -96,6 +103,25 @@ func New[T any](options ...Option[T]) *Batcher[T] {
return b
}

// validateConfig rejects configurations whose guarantees contradict each other.
//
// This panics rather than returning an error because New has never been able to
// fail, and adding an error return would break every existing call site for a
// mistake that is always a programming error: the combination is fixed at
// construction, so it is caught by the first test run rather than in production.
func (b *Batcher[T]) validateConfig() {
if b.config.Concurrency > 1 && !b.config.UnorderedProcessingAcknowledged {
panic(fmt.Sprintf(
"batcher: WithConcurrency(%d) requires WithoutOrderedProcessing(). "+
"Concurrent processing gives up cross-batch ordering and lets the "+
"processor be invoked concurrently, so the processor must be "+
"goroutine-safe. Add WithoutOrderedProcessing() to acknowledge this, "+
"or keep WithConcurrency(1).",
b.config.Concurrency,
))
}
}

// Start begins processing. It is idempotent: repeated or concurrent calls start
// exactly one processing lifecycle.
//
Expand Down Expand Up @@ -197,6 +223,7 @@ func (b *Batcher[T]) Stats() Stats {
IntakePending: b.counters.intakePending.Load(),
PublishersInGate: b.gate.inGate(),
Queued: int64(b.input.length()),
InFlight: b.counters.inFlight.Load(),
Accepted: b.counters.accepted.Load(),
Completed: b.counters.completed.Load(),
Failed: b.counters.failed.Load(),
Expand Down Expand Up @@ -245,30 +272,55 @@ func (b *Batcher[T]) Errors() <-chan error {
// - Empty batches are never emitted.
// - Shutdown flushes the pending partial batch rather than discarding it.
//
// Aggregation runs here and processing runs in its own goroutine, connected by an
// unbuffered channel. That separation is load-bearing: it lets the next batch
// accumulate while the current one is being processed. Merging the two measurably
// changes latency — with a 50ms processor at 10k items/s a merged loop made a 5ms
// window faster than a 100ms one, inverting the documented baseline — and that is
// a Phase 3 decision, not this milestone's.
// Aggregation runs here and processing runs in Concurrency separate goroutines,
// connected by an UNBUFFERED channel. Two properties follow, and both are
// deliberate:
//
// - The separation lets the next batch accumulate while the current one is being
// processed. Merging them measurably changes latency: with a 50ms processor at
// 10k items/s a merged loop made a 5ms window faster than a 100ms one,
// inverting the documented baseline.
// - The channel is unbuffered, so back-pressure stays at admission. A buffered
// dispatch queue would silently void the MaxQueueSize contract by holding
// batches nobody counted, so accepted-but-unfinished work stays bounded by
// MaxQueueSize + (1 + Concurrency) × BatchSize + publishers in the gate. The
// leading batch is the one the aggregator holds after it leaves the queue; the
// Concurrency term is the batches inside processors.
//
// At Concurrency 1 there is exactly one worker, so the processor is never invoked
// concurrently and batches are processed in publication order. Above 1 the caller
// has acknowledged, via WithoutOrderedProcessing, that cross-batch ordering and
// processor mutual exclusion are given up.
func (b *Batcher[T]) run() {
defer close(b.stopped)

batches := make(chan []T)

processing := make(chan struct{})
workers := b.config.Concurrency
if workers < 1 {
workers = 1
}

go func() {
defer close(processing)
var processing sync.WaitGroup

for items := range batches {
b.process(items)
}
}()
processing.Add(workers)

for range workers {
go func() {
defer processing.Done()

for items := range batches {
b.process(items)
}
}()
}

// Closing the dispatch channel stops the workers, and waiting for them is what
// makes the drain complete: shutdown must not report success while a worker is
// still inside the processor.
defer func() {
close(batches)
<-processing
processing.Wait()
}()

var (
Expand Down Expand Up @@ -378,12 +430,16 @@ func (b *Batcher[T]) run() {
// category is counted and the drain obligation is released exactly once. Leaving it
// unreleased would inflate Pending permanently and hang shutdown forever.
func (b *Batcher[T]) process(items []T) {
b.counters.inFlight.Add(int64(len(items)))

outcome := outcomeCompleted

defer func() {
b.counters.terminal(outcome, len(items))
}()

defer b.counters.inFlight.Add(int64(-len(items)))

defer func() {
recovered := recover()
if recovered == nil {
Expand Down
Loading
Loading