use super::{ProviderRunOptions, bounded_auto_compaction_diagnostic, reborrow_output_sink};
use crate::{
agent::steering::{AgentSteering, SteeringReservation},
agent::{
AgentOutputSink, AgentRunOutput, AgentRunRequest, AgentSession, ManualCompactionRequest,
PreflightRequestProjection,
},
cancellation::{AgentCancellation, is_run_canceled},
compaction::CompactionResult,
config::{AutoCompactionSettings, EffectiveConfig, Settings},
output::{InvocationMode, OutputEvent, UserPromptOrigin},
sessions::Session,
tools::ToolRuntime,
};
use anyhow::Result;
use std::sync::{Arc, atomic::AtomicU64};
use super::preparation::PreparedRun;
fn manual_compaction_instructions(
request: &crate::agent::ManualCompactionRequest,
) -> Option<String> {
request.context.as_ref().map(|context| {
format!("The requesting agent identified this context as important to preserve:\n{context}")
})
}
pub(super) fn run(
prepared: PreparedRun<'_>,
instructions: &[crate::instructions::InstructionFile],
options: &mut ProviderRunOptions<'_, '_>,
steering: Option<AgentSteering>,
) -> Result<AgentRunOutput> {
let PreparedRun {
active_config,
settings,
context_budget,
cancellation,
parent_agent_for_provider,
provider,
hooks,
tools,
mut title_job,
auto,
auto_eligible,
auto_policy,
herdr_reporter,
} = prepared;
let request_sequence = Arc::new(AtomicU64::new(0));
let scope = CompactionScope {
config: active_config.as_ref(),
settings: &settings,
cancellation: &cancellation,
agent: &parent_agent_for_provider,
tools: &tools,
auto: &auto,
request_sequence: &request_sequence,
};
let compaction_limit = auto.compaction_limit();
let mut prompt = options.prompt.to_string();
let mut prompt_origin = UserPromptOrigin::User;
let mut effective_prompt: Option<String> = None;
let mut preflight_projection = None;
let mut combined = AgentRunOutput::default();
let mut auto_compaction_count = 0_usize;
let mut preflight_compaction_used = false;
let mut provider_rejection_compacted = false;
let mut projected_candidate: Option<(String, String)> = None;
let mut steering_reservation = None;
if auto_eligible {
let (effective_current_prompt, outcome) = scope.preflight_compaction(
&prompt,
context_budget.threshold_tokens(),
compaction_limit.allows(auto_compaction_count),
&mut combined,
options,
)?;
effective_prompt = Some(effective_current_prompt);
match outcome {
PreflightOutcome::NotAttempted => {}
PreflightOutcome::Projected(projected) => preflight_projection = Some(*projected),
PreflightOutcome::Compacted => {
preflight_compaction_used = true;
auto_compaction_count = auto_compaction_count.saturating_add(1);
}
}
}
loop {
let request = AgentRunRequest {
prompt: &prompt,
prompt_origin,
effective_prompt: effective_prompt.as_deref(),
tools: Some(&tools),
hooks: Some(&hooks),
session: options.session,
cwd: options.cwd,
output_sink: reborrow_output_sink(&mut options.output_sink),
cancellation: cancellation.clone(),
session_title_job: title_job.take(),
semantic_progress_timeout: Some(settings.provider_stream.semantic_progress_timeout()),
invocation_mode: options.invocation_mode,
agent_id: options
.selected_primary_agent
.as_ref()
.map(|profile| profile.id.clone()),
initial_instructions: instructions,
herdr_reporter: herdr_reporter.clone(),
continuation_auto_compaction_policy: (auto_eligible
&& compaction_limit.allows(auto_compaction_count))
.then(|| auto_policy.clone())
.flatten(),
};
let output_result = parent_agent_for_provider.run_turn(
provider.as_ref(),
request,
steering.clone(),
steering.is_some() && auto_eligible,
Arc::clone(&request_sequence),
preflight_projection.take(),
);
drop(steering_reservation.take());
if output_result.as_ref().err().is_some_and(|error| {
error
.downcast_ref::<crate::agent::ContextBudgetError>()
.is_some_and(crate::agent::ContextBudgetError::provider_rejected)
}) {
if provider_rejection_compacted {
return output_result;
}
provider_rejection_compacted = true;
}
let (output, continuation_overflow) = recover_continuation_overflow(
output_result,
auto_eligible && compaction_limit.allows(auto_compaction_count),
auto_policy.as_ref(),
preflight_compaction_used,
auto_compaction_count > 0,
&mut options.output_sink,
)?;
let signals = merge_turn_output(&mut combined, output);
let manual_compaction_requested = signals.manual_compaction_request.is_some();
if !manual_compaction_requested
&& (!auto_eligible || !compaction_limit.allows(auto_compaction_count))
{
return Ok(combined);
}
if signals.persistence_degraded
|| (!manual_compaction_requested && signals.auto_compaction_blocked_by_recovery)
{
if let Some(overflow) = continuation_overflow {
return Err(anyhow::anyhow!(overflow.message));
}
return Ok(combined);
}
cancellation.check()?;
let session = options
.session
.ok_or_else(|| anyhow::anyhow!("magi_control compaction requires an active session"))?;
let observed_steering = steering.as_ref().and_then(AgentSteering::reserve_collapsed);
let effective_candidate = continuation_prompt_projection(
observed_steering.as_ref(),
&mut projected_candidate,
&tools,
options.invocation_mode,
);
let current_tokens = parent_agent_for_provider.project_prompt_input_tokens(
&effective_candidate,
session,
Some(&tools),
)?;
let max_tokens = parent_agent_for_provider.context_max_tokens();
if !manual_compaction_requested
&& continuation_overflow.is_none()
&& !auto.triggered(current_tokens as u64, max_tokens as u64)
{
if let Some(batch) = observed_steering {
prompt = batch.text().to_string();
steering_reservation = Some(batch);
effective_prompt = Some(effective_candidate);
prompt_origin = UserPromptOrigin::Steering;
continue;
}
return Ok(combined);
}
drop(observed_steering);
let (trigger_tokens, trigger_threshold) =
match (manual_compaction_requested, continuation_overflow) {
(true, _) => (current_tokens, "agent-requested compaction".to_string()),
(false, Some(overflow)) => (overflow.estimated_tokens, overflow.threshold_display),
(false, None) => (
current_tokens,
auto.threshold_display()
.unwrap_or_else(|| "invalid threshold".to_string()),
),
};
emit_compaction_started(
&mut options.output_sink,
trigger_tokens,
max_tokens,
trigger_threshold,
)?;
let result = scope.compact_for_continuation(
session,
signals.manual_compaction_request.as_ref(),
&mut combined,
options,
)?;
auto_compaction_count = auto_compaction_count.saturating_add(1);
let observed_steering = steering.as_ref().and_then(AgentSteering::reserve_collapsed);
if let Some(warning) = result.rotation_warning.as_ref() {
emit_diagnostic(&mut options.output_sink, "warning", warning)?;
}
let effective_candidate = continuation_prompt_projection(
observed_steering.as_ref(),
&mut projected_candidate,
&tools,
options.invocation_mode,
);
scope.verify_compacted_continuation(
session,
&effective_candidate,
max_tokens,
result.summary,
manual_compaction_requested,
options,
)?;
match observed_steering {
Some(batch) => {
prompt = batch.text().to_string();
steering_reservation = Some(batch);
effective_prompt = Some(effective_candidate);
prompt_origin = UserPromptOrigin::Steering;
}
None => {
prompt = "continue".to_string();
effective_prompt = None;
prompt_origin = UserPromptOrigin::AutomaticCompaction;
}
}
}
}
struct CompactionScope<'r> {
config: &'r EffectiveConfig,
settings: &'r Settings,
cancellation: &'r AgentCancellation,
agent: &'r AgentSession,
tools: &'r ToolRuntime,
auto: &'r AutoCompactionSettings,
request_sequence: &'r Arc<AtomicU64>,
}
enum PreflightOutcome {
NotAttempted,
Projected(Box<PreflightRequestProjection>),
Compacted,
}
struct ContinuationOverflow {
estimated_tokens: usize,
threshold_display: String,
message: String,
}
struct TurnCompactionSignals {
persistence_degraded: bool,
auto_compaction_blocked_by_recovery: bool,
manual_compaction_request: Option<ManualCompactionRequest>,
}
impl CompactionScope<'_> {
fn compact(
&self,
session: &Session,
additional_instructions: Option<String>,
combined: &mut AgentRunOutput,
options: &mut ProviderRunOptions<'_, '_>,
) -> Result<Option<CompactionResult>> {
compact_accounted(
crate::compaction::CompactSessionJob {
active_config: self.config.clone(),
settings: self.settings.clone(),
session: session.clone(),
cwd: options.cwd.to_path_buf(),
cancellation: self.cancellation.clone(),
custom_instructions: None,
additional_instructions,
},
combined,
self.request_sequence,
&mut options.output_sink,
)
}
fn preflight_compaction(
&self,
prompt: &str,
hard_threshold: usize,
compaction_allowed: bool,
combined: &mut AgentRunOutput,
options: &mut ProviderRunOptions<'_, '_>,
) -> Result<(String, PreflightOutcome)> {
self.cancellation.check()?;
let session = options
.session
.expect("auto-compaction eligibility requires session");
let effective_prompt = AgentSession::effective_prompt_for_projection(
prompt,
Some(self.tools),
options.invocation_mode,
);
let static_tokens = self
.agent
.project_prompt_input_tokens_without_session(&effective_prompt, Some(self.tools))?;
if static_tokens > hard_threshold || !compaction_allowed {
return Ok((effective_prompt, PreflightOutcome::NotAttempted));
}
let projected = self.agent.project_prompt_for_preflight(
&effective_prompt,
session,
Some(self.tools),
)?;
let current_tokens = projected.projection.tokens;
let max_tokens = self.agent.context_max_tokens();
let early_reason = if current_tokens <= hard_threshold {
early_compaction_reason(
self.settings,
prompt,
&projected,
max_tokens,
&mut options.output_sink,
)?
} else {
None
};
if current_tokens <= hard_threshold && early_reason.is_none() {
return Ok((
effective_prompt,
PreflightOutcome::Projected(Box::new(projected)),
));
}
drop(projected);
let threshold = early_reason
.unwrap_or_else(|| format!("preflight hard budget ({hard_threshold} tokens)"));
emit_compaction_started(
&mut options.output_sink,
current_tokens,
max_tokens,
threshold,
)?;
self.compact_before_first_request(session, &effective_prompt, combined, options)?;
Ok((effective_prompt, PreflightOutcome::Compacted))
}
fn compact_before_first_request(
&self,
session: &Session,
effective_prompt: &str,
combined: &mut AgentRunOutput,
options: &mut ProviderRunOptions<'_, '_>,
) -> Result<()> {
let summary = match self.compact(session, None, combined, options) {
Ok(Some(result)) => {
super::emit_compaction_fast_observation(
&mut options.output_sink,
&result,
self.request_sequence,
)?;
if let Some(warning) = result.rotation_warning.as_ref() {
emit_diagnostic(&mut options.output_sink, "warning", warning)?;
}
result.summary
}
Ok(None) => {
let message = bounded_auto_compaction_diagnostic(
"automatic preflight compaction produced no usable summary; old primary remains authoritative; no provider request was sent",
);
emit_compaction_failed(&mut options.output_sink, &message, false)?;
return Err(anyhow::anyhow!(message));
}
Err(error) if is_run_canceled(&error) => {
let message = "automatic preflight compaction canceled; old primary remains authoritative; no provider request was sent";
emit_compaction_failed(&mut options.output_sink, message, true)?;
return Err(error);
}
Err(error) => {
let authority = error
.downcast_ref::<crate::sessions::CompactionRotationError>()
.map(|_| "")
.unwrap_or(" old primary remains authoritative;");
let message = bounded_auto_compaction_diagnostic(format!(
"automatic preflight compaction failed;{authority} no provider request was sent: {error}"
));
emit_compaction_failed(&mut options.output_sink, &message, false)?;
return Err(anyhow::anyhow!(message));
}
};
let projected_tokens =
self.agent
.project_prompt_input_tokens(effective_prompt, session, Some(self.tools))?;
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::CompactionCompleted {
current_tokens: projected_tokens,
max_tokens: self.agent.context_max_tokens(),
summary,
})?;
}
if let Err(error) =
self.agent
.ensure_prompt_context_fits(effective_prompt, session, Some(self.tools))
{
let message = bounded_auto_compaction_diagnostic(format!(
"automatic preflight compaction completed, but current prompt still exceeds hard context budget; new checkpoint remains authoritative; no provider request was sent: {error}"
));
emit_diagnostic(&mut options.output_sink, "error", &message)?;
return Err(anyhow::anyhow!(message));
}
Ok(())
}
fn compact_for_continuation(
&self,
session: &Session,
manual_compaction_request: Option<&ManualCompactionRequest>,
combined: &mut AgentRunOutput,
options: &mut ProviderRunOptions<'_, '_>,
) -> Result<CompactionResult> {
let additional_instructions =
manual_compaction_request.and_then(manual_compaction_instructions);
let result = match self.compact(session, additional_instructions, combined, options) {
Ok(Some(result)) => result,
Ok(None) => {
let message = "automatic compaction produced no usable summary; old primary remains authoritative; no continuation was sent";
emit_compaction_failed(&mut options.output_sink, message, false)?;
anyhow::bail!(message)
}
Err(error) if is_run_canceled(&error) => {
let message = "automatic compaction canceled; old primary remains authoritative; no continuation was sent";
emit_compaction_failed(&mut options.output_sink, message, true)?;
return Err(error);
}
Err(error) => {
let message = if let Some(rotation_error) =
error.downcast_ref::<crate::sessions::CompactionRotationError>()
{
bounded_auto_compaction_diagnostic(format!(
"automatic compaction failed; no continuation was sent: {rotation_error}"
))
} else {
bounded_auto_compaction_diagnostic(format!(
"automatic compaction failed; old primary remains authoritative; no continuation was sent: {error}"
))
};
emit_compaction_failed(&mut options.output_sink, &message, false)?;
return Err(anyhow::anyhow!(message));
}
};
super::emit_compaction_fast_observation(
&mut options.output_sink,
&result,
self.request_sequence,
)?;
Ok(result)
}
fn verify_compacted_continuation(
&self,
session: &Session,
effective_candidate: &str,
max_tokens: usize,
summary: String,
manual_compaction_requested: bool,
options: &mut ProviderRunOptions<'_, '_>,
) -> Result<()> {
let projected_tokens = self.agent.project_prompt_input_tokens(
effective_candidate,
session,
Some(self.tools),
)?;
if let Some(sink) = options.output_sink.as_deref_mut() {
sink.output_event(OutputEvent::CompactionCompleted {
current_tokens: projected_tokens,
max_tokens,
summary,
})?;
}
let compacted_tokens = match self.agent.ensure_prompt_context_fits(
effective_candidate,
session,
Some(self.tools),
) {
Ok(tokens) => tokens,
Err(error) => {
let message = bounded_auto_compaction_diagnostic(format!(
"automatic compaction completed, but projected continuation exceeds normal context budget; new checkpoint remains authoritative; no provider continuation was sent: {error}"
));
emit_diagnostic(&mut options.output_sink, "error", &message)?;
return Err(anyhow::anyhow!(message));
}
};
if !manual_compaction_requested
&& self
.auto
.triggered(compacted_tokens as u64, max_tokens as u64)
{
let message = "automatic compaction completed, but projected context remains at or above threshold; new checkpoint remains authoritative; no continuation was sent";
emit_diagnostic(&mut options.output_sink, "error", message)?;
anyhow::bail!(message)
}
if let Err(error) = self.cancellation.check() {
let message = "automatic continuation canceled; new checkpoint remains authoritative; no provider continuation was sent";
emit_diagnostic(&mut options.output_sink, "warning", message)?;
return Err(error);
}
Ok(())
}
}
fn early_compaction_reason(
settings: &Settings,
prompt: &str,
projected: &PreflightRequestProjection,
max_tokens: usize,
output_sink: &mut Option<&mut dyn AgentOutputSink>,
) -> Result<Option<String>> {
let timing = &settings.jev.compaction_timing;
if !timing.enabled
|| projected.projection.tokens
<= crate::agent::compaction_timing::soft_threshold_tokens(timing, max_tokens)
{
return Ok(None);
}
let conversation = projected
.request
.conversation_items_iter()
.collect::<Vec<_>>();
match crate::agent::compaction_timing::unrelated_work_probability(prompt, &conversation) {
Ok(probability) if probability >= timing.threshold => Ok(Some(format!(
"Jev compaction timing: new request starts unrelated work (p={probability:.2})"
))),
Ok(_) => Ok(None),
Err(error) => {
if settings.jev.show_debug_messages {
emit_diagnostic(
output_sink,
"warning",
&format!("Jev compaction timing skipped: {error}"),
)?;
}
Ok(None)
}
}
}
fn recover_continuation_overflow(
output_result: Result<AgentRunOutput>,
continuation_compaction_allowed: bool,
auto_policy: Option<&(usize, String)>,
preflight_compaction_used: bool,
auto_compaction_used: bool,
output_sink: &mut Option<&mut dyn AgentOutputSink>,
) -> Result<(AgentRunOutput, Option<ContinuationOverflow>)> {
let error = match output_result {
Ok(output) => return Ok((output, None)),
Err(error) => error,
};
if continuation_compaction_allowed
&& error
.downcast_ref::<crate::agent::ContextBudgetError>()
.is_some_and(|budget| budget.phase() == crate::agent::ContextBudgetPhase::Continuation)
{
let budget = error
.downcast::<crate::agent::ContextBudgetError>()
.expect("context budget error checked above");
let threshold = budget.threshold_tokens();
let threshold_display = if budget.provider_rejected() {
"provider rejected prompt as too long".to_string()
} else {
auto_policy
.filter(|(soft_threshold, _)| *soft_threshold == threshold)
.map(|(_, display)| display.clone())
.unwrap_or_else(|| format!("continuation hard budget ({threshold} tokens)"))
};
let overflow = ContinuationOverflow {
estimated_tokens: budget.estimated_tokens(),
threshold_display,
message: budget.to_string(),
};
return Ok((
budget.into_partial_output().unwrap_or_default(),
Some(overflow),
));
}
let input_persistence_failed = error
.downcast_ref::<crate::agent::RequiredUserInputPersistenceError>()
.is_some();
let message = if input_persistence_failed && preflight_compaction_used {
format!(
"current user input persistence failed after preflight compaction; new checkpoint remains authoritative; no provider request was sent: {error}"
)
} else if input_persistence_failed && auto_compaction_used {
format!(
"automatic continuation persistence failed; new checkpoint remains authoritative; no provider continuation was sent: {error}"
)
} else {
return Err(error);
};
let message = bounded_auto_compaction_diagnostic(message);
emit_diagnostic(output_sink, "error", &message)?;
Err(error.context(message))
}
fn merge_turn_output(
combined: &mut AgentRunOutput,
output: AgentRunOutput,
) -> TurnCompactionSignals {
if !combined.text.is_empty() && !output.text.is_empty() {
combined.text.push('\n');
}
combined.text.push_str(&output.text);
combined.usage = output.usage;
if !matches!(output.fast_outcome, crate::fast::FastOutcome::NotRequested) {
combined.fast_outcome = output.fast_outcome;
}
combined.total_tokens = match (combined.total_tokens, output.total_tokens) {
(Some(total), Some(next)) => Some(total.saturating_add(next)),
(None, next) => next,
(total, None) => total,
};
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;
combined.manual_compaction_request = None;
TurnCompactionSignals {
persistence_degraded: output.persistence_degraded,
auto_compaction_blocked_by_recovery: output.auto_compaction_blocked_by_recovery,
manual_compaction_request: output.manual_compaction_request,
}
}
fn continuation_prompt_projection(
observed_steering: Option<&SteeringReservation>,
projected_candidate: &mut Option<(String, String)>,
tools: &ToolRuntime,
invocation_mode: InvocationMode,
) -> String {
let Some(reservation) = observed_steering else {
return "continue".to_string();
};
let candidate = reservation.text();
if let Some((raw, effective)) = projected_candidate.as_ref()
&& raw == candidate
{
return effective.clone();
}
let effective =
AgentSession::effective_prompt_for_projection(candidate, Some(tools), invocation_mode);
*projected_candidate = Some((candidate.to_string(), effective.clone()));
effective
}
fn emit_diagnostic(
output_sink: &mut Option<&mut dyn AgentOutputSink>,
level: &str,
message: &str,
) -> Result<()> {
if let Some(sink) = output_sink.as_deref_mut() {
sink.output_event(OutputEvent::Diagnostic {
level: level.to_string(),
message: message.to_string(),
})?;
}
Ok(())
}
fn emit_compaction_started(
output_sink: &mut Option<&mut dyn AgentOutputSink>,
current_tokens: usize,
max_tokens: usize,
threshold: String,
) -> Result<()> {
if let Some(sink) = output_sink.as_deref_mut() {
sink.output_event(OutputEvent::CompactionTriggered {
current_tokens,
max_tokens,
threshold,
})?;
sink.output_event(OutputEvent::CompactionStarted)?;
}
Ok(())
}
fn emit_compaction_failed(
output_sink: &mut Option<&mut dyn AgentOutputSink>,
message: &str,
canceled: bool,
) -> Result<()> {
if let Some(sink) = output_sink.as_deref_mut() {
sink.output_event(OutputEvent::CompactionFailed {
message: message.to_string(),
canceled,
})?;
sink.output_event(OutputEvent::Diagnostic {
level: if canceled { "warning" } else { "error" }.to_string(),
message: message.to_string(),
})?;
}
Ok(())
}
fn compact_accounted(
job: crate::compaction::CompactSessionJob,
output: &mut AgentRunOutput,
sequence: &AtomicU64,
sink: &mut Option<&mut dyn crate::agent::AgentOutputSink>,
) -> Result<Option<crate::compaction::CompactionResult>> {
let sequence = crate::agent::next_request_sequence(sequence);
let mut usage = crate::agent::provider_stream::CompactionUsage::default();
let result = crate::compaction::compact_session_observed(job, &mut |provider, event| {
usage.observe(provider, event, sequence, sink)
});
if let Some(total) = usage.total {
output.total_tokens = Some(output.total_tokens.unwrap_or(0).saturating_add(total));
}
result
}