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
248
249
250
251
252
253
254
255
256
257
258
259
260
//! Task sets, joining, and bounded-concurrency fan-out.
//!
//! This module deliberately exports **no function that hands back a handle to a
//! running task**. There is no `spawn`. A `JoinHandle` is droppable, and a
//! dropped handle is a task whose outcome nobody ever sees: the task keeps
//! running, the caller never learns whether it succeeded, and the only record
//! that it existed is gone. That is the failure the crate's background-work
//! rule exists to prevent, and a convention ("remember to join it") is not a
//! guarantee. The hole is closed by removing the constructor, not by
//! documenting it.
//!
//! Concurrent work is started two ways, and both own what they start:
//!
//! - [`Supervisor`] — a bounded set of tasks
//! that reports every terminal outcome ([`TaskOutcome`]) and stops its tasks
//! when it is dropped or cancelled. This is the module a bot's background work
//! belongs in, and the only place a subprocess is started (`spawn_process`,
//! behind the `process` feature).
//! - [`JoinSet`] — the tracked envelope itself. Its constructor takes the tasks,
//! its drop aborts them, and `join_next` yields [`JoinError`]-carrying results
//! to the caller, so nothing is started that nothing owns.
//!
//! [`Supervisor`]: crate::rt::supervise::Supervisor
//! [`TaskOutcome`]: crate::rt::supervise::TaskOutcome
//!
//! What remains here besides the set is the fan-out an agent SDK actually
//! needs: [`join_all_bounded`] runs many futures, never exceeds a concurrency
//! ceiling, and returns results in input order regardless of completion order.
//!
//! # Starting work, at a glance
//!
//! - **One future, and you need its value.** Run it: `join_all_bounded(1,
//! [work()]).await` returns a one-element vector, or await the future
//! directly. There is no handle in either form, so there is nothing to drop.
//! - **One task now, read later.** Build a [`JoinSet`], [`JoinSet::spawn`] into
//! it, and `join_next` when you want the result. Dropping the set aborts
//! whatever has not finished.
//! - **Background work that must be cancellable and accounted for.** Use
//! [`Supervisor`]: it awaits a permit
//! before it starts anything, reaps finished tasks, and reports how each one
//! ended.
//!
//! # Non-`Send` futures
//!
//! Every verb in this crate is deliberately **not** `Send`, so a bot's own
//! futures cannot be placed on another worker thread — that is why
//! [`BoxFuture`](crate::BoxFuture) is unconstrained. A non-`Send` future is
//! driven by awaiting it, or by
//! [`lgwks_std::task::join_all`], which polls many
//! futures on the calling thread and returns their outputs in input order.
//!
//! There is no public API here that *spawns* a non-`Send` future. The
//! single-threaded spawn was the same un-owned handle as the multi-threaded one
//! — a caller could start a local task and forget it just as easily, with the
//! added trap that the handle's type could not even be named at the call site —
//! so it is gone for the same reason: nothing may be started that nothing owns.
use Id;
pub use ;
use Future;
/// A task owner with a deliberately narrow API.
///
/// The engine's raw `JoinSet` also exposes `detach_all`, which removes tasks
/// without aborting them. This wrapper does not: every task remains owned by
/// this set until joined, aborted, or the set is dropped (which aborts its
/// remaining tasks). Blocking and local tasks are likewise absent; blocking
/// work needs an explicit non-preemptible owner, not a method on this async
/// owner.
/// Run `futures` with at most `limit` in flight at once, returning their
/// outputs in input order.
///
/// Every input is polled to completion; `limit` bounds concurrency, not the
/// number of futures accepted. `limit` of zero is treated as one, and a `limit`
/// above [`Semaphore::MAX_PERMITS`][max] is clamped to it. An empty input
/// resolves immediately.
///
/// [max]: lgwks_deps::tokio::sync::Semaphore::MAX_PERMITS
///
/// # Ordering
///
/// Output `i` is the output of input `i`, whichever completes first. This is
/// the property a plain [`JoinSet`] does not give: `JoinSet::join_next` yields
/// completion order.
///
/// # Concurrency and memory bound
///
/// The bound is enforced by replenishment: at most `limit` tasks are spawned
/// and awaited at once, and each completion spawns the next pending input. The
/// retained [`JoinSet`] therefore never exceeds `limit` entries, and only the
/// output vector grows with input length, so a caller fanning out over thousands
/// of inputs needs no manual chunking to keep task memory bounded.
///
/// # Cancellation and failure
///
/// Dropping the returned future drops the [`JoinSet`], which aborts every task
/// that has not finished. A panicking input is resumed on the *awaiting* task,
/// matching `join_all`: it is not converted into a [`JoinError`], and it does
/// not abort the process. The remaining tasks are then aborted as the set
/// drops. A completed input is never silently dropped from the result vector.
pub async