use std::path::{Path, PathBuf};
use std::sync::Mutex;
use tokio::sync::{broadcast, watch};
use tokio_util::sync::CancellationToken;
use uuid::Uuid;
use crate::event::{Event, EventSink, FlowRunId, TurnId};
use crate::event_writer::EventWriter;
use crate::injection::{Injection, InjectionId, InjectionState};
use crate::message::{Message, MessageRole};
use crate::stream::StreamFrame;
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct SessionId(pub Uuid);
impl SessionId {
pub fn now() -> Self {
Self(Uuid::new_v4())
}
pub fn parse(s: &str) -> Result<Self, uuid::Error> {
Uuid::parse_str(s).map(Self)
}
}
impl std::fmt::Display for SessionId {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
self.0.fmt(f)
}
}
type WatchKeepalive = (
watch::Receiver<ContextSnapshot>,
watch::Receiver<Option<String>>,
watch::Receiver<usize>,
watch::Receiver<Vec<crate::memory::todo::Todo>>,
watch::Receiver<Vec<crate::memory::plan::Plan>>,
);
pub struct Session {
id: SessionId,
dir: PathBuf,
writer: Option<EventWriter>,
sink: EventSink,
messages: std::sync::Arc<std::sync::Mutex<Vec<Message>>>,
current_turn: Mutex<Option<TurnId>>,
injection_queue: Mutex<Vec<Injection>>,
injection_tx: broadcast::Sender<Injection>,
stream_tx: broadcast::Sender<StreamFrame>,
flow_cancel: Mutex<CancellationToken>,
context_watch: watch::Sender<ContextSnapshot>,
goal_watch: watch::Sender<Option<String>>,
attach_watch: watch::Sender<usize>,
todos_watch: watch::Sender<Vec<crate::memory::todo::Todo>>,
plans_watch: watch::Sender<Vec<crate::memory::plan::Plan>>,
_watch_keepalive: WatchKeepalive,
streamed_this_turn: std::sync::atomic::AtomicBool,
manual_compact_pending: std::sync::atomic::AtomicBool,
last_input_tokens: std::sync::atomic::AtomicU64,
compact_review_mode: Mutex<CompactReviewMode>,
compact_lock: std::sync::Arc<tokio::sync::Mutex<()>>,
last_image_user_msg: Mutex<Option<LastImageUserMsg>>,
read_files: std::sync::Arc<std::sync::Mutex<std::collections::HashSet<std::path::PathBuf>>>,
approval: std::sync::Arc<ApprovalRegistry>,
compact_reviews: std::sync::Arc<CompactReviewRegistry>,
forms: std::sync::Arc<FormRegistry>,
fs_access_mode: Mutex<Option<crate::fs_access::FsAccessMode>>,
project_index: Option<std::sync::Arc<crate::index::AnchorIndex>>,
}
#[derive(Debug, Clone)]
pub struct PendingCompactReview {
pub review_id: String,
pub summary: String,
pub slice_preview: String,
pub slice_count: usize,
pub range_start: usize,
pub range_end: usize,
pub tokens_before: u64,
pub emitted_at: chrono::DateTime<chrono::Utc>,
}
#[derive(Debug, Clone)]
pub enum CompactReviewDecision {
AcceptAsIs,
AcceptEdited { summary: String },
Reject,
}
pub struct CompactReviewRegistry {
entry: std::sync::Mutex<Option<CompactReviewEntry>>,
watch_tx: watch::Sender<Option<PendingCompactReview>>,
}
struct CompactReviewEntry {
pending: PendingCompactReview,
responder: tokio::sync::oneshot::Sender<CompactReviewDecision>,
}
impl Default for CompactReviewRegistry {
fn default() -> Self {
Self::new()
}
}
impl CompactReviewRegistry {
pub fn new() -> Self {
let (watch_tx, _) = watch::channel(None);
Self {
entry: std::sync::Mutex::new(None),
watch_tx,
}
}
pub fn subscribe(&self) -> watch::Receiver<Option<PendingCompactReview>> {
self.watch_tx.subscribe()
}
pub fn list_pending(&self) -> Option<PendingCompactReview> {
self.entry
.lock()
.unwrap()
.as_ref()
.map(|e| e.pending.clone())
}
pub fn subscriber_count(&self) -> usize {
self.watch_tx.receiver_count()
}
pub fn request(
&self,
pending: PendingCompactReview,
) -> tokio::sync::oneshot::Receiver<CompactReviewDecision> {
let (tx, rx) = tokio::sync::oneshot::channel();
if self.watch_tx.receiver_count() == 0 {
let _ = tx.send(CompactReviewDecision::AcceptAsIs);
return rx;
}
{
let mut slot = self.entry.lock().unwrap();
if let Some(prev) = slot.take() {
let _ = prev.responder.send(CompactReviewDecision::Reject);
}
*slot = Some(CompactReviewEntry {
pending: pending.clone(),
responder: tx,
});
}
let _ = self.watch_tx.send(Some(pending));
rx
}
pub fn decide(&self, review_id: &str, decision: CompactReviewDecision) -> bool {
let entry = {
let mut slot = self.entry.lock().unwrap();
match slot.as_ref() {
Some(e) if e.pending.review_id == review_id => slot.take(),
_ => None,
}
};
match entry {
Some(e) => {
let _ = e.responder.send(decision);
let _ = self.watch_tx.send(None);
true
}
None => false,
}
}
}
#[derive(Debug, Clone)]
pub struct PendingApproval {
pub tool_use_id: String,
pub tool_name: String,
pub args_preview: String,
pub preview: Option<String>,
pub level: crate::tool::ApprovalLevel,
pub run_id: FlowRunId,
pub emitted_at: chrono::DateTime<chrono::Utc>,
pub bypass_auto_ceiling: bool,
}
#[derive(Debug, Clone)]
pub enum ApprovalDecision {
Approve,
Deny { reason: String },
}
pub struct FormRegistry {
entries: std::sync::Mutex<Vec<FormEntry>>,
watch_tx: watch::Sender<Vec<crate::form::PendingForm>>,
}
struct FormEntry {
pending: crate::form::PendingForm,
responder: tokio::sync::oneshot::Sender<crate::form::FormAnswer>,
}
impl Default for FormRegistry {
fn default() -> Self {
Self::new()
}
}
impl FormRegistry {
pub fn new() -> Self {
let (watch_tx, _) = watch::channel(Vec::new());
Self {
entries: std::sync::Mutex::new(Vec::new()),
watch_tx,
}
}
pub fn subscribe(&self) -> watch::Receiver<Vec<crate::form::PendingForm>> {
self.watch_tx.subscribe()
}
pub fn list_pending(&self) -> Vec<crate::form::PendingForm> {
self.entries
.lock()
.unwrap()
.iter()
.map(|e| e.pending.clone())
.collect()
}
pub fn subscriber_count(&self) -> usize {
self.watch_tx.receiver_count()
}
pub fn request(
&self,
pending: crate::form::PendingForm,
) -> tokio::sync::oneshot::Receiver<crate::form::FormAnswer> {
let (tx, rx) = tokio::sync::oneshot::channel();
if self.watch_tx.receiver_count() == 0 {
let _ = tx.send(crate::form::FormAnswer::Cancelled);
return rx;
}
{
let mut entries = self.entries.lock().unwrap();
entries.push(FormEntry {
pending: pending.clone(),
responder: tx,
});
}
self.broadcast_snapshot();
rx
}
pub fn submit(&self, form_id: &str, answer: crate::form::FormAnswer) -> bool {
let entry = {
let mut entries = self.entries.lock().unwrap();
let pos = entries.iter().position(|e| e.pending.form_id == form_id);
pos.map(|p| entries.remove(p))
};
match entry {
Some(e) => {
let _ = e.responder.send(answer);
self.broadcast_snapshot();
true
}
None => false,
}
}
pub fn cancel_all(&self) {
let drained: Vec<FormEntry> = {
let mut entries = self.entries.lock().unwrap();
std::mem::take(&mut *entries)
};
for e in drained {
let _ = e.responder.send(crate::form::FormAnswer::Cancelled);
}
self.broadcast_snapshot();
}
pub fn promote(&self, form_id: &str) {
let mut entries = self.entries.lock().unwrap();
if let Some(pos) = entries.iter().position(|e| e.pending.form_id == form_id) {
if pos == 0 {
return;
}
let entry = entries.remove(pos);
entries.insert(0, entry);
}
drop(entries);
self.broadcast_snapshot();
}
fn broadcast_snapshot(&self) {
let snap = self
.entries
.lock()
.unwrap()
.iter()
.map(|e| e.pending.clone())
.collect();
let _ = self.watch_tx.send(snap);
}
}
pub struct ApprovalRegistry {
entries: std::sync::Mutex<Vec<ApprovalEntry>>,
auto_ceiling: std::sync::Mutex<crate::tool::ApprovalLevel>,
watch_tx: watch::Sender<Vec<PendingApproval>>,
}
struct ApprovalEntry {
pending: PendingApproval,
responder: tokio::sync::oneshot::Sender<ApprovalDecision>,
}
impl Default for ApprovalRegistry {
fn default() -> Self {
Self::new()
}
}
impl ApprovalRegistry {
pub fn new() -> Self {
let (watch_tx, _) = watch::channel(Vec::new());
Self {
entries: std::sync::Mutex::new(Vec::new()),
auto_ceiling: std::sync::Mutex::new(crate::tool::ApprovalLevel::Approve),
watch_tx,
}
}
pub fn subscribe(&self) -> watch::Receiver<Vec<PendingApproval>> {
self.watch_tx.subscribe()
}
pub fn list_pending(&self) -> Vec<PendingApproval> {
self.entries
.lock()
.unwrap()
.iter()
.map(|e| e.pending.clone())
.collect()
}
pub fn set_auto_ceiling(&self, level: crate::tool::ApprovalLevel) {
*self.auto_ceiling.lock().unwrap() = level;
}
pub fn request(
&self,
pending: PendingApproval,
) -> tokio::sync::oneshot::Receiver<ApprovalDecision> {
let (tx, rx) = tokio::sync::oneshot::channel();
if !pending.bypass_auto_ceiling && pending.level <= *self.auto_ceiling.lock().unwrap() {
let _ = tx.send(ApprovalDecision::Approve);
return rx;
}
{
let mut entries = self.entries.lock().unwrap();
entries.push(ApprovalEntry {
pending,
responder: tx,
});
}
self.broadcast_snapshot();
rx
}
pub fn decide(&self, tool_use_id: &str, decision: ApprovalDecision) -> bool {
let mut entries = self.entries.lock().unwrap();
if let Some(pos) = entries
.iter()
.position(|e| e.pending.tool_use_id == tool_use_id)
{
let entry = entries.remove(pos);
let _ = entry.responder.send(decision);
drop(entries);
self.broadcast_snapshot();
true
} else {
false
}
}
pub fn decide_all(&self, decision: ApprovalDecision) -> usize {
let mut entries = self.entries.lock().unwrap();
let count = entries.len();
for entry in entries.drain(..) {
let _ = entry.responder.send(decision.clone());
}
drop(entries);
self.broadcast_snapshot();
count
}
fn broadcast_snapshot(&self) {
let snapshot = self
.entries
.lock()
.unwrap()
.iter()
.map(|e| e.pending.clone())
.collect();
let _ = self.watch_tx.send(snapshot);
}
}
type ImagePart = (usize, String);
#[derive(Debug, Clone)]
struct LastImageUserMsg {
message_seq: u64,
images: Vec<ImagePart>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum CompactReviewMode {
Always,
#[default]
ManualOnly,
Never,
}
impl CompactReviewMode {
pub fn parse(s: &str) -> Option<Self> {
match s.trim() {
"always" => Some(Self::Always),
"manual-only" | "manual_only" => Some(Self::ManualOnly),
"never" => Some(Self::Never),
_ => None,
}
}
pub fn should_review(self, forced: bool) -> bool {
match self {
Self::Always => true,
Self::ManualOnly => forced,
Self::Never => false,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CompactResult {
pub before_tokens: u64,
pub after_tokens: u64,
pub compacted_start: usize,
pub compacted_end: usize,
}
#[derive(Debug, Clone, Default, PartialEq)]
pub struct ContextSnapshot {
pub model: String,
pub tokens_in: u64,
pub tokens_out: u64,
pub cost_usd: f64,
pub mcp_ok: u16,
pub mcp_total: u16,
pub memory_recent_count: u16,
pub window_tokens: u64,
pub window_budget: u64,
pub cache_read: u64,
pub cache_write: u64,
pub last_ttft_ms: u64,
pub last_tokens_per_sec: f64,
}
#[derive(Debug, thiserror::Error)]
pub enum SessionOpenError {
#[error("invalid session id `{sid}` (want a UUID)")]
InvalidId { sid: String },
#[error("session `{sid}` not found at {}", dir.display())]
NotFound { sid: String, dir: PathBuf },
#[error("session writer init: {0}")]
WriterInit(#[source] std::io::Error),
#[error("replay {}: {source}", path.display())]
Replay {
path: PathBuf,
#[source]
source: std::io::Error,
},
}
fn load_goal(dir: &Path) -> Option<String> {
if dir.as_os_str().is_empty() {
return None;
}
let store = crate::memory::goal::GoalStore::at(dir);
match store.get() {
Ok(s) if !s.is_empty() => Some(s),
_ => None,
}
}
#[derive(Debug, Clone)]
pub enum TranscriptEntry {
Message {
message: Message,
flow_run_id: Option<String>,
},
CompactionSummary {
range_start: usize,
range_end: usize,
compacted_count: usize,
before_tokens: u64,
after_tokens: u64,
summary: String,
ts: Option<chrono::DateTime<chrono::Utc>>,
},
DiffPreview {
title: String,
old_content: Option<String>,
new_content: Option<String>,
unified_diff: Option<String>,
},
FlowGraph {
run_id: String,
flow_name: String,
graph: crate::nodegraph::FlowGraph,
ts: Option<chrono::DateTime<chrono::Utc>>,
},
FlowStart {
run_id: String,
flow_name: String,
parent_run_id: Option<String>,
parent_node_id: Option<String>,
ts: Option<chrono::DateTime<chrono::Utc>>,
},
FlowNodeStart {
run_id: String,
node_id: String,
kind: crate::nodegraph::NodeKind,
label: String,
parent_node_id: Option<String>,
ts: Option<chrono::DateTime<chrono::Utc>>,
},
FlowNodeEnd {
run_id: String,
node_id: String,
status: crate::event::FlowNodeStatus,
output_preview: Option<String>,
ts: Option<chrono::DateTime<chrono::Utc>>,
},
ToolNode {
run_id: String,
parent_node_id: String,
tool_use_id: String,
tool_name: String,
args_preview: String,
ts: Option<chrono::DateTime<chrono::Utc>>,
},
FlowDone {
run_id: String,
ok: bool,
cancelled: bool,
ts: Option<chrono::DateTime<chrono::Utc>>,
},
LlmCall {
model: String,
usage: crate::provider::TokenUsage,
wallclock_ms: u64,
ttft_ms: Option<u64>,
tokens_per_second: Option<f64>,
run_id: Option<crate::event::FlowRunId>,
node_id: Option<String>,
ts: Option<chrono::DateTime<chrono::Utc>>,
},
}
fn replay_context_snapshot_from(path: &Path) -> ContextSnapshot {
let mut snap = ContextSnapshot::default();
let text = match std::fs::read_to_string(path) {
Ok(t) => t,
Err(_) => return snap,
};
for value in parse_json_lines(&text) {
if value["type"].as_str() != Some("llm_call") {
continue;
}
if let Some(model) = value["model"].as_str() {
snap.model = model.to_string();
}
let usage = &value["usage"];
let input = usage["input"].as_u64().unwrap_or(0);
let cached = usage["cached_input"].as_u64().unwrap_or(0);
let output = usage["output"].as_u64().unwrap_or(0);
snap.tokens_in = snap.tokens_in.saturating_add(input).saturating_add(cached);
snap.tokens_out = snap.tokens_out.saturating_add(output);
}
snap
}
fn replay_messages_from(path: &Path) -> Result<Vec<Message>, SessionOpenError> {
if let Some((checkpoint_seq, messages)) = load_last_checkpoint(path)? {
let values = read_jsonl_values(path)?;
let patches = collect_attachment_patches(&values);
let mut message_seqs: Vec<(u64, Message)> =
messages.into_iter().map(|message| (0, message)).collect();
for v in &values {
let ty = v["type"].as_str().unwrap_or("");
let seq = v["seq"].as_u64().unwrap_or(0);
if seq <= checkpoint_seq {
continue;
}
if let "user_msg" | "assistant_msg" | "tool_result_msg" | "system_msg" = ty {
if let Some(m) = v.get("message")
&& let Ok(mut msg) = serde_json::from_value::<Message>(m.clone())
{
if let Some(seq) = v["seq"].as_u64()
&& let Some(ps) = patches.get(&seq)
{
apply_attachment_patches(&mut msg, ps);
}
message_seqs.push((seq, msg));
}
} else if ty == "context_compact" {
let Some(event) = parse_context_compact_event(v) else {
continue;
};
if event.range_start > event.range_end || event.range_end >= message_seqs.len() {
continue;
}
let Some(replacement_seq) = event.replacement_msg_seq else {
continue;
};
let Some(replacement_idx) = message_seqs
.iter()
.position(|(msg_seq, _)| *msg_seq == replacement_seq)
else {
continue;
};
if event.after_tokens >= event.before_tokens {
continue;
}
let replacement = message_seqs.remove(replacement_idx);
let removed_count = event.range_end - event.range_start + 1;
for _ in 0..removed_count {
message_seqs.remove(event.range_start);
}
let insertion_idx = event.range_start.min(message_seqs.len());
if let Some(summary) = event.summary_text {
message_seqs.insert(
insertion_idx,
(
replacement_seq,
Message::system_compact_summary(
TurnId::now(),
summary,
event.range_start as u64,
event.range_end as u64,
removed_count,
),
),
);
} else {
message_seqs.insert(insertion_idx, replacement);
}
}
}
return Ok(message_seqs.into_iter().map(|(_, msg)| msg).collect());
}
let entries = replay_transcript_from(path)?;
let mut out = Vec::new();
for entry in entries {
if let TranscriptEntry::Message { message, .. } = entry {
out.push(message);
}
}
Ok(out)
}
fn read_jsonl_values(path: &Path) -> Result<Vec<serde_json::Value>, SessionOpenError> {
let text = match std::fs::read_to_string(path) {
Ok(t) => t,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
Err(e) => {
return Err(SessionOpenError::Replay {
path: path.to_path_buf(),
source: e,
});
}
};
Ok(parse_json_lines(&text))
}
fn load_last_checkpoint(path: &Path) -> Result<Option<(u64, Vec<Message>)>, SessionOpenError> {
let values = read_jsonl_values(path)?;
for v in values.iter().rev() {
if v["type"].as_str() == Some("checkpoint") {
let seq = v["seq"].as_u64().unwrap_or(0);
let messages = v
.get("messages")
.and_then(|m| serde_json::from_value::<Vec<Message>>(m.clone()).ok())
.unwrap_or_default();
return Ok(Some((seq, messages)));
}
}
Ok(None)
}
fn find_last_seq(path: &Path) -> Result<Option<u64>, SessionOpenError> {
let values = read_jsonl_values(path)?;
Ok(values.iter().rev().find_map(|v| v["seq"].as_u64()))
}
#[derive(Debug, Clone)]
struct AttachmentPatch {
part_index: usize,
file_basename: String,
reason: String,
}
fn parse_ts(v: &serde_json::Value) -> Option<chrono::DateTime<chrono::Utc>> {
v.get("ts")?
.as_str()
.and_then(|s| chrono::DateTime::parse_from_rfc3339(s).ok())
.map(|dt| dt.with_timezone(&chrono::Utc))
}
fn parse_context_compact_event(v: &serde_json::Value) -> Option<CompactReplayEvent> {
if v["type"].as_str() != Some("context_compact") {
return None;
}
Some(CompactReplayEvent {
range_start: v["compacted_range_start"].as_u64().unwrap_or(0) as usize,
range_end: v["compacted_range_end"].as_u64().unwrap_or(0) as usize,
before_tokens: v["before_tokens"].as_u64().unwrap_or(0),
after_tokens: v["after_tokens"].as_u64().unwrap_or(0),
summary_text: v
.get("summary_text")
.and_then(|s| s.as_str())
.map(String::from),
replacement_msg_seq: v["replacement_msg_seq"].as_u64(),
})
}
#[derive(Debug, Clone)]
struct CompactReplayEvent {
range_start: usize,
range_end: usize,
before_tokens: u64,
after_tokens: u64,
summary_text: Option<String>,
replacement_msg_seq: Option<u64>,
}
fn parse_json_lines(text: &str) -> Vec<serde_json::Value> {
text.lines()
.filter_map(|line| {
let t = line.trim();
if t.is_empty() {
None
} else {
serde_json::from_str::<serde_json::Value>(t).ok()
}
})
.collect()
}
fn collect_attachment_patches(
values: &[serde_json::Value],
) -> std::collections::HashMap<u64, Vec<AttachmentPatch>> {
let mut map: std::collections::HashMap<u64, Vec<AttachmentPatch>> =
std::collections::HashMap::new();
for v in values {
if v["type"].as_str() == Some("attachment_degraded") {
let Some(msg_seq) = v["message_seq"].as_u64() else {
continue;
};
let Some(part_index) = v["part_index"].as_u64() else {
continue;
};
let file_basename = v["file_basename"].as_str().unwrap_or("").to_string();
let reason = v["reason"].as_str().unwrap_or("degraded").to_string();
map.entry(msg_seq).or_default().push(AttachmentPatch {
part_index: part_index as usize,
file_basename,
reason,
});
}
}
map
}
fn apply_attachment_patches(msg: &mut Message, patches: &[AttachmentPatch]) {
for p in patches {
if let Some(part) = msg.parts.get_mut(p.part_index) {
*part = crate::message::MessagePart::Text {
text: format!(
"[attachment unavailable: {} — {}]",
p.file_basename, p.reason
),
};
}
}
}
pub fn replay_transcript_from(path: &Path) -> Result<Vec<TranscriptEntry>, SessionOpenError> {
let text = match std::fs::read_to_string(path) {
Ok(t) => t,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
Err(e) => {
return Err(SessionOpenError::Replay {
path: path.to_path_buf(),
source: e,
});
}
};
let values = parse_json_lines(&text);
let patches = collect_attachment_patches(&values);
let mut out = Vec::new();
let mut msg_indices: Vec<usize> = Vec::new();
let mut msg_seqs: Vec<u64> = Vec::new();
for v in &values {
let ty = v["type"].as_str().unwrap_or("");
match ty {
"user_msg" | "assistant_msg" | "tool_result_msg" | "system_msg" => {
if let Some(m) = v.get("message")
&& let Ok(mut msg) = serde_json::from_value::<Message>(m.clone())
{
let seq = v["seq"].as_u64().unwrap_or(0);
if let Some(ps) = patches.get(&seq) {
apply_attachment_patches(&mut msg, ps);
}
let flow_run_id = v["flow_run_id"].as_str().map(String::from);
msg_indices.push(out.len());
msg_seqs.push(seq);
out.push(TranscriptEntry::Message {
message: msg,
flow_run_id,
});
}
}
"context_compact" => {
let Some(event) = parse_context_compact_event(v) else {
continue;
};
if event.range_start > event.range_end || event.range_end >= msg_indices.len() {
continue;
}
let Some(replacement_seq) = event.replacement_msg_seq else {
continue;
};
let Some(replacement_pos) = msg_seqs.iter().position(|seq| *seq == replacement_seq)
else {
continue;
};
let replacement_out_idx = msg_indices[replacement_pos];
let replacement_entry = out.remove(replacement_out_idx);
let removed_out_start = msg_indices[event.range_start];
let removed_count = event.range_end - event.range_start + 1;
for _ in 0..removed_count {
out.remove(removed_out_start);
}
msg_indices.drain(event.range_start..=event.range_end);
msg_seqs.drain(event.range_start..=event.range_end);
out.insert(removed_out_start, replacement_entry);
msg_indices.insert(event.range_start, removed_out_start);
msg_seqs.insert(event.range_start, replacement_seq);
for (i, ordinal_out_idx) in msg_indices.iter_mut().enumerate() {
if i > event.range_start {
*ordinal_out_idx =
ordinal_out_idx.saturating_sub(removed_count.saturating_sub(1));
}
}
}
"compaction_summary" => {
out.push(TranscriptEntry::CompactionSummary {
range_start: v["range_start"].as_u64().unwrap_or(0) as usize,
range_end: v["range_end"].as_u64().unwrap_or(0) as usize,
compacted_count: v["compacted_count"].as_u64().unwrap_or(0) as usize,
before_tokens: v["before_tokens"].as_u64().unwrap_or(0),
after_tokens: v["after_tokens"].as_u64().unwrap_or(0),
summary: v["summary"].as_str().unwrap_or("").to_string(),
ts: parse_ts(v),
});
}
"diff_preview" => {
out.push(TranscriptEntry::DiffPreview {
title: v["title"].as_str().unwrap_or("").to_string(),
old_content: v["old_content"].as_str().map(String::from),
new_content: v["new_content"].as_str().map(String::from),
unified_diff: v["unified_diff"].as_str().map(String::from),
});
}
"flow_graph" => {
let run_id = v["run_id"].as_str().unwrap_or("").to_string();
let flow_name = v
.get("graph")
.and_then(|g| g["flow_name"].as_str())
.unwrap_or("")
.to_string();
let ts = parse_ts(v);
if let Some(g) = v.get("graph")
&& let Ok(graph) =
serde_json::from_value::<crate::nodegraph::FlowGraph>(g.clone())
{
out.push(TranscriptEntry::FlowGraph {
run_id,
flow_name,
graph,
ts,
});
}
}
"flow_start" => {
let run_id = v["run_id"].as_str().unwrap_or("").to_string();
let flow_name = v["flow_name"].as_str().unwrap_or("").to_string();
let parent_run_id = v["parent_run_id"].as_str().map(String::from);
let parent_node_id = v["parent_node_id"].as_str().map(String::from);
let ts = parse_ts(v);
out.push(TranscriptEntry::FlowStart {
run_id,
flow_name,
parent_run_id,
parent_node_id,
ts,
});
}
"flow_node_start" => {
let run_id = v["run_id"].as_str().unwrap_or("").to_string();
let node_id = v["node_id"].as_str().unwrap_or("").to_string();
let label = v["label"].as_str().unwrap_or(&node_id).to_string();
let parent_node_id = v["parent_node_id"].as_str().map(String::from);
let kind = v
.get("kind")
.and_then(|k| serde_json::from_value(k.clone()).ok())
.unwrap_or(crate::nodegraph::NodeKind::UserConfirm);
let ts = parse_ts(v);
out.push(TranscriptEntry::FlowNodeStart {
run_id,
node_id,
kind,
label,
parent_node_id,
ts,
});
}
"flow_node_end" => {
let run_id = v["run_id"].as_str().unwrap_or("").to_string();
let node_id = v["node_id"].as_str().unwrap_or("").to_string();
let status: crate::event::FlowNodeStatus = v
.get("status")
.and_then(|s| serde_json::from_value(s.clone()).ok())
.unwrap_or(crate::event::FlowNodeStatus::Ok);
let output_preview = v["output_preview"].as_str().map(String::from);
let ts = parse_ts(v);
out.push(TranscriptEntry::FlowNodeEnd {
run_id,
node_id,
status,
output_preview,
ts,
});
}
"tool_node" => {
let run_id = v["run_id"].as_str().unwrap_or("").to_string();
let parent_node_id = v["parent_node_id"].as_str().unwrap_or("").to_string();
let tool_use_id = v["tool_use_id"].as_str().unwrap_or("").to_string();
let tool_name = v["tool_name"].as_str().unwrap_or("").to_string();
let args_preview = v["args_preview"].as_str().unwrap_or("").to_string();
let ts = parse_ts(v);
out.push(TranscriptEntry::ToolNode {
run_id,
parent_node_id,
tool_use_id,
tool_name,
args_preview,
ts,
});
}
"flow_end" => {
let run_id = v["run_id"].as_str().unwrap_or("").to_string();
let ok = v["status"]["kind"].as_str() == Some("ok");
let cancelled = v["status"]["kind"].as_str() == Some("cancelled");
let ts = parse_ts(v);
out.push(TranscriptEntry::FlowDone {
run_id,
ok,
cancelled,
ts,
});
}
"llm_call" => {
let model = v["model"].as_str().unwrap_or("").to_string();
let usage: crate::provider::TokenUsage = v
.get("usage")
.and_then(|u| serde_json::from_value(u.clone()).ok())
.unwrap_or_default();
let wallclock_ms = v["wallclock_ms"].as_u64().unwrap_or(0);
let ttft_ms = v["ttft_ms"].as_u64();
let tokens_per_second = v["tokens_per_second"].as_f64();
let run_id = v["run_id"]
.as_str()
.and_then(|s| uuid::Uuid::parse_str(s).ok())
.map(crate::event::FlowRunId);
let node_id = v["node_id"].as_str().map(String::from);
let ts = parse_ts(v);
out.push(TranscriptEntry::LlmCall {
model,
usage,
wallclock_ms,
ttft_ms,
tokens_per_second,
run_id,
node_id,
ts,
});
}
_ => {}
}
}
Ok(out)
}
fn default_project_index(root: &Path) -> Option<std::sync::Arc<crate::index::AnchorIndex>> {
match crate::index::AnchorIndex::open_project(root) {
Ok(idx) => Some(std::sync::Arc::new(idx)),
Err(e) => {
eprintln!(
"[atman] project index unavailable at {} — history search disabled: {e}",
root.display()
);
None
}
}
}
impl Session {
pub fn open(root: impl AsRef<Path>) -> std::io::Result<Self> {
Self::open_with_redactor(root, None)
}
pub fn open_with_redactor(
root: impl AsRef<Path>,
redactor: Option<std::sync::Arc<crate::redact::Redactor>>,
) -> std::io::Result<Self> {
let root_ref = root.as_ref();
let project_index = default_project_index(root_ref);
Self::open_with_context(root_ref, redactor, project_index)
}
pub fn open_with_context(
root: impl AsRef<Path>,
redactor: Option<std::sync::Arc<crate::redact::Redactor>>,
project_index: Option<std::sync::Arc<crate::index::AnchorIndex>>,
) -> std::io::Result<Self> {
let id = SessionId::now();
let dir = root.as_ref().join("sessions").join(id.to_string());
let writer = EventWriter::spawn_full(
&dir,
redactor.clone(),
project_index.clone(),
Some(id.to_string()),
)?;
if let Err(e) = crate::session_meta::SessionMeta::from_cwd().save(&dir) {
eprintln!("[atman] session meta write failed: {e}");
}
let mut sink = EventSink::new().with_forwarder(writer.sender());
if let Some(r) = redactor {
sink = sink.with_redactor(r);
}
let (injection_tx, _) = broadcast::channel(32);
let (stream_tx, _) = broadcast::channel(1024);
let (context_watch, context_rx) = watch::channel(ContextSnapshot::default());
let (goal_watch, goal_rx) = watch::channel(None);
let (attach_watch, attach_rx) = watch::channel(0);
let (todos_watch, todos_rx) = watch::channel(Vec::new());
let (plans_watch, plans_rx) = watch::channel(Vec::new());
Ok(Self {
id,
dir,
writer: Some(writer),
sink,
messages: std::sync::Arc::new(std::sync::Mutex::new(Vec::new())),
current_turn: Mutex::new(None),
injection_queue: Mutex::new(Vec::new()),
injection_tx,
stream_tx,
flow_cancel: Mutex::new(CancellationToken::new()),
context_watch,
goal_watch,
attach_watch,
todos_watch,
plans_watch,
_watch_keepalive: (context_rx, goal_rx, attach_rx, todos_rx, plans_rx),
streamed_this_turn: std::sync::atomic::AtomicBool::new(false),
manual_compact_pending: std::sync::atomic::AtomicBool::new(false),
last_input_tokens: std::sync::atomic::AtomicU64::new(0),
compact_review_mode: Mutex::new(CompactReviewMode::default()),
compact_lock: std::sync::Arc::new(tokio::sync::Mutex::new(())),
last_image_user_msg: Mutex::new(None),
read_files: std::sync::Arc::new(
std::sync::Mutex::new(std::collections::HashSet::new()),
),
approval: std::sync::Arc::new(ApprovalRegistry::new()),
compact_reviews: std::sync::Arc::new(CompactReviewRegistry::new()),
forms: std::sync::Arc::new(FormRegistry::new()),
fs_access_mode: Mutex::new(None),
project_index,
})
}
pub fn open_existing(root: impl AsRef<Path>, sid: &str) -> Result<Self, SessionOpenError> {
Self::open_existing_with_redactor(root, sid, None)
}
pub fn open_existing_with_redactor(
root: impl AsRef<Path>,
sid: &str,
redactor: Option<std::sync::Arc<crate::redact::Redactor>>,
) -> Result<Self, SessionOpenError> {
let project_index = default_project_index(root.as_ref());
Self::open_existing_with_context(root, sid, redactor, project_index)
}
pub fn open_existing_with_context(
root: impl AsRef<Path>,
sid: &str,
redactor: Option<std::sync::Arc<crate::redact::Redactor>>,
project_index: Option<std::sync::Arc<crate::index::AnchorIndex>>,
) -> Result<Self, SessionOpenError> {
let id = SessionId::parse(sid).map_err(|_| SessionOpenError::InvalidId {
sid: sid.to_string(),
})?;
let dir = root.as_ref().join("sessions").join(id.to_string());
if !dir.exists() {
return Err(SessionOpenError::NotFound {
sid: sid.to_string(),
dir: dir.clone(),
});
}
let writer = EventWriter::spawn_full(
&dir,
redactor.clone(),
project_index.clone(),
Some(id.to_string()),
)
.map_err(SessionOpenError::WriterInit)?;
let mut sink = EventSink::new().with_forwarder(writer.sender());
if let Some(r) = redactor {
sink = sink.with_redactor(r);
}
let events_path = dir.join("events.jsonl");
let messages = replay_messages_from(&events_path)?;
if let Some(last_seq) = find_last_seq(&events_path)? {
sink.restore_seq(last_seq);
}
let mut initial_context = replay_context_snapshot_from(&events_path);
initial_context.window_tokens = crate::compaction::estimate_tokens_for_messages(&messages);
initial_context.window_budget =
crate::model_registry::model_info(&initial_context.model).context_budget;
let initial_goal = load_goal(&dir);
let (injection_tx, _) = broadcast::channel(32);
let (stream_tx, _) = broadcast::channel(1024);
let (context_watch, context_rx) = watch::channel(initial_context);
let (goal_watch, goal_rx) = watch::channel(initial_goal);
let (attach_watch, attach_rx) = watch::channel(0);
let (todos_watch, todos_rx) = watch::channel(Vec::new());
let (plans_watch, plans_rx) = watch::channel(Vec::new());
Ok(Self {
id,
dir,
writer: Some(writer),
sink,
messages: std::sync::Arc::new(std::sync::Mutex::new(messages)),
current_turn: Mutex::new(None),
injection_queue: Mutex::new(Vec::new()),
injection_tx,
stream_tx,
flow_cancel: Mutex::new(CancellationToken::new()),
context_watch,
goal_watch,
attach_watch,
todos_watch,
plans_watch,
_watch_keepalive: (context_rx, goal_rx, attach_rx, todos_rx, plans_rx),
streamed_this_turn: std::sync::atomic::AtomicBool::new(false),
manual_compact_pending: std::sync::atomic::AtomicBool::new(false),
last_input_tokens: std::sync::atomic::AtomicU64::new(0),
compact_review_mode: Mutex::new(CompactReviewMode::default()),
compact_lock: std::sync::Arc::new(tokio::sync::Mutex::new(())),
last_image_user_msg: Mutex::new(None),
read_files: std::sync::Arc::new(
std::sync::Mutex::new(std::collections::HashSet::new()),
),
approval: std::sync::Arc::new(ApprovalRegistry::new()),
compact_reviews: std::sync::Arc::new(CompactReviewRegistry::new()),
forms: std::sync::Arc::new(FormRegistry::new()),
fs_access_mode: Mutex::new(None),
project_index,
})
}
pub fn open_ephemeral() -> Self {
let (injection_tx, _) = broadcast::channel(32);
let (stream_tx, _) = broadcast::channel(1024);
let (context_watch, context_rx) = watch::channel(ContextSnapshot::default());
let (goal_watch, goal_rx) = watch::channel(None);
let (attach_watch, attach_rx) = watch::channel(0);
let (todos_watch, todos_rx) = watch::channel(Vec::new());
let (plans_watch, plans_rx) = watch::channel(Vec::new());
Self {
id: SessionId::now(),
dir: PathBuf::new(),
writer: None,
sink: EventSink::new(),
messages: std::sync::Arc::new(std::sync::Mutex::new(Vec::new())),
current_turn: Mutex::new(None),
injection_queue: Mutex::new(Vec::new()),
injection_tx,
stream_tx,
flow_cancel: Mutex::new(CancellationToken::new()),
context_watch,
goal_watch,
attach_watch,
todos_watch,
plans_watch,
_watch_keepalive: (context_rx, goal_rx, attach_rx, todos_rx, plans_rx),
streamed_this_turn: std::sync::atomic::AtomicBool::new(false),
manual_compact_pending: std::sync::atomic::AtomicBool::new(false),
last_input_tokens: std::sync::atomic::AtomicU64::new(0),
compact_review_mode: Mutex::new(CompactReviewMode::default()),
compact_lock: std::sync::Arc::new(tokio::sync::Mutex::new(())),
last_image_user_msg: Mutex::new(None),
read_files: std::sync::Arc::new(
std::sync::Mutex::new(std::collections::HashSet::new()),
),
approval: std::sync::Arc::new(ApprovalRegistry::new()),
compact_reviews: std::sync::Arc::new(CompactReviewRegistry::new()),
forms: std::sync::Arc::new(FormRegistry::new()),
fs_access_mode: Mutex::new(None),
project_index: None,
}
}
pub fn project_index(&self) -> Option<std::sync::Arc<crate::index::AnchorIndex>> {
self.project_index.clone()
}
pub fn approval(&self) -> std::sync::Arc<ApprovalRegistry> {
self.approval.clone()
}
pub fn compact_reviews(&self) -> std::sync::Arc<CompactReviewRegistry> {
self.compact_reviews.clone()
}
pub fn forms(&self) -> std::sync::Arc<FormRegistry> {
self.forms.clone()
}
pub fn fs_access_mode(&self) -> Option<crate::fs_access::FsAccessMode> {
*self.fs_access_mode.lock().unwrap()
}
pub fn set_fs_access_mode(&self, mode: crate::fs_access::FsAccessMode) {
*self.fs_access_mode.lock().unwrap() = Some(mode);
}
pub fn compact_review_mode(&self) -> CompactReviewMode {
*self.compact_review_mode.lock().unwrap()
}
pub fn set_compact_review_mode(&self, mode: CompactReviewMode) {
*self.compact_review_mode.lock().unwrap() = mode;
}
pub fn read_files(
&self,
) -> std::sync::Arc<std::sync::Mutex<std::collections::HashSet<std::path::PathBuf>>> {
self.read_files.clone()
}
pub fn mark_file_read(&self, path: &std::path::Path) {
if let Ok(mut set) = self.read_files.lock() {
set.insert(path.to_path_buf());
if let Ok(canonical) = std::fs::canonicalize(path) {
set.insert(canonical);
}
}
}
pub fn stream_tx(&self) -> broadcast::Sender<StreamFrame> {
self.stream_tx.clone()
}
pub fn stream_subscribe(&self) -> broadcast::Receiver<StreamFrame> {
self.stream_tx.subscribe()
}
pub fn id(&self) -> &SessionId {
&self.id
}
pub fn dir(&self) -> &Path {
&self.dir
}
pub fn transcript_replay(&self) -> Vec<TranscriptEntry> {
let Some(path) = self.events_path() else {
return Vec::new();
};
replay_transcript_from(path).unwrap_or_default()
}
pub fn events_path(&self) -> Option<&Path> {
self.writer.as_ref().map(|w| w.events_path())
}
pub async fn plan_system_prompt(&self) -> Option<String> {
let store = crate::memory::plan::PlanStore::at(&self.dir);
let plan = store.latest().await.ok().flatten()?;
Some(crate::tools::plan::render_plan(&plan))
}
pub fn goal(&self) -> Option<String> {
if let Some(cached) = self.goal_watch.borrow().clone() {
return Some(cached);
}
load_goal(&self.dir)
}
pub fn subscribe_goal(&self) -> watch::Receiver<Option<String>> {
self.goal_watch.subscribe()
}
pub fn goal_watch(&self) -> &watch::Sender<Option<String>> {
&self.goal_watch
}
pub fn subscribe_context(&self) -> watch::Receiver<ContextSnapshot> {
self.context_watch.subscribe()
}
pub fn subscribe_attach(&self) -> watch::Receiver<usize> {
self.attach_watch.subscribe()
}
pub fn subscribe_pending_approvals(&self) -> watch::Receiver<Vec<PendingApproval>> {
self.approval.subscribe()
}
pub fn meta(&self) -> Option<crate::session_meta::SessionMeta> {
crate::session_meta::SessionMeta::load(&self.dir)
}
pub fn request_manual_compact(&self) {
self.manual_compact_pending
.store(true, std::sync::atomic::Ordering::SeqCst);
}
pub fn take_manual_compact_request(&self) -> bool {
self.manual_compact_pending
.swap(false, std::sync::atomic::Ordering::SeqCst)
}
pub fn set_goal(&self, goal: Option<String>) {
let _ = self.goal_watch.send(goal);
}
pub fn set_attach_count(&self, count: usize) {
let _ = self.attach_watch.send(count);
}
#[allow(clippy::too_many_arguments)]
pub fn record_llm_call(
&self,
model: &str,
tokens_in: u64,
tokens_out: u64,
cache_read: u64,
cache_write: u64,
ttft_ms: Option<u64>,
tokens_per_sec: Option<f64>,
) {
self.last_input_tokens
.store(tokens_in, std::sync::atomic::Ordering::Relaxed);
self.context_watch.send_modify(|snap| {
snap.model = model.to_string();
snap.tokens_in = snap.tokens_in.saturating_add(tokens_in);
snap.tokens_out = snap.tokens_out.saturating_add(tokens_out);
snap.cache_read = snap.cache_read.saturating_add(cache_read);
snap.cache_write = snap.cache_write.saturating_add(cache_write);
snap.last_ttft_ms = ttft_ms.unwrap_or(0);
snap.last_tokens_per_sec = tokens_per_sec.unwrap_or(0.0);
});
self.refresh_window_snapshot();
}
pub fn last_input_tokens(&self) -> u64 {
self.last_input_tokens
.load(std::sync::atomic::Ordering::Relaxed)
}
pub async fn acquire_compact_lock(&self) -> tokio::sync::MutexGuard<'_, ()> {
self.compact_lock.lock().await
}
pub async fn acquire_compact_lock_owned(&self) -> tokio::sync::OwnedMutexGuard<()> {
self.compact_lock.clone().lock_owned().await
}
pub fn compact_lock_handle(&self) -> std::sync::Arc<tokio::sync::Mutex<()>> {
self.compact_lock.clone()
}
fn clear_last_input_tokens(&self) {
self.last_input_tokens
.store(0, std::sync::atomic::Ordering::Relaxed);
}
pub fn refresh_window_snapshot(&self) {
let provider_tokens = self.last_input_tokens();
let estimated = crate::compaction::estimate_tokens_for_messages(&self.messages());
let window = if provider_tokens > 0 {
provider_tokens
} else {
estimated
};
let budget = crate::model_registry::model_info(&self.last_model()).context_budget;
self.context_watch.send_modify(|snap| {
snap.window_tokens = window;
snap.window_budget = budget;
});
}
pub fn cumulative_input_tokens(&self) -> u64 {
self.context_watch.borrow().tokens_in
}
pub fn reset_input_tokens_to(&self, tokens: u64) {
self.context_watch.send_modify(|snap| {
snap.tokens_in = tokens;
});
}
pub fn last_model(&self) -> String {
self.context_watch.borrow().model.clone()
}
pub fn set_mcp_totals(&self, ok: u16, total: u16) {
self.context_watch.send_modify(|snap| {
snap.mcp_ok = ok;
snap.mcp_total = total;
});
}
pub fn set_memory_recent_count(&self, count: u16) {
self.context_watch.send_modify(|snap| {
snap.memory_recent_count = count;
});
}
pub fn subscribe_todos(&self) -> watch::Receiver<Vec<crate::memory::todo::Todo>> {
self.todos_watch.subscribe()
}
pub fn todos_watch(&self) -> &watch::Sender<Vec<crate::memory::todo::Todo>> {
&self.todos_watch
}
pub fn subscribe_plans(&self) -> watch::Receiver<Vec<crate::memory::plan::Plan>> {
self.plans_watch.subscribe()
}
pub fn plans_watch(&self) -> &watch::Sender<Vec<crate::memory::plan::Plan>> {
&self.plans_watch
}
pub async fn refresh_plans_from_store_async(&self) {
if self.dir.as_os_str().is_empty() {
return;
}
let store = crate::memory::plan::PlanStore::at(&self.dir);
match store.list().await {
Ok(list) => {
let _ = self.plans_watch.send(list);
}
Err(e) => {
eprintln!("[atman] refresh_plans_from_store_async: {e}");
}
}
}
pub fn refresh_todos_from_store(&self) {
if self.dir.as_os_str().is_empty() {
return;
}
let store = crate::memory::todo::TodoStore::at(&self.dir);
match tokio::task::block_in_place(|| {
tokio::runtime::Handle::try_current()
.ok()
.map(|h| h.block_on(store.list()))
}) {
Some(Ok(list)) => {
let _ = self.todos_watch.send(list);
}
Some(Err(e)) => {
eprintln!("[atman] refresh_todos_from_store: {e}");
}
None => {}
}
}
pub async fn refresh_todos_from_store_async(&self) {
if self.dir.as_os_str().is_empty() {
return;
}
let store = crate::memory::todo::TodoStore::at(&self.dir);
match store.list().await {
Ok(list) => {
let _ = self.todos_watch.send(list);
}
Err(e) => {
eprintln!("[atman] refresh_todos_from_store_async: {e}");
}
}
}
pub fn sink(&self) -> &EventSink {
&self.sink
}
pub fn append_message(&self, msg: Message, flow_run_id: Option<FlowRunId>) {
let ts = chrono::Utc::now();
let flow_run_id_str = flow_run_id.as_ref().map(|r| r.0.to_string());
let event = match msg.role {
MessageRole::User => Event::UserMsg {
seq: 0,
turn_id: msg.turn_id.clone(),
message: msg.clone(),
ts,
},
MessageRole::Assistant => {
let _ = self
.stream_tx
.send(crate::stream::StreamFrame::AssistantMsg {
flow_run_id: flow_run_id_str.clone(),
message: msg.clone(),
});
Event::AssistantMsg {
seq: 0,
turn_id: msg.turn_id.clone(),
flow_run_id,
message: msg.clone(),
ts,
}
}
MessageRole::Tool => {
let _ = self
.stream_tx
.send(crate::stream::StreamFrame::ToolResultMsg {
flow_run_id: flow_run_id_str.clone(),
message: msg.clone(),
});
Event::ToolResultMsg {
seq: 0,
turn_id: msg.turn_id.clone(),
flow_run_id,
message: msg.clone(),
ts,
}
}
MessageRole::System => Event::SystemMsg {
seq: 0,
turn_id: msg.turn_id.clone(),
message: msg.clone(),
ts,
},
};
let seq = self.sink.emit_returning_seq(event);
if matches!(msg.role, MessageRole::User) {
let images: Vec<(usize, String)> = msg
.parts
.iter()
.enumerate()
.filter_map(|(i, p)| match p {
crate::message::MessagePart::Image { source } => {
let basename = match &source.data {
crate::message::ImageData::Path { path } => path
.file_name()
.and_then(|n| n.to_str())
.unwrap_or("unknown")
.to_string(),
crate::message::ImageData::Base64 { .. } => "base64".into(),
};
Some((i, basename))
}
_ => None,
})
.collect();
if !images.is_empty() {
*self.last_image_user_msg.lock().unwrap() = Some(LastImageUserMsg {
message_seq: seq,
images,
});
}
}
self.messages.lock().unwrap().push(msg);
}
pub fn emit_attachment_degrade(
&self,
message_seq: u64,
part_index: usize,
file_basename: String,
reason: String,
) {
self.sink.emit(Event::AttachmentDegraded {
seq: 0,
turn_id: None,
flow_run_id: None,
message_seq,
part_index,
file_basename,
reason,
ts: chrono::Utc::now(),
});
}
pub fn record_attachment_degrade(&self, reason: &str) -> usize {
let target = self.last_image_user_msg.lock().unwrap().take();
let Some(entry) = target else {
return 0;
};
let turn_id = self.current_turn.lock().unwrap().clone();
let now = chrono::Utc::now();
for (part_index, basename) in &entry.images {
self.sink.emit(Event::AttachmentDegraded {
seq: 0,
turn_id: turn_id.clone(),
flow_run_id: None,
message_seq: entry.message_seq,
part_index: *part_index,
file_basename: basename.clone(),
reason: reason.into(),
ts: now,
});
}
if let Ok(mut msgs) = self.messages.lock() {
for m in msgs.iter_mut() {
for (part_index, basename) in &entry.images {
if let Some(part) = m.parts.get_mut(*part_index)
&& matches!(part, crate::message::MessagePart::Image { .. })
{
*part = crate::message::MessagePart::Text {
text: format!("[attachment unavailable: {basename} — {reason}]"),
};
}
}
}
}
entry.images.len()
}
pub fn messages(&self) -> Vec<Message> {
self.messages.lock().unwrap().clone()
}
pub fn messages_handle(&self) -> std::sync::Arc<std::sync::Mutex<Vec<Message>>> {
self.messages.clone()
}
pub fn message_count(&self) -> usize {
self.messages.lock().unwrap().len()
}
pub fn user_message_count(&self) -> usize {
self.messages
.lock()
.unwrap()
.iter()
.filter(|m| matches!(m.role, MessageRole::User))
.count()
}
pub fn push_system_note(&self, text: String) {
let _ = self.stream_tx.send(crate::stream::StreamFrame::Note(text));
}
pub fn approval_cooldown_ok_for_compact(&self) -> bool {
self.sink.last_compact_ago_seconds().is_none_or(|s| s >= 60)
}
pub fn emit_compact_warning(
&self,
model: &str,
current_tokens: u64,
threshold: u64,
budget: u64,
reason: &str,
) {
let message = format!(
"context {current_tokens} > threshold {threshold} (budget {budget}, model {model}); skipping compaction: {reason}"
);
self.sink.emit(Event::WatchWarn {
seq: 0,
turn_id: self.current_turn.lock().unwrap().clone(),
flow_run_id: None,
target: "context.compaction".into(),
trigger: "auto_compact".into(),
message,
ts: chrono::Utc::now(),
});
self.push_system_note(format!("[warn] compaction skipped: {reason}"));
}
pub fn compact_messages(&self, summary: String) -> Option<CompactResult> {
use crate::compaction::{
estimate_tokens_for_messages, find_compact_range, replace_range_with_summary,
};
let mut guard = self.messages.lock().unwrap();
let before = guard.clone();
let info = crate::model_registry::model_info(&self.last_model());
let threshold = info.compact_threshold_tokens();
let range = find_compact_range(&before, threshold)?;
let turn_id = before
.get(range.start)
.map(|m| m.turn_id.clone())
.unwrap_or_else(TurnId::now);
let before_tokens = estimate_tokens_for_messages(&before);
let after = replace_range_with_summary(&before, &range, summary.clone(), turn_id.clone());
let after_tokens = estimate_tokens_for_messages(&after);
if after_tokens >= before_tokens {
drop(guard);
self.push_system_note(format!(
"compaction skipped: summary would not shrink transcript ({} >= {} tokens)",
after_tokens, before_tokens
));
return None;
}
let replacement_msg = after.get(range.start).cloned().unwrap_or_else(|| {
Message::system_compact_summary(
turn_id.clone(),
summary.clone(),
range.start as u64,
range.end.saturating_sub(1) as u64,
range.end - range.start,
)
});
*guard = after;
drop(guard);
self.sink.mark_compacted();
let replacement_seq = self.sink.next_seq_peek();
let ts = chrono::Utc::now();
self.sink.emit(Event::SystemMsg {
seq: 0,
turn_id: turn_id.clone(),
message: replacement_msg,
ts,
});
self.sink.emit(Event::ContextCompact {
seq: 0,
session_id: self.id.to_string(),
before_tokens,
after_tokens,
compacted_range_start: range.start as u64,
compacted_range_end: range.end.saturating_sub(1) as u64,
summary_text: Some(summary.clone()),
replacement_msg_seq: Some(replacement_seq),
ts,
});
self.sink.emit(Event::CompactionSummary {
seq: 0,
session_id: self.id.to_string(),
range_start: range.start as u64,
range_end: range.end.saturating_sub(1) as u64,
compacted_count: range.end - range.start,
before_tokens,
after_tokens,
summary: summary.clone(),
ts,
});
let _ = self
.stream_tx
.send(crate::stream::StreamFrame::CompactionSummary {
phase: crate::stream::CompactionPhase::Finished,
range_start: range.start,
range_end: range.end.saturating_sub(1),
summary,
before_tokens,
after_tokens,
compacted_count: range.end - range.start,
});
self.clear_last_input_tokens();
self.refresh_window_snapshot();
let checkpoint_messages = self.messages.lock().unwrap().clone();
let window_tokens = estimate_tokens_for_messages(&checkpoint_messages);
self.sink.emit(Event::Checkpoint {
seq: 0,
session_id: self.id.to_string(),
messages: checkpoint_messages,
window_tokens,
ts: chrono::Utc::now(),
});
Some(CompactResult {
before_tokens,
after_tokens,
compacted_start: range.start,
compacted_end: range.end,
})
}
pub fn begin_turn(&self, user_msg: Message) -> TurnId {
let turn_id = user_msg.turn_id.clone();
*self.current_turn.lock().unwrap() = Some(turn_id.clone());
*self.flow_cancel.lock().unwrap() = CancellationToken::new();
self.sink.emit(Event::TurnStart {
seq: 0,
turn_id: turn_id.clone(),
ts: chrono::Utc::now(),
});
self.append_message(user_msg, None);
turn_id
}
pub fn mark_streamed(&self) {
self.streamed_this_turn
.store(true, std::sync::atomic::Ordering::Relaxed);
}
pub fn take_streamed_flag(&self) -> bool {
self.streamed_this_turn
.swap(false, std::sync::atomic::Ordering::Relaxed)
}
pub fn end_turn(&self) {
self.streamed_this_turn
.store(false, std::sync::atomic::Ordering::Relaxed);
let turn_id = self.current_turn.lock().unwrap().take();
if let Some(turn_id) = turn_id {
let now = chrono::Utc::now();
let mut q = self.injection_queue.lock().unwrap();
for inj in q.iter_mut() {
if inj.state == InjectionState::Pending && inj.turn_id == turn_id {
inj.state = InjectionState::Cancelled;
let _ = self.injection_tx.send(inj.clone());
}
}
drop(q);
self.sink.emit(Event::TurnEnd {
seq: 0,
turn_id,
ts: now,
});
}
}
pub fn current_turn(&self) -> Option<TurnId> {
self.current_turn.lock().unwrap().clone()
}
pub fn enqueue_injection(&self, text: impl Into<String>) -> Result<InjectionId, EnqueueError> {
self.enqueue_injection_with_level(text, crate::injection::InjectionLevel::L1Nudge, None)
}
pub fn enqueue_injection_with_level(
&self,
text: impl Into<String>,
level: crate::injection::InjectionLevel,
redirect_target: Option<String>,
) -> Result<InjectionId, EnqueueError> {
let turn_id = self
.current_turn
.lock()
.unwrap()
.clone()
.ok_or(EnqueueError::NoActiveTurn)?;
let inj = Injection::with_level(turn_id.clone(), text, level, redirect_target);
let id = inj.id.clone();
self.sink.emit(Event::UserInject {
seq: 0,
turn_id,
injection: inj.clone(),
ts: inj.created_at,
});
self.injection_queue.lock().unwrap().push(inj.clone());
let _ = self.injection_tx.send(inj);
Ok(id)
}
pub fn subscribe_injections(&self) -> broadcast::Receiver<Injection> {
self.injection_tx.subscribe()
}
pub fn mark_injection_consumed(&self, id: &InjectionId) {
let mut q = self.injection_queue.lock().unwrap();
for inj in q.iter_mut() {
if inj.id == *id && inj.state == InjectionState::Pending {
inj.state = InjectionState::Injected;
let _ = self.injection_tx.send(inj.clone());
return;
}
}
}
pub fn peek_pending_l2_or_higher(&self, turn_id: &TurnId) -> Option<Injection> {
let q = self.injection_queue.lock().unwrap();
q.iter()
.find(|i| {
i.state == InjectionState::Pending
&& i.turn_id == *turn_id
&& !matches!(i.level, crate::injection::InjectionLevel::L1Nudge)
})
.cloned()
}
pub fn drain_injections(&self, turn_id: &TurnId) -> Vec<Injection> {
let mut q = self.injection_queue.lock().unwrap();
let mut out = Vec::new();
for inj in q.iter_mut() {
if inj.state == InjectionState::Pending && inj.turn_id == *turn_id {
inj.state = InjectionState::Injected;
let _ = self.injection_tx.send(inj.clone());
out.push(inj.clone());
}
}
out
}
pub fn list_pending_injections(&self) -> Vec<Injection> {
self.injection_queue
.lock()
.unwrap()
.iter()
.filter(|i| i.state == InjectionState::Pending)
.cloned()
.collect()
}
pub fn cancel_flow(&self) {
self.flow_cancel.lock().unwrap().cancel();
}
pub fn flow_cancel_token(&self) -> CancellationToken {
self.flow_cancel.lock().unwrap().clone()
}
pub async fn shutdown(mut self) {
if let Some(writer) = self.writer.take() {
writer.shutdown().await;
}
}
pub async fn flush_writer(&self) {
let Some(writer) = self.writer.as_ref() else {
return;
};
writer.flush().await;
}
}
#[derive(Debug, thiserror::Error)]
pub enum EnqueueError {
#[error("enqueue_injection called with no active turn")]
NoActiveTurn,
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::TempDir;
fn write_events(dir: &Path, lines: &[&str]) {
let path = dir.join("events.jsonl");
std::fs::write(&path, lines.join("\n") + "\n").unwrap();
}
#[test]
fn replay_applies_attachment_degraded_patch() {
let dir = TempDir::new().unwrap();
let user_msg = r#"{"type":"user_msg","seq":5,"turn_id":"019f0000-0000-7000-0000-000000000001","message":{"role":"user","parts":[{"type":"image","source":{"media_type":"image/png","data":{"kind":"path","path":"/tmp/photo.png"}}},{"type":"text","text":"describe"}],"turn_id":"019f0000-0000-7000-0000-000000000001"},"ts":"2026-07-07T00:00:00Z"}"#;
let degrade = r#"{"type":"attachment_degraded","seq":6,"turn_id":null,"flow_run_id":null,"message_seq":5,"part_index":0,"file_basename":"photo.png","reason":"image_too_large","ts":"2026-07-07T00:00:01Z"}"#;
write_events(dir.path(), &[user_msg, degrade]);
let entries = replay_transcript_from(&dir.path().join("events.jsonl")).unwrap();
let msg = entries
.into_iter()
.find_map(|e| match e {
TranscriptEntry::Message { message, .. } => Some(message),
_ => None,
})
.unwrap();
assert_eq!(msg.parts.len(), 2);
match &msg.parts[0] {
crate::message::MessagePart::Text { text } => {
assert!(text.contains("photo.png"), "expected basename: {text}");
assert!(text.contains("image_too_large"), "expected reason: {text}");
assert!(text.starts_with("[attachment unavailable"));
}
other => panic!("expected Text stub, got {other:?}"),
}
assert!(matches!(
msg.parts[1],
crate::message::MessagePart::Text { .. }
));
}
#[test]
fn approval_registry_auto_approves_when_level_leq_ceiling() {
let reg = ApprovalRegistry::new();
reg.set_auto_ceiling(crate::tool::ApprovalLevel::Approve);
let pending = PendingApproval {
tool_use_id: "tu1".into(),
tool_name: "fs.read".into(),
args_preview: "{}".into(),
preview: None,
level: crate::tool::ApprovalLevel::Auto,
run_id: FlowRunId::now(),
emitted_at: chrono::Utc::now(),
bypass_auto_ceiling: false,
};
let rx = reg.request(pending);
let got = rx.blocking_recv().unwrap();
assert!(matches!(got, ApprovalDecision::Approve));
assert!(reg.list_pending().is_empty());
}
#[test]
fn approval_registry_queues_when_level_above_ceiling() {
let reg = std::sync::Arc::new(ApprovalRegistry::new());
reg.set_auto_ceiling(crate::tool::ApprovalLevel::Auto);
let pending = PendingApproval {
tool_use_id: "tu42".into(),
tool_name: "fs.write".into(),
args_preview: "{}".into(),
preview: None,
level: crate::tool::ApprovalLevel::Approve,
run_id: FlowRunId::now(),
emitted_at: chrono::Utc::now(),
bypass_auto_ceiling: false,
};
let mut rx = reg.request(pending);
assert_eq!(reg.list_pending().len(), 1);
assert!(rx.try_recv().is_err(), "should still be queued");
assert!(reg.decide("tu42", ApprovalDecision::Approve));
let got = rx.blocking_recv().unwrap();
assert!(matches!(got, ApprovalDecision::Approve));
assert!(reg.list_pending().is_empty());
}
#[test]
fn approval_registry_decide_all_flushes_queue() {
let reg = ApprovalRegistry::new();
reg.set_auto_ceiling(crate::tool::ApprovalLevel::Auto);
let mut rxs = Vec::new();
for i in 0..3 {
rxs.push(reg.request(PendingApproval {
tool_use_id: format!("tu{i}"),
tool_name: "bash.exec".into(),
args_preview: "{}".into(),
preview: None,
level: crate::tool::ApprovalLevel::Dangerous,
run_id: FlowRunId::now(),
emitted_at: chrono::Utc::now(),
bypass_auto_ceiling: false,
}));
}
assert_eq!(reg.list_pending().len(), 3);
assert_eq!(
reg.decide_all(ApprovalDecision::Deny {
reason: "user cancelled".into()
}),
3
);
assert!(reg.list_pending().is_empty());
}
#[test]
fn compact_review_registry_auto_accepts_when_no_subscriber() {
let reg = CompactReviewRegistry::new();
let pending = PendingCompactReview {
review_id: "r1".into(),
summary: "gist".into(),
slice_preview: String::new(),
slice_count: 0,
range_start: 0,
range_end: 0,
tokens_before: 0,
emitted_at: chrono::Utc::now(),
};
let rx = reg.request(pending);
let got = rx.blocking_recv().unwrap();
assert!(matches!(got, CompactReviewDecision::AcceptAsIs));
assert!(reg.list_pending().is_none());
}
#[test]
fn compact_review_registry_holds_pending_and_decides() {
let reg = std::sync::Arc::new(CompactReviewRegistry::new());
let _sub = reg.subscribe();
let pending = PendingCompactReview {
review_id: "r2".into(),
summary: "old".into(),
slice_preview: "slice".into(),
slice_count: 3,
range_start: 1,
range_end: 4,
tokens_before: 500,
emitted_at: chrono::Utc::now(),
};
let mut rx = reg.request(pending);
assert!(rx.try_recv().is_err(), "should be queued");
assert!(reg.list_pending().is_some());
assert!(reg.decide(
"r2",
CompactReviewDecision::AcceptEdited {
summary: "new".into()
}
));
let got = rx.blocking_recv().unwrap();
match got {
CompactReviewDecision::AcceptEdited { summary } => assert_eq!(summary, "new"),
other => panic!("unexpected decision: {other:?}"),
}
assert!(reg.list_pending().is_none());
}
#[test]
fn compact_review_registry_reject_flushes() {
let reg = std::sync::Arc::new(CompactReviewRegistry::new());
let _sub = reg.subscribe();
let rx = reg.request(PendingCompactReview {
review_id: "r3".into(),
summary: String::new(),
slice_preview: String::new(),
slice_count: 0,
range_start: 0,
range_end: 0,
tokens_before: 0,
emitted_at: chrono::Utc::now(),
});
assert!(reg.decide("r3", CompactReviewDecision::Reject));
let got = rx.blocking_recv().unwrap();
assert!(matches!(got, CompactReviewDecision::Reject));
}
#[test]
fn replay_context_snapshot_accumulates_llm_call_usage() {
let dir = TempDir::new().unwrap();
let events = [
r#"{"type":"llm_call","seq":1,"model":"anthropic/claude-4","provider":"anthropic","usage":{"input":100,"cached_input":10,"output":50,"cache_write":0},"wallclock_ms":1000,"status":{"kind":"ok"},"ts":"2026-07-08T00:00:00Z"}"#,
r#"{"type":"user_msg","seq":2,"turn_id":"019f0000-0000-7000-0000-000000000002","message":{"role":"user","parts":[{"type":"text","text":"hi"}],"turn_id":"019f0000-0000-7000-0000-000000000002"},"ts":"2026-07-08T00:00:00Z"}"#,
r#"{"type":"llm_call","seq":3,"model":"anthropic/claude-4","provider":"anthropic","usage":{"input":200,"cached_input":0,"output":80,"cache_write":0},"wallclock_ms":1000,"status":{"kind":"ok"},"ts":"2026-07-08T00:00:01Z"}"#,
];
write_events(dir.path(), &events);
let snap = replay_context_snapshot_from(&dir.path().join("events.jsonl"));
assert_eq!(snap.model, "anthropic/claude-4");
assert_eq!(snap.tokens_in, 310);
assert_eq!(snap.tokens_out, 130);
}
#[test]
fn compact_review_mode_parses_all_variants() {
assert_eq!(
CompactReviewMode::parse("always"),
Some(CompactReviewMode::Always)
);
assert_eq!(
CompactReviewMode::parse("manual-only"),
Some(CompactReviewMode::ManualOnly)
);
assert_eq!(
CompactReviewMode::parse("manual_only"),
Some(CompactReviewMode::ManualOnly)
);
assert_eq!(
CompactReviewMode::parse("never"),
Some(CompactReviewMode::Never)
);
assert_eq!(CompactReviewMode::parse(" bogus "), None);
}
#[test]
fn compact_review_mode_should_review_matrix() {
assert!(CompactReviewMode::Always.should_review(false));
assert!(CompactReviewMode::Always.should_review(true));
assert!(!CompactReviewMode::ManualOnly.should_review(false));
assert!(CompactReviewMode::ManualOnly.should_review(true));
assert!(!CompactReviewMode::Never.should_review(false));
assert!(!CompactReviewMode::Never.should_review(true));
}
#[test]
fn compact_review_registry_new_request_rejects_previous() {
let reg = std::sync::Arc::new(CompactReviewRegistry::new());
let _sub = reg.subscribe();
let rx_a = reg.request(PendingCompactReview {
review_id: "rA".into(),
summary: String::new(),
slice_preview: String::new(),
slice_count: 0,
range_start: 0,
range_end: 0,
tokens_before: 0,
emitted_at: chrono::Utc::now(),
});
let _rx_b = reg.request(PendingCompactReview {
review_id: "rB".into(),
summary: String::new(),
slice_preview: String::new(),
slice_count: 0,
range_start: 0,
range_end: 0,
tokens_before: 0,
emitted_at: chrono::Utc::now(),
});
let got = rx_a.blocking_recv().unwrap();
assert!(matches!(got, CompactReviewDecision::Reject));
}
fn mk_form(form_id: &str, prompt: &str) -> crate::form::PendingForm {
crate::form::PendingForm {
form_id: form_id.into(),
run_id: crate::event::FlowRunId::now(),
tool_use_id: "tu".into(),
kind: crate::form::FormKind::Confirm {
prompt: prompt.into(),
},
emitted_at: chrono::Utc::now(),
}
}
#[test]
fn form_registry_auto_cancels_without_subscriber() {
let reg = FormRegistry::new();
let rx = reg.request(mk_form("f1", "sure?"));
let got = rx.blocking_recv().unwrap();
assert_eq!(got, crate::form::FormAnswer::Cancelled);
assert!(reg.list_pending().is_empty());
}
#[test]
fn form_registry_delivers_answer_by_form_id() {
let reg = std::sync::Arc::new(FormRegistry::new());
let _sub = reg.subscribe();
let rx = reg.request(mk_form("fA", "?"));
assert_eq!(reg.list_pending().len(), 1);
let ok = reg.submit("fA", crate::form::FormAnswer::Confirmed { value: true });
assert!(ok);
let got = rx.blocking_recv().unwrap();
assert_eq!(got, crate::form::FormAnswer::Confirmed { value: true });
assert!(reg.list_pending().is_empty());
}
#[test]
fn form_registry_submit_unknown_id_is_noop() {
let reg = std::sync::Arc::new(FormRegistry::new());
let _sub = reg.subscribe();
let _rx = reg.request(mk_form("real", "?"));
assert!(!reg.submit("ghost", crate::form::FormAnswer::Cancelled));
assert_eq!(reg.list_pending().len(), 1);
}
#[test]
fn form_registry_cancel_all_flushes_pending() {
let reg = std::sync::Arc::new(FormRegistry::new());
let _sub = reg.subscribe();
let rx_a = reg.request(mk_form("a", "?"));
let rx_b = reg.request(mk_form("b", "?"));
reg.cancel_all();
assert_eq!(
rx_a.blocking_recv().unwrap(),
crate::form::FormAnswer::Cancelled
);
assert_eq!(
rx_b.blocking_recv().unwrap(),
crate::form::FormAnswer::Cancelled
);
assert!(reg.list_pending().is_empty());
}
#[test]
fn form_registry_queues_multiple_pending() {
let reg = std::sync::Arc::new(FormRegistry::new());
let _sub = reg.subscribe();
let _rx1 = reg.request(mk_form("1", "?"));
let _rx2 = reg.request(mk_form("2", "?"));
let pending = reg.list_pending();
assert_eq!(pending.len(), 2);
assert_eq!(pending[0].form_id, "1");
assert_eq!(pending[1].form_id, "2");
}
#[test]
fn replay_without_degraded_events_preserves_image_parts() {
let dir = TempDir::new().unwrap();
let user_msg = r#"{"type":"user_msg","seq":1,"turn_id":"019f0000-0000-7000-0000-000000000002","message":{"role":"user","parts":[{"type":"image","source":{"media_type":"image/png","data":{"kind":"path","path":"/tmp/x.png"}}}],"turn_id":"019f0000-0000-7000-0000-000000000002"},"ts":"2026-07-07T00:00:00Z"}"#;
write_events(dir.path(), &[user_msg]);
let entries = replay_transcript_from(&dir.path().join("events.jsonl")).unwrap();
let msg = entries
.into_iter()
.find_map(|e| match e {
TranscriptEntry::Message { message, .. } => Some(message),
_ => None,
})
.unwrap();
assert!(matches!(
msg.parts[0],
crate::message::MessagePart::Image { .. }
));
}
}