From cf45e4e11e89a4caf89e1de41575e9d722315a0f Mon Sep 17 00:00:00 2001 From: tr1v3r Date: Mon, 14 Sep 2026 10:12:46 +0800 Subject: [PATCH] fix: preserve Ordered() set before Parallel() Parallel() unconditionally reset s.ordered=false, contradicting the Ordered() contract ("the current parallel section ... or the next one if called before it"): Ordered().Parallel(n) silently lost encounter order. The new section's default-unordered comes from the zero value, so an explicit Ordered() now carries into the section regardless of call order. Adds TestParallelV2_OrderedBeforeParallel. Audit finding M1 (stacked on fix/midchain-parallel-section). --- parallel_v2_test.go | 18 ++++++++++++++++++ stream.go | 12 ++++++------ 2 files changed, 24 insertions(+), 6 deletions(-) diff --git a/parallel_v2_test.go b/parallel_v2_test.go index cf0ecbf..433ee24 100644 --- a/parallel_v2_test.go +++ b/parallel_v2_test.go @@ -371,3 +371,21 @@ func TestParallelV2_MidChainShortCircuitNoLeak(t *testing.T) { t.Fatalf("goroutine leak across chained sections: before=%d after=%d", before, n) } } + +// M1 regression: Ordered() before Parallel() applies to the section that +// Parallel opens — export.go: "the current (or next) parallel section". +func TestParallelV2_OrderedBeforeParallel(t *testing.T) { + const n = 600 + src := make([]int, n) + for i := range src { + src[i] = i + } + jitter := func(v int) int { + time.Sleep(time.Duration(n-v) * 60 * time.Microsecond) + return v + } + got := stream.SliceOf(src...).Ordered().Parallel(4).Map(jitter).ToSlice() + if !slices.Equal(got, src) { + t.Fatalf("Ordered() before Parallel(): output must reproduce encounter order") + } +} diff --git a/stream.go b/stream.go index 877fb2c..69a5ea9 100644 --- a/stream.go +++ b/stream.go @@ -676,11 +676,12 @@ func (s *streamer[T]) Execute() Streamer[T] { // 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 (unordered; follow -// with Ordered() to preserve encounter order). Each section runs on its own -// pool sized by its own Parallel call; adjacent sections overlap because -// the downstream section's feeder pulls the upstream section's output -// lazily. Consecutive stateless ops inside the section fuse into one pool +// section (if any) and opens a new one with n workers. Each section runs on +// its own pool sized by its own Parallel call; adjacent sections overlap +// because the downstream section's feeder pulls the upstream section's +// output lazily. Sections run unordered unless Ordered() was set on the +// stream (before or after the Parallel call — the flag carries into the +// section). 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 { @@ -690,7 +691,6 @@ func (s streamer[T]) Parallel(n int) Streamer[T] { // (the section's output becomes the new upstream) and drop fused stages. s = *s.ensureFlushed() s.parallelSize = n - s.ordered = false // new section starts unordered; Ordered() opts back in return &s }