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, ModelId, ModelProvider};
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>,
system: String,
model_role: (ModelId, Option<u32>, Option<crate::model::ReasoningEffort>),
egress: Option<crate::core::Sensitivity>,
granted: Vec<crate::manifest::ToolGrant>,
formation: Option<crate::manifest::MemoryFormation>,
) -> Result<Outcome, SkillError> {
let (model, max_output_tokens, reasoning_effort) = model_role;
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: 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();
let prompt = Tainted::object([
("system".to_owned(), Tainted::trusted(json!(system))),
("input".to_owned(), input),
]);
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 mut call = ModelCall::new(
Arc::clone(&self.provider),
model.clone(),
prompt.peek().clone(),
)
.with_tools(declared.clone())
.continuing(exchanges.clone())
.with_output_sensitivity(conversation_label.sensitivity);
if let Some(state) = continuation.take() {
call = call.with_continuation(state);
}
if let Some(max_output_tokens) = max_output_tokens {
call = call.with_max_output_tokens(max_output_tokens);
}
if let Some(effort) = reasoning_effort {
call = call.with_reasoning_effort(effort);
}
if let Some(ceiling) = egress {
call = call.with_max_sensitivity(ceiling);
}
let outbound = prompt.with_joined_label(&conversation_label);
let completion = cx.sink(call, &outbound).await?;
let label = completion.label().join(&conversation_label);
let completion = Tainted::with_label(completion.into_unlabelled(), label);
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(),
);
self.form_answer(cx, formation.as_ref(), formed_source, &model)
.await?;
let answer =
completion.map(|c| c.structured.unwrap_or_else(|| json!({ "text": c.text })));
return Ok(Outcome::done(answer));
}
continuation.clone_from(&completion.peek().continuation);
exchanges.clear();
for asked in completion.peek().tool_calls.clone() {
let Some(id) = catalog
.resolve(&asked.name)
.filter(|id| offered.iter().any(|(o, _)| o == id))
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 prepared = crate::tools::ToolCall::prepare(
&catalog,
Arc::clone(&client),
id,
asked.arguments.clone(),
)
.map_err(|e| SkillError::Other(e.to_string()))?;
let args = crate::core::Tainted::with_label(
asked.arguments.clone(),
completion.label().clone(),
);
match cx.sink(prepared, &args).await {
Ok(result) => {
conversation_label = conversation_label.join(result.label());
exchanges
.push(crate::model::ToolExchange::ok(asked, result.peek().clone()));
}
Err(e) => match model_facing(&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
)))
}
async fn form_answer(
&self,
cx: &mut StepCtx<'_>,
declaration: Option<&crate::manifest::MemoryFormation>,
answer: Tainted<Value>,
model: &ModelId,
) -> Result<(), SkillError> {
let Some(declaration) = declaration else {
return Ok(());
};
let expires_at = if let Some(seconds) = declaration.retention_seconds {
let now = cx.now().await?;
Some(now + time::Duration::seconds(i64::try_from(seconds).unwrap_or(i64::MAX)))
} else {
None
};
cx.form_memories(
crate::memory::Formation {
subject: declaration.subject.clone(),
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),
model.clone(),
)
.await?;
Ok(())
}
}
#[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,
model,
max_output_tokens,
reasoning_effort,
schema,
egress,
oversight,
granted,
formation,
) = {
let m = cx.manifest().ok_or_else(|| {
SkillError::Other(
"a declarative agent ran without a manifest — it has nothing to be".into(),
)
})?;
let (model, max_output_tokens, reasoning_effort) = 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(),
model,
max_output_tokens,
reasoning_effort,
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_formation.clone(),
)
};
match self.kind {
ExecutionKind::Completion => {
let prompt = Tainted::object([
("system".to_owned(), Tainted::trusted(json!(system))),
("input".to_owned(), input),
]);
let mut call = ModelCall::new(
Arc::clone(&self.provider),
model.clone(),
prompt.peek().clone(),
)
.with_output_sensitivity(prompt.label().sensitivity);
if let Some(max_output_tokens) = max_output_tokens {
call = call.with_max_output_tokens(max_output_tokens);
}
if let Some(effort) = reasoning_effort {
call = call.with_reasoning_effort(effort);
}
if let Some(schema) = schema {
call = call.expecting(schema);
}
if let Some(ceiling) = egress {
call = call.with_max_sensitivity(ceiling);
}
let completion = cx.sink(call, &prompt).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(),
);
self.form_answer(cx, formation.as_ref(), formed_source, &model)
.await?;
let answer =
completion.map(|c| c.structured.unwrap_or_else(|| json!({ "text": c.text })));
let Some(spec) = oversight else {
return Ok(Outcome::done(answer));
};
let decision = cx.task(&spec.with_action(answer.peek().clone())).await?;
if decision.approved {
Ok(Outcome::done(answer))
} else {
Ok(Outcome::fail(format!(
"{} refused this answer: {}",
decision.actor, decision.reason
)))
}
}
ExecutionKind::ToolCalling => {
self.tool_loop(
cx,
input,
system,
(model, max_output_tokens, reasoning_effort),
egress,
granted,
formation,
)
.await
}
}
}
}
#[derive(Debug, Clone)]
struct Proposal {
approvers: Vec<String>,
deadline: String,
on_expiry: crate::core::OnExpiry,
allow_unattended: bool,
}
impl Proposal {
fn from_manifest(o: &crate::manifest::Oversight) -> Self {
use crate::manifest::Expiry;
Self {
approvers: o.approvers.clone(),
deadline: o.deadline.clone(),
on_expiry: match o.on_expiry {
Expiry::Deny => crate::core::OnExpiry::Deny,
Expiry::Escalate => crate::core::OnExpiry::Escalate,
Expiry::Proceed => crate::core::OnExpiry::Proceed,
},
allow_unattended: o.allow_unattended,
}
}
fn with_action(&self, action: Value) -> crate::core::TaskSpec {
let mut spec = crate::core::TaskSpec::new(
"agent.approve",
crate::core::Justification::new("approve this agent's answer", action),
self.deadline.clone(),
);
spec.candidate_roles.clone_from(&self.approvers);
spec.on_expiry = self.on_expiry;
spec.allow_unattended = self.allow_unattended;
spec
}
}
fn privileged(
m: &Manifest,
) -> Option<(ModelId, Option<u32>, Option<crate::model::ReasoningEffort>)> {
let r = m.spec.models.as_ref()?.privileged.as_ref()?;
Some((
ModelId::new(&r.provider, &r.model),
r.max_tokens,
r.reasoning_effort,
))
}
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,
}
}