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 compact_messages(
1307 &self,
1308 summary: String,
1309 range: crate::compaction::CompactRange,
1310 before_tokens: u64,
1311 ) -> Option<CompactResult> {
1312 use crate::compaction::{estimate_tokens_for_messages, replace_range_with_summary};
1313 let msgs = self.messages();
1314 let turn_id = msgs
1315 .get(range.start)
1316 .map(|m| m.turn_id.clone())
1317 .unwrap_or_else(TurnId::now);
1318 let after = replace_range_with_summary(&msgs, &range, summary.clone(), turn_id.clone());
1319 let after_tokens = estimate_tokens_for_messages(&after);
1320 if after_tokens >= before_tokens {
1321 self.push_system_note(format!(
1322 "compaction skipped: summary would not shrink transcript ({} >= {} tokens)",
1323 after_tokens, before_tokens
1324 ));
1325 return None;
1326 }
1327 let replacement_msg = after.first().cloned().unwrap_or_else(|| {
1328 Message::system_compact_summary(
1329 turn_id.clone(),
1330 summary.clone(),
1331 range.start as u64,
1332 range.end.saturating_sub(1) as u64,
1333 range.end - range.start,
1334 )
1335 });
1336 self.sink.mark_compacted();
1337 let replacement_seq = self.sink.next_seq_peek();
1338 self.sink.emit(Event::SystemMsg {
1339 turn_id: turn_id.clone(),
1340 message: replacement_msg,
1341 });
1342 self.sink.emit(Event::ContextCompact {
1343 session_id: self.id.to_string(),
1344 before_tokens,
1345 after_tokens,
1346 compacted_range_start: range.start as u64,
1347 compacted_range_end: range.end.saturating_sub(1) as u64,
1348 summary_text: Some(summary.clone()),
1349 replacement_msg_seq: Some(replacement_seq),
1350 });
1351 self.sink.emit(Event::CompactionSummary {
1352 session_id: self.id.to_string(),
1353 range_start: range.start as u64,
1354 range_end: range.end.saturating_sub(1) as u64,
1355 compacted_count: range.end - range.start,
1356 before_tokens,
1357 after_tokens,
1358 summary: summary.clone(),
1359 });
1360 let _ = self
1361 .watch
1362 .stream_tx
1363 .send(crate::stream::StreamFrame::CompactionSummary {
1364 phase: crate::stream::CompactionPhase::Finished,
1365 range_start: range.start,
1366 range_end: range.end.saturating_sub(1),
1367 summary,
1368 before_tokens,
1369 after_tokens,
1370 compacted_count: range.end - range.start,
1371 });
1372 let checkpoint_messages = self.messages();
1373 let window_tokens = estimate_tokens_for_messages(&checkpoint_messages);
1374 self.compaction
1375 .last_input_tokens
1376 .store(window_tokens, std::sync::atomic::Ordering::Relaxed);
1377 let window_owned = checkpoint_messages.to_vec();
1381 if let Ok(mut vec) = self.messages.lock() {
1382 *vec = window_owned;
1383 }
1384 self.refresh_window_snapshot();
1385 self.sink.emit(Event::Checkpoint {
1386 session_id: self.id.to_string(),
1387 messages: checkpoint_messages.to_vec(),
1388 window_tokens,
1389 });
1390 Some(CompactResult {
1391 before_tokens,
1392 after_tokens,
1393 compacted_start: range.start,
1394 compacted_end: range.end,
1395 })
1396 }
1397
1398 pub fn begin_turn(&self, user_msg: Message) -> TurnId {
1399 BeginTurnCommand { user_msg }.execute(self)
1400 }
1401
1402 pub fn mark_streamed(&self) {
1403 self.turn
1404 .streamed
1405 .store(true, std::sync::atomic::Ordering::Relaxed);
1406 }
1407
1408 pub fn take_streamed_flag(&self) -> bool {
1409 self.turn
1410 .streamed
1411 .swap(false, std::sync::atomic::Ordering::Relaxed)
1412 }
1413
1414 pub fn end_turn(&self) {
1415 self.turn
1416 .streamed
1417 .store(false, std::sync::atomic::Ordering::Relaxed);
1418 let turn_id = self.turn.current_turn.lock().unwrap().take();
1419 if let Some(turn_id) = turn_id {
1420 let mut q = self.injection_queue.lock().unwrap();
1421 for inj in q.iter_mut() {
1422 if inj.state == InjectionState::Pending && inj.turn_id == turn_id {
1423 inj.state = InjectionState::Cancelled;
1424 let _ = self.injection_tx.send(inj.clone());
1425 }
1426 }
1427 drop(q);
1428 self.sink.emit(Event::TurnEnd { turn_id });
1429 }
1430 }
1431
1432 pub fn current_turn(&self) -> Option<TurnId> {
1433 self.turn.current_turn.lock().unwrap().clone()
1434 }
1435
1436 pub fn enqueue_injection(&self, text: impl Into<String>) -> Result<InjectionId, EnqueueError> {
1437 self.enqueue_injection_with_level(text, crate::injection::InjectionLevel::L1Nudge, None)
1438 }
1439
1440 pub fn enqueue_injection_with_level(
1441 &self,
1442 text: impl Into<String>,
1443 level: crate::injection::InjectionLevel,
1444 redirect_target: Option<String>,
1445 ) -> Result<InjectionId, EnqueueError> {
1446 let turn_id = self
1447 .turn
1448 .current_turn
1449 .lock()
1450 .unwrap()
1451 .clone()
1452 .ok_or(EnqueueError::NoActiveTurn)?;
1453 let inj = Injection::with_level(turn_id.clone(), text, level, redirect_target);
1454 let id = inj.id.clone();
1455 self.sink.emit(Event::UserInject {
1456 turn_id,
1457 injection: inj.clone(),
1458 });
1459 self.injection_queue.lock().unwrap().push(inj.clone());
1460 let _ = self.injection_tx.send(inj);
1461 Ok(id)
1462 }
1463
1464 pub fn subscribe_injections(&self) -> broadcast::Receiver<Injection> {
1465 self.injection_tx.subscribe()
1466 }
1467
1468 pub fn mark_injection_consumed(&self, id: &InjectionId) {
1469 let mut q = self.injection_queue.lock().unwrap();
1470 for inj in q.iter_mut() {
1471 if inj.id == *id && inj.state == InjectionState::Pending {
1472 inj.state = InjectionState::Injected;
1473 let _ = self.injection_tx.send(inj.clone());
1474 return;
1475 }
1476 }
1477 }
1478
1479 pub fn peek_pending_l2_or_higher(&self, turn_id: &TurnId) -> Option<Injection> {
1480 let q = self.injection_queue.lock().unwrap();
1481 q.iter()
1482 .find(|i| {
1483 i.state == InjectionState::Pending
1484 && i.turn_id == *turn_id
1485 && !matches!(i.level, crate::injection::InjectionLevel::L1Nudge)
1486 })
1487 .cloned()
1488 }
1489
1490 pub fn drain_injections(&self, turn_id: &TurnId) -> Vec<Injection> {
1493 let mut q = self.injection_queue.lock().unwrap();
1494 let mut out = Vec::new();
1495 for inj in q.iter_mut() {
1496 if inj.state == InjectionState::Pending && inj.turn_id == *turn_id {
1497 inj.state = InjectionState::Injected;
1498 let _ = self.injection_tx.send(inj.clone());
1499 out.push(inj.clone());
1500 }
1501 }
1502 out
1503 }
1504
1505 pub fn list_pending_injections(&self) -> Vec<Injection> {
1506 self.injection_queue
1507 .lock()
1508 .unwrap()
1509 .iter()
1510 .filter(|i| i.state == InjectionState::Pending)
1511 .cloned()
1512 .collect()
1513 }
1514
1515 pub fn cancel_flow(&self) {
1516 self.turn.flow_cancel.lock().unwrap().cancel();
1517 }
1518
1519 pub fn flow_cancel_token(&self) -> CancellationToken {
1520 self.turn.flow_cancel.lock().unwrap().clone()
1521 }
1522
1523 pub async fn shutdown(&self) {
1524 let writer = self.writer.lock().unwrap().take();
1525 if let Some(writer) = writer {
1526 writer.shutdown().await;
1527 }
1528 }
1529
1530 #[allow(clippy::await_holding_lock)]
1533 pub async fn flush_writer(&self) {
1534 let guard = self.writer.lock().unwrap();
1535 let Some(ref writer) = *guard else {
1536 return;
1537 };
1538 writer.flush().await;
1539 }
1540}
1541
1542#[derive(Debug, thiserror::Error)]
1543pub enum EnqueueError {
1544 #[error("enqueue_injection called with no active turn")]
1545 NoActiveTurn,
1546}
1547
1548pub struct AppendMessageCommand {
1549 pub msg: Message,
1550 pub flow_run_id: Option<FlowRunId>,
1551}
1552
1553impl AppendMessageCommand {
1554 pub fn execute(&self, session: &Session) -> u64 {
1555 let flow_run_id_str = self.flow_run_id.as_ref().map(|r| r.0.to_string());
1556 let msg =
1557 crate::tools::tool_output::maybe_truncate_tool_message(&self.msg, Some(&session.dir));
1558 let event = match msg.role {
1559 MessageRole::User => Event::UserMsg {
1560 turn_id: msg.turn_id.clone(),
1561 flow_run_id: self.flow_run_id.clone(),
1562 message: msg.clone(),
1563 },
1564 MessageRole::Assistant => {
1565 let _ = session
1566 .watch
1567 .stream_tx
1568 .send(crate::stream::StreamFrame::AssistantMsg {
1569 flow_run_id: flow_run_id_str.clone(),
1570 message: msg.clone(),
1571 });
1572 for source in extract_mermaid_blocks(&msg) {
1573 let _ =
1574 session
1575 .watch
1576 .stream_tx
1577 .send(crate::stream::StreamFrame::MermaidDiagram {
1578 source: source.clone(),
1579 });
1580 session
1581 .sink
1582 .emit(crate::event::Event::MermaidDiagram { source });
1583 }
1584 Event::AssistantMsg {
1585 turn_id: msg.turn_id.clone(),
1586 flow_run_id: self.flow_run_id.clone(),
1587 message: msg.clone(),
1588 }
1589 }
1590 MessageRole::Tool => {
1591 let _ = session
1592 .watch
1593 .stream_tx
1594 .send(crate::stream::StreamFrame::ToolResultMsg {
1595 flow_run_id: flow_run_id_str.clone(),
1596 message: msg.clone(),
1597 });
1598 Event::ToolResultMsg {
1599 turn_id: msg.turn_id.clone(),
1600 flow_run_id: self.flow_run_id.clone(),
1601 message: msg.clone(),
1602 }
1603 }
1604 MessageRole::System => Event::SystemMsg {
1605 turn_id: msg.turn_id.clone(),
1606 message: msg.clone(),
1607 },
1608 };
1609 let seq = session.sink.emit_returning_seq(event);
1610 if matches!(msg.role, MessageRole::User) {
1611 let images: Vec<(usize, String)> = msg
1612 .parts
1613 .iter()
1614 .enumerate()
1615 .filter_map(|(i, p)| match p {
1616 crate::message::MessagePart::Image { source } => {
1617 let basename = match &source.data {
1618 crate::message::ImageData::Path { path } => path
1619 .file_name()
1620 .and_then(|n| n.to_str())
1621 .unwrap_or("unknown")
1622 .to_string(),
1623 crate::message::ImageData::Base64 { .. } => "base64".into(),
1624 };
1625 Some((i, basename))
1626 }
1627 _ => None,
1628 })
1629 .collect();
1630 if !images.is_empty() {
1631 *session.last_image_user_msg.lock().unwrap() = Some(LastImageUserMsg {
1632 message_seq: seq,
1633 images,
1634 });
1635 }
1636 }
1637 session.messages.lock().unwrap().push(msg.clone());
1638 seq
1639 }
1640}
1641
1642pub struct BeginTurnCommand {
1643 pub user_msg: Message,
1644}
1645
1646impl BeginTurnCommand {
1647 pub fn execute(&self, session: &Session) -> TurnId {
1648 let turn_id = self.user_msg.turn_id.clone();
1649 *session.turn.current_turn.lock().unwrap() = Some(turn_id.clone());
1650 *session.turn.flow_cancel.lock().unwrap() = tokio_util::sync::CancellationToken::new();
1651 session.sink.emit(Event::TurnStart {
1652 turn_id: turn_id.clone(),
1653 });
1654 AppendMessageCommand {
1655 msg: self.user_msg.clone(),
1656 flow_run_id: None,
1657 }
1658 .execute(session);
1659 turn_id
1660 }
1661}
1662
1663fn extract_mermaid_blocks(msg: &crate::message::Message) -> Vec<String> {
1664 let text = msg.text_concat();
1665 let mut blocks = Vec::new();
1666 let mut lines = text.lines().peekable();
1667 while let Some(line) = lines.next() {
1668 let trimmed = line.trim();
1669 if trimmed.starts_with("```") {
1670 let lang = trimmed.trim_start_matches("```").trim();
1671 if lang == "mermaid" {
1672 let mut source = String::new();
1673 for inner in lines.by_ref() {
1674 if inner.trim() == "```" {
1675 break;
1676 }
1677 if !source.is_empty() {
1678 source.push('\n');
1679 }
1680 source.push_str(inner);
1681 }
1682 if !source.is_empty() {
1683 blocks.push(source);
1684 }
1685 } else {
1686 for inner in lines.by_ref() {
1687 if inner.trim() == "```" {
1688 break;
1689 }
1690 }
1691 }
1692 }
1693 }
1694 blocks
1695}
1696
1697#[cfg(test)]
1698mod tests {
1699 use super::*;
1700 use tempfile::TempDir;
1701
1702 fn write_events(dir: &Path, lines: &[&str]) {
1703 let path = dir.join("events.jsonl");
1704 std::fs::write(&path, lines.join("\n") + "\n").unwrap();
1705 }
1706
1707 #[test]
1708 fn replay_applies_attachment_degraded_patch() {
1709 let dir = TempDir::new().unwrap();
1710 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"}"#;
1711 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"}"#;
1712 write_events(dir.path(), &[user_msg, degrade]);
1713 let entries = replay_transcript_from(&dir.path().join("events.jsonl")).unwrap();
1714 let msg = entries
1715 .into_iter()
1716 .find_map(|e| match e {
1717 TranscriptEntry::Message { message, .. } => Some(message),
1718 _ => None,
1719 })
1720 .unwrap();
1721 assert_eq!(msg.parts.len(), 2);
1722 match &msg.parts[0] {
1723 crate::message::MessagePart::Text { text } => {
1724 assert!(text.contains("photo.png"), "expected basename: {text}");
1725 assert!(text.contains("image_too_large"), "expected reason: {text}");
1726 assert!(text.starts_with("[attachment unavailable"));
1727 }
1728 other => panic!("expected Text stub, got {other:?}"),
1729 }
1730 assert!(matches!(
1731 msg.parts[1],
1732 crate::message::MessagePart::Text { .. }
1733 ));
1734 }
1735
1736 #[test]
1737 fn approval_registry_auto_approves_when_level_leq_ceiling() {
1738 let reg = ApprovalRegistry::new();
1739 reg.set_auto_ceiling(crate::tool::ApprovalLevel::Approve);
1740 let pending = PendingApproval {
1741 tool_use_id: "tu1".into(),
1742 tool_name: "fs.read".into(),
1743 args_preview: "{}".into(),
1744 preview: None,
1745 level: crate::tool::ApprovalLevel::Auto,
1746 run_id: FlowRunId::now(),
1747 emitted_at: chrono::Utc::now(),
1748 bypass_auto_ceiling: false,
1749 };
1750 let rx = reg.request(pending);
1751 let got = rx.blocking_recv().unwrap();
1752 assert!(matches!(got, ApprovalDecision::Approve));
1753 assert!(reg.list_pending().is_empty());
1754 }
1755
1756 #[test]
1757 fn approval_registry_queues_when_level_above_ceiling() {
1758 let reg = std::sync::Arc::new(ApprovalRegistry::new());
1759 reg.set_auto_ceiling(crate::tool::ApprovalLevel::Auto);
1760 let pending = PendingApproval {
1761 tool_use_id: "tu42".into(),
1762 tool_name: "fs.write".into(),
1763 args_preview: "{}".into(),
1764 preview: None,
1765 level: crate::tool::ApprovalLevel::Approve,
1766 run_id: FlowRunId::now(),
1767 emitted_at: chrono::Utc::now(),
1768 bypass_auto_ceiling: false,
1769 };
1770 let mut rx = reg.request(pending);
1771 assert_eq!(reg.list_pending().len(), 1);
1772 assert!(rx.try_recv().is_err(), "should still be queued");
1773 assert!(reg.decide("tu42", ApprovalDecision::Approve));
1774 let got = rx.blocking_recv().unwrap();
1775 assert!(matches!(got, ApprovalDecision::Approve));
1776 assert!(reg.list_pending().is_empty());
1777 }
1778
1779 #[test]
1780 fn approval_registry_decide_all_flushes_queue() {
1781 let reg = ApprovalRegistry::new();
1782 reg.set_auto_ceiling(crate::tool::ApprovalLevel::Auto);
1783 let mut rxs = Vec::new();
1784 for i in 0..3 {
1785 rxs.push(reg.request(PendingApproval {
1786 tool_use_id: format!("tu{i}"),
1787 tool_name: "bash.exec".into(),
1788 args_preview: "{}".into(),
1789 preview: None,
1790 level: crate::tool::ApprovalLevel::Dangerous,
1791 run_id: FlowRunId::now(),
1792 emitted_at: chrono::Utc::now(),
1793 bypass_auto_ceiling: false,
1794 }));
1795 }
1796 assert_eq!(reg.list_pending().len(), 3);
1797 assert_eq!(
1798 reg.decide_all(ApprovalDecision::Deny {
1799 reason: "user cancelled".into()
1800 }),
1801 3
1802 );
1803 assert!(reg.list_pending().is_empty());
1804 }
1805
1806 #[test]
1807 fn compact_review_registry_auto_accepts_when_no_subscriber() {
1808 let reg = CompactReviewRegistry::new();
1809 let pending = PendingCompactReview {
1810 review_id: "r1".into(),
1811 summary: "gist".into(),
1812 slice_preview: String::new(),
1813 slice_count: 0,
1814 range_start: 0,
1815 range_end: 0,
1816 tokens_before: 0,
1817 emitted_at: chrono::Utc::now(),
1818 };
1819 let rx = reg.request(pending);
1820 let got = rx.blocking_recv().unwrap();
1821 assert!(matches!(got, CompactReviewDecision::AcceptAsIs));
1822 assert!(reg.list_pending().is_none());
1823 }
1824
1825 #[test]
1826 fn compact_review_registry_holds_pending_and_decides() {
1827 let reg = std::sync::Arc::new(CompactReviewRegistry::new());
1828 let _sub = reg.subscribe();
1829 let pending = PendingCompactReview {
1830 review_id: "r2".into(),
1831 summary: "old".into(),
1832 slice_preview: "slice".into(),
1833 slice_count: 3,
1834 range_start: 1,
1835 range_end: 4,
1836 tokens_before: 500,
1837 emitted_at: chrono::Utc::now(),
1838 };
1839 let mut rx = reg.request(pending);
1840 assert!(rx.try_recv().is_err(), "should be queued");
1841 assert!(reg.list_pending().is_some());
1842 assert!(reg.decide(
1843 "r2",
1844 CompactReviewDecision::AcceptEdited {
1845 summary: "new".into()
1846 }
1847 ));
1848 let got = rx.blocking_recv().unwrap();
1849 match got {
1850 CompactReviewDecision::AcceptEdited { summary } => assert_eq!(summary, "new"),
1851 other => panic!("unexpected decision: {other:?}"),
1852 }
1853 assert!(reg.list_pending().is_none());
1854 }
1855
1856 #[test]
1857 fn compact_review_registry_reject_flushes() {
1858 let reg = std::sync::Arc::new(CompactReviewRegistry::new());
1859 let _sub = reg.subscribe();
1860 let rx = reg.request(PendingCompactReview {
1861 review_id: "r3".into(),
1862 summary: String::new(),
1863 slice_preview: String::new(),
1864 slice_count: 0,
1865 range_start: 0,
1866 range_end: 0,
1867 tokens_before: 0,
1868 emitted_at: chrono::Utc::now(),
1869 });
1870 assert!(reg.decide("r3", CompactReviewDecision::Reject));
1871 let got = rx.blocking_recv().unwrap();
1872 assert!(matches!(got, CompactReviewDecision::Reject));
1873 }
1874
1875 #[test]
1876 fn replay_context_snapshot_accumulates_llm_call_usage() {
1877 let dir = TempDir::new().unwrap();
1878 let events = [
1879 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"}"#,
1880 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"}"#,
1881 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"}"#,
1882 ];
1883 write_events(dir.path(), &events);
1884 let snap = replay_context_snapshot_from(&dir.path().join("events.jsonl"));
1885 assert_eq!(snap.model, "anthropic/claude-4");
1886 assert_eq!(snap.tokens_in, 310);
1887 assert_eq!(snap.tokens_out, 130);
1888 }
1889
1890 #[test]
1891 fn replay_context_snapshot_skips_subagent_llm_calls() {
1892 let dir = TempDir::new().unwrap();
1893 let events = [
1894 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"}"#,
1895 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"}"#,
1896 ];
1897 write_events(dir.path(), &events);
1898 let snap = replay_context_snapshot_from(&dir.path().join("events.jsonl"));
1899 assert_eq!(snap.model, "zhipuai/glm-5.2");
1900 assert_eq!(snap.tokens_in, 100);
1901 assert_eq!(snap.tokens_out, 50);
1902 }
1903
1904 #[test]
1905 fn compact_review_mode_parses_all_variants() {
1906 assert_eq!(
1907 CompactReviewMode::parse("always"),
1908 Some(CompactReviewMode::Always)
1909 );
1910 assert_eq!(
1911 CompactReviewMode::parse("manual-only"),
1912 Some(CompactReviewMode::ManualOnly)
1913 );
1914 assert_eq!(
1915 CompactReviewMode::parse("manual_only"),
1916 Some(CompactReviewMode::ManualOnly)
1917 );
1918 assert_eq!(
1919 CompactReviewMode::parse("never"),
1920 Some(CompactReviewMode::Never)
1921 );
1922 assert_eq!(CompactReviewMode::parse(" bogus "), None);
1923 }
1924
1925 #[test]
1926 fn compact_review_mode_should_review_matrix() {
1927 assert!(CompactReviewMode::Always.should_review(false));
1928 assert!(CompactReviewMode::Always.should_review(true));
1929 assert!(!CompactReviewMode::ManualOnly.should_review(false));
1930 assert!(CompactReviewMode::ManualOnly.should_review(true));
1931 assert!(!CompactReviewMode::Never.should_review(false));
1932 assert!(!CompactReviewMode::Never.should_review(true));
1933 }
1934
1935 #[test]
1936 fn compact_review_registry_new_request_rejects_previous() {
1937 let reg = std::sync::Arc::new(CompactReviewRegistry::new());
1938 let _sub = reg.subscribe();
1939 let rx_a = reg.request(PendingCompactReview {
1940 review_id: "rA".into(),
1941 summary: String::new(),
1942 slice_preview: String::new(),
1943 slice_count: 0,
1944 range_start: 0,
1945 range_end: 0,
1946 tokens_before: 0,
1947 emitted_at: chrono::Utc::now(),
1948 });
1949 let _rx_b = reg.request(PendingCompactReview {
1950 review_id: "rB".into(),
1951 summary: String::new(),
1952 slice_preview: String::new(),
1953 slice_count: 0,
1954 range_start: 0,
1955 range_end: 0,
1956 tokens_before: 0,
1957 emitted_at: chrono::Utc::now(),
1958 });
1959 let got = rx_a.blocking_recv().unwrap();
1960 assert!(matches!(got, CompactReviewDecision::Reject));
1961 }
1962
1963 fn mk_form(form_id: &str, prompt: &str) -> crate::form::PendingForm {
1964 crate::form::PendingForm {
1965 form_id: form_id.into(),
1966 run_id: crate::event::FlowRunId::now(),
1967 tool_use_id: "tu".into(),
1968 kind: crate::form::FormKind::Confirm {
1969 prompt: prompt.into(),
1970 },
1971 emitted_at: chrono::Utc::now(),
1972 }
1973 }
1974
1975 #[test]
1976 fn form_registry_auto_cancels_without_subscriber() {
1977 let reg = FormRegistry::new();
1978 let rx = reg.request(mk_form("f1", "sure?"));
1979 let got = rx.blocking_recv().unwrap();
1980 assert_eq!(got, crate::form::FormAnswer::Cancelled);
1981 assert!(reg.list_pending().is_empty());
1982 }
1983
1984 #[test]
1985 fn form_registry_delivers_answer_by_form_id() {
1986 let reg = std::sync::Arc::new(FormRegistry::new());
1987 let _sub = reg.subscribe();
1988 let rx = reg.request(mk_form("fA", "?"));
1989 assert_eq!(reg.list_pending().len(), 1);
1990 let ok = reg.submit("fA", crate::form::FormAnswer::Confirmed { value: true });
1991 assert!(ok);
1992 let got = rx.blocking_recv().unwrap();
1993 assert_eq!(got, crate::form::FormAnswer::Confirmed { value: true });
1994 assert!(reg.list_pending().is_empty());
1995 }
1996
1997 #[test]
1998 fn form_registry_submit_unknown_id_is_noop() {
1999 let reg = std::sync::Arc::new(FormRegistry::new());
2000 let _sub = reg.subscribe();
2001 let _rx = reg.request(mk_form("real", "?"));
2002 assert!(!reg.submit("ghost", crate::form::FormAnswer::Cancelled));
2003 assert_eq!(reg.list_pending().len(), 1);
2004 }
2005
2006 #[test]
2007 fn form_registry_cancel_all_flushes_pending() {
2008 let reg = std::sync::Arc::new(FormRegistry::new());
2009 let _sub = reg.subscribe();
2010 let rx_a = reg.request(mk_form("a", "?"));
2011 let rx_b = reg.request(mk_form("b", "?"));
2012 reg.cancel_all();
2013 assert_eq!(
2014 rx_a.blocking_recv().unwrap(),
2015 crate::form::FormAnswer::Cancelled
2016 );
2017 assert_eq!(
2018 rx_b.blocking_recv().unwrap(),
2019 crate::form::FormAnswer::Cancelled
2020 );
2021 assert!(reg.list_pending().is_empty());
2022 }
2023
2024 #[test]
2025 fn form_registry_queues_multiple_pending() {
2026 let reg = std::sync::Arc::new(FormRegistry::new());
2027 let _sub = reg.subscribe();
2028 let _rx1 = reg.request(mk_form("1", "?"));
2029 let _rx2 = reg.request(mk_form("2", "?"));
2030 let pending = reg.list_pending();
2031 assert_eq!(pending.len(), 2);
2032 assert_eq!(pending[0].form_id, "1");
2033 assert_eq!(pending[1].form_id, "2");
2034 }
2035
2036 #[test]
2037 fn replay_without_degraded_events_preserves_image_parts() {
2038 let dir = TempDir::new().unwrap();
2039 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"}"#;
2040 write_events(dir.path(), &[user_msg]);
2041 let entries = replay_transcript_from(&dir.path().join("events.jsonl")).unwrap();
2042 let msg = entries
2043 .into_iter()
2044 .find_map(|e| match e {
2045 TranscriptEntry::Message { message, .. } => Some(message),
2046 _ => None,
2047 })
2048 .unwrap();
2049 assert!(matches!(
2050 msg.parts[0],
2051 crate::message::MessagePart::Image { .. }
2052 ));
2053 }
2054
2055 #[test]
2056 fn replay_messages_from_old_format_no_seq_no_ts() {
2057 let dir = TempDir::new().unwrap();
2058 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"}}"#;
2060 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}"#;
2061 write_events(dir.path(), &[user_json, asst_json]);
2062 let msgs = replay_messages_from(&dir.path().join("events.jsonl")).unwrap();
2063 assert_eq!(msgs.len(), 2, "should load both messages from old format");
2064 assert_eq!(msgs[0].text_concat(), "hello");
2065 assert_eq!(msgs[1].text_concat(), "hi there");
2066 }
2067
2068 #[test]
2069 fn replay_messages_from_old_format_with_null_fields() {
2070 let dir = TempDir::new().unwrap();
2071 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"}}"#;
2073 write_events(dir.path(), &[sys_json]);
2074 let msgs = replay_messages_from(&dir.path().join("events.jsonl")).unwrap();
2075 assert_eq!(msgs.len(), 1);
2076 assert_eq!(msgs[0].text_concat(), "note");
2077 }
2078
2079 #[test]
2080 fn replay_messages_from_applies_attachment_degrade() {
2081 let dir = TempDir::new().unwrap();
2082 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"}"#;
2083 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"}"#;
2084 write_events(dir.path(), &[user_msg, degrade]);
2085 let msgs = replay_messages_from(&dir.path().join("events.jsonl")).unwrap();
2086 assert_eq!(msgs.len(), 1, "only the user message (patched)");
2087 assert_eq!(msgs[0].parts.len(), 2);
2088 match &msgs[0].parts[0] {
2089 crate::message::MessagePart::Text { text } => {
2090 assert!(text.contains("photo.png"), "expected basename: {text}");
2091 assert!(text.contains("image_too_large"), "expected reason: {text}");
2092 assert!(
2093 text.starts_with("[attachment unavailable"),
2094 "expected stub prefix: {text}"
2095 );
2096 }
2097 other => panic!("expected Text stub, got {other:?}"),
2098 }
2099 assert!(
2100 matches!(msgs[0].parts[1], crate::message::MessagePart::Text { .. }),
2101 "second part should remain text"
2102 );
2103 }
2104
2105 #[test]
2106 fn replay_messages_from_degrade_before_message_is_noop() {
2107 let dir = TempDir::new().unwrap();
2108 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"}"#;
2110 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"}"#;
2111 write_events(dir.path(), &[degrade, user_msg]);
2112 let msgs = replay_messages_from(&dir.path().join("events.jsonl")).unwrap();
2113 assert_eq!(msgs.len(), 1);
2114 assert!(
2116 matches!(msgs[0].parts[0], crate::message::MessagePart::Image { .. }),
2117 "image should remain when degrade targets unknown seq"
2118 );
2119 }
2120
2121 #[test]
2122 fn replay_messages_from_degrade_wrong_seq_leaves_image() {
2123 let dir = TempDir::new().unwrap();
2124 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"}"#;
2125 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"}"#;
2127 write_events(dir.path(), &[user_msg, degrade]);
2128 let msgs = replay_messages_from(&dir.path().join("events.jsonl")).unwrap();
2129 assert_eq!(msgs.len(), 1);
2130 assert!(
2131 matches!(msgs[0].parts[0], crate::message::MessagePart::Image { .. }),
2132 "image should remain when degrade targets wrong seq"
2133 );
2134 }
2135
2136 #[test]
2137 fn replay_messages_from_applies_context_compact() {
2138 let dir = TempDir::new().unwrap();
2139 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"}"#;
2140 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"}"#;
2141 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"}"#;
2142 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"}"#;
2144 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"}"#;
2145 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"}"#;
2146 write_events(dir.path(), &[user1, asst1, user2, summary, compact, after]);
2147 let msgs = replay_messages_from(&dir.path().join("events.jsonl")).unwrap();
2148 assert_eq!(msgs.len(), 3, "compact summary + user2 + after compact");
2150 assert!(
2151 matches!(
2152 msgs[0].parts[0],
2153 crate::message::MessagePart::CompactSummary { .. }
2154 ),
2155 "first should be compact summary"
2156 );
2157 if let crate::message::MessagePart::CompactSummary { summary, .. } = &msgs[0].parts[0] {
2158 assert_eq!(summary, "two messages compacted");
2159 }
2160 assert_eq!(msgs[1].text_concat(), "old u2");
2161 assert_eq!(msgs[2].text_concat(), "after compact");
2162 }
2163
2164 #[test]
2165 fn replay_messages_from_no_replacement_seq_ignores_compact() {
2166 let dir = TempDir::new().unwrap();
2167 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"}"#;
2168 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"}"#;
2170 write_events(dir.path(), &[user1, compact]);
2171 let msgs = replay_messages_from(&dir.path().join("events.jsonl")).unwrap();
2172 assert_eq!(msgs.len(), 1, "compact without replacement seq is ignored");
2173 assert_eq!(msgs[0].text_concat(), "hello");
2174 }
2175
2176 #[test]
2177 fn replay_messages_from_compact_after_no_change_ignored() {
2178 let dir = TempDir::new().unwrap();
2179 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"}"#;
2180 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"}"#;
2181 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"}"#;
2183 write_events(dir.path(), &[user1, summary, compact]);
2184 let msgs = replay_messages_from(&dir.path().join("events.jsonl")).unwrap();
2185 assert_eq!(msgs.len(), 2, "compact with after>=before is ignored");
2186 }
2187
2188 #[test]
2189 fn replay_messages_from_missing_file_returns_empty() {
2190 let dir = TempDir::new().unwrap();
2191 let msgs = replay_messages_from(&dir.path().join("events.jsonl")).unwrap();
2192 assert!(msgs.is_empty());
2193 }
2194
2195 #[test]
2196 fn replay_messages_from_empty_file_returns_empty() {
2197 let dir = TempDir::new().unwrap();
2198 write_events(dir.path(), &[]);
2199 let msgs = replay_messages_from(&dir.path().join("events.jsonl")).unwrap();
2200 assert!(msgs.is_empty());
2201 }
2202
2203 #[test]
2204 fn replay_all_messages_with_seq_includes_compacted() {
2205 let dir = TempDir::new().unwrap();
2206 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"}"#;
2207 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"}"#;
2208 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"}"#;
2209 write_events(dir.path(), &[user1, summary, compact]);
2210 let all = replay_all_messages_with_seq(&dir.path().join("events.jsonl")).unwrap();
2211 assert_eq!(all.len(), 2, "all messages preserved (no compaction)");
2213 assert_eq!(all[0].1.text_concat(), "old");
2214 }
2215}