use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use foundation_core::valtron::{
BoxedSendExecutionAction, CancelOutcome, CancellableFutureTask, DrivenTaskIterator, Stream,
TaskIterator, TaskStatus,
};
use foundation_db::traits::DocumentStore;
use crate::agentic::context::{AgentContext, ContextProvider};
use crate::agentic::errors::{AgentAction, AgenticError, CircuitBreaker, ErrorPolicy};
use crate::agentic::loop_detection::{
is_vacuous_answer, Escalation, LoopDetection, LoopDetector, LoopDetectorConfig,
};
use crate::agentic::memory::{MemoryAction, MemoryHierarchy};
use crate::agentic::memory_store::MemoryStore;
use crate::agentic::message_api::MessageApi;
use crate::agentic::progress::{lift_model_item, AgentProgress, MemoryKind};
use crate::agentic::steering::SteeringQueues;
use crate::agentic::token_ledger::TokenLedger;
use crate::agentic::tool_impl::{ToolCallManager, ToolCallRequest, ToolCallResult, ToolError};
use crate::types::{
MessageRole, Messages, ModelId, ModelInteraction, ModelOutput, ModelParams, ModelState,
SessionId, SessionRecord, TextContent, UserModelContent,
};
#[derive(Debug, Clone)]
pub struct AgentConfig {
pub primary_model: ModelId,
pub fallback_models: Vec<ModelId>,
pub memory_model: Option<ModelId>,
pub max_inner_iterations: usize,
pub max_outer_iterations: usize,
pub circuit_breaker_threshold: u32,
pub preflight_compression_threshold: f32,
pub context_pressure_threshold: f32,
pub model_params: ModelParams,
}
impl Default for AgentConfig {
fn default() -> Self {
Self {
primary_model: ModelId::Name(String::new(), None),
fallback_models: Vec::new(),
memory_model: None,
max_inner_iterations: 25,
max_outer_iterations: 10,
circuit_breaker_threshold: 3,
preflight_compression_threshold: 0.85,
context_pressure_threshold: 0.70,
model_params: ModelParams::default(),
}
}
}
type ToolDrivenIterator = DrivenTaskIterator<
CancellableFutureTask<
core::pin::Pin<
Box<dyn core::future::Future<Output = Result<ToolCallResult, ToolError>> + Send>,
>,
>,
>;
pub enum AgentLoopState {
Initializing,
OuterBoundary,
InnerAssemble,
InnerGenerate {
stream: Box<dyn Iterator<Item = Stream<Messages, ModelState>> + Send>,
collected: Vec<Messages>,
},
InnerToolCalls { calls: Vec<ToolCallRequest> },
InnerExecuting {
calls: Vec<ToolCallRequest>,
results: Vec<(ToolCallRequest, Result<ToolCallResult, ToolError>)>,
idx: usize,
active: Option<ToolDrivenIterator>,
cancel_signals: Vec<Arc<AtomicBool>>,
},
InnerEmitResults { results: Vec<Messages>, idx: usize },
OutputProcessing,
Ending,
Done,
}
impl AgentLoopState {
pub(crate) const fn variant_name(&self) -> &'static str {
match self {
Self::Initializing => "Initializing",
Self::OuterBoundary => "OuterBoundary",
Self::InnerAssemble => "InnerAssemble",
Self::InnerGenerate { .. } => "InnerGenerate",
Self::InnerToolCalls { .. } => "InnerToolCalls",
Self::InnerExecuting { .. } => "InnerExecuting",
Self::InnerEmitResults { .. } => "InnerEmitResults",
Self::OutputProcessing => "OutputProcessing",
Self::Ending => "Ending",
Self::Done => "Done",
}
}
}
pub struct AgentLoop<D, M> {
session_id: SessionId,
context_provider: ContextProvider<D, M>,
tool_manager: ToolCallManager,
queues: SteeringQueues,
memory: MemoryHierarchy<M, D>,
message_api: MessageApi<D>,
ledger: TokenLedger,
detector: LoopDetector,
policy: ErrorPolicy,
breaker: CircuitBreaker,
router: crate::types::ProviderRouter,
state: AgentLoopState,
current_model: ModelId,
config: AgentConfig,
inner_iteration: usize,
outer_iteration: usize,
last_user_prompt: String,
message_count: u64,
}
fn assembled_answer(collected: &[Messages]) -> String {
let mut answer = String::new();
for message in collected {
if let Messages::Assistant {
content: ModelOutput::Text(text),
..
} = message
{
answer.push_str(&text.content);
}
}
answer
}
fn vacuous_answer_redirect() -> Messages {
Messages::User {
id: foundation_compact::ids::new_scru128(),
role: MessageRole::System,
content: UserModelContent::Text(TextContent {
content: "Your last reply contained no answer - only punctuation or \
whitespace. Answer the user's question directly, in plain \
words."
.into(),
signature: None,
}),
signature: None,
}
}
impl<D: DocumentStore, M: MemoryStore> AgentLoop<D, M> {
#[allow(clippy::too_many_arguments)]
#[must_use]
pub fn new(
session_id: SessionId,
context_provider: ContextProvider<D, M>,
tool_manager: ToolCallManager,
queues: SteeringQueues,
memory: MemoryHierarchy<M, D>,
message_api: MessageApi<D>,
ledger: TokenLedger,
policy: ErrorPolicy,
router: crate::types::ProviderRouter,
config: AgentConfig,
) -> Self {
let breaker = CircuitBreaker::new(
config.circuit_breaker_threshold,
config.fallback_models.clone(),
);
let current_model = config.primary_model.clone();
Self {
session_id,
context_provider,
tool_manager,
queues,
memory,
message_api,
ledger,
detector: LoopDetector::new(LoopDetectorConfig::default()),
policy,
breaker,
router,
state: AgentLoopState::Initializing,
current_model,
config,
inner_iteration: 0,
outer_iteration: 0,
last_user_prompt: String::new(),
message_count: 0,
}
}
pub fn push_user_message(&mut self, msg: Messages) {
let _ = self
.message_api
.append(SessionRecord::Conversation { message: msg });
}
#[must_use]
pub fn state_label(&self) -> &'static str {
match &self.state {
AgentLoopState::Initializing => "initializing",
AgentLoopState::OuterBoundary => "outer_boundary",
AgentLoopState::InnerAssemble => "inner_assemble",
AgentLoopState::InnerGenerate { .. } => "inner_generate",
AgentLoopState::InnerToolCalls { .. } => "inner_tool_calls",
AgentLoopState::InnerExecuting { .. } => "inner_executing",
AgentLoopState::InnerEmitResults { .. } => "inner_emit_results",
AgentLoopState::OutputProcessing => "output_processing",
AgentLoopState::Ending => "ending",
AgentLoopState::Done => "done",
}
}
#[must_use]
pub fn current_model(&self) -> &ModelId {
&self.current_model
}
fn transition_outer_boundary(
&mut self,
) -> TaskStatus<SessionRecord, AgentProgress, BoxedSendExecutionAction> {
if self.queues.is_aborted() {
self.queues.reset_cancel();
self.state = AgentLoopState::Ending;
return TaskStatus::Pending(AgentProgress::SessionEnding);
}
self.outer_iteration += 1;
if self.outer_iteration > self.config.max_outer_iterations {
self.state = AgentLoopState::Ending;
return TaskStatus::Pending(AgentProgress::SessionEnding);
}
let priority_msgs = self.queues.drain_priority();
if !priority_msgs.is_empty() {
self.queues.reset_cancel();
for msg in priority_msgs {
self.push_user_message(msg);
}
self.inner_iteration = 0;
self.state = AgentLoopState::InnerAssemble;
return TaskStatus::Pending(AgentProgress::Steering {
source: std::borrow::Cow::Borrowed("priority_queue"),
});
}
let follow_up_msgs = self.queues.drain_follow_up();
tracing::trace!(
drained = follow_up_msgs.len(),
outer_iteration = self.outer_iteration,
"outer_boundary: drained follow_up queue"
);
if !follow_up_msgs.is_empty() {
for msg in follow_up_msgs {
if let Messages::User {
role: MessageRole::User,
content: UserModelContent::Text(ref text),
..
} = msg
{
self.last_user_prompt.clone_from(&text.content);
}
self.push_user_message(msg);
}
self.inner_iteration = 0;
self.state = AgentLoopState::InnerAssemble;
return TaskStatus::Pending(AgentProgress::Steering {
source: std::borrow::Cow::Borrowed("follow_up_queue"),
});
}
self.state = AgentLoopState::Ending;
TaskStatus::Pending(AgentProgress::SessionEnding)
}
fn transition_inner_assemble(
&mut self,
) -> TaskStatus<SessionRecord, AgentProgress, BoxedSendExecutionAction> {
if self.queues.is_aborted() {
self.queues.reset_cancel();
self.state = AgentLoopState::Ending;
return TaskStatus::Pending(AgentProgress::SessionEnding);
}
if self.queues.has_priority() {
let msgs = self.queues.drain_priority();
self.queues.reset_cancel();
for msg in msgs {
self.push_user_message(msg);
}
return TaskStatus::Pending(AgentProgress::Steering {
source: std::borrow::Cow::Borrowed("priority_interrupt"),
});
}
if self.ledger.is_exhausted() {
let snapshot = self.ledger.snapshot();
self.state = AgentLoopState::Ending;
return TaskStatus::Ready(
AgenticError::BudgetExhausted { snapshot }.into_failed_action(),
);
}
let model = match self.router.get_model(&self.current_model) {
Ok(m) => m,
Err(e) => {
tracing::error!("router failed to get model: {e:?}");
let err = AgenticError::from(e);
return self.handle_error(err);
}
};
let memory = self
.context_provider
.memory_store()
.hydrate_sync(&self.session_id)
.unwrap_or_default();
let mut ctx = self.context_provider.assemble_from_memory(&memory);
self.apply_preflight_compression(&mut ctx);
let system_prompt = self.apply_context_pressure(&ctx);
let toolshed = self.tool_manager.build_toolshed();
let interaction = ModelInteraction {
system_prompt: system_prompt.or(ctx.system_prompt),
soul: None,
tools_shed: toolshed,
messages: ctx.messages,
chat_template: None,
tool_choice: None,
};
let params = self.config.model_params.clone();
let effective_max = self.ledger.effective_max_tokens(¶ms);
let params = ModelParams {
max_tokens: effective_max,
..params
};
tracing::trace!(
messages = interaction.messages.len(),
has_system = interaction.system_prompt.is_some(),
"Sending interactions to model for generation"
);
match model.stream(interaction, Some(params)) {
Ok(stream) => {
self.state = AgentLoopState::InnerGenerate {
stream: Box::new(stream),
collected: Vec::new(),
};
TaskStatus::Pending(AgentProgress::Generating {
model: self.current_model.clone(),
tokens_so_far: None,
})
}
Err(gen_err) => {
let err = AgenticError::from_generation(&gen_err, None, 0);
self.handle_error(err)
}
}
}
fn transition_inner_generate(
&mut self,
) -> TaskStatus<SessionRecord, AgentProgress, BoxedSendExecutionAction> {
let previous = std::mem::replace(&mut self.state, AgentLoopState::Done);
let actual = previous.variant_name();
let AgentLoopState::InnerGenerate {
mut stream,
mut collected,
} = previous
else {
panic!("transition_inner_generate requires AgentLoopState::InnerGenerate, found {actual}")
};
let Some(item) = stream.next() else {
self.breaker.on_success();
return self.on_generation_complete(&collected);
};
if self.queues.has_priority() {
self.queues.reset_cancel();
let msgs = self.queues.drain_priority();
for msg in msgs {
self.push_user_message(msg);
}
self.state = AgentLoopState::InnerAssemble;
return TaskStatus::Pending(AgentProgress::Steering {
source: std::borrow::Cow::Borrowed("mid_gen_priority"),
});
}
let lifted = lift_model_item(item, &self.current_model);
match lifted {
Stream::Next(record) => {
if let SessionRecord::Conversation { ref message } = record {
collected.push(message.clone());
self.message_count += 1;
}
self.state = AgentLoopState::InnerGenerate { stream, collected };
TaskStatus::Ready(record)
}
Stream::Pending(progress) => {
self.state = AgentLoopState::InnerGenerate { stream, collected };
TaskStatus::Pending(progress)
}
Stream::Init | Stream::Ignore | Stream::Wait => {
self.state = AgentLoopState::InnerGenerate { stream, collected };
TaskStatus::Ignore
}
Stream::Delayed(d) => {
self.state = AgentLoopState::InnerGenerate { stream, collected };
TaskStatus::Delayed(d)
}
Stream::Spread(items) => {
use foundation_core::valtron::StreamSpread;
for s in &items {
if let StreamSpread::Done(SessionRecord::Conversation { ref message }) = s {
collected.push(message.clone());
self.message_count += 1;
}
}
self.state = AgentLoopState::InnerGenerate { stream, collected };
if let Some(first) = items.into_iter().next() {
match first {
StreamSpread::Done(rec) => TaskStatus::Ready(rec),
StreamSpread::Pending(p) => TaskStatus::Pending(p),
}
} else {
TaskStatus::Ignore
}
}
}
}
fn on_generation_complete(
&mut self,
collected: &[Messages],
) -> TaskStatus<SessionRecord, AgentProgress, BoxedSendExecutionAction> {
for msg in collected {
if let Messages::Assistant { ref content, .. } = msg {
let detection = self.detector.check(content);
if detection != LoopDetection::NoLoop {
let escalation = self.detector.escalate();
match escalation {
Escalation::Redirect => {
let redirect = Messages::User {
id: foundation_compact::ids::new_scru128(),
role: MessageRole::System,
content: UserModelContent::Text(TextContent {
content: "Loop detected. Please try a different approach — \
avoid repeating the same actions or responses."
.into(),
signature: None,
}),
signature: None,
};
self.push_user_message(redirect);
self.state = AgentLoopState::InnerAssemble;
return TaskStatus::Pending(AgentProgress::Steering {
source: std::borrow::Cow::Borrowed("loop_redirect"),
});
}
Escalation::SwitchModelOrTemperature { .. } => {
if let Some(fallback) = self.breaker.on_failure() {
self.current_model = fallback;
}
let redirect = Messages::User {
id: foundation_compact::ids::new_scru128(),
role: MessageRole::System,
content: UserModelContent::Text(TextContent {
content: "Loop detected after redirect. Switching approach."
.into(),
signature: None,
}),
signature: None,
};
self.push_user_message(redirect);
self.state = AgentLoopState::InnerAssemble;
return TaskStatus::Pending(AgentProgress::Steering {
source: std::borrow::Cow::Borrowed("loop_model_switch"),
});
}
Escalation::Terminate => {
if is_vacuous_answer(&assembled_answer(collected)) {
tracing::debug!(
"agent: repeated empty answer, passing it \
through rather than failing the turn"
);
break;
}
#[allow(clippy::cast_possible_truncation)]
let occurrences = self.detector.redirect_count() as u32;
let err =
AgenticError::LoopDetected(crate::agentic::errors::LoopDetection {
kind: format!("{detection:?}"),
occurrences,
});
self.state = AgentLoopState::Ending;
return TaskStatus::Ready(err.into_failed_action());
}
}
}
}
}
let called_a_tool = collected.iter().any(|message| {
matches!(
message,
Messages::Assistant {
content: ModelOutput::ToolCall { .. },
..
}
)
});
let mut gave_up_on_a_bad_answer = false;
if !called_a_tool {
let answer = assembled_answer(collected);
if self.detector.check_answer(&answer, &self.last_user_prompt) != LoopDetection::NoLoop
{
match self.detector.escalate() {
Escalation::Redirect | Escalation::SwitchModelOrTemperature { .. } => {
tracing::debug!(
answer = %answer,
attempt = self.detector.redirect_count(),
"agent: turn produced no usable answer, asking again"
);
self.push_user_message(vacuous_answer_redirect());
self.state = AgentLoopState::InnerAssemble;
return TaskStatus::Ready(SessionRecord::Retracted {
id: foundation_compact::ids::new_scru128(),
reason: format!(
"turn produced no usable answer ({answer:?}), retrying"
),
timestamp: std::time::SystemTime::now(),
});
}
Escalation::Terminate => {
gave_up_on_a_bad_answer = true;
tracing::debug!(
answer = %answer,
"agent: still no usable answer after retries, passing it through"
);
}
}
}
}
if !gave_up_on_a_bad_answer {
self.detector.reset();
}
for msg in collected {
if matches!(msg, Messages::Assistant { .. }) {
let _ = self.message_api.append(SessionRecord::Conversation {
message: msg.clone(),
});
}
}
if let Some(Messages::Assistant { usage, .. }) = collected
.iter()
.rev()
.find(|m| matches!(m, Messages::Assistant { .. }))
{
self.ledger.record(usage);
}
let tool_calls = Self::extract_tool_calls(collected);
if tool_calls.is_empty() {
self.state = AgentLoopState::OutputProcessing;
TaskStatus::Ignore
} else {
let total = tool_calls.len();
self.state = AgentLoopState::InnerToolCalls { calls: tool_calls };
TaskStatus::Pending(AgentProgress::ExecutingTools {
total,
completed: 0,
})
}
}
fn extract_tool_calls(messages: &[Messages]) -> Vec<ToolCallRequest> {
let mut calls = Vec::new();
for msg in messages {
if let Messages::Assistant {
content:
ModelOutput::ToolCall {
id,
name,
arguments,
depends_on,
execution_hint,
..
},
..
} = msg
{
calls.push(ToolCallRequest {
id: id.clone(),
name: name.clone(),
arguments: arguments.clone().unwrap_or_default(),
depends_on: depends_on.clone(),
execution_hint: *execution_hint,
});
}
}
calls
}
fn transition_inner_tool_calls(
&mut self,
) -> TaskStatus<SessionRecord, AgentProgress, BoxedSendExecutionAction> {
let previous = std::mem::replace(&mut self.state, AgentLoopState::Done);
let actual = previous.variant_name();
let AgentLoopState::InnerToolCalls { calls } = previous else {
panic!(
"transition_inner_tool_calls requires AgentLoopState::InnerToolCalls, found {actual}"
)
};
let workflow = match self.tool_manager.build_workflow(&calls) {
Ok(w) => w,
Err(e) => {
let err = AgenticError::ToolCall {
tool_name: "workflow".into(),
reason: e.to_string(),
};
self.state = AgentLoopState::OutputProcessing;
return TaskStatus::Ready(err.into_failed_action());
}
};
let mut flat_calls = Vec::new();
for stage in workflow.stages {
match stage {
crate::agentic::tool_impl::ToolCallStage::Parallel { calls: sc, .. }
| crate::agentic::tool_impl::ToolCallStage::Sequential { calls: sc, .. } => {
flat_calls.extend(sc);
}
}
}
if flat_calls.is_empty() {
self.state = AgentLoopState::OutputProcessing;
return TaskStatus::Ignore;
}
self.state = AgentLoopState::InnerExecuting {
calls: flat_calls,
results: Vec::new(),
idx: 0,
active: None,
cancel_signals: Vec::new(),
};
TaskStatus::Ignore
}
fn transition_inner_executing(
&mut self,
) -> TaskStatus<SessionRecord, AgentProgress, BoxedSendExecutionAction> {
let previous = std::mem::replace(&mut self.state, AgentLoopState::Done);
let actual = previous.variant_name();
let AgentLoopState::InnerExecuting {
calls,
mut results,
idx,
mut active,
mut cancel_signals,
} = previous
else {
panic!(
"transition_inner_executing requires AgentLoopState::InnerExecuting, found {actual}"
)
};
if self.queues.has_priority() {
for sig in &cancel_signals {
sig.store(true, Ordering::Release);
}
let msgs = self.queues.drain_priority();
self.queues.reset_cancel();
for msg in msgs {
self.push_user_message(msg);
}
self.state = AgentLoopState::InnerAssemble;
return TaskStatus::Pending(AgentProgress::Steering {
source: std::borrow::Cow::Borrowed("cancel_executing"),
});
}
if idx >= calls.len() {
let result_messages: Vec<Messages> = results
.into_iter()
.map(|(req, result)| {
let (content, error_detail) = match result {
Ok(r) => (r.content, r.error_detail),
Err(e) => (
UserModelContent::Text(TextContent {
content: e.to_string(),
signature: None,
}),
Some(e.to_string()),
),
};
Messages::ToolResult {
id: foundation_compact::ids::new_scru128(),
tool_call_id: req.id,
name: req.name,
timestamp: foundation_compact::SystemTime::now(),
details: None,
content,
error_detail,
signature: None,
}
})
.collect();
if result_messages.is_empty() {
self.state = AgentLoopState::OutputProcessing;
return TaskStatus::Ignore;
}
self.state = AgentLoopState::InnerEmitResults {
results: result_messages,
idx: 0,
};
return TaskStatus::Ignore;
}
if active.is_none() {
let call = calls[idx].clone();
let mgr = self.tool_manager.clone();
let retry_config = mgr.retry_config(&call.name);
let fut: core::pin::Pin<
Box<dyn core::future::Future<Output = Result<ToolCallResult, ToolError>> + Send>,
> = Box::pin(async move { mgr.execute_with_retry(&call, &retry_config).await });
let signal = Arc::new(AtomicBool::new(false));
cancel_signals.push(Arc::clone(&signal));
let task = CancellableFutureTask::new(fut, signal);
active = Some(foundation_core::valtron::drive_iterator(task));
}
let total = calls.len();
let driven = active.as_mut().expect("just created");
match driven.next() {
Some(TaskStatus::Ready(Ok(result))) => {
results.push((calls[idx].clone(), result));
self.state = AgentLoopState::InnerExecuting {
calls,
results,
idx: idx + 1,
active: None,
cancel_signals,
};
TaskStatus::Pending(AgentProgress::ExecutingTools {
total,
completed: idx + 1,
})
}
Some(TaskStatus::Ready(Err(CancelOutcome::Cancelled))) => {
let err = ToolError::Cancelled(calls[idx].name.clone());
results.push((calls[idx].clone(), Err(err)));
self.state = AgentLoopState::InnerExecuting {
calls,
results,
idx: idx + 1,
active: None,
cancel_signals,
};
TaskStatus::Pending(AgentProgress::ExecutingTools {
total,
completed: idx + 1,
})
}
Some(TaskStatus::Pending(_)) => {
self.state = AgentLoopState::InnerExecuting {
calls,
results,
idx,
active,
cancel_signals,
};
TaskStatus::Pending(AgentProgress::ExecutingTools {
total,
completed: idx,
})
}
None => {
let err = ToolError::Execution {
tool: calls[idx].name.clone(),
reason: "future completed without result".into(),
};
results.push((calls[idx].clone(), Err(err)));
self.state = AgentLoopState::InnerExecuting {
calls,
results,
idx: idx + 1,
active: None,
cancel_signals,
};
TaskStatus::Ignore
}
_ => {
self.state = AgentLoopState::InnerExecuting {
calls,
results,
idx,
active,
cancel_signals,
};
TaskStatus::Ignore
}
}
}
fn transition_inner_emit_results(
&mut self,
) -> TaskStatus<SessionRecord, AgentProgress, BoxedSendExecutionAction> {
let previous = std::mem::replace(&mut self.state, AgentLoopState::Done);
let actual = previous.variant_name();
let AgentLoopState::InnerEmitResults { results, idx } = previous else {
panic!(
"transition_inner_emit_results requires AgentLoopState::InnerEmitResults, found {actual}"
)
};
if idx >= results.len() {
tracing::trace!(
"Loop index greater tan results len: idx={}, len={}",
idx,
results.len()
);
self.inner_iteration += 1;
if self.inner_iteration >= self.config.max_inner_iterations {
tracing::trace!("Max loop iteration method");
self.state = AgentLoopState::OutputProcessing;
return TaskStatus::Pending(AgentProgress::SessionEnding);
}
self.state = AgentLoopState::InnerAssemble;
return TaskStatus::Ignore;
}
let msg = results[idx].clone();
tracing::trace!("Adding new message from model: {:?}", &msg);
let _ = self.message_api.append(SessionRecord::Conversation {
message: msg.clone(),
});
self.message_count += 1;
self.state = AgentLoopState::InnerEmitResults {
results,
idx: idx + 1,
};
TaskStatus::Ready(SessionRecord::Conversation { message: msg })
}
fn transition_output_processing(
&mut self,
) -> TaskStatus<SessionRecord, AgentProgress, BoxedSendExecutionAction> {
let action = self.memory.check_triggers();
self.state = AgentLoopState::OuterBoundary;
match action {
MemoryAction::GenerateObservation => {
TaskStatus::Pending(AgentProgress::ProcessingMemory {
kind: MemoryKind::Observation,
})
}
MemoryAction::GenerateReflection => {
TaskStatus::Pending(AgentProgress::ProcessingMemory {
kind: MemoryKind::Reflection,
})
}
MemoryAction::None => TaskStatus::Ignore,
}
}
fn transition_ending(
&mut self,
) -> TaskStatus<SessionRecord, AgentProgress, BoxedSendExecutionAction> {
let snapshot = self.ledger.snapshot();
let summary = SessionRecord::Summary {
message_count: self.message_count,
usage: snapshot,
};
self.state = AgentLoopState::Done;
TaskStatus::Ready(summary)
}
fn handle_error(
&mut self,
error: AgenticError,
) -> TaskStatus<SessionRecord, AgentProgress, BoxedSendExecutionAction> {
let action = self.policy.classify(error.clone());
match action {
AgentAction::Continue => TaskStatus::Ignore,
AgentAction::RetryWithReducedContext => {
self.state = AgentLoopState::InnerAssemble;
TaskStatus::Pending(AgentProgress::Steering {
source: std::borrow::Cow::Borrowed("retry_reduced_context"),
})
}
AgentAction::SwitchModel => {
if let Some(fallback) = self.breaker.on_failure() {
self.current_model = fallback;
self.state = AgentLoopState::InnerAssemble;
TaskStatus::Pending(AgentProgress::Steering {
source: std::borrow::Cow::Borrowed("model_switch"),
})
} else {
self.state = AgentLoopState::Ending;
TaskStatus::Ready(error.into_failed_action())
}
}
AgentAction::Terminate(err) => {
self.state = AgentLoopState::Ending;
TaskStatus::Ready(err.into_failed_action())
}
}
}
fn apply_preflight_compression(&self, ctx: &mut AgentContext) {
if self.config.preflight_compression_threshold <= 0.0 {
return;
}
let Some(budget) = self.ledger.budget() else {
return; };
if budget == 0 {
return;
}
#[allow(clippy::cast_precision_loss)]
let limit = f64::from(self.config.preflight_compression_threshold) * budget as f64;
while ctx.messages.len() > 1 {
#[allow(clippy::cast_precision_loss)]
let ratio = ctx.token_estimate as f64;
if ratio <= limit {
break;
}
let dropped = ctx.messages.remove(0);
let est = crate::agentic::context::estimate_tokens_pub(&dropped);
ctx.token_estimate = ctx.token_estimate.saturating_sub(est);
}
}
fn apply_context_pressure(&self, ctx: &AgentContext) -> Option<String> {
if self.config.context_pressure_threshold <= 0.0 {
return ctx.system_prompt.clone();
}
let budget = self.ledger.budget().unwrap_or(u64::MAX);
#[allow(clippy::cast_precision_loss)]
let usage_ratio = ctx.token_estimate as f64 / budget as f64;
if usage_ratio >= f64::from(self.config.context_pressure_threshold) {
#[allow(clippy::cast_possible_truncation, clippy::cast_sign_loss)]
let pct = (usage_ratio * 100.0) as u32;
let pressure_note = format!(
"Context is at {pct}% capacity — prefer concise responses \
and avoid requesting large tool outputs."
);
let base = ctx.system_prompt.as_deref().unwrap_or("");
Some(format!("{base}\n\n[{pressure_note}]"))
} else {
ctx.system_prompt.clone()
}
}
}
impl<D, M> TaskIterator for AgentLoop<D, M>
where
D: DocumentStore + 'static,
M: MemoryStore + 'static,
{
type Ready = SessionRecord;
type Pending = AgentProgress;
type Spawner = BoxedSendExecutionAction;
fn next_status(&mut self) -> Option<TaskStatus<Self::Ready, Self::Pending, Self::Spawner>> {
match &self.state {
AgentLoopState::Initializing => {
self.state = AgentLoopState::OuterBoundary;
Some(TaskStatus::Init)
}
AgentLoopState::OuterBoundary => Some(self.transition_outer_boundary()),
AgentLoopState::InnerAssemble => Some(self.transition_inner_assemble()),
AgentLoopState::InnerGenerate { .. } => Some(self.transition_inner_generate()),
AgentLoopState::InnerToolCalls { .. } => Some(self.transition_inner_tool_calls()),
AgentLoopState::InnerExecuting { .. } => Some(self.transition_inner_executing()),
AgentLoopState::InnerEmitResults { .. } => Some(self.transition_inner_emit_results()),
AgentLoopState::OutputProcessing => Some(self.transition_output_processing()),
AgentLoopState::Ending => Some(self.transition_ending()),
AgentLoopState::Done => None,
}
}
}