use std::cell::Cell;
use std::future::Future;
use std::path::{Path, PathBuf};
use std::pin::Pin;
use std::sync::Arc;
use serde_json::json;
use tracing::info;
use crate::approve::{ApproveAll, Approver, Decision, Request};
use crate::containment::{Containment, Draw, Ledger};
use crate::context::{
assemble, bound, entry_cap_chars, Assembly, Ledger as ContextLedger, ObsKind, Observation,
};
use crate::contract::TaskContract;
use crate::error::{Error, Result};
use crate::mcp::McpSession;
use crate::net::{self, NetGuard};
use crate::observe::{EventKind, Ignore, Observer, RunEvent};
use crate::policy::{Act, Effect, Policy, Rule};
use crate::provider::{CompletionRequest, CompletionResponse, Provider, ToolCall, ToolSpec};
use crate::resilience::{Progress, Progressing};
use crate::skills::Skills;
use crate::state::PolicyEvent;
use crate::state::{AgentEvent, ContextEvent, RunStatus, StepRecord, Store};
use crate::tools::{
FsTool, Toolbox, Workspace, FIND_TOOL, GREP_TOOL, READ_FILE_TOOL, READ_SKILL_TOOL,
REMEMBER_TOOL, WRITE_FILE_TOOL,
};
use crate::verify::{ExecGuard, Verification};
pub const SPAWN_TOOL: &str = "spawn_agent";
const OBS_GREP_CAP: usize = 50;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RunOutcome {
Success { steps: u32 },
StepCapReached { steps: u32 },
TimeBudgetExceeded { steps: u32 },
CostBudgetExceeded { steps: u32 },
Denied { steps: u32 },
AwaitingApproval { request_id: i64, steps: u32 },
Stalled { steps: u32 },
Escalated { steps: u32, retryable: bool },
BudgetCeilingReached { steps: u32 },
Refused { steps: u32 },
Cancelled { steps: u32 },
}
#[derive(Debug, Clone)]
pub struct RunResult {
pub outcome: RunOutcome,
pub run_id: i64,
pub remembered: Vec<Rule>,
}
impl RunResult {
pub fn summary(&self, store: &Store) -> Result<Option<crate::RunSummary>> {
store.run_summary(self.run_id)
}
}
impl RunResult {
fn new(outcome: RunOutcome, run_id: i64) -> Self {
Self {
outcome,
run_id,
remembered: Vec::new(),
}
}
fn with_remembered(mut self, remembered: Vec<Rule>) -> Self {
self.remembered = remembered;
self
}
}
pub async fn run<P: Provider>(
contract: &TaskContract,
provider: &P,
store: &Store,
) -> Result<RunResult> {
run_observed(contract, provider, store, &Ignore).await
}
pub async fn run_observed<P: Provider>(
contract: &TaskContract,
provider: &P,
store: &Store,
observer: &dyn Observer,
) -> Result<RunResult> {
run_with_observed(
contract,
provider,
store,
&Policy::permissive(),
&ApproveAll,
observer,
)
.await
}
pub async fn run_with<P: Provider>(
contract: &TaskContract,
provider: &P,
store: &Store,
policy: &Policy,
approver: &dyn Approver,
) -> Result<RunResult> {
run_with_observed(contract, provider, store, policy, approver, &Ignore).await
}
pub async fn run_with_observed<P: Provider>(
contract: &TaskContract,
provider: &P,
store: &Store,
policy: &Policy,
approver: &dyn Approver,
observer: &dyn Observer,
) -> Result<RunResult> {
contract.tools.validate()?;
let skills = contract.discover_skills()?;
let file_str = contract.file.display().to_string();
let run_id = store.start_run(&contract.goal, &file_str)?;
store.set_provider(run_id, provider.name())?;
store.record_run_policy(run_id, policy)?;
let watch = &Watch::new(observer);
watch.emit(RunEvent::new(
run_id,
0,
EventKind::Started {
goal: contract.goal.clone(),
provider: provider.name().to_string(),
},
));
let caller_enforces = !policy.is_permissive();
let policy = &match authorize_provider(provider, policy, store, run_id, approver, watch).await?
{
ProviderAccess::Granted(p) => p,
ProviderAccess::Pending(request_id) => {
return Ok(RunResult::new(
RunOutcome::AwaitingApproval {
request_id,
steps: 0,
},
run_id,
))
}
};
match contract.root.clone() {
Some(root) => {
let mcp = McpSession::connect(&contract.mcp, policy, store, run_id, watch).await?;
let result = run_workspace_from(
contract, provider, store, run_id, &root, 1, policy, approver, &mcp, &skills, watch,
)
.await;
mcp.shutdown(store, run_id, watch).await;
result
}
None if caller_enforces => Err(crate::error::Error::Config(
"a permission policy requires workspace mode — build the contract \
with TaskContract::workspace(goal, root, verify). Single-file \
contracts are not policy-enforced in 0.4.0."
.into(),
)),
None => run_from(contract, provider, store, run_id, 1, watch).await,
}
}
pub async fn resume<P: Provider>(
contract: &TaskContract,
provider: &P,
store: &Store,
run_id: i64,
) -> Result<RunResult> {
resume_observed(contract, provider, store, run_id, &Ignore).await
}
pub async fn resume_observed<P: Provider>(
contract: &TaskContract,
provider: &P,
store: &Store,
run_id: i64,
observer: &dyn Observer,
) -> Result<RunResult> {
store.check_resumable(run_id)?;
if let Some(o) = finished_outcome(store, run_id)? {
return Ok(RunResult::new(o, run_id));
}
if let Some(recorded) = store.run_policy(run_id)? {
if !recorded.is_permissive() {
return Err(crate::error::Error::Resume {
reason: format!(
"run {run_id} was started under a permission policy; resume it with \
resume_with (or resume_with_observed), supplying that policy — resuming \
here would drop the boundary the run was executing under"
),
});
}
}
resume_with_observed(
contract,
provider,
store,
run_id,
&Policy::permissive(),
&ApproveAll,
observer,
)
.await
}
pub async fn resume_with<P: Provider>(
contract: &TaskContract,
provider: &P,
store: &Store,
run_id: i64,
policy: &Policy,
approver: &dyn Approver,
) -> Result<RunResult> {
resume_with_observed(contract, provider, store, run_id, policy, approver, &Ignore).await
}
#[allow(clippy::too_many_arguments)]
pub async fn resume_with_observed<P: Provider>(
contract: &TaskContract,
provider: &P,
store: &Store,
run_id: i64,
policy: &Policy,
approver: &dyn Approver,
observer: &dyn Observer,
) -> Result<RunResult> {
contract.tools.validate()?;
let skills = contract.discover_skills()?;
store.check_resumable(run_id)?;
if let Some(o) = finished_outcome(store, run_id)? {
return Ok(RunResult::new(o, run_id));
}
let caller_enforces = !policy.is_permissive();
let start_step = record_resume_markers(store, run_id)?;
store.set_provider(run_id, provider.name())?;
store.record_run_policy(run_id, policy)?;
let watch = &Watch::new(observer);
watch.emit(RunEvent::new(
run_id,
start_step.saturating_sub(1),
EventKind::Started {
goal: contract.goal.clone(),
provider: provider.name().to_string(),
},
));
match contract.root.clone() {
Some(root) => {
let policy = &match authorize_provider(provider, policy, store, run_id, approver, watch)
.await?
{
ProviderAccess::Granted(p) => p,
ProviderAccess::Pending(request_id) => {
return Ok(RunResult::new(
RunOutcome::AwaitingApproval {
request_id,
steps: start_step.saturating_sub(1),
},
run_id,
))
}
};
let mcp = McpSession::connect(&contract.mcp, policy, store, run_id, watch).await?;
let result = run_workspace_from(
contract, provider, store, run_id, &root, start_step, policy, approver, &mcp,
&skills, watch,
)
.await;
mcp.shutdown(store, run_id, watch).await;
result
}
None if caller_enforces => Err(crate::error::Error::Config(
"a permission policy requires workspace mode — build the contract \
with TaskContract::workspace(goal, root, verify). Single-file \
contracts are not policy-enforced."
.into(),
)),
None => run_from(contract, provider, store, run_id, start_step, watch).await,
}
}
#[allow(clippy::too_many_arguments)]
pub async fn resume_with_decision<P: Provider>(
contract: &TaskContract,
provider: &P,
store: &Store,
run_id: i64,
request_id: i64,
decision: Decision,
policy: &Policy,
approver: &dyn Approver,
) -> Result<RunResult> {
resume_with_decision_observed(
contract, provider, store, run_id, request_id, decision, policy, approver, &Ignore,
)
.await
}
#[allow(clippy::too_many_arguments)]
pub async fn resume_with_decision_observed<P: Provider>(
contract: &TaskContract,
provider: &P,
store: &Store,
run_id: i64,
request_id: i64,
decision: Decision,
policy: &Policy,
approver: &dyn Approver,
observer: &dyn Observer,
) -> Result<RunResult> {
contract.tools.validate()?;
let skills = contract.discover_skills()?;
let pending = store
.pending(request_id)?
.ok_or_else(|| crate::error::Error::Config(format!("no pending request {request_id}")))?;
if pending.run_id != run_id {
return Err(crate::error::Error::Config(format!(
"request {request_id} belongs to run {}, not {run_id}",
pending.run_id
)));
}
let root = contract.root.clone().ok_or_else(|| {
crate::error::Error::Config("resume_with_decision needs a workspace".into())
})?;
let step = pending.step;
let watch = &Watch::new(observer);
watch.emit(RunEvent::new(
run_id,
step,
EventKind::Started {
goal: contract.goal.clone(),
provider: provider.name().to_string(),
},
));
match decision {
Decision::Defer => Ok(RunResult::new(
RunOutcome::AwaitingApproval {
request_id,
steps: step,
},
run_id,
)),
Decision::Deny { reason } => {
store.resolve_pending(request_id, "deny")?;
store.record_event(
run_id,
&PolicyEvent::decision(
step,
&pending.act,
&pending.target,
"deny",
format!("resumed:{request_id}"),
),
)?;
info!(run_id, request_id, %reason, "deferred action denied");
finish(store, watch, run_id, 0, step, "denied")?;
Ok(RunResult::new(RunOutcome::Denied { steps: step }, run_id))
}
Decision::Approve { ref remember, .. } if pending.act == "net" => {
let effective = policy
.clone()
.merge(net::provider_layer(&pending.target))
.merge(remembered_layer(remember));
store.resolve_pending(request_id, "approve")?;
store.record_event(
run_id,
&PolicyEvent::decision(
step,
"net",
&pending.target,
"approve",
format!("resumed:{request_id}"),
),
)?;
let remember = remember.clone();
let mcp = McpSession::connect(&contract.mcp, &effective, store, run_id, watch).await?;
let result = run_workspace_from(
contract,
provider,
store,
run_id,
&root,
step + 1,
&effective,
approver,
&mcp,
&skills,
watch,
)
.await;
mcp.shutdown(store, run_id, watch).await;
result.map(|r| r.with_remembered(remember))
}
Decision::Approve { modified, remember } => {
let target = modified
.as_ref()
.map(|m| m.target.clone())
.unwrap_or_else(|| pending.target.clone());
let content = modified
.as_ref()
.and_then(|m| m.content.clone())
.or_else(|| pending.content.clone());
let mut effective = policy.clone();
if !remember.is_empty() {
let mut layer = Policy::permissive().layer("remembered");
for r in &remember {
layer = layer.rule(r.act, r.effect, r.pattern.clone());
}
effective = effective.merge(layer);
}
let ws = Workspace::with_policy(&root, effective.clone());
let act = if pending.act == "read" {
Act::Read
} else {
Act::Write
};
let recheck = ws.check_path(act, &target);
if recheck.effect == Effect::Deny {
let mut ev = PolicyEvent::refusal(step, &pending.act, &target);
ev.rule = recheck.rule.clone();
ev.layer = recheck.layer.clone();
store.record_event(run_id, &ev)?;
refused(watch, run_id, 0, &ev);
store.resolve_pending(request_id, "deny")?;
finish(store, watch, run_id, 0, step, "denied")?;
return Ok(RunResult::new(RunOutcome::Denied { steps: step }, run_id));
}
if act == Act::Write {
ws.write_file(&target, content.as_deref().unwrap_or_default())?;
}
store.resolve_pending(request_id, "approve")?;
let mut ev = PolicyEvent::decision(
step,
&pending.act,
&pending.target,
"approve",
format!("resumed:{request_id}"),
);
if target != pending.target {
ev = ev.with_performed(&target);
}
store.record_event(run_id, &ev)?;
let mcp = McpSession::connect(&contract.mcp, &effective, store, run_id, watch).await?;
let result = run_workspace_from(
contract,
provider,
store,
run_id,
&root,
step + 1,
&effective,
approver,
&mcp,
&skills,
watch,
)
.await;
mcp.shutdown(store, run_id, watch).await;
result.map(|r| r.with_remembered(remember))
}
}
}
#[allow(clippy::too_many_arguments)]
pub async fn resume_tree_with_decision<P: Provider>(
contract: &TaskContract,
provider: &P,
store: &Store,
run_id: i64,
request_id: i64,
decision: Decision,
policy: &Policy,
approver: &dyn Approver,
containment: &Containment,
) -> Result<RunResult> {
resume_tree_with_decision_observed(
contract,
provider,
store,
run_id,
request_id,
decision,
policy,
approver,
containment,
&Ignore,
)
.await
}
#[allow(clippy::too_many_arguments)]
pub async fn resume_tree_with_decision_observed<P: Provider>(
contract: &TaskContract,
provider: &P,
store: &Store,
run_id: i64,
request_id: i64,
decision: Decision,
policy: &Policy,
approver: &dyn Approver,
containment: &Containment,
observer: &dyn Observer,
) -> Result<RunResult> {
store.check_resumable(run_id)?;
contract.tools.validate()?;
let skills = contract.discover_skills()?;
let pending = store
.pending(request_id)?
.ok_or_else(|| crate::error::Error::Config(format!("no pending request {request_id}")))?;
if !store.tree_run_ids(run_id)?.contains(&pending.run_id) {
return Err(crate::error::Error::Config(format!(
"request {request_id} belongs to run {}, which is not in the tree rooted at {run_id}",
pending.run_id
)));
}
let root = contract.root.clone().ok_or_else(|| {
crate::error::Error::Config("resume_tree_with_decision needs a workspace".into())
})?;
let step = pending.step;
let watch = &Watch::new(observer);
watch.emit(RunEvent::new(
run_id,
step,
EventKind::Started {
goal: contract.goal.clone(),
provider: provider.name().to_string(),
},
));
match decision {
Decision::Defer => Ok(RunResult::new(
RunOutcome::AwaitingApproval {
request_id,
steps: step,
},
run_id,
)),
Decision::Deny { reason } => {
store.resolve_pending(request_id, "deny")?;
store.record_event(
pending.run_id,
&PolicyEvent::decision(
step,
&pending.act,
&pending.target,
"deny",
format!("resumed:{request_id}"),
),
)?;
info!(run_id, request_id, %reason, "deferred tree action denied; tree stops");
finish(store, watch, run_id, 0, step, "denied")?;
Ok(RunResult::new(RunOutcome::Denied { steps: step }, run_id))
}
Decision::Approve { ref remember, .. } if pending.act == "net" => {
let effective = policy
.clone()
.merge(net::provider_layer(&pending.target))
.merge(remembered_layer(remember));
store.resolve_pending(request_id, "approve")?;
store.record_event(
pending.run_id,
&PolicyEvent::decision(
step,
"net",
&pending.target,
"approve",
format!("resumed:{request_id}"),
),
)?;
let ledger = Arc::new(Ledger::from_state(
containment,
store.spent_tokens_tree(run_id)?,
store.agent_count_tree(run_id)?,
));
let start_step = record_resume_markers(store, run_id)?;
store.set_provider(run_id, provider.name())?;
let mcp = McpSession::connect(&contract.mcp, &effective, store, run_id, watch).await?;
let tree = Tree {
mcp: &mcp,
tools: &contract.tools,
skills: &skills,
provider,
store,
approver,
watch,
ledger,
containment,
root,
root_run_id: run_id,
};
let outcome = run_agent(&tree, contract, run_id, 0, &effective, start_step).await;
mcp.shutdown(store, run_id, watch).await;
Ok(RunResult::new(outcome?, run_id).with_remembered(remember.clone()))
}
Decision::Approve { modified, remember } => {
let target = modified
.as_ref()
.map(|m| m.target.clone())
.unwrap_or_else(|| pending.target.clone());
let content = modified
.as_ref()
.and_then(|m| m.content.clone())
.or_else(|| pending.content.clone());
let ws = Workspace::with_policy(&root, policy.clone());
let act = if pending.act == "read" {
Act::Read
} else {
Act::Write
};
if ws.check_path(act, &target).effect == Effect::Deny {
store.resolve_pending(request_id, "deny")?;
finish(store, watch, run_id, 0, step, "denied")?;
return Ok(RunResult::new(RunOutcome::Denied { steps: step }, run_id));
}
if act == Act::Write {
ws.write_file(&target, content.as_deref().unwrap_or_default())?;
}
store.resolve_pending(request_id, "approve")?;
store.record_event(
pending.run_id,
&PolicyEvent::decision(
step,
&pending.act,
&pending.target,
"approve",
format!("resumed:{request_id}"),
),
)?;
let mut effective = policy.clone();
if !remember.is_empty() {
let mut layer = Policy::permissive().layer("remembered");
for r in &remember {
layer = layer.rule(r.act, r.effect, r.pattern.clone());
}
effective = effective.merge(layer);
}
let ledger = Arc::new(Ledger::from_state(
containment,
store.spent_tokens_tree(run_id)?,
store.agent_count_tree(run_id)?,
));
let start_step = record_resume_markers(store, run_id)?;
store.set_provider(run_id, provider.name())?;
let mcp = McpSession::connect(&contract.mcp, &effective, store, run_id, watch).await?;
let tree = Tree {
mcp: &mcp,
tools: &contract.tools,
skills: &skills,
provider,
store,
approver,
watch,
ledger,
containment,
root,
root_run_id: run_id,
};
let outcome = run_agent(&tree, contract, run_id, 0, &effective, start_step).await;
mcp.shutdown(store, run_id, watch).await;
Ok(RunResult::new(outcome?, run_id).with_remembered(remember))
}
}
}
fn record_resume_markers(store: &Store, run_id: i64) -> Result<u32> {
let last = store.last_step(run_id)?;
let start_step = last + 1;
store.record_checkpoint_event(&crate::state::CheckpointEvent::resume(
run_id,
start_step,
format!("resuming at step {start_step}, {last} committed step(s) skipped"),
))?;
for s in 1..=last {
store.record_checkpoint_event(&crate::state::CheckpointEvent::skipped(run_id, s))?;
}
Ok(start_step)
}
fn restore_ledger(store: &Store, run_id: i64) -> Result<(ContextLedger, usize)> {
let mut ledger = ContextLedger::new();
for obs in store.observations(run_id)? {
ledger.push(obs);
}
let written = ledger.len();
Ok((ledger, written))
}
fn persist_ledger(
store: &Store,
run_id: i64,
ledger: &ContextLedger,
written: usize,
) -> Result<usize> {
store.record_observations(run_id, &ledger.entries()[written..])?;
Ok(ledger.len())
}
fn finished_outcome(store: &Store, run_id: i64) -> Result<Option<RunOutcome>> {
if store.run_status(run_id)? != Some(RunStatus::Completed) {
return Ok(None);
}
terminal_outcome(store, run_id)
}
fn terminal_outcome(store: &Store, run_id: i64) -> Result<Option<RunOutcome>> {
let last = store.last_step(run_id)?;
Ok(store.outcome(run_id)?.and_then(|o| match o.as_str() {
"success" => Some(RunOutcome::Success { steps: last }),
"denied" => Some(RunOutcome::Denied { steps: last }),
"budget_ceiling_reached" => Some(RunOutcome::BudgetCeilingReached { steps: last }),
"stalled" => Some(RunOutcome::Stalled { steps: last }),
"escalated_retryable" => Some(RunOutcome::Escalated {
steps: last,
retryable: true,
}),
"escalated_terminal" | "escalated" => Some(RunOutcome::Escalated {
steps: last,
retryable: false,
}),
"refused" => Some(RunOutcome::Refused { steps: last }),
"cancelled" => Some(RunOutcome::Cancelled { steps: last }),
_ => None,
}))
}
pub(crate) struct Watch<'a> {
observer: &'a dyn Observer,
cancelled: Cell<bool>,
}
impl<'a> Watch<'a> {
fn new(observer: &'a dyn Observer) -> Self {
Self {
observer,
cancelled: Cell::new(false),
}
}
pub(crate) fn emit(&self, event: RunEvent) {
if self.observer.event(&event).is_cancel() {
self.cancelled.set(true);
}
}
fn cancelled(&self) -> bool {
self.cancelled.get()
}
}
pub(crate) fn refused(watch: &Watch<'_>, run_id: i64, depth: u32, ev: &PolicyEvent) {
watch.emit(RunEvent::at_depth(
run_id,
ev.step,
depth,
EventKind::Refused {
act: ev.act.clone(),
target: ev.target.clone(),
rule: ev.rule.clone(),
layer: ev.layer.clone(),
},
));
}
fn decided(watch: &Watch<'_>, run_id: i64, depth: u32, ev: &PolicyEvent) {
watch.emit(RunEvent::at_depth(
run_id,
ev.step,
depth,
EventKind::ApprovalDecided {
act: ev.act.clone(),
target: ev.target.clone(),
decision: ev.decision.clone().unwrap_or_default(),
},
));
}
fn commit_step(
store: &Store,
watch: &Watch<'_>,
run_id: i64,
depth: u32,
record: StepRecord,
changed: bool,
commit: bool,
) -> Result<()> {
if !commit {
info!(
run_id,
depth,
step = record.step,
"tree paused for a child's approval (step left uncommitted for replay)"
);
return Ok(());
}
store.checkpoint_step(run_id, &record)?;
info!(
run_id,
depth,
step = record.step,
decision = %record.decision,
tokens = record.tokens,
changed,
"step"
);
watch.emit(RunEvent::at_depth(
run_id,
record.step,
depth,
EventKind::Step {
decision: record.decision,
tool_call: record.tool_call,
tokens: record.tokens,
changed,
},
));
Ok(())
}
fn cancelled(
store: &Store,
watch: &Watch<'_>,
run_id: i64,
depth: u32,
steps: u32,
) -> Result<Option<RunOutcome>> {
if !watch.cancelled() {
return Ok(None);
}
finish(store, watch, run_id, depth, steps, "cancelled")?;
info!(run_id, depth, steps, "run cancelled by its observer");
Ok(Some(RunOutcome::Cancelled { steps }))
}
fn finish(
store: &Store,
watch: &Watch<'_>,
run_id: i64,
depth: u32,
steps: u32,
outcome: &str,
) -> Result<()> {
store.finish_run(run_id, outcome)?;
watch.emit(RunEvent::at_depth(
run_id,
steps,
depth,
EventKind::Finished {
outcome: outcome.to_string(),
steps,
tokens: store.spent_tokens(run_id)?,
},
));
Ok(())
}
async fn run_from<P: Provider>(
contract: &TaskContract,
provider: &P,
store: &Store,
run_id: i64,
start_step: u32,
watch: &Watch<'_>,
) -> Result<RunResult> {
let fs = FsTool::new(&contract.file);
let system = system_prompt();
let tool = write_file_tool();
let mut tokens_used: u64 = store.spent_tokens(run_id)?;
let permissive = Policy::permissive();
for step in start_step..=contract.max_steps {
if let Some(o) = cancelled(store, watch, run_id, 0, step - 1)? {
return Ok(RunResult::new(o, run_id));
}
if let Some(max) = contract.max_duration {
if store.elapsed_secs(run_id)? > max.as_secs_f64() {
finish(store, watch, run_id, 0, step - 1, "time_budget_exceeded")?;
return Ok(RunResult::new(
RunOutcome::TimeBudgetExceeded { steps: step - 1 },
run_id,
));
}
}
let current = fs.read().await?;
let user =
user_prompt(
contract,
&bound(
¤t,
entry_cap_chars(contract.context.effective_tokens(
contract.max_tokens.map(|m| m.saturating_sub(tokens_used)),
)),
ObsKind::Read,
),
);
let request = CompletionRequest {
system: system.clone(),
user: user.clone(),
tools: vec![tool.clone()],
};
let response =
complete_with_retry(provider, &request, contract, store, run_id, step, watch, 0)
.await?;
if let Some(served) = provider.last_served() {
store.record_context_event(run_id, &ContextEvent::served(step, served.clone()))?;
watch.emit(RunEvent::new(
run_id,
step,
EventKind::FellBackTo { provider: served },
));
}
let step_tokens = response.usage.map(|u| u.total_tokens).unwrap_or(0);
tokens_used += step_tokens;
let call = response
.tool_calls
.iter()
.find(|c| c.name == WRITE_FILE_TOOL);
let tool_call_json = call.map(|c| c.arguments.to_string()).unwrap_or_default();
let write = call.and_then(|c| c.arguments.get("content").and_then(|v| v.as_str()));
let (decision, result_text) = match write {
Some(content) => {
fs.write(content).await?;
("wrote file", content.to_string())
}
None => ("no tool call", response.text.clone().unwrap_or_default()),
};
commit_step(
store,
watch,
run_id,
0,
StepRecord::new(step, decision, result_text).with_trace(
user,
tool_call_json,
step_tokens,
),
write.is_some(),
true,
)?;
if let Some(max) = contract.max_tokens {
if tokens_used > max {
finish(store, watch, run_id, 0, step, "cost_budget_exceeded")?;
return Ok(RunResult::new(
RunOutcome::CostBudgetExceeded { steps: step },
run_id,
));
}
}
let contents = fs.read().await?;
let guard = ExecGuard::new(&permissive)
.tracing(store, run_id, step)
.watching(watch, 0);
if contract
.verify
.passes_guarded(&contract.file, &contents, &guard)
.await?
{
finish(store, watch, run_id, 0, step, "success")?;
return Ok(RunResult::new(RunOutcome::Success { steps: step }, run_id));
}
}
finish(
store,
watch,
run_id,
0,
contract.max_steps,
"step_cap_reached",
)?;
Ok(RunResult::new(
RunOutcome::StepCapReached {
steps: contract.max_steps,
},
run_id,
))
}
#[allow(clippy::too_many_arguments)]
async fn run_workspace_from<P: Provider>(
contract: &TaskContract,
provider: &P,
store: &Store,
run_id: i64,
root: &Path,
start_step: u32,
policy: &Policy,
approver: &dyn Approver,
mcp: &McpSession,
skills: &Skills,
watch: &Watch<'_>,
) -> Result<RunResult> {
let mut effective = policy.clone();
let mut remembered: Vec<Rule> = Vec::new();
let mut ws = Workspace::with_policy(root, effective.clone());
let mut extra = contract.tools.specs();
extra.extend(mcp.tool_specs());
extra.extend(skill_tool(skills));
let system = with_skill_catalog(with_extra_tools(workspace_system_prompt(), &extra), skills);
let mut tools = workspace_tools();
tools.extend(extra);
let mut tokens_used: u64 = store.spent_tokens(run_id)?;
let (mut ledger, mut written) = restore_ledger(store, run_id)?;
let mut progress = Progress::new();
let mem_key = memory_key(root);
for step in start_step..=contract.max_steps {
if let Some(o) = cancelled(store, watch, run_id, 0, step - 1)? {
return Ok(RunResult::new(o, run_id).with_remembered(remembered));
}
if let Some(max) = contract.max_duration {
if store.elapsed_secs(run_id)? > max.as_secs_f64() {
finish(store, watch, run_id, 0, step - 1, "time_budget_exceeded")?;
return Ok(RunResult::new(
RunOutcome::TimeBudgetExceeded { steps: step - 1 },
run_id,
)
.with_remembered(remembered));
}
}
let budget_tokens = contract
.context
.effective_tokens(contract.max_tokens.map(|m| m.saturating_sub(tokens_used)));
let entry_cap = entry_cap_chars(budget_tokens);
let notes = store.memory_list(&mem_key)?;
let assembled = assemble(
&ledger,
budget_tokens,
¬es,
Assembly {
ws: Some(&ws),
policy: &effective,
store,
run_id,
step,
},
)
.await?;
let user = workspace_user_prompt(contract, &assembled.text);
let request = CompletionRequest {
system: system.clone(),
user: user.clone(),
tools: tools.clone(),
};
let response =
complete_with_retry(provider, &request, contract, store, run_id, step, watch, 0)
.await?;
if let Some(served) = provider.last_served() {
store.record_context_event(run_id, &ContextEvent::served(step, served.clone()))?;
watch.emit(RunEvent::new(
run_id,
step,
EventKind::FellBackTo { provider: served },
));
}
let step_tokens = response.usage.map(|u| u.total_tokens).unwrap_or(0);
tokens_used += step_tokens;
if step_tokens > 0 {
store.record_context_reported(run_id, step, step_tokens)?;
}
let mut decisions: Vec<String> = Vec::new();
let mut calls_json: Vec<String> = Vec::new();
let mut step_changed = false;
if response.tool_calls.is_empty() {
let said = response.text.clone().unwrap_or_default();
ledger.push(Observation::new(
step,
ObsKind::Message,
None,
bound(
&format!("\n[step {step}] (no tool call) {said}\n"),
entry_cap,
ObsKind::Message,
),
));
decisions.push("no tool call".into());
}
let mut paused: Option<i64> = None;
let mut new_rules: Vec<Rule> = Vec::new();
for call in &response.tool_calls {
calls_json.push(format!("{}:{}", call.name, call.arguments));
match dispatch(
&ws,
call,
approver,
store,
run_id,
step,
mcp,
&contract.tools,
skills,
entry_cap,
&mem_key,
watch,
0,
)
.await?
{
Dispatched::Continue {
decision,
obs,
kind,
target,
changed,
remember,
} => {
step_changed |= changed;
ledger.push(Observation::new(step, kind, target, obs));
decisions.push(decision);
new_rules.extend(remember);
}
Dispatched::Pause { request_id } => {
decisions.push(format!("awaiting approval (request {request_id})"));
paused = Some(request_id);
break;
}
}
}
commit_step(
store,
watch,
run_id,
0,
StepRecord::new(step, decisions.join("; "), ledger.text_for_step(step)).with_trace(
user,
calls_json.join(" | "),
step_tokens,
),
step_changed,
true,
)?;
written = persist_ledger(store, run_id, &ledger, written)?;
let signature = calls_json.join(" | ");
match progress.step(contract.stall, step_changed, &signature) {
Progressing::Fine => {}
Progressing::Replan => {
store.record_context_event(
run_id,
&ContextEvent::replan(
step,
format!(
"{} steps without progress; replanning",
contract.stall.window
),
),
)?;
ledger.push(Observation::new(
step,
ObsKind::Message,
None,
bound(
&progress.replan_directive(contract.stall.window, &decisions),
entry_cap,
ObsKind::Message,
),
));
info!(run_id, step, "agent told to change approach");
watch.emit(RunEvent::new(
run_id,
step,
EventKind::Replan {
window: contract.stall.window,
},
));
}
Progressing::Stalled => {
store.record_context_event(
run_id,
&ContextEvent::stalled(step, "still no progress after replanning"),
)?;
info!(run_id, step, "run stopped: stalled");
watch.emit(RunEvent::new(run_id, step, EventKind::Stalled));
finish(store, watch, run_id, 0, step, "stalled")?;
return Ok(RunResult::new(RunOutcome::Stalled { steps: step }, run_id)
.with_remembered(remembered));
}
}
if let Some(request_id) = paused {
finish(store, watch, run_id, 0, step, "awaiting_approval")?;
return Ok(RunResult::new(
RunOutcome::AwaitingApproval {
request_id,
steps: step,
},
run_id,
)
.with_remembered(remembered));
}
if !new_rules.is_empty() {
let mut layer = Policy::permissive().layer("remembered");
for r in &new_rules {
layer = layer.rule(r.act, r.effect, r.pattern.clone());
}
effective = effective.merge(layer);
ws = Workspace::with_policy(root, effective.clone());
remembered.extend(new_rules);
}
if let Some(max) = contract.max_tokens {
if tokens_used > max {
finish(store, watch, run_id, 0, step, "cost_budget_exceeded")?;
return Ok(
RunResult::new(RunOutcome::CostBudgetExceeded { steps: step }, run_id)
.with_remembered(remembered),
);
}
}
if contract
.verify
.passes_in_guarded(
root,
&ExecGuard::new(&effective)
.tracing(store, run_id, step)
.watching(watch, 0),
)
.await?
{
finish(store, watch, run_id, 0, step, "success")?;
return Ok(RunResult::new(RunOutcome::Success { steps: step }, run_id)
.with_remembered(remembered));
}
}
finish(
store,
watch,
run_id,
0,
contract.max_steps,
"step_cap_reached",
)?;
Ok(RunResult::new(
RunOutcome::StepCapReached {
steps: contract.max_steps,
},
run_id,
))
}
struct Tree<'a, P: Provider> {
mcp: &'a McpSession,
tools: &'a Toolbox,
skills: &'a Skills,
provider: &'a P,
store: &'a Store,
approver: &'a dyn Approver,
watch: &'a Watch<'a>,
ledger: Arc<Ledger>,
containment: &'a Containment,
root: PathBuf,
root_run_id: i64,
}
pub async fn run_tree<P: Provider>(
contract: &TaskContract,
provider: &P,
store: &Store,
policy: &Policy,
approver: &dyn Approver,
containment: &Containment,
) -> Result<RunResult> {
run_tree_observed(
contract,
provider,
store,
policy,
approver,
containment,
&Ignore,
)
.await
}
#[allow(clippy::too_many_arguments)]
pub async fn run_tree_observed<P: Provider>(
contract: &TaskContract,
provider: &P,
store: &Store,
policy: &Policy,
approver: &dyn Approver,
containment: &Containment,
observer: &dyn Observer,
) -> Result<RunResult> {
contract.tools.validate()?;
let skills = contract.discover_skills()?;
let root = contract.root.clone().ok_or_else(|| {
crate::error::Error::Config(
"run_tree needs a workspace contract — build it with TaskContract::workspace".into(),
)
})?;
let ledger = Arc::new(Ledger::new(containment));
let run_id = store.start_run(&contract.goal, &root.display().to_string())?;
store.set_provider(run_id, provider.name())?;
store.record_run_policy(run_id, policy)?;
let watch = &Watch::new(observer);
watch.emit(RunEvent::new(
run_id,
0,
EventKind::Started {
goal: contract.goal.clone(),
provider: provider.name().to_string(),
},
));
let policy = &match authorize_provider(provider, policy, store, run_id, approver, watch).await?
{
ProviderAccess::Granted(p) => p,
ProviderAccess::Pending(request_id) => {
return Ok(RunResult::new(
RunOutcome::AwaitingApproval {
request_id,
steps: 0,
},
run_id,
))
}
};
let mcp = McpSession::connect(&contract.mcp, policy, store, run_id, watch).await?;
let tree = Tree {
mcp: &mcp,
tools: &contract.tools,
skills: &skills,
provider,
store,
approver,
watch,
ledger,
containment,
root,
root_run_id: run_id,
};
let outcome = run_agent(&tree, contract, run_id, 0, policy, 1).await;
mcp.shutdown(store, run_id, watch).await;
Ok(RunResult::new(outcome?, run_id))
}
#[allow(clippy::too_many_arguments)]
pub async fn resume_tree<P: Provider>(
contract: &TaskContract,
provider: &P,
store: &Store,
run_id: i64,
policy: &Policy,
approver: &dyn Approver,
containment: &Containment,
) -> Result<RunResult> {
resume_tree_observed(
contract,
provider,
store,
run_id,
policy,
approver,
containment,
&Ignore,
)
.await
}
#[allow(clippy::too_many_arguments)]
pub async fn resume_tree_observed<P: Provider>(
contract: &TaskContract,
provider: &P,
store: &Store,
run_id: i64,
policy: &Policy,
approver: &dyn Approver,
containment: &Containment,
observer: &dyn Observer,
) -> Result<RunResult> {
contract.tools.validate()?;
let skills = contract.discover_skills()?;
store.check_resumable(run_id)?;
if store.run_status(run_id)? == Some(RunStatus::Completed) {
if let Some(o) = terminal_outcome(store, run_id)? {
return Ok(RunResult::new(o, run_id));
}
}
let root = contract.root.clone().ok_or_else(|| {
crate::error::Error::Config(
"resume_tree needs a workspace contract — build it with TaskContract::workspace".into(),
)
})?;
let ledger = Arc::new(Ledger::from_state(
containment,
store.spent_tokens_tree(run_id)?,
store.agent_count_tree(run_id)?,
));
let start_step = record_resume_markers(store, run_id)?;
store.set_provider(run_id, provider.name())?;
let watch = &Watch::new(observer);
watch.emit(RunEvent::new(
run_id,
start_step.saturating_sub(1),
EventKind::Started {
goal: contract.goal.clone(),
provider: provider.name().to_string(),
},
));
let policy = &match authorize_provider(provider, policy, store, run_id, approver, watch).await?
{
ProviderAccess::Granted(p) => p,
ProviderAccess::Pending(request_id) => {
return Ok(RunResult::new(
RunOutcome::AwaitingApproval {
request_id,
steps: start_step.saturating_sub(1),
},
run_id,
))
}
};
let mcp = McpSession::connect(&contract.mcp, policy, store, run_id, watch).await?;
let tree = Tree {
mcp: &mcp,
tools: &contract.tools,
skills: &skills,
provider,
store,
approver,
watch,
ledger,
containment,
root,
root_run_id: run_id,
};
let outcome = run_agent(&tree, contract, run_id, 0, policy, start_step).await;
mcp.shutdown(store, run_id, watch).await;
Ok(RunResult::new(outcome?, run_id))
}
fn run_agent<'f, P: Provider>(
tree: &'f Tree<'_, P>,
contract: &'f TaskContract,
run_id: i64,
depth: u32,
policy: &'f Policy,
start_step: u32,
) -> Pin<Box<dyn Future<Output = Result<RunOutcome>> + 'f>> {
Box::pin(async move {
let ws = Workspace::with_policy(&tree.root, policy.clone());
let mut extra = tree.tools.specs();
extra.extend(tree.mcp.tool_specs());
extra.extend(skill_tool(tree.skills));
let system =
with_skill_catalog(with_extra_tools(tree_system_prompt(), &extra), tree.skills);
let mut tools = tree_tools();
tools.extend(extra);
let token_cap = tree.ledger.effective_token_budget(contract.max_tokens);
let mut tokens_used: u64 = tree.store.spent_tokens(run_id)?;
let (mut ledger, mut written) = restore_ledger(tree.store, run_id)?;
let mut progress = Progress::new();
let mem_key = memory_key(&tree.root);
for step in start_step..=contract.max_steps {
if let Some(o) = cancelled(tree.store, tree.watch, run_id, depth, step - 1)? {
return Ok(o);
}
if let Some(max) = contract.max_duration {
if tree.store.elapsed_secs(run_id)? > max.as_secs_f64() {
finish(
tree.store,
tree.watch,
run_id,
depth,
step - 1,
"time_budget_exceeded",
)?;
return Ok(RunOutcome::TimeBudgetExceeded { steps: step - 1 });
}
}
let budget_tokens = contract
.context
.effective_tokens(Some(token_cap.saturating_sub(tokens_used)));
let entry_cap = entry_cap_chars(budget_tokens);
let notes = tree.store.memory_list(&mem_key)?;
let assembled = assemble(
&ledger,
budget_tokens,
¬es,
Assembly {
ws: Some(&ws),
policy,
store: tree.store,
run_id,
step,
},
)
.await?;
let user = workspace_user_prompt(contract, &assembled.text);
let request = CompletionRequest {
system: system.clone(),
user: user.clone(),
tools: tools.clone(),
};
let response = complete_with_retry(
tree.provider,
&request,
contract,
tree.store,
run_id,
step,
tree.watch,
depth,
)
.await?;
if let Some(served) = tree.provider.last_served() {
tree.store
.record_context_event(run_id, &ContextEvent::served(step, served.clone()))?;
tree.watch.emit(RunEvent::at_depth(
run_id,
step,
depth,
EventKind::FellBackTo { provider: served },
));
}
let step_tokens = response.usage.map(|u| u.total_tokens).unwrap_or(0);
tokens_used += step_tokens;
if step_tokens > 0 {
tree.store
.record_context_reported(run_id, step, step_tokens)?;
}
let mut decisions: Vec<String> = Vec::new();
let mut calls_json: Vec<String> = Vec::new();
let mut step_changed = false;
if response.tool_calls.is_empty() {
let said = response.text.clone().unwrap_or_default();
ledger.push(Observation::new(
step,
ObsKind::Message,
None,
bound(
&format!("\n[step {step}] (no tool call) {said}\n"),
entry_cap,
ObsKind::Message,
),
));
decisions.push("no tool call".into());
}
let mut paused: Option<i64> = None;
let mut paused_by_child = false;
let mut spawn_calls: Vec<&ToolCall> = Vec::new();
for call in &response.tool_calls {
calls_json.push(format!("{}:{}", call.name, call.arguments));
if call.name == SPAWN_TOOL {
spawn_calls.push(call);
continue;
}
match dispatch(
&ws,
call,
tree.approver,
tree.store,
run_id,
step,
tree.mcp,
tree.tools,
tree.skills,
entry_cap,
&mem_key,
tree.watch,
depth,
)
.await?
{
Dispatched::Continue {
decision,
obs,
kind,
target,
changed,
..
} => {
step_changed |= changed;
ledger.push(Observation::new(step, kind, target, obs));
decisions.push(decision);
}
Dispatched::Pause { request_id } => {
decisions.push(format!("awaiting approval (request {request_id})"));
paused = Some(request_id);
break;
}
}
}
if paused.is_none() && !spawn_calls.is_empty() {
use futures_util::stream::{self, StreamExt};
let max_c = tree.containment.max_concurrent.max(1) as usize;
let results: Vec<Result<SpawnResult>> = stream::iter(
spawn_calls
.into_iter()
.map(|c| spawn_child(tree, c, run_id, depth, policy, step)),
)
.buffered(max_c)
.collect()
.await;
for r in results {
match r? {
SpawnResult::Composed { decision, obs } => {
ledger.push(Observation::new(
step,
ObsKind::Child,
None,
bound(&obs, entry_cap, ObsKind::Child),
));
decisions.push(decision);
step_changed = true;
}
SpawnResult::Paused { request_id } => {
decisions
.push(format!("child awaiting approval (request {request_id})"));
paused = Some(request_id);
paused_by_child = true;
}
}
}
}
let committed = !(paused.is_some() && paused_by_child);
commit_step(
tree.store,
tree.watch,
run_id,
depth,
StepRecord::new(step, decisions.join("; "), ledger.text_for_step(step)).with_trace(
user,
calls_json.join(" | "),
step_tokens,
),
step_changed,
committed,
)?;
if committed {
written = persist_ledger(tree.store, run_id, &ledger, written)?;
}
let signature = calls_json.join(" | ");
match progress.step(contract.stall, step_changed, &signature) {
Progressing::Fine => {}
Progressing::Replan => {
tree.store.record_context_event(
run_id,
&ContextEvent::replan(
step,
format!(
"{} steps without progress; replanning",
contract.stall.window
),
),
)?;
ledger.push(Observation::new(
step,
ObsKind::Message,
None,
bound(
&progress.replan_directive(contract.stall.window, &decisions),
entry_cap,
ObsKind::Message,
),
));
info!(run_id, depth, step, "agent told to change approach");
tree.watch.emit(RunEvent::at_depth(
run_id,
step,
depth,
EventKind::Replan {
window: contract.stall.window,
},
));
}
Progressing::Stalled => {
tree.store.record_context_event(
run_id,
&ContextEvent::stalled(step, "still no progress after replanning"),
)?;
info!(run_id, depth, step, "agent stopped: stalled");
tree.watch
.emit(RunEvent::at_depth(run_id, step, depth, EventKind::Stalled));
finish(tree.store, tree.watch, run_id, depth, step, "stalled")?;
return Ok(RunOutcome::Stalled { steps: step });
}
}
if let Some(request_id) = paused {
finish(
tree.store,
tree.watch,
run_id,
depth,
step,
"awaiting_approval",
)?;
return Ok(RunOutcome::AwaitingApproval {
request_id,
steps: step,
});
}
let draw = tree.ledger.draw_tokens(step_tokens);
let remaining = tree.ledger.remaining_tokens();
tree.store.record_agent_event(&AgentEvent::budget_draw(
run_id,
step,
step_tokens,
remaining,
))?;
tree.watch.emit(RunEvent::at_depth(
run_id,
step,
depth,
EventKind::SpendDraw {
tokens: step_tokens,
remaining: Some(remaining),
},
));
if draw == Draw::Halted {
finish(
tree.store,
tree.watch,
run_id,
depth,
step,
"budget_ceiling_reached",
)?;
return Ok(RunOutcome::BudgetCeilingReached { steps: step });
}
if let Some(max) = tree.containment.max_total_duration {
if tree.store.elapsed_secs(tree.root_run_id)? > max.as_secs_f64() {
finish(
tree.store,
tree.watch,
run_id,
depth,
step,
"budget_ceiling_reached",
)?;
info!(run_id, depth, step, "tree stopped: duration ceiling");
return Ok(RunOutcome::BudgetCeilingReached { steps: step });
}
}
if tokens_used > token_cap {
finish(
tree.store,
tree.watch,
run_id,
depth,
step,
"cost_budget_exceeded",
)?;
return Ok(RunOutcome::CostBudgetExceeded { steps: step });
}
if contract
.verify
.passes_in_guarded(
&tree.root,
&ExecGuard::new(policy)
.tracing(tree.store, run_id, step)
.watching(tree.watch, depth),
)
.await?
{
finish(tree.store, tree.watch, run_id, depth, step, "success")?;
return Ok(RunOutcome::Success { steps: step });
}
}
finish(
tree.store,
tree.watch,
run_id,
depth,
contract.max_steps,
"step_cap_reached",
)?;
Ok(RunOutcome::StepCapReached {
steps: contract.max_steps,
})
})
}
enum SpawnResult {
Composed { decision: String, obs: String },
Paused { request_id: i64 },
}
async fn spawn_child<P: Provider>(
tree: &Tree<'_, P>,
call: &ToolCall,
parent_run_id: i64,
depth: u32,
parent_policy: &Policy,
step: u32,
) -> Result<SpawnResult> {
let a = &call.arguments;
let goal = a.get("goal").and_then(|v| v.as_str()).unwrap_or_default();
let file = a
.get("verify_file")
.and_then(|v| v.as_str())
.unwrap_or_default();
let needle = a
.get("verify_contains")
.and_then(|v| v.as_str())
.unwrap_or_default();
if goal.is_empty() || file.is_empty() {
return Ok(SpawnResult::Composed {
decision: "spawn missing fields".into(),
obs: "\n[spawn error] spawn_agent needs \"goal\" and \"verify_file\"\n".into(),
});
}
let child_depth = depth + 1;
let mut overlay = Policy::permissive().layer("child");
if let Some(denies) = a.get("deny_write").and_then(|v| v.as_array()) {
for d in denies.iter().filter_map(|v| v.as_str()) {
overlay = overlay.deny_write(d);
}
}
if let Some(denies) = a.get("deny_net").and_then(|v| v.as_array()) {
for d in denies.iter().filter_map(|v| v.as_str()) {
overlay = overlay.deny_net(d);
}
}
let child_policy = parent_policy.contain(&overlay);
let verify = Verification::WorkspaceFileContains {
file: file.into(),
needle: needle.into(),
};
let mut child_contract = TaskContract::workspace(goal, &tree.root, verify);
if let Some(n) = a.get("max_steps").and_then(|v| v.as_u64()) {
child_contract = child_contract.with_max_steps(n as u32);
}
let (child_run, child_start) = match tree.store.find_spawn(parent_run_id, step, goal)? {
Some(row) => {
if let Some(o) = terminal_outcome(tree.store, row.child_run_id)? {
return Ok(compose_child(row.child_run_id, goal, o));
}
(
row.child_run_id,
tree.store.last_step(row.child_run_id)? + 1,
)
}
None => {
if let Err(refusal) = tree.ledger.register_agent(child_depth) {
tree.store.record_agent_event(&AgentEvent::spawn_refused(
parent_run_id,
step,
refusal.cap(),
))?;
tree.watch.emit(RunEvent::at_depth(
parent_run_id,
step,
depth,
EventKind::SpawnRefused {
cap: refusal.cap().to_string(),
},
));
return Ok(SpawnResult::Composed {
decision: format!("spawn refused ({})", refusal.cap()),
obs: format!(
"\n[spawn refused] {refusal} — adapt or finish with what you have\n"
),
});
}
let child_run = tree.store.start_child_run(
goal,
&tree.root.display().to_string(),
parent_run_id,
child_depth,
)?;
tree.store.record_agent_event(&AgentEvent::spawn(
parent_run_id,
step,
child_run,
goal,
))?;
tree.watch.emit(RunEvent::at_depth(
parent_run_id,
step,
depth,
EventKind::Spawned {
child_run_id: child_run,
goal: goal.to_string(),
},
));
let deny_json = a
.get("deny_write")
.map(|v| v.to_string())
.unwrap_or_else(|| "[]".into());
tree.store.record_spawn(
parent_run_id,
step,
child_run,
goal,
file,
needle,
a.get("max_steps")
.and_then(|v| v.as_u64())
.map(|n| n as u32),
&deny_json,
)?;
(child_run, 1)
}
};
let outcome = run_agent(
tree,
&child_contract,
child_run,
child_depth,
&child_policy,
child_start,
)
.await?;
if let RunOutcome::AwaitingApproval { request_id, .. } = outcome {
return Ok(SpawnResult::Paused { request_id });
}
Ok(compose_child(child_run, goal, outcome))
}
fn compose_child(child_run: i64, goal: &str, outcome: RunOutcome) -> SpawnResult {
SpawnResult::Composed {
decision: format!("spawned child {child_run}: {outcome:?}"),
obs: format!("\n[child {child_run} \"{goal}\" -> {outcome:?}]\n"),
}
}
fn remembered_layer(rules: &[Rule]) -> Policy {
let mut layer = Policy::permissive().layer("remembered");
for r in rules {
layer = layer.rule(r.act, r.effect, r.pattern.clone());
}
layer
}
enum ProviderAccess {
Granted(Policy),
Pending(i64),
}
async fn authorize_provider<P: Provider>(
provider: &P,
policy: &Policy,
store: &Store,
run_id: i64,
approver: &dyn Approver,
watch: &Watch<'_>,
) -> Result<ProviderAccess> {
let urls = provider.endpoints();
if urls.is_empty() {
return Ok(ProviderAccess::Granted(policy.clone()));
}
let mut effective = policy.clone();
let mut ask: Option<String> = None;
for url in urls {
let Some(target) = net::target(url) else {
return Err(crate::error::Error::Refused {
act: "net".into(),
target: url.to_string(),
rule: None,
layer: None,
});
};
effective = effective.merge(net::provider_layer(&target));
let verdict = NetGuard::new(&effective)
.tracing(store, run_id, 0)
.watching(watch, 0)
.check_target(&target)?;
if verdict.effect == Effect::Ask {
ask = Some(target.clone());
}
}
let Some(target) = ask else {
return Ok(ProviderAccess::Granted(effective));
};
watch.emit(RunEvent::new(
run_id,
0,
EventKind::ApprovalRequested {
act: "net".into(),
target: target.clone(),
},
));
match approver.decide(&Request::new(Act::Net, &target)).await {
Decision::Approve { .. } => {
let ev = PolicyEvent::decision(0, "net", &target, "approve", "approver");
store.record_event(run_id, &ev)?;
decided(watch, run_id, 0, &ev);
Ok(ProviderAccess::Granted(effective))
}
Decision::Deny { reason } => {
let ev = PolicyEvent::decision(0, "net", &target, "deny", "approver");
store.record_event(run_id, &ev)?;
decided(watch, run_id, 0, &ev);
finish(store, watch, run_id, 0, 0, "refused")?;
Err(crate::error::Error::Refused {
act: "net".into(),
target: format!("{target} — {reason}"),
rule: None,
layer: None,
})
}
Decision::Defer => {
let ev = PolicyEvent::decision(0, "net", &target, "defer", "approver");
store.record_event(run_id, &ev)?;
decided(watch, run_id, 0, &ev);
let request_id = store.put_pending(run_id, 0, "net", &target, None)?;
finish(store, watch, run_id, 0, 0, "awaiting_approval")?;
Ok(ProviderAccess::Pending(request_id))
}
}
}
enum Dispatched {
Continue {
decision: String,
obs: String,
kind: ObsKind,
target: Option<String>,
changed: bool,
remember: Vec<Rule>,
},
Pause { request_id: i64 },
}
impl Dispatched {
fn seen(
decision: impl Into<String>,
obs: impl Into<String>,
kind: ObsKind,
target: Option<String>,
) -> Self {
Dispatched::Continue {
decision: decision.into(),
obs: obs.into(),
kind,
target,
changed: false,
remember: Vec::new(),
}
}
fn go(decision: impl Into<String>, obs: impl Into<String>) -> Self {
Self::seen(decision, obs, ObsKind::Error, None)
}
}
#[allow(clippy::too_many_arguments)]
async fn dispatch(
ws: &Workspace,
call: &ToolCall,
approver: &dyn Approver,
store: &Store,
run_id: i64,
step: u32,
mcp: &McpSession,
custom: &Toolbox,
skills: &Skills,
cap: usize,
memory_key: &str,
watch: &Watch<'_>,
depth: u32,
) -> Result<Dispatched> {
let a = &call.arguments;
let s = |k: &str| a.get(k).and_then(|v| v.as_str());
watch.emit(RunEvent::at_depth(
run_id,
step,
depth,
EventKind::ToolCall {
name: call.name.clone(),
target: ["path", "pattern", "name_glob", "glob", "key", "name"]
.into_iter()
.find_map(s)
.unwrap_or(&call.name)
.to_string(),
},
));
Ok(match call.name.as_str() {
GREP_TOOL => {
let pattern = s("pattern").unwrap_or_default();
match ws.grep(pattern, s("path_glob")) {
Ok(hits) => {
let shown: Vec<String> = hits
.iter()
.take(OBS_GREP_CAP)
.map(|m| format!("{}:{}: {}", m.path, m.line, m.text))
.collect();
Dispatched::seen(
format!("grep {pattern:?} ({} hits)", hits.len()),
bound(
&format!("\n[grep {pattern:?}]\n{}\n", shown.join("\n")),
cap,
ObsKind::Grep,
),
ObsKind::Grep,
Some(pattern.to_string()),
)
}
Err(e) => Dispatched::go("grep error", format!("\n[grep error] {e}\n")),
}
}
FIND_TOOL => {
let glob = s("name_glob").or_else(|| s("glob")).unwrap_or_default();
match ws.find(glob) {
Ok(paths) => Dispatched::seen(
format!("find {glob:?} ({} paths)", paths.len()),
bound(
&format!("\n[find {glob:?}]\n{}\n", paths.join("\n")),
cap,
ObsKind::Find,
),
ObsKind::Find,
Some(glob.to_string()),
),
Err(e) => Dispatched::go("find error", format!("\n[find error] {e}\n")),
}
}
REMEMBER_TOOL => {
let key = s("key").unwrap_or_default();
let value = s("value").unwrap_or_default();
if key.is_empty() || value.is_empty() {
return Ok(Dispatched::go(
"remember error",
"\n[remember error] both key and value are required\n",
));
}
let evicted = store.memory_put(memory_key, key, value, run_id, step)?;
store.record_context_event(
run_id,
&ContextEvent::memory_write(
step,
format!("{key} ({} chars)", value.chars().count()),
),
)?;
watch.emit(RunEvent::at_depth(
run_id,
step,
depth,
EventKind::MemoryWrote {
key: key.to_string(),
},
));
for gone in &evicted {
store.record_context_event(
run_id,
&ContextEvent::memory_evict(step, format!("{gone} (evicted to hold the cap)")),
)?;
}
info!(run_id, step, key, evicted = evicted.len(), "remembered");
Dispatched::seen(
format!("remembered {key}"),
format!("\n[remember {key}]\n"),
ObsKind::Tool,
None,
)
}
READ_FILE_TOOL => {
let path = s("path").unwrap_or_default();
match gate(
ws,
approver,
store,
run_id,
step,
Act::Read,
path,
None,
watch,
depth,
)
.await?
{
Gated::Refused { decision, obs } => Dispatched::go(decision, obs),
Gated::Paused { request_id } => Dispatched::Pause { request_id },
Gated::Go {
target, remember, ..
} => match ws.read_file(&target) {
Ok(c) => Dispatched::Continue {
decision: format!("read {target}"),
obs: format!("\n[read {target}]\n{}\n", bound(&c, cap, ObsKind::Read)),
kind: ObsKind::Read,
target: Some(target.clone()),
changed: false,
remember,
},
Err(e) => Dispatched::go("read error", format!("\n[read error] {e}\n")),
},
}
}
WRITE_FILE_TOOL => {
let path = s("path").unwrap_or_default();
let content = s("content").unwrap_or_default();
if path.is_empty() {
return Ok(Dispatched::go(
"write missing path",
"\n[write error] write_file needs a \"path\" in workspace mode\n",
));
}
match gate(
ws,
approver,
store,
run_id,
step,
Act::Write,
path,
Some(content),
watch,
depth,
)
.await?
{
Gated::Refused { decision, obs } => Dispatched::go(decision, obs),
Gated::Paused { request_id } => Dispatched::Pause { request_id },
Gated::Go {
target,
content,
remember,
} => {
let body = content.unwrap_or_default();
match ws.write_file(&target, &body) {
Ok(wrote) => Dispatched::Continue {
decision: format!("wrote {target}"),
obs: bound(
&format!(
"\n[wrote {target}] ({} chars{})\n",
body.chars().count(),
if wrote.moved_the_workspace() {
""
} else {
", identical to what was already there — the \
workspace did not change"
}
),
cap,
ObsKind::Write,
),
kind: ObsKind::Write,
target: Some(target.clone()),
changed: wrote.moved_the_workspace(),
remember,
},
Err(e) => Dispatched::go("write error", format!("\n[write error] {e}\n")),
}
}
}
}
READ_SKILL_TOOL if !skills.is_empty() => {
let name = s("name").unwrap_or_default();
let Some(skill) = skills.get(name) else {
return Ok(Dispatched::go(
format!("unknown skill {name}"),
format!(
"\n[read_skill] there is no skill named {name:?}. Available: {}\n",
skills.names().join(", ")
),
));
};
let path = skill.path.display().to_string();
match gate(
ws,
approver,
store,
run_id,
step,
Act::Read,
&path,
None,
watch,
depth,
)
.await?
{
Gated::Refused { decision, obs } => Dispatched::go(decision, obs),
Gated::Paused { request_id } => Dispatched::Pause { request_id },
Gated::Go {
target, remember, ..
} => match std::fs::read_to_string(&target) {
Ok(body) => {
let (body, truncated) = crate::tools::cap_result(body, cap);
info!(run_id, step, skill = name, truncated, "skill read");
Dispatched::Continue {
decision: format!("read skill {name}"),
obs: format!("\n[skill {name}]\n{body}\n"),
kind: ObsKind::Skill,
target: Some(name.to_string()),
changed: false,
remember,
}
}
Err(e) => Dispatched::go(
format!("skill {name} read error"),
format!("\n[skill {name} error] {e}\n"),
),
},
}
}
name if custom.owns(name) => {
match gate(
ws,
approver,
store,
run_id,
step,
Act::Exec,
name,
None,
watch,
depth,
)
.await?
{
Gated::Refused { decision, obs } => Dispatched::go(decision, obs),
Gated::Paused { request_id } => Dispatched::Pause { request_id },
Gated::Go { remember, .. } => {
let tool = custom.get(name).expect("owns() and get() agree");
match tool.invoke(&call.arguments).await {
Ok(out) => {
let (out, truncated) = crate::tools::cap_result(out, cap);
info!(run_id, step, tool = name, truncated, "registered tool call");
Dispatched::Continue {
decision: format!("called {name}"),
obs: format!("\n[{name}]\n{out}\n"),
kind: ObsKind::Tool,
target: Some(name.to_string()),
changed: false,
remember,
}
}
Err(e) => {
info!(run_id, step, tool = name, error = %e, "registered tool failed");
Dispatched::Continue {
decision: format!("{name} failed"),
obs: format!("\n[{name} error] {e}\n"),
kind: ObsKind::Error,
target: None,
changed: false,
remember,
}
}
}
}
}
}
name if mcp.owns(name) => {
let verdict = ws.policy().check(Act::Exec, name);
if verdict.effect != Effect::Allow {
let mut ev = PolicyEvent::refusal(step, "exec", name);
ev.rule = verdict.rule.clone();
ev.layer = verdict.layer.clone();
store.record_event(run_id, &ev)?;
refused(watch, run_id, depth, &ev);
let why = verdict
.rule
.as_deref()
.map(|r| format!(" (rule {r})"))
.unwrap_or_default();
return Ok(Dispatched::go(
format!("{name} refused"),
format!("\n[{name} refused]{why} — the policy forbids calling this tool\n"),
));
}
let out = mcp
.call(
name,
&call.arguments,
store,
run_id,
step,
cap,
watch,
depth,
)
.await?;
Dispatched::seen(
format!("called {name}"),
format!("\n[{name}]\n{out}\n"),
ObsKind::Mcp,
Some(name.to_string()),
)
}
other => Dispatched::go(
format!("unknown tool {other}"),
format!("\n[unknown tool {other}]\n"),
),
})
}
enum Gated {
Go {
target: String,
content: Option<String>,
remember: Vec<Rule>,
},
Refused { decision: String, obs: String },
Paused { request_id: i64 },
}
#[allow(clippy::too_many_arguments)]
async fn gate(
ws: &Workspace,
approver: &dyn Approver,
store: &Store,
run_id: i64,
step: u32,
act: Act,
target: &str,
content: Option<&str>,
watch: &Watch<'_>,
depth: u32,
) -> Result<Gated> {
let kind = format!("{act:?}").to_lowercase();
let check = |act: Act, target: &str| match act {
Act::Exec | Act::Net => ws.policy().check(act, target),
Act::Read | Act::Write if Path::new(target).is_absolute() => ws.policy().check(act, target),
Act::Read | Act::Write => ws.check_path(act, target),
};
let verdict = check(act, target);
match verdict.effect {
Effect::Deny => {
let mut ev = PolicyEvent::refusal(step, &kind, target);
if let (Some(rule), layer) = (verdict.rule.clone(), verdict.layer.clone()) {
ev.rule = Some(rule);
ev.layer = layer;
}
store.record_event(run_id, &ev)?;
refused(watch, run_id, depth, &ev);
let why = verdict
.rule
.as_deref()
.map(|r| format!(" (rule {r})"))
.unwrap_or_default();
Ok(Gated::Refused {
decision: format!("{kind} refused"),
obs: format!("\n[{kind} refused] {target}{why} — the policy forbids this; try another path\n"),
})
}
Effect::Allow => Ok(Gated::Go {
target: target.to_string(),
content: content.map(str::to_string),
remember: Vec::new(),
}),
Effect::Ask => {
let mut request = Request::new(act, target);
if let Some(c) = content {
request = request.with_content(c);
}
watch.emit(RunEvent::at_depth(
run_id,
step,
depth,
EventKind::ApprovalRequested {
act: kind.clone(),
target: target.to_string(),
},
));
match approver.decide(&request).await {
Decision::Approve { modified, remember } => {
let performed = modified.unwrap_or_else(|| request.clone());
let recheck = check(act, &performed.target);
if recheck.effect == Effect::Deny {
let mut ev = PolicyEvent::refusal(step, &kind, &performed.target);
ev.rule = recheck.rule.clone();
ev.layer = recheck.layer.clone();
store.record_event(run_id, &ev)?;
refused(watch, run_id, depth, &ev);
return Ok(Gated::Refused {
decision: format!("{kind} refused after approval"),
obs: format!(
"\n[{kind} refused] {} — an approved change may not cross a deny\n",
performed.target
),
});
}
let mut ev = PolicyEvent::decision(step, &kind, target, "approve", "approver");
if performed.target != target {
ev = ev.with_performed(&performed.target);
}
store.record_event(run_id, &ev)?;
decided(watch, run_id, depth, &ev);
Ok(Gated::Go {
target: performed.target,
content: performed.content,
remember,
})
}
Decision::Deny { reason } => {
let ev = PolicyEvent::decision(step, &kind, target, "deny", "approver");
store.record_event(run_id, &ev)?;
decided(watch, run_id, depth, &ev);
Ok(Gated::Refused {
decision: format!("{kind} denied"),
obs: format!("\n[{kind} denied] {target} — {reason}\n"),
})
}
Decision::Defer => {
let ev = PolicyEvent::decision(step, &kind, target, "defer", "approver");
store.record_event(run_id, &ev)?;
decided(watch, run_id, depth, &ev);
let request_id = store.put_pending(run_id, step, &kind, target, content)?;
Ok(Gated::Paused { request_id })
}
}
}
}
}
fn memory_key(root: &Path) -> String {
std::fs::canonicalize(root)
.unwrap_or_else(|_| root.to_path_buf())
.to_string_lossy()
.into_owned()
}
#[allow(clippy::too_many_arguments)]
async fn complete_with_retry<P: Provider>(
provider: &P,
request: &CompletionRequest,
contract: &TaskContract,
store: &Store,
run_id: i64,
step: u32,
watch: &Watch<'_>,
depth: u32,
) -> Result<CompletionResponse> {
let max_retries = contract.max_retries;
let retry = contract.retry;
let max_duration = contract.max_duration;
let mut attempt = 0;
loop {
match provider.complete(request.clone()).await {
Ok(response) => return Ok(response),
Err(e) if attempt < max_retries && retryable(&e) => {
attempt += 1;
let wait = retry.wait(attempt, retry_after(&e));
if let Some(max) = max_duration {
let elapsed = store.elapsed_secs(run_id)?;
if elapsed + wait.as_secs_f64() > max.as_secs_f64() {
store.record(
run_id,
&StepRecord::new(
step,
format!(
"escalated after {} (a retry would outlast the time budget)",
kind_of(&e)
),
e.to_string(),
),
)?;
let steps = store.last_step(run_id)?;
finish(store, watch, run_id, depth, steps, escalation_outcome(&e))?;
return Err(e);
}
}
store.record(
run_id,
&StepRecord::new(
step,
format!("retry {attempt} after {} in {:?}", kind_of(&e), wait),
e.to_string(),
),
)?;
watch.emit(RunEvent::at_depth(
run_id,
step,
depth,
EventKind::Retry {
kind: kind_of(&e),
attempt,
delay_ms: wait.as_millis() as u64,
},
));
if !wait.is_zero() {
tokio::time::sleep(wait).await;
}
}
Err(e) => {
store.record(
run_id,
&StepRecord::new(
step,
format!("escalated after {}", kind_of(&e)),
e.to_string(),
),
)?;
let steps = store.last_step(run_id)?;
finish(store, watch, run_id, depth, steps, escalation_outcome(&e))?;
return Err(e);
}
}
}
}
fn retryable(e: &Error) -> bool {
matches!(e, Error::Provider { kind, .. } if kind.is_retryable())
}
fn retry_after(e: &Error) -> Option<std::time::Duration> {
match e {
Error::Provider { retry_after, .. } => *retry_after,
_ => None,
}
}
fn kind_of(e: &Error) -> String {
match e {
Error::Provider { kind, status, .. } => match status {
Some(s) => format!("{kind:?} (HTTP {s})"),
None => format!("{kind:?}"),
},
other => format!("{other}"),
}
}
fn escalation_outcome(e: &Error) -> &'static str {
if retryable(e) {
"escalated_retryable"
} else {
"escalated_terminal"
}
}
fn system_prompt() -> String {
"You are an agent that edits exactly one file to meet a stated specification. \
Call the `write_file` tool with the file's full new contents. Do not explain; \
make the edit. The file will be checked against the success criterion after \
each write."
.to_string()
}
fn user_prompt(contract: &TaskContract, current: &str) -> String {
let constraints = if contract.constraints.is_empty() {
"(none)".to_string()
} else {
contract.constraints.join("; ")
};
format!(
"Goal: {goal}\nConstraints: {constraints}\nSuccess criterion: {criterion}\n\n\
Current file contents:\n---\n{current}\n---\n\n\
Call write_file with the full new contents that satisfy the success criterion.",
goal = contract.goal,
criterion = contract.verify.describe(),
)
}
fn write_file_tool() -> ToolSpec {
ToolSpec {
name: WRITE_FILE_TOOL.to_string(),
description: "Write the full new contents of the target file.".to_string(),
parameters: json!({
"type": "object",
"properties": {
"content": { "type": "string", "description": "Full new file contents." }
},
"required": ["content"]
}),
}
}
fn workspace_system_prompt() -> String {
"You are an agent working across a repository to meet a stated specification. \
Use `grep` to search file contents and `find` to locate files by name, then \
`read_file` to inspect a file before changing it, and `write_file` with the \
file's path and full new contents to edit it. You may edit several files. \
Work in small steps; after each of your steps the whole set is checked \
against the success criterion. Do not explain; call tools."
.to_string()
}
fn with_extra_tools(base: String, extra: &[ToolSpec]) -> String {
if extra.is_empty() {
return base;
}
let names: Vec<&str> = extra.iter().map(|t| t.name.as_str()).collect();
format!(
"{base} These extra tools are also available and work the same way: {}. \
Each tool's result appears in the observations below; once a tool has \
returned what you asked for, move on rather than calling it again.",
names.join(", ")
)
}
fn skill_tool(skills: &Skills) -> Option<ToolSpec> {
if skills.is_empty() {
return None;
}
Some(ToolSpec {
name: READ_SKILL_TOOL.to_string(),
description: "Load one skill's full instructions into your observations, by the name it \
is listed under. Read a skill when its description says it covers what you \
are about to do."
.to_string(),
parameters: json!({
"type": "object",
"properties": {
"name": { "type": "string", "description": "The skill's name, as listed in the system prompt." }
},
"required": ["name"]
}),
})
}
fn with_skill_catalog(base: String, skills: &Skills) -> String {
if skills.is_empty() {
return base;
}
format!(
"{base}\n\nSkills available to you — instructions written for this repository. Only each \
skill's name and description is shown; call `{READ_SKILL_TOOL}` with a name to read that \
skill's full text when its description matches what you are doing.\n{}",
skills.catalog()
)
}
fn workspace_user_prompt(contract: &TaskContract, observations: &str) -> String {
let constraints = if contract.constraints.is_empty() {
"(none)".to_string()
} else {
contract.constraints.join("; ")
};
let obs = if observations.is_empty() {
"(nothing yet — start by grepping or finding)".to_string()
} else {
observations.to_string()
};
format!(
"Goal: {goal}\nConstraints: {constraints}\nSuccess criterion: {criterion}\n\n\
Observations so far (results of your tool calls):\n{obs}\n\n\
Call a tool to make progress toward the success criterion.",
goal = contract.goal,
criterion = contract.verify.describe(),
)
}
fn tree_system_prompt() -> String {
"You are an agent working across a repository to meet a stated specification. \
Use `grep`, `find`, `read_file`, and `write_file` as in a normal run. You may \
also decompose the work: call `spawn_agent` to launch a sub-agent that pursues \
a smaller goal over the same workspace, and its result is reported back to you. \
A sub-agent inherits your permissions and can only be more restricted, never \
less. Prefer spawning when parts of the task are independent. Work in small \
steps; the whole set is checked against the success criterion after each. Do \
not explain; call tools."
.to_string()
}
fn tree_tools() -> Vec<ToolSpec> {
let mut tools = workspace_tools();
tools.push(ToolSpec {
name: SPAWN_TOOL.to_string(),
description: "Spawn a contained sub-agent to pursue a smaller goal over the same \
workspace. The sub-agent inherits your permissions (it can only be \
further restricted) and its outcome is reported back to you."
.to_string(),
parameters: json!({
"type": "object",
"properties": {
"goal": { "type": "string", "description": "The sub-agent's goal." },
"verify_file": { "type": "string", "description": "File (relative to the workspace root) whose contents decide the sub-agent's success." },
"verify_contains": { "type": "string", "description": "Text that file must contain for the sub-agent to succeed." },
"deny_write": { "type": "array", "items": { "type": "string" }, "description": "Optional globs the sub-agent must not write — tightens its inherited policy." },
"deny_net": { "type": "array", "items": { "type": "string" }, "description": "Optional host globs (host or host:port) the sub-agent must not reach — tightens its inherited policy." },
"max_steps": { "type": "integer", "description": "Optional step budget for the sub-agent." }
},
"required": ["goal", "verify_file", "verify_contains"]
}),
});
tools
}
fn workspace_tools() -> Vec<ToolSpec> {
vec![
ToolSpec {
name: GREP_TOOL.to_string(),
description: "Search file contents by regex (a plain substring is valid). Returns file:line: matches.".to_string(),
parameters: json!({
"type": "object",
"properties": {
"pattern": { "type": "string", "description": "Regex or substring to search for." },
"path_glob": { "type": "string", "description": "Optional glob limiting which files are searched, e.g. src/*.rs." }
},
"required": ["pattern"]
}),
},
ToolSpec {
name: FIND_TOOL.to_string(),
description: "List files whose name or relative path matches a glob (* and ?).".to_string(),
parameters: json!({
"type": "object",
"properties": {
"name_glob": { "type": "string", "description": "Glob to match, e.g. *.rs or src/*.rs." }
},
"required": ["name_glob"]
}),
},
ToolSpec {
name: READ_FILE_TOOL.to_string(),
description: "Read a file (path relative to the workspace root) into context.".to_string(),
parameters: json!({
"type": "object",
"properties": {
"path": { "type": "string", "description": "File path relative to the workspace root." }
},
"required": ["path"]
}),
},
ToolSpec {
name: REMEMBER_TOOL.to_string(),
description: "Record a short fact or decision worth keeping for a later run over this \
workspace — a build command, a layout you had to discover, a decision and \
why. Notes are yours, not instructions, and are recalled at the start of \
later runs so you do not rediscover the same thing twice."
.to_string(),
parameters: json!({
"type": "object",
"properties": {
"key": { "type": "string", "description": "Short name to recall it by; writing the same key again replaces it." },
"value": { "type": "string", "description": "The fact, in one or two sentences." }
},
"required": ["key", "value"]
}),
},
ToolSpec {
name: WRITE_FILE_TOOL.to_string(),
description: "Write the full new contents of a file (path relative to the workspace root); creates it if absent.".to_string(),
parameters: json!({
"type": "object",
"properties": {
"path": { "type": "string", "description": "File path relative to the workspace root." },
"content": { "type": "string", "description": "Full new file contents." }
},
"required": ["path", "content"]
}),
},
]
}