use std::fmt::Write as _;
use std::{collections::HashSet, future::Future, pin::Pin};
use async_trait::async_trait;
use daat_locus_macros::model_schema;
use miette::{Result, miette};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use serde_json::{Value, json};
use crate::{
activity_event::{TextActivityDescriptor, ToolCallActivityEvent, glyph},
app::{AppId, AppManager, AppStateRender, AppToolExecutionContext},
context::Context,
context_budget::truncate_text_to_token_budget_with_notice,
dashboard::SessionActivityEvent,
live_progress::TelegramLiveStatus,
reasoning::{
episode::EpisodeActionRecord,
runtime::{AgentContentPart, AgentToolCall, AgentToolInputSpec, AgentToolSpec},
},
schema_utils::{model_schema, model_schema_for, validate_model_facing_schema},
workflow::{WorkflowInvocation, invoke as invoke_workflow},
};
mod files;
mod view_image;
mod work;
pub type ToolFuture<'a> =
Pin<Box<dyn Future<Output = miette::Result<ToolExecutionResult>> + Send + 'a>>;
type ToolExecutor = for<'a> fn(&'a mut Context, &'a AgentToolCall) -> ToolFuture<'a>;
type ToolSummarizer = fn(&AgentToolCall) -> miette::Result<EpisodeActionRecord>;
type ToolCallActivityBuilder = fn(&AgentToolCall) -> miette::Result<ToolCallActivityEvent>;
type ToolAvailability = fn(&Context) -> bool;
pub fn parse_tool_args<T: for<'de> serde::Deserialize<'de>>(
call: &AgentToolCall,
) -> miette::Result<T> {
serde_json::from_value(call.arguments.clone()).map_err(|err| {
miette!(
"invalid arguments for tool `{}`: {}; arguments={}",
call.name,
err,
call.arguments
)
})
}
pub fn summarize_inline_text(text: &str) -> String {
const MAX_CHARS: usize = 120;
let compact = text.replace('\n', "\\n");
let mut chars = compact.chars();
let summary = chars.by_ref().take(MAX_CHARS).collect::<String>();
if chars.next().is_some() {
format!("{summary}...")
} else {
summary
}
}
#[derive(Clone, Debug)]
pub struct ToolExecutionResult {
pub summary: String,
pub payload: Value,
pub model_content_override: Option<String>,
pub model_image_parts: Vec<AgentContentPart>,
pub activity_event: Option<SessionActivityEvent>,
pub skip_source_elision: bool,
}
impl ToolExecutionResult {
pub fn from_activity_event(
summary: impl Into<String>,
payload: Value,
activity_event: Option<SessionActivityEvent>,
) -> Self {
Self {
summary: summary.into(),
payload,
model_content_override: None,
model_image_parts: Vec::new(),
activity_event,
skip_source_elision: false,
}
}
pub fn with_model_content(mut self, model_content: impl Into<String>) -> Self {
self.model_content_override = Some(model_content.into());
self
}
pub fn with_model_image_part(mut self, image: AgentContentPart) -> Self {
self.model_image_parts.push(image);
self
}
pub const fn with_skip_source_elision(mut self, skip: bool) -> Self {
self.skip_source_elision = skip;
self
}
pub fn model_content(&self) -> String {
if let Some(model_content) = &self.model_content_override {
return model_content.clone();
}
self.default_content_for_payload(&self.payload)
}
pub fn history_content(&self, tool_call_id: &str, tool_name: &str) -> String {
format!(
"tool_call_id={tool_call_id}\nname={tool_name}\n{}",
self.default_content_for_payload(&self.payload)
)
}
pub fn history_content_with_budget(
&self,
tool_call_id: &str,
tool_name: &str,
max_tokens: usize,
) -> String {
truncate_text_to_token_budget_with_notice(
&self.history_content(tool_call_id, tool_name),
max_tokens.max(1),
"... [tool output too long; history content truncated]",
)
}
fn default_content_for_payload(&self, payload: &Value) -> String {
if payload.is_null() {
format!("summary={}", self.summary)
} else {
format!(
"summary={}\npayload=\n{}",
self.summary,
serde_json::to_string_pretty(payload).unwrap_or_else(|_| payload.to_string())
)
}
}
fn ensure_model_content_with_budget(mut self, max_tokens: usize) -> Self {
if self.model_content_override.is_none() {
self.model_content_override = Some(truncate_text_to_token_budget_with_notice(
&self.default_content_for_payload(&self.payload),
max_tokens,
"... [tool output too long; model content truncated]",
));
}
self
}
}
#[async_trait]
pub trait RuntimeTool: Send + Sync {
fn name(&self) -> &str;
fn description(&self) -> &str;
fn input_spec(&self) -> AgentToolInputSpec;
fn app_tool_name(&self) -> Option<&str> {
None
}
fn is_available(&self, _: &Context) -> bool {
true
}
fn spec(&self) -> AgentToolSpec {
AgentToolSpec {
name: self.name().to_string(),
description: self.description().to_string(),
input_spec: self.input_spec(),
}
}
fn summarize_action(&self, call: &AgentToolCall) -> miette::Result<EpisodeActionRecord>;
fn call_activity_event(&self, call: &AgentToolCall) -> miette::Result<ToolCallActivityEvent>;
async fn execute(
&self,
context: &mut Context,
call: &AgentToolCall,
) -> miette::Result<ToolExecutionResult>;
}
struct StaticRuntimeTool {
name: &'static str,
description: &'static str,
input_spec: AgentToolInputSpec,
availability: Option<ToolAvailability>,
summarize: ToolSummarizer,
call_ui: ToolCallActivityBuilder,
execute: ToolExecutor,
}
impl StaticRuntimeTool {
fn new_with_schema(
name: &'static str,
description: &'static str,
schema: serde_json::Value,
summarize: ToolSummarizer,
call_ui: ToolCallActivityBuilder,
execute: ToolExecutor,
) -> Self {
Self {
name,
description,
input_spec: AgentToolInputSpec::JsonSchema {
schema: model_schema(schema),
},
availability: None,
summarize,
call_ui,
execute,
}
}
fn new_with_schema_and_availability(
name: &'static str,
description: &'static str,
schema: serde_json::Value,
availability: ToolAvailability,
summarize: ToolSummarizer,
call_ui: ToolCallActivityBuilder,
execute: ToolExecutor,
) -> Self {
Self {
name,
description,
input_spec: AgentToolInputSpec::JsonSchema {
schema: model_schema(schema),
},
availability: Some(availability),
summarize,
call_ui,
execute,
}
}
}
#[async_trait]
impl RuntimeTool for StaticRuntimeTool {
fn name(&self) -> &str {
self.name
}
fn description(&self) -> &str {
self.description
}
fn input_spec(&self) -> AgentToolInputSpec {
self.input_spec.clone()
}
fn is_available(&self, context: &Context) -> bool {
self.availability
.is_none_or(|availability| availability(context))
}
fn summarize_action(&self, call: &AgentToolCall) -> miette::Result<EpisodeActionRecord> {
(self.summarize)(call)
}
fn call_activity_event(&self, call: &AgentToolCall) -> miette::Result<ToolCallActivityEvent> {
(self.call_ui)(call)
}
async fn execute(
&self,
context: &mut Context,
call: &AgentToolCall,
) -> miette::Result<ToolExecutionResult> {
(self.execute)(context, call).await
}
}
struct AppRuntimeTool {
owner_app_id: AppId,
exposed_name: String,
app_tool_name: String,
description: String,
input_spec: AgentToolInputSpec,
}
const APP_GET_STATE_TOOL_NAME: &str = "get_state";
#[model_schema]
#[derive(Clone, Copy, Debug, Default, Deserialize, Serialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
enum AppStateDetail {
#[default]
Summary,
Full,
}
#[model_schema]
#[derive(Clone, Debug, Default, Deserialize, Serialize, JsonSchema)]
struct AppGetStateArgs {
detail: Option<AppStateDetail>,
}
struct AppGetStateRuntimeTool {
owner_app_id: AppId,
exposed_name: String,
input_spec: AgentToolInputSpec,
}
impl AppGetStateRuntimeTool {
fn new(owner_app_id: AppId, exposed_name: String) -> Self {
Self {
owner_app_id,
exposed_name,
input_spec: AgentToolInputSpec::JsonSchema {
schema: model_schema_for::<AppGetStateArgs>(),
},
}
}
}
fn app_state_payload(app_id: &AppId, state: &AppStateRender, detail: AppStateDetail) -> Value {
json!({
"app": app_id.to_string(),
"detail": match detail {
AppStateDetail::Summary => "summary",
AppStateDetail::Full => "full",
},
"state": state,
})
}
fn render_app_get_state_model_content(app_id: &AppId, state: &AppStateRender) -> String {
let mut out = String::new();
let _ = writeln!(out, "app={app_id}");
let _ = writeln!(out, "state_title={}", state.title);
out.push_str("state:\n");
if state.lines.is_empty() {
out.push_str("- no visible state\n");
} else {
for line in &state.lines {
out.push_str("- ");
out.push_str(line);
out.push('\n');
}
}
out.trim_end().to_string()
}
#[async_trait]
impl RuntimeTool for AppGetStateRuntimeTool {
fn name(&self) -> &str {
&self.exposed_name
}
fn app_tool_name(&self) -> Option<&str> {
Some(APP_GET_STATE_TOOL_NAME)
}
fn description(&self) -> &'static str {
"Read the current state for this app capability domain."
}
fn input_spec(&self) -> AgentToolInputSpec {
self.input_spec.clone()
}
fn summarize_action(&self, call: &AgentToolCall) -> miette::Result<EpisodeActionRecord> {
let args: AppGetStateArgs = parse_tool_args(call)?;
Ok(EpisodeActionRecord {
kind: self.exposed_name.clone(),
summary: format!(
"detail={}",
match args.detail.unwrap_or_default() {
AppStateDetail::Summary => "summary",
AppStateDetail::Full => "full",
}
),
})
}
fn call_activity_event(&self, call: &AgentToolCall) -> miette::Result<ToolCallActivityEvent> {
let args: AppGetStateArgs = parse_tool_args(call)?;
let lines = vec![format!(
"detail={}",
match args.detail.unwrap_or_default() {
AppStateDetail::Summary => "summary",
AppStateDetail::Full => "full",
}
)];
Ok(ToolCallActivityEvent::app(
AppId::render_exposed_tool_name(&self.exposed_name),
lines,
))
}
async fn execute(
&self,
context: &mut Context,
call: &AgentToolCall,
) -> miette::Result<ToolExecutionResult> {
let args: AppGetStateArgs = parse_tool_args(call)?;
let state = context
.apps
.state_render_for(&self.owner_app_id)
.ok_or_else(|| miette!("app missing for state read: {}", self.owner_app_id))?;
let payload =
app_state_payload(&self.owner_app_id, &state, args.detail.unwrap_or_default());
let model_content = render_app_get_state_model_content(&self.owner_app_id, &state);
Ok(ToolExecutionResult::from_activity_event(
format!("read {} state", self.owner_app_id),
payload,
Some(SessionActivityEvent::GenericApp(
TextActivityDescriptor {
title: AppId::render_exposed_tool_name(&self.exposed_name),
body_lines: state.lines,
}
.into(),
)),
)
.with_model_content(model_content))
}
}
#[async_trait]
impl RuntimeTool for AppRuntimeTool {
fn name(&self) -> &str {
&self.exposed_name
}
fn app_tool_name(&self) -> Option<&str> {
Some(&self.app_tool_name)
}
fn description(&self) -> &str {
&self.description
}
fn input_spec(&self) -> AgentToolInputSpec {
self.input_spec.clone()
}
fn is_available(&self, _context: &Context) -> bool {
true
}
fn summarize_action(&self, _call: &AgentToolCall) -> miette::Result<EpisodeActionRecord> {
context_free_error()?;
unreachable!()
}
fn call_activity_event(&self, _call: &AgentToolCall) -> miette::Result<ToolCallActivityEvent> {
context_free_error()?;
unreachable!()
}
async fn execute(
&self,
context: &mut Context,
call: &AgentToolCall,
) -> miette::Result<ToolExecutionResult> {
let app_context = AppToolExecutionContext {
execution_cwd: context.execution_cwd.clone(),
sandbox_policy: context.sandbox_policy.clone(),
dashboard_tx: context.dashboard_tx.clone(),
tool_output_max_tokens: context
.config
.main_model_config()
.tool_output_max_tokens
.max(1),
turn_epoch: context.runtime_turn_epoch,
};
let app_call = call.with_name(self.app_tool_name.clone());
let result = context
.apps
.execute_tool_for_app(&self.owner_app_id, &app_call, &app_context)
.await?;
let mut output = ToolExecutionResult::from_activity_event(
result.summary.clone(),
result.payload,
result.activity_event,
);
if let Some(model_content) = result.model_content {
output = output.with_model_content(model_content);
}
Ok(output)
}
}
struct WorkflowRuntimeTool {
workflow_id: String,
name: String,
input_spec: AgentToolInputSpec,
}
impl WorkflowRuntimeTool {
fn new(definition: &crate::workflow::WorkflowDefinition) -> Self {
Self {
workflow_id: definition.id.clone(),
name: definition.tool_name(),
input_spec: AgentToolInputSpec::JsonSchema {
schema: definition.input_schema.clone(),
},
}
}
}
#[async_trait]
impl RuntimeTool for WorkflowRuntimeTool {
fn name(&self) -> &str {
&self.name
}
fn description(&self) -> &'static str {
"Run this typed Lua workflow. It orchestrates isolated workers and returns the workflow's declared output."
}
fn input_spec(&self) -> AgentToolInputSpec {
self.input_spec.clone()
}
fn summarize_action(&self, call: &AgentToolCall) -> miette::Result<EpisodeActionRecord> {
Ok(EpisodeActionRecord {
kind: self.name.clone(),
summary: summarize_inline_text(&call.arguments.to_string()),
})
}
fn call_activity_event(&self, call: &AgentToolCall) -> miette::Result<ToolCallActivityEvent> {
Ok(ToolCallActivityEvent::app(
self.name.clone(),
vec![format!(
"input={}",
summarize_inline_text(&call.arguments.to_string())
)],
))
}
async fn execute(
&self,
context: &mut Context,
call: &AgentToolCall,
) -> miette::Result<ToolExecutionResult> {
let result = invoke_workflow(
context,
WorkflowInvocation {
workflow_id: self.workflow_id.clone(),
input: call.arguments.clone(),
},
)
.await?;
let payload = json!({
"workflow_id": result.workflow_id,
"status": result.status,
"output": result.output,
"message": result.message,
});
Ok(ToolExecutionResult::from_activity_event(
format!("workflow {}", self.workflow_id),
payload,
Some(workflow_activity_event(&result)),
))
}
}
fn workflow_activity_event(
result: &crate::workflow::WorkflowInvocationResult,
) -> SessionActivityEvent {
SessionActivityEvent::Workflow(crate::dashboard::WorkflowActivityData {
workflow_id: result.workflow_id.clone(),
status: result.status.clone(),
output: result.output.clone(),
message: result.message.clone(),
snapshot: Some(result.snapshot.clone()),
})
}
fn build_workflow_runtime_tools(
context: &Context,
reserved_names: &HashSet<String>,
) -> Vec<Box<dyn RuntimeTool>> {
let mut seen_names = reserved_names.clone();
let mut tools = Vec::new();
for definition in context.workflows.definitions() {
let tool = WorkflowRuntimeTool::new(definition);
if seen_names.insert(tool.name.clone()) {
tools.push(Box::new(tool) as Box<dyn RuntimeTool>);
} else {
tracing::warn!(
"skipping workflow `{}` because exposed tool `{}` conflicts with another runtime tool",
definition.id,
definition.tool_name()
);
}
}
tools
}
pub fn worker_finish_and_send_tool(output_schema: Value) -> Box<dyn RuntimeTool> {
Box::new(StaticRuntimeTool::new_with_schema(
"finish_and_send",
"Finish this isolated workflow worker by returning its declared typed output.",
output_schema,
|call| Ok(summarize_worker_finish_and_send_tool(call)),
|call| Ok(render_worker_finish_and_send_tool(call)),
execute_worker_finish_and_send_tool,
))
}
fn summarize_worker_finish_and_send_tool(call: &AgentToolCall) -> EpisodeActionRecord {
EpisodeActionRecord {
kind: "finish_and_send".to_string(),
summary: summarize_inline_text(&call.arguments.to_string()),
}
}
fn render_worker_finish_and_send_tool(call: &AgentToolCall) -> ToolCallActivityEvent {
ToolCallActivityEvent::app(
"finish_and_send",
vec![summarize_inline_text(&call.arguments.to_string())],
)
}
fn execute_worker_finish_and_send_tool<'a>(
_context: &'a mut Context,
call: &'a AgentToolCall,
) -> ToolFuture<'a> {
Box::pin(async move {
Ok(ToolExecutionResult::from_activity_event(
"worker completed",
call.arguments.clone(),
None,
))
})
}
fn build_worker_runtime_tools_for_apps(
apps: &AppManager,
output_schema: Value,
) -> Vec<Box<dyn RuntimeTool>> {
let mut tools = build_static_runtime_tools();
tools.retain(|tool| tool.name() != "finish_and_send");
tools.push(worker_finish_and_send_tool(output_schema));
let reserved_names = tools
.iter()
.map(|tool| tool.name().to_string())
.collect::<HashSet<_>>();
tools.extend(build_app_runtime_tools_for_apps(apps, &reserved_names));
tools
}
pub fn build_worker_runtime_tool_specs_for_apps(
apps: &AppManager,
output_schema: Value,
supports_vision: bool,
) -> Vec<AgentToolSpec> {
build_worker_runtime_tools_for_apps(apps, output_schema)
.into_iter()
.filter(|tool| worker_runtime_tool_is_available(tool.as_ref(), supports_vision))
.map(|tool| tool.spec())
.collect()
}
pub struct WorkerRuntimeToolCallContext<'a> {
pub(crate) execution_cwd: &'a std::path::Path,
pub(crate) sandbox_policy: &'a crate::sandbox::RuntimeSandboxPolicy,
pub(crate) tool_output_max_tokens: usize,
pub(crate) supports_vision: Option<bool>,
pub(crate) image_state_dir: &'a std::path::Path,
pub(crate) turn_epoch: u64,
pub(crate) output_schema: &'a Value,
pub(crate) worker_plan: &'a mut crate::plan::Plan,
}
pub async fn execute_worker_runtime_tool_call_for_apps(
apps: &mut AppManager,
call: &AgentToolCall,
context: WorkerRuntimeToolCallContext<'_>,
) -> Result<ToolExecutionResult> {
let WorkerRuntimeToolCallContext {
execution_cwd,
sandbox_policy,
tool_output_max_tokens,
supports_vision,
image_state_dir,
turn_epoch,
output_schema,
worker_plan,
} = context;
let tools = build_worker_runtime_tools_for_apps(apps, output_schema.clone());
let tool = find_runtime_tool(&tools, &call.name)?;
let worker_model_supports_vision = supports_vision.unwrap_or_else(|| {
crate::model_catalog::catalog_model_capacity("workflow-worker")
.is_none_or(|capacity| capacity.supports_vision)
});
if !worker_runtime_tool_is_available(tool, worker_model_supports_vision) {
return Ok(unavailable_worker_tool_result(
call,
worker_model_supports_vision,
));
}
let app_context = AppToolExecutionContext {
execution_cwd: execution_cwd.to_path_buf(),
sandbox_policy: sandbox_policy.clone(),
dashboard_tx: None,
tool_output_max_tokens: tool_output_max_tokens.max(1),
turn_epoch,
};
apps.before_runtime_tool_call(call, &app_context)?;
execute_worker_runtime_tool(
tool,
apps,
&app_context,
call,
worker_plan,
image_state_dir,
worker_model_supports_vision,
)
.await
.map(|result| result.ensure_model_content_with_budget(tool_output_max_tokens.max(1)))
}
fn build_app_runtime_tools_for_apps(
apps: &AppManager,
reserved_names: &HashSet<String>,
) -> Vec<Box<dyn RuntimeTool>> {
let mut tools: Vec<Box<dyn RuntimeTool>> = Vec::new();
let mut seen_names = reserved_names.clone();
for (owner_app_id, app_tools) in apps.all_tool_specs() {
let get_state_exposed_name = owner_app_id.mangle_tool_name(APP_GET_STATE_TOOL_NAME);
if seen_names.insert(get_state_exposed_name.clone()) {
tools.push(Box::new(AppGetStateRuntimeTool::new(
owner_app_id.clone(),
get_state_exposed_name,
)));
} else {
tracing::warn!(
"skipping generated state tool for app `{}` because exposed name `{}` conflicts with another runtime tool",
owner_app_id,
get_state_exposed_name
);
}
for tool in &app_tools {
if !is_valid_dynamic_tool_name(&tool.name) {
tracing::warn!(
"skipping app tool `{}` from app `{}` because its name must match [A-Za-z0-9_-]+",
tool.name,
owner_app_id
);
continue;
}
let exposed_name = owner_app_id.mangle_tool_name(&tool.name);
if !seen_names.insert(exposed_name.clone()) {
tracing::warn!(
"skipping app tool `{}` from app `{}` because exposed name `{}` conflicts with another runtime tool",
tool.name,
owner_app_id,
exposed_name
);
continue;
}
if let Err(err) = validate_model_facing_schema(&tool.input_schema) {
tracing::warn!(
"skipping app tool `{}` from app `{}` because its input schema is invalid: {}",
tool.name,
owner_app_id,
err
);
continue;
}
tools.push(Box::new(AppRuntimeTool {
owner_app_id: owner_app_id.clone(),
exposed_name,
app_tool_name: tool.name.clone(),
description: tool.description.clone(),
input_spec: AgentToolInputSpec::JsonSchema {
schema: tool.input_schema.clone(),
},
}));
}
}
tools
}
fn worker_runtime_tool_is_available(tool: &dyn RuntimeTool, supports_vision: bool) -> bool {
tool.name() != "view_image" || supports_vision
}
fn unavailable_worker_tool_result(
call: &AgentToolCall,
supports_vision: bool,
) -> ToolExecutionResult {
let reason = if call.name == "view_image" && !supports_vision {
"`view_image` is unavailable because this worker's selected model does not support image inputs."
} else {
"This worker tool is unavailable."
};
ToolExecutionResult::from_activity_event(
format!("{} unavailable", AppId::render_exposed_tool_name(&call.name)),
json!({
"available": false,
"tool": call.name,
"reason": reason,
"allowed_next_action": "Use another available worker tool.",
}),
None,
)
.with_model_content(format!(
"Tool unavailable: `{}`\nReason: {reason}\nAllowed next action: Use another available worker tool.",
AppId::render_exposed_tool_name(&call.name)
))
}
async fn execute_worker_runtime_tool(
tool: &dyn RuntimeTool,
apps: &mut AppManager,
app_context: &AppToolExecutionContext,
call: &AgentToolCall,
worker_plan: &mut crate::plan::Plan,
image_state_dir: &std::path::Path,
supports_vision: bool,
) -> Result<ToolExecutionResult> {
match tool.name() {
"read_file" => {
return files::execute_worker_read_file(
app_context.execution_cwd.as_path(),
&app_context.sandbox_policy,
call,
)
.await;
}
"edit_file" => {
return files::execute_worker_edit_file(
app_context.execution_cwd.as_path(),
&app_context.sandbox_policy,
call,
);
}
"update_plan" => return work::execute_worker_update_plan(worker_plan, call),
"view_image" => {
return view_image::execute_worker_view_image(
app_context.execution_cwd.as_path(),
&app_context.sandbox_policy,
image_state_dir,
supports_vision,
call,
);
}
"finish_and_send" => {
return Ok(ToolExecutionResult::from_activity_event(
"worker completed",
call.arguments.clone(),
None,
));
}
_ => {}
}
if let Some(app_tool_name) = tool.app_tool_name() {
let app_call = call.with_name(app_tool_name.to_string());
let app_id = if app_tool_name == APP_GET_STATE_TOOL_NAME {
apps.all_tool_specs()
.into_iter()
.find_map(|(app_id, _)| {
(app_id.mangle_tool_name(APP_GET_STATE_TOOL_NAME) == tool.name())
.then_some(app_id)
})
.ok_or_else(|| {
miette!("worker app state tool owner missing for `{}`", tool.name())
})?
} else {
apps.all_tool_specs()
.into_iter()
.find_map(|(app_id, specs)| {
specs
.iter()
.any(|spec| app_id.mangle_tool_name(&spec.name) == tool.name())
.then_some(app_id)
})
.ok_or_else(|| miette!("worker app tool owner missing for `{}`", tool.name()))?
};
if app_tool_name == APP_GET_STATE_TOOL_NAME {
let args: AppGetStateArgs = parse_tool_args(call)?;
let state = apps
.state_render_for(&app_id)
.ok_or_else(|| miette!("worker app state missing for {app_id}"))?;
let payload = app_state_payload(&app_id, &state, args.detail.unwrap_or_default());
let model_content = render_app_get_state_model_content(&app_id, &state);
return Ok(ToolExecutionResult::from_activity_event(
format!("read {app_id} state"),
payload,
None,
)
.with_model_content(model_content));
}
let result = apps
.execute_tool_for_app(&app_id, &app_call, app_context)
.await?;
let mut output = ToolExecutionResult::from_activity_event(
result.summary,
result.payload,
result.activity_event,
);
if let Some(model_content) = result.model_content {
output = output.with_model_content(model_content);
}
return Ok(output);
}
Err(miette!("unknown worker runtime tool `{}`", tool.name()))
}
fn build_static_runtime_tools() -> Vec<Box<dyn RuntimeTool>> {
let mut tools: Vec<Box<dyn RuntimeTool>> = Vec::new();
tools.extend(files::register_tools());
tools.extend(view_image::register_tools());
tools.extend(work::register_tools());
tools
}
fn build_app_runtime_tools(
context: &Context,
reserved_names: &HashSet<String>,
) -> Vec<Box<dyn RuntimeTool>> {
let mut tools: Vec<Box<dyn RuntimeTool>> = Vec::new();
let mut seen_names = reserved_names.clone();
for (owner_app_id, app_tools) in context.apps.all_tool_specs() {
let get_state_exposed_name = owner_app_id.mangle_tool_name(APP_GET_STATE_TOOL_NAME);
if seen_names.insert(get_state_exposed_name.clone()) {
tools.push(Box::new(AppGetStateRuntimeTool::new(
owner_app_id.clone(),
get_state_exposed_name,
)));
} else {
tracing::warn!(
"skipping generated state tool for app `{}` because exposed name `{}` conflicts with another runtime tool",
owner_app_id,
get_state_exposed_name
);
}
for tool in &app_tools {
if !is_valid_dynamic_tool_name(&tool.name) {
tracing::warn!(
"skipping app tool `{}` from app `{}` because its name must match [A-Za-z0-9_-]+",
tool.name,
owner_app_id
);
continue;
}
let exposed_name = owner_app_id.mangle_tool_name(&tool.name);
if !seen_names.insert(exposed_name.clone()) {
tracing::warn!(
"skipping app tool `{}` from app `{}` because exposed name `{}` conflicts with another runtime tool",
tool.name,
owner_app_id,
exposed_name
);
continue;
}
if let Err(err) = validate_model_facing_schema(&tool.input_schema) {
tracing::warn!(
"skipping app tool `{}` from app `{}` because its input schema is invalid: {}",
tool.name,
owner_app_id,
err
);
continue;
}
tools.push(Box::new(AppRuntimeTool {
owner_app_id: owner_app_id.clone(),
exposed_name,
app_tool_name: tool.name.clone(),
description: tool.description.clone(),
input_spec: AgentToolInputSpec::JsonSchema {
schema: tool.input_schema.clone(),
},
}));
}
}
tools
}
pub fn build_runtime_tools(context: &Context) -> Vec<Box<dyn RuntimeTool>> {
let mut tools = build_static_runtime_tools();
let mut reserved_names = tools
.iter()
.map(|tool| tool.name().to_string())
.collect::<HashSet<_>>();
let app_tools = build_app_runtime_tools(context, &reserved_names);
reserved_names.extend(app_tools.iter().map(|tool| tool.name().to_string()));
tools.extend(app_tools);
tools.extend(build_workflow_runtime_tools(context, &reserved_names));
tools
}
fn is_valid_dynamic_tool_name(name: &str) -> bool {
!name.is_empty()
&& name.len() <= 64
&& name
.chars()
.all(|ch| ch.is_ascii_alphanumeric() || matches!(ch, '_' | '-'))
}
fn context_free_error<T>() -> miette::Result<T> {
Err(miette!(
"app runtime tools require app-owned summarize/call-ui dispatch"
))
}
fn find_runtime_tool<'a>(
tools: &'a [Box<dyn RuntimeTool>],
name: &str,
) -> miette::Result<&'a dyn RuntimeTool> {
let canonical_name = canonical_runtime_tool_name(name);
tools
.iter()
.find(|tool| tool.name() == canonical_name)
.map(std::convert::AsRef::as_ref)
.ok_or_else(|| miette!("unknown runtime tool: {name}"))
}
pub fn build_runtime_tool_specs(context: &Context) -> Vec<AgentToolSpec> {
build_runtime_tools(context)
.into_iter()
.filter(|tool| tool.is_available(context))
.map(|tool| tool.spec())
.collect()
}
fn runtime_availability_denial(
context: &Context,
tool: &dyn RuntimeTool,
) -> Option<(String, String)> {
if tool.is_available(context) {
return None;
}
let name = tool.name();
Some((
format!("`{name}` is disabled by the current runtime availability policy."),
"Use a currently allowed tool, or satisfy the tool's required runtime state before retrying."
.to_string(),
))
}
fn unavailable_tool_result(
call: &AgentToolCall,
reason: String,
allowed_next_action: String,
) -> ToolExecutionResult {
let display_tool_name = AppId::render_exposed_tool_name(&call.name);
let model_content = format!(
"Tool unavailable: `{display_tool_name}`\nReason: {reason}\nAllowed next action: {allowed_next_action}"
);
ToolExecutionResult::from_activity_event(
format!("{display_tool_name} unavailable"),
json!({
"available": false,
"tool": call.name,
"reason": reason,
"allowed_next_action": allowed_next_action,
}),
Some(SessionActivityEvent::Error(
TextActivityDescriptor {
title: format!("{display_tool_name} unavailable"),
body_lines: vec![reason, allowed_next_action],
}
.into(),
)),
)
.with_model_content(model_content)
}
pub fn summarize_action_from_tool_call(
context: &Context,
call: &AgentToolCall,
) -> Result<EpisodeActionRecord> {
let tools = build_runtime_tools(context);
let tool = find_runtime_tool(&tools, &call.name)?;
let tool_call = tool_call_for_runtime_tool(tool, call);
tool.summarize_action(&tool_call)
.map_or_else(|_| context.apps.summarize_tool_call(call), Ok)
}
pub fn build_tool_call_activity_event(
context: &Context,
call: &AgentToolCall,
) -> Result<ToolCallActivityEvent> {
build_tool_call_activity_event_from_tools(&build_runtime_tools(context), call, &context.apps)
}
fn build_tool_call_activity_event_from_tools(
tools: &[Box<dyn RuntimeTool>],
call: &AgentToolCall,
apps: &AppManager,
) -> Result<ToolCallActivityEvent> {
let tool = find_runtime_tool(tools, &call.name)?;
let tool_call = tool_call_for_runtime_tool(tool, call);
tool.call_activity_event(&tool_call)
.map_or_else(|_| apps.tool_call_activity_event(call), Ok)
}
fn tool_call_for_runtime_tool(tool: &dyn RuntimeTool, call: &AgentToolCall) -> AgentToolCall {
tool.app_tool_name().map_or_else(
|| call.clone(),
|app_tool_name| call.with_name(app_tool_name),
)
}
pub fn render_telegram_tool_result_status(
call: &AgentToolCall,
result: &ToolExecutionResult,
) -> Option<TelegramLiveStatus> {
let tool_name = demangle_known_app_tool_name(&call.name);
if telegram_status_ignored_tool(tool_name) {
return None;
}
if matches!(result.activity_event, Some(SessionActivityEvent::Error(_))) {
return telegram_tool_failure_status(tool_name);
}
match tool_name {
"update_plan" => Some(telegram_status(glyph::PLAN, "Plan Updated")),
"edit_file" => match &result.activity_event {
Some(SessionActivityEvent::Patch(event)) => Some(telegram_status(
glyph::PATCH,
format!(
"Edited {} {}",
event.files.len(),
plural_noun(event.files.len(), "File", "Files")
),
)),
Some(SessionActivityEvent::CodingEdit(event)) => Some(telegram_status(
glyph::PATCH,
format!(
"Edited {} {}",
event.diff_files.len(),
plural_noun(event.diff_files.len(), "File", "Files")
),
)),
_ => Some(telegram_status(glyph::PATCH, "Edited Files")),
},
"terminal_exec" => {
if result
.payload
.get("running")
.and_then(Value::as_bool)
.unwrap_or(false)
{
Some(telegram_status(glyph::EXEC, "Command Running"))
} else {
Some(telegram_status(glyph::EXEC, "Command Ran"))
}
}
"terminal_write_stdin" => Some(telegram_status(glyph::EXEC, "Terminal Continued")),
"terminal_terminate" => Some(telegram_status(glyph::EXEC, "Terminal Stopped")),
"browser_open_page" => Some(telegram_status(glyph::BROWSER, "Browser Opened")),
"browser_snapshot" => Some(telegram_status(glyph::BROWSER, "Browser Read")),
"browser_wait" => Some(telegram_status(glyph::BROWSER, "Browser Waited")),
"browser_click" | "browser_fill" => Some(telegram_status(glyph::BROWSER, "Browser Acted")),
"browser_back" | "browser_forward" => {
Some(telegram_status(glyph::BROWSER, "Browser Navigated"))
}
"browser_reload" => Some(telegram_status(glyph::BROWSER, "Browser Reloaded")),
"browser_close_page" => Some(telegram_status(glyph::BROWSER, "Browser Closed")),
_ => Some(telegram_status(glyph::EXEC, "App Updated")),
}
}
fn demangle_known_app_tool_name(tool_name: &str) -> &str {
if let Some((app_id, app_tool_name)) = tool_name.split_once(AppId::TOOL_NAME_SEPARATOR)
&& AppId::is_valid_name(app_id)
{
return canonical_runtime_tool_name(app_tool_name);
}
canonical_runtime_tool_name(tool_name)
}
const fn canonical_runtime_tool_name(tool_name: &str) -> &str {
tool_name
}
fn telegram_status_ignored_tool(tool_name: &str) -> bool {
matches!(tool_name, "finish_and_send")
}
fn telegram_tool_failure_status(tool_name: &str) -> Option<TelegramLiveStatus> {
match tool_name {
"finish_and_send" => None,
"update_plan" => Some(telegram_status(glyph::ERROR, "Plan Update Failed")),
"edit_file" => Some(telegram_status(glyph::ERROR, "File Edit Failed")),
"terminal_exec" => Some(telegram_status(glyph::ERROR, "Command Failed")),
"terminal_write_stdin" => Some(telegram_status(glyph::ERROR, "Terminal Write Failed")),
"terminal_terminate" => Some(telegram_status(glyph::ERROR, "Terminal Stop Failed")),
"browser_open_page" | "browser_snapshot" | "browser_wait" | "browser_click"
| "browser_fill" | "browser_back" | "browser_forward" | "browser_reload"
| "browser_close_page" => Some(telegram_status(glyph::ERROR, "Browser Action Failed")),
_ => Some(telegram_status(glyph::ERROR, "App Failed")),
}
}
fn telegram_status(icon: impl Into<String>, text: impl Into<String>) -> TelegramLiveStatus {
TelegramLiveStatus {
icon: icon.into(),
text: text.into(),
}
}
const fn plural_noun(count: usize, singular: &'static str, plural: &'static str) -> &'static str {
if count == 1 { singular } else { plural }
}
pub async fn execute_agent_tool_call(
context: &mut Context,
call: &AgentToolCall,
) -> Result<ToolExecutionResult> {
let tools = build_runtime_tools(context);
let tool = find_runtime_tool(&tools, &call.name)?;
if let Some((reason, allowed_next_action)) = runtime_availability_denial(context, tool) {
return Ok(unavailable_tool_result(call, reason, allowed_next_action));
}
let app_context = AppToolExecutionContext {
execution_cwd: context.execution_cwd.clone(),
sandbox_policy: context.sandbox_policy.clone(),
dashboard_tx: context.dashboard_tx.clone(),
tool_output_max_tokens: context
.config
.main_model_config()
.tool_output_max_tokens
.max(1),
turn_epoch: context.runtime_turn_epoch,
};
context.apps.before_runtime_tool_call(call, &app_context)?;
let result = tool.execute(context, call).await?;
Ok(result.ensure_model_content_with_budget(
context
.config
.main_model_config()
.tool_output_max_tokens
.max(1),
))
}
#[cfg(test)]
mod tests {
use super::*;
use std::collections::HashMap;
use async_trait::async_trait;
use tempfile::TempDir;
use crate::{
app::{App, AppManager},
browser_app::BrowserApp,
coding_app::CodingApp,
config::Config,
context::Context,
context_budget::TokenEstimateBaseline,
core::{ModelProvider, ModelRequestOptions},
events::EventStore,
memory::Memory,
openskills::OpenSkillsCatalog,
pending_work::PendingWorkQueue,
plan::Plan,
reasoning::{
compiled::CompiledPromptStore,
runtime::{AgentTurnRequest, AgentTurnStreamResult, PromptRequest},
},
runtime::bootstrap::DaatLocusHomeOverride,
sandbox::RuntimeSandboxPolicy,
telegram_acl::TelegramAclHandle,
telegram_transport::state::TelegramTransportState,
terminal_app::TerminalApp,
workspace_app::WorkspaceAppRegistry,
};
#[test]
fn worker_finish_and_send_uses_declared_output_schema() {
let output_schema = json!({
"type": "object",
"properties": {
"summary": { "type": "string" }
},
"required": ["summary"],
"additionalProperties": false,
});
let AgentToolInputSpec::JsonSchema { schema } =
worker_finish_and_send_tool(output_schema.clone()).input_spec()
else {
panic!("worker finish_and_send should use a JSON schema");
};
assert_eq!(schema, output_schema);
assert!(!json_contains_key(&schema, "disposition"));
assert!(!json_contains_key(&schema, "reply_message"));
}
#[test]
fn worker_runtime_tools_exclude_main_finish_and_workflow_tools() {
let output_schema = json!({
"type": "object",
"properties": {},
"required": [],
"additionalProperties": false,
});
let apps: Vec<Box<dyn App>> = vec![Box::new(BrowserApp::new())];
let apps = AppManager::new(apps).expect("app manager");
let names = build_worker_runtime_tool_specs_for_apps(&apps, output_schema, false)
.into_iter()
.map(|tool| tool.name)
.collect::<Vec<_>>();
assert!(names.iter().any(|name| name == "finish_and_send"));
assert!(names.iter().any(|name| name == "read_file"));
assert!(names.iter().any(|name| name == "edit_file"));
assert!(names.iter().any(|name| name == "update_plan"));
assert!(names.iter().any(|name| name == "browser__get_state"));
assert!(!names.iter().any(|name| name.starts_with("workflow__")));
assert!(!names.iter().any(|name| name == "view_image"));
let vision_names = build_worker_runtime_tool_specs_for_apps(
&apps,
json!({
"type": "object",
"properties": {},
"required": [],
"additionalProperties": false,
}),
true,
)
.into_iter()
.map(|tool| tool.name)
.collect::<Vec<_>>();
assert!(vision_names.iter().any(|name| name == "view_image"));
}
#[tokio::test]
async fn worker_static_tools_execute_without_main_context_state() {
let home = tempfile::tempdir().expect("home");
let _home = DaatLocusHomeOverride::set(home.path().to_path_buf()).await;
let execution = tempfile::tempdir().expect("execution cwd");
std::fs::write(execution.path().join("notes.txt"), "alpha\nbeta\n").expect("write fixture");
let mut apps = AppManager::new(vec![Box::new(BrowserApp::new()) as Box<dyn App>])
.expect("worker apps");
let mut worker_plan = Plan::default();
let output_schema = json!({
"type": "object",
"properties": {},
"required": [],
"additionalProperties": false,
});
let read_call = AgentToolCall {
id: "worker-read".to_string(),
name: "read_file".to_string(),
arguments: json!({
"path": "notes.txt",
"start_line": 2,
"line_count": 1,
"force_no_elide": false,
}),
};
let image_state_dir = execution.path().join("images");
let result = execute_worker_runtime_tool_call_for_apps(
&mut apps,
&read_call,
WorkerRuntimeToolCallContext {
execution_cwd: execution.path(),
sandbox_policy: &RuntimeSandboxPolicy::disabled(),
tool_output_max_tokens: 1024,
supports_vision: Some(false),
image_state_dir: &image_state_dir,
turn_epoch: 1,
output_schema: &output_schema,
worker_plan: &mut worker_plan,
},
)
.await
.expect("worker read_file");
assert_eq!(
result.model_content(),
format!("2#{}|beta", scope_engine::patch::line_hash("beta"))
);
let plan_call = AgentToolCall {
id: "worker-plan".to_string(),
name: "update_plan".to_string(),
arguments: json!({
"explanation": "Track the worker task.",
"plan": [{ "step": "Inspect fixture", "status": "in_progress" }],
}),
};
execute_worker_runtime_tool_call_for_apps(
&mut apps,
&plan_call,
WorkerRuntimeToolCallContext {
execution_cwd: execution.path(),
sandbox_policy: &RuntimeSandboxPolicy::disabled(),
tool_output_max_tokens: 1024,
supports_vision: Some(false),
image_state_dir: &image_state_dir,
turn_epoch: 2,
output_schema: &output_schema,
worker_plan: &mut worker_plan,
},
)
.await
.expect("worker update_plan");
assert_eq!(worker_plan.steps().len(), 1);
assert_eq!(worker_plan.steps()[0].step, "Inspect fixture");
}
struct UnusedModelProvider;
#[async_trait]
impl ModelProvider for UnusedModelProvider {
async fn complete_json(
&self,
_request: PromptRequest,
_options: ModelRequestOptions,
) -> Result<serde_json::Value> {
Err(miette!("unused test model provider"))
}
async fn complete_agent_turn(
&self,
_request: AgentTurnRequest,
_options: ModelRequestOptions,
) -> Result<AgentTurnStreamResult> {
Err(miette!("unused test model provider"))
}
fn request_budget_limits(&self) -> crate::context_budget::RequestBudgetLimits {
crate::context_budget::RequestBudgetLimits {
context_window_tokens: crate::context_budget::DEFAULT_CONTEXT_WINDOW_TOKENS,
auto_compact_threshold_tokens: crate::context_budget::DEFAULT_CONTEXT_WINDOW_TOKENS,
reserved_output_tokens: crate::context_budget::DEFAULT_MAX_COMPLETION_TOKENS,
}
}
fn token_usage_info(&self) -> crate::core::TokenUsageInfo {
crate::core::TokenUsageInfo::default()
}
fn model_name(&self) -> String {
"unused-test-model-provider".to_string()
}
}
struct IsolatedTestContext {
context: Context,
_home_override: DaatLocusHomeOverride,
_home: TempDir,
_execution: TempDir,
}
impl IsolatedTestContext {
async fn new() -> Self {
let home = tempfile::tempdir().expect("test home");
let execution = tempfile::tempdir().expect("test execution cwd");
let home_override = DaatLocusHomeOverride::set(home.path().to_path_buf()).await;
let config = Config::default();
let telegram = TelegramTransportState::new();
let (daemon_control_tx, _daemon_control_rx) = tokio::sync::mpsc::unbounded_channel();
let apps: Vec<Box<dyn App>> = vec![
Box::new(BrowserApp::new()),
Box::new(TerminalApp::new()),
Box::new(CodingApp::new()),
];
let apps = AppManager::new(apps).expect("app manager");
let context = Context {
session_id: None,
model_provider: Box::new(UnusedModelProvider),
efficient_model_provider: std::sync::Arc::new(UnusedModelProvider),
config,
memory: Memory::new().await,
plan: Plan::new().await,
events: EventStore::new().await,
pending_work: PendingWorkQueue::new().await,
openskills: OpenSkillsCatalog::default(),
workflows: crate::workflow::WorkflowCatalog::load(),
workflow_cancellation: crate::workflow::WorkflowCancellationRegistry::default(),
active_skill_run: None,
pending_skill_run_flushes: Vec::new(),
current_work_origin: None,
apps,
workspace_apps: WorkspaceAppRegistry::default(),
telegram: telegram.handle(),
telegram_acl: TelegramAclHandle::load().await,
compiled_prompts: CompiledPromptStore::from_entries(Vec::new()),
execution_cwd: execution.path().to_path_buf(),
coding_project_dir: None,
sandbox_policy: RuntimeSandboxPolicy::disabled(),
dashboard_tx: None,
dashboard_history: None,
daemon_control_tx,
latest_context_composition: None,
active_runtime_turn: false,
active_runtime_phase: None,
runtime_turn_started_at: None,
runtime_turn_started_at_ms: None,
runtime_turn_epoch: 0,
runtime_overflow_failures: std::sync::Arc::new(parking_lot::Mutex::new(
HashMap::new(),
)),
runtime_model_request_failures: std::sync::Arc::new(parking_lot::Mutex::new(
HashMap::new(),
)),
live_progress_tx: std::sync::Arc::new(parking_lot::Mutex::new(None)),
telegram_live_drafts: std::sync::Arc::new(parking_lot::Mutex::new(HashMap::new())),
claimed_event_ids: Vec::new(),
afterclaim_context_fingerprint: None,
visible_source_lines: HashSet::new(),
delivered_root_instruction_fingerprint: None,
idle_since: None,
last_idle_sleep_at: None,
session_title: crate::runtime::session_title::SessionTitleState::default(),
token_estimate_baseline: TokenEstimateBaseline::default(),
};
Self {
context,
_home_override: home_override,
_home: home,
_execution: execution,
}
}
}
fn json_contains_key(value: &Value, needle: &str) -> bool {
match value {
Value::Object(object) => {
object.contains_key(needle)
|| object
.values()
.any(|value| json_contains_key(value, needle))
}
Value::Array(values) => values.iter().any(|value| json_contains_key(value, needle)),
_ => false,
}
}
fn tool_result(
tool_name: &str,
payload: Value,
activity_event: Option<SessionActivityEvent>,
) -> ToolExecutionResult {
ToolExecutionResult::from_activity_event(
format!("{tool_name} summary"),
payload,
activity_event,
)
}
#[test]
fn telegram_tool_status_renders_plan_update_without_steps() {
let call = AgentToolCall {
id: "call_1".to_string(),
name: "update_plan".to_string(),
arguments: serde_json::json!({}),
};
let result = tool_result("update_plan", serde_json::json!({}), None);
let status = render_telegram_tool_result_status(&call, &result).unwrap();
assert_eq!(status.icon, glyph::PLAN);
assert_eq!(status.text, "Plan Updated");
}
#[test]
fn telegram_tool_status_hides_final_reply_tool() {
let call = AgentToolCall {
id: "call_1".to_string(),
name: "finish_and_send".to_string(),
arguments: serde_json::json!({
"disposition": "resolved",
"reply_message": "done",
}),
};
let result = tool_result("finish_and_send", serde_json::json!({}), None);
assert!(render_telegram_tool_result_status(&call, &result).is_none());
}
#[test]
fn telegram_tool_status_renders_terminal_running_and_finished() {
let call = AgentToolCall {
id: "call_1".to_string(),
name: "terminal__terminal_exec".to_string(),
arguments: serde_json::json!({}),
};
let running = tool_result(
"terminal_exec",
serde_json::json!({ "running": true }),
None,
);
let finished = tool_result(
"terminal_exec",
serde_json::json!({ "running": false }),
None,
);
assert_eq!(
render_telegram_tool_result_status(&call, &running)
.unwrap()
.text,
"Command Running"
);
assert_eq!(
render_telegram_tool_result_status(&call, &finished)
.unwrap()
.text,
"Command Ran"
);
}
#[tokio::test]
async fn terminal_write_stdin_tool_schema_does_not_use_schema_composition() {
let isolated = IsolatedTestContext::new().await;
let spec = build_runtime_tool_specs(&isolated.context)
.into_iter()
.find(|tool| tool.name == "terminal__terminal_write_stdin")
.expect("terminal write stdin tool");
let AgentToolInputSpec::JsonSchema { schema } = spec.input_spec else {
panic!("terminal_write_stdin should use json schema");
};
drop(isolated);
crate::schema_utils::validate_model_facing_schema(&schema).unwrap();
for key in ["oneOf", "anyOf", "allOf"] {
assert!(!json_contains_key(&schema, key), "{schema:#}");
}
}
#[tokio::test]
async fn coding_read_code_tool_schema_does_not_use_schema_composition() {
let isolated = IsolatedTestContext::new().await;
let spec = build_runtime_tool_specs(&isolated.context)
.into_iter()
.find(|tool| tool.name == "coding__read_code")
.expect("coding read code tool");
let AgentToolInputSpec::JsonSchema { schema } = spec.input_spec else {
panic!("coding_read_code should use json schema");
};
drop(isolated);
crate::schema_utils::validate_model_facing_schema(&schema).unwrap();
for key in ["oneOf", "anyOf", "allOf"] {
assert!(!json_contains_key(&schema, key), "{schema:#}");
}
assert!(!json_contains_key(&schema, "ref"), "{schema:#}");
assert!(!json_contains_key(&schema, "handle"), "{schema:#}");
assert!(json_contains_key(&schema, "path"), "{schema:#}");
assert!(json_contains_key(&schema, "anchor"), "{schema:#}");
assert!(json_contains_key(&schema, "mode"), "{schema:#}");
assert!(!json_contains_key(&schema, "start_line"), "{schema:#}");
assert!(!json_contains_key(&schema, "line_count"), "{schema:#}");
}
#[tokio::test]
async fn coding_search_code_tool_schema_exposes_rg_aligned_options() {
let isolated = IsolatedTestContext::new().await;
let spec = build_runtime_tool_specs(&isolated.context)
.into_iter()
.find(|tool| tool.name == "coding__search_code")
.expect("coding search code tool");
let AgentToolInputSpec::JsonSchema { schema } = spec.input_spec else {
panic!("coding_search_code should use json schema");
};
drop(isolated);
crate::schema_utils::validate_model_facing_schema(&schema).unwrap();
let properties = schema
.get("properties")
.and_then(serde_json::Value::as_object)
.unwrap_or_else(|| panic!("schema should have object properties: {schema:#}"));
for key in [
"query",
"mode",
"path",
"include",
"exclude",
"types",
"type_not",
"case",
"word",
"whole_line",
"hidden",
"respect_ignore",
"follow",
"limit",
] {
assert!(properties.contains_key(key), "missing {key}: {schema:#}");
}
assert!(
!properties.contains_key("case_mode"),
"schema should expose `case`, not internal field name: {schema:#}"
);
let required = schema
.get("required")
.and_then(serde_json::Value::as_array)
.unwrap_or_else(|| panic!("schema should have required fields: {schema:#}"));
assert_eq!(required.len(), properties.len(), "{schema:#}");
for key in ["include", "exclude", "types", "type_not"] {
assert_eq!(
properties
.get(key)
.and_then(|value| value.get("type"))
.and_then(serde_json::Value::as_str),
Some("array"),
"{key} should be an array: {schema:#}"
);
}
}
#[tokio::test]
async fn structured_edit_tool_schemas_do_not_use_schema_composition() {
let isolated = IsolatedTestContext::new().await;
let specs = build_runtime_tool_specs(&isolated.context);
drop(isolated);
for tool_name in ["edit_file", "coding__edit_code"] {
let spec = specs
.iter()
.find(|tool| tool.name == tool_name)
.unwrap_or_else(|| panic!("{tool_name} tool"));
let AgentToolInputSpec::JsonSchema { schema } = &spec.input_spec else {
panic!("{tool_name} should use json schema");
};
crate::schema_utils::validate_model_facing_schema(schema).unwrap();
for key in ["oneOf", "anyOf", "allOf"] {
assert!(
!json_contains_key(schema, key),
"tool={tool_name} schema={schema:#}"
);
}
}
}
#[tokio::test]
async fn exposed_runtime_tool_schemas_follow_model_facing_dialect() {
let isolated = IsolatedTestContext::new().await;
for spec in build_runtime_tool_specs(&isolated.context) {
match spec.input_spec {
AgentToolInputSpec::JsonSchema { schema } => {
crate::schema_utils::validate_model_facing_schema(&schema).unwrap_or_else(
|err| panic!("tool={} schema={schema:#}\n{err}", spec.name),
);
}
AgentToolInputSpec::FreeformGrammar {
fallback_schema, ..
} => {
crate::schema_utils::validate_model_facing_schema(&fallback_schema)
.unwrap_or_else(|err| {
panic!(
"tool={} fallback_schema={fallback_schema:#}\n{err}",
spec.name
)
});
}
}
}
}
#[tokio::test]
async fn read_file_returns_line_hash_anchored_lines() {
let mut isolated = IsolatedTestContext::new().await;
let root = isolated.context.execution_cwd.clone();
std::fs::write(root.join("notes.txt"), "alpha\nbeta\ngamma\n").expect("write fixture");
let call = AgentToolCall {
id: "call_read".to_string(),
name: "read_file".to_string(),
arguments: json!({
"path": "notes.txt",
"start_line": 2,
"line_count": 1,
}),
};
let result = execute_agent_tool_call(&mut isolated.context, &call)
.await
.expect("read file");
drop(isolated);
assert_eq!(
result.model_content(),
format!("2#{}|beta", scope_engine::patch::line_hash("beta"))
);
let Some(SessionActivityEvent::Explored(ui_event)) = &result.activity_event else {
panic!("read_file should render as explored activity");
};
assert_eq!(
ui_event.stable_id,
crate::activity_event::EXPLORED_STABLE_ID
);
assert_eq!(ui_event.calls.len(), 1);
assert_eq!(ui_event.calls[0].tool_name, "Read");
assert_eq!(ui_event.calls[0].summary, "notes.txt#L2-L2");
}
#[tokio::test]
async fn edit_file_applies_structured_line_hash_edits() {
let mut isolated = IsolatedTestContext::new().await;
let root = isolated.context.execution_cwd.clone();
std::fs::write(root.join("README.md"), "old\n").expect("write markdown fixture");
let hash = scope_engine::patch::line_hash("old");
let call = AgentToolCall {
id: "call_edit".to_string(),
name: "edit_file".to_string(),
arguments: json!({
"edits": [{
"path": "README.md",
"op": "replace",
"start": format!("1#{hash}"),
"end": format!("1#{hash}"),
"content": "new"
}]
}),
};
let result = execute_agent_tool_call(&mut isolated.context, &call)
.await
.expect("edit file");
assert_eq!(
std::fs::read_to_string(root.join("README.md")).expect("read markdown fixture"),
"new\n"
);
drop(isolated);
let Some(SessionActivityEvent::CodingEdit(ui_event)) = &result.activity_event else {
panic!("edit_file should render as coding edit activity");
};
assert_eq!(ui_event.title, "Edited File");
assert_eq!(ui_event.tool_name.as_deref(), Some("edit_file"));
assert_eq!(ui_event.tool_app.as_deref(), Some("Workspace"));
assert_eq!(ui_event.file.as_deref(), Some("README.md"));
assert_eq!(ui_event.added_lines, 1);
assert_eq!(ui_event.removed_lines, 1);
assert_eq!(ui_event.propagation_count, 0);
assert_eq!(ui_event.diff_files.len(), 1);
}
#[tokio::test]
async fn all_app_tools_are_exposed_by_namespace() {
let isolated = IsolatedTestContext::new().await;
let specs = build_runtime_tool_specs(&isolated.context);
drop(isolated);
let names = specs
.into_iter()
.map(|tool| tool.name)
.collect::<HashSet<_>>();
assert!(names.contains("read_file"));
assert!(names.contains("edit_file"));
assert!(names.contains("view_image"));
assert!(names.contains("browser__get_state"));
assert!(names.contains("browser__browser_open_page"));
assert!(names.contains("browser__browser_snapshot"));
assert!(names.contains("coding__get_state"));
assert!(names.contains("coding__open_project"));
assert!(names.contains("coding__search_code"));
assert!(names.contains("coding__read_code"));
assert!(names.contains("coding__edit_code"));
assert!(names.contains("coding__next_review"));
assert!(names.contains("terminal__get_state"));
assert!(names.contains("terminal__terminal_exec"));
assert!(names.contains("terminal__terminal_write_stdin"));
assert!(names.contains("terminal__terminal_terminate"));
assert!(!names.contains("apply_patch"));
assert!(!names.contains("coding__grep"));
assert!(!names.contains("coding__glob"));
assert!(!names.contains("open_project"));
assert!(!names.contains("coding__coding_open_project"));
assert!(!names.contains("coding_open_project"));
assert!(!names.contains("terminal_exec"));
assert!(!names.contains("browser_open_page"));
}
#[tokio::test]
async fn view_image_attaches_a_supported_local_file() {
let mut isolated = IsolatedTestContext::new().await;
let root = isolated.context.execution_cwd.clone();
let model_name = isolated.context.config.main_model.clone();
isolated
.context
.config
.models
.get_mut(&model_name)
.expect("default model")
.supports_vision = Some(true);
let mut bytes = b"\x89PNG\r\n\x1a\nfixture".to_vec();
bytes.extend_from_slice(b"image-data");
std::fs::write(root.join("diagram.png"), &bytes).expect("write image fixture");
let call = AgentToolCall {
id: "call_view_image".to_string(),
name: "view_image".to_string(),
arguments: json!({ "path": "diagram.png" }),
};
let result = execute_agent_tool_call(&mut isolated.context, &call)
.await
.expect("view image");
drop(isolated);
assert_eq!(result.payload["media_type"], "image/png");
assert_eq!(result.payload["bytes"], bytes.len());
assert_eq!(result.model_image_parts.len(), 1);
assert!(matches!(
&result.model_image_parts[0],
crate::reasoning::runtime::AgentContentPart::Image { media_type, .. }
if media_type == "image/png"
));
}
#[tokio::test]
async fn view_image_rejects_a_non_vision_model() {
let mut isolated = IsolatedTestContext::new().await;
let root = isolated.context.execution_cwd.clone();
std::fs::write(root.join("diagram.png"), b"\x89PNG\r\n\x1a\nfixture")
.expect("write image fixture");
let model_name = isolated.context.config.main_model.clone();
isolated
.context
.config
.models
.get_mut(&model_name)
.expect("default model")
.supports_vision = Some(false);
let call = AgentToolCall {
id: "call_view_image".to_string(),
name: "view_image".to_string(),
arguments: json!({ "path": "diagram.png" }),
};
let result = execute_agent_tool_call(&mut isolated.context, &call)
.await
.expect("unavailable tool result");
drop(isolated);
assert_eq!(result.payload["available"], false);
assert!(
result
.model_content()
.contains("disabled by the current runtime availability policy")
);
}
#[test]
fn tool_output_images_are_separate_from_the_tool_protocol_message() {
let result = ToolExecutionResult::from_activity_event("image", json!({}), None)
.with_model_image_part(crate::reasoning::runtime::AgentContentPart::Image {
path: "/tmp/diagram.png".to_string(),
media_type: "image/png".to_string(),
description: Some("diagram".to_string()),
});
let tool_message = crate::reasoning::runtime::AgentMessage::tool(
"call_image",
"view_image",
result.model_content(),
);
let attachment_message = crate::reasoning::runtime::AgentMessage::user_content(
crate::reasoning::runtime::AgentContent::multimodal(
"The `view_image` tool attached image content for visual inspection.",
result.model_image_parts,
),
);
assert!(matches!(
tool_message,
crate::reasoning::runtime::AgentMessage::Tool { .. }
));
let crate::reasoning::runtime::AgentMessage::User { content } = attachment_message else {
panic!("expected image attachment user message");
};
assert_eq!(content.parts().len(), 1);
assert!(matches!(
&content.parts()[0],
crate::reasoning::runtime::AgentContentPart::Image { media_type, .. }
if media_type == "image/png"
));
}
#[tokio::test]
async fn namespaced_terminal_tool_executes_directly() {
let mut isolated = IsolatedTestContext::new().await;
let call = AgentToolCall {
id: "call_1".to_string(),
name: "terminal__terminal_exec".to_string(),
arguments: json!({
"command": "printf '%s\\n' delegated-terminal",
"yield_time_ms": 100,
}),
};
let result = execute_agent_tool_call(&mut isolated.context, &call)
.await
.unwrap();
drop(isolated);
assert!(
result.model_content().contains("delegated-terminal"),
"model content was: {}",
result.model_content()
);
assert!(matches!(
result.activity_event,
Some(SessionActivityEvent::ExecResult(_))
));
}
#[tokio::test]
async fn project_scope_rejects_file_tool_for_scope_owned_source() {
let mut isolated = IsolatedTestContext::new().await;
let root = isolated.context.execution_cwd.clone();
std::fs::write(root.join("lib.rs"), "pub fn value() -> i32 {\n 1\n}\n")
.expect("write rust fixture");
let open_call = AgentToolCall {
id: "call_open".to_string(),
name: "coding__open_project".to_string(),
arguments: json!({
"project_root": root,
}),
};
execute_agent_tool_call(&mut isolated.context, &open_call)
.await
.expect("open project");
let hash = scope_engine::patch::line_hash(" 1");
let edit_call = AgentToolCall {
id: "call_edit".to_string(),
name: "edit_file".to_string(),
arguments: json!({
"edits": [{
"path": "lib.rs",
"op": "replace",
"start": format!("2#{hash}"),
"end": format!("2#{hash}"),
"content": " 2"
}]
}),
};
let err = execute_agent_tool_call(&mut isolated.context, &edit_call)
.await
.expect_err("SCOPE-owned source edit should be rejected");
assert!(
err.to_string()
.contains("edit_file is forbidden for SCOPE-owned source files"),
"unexpected error: {err}"
);
assert_eq!(
std::fs::read_to_string(root.join("lib.rs")).expect("read rust fixture"),
"pub fn value() -> i32 {\n 1\n}\n"
);
drop(isolated);
}
#[tokio::test]
async fn project_scope_allows_file_tool_for_non_scope_file() {
let mut isolated = IsolatedTestContext::new().await;
let root = isolated.context.execution_cwd.clone();
std::fs::write(root.join("README.md"), "old\n").expect("write markdown fixture");
let open_call = AgentToolCall {
id: "call_open".to_string(),
name: "coding__open_project".to_string(),
arguments: json!({
"project_root": root,
}),
};
execute_agent_tool_call(&mut isolated.context, &open_call)
.await
.expect("open project");
let hash = scope_engine::patch::line_hash("old");
let edit_call = AgentToolCall {
id: "call_edit".to_string(),
name: "edit_file".to_string(),
arguments: json!({
"edits": [{
"path": "README.md",
"op": "replace",
"start": format!("1#{hash}"),
"end": format!("1#{hash}"),
"content": "new"
}]
}),
};
execute_agent_tool_call(&mut isolated.context, &edit_call)
.await
.expect("non-SCOPE edit should be allowed");
assert_eq!(
std::fs::read_to_string(root.join("README.md")).expect("read markdown fixture"),
"new\n"
);
drop(isolated);
}
#[tokio::test]
async fn coding_edit_code_falls_back_for_non_scope_file() {
let mut isolated = IsolatedTestContext::new().await;
let root = isolated.context.execution_cwd.clone();
std::fs::write(root.join("README.md"), "old\n").expect("write markdown fixture");
let open_call = AgentToolCall {
id: "call_open".to_string(),
name: "coding__open_project".to_string(),
arguments: json!({
"project_root": root,
}),
};
execute_agent_tool_call(&mut isolated.context, &open_call)
.await
.expect("open project");
let hash = scope_engine::patch::line_hash("old");
let edit_call = AgentToolCall {
id: "call_edit".to_string(),
name: "coding__edit_code".to_string(),
arguments: json!({
"edits": [{
"path": "README.md",
"op": "replace",
"start": format!("1#{hash}"),
"end": format!("1#{hash}"),
"content": "new"
}]
}),
};
let result = execute_agent_tool_call(&mut isolated.context, &edit_call)
.await
.expect("unsupported file should fall back to a plain edit");
assert_eq!(result.payload["propagation_results"], json!([]));
assert_eq!(
std::fs::read_to_string(root.join("README.md")).expect("read markdown fixture"),
"new\n"
);
drop(isolated);
}
#[tokio::test]
async fn generated_get_state_tool_reads_app_state() {
let mut isolated = IsolatedTestContext::new().await;
let call = AgentToolCall {
id: "call_state".to_string(),
name: "terminal__get_state".to_string(),
arguments: json!({}),
};
let result = execute_agent_tool_call(&mut isolated.context, &call)
.await
.unwrap();
drop(isolated);
assert!(result.model_content().contains("app=terminal"));
assert_eq!(result.payload["app"], "terminal");
assert!(result.payload.get("usage").is_none());
assert!(result.payload.get("docs").is_none());
assert!(matches!(
result.activity_event,
Some(SessionActivityEvent::GenericApp(_))
));
}
#[tokio::test]
async fn app_tools_do_not_expose_unscoped_aliases() {
let isolated = IsolatedTestContext::new().await;
let names = build_runtime_tool_specs(&isolated.context)
.into_iter()
.map(|tool| tool.name)
.collect::<HashSet<_>>();
drop(isolated);
assert!(names.contains("read_file"));
assert!(names.contains("edit_file"));
assert!(names.contains("browser__browser_open_page"));
assert!(names.contains("coding__open_project"));
assert!(names.contains("terminal__terminal_exec"));
assert!(!names.contains("apply_patch"));
assert!(!names.contains("terminal_exec"));
assert!(!names.contains("open_project"));
assert!(!names.contains("coding__coding_open_project"));
assert!(!names.contains("coding_open_project"));
assert!(!names.contains("browser_open_page"));
}
}