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
4 changes: 3 additions & 1 deletion .github/workflows/bench.yml
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,9 @@ jobs:
- uses: actions/checkout@v4
- uses: actions/setup-go@v6
with:
go-version: "stable"
# Match the module floor exactly; "stable" makes benchmark history
# change when a new Go release appears without a repository change.
go-version-file: go.mod
- name: Run benchmark
run: make ensure-build && make bench-run | tee build/output.txt
- name: Store benchmark result
Expand Down
90 changes: 84 additions & 6 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -39,8 +39,12 @@ func main() {
return nil
}),
)
// stop the batcher
defer batcher.Close()
// stop the batcher, and report if the drain did not finish in time
defer func() {
if err := batcher.Close(); err != nil {
fmt.Printf("shutdown incomplete: %v\n", err)
}
}()

// add operations to the batcher
for i := 0; i < 1000; i++ {
Expand Down Expand Up @@ -171,15 +175,89 @@ if err := batcher.Join(timeout); err != nil {

### Stopping the batcher

To stop the batcher you can use the `StopProcessing` function:
To stop the batcher, use `Close`:

```go
defer batcher.Close()
defer func() {
if err := batcher.Close(); err != nil {
// The drain did not finish within the grace period. Work is still being
// processed in the background, so decide deliberately: wait longer with
// Shutdown, or accept that the process is about to exit with work pending.
log.Printf("batcher shutdown incomplete: %v", err)
}
}()

// batcher.IsClosed() == true once the drain has completed
```

`Close` seals admission and drains work that has already been accepted, waiting up
to 30 seconds by default (configurable with `WithCloseGrace`). It is safe to call
multiple times and from multiple goroutines.

If the grace period expires, `Close` reports that the drain is incomplete — it does
**not** discard the remaining work, which keeps being processed in the background.
Comment thread
heynemann marked this conversation as resolved.
Do not write a bare `defer batcher.Close()`: it discards that report, and if the
process exits immediately afterwards, accepted work is lost without any signal.

When you need to control the wait, or to keep waiting after a timeout, use
`Shutdown`:

```go
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()

if err := b.Shutdown(ctx); err != nil {
var incomplete *batcher.ShutdownIncompleteError
if errors.As(err, &incomplete) {
// Still draining. Nothing was lost, and we can wait longer.
log.Printf("%d items still pending", incomplete.Pending)

err = b.Shutdown(context.Background())
}
}
```
Comment thread
coderabbitai[bot] marked this conversation as resolved.

`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.

// batcher.IsClosed() == true after this point
### Back-pressure and rejection

By default the queue is unbounded, which absorbs bursts well but turns a sustained
overload into unbounded memory growth. Note that shrinking the batch interval does
**not** bound queued work — only bounding the queue does.

To get back-pressure instead, bound the queue and use `Enqueue`, which reports why
an item was refused:

```go
b := batcher.New(
batcher.WithProcessor(processor.Process),
batcher.WithMaxQueueSize[Item](10_000),
)
Comment thread
coderabbitai[bot] marked this conversation as resolved.

if err := b.Enqueue(ctx, item); err != nil {
// batcher.ErrClosing -> shutting down
// context.DeadlineExceeded -> queue stayed full until the deadline passed
// context.Canceled -> queue was full and the caller cancelled
return err
Comment thread
heynemann marked this conversation as resolved.
}
```

This function is safe to be called multiple times as it will only stop the processor once.
`Add` remains available as the fast path for best-effort use: it returns no error,
and after shutdown it is a counted no-op rather than a panic.

### Observing a batcher

`Stats` returns an allocation-free snapshot suitable for frequent scraping:

```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
```

### Handling Errors

Expand Down
63 changes: 50 additions & 13 deletions docs/improvements/plan-perf.md
Original file line number Diff line number Diff line change
Expand Up @@ -373,11 +373,15 @@ conservation check, never an instantaneous API guarantee.
is bounded by:

```text
N + BatchSize + (concurrency × BatchSize) + P
N + (1 + concurrency) × BatchSize + P
```

where `P` is the number of publishers currently inside the gate, exposed as
`Stats().PublishersInGate`. At `n = 1` there is no dispatch stage, so the
`Stats().PublishersInGate`. The `(1 + concurrency)` term is deliberate: one batch
is held by the aggregator after leaving the queue, and `concurrency` batches can
be inside processors. Writing it as `BatchSize + concurrency × BatchSize`
understates the peak by one batch, which at `n = 1` is a factor-of-two error on
that term. At `n = 1` there is no dispatch stage, so the
`concurrency × BatchSize` term collapses to the single inline batch already
counted by `BatchSize`; the formula is therefore conservative rather than wrong
at `n = 1`. `P` is bounded by the caller's own concurrency, not
Expand Down Expand Up @@ -567,9 +571,32 @@ reviewable behaviour-preserving change.
- Characterization tests written pre-swap pass unchanged post-swap: batch
boundaries, timer-armed-on-first-item behaviour, ordering, error propagation,
no-empty-batch guarantee.
- Goroutines per *constructed, auto-started* batcher drop from 6 to 3, verified by
test. (2.2 removes the remaining two by replacing `chann`, reaching 1 per
running batcher and 0 for an unstarted one.)
- Goroutines per *constructed, auto-started* batcher drop from 6 to **5**,
verified by test and by direct measurement (50 batchers: +250 goroutines, 0
leaked after `Close`).

The plan first predicted 3, then 4. Both were wrong, and the reasons are worth
recording because they constrain the design:

1. **Merging aggregation into the processing loop changes latency.** Measured
with a 50ms processor at 10k items/s, a merged loop inverted the baseline
finding (5ms window p50 28ms versus 100ms window p50 75ms), because the next
batch only began accumulating after the processor returned. Aggregation and
processing must stay separate goroutines here; decoupling them differently
is a Phase 3 decision, not one this milestone may smuggle in.
2. **Input draining must be isolated from batch bookkeeping.** With a single
goroutine doing both, sequential `Add` regressed 39-50% against the stored
baseline, breaching the ≤+10% gate: the `chann` relay has a small bounded
ingress, and its output was not being read promptly while the aggregator was
busy with timer and slice work. Adding a dedicated forwarding stage restored
parity (no regression; one case 6.7% faster). This is what `rill`'s
`ToChans`/`Batch`/`FromChans` pipeline was buying.

The five are: `chann` input relay, input forwarder, aggregator, processor loop,
and `chann` error relay. Phase 2.2 removes both `chann` relays and the
forwarder along with them, because an owned queue needs no relay and can be
read directly. The single-goroutine target in the threshold table applies only
once Phase 3 owns the worker model.
- Enqueue benchmarks from 1.1 show no regression beyond predeclared thresholds.
- `go.mod`/`go.sum` no longer reference `rill`; no other dependency added.
- No public API change in this commit.
Expand Down Expand Up @@ -1032,7 +1059,9 @@ rather than intuition.
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; and any `ProvideBatcherInFX` signature change.
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.
Comment thread
heynemann marked this conversation as resolved.
- **`Config()` must stop exposing live mutable state.** It currently returns the
internal pointer (`batcher.go:76-78`), letting callers mutate `ProcessorFunc`
or `BatchSize` during processing — an exported data race that can invalidate
Expand Down Expand Up @@ -1135,14 +1164,22 @@ commit leaves the library in a consistent state, so any prefix of the plan is a
valid stopping point:

- Stopping after **1.x**: measurement only, no behaviour change.
- Stopping after **2.1**: same semantics, `rill` removed, goroutines 6 → 3.
- Stopping after **2.1**: same semantics, `rill` removed, goroutines 6 → 5,
and ~47% less memory per enqueued item.
- Stopping after **2.2**: admission is panic-free, accounting is sound, `chann`
is gone (1 goroutine per running batcher, 0 idle), and the internal
seal/gate/`noInput` coordinator is complete. `Close()` keeps its existing
signature, no longer abandons the drain, and a minimal `Stats()` is public; only
the resumable API and typed error are absent.
- Stopping after **2.3**: no accepted work is ever silently discarded.
- Stopping after **2.4**: diagnostics bounded, panics survivable.
is gone (**measured 2 goroutines per running batcher**, down from 6, and enqueue
54% faster), and the internal seal/gate/`noInput` coordinator is complete.
`Close()` keeps its existing signature, no longer abandons the drain, and a
minimal `Stats()` is public; only the resumable API and typed error are absent.
- Stopping after **2.3**: no accepted work is ever silently discarded. Verified
directly: a partial batch with a 30s interval and a shorter grace now reports 50
accepted / 50 processed / 0 pending, where it previously reported 50 accepted /
0 processed / 50 phantom pending.
- Stopping after **2.4**: diagnostics bounded, panics survivable. Phase 2 as a
whole is verified against the original defects: `Add` after `Close` is a counted
rejection instead of a process panic, bounded mode caps the queue and reports
rejections instead of growing the heap without limit, and goroutines per batcher
are 2 with none leaked.
- Stopping after **3.x**: small windows are honoured under slow processors.
- Stopping after **4.x**: operators can see overload developing.

Expand Down
87 changes: 63 additions & 24 deletions docs/improvements/thresholds.md
Original file line number Diff line number Diff line change
Expand Up @@ -12,9 +12,17 @@ reported by the scenario matrix and compared as a trend; it never fails a PR.

Thresholds are keyed to the environment. Re-baseline when any of these change.

The Go version is whatever `go.mod` declares: CI resolves it with
`go-version-file: go.mod` rather than a hardcoded string, so this row and the
toolchain cannot drift apart. Phase 2 raised the floor from 1.22.4 to 1.25.0 when
`golang.org/x/sys` was updated.

Stored baselines record the toolchain they were captured on in their own header,
which is what makes them comparable or not; see Baselines below.

| Field | Value |
| ---------- | ---------------------------------------- |
| Go version | 1.22.4 |
| Go version | 1.25.0 (from `go.mod`) |
| Runner | `ubuntu-latest` (GitHub-hosted, mutable) |
| Arch | `amd64` |

Expand All @@ -27,25 +35,28 @@ orientation only and are not gates.
Allocation counts are stable even on noisy shared runners, which is why they are
the blocking signal rather than timings.

| Gate | Threshold | Enforced by |
| --------------------------------------- | --------- | ---------------------------------------- |
| `Add` allocations, unbounded path | exactly 0 | `TestAddAllocatesNothingPerCall` |
| `Stats()` allocations | exactly 0 | added with `Stats()` in Phase 2 |
| Recovery wrapper allocations, non-panic | exactly 0 | added with panic recovery in Phase 2 |
| Gate | Threshold | Enforced by |
| --------------------------------------- | -------------------------------------------- | ------------------------------------------------------------- |
| `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` |

`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.
Comment thread
coderabbitai[bot] marked this conversation as resolved.

The recorder threshold is not "exactly 0" because `AllocsPerItem` measures the
whole pipeline, including Batcher's own per-batch allocations, not just the
recorder. Measured values are 0.03-0.04 allocations per item and do not grow with
run length, which is the property that matters: a recorder that allocated per item
would make every allocation figure it reports a measurement of itself.

`Add` currently allocates 0 per call in the timed region because the item is
constructed by the caller. Phase 2 replaces `chann` with a slice-backed unbounded
queue, at which point the gate becomes "zero allocations per `Add` in steady
state": `append` must allocate when it grows its backing array, so the gate is
measured after a stated warmup with queue capacity retained across drains, and
growth-path allocations are the one named exemption.
`Add` allocates 0 per call in the timed region because the item is constructed by
the caller. Phase 2 replaced `chann` with a slice-backed unbounded queue, so the
gate is now "zero allocations per `Add` in steady state": `append` must allocate
when it grows its backing array, so the gate is measured after a stated warmup
with queue capacity retained across drains, and growth-path allocations are the
one named exemption.

## Throughput gate (advisory until CI baseline exists)

Expand All @@ -55,18 +66,46 @@ growth-path allocations are the one named exemption.

Compare with `benchstat` over `-count=10`. A single run is not evidence. This is advisory until a scheduled run stores an ubuntu-latest/amd64 baseline; the current workflow uploads raw input but intentionally does not fail on this threshold.

## Goroutine gates (blocking, from Phase 2 onward)

| Gate | Threshold |
| ----------------------------------- | ---------------------------------- |
| Goroutines per idle batcher | exactly 0 |
| Goroutines per running `n=1` | exactly 1 (aggregator) |
| Goroutines per running `n>1` | exactly 1 + n |
| Goroutines after terminal `closed` | equal to pre-construction baseline |

Current `main` owns 6 goroutines per batcher: 1 aggregator, 2 `rill` pipeline
relays, and 2 `chann` relays plus the pipeline's batch goroutine. Phase 2.1
removes `rill` (6 → 3) and Phase 2.2 removes `chann` (3 → 1).
## 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.

| 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.

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
forwarder they required (**measured 5 → 2**: aggregator plus processor), enforced
by `TestGoroutineBudgetPerRunningBatcher`.

Removing the relays also removed a channel hop and a goroutine handoff per item.
Measured against the stored baseline: **-54% geomean sec/op** on the enqueue
microbenchmarks (`Add` 229.8ns → 63.4ns at small batch sizes) and -49% to -90%
bytes/op. This is the one place in the plan where a safety change also made the
hot path materially faster.

The 6 → 5 figure supersedes earlier 6 → 3 and 6 → 4 estimates. Fewer goroutines
were unreachable without changing observable behaviour: merging aggregation into
the processing loop inverted the documented latency baseline, and merging input
draining into aggregation regressed sequential `Add` by 39-50% because the
`chann` relay's bounded ingress was not drained promptly.

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.

## Conditional gates (Milestone 4.2 only)

Expand Down
13 changes: 5 additions & 8 deletions go.mod
Original file line number Diff line number Diff line change
@@ -1,21 +1,18 @@
module github.com/NSXBet/batcher

go 1.22.4
go 1.25.0

require (
github.com/destel/rill v0.1.2
github.com/stretchr/testify v1.9.0
go.uber.org/atomic v1.11.0
github.com/stretchr/testify v1.11.1
go.uber.org/fx v1.24.0
go.uber.org/zap v1.27.0
golang.design/x/chann v0.1.2
go.uber.org/zap v1.28.0
)

require (
github.com/davecgh/go-spew v1.1.1 // indirect
github.com/pmezard/go-difflib v1.0.0 // indirect
go.uber.org/dig v1.19.0 // indirect
go.uber.org/multierr v1.10.0 // indirect
golang.org/x/sys v0.0.0-20220412211240-33da011f77ad // indirect
go.uber.org/multierr v1.11.0 // indirect
golang.org/x/sys v0.46.0 // indirect
gopkg.in/yaml.v3 v3.0.1 // indirect
)
Loading
Loading