use std::fmt::Write;
use std::path::Path;
use std::sync::Arc;
use std::time::Instant;
use anyhow::Context;
use tracing::Instrument;
use crate::providers::reasoning::assistant_replay_payload;
use crate::session::Session;
use crate::tools::{
ImagePayload, ToolExecutionOutcome, find_tool, format_tool_failure_feedback,
normalize_tool_call, scrub_tool_output,
};
use crate::util::{
MEDIA_MARKER_RE, MediaMarkerKind, UnwrapPoison, parse_media_marker, scrub_credentials,
};
use crate::{Agent, ChatMessage, ChatRequest, ChatResponse, Tool, ToolCall};
pub(crate) mod extraction;
pub mod maintainer;
pub mod message_router;
pub mod registry;
pub(crate) mod role;
pub(crate) mod skills;
pub(crate) fn tool_user_name() -> String {
CURRENT_TOOL_USER_NAME
.try_with(String::clone)
.unwrap_or_default()
}
pub(crate) fn tool_identity() -> anyhow::Result<(String, String)> {
let user_name = tool_user_name();
let agent_id = CURRENT_TOOL_AGENT_ID
.try_with(Option::clone)
.ok()
.flatten()
.unwrap_or_default();
anyhow::ensure!(
!agent_id.is_empty() && !user_name.is_empty(),
"No agent identity available — the tool must run in a user session."
);
Ok((agent_id, user_name))
}
pub(crate) fn tool_record_attribution() -> (String, String) {
CURRENT_TOOL_AGENT_TRACKING
.try_with(|t| t.as_ref().map(|t| (t.agent_id.clone(), t.role.clone())))
.ok()
.flatten()
.unwrap_or_default()
}
pub(crate) fn tool_channel() -> String {
CURRENT_TOOL_CHANNEL
.try_with(String::clone)
.unwrap_or_default()
}
tokio::task_local! {
pub(crate) static CURRENT_TOOL_USER_NAME: String;
pub(crate) static CURRENT_TOOL_CHANNEL: String;
pub(crate) static CURRENT_TOOL_PARENT_KEY: Option<crate::agent::registry::ParentKey>;
pub(crate) static CURRENT_TOOL_PARENT_LABEL: Option<String>;
pub(crate) static CURRENT_TOOL_BACKGROUND_SESSIONS:
Option<std::sync::Arc<crate::tools::shell::BackgroundSessions>>;
pub(crate) static CURRENT_TOOL_AGENT_ID: Option<String>;
pub(crate) static CURRENT_TOOL_AGENT_TRACKING:
Option<crate::agent::registry::AgentTracking>;
}
const MAX_TOOL_ROUNDS: usize = 1000;
const MAX_STATS_ARG_LENGTH: usize = 500;
pub(crate) const RETRY_EXHAUSTION_MARKER: &str = "exhausted retry budget";
const SUMMARIZE_ATTEMPTS: u32 = 3;
type CompletePendingToolCallsFuture<'a> =
std::pin::Pin<Box<dyn std::future::Future<Output = anyhow::Result<bool>> + Send + 'a>>;
fn extract_media_from_outcomes(
tools: &[Box<dyn Tool>],
tool_calls: &[ToolCall],
outcomes: &[ToolExecutionOutcome],
) -> Vec<(&'static str, String)> {
let mut paths = Vec::new();
for (call, outcome) in tool_calls.iter().zip(outcomes.iter()) {
if outcome.success
&& let Some(marker_prefix) = find_tool(tools, &call.name).and_then(Tool::media_marker)
{
let marker = &marker_prefix[1..marker_prefix.len() - 1];
let mut matched = false;
for caps in MEDIA_MARKER_RE.captures_iter(&outcome.output) {
let (captured_kind, path) = parse_media_marker(&caps);
if captured_kind.token() == marker {
matched = true;
paths.push((marker_prefix, path.to_string()));
}
}
if !matched {
tracing::warn!(
media_tool = %call.name,
marker = %marker_prefix,
"Could not parse media path from tool output — skipping media marker",
);
}
}
}
paths
}
fn user_image_marker_values(history: &[ChatMessage]) -> impl Iterator<Item = &str> {
history
.iter()
.filter(|msg| msg.role == crate::ChatRole::User)
.flat_map(|msg| {
MEDIA_MARKER_RE
.captures_iter(&msg.content)
.filter_map(|caps| {
let (kind, path) = parse_media_marker(&caps);
(kind == MediaMarkerKind::Image).then_some(path)
})
})
}
fn existing_image_marker_values(history: &[ChatMessage]) -> std::collections::HashSet<String> {
user_image_marker_values(history)
.map(str::to_string)
.collect()
}
fn session_image_marker_count(history: &[ChatMessage]) -> usize {
user_image_marker_values(history)
.filter(|value| value.starts_with("data:image/"))
.collect::<std::collections::HashSet<_>>()
.len()
}
async fn derive_image_payload_from_marker(output: &str) -> Option<ImagePayload> {
if !output.contains("[IMAGE:") {
return None;
}
for caps in MEDIA_MARKER_RE.captures_iter(output) {
let (kind, path) = parse_media_marker(&caps);
if kind != MediaMarkerKind::Image {
continue;
}
let p = Path::new(path);
if !p.is_absolute() {
continue;
}
let Ok(meta) = crate::util::local_image_to_compressed_data_uri_with_meta(p).await else {
continue;
};
return Some(ImagePayload::from_compressed_meta(
p,
meta,
None,
crate::tools::ImagePayloadSource::Generated,
));
}
None
}
async fn tool_result_content(
tools: &[Box<dyn Tool>],
call_name: &str,
outcome: ToolExecutionOutcome,
seen: &mut Option<std::collections::HashSet<String>>,
history: &[ChatMessage],
) -> (String, Vec<ImagePayload>) {
let formatted = match find_tool(tools, call_name) {
Some(t) => t.format_output(&outcome.output),
None => crate::util::truncate_tool_output(&outcome.output),
};
let mut payloads: Vec<ImagePayload> = outcome.image_payloads;
if payloads.is_empty() && !outcome.text_is_content && outcome.success {
payloads.extend(derive_image_payload_from_marker(&outcome.output).await);
}
if payloads.is_empty() {
return (formatted, Vec::new());
}
let seen = seen.get_or_insert_with(|| existing_image_marker_values(history));
let mut fresh = Vec::with_capacity(payloads.len());
let mut annotations = Vec::with_capacity(payloads.len());
for payload in payloads {
if seen.insert(payload.data_uri.clone()) {
annotations.push(payload.attached_annotation());
fresh.push(payload);
} else {
annotations.push(payload.already_attached_annotation());
}
}
let annotation_block = annotations.join("\n");
let output = if outcome.text_is_content {
format!("{formatted}\n{annotation_block}")
} else {
annotation_block
};
(output, fresh)
}
fn injected_images_message(
agent_id: &str,
role: crate::Role,
payloads: Vec<ImagePayload>,
) -> Option<ChatMessage> {
if payloads.is_empty() {
return None;
}
let mut data_uris = Vec::with_capacity(payloads.len());
for payload in payloads {
tracing::debug!(
agent_id,
%role,
path = %payload.path,
"Injecting tool image as a synthetic user message"
);
data_uris.push(payload.data_uri);
}
Some(ChatMessage::user(crate::util::injected_image_user_message(
&data_uris,
)))
}
#[must_use]
pub(crate) fn role_tools_and_specs(
role: crate::Role,
ws: &crate::Workspace,
is_admin: bool,
chrome_sessions: std::sync::Arc<crate::tools::chrome::ChromeRunSessions>,
) -> (Vec<Box<dyn Tool>>, Vec<crate::ToolSpec>) {
let tools: Vec<Box<dyn Tool>> = role
.tools(ws, is_admin, chrome_sessions)
.into_iter()
.filter(|t| t.is_advertised())
.collect();
let tool_specs = tools.iter().map(|t| t.spec()).collect();
(tools, tool_specs)
}
#[must_use]
pub(crate) fn role_tool_specs(
role: crate::Role,
ws: &crate::Workspace,
is_admin: bool,
) -> Vec<crate::ToolSpec> {
role_tools_and_specs(
role,
ws,
is_admin,
std::sync::Arc::new(crate::tools::chrome::ChromeRunSessions::default()),
)
.1
}
pub(crate) struct RoleChatParams {
pub model: String,
pub provider_order: Option<String>,
pub reasoning_effort: Option<String>,
}
#[must_use]
pub(crate) fn role_chat_params(role: crate::Role) -> RoleChatParams {
let model = crate::config::CONFIG.role_model(role);
let routing = crate::config::CONFIG.model_routing(&model);
RoleChatParams {
model,
provider_order: routing.provider_order,
reasoning_effort: Some(
crate::agent::role::role_info(&role)
.default_reasoning_effort
.to_string(),
),
}
}
#[must_use]
pub(crate) fn chat_request(
role: crate::Role,
tool_specs: Option<Vec<crate::ToolSpec>>,
messages: Vec<ChatMessage>,
) -> ChatRequest {
let RoleChatParams {
model,
provider_order,
reasoning_effort,
} = role_chat_params(role);
ChatRequest {
messages,
tools: tool_specs,
model,
max_tokens: Some(crate::DEFAULT_MAX_TOKENS),
reasoning_effort,
provider_order,
meta: None,
}
}
impl Agent {
#[must_use]
#[expect(clippy::too_many_arguments)] pub fn new(
agent_id: String,
role: crate::Role,
ws: &crate::Workspace,
ticket: Option<crate::pipeline::board::Ticket>,
user_name: String,
channel: String,
is_admin: bool,
parent_key: Option<crate::agent::registry::ParentKey>,
parent_label: Option<String>,
) -> Self {
let chrome_sessions = crate::tools::chrome::ChromeRunSessions::for_run(&agent_id);
let (tools, tool_specs) =
role_tools_and_specs(role, ws, is_admin, std::sync::Arc::clone(&chrome_sessions));
let cancel_token = tokio_util::sync::CancellationToken::new();
let pause_stop = Arc::new(std::sync::atomic::AtomicBool::new(false));
let parent_key = parent_key.or_else(|| {
ticket
.as_ref()
.map(|t| crate::agent::registry::ParentKey::Ticket(t.id.clone()))
});
let parent_label = parent_label.or_else(|| match parent_key {
Some(crate::agent::registry::ParentKey::Ticket(_)) => {
ticket.as_ref().map(|t| t.title.clone())
}
_ => None,
});
if let Some(crate::agent::registry::ParentKey::Research(run_id)) = &parent_key
&& crate::research_cancel::is_cancelled(run_id)
{
cancel_token.cancel();
}
let generation = crate::agent::registry::AGENT_REGISTRY.register(
agent_id.clone(),
role.to_string(),
ticket.as_ref().map(|t| t.id.clone()),
ws,
cancel_token.clone(),
parent_key.clone(),
parent_label.clone(),
pause_stop.clone(),
);
let mut session = Session::default();
session.attach_transcript(
crate::session::TRANSCRIPT_REGISTRY.register(agent_id.clone(), generation),
);
Self {
agent_id,
role,
session,
workspace: Arc::new(ws.clone()),
tools,
tool_specs,
cancel_token,
pause_stop,
paused_frozen: std::sync::atomic::AtomicBool::new(false),
ticket,
generation,
tool_stats: std::sync::Mutex::new(Vec::new()),
user_name,
channel,
is_admin,
parent_key,
parent_label,
incoming_rx: None,
round_ts: None,
first_call_notify: None,
failure: None,
failure_class: None,
sleep_ended: false,
background_sessions: std::sync::Arc::new(
crate::tools::shell::BackgroundSessions::default(),
),
chrome_sessions,
}
}
}
impl Drop for Agent {
fn drop(&mut self) {
if self.generation > 0 {
crate::agent::registry::AGENT_REGISTRY.deregister(&self.agent_id, self.generation);
crate::session::TRANSCRIPT_REGISTRY.deregister(&self.agent_id, self.generation);
}
self.background_sessions.terminate_all();
}
}
fn format_comment_message(user_name: &str, content: &str) -> String {
format!("[Comment from {user_name} on ticket]: {content}")
}
impl Agent {
pub async fn finalize_session(&mut self) -> anyhow::Result<()> {
let stats = {
let mut guard = self.tool_stats.lock().unwrap_poison();
std::mem::take(&mut *guard)
};
if !stats.is_empty()
&& let Some(store) = crate::logs::LOG_STORE.get()
&& let Err(e) = store
.flush_batch(
&self.agent_id,
self.role.as_str(),
&self.workspace.path,
&stats,
)
.await
{
tracing::warn!(
agent_id = %self.agent_id,
role = %self.role.as_str(),
error = %e,
"Failed to flush tool usage stats"
);
}
if self.cancel_token.is_cancelled() || crate::shutdown::shutdown_token().is_cancelled() {
tracing::debug!(
agent_id = %self.agent_id,
role = %self.role,
workspace = %self.workspace.name,
ticket = self.ticket.as_ref().map(|t| t.id.as_str()),
"Session finalize skipped (agent cancelled or shutdown)"
);
return Ok(());
}
match self.session.finalize(&self.agent_id).await? {
crate::session::FinalizeOutcome::Flushed => {}
crate::session::FinalizeOutcome::NoUnpersistedTail => {
if crate::shutdown::is_draining() {
tracing::info!(
agent_id = %self.agent_id,
role = %self.role,
workspace = %self.workspace.name,
ticket = self.ticket.as_ref().map(|t| t.id.as_str()),
"Session finalize no-op: turn cut by graceful drain — \
committed frames are durable; resumes at boot or on the next user message"
);
} else if self.sleep_ended {
tracing::debug!(
agent_id = %self.agent_id,
role = %self.role,
workspace = %self.workspace.name,
ticket = self.ticket.as_ref().map(|t| t.id.as_str()),
"Session finalize no-op: turn ended via sleep — waiting for new input"
);
} else {
tracing::info!(
agent_id = %self.agent_id,
role = %self.role,
workspace = %self.workspace.name,
ticket = self.ticket.as_ref().map(|t| t.id.as_str()),
"finalize called but no new assistant message in history"
);
}
}
}
Ok(())
}
#[must_use]
pub fn is_cancelled(&self) -> bool {
self.cancel_token.is_cancelled()
}
#[must_use]
fn is_cancelled_by_pause(&self) -> bool {
self.pause_stop.load(std::sync::atomic::Ordering::SeqCst)
}
#[must_use]
pub(crate) fn is_paused_frozen(&self) -> bool {
self.paused_frozen.load(std::sync::atomic::Ordering::SeqCst)
}
#[must_use]
pub(crate) fn failure_reason(&self, fallback: &str) -> String {
if crate::shutdown::shutdown_token().is_cancelled() {
"service shutting down".to_string()
} else if self.is_cancelled() {
"agent cancelled".to_string()
} else {
self.failure.clone().unwrap_or_else(|| fallback.to_string())
}
}
fn bail_on_shutdown() -> anyhow::Result<()> {
if crate::shutdown::aborting() {
anyhow::bail!(
"Agent round cut short by shutdown/drain — session is durable; resumes when re-driven"
);
}
Ok(())
}
fn log_drain_cut(tool: &str) {
tracing::info!(tool = %tool, "Drain-cut during resume-completion — leaving calls dangling");
}
fn activity_guard(&self, label: &'static str) -> crate::agent::registry::ActivityGuard {
crate::agent::registry::AGENT_REGISTRY.activity_started(
&self.agent_id,
self.generation,
label,
)
}
pub async fn work(&mut self, msg: &str, resume: bool) -> anyhow::Result<String> {
self.session.init(&self.agent_id).await?;
if matches!(self.role, crate::Role::Assistant | crate::Role::Manager)
&& crate::session::store()
.get_sleep_ended(&self.agent_id)
.await
&& let Err(e) = crate::session::store()
.set_sleep_ended(&self.agent_id, false)
.await
{
tracing::warn!(
agent_id = %self.agent_id,
error = %e,
"Failed to clear sleep-ended flag"
);
}
let settled = if crate::shutdown::aborting() {
false
} else {
self.complete_pending_tool_calls().await?
};
if !msg.is_empty() {
self.session
.append_turn_message(
&self.agent_id,
msg,
&self.workspace,
&self.role,
self.ticket.as_ref(),
&self.channel,
&self.user_name,
self.round_ts.as_deref(),
self.is_admin,
)
.await?;
}
Self::bail_on_shutdown()?;
self.inject_outstanding_comments(resume, msg).await;
if !resume && !settled {
self.maybe_summarize().await;
}
let shutdown = crate::shutdown::shutdown_token();
let response_result = tokio::select! {
() = shutdown.cancelled() => {
Err(anyhow::anyhow!("Shutting down"))
}
result = self.llm_loop() => result,
};
if let Err(e) = self.finalize_session().await {
tracing::error!(error = %e, "Session finalize failed");
}
let response = response_result?;
Ok(response)
}
#[expect(clippy::too_many_lines)] fn complete_pending_tool_calls(&mut self) -> CompletePendingToolCallsFuture<'_> {
Box::pin(async move {
let Some(frame) = self.session.pending_tool_frame() else {
return Ok(false);
};
let calls = frame.calls;
tracing::info!(
agent_id = %self.agent_id,
role = %self.role,
count = calls.len(),
"Completing pending tool calls before the LLM call"
);
let conn = &crate::session::store().conn;
let has_durable = calls
.iter()
.any(|c| crate::tools::SyncDurableCore::from_tool_name(&c.name).is_some());
let mut owned = if has_durable {
crate::jobs::find_owned_launched_jobs(conn, &self.agent_id).await?
} else {
Vec::new()
};
let mut pairs: Vec<(String, String)> = Vec::with_capacity(calls.len());
let mut follow_up: Vec<ChatMessage> = Vec::new();
let mut terminalize: Vec<String> = Vec::new();
let mut fresh_images: Vec<ImagePayload> = Vec::new();
let mut seen_data_uris: Option<std::collections::HashSet<String>> = None;
for call in calls {
let mut resumed_content: Option<String> = None;
let core = crate::tools::SyncDurableCore::from_tool_name(&call.name);
if let (Some(core), Some(pos)) =
(core, owned.iter().rposition(|j| j.kind == call.name))
{
let job = owned.swap_remove(pos);
let outcome = core.resume_sync_core(&self.workspace, &job.id, true).await;
match outcome {
Ok(crate::jobs::SyncResumeOutcome::Terminal(_, _, res)) => {
terminalize.push(job.id);
resumed_content = Some(match res {
Ok(text) => text,
Err(e) => e.to_string(),
});
}
Ok(crate::jobs::SyncResumeOutcome::DrainCut) => {
Self::log_drain_cut(&call.name);
return Ok(false);
}
Ok(crate::jobs::SyncResumeOutcome::Gone) => {
tracing::info!(
job = %job.id,
"Job row gone (explicitly abandoned) — surfacing note instead of re-executing"
);
resumed_content = Some(
"The durable job behind this call was explicitly abandoned — \
no result is available. Re-issue the call if the work is still needed."
.to_string(),
);
}
Err(e) => {
tracing::error!(
job = %job.id,
error = %e,
"Resume-core infra failure — failing the round"
);
return Err(e).context("resume-core infra failure");
}
}
}
let formatted = if let Some(content) = resumed_content {
find_tool(&self.tools, &call.name)
.map(|t| t.format_output(&content))
.unwrap_or(content)
} else if call.name == "sleep"
&& let Some(reason) =
sleep_rejection_reason(frame.has_text, frame.call_count > 1)
{
tracing::warn!(
agent_id = %self.agent_id,
tool = %call.name,
reason,
"Sleep call rejected during resume-completion"
);
self.push_tool_call_record(
call.name.clone(),
&call.arguments,
0,
false,
Some(reason.to_string()),
);
Self::failure_outcome(&call.name, &call.arguments, reason)
.0
.output
} else {
let outcome = self
.in_tool_context(async {
self.execute_tool(&call.name, call.arguments.clone()).await
})
.await;
if outcome.suspended {
Self::log_drain_cut(&call.name);
return Ok(false);
}
let (formatted, fresh_payloads) = tool_result_content(
&self.tools,
&call.name,
outcome,
&mut seen_data_uris,
self.session.history(),
)
.await;
fresh_images.extend(fresh_payloads);
formatted
};
pairs.push((call.id, formatted));
}
if let Some(message) = injected_images_message(&self.agent_id, self.role, fresh_images)
{
follow_up.push(message);
}
if !pairs.is_empty() || !follow_up.is_empty() {
let tx = conn
.begin_tx()
.await
.context("begin settle+terminalize tx")?;
let settled = crate::session::store()
.settle_tool_results_tx(&tx, &self.agent_id, &pairs, &follow_up)
.await
.context("settle resumed results")?;
for id in terminalize {
tx.execute(
"DELETE FROM jobs WHERE id = ?1",
crate::db::params![id.as_str()],
)
.await
.with_context(|| format!("terminalize job {id}"))?;
}
tx.commit().await.context("commit settle+terminalize tx")?;
self.session
.mirror_settled(settled)
.context("mirror settled results")?;
}
Ok(!pairs.is_empty() || !follow_up.is_empty())
})
}
async fn llm_loop(&mut self) -> anyhow::Result<String> {
let span = tracing::info_span!("agent", agent_id = %self.agent_id, role = %self.role, workspace = %self.workspace.path);
async {
let mut tool_rounds = 0usize;
let mut accumulated_media_paths: Vec<(&'static str, String)> = Vec::new();
loop {
if self.cancel_token.is_cancelled() {
anyhow::bail!("Agent cancelled");
}
if self.is_cancelled_by_pause() {
self.paused_frozen
.store(true, std::sync::atomic::Ordering::SeqCst);
anyhow::bail!("Agent frozen by workspace pause — resumes at unpause");
}
Self::bail_on_shutdown()?;
if tool_rounds >= MAX_TOOL_ROUNDS {
anyhow::bail!(
"Agent exceeded maximum of {MAX_TOOL_ROUNDS} tool rounds \
— model may be stuck in a tool-calling loop"
);
}
self.drain_incoming_messages().await;
let llm_result = self.llm_call().await;
if tool_rounds == 0
&& let Some(notify) = &self.first_call_notify
{
notify.notify_one();
}
let image_rejection = llm_result.as_ref().err().and_then(|e| {
e.chain()
.find_map(|cause| cause.downcast_ref::<crate::retry::RetryExhausted>())
});
if let Some(exhausted) = image_rejection
&& self.strip_rejected_input_image(exhausted).await
{
tracing::info!(
agent_id = %self.agent_id,
role = %self.role,
tool_rounds,
"Stripped provider-rejected input image from the most recent user \
message — continuing the normal loop"
);
continue;
}
let PreparedAssistantTurn {
mut display_text,
tool_calls,
history_content,
raw_text,
} = prepare_assistant_turn(
llm_result.with_context(|| {
llm_step_failure_context(&self.session, tool_rounds)
})?,
);
if tool_calls.is_empty() {
self.session.push_assistant(history_content);
for (marker_prefix, path) in &accumulated_media_paths {
let marker = format!("{marker_prefix}{path}]");
if !display_text.contains(&marker) {
let _ = write!(display_text, "\n{marker}");
}
}
return Ok(display_text);
}
self.session
.persist_messages(
&self.agent_id,
&[ChatMessage::assistant(history_content.clone())],
)
.await?; let has_sleep = tool_calls.iter().any(|c| c.name == "sleep");
let sleep_rejection = has_sleep
.then(|| sleep_rejection_reason(!raw_text.trim().is_empty(), tool_calls.len() > 1))
.flatten();
let all_outcomes = self.execute_tool_round(&tool_calls, sleep_rejection).await;
accumulated_media_paths.extend(extract_media_from_outcomes(
&self.tools,
&tool_calls,
&all_outcomes,
));
let sleep_ended = all_outcomes.iter().any(|o| o.success && o.ends_turn);
self.commit_tool_results(&tool_calls, all_outcomes).await?;
if sleep_ended {
self.sleep_ended = true;
if let Err(e) = crate::session::store().set_sleep_ended(&self.agent_id, true).await {
tracing::warn!(
agent_id = %self.agent_id,
error = %e,
"Failed to persist sleep-ended flag — a stale tool-tail session may be re-recovered (accepted)"
);
}
return Ok(String::new());
}
tool_rounds += 1;
}
}
.instrument(span)
.await
}
async fn execute_tool_group(&self, tool_calls: &[ToolCall]) -> Vec<ToolExecutionOutcome> {
let side_flags: Vec<bool> = tool_calls
.iter()
.map(|call| find_tool(&self.tools, &call.name).is_none_or(super::Tool::side_effects))
.collect();
let mut outcomes: Vec<ToolExecutionOutcome> = Vec::with_capacity(tool_calls.len());
let mut i = 0usize;
self.in_tool_context(async {
while i < tool_calls.len() {
if side_flags[i] {
let outcome = self
.execute_tool(&tool_calls[i].name, tool_calls[i].arguments.clone())
.await;
outcomes.push(outcome);
i += 1;
} else {
let group_start = i;
while i < tool_calls.len() && !side_flags[i] {
i += 1;
}
let group_calls = &tool_calls[group_start..i];
let group_outcomes: Vec<_> = futures_util::future::join_all(
group_calls
.iter()
.map(|call| self.execute_tool(&call.name, call.arguments.clone())),
)
.await;
outcomes.extend(group_outcomes);
}
}
})
.await;
outcomes
}
async fn execute_tool_round(
&self,
tool_calls: &[ToolCall],
sleep_rejection: Option<&'static str>,
) -> Vec<ToolExecutionOutcome> {
let Some(reason) = sleep_rejection else {
return self.execute_tool_group(tool_calls).await;
};
let rejected = tool_calls.iter().filter(|c| c.name == "sleep").count();
tracing::info!(
agent_id = %self.agent_id,
role = %self.role,
rejected,
reason,
"Sleep call rejected — must run alone in a text-free round"
);
let sibling_calls: Vec<ToolCall> = tool_calls
.iter()
.filter(|c| c.name != "sleep")
.cloned()
.collect();
let mut sibling_outcomes = self.execute_tool_group(&sibling_calls).await.into_iter();
let mut outcomes = Vec::with_capacity(tool_calls.len());
for call in tool_calls {
if call.name == "sleep" {
self.push_tool_call_record(
call.name.clone(),
&call.arguments,
0,
false,
Some(reason.to_string()),
);
outcomes.push(Self::failure_outcome(&call.name, &call.arguments, reason).0);
} else {
outcomes.push(
sibling_outcomes
.next()
.expect("non-sleep outcomes align one-to-one"),
);
}
}
outcomes
}
async fn in_tool_context<T>(&self, fut: impl std::future::Future<Output = T>) -> T {
let user_name = self.user_name.clone();
let channel = self.channel.clone();
let parent_key = self.parent_key.clone();
let parent_label = self.parent_label.clone();
let background_sessions = Some(self.background_sessions.clone());
let agent_id = Some(self.agent_id.clone());
let agent_tracking = Some(crate::agent::registry::AgentTracking {
agent_id: self.agent_id.clone(),
generation: self.generation,
role: self.role.as_str().to_string(),
workspace: self.workspace.name.clone(),
});
CURRENT_TOOL_USER_NAME
.scope(user_name, async {
CURRENT_TOOL_CHANNEL
.scope(channel, async {
CURRENT_TOOL_PARENT_KEY
.scope(parent_key, async {
CURRENT_TOOL_PARENT_LABEL
.scope(parent_label, async {
CURRENT_TOOL_BACKGROUND_SESSIONS
.scope(background_sessions, async {
CURRENT_TOOL_AGENT_ID
.scope(agent_id, async {
CURRENT_TOOL_AGENT_TRACKING
.scope(agent_tracking, fut)
.await
})
.await
})
.await
})
.await
})
.await
})
.await
})
.await
}
#[must_use]
fn failure_outcome(
call_name: &str,
call_arguments: &serde_json::Value,
reason: &str,
) -> (ToolExecutionOutcome, String) {
let reason = scrub_credentials(reason);
(
ToolExecutionOutcome {
output: format_tool_failure_feedback(call_name, call_arguments, &reason),
success: false,
image_payloads: Vec::new(),
text_is_content: false,
suspended: false,
ends_turn: false,
},
reason,
)
}
fn push_tool_call_record(
&self,
tool_name: String,
arguments: &serde_json::Value,
duration_ms: i64,
success: bool,
error_message: Option<String>,
) {
let args_str = serde_json::to_string(arguments).expect("Value is always serializable");
let args_scrubbed = scrub_credentials(&args_str);
let arguments =
crate::util::truncate_bytes(&args_scrubbed, MAX_STATS_ARG_LENGTH).to_string();
let mut guard = self.tool_stats.lock().unwrap_poison();
guard.push(crate::ToolCallRecord {
tool_name,
arguments,
duration_ms,
success,
error_message,
});
}
#[expect(clippy::too_many_lines)]
async fn execute_tool(
&self,
call_name: &str,
call_arguments: serde_json::Value,
) -> ToolExecutionOutcome {
if self.cancel_token.is_cancelled() {
let reason = "Agent cancelled — tool execution skipped";
tracing::debug!(
tool = %call_name,
"Agent cancelled — skipping tool execution"
);
return Self::failure_outcome(call_name, &call_arguments, reason).0;
}
let start = Instant::now();
let (tool_name, tool_arguments) = normalize_tool_call(call_name, call_arguments);
if tool_name != call_name {
tracing::debug!(
original = %call_name,
normalized = %tool_name,
"Repaired tool call name"
);
}
let (outcome, error_reason) = match find_tool(&self.tools, &tool_name) {
None => {
let reason = format!("Unknown tool: {tool_name}");
let duration = start.elapsed();
tracing::debug!(
tool = %tool_name,
duration_ms = duration.as_millis(),
success = false,
"Unknown tool call"
);
Self::failure_outcome(&tool_name, &tool_arguments, &reason)
}
Some(tool) => {
let exec_result = tool
.execute_with_payloads(&self.workspace, tool_arguments.clone())
.await;
let duration = start.elapsed();
match exec_result {
Ok(result) => {
let output_text = if result.text.is_empty() {
String::from("(no output)")
} else {
result.text
};
tracing::debug!(
tool = %tool_name,
duration_ms = duration.as_millis(),
"Tool execution completed"
);
(
ToolExecutionOutcome {
output: scrub_tool_output(tool, &tool_arguments, &output_text),
success: true,
image_payloads: result.image_payloads,
text_is_content: result.text_is_content,
suspended: false,
ends_turn: tool.ends_turn_on_success(),
},
String::new(),
)
}
Err(e) => {
if e.downcast_ref::<crate::tools::CallSuspended>().is_some() {
tracing::debug!(
tool = %tool_name,
duration_ms = duration.as_millis(),
success = false,
suspended = true,
"Tool execution drain-cut — durable work left launched"
);
(
ToolExecutionOutcome {
output: String::new(),
success: false,
image_payloads: Vec::new(),
text_is_content: false,
suspended: true,
ends_turn: false,
},
String::new(),
)
} else {
let (outcome, error_reason) = Self::failure_outcome(
&tool_name,
&tool_arguments,
&format!("Error executing {tool_name}: {e}"),
);
tracing::debug!(
tool = %tool_name,
duration_ms = duration.as_millis(),
success = false,
"Tool execution error: {error_reason}"
);
(outcome, error_reason)
}
}
}
}
};
let elapsed_ms = start.elapsed().as_millis();
let duration_ms = i64::try_from(elapsed_ms).unwrap_or(0);
self.push_tool_call_record(
tool_name,
&tool_arguments,
duration_ms,
outcome.success,
(!error_reason.is_empty()).then_some(error_reason),
);
outcome
}
async fn llm_call(&mut self) -> anyhow::Result<ChatResponse> {
let messages = self.session.history().to_vec();
let request = self.build_chat_request(messages.clone(), "agent");
let policy = crate::retry::RetryPolicy::current();
let response = crate::retry::agent_chat(request, &policy)
.await
.with_context(|| format!("LLM call {RETRY_EXHAUSTION_MARKER}"))?;
let response = self
.recover_if_reasoning_only_stop(messages, response, Self::AGENT_REASONING_RECOVERY)
.await?;
self.record_session_usage(&response).await;
Ok(response)
}
async fn strip_rejected_input_image(
&mut self,
exhausted: &crate::retry::RetryExhausted,
) -> bool {
let Some(idx) = crate::session::image_strip::detect_input_image_rejection(
exhausted,
self.session.history(),
) else {
return false;
};
let reason = crate::session::image_strip::extract_provider_reason(exhausted);
let content = {
let original = &self.session.history()[idx].content;
let stripped =
crate::session::image_strip::strip_image_markers(original, reason.as_deref());
if stripped == *original {
tracing::warn!(
agent_id = %self.agent_id,
role = %self.role,
"Input-image rejection detected but stripping produced no change — \
treating as a non-strip (normal failure path)"
);
return false;
}
stripped
};
match self
.session
.rewrite_last_user_message(&self.agent_id, content)
.await
{
Ok(crate::session::RewriteOutcome::Rewritten) => true,
Ok(crate::session::RewriteOutcome::UnpersistedTailNoop) => {
tracing::info!(
agent_id = %self.agent_id,
role = %self.role,
"Input-image rejection detected but the most recent user message is in the \
unpersisted tail — conservative no-op, normal failure path applies"
);
false
}
Err(e) => {
tracing::error!(
agent_id = %self.agent_id,
role = %self.role,
error = %e,
"Failed to persist stripped user message — keeping the original error path"
);
false
}
}
}
async fn recover_reasoning_only_stop(
&self,
base: Vec<ChatMessage>,
first: ChatResponse,
purpose: &'static str,
) -> Result<ChatResponse, crate::retry::RetryExhausted> {
let policy = crate::retry::RetryPolicy::continuation();
let operation_started = Instant::now();
let nudge = crate::prompt::load_prompt("resume_unfinished_turn.md")
.trim()
.to_string();
let mut failures: Vec<crate::retry::RetryFailureRecord> = Vec::new();
let mut last_request: Option<ChatRequest> = None;
let mut operation: Option<crate::stats::LlmOperationCtx> = None;
let mut tail: Vec<ChatMessage> = vec![
ChatMessage::assistant(
assistant_replay_payload(None, &[], first.reasoning.as_ref()).to_string(),
),
ChatMessage::user(nudge.clone()),
];
let mut tail_grew = true;
for attempt in 1..=policy.max_attempts {
let attempt_started = Instant::now();
if crate::shutdown::aborting() {
failures.push(crate::retry::RetryFailureRecord::new_simple(
crate::retry::FailureClass::Shutdown,
&anyhow::anyhow!("global shutdown or drain during continuation recovery"),
None,
));
break;
}
if self.cancel_token.is_cancelled() {
break;
}
let request = if tail_grew {
tail_grew = false;
let mut messages = base.clone();
messages.extend(tail.iter().cloned());
let built = self.build_chat_request(messages, purpose);
if operation.is_none() {
operation = crate::stats::LlmOperationCtx::from_request(&built);
}
last_request = Some(built.clone());
built
} else {
last_request
.clone()
.expect("the first iteration always builds a request")
};
match crate::providers::chat_scoped(request.clone()).await {
Ok(resp) if is_reasoning_only_stop(&resp) => {
let rec = crate::retry::RetryFailureRecord::with_metadata(
crate::retry::FailureClass::NoResponse,
&anyhow::anyhow!(
"model returned only reasoning with no answer \
(continuation attempt {attempt})"
),
resp.finish_reason.clone(),
None,
);
record_continuation_failure(operation.as_ref(), attempt, &rec, attempt_started)
.await;
failures.push(rec);
tail.push(ChatMessage::assistant(
assistant_replay_payload(None, &[], resp.reasoning.as_ref()).to_string(),
));
tail.push(ChatMessage::user(nudge.clone()));
tail_grew = true;
}
Ok(resp) => {
crate::stats::record_llm_success(&request, operation_started, attempt, &resp)
.await;
return Ok(resp);
}
Err(err) => {
record_continuation_failure(
operation.as_ref(),
attempt,
&err.record,
attempt_started,
)
.await;
failures.push(err.record);
if !err.class.is_retryable() {
break;
}
}
}
}
let final_class = failures
.last()
.map_or(crate::retry::FailureClass::NoResponse, |r| r.class);
let exhausted = crate::retry::RetryExhausted::with_last_raw(failures, final_class, None);
match last_request {
Some(request) => {
crate::retry::fail_exhausted(&request, operation_started, exhausted).await
}
None => Err(exhausted),
}
}
const AGENT_REASONING_RECOVERY: ReasoningOnlyStopRecovery = ReasoningOnlyStopRecovery {
purpose: "agent-continuation",
exhausted_ctx: "model returned only reasoning without an answer after continuation attempts",
};
const SUMMARIZE_REASONING_RECOVERY: ReasoningOnlyStopRecovery = ReasoningOnlyStopRecovery {
purpose: "summarize-continuation",
exhausted_ctx: "summarization continuation exhausted — failing open with full history",
};
async fn recover_if_reasoning_only_stop(
&self,
messages: Vec<ChatMessage>,
response: ChatResponse,
recovery: ReasoningOnlyStopRecovery,
) -> anyhow::Result<ChatResponse> {
if !is_reasoning_only_stop(&response) {
return Ok(response);
}
self.recover_reasoning_only_stop(messages, response, recovery.purpose)
.await
.map_err(|e| anyhow::Error::new(e).context(recovery.exhausted_ctx))
}
async fn record_session_usage(&mut self, response: &ChatResponse) {
let Some(usage) = &response.usage else {
return;
};
let (Some(input), Some(output)) = (usage.input_tokens, usage.output_tokens) else {
return;
};
let token_length = input.saturating_add(output);
self.session.set_token_length(Some(token_length));
if let Err(e) = crate::session::store()
.set_token_length(&self.agent_id, Some(token_length))
.await
{
tracing::warn!(
agent_id = %self.agent_id,
error = %e,
"Failed to persist session token length — in-memory value may drift from the store until the next successful call"
);
}
}
async fn commit_tool_results(
&mut self,
tool_calls: &[ToolCall],
outcomes: Vec<ToolExecutionOutcome>,
) -> anyhow::Result<()> {
let tools = &self.tools;
let mut seen_data_uris: Option<std::collections::HashSet<String>> = None;
let mut fresh_image_payloads: Vec<ImagePayload> = Vec::new();
let mut db_messages = Vec::with_capacity(outcomes.len());
for (call, outcome) in tool_calls.iter().zip(outcomes) {
if outcome.suspended {
continue;
}
let (output, fresh_payloads) = tool_result_content(
tools,
&call.name,
outcome,
&mut seen_data_uris,
self.session.history(),
)
.await;
fresh_image_payloads.extend(fresh_payloads);
db_messages.push(ChatMessage::tool_result(&call.id, &output));
}
if let Some(message) =
injected_images_message(&self.agent_id, self.role, fresh_image_payloads)
{
db_messages.push(message);
}
if db_messages.is_empty() {
return Ok(());
}
self.session
.persist_messages(&self.agent_id, &db_messages)
.await
.map_err(|e| anyhow::anyhow!("Failed to persist tool results: {e}"))?;
Ok(())
}
async fn drain_incoming_messages(&mut self) {
let Some(rx) = &mut self.incoming_rx else {
return;
};
let mut messages = Vec::new();
loop {
match rx.try_recv() {
Ok(job) => {
let content = match job.kind {
crate::agent::message_router::MessageKind::TicketComment => {
format_comment_message(&job.user_name, &job.content)
}
_ => job.content,
};
messages.push(crate::session::user_msg_with_ts(&content, None));
}
Err(tokio::sync::mpsc::error::TryRecvError::Empty) => break,
Err(tokio::sync::mpsc::error::TryRecvError::Disconnected) => {
self.incoming_rx = None;
break;
}
}
}
if messages.is_empty() {
return;
}
self.persist_messages_fail_open(&messages).await;
}
async fn persist_messages_fail_open(&mut self, messages: &[ChatMessage]) {
if let Err(e) = self
.session
.persist_messages(&self.agent_id, messages)
.await
{
tracing::warn!(
agent_id = %self.agent_id,
error = %e,
"Failed to persist messages to session DB — continuing without persistence",
);
self.session.push_messages_unpersisted(messages);
}
}
async fn inject_outstanding_comments(&mut self, resume: bool, msg: &str) {
if !resume || !msg.is_empty() {
return;
}
let Some(ticket_id) = self.ticket.as_ref().map(|t| t.id.clone()) else {
return;
};
let comments = match crate::pipeline::board::store()
.get_comments(&ticket_id)
.await
{
Ok(comments) => comments,
Err(e) => {
tracing::warn!(
agent_id = %self.agent_id,
ticket = %ticket_id,
error = %e,
"Failed to read board comments for resume catch-up — skipping this round",
);
return;
}
};
if comments.is_empty() {
return;
}
self.drain_incoming_messages().await;
let new_messages: Vec<_> = {
let history = self.session.history();
comments
.iter()
.filter(|c| {
!c.content.trim().is_empty()
&& !history.iter().any(|m| m.content.contains(&c.content))
})
.map(|c| {
crate::session::user_msg_with_ts(
&format_comment_message(&c.role, &c.content),
None,
)
})
.collect()
};
if new_messages.is_empty() {
return;
}
self.persist_messages_fail_open(&new_messages).await;
tracing::info!(
agent_id = %self.agent_id,
role = %self.role,
ticket = %ticket_id,
count = new_messages.len(),
"Injected outstanding board comments into resumed stage-agent session",
);
}
fn build_chat_request(&self, messages: Vec<ChatMessage>, purpose: &'static str) -> ChatRequest {
ChatRequest {
meta: Some(crate::ChatRequestMeta {
purpose,
agent_id: self.agent_id.clone(),
role: self.role.as_str().to_string(),
workspace: self.workspace.name.clone(),
ticket_id: self.ticket.as_ref().map(|t| t.id.clone()),
}),
..chat_request(self.role, Some(self.tool_specs.clone()), messages)
}
}
pub(crate) async fn extract_verdict<T: serde::de::DeserializeOwned>(
&self,
extraction_prompt: &str,
validate: Option<&crate::ExtractionValidator<T>>,
policy_override: Option<&crate::retry::RetryPolicy>,
) -> Result<T, crate::retry::RetryExhausted> {
let _activity = self.activity_guard("extracting");
let params = self.build_chat_request(vec![], "extraction");
crate::agent::extraction::retry_extract_structured_scoped(
self.session.history(),
extraction_prompt,
¶ms,
validate,
policy_override,
)
.await
}
pub(crate) async fn summarize(&self) -> anyhow::Result<String> {
let _activity = self.activity_guard("summarizing");
let mut history = self.session.history().to_vec();
history.push(crate::ChatMessage::user(self.role.summary_prompt()));
let request = self.build_chat_request(history.clone(), "summarize");
let policy = crate::retry::RetryPolicy::current();
for _ in 1..=SUMMARIZE_ATTEMPTS {
let chat_resp = crate::retry::agent_chat(request.clone(), &policy)
.await
.with_context(|| format!("summarization LLM call {RETRY_EXHAUSTION_MARKER}"))?;
let chat_resp = self
.recover_if_reasoning_only_stop(
history.clone(),
chat_resp,
Self::SUMMARIZE_REASONING_RECOVERY,
)
.await?;
if let Some(ref u) = chat_resp.usage {
tracing::debug!(
input_tokens = u.input_tokens,
cached_input_tokens = u.cached_input_tokens,
output_tokens = u.output_tokens,
"Summarization token usage",
);
}
if let Some(summary_text) = chat_resp.text.filter(|t| !t.trim().is_empty()) {
return Ok(crate::util::truncate(&summary_text, 32_000));
}
}
anyhow::bail!("summarization produced empty response")
}
async fn maybe_summarize(&mut self) {
let token_length = self.session.token_length();
let image_count = session_image_marker_count(self.session.history());
tracing::debug!(
agent_id = %self.agent_id,
role = %self.role,
token_length,
image_count,
"Session token length / image count",
);
let Some(trigger) = summarization_trigger(token_length, image_count) else {
return;
};
tracing::info!(
agent_id = %self.agent_id,
role = %self.role,
trigger,
token_length,
image_count,
"Session exceeded summarization trigger ({trigger})",
);
match self.summarize().await {
Ok(summary) => {
self.session
.apply_summary(
&self.agent_id,
&summary,
&self.workspace,
&self.role,
self.ticket.as_ref(),
self.is_admin,
&self.user_name,
)
.await;
}
Err(e) => {
tracing::warn!(
error_chain = %crate::util::failure_detail(&format!("{e:#}"), "summarization failure"),
"Summarization failed — continuing with full history"
);
}
}
}
}
#[must_use]
fn summarization_trigger(token_length: Option<u64>, image_count: usize) -> Option<&'static str> {
if token_length.is_some_and(|t| t > crate::session::SUMMARIZATION_THRESHOLD) {
Some("token threshold")
} else if image_count >= crate::session::SUMMARIZATION_IMAGE_COUNT {
Some("image count")
} else {
None
}
}
struct ReasoningOnlyStopRecovery {
purpose: &'static str,
exhausted_ctx: &'static str,
}
struct PreparedAssistantTurn {
display_text: String,
tool_calls: Vec<ToolCall>,
history_content: String,
raw_text: String,
}
#[must_use]
fn is_reasoning_only_stop(response: &ChatResponse) -> bool {
response.tool_calls.is_empty() && response.text.as_deref().is_none_or(|t| t.trim().is_empty())
}
async fn record_continuation_failure(
operation: Option<&crate::stats::LlmOperationCtx>,
attempt: u32,
rec: &crate::retry::RetryFailureRecord,
attempt_started: Instant,
) {
if let Some(op) = operation {
crate::stats::record_llm_attempt_failure(op, attempt, rec, Some(attempt_started)).await;
}
}
fn llm_step_failure_context(session: &Session, tool_rounds: usize) -> String {
let tokens = session
.token_length()
.map_or_else(String::new, |t| format!(", {t} tokens"));
format!(
"LLM step failed at tool round {tool_rounds} (run-local counter; \
session history: {} messages{tokens})",
session.history().len(),
)
}
fn prepare_assistant_turn(response: ChatResponse) -> PreparedAssistantTurn {
let response_text = response.text_or_empty().to_string();
let raw_text = response_text.clone();
let tool_calls = response.tool_calls;
let reasoning = response.reasoning.as_ref();
let json_payload =
assistant_replay_payload(Some(&response_text), &tool_calls, reasoning).to_string();
let (display_text, history_content) = match (tool_calls.is_empty(), reasoning.is_some()) {
(true, false) => (response_text.clone(), response_text),
(true, true) => (response_text, json_payload),
(false, _) => (String::new(), json_payload),
};
PreparedAssistantTurn {
display_text,
tool_calls,
history_content,
raw_text,
}
}
const SLEEP_WITH_TEXT_REJECTION: &str = "sleep must be called alone: your text was NOT delivered — text written in any round that contains a tool call never reaches the user. If that text was meant for the user, resend it as a plain-text round with no tool calls and no sleep call: that round ends your turn and delivers it. If it was not meant for the user, call sleep alone, with no text and no other tool call.";
const SLEEP_BUNDLED_REJECTION: &str = "sleep must be called alone: it was bundled with other tool calls. Call sleep in its own round with no other tool calls — and no text, since text written in a tool round is never delivered to the user.";
#[must_use]
fn sleep_rejection_reason(has_text: bool, bundled: bool) -> Option<&'static str> {
if has_text {
Some(SLEEP_WITH_TEXT_REJECTION)
} else if bundled {
Some(SLEEP_BUNDLED_REJECTION)
} else {
None
}
}
#[derive(Clone)]
pub(crate) struct RoundOpts {
pub(crate) round_ts: String,
pub(crate) first_call_notify: Option<std::sync::Arc<tokio::sync::Notify>>,
}
const DEFAULT_STAGGER_WAIT_SECS: u64 = 8;
fn leader_stagger_wait() -> std::time::Duration {
crate::util::env_duration_secs("MAHBOT_STAGGER_WAIT_SECS", DEFAULT_STAGGER_WAIT_SECS)
}
pub(crate) async fn spawn_staggered_round<T, Fut, F>(
members: Vec<F>,
resume: bool,
) -> Vec<tokio::task::JoinHandle<T>>
where
T: Send + 'static,
Fut: std::future::Future<Output = T> + Send + 'static,
F: FnOnce(RoundOpts) -> Fut + Send + 'static,
{
let mut members = members.into_iter();
let Some(leader) = members.next() else {
return Vec::new();
};
let followers: Vec<F> = members.collect();
let round_ts = crate::session::render_timestamp();
if resume || followers.is_empty() {
let opts = RoundOpts {
round_ts,
first_call_notify: None,
};
let mut handles = Vec::with_capacity(followers.len() + 1);
handles.push(tokio::spawn(leader(opts.clone())));
handles.extend(followers.into_iter().map(|m| tokio::spawn(m(opts.clone()))));
return handles;
}
let notify = std::sync::Arc::new(tokio::sync::Notify::new());
let mut handles = vec![tokio::spawn(leader(RoundOpts {
round_ts: round_ts.clone(),
first_call_notify: Some(notify.clone()),
}))];
let follower_opts = RoundOpts {
round_ts,
first_call_notify: None,
};
let shutdown = crate::shutdown::shutdown_token();
tokio::select! {
() = notify.notified() => {}
() = tokio::time::sleep(leader_stagger_wait()) => {}
() = shutdown.cancelled() => {}
() = crate::shutdown::drain_wait() => {}
}
handles.extend(
followers
.into_iter()
.map(|m| tokio::spawn(m(follower_opts.clone()))),
);
handles
}
struct RunEndCleanup {
agent_id: String,
chrome: std::sync::Arc<crate::tools::chrome::ChromeRunSessions>,
held: bool,
}
impl RunEndCleanup {
fn new(
agent_id: String,
chrome: std::sync::Arc<crate::tools::chrome::ChromeRunSessions>,
) -> Self {
Self {
agent_id,
chrome,
held: false,
}
}
fn ended(&mut self, classification: &str) {
self.held = crate::tools::chrome_release::holds_run_end(classification);
}
}
impl Drop for RunEndCleanup {
fn drop(&mut self) {
crate::tools::shell::cleanup_agent_spills(&self.agent_id);
#[cfg(not(target_os = "macos"))]
crate::tools::computer::cleanup_agent_state(&self.agent_id);
crate::tools::chrome_release::queue_run_session_release(&self.chrome, self.held);
}
}
#[expect(clippy::too_many_arguments)]
pub(crate) async fn run_agent(
agent_id: String,
role: crate::Role,
ws: &crate::Workspace,
ticket: Option<&crate::pipeline::board::Ticket>,
message: &str,
user_name: String,
channel: String,
is_admin: bool,
incoming_rx: Option<
tokio::sync::mpsc::UnboundedReceiver<crate::agent::message_router::AgentJob>,
>,
resume: bool,
round: Option<RoundOpts>,
parent_key: Option<crate::agent::registry::ParentKey>,
parent_label: Option<String>,
) -> (Agent, Option<String>) {
struct UnregisterOnDrop(String);
impl Drop for UnregisterOnDrop {
fn drop(&mut self) {
crate::agent::message_router::unregister_agent(&self.0);
}
}
let _router_guard = incoming_rx
.is_some()
.then(|| UnregisterOnDrop(agent_id.clone()));
let mut agent = Agent::new(
agent_id,
role,
ws,
ticket.cloned(),
user_name,
channel,
is_admin,
parent_key,
parent_label,
);
agent.incoming_rx = incoming_rx;
if let Some(round) = round {
agent.round_ts = Some(round.round_ts);
agent.first_call_notify = round.first_call_notify;
}
let mut run_end = RunEndCleanup::new(
agent.agent_id.clone(),
std::sync::Arc::clone(&agent.chrome_sessions),
);
let result = agent.work(message, resume).await;
match result {
Ok(response) => (agent, Some(response)),
Err(e) => {
agent.failure = Some(format!("{e:#}"));
agent.failure_class = failure_class_from_error(&e);
let classification = failure_classification(&agent, Some(&e));
run_end.ended(classification);
let error_chain = crate::util::failure_detail(&format!("{e:#}"), "agent failure log");
let non_failure = matches!(
classification,
"drain" | "shutdown" | "pause" | "internal_cancel"
);
if non_failure {
tracing::debug!(
agent_id = %agent.agent_id,
workspace = %ws.name,
role = %role,
ticket = ticket.map(|t| t.id.as_str()),
classification,
error_chain,
"Agent stopped by a non-failure cancellation"
);
} else {
tracing::error!(
agent_id = %agent.agent_id,
workspace = %ws.name,
role = %role,
ticket = ticket.map(|t| t.id.as_str()),
classification,
error_chain,
"Agent failed"
);
}
(agent, None)
}
}
}
#[expect(clippy::too_many_arguments)]
pub(crate) async fn run_default_agent(
agent_id: &str,
role: crate::Role,
ws: &crate::Workspace,
message: &str,
resume: bool,
round: Option<RoundOpts>,
parent_key: Option<crate::agent::registry::ParentKey>,
parent_label: Option<String>,
) -> (Agent, Option<String>) {
run_agent(
agent_id.to_string(),
role,
ws,
None,
message,
String::new(),
String::new(),
false,
None,
resume,
round,
parent_key,
parent_label,
)
.await
}
pub(crate) fn failure_class_from_error(
error: &anyhow::Error,
) -> Option<crate::retry::FailureClass> {
error
.chain()
.find_map(|cause| cause.downcast_ref::<crate::retry::RetryExhausted>())
.map(|exhausted| exhausted.final_class)
}
fn failure_classification(agent: &Agent, error: Option<&anyhow::Error>) -> &'static str {
if crate::shutdown::is_draining() {
"drain"
} else if crate::shutdown::shutdown_token().is_cancelled() {
"shutdown"
} else if agent.is_paused_frozen() {
"pause"
} else if agent.is_cancelled() {
"internal_cancel"
} else if let Some(class) = error.and_then(failure_class_from_error) {
class.label()
} else {
"runtime"
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::Tool;
use crate::tools::chrome_release;
use crate::util::test::{FakeProvider, install_fake_provider, install_test_retry_policy};
use async_trait::async_trait;
use tokio_util::sync::CancellationToken;
struct TestTool {
output: String,
scrub: bool,
}
#[async_trait]
impl Tool for TestTool {
fn name(&self) -> &'static str {
if self.scrub {
"always_scrub"
} else {
"never_scrub"
}
}
fn description(&self) -> String {
"test".into()
}
fn parameters_schema(&self) -> serde_json::Value {
serde_json::json!({})
}
async fn execute(
&self,
_ws: &crate::Workspace,
_args: serde_json::Value,
) -> anyhow::Result<String> {
Ok(self.output.clone())
}
fn should_scrub_output(&self, _args: &serde_json::Value) -> bool {
self.scrub
}
}
const SCRUBBABLE_LINE: &str = "API_KEY=sk-1234567890abcdef";
fn make_agent(tools: Vec<Box<dyn Tool>>) -> Agent {
make_agent_with_role(tools, crate::Role::Engineer)
}
fn make_agent_with_role(tools: Vec<Box<dyn Tool>>, role: crate::Role) -> Agent {
let tool_specs = tools.iter().map(|t| t.spec()).collect();
Agent {
agent_id: "test-agent".into(),
role,
session: Session::default(),
workspace: std::sync::Arc::new(crate::Workspace::default()),
tools,
tool_specs,
cancel_token: CancellationToken::new(),
pause_stop: std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)),
paused_frozen: std::sync::atomic::AtomicBool::new(false),
ticket: None,
generation: 0,
tool_stats: std::sync::Mutex::new(Vec::new()),
user_name: String::new(),
channel: String::new(),
is_admin: false,
parent_key: None,
parent_label: None,
incoming_rx: None,
round_ts: None,
first_call_notify: None,
failure: None,
failure_class: None,
sleep_ended: false,
background_sessions: std::sync::Arc::new(
crate::tools::shell::BackgroundSessions::default(),
),
chrome_sessions: std::sync::Arc::new(crate::tools::chrome::ChromeRunSessions::default()),
}
}
fn make_agent_on(tools: Vec<Box<dyn Tool>>, agent_id: &str, ws: crate::Workspace) -> Agent {
let mut agent = make_agent_with_role(tools, crate::Role::Engineer);
agent.agent_id = agent_id.to_string();
agent.workspace = std::sync::Arc::new(ws);
agent
}
#[tokio::test]
async fn tool_with_scrub_disabled_preserves_output() {
assert_scrubbed(false).await;
}
#[test]
#[serial_test::serial(drain)] fn failure_classification_recovers_retry_exhaustion() {
let exhausted = crate::retry::RetryExhausted::with_last_raw(
vec![],
crate::retry::FailureClass::NonRetryable,
None,
);
let err = anyhow::Error::new(exhausted).context("LLM call exhausted retry budget");
assert!(format!("{err:#}").contains("exhausted retry budget"));
let agent = make_agent(vec![]);
assert_eq!(
failure_classification(&agent, Some(&err)),
"non_retryable",
"RetryExhausted final_class must be recovered from the chain",
);
let runtime_err = anyhow::anyhow!("tool panicked");
assert_eq!(
failure_classification(&agent, Some(&runtime_err)),
"runtime"
);
let internally_cancelled = make_agent(vec![]);
internally_cancelled.cancel_token.cancel();
assert_eq!(
failure_classification(&internally_cancelled, Some(&runtime_err)),
"internal_cancel"
);
assert_eq!(
failure_classification(&internally_cancelled, None),
"internal_cancel"
);
}
#[tokio::test]
#[serial_test::serial(drain)] async fn failure_classification_recognizes_drain() {
let agent = make_agent(vec![]);
assert_eq!(failure_classification(&agent, None), "runtime");
crate::shutdown::drain_begin();
assert_eq!(
failure_classification(&agent, None),
"drain",
"a drained agent must classify as drain, never failure",
);
agent.cancel_token.cancel();
assert_eq!(failure_classification(&agent, None), "drain");
crate::shutdown::drain_clear();
}
async fn assert_scrubbed(should_scrub: bool) {
let tool: Box<dyn Tool> = Box::new(TestTool {
output: SCRUBBABLE_LINE.into(),
scrub: should_scrub,
});
let name = tool.name();
let agent = make_agent(vec![tool]);
let out = agent.execute_tool(name, serde_json::json!({})).await;
assert!(out.success, "{name} should succeed");
if should_scrub {
assert!(out.output.contains("[REDACTED]"), "{name} should redact");
assert!(
!out.output.contains("abcdef"),
"{name} should not leak original"
);
} else {
assert!(
!out.output.contains("[REDACTED]"),
"{name} should not redact"
);
assert!(
out.output.contains(SCRUBBABLE_LINE),
"{name} should preserve output"
);
}
}
#[tokio::test]
async fn tool_with_scrub_enabled_scrubs_sensitive_output() {
assert_scrubbed(true).await;
}
struct MediaTestTool {
name: &'static str,
marker: &'static str,
}
#[async_trait]
impl Tool for MediaTestTool {
fn name(&self) -> &'static str {
self.name
}
fn description(&self) -> String {
"media test tool".into()
}
fn parameters_schema(&self) -> serde_json::Value {
serde_json::json!({})
}
async fn execute(
&self,
_ws: &crate::Workspace,
_args: serde_json::Value,
) -> anyhow::Result<String> {
Ok(String::new())
}
fn media_marker(&self) -> Option<&'static str> {
Some(self.marker)
}
}
#[test]
#[expect(clippy::too_many_lines)]
fn extract_media_outcomes_consolidated() {
enum ToolDef {
Media {
name: &'static str,
marker: &'static str,
},
NonMedia,
}
struct OutcomeDef {
output: &'static str,
success: bool,
}
struct TestCase {
name: &'static str,
msg: &'static str,
tools: Vec<ToolDef>,
outcomes: Vec<OutcomeDef>,
expected: Vec<(&'static str, &'static str)>,
}
let cases = vec![
TestCase {
name: "parses_valid_marker",
msg: "valid marker with success=true should extract the path",
tools: vec![ToolDef::Media {
name: "image_gen",
marker: "[IMAGE:",
}],
outcomes: vec![OutcomeDef {
output: "[IMAGE:/tmp/img.png]",
success: true,
}],
expected: vec![("[IMAGE:", "/tmp/img.png")],
},
TestCase {
name: "skips_malformed_marker",
msg: "malformed marker should be skipped",
tools: vec![ToolDef::Media {
name: "image_gen",
marker: "[IMAGE:",
}],
outcomes: vec![OutcomeDef {
output: "description text [IMAGE:",
success: true,
}],
expected: vec![],
},
TestCase {
name: "skips_empty_marker",
msg: "empty marker '[IMAGE:]' should be skipped",
tools: vec![ToolDef::Media {
name: "image_gen",
marker: "[IMAGE:",
}],
outcomes: vec![OutcomeDef {
output: "[IMAGE:]",
success: true,
}],
expected: vec![],
},
TestCase {
name: "skips_no_closing_bracket_non_empty_path",
msg: "output with '[IMAGE:bogus' (no closing bracket, non-empty path) should be skipped",
tools: vec![ToolDef::Media {
name: "image_gen",
marker: "[IMAGE:",
}],
outcomes: vec![OutcomeDef {
output: "oops [IMAGE:bogus",
success: true,
}],
expected: vec![],
},
TestCase {
name: "skips_non_media_tool",
msg: "non-media tool should not be inspected for media markers",
tools: vec![ToolDef::NonMedia],
outcomes: vec![OutcomeDef {
output: "[IMAGE:path]",
success: true,
}],
expected: vec![],
},
TestCase {
name: "skips_failed_outcome",
msg: "failed outcomes should not produce media paths",
tools: vec![ToolDef::Media {
name: "image_gen",
marker: "[IMAGE:",
}],
outcomes: vec![OutcomeDef {
output: "[IMAGE:/tmp/img.png]",
success: false,
}],
expected: vec![],
},
TestCase {
name: "handles_mixed_tools",
msg: "mixed tools with valid outcomes should extract only media paths",
tools: vec![
ToolDef::Media {
name: "image_gen",
marker: "[IMAGE:",
},
ToolDef::NonMedia,
ToolDef::Media {
name: "video_gen",
marker: "[VIDEO:",
},
],
outcomes: vec![
OutcomeDef {
output: "[IMAGE:/tmp/img.png]",
success: true,
},
OutcomeDef {
output: "non-media output",
success: true,
},
OutcomeDef {
output: "[VIDEO:/tmp/vid.mp4]",
success: true,
},
],
expected: vec![("[IMAGE:", "/tmp/img.png"), ("[VIDEO:", "/tmp/vid.mp4")],
},
];
for case in cases {
let tools: Vec<Box<dyn Tool>> = case
.tools
.iter()
.map(|t| match t {
ToolDef::Media { name, marker } => {
Box::new(MediaTestTool { name, marker }) as Box<dyn Tool>
}
ToolDef::NonMedia => Box::new(TestTool {
output: String::new(),
scrub: false,
}) as Box<dyn Tool>,
})
.collect();
let calls: Vec<ToolCall> = case
.tools
.iter()
.enumerate()
.map(|(i, t)| {
let name = match t {
ToolDef::Media { name, .. } => *name,
ToolDef::NonMedia => "never_scrub",
};
ToolCall {
id: (i + 1).to_string(),
name: name.to_string(),
arguments: serde_json::json!({}),
}
})
.collect();
let outcomes: Vec<ToolExecutionOutcome> = case
.outcomes
.iter()
.map(|o| ToolExecutionOutcome {
output: o.output.to_string(),
success: o.success,
image_payloads: Vec::new(),
text_is_content: false,
suspended: false,
ends_turn: false,
})
.collect();
let expected: Vec<(&'static str, String)> = case
.expected
.iter()
.map(|(n, p)| (*n, p.to_string()))
.collect();
let result = extract_media_from_outcomes(&tools, &calls, &outcomes);
assert_eq!(result, expected, "case '{}': {}", case.name, case.msg);
}
}
#[tokio::test]
async fn finalize_session_skipped_when_cancelled() {
let mut agent = make_agent(vec![]);
agent.cancel_token.cancel();
let result = agent.finalize_session().await;
assert!(
result.is_ok(),
"finalize_session should return Ok when cancelled, \
skipping the 'no assistant message' warning"
);
}
#[tokio::test]
async fn execute_tool_skips_when_cancelled() {
let tool: Box<dyn Tool> = Box::new(TestTool {
output: "should not run".into(),
scrub: false,
});
let name = tool.name();
let agent = make_agent(vec![tool]);
agent.cancel_token.cancel();
let out = agent.execute_tool(name, serde_json::json!({})).await;
assert!(!out.success, "cancelled agent should not execute tool");
assert!(
out.output.contains("Agent cancelled"),
"failure output should mention cancellation: {}",
out.output
);
}
#[tokio::test]
async fn task_locals_propagate_to_parallel_tool_execution() {
struct ReadTaskLocalsTool;
#[async_trait]
impl Tool for ReadTaskLocalsTool {
fn name(&self) -> &'static str {
"read_task_locals"
}
fn description(&self) -> String {
"test tool that reads task-local user context".into()
}
fn parameters_schema(&self) -> serde_json::Value {
serde_json::json!({})
}
async fn execute(
&self,
_ws: &crate::Workspace,
_args: serde_json::Value,
) -> anyhow::Result<String> {
let user_name = tool_user_name();
let channel = tool_channel();
Ok(format!("user={user_name},channel={channel}"))
}
}
let tool: Box<dyn Tool> = Box::new(ReadTaskLocalsTool);
let name = tool.name();
let mut agent = make_agent(vec![tool]);
agent.user_name = "alice".into();
agent.channel = "gui".into();
let call = ToolCall {
id: "call_1".into(),
name: name.to_string(),
arguments: serde_json::json!({}),
};
let outcomes = agent.execute_tool_group(&[call]).await;
assert_eq!(
outcomes.len(),
1,
"one tool call should produce one outcome"
);
assert!(outcomes[0].success, "ReadTaskLocalsTool should succeed");
assert!(
outcomes[0].output.contains("user=alice"),
"should propagate user_name: {}",
outcomes[0].output
);
assert!(
outcomes[0].output.contains("channel=gui"),
"should propagate channel: {}",
outcomes[0].output
);
}
#[tokio::test]
async fn test_drain_incoming_messages_injects_ticket_comment() {
crate::util::test::init_management_test_stores().await;
let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
let mut agent = make_agent(vec![]);
agent.incoming_rx = Some(rx);
agent.agent_id = "_test_drain_ticket_comment".into();
let job = crate::agent::message_router::AgentJob {
content: "Please fix the formatting".to_string(),
workspace_name: "test_ws".to_string(),
user_name: "manager".to_string(),
channel: String::new(),
kind: crate::agent::message_router::MessageKind::TicketComment,
role: crate::Role::Manager,
reply_target: None,
pending_job_id: None,
originating_workspace: None,
};
let _ = tx.send(job);
agent.drain_incoming_messages().await;
let history = agent.session.history();
assert!(!history.is_empty(), "should have at least one message");
let last = history.last().unwrap();
assert_eq!(last.role, crate::ChatRole::User, "should be a user message");
assert!(
last.content.contains("Please fix the formatting"),
"should contain the comment content: {}",
last.content,
);
assert!(
last.content.contains("[Comment from manager on ticket]"),
"should include comment prefix: {}",
last.content,
);
}
#[tokio::test]
async fn test_drain_incoming_messages_non_comment() {
crate::util::test::init_management_test_stores().await;
let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
let mut agent = make_agent(vec![]);
agent.incoming_rx = Some(rx);
agent.agent_id = "_test_drain_non_comment".into();
let job = crate::agent::message_router::AgentJob {
content: "Hello agent".to_string(),
workspace_name: "test_ws".to_string(),
user_name: "user".to_string(),
channel: String::new(),
kind: crate::agent::message_router::MessageKind::UserMessage,
role: crate::Role::Assistant,
reply_target: None,
pending_job_id: None,
originating_workspace: None,
};
let _ = tx.send(job);
agent.drain_incoming_messages().await;
let history = agent.session.history();
assert!(!history.is_empty(), "should have at least one message");
let last = history.last().unwrap();
assert!(
last.content.contains("Hello agent"),
"should contain the raw content: {}",
last.content,
);
}
#[tokio::test]
async fn test_drain_incoming_messages_disconnected() {
let (tx, rx) =
tokio::sync::mpsc::unbounded_channel::<crate::agent::message_router::AgentJob>();
let mut agent = make_agent(vec![]);
agent.incoming_rx = Some(rx);
agent.agent_id = "_test_drain_disconnected".into();
drop(tx);
agent.drain_incoming_messages().await;
assert!(
agent.incoming_rx.is_none(),
"incoming_rx should be set to None after disconnect",
);
}
async fn make_inject_agent(
agent_id: &str,
role: crate::Role,
ticket: crate::pipeline::board::Ticket,
seed_content: Option<&str>,
) -> Agent {
if let Some(content) = seed_content {
let now = crate::db::now();
crate::session::store()
.conn
.execute(
"INSERT INTO sessions (agent_id, role, content, created_at) \
VALUES (?1, 'user', ?2, ?3)",
crate::db::params![agent_id, content, now],
)
.await
.unwrap();
}
let mut agent = make_agent_with_role(vec![], role);
agent.agent_id = agent_id.to_string();
agent.session.init(agent_id).await.unwrap();
agent.ticket = Some(ticket);
agent
}
#[tokio::test]
async fn test_inject_outstanding_comments_injects_for_stage_roles() {
crate::util::test::init_management_test_stores().await;
let cases: &[(
&str,
crate::Role,
&str,
&str,
crate::pipeline::board::TicketPhase,
)] = &[
(
"inject-new-agent",
crate::Role::Engineer,
"inject_new",
"please fix the formatting",
crate::pipeline::board::TicketPhase::InDevelopment,
),
(
"inject-sanitation-agent",
crate::Role::Sanitation,
"inject_sanitation",
"please verify the cleanup",
crate::pipeline::board::TicketPhase::InSanitation,
),
];
for &(agent_id, role, ws_name, comment, phase) in cases {
let ws = crate::workspace::test_ws_named("/tmp/test", ws_name);
let ticket_id = crate::util::test::make_ticket(
crate::pipeline::board::store(),
&ws,
"Inject Comment",
phase,
)
.await;
crate::pipeline::board::store()
.add_comment(&ticket_id, "manager", comment)
.await
.unwrap();
let ticket =
crate::util::test::expect_ticket(crate::pipeline::board::store(), &ticket_id).await;
let mut agent = make_inject_agent(agent_id, role, ticket, Some("round 1 task")).await;
agent.inject_outstanding_comments(true, "").await;
let history = agent.session.history();
assert!(
history.iter().any(|m| {
m.content.contains("[Comment from manager on ticket]")
&& m.content.contains(comment)
}),
"a new board comment should be injected into a resumed {role} session",
);
}
}
#[tokio::test]
async fn test_inject_outstanding_comments_dedups_existing_content() {
crate::util::test::init_management_test_stores().await;
let ws = crate::workspace::test_ws_named("/tmp/test", "inject_dedup");
let ticket_id = crate::util::test::make_ticket(
crate::pipeline::board::store(),
&ws,
"Inject Dedup",
crate::pipeline::board::TicketPhase::InDevelopment,
)
.await;
let content = "the comment content already seen";
crate::pipeline::board::store()
.add_comment(&ticket_id, "manager", content)
.await
.unwrap();
let ticket =
crate::util::test::expect_ticket(crate::pipeline::board::store(), &ticket_id).await;
let seed = format!("[Comment from manager on ticket]: {content}");
let mut agent = make_inject_agent(
"inject-dedup-agent",
crate::Role::Engineer,
ticket,
Some(&seed),
)
.await;
agent.inject_outstanding_comments(true, "").await;
let history = agent.session.history();
assert_eq!(
history.len(),
1,
"an already-seen comment must not be re-injected",
);
}
#[tokio::test]
async fn test_inject_outstanding_comments_does_not_special_case_own_role() {
crate::util::test::init_management_test_stores().await;
let ws = crate::workspace::test_ws_named("/tmp/test", "inject_own_role");
let ticket_id = crate::util::test::make_ticket(
crate::pipeline::board::store(),
&ws,
"Inject Own Role",
crate::pipeline::board::TicketPhase::InDevelopment,
)
.await;
crate::pipeline::board::store()
.add_comment(&ticket_id, "engineer", "note from the engineer")
.await
.unwrap();
let ticket =
crate::util::test::expect_ticket(crate::pipeline::board::store(), &ticket_id).await;
let mut agent = make_inject_agent(
"inject-own-role-agent",
crate::Role::Engineer,
ticket,
Some("round 1 task"),
)
.await;
agent.inject_outstanding_comments(true, "").await;
let history = agent.session.history();
assert!(
history
.iter()
.any(|m| m.content.contains("[Comment from engineer on ticket]")),
"an own-role comment absent from history must still be injected",
);
}
#[tokio::test]
async fn test_inject_outstanding_comments_noop_for_feedback_and_fresh() {
crate::util::test::init_management_test_stores().await;
let ws = crate::workspace::test_ws_named("/tmp/test", "inject_noop");
let ticket_id = crate::util::test::make_ticket(
crate::pipeline::board::store(),
&ws,
"Inject Noop",
crate::pipeline::board::TicketPhase::InDevelopment,
)
.await;
crate::pipeline::board::store()
.add_comment(&ticket_id, "manager", "please fix the formatting")
.await
.unwrap();
let ticket =
crate::util::test::expect_ticket(crate::pipeline::board::store(), &ticket_id).await;
let mut agent = make_inject_agent(
"inject-noop-agent",
crate::Role::Engineer,
ticket,
Some("round 1 task"),
)
.await;
agent
.inject_outstanding_comments(true, "round 2 feedback")
.await;
agent.inject_outstanding_comments(false, "").await;
let history = agent.session.history();
assert!(
!history
.iter()
.any(|m| m.content.contains("[Comment from manager on ticket]")),
"feedback path and fresh round must not inject comments",
);
}
#[tokio::test]
async fn test_inject_outstanding_comments_injects_for_non_stage_role() {
crate::util::test::init_management_test_stores().await;
let ws = crate::workspace::test_ws_named("/tmp/test", "inject_analyst");
let ticket_id = crate::util::test::make_ticket(
crate::pipeline::board::store(),
&ws,
"Inject Analyst",
crate::pipeline::board::TicketPhase::InDevelopment,
)
.await;
crate::pipeline::board::store()
.add_comment(&ticket_id, "manager", "please fix the formatting")
.await
.unwrap();
let ticket =
crate::util::test::expect_ticket(crate::pipeline::board::store(), &ticket_id).await;
let mut agent = make_inject_agent(
"inject-analyst-agent",
crate::Role::Analyst,
ticket,
Some("round 1 task"),
)
.await;
agent.inject_outstanding_comments(true, "").await;
let history = agent.session.history();
assert!(
history
.iter()
.any(|m| m.content.contains("[Comment from manager on ticket]")),
"injection is no longer role-gated — non-stage roles reaching the resume path inject",
);
}
#[tokio::test]
async fn spawn_staggered_round_single_member_is_noop() {
let handles = crate::agent::spawn_staggered_round(
vec![move |round: crate::agent::RoundOpts| async move {
assert!(
round.first_call_notify.is_none(),
"sole member must not receive a signal"
);
1u8
}],
false,
)
.await;
assert_eq!(handles.len(), 1);
assert_eq!(handles.into_iter().next().unwrap().await.unwrap(), 1);
}
#[tokio::test]
async fn spawn_staggered_round_resume_skips_stagger() {
let handles = crate::agent::spawn_staggered_round(
(0..3)
.map(|i| {
move |round: crate::agent::RoundOpts| async move {
assert!(
round.first_call_notify.is_none(),
"resume must not stagger (member {i})"
);
i
}
})
.collect(),
true,
)
.await;
assert_eq!(handles.len(), 3);
let mut out = Vec::new();
for h in handles {
out.push(h.await.unwrap());
}
assert_eq!(out, vec![0, 1, 2]);
}
#[tokio::test]
async fn spawn_staggered_round_leader_notify_releases_followers() {
let handles = crate::agent::spawn_staggered_round(
(0..3)
.map(|i| {
move |round: crate::agent::RoundOpts| async move {
let tag = if i == 0 {
round
.first_call_notify
.expect("leader receives the signal")
.notify_one();
"leader".to_string()
} else {
assert!(
round.first_call_notify.is_none(),
"followers must not receive the signal"
);
"follower".to_string()
};
(tag, round.round_ts)
}
})
.collect(),
false,
)
.await;
assert_eq!(handles.len(), 3);
let mut out = Vec::new();
for h in handles {
out.push(h.await.unwrap());
}
assert!(out.iter().any(|(tag, _)| tag == "leader"));
assert_eq!(out.iter().filter(|(tag, _)| tag == "follower").count(), 2);
assert_eq!(
out.iter()
.map(|(_, ts)| ts)
.collect::<std::collections::HashSet<_>>()
.len(),
1,
"all members share one round-fixed timestamp"
);
}
#[tokio::test]
async fn spawn_staggered_round_leader_timeout_fail_open() {
let _guard = crate::util::test::set_env_var("MAHBOT_STAGGER_WAIT_SECS", Some("0"));
let handles = crate::agent::spawn_staggered_round(
(0..2)
.map(|i| {
move |round: crate::agent::RoundOpts| async move {
if i == 0 {
let _ = round; "stuck-leader".to_string()
} else {
assert!(
round.first_call_notify.is_none(),
"released follower gets no signal"
);
"released".to_string()
}
}
})
.collect(),
false,
)
.await;
assert_eq!(
handles.len(),
2,
"followers must be released despite the stuck leader"
);
for h in handles {
h.await.unwrap();
}
}
#[tokio::test]
#[serial_test::serial(provider, drain)]
async fn leader_first_call_failure_still_fires_signal() {
crate::util::test::init_test_stores().await;
let _policy_guard = install_test_retry_policy(crate::retry::tiny_test_policy());
let ws = crate::workspace::test_ws_named("/tmp/ws_leader_fail", "leader_fail");
let notify = std::sync::Arc::new(tokio::sync::Notify::new());
let provider = FakeProvider::new()
.err(crate::retry::FailureClass::Transport, "boom")
.err(crate::retry::FailureClass::Transport, "boom")
.err(crate::retry::FailureClass::Transport, "boom");
let _provider = install_fake_provider(std::sync::Arc::new(provider));
let (_agent, response) = run_agent(
"leader_fail_agent".to_string(),
crate::Role::Analyst,
&ws,
None,
"task",
String::new(),
String::new(),
false,
None,
false,
Some(RoundOpts {
round_ts: crate::session::render_timestamp(),
first_call_notify: Some(notify.clone()),
}),
None,
None,
)
.await;
assert!(
response.is_none(),
"leader round must fail with an exhausted budget"
);
tokio::time::timeout(std::time::Duration::from_secs(5), notify.notified())
.await
.expect("first-call signal must fire even when the call fails");
}
#[tokio::test]
#[serial_test::serial(provider, drain)]
async fn sleep_call_ends_the_run_gracefully_and_flag_resets_on_reengage() {
crate::util::test::init_test_stores().await;
let fake = std::sync::Arc::new(FakeProvider::new().ok_tool_call("sleep").ok("back online"));
let provider: std::sync::Arc<dyn crate::Provider> = fake.clone();
let _provider_guard = install_fake_provider(provider);
let ws = crate::workspace::test_ws_named("/tmp/ws_sleep_stop", "sleep_stop");
let agent_id = "sleep_stop_agent".to_string();
let (agent, response) = run_agent(
agent_id.clone(),
crate::Role::Assistant,
&ws,
None,
"hello",
"sleep_user".to_string(),
"gui".to_string(),
false,
None,
false,
None,
None,
None,
)
.await;
assert!(agent.sleep_ended, "sleep must mark the turn sleep-ended");
assert_eq!(
response.as_deref(),
Some(""),
"the turn must end silently — any text would mean the loop continued"
);
assert_eq!(
fake.request_messages.lock().unwrap().len(),
1,
"no LLM call may follow the committed sleep round"
);
let (tail_role, tail_content) = crate::session::store()
.get_last_message_tail(&agent_id)
.await
.expect("session must have history");
assert_eq!(tail_role, crate::ChatRole::Tool);
assert!(
tail_content.contains("\"tool_call_id\":\"call_test\"")
&& tail_content.contains("\"content\":\"Zzz...\""),
"sleep tool result must be committed as a normal tool result: {tail_content}"
);
assert!(
crate::session::store().get_sleep_ended(&agent_id).await,
"the durable sleep-ended flag must be set"
);
let (agent2, response2) = run_agent(
agent_id.clone(),
crate::Role::Assistant,
&ws,
None,
"wake up",
"sleep_user".to_string(),
"gui".to_string(),
false,
None,
false,
None,
None,
None,
)
.await;
assert_eq!(response2.as_deref(), Some("back online"));
assert!(!agent2.sleep_ended);
assert!(
!crate::session::store().get_sleep_ended(&agent_id).await,
"the flag must be cleared when the session re-engages"
);
}
#[tokio::test]
#[serial_test::serial(provider, drain)]
async fn sleep_with_text_is_rejected_and_resend_is_delivered() {
crate::util::test::init_test_stores().await;
let fake = std::sync::Arc::new(
FakeProvider::new()
.ok_text_and_tool_calls("here is the answer", &[("sleep", serde_json::json!({}))])
.ok("here is the answer"),
);
let provider: std::sync::Arc<dyn crate::Provider> = fake.clone();
let _provider_guard = install_fake_provider(provider);
let ws = crate::workspace::test_ws_named("/tmp/ws_sleep_text", "sleep_text");
let agent_id = "sleep_text_agent".to_string();
let (agent, response) = run_agent(
agent_id.clone(),
crate::Role::Assistant,
&ws,
None,
"hello",
"sleep_text_user".to_string(),
"gui".to_string(),
false,
None,
false,
None,
None,
None,
)
.await;
assert_eq!(
response.as_deref(),
Some("here is the answer"),
"the text must be resent as a plain round after the rejection"
);
assert!(!agent.sleep_ended, "a rejected sleep must not end the run");
assert_eq!(
fake.request_messages.lock().unwrap().len(),
2,
"the rejection round plus the plain resend round"
);
let history = agent.session.history();
let rejection_tools: Vec<String> = history
.iter()
.filter(|m| m.role == crate::ChatRole::Tool)
.map(|m| {
serde_json::from_str::<crate::ToolResultPayload>(&m.content)
.expect("tool message must be a ToolResultPayload")
.content
})
.collect();
assert!(
rejection_tools
.iter()
.any(|c| c.contains("NOT delivered") && !c.contains("Zzz...")),
"the sleep rejection must be committed without a Zzz... result: {rejection_tools:?}"
);
assert!(
!crate::session::store().get_sleep_ended(&agent_id).await,
"the durable sleep-ended flag must stay clear"
);
}
#[tokio::test]
#[serial_test::serial(provider, drain)]
async fn sleep_bundled_with_other_tool_is_rejected_but_sibling_executes() {
crate::util::test::init_test_stores().await;
let dir = tempfile::tempdir().unwrap();
std::fs::write(dir.path().join("seed.txt"), "bundle-marker\n").unwrap();
let ws = crate::Workspace::from_path(dir.path());
let fake = std::sync::Arc::new(
FakeProvider::new()
.ok_text_and_tool_calls(
"",
&[
("read", serde_json::json!({"path": "seed.txt"})),
("sleep", serde_json::json!({})),
],
)
.ok("done"),
);
let provider: std::sync::Arc<dyn crate::Provider> = fake.clone();
let _provider_guard = install_fake_provider(provider);
let agent_id = "sleep_bundle_agent".to_string();
let (agent, response) = run_agent(
agent_id.clone(),
crate::Role::Assistant,
&ws,
None,
"hello",
"sleep_bundle_user".to_string(),
"gui".to_string(),
false,
None,
false,
None,
None,
None,
)
.await;
assert_eq!(response.as_deref(), Some("done"));
assert!(!agent.sleep_ended, "a bundled sleep must not end the run");
assert_eq!(fake.request_messages.lock().unwrap().len(), 2);
let history = agent.session.history();
let tool_results: Vec<String> = history
.iter()
.filter(|m| m.role == crate::ChatRole::Tool)
.map(|m| {
serde_json::from_str::<crate::ToolResultPayload>(&m.content)
.expect("tool message must be a ToolResultPayload")
.content
})
.collect();
assert!(
tool_results.iter().any(|c| c.contains("bundle-marker")),
"the read sibling must execute normally: {tool_results:?}"
);
assert!(
tool_results
.iter()
.any(|c| c.contains("bundled with other tool calls")),
"the bundled sleep must be rejected: {tool_results:?}"
);
}
#[tokio::test]
#[serial_test::serial(provider, drain)]
async fn resume_completion_rejects_text_sleep_frame() {
crate::util::test::init_test_stores().await;
let dir = tempfile::tempdir().unwrap();
let ws = crate::Workspace::from_path(dir.path());
let mut session = Session::default();
let frame = crate::providers::reasoning::assistant_replay_payload(
Some("the answer"),
&[crate::ToolCall {
id: "call_sleep_r".to_string(),
name: "sleep".to_string(),
arguments: serde_json::json!({}),
}],
None,
)
.to_string();
session
.persist_messages(
"test-sleep-resume",
&[ChatMessage::user("wake"), ChatMessage::assistant(frame)],
)
.await
.unwrap();
let mut agent = make_agent_on(
vec![Box::new(crate::tools::SleepTool)],
"test-sleep-resume",
ws,
);
agent.session = session;
let settled = agent.complete_pending_tool_calls().await.unwrap();
assert!(
settled,
"a text-accompanied sleep frame settles the rejection"
);
let history = agent.session.history();
let roles: Vec<crate::ChatRole> = history.iter().map(|m| m.role).collect();
assert_eq!(
roles,
vec![
crate::ChatRole::User,
crate::ChatRole::Assistant,
crate::ChatRole::Tool
],
"result contiguous after the frame: {roles:?}"
);
let payload: crate::ToolResultPayload = serde_json::from_str(&history[2].content).unwrap();
assert_eq!(
payload.tool_call_id, "call_sleep_r",
"rejection keeps the ORIGINAL call id"
);
assert!(
payload.content.contains("NOT delivered"),
"sleep rejection content: {}",
payload.content
);
assert!(!payload.content.contains("Zzz..."));
assert!(agent.session.pending_tool_frame().is_none());
}
#[test]
fn reasoning_only_stop_classification() {
let reasoning = || {
Some(crate::Reasoning {
reasoning: Some("thinking".into()),
reasoning_content: Some("thinking".into()),
reasoning_details: None,
})
};
for finish in [None, Some("stop"), Some("length"), Some("tool_calls")] {
let resp = crate::ChatResponse {
text: None,
reasoning: reasoning(),
finish_reason: finish.map(str::to_string),
..crate::ChatResponse::default()
};
assert!(is_reasoning_only_stop(&resp), "finish_reason={finish:?}");
}
let resp = crate::ChatResponse {
text: Some(String::new()),
reasoning: reasoning(),
..crate::ChatResponse::default()
};
assert!(is_reasoning_only_stop(&resp));
let resp = crate::ChatResponse {
text: Some(" \n ".into()),
..crate::ChatResponse::default()
};
assert!(is_reasoning_only_stop(&resp));
let resp = crate::ChatResponse::default();
assert!(is_reasoning_only_stop(&resp));
let resp = crate::ChatResponse {
text: None,
tool_calls: vec![crate::ToolCall {
id: "t1".into(),
name: "read".into(),
arguments: serde_json::json!({}),
}],
..crate::ChatResponse::default()
};
assert!(!is_reasoning_only_stop(&resp));
let resp = crate::ChatResponse {
text: Some("real answer".into()),
reasoning: reasoning(),
..crate::ChatResponse::default()
};
assert!(!is_reasoning_only_stop(&resp));
}
#[tokio::test]
#[serial_test::serial(provider, drain)]
async fn llm_call_recovers_reasoning_only_stop_via_continuation() {
crate::util::test::init_test_stores().await;
let _policy_guard = install_test_retry_policy(crate::retry::tiny_test_policy());
let fake = std::sync::Arc::new(
FakeProvider::new()
.ok_reasoning_only("draft plan then execute tool", Some("stop"))
.ok("final answer"),
);
let provider: std::sync::Arc<dyn crate::Provider> = fake.clone();
let _provider_guard = install_fake_provider(provider);
let mut agent = make_agent(vec![]);
let resp = agent.llm_call().await.expect("continuation must resolve");
assert_eq!(resp.text_or_empty(), "final answer");
let fingerprints = fake.request_fingerprints.lock().unwrap().clone();
assert_eq!(fingerprints.len(), 2, "original call + one continuation");
assert!(
!fingerprints[0].contains("Resume your unfinished turn"),
"original request must not carry the continuation tail"
);
assert!(
fingerprints[1].contains("Resume your unfinished turn"),
"continuation request must carry the appended nudge"
);
assert!(
fingerprints[1].contains("draft plan then execute tool"),
"continuation must echo the previous reasoning as the assistant turn"
);
}
#[tokio::test]
#[serial_test::serial(provider, drain)]
async fn llm_call_continuation_accumulates_tail_until_answer() {
crate::util::test::init_test_stores().await;
let _policy_guard = install_test_retry_policy(crate::retry::tiny_test_policy());
let fake = std::sync::Arc::new(
FakeProvider::new()
.ok_reasoning_only("thinking 1", Some("stop"))
.ok_reasoning_only("thinking 2", Some("stop"))
.ok("answer after two continuations"),
);
let provider: std::sync::Arc<dyn crate::Provider> = fake.clone();
let _provider_guard = install_fake_provider(provider);
let mut agent = make_agent(vec![]);
let resp = agent.llm_call().await.expect("continuation must resolve");
assert_eq!(resp.text_or_empty(), "answer after two continuations");
let fingerprints = fake.request_fingerprints.lock().unwrap().clone();
assert_eq!(fingerprints.len(), 3);
assert_eq!(
fingerprints[1]
.matches("Resume your unfinished turn")
.count(),
1,
"attempt 2 carries exactly the first appended pair"
);
assert_eq!(
fingerprints[2]
.matches("Resume your unfinished turn")
.count(),
2,
"attempt 3 carries both appended pairs"
);
assert!(fingerprints[2].contains("thinking 1"));
assert!(fingerprints[2].contains("thinking 2"));
let messages = fake.request_messages.lock().unwrap().clone();
assert_eq!(messages.len(), 3);
assert!(
messages[2].starts_with(&messages[1]),
"byte-stable prefix: attempt 3's messages begin with attempt 2's verbatim"
);
}
#[tokio::test]
#[serial_test::serial(provider, drain)]
async fn llm_call_continuation_exhaustion_fails_safely_without_leaking() {
crate::util::test::init_test_stores().await;
let _policy_guard = install_test_retry_policy(crate::retry::tiny_test_policy());
let fake = std::sync::Arc::new(
FakeProvider::new()
.ok_reasoning_only("secret thinking alpha", Some("stop"))
.ok_reasoning_only("secret thinking beta", Some("stop"))
.ok_reasoning_only("secret thinking gamma", Some("stop"))
.ok_reasoning_only("secret thinking delta", Some("stop")),
);
let provider: std::sync::Arc<dyn crate::Provider> = fake.clone();
let _provider_guard = install_fake_provider(provider);
let mut agent = make_agent(vec![]);
let err = agent
.llm_call()
.await
.expect_err("continuation must exhaust");
let exhausted = err
.chain()
.find_map(|c| c.downcast_ref::<crate::retry::RetryExhausted>())
.expect("RetryExhausted must survive in the error chain");
assert_eq!(
exhausted.final_class,
crate::retry::FailureClass::NoResponse,
"granular no-response classification"
);
assert_eq!(
exhausted.last_raw, None,
"no raw text on the exhausted error"
);
let last_failure = exhausted
.failures
.last()
.expect("failure trail is non-empty");
assert_eq!(
last_failure.finish_reason.as_deref(),
Some("stop"),
"in-class NoResponse records carry the response finish_reason into the telemetry trail"
);
let rendered = format!("{err:#}");
assert!(
!rendered.contains("secret thinking"),
"the thinking must never leak into the failure error"
);
assert!(
!rendered.contains(RETRY_EXHAUSTION_MARKER),
"must not be misclassified as LLM provider retry exhaustion"
);
assert!(
agent.session.history().is_empty(),
"the continuation tail must never reach the session transcript"
);
assert_eq!(failure_classification(&agent, Some(&err)), "no_response");
}
#[tokio::test]
#[serial_test::serial(provider, drain)]
async fn llm_call_continuation_transport_error_does_not_duplicate_tail() {
crate::util::test::init_test_stores().await;
let _policy_guard = install_test_retry_policy(crate::retry::tiny_test_policy());
let fake = std::sync::Arc::new(
FakeProvider::new()
.ok_reasoning_only("thinking 1", Some("stop"))
.err(crate::retry::FailureClass::Transport, "connection reset")
.err(crate::retry::FailureClass::Transport, "connection reset")
.err(crate::retry::FailureClass::Transport, "connection reset"),
);
let provider: std::sync::Arc<dyn crate::Provider> = fake.clone();
let _provider_guard = install_fake_provider(provider);
let mut agent = make_agent(vec![]);
let err = agent
.llm_call()
.await
.expect_err("continuation must exhaust");
let exhausted = err
.chain()
.find_map(|c| c.downcast_ref::<crate::retry::RetryExhausted>())
.expect("RetryExhausted must survive in the error chain");
assert_eq!(
exhausted.final_class,
crate::retry::FailureClass::Transport,
"final class derives from the last recorded failure, not NoResponse"
);
let fingerprints = fake.request_fingerprints.lock().unwrap().clone();
assert_eq!(
fingerprints.len(),
4,
"original call + 3 continuation attempts"
);
assert_eq!(
fingerprints[1], fingerprints[2],
"attempt 2 re-sends the byte-identical request after a transport error"
);
assert_eq!(fingerprints[2], fingerprints[3]);
let rendered = format!("{err:#}");
assert!(
!rendered.contains("thinking"),
"the thinking must never leak into the failure error"
);
}
#[tokio::test]
#[serial_test::serial(provider, drain)]
async fn llm_call_continuation_non_retryable_error_breaks_immediately() {
crate::util::test::init_test_stores().await;
let _policy_guard = install_test_retry_policy(crate::retry::tiny_test_policy());
let fake = std::sync::Arc::new(
FakeProvider::new()
.ok_reasoning_only("thinking 1", Some("stop"))
.err(crate::retry::FailureClass::NonRetryable, "invalid model"),
);
let provider: std::sync::Arc<dyn crate::Provider> = fake.clone();
let _provider_guard = install_fake_provider(provider);
let mut agent = make_agent(vec![]);
let err = agent.llm_call().await.expect_err("continuation must fail");
let exhausted = err
.chain()
.find_map(|c| c.downcast_ref::<crate::retry::RetryExhausted>())
.expect("RetryExhausted must survive in the error chain");
assert_eq!(
exhausted.final_class,
crate::retry::FailureClass::NonRetryable,
"non-retryable class survives to the terminal error"
);
assert_eq!(
fake.request_fingerprints.lock().unwrap().len(),
2,
"original call + exactly one continuation attempt (no budget burn)"
);
}
#[tokio::test]
#[serial_test::serial(drain)] async fn recover_reasoning_only_stop_abort_classifies_as_shutdown() {
let agent = make_agent(vec![]);
let first = crate::ChatResponse {
text: None,
reasoning: Some(crate::Reasoning {
reasoning: Some("thinking".into()),
reasoning_content: Some("thinking".into()),
reasoning_details: None,
}),
finish_reason: Some("stop".into()),
..crate::ChatResponse::default()
};
crate::shutdown::drain_begin();
let exhausted = agent
.recover_reasoning_only_stop(vec![], first, "agent-continuation")
.await
.expect_err("the drain must break the continuation immediately");
crate::shutdown::drain_clear();
assert_eq!(
exhausted.final_class,
crate::retry::FailureClass::Shutdown,
"a global abort must classify as shutdown, never no_response"
);
}
#[tokio::test]
#[serial_test::serial(provider, drain)]
async fn llm_call_continuation_exhaustion_does_not_update_session_length() {
crate::util::test::init_test_stores().await;
let _policy_guard = install_test_retry_policy(crate::retry::tiny_test_policy());
let fake = std::sync::Arc::new(
FakeProvider::new()
.ok_reasoning_only_with_usage("thinking a", Some("stop"), 1_000, 500)
.ok_reasoning_only_with_usage("thinking b", Some("stop"), 1_000, 500)
.ok_reasoning_only_with_usage("thinking c", Some("stop"), 1_000, 500)
.ok_reasoning_only_with_usage("thinking d", Some("stop"), 1_000, 500),
);
let provider: std::sync::Arc<dyn crate::Provider> = fake.clone();
let _provider_guard = install_fake_provider(provider);
let mut agent = make_agent(vec![]);
assert_eq!(agent.session.token_length(), None);
let err = agent
.llm_call()
.await
.expect_err("continuation must exhaust");
assert_eq!(
agent.session.token_length(),
None,
"a failed turn (continuation exhaustion) must not update the session length"
);
let rendered = format!("{err:#}");
assert!(
!rendered.contains("thinking"),
"the thinking must never leak into the failure error"
);
}
#[tokio::test]
#[serial_test::serial(provider, drain)]
async fn llm_call_continuation_success_records_only_resolving_usage() {
crate::util::test::init_test_stores().await;
let _policy_guard = install_test_retry_policy(crate::retry::tiny_test_policy());
let fake = std::sync::Arc::new(
FakeProvider::new()
.ok_reasoning_only_with_usage("thinking a", Some("stop"), 1_000, 500)
.ok_with_usage("final answer", 200, 300),
);
let provider: std::sync::Arc<dyn crate::Provider> = fake.clone();
let _provider_guard = install_fake_provider(provider);
let mut agent = make_agent(vec![]);
let resp = agent.llm_call().await.expect("continuation must resolve");
assert_eq!(resp.text_or_empty(), "final answer");
assert_eq!(
agent.session.token_length(),
Some(500),
"only the resolving continuation response (200 + 300) updates the session length"
);
}
#[tokio::test]
#[serial_test::serial(provider, drain)]
async fn llm_call_skips_continuation_for_normal_and_tool_call_turns() {
crate::util::test::init_test_stores().await;
let _policy_guard = install_test_retry_policy(crate::retry::tiny_test_policy());
{
let fake = std::sync::Arc::new(FakeProvider::new().ok("normal answer"));
let provider: std::sync::Arc<dyn crate::Provider> = fake.clone();
let _provider_guard = install_fake_provider(provider);
let mut agent = make_agent(vec![]);
let resp = agent.llm_call().await.expect("normal answer");
assert_eq!(resp.text_or_empty(), "normal answer");
assert_eq!(
fake.request_fingerprints.lock().unwrap().len(),
1,
"normal answer must not trigger continuation"
);
}
{
let fake = std::sync::Arc::new(FakeProvider::new().ok_tool_call("read"));
let provider: std::sync::Arc<dyn crate::Provider> = fake.clone();
let _provider_guard = install_fake_provider(provider);
let mut agent = make_agent(vec![]);
let resp = agent.llm_call().await.expect("tool-call turn");
assert!(resp.text_or_empty().is_empty());
assert_eq!(resp.tool_calls.len(), 1);
assert_eq!(
fake.request_fingerprints.lock().unwrap().len(),
1,
"tool-call turn (empty text) must not trigger continuation"
);
}
}
#[tokio::test]
#[serial_test::serial(provider, drain)]
async fn summarize_recovers_reasoning_only_stop_via_continuation() {
crate::util::test::init_test_stores().await;
let _policy_guard = install_test_retry_policy(crate::retry::tiny_test_policy());
let fake = std::sync::Arc::new(
FakeProvider::new()
.ok_reasoning_only("thinking about the summary", Some("stop"))
.ok("the summary"),
);
let provider: std::sync::Arc<dyn crate::Provider> = fake.clone();
let _provider_guard = install_fake_provider(provider);
let agent = make_agent(vec![]);
let summary = agent.summarize().await.expect("continuation must resolve");
assert_eq!(summary, "the summary");
let fingerprints = fake.request_fingerprints.lock().unwrap().clone();
assert_eq!(
fingerprints.len(),
2,
"original summary call + one continuation"
);
assert!(fingerprints[1].contains("Resume your unfinished turn"));
}
#[tokio::test]
#[serial_test::serial(provider, drain)]
async fn summarize_continuation_exhaustion_fails_open_without_leaking() {
crate::util::test::init_test_stores().await;
let _policy_guard = install_test_retry_policy(crate::retry::tiny_test_policy());
let fake = std::sync::Arc::new(
FakeProvider::new()
.ok_reasoning_only("summary thinking", Some("stop"))
.ok_reasoning_only("summary thinking", Some("stop"))
.ok_reasoning_only("summary thinking", Some("stop"))
.ok_reasoning_only("summary thinking", Some("stop")),
);
let provider: std::sync::Arc<dyn crate::Provider> = fake.clone();
let _provider_guard = install_fake_provider(provider);
let agent = make_agent(vec![]);
let err = agent
.summarize()
.await
.expect_err("continuation must exhaust");
let rendered = format!("{err:#}");
assert!(
rendered.contains("summarization"),
"fail-open path surfaces the summarization error for maybe_summarize"
);
assert!(
!rendered.contains("summary thinking"),
"the thinking must never leak into the summarization error"
);
}
#[tokio::test]
#[serial_test::serial(provider, drain)]
async fn summarize_reasks_a_tool_call_answer_byte_identically() {
crate::util::test::init_test_stores().await;
let _policy_guard = install_test_retry_policy(crate::retry::tiny_test_policy());
let fake = std::sync::Arc::new(FakeProvider::new().ok_tool_call("read").ok("the summary"));
let provider: std::sync::Arc<dyn crate::Provider> = fake.clone();
let _provider_guard = install_fake_provider(provider);
let agent = make_agent(vec![]);
let summary = agent.summarize().await.expect("the re-ask must resolve");
assert_eq!(summary, "the summary");
let fingerprints = fake.request_fingerprints.lock().unwrap().clone();
assert_eq!(fingerprints.len(), 2, "tool-call answer + one re-ask");
assert_eq!(
fingerprints[0], fingerprints[1],
"the re-ask must re-send the identical request"
);
}
#[tokio::test]
#[serial_test::serial(provider, drain)]
async fn summarize_tool_call_exhaustion_keeps_the_empty_response_error() {
crate::util::test::init_test_stores().await;
let _policy_guard = install_test_retry_policy(crate::retry::tiny_test_policy());
let fake = std::sync::Arc::new(
FakeProvider::new()
.ok_tool_call("read")
.ok_tool_call("read")
.ok_tool_call("read"),
);
let provider: std::sync::Arc<dyn crate::Provider> = fake.clone();
let _provider_guard = install_fake_provider(provider);
let agent = make_agent(vec![]);
let err = agent
.summarize()
.await
.expect_err("the budget must exhaust");
assert!(
format!("{err:#}").contains("summarization produced empty response"),
"the fail-open error is unchanged: {err:#}"
);
assert_eq!(
fake.request_fingerprints.lock().unwrap().len(),
SUMMARIZE_ATTEMPTS as usize,
"the re-ask must stop at the attempt bound"
);
}
#[tokio::test]
async fn record_session_usage_keeps_last_on_missing_or_partial_usage() {
crate::util::test::init_test_stores().await;
let mut agent = make_agent(vec![]);
agent.session.set_token_length(Some(7_000));
agent
.record_session_usage(&crate::ChatResponse::default())
.await;
assert_eq!(agent.session.token_length(), Some(7_000));
let partial_input = crate::ChatResponse {
usage: Some(crate::ProviderUsage {
input_tokens: Some(1_000),
..crate::ProviderUsage::default()
}),
..crate::ChatResponse::default()
};
agent.record_session_usage(&partial_input).await;
assert_eq!(agent.session.token_length(), Some(7_000));
let partial_output = crate::ChatResponse {
usage: Some(crate::ProviderUsage {
output_tokens: Some(1_000),
..crate::ProviderUsage::default()
}),
..crate::ChatResponse::default()
};
agent.record_session_usage(&partial_output).await;
assert_eq!(agent.session.token_length(), Some(7_000));
let full = crate::ChatResponse {
usage: Some(crate::ProviderUsage {
input_tokens: Some(12_000),
output_tokens: Some(300),
..crate::ProviderUsage::default()
}),
..crate::ChatResponse::default()
};
agent.record_session_usage(&full).await;
assert_eq!(agent.session.token_length(), Some(12_300));
agent
.record_session_usage(&crate::ChatResponse::default())
.await;
assert_eq!(
agent.session.token_length(),
Some(12_300),
"a usage-less response must never reset the length to zero"
);
}
#[tokio::test]
async fn record_session_usage_overflow_saturates() {
crate::util::test::init_test_stores().await;
let mut agent = make_agent(vec![]);
let huge = crate::ChatResponse {
usage: Some(crate::ProviderUsage {
input_tokens: Some(u64::MAX),
output_tokens: Some(u64::MAX),
..crate::ProviderUsage::default()
}),
..crate::ChatResponse::default()
};
agent.record_session_usage(&huge).await;
assert_eq!(agent.session.token_length(), Some(u64::MAX));
}
#[tokio::test]
async fn llm_step_failure_context_reports_run_round_and_session_depth() {
crate::util::test::init_test_stores().await;
let mut agent = make_agent(vec![]);
agent
.session
.push_messages_unpersisted(&[ChatMessage::user("hello"), ChatMessage::assistant("hi")]);
agent.session.set_token_length(Some(42_110));
let ctx = llm_step_failure_context(&agent.session, 0);
assert!(ctx.contains("tool round 0"), "run-local round: {ctx}");
assert!(ctx.contains("2 messages"), "session depth: {ctx}");
assert!(ctx.contains("42110 tokens"), "token segment: {ctx}");
agent.session.set_token_length(None);
let ctx = llm_step_failure_context(&agent.session, 3);
assert!(ctx.contains("tool round 3"), "round after increment: {ctx}");
assert!(ctx.contains("2 messages"), "session depth: {ctx}");
assert!(
!ctx.contains("tokens"),
"None omits the token segment: {ctx}"
);
}
#[tokio::test]
async fn maybe_summarize_none_and_below_threshold_are_noops() {
crate::util::test::init_test_stores().await;
let mut agent = make_agent(vec![]);
agent.maybe_summarize().await;
assert_eq!(agent.session.token_length(), None);
agent
.session
.set_token_length(Some(crate::session::SUMMARIZATION_THRESHOLD / 2));
agent.maybe_summarize().await;
}
const IMAGE_REJECTION_BODY: &str = r#"{"error":{"message":"Input image data may contain inappropriate content.","code":"data_inspection_failed","type":"invalid_request_error"}}"#;
const TEXT_REJECTION_BODY: &str = r#"{"error":{"message":"Input data may contain inappropriate content.","code":"data_inspection_failed","type":"invalid_request_error"}}"#;
fn seed_empty_catalogs() {
crate::tools::media_catalog::image::seed_cache(Some(std::sync::Arc::new(
crate::tools::media_catalog::image::ImageCatalog::default(),
)));
crate::tools::media_catalog::video::seed_cache(Some(std::sync::Arc::new(
crate::tools::media_catalog::video::VideoCatalog::default(),
)));
}
#[tokio::test]
#[serial_test::serial(active_models, provider, drain)] async fn rejected_input_image_is_stripped_and_run_continues() {
crate::util::test::init_test_stores().await;
seed_empty_catalogs();
let _policy_guard = install_test_retry_policy(crate::retry::tiny_test_policy());
let agent_id = "e2e_assistant_image_reject";
let fake = std::sync::Arc::new(
FakeProvider::new()
.err_http(400, IMAGE_REJECTION_BODY)
.ok("The image was rejected by the provider's content check."),
);
let provider: std::sync::Arc<dyn crate::Provider> = fake.clone();
let _provider_guard = install_fake_provider(provider);
let mut agent = make_agent_with_role(vec![], crate::Role::Assistant);
agent.agent_id = agent_id.to_string();
let resp = agent
.work("[IMAGE:/tmp/photo.png] describe this photo", false)
.await
.expect("the run must not fail wholesale on a rejected input image");
assert_eq!(
resp,
"The image was rejected by the provider's content check."
);
let messages = fake.request_messages.lock().unwrap().clone();
assert_eq!(messages.len(), 2, "rejected attempt + normal-loop retry");
assert!(
messages[0].contains("[IMAGE:/tmp/photo.png]"),
"first request carried the image"
);
let retried_user = messages[1]
.split('\u{0}')
.next_back()
.expect("retried request has a user segment");
assert!(
!retried_user.contains("[IMAGE:"),
"retried user message no longer carries the image"
);
assert!(
retried_user.contains("rejected by the provider's content-inspection check"),
"retried user message carries the explanatory phrase"
);
assert!(
retried_user.contains("Input image data may contain inappropriate content."),
"phrase embeds the provider reason"
);
let history = crate::session::store().load(agent_id).await;
let last_user = history
.iter()
.rev()
.find(|m| m.role == crate::ChatRole::User)
.expect("user message exists");
assert!(
!last_user.content.contains("[IMAGE:"),
"rejected image durably removed from the session"
);
assert!(
last_user.content.contains("describe this photo"),
"user's accompanying text preserved"
);
assert_eq!(
history
.iter()
.filter(|m| m.role == crate::ChatRole::User)
.count(),
1,
"no separate notification message was added"
);
}
#[tokio::test]
#[serial_test::serial(active_models, provider, drain)] async fn text_content_rejection_follows_normal_failure_path() {
crate::util::test::init_test_stores().await;
seed_empty_catalogs();
let _policy_guard = install_test_retry_policy(crate::retry::tiny_test_policy());
let agent_id = "e2e_assistant_text_reject";
let fake = std::sync::Arc::new(FakeProvider::new().err_http(400, TEXT_REJECTION_BODY));
let provider: std::sync::Arc<dyn crate::Provider> = fake.clone();
let _provider_guard = install_fake_provider(provider);
let mut agent = make_agent_with_role(vec![], crate::Role::Assistant);
agent.agent_id = agent_id.to_string();
let result = agent
.work("[IMAGE:/tmp/photo.png] describe this photo", false)
.await;
assert!(
result.is_err(),
"text-content rejection must fail the run normally"
);
let history = crate::session::store().load(agent_id).await;
let last_user = history
.iter()
.rev()
.find(|m| m.role == crate::ChatRole::User)
.expect("user message exists");
assert!(
last_user.content.contains("[IMAGE:/tmp/photo.png]"),
"image untouched on a text-content rejection"
);
}
#[tokio::test]
#[serial_test::serial(active_models, provider, drain)] async fn subsequent_failure_after_strip_takes_normal_failure_path() {
crate::util::test::init_test_stores().await;
seed_empty_catalogs();
let _policy_guard = install_test_retry_policy(crate::retry::tiny_test_policy());
let agent_id = "e2e_assistant_second_failure";
let fake = std::sync::Arc::new(
FakeProvider::new()
.err_http(400, IMAGE_REJECTION_BODY)
.err_http(400, IMAGE_REJECTION_BODY),
);
let provider: std::sync::Arc<dyn crate::Provider> = fake.clone();
let _provider_guard = install_fake_provider(provider);
let mut agent = make_agent_with_role(vec![], crate::Role::Assistant);
agent.agent_id = agent_id.to_string();
let result = agent
.work("[IMAGE:/tmp/photo.png] describe this photo", false)
.await;
assert!(
result.is_err(),
"a second failure after the strip must fail normally"
);
let messages = fake.request_messages.lock().unwrap().clone();
assert_eq!(
messages.len(),
2,
"exactly one strip, then the normal failure"
);
let retried_user = messages[1]
.split('\u{0}')
.next_back()
.expect("retried request has a user segment");
assert!(
!retried_user.contains("[IMAGE:"),
"retried user message is already stripped"
);
let history = crate::session::store().load(agent_id).await;
let last_user = history
.iter()
.rev()
.find(|m| m.role == crate::ChatRole::User)
.expect("user message exists");
assert!(!last_user.content.contains("[IMAGE:"));
assert!(
last_user
.content
.contains("rejected by the provider's content-inspection check"),
"phrase present after a single strip"
);
}
#[test]
fn existing_image_marker_values_extracts_image_markers() {
let history = vec![
ChatMessage::system("sys"),
ChatMessage::user("[IMAGE:data:image/jpeg;base64,aaa] see this"),
ChatMessage::assistant("ok"),
ChatMessage::user("plain text, no marker"),
ChatMessage::user("[AUDIO:data:audio/ogg;base64,zzz]"),
];
let set = existing_image_marker_values(&history);
assert_eq!(
set.len(),
1,
"only IMAGE markers are collected, got: {set:?}"
);
assert!(set.contains("data:image/jpeg;base64,aaa"));
}
#[test]
fn session_image_marker_count_prefix_gate() {
let png = crate::util::test::tiny_png_data_uri();
let history = vec![
ChatMessage::system("sys"),
ChatMessage::user(format!("[IMAGE:{png}] first")),
ChatMessage::user(format!("[IMAGE:{png}] duplicate — same URI")),
ChatMessage::user("[IMAGE:/local/path.png] file path"),
ChatMessage::user("[AUDIO:data:audio/ogg;base64,zzz]"),
ChatMessage::assistant(format!("[IMAGE:{png}] assistant role must not count")),
];
assert_eq!(
session_image_marker_count(&history),
1,
"only the distinct native data-URI IMAGE marker in User messages counts"
);
}
#[test]
fn summarization_trigger_decides_reason() {
use crate::session::{SUMMARIZATION_IMAGE_COUNT, SUMMARIZATION_THRESHOLD};
let nine = SUMMARIZATION_IMAGE_COUNT - 1;
assert_eq!(summarization_trigger(None, nine), None);
assert_eq!(
summarization_trigger(None, SUMMARIZATION_IMAGE_COUNT),
Some("image count")
);
assert_eq!(
summarization_trigger(Some(SUMMARIZATION_THRESHOLD), 0),
None
);
assert_eq!(
summarization_trigger(Some(SUMMARIZATION_THRESHOLD + 1), 0),
Some("token threshold")
);
assert_eq!(
summarization_trigger(Some(SUMMARIZATION_THRESHOLD + 1), SUMMARIZATION_IMAGE_COUNT,),
Some("token threshold")
);
assert_eq!(summarization_trigger(Some(1), nine), None);
assert_eq!(
summarization_trigger(Some(1), SUMMARIZATION_IMAGE_COUNT),
Some("image count")
);
}
#[tokio::test]
#[expect(clippy::too_many_lines)]
async fn commit_tool_results_injects_image_after_tool_results_with_dedup() {
crate::util::test::init_test_stores().await;
let mut agent = make_agent(vec![]);
agent.session.push_messages_unpersisted(&[ChatMessage::user(
"[IMAGE:data:image/jpeg;base64,prior]",
)]);
agent
.session
.push_assistant("assistant with tool_calls".to_string());
let tool_calls = vec![
ToolCall {
id: "call1".into(),
name: "read".into(),
arguments: serde_json::json!({"path": "a.png"}),
},
ToolCall {
id: "call2".into(),
name: "read".into(),
arguments: serde_json::json!({"path": "b.png"}),
},
];
let outcomes = vec![
ToolExecutionOutcome {
output: "Read image file /a.png (1x1, PNG).".into(),
success: true,
image_payloads: vec![ImagePayload {
path: "/a.png".into(),
data_uri: "data:image/jpeg;base64,prior".into(),
width: 1,
height: 1,
format: "PNG".into(),
recovery_note: None,
source: crate::tools::ImagePayloadSource::Read,
}],
text_is_content: false,
suspended: false,
ends_turn: false,
},
ToolExecutionOutcome {
output: "Read image file /b.png (1x1, PNG).".into(),
success: true,
image_payloads: vec![ImagePayload {
path: "/b.png".into(),
data_uri: "data:image/jpeg;base64,fresh".into(),
width: 1,
height: 1,
format: "PNG".into(),
recovery_note: None,
source: crate::tools::ImagePayloadSource::Read,
}],
text_is_content: false,
suspended: false,
ends_turn: false,
},
];
agent
.commit_tool_results(&tool_calls, outcomes)
.await
.expect("commit_tool_results must succeed");
let history = agent.session.history();
let roles: Vec<crate::ChatRole> = history.iter().map(|m| m.role).collect();
assert_eq!(
roles,
vec![
crate::ChatRole::User,
crate::ChatRole::Assistant,
crate::ChatRole::Tool,
crate::ChatRole::Tool,
crate::ChatRole::User,
],
"unexpected message ordering: {roles:?}"
);
let image_users: Vec<&str> = history
.iter()
.filter(|m| m.role == crate::ChatRole::User && m.content.contains("[IMAGE:data:"))
.map(|m| m.content.as_str())
.collect();
assert_eq!(
image_users.len(),
2,
"prior + fresh image messages, got: {image_users:?}"
);
assert!(
image_users.contains(&"[IMAGE:data:image/jpeg;base64,prior]"),
"prior image survives in history: {image_users:?}"
);
assert!(
image_users
.iter()
.any(|s| s.contains("[IMAGE:data:image/jpeg;base64,fresh]")),
"fresh image injected once: {image_users:?}"
);
let fresh_msg = image_users
.iter()
.find(|s| s.contains("base64,fresh"))
.expect("fresh image message present");
assert!(
fresh_msg.starts_with("<injected-tool-result-image>\n[IMAGE:data:image/jpeg;base64,"),
"injected image carries the provenance tag: {fresh_msg:?}"
);
let fresh_count = history
.iter()
.filter(|m| m.role == crate::ChatRole::User && m.content.contains("base64,fresh"))
.count();
assert_eq!(fresh_count, 1, "a fresh image is injected exactly once");
let tool_results: Vec<String> = history
.iter()
.filter(|m| m.role == crate::ChatRole::Tool)
.map(|m| {
let v: serde_json::Value =
serde_json::from_str(&m.content).expect("tool result is valid JSON");
crate::util::json::get_str(&v, "content")
.expect("tool result has content")
.to_string()
})
.collect();
assert_eq!(
tool_results,
vec![
"Read image file /a.png (1x1, PNG). Image content is already attached to the conversation as a native image.",
"Read image file /b.png (1x1, PNG). Image content attached to the conversation as a native image.",
],
"unexpected tool-result text: {tool_results:?}"
);
}
#[tokio::test]
async fn commit_tool_results_dedups_identical_reads_in_same_round() {
crate::util::test::init_test_stores().await;
let mut agent = make_agent(vec![]);
let tool_calls = vec![
ToolCall {
id: "callA".into(),
name: "read".into(),
arguments: serde_json::json!({"path": "a.png"}),
},
ToolCall {
id: "callB".into(),
name: "read".into(),
arguments: serde_json::json!({"path": "a.png"}),
},
];
let outcomes = vec![
ToolExecutionOutcome {
output: "Read image file /a.png (PNG).".into(),
success: true,
image_payloads: vec![ImagePayload {
path: "/a.png".into(),
data_uri: "data:image/jpeg;base64,same".into(),
width: 1,
height: 1,
format: "PNG".into(),
recovery_note: None,
source: crate::tools::ImagePayloadSource::Read,
}],
text_is_content: false,
suspended: false,
ends_turn: false,
},
ToolExecutionOutcome {
output: "Read image file /a.png (PNG).".into(),
success: true,
image_payloads: vec![ImagePayload {
path: "/a.png".into(),
data_uri: "data:image/jpeg;base64,same".into(),
width: 1,
height: 1,
format: "PNG".into(),
recovery_note: None,
source: crate::tools::ImagePayloadSource::Read,
}],
text_is_content: false,
suspended: false,
ends_turn: false,
},
];
agent
.commit_tool_results(&tool_calls, outcomes)
.await
.expect("commit_tool_results must succeed");
let history = agent.session.history();
let image_users: Vec<&str> = history
.iter()
.filter(|m| {
m.role == crate::ChatRole::User
&& m.content.contains("[IMAGE:data:image/jpeg;base64,same]")
})
.map(|m| m.content.as_str())
.collect();
assert_eq!(
image_users.len(),
1,
"identical reads in one round inject exactly once, got: {image_users:?}"
);
let tool_results: Vec<String> = history
.iter()
.filter(|m| m.role == crate::ChatRole::Tool)
.map(|m| {
let v: serde_json::Value =
serde_json::from_str(&m.content).expect("tool result is valid JSON");
crate::util::json::get_str(&v, "content")
.expect("tool result has content")
.to_string()
})
.collect();
assert_eq!(tool_results.len(), 2, "two tool results persisted");
assert!(
tool_results[0].contains("attached to the conversation as a native image")
&& !tool_results[0].contains("already"),
"first read claims a fresh attachment: {tool_results:?}"
);
assert!(
tool_results[1].contains("already attached to the conversation"),
"second read returns a reference: {tool_results:?}"
);
}
#[tokio::test]
async fn commit_tool_results_injects_one_image_message_per_round() {
crate::util::test::init_test_stores().await;
let mut agent = make_agent(vec![]);
let tool_calls = vec![
ToolCall {
id: "callA".into(),
name: "read".into(),
arguments: serde_json::json!({"path": "a.pdf"}),
},
ToolCall {
id: "callB".into(),
name: "read".into(),
arguments: serde_json::json!({"path": "b.pdf"}),
},
];
let payload = |suffix: &str| ImagePayload {
path: format!("/tmp/{suffix}.jpg"),
data_uri: format!("data:image/jpeg;base64,{suffix}"),
width: 1,
height: 1,
format: "JPEG".into(),
recovery_note: None,
source: crate::tools::ImagePayloadSource::Read,
};
let outcomes = vec![
ToolExecutionOutcome {
output: "Read image file /tmp/a.jpg.".into(),
success: true,
image_payloads: vec![payload("a1"), payload("a2")],
text_is_content: false,
suspended: false,
ends_turn: false,
},
ToolExecutionOutcome {
output: "Read image file /tmp/b.jpg.".into(),
success: true,
image_payloads: vec![payload("b1")],
text_is_content: false,
suspended: false,
ends_turn: false,
},
];
agent
.commit_tool_results(&tool_calls, outcomes)
.await
.expect("commit_tool_results must succeed");
let history = agent.session.history();
let image_messages: Vec<&str> = history
.iter()
.filter(|m| m.role == crate::ChatRole::User && m.content.contains("[IMAGE:"))
.map(|m| m.content.as_str())
.collect();
assert_eq!(
image_messages.len(),
1,
"one message for the whole round, got: {image_messages:?}"
);
let message = image_messages[0];
assert!(
message.starts_with("<injected-tool-result-image>"),
"the batch keeps the provenance tag: {message:?}"
);
for suffix in ["a1", "a2", "b1"] {
assert!(
message.contains(&format!("[IMAGE:data:image/jpeg;base64,{suffix}]")),
"every fresh image of the round is in the batch: {message:?}"
);
}
}
#[tokio::test]
async fn commit_tool_results_leaves_a_marker_in_content_text_alone() {
crate::util::test::init_test_stores().await;
let dir = tempfile::TempDir::new().expect("tempdir");
let png_path = dir.path().join("page.png");
std::fs::write(&png_path, crate::util::test::noisy_png(4, 4)).expect("write png");
let abs = std::fs::canonicalize(&png_path).expect("canonicalize");
let mut agent = make_agent(vec![]);
let marker = format!("[IMAGE:{}]", abs.display());
let body = format!("[scan.pdf: extracted text follows]\n\nsee {marker} in the report");
let tool_calls = vec![ToolCall {
id: "call1".into(),
name: "read".into(),
arguments: serde_json::json!({"path": "scan.pdf"}),
}];
let outcomes = vec![ToolExecutionOutcome {
output: body.clone(),
success: true,
image_payloads: Vec::new(),
text_is_content: true,
suspended: false,
ends_turn: false,
}];
agent
.commit_tool_results(&tool_calls, outcomes)
.await
.expect("commit_tool_results must succeed");
let history = agent.session.history();
let tool_result = history
.iter()
.find(|m| m.role == crate::ChatRole::Tool)
.map(|m| {
let v: serde_json::Value =
serde_json::from_str(&m.content).expect("tool result is valid JSON");
crate::util::json::get_str(&v, "content")
.expect("tool result has content")
.to_string()
})
.expect("the tool result is persisted");
assert!(
tool_result.contains(&marker),
"a marker inside content text stays literal: {tool_result:?}"
);
assert!(
!history
.iter()
.any(|m| m.role == crate::ChatRole::User && m.content.contains("[IMAGE:")),
"nothing is derived from a marker in content text: {history:?}"
);
}
#[tokio::test]
async fn commit_tool_results_appends_annotations_for_content_text() {
crate::util::test::init_test_stores().await;
let mut agent = make_agent(vec![]);
let tool_calls = vec![ToolCall {
id: "call1".into(),
name: "read".into(),
arguments: serde_json::json!({"path": "scan.pdf"}),
}];
let outcomes = vec![ToolExecutionOutcome {
output: "[scan.pdf: no text could be extracted — the pages were provided as images]"
.into(),
success: true,
image_payloads: vec![ImagePayload {
path: "/tmp/read_1/page_1.jpg".into(),
data_uri: "data:image/jpeg;base64,page1".into(),
width: 800,
height: 1000,
format: "JPEG".into(),
recovery_note: None,
source: crate::tools::ImagePayloadSource::Generated,
}],
text_is_content: true,
suspended: false,
ends_turn: false,
}];
agent
.commit_tool_results(&tool_calls, outcomes)
.await
.expect("commit_tool_results must succeed");
let history = agent.session.history();
let tool_result = history
.iter()
.find(|m| m.role == crate::ChatRole::Tool)
.map(|m| {
let v: serde_json::Value =
serde_json::from_str(&m.content).expect("tool result is valid JSON");
crate::util::json::get_str(&v, "content")
.expect("tool result has content")
.to_string()
})
.expect("the tool result is persisted");
assert_eq!(
tool_result,
"[scan.pdf: no text could be extracted — the pages were provided as images]\n\
Generated image file /tmp/read_1/page_1.jpg (800x1000, JPEG). Image content attached \
to the conversation as a native image.",
"the answer keeps its text and gains the image annotation: {tool_result:?}"
);
assert!(
history.iter().any(|m| m.role == crate::ChatRole::User
&& m.content.contains("[IMAGE:data:image/jpeg;base64,page1]")),
"the page is injected as a synthetic user message"
);
}
#[tokio::test]
async fn commit_tool_results_derives_generated_image_and_tags_it() {
crate::util::test::init_test_stores().await;
let dir = tempfile::TempDir::new().expect("tempdir");
let png_path = dir.path().join("generated.png");
std::fs::write(&png_path, crate::util::test::noisy_png(4, 4)).expect("write png");
let abs = std::fs::canonicalize(&png_path).expect("canonicalize");
let tool = Box::new(MediaTestTool {
name: "image_gen",
marker: "[IMAGE:",
}) as Box<dyn Tool>;
let mut agent = make_agent(vec![tool]);
let marker = format!("[IMAGE:{}]", abs.display());
let tool_calls = vec![ToolCall {
id: "callgen".into(),
name: "image_gen".into(),
arguments: serde_json::json!({}),
}];
let outcomes = vec![ToolExecutionOutcome {
output: marker.clone(),
success: true,
image_payloads: Vec::new(),
text_is_content: false,
suspended: false,
ends_turn: false,
}];
let media = extract_media_from_outcomes(&agent.tools, &tool_calls, &outcomes);
assert_eq!(
media,
vec![("[IMAGE:", abs.display().to_string())],
"marker preserved for user delivery: {media:?}"
);
agent
.commit_tool_results(&tool_calls, outcomes)
.await
.expect("commit_tool_results must succeed");
let history = agent.session.history();
let image_users: Vec<&str> = history
.iter()
.filter(|m| {
m.role == crate::ChatRole::User
&& m.content.contains("[IMAGE:data:image/jpeg;base64,")
})
.map(|m| m.content.as_str())
.collect();
assert_eq!(
image_users.len(),
1,
"generated image injected once: {image_users:?}"
);
assert!(
image_users[0]
.starts_with("<injected-tool-result-image>\n[IMAGE:data:image/jpeg;base64,"),
"injected image carries the provenance tag: {image_users:?}"
);
let tool_results: Vec<String> = history
.iter()
.filter(|m| m.role == crate::ChatRole::Tool)
.map(|m| {
let v: serde_json::Value =
serde_json::from_str(&m.content).expect("tool result JSON");
crate::util::json::get_str(&v, "content")
.expect("content")
.to_string()
})
.collect();
assert_eq!(tool_results.len(), 1);
assert!(
tool_results[0].starts_with("Generated image file"),
"generated annotation: {tool_results:?}"
);
assert!(
!tool_results[0].starts_with("Read image file"),
"must not be a Read annotation: {tool_results:?}"
);
}
#[tokio::test]
async fn complete_pending_tool_calls_injects_one_image_message_per_round() {
crate::util::test::init_test_stores().await;
let dir = tempfile::tempdir().unwrap();
for (page, size) in [("page-a.png", 4), ("page-b.png", 5)] {
std::fs::write(
dir.path().join(page),
crate::util::test::noisy_png(size, size),
)
.unwrap();
}
let ws = crate::Workspace::from_path(dir.path());
let calls: Vec<crate::ToolCall> = ["page-a.png", "page-b.png"]
.iter()
.enumerate()
.map(|(index, page)| crate::ToolCall {
id: format!("call_img_{index}"),
name: "read".to_string(),
arguments: serde_json::json!({ "path": page }),
})
.collect();
let frame = crate::providers::reasoning::assistant_replay_payload(Some(""), &calls, None)
.to_string();
let mut session = Session::default();
session
.persist_messages(
"test-resume-image-batch",
&[
ChatMessage::user("read both images"),
ChatMessage::assistant(frame),
],
)
.await
.unwrap();
let mut agent = make_agent_on(
vec![Box::new(crate::tools::ReadTool::general())],
"test-resume-image-batch",
ws,
);
agent.session = session;
assert!(
agent.complete_pending_tool_calls().await.unwrap(),
"both dangling calls settle a result"
);
let history = agent.session.history();
let image_messages: Vec<&str> = history
.iter()
.filter(|m| m.role == crate::ChatRole::User && m.content.contains("[IMAGE:"))
.map(|m| m.content.as_str())
.collect();
assert_eq!(
image_messages.len(),
1,
"one message for the whole re-executed round, got: {image_messages:?}"
);
assert_eq!(
image_messages[0]
.matches("[IMAGE:data:image/jpeg;base64,")
.count(),
2,
"both pages of the round are in the batch: {:?}",
image_messages[0]
);
}
#[tokio::test]
async fn complete_pending_tool_calls_reexecutes_non_durable_call() {
crate::util::test::init_test_stores().await;
let dir = tempfile::tempdir().unwrap();
std::fs::write(dir.path().join("seed.txt"), "completion-marker\n").unwrap();
let ws = crate::Workspace::from_path(dir.path());
let mut session = Session::default();
let frame = crate::providers::reasoning::assistant_replay_payload(
Some(""),
&[crate::ToolCall {
id: "call_read_c".to_string(),
name: "read".to_string(),
arguments: serde_json::json!({"path": "seed.txt"}),
}],
None,
)
.to_string();
session
.persist_messages(
"test-read-reexec",
&[
ChatMessage::user("read the file"),
ChatMessage::assistant(frame),
],
)
.await
.unwrap();
let mut agent = make_agent_on(
vec![Box::new(crate::tools::ReadTool::general())],
"test-read-reexec",
ws,
);
agent.session = session;
let settled = agent.complete_pending_tool_calls().await.unwrap();
assert!(settled, "a re-executed non-durable call settles a result");
let history = agent.session.history();
let roles: Vec<crate::ChatRole> = history.iter().map(|m| m.role).collect();
assert_eq!(
roles,
vec![
crate::ChatRole::User,
crate::ChatRole::Assistant,
crate::ChatRole::Tool
],
"result contiguous after the frame: {roles:?}"
);
let payload: crate::ToolResultPayload = serde_json::from_str(&history[2].content).unwrap();
assert_eq!(
payload.tool_call_id, "call_read_c",
"re-execution keeps the ORIGINAL call id"
);
assert!(
payload.content.contains("completion-marker"),
"real tool output recorded: {}",
payload.content
);
assert!(agent.session.pending_tool_frame().is_none());
}
#[tokio::test]
#[serial_test::serial(provider, drain)] async fn drain_mid_sync_analyze_leaves_call_dangling_then_completion_resumes() {
let fake = std::sync::Arc::new(crate::util::test::FakeProvider::new());
let _seam = crate::util::test::install_retry_seam_dyn(fake.clone());
crate::util::test::init_test_stores().await;
let ws = crate::workspace::test_ws("/tmp/test_ws_drain_analyze");
let pin = "drain_sync_analyze_pin";
let conn = &crate::session::store().conn;
let mut session = Session::default();
let frame = crate::providers::reasoning::assistant_replay_payload(
Some(""),
&[crate::ToolCall {
id: "call_analyze_d".to_string(),
name: "analyze".to_string(),
arguments: serde_json::json!({"analyze": "analyze this"}),
}],
None,
)
.to_string();
session
.persist_messages(
pin,
&[
ChatMessage::user("analyze this"),
ChatMessage::assistant(frame),
],
)
.await
.unwrap();
crate::shutdown::drain_begin();
let tool = crate::tools::analyze::AnalyzeTool::new(
crate::tools::analyze::DispatchMode::Sync,
crate::Role::Engineer,
);
let res = crate::agent::CURRENT_TOOL_AGENT_ID
.scope(Some(pin.to_string()), async {
tool.execute(&ws, serde_json::json!({"analyze": "analyze this"}))
.await
})
.await;
crate::shutdown::drain_clear();
let err = res.expect_err("drain must cut the sync analyze dispatch");
assert!(
err.downcast_ref::<crate::tools::CallSuspended>().is_some(),
"CallSuspended carrier expected: {err:#}"
);
let jobs = conn
.query(
"SELECT id, status, caller_agent_id FROM jobs WHERE caller_agent_id = ?1 AND kind = 'analyze'",
crate::db::params![pin],
)
.await
.unwrap();
assert_eq!(jobs.len(), 1, "one launched analyze job");
assert_eq!(jobs[0].get::<String>(1).unwrap(), "launched");
assert_eq!(
jobs[0].get::<String>(2).unwrap(),
pin,
"job is caller-owned by the session pin"
);
let job_id = jobs[0].get::<String>(0).unwrap();
let session_rows = conn
.query(
"SELECT role FROM sessions WHERE agent_id = ?1 ORDER BY id",
crate::db::params![pin],
)
.await
.unwrap();
let roles: Vec<String> = session_rows
.iter()
.map(|r| r.get::<String>(0).unwrap())
.collect();
assert_eq!(
roles,
vec!["user", "assistant"],
"frame persisted with NO tool-result row: {roles:?}"
);
let roster = crate::jobs::list_agents_for_job(conn, &job_id)
.await
.unwrap();
assert!(
roster.len() >= 2,
"analyze round spawns multiple analysts: {}",
roster.len()
);
for (i, row) in roster.iter().enumerate() {
let outcome = if i == 0 { "ANALYST_RAW" } else { "" };
crate::jobs::write_agent_outcome(
conn,
&job_id,
&row.agent_id,
crate::jobs::RowStatus::Done,
Some(outcome),
)
.await
.unwrap();
}
let mut agent = make_agent_on(vec![], pin, ws);
agent.session = session;
let settled = agent.complete_pending_tool_calls().await.unwrap();
assert!(settled, "resumed durable job settles a result");
let jobs = conn
.query(
"SELECT id FROM jobs WHERE id = ?1",
crate::db::params![job_id.clone()],
)
.await
.unwrap();
assert!(jobs.is_empty(), "resumed job must be terminalized");
assert!(
fake.request_fingerprints.lock().unwrap().is_empty(),
"replaying a Done slot must not call the model"
);
let history = agent.session.history();
let roles: Vec<crate::ChatRole> = history.iter().map(|m| m.role).collect();
assert_eq!(
roles,
vec![
crate::ChatRole::User,
crate::ChatRole::Assistant,
crate::ChatRole::Tool
],
"tool result after the frame: {roles:?}"
);
let payload: crate::ToolResultPayload = serde_json::from_str(&history[2].content).unwrap();
assert_eq!(payload.tool_call_id, "call_analyze_d");
assert!(
payload.content.contains("ANALYST_RAW"),
"{}",
payload.content
);
let pending_jobs = conn
.query(
"SELECT id FROM pending_jobs WHERE id = ?1",
crate::db::params![job_id],
)
.await
.unwrap();
assert!(pending_jobs.is_empty(), "no envelope pending after resume");
}
#[tokio::test]
#[serial_test::serial(provider, drain)]
async fn complete_pending_tool_calls_binds_same_kind_calls_to_distinct_jobs() {
let fake = std::sync::Arc::new(crate::util::test::FakeProvider::new());
let _seam = crate::util::test::install_retry_seam_dyn(fake.clone());
crate::util::test::init_test_stores().await;
let ws = crate::workspace::test_ws("/tmp/test_ws_two_analyze_jobs");
let pin = "two_analyze_jobs_pin";
let conn = &crate::session::store().conn;
let mut session = Session::default();
let frame = assistant_replay_payload(
Some(""),
&[
crate::ToolCall {
id: "call_analyze_g1".to_string(),
name: "analyze".to_string(),
arguments: serde_json::json!({"analyze": "analyze task one"}),
},
crate::ToolCall {
id: "call_analyze_g2".to_string(),
name: "analyze".to_string(),
arguments: serde_json::json!({"analyze": "analyze task two"}),
},
],
None,
)
.to_string();
session
.persist_messages(
pin,
&[
ChatMessage::user("analyze this"),
ChatMessage::assistant(frame),
],
)
.await
.unwrap();
let jobs: [(&str, &str); 3] = [
("two_analyze_job_0", "RESULT_ONE"),
("two_analyze_job_1", "RESULT_TWO"),
("two_analyze_job_2", "RESULT_THREE"),
];
for (i, (job_id, outcome)) in jobs.iter().copied().enumerate() {
let task = format!("analyze task {i}");
let analyst_id = format!("{job_id}_analyst");
crate::jobs::spawn_job(
conn,
job_id,
&task,
&ws.name,
"caller-user",
"telegram",
crate::Role::Engineer,
&[crate::jobs::NewAgent {
agent_id: analyst_id.clone(),
kind: crate::jobs::AgentKind::Analyst,
idx: Some(0),
task: task.clone(),
}],
&crate::jobs::SpawnChild::Analyze,
Some(pin),
)
.await
.unwrap();
crate::jobs::write_agent_outcome(
conn,
job_id,
&analyst_id,
crate::jobs::RowStatus::Done,
Some(outcome),
)
.await
.unwrap();
let created = format!("2024-01-0{}T00:00:00+00:00", i + 2);
conn.execute(
"UPDATE jobs SET created_at = ?1 WHERE id = ?2",
crate::db::params![created, job_id],
)
.await
.unwrap();
}
let mut agent = make_agent_on(vec![], pin, ws);
agent.session = session;
let settled = agent.complete_pending_tool_calls().await.unwrap();
assert!(settled, "two resumed jobs settle results");
assert!(
fake.request_fingerprints.lock().unwrap().is_empty(),
"replaying Done slots must not call the model"
);
assert!(agent.session.pending_tool_frame().is_none());
let jobs = conn
.query(
"SELECT id FROM jobs WHERE caller_agent_id = ?1 AND kind = 'analyze' AND status = 'launched'",
crate::db::params![pin],
)
.await
.unwrap();
assert_eq!(jobs.len(), 1, "only the unbound job remains launched");
assert_eq!(jobs[0].get::<String>(0).unwrap(), "two_analyze_job_2");
let history = agent.session.history();
let mut per_call: std::collections::HashMap<String, String> =
std::collections::HashMap::new();
for msg in history {
if msg.role != crate::ChatRole::Tool {
continue;
}
let payload: crate::ToolResultPayload = serde_json::from_str(&msg.content).unwrap();
per_call.insert(payload.tool_call_id, payload.content);
}
assert_eq!(per_call.len(), 2, "both calls settled");
assert_eq!(
per_call.get("call_analyze_g1").map(String::as_str),
Some("RESULT_ONE"),
"first frame call binds the OLDEST job"
);
assert_eq!(
per_call.get("call_analyze_g2").map(String::as_str),
Some("RESULT_TWO"),
"second frame call binds the newer job"
);
}
async fn seed_run_owned_state(
agent_id: &str,
) -> (
tempfile::TempDir,
std::path::PathBuf,
std::sync::Arc<crate::tools::chrome::ChromeRunSessions>,
) {
let dir = tempfile::tempdir().unwrap();
let spill = dir.path().join("spill.txt");
std::fs::write(&spill, "spilled output").unwrap();
CURRENT_TOOL_AGENT_ID
.scope(Some(agent_id.to_string()), async {
crate::tools::shell::record_spill_owner(spill.clone());
})
.await;
let sessions = crate::tools::chrome::ChromeRunSessions::for_run(agent_id);
sessions.track(&format!("{}default", sessions.namespace()));
(dir, spill, sessions)
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn run_end_guard_cleans_up_a_dropped_run() {
let agent_id = format!("run-end-guard-{:016x}", rand::random::<u64>());
let (_dir, spill, sessions) = seed_run_owned_state(&agent_id).await;
chrome_release::clear_pending_releases();
let _no_store = chrome_release::no_release_store();
drop(RunEndCleanup::new(
agent_id,
std::sync::Arc::clone(&sessions),
));
assert!(!spill.exists(), "the run's spill file is reclaimed");
assert_eq!(
chrome_release::pending_names_and_attempts(),
vec![(vec![format!("{}default", sessions.namespace())], 0)],
"the ended run handed its session over for release"
);
chrome_release::clear_pending_releases();
}
#[tokio::test]
#[serial_test::serial(chrome_release)]
async fn an_ended_run_hands_its_sessions_over_and_holds_only_a_resumed_run() {
use crate::tools::chrome_release::RELEASE_HOLD_REDRIVEN_RUN;
chrome_release::clear_pending_releases();
let _no_store = chrome_release::no_release_store();
let cut_id = format!("run-end-keep-{:016x}", rand::random::<u64>());
let (_cut_dir, _cut_spill, cut_sessions) = seed_run_owned_state(&cut_id).await;
let mut cut = RunEndCleanup::new(cut_id, std::sync::Arc::clone(&cut_sessions));
cut.ended("drain");
drop(cut);
assert!(
chrome_release::pending_release_delays()[0]
>= RELEASE_HOLD_REDRIVEN_RUN.saturating_sub(std::time::Duration::from_secs(5)),
"a run the daemon cut off holds its release back for its resumed segment"
);
chrome_release::clear_pending_releases();
let pause_id = format!("run-end-pause-{:016x}", rand::random::<u64>());
let (_pause_dir, _pause_spill, paused_sessions) = seed_run_owned_state(&pause_id).await;
let mut paused = RunEndCleanup::new(pause_id, std::sync::Arc::clone(&paused_sessions));
paused.ended("pause");
drop(paused);
assert!(
chrome_release::pending_release_delays()[0]
>= RELEASE_HOLD_REDRIVEN_RUN.saturating_sub(std::time::Duration::from_secs(5)),
"a run frozen by a workspace pause resumes at unpause under its own id"
);
chrome_release::clear_pending_releases();
let final_id = format!("run-end-release-{:016x}", rand::random::<u64>());
let (_dir, _spill, sessions) = seed_run_owned_state(&final_id).await;
let mut ended = RunEndCleanup::new(final_id, std::sync::Arc::clone(&sessions));
ended.ended("transport");
drop(ended);
assert!(
chrome_release::pending_release_delays()[0] < std::time::Duration::from_secs(1),
"an end that is final releases at once"
);
chrome_release::clear_pending_releases();
}
}