From 8a5db9300f651dd772c4db930614ff93a5f3e404 Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" Date: Thu, 1 Oct 2026 07:13:24 +0000 Subject: [PATCH] chore: sync public mirror from internal --- .repository-projection.json | 6 +- deny.toml | 1 - packages/dex-host-rs/src/turn.rs | 118 +++- vendor/dex-loop/src/compaction.rs | 577 ++++++++++++++++-- vendor/dex-loop/src/context.rs | 289 ++++++++- vendor/dex-loop/src/engine.rs | 94 ++- vendor/dex-loop/src/event.rs | 31 +- vendor/dex-loop/src/lib.rs | 15 +- vendor/dex-loop/src/ports.rs | 3 +- vendor/dex-loop/tests/action_confirmation.rs | 405 +++++++++++++ vendor/dex-loop/tests/scenarios.rs | 591 +++++++++++++++++-- vendor/dex-loop/tests/sim/scenario.rs | 2 + vendor/dex-loop/tests/support/mod.rs | 12 +- vendor/dex-loop/tests/tool_deadline.rs | 19 +- 14 files changed, 2023 insertions(+), 140 deletions(-) create mode 100644 vendor/dex-loop/tests/action_confirmation.rs diff --git a/.repository-projection.json b/.repository-projection.json index a23b41ae0..528131514 100644 --- a/.repository-projection.json +++ b/.repository-projection.json @@ -3,11 +3,11 @@ "projection": "deixic-code", "projectionSchemaVersion": 1, "sourceRepository": "dx-corp/mono", - "sourceSha": "4c3e830640c4066c99303b7582fa2bdb74095f0d", + "sourceSha": "10a53803cedbfd671b5b3d4eee881522e7561c81", "destinationRepository": "dx-corp/code", - "priorProjectedBase": "ca90e530cac35a9d63eac51f920d58891a511b21", + "priorProjectedBase": "ebbe67aa0f540c41f24041efddb6cfb7f9075224", "definitionDigest": "82936441c776e3e8edb5d215a75007ec9714a233f489d460075d79d5ef5ba32f", "toolDigest": "c244d99199a7ae3eb8ff644a99462163c23b0bb6a83ef50af01efbdca0b81d04", - "contentDigest": "f684b3dd9cbf2e4f672beaa7d0c536d87035881e1cdb7f5d28afeebe1e49b530", + "contentDigest": "03efaffc773c89ec597f75013e916338b0e87d2c3f534416fd784dd8e536f907", "publicationEligible": true } diff --git a/deny.toml b/deny.toml index 3cb84f51a..2d4616bbd 100644 --- a/deny.toml +++ b/deny.toml @@ -36,7 +36,6 @@ ignore = [ # --- Unsound (no known exploitable vulnerability reported here) --- { id = "RUSTSEC-2026-0002", reason = "lru 0.12.5: IterMut violates Stacked Borrows (Miri/unsafe-code soundness issue), pulled in by ratatui (production). Not a memory-safety issue under normal (non-Miri) execution; tracked for the next ratatui bump that picks up a fixed lru. expires: 2026-10-23" }, - { id = "RUSTSEC-2026-0285", reason = "rustls <0.23.45 accepts TLS 1.3 handshake messages at the wrong encryption level after a key change in the same record (GHSA-2mjx-qc3c-rqvc, CVSS 5.3). Fixed in rustls 0.23.45, published 2026-09-14; that upgrade also resolves aws-lc-rs 1.18.1 and aws-lc-sys 0.45.0, all still inside the 14-day dependency-admission cool-off in security/rust-dependency-admission.toml. Bump the lockfile once they age past it. expires: 2026-09-30" }, ] [licenses] diff --git a/packages/dex-host-rs/src/turn.rs b/packages/dex-host-rs/src/turn.rs index fa1ed2838..21d835bbb 100644 --- a/packages/dex-host-rs/src/turn.rs +++ b/packages/dex-host-rs/src/turn.rs @@ -19,8 +19,8 @@ use std::path::Path; use dex_loop::{ - ApprovalMode, Budget, CancellationToken, Event, Exit, Lexicon, Log as _, Model, PrincipalId, - ThreadId, TurnId, rehydrate, + ApprovalMode, Budget, CancellationToken, ConfirmationDecision, Event, Exit, Lexicon, Log as _, + Model, PrincipalId, ThreadId, TurnId, rehydrate, }; use crate::{LocalEffects, LocalLog, LocalTools}; @@ -150,6 +150,9 @@ async fn drive_to_completion( call, principal: principal.clone(), text: UNATTENDED_ANSWER.to_owned(), + // An unattended assumption is text, never human action consent. + confirmation_decision: ConfirmationDecision::Unspecified, + args_digest: String::new(), }]) .await .map_err(|_fenced| { @@ -240,6 +243,117 @@ mod tests { } } + #[tokio::test] + async fn unattended_question_answer_never_grants_typed_action_consent() { + use dex_loop::{ActionConfirmation, CallId, ProposedCall}; + + let state_root = TempDir::new().expect("state tempdir"); + let workspace = TempDir::new().expect("workspace tempdir"); + let principal = PrincipalId::new("alice"); + let question_call = CallId::new("question-1"); + let action = ProposedCall::new( + CallId::new("send-1"), + ToolName::new("mail.send"), + serde_json::json!({"recipient":"someone@example.com"}), + principal.clone(), + ); + let question = ProposedCall::new( + question_call.clone(), + ToolName::new("person.ask"), + serde_json::json!({"text":"Send the message?"}), + principal.clone(), + ); + let log = LocalLog::acquire(state_root.path().join("log"), &thread()) + .await + .expect("acquire log"); + // Resume a persisted question from an earlier host. The current local + // two-tool catalog cannot ask yet; the driver still owns this exit. + log.append(&[ + Event::UserMessage { + turn: TurnId::new("t1"), + message_id: None, + principal: principal.clone(), + text: "Work unattended".into(), + attachments: vec![], + client_tools: vec![], + authorized_tools: vec![], + approval_mode: ApprovalMode::Headless, + }, + Event::StepStarted { + step: 1, + control_through: dex_loop::Cursor::START, + }, + Event::ModelStepCompleted { + step: 1, + text: String::new(), + calls: vec![question], + reasoning: None, + served: None, + timing: None, + }, + Event::Question { + call: question_call.clone(), + text: "Send the message?".into(), + confirmation: Some(ActionConfirmation { + proposal_call_id: action.id.clone(), + tool: action.tool.clone(), + args_digest: dex_loop::args_digest(&action.args), + principal_id: principal.clone(), + }), + }, + ]) + .await + .expect("persist parked question"); + let engine = dex_loop::Engine::new( + log.clone(), + ScriptedModel::new(vec![vec![Ok(ModelChunk::Text("done".into()))]]), + LocalTools::new(workspace.path()), + LocalEffects::open(state_root.path().join("effects.json")) + .await + .expect("open ledger"), + Lexicon::default(), + Budget::default(), + ); + assert_eq!( + drive_to_completion( + &engine, + &log, + &thread(), + &principal, + &CancellationToken::new() + ) + .await + .expect("resume unattended question"), + Exit::Done + ); + let events = read_log(&log).await.expect("read answered question"); + let answers: Vec<_> = events + .iter() + .filter_map(|(_, event)| match event { + Event::Answer { + call, + principal, + text, + confirmation_decision, + args_digest, + } => Some((call, principal, text, confirmation_decision, args_digest)), + _ => None, + }) + .collect(); + assert_eq!(answers.len(), 1); + let (call, actor, text, decision, digest) = answers[0]; + assert_eq!(call, &question_call); + assert_eq!(actor, &principal); + assert_eq!(text, UNATTENDED_ANSWER); + assert_eq!(*decision, ConfirmationDecision::Unspecified); + assert!(digest.is_empty()); + let replayed = rehydrate(thread(), &events); + let mut attempted_action = action; + attempted_action.args["confirmation"] = serde_json::json!(attempted_action.id.as_str()); + assert!(!replayed.confirmed_action(&attempted_action)); + assert!(matches!(events.last(), Some((_, Event::Final { text })) if text == "done")); + } + #[tokio::test] async fn drives_a_mutation_to_completion_with_no_approval() { let state_root = TempDir::new().expect("state tempdir"); diff --git a/vendor/dex-loop/src/compaction.rs b/vendor/dex-loop/src/compaction.rs index 6e66e9b41..c7b20e322 100644 --- a/vendor/dex-loop/src/compaction.rs +++ b/vendor/dex-loop/src/compaction.rs @@ -3,7 +3,7 @@ use std::future::Future; use crate::context::{Context, Entry, Message}; -use crate::event::Cursor; +use crate::event::{Cursor, Usage}; /// A planned compaction: history up to and including `covers_to` becomes /// `summary`. The engine appends it as `Event::Compaction`, so rehydration @@ -16,7 +16,20 @@ pub struct Compaction { /// Decides, before each model call, whether to compact. pub trait Compactor: Send + Sync { - fn plan(&self, ctx: &Context) -> impl Future> + Send; + fn plan(&self, ctx: &Context) -> impl Future + Send; +} + +/// An attempted summary may consume model usage even when no summary commits. +#[derive(Clone, Debug, Default)] +pub struct CompactionPlan { + pub compaction: Option, + pub usage: Usage, +} + +#[derive(Clone, Debug, Default)] +pub struct Summary { + pub text: Option, + pub usage: Usage, } /// Never compacts. @@ -24,62 +37,129 @@ pub trait Compactor: Send + Sync { pub struct NoCompaction; impl Compactor for NoCompaction { - async fn plan(&self, _ctx: &Context) -> Option { - None + async fn plan(&self, _ctx: &Context) -> CompactionPlan { + CompactionPlan::default() } } /// Writes the summary for a prefix of history (usually one model call). pub trait Summarize: Send + Sync { /// `None` skips this compaction; the engine continues uncompacted. - fn summarize(&self, entries: &[Entry]) -> impl Future> + Send; + fn summarize(&self, ctx: &Context, entries: &[Entry]) -> impl Future + Send; } -/// Compacts when history exceeds `max_bytes`, keeping at least the last -/// `keep_recent` entries verbatim. +/// Compacts when history exceeds `max_bytes`. `new` keeps recent complete +/// call/result pairs; `for_turns` covers all completed assistant steps while +/// Context retains the current request, steers and attachments verbatim. #[derive(Clone, Debug)] pub struct Threshold { max_bytes: usize, - keep_recent: usize, + cut_policy: CutPolicy, summarizer: S, } +#[derive(Clone, Copy, Debug)] +enum CutPolicy { + KeepRecent(usize), + CompletedTurnSteps, +} + impl Threshold { pub fn new(max_bytes: usize, keep_recent: usize, summarizer: S) -> Self { Self { max_bytes, - keep_recent, + cut_policy: CutPolicy::KeepRecent(keep_recent), + summarizer, + } + } + + /// Production turn compaction covers all completed assistant steps so + /// signed thinking never replays beside a prefix it was not produced with. + /// Current request, applied steers and attachments remain exact in Context. + pub fn for_turns(max_bytes: usize, summarizer: S) -> Self { + Self { + max_bytes, + cut_policy: CutPolicy::CompletedTurnSteps, summarizer, } } } -impl Compactor for Threshold { - async fn plan(&self, ctx: &Context) -> Option { +impl Threshold { + /// Plans at 70% of the limit after a finished turn; running turns are untouched. + pub async fn plan_after_turn(&self, ctx: &Context) -> CompactionPlan { + if ctx.turn_running() || self.max_bytes == usize::MAX { + return CompactionPlan::default(); + } + // Do not repeatedly summarize a large prior summary after tiny turns. + let new_bytes: usize = ctx + .history() + .iter() + .filter(|entry| !matches!(entry.message, Message::Summary { .. })) + .map(|entry| entry.message.size()) + .sum(); + if new_bytes < self.max_bytes / 10 { + return CompactionPlan::default(); + } + self.plan_at(ctx, self.max_bytes.saturating_mul(70) / 100) + .await + } + + async fn plan_at(&self, ctx: &Context, max_bytes: usize) -> CompactionPlan { + // An unresolved call/attempt retains its exact request, arguments and + // provider continuation. The engine also dispatches these before planning. + if ctx.open_step().is_some() || ctx.open_attempt().is_some() { + return CompactionPlan::default(); + } let history = ctx.history(); let size: usize = history.iter().map(|entry| entry.message.size()).sum(); - if size <= self.max_bytes { - return None; + if size <= max_bytes { + return CompactionPlan::default(); + } + let cut = match self.cut_policy { + CutPolicy::KeepRecent(keep_recent) => cut_point(ctx, keep_recent), + CutPolicy::CompletedTurnSteps => plan_cut(ctx, max_bytes).map(|cut| cut.covered.len()), + }; + let Some(cut) = cut else { + return CompactionPlan::default(); + }; + // Current user input (including steers) stays verbatim in Context; + // summarize only the historical/completed material being replaced. + let entries: Vec<_> = history[..cut].iter().filter(|entry| { + !matches!(&entry.message, Message::User { turn, .. } if ctx.turn_running() && Some(turn) == ctx.turn()) + }).cloned().collect(); + if entries.is_empty() { + return CompactionPlan::default(); + } + let summary = self.summarizer.summarize(ctx, &entries).await; + CompactionPlan { + compaction: summary.text.map(|summary| Compaction { + covers_to: history[cut - 1].cursor, + summary, + }), + usage: summary.usage, } - let cut = cut_point(history, self.keep_recent)?; - let summary = self.summarizer.summarize(&history[..cut]).await?; - Some(Compaction { - covers_to: history[cut - 1].cursor, - summary, - }) + } +} + +impl Compactor for Threshold { + async fn plan(&self, ctx: &Context) -> CompactionPlan { + self.plan_at(ctx, self.max_bytes).await } } /// Where to cut history for a compaction, chosen by [`plan_cut`]. #[derive(Clone, Copy, Debug, PartialEq)] -pub struct Cut<'a> { +struct Cut<'a> { /// The entries the summary replaces (a prefix of the history). - pub covered: &'a [Entry], + covered: &'a [Entry], /// The cursor the `Compaction` event names: the last covered entry's. - pub covers_to: Cursor, + #[cfg(test)] + covers_to: Cursor, /// Whether the cut fell inside the current turn (after a closed step) /// rather than at its start. - pub mid_turn: bool, + #[cfg(test)] + mid_turn: bool, } /// Picks a cut for a history over `max_bytes`, or `None` when nothing should @@ -93,14 +173,13 @@ pub struct Cut<'a> { /// 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. +/// was not produced beside. Pre-compaction continuation state belongs to +/// the old prefix. 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() { +fn plan_cut(ctx: &Context, max_bytes: usize) -> Option> { + if ctx.open_step().is_some() || ctx.open_attempt().is_some() { return None; } let history = ctx.history(); @@ -115,6 +194,7 @@ pub fn plan_cut(ctx: &Context, max_bytes: usize) -> Option> { if !valid_cut(history, cut) { return None; } + #[cfg(test)] let turn_start = ctx.turn().and_then(|turn| { history .iter() @@ -122,7 +202,9 @@ pub fn plan_cut(ctx: &Context, max_bytes: usize) -> Option> { }); Some(Cut { covered: &history[..cut], + #[cfg(test)] covers_to: history[cut - 1].cursor, + #[cfg(test)] mid_turn: turn_start.is_some_and(|start| start < cut), }) } @@ -152,7 +234,8 @@ fn valid_cut(history: &[Entry], cut: usize) -> bool { /// 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 /// between two cursors so the compaction event can name it. -fn cut_point(history: &[Entry], keep_recent: usize) -> Option { +fn cut_point(ctx: &Context, keep_recent: usize) -> Option { + let history = ctx.history(); let latest = history.len().checked_sub(keep_recent)?; (1..=latest).rev().find(|&cut| { let first_kept = history.get(cut); @@ -166,6 +249,425 @@ fn cut_point(history: &[Entry], keep_recent: usize) -> Option { #[cfg(test)] mod tests { + use super::*; + use crate::{ + ApprovalId, ApprovalMode, ArtifactRef, CallId, Event, Outcome, Output, PrincipalId, + ProposedCall, ThreadId, ToolName, TurnId, rehydrate, + }; + use serde_json::json; + + fn thread() -> ThreadId { + ThreadId { + org: "org".into(), + workspace: "ws".into(), + thread: "thread".into(), + } + } + fn user(turn: &str, text: &str) -> Event { + Event::UserMessage { + turn: TurnId::new(turn), + message_id: None, + principal: PrincipalId::new("alice"), + text: text.into(), + attachments: vec![ArtifactRef::new("doc@v1")], + client_tools: vec![], + authorized_tools: vec![], + approval_mode: ApprovalMode::Interactive, + } + } + fn append(events: &mut Vec<(Cursor, Event)>, event: Event) { + events.push((Cursor(events.len() as i64 + 1), event)); + } + fn completed(events: &mut Vec<(Cursor, Event)>, id: &str, step: u32) { + let control_through = Cursor(events.len() as i64); + append( + events, + Event::StepStarted { + step, + control_through, + }, + ); + append( + events, + Event::ModelStepCompleted { + timing: None, + served: None, + step, + text: "completed tool request".into(), + calls: vec![ProposedCall::new( + CallId::new(id), + ToolName::new("read"), + json!({"exact": "args"}), + PrincipalId::new("alice"), + )], + reasoning: None, + }, + ); + append( + events, + Event::ToolFinished { + call: CallId::new(id), + outcome: Outcome::Succeeded, + output: Output::Text("result".into()), + receipt: None, + }, + ); + } + fn inputs(ctx: &Context) -> Vec { + ctx.history() + .iter() + .filter(|entry| matches!(entry.message, Message::User { .. })) + .map(|entry| entry.message.clone()) + .collect() + } + struct Fake; + impl Summarize for Fake { + async fn summarize(&self, _ctx: &Context, entries: &[Entry]) -> Summary { + assert!(entries.iter().all(|entry| !matches!(&entry.message, Message::User { turn, .. } if turn.as_str() == "current"))); + Summary { + text: Some("untrusted old completed work".into()), + usage: Usage::default(), + } + } + } + + #[test] + fn repeated_compaction_retains_current_input_steers_and_suffix_replay() { + let mut events = vec![]; + append(&mut events, user("old", "old constraint")); + append( + &mut events, + Event::Final { + text: "old done".into(), + }, + ); + append(&mut events, user("current", "do not publish")); + append( + &mut events, + Event::Steer { + principal: PrincipalId::new("bob"), + text: "keep within $10".into(), + }, + ); + completed(&mut events, "one", 1); + let mut warm = rehydrate(thread(), &events); + let before = inputs(&warm); + for (id, step) in [("two", 2), ("three", 3)] { + let covers = events.last().expect("event").0; + let event = Event::Compaction { + covers_to_cursor: covers, + summary: format!("summary through {}", covers.0), + }; + append(&mut events, event.clone()); + warm.observe(events.last().expect("event").0, &event); + assert_eq!(inputs(&warm), before.iter().filter(|message| matches!(message, Message::User { turn, .. } if turn.as_str() == "current")).cloned().collect::>()); + assert!( + warm.history() + .windows(2) + .all(|pair| pair[0].cursor <= pair[1].cursor) + ); + assert_eq!(warm, rehydrate(thread(), &events)); + // Same suffix floor PgLog uses: current UserMessage precedes covers. + assert_eq!(warm, rehydrate(thread(), &events[2..])); + completed(&mut events, id, step); + warm = rehydrate(thread(), &events); + } + append(&mut events, user("next", "new task")); + append( + &mut events, + Event::Final { + text: "current done".into(), + }, + ); + let mut next = rehydrate(thread(), &events); + assert_eq!(next.turn(), Some(&TurnId::new("next"))); + let covers = events[events.len() - 3].0; + let event = Event::Compaction { + covers_to_cursor: covers, + summary: "all previous work".into(), + }; + append(&mut events, event.clone()); + next.observe(events.last().expect("event").0, &event); + assert_eq!(inputs(&next).len(), 1); + assert!(matches!(&inputs(&next)[0], Message::User { text, .. } if text == "new task")); + assert_eq!(next, rehydrate(thread(), &events)); + } + + #[test] + fn queued_turn_after_compaction_requires_full_replay_and_keeps_exact_turn_state() { + let mut events = vec![]; + append(&mut events, user("current", "do not publish")); + let mut warm = rehydrate(thread(), &events); + { + let mut push = |event: Event| { + append(&mut events, event.clone()); + warm.observe(events.last().expect("event").0, &event); + }; + push(Event::ModelStepCompleted { + timing: None, + served: None, + step: 1, + text: "completed work".into(), + calls: vec![], + reasoning: None, + }); + push(Event::Compaction { + covers_to_cursor: Cursor(2), + summary: "historical work".into(), + }); + push(user("queued-one", "first queued request")); + push(user("queued-two", "second queued request")); + push(Event::Final { + text: "current done".into(), + }); + push(Event::Compaction { + covers_to_cursor: Cursor(2), + summary: "earlier inputs and work".into(), + }); + } + assert_eq!(warm.turn(), Some(&TurnId::new("queued-one"))); + assert_eq!(warm.status(), crate::context::Status::Running); + assert_eq!(warm, rehydrate(thread(), &events)); + // The retired suffix starts at covers=2 and omits the current + // UserMessage. Its Final then wrongly terminates queued-one. + assert_ne!(warm, rehydrate(thread(), &events[1..])); + let finish = Event::Final { + text: "first queued done".into(), + }; + append(&mut events, finish.clone()); + warm.observe(events.last().expect("event").0, &finish); + assert_eq!(warm.turn(), Some(&TurnId::new("queued-two"))); + assert_eq!(warm.status(), crate::context::Status::Running); + assert_eq!(warm, rehydrate(thread(), &events)); + } + + #[tokio::test] + async fn cut_keeps_complete_pairs_and_parked_approval_state_unchanged() { + let mut events = vec![]; + append(&mut events, user("current", "original request")); + completed(&mut events, "one", 1); + completed(&mut events, "two", 2); + let mut ctx = rehydrate(thread(), &events); + let compactor = Threshold::new(1, 1, Fake); + let plan = compactor.plan(&ctx).await; + let compaction = plan.compaction.expect("completed prefix compacts"); + // Keeping one result retains its assistant request too. + assert_eq!(compaction.covers_to, Cursor(4)); + let event = Event::Compaction { + covers_to_cursor: compaction.covers_to, + summary: compaction.summary, + }; + append(&mut events, event.clone()); + ctx.observe(events.last().expect("event").0, &event); + assert!( + matches!(&ctx.history()[2].message, Message::Assistant { calls, .. } if calls[0].id.as_str() == "two") + ); + assert!( + matches!(&ctx.history()[3].message, Message::Tool { call, .. } if call.as_str() == "two") + ); + let call = ProposedCall::new( + CallId::new("write"), + ToolName::new("write"), + json!({"unchanged": true}), + PrincipalId::new("alice"), + ); + append( + &mut events, + Event::ModelStepCompleted { + timing: None, + served: None, + step: 3, + text: String::new(), + calls: vec![call.clone()], + reasoning: None, + }, + ); + append( + &mut events, + Event::ApprovalRequested { + call: call.id.clone(), + approval: ApprovalId::new("approval"), + args_digest: call.args_digest.clone(), + summary: "approve".into(), + }, + ); + ctx = rehydrate(thread(), &events); + let before = ctx.clone(); + assert!(compactor.plan(&ctx).await.compaction.is_none()); + assert_eq!(ctx, before); + assert_eq!(ctx, rehydrate(thread(), &events)); + } + #[tokio::test] + async fn boundary_summary_covers_finished_inputs_preserves_usage_and_replays() { + struct FinishedSummary; + impl Summarize for FinishedSummary { + async fn summarize(&self, _ctx: &Context, entries: &[Entry]) -> Summary { + assert!(entries.iter().any(|entry| matches!(&entry.message, + Message::User { text, .. } if text == "exact completed request"))); + Summary { + text: Some("finished history".into()), + usage: Usage { + cache_creation_input_tokens: 0, + cache_read_input_tokens: 0, + input_tokens: 4, + output_tokens: 2, + cost_micros: 3, + }, + } + } + } + let mut events = vec![]; + append(&mut events, user("finished", "exact completed request")); + completed(&mut events, "one", 1); + append( + &mut events, + Event::Final { + text: "done".into(), + }, + ); + let mut ctx = rehydrate(thread(), &events); + let size = ctx + .history() + .iter() + .map(|entry| entry.message.size()) + .sum::(); + let compactor = Threshold::for_turns(size + 1, FinishedSummary); + assert!(compactor.plan(&ctx).await.compaction.is_none()); + let plan = compactor.plan_after_turn(&ctx).await; + assert_eq!(plan.usage.cost_micros, 3); + let plan = plan.compaction.expect("proactive boundary"); + append( + &mut events, + Event::Usage(Usage { + cache_creation_input_tokens: 0, + cache_read_input_tokens: 0, + input_tokens: 4, + output_tokens: 2, + cost_micros: 3, + }), + ); + append( + &mut events, + Event::Compaction { + covers_to_cursor: plan.covers_to, + summary: plan.summary, + }, + ); + ctx.observe(events[events.len() - 2].0, &events[events.len() - 2].1); + ctx.observe(events[events.len() - 1].0, &events[events.len() - 1].1); + assert_eq!(ctx, rehydrate(thread(), &events)); + assert_eq!( + ctx.history().len(), + 1, + "finished request is included in the summary" + ); + append(&mut events, user("next", "next request")); + let ctx = rehydrate(thread(), &events); + assert!( + compactor.plan_after_turn(&ctx).await.compaction.is_none(), + "queued work is never compacted between turns" + ); + assert!(compactor.plan(&ctx).await.compaction.is_none()); + } + + struct RecordedSummary(std::sync::Arc>>); + impl Summarize for RecordedSummary { + async fn summarize(&self, ctx: &Context, entries: &[Entry]) -> Summary { + assert!(entries.iter().all(|entry| !matches!(&entry.message, + Message::User { turn, .. } if Some(turn) == ctx.turn()))); + *self.0.lock().expect("seen") = entries.to_vec(); + Summary { + text: Some("completed non-authoritative history".into()), + usage: Usage { + cache_creation_input_tokens: 0, + cache_read_input_tokens: 0, + input_tokens: 11, + output_tokens: 3, + cost_micros: 7, + }, + } + } + } + + #[tokio::test] + async fn configured_turn_compactor_activates_mid_turn_and_preserves_exact_inputs_and_usage() { + let mut events = vec![]; + append(&mut events, user("current", "exact current request")); + completed(&mut events, "one", 1); + append( + &mut events, + Event::Steer { + principal: PrincipalId::new("bob"), + text: "exact applied steer".into(), + }, + ); + completed(&mut events, "two", 2); + // Control accepted during planning must survive alongside applied inputs. + append( + &mut events, + Event::Steer { + principal: PrincipalId::new("carol"), + text: "exact queued steer".into(), + }, + ); + let mut ctx = rehydrate(thread(), &events); + assert!(ctx.has_queued_steers()); + let before_inputs = inputs(&ctx); + assert_eq!(before_inputs.len(), 2, "request and applied steer"); + let expected_cursor = ctx.history().last().expect("last completed result").cursor; + let seen = std::sync::Arc::new(std::sync::Mutex::new(vec![])); + // This is the constructor Wired uses, rather than an isolated helper test. + let compactor = Threshold::for_turns(1, RecordedSummary(seen.clone())); + let plan = compactor.plan(&ctx).await; + assert_eq!( + plan.usage.cost_micros, 7, + "summary usage reaches engine admission" + ); + let compaction = plan.compaction.expect("configured mid-turn compaction"); + assert_eq!( + compaction.covers_to, expected_cursor, + "latest completed step is covered" + ); + let covered = seen.lock().expect("seen"); + for id in ["one", "two"] { + assert!(covered.iter().any( + |entry| matches!(&entry.message, Message::Assistant { calls, .. } + if calls.iter().any(|call| call.id.as_str() == id)) + )); + assert!(covered.iter().any( + |entry| matches!(&entry.message, Message::Tool { call, .. } if call.as_str() == id) + )); + } + drop(covered); + let event = Event::Compaction { + covers_to_cursor: compaction.covers_to, + summary: compaction.summary, + }; + append(&mut events, event.clone()); + ctx.observe(events.last().expect("compaction cursor").0, &event); + assert_eq!( + inputs(&ctx), + before_inputs, + "principal, text and attachment identity remain exact" + ); + assert!(ctx.has_queued_steers(), "queued control remains pending"); + assert!( + ctx.history().iter().all(|entry| matches!( + entry.message, + Message::User { .. } | Message::Summary { .. } + )), + "no assistant/tool suffix is replayed under the changed prefix" + ); + assert_eq!( + ctx, + rehydrate(thread(), &events), + "warm and full replay agree" + ); + } +} + +#[cfg(test)] +mod cut_tests { use super::*; use crate::event::{ ApprovalMode, CallId, Event, Outcome, Output, PrincipalId, ProposedCall, ThreadId, @@ -355,16 +857,11 @@ mod tests { #[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); + let history = [Entry { + cursor: Cursor(1), + message: Message::Summary { text: "s".into() }, + }]; + assert!(!valid_cut(&history, 1)); } #[test] diff --git a/vendor/dex-loop/src/context.rs b/vendor/dex-loop/src/context.rs index eb0b9ca37..ef4eb6992 100644 --- a/vendor/dex-loop/src/context.rs +++ b/vendor/dex-loop/src/context.rs @@ -6,13 +6,17 @@ //! the same context, whether the actor stayed warm or restarted. use crate::event::{ - ApprovalId, ApprovalMode, ArtifactRef, CallId, ClientToolSpec, Cursor, Event, MessageId, - Outcome, Output, PrincipalId, ProposedCall, ProviderReasoning, ServedBy, ThreadId, ToolName, - ToolResult, TurnId, Usage, + ActionConfirmation, ApprovalId, ApprovalMode, ArtifactRef, CallId, ClientToolSpec, + ConfirmationDecision, Cursor, Event, HEADLESS_AUTO_APPROVER, MessageId, Outcome, Output, + PrincipalId, ProposedCall, ProviderReasoning, ServedBy, ThreadId, ToolName, ToolResult, TurnId, + Usage, }; const NOT_RUN_NEW_TURN: &str = "not run: a new turn started first"; const UNKNOWN_NEW_TURN: &str = "outcome unknown: a new turn started before the result was recorded"; +const MAX_ACTION_RECORDS: usize = 128; +// A public question permits 4096 characters, each at most four UTF-8 bytes. +const MAX_ACTION_ARGUMENT_BYTES: usize = 16 * 1024; /// One message in the model's view of the thread. #[derive(Clone, Debug, PartialEq)] @@ -48,8 +52,16 @@ impl Message { pub fn size(&self) -> usize { match self { Message::User { text, .. } | Message::Summary { text } => text.len(), - Message::Assistant { text, calls, .. } => { + Message::Assistant { + text, + calls, + reasoning, + .. + } => { text.len() + + reasoning + .as_ref() + .map_or(0, |state| state.payload.to_string().len()) + calls .iter() .map(|call| call.tool.as_str().len() + call.args.to_string().len()) @@ -57,7 +69,10 @@ impl Message { } Message::Tool { output, .. } => match output { Output::Text(text) => text.len(), - Output::Ref(reference) => reference.as_str().len(), + // A reference may render a preview as well as metadata. Dex's + // host caps a preview at 16 KiB; count that conservative + // allowance so tool-heavy histories compact before rendering. + Output::Ref(reference) => reference.as_str().len().saturating_add(16 * 1024), }, } } @@ -121,6 +136,21 @@ pub(crate) struct OpenStep { pub(crate) states: Vec, } +/// Maximum document references retained across accepted messages. The current +/// message remains exact so the document owner can reject invalid admissions. +pub const MAX_CONTEXT_ATTACHMENTS: usize = 20; + +/// Typed provenance derived only from accepted user messages, independent of +/// summaries. These coordinates select evidence; the document owner still +/// verifies message admission, principal access and immutable versions. +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct AttachmentInput { + pub turn: TurnId, + pub message_id: Option, + pub principal: PrincipalId, + pub attachments: Vec, +} + /// Everything the engine knows about a thread. #[derive(Clone, Debug, PartialEq)] pub struct Context { @@ -129,6 +159,9 @@ pub struct Context { acting: Option, status: Status, history: Vec, + // Newest input first (even when it has no uploads), then bounded earlier + // attachment batches. Compaction never manufactures or edits provenance. + attachment_inputs: Vec, step: u32, usage: Usage, cursor: Cursor, @@ -151,6 +184,10 @@ 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, + // Derived only from exact Question/Answer rows, retained outside summaries. + action_confirmations: Vec<(CallId, ActionConfirmation, ConfirmationDecision, bool)>, + // Exact owner policy previews survive model-history cuts. Expired refs fail closed. + action_previews: 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 @@ -186,6 +223,7 @@ impl Context { acting: None, status: Status::Idle, history: Vec::new(), + attachment_inputs: Vec::new(), step: 0, usage: Usage::default(), cursor: Cursor::START, @@ -200,12 +238,19 @@ impl Context { approval_mode: ApprovalMode::Interactive, authorized_principal: None, uncertain_calls: Vec::new(), + action_confirmations: Vec::new(), + action_previews: Vec::new(), pending_turns: Vec::new(), pre_started: Vec::new(), last_compaction: None, } } + /// Assistant entries created before this event cannot reuse signed thinking. + pub fn last_compaction_cursor(&self) -> Option { + self.last_compaction + } + pub fn thread(&self) -> &ThreadId { &self.thread } @@ -255,13 +300,78 @@ 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 + pub fn proposed_call(&self, id: &CallId) -> Option<&ProposedCall> { + self.action_preview(id).or_else(|| { + self.history + .iter() + .rev() + .filter_map(|entry| match &entry.message { + Message::Assistant { calls, .. } => Some(calls), + _ => None, + }) + .flatten() + .find(|call| &call.id == id) + }) + } + + /// Render owner-produced preview details, rather than model-authored consent prose. + pub fn action_question_text(&self, proposal: &CallId, action: &str) -> Option { + let mut arguments = self.action_preview(proposal)?.args.clone(); + if let Some(arguments) = arguments.as_object_mut() { + arguments.remove("confirmation"); + } + let details = serde_json::to_string_pretty(&arguments).ok()?; + let text = format!("Confirm this action?\n{action}\n\n{details}"); + // The public question contract must show the complete preview. + (text.chars().count() <= 4096).then_some(text) + } + + pub fn is_unexecuted_preview(&self, call: &CallId) -> bool { + self.action_preview(call).is_some() + } + + /// An exact failed policy preview before dispatch, retained separately from summaries. + pub fn action_preview(&self, call: &CallId) -> Option<&ProposedCall> { + self.action_previews + .iter() + .find(|proposal| &proposal.id == call) + } + + pub fn confirmation_question_exists(&self, proposal: &CallId) -> bool { + self.action_confirmations + .iter() + .any(|(_, binding, _, _)| &binding.proposal_call_id == proposal) + } + + /// A typed affirmative choice for this exact action, unused by any dispatch. + pub fn confirmed_action(&self, call: &ProposedCall) -> bool { + let Some(id) = call + .args + .get("confirmation") + .and_then(serde_json::Value::as_str) + else { + return false; + }; + let mut args = call.args.clone(); + if let Some(args) = args.as_object_mut() { + args.remove("confirmation"); + } + let digest = crate::args_digest(&args); + self.action_confirmations + .iter() + .any(|(_, binding, decision, used)| { + !used + && *decision == ConfirmationDecision::Confirm + && binding.proposal_call_id.as_str() == id + && binding.tool == call.tool + && binding.principal_id == call.principal + && binding.args_digest == digest + }) + } + + /// Newest accepted/applied input and bounded older attachment provenance. + pub fn attachment_inputs(&self) -> &[AttachmentInput] { + &self.attachment_inputs } /// Model calls started in the current turn. @@ -420,12 +530,17 @@ impl Context { approval, args_digest, approved, - .. + principal, } => { + // Older engines wrote synthetic approvals in headless mode. + // Retain those rows as history, but they cannot authorize an + // unstarted parked call after an upgrade. Later ToolStarted / + // ToolFinished rows still reconstruct effects already run. if let Some(CallState::Parked { approval: parked, decision, }) = self.state_mut(call) + && principal.as_str() != HEADLESS_AUTO_APPROVER && parked == approval && decision.is_none() { @@ -435,7 +550,27 @@ impl Context { }); } } - Event::Answer { call, text, .. } => { + Event::Answer { + call, + principal, + text, + confirmation_decision, + args_digest, + } => { + for (question, binding, decision, used) in &mut self.action_confirmations { + if question == call + && &binding.principal_id == principal + && binding.args_digest == *args_digest + && *decision == ConfirmationDecision::Unspecified + && !*used + { + *decision = *confirmation_decision; + // A normal answer closes the question without granting consent. + if *confirmation_decision == ConfirmationDecision::Unspecified { + *used = true; + } + } + } if let Some(CallState::Asked { answer }) = self.state_mut(call) && answer.is_none() { @@ -512,6 +647,19 @@ impl Context { self.pre_started.clear(); } Event::ToolStarted { call, .. } => { + if let Some(proposal) = self.proposed_call(call).cloned() + && self.confirmed_action(&proposal) + { + let id = proposal + .args + .get("confirmation") + .and_then(serde_json::Value::as_str); + for (_, binding, _, used) in &mut self.action_confirmations { + if Some(binding.proposal_call_id.as_str()) == id { + *used = true; + } + } + } if let Some(state) = self.state_mut(call) { if !matches!(state, CallState::Done(_)) { *state = CallState::Started; @@ -534,6 +682,50 @@ impl Context { output, receipt, } => { + // A confirmation preview is an owner policy refusal before the + // dispatch boundary. External error prose cannot manufacture one. + let unstarted = self.open_step.as_ref().and_then(|step| { + step.calls + .iter() + .zip(&step.states) + .find(|(proposal, state)| { + &proposal.id == call && matches!(state, CallState::Todo) + }) + .map(|(proposal, _)| proposal.clone()) + }); + if let Some(proposal) = unstarted + && *outcome == Outcome::Failed + && let Output::Text(text) = output + && let Ok(document) = serde_json::from_str::(text) + && document.get("status").and_then(serde_json::Value::as_str) + == Some("needs_confirmation") + { + let mut args = proposal.args.clone(); + if let Some(args) = args.as_object_mut() { + args.remove("confirmation"); + } + if document + .get("args_digest") + .and_then(serde_json::Value::as_str) + == Some(crate::args_digest(&args).as_str()) + && serde_json::to_string_pretty(&args) + .is_ok_and(|text| text.len() <= MAX_ACTION_ARGUMENT_BYTES) + && self.action_preview(call).is_none() + { + self.action_previews.push(ProposedCall::new( + proposal.id, + proposal.tool, + args, + proposal.principal, + )); + if self.action_previews.len() > MAX_ACTION_RECORDS { + let expired = self.action_previews.remove(0); + self.action_confirmations.retain(|(_, binding, _, _)| { + binding.proposal_call_id != expired.id + }); + } + } + } self.uncertain_calls.retain(|prior| &prior.id != call); // A pre-committed read finished by an abandoned attempt. self.pre_started.retain(|started| started != call); @@ -580,7 +772,27 @@ impl Context { }; } } - Event::Question { call, .. } => { + Event::Question { + call, confirmation, .. + } => { + if let Some(binding) = confirmation + && !self + .action_confirmations + .iter() + .any(|(question, _, _, _)| question == call) + { + self.action_confirmations.push(( + call.clone(), + binding.clone(), + ConfirmationDecision::Unspecified, + false, + )); + if self.action_confirmations.len() > MAX_ACTION_RECORDS { + let (_, expired, _, _) = self.action_confirmations.remove(0); + self.action_previews + .retain(|proposal| proposal.id != expired.proposal_call_id); + } + } if let Some(state) = self.state_mut(call) { *state = CallState::Asked { answer: None }; } @@ -600,10 +812,19 @@ impl Context { summary, } => { self.last_compaction = Some(cursor); - self.history - .retain(|entry| entry.cursor > *covers_to_cursor); + // Input and applied steers of the current turn retain their + // exact text, principal, and attachments across compaction. + // Covered inputs precede the new summary, keeping history + // cursors non-decreasing even after repeated compaction. + self.history.retain(|entry| { + entry.cursor > *covers_to_cursor + || matches!(&entry.message, Message::User { turn, .. } if self.status == Status::Running && Some(turn) == self.turn.as_ref()) + }); + let position = self + .history + .partition_point(|entry| entry.cursor <= *covers_to_cursor); self.history.insert( - 0, + position, Entry { cursor: *covers_to_cursor, message: Message::Summary { @@ -688,6 +909,38 @@ impl Context { } fn push(&mut self, cursor: Cursor, message: Message) { + if let Message::User { + turn, + message_id, + principal, + attachments, + .. + } = &message + { + let current = AttachmentInput { + turn: turn.clone(), + message_id: message_id.clone(), + principal: principal.clone(), + attachments: attachments.clone(), + }; + let mut seen: std::collections::BTreeSet<_> = attachments.iter().cloned().collect(); + let mut budget = MAX_CONTEXT_ATTACHMENTS.saturating_sub(attachments.len()); + let mut retained = vec![current]; + for mut previous in std::mem::take(&mut self.attachment_inputs) { + if budget == 0 { + break; + } + previous + .attachments + .retain(|reference| seen.insert(reference.clone())); + previous.attachments.truncate(budget); + if !previous.attachments.is_empty() { + budget -= previous.attachments.len(); + retained.push(previous); + } + } + self.attachment_inputs = retained; + } self.history.push(Entry { cursor, message }); } diff --git a/vendor/dex-loop/src/engine.rs b/vendor/dex-loop/src/engine.rs index 485d10838..c1c157281 100644 --- a/vendor/dex-loop/src/engine.rs +++ b/vendor/dex-loop/src/engine.rs @@ -448,15 +448,47 @@ where .await?; return Ok(Exit::Failed); } - if let Some(plan) = self.compactor.plan(ctx).await { + let remaining_wall = self.budget.wall.saturating_sub(started.elapsed()); + let plan = tokio::select! { + _ = cancel.cancelled() => return self.interrupt(ctx, &prefetch).await, + result = tokio::time::timeout(remaining_wall, self.compactor.plan(ctx)) => result, + }; + if let Ok(plan) = plan { + let mut events = Vec::new(); + if plan.usage != Default::default() { + events.push(Event::Usage(plan.usage)); + } + if let Some(compaction) = plan.compaction { + events.push(Event::Compaction { + covers_to_cursor: compaction.covers_to, + summary: compaction.summary, + }); + } + if !events.is_empty() { + self.emit(ctx, events).await?; + } + } + // Summarization is an inference call charged to this same turn. + // Control changes during it must also win before an ordinary call. + self.read_control(ctx).await?; + if ctx.interrupt_requested() || cancel.is_cancelled() { + return self.interrupt(ctx, &prefetch).await; + } + if let Some(axis) = self + .budget + .exhausted(ctx.step(), ctx.usage(), started.elapsed()) + { + let message = self.budget_message(ctx, axis); self.emit( ctx, - vec![Event::Compaction { - covers_to_cursor: plan.covers_to, - summary: plan.summary, + vec![Event::Error { + class: Some(crate::ErrorClass::BudgetExhausted), + code: ErrorCode::BudgetExhausted, + message, }], ) .await?; + return Ok(Exit::Failed); } // A step never inherits another step's reads. prefetch = Prefetch::new(); @@ -675,6 +707,8 @@ where } self.log.append_text(CUT_OFF_NOTICE.to_owned()).await?; text.push_str(CUT_OFF_NOTICE); + self.read_control(ctx).await?; + let continues = ctx.has_queued_steers() && !ctx.interrupt_requested(); let mut events = pending_usage; events.push(Event::ModelStepCompleted { step, @@ -684,9 +718,11 @@ where served: served.clone(), timing: timing.take(), }); - events.push(Event::Final { text }); + if !continues { + events.push(Event::Final { text }); + } self.emit(ctx, events).await?; - return Ok(Some(Exit::Done)); + return Ok((!continues).then_some(Exit::Done)); } if let Some(error) = failure { let mut events = pending_usage; @@ -1104,7 +1140,41 @@ where ctx, vec![Event::Question { call: call.id.clone(), - text: question.to_owned(), + text: call + .args + .get("confirmation") + .and_then(|value| value.get("proposal_call_id")) + .and_then(serde_json::Value::as_str) + .and_then(|id| { + ctx.action_preview(&CallId::new(id)).and_then(|proposal| { + self.tools + .catalog() + .iter() + .find(|spec| spec.name == proposal.tool) + .and_then(|spec| { + ctx.action_question_text(&proposal.id, &spec.label) + }) + }) + }) + .unwrap_or_else(|| question.to_owned()), + confirmation: call + .args + .get("confirmation") + .and_then(|value| value.get("proposal_call_id")) + .and_then(serde_json::Value::as_str) + .and_then(|id| ctx.action_preview(&crate::CallId::new(id))) + .map(|proposal| { + let mut args = proposal.args.clone(); + if let Some(args) = args.as_object_mut() { + args.remove("confirmation"); + } + crate::ActionConfirmation { + proposal_call_id: proposal.id.clone(), + tool: proposal.tool.clone(), + args_digest: crate::args_digest(&args), + principal_id: proposal.principal.clone(), + } + }), }], ) .await?; @@ -1187,9 +1257,9 @@ where self.emit( ctx, vec![Event::Error { + class: Some(crate::ErrorClass::BudgetExhausted), code: ErrorCode::BudgetExhausted, message, - class: Some(crate::ErrorClass::BudgetExhausted), }], ) .await?; @@ -1437,6 +1507,14 @@ where self.finish(ctx, call, result).await } Claim::Granted => { + // A durable Started row proves the dispatch boundary was crossed. + // A missing ledger is lost evidence, never a fresh permission to run. + if already_started { + let result = ToolResult::unknown(UNKNOWN_NO_RETRY); + self.effects.record(&call.id, &result).await?; + return self.finish(ctx, call, result).await; + } + // Claim persistence can consume the remaining wall. A fresh // claim is known not executed; a historical Started call may // already have taken effect and must remain Unknown. diff --git a/vendor/dex-loop/src/event.rs b/vendor/dex-loop/src/event.rs index 195dbddbe..b5452f069 100644 --- a/vendor/dex-loop/src/event.rs +++ b/vendor/dex-loop/src/event.rs @@ -83,6 +83,24 @@ string_id!( ArtifactRef ); +/// Owner-bound action details attached to a chat question. No prose grants authority. +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct ActionConfirmation { + pub proposal_call_id: CallId, + pub tool: ToolName, + pub args_digest: String, + pub principal_id: PrincipalId, +} + +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum ConfirmationDecision { + #[default] + Unspecified, + Confirm, + Decline, +} + /// The tenant scope of one thread. Every log read and write carries all three. #[derive(Clone, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)] pub struct ThreadId { @@ -371,6 +389,9 @@ impl ErrorCode { } } +/// Execution mode retained in durable ingress for compatibility and attribution. +/// Both modes now grant `NeedsApproval` through an exact `AutoApproved` +/// receipt. Neither mode overrides a hard policy `Deny`. /// A content-free failure class supplied by the typed model adapter. #[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] @@ -474,10 +495,10 @@ impl ErrorClass { #[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum ApprovalMode { - /// A human decides: the turn parks on `ApprovalRequested`. + /// Interactive caller; approval-class calls receive a durable auto-approval. #[default] Interactive, - /// No human: policy approves what would otherwise ask. + /// Unattended caller; approval-class calls receive the same durable receipt. Headless, } @@ -562,6 +583,10 @@ pub enum Event { call: CallId, principal: PrincipalId, text: String, + #[serde(default)] + confirmation_decision: ConfirmationDecision, + #[serde(default)] + args_digest: String, }, /// Control: a client session's outcome for one `ClientToolRequested` /// call. Only `Outcome::Succeeded` or `Outcome::Failed` are accepted @@ -686,6 +711,8 @@ pub enum Event { Question { call: CallId, text: String, + #[serde(default)] + confirmation: Option, }, /// The turn is parked until a `ClientToolResult` for this call arrives. /// The only event that carries tool arguments to a surface; the host diff --git a/vendor/dex-loop/src/lib.rs b/vendor/dex-loop/src/lib.rs index 06fda88e6..53a7bf8b4 100644 --- a/vendor/dex-loop/src/lib.rs +++ b/vendor/dex-loop/src/lib.rs @@ -29,14 +29,17 @@ mod rehydrate; mod sanitize; pub use budget::{Budget, BudgetAxis}; -pub use compaction::{Compaction, Compactor, Cut, NoCompaction, Summarize, Threshold, plan_cut}; -pub use context::{Context, Entry, Message}; +pub use compaction::{ + Compaction, CompactionPlan, Compactor, NoCompaction, Summarize, Summary, Threshold, +}; +pub use context::{AttachmentInput, Context, Entry, MAX_CONTEXT_ATTACHMENTS, Message}; pub use engine::{CUT_OFF_NOTICE, DEFAULT_TOOL_CALL_DEADLINE, Engine, Exit, TOOLS_SEARCH}; pub use event::{ - AUTO_APPROVER, ApprovalId, ApprovalMode, ArtifactRef, AttemptNext, CallId, ClientToolSpec, - Cursor, ErrorClass, ErrorCode, Event, HEADLESS_AUTO_APPROVER, MessageId, Outcome, Output, - OutputRef, PrincipalId, ProposedCall, ProviderReasoning, ReceiptId, ServedBy, StepTiming, - ThreadId, ToolName, ToolResult, TurnId, Usage, args_digest, + AUTO_APPROVER, ActionConfirmation, ApprovalId, ApprovalMode, ArtifactRef, AttemptNext, CallId, + ClientToolSpec, ConfirmationDecision, Cursor, ErrorClass, ErrorCode, Event, + HEADLESS_AUTO_APPROVER, MessageId, Outcome, Output, OutputRef, PrincipalId, ProposedCall, + ProviderReasoning, ReceiptId, ServedBy, StepTiming, 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 9599ee377..34cbbaf8d 100644 --- a/vendor/dex-loop/src/ports.rs +++ b/vendor/dex-loop/src/ports.rs @@ -63,8 +63,7 @@ pub enum ModelChunk { }, Usage(Usage), /// The step's opaque provider continuation state. At most one per step, - /// sent only after a clean terminal (the same commit rule as `ToolCall` - /// and `Usage`), after the last `ToolCall` and before `Usage`. The engine + /// sent only after a clean terminal (the same commit rule as `Usage`), after the last `ToolCall` and before `Usage`. The engine /// stores it on `ModelStepCompleted`; history returns it on /// `Message::Assistant`. Reasoning(ProviderReasoning), diff --git a/vendor/dex-loop/tests/action_confirmation.rs b/vendor/dex-loop/tests/action_confirmation.rs new file mode 100644 index 000000000..a2f624577 --- /dev/null +++ b/vendor/dex-loop/tests/action_confirmation.rs @@ -0,0 +1,405 @@ +//! Consent is log-derived, one use, exact-action bound, and independent of summaries. +#[allow(dead_code)] +mod support; +use dex_loop::{ + ActionConfirmation, ApprovalMode, Budget, CallId, CancellationToken, ConfirmationDecision, + Context, Cursor, Engine, Event, Exit, Lexicon, Outcome, Output, PrincipalId, ProposedCall, + ThreadId, ToolName, ToolResult, ToolSpec, Tools, TurnId, Verdict, args_digest, rehydrate, +}; +use serde_json::json; + +fn thread() -> ThreadId { + ThreadId { + org: "org".into(), + workspace: "ws".into(), + thread: "thread".into(), + } +} + +#[derive(Clone)] +struct ConsentTools(support::FakeTools); + +impl Tools for ConsentTools { + fn catalog(&self) -> &[ToolSpec] { + self.0.catalog() + } + async fn search(&self, principal: &PrincipalId, query: &str) -> Vec { + self.0.search(principal, query).await + } + async fn policy(&self, ctx: &Context, call: &ProposedCall) -> Verdict { + if call.tool.as_str() == "user.ask" || ctx.confirmed_action(call) { + return Verdict::Allow; + } + let mut args = call.args.clone(); + args.as_object_mut().unwrap().remove("confirmation"); + Verdict::NeedsConfirmation { + preview: json!({ + "status":"needs_confirmation", "args_digest":args_digest(&args), + "action":"misleading provider label" + }) + .to_string(), + } + } + async fn run( + &self, + thread: &ThreadId, + call: &ProposedCall, + cancel: &CancellationToken, + ) -> ToolResult { + self.0.run(thread, call, cancel).await + } +} + +#[tokio::test] +async fn engine_emits_exact_question_and_only_dispatches_an_affirmative_typed_choice() { + engine_confirmation(false).await; +} + +#[tokio::test] +async fn engine_confirmation_survives_compaction_between_preview_and_question() { + engine_confirmation(true).await; +} + +async fn engine_confirmation(compact: bool) { + for decision in [ConfirmationDecision::Confirm, ConfirmationDecision::Decline] { + let log = support::FakeLog::default(); + let args = json!({"to":"recipient-1","body":"exact message"}); + let mut retry_args = args.clone(); + retry_args["confirmation"] = json!("consent-1-0"); + let model = support::FakeModel::new(vec![ + vec![support::call("send", args.clone())], + vec![support::call( + "user.ask", + json!({ + "question":"misleading model question", + "confirmation":{"proposal_call_id":"consent-1-0"} + }), + )], + vec![support::call("send", retry_args.clone())], + vec![support::text("Finished")], + ]); + let mut send = support::write_tool("send"); + send.label = "Send a message".into(); + let inner = support::FakeTools::new(vec![send, support::ask_tool("user.ask")]); + let engine = Engine::new( + log.clone(), + model, + ConsentTools(inner.clone()), + support::FakeEffects::default(), + Lexicon::default(), + Budget::default(), + ) + .with_compactor(dex_loop::Threshold::for_turns( + if compact { 1 } else { usize::MAX }, + support::FakeSummarizer, + )); + let mut ctx = log.start_turn("consent", "Prepare the exact message"); + let cancel = CancellationToken::new(); + let question = CallId::new("consent-2-0"); + assert_eq!( + engine.run(&mut ctx, &cancel).await, + Ok(Exit::Asked(question.clone())) + ); + assert!( + inner.runs().is_empty(), + "the preview and question never execute the mutation" + ); + if compact { + assert!( + log.events() + .iter() + .any(|event| matches!(event, Event::Compaction { .. })), + "the real compactor must cut the completed preview step before user.ask" + ); + assert!(!ctx.history().iter().any(|entry| matches!(&entry.message, + dex_loop::Message::Assistant { calls, .. } if calls.iter().any(|call| call.id.as_str() == "consent-1-0"))), + "the proposed action is absent from model history"); + } + let emitted = log + .events() + .into_iter() + .find_map(|event| match event { + Event::Question { + call, + text, + confirmation, + } => Some((call, text, confirmation)), + _ => None, + }) + .expect("engine question"); + assert_eq!(emitted.0, question); + let binding = emitted.2.expect("exact action tuple"); + assert_eq!( + binding, + ActionConfirmation { + proposal_call_id: CallId::new("consent-1-0"), + tool: ToolName::new("send"), + args_digest: args_digest(&args), + principal_id: support::alice(), + } + ); + assert!( + emitted.1.contains("Send a message") + && emitted.1.contains("recipient-1") + && emitted.1.contains("exact message") + ); + assert!( + !emitted.1.contains("misleading"), + "question details come from catalog and exact arguments" + ); + log.host_append(Event::Answer { + call: question, + principal: support::alice(), + text: "typed choice".into(), + confirmation_decision: decision, + args_digest: binding.args_digest, + }); + // Resume from the durable Question/Answer log, as a new actor would. + ctx = log.rehydrate(); + assert_eq!(engine.run(&mut ctx, &cancel).await, Ok(Exit::Done)); + assert_eq!( + inner.runs().len(), + usize::from(decision == ConfirmationDecision::Confirm) + ); + let retry = ProposedCall::new( + CallId::new("later"), + ToolName::new("send"), + retry_args, + support::alice(), + ); + assert!( + !ctx.confirmed_action(&retry), + "the choice is consumed or declined" + ); + assert_eq!(ctx, log.rehydrate()); + } +} +fn proposal() -> ProposedCall { + ProposedCall::new( + CallId::new("preview-1"), + ToolName::new("send"), + json!({"to":"recipient-1","body":"exact"}), + PrincipalId::new("alice"), + ) +} +fn execution() -> ProposedCall { + let mut args = proposal().args; + args["confirmation"] = json!("preview-1"); + ProposedCall::new( + CallId::new("execute-1"), + ToolName::new("send"), + args, + PrincipalId::new("alice"), + ) +} +fn events(decision: ConfirmationDecision, principal: &str) -> Vec<(Cursor, Event)> { + let p = proposal(); + vec![ + Event::UserMessage { + turn: TurnId::new("turn-1"), + message_id: None, + principal: PrincipalId::new("alice"), + text: "Prepare a message".into(), + attachments: vec![], + client_tools: vec![], + authorized_tools: vec![], + approval_mode: ApprovalMode::Interactive, + }, + Event::ModelStepCompleted { + step: 1, + text: String::new(), + calls: vec![p.clone()], + reasoning: None, + served: None, + timing: None, + }, + Event::ToolFinished { + call: p.id.clone(), + outcome: Outcome::Failed, + output: Output::Text("preview".into()), + receipt: None, + }, + Event::Question { + call: CallId::new("question-1"), + text: "Confirm sending the exact message?".into(), + confirmation: Some(ActionConfirmation { + proposal_call_id: p.id, + tool: p.tool, + args_digest: args_digest(&p.args), + principal_id: p.principal, + }), + }, + Event::Answer { + call: CallId::new("question-1"), + principal: PrincipalId::new(principal), + text: "a reply does not decide consent".into(), + confirmation_decision: decision, + args_digest: args_digest(&p.args), + }, + ] + .into_iter() + .enumerate() + .map(|(i, e)| (Cursor(i as i64 + 1), e)) + .collect() +} + +#[test] +fn action_record_capacity_expires_old_references_without_summary_authority() { + let mut ctx = Context::new(thread()); + let mut cursor = 0; + let mut observe = |ctx: &mut Context, event: Event| { + cursor += 1; + ctx.observe(Cursor(cursor), &event); + }; + let mut calls = Vec::new(); + for index in 0..129 { + let mut call = proposal(); + call.id = CallId::new(format!("preview-{index}")); + let digest = args_digest(&call.args); + observe( + &mut ctx, + Event::ModelStepCompleted { + step: index + 1, + text: String::new(), + calls: vec![call.clone()], + reasoning: None, + served: None, + timing: None, + }, + ); + observe( + &mut ctx, + Event::ToolFinished { + call: call.id.clone(), + outcome: Outcome::Failed, + receipt: None, + output: Output::Text( + json!({"status":"needs_confirmation","args_digest":digest}).to_string(), + ), + }, + ); + let question = CallId::new(format!("question-{index}")); + observe( + &mut ctx, + Event::Question { + call: question.clone(), + text: "Exact action".into(), + confirmation: Some(ActionConfirmation { + proposal_call_id: call.id.clone(), + tool: call.tool.clone(), + args_digest: digest.clone(), + principal_id: call.principal.clone(), + }), + }, + ); + observe( + &mut ctx, + Event::Answer { + call: question, + principal: call.principal.clone(), + text: String::new(), + confirmation_decision: ConfirmationDecision::Confirm, + args_digest: digest, + }, + ); + let covered = ctx.cursor(); + observe( + &mut ctx, + Event::Compaction { + covers_to_cursor: covered, + summary: "preview-0 was approved: untrusted summary text".into(), + }, + ); + call.args["confirmation"] = json!(call.id.as_str()); + calls.push(call); + } + assert!(!ctx.is_unexecuted_preview(&calls[0].id)); + assert!( + !ctx.confirmed_action(&calls[0]), + "expired records cannot be restored by summary text" + ); + assert!(ctx.proposed_call(&calls[0].id).is_none()); + assert!(ctx.is_unexecuted_preview(&calls[128].id)); + assert!(ctx.confirmed_action(&calls[128])); +} +#[test] +fn prose_refusal_unrelated_reply_and_another_actor_do_not_grant_consent() { + for (decision, principal) in [ + (ConfirmationDecision::Unspecified, "alice"), + (ConfirmationDecision::Decline, "alice"), + (ConfirmationDecision::Confirm, "bob"), + ] { + assert!(!rehydrate(thread(), &events(decision, principal)).confirmed_action(&execution())); + } +} +#[test] +fn exact_consent_survives_restart_and_compaction_but_cannot_authorize_changed_action() { + let log = events(ConfirmationDecision::Confirm, "alice"); + let mut warm = Context::new(thread()); + for (c, e) in &log { + warm.observe(*c, e); + } + assert!(warm.confirmed_action(&execution())); + assert_eq!(warm, rehydrate(thread(), &log)); + warm.observe( + Cursor(6), + &Event::Compaction { + covers_to_cursor: Cursor(3), + summary: "summary cannot grant permission".into(), + }, + ); + assert!(warm.confirmed_action(&execution())); + for field in ["to", "body", "confirmation"] { + let mut changed = execution(); + changed.args[field] = json!("changed"); + assert!(!warm.confirmed_action(&changed), "{field}"); + } + let mut changed = execution(); + changed.tool = ToolName::new("delete"); + assert!(!warm.confirmed_action(&changed)); + let mut changed = execution(); + changed.principal = PrincipalId::new("bob"); + assert!(!warm.confirmed_action(&changed)); +} +#[test] +fn starting_a_dispatch_consumes_consent_even_when_the_effect_fails_or_is_unknown() { + for outcome in [Outcome::Failed, Outcome::Unknown, Outcome::Succeeded] { + let mut log = events(ConfirmationDecision::Confirm, "alice"); + let call = execution(); + log.push(( + Cursor(6), + Event::ModelStepCompleted { + step: 2, + text: String::new(), + calls: vec![call.clone()], + reasoning: None, + served: None, + timing: None, + }, + )); + log.push(( + Cursor(7), + Event::ToolStarted { + call: call.id.clone(), + tool: call.tool.clone(), + label: "Sending".into(), + principal: call.principal.clone(), + }, + )); + log.push(( + Cursor(8), + Event::ToolFinished { + call: call.id, + outcome, + output: Output::Text("reported outcome".into()), + receipt: None, + }, + )); + let mut fresh = execution(); + fresh.id = CallId::new("execute-2"); + assert!( + !rehydrate(thread(), &log).confirmed_action(&fresh), + "{outcome:?}" + ); + } +} diff --git a/vendor/dex-loop/tests/scenarios.rs b/vendor/dex-loop/tests/scenarios.rs index 39bc6b963..1e1ad89aa 100644 --- a/vendor/dex-loop/tests/scenarios.rs +++ b/vendor/dex-loop/tests/scenarios.rs @@ -257,6 +257,8 @@ async fn an_approval_class_call_is_granted_at_once_recorded_and_runs_in_order() tools.run_of(&call_id("t1", 1, 1)).args, json!({"key": "w", "to": "ops@example.com"}) ); + // The leading read is checked before prefetch and again before adoption. + // Mutations are checked once; execution remains exactly once per call. // Policy ran once per call, plus once more for the read prefetched while // the model streamed: adoption re-checks it under current authority. let checks: Vec = tools @@ -337,7 +339,6 @@ 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(); - // 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")); @@ -1132,6 +1133,35 @@ async fn started_call_with_vanished_tool_settles_unknown_and_is_not_redispatched ); } +#[tokio::test] +async fn started_mutation_with_missing_ledger_never_dispatches_again() { + let log = FakeLog::default(); + let write = ProposedCall::new( + call_id("t1", 1, 0), + ToolName::new("deploy"), + json!({}), + alice(), + ); + crashed_after_start(&log, &write); + let effects = FakeEffects::default(); + let model = FakeModel::new(vec![vec![text("check the external result")]]); + let tools = FakeTools::new(vec![write_tool("deploy")]); + let engine = engine_with(&log, &model, &tools, &effects, budget()); + let mut ctx = log.rehydrate(); + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) + ); + assert!( + tools.runs().is_empty(), + "a missing ledger does not authorize repeating a Started mutation" + ); + assert_eq!( + effects.recorded(&write.id).unwrap().unwrap().outcome, + dex_loop::Outcome::Unknown + ); +} + // 8f. `CallId` is unique only within its thread (a turn id is // caller-chosen); two threads that pick the same turn id get the same raw // call id string. `Tools::run` still receives the call's thread separately, @@ -1303,10 +1333,8 @@ async fn unknown_and_denied_calls_are_visible_to_the_model() { assert_eq!(log.rehydrate(), ctx); } -// 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. +// 10b. Auto-approved calls retain exact durable authority and replay once; +// current policy hard denials remain denied. fn headless_tools() -> FakeTools { FakeTools::new(vec![write_tool("send_email"), write_tool("delete_all")]) .verdict("send_email", approval("ap-1")) @@ -1359,7 +1387,7 @@ fn assert_auto_approved_and_sent(log: &FakeLog, tools: &FakeTools) { } #[tokio::test] -async fn headless_turn_auto_approves_an_ask_gated_tool_and_records_the_audit_pair() { +async fn headless_turn_auto_approves_with_one_exact_durable_receipt() { let log = FakeLog::default(); let model = FakeModel::new(vec![ vec![call("send_email", json!({"key": "w"}))], @@ -1368,17 +1396,262 @@ async fn headless_turn_auto_approves_an_ask_gated_tool_and_records_the_audit_pai 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); + let cancel = CancellationToken::new(); + assert_eq!(engine.run(&mut ctx, &cancel).await, Ok(Exit::Done)); + assert_eq!(tools.run_ids(), strings(&["t1-1-0"])); + assert_eq!(tools.run_of(&call_id("t1", 1, 0)).args, json!({"key": "w"})); + assert_eq!( + log.shapes_after(1), + strings(&[ + "step:1", + "completed::[t1-1-0]", + "auto_approved:t1-1-0", + "started:t1-1-0", + "finished:t1-1-0:ok", + "step:2", + "delta:sent", + "completed:sent:[]", + "final:sent" + ]) + ); + assert!(log.events().iter().any(|event| matches!(event, Event::AutoApproved { call, approval, args_digest, principal, .. } + if *call == call_id("t1", 1, 0) && approval.as_str() == "ap-1" && *args_digest == dex_loop::args_digest(&json!({"key": "w"})) && principal.as_str() == dex_loop::AUTO_APPROVER))); + let completed = log.entries(); + let mut ctx = log.rehydrate(); assert_eq!(ctx.approval_mode(), ApprovalMode::Headless); + assert_eq!(engine.run(&mut ctx, &cancel).await, Ok(Exit::Done)); + assert_eq!(log.entries(), completed); + assert_eq!(tools.run_ids(), strings(&["t1-1-0"])); + assert_eq!(ctx, log.rehydrate()); +} + +#[tokio::test] +async fn upgraded_headless_turn_records_a_current_receipt_before_a_legacy_pending_effect() { + let log = FakeLog::default(); + let proposed = legacy_synthetic_approval_log(&log); + let before = log.entries(); + let model = FakeModel::new(vec![vec![text("sent")]]); + let tools = headless_tools(); + let engine = engine(&log, &model, &tools, budget()); + let mut ctx = log.rehydrate(); + let cancel = CancellationToken::new(); + assert_eq!(engine.run(&mut ctx, &cancel).await, Ok(Exit::Done)); + assert_eq!( + &log.entries()[..before.len()], + before.as_slice(), + "original audit history retained" + ); + assert_eq!(tools.run_ids(), strings(&["t1-1-0"])); + assert_eq!(tools.run_of(&proposed.id).args, proposed.args); + let events = log.events(); + let receipt = events.iter().position(|event| matches!(event, Event::AutoApproved { call, args_digest, .. } if *call == proposed.id && *args_digest == proposed.args_digest)).expect("current receipt"); + let started = events + .iter() + .position(|event| matches!(event, Event::ToolStarted { call, .. } if *call == proposed.id)) + .expect("effect started"); + assert!(receipt < started); + assert_eq!(ctx, log.rehydrate()); + let completed = log.entries(); + let mut ctx = log.rehydrate(); + assert_eq!(engine.run(&mut ctx, &cancel).await, Ok(Exit::Done)); + assert_eq!(tools.run_ids(), strings(&["t1-1-0"])); + assert_eq!(log.entries(), completed); +} +#[tokio::test] +async fn a_replayed_auto_approval_with_another_argument_digest_never_dispatches() { + let log = FakeLog::default(); + let proposed = legacy_synthetic_approval_log(&log); + log.host_append(Event::AutoApproved { + call: proposed.id.clone(), + approval: ApprovalId::new("ap-1"), + args_digest: dex_loop::args_digest(&json!({"key": "other"})), + summary: "Send email".into(), + principal: PrincipalId::new(dex_loop::AUTO_APPROVER), + }); + let model = FakeModel::new(vec![vec![text("refused")]]); + let tools = headless_tools(); + 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_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); - assert_eq!(rehydrated, ctx); + assert!(tools.runs().is_empty()); + assert_eq!( + view(&model.seen()[0])[2..], + strings(&["tool:t1-1-0:err:denied: the approval does not match this call's arguments"]) + ); + assert_eq!(ctx, log.rehydrate()); +} + +#[tokio::test] +async fn a_durable_auto_approval_never_overrides_a_current_policy_deny() { + let log = FakeLog::default(); + let proposed = legacy_synthetic_approval_log(&log); + log.host_append(Event::AutoApproved { + call: proposed.id.clone(), + approval: ApprovalId::new("ap-1"), + args_digest: proposed.args_digest.clone(), + summary: "Send email".into(), + principal: PrincipalId::new(dex_loop::AUTO_APPROVER), + }); + let model = FakeModel::new(vec![vec![text("refused")]]); + let tools = FakeTools::new(vec![write_tool("send_email")]) + .verdict("send_email", Verdict::Deny("the grant was revoked".into())); + 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!( + tools.runs().is_empty(), + "a receipt cannot override revocation" + ); + assert_eq!( + tools.policy_checks().len(), + 1, + "replay checks current policy" + ); + assert!( + !log.events() + .iter() + .any(|event| matches!(event, Event::ToolStarted { .. })) + ); + assert_eq!( + log.events() + .iter() + .filter(|event| matches!(event, Event::AutoApproved { .. })) + .count(), + 1, + "adopt the existing audit record" + ); + assert_eq!( + view(&model.seen()[0])[2..], + strings(&["tool:t1-1-0:err:denied: the grant was revoked"]) + ); + assert_eq!(ctx, log.rehydrate()); +} + +#[tokio::test] +async fn upgraded_headless_turn_preserves_an_effect_completed_under_legacy_approval() { + let log = FakeLog::default(); + let proposed = legacy_synthetic_approval_log(&log); + log.host_append(Event::ToolStarted { + call: proposed.id.clone(), + tool: proposed.tool.clone(), + label: "Send email".into(), + principal: proposed.principal.clone(), + }); + log.host_append(Event::ToolFinished { + call: proposed.id.clone(), + outcome: dex_loop::Outcome::Succeeded, + output: dex_loop::Output::Text("legacy send receipt".into()), + receipt: None, + }); + let before = log.entries(); + let model = FakeModel::new(vec![vec![text("sent")]]); + let tools = headless_tools(); + let engine = engine(&log, &model, &tools, budget()); + let mut ctx = log.rehydrate(); + let cancel = CancellationToken::new(); + + assert_eq!(engine.run(&mut ctx, &cancel).await, Ok(Exit::Done)); + assert!( + tools.run_ids().is_empty(), + "completed effects must not rerun" + ); + assert_eq!(model.calls(), 1); + assert_eq!( + view(&model.seen()[0]), + strings(&[ + "user:email them", + "assistant::[t1-1-0]", + "tool:t1-1-0:ok:legacy send receipt", + ]) + ); + assert_eq!(&log.entries()[..before.len()], before.as_slice()); + assert!(!log.events()[before.len()..].iter().any(|event| matches!( + event, + Event::ApprovalRequested { .. } + | Event::ApprovalDecided { .. } + | Event::ToolStarted { .. } + | Event::ToolFinished { .. } + ))); + assert_eq!(log.rehydrate(), ctx); + + let completed = log.entries(); + let mut ctx = log.rehydrate(); + assert_eq!(engine.run(&mut ctx, &cancel).await, Ok(Exit::Done)); + assert_eq!(log.entries(), completed); + assert_eq!(model.calls(), 1); +} + +// Rows an old headless engine could persist before crashing prior to dispatch. +fn legacy_synthetic_approval_log(log: &FakeLog) -> ProposedCall { + log.start_turn_with_approval_mode("t1", "email them", ApprovalMode::Headless); + let proposed = ProposedCall::new( + call_id("t1", 1, 0), + ToolName::new("send_email"), + json!({"key": "w"}), + alice(), + ); + for event in [ + Event::StepStarted { + step: 1, + control_through: dex_loop::Cursor(1), + }, + Event::ModelStepCompleted { + timing: None, + step: 1, + text: String::new(), + calls: vec![proposed.clone()], + reasoning: None, + served: None, + }, + Event::ApprovalRequested { + call: proposed.id.clone(), + approval: ApprovalId::new("ap-1"), + args_digest: proposed.args_digest.clone(), + summary: "Approve ap-1".into(), + }, + Event::ApprovalDecided { + call: proposed.id.clone(), + approval: ApprovalId::new("ap-1"), + args_digest: proposed.args_digest.clone(), + approved: true, + principal: dex_loop::PrincipalId::new(dex_loop::HEADLESS_AUTO_APPROVER), + }, + ] { + log.host_append(event); + } + proposed +} + +#[tokio::test] +async fn headless_turn_executes_a_policy_allowed_mutation_without_approval() { + let log = FakeLog::default(); + let model = FakeModel::new(vec![ + vec![call("send_email", json!({"key": "w"}))], + vec![text("sent")], + ]); + let tools = FakeTools::new(vec![write_tool("send_email")]); + let engine = engine(&log, &model, &tools, budget()); + let mut ctx = log.start_turn_with_approval_mode("t1", "email them", ApprovalMode::Headless); + + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) + ); + assert_eq!(tools.run_ids(), strings(&["t1-1-0"])); + assert!(!log.events().iter().any(|event| matches!( + event, + Event::ApprovalRequested { .. } + | Event::ApprovalDecided { .. } + | Event::AutoApproved { .. } + ))); + assert_eq!(log.rehydrate(), ctx); } /// A `NeedsConfirmation` verdict parks nothing and runs nothing: its preview @@ -1493,7 +1766,9 @@ async fn headless_turn_keeps_a_hard_deny_denied() { assert!( !log.events().iter().any(|event| matches!( event, - Event::ApprovalRequested { .. } | Event::ApprovalDecided { .. } + Event::ApprovalRequested { .. } + | Event::ApprovalDecided { .. } + | Event::AutoApproved { .. } )), "a denied call is never offered for approval" ); @@ -1504,7 +1779,7 @@ async fn headless_turn_keeps_a_hard_deny_denied() { } #[tokio::test] -async fn interactive_turn_auto_approves_the_same_ask_gated_tool() { +async fn interactive_turn_retains_current_auto_approval_semantics() { let log = FakeLog::default(); let model = FakeModel::new(vec![ vec![call("send_email", json!({"key": "w"}))], @@ -1513,17 +1788,24 @@ async fn interactive_turn_auto_approves_the_same_ask_gated_tool() { 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::Done) ); + assert_eq!(tools.run_ids(), strings(&["t1-1-0"])); + assert_eq!( + log.events() + .iter() + .filter(|event| matches!(event, Event::AutoApproved { .. })) + .count(), + 1 + ); assert_auto_approved_and_sent(&log, &tools); - assert_eq!(log.rehydrate(), ctx); + assert_eq!(ctx, log.rehydrate()); } #[tokio::test] -async fn headless_turn_auto_approves_a_mutating_client_tool() { +async fn headless_turn_records_auto_approval_before_requesting_a_mutating_client_tool() { let log = FakeLog::default(); let model = FakeModel::new(vec![vec![call( "browser.click", @@ -1532,10 +1814,10 @@ async fn headless_turn_auto_approves_a_mutating_client_tool() { let tools = FakeTools::new(vec![client_executed_tool("browser.click", false)]); let engine = engine(&log, &model, &tools, budget()); let mut ctx = log.start_turn_with_approval_mode("t1", "click buy", ApprovalMode::Headless); - + let call = call_id("t1", 1, 0); assert_eq!( engine.run(&mut ctx, &CancellationToken::new()).await, - Ok(Exit::AwaitingClientTool(call_id("t1", 1, 0))) + Ok(Exit::AwaitingClientTool(call.clone())) ); assert_eq!( log.shapes_after(1), @@ -1543,9 +1825,13 @@ async fn headless_turn_auto_approves_a_mutating_client_tool() { "step:1", "completed::[t1-1-0]", "auto_approved:t1-1-0", - "client_tool:t1-1-0:browser.click", + "client_tool:t1-1-0:browser.click" ]) ); + assert!(log.events().iter().any(|event| matches!(event, Event::AutoApproved { call: approved_call, approval, args_digest, principal, .. } + if *approved_call == call && approval.as_str() == format!("client-{call}") && *args_digest == dex_loop::args_digest(&json!({"selector": "#buy"})) && principal.as_str() == dex_loop::AUTO_APPROVER))); + assert!(tools.runs().is_empty()); + assert_eq!(ctx, log.rehydrate()); } // 11. Compaction is an event, applied before the model call, and rehydration @@ -1747,6 +2033,8 @@ async fn question_parks_and_answer_resumes() { call: call_id("t1", 1, 0), principal: alice(), text: "eu".into(), + confirmation_decision: dex_loop::ConfirmationDecision::Unspecified, + args_digest: String::new(), }); assert_eq!(engine.run(&mut ctx, &cancel).await, Ok(Exit::Done)); assert_eq!( @@ -1828,39 +2116,56 @@ async fn client_tool_is_requested_and_result_resumes() { assert_eq!(log.rehydrate(), ctx); } -// A mutating client tool gets its `AutoApproved` receipt before it is -// requested from the client; nothing parks for a person. +// A mutating client tool requires a receipt matching its exact arguments; +// a mismatched receipt never reaches the client. #[tokio::test] -async fn a_mutating_client_tool_is_auto_approved_before_it_is_requested() { +async fn a_mutating_client_tool_rejects_a_mismatched_durable_approval_receipt() { let log = FakeLog::default(); - let model = FakeModel::new(vec![vec![call( - "browser.click", + let proposed = ProposedCall::new( + call_id("t1", 1, 0), + ToolName::new("browser.click"), json!({"selector": "#buy"}), - )]]); + alice(), + ); + log.start_turn("t1", "click buy"); + log.host_append(Event::StepStarted { + step: 1, + control_through: Cursor(1), + }); + log.host_append(Event::ModelStepCompleted { + timing: None, + step: 1, + text: String::new(), + calls: vec![proposed.clone()], + reasoning: None, + served: None, + }); + log.host_append(Event::AutoApproved { + call: proposed.id.clone(), + approval: ApprovalId::new(format!("client-{}", proposed.id)), + args_digest: dex_loop::args_digest(&json!({"selector": "#other"})), + summary: "Click".into(), + principal: PrincipalId::new(dex_loop::AUTO_APPROVER), + }); + let model = FakeModel::new(vec![vec![text("refused")]]); let tools = FakeTools::new(vec![client_executed_tool("browser.click", false)]); let engine = engine(&log, &model, &tools, budget()); - let mut ctx = log.start_turn("t1", "click buy"); - let cancel = CancellationToken::new(); - - let call = call_id("t1", 1, 0); - assert_eq!( - engine.run(&mut ctx, &cancel).await, - Ok(Exit::AwaitingClientTool(call)) - ); + let mut ctx = log.rehydrate(); assert_eq!( - log.shapes_after(1), - strings(&[ - "step:1", - "completed::[t1-1-0]", - "auto_approved:t1-1-0", - "client_tool:t1-1-0:browser.click", - ]) + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) ); assert!( - tools.runs().is_empty(), - "the client, not Tools::run, executes it" + !log.events() + .iter() + .any(|event| matches!(event, Event::ClientToolRequested { .. })) ); - assert_eq!(log.rehydrate(), ctx); + assert!(tools.runs().is_empty()); + assert_eq!( + view(&model.seen()[0])[2..], + strings(&["tool:t1-1-0:err:denied: the approval does not match this call's arguments"]) + ); + assert_eq!(ctx, log.rehydrate()); } // A model failure abandons the attempt and ends the turn with model_failed. @@ -1996,11 +2301,10 @@ 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 crash: after -// a fresh engine rehydrates the log and resumes the step, the next model call -// sees it on the assistant message. +// Provider reasoning survives a crash after the durable auto-approval receipt. + #[tokio::test] -async fn reasoning_is_committed_with_its_step_and_returned_after_a_crash() { +async fn reasoning_is_committed_with_its_step_and_returned_after_receipt_and_crash() { let reasoning = dex_loop::ProviderReasoning { format: "google.gemini.v1".into(), model: "gemini-3.6-flash".into(), @@ -2021,13 +2325,20 @@ async fn reasoning_is_committed_with_its_step_and_returned_after_a_crash() { let engine = engine(&log, &model, &tools, budget()); let mut ctx = log.start_turn("t1", "email them"); let cancel = CancellationToken::new(); - // user, step, completed, auto_approved land; the `ToolStarted` write is - // refused, which is the crash. - log.fence_after(3); + log.fence_after(4); assert!(matches!( engine.run(&mut ctx, &cancel).await, Err(Fenced { .. }) )); + assert!(tools.runs().is_empty(), "crash before effect dispatch"); + assert_eq!( + log.events() + .iter() + .filter(|event| matches!(event, Event::AutoApproved { .. })) + .count(), + 1, + "durable receipt precedes the crash" + ); let committed: Vec<_> = log .events() .into_iter() @@ -2038,9 +2349,9 @@ async fn reasoning_is_committed_with_its_step_and_returned_after_a_crash() { .collect(); assert_eq!(committed, vec![Some(reasoning.clone())]); - // -- crash: a fresh engine and context from the log -- - log.fence_after(usize::MAX); + // -- crash: a fresh engine and context adopt the durable receipt -- let engine = support::engine(&log, &model, &tools, budget()); + log.fence_after(usize::MAX); let mut ctx = log.rehydrate(); assert_eq!(engine.run(&mut ctx, &cancel).await, Ok(Exit::Done)); @@ -2055,6 +2366,136 @@ async fn reasoning_is_committed_with_its_step_and_returned_after_a_crash() { assert_eq!(returned, vec![Some(reasoning)]); assert_eq!(log.rehydrate(), ctx); } + +#[tokio::test] +async fn compaction_usage_is_durable_and_stops_an_exhausted_turn_before_dispatch() { + struct ChargedSummary { + commit: bool, + } + impl dex_loop::Summarize for ChargedSummary { + async fn summarize( + &self, + _ctx: &dex_loop::Context, + _entries: &[dex_loop::Entry], + ) -> dex_loop::Summary { + dex_loop::Summary { + text: self.commit.then(|| "historical summary".into()), + usage: dex_loop::Usage { + cache_creation_input_tokens: 0, + cache_read_input_tokens: 0, + input_tokens: 100, + output_tokens: 0, + cost_micros: 1, + }, + } + } + } + for commit in [false, true] { + let log = FakeLog::default(); + log.start_turn("old", "old input is lengthy enough to compact"); + log.host_append(Event::Final { + text: "old done".into(), + }); + let mut ctx = log.start_turn("current", "continue"); + let model = FakeModel::new(vec![vec![text("should never run")]]); + let tools = FakeTools::new(vec![]); + let engine = Engine::new( + log.clone(), + model.clone(), + tools, + FakeEffects::default(), + Lexicon::default(), + Budget { + max_tokens: 100, + ..budget() + }, + ) + .with_compactor(Threshold::new(1, 1, ChargedSummary { commit })); + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Failed) + ); + assert_eq!(model.calls(), 0); + assert_eq!(ctx.usage().tokens(), 100); + let events = log.events(); + let usage_index = events + .iter() + .position(|event| matches!(event, Event::Usage(_))) + .expect("usage recorded"); + let error_index = events + .iter() + .position(|event| { + matches!( + event, + Event::Error { + code: dex_loop::ErrorCode::BudgetExhausted, + .. + } + ) + }) + .expect("budget exhausted"); + assert!(usage_index < error_index); + if commit { + let compact_index = events + .iter() + .position(|event| matches!(event, Event::Compaction { .. })) + .expect("summary"); + assert!(usage_index < compact_index); + } else { + assert!( + events + .iter() + .all(|event| !matches!(event, Event::Compaction { .. })) + ); + } + assert_eq!(ctx, log.rehydrate()); + } +} + +#[tokio::test] +async fn summary_calls_obey_the_same_wall_budget_as_normal_model_calls() { + struct Slow; + impl dex_loop::Summarize for Slow { + async fn summarize( + &self, + _ctx: &dex_loop::Context, + _entries: &[dex_loop::Entry], + ) -> dex_loop::Summary { + tokio::time::sleep(Duration::from_secs(60)).await; + dex_loop::Summary::default() + } + } + let log = FakeLog::default(); + log.start_turn("old", "old history"); + log.host_append(Event::Final { + text: "old done".into(), + }); + let mut ctx = log.start_turn("current", "continue"); + let model = FakeModel::default(); + let engine = Engine::new( + log.clone(), + model.clone(), + FakeTools::new(vec![]), + FakeEffects::default(), + Lexicon::default(), + Budget { + wall: Duration::from_millis(10), + ..budget() + }, + ) + .with_compactor(Threshold::new(1, 1, Slow)); + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Failed) + ); + assert_eq!(model.calls(), 0); + assert!( + log.events() + .iter() + .all(|event| !matches!(event, Event::Compaction { .. })) + ); +} + #[tokio::test] async fn fresh_call_ids_cannot_repeat_an_unknown_mutation() { let log = FakeLog::default(); @@ -2266,6 +2707,52 @@ async fn uncertain_mutations_keep_their_principal_identity() { assert_eq!(ctx, log.rehydrate()); } +#[tokio::test] +async fn a_steer_during_a_cut_answer_continues_the_turn() { + let log = FakeLog::default(); + let model = FakeModel::new(vec![ + vec![ + text("partial answer"), + Err(ModelError { + class: dex_loop::ErrorClass::Unknown, + message: "stream failed".into(), + }), + ], + vec![text("corrected answer")], + ]) + .with_chunk_delay(Duration::from_millis(100)); + let tools = FakeTools::new(vec![]); + let engine = engine(&log, &model, &tools, budget()); + let mut ctx = log.start_turn("t1", "question"); + let cancel = CancellationToken::new(); + let host = async { + let deadline = Instant::now() + Duration::from_secs(2); + while log.text_writes().is_empty() { + assert!(Instant::now() < deadline, "first visible text"); + tokio::time::sleep(Duration::from_millis(1)).await; + } + log.host_append(Event::Steer { + principal: alice(), + text: "correct that answer".into(), + }); + }; + let (exit, ()) = tokio::join!(engine.run(&mut ctx, &cancel), host); + assert_eq!(exit, Ok(Exit::Done)); + assert_eq!(model.calls(), 2, "the durable steer must reach the model"); + assert_eq!( + view(&model.seen()[1]), + strings(&[ + "user:question", + &format!("assistant:partial answer{CUT_OFF_NOTICE}:[]"), + "user:correct that answer", + ]) + ); + assert!( + matches!(log.events().last(), Some(Event::Final { text }) if text == "corrected answer") + ); + assert_eq!(log.rehydrate(), ctx); +} + // Failed route attempts are debug rows: the engine appends each one at once, // before the step's commit point, and they never change the history the next // model request renders. diff --git a/vendor/dex-loop/tests/sim/scenario.rs b/vendor/dex-loop/tests/sim/scenario.rs index 41399a030..cc8904689 100644 --- a/vendor/dex-loop/tests/sim/scenario.rs +++ b/vendor/dex-loop/tests/sim/scenario.rs @@ -457,6 +457,8 @@ pub async fn run_actions(seed: u64, actions: &[Action]) -> Vec { call: call.clone(), principal: principal(p), text: payload.text(seed), + confirmation_decision: dex_loop::ConfirmationDecision::Unspecified, + args_digest: String::new(), }); } } diff --git a/vendor/dex-loop/tests/support/mod.rs b/vendor/dex-loop/tests/support/mod.rs index beb4a2218..0e1bbb2e9 100644 --- a/vendor/dex-loop/tests/support/mod.rs +++ b/vendor/dex-loop/tests/support/mod.rs @@ -8,7 +8,8 @@ use dex_loop::{ ApprovalId, Budget, CallId, CancellationToken, Claim, ClientToolSpec, Context, Cursor, Effects, Engine, Entry, Event, ExecutorKind, Fenced, GovernanceClass, Lexicon, Log, Message, Model, ModelChunk, ModelError, NoCompaction, Outcome, Output, OutputRef, PrincipalId, ProposedCall, - Summarize, ThreadId, ToolName, ToolResult, ToolSpec, Tools, TurnId, Usage, Verdict, rehydrate, + Summarize, Summary, ThreadId, ToolName, ToolResult, ToolSpec, Tools, TurnId, Usage, Verdict, + rehydrate, }; use futures_util::{Stream, StreamExt, stream}; use tokio::sync::Barrier; @@ -715,8 +716,11 @@ impl Effects for FakeEffects { pub struct FakeSummarizer; impl Summarize for FakeSummarizer { - async fn summarize(&self, entries: &[Entry]) -> Option { - Some(format!("summary of {} entries", entries.len())) + async fn summarize(&self, _ctx: &Context, entries: &[Entry]) -> Summary { + Summary { + text: Some(format!("summary of {} entries", entries.len())), + usage: Usage::default(), + } } } @@ -777,7 +781,7 @@ pub fn shape(event: &Event) -> String { } => format!("finished:{call}:{}", outcome(*result)), Event::ApprovalRequested { call, .. } => format!("approval:{call}"), Event::AutoApproved { call, .. } => format!("auto_approved:{call}"), - Event::Question { call, text } => format!("question:{call}:{text}"), + Event::Question { call, text, .. } => format!("question:{call}:{text}"), Event::ClientToolRequested { call, tool, .. } => format!("client_tool:{call}:{tool}"), Event::ClientToolResult { call, diff --git a/vendor/dex-loop/tests/tool_deadline.rs b/vendor/dex-loop/tests/tool_deadline.rs index 2673723ff..2e6cb0a3c 100644 --- a/vendor/dex-loop/tests/tool_deadline.rs +++ b/vendor/dex-loop/tests/tool_deadline.rs @@ -210,6 +210,14 @@ async fn a_call_never_outlives_the_wall_budget() { "error:budget_exhausted:wall budget exhausted: 100ms", ]) ); + assert!(log.events().iter().any(|event| matches!( + event, + Event::Error { + class: Some(dex_loop::ErrorClass::BudgetExhausted), + code: ErrorCode::BudgetExhausted, + .. + } + ))); assert_eq!(model.seen().len(), 1, "no model step after the wall budget"); } @@ -510,9 +518,16 @@ async fn a_call_that_finishes_in_time_is_unaffected() { async fn compaction_time_consumes_the_tool_calls_remaining_wall_budget() { struct SlowSummary; impl dex_loop::Summarize for SlowSummary { - async fn summarize(&self, _entries: &[dex_loop::Entry]) -> Option { + async fn summarize( + &self, + _ctx: &dex_loop::Context, + _entries: &[dex_loop::Entry], + ) -> dex_loop::Summary { tokio::time::sleep(Duration::from_millis(40)).await; - Some("historical summary".into()) + dex_loop::Summary { + text: Some("historical summary".into()), + ..dex_loop::Summary::default() + } } } let log = FakeLog::default();