1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
//! Concurrent stream combinators dispatched to the unified hybrid scheduler.
//!
//! The caller expresses **how much concurrency the work warrants** via `limit`,
//! and the combinator does no more than that:
//!
//! - `limit == 1` — the work does not warrant concurrency: each item future runs
//! **inline and sequentially** on the consuming thread, with no spawn, no
//! cross-thread hop, and no channel. This is the zero-overhead path.
//! - `limit > 1` — items are **distributed across the scheduler's worker
//! threads**, up to `limit` in flight, so the hybrid `ThreadScheduler` runs
//! them in parallel.
//!
//! The API says `concurrent`, not `parallel`, because at `limit > 1` an item
//! future may still be I/O-bound and never saturate a core; concurrency is the
//! contract, the execution mechanism is the scheduler's to optimize.
//!
//! Two ordering disciplines, same dispatch:
//!
//! - [`concurrent_map`](ConcurrentStreamExt::concurrent_map) yields in
//! **completion order** ([`StreamExt::buffer_unordered`]) — no head-of-line
//! blocking, maximum throughput.
//! - [`concurrent_map_ordered`](ConcurrentStreamExt::concurrent_map_ordered)
//! yields in **input order** through retained bounded future slots — a slow
//! early item delays later-completed items, but order is preserved without a
//! heap task node per input.
//!
//! Only operations whose per-item work is heavy enough to outweigh a thread hop
//! belong here. **Cheap, sequential operations — filtering on a simple
//! predicate, light maps — should use the standard [`StreamExt`] combinators**
//! (`map`, `filter`, `filter_map`), which run inline; they compose directly with
//! `concurrent_map`, e.g. `stream.concurrent_map(n, heavy).filter_map(cheap)`.
//!
//! Design, building on the lessons of the [`parallel-stream`](https://docs.rs/parallel-stream)
//! crate but routed through moirai's own infrastructure:
//!
//! - **Unified scheduler.** At `limit > 1` each item future is spawned on
//! [`moirai_executor::global()`] — the same work-stealing scheduler that backs
//! `spawn_async` and the parallel iterators. The result is handed back through
//! a one-shot channel ([`ScheduledItem`]) so the consuming stream awaits it
//! *cooperatively*; it never blocks a worker the way `TaskHandle::join` would.
//! - **Bounded by construction.** `limit` caps in-flight item futures — the
//! central lesson from `parallel-stream`: stream fan-out must be bounded,
//! never unbounded.
//! - **Monomorphized, zero-cost.** Generic over the stream, the closure, and the
//! item future; the spawned result is a named [`ScheduledItem`] future and the
//! sequential/distributed split is a [`Either`], so there is no `Box<dyn>` on
//! the data path.
//!
//! # Performance — sizing `limit`
//!
//! The distributed path (`limit > 1`) pays a per-item dispatch cost: the spawn,
//! the one-shot hand-back, and a cross-thread wake. Measured on a 24-core
//! x86-64 box with identity item work (`tests/stream_overhead.rs`): ~0.8 µs per
//! item distributed, versus ~70 ns per item on the `limit == 1` inline path
//! (~12× cheaper).
//!
//! Of that ~0.8 µs the one-shot is only ~0.1 µs (~1/8); the scheduler dispatch —
//! task allocation, enqueue, and cross-thread wake — dominates the rest. So the
//! per-item overhead is intrinsic to handing one item to another worker, not an
//! artifact this layer can cheaply shave.
//!
//! Consequence: distribution is a net win only when **per-item work exceeds
//! roughly a microsecond**. Below that the dispatch overhead dominates, and the
//! right tool depends on the shape of the work:
//!
//! - **Lazy, item-at-a-time, latency-bound (I/O):** prefer `limit == 1` or the
//! inline [`StreamExt`] combinators — the work, when it is real, amortizes the
//! hop itself.
//! - **Eager bulk of many cheap CPU items:** use the parallel iterators
//! ([`crate::moirai_iter_parallel`]), which **chunk** the data and spawn once
//! per chunk — amortizing the same dispatch cost over hundreds of items. This
//! stream deliberately trades that batching away for per-item laziness; it is
//! not the tool for bulk cheap compute, and does not duplicate it.
//!
//! ```no_run
//! use futures::StreamExt;
//! use moirai_iter::stream::ConcurrentStreamExt;
//!
//! # async fn demo() {
//! let source = futures::stream::iter(0..1_000u64);
//! // Up to 16 item futures in flight; the scheduler distributes them.
//! let doubled: Vec<u64> = source.concurrent_map(16, |x| async move { x * 2 }).collect().await;
//! # let _ = doubled;
//! # }
//! ```
use Future;
use Pin;
use ;
use oneshot;
use Either;
use ;
use TaskSpawner;
pub use ;
/// A future resolving to the output of an item spawned on the unified scheduler.
///
/// Awaiting it is *cooperative*: it registers the task waker on a one-shot
/// channel and yields the worker, rather than blocking it the way
/// `TaskHandle::join` would. Zero-cost — a single channel receiver, no `Box`.
/// Spawn `fut` on the global unified scheduler, yielding a [`ScheduledItem`] the
/// consuming stream can await cooperatively.
/// Shared dispatch: turn a stream of items into a stream of per-item futures,
/// each either run inline (`limit == 1`, no concurrency requested) or spawned on
/// the scheduler (`limit > 1`). The caller applies the bound through its
/// ordered or completion-order buffering policy; this is the single place the
/// inline/distributed choice is made, so both combinators share it.
+ Send
where
S: Stream + Send + 'static,
Item: Send + 'static,
F: FnMut + Send + 'static,
Fut: + Send + 'static,
R: Send + 'static,
/// Concurrent [`Stream`] combinators dispatched through the unified hybrid
/// scheduler.
///
/// Implemented for every [`Stream`]; bring it into scope to call the
/// `concurrent_*` methods on any stream. See the [module docs](self) for when to
/// reach for these versus the inline [`StreamExt`] combinators, and for the
/// ordered/unordered distinction.