From ea94697740f5f9392142938fb24acff4c919bded Mon Sep 17 00:00:00 2001 From: tr1v3r Date: Mon, 14 Sep 2026 10:11:40 +0800 Subject: [PATCH] fix: close parallel section on mid-chain Parallel MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Parallel() invoked ensureFlushed() but discarded the returned streamer, so a mid-chain call never closed the open section: fused stages from both sections merged into one pool sized by the second Parallel's n, and an Ordered() section silently lost its guarantee — Parallel(2).Ordered().Map(f).Parallel(4).Limit(3) selected the wrong elements ([192 193 194] instead of [0 1 2]). Adopt the flushed streamer so the section's output becomes the new upstream; adjacent sections overlap lazily as documented. Adds three H2 regressions (separate pools, ordered-section survival, chained-section short- circuit without goroutine leaks). Audit finding H2. --- parallel_v2_test.go | 102 ++++++++++++++++++++++++++++++++++++++++++++ stream.go | 12 ++++-- 2 files changed, 110 insertions(+), 4 deletions(-) diff --git a/parallel_v2_test.go b/parallel_v2_test.go index fa82709..cf0ecbf 100644 --- a/parallel_v2_test.go +++ b/parallel_v2_test.go @@ -3,6 +3,7 @@ package stream_test import ( "context" "math/rand/v2" + "runtime" "slices" "sync" "testing" @@ -269,3 +270,104 @@ func TestParallelV2_OrderedAllFiltered(t *testing.T) { t.Fatal("all-hole ordered pipeline hung") } } + +// H2 regression: a mid-chain Parallel(n) must actually close the current +// section — each section runs on its own pool sized by its own Parallel +// call, never by a later one (export.go: "A mid-chain call closes the +// current section and opens a new one"). +func TestParallelV2_MidChainParallelSeparatePools(t *testing.T) { + data := make([]int, 400) + for i := range data { + data[i] = i + } + var mu sync.Mutex + active, peak := 0, 0 + slow := func(v int) int { + mu.Lock() + active++ + if active > peak { + peak = active + } + mu.Unlock() + time.Sleep(2 * time.Millisecond) + mu.Lock() + active-- + mu.Unlock() + return v + } + got := stream.SliceOf(data...). + Parallel(2). // section A: slow() must run on exactly 2 workers + Map(slow). + Parallel(8). // closes A, opens section B on A's output + Map(func(v int) int { return v * 10 }). + ToSlice() + mu.Lock() + p := peak + mu.Unlock() + if len(got) != 400 { + t.Fatalf("two-section pipeline lost elements: %d elems", len(got)) + } + slices.Sort(got) + if got[0] != 0 || got[399] != 3990 { + t.Fatalf("two-section pipeline corrupted data: [%d..%d]", got[0], got[399]) + } + if p > 2 { + t.Fatalf("section A must stay on its own 2-worker pool: peak concurrency %d > 2 (sections merged)", p) + } + if p < 2 { + t.Fatalf("section A pool of 2 never overlapped: peak %d", p) + } +} + +// H2 regression: an Ordered() section keeps its ordering guarantee across a +// mid-chain Parallel — a Limit downstream of the second section must select +// the encounter-order head, not whichever elements finished first. +func TestParallelV2_MidChainParallelKeepsOrderedSection(t *testing.T) { + const n = 600 + src := make([]int, n) + for i := range src { + src[i] = i + } + // reversed cost: early elements slow, late fast — any ordering lapse + // lets late elements overtake and become visible in the first three. + jitter := func(v int) int { + time.Sleep(time.Duration(n-v) * 60 * time.Microsecond) + return v + } + got := stream.SliceOf(src...). + Parallel(2).Ordered(). + Map(jitter). + Parallel(4). + Limit(3). + ToSlice() + if !slices.Equal(got, []int{0, 1, 2}) { + t.Fatalf("ordered section + mid-chain Parallel + Limit: got %v, want [0 1 2]", got) + } +} + +// H2 regression: short-circuiting downstream of two chained sections must +// unwind both pools without leaking goroutines — the downstream section's +// cancellation propagates through the upstream section's consumer loop. +func TestParallelV2_MidChainShortCircuitNoLeak(t *testing.T) { + runtime.GC() + time.Sleep(50 * time.Millisecond) + before := runtime.NumGoroutine() + + if got := stream.Repeat(1). + Parallel(2). + Map(func(n int) int { time.Sleep(5 * time.Millisecond); return n + 1 }). + Parallel(4). + Map(func(n int) int { time.Sleep(2 * time.Millisecond); return n * 2 }). + Limit(3).ToSlice(); len(got) != 3 { + t.Fatalf("expected 3 elements, got %d", len(got)) + } + + deadline := time.Now().Add(5 * time.Second) + for runtime.NumGoroutine() > before && time.Now().Before(deadline) { + time.Sleep(10 * time.Millisecond) + runtime.GC() + } + if n := runtime.NumGoroutine(); n > before { + t.Fatalf("goroutine leak across chained sections: before=%d after=%d", before, n) + } +} diff --git a/stream.go b/stream.go index 4232914..319b776 100644 --- a/stream.go +++ b/stream.go @@ -656,14 +656,18 @@ 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). Consecutive stateless ops -// inside the section fuse into one pool (proposal -// docs/proposals/parallel-v2.md). +// 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 +// (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) + // Close the current section if one is open: adopt the flushed pipeline + // (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