mod completion;
mod context_budget;
mod continuation;
mod output;
mod preparation;
mod prompt_cache_key;
mod request;
mod run_request;
mod session;
mod titles;
mod tool_advertising;
mod turn_run;
pub(crate) use context_budget::ContextBudgetError;
pub(crate) use context_budget::ContextBudgetPhase;
pub(crate) use continuation::RequiredUserInputPersistenceError;
use continuation::inject_steering_update_at_continuation_boundary;
pub use output::AgentOutputSink;
pub use output::AgentRunOutput;
use prompt_cache_key::prompt_cache_key_for_run;
pub(crate) use run_request::AgentRunRequest;
pub(crate) use run_request::ManualCompactionRequest;
pub(crate) use run_request::filtered_herdr_reporter;
pub(crate) use run_request::report_preflight_blocked;
pub(crate) use run_request::report_preflight_cancelled;
pub(crate) use run_request::run_with_herdr_turn;
pub use session::AgentSession;
pub(crate) use session::AgentSessionConfig;
use titles::TitleGenerationGuard;
use tool_advertising::ToolRequestConfiguration;
use tool_advertising::tool_request_configuration;
mod code_mode;
mod compaction_timing;
mod completion_verification;
mod prompt;
pub(crate) mod prompt_injections;
pub(crate) mod provider_stream;
pub(crate) mod recovery;
pub(crate) mod runner;
mod session_persistence;
pub(crate) mod skill_suggestion;
pub(crate) mod steering;
mod subdir_instructions;
mod tool_continuation;
mod tool_lifecycle;
mod tool_result_diagnostics;
mod turn_state;
mod util;
use self::completion_verification::{CompletionVerificationDecision, verify_completion};
use self::{
prompt::build_system_prompt_with_prompt_dir_and_subagents,
provider_stream::{
ProviderEventCollector, ProviderStreamFailure, ProviderStreamResult,
collect_provider_events,
},
recovery::{
DANGLING_TOOL_INTENT_PROMPT, DANGLING_TOOL_INTENT_REASON, ToolExecutionExpectation,
completion_needs_tool_recovery_for_expectation,
},
session_persistence::{SessionPersistence, emit_replay_diagnostics},
steering::{AgentSteering, SteeringBatch},
subdir_instructions::SubdirInstructionState,
tool_lifecycle::{
ToolLifecycleRun, record_provider_context_item, run_message_phase_hooks, run_tool_lifecycle,
},
turn_state::AgentTurnState,
util::{AssistantChunkBatch, normalize_agent_identifier},
};
use crate::{
cancellation::{AgentCancellation, is_run_canceled},
config::{CompletionVerificationSettings, TextVerbosity},
context::{
ContextBudget, ContextTokenCount, ConversationReplayCache,
project_provider_conversation_items_tokens, project_provider_request_input_tokens,
},
herdr::{HerdrReporter, HerdrTurnReporter},
hex::lower_hex,
hooks::{HookPhase, HookPolicyError, HookRuntime},
instructions::InstructionFile,
output::tool_display::format_tool_block,
output::{
ActivityEvent, ActivityId, ActivitySender, ContextUsageSource, HookContextMetadata,
InvocationMode, OutputEvent,
},
providers::{
ChatMessage, Provider, ProviderConversationItem, ProviderRequest, ToolCall, Usage,
},
sessions::{Session, SessionEventKind, SessionReadDiagnostic, TurnStatus},
skills::SkillDiscovery,
tools::{ToolResult, ToolRuntime},
};
use serde_json::json;
use sha2::{Digest, Sha256};
use std::{
collections::HashSet,
path::{Path, PathBuf},
sync::{
Arc, Mutex,
atomic::{AtomicU64, Ordering},
},
time::Duration,
};
pub(crate) use self::prompt::{
SystemPromptSection, build_system_prompt_sections_with_prompt_dir_and_subagents,
load_compact_prompt,
};
struct InitialConversation {
conversation: Vec<ProviderConversationItem>,
session_read_diagnostics: Vec<SessionReadDiagnostic>,
replay_warnings: Vec<String>,
}
struct PreparedInitialRun<'run> {
provider_prompt: String,
session_read_diagnostics: Vec<SessionReadDiagnostic>,
replay_warnings: Vec<String>,
base_conversation: Arc<[ProviderConversationItem]>,
prompt_cache_key: Option<String>,
conversation_id: Option<String>,
request_projection_cache: RequestTokenProjectionCache,
session_persistence: SessionPersistence<'run>,
tool_request_configuration: ToolRequestConfiguration,
}
struct PreflightRequestProjection {
provider_id: String,
session_id: String,
session_path: PathBuf,
replay_generation: u64,
request: ProviderRequest,
projection: ContextTokenCount,
}
impl PreflightRequestProjection {
fn matching_projection(
self,
provider_id: &str,
request: &ProviderRequest,
session: Option<&Session>,
) -> Option<ContextTokenCount> {
let session = session?;
(self.provider_id == provider_id
&& self.session_id == session.id()
&& self.session_path == session.path()
&& self.replay_generation == session.replay_generation()
&& self.request == *request)
.then_some(self.projection)
}
}
enum TurnLoopStep {
Continue,
Break,
}
struct TurnSteering<'s> {
steering: Option<&'s AgentSteering>,
defer_final: bool,
expand: &'s dyn Fn(&str) -> String,
}
struct PrintTurnState<'run> {
base_conversation: Arc<[ProviderConversationItem]>,
tool_request_configuration: ToolRequestConfiguration,
prompt_cache_key: Option<String>,
conversation_id: Option<String>,
codex_turn_context: crate::providers::CodexTurnContext,
session_persistence: SessionPersistence<'run>,
subdir_instruction_state: SubdirInstructionState,
assistant_chunk_batch: AssistantChunkBatch<'run>,
turn_state: AgentTurnState,
output: AgentRunOutput,
request_projection_cache: RequestTokenProjectionCache,
request_sequence: Arc<AtomicU64>,
hook_context: HookContextMetadata,
herdr_reporter: Option<HerdrReporter>,
title_guard: TitleGenerationGuard,
}
pub(crate) fn next_request_sequence(owner: &AtomicU64) -> u64 {
owner
.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| {
Some(current.saturating_add(1))
})
.expect("request sequence update always succeeds")
}
fn should_expand_prompt_context(invocation_mode: InvocationMode) -> bool {
matches!(
invocation_mode,
InvocationMode::Print
| InvocationMode::MissionControl
| InvocationMode::LocalAgent
| InvocationMode::Side
)
}
struct RequestTokenProjectionCache {
provider_id: String,
model: String,
base_projection: ContextTokenCount,
deterministic_projected_turn_items: usize,
deterministic_projected_turn_tokens: usize,
calibrated_input_tokens: Option<(usize, usize)>,
}
impl RequestTokenProjectionCache {
fn new_for_request(provider_id: &str, request: &ProviderRequest) -> Self {
Self::new_with_base_projection(
provider_id,
&request.model,
project_provider_request_input_tokens(provider_id, request),
)
}
fn new_with_base_projection(
provider_id: &str,
model: &str,
base_projection: ContextTokenCount,
) -> Self {
Self {
provider_id: provider_id.to_string(),
model: model.to_string(),
base_projection,
deterministic_projected_turn_items: 0,
deterministic_projected_turn_tokens: 0,
calibrated_input_tokens: None,
}
}
}
impl RequestTokenProjectionCache {
fn initial_projection(&self) -> ContextTokenCount {
self.base_projection
}
fn apply_usage_calibration(&self, projection: ContextTokenCount) -> ContextTokenCount {
let Some((estimated, observed)) = self.calibrated_input_tokens else {
return projection;
};
let tokens = if projection.tokens >= estimated {
observed.saturating_add(projection.tokens - estimated)
} else {
observed.saturating_sub(estimated - projection.tokens)
};
ContextTokenCount {
tokens,
source: ContextUsageSource::LastProviderUsage,
}
}
fn project_deterministic_turn_items(
&mut self,
turn_items: &[ProviderConversationItem],
temporary_suffix: &[ProviderConversationItem],
) -> ContextTokenCount {
if turn_items.len() < self.deterministic_projected_turn_items {
self.invalidate_deterministic_turn_items();
}
if self.deterministic_projected_turn_items < turn_items.len() {
let new_items = &turn_items[self.deterministic_projected_turn_items..];
let new_projection = project_provider_conversation_items_tokens(
&self.provider_id,
&self.model,
new_items,
);
self.deterministic_projected_turn_tokens = self
.deterministic_projected_turn_tokens
.saturating_add(new_projection.tokens);
self.deterministic_projected_turn_items = turn_items.len();
}
let temporary_projection = project_provider_conversation_items_tokens(
&self.provider_id,
&self.model,
temporary_suffix,
);
ContextTokenCount {
tokens: self
.base_projection
.tokens
.saturating_add(self.deterministic_projected_turn_tokens)
.saturating_add(temporary_projection.tokens),
source: self.base_projection.source,
}
}
fn invalidate_deterministic_turn_items(&mut self) {
self.deterministic_projected_turn_items = 0;
self.deterministic_projected_turn_tokens = 0;
self.calibrated_input_tokens = None;
}
fn calibrate_anthropic_input_tokens(&mut self, estimated_tokens: usize, input_tokens: usize) {
if crate::providers::reports_anthropic_usage(&self.provider_id) {
self.calibrated_input_tokens = Some((estimated_tokens, input_tokens));
}
}
}