diff --git a/.repository-projection.json b/.repository-projection.json index 5a17de914..825767348 100644 --- a/.repository-projection.json +++ b/.repository-projection.json @@ -3,11 +3,11 @@ "projection": "deixic-code", "projectionSchemaVersion": 1, "sourceRepository": "dx-corp/mono", - "sourceSha": "eb0ec94780d951b21cd54c7de3380ac3624791ac", + "sourceSha": "518c4d79c54135d51a4d885f38cf4f4e701e44a0", "destinationRepository": "dx-corp/code", - "priorProjectedBase": "504ba52e32b5b99544a3c8a5c1a546f299ef9cf0", + "priorProjectedBase": "bac6ed17bdc31b56dfdc9d7740103300965fec5c", "definitionDigest": "cb9d429542ebb0a2de9b42a7aad60d9d8696a648ceba47c30f05c0b285ca0db7", "toolDigest": "f8cb071b0f27267120ccf45a00d0982f45113bd23535bef6a1555b4933f99f13", - "contentDigest": "cd768ddf8f8c6ff66d5b037a18b20f30adebed7ee3ef2169c0b68843d84b839f", + "contentDigest": "00d6e14ad707b1d000be9e241309bd5d6a63d5aaa0c95f61c30df1be063723ef", "publicationEligible": true } diff --git a/Cargo.lock b/Cargo.lock index 84e94b3c2..d94896bf8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3967,6 +3967,9 @@ dependencies = [ "fd-lock", "futures-util", "maestro-ai", + "maestro-local-host", + "maestro-runtime", + "maestro-runtime-contracts", "serde", "serde_json", "tempfile", diff --git a/packages/dex-host-rs/Cargo.toml b/packages/dex-host-rs/Cargo.toml index 2c850aea5..8711a92d8 100644 --- a/packages/dex-host-rs/Cargo.toml +++ b/packages/dex-host-rs/Cargo.toml @@ -18,6 +18,11 @@ dex-loop = { path = "../../vendor/dex-loop" } # `dex_loop::Model` -- default features off, matching every other consumer's # workspace alias, so this stays Bedrock-free. maestro-ai.workspace = true +# `HostTools` (`src/host_tools.rs`) offers a Maestro execution host's tools to +# the kernel and gates them with the native actor's own classification +# (`maestro_runtime::agent::loop_policy`). +maestro-runtime.workspace = true +maestro-runtime-contracts.workspace = true anyhow.workspace = true fd-lock.workspace = true @@ -29,5 +34,6 @@ tokio = { workspace = true, features = ["fs", "rt", "sync"] } tokio-stream.workspace = true [dev-dependencies] +maestro-local-host.workspace = true tempfile.workspace = true tokio = { workspace = true, features = ["rt-multi-thread", "macros", "time"] } diff --git a/packages/dex-host-rs/src/host_tools.rs b/packages/dex-host-rs/src/host_tools.rs new file mode 100644 index 000000000..b7f537e66 --- /dev/null +++ b/packages/dex-host-rs/src/host_tools.rs @@ -0,0 +1,402 @@ +//! The `Tools` port over a Maestro `NativeExecutionHost`. +//! +//! Every tool the host registers (bash, file edits, search, MCP, inline +//! tools) is offered to the dex-loop kernel as it is. The host still owns +//! execution, hooks, sandbox and the action firewall; this module only +//! translates. Approval and read-only classification call the native actor's +//! own functions (`maestro_runtime::agent::loop_policy`), so a call is gated +//! the same way on both loops. +//! +//! A call that would have prompted under the native actor does not run. Its +//! result is a `needs_confirmation` preview; the model asks the person with +//! `user.ask` bound to that preview, and once they choose Confirm it calls +//! again with `confirmation` set to the preview's call id. That is the +//! kernel's one human gate (see `dex-tools`' `confirm` module, which this +//! mirrors for Maestro's tools). + +use std::sync::{Arc, Mutex, PoisonError}; +use std::time::Instant; + +use dex_loop::{ + ApprovalId, CallId, CancellationToken, Context, ExecutorKind, GovernanceClass, PrincipalId, + ProposedCall, ThreadId, ToolName, ToolResult, ToolSpec, Tools, Verdict, args_digest, +}; +use maestro_runtime::agent::loop_policy; +use maestro_runtime::agent::native_host::{ + ApprovalMode, NativeExecutionHostHandle, NativeFirewallVerdict, NativeHookResult, + NativeToolExecutionOptions, +}; +use maestro_runtime::agent::workflow_state::{WorkflowStateTracker, apply_workflow_state_hooks}; +use maestro_runtime_contracts::contracts::ToolOutcome; +use serde_json::{Map, Value, json}; + +/// The kernel's question tool. The engine parks on it (`ExecutorKind::User`) +/// and binds a `confirmation` reference to the preview it names. Here it only +/// confirms a previewed action: Maestro's own `ask_user` stays the way to ask +/// an open question, which the person answers in their next message. +pub const USER_ASK: &str = "user.ask"; + +/// The argument that carries a confirmed preview's call id. +pub const CONFIRMATION_FIELD: &str = "confirmation"; + +/// Longest string argument value shown in a preview. +const PREVIEW_VALUE_CHARS: usize = 400; + +/// Refusal for a gated call in a turn nobody is watching. +pub const HEADLESS_GATED: &str = + "this action needs a person's confirmation and nobody is attached to this turn"; + +fn user_ask() -> ToolSpec { + ToolSpec { + name: ToolName::new(USER_ASK), + label: "Confirm an action".into(), + description: "Ask the person to confirm one previewed action and wait for their choice. Only for a needs_confirmation preview; ask open questions with ask_user.".into(), + schema: json!({ + "type": "object", + "additionalProperties": false, + "required": ["question", CONFIRMATION_FIELD], + "properties": { + "question": {"type": "string", "minLength": 1, "maxLength": 2000}, + CONFIRMATION_FIELD: { + "type": "object", + "additionalProperties": false, + "required": ["proposal_call_id"], + "properties": { + "proposal_call_id": {"type": "string", "minLength": 1, "maxLength": 256}, + }, + }, + }, + }), + read_only: true, + core: true, + governance: GovernanceClass::Plain, + executor: ExecutorKind::User, + } +} + +/// Adds the optional confirmation reference to a mutation's input schema, so +/// a strict schema still admits it. +fn declare_confirmation(schema: &mut Value) { + let Some(root) = schema.as_object_mut() else { + return; + }; + let properties = root + .entry("properties") + .or_insert_with(|| Value::Object(Map::new())); + if let Some(properties) = properties.as_object_mut() { + properties.insert( + CONFIRMATION_FIELD.to_owned(), + json!({ + "type": "string", + "description": "Only after the person chose Confirm in a user.ask question bound to this action's preview: that preview's proposal_call_id. Single use." + }), + ); + } +} + +/// The arguments without `confirmation`: what the host executes. +fn strip(args: &Value) -> Value { + match args { + Value::Object(map) => Value::Object( + map.iter() + .filter(|(key, _)| key.as_str() != CONFIRMATION_FIELD) + .map(|(key, value)| (key.clone(), value.clone())) + .collect(), + ), + other => other.clone(), + } +} + +fn truncate(value: &Value) -> Value { + match value { + Value::String(text) if text.chars().count() > PREVIEW_VALUE_CHARS => { + let cut: String = text.chars().take(PREVIEW_VALUE_CHARS).collect(); + Value::String(format!( + "{cut}... [{} characters in all]", + text.chars().count() + )) + } + Value::Array(items) => Value::Array(items.iter().map(truncate).collect()), + Value::Object(map) => Value::Object( + map.iter() + .map(|(key, value)| (key.clone(), truncate(value))) + .collect(), + ), + other => other.clone(), + } +} + +/// The `needs_confirmation` document a gated call returns. `Context` records +/// it as a preview only when `args_digest` matches the call's arguments. +fn preview(spec: &ToolSpec, call: &ProposedCall, reason: Option<&str>) -> String { + let arguments = strip(&call.args); + json!({ + "status": "needs_confirmation", + "tool": spec.name.as_str(), + "action": spec.label, + "reason": reason, + "arguments": truncate(&arguments), + "args_digest": args_digest(&arguments), + "proposal_call_id": call.id.as_str(), + "instructions": format!( + "Nothing was done. Call {USER_ASK} with a plain question showing the action ({}) and its key details, and confirmation={{\"proposal_call_id\":\"{}\"}}. The person must choose Confirm in that question. Ordinary chat replies do not authorize execution. After Confirm, call this tool with exactly the same arguments and confirmation=\"{}\". Changed arguments need a new preview and question; every decision is single use.", + spec.label, call.id, call.id + ), + }) + .to_string() +} + +fn approval_id(call: &CallId) -> ApprovalId { + ApprovalId::new(format!("approval-{call}")) +} + +/// A `user.ask` call must carry an exact-action reference naming an +/// unexecuted preview, by the same person, not already asked about. Unlike `dex-tools`, +/// a read can be gated too (the firewall holds some reads), so a read's +/// preview is confirmable. +fn question_binding(ctx: &Context, ask: &ProposedCall, tools: &HostTools) -> Result<(), String> { + let reference = ask + .args + .get(CONFIRMATION_FIELD) + .ok_or("user.ask only confirms a previewed action; ask open questions with ask_user")?; + let id = reference + .get("proposal_call_id") + .and_then(Value::as_str) + .ok_or("confirmation requires proposal_call_id")?; + let prior = ctx + .action_preview(&CallId::new(id)) + .ok_or("confirmation proposal is unavailable")?; + let spec = tools + .spec(&prior.tool) + .ok_or("confirmation tool is unavailable")?; + if !ctx.is_unexecuted_preview(&prior.id) { + return Err("confirmation requires a preview before execution".into()); + } + if prior.principal != ask.principal { + return Err("confirmation principal does not match".into()); + } + if ctx.confirmation_question_exists(&prior.id) { + return Err("confirmation proposal already has a question; create a fresh preview".into()); + } + if ctx.action_question_text(&prior.id, &spec.label).is_none() { + return Err( + "confirmation details cannot fit the chat question; split or reduce the action".into(), + ); + } + Ok(()) +} + +fn render(outcome: &ToolOutcome) -> ToolResult { + match outcome { + ToolOutcome::Succeeded { output } => ToolResult::text(output.as_str()), + ToolOutcome::Failed { + error, + partial_output, + } => match partial_output { + Some(partial) if !partial.as_str().is_empty() => { + ToolResult::error(format!("{}\n\n{}", error.message(), partial.as_str())) + } + _ => ToolResult::error(error.message()), + }, + ToolOutcome::Denied { reason } => { + ToolResult::error(format!("denied: {}", reason.message())) + } + ToolOutcome::Cancelled { .. } => ToolResult::error("cancelled"), + ToolOutcome::Indeterminate { reason } => ToolResult::unknown(reason.clone()), + } +} + +/// The `Tools` port over one Maestro execution host. +#[derive(Clone)] +pub struct HostTools { + host: NativeExecutionHostHandle, + mode: ApprovalMode, + catalog: Arc<[ToolSpec]>, + workflow: Arc>, +} + +impl HostTools { + /// Offers every tool `host` registers, gated as the native actor gates + /// them under `mode`. + pub fn new(host: NativeExecutionHostHandle, mode: ApprovalMode) -> Self { + let mut catalog = Vec::new(); + for definition in host.tool_definitions() { + let name = definition.tool.name.clone(); + if name == USER_ASK { + continue; + } + let annotations = host.tool_annotations(&name); + let read_only = loop_policy::parallel_read_only( + &name, + definition.requires_approval, + annotations.as_ref(), + host.is_explicit_inline_read_only_tool(&name), + ); + let mut schema = definition.tool.input_schema.clone(); + declare_confirmation(&mut schema); + catalog.push(ToolSpec { + name: ToolName::new(&name), + label: name.clone(), + description: definition.tool.description.clone(), + schema, + read_only, + core: true, + governance: if definition.requires_approval { + GovernanceClass::Approval + } else { + GovernanceClass::Plain + }, + executor: ExecutorKind::InProcess, + }); + } + catalog.push(user_ask()); + Self { + host, + mode, + catalog: Arc::from(catalog), + workflow: Arc::new(Mutex::new(WorkflowStateTracker::default())), + } + } + + fn firewall(&self, name: &str, args: &Value) -> NativeFirewallVerdict { + let snapshot = self + .workflow + .lock() + .unwrap_or_else(PoisonError::into_inner) + .snapshot(); + loop_policy::firewall_verdict( + &self.host, + name, + args, + &snapshot, + self.host.tool_annotations(name).as_ref(), + false, + ) + } +} + +impl Tools for HostTools { + fn catalog(&self) -> &[ToolSpec] { + &self.catalog + } + + async fn search(&self, _principal: &PrincipalId, query: &str) -> Vec { + let query = query.to_ascii_lowercase(); + if query.trim().is_empty() { + return Vec::new(); + } + self.catalog + .iter() + .filter(|spec| { + spec.name.as_str().to_ascii_lowercase().contains(&query) + || spec.description.to_ascii_lowercase().contains(&query) + }) + .map(|spec| spec.name.clone()) + .collect() + } + + async fn policy(&self, ctx: &Context, call: &ProposedCall) -> Verdict { + let Some(spec) = self.spec(&call.tool) else { + return Verdict::Deny(format!("unknown tool: {}", call.tool)); + }; + if call.tool.as_str() == USER_ASK { + return match question_binding(ctx, call, self) { + Ok(()) => Verdict::Allow, + Err(reason) => Verdict::Deny(reason), + }; + } + let name = call.tool.as_str(); + let args = strip(&call.args); + let missing = self.host.missing_required(name, &args); + if !missing.is_empty() { + return Verdict::Deny(format!( + "missing required arguments: {}", + missing.join(", ") + )); + } + let firewall = self.firewall(name, &args); + if let NativeFirewallVerdict::Block { reason } = &firewall { + return Verdict::Deny(reason.clone()); + } + let required = + loop_policy::approval_required(self.mode, false, &firewall, &self.host, name, &args); + if !required { + return Verdict::Allow; + } + if ctx.approval_mode() == dex_loop::ApprovalMode::Headless { + return Verdict::Deny(HEADLESS_GATED.into()); + } + if ctx.confirmed_action(call) { + return Verdict::Confirmed { + approval: approval_id(&call.id), + summary: format!( + "decision=policy_approved_after_action_confirmation; {}", + spec.label + ), + }; + } + let reason = match &firewall { + NativeFirewallVerdict::RequireApproval { reason } => Some(reason.as_str()), + _ => None, + }; + Verdict::NeedsConfirmation { + preview: preview(spec, call, reason), + } + } + + async fn run( + &self, + _thread: &ThreadId, + call: &ProposedCall, + cancel: &CancellationToken, + ) -> ToolResult { + let name = call.tool.as_str(); + let call_id = call.id.as_str(); + let mut args = strip(&call.args); + match self.host.hook_pre_tool_use(name, call_id, &args).await { + NativeHookResult::Block { reason } => { + return ToolResult::error(format!("blocked by hook: {reason}")); + } + NativeHookResult::ModifyInput { new_input } => args = new_input, + NativeHookResult::Continue | NativeHookResult::InjectContext { .. } => {} + } + let started = Instant::now(); + let execution = self + .host + .execute_tool( + name, + &args, + None, + call_id, + NativeToolExecutionOptions { + cancel: cancel.clone(), + approved_inline_env: None, + }, + ) + .await; + let result = render(&execution.outcome); + let is_error = result.outcome != dex_loop::Outcome::Succeeded; + // A PII tracking error leaves the tracker as it was; the firewall + // keeps judging later calls against the last good snapshot. + let _ = apply_workflow_state_hooks( + name, + call_id, + &args, + &mut self.workflow.lock().unwrap_or_else(PoisonError::into_inner), + is_error, + ); + let output = match &result.output { + dex_loop::Output::Text(text) => text.as_str(), + _ => "", + }; + let elapsed = u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX); + if let NativeHookResult::Block { reason } = self + .host + .hook_post_tool_use(name, call_id, &args, output, is_error, elapsed) + .await + { + return ToolResult::error(format!("{output}\n\nblocked by hook: {reason}")); + } + result + } +} diff --git a/packages/dex-host-rs/src/lib.rs b/packages/dex-host-rs/src/lib.rs index c0c38f5a7..511b46286 100644 --- a/packages/dex-host-rs/src/lib.rs +++ b/packages/dex-host-rs/src/lib.rs @@ -1,7 +1,10 @@ //! A local `dex_loop` host: `Log`, `Tools` and `Effects` ports backed by the //! workspace filesystem, with no database and no HTTP server, plus a `Model` //! adapter over `maestro-ai` (`model.rs`) and a consumer that drives one -//! turn to completion unattended (`turn.rs`, `run_local_turn`). +//! turn to completion unattended (`turn.rs`, `run_local_turn`), and +//! `HostTools` (`host_tools.rs`): the `Tools` port over a Maestro +//! `NativeExecutionHost`, which offers Maestro's real tool registry to the +//! kernel and gates it the way the native actor does. //! //! This is the first slice of "Maestro becomes another host with a local //! Log" (see `docs/design/maestro-on-dex-loop.md` at the repository root). @@ -18,6 +21,7 @@ //! boundary, the same shape `maestro-swarm` already uses. mod effects; +mod host_tools; mod lease; mod log; mod model; @@ -25,6 +29,7 @@ mod tools; mod turn; pub use effects::LocalEffects; +pub use host_tools::{CONFIRMATION_FIELD, HEADLESS_GATED, HostTools, USER_ASK}; pub use log::{LocalLog, LogError}; pub use model::AiRsModel; pub use tools::{LocalTools, READ_FILE, WRITE_FILE}; diff --git a/packages/dex-host-rs/tests/host_tools.rs b/packages/dex-host-rs/tests/host_tools.rs new file mode 100644 index 000000000..1e3c6f75f --- /dev/null +++ b/packages/dex-host-rs/tests/host_tools.rs @@ -0,0 +1,436 @@ +//! Runs Maestro's real local execution host (`ToolExecutor`, hooks, sandbox, +//! action firewall) on the dex-loop kernel through `HostTools`, with a +//! scripted model. +//! +//! A command the native actor would have prompted for does not run: it +//! returns a `needs_confirmation` preview, the model asks with `user.ask` +//! bound to it, and only the person's Confirm lets the identical call run. + +use std::collections::VecDeque; +use std::path::Path; +use std::sync::{Arc, Mutex, PoisonError}; + +use dex_loop::{ + ApprovalMode as TurnMode, Budget, CallId, CancellationToken, ConfirmationDecision, Context, + Engine, Event, Exit, Lexicon, Log, Model, ModelChunk, ModelError, Outcome, PrincipalId, + ProposedCall, ThreadId, ToolName, ToolSpec, Tools, TurnId, Verdict, args_digest, +}; +use futures_util::{Stream, stream}; +use maestro_dex_host::{ + CONFIRMATION_FIELD, HEADLESS_GATED, HostTools, LocalEffects, LocalLog, USER_ASK, +}; +use maestro_local_host::agent::{NativeAgentConfig, dex_loop_execution_host}; +use maestro_runtime::agent::CredentialVault; +use maestro_runtime::agent::native_host::ApprovalMode; +use serde_json::{Value, json}; +use tempfile::TempDir; + +type Script = Vec>; +type Step = Box Script + Send + Sync>; + +/// One scripted response per model call; a step reads the context so it can +/// name a call id the engine assigned. +#[derive(Clone, Default)] +struct ScriptedModel { + steps: Arc>>, +} + +impl ScriptedModel { + fn new(steps: Vec) -> Self { + Self { + steps: Arc::new(Mutex::new(steps.into())), + } + } +} + +impl Model for ScriptedModel { + fn stream<'a>( + &'a self, + ctx: &'a Context, + _tools: &'a [&'a ToolSpec], + ) -> impl Stream> + Send + 'a { + let step = self + .steps + .lock() + .unwrap_or_else(PoisonError::into_inner) + .pop_front(); + let script = match step { + Some(step) => step(ctx), + None => vec![Err(ModelError { + class: dex_loop::ErrorClass::Unknown, + message: "no script left".into(), + })], + }; + stream::iter(script) + } +} + +fn call(name: &str, args: Value) -> Step { + let name = name.to_owned(); + Box::new(move |_| { + vec![Ok(ModelChunk::ToolCall { + name: ToolName::new(&name), + args: args.clone(), + })] + }) +} + +fn answer(text: &str) -> Step { + let text = text.to_owned(); + Box::new(move |_| vec![Ok(ModelChunk::Text(text.clone()))]) +} + +/// The newest unexecuted preview's call id. +fn preview_id(ctx: &Context) -> String { + ctx.tool_evidence() + .iter() + .rev() + .find(|record| ctx.is_unexecuted_preview(&record.call.id)) + .map(|record| record.call.id.as_str().to_owned()) + .expect("a needs_confirmation preview") +} + +fn ask_about_preview() -> Step { + Box::new(|ctx| { + vec![Ok(ModelChunk::ToolCall { + name: ToolName::new(USER_ASK), + args: json!({ + "question": "Run this command?", + CONFIRMATION_FIELD: {"proposal_call_id": preview_id(ctx)}, + }), + })] + }) +} + +fn confirmed(name: &str, args: Value) -> Step { + let name = name.to_owned(); + Box::new(move |ctx| { + let mut args = args.clone(); + args[CONFIRMATION_FIELD] = Value::String(confirmed_preview(ctx)); + vec![Ok(ModelChunk::ToolCall { + name: ToolName::new(&name), + args, + })] + }) +} + +/// The preview the last `user.ask` was bound to. +fn confirmed_preview(ctx: &Context) -> String { + ctx.history() + .iter() + .rev() + .find_map(|entry| match &entry.message { + dex_loop::Message::Assistant { calls, .. } => calls + .iter() + .find(|call| call.tool.as_str() == USER_ASK) + .and_then(|call| call.args[CONFIRMATION_FIELD]["proposal_call_id"].as_str()) + .map(str::to_owned), + _ => None, + }) + .expect("a bound user.ask") +} + +fn thread() -> ThreadId { + ThreadId { + org: "local".into(), + workspace: "local".into(), + thread: "thread-1".into(), + } +} + +fn alice() -> PrincipalId { + PrincipalId::new("alice") +} + +fn host_tools(workspace: &Path, mode: ApprovalMode) -> HostTools { + let config = NativeAgentConfig { + cwd: workspace.to_string_lossy().into_owned(), + ..NativeAgentConfig::default() + }; + let host = dex_loop_execution_host(&config, CredentialVault::new()).expect("compose host"); + HostTools::new(host, mode) +} + +struct Turn { + _state: TempDir, + log: LocalLog, + engine: Engine, + cancel: CancellationToken, +} + +impl Turn { + async fn start(tools: HostTools, mode: TurnMode, steps: Vec) -> Self { + let state = TempDir::new().expect("state tempdir"); + let log = LocalLog::acquire(state.path().join("log"), &thread()) + .await + .expect("acquire log"); + log.append(&[Event::UserMessage { + turn: TurnId::new("t1"), + message_id: None, + model_binding: None, + principal: alice(), + text: "make the marker".into(), + attachments: Vec::new(), + client_tools: Vec::new(), + authorized_tools: Vec::new(), + approval_mode: mode, + }]) + .await + .expect("append user message"); + let effects = LocalEffects::open(state.path().join("effects.json")) + .await + .expect("open effects ledger"); + let engine = Engine::new( + log.clone(), + ScriptedModel::new(steps), + tools, + effects, + Lexicon::default(), + Budget::default(), + ); + Self { + _state: state, + log, + engine, + cancel: CancellationToken::new(), + } + } + + async fn run(&self) -> Exit { + let entries = self.log.read_all().await.expect("read for rehydrate"); + let mut ctx = dex_loop::rehydrate(thread(), &entries); + self.engine.run(&mut ctx, &self.cancel).await.expect("run") + } + + async fn events(&self) -> Vec { + self.log + .read_all() + .await + .expect("read log") + .into_iter() + .map(|(_, event)| event) + .collect() + } + + /// Answers the parked question the way the desktop's card does. + async fn decide(&self, asked: &CallId, decision: ConfirmationDecision) { + let binding = self + .events() + .await + .into_iter() + .find_map(|event| match event { + Event::Question { + call, + confirmation: Some(binding), + .. + } if &call == asked => Some(binding), + _ => None, + }) + .expect("the question carries the action binding"); + self.log + .append(&[Event::Answer { + call: asked.clone(), + principal: alice(), + text: format!("{decision:?}"), + confirmation_decision: decision, + args_digest: binding.args_digest, + }]) + .await + .expect("append answer"); + } +} + +/// `touch` writes, so the native actor's bash classifier asks first. +fn touch(workspace: &Path) -> Value { + json!({"command": format!("touch {}", workspace.join("marker").display())}) +} + +#[tokio::test] +async fn the_catalog_offers_the_host_tools_and_the_kernel_question() { + let workspace = TempDir::new().expect("workspace"); + let tools = host_tools(workspace.path(), ApprovalMode::Selective); + let spec = |name: &str| tools.spec(&ToolName::new(name)).cloned(); + assert!(spec("read").expect("read").read_only); + assert!(!spec("bash").expect("bash").read_only); + assert!(spec("bash").expect("bash").schema["properties"][CONFIRMATION_FIELD].is_object()); + assert_eq!( + spec(USER_ASK).expect("user.ask").executor, + dex_loop::ExecutorKind::User + ); + assert!( + spec("ask_user").is_some(), + "open questions stay on Maestro's own ask_user" + ); +} + +#[tokio::test] +async fn user_ask_only_confirms_a_preview() { + let workspace = TempDir::new().expect("workspace"); + let tools = host_tools(workspace.path(), ApprovalMode::Selective); + let ctx = dex_loop::rehydrate(thread(), &[]); + let open = ProposedCall::new( + CallId::new("q1"), + ToolName::new(USER_ASK), + json!({"question": "Which branch?"}), + alice(), + ); + assert!(matches!(tools.policy(&ctx, &open).await, Verdict::Deny(_))); + let unbound = ProposedCall::new( + CallId::new("q2"), + ToolName::new(USER_ASK), + json!({"question": "Run it?", CONFIRMATION_FIELD: {"proposal_call_id": "nope"}}), + alice(), + ); + assert!(matches!( + tools.policy(&ctx, &unbound).await, + Verdict::Deny(_) + )); +} + +#[tokio::test] +async fn a_gated_command_runs_only_after_the_person_confirms_it() { + let workspace = TempDir::new().expect("workspace"); + let marker = workspace.path().join("marker"); + let args = touch(workspace.path()); + let tools = host_tools(workspace.path(), ApprovalMode::Selective); + let turn = Turn::start( + tools, + TurnMode::Interactive, + vec![ + call("bash", args.clone()), + ask_about_preview(), + confirmed("bash", args.clone()), + answer("created the marker"), + ], + ) + .await; + + let Exit::Asked(asked) = turn.run().await else { + panic!("the turn parks on the confirmation question"); + }; + assert!(!marker.exists(), "nothing runs before the person decides"); + let events = turn.events().await; + assert!(events.iter().any(|event| matches!( + event, + Event::ToolFinished { outcome: Outcome::Failed, output: dex_loop::Output::Text(text), .. } + if text.contains("\"needs_confirmation\"") + ))); + + turn.decide(&asked, ConfirmationDecision::Confirm).await; + assert_eq!(turn.run().await, Exit::Done); + assert!(marker.exists(), "the confirmed command ran"); + assert!(matches!( + turn.events().await.last(), + Some(Event::Final { text }) if text == "created the marker" + )); +} + +#[tokio::test] +async fn a_declined_command_never_runs() { + let workspace = TempDir::new().expect("workspace"); + let marker = workspace.path().join("marker"); + let args = touch(workspace.path()); + let tools = host_tools(workspace.path(), ApprovalMode::Selective); + let turn = Turn::start( + tools, + TurnMode::Interactive, + vec![ + call("bash", args.clone()), + ask_about_preview(), + // A model that ignores the decline and retries the bound call + // gets another preview, not an execution. + confirmed("bash", args.clone()), + answer("left it alone"), + ], + ) + .await; + + let Exit::Asked(asked) = turn.run().await else { + panic!("the turn parks on the confirmation question"); + }; + turn.decide(&asked, ConfirmationDecision::Decline).await; + assert_eq!(turn.run().await, Exit::Done); + assert!(!marker.exists(), "a declined command never runs"); +} + +#[tokio::test] +async fn an_ungated_read_runs_at_once() { + let workspace = TempDir::new().expect("workspace"); + std::fs::write(workspace.path().join("notes.txt"), "shopping list").expect("seed"); + let tools = host_tools(workspace.path(), ApprovalMode::Selective); + let turn = Turn::start( + tools, + TurnMode::Interactive, + vec![ + call( + "read", + json!({"path": workspace.path().join("notes.txt").display().to_string()}), + ), + answer("read it"), + ], + ) + .await; + assert_eq!(turn.run().await, Exit::Done); + assert!(turn.events().await.iter().any(|event| matches!( + event, + Event::ToolFinished { outcome: Outcome::Succeeded, output: dex_loop::Output::Text(text), .. } + if text.contains("shopping list") + ))); +} + +#[tokio::test] +async fn policy_matches_the_native_actor() { + let workspace = TempDir::new().expect("workspace"); + let selective = host_tools(workspace.path(), ApprovalMode::Selective); + let yolo = host_tools(workspace.path(), ApprovalMode::Yolo); + let interactive = dex_loop::rehydrate(thread(), &[]); + let proposed = + |args: Value| ProposedCall::new(CallId::new("c1"), ToolName::new("bash"), args, alice()); + + // Gated under Selective; Yolo runs it, as the native actor does. + let gated = proposed(touch(workspace.path())); + assert!(matches!( + selective.policy(&interactive, &gated).await, + Verdict::NeedsConfirmation { .. } + )); + assert_eq!(yolo.policy(&interactive, &gated).await, Verdict::Allow); + + // A safe command is not gated. + assert_eq!( + selective + .policy(&interactive, &proposed(json!({"command": "ls"}))) + .await, + Verdict::Allow + ); + + // A preview's digest is the one `Context` binds a confirmation to. + let Verdict::NeedsConfirmation { preview } = selective.policy(&interactive, &gated).await + else { + panic!("gated"); + }; + let document: Value = serde_json::from_str(&preview).expect("preview is JSON"); + assert_eq!( + document["args_digest"], + json!(args_digest(&touch(workspace.path()))) + ); +} + +#[tokio::test] +async fn a_headless_turn_refuses_a_gated_command() { + let workspace = TempDir::new().expect("workspace"); + let marker = workspace.path().join("marker"); + let tools = host_tools(workspace.path(), ApprovalMode::Selective); + let turn = Turn::start( + tools, + TurnMode::Headless, + vec![call("bash", touch(workspace.path())), answer("could not")], + ) + .await; + assert_eq!(turn.run().await, Exit::Done); + assert!(!marker.exists()); + assert!(turn.events().await.iter().any(|event| matches!( + event, + Event::ToolFinished { outcome: Outcome::Failed, output: dex_loop::Output::Text(text), .. } + if text.contains(HEADLESS_GATED) + ))); +} diff --git a/packages/local-host-rs/src/agent/mod.rs b/packages/local-host-rs/src/agent/mod.rs index c0d330b34..30b8d1b42 100644 --- a/packages/local-host-rs/src/agent/mod.rs +++ b/packages/local-host-rs/src/agent/mod.rs @@ -601,6 +601,25 @@ fn resolve_native_client( )) } +/// The local execution host, composed exactly as [`NativeAgent`] composes it, +/// for a turn that runs on the dex-loop kernel (`maestro-dex-host`) instead of +/// this crate's actor. The same `ToolExecutor`, hooks, sandbox and firewall +/// serve both loops; resolve the turn's model client with +/// [`NativeExecutionHostHandle::resolve_model`]. +pub fn dex_loop_execution_host( + config: &NativeAgentConfig, + credential_vault: CredentialVault, +) -> Result { + build_local_host( + config, + credential_vault, + None, + None, + None, + Arc::new(RwLock::new(None)), + ) +} + fn build_local_host( config: &NativeAgentConfig, credential_vault: CredentialVault, diff --git a/packages/local-host-rs/src/headless/supervisor/tests.rs b/packages/local-host-rs/src/headless/supervisor/tests.rs index da02d2074..f4c81969d 100644 --- a/packages/local-host-rs/src/headless/supervisor/tests.rs +++ b/packages/local-host-rs/src/headless/supervisor/tests.rs @@ -1727,82 +1727,6 @@ done Ok(script_path) } -#[cfg(unix)] -fn create_delayed_ack_headless_script(dir: &Path) -> std::io::Result { - let script_path = dir.join("fake-maestro-delayed-ack.sh"); - fs::write( - &script_path, - r#"#!/bin/sh -log_file="${MAESTRO_TEST_LOG:-}" -: > "$log_file" -printf '{"type":"ready","model":"test","provider":"test"}\n' -while IFS= read -r line; do - printf '%s\n' "$line" >> "$log_file" - case "$line" in - *'"type":"tool_response"'*) - sleep 0.65 - printf '{"type":"response_accepted","request_id":"gated-call"}\n' - ;; - esac -done -"#, - )?; - let mut permissions = fs::metadata(&script_path)?.permissions(); - permissions.set_mode(0o755); - fs::set_permissions(&script_path, permissions)?; - Ok(script_path) -} - -#[cfg(unix)] -#[tokio::test(flavor = "multi_thread", worker_threads = 2)] -async fn delayed_native_tool_consumption_is_queued_success_with_one_dispatch() { - let temp = tempfile::tempdir().expect("tempdir"); - let script_path = create_delayed_ack_headless_script(temp.path()).expect("script"); - let log_path = temp.path().join("messages.log"); - let mut config = SupervisorConfig::default(); - config.transport.cli_path = script_path.to_string_lossy().into_owned(); - config.transport.env.push(( - "MAESTRO_TEST_LOG".to_string(), - log_path.display().to_string(), - )); - config.auto_reconnect = false; - let mut supervisor = AgentSupervisor::new(config); - supervisor.connect().await.expect("connect"); - let _ = supervisor.recv().await.expect("connected"); - let _ = supervisor.recv().await.expect("healthy"); - let _ = supervisor.recv().await.expect("ready"); - - let (_messages, acknowledgement) = supervisor - .send_and_drain_agent_messages_with_ack(ToAgentMessage::ToolResponse { - call_id: "gated-call".to_string(), - tool_execution_id: Some("gated-execution".to_string()), - approved: true, - result: None, - }) - .expect("queue gated tool response"); - - assert_eq!(acknowledgement, ResponseAcknowledgement::Queued); - let supervisor = Arc::new(std::sync::Mutex::new(supervisor)); - let (_messages, acknowledgement) = AgentSupervisor::wait_for_response_acknowledgement_async( - Arc::clone(&supervisor), - "gated-call".to_string(), - Duration::from_millis(800), - ) - .await; - assert_eq!(acknowledgement, ResponseAcknowledgement::Consumed); - assert_eq!( - fs::read_to_string(log_path) - .expect("message log") - .lines() - .filter(|line| line.contains("\"type\":\"tool_response\"")) - .count(), - 1, - "the queued response must not be dispatched a second time while awaiting consumption" - ); - supervisor.lock().expect("supervisor").shutdown(); - tokio::time::sleep(Duration::from_millis(50)).await; -} - #[cfg(unix)] fn create_streaming_headless_script(dir: &Path) -> std::io::Result { let script_path = dir.join("fake-maestro-headless-stream.sh"); @@ -3752,3 +3676,5 @@ fn supervisor_drain_preserves_client_and_server_requests() { } mod managed_authorization_delivery; + +mod delayed_ack; diff --git a/packages/local-host-rs/src/headless/supervisor/tests/delayed_ack.rs b/packages/local-host-rs/src/headless/supervisor/tests/delayed_ack.rs new file mode 100644 index 000000000..40c697a3e --- /dev/null +++ b/packages/local-host-rs/src/headless/supervisor/tests/delayed_ack.rs @@ -0,0 +1,85 @@ +use super::*; + +#[cfg(unix)] +fn create_delayed_ack_headless_script(dir: &Path) -> std::io::Result { + let script_path = dir.join("fake-maestro-delayed-ack.sh"); + fs::write( + &script_path, + r#"#!/bin/sh +log_file="${MAESTRO_TEST_LOG:-}" +: > "$log_file" +printf '{"type":"ready","model":"test","provider":"test"}\n' +while IFS= read -r line; do + printf '%s\n' "$line" >> "$log_file" + case "$line" in + *'"type":"tool_response"'*) + while [ ! -f "$MAESTRO_TEST_ACK_GATE" ]; do sleep 0.01; done + printf '{"type":"response_accepted","request_id":"gated-call"}\n' + ;; + esac +done +"#, + )?; + let mut permissions = fs::metadata(&script_path)?.permissions(); + permissions.set_mode(0o755); + fs::set_permissions(&script_path, permissions)?; + Ok(script_path) +} + +#[cfg(unix)] +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn delayed_native_tool_consumption_is_queued_success_with_one_dispatch() { + let temp = tempfile::tempdir().expect("tempdir"); + let script_path = create_delayed_ack_headless_script(temp.path()).expect("script"); + let log_path = temp.path().join("messages.log"); + let mut config = SupervisorConfig::default(); + config.transport.cli_path = script_path.to_string_lossy().into_owned(); + config.transport.env.push(( + "MAESTRO_TEST_LOG".to_string(), + log_path.display().to_string(), + )); + let gate_path = temp.path().join("ack-release"); + config.transport.env.push(( + "MAESTRO_TEST_ACK_GATE".to_owned(), + gate_path.display().to_string(), + )); + config.auto_reconnect = false; + let mut supervisor = AgentSupervisor::new(config); + supervisor.connect().await.expect("connect"); + let _ = supervisor.recv().await.expect("connected"); + let _ = supervisor.recv().await.expect("healthy"); + let _ = supervisor.recv().await.expect("ready"); + + let (_messages, acknowledgement) = supervisor + .send_and_drain_agent_messages_with_ack(ToAgentMessage::ToolResponse { + call_id: "gated-call".to_string(), + tool_execution_id: Some("gated-execution".to_string()), + approved: true, + result: None, + }) + .expect("queue gated tool response"); + + assert_eq!(acknowledgement, ResponseAcknowledgement::Queued); + // Release the child only after observing queue admission. A fixed sleep + // raced full-workspace scheduling against the unchanged 800ms wait budget. + fs::write(&gate_path, b"release").expect("release response acknowledgement"); + let supervisor = Arc::new(std::sync::Mutex::new(supervisor)); + let (_messages, acknowledgement) = AgentSupervisor::wait_for_response_acknowledgement_async( + Arc::clone(&supervisor), + "gated-call".to_string(), + Duration::from_millis(800), + ) + .await; + assert_eq!(acknowledgement, ResponseAcknowledgement::Consumed); + assert_eq!( + fs::read_to_string(log_path) + .expect("message log") + .lines() + .filter(|line| line.contains("\"type\":\"tool_response\"")) + .count(), + 1, + "the queued response must not be dispatched a second time while awaiting consumption" + ); + supervisor.lock().expect("supervisor").shutdown(); + tokio::time::sleep(Duration::from_millis(50)).await; +} diff --git a/packages/local-host-rs/src/headless_server.rs b/packages/local-host-rs/src/headless_server.rs index 4306eb0ab..fccf49b93 100644 --- a/packages/local-host-rs/src/headless_server.rs +++ b/packages/local-host-rs/src/headless_server.rs @@ -1562,6 +1562,17 @@ fn init_headless_tracing() -> Option { ))) } +fn headless_startup_error_message(error: &anyhow::Error) -> String { + if error + .chain() + .any(|cause| cause.to_string().contains("invalid_grant")) + { + "Your Deixic Identity session expired or was revoked. Run `maestro login`, then retry this conversation.".to_owned() + } else { + format!("Failed to start Deixic Code: {error:#}") + } +} + pub async fn run_headless_server(model_override: Option) -> Result { crate::safety::disconnected_policy().map_err(anyhow::Error::msg)?; let _telemetry = init_headless_tracing(); @@ -1585,7 +1596,16 @@ pub async fn run_headless_server(model_override: Option) -> Result // Provider construction resolves the exact model route and validates its // credential before `ensure_agent` emits the first Ready boundary. - state.ensure_agent()?; + if let Err(error) = state.ensure_agent() { + emit(&FromAgentMessage::Error { + request_id: None, + message: headless_startup_error_message(&error), + fatal: true, + terminal: true, + error_type: Some(HeadlessErrorType::Fatal), + })?; + return Err(error); + } // stdin reader on a blocking thread → channel let (stdin_tx, mut stdin_rx) = mpsc::unbounded_channel::(); @@ -5351,6 +5371,15 @@ mod tests { ); } + #[test] + fn startup_identity_failure_has_a_safe_recovery_action() { + let error = anyhow::anyhow!("invalid_grant: private diagnostic") + .context("Failed to create native agent for headless server"); + let message = headless_startup_error_message(&error); + assert!(message.contains("maestro login")); + assert!(!message.contains("private diagnostic")); + } + #[tokio::test] async fn managed_evalops_headless_route_requires_identity() { if std::env::var_os("MAESTRO_HEADLESS_MANAGED_IDENTITY_REQUIRED_FIXTURE").is_some() { @@ -5385,6 +5414,19 @@ mod tests { "fixture failed: {}", String::from_utf8_lossy(&output.stderr) ); + let event = String::from_utf8_lossy(&output.stdout) + .lines() + .filter_map(|line| serde_json::from_str::(line).ok()) + .find(|event| event["type"] == "error") + .expect("fatal startup protocol event"); + assert_eq!(event["fatal"], true); + assert_eq!(event["terminal"], true); + assert!( + event["message"] + .as_str() + .unwrap() + .contains(crate::credential_mode::IDENTITY_REQUIRED_MESSAGE) + ); assert!( String::from_utf8_lossy(&output.stderr) .contains(crate::credential_mode::IDENTITY_REQUIRED_MESSAGE), diff --git a/packages/runtime-gateway-rs/src/chat.rs b/packages/runtime-gateway-rs/src/chat.rs index f1421e7a6..3133f8420 100644 --- a/packages/runtime-gateway-rs/src/chat.rs +++ b/packages/runtime-gateway-rs/src/chat.rs @@ -273,6 +273,7 @@ async fn handle_codex_app_server_chat_transport_scoped( let assistant_output = match assistant_output_result { Ok(output) => output, Err(error) => { + record_chat_error(state, session_id, error.clone()).await; send_codex_bridge_event( stream, transport, @@ -405,6 +406,21 @@ pub(crate) async fn record_chat_user_message( ))) } +pub(crate) async fn record_chat_error(state: &AppState, session_id: Option<&str>, message: String) { + let Some(session_id) = session_id else { + return; + }; + { + let mut store = state.sessions.lock().await; + let Some(session) = store.sessions.get_mut(session_id) else { + return; + }; + session.last_turn_error = Some(message); + session.updated_at = now_rfc3339(); + } + persist_session_store(state).await; +} + async fn record_chat_assistant_message(state: &AppState, session_id: Option<&str>, message: Value) { let Some(session_id) = session_id else { return; @@ -479,6 +495,9 @@ async fn append_session_message( session.title = title; } } + if message.get("role").and_then(Value::as_str) == Some("user") { + session.last_turn_error = None; + } session.messages.push(message); session.message_count = session.messages.len() as u64; let message_count = session.message_count; @@ -699,6 +718,7 @@ pub(crate) async fn handle_chat_endpoint( { Ok(session) => session.into_parts(), Err(error) => { + record_chat_error(&state, session_id.as_deref(), error.to_string()).await; send_sse( &mut stream, &serde_json::json!({ "type": "error", "message": error.to_string() }), @@ -1215,6 +1235,7 @@ pub(crate) async fn handle_chat_endpoint( } => { let request_ended = fatal || terminal; if request_ended { + record_chat_error(&state, session_id.as_deref(), message.clone()).await; let usage = last_usage.take(); if usage.is_some() { record_usage_entry( @@ -1376,6 +1397,7 @@ pub(crate) async fn handle_chat_endpoint( break; } FromAgent::TurnInterrupted { reason, .. } => { + record_chat_error(&state, session_id.as_deref(), reason.clone()).await; send_sse( &mut stream, &serde_json::json!({ "type": "error", "message": reason }), @@ -1386,6 +1408,7 @@ pub(crate) async fn handle_chat_endpoint( break; } FromAgent::ProviderError { kind, message } => { + record_chat_error(&state, session_id.as_deref(), message.clone()).await; send_sse( &mut stream, &serde_json::json!({ @@ -1403,6 +1426,7 @@ pub(crate) async fn handle_chat_endpoint( } if !terminal_sent { + record_chat_error(&state, session_id.as_deref(), "Agent stream closed before response completed".to_owned()).await; send_sse( &mut stream, &serde_json::json!({ @@ -1677,6 +1701,7 @@ pub(crate) async fn handle_chat_websocket_endpoint( { Ok(session) => session.into_parts(), Err(error) => { + record_chat_error(&state, session_id.as_deref(), error.to_string()).await; send_ws_json( &mut stream, &serde_json::json!({ "type": "error", "message": error.to_string() }), @@ -2171,6 +2196,7 @@ pub(crate) async fn handle_chat_websocket_endpoint( } => { let request_ended = fatal || terminal; if request_ended { + record_chat_error(&state, session_id.as_deref(), message.clone()).await; let usage = last_usage.take(); if usage.is_some() { record_usage_entry( @@ -2323,6 +2349,7 @@ pub(crate) async fn handle_chat_websocket_endpoint( break; } FromAgent::TurnInterrupted { reason, .. } => { + record_chat_error(&state, session_id.as_deref(), reason.clone()).await; send_ws_json( &mut stream, &serde_json::json!({ "type": "error", "message": reason }), @@ -2333,6 +2360,7 @@ pub(crate) async fn handle_chat_websocket_endpoint( break; } FromAgent::ProviderError { kind, message } => { + record_chat_error(&state, session_id.as_deref(), message.clone()).await; send_ws_json( &mut stream, &serde_json::json!({ @@ -2350,6 +2378,7 @@ pub(crate) async fn handle_chat_websocket_endpoint( } if !terminal_sent { + record_chat_error(&state, session_id.as_deref(), "Agent stream closed before response completed".to_owned()).await; send_ws_json( &mut stream, &serde_json::json!({ diff --git a/packages/runtime-gateway-rs/src/codex_bridge.rs b/packages/runtime-gateway-rs/src/codex_bridge.rs index 70e92cf05..f89981b7d 100644 --- a/packages/runtime-gateway-rs/src/codex_bridge.rs +++ b/packages/runtime-gateway-rs/src/codex_bridge.rs @@ -1446,6 +1446,16 @@ async fn cleanup_codex_headless_approvals( } } +pub(crate) fn codex_headless_hello() -> Value { + serde_json::json!({ + "type": "hello", + "protocol_version": maestro_local_host::headless::HEADLESS_PROTOCOL_VERSION, + "client_info": { "name": "maestro-rust-control-plane" }, + "capabilities": { "server_requests": ["approval"] }, + "role": "controller" + }) +} + #[allow(clippy::too_many_arguments)] pub(crate) async fn run_codex_app_server_headless_cli( stream: &mut TcpStream, @@ -1501,8 +1511,9 @@ pub(crate) async fn run_codex_app_server_headless_cli( .take() .ok_or_else(|| "failed to capture Codex headless stderr".to_string())?; let stderr_task = tokio::spawn(async move { - let mut bytes = Vec::new(); - stderr.read_to_end(&mut bytes).await.map(|_| bytes) + // The protocol carries actionable errors. Drain diagnostics without + // retaining an unbounded buffer or disclosing provider credentials. + tokio::io::copy(&mut stderr, &mut tokio::io::sink()).await }); let pending_approval_ids = Arc::new(Mutex::new(Vec::new())); let approval_run_id = codex_headless_run_id(); @@ -1511,21 +1522,7 @@ pub(crate) async fn run_codex_app_server_headless_cli( let mut lines = BufReader::new(stdout).lines(); let request_result = async { - write_codex_headless_message( - &mut stdin, - &serde_json::json!({ - "type": "hello", - "protocol_version": "2026-08-01", - "client_info": { - "name": "maestro-rust-control-plane" - }, - "capabilities": { - "server_requests": ["approval"] - }, - "role": "controller" - }), - ) - .await?; + write_codex_headless_message(&mut stdin, &codex_headless_hello()).await?; write_codex_headless_message( &mut stdin, &serde_json::json!({ diff --git a/packages/runtime-gateway-rs/src/sessions.rs b/packages/runtime-gateway-rs/src/sessions.rs index 8dfe7edc2..01cafb97d 100644 --- a/packages/runtime-gateway-rs/src/sessions.rs +++ b/packages/runtime-gateway-rs/src/sessions.rs @@ -56,6 +56,8 @@ pub(super) struct SessionRecord { pub(super) log_group_id: Option, #[serde(default)] pub(super) messages: Vec, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub(super) last_turn_error: Option, } #[derive(Clone, Serialize, Deserialize)] @@ -522,6 +524,7 @@ pub(super) fn create_session_record(title: Option, owner: Option tags: Vec::new(), log_group_id: None, messages: Vec::new(), + last_turn_error: None, } } @@ -1888,6 +1891,9 @@ pub(super) fn session_summary_value(session: &SessionRecord) -> Value { pub(super) fn session_full_value(session: &SessionRecord) -> Value { let mut value = session_summary_value(session); + if let Some(error) = &session.last_turn_error { + value["lastTurnError"] = Value::String(error.clone()); + } value["messages"] = Value::Array( session .messages diff --git a/packages/runtime-gateway-rs/src/tests.rs b/packages/runtime-gateway-rs/src/tests.rs index 1f9c2ede5..7e5424603 100644 --- a/packages/runtime-gateway-rs/src/tests.rs +++ b/packages/runtime-gateway-rs/src/tests.rs @@ -2711,6 +2711,7 @@ fn test_session_record(id: &str) -> SessionRecord { tags: Vec::new(), log_group_id: None, messages: Vec::new(), + last_turn_error: None, } } @@ -13227,6 +13228,7 @@ async fn delete_session_subpath_returns_404_without_removing_session() { tags: Vec::new(), log_group_id: None, messages: Vec::new(), + last_turn_error: None, }; let state = AppState { config: Arc::new(Config { @@ -14530,3 +14532,53 @@ async fn platform_a2a_push_evicts_terminal_payloads_and_replay_history() { #[path = "tests/turn_diffs.rs"] mod turn_diff_tests; + +#[tokio::test] +async fn failed_chat_turn_survives_session_store_reload_and_retry_clears_it() { + let session = test_session_record("failed-turn"); + let state = test_app_state_with_sessions(HashMap::from([(session.id.clone(), session)])); + chat::record_chat_error( + &state, + Some("failed-turn"), + "Run `maestro login`".to_owned(), + ) + .await; + let store = state.sessions.lock().await.clone(); + let reloaded: SessionStore = + serde_json::from_slice(&serde_json::to_vec(&store).unwrap()).unwrap(); + let value = sessions::session_full_value(&reloaded.sessions["failed-turn"]); + assert_eq!(value["lastTurnError"], "Run `maestro login`"); + assert_eq!( + value["messageCount"], 0, + "a failure is not an assistant answer" + ); + let auth = AuthContext { + unrestricted: true, + ..AuthContext::default() + }; + let chat: ChatRequest = serde_json::from_value(serde_json::json!({ + "sessionId": "failed-turn", "messages": [{"role": "user", "content": "Retry"}] + })) + .unwrap(); + chat::record_chat_user_message(&state, &chat, &auth) + .await + .unwrap(); + assert!( + state.sessions.lock().await.sessions["failed-turn"] + .last_turn_error + .is_none() + ); +} + +#[test] +fn desktop_headless_handshake_uses_the_native_owner_protocol() { + let hello = codex_bridge::codex_headless_hello(); + let _: maestro_local_host::headless::ToAgentMessage = + serde_json::from_value(hello.clone()).unwrap(); + assert!( + maestro_local_host::headless::messages::SUPPORTED_CLIENT_PROTOCOL_VERSIONS + .contains(&hello["protocol_version"].as_str().unwrap()) + ); + assert_eq!(hello["role"], "controller"); + assert_eq!(hello["capabilities"]["server_requests"][0], "approval"); +} diff --git a/packages/runtime-rs/src/agent/mod.rs b/packages/runtime-rs/src/agent/mod.rs index 8dda29301..d5663bba2 100644 --- a/packages/runtime-rs/src/agent/mod.rs +++ b/packages/runtime-rs/src/agent/mod.rs @@ -54,6 +54,7 @@ pub use message_queue::{ pub use model_dynamics::{ BoostStatus, ModelChoice, ModelDynamicsConfig, TaskDifficulty, ThinkingLevel, }; +pub use native::loop_policy; pub use native::{ ExternalToolSchemaPolicy, MaxTokensSource, NativeAgent, NativeAgentConfig, REQUEST_CONTEXT_SAFETY_TOKENS, RuntimeAuditSnapshot, ToolResponseConsumption, diff --git a/packages/runtime-rs/src/agent/native.rs b/packages/runtime-rs/src/agent/native.rs index 4a13ad256..823122d5c 100644 --- a/packages/runtime-rs/src/agent/native.rs +++ b/packages/runtime-rs/src/agent/native.rs @@ -453,6 +453,7 @@ mod model_dynamics; mod provider_history; mod read_only_results; use provider_history::sanitize_semantic_conversation; +pub mod loop_policy; mod provider_admission; mod provider_loop; mod provider_payload; diff --git a/packages/runtime-rs/src/agent/native/loop_policy.rs b/packages/runtime-rs/src/agent/native/loop_policy.rs new file mode 100644 index 000000000..8e7a74f97 --- /dev/null +++ b/packages/runtime-rs/src/agent/native/loop_policy.rs @@ -0,0 +1,75 @@ +//! The native actor's call classification, exposed for hosts that run +//! Maestro tools on the dex-loop kernel (`maestro-dex-host`). +//! +//! Each function wraps the one the native actor itself calls, so a tool call +//! is gated, firewalled and classified as read-only identically on both +//! loops. Turn-local refusal memory is the calling loop's own concern: the +//! kernel refuses an exact repeat of an uncertain or refused call itself. + +use serde_json::Value; + +use super::read_only_tools::is_native_parallel_read_only_tool_call; +use super::tool_execution::{deferred_firewall_verdict, tool_requires_approval}; +use crate::agent::denial_memory::DenialMemory; +use crate::agent::native_host::{ + ApprovalMode, NativeExecutionHostHandle, NativeFirewallVerdict, NativeToolAnnotations, +}; +use crate::agent::safety::WorkflowStateSnapshot; + +/// The action firewall's verdict for one call, as the native actor computes it. +#[must_use] +pub fn firewall_verdict( + host: &NativeExecutionHostHandle, + tool_name: &str, + args: &Value, + workflow: &WorkflowStateSnapshot, + annotations: Option<&NativeToolAnnotations>, + is_external_tool: bool, +) -> NativeFirewallVerdict { + deferred_firewall_verdict( + host, + tool_name, + args, + workflow, + annotations, + is_external_tool, + ) +} + +/// Whether the call needs a human decision before it runs under `mode`. +#[must_use] +pub fn approval_required( + mode: ApprovalMode, + is_external_tool: bool, + firewall: &NativeFirewallVerdict, + host: &NativeExecutionHostHandle, + tool_name: &str, + args: &Value, +) -> bool { + tool_requires_approval( + mode, + is_external_tool, + firewall, + host, + tool_name, + args, + &DenialMemory::new(), + ) + .requires_approval() +} + +/// Whether the call may run in a parallel read-only wave. +#[must_use] +pub fn parallel_read_only( + tool_name: &str, + requires_approval: bool, + annotations: Option<&NativeToolAnnotations>, + explicit_inline_read_only: bool, +) -> bool { + is_native_parallel_read_only_tool_call( + tool_name, + requires_approval, + annotations, + explicit_inline_read_only, + ) +} diff --git a/packages/tui-rs/src/app.rs b/packages/tui-rs/src/app.rs index 6a190919c..ab72c118f 100644 --- a/packages/tui-rs/src/app.rs +++ b/packages/tui-rs/src/app.rs @@ -3094,7 +3094,7 @@ Always use tools when they would be helpful. Be concise and direct in your respo let mut message = self .state .locale - .format("Failed to create agent: {0}", &[(e).to_string()]); + .format("Failed to create agent: {0}", &[format!("{e:#}")]); if crate::codex_auth::read_codex_auth().is_none() && std::env::var_os("OPENAI_API_KEY").is_none() && std::env::var_os("OPENAI_CODEX_TOKEN").is_none() diff --git a/packages/tui-rs/src/app/prompt_queue.rs b/packages/tui-rs/src/app/prompt_queue.rs index 77a1dc1ad..fd0c95ace 100644 --- a/packages/tui-rs/src/app/prompt_queue.rs +++ b/packages/tui-rs/src/app/prompt_queue.rs @@ -1,6 +1,12 @@ use super::*; impl App { + fn report_agent_unavailable(&mut self) { + if self.state.error.is_none() { + self.state.error = Some("Deixic Code could not start. Run `maestro setup --live` to check sign-in and provider access, then restart.".to_owned()); + } + } + pub(super) async fn handle_side_question(&mut self, question: String) -> Result { if question.trim().is_empty() { return Ok(false); @@ -14,12 +20,7 @@ impl App { .await; } let Some(agent) = &self.native_agent else { - self.state.error = Some( - self.state - .locale - .translate("Agent not initialized") - .to_string(), - ); + self.report_agent_unavailable(); return Ok(false); }; self.state.busy = true; @@ -157,12 +158,7 @@ impl App { self.auto_activate_skills_for_queued_prompt(&content, queue_id); } let Some(agent) = &self.native_agent else { - self.state.error = Some( - self.state - .locale - .translate("Agent not initialized") - .to_string(), - ); + self.report_agent_unavailable(); return Ok(false); }; if let Err(e) = agent @@ -593,9 +589,14 @@ impl App { /// Submit a prompt to the agent pub(super) async fn submit_prompt(&mut self, content: String) -> Result<()> { - let _ = self + let retry = content.clone(); + if !self .submit_prompt_with_kind(content, PromptKind::Prompt) - .await?; + .await? + && self.state.input().is_empty() + { + self.state.set_input(&retry); + } Ok(()) } @@ -655,6 +656,11 @@ impl App { return Ok(false); } + if self.native_agent.is_none() { + self.report_agent_unavailable(); + return Ok(false); + } + self.auto_activate_skills_for_prompt(&content); // Snapshot the worktree so `/rewind files` can restore what this turn changes. @@ -698,12 +704,7 @@ impl App { } return Ok(true); } - self.state.error = Some( - self.state - .locale - .translate("Agent not initialized") - .to_string(), - ); + self.report_agent_unavailable(); self.state.busy = false; Ok(false) } diff --git a/packages/tui-rs/src/app/tests.rs b/packages/tui-rs/src/app/tests.rs index 28d795fe3..1528a41a5 100644 --- a/packages/tui-rs/src/app/tests.rs +++ b/packages/tui-rs/src/app/tests.rs @@ -2351,21 +2351,6 @@ async fn test_handle_config_event_reloads_keybindings() { ); } -#[tokio::test] -async fn test_tab_submits_when_idle_with_non_shell_input() { - let mut app = new_test_app(); - app.state.set_input("ship it"); - - app.handle_key(KeyCode::Tab, CrosstermModifiers::NONE) - .await - .unwrap(); - - assert_eq!(app.state.input(), ""); - let last = app.state.messages.last().expect("user message"); - assert_eq!(last.role, MessageRole::User); - assert_eq!(last.content, "ship it"); -} - #[tokio::test] async fn test_tab_does_not_submit_idle_shell_draft() { let mut app = new_test_app(); @@ -6696,3 +6681,5 @@ mod settings_experiments; mod session_restore; mod checkpoints; + +mod startup_recovery; diff --git a/packages/tui-rs/src/app/tests/startup_recovery.rs b/packages/tui-rs/src/app/tests/startup_recovery.rs new file mode 100644 index 000000000..c40f66eaa --- /dev/null +++ b/packages/tui-rs/src/app/tests/startup_recovery.rs @@ -0,0 +1,66 @@ +use super::*; + +#[tokio::test] +async fn test_tab_submits_when_idle_with_non_shell_input() { + let mut app = new_test_app(); + // This binding tests an admitted turn; startup failures have a separate + // regression below and must never invent a user-history entry. + let temp = tempdir().unwrap(); + let (agent, _events) = crate::agent::NativeAgent::new_with_test_client( + crate::agent::NativeAgentConfig { + model: "openai/gpt-4o".into(), + cwd: temp.path().display().to_string(), + ..Default::default() + }, + crate::ai::UnifiedClient::OpenAI( + crate::ai::OpenAiClient::with_base_url("fixture", "http://127.0.0.1:1/v1").unwrap(), + ), + ) + .unwrap(); + app.native_agent = Some(agent); + app.state.set_input("ship it"); + + app.handle_key(KeyCode::Tab, CrosstermModifiers::NONE) + .await + .unwrap(); + + assert_eq!(app.state.input(), ""); + let last = app.state.messages.last().expect("user message"); + assert_eq!(last.role, MessageRole::User); + assert_eq!(last.content, "ship it"); +} + +#[tokio::test] +async fn failed_startup_keeps_recovery_error_and_does_not_record_phantom_turn() { + let mut app = new_test_app(); + app.state.error = Some("Identity session revoked. Run `maestro login`.".to_owned()); + let count = app.state.messages.len(); + assert!( + !app.submit_prompt_with_kind("Hello".to_owned(), PromptKind::Prompt) + .await + .unwrap() + ); + assert_eq!( + app.state.error.as_deref(), + Some("Identity session revoked. Run `maestro login`.") + ); + assert_eq!(app.state.messages.len(), count); + assert!(!app.state.busy); +} + +#[tokio::test] +async fn tab_after_failed_startup_keeps_input_and_recovery_action() { + let mut app = new_test_app(); + app.state.error = Some("Identity expired. Run `maestro login`.".to_owned()); + app.state.set_input("retry this request"); + let count = app.state.messages.len(); + app.handle_key(KeyCode::Tab, CrosstermModifiers::NONE) + .await + .unwrap(); + assert_eq!(app.state.input(), "retry this request"); + assert_eq!(app.state.messages.len(), count); + assert_eq!( + app.state.error.as_deref(), + Some("Identity expired. Run `maestro login`.") + ); +} diff --git a/packages/tui-rs/src/setup_cli.rs b/packages/tui-rs/src/setup_cli.rs index ec48195c3..d705aa26d 100644 --- a/packages/tui-rs/src/setup_cli.rs +++ b/packages/tui-rs/src/setup_cli.rs @@ -37,7 +37,7 @@ struct SetupOptions { fn parse_options(args: &[String]) -> Result { let mut options = SetupOptions { json: false, - live: false, + live: true, model: None, platform: false, byok: false, @@ -47,6 +47,7 @@ fn parse_options(args: &[String]) -> Result { match args[index].as_str() { "--json" => options.json = true, "--live" => options.live = true, + "--offline" => options.live = false, "--platform" => options.platform = true, "--byok" => options.byok = true, "--model" => { @@ -135,9 +136,15 @@ pub fn build_setup_report(doctor: DoctorReport) -> SetupReport { } } + if !doctor.live_requested && next_steps.is_empty() { + push_step(&mut next_steps, "verify-credentials", "deixic-code setup --live", + "Stored credentials have not been verified. Check the Identity session and provider before starting.".to_owned()); + } + SetupReport { schema_version: SETUP_SCHEMA_VERSION, - ready: doctor.ok + ready: doctor.live_requested + && doctor.ok && !doctor.checks.iter().any(|check| { check.status == CheckStatus::Fail || (check.id == "managed_inference" @@ -173,7 +180,7 @@ pub async fn run_setup(args: &[String]) -> Result { let options = match parse_options(args) { Ok(options) => options, Err(error) if error.to_string() == "help" => { - println!("{}", crate::localization::cli_locale().format("Usage: deixic-code setup [--json] [--live] [--model ] [--platform|--byok]", &[])); + println!("{}", crate::localization::cli_locale().format("Usage: deixic-code setup [--json] [--live|--offline] [--model ] [--platform|--byok]", &[])); return Ok(0); } Err(error) => return Err(error), @@ -445,7 +452,7 @@ mod tests { DoctorReport { schema_version: crate::doctor::REPORT_SCHEMA_VERSION, ok: true, - live_requested: false, + live_requested: true, selected_model: SelectedModelReport { requested: format!("{provider}/model"), provider: provider.to_owned(), @@ -531,6 +538,25 @@ mod tests { assert_eq!(result.next_steps[0].command, "deixic-code codex login"); } + #[test] + fn offline_stored_credentials_do_not_claim_live_readiness() { + let mut doctor = report( + "evalops", + vec![check( + "credential_mode", + CheckStatus::Pass, + "platform: org org_1 via EvalOps identity", + None, + )], + ); + doctor.live_requested = false; + let result = build_setup_report(doctor); + assert!(!result.ready); + assert_eq!(result.next_steps[0].command, "deixic-code setup --live"); + assert!(parse_options(&["--json".to_owned()]).unwrap().live); + assert!(!parse_options(&["--offline".to_owned()]).unwrap().live); + } + #[test] fn platform_session_makes_setup_ready_without_local_keys() { let result = build_setup_report(report( diff --git a/packages/ui-rs/src/translations.rs b/packages/ui-rs/src/translations.rs index fbdd5b306..85599cc99 100644 --- a/packages/ui-rs/src/translations.rs +++ b/packages/ui-rs/src/translations.rs @@ -23031,14 +23031,14 @@ pub(crate) const MESSAGES: &[(&str, [&str; 6])] = &[ ], ), ( - "Usage: deixic-code setup [--json] [--live] [--model ] [--platform|--byok]", - [ - "Usage: deixic-code setup [--json] [--live] [--model ] [--platform|--byok]", - "Utilisation: deixic-code setup [--json] [--live] [--model ] [--platform|--byok]", - "Verwendung: deixic-code setup [--json] [--live] [--model ] [--platform|--byok]", - "使用方法: deixic-code setup [--json] [--live] [--model ] [--platform|--byok]", - "제품 정보: deixic-code setup [--json] [--live] [--model ] [--platform|--byok]", - "使用量: deixic-code setup [--json] [--live] [--model ] [--platform|--byok]", + "Usage: deixic-code setup [--json] [--live|--offline] [--model ] [--platform|--byok]", + [ + "Usage: deixic-code setup [--json] [--live|--offline] [--model ] [--platform|--byok]", + "Utilisation: deixic-code setup [--json] [--live|--offline] [--model ] [--platform|--byok]", + "Verwendung: deixic-code setup [--json] [--live|--offline] [--model ] [--platform|--byok]", + "使用方法: deixic-code setup [--json] [--live|--offline] [--model ] [--platform|--byok]", + "제품 정보: deixic-code setup [--json] [--live|--offline] [--model ] [--platform|--byok]", + "使用量: deixic-code setup [--json] [--live|--offline] [--model ] [--platform|--byok]", ], ), (