mod child_agent;
mod compaction;
mod results;
mod schema;
mod setup;
use setup::ChildRunSetup;
use setup::PreparedChildRun;
use compaction::BoundedPartialOutput;
use compaction::ChildAutoCompactionState;
use compaction::ChildCompactionFailure;
use compaction::ChildTurnRun;
use compaction::run_child_turn_with_auto_compaction;
pub(super) use child_agent::ChildAgentRun;
pub(super) use child_agent::child_agent_for_task;
pub(super) use child_agent::subagent_prompt;
use schema::resolve_output_schema_for_task;
use schema::schema_validation_exhausted_error;
use schema::schema_validation_feedback_prompt;
use schema::validate_schema_bound_output;
pub(super) use results::FailedSubagentInput;
use results::completed_subagent_result;
pub(super) use results::failed_result_with_session;
pub(super) use results::failed_result_with_session_and_usage;
pub(super) use results::finish_failed_subagent;
use results::subagent_failure_message;
use super::{
activity::{SubagentActivitySink, subagent_task_activity_metadata},
config,
config::SubagentRunConfig,
dispatch_subagents,
dto::{SubagentStatus, SubagentTask, SubagentTaskResult, SubagentsOutput},
output_schema::{self, SchemaValidationError, SubagentOutputSchemaRef},
profiles,
scheduler::{SchedulerTaskReporter, TaskActivityFinisher},
};
use crate::{
agent::{
AgentOutputSink, AgentRunOutput, AgentRunRequest, AgentSession, ContextBudgetError,
ContextBudgetPhase,
},
cancellation::AgentCancellation,
config::AutoCompactionLimit,
hooks::{HookPolicyError, HookRuntime},
output::{
ActivityEvent, ActivityId, ActivityKind, ActivityStatus, OutputEvent, redact_sensitive_text,
},
providers::{Provider, ProviderSelection},
sessions::{Session, SessionEventKind, SessionManager},
tools::ToolRuntime,
};
pub(super) use results::format_duration;
pub(super) use results::truncate_string_field;
pub(super) use results::truncate_subagents_output;
use serde::Serialize;
use serde_json::{Value, json};
use std::{
collections::BTreeSet,
path::{Path, PathBuf},
sync::{Arc, atomic::AtomicU64},
time::Duration,
};
pub(super) const SUBAGENT_RESULT_OUTPUT_CHAR_LIMIT: usize = 24_000;
pub(super) const SUBAGENT_PROVIDER_RETRY_MAX_ATTEMPTS: u64 = 1;
pub(super) const SUBAGENT_PROVIDER_RETRY_BACKOFF: Duration = Duration::from_secs(2);
const SUBAGENT_PROVIDER_RETRY_POLL_INTERVAL: Duration = Duration::from_millis(100);
pub(super) const SUBAGENT_RESULT_ERROR_CHAR_LIMIT: usize = 8_000;
pub(super) const SUBAGENT_TRUNCATION_MARKER: &str = "\n[subagent result truncated: original_chars=";
const SNAPSHOT_EVENT_LIMIT: usize = 200;
pub(super) const SNAPSHOT_BYTE_LIMIT: usize = 256 * 1024;
pub(super) struct SubagentRunInput<'a> {
pub(super) id: String,
pub(super) task: SubagentTask,
pub(super) cwd: PathBuf,
pub(super) config: &'a SubagentRunConfig,
pub(super) cancellation: AgentCancellation,
pub(super) batch_id: &'a ActivityId,
pub(super) reporter: Option<SchedulerTaskReporter>,
pub(super) finisher: TaskActivityFinisher,
}
pub(super) fn run_one_subagent(input: SubagentRunInput<'_>) -> SubagentTaskResult {
let SubagentRunInput {
id,
task,
cwd,
config,
cancellation,
batch_id,
reporter,
finisher,
} = input;
let task_activity_id = batch_id.child(&id);
let (cancellation, cancel_handle) = cancellation.child_token();
let steering = crate::agent::steering::AgentSteering::new();
let _control_registration = config.subagent_controls.as_ref().and_then(|controls| {
controls.register(
task_activity_id.clone(),
steering.clone(),
cancel_handle,
cancellation.clone(),
)
});
finisher.try_start(config, || ActivityEvent::Started {
id: task_activity_id.clone(),
parent_id: Some(batch_id.clone()),
kind: ActivityKind::SubagentTask,
status: ActivityStatus::Running,
metadata: subagent_task_activity_metadata(&id, &task, config.depth + 1, None),
});
let setup = ChildRunSetup {
id: &id,
task: &task,
cwd: &cwd,
config,
task_activity_id: &task_activity_id,
batch_id,
finisher: &finisher,
reporter: reporter.as_ref(),
cancellation: &cancellation,
};
let PreparedChildRun {
run: child_run,
request_agent,
session,
tools,
hooks: child_hooks,
} = match setup.prepare() {
Ok(prepared) => prepared,
Err(failure) => {
return finish_failed_subagent(FailedSubagentInput {
config,
task_activity_id,
id,
task,
cwd,
error: failure.error,
finisher: &finisher,
session: failure.session.as_ref(),
partial_output: None,
usage: None,
});
}
};
let output_schema = resolve_output_schema_for_task(&task, config);
let prompt = subagent_prompt(&id, &task, output_schema.as_ref());
let request_sequence = Arc::new(AtomicU64::new(0));
let mut child_sink = SubagentActivitySink {
parent_id: task_activity_id.clone(),
activity_sender: config.activity_sender.clone(),
assistant_id: task_activity_id.child("assistant"),
assistant_started: false,
reasoning_summary_group_count: 0,
reasoning_summary_group_lines: Vec::new(),
progress_reporter: reporter
.map(|reporter| Arc::new(move || reporter.progress()) as Arc<dyn Fn() + Send + Sync>),
cancellation: cancellation.clone(),
changed_files: BTreeSet::new(),
compaction_sequence: 0,
usage_by_request: super::activity::UsageAccumulator::default(),
active_compaction: None,
};
let max_schema_retries = config.schema_validation_max_retries.min(5);
let mut attempt = 0_u64;
let mut provider_retry_attempt = 0_u64;
let mut auto_compaction_state = ChildAutoCompactionState::default();
let mut prompt = prompt;
let mut prompt_origin = crate::output::UserPromptOrigin::User;
let mut steering_reservation = None;
let mut result = loop {
let run_result = run_child_turn_with_auto_compaction(ChildTurnRun {
agent: &request_agent,
original_task: &task,
provider: child_run.provider.as_ref(),
prompt: &prompt,
prompt_origin,
steering: steering.clone(),
tools: &tools,
hooks: child_hooks.as_ref(),
session: session.as_ref(),
cwd: &cwd,
sink: &mut child_sink,
cancellation: &cancellation,
semantic_progress_timeout: config.semantic_progress_timeout,
agent_id: task.identity.clone().or_else(|| task.agent.clone()),
request_sequence: Arc::clone(&request_sequence),
compaction: child_run.compaction.as_ref(),
auto_compaction_state: &mut auto_compaction_state,
});
match run_result {
Ok(mut output) => {
output.total_tokens = child_sink.usage_by_request.totals.total();
drop(steering_reservation.take());
let schema_result = output_schema
.as_ref()
.map(|schema| validate_schema_bound_output(schema, &output.text));
let final_output = schema_result.as_ref().is_none_or(Result::is_ok);
if final_output && let Some(reservation) = steering.reserve_collapsed_or_close() {
prompt = reservation.text().to_string();
prompt_origin = crate::output::UserPromptOrigin::Steering;
steering_reservation = Some(reservation);
continue;
}
if let Some(schema_result) = schema_result {
match schema_result {
Ok(structured_output) => {
child_sink.finish_assistant(ActivityStatus::Success);
finisher.finish(config, task_activity_id, ActivityStatus::Success);
break completed_subagent_result(
id,
task,
cwd,
session.as_ref(),
output,
Some(structured_output),
child_sink.usage_metrics(),
);
}
Err(errors) if attempt < max_schema_retries => {
attempt = attempt.saturating_add(1);
prompt_origin = crate::output::UserPromptOrigin::User;
prompt = schema_validation_feedback_prompt(
attempt,
max_schema_retries,
&errors,
);
continue;
}
Err(errors) => {
child_sink.finish_assistant(ActivityStatus::Failed);
let mut result = failed_result_with_session_and_usage(
id,
task,
cwd,
session.as_ref().map(|session| session.id().to_string()),
session.as_ref().map(|session| session.path().to_path_buf()),
schema_validation_exhausted_error(&errors),
child_sink.usage_metrics(),
);
result.total_tokens = child_sink.usage_by_request.totals.total();
finisher.finish(config, task_activity_id, ActivityStatus::Failed);
break result;
}
}
}
child_sink.finish_assistant(ActivityStatus::Success);
finisher.finish(config, task_activity_id, ActivityStatus::Success);
break completed_subagent_result(
id,
task,
cwd,
session.as_ref(),
output,
None,
child_sink.usage_metrics(),
);
}
Err(error) => {
let (error, partial_output) = match error.downcast::<ChildCompactionFailure>() {
Ok(error) => error.into_parts(),
Err(error) => (error, None),
};
let retryable = crate::providers::error::retryable_provider_error(&error);
let can_retry = can_retry_subagent_provider(
provider_retry_attempt,
&cancellation,
retryable,
child_sink.has_filesystem_changes(),
auto_compaction_state.checkpoint_committed,
);
if can_retry && wait_for_provider_retry(&cancellation) {
provider_retry_attempt = provider_retry_attempt.saturating_add(1);
continue;
}
let finish_status = if cancellation.is_canceled() {
ActivityStatus::Canceled
} else {
ActivityStatus::Failed
};
let _ = child_sink.finish_compaction(finish_status);
child_sink.finish_assistant(finish_status);
let error = subagent_failure_message(&error, retryable, provider_retry_attempt);
finisher.finish(config, task_activity_id.clone(), finish_status);
let mut result = finish_failed_subagent(FailedSubagentInput {
config,
task_activity_id,
id,
task,
cwd,
error,
finisher: &finisher,
session: session.as_ref(),
partial_output,
usage: child_sink.usage_metrics(),
});
result.total_tokens = child_sink.usage_by_request.totals.total();
break result;
}
}
};
result.changed_files = child_sink.changed_files.into_iter().collect();
result
}
pub(super) fn can_retry_subagent_provider(
retry_attempt: u64,
cancellation: &AgentCancellation,
retryable: bool,
has_filesystem_changes: bool,
compaction_checkpoint_committed: bool,
) -> bool {
retry_attempt < SUBAGENT_PROVIDER_RETRY_MAX_ATTEMPTS
&& !cancellation.is_canceled()
&& retryable
&& !has_filesystem_changes
&& !compaction_checkpoint_committed
}
fn wait_for_provider_retry(cancellation: &AgentCancellation) -> bool {
let deadline = std::time::Instant::now() + SUBAGENT_PROVIDER_RETRY_BACKOFF;
loop {
if cancellation.is_canceled() {
return false;
}
let remaining = deadline.saturating_duration_since(std::time::Instant::now());
if remaining.is_zero() {
return true;
}
std::thread::sleep(remaining.min(SUBAGENT_PROVIDER_RETRY_POLL_INTERVAL));
}
}