diff --git a/.repository-projection.json b/.repository-projection.json index 0ba123d70..cb4e32c18 100644 --- a/.repository-projection.json +++ b/.repository-projection.json @@ -3,11 +3,11 @@ "projection": "deixic-code", "projectionSchemaVersion": 1, "sourceRepository": "dx-corp/mono", - "sourceSha": "ac81c8d7cff3a1002c8c5dbfac909bbaea58a854", + "sourceSha": "61ea3dea08e87af097277910b63a97fe6e216cf2", "destinationRepository": "dx-corp/code", - "priorProjectedBase": "88f93e94824e24ef8cd97b3b9a3c4c02777b9549", + "priorProjectedBase": "11ca01744d12603c659827f5cec9e86e834b5da8", "definitionDigest": "82936441c776e3e8edb5d215a75007ec9714a233f489d460075d79d5ef5ba32f", "toolDigest": "c244d99199a7ae3eb8ff644a99462163c23b0bb6a83ef50af01efbdca0b81d04", - "contentDigest": "e57247ac09795028aad4cdee2023bcf4ad7a30fd7afbe126c5fad2e5bb70cb14", + "contentDigest": "9d215910206a5cea42c6e297cd0553a1ae73a53950fad655a696861809f80675", "publicationEligible": true } diff --git a/packages/ai-rs/src/openai.rs b/packages/ai-rs/src/openai.rs index 048c295ca..69c1eaa69 100644 --- a/packages/ai-rs/src/openai.rs +++ b/packages/ai-rs/src/openai.rs @@ -236,6 +236,9 @@ async fn send_with_response_open_timeout( #[path = "openai/managed_gateway.rs"] mod managed_gateway; +#[cfg(test)] +#[path = "openai/response_open_timeout_tests.rs"] +mod response_open_timeout_tests; use managed_gateway::{ managed_gateway_error_retry_after, managed_gateway_receipt, managed_provider_tools_evidence, }; @@ -6899,74 +6902,6 @@ data: {"type":"response.completed","response":{"output":[{"type":"message","cont server.await.unwrap(); } - #[tokio::test] - async fn managed_gateway_response_open_timeout_stops_stalled_headers() { - use std::io::{Read, Write}; - use std::net::TcpListener; - - let listener = TcpListener::bind("127.0.0.1:0").expect("bind mock gateway"); - let address = listener.local_addr().expect("mock gateway address"); - std::thread::spawn(move || { - let (mut stream, _) = listener.accept().expect("accept gateway request"); - let mut request = [0_u8; 4096]; - let _ = stream.read(&mut request).expect("read gateway request"); - std::thread::sleep(std::time::Duration::from_millis(150)); - let _ = stream - .write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"); - }); - - let request = reqwest::Client::new().get(format!("http://{address}/responses")); - let error = send_with_response_open_timeout( - request, - Some(std::time::Duration::from_millis(25)), - "managed gateway", - ) - .await - .expect_err("stalled response opening must time out"); - - assert!( - error - .to_string() - .contains("managed gateway response headers timed out"), - "unexpected error: {error:#}" - ); - } - - #[tokio::test] - async fn managed_gateway_response_open_timeout_does_not_cover_stream_body() { - use std::io::{Read, Write}; - use std::net::TcpListener; - - let listener = TcpListener::bind("127.0.0.1:0").expect("bind mock gateway"); - let address = listener.local_addr().expect("mock gateway address"); - std::thread::spawn(move || { - let (mut stream, _) = listener.accept().expect("accept gateway request"); - let mut request = [0_u8; 4096]; - let _ = stream.read(&mut request).expect("read gateway request"); - stream - .write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 5\r\nConnection: close\r\n\r\n") - .expect("write response headers"); - stream.flush().expect("flush response headers"); - std::thread::sleep(std::time::Duration::from_millis(150)); - stream.write_all(b"hello").expect("write delayed body"); - }); - - let request = reqwest::Client::new().get(format!("http://{address}/responses")); - let response = send_with_response_open_timeout( - request, - Some(std::time::Duration::from_millis(50)), - "managed gateway", - ) - .await - .expect("response headers should arrive within the open timeout"); - let body = response - .text() - .await - .expect("body may continue past the response-open timeout"); - - assert_eq!(body, "hello"); - } - #[test] fn managed_gateway_scope_adds_workspace_header_without_changing_provider_reference() { let client = diff --git a/packages/ai-rs/src/openai/response_open_timeout_tests.rs b/packages/ai-rs/src/openai/response_open_timeout_tests.rs new file mode 100644 index 000000000..8d6ebd69f --- /dev/null +++ b/packages/ai-rs/src/openai/response_open_timeout_tests.rs @@ -0,0 +1,143 @@ +//! Response-header deadlines preserve delayed-body streaming. +use super::send_with_response_open_timeout; + +// Real loopback I/O must finish without Tokio advancing a paused clock to +// the open deadline just because the socket is briefly idle. A runnable +// task keeps time under the test's explicit control until it is aborted. +fn prevent_network_clock_auto_advance() -> tokio::task::JoinHandle<()> { + tokio::spawn(async { + loop { + tokio::task::yield_now().await; + } + }) +} + +#[tokio::test(start_paused = true)] +async fn managed_gateway_response_open_timeout_stops_stalled_headers() { + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + + let clock = prevent_network_clock_auto_advance(); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0") + .await + .expect("bind mock gateway"); + let address = listener.local_addr().expect("mock gateway address"); + let (accepted, received) = tokio::sync::oneshot::channel(); + let server = tokio::spawn(async move { + let (mut stream, _) = listener.accept().await.expect("accept gateway request"); + let mut request = [0_u8; 4096]; + assert!( + stream + .read(&mut request) + .await + .expect("read gateway request") + > 0 + ); + accepted.send(()).expect("request readiness"); + tokio::time::sleep(std::time::Duration::from_millis(150)).await; + let _ = stream + .write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 0\r\nConnection: close\r\n\r\n") + .await; + }); + + let request = reqwest::Client::builder() + .no_proxy() + .build() + .unwrap() + .get(format!("http://{address}/responses")); + let client = tokio::spawn(send_with_response_open_timeout( + request, + Some(std::time::Duration::from_millis(25)), + "managed gateway", + )); + received.await.expect("request reached gateway"); + assert!( + !client.is_finished(), + "headers are still pending before the deadline" + ); + tokio::time::advance(std::time::Duration::from_millis(25)).await; + let error = client + .await + .unwrap() + .expect_err("stalled response opening must time out"); + + assert!( + error + .to_string() + .contains("managed gateway response headers timed out"), + "unexpected error: {error:#}" + ); + tokio::time::advance(std::time::Duration::from_millis(125)).await; + server.await.expect("mock gateway server"); + clock.abort(); +} + +#[tokio::test(start_paused = true)] +async fn managed_gateway_response_open_timeout_does_not_cover_stream_body() { + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + + let clock = prevent_network_clock_auto_advance(); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0") + .await + .expect("bind mock gateway"); + let address = listener.local_addr().expect("mock gateway address"); + let (headers_written, received) = tokio::sync::oneshot::channel(); + let server = tokio::spawn(async move { + let (mut stream, _) = listener.accept().await.expect("accept gateway request"); + let mut request = [0_u8; 4096]; + assert!( + stream + .read(&mut request) + .await + .expect("read gateway request") + > 0 + ); + stream + .write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 5\r\nConnection: close\r\n\r\n") + .await + .expect("write response headers"); + stream.flush().await.expect("flush response headers"); + headers_written.send(()).expect("header readiness"); + tokio::time::sleep(std::time::Duration::from_millis(150)).await; + stream + .write_all(b"hello") + .await + .expect("write delayed body"); + }); + + let request = reqwest::Client::builder() + .no_proxy() + .build() + .unwrap() + .get(format!("http://{address}/responses")); + let client = tokio::spawn(send_with_response_open_timeout( + request, + Some(std::time::Duration::from_millis(50)), + "managed gateway", + )); + received.await.expect("gateway wrote headers"); + let response = client + .await + .unwrap() + .expect("response headers should arrive within the open timeout"); + let body = response.text(); + tokio::pin!(body); + assert!(futures::poll!(&mut body).is_pending()); + tokio::time::advance(std::time::Duration::from_millis(50)).await; + assert!( + futures::poll!(&mut body).is_pending(), + "body must outlive the 50 ms open timeout" + ); + tokio::time::advance(std::time::Duration::from_millis(99)).await; + assert!( + futures::poll!(&mut body).is_pending(), + "body remains pending until its 150 ms delay" + ); + tokio::time::advance(std::time::Duration::from_millis(1)).await; + let body = body + .await + .expect("body may continue past the response-open timeout"); + + assert_eq!(body, "hello"); + server.await.expect("mock gateway server"); + clock.abort(); +} diff --git a/packages/local-host-rs/src/native_credentials.rs b/packages/local-host-rs/src/native_credentials.rs index 3882c31d1..9fc8da352 100644 --- a/packages/local-host-rs/src/native_credentials.rs +++ b/packages/local-host-rs/src/native_credentials.rs @@ -1,14 +1,17 @@ //! Native credential access shared by OAuth, connections, and MCP. //! -//! Unit tests cannot open the developer's credential store. Process-based -//! fixtures use MAESTRO_DISABLE_KEYCHAIN=1 to enforce the same boundary. +//! Unit tests and opt-in test-support dependency builds cannot open the +//! developer's credential store. Other process-based fixtures use +//! MAESTRO_DISABLE_KEYCHAIN=1 to enforce the same boundary. use anyhow::{Result, bail}; pub(crate) fn entry(service: &str, account: &str) -> Result { let disabled = std::env::var("MAESTRO_DISABLE_KEYCHAIN").ok(); - open_with_policy(cfg!(test), disabled.as_deref(), || { - keyring::Entry::new(service, account).map_err(Into::into) - }) + open_with_policy( + cfg!(any(test, feature = "test-support")), + disabled.as_deref(), + || keyring::Entry::new(service, account).map_err(Into::into), + ) } fn open_with_policy( @@ -17,7 +20,9 @@ fn open_with_policy( open: impl FnOnce() -> Result, ) -> Result { if unit_test { - bail!("native credential access is disabled in unit tests; inject a test secret backend"); + bail!( + "native credential access is disabled in unit tests and test-support builds; inject a test secret backend" + ); } if disabled == Some("1") { bail!("native credential access is disabled by MAESTRO_DISABLE_KEYCHAIN=1"); diff --git a/packages/local-host-rs/tests/embedding.rs b/packages/local-host-rs/tests/embedding.rs index 65d98bd47..00931364c 100644 --- a/packages/local-host-rs/tests/embedding.rs +++ b/packages/local-host-rs/tests/embedding.rs @@ -12,6 +12,56 @@ use maestro_local_host::embedding::{ }; use maestro_local_host::state::ApprovalMode; +#[test] +fn test_support_dependency_cannot_open_native_credentials() { + const CHILD: &str = "MAESTRO_TEST_NATIVE_CREDENTIAL_BOUNDARY"; + if std::env::var_os(CHILD).is_some() { + // This integration binary links local-host without cfg(test), just as + // TUI tests do. Force the real public storage path without the runtime + // disable flag so the compiled test-support boundary must reject it. + let error = match maestro_local_host::init_cli::load_evalops_snapshot() { + Err(error) => error, + Ok(_) => panic!("test-support dependency opened native credential storage"), + }; + assert!(format!("{error:#}").contains("disabled in unit tests and test-support builds")); + return; + } + let home = tempfile::tempdir().expect("isolated credential home"); + let mut child = std::process::Command::new(std::env::current_exe().unwrap()) + .args([ + "--exact", + "test_support_dependency_cannot_open_native_credentials", + "--nocapture", + ]) + .env(CHILD, "1") + .env("MAESTRO_HOME", home.path()) + .env("MAESTRO_OAUTH_STORAGE_MODE", "keychain") + .env_remove("MAESTRO_DISABLE_KEYCHAIN") + .spawn() + .expect("spawn isolated credential boundary fixture"); + let deadline = std::time::Instant::now() + Duration::from_secs(5); + loop { + if let Some(status) = child + .try_wait() + .expect("wait for credential boundary fixture") + { + assert!( + status.success(), + "credential boundary fixture failed: {status}" + ); + break; + } + if std::time::Instant::now() >= deadline { + child + .kill() + .expect("stop stalled credential boundary fixture"); + child.wait().expect("reap credential boundary fixture"); + panic!("credential boundary fixture did not reject native storage promptly"); + } + std::thread::sleep(Duration::from_millis(10)); + } +} + async fn next_event_matching( session: &mut EmbeddedAgentSession, predicate: impl Fn(&FromAgent) -> bool, diff --git a/packages/tui-rs/src/app/tests.rs b/packages/tui-rs/src/app/tests.rs index ecf93b630..d2ad53d52 100644 --- a/packages/tui-rs/src/app/tests.rs +++ b/packages/tui-rs/src/app/tests.rs @@ -6449,9 +6449,49 @@ async fn interactive_resume_restores_exact_provider_history() { assert_session_restore_provider_history(true).await; } +struct RestoreEnv(Vec<(String, Option)>); +impl Drop for RestoreEnv { + fn drop(&mut self) { + for (name, value) in &self.0 { + match value { + Some(value) => std::env::set_var(name, value), + None => std::env::remove_var(name), + } + } + } +} + async fn assert_session_restore_provider_history(interactive: bool) { use crate::agent::{NativeAgent, NativeAgentConfig}; + let _env_lock = crate::config::test_process_env_lock_async().await; let temp = tempfile::tempdir().unwrap(); + let names = crate::credential_mode::TEST_IDENTITY_ENV_VARS + .iter() + .copied() + .chain([ + "MAESTRO_HOME", + "MAESTRO_OAUTH_STORAGE_MODE", + "MAESTRO_DISABLE_KEYCHAIN", + "OPENAI_API_KEY", + "OPENAI_BASE_URL", + ]); + let restore = RestoreEnv( + names + .map(|name| (name.to_string(), std::env::var_os(name))) + .collect(), + ); + for (name, _) in &restore.0 { + std::env::remove_var(name); + } + let _restore = restore; + // Resuming also resolves Identity credentials; an injected provider client + // does not isolate that owner or the host's native credential store. + std::env::set_var("MAESTRO_HOME", temp.path()); + std::env::set_var("MAESTRO_OAUTH_STORAGE_MODE", "file"); + std::env::set_var("MAESTRO_DISABLE_KEYCHAIN", "1"); + crate::credential_mode::install_test_identity_env(); + std::env::set_var("OPENAI_API_KEY", "fixture"); + std::env::set_var("OPENAI_BASE_URL", "http://127.0.0.1:1/v1"); let mut app = new_test_app(); app.session_manager = SessionManager::with_sessions_dir("/tmp", temp.path()); app.current_model = "gpt-6-astra".into(); @@ -6520,17 +6560,6 @@ async fn resumed_continuation_survives_the_next_compaction() { use crate::agent::{NativeAgent, NativeAgentConfig}; let temp = tempfile::tempdir().unwrap(); let _env_lock = crate::config::test_process_env_lock_async().await; - struct RestoreEnv(Vec<(String, Option)>); - impl Drop for RestoreEnv { - fn drop(&mut self) { - for (name, value) in &self.0 { - match value { - Some(value) => std::env::set_var(name, value), - None => std::env::remove_var(name), - } - } - } - } let names = crate::credential_mode::TEST_IDENTITY_ENV_VARS .iter() .copied() diff --git a/vendor/dex-loop/src/budget.rs b/vendor/dex-loop/src/budget.rs index 5f85344db..5afc17694 100644 --- a/vendor/dex-loop/src/budget.rs +++ b/vendor/dex-loop/src/budget.rs @@ -34,6 +34,26 @@ impl Default for Budget { } } +/// Capacity for the current model attempt. Advisory only: the engine still +/// enforces every limit. Unlimited policy axes are represented by `None`. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct RemainingBudget { + pub tool_steps: u32, + pub tokens: Option, + pub cost_micros: Option, + pub wall_ms: u64, + pub answer_only: bool, +} + +impl RemainingBudget { + pub fn guidance(&self) -> String { + format!( + "Turn capacity for this attempt: tool steps remaining (including this attempt)={}, tokens remaining={:?}, cost micros remaining={:?}, time remaining={} ms, answer_only={}. None means no policy cap. Reserve capacity for a truthful final answer. Prioritize the user's remaining work; distinguish verified completion, partial progress, and blockers. When answer_only=true, use no tools and report only established results.", + self.tool_steps, self.tokens, self.cost_micros, self.wall_ms, self.answer_only, + ) + } +} + /// The limit a turn ran out of. #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub enum BudgetAxis { @@ -55,6 +75,24 @@ impl fmt::Display for BudgetAxis { } impl Budget { + pub fn remaining( + &self, + completed_steps: u32, + usage: Usage, + elapsed: Duration, + ) -> RemainingBudget { + RemainingBudget { + tool_steps: self.max_steps.saturating_sub(completed_steps), + tokens: (self.max_tokens != u64::MAX) + .then(|| self.max_tokens.saturating_sub(usage.tokens())), + cost_micros: (self.max_cost_micros != u64::MAX) + .then(|| self.max_cost_micros.saturating_sub(usage.cost_micros)), + wall_ms: u64::try_from(self.wall.saturating_sub(elapsed).as_millis()) + .unwrap_or(u64::MAX), + answer_only: self.answer_only(completed_steps), + } + } + /// Whether the next model call is the answer-only call: `steps` calls /// have already been made and the tool-offering allowance is spent. pub fn answer_only(&self, steps: u32) -> bool { diff --git a/vendor/dex-loop/src/context.rs b/vendor/dex-loop/src/context.rs index ef4eb6992..cda1eec57 100644 --- a/vendor/dex-loop/src/context.rs +++ b/vendor/dex-loop/src/context.rs @@ -136,6 +136,28 @@ pub(crate) struct OpenStep { pub(crate) states: Vec, } +/// The most recent consecutive failed operation, derived from tool events. +#[derive(Clone, Debug, PartialEq)] +struct FailedCall { + principal: PrincipalId, + tool: ToolName, + args: serde_json::Value, + count: u32, +} + +/// The latest owner tool completions retained independently of model history. +/// Proposals and results are paired only by the typed event fold. +#[derive(Clone, Debug, PartialEq)] +pub struct ToolEvidence { + pub cursor: Cursor, + pub call: ProposedCall, + pub result: ToolResult, +} + +/// Evidence retention is bounded. Evicted results no longer authorize follow-up +/// verification; callers must run the tool again to establish fresh evidence. +pub const TOOL_EVIDENCE_LIMIT: usize = 128; + /// 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; @@ -159,6 +181,7 @@ pub struct Context { acting: Option, status: Status, history: Vec, + tool_evidence: Vec, // Newest input first (even when it has no uploads), then bounded earlier // attachment batches. Compaction never manufactures or edits provenance. attachment_inputs: Vec, @@ -198,6 +221,11 @@ pub struct Context { /// covers to). Assistant entries at or before it were produced before /// the summary existed. last_compaction: Option, + /// Derived from tool events, outside compactable history. This tracks a + /// consecutive failure, not a permanent blacklist or an automatic retry. + failed_call: Option, + /// Set only on the model-attempt copy, never on durable thread context. + remaining_budget: Option, } #[derive(Clone, Debug, PartialEq)] @@ -223,6 +251,7 @@ impl Context { acting: None, status: Status::Idle, history: Vec::new(), + tool_evidence: Vec::new(), attachment_inputs: Vec::new(), step: 0, usage: Usage::default(), @@ -243,9 +272,17 @@ impl Context { pending_turns: Vec::new(), pre_started: Vec::new(), last_compaction: None, + failed_call: None, + remaining_budget: None, } } + /// Owner results survive summaries; model-authored summary text is never + /// folded into this evidence. The records remain in event cursor order. + pub fn tool_evidence(&self) -> &[ToolEvidence] { + &self.tool_evidence + } + /// Assistant entries created before this event cannot reuse signed thinking. pub fn last_compaction_cursor(&self) -> Option { self.last_compaction @@ -295,6 +332,34 @@ impl Context { self.approval_mode } + /// Advisory capacity, refreshed before every model stream. It does not + /// grant permission or change the engine's caps. + pub fn remaining_budget(&self) -> Option { + self.remaining_budget + } + + pub(crate) fn for_model(&self, remaining: crate::RemainingBudget) -> Self { + let mut ctx = self.clone(); + ctx.remaining_budget = Some(remaining); + ctx + } + + /// A fresh identical failed call needs new information or a new approach. + pub fn recovery_guidance(&self) -> Option { + self.failed_call.as_ref().filter(|failure| failure.count >= 3).map(|failure| { + format!("The same call to {} failed {} consecutive times without progress. Do not repeat the identical call. Inspect the failure, change the inputs or approach, use another available capability, or report the concrete blocker. A successful intervening call or new user input resets this guard. Never retry a mutation with an unknown outcome.", failure.tool, failure.count) + }) + } + + pub(crate) fn has_stalled_call(&self, call: &ProposedCall) -> bool { + self.failed_call.as_ref().is_some_and(|failure| { + failure.count >= 3 + && failure.principal == call.principal + && failure.tool == call.tool + && failure.args == call.args + }) + } + /// The model's view of the thread. pub fn history(&self) -> &[Entry] { &self.history @@ -682,6 +747,29 @@ impl Context { output, receipt, } => { + if let Some(proposal) = self + .open_step + .as_ref() + .and_then(|step| step.calls.iter().find(|proposal| &proposal.id == call)) + .cloned() + && !self + .tool_evidence + .iter() + .any(|record| record.call.id == *call) + { + self.tool_evidence.push(ToolEvidence { + cursor, + call: proposal, + result: ToolResult { + outcome: *outcome, + output: output.clone(), + receipt: receipt.clone(), + }, + }); + if self.tool_evidence.len() > TOOL_EVIDENCE_LIMIT { + self.tool_evidence.remove(0); + } + } // 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| { @@ -737,6 +825,31 @@ impl Context { { self.uncertain_calls.push(proposal.clone()); } + if let Some(proposal) = self + .open_step + .as_ref() + .and_then(|step| step.calls.iter().find(|proposal| &proposal.id == call)) + { + if *outcome == Outcome::Failed { + let count = self + .failed_call + .as_ref() + .filter(|failure| { + failure.principal == proposal.principal + && failure.tool == proposal.tool + && failure.args == proposal.args + }) + .map_or(1, |failure| failure.count.saturating_add(1)); + self.failed_call = Some(FailedCall { + principal: proposal.principal.clone(), + tool: proposal.tool.clone(), + args: proposal.args.clone(), + count, + }); + } else { + self.failed_call = None; + } + } if let Some(state) = self.state_mut(call) { *state = CallState::Done(ToolResult { outcome: *outcome, @@ -880,6 +993,7 @@ impl Context { self.interrupt_requested = false; self.exposed.clear(); self.uncertain_calls.clear(); + self.failed_call = None; self.pre_started.clear(); self.client_tools = client_tools; self.authorized_tools = authorized_tools; @@ -971,6 +1085,7 @@ impl Context { .turn .clone() .unwrap_or_else(|| TurnId::new(String::new())); + self.failed_call = None; for (_, principal, text) in ready { self.acting = Some(principal.clone()); self.push( diff --git a/vendor/dex-loop/src/engine.rs b/vendor/dex-loop/src/engine.rs index c1c157281..6c16d728f 100644 --- a/vendor/dex-loop/src/engine.rs +++ b/vendor/dex-loop/src/engine.rs @@ -20,9 +20,10 @@ use std::future::Future; use std::pin::{Pin, pin}; use std::sync::{Arc, Mutex, PoisonError}; use std::task::Poll; -use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; +use std::time::{Duration, SystemTime, UNIX_EPOCH}; use futures_util::StreamExt; +use tokio::time::Instant; use tokio_util::sync::CancellationToken; use crate::budget::{Budget, BudgetAxis}; @@ -382,6 +383,18 @@ where .min(self.budget.wall.saturating_sub(run_started.elapsed())) } + /// Anchor both the remaining wall time and the timer to one clock sample. + /// Sampling the timer before a later elapsed-time read can expire early + /// if the task is descheduled between those reads. + fn call_expires_at(&self, run_started: Instant) -> Instant { + let now = Instant::now(); + let remaining_wall = self + .budget + .wall + .saturating_sub(now.saturating_duration_since(run_started)); + now + self.tool_call_deadline.min(remaining_wall) + } + /// Runs the current turn from wherever `ctx` stands until it finishes, /// parks, or is interrupted. `ctx` is usually fresh from `rehydrate`; the /// same call resumes after a restart, an approval, or an answer. @@ -563,7 +576,12 @@ where self.offered(ctx) }; let offered: Vec<&ToolSpec> = owned.iter().collect(); - let mut stream = pin!(self.model.stream(ctx, &offered)); + let model_ctx = ctx.for_model(self.budget.remaining( + step.saturating_sub(1), + ctx.usage(), + started.elapsed(), + )); + let mut stream = pin!(self.model.stream(&model_ctx, &offered)); // Dropping the stream on cancel is safe: a model call has no // effects to wait for. A stream that never yields another chunk // and is never cancelled would otherwise hang here forever, past @@ -827,11 +845,12 @@ where && !ctx.client_tools().iter().any(|tool| tool.name == call.tool) && !matches!(spec.executor, ExecutorKind::User | ExecutorKind::Client) && validate_args(&spec, &call.args).is_ok() - && !ctx.has_uncertain_call(call); + && !ctx.has_uncertain_call(call) + && !ctx.has_stalled_call(call); if !eligible { return Ok(()); } - let deadline = tokio::time::Instant::now() + self.call_deadline(run_started); + let deadline = self.call_expires_at(run_started); let verdict = tokio::select! { biased; () = cancel.cancelled() => return Ok(()), @@ -920,6 +939,7 @@ where .open_step() .and_then(|step| step.states.get(index)) .cloned(); + let already_started = matches!(&state, Some(CallState::Started)); let decision = match state { None | Some(CallState::Done(_)) => continue, Some(CallState::Asked { answer: None }) => { @@ -1010,6 +1030,15 @@ where Some(CallState::Todo) => None, }; + // A historical Started mutation must settle through the ledger + // above. This guard applies only before a fresh call can run. + if !already_started && ctx.has_stalled_call(call) { + prefetch.reads.discard(index); + self.finish(ctx, call, ToolResult::error( + "not executed: the identical call failed three times without progress; change the inputs or approach, or report the blocker", + )).await?; + continue; + } if call.tool.as_str() == TOOLS_SEARCH { self.search_tools(ctx, call).await?; continue; @@ -1405,7 +1434,7 @@ where Some(started(call, &spec)) }) .collect(); - let deadline = tokio::time::Instant::now() + self.call_deadline(run_started); + let deadline = self.call_expires_at(run_started); self.emit(ctx, starts).await?; let thread = ctx.thread().clone(); // One deadline for the wave: its reads run concurrently, so each @@ -1533,8 +1562,8 @@ where // under the claim, so a later resume of this call adopts it // instead of dispatching the mutation a second time. Interrupts // reach cancellable process tools; settlement is still awaited. - let deadline = self.call_deadline(run_started); - let result = if deadline.is_zero() { + let deadline = self.call_expires_at(run_started); + let result = if Instant::now() >= deadline { if already_started { ToolResult::unknown(UNKNOWN_NO_RETRY) } else { @@ -1542,7 +1571,7 @@ where } } else { let run = self.tools.run(ctx.thread(), call, cancel); - match tokio::time::timeout(deadline, run).await { + match tokio::time::timeout_at(deadline, run).await { Ok(result) => result, Err(_elapsed) => ToolResult::unknown(DEADLINE_MUTATION), } diff --git a/vendor/dex-loop/src/event.rs b/vendor/dex-loop/src/event.rs index b5452f069..b9739d3c7 100644 --- a/vendor/dex-loop/src/event.rs +++ b/vendor/dex-loop/src/event.rs @@ -486,12 +486,9 @@ impl ErrorClass { } } -/// Who can answer a turn's approval requests. -/// -/// `Headless` turns come from callers with no human to click Approve (service -/// accounts, workloads, agents, synthetic canaries, API automation). For them -/// the engine resolves a policy `NeedsApproval` verdict itself, recording the -/// request and the decision in the log. A `Deny` verdict is never affected. +/// 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`. #[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum ApprovalMode { @@ -866,64 +863,6 @@ mod tests { } } - #[test] - fn auto_approval_receipts_decode_losslessly_and_never_become_control_events() { - let event = Event::AutoApproved { - call: CallId::new("t1-1-0"), - approval: ApprovalId::new("ap-1"), - args_digest: args_digest(&serde_json::json!({"key": "w"})), - summary: "Send\0email".into(), - principal: PrincipalId::new(AUTO_APPROVER), - }; - let exact = serde_json::to_string(&event).expect("serialize receipt"); - let envelope = serde_json::json!({ - "type": Event::STORED_JSON_V1_TYPE, - "_dex_event_json_v1": exact, - "summary": "non-authoritative projection", - }); - assert_eq!( - Event::from_stored_json(&envelope, Some("auto_approved")).expect("exact receipt"), - event - ); - assert!(Event::from_stored_json(&envelope, Some("approval_decided")).is_err()); - assert!( - !event.is_control(), - "engine receipts must not advance the host control cursor" - ); - } - - #[test] - fn stored_event_decoder_refuses_corrupt_envelopes_and_call_digests() { - let event = Event::ModelStepCompleted { - step: 1, - text: String::new(), - calls: vec![ProposedCall::new( - CallId::new("turn-1-0"), - ToolName::new("dex.report_feedback"), - serde_json::json!({"diagnosis": "exact\0value"}), - PrincipalId::new("alice"), - )], - reasoning: None, - served: None, - timing: None, - }; - let mut envelope = serde_json::json!({ - "type": Event::STORED_JSON_V1_TYPE, - "_dex_event_json_v1": serde_json::to_string(&event).unwrap(), - "calls": [], - }); - assert_eq!( - Event::from_stored_json(&envelope, Some("model_step_completed")).unwrap(), - event - ); - assert!(Event::from_stored_json(&envelope, Some("tool_started")).is_err()); - envelope[Event::STORED_JSON_V1_KEY] = serde_json::json!("invalid JSON"); - assert!(Event::from_stored_json(&envelope, Some("model_step_completed")).is_err()); - let mut tampered = serde_json::to_value(&event).unwrap(); - tampered["calls"][0]["args"]["diagnosis"] = serde_json::json!("changed"); - assert!(Event::from_stored_json(&tampered, Some("model_step_completed")).is_err()); - } - #[test] fn events_round_trip_through_json() { let events = vec![ @@ -1035,32 +974,6 @@ mod tests { } } - #[test] - fn approval_mode_defaults_to_interactive_for_old_rows_and_round_trips() { - let old_row = serde_json::json!({ - "type": "user_message", - "turn": "t1", - "principal": "alice", - "text": "hi", - "attachments": [], - }); - match serde_json::from_value::(old_row).expect("old row deserializes") { - Event::UserMessage { approval_mode, .. } => { - assert_eq!(approval_mode, ApprovalMode::Interactive); - } - other => panic!("expected UserMessage, got {other:?}"), - } - assert_eq!( - serde_json::to_value(ApprovalMode::Headless).expect("serialize"), - serde_json::json!("headless") - ); - assert_eq!( - serde_json::from_value::(serde_json::json!("interactive")) - .expect("deserialize"), - ApprovalMode::Interactive - ); - } - #[test] fn user_message_without_message_id_round_trips_to_none() { let event = Event::UserMessage { diff --git a/vendor/dex-loop/src/lib.rs b/vendor/dex-loop/src/lib.rs index 53a7bf8b4..9845a219b 100644 --- a/vendor/dex-loop/src/lib.rs +++ b/vendor/dex-loop/src/lib.rs @@ -28,11 +28,14 @@ mod ports; mod rehydrate; mod sanitize; -pub use budget::{Budget, BudgetAxis}; +pub use budget::{Budget, BudgetAxis, RemainingBudget}; pub use compaction::{ Compaction, CompactionPlan, Compactor, NoCompaction, Summarize, Summary, Threshold, }; -pub use context::{AttachmentInput, Context, Entry, MAX_CONTEXT_ATTACHMENTS, Message}; +pub use context::{ + AttachmentInput, Context, Entry, MAX_CONTEXT_ATTACHMENTS, Message, TOOL_EVIDENCE_LIMIT, + ToolEvidence, +}; pub use engine::{CUT_OFF_NOTICE, DEFAULT_TOOL_CALL_DEADLINE, Engine, Exit, TOOLS_SEARCH}; pub use event::{ AUTO_APPROVER, ActionConfirmation, ApprovalId, ApprovalMode, ArtifactRef, AttemptNext, CallId, diff --git a/vendor/dex-loop/tests/budget_contract.rs b/vendor/dex-loop/tests/budget_contract.rs new file mode 100644 index 000000000..287d87cc0 --- /dev/null +++ b/vendor/dex-loop/tests/budget_contract.rs @@ -0,0 +1,102 @@ +use std::time::Duration; + +use dex_loop::{Budget, BudgetAxis, RemainingBudget, Usage}; + +#[test] +fn each_axis_is_reported() { + let budget = Budget { + max_steps: 2, + max_tokens: 100, + max_cost_micros: 50, + wall: Duration::from_secs(1), + }; + let usage = |tokens, cost| Usage { + input_tokens: tokens, + output_tokens: 0, + cost_micros: cost, + ..Usage::default() + }; + let short = Duration::from_millis(1); + assert_eq!(budget.exhausted(1, usage(10, 1), short), None); + // The answer-only call after `max_steps` is still allowed. + assert_eq!(budget.exhausted(2, usage(10, 1), short), None); + assert!(!budget.answer_only(1)); + assert!(budget.answer_only(2)); + assert_eq!( + budget.exhausted(3, usage(10, 1), short), + Some(BudgetAxis::Steps) + ); + assert_eq!( + budget.exhausted(1, usage(100, 1), short), + Some(BudgetAxis::Tokens) + ); + assert_eq!( + budget.exhausted(1, usage(10, 50), short), + Some(BudgetAxis::Cost) + ); + assert_eq!( + budget.exhausted(1, usage(10, 1), Duration::from_secs(1)), + Some(BudgetAxis::Wall) + ); +} + +#[test] +fn remaining_is_saturating_and_unlimited_axes_are_explicit() { + let budget = Budget { + max_steps: 2, + max_tokens: 100, + max_cost_micros: 50, + wall: Duration::from_millis(100), + }; + let remaining = budget.remaining( + 1, + Usage { + input_tokens: 20, + output_tokens: 10, + cost_micros: 7, + ..Usage::default() + }, + Duration::from_millis(25), + ); + assert_eq!( + remaining, + RemainingBudget { + tool_steps: 1, + tokens: Some(70), + cost_micros: Some(43), + wall_ms: 75, + answer_only: false + } + ); + let spent = budget.remaining( + 3, + Usage { + input_tokens: 200, + cost_micros: 90, + ..Usage::default() + }, + Duration::from_secs(1), + ); + assert_eq!( + spent, + RemainingBudget { + tool_steps: 0, + tokens: Some(0), + cost_micros: Some(0), + wall_ms: 0, + answer_only: true + } + ); + assert_eq!( + Budget::default() + .remaining(0, Usage::default(), Duration::ZERO) + .tokens, + None + ); + assert_eq!( + Budget::default() + .remaining(0, Usage::default(), Duration::ZERO) + .cost_micros, + None + ); +} diff --git a/vendor/dex-loop/tests/client_tool_replay.rs b/vendor/dex-loop/tests/client_tool_replay.rs index 44c9542af..294568e73 100644 --- a/vendor/dex-loop/tests/client_tool_replay.rs +++ b/vendor/dex-loop/tests/client_tool_replay.rs @@ -387,3 +387,56 @@ async fn interrupting_a_turn_parked_on_a_client_tool_ends_it_promptly() { ); assert_eq!(log.rehydrate(), ctx); } + +#[tokio::test] +async fn configured_client_timeout_survives_compaction_and_is_persisted_for_restart() { + let log = FakeLog::default(); + let model = FakeModel::new(vec![vec![call("browser.read_tab", json!({}))]]); + let tools = FakeTools::new(vec![client_executed_tool("browser.read_tab", true)]); + let configured = Duration::from_secs(20); + let engine = engine(&log, &model, &tools, budget()) + .with_client_timeout(configured) + .with_compactor(dex_loop::Threshold::new(usize::MAX, 1, FakeSummarizer)); + let mut ctx = log.start_turn_with_client_tools( + "t1", + "read the tab", + vec![client_tool("browser.read_tab", true)], + ); + let before = now_ms(); + let exit = engine.run(&mut ctx, &CancellationToken::new()).await; + let after = now_ms(); + assert_eq!(exit, Ok(Exit::AwaitingClientTool(call_id("t1", 1, 0)))); + let deadline = log + .events() + .into_iter() + .find_map(|event| match event { + Event::ClientToolRequested { deadline_ms, .. } => Some(deadline_ms), + _ => None, + }) + .expect("client request has a durable deadline"); + assert!( + deadline >= before + 20_000 && deadline <= after + 20_000, + "configured wait must replace the 120-second default: {deadline}, {before}, {after}" + ); + let restarted_model = FakeModel::new(vec![]); + let restarted = support::engine(&log, &restarted_model, &tools, budget()); + let mut replay = log.rehydrate(); + assert_eq!( + restarted.run(&mut replay, &CancellationToken::new()).await, + Ok(Exit::AwaitingClientTool(call_id("t1", 1, 0))) + ); + let replay_deadlines: Vec<_> = log + .events() + .into_iter() + .filter_map(|event| match event { + Event::ClientToolRequested { deadline_ms, .. } => Some(deadline_ms), + _ => None, + }) + .collect(); + assert_eq!( + replay_deadlines, + vec![deadline], + "restart cannot extend the configured wait" + ); + assert_eq!(replay, log.rehydrate()); +} diff --git a/vendor/dex-loop/tests/evidence_contract.rs b/vendor/dex-loop/tests/evidence_contract.rs new file mode 100644 index 000000000..21384a8e9 --- /dev/null +++ b/vendor/dex-loop/tests/evidence_contract.rs @@ -0,0 +1,104 @@ +use dex_loop::{ + ApprovalMode, CallId, Cursor, Event, Message, Outcome, Output, PrincipalId, ProposedCall, + TOOL_EVIDENCE_LIMIT, ThreadId, ToolName, TurnId, rehydrate, +}; + +#[test] +fn typed_owner_evidence_is_bounded_and_survives_compaction_and_full_replay() { + let thread = ThreadId { + org: "org".into(), + workspace: "ws".into(), + thread: "thread".into(), + }; + let mut events = Vec::new(); + fn push(events: &mut Vec<(Cursor, Event)>, event: Event) { + events.push((Cursor(events.len() as i64 + 1), event)); + } + // A real thread starts from accepted input before owner completions. + push( + &mut events, + Event::UserMessage { + turn: TurnId::new("turn-1"), + message_id: None, + principal: PrincipalId::new("alice"), + text: "Read owner results".into(), + attachments: vec![], + client_tools: vec![], + authorized_tools: vec![], + approval_mode: ApprovalMode::Interactive, + }, + ); + for index in 0..=TOOL_EVIDENCE_LIMIT { + let call = ProposedCall::new( + CallId::new(format!("call-{index}")), + ToolName::new("owner.read"), + serde_json::json!({"id":index}), + PrincipalId::new("alice"), + ); + push( + &mut events, + Event::ModelStepCompleted { + step: index as u32, + text: String::new(), + calls: vec![call.clone()], + reasoning: None, + served: None, + timing: None, + }, + ); + push( + &mut events, + Event::ToolFinished { + call: call.id, + outcome: Outcome::Succeeded, + output: Output::Text(format!("owner-result-{index}")), + receipt: None, + }, + ); + } + // A turn-boundary compaction covers completed history. The active + // turn's exact input remains protected until its terminal event. + push( + &mut events, + Event::Final { + text: String::new(), + }, + ); + let covers_to = Cursor(events.len() as i64); + push( + &mut events, + Event::Compaction { + covers_to_cursor: covers_to, + summary: "model claims any call passed".into(), + }, + ); + // Production first rehydrates the accepted input, establishing the + // control replay floor. UserMessage itself is not a control event. + // Then its warm actor observes subsequent owner events and compaction. + let mut warm = rehydrate(thread.clone(), &events[..1]); + for (cursor, event) in &events[1..] { + warm.observe(*cursor, event); + } + let restarted = rehydrate(thread, &events); + assert_eq!(warm, restarted); + assert_eq!(warm.tool_evidence().len(), TOOL_EVIDENCE_LIMIT); + assert_eq!(warm.tool_evidence()[0].call.id.as_str(), "call-1"); + assert_eq!(warm.tool_evidence()[0].cursor, Cursor(5)); + assert_eq!(warm.tool_evidence().last().unwrap().cursor, Cursor(259)); + assert!( + warm.history() + .iter() + .all(|entry| matches!(entry.message, Message::Summary { .. })) + ); + // Unmatched completions and summary prose cannot create owner evidence. + warm.observe( + Cursor(covers_to.0 + 2), + &Event::ToolFinished { + call: CallId::new("forged"), + outcome: Outcome::Succeeded, + output: Output::Text("passed".into()), + receipt: None, + }, + ); + assert_eq!(warm.tool_evidence(), restarted.tool_evidence()); +} diff --git a/vendor/dex-loop/tests/loop_recovery.rs b/vendor/dex-loop/tests/loop_recovery.rs new file mode 100644 index 000000000..0a8b59dfd --- /dev/null +++ b/vendor/dex-loop/tests/loop_recovery.rs @@ -0,0 +1,274 @@ +//! Budget visibility and recovery exercise the real Engine and model boundary. +// Shared fixtures include helpers used by other integration-test binaries. +#[allow(dead_code)] +mod support; + +use dex_loop::{ + Budget, CancellationToken, Context, Engine, Event, Exit, Lexicon, Model, ModelChunk, + ModelError, PrincipalId, ProposedCall, RemainingBudget, ThreadId, ToolName, ToolResult, + ToolSpec, Tools, Verdict, +}; +use futures_util::{Stream, stream}; +use serde_json::json; +use std::sync::{Arc, Mutex}; +use std::time::Duration; +use support::{FakeEffects, FakeLog, read_tool}; + +type ModelObservation = (RemainingBudget, Option); + +#[derive(Clone, Default)] +struct ObservingModel { + seen: Arc>>, + stubborn: bool, +} +impl Model for ObservingModel { + fn stream<'a>( + &'a self, + ctx: &'a Context, + _: &'a [&'a ToolSpec], + ) -> impl Stream> + Send + 'a { + let remaining = ctx + .remaining_budget() + .expect("every engine attempt has fresh capacity"); + let recovery = ctx.recovery_guidance(); + self.seen + .lock() + .expect("seen") + .push((remaining, recovery.clone())); + let chunk = if (recovery.is_some() || remaining.answer_only) && !self.stubborn { + ModelChunk::Text("Blocked: the read failed; I could not verify completion.".into()) + } else { + ModelChunk::ToolCall { + name: ToolName::new("read"), + args: json!({"key":"same"}), + } + }; + stream::iter([Ok(chunk)]) + } +} +#[derive(Clone)] +struct FailingTools { + spec: Vec, + calls: Arc>, +} +impl FailingTools { + fn new() -> Self { + Self { + spec: vec![read_tool("read")], + calls: Arc::default(), + } + } +} +impl Tools for FailingTools { + fn catalog(&self) -> &[ToolSpec] { + &self.spec + } + async fn search(&self, _: &PrincipalId, _: &str) -> Vec { + vec![] + } + async fn policy(&self, _: &Context, _: &ProposedCall) -> Verdict { + Verdict::Allow + } + async fn run(&self, _: &ThreadId, _: &ProposedCall, _: &CancellationToken) -> ToolResult { + *self.calls.lock().expect("calls") += 1; + ToolResult::error("recorded read failure") + } +} +fn budget(steps: u32) -> Budget { + Budget { + max_steps: steps, + max_tokens: 1234, + max_cost_micros: 5678, + wall: Duration::from_secs(10), + } +} + +#[tokio::test] +async fn three_failures_offer_recovery_before_spending_the_whole_budget() { + let log = FakeLog::default(); + let model = ObservingModel::default(); + let tools = FailingTools::new(); + let engine = Engine::new( + log.clone(), + model.clone(), + tools.clone(), + FakeEffects::default(), + Lexicon::default(), + budget(20), + ); + let mut ctx = log.start_turn("t1", "verify the record"); + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) + ); + assert_eq!(*tools.calls.lock().expect("calls"), 3); + let seen = model.seen.lock().expect("seen"); + assert_eq!(seen.len(), 4); + assert_eq!(seen[0].0.tool_steps, 20); + assert_eq!(seen[3].0.tool_steps, 17); + assert_eq!(seen[0].0.tokens, Some(1234)); + assert_eq!(seen[0].0.cost_micros, Some(5678)); + assert!( + seen[3] + .1 + .as_ref() + .is_some_and(|text| text.contains("failed 3")) + ); + assert_eq!( + ctx, + log.rehydrate(), + "advisory attempt state must not pollute durable context" + ); +} + +#[tokio::test] +async fn stubborn_model_cannot_redispatch_stalled_reads_or_cross_answer_only_cap() { + let log = FakeLog::default(); + let model = ObservingModel { + stubborn: true, + ..Default::default() + }; + let tools = FailingTools::new(); + let engine = Engine::new( + log.clone(), + model.clone(), + tools.clone(), + FakeEffects::default(), + Lexicon::default(), + budget(5), + ); + let mut ctx = log.start_turn("t1", "read"); + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Failed) + ); + assert_eq!(*tools.calls.lock().expect("calls"), 3); + let seen = model.seen.lock().expect("seen"); + assert_eq!(seen.len(), 6); + assert!(seen[5].0.answer_only); + assert_eq!(seen[5].0.tool_steps, 0); + assert_eq!(ctx, log.rehydrate()); +} + +#[tokio::test] +async fn answer_only_remains_available_and_zero_tool_budget_is_visible() { + let log = FakeLog::default(); + let model = ObservingModel::default(); + let tools = FailingTools::new(); + let engine = Engine::new( + log.clone(), + model.clone(), + tools.clone(), + FakeEffects::default(), + Lexicon::default(), + budget(0), + ); + let mut ctx = log.start_turn("t1", "report existing evidence"); + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) + ); + assert_eq!(*tools.calls.lock().expect("calls"), 0); + assert!(model.seen.lock().expect("seen")[0].0.answer_only); +} + +#[test] +fn failure_recovery_survives_restart_and_compaction_but_new_input_resets_it() { + let log = FakeLog::default(); + let mut ctx = log.start_turn("t1", "verify"); + for step in 1..=3 { + let call = ProposedCall::new( + dex_loop::CallId::new(format!("call{step}")), + ToolName::new("read"), + json!({"key":"same"}), + support::alice(), + ); + for event in [ + Event::StepStarted { + step, + control_through: ctx.control_cursor(), + }, + Event::ModelStepCompleted { + step, + text: String::new(), + calls: vec![call.clone()], + reasoning: None, + served: None, + timing: None, + }, + Event::ToolFinished { + call: call.id, + outcome: dex_loop::Outcome::Failed, + output: dex_loop::Output::Text("failed".into()), + receipt: None, + }, + ] { + let cursor = log.host_append(event.clone()); + ctx.observe(cursor, &event); + } + } + let event = Event::Compaction { + covers_to_cursor: ctx.cursor(), + summary: "reads failed".into(), + }; + let cursor = log.host_append(event.clone()); + ctx.observe(cursor, &event); + assert_eq!(ctx, log.rehydrate()); + assert!(ctx.recovery_guidance().is_some()); + let event = Event::Steer { + principal: support::alice(), + text: "try this newly corrected key".into(), + }; + let cursor = log.host_append(event.clone()); + ctx.observe(cursor, &event); + let event = Event::StepStarted { + step: 4, + control_through: cursor, + }; + let cursor = log.host_append(event.clone()); + ctx.observe(cursor, &event); + assert!(ctx.recovery_guidance().is_none()); + assert_eq!(ctx, log.rehydrate()); +} + +#[tokio::test] +async fn recovery_never_discards_results_of_reads_that_already_started() { + let log = FakeLog::default(); + let model = support::FakeModel::new(vec![ + (0..4) + .map(|_| support::call("read", json!({"key":"same"}))) + .collect(), + vec![support::text("Blocked: all four attempted reads failed.")], + ]); + let tools = FailingTools::new(); + let engine = Engine::new( + log.clone(), + model, + tools.clone(), + FakeEffects::default(), + Lexicon::default(), + budget(5), + ); + let mut ctx = log.start_turn("t1", "read"); + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) + ); + assert_eq!(*tools.calls.lock().expect("calls"), 4); + let finishes = log + .events() + .into_iter() + .filter_map(|event| { + if let Event::ToolFinished { output, .. } = event { + Some(output) + } else { + None + } + }) + .collect::>(); + assert_eq!(finishes.len(), 4); + assert!(finishes.iter().all( + |output| matches!(output, dex_loop::Output::Text(text) if text == "recorded read failure") + )); + assert_eq!(ctx, log.rehydrate()); +} diff --git a/vendor/dex-loop/tests/scenarios.rs b/vendor/dex-loop/tests/scenarios.rs index 1e1ad89aa..541981c44 100644 --- a/vendor/dex-loop/tests/scenarios.rs +++ b/vendor/dex-loop/tests/scenarios.rs @@ -497,8 +497,6 @@ async fn steer_from_another_principal_is_checked_under_that_principal() { "user:check prod", "step:1", "started:t1-1-0", - // The read is prefetched while the model streams, so the steer it - // triggers lands before the model step completes. "steer:also update staging", "completed::[t1-1-0]", "finished:t1-1-0:ok", @@ -1944,6 +1942,7 @@ async fn tools_search_exposes_schemas_for_the_next_step() { .search_result("crm", &["crm.lookup"]); let engine = engine(&log, &model, &tools, budget()); let mut ctx = log.start_turn("t1", "find acme in the crm"); + assert!(ctx.exposed_tools().is_empty(), "hidden before discovery"); assert_eq!( engine.run(&mut ctx, &CancellationToken::new()).await, @@ -1967,6 +1966,7 @@ async fn tools_search_exposes_schemas_for_the_next_step() { "final:found acme", ]) ); + assert_eq!(ctx.exposed_tools(), &[ToolName::new("crm.lookup")]); let offered = model.offered(); assert_eq!(offered[0], strings(&["tools.search", "search"])); assert_eq!( @@ -1999,11 +1999,17 @@ async fn exposed_tools_are_appended_in_exposure_order() { .search_result("crm", &["crm.zed", "crm.alpha"]); let engine = engine(&log, &model, &tools, budget()); let mut ctx = log.start_turn("t1", "find acme in the crm"); + assert!(ctx.exposed_tools().is_empty(), "hidden before discovery"); assert_eq!( engine.run(&mut ctx, &CancellationToken::new()).await, Ok(Exit::Done) ); + assert_eq!( + ctx.exposed_tools(), + &[ToolName::new("crm.zed"), ToolName::new("crm.alpha")] + ); + assert_eq!(log.rehydrate().exposed_tools(), ctx.exposed_tools()); let offered = model.offered(); assert_eq!(offered[0], strings(&["tools.search", "search"])); assert_eq!( diff --git a/vendor/dex-loop/tests/storage_contract.rs b/vendor/dex-loop/tests/storage_contract.rs new file mode 100644 index 000000000..33265c8d5 --- /dev/null +++ b/vendor/dex-loop/tests/storage_contract.rs @@ -0,0 +1,131 @@ +use dex_loop::{ + AUTO_APPROVER, ApprovalId, ApprovalMode, CallId, Event, HEADLESS_AUTO_APPROVER, PrincipalId, + ProposedCall, ToolName, args_digest, +}; + +#[test] +fn auto_approval_receipts_decode_losslessly_and_never_become_control_events() { + let event = Event::AutoApproved { + call: CallId::new("t1-1-0"), + approval: ApprovalId::new("ap-1"), + args_digest: args_digest(&serde_json::json!({"key": "w"})), + summary: "Send\0email".into(), + principal: PrincipalId::new(AUTO_APPROVER), + }; + let exact = serde_json::to_string(&event).expect("serialize receipt"); + let envelope = serde_json::json!({ + "type": Event::STORED_JSON_V1_TYPE, + "_dex_event_json_v1": exact, + "summary": "non-authoritative projection", + }); + assert_eq!( + Event::from_stored_json(&envelope, Some("auto_approved")).expect("exact receipt"), + event + ); + assert!(Event::from_stored_json(&envelope, Some("approval_decided")).is_err()); + assert!( + !event.is_control(), + "engine receipts must not advance the host control cursor" + ); +} + +#[test] +fn stored_event_decoder_refuses_corrupt_envelopes_and_call_digests() { + let event = Event::ModelStepCompleted { + step: 1, + text: String::new(), + calls: vec![ProposedCall::new( + CallId::new("turn-1-0"), + ToolName::new("dex.report_feedback"), + serde_json::json!({"diagnosis": "exact\0value"}), + PrincipalId::new("alice"), + )], + reasoning: None, + served: None, + timing: None, + }; + let mut envelope = serde_json::json!({ + "type": Event::STORED_JSON_V1_TYPE, + "_dex_event_json_v1": serde_json::to_string(&event).unwrap(), + "calls": [], + }); + assert_eq!( + Event::from_stored_json(&envelope, Some("model_step_completed")).unwrap(), + event + ); + assert!(Event::from_stored_json(&envelope, Some("tool_started")).is_err()); + envelope[Event::STORED_JSON_V1_KEY] = serde_json::json!("invalid JSON"); + assert!(Event::from_stored_json(&envelope, Some("model_step_completed")).is_err()); + let mut tampered = serde_json::to_value(&event).unwrap(); + tampered["calls"][0]["args"]["diagnosis"] = serde_json::json!("changed"); + assert!(Event::from_stored_json(&tampered, Some("model_step_completed")).is_err()); +} + +#[test] +fn approval_mode_defaults_to_interactive_for_old_rows_and_round_trips() { + let old_row = serde_json::json!({ + "type": "user_message", + "turn": "t1", + "principal": "alice", + "text": "hi", + "attachments": [], + }); + match serde_json::from_value::(old_row).expect("old row deserializes") { + Event::UserMessage { approval_mode, .. } => { + assert_eq!(approval_mode, ApprovalMode::Interactive); + } + other => panic!("expected UserMessage, got {other:?}"), + } + assert_eq!( + serde_json::to_value(ApprovalMode::Headless).expect("serialize"), + serde_json::json!("headless") + ); + assert_eq!( + serde_json::from_value::(serde_json::json!("interactive")) + .expect("deserialize"), + ApprovalMode::Interactive + ); +} + +#[test] +fn public_approval_labels_decode_as_the_same_recorded_turn_mode() { + for (mode, wire_name) in [ + (ApprovalMode::Interactive, "interactive"), + (ApprovalMode::Headless, "headless"), + ] { + let row = serde_json::json!({"type": "user_message", "turn": "t1", "principal": "alice", "text": "hi", "attachments": [], "approval_mode": mode.as_str()}); + assert_eq!(row["approval_mode"], wire_name); + let decoded = Event::from_stored_json(&row, Some("user_message")) + .expect("public wire label must decode"); + match decoded { + Event::UserMessage { approval_mode, .. } => assert_eq!(approval_mode, mode), + other => panic!("expected UserMessage, got {other:?}"), + } + } +} + +#[test] +fn legacy_headless_approval_rows_keep_the_principal_and_control_semantics() { + let digest = "a".repeat(64); + let row = serde_json::json!({ + "type": "approval_decided", "call": "t1-1-0", "approval": "ap-1", + "args_digest": digest, "approved": true, + "principal": "policy:headless_auto_approve", + }); + let decoded = Event::from_stored_json(&row, Some("approval_decided")).expect("legacy row"); + assert_eq!( + decoded, + Event::ApprovalDecided { + call: CallId::new("t1-1-0"), + approval: ApprovalId::new("ap-1"), + args_digest: digest, + approved: true, + principal: PrincipalId::new(HEADLESS_AUTO_APPROVER), + } + ); + assert!(decoded.is_control()); + assert_eq!( + serde_json::to_value(decoded).expect("legacy round trip"), + row + ); +} diff --git a/vendor/dex-loop/tests/tool_deadline.rs b/vendor/dex-loop/tests/tool_deadline.rs index 2e6cb0a3c..dab8d6aa6 100644 --- a/vendor/dex-loop/tests/tool_deadline.rs +++ b/vendor/dex-loop/tests/tool_deadline.rs @@ -168,10 +168,9 @@ async fn a_mutation_that_overruns_the_deadline_is_recorded_unknown() { assert_eq!(log.rehydrate(), ctx); } -// Real time, not paused: the loop's between-steps wall check reads a -// `std::time::Instant`, which tokio's paused clock does not advance. 100ms -// of real time is the whole cost; the stuck call is dropped, never awaited. -#[tokio::test(flavor = "current_thread")] +// The tool timer and between-step wall check must share the same clock. +// Advancing virtual time must exhaust the wall before another model step. +#[tokio::test(flavor = "current_thread", start_paused = true)] async fn a_call_never_outlives_the_wall_budget() { // The per-call deadline is generous; the wall budget is not. The call is // cut at the wall, its result appended, and the turn ends on the wall