use std::sync::Arc;
use async_trait::async_trait;
use serde_json::{Value, json};
use crate::core::{Outcome, Skill, SkillDescriptor, SkillError, Tainted};
use crate::manifest::{ExecutionKind, Identity, Manifest};
use crate::model::{ModelCall, ModelProvider, ModelRole};
use super::StepCtx;
#[derive(Debug)]
pub(super) struct Declarative {
kind: ExecutionKind,
capability: String,
name: String,
provider: Arc<dyn ModelProvider>,
tools: Option<(
Arc<crate::tools::ToolCatalog>,
Arc<dyn crate::tools::ToolClient>,
)>,
max_turns: u32,
}
impl Declarative {
pub(super) fn new(
kind: ExecutionKind,
capability: String,
name: String,
provider: Arc<dyn ModelProvider>,
tools: Option<(
Arc<crate::tools::ToolCatalog>,
Arc<dyn crate::tools::ToolClient>,
)>,
max_turns: u32,
) -> Self {
Self {
kind,
capability,
name,
provider,
tools,
max_turns,
}
}
#[allow(clippy::too_many_arguments, clippy::too_many_lines)]
async fn tool_loop(
&self,
cx: &mut StepCtx<'_>,
input: Tainted<Value>,
remembered: Option<Tainted<Value>>,
system: String,
role: ModelRole,
egress: Option<crate::core::Sensitivity>,
granted: Vec<crate::manifest::ToolGrant>,
oversight: Option<Proposal>,
formation: Option<crate::manifest::MemoryFormation>,
output_schema: Option<Value>,
) -> Result<Outcome, SkillError> {
let (catalog, client) = self.tools.clone().ok_or_else(|| {
SkillError::Other(
"this agent declares `tool-calling` but the plane has no tool \
catalogue — `RuntimeBuilder::tools` is what lets a declarative \
agent reach one"
.into(),
)
})?;
let (offered, declared) = offered_tools(&catalog, &granted);
let bindable_input = input.clone();
let prompt = prompt_object(&system, input, remembered);
let mut exchanges: Vec<crate::model::ToolExchange> = Vec::new();
let mut continuation: Option<crate::model::ProviderContinuation> = None;
let mut conversation_label = prompt.label().clone();
for _turn in 0..self.max_turns {
let outbound = prompt.with_joined_label(&conversation_label);
let completion = cx
.sink_with(&outbound, |value| {
let mut call =
ModelCall::new(Arc::clone(&self.provider), role.model.clone(), value)
.with_tools(declared.clone())
.continuing(exchanges.clone());
if let Some(schema) = output_schema.clone() {
call = call.expecting(schema);
}
if let Some(state) = continuation.take() {
call = call.with_continuation(state);
}
call = role.applied_to(call);
if let Some(ceiling) = egress {
call = call.with_max_sensitivity(ceiling);
}
call
})
.await?;
let label = completion.label().join(&conversation_label);
let completion = Tainted::with_label(completion.into_unlabelled(), label);
if completion.peek().truncated {
let reason = if completion.peek().tool_calls.is_empty() {
"the model ran out of output budget mid-answer, so this is a partial answer and not the agent's answer — raise `max_output_tokens` for this role, or narrow what the agent is asked to produce"
} else {
"the model ran out of output budget while it was still asking for tools, so the last call's arguments are whatever survived the cut — running them would act on a request the model never finished writing. Raise `max_output_tokens` for this role"
};
return Ok(Outcome::fail(reason));
}
if completion.peek().tool_calls.is_empty() {
let formed_source = Tainted::with_label(
completion
.peek()
.structured
.clone()
.unwrap_or_else(|| json!({ "text": completion.peek().text.clone() })),
completion.label().clone(),
);
let answer =
completion.map(|c| c.structured.unwrap_or_else(|| json!({ "text": c.text })));
if let Some(schema) = output_schema.as_ref()
&& let Err(detail) = crate::model::validate_schema(schema, answer.peek())
{
return Ok(Outcome::fail(format!(
"the answer does not satisfy the declared output shape: {detail}"
)));
}
return self
.settle(
cx,
answer,
formed_source,
&bindable_input,
oversight,
formation.as_ref(),
&role,
)
.await;
}
continuation.clone_from(&completion.peek().continuation);
exchanges.clear();
for asked in completion.peek().tool_calls.clone() {
let Some((id, grant)) = catalog.resolve(&asked.name).and_then(|id| {
offered
.iter()
.find(|(offered, _)| *offered == id)
.map(|(_, grant)| (id, *grant))
}) else {
exchanges.push(crate::model::ToolExchange::failed(
asked,
"no tool of that name is granted to this agent",
));
continue;
};
let Some(declaration) = declared.iter().find(|tool| tool.name == asked.name) else {
exchanges.push(crate::model::ToolExchange::failed(
asked,
"the selected tool has no model-facing declaration",
));
continue;
};
if let Err(detail) =
crate::model::validate_schema(&declaration.parameters, &asked.arguments)
{
exchanges.push(crate::model::ToolExchange::failed(asked, detail));
continue;
}
let reference = id.reference();
let mut args = crate::core::Tainted::with_label(
asked.arguments.clone(),
completion.label().clone(),
);
if grant.requires_approval {
let Some(spec) = oversight.as_ref() else {
return Err(SkillError::Other(
"a tool grant requires approval but the agent declares no oversight policy — there is nobody to ask"
.into(),
));
};
cx.deadline(spec.deadline.name.clone(), &spec.deadline.spec(), None)
.await?;
let mut task = spec.approve_call(&reference, &asked.arguments);
if let Some(preview) = grant.preview.as_deref() {
let evidence =
preview_evidence(cx, &catalog, &client, preview, &args).await;
task.justification.evidence.push(evidence);
}
let decision = cx.task(&task).await?;
if !decision.approved {
exchanges.push(crate::model::ToolExchange::failed(
asked,
"a reviewer did not approve this call",
));
continue;
}
args = match approved_arguments(&decision, &declaration.parameters, args) {
Ok(args) => args,
Err(detail) => {
exchanges.push(crate::model::ToolExchange::failed(asked, detail));
continue;
}
};
}
if id.server == crate::tools::AGENT_SERVER {
match cx.commission(&id.tool, args).await {
Ok(answer) => {
conversation_label = conversation_label.join(answer.label());
exchanges
.push(crate::model::ToolExchange::ok(asked, answer.peek().clone()));
}
Err(e) => match relayed_failure(cx, &mut conversation_label, &e) {
Some(detail) => {
exchanges.push(crate::model::ToolExchange::failed(asked, detail));
}
None => return Err(e.into()),
},
}
continue;
}
if let Some(peer) = cx.peer_named(&id.server) {
match cx.call_peer(&peer, &id.tool, &args).await {
Ok(answer) => {
conversation_label = conversation_label.join(answer.label());
exchanges
.push(crate::model::ToolExchange::ok(asked, answer.peek().clone()));
}
Err(e) => match relayed_failure(cx, &mut conversation_label, &e) {
Some(detail) => {
exchanges.push(crate::model::ToolExchange::failed(asked, detail));
}
None => return Err(e.into()),
},
}
continue;
}
let egress = cx.tool_egress();
let dispatched = cx
.sink_with(&args, |value| {
crate::tools::ToolCall::prepare(
&catalog,
Arc::clone(&client),
id,
value,
egress.as_deref(),
)
.map_err(|e| {
crate::core::StepError::Effect(crate::core::EffectError::Rejected(
e.to_string(),
))
})
})
.await;
match dispatched {
Ok(result) => {
conversation_label = conversation_label.join(result.label());
exchanges
.push(crate::model::ToolExchange::ok(asked, result.peek().clone()));
}
Err(e) => match relayed_failure(cx, &mut conversation_label, &e) {
Some(detail) => {
exchanges.push(crate::model::ToolExchange::failed(asked, detail));
}
None => return Err(e.into()),
},
}
}
}
Ok(Outcome::fail(format!(
"'{}' did not finish within {} model turns — it was still asking for tools",
self.name, self.max_turns
)))
}
#[allow(clippy::too_many_arguments)]
async fn settle(
&self,
cx: &mut StepCtx<'_>,
answer: Tainted<Value>,
formed_source: Tainted<Value>,
input: &Tainted<Value>,
oversight: Option<Proposal>,
formation: Option<&crate::manifest::MemoryFormation>,
role: &ModelRole,
) -> Result<Outcome, SkillError> {
if let Some(spec) = oversight.as_ref().filter(|s| s.gates_the_answer()) {
cx.deadline(spec.deadline.name.clone(), &spec.deadline.spec(), None)
.await?;
let decision = cx.task(&spec.approve_answer(answer.peek().clone())).await?;
if !decision.approved {
return Ok(Outcome::fail(format!(
"{} refused this answer: {}",
decision.decided, decision.reason
)));
}
}
self.form_answer(cx, formation, formed_source, input, role)
.await?;
if let Some(spec) = oversight.as_ref() {
self.triage(cx, spec, &answer).await?;
}
Ok(Outcome::done(answer))
}
async fn triage(
&self,
cx: &mut StepCtx<'_>,
oversight: &Proposal,
answer: &Tainted<Value>,
) -> Result<(), SkillError> {
for rule in oversight.triage.iter().filter(|r| r.matches(answer.peek())) {
cx.deadline(rule.deadline.name.clone(), &rule.deadline.spec(), None)
.await?;
let mut spec = crate::core::TaskSpec::new(
rule.task_kind(),
crate::core::Justification::new(
Tainted::trusted(rule.summary.clone()),
answer.peek().clone(),
),
rule.deadline.name.clone(),
);
spec.candidate_roles.clone_from(&rule.audience);
spec.priority = rule.priority.into();
spec.on_expiry = match &oversight.on_expiry {
escalate @ crate::core::Expiry::Escalate { .. } => escalate.clone(),
_ => crate::core::Expiry::Deny,
};
cx.open_task(&spec).await?;
}
Ok(())
}
async fn form_answer(
&self,
cx: &mut StepCtx<'_>,
declaration: Option<&crate::manifest::MemoryFormation>,
answer: Tainted<Value>,
input: &Tainted<Value>,
role: &ModelRole,
) -> Result<(), SkillError> {
let Some(declaration) = declaration else {
return Ok(());
};
let subject = resolve_subject(cx, "memory.formation", &declaration.subject, input)?;
let formation_role = match cx.manifest() {
Some(m) => untrusted_contact_model(m, role),
None => role.clone(),
};
let expires_at = if let Some(seconds) = declaration.retention_seconds {
let now = cx.now().await?;
Some(crate::core::seconds_after(now, seconds).ok_or_else(|| {
SkillError::Other(format!(
"a retention window of {seconds}s from {} is past the last instant \
this runtime can name",
crate::core::format_timestamp(now)
))
})?)
} else {
None
};
cx.form_memories(
crate::memory::Formation {
subject,
purpose: declaration.purpose.clone(),
instruction: declaration.instruction.clone(),
max_items: declaration.max_items,
expires_at,
access_retention_seconds: declaration.access_retention_seconds,
max_sensitivity: declaration.max_sensitivity,
},
answer,
Arc::clone(&self.provider),
formation_role,
)
.await?;
Ok(())
}
async fn recall_into(
&self,
cx: &mut StepCtx<'_>,
memory: Option<&crate::manifest::Memory>,
input: &Tainted<Value>,
) -> Result<Option<Tainted<Value>>, SkillError> {
let Some(declaration) = memory.and_then(|m| m.recall.as_ref()) else {
return Ok(None);
};
let subject = resolve_subject(cx, "memory.recall", &declaration.subject, input)?;
let mut query = crate::memory::Recall::about(subject).limit(declaration.limit);
if let Some(purpose) = &declaration.purpose {
query = query.for_purpose(purpose.clone());
}
if declaration.refresh_access {
query = query.refresh_access();
}
let recalled = cx.recall(query).await?;
Ok(Some(Tainted::array(recalled.into_iter().map(|item| {
item.map(|item| {
json!({
"id": item.id,
"purpose": item.purpose,
"content": item.content,
"written_at": crate::core::format_timestamp(item.created_at),
})
})
}))))
}
#[allow(clippy::too_many_arguments, clippy::too_many_lines)]
async fn planned(
&self,
cx: &mut StepCtx<'_>,
input: Tainted<Value>,
system: String,
role: ModelRole,
egress: Option<crate::core::Sensitivity>,
granted: Vec<crate::manifest::ToolGrant>,
oversight: Option<Proposal>,
formation: Option<crate::manifest::MemoryFormation>,
output_schema: Option<Value>,
) -> Result<Outcome, SkillError> {
if input.label().trust != crate::core::Trust::Trusted {
return Err(SkillError::Other(
"a `planned` agent refuses untrusted input: the plan is compiled from \
what the planner reads, and untrusted input authoring a plan is the \
attacker choosing the control flow. Hand hostile content to this \
agent through a tool or a parse step, or use `tool-calling`"
.into(),
));
}
let tools = match (granted.is_empty(), self.tools.clone()) {
(true, _) => None,
(false, Some(wired)) => Some(wired),
(false, None) => {
return Err(SkillError::Other(
"this agent declares `planned` with tool grants but the plane has \
no tool catalogue — `RuntimeBuilder::tools` is what lets a \
declarative agent reach one"
.into(),
));
}
};
let (offered, declared) = tools
.as_ref()
.map(|(catalog, _)| offered_tools(catalog, &granted))
.unwrap_or_default();
let surface: Vec<Value> = declared
.iter()
.map(|t| {
json!({
"tool": t.name,
"description": t.description,
"parameters": t.parameters,
})
})
.collect();
let prompt = Tainted::object([
("system".to_owned(), Tainted::trusted(json!(system))),
("input".to_owned(), input.clone()),
("tools".to_owned(), Tainted::trusted(json!(surface))),
]);
let completion = cx
.sink_with(&prompt, |value| {
let mut call =
ModelCall::new(Arc::clone(&self.provider), role.model.clone(), value)
.expecting(plan_schema(self.max_turns));
call = role.applied_to(call);
if let Some(ceiling) = egress {
call = call.with_max_sensitivity(ceiling);
}
call
})
.await?;
let plan_label = completion.label().join(prompt.label()).clone();
let Some(plan_value) = completion.peek().structured.clone() else {
return Ok(Outcome::fail("the planner returned no structured plan"));
};
let plan: PlanDoc = match serde_json::from_value(plan_value) {
Ok(plan) => plan,
Err(e) => {
return Ok(Outcome::fail(format!(
"the planner's output is not a plan: {e}"
)));
}
};
if plan.steps.is_empty() || plan.steps.len() > self.max_turns as usize {
return Ok(Outcome::fail(format!(
"the plan has {} steps and this agent is bounded to {}",
plan.steps.len(),
self.max_turns
)));
}
let mut outputs: Vec<Tainted<Value>> = Vec::new();
for (index, step) in plan.steps.iter().enumerate() {
match (&step.tool, &step.parse) {
(Some(name), None) => {
let Some((catalog, client)) = tools.as_ref() else {
return Ok(Outcome::fail(format!(
"plan step {index} calls '{name}' but this agent grants no tools"
)));
};
let Some((id, grant)) = catalog.resolve(name).and_then(|id| {
offered
.iter()
.find(|(offered, _)| *offered == id)
.map(|(_, grant)| (id, *grant))
}) else {
let near_miss = offered.iter().find_map(|(id, _)| {
(id.reference() == *name
|| format!("{}__{}", id.server, id.tool) == *name)
.then(|| (id.reference(), id.wire_name()))
});
return Ok(Outcome::fail(match near_miss {
Some((reference, wire)) => format!(
"plan step {index} calls '{name}', which is not granted to \
this agent — but '{reference}' is, and a plan step names a \
tool by its wire name: did you mean '{wire}'?"
),
None => format!(
"plan step {index} calls '{name}', which is not granted to \
this agent"
),
}));
};
let Some(declaration) = declared.iter().find(|tool| &tool.name == name) else {
return Ok(Outcome::fail(format!(
"plan step {index}: '{name}' has no model-facing declaration"
)));
};
let args = match step_arguments(step.args.as_ref()) {
Ok(args) => args,
Err(why) => {
return Ok(Outcome::fail(format!("plan step {index}: {why}")));
}
};
let mut assembled = match assemble_arguments(
&Value::Object(args),
&plan_label,
&input,
&outputs,
) {
Ok(assembled) => assembled,
Err(why) => {
return Ok(Outcome::fail(format!("plan step {index}: {why}")));
}
};
if let Err(detail) =
crate::model::validate_schema(&declaration.parameters, assembled.peek())
{
return Ok(Outcome::fail(format!("plan step {index}: {detail}")));
}
let reference = id.reference();
if grant.requires_approval {
let Some(spec) = oversight.as_ref() else {
return Err(SkillError::Other(
"a tool grant requires approval but the agent declares no \
oversight policy — there is nobody to ask"
.into(),
));
};
cx.deadline(spec.deadline.name.clone(), &spec.deadline.spec(), None)
.await?;
let mut task = spec.approve_call(&reference, assembled.peek());
if let Some(preview) = grant.preview.as_deref() {
let evidence =
preview_evidence(cx, catalog, client, preview, &assembled).await;
task.justification.evidence.push(evidence);
}
let decision = cx.task(&task).await?;
if !decision.approved {
return Ok(Outcome::fail(format!(
"{} refused the call to {reference}: {}",
decision.decided, decision.reason
)));
}
assembled =
match approved_arguments(&decision, &declaration.parameters, assembled)
{
Ok(assembled) => assembled,
Err(detail) => {
return Ok(Outcome::fail(format!(
"plan step {index}: {detail}"
)));
}
};
}
let out = if id.server == crate::tools::AGENT_SERVER {
cx.commission(&id.tool, assembled).await?
} else if let Some(peer) = cx.peer_named(&id.server) {
cx.call_peer(&peer, &id.tool, &assembled).await?
} else {
let egress = cx.tool_egress();
cx.sink_with(&assembled, |value| {
crate::tools::ToolCall::prepare(
catalog,
Arc::clone(client),
id,
value,
egress.as_deref(),
)
.map_err(|e| {
crate::core::StepError::Effect(crate::core::EffectError::Rejected(
e.to_string(),
))
})
})
.await?
};
outputs.push(out);
}
(None, Some(parse)) => {
match step_arguments(step.args.as_ref()) {
Ok(args) if args.is_empty() => {}
_ => {
return Ok(Outcome::fail(format!(
"plan step {index} is a parse and carries `args` — a parse \
takes `from` and `schema`, and arguments nothing executes \
would be accepted prose"
)));
}
}
let source = match resolve_reference(&parse.from, &input, &outputs) {
Ok(source) => source,
Err(why) => {
return Ok(Outcome::fail(format!("plan step {index}: {why}")));
}
};
let declared = serde_json::from_str::<Value>(&parse.schema);
let Some(schema) = declared.ok().as_ref().and_then(bounded_parse_schema) else {
return Ok(Outcome::fail(format!(
"plan step {index}: a parse schema must be JSON text describing an \
object schema"
)));
};
let parse_role = match cx.manifest() {
Some(m) => untrusted_contact_model(m, &role),
None => role.clone(),
};
let prompt = Tainted::object([
(
"system".to_owned(),
Tainted::trusted(json!(PARSE_INSTRUCTION)),
),
("source".to_owned(), source.clone()),
]);
let completion = cx
.sink_with(&prompt, |value| {
let mut call = ModelCall::new(
Arc::clone(&self.provider),
parse_role.model.clone(),
value,
)
.expecting(schema);
call = parse_role.applied_to(call);
if let Some(ceiling) = egress {
call = call.with_max_sensitivity(ceiling);
}
call
})
.await?;
let label = completion.label().join(prompt.label()).clone();
let Some(mut value) = completion.peek().structured.clone() else {
return Ok(Outcome::fail(format!(
"plan step {index}: the parse returned nothing structured"
)));
};
let enough = value
.get("have_enough_information")
.and_then(Value::as_bool)
.unwrap_or(false);
if !enough {
return Ok(Outcome::fail(format!(
"plan step {index}: the parse declared the source does not \
contain enough information — the plan must hand it more of \
the source, not let a guess stand"
)));
}
if let Some(map) = value.as_object_mut() {
map.remove("have_enough_information");
}
outputs.push(Tainted::with_label(value, label));
}
_ => {
return Ok(Outcome::fail(format!(
"plan step {index} must name exactly one of `tool` or `parse`"
)));
}
}
}
let answer = match plan.answer.as_deref() {
Some(reference) => match resolve_reference(reference, &input, &outputs) {
Ok(answer) => answer,
Err(why) => return Ok(Outcome::fail(format!("plan answer: {why}"))),
},
None => match outputs.last() {
Some(last) => last.clone(),
None => return Ok(Outcome::fail("the plan produced nothing to answer with")),
},
};
if let Some(schema) = output_schema
&& let Err(detail) = crate::model::validate_schema(&schema, answer.peek())
{
return Ok(Outcome::fail(format!(
"the answer does not satisfy the declared output shape: {detail}"
)));
}
let formed_source = answer.clone();
self.settle(
cx,
answer,
formed_source,
&input,
oversight,
formation.as_ref(),
&role,
)
.await
}
}
#[allow(clippy::too_many_lines)]
#[async_trait]
impl Skill for Declarative {
fn descriptor(&self) -> SkillDescriptor {
SkillDescriptor::new(self.name.clone())
.provides(crate::core::Capability::new(self.capability.clone()))
}
async fn invoke(
&self,
cx: &mut StepCtx<'_>,
input: Tainted<Value>,
) -> Result<Outcome, SkillError> {
let (system, role, schema, egress, oversight, granted, memory) = {
let m = cx.manifest().ok_or_else(|| {
SkillError::Other(
"a declarative agent ran without a manifest — it has nothing to be".into(),
)
})?;
let role = privileged(m).ok_or_else(|| {
SkillError::Other(format!(
"manifest '{}' declares execution but no privileged model — a \
declarative agent has nothing to call",
m.metadata.name
))
})?;
(
m.spec
.identity
.as_ref()
.map(Identity::system_prompt)
.unwrap_or_default(),
role,
m.output_schema().cloned(),
m.spec.security.max_sensitivity_egress,
m.spec.oversight.as_ref().map(Proposal::from_manifest),
m.spec.tools.clone(),
m.spec.memory.clone(),
)
};
let formation = memory.as_ref().and_then(|m| m.formation.clone());
let remembered = match self.kind {
ExecutionKind::Planned => None,
_ => self.recall_into(cx, memory.as_ref(), &input).await?,
};
match self.kind {
ExecutionKind::Completion => {
let bindable_input = input.clone();
let prompt = prompt_object(&system, input, remembered);
let completion = cx
.sink_with(&prompt, |value| {
let mut call =
ModelCall::new(Arc::clone(&self.provider), role.model.clone(), value);
call = role.applied_to(call);
if let Some(schema) = schema {
call = call.expecting(schema);
}
if let Some(ceiling) = egress {
call = call.with_max_sensitivity(ceiling);
}
call
})
.await?;
let label = completion.label().join(prompt.label());
let completion = Tainted::with_label(completion.into_unlabelled(), label);
let formed_source = Tainted::with_label(
completion
.peek()
.structured
.clone()
.unwrap_or_else(|| json!({ "text": completion.peek().text.clone() })),
completion.label().clone(),
);
let answer =
completion.map(|c| c.structured.unwrap_or_else(|| json!({ "text": c.text })));
self.settle(
cx,
answer,
formed_source,
&bindable_input,
oversight,
formation.as_ref(),
&role,
)
.await
}
ExecutionKind::ToolCalling => {
self.tool_loop(
cx, input, remembered, system, role, egress, granted, oversight, formation,
schema,
)
.await
}
ExecutionKind::Planned => {
self.planned(
cx, input, system, role, egress, granted, oversight, formation, schema,
)
.await
}
}
}
}
#[derive(Debug, Clone)]
struct Proposal {
approval: crate::manifest::Approval,
approvers: Vec<String>,
deadline: crate::manifest::OversightDeadline,
on_expiry: crate::core::Expiry,
triage: Vec<crate::manifest::TriageRule>,
}
impl Proposal {
fn from_manifest(o: &crate::manifest::Oversight) -> Self {
use crate::manifest::Expiry;
Self {
approval: o.approval,
approvers: o.approvers.clone(),
deadline: o.deadline.clone(),
on_expiry: match o.on_expiry {
Expiry::Escalate => crate::core::Expiry::escalate_to(o.escalate_to.iter().cloned()),
Expiry::Proceed if o.allow_unattended => crate::core::Expiry::ProceedUnattended,
Expiry::Deny | Expiry::Proceed => crate::core::Expiry::Deny,
},
triage: o.triage.clone(),
}
}
const fn gates_the_answer(&self) -> bool {
matches!(self.approval, crate::manifest::Approval::Required)
}
fn approve_answer(&self, answer: Value) -> crate::core::TaskSpec {
self.task("agent.approve", "approve this agent's answer", answer)
}
fn approve_call(&self, reference: &str, arguments: &Value) -> crate::core::TaskSpec {
self.task(
APPROVE_CALL_KIND,
format!("approve this agent's call to {reference}"),
json!({ "tool": reference, "arguments": arguments }),
)
}
fn task(&self, kind: &str, summary: impl Into<String>, action: Value) -> crate::core::TaskSpec {
let mut spec = crate::core::TaskSpec::new(
kind,
crate::core::Justification::new(Tainted::trusted(summary.into()), action),
self.deadline.name.clone(),
);
spec.candidate_roles.clone_from(&self.approvers);
spec.on_expiry = self.on_expiry.clone();
spec
}
}
pub(crate) const APPROVE_CALL_KIND: &str = "agent.approve_call";
async fn preview_evidence(
cx: &mut StepCtx<'_>,
catalog: &Arc<crate::tools::ToolCatalog>,
client: &Arc<dyn crate::tools::ToolClient>,
preview: &str,
args: &Tainted<Value>,
) -> Tainted<String> {
let Some(id) = crate::tools::ToolId::parse(preview) else {
return Tainted::trusted(format!(
"preview '{preview}' is not a tool reference, so none was computed"
));
};
let dispatched = if id.server == crate::tools::AGENT_SERVER {
cx.commission(&id.tool, args.clone()).await
} else if let Some(peer) = cx.peer_named(&id.server) {
cx.call_peer(&peer, &id.tool, args).await
} else {
let egress = cx.tool_egress();
cx.sink_with(args, |value| {
crate::tools::ToolCall::prepare(
catalog,
Arc::clone(client),
id.clone(),
value,
egress.as_deref(),
)
.map_err(|e| {
crate::core::StepError::Effect(crate::core::EffectError::Rejected(e.to_string()))
})
})
.await
};
match dispatched {
Ok(answer) => {
let rendered = bounded_evidence(preview, &answer.peek().to_string());
answer.map(|_| rendered)
}
Err(why) => Tainted::trusted(format!(
"preview from {preview} could not be produced ({why}) — this task is being \
decided without it"
)),
}
}
pub const PREVIEW_EVIDENCE_BYTES: usize = 64 * 1024;
fn bounded_evidence(preview: &str, rendered: &str) -> String {
if rendered.len() <= PREVIEW_EVIDENCE_BYTES {
return format!("preview from {preview}: {rendered}");
}
let mut cut = PREVIEW_EVIDENCE_BYTES;
while !rendered.is_char_boundary(cut) {
cut -= 1;
}
format!(
"preview from {preview}: {} … [truncated: showing {cut} of {} bytes; the journaled \
effect output holds the whole answer, sha256 {}]",
&rendered[..cut],
rendered.len(),
crate::core::Digest::of(rendered.as_bytes())
)
}
fn approved_arguments(
decision: &crate::core::Decision,
parameters: &Value,
original: Tainted<Value>,
) -> Result<Tainted<Value>, String> {
if decision.amendment.is_null() {
return Ok(original);
}
crate::model::validate_schema(parameters, &decision.amendment).map_err(|detail| {
format!("the reviewer's amendment does not fit the tool's declared arguments: {detail}")
})?;
let mut label = crate::core::Label::trusted();
label.provenance.insert(crate::core::SourceId::new(format!(
"task:{APPROVE_CALL_KIND}"
)));
label.sensitivity = original.label().sensitivity;
Ok(Tainted::with_label(decision.amendment.clone(), label))
}
const PARSE_INSTRUCTION: &str = "Extract the requested fields from the source. Record only \
what the source literally states: do not infer or invent email addresses, dates, \
identifiers, names or amounts that are not present. If the source does not contain \
enough information, set `have_enough_information` to false and every other field to \
an empty or zero value.";
fn plan_schema(max_steps: u32) -> Value {
json!({
"type": "object",
"additionalProperties": false,
"required": ["steps", "answer"],
"properties": {
"steps": {
"type": "array",
"minItems": 1,
"maxItems": max_steps,
"items": {
"type": "object",
"additionalProperties": false,
"required": ["tool", "args", "parse"],
"properties": {
"tool": {
"type": ["string", "null"],
"description": "a granted tool to call, named exactly as \
offered; null on a parse step"
},
"args": {
"type": ["string", "null"],
"description": "the tool's arguments as a JSON object, \
written as text, e.g. {\"id\": \"$input/customer\"}. \
A string beginning with '$' is a reference to earlier \
data, not a literal: '$input' is the run's input and \
'$step0' is the first step's output, and a JSON Pointer \
may follow the head, e.g. '$step1/customer/email'. \
Escape a literal leading '$' as '$$'. Prefer references \
over copying values: a reference carries the data's \
provenance, a copy does not. Null on a parse step"
},
"parse": {
"type": ["object", "null"],
"additionalProperties": false,
"required": ["from", "schema"],
"properties": {
"from": {
"type": "string",
"description": "reference to the value to extract from, \
e.g. '$step0/body'"
},
"schema": {
"type": "string",
"description": "a JSON Schema of type 'object' naming \
the fields to extract, written as text, e.g. \
{\"type\": \"object\", \"properties\": \
{\"order\": {\"type\": \"string\"}}}"
}
},
"description": "extract structured fields from a prior output \
instead of calling a tool; null on a tool step"
}
}
}
},
"answer": {
"type": ["string", "null"],
"description": "reference selecting the run's answer, e.g. '$step1/summary'; \
null means the last step's output"
}
}
})
}
#[derive(Debug, serde::Deserialize)]
#[serde(deny_unknown_fields)]
struct PlanDoc {
steps: Vec<PlanStep>,
#[serde(default)]
answer: Option<String>,
}
#[derive(Debug, serde::Deserialize)]
#[serde(deny_unknown_fields)]
struct PlanStep {
#[serde(default)]
tool: Option<String>,
#[serde(default)]
args: Option<String>,
#[serde(default)]
parse: Option<ParseStep>,
}
#[derive(Debug, serde::Deserialize)]
#[serde(deny_unknown_fields)]
struct ParseStep {
from: String,
schema: String,
}
fn step_arguments(args: Option<&String>) -> Result<serde_json::Map<String, Value>, String> {
let Some(text) = args else {
return Ok(serde_json::Map::new());
};
match serde_json::from_str::<Value>(text) {
Ok(Value::Object(map)) => Ok(map),
Ok(_) => Err("`args` is JSON, but not a JSON object".to_owned()),
Err(e) => Err(format!("`args` is not JSON: {e}")),
}
}
fn resolve_reference(
reference: &str,
input: &Tainted<Value>,
outputs: &[Tainted<Value>],
) -> Result<Tainted<Value>, String> {
let (head, pointer) = match reference.find('/') {
Some(split) => (&reference[..split], &reference[split..]),
None => (reference, ""),
};
let base = if head == "$input" {
input
} else {
let step = head
.strip_prefix("$step")
.and_then(|n| n.parse::<usize>().ok())
.ok_or_else(|| {
format!(
"'{reference}' is not a reference this plan can hold — use $input \
or $step<N>"
)
})?;
outputs
.get(step)
.ok_or_else(|| format!("'{reference}' points at a step that has not run yet"))?
};
base.project_pointer(pointer)
.ok_or_else(|| format!("'{reference}' selects nothing in the value it points at"))
}
fn assemble_arguments(
value: &Value,
plan_label: &crate::core::Label,
input: &Tainted<Value>,
outputs: &[Tainted<Value>],
) -> Result<Tainted<Value>, String> {
match value {
Value::String(s) if s.starts_with("$$") => Ok(Tainted::with_label(
Value::String(s[1..].to_owned()),
plan_label.clone(),
)),
Value::String(s) if s.starts_with('$') => resolve_reference(s, input, outputs),
Value::Object(map) => {
let mut fields = Vec::with_capacity(map.len());
for (name, nested) in map {
fields.push((
name.clone(),
assemble_arguments(nested, plan_label, input, outputs)?,
));
}
Ok(Tainted::object(fields))
}
Value::Array(items) => {
let mut elements = Vec::with_capacity(items.len());
for nested in items {
elements.push(assemble_arguments(nested, plan_label, input, outputs)?);
}
Ok(Tainted::array(elements))
}
other => Ok(Tainted::with_label(other.clone(), plan_label.clone())),
}
}
fn bounded_parse_schema(declared: &Value) -> Option<Value> {
if declared.get("type") != Some(&json!("object")) {
return None;
}
let mut schema = declared.clone();
let map = schema.as_object_mut()?;
map.insert("additionalProperties".to_owned(), json!(false));
let properties = map
.entry("properties")
.or_insert_with(|| json!({}))
.as_object_mut()?;
properties.insert(
"have_enough_information".to_owned(),
json!({
"type": "boolean",
"description": "Whether the source provided enough information. Set false \
rather than inventing any value."
}),
);
close(&mut schema);
Some(schema)
}
fn close(node: &mut Value) {
if let Some(list) = node.get_mut("properties").and_then(Value::as_object_mut) {
let names: Vec<Value> = list.keys().map(|k| json!(k)).collect();
for nested in list.values_mut() {
close(nested);
}
let is_object = match node.get("type") {
Some(Value::String(name)) => name == "object",
Some(Value::Array(names)) => names.iter().any(|n| n == "object"),
_ => false,
};
if is_object {
let map = node.as_object_mut().expect("an object schema is an object");
map.insert("additionalProperties".to_owned(), json!(false));
map.insert("required".to_owned(), Value::Array(names));
}
}
if let Some(items) = node.get_mut("items") {
close(items);
}
}
fn prompt_object(
system: &str,
input: Tainted<Value>,
remembered: Option<Tainted<Value>>,
) -> Tainted<Value> {
let mut parts = vec![
("system".to_owned(), Tainted::trusted(json!(system))),
("input".to_owned(), input),
];
if let Some(remembered) = remembered {
parts.push(("memory".to_owned(), remembered));
}
Tainted::object(parts)
}
fn resolve_subject(
cx: &StepCtx<'_>,
field: &str,
subject: &crate::manifest::MemorySubject,
input: &Tainted<Value>,
) -> Result<String, SkillError> {
use crate::manifest::MemorySubject;
let refuse = SkillError::Other;
match subject {
MemorySubject::Literal(literal) => Ok(literal.clone()),
MemorySubject::Case => cx.case_id().map(|id| id.to_string()).ok_or_else(|| {
refuse(format!(
"`{field}.subject: $case` needs the run to belong to a case, and \
this run has none. Admit it with `run_correlated(..)` or `run_in_case(..)` \
— filing the memory under a constant instead would pool every matter's \
facts under one key"
))
}),
MemorySubject::Correlation(namespace) => cx
.correlation_value(namespace)
.map(ToOwned::to_owned)
.ok_or_else(|| {
let held: Vec<&str> = cx
.correlation()
.iter()
.map(|key| key.namespace.as_str())
.collect();
refuse(format!(
"`{field}.subject: $correlation/{namespace}` found no key in \
that namespace; this run's case is keyed by {held:?}. A memory filed \
under the wrong scope is recalled into another subject's run and \
survives that subject's erasure, so the run fails rather than \
guessing"
))
}),
MemorySubject::Input(pointer) => {
let selected = input.project_pointer(pointer).ok_or_else(|| {
refuse(format!(
"`{field}.subject: $input{pointer}` selects nothing in this \
run's input"
))
})?;
if selected.label().trust != crate::core::Trust::Trusted {
return Err(refuse(format!(
"`{field}.subject: $input{pointer}` names an **untrusted** \
field, and a subject taken from untrusted input lets whoever supplied \
it choose whose memories this run writes into. Bind the subject to a \
correlation key instead — those are settled at admission by a \
deterministic lookup — or release the field explicitly if it really is \
the plane's own"
)));
}
match selected.peek() {
Value::String(value) if !value.trim().is_empty() => Ok(value.clone()),
Value::Number(value) => Ok(value.to_string()),
other => Err(refuse(format!(
"`{field}.subject: $input{pointer}` selects {}, and a subject \
is a scope name — it must be a non-empty string or a number",
if other.is_string() {
"an empty string"
} else {
crate::core::canon::json_kind(other)
}
))),
}
}
}
}
fn untrusted_contact_model(m: &Manifest, fallback: &ModelRole) -> ModelRole {
m.quarantined_role().unwrap_or_else(|| fallback.clone())
}
fn offered_tools<'g>(
catalog: &crate::tools::ToolCatalog,
granted: &'g [crate::manifest::ToolGrant],
) -> (
Vec<(crate::tools::ToolId, &'g crate::manifest::ToolGrant)>,
Vec<crate::model::ToolDeclaration>,
) {
let offered: Vec<(crate::tools::ToolId, &crate::manifest::ToolGrant)> = granted
.iter()
.filter_map(|g| catalog.resolve_reference(&g.reference).map(|id| (id, g)))
.collect();
let declared: Vec<crate::model::ToolDeclaration> = offered
.iter()
.map(|(id, grant)| {
let (description, arguments) = catalog.declaration(id).map_or_else(
|| {
(
grant.description.clone().unwrap_or_default(),
grant
.arguments
.clone()
.unwrap_or_else(|| json!({ "type": "object" })),
)
},
|(description, arguments)| (description.to_owned(), arguments.clone()),
);
crate::model::ToolDeclaration::new(id.wire_name(), description, arguments)
})
.collect();
(offered, declared)
}
fn privileged(m: &Manifest) -> Option<ModelRole> {
m.privileged_role()
}
fn model_facing(e: &crate::core::StepError) -> Option<String> {
match e {
crate::core::StepError::Policy(p) => Some(p.for_model().to_owned()),
crate::core::StepError::Effect(inner)
if inner.disposition() != crate::core::Disposition::InDoubt =>
{
Some(inner.to_string())
}
_ => None,
}
}
fn relayed_failure(
cx: &mut StepCtx<'_>,
conversation_label: &mut crate::core::Label,
e: &crate::core::StepError,
) -> Option<String> {
let detail = model_facing(e)?;
if matches!(e, crate::core::StepError::Effect(_))
&& let Some(failed) = cx.take_failed_output_label()
{
*conversation_label = conversation_label.join(&failed);
}
Some(detail)
}
#[cfg(all(test, feature = "providers"))]
mod tests {
use super::*;
#[test]
fn the_plan_format_survives_constrained_decoding() {
assert_eq!(crate::model::strict_schema_problem(&plan_schema(4)), None);
}
#[test]
fn a_loosely_written_parse_schema_is_closed_rather_than_refused() {
let loose = json!({
"type": "object",
"properties": {
"order": { "type": "string" },
"line": {
"type": "object",
"properties": { "sku": { "type": "string" } }
}
}
});
let bounded = bounded_parse_schema(&loose).expect("an object schema is bounded");
assert_eq!(crate::model::strict_schema_problem(&bounded), None);
let mut required: Vec<&str> = bounded["required"]
.as_array()
.expect("required is a list")
.iter()
.filter_map(Value::as_str)
.collect();
required.sort_unstable();
assert_eq!(
required,
["have_enough_information", "line", "order"],
"a field the planner named must be answered — absence is what the \
escape bit is for, not an omitted key"
);
}
}