Oregami
Repositories/oxedyne/daimond

oxedyne/daimond/src/batch.rs

14.8 KiB, 1 run

created by r2519314175:935, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1//! Which of a round's tool calls may run at the same time.
2//!
3//! A model routinely asks for several tools in one reply, and Daimond ran them one after
4//! another: the round's wall clock was the SUM of the calls where it could be the MAXIMUM.
5//! Three reads of 400 ms were 1.2 seconds instead of 400 ms, every round, for the whole loop.
6//!
7//! ## THE RULE
8//!
9//! **A batch may hold only calls that read, and nothing else.** A call joins the batch running
10//! beside it only when all four are true of it:
11//!
12//! 1. It touches no network, so it cannot reach a destination the model chose.
13//! 2. It puts no question to the user, so two of them cannot race one dialog.
14//! 3. It changes nothing -- no file, no crystal, no store, no ledger, no open page, no panel.
15//! 4. It reads no flag that another call in the round can set.
16//!
17//! **Everything else runs ALONE, in the order the model gave.** That is a narrowing of what runs
18//! together, not a rewrite of what runs: a call this cannot confidently place in a batch keeps
19//! exactly the behaviour it has today, which is a batch of one.
20//!
21//! ## WHY THE FOUR, EACH IN TURN
22//!
23//! **Writes to one target must not race.** Two `file_edit`s of one file, or a write and a read of
24//! one path, have an order the model intended and expressed by emitting them together. No
25//! intersection of paths is computed here and none is needed: a write is never in a batch, so a
26//! write always separates the reads before it from the reads after it, and the model's order
27//! across a write is preserved by construction. A path-collision test would be a second, weaker
28//! statement of a rule the batching already enforces -- and one that would have to be extended
29//! every time a tool learned to touch a path it does not name.
30//!
31//! **Publishing stays serial and stays exactly-once.** `social_send` is the reachable tool on the
32//! wrong side of that line, and it is already excluded from road-failure retries for the same
33//! reason (see [`crate::tools::Tool::road_retryable`]): sending in somebody's name twice because
34//! of an infrastructure event is not something a budget may buy. Nothing that publishes is ever
35//! in a batch.
36//!
37//! **THE CONSENT SURFACE IS WHY THE WEB READS ARE NOT HERE**, and it is the trap this module was
38//! written around. `web_fetch`, `web_search` and `web_read` look like the ideal batch -- pure
39//! reads, hundreds of milliseconds each, no shared file. They are not, and the reason is
40//! [`crate::tools::ToolContext::wrap_untrusted`]: every one of them marks the turn as having read
41//! a stranger's words, and [`crate::tools::web_step`] consults that very flag to decide whether to
42//! ask the user. Run serially -- which is what happens today -- the first fetch is free, taints
43//! the turn, and the second fetch therefore asks. Run together, all of them read the flag before
44//! any of them sets it, every one comes out free, and a question the user is owed is never put.
45//! That is not slowness traded for a race; it is a consent gate silently skipped. So a web read is
46//! a batch of one until something has settled what a batch of them should ask, and the one ask per
47//! conversation that `web_step` now makes is not on its own an answer to it.
48//!
49//! **The batch is capped** at [`AT_ONCE`]. Nothing about correctness needs a cap; the turn's
50//! output budget does. Each result is charged as it lands ([`crate::tools::ToolRegistry::charge`]),
51//! so calls already in flight when the budget runs short are not cut by it, and the overshoot a
52//! turn can take is bounded by the size of a batch and by nothing else.
53
54use oxedyne_fe2o3_core::prelude::*;
55
56use std::future::Future;
57use std::ops::Range;
58use std::pin::Pin;
59use std::task::{Context, Poll};
60
61use crate::tools::Tool;
62
63// The most calls that ever run at once, whatever the model asked for. See the module note.
64pub const AT_ONCE: usize = 8;
65
66/// May this tool run beside another call in the same round?
67///
68/// The allow-list is small and stays small. A tool added to Daimond later is serial until
69/// somebody reads the four conditions in the module note and decides otherwise about it, which is
70/// the whole point of writing it as an allow-list: a new tool cannot become concurrent by
71/// accident, and a tool that grows a side effect does not quietly keep a permission it was granted
72/// when it had none.
73///
74/// `file_show` is the near miss worth naming. It reads a file, so it passes the first two
75/// conditions and looks like the others -- but what it does with what it reads is put it in the
76/// document panel, and two of them racing fight over one panel. Reading is not the test; changing
77/// nothing is.
78pub fn may_run_beside(name: &str) -> bool {
79 match Tool::from_name(name) {
80 Some(t) => matches!(t,
81 Tool::FileRead
82 | Tool::FileList
83 | Tool::FileSearch
84 | Tool::FileGlob
85 | Tool::SheetRead),
86 None => false,
87 }
88}
89
90/// Split a round's calls into consecutive batches, in the order the model gave them.
91///
92/// Every returned range is non-empty, they are contiguous, and they cover the whole slice -- so
93/// running them in order and recording each batch's results in order reproduces the model's own
94/// sequence exactly. A range of length one is today's behaviour and is what everything that is
95/// not in [`may_run_beside`] gets.
96///
97/// # Arguments
98/// * `names` - The wire name of each call, in the order the model emitted them.
99pub fn batches<S: AsRef<str>>(names: &[S]) -> Vec<Range<usize>> {
100 let mut out = Vec::new();
101 let mut i = 0;
102 while i < names.len() {
103 if !may_run_beside(names[i].as_ref()) {
104 out.push(i..i + 1);
105 i += 1;
106 continue;
107 }
108 let mut j = i + 1;
109 while j < names.len()
110 && j - i < AT_ONCE
111 && may_run_beside(names[j].as_ref())
112 {
113 j += 1;
114 }
115 out.push(i..j);
116 i = j;
117 }
118 out
119}
120
121/// Several futures of one type, run together, yielding their outputs in the order they were given.
122///
123/// Hand-written rather than taken from a futures crate, and generic rather than over
124/// `dyn Future`: Daimond depends on no async runtime beyond the one the browser already is, and
125/// the whole combinator is the `poll` below.
126///
127/// **It forwards the caller's own waker to every child.** There is no waker of its own to build,
128/// so there is no `RawWakerVTable` and no `unsafe`. A child that wakes wakes this task, which
129/// polls every child still running -- a few spurious polls for a batch bounded at
130/// [`AT_ONCE`], and no bookkeeping that could get the accounting wrong.
131///
132/// Each future is boxed, so the whole thing is `Unpin` and needs no pin projection. One
133/// allocation against a call that is about to wait on a disk or a network is not a cost worth
134/// avoiding.
135pub struct AllOf<F: Future> {
136 work: Vec<Option<Pin<Box<F>>>>,
137 done: Vec<Option<F::Output>>,
138}
139
140/// Run every future together. The outputs come back in the order the futures were given, not the
141/// order they finished.
142pub fn all_of<F: Future>(futs: Vec<F>) -> AllOf<F> {
143 let n = futs.len();
144 AllOf {
145 work: futs.into_iter().map(|f| Some(Box::pin(f))).collect(),
146 done: (0..n).map(|_| None).collect(),
147 }
148}
149
150// `F::Output: Unpin` because the finished outputs are held in a `Vec` beside the futures, and a
151// `Vec<T>` is `Unpin` only where `T` is. Every caller's output is an `Outcome`, which is.
152impl<F: Future> Future for AllOf<F>
153where
154 F::Output: Unpin,
155{
156 type Output = Vec<F::Output>;
157
158 fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
159 let me = self.get_mut();
160 let mut waiting = false;
161 for (i, slot) in me.work.iter_mut().enumerate() {
162 let ready = match slot {
163 Some(f) => match f.as_mut().poll(cx) {
164 Poll::Ready(v) => Some(v),
165 Poll::Pending => { waiting = true; None },
166 },
167 None => None,
168 };
169 if let Some(v) = ready {
170 me.done[i] = Some(v);
171 *slot = None;
172 }
173 }
174 if waiting {
175 return Poll::Pending;
176 }
177 // Every slot is filled: the loop above only leaves `waiting` false when no future is
178 // still running, and a future that finished wrote its output before its slot was cleared.
179 Poll::Ready(std::mem::take(&mut me.done).into_iter().flatten().collect())
180 }
181}
182
183
184// ┌───────────────────────────────────────────────────────────────┐
185// │ Tests │
186// └───────────────────────────────────────────────────────────────┘
187
188#[cfg(test)]
189mod tests {
190 use super::*;
191
192 #[test]
193 fn test_reads_batch_and_writes_do_not() {
194 let names = ["file_read", "file_read", "file_read"];
195 assert_eq!(vec![0..3], batches(&names), "three reads did not become one batch");
196
197 let names = ["file_write", "file_write"];
198 assert_eq!(vec![0..1, 1..2], batches(&names), "two writes were put in one batch");
199 }
200
201 #[test]
202 fn test_a_write_separates_the_reads_around_it() {
203 // The model's order across the write is the whole point: batch, write, batch.
204 let names = ["file_read", "file_read", "file_edit", "file_read"];
205 assert_eq!(vec![0..2, 2..3, 3..4], batches(&names),
206 "a write did not break the batch around it");
207 }
208
209 #[test]
210 fn test_two_writes_to_one_path_keep_the_model_s_order() {
211 // The property the batching has to preserve, stated as the batching sees it: each write
212 // is alone, and they are in the order given, so nothing can reorder them downstream.
213 let names = ["file_edit", "file_edit"];
214 let spans = batches(&names);
215 assert_eq!(2, spans.len(), "two writes did not stay two batches");
216 for (n, span) in spans.iter().enumerate() {
217 assert_eq!(n..n + 1, *span, "write {} was not alone and in place", n);
218 }
219 }
220
221 #[test]
222 fn test_nothing_that_reaches_the_network_or_asks_is_batched() {
223 // Named one by one rather than as a list, so a tool that changes side stands out in the
224 // diff. `web_fetch` is the one the module note is mostly about.
225 for name in ["web_fetch", "web_search", "web_read", "web_open", "web_click",
226 "social_send", "social_read", "ask", "run", "shell", "spawn_agent",
227 "file_show", "file_fetch", "typst_compile", "link_add"] {
228 assert!(!may_run_beside(name), "'{}' was allowed into a batch", name);
229 }
230 }
231
232 #[test]
233 fn test_an_unknown_tool_is_never_batched() {
234 assert!(!may_run_beside("no_such_tool"), "an unknown name was allowed into a batch");
235 let names = ["file_read", "no_such_tool", "file_read"];
236 assert_eq!(vec![0..1, 1..2, 2..3], batches(&names),
237 "an unknown name did not break the batch");
238 }
239
240 #[test]
241 fn test_a_batch_is_capped() {
242 let names = vec!["file_read"; AT_ONCE + 3];
243 assert_eq!(vec![0..AT_ONCE, AT_ONCE..AT_ONCE + 3], batches(&names),
244 "the cap did not bound the batch");
245 }
246
247 #[test]
248 fn test_the_batches_cover_every_call_exactly_once() {
249 let names = ["file_read", "file_edit", "file_list", "file_glob", "run", "sheet_read"];
250 let spans = batches(&names);
251 let mut at = 0;
252 for span in &spans {
253 assert_eq!(at, span.start, "the batches are not contiguous");
254 assert!(span.end > span.start, "an empty batch was produced");
255 at = span.end;
256 }
257 assert_eq!(names.len(), at, "the batches did not cover every call");
258 }
259
260 #[test]
261 fn test_nothing_is_batched_when_there_is_nothing() {
262 let names: [&str; 0] = [];
263 assert!(batches(&names).is_empty(), "an empty round produced a batch");
264 }
265
266 #[test]
267 fn test_all_of_returns_outputs_in_the_order_given_not_the_order_finished() {
268 // Deliberately finishing backwards: the last future is ready first.
269 let futs = vec![
270 ready_after(2, "a"),
271 ready_after(1, "b"),
272 ready_after(0, "c"),
273 ];
274 let got = block_on(all_of(futs));
275 assert_eq!(vec!["a", "b", "c"], got, "outputs came back in completion order");
276 }
277
278 /// THE WHOLE POINT, ON A CLOCK: three waits of one length take one wait's time together and
279 /// three waits' time in a queue.
280 ///
281 /// Both halves are measured in the one test rather than one of them being asserted from
282 /// memory, because the claim is a COMPARISON -- a machine under load makes any single figure
283 /// here meaningless, and the ratio survives what the absolute number does not.
284 #[tokio::test]
285 async fn test_waits_run_together_rather_than_one_after_another() {
286 let each = std::time::Duration::from_millis(200);
287
288 let began = std::time::Instant::now();
289 let got = all_of(vec![waited(each, 'a'), waited(each, 'b'), waited(each, 'c')]).await;
290 let batch = began.elapsed();
291 assert_eq!(vec!['a', 'b', 'c'], got, "the batch did not answer in the order given");
292
293 let began = std::time::Instant::now();
294 for c in ['a', 'b', 'c'] {
295 let _ = waited(each, c).await;
296 }
297 let queue = began.elapsed();
298
299 assert!(batch < each * 2,
300 "three {:?} waits took {:?} together, which is a queue and not a batch", each, batch);
301 assert!(queue > batch * 2,
302 "the queue took {:?} and the batch {:?}, which is not the difference this is for",
303 queue, batch);
304 }
305
306 /// A wait of a known length, so the test above measures a clock and not a poll count.
307 async fn waited(d: std::time::Duration, v: char) -> char {
308 tokio::time::sleep(d).await;
309 v
310 }
311
312 /// A future that answers after `n` further polls, so a test can order completions without a
313 /// clock.
314 fn ready_after(n: usize, v: &'static str) -> impl Future<Output = &'static str> {
315 let mut left = n;
316 std::future::poll_fn(move |cx| {
317 if left == 0 {
318 Poll::Ready(v)
319 } else {
320 left -= 1;
321 cx.waker().wake_by_ref();
322 Poll::Pending
323 }
324 })
325 }
326
327 /// Drive a future to completion on this thread, with the no-op waker the standard library
328 /// supplies -- so the test needs neither a runtime nor any `unsafe` of its own.
329 fn block_on<F: Future>(f: F) -> F::Output {
330 let mut f = Box::pin(f);
331 let waker = std::task::Waker::noop();
332 let mut cx = Context::from_waker(waker);
333 loop {
334 if let Poll::Ready(v) = f.as_mut().poll(&mut cx) {
335 return v;
336 }
337 }
338 }
339}