use zeph_common::text::estimate_tokens;
use zeph_llm::provider::{Message, MessagePart, Role};
use crate::agent::Agent;
use crate::channel::Channel;
impl<C: Channel> Agent<C> {
pub(crate) fn handle_focus_tool(
&mut self,
tool_name: &str,
input: &serde_json::Value,
) -> (String, Option<zeph_llm::provider::Message>) {
match tool_name {
"start_focus" => self.start_focus_tool(input),
"complete_focus" => self.complete_focus_tool(input),
other => (format!("[error] Unknown focus tool: {other}"), None),
}
}
fn start_focus_tool(
&mut self,
input: &serde_json::Value,
) -> (String, Option<zeph_llm::provider::Message>) {
let scope = input
.get("scope")
.and_then(|v| v.as_str())
.unwrap_or("(unspecified)")
.to_string();
if self.services.focus.is_active() {
return (
"[error] A focus session is already active. Call complete_focus first.".to_string(),
None,
);
}
let marker = self.services.focus.start(scope.clone());
let checkpoint_msg = zeph_llm::provider::Message {
role: zeph_llm::provider::Role::System,
content: format!("[focus checkpoint: {scope}]"),
parts: vec![],
metadata: zeph_llm::provider::MessageMetadata {
focus_pinned: true,
focus_marker_id: Some(marker),
..zeph_llm::provider::MessageMetadata::agent_only()
},
};
(
format!("Focus session started. Checkpoint ID: {marker}. Scope: {scope}"),
Some(checkpoint_msg),
)
}
fn complete_focus_tool(
&mut self,
input: &serde_json::Value,
) -> (String, Option<zeph_llm::provider::Message>) {
let summary = input
.get("summary")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
if !self.services.focus.is_active() {
return (
"[error] No active focus session. Call start_focus first.".to_string(),
None,
);
}
let Some(marker) = self.services.focus.active_marker else {
return (
"[error] Internal error: active_marker is None.".to_string(),
None,
);
};
let checkpoint_pos = self
.msg
.messages
.iter()
.position(|m| m.metadata.focus_marker_id == Some(marker));
let Some(checkpoint_pos) = checkpoint_pos else {
return (
format!(
"[error] Checkpoint marker {marker} not found in message history. \
The focus session may have been evicted by compaction."
),
None,
);
};
let _ = self.msg.messages[checkpoint_pos + 1..].to_vec();
let sanitized_summary = self
.services
.security
.sanitizer
.sanitize(
&summary,
zeph_sanitizer::ContentSource::new(zeph_sanitizer::ContentSourceKind::WebScrape),
)
.body;
self.services
.focus
.append_llm_knowledge(sanitized_summary.clone());
if let Some(ref d) = self.runtime.debug.debug_dumper {
let kb = self
.services
.focus
.knowledge_blocks
.iter()
.map(|b| b.content.as_str())
.collect::<Vec<_>>()
.join("\n---\n");
d.dump_focus_knowledge(&kb);
}
self.services.focus.complete();
let current_turn_assistant = {
let last_idx = self.msg.messages.len().saturating_sub(1);
if last_idx >= checkpoint_pos {
self.msg.messages.last().and_then(|m| {
if m.role == Role::Assistant
&& m.parts
.iter()
.any(|p| matches!(p, MessagePart::ToolUse { .. }))
{
Some(m.clone())
} else {
None
}
})
} else {
None
}
};
self.msg.messages.truncate(checkpoint_pos);
if let Some(assistant_msg) = current_turn_assistant {
self.msg.messages.push(assistant_msg);
}
self.msg.recompute_non_system_count();
self.recompute_prompt_tokens();
self.context_manager.set_compaction_state(
crate::agent::context_manager::CompactionState::CompactedThisTurn { cooldown: 0 },
);
self.rebuild_knowledge_block();
(
format!("Focus session complete. Knowledge block updated with: {sanitized_summary}"),
None,
)
}
pub(crate) fn rebuild_knowledge_block(&mut self) {
self.msg
.messages
.retain(|m| !(m.metadata.focus_pinned && m.metadata.focus_marker_id.is_none()));
if let Some(kb_msg) = self.services.focus.build_knowledge_message() {
if self.msg.messages.is_empty() {
self.msg.messages.push(kb_msg);
} else {
self.msg.messages.insert(1, kb_msg);
}
}
self.msg.recompute_non_system_count();
self.recompute_prompt_tokens();
}
#[tracing::instrument(name = "core.tool.handle_compress_context", skip_all, level = "debug")]
pub(crate) async fn handle_compress_context(&mut self) -> String {
use zeph_llm::provider::LlmProvider as _;
if self.services.focus.is_active() {
return "[error] Cannot compress context while a focus session is active. \
Call complete_focus first."
.to_string();
}
if !self.services.focus.try_acquire_compression() {
return "[error] A context compression is already in progress.".to_string();
}
let preserve_tail = self.context_manager.compaction_preserve_tail;
let (to_remove_indices, to_compress) =
match self.select_messages_for_compression(preserve_tail) {
Ok(pair) => pair,
Err(total) => {
self.services.focus.release_compression();
return format!(
"Not enough messages to compress (found {total}, need at least {}).",
preserve_tail + 4
);
}
};
let compress_total = to_compress.len();
let summary_messages = build_compression_prompt(&to_compress);
let compress_provider = self
.runtime
.providers
.compress_provider
.as_ref()
.unwrap_or(&self.provider);
let summary = match tokio::time::timeout(
std::time::Duration::from_secs(30),
compress_provider.chat(&summary_messages),
)
.await
{
Ok(Ok(s)) => s,
Ok(Err(e)) => {
self.services.focus.release_compression();
return format!("[error] Compression LLM call failed: {e}");
}
Err(_) => {
self.services.focus.release_compression();
return "[error] Compression LLM call timed out.".to_string();
}
};
if summary.trim().is_empty() {
self.services.focus.release_compression();
return "[error] Compression produced an empty summary.".to_string();
}
let tokens_freed = to_compress
.iter()
.map(|m| estimate_tokens(&m.content))
.sum::<usize>();
let sanitized_summary = self
.services
.security
.sanitizer
.sanitize(
summary.trim(),
zeph_sanitizer::ContentSource::new(zeph_sanitizer::ContentSourceKind::WebScrape),
)
.body;
self.services.focus.append_llm_knowledge(sanitized_summary);
self.apply_compression_removals(to_remove_indices);
self.context_manager.set_compaction_state(
crate::agent::context_manager::CompactionState::CompactedThisTurn { cooldown: 0 },
);
self.services.focus.release_compression();
format!(
"Compressed {compress_total} messages into a summary (~{tokens_freed} tokens freed). \
Knowledge block updated."
)
}
pub(crate) fn select_messages_for_compression(
&self,
preserve_tail: usize,
) -> Result<
(
std::collections::HashSet<usize>,
Vec<zeph_llm::provider::Message>,
),
usize,
> {
let compressible_indices: Vec<usize> = self
.msg
.messages
.iter()
.enumerate()
.filter(|(_, m)| !m.metadata.focus_pinned && m.role != zeph_llm::provider::Role::System)
.map(|(i, _)| i)
.collect();
let total = compressible_indices.len();
if total <= preserve_tail + 3 {
return Err(total);
}
let to_remove_indices: std::collections::HashSet<usize> = compressible_indices
[..total.saturating_sub(preserve_tail)]
.iter()
.copied()
.collect();
let mut sorted_indices: Vec<usize> = to_remove_indices.iter().copied().collect();
sorted_indices.sort_unstable();
let to_compress: Vec<zeph_llm::provider::Message> = sorted_indices
.iter()
.map(|&i| self.msg.messages[i].clone())
.collect();
Ok((to_remove_indices, to_compress))
}
fn apply_compression_removals(&mut self, to_remove_indices: std::collections::HashSet<usize>) {
let mut remove_idx = to_remove_indices.into_iter().collect::<Vec<_>>();
remove_idx.sort_unstable_by(|a, b| b.cmp(a));
for idx in remove_idx {
if idx < self.msg.messages.len() {
self.msg.messages.remove(idx);
}
}
self.msg.recompute_non_system_count();
self.rebuild_knowledge_block();
}
#[tracing::instrument(
name = "core.tool.persist_cancelled_tool_results",
skip_all,
level = "debug"
)]
pub(crate) async fn persist_cancelled_tool_results(
&mut self,
tool_calls: &[zeph_llm::provider::ToolUseRequest],
insert_at: Option<usize>,
) {
let turn_start = self
.msg
.messages
.iter()
.rposition(|m| m.role == Role::Assistant)
.unwrap_or(0);
let already_resolved: std::collections::HashSet<&str> = self.msg.messages[turn_start..]
.iter()
.flat_map(|m| m.parts.iter())
.filter_map(|p| {
if let MessagePart::ToolResult { tool_use_id, .. } = p {
Some(tool_use_id.as_str())
} else {
None
}
})
.collect();
let result_parts: Vec<MessagePart> = tool_calls
.iter()
.filter(|tc| !already_resolved.contains(tc.id.as_str()))
.map(|tc| MessagePart::ToolResult {
tool_use_id: tc.id.clone(),
content: "[Cancelled]".to_owned(),
is_error: true,
})
.collect();
if result_parts.is_empty() {
return;
}
let user_msg = Message::from_parts(Role::User, result_parts);
self.persist_message(Role::User, &user_msg.content, &user_msg.parts, false)
.await;
match insert_at {
Some(index) => self.insert_message(index, user_msg),
None => self.push_message(user_msg),
}
}
#[tracing::instrument(name = "agent.request_compaction", skip_all, level = "info")]
pub(crate) async fn handle_request_compaction(&mut self, input: &serde_json::Value) -> String {
let raw_reason = input
.get("reason")
.and_then(|v| v.as_str())
.unwrap_or("(no reason provided)");
let reason = &raw_reason[..raw_reason.floor_char_boundary(256.min(raw_reason.len()))];
if self
.context_manager
.compaction_state()
.is_compacted_this_turn()
{
return "[error] Compaction already performed this turn. Try again next turn."
.to_string();
}
let cached = self.runtime.providers.cached_prompt_tokens;
let tier = self.context_manager.compaction_tier(cached);
if matches!(tier, zeph_context::manager::CompactionTier::None) {
return format!(
"Context usage is below the compaction threshold. \
No compaction needed at this time ({cached} tokens cached)."
);
}
tracing::info!(reason, "agent requested compaction (ARC)");
self.handle_compress_context().await
}
}
fn build_compression_prompt(
to_compress: &[zeph_llm::provider::Message],
) -> Vec<zeph_llm::provider::Message> {
let role_label = |role: &zeph_llm::provider::Role| match role {
zeph_llm::provider::Role::Assistant => "assistant",
zeph_llm::provider::Role::System => "system",
zeph_llm::provider::Role::User | _ => "user",
};
let bullet_list: String = to_compress
.iter()
.enumerate()
.map(|(i, m)| {
let truncated: String = m.content.chars().take(500).collect();
let content = repair_truncated_spotlight_wrapper(truncated, m.metadata.trust_level);
format!("{}. [{}] {}", i + 1, role_label(&m.role), content)
})
.collect::<Vec<_>>()
.join("\n");
let total = to_compress.len();
let system_content = "You are a context compression agent. \
Summarize the following conversation messages into a concise, information-dense summary. \
Preserve key facts, decisions, and context. Strip filler and small talk. \
Output ONLY the summary — no headers, no preamble.";
vec![
zeph_llm::provider::Message {
role: zeph_llm::provider::Role::System,
content: system_content.to_owned(),
parts: vec![],
metadata: zeph_llm::provider::MessageMetadata::default(),
},
zeph_llm::provider::Message {
role: zeph_llm::provider::Role::User,
content: format!("Summarize these {total} conversation messages:\n\n{bullet_list}"),
parts: vec![],
metadata: zeph_llm::provider::MessageMetadata::default(),
},
]
}
fn repair_truncated_spotlight_wrapper(mut content: String, trust_level: Option<u8>) -> String {
if matches!(trust_level, None | Some(0)) {
return content;
}
if content.matches("<tool-output").count() > content.matches("</tool-output>").count() {
content.push_str("\n\n[END OF TOOL OUTPUT]\n</tool-output>");
}
if content.matches("<external-data").count() > content.matches("</external-data>").count() {
content.push_str("\n\n[END OF EXTERNAL DATA]\n</external-data>");
}
content
}
#[cfg(test)]
mod tests {
use super::repair_truncated_spotlight_wrapper;
#[test]
fn repairs_severed_tool_output_close() {
let content = "<tool-output source=\"tool_result\" name=\"shell\" trust=\"local\">\
\n[NOTE: ...]\n\nsome truncated shell output"
.to_owned();
let result = repair_truncated_spotlight_wrapper(content, Some(1));
assert!(result.ends_with("</tool-output>"));
assert_eq!(result.matches("<tool-output").count(), 1);
assert_eq!(result.matches("</tool-output>").count(), 1);
}
#[test]
fn repairs_severed_external_data_close() {
let content = "<external-data source=\"web_scrape\" ref=\"http://example.com\" \
trust=\"untrusted\">\n[IMPORTANT: ...]\n\nsome truncated page content"
.to_owned();
let result = repair_truncated_spotlight_wrapper(content, Some(2));
assert!(result.ends_with("</external-data>"));
assert_eq!(result.matches("<external-data").count(), 1);
assert_eq!(result.matches("</external-data>").count(), 1);
}
#[test]
fn repairs_only_the_severed_kind_in_a_mixed_batch() {
let content = "<tool-output source=\"tool_result\" name=\"shell\" trust=\"local\">\
\n\nls output\n\n[END OF TOOL OUTPUT]\n</tool-output>\
<external-data source=\"web_scrape\" ref=\"http://example.com\" \
trust=\"untrusted\">\n[IMPORTANT: ...]\n\ntruncated page"
.to_owned();
let result = repair_truncated_spotlight_wrapper(content, Some(2));
assert_eq!(result.matches("<tool-output").count(), 1);
assert_eq!(result.matches("</tool-output>").count(), 1);
assert_eq!(result.matches("<external-data").count(), 1);
assert_eq!(result.matches("</external-data>").count(), 1);
assert!(result.ends_with("</external-data>"));
}
#[test]
fn trusted_content_is_a_short_circuit_noop_even_with_tag_like_text() {
let content = "<tool-output> looks like a wrapper but isn't, and is unbalanced".to_owned();
assert_eq!(
repair_truncated_spotlight_wrapper(content.clone(), None),
content
);
assert_eq!(
repair_truncated_spotlight_wrapper(content.clone(), Some(0)),
content
);
}
#[test]
fn intact_wrapper_under_500_chars_is_not_double_appended() {
let content = "<tool-output source=\"tool_result\" name=\"shell\" trust=\"local\">\
\n\nshort output\n\n[END OF TOOL OUTPUT]\n</tool-output>"
.to_owned();
let result = repair_truncated_spotlight_wrapper(content.clone(), Some(1));
assert_eq!(
result, content,
"already-balanced wrapper must be left untouched"
);
}
}