use super::*;
#[derive(Default)]
pub(super) struct ChildAutoCompactionState {
pub(super) count: usize,
pub(super) checkpoint_committed: bool,
}
#[derive(Debug)]
pub(super) struct BoundedPartialOutput {
pub(super) text: String,
pub(super) truncated: bool,
}
#[derive(Debug)]
pub(super) struct ChildCompactionFailure {
pub(super) error: anyhow::Error,
pub(super) partial_output: Option<BoundedPartialOutput>,
}
impl ChildCompactionFailure {
pub(super) fn new(error: anyhow::Error, partial_output: &mut String) -> Self {
let mut text = std::mem::take(partial_output);
let truncated = truncate_string_field(&mut text, SUBAGENT_RESULT_OUTPUT_CHAR_LIMIT);
let partial_output = (!text.is_empty()).then_some(BoundedPartialOutput { text, truncated });
Self {
error,
partial_output,
}
}
pub(super) fn into_parts(self) -> (anyhow::Error, Option<BoundedPartialOutput>) {
(self.error, self.partial_output)
}
}
impl std::fmt::Display for ChildCompactionFailure {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
self.error.fmt(formatter)
}
}
impl std::error::Error for ChildCompactionFailure {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
Some(self.error.as_ref())
}
}
pub(super) struct ChildTurnRun<'a> {
pub(super) agent: &'a AgentSession,
pub(super) provider: &'a dyn Provider,
pub(super) original_task: &'a SubagentTask,
pub(super) prompt: &'a str,
pub(super) prompt_origin: crate::output::UserPromptOrigin,
pub(super) steering: crate::agent::steering::AgentSteering,
pub(super) tools: &'a ToolRuntime,
pub(super) hooks: Option<&'a HookRuntime>,
pub(super) session: Option<&'a Session>,
pub(super) cwd: &'a Path,
pub(super) sink: &'a mut SubagentActivitySink,
pub(super) cancellation: &'a AgentCancellation,
pub(super) semantic_progress_timeout: Duration,
pub(super) agent_id: Option<String>,
pub(super) compaction: Option<&'a super::config::SubagentCompactionConfig>,
pub(super) auto_compaction_state: &'a mut ChildAutoCompactionState,
pub(super) request_sequence: Arc<AtomicU64>,
}
pub(super) fn run_child_turn_with_auto_compaction(
mut run: ChildTurnRun<'_>,
) -> anyhow::Result<AgentRunOutput> {
let initial_prompt = run.prompt.to_string();
let initial_origin = run.prompt_origin;
let Some(session) = run.session else {
return run_child_agent_once(&mut run, &initial_prompt, initial_origin, None);
};
let Some(compaction) = run.compaction else {
return run_child_agent_once(&mut run, &initial_prompt, initial_origin, None);
};
let auto = compaction
.settings
.compaction
.auto_for_provider(run.agent.provider_id(), run.agent.context_max_tokens());
let Some(auto_policy) = (auto.is_enabled() && run.agent.context_enabled())
.then(|| {
crate::agent::runner::auto_compaction_policy(&auto, run.agent.context_max_tokens())
})
.flatten()
else {
return run_child_agent_once(&mut run, &initial_prompt, initial_origin, None);
};
let compaction_limit: AutoCompactionLimit = auto.compaction_limit();
let mut prompt = initial_prompt;
let mut prompt_origin = initial_origin;
let mut combined = AgentRunOutput::default();
let mut combined_segments = 0_usize;
let mut steering_reservation = None;
let mut provider_rejection_compacted = false;
run.cancellation.check()?;
let static_tokens = run
.agent
.project_prompt_input_tokens_without_session(&prompt, Some(run.tools))?;
let current_tokens =
run.agent
.project_prompt_input_tokens(&prompt, session, Some(run.tools))?;
let hard_threshold = run.agent.context_threshold_tokens();
if static_tokens <= hard_threshold
&& current_tokens > hard_threshold
&& compaction_limit.allows(run.auto_compaction_state.count)
{
let (summary, rotation_warning) = compact_child_session_with_partial_output(
&mut run,
session,
compaction,
current_tokens,
format!("preflight hard budget ({hard_threshold} tokens)"),
&mut combined.text,
)?;
run.auto_compaction_state.count = run.auto_compaction_state.count.saturating_add(1);
let projected_tokens =
run.agent
.project_prompt_input_tokens(&prompt, session, Some(run.tools))?;
run.sink.output_event(OutputEvent::CompactionCompleted {
current_tokens: projected_tokens,
max_tokens: run.agent.context_max_tokens(),
summary,
})?;
if let Some(warning) = rotation_warning {
run.sink.output_event(OutputEvent::Diagnostic {
level: "warning".to_string(),
message: warning,
})?;
}
run.auto_compaction_state.checkpoint_committed = true;
run.agent
.ensure_prompt_context_fits(&prompt, session, Some(run.tools))?;
}
loop {
run.cancellation.check()?;
let policy = (compaction_limit.allows(run.auto_compaction_state.count))
.then_some(auto_policy.clone());
let output_result = run_child_agent_once(&mut run, &prompt, prompt_origin, policy);
drop(steering_reservation.take());
let mut continuation_overflow = None;
let output = match output_result {
Ok(output) => output,
Err(error)
if compaction_limit.allows(run.auto_compaction_state.count)
&& error
.downcast_ref::<ContextBudgetError>()
.is_some_and(|budget| {
budget.phase() == ContextBudgetPhase::Continuation
}) =>
{
let budget = error
.downcast::<ContextBudgetError>()
.expect("context budget error checked above");
if budget.provider_rejected() {
if provider_rejection_compacted {
return Err(budget.into());
}
provider_rejection_compacted = true;
}
let threshold = budget.threshold_tokens();
let threshold_display = if budget.provider_rejected() {
"provider rejected prompt as too long".to_string()
} else if auto_policy.0 == threshold {
auto_policy.1.clone()
} else {
format!("continuation hard budget ({threshold} tokens)")
};
continuation_overflow = Some((
budget.estimated_tokens(),
threshold_display,
budget.to_string(),
));
budget.into_partial_output().unwrap_or_default()
}
Err(error) => return Err(error),
};
let persistence_degraded = output.persistence_degraded;
let recovery_blocked = output.auto_compaction_blocked_by_recovery;
merge_agent_output(&mut combined, output, combined_segments == 0);
combined_segments = combined_segments.saturating_add(1);
if persistence_degraded || recovery_blocked {
if let Some((_, _, message)) = continuation_overflow {
return Err(anyhow::anyhow!(message));
}
return Ok(combined);
}
if !compaction_limit.allows(run.auto_compaction_state.count) {
return Ok(combined);
}
let continuation_prompt = "continue";
let current_tokens =
run.agent
.project_prompt_input_tokens(continuation_prompt, session, Some(run.tools))?;
if continuation_overflow.is_none()
&& !auto.triggered(current_tokens as u64, run.agent.context_max_tokens() as u64)
{
return Ok(combined);
}
let (trigger_tokens, trigger_threshold) = continuation_overflow
.map(|(tokens, threshold, _)| (tokens, threshold))
.unwrap_or_else(|| (current_tokens, auto_policy.1.clone()));
let (summary, rotation_warning) = compact_child_session_with_partial_output(
&mut run,
session,
compaction,
trigger_tokens,
trigger_threshold,
&mut combined.text,
)?;
run.auto_compaction_state.count = run.auto_compaction_state.count.saturating_add(1);
let reserved = run.steering.reserve_collapsed();
let continuation_prompt = reserved.as_ref().map_or("continue", |batch| batch.text());
let projected_tokens =
run.agent
.project_prompt_input_tokens(continuation_prompt, session, Some(run.tools))?;
run.sink.output_event(OutputEvent::CompactionCompleted {
current_tokens: projected_tokens,
max_tokens: run.agent.context_max_tokens(),
summary,
})?;
if let Some(warning) = rotation_warning {
run.sink.output_event(OutputEvent::Diagnostic {
level: "warning".to_string(),
message: warning,
})?;
}
run.auto_compaction_state.checkpoint_committed = true;
let compacted_tokens =
run.agent
.ensure_prompt_context_fits(continuation_prompt, session, Some(run.tools))?;
if auto.triggered(
compacted_tokens as u64,
run.agent.context_max_tokens() as u64,
) {
anyhow::bail!(
"automatic subagent compaction completed, but the child context remains at or above threshold; no child continuation was sent"
);
}
run.cancellation.check()?;
prompt = continuation_prompt.to_string();
prompt_origin = if reserved.is_some() {
crate::output::UserPromptOrigin::Steering
} else {
crate::output::UserPromptOrigin::AutomaticCompaction
};
steering_reservation = reserved;
}
}
fn run_child_agent_once(
run: &mut ChildTurnRun<'_>,
prompt: &str,
prompt_origin: crate::output::UserPromptOrigin,
continuation_auto_compaction_policy: Option<(usize, String)>,
) -> anyhow::Result<AgentRunOutput> {
run.agent.run_turn_untracked(
run.provider,
AgentRunRequest {
prompt,
prompt_origin,
effective_prompt: None,
tools: Some(run.tools),
hooks: run.hooks,
session: run.session,
cwd: run.cwd,
output_sink: Some(&mut *run.sink),
cancellation: run.cancellation.clone(),
session_title_job: None,
semantic_progress_timeout: Some(run.semantic_progress_timeout),
invocation_mode: crate::output::InvocationMode::Subagent,
agent_id: run.agent_id.clone(),
herdr_reporter: None,
initial_instructions: &[],
continuation_auto_compaction_policy,
},
Some(run.steering.clone()),
false,
Arc::clone(&run.request_sequence),
)
}
pub(super) const SUBAGENT_COMPACTION_SCOPE_PREFIX: &str = "The following single JSON document extends to the end of this instruction. The document contains the authoritative scope for this compaction:\n";
pub(super) const SUBAGENT_COMPACTION_SCOPE_DIRECTIVE: &str = "The original task payload is authoritative; summarize only that task; exclude unrelated work or tasks.";
#[derive(Serialize)]
pub(super) struct SubagentCompactionTaskPayload<'a> {
pub(super) intent: &'a str,
pub(super) mode: crate::subagents::SubagentMode,
pub(super) agent: Option<&'a str>,
pub(super) identity: Option<&'a str>,
pub(super) context: Option<&'a str>,
pub(super) cwd: Option<&'a Path>,
}
impl<'a> From<&'a SubagentTask> for SubagentCompactionTaskPayload<'a> {
fn from(task: &'a SubagentTask) -> Self {
Self {
intent: &task.intent,
mode: task.mode,
agent: task.agent.as_deref(),
identity: task.identity.as_deref(),
context: task.context.as_deref(),
cwd: task.cwd.as_deref(),
}
}
}
#[derive(Serialize)]
pub(super) struct SubagentCompactionScope<'a> {
directive: &'static str,
pub(super) original_task: SubagentCompactionTaskPayload<'a>,
}
pub(in crate::subagents) fn subagent_compaction_instructions(
task: &SubagentTask,
) -> anyhow::Result<String> {
let scope = SubagentCompactionScope {
directive: SUBAGENT_COMPACTION_SCOPE_DIRECTIVE,
original_task: SubagentCompactionTaskPayload::from(task),
};
let serialized_scope = serde_json::to_string(&scope).map_err(|error| {
anyhow::anyhow!("failed to serialize original subagent task payload: {error}")
})?;
Ok(format!(
"{SUBAGENT_COMPACTION_SCOPE_PREFIX}{serialized_scope}"
))
}
fn compact_child_session(
run: &mut ChildTurnRun<'_>,
session: &Session,
compaction: &super::config::SubagentCompactionConfig,
current_tokens: usize,
threshold: String,
) -> anyhow::Result<(String, Option<String>)> {
run.cancellation.check()?;
let additional_instructions = subagent_compaction_instructions(run.original_task)?;
run.sink.output_event(OutputEvent::CompactionTriggered {
current_tokens,
max_tokens: run.agent.context_max_tokens(),
threshold,
})?;
run.sink.output_event(OutputEvent::CompactionStarted)?;
let sequence = crate::agent::next_request_sequence(&run.request_sequence);
let mut usage = crate::agent::provider_stream::CompactionUsage::default();
match crate::compaction::compact_session_observed(
crate::compaction::CompactSessionJob {
active_config: compaction.active_config.clone(),
settings: compaction.settings.clone(),
session: session.clone(),
cwd: run.cwd.to_path_buf(),
cancellation: run.cancellation.clone(),
custom_instructions: None,
additional_instructions: Some(additional_instructions),
},
&mut |provider, event| usage.observe(provider, event, sequence, &mut Some(run.sink)),
) {
Ok(Some(result)) => Ok((result.summary, result.rotation_warning)),
Ok(None) => {
let message = format!(
"automatic subagent compaction produced no usable summary; {}; partial child output is preserved in the failed child result; no child continuation was sent",
child_session_recovery_context(session),
);
Err(anyhow::anyhow!(
crate::agent::runner::bounded_auto_compaction_diagnostic(message)
))
}
Err(error) if crate::cancellation::is_run_canceled(&error) => Err(error),
Err(error) => {
if error
.downcast_ref::<crate::sessions::CompactionRotationError>()
.is_some()
{
run.auto_compaction_state.count = run.auto_compaction_state.count.saturating_add(1);
run.auto_compaction_state.checkpoint_committed = true;
}
let message = format!(
"automatic subagent compaction failed; {}; partial child output is preserved in the failed child result; no child continuation was sent: {}",
child_session_recovery_context(session),
crate::agent::runner::bounded_auto_compaction_diagnostic(error),
);
Err(anyhow::anyhow!(
crate::agent::runner::bounded_auto_compaction_diagnostic(message)
))
}
}
}
fn compact_child_session_with_partial_output(
run: &mut ChildTurnRun<'_>,
session: &Session,
compaction: &super::config::SubagentCompactionConfig,
current_tokens: usize,
threshold: String,
partial_output: &mut String,
) -> anyhow::Result<(String, Option<String>)> {
match compact_child_session(run, session, compaction, current_tokens, threshold) {
Ok(result) => Ok(result),
Err(error) if crate::cancellation::is_run_canceled(&error) => Err(error),
Err(error) => Err(ChildCompactionFailure::new(error, partial_output).into()),
}
}
fn child_session_recovery_context(session: &Session) -> String {
format!(
"child session id={} path={}",
redact_sensitive_text(session.id()),
redact_sensitive_text(&session.path().display().to_string()),
)
}
fn merge_agent_output(combined: &mut AgentRunOutput, output: AgentRunOutput, first_segment: bool) {
if !combined.text.is_empty() && !output.text.is_empty() {
combined.text.push('\n');
}
combined.text.push_str(&output.text);
combined.usage = output.usage;
combined.total_tokens = if first_segment {
output.total_tokens
} else {
match (combined.total_tokens, output.total_tokens) {
(Some(total), Some(next)) => Some(total.saturating_add(next)),
_ => None,
}
};
combined.tool_results.extend(output.tool_results);
combined.persistence_degraded |= output.persistence_degraded;
combined.recovered_incomplete_stream |= output.recovered_incomplete_stream;
combined.auto_compaction_blocked_by_recovery |= output.auto_compaction_blocked_by_recovery;
}