use super::{
AgentOutputSink,
session_persistence::SessionPersistence,
subdir_instructions::SubdirInstructionState,
tool_lifecycle::{self, AgentDispatchOutcome, ToolExecutionOutcome},
};
use crate::{
cancellation::AgentCancellation,
code_mode::{
catalog::{Catalog, ToolCatalogSurface, effective_catalog},
engine::{CodeModeEngine, QuickJsEngine},
},
hooks::{HookPolicyError, HookRuntime},
output::{ActivityEvent, ActivityId, ActivitySender, OutputEvent, ToolDispatchContext},
providers::{ChatMessage, ProviderConversationItem, ToolCall},
sessions::SessionEventKind,
tools::{ToolResult, ToolRuntime, dispatch::ToolDispatchOutcome},
};
use serde_json::{Value, json};
use std::time::Instant;
pub(super) fn assesses_structured_content(
persistence: &SessionPersistence<'_>,
result: &ToolResult,
) -> bool {
persistence.code_mode_parent.is_some() && result.structured_content().is_some()
}
pub(super) fn protected_result(
tools: Option<&ToolRuntime>,
persistence: &SessionPersistence<'_>,
result: &ToolResult,
provider_output: &str,
protection_context: &Value,
cancellation: &AgentCancellation,
memory: &std::sync::Arc<crate::code_mode::engine::MemoryBudget>,
) -> anyhow::Result<Value> {
let content_charge = result
.structured_content()
.map(|value| crate::code_mode::engine::value_reservation(value, cancellation))
.transpose()?
.unwrap_or_default();
let _projection = memory.reserve(
content_charge
.saturating_mul(2)
.saturating_add(provider_output.len().saturating_mul(4))
.saturating_add(4096),
)?;
let structured = result.structured_result(provider_output);
let limit = tools.map_or(1024 * 1024, |tools| {
tools.code_mode.limits.max_inner_result_bytes
});
let size = crate::code_mode::engine::serialized_size(&structured, limit, cancellation)
.map_err(|error| {
if error
.downcast_ref::<crate::code_mode::engine::ResourceLimitExceeded>()
.is_some()
{
anyhow::Error::from(crate::code_mode::engine::ResourceLimitExceeded(
"inner result bytes",
))
.context(format!("Protected SDK result exceeded {limit} bytes"))
} else {
error
}
})?;
let _encoded = memory.reserve(size)?;
let structured_assessment =
match tools.filter(|_| assesses_structured_content(persistence, result)) {
Some(tools) => {
let encoded = serde_json::to_string(&structured)?;
tools
.assess_tool_result_for_prompt_injection(
&result.tool_name,
&encoded,
protection_context,
)
.map(|(assessment, enforced)| {
(
assessment.metadata(enforced),
assessment.provider_projection(&encoded, enforced),
)
})
}
None => None,
};
cancellation.check()?;
if let Some((metadata, _)) = &structured_assessment {
persistence.record_required(
SessionEventKind::CodeModeAssessment,
json!({"tool_name": result.tool_name, "prompt_injection_protection": metadata}),
)?;
}
let (assessment, withheld) = structured_assessment.unwrap_or_else(|| {
(
result.metadata["prompt_injection_protection"].clone(),
provider_output.to_owned(),
)
});
let action = assessment["action"].as_str().unwrap_or("allow");
if assessment["enforced"] == true && action != "allow" {
persistence.record_required(SessionEventKind::CodeModeWarning, json!({"warning":format!("{}: prompt-injection protection {action}; embedded instructions are untrusted data", result.tool_name)}))?;
if matches!(action, "quarantine" | "escalate") {
return Ok(
json!({"success":false,"error":{"code":"protected_content","message":withheld}}),
);
}
}
Ok(structured)
}
pub(super) fn execute(
tools: &ToolRuntime,
hooks: Option<&HookRuntime>,
persistence: &SessionPersistence<'_>,
sink: &mut Option<&mut dyn AgentOutputSink>,
mut instructions: Option<&mut SubdirInstructionState>,
call: &ToolCall,
context: &ToolDispatchContext,
) -> AgentDispatchOutcome {
let started = Instant::now();
let cancellation = &context.cancellation;
let mut nested = persistence.nested(&call.id);
let mut nested_sink = NestedSink {
sink,
parent: context.parent_activity_id.clone(),
changed_paths: Default::default(),
};
let catalog = Catalog(effective_catalog(
tools,
ToolCatalogSurface::CodeModeInternal,
));
let mut fatal = None;
let execution = (|| -> anyhow::Result<Value> {
anyhow::ensure!(
crate::code_mode::catalog::code_mode_available(tools),
"Code Mode is disabled"
);
anyhow::ensure!(
persistence.code_mode_parent.is_none(),
"Recursive Code Mode is forbidden"
);
anyhow::ensure!(
call.arguments.as_object().is_some_and(|args| args
.keys()
.all(|key| matches!(key.as_str(), "code" | "intent"))),
"code_mode accepts code and optional intent"
);
anyhow::ensure!(
call.arguments.get("intent").is_none_or(|intent| {
intent.as_str().is_some_and(|text| !text.trim().is_empty())
}),
"code_mode intent must be a nonblank string"
);
let source = call.arguments["code"]
.as_str()
.ok_or_else(|| anyhow::anyhow!("code_mode requires JavaScript code"))?;
let internal_tools = tools.for_code_mode();
QuickJsEngine.execute(source, &tools.code_mode.limits, cancellation, &mut |request| {
cancellation.check()?;
if request.operation == "budget" {
return nested.code_mode_budget();
}
if request.operation != "call" {
return Ok(catalog.discover(&request.operation, &request.name));
}
let canonical_name = catalog.canonical_name(&request.name);
nested.code_mode_sequence += 1;
let sequence = nested.code_mode_sequence;
let available = canonical_name.is_some();
let inner = ToolCall { id: format!("{}:inner:{sequence}", call.id), name: canonical_name.unwrap_or(request.name), arguments: request.arguments };
let run = (|| -> anyhow::Result<Value> {
nested.record_required(SessionEventKind::ToolCall, json!({"id":inner.id,"name":inner.name,"arguments":inner.arguments}))?;
let mut inner_sink: Option<&mut dyn AgentOutputSink> = Some(&mut nested_sink);
let activity = tool_lifecycle::emit_tool_started(&mut inner_sink, 0, sequence as usize, &inner)?;
if !available {
let result = ToolResult {
tool_name: inner.name.clone(),
success: false,
content: "Capability unavailable in Code Mode".into(),
metadata: json!({}),
display: Default::default(),
};
nested.record_required(SessionEventKind::ToolResult, json!({"call_id":inner.id,"result":result}))?;
if let Some(sink) = inner_sink.as_deref_mut() {
sink.output_event(OutputEvent::ToolResult {
call: Box::new(inner.clone()), result: Box::new(result.clone()),
changed_paths: Vec::new(),
summary: Box::new(crate::output::tool_display_summary(&inner, &result)),
})?;
}
tool_lifecycle::emit_tool_finished(&mut inner_sink, activity.activity_id, &inner, &result)?;
return Ok(json!({"success":false,"error":{"code":"unavailable_tool","message":result.content}}));
}
let mut inner_context = context.clone();
inner_context.parent_activity_id = Some(activity.activity_id.clone());
let result = tool_lifecycle::execute_tool_call(Some(&internal_tools), hooks, &mut nested, &mut inner_sink,
instructions.as_deref_mut(), inner.clone(), inner_context);
let outcome = match result {
Ok(outcome) => outcome,
Err(error) => {
if let Some(sink) = inner_sink.as_deref_mut() {
sink.activity_event(ActivityEvent::Finished {
id: activity.activity_id,
status: if crate::cancellation::is_run_canceled(&error) {
crate::output::ActivityStatus::Canceled
} else { crate::output::ActivityStatus::Failed },
metadata: None,
})?;
}
return Err(error);
}
};
match outcome {
ToolExecutionOutcome::Dispatched { tool_result, provider_result, after_hook_failure, after_hook_context_items, .. } => {
tool_lifecycle::emit_tool_finished(&mut inner_sink, activity.activity_id, &inner, &tool_result)?;
for item in &after_hook_context_items {
tool_lifecycle::record_provider_context_item(&mut nested, &mut inner_sink, item)?;
}
if let Some(diagnostic) = after_hook_failure { return Err(HookPolicyError::new(diagnostic).into()); }
cancellation.check()?;
protected_result(Some(tools), &nested, &tool_result,
&provider_result.output, &json!({
"user_context": context.request_context,
"tool_arguments": crate::typesafe::evidence::excerpt(&inner.arguments.to_string(), 4096),
}), cancellation, &request.memory)
}
ToolExecutionOutcome::BeforeHookFailed { diagnostic, tool_result } => {
tool_lifecycle::emit_tool_finished(&mut inner_sink, activity.activity_id, &inner, &tool_result)?;
Err(HookPolicyError::new(diagnostic).into())
}
}
})();
match run {
Ok(value) => Ok(value),
Err(error) if error.downcast_ref::<crate::code_mode::engine::ResourceLimitExceeded>().is_some() => Err(error),
Err(error) => {
fatal = Some(error);
anyhow::bail!("Code Mode guarded host execution failed")
}
}
})
})();
let audit = nested.code_mode_audit.into_inner();
let mut summary = audit.summary();
summary["elapsed_ms"] = (started.elapsed().as_millis() as u64).into();
let success = execution.is_ok();
let content = match execution {
Ok(value) => json!({"value":value,"status":"completed","execution":summary}),
Err(error) => {
let recovery = error.downcast_ref::<crate::code_mode::engine::ResourceLimitExceeded>()
.map_or("Inspect error and execution evidence; re-read affected files before retrying. Completed effects remain applied; never replay the whole batch automatically.", |limit| limit.recovery());
json!({"status":"terminated","error":crate::output::redact_sensitive_text(&format!("{error:#}")),"recovery":recovery,"execution":summary})
}
};
if context.cancellation.is_canceled() {
fatal = Some(crate::cancellation::AgentRunCanceled.into());
}
AgentDispatchOutcome {
dispatch: ToolDispatchOutcome {
result: ToolResult {
tool_name: crate::code_mode::TOOL_NAME.into(),
success,
content: content.to_string(),
metadata: json!({"code_mode_execution":summary}),
display: Default::default(),
},
touched_paths: Vec::new(),
changed_paths: nested_sink.changed_paths.into_iter().collect(),
},
context_payloads: audit.contexts,
fatal,
}
}
pub(super) fn context_item(payload: &Value) -> Option<ProviderConversationItem> {
payload["content"]
.as_str()
.map(|text| ProviderConversationItem::Message(ChatMessage::user(text)))
}
struct NestedSink<'a, 'sink> {
sink: &'a mut Option<&'sink mut dyn AgentOutputSink>,
parent: Option<ActivityId>,
changed_paths: std::collections::BTreeSet<std::path::PathBuf>,
}
impl AgentOutputSink for NestedSink<'_, '_> {
fn assistant_delta(&mut self, _: &str) -> anyhow::Result<()> {
Ok(())
}
fn tool_block(&mut self, _: &str) -> anyhow::Result<()> {
Ok(())
}
fn current_parent_activity_id(&self) -> Option<ActivityId> {
self.parent.clone()
}
fn activity_sender(&self) -> Option<ActivitySender> {
self.sink
.as_deref()
.and_then(AgentOutputSink::activity_sender)
}
fn activity_event(&mut self, event: ActivityEvent) -> anyhow::Result<()> {
if let Some(sink) = self.sink.as_deref_mut() {
sink.activity_event(event)?;
}
Ok(())
}
fn output_event(&mut self, event: OutputEvent) -> anyhow::Result<()> {
match event {
OutputEvent::ToolStarted { call, label } => {
self.activity_event(ActivityEvent::ToolStartedDetail {
id: ActivityId::new(call.id.clone()),
detail: crate::output::tool_display::tool_activity_detail_pending(
&call,
label,
crate::output::pending_activity_status(&call),
),
})
}
OutputEvent::ToolResult {
call,
result,
changed_paths,
summary,
} => {
self.changed_paths.extend(changed_paths);
self.activity_event(ActivityEvent::ToolResultDetail {
id: ActivityId::new(call.id.clone()),
detail: crate::output::tool_display::tool_activity_detail(
&call,
&result,
summary.label.clone(),
if result.success {
crate::output::ActivityStatus::Success
} else {
crate::output::ActivityStatus::Failed
},
),
})
}
OutputEvent::Diagnostic { .. } | OutputEvent::HookDiagnostic { .. } => {
if let Some(sink) = self.sink.as_deref_mut() {
sink.output_event(event)?;
}
Ok(())
}
_ => Ok(()),
}
}
fn request_bash_approval(
&mut self,
request: crate::protection::bash::BashApprovalRequest,
cancellation: &AgentCancellation,
) -> anyhow::Result<bool> {
self.sink.as_deref_mut().map_or(Ok(false), |sink| {
sink.request_bash_approval(request, cancellation)
})
}
}