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
18 changes: 18 additions & 0 deletions parallel_v2_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
}
}
12 changes: 6 additions & 6 deletions stream.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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
}

Expand Down
Loading