magi-code 0.96.1

Repository-aware CLI coding agent for terminal work
Documentation
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;

/// Nested structured results are assessed once, by `protected_result`, over the
/// encoded projection the script receives; text-only results use the lifecycle assessment.
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)
        })
    }
}