#[cfg(test)]
mod tests;
use crate::{
cancellation::{AgentCancellation, AgentCancellationHandle},
config::{EffectiveConfig, Settings, load_effective_provider_selection},
providers::{
ChatMessage, ProviderEvent, ProviderRequest, ProviderSelection,
provider_from_selection_with_settings_for_workload,
},
sessions::{Session, SessionEvent, SessionManager},
};
use serde::{Deserialize, Serialize};
use std::{
collections::BTreeMap,
path::Path,
sync::{Arc, Mutex},
thread,
time::Duration,
};
const MAX_ENTRIES: usize = 200;
const RECENT_ENTRIES: usize = 12;
const MAX_SOURCES: usize = 33;
const WINDOW_EVENTS: usize = 128;
const WINDOW_BYTES: usize = 256 * 1024;
const BATCH_CHARS: usize = 24_000;
const POLL_INTERVAL: Duration = Duration::from_secs(3);
const PROMPT: &str = include_str!("../prompts/summarizer.md");
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct SummarySnapshot {
pub(crate) session_id: String,
pub(crate) enabled: bool,
pub(crate) running: bool,
pub(crate) entries: Vec<String>,
pub(crate) error: Option<String>,
}
struct Command {
generation: u64,
session: Option<Session>,
enabled: Option<bool>,
cancellation: AgentCancellation,
config: Arc<(EffectiveConfig, Settings)>,
}
type Mailbox = Arc<Mutex<Option<(u64, SummarySnapshot)>>>;
pub(crate) struct SummarizerRuntime {
sender: Option<crossbeam_channel::Sender<Command>>,
mailbox: Mailbox,
session: Option<Session>,
generation: u64,
cancellation: Option<AgentCancellationHandle>,
config: Arc<(EffectiveConfig, Settings)>,
}
impl SummarizerRuntime {
pub(crate) fn new(active_config: EffectiveConfig, settings: Settings) -> Self {
let (sender, receiver) = crossbeam_channel::bounded(32);
let mailbox = Arc::new(Mutex::new(None));
let output = Arc::clone(&mailbox);
let config = Arc::new((active_config, settings));
let spawned = thread::Builder::new()
.name("session-summarizer".into())
.spawn(move || {
worker(receiver, output);
});
Self {
sender: spawned.ok().map(|_| sender),
mailbox,
session: None,
generation: 0,
cancellation: None,
config,
}
}
pub(crate) fn refresh_config(&mut self, config: &EffectiveConfig, settings: &Settings) {
if &self.config.0 == config && &self.config.1 == settings {
return;
}
self.config = Arc::new((config.clone(), settings.clone()));
self.submit(self.session.clone(), None);
}
pub(crate) fn set_session(&mut self, session: Option<Session>) {
if self.session == session {
return;
}
self.submit(session, None);
}
pub(crate) fn set_enabled(&mut self, enabled: bool) {
if self.session.is_some() {
self.submit(self.session.clone(), Some(enabled));
}
}
fn submit(&mut self, session: Option<Session>, enabled: Option<bool>) {
let (cancellation, handle) = AgentCancellation::default().child_token();
let generation = self.generation.wrapping_add(1);
let command = Command {
generation,
session: session.clone(),
enabled,
cancellation,
config: Arc::clone(&self.config),
};
if let Some(old) = self.cancellation.take() {
old.cancel();
}
self.generation = generation;
self.session = session;
if self
.sender
.as_ref()
.is_some_and(|sender| sender.try_send(command).is_ok())
{
self.cancellation = Some(handle);
} else if let Some(session) = self.session.as_ref() {
publish(
&self.mailbox,
self.generation,
SummarySnapshot {
session_id: session.id().into(),
enabled: false,
running: false,
entries: vec![],
error: Some(
"Summary worker unavailable or command queue full; try again.".into(),
),
},
);
}
}
pub(crate) fn poll(&mut self) -> Option<SummarySnapshot> {
let (generation, snapshot) = self.mailbox.lock().ok()?.take()?;
(generation == self.generation).then_some(snapshot)
}
pub(crate) fn shutdown(&mut self) {
if let Some(handle) = self.cancellation.take() {
handle.cancel();
}
self.sender.take();
}
}
impl Drop for SummarizerRuntime {
fn drop(&mut self) {
self.shutdown();
}
}
fn publish(mailbox: &Mailbox, generation: u64, snapshot: SummarySnapshot) {
if let Ok(mut slot) = mailbox.lock()
&& slot
.as_ref()
.is_none_or(|(current, _)| *current <= generation)
{
*slot = Some((generation, snapshot));
}
}
#[derive(Clone, Default, Serialize, Deserialize)]
struct Checkpoint {
enabled: bool,
cursors: BTreeMap<String, u64>,
children: Vec<String>,
entries: Vec<String>,
}
struct Active {
command: Command,
log: Session,
checkpoint: Checkpoint,
error: Option<String>,
_lease: crate::persistence::CrossProcessFileLock,
}
impl Active {
fn snapshot(&self, running: bool) -> SummarySnapshot {
SummarySnapshot {
session_id: self.log.id().into(),
enabled: self.checkpoint.enabled,
running,
entries: self.checkpoint.entries.clone(),
error: self.error.clone(),
}
}
fn save(&self) -> anyhow::Result<()> {
self.save_checkpoint(&self.checkpoint)
}
fn commit(&mut self, checkpoint: Checkpoint) -> anyhow::Result<()> {
self.save_checkpoint(&checkpoint)?;
self.checkpoint = checkpoint;
Ok(())
}
fn save_checkpoint(&self, checkpoint: &Checkpoint) -> anyhow::Result<()> {
self.log.append(&SessionEvent::new(
"summary_checkpoint",
self.log.id().into(),
Path::new(".").to_path_buf(),
serde_json::to_value(checkpoint)?,
))
}
}
fn load(log: &Session, auto_start: bool) -> anyhow::Result<Checkpoint> {
let history = log.read_recent_events_tolerant_tail(1, 1024 * 1024, 1024 * 1024)?;
match history.events.last() {
Some(event) if event.event_type == "summary_checkpoint" => {
let checkpoint: Checkpoint = serde_json::from_value(event.payload.clone())?;
anyhow::ensure!(
checkpoint.entries.len() <= MAX_ENTRIES
&& checkpoint.cursors.len() <= MAX_SOURCES
&& checkpoint.children.len() < MAX_SOURCES,
"summary checkpoint exceeds limits"
);
anyhow::ensure!(
checkpoint
.entries
.iter()
.all(|entry| entry.chars().count() <= 600),
"summary entry exceeds limit"
);
Ok(checkpoint)
}
Some(_) => anyhow::bail!("unexpected summary log record"),
None if !history.diagnostics.is_empty() => anyhow::bail!("summary log cannot be read"),
None => Ok(Checkpoint {
enabled: auto_start,
..Checkpoint::default()
}),
}
}
fn worker(receiver: crossbeam_channel::Receiver<Command>, mailbox: Mailbox) {
let mut active: Option<Active> = None;
loop {
match receiver.recv_timeout(POLL_INTERVAL) {
Ok(command) => {
let (config, settings) = &*command.config;
let Some(session) = command.session.as_ref() else {
active = None;
continue;
};
drop(active.take());
let opened = (|| -> anyhow::Result<Active> {
let root = config.paths.sessions.join("summaries");
crate::sessions::prepare_session_root(&root)?;
crate::sessions::validate_session_id(session.id().into())?;
anyhow::ensure!(
session.path()
== config
.paths
.sessions
.join(format!("{}.jsonl", session.id())),
"summary observer requires a primary session"
);
let lease = crate::persistence::CrossProcessFileLock::try_acquire(
&root.join(format!("{}.observer", session.id())),
)?
.ok_or_else(|| {
anyhow::anyhow!("summary session is already observed by another runtime")
})?;
let log = SessionManager::new(root).open(session.id())?;
let mut checkpoint = load(&log, settings.summarizer.auto_start)?;
if let Some(enabled) = command.enabled {
checkpoint.enabled = enabled;
}
Ok(Active {
command: Command {
generation: command.generation,
session: command.session.clone(),
enabled: command.enabled,
cancellation: command.cancellation.clone(),
config: Arc::clone(&command.config),
},
log,
checkpoint,
error: None,
_lease: lease,
})
})();
match opened {
Ok(state) => {
let mut state = state;
if let Err(error) = state.save() {
state.error = Some(safe_error(&error));
state.checkpoint.enabled = false;
}
publish(&mailbox, state.command.generation, state.snapshot(false));
active = Some(state);
}
Err(error) => {
publish(
&mailbox,
command.generation,
SummarySnapshot {
session_id: session.id().into(),
enabled: false,
running: false,
entries: vec![],
error: Some(safe_error(&error)),
},
);
active = None;
}
}
continue;
}
Err(crossbeam_channel::RecvTimeoutError::Disconnected) => break,
Err(crossbeam_channel::RecvTimeoutError::Timeout) => {}
}
let Some(state) = active.as_mut() else {
continue;
};
if !state.checkpoint.enabled
|| state.command.cancellation.is_canceled()
|| state.error.is_some()
{
continue;
}
let result = (|| -> anyhow::Result<()> {
let (config, settings) = &*state.command.config;
let session = state
.command
.session
.as_ref()
.ok_or_else(|| anyhow::anyhow!("missing summary session"))?;
let (activity, cursors, children) =
collect_activity(session, config, &state.checkpoint)?;
if cursors == state.checkpoint.cursors && children == state.checkpoint.children {
return Ok(());
}
state.command.cancellation.check()?;
let entries = if activity.is_empty() {
vec![]
} else {
publish(&mailbox, state.command.generation, state.snapshot(true));
generate(
config,
settings,
&activity,
&state.checkpoint.entries,
&state.command.cancellation,
)?
};
state.command.cancellation.check()?;
let mut checkpoint = state.checkpoint.clone();
checkpoint.cursors = cursors;
checkpoint.children = children;
append_unique_entries(&mut checkpoint.entries, entries);
let excess = checkpoint.entries.len().saturating_sub(MAX_ENTRIES);
checkpoint.entries.drain(..excess);
state.commit(checkpoint)?;
Ok(())
})();
if let Err(error) = result
&& !state.command.cancellation.is_canceled()
{
state.error = Some(safe_error(&error));
}
publish(&mailbox, state.command.generation, state.snapshot(false));
}
}
fn append_unique_entries(previous: &mut Vec<String>, entries: Vec<String>) {
let recent_start = previous.len().saturating_sub(RECENT_ENTRIES);
for entry in entries {
if !previous[recent_start..].contains(&entry) {
previous.push(entry);
}
}
}
fn safe_error(error: &anyhow::Error) -> String {
crate::output::redact_sensitive_text(&error.to_string())
.chars()
.take(300)
.collect()
}
fn fingerprint(event: &SessionEvent) -> anyhow::Result<u64> {
Ok(crate::hash::stable_hash(&serde_json::to_vec(event)?))
}
type ActivityBatch = (String, BTreeMap<String, u64>, Vec<String>);
fn collect_activity(
primary: &Session,
config: &EffectiveConfig,
checkpoint: &Checkpoint,
) -> anyhow::Result<ActivityBatch> {
let mut children = checkpoint.children.clone();
let mut cursors = checkpoint.cursors.clone();
let mut activity = String::new();
let mut sources = vec![primary.clone()];
let manager = SessionManager::new(config.paths.sessions.join("subagents"));
for child in &children {
sources.push(manager.open(child)?);
}
let mut index = 0;
let mut remaining = BATCH_CHARS;
while index < sources.len() && remaining > 0 {
let source = &sources[index];
let mut source_remaining = remaining / (sources.len() - index);
let history =
source.read_recent_events_tolerant_tail(WINDOW_EVENTS, WINDOW_BYTES, WINDOW_BYTES)?;
let events = history.events;
let previous = cursors.get(source.id()).copied();
let mut start = 0;
if let Some(previous) = previous {
if let Some(position) = events
.iter()
.rposition(|event| fingerprint(event).ok() == Some(previous))
{
start = position + 1;
} else if !events.is_empty() {
let marker = "[Earlier activity omitted or session rewritten.]\n";
if source_remaining >= marker.len() {
activity.push_str(marker);
remaining -= marker.len();
source_remaining -= marker.len();
}
}
}
let mut discovered = vec![];
for event in events.iter().skip(start) {
if remaining == 0 || source_remaining == 0 {
break;
}
discover_children(event, &mut children, &mut discovered);
if let Some(text) = activity_text(event) {
let source_label = if index == 0 { "Primary" } else { "Subagent" };
let header = format!("{source_label} {} {}: ", source.id(), event.event_type);
let overhead = header.chars().count() + 1;
let available = remaining.min(source_remaining).saturating_sub(overhead);
if available == 0 {
break;
}
let text: String = text.chars().take(available.min(2000)).collect();
let used = text.chars().count() + overhead;
remaining -= used;
source_remaining -= used;
activity.push_str(&format!("{header}{text}\n"));
}
cursors.insert(source.id().into(), fingerprint(event)?);
}
for child in discovered {
sources.push(manager.open(child)?);
}
index += 1;
}
Ok((activity, cursors, children))
}
fn discover_children(
event: &SessionEvent,
children: &mut Vec<String>,
discovered: &mut Vec<String>,
) {
if event.event_type == crate::sessions::SessionEventKind::SubagentSession.as_str() {
if let Some(id) = event.payload["session_id"].as_str()
&& children.len() < MAX_SOURCES - 1
&& crate::sessions::validate_session_id(id.into()).is_ok()
&& !children.iter().any(|child| child == id)
{
children.push(id.into());
discovered.push(id.into());
}
return;
}
if event.event_type != "tool_result" {
return;
}
let result = &event.payload["result"];
if result["tool_name"] != "subagents" {
return;
}
let Some(content) = result["content"].as_str() else {
return;
};
let Ok(output) = serde_json::from_str::<serde_json::Value>(content) else {
return;
};
if let Some(results) = output["results"].as_array() {
for result in results {
let Some(id) = result["session_id"].as_str() else {
continue;
};
if children.len() >= MAX_SOURCES - 1 {
break;
}
if crate::sessions::validate_session_id(id.into()).is_ok()
&& !children.iter().any(|child| child == id)
{
children.push(id.into());
discovered.push(id.into());
}
}
}
}
fn activity_text(event: &SessionEvent) -> Option<String> {
let text = match event.event_type.as_str() {
"user_input" | "assistant_output" | "assistant_chunk" => {
event.payload["text"].as_str()?.to_string()
}
"tool_call" | "tool_result" | "turn_status" => event.payload.to_string(),
_ => return None,
};
Some(crate::output::redact_sensitive_text(&text))
}
fn summary_selection(
config: &EffectiveConfig,
settings: &Settings,
) -> anyhow::Result<(String, String)> {
fn identifier(value: Option<&str>) -> Option<&str> {
value.map(str::trim).filter(|value| !value.is_empty())
}
let active_provider = identifier(config.provider.as_deref());
let provider = identifier(settings.summarizer.provider.as_deref())
.or(active_provider)
.ok_or_else(|| anyhow::anyhow!("summary provider is not configured"))?;
let inherited_model = (Some(provider) == active_provider)
.then(|| identifier(config.model.as_deref()))
.flatten();
let model = identifier(settings.summarizer.model.as_deref())
.or(inherited_model)
.or_else(|| match provider {
crate::providers::ANTHROPIC_PROVIDER | crate::providers::OPENAI_CODEX_PROVIDER => {
Some(crate::providers::default_model_for_provider(provider))
}
_ => None,
})
.ok_or_else(|| anyhow::anyhow!("summary model must be configured for this provider"))?;
Ok((provider.into(), model.into()))
}
fn summary_prompt(settings: &Settings) -> String {
settings
.summarizer
.prompt
.as_deref()
.unwrap_or(PROMPT)
.chars()
.take(8000)
.collect()
}
fn generate(
config: &EffectiveConfig,
settings: &Settings,
activity: &str,
previous: &[String],
cancellation: &AgentCancellation,
) -> anyhow::Result<Vec<String>> {
cancellation.check()?;
let (provider_id, model) = summary_selection(config, settings)?;
let mut selected = if config.provider.as_deref().map(str::trim) == Some(provider_id.as_str()) {
config.clone()
} else {
load_effective_provider_selection(&config.paths, &provider_id, &model)?
};
selected.provider = Some(provider_id.clone());
selected.model = Some(model);
selected.thinking_level = settings
.summarizer
.reasoning
.unwrap_or(config.thinking_level);
let selection = ProviderSelection::from_config(&selected)?;
let provider = provider_from_selection_with_settings_for_workload(
&selected,
&selection,
Path::new("."),
settings,
crate::fast::FastWorkload::SessionTitle,
)?;
let system = summary_prompt(settings);
let previous = previous
.iter()
.rev()
.take(RECENT_ENTRIES)
.rev()
.cloned()
.collect::<Vec<_>>();
let request = ProviderRequest::new_without_tools(
selection.model,
vec![
ChatMessage::system(system),
ChatMessage::user(format!(
"Recent entries (avoid repetition):\n{}\n\nNew observed activity:\n{activity}",
serde_json::to_string(&previous)?
)),
],
)
.with_thinking_level(selected.thinking_level)
.with_text_verbosity(settings.text_verbosity_for(&provider_id));
stream_entries(
provider.as_ref(),
request,
cancellation,
Duration::from_secs(60),
)
}
fn stream_entries(
provider: &dyn crate::providers::Provider,
request: ProviderRequest,
cancellation: &AgentCancellation,
timeout_duration: Duration,
) -> anyhow::Result<Vec<String>> {
cancellation.check()?;
let (token, timeout) = cancellation.child_token();
let (done_sender, done_receiver) = crossbeam_channel::bounded::<()>(1);
let watchdog = thread::Builder::new()
.name("summary-timeout".into())
.spawn(move || {
if done_receiver.recv_timeout(timeout_duration)
== Err(crossbeam_channel::RecvTimeoutError::Timeout)
{
timeout.cancel();
true
} else {
false
}
})?;
let mut raw = String::new();
let result = provider.stream_cancellable(request, &token, &mut |event| {
token.check()?;
if let ProviderEvent::TextDelta(delta) = event {
anyhow::ensure!(
raw.len().saturating_add(delta.len()) <= 16_384,
"summary response exceeds limit"
);
raw.push_str(&delta);
}
Ok(())
});
let _ = done_sender.send(());
let timed_out = watchdog.join().unwrap_or(true);
cancellation.check()?;
anyhow::ensure!(!timed_out, "summary provider timed out");
result?;
parse_entries(&raw)
}
fn parse_entries(raw: &str) -> anyhow::Result<Vec<String>> {
let entries: Vec<String> = serde_json::from_str(raw.trim())
.map_err(|_| anyhow::anyhow!("summary provider must return a JSON array of strings"))?;
anyhow::ensure!(
entries.len() <= 6,
"summary provider returned too many entries"
);
entries
.into_iter()
.map(|entry| {
let entry =
crate::output::sanitize_display_text(&crate::output::redact_sensitive_text(&entry));
let entry = entry.split_whitespace().collect::<Vec<_>>().join(" ");
anyhow::ensure!(
!entry.is_empty() && entry.chars().count() <= 600,
"summary entry is empty or too long"
);
Ok(entry)
})
.collect()
}