use std::collections::HashSet;
use std::path::Path;
use std::sync::Arc;
use std::sync::atomic::AtomicBool;
use crate::agent::turn_budget::TurnBudgetCounter;
use crate::workspace::changelist::WorkspaceChangelist;
use super::errors::AgentTurnSubPhase;
use super::turn_sink::TurnControlSink;
use crate::PerTurnFlight;
use crate::WebRequestAudit;
use crate::agent::agent_turn::messages::{
insert_separator_after_last_user_for_turn, push_assistant_merging_trailing_empty_placeholder,
};
use crate::agent::plan_artifact::PlanStepExecutorKind;
use crate::config::AgentConfig;
use crate::memory::long_term_memory::LongTermMemoryRuntime;
use crate::tool_registry;
use crate::types::{LlmSeedOverride, Message};
pub(crate) struct RunLoopCore<'a> {
pub llm_backend: &'a (dyn crate::llm::ChatCompletionsBackend + 'static),
pub client: &'a reqwest::Client,
pub api_key: &'a str,
pub cfg: &'a Arc<AgentConfig>,
pub tools_defs: &'a [crate::types::Tool],
pub effective_working_dir: &'a Path,
pub workspace_is_set: bool,
}
pub(crate) struct RunLoopIo<'a> {
pub no_stream: bool,
pub cancel: Option<Arc<AtomicBool>>,
pub control: TurnControlSink<'a>,
}
pub(crate) struct RunLoopAttach<'a> {
pub web_tool_ctx: Option<&'a tool_registry::WebToolRuntime>,
pub per_flight: Option<Arc<PerTurnFlight>>,
pub long_term_memory: Option<Arc<LongTermMemoryRuntime>>,
pub long_term_memory_scope_id: Option<String>,
pub mcp_turn: Option<crate::mcp::McpTurnHandle>,
pub read_file_turn_cache: Option<Arc<crate::read_file_turn_cache::ReadFileTurnCache>>,
pub workspace_changelist: Option<Arc<WorkspaceChangelist>>,
pub turn_allowed_tool_names: Option<Arc<HashSet<String>>>,
pub session_mode: crate::cm_types::SessionMode,
}
pub(crate) struct RunLoopObs {
pub request_chrome_trace: Option<std::sync::Arc<crate::request_chrome_trace::RequestTurnTrace>>,
pub tracing_chat_turn: Option<Arc<crate::observability::TracingChatTurn>>,
pub request_audit: Option<Arc<WebRequestAudit>>,
pub process_handles: Arc<crate::process_handles::TurnProcessHandles>,
pub trace_sink: Option<Arc<dyn crate::cm_llm::TraceSink>>,
pub tool_job_registry: Option<std::sync::Arc<crate::cm_internal::tool_jobs::ToolJobRegistry>>,
}
pub(crate) struct RunLoopCtx<'a> {
pub core: RunLoopCore<'a>,
pub io: RunLoopIo<'a>,
pub attach: RunLoopAttach<'a>,
pub obs: RunLoopObs,
}
#[derive(Debug, Clone, Default)]
pub(crate) struct TurnPlannerHints {
pub(crate) execution_constraint_hint: Option<String>,
pub(crate) step_executor_constraint: Option<PlanStepExecutorKind>,
pub(crate) turn_start_snapshot: Option<crate::cm_agent::agent_turn::TurnStartSnapshot>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum OuterLoopPlanCallModelRole {
PlannerRound,
ExecutorRound,
}
impl OuterLoopPlanCallModelRole {
#[inline]
pub(crate) fn from_outer_loop_iteration(iteration_count: u32) -> Self {
if iteration_count <= 1 {
Self::PlannerRound
} else {
Self::ExecutorRound
}
}
#[inline]
pub(crate) fn sets_use_executor_model(self) -> bool {
matches!(self, Self::ExecutorRound)
}
#[inline]
pub(crate) fn as_trace_str(self) -> &'static str {
match self {
Self::PlannerRound => "planner_round",
Self::ExecutorRound => "executor_round",
}
}
}
impl TurnPlannerHints {
pub(crate) fn take_execution_constraint_hint(&mut self) -> Option<String> {
self.execution_constraint_hint.take()
}
}
pub(crate) struct RunLoopTurnState<'a> {
pub(crate) messages_buf: &'a mut Vec<Message>,
pub(crate) messages_revision: u64,
pub sub_phase: AgentTurnSubPhase,
pub(crate) turn_planner_hints: TurnPlannerHints,
pub temperature_override: Option<f32>,
pub model_override: Option<String>,
pub use_executor_model: bool,
pub executor_model_override: Option<String>,
pub executor_api_base: Option<String>,
pub executor_api_key: Option<String>,
pub seed_override: LlmSeedOverride,
pub turn_budget: Arc<TurnBudgetCounter>,
pub(crate) context_timeline: crate::cm_agent::context_timeline::ContextTimelineAcc,
pub(crate) model_context_artifacts:
Vec<crate::agent::model_context_view::ModelContextArtifact>,
pub(crate) provider_usage:
Arc<std::sync::Mutex<Option<crate::cm_types::Usage>>>,
}
impl<'a> RunLoopTurnState<'a> {
#[inline]
fn bump_messages_revision(&mut self) {
self.messages_revision = self.messages_revision.wrapping_add(1);
}
#[inline]
pub(crate) fn messages(&self) -> &[Message] {
self.messages_buf
}
#[inline]
pub(crate) fn messages_buffer_mut(&mut self) -> &mut Vec<Message> {
self.messages_buf
}
#[inline]
pub(crate) fn messages_buffer_revision(&self) -> u64 {
self.messages_revision
}
pub(crate) fn push_message(&mut self, msg: Message) {
self.messages_buf.push(msg);
self.bump_messages_revision();
}
pub(crate) fn pop_message(&mut self) -> Option<Message> {
let r = self.messages_buf.pop();
if r.is_some() {
self.bump_messages_revision();
}
r
}
#[cfg(test)]
pub(crate) fn truncate_messages(&mut self, len: usize) {
if self.messages_buf.len() != len {
self.messages_buf.truncate(len);
self.bump_messages_revision();
}
}
pub(crate) fn retain_messages(&mut self, mut keep: impl FnMut(&Message) -> bool) {
let before = self.messages_buf.len();
self.messages_buf.retain(|m| keep(m));
if self.messages_buf.len() != before {
self.bump_messages_revision();
}
}
pub(crate) fn push_assistant_merging_trailing_empty(&mut self, msg: Message) {
push_assistant_merging_trailing_empty_placeholder(self.messages_buf, msg);
self.bump_messages_revision();
}
pub(crate) fn flush_context_timeline_markers(&mut self) {
let markers = self.context_timeline.persist_markers();
if markers.is_empty() {
return;
}
crate::cm_agent::context_timeline::strip_context_window_timeline_markers(self.messages_buf);
self.messages_buf.extend(markers);
self.bump_messages_revision();
}
pub(crate) fn insert_separator_after_last_user_for_turn(&mut self) {
let n = self.messages_buf.len();
insert_separator_after_last_user_for_turn(self.messages_buf);
if self.messages_buf.len() != n {
self.bump_messages_revision();
}
}
pub(crate) fn take_execution_constraint_hint(&mut self) -> Option<String> {
self.turn_planner_hints.take_execution_constraint_hint()
}
}
pub(crate) struct RunLoopParams<'a> {
pub ctx: RunLoopCtx<'a>,
pub turn: RunLoopTurnState<'a>,
}
impl RunLoopParams<'_> {
#[inline]
pub(crate) fn apply_outer_loop_plan_call_model_role(
&mut self,
role: OuterLoopPlanCallModelRole,
) {
self.turn.use_executor_model = role.sets_use_executor_model();
}
#[inline]
pub(crate) fn plan_call_executor_endpoint_cloned(&self) -> (Option<String>, Option<String>) {
if self.turn.use_executor_model {
(
self.turn.executor_api_base.clone(),
self.turn.executor_api_key.clone(),
)
} else {
(None, None)
}
}
#[inline]
pub(crate) fn effective_model(&self) -> Option<&str> {
if self.turn.use_executor_model {
self.turn
.executor_model_override
.as_deref()
.or_else(|| self.ctx.core.cfg.llm.executor_model.as_deref())
} else {
self.turn
.model_override
.as_deref()
.or_else(|| self.ctx.core.cfg.llm.planner_model.as_deref())
}
}
pub(crate) async fn prepare_turn_messages_for_model(
&mut self,
tools_defs: &[crate::types::Tool],
per_coord_layer_cache: Option<&mut crate::agent::per_coord::PerCoordinator>,
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
let model_override = self.effective_model().map(str::to_owned);
let delta = crate::agent::context_window::prepare_messages_for_model(
self.ctx.core.llm_backend,
self.ctx.core.client,
self.ctx.core.api_key,
self.ctx.core.cfg.as_ref(),
self.turn.messages_buf,
crate::agent::context_window::PrepareMessagesForModelHooks {
tools: tools_defs,
model_override: model_override.as_deref(),
workspace_changelist: self
.ctx
.attach
.workspace_changelist
.as_ref()
.map(|a| a.as_ref()),
per_coord_layer_cache,
run_loop_messages_revision: Some(&mut self.turn.messages_revision),
turn_budget: Some(&self.turn.turn_budget),
},
)
.await?;
let model_call_sequence = self.turn.model_context_artifacts.len().saturating_add(1);
self.turn.model_context_artifacts.push(
crate::agent::model_context_view::ModelContextArtifact::capture(
model_call_sequence,
0,
self.turn.messages_buf,
delta.compaction,
delta.summary_tail_kept,
),
);
let events = self.turn.context_timeline.merge(
crate::cm_agent::context_timeline::ContextTimelineSnapshot {
messages: self.turn.messages_buf,
pipeline: delta.pipeline,
summarized: delta.summarized,
summary_tail_kept: delta.summary_tail_kept,
compaction:
crate::cm_agent::context_timeline::ContextCompactionTimelineSnapshot {
before_tokens: delta.compaction.before.used_input_tokens,
after_tokens: delta.compaction.after.used_input_tokens,
max_input_tokens: delta
.compaction
.budget
.map(|budget| budget.max_input_tokens),
reserved_output_tokens: delta
.compaction
.budget
.map(|budget| budget.reserved_output_tokens),
message_tokens: delta.compaction.after.message_tokens,
tool_schema_tokens: delta.compaction.after.tool_schema_tokens,
attachment_tokens: delta.compaction.after.attachment_tokens,
counting_source: delta
.compaction
.after
.counting_source
.map(crate::agent::context_compaction::ContextTokenCountingSource::as_str),
token_triggered: delta.compaction.token_triggered,
removed_turn_groups: delta.compaction.removed_turn_groups,
removed_messages: delta.compaction.removed_messages,
compaction_reason: delta.compaction.compaction_reason(),
},
},
);
crate::agent::agent_turn::turn_loop::context_timeline_sse::emit_context_timeline_sse(
&self.ctx.io.control,
&events,
)
.await;
Ok(())
}
}
#[cfg(test)]
mod turn_planner_hints_tests {
use super::{OuterLoopPlanCallModelRole, TurnPlannerHints};
#[test]
fn take_execution_constraint_hint_drains_once() {
let mut h = TurnPlannerHints {
execution_constraint_hint: Some("hint".into()),
..Default::default()
};
assert_eq!(h.take_execution_constraint_hint().as_deref(), Some("hint"));
assert!(h.take_execution_constraint_hint().is_none());
}
#[test]
fn outer_loop_plan_role_matches_iteration_and_trace() {
assert_eq!(
OuterLoopPlanCallModelRole::from_outer_loop_iteration(1),
OuterLoopPlanCallModelRole::PlannerRound
);
assert!(!OuterLoopPlanCallModelRole::PlannerRound.sets_use_executor_model());
assert_eq!(
OuterLoopPlanCallModelRole::PlannerRound.as_trace_str(),
"planner_round"
);
assert_eq!(
OuterLoopPlanCallModelRole::from_outer_loop_iteration(2),
OuterLoopPlanCallModelRole::ExecutorRound
);
assert!(OuterLoopPlanCallModelRole::ExecutorRound.sets_use_executor_model());
assert_eq!(
OuterLoopPlanCallModelRole::ExecutorRound.as_trace_str(),
"executor_round"
);
}
#[test]
fn messages_revision_increments_on_buffer_mutations() {
use crate::agent::agent_turn::errors::AgentTurnSubPhase;
use crate::types::{LlmSeedOverride, Message};
let mut storage = vec![Message::user_only("u")];
let mut turn = super::RunLoopTurnState {
messages_buf: &mut storage,
messages_revision: 0,
sub_phase: AgentTurnSubPhase::Planner,
turn_planner_hints: TurnPlannerHints::default(),
temperature_override: None,
model_override: None,
use_executor_model: false,
executor_model_override: None,
executor_api_base: None,
executor_api_key: None,
seed_override: LlmSeedOverride::FromConfig,
turn_budget: crate::agent::turn_budget::TurnBudgetCounter::new_shared(),
context_timeline: Default::default(),
model_context_artifacts: Vec::new(),
provider_usage: std::sync::Arc::new(std::sync::Mutex::new(None)),
};
assert_eq!(turn.messages_buffer_revision(), 0);
turn.push_message(Message::assistant_only("a"));
assert_eq!(turn.messages_buffer_revision(), 1);
turn.truncate_messages(1);
assert_eq!(turn.messages_buffer_revision(), 2);
turn.retain_messages(|_| true);
assert_eq!(turn.messages_buffer_revision(), 2);
turn.retain_messages(|m| m.role != "tool");
assert_eq!(turn.messages_buffer_revision(), 2);
}
}