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
35 changes: 31 additions & 4 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -291,12 +291,39 @@ and after shutdown it is a counted no-op rather than a panic.
```go
s := b.Stats()

// s.Queued -> queue depth; the value to alert on
// s.Pending -> accepted work not yet finished, including in-flight batches
// s.Rejected -> refused enqueues
// s.DroppedErrors -> diagnostics lost because Errors() was not drained
// Where work currently is. These three are disjoint, so together they show
// whether a backlog is waiting on the queue, on batching, or on the processor:
// s.Queued -> published, not yet taken by the aggregator (alert on this)
// s.BatchHeld -> held by the aggregator: filling, or waiting for a free worker
// s.InFlight -> inside a processor call

// Totals.
// s.Pending -> conservative drain obligation: accepted-or-reserved work,
// including in-flight batches. It counts publishers that have
// reserved but not yet published, so it equals
// accepted-but-unfinished work only once PublishersInGate == 0
// s.Accepted -> successful enqueues
Comment thread
heynemann marked this conversation as resolved.
// s.Rejected -> refused enqueues
// s.Completed / s.Failed / s.Panicked -> mutually exclusive terminal outcomes
// s.BatchesFlushed -> batches emitted. After a terminal drain,
// (Completed+Failed+Panicked)/BatchesFlushed is mean batch size
// s.DroppedErrors -> diagnostics lost because Errors() was not drained
```

`BatchesFlushed` is the coalescing signal when tuning `WithBatchInterval`: a mean
batch size well below `BatchSize` means windows are closing on the timer rather than
filling, so the interval is costing latency without buying batching. Include every
terminal outcome in that mean — a failed or panicked batch was still flushed, so
dividing by `Completed` alone undercounts whenever the processor errors.

A rising `BatchHeld` with `InFlight` at its ceiling means batches are ready but every
worker is busy — that is the signal to raise `WithConcurrency`, not to shrink the
window.

The snapshot is eventually consistent, not transactional: each field is read
independently, so use it to observe where work sits, not as an accounting identity
while load is in flight.

### Handling Errors

Whenever the processor function returns an error, the batcher will send the error in the `Errors()` channel:
Expand Down
6 changes: 4 additions & 2 deletions docs/improvements/plan-perf.md
Original file line number Diff line number Diff line change
Expand Up @@ -1080,8 +1080,10 @@ rather than intuition.
sealing instead of panicking; `Close()` no longer abandoning the drain at its
deadline; **`Len()` now counting accepted-but-unfinished work including
in-flight batches, so it can be non-zero with an empty queue**; `Join`'s
clarified quiescence-snapshot contract; the removal of the `rill` and `chann`
dependencies; **the module's minimum Go version moving from 1.22.4 to 1.25.0**
clarified quiescence-snapshot contract; **`Config()` returning a value snapshot
instead of a live mutable pointer**, which removes the ability to reconfigure a
running batcher and the data race that came with it; the removal of the `rill` and
`chann` dependencies; **the module's minimum Go version moving from 1.22.4 to 1.25.0**
(required by `golang.org/x/sys` v0.46.0, pulled in when dependencies were
modernised in 2.1); and any `ProvideBatcherInFX` signature change.
- **`Config()` must stop exposing live mutable state.** It currently returns the
Expand Down
49 changes: 43 additions & 6 deletions docs/improvements/thresholds.md
Original file line number Diff line number Diff line change
Expand Up @@ -41,9 +41,12 @@ the blocking signal rather than timings.
| 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, `n=1` | exactly 2 | `TestGoroutineBudgetPerRunningBatcher` |
| `Stats()` allocations | exactly 0 | `TestStatsIsAllocationFree` |

`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.
`Stats()` returns a value type built from atomic loads plus one queue length check
under the queue's existing mutex. It now carries its own allocation gate, because
metrics scraping runs continuously in production and generating garbage per scrape
would be a real cost.

The recorder threshold is not "exactly 0" because `AllocsPerItem` measures the
whole pipeline, including Batcher's own per-batch allocations, not just the
Expand Down Expand Up @@ -104,10 +107,44 @@ zero leaked.

## Conditional gates (Milestone 4.2 only)

| Gate | Threshold |
| ------------------------------------------ | ---------------------------- |
| Sparse-window allocated bytes/flush | > 2 KB/flush to justify work |
| Allocation regression ceiling, any scenario | ≤ +2% allocated bytes |
| Gate | Threshold | Measured |
| ------------------------------------------- | ---------------------------- | --------------- |
| Sparse-window allocated bytes/flush | > 2 KB/flush to justify work | 55,944 B/flush |
| Allocation regression ceiling, any scenario | ≤ +2% allocated bytes | +1.12% (worst) |

Milestone 4.2 was **triggered and implemented**. Sparse-window waste measured far
above the justification threshold: at a 1ms window with `BatchSize=1000` the
aggregator reserved 55,944 B/flush to hold a single item, and 559,944 B/flush at
`BatchSize=10000`.

Results per workload, adaptive versus the previous full-capacity strategy:

| Workload | Change |
| ----------------------- | ------- |
| steady sparse | -97.9% |
| small batches | -93.3% |
| full batches (control) | -0.00% |
| alternating sparse/full | +1.12% |
| burst after idle | -72.2% |
| bimodal | -3.45% |

The alternating case is the one that killed the rejected EWMA estimator, which
allocated *more* than doing nothing there. Recent-max keeps it inside the +2%
budget.

These figures are **measured evidence, not an enforced gate.** `BenchmarkCapacity*`
reports allocated bytes for each workload but does not compare them against the +2%
limit, and `capacity_test.go` asserts estimator behaviour (bounds, adaptation,
decay) rather than allocation deltas. Re-check the numbers with:

```sh
go test -run='^$' -bench='BenchmarkCapacity' -benchmem -count=1 ./pkg/batcher
```

Turning this into a gate needs a stored allocation baseline for the reference
runner, which does not exist yet for the same reason the throughput gate is still
advisory. Until then, a regression here is caught by reading the benchmark output,
not by CI.

The second gate must hold for alternating sparse/full, burst-after-idle, and
bimodal workloads, not just the sparse case the optimisation targets. An
Expand Down
149 changes: 127 additions & 22 deletions pkg/batcher/batcher.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import (
"fmt"
"runtime/debug"
"sync"
"sync/atomic"
"time"
)

Expand Down Expand Up @@ -41,6 +42,17 @@ type Config[T any] struct {
type Batcher[T any] struct {
config *Config[T]

// configFrozen becomes true before any processing goroutine starts. Options
// are construction-time configuration: applying one to an existing batcher is
// a no-op, so a caller cannot race the aggregator or processor by mutating
// BatchSize, BatchInterval or ProcessorFunc after New returns.
configFrozen atomic.Bool

// runtime is the immutable snapshot the processing goroutines use. It is taken
// in New, before any goroutine exists, so no worker ever reads the Config
// struct a caller may still hold a reference to.
runtime runtimeConfig[T]

gate *admissionGate
counters counters

Expand Down Expand Up @@ -69,6 +81,19 @@ type Batcher[T any] struct {
shutdownDone chan struct{}
}

// runtimeConfig is the frozen view of Config that processing goroutines read.
//
// Config is a plain struct and callers can hold a reference to one, so reading it
// from a worker is a data race by construction. Copying the fields the pipeline
// needs, once, before any goroutine starts, removes that whole class rather than
// relying on every read site to remember to use a local.
type runtimeConfig[T any] struct {
batchSize int
batchInterval time.Duration
workers int
processor Processor[T]
}

// New creates a new Batcher with the given options.
func New[T any](options ...Option[T]) *Batcher[T] {
b := &Batcher[T]{
Expand All @@ -93,6 +118,30 @@ func New[T any](options ...Option[T]) *Batcher[T] {

b.validateConfig()

// No option may change configuration after this point, including a batcher
// constructed with WithSkipAutoStart. SkipAutoStart controls lifecycle, not
// whether the configuration is still mutable.
b.configFrozen.Store(true)

b.runtime = runtimeConfig[T]{
batchSize: b.config.BatchSize,
batchInterval: b.config.BatchInterval,
workers: b.config.Concurrency,
processor: b.config.ProcessorFunc,
}

if b.runtime.batchSize < 1 {
b.runtime.batchSize = DefaultBatchSize
}

if b.runtime.workers < 1 {
b.runtime.workers = 1
}

if b.runtime.processor == nil {
b.runtime.processor = NoOpProcessor[T]
}

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

Expand Down Expand Up @@ -138,9 +187,17 @@ func (b *Batcher[T]) Start() {
})
}

// Config returns the batcher's configuration.
func (b *Batcher[T]) Config() *Config[T] {
return b.config
// Config returns a copy of the batcher's configuration.
//
// It is a snapshot, not a handle. Returning the live pointer let a caller mutate a
// running batcher's batch size, interval or processor from another goroutine, which
// is a data race against the aggregation loop and could change batching semantics
// mid-batch. Options are meant to be applied at construction; mutating them
// afterwards was never coherent, so this closes the hole rather than documenting it.
//
// Callers that need to change configuration should construct a new Batcher.
func (b *Batcher[T]) Config() Config[T] {
return *b.config
}

// Add enqueues an item. It is the compatibility fast path and returns no error.
Expand Down Expand Up @@ -223,11 +280,13 @@ func (b *Batcher[T]) Stats() Stats {
IntakePending: b.counters.intakePending.Load(),
PublishersInGate: b.gate.inGate(),
Queued: int64(b.input.length()),
BatchHeld: b.counters.batchHeld.Load(),
InFlight: b.counters.inFlight.Load(),
Accepted: b.counters.accepted.Load(),
Completed: b.counters.completed.Load(),
Failed: b.counters.failed.Load(),
Panicked: b.counters.panicked.Load(),
BatchesFlushed: b.counters.batchesFlushed.Load(),
Rejected: b.counters.rejected.Load(),
DroppedErrors: b.counters.droppedErrors.Load(),
}
Expand Down Expand Up @@ -294,12 +353,18 @@ func (b *Batcher[T]) Errors() <-chan error {
func (b *Batcher[T]) run() {
defer close(b.stopped)

batches := make(chan []T)
// Use the immutable runtime snapshot taken by New before any goroutine starts.
// In particular, ProcessorFunc is passed to workers rather than dereferenced
// from Config per batch. This keeps every processing read off the mutable Config
// struct, so a caller with a stale reference cannot race the pipeline.
var (
batchSize = b.runtime.batchSize
batchInterval = b.runtime.batchInterval
workers = b.runtime.workers
processor = b.runtime.processor
)

workers := b.config.Concurrency
if workers < 1 {
workers = 1
}
batches := make(chan []T)

var processing sync.WaitGroup

Expand All @@ -310,7 +375,7 @@ func (b *Batcher[T]) run() {
defer processing.Done()

for items := range batches {
b.process(items)
b.process(processor, items)
}
}()
}
Expand All @@ -324,9 +389,10 @@ func (b *Batcher[T]) run() {
}()

var (
batch []T
timer *time.Timer
timerC <-chan time.Time
batch []T
timer *time.Timer
timerC <-chan time.Time
capacities = newCapacityEstimator(batchSize)
)

stopTimer := func() {
Expand Down Expand Up @@ -356,40 +422,79 @@ func (b *Batcher[T]) run() {

stopTimer()

// The send blocks until a worker is free. Until it returns, the batch is
// still aggregator-held, which is the state BatchHeld exists to expose:
// releasing the counter before the send would hide a batch waiting on
// saturated workers.
batches <- items

b.counters.dispatched(len(items))
capacities.observe(len(items))
}

take := func(item T) {
if len(batch) == 0 {
batch = make([]T, 0, b.config.BatchSize)
// The capacity estimator is owned by this aggregation goroutine. It
// reduces timer-flush allocation pressure without pooling or reusing a
// slice after the processor receives it, so callers may still retain the
// batch slice exactly as before.
batch = make([]T, 0, capacities.capacity())

if timer == nil {
timer = time.NewTimer(b.config.BatchInterval)
timer = time.NewTimer(batchInterval)
} else {
timer.Reset(b.config.BatchInterval)
timer.Reset(batchInterval)
}

timerC = timer.C
}

batch = append(batch, item)

if len(batch) >= b.config.BatchSize {
if len(batch) >= batchSize {
flush()
}
}

// drainReady empties whatever is currently queued. The aggregator drains
// greedily rather than one item per wakeup, so a burst costs one signal.
//
// It transfers a run of items under a single queue lock rather than locking per
// item. The aggregator is the only consumer and keeps no invariant between items,
// so a block transfer is indistinguishable from N pops to it, while per-item
// locking was what serialised producers: profiling attributed 86% of mutex delay
// and 61% of CPU to the push path contending with this loop.
//
// Each transfer takes only what the batch being built can still accept. That
// bounds how long the queue lock is held, and it keeps the capacity contract
// exact: transferred items move straight into batch, so no additional
// batch-sized buffer of accepted work exists outside the queue. Draining a full
// batchSize regardless of the batch's remaining space added exactly that third
// buffer and pushed accepted work past N + 2*BatchSize + gate.
//
// drained is reused across wakeups to keep steady-state draining allocation-free.
// It is deliberately separate from batch: batch is handed to a worker and may be
// retained by the processor, whereas drained never leaves this goroutine.
drained := make([]T, 0, batchSize)

drainReady := func() {
for {
item, ok := b.input.pop()
if !ok {
room := batchSize - len(batch)
if room <= 0 {
// take flushes at batchSize, so this only happens if a flush could not
// complete. Fall back to one item and let take drive the flush.
room = 1
}

drained = b.input.popBatch(drained[:0], room)
if len(drained) == 0 {
return
}

b.counters.received(1)
take(item)
for _, item := range drained {
b.counters.received(1)
take(item)
}
}
}

Expand Down Expand Up @@ -429,7 +534,7 @@ func (b *Batcher[T]) run() {
// The 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.
func (b *Batcher[T]) process(items []T) {
func (b *Batcher[T]) process(processor Processor[T], items []T) {
b.counters.inFlight.Add(int64(len(items)))

outcome := outcomeCompleted
Expand All @@ -454,7 +559,7 @@ func (b *Batcher[T]) process(items []T) {
})
}()

if err := b.config.ProcessorFunc(items); err != nil {
if err := processor(items); err != nil {
outcome = outcomeFailed

b.publishError(err)
Expand Down
Loading
Loading