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
2 changes: 1 addition & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -101,5 +101,5 @@ When adding new stream operations:
2. Implement in `streamer[T]` in `stream.go` (compose an `iter.Seq` closure; keep lazy)
3. Update `sizeHint` propagation deliberately (see table above)
4. For parallel support, accumulate into the fused section (`fused` + `thenFused`) and let `ensureFlushed`/`effectiveSeq` close it. Leak-freedom pattern (keep it): `flushFused`/`orderedSeq` derive a cancellable child ctx; the feeder has a single cancellation exit (top-of-loop check, plain blocking send — workers always drain `in`, discarding after cancel); the consumer drains `out` on any exit. Avoid select-with-Done on the feeder send: when both cases are ready Go picks randomly, which made coverage and exit paths nondeterministic.
5. Add assertion-based tests: `factory_test.go` / `ops_test.go` / `terminal_test.go` / `parallel_test.go` / `branch_test.go` hold the suite (statement coverage 100%); parallel results must be compared as sorted multisets (order is not preserved). `TestParallel_ShortCircuitNoLeak` and `TestParallel_ConcurrentTakeNoRace` guard the concurrency fixes — always run `-race` before shipping parallel changes. `branch_test.go` covers defensive branches (mid-stream cancellation, downstream short-circuit per op, Pick materialize path, foreign Streamer fallback) — extend it when adding new branches. `benchmark_gate_test.go` is the CI performance gate: parallel-vs-serial ratio thresholds (G1 unordered <=6x, G2 ordered <=8x, G3 heavy speedup >=1.5x; thresholds calibrated against measured healthy-v2 ranges on both laptop and shared-CI hardware — see the calibration table in the file) that are machine-independent by construction; if a gate fails on real regressions do not loosen the threshold without a proposal-level justification.
5. Add assertion-based tests: `factory_test.go` / `ops_test.go` / `terminal_test.go` / `parallel_test.go` / `branch_test.go` hold the suite (statement coverage 100%); parallel results must be compared as sorted multisets (order is not preserved). `TestParallel_ShortCircuitNoLeak` and `TestParallel_ConcurrentTakeNoRace` guard the concurrency fixes — always run `-race` before shipping parallel changes. `branch_test.go` covers defensive branches (mid-stream cancellation, downstream short-circuit per op, Pick materialize path, foreign Streamer fallback) — extend it when adding new branches. `benchmark_gate_test.go` is the CI performance gate (G1-G7), machine-independent by construction: same-process ratios (G1 unordered <=6x, G2 ordered <=8x, G3 heavy speedup >=1.5x, G4 sort <=2.5x direct stdlib, G5 DistinctBy <=0.6x Distinct time and <=5% of its allocs) plus deterministic allocation gates (G6 Take <=50 allocs, G7 lazy Limit <=200 allocs on 100k). Thresholds are calibrated against measured healthy ranges on laptop and shared-CI hardware — see the calibration table in the file; if a gate fails on real regressions do not loosen the threshold without a proposal-level justification. Proposal A6 (multi-section speedup) is deliberately not a time gate: it is core-count dependent; its correctness is pinned by TestParallelV2_MidChainParallelReopens.
6. Update `README.md` and `doc.go` documentation
106 changes: 102 additions & 4 deletions benchmark_gate_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,33 +2,51 @@ package stream_test

import (
"runtime"
"slices"
"testing"
"time"

"github.com/tr1v3r/stream"
)

// Machine-independent performance gates for the parallel architecture
// (docs/proposals/parallel-v2.md acceptance A1/A2). Every check compares
// parallel against serial execution IN THE SAME PROCESS, so the ratio — not
// absolute time — is the contract; identical thresholds hold on any runner.
// (docs/proposals/parallel-v2.md acceptance A1/A2) and the other
// performance contracts of the library. Every check is either a RATIO of
// two same-process measurements (hardware cancels out) or an ALLOCATION
// COUNT (deterministic by construction) — never an absolute wall time.
//
// Gates (best-of-3 timings, GC settled before each run):
//
// G1 near-free Filter+Map, Parallel(4) unordered <= 6x serial
// G2 near-free, Parallel(4).Ordered() <= 8x serial
// G3 heavy work, Parallel(4) >= 1.5x faster than serial
// G4 Sort pipeline <= 2.5x direct slices.SortFunc on the same data
// G5 DistinctBy <= 0.6x Distinct time AND <= 5% of its allocations
// G6 Take on 100k elements: <= 50 allocations (reservoir O(1) memory)
// G7 Limit(2).ToSlice() on 100k elements: <= 200 allocations (lazy
// short-circuit must not traverse the source)
//
// Threshold calibration (measured; do not retune without cross-machine data):
//
// G1: Apple M3 Pro 1.1-2.2x | GitHub CI shared runner 4.2x -> gate 6.0
// G2: Apple M3 Pro 0.8-3.2x | GitHub CI shared runner 5.0x -> gate 8.0
// G3: Apple M3 Pro 2.9-3.2x | GitHub CI shared runner 3.0x -> gate 1.5
// G4: Apple M3 Pro 1.17-1.22x (materialize overhead only) -> gate 2.5
// G5: Apple M3 Pro 0.24-0.27x time, 0.3% of Distinct allocs -> gate 0.6 / 0.05
// G6: Apple M3 Pro 6 allocs/op -> gate 50
// G7: Apple M3 Pro 12 allocs/op -> gate 200
//
// Healthy-v2 machinery overhead spans 1-5x across hardware (shared CI
// runners schedule channel machinery far worse than a laptop core); the v1
// per-op-pool regression sat at 14-21x and a fusion failure collapses G3 to
// ~1x — both remain far outside these bounds.
// ~1x — both remain far outside these bounds. Allocation gates are exact and
// hardware-free; a regression to materialize-all Take or eager traversal
// overshoots them by orders of magnitude.
//
// Deliberately NOT a time gate: proposal A6 (multi-section 16+2 vs single
// pool on IO-sim) — its speedup is core-count dependent (16 workers need
// 16 cores to beat 4), so it cannot be machine-independent; multi-section
// correctness is pinned by TestParallelV2_MidChainParallelReopens instead.
//
// Skipped under -race (the detector instruments channel operations and would
// skew machinery ratios) and under -short. CI runs it via the dedicated
Expand Down Expand Up @@ -125,4 +143,84 @@ func TestBenchmarkGate(t *testing.T) {
} else {
t.Logf("G3 pass: heavy speedup %.2fx", float64(serialHeavy)/float64(parallelHeavy))
}

// G4: the Sort pipeline is materialize + slices.SortFunc; it must track
// a hand-rolled clone+sort of the same data. A reintroduced
// sort.Interface adapter (v0.x) or an accidental O(n^2) comparator path
// blows past this ratio.
sortSrc := makeUnsorted(100000)
cmpInt := func(a, b int) int { return a - b }
pipelineSort := gateBestOf(func() {
sink := 0
for _, v := range stream.SliceOf(sortSrc...).Sort(cmpInt).ToSlice() {
sink += v
}
_ = sink
})
directSort := gateBestOf(func() {
data := slices.Clone(sortSrc)
slices.SortFunc(data, cmpInt)
sink := 0
for _, v := range data {
sink += v
}
_ = sink
})
if r := float64(pipelineSort) / float64(directSort); r > 2.5 {
t.Errorf("G4 FAIL: Sort pipeline %.2fx direct slices.SortFunc (gate <= 2.5x)", r)
} else {
t.Logf("G4 pass: sort %.2fx direct", float64(pipelineSort)/float64(directSort))
}

// G5: DistinctBy (comparable keys) vs Distinct (fmt.Sprint keys) — both
// time ratio and allocation ratio on the same input.
distinctSrc := makeCyclic(5000, 1000)
distinctBy := gateBestOf(func() {
sink := 0
for _, v := range stream.DistinctBy(stream.SliceOf(distinctSrc...), func(n int) int { return n }).ToSlice() {
sink += v
}
_ = sink
})
distinct := gateBestOf(func() {
sink := 0
for _, v := range stream.SliceOf(distinctSrc...).Distinct().ToSlice() {
sink += v
}
_ = sink
})
byAllocs := testing.AllocsPerRun(3, func() {
stream.DistinctBy(stream.SliceOf(distinctSrc...), func(n int) int { return n }).ToSlice()
})
defAllocs := testing.AllocsPerRun(3, func() {
stream.SliceOf(distinctSrc...).Distinct().ToSlice()
})
if r := float64(distinctBy) / float64(distinct); r > 0.6 {
t.Errorf("G5 FAIL: DistinctBy %.2fx Distinct time (gate <= 0.6x; measured ~0.2x when healthy)", r)
} else {
t.Logf("G5 pass: DistinctBy %.2fx Distinct time", float64(distinctBy)/float64(distinct))
}
if r := byAllocs / defAllocs; r > 0.05 {
t.Errorf("G5 FAIL: DistinctBy %.1f%% of Distinct allocations (%.0f vs %.0f; gate <= 5%%)", r*100, byAllocs, defAllocs)
} else {
t.Logf("G5 pass: DistinctBy %.1f%% of Distinct allocs", r*100)
}

// G6: Take is reservoir-sampled — O(1) memory regardless of source
// size. A regression to materialize-all costs ~800KB / thousands of
// allocs on this input.
takeSrc := makeUnsorted(100000)
if allocs := testing.AllocsPerRun(3, func() { stream.SliceOf(takeSrc...).Take() }); allocs > 50 {
t.Errorf("G6 FAIL: Take on 100k elements allocated %.0f times (gate <= 50; reservoir is O(1) memory)", allocs)
} else {
t.Logf("G6 pass: Take %.0f allocs", allocs)
}

// G7: lazy short-circuit — Limit(2) must not traverse (let alone
// materialize) the 100k source.
if allocs := testing.AllocsPerRun(3, func() { stream.SliceOf(takeSrc...).Limit(2).ToSlice() }); allocs > 200 {
t.Errorf("G7 FAIL: Limit(2).ToSlice() on 100k elements allocated %.0f times (gate <= 200; eager traversal would be orders more)", allocs)
} else {
t.Logf("G7 pass: lazy short-circuit %.0f allocs", allocs)
}
}
Loading