From fa55c2ce2d66f1df6e16830668ef971c8e2da715 Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" Date: Wed, 30 Sep 2026 00:37:07 +0000 Subject: [PATCH] chore: sync public mirror from internal --- .repository-projection.json | 6 +- CHANGELOG.md | 40 ++++ Cargo.lock | 2 +- package-lock.json | 4 +- package.json | 2 +- packages/maestro-rs/Cargo.toml | 2 +- scripts/smoke-release-native-only.mjs | 22 +-- scripts/smoke-release-process.mjs | 30 +++ scripts/smoke-release-process.test.mjs | 47 +++++ vendor/dex-loop/src/context.rs | 24 +++ vendor/dex-loop/src/engine.rs | 54 +++++- vendor/dex-loop/src/lib.rs | 2 +- vendor/dex-loop/tests/scenarios.rs | 249 ++++++++++++++++++++++++- 13 files changed, 457 insertions(+), 27 deletions(-) create mode 100644 scripts/smoke-release-process.mjs create mode 100644 scripts/smoke-release-process.test.mjs diff --git a/.repository-projection.json b/.repository-projection.json index c6bc36bcf..0394a2a90 100644 --- a/.repository-projection.json +++ b/.repository-projection.json @@ -3,11 +3,11 @@ "projection": "deixic-code", "projectionSchemaVersion": 1, "sourceRepository": "dx-corp/mono", - "sourceSha": "e3696732f3a134d547cb39da1ae758c94843bc83", + "sourceSha": "5c0ef43f86d41311f868236ebbc050f14969e263", "destinationRepository": "dx-corp/code", - "priorProjectedBase": "c0f08c593b1f58f7802b018315b296e2a1268912", + "priorProjectedBase": "931b3d60176d650c886e2ec3802031569808b470", "definitionDigest": "82936441c776e3e8edb5d215a75007ec9714a233f489d460075d79d5ef5ba32f", "toolDigest": "c244d99199a7ae3eb8ff644a99462163c23b0bb6a83ef50af01efbdca0b81d04", - "contentDigest": "bfa2803b378a3e7f65c5998f3de8f85c5490f0116c8c9a650fc9254777d3e49a", + "contentDigest": "11f79bfdda6fbb50a86620b4528c84bbabc195dc16b6c97058a812bcfadeb395", "publicationEligible": true } diff --git a/CHANGELOG.md b/CHANGELOG.md index 080ad5293..d771eb7f2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -51,6 +51,46 @@ versioning when releases are cut. and keep scheduled public runs inert so public publishing stays downstream of the internal source-of-truth release. +## [0.10.111] - 2026-09-30 + +### Added + +- Make dex-journey, one-command user-journey harness with latency table (#11440). +- Serve Claude on Vertex AI as provider vertex-anthropic (#11427). +- Run Claude on Vertex through the Messages loop with effort and prompt caching (#11424). +- Add protected bounded public conversation fuzz (#11417). +- Dispatch a Dex-approved call without parking for its own resume (#11373). + +### Changed + +- Delete the test-only process runner and proactive investigation worker (#11437). +- Web e2e caller is a person; history stubs are not exposure (#11438). +- Remove retired computer.coding_* capabilities and dead coding-workspace CI (#11442). +- Run desktop-lane main verdicts on arc-main-heavy (#11439). +- Path-scope the Merlin, World and Computer iOS workflows at the trigger (#11435). +- Make the parked-turn index probe discriminate (#11431). +- Repair solution provenance and vocabulary (#11428). +- Repair governed CRM solution metadata (#11416). +- Tighten protected fuzz cancellation and receipt proof (#11420). +- Keep public mirror staging proof current (#11414). +- Restore Postgres runtime fixture coverage (#11412). +- Exercise public web turn and receipt ownership in fuzz (#11409). + +### Fixed + +- Bound Rosetta release smoke with useful diagnostics. +- Keep each turn's context in its own user message for Claude (#11447). +- Reach Claude on the Vertex us/eu multi-region endpoints (#11448). +- Declare the GetOperatingReceipt codec once (#11436). +- Derive PartialEq for WorkloadProviderRef (#11441). +- Keep Dex final visible until message hydration (#11434). +- Keep a cut-off answer instead of withdrawing it (#11433). +- Recall team memory once per turn, not once per model step (#11432). +- Resolve public owner receipt details before fuzz proof (#11426). +- Forward long Dex answers over Watch (#11430). +- Read one-loop receipt details behind thread summaries (#11425). +- Bound provider streams by silence, not total length (#11429). + ## [0.10.110] - 2026-09-29 ### Added diff --git a/Cargo.lock b/Cargo.lock index 530c81c89..cec8d0042 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3752,7 +3752,7 @@ checksum = "dae608c151f68243f2b000364e1f7b186d9c29845f7d2d85bd31b9ad77ad552b" [[package]] name = "maestro" -version = "0.10.110" +version = "0.10.111" dependencies = [ "anyhow", "ctor", diff --git a/package-lock.json b/package-lock.json index 42147dc67..c89916399 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1,12 +1,12 @@ { "name": "@evalops/maestro", - "version": "0.10.110", + "version": "0.10.111", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "@evalops/maestro", - "version": "0.10.110", + "version": "0.10.111", "license": "BUSL-1.1", "bin": { "deixic-code": "bin/deixic-code", diff --git a/package.json b/package.json index 5e5c17752..0e62fc0c9 100644 --- a/package.json +++ b/package.json @@ -1,7 +1,7 @@ { "name": "@evalops/deixic-code", "description": "Deixic Code — native Rust coding agent, CLI, TUI, and runtime gateway", - "version": "0.10.110", + "version": "0.10.111", "private": false, "type": "module", "bin": { diff --git a/packages/maestro-rs/Cargo.toml b/packages/maestro-rs/Cargo.toml index bff09ca48..b8b9bb436 100644 --- a/packages/maestro-rs/Cargo.toml +++ b/packages/maestro-rs/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "maestro" -version = "0.10.110" +version = "0.10.111" edition = "2021" license = "MIT" description = "Canonical native Rust CLI for Deixic Code" diff --git a/scripts/smoke-release-native-only.mjs b/scripts/smoke-release-native-only.mjs index f5c9dddaf..4e1779ecf 100755 --- a/scripts/smoke-release-native-only.mjs +++ b/scripts/smoke-release-native-only.mjs @@ -9,6 +9,7 @@ import { } from "node:fs"; import { tmpdir } from "node:os"; import { join, resolve } from "node:path"; +import { firstLaunchTimeoutMs, runNativeSmokeCommand } from "./smoke-release-process.mjs"; // The agent only serves protocol versions it implements, so the smoke has to // announce the version this build speaks rather than a placeholder. @@ -39,22 +40,13 @@ const env = { TERM: "xterm-256color", }; -function run(args, input) { - const result = spawnSync(binary, args, { - encoding: "utf8", - env, - input, - timeout: 30_000, - }); - if (result.status !== 0) { - throw new Error(`${args.join(" ")} failed:\n${result.stderr}\n${result.stdout}`); - } - return result.stdout; -} - -if (!run(["--version"]).startsWith("deixic-code ")) +if (!runNativeSmokeCommand(binary, ["--version"], { + env, + timeoutMs: firstLaunchTimeoutMs(process.platform, process.env.RELEASE_PLATFORM), +}).startsWith("deixic-code ")) throw new Error("version smoke failed"); -if (!run(["--help"]).includes("Usage:")) throw new Error("help smoke failed"); +if (!runNativeSmokeCommand(binary, ["--help"], { env }).includes("Usage:")) + throw new Error("help smoke failed"); const headlessInput = `${JSON.stringify({ type: "hello", protocol_version: protocolVersion, client_info: { name: "native-smoke", version: "1" }, role: "controller" })}\n${JSON.stringify({ type: "shutdown" })}\n`; const headlessResult = spawnSync(binary, ["--headless"], { encoding: "utf8", diff --git a/scripts/smoke-release-process.mjs b/scripts/smoke-release-process.mjs new file mode 100644 index 000000000..ccea8d6a2 --- /dev/null +++ b/scripts/smoke-release-process.mjs @@ -0,0 +1,30 @@ +import { spawnSync } from "node:child_process"; + +const DEFAULT_TIMEOUT_MS = 30_000; + +export function firstLaunchTimeoutMs(platform, releasePlatform) { + // Cold Rosetta translation of the signed Intel release can exceed the + // ordinary smoke budget. Keep the launch bounded and test every exit. + return platform === "darwin" && releasePlatform === "darwin-x64" + ? 120_000 + : DEFAULT_TIMEOUT_MS; +} + +export function runNativeSmokeCommand(binary, args, { env, input, timeoutMs = DEFAULT_TIMEOUT_MS }) { + const result = spawnSync(binary, args, { + encoding: "utf8", + env, + input, + timeout: timeoutMs, + }); + if (result.status !== 0) { + const details = [ + `status=${result.status ?? "none"}`, + `signal=${result.signal ?? "none"}`, + `error=${result.error?.code ?? "none"}`, + `timeout=${timeoutMs}ms`, + ].join(", "); + throw new Error(`${args.join(" ")} failed (${details}):\n${result.stderr ?? ""}\n${result.stdout ?? ""}`); + } + return result.stdout; +} diff --git a/scripts/smoke-release-process.test.mjs b/scripts/smoke-release-process.test.mjs new file mode 100644 index 000000000..198c105fd --- /dev/null +++ b/scripts/smoke-release-process.test.mjs @@ -0,0 +1,47 @@ +import assert from "node:assert/strict"; +import { chmodSync, mkdtempSync, rmSync, writeFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import test from "node:test"; +import { firstLaunchTimeoutMs, runNativeSmokeCommand } from "./smoke-release-process.mjs"; + +test("only Darwin x64 gets a longer first-launch smoke budget", () => { + assert.equal(firstLaunchTimeoutMs("darwin", "darwin-x64"), 120_000); + assert.equal(firstLaunchTimeoutMs("darwin", "darwin-arm64"), 30_000); + assert.equal(firstLaunchTimeoutMs("linux", "darwin-x64"), 30_000); +}); + +test("a timed-out native binary reports the timeout, signal, and spawn error", () => { + const directory = mkdtempSync(join(tmpdir(), "maestro-smoke-process-")); + try { + const binary = join(directory, "slow-binary"); + writeFileSync(binary, "#!/bin/sh\n/bin/sleep 1\n"); + chmodSync(binary, 0o755); + assert.throws( + () => runNativeSmokeCommand(binary, ["--version"], { env: process.env, timeoutMs: 20 }), + (error) => { + assert.match(error.message, /--version failed/); + assert.match(error.message, /error=ETIMEDOUT/); + assert.match(error.message, /timeout=20ms/); + return true; + }, + ); + } finally { + rmSync(directory, { recursive: true, force: true }); + } +}); + +test("a native binary with a successful exit returns its output", () => { + const directory = mkdtempSync(join(tmpdir(), "maestro-smoke-process-")); + try { + const binary = join(directory, "version-binary"); + writeFileSync(binary, "#!/bin/sh\nprintf 'deixic-code 0.10.110\\n'\n"); + chmodSync(binary, 0o755); + assert.equal( + runNativeSmokeCommand(binary, ["--version"], { env: process.env }), + "deixic-code 0.10.110\n", + ); + } finally { + rmSync(directory, { recursive: true, force: true }); + } +}); diff --git a/vendor/dex-loop/src/context.rs b/vendor/dex-loop/src/context.rs index 04cb03682..836d66d04 100644 --- a/vendor/dex-loop/src/context.rs +++ b/vendor/dex-loop/src/context.rs @@ -146,6 +146,9 @@ pub struct Context { /// `Final`) are never attributed to the turn that raced in ahead of them. pending_turns: Vec, authorized_principal: Option, + /// 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, } #[derive(Clone, Debug, PartialEq)] @@ -184,6 +187,7 @@ impl Context { authorized_tools: Vec::new(), approval_mode: ApprovalMode::Interactive, authorized_principal: None, + uncertain_calls: Vec::new(), pending_turns: Vec::new(), } } @@ -302,6 +306,16 @@ impl Context { self.open_step.as_ref() } + /// A fresh ID does not make an unresolved operation safe to repeat. + pub(crate) fn has_uncertain_call(&self, call: &ProposedCall) -> bool { + self.uncertain_calls.iter().any(|prior| { + prior.id != call.id + && prior.tool == call.tool + && prior.args_digest == call.args_digest + && prior.principal == call.principal + }) + } + /// The step of a model attempt that started but never completed. pub(crate) fn open_attempt(&self) -> Option { self.attempt @@ -470,6 +484,15 @@ impl Context { output, receipt, } => { + self.uncertain_calls.retain(|prior| &prior.id != call); + if *outcome == Outcome::Unknown + && let Some(proposal) = self + .open_step + .as_ref() + .and_then(|step| step.calls.iter().find(|proposal| &proposal.id == call)) + { + self.uncertain_calls.push(proposal.clone()); + } if let Some(state) = self.state_mut(call) { *state = CallState::Done(ToolResult { outcome: *outcome, @@ -564,6 +587,7 @@ impl Context { self.attempt = None; self.interrupt_requested = false; self.exposed.clear(); + self.uncertain_calls.clear(); self.client_tools = client_tools; self.authorized_tools = authorized_tools; self.approval_mode = approval_mode; diff --git a/vendor/dex-loop/src/engine.rs b/vendor/dex-loop/src/engine.rs index ee514ba7a..375f4817e 100644 --- a/vendor/dex-loop/src/engine.rs +++ b/vendor/dex-loop/src/engine.rs @@ -36,6 +36,7 @@ const UNKNOWN_INTERRUPTED: &str = /// a `Running` outcome that nothing will update: both invite the model to /// retry, which mints a new `CallId` and can run the mutation twice. const UNKNOWN_NO_RETRY: &str = "outcome unknown: this call already started; do not retry it without first checking whether it took effect"; +const UNCERTAIN_REPEAT: &str = "not run: the same operation has an unknown outcome in this turn; check whether it took effect before trying again"; const APPROVER_DECLINED: &str = "denied: the approver declined this call"; const APPROVAL_MISMATCH: &str = "denied: the approval does not match this call's arguments"; const MISSING_QUESTION: &str = "invalid call: args.question must be a non-empty string"; @@ -382,6 +383,30 @@ where self.emit(ctx, events).await?; return Ok(Some(Exit::Failed)); } + if failure.is_some() && !text.is_empty() && !cancel.is_cancelled() { + // The customer already read this text. Keep it as the + // answer, visibly marked as cut off, instead of withdrawing + // it: a long answer that loses its stream near the end (a + // provider or gateway limit) must not vanish. + let mut text = text; + let tail = filter.finish(); + if !tail.is_empty() { + text.push_str(&tail); + self.log.append_text(tail).await?; + } + self.log.append_text(CUT_OFF_NOTICE.to_owned()).await?; + text.push_str(CUT_OFF_NOTICE); + let mut events = pending_usage; + events.push(Event::ModelStepCompleted { + step, + text: text.clone(), + calls: Vec::new(), + reasoning: None, + }); + events.push(Event::Final { text }); + self.emit(ctx, events).await?; + return Ok(Some(Exit::Done)); + } if let Some(message) = failure { let mut events = pending_usage; events.push(Event::ModelAttemptAbandoned { step }); @@ -613,6 +638,9 @@ where continue; } } + if self.refuse_uncertain_repeat(ctx, call, &spec).await? { + continue; + } if let (None, Verdict::NeedsApproval { approval, summary }) = (decision, verdict) { if self .flush(ctx, &calls, &mut wave, cancel, run_started) @@ -670,6 +698,22 @@ where Ok(None) } + /// Refusal does not claim or dispatch a second effect. Reads remain safe + /// to retry; an existing call ID still resolves through its effect ledger. + async fn refuse_uncertain_repeat( + &self, + ctx: &mut Context, + call: &ProposedCall, + spec: &ToolSpec, + ) -> Result { + if spec.read_only || !ctx.has_uncertain_call(call) { + return Ok(false); + } + self.finish(ctx, call, ToolResult::error(UNCERTAIN_REPEAT)) + .await?; + Ok(true) + } + /// One call to a `Client`-executor tool: approval (if the host's /// allowlist marked it a mutation), then `ClientToolRequested`, mirroring /// how the main `dispatch` loop handles `NeedsApproval` and `User`. @@ -701,7 +745,11 @@ where .await?; return Ok(ClientToolOutcome::Continue); } - } else if spec.governance == GovernanceClass::Approval { + } + if self.refuse_uncertain_repeat(ctx, call, &spec).await? { + return Ok(ClientToolOutcome::Continue); + } + if decision.is_none() && spec.governance == GovernanceClass::Approval { if self.flush(ctx, calls, wave, cancel, run_started).await? { return Ok(ClientToolOutcome::Break); } @@ -1150,6 +1198,10 @@ fn search_spec() -> ToolSpec { } } +/// Appended to an answer whose model stream failed after text was shown. +pub const CUT_OFF_NOTICE: &str = + "\n\n_This answer was cut off before it finished. Ask me to continue from here._"; + fn call_id(turn: &TurnId, step: u32, index: usize) -> CallId { CallId(format!("{turn}-{step}-{index}")) } diff --git a/vendor/dex-loop/src/lib.rs b/vendor/dex-loop/src/lib.rs index 535a91c8a..18ae89167 100644 --- a/vendor/dex-loop/src/lib.rs +++ b/vendor/dex-loop/src/lib.rs @@ -31,7 +31,7 @@ mod sanitize; pub use budget::{Budget, BudgetAxis}; pub use compaction::{Compaction, Compactor, NoCompaction, Summarize, Threshold}; pub use context::{Context, Entry, Message}; -pub use engine::{DEFAULT_TOOL_CALL_DEADLINE, Engine, Exit, TOOLS_SEARCH}; +pub use engine::{CUT_OFF_NOTICE, DEFAULT_TOOL_CALL_DEADLINE, Engine, Exit, TOOLS_SEARCH}; pub use event::{ ApprovalId, ApprovalMode, ArtifactRef, CallId, ClientToolSpec, Cursor, ErrorCode, Event, HEADLESS_AUTO_APPROVER, MessageId, Outcome, Output, OutputRef, PrincipalId, ProposedCall, diff --git a/vendor/dex-loop/tests/scenarios.rs b/vendor/dex-loop/tests/scenarios.rs index 7ca1e6c01..3352959d5 100644 --- a/vendor/dex-loop/tests/scenarios.rs +++ b/vendor/dex-loop/tests/scenarios.rs @@ -7,8 +7,8 @@ mod support; use std::time::{Duration, Instant}; use dex_loop::{ - ApprovalId, ApprovalMode, Budget, CancellationToken, Engine, Event, Exit, Lexicon, ModelError, - OutputRef, ProposedCall, Threshold, ToolName, ToolResult, TurnId, Verdict, + ApprovalId, ApprovalMode, Budget, CUT_OFF_NOTICE, CancellationToken, Engine, Event, Exit, + Lexicon, ModelError, OutputRef, ProposedCall, Threshold, ToolName, ToolResult, TurnId, Verdict, }; use serde_json::json; use support::*; @@ -981,6 +981,41 @@ async fn usage_reported_before_a_failed_attempt_still_lands_with_the_abandon() { assert_eq!(model.calls(), 1); } +// 8d3. A stream that fails after the customer saw text keeps that text as +// the answer, marked as cut off: a long answer that loses its stream near +// the end must not vanish. +#[tokio::test] +async fn a_stream_that_fails_after_text_keeps_the_answer_marked_cut_off() { + let (log, model, exit) = run_to_budget( + budget(), + FakeModel::new(vec![vec![ + text("A long answer, "), + text("nearly done"), + usage(10, 5, 100), + Err(ModelError { + message: "provider_stream_timeout: provider stream timed out".into(), + }), + ]]), + ) + .await; + assert_eq!(exit, Exit::Done); + let kept = format!("A long answer, nearly done{CUT_OFF_NOTICE}"); + assert_eq!( + log.shapes_after(2), + strings(&[ + &format!("delta:{kept}"), + "usage:15", + &format!("completed:{kept}:[]"), + &format!("final:{kept}"), + ]) + ); + assert_eq!( + model.calls(), + 1, + "text was shown, so the step is not retried" + ); +} + // 8e. A `Started` call whose tool vanished from the offered catalog (deploy, // grant revoke) resolves through the ledger, never as "unknown tool": the // model must not be told to retry a call that may have already run under @@ -1819,3 +1854,213 @@ async fn reasoning_is_committed_with_its_step_and_returned_after_park_and_crash( assert_eq!(returned, vec![Some(reasoning)]); assert_eq!(log.rehydrate(), ctx); } +#[tokio::test] +async fn fresh_call_ids_cannot_repeat_an_unknown_mutation() { + let log = FakeLog::default(); + let write = ProposedCall::new( + call_id("t1", 1, 0), + ToolName::new("update"), + json!({"key":"w"}), + alice(), + ); + crashed_after_start(&log, &write); + let effects = FakeEffects::default().seed( + write.id.clone(), + Some(ToolResult::unknown("executor lost the result")), + ); + let model = FakeModel::new(vec![ + vec![call("update", json!({"key":"w"}))], + vec![call("update", json!({"key":"w"}))], + vec![text("check the original")], + ]); + let tools = FakeTools::new(vec![write_tool("update")]); + 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 fresh ID repeated an unknown mutation" + ); + assert_eq!( + tools.policy_checks().len(), + 2, + "current policy still applies" + ); + for step in [2, 3] { + assert_eq!( + effects.recorded(&call_id("t1", step, 0)), + None, + "refusal must not claim a second effect" + ); + } + assert_eq!( + effects.recorded(&write.id).unwrap().unwrap().outcome, + dex_loop::Outcome::Unknown + ); + assert_eq!(ctx, log.rehydrate()); +} + +#[tokio::test] +async fn uncertain_repeat_guard_preserves_reads_distinct_operations_and_known_retries() { + for (prior_outcome, next_tool, next_args) in [ + ( + ToolResult::unknown("lost"), + "update", + json!({"key":"different"}), + ), + (ToolResult::unknown("lost"), "other", json!({"key":"w"})), + (ToolResult::unknown("lost"), "search", json!({"key":"w"})), + (ToolResult::error("not run"), "update", json!({"key":"w"})), + ] { + let log = FakeLog::default(); + let write = ProposedCall::new( + call_id("t1", 1, 0), + ToolName::new("update"), + json!({"key":"w"}), + alice(), + ); + crashed_after_start(&log, &write); + let effects = FakeEffects::default().seed(write.id.clone(), Some(prior_outcome)); + let model = FakeModel::new(vec![vec![call(next_tool, next_args)], vec![text("done")]]); + let tools = FakeTools::new(vec![ + write_tool("update"), + write_tool("other"), + read_tool("search"), + ]); + 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_eq!(tools.runs().len(), 1); + assert_eq!(ctx, log.rehydrate()); + } +} + +#[tokio::test] +async fn uncertain_mutation_guard_survives_compaction_and_a_new_turn_is_explicit() { + let log = FakeLog::default(); + let write = ProposedCall::new( + call_id("t1", 1, 0), + ToolName::new("update"), + json!({"key":"w"}), + alice(), + ); + crashed_after_start(&log, &write); + let cursor = log.host_append(Event::ToolFinished { + call: write.id.clone(), + outcome: dex_loop::Outcome::Unknown, + output: dex_loop::Output::Text("lost result".into()), + receipt: None, + }); + log.host_append(Event::Compaction { + covers_to_cursor: cursor, + summary: "Earlier work requires reconciliation".into(), + }); + let effects = + FakeEffects::default().seed(write.id.clone(), Some(ToolResult::unknown("lost result"))); + let model = FakeModel::new(vec![ + vec![call("update", json!({"key":"w"}))], + vec![text("await a decision")], + vec![call("update", json!({"key":"w"}))], + vec![text("done")], + ]); + let tools = FakeTools::new(vec![write_tool("update")]); + 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()); + assert_eq!(ctx, log.rehydrate()); + let mut ctx = log.start_turn("t2", "I checked; try it again"); + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) + ); + assert_eq!(tools.runs().len(), 1); + assert_eq!(tools.policy_checks().len(), 2); + assert_eq!(ctx, log.rehydrate()); +} + +#[tokio::test] +async fn unknown_client_mutations_are_not_requested_again_and_unknown_reads_can_retry() { + for read_only in [false, true] { + let log = FakeLog::default(); + let prior = ProposedCall::new( + call_id("t1", 1, 0), + ToolName::new("client.operation"), + json!({"key":"w"}), + alice(), + ); + crashed_after_start(&log, &prior); + log.host_append(Event::ToolFinished { + call: prior.id.clone(), + outcome: dex_loop::Outcome::Unknown, + output: dex_loop::Output::Text("lost".into()), + receipt: None, + }); + let model = FakeModel::new(vec![ + vec![call("client.operation", json!({"key":"w"}))], + vec![text("check before retry")], + ]); + let tools = FakeTools::new(vec![if read_only { + read_tool("client.operation") + } else { + client_executed_tool("client.operation", false) + }]); + let engine = engine(&log, &model, &tools, budget()); + let mut ctx = log.rehydrate(); + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) + ); + assert_eq!(tools.runs().len(), usize::from(read_only)); + assert!(!log.events().iter().any(|event| matches!( + event, + Event::ClientToolRequested { .. } | Event::ApprovalRequested { .. } + ))); + assert_eq!(ctx, log.rehydrate()); + } +} + +#[tokio::test] +async fn uncertain_mutations_keep_their_principal_identity() { + let log = FakeLog::default(); + let prior = ProposedCall::new( + call_id("t1", 1, 0), + ToolName::new("update"), + json!({"key":"w"}), + alice(), + ); + crashed_after_start(&log, &prior); + log.host_append(Event::ToolFinished { + call: prior.id.clone(), + outcome: dex_loop::Outcome::Unknown, + output: dex_loop::Output::Text("lost".into()), + receipt: None, + }); + log.host_append(Event::Steer { + principal: bob(), + text: "I checked my operation".into(), + }); + let model = FakeModel::new(vec![ + vec![call("update", json!({"key":"w"}))], + vec![text("done")], + ]); + let tools = FakeTools::new(vec![write_tool("update")]); + let engine = engine(&log, &model, &tools, budget()); + let mut ctx = log.rehydrate(); + assert_eq!( + engine.run(&mut ctx, &CancellationToken::new()).await, + Ok(Exit::Done) + ); + assert_eq!(tools.runs().len(), 1); + assert_eq!(tools.policy_checks()[0].1, "bob"); + assert_eq!(ctx, log.rehydrate()); +}