diff --git a/.repository-projection.json b/.repository-projection.json index 206f82278..dd9c18c2a 100644 --- a/.repository-projection.json +++ b/.repository-projection.json @@ -3,11 +3,11 @@ "projection": "deixic-code", "projectionSchemaVersion": 1, "sourceRepository": "dx-corp/mono", - "sourceSha": "fac55f7217fbc6e14c155ecfc00aed5c8c31adca", + "sourceSha": "e93942ba70f7585b80a55a24065ebbd31e57708d", "destinationRepository": "dx-corp/code", - "priorProjectedBase": "e9728f9fd83065904e209c56238a89291741bd74", + "priorProjectedBase": "a6d447c5960e7b738a3c3a27fa564f8dcae7f448", "definitionDigest": "82936441c776e3e8edb5d215a75007ec9714a233f489d460075d79d5ef5ba32f", "toolDigest": "c244d99199a7ae3eb8ff644a99462163c23b0bb6a83ef50af01efbdca0b81d04", - "contentDigest": "40c32f0bd494617b1e7297f7daf817ddd6861452f132690578f73452179f2900", + "contentDigest": "13746ab7cd6f508ce691afd179ce7d3f11b2bb8d901b5db0ba91153d55c25366", "publicationEligible": true } diff --git a/Cargo.lock b/Cargo.lock index cec8d0042..456ef8889 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -55,7 +55,9 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5a15f179cd60c4584b8a8c596927aadc462e27f2ca70c04e0071964a73ba7a75" dependencies = [ "cfg-if", + "getrandom 0.3.4", "once_cell", + "serde", "version_check", "zerocopy", ] @@ -1065,6 +1067,12 @@ dependencies = [ "piper", ] +[[package]] +name = "borrow-or-share" +version = "0.2.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dc0b364ead1874514c8c2855ab558056ebfeb775653e7ae45ff72f28f8f3166c" + [[package]] name = "bstr" version = "1.13.0" @@ -1091,6 +1099,12 @@ version = "1.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "64fa3c856b712db6612c019f14756e64e4bcea13337a6b33b696333a9eaa2d06" +[[package]] +name = "bytecount" +version = "0.6.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "175812e0be2bccb6abe50bb8d566126198344f707e304f45c648fd8f2cc0365e" + [[package]] name = "bytemuck" version = "1.25.2" @@ -2010,6 +2024,7 @@ version = "0.1.0" dependencies = [ "futures-util", "hex", + "jsonschema", "serde", "serde_json", "sha2 0.10.9", @@ -2229,6 +2244,15 @@ dependencies = [ "zeroize", ] +[[package]] +name = "email_address" +version = "0.2.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e079f19b08ca6239f47f8ba8509c11cf3ea30095831f7fed61441475edd8c449" +dependencies = [ + "serde", +] + [[package]] name = "embedded-io" version = "0.4.0" @@ -2403,6 +2427,17 @@ dependencies = [ "regex-syntax", ] +[[package]] +name = "fancy-regex" +version = "0.19.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d301f5bf187b3c295fce6468d3875037a0bccc5f6b151c63cac2f85babf21912" +dependencies = [ + "bit-set 0.8.0", + "regex-automata", + "regex-syntax", +] + [[package]] name = "fastrand" version = "2.5.0" @@ -2503,6 +2538,17 @@ dependencies = [ "zlib-rs", ] +[[package]] +name = "fluent-uri" +version = "0.4.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bc74ac4d8359ae70623506d512209619e5cf8f347124910440dbc221714b328e" +dependencies = [ + "borrow-or-share", + "ref-cast", + "serde", +] + [[package]] name = "fnv" version = "1.0.7" @@ -2530,6 +2576,16 @@ dependencies = [ "percent-encoding", ] +[[package]] +name = "fraction" +version = "0.15.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e076045bb43dac435333ed5f04caf35c7463631d0dae2deb2638d94dd0a5b872" +dependencies = [ + "lazy_static", + "num", +] + [[package]] name = "fs2" version = "0.4.3" @@ -2720,9 +2776,11 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "899def5c37c4fd7b2664648c28120ecec138e4d395b459e5ca34f9cce2dd77fd" dependencies = [ "cfg-if", + "js-sys", "libc", "r-efi 5.3.0", "wasip2", + "wasm-bindgen", ] [[package]] @@ -3441,6 +3499,58 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "jsonschema" +version = "0.49.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "59ec8a241beed129f06114aa68007e905ca350e7baeb6e17a7631bb7978d91b2" +dependencies = [ + "ahash", + "bytecount", + "data-encoding", + "email_address", + "fancy-regex 0.19.2", + "fraction", + "getrandom 0.3.4", + "idna", + "itoa", + "jsonschema-regex", + "jsonschema-value", + "num-cmp", + "num-traits", + "percent-encoding", + "referencing", + "regex", + "serde", + "serde_json", + "strum", + "unicode-general-category", + "uuid-simd", +] + +[[package]] +name = "jsonschema-regex" +version = "0.49.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "91994f45017ed5e66aa8e59b8415f4cb033a6380d7200387b7cf117595fbdf85" +dependencies = [ + "regex-syntax", +] + +[[package]] +name = "jsonschema-value" +version = "0.49.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7ec7637f83e510868ae6ed625f7ebfbbde4554ee8ce49854caa5126a8b9b9ecb" +dependencies = [ + "ahash", + "bytecount", + "fraction", + "num-cmp", + "num-traits", + "serde_json", +] + [[package]] name = "jsonwebtoken" version = "10.4.0" @@ -4314,6 +4424,12 @@ dependencies = [ "autocfg", ] +[[package]] +name = "micromap" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c2a86d3146ed3995b5913c414f6664344b9617457320782e64f0bb44afd49d74" + [[package]] name = "minimal-lexical" version = "0.2.1" @@ -4556,6 +4672,12 @@ dependencies = [ "zeroize", ] +[[package]] +name = "num-cmp" +version = "0.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "63335b2e2c34fae2fb0aa2cecfd9f0832a1e24b3b32ecec612c3426d46dc8aaa" + [[package]] name = "num-complex" version = "0.4.6" @@ -6050,6 +6172,43 @@ dependencies = [ "thiserror 2.0.20", ] +[[package]] +name = "ref-cast" +version = "1.0.27" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7e440fb4e4b4147295338efb76001ab9e4efc0e5839df2c47fc5ac2381d365c3" +dependencies = [ + "ref-cast-impl", +] + +[[package]] +name = "ref-cast-impl" +version = "1.0.27" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "92ecd8964f8453721699a1ed72037b0db49ce2f5a5138486ee89bed6f67cdf3a" +dependencies = [ + "proc-macro2", + "quote", + "syn 3.0.3", +] + +[[package]] +name = "referencing" +version = "0.49.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6efa2154ea6f5ce0fdecdd2a8d18f2fa1a39a8fbba91564f555a592e4dce8278" +dependencies = [ + "ahash", + "fluent-uri", + "getrandom 0.3.4", + "hashbrown 0.17.1", + "itoa", + "micromap", + "parking_lot", + "percent-encoding", + "serde_json", +] + [[package]] name = "regalloc2" version = "0.15.2" @@ -7755,6 +7914,12 @@ version = "0.3.18" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5c1cb5db39152898a79168971543b1cb5020dff7fe43c8dc468b0885f5e29df5" +[[package]] +name = "unicode-general-category" +version = "1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0b993bddc193ae5bd0d623b49ec06ac3e9312875fdae725a975c51db1cc1677f" + [[package]] name = "unicode-ident" version = "1.0.24" @@ -7904,6 +8069,16 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "uuid-simd" +version = "0.8.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "23b082222b4f6619906941c17eb2297fff4c2fb96cb60164170522942a200bd8" +dependencies = [ + "outref", + "vsimd", +] + [[package]] name = "valuable" version = "0.1.1" diff --git a/packages/dex-host-rs/src/log.rs b/packages/dex-host-rs/src/log.rs index b274eef1f..bdf7ea7c9 100644 --- a/packages/dex-host-rs/src/log.rs +++ b/packages/dex-host-rs/src/log.rs @@ -291,7 +291,7 @@ mod tests { attachments: Vec::new(), client_tools: Vec::new(), authorized_tools: Vec::new(), - approval_mode: dex_loop::ApprovalMode::Interactive, + approval_mode: dex_loop::ApprovalMode::Headless, } } diff --git a/packages/dex-host-rs/src/model.rs b/packages/dex-host-rs/src/model.rs index b724dcfec..560e241c1 100644 --- a/packages/dex-host-rs/src/model.rs +++ b/packages/dex-host-rs/src/model.rs @@ -286,7 +286,7 @@ mod tests { attachments: Vec::new(), client_tools: Vec::new(), authorized_tools: Vec::new(), - approval_mode: dex_loop::ApprovalMode::Interactive, + approval_mode: dex_loop::ApprovalMode::Headless, }, )], ); diff --git a/packages/dex-host-rs/src/tools.rs b/packages/dex-host-rs/src/tools.rs index 395c61a3c..0fb441224 100644 --- a/packages/dex-host-rs/src/tools.rs +++ b/packages/dex-host-rs/src/tools.rs @@ -2,7 +2,8 @@ //! //! Two tools, standing in for `local-host-rs`'s much larger registry //! (`packages/local-host-rs/src/tools/`): one read (`fs.read_file`) and one -//! mutation that needs approval (`fs.write_file`). A real cutover ports the +//! mutation in the `Approval` governance class (`fs.write_file`), which the +//! headless host allows without a prompt. A real cutover ports the //! rest of that registry the same way — each tool becomes a `ToolSpec` plus //! a `run` arm, and `NativeHost::requires_approval`'s per-call decision //! becomes this module's `policy` — not a rewrite of the tools themselves. @@ -11,8 +12,8 @@ use std::path::{Path, PathBuf}; use std::sync::Arc; use dex_loop::{ - ApprovalId, CancellationToken, Context, ExecutorKind, GovernanceClass, PrincipalId, - ProposedCall, ThreadId, ToolName, ToolResult, ToolSpec, Tools, Verdict, + CancellationToken, Context, ExecutorKind, GovernanceClass, PrincipalId, ProposedCall, ThreadId, + ToolName, ToolResult, ToolSpec, Tools, Verdict, }; use serde_json::Value; @@ -134,13 +135,6 @@ async fn write_file(root: &Path, args: &Value) -> ToolResult { } } -fn approval_summary(call: &ProposedCall) -> String { - match call.args.get("path").and_then(Value::as_str) { - Some(path) => format!("Write to {path}"), - None => "Write to a workspace file".to_owned(), - } -} - /// The `Tools` port over the workspace filesystem at `root`. #[derive(Clone)] pub struct LocalTools { @@ -175,13 +169,13 @@ impl Tools for LocalTools { } async fn policy(&self, _ctx: &Context, call: &ProposedCall) -> Verdict { - match self.spec(&call.tool).map(|spec| spec.governance) { - Some(GovernanceClass::Approval) => Verdict::NeedsApproval { - approval: ApprovalId::new(format!("approve-{}", call.id)), - summary: approval_summary(call), - }, - _ => Verdict::Allow, - } + // Maestro turns are headless: no human approves a tool call, so this + // host never returns `NeedsApproval`. An `Approval`-class tool (the + // mutation `fs.write_file`) is allowed outright; its + // `ToolStarted`/`ToolFinished` pair in the log is the audit record. + // A hard deny would be `Verdict::Deny`; no local tool has one today. + let _ = call; + Verdict::Allow } async fn run( @@ -275,7 +269,7 @@ mod tests { } #[tokio::test] - async fn write_file_requires_approval_and_read_file_does_not() { + async fn every_tool_is_allowed_because_turns_are_headless() { let tools = LocalTools::new("."); let read_verdict = tools .policy( @@ -290,7 +284,7 @@ mod tests { &call(WRITE_FILE, serde_json::json!({"path": "a", "content": "x"})), ) .await; - assert!(matches!(write_verdict, Verdict::NeedsApproval { .. })); + assert_eq!(write_verdict, Verdict::Allow); } /// `Tools::policy` here does not read `ctx`, so an empty rehydrated diff --git a/packages/dex-host-rs/src/turn.rs b/packages/dex-host-rs/src/turn.rs index 231ac8296..bb35cdf5a 100644 --- a/packages/dex-host-rs/src/turn.rs +++ b/packages/dex-host-rs/src/turn.rs @@ -5,8 +5,9 @@ //! //! This is Maestro's `MAESTRO_DEX_LOOP=1` cutover step: a local turn now has //! a consumer of this crate's ports, not just `tests/turn.rs`'s own -//! hand-driven approve/resume. Parked approvals are auto-approved and parked -//! questions are auto-answered with [`UNATTENDED_ANSWER`] -- the same +//! hand-driven approve/resume. Turns are headless (`ApprovalMode::Headless`, +//! `Verdict::Allow` for every local tool), so there is no approval step at +//! all, and a `Parked` exit is an error. Parked questions are auto-answered with [`UNATTENDED_ANSWER`] -- the same //! unattended-turn policy `print_mode.rs` (Maestro's existing non-interactive //! "auto-approves tools" entry point) and `cloud_cli.rs`'s attached REPL //! already use. This does not change trust posture: it matches the mode @@ -47,7 +48,7 @@ pub struct LocalTurnOutcome { /// Runs `request` to completion against a fresh `LocalLog`/`LocalTools`/ /// `LocalEffects` rooted at `state_root`/`workspace_root`, auto-approving -/// every parked approval and auto-answering every parked question. +/// auto-answering every parked question and failing on a parked approval. /// /// # Errors /// Returns an error if the log cannot be acquired (e.g. its lease is held by @@ -70,7 +71,7 @@ pub async fn run_local_turn( attachments: Vec::new(), client_tools: Vec::new(), authorized_tools: Vec::new(), - approval_mode: ApprovalMode::Interactive, + approval_mode: ApprovalMode::Headless, }]) .await .map_err(|_fenced| { @@ -103,8 +104,8 @@ pub async fn run_local_turn( type LocalEngine = dex_loop::Engine; -/// Runs `engine` until it reaches a terminal `Exit`, auto-resolving every -/// `Parked` (approval) and `Asked` (question) it hits along the way. +/// Runs `engine` until it reaches a terminal `Exit`, auto-answering every +/// `Asked` (question) and rejecting `Parked` (approval) it hits along the way. /// `AwaitingClientTool` has no local consumer yet -- client-side tools are /// not part of this crate's two-tool slice (`tools.rs`) -- so it is reported /// as an error instead of hanging forever. @@ -129,34 +130,15 @@ async fn drive_to_completion( consumer has no client session to answer it" ); } + // Headless turns never wait on a human: `LocalTools::policy` + // returns `Allow`, and `ApprovalMode::Headless` makes the engine + // grant any `NeedsApproval` itself. Reaching `Parked` is a bug, + // so fail loudly instead of parking or auto-deciding here. Exit::Parked(approval) => { - let entries = read_log(log).await?; - let pending = entries.iter().find_map(|(_, event)| match event { - Event::ApprovalRequested { - call, - approval: pending_approval, - args_digest, - .. - } if *pending_approval == approval => Some((call.clone(), args_digest.clone())), - _ => None, - }); - let Some((call, args_digest)) = pending else { - anyhow::bail!("no ApprovalRequested event found for {approval}"); - }; - log.append(&[Event::ApprovalDecided { - call, - approval, - args_digest, - approved: true, - principal: principal.clone(), - }]) - .await - .map_err(|_fenced| { - anyhow::anyhow!( - "the local dex-loop log's lease was superseded while auto-approving" - ) - })?; - ctx = rehydrate(thread.clone(), &read_log(log).await?); + anyhow::bail!( + "headless turn parked on approval {approval}; Maestro turns run with no \ + approval step, so this must not happen" + ); } // `LocalTools`' two-tool catalog (`tools.rs`) has no // `ExecutorKind::User` tool today, so the engine cannot actually @@ -258,7 +240,7 @@ mod tests { } #[tokio::test] - async fn drives_an_approval_and_a_question_to_completion_unattended() { + async fn drives_a_mutation_to_completion_with_no_approval() { let state_root = TempDir::new().expect("state tempdir"); let workspace = TempDir::new().expect("workspace tempdir"); std::fs::write(workspace.path().join("notes.txt"), "shopping list").expect("seed file"); @@ -268,7 +250,7 @@ mod tests { name: ToolName::new(READ_FILE), args: serde_json::json!({"path": "notes.txt"}), })], - // A mutation the consumer must auto-approve before it can run. + // A mutation: runs at once, with no approval step. vec![Ok(ModelChunk::ToolCall { name: ToolName::new(WRITE_FILE), args: serde_json::json!({"path": "out.txt", "content": "hello"}), @@ -296,6 +278,17 @@ mod tests { std::fs::read_to_string(workspace.path().join("out.txt")).expect("read written file"), "hello" ); + let log = LocalLog::acquire(state_root.path().join("log"), &thread()) + .await + .expect("reopen log"); + let events = read_log(&log).await.expect("read log"); + assert!( + !events.iter().any(|(_, event)| matches!( + event, + Event::ApprovalRequested { .. } | Event::ApprovalDecided { .. } + )), + "a headless Maestro turn must not request or decide an approval" + ); } #[tokio::test] @@ -322,7 +315,7 @@ mod tests { #[tokio::test] async fn a_denied_write_reports_failure_through_the_tool_result() { - // `run_local_turn` always auto-*approves*; this proves the turn + // `run_local_turn` never asks for approval; this proves the turn // still reaches `Exit::Done` (not stuck) when the mutation itself // fails for a reason unrelated to approval, e.g. a path outside the // workspace root. diff --git a/packages/dex-host-rs/tests/turn.rs b/packages/dex-host-rs/tests/turn.rs index ec005a6a7..57f4f243d 100644 --- a/packages/dex-host-rs/tests/turn.rs +++ b/packages/dex-host-rs/tests/turn.rs @@ -1,9 +1,10 @@ //! Runs one full `dex_loop::Engine` turn against `maestro_dex_host`'s local //! ports and a scripted fake model: a read tool that runs immediately, a -//! mutation that parks for approval, an approval, and a final answer. +//! mutation that runs straight through (headless: no approval step), and a +//! final answer. //! -//! This proves the kernel's park/approve/resume semantics run correctly -//! against a purely local host — no database, no HTTP — before any of +//! This proves a headless turn completes with no approval request against a +//! purely local host — no database, no HTTP — before any of //! Maestro's real turn logic is touched. See //! `docs/design/maestro-on-dex-loop.md`. @@ -11,9 +12,8 @@ use std::collections::VecDeque; use std::sync::{Arc, Mutex, PoisonError}; use dex_loop::{ - ApprovalId, Budget, CallId, CancellationToken, Context, Engine, Event, Exit, Lexicon, Log, - Model, ModelChunk, ModelError, Outcome, PrincipalId, ThreadId, ToolName, ToolSpec, TurnId, - args_digest, + Budget, CancellationToken, Context, Engine, Event, Exit, Lexicon, Log, Model, ModelChunk, + ModelError, Outcome, PrincipalId, ThreadId, ToolName, ToolSpec, TurnId, }; use futures_util::Stream; use futures_util::stream; @@ -73,7 +73,7 @@ fn alice() -> PrincipalId { } #[tokio::test] -async fn read_then_approved_write_then_done() { +async fn read_then_write_then_done_with_no_approval() { let state_root = TempDir::new().expect("state tempdir"); let workspace = TempDir::new().expect("workspace tempdir"); std::fs::write(workspace.path().join("notes.txt"), "shopping list").expect("seed file"); @@ -89,7 +89,7 @@ async fn read_then_approved_write_then_done() { attachments: Vec::new(), client_tools: Vec::new(), authorized_tools: Vec::new(), - approval_mode: dex_loop::ApprovalMode::Interactive, + approval_mode: dex_loop::ApprovalMode::Headless, }]) .await .expect("append user message"); @@ -123,52 +123,8 @@ async fn read_then_approved_write_then_done() { let entries = log.read_all().await.expect("read for rehydrate"); let mut ctx = dex_loop::rehydrate(thread(), &entries); - let exit = engine.run(&mut ctx, &cancel).await.expect("first run"); - // The read call ran with no approval; the write call parked on one. - let expected_write_call = CallId::new("t1-2-0"); - let expected_approval = ApprovalId::new(format!("approve-{expected_write_call}")); - assert_eq!(exit, Exit::Parked(expected_approval.clone())); - - // The mutation has not run yet: no file, no approval requested for the - // read call. - assert!(!workspace.path().join("out.txt").exists()); - let events: Vec = log - .read_all() - .await - .expect("read log") - .into_iter() - .map(|(_, event)| event) - .collect(); - assert!( - events.iter().any(|event| matches!( - event, - Event::ToolFinished { - outcome: Outcome::Succeeded, - .. - } - )), - "the read call must have finished before the write call parked" - ); - assert!( - !events - .iter() - .any(|event| matches!(event, Event::ApprovalRequested { call, .. } if call.as_str() != expected_write_call.as_str())), - "only the mutation should have requested approval" - ); - - log.append(&[Event::ApprovalDecided { - call: expected_write_call.clone(), - approval: expected_approval, - args_digest: args_digest(&write_args), - approved: true, - principal: alice(), - }]) - .await - .expect("append approval decision"); - - let entries = log.read_all().await.expect("read for resume"); - let mut ctx = dex_loop::rehydrate(thread(), &entries); - let exit = engine.run(&mut ctx, &cancel).await.expect("second run"); + // One run: the mutation goes straight through, nothing parks. + let exit = engine.run(&mut ctx, &cancel).await.expect("run"); assert_eq!(exit, Exit::Done); assert_eq!( @@ -182,5 +138,25 @@ async fn read_then_approved_write_then_done() { .into_iter() .map(|(_, event)| event) .collect(); + assert!( + !events.iter().any(|event| matches!( + event, + Event::ApprovalRequested { .. } | Event::ApprovalDecided { .. } + )), + "a headless turn must not request or decide any approval" + ); + let finished = events + .iter() + .filter(|event| { + matches!( + event, + Event::ToolFinished { + outcome: Outcome::Succeeded, + .. + } + ) + }) + .count(); + assert_eq!(finished, 2, "the read and the write both ran to completion"); assert!(matches!(events.last(), Some(Event::Final { text }) if text == "done")); } diff --git a/packages/tui-rs/src/dex_loop_local.rs b/packages/tui-rs/src/dex_loop_local.rs index 038f92875..0675dba4b 100644 --- a/packages/tui-rs/src/dex_loop_local.rs +++ b/packages/tui-rs/src/dex_loop_local.rs @@ -2,8 +2,8 @@ //! `MAESTRO_DEX_LOOP=1`. //! //! `products/maestro/packages/dex-host-rs` (`maestro_dex_host`) had no -//! consumer before this: `tests/turn.rs` proved the kernel's park/approve/ -//! resume semantics against a scripted model, but nothing in Maestro's own +//! consumer before this: `tests/turn.rs` proved the kernel's headless +//! (no approval step) semantics against a scripted model, but nothing in Maestro's own //! binaries ever ran a real local turn through it. This module is that //! consumer, wired into `print_mode.rs`'s non-interactive single-shot entry //! point -- itself already documented as "auto-approves tools" -- so the diff --git a/vendor/dex-loop/Cargo.toml b/vendor/dex-loop/Cargo.toml index 9bcba04cf..25db88718 100644 --- a/vendor/dex-loop/Cargo.toml +++ b/vendor/dex-loop/Cargo.toml @@ -9,6 +9,8 @@ rust-version = "1.95" [dependencies] futures-util = "0.3" hex = "0.4" +# Tool arguments are checked against the tool's schema before a call runs. +jsonschema = { version = "0.49", default-features = false } serde = { version = "1.0", features = ["derive"] } serde_json = "1.0" sha2 = "0.10" diff --git a/vendor/dex-loop/src/compaction.rs b/vendor/dex-loop/src/compaction.rs index 9a0447f9b..bf11c459c 100644 --- a/vendor/dex-loop/src/compaction.rs +++ b/vendor/dex-loop/src/compaction.rs @@ -70,6 +70,84 @@ impl Compactor for Threshold { } } +/// Where to cut history for a compaction, chosen by [`plan_cut`]. +#[derive(Clone, Copy, Debug, PartialEq)] +pub struct Cut<'a> { + /// The entries the summary replaces (a prefix of the history). + pub covered: &'a [Entry], + /// The cursor the `Compaction` event names: the last covered entry's. + pub covers_to: Cursor, + /// Whether the cut fell inside the current turn (after a closed step) + /// rather than at its start. + pub mid_turn: bool, +} + +/// Picks a cut for a history over `max_bytes`, or `None` when nothing should +/// be compacted now. +/// +/// - Never while a step's tool calls are unresolved (`open_step`): a cut +/// there would separate a tool call from its result. +/// - The kept part is only the user messages that follow the last assistant +/// or tool entry (the current turn's request and any steer), so the summary +/// is followed by what the model still has to answer. A cut at the start of +/// a turn keeps that turn's user message; a cut inside a turn (`mid_turn`) +/// falls right after a closed step and covers the turn's own steps so far. +/// - No assistant step is ever kept, so none is replayed under a summary it +/// was not produced beside. Claude requires the last assistant tool step to +/// carry its own unmodified thinking blocks (the request is rejected +/// otherwise), and a step produced before a compaction cannot replay them. +/// The next step starts a fresh assistant message after the summary. +/// - Never leaves a `Message::Tool` as the first kept entry and never splits +/// entries that share a cursor. +pub fn plan_cut(ctx: &Context, max_bytes: usize) -> Option> { + if ctx.open_step().is_some() { + return None; + } + let history = ctx.history(); + let total: usize = history.iter().map(|entry| entry.message.size()).sum(); + if total <= max_bytes { + return None; + } + let cut = history + .iter() + .rposition(|entry| !matches!(entry.message, Message::User { .. }))? + + 1; + if !valid_cut(history, cut) { + return None; + } + let turn_start = ctx.turn().and_then(|turn| { + history + .iter() + .position(|entry| matches!(&entry.message, Message::User { turn: t, .. } if t == turn)) + }); + Some(Cut { + covered: &history[..cut], + covers_to: history[cut - 1].cursor, + mid_turn: turn_start.is_some_and(|start| start < cut), + }) +} + +/// Whether history can be split before index `cut` (`cut == len` keeps +/// nothing): something precedes it that is more than a lone summary, the +/// last covered entry is a finished step (a tool result, or an assistant +/// message that made no calls), the first kept entry is not a tool result, +/// and the cut falls between two cursors. +fn valid_cut(history: &[Entry], cut: usize) -> bool { + if cut == 0 || cut > history.len() { + return false; + } + let only_summary = cut == 1 && matches!(history[0].message, Message::Summary { .. }); + let finished = match &history[cut - 1].message { + Message::Tool { .. } => true, + Message::Assistant { calls, .. } => calls.is_empty(), + Message::Summary { .. } | Message::User { .. } => false, + }; + let kept = history.get(cut); + let splits_cursor = kept.is_some_and(|kept| kept.cursor == history[cut - 1].cursor); + let orphans_result = kept.is_some_and(|kept| matches!(kept.message, Message::Tool { .. })); + !only_summary && finished && !splits_cursor && !orphans_result +} + /// The latest index `cut` with at least `keep_recent` entries after it where /// history can be split: the kept part must not start with a tool result /// (it belongs to the assistant message before it), and the cut must fall @@ -85,3 +163,219 @@ fn cut_point(history: &[Entry], keep_recent: usize) -> Option { !only_summary && !splits_cursor && !orphans_result }) } + +#[cfg(test)] +mod tests { + use super::*; + use crate::event::{ + ApprovalMode, CallId, Event, Outcome, Output, PrincipalId, ProposedCall, ThreadId, + ToolName, TurnId, + }; + use crate::rehydrate::rehydrate; + + const BIG: usize = 1_000; + + fn thread() -> ThreadId { + ThreadId { + org: "o".into(), + workspace: "w".into(), + thread: "t".into(), + } + } + + fn user(turn: &str) -> Event { + Event::UserMessage { + turn: TurnId::new(turn), + message_id: None, + principal: PrincipalId::new("alice"), + text: "x".repeat(BIG), + attachments: Vec::new(), + client_tools: Vec::new(), + authorized_tools: Vec::new(), + approval_mode: ApprovalMode::Interactive, + } + } + + fn step(n: u32, call: Option<&str>) -> Vec { + let calls = call + .map(|id| { + vec![ProposedCall::new( + CallId::new(id), + ToolName::new("tool"), + serde_json::json!({}), + PrincipalId::new("alice"), + )] + }) + .unwrap_or_default(); + vec![ + Event::StepStarted { + step: n, + control_through: Cursor::START, + }, + Event::ModelStepCompleted { + step: n, + text: "y".repeat(BIG), + calls, + reasoning: None, + served: None, + }, + ] + } + + fn result(call: &str) -> Event { + Event::ToolFinished { + call: CallId::new(call), + outcome: Outcome::Succeeded, + output: Output::Text("z".repeat(BIG)), + receipt: None, + } + } + + fn context(events: Vec) -> Context { + let numbered: Vec<_> = events + .into_iter() + .enumerate() + .map(|(index, event)| (Cursor(index as i64 + 1), event)) + .collect(); + rehydrate(thread(), &numbered) + } + + /// Turn 1 answered; turn 2 with `rounds` closed tool rounds and, when + /// `open`, a final step whose call has no result yet. + fn two_turns(rounds: u32, open: bool) -> Context { + let mut events = vec![user("t1")]; + events.extend(step(1, None)); + events.push(Event::Final { + text: "done".into(), + }); + events.push(user("t2")); + for round in 1..=rounds { + let call = format!("c{round}"); + events.extend(step(round, Some(&call))); + events.push(result(&call)); + } + if open { + events.extend(step(rounds + 1, Some("open"))); + } + context(events) + } + + #[test] + fn under_the_limit_nothing_is_cut() { + assert_eq!(plan_cut(&two_turns(1, false), 1_000_000), None); + } + + #[test] + fn cuts_at_the_start_of_the_current_turn() { + // Turn 2 has only its user message so far: the summary covers turn 1 + // and turn 2's request stays verbatim. + let ctx = two_turns(0, false); + let cut = plan_cut(&ctx, 2 * BIG).expect("cut"); + // Turn 1: user + assistant. Turn 2 starts at index 2. + assert_eq!(cut.covered.len(), 2); + assert!(!cut.mid_turn); + assert_eq!(cut.covers_to, ctx.history()[1].cursor); + assert!(matches!(ctx.history()[2].message, Message::User { .. })); + assert_eq!(ctx.history().len(), 3); + } + + #[test] + fn cuts_inside_the_turn_after_the_latest_closed_step() { + let ctx = two_turns(3, false); + let cut = plan_cut(&ctx, 4 * BIG).expect("cut"); + assert!(cut.mid_turn); + // Everything up to the last tool result is covered: the summary + // stands alone and no assistant step is replayed under it. + assert_eq!(cut.covered.len(), ctx.history().len()); + assert!(matches!( + cut.covered.last().expect("covered").message, + Message::Tool { .. } + )); + assert_eq!(cut.covers_to, ctx.history().last().expect("last").cursor); + } + + #[test] + fn never_keeps_an_assistant_step_after_the_cut() { + // Tool rounds inside a turn small enough to fit on its own: the + // turn-boundary cut is not taken because it would keep steps. + for rounds in 0..=4 { + let ctx = two_turns(rounds, false); + let Some(cut) = plan_cut(&ctx, BIG) else { + continue; + }; + let kept = &ctx.history()[cut.covered.len()..]; + assert!( + kept.iter() + .all(|entry| matches!(entry.message, Message::User { .. })), + "rounds={rounds}: only user messages are kept" + ); + } + } + + #[test] + fn refuses_while_a_tool_round_is_open() { + assert_eq!(plan_cut(&two_turns(3, true), BIG), None); + assert_eq!(plan_cut(&two_turns(0, true), BIG), None); + } + + #[test] + fn a_tool_result_never_leads_the_kept_part() { + let ctx = two_turns(3, false); + let history = ctx.history(); + for cut in 1..=history.len() { + if matches!( + history.get(cut).map(|e| &e.message), + Some(Message::Tool { .. }) + ) { + assert!(!valid_cut(history, cut), "cut {cut} would orphan a result"); + } + } + } + + #[test] + fn never_cuts_after_an_unfinished_assistant_message() { + let ctx = two_turns(2, true); + let history = ctx.history(); + let last = history.len() - 1; + assert!(matches!(history[last].message, Message::Assistant { .. })); + assert!(!valid_cut(history, last + 1)); + } + + #[test] + fn never_cuts_between_entries_of_one_cursor() { + let mut history = two_turns(2, false).history().to_vec(); + // Force the assistant after the first Tool group to share its cursor. + let tool_cursor = history[4].cursor; + history[5].cursor = tool_cursor; + assert!(!valid_cut(&history, 5)); + assert!(!valid_cut(&history, 0)); + assert!(!valid_cut(&history, history.len() + 1)); + } + + #[test] + fn a_lone_summary_is_not_worth_recompacting() { + let mut events = vec![user("t1")]; + events.extend(step(1, None)); + events.push(Event::Compaction { + covers_to_cursor: Cursor(3), + summary: "s".repeat(BIG), + }); + events.push(user("t2")); + let ctx = context(events); + // History: Summary + the new user message; the turn starts at 1. + assert_eq!(plan_cut(&ctx, BIG), None); + } + + #[test] + fn the_compaction_event_cursor_is_recorded() { + let mut events = vec![user("t1")]; + events.extend(step(1, None)); + events.push(Event::Compaction { + covers_to_cursor: Cursor(3), + summary: "s".into(), + }); + let ctx = context(events); + assert_eq!(ctx.last_compaction_cursor(), Some(Cursor(4))); + assert_eq!(context(vec![user("t1")]).last_compaction_cursor(), None); + } +} diff --git a/vendor/dex-loop/src/context.rs b/vendor/dex-loop/src/context.rs index 597918f79..7a5f4740a 100644 --- a/vendor/dex-loop/src/context.rs +++ b/vendor/dex-loop/src/context.rs @@ -7,8 +7,8 @@ use crate::event::{ ApprovalId, ApprovalMode, ArtifactRef, CallId, ClientToolSpec, Cursor, Event, MessageId, - Outcome, Output, PrincipalId, ProposedCall, ProviderReasoning, ThreadId, ToolName, ToolResult, - TurnId, Usage, + Outcome, Output, PrincipalId, ProposedCall, ProviderReasoning, ServedBy, ThreadId, ToolName, + ToolResult, TurnId, Usage, }; const NOT_RUN_NEW_TURN: &str = "not run: a new turn started first"; @@ -30,6 +30,8 @@ pub enum Message { /// The step's provider continuation state, as the `Model` port /// wrote it; `None` for steps logged without one. reasoning: Option, + /// The provider and model that served the step, when known. + served: Option, }, Tool { call: CallId, @@ -149,6 +151,16 @@ pub struct Context { /// Unknown call outcomes in this turn, derived from the durable log. /// Kept outside model history so compaction cannot permit a fresh retry. uncertain_calls: Vec, + /// Calls whose `ToolStarted` landed before their step's + /// `ModelStepCompleted`: reads the engine started while the model was + /// still streaming. They begin the step as `Started`, so a warm engine + /// adopts their results and a rehydrated one runs them again, the same + /// as any other read that started before a restart. + pre_started: Vec, + /// Cursor of the latest `Compaction` event itself (not the cursor it + /// covers to). Assistant entries at or before it were produced before + /// the summary existed. + last_compaction: Option, } #[derive(Clone, Debug, PartialEq)] @@ -189,6 +201,8 @@ impl Context { authorized_principal: None, uncertain_calls: Vec::new(), pending_turns: Vec::new(), + pre_started: Vec::new(), + last_compaction: None, } } @@ -241,6 +255,15 @@ impl Context { &self.history } + /// The cursor of the latest `Compaction` event, if history was compacted. + /// A history entry with `cursor <= last_compaction_cursor()` was produced + /// before that compaction; the model never saw it under the summary that + /// now precedes it, so provider continuation state bound to the old + /// prefix (signed thinking) cannot be replayed for it. + pub fn last_compaction_cursor(&self) -> Option { + self.last_compaction + } + /// Model calls started in the current turn. pub fn step(&self) -> u32 { self.step @@ -283,13 +306,6 @@ impl Context { self.control = self.control.max(floor); } - /// Puts the control cursor back after the engine appended a control event - /// of its own, so a control event from another writer that landed just - /// before it is still read. - pub(crate) fn rewind_control(&mut self, to: Cursor) { - self.control = to; - } - pub(crate) fn status(&self) -> Status { self.status } @@ -316,6 +332,12 @@ impl Context { }) } + /// Reads a cut attempt started ahead of its step's commit and never + /// finished, in log order. + pub(crate) fn pre_started_calls(&self) -> &[CallId] { + &self.pre_started + } + /// The step of a model attempt that started but never completed. pub(crate) fn open_attempt(&self) -> Option { self.attempt @@ -439,12 +461,13 @@ impl Context { self.flush_steers(cursor, *control_through); } // Text reaches history through `ModelStepCompleted`. - Event::TextDelta { .. } | Event::ToolProgress { .. } => {} + Event::TextDelta { .. } | Event::ThinkingDelta { .. } | Event::ToolProgress { .. } => {} Event::Usage(usage) => self.usage += *usage, Event::ModelStepCompleted { text, calls, reasoning, + served, .. } => { self.attempt = None; @@ -454,21 +477,39 @@ impl Context { text: text.clone(), calls: calls.clone(), reasoning: reasoning.clone(), + served: served.clone(), }, ); if !calls.is_empty() { + let pre_started = std::mem::take(&mut self.pre_started); self.open_step = Some(OpenStep { calls: calls.clone(), - states: vec![CallState::Todo; calls.len()], + states: calls + .iter() + .map(|call| { + if pre_started.contains(&call.id) { + CallState::Started + } else { + CallState::Todo + } + }) + .collect(), }); } + self.pre_started.clear(); + } + Event::ModelAttemptAbandoned { .. } => { + self.attempt = None; + self.pre_started.clear(); } - Event::ModelAttemptAbandoned { .. } => self.attempt = None, Event::ToolStarted { call, .. } => { - if let Some(state) = self.state_mut(call) - && !matches!(state, CallState::Done(_)) - { - *state = CallState::Started; + if let Some(state) = self.state_mut(call) { + if !matches!(state, CallState::Done(_)) { + *state = CallState::Started; + } + } else if self.open_step.is_none() && !self.pre_started.contains(call) { + // Started ahead of its step's commit point. + self.pre_started.push(call.clone()); } } Event::ToolsExposed { tools, .. } => { @@ -485,6 +526,8 @@ impl Context { receipt, } => { self.uncertain_calls.retain(|prior| &prior.id != call); + // A pre-committed read finished by an abandoned attempt. + self.pre_started.retain(|started| started != call); if *outcome == Outcome::Unknown && let Some(proposal) = self .open_step @@ -547,6 +590,7 @@ impl Context { covers_to_cursor, summary, } => { + self.last_compaction = Some(cursor); self.history .retain(|entry| entry.cursor > *covers_to_cursor); self.history.insert( @@ -606,6 +650,7 @@ impl Context { self.interrupt_requested = false; self.exposed.clear(); self.uncertain_calls.clear(); + self.pre_started.clear(); self.client_tools = client_tools; self.authorized_tools = authorized_tools; self.approval_mode = approval_mode; diff --git a/vendor/dex-loop/src/engine.rs b/vendor/dex-loop/src/engine.rs index b0eaffa27..6fb68322e 100644 --- a/vendor/dex-loop/src/engine.rs +++ b/vendor/dex-loop/src/engine.rs @@ -6,8 +6,18 @@ //! `ToolStarted` → effect → `ToolFinished`. The turn ends when a step //! proposes no calls. No call ever parks for a human: a `NeedsApproval` //! verdict is granted at once and recorded as `AutoApproved`. +//! +//! One exception to "after the commit point": a read-only call whose +//! arguments validate and whose policy allows it starts while the model is +//! still streaming (`ToolStarted` before `ModelStepCompleted`, the same +//! way Codex starts a call when its output item closes). A read has no +//! effect to wait for, so an attempt that never commits simply finishes +//! those calls as not run; a committed step adopts their results in place +//! of running them again. -use std::pin::pin; +use std::collections::{HashMap, HashSet}; +use std::future::Future; +use std::pin::{Pin, pin}; use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; use futures_util::StreamExt; @@ -18,8 +28,8 @@ use crate::budget::{Budget, BudgetAxis}; use crate::compaction::{Compactor, NoCompaction}; use crate::context::{CallState, Context, Decision, Status}; use crate::event::{ - AUTO_APPROVER, ApprovalId, CallId, ErrorCode, Event, Outcome, PrincipalId, ProposedCall, - ToolName, ToolResult, TurnId, + AUTO_APPROVER, ApprovalId, CallId, Cursor, ErrorCode, Event, Outcome, PrincipalId, + ProposedCall, ToolName, ToolResult, TurnId, }; use crate::ports::{ Claim, Effects, ExecutorKind, Fenced, GovernanceClass, Log, Model, ModelChunk, ModelError, @@ -47,6 +57,10 @@ const CLIENT_TIMED_OUT_READ: &str = "not completed: the client did not report a result in time; it is safe to try again"; const CLIENT_TIMED_OUT_MUTATION: &str = "outcome unknown: the client did not report a result in time; check before trying again"; +/// A read started while the model streamed, whose attempt never committed +/// (the stream failed or was cut): it is finished so its `ToolStarted` does +/// not dangle; the next attempt proposes and runs it afresh. +const NOT_RUN_ATTEMPT_ABANDONED: &str = "not run: the model attempt was abandoned"; /// How long a call waits for `Event::ClientToolResult` before the engine /// gives up on it. Matches `dex_tools::Timeouts::default().client`; a host @@ -105,8 +119,8 @@ struct ClientCall<'a> { decision: Option, } -/// One result of racing the model stream against `cancel` and the wall -/// budget in `model_step`. +/// One result of racing the model stream against `cancel`, the wall budget +/// and the reads started during the stream in `model_step`. enum StreamStep { Chunk(Result), /// The stream ended on its own (a truncated or otherwise finite stream). @@ -115,6 +129,39 @@ enum StreamStep { /// `budget.wall` elapsed while waiting for the next chunk: the stream /// itself never errored or ended, so nothing else would have caught this. WallExceeded, + /// A read started during the stream returned. + Prefetched(usize, ToolResult), +} + +/// One read-only call running ahead of its step's commit point, by index +/// in the step's calls. +type PrefetchFuture<'e> = Pin + Send + 'e>>; + +/// Reads started while the model streamed (see the module doc). Lives in +/// `Engine::run` for one step: filled by `model_step`, drained by +/// `dispatch`, and dropped whole when the attempt does not commit. +struct Prefetch<'e> { + /// Still running. Each carries its own deadline, the one a wave would + /// have given it. + pending: FuturesUnordered>, + /// Returned before the step committed. + done: HashMap, + /// Indices whose `ToolStarted` is already on the log. + started: HashSet, + /// The `ToolStarted` rows appended mid-stream. The model stream borrows + /// `ctx` until it is dropped, so they are observed then, in log order. + unobserved: Vec<(Cursor, Event)>, +} + +impl Prefetch<'_> { + fn new() -> Self { + Self { + pending: FuturesUnordered::new(), + done: HashMap::new(), + started: HashSet::new(), + unobserved: Vec::new(), + } + } } /// The agent loop. Holds the host's ports and no state of its own, so any @@ -208,8 +255,13 @@ where /// /// The host appends `Interrupt` and then cancels `cancel`. `Err(Fenced)` /// means a write was refused and nothing more was appended. - pub async fn run(&self, ctx: &mut Context, cancel: &CancellationToken) -> Result { + pub async fn run<'e>( + &'e self, + ctx: &mut Context, + cancel: &'e CancellationToken, + ) -> Result { let started = Instant::now(); + let mut prefetch = Prefetch::new(); loop { self.read_control(ctx).await?; match ctx.status() { @@ -219,15 +271,30 @@ where Status::Running => {} } if let Some(step) = ctx.open_attempt() { - // A crash mid-stream: the attempt's text never committed. - self.emit(ctx, vec![Event::ModelAttemptAbandoned { step }]) - .await?; + // A crash mid-stream: the attempt's text never committed, and + // neither did the reads it had started; close those so no + // `ToolStarted` stays open. + let mut events: Vec = ctx + .pre_started_calls() + .iter() + .map(|call| { + let result = ToolResult::error(NOT_RUN_ATTEMPT_ABANDONED); + Event::ToolFinished { + call: call.clone(), + outcome: result.outcome, + output: result.output, + receipt: result.receipt, + } + }) + .collect(); + events.push(Event::ModelAttemptAbandoned { step }); + self.emit(ctx, events).await?; } if ctx.interrupt_requested() || cancel.is_cancelled() { return self.interrupt(ctx).await; } if ctx.open_step().is_some() { - if let Some(exit) = self.dispatch(ctx, cancel, started).await? { + if let Some(exit) = self.dispatch(ctx, cancel, started, &mut prefetch).await? { return Ok(exit); } continue; @@ -257,7 +324,9 @@ where ) .await?; } - if let Some(exit) = self.model_step(ctx, cancel, started).await? { + // A step never inherits another step's reads. + prefetch = Prefetch::new(); + if let Some(exit) = self.model_step(ctx, cancel, started, &mut prefetch).await? { return Ok(exit); } } @@ -275,11 +344,12 @@ where /// (including a stream that never completes: `budget.wall` bounds the /// whole `Engine::run` call, streaming included, not just the time /// between steps). - async fn model_step( - &self, + async fn model_step<'e>( + &'e self, ctx: &mut Context, - cancel: &CancellationToken, + cancel: &'e CancellationToken, started: Instant, + prefetch: &mut Prefetch<'e>, ) -> Result, Fenced> { let step = ctx.step().saturating_add(1); self.emit( @@ -298,6 +368,7 @@ where }; let mut filter = self.sanitizer.filter(); + let mut thinking = Thinking::new(self.sanitizer.filter()); let mut text = String::new(); let mut calls = Vec::new(); let mut failure = None; @@ -313,9 +384,13 @@ where // (the model sends it only after a clean terminal). A second chunk // replaces the first: one step has one. let mut reasoning = None; + // The route and model that served this attempt, kept for every + // `ModelStepCompleted` below, including a cut-off or cancelled one. + let mut served = None; { // The answer-only call offers nothing, not even `tools.search`. - let owned = if self.budget.answer_only(step.saturating_sub(1)) { + let answer_only = self.budget.answer_only(step.saturating_sub(1)); + let owned = if answer_only { Vec::new() } else { self.offered(ctx) @@ -329,6 +404,10 @@ where // not only against `cancel`. loop { let remaining = self.budget.wall.saturating_sub(started.elapsed()); + // `FuturesUnordered::next` on an empty set resolves at once + // with `None`; the guard keeps it out of the race until a + // read is actually running. + let has_pending = !prefetch.pending.is_empty(); let outcome = tokio::select! { biased; () = cancel.cancelled() => StreamStep::Cancelled, @@ -337,9 +416,25 @@ where Some(chunk) => StreamStep::Chunk(chunk), None => StreamStep::Ended, }, + done = prefetch.pending.next(), if has_pending => match done { + Some((index, result)) => StreamStep::Prefetched(index, result), + None => continue, + }, }; match outcome { + StreamStep::Chunk(Ok(ModelChunk::Thinking(delta))) => { + if let Some(summary) = thinking.push(&delta, !text.is_empty()) { + self.log + .append(&[Event::ThinkingDelta { text: summary }]) + .await?; + } + } StreamStep::Chunk(Ok(ModelChunk::Text(delta))) => { + if let Some(summary) = thinking.flush() { + self.log + .append(&[Event::ThinkingDelta { text: summary }]) + .await?; + } let safe = filter.push(&delta); if !safe.is_empty() { text.push_str(&safe); @@ -347,8 +442,19 @@ where } } StreamStep::Chunk(Ok(ModelChunk::ToolCall { name, args })) => { - let id = call_id(&turn, step, calls.len()); - calls.push(ProposedCall::new(id, name, args, principal.clone())); + let index = calls.len(); + let id = call_id(&turn, step, index); + let call = ProposedCall::new(id, name, args, principal.clone()); + // The answer-only call offers no tools, so nothing may + // start on it. + if !answer_only { + self.prefetch(ctx, &call, index, cancel, started, prefetch) + .await?; + } + calls.push(call); + } + StreamStep::Prefetched(index, result) => { + prefetch.done.insert(index, result); } StreamStep::Chunk(Ok(ModelChunk::Usage(usage))) => { pending_usage.push(Event::Usage(usage)); @@ -356,6 +462,9 @@ where StreamStep::Chunk(Ok(ModelChunk::Reasoning(state))) => { reasoning = Some(state); } + StreamStep::Chunk(Ok(ModelChunk::Served(by))) => { + served = Some(by); + } StreamStep::Chunk(Err(error)) => { failure = Some(error.message); break; @@ -368,6 +477,11 @@ where } } if failure.is_none() && !wall_exceeded { + if let Some(summary) = thinking.flush() { + self.log + .append(&[Event::ThinkingDelta { text: summary }]) + .await?; + } let tail = filter.finish(); if !tail.is_empty() { text.push_str(&tail); @@ -376,6 +490,15 @@ where } } + for (cursor, event) in std::mem::take(&mut prefetch.unobserved) { + ctx.observe(cursor, &event); + } + // Every path below that does not commit `calls` first closes the + // reads that already started, so no `ToolStarted` dangles. + let committing = failure.is_none() && !wall_exceeded && !cancel.is_cancelled(); + if !committing { + pending_usage.splice(0..0, Self::abandon_prefetch(&calls, prefetch)); + } if wall_exceeded { let message = self.budget_message(ctx, BudgetAxis::Wall); let mut events = pending_usage; @@ -406,6 +529,7 @@ where text: text.clone(), calls: Vec::new(), reasoning: None, + served: served.clone(), }); events.push(Event::Final { text }); self.emit(ctx, events).await?; @@ -430,13 +554,15 @@ where text, calls: Vec::new(), reasoning: None, + served: served.clone(), }); self.emit(ctx, events).await?; return self.interrupt(ctx).await.map(Some); } if !calls.is_empty() && self.budget.answer_only(step.saturating_sub(1)) { // Asked for a tool on the call that offered none. Nothing can run - // it, so the turn ends here instead of looping. + // it, so the turn ends here instead of looping. (Nothing was + // prefetched: the answer-only call offers no tools.) let message = self.budget_message(ctx, BudgetAxis::Steps); let mut events = pending_usage; events.push(Event::ModelAttemptAbandoned { step }); @@ -457,6 +583,7 @@ where text: text.clone(), calls: Vec::new(), reasoning, + served, }); if !continues { events.push(Event::Final { text }); @@ -470,19 +597,93 @@ where text, calls, reasoning, + served, }); self.emit(ctx, events).await?; Ok(None) } + /// Starts `call` now, ahead of its step's commit, when it is a read the + /// model was offered, its arguments fit the tool's schema, its policy + /// allows it, and it runs on the host (not a person, not a client + /// session). Anything else waits for `dispatch`, which re-derives the + /// same verdicts and finishes an invalid or denied call there. + async fn prefetch<'e>( + &'e self, + ctx: &Context, + call: &ProposedCall, + index: usize, + cancel: &'e CancellationToken, + run_started: Instant, + prefetch: &mut Prefetch<'e>, + ) -> Result<(), Fenced> { + // Only an unbroken run of reads from the first call on starts early. + // `dispatch` keeps the model's order around effects, so a read that + // follows a mutation (or any call held back) must not observe the + // world before that call has run. + if prefetch.started.len() != index || call.tool.as_str() == TOOLS_SEARCH { + return Ok(()); + } + let Some(spec) = self.offered_spec(ctx, &call.tool) else { + return Ok(()); + }; + let eligible = spec.read_only + && !matches!(spec.executor, ExecutorKind::User | ExecutorKind::Client) + && validate_args(&spec, &call.args).is_ok() + && !ctx.has_uncertain_call(call); + if !eligible || self.tools.policy(ctx, call).await != Verdict::Allow { + return Ok(()); + } + // Appended without `ctx.observe`: the stream still borrows `ctx`. + let event = started(call, &spec); + let cursors = self.log.append(std::slice::from_ref(&event)).await?; + let [cursor] = cursors[..] else { + return Err(Fenced::new(format!( + "log returned {} cursors for 1 event", + cursors.len() + ))); + }; + prefetch.unobserved.push((cursor, event)); + prefetch.started.insert(index); + let deadline = self.call_deadline(run_started); + let thread = ctx.thread().clone(); + let call = call.clone(); + prefetch.pending.push(Box::pin(async move { + let run = self.tools.run(&thread, &call, cancel); + let result = match tokio::time::timeout(deadline, run).await { + Ok(result) => result, + Err(_elapsed) => ToolResult::error(DEADLINE_READ), + }; + (index, result) + })); + Ok(()) + } + + /// The `ToolFinished` rows for reads that started under an attempt + /// that will not commit, in call order; the reads themselves are + /// dropped (a read has no effect to wait for). Leaves `prefetch` empty. + fn abandon_prefetch(calls: &[ProposedCall], prefetch: &mut Prefetch<'_>) -> Vec { + let events = calls + .iter() + .enumerate() + .filter(|(index, _)| prefetch.started.contains(index)) + .map(|(_, call)| finished(call, ToolResult::error(NOT_RUN_ATTEMPT_ABANDONED))) + .collect(); + *prefetch = Prefetch::new(); + events + } + /// Dispatches the open step's calls in the model's order. Allowed /// read-only calls collect into a wave that runs in parallel; anything /// else runs the pending wave first, so effects keep the model's order. - async fn dispatch( - &self, + /// Reads that already started during the stream (`prefetch`) are + /// adopted by the wave instead of running again. + async fn dispatch<'e>( + &'e self, ctx: &mut Context, - cancel: &CancellationToken, + cancel: &'e CancellationToken, run_started: Instant, + prefetch: &mut Prefetch<'e>, ) -> Result, Fenced> { let Some(step) = ctx.open_step() else { return Ok(None); @@ -502,7 +703,7 @@ where None | Some(CallState::Done(_)) => continue, Some(CallState::Asked { answer: None }) => { if self - .flush(ctx, &calls, &mut wave, cancel, run_started) + .flush(ctx, &calls, &mut wave, cancel, run_started, prefetch) .await? { break; @@ -524,7 +725,7 @@ where continue; } if self - .flush(ctx, &calls, &mut wave, cancel, run_started) + .flush(ctx, &calls, &mut wave, cancel, run_started, prefetch) .await? { break; @@ -547,7 +748,7 @@ where // under the same approval id, so the thread resumes // instead of waiting forever. if self - .flush(ctx, &calls, &mut wave, cancel, run_started) + .flush(ctx, &calls, &mut wave, cancel, run_started, prefetch) .await? { break; @@ -573,12 +774,12 @@ where Some(spec) if spec.read_only => wave.push(index), Some(_) => { if self - .flush(ctx, &calls, &mut wave, cancel, run_started) + .flush(ctx, &calls, &mut wave, cancel, run_started, prefetch) .await? { break; } - self.run_mutation(ctx, call, run_started).await?; + self.run_mutation(ctx, call, cancel, run_started).await?; } // The tool is no longer offered (deploy, grant // revoke), but this call already started: it may @@ -601,6 +802,13 @@ where self.finish(ctx, call, unknown_tool(&call.tool)).await?; continue; }; + // Arguments that do not fit the tool's schema never reach an + // executor: the model sees why and can call again. (A call the + // stream already started passed this check before it ran.) + if let Err(reason) = validate_args(&spec, &call.args) { + self.finish(ctx, call, ToolResult::error(reason)).await?; + continue; + } // Client-executor tools are not in `self.tools`'s catalog and // carry their own governance (set by the host's allowlist when // it resolved the client's declaration), so they skip @@ -615,6 +823,7 @@ where &mut wave, cancel, run_started, + prefetch, ClientCall { proposal: call, decision, @@ -657,7 +866,7 @@ where // effect; the pending wave runs first so effects keep order. if let (None, Verdict::NeedsApproval { approval, summary }) = (decision, verdict) { if self - .flush(ctx, &calls, &mut wave, cancel, run_started) + .flush(ctx, &calls, &mut wave, cancel, run_started, prefetch) .await? { break; @@ -667,7 +876,7 @@ where if spec.executor == ExecutorKind::User { if self - .flush(ctx, &calls, &mut wave, cancel, run_started) + .flush(ctx, &calls, &mut wave, cancel, run_started, prefetch) .await? { break; @@ -691,15 +900,15 @@ where wave.push(index); } else { if self - .flush(ctx, &calls, &mut wave, cancel, run_started) + .flush(ctx, &calls, &mut wave, cancel, run_started, prefetch) .await? { break; } - self.run_mutation(ctx, call, run_started).await?; + self.run_mutation(ctx, call, cancel, run_started).await?; } } - self.run_wave(ctx, &calls, wave, cancel, run_started) + self.run_wave(ctx, &calls, wave, cancel, run_started, prefetch) .await?; if cancel.is_cancelled() { return self.interrupt(ctx).await.map(Some); @@ -726,13 +935,15 @@ where /// One call to a `Client`-executor tool: approval (if the host's /// allowlist marked it a mutation), then `ClientToolRequested`, mirroring /// how the main `dispatch` loop handles `NeedsApproval` and `User`. - async fn dispatch_client_tool( - &self, + #[allow(clippy::too_many_arguments)] + async fn dispatch_client_tool<'e>( + &'e self, ctx: &mut Context, calls: &[ProposedCall], wave: &mut Vec, - cancel: &CancellationToken, + cancel: &'e CancellationToken, run_started: Instant, + prefetch: &mut Prefetch<'e>, client_call: ClientCall<'_>, ) -> Result { let ClientCall { @@ -759,14 +970,20 @@ where return Ok(ClientToolOutcome::Continue); } if decision.is_none() && spec.governance == GovernanceClass::Approval { - if self.flush(ctx, calls, wave, cancel, run_started).await? { + if self + .flush(ctx, calls, wave, cancel, run_started, prefetch) + .await? + { return Ok(ClientToolOutcome::Break); } let approval = ApprovalId::new(format!("client-{}", call.id)); let summary = format!("Run {} in your browser", spec.label); self.auto_approve(ctx, call, approval, summary).await?; } - if self.flush(ctx, calls, wave, cancel, run_started).await? { + if self + .flush(ctx, calls, wave, cancel, run_started, prefetch) + .await? + { return Ok(ClientToolOutcome::Break); } let deadline_ms = now_ms().saturating_add(self.client_timeout_millis()); @@ -791,35 +1008,47 @@ where /// Runs the pending wave before a call that must not overlap it. Returns /// true when an interrupt arrived meanwhile: the next effect must not /// start. - async fn flush( - &self, + async fn flush<'e>( + &'e self, ctx: &mut Context, calls: &[ProposedCall], wave: &mut Vec, - cancel: &CancellationToken, + cancel: &'e CancellationToken, run_started: Instant, + prefetch: &mut Prefetch<'e>, ) -> Result { - self.run_wave(ctx, calls, std::mem::take(wave), cancel, run_started) - .await?; + self.run_wave( + ctx, + calls, + std::mem::take(wave), + cancel, + run_started, + prefetch, + ) + .await?; Ok(cancel.is_cancelled()) } /// Runs read-only calls concurrently. Each `ToolFinished` is appended as /// its call returns; history receives the results in call order when the - /// step closes. - async fn run_wave( - &self, + /// step closes. A call that already started during the stream keeps its + /// `ToolStarted` and its run: a finished result is adopted at once and a + /// pending one joins the wave. + async fn run_wave<'e>( + &'e self, ctx: &mut Context, calls: &[ProposedCall], wave: Vec, - cancel: &CancellationToken, + cancel: &'e CancellationToken, run_started: Instant, + prefetch: &mut Prefetch<'e>, ) -> Result<(), Fenced> { if wave.is_empty() || cancel.is_cancelled() { return Ok(()); } let starts: Vec = wave .iter() + .filter(|index| !prefetch.started.contains(index)) .filter_map(|&index| { let call = &calls[index]; let spec = self.offered_spec(ctx, &call.tool)?; @@ -828,29 +1057,48 @@ where .collect(); self.emit(ctx, starts).await?; let thread = ctx.thread().clone(); - let thread = &thread; // One deadline for the wave: its reads run concurrently, so each // gets the full time. A read that overruns is dropped and finished // `Failed`; a read has no effect to wait for, so retrying is safe. let deadline = self.call_deadline(run_started); - let mut running: FuturesUnordered<_> = wave - .iter() - .map(|&index| { - let call = &calls[index]; - async move { - let run = self.tools.run(thread, call, cancel); - let result = match tokio::time::timeout(deadline, run).await { - Ok(result) => result, - Err(_elapsed) => ToolResult::error(DEADLINE_READ), - }; - (call, result) - } - }) - .collect(); + let mut running: FuturesUnordered> = FuturesUnordered::new(); + for &index in &wave { + if let Some(result) = prefetch.done.remove(&index) { + self.finish(ctx, &calls[index], result).await?; + continue; + } + if prefetch.started.contains(&index) { + // Still running from the stream; it arrives through + // `prefetch.pending` below. + continue; + } + let call = calls[index].clone(); + let thread = thread.clone(); + running.push(Box::pin(async move { + let run = self.tools.run(&thread, &call, cancel); + let result = match tokio::time::timeout(deadline, run).await { + Ok(result) => result, + Err(_elapsed) => ToolResult::error(DEADLINE_READ), + }; + (index, result) + })); + } + for future in std::mem::take(&mut prefetch.pending) { + running.push(future); + } // On `Fenced` the remaining reads are dropped: a stale owner must not // append, and the new owner runs them again. - while let Some((call, result)) = running.next().await { - self.finish(ctx, call, result).await?; + while let Some((index, result)) = running.next().await { + let still_open = ctx + .open_step() + .and_then(|step| step.states.get(index)) + .is_some_and(|state| !matches!(state, CallState::Done(_))); + // A prefetched read whose call `dispatch` already finished + // (policy denied it on re-check) has nothing left to report. + if !still_open || !(wave.contains(&index) || prefetch.started.contains(&index)) { + continue; + } + self.finish(ctx, &calls[index], result).await?; } Ok(()) } @@ -859,11 +1107,14 @@ where /// its recorded outcome is adopted instead, with `Running` settled to /// `Unknown` first — nothing ever revisits a `Running` report, so /// showing it as final would leave the model unable to tell whether to - /// retry. Interrupt does not cancel a mutation that has started. + /// retry. Interrupt reaches a running mutation through `cancel`, but the + /// engine still awaits the run and records its result: the tool decides + /// what cancel means, and the mutation's future is never dropped. async fn run_mutation( &self, ctx: &mut Context, call: &ProposedCall, + cancel: &CancellationToken, run_started: Instant, ) -> Result<(), Fenced> { let Some(spec) = self.offered_spec(ctx, &call.tool) else { @@ -881,12 +1132,12 @@ where } Claim::Granted => { self.emit(ctx, vec![started(call, &spec)]).await?; - let never = CancellationToken::new(); // A mutation that overruns its deadline is dropped, not // cancelled: the effect may still land. `Unknown` is recorded // under the claim, so a later resume of this call adopts it - // instead of dispatching the mutation a second time. - let run = self.tools.run(ctx.thread(), call, &never); + // instead of dispatching the mutation a second time. An + // interrupt only fires `cancel`; the run is still awaited. + let run = self.tools.run(ctx.thread(), call, cancel); let result = match tokio::time::timeout(self.call_deadline(run_started), run).await { Ok(result) => result, @@ -1188,6 +1439,60 @@ fn search_spec() -> ToolSpec { } } +/// Most thinking summary one attempt shows; the rest is dropped. +const MAX_THINKING_BYTES: usize = 16 * 1024; +/// Held thinking is written once it reaches this size (the first piece is +/// written at once, so progress appears as soon as the model starts). +const THINKING_FLUSH_BYTES: usize = 240; + +/// One attempt's thinking summary on its way to the log: sanitized like +/// answer text, written in bounded pieces, and never after the answer began. +struct Thinking { + filter: F, + held: String, + written: usize, +} + +impl Thinking { + fn new(filter: F) -> Self { + Self { + filter, + held: String::new(), + written: 0, + } + } + + /// Takes a thinking delta; returns a piece to write now, if any. + fn push(&mut self, delta: &str, answering: bool) -> Option { + if answering || self.written + self.held.len() >= MAX_THINKING_BYTES { + return None; + } + self.held.push_str(&self.filter.push(delta)); + if self.written == 0 || self.held.len() >= THINKING_FLUSH_BYTES { + return self.flush(); + } + None + } + + /// Whatever is held, bounded, once. + fn flush(&mut self) -> Option { + if self.held.is_empty() { + return None; + } + let room = MAX_THINKING_BYTES.saturating_sub(self.written); + let mut piece = std::mem::take(&mut self.held); + if piece.len() > room { + let mut end = room; + while !piece.is_char_boundary(end) { + end -= 1; + } + piece.truncate(end); + } + self.written += piece.len(); + (!piece.is_empty()).then_some(piece) + } +} + /// Appended to an answer whose model stream failed after text was shown. pub const CUT_OFF_NOTICE: &str = "\n\n_This answer was cut off before it finished. Ask me to continue from here._"; @@ -1214,6 +1519,23 @@ fn non_empty_str<'a>(call: &'a ProposedCall, key: &str) -> Option<&'a str> { .filter(|value| !value.trim().is_empty()) } +/// `args` against `spec.schema`. A schema that is absent, not an object, or +/// does not compile validates nothing (the executor still checks what it +/// needs); a schema violation names the first error so the model can call +/// again with arguments that fit. +fn validate_args(spec: &ToolSpec, args: &serde_json::Value) -> Result<(), String> { + if !spec.schema.is_object() { + return Ok(()); + } + let Ok(validator) = jsonschema::validator_for(&spec.schema) else { + return Ok(()); + }; + match validator.iter_errors(args).next() { + None => Ok(()), + Some(error) => Err(format!("invalid arguments: {error}")), + } +} + fn unknown_tool(name: &ToolName) -> ToolResult { ToolResult::error(format!( "unknown tool: {name}; use {TOOLS_SEARCH} to find tools" diff --git a/vendor/dex-loop/src/event.rs b/vendor/dex-loop/src/event.rs index 90f9de25b..6d8f5256a 100644 --- a/vendor/dex-loop/src/event.rs +++ b/vendor/dex-loop/src/event.rs @@ -138,6 +138,17 @@ impl fmt::Debug for ProviderReasoning { } } +/// The provider route and model that actually served one model step. With +/// ordered failover the serving route can differ from the configured +/// primary, so surfaces and metering read it from the step, not config. +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct ServedBy { + /// The model-gateway provider id, e.g. `vertex-anthropic`, `vertex-ai`. + pub provider: String, + /// The provider model id, e.g. `claude-opus-5-5`. + pub model: String, +} + /// Model spend reported by one model response. #[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize)] pub struct Usage { @@ -425,6 +436,13 @@ pub enum Event { TextDelta { text: String, }, + /// Customer-safe summary of the model's thinking in the current attempt + /// (through the `Sanitizer`), for surfaces to show as live progress until + /// the answer starts. Never part of the answer or of model history; the + /// engine bounds how much one attempt emits. + ThinkingDelta { + text: String, + }, Usage(Usage), /// The commit point of a model attempt: its full text and every proposed /// call with full arguments, appended before any policy check or @@ -440,6 +458,11 @@ pub enum Event { /// written before this field existed still deserialize, to `None`. #[serde(default, skip_serializing_if = "Option::is_none")] reasoning: Option, + /// The provider and model that served the step + /// (`ModelChunk::Served`); `None` for steps logged before the field + /// existed or from a model port that does not report it. + #[serde(default, skip_serializing_if = "Option::is_none")] + served: Option, }, /// A model attempt with no `ModelStepCompleted` (a crash mid-stream or a /// model failure). Its streamed text is dropped from model context; @@ -612,6 +635,7 @@ mod tests { PrincipalId::new("alice"), )], reasoning: None, + served: None, }, Event::ModelStepCompleted { step: 2, @@ -622,6 +646,10 @@ mod tests { model: "gemini-3.6-flash".into(), payload: serde_json::json!({"calls": [{"thought_signature": "c2ln"}]}), }), + served: Some(ServedBy { + provider: "vertex-ai".into(), + model: "gemini-3.6-flash".into(), + }), }, Event::ToolFinished { call: CallId::new("t1-1-0"), diff --git a/vendor/dex-loop/src/lib.rs b/vendor/dex-loop/src/lib.rs index fdc10aff5..014730014 100644 --- a/vendor/dex-loop/src/lib.rs +++ b/vendor/dex-loop/src/lib.rs @@ -29,14 +29,14 @@ mod rehydrate; mod sanitize; pub use budget::{Budget, BudgetAxis}; -pub use compaction::{Compaction, Compactor, NoCompaction, Summarize, Threshold}; +pub use compaction::{Compaction, Compactor, Cut, NoCompaction, Summarize, Threshold, plan_cut}; pub use context::{Context, Entry, Message}; pub use engine::{CUT_OFF_NOTICE, DEFAULT_TOOL_CALL_DEADLINE, Engine, Exit, TOOLS_SEARCH}; pub use event::{ AUTO_APPROVER, ApprovalId, ApprovalMode, ArtifactRef, CallId, ClientToolSpec, Cursor, ErrorCode, Event, HEADLESS_AUTO_APPROVER, MessageId, Outcome, Output, OutputRef, PrincipalId, - ProposedCall, ProviderReasoning, ReceiptId, ThreadId, ToolName, ToolResult, TurnId, Usage, - args_digest, + ProposedCall, ProviderReasoning, ReceiptId, ServedBy, ThreadId, ToolName, ToolResult, TurnId, + Usage, args_digest, }; pub use ports::{ Claim, Effects, ExecutorKind, Fenced, GovernanceClass, Log, Model, ModelChunk, ModelError, diff --git a/vendor/dex-loop/src/ports.rs b/vendor/dex-loop/src/ports.rs index ba5e7f64e..1759c55ed 100644 --- a/vendor/dex-loop/src/ports.rs +++ b/vendor/dex-loop/src/ports.rs @@ -9,8 +9,8 @@ use tokio_util::sync::CancellationToken; use crate::context::Context; use crate::event::{ - ApprovalId, CallId, Cursor, Event, PrincipalId, ProposedCall, ProviderReasoning, ThreadId, - ToolName, ToolResult, Usage, + ApprovalId, CallId, Cursor, Event, PrincipalId, ProposedCall, ProviderReasoning, ServedBy, + ThreadId, ToolName, ToolResult, Usage, }; /// The log or the effect ledger refused a write. The engine stops at once and @@ -53,6 +53,9 @@ pub trait Log: Send + Sync { #[derive(Clone, Debug, PartialEq)] pub enum ModelChunk { Text(String), + /// A summary of the model's thinking, streamed while it thinks. Shown as + /// progress only: never answer text, never model history, not usage. + Thinking(String), /// The engine assigns the `CallId`; provider call ids are not used. ToolCall { name: ToolName, @@ -65,6 +68,10 @@ pub enum ModelChunk { /// stores it on `ModelStepCompleted`; history returns it on /// `Message::Assistant`. Reasoning(ProviderReasoning), + /// Which provider route and model served this response. Sent at most + /// once, before the first other chunk; the engine stores it on + /// `ModelStepCompleted`. + Served(ServedBy), } /// The model call failed after the `Model` port's own retries. @@ -170,9 +177,13 @@ pub trait Tools: Send + Sync { /// caller-chosen, so the same `CallId` string can occur in two different /// threads. Where the downstream system needs a globally unique /// idempotency key, build it from `thread` and `call.id` together, not - /// `call.id` alone. For read-only calls `cancel` fires on interrupt and - /// the engine waits for the call to return; for mutations it never - /// fires. + /// `call.id` alone. `cancel` fires on interrupt for reads and mutations + /// alike, and the engine keeps awaiting the call to return either way: + /// it never drops a mutation's future, and it records whatever result + /// comes back in the effect ledger. What cancel means is the tool's to + /// decide. A read abandons its work; a mutation that can stop cleanly + /// (a background process, say) stops it and returns the partial outcome, + /// and one that cannot simply finishes. fn run( &self, thread: &ThreadId, diff --git a/vendor/dex-loop/tests/client_tool_replay.rs b/vendor/dex-loop/tests/client_tool_replay.rs index 60945f665..8b8d96870 100644 --- a/vendor/dex-loop/tests/client_tool_replay.rs +++ b/vendor/dex-loop/tests/client_tool_replay.rs @@ -65,6 +65,7 @@ fn log_up_to_model_step(log: &FakeLog, call: &ProposedCall) { text: String::new(), calls: vec![call.clone()], reasoning: None, + served: None, }); } diff --git a/vendor/dex-loop/tests/prefetch.rs b/vendor/dex-loop/tests/prefetch.rs new file mode 100644 index 000000000..d7dbca0d0 --- /dev/null +++ b/vendor/dex-loop/tests/prefetch.rs @@ -0,0 +1,497 @@ +//! Reads that start while the model is still streaming. +//! +//! A read-only call whose arguments fit its schema, whose policy allows it +//! and which runs on the host starts when its `ToolCall` chunk arrives: +//! `ToolStarted` lands before `ModelStepCompleted`. Everything else waits for +//! the commit point, and an attempt that never commits closes the reads it +//! started as not run, so no `ToolStarted` dangles and no result reaches +//! history. + +// Each `tests/*.rs` file is its own crate; this one uses a subset of the +// shared helpers. +#[allow(dead_code)] +mod support; + +use std::time::Duration; + +use dex_loop::{Budget, CancellationToken, Event, Exit, ModelError, Verdict}; +use serde_json::json; +use support::*; + +const CHUNK: Duration = Duration::from_millis(30); + +fn budget() -> Budget { + Budget { + max_steps: 10, + max_tokens: 1_000_000, + max_cost_micros: 1_000_000, + wall: Duration::from_secs(30), + } +} + +/// How many `ToolStarted` and `ToolFinished` rows the log holds for `call`. +fn started_and_finished(log: &FakeLog, call: &str) -> (usize, usize) { + let events = log.events(); + let started = events + .iter() + .filter(|event| matches!(event, Event::ToolStarted { call: id, .. } if id.as_str() == call)) + .count(); + let finished = events + .iter() + .filter( + |event| matches!(event, Event::ToolFinished { call: id, .. } if id.as_str() == call), + ) + .count(); + (started, finished) +} + +#[tokio::test] +async fn reads_start_while_the_model_streams_and_the_step_adopts_their_results() { + let log = FakeLog::default(); + let model = FakeModel::new(vec![ + vec![ + call("search", json!({"key": "a"})), + call("search", json!({"key": "b"})), + text("tail"), + ], + vec![text("done")], + ]) + .with_chunk_delay(CHUNK); + let tools = FakeTools::new(vec![read_tool("search")]) + .delay("a", Duration::from_millis(10)) + .delay("b", Duration::from_millis(10)); + let engine = engine(&log, &model, &tools, budget()); + let mut ctx = log.start_turn("t1", "look"); + + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) + ); + + assert_eq!( + log.shapes_after(1), + strings(&[ + "step:1", + "started:t1-1-0", + "started:t1-1-1", + "delta:tail", + "completed:tail:[t1-1-0,t1-1-1]", + "finished:t1-1-0:ok", + "finished:t1-1-1:ok", + "step:2", + "delta:done", + "completed:done:[]", + "final:done", + ]) + ); + // Adopted, not run again. + let mut runs = tools.run_ids(); + runs.sort(); + assert_eq!(runs, strings(&["t1-1-0", "t1-1-1"])); + assert_eq!( + view(&model.seen()[1]), + strings(&[ + "user:look", + "assistant:tail:[t1-1-0,t1-1-1]", + "tool:t1-1-0:ok:out/t1-1-0", + "tool:t1-1-1:ok:out/t1-1-1", + ]) + ); + assert_eq!(log.rehydrate(), ctx); +} + +#[tokio::test] +async fn a_read_that_is_still_running_at_the_commit_joins_the_wave() { + let log = FakeLog::default(); + let model = FakeModel::new(vec![ + vec![call("search", json!({"key": "a"}))], + vec![text("done")], + ]) + .with_chunk_delay(Duration::from_millis(5)); + let tools = FakeTools::new(vec![read_tool("search")]).delay("a", Duration::from_millis(150)); + let engine = engine(&log, &model, &tools, budget()); + let mut ctx = log.start_turn("t1", "look"); + + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) + ); + + assert_eq!( + log.shapes_after(1)[..4], + strings(&[ + "step:1", + "started:t1-1-0", + "completed::[t1-1-0]", + "finished:t1-1-0:ok", + ]) + ); + assert_eq!(tools.run_ids(), strings(&["t1-1-0"])); + assert_eq!(started_and_finished(&log, "t1-1-0"), (1, 1)); + assert_eq!(log.rehydrate(), ctx); +} + +#[tokio::test] +async fn a_mutation_never_starts_before_the_commit_and_holds_back_the_reads_after_it() { + let log = FakeLog::default(); + let model = FakeModel::new(vec![ + vec![ + call("search", json!({"key": "a"})), + call("update", json!({"key": "w"})), + call("search", json!({"key": "b"})), + ], + vec![text("ok")], + ]) + .with_chunk_delay(Duration::from_millis(5)); + let tools = FakeTools::new(vec![read_tool("search"), write_tool("update")]); + let engine = engine(&log, &model, &tools, budget()); + let mut ctx = log.start_turn("t1", "change it"); + + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) + ); + + assert_eq!( + log.shapes_after(1)[..8], + strings(&[ + "step:1", + "started:t1-1-0", + "completed::[t1-1-0,t1-1-1,t1-1-2]", + "finished:t1-1-0:ok", + "started:t1-1-1", + "finished:t1-1-1:ok", + "started:t1-1-2", + "finished:t1-1-2:ok", + ]) + ); + // The read after the mutation ran after it, as the model ordered them. + let write = tools.run_of(&call_id("t1", 1, 1)); + let later_read = tools.run_of(&call_id("t1", 1, 2)); + assert!(later_read.started >= write.finished); +} + +#[tokio::test] +async fn a_lone_mutation_starts_only_after_the_commit() { + let log = FakeLog::default(); + let model = FakeModel::new(vec![ + vec![call("update", json!({"key": "w"}))], + vec![text("ok")], + ]) + .with_chunk_delay(Duration::from_millis(5)); + let tools = FakeTools::new(vec![write_tool("update")]); + let engine = engine(&log, &model, &tools, budget()); + let mut ctx = log.start_turn("t1", "change it"); + + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) + ); + assert_eq!( + log.shapes_after(1)[..4], + strings(&[ + "step:1", + "completed::[t1-1-0]", + "started:t1-1-0", + "finished:t1-1-0:ok", + ]) + ); +} + +#[tokio::test] +async fn calls_that_wait_for_a_person_or_a_client_never_start_early() { + let log = FakeLog::default(); + let model = FakeModel::new(vec![vec![ + call("ask", json!({"question": "which one?"})), + call("search", json!({"key": "a"})), + ]]) + .with_chunk_delay(Duration::from_millis(5)); + let tools = FakeTools::new(vec![ask_tool("ask"), read_tool("search")]); + let engine = engine(&log, &model, &tools, budget()); + let mut ctx = log.start_turn("t1", "go"); + + let exit = engine.run(&mut ctx, &CancellationToken::new()).await; + + assert_eq!(exit, Ok(Exit::Asked(call_id("t1", 1, 0)))); + assert_eq!( + log.shapes_after(1)[..3], + strings(&[ + "step:1", + "completed::[t1-1-0,t1-1-1]", + "question:t1-1-0:which one?" + ]) + ); + assert!(tools.run_ids().is_empty()); + + let log = FakeLog::default(); + let model = FakeModel::new(vec![vec![call("browser.read", json!({}))]]) + .with_chunk_delay(Duration::from_millis(5)); + let tools = FakeTools::new(vec![client_executed_tool("browser.read", true)]); + let engine = support::engine(&log, &model, &tools, budget()); + let mut ctx = log.start_turn("t1", "go"); + + let exit = engine.run(&mut ctx, &CancellationToken::new()).await; + + assert_eq!(exit, Ok(Exit::AwaitingClientTool(call_id("t1", 1, 0)))); + assert_eq!( + log.shapes_after(1)[..2], + strings(&["step:1", "completed::[t1-1-0]"]) + ); + assert!(tools.run_ids().is_empty()); +} + +#[tokio::test] +async fn a_read_policy_does_not_allow_never_starts_early() { + for verdict in [Verdict::Deny("no grant".into()), approval("ap-1")] { + let denied = matches!(verdict, Verdict::Deny(_)); + let log = FakeLog::default(); + let model = FakeModel::new(vec![ + vec![call("search", json!({"key": "a"}))], + vec![text("ok")], + ]) + .with_chunk_delay(Duration::from_millis(5)); + let tools = FakeTools::new(vec![read_tool("search")]).verdict("search", verdict); + let engine = engine(&log, &model, &tools, budget()); + let mut ctx = log.start_turn("t1", "go"); + + engine.run(&mut ctx, &CancellationToken::new()).await.ok(); + + let shapes = log.shapes_after(1); + assert_eq!(shapes[..2], strings(&["step:1", "completed::[t1-1-0]"])); + if denied { + assert_eq!(shapes[2], "finished:t1-1-0:err"); + assert!(tools.run_ids().is_empty()); + } + } +} + +#[tokio::test] +async fn arguments_that_fail_the_schema_never_start_and_the_model_sees_why() { + let log = FakeLog::default(); + let model = FakeModel::new(vec![ + vec![ + call("search", json!({"key": 7})), + call("search", json!({"key": "ok"})), + ], + vec![text("done")], + ]) + .with_chunk_delay(Duration::from_millis(5)); + let tools = FakeTools::new(vec![strict_read_tool("search")]); + let engine = engine(&log, &model, &tools, budget()); + let mut ctx = log.start_turn("t1", "go"); + + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) + ); + + // The invalid call is the first call, so nothing after it starts early + // either (only an unbroken run of reads from the first call does). + assert_eq!( + log.shapes_after(1)[..5], + strings(&[ + "step:1", + "completed::[t1-1-0,t1-1-1]", + "finished:t1-1-0:err", + "started:t1-1-1", + "finished:t1-1-1:ok", + ]) + ); + assert_eq!(tools.run_ids(), strings(&["t1-1-1"])); + let seen = view(&model.seen()[1]); + assert!( + seen[2].starts_with("tool:t1-1-0:err:invalid arguments: "), + "{seen:?}" + ); + assert_eq!(log.rehydrate(), ctx); +} + +#[tokio::test] +async fn a_failed_stream_closes_the_reads_it_started_and_records_no_result() { + let log = FakeLog::default(); + let model = FakeModel::new(vec![vec![ + call("search", json!({"key": "a"})), + Err(ModelError { + message: "connection reset".into(), + }), + ]]) + .with_chunk_delay(Duration::from_millis(5)); + let tools = FakeTools::new(vec![read_tool("search")]).delay("a", Duration::from_millis(1)); + let engine = engine(&log, &model, &tools, budget()); + let mut ctx = log.start_turn("t1", "go"); + + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Failed) + ); + + let shapes = log.shapes_after(1); + assert_eq!( + shapes[..4], + strings(&[ + "step:1", + "started:t1-1-0", + "finished:t1-1-0:err", + "abandoned:1", + ]) + ); + assert!(shapes[4].starts_with("error:"), "{shapes:?}"); + assert_eq!(started_and_finished(&log, "t1-1-0"), (1, 1)); + // Nothing committed: the finished row is not a tool result in history. + assert_eq!(history(&ctx), strings(&["user:go"])); + assert_eq!(log.rehydrate(), ctx); +} + +#[tokio::test] +async fn a_new_turn_after_a_failed_stream_starts_its_own_calls_once() { + let log = FakeLog::default(); + let model = FakeModel::new(vec![ + vec![ + call("search", json!({"key": "a"})), + Err(ModelError { + message: "connection reset".into(), + }), + ], + vec![call("search", json!({"key": "a"}))], + vec![text("done")], + ]) + .with_chunk_delay(Duration::from_millis(5)); + let tools = FakeTools::new(vec![read_tool("search")]); + let engine = engine(&log, &model, &tools, budget()); + let mut ctx = log.start_turn("t1", "go"); + let cancel = CancellationToken::new(); + assert_eq!(engine.run(&mut ctx, &cancel).await, Ok(Exit::Failed)); + + let mut ctx = log.start_turn("t2", "again"); + assert_eq!(engine.run(&mut ctx, &cancel).await, Ok(Exit::Done)); + + for id in ["t1-1-0", "t2-1-0"] { + assert_eq!(started_and_finished(&log, id), (1, 1), "{id}"); + } + // The first turn's read ran (a read has no effect), but its result was + // never committed; the second turn's call ran once. + assert_eq!(tools.run_ids(), strings(&["t1-1-0", "t2-1-0"])); + assert!( + history(&ctx).iter().all(|row| !row.contains("t1-1-0")), + "{:?}", + history(&ctx) + ); + assert_eq!(log.rehydrate(), ctx); +} + +#[tokio::test] +async fn an_interrupt_during_the_stream_ends_cleanly_with_a_read_in_flight() { + let log = FakeLog::default(); + let model = FakeModel::new(vec![vec![call("search", json!({"key": "a"}))]]) + .with_chunk_delay(Duration::from_millis(5)) + .hanging(); + let tools = FakeTools::new(vec![read_tool("search")]).delay("a", Duration::from_secs(10)); + let engine = engine(&log, &model, &tools, budget()); + let mut ctx = log.start_turn("t1", "go"); + let cancel = CancellationToken::new(); + + let host = async { + tokio::time::sleep(Duration::from_millis(100)).await; + log.host_append(Event::Interrupt { principal: alice() }); + cancel.cancel(); + }; + let started = std::time::Instant::now(); + let (exit, ()) = tokio::join!(engine.run(&mut ctx, &cancel), host); + + assert_eq!(exit, Ok(Exit::Interrupted)); + assert!(started.elapsed() < Duration::from_secs(2)); + assert_eq!( + log.shapes_after(1), + strings(&[ + "step:1", + "started:t1-1-0", + "interrupt", + "finished:t1-1-0:err", + "completed::[]", + "interrupted", + ]) + ); + assert_eq!(started_and_finished(&log, "t1-1-0"), (1, 1)); + assert_eq!(history(&ctx), strings(&["user:go", "assistant::[]"])); +} + +#[tokio::test] +async fn a_stream_that_outlives_the_wall_budget_closes_its_reads() { + let log = FakeLog::default(); + let model = FakeModel::new(vec![vec![call("search", json!({"key": "a"}))]]) + .with_chunk_delay(Duration::from_millis(5)) + .hanging(); + let tools = FakeTools::new(vec![read_tool("search")]).delay("a", Duration::from_secs(10)); + let engine = engine( + &log, + &model, + &tools, + Budget { + wall: Duration::from_millis(150), + ..budget() + }, + ); + let mut ctx = log.start_turn("t1", "go"); + + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Failed) + ); + + let shapes = log.shapes_after(1); + assert_eq!( + shapes[..4], + strings(&[ + "step:1", + "started:t1-1-0", + "finished:t1-1-0:err", + "abandoned:1", + ]) + ); + assert_eq!(started_and_finished(&log, "t1-1-0"), (1, 1)); +} + +#[tokio::test] +async fn a_restart_after_a_read_started_mid_stream_closes_it_and_runs_the_step_once() { + let log = FakeLog::default(); + // The first engine died mid-stream: the read's `ToolStarted` is on the + // log, the step never committed. + log.start_turn("t1", "look"); + log.host_append(Event::StepStarted { + step: 1, + control_through: dex_loop::Cursor::START, + }); + log.host_append(Event::ToolStarted { + call: call_id("t1", 1, 0), + tool: dex_loop::ToolName::new("search"), + label: "Label for search".into(), + principal: alice(), + }); + let model = FakeModel::new(vec![ + vec![call("search", json!({"key": "a"}))], + vec![text("done")], + ]) + .with_chunk_delay(Duration::from_millis(5)); + let tools = FakeTools::new(vec![read_tool("search")]); + let engine = engine(&log, &model, &tools, budget()); + let mut ctx = log.rehydrate(); + + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) + ); + + assert_eq!( + log.shapes_after(3)[..6], + strings(&[ + "finished:t1-1-0:err", + "abandoned:1", + "step:2", + "started:t1-2-0", + "completed::[t1-2-0]", + "finished:t1-2-0:ok", + ]) + ); + assert_eq!(tools.run_ids(), strings(&["t1-2-0"])); + assert_eq!(log.rehydrate(), ctx); +} diff --git a/vendor/dex-loop/tests/scenarios.rs b/vendor/dex-loop/tests/scenarios.rs index 388ced1b4..e7c5eef60 100644 --- a/vendor/dex-loop/tests/scenarios.rs +++ b/vendor/dex-loop/tests/scenarios.rs @@ -8,8 +8,8 @@ use std::time::{Duration, Instant}; use dex_loop::{ ApprovalId, ApprovalMode, Budget, CUT_OFF_NOTICE, CancellationToken, Cursor, Engine, Event, - Exit, Fenced, Lexicon, ModelError, OutputRef, PrincipalId, ProposedCall, Threshold, ToolName, - ToolResult, TurnId, Verdict, + Exit, Fenced, Lexicon, ModelError, Output, OutputRef, PrincipalId, ProposedCall, Threshold, + ToolName, ToolResult, TurnId, Verdict, }; use serde_json::json; use support::*; @@ -100,10 +100,10 @@ async fn read_only_wave_overlaps_and_keeps_call_order() { strings(&[ "user:look", "step:1", - "completed::[t1-1-0,t1-1-1,t1-1-2]", "started:t1-1-0", "started:t1-1-1", "started:t1-1-2", + "completed::[t1-1-0,t1-1-1,t1-1-2]", "finished:t1-1-2:ok", "finished:t1-1-1:ok", "finished:t1-1-0:ok", @@ -157,9 +157,9 @@ async fn mutating_call_runs_serially_after_the_wave() { strings(&[ "user:change it", "step:1", - "completed::[t1-1-0,t1-1-1,t1-1-2,t1-1-3]", "started:t1-1-0", "started:t1-1-1", + "completed::[t1-1-0,t1-1-1,t1-1-2,t1-1-3]", "finished:t1-1-0:ok", "finished:t1-1-1:ok", "started:t1-1-2", @@ -217,8 +217,8 @@ async fn an_approval_class_call_is_granted_at_once_recorded_and_runs_in_order() log.shapes_after(1), strings(&[ "step:1", - "completed::[t1-1-0,t1-1-1,t1-1-2]", "started:t1-1-0", + "completed::[t1-1-0,t1-1-1,t1-1-2]", "finished:t1-1-0:ok", "auto_approved:t1-1-1", "started:t1-1-1", @@ -336,7 +336,8 @@ async fn a_receipt_survives_a_crash_and_the_call_runs_once_after_rehydrate() { #[tokio::test] async fn a_legacy_parked_call_is_granted_on_rehydrate_and_the_step_continues() { let log = FakeLog::default(); - let model = FakeModel::new(vec![vec![], vec![text("sent")]]); + // Step 1 is already committed in the log; the only model call is step 2. + let model = FakeModel::new(vec![vec![text("sent")]]); let tools = FakeTools::new(vec![read_tool("search"), write_tool("send_email")]) .verdict("send_email", approval("ap-1")); let send = ProposedCall::new( @@ -371,6 +372,7 @@ async fn a_legacy_parked_call_is_granted_on_rehydrate_and_the_step_continues() { text: String::new(), calls: vec![send.clone(), search], reasoning: None, + served: None, }, Event::ApprovalRequested { call: send.id.clone(), @@ -491,8 +493,8 @@ async fn steer_from_another_principal_is_checked_under_that_principal() { strings(&[ "user:check prod", "step:1", - "completed::[t1-1-0]", "started:t1-1-0", + "completed::[t1-1-0]", "steer:also update staging", "finished:t1-1-0:ok", "step:2", @@ -606,9 +608,9 @@ async fn interrupt_mid_wave_cancels_reads_and_emits_interrupted() { strings(&[ "user:go", "step:1", - "completed::[t1-1-0,t1-1-1,t1-1-2]", "started:t1-1-0", "started:t1-1-1", + "completed::[t1-1-0,t1-1-1,t1-1-2]", "interrupt", ]) ); @@ -631,17 +633,20 @@ async fn interrupt_mid_wave_cancels_reads_and_emits_interrupted() { ); } -// 6b. Interrupt during a mutation lets that mutation complete, then stops -// before the next effect. +// 6b. Interrupt during a mutation reaches the running tool through `cancel`, +// but the engine still awaits the run to completion, records its result in +// the effect ledger, finishes the call exactly once, and stops before the +// next effect. #[tokio::test] -async fn interrupt_during_a_mutation_completes_it_then_stops() { +async fn interrupt_during_a_mutation_cancels_it_records_its_result_then_stops() { let log = FakeLog::default(); let model = FakeModel::new(vec![vec![ call("update", json!({"key": "w"})), call("update", json!({"key": "x"})), ]]); - let tools = FakeTools::new(vec![write_tool("update")]).delay("w", Duration::from_millis(300)); - let engine = engine(&log, &model, &tools, budget()); + let tools = FakeTools::new(vec![write_tool("update")]).delay("w", Duration::from_secs(30)); + let effects = FakeEffects::default(); + let engine = engine_with(&log, &model, &tools, &effects, budget()); let mut ctx = log.start_turn("t1", "go"); let cancel = CancellationToken::new(); @@ -650,14 +655,18 @@ async fn interrupt_during_a_mutation_completes_it_then_stops() { log.host_append(Event::Interrupt { principal: alice() }); cancel.cancel(); }; + let started = Instant::now(); let (exit, ()) = tokio::join!(engine.run(&mut ctx, &cancel), host); assert_eq!(exit, Ok(Exit::Interrupted)); + assert!(started.elapsed() < Duration::from_secs(5)); assert_eq!(tools.run_ids(), strings(&["t1-1-0"])); assert!( - !tools.runs()[0].cancelled, - "a started mutation was cancelled" + tools.runs()[0].cancelled, + "the running mutation did not observe the interrupt" ); + // Exactly one ToolFinished for the mutation, then the unstarted call is + // closed and the turn ends. assert_eq!( log.shapes(), strings(&[ @@ -666,11 +675,17 @@ async fn interrupt_during_a_mutation_completes_it_then_stops() { "completed::[t1-1-0,t1-1-1]", "started:t1-1-0", "interrupt", - "finished:t1-1-0:ok", + "finished:t1-1-0:err", "finished:t1-1-1:err", "interrupted", ]) ); + // The engine awaited the run and recorded what the tool returned. + let recorded = effects + .recorded(&call_id("t1", 1, 0)) + .flatten() + .expect("the cancelled mutation's result is recorded under its claim"); + assert_eq!(recorded.output, Output::Text("cancelled".into())); } // 7. Each budget axis stops the turn with budget_exhausted. @@ -705,8 +720,8 @@ async fn step_cap_gets_an_answer_only_step_that_ends_in_final() { log.shapes_after(1), strings(&[ "step:1", - "completed::[t1-1-0]", "started:t1-1-0", + "completed::[t1-1-0]", "finished:t1-1-0:ok", "step:2", "delta:here is what I found", @@ -1034,6 +1049,41 @@ async fn a_stream_that_fails_after_text_keeps_the_answer_marked_cut_off() { ); } +// 8d4. Thinking summaries stream as progress before the answer: sanitized, +// never part of the answer text, the committed step, or the history. +#[tokio::test] +async fn thinking_streams_as_progress_and_never_joins_the_answer() { + let log = FakeLog::default(); + let model = FakeModel::new(vec![vec![ + thinking("Checking "), + thinking(&"x".repeat(300)), + thinking("tail"), + text("Hello"), + thinking("after the answer began"), + usage(3, 2, 7), + ]]); + let tools = FakeTools::new(vec![]); + let engine = engine(&log, &model, &tools, budget()); + let mut ctx = log.start_turn("t1", "hi"); + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) + ); + assert_eq!( + log.shapes_after(2), + strings(&[ + "thinking:Checking ", + &format!("thinking:{}", "x".repeat(300)), + "thinking:tail", + "delta:Hello", + "usage:5", + "completed:Hello:[]", + "final:Hello", + ]) + ); + assert_eq!(log.rehydrate(), ctx); +} + // 8e. A `Started` call whose tool vanished from the offered catalog (deploy, // grant revoke) resolves through the ledger, never as "unknown tool": the // model must not be told to retry a call that may have already run under @@ -1166,9 +1216,9 @@ async fn fenced_during_a_wave_drops_the_remaining_results() { let tools = FakeTools::new(vec![read_tool("search")]).delay("b", Duration::from_millis(50)); let engine = engine(&log, &model, &tools, budget()); let mut ctx = log.start_turn("t1", "go"); - // StepStarted, ModelStepCompleted and the ToolStarted batch land; the + // StepStarted, both ToolStarted rows and ModelStepCompleted land; the // first ToolFinished is refused. - log.fence_after(3); + log.fence_after(4); let result = engine.run(&mut ctx, &CancellationToken::new()).await; @@ -1179,9 +1229,9 @@ async fn fenced_during_a_wave_drops_the_remaining_results() { strings(&[ "user:go", "step:1", - "completed::[t1-1-0,t1-1-1]", "started:t1-1-0", "started:t1-1-1", + "completed::[t1-1-0,t1-1-1]", ]) ); } @@ -1245,10 +1295,10 @@ async fn unknown_and_denied_calls_are_visible_to_the_model() { assert_eq!(log.rehydrate(), ctx); } -// 10b. Headless turns: no human can answer, so a `NeedsApproval` verdict is -// approved by policy. The log keeps the request and the decision, attributed -// to the policy principal; a `Deny` verdict stays denied; an interactive turn -// still parks. +// 10b. No human approves a Dex call (#11467): a `NeedsApproval` verdict is +// granted at once on every turn, headless or interactive, and the log keeps +// one `AutoApproved` receipt attributed to `AUTO_APPROVER`. A `Deny` verdict +// stays denied. fn headless_tools() -> FakeTools { FakeTools::new(vec![write_tool("send_email"), write_tool("delete_all")]) .verdict("send_email", approval("ap-1")) @@ -1258,30 +1308,17 @@ fn headless_tools() -> FakeTools { ) } -#[tokio::test] -async fn headless_turn_auto_approves_an_ask_gated_tool_and_records_the_audit_pair() { - let log = FakeLog::default(); - let model = FakeModel::new(vec![ - vec![call("send_email", json!({"key": "w"}))], - vec![text("sent")], - ]); - let tools = headless_tools(); - let engine = engine(&log, &model, &tools, budget()); - let mut ctx = log.start_turn_with_approval_mode("t1", "email them", ApprovalMode::Headless); - assert_eq!(ctx.approval_mode(), ApprovalMode::Headless); - - assert_eq!( - engine.run(&mut ctx, &CancellationToken::new()).await, - Ok(Exit::Done) - ); +/// The shared outcome of an ask-gated `send_email` on any turn: one +/// `AutoApproved` receipt under policy's approval id and the call's digest, +/// the call runs once, and the next step answers. +fn assert_auto_approved_and_sent(log: &FakeLog, tools: &FakeTools) { assert_eq!(tools.run_ids(), strings(&["t1-1-0"])); assert_eq!( log.shapes_after(1), strings(&[ "step:1", "completed::[t1-1-0]", - "approval:t1-1-0", - "decided:t1-1-0:true", + "auto_approved:t1-1-0", "started:t1-1-0", "finished:t1-1-0:ok", "step:2", @@ -1291,31 +1328,45 @@ async fn headless_turn_auto_approves_an_ask_gated_tool_and_records_the_audit_pai ]) ); let call = call_id("t1", 1, 0); + let digest = dex_loop::args_digest(&tools.run_of(&call).args); let events = log.events(); - let requested = events - .iter() - .find_map(|event| match event { - Event::ApprovalRequested { - approval, - args_digest, - .. - } => Some((approval.clone(), args_digest.clone())), - _ => None, - }) - .expect("approval_requested is recorded"); - assert!(events.iter().any(|event| matches!( - event, - Event::ApprovalDecided { - call: decided_call, - approval, - args_digest, - approved: true, - principal, - } if *decided_call == call - && *approval == requested.0 - && *args_digest == requested.1 - && principal.as_str() == dex_loop::HEADLESS_AUTO_APPROVER - ))); + assert!( + events.iter().any(|event| matches!( + event, + Event::AutoApproved { call: receipt_call, approval, args_digest, principal, .. } + if *receipt_call == call + && approval.as_str() == "ap-1" + && *args_digest == digest + && principal.as_str() == dex_loop::AUTO_APPROVER + )), + "{events:?}" + ); + assert!( + !events.iter().any(|event| matches!( + event, + Event::ApprovalRequested { .. } | Event::ApprovalDecided { .. } + )), + "no human prompt is written" + ); +} + +#[tokio::test] +async fn headless_turn_auto_approves_an_ask_gated_tool_and_records_the_audit_pair() { + let log = FakeLog::default(); + let model = FakeModel::new(vec![ + vec![call("send_email", json!({"key": "w"}))], + vec![text("sent")], + ]); + let tools = headless_tools(); + let engine = engine(&log, &model, &tools, budget()); + let mut ctx = log.start_turn_with_approval_mode("t1", "email them", ApprovalMode::Headless); + assert_eq!(ctx.approval_mode(), ApprovalMode::Headless); + + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) + ); + assert_auto_approved_and_sent(&log, &tools); // Resumes keep the mode: it is on the logged UserMessage. let rehydrated = log.rehydrate(); assert_eq!(rehydrated.approval_mode(), ApprovalMode::Headless); @@ -1352,22 +1403,22 @@ async fn headless_turn_keeps_a_hard_deny_denied() { } #[tokio::test] -async fn interactive_turn_still_parks_on_the_same_ask_gated_tool() { +async fn interactive_turn_auto_approves_the_same_ask_gated_tool() { let log = FakeLog::default(); - let model = FakeModel::new(vec![vec![call("send_email", json!({"key": "w"}))]]); + let model = FakeModel::new(vec![ + vec![call("send_email", json!({"key": "w"}))], + vec![text("sent")], + ]); let tools = headless_tools(); let engine = engine(&log, &model, &tools, budget()); let mut ctx = log.start_turn_with_approval_mode("t1", "email them", ApprovalMode::Interactive); assert_eq!( engine.run(&mut ctx, &CancellationToken::new()).await, - Ok(Exit::Parked(ApprovalId::new("ap-1"))) - ); - assert!(tools.run_ids().is_empty()); - assert_eq!( - log.shapes_after(1), - strings(&["step:1", "completed::[t1-1-0]", "approval:t1-1-0"]) + Ok(Exit::Done) ); + assert_auto_approved_and_sent(&log, &tools); + assert_eq!(log.rehydrate(), ctx); } #[tokio::test] @@ -1390,8 +1441,7 @@ async fn headless_turn_auto_approves_a_mutating_client_tool() { strings(&[ "step:1", "completed::[t1-1-0]", - "approval:t1-1-0", - "decided:t1-1-0:true", + "auto_approved:t1-1-0", "client_tool:t1-1-0:browser.click", ]) ); @@ -1521,8 +1571,8 @@ async fn tools_search_exposes_schemas_for_the_next_step() { "exposed:[crm.lookup]", "finished:t1-1-0:ok", "step:2", - "completed::[t1-2-0]", "started:t1-2-0", + "completed::[t1-2-0]", "finished:t1-2-0:ok", "step:3", "delta:found acme", @@ -1646,10 +1696,10 @@ async fn client_tool_is_requested_and_result_resumes() { assert_eq!(log.rehydrate(), ctx); } -// A mutating client tool parks for approval before it is requested; denying -// it never reaches the client. +// A mutating client tool gets its `AutoApproved` receipt before it is +// requested from the client; nothing parks for a person. #[tokio::test] -async fn a_mutating_client_tool_needs_approval_first() { +async fn a_mutating_client_tool_is_auto_approved_before_it_is_requested() { let log = FakeLog::default(); let model = FakeModel::new(vec![vec![call( "browser.click", @@ -1663,23 +1713,21 @@ async fn a_mutating_client_tool_needs_approval_first() { let call = call_id("t1", 1, 0); assert_eq!( engine.run(&mut ctx, &cancel).await, - Ok(Exit::Parked(ApprovalId::new(format!("client-{call}")))) + Ok(Exit::AwaitingClientTool(call)) ); - log.decide(&call, &format!("client-{call}"), false); - assert_eq!(engine.run(&mut ctx, &cancel).await, Ok(Exit::Failed)); assert_eq!( log.shapes_after(1), strings(&[ "step:1", "completed::[t1-1-0]", - "approval:t1-1-0", - "decided:t1-1-0:false", - "finished:t1-1-0:err", - "step:2", - "abandoned:2", - "error:model_failed:no script left", + "auto_approved:t1-1-0", + "client_tool:t1-1-0:browser.click", ]) ); + assert!( + tools.runs().is_empty(), + "the client, not Tools::run, executes it" + ); assert_eq!(log.rehydrate(), ctx); } @@ -1816,11 +1864,11 @@ async fn a_steer_from_before_the_rehydrate_point_is_not_carried_into_a_later_tur ); } -// Provider reasoning is committed with its step and survives a park and a -// crash: after a fresh engine rehydrates the log and the approval resumes -// the step, the next model call sees it on the assistant message. +// Provider reasoning is committed with its step and survives a crash: after +// a fresh engine rehydrates the log and resumes the step, the next model call +// sees it on the assistant message. #[tokio::test] -async fn reasoning_is_committed_with_its_step_and_returned_after_park_and_crash() { +async fn reasoning_is_committed_with_its_step_and_returned_after_a_crash() { let reasoning = dex_loop::ProviderReasoning { format: "google.gemini.v1".into(), model: "gemini-3.6-flash".into(), @@ -1841,10 +1889,13 @@ async fn reasoning_is_committed_with_its_step_and_returned_after_park_and_crash( let engine = engine(&log, &model, &tools, budget()); let mut ctx = log.start_turn("t1", "email them"); let cancel = CancellationToken::new(); - assert_eq!( + // user, step, completed, auto_approved land; the `ToolStarted` write is + // refused, which is the crash. + log.fence_after(3); + assert!(matches!( engine.run(&mut ctx, &cancel).await, - Ok(Exit::Parked(ApprovalId::new("ap-1"))) - ); + Err(Fenced { .. }) + )); let committed: Vec<_> = log .events() .into_iter() @@ -1855,9 +1906,9 @@ async fn reasoning_is_committed_with_its_step_and_returned_after_park_and_crash( .collect(); assert_eq!(committed, vec![Some(reasoning.clone())]); - // -- crash: a fresh engine and context from the log, then the approval -- + // -- crash: a fresh engine and context from the log -- + log.fence_after(usize::MAX); let engine = support::engine(&log, &model, &tools, budget()); - log.decide(&call_id("t1", 1, 0), "ap-1", true); let mut ctx = log.rehydrate(); assert_eq!(engine.run(&mut ctx, &cancel).await, Ok(Exit::Done)); diff --git a/vendor/dex-loop/tests/sim/invariants.rs b/vendor/dex-loop/tests/sim/invariants.rs index 31717b4db..99506aea7 100644 --- a/vendor/dex-loop/tests/sim/invariants.rs +++ b/vendor/dex-loop/tests/sim/invariants.rs @@ -65,6 +65,7 @@ fn kind_str(event: &Event) -> &'static str { Event::ClientToolResult { .. } => "client_tool_result", Event::StepStarted { .. } => "step_started", Event::TextDelta { .. } => "text_delta", + Event::ThinkingDelta { .. } => "thinking_delta", Event::Usage(_) => "usage", Event::ModelStepCompleted { .. } => "model_step_completed", Event::ModelAttemptAbandoned { .. } => "model_attempt_abandoned", @@ -357,6 +358,32 @@ pub fn check_log( } } + // (7) Only reads start ahead of their step's commit point: a mutation + // or client call whose `ToolStarted` precedes the `ModelStepCompleted` + // that proposes it began before the model finished. + let mut committed_at: HashMap<&CallId, Cursor> = HashMap::new(); + for (cursor, event) in events { + if let Event::ModelStepCompleted { calls, .. } = event { + for call in calls { + committed_at.insert(&call.id, *cursor); + } + } + } + for (call, history) in &calls { + let (Some(started), Some(committed)) = (history.started_at, committed_at.get(call)) else { + continue; + }; + let is_read = !mutation_names.contains(history.tool.as_str()) + && !client_names.contains(history.tool.as_str()); + if started < *committed && !is_read { + let violation = Violation::new(format!( + "call {call} ({}) started before its step committed, but is not a read", + history.tool + )); + violations.push(tag_known_pending(violation, &history.tool, client_names)); + } + } + // (6) Interrupt never turns a started call into "not run". for (call, history) in &calls { if history.started_at.is_none() { diff --git a/vendor/dex-loop/tests/support/mod.rs b/vendor/dex-loop/tests/support/mod.rs index 13daec34a..45fa76f91 100644 --- a/vendor/dex-loop/tests/support/mod.rs +++ b/vendor/dex-loop/tests/support/mod.rs @@ -198,32 +198,6 @@ impl FakeLog { pub fn refused(&self) -> usize { lock(&self.state).refused } - - /// The digest the engine put on the approval request for `call`. - pub fn requested_digest(&self, call: &CallId) -> String { - self.events() - .into_iter() - .find_map(|event| match event { - Event::ApprovalRequested { - call: requested, - args_digest, - .. - } if &requested == call => Some(args_digest), - _ => None, - }) - .unwrap_or_else(|| panic!("no approval requested for {call}")) - } - - pub fn decide(&self, call: &CallId, approval: &str, approved: bool) { - let args_digest = self.requested_digest(call); - self.host_append(Event::ApprovalDecided { - call: call.clone(), - approval: ApprovalId::new(approval), - args_digest, - approved, - principal: alice(), - }); - } } impl Log for FakeLog { @@ -264,6 +238,10 @@ pub fn text(text: &str) -> Result { Ok(ModelChunk::Text(text.into())) } +pub fn thinking(text: &str) -> Result { + Ok(ModelChunk::Thinking(text.into())) +} + pub fn call(name: &str, args: serde_json::Value) -> Result { Ok(ModelChunk::ToolCall { name: ToolName::new(name), @@ -287,6 +265,8 @@ pub fn usage( struct ModelState { scripts: VecDeque>>, chunk_delay: Duration, + /// After its script, a call's stream never yields again or ends. + hang: bool, seen: Vec>, offered: Vec>, } @@ -309,6 +289,14 @@ impl FakeModel { self } + /// Each call's stream stays open after its script: no further chunk and + /// no end, as a provider that stalls mid-response. + #[allow(dead_code)] // used by tests/prefetch.rs only + pub fn hanging(self) -> Self { + lock(&self.state).hang = true; + self + } + /// The history of every model call, in order. pub fn seen(&self) -> Vec> { lock(&self.state).seen.clone() @@ -346,12 +334,21 @@ impl Model for FakeModel { })] }); let delay = state.chunk_delay; - stream::iter(script).then(move |chunk| async move { - if !delay.is_zero() { - tokio::time::sleep(delay).await; + let hang = state.hang; + let hold = stream::once(async move { + if hang { + std::future::pending::<()>().await; } - chunk }) + .filter_map(|()| async { None }); + stream::iter(script) + .then(move |chunk| async move { + if !delay.is_zero() { + tokio::time::sleep(delay).await; + } + chunk + }) + .chain(hold) } } @@ -361,6 +358,19 @@ pub fn read_tool(name: &str) -> ToolSpec { spec(name, true, true, ExecutorKind::InProcess) } +/// A read whose arguments must be an object with a string `key`. +#[allow(dead_code)] // used by tests/prefetch.rs only +pub fn strict_read_tool(name: &str) -> ToolSpec { + ToolSpec { + schema: serde_json::json!({ + "type": "object", + "properties": {"key": {"type": "string"}}, + "required": ["key"], + }), + ..read_tool(name) + } +} + pub fn write_tool(name: &str) -> ToolSpec { spec(name, false, true, ExecutorKind::ToolExecutor) } @@ -734,6 +744,7 @@ pub fn shape(event: &Event) -> String { Event::Answer { call, text, .. } => format!("answer:{call}:{text}"), Event::StepStarted { step, .. } => format!("step:{step}"), Event::TextDelta { text } => format!("delta:{text}"), + Event::ThinkingDelta { text } => format!("thinking:{text}"), Event::Usage(usage) => format!("usage:{}", usage.tokens()), Event::ModelStepCompleted { text, calls, .. } => { format!("completed:{text}:[{}]", ids(calls)) @@ -838,6 +849,7 @@ pub fn crashed_after_start(log: &FakeLog, call: &ProposedCall) { text: String::new(), calls: vec![call.clone()], reasoning: None, + served: None, }, Event::ToolStarted { call: call.id.clone(), diff --git a/vendor/dex-loop/tests/tool_deadline.rs b/vendor/dex-loop/tests/tool_deadline.rs index dd0dbebca..84836f90f 100644 --- a/vendor/dex-loop/tests/tool_deadline.rs +++ b/vendor/dex-loop/tests/tool_deadline.rs @@ -58,8 +58,8 @@ async fn a_read_that_overruns_the_deadline_is_finished_failed_and_the_turn_conti log.shapes_after(1), strings(&[ "step:1", - "completed::[t1-1-0]", "started:t1-1-0", + "completed::[t1-1-0]", "finished:t1-1-0:err", "step:2", "delta:done", @@ -170,8 +170,8 @@ async fn a_call_never_outlives_the_wall_budget() { log.shapes_after(1), strings(&[ "step:1", - "completed::[t1-1-0]", "started:t1-1-0", + "completed::[t1-1-0]", "finished:t1-1-0:err", "error:budget_exhausted:wall budget exhausted: 100ms", ]) @@ -199,8 +199,8 @@ async fn a_call_that_finishes_in_time_is_unaffected() { log.shapes_after(1), strings(&[ "step:1", - "completed::[t1-1-0]", "started:t1-1-0", + "completed::[t1-1-0]", "finished:t1-1-0:ok", "step:2", "delta:done",