Skip to main content

atman_runtime/
session.rs

1use std::path::{Path, PathBuf};
2use std::sync::Mutex;
3
4use tokio::sync::{broadcast, watch};
5use tokio_util::sync::CancellationToken;
6use uuid::Uuid;
7
8use crate::event::{Event, EventSink, FlowRunId, TurnId};
9use crate::event_log::reader::{find_last_seq, replay_context_snapshot_from};
10use crate::event_writer::EventWriter;
11use crate::injection::{Injection, InjectionId, InjectionState};
12use crate::message::{Message, MessageRole};
13use crate::projection::message_window::{
14    TranscriptEntry, replay_all_messages_with_seq, replay_messages_from, replay_messages_with_seq,
15    replay_transcript_from,
16};
17use crate::stream::StreamFrame;
18
19#[derive(Debug, Clone, PartialEq, Eq, Hash)]
20pub struct SessionId(pub Uuid);
21
22fn is_auto_name_threshold(mut count: u64) -> bool {
23    if count < 3 {
24        return false;
25    }
26    while count % 3 == 0 {
27        count /= 3;
28    }
29    count == 1
30}
31
32impl SessionId {
33    pub fn now() -> Self {
34        Self(Uuid::new_v4())
35    }
36
37    pub fn parse(s: &str) -> Result<Self, uuid::Error> {
38        Uuid::parse_str(s).map(Self)
39    }
40}
41
42impl std::fmt::Display for SessionId {
43    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
44        self.0.fmt(f)
45    }
46}
47
48type WatchKeepalive = (
49    watch::Receiver<ContextSnapshot>,
50    watch::Receiver<Option<String>>,
51    watch::Receiver<usize>,
52    watch::Receiver<Vec<crate::memory::todo::Todo>>,
53    watch::Receiver<Vec<crate::memory::plan::Plan>>,
54);
55
56#[derive(Debug)]
57pub struct TurnState {
58    pub current_turn: Mutex<Option<TurnId>>,
59    pub flow_cancel: Mutex<CancellationToken>,
60    pub streamed: std::sync::atomic::AtomicBool,
61}
62
63impl TurnState {
64    fn new() -> Self {
65        Self {
66            current_turn: Mutex::new(None),
67            flow_cancel: Mutex::new(CancellationToken::new()),
68            streamed: std::sync::atomic::AtomicBool::new(false),
69        }
70    }
71}
72
73pub struct WatchHub {
74    pub stream_tx: broadcast::Sender<StreamFrame>,
75    pub context: watch::Sender<ContextSnapshot>,
76    pub goal: watch::Sender<Option<String>>,
77    pub attach: watch::Sender<usize>,
78    pub todos: watch::Sender<Vec<crate::memory::todo::Todo>>,
79    pub plans: watch::Sender<Vec<crate::memory::plan::Plan>>,
80    _keepalive: WatchKeepalive,
81}
82
83pub struct CompactionState {
84    pub manual_pending: std::sync::atomic::AtomicBool,
85    pub last_input_tokens: std::sync::atomic::AtomicU64,
86    pub review_mode: Mutex<CompactReviewMode>,
87    pub lock: std::sync::Arc<tokio::sync::Mutex<()>>,
88}
89
90impl CompactionState {
91    fn new() -> Self {
92        Self {
93            manual_pending: std::sync::atomic::AtomicBool::new(false),
94            last_input_tokens: std::sync::atomic::AtomicU64::new(0),
95            review_mode: Mutex::new(CompactReviewMode::default()),
96            lock: std::sync::Arc::new(tokio::sync::Mutex::new(())),
97        }
98    }
99}
100
101pub struct InteractionServices {
102    pub approval: std::sync::Arc<ApprovalRegistry>,
103    pub compact_reviews: std::sync::Arc<CompactReviewRegistry>,
104    pub forms: std::sync::Arc<FormRegistry>,
105}
106
107impl InteractionServices {
108    fn new() -> Self {
109        Self {
110            approval: std::sync::Arc::new(ApprovalRegistry::new()),
111            compact_reviews: std::sync::Arc::new(CompactReviewRegistry::new()),
112            forms: std::sync::Arc::new(FormRegistry::new()),
113        }
114    }
115}
116
117pub struct Session {
118    id: SessionId,
119    dir: PathBuf,
120    writer: std::sync::Mutex<Option<EventWriter>>,
121    sink: EventSink,
122    message_stream: crate::message_stream::MessageStream,
123    messages: std::sync::Arc<std::sync::Mutex<Vec<Message>>>,
124    pub turn: TurnState,
125    pub watch: WatchHub,
126    pub watch_hub: std::sync::Arc<crate::watch::WatchHub>,
127    pub flow_registry: std::sync::Arc<crate::tools::agent_ctrl::FlowRegistry>,
128    /// Handle of the current root FlowRun; set per turn.
129    current_root: std::sync::Mutex<Option<String>>,
130    successful_flow_count: std::sync::atomic::AtomicU64,
131    pub compaction: CompactionState,
132    pub interactions: InteractionServices,
133    injection_queue: Mutex<Vec<Injection>>,
134    injection_tx: broadcast::Sender<Injection>,
135    last_image_user_msg: Mutex<Option<LastImageUserMsg>>,
136    read_files: std::sync::Arc<std::sync::Mutex<std::collections::HashSet<std::path::PathBuf>>>,
137    output_store: std::sync::Arc<crate::tools::tool_output::OutputStore>,
138    fs_access_mode: Mutex<Option<crate::fs_access::FsAccessMode>>,
139    project_index: Option<std::sync::Arc<crate::index::AnchorIndex>>,
140}
141
142#[derive(Debug, Clone)]
143pub struct PendingCompactReview {
144    pub review_id: String,
145    pub summary: String,
146    pub slice_preview: String,
147    pub slice_count: usize,
148    pub range_start: usize,
149    pub range_end: usize,
150    pub tokens_before: u64,
151    pub emitted_at: chrono::DateTime<chrono::Utc>,
152}
153
154#[derive(Debug, Clone)]
155pub enum CompactReviewDecision {
156    AcceptAsIs,
157    AcceptEdited { summary: String },
158    Reject,
159}
160
161pub struct CompactReviewRegistry {
162    entry: std::sync::Mutex<Option<CompactReviewEntry>>,
163    watch_tx: watch::Sender<Option<PendingCompactReview>>,
164}
165
166struct CompactReviewEntry {
167    pending: PendingCompactReview,
168    responder: tokio::sync::oneshot::Sender<CompactReviewDecision>,
169}
170
171impl Default for CompactReviewRegistry {
172    fn default() -> Self {
173        Self::new()
174    }
175}
176
177impl CompactReviewRegistry {
178    pub fn new() -> Self {
179        let (watch_tx, _) = watch::channel(None);
180        Self {
181            entry: std::sync::Mutex::new(None),
182            watch_tx,
183        }
184    }
185
186    pub fn subscribe(&self) -> watch::Receiver<Option<PendingCompactReview>> {
187        self.watch_tx.subscribe()
188    }
189
190    pub fn list_pending(&self) -> Option<PendingCompactReview> {
191        self.entry
192            .lock()
193            .unwrap()
194            .as_ref()
195            .map(|e| e.pending.clone())
196    }
197
198    pub fn subscriber_count(&self) -> usize {
199        self.watch_tx.receiver_count()
200    }
201
202    pub fn request(
203        &self,
204        pending: PendingCompactReview,
205    ) -> tokio::sync::oneshot::Receiver<CompactReviewDecision> {
206        let (tx, rx) = tokio::sync::oneshot::channel();
207        if self.watch_tx.receiver_count() == 0 {
208            let _ = tx.send(CompactReviewDecision::AcceptAsIs);
209            return rx;
210        }
211        {
212            let mut slot = self.entry.lock().unwrap();
213            if let Some(prev) = slot.take() {
214                let _ = prev.responder.send(CompactReviewDecision::Reject);
215            }
216            *slot = Some(CompactReviewEntry {
217                pending: pending.clone(),
218                responder: tx,
219            });
220        }
221        let _ = self.watch_tx.send(Some(pending));
222        rx
223    }
224
225    pub fn decide(&self, review_id: &str, decision: CompactReviewDecision) -> bool {
226        let entry = {
227            let mut slot = self.entry.lock().unwrap();
228            match slot.as_ref() {
229                Some(e) if e.pending.review_id == review_id => slot.take(),
230                _ => None,
231            }
232        };
233        match entry {
234            Some(e) => {
235                let _ = e.responder.send(decision);
236                let _ = self.watch_tx.send(None);
237                true
238            }
239            None => false,
240        }
241    }
242}
243
244#[derive(Debug, Clone)]
245pub struct PendingApproval {
246    pub tool_use_id: String,
247    pub tool_name: String,
248    pub args_preview: String,
249    pub preview: Option<String>,
250    pub level: crate::tool::ApprovalLevel,
251    pub run_id: FlowRunId,
252    pub emitted_at: chrono::DateTime<chrono::Utc>,
253    pub bypass_auto_ceiling: bool,
254}
255
256#[derive(Debug, Clone)]
257pub enum ApprovalDecision {
258    Approve,
259    Deny { reason: String },
260}
261
262pub struct FormRegistry {
263    entries: std::sync::Mutex<Vec<FormEntry>>,
264    watch_tx: watch::Sender<Vec<crate::form::PendingForm>>,
265}
266
267struct FormEntry {
268    pending: crate::form::PendingForm,
269    responder: tokio::sync::oneshot::Sender<crate::form::FormAnswer>,
270}
271
272impl Default for FormRegistry {
273    fn default() -> Self {
274        Self::new()
275    }
276}
277
278impl FormRegistry {
279    pub fn new() -> Self {
280        let (watch_tx, _) = watch::channel(Vec::new());
281        Self {
282            entries: std::sync::Mutex::new(Vec::new()),
283            watch_tx,
284        }
285    }
286
287    pub fn subscribe(&self) -> watch::Receiver<Vec<crate::form::PendingForm>> {
288        self.watch_tx.subscribe()
289    }
290
291    pub fn list_pending(&self) -> Vec<crate::form::PendingForm> {
292        self.entries
293            .lock()
294            .unwrap()
295            .iter()
296            .map(|e| e.pending.clone())
297            .collect()
298    }
299
300    pub fn subscriber_count(&self) -> usize {
301        self.watch_tx.receiver_count()
302    }
303
304    // No TUI attached → auto-cancel so flows don't hang forever. Otherwise
305    // enqueue and hand a receiver back to the caller.
306    pub fn request(
307        &self,
308        pending: crate::form::PendingForm,
309    ) -> tokio::sync::oneshot::Receiver<crate::form::FormAnswer> {
310        let (tx, rx) = tokio::sync::oneshot::channel();
311        if self.watch_tx.receiver_count() == 0 {
312            let _ = tx.send(crate::form::FormAnswer::Cancelled);
313            return rx;
314        }
315        {
316            let mut entries = self.entries.lock().unwrap();
317            entries.push(FormEntry {
318                pending: pending.clone(),
319                responder: tx,
320            });
321        }
322        self.broadcast_snapshot();
323        rx
324    }
325
326    pub fn submit(&self, form_id: &str, answer: crate::form::FormAnswer) -> bool {
327        let entry = {
328            let mut entries = self.entries.lock().unwrap();
329            let pos = entries.iter().position(|e| e.pending.form_id == form_id);
330            pos.map(|p| entries.remove(p))
331        };
332        match entry {
333            Some(e) => {
334                let _ = e.responder.send(answer);
335                self.broadcast_snapshot();
336                true
337            }
338            None => false,
339        }
340    }
341
342    pub fn cancel_all(&self) {
343        let drained: Vec<FormEntry> = {
344            let mut entries = self.entries.lock().unwrap();
345            std::mem::take(&mut *entries)
346        };
347        for e in drained {
348            let _ = e.responder.send(crate::form::FormAnswer::Cancelled);
349        }
350        self.broadcast_snapshot();
351    }
352
353    pub fn promote(&self, form_id: &str) {
354        let mut entries = self.entries.lock().unwrap();
355        if let Some(pos) = entries.iter().position(|e| e.pending.form_id == form_id) {
356            if pos == 0 {
357                return;
358            }
359            let entry = entries.remove(pos);
360            entries.insert(0, entry);
361        }
362        drop(entries);
363        self.broadcast_snapshot();
364    }
365
366    fn broadcast_snapshot(&self) {
367        let snap = self
368            .entries
369            .lock()
370            .unwrap()
371            .iter()
372            .map(|e| e.pending.clone())
373            .collect();
374        let _ = self.watch_tx.send(snap);
375    }
376}
377
378pub struct ApprovalRegistry {
379    entries: std::sync::Mutex<Vec<ApprovalEntry>>,
380    auto_ceiling: std::sync::Mutex<crate::tool::ApprovalLevel>,
381    watch_tx: watch::Sender<Vec<PendingApproval>>,
382}
383
384struct ApprovalEntry {
385    pending: PendingApproval,
386    responder: tokio::sync::oneshot::Sender<ApprovalDecision>,
387}
388
389impl Default for ApprovalRegistry {
390    fn default() -> Self {
391        Self::new()
392    }
393}
394
395impl ApprovalRegistry {
396    pub fn new() -> Self {
397        let (watch_tx, _) = watch::channel(Vec::new());
398        Self {
399            entries: std::sync::Mutex::new(Vec::new()),
400            auto_ceiling: std::sync::Mutex::new(crate::tool::ApprovalLevel::Approve),
401            watch_tx,
402        }
403    }
404
405    pub fn subscribe(&self) -> watch::Receiver<Vec<PendingApproval>> {
406        self.watch_tx.subscribe()
407    }
408
409    pub fn list_pending(&self) -> Vec<PendingApproval> {
410        self.entries
411            .lock()
412            .unwrap()
413            .iter()
414            .map(|e| e.pending.clone())
415            .collect()
416    }
417
418    pub fn set_auto_ceiling(&self, level: crate::tool::ApprovalLevel) {
419        *self.auto_ceiling.lock().unwrap() = level;
420    }
421
422    pub fn request(
423        &self,
424        pending: PendingApproval,
425    ) -> tokio::sync::oneshot::Receiver<ApprovalDecision> {
426        let (tx, rx) = tokio::sync::oneshot::channel();
427        if !pending.bypass_auto_ceiling && pending.level <= *self.auto_ceiling.lock().unwrap() {
428            let _ = tx.send(ApprovalDecision::Approve);
429            return rx;
430        }
431        {
432            let mut entries = self.entries.lock().unwrap();
433            entries.push(ApprovalEntry {
434                pending,
435                responder: tx,
436            });
437        }
438        self.broadcast_snapshot();
439        rx
440    }
441
442    pub fn decide(&self, tool_use_id: &str, decision: ApprovalDecision) -> bool {
443        let mut entries = self.entries.lock().unwrap();
444        if let Some(pos) = entries
445            .iter()
446            .position(|e| e.pending.tool_use_id == tool_use_id)
447        {
448            let entry = entries.remove(pos);
449            let _ = entry.responder.send(decision);
450            drop(entries);
451            self.broadcast_snapshot();
452            true
453        } else {
454            false
455        }
456    }
457
458    pub fn decide_all(&self, decision: ApprovalDecision) -> usize {
459        let mut entries = self.entries.lock().unwrap();
460        let count = entries.len();
461        for entry in entries.drain(..) {
462            let _ = entry.responder.send(decision.clone());
463        }
464        drop(entries);
465        self.broadcast_snapshot();
466        count
467    }
468
469    fn broadcast_snapshot(&self) {
470        let snapshot = self
471            .entries
472            .lock()
473            .unwrap()
474            .iter()
475            .map(|e| e.pending.clone())
476            .collect();
477        let _ = self.watch_tx.send(snapshot);
478    }
479}
480type ImagePart = (usize, String);
481
482#[derive(Debug, Clone)]
483struct LastImageUserMsg {
484    message_seq: u64,
485    images: Vec<ImagePart>,
486}
487
488#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
489pub enum CompactReviewMode {
490    Always,
491    #[default]
492    ManualOnly,
493    Never,
494}
495
496impl CompactReviewMode {
497    pub fn parse(s: &str) -> Option<Self> {
498        match s.trim() {
499            "always" => Some(Self::Always),
500            "manual-only" | "manual_only" => Some(Self::ManualOnly),
501            "never" => Some(Self::Never),
502            _ => None,
503        }
504    }
505
506    pub fn should_review(self, forced: bool) -> bool {
507        match self {
508            Self::Always => true,
509            Self::ManualOnly => forced,
510            Self::Never => false,
511        }
512    }
513}
514
515#[derive(Debug, Clone, PartialEq, Eq)]
516pub struct CompactResult {
517    pub before_tokens: u64,
518    pub after_tokens: u64,
519    pub compacted_start: usize,
520    pub compacted_end: usize,
521}
522
523#[derive(Debug, Clone, Default, PartialEq)]
524pub struct ContextSnapshot {
525    pub model: String,
526    pub tokens_in: u64,
527    pub tokens_out: u64,
528    pub cost_usd: f64,
529    pub mcp_servers: Vec<crate::mcp::McpServerStatus>,
530    pub memory_recent_count: u16,
531    pub window_tokens: u64,
532    pub window_budget: u64,
533    pub cache_read: u64,
534    pub cache_write: u64,
535    pub last_ttft_ms: u64,
536    pub last_tokens_per_sec: f64,
537}
538
539#[derive(Debug, thiserror::Error)]
540pub enum SessionOpenError {
541    #[error("invalid session id `{sid}` (want a UUID)")]
542    InvalidId { sid: String },
543    #[error("session `{sid}` not found at {}", dir.display())]
544    NotFound { sid: String, dir: PathBuf },
545    #[error("session writer init: {0}")]
546    WriterInit(#[source] std::io::Error),
547    #[error("replay {}: {source}", path.display())]
548    Replay {
549        path: PathBuf,
550        #[source]
551        source: std::io::Error,
552    },
553}
554
555fn load_goal(dir: &Path) -> Option<String> {
556    if dir.as_os_str().is_empty() {
557        return None;
558    }
559    let store = crate::memory::goal::GoalStore::at(dir);
560    match store.get() {
561        Ok(s) if !s.is_empty() => Some(s),
562        _ => None,
563    }
564}
565
566#[derive(serde::Serialize, serde::Deserialize, Default)]
567struct PersistedContextState {
568    #[serde(default)]
569    model: String,
570    #[serde(default)]
571    window_tokens: u64,
572    #[serde(default)]
573    window_budget: u64,
574}
575
576impl PersistedContextState {
577    fn path(dir: &Path) -> PathBuf {
578        dir.join("context_state.json")
579    }
580
581    fn load(dir: &Path) -> Self {
582        match std::fs::read_to_string(Self::path(dir)) {
583            Ok(text) => serde_json::from_str(&text).unwrap_or_default(),
584            Err(_) => Self::default(),
585        }
586    }
587
588    fn save(&self, dir: &Path) {
589        if dir.as_os_str().is_empty() {
590            return;
591        }
592        if let Ok(json) = serde_json::to_string_pretty(self) {
593            let _ = std::fs::write(Self::path(dir), &json);
594        }
595    }
596}
597
598fn default_project_index(root: &Path) -> Option<std::sync::Arc<crate::index::AnchorIndex>> {
599    match crate::index::AnchorIndex::open_project(root) {
600        Ok(idx) => Some(std::sync::Arc::new(idx)),
601        Err(e) => {
602            crate::notify!(
603                warn,
604                "project index unavailable at {} — history search disabled: {e}",
605                root.display()
606            );
607            None
608        }
609    }
610}
611
612impl Session {
613    pub fn open(root: impl AsRef<Path>) -> std::io::Result<Self> {
614        Self::open_with_redactor(root, None)
615    }
616
617    pub fn open_with_redactor(
618        root: impl AsRef<Path>,
619        redactor: Option<std::sync::Arc<crate::redact::Redactor>>,
620    ) -> std::io::Result<Self> {
621        let root_ref = root.as_ref();
622        let project_index = default_project_index(root_ref);
623        Self::open_with_context(root_ref, redactor, project_index)
624    }
625
626    pub fn open_with_context(
627        root: impl AsRef<Path>,
628        redactor: Option<std::sync::Arc<crate::redact::Redactor>>,
629        project_index: Option<std::sync::Arc<crate::index::AnchorIndex>>,
630    ) -> std::io::Result<Self> {
631        let id = SessionId::now();
632        let dir = root.as_ref().join("sessions").join(id.to_string());
633        if let Some(ls) = crate::notify::log_sink() {
634            ls.set_session_id(Some(id.to_string()));
635        }
636        let writer = EventWriter::spawn_full(
637            &dir,
638            redactor.clone(),
639            project_index.clone(),
640            Some(id.to_string()),
641        )?;
642        if let Err(e) = crate::session_meta::SessionMeta::from_cwd().save(&dir) {
643            crate::notify!(error, "session meta write failed: {e}");
644        }
645        let mut sink = EventSink::new().with_forwarder(writer.sender());
646        if let Some(r) = redactor {
647            sink = sink.with_redactor(r);
648        }
649        let (injection_tx, _) = broadcast::channel(32);
650        let (stream_tx, _) = broadcast::channel(2048);
651        let (context_watch, context_rx) = watch::channel(ContextSnapshot::default());
652        let (goal_watch, goal_rx) = watch::channel(None);
653        let (attach_watch, attach_rx) = watch::channel(0);
654        let (todos_watch, todos_rx) = watch::channel(Vec::new());
655        let (plans_watch, plans_rx) = watch::channel(Vec::new());
656        let events_handle = sink.events_handle();
657        let output_store = std::sync::Arc::new(crate::tools::tool_output::OutputStore::at(&dir));
658        Ok(Self {
659            id,
660            dir,
661            writer: std::sync::Mutex::new(Some(writer)),
662            sink,
663            message_stream: crate::message_stream::MessageStream::new(events_handle),
664            messages: std::sync::Arc::new(std::sync::Mutex::new(Vec::new())),
665            output_store: output_store.clone(),
666            turn: TurnState::new(),
667            watch: WatchHub {
668                stream_tx,
669                context: context_watch,
670                goal: goal_watch,
671                attach: attach_watch,
672                todos: todos_watch,
673                plans: plans_watch,
674                _keepalive: (context_rx, goal_rx, attach_rx, todos_rx, plans_rx),
675            },
676            watch_hub: std::sync::Arc::new(crate::watch::WatchHub::new()),
677            flow_registry: std::sync::Arc::new(crate::tools::agent_ctrl::FlowRegistry::new()),
678            current_root: std::sync::Mutex::new(None),
679            successful_flow_count: std::sync::atomic::AtomicU64::new(0),
680            compaction: CompactionState::new(),
681            interactions: InteractionServices::new(),
682            injection_queue: Mutex::new(Vec::new()),
683            injection_tx,
684            last_image_user_msg: Mutex::new(None),
685            read_files: std::sync::Arc::new(
686                std::sync::Mutex::new(std::collections::HashSet::new()),
687            ),
688            fs_access_mode: Mutex::new(None),
689            project_index,
690        })
691    }
692
693    pub fn open_existing(root: impl AsRef<Path>, sid: &str) -> Result<Self, SessionOpenError> {
694        Self::open_existing_with_redactor(root, sid, None)
695    }
696
697    pub fn open_existing_with_redactor(
698        root: impl AsRef<Path>,
699        sid: &str,
700        redactor: Option<std::sync::Arc<crate::redact::Redactor>>,
701    ) -> Result<Self, SessionOpenError> {
702        let project_index = default_project_index(root.as_ref());
703        Self::open_existing_with_context(root, sid, redactor, project_index)
704    }
705
706    pub fn open_existing_with_context(
707        root: impl AsRef<Path>,
708        sid: &str,
709        redactor: Option<std::sync::Arc<crate::redact::Redactor>>,
710        project_index: Option<std::sync::Arc<crate::index::AnchorIndex>>,
711    ) -> Result<Self, SessionOpenError> {
712        let id = SessionId::parse(sid).map_err(|_| SessionOpenError::InvalidId {
713            sid: sid.to_string(),
714        })?;
715        let dir = root.as_ref().join("sessions").join(id.to_string());
716        if let Some(ls) = crate::notify::log_sink() {
717            ls.set_session_id(Some(id.to_string()));
718        }
719        if !dir.exists() {
720            return Err(SessionOpenError::NotFound {
721                sid: sid.to_string(),
722                dir: dir.clone(),
723            });
724        }
725        let writer = EventWriter::spawn_full(
726            &dir,
727            redactor.clone(),
728            project_index.clone(),
729            Some(id.to_string()),
730        )
731        .map_err(SessionOpenError::WriterInit)?;
732        let mut sink = EventSink::new().with_forwarder(writer.sender());
733        if let Some(r) = redactor {
734            sink = sink.with_redactor(r);
735        }
736        let events_path = dir.join("events.jsonl");
737        let messages = replay_messages_from(&events_path)?;
738        let initial_msgs = replay_messages_with_seq(&events_path)?;
739        let all_msgs = replay_all_messages_with_seq(&events_path)?;
740        if let Some(last_seq) = find_last_seq(&events_path)? {
741            sink.restore_seq(last_seq);
742        }
743        let mut initial_context = replay_context_snapshot_from(&events_path);
744        let persisted = PersistedContextState::load(&dir);
745        if !persisted.model.is_empty() {
746            initial_context.model = persisted.model;
747        }
748        initial_context.window_tokens = persisted.window_tokens;
749        initial_context.window_budget = persisted.window_budget;
750        let initial_goal = load_goal(&dir);
751        let (injection_tx, _) = broadcast::channel(32);
752        let (stream_tx, _) = broadcast::channel(2048);
753        let (context_watch, context_rx) = watch::channel(initial_context);
754        let (goal_watch, goal_rx) = watch::channel(initial_goal);
755        let (attach_watch, attach_rx) = watch::channel(0);
756        let (todos_watch, todos_rx) = watch::channel(Vec::new());
757        let (plans_watch, plans_rx) = watch::channel(Vec::new());
758        let events_handle = sink.events_handle();
759        let output_store = std::sync::Arc::new(crate::tools::tool_output::OutputStore::at(&dir));
760        Ok(Self {
761            id,
762            dir,
763            writer: std::sync::Mutex::new(Some(writer)),
764            sink,
765            message_stream: crate::message_stream::MessageStream::with_initial(
766                events_handle,
767                initial_msgs,
768                all_msgs,
769            ),
770            messages: std::sync::Arc::new(std::sync::Mutex::new(messages)),
771            output_store: output_store.clone(),
772            turn: TurnState::new(),
773            watch: WatchHub {
774                stream_tx,
775                context: context_watch,
776                goal: goal_watch,
777                attach: attach_watch,
778                todos: todos_watch,
779                plans: plans_watch,
780                _keepalive: (context_rx, goal_rx, attach_rx, todos_rx, plans_rx),
781            },
782            watch_hub: std::sync::Arc::new(crate::watch::WatchHub::new()),
783            flow_registry: std::sync::Arc::new(crate::tools::agent_ctrl::FlowRegistry::new()),
784            current_root: std::sync::Mutex::new(None),
785            successful_flow_count: std::sync::atomic::AtomicU64::new(0),
786            compaction: {
787                let c = CompactionState::new();
788                if persisted.window_tokens > 0 {
789                    c.last_input_tokens.store(
790                        persisted.window_tokens,
791                        std::sync::atomic::Ordering::Relaxed,
792                    );
793                }
794                c
795            },
796            interactions: InteractionServices::new(),
797            injection_queue: Mutex::new(Vec::new()),
798            injection_tx,
799            last_image_user_msg: Mutex::new(None),
800            read_files: std::sync::Arc::new(
801                std::sync::Mutex::new(std::collections::HashSet::new()),
802            ),
803            fs_access_mode: Mutex::new(None),
804            project_index,
805        })
806    }
807
808    pub fn open_ephemeral() -> Self {
809        let (injection_tx, _) = broadcast::channel(32);
810        let (stream_tx, _) = broadcast::channel(2048);
811        let (context_watch, context_rx) = watch::channel(ContextSnapshot::default());
812        let (goal_watch, goal_rx) = watch::channel(None);
813        let (attach_watch, attach_rx) = watch::channel(0);
814        let (todos_watch, todos_rx) = watch::channel(Vec::new());
815        let (plans_watch, plans_rx) = watch::channel(Vec::new());
816        let sink = EventSink::new();
817        let events_handle = sink.events_handle();
818        let output_store = std::sync::Arc::new(crate::tools::tool_output::OutputStore::default());
819        Self {
820            id: SessionId::now(),
821            dir: PathBuf::new(),
822            writer: std::sync::Mutex::new(None),
823            sink,
824            message_stream: crate::message_stream::MessageStream::new(events_handle),
825            messages: std::sync::Arc::new(std::sync::Mutex::new(Vec::new())),
826            output_store: output_store.clone(),
827            turn: TurnState::new(),
828            watch: WatchHub {
829                stream_tx,
830                context: context_watch,
831                goal: goal_watch,
832                attach: attach_watch,
833                todos: todos_watch,
834                plans: plans_watch,
835                _keepalive: (context_rx, goal_rx, attach_rx, todos_rx, plans_rx),
836            },
837            watch_hub: std::sync::Arc::new(crate::watch::WatchHub::new()),
838            flow_registry: std::sync::Arc::new(crate::tools::agent_ctrl::FlowRegistry::new()),
839            current_root: std::sync::Mutex::new(None),
840            successful_flow_count: std::sync::atomic::AtomicU64::new(0),
841            compaction: CompactionState::new(),
842            interactions: InteractionServices::new(),
843            injection_queue: Mutex::new(Vec::new()),
844            injection_tx,
845            last_image_user_msg: Mutex::new(None),
846            read_files: std::sync::Arc::new(
847                std::sync::Mutex::new(std::collections::HashSet::new()),
848            ),
849            fs_access_mode: Mutex::new(None),
850            project_index: None,
851        }
852    }
853
854    pub fn project_index(&self) -> Option<std::sync::Arc<crate::index::AnchorIndex>> {
855        self.project_index.clone()
856    }
857
858    pub fn approval(&self) -> std::sync::Arc<ApprovalRegistry> {
859        self.interactions.approval.clone()
860    }
861
862    pub fn compact_reviews(&self) -> std::sync::Arc<CompactReviewRegistry> {
863        self.interactions.compact_reviews.clone()
864    }
865
866    pub fn forms(&self) -> std::sync::Arc<FormRegistry> {
867        self.interactions.forms.clone()
868    }
869
870    pub fn fs_access_mode(&self) -> Option<crate::fs_access::FsAccessMode> {
871        *self.fs_access_mode.lock().unwrap()
872    }
873
874    pub fn set_fs_access_mode(&self, mode: crate::fs_access::FsAccessMode) {
875        *self.fs_access_mode.lock().unwrap() = Some(mode);
876    }
877
878    pub fn compact_review_mode(&self) -> CompactReviewMode {
879        *self.compaction.review_mode.lock().unwrap()
880    }
881
882    pub fn set_compact_review_mode(&self, mode: CompactReviewMode) {
883        *self.compaction.review_mode.lock().unwrap() = mode;
884    }
885
886    pub fn read_files(
887        &self,
888    ) -> std::sync::Arc<std::sync::Mutex<std::collections::HashSet<std::path::PathBuf>>> {
889        self.read_files.clone()
890    }
891
892    pub fn output_store(&self) -> std::sync::Arc<crate::tools::tool_output::OutputStore> {
893        self.output_store.clone()
894    }
895
896    pub fn mark_file_read(&self, path: &std::path::Path) {
897        if let Ok(mut set) = self.read_files.lock() {
898            set.insert(path.to_path_buf());
899            if let Ok(canonical) = std::fs::canonicalize(path) {
900                set.insert(canonical);
901            }
902        }
903    }
904
905    pub fn stream_tx(&self) -> broadcast::Sender<StreamFrame> {
906        self.watch.stream_tx.clone()
907    }
908
909    pub fn set_current_root(&self, handle: String) {
910        *self.current_root.lock().unwrap() = Some(handle);
911    }
912
913    pub fn current_root(&self) -> Option<String> {
914        self.current_root.lock().unwrap().clone()
915    }
916
917    pub fn clear_current_root(&self) {
918        *self.current_root.lock().unwrap() = None;
919    }
920
921    pub fn record_successful_flow(&self) -> Option<u64> {
922        let count = self
923            .successful_flow_count
924            .fetch_add(1, std::sync::atomic::Ordering::Relaxed)
925            + 1;
926        is_auto_name_threshold(count).then_some(count)
927    }
928
929    pub fn successful_flow_count(&self) -> u64 {
930        self.successful_flow_count
931            .load(std::sync::atomic::Ordering::Relaxed)
932    }
933
934    pub fn stream_subscribe(&self) -> broadcast::Receiver<StreamFrame> {
935        self.watch.stream_tx.subscribe()
936    }
937
938    pub fn id(&self) -> &SessionId {
939        &self.id
940    }
941
942    pub fn dir(&self) -> &Path {
943        &self.dir
944    }
945
946    pub fn transcript_replay(&self) -> Vec<TranscriptEntry> {
947        let Some(path) = self.events_path() else {
948            return Vec::new();
949        };
950        replay_transcript_from(&path).unwrap_or_default()
951    }
952
953    pub fn events_path(&self) -> Option<std::path::PathBuf> {
954        self.writer
955            .lock()
956            .unwrap()
957            .as_ref()
958            .map(|w| w.events_path().to_path_buf())
959    }
960
961    pub async fn plan_system_prompt(&self) -> Option<String> {
962        let store = crate::memory::plan::PlanStore::at(&self.dir);
963        let plan = store.latest().await.ok().flatten()?;
964        Some(crate::tools::plan::render_plan(&plan))
965    }
966
967    pub fn goal(&self) -> Option<String> {
968        if let Some(cached) = self.watch.goal.borrow().clone() {
969            return Some(cached);
970        }
971        load_goal(&self.dir)
972    }
973
974    pub fn subscribe_goal(&self) -> watch::Receiver<Option<String>> {
975        self.watch.goal.subscribe()
976    }
977
978    pub fn goal_watch(&self) -> &watch::Sender<Option<String>> {
979        &self.watch.goal
980    }
981
982    pub fn subscribe_context(&self) -> watch::Receiver<ContextSnapshot> {
983        self.watch.context.subscribe()
984    }
985
986    pub fn subscribe_attach(&self) -> watch::Receiver<usize> {
987        self.watch.attach.subscribe()
988    }
989
990    pub fn subscribe_pending_approvals(&self) -> watch::Receiver<Vec<PendingApproval>> {
991        self.interactions.approval.subscribe()
992    }
993
994    pub fn meta(&self) -> Option<crate::session_meta::SessionMeta> {
995        crate::session_meta::SessionMeta::load(&self.dir)
996    }
997
998    pub fn request_manual_compact(&self) {
999        self.compaction
1000            .manual_pending
1001            .store(true, std::sync::atomic::Ordering::SeqCst);
1002    }
1003
1004    pub fn take_manual_compact_request(&self) -> bool {
1005        self.compaction
1006            .manual_pending
1007            .swap(false, std::sync::atomic::Ordering::SeqCst)
1008    }
1009
1010    pub fn set_goal(&self, goal: Option<String>) {
1011        let _ = self.watch.goal.send(goal);
1012    }
1013
1014    pub fn set_attach_count(&self, count: usize) {
1015        let _ = self.watch.attach.send(count);
1016    }
1017
1018    #[allow(clippy::too_many_arguments)]
1019    pub fn record_llm_call(
1020        &self,
1021        model: &str,
1022        tokens_in: u64,
1023        tokens_out: u64,
1024        cache_read: u64,
1025        cache_write: u64,
1026        ttft_ms: Option<u64>,
1027        tokens_per_sec: Option<f64>,
1028    ) {
1029        if tokens_in > 0 {
1030            self.compaction
1031                .last_input_tokens
1032                .store(tokens_in, std::sync::atomic::Ordering::Relaxed);
1033        }
1034        self.watch.context.send_modify(|snap| {
1035            snap.model = model.to_string();
1036            snap.tokens_in = snap.tokens_in.saturating_add(tokens_in);
1037            snap.tokens_out = snap.tokens_out.saturating_add(tokens_out);
1038            snap.cache_read = snap.cache_read.saturating_add(cache_read);
1039            snap.cache_write = snap.cache_write.saturating_add(cache_write);
1040            snap.last_ttft_ms = ttft_ms.unwrap_or(0);
1041            snap.last_tokens_per_sec = tokens_per_sec.unwrap_or(0.0);
1042        });
1043        self.refresh_window_snapshot();
1044    }
1045
1046    pub fn last_input_tokens(&self) -> u64 {
1047        self.compaction
1048            .last_input_tokens
1049            .load(std::sync::atomic::Ordering::Relaxed)
1050    }
1051
1052    pub async fn acquire_compact_lock(&self) -> tokio::sync::MutexGuard<'_, ()> {
1053        self.compaction.lock.lock().await
1054    }
1055
1056    pub async fn acquire_compact_lock_owned(&self) -> tokio::sync::OwnedMutexGuard<()> {
1057        self.compaction.lock.clone().lock_owned().await
1058    }
1059
1060    pub fn compact_lock_handle(&self) -> std::sync::Arc<tokio::sync::Mutex<()>> {
1061        self.compaction.lock.clone()
1062    }
1063
1064    pub fn refresh_window_snapshot(&self) {
1065        let provider_tokens = self.last_input_tokens();
1066        let estimated = crate::compaction::estimate_tokens_for_messages(&self.messages());
1067        let window = if provider_tokens > 0 {
1068            provider_tokens
1069        } else {
1070            estimated
1071        };
1072        let model = self.last_model();
1073        let budget = crate::model_registry::model_info(&model).context_budget;
1074        self.watch.context.send_modify(|snap| {
1075            snap.window_tokens = window;
1076            if budget > 0 {
1077                snap.window_budget = budget;
1078            }
1079        });
1080        let snap = self.watch.context.borrow();
1081        PersistedContextState {
1082            model,
1083            window_tokens: snap.window_tokens,
1084            window_budget: snap.window_budget,
1085        }
1086        .save(&self.dir);
1087    }
1088
1089    pub fn cumulative_input_tokens(&self) -> u64 {
1090        self.watch.context.borrow().tokens_in
1091    }
1092
1093    pub fn reset_input_tokens_to(&self, tokens: u64) {
1094        self.watch.context.send_modify(|snap| {
1095            snap.tokens_in = tokens;
1096        });
1097    }
1098
1099    pub fn last_model(&self) -> String {
1100        self.watch.context.borrow().model.clone()
1101    }
1102
1103    pub fn set_current_model(&self, model: impl Into<String>) {
1104        let model = model.into();
1105        let budget = crate::model_registry::model_info(&model).context_budget;
1106        self.watch.context.send_modify(|snap| {
1107            snap.model = model.clone();
1108            if budget > 0 {
1109                snap.window_budget = budget;
1110            }
1111        });
1112        let snap = self.watch.context.borrow();
1113        PersistedContextState {
1114            model,
1115            window_tokens: snap.window_tokens,
1116            window_budget: snap.window_budget,
1117        }
1118        .save(&self.dir);
1119    }
1120
1121    pub fn update_mcp_server(&self, status: crate::mcp::McpServerStatus) {
1122        self.watch.context.send_modify(|snap| {
1123            if let Some(existing) = snap.mcp_servers.iter_mut().find(|s| s.name == status.name) {
1124                *existing = status;
1125            } else {
1126                snap.mcp_servers.push(status);
1127            }
1128        });
1129    }
1130
1131    pub fn set_memory_recent_count(&self, count: u16) {
1132        self.watch.context.send_modify(|snap| {
1133            snap.memory_recent_count = count;
1134        });
1135    }
1136
1137    pub fn subscribe_todos(&self) -> watch::Receiver<Vec<crate::memory::todo::Todo>> {
1138        self.watch.todos.subscribe()
1139    }
1140
1141    pub fn todos_watch(&self) -> &watch::Sender<Vec<crate::memory::todo::Todo>> {
1142        &self.watch.todos
1143    }
1144
1145    pub fn subscribe_plans(&self) -> watch::Receiver<Vec<crate::memory::plan::Plan>> {
1146        self.watch.plans.subscribe()
1147    }
1148
1149    pub fn plans_watch(&self) -> &watch::Sender<Vec<crate::memory::plan::Plan>> {
1150        &self.watch.plans
1151    }
1152
1153    pub async fn refresh_plans_from_store_async(&self) {
1154        if self.dir.as_os_str().is_empty() {
1155            return;
1156        }
1157        let store = crate::memory::plan::PlanStore::at(&self.dir);
1158        match store.list().await {
1159            Ok(list) => {
1160                let _ = self.watch.plans.send(list);
1161            }
1162            Err(e) => {
1163                crate::notify!(
1164                    warn,
1165                    location = Log,
1166                    stack = dedupe("memory.refresh_plans_async", 60_000),
1167                    "refresh_plans_from_store_async: {e}"
1168                );
1169            }
1170        }
1171    }
1172
1173    pub fn refresh_todos_from_store(&self) {
1174        if self.dir.as_os_str().is_empty() {
1175            return;
1176        }
1177        let store = crate::memory::todo::TodoStore::at(&self.dir);
1178        match tokio::task::block_in_place(|| {
1179            tokio::runtime::Handle::try_current()
1180                .ok()
1181                .map(|h| h.block_on(store.list()))
1182        }) {
1183            Some(Ok(list)) => {
1184                let _ = self.watch.todos.send(list);
1185            }
1186            Some(Err(e)) => {
1187                crate::notify!(
1188                    warn,
1189                    location = Log,
1190                    stack = dedupe("memory.refresh_todos", 60_000),
1191                    "refresh_todos_from_store: {e}"
1192                );
1193            }
1194            None => {}
1195        }
1196    }
1197
1198    pub async fn refresh_todos_from_store_async(&self) {
1199        if self.dir.as_os_str().is_empty() {
1200            return;
1201        }
1202        let store = crate::memory::todo::TodoStore::at(&self.dir);
1203        match store.list().await {
1204            Ok(list) => {
1205                let _ = self.watch.todos.send(list);
1206            }
1207            Err(e) => {
1208                crate::notify!(
1209                    warn,
1210                    location = Log,
1211                    stack = dedupe("memory.refresh_todos_async", 60_000),
1212                    "refresh_todos_from_store_async: {e}"
1213                );
1214            }
1215        }
1216    }
1217
1218    pub fn sink(&self) -> &EventSink {
1219        &self.sink
1220    }
1221
1222    /// Single-writer append. Emits the matching event before the in-memory push
1223    /// so events.jsonl remains the authority (§I5).
1224    pub fn append_message(&self, msg: Message, flow_run_id: Option<FlowRunId>) {
1225        AppendMessageCommand { msg, flow_run_id }.execute(self);
1226    }
1227
1228    pub fn emit_attachment_degrade(
1229        &self,
1230        message_seq: u64,
1231        part_index: usize,
1232        file_basename: String,
1233        reason: String,
1234    ) {
1235        self.sink.emit(Event::AttachmentDegraded {
1236            turn_id: None,
1237            flow_run_id: None,
1238            message_seq,
1239            part_index,
1240            file_basename,
1241            reason,
1242        });
1243    }
1244
1245    pub fn record_attachment_degrade(&self, reason: &str) -> usize {
1246        let target = self.last_image_user_msg.lock().unwrap().take();
1247        let Some(entry) = target else {
1248            return 0;
1249        };
1250        let turn_id = self.turn.current_turn.lock().unwrap().clone();
1251        for (part_index, basename) in &entry.images {
1252            self.sink.emit(Event::AttachmentDegraded {
1253                turn_id: turn_id.clone(),
1254                flow_run_id: None,
1255                message_seq: entry.message_seq,
1256                part_index: *part_index,
1257                file_basename: basename.clone(),
1258                reason: reason.into(),
1259            });
1260        }
1261        if let Ok(mut msgs) = self.messages.lock() {
1262            for m in msgs.iter_mut() {
1263                for (part_index, basename) in &entry.images {
1264                    if let Some(part) = m.parts.get_mut(*part_index)
1265                        && matches!(part, crate::message::MessagePart::Image { .. })
1266                    {
1267                        *part = crate::message::MessagePart::Text {
1268                            text: format!("[attachment unavailable: {basename} — {reason}]"),
1269                        };
1270                    }
1271                }
1272            }
1273        }
1274        entry.images.len()
1275    }
1276
1277    pub fn messages(&self) -> crate::message_stream::MessageWindow {
1278        self.message_stream.window()
1279    }
1280
1281    pub fn messages_full(&self) -> std::sync::Arc<Vec<Message>> {
1282        self.message_stream.full_messages()
1283    }
1284
1285    pub fn messages_handle(&self) -> std::sync::Arc<std::sync::Mutex<Vec<Message>>> {
1286        self.messages.clone()
1287    }
1288
1289    pub fn message_count(&self) -> usize {
1290        self.messages().len()
1291    }
1292
1293    pub fn user_message_count(&self) -> usize {
1294        self.messages()
1295            .iter()
1296            .filter(|m| matches!(m.role, MessageRole::User))
1297            .count()
1298    }
1299
1300    pub fn push_system_note(&self, text: String) {
1301        let _ = self
1302            .watch
1303            .stream_tx
1304            .send(crate::stream::StreamFrame::Note(text));
1305    }
1306
1307    pub fn approval_cooldown_ok_for_compact(&self) -> bool {
1308        self.sink.last_compact_ago_seconds().is_none_or(|s| s >= 60)
1309    }
1310
1311    pub fn emit_compact_warning(
1312        &self,
1313        model: &str,
1314        current_tokens: u64,
1315        threshold: u64,
1316        budget: u64,
1317        reason: &str,
1318    ) {
1319        let message = format!(
1320            "context {current_tokens} > threshold {threshold} (budget {budget}, model {model}); skipping compaction: {reason}"
1321        );
1322        self.sink.emit(Event::WatchWarn {
1323            turn_id: self.turn.current_turn.lock().unwrap().clone(),
1324            flow_run_id: None,
1325            target: "context.compaction".into(),
1326            trigger: "auto_compact".into(),
1327            message,
1328        });
1329        self.push_system_note(format!("[warn] compaction skipped: {reason}"));
1330    }
1331
1332    /// Convenience wrapper that computes the compact range and token count
1333    /// from the current message window. Used by tests and internal callers
1334    /// that don't already have a pre-computed range.
1335    pub fn compact_messages_auto(&self, summary: String) -> Option<CompactResult> {
1336        let msgs = self.messages();
1337        let tokens = crate::compaction::estimate_tokens_for_messages(&msgs);
1338        let info = crate::model_registry::model_info(&self.last_model());
1339        let target = info.compaction_target_after();
1340        let range = crate::compaction::find_compact_range(&msgs, target)?;
1341        self.compact_messages(summary, range, tokens)
1342    }
1343
1344    pub fn commit_rewritten_window(
1345        &self,
1346        replacement: Vec<Message>,
1347        before_tokens: u64,
1348        before_window_tokens: u64,
1349        rewritten_count: usize,
1350    ) -> Option<CompactResult> {
1351        let after_tokens = crate::compaction::estimate_tokens_for_messages(&replacement);
1352        if rewritten_count == 0 || after_tokens >= before_window_tokens {
1353            return None;
1354        }
1355        let summary = format!(
1356            "[atman: persistently compacted output from {rewritten_count} retained messages]"
1357        );
1358        self.sink.mark_compacted();
1359        self.sink.emit(Event::ContextCompact {
1360            session_id: self.id.to_string(),
1361            before_tokens,
1362            after_tokens,
1363            compacted_range_start: 0,
1364            compacted_range_end: 0,
1365            summary_text: Some(summary.clone()),
1366            replacement_msg_seq: None,
1367        });
1368        self.sink.emit(Event::CompactionSummary {
1369            session_id: self.id.to_string(),
1370            range_start: 0,
1371            range_end: 0,
1372            compacted_count: rewritten_count,
1373            before_tokens,
1374            after_tokens,
1375            summary: summary.clone(),
1376        });
1377        let _ = self
1378            .watch
1379            .stream_tx
1380            .send(crate::stream::StreamFrame::CompactionSummary {
1381                phase: crate::stream::CompactionPhase::Finished,
1382                range_start: 0,
1383                range_end: 0,
1384                summary,
1385                before_tokens,
1386                after_tokens,
1387                compacted_count: rewritten_count,
1388            });
1389        self.compaction
1390            .last_input_tokens
1391            .store(after_tokens, std::sync::atomic::Ordering::Relaxed);
1392        if let Ok(mut messages) = self.messages.lock() {
1393            *messages = replacement.clone();
1394        }
1395        self.sink.emit(Event::Checkpoint {
1396            session_id: self.id.to_string(),
1397            messages: replacement,
1398            window_tokens: after_tokens,
1399        });
1400        self.refresh_window_snapshot();
1401        Some(CompactResult {
1402            before_tokens,
1403            after_tokens,
1404            compacted_start: 0,
1405            compacted_end: 0,
1406        })
1407    }
1408
1409    pub fn commit_compacted_window(
1410        &self,
1411        summary: String,
1412        replacement: Vec<Message>,
1413        range: crate::compaction::CompactRange,
1414        before_tokens: u64,
1415        before_window_tokens: u64,
1416    ) -> Option<CompactResult> {
1417        let after_tokens = crate::compaction::estimate_tokens_for_messages(&replacement);
1418        if after_tokens >= before_window_tokens {
1419            self.push_system_note(format!(
1420                "compaction skipped: replacement would not shrink transcript ({} >= {} tokens)",
1421                after_tokens, before_window_tokens
1422            ));
1423            return None;
1424        }
1425        self.sink.mark_compacted();
1426        self.sink.emit(Event::ContextCompact {
1427            session_id: self.id.to_string(),
1428            before_tokens,
1429            after_tokens,
1430            compacted_range_start: range.start as u64,
1431            compacted_range_end: range.end.saturating_sub(1) as u64,
1432            summary_text: Some(summary.clone()),
1433            replacement_msg_seq: None,
1434        });
1435        self.sink.emit(Event::CompactionSummary {
1436            session_id: self.id.to_string(),
1437            range_start: range.start as u64,
1438            range_end: range.end.saturating_sub(1) as u64,
1439            compacted_count: range.end - range.start,
1440            before_tokens,
1441            after_tokens,
1442            summary: summary.clone(),
1443        });
1444        let _ = self
1445            .watch
1446            .stream_tx
1447            .send(crate::stream::StreamFrame::CompactionSummary {
1448                phase: crate::stream::CompactionPhase::Finished,
1449                range_start: range.start,
1450                range_end: range.end.saturating_sub(1),
1451                summary,
1452                before_tokens,
1453                after_tokens,
1454                compacted_count: range.end - range.start,
1455            });
1456        self.compaction
1457            .last_input_tokens
1458            .store(after_tokens, std::sync::atomic::Ordering::Relaxed);
1459        if let Ok(mut messages) = self.messages.lock() {
1460            *messages = replacement.clone();
1461        }
1462        self.sink.emit(Event::Checkpoint {
1463            session_id: self.id.to_string(),
1464            messages: replacement,
1465            window_tokens: after_tokens,
1466        });
1467        self.refresh_window_snapshot();
1468        Some(CompactResult {
1469            before_tokens,
1470            after_tokens,
1471            compacted_start: range.start,
1472            compacted_end: range.end,
1473        })
1474    }
1475
1476    pub fn compact_messages(
1477        &self,
1478        summary: String,
1479        range: crate::compaction::CompactRange,
1480        before_tokens: u64,
1481    ) -> Option<CompactResult> {
1482        use crate::compaction::{estimate_tokens_for_messages, replace_range_with_summary};
1483        let msgs = self.messages();
1484        let turn_id = msgs
1485            .get(range.start)
1486            .map(|m| m.turn_id.clone())
1487            .unwrap_or_else(TurnId::now);
1488        let after = replace_range_with_summary(&msgs, &range, summary.clone(), turn_id.clone());
1489        let after_tokens = estimate_tokens_for_messages(&after);
1490        if after_tokens >= before_tokens {
1491            self.push_system_note(format!(
1492                "compaction skipped: summary would not shrink transcript ({} >= {} tokens)",
1493                after_tokens, before_tokens
1494            ));
1495            return None;
1496        }
1497        let replacement_msg = after.first().cloned().unwrap_or_else(|| {
1498            Message::system_compact_summary(
1499                turn_id.clone(),
1500                summary.clone(),
1501                range.start as u64,
1502                range.end.saturating_sub(1) as u64,
1503                range.end - range.start,
1504            )
1505        });
1506        self.sink.mark_compacted();
1507        let replacement_seq = self.sink.next_seq_peek();
1508        self.sink.emit(Event::SystemMsg {
1509            turn_id: turn_id.clone(),
1510            message: replacement_msg,
1511        });
1512        self.sink.emit(Event::ContextCompact {
1513            session_id: self.id.to_string(),
1514            before_tokens,
1515            after_tokens,
1516            compacted_range_start: range.start as u64,
1517            compacted_range_end: range.end.saturating_sub(1) as u64,
1518            summary_text: Some(summary.clone()),
1519            replacement_msg_seq: Some(replacement_seq),
1520        });
1521        self.sink.emit(Event::CompactionSummary {
1522            session_id: self.id.to_string(),
1523            range_start: range.start as u64,
1524            range_end: range.end.saturating_sub(1) as u64,
1525            compacted_count: range.end - range.start,
1526            before_tokens,
1527            after_tokens,
1528            summary: summary.clone(),
1529        });
1530        let _ = self
1531            .watch
1532            .stream_tx
1533            .send(crate::stream::StreamFrame::CompactionSummary {
1534                phase: crate::stream::CompactionPhase::Finished,
1535                range_start: range.start,
1536                range_end: range.end.saturating_sub(1),
1537                summary,
1538                before_tokens,
1539                after_tokens,
1540                compacted_count: range.end - range.start,
1541            });
1542        let checkpoint_messages = self.messages();
1543        let window_tokens = estimate_tokens_for_messages(&checkpoint_messages);
1544        self.compaction
1545            .last_input_tokens
1546            .store(window_tokens, std::sync::atomic::Ordering::Relaxed);
1547        // Sync the messages Vec (root's messages_handle) with the compacted
1548        // windowed view so root's llm context via messages_handle respects
1549        // compaction (branch2 retired → unified on messages_handle).
1550        let window_owned = checkpoint_messages.to_vec();
1551        if let Ok(mut vec) = self.messages.lock() {
1552            *vec = window_owned;
1553        }
1554        self.refresh_window_snapshot();
1555        self.sink.emit(Event::Checkpoint {
1556            session_id: self.id.to_string(),
1557            messages: checkpoint_messages.to_vec(),
1558            window_tokens,
1559        });
1560        Some(CompactResult {
1561            before_tokens,
1562            after_tokens,
1563            compacted_start: range.start,
1564            compacted_end: range.end,
1565        })
1566    }
1567
1568    pub fn begin_turn(&self, user_msg: Message) -> TurnId {
1569        BeginTurnCommand { user_msg }.execute(self)
1570    }
1571
1572    pub fn mark_streamed(&self) {
1573        self.turn
1574            .streamed
1575            .store(true, std::sync::atomic::Ordering::Relaxed);
1576    }
1577
1578    pub fn take_streamed_flag(&self) -> bool {
1579        self.turn
1580            .streamed
1581            .swap(false, std::sync::atomic::Ordering::Relaxed)
1582    }
1583
1584    pub fn end_turn(&self) {
1585        self.turn
1586            .streamed
1587            .store(false, std::sync::atomic::Ordering::Relaxed);
1588        let turn_id = self.turn.current_turn.lock().unwrap().take();
1589        if let Some(turn_id) = turn_id {
1590            let mut q = self.injection_queue.lock().unwrap();
1591            for inj in q.iter_mut() {
1592                if inj.state == InjectionState::Pending && inj.turn_id == turn_id {
1593                    inj.state = InjectionState::Cancelled;
1594                    let _ = self.injection_tx.send(inj.clone());
1595                }
1596            }
1597            drop(q);
1598            self.sink.emit(Event::TurnEnd { turn_id });
1599        }
1600    }
1601
1602    pub fn current_turn(&self) -> Option<TurnId> {
1603        self.turn.current_turn.lock().unwrap().clone()
1604    }
1605
1606    pub fn enqueue_injection(&self, text: impl Into<String>) -> Result<InjectionId, EnqueueError> {
1607        self.enqueue_injection_with_level(text, crate::injection::InjectionLevel::L1Nudge, None)
1608    }
1609
1610    pub fn enqueue_injection_with_level(
1611        &self,
1612        text: impl Into<String>,
1613        level: crate::injection::InjectionLevel,
1614        redirect_target: Option<String>,
1615    ) -> Result<InjectionId, EnqueueError> {
1616        let turn_id = self
1617            .turn
1618            .current_turn
1619            .lock()
1620            .unwrap()
1621            .clone()
1622            .ok_or(EnqueueError::NoActiveTurn)?;
1623        let inj = Injection::with_level(turn_id.clone(), text, level, redirect_target);
1624        let id = inj.id.clone();
1625        self.sink.emit(Event::UserInject {
1626            turn_id,
1627            injection: inj.clone(),
1628        });
1629        self.injection_queue.lock().unwrap().push(inj.clone());
1630        let _ = self.injection_tx.send(inj);
1631        Ok(id)
1632    }
1633
1634    pub fn subscribe_injections(&self) -> broadcast::Receiver<Injection> {
1635        self.injection_tx.subscribe()
1636    }
1637
1638    pub fn mark_injection_consumed(&self, id: &InjectionId) {
1639        let mut q = self.injection_queue.lock().unwrap();
1640        for inj in q.iter_mut() {
1641            if inj.id == *id && inj.state == InjectionState::Pending {
1642                inj.state = InjectionState::Injected;
1643                let _ = self.injection_tx.send(inj.clone());
1644                return;
1645            }
1646        }
1647    }
1648
1649    pub fn peek_pending_l2_or_higher(&self, turn_id: &TurnId) -> Option<Injection> {
1650        let q = self.injection_queue.lock().unwrap();
1651        q.iter()
1652            .find(|i| {
1653                i.state == InjectionState::Pending
1654                    && i.turn_id == *turn_id
1655                    && !matches!(i.level, crate::injection::InjectionLevel::L1Nudge)
1656            })
1657            .cloned()
1658    }
1659
1660    /// Drain all Pending injections for `turn_id`. Marks them Injected.
1661    /// Returns them in creation order.
1662    pub fn drain_injections(&self, turn_id: &TurnId) -> Vec<Injection> {
1663        let mut q = self.injection_queue.lock().unwrap();
1664        let mut out = Vec::new();
1665        for inj in q.iter_mut() {
1666            if inj.state == InjectionState::Pending && inj.turn_id == *turn_id {
1667                inj.state = InjectionState::Injected;
1668                let _ = self.injection_tx.send(inj.clone());
1669                out.push(inj.clone());
1670            }
1671        }
1672        out
1673    }
1674
1675    pub fn list_pending_injections(&self) -> Vec<Injection> {
1676        self.injection_queue
1677            .lock()
1678            .unwrap()
1679            .iter()
1680            .filter(|i| i.state == InjectionState::Pending)
1681            .cloned()
1682            .collect()
1683    }
1684
1685    pub fn cancel_flow(&self) {
1686        self.turn.flow_cancel.lock().unwrap().cancel();
1687    }
1688
1689    pub fn flow_cancel_token(&self) -> CancellationToken {
1690        self.turn.flow_cancel.lock().unwrap().clone()
1691    }
1692
1693    pub async fn shutdown(&self) {
1694        let writer = self.writer.lock().unwrap().take();
1695        if let Some(writer) = writer {
1696            writer.shutdown().await;
1697        }
1698    }
1699
1700    // Rides FIFO queue ordering: once flush's own barrier is written,
1701    // every earlier sink.emit is on disk too.
1702    #[allow(clippy::await_holding_lock)]
1703    pub async fn flush_writer(&self) {
1704        let guard = self.writer.lock().unwrap();
1705        let Some(ref writer) = *guard else {
1706            return;
1707        };
1708        writer.flush().await;
1709    }
1710}
1711
1712#[derive(Debug, thiserror::Error)]
1713pub enum EnqueueError {
1714    #[error("enqueue_injection called with no active turn")]
1715    NoActiveTurn,
1716}
1717
1718pub struct AppendMessageCommand {
1719    pub msg: Message,
1720    pub flow_run_id: Option<FlowRunId>,
1721}
1722
1723impl AppendMessageCommand {
1724    pub fn execute(&self, session: &Session) -> u64 {
1725        let flow_run_id_str = self.flow_run_id.as_ref().map(|r| r.0.to_string());
1726        let msg = crate::tools::tool_output::maybe_truncate_tool_message(
1727            &self.msg,
1728            Some(&session.output_store),
1729        );
1730        let is_internal = msg.origin == crate::message::MessageOrigin::Internal;
1731        let event =
1732            match msg.role {
1733                MessageRole::User => Event::UserMsg {
1734                    turn_id: msg.turn_id.clone(),
1735                    flow_run_id: self.flow_run_id.clone(),
1736                    message: msg.clone(),
1737                },
1738                MessageRole::Assistant => {
1739                    if !is_internal {
1740                        let _ = session.watch.stream_tx.send(
1741                            crate::stream::StreamFrame::AssistantMsg {
1742                                flow_run_id: flow_run_id_str.clone(),
1743                                message: msg.clone(),
1744                            },
1745                        );
1746                        for source in extract_mermaid_blocks(&msg) {
1747                            let _ = session.watch.stream_tx.send(
1748                                crate::stream::StreamFrame::MermaidDiagram {
1749                                    source: source.clone(),
1750                                },
1751                            );
1752                            session
1753                                .sink
1754                                .emit(crate::event::Event::MermaidDiagram { source });
1755                        }
1756                    }
1757                    Event::AssistantMsg {
1758                        turn_id: msg.turn_id.clone(),
1759                        flow_run_id: self.flow_run_id.clone(),
1760                        message: msg.clone(),
1761                    }
1762                }
1763                MessageRole::Tool => {
1764                    if !is_internal {
1765                        let _ = session.watch.stream_tx.send(
1766                            crate::stream::StreamFrame::ToolResultMsg {
1767                                flow_run_id: flow_run_id_str.clone(),
1768                                message: msg.clone(),
1769                            },
1770                        );
1771                    }
1772                    Event::ToolResultMsg {
1773                        turn_id: msg.turn_id.clone(),
1774                        flow_run_id: self.flow_run_id.clone(),
1775                        message: msg.clone(),
1776                    }
1777                }
1778                MessageRole::System => Event::SystemMsg {
1779                    turn_id: msg.turn_id.clone(),
1780                    message: msg.clone(),
1781                },
1782            };
1783        let seq = session.sink.emit_returning_seq(event);
1784        if matches!(msg.role, MessageRole::User) {
1785            let images: Vec<(usize, String)> = msg
1786                .parts
1787                .iter()
1788                .enumerate()
1789                .filter_map(|(i, p)| match p {
1790                    crate::message::MessagePart::Image { source } => {
1791                        let basename = match &source.data {
1792                            crate::message::ImageData::Path { path } => path
1793                                .file_name()
1794                                .and_then(|n| n.to_str())
1795                                .unwrap_or("unknown")
1796                                .to_string(),
1797                            crate::message::ImageData::Base64 { .. } => "base64".into(),
1798                        };
1799                        Some((i, basename))
1800                    }
1801                    _ => None,
1802                })
1803                .collect();
1804            if !images.is_empty() {
1805                *session.last_image_user_msg.lock().unwrap() = Some(LastImageUserMsg {
1806                    message_seq: seq,
1807                    images,
1808                });
1809            }
1810        }
1811        session.messages.lock().unwrap().push(msg.clone());
1812        seq
1813    }
1814}
1815
1816pub struct BeginTurnCommand {
1817    pub user_msg: Message,
1818}
1819
1820impl BeginTurnCommand {
1821    pub fn execute(&self, session: &Session) -> TurnId {
1822        let turn_id = self.user_msg.turn_id.clone();
1823        *session.turn.current_turn.lock().unwrap() = Some(turn_id.clone());
1824        *session.turn.flow_cancel.lock().unwrap() = tokio_util::sync::CancellationToken::new();
1825        session.sink.emit(Event::TurnStart {
1826            turn_id: turn_id.clone(),
1827        });
1828        AppendMessageCommand {
1829            msg: self.user_msg.clone(),
1830            flow_run_id: None,
1831        }
1832        .execute(session);
1833        turn_id
1834    }
1835}
1836
1837fn extract_mermaid_blocks(msg: &crate::message::Message) -> Vec<String> {
1838    let text = msg.text_concat();
1839    let mut blocks = Vec::new();
1840    let mut lines = text.lines().peekable();
1841    while let Some(line) = lines.next() {
1842        let trimmed = line.trim();
1843        if trimmed.starts_with("```") {
1844            let lang = trimmed.trim_start_matches("```").trim();
1845            if lang == "mermaid" {
1846                let mut source = String::new();
1847                for inner in lines.by_ref() {
1848                    if inner.trim() == "```" {
1849                        break;
1850                    }
1851                    if !source.is_empty() {
1852                        source.push('\n');
1853                    }
1854                    source.push_str(inner);
1855                }
1856                if !source.is_empty() {
1857                    blocks.push(source);
1858                }
1859            } else {
1860                for inner in lines.by_ref() {
1861                    if inner.trim() == "```" {
1862                        break;
1863                    }
1864                }
1865            }
1866        }
1867    }
1868    blocks
1869}
1870
1871#[cfg(test)]
1872mod tests {
1873    use super::*;
1874    use tempfile::TempDir;
1875
1876    fn write_events(dir: &Path, lines: &[&str]) {
1877        let path = dir.join("events.jsonl");
1878        std::fs::write(&path, lines.join("\n") + "\n").unwrap();
1879    }
1880
1881    #[test]
1882    fn successful_flow_count_triggers_at_powers_of_three_and_resets_per_session() {
1883        let session = Session::open_ephemeral();
1884        let mut hits = Vec::new();
1885        for _ in 0..27 {
1886            if let Some(count) = session.record_successful_flow() {
1887                hits.push(count);
1888            }
1889        }
1890        assert_eq!(hits, vec![3, 9, 27]);
1891        assert_eq!(session.successful_flow_count(), 27);
1892        assert_eq!(Session::open_ephemeral().successful_flow_count(), 0);
1893    }
1894
1895    #[test]
1896    fn commit_rewritten_window_uses_checkpoint_without_legacy_range_replay() {
1897        let session = Session::open_ephemeral();
1898        let original = vec![
1899            Message::user_text(TurnId::now(), "first user"),
1900            Message::assistant_text(TurnId::now(), "large output".repeat(2_000)),
1901            Message::user_text(TurnId::now(), "current user"),
1902        ];
1903        for message in original.clone() {
1904            session.append_message(message, None);
1905        }
1906        let replacement = vec![
1907            original[0].clone(),
1908            Message::assistant_text(TurnId::now(), "persisted omission"),
1909            original[2].clone(),
1910        ];
1911        let before_tokens = crate::compaction::estimate_tokens_for_messages(&original);
1912
1913        session
1914            .commit_rewritten_window(replacement.clone(), before_tokens, before_tokens, 1)
1915            .expect("rewrite commit");
1916
1917        assert_eq!(session.messages().as_ref(), replacement.as_slice());
1918        let events = session.sink().snapshot();
1919        assert!(events.iter().any(|event| matches!(
1920            event,
1921            Event::ContextCompact {
1922                replacement_msg_seq: None,
1923                ..
1924            }
1925        )));
1926        assert!(
1927            !events
1928                .iter()
1929                .any(|event| matches!(event, Event::SystemMsg { .. }))
1930        );
1931        let checkpoint = events
1932            .into_iter()
1933            .find(|event| matches!(event, Event::Checkpoint { .. }))
1934            .expect("checkpoint");
1935        let replay = crate::message_stream::MessageStream::new(std::sync::Arc::new(
1936            std::sync::Mutex::new(vec![crate::event::EventEnvelope::new(1, checkpoint)]),
1937        ));
1938        assert_eq!(&*replay.window(), replacement.as_slice());
1939    }
1940
1941    #[test]
1942    fn commit_compacted_window_updates_live_handle_and_checkpoint_replay() {
1943        let session = Session::open_ephemeral();
1944        let old = vec![
1945            Message::user_text(TurnId::now(), "old user".repeat(2_000)),
1946            Message::assistant_text(TurnId::now(), "old assistant".repeat(2_000)),
1947            Message::user_text(TurnId::now(), "current user"),
1948        ];
1949        for message in old.clone() {
1950            session.append_message(message, None);
1951        }
1952        let replacement = vec![
1953            Message::system_compact_summary(TurnId::now(), "anchor", 0, 1, 2),
1954            old[2].clone(),
1955            Message::assistant_text(TurnId::now(), "persisted omission"),
1956        ];
1957        let before_tokens = crate::compaction::estimate_tokens_for_messages(&old);
1958        let range = crate::compaction::CompactRange {
1959            start: 0,
1960            end: 2,
1961            tokens_saved_estimate: before_tokens,
1962        };
1963
1964        session
1965            .commit_compacted_window(
1966                "anchor".into(),
1967                replacement.clone(),
1968                range,
1969                before_tokens,
1970                before_tokens,
1971            )
1972            .expect("commit");
1973
1974        assert_eq!(session.messages().as_ref(), replacement.as_slice());
1975        assert_eq!(
1976            session.messages_handle().lock().unwrap().as_slice(),
1977            replacement.as_slice()
1978        );
1979        let checkpoint = session
1980            .sink()
1981            .snapshot()
1982            .into_iter()
1983            .find(|event| matches!(event, Event::Checkpoint { .. }))
1984            .expect("checkpoint");
1985        let replay = crate::message_stream::MessageStream::new(std::sync::Arc::new(
1986            std::sync::Mutex::new(vec![crate::event::EventEnvelope::new(1, checkpoint)]),
1987        ));
1988        assert_eq!(&*replay.window(), replacement.as_slice());
1989    }
1990
1991    #[test]
1992    fn replay_applies_attachment_degraded_patch() {
1993        let dir = TempDir::new().unwrap();
1994        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"}"#;
1995        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"}"#;
1996        write_events(dir.path(), &[user_msg, degrade]);
1997        let entries = replay_transcript_from(&dir.path().join("events.jsonl")).unwrap();
1998        let msg = entries
1999            .into_iter()
2000            .find_map(|e| match e {
2001                TranscriptEntry::Message { message, .. } => Some(message),
2002                _ => None,
2003            })
2004            .unwrap();
2005        assert_eq!(msg.parts.len(), 2);
2006        match &msg.parts[0] {
2007            crate::message::MessagePart::Text { text } => {
2008                assert!(text.contains("photo.png"), "expected basename: {text}");
2009                assert!(text.contains("image_too_large"), "expected reason: {text}");
2010                assert!(text.starts_with("[attachment unavailable"));
2011            }
2012            other => panic!("expected Text stub, got {other:?}"),
2013        }
2014        assert!(matches!(
2015            msg.parts[1],
2016            crate::message::MessagePart::Text { .. }
2017        ));
2018    }
2019
2020    #[test]
2021    fn approval_registry_auto_approves_when_level_leq_ceiling() {
2022        let reg = ApprovalRegistry::new();
2023        reg.set_auto_ceiling(crate::tool::ApprovalLevel::Approve);
2024        let pending = PendingApproval {
2025            tool_use_id: "tu1".into(),
2026            tool_name: "fs.read".into(),
2027            args_preview: "{}".into(),
2028            preview: None,
2029            level: crate::tool::ApprovalLevel::Auto,
2030            run_id: FlowRunId::now(),
2031            emitted_at: chrono::Utc::now(),
2032            bypass_auto_ceiling: false,
2033        };
2034        let rx = reg.request(pending);
2035        let got = rx.blocking_recv().unwrap();
2036        assert!(matches!(got, ApprovalDecision::Approve));
2037        assert!(reg.list_pending().is_empty());
2038    }
2039
2040    #[test]
2041    fn approval_registry_queues_when_level_above_ceiling() {
2042        let reg = std::sync::Arc::new(ApprovalRegistry::new());
2043        reg.set_auto_ceiling(crate::tool::ApprovalLevel::Auto);
2044        let pending = PendingApproval {
2045            tool_use_id: "tu42".into(),
2046            tool_name: "fs.write".into(),
2047            args_preview: "{}".into(),
2048            preview: None,
2049            level: crate::tool::ApprovalLevel::Approve,
2050            run_id: FlowRunId::now(),
2051            emitted_at: chrono::Utc::now(),
2052            bypass_auto_ceiling: false,
2053        };
2054        let mut rx = reg.request(pending);
2055        assert_eq!(reg.list_pending().len(), 1);
2056        assert!(rx.try_recv().is_err(), "should still be queued");
2057        assert!(reg.decide("tu42", ApprovalDecision::Approve));
2058        let got = rx.blocking_recv().unwrap();
2059        assert!(matches!(got, ApprovalDecision::Approve));
2060        assert!(reg.list_pending().is_empty());
2061    }
2062
2063    #[test]
2064    fn approval_registry_decide_all_flushes_queue() {
2065        let reg = ApprovalRegistry::new();
2066        reg.set_auto_ceiling(crate::tool::ApprovalLevel::Auto);
2067        let mut rxs = Vec::new();
2068        for i in 0..3 {
2069            rxs.push(reg.request(PendingApproval {
2070                tool_use_id: format!("tu{i}"),
2071                tool_name: "bash.exec".into(),
2072                args_preview: "{}".into(),
2073                preview: None,
2074                level: crate::tool::ApprovalLevel::Dangerous,
2075                run_id: FlowRunId::now(),
2076                emitted_at: chrono::Utc::now(),
2077                bypass_auto_ceiling: false,
2078            }));
2079        }
2080        assert_eq!(reg.list_pending().len(), 3);
2081        assert_eq!(
2082            reg.decide_all(ApprovalDecision::Deny {
2083                reason: "user cancelled".into()
2084            }),
2085            3
2086        );
2087        assert!(reg.list_pending().is_empty());
2088    }
2089
2090    #[test]
2091    fn compact_review_registry_auto_accepts_when_no_subscriber() {
2092        let reg = CompactReviewRegistry::new();
2093        let pending = PendingCompactReview {
2094            review_id: "r1".into(),
2095            summary: "gist".into(),
2096            slice_preview: String::new(),
2097            slice_count: 0,
2098            range_start: 0,
2099            range_end: 0,
2100            tokens_before: 0,
2101            emitted_at: chrono::Utc::now(),
2102        };
2103        let rx = reg.request(pending);
2104        let got = rx.blocking_recv().unwrap();
2105        assert!(matches!(got, CompactReviewDecision::AcceptAsIs));
2106        assert!(reg.list_pending().is_none());
2107    }
2108
2109    #[test]
2110    fn compact_review_registry_holds_pending_and_decides() {
2111        let reg = std::sync::Arc::new(CompactReviewRegistry::new());
2112        let _sub = reg.subscribe();
2113        let pending = PendingCompactReview {
2114            review_id: "r2".into(),
2115            summary: "old".into(),
2116            slice_preview: "slice".into(),
2117            slice_count: 3,
2118            range_start: 1,
2119            range_end: 4,
2120            tokens_before: 500,
2121            emitted_at: chrono::Utc::now(),
2122        };
2123        let mut rx = reg.request(pending);
2124        assert!(rx.try_recv().is_err(), "should be queued");
2125        assert!(reg.list_pending().is_some());
2126        assert!(reg.decide(
2127            "r2",
2128            CompactReviewDecision::AcceptEdited {
2129                summary: "new".into()
2130            }
2131        ));
2132        let got = rx.blocking_recv().unwrap();
2133        match got {
2134            CompactReviewDecision::AcceptEdited { summary } => assert_eq!(summary, "new"),
2135            other => panic!("unexpected decision: {other:?}"),
2136        }
2137        assert!(reg.list_pending().is_none());
2138    }
2139
2140    #[test]
2141    fn compact_review_registry_reject_flushes() {
2142        let reg = std::sync::Arc::new(CompactReviewRegistry::new());
2143        let _sub = reg.subscribe();
2144        let rx = reg.request(PendingCompactReview {
2145            review_id: "r3".into(),
2146            summary: String::new(),
2147            slice_preview: String::new(),
2148            slice_count: 0,
2149            range_start: 0,
2150            range_end: 0,
2151            tokens_before: 0,
2152            emitted_at: chrono::Utc::now(),
2153        });
2154        assert!(reg.decide("r3", CompactReviewDecision::Reject));
2155        let got = rx.blocking_recv().unwrap();
2156        assert!(matches!(got, CompactReviewDecision::Reject));
2157    }
2158
2159    #[test]
2160    fn replay_context_snapshot_accumulates_llm_call_usage() {
2161        let dir = TempDir::new().unwrap();
2162        let events = [
2163            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"},"run_id":"019f0000-0000-7000-0000-000000000099","ts":"2026-07-08T00:00:00Z"}"#,
2164            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"}"#,
2165            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"},"run_id":"019f0000-0000-7000-0000-000000000099","ts":"2026-07-08T00:00:01Z"}"#,
2166        ];
2167        write_events(dir.path(), &events);
2168        let snap = replay_context_snapshot_from(&dir.path().join("events.jsonl"));
2169        assert_eq!(snap.model, "anthropic/claude-4");
2170        assert_eq!(snap.tokens_in, 310);
2171        assert_eq!(snap.tokens_out, 130);
2172    }
2173
2174    #[test]
2175    fn replay_context_snapshot_skips_subagent_llm_calls() {
2176        let dir = TempDir::new().unwrap();
2177        let events = [
2178            r#"{"type":"llm_call","seq":1,"model":"zhipuai/glm-5.2","provider":"zhipu","usage":{"input":100,"cached_input":0,"output":50,"cache_write":0},"wallclock_ms":1000,"status":{"kind":"ok"},"run_id":"019f0000-0000-7000-0000-000000000099","ts":"2026-07-08T00:00:00Z"}"#,
2179            r#"{"type":"llm_call","seq":2,"model":"gpt-4o-mini","provider":"openai","usage":{"input":200,"cached_input":0,"output":80,"cache_write":0},"wallclock_ms":1000,"status":{"kind":"ok"},"run_id":null,"ts":"2026-07-08T00:00:01Z"}"#,
2180        ];
2181        write_events(dir.path(), &events);
2182        let snap = replay_context_snapshot_from(&dir.path().join("events.jsonl"));
2183        assert_eq!(snap.model, "zhipuai/glm-5.2");
2184        assert_eq!(snap.tokens_in, 100);
2185        assert_eq!(snap.tokens_out, 50);
2186    }
2187
2188    #[test]
2189    fn compact_review_mode_parses_all_variants() {
2190        assert_eq!(
2191            CompactReviewMode::parse("always"),
2192            Some(CompactReviewMode::Always)
2193        );
2194        assert_eq!(
2195            CompactReviewMode::parse("manual-only"),
2196            Some(CompactReviewMode::ManualOnly)
2197        );
2198        assert_eq!(
2199            CompactReviewMode::parse("manual_only"),
2200            Some(CompactReviewMode::ManualOnly)
2201        );
2202        assert_eq!(
2203            CompactReviewMode::parse("never"),
2204            Some(CompactReviewMode::Never)
2205        );
2206        assert_eq!(CompactReviewMode::parse(" bogus "), None);
2207    }
2208
2209    #[test]
2210    fn compact_review_mode_should_review_matrix() {
2211        assert!(CompactReviewMode::Always.should_review(false));
2212        assert!(CompactReviewMode::Always.should_review(true));
2213        assert!(!CompactReviewMode::ManualOnly.should_review(false));
2214        assert!(CompactReviewMode::ManualOnly.should_review(true));
2215        assert!(!CompactReviewMode::Never.should_review(false));
2216        assert!(!CompactReviewMode::Never.should_review(true));
2217    }
2218
2219    #[test]
2220    fn compact_review_registry_new_request_rejects_previous() {
2221        let reg = std::sync::Arc::new(CompactReviewRegistry::new());
2222        let _sub = reg.subscribe();
2223        let rx_a = reg.request(PendingCompactReview {
2224            review_id: "rA".into(),
2225            summary: String::new(),
2226            slice_preview: String::new(),
2227            slice_count: 0,
2228            range_start: 0,
2229            range_end: 0,
2230            tokens_before: 0,
2231            emitted_at: chrono::Utc::now(),
2232        });
2233        let _rx_b = reg.request(PendingCompactReview {
2234            review_id: "rB".into(),
2235            summary: String::new(),
2236            slice_preview: String::new(),
2237            slice_count: 0,
2238            range_start: 0,
2239            range_end: 0,
2240            tokens_before: 0,
2241            emitted_at: chrono::Utc::now(),
2242        });
2243        let got = rx_a.blocking_recv().unwrap();
2244        assert!(matches!(got, CompactReviewDecision::Reject));
2245    }
2246
2247    fn mk_form(form_id: &str, prompt: &str) -> crate::form::PendingForm {
2248        crate::form::PendingForm {
2249            form_id: form_id.into(),
2250            run_id: crate::event::FlowRunId::now(),
2251            tool_use_id: "tu".into(),
2252            kind: crate::form::FormKind::Confirm {
2253                prompt: prompt.into(),
2254            },
2255            emitted_at: chrono::Utc::now(),
2256        }
2257    }
2258
2259    #[test]
2260    fn form_registry_auto_cancels_without_subscriber() {
2261        let reg = FormRegistry::new();
2262        let rx = reg.request(mk_form("f1", "sure?"));
2263        let got = rx.blocking_recv().unwrap();
2264        assert_eq!(got, crate::form::FormAnswer::Cancelled);
2265        assert!(reg.list_pending().is_empty());
2266    }
2267
2268    #[test]
2269    fn form_registry_delivers_answer_by_form_id() {
2270        let reg = std::sync::Arc::new(FormRegistry::new());
2271        let _sub = reg.subscribe();
2272        let rx = reg.request(mk_form("fA", "?"));
2273        assert_eq!(reg.list_pending().len(), 1);
2274        let ok = reg.submit("fA", crate::form::FormAnswer::Confirmed { value: true });
2275        assert!(ok);
2276        let got = rx.blocking_recv().unwrap();
2277        assert_eq!(got, crate::form::FormAnswer::Confirmed { value: true });
2278        assert!(reg.list_pending().is_empty());
2279    }
2280
2281    #[test]
2282    fn form_registry_submit_unknown_id_is_noop() {
2283        let reg = std::sync::Arc::new(FormRegistry::new());
2284        let _sub = reg.subscribe();
2285        let _rx = reg.request(mk_form("real", "?"));
2286        assert!(!reg.submit("ghost", crate::form::FormAnswer::Cancelled));
2287        assert_eq!(reg.list_pending().len(), 1);
2288    }
2289
2290    #[test]
2291    fn form_registry_cancel_all_flushes_pending() {
2292        let reg = std::sync::Arc::new(FormRegistry::new());
2293        let _sub = reg.subscribe();
2294        let rx_a = reg.request(mk_form("a", "?"));
2295        let rx_b = reg.request(mk_form("b", "?"));
2296        reg.cancel_all();
2297        assert_eq!(
2298            rx_a.blocking_recv().unwrap(),
2299            crate::form::FormAnswer::Cancelled
2300        );
2301        assert_eq!(
2302            rx_b.blocking_recv().unwrap(),
2303            crate::form::FormAnswer::Cancelled
2304        );
2305        assert!(reg.list_pending().is_empty());
2306    }
2307
2308    #[test]
2309    fn form_registry_queues_multiple_pending() {
2310        let reg = std::sync::Arc::new(FormRegistry::new());
2311        let _sub = reg.subscribe();
2312        let _rx1 = reg.request(mk_form("1", "?"));
2313        let _rx2 = reg.request(mk_form("2", "?"));
2314        let pending = reg.list_pending();
2315        assert_eq!(pending.len(), 2);
2316        assert_eq!(pending[0].form_id, "1");
2317        assert_eq!(pending[1].form_id, "2");
2318    }
2319
2320    #[test]
2321    fn replay_without_degraded_events_preserves_image_parts() {
2322        let dir = TempDir::new().unwrap();
2323        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"}"#;
2324        write_events(dir.path(), &[user_msg]);
2325        let entries = replay_transcript_from(&dir.path().join("events.jsonl")).unwrap();
2326        let msg = entries
2327            .into_iter()
2328            .find_map(|e| match e {
2329                TranscriptEntry::Message { message, .. } => Some(message),
2330                _ => None,
2331            })
2332            .unwrap();
2333        assert!(matches!(
2334            msg.parts[0],
2335            crate::message::MessagePart::Image { .. }
2336        ));
2337    }
2338
2339    #[test]
2340    fn replay_messages_from_old_format_no_seq_no_ts() {
2341        let dir = TempDir::new().unwrap();
2342        // Old-style JSONL: no seq, no ts on events
2343        let user_json = r#"{"type":"user_msg","turn_id":"019f0000-0000-7000-0000-000000000001","message":{"role":"user","parts":[{"type":"text","text":"hello"}],"turn_id":"019f0000-0000-7000-0000-000000000001"}}"#;
2344        let asst_json = r#"{"type":"assistant_msg","turn_id":"019f0000-0000-7000-0000-000000000001","message":{"role":"assistant","parts":[{"type":"text","text":"hi there"}],"turn_id":"019f0000-0000-7000-0000-000000000001"},"flow_run_id":null}"#;
2345        write_events(dir.path(), &[user_json, asst_json]);
2346        let msgs = replay_messages_from(&dir.path().join("events.jsonl")).unwrap();
2347        assert_eq!(msgs.len(), 2, "should load both messages from old format");
2348        assert_eq!(msgs[0].text_concat(), "hello");
2349        assert_eq!(msgs[1].text_concat(), "hi there");
2350    }
2351
2352    #[test]
2353    fn replay_messages_from_old_format_with_null_fields() {
2354        let dir = TempDir::new().unwrap();
2355        // Old JSON with null turn_id / flow_run_id (graceful parse)
2356        let sys_json = r#"{"type":"system_msg","turn_id":"019f0000-0000-7000-0000-000000000003","message":{"role":"system","parts":[{"type":"text","text":"note"}],"turn_id":"019f0000-0000-7000-0000-000000000003"}}"#;
2357        write_events(dir.path(), &[sys_json]);
2358        let msgs = replay_messages_from(&dir.path().join("events.jsonl")).unwrap();
2359        assert_eq!(msgs.len(), 1);
2360        assert_eq!(msgs[0].text_concat(), "note");
2361    }
2362
2363    #[test]
2364    fn replay_messages_from_applies_attachment_degrade() {
2365        let dir = TempDir::new().unwrap();
2366        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"}"#;
2367        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"}"#;
2368        write_events(dir.path(), &[user_msg, degrade]);
2369        let msgs = replay_messages_from(&dir.path().join("events.jsonl")).unwrap();
2370        assert_eq!(msgs.len(), 1, "only the user message (patched)");
2371        assert_eq!(msgs[0].parts.len(), 2);
2372        match &msgs[0].parts[0] {
2373            crate::message::MessagePart::Text { text } => {
2374                assert!(text.contains("photo.png"), "expected basename: {text}");
2375                assert!(text.contains("image_too_large"), "expected reason: {text}");
2376                assert!(
2377                    text.starts_with("[attachment unavailable"),
2378                    "expected stub prefix: {text}"
2379                );
2380            }
2381            other => panic!("expected Text stub, got {other:?}"),
2382        }
2383        assert!(
2384            matches!(msgs[0].parts[1], crate::message::MessagePart::Text { .. }),
2385            "second part should remain text"
2386        );
2387    }
2388
2389    #[test]
2390    fn replay_messages_from_degrade_before_message_is_noop() {
2391        let dir = TempDir::new().unwrap();
2392        // Degrade event appears BEFORE the message it references (should not crash)
2393        let degrade = r#"{"type":"attachment_degraded","seq":1,"turn_id":null,"flow_run_id":null,"message_seq":99,"part_index":0,"file_basename":"x.png","reason":"test","ts":"2026-07-07T00:00:00Z"}"#;
2394        let user_msg = r#"{"type":"user_msg","seq":2,"turn_id":"019f0000-0000-7000-0000-000000000001","message":{"role":"user","parts":[{"type":"image","source":{"media_type":"image/png","data":{"kind":"path","path":"/tmp/x.png"}}}],"turn_id":"019f0000-0000-7000-0000-000000000001"},"ts":"2026-07-07T00:00:01Z"}"#;
2395        write_events(dir.path(), &[degrade, user_msg]);
2396        let msgs = replay_messages_from(&dir.path().join("events.jsonl")).unwrap();
2397        assert_eq!(msgs.len(), 1);
2398        // Image part preserved — degrade referenced unknown message_seq
2399        assert!(
2400            matches!(msgs[0].parts[0], crate::message::MessagePart::Image { .. }),
2401            "image should remain when degrade targets unknown seq"
2402        );
2403    }
2404
2405    #[test]
2406    fn replay_messages_from_degrade_wrong_seq_leaves_image() {
2407        let dir = TempDir::new().unwrap();
2408        let user_msg = r#"{"type":"user_msg","seq":1,"turn_id":"019f0000-0000-7000-0000-000000000001","message":{"role":"user","parts":[{"type":"image","source":{"media_type":"image/png","data":{"kind":"path","path":"/tmp/x.png"}}}],"turn_id":"019f0000-0000-7000-0000-000000000001"},"ts":"2026-07-07T00:00:00Z"}"#;
2409        // Degrade references wrong message_seq (2, but message has seq 1)
2410        let degrade = r#"{"type":"attachment_degraded","seq":2,"turn_id":null,"flow_run_id":null,"message_seq":2,"part_index":0,"file_basename":"x.png","reason":"test","ts":"2026-07-07T00:00:01Z"}"#;
2411        write_events(dir.path(), &[user_msg, degrade]);
2412        let msgs = replay_messages_from(&dir.path().join("events.jsonl")).unwrap();
2413        assert_eq!(msgs.len(), 1);
2414        assert!(
2415            matches!(msgs[0].parts[0], crate::message::MessagePart::Image { .. }),
2416            "image should remain when degrade targets wrong seq"
2417        );
2418    }
2419
2420    #[test]
2421    fn replay_messages_from_applies_context_compact() {
2422        let dir = TempDir::new().unwrap();
2423        let user1 = r#"{"type":"user_msg","seq":1,"turn_id":"019f0000-0000-7000-0000-000000000001","message":{"role":"user","parts":[{"type":"text","text":"old u1"}],"turn_id":"019f0000-0000-7000-0000-000000000001"},"ts":"2026-07-07T00:00:00Z"}"#;
2424        let asst1 = r#"{"type":"assistant_msg","seq":2,"turn_id":"019f0000-0000-7000-0000-000000000001","message":{"role":"assistant","parts":[{"type":"text","text":"old a1"}],"turn_id":"019f0000-0000-7000-0000-000000000001"},"flow_run_id":null,"ts":"2026-07-07T00:00:01Z"}"#;
2425        let user2 = r#"{"type":"user_msg","seq":3,"turn_id":"019f0000-0000-7000-0000-000000000002","message":{"role":"user","parts":[{"type":"text","text":"old u2"}],"turn_id":"019f0000-0000-7000-0000-000000000002"},"ts":"2026-07-07T00:00:02Z"}"#;
2426        // replacement message (compact summary)
2427        let summary = r#"{"type":"system_msg","seq":4,"turn_id":"019f0000-0000-7000-0000-000000000002","message":{"role":"system","parts":[{"type":"compact_summary","summary":"two messages compacted","seq_start":0,"seq_end":1,"count":2}],"turn_id":"019f0000-0000-7000-0000-000000000002"},"ts":"2026-07-07T00:00:03Z"}"#;
2428        let compact = r#"{"type":"context_compact","seq":5,"session_id":"sess","before_tokens":200,"after_tokens":50,"compacted_range_start":0,"compacted_range_end":1,"summary_text":"two messages compacted","replacement_msg_seq":4,"ts":"2026-07-07T00:00:04Z"}"#;
2429        let after = r#"{"type":"user_msg","seq":6,"turn_id":"019f0000-0000-7000-0000-000000000003","message":{"role":"user","parts":[{"type":"text","text":"after compact"}],"turn_id":"019f0000-0000-7000-0000-000000000003"},"ts":"2026-07-07T00:00:05Z"}"#;
2430        write_events(dir.path(), &[user1, asst1, user2, summary, compact, after]);
2431        let msgs = replay_messages_from(&dir.path().join("events.jsonl")).unwrap();
2432        // compact range 0-1 removes user1+asst1; user2 (outside range) + summary + after = 3
2433        assert_eq!(msgs.len(), 3, "compact summary + user2 + after compact");
2434        assert!(
2435            matches!(
2436                msgs[0].parts[0],
2437                crate::message::MessagePart::CompactSummary { .. }
2438            ),
2439            "first should be compact summary"
2440        );
2441        if let crate::message::MessagePart::CompactSummary { summary, .. } = &msgs[0].parts[0] {
2442            assert_eq!(summary, "two messages compacted");
2443        }
2444        assert_eq!(msgs[1].text_concat(), "old u2");
2445        assert_eq!(msgs[2].text_concat(), "after compact");
2446    }
2447
2448    #[test]
2449    fn replay_messages_from_no_replacement_seq_ignores_compact() {
2450        let dir = TempDir::new().unwrap();
2451        let user1 = r#"{"type":"user_msg","seq":1,"turn_id":"019f0000-0000-7000-0000-000000000001","message":{"role":"user","parts":[{"type":"text","text":"hello"}],"turn_id":"019f0000-0000-7000-0000-000000000001"},"ts":"2026-07-07T00:00:00Z"}"#;
2452        // context_compact with no replacement_msg_seq → should be ignored
2453        let compact = r#"{"type":"context_compact","seq":2,"session_id":"sess","before_tokens":200,"after_tokens":50,"compacted_range_start":0,"compacted_range_end":0,"summary_text":"ignored","replacement_msg_seq":null,"ts":"2026-07-07T00:00:01Z"}"#;
2454        write_events(dir.path(), &[user1, compact]);
2455        let msgs = replay_messages_from(&dir.path().join("events.jsonl")).unwrap();
2456        assert_eq!(msgs.len(), 1, "compact without replacement seq is ignored");
2457        assert_eq!(msgs[0].text_concat(), "hello");
2458    }
2459
2460    #[test]
2461    fn replay_messages_from_compact_after_no_change_ignored() {
2462        let dir = TempDir::new().unwrap();
2463        let user1 = r#"{"type":"user_msg","seq":1,"turn_id":"019f0000-0000-7000-0000-000000000001","message":{"role":"user","parts":[{"type":"text","text":"hello"}],"turn_id":"019f0000-0000-7000-0000-000000000001"},"ts":"2026-07-07T00:00:00Z"}"#;
2464        let summary = r#"{"type":"system_msg","seq":2,"turn_id":"019f0000-0000-7000-0000-000000000001","message":{"role":"system","parts":[{"type":"compact_summary","summary":"no change","seq_start":0,"seq_end":0,"count":1}],"turn_id":"019f0000-0000-7000-0000-000000000001"},"ts":"2026-07-07T00:00:01Z"}"#;
2465        // after_tokens >= before_tokens → compaction is no-op
2466        let compact = r#"{"type":"context_compact","seq":3,"session_id":"sess","before_tokens":50,"after_tokens":100,"compacted_range_start":0,"compacted_range_end":0,"summary_text":"no change","replacement_msg_seq":2,"ts":"2026-07-07T00:00:02Z"}"#;
2467        write_events(dir.path(), &[user1, summary, compact]);
2468        let msgs = replay_messages_from(&dir.path().join("events.jsonl")).unwrap();
2469        assert_eq!(msgs.len(), 2, "compact with after>=before is ignored");
2470    }
2471
2472    #[test]
2473    fn replay_messages_from_missing_file_returns_empty() {
2474        let dir = TempDir::new().unwrap();
2475        let msgs = replay_messages_from(&dir.path().join("events.jsonl")).unwrap();
2476        assert!(msgs.is_empty());
2477    }
2478
2479    #[test]
2480    fn replay_messages_from_empty_file_returns_empty() {
2481        let dir = TempDir::new().unwrap();
2482        write_events(dir.path(), &[]);
2483        let msgs = replay_messages_from(&dir.path().join("events.jsonl")).unwrap();
2484        assert!(msgs.is_empty());
2485    }
2486
2487    #[test]
2488    fn replay_all_messages_with_seq_includes_compacted() {
2489        let dir = TempDir::new().unwrap();
2490        let user1 = r#"{"type":"user_msg","seq":1,"turn_id":"019f0000-0000-7000-0000-000000000001","message":{"role":"user","parts":[{"type":"text","text":"old"}],"turn_id":"019f0000-0000-7000-0000-000000000001"},"ts":"2026-07-07T00:00:00Z"}"#;
2491        let summary = r#"{"type":"system_msg","seq":4,"turn_id":"019f0000-0000-7000-0000-000000000002","message":{"role":"system","parts":[{"type":"compact_summary","summary":"s","seq_start":0,"seq_end":0,"count":1}],"turn_id":"019f0000-0000-7000-0000-000000000002"},"ts":"2026-07-07T00:00:01Z"}"#;
2492        let compact = r#"{"type":"context_compact","seq":5,"session_id":"sess","before_tokens":200,"after_tokens":50,"compacted_range_start":0,"compacted_range_end":0,"summary_text":"s","replacement_msg_seq":4,"ts":"2026-07-07T00:00:02Z"}"#;
2493        write_events(dir.path(), &[user1, summary, compact]);
2494        let all = replay_all_messages_with_seq(&dir.path().join("events.jsonl")).unwrap();
2495        // all includes both user1 and summary — compaction NOT applied
2496        assert_eq!(all.len(), 2, "all messages preserved (no compaction)");
2497        assert_eq!(all[0].1.text_concat(), "old");
2498    }
2499}