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 | |
| 54 | use oxedyne_fe2o3_core::prelude::*; |
| 55 | |
| 56 | use std::future::Future; |
| 57 | use std::ops::Range; |
| 58 | use std::pin::Pin; |
| 59 | use std::task::{Context, Poll}; |
| 60 | |
| 61 | use crate::tools::Tool; |
| 62 | |
| 63 | // The most calls that ever run at once, whatever the model asked for. See the module note. |
| 64 | pub 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. |
| 78 | pub 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. |
| 99 | pub 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. |
| 135 | pub 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. |
| 142 | pub 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. |
| 152 | impl<F: Future> Future for AllOf<F> |
| 153 | where |
| 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)] |
| 189 | mod 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 | } |