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