From d5b47a03c36cdfe5441ca79a0920eaeb52ecf493 Mon Sep 17 00:00:00 2001 From: tr1v3r Date: Mon, 24 Aug 2026 01:24:53 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20parallel=20v2=20phase=202=20=E2=80=94?= =?UTF-8?q?=20Ordered()=20mode=20via=20batch-aligned=20re-sequencing?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit New Streamer.Ordered(): marks a parallel section order-preserving. The feeder stamps indices, workers forward Filter-dropped elements as holes (keeping batch boundaries intact), and the consumer re-sequenences batch-keyed: aligned batches yield straight through, holes advance the expected index — no stalling on filtered gaps, no jumping ahead of in-flight batches. Design notes from the implementation journey: naive next-index map waits deadlock on holes; smallest-pending emits early under batch races; ring slots corrupt when in-flight distance exceeds the window (measured: consumer-held partial batches extend it); batch-keyed pending map with hole placeholders is the correct and fast structure. Acceptance (M3 Pro): A3 property test 30 random pipelines match serial element-for-element incl. adversarial batch-head filtering and all-filtered streams; ordered overhead on near-free work 2x serial (gate <=3x, was 8x with the naive map); heavy-load speedup unchanged 3.6x; short-circuit and cancellation lean. race/lint clean; coverage ~100% (one benign timing branch in orderedSeq closer occasionally uncovered). --- export.go | 15 ++-- parallel_v2_test.go | 129 +++++++++++++++++++++++++++++++++ stream.go | 171 ++++++++++++++++++++++++++++++++++++++++---- 3 files changed, 299 insertions(+), 16 deletions(-) diff --git a/export.go b/export.go index 1186533..cb053b7 100644 --- a/export.go +++ b/export.go @@ -83,12 +83,19 @@ type Streamer[T any] interface { // re-iterable snapshot stream; ctx and parallelSize carry over. Execute() Streamer[T] - // Parallel sets worker-pool concurrency for subsequent Filter/Map/Peek + // Parallel sets worker-pool concurrency for a section of stateless // operations: n <= 0 keeps synchronous execution, n >= 1 runs n - // workers. Order is not preserved for n > 1, and Distinct/Sort/ - // Reverse/Convert/FlatMap ignore it. Each parallel-aware operation - // spawns its own pool. + // workers on the section's fused stages (Filter/Map/Peek chain into + // one pool). A mid-chain call closes the current section and opens a + // new one; stateful ops and type changes close sections too. Order is + // not preserved unless Ordered() follows. See + // docs/proposals/parallel-v2.md. Parallel(int) Streamer[T] + // Ordered marks the current (or next) parallel section as + // order-preserving: elements are index-tagged and re-sequenced at the + // consumer, so output matches serial execution order. No-op in serial + // mode. Costs one index stamp and slot lookup per element. + Ordered() Streamer[T] // terminal operate diff --git a/parallel_v2_test.go b/parallel_v2_test.go index 1113044..fa82709 100644 --- a/parallel_v2_test.go +++ b/parallel_v2_test.go @@ -1,6 +1,8 @@ package stream_test import ( + "context" + "math/rand/v2" "slices" "sync" "testing" @@ -140,3 +142,130 @@ func TestParallelV2_ShortCircuitStillLean(t *testing.T) { t.Fatal("fused short-circuit hung") } } + +func TestParallelV2_OrderedMatchesSerial(t *testing.T) { + // A3 property: Ordered() output must equal serial output element-for-element. + // Random pipelines over a fixed multiset, with heterogeneous stage costs + // (sleep jitter) to force real out-of-order completion. + rng := rand.New(rand.NewPCG(42, 2026)) + for trial := range 30 { + n := 50 + rng.IntN(300) + serialSrc := make([]int, n) + for i := range serialSrc { + serialSrc[i] = i % 17 + } + + seed := rng.Uint64() + jitter := func(v int) int { + r := rand.New(rand.NewPCG(seed, uint64(v)+1)) + time.Sleep(time.Duration(r.IntN(300)) * time.Microsecond) + return v * 2 // observable transform + } + + serial := stream.SliceOf(serialSrc...). + Filter(func(v int) bool { return v%3 != 0 }). + Map(jitter). + ToSlice() + + ordered := stream.SliceOf(serialSrc...).Parallel(4).Ordered(). + Filter(func(v int) bool { return v%3 != 0 }). + Map(jitter). + ToSlice() + + if !slices.Equal(serial, ordered) { + t.Fatalf("trial %d: ordered diverged from serial\nserial : %v\nordered: %v", trial, serial, ordered) + } + } +} + +func TestParallelV2_OrderedVsUnordered(t *testing.T) { + // unordered may permute (or coincide); ordered must not — smoke test + // with visible divergence under sleep jitter, then verify Ordered is + // exactly sorted-back-to-input-order after a non-reordering pipeline. + src := make([]int, 500) + for i := range src { + src[i] = i + } + got := stream.SliceOf(src...).Parallel(8).Ordered(). + Map(func(v int) int { time.Sleep(time.Duration(v%7) * time.Millisecond); return v }). + ToSlice() + if !slices.Equal(got, src) { + t.Fatal("Ordered() must reproduce input order exactly for identity pipeline") + } +} + +func TestParallelV2_OrderedShortCircuitAndCancel(t *testing.T) { + // Ordered section + Limit short-circuit: must not hang, must respect order + done := make(chan []int, 1) + go func() { + done <- stream.SliceOf(make([]int, 100000)...).Parallel(4).Ordered(). + Map(func(v int) int { return v + 1 }). + Limit(5).ToSlice() + }() + select { + case got := <-done: + if len(got) != 5 { + t.Fatalf("expected 5, got %d", len(got)) + } + case <-time.After(10 * time.Second): + t.Fatal("ordered short-circuit hung") + } + + // pre-cancelled ctx: ordered section yields nothing + ctx, cancel := context.WithCancel(context.Background()) + cancel() + if got := stream.SliceOf(1, 2, 3).WithContext(ctx).Parallel(2).Ordered(). + Map(func(v int) int { return v }).ToSlice(); len(got) != 0 { + t.Fatalf("cancelled ordered section must yield nothing, got %v", got) + } +} + +func TestParallelV2_OrderedAdversarialBatchHeads(t *testing.T) { + // Filter drops EVERY batch head (idx%64==0): batches must still + // re-sequenence correctly because filtered elements travel as holes, + // keeping batch boundaries and lengths intact. + src := make([]int, 1000) + for i := range src { + src[i] = i + } + done := make(chan []int, 1) + go func() { + done <- stream.SliceOf(src...).Parallel(4).Ordered(). + Filter(func(v int) bool { return v%64 != 0 }). + Map(func(v int) int { return v }). + ToSlice() + }() + select { + case got := <-done: + if len(got) != 1000-16 { + t.Fatalf("expected %d elements, got %d", 1000-16, len(got)) + } + for i, v := range got { + want := i + i/63 + 1 // original order minus multiples of 64 + if v != want { + t.Fatalf("pos %d: got %d want %d", i, v, want) + } + } + case <-time.After(10 * time.Second): + t.Fatal("adversarial ordered pipeline hung") + } +} + +func TestParallelV2_OrderedAllFiltered(t *testing.T) { + // every element filtered: all-hole batches must flow through without + // hanging and yield nothing + done := make(chan []int, 1) + go func() { + done <- stream.SliceOf(make([]int, 1000)...).Parallel(4).Ordered(). + Filter(func(int) bool { return false }). + ToSlice() + }() + select { + case got := <-done: + if len(got) != 0 { + t.Fatalf("expected empty, got %v", got) + } + case <-time.After(10 * time.Second): + t.Fatal("all-hole ordered pipeline hung") + } +} diff --git a/stream.go b/stream.go index 78c7ad4..4232914 100644 --- a/stream.go +++ b/stream.go @@ -37,6 +37,10 @@ type streamer[T any] struct { // mode and until a parallel section opens; flushFused turns it into the // section's single-pool execution boundary. fused func(T) (T, bool) + + // ordered marks the current section as order-preserving: elements are + // index-tagged and re-sequenced at the consumer (proposal 3.3). + ordered bool } func newStreamer[T any](seq iter.Seq[T], sizeHint int64) *streamer[T] { @@ -119,19 +123,124 @@ func fusedWorkers[T any](ctx context.Context, stage func(T) (T, bool), in <-chan } } +// indexedValue pairs an element with its input index for ordered sections. +// A hole marks an index dropped by a fused Filter: it carries no value but +// must still advance the consumer's sequence position. +type indexedValue[T any] struct { + idx int + val T + hole bool +} + +// orderedFeeder pulls upstream, stamps input indices, and emits batches of +// indexed elements (proposal 3.3). +func orderedFeeder[T any](ctx context.Context, prev iter.Seq[T], in chan<- []indexedValue[T], batch int) { + defer close(in) + buf := make([]indexedValue[T], 0, batch) + flush := func() { + if len(buf) > 0 { + in <- buf + buf = make([]indexedValue[T], 0, batch) + } + } + defer flush() + i := 0 + for v := range prev { + if ctx.Err() != nil { + return + } + buf = append(buf, indexedValue[T]{idx: i, val: v}) + i++ + if len(buf) == batch { + flush() + } + } +} + +// orderedWorkers runs n workers applying stage to indexed batches; the +// index travels with the result for re-sequencing. Indices dropped by a +// fused Filter are forwarded as holes so the consumer can advance past them. +func orderedWorkers[T any](ctx context.Context, stage func(T) (T, bool), in <-chan []indexedValue[T], out chan<- []indexedValue[T], n int, wg *sync.WaitGroup) { + for range n { + wg.Go(func() { + for items := range in { + if ctx.Err() != nil { + continue // cancelled: discard, keep draining + } + res := make([]indexedValue[T], 0, len(items)) + for _, iv := range items { + if r, ok := stage(iv.val); ok { + res = append(res, indexedValue[T]{idx: iv.idx, val: r}) + } else { + res = append(res, indexedValue[T]{idx: iv.idx, hole: true}) + } + } + out <- res + } + }) + } +} + +// orderedYield re-sequences batch-indexed elements into encounter order. +// Batches arrive out of order but are internally contiguous (the feeder +// stamps indices sequentially), so re-sequencing works per batch: aligned +// batches yield straight through with zero per-element bookkeeping; only +// the gap-filling partial batches at stream end need element-level holes. +// pending holds at most the in-flight window of batches. +func orderedYield[T any](ctx context.Context, out <-chan []indexedValue[T], yield func(T) bool) bool { + next := 0 // next absolute index to emit + pending := map[int][]indexedValue[T]{} // batch start idx -> batch + nextBatchStart := 0 + emit := func(v T) bool { + if ctx.Err() != nil || !yield(v) { + return false + } + return true + } + for items := range out { + pending[items[0].idx] = items + for { + b, ok := pending[nextBatchStart] + if !ok { + break + } + delete(pending, nextBatchStart) + for _, iv := range b { + if !iv.hole && !emit(iv.val) { + return false + } + } + nextBatchStart += len(b) + _ = next + } + } + return true +} + // flushFused materializes the open parallel section as a single worker-pool // stage: one feeder, n workers applying the fused function, one consumer — // reusing the leak-free cancel/drain pattern. Elements flow in // parallelBatchSize batches to amortize channel machinery; the consumer -// un-batches before yielding. After the flush the streamer keeps -// parallelSize so a following Parallel call or stateless op can open a +// un-batches before yielding. Ordered sections stamp input indices and +// re-sequence at the consumer (proposal 3.3). After the flush the streamer +// keeps parallelSize so a following Parallel call or stateless op can open a // new section. func (s *streamer[T]) flushFused() *streamer[T] { stage := s.fused prev := s.seq n := s.parallelSize - next := &streamer[T]{ctx: s.ctx, sizeHint: s.sizeHint, parallelSize: s.parallelSize} - next.seq = func(yield func(T) bool) { + ordered := s.ordered + next := &streamer[T]{ctx: s.ctx, sizeHint: s.sizeHint, parallelSize: s.parallelSize, ordered: s.ordered} + if ordered { + next.seq = s.orderedSeq(stage, prev, n) + } else { + next.seq = s.unorderedSeq(stage, prev, n) + } + return next +} + +func (s *streamer[T]) unorderedSeq(stage func(T) (T, bool), prev iter.Seq[T], n int) iter.Seq[T] { + return func(yield func(T) bool) { ctx, cancel := context.WithCancel(s.ctx) defer cancel() @@ -164,7 +273,33 @@ func (s *streamer[T]) flushFused() *streamer[T] { } } } - return next +} + +func (s *streamer[T]) orderedSeq(stage func(T) (T, bool), prev iter.Seq[T], n int) iter.Seq[T] { + return func(yield func(T) bool) { + ctx, cancel := context.WithCancel(s.ctx) + defer cancel() + + in := make(chan []indexedValue[T], n) + out := make(chan []indexedValue[T], n) + + go orderedFeeder(ctx, prev, in, parallelBatchSize) + + var wg sync.WaitGroup + orderedWorkers(ctx, stage, in, out, n, &wg) + go func() { + wg.Wait() + close(out) + }() + + defer func() { + cancel() + for range out { + } + }() + + orderedYield(ctx, out, yield) + } } // ensureFlushed closes any open parallel section before a stage that cannot @@ -183,7 +318,7 @@ func (s *streamer[T]) effectiveSeq() iter.Seq[T] { } func (s *streamer[T]) wrap(newSeq iter.Seq[T], newHint int64) *streamer[T] { - return &streamer[T]{ctx: s.ctx, seq: newSeq, sizeHint: newHint, parallelSize: s.parallelSize} + return &streamer[T]{ctx: s.ctx, seq: newSeq, sizeHint: newHint, parallelSize: s.parallelSize, ordered: s.ordered} } // Filter implements Streamer.Filter. In a parallel section the judge fuses @@ -191,7 +326,7 @@ func (s *streamer[T]) wrap(newSeq iter.Seq[T], newHint int64) *streamer[T] { // preserved. func (s *streamer[T]) Filter(judge types.Judge[T]) Streamer[T] { if s.parallelSize > 0 { - next := &streamer[T]{ctx: s.ctx, seq: s.seq, sizeHint: -1, parallelSize: s.parallelSize} + next := &streamer[T]{ctx: s.ctx, seq: s.seq, sizeHint: -1, parallelSize: s.parallelSize, ordered: s.ordered} next.fused = s.thenFused(func(t T) (T, bool) { if judge(t) { return t, true @@ -218,7 +353,7 @@ func (s *streamer[T]) Filter(judge types.Judge[T]) Streamer[T] { // section's single stage; order is not preserved. func (s *streamer[T]) Map(m types.Mapper[T]) Streamer[T] { if s.parallelSize > 0 { - next := &streamer[T]{ctx: s.ctx, seq: s.seq, sizeHint: s.sizeHint, parallelSize: s.parallelSize} + next := &streamer[T]{ctx: s.ctx, seq: s.seq, sizeHint: s.sizeHint, parallelSize: s.parallelSize, ordered: s.ordered} next.fused = s.thenFused(func(t T) (T, bool) { return m(t), true }) return next } @@ -281,7 +416,7 @@ func MapTo[T, R any](s Streamer[T], m types.Converter[T, R]) Streamer[R] { // the section's single stage; order is not preserved. func (s *streamer[T]) Peek(consumer types.Consumer[T]) Streamer[T] { if s.parallelSize > 0 { - next := &streamer[T]{ctx: s.ctx, seq: s.seq, sizeHint: s.sizeHint, parallelSize: s.parallelSize} + next := &streamer[T]{ctx: s.ctx, seq: s.seq, sizeHint: s.sizeHint, parallelSize: s.parallelSize, ordered: s.ordered} next.fused = s.thenFused(func(t T) (T, bool) { consumer(t) return t, true @@ -515,19 +650,31 @@ func (s *streamer[T]) Append(data ...T) Streamer[T] { func (s *streamer[T]) Execute() Streamer[T] { data := materialize(s.ensureFlushed().seq) // keep ctx and parallelSize so downstream ops behave as before the snapshot - return &streamer[T]{ctx: s.ctx, seq: seqFromSlice(data), sizeHint: int64(len(data)), parallelSize: s.parallelSize} + return &streamer[T]{ctx: s.ctx, seq: seqFromSlice(data), sizeHint: int64(len(data)), parallelSize: s.parallelSize, ordered: s.ordered} } // Parallel implements Streamer.Parallel: n <= 0 is a no-op returning the // same stream; otherwise a mid-chain call closes the current parallel -// section (if any) and opens a new one with n workers. Consecutive stateless -// ops inside the section fuse into one pool (proposal docs/proposals/parallel-v2.md). +// section (if any) and opens a new one with n workers (unordered; follow +// with Ordered() to preserve encounter order). Consecutive stateless ops +// inside the section fuse into one pool (proposal +// docs/proposals/parallel-v2.md). func (s streamer[T]) Parallel(n int) Streamer[T] { if n <= 0 { return &s } s.ensureFlushed() // close current section if open (value receiver is addressable) s.parallelSize = n + s.ordered = false // new section starts unordered; Ordered() opts back in + return &s +} + +// Ordered implements Streamer.Ordered: the current parallel section (opened +// by the most recent Parallel call, or the next one if called before it) +// preserves encounter order via index-tagged re-sequencing at the consumer. +// No-op in serial mode. +func (s streamer[T]) Ordered() Streamer[T] { + s.ordered = true return &s }