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