#![allow(clippy::result_large_err)]
use std::marker::PhantomData;
use crate::reasoning::circuit_breaker::CircuitBreakerRegistry;
use crate::reasoning::context_manager::ContextManager;
use crate::reasoning::conversation::Conversation;
use crate::reasoning::executor::ActionExecutor;
use crate::reasoning::inference::{InferenceProvider, ToolDefinition};
use crate::reasoning::loop_types::*;
use crate::reasoning::policy_bridge::ReasoningPolicyGate;
use crate::reasoning::prepared::AuthorizedAction;
pub struct Reasoning;
pub struct PolicyCheck;
pub struct ToolDispatching;
pub struct Observing;
pub trait AgentPhase {}
impl AgentPhase for Reasoning {}
impl AgentPhase for PolicyCheck {}
impl AgentPhase for ToolDispatching {}
impl AgentPhase for Observing {}
pub struct ReasoningOutput {
pub proposed_actions: Vec<ProposedAction>,
}
pub struct PolicyOutput {
pub approved_actions: Vec<AuthorizedAction>,
pub denied_reasons: Vec<(ProposedAction, String)>,
pub has_terminal_action: bool,
pub terminal_output: Option<String>,
}
pub struct DispatchOutput {
pub observations: Vec<Observation>,
pub should_terminate: bool,
pub terminal_output: Option<String>,
}
pub struct AgentLoop<Phase: AgentPhase> {
pub state: LoopState,
pub config: LoopConfig,
phase_data: Option<PhaseData>,
_phase: PhantomData<Phase>,
}
enum PhaseData {
Reasoning(ReasoningOutput),
Policy(PolicyOutput),
Dispatch(DispatchOutput),
}
impl AgentLoop<Reasoning> {
pub fn new(state: LoopState, mut config: LoopConfig) -> Self {
config
.shared_budget
.get_or_insert_with(|| super::budget::SharedBudget::new(config.max_total_tokens));
Self {
state,
config,
phase_data: None,
_phase: PhantomData,
}
}
pub async fn produce_output(
mut self,
provider: &dyn InferenceProvider,
context_manager: &dyn ContextManager,
delegation_available: bool,
journal: &dyn JournalWriter,
) -> Result<AgentLoop<PolicyCheck>, LoopTermination> {
self.state.current_phase = "reasoning".into();
if self.state.iteration >= self.config.max_iterations {
return Err(LoopTermination {
reason: LoopTerminationReason::MaxIterations {
iterations: self.state.iteration,
},
state: self.state,
});
}
if self.state.total_usage.total_tokens >= self.config.max_total_tokens {
return Err(LoopTermination {
reason: LoopTerminationReason::MaxTokens {
tokens: self.state.total_usage.total_tokens,
},
state: self.state,
});
}
let before_len = self.state.conversation.len();
let before_tokens = self.state.conversation.estimate_tokens();
context_manager.manage_context(
&mut self.state.conversation,
self.config.context_token_budget,
);
let after_len = self.state.conversation.len();
if after_len != before_len {
tracing::debug!(
iter = self.state.iteration,
before_len,
after_len,
before_tokens,
budget = self.config.context_token_budget,
"context_manager truncated conversation"
);
}
self.state.pending_observations.clear();
let mut options = crate::reasoning::inference::InferenceOptions {
max_tokens: self
.config
.max_total_tokens
.saturating_sub(self.state.total_usage.total_tokens)
.min(self.config.max_output_tokens),
temperature: self.config.temperature,
tool_definitions: self.config.tool_definitions.clone(),
tool_choice: self.config.tool_choice.clone(),
..Default::default()
};
let input_reservation =
match provider.input_token_reservation(&self.state.conversation, &options) {
Ok(tokens) => tokens,
Err(error) => {
return Err(LoopTermination {
reason: LoopTerminationReason::Error {
message: format!("Inference accounting unavailable: {error}"),
},
state: self.state,
})
}
};
let budget = self
.config
.shared_budget
.as_ref()
.expect("initialized loop budget");
let proof = |maximum| {
let mut exact_options = options.clone();
exact_options.max_tokens = maximum;
let request =
serde_json::json!({"conversation":self.state.conversation,"options":exact_options});
let encoded = crate::reasoning::prepared::canonical_json(&request)?;
if encoded.len() > 4 * 1024 * 1024 {
return Err("inference accounting contract exceeds its byte limit".into());
}
use sha2::{Digest, Sha256};
Ok(super::budget::journal::InferenceReservation {
id: uuid::Uuid::nil(),
root_id: uuid::Uuid::nil(),
ancestors: Vec::new(),
input_tokens: 0,
output_tokens: 0,
agent_id: self.state.agent_id,
audit: journal.audit_reference(),
iteration: self.state.iteration,
provider: provider.provider_name().into(),
model: options
.model
.clone()
.unwrap_or_else(|| provider.default_model().into()),
request_hash: format!("sha256:{}", hex::encode(Sha256::digest(encoded.as_bytes()))),
request_bytes: encoded.len() as u64,
})
};
let reservation = match budget
.reserve_audited(input_reservation, options.max_tokens, proof)
.await
{
Ok(reservation) => reservation,
Err(message) => {
return Err(LoopTermination {
reason: if message == "shared token budget exhausted or already reserved" {
LoopTerminationReason::MaxTokens {
tokens: budget.snapshot().usage.total_tokens,
}
} else {
LoopTerminationReason::Error { message }
},
state: self.state,
})
}
};
options.max_tokens = reservation.output_tokens();
let response = match provider.complete(&self.state.conversation, &options).await {
Ok(r) => r,
Err(e) => {
let settlement = reservation
.settle_audited(
&crate::reasoning::inference::Usage::default(),
self.state.iteration,
)
.await;
let message = match settlement {
Ok(()) => format!("Inference failed: {e}"),
Err(settle) => {
format!("Inference failed: {e}; reservation settlement failed: {settle}")
}
};
return Err(LoopTermination {
reason: LoopTerminationReason::Error { message },
state: self.state,
});
}
};
self.state.add_usage(&response.usage);
if let Err(message) = reservation
.settle_audited(&response.usage, self.state.iteration)
.await
{
return Err(LoopTermination {
reason: LoopTerminationReason::Error { message },
state: self.state,
});
}
if response.finish_reason == crate::reasoning::inference::FinishReason::Refusal {
return Err(LoopTermination {
reason: LoopTerminationReason::Error {
message: "model refused the request (stop_reason=refusal)".to_string(),
},
state: self.state,
});
}
if !response.has_tool_calls() && response.content.trim().is_empty() {
return Err(LoopTermination {
reason: LoopTerminationReason::Error {
message: "model produced no tool calls and no text (no-progress turn)"
.to_string(),
},
state: self.state,
});
}
let proposed_actions = if response.has_tool_calls() {
let tool_calls: Vec<crate::reasoning::conversation::ToolCall> = response
.tool_calls
.iter()
.map(|tc| crate::reasoning::conversation::ToolCall {
id: tc.id.clone(),
name: tc.name.clone(),
arguments: tc.arguments.clone(),
})
.collect();
self.state.conversation.push(
crate::reasoning::conversation::ConversationMessage::assistant_tool_calls(
tool_calls,
),
);
response
.tool_calls
.into_iter()
.map(|tc| tool_call_to_action(tc.id, tc.name, tc.arguments, delegation_available))
.collect()
} else {
self.state.conversation.push(
crate::reasoning::conversation::ConversationMessage::assistant(&response.content),
);
vec![ProposedAction::Respond {
content: response.content,
}]
};
self.state.iteration += 1;
Ok(AgentLoop {
state: self.state,
config: self.config,
phase_data: Some(PhaseData::Reasoning(ReasoningOutput { proposed_actions })),
_phase: PhantomData,
})
}
}
impl AgentLoop<PolicyCheck> {
pub fn proposed_actions(&self) -> Vec<ProposedAction> {
match &self.phase_data {
Some(PhaseData::Reasoning(output)) => output.proposed_actions.clone(),
_ => Vec::new(),
}
}
pub async fn check_policy(
mut self,
gate: &dyn ReasoningPolicyGate,
executor: &dyn ActionExecutor,
) -> Result<AgentLoop<ToolDispatching>, LoopTermination> {
self.state.current_phase = "policy_check".into();
let reasoning_output = match self.phase_data {
Some(PhaseData::Reasoning(output)) => output,
_ => {
return Err(LoopTermination {
reason: LoopTerminationReason::Error {
message: "Invalid phase data: expected ReasoningOutput".into(),
},
state: self.state,
});
}
};
let mut approved = Vec::new();
let mut denied = Vec::new();
let mut has_terminal = false;
let mut terminal_output = None;
let valid_ids = super::dispatch::valid_call_identities(&reasoning_output.proposed_actions);
for original in reasoning_output.proposed_actions {
let decision = if valid_ids {
super::dispatch::authorize_action(
&original,
&self.state,
&self.config,
executor,
gate,
)
.await
} else {
Err("action batch requires unique, nonempty call identities".into())
};
let denial = match decision {
Ok(authorized) => {
match authorized.action() {
ProposedAction::Respond { content } => {
has_terminal = true;
terminal_output = Some(content.clone());
}
ProposedAction::Terminate { output, .. } => {
has_terminal = true;
terminal_output = Some(output.clone());
}
_ => {}
}
approved.push(authorized);
continue;
}
Err(reason) => reason,
};
push_denial_tool_result(&mut self.state.conversation, &original, &denial);
self.state
.pending_observations
.push(Observation::policy_denial(&denial));
denied.push((original, denial));
}
Ok(AgentLoop {
state: self.state,
config: self.config,
phase_data: Some(PhaseData::Policy(PolicyOutput {
approved_actions: approved,
denied_reasons: denied,
has_terminal_action: has_terminal,
terminal_output,
})),
_phase: PhantomData,
})
}
}
pub(crate) fn push_denial_tool_result(
conversation: &mut Conversation,
action: &ProposedAction,
reason: &str,
) {
let (call_id, name) = match action {
ProposedAction::ToolCall { call_id, name, .. } => (call_id, name.clone()),
ProposedAction::Delegate {
call_id, target, ..
} => (call_id, format!("delegate:{}", target)),
_ => return,
};
conversation.push(
crate::reasoning::conversation::ConversationMessage::tool_result(
call_id,
&name,
format!("[Policy denied] {}", reason),
),
);
}
pub(crate) const DELEGATE_TOOL_NAME: &str = "delegate";
pub(crate) fn tool_call_to_action(
call_id: String,
name: String,
arguments: String,
delegation_available: bool,
) -> ProposedAction {
if delegation_available && name == DELEGATE_TOOL_NAME {
if let Ok(serde_json::Value::Object(map)) =
serde_json::from_str::<serde_json::Value>(&arguments)
{
if let (Some(agent), Some(task)) = (
map.get("agent").and_then(|v| v.as_str()),
map.get("task").and_then(|v| v.as_str()),
) {
return ProposedAction::Delegate {
call_id,
target: agent.to_string(),
message: task.to_string(),
};
}
}
}
ProposedAction::ToolCall {
call_id,
name,
arguments,
}
}
#[cfg(test)]
pub(crate) async fn dispatch_delegations(
actions: &[ProposedAction],
delegation: Option<&dyn crate::reasoning::delegation::DelegationExecutor>,
config: &LoopConfig,
) -> Vec<Observation> {
use crate::reasoning::delegation::DelegationContext;
let depth = config.delegation_depth;
let chain = &config.delegation_chain;
let mut observations = Vec::new();
for action in actions {
if let ProposedAction::Delegate {
call_id,
target,
message,
} = action
{
let source = format!("delegate:{}", target);
let obs = match delegation {
None => Observation::tool_error(
source,
"agent delegation is not available in this runner",
),
Some(d) => {
let ctx = DelegationContext {
depth,
chain: chain.to_vec(),
max_iterations: config.max_iterations,
max_total_tokens: config.max_total_tokens,
shared_budget: config.shared_budget.clone(),
timeout: config.timeout,
};
match d.delegate(target, message, ctx).await {
Ok(output) => Observation::tool_result(source, output),
Err(e) => Observation::tool_error(source, e.to_string()),
}
}
};
observations.push(obs.with_call_id(call_id.clone()));
}
}
observations
}
impl AgentLoop<ToolDispatching> {
pub fn policy_summary(&self) -> (usize, usize) {
match &self.phase_data {
Some(PhaseData::Policy(output)) => (
output.approved_actions.len() + output.denied_reasons.len(),
output.denied_reasons.len(),
),
_ => (0, 0),
}
}
pub fn approved_calls(&self) -> Vec<serde_json::Value> {
match &self.phase_data {
Some(PhaseData::Policy(output)) => output
.approved_actions
.iter()
.map(|grant| grant.audit_context())
.collect(),
_ => Vec::new(),
}
}
pub fn denied_calls(&self) -> Vec<serde_json::Value> {
match &self.phase_data {
Some(PhaseData::Policy(output)) => output
.denied_reasons
.iter()
.map(|(action, reason)| serde_json::json!({"action": action, "reason": reason}))
.collect(),
_ => Vec::new(),
}
}
pub async fn dispatch_tools(
self,
executor: &dyn ActionExecutor,
circuit_breakers: &CircuitBreakerRegistry,
delegation: Option<&dyn crate::reasoning::delegation::DelegationExecutor>,
) -> Result<AgentLoop<Observing>, LoopTermination> {
self.dispatch_tools_audited(executor, circuit_breakers, delegation, None)
.await
}
pub(crate) async fn dispatch_tools_audited(
mut self,
executor: &dyn ActionExecutor,
circuit_breakers: &CircuitBreakerRegistry,
delegation: Option<&dyn crate::reasoning::delegation::DelegationExecutor>,
journal: Option<&dyn super::loop_types::JournalWriter>,
) -> Result<AgentLoop<Observing>, LoopTermination> {
self.state.current_phase = "tool_dispatching".into();
let policy_output = match self.phase_data {
Some(PhaseData::Policy(output)) => output,
_ => {
return Err(LoopTermination {
reason: LoopTerminationReason::Error {
message: "Invalid phase data: expected PolicyOutput".into(),
},
state: self.state,
});
}
};
if policy_output.has_terminal_action {
for grant in policy_output.approved_actions {
if !matches!(
grant.action(),
ProposedAction::Respond { .. } | ProposedAction::Terminate { .. }
) {
continue;
}
let result = match grant.check_binding(&self.state, &self.config) {
Ok(()) => super::response_delivery::dispatch(grant, journal).await,
Err(error) => Err(error),
};
if let Err(message) = result {
return Err(LoopTermination {
reason: LoopTerminationReason::Error { message },
state: self.state,
});
}
}
return Ok(AgentLoop {
state: self.state,
config: self.config,
phase_data: Some(PhaseData::Dispatch(DispatchOutput {
observations: Vec::new(),
should_terminate: true,
terminal_output: policy_output.terminal_output,
})),
_phase: PhantomData,
});
}
let mut authorized = Vec::new();
let mut observations = Vec::new();
for grant in policy_output.approved_actions {
match grant.check_binding(&self.state, &self.config) {
Ok(()) => authorized.push(grant),
Err(reason) => {
push_denial_tool_result(&mut self.state.conversation, grant.action(), &reason);
observations.push(Observation::policy_denial(&reason));
}
}
}
let (delegates, tools): (Vec<_>, Vec<_>) = authorized
.into_iter()
.partition(|grant| matches!(grant.action(), ProposedAction::Delegate { .. }));
match super::dispatch::execute_tool_grants(
tools,
&self.config,
executor,
circuit_breakers,
journal,
)
.await
{
Ok(results) => observations.extend(results),
Err(error) => {
return Err(LoopTermination {
reason: LoopTerminationReason::Error {
message: format!("required tool effect journal failed: {error}"),
},
state: self.state,
})
}
}
for grant in delegates {
if observations.iter().any(Observation::has_unconfirmed_effect) {
if let ProposedAction::Delegate {
call_id, target, ..
} = grant.action()
{
observations.push(
Observation::tool_error(
format!("delegate:{target}"),
"delegation not started after an unconfirmed effect",
)
.with_call_id(call_id),
);
}
continue;
}
match grant.check_binding(&self.state, &self.config) {
Ok(()) => {
let ProposedAction::Delegate {
call_id, target, ..
} = grant.action()
else {
unreachable!()
};
let ctx = super::delegation::DelegationContext {
depth: self.config.delegation_depth,
chain: self.config.delegation_chain.clone(),
max_iterations: self.config.max_iterations,
max_total_tokens: self.config.max_total_tokens,
shared_budget: self.config.shared_budget.clone(),
timeout: self.config.timeout,
};
let result = match delegation {
Some(delegation) => {
delegation.delegate_authorized(&grant, ctx, journal).await
}
None => Err(super::delegation::DelegationError::Failed(
"agent delegation is not available in this runner".into(),
)),
};
let source = format!("delegate:{target}");
let observation = match result {
Ok(output) => Observation::tool_result(source, output),
Err(super::delegation::DelegationError::Audit(message)) => {
return Err(LoopTermination {
reason: LoopTerminationReason::Error {
message: format!("required delegation audit failed: {message}"),
},
state: self.state,
});
}
Err(super::delegation::DelegationError::Unconfirmed(message)) => {
let mut observation = Observation::tool_error(source, message);
observation.mark_unconfirmed_effect();
observation
}
Err(error) => Observation::tool_error(source, error.to_string()),
};
observations.push(observation.with_call_id(call_id.clone()));
}
Err(reason) => observations.push(Observation::policy_denial(&reason)),
}
}
for obs in &observations {
let tool_call_id = obs.call_id.as_deref().unwrap_or(&obs.source);
if !obs.is_error {
self.state.conversation.push(
crate::reasoning::conversation::ConversationMessage::tool_result(
tool_call_id,
&obs.source,
&obs.content,
),
);
} else {
self.state.conversation.push(
crate::reasoning::conversation::ConversationMessage::tool_result(
tool_call_id,
&obs.source,
format!("[Error] {}", obs.content),
),
);
}
}
Ok(AgentLoop {
state: self.state,
config: self.config,
phase_data: Some(PhaseData::Dispatch(DispatchOutput {
observations,
should_terminate: false,
terminal_output: None,
})),
_phase: PhantomData,
})
}
}
pub enum LoopContinuation {
Continue(Box<AgentLoop<Reasoning>>),
Complete(LoopResult),
}
impl AgentLoop<Observing> {
pub fn observations(&self) -> Vec<Observation> {
match &self.phase_data {
Some(PhaseData::Dispatch(output)) => output.observations.clone(),
_ => Vec::new(),
}
}
pub fn observation_count(&self) -> usize {
match &self.phase_data {
Some(PhaseData::Dispatch(output)) => output.observations.len(),
_ => 0,
}
}
pub fn observe_results(mut self) -> LoopContinuation {
self.state.current_phase = "observing".into();
let dispatch_output = match self.phase_data {
Some(PhaseData::Dispatch(output)) => output,
_ => {
return LoopContinuation::Complete(LoopResult {
output: String::new(),
iterations: self.state.iteration,
total_usage: self.state.total_usage.clone(),
budget: self
.config
.shared_budget
.as_ref()
.map(|budget| budget.snapshot()),
termination_reason: TerminationReason::Error {
message: "Invalid phase data".into(),
},
duration: self.state.elapsed().to_std().unwrap_or_default(),
conversation: self.state.conversation,
});
}
};
if dispatch_output
.observations
.iter()
.any(Observation::has_unconfirmed_effect)
{
return LoopContinuation::Complete(
LoopTermination {
reason: LoopTerminationReason::UnconfirmedEffects,
state: self.state,
}
.into_result(),
);
}
if dispatch_output.should_terminate {
return LoopContinuation::Complete(LoopResult {
output: dispatch_output.terminal_output.unwrap_or_default(),
iterations: self.state.iteration,
total_usage: self.state.total_usage.clone(),
budget: self
.config
.shared_budget
.as_ref()
.map(|budget| budget.snapshot()),
termination_reason: TerminationReason::Completed,
duration: self.state.elapsed().to_std().unwrap_or_default(),
conversation: self.state.conversation,
});
}
self.state
.pending_observations
.extend(dispatch_output.observations);
LoopContinuation::Continue(Box::new(AgentLoop {
state: self.state,
config: self.config,
phase_data: None,
_phase: PhantomData,
}))
}
}
#[derive(Debug)]
pub struct LoopTermination {
pub reason: LoopTerminationReason,
pub state: LoopState,
}
#[derive(Debug)]
pub enum LoopTerminationReason {
MaxIterations { iterations: u32 },
MaxTokens { tokens: u32 },
Timeout,
UnconfirmedEffects,
Error { message: String },
}
impl LoopTermination {
pub fn into_result(self) -> LoopResult {
let reason = match &self.reason {
LoopTerminationReason::MaxIterations { .. } => TerminationReason::MaxIterations,
LoopTerminationReason::MaxTokens { .. } => TerminationReason::MaxTokens,
LoopTerminationReason::Timeout => TerminationReason::Timeout,
LoopTerminationReason::UnconfirmedEffects => TerminationReason::UnconfirmedEffects,
LoopTerminationReason::Error { message } => TerminationReason::Error {
message: message.clone(),
},
};
LoopResult {
output: String::new(),
iterations: self.state.iteration,
total_usage: self.state.total_usage.clone(),
budget: None,
termination_reason: reason,
duration: self.state.elapsed().to_std().unwrap_or_default(),
conversation: self.state.conversation,
}
}
}
pub(crate) fn validate_tool_call_arguments(
name: &str,
arguments: &str,
tool_definitions: &[ToolDefinition],
) -> Result<(), String> {
let parsed: serde_json::Value = serde_json::from_str(arguments)
.map_err(|e| format!("tool '{}' arguments are not valid JSON: {}", name, e))?;
if !parsed.is_object() {
return Err(format!(
"tool '{}' arguments must be a JSON object, got {}",
name,
json_type_of(&parsed)
));
}
let def = tool_definitions
.iter()
.find(|d| d.name == name)
.ok_or_else(|| format!("tool '{name}' was not advertised for this run"))?;
if let Some(properties) = def.parameters.get("properties").and_then(|p| p.as_object()) {
if let Some(unknown) = parsed
.as_object()
.unwrap()
.keys()
.find(|key| !properties.contains_key(*key))
{
return Err(format!("tool '{name}' has unknown argument '{unknown}'"));
}
}
{
if !def.parameters.is_null() {
match jsonschema::validator_for(&def.parameters) {
Ok(validator) => {
let errors: Vec<String> = validator
.iter_errors(&parsed)
.map(|e| {
let path = e.instance_path.to_string();
if path.is_empty() {
e.to_string()
} else {
format!("at '{}': {}", path, e)
}
})
.collect();
if !errors.is_empty() {
return Err(format!(
"tool '{}' arguments failed schema validation: {}",
name,
errors.join("; ")
));
}
}
Err(e) => {
return Err(format!(
"tool '{}' has an invalid declared schema (refusing to dispatch): {}",
name, e
));
}
}
}
}
Ok(())
}
fn json_type_of(v: &serde_json::Value) -> &'static str {
match v {
serde_json::Value::Null => "null",
serde_json::Value::Bool(_) => "boolean",
serde_json::Value::Number(_) => "number",
serde_json::Value::String(_) => "string",
serde_json::Value::Array(_) => "array",
serde_json::Value::Object(_) => "object",
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::reasoning::conversation::Conversation;
use crate::types::AgentId;
#[test]
fn test_agent_loop_creation() {
let state = LoopState::new(AgentId::new(), Conversation::with_system("test"));
let config = LoopConfig::default();
let loop_instance = AgentLoop::<Reasoning>::new(state, config);
assert_eq!(loop_instance.state.iteration, 0);
}
#[test]
fn test_loop_termination_into_result() {
let state = LoopState::new(AgentId::new(), Conversation::new());
let termination = LoopTermination {
reason: LoopTerminationReason::MaxIterations { iterations: 25 },
state,
};
let result = termination.into_result();
assert!(matches!(
result.termination_reason,
TerminationReason::MaxIterations
));
}
fn _prove_reasoning_to_policy(_loop: AgentLoop<Reasoning>) {
}
fn _prove_policy_to_dispatch(_loop: AgentLoop<PolicyCheck>) {
}
fn _prove_dispatch_to_observing(_loop: AgentLoop<ToolDispatching>) {
}
fn _prove_observing_to_continuation(_loop: AgentLoop<Observing>) {
}
struct FakeDelegation;
#[async_trait::async_trait]
impl crate::reasoning::delegation::DelegationExecutor for FakeDelegation {
async fn delegate(
&self,
target: &str,
message: &str,
ctx: crate::reasoning::delegation::DelegationContext,
) -> Result<String, crate::reasoning::delegation::DelegationError> {
Ok(format!(
"handled '{}' for '{}' at depth {}",
message, target, ctx.depth
))
}
}
#[tokio::test]
async fn dispatch_routes_delegate_through_delegation_handle() {
let obs = super::dispatch_delegations(
&[ProposedAction::Delegate {
call_id: "test-call".into(),
target: "reviewer".into(),
message: "check this".into(),
}],
Some(&FakeDelegation),
&LoopConfig {
delegation_depth: 2,
delegation_chain: vec!["planner".into()],
..Default::default()
},
)
.await;
assert_eq!(obs.len(), 1);
assert!(!obs[0].is_error, "delegation should succeed: {:?}", obs[0]);
assert!(obs[0].content.contains("check this"));
assert!(obs[0].content.contains("depth 2"));
}
#[tokio::test]
async fn dispatch_delegate_without_handle_is_honest_error() {
let obs = super::dispatch_delegations(
&[ProposedAction::Delegate {
call_id: "test-call".into(),
target: "reviewer".into(),
message: "hi".into(),
}],
None,
&LoopConfig::default(),
)
.await;
assert_eq!(obs.len(), 1);
assert!(obs[0].is_error);
assert!(obs[0].content.to_lowercase().contains("delegation"));
}
#[tokio::test]
async fn delegate_observations_carry_the_originating_call_id() {
let obs = super::dispatch_delegations(
&[ProposedAction::Delegate {
call_id: "toolu_abc123".into(),
target: "reviewer".into(),
message: "check this".into(),
}],
Some(&FakeDelegation),
&LoopConfig::default(),
)
.await;
assert_eq!(obs.len(), 1);
assert_eq!(obs[0].call_id.as_deref(), Some("toolu_abc123"));
let err_obs = super::dispatch_delegations(
&[ProposedAction::Delegate {
call_id: "toolu_xyz789".into(),
target: "reviewer".into(),
message: "check this".into(),
}],
None,
&LoopConfig::default(),
)
.await;
assert_eq!(err_obs.len(), 1);
assert!(err_obs[0].is_error);
assert_eq!(err_obs[0].call_id.as_deref(), Some("toolu_xyz789"));
}
#[test]
fn delegate_tool_call_converts_to_a_delegate_action() {
let action = super::tool_call_to_action(
"toolu_1".to_string(),
"delegate".to_string(),
r#"{"agent":"reviewer","task":"check src/main.rs"}"#.to_string(),
true,
);
match action {
ProposedAction::Delegate {
call_id,
target,
message,
} => {
assert_eq!(call_id, "toolu_1");
assert_eq!(target, "reviewer");
assert_eq!(message, "check src/main.rs");
}
other => panic!("expected Delegate, got {other:?}"),
}
}
#[test]
fn delegate_tool_call_is_left_alone_when_the_runner_has_no_delegation_handle() {
let action = super::tool_call_to_action(
"toolu_shell".to_string(),
"delegate".to_string(),
r#"{"agent":"reviewer","task":"check src/main.rs"}"#.to_string(),
false,
);
match action {
ProposedAction::ToolCall {
call_id,
name,
arguments,
} => {
assert_eq!(call_id, "toolu_shell");
assert_eq!(name, "delegate");
assert!(arguments.contains("reviewer"));
}
other => panic!("expected the call to stay a ToolCall, got {other:?}"),
}
}
#[test]
fn non_delegate_tool_call_stays_a_tool_call() {
let action = super::tool_call_to_action(
"toolu_2".to_string(),
"list_agents".to_string(),
"{}".to_string(),
true,
);
assert!(matches!(action, ProposedAction::ToolCall { .. }));
}
#[test]
fn malformed_delegate_args_fall_through_to_a_tool_call() {
for bad in [
"not json",
r#"{"agent":"reviewer"}"#,
r#"{"task":"do it"}"#,
r#"{"agent":42,"task":"do it"}"#,
] {
let action = super::tool_call_to_action(
"toolu_3".to_string(),
"delegate".to_string(),
bad.to_string(),
true,
);
assert!(
matches!(action, ProposedAction::ToolCall { .. }),
"malformed delegate args {bad:?} must stay a ToolCall"
);
}
}
#[test]
fn denied_delegate_pushes_a_correlated_tool_result() {
let mut conversation = Conversation::new();
super::push_denial_tool_result(
&mut conversation,
&ProposedAction::Delegate {
call_id: "toolu_deny".into(),
target: "reviewer".into(),
message: "check".into(),
},
"no cedar permit",
);
let rendered = format!("{:?}", conversation);
assert!(
rendered.contains("toolu_deny"),
"denial must be correlated to the originating call id: {rendered}"
);
assert!(rendered.contains("no cedar permit"));
}
}