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 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 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 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 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 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 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 #[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 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 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 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 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 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 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 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 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 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 assert_eq!(all.len(), 2, "all messages preserved (no compaction)");
2443 assert_eq!(all[0].1.text_concat(), "old");
2444 }
2445}