Merge latest origin/main into feat/sandbox-bash-contract

This commit is contained in:
Bryan Helmkamp 2026-07-24 22:45:35 -04:00
commit 2f84b67558
No known key found for this signature in database
26 changed files with 1134 additions and 148 deletions

View file

@ -1157,6 +1157,45 @@ Emitted when a tool call finishes.
| `output` | any | Tool output (string or structured) |
| `is_error` | boolean | Whether the tool returned an error |
### `agent.tool.process.completed`
Subordinate diagnostic for a tool call that ran a process, emitted between
`agent.tool.started` and `agent.tool.completed`. It explains the underlying
process outcome; `agent.tool.completed.is_error` remains the protocol and UI
truth. Absent when the tool never produced a process result (setup, transport,
or launch failure) and when the tool ran without a session-bound emitter.
```json
{
"id": "...", "ts": "...", "run_id": "...",
"event": "agent.tool.process.completed",
"node_id": "code", "node_label": "code",
"session_id": "ses_abc",
"tool_call_id": "call_abc123",
"properties": {
"exit_code": 7,
"termination": "exited",
"duration_ms": 812,
"streams_separated": true,
"exec_output_tail": {"stdout": "...", "stderr": "..."},
"visit": 1
}
}
```
| Property | Type | Description |
|----------|------|-------------|
| `exit_code` | integer | Process exit code; omitted for timeout and cancellation |
| `termination` | string | `exited`, `timed_out`, or `cancelled` |
| `duration_ms` | integer | Process duration |
| `streams_separated` | boolean | `false` when the provider could not separate stdout from stderr; the combined output is then in `exec_output_tail.stdout` |
| `exec_output_tail` | object | Bounded, redacted output tails; omitted when both streams were empty |
| `exec_output_tail.stdout` | string | Bounded stdout tail, or combined-output tail when `streams_separated` is `false`; omitted when empty |
| `exec_output_tail.stderr` | string | Bounded stderr tail; omitted when empty |
| `exec_output_tail.stdout_truncated` | boolean | `true` when earlier stdout bytes were omitted; omitted when `false` |
| `exec_output_tail.stderr_truncated` | boolean | `true` when earlier stderr bytes were omitted; omitted when `false` |
| `visit` | integer | Stage visit |
### `agent.error`
Emitted when the agent encounters an error.

View file

@ -54,7 +54,7 @@ When a run resumes after a node was cancelled or lost mid-flight, the replay now
</Accordion>
<Accordion title="Improvements">
- Added `claude-opus-5` to the first-party Anthropic model catalog; the `opus` and `claude-opus` aliases now resolve to Opus 5
- Added `claude-opus-5` to the first-party Anthropic and optional OpenRouter model catalogs; the `opus` and `claude-opus` aliases now resolve to Opus 5
- Added `gpt-sol`, `gpt-terra`, and `gpt-luna` aliases for GPT-5.6 offerings
- Added portable `glm`, `glm52`, `glm5.2`, `deepseek`, and `deepseek-flash` aliases across direct and OpenRouter offerings
</Accordion>

View file

@ -48,7 +48,7 @@ The built-in catalog gives OpenRouter offerings the same human-facing model slug
| Fabro model slug | OpenRouter API ID / notes |
| --- | --- |
| `claude-fable-5`, `claude-opus-4-8`, `claude-opus-4-7` | Matching `anthropic/...` API IDs; Anthropic-style cache billing |
| `claude-fable-5`, `claude-opus-5`, `claude-opus-4-8`, `claude-opus-4-7` | Matching `anthropic/...` API IDs; Anthropic-style cache billing |
| `claude-sonnet-4-6` | `anthropic/claude-sonnet-4.6`; provider default |
| `claude-haiku-4-5` | `anthropic/claude-haiku-4.5`; provider small default |
| `gpt-5.6-sol`, `gpt-5.6-terra`, `gpt-5.6-luna`, `gpt-5.4`, `gpt-5.5` | Matching `openai/...` API IDs |

View file

@ -30,8 +30,8 @@ use crate::native_tool::NativeTool;
use crate::sandbox::{GrepOptions, format_lines_numbered};
use crate::tool_registry::{RegisteredTool, ToolSource};
use crate::tools::{
DEFAULT_READ_LINES, execute_grep, execute_shell_command, grep_result_path, make_edit_file_tool,
optional_usize_arg, required_str,
DEFAULT_READ_LINES, emit_shell_process_completed, execute_grep, execute_shell_command,
grep_result_path, make_edit_file_tool, optional_usize_arg, required_str,
};
const DEFAULT_GREP_RESULTS: usize = 250;
@ -120,7 +120,8 @@ explicitly asked. Never run commands requiring superuser privileges unless expli
None => default_timeout_ms,
};
let result = execute_shell_command(&ctx, command, timeout_ms, cwd).await?;
let streaming = execute_shell_command(&ctx, command, timeout_ms, cwd).await?;
let result = &streaming.result;
let mut out = String::new();
if result.is_timed_out() {
@ -141,7 +142,9 @@ explicitly asked. Never run commands requiring superuser privileges unless expli
}
let _ = write!(out, "Command failed with exit code: {code}");
}
Ok(out)
let is_success = result.is_success();
emit_shell_process_completed(&ctx, streaming).await;
if is_success { Ok(out) } else { Err(out) }
})
}),
source: ToolSource::Native,
@ -699,7 +702,7 @@ mod tests {
tool_ctx,
)
.await
.unwrap();
.expect_err("a timeout is a failed tool result");
assert!(output.starts_with("Command timed out.\n"), "{output}");
assert_eq!(*env.captured_timeout.lock().unwrap(), Some(7_000));
@ -707,12 +710,9 @@ mod tests {
Some("/repo".to_string())
]);
assert_eq!(*env.captured_env_vars.lock().unwrap(), Some(tool_env));
assert!(
env.captured_command
.lock()
.unwrap()
.as_deref()
.is_some_and(|command| command.starts_with("exec 2>&1\n"))
assert_eq!(
env.captured_command.lock().unwrap().as_deref(),
Some("echo $TOKEN")
);
}
}

View file

@ -558,6 +558,7 @@ mod tests {
use async_trait::async_trait;
use fabro_llm::types::{ToolCall, ToolDefinition};
use fabro_model::AgentProfileKind;
use tokio::sync::broadcast;
use super::*;
use crate::config::{
@ -570,11 +571,13 @@ mod tests {
AgentToolRuntime, register_question_tools,
};
use crate::read_before_write_sandbox::ReadBeforeWriteSandbox;
use crate::test_support::MutableMockSandbox;
use crate::test_support::{MockSandbox, MutableMockSandbox};
use crate::tool_registry::{RegisteredTool, ToolContext, ToolRegistry, ToolSource};
use crate::tools::{
make_edit_file_tool, make_grep_tool, make_read_file_tool, make_write_file_tool,
make_edit_file_tool, make_grep_tool, make_read_file_tool, make_shell_tool,
make_write_file_tool,
};
use crate::types::SessionEvent;
struct NamedPolicy {
decisions: HashMap<String, ToolAccess>,
@ -1294,4 +1297,173 @@ mod tests {
assert!(!result.is_error);
}
fn shell_sandbox(result: fabro_sandbox::ExecResult) -> Arc<dyn Sandbox> {
Arc::new(MockSandbox {
exec_result: result,
..Default::default()
})
}
fn exited(exit_code: i32) -> fabro_sandbox::ExecResult {
fabro_sandbox::ExecResult {
stdout: "out".into(),
stderr: "err".into(),
exit_code: Some(exit_code),
termination: fabro_types::CommandTermination::Exited,
duration_ms: 12,
}
}
fn cancelled() -> fabro_sandbox::ExecResult {
fabro_sandbox::ExecResult {
stdout: "out".into(),
stderr: String::new(),
exit_code: None,
termination: fabro_types::CommandTermination::Cancelled,
duration_ms: 12,
}
}
async fn run_shell_tool(
exec_result: fabro_sandbox::ExecResult,
hooks: Option<&Arc<dyn ToolHookCallback>>,
emitter: &Emitter,
) -> ToolResult {
let mut registry = ToolRegistry::new();
registry.register(make_shell_tool());
let tc = make_tool_call(
"shell",
"call_1",
serde_json::json!({"command": "make test"}),
);
execute_and_emit_one_tool(
&tc,
&registry,
shell_sandbox(exec_result),
hooks,
CancellationToken::new(),
&SessionOptions::default(),
emitter,
"test-session",
"test-session",
None,
)
.await
}
fn drain(receiver: &mut broadcast::Receiver<SessionEvent>) -> Vec<SessionEvent> {
let mut events = Vec::new();
while let Ok(event) = receiver.try_recv() {
events.push(event);
}
events
}
#[tokio::test]
async fn shell_nonzero_exit_becomes_an_error_tool_result() {
let emitter = Emitter::new();
let result = run_shell_tool(exited(7), None, &emitter).await;
assert!(result.is_error);
assert!(
result.content.as_str().unwrap().contains("Exit code: 7"),
"got: {}",
result.content
);
}
#[tokio::test]
async fn shell_exit_zero_remains_a_successful_tool_result() {
let emitter = Emitter::new();
let result = run_shell_tool(exited(0), None, &emitter).await;
assert!(!result.is_error);
}
#[tokio::test]
async fn shell_failure_emits_started_then_process_then_completed() {
let emitter = Emitter::new();
let mut receiver = emitter.subscribe();
run_shell_tool(exited(7), None, &emitter).await;
let events = drain(&mut receiver);
let names: Vec<&str> = events
.iter()
.filter_map(|event| match &event.event {
AgentEvent::ToolCallStarted { .. } => Some("started"),
AgentEvent::ToolProcessCompleted { .. } => Some("process"),
AgentEvent::ToolCallCompleted { .. } => Some("completed"),
_ => None,
})
.collect();
assert_eq!(names, vec!["started", "process", "completed"]);
for event in &events {
assert_eq!(event.session_id, "test-session");
}
let process = events
.iter()
.find(|event| matches!(event.event, AgentEvent::ToolProcessCompleted { .. }))
.expect("process event");
assert_eq!(process.tool_call_id.as_deref(), Some("call_1"));
match &process.event {
AgentEvent::ToolProcessCompleted {
exit_code,
termination,
..
} => {
assert_eq!(*exit_code, Some(7));
assert_eq!(*termination, fabro_types::CommandTermination::Exited);
}
other => panic!("expected a process event, got {other:?}"),
}
let completed = events
.iter()
.find_map(|event| match &event.event {
AgentEvent::ToolCallCompleted {
tool_call_id,
is_error,
..
} => Some((tool_call_id.clone(), *is_error)),
_ => None,
})
.expect("tool completed event");
assert_eq!(completed, ("call_1".to_string(), true));
}
#[tokio::test]
async fn shell_failure_runs_only_the_failure_hook() {
for exec_result in [exited(7), cancelled()] {
let mock = Arc::new(MockHookCallback::new(ToolHookDecision::Proceed));
let hooks: Arc<dyn ToolHookCallback> = mock.clone();
run_shell_tool(exec_result, Some(&hooks), &Emitter::new()).await;
assert_eq!(mock.post_failure_calls.lock().unwrap().len(), 1);
assert!(mock.post_calls.lock().unwrap().is_empty());
}
}
#[tokio::test]
async fn shell_success_runs_only_the_success_hook() {
let mock = Arc::new(MockHookCallback::new(ToolHookDecision::Proceed));
let hooks: Arc<dyn ToolHookCallback> = mock.clone();
run_shell_tool(exited(0), Some(&hooks), &Emitter::new()).await;
assert_eq!(mock.post_calls.lock().unwrap().len(), 1);
assert!(mock.post_failure_calls.lock().unwrap().is_empty());
}
#[test]
fn truncation_preserves_tool_call_id_and_error_state() {
let result = ToolResult::error("call_1", "x".repeat(60_000));
let truncated = truncate_tool_result(&result, "shell", &SessionOptions::default());
assert_eq!(truncated.tool_call_id, "call_1");
assert!(truncated.is_error);
assert!(truncated.content.as_str().unwrap().len() < 60_000);
}
}

View file

@ -8,10 +8,12 @@ use fabro_model::ModelHandle;
#[cfg(test)]
use fabro_static::EnvVars;
use futures::{StreamExt, stream};
use tokio::task;
use crate::config::NativeToolOptions;
use crate::sandbox::{ExecResult, GrepOptions};
use crate::sandbox::{ExecStreamingResult, GrepOptions};
use crate::tool_registry::{RegisteredTool, ToolContext, ToolRegistry, ToolSource};
use crate::types::AgentEvent;
const MAX_WEB_FETCH_BYTES: usize = 100 * 1024;
const MAX_READ_MANY_FILES_CONCURRENCY: usize = 8;
@ -259,29 +261,27 @@ pub fn make_shell_tool_with_options(options: &NativeToolOptions) -> RegisteredTo
.unwrap_or(default_timeout)
.min(max_timeout);
let result = execute_shell_command(&ctx, command, timeout_ms, None).await?;
let streaming = execute_shell_command(&ctx, command, timeout_ms, None).await?;
let mut output = String::new();
if result.is_timed_out() {
output.push_str("Command timed out.\n");
} else if result.is_cancelled() {
output.push_str("Command cancelled.\n");
let text = render_shell_result(&streaming);
let is_success = streaming.result.is_success();
emit_shell_process_completed(&ctx, streaming).await;
if is_success {
Ok(text)
} else {
Err(text)
}
let _ = write!(
output,
"Exit code: {}\noutput:\n{}",
result
.exit_code
.map_or_else(|| "none".to_string(), |code| code.to_string()),
result.stdout
);
Ok(output)
})
}),
source: ToolSource::Native,
}
}
/// Prefix for shell failures that never produced an `ExecResult`, so the model
/// can distinguish missing process diagnostics from a reported process failure.
const SHELL_NO_PROCESS_RESULT: &str = "Shell command produced no process result";
/// Execute a shell command with the session's environment and cancellation
/// plumbing. Provider profiles can vary their wire schema and result
/// rendering without accidentally bypassing those shared semantics.
@ -290,23 +290,88 @@ pub(crate) async fn execute_shell_command(
command: &str,
timeout_ms: u64,
cwd: Option<&str>,
) -> Result<ExecResult, String> {
let command = format!("exec 2>&1\n{command}");
let tool_env = ctx.resolve_tool_env().await.map_err(|e| format!("{e:#}"))?;
) -> Result<ExecStreamingResult, String> {
let tool_env = ctx
.resolve_tool_env()
.await
.map_err(|e| format!("{SHELL_NO_PROCESS_RESULT}: {e:#}"))?;
tracing::debug!(
env_var_count = tool_env.as_ref().map_or(0, std::collections::HashMap::len),
"Injecting sandbox env vars into tool execution"
);
ctx.env
.exec_command(
&command,
timeout_ms,
.exec_command_streaming(
command,
Some(timeout_ms),
cwd,
tool_env.as_ref(),
Some(ctx.cancel.clone()),
None,
)
.await
.map_err(|e| e.display_with_causes())
.map_err(|e| format!("{SHELL_NO_PROCESS_RESULT}: {}", e.display_with_causes()))
}
/// Emit the subordinate process outcome after model-facing output has been
/// rendered. Consumes the raw result so redaction does not require cloning
/// potentially large process output.
pub(crate) async fn emit_shell_process_completed(
ctx: &ToolContext,
streaming: ExecStreamingResult,
) {
if ctx.agent_event_emitter.is_none() {
return;
}
let exit_code = streaming.result.exit_code;
let termination = streaming.result.termination;
let duration_ms = streaming.result.duration_ms;
let streams_separated = streaming.streams_separated;
let result = streaming.result;
let exec_output_tail =
match task::spawn_blocking(move || result.default_redacted_output_tail()).await {
Ok(exec_output_tail) => exec_output_tail,
Err(err) => {
tracing::warn!(
error = ?err,
"Failed to redact shell process output tail"
);
None
}
};
ctx.emit_agent_event(AgentEvent::ToolProcessCompleted {
exit_code,
termination,
duration_ms,
streams_separated,
exec_output_tail,
});
}
/// Renders the model-facing shell result: termination, exit code, duration,
/// and provider-honest output sections. Metadata stays at the head and
/// `stderr` at the tail so head/tail truncation preserves both.
fn render_shell_result(streaming: &ExecStreamingResult) -> String {
let result = &streaming.result;
let mut output = format!(
"Termination: {}\nExit code: {}\nDuration: {}ms\n",
result.termination.as_str(),
result
.exit_code
.map_or_else(|| "none".to_string(), |code| code.to_string()),
result.duration_ms,
);
if streaming.streams_separated {
if !result.stdout.is_empty() {
let _ = write!(output, "stdout:\n{}\n", result.stdout);
}
if !result.stderr.is_empty() {
let _ = write!(output, "stderr:\n{}\n", result.stderr);
}
} else if !result.stdout.is_empty() {
let _ = write!(output, "output (combined):\n{}\n", result.stdout);
}
output
}
#[must_use]
@ -746,13 +811,18 @@ mod tests {
use fabro_llm::provider::ProviderAdapter;
use fabro_model::ProviderId;
use fabro_types::CommandTermination;
use tokio::sync::broadcast;
use tokio_util::sync::CancellationToken;
use super::*;
use crate::config::{NativeToolOptions, ToolSecrets};
use crate::config::{NativeToolOptions, SessionOptions, ToolSecrets};
use crate::event::{Emitter, SessionBoundEmitter};
use crate::local_sandbox::LocalSandbox;
use crate::sandbox::*;
use crate::test_support::MockSandbox;
use crate::tool_registry::ToolContext;
use crate::truncation;
use crate::types::SessionEvent;
#[test]
fn core_tool_descriptions_include_actionable_guidance() {
@ -1093,20 +1163,8 @@ mod tests {
assert_eq!(written[0].1, "1 | keep this literal\ngoodbye");
}
#[tokio::test]
async fn shell_basic_command() {
let tool = make_shell_tool();
let env: Arc<dyn Sandbox> = Arc::new(MockSandbox {
exec_result: ExecResult {
stdout: "hello".into(),
stderr: String::new(),
exit_code: Some(0),
termination: CommandTermination::Exited,
duration_ms: 10,
},
..Default::default()
});
let result = (tool.executor)(serde_json::json!({"command": "echo hello"}), ToolContext {
fn shell_context(env: Arc<dyn Sandbox>) -> ToolContext {
ToolContext {
env,
cancel: CancellationToken::new(),
tool_env_provider: None,
@ -1114,11 +1172,87 @@ mod tests {
root_session_id: None,
tool_call_id: None,
agent_event_emitter: None,
}
}
fn shell_context_with_emitter(env: Arc<dyn Sandbox>, emitter: &Emitter) -> ToolContext {
ToolContext {
session_id: Some("test-session".to_string()),
root_session_id: Some("test-session".to_string()),
tool_call_id: Some("call_1".to_string()),
agent_event_emitter: Some(Arc::new(SessionBoundEmitter {
emitter: emitter.clone(),
session_id: "test-session".to_string(),
tool_call_id: Some("call_1".to_string()),
})),
..shell_context(env)
}
}
fn only_process_event(receiver: &mut broadcast::Receiver<SessionEvent>) -> AgentEvent {
let event = receiver.try_recv().expect("one process event");
assert_eq!(event.session_id, "test-session");
assert_eq!(event.tool_call_id.as_deref(), Some("call_1"));
assert!(matches!(
receiver.try_recv(),
Err(broadcast::error::TryRecvError::Empty)
));
event.event
}
fn mock_sandbox_with(result: ExecResult) -> Arc<MockSandbox> {
Arc::new(MockSandbox {
exec_result: result,
..Default::default()
})
}
#[tokio::test]
async fn shell_success_returns_ok_with_metadata_and_separate_streams() {
let tool = make_shell_tool();
let env: Arc<dyn Sandbox> = mock_sandbox_with(ExecResult {
stdout: "hello".into(),
stderr: "a warning".into(),
exit_code: Some(0),
termination: CommandTermination::Exited,
duration_ms: 10,
});
let output = (tool.executor)(
serde_json::json!({"command": "echo hello"}),
shell_context(env),
)
.await
.expect("exit 0 is a successful tool result");
assert_eq!(
output,
"Termination: exited\nExit code: 0\nDuration: 10ms\nstdout:\nhello\nstderr:\na \
warning\n"
);
}
#[tokio::test]
async fn shell_forwards_command_without_stream_redirection_wrapper() {
let tool = make_shell_tool();
let env = mock_sandbox_with(ExecResult {
stdout: String::new(),
stderr: String::new(),
exit_code: Some(0),
termination: CommandTermination::Exited,
duration_ms: 1,
});
let _ = (tool.executor)(
serde_json::json!({"command": "make test"}),
shell_context(env.clone()),
)
.await;
let output = result.unwrap();
assert!(output.contains("Exit code: 0"));
assert!(output.contains("hello"));
let captured = env
.captured_command
.lock()
.expect("captured_command lock poisoned")
.clone();
assert_eq!(captured.as_deref(), Some("make test"));
}
#[tokio::test]
@ -1155,46 +1289,243 @@ mod tests {
},
..Default::default()
});
let result = (tool.executor)(serde_json::json!({"command": "false"}), ToolContext {
env,
cancel: CancellationToken::new(),
tool_env_provider: None,
session_id: None,
root_session_id: None,
tool_call_id: None,
agent_event_emitter: None,
})
.await;
let output = result.unwrap();
assert!(output.contains("Exit code: 1"));
assert!(output.contains("error"));
let output = (tool.executor)(serde_json::json!({"command": "false"}), shell_context(env))
.await
.expect_err("a nonzero exit is a failed tool result");
assert!(output.contains("Termination: exited"), "got: {output}");
assert!(output.contains("Exit code: 1"), "got: {output}");
assert!(output.contains("stdout:\nerror"), "got: {output}");
assert!(!output.contains("stderr:"), "got: {output}");
}
#[tokio::test]
async fn shell_timeout_output() {
async fn shell_timeout_returns_error_with_partial_output() {
let tool = make_shell_tool();
let env: Arc<dyn Sandbox> = mock_sandbox_with(ExecResult {
stdout: "partial".into(),
stderr: String::new(),
exit_code: None,
termination: CommandTermination::TimedOut,
duration_ms: 10000,
});
let output = (tool.executor)(
serde_json::json!({"command": "sleep 100"}),
shell_context(env),
)
.await
.expect_err("a timeout is a failed tool result");
assert!(output.contains("Termination: timed_out"), "got: {output}");
assert!(output.contains("Exit code: none"), "got: {output}");
assert!(output.contains("stdout:\npartial"), "got: {output}");
}
#[tokio::test]
async fn shell_cancellation_returns_error_with_partial_output() {
let tool = make_shell_tool();
let env: Arc<dyn Sandbox> = mock_sandbox_with(ExecResult {
stdout: "partial".into(),
stderr: String::new(),
exit_code: None,
termination: CommandTermination::Cancelled,
duration_ms: 42,
});
let output = (tool.executor)(
serde_json::json!({"command": "sleep 100"}),
shell_context(env),
)
.await
.expect_err("a cancellation is a failed tool result");
assert!(output.contains("Termination: cancelled"), "got: {output}");
assert!(output.contains("Exit code: none"), "got: {output}");
assert!(output.contains("stdout:\npartial"), "got: {output}");
}
#[tokio::test]
async fn shell_sandbox_failure_returns_error_without_a_process_outcome() {
let tool = make_shell_tool();
let env: Arc<dyn Sandbox> = Arc::new(MockSandbox {
exec_error: Some("sandbox transport is down".into()),
..Default::default()
});
let emitter = Emitter::new();
let mut receiver = emitter.subscribe();
let output = (tool.executor)(
serde_json::json!({"command": "make test"}),
shell_context_with_emitter(env, &emitter),
)
.await
.expect_err("a sandbox transport failure is a failed tool result");
assert!(
output.contains("Shell command produced no process result"),
"got: {output}"
);
assert!(
output.contains("sandbox transport is down"),
"got: {output}"
);
assert!(!output.contains("Exit code"), "got: {output}");
assert!(matches!(
receiver.try_recv(),
Err(broadcast::error::TryRecvError::Empty)
));
}
#[tokio::test]
async fn shell_emits_process_event_with_typed_outcome_and_redacted_tails() {
let tool = make_shell_tool();
let env: Arc<dyn Sandbox> = mock_sandbox_with(ExecResult {
stdout: "out".into(),
stderr: "boom key=AKIAYRWQG5EJLPZLBYNP".into(),
exit_code: Some(7),
termination: CommandTermination::Exited,
duration_ms: 12,
});
let emitter = Emitter::new();
let mut receiver = emitter.subscribe();
let _ = (tool.executor)(
serde_json::json!({"command": "printf out; printf err >&2; exit 7"}),
shell_context_with_emitter(env, &emitter),
)
.await;
match only_process_event(&mut receiver) {
AgentEvent::ToolProcessCompleted {
exit_code,
termination,
duration_ms,
streams_separated,
exec_output_tail,
} => {
assert_eq!(exit_code, Some(7));
assert_eq!(termination, CommandTermination::Exited);
assert_eq!(duration_ms, 12);
assert!(streams_separated);
let tail = exec_output_tail.expect("output tail");
assert_eq!(tail.stdout.as_deref(), Some("out"));
let stderr = tail.stderr.expect("stderr tail");
assert!(stderr.contains("boom"), "got: {stderr}");
assert!(!stderr.contains("AKIAYRWQG5EJLPZLBYNP"), "got: {stderr}");
}
other => panic!("expected a process event, got {other:?}"),
}
}
#[tokio::test]
async fn shell_renders_combined_output_when_streams_are_not_separated() {
let tool = make_shell_tool();
let env: Arc<dyn Sandbox> = Arc::new(MockSandbox {
exec_result: ExecResult {
stdout: String::new(),
stdout: "interleaved".into(),
stderr: String::new(),
exit_code: None,
termination: CommandTermination::TimedOut,
duration_ms: 10000,
exit_code: Some(0),
termination: CommandTermination::Exited,
duration_ms: 5,
},
streams_separated: false,
..Default::default()
});
let result = (tool.executor)(serde_json::json!({"command": "sleep 100"}), ToolContext {
env,
cancel: CancellationToken::new(),
tool_env_provider: None,
session_id: None,
root_session_id: None,
tool_call_id: None,
agent_event_emitter: None,
})
.await;
let output = result.unwrap();
assert!(output.starts_with("Command timed out.\n"));
let emitter = Emitter::new();
let mut receiver = emitter.subscribe();
let output = (tool.executor)(
serde_json::json!({"command": "echo interleaved"}),
shell_context_with_emitter(env, &emitter),
)
.await
.expect("exit 0 is a successful tool result");
assert!(
output.contains("output (combined):\ninterleaved"),
"got: {output}"
);
assert!(!output.contains("stderr:"), "got: {output}");
match only_process_event(&mut receiver) {
AgentEvent::ToolProcessCompleted {
streams_separated, ..
} => assert!(!streams_separated),
other => panic!("expected a process event, got {other:?}"),
}
}
#[tokio::test]
async fn shell_truncation_preserves_exit_metadata_and_stderr_tail() {
let tool = make_shell_tool();
let stdout = (0..400)
.map(|line| format!("{line}: {}", "x".repeat(100)))
.collect::<Vec<_>>()
.join("\n");
assert!(stdout.len() > 30_000);
let env: Arc<dyn Sandbox> = mock_sandbox_with(ExecResult {
stdout,
stderr: "the build failed".into(),
exit_code: Some(2),
termination: CommandTermination::Exited,
duration_ms: 900,
});
let output = (tool.executor)(
serde_json::json!({"command": "make build"}),
shell_context(env),
)
.await
.expect_err("a nonzero exit is a failed tool result");
let truncated =
truncation::truncate_tool_output(&output, "shell", &SessionOptions::default());
assert!(truncated.len() < output.len());
assert!(truncated.starts_with("Termination: exited\nExit code: 2\n"));
assert!(
truncated.contains("stderr:\nthe build failed"),
"stderr tail did not survive truncation"
);
}
/// End-to-end against a real process: the local provider separates the
/// streams and reports the real exit code, and none of it is laundered
/// into a successful tool result.
#[tokio::test]
async fn shell_reports_real_local_process_outcome() {
let tool = make_shell_tool();
let env: Arc<dyn Sandbox> = Arc::new(LocalSandbox::new(
std::env::current_dir().expect("current dir"),
));
let emitter = Emitter::new();
let mut receiver = emitter.subscribe();
let output = (tool.executor)(
serde_json::json!({"command": "printf 'out'; printf 'err' >&2; exit 7"}),
shell_context_with_emitter(env, &emitter),
)
.await
.expect_err("exit 7 is a failed tool result");
assert!(output.contains("Termination: exited"), "got: {output}");
assert!(output.contains("Exit code: 7"), "got: {output}");
assert!(output.contains("stdout:\nout"), "got: {output}");
assert!(output.contains("stderr:\nerr"), "got: {output}");
match only_process_event(&mut receiver) {
AgentEvent::ToolProcessCompleted {
exit_code,
termination,
streams_separated,
exec_output_tail,
..
} => {
assert_eq!(exit_code, Some(7));
assert_eq!(termination, CommandTermination::Exited);
assert!(streams_separated);
let tail = exec_output_tail.expect("output tail");
assert_eq!(tail.stdout.as_deref(), Some("out"));
assert_eq!(tail.stderr.as_deref(), Some("err"));
}
other => panic!("expected a process event, got {other:?}"),
}
}
#[tokio::test]

View file

@ -4,7 +4,10 @@ use chrono::{DateTime, Utc};
use fabro_llm::Error as LlmError;
use fabro_llm::types::{ContentPart, ThinkingData, TokenCounts, ToolCall, ToolResult};
use fabro_model::{CostSource, ModelRef};
use fabro_types::{ReasoningOutput, SessionMessage, StageContextWindowProjection};
use fabro_types::{
CommandTermination, ExecOutputTail, ReasoningOutput, SessionMessage,
StageContextWindowProjection,
};
use serde::de::DeserializeOwned;
use serde::{Deserialize, Serialize};
@ -280,6 +283,19 @@ pub enum AgentEvent {
output: serde_json::Value,
is_error: bool,
},
/// Subordinate process outcome for a tool call that ran a command.
/// Emitted before the owning `ToolCallCompleted`, which stays the single
/// tool-protocol completion and the authoritative owner of `is_error`.
/// Session and tool-call identity come from the emitting envelope.
ToolProcessCompleted {
#[serde(default, skip_serializing_if = "Option::is_none")]
exit_code: Option<i32>,
termination: CommandTermination,
duration_ms: u64,
streams_separated: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
exec_output_tail: Option<ExecOutputTail>,
},
Error {
error: Error,
},
@ -457,6 +473,28 @@ impl AgentEvent {
"Tool call completed"
);
}
Self::ToolProcessCompleted {
exit_code,
termination,
duration_ms,
streams_separated,
exec_output_tail,
} => {
let tail = ExecOutputTail::trace_summary(exec_output_tail.as_ref());
debug!(
session_id,
exit_code = ?exit_code,
termination = termination.as_str(),
duration_ms,
streams_separated,
output_tail_present = tail.present,
stdout_bytes = tail.stdout_bytes,
stderr_bytes = tail.stderr_bytes,
stdout_truncated = tail.stdout_truncated,
stderr_truncated = tail.stderr_truncated,
"Tool process completed"
);
}
Self::Error { error } => {
error!(session_id, error = %error, "Agent error");
}

View file

@ -0,0 +1,95 @@
//! Proves the agent shell tool reports real process outcomes through the
//! Docker provider's streaming path, which uses a `bash -lc` supervisor and
//! separate stdout/stderr channels.
use std::sync::Arc;
use fabro_agent::event::SessionBoundEmitter;
use fabro_agent::sandbox::Sandbox;
use fabro_agent::tool_registry::ToolContext;
use fabro_agent::tools::make_shell_tool;
use fabro_agent::types::AgentEvent;
use fabro_agent::{DockerSandbox, DockerSandboxOptions, Emitter};
use fabro_types::CommandTermination;
use tokio::sync::broadcast;
use tokio_util::sync::CancellationToken;
#[tokio::test]
#[ignore = "requires real Docker container lifecycle; run explicitly when changing shell tool exec integration"]
async fn shell_reports_real_docker_process_outcome() {
let Ok(sandbox) = DockerSandbox::new(
DockerSandboxOptions {
image: "buildpack-deps:noble".to_string(),
auto_pull: false,
skip_clone: true,
..DockerSandboxOptions::default()
},
None,
None,
None,
None,
) else {
return;
};
// No Docker daemon or no local image: the integration precondition is not met.
if sandbox.initialize().await.is_err() {
return;
}
let sandbox = Arc::new(sandbox);
let emitter = Emitter::new();
let mut receiver = emitter.subscribe();
let tool = make_shell_tool();
let result = (tool.executor)(
serde_json::json!({"command": "printf 'out'; printf 'err' >&2; exit 7"}),
ToolContext {
env: sandbox.clone() as Arc<dyn Sandbox>,
cancel: CancellationToken::new(),
tool_env_provider: None,
session_id: Some("test-session".to_string()),
root_session_id: Some("test-session".to_string()),
tool_call_id: Some("call_1".to_string()),
agent_event_emitter: Some(Arc::new(SessionBoundEmitter {
emitter: emitter.clone(),
session_id: "test-session".to_string(),
tool_call_id: Some("call_1".to_string()),
})),
},
)
.await;
sandbox
.cleanup()
.await
.expect("docker cleanup should succeed");
let output = result.expect_err("exit 7 is a failed tool result");
assert!(output.contains("Termination: exited"), "got: {output}");
assert!(output.contains("Exit code: 7"), "got: {output}");
assert!(output.contains("stdout:\nout"), "got: {output}");
assert!(output.contains("stderr:\nerr"), "got: {output}");
let event = receiver.try_recv().expect("one process event");
assert_eq!(event.session_id, "test-session");
assert_eq!(event.tool_call_id.as_deref(), Some("call_1"));
assert!(matches!(
receiver.try_recv(),
Err(broadcast::error::TryRecvError::Empty)
));
match event.event {
AgentEvent::ToolProcessCompleted {
exit_code,
termination,
streams_separated,
exec_output_tail,
..
} => {
assert_eq!(exit_code, Some(7));
assert_eq!(termination, CommandTermination::Exited);
assert!(streams_separated);
let tail = exec_output_tail.expect("output tail");
assert_eq!(tail.stdout.as_deref(), Some("out"));
assert_eq!(tail.stderr.as_deref(), Some("err"));
}
other => panic!("expected a process event, got {other:?}"),
}
}

View file

@ -1,3 +1,5 @@
mod compaction;
#[cfg(feature = "docker")]
mod docker_shell;
mod guardrails;
mod parity_matrix;

View file

@ -1707,7 +1707,7 @@ impl Sandbox for DaytonaSandbox {
working_dir: Option<&str>,
env_vars: Option<&HashMap<String, String>>,
cancel_token: Option<CancellationToken>,
output_callback: CommandOutputCallback,
output_callback: Option<CommandOutputCallback>,
) -> crate::Result<ExecStreamingResult> {
let sandbox = self.sandbox()?;
let start = Instant::now();
@ -1765,9 +1765,11 @@ impl Sandbox for DaytonaSandbox {
if !bytes.is_empty() {
saw_live_chunk.store(true, Ordering::Relaxed);
stdout_seen.lock().await.extend_from_slice(&bytes);
callback(CommandOutputStream::Stdout, bytes)
.await
.map_err(|err| daytona_callback_error(&err))?;
if let Some(callback) = callback {
callback(CommandOutputStream::Stdout, bytes)
.await
.map_err(|err| daytona_callback_error(&err))?;
}
}
Ok(())
}
@ -1781,9 +1783,11 @@ impl Sandbox for DaytonaSandbox {
if !bytes.is_empty() {
saw_live_chunk.store(true, Ordering::Relaxed);
stderr_seen.lock().await.extend_from_slice(&bytes);
callback(CommandOutputStream::Stderr, bytes)
.await
.map_err(|err| daytona_callback_error(&err))?;
if let Some(callback) = callback {
callback(CommandOutputStream::Stderr, bytes)
.await
.map_err(|err| daytona_callback_error(&err))?;
}
}
Ok(())
}
@ -1840,14 +1844,14 @@ impl Sandbox for DaytonaSandbox {
CommandOutputStream::Stdout,
logs.stdout.as_bytes(),
&stdout_seen,
&output_callback,
output_callback.as_ref(),
)
.await?;
append_missing_log_suffix(
CommandOutputStream::Stderr,
logs.stderr.as_bytes(),
&stderr_seen,
&output_callback,
output_callback.as_ref(),
)
.await?;
}
@ -2252,7 +2256,7 @@ async fn append_missing_log_suffix(
stream: CommandOutputStream,
final_bytes: &[u8],
seen: &Arc<Mutex<Vec<u8>>>,
output_callback: &CommandOutputCallback,
output_callback: Option<&CommandOutputCallback>,
) -> crate::Result<()> {
if final_bytes.is_empty() {
return Ok(());
@ -2267,7 +2271,10 @@ async fn append_missing_log_suffix(
let missing = final_bytes[offset..].to_vec();
seen.extend_from_slice(&missing);
drop(seen);
output_callback(stream, missing).await
match output_callback {
Some(output_callback) => output_callback(stream, missing).await,
None => Ok(()),
}
}
fn missing_log_suffix_offset(seen: &[u8], final_bytes: &[u8]) -> usize {

View file

@ -370,7 +370,7 @@ impl DockerSandbox {
cmd: Vec<String>,
working_dir: Option<String>,
env: Option<Vec<String>>,
output_callback: CommandOutputCallback,
output_callback: Option<CommandOutputCallback>,
) -> crate::Result<(Vec<u8>, Vec<u8>, i32)> {
let exec_opts = CreateExecOptions {
cmd: Some(cmd),
@ -399,11 +399,15 @@ impl DockerSandbox {
match chunk {
Ok(LogOutput::StdOut { message }) => {
stdout.extend_from_slice(&message);
output_callback(CommandOutputStream::Stdout, message.to_vec()).await?;
if let Some(output_callback) = output_callback.as_ref() {
output_callback(CommandOutputStream::Stdout, message.to_vec()).await?;
}
}
Ok(LogOutput::StdErr { message }) => {
stderr.extend_from_slice(&message);
output_callback(CommandOutputStream::Stderr, message.to_vec()).await?;
if let Some(output_callback) = output_callback.as_ref() {
output_callback(CommandOutputStream::Stderr, message.to_vec()).await?;
}
}
Ok(_) => {}
Err(e) => {
@ -489,7 +493,7 @@ impl DockerSandbox {
working_dir: Option<&str>,
env_vars: Option<&HashMap<String, String>>,
cancel_token: Option<CancellationToken>,
output_callback: CommandOutputCallback,
output_callback: Option<CommandOutputCallback>,
) -> crate::Result<ExecStreamingResult> {
let start = Instant::now();
let effective_dir = working_dir
@ -1579,7 +1583,7 @@ impl Sandbox for DockerSandbox {
working_dir: Option<&str>,
env_vars: Option<&HashMap<String, String>>,
cancel_token: Option<CancellationToken>,
output_callback: CommandOutputCallback,
output_callback: Option<CommandOutputCallback>,
) -> crate::Result<ExecStreamingResult> {
let dir = working_dir.map(|path| self.resolve_container_path(path));
self.docker_exec_shell_streaming(

View file

@ -508,7 +508,7 @@ impl Sandbox for LocalSandbox {
working_dir: Option<&str>,
env_vars: Option<&std::collections::HashMap<String, String>>,
cancel_token: Option<CancellationToken>,
output_callback: CommandOutputCallback,
output_callback: Option<CommandOutputCallback>,
) -> crate::Result<ExecStreamingResult> {
let start = Instant::now();
@ -950,7 +950,7 @@ async fn sigterm_then_kill(child: &mut Child) {
async fn drain_command_pipe<R>(
mut reader: Option<R>,
stream: CommandOutputStream,
output_callback: CommandOutputCallback,
output_callback: Option<CommandOutputCallback>,
) -> crate::Result<Vec<u8>>
where
R: AsyncRead + Unpin,
@ -970,7 +970,9 @@ where
return Ok(output);
}
output.extend_from_slice(&buf[..read]);
output_callback(stream, buf[..read].to_vec()).await?;
if let Some(output_callback) = output_callback.as_ref() {
output_callback(stream, buf[..read].to_vec()).await?;
}
}
}
@ -1256,7 +1258,7 @@ mod tests {
None,
None,
None,
Arc::new(|_, _| Box::pin(async { Ok(()) })),
Some(Arc::new(|_, _| Box::pin(async { Ok(()) }))),
)
.await
.unwrap();
@ -1327,7 +1329,7 @@ mod tests {
None,
Some(&env_vars),
None,
Arc::new(|_, _| Box::pin(async { Ok(()) })),
Some(Arc::new(|_, _| Box::pin(async { Ok(()) }))),
)
.await
.unwrap();

View file

@ -184,7 +184,7 @@ macro_rules! delegate_sandbox {
working_dir: Option<&str>,
env_vars: Option<&std::collections::HashMap<String, String>>,
cancel_token: Option<tokio_util::sync::CancellationToken>,
output_callback: $crate::CommandOutputCallback,
output_callback: Option<$crate::CommandOutputCallback>,
) -> $crate::Result<$crate::ExecStreamingResult> {
self.$field
.exec_command_streaming(
@ -753,6 +753,34 @@ pub type CommandOutputCallback = Arc<
+ Sync,
>;
pub(crate) async fn replay_exec_result(
result: ExecResult,
streams_separated: bool,
output_callback: Option<&CommandOutputCallback>,
) -> crate::Result<ExecStreamingResult> {
if let Some(output_callback) = output_callback {
if !result.stdout.is_empty() {
output_callback(
CommandOutputStream::Stdout,
result.stdout.as_bytes().to_vec(),
)
.await?;
}
if !result.stderr.is_empty() {
output_callback(
CommandOutputStream::Stderr,
result.stderr.as_bytes().to_vec(),
)
.await?;
}
}
Ok(ExecStreamingResult {
result,
streams_separated,
live_streaming: false,
})
}
pub struct StdioProcess {
pub stdin: Pin<Box<dyn AsyncWrite + Send>>,
pub stdout: Pin<Box<dyn AsyncRead + Send>>,
@ -949,11 +977,13 @@ pub trait Sandbox: Send + Sync {
///
/// **Production sandboxes must override this.** The default falls back to
/// the non-streaming [`exec_command`](Self::exec_command) and replays its
/// output through `output_callback` at the end, marking
/// `live_streaming: false`. That's the right behavior for test mocks but
/// silently drops live output for any real sandbox that wraps another —
/// decorators in particular must forward to the inner sandbox's streaming
/// implementation rather than relying on this default.
/// output through `output_callback` at the end when one is supplied,
/// marking `live_streaming: false`. Passing `None` captures the final
/// result without paying per-chunk callback costs. That's the right
/// behavior for test mocks but silently drops live output for any real
/// sandbox that wraps another — decorators in particular must forward to
/// the inner sandbox's streaming implementation rather than relying on
/// this default.
async fn exec_command_streaming(
&self,
command: &str,
@ -961,7 +991,7 @@ pub trait Sandbox: Send + Sync {
working_dir: Option<&str>,
env_vars: Option<&std::collections::HashMap<String, String>>,
cancel_token: Option<CancellationToken>,
output_callback: CommandOutputCallback,
output_callback: Option<CommandOutputCallback>,
) -> crate::Result<ExecStreamingResult> {
let fallback_timeout_ms = timeout_ms.unwrap_or(u64::MAX);
let result = self
@ -973,25 +1003,7 @@ pub trait Sandbox: Send + Sync {
cancel_token,
)
.await?;
if !result.stdout.is_empty() {
output_callback(
CommandOutputStream::Stdout,
result.stdout.as_bytes().to_vec(),
)
.await?;
}
if !result.stderr.is_empty() {
output_callback(
CommandOutputStream::Stderr,
result.stderr.as_bytes().to_vec(),
)
.await?;
}
Ok(ExecStreamingResult {
result,
streams_separated: true,
live_streaming: false,
})
replay_exec_result(result, true, output_callback.as_ref()).await
}
/// Launch a long-lived process with bidirectional stdio attached.

View file

@ -9,7 +9,7 @@ use tokio::io::{DuplexStream, duplex};
use tokio::time::sleep;
use tokio_util::sync::CancellationToken;
use crate::sandbox::StdioProcessControl;
use crate::sandbox::{self, StdioProcessControl};
use crate::{
DEFAULT_EXEC_OUTPUT_TAIL_BYTES, DirEntry, ExecResult, GrepOptions, Sandbox, SandboxEvent,
SandboxEventCallback, StderrCollector, StdioProcess, StdioProcessHandle,
@ -44,6 +44,12 @@ pub struct MockSandbox {
pub event_callback: Option<SandboxEventCallback>,
pub stdio_process_error: Option<String>,
pub stdio_process: Mutex<Option<MockStdioProcess>>,
/// Fails `exec_command` and `exec_command_streaming` before any process
/// runs, so callers see a transport error rather than an `ExecResult`.
pub exec_error: Option<String>,
/// Reported by `exec_command_streaming`. Set to `false` to model a
/// provider that cannot separate stdout from stderr.
pub streams_separated: bool,
}
impl MockSandbox {
@ -116,6 +122,8 @@ impl Default for MockSandbox {
event_callback: None,
stdio_process_error: None,
stdio_process: Mutex::new(None),
exec_error: None,
streams_separated: true,
}
}
}
@ -233,7 +241,31 @@ impl Sandbox for MockSandbox {
.captured_env_vars
.lock()
.expect("captured_env_vars lock poisoned") = env_vars.cloned();
Ok(self.exec_result.clone())
match &self.exec_error {
Some(error) => Err(crate::Error::message(error.clone())),
None => Ok(self.exec_result.clone()),
}
}
async fn exec_command_streaming(
&self,
command: &str,
timeout_ms: Option<u64>,
working_dir: Option<&str>,
env_vars: Option<&std::collections::HashMap<String, String>>,
cancel_token: Option<CancellationToken>,
output_callback: Option<crate::CommandOutputCallback>,
) -> crate::Result<crate::ExecStreamingResult> {
let result = self
.exec_command(
command,
timeout_ms.unwrap_or(u64::MAX),
working_dir,
env_vars,
cancel_token,
)
.await?;
sandbox::replay_exec_result(result, self.streams_separated, output_callback.as_ref()).await
}
async fn spawn_stdio_process(

View file

@ -335,7 +335,7 @@ mod daytona_streaming_live {
None,
None,
Some(cancel_for_exec),
callback,
Some(callback),
)
.await
});
@ -460,7 +460,7 @@ mod daytona_streaming_live {
None,
None,
cancel_token,
callback,
Some(callback),
)
.await?;
let chunks = chunks.lock().await.clone();

View file

@ -55,7 +55,7 @@ async fn streaming_timeout_terminates_docker_exec_before_returning() {
None,
None,
None,
capture_bytes(Arc::clone(&chunks)),
Some(capture_bytes(Arc::clone(&chunks))),
)
.await
.expect("streaming command should return a timeout result");
@ -218,7 +218,7 @@ async fn docker_runs_clean_bash_through_both_command_paths() {
None,
None,
None,
capture_bytes(Arc::clone(&chunks)),
Some(capture_bytes(Arc::clone(&chunks))),
)
.await
.expect("streaming command should run");

View file

@ -655,6 +655,22 @@ fn event_body_from_event(event: &Event) -> EventBody {
tool_result: None,
turn_id: None,
}),
AgentEvent::ToolProcessCompleted {
exit_code,
termination,
duration_ms,
streams_separated,
exec_output_tail,
} => EventBody::AgentToolProcessCompleted(
fabro_types::AgentToolProcessCompletedProps {
exit_code: *exit_code,
termination: *termination,
duration_ms: *duration_ms,
streams_separated: *streams_separated,
exec_output_tail: exec_output_tail.clone(),
visit: *visit,
},
),
AgentEvent::Error { error } => EventBody::AgentError(fabro_types::AgentErrorProps {
error: serde_json::to_value(error).expect("agent Error derives Serialize with no custom logic that can fail"),
visit: *visit,
@ -1517,6 +1533,47 @@ mod tests {
assert_eq!(properties["visit"], 2);
}
#[test]
fn run_event_agent_tool_process_completed_carries_stage_session_and_actor() {
let stored = to_run_event(&fixtures::RUN_4, &Event::Agent {
stage: "code".to_string(),
visit: 2,
event: AgentEvent::ToolProcessCompleted {
exit_code: Some(7),
termination: ::fabro_types::CommandTermination::Exited,
duration_ms: 12,
streams_separated: true,
exec_output_tail: Some(exec_tail()),
},
session_id: Some("ses_child".to_string()),
parent_session_id: Some("ses_parent".to_string()),
tool_call_id: Some("call_1".to_string()),
});
assert_eq!(stored.event_name(), "agent.tool.process.completed");
assert_eq!(stored.node_id.as_deref(), Some("code"));
assert_eq!(stored.stage_id, Some(StageId::new("code", 2)));
assert_eq!(stored.session_id.as_deref(), Some("ses_child"));
assert_eq!(stored.parent_session_id.as_deref(), Some("ses_parent"));
assert_eq!(stored.tool_call_id.as_deref(), Some("call_1"));
assert_eq!(
stored.actor,
Some(::fabro_types::Principal::Agent {
session_id: Some("ses_child".to_string()),
parent_session_id: Some("ses_parent".to_string()),
model: None,
})
);
let properties = stored.properties().unwrap();
assert_eq!(properties["exit_code"], 7);
assert_eq!(properties["termination"], "exited");
assert_eq!(properties["duration_ms"], 12);
assert_eq!(properties["streams_separated"], true);
assert_eq!(properties["exec_output_tail"]["stdout"], "last stdout line");
assert_eq!(properties["visit"], 2);
}
#[test]
fn run_event_agent_tools_available_moves_session_and_stage_metadata_to_header() {
let stored = to_run_event(&fixtures::RUN_4, &Event::AgentToolsAvailable {

View file

@ -75,6 +75,7 @@ pub fn event_name(event: &Event) -> &'static str {
AgentEvent::ToolCallStarted { .. } => "agent.tool.started",
AgentEvent::ToolCallOutputDelta { .. } => "agent.tool.output.delta",
AgentEvent::ToolCallCompleted { .. } => "agent.tool.completed",
AgentEvent::ToolProcessCompleted { .. } => "agent.tool.process.completed",
AgentEvent::Error { .. } => "agent.error",
AgentEvent::Warning { .. } => "agent.warning",
AgentEvent::LoopDetected => "agent.loop.detected",

View file

@ -76,6 +76,41 @@ mod tests {
);
}
#[test]
fn build_redacted_event_payload_redacts_tool_process_output_tails() {
let secret = "sk-ant-api03-xK9mZ2vL8nQ5rT1wY4bC7dF0gH3jE6pA";
let stored = to_run_event(&fixtures::RUN_8, &Event::Agent {
stage: "code".to_string(),
visit: 1,
event: AgentEvent::ToolProcessCompleted {
exit_code: Some(7),
termination: ::fabro_types::CommandTermination::Exited,
duration_ms: 12,
streams_separated: true,
exec_output_tail: Some(fabro_types::ExecOutputTail {
stdout: Some(format!("stdout {secret}")),
stderr: Some("plain stderr".to_string()),
stdout_truncated: false,
stderr_truncated: false,
}),
},
session_id: Some("ses_child".to_string()),
parent_session_id: None,
tool_call_id: Some("call_1".to_string()),
});
let payload = build_redacted_event_payload(&stored, &fixtures::RUN_8).unwrap();
let payload_text = serde_json::to_string(payload.as_value()).unwrap();
assert!(!payload_text.contains(secret));
assert!(payload_text.contains("REDACTED"));
assert_eq!(payload.as_value()["event"], "agent.tool.process.completed");
assert_eq!(
payload.as_value()["properties"]["exec_output_tail"]["stderr"],
"plain stderr"
);
}
/// Reasoning is model-authored text like any other, so it goes through
/// the same canonical redaction pass as assistant output.
#[test]

View file

@ -328,7 +328,8 @@ fn agent_actor_for_event(
}),
AgentEvent::ToolCallStarted { .. }
| AgentEvent::ToolCallOutputDelta { .. }
| AgentEvent::ToolCallCompleted { .. } => Some(Principal::Agent {
| AgentEvent::ToolCallCompleted { .. }
| AgentEvent::ToolProcessCompleted { .. } => Some(Principal::Agent {
session_id: session_id.map(str::to_string),
parent_session_id: parent_session_id.map(str::to_string),
model: None,

View file

@ -128,7 +128,7 @@ impl Handler for CommandHandler {
None,
env_vars,
Some(cancel_token.clone()),
output_callback,
Some(output_callback),
)
.await;
cancel_token.cancel();

View file

@ -3083,6 +3083,19 @@ enabled = true
false,
BillingPolicy::OpenAi,
),
(
"claude-opus-5",
"anthropic/claude-opus-5",
"claude-5",
1_000_000,
5.0,
25.0,
0.5,
ReasoningEffortFeature::Levels,
false,
true,
BillingPolicy::Anthropic,
),
(
"claude-opus-4-8",
"anthropic/claude-opus-4.8",
@ -3161,6 +3174,13 @@ enabled = true
"{id}"
);
}
for alias in ["opus", "claude-opus"] {
let model = catalog
.resolve_on_provider(&ProviderId::new("openrouter"), alias)
.unwrap_or_else(|error| panic!("{alias} should resolve on OpenRouter: {error}"));
assert_eq!(model.id, "claude-opus-5", "{alias}");
}
}
#[test]

View file

@ -58,6 +58,33 @@ input_cost_per_mtok = 10.0
output_cost_per_mtok = 50.0
cache_input_cost_per_mtok = 1.0
[providers.openrouter.models."claude-opus-5"]
api_id = "anthropic/claude-opus-5"
display_name = "Claude Opus 5 (via OpenRouter)"
family = "claude-5"
billing_policy = "anthropic"
training = "2026-05-01"
knowledge_cutoff = "May 2026"
aliases = ["opus", "claude-opus"]
[providers.openrouter.models."claude-opus-5".limits]
context_window = 1000000
max_output = 128000
[providers.openrouter.models."claude-opus-5".features]
tools = true
vision = true
reasoning = true
reasoning_effort = "levels"
prompt_cache = true
cache_control_breakpoints = true
sampling_params = false
[providers.openrouter.models."claude-opus-5".costs]
input_cost_per_mtok = 5.0
output_cost_per_mtok = 25.0
cache_input_cost_per_mtok = 0.5
[providers.openrouter.models."claude-opus-4-8"]
api_id = "anthropic/claude-opus-4.8"
display_name = "Claude Opus 4.8 (via OpenRouter)"
@ -65,7 +92,6 @@ family = "claude-4"
billing_policy = "anthropic"
training = "2026-01-01"
knowledge_cutoff = "Jan 2026"
aliases = ["opus", "claude-opus"]
[providers.openrouter.models."claude-opus-4-8".limits]
context_window = 1000000

View file

@ -3,11 +3,11 @@ use serde::{Deserialize, Serialize};
use serde_json::Value;
use strum::{Display, EnumString, IntoStaticStr};
use super::BilledTokenCounts;
use super::{BilledTokenCounts, ExecOutputTail};
use crate::transcript::{ToolCall, ToolResult, TranscriptMessage};
use crate::{
MessageId, ModelRef, PairId, PairMessageId, PairSystemMessageKind, PermissionLevel,
ReasoningOutput, StageContextWindowProjection, TurnId,
CommandTermination, MessageId, ModelRef, PairId, PairMessageId, PairSystemMessageKind,
PermissionLevel, ReasoningOutput, StageContextWindowProjection, TurnId,
};
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
@ -181,6 +181,25 @@ pub struct AgentToolCompletedProps {
pub turn_id: Option<TurnId>,
}
/// Subordinate diagnostic for a tool call that ran a process: the real
/// termination, exit code, duration, and bounded redacted output tails.
///
/// This never replaces `agent.tool.completed`, which remains the single
/// tool-protocol completion and the authoritative owner of `is_error`.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AgentToolProcessCompletedProps {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub exit_code: Option<i32>,
pub termination: CommandTermination,
pub duration_ms: u64,
/// `false` when the provider could not separate stdout from stderr. The
/// combined output is then carried in `exec_output_tail.stdout`.
pub streams_separated: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub exec_output_tail: Option<ExecOutputTail>,
pub visit: u32,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AgentErrorProps {
pub error: Value,

View file

@ -401,3 +401,30 @@ pub struct CliEnsureFailedProps {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub exec_output_tail: Option<ExecOutputTail>,
}
#[cfg(test)]
mod tests {
use super::ExecOutputTail;
/// The trace summary is expanded into tracing fields, so it must carry
/// sizes and truncation flags only.
#[test]
fn exec_output_tail_trace_summary_exposes_sizes_not_content() {
let tail = ExecOutputTail {
stdout: Some("secret stdout bytes".to_string()),
stderr: Some("secret stderr".to_string()),
stdout_truncated: true,
stderr_truncated: false,
};
let summary = ExecOutputTail::trace_summary(Some(&tail));
assert!(summary.present);
assert_eq!(summary.stdout_bytes, 19);
assert_eq!(summary.stderr_bytes, 13);
assert!(summary.stdout_truncated);
assert!(!summary.stderr_truncated);
let rendered = format!("{summary:?}");
assert!(!rendered.contains("secret"), "got: {rendered}");
}
}

View file

@ -208,6 +208,8 @@ pub enum EventBody {
AgentToolStarted(AgentToolStartedProps),
#[serde(rename = "agent.tool.completed")]
AgentToolCompleted(AgentToolCompletedProps),
#[serde(rename = "agent.tool.process.completed")]
AgentToolProcessCompleted(AgentToolProcessCompletedProps),
#[serde(rename = "agent.error")]
AgentError(AgentErrorProps),
#[serde(rename = "agent.warning")]
@ -487,6 +489,7 @@ impl EventBody {
Self::AgentMessage(_) => "agent.message",
Self::AgentToolStarted(_) => "agent.tool.started",
Self::AgentToolCompleted(_) => "agent.tool.completed",
Self::AgentToolProcessCompleted(_) => "agent.tool.process.completed",
Self::AgentError(_) => "agent.error",
Self::AgentWarning(_) => "agent.warning",
Self::AgentLoopDetected(_) => "agent.loop.detected",
@ -656,6 +659,7 @@ fn is_known_event_name(event: &str) -> bool {
| "agent.message"
| "agent.tool.started"
| "agent.tool.completed"
| "agent.tool.process.completed"
| "agent.error"
| "agent.warning"
| "agent.loop.detected"
@ -920,8 +924,8 @@ mod tests {
use super::*;
use crate::{
AuthMethod, Edge, Graph, IdpIdentity, Node, PendingReason, RunBlobId, WorkflowSettings,
fixtures, test_support,
AuthMethod, CommandTermination, Edge, Graph, IdpIdentity, Node, PendingReason, RunBlobId,
WorkflowSettings, fixtures, test_support,
};
fn user_principal(login: &str) -> Principal {
@ -2409,6 +2413,68 @@ mod tests {
assert_eq!(parsed, body);
}
#[test]
fn agent_tool_process_completed_round_trips_as_a_known_typed_event() {
let value = json!({
"id": "evt_process",
"ts": "2026-04-08T16:21:11.106Z",
"run_id": fixtures::RUN_1,
"event": "agent.tool.process.completed",
"node_id": "code",
"session_id": "ses_child",
"tool_call_id": "call_1",
"properties": {
"exit_code": 7,
"termination": "exited",
"duration_ms": 12,
"streams_separated": true,
"exec_output_tail": {"stdout": "out", "stderr": "err"},
"visit": 1
}
});
let parsed = RunEvent::from_value(value.clone()).unwrap();
assert_eq!(parsed.event_name(), "agent.tool.process.completed");
assert_eq!(parsed.tool_call_id.as_deref(), Some("call_1"));
let EventBody::AgentToolProcessCompleted(props) = &parsed.body else {
panic!("expected a typed process event, got {:?}", parsed.body);
};
assert_eq!(props.exit_code, Some(7));
assert_eq!(props.termination, CommandTermination::Exited);
assert_eq!(props.duration_ms, 12);
assert!(props.streams_separated);
assert_eq!(
props.exec_output_tail.as_ref().unwrap().stdout.as_deref(),
Some("out")
);
assert_eq!(parsed.to_value().unwrap(), value);
}
#[test]
fn agent_tool_process_completed_omits_absent_exit_code_and_output_tail() {
let body = EventBody::AgentToolProcessCompleted(AgentToolProcessCompletedProps {
exit_code: None,
termination: CommandTermination::TimedOut,
duration_ms: 10_000,
streams_separated: false,
exec_output_tail: None,
visit: 1,
});
let value = serde_json::to_value(&body).unwrap();
assert_eq!(value["event"], "agent.tool.process.completed");
assert_eq!(value["properties"]["termination"], "timed_out");
assert_eq!(value["properties"]["streams_separated"], false);
let properties = value["properties"].as_object().unwrap();
assert!(!properties.contains_key("exit_code"));
assert!(!properties.contains_key("exec_output_tail"));
let parsed: EventBody = serde_json::from_value(value).unwrap();
assert_eq!(parsed, body);
}
#[test]
fn agent_tool_source_and_category_use_public_json_shape() {
assert_eq!(