1mod config_watch;
35mod turn_scope;
36
37use std::collections::HashMap;
38use std::collections::VecDeque;
39use std::path::PathBuf;
40use std::sync::Arc;
41use std::sync::Mutex;
42
43use tokio::sync::mpsc;
44
45use crate::providers::ctx::{ExecContext, StreamContext};
46use crate::providers::{ProviderFactory, StreamEvent, ToolRegistry};
47use mermaid_domain::{
48 Cmd, CompactionRequest, CompactionResult, CompactionTrigger, Msg, Query, QueryResult, TurnId,
49};
50use mermaid_domain::{Config, MemoryConfig};
51
52pub use turn_scope::TurnScope;
53
54#[cfg(not(test))]
55const CANCEL_DRAIN_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(2);
56#[cfg(test)]
57const CANCEL_DRAIN_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(50);
58
59const CANCELLED_TOMBSTONE_CAP: usize = 256;
65
66pub type MsgSender = mpsc::Sender<Msg>;
72
73pub const MSG_CHANNEL_CAPACITY: usize = 512;
79
80fn compaction_row(
92 record: &mermaid_domain::CompactionEvent,
93 archive_path: &std::path::Path,
94 task_id: Option<String>,
95 session_id: String,
96) -> mermaid_runtime::NewCompaction {
97 mermaid_runtime::NewCompaction {
98 id: Some(record.id.clone()),
99 task_id,
100 session_id: Some(session_id),
101 source_token_estimate: Some(record.before_tokens as i64),
102 summary_token_count: Some(record.summary_tokens as i64),
103 preserved_turns: Some(record.preserved_turn_count as i64),
104 archive_path: Some(archive_path.display().to_string()),
105 verification_status: Some(record.review_status.as_str().to_string()),
106 }
107}
108
109fn upsert_session_index(
117 manager: &crate::session::ConversationManager,
118 snapshot: &mermaid_domain::ConversationHistory,
119) {
120 let row = crate::session::session_row(manager.conversations_dir(), snapshot);
121 let _ = mermaid_runtime::with_shared_store(|store| store.sessions().upsert(row));
122}
123
124#[derive(Clone)]
125enum PersistenceJob {
126 Conversation {
127 snapshot: Box<mermaid_domain::ConversationHistory>,
128 events: Vec<mermaid_domain::SessionEvent>,
129 },
130 Compaction(Box<PendingCompactionSave>),
131}
132
133#[derive(Clone)]
134struct PendingCompactionSave {
135 record: mermaid_domain::CompactionEvent,
136 conversation: mermaid_domain::ConversationHistory,
137 events: Vec<mermaid_domain::SessionEvent>,
138 events_appended: bool,
142 task_id: Option<String>,
143}
144
145struct PersistedCompaction {
146 id: String,
147 task_id: Option<String>,
148 session_id: String,
149 archive_path: PathBuf,
150}
151
152const CHECKPOINT_EVERY_EVENTS: usize = 200;
161
162struct PersistenceState {
163 workdir: PathBuf,
164 manager: Option<crate::session::ConversationManager>,
165 blocked: HashMap<String, VecDeque<PendingCompactionSave>>,
166 dirty: HashMap<String, DirtySession>,
170 unappended: HashMap<String, Vec<mermaid_domain::SessionEvent>>,
175}
176
177struct DirtySession {
178 snapshot: mermaid_domain::ConversationHistory,
179 events_since_checkpoint: usize,
180}
181
182impl PersistenceState {
183 fn new(workdir: PathBuf) -> Self {
184 Self {
185 workdir,
186 manager: None,
187 blocked: HashMap::new(),
188 dirty: HashMap::new(),
189 unappended: HashMap::new(),
190 }
191 }
192
193 fn manager(&mut self) -> anyhow::Result<&crate::session::ConversationManager> {
194 if self.manager.is_none() {
195 self.manager = Some(crate::session::ConversationManager::new(&self.workdir)?);
196 }
197 Ok(self.manager.as_ref().expect("manager initialized"))
198 }
199
200 fn process(&mut self, job: PersistenceJob) -> (Vec<PersistedCompaction>, anyhow::Result<()>) {
205 match job {
206 PersistenceJob::Conversation { snapshot, events } => {
207 let (persisted, retried) = self.retry_blocked(&snapshot.id);
210 if retried.is_err() {
211 return (persisted, retried);
212 }
213 let saved = self.save_session(*snapshot, events);
214 (persisted, saved)
215 },
216 PersistenceJob::Compaction(save) => {
217 let conversation_id = save.conversation.id.clone();
223 self.blocked
224 .entry(conversation_id.clone())
225 .or_default()
226 .push_back(*save);
227 self.retry_blocked(&conversation_id)
228 },
229 }
230 }
231
232 fn retry_blocked(
233 &mut self,
234 conversation_id: &str,
235 ) -> (Vec<PersistedCompaction>, anyhow::Result<()>) {
236 let mut persisted = Vec::new();
237 if !self.blocked.contains_key(conversation_id) {
238 return (persisted, Ok(()));
239 }
240 if let Err(error) = self.manager() {
241 return (persisted, Err(error));
242 }
243 let manager = self.manager.as_ref().expect("manager initialized");
247 let dirty = &mut self.dirty;
248 let queue = self
249 .blocked
250 .get_mut(conversation_id)
251 .expect("checked above");
252 while let Some(save) = queue.front_mut() {
253 match Self::persist_compaction(manager, dirty, save) {
254 Ok(event) => {
258 persisted.push(event);
259 queue.pop_front();
260 },
261 Err(error) => return (persisted, Err(error)),
262 }
263 }
264 self.blocked.remove(conversation_id);
265 (persisted, Ok(()))
266 }
267
268 fn save_session(
277 &mut self,
278 snapshot: mermaid_domain::ConversationHistory,
279 events: Vec<mermaid_domain::SessionEvent>,
280 ) -> anyhow::Result<()> {
281 let id = snapshot.id.clone();
282 let mut batch = self.unappended.remove(&id).unwrap_or_default();
284 batch.extend(events);
285
286 let manager = self.manager()?;
287 if let Err(error) = manager.append_session_events(&snapshot, &batch) {
288 tracing::warn!(
289 id = %id,
290 pending = batch.len(),
291 %error,
292 "session event append failed; holding the events for the next save"
293 );
294 self.unappended.insert(id, batch);
295 return Err(error);
296 }
297 upsert_session_index(manager, &snapshot);
298
299 let entry = self
300 .dirty
301 .entry(id.clone())
302 .or_insert_with(|| DirtySession {
303 snapshot: snapshot.clone(),
304 events_since_checkpoint: 0,
305 });
306 entry.snapshot = snapshot;
307 entry.events_since_checkpoint += batch.len();
308 if entry.events_since_checkpoint >= CHECKPOINT_EVERY_EVENTS {
309 return self.write_checkpoint(&id);
310 }
311 Ok(())
312 }
313
314 fn write_checkpoint(&mut self, id: &str) -> anyhow::Result<()> {
317 let Some(dirty) = self.dirty.remove(id) else {
318 return Ok(());
319 };
320 let manager = self.manager()?;
321 if let Err(error) = manager.save_conversation(&dirty.snapshot) {
322 self.dirty.insert(id.to_string(), dirty);
326 return Err(error);
327 }
328 Ok(())
329 }
330
331 fn flush_checkpoints(&mut self) -> anyhow::Result<()> {
335 let ids: Vec<String> = self.dirty.keys().cloned().collect();
336 let mut first_error = None;
337 for id in ids {
338 if let Err(error) = self.write_checkpoint(&id) {
339 first_error.get_or_insert(error);
340 }
341 }
342 first_error.map_or(Ok(()), Err)
343 }
344
345 fn retry_all_blocked(&mut self) -> (Vec<PersistedCompaction>, anyhow::Result<()>) {
346 let ids: Vec<String> = self.blocked.keys().cloned().collect();
347 let mut persisted = Vec::new();
348 let mut first_error = None;
349 for id in ids {
350 let (events, result) = self.retry_blocked(&id);
353 persisted.extend(events);
354 if let Err(error) = result {
355 first_error.get_or_insert(error);
356 }
357 }
358 match first_error {
359 None => (persisted, Ok(())),
360 Some(error) => (persisted, Err(error)),
361 }
362 }
363
364 fn persist_compaction(
365 manager: &crate::session::ConversationManager,
366 dirty: &mut HashMap<String, DirtySession>,
367 save: &mut PendingCompactionSave,
368 ) -> anyhow::Result<PersistedCompaction> {
369 if !save.events_appended {
378 manager.append_session_events(&save.conversation, &save.events)?;
379 save.events_appended = true;
380 }
381 manager.save_conversation(&save.conversation)?;
382 upsert_session_index(manager, &save.conversation);
383 dirty.remove(&save.conversation.id);
387
388 let log_path = manager.event_log_path(&save.conversation.id);
389 let _ = mermaid_runtime::with_shared_store(|store| {
390 store.compactions().create(compaction_row(
391 &save.record,
392 &log_path,
393 save.task_id.clone(),
394 save.conversation.id.clone(),
395 ))
396 });
397
398 Ok(PersistedCompaction {
399 id: save.record.id.clone(),
400 task_id: save.task_id.clone(),
401 session_id: save.conversation.id.clone(),
402 archive_path: log_path,
403 })
404 }
405}
406
407async fn fire_compaction_hook(event: &PersistedCompaction) {
409 fire_plugin_hooks(
410 "compaction",
411 serde_json::json!({
412 "id": event.id,
413 "task_id": event.task_id,
414 "session_id": event.session_id,
415 "archive_path": event.archive_path.display().to_string(),
416 }),
417 )
418 .await;
419}
420
421pub struct EffectRunner {
424 msg_tx: MsgSender,
425 scopes: HashMap<TurnId, TurnScope>,
431 cancelled_turns: VecDeque<TurnId>,
437 detached: tokio::task::JoinSet<()>,
441 persistence_state: Arc<Mutex<PersistenceState>>,
445 persistence_tail: Option<tokio::task::JoinHandle<()>>,
446 workdir: PathBuf,
450 providers: Option<Arc<ProviderFactory>>,
455 tools: Option<Arc<ToolRegistry>>,
458 task_id: Option<String>,
460 terminal_title_enabled: bool,
464 owns_global_mcp: bool,
470 approval: Option<crate::providers::ApprovalBroker>,
474 questions: Option<crate::providers::QuestionBroker>,
478 tasks: crate::providers::TaskBroker,
483 config_watch: Option<tokio::task::AbortHandle>,
487}
488
489impl EffectRunner {
490 #[must_use]
492 pub fn new(msg_tx: MsgSender, workdir: PathBuf) -> Self {
493 let persistence_state = Arc::new(Mutex::new(PersistenceState::new(workdir.clone())));
494 Self {
495 tasks: crate::providers::TaskBroker::new(msg_tx.clone()),
496 msg_tx,
497 scopes: HashMap::new(),
498 cancelled_turns: VecDeque::new(),
499 detached: tokio::task::JoinSet::new(),
500 persistence_state,
501 persistence_tail: None,
502 workdir,
503 providers: None,
504 tools: None,
505 task_id: None,
506 terminal_title_enabled: true,
507 owns_global_mcp: true,
508 approval: None,
509 questions: None,
510 config_watch: None,
511 }
512 }
513
514 #[must_use]
520 pub fn sender(&self) -> MsgSender {
521 self.msg_tx.clone()
522 }
523
524 #[must_use]
528 pub fn with_interactive_approvals(mut self) -> Self {
529 self.approval = Some(crate::providers::ApprovalBroker::new(self.msg_tx.clone()));
530 self
531 }
532
533 #[must_use]
537 pub fn with_interactive_questions(mut self) -> Self {
538 self.questions = Some(crate::providers::QuestionBroker::new(self.msg_tx.clone()));
539 self
540 }
541
542 pub fn spawn_config_watcher(&mut self, cwd: PathBuf, memory: MemoryConfig) {
548 let handle = self.detached.spawn(config_watch::config_watcher(
549 self.msg_tx.clone(),
550 cwd,
551 memory,
552 ));
553 self.config_watch = Some(handle);
554 }
555
556 #[must_use]
559 pub fn with_task_id(mut self, task_id: Option<String>) -> Self {
560 self.task_id = task_id;
561 self
562 }
563
564 #[must_use]
566 pub fn without_terminal_title(mut self) -> Self {
567 self.terminal_title_enabled = false;
568 self
569 }
570
571 #[must_use]
574 pub fn without_global_mcp_shutdown(mut self) -> Self {
575 self.owns_global_mcp = false;
576 self
577 }
578
579 pub fn with_bindings(
584 mut self,
585 providers: Arc<ProviderFactory>,
586 tools: Arc<ToolRegistry>,
587 ) -> Self {
588 self.providers = Some(providers);
589 self.tools = Some(tools);
590 self
591 }
592
593 #[must_use]
597 pub fn pair(workdir: PathBuf) -> (Self, mpsc::Receiver<Msg>) {
598 let (tx, rx) = mpsc::channel(MSG_CHANNEL_CAPACITY);
599 (Self::new(tx, workdir), rx)
600 }
601
602 #[must_use]
605 pub fn pair_with_bindings(
606 workdir: PathBuf,
607 config: Config,
608 tools: Arc<ToolRegistry>,
609 ) -> (Self, mpsc::Receiver<Msg>) {
610 let providers = Arc::new(ProviderFactory::new(config));
611 Self::pair_from(workdir, providers, tools)
612 }
613
614 pub fn pair_from(
619 workdir: PathBuf,
620 providers: Arc<ProviderFactory>,
621 tools: Arc<ToolRegistry>,
622 ) -> (Self, mpsc::Receiver<Msg>) {
623 let (tx, rx) = mpsc::channel(MSG_CHANNEL_CAPACITY);
624 (Self::new(tx, workdir).with_bindings(providers, tools), rx)
625 }
626
627 pub fn pair_from_with_task(
628 workdir: PathBuf,
629 providers: Arc<ProviderFactory>,
630 tools: Arc<ToolRegistry>,
631 task_id: Option<String>,
632 ) -> (Self, mpsc::Receiver<Msg>) {
633 let (runner, rx) = Self::pair_from(workdir, providers, tools);
634 (runner.with_task_id(task_id), rx)
635 }
636
637 pub fn new_child(
642 msg_tx: MsgSender,
643 workdir: PathBuf,
644 providers: Arc<ProviderFactory>,
645 tools: Arc<ToolRegistry>,
646 ) -> Self {
647 Self::new(msg_tx, workdir)
655 .with_bindings(providers, tools)
656 .without_terminal_title()
657 .without_global_mcp_shutdown()
658 }
659
660 fn scope_mut(&mut self, turn: TurnId) -> &mut TurnScope {
664 self.scopes
665 .entry(turn)
666 .or_insert_with(|| TurnScope::new(turn))
667 }
668
669 fn tombstone_turn(&mut self, turn: TurnId) {
673 if self.cancelled_turns.contains(&turn) {
674 return;
675 }
676 if self.cancelled_turns.len() >= CANCELLED_TOMBSTONE_CAP {
677 self.cancelled_turns.pop_front();
678 }
679 self.cancelled_turns.push_back(turn);
680 }
681
682 fn is_tombstoned(&self, turn: TurnId) -> bool {
685 self.cancelled_turns.contains(&turn)
686 }
687
688 fn dispatch_query(&mut self, query: Query) {
694 let tx = self.msg_tx.clone();
695 match query {
696 Query::LoadConversation { id } => self.query_load_conversation(id, tx),
697 Query::ListConversations => self.query_list_conversations(tx),
698 Query::ListAvailableModels => {
699 let providers = self.providers.clone();
700 self.detached.spawn(async move {
701 let choices = discover_available_models(providers).await;
702 let _ = tx
703 .send(Msg::QueryResult(QueryResult::AvailableModelsListed(
704 choices,
705 )))
706 .await;
707 });
708 },
709 Query::ListProjectFiles => {
710 let workdir = self.workdir.clone();
711 self.send_blocking_query(move || {
712 QueryResult::ProjectFilesListed(walk_project_files(&workdir))
713 });
714 },
715 Query::ListOutputStyles => self.dispatch_list_output_styles(),
716 Query::LoadOutputStyle { name, project } => {
717 self.dispatch_load_output_style(name, project);
718 },
719 Query::ListRuntimeTasks { limit } => self.send_blocking_query(move || {
720 QueryResult::RuntimeTasksListed(
721 crate::runtime_client::RuntimeClient::auto()
722 .list_tasks(limit)
723 .map(|read| read.value)
724 .unwrap_or_default(),
725 )
726 }),
727 Query::LoadRuntimeTask { id } => self.send_blocking_query(move || {
728 let (task, events) = crate::runtime_client::RuntimeClient::auto()
729 .task_detail(&id)
730 .map(|read| (Some(Box::new(read.value.task)), read.value.events))
731 .unwrap_or((None, Vec::new()));
732 QueryResult::RuntimeTaskLoaded { task, events }
733 }),
734 Query::ListRuntimeProcesses { limit } => self.send_blocking_query(move || {
735 QueryResult::RuntimeProcessesListed(
736 crate::runtime_client::RuntimeClient::auto()
737 .list_processes(limit)
738 .map(|read| read.value)
739 .unwrap_or_default(),
740 )
741 }),
742 Query::ListRuntimeApprovals => self.send_blocking_query(move || {
743 QueryResult::RuntimeApprovalsListed(
744 crate::runtime_client::RuntimeClient::auto()
745 .list_approvals()
746 .map(|read| read.value)
747 .unwrap_or_default(),
748 )
749 }),
750 Query::ListRuntimeCheckpoints { limit } => self.send_blocking_query(move || {
751 QueryResult::RuntimeCheckpointsListed(
752 crate::runtime_client::RuntimeClient::auto()
753 .list_checkpoints(limit)
754 .map(|read| read.value)
755 .unwrap_or_default(),
756 )
757 }),
758 Query::ListForkCheckpoints {
759 session_id,
760 message_index,
761 } => self.send_blocking_query(move || {
762 QueryResult::ForkCheckpointsFound(
763 mermaid_runtime::with_shared_store(|store| {
764 store
765 .checkpoints()
766 .list_for_session(&session_id, message_index as i64)
767 })
768 .unwrap_or_default(),
769 )
770 }),
771 Query::ListRuntimePlugins => self.send_blocking_query(move || {
772 QueryResult::RuntimePluginsListed(
773 crate::runtime_client::RuntimeClient::auto()
774 .list_plugins()
775 .map(|read| read.value)
776 .unwrap_or_default(),
777 )
778 }),
779 }
780 }
781
782 fn dispatch_list_output_styles(&mut self) {
785 let workdir = self.workdir.clone();
786 self.send_blocking_query(move || {
787 QueryResult::OutputStylesListed(crate::app::output_styles::list_styles(&workdir))
788 });
789 }
790
791 fn dispatch_load_output_style(&mut self, name: String, project: bool) {
795 let workdir = self.workdir.clone();
796 self.send_blocking_query(move || {
797 let loaded = crate::app::output_styles::load_style_body(&workdir, &name);
798 let (found, body, keep, custom, source) = match loaded {
799 Some((body, keep, custom, source)) => {
800 (true, body, keep, custom, source.to_string())
801 },
802 None => (false, String::new(), true, false, String::new()),
803 };
804 QueryResult::OutputStyleLoaded {
805 name,
806 project,
807 found,
808 body,
809 keep_coding_instructions: keep,
810 custom,
811 source,
812 }
813 });
814 }
815
816 fn query_load_conversation(&mut self, id: String, tx: MsgSender) {
820 let workdir = self.workdir.clone();
821 self.detached.spawn(async move {
822 match crate::session::ConversationManager::new(&workdir) {
823 Ok(mgr) => match mgr.load_conversation(&id) {
824 Ok(history) => {
825 let _ = tx
826 .send(Msg::QueryResult(QueryResult::ConversationLoaded(Box::new(
827 history,
828 ))))
829 .await;
830 },
831 Err(e) => {
832 tracing::warn!(id = %id, error = %e, "LoadConversation failed");
833 },
834 },
835 Err(e) => {
836 tracing::warn!(error = %e, "ConversationManager init failed");
837 },
838 }
839 });
840 }
841
842 fn query_list_conversations(&mut self, tx: MsgSender) {
845 let workdir = self.workdir.clone();
846 self.detached.spawn(async move {
847 let summaries = match crate::session::ConversationManager::new(&workdir) {
848 Ok(mgr) => mgr
849 .list_conversation_metas()
850 .unwrap_or_default()
851 .into_iter()
852 .map(|m| mermaid_domain::ConversationSummary {
853 id: m.id,
854 title: m.title,
855 message_count: m.message_count,
856 updated_at: m.updated_at.to_rfc3339(),
857 })
858 .collect(),
859 Err(_) => Vec::new(),
860 };
861 let _ = tx
862 .send(Msg::QueryResult(QueryResult::ConversationsListed(
863 summaries,
864 )))
865 .await;
866 });
867 }
868
869 fn send_blocking_query(&mut self, run: impl FnOnce() -> QueryResult + Send + 'static) {
874 let tx = self.msg_tx.clone();
875 self.detached.spawn_blocking(move || {
876 let _ = tx.blocking_send(Msg::QueryResult(run()));
877 });
878 }
879
880 fn drop_scope(&mut self, turn: TurnId) {
890 self.tombstone_turn(turn);
895 if let Some(mut scope) = self.scopes.remove(&turn) {
896 scope.cancel();
897 let tx = self.msg_tx.clone();
898 self.detached.spawn(async move {
899 if tokio::time::timeout(CANCEL_DRAIN_TIMEOUT, scope.drain())
900 .await
901 .is_err()
902 {
903 tracing::warn!(
904 turn = %turn,
905 timeout_ms = CANCEL_DRAIN_TIMEOUT.as_millis(),
906 "cancel drain timed out; aborting remaining scoped tasks"
907 );
908 }
909 let _ = tx.send(Msg::TurnCancelled(turn)).await;
910 });
911 } else {
912 let tx = self.msg_tx.clone();
919 self.detached.spawn(async move {
920 let _ = tx.send(Msg::TurnCancelled(turn)).await;
921 });
922 }
923 }
924
925 #[must_use]
928 pub fn scope_count(&self) -> usize {
929 self.scopes.len()
930 }
931
932 fn reap_empty_scopes(&mut self) {
941 self.reap_detached();
942 self.scopes.retain(|_, scope| {
943 scope.drain_completed();
944 !scope.is_empty()
945 });
946 }
947
948 fn reap_detached(&mut self) {
953 while let Some(result) = self.detached.try_join_next() {
954 if let Err(e) = result
955 && !e.is_cancelled()
956 {
957 tracing::warn!(error = %e, "effect: detached task panicked");
958 }
959 }
960 }
961
962 #[expect(
966 clippy::too_many_lines,
967 reason = "the effect router: one arm per Cmd variant, each spawning or calling the \
968 handler that owns that effect; the arms are short and the routing table is the point, so \
969 the length is the Cmd count and drops only as commands are retired"
970 )]
971 pub fn dispatch(&mut self, cmd: Cmd) {
972 self.reap_empty_scopes();
975 tracing::trace!(cmd = %cmd.summary(), "effect: dispatch");
976
977 if let Some(turn) = cmd.scope_turn()
985 && self.is_tombstoned(turn)
986 {
987 tracing::debug!(
988 cmd = %cmd.summary(),
989 turn = %turn,
990 "effect: dropping turn-scoped cmd for an already-cancelled turn"
991 );
992 return;
993 }
994
995 match cmd {
996 Cmd::CallModel { turn, mut request } => {
997 let tx = self.msg_tx.clone();
998 let providers = self.providers.clone();
999 if let Some(tools) = &self.tools
1008 && request.output_schema.is_none()
1009 {
1010 let mut enriched =
1011 filter_suppressed(tools.describe_all(), &request.suppressed_builtin_tools);
1012 let builtin_tokens = mermaid_domain::estimate_tool_schema_tokens(&enriched);
1017 if let Err(e) = tx.try_send(Msg::BuiltinToolSchemaTokens(builtin_tokens)) {
1023 tracing::debug!(
1024 error = %e,
1025 "effect: dropped builtin tool-schema token estimate (channel full); \
1026 /context preview may be briefly stale"
1027 );
1028 }
1029 enriched.append(&mut request.tools);
1030 request.tools = enriched;
1031 }
1032 self.detached.spawn(fire_plugin_hooks(
1035 "prompt_submit",
1036 serde_json::json!({
1037 "turn_id": turn.0,
1038 "model_id": request.model_id.clone(),
1039 "message_count": request.messages.len(),
1040 "tool_count": request.tools.len(),
1041 }),
1042 ));
1043 let task_usage = self.tasks.clone();
1046 let scope = self.scope_mut(turn);
1047 let token = scope.token();
1048 scope.spawn(async move {
1049 use futures::FutureExt;
1050 let fallback_tx = tx.clone();
1051 if std::panic::AssertUnwindSafe(dispatch_call_model(
1052 tx, providers, turn, request, token, task_usage,
1053 ))
1054 .catch_unwind()
1055 .await
1056 .is_err()
1057 {
1058 tracing::error!(turn = %turn, "dispatch_call_model panicked");
1063 let _ = fallback_tx
1064 .send(Msg::UpstreamError {
1065 turn,
1066 error: mermaid_model::models::UserFacingError {
1067 summary: "Internal error".to_string(),
1068 message: "The model dispatch task panicked unexpectedly."
1069 .to_string(),
1070 suggestion: "This is a bug. Please retry; if it persists, \
1071 check the logs."
1072 .to_string(),
1073 category: mermaid_model::models::ErrorCategory::Internal,
1074 recoverable: true,
1075 },
1076 })
1077 .await;
1078 }
1079 });
1080 },
1081 Cmd::CompactConversation { turn, mut request } => {
1082 let tx = self.msg_tx.clone();
1083 let providers = self.providers.clone();
1084 if let Some(tools) = &self.tools {
1085 let mut enriched = tools.describe_all();
1086 enriched.append(&mut request.chat.tools);
1087 request.chat.tools = enriched;
1088 }
1089 let trigger = request.trigger;
1092 let scope = self.scope_mut(turn);
1093 let token = scope.token();
1094 scope.spawn(async move {
1095 use futures::FutureExt;
1096 let fallback_tx = tx.clone();
1097 if std::panic::AssertUnwindSafe(dispatch_compact_conversation(
1098 tx, providers, turn, request, token,
1099 ))
1100 .catch_unwind()
1101 .await
1102 .is_err()
1103 {
1104 tracing::error!(turn = %turn, "dispatch_compact_conversation panicked");
1110 let _ = fallback_tx
1111 .send(Msg::CompactionFailed {
1112 turn,
1113 trigger,
1114 message: "the compaction task panicked unexpectedly".to_string(),
1115 kind: mermaid_domain::StatusKind::Error,
1116 })
1117 .await;
1118 }
1119 });
1120 },
1121 Cmd::ExecuteTool {
1122 turn,
1123 call_id,
1124 source,
1125 dispatch,
1126 } => {
1127 let tx = self.msg_tx.clone();
1128 let tools = self.tools.clone();
1129 let workdir = self.workdir.clone();
1130 let config = self
1135 .providers
1136 .as_ref()
1137 .map(|p| Arc::new(p.config().clone()))
1138 .unwrap_or_else(|| Arc::new(mermaid_domain::Config::default()));
1139 let classifier: Option<Arc<dyn crate::providers::AutoClassifier>> =
1147 if dispatch.safety_mode == mermaid_runtime::SafetyMode::Auto
1148 || dispatch.plan_file.is_some()
1149 {
1150 self.providers.as_ref().map(|p| {
1151 let model = config
1152 .safety
1153 .auto_classifier_model
1154 .clone()
1155 .unwrap_or_else(|| dispatch.model_id.clone());
1156 Arc::new(crate::providers::ModelAutoClassifier::new(p.clone(), model))
1157 as Arc<dyn crate::providers::AutoClassifier>
1158 })
1159 } else {
1160 None
1161 };
1162 let services = crate::providers::ctx::ToolServices {
1163 workdir,
1164 config,
1165 task_id: self.task_id.clone(),
1166 notify: Some(self.msg_tx.clone()),
1170 classifier,
1171 approval: self.approval.clone(),
1172 questions: self.questions.clone(),
1173 tasks: Some(self.tasks.clone()),
1174 };
1175 let scope = self.scope_mut(turn);
1176 let signals = crate::providers::ctx::TurnSignals {
1177 token: scope.token(),
1178 background: scope.background_token(),
1179 web_bytes: scope.web_bytes(),
1180 };
1181 scope.spawn(async move {
1182 use futures::FutureExt;
1183 let fallback_tx = tx.clone();
1184 if std::panic::AssertUnwindSafe(dispatch_execute_tool(
1185 tx, tools, turn, call_id, source, signals, dispatch, services,
1186 ))
1187 .catch_unwind()
1188 .await
1189 .is_err()
1190 {
1191 tracing::error!(
1196 turn = %turn,
1197 call_id = call_id.0,
1198 "dispatch_execute_tool panicked"
1199 );
1200 let _ = fallback_tx
1201 .send(Msg::ToolFinished {
1202 turn,
1203 call_id,
1204 outcome: mermaid_domain::ToolOutcome::error(
1205 "internal error: the tool execution task panicked".to_string(),
1206 0.0,
1207 ),
1208 })
1209 .await;
1210 }
1211 });
1212 },
1213 Cmd::ResolveApproval { call_id, decision } => {
1214 if let Some(broker) = &self.approval {
1217 broker.resolve(call_id, decision.into());
1218 }
1219 },
1220 Cmd::ResolveQuestion {
1221 call_id,
1222 resolution,
1223 } => {
1224 if let Some(broker) = &self.questions {
1227 broker.resolve(call_id, resolution);
1228 }
1229 },
1230 Cmd::SyncTaskStore(store) => {
1231 self.tasks.seed(store);
1235 },
1236 Cmd::EnsureScratchpad { session_id } => {
1237 let tx = self.msg_tx.clone();
1238 let workdir = self.workdir.clone();
1239 self.detached.spawn(async move {
1240 match crate::session::scratchpad::ensure(&workdir, &session_id) {
1241 Ok(path) => {
1242 let _ = tx.send(Msg::ScratchpadReady { session_id, path }).await;
1243 },
1244 Err(err) => {
1245 tracing::warn!(error = %err, "failed to create session scratchpad");
1248 },
1249 }
1250 if let Err(err) = crate::session::scratchpad::sweep_stale(
1253 crate::session::scratchpad::RETENTION_DAYS,
1254 ) {
1255 tracing::warn!(error = %err, "scratchpad sweep failed");
1256 }
1257 });
1258 },
1259 Cmd::ListScratchpad { path } => {
1260 let tx = self.msg_tx.clone();
1263 self.detached.spawn(async move {
1264 let text = tokio::task::spawn_blocking(move || {
1265 crate::session::scratchpad::list_text(&path)
1266 })
1267 .await
1268 .unwrap_or_else(|e| format!("Couldn't list the scratchpad: {e}"));
1269 let _ = tx.send(Msg::RuntimeText(text)).await;
1270 });
1271 },
1272 Cmd::UserTaskEdit(edit) => {
1273 let broker = self.tasks.clone();
1278 let tx = self.msg_tx.clone();
1279 self.detached.spawn(async move {
1280 let (line, _snapshot) = broker.user_edit(edit).await;
1281 let _ = tx
1286 .send(Msg::TaskNotice {
1287 text: format!(
1288 "The user edited the task checklist: {line}. Acknowledge and \
1289 incorporate this into your plan."
1290 ),
1291 })
1292 .await;
1293 let _ = tx.send(Msg::TransientStatus { text: line }).await;
1294 });
1295 },
1296 Cmd::NotifyTaskCompleted {
1297 task,
1298 completed,
1299 total,
1300 } => {
1301 let payload = serde_json::json!({
1309 "task_id": task.id,
1310 "subject": task.subject,
1311 "description": task.description,
1312 "evidence": task.evidence,
1313 "completed": completed,
1314 "total": total,
1315 });
1316 let broker = self.tasks.clone();
1317 let tx = self.msg_tx.clone();
1318 self.detached.spawn(async move {
1319 let gate = run_plugin_hooks_gated("task_completed", payload).await;
1320 let Some((plugin, reason)) = gate.deny else {
1321 return;
1322 };
1323 let reason = mermaid_model::utils::redact_secrets(&reason);
1324 let _ = broker
1325 .update(vec![mermaid_domain::ChecklistEdit {
1326 id: task.id,
1327 status: Some(mermaid_domain::ChecklistStatus::InProgress),
1328 ..mermaid_domain::ChecklistEdit::default()
1329 }])
1330 .await;
1331 let _ = tx
1332 .send(Msg::TaskNotice {
1333 text: format!(
1334 "Completion of task #{} '{}' was vetoed by the {plugin} hook: \
1335 {reason}. The task is back in_progress; address the reason \
1336 before completing it again.",
1337 task.id, task.subject
1338 ),
1339 })
1340 .await;
1341 let _ = tx
1342 .send(Msg::TransientStatus {
1343 text: format!(
1344 "task #{} completion vetoed by {plugin}: {reason}",
1345 task.id
1346 ),
1347 })
1348 .await;
1349 });
1350 },
1351 Cmd::CancelScope(turn) => {
1352 self.drop_scope(turn);
1353 },
1354 Cmd::BackgroundScope(turn) => {
1355 self.scope_mut(turn).background();
1359 },
1360 Cmd::SaveConversation { snapshot, events } => {
1361 self.queue_persistence(PersistenceJob::Conversation {
1362 snapshot: Box::new(snapshot),
1363 events,
1364 });
1365 },
1366 Cmd::SaveCompaction {
1367 record,
1368 conversation,
1369 events,
1370 } => {
1371 self.queue_persistence(PersistenceJob::Compaction(Box::new(
1372 PendingCompactionSave {
1373 record,
1374 conversation,
1375 events,
1376 events_appended: false,
1377 task_id: self.task_id.clone(),
1378 },
1379 )));
1380 },
1381 Cmd::SaveProcess(process) => {
1382 let task_id = self.task_id.clone();
1383 self.detached.spawn(async move {
1384 let status = process.status;
1385 let _ = mermaid_runtime::with_shared_store(|store| {
1386 store.processes().upsert(mermaid_runtime::NewProcess {
1387 id: Some(process.id),
1388 task_id,
1389 pid: process.pid,
1390 command: process.command,
1391 cwd: process.cwd,
1392 log_path: Some(process.log_path),
1393 detected_url: process.detected_url,
1394 status,
1395 health: None,
1396 })
1397 });
1398 });
1399 },
1400 Cmd::PersistPlanConfig(plan) => {
1401 self.detached.spawn(async move {
1402 if let Err(err) = crate::app::persist_plan_config(&plan) {
1403 tracing::warn!(error = %err, "failed to persist [plan] config");
1404 }
1405 });
1406 },
1407 Cmd::PersistLastModel(model) => {
1408 self.detached.spawn(async move {
1409 if let Err(err) = crate::app::persist_last_model(&model) {
1410 tracing::warn!(error = %err, "failed to persist last-used model");
1411 }
1412 });
1413 },
1414 Cmd::PersistReasoningFor { model_id, level } => {
1415 self.detached.spawn(async move {
1416 if let Err(err) = crate::app::persist_reasoning_for_model(&model_id, level) {
1417 tracing::warn!(error = %err, "failed to persist reasoning level for model");
1418 }
1419 });
1420 },
1421 Cmd::PersistOllamaNumCtxFor { model_id, num_ctx } => {
1422 self.detached.spawn(async move {
1423 if let Err(err) =
1424 crate::app::persist_ollama_num_ctx_for_model(&model_id, num_ctx)
1425 {
1426 tracing::warn!(error = %err, "failed to persist Ollama num_ctx for model");
1427 }
1428 });
1429 },
1430 Cmd::PersistOllamaOffload(enabled) => {
1431 self.detached.spawn(async move {
1432 if let Err(err) = crate::app::persist_ollama_allow_ram_offload(enabled) {
1433 tracing::warn!(error = %err, "failed to persist Ollama RAM-offload setting");
1434 }
1435 });
1436 },
1437 Cmd::PersistUiTheme(theme) => {
1438 self.detached.spawn(async move {
1439 if let Err(err) = crate::app::persist_ui_theme(theme) {
1440 tracing::warn!(error = %err, "failed to persist theme");
1441 }
1442 });
1443 },
1444 Cmd::PersistOutputStyle { style } => {
1445 self.detached.spawn(async move {
1446 if let Err(err) = crate::app::persist_output_style(&style) {
1447 tracing::warn!(error = %err, "failed to persist output style");
1448 }
1449 });
1450 },
1451 Cmd::PersistProjectOutputStyle { style } => {
1452 let tx = self.msg_tx.clone();
1453 let workdir = self.workdir.clone();
1454 self.detached.spawn(async move {
1455 if let Err(err) = crate::app::persist_project_output_style(&workdir, &style) {
1456 let _ = tx
1457 .send(Msg::TransientStatus {
1458 text: format!("Couldn't save the project output style: {err}"),
1459 })
1460 .await;
1461 }
1462 });
1463 },
1464 Cmd::ComposeInEditor { .. } => {
1465 tracing::warn!("compose_in_editor is unavailable outside the interactive TUI");
1469 },
1470 Cmd::ListMemory => {
1471 let tx = self.msg_tx.clone();
1472 let workdir = self.workdir.clone();
1473 self.detached.spawn(async move {
1474 let cfg = crate::app::load_project_scoped_config(&workdir).memory;
1475 let text = match crate::app::memory::load(&workdir, &cfg) {
1476 Some(mem) => mem.index,
1477 None => "No memories saved yet. Durable facts (yours or mine) show up here — use `/remember <fact>` or just ask me to remember something.".to_string(),
1478 };
1479 let _ = tx.send(Msg::RuntimeText(text)).await;
1480 });
1481 },
1482 Cmd::RememberMemory { text } => {
1483 let tx = self.msg_tx.clone();
1484 let workdir = self.workdir.clone();
1485 self.detached.spawn(async move {
1486 let cfg = crate::app::load_project_scoped_config(&workdir).memory;
1487 let name = memory_title_from_text(&text);
1488 let status = match crate::app::memory::write_memory(
1489 &workdir,
1490 mermaid_domain::MemoryScope::ProjectPrivate,
1491 &name,
1492 &text,
1493 &[],
1494 &text,
1495 ) {
1496 Ok(_) => format!("Remembered: {name}"),
1497 Err(e) => format!("Couldn't save memory: {e}"),
1498 };
1499 let (loaded, _) = crate::app::memory::refresh(None, &workdir, &cfg);
1500 let _ = tx.send(Msg::MemoryChanged(loaded)).await;
1501 let _ = tx.send(Msg::TransientStatus { text: status }).await;
1502 });
1503 },
1504 Cmd::ForgetMemory { id } => {
1505 let tx = self.msg_tx.clone();
1506 let workdir = self.workdir.clone();
1507 self.detached.spawn(async move {
1508 let cfg = crate::app::load_project_scoped_config(&workdir).memory;
1509 let status = match crate::app::memory::delete_memory(&workdir, &id) {
1510 Ok(Some(_)) => format!("Forgot: {id}"),
1511 Ok(None) => format!("No memory named '{id}'"),
1512 Err(e) => format!("Couldn't forget memory: {e}"),
1513 };
1514 let (loaded, _) = crate::app::memory::refresh(None, &workdir, &cfg);
1515 let _ = tx.send(Msg::MemoryChanged(loaded)).await;
1516 let _ = tx.send(Msg::TransientStatus { text: status }).await;
1517 });
1518 },
1519 Cmd::ConsolidateMemory { model_id } => {
1520 let tx = self.msg_tx.clone();
1521 let workdir = self.workdir.clone();
1522 let providers = self.providers.clone();
1523 self.detached.spawn(async move {
1524 consolidate_memory(tx, providers, workdir, model_id).await;
1525 });
1526 },
1527 Cmd::Query(query) => self.dispatch_query(query),
1528 Cmd::ShowRuntimeProcessLogs { id } => {
1529 let tx = self.msg_tx.clone();
1530 self.detached.spawn_blocking(move || {
1531 let text = crate::runtime_client::RuntimeClient::auto()
1532 .process_log(&id, None)
1533 .map(|log| format!("Process log {}\n\n{}", id, log.content))
1534 .unwrap_or_else(|err| format!("Process log error: {err}"));
1535 let _ = tx.blocking_send(Msg::RuntimeText(text));
1536 });
1537 },
1538 Cmd::StopRuntimeProcess { id } => {
1539 let tx = self.msg_tx.clone();
1540 self.detached.spawn_blocking(move || {
1541 let msg = match crate::runtime_client::RuntimeClient::auto().stop_process(&id) {
1542 Ok(response) => Msg::TransientStatus {
1543 text: format!("Stopped process {} (pid {})", id, response.item.pid),
1544 },
1545 Err(err) => Msg::TransientStatus {
1546 text: format!("Process stop failed: {err}"),
1547 },
1548 };
1549 let _ = tx.blocking_send(msg);
1550 });
1551 },
1552 Cmd::KillBackgroundAgent { agent_id } => {
1553 let spawner = self.tools.as_ref().and_then(|t| t.subagent_spawner());
1557 if let Some(spawner) = spawner {
1558 match agent_id {
1559 Some(id) => {
1560 if let crate::providers::tool::subagent::KillResult::Evicted(
1565 workspace,
1566 ) = spawner.kill_detached(&id)
1567 {
1568 self.detached.spawn(async move {
1572 workspace.discard().await;
1573 });
1574 }
1575 },
1576 None => {
1577 spawner.kill_all_detached();
1578 },
1579 }
1580 }
1581 },
1582 Cmd::RestartRuntimeProcess { id } => {
1583 let tx = self.msg_tx.clone();
1584 self.detached.spawn_blocking(move || {
1585 let msg = match crate::runtime_client::RuntimeClient::auto()
1586 .restart_process(&id)
1587 {
1588 Ok(response) => Msg::TransientStatus {
1589 text: format!("Restarted process {} (pid {})", id, response.item.pid),
1590 },
1591 Err(err) => Msg::TransientStatus {
1592 text: format!("Process restart failed: {err}"),
1593 },
1594 };
1595 let _ = tx.blocking_send(msg);
1596 });
1597 },
1598 Cmd::OpenRuntimeTarget { target } => {
1599 self.detached.spawn_blocking(move || {
1600 let resolved = crate::runtime_client::RuntimeService::open_default()
1601 .and_then(|service| service.resolve_open_target(&target))
1602 .unwrap_or(target);
1603 if let Err(err) = crate::runtime_client::validate_open_target(&resolved) {
1607 tracing::warn!(error = %err, "refusing to open runtime target");
1608 return;
1609 }
1610 mermaid_model::utils::open_file(resolved);
1611 });
1612 },
1613 Cmd::ShowRuntimePorts => {
1614 let tx = self.msg_tx.clone();
1615 self.detached.spawn_blocking(move || {
1616 let text = crate::runtime_client::RuntimeClient::auto()
1617 .ports()
1618 .map(|ports| format!("Listening TCP ports\n\n{}", ports.ports))
1619 .unwrap_or_else(|err| format!("Port inspection failed: {err}"));
1620 let _ = tx.blocking_send(Msg::RuntimeText(text));
1621 });
1622 },
1623 Cmd::DecideRuntimeApproval { id, decision } => {
1624 let tx = self.msg_tx.clone();
1625 self.detached.spawn_blocking(move || {
1626 let result = if decision == "approved" {
1627 crate::runtime_client::RuntimeClient::auto().approve(&id)
1628 } else {
1629 crate::runtime_client::RuntimeClient::auto().deny(&id)
1630 };
1631 let msg = match result {
1632 Ok(result) => Msg::TransientStatus {
1633 text: if result.replayed {
1634 format!("Approval {} {}: {}", id, decision, result.summary)
1635 } else {
1636 format!("Approval {id} {decision}")
1637 },
1638 },
1639 Err(err) => Msg::TransientStatus {
1640 text: format!("Approval update failed: {err}"),
1641 },
1642 };
1643 let _ = tx.blocking_send(msg);
1644 });
1645 },
1646 Cmd::UpdateRuntimeTaskStatus {
1647 id,
1648 status,
1649 final_report,
1650 } => {
1651 let tx = self.msg_tx.clone();
1652 self.detached.spawn_blocking(move || {
1653 let msg = match mermaid_runtime::with_shared_store(|store| {
1654 store
1655 .tasks()
1656 .update_status(&id, status, final_report.as_deref())
1657 }) {
1658 Ok(()) => Msg::TransientStatus {
1659 text: format!("Task {id} -> {status}"),
1660 },
1661 Err(err) => Msg::TransientStatus {
1662 text: format!("Task update failed: {err}"),
1663 },
1664 };
1665 let _ = tx.blocking_send(msg);
1666 });
1667 },
1668 Cmd::CreateRuntimeCheckpoint { paths } => {
1669 let tx = self.msg_tx.clone();
1670 let workdir = self.workdir.clone();
1671 self.detached.spawn_blocking(move || {
1672 let pending_action = Some(serde_json::json!({
1673 "source": "tui",
1674 "command": "checkpoint",
1675 }));
1676 let msg = match mermaid_runtime::create_checkpoint(
1677 &workdir,
1678 &paths,
1679 pending_action,
1680 ) {
1681 Ok(manifest) => Msg::TransientStatus {
1682 text: format!(
1683 "Checkpoint {} created for {} path(s)",
1684 manifest.id,
1685 manifest.files.len()
1686 ),
1687 },
1688 Err(err) => Msg::TransientStatus {
1689 text: format!("Checkpoint failed: {err}"),
1690 },
1691 };
1692 let _ = tx.blocking_send(msg);
1693 });
1694 },
1695 Cmd::RestoreRuntimeCheckpoint { id } => {
1696 let tx = self.msg_tx.clone();
1697 self.detached.spawn_blocking(move || {
1698 let msg = match crate::runtime_client::RuntimeClient::auto()
1699 .restore_checkpoint(&id)
1700 {
1701 Ok(result) => Msg::TransientStatus {
1702 text: format!(
1703 "Restored checkpoint {} ({} file(s)){}",
1704 result.checkpoint.id,
1705 result.checkpoint.files.len(),
1706 if result.checkpoint.pending_action.is_some() {
1707 "; pending action available in checkpoint manifest"
1708 } else {
1709 ""
1710 }
1711 ),
1712 },
1713 Err(err) => Msg::TransientStatus {
1714 text: format!("Restore failed: {err}"),
1715 },
1716 };
1717 let _ = tx.blocking_send(msg);
1718 });
1719 },
1720 Cmd::ShowRuntimeModelInfo { model } => {
1721 let tx = self.msg_tx.clone();
1722 self.detached.spawn_blocking(move || {
1723 let text = runtime_model_info_text(&model);
1724 let _ = tx.blocking_send(Msg::RuntimeText(text));
1725 });
1726 },
1727 Cmd::InitMcpServers(configs) => {
1728 let tx = self.msg_tx.clone();
1729 self.detached
1730 .spawn(async move { dispatch_init_mcp_servers(configs, tx).await });
1731 },
1732 Cmd::StopMcpServer { name } => {
1733 let tx = self.msg_tx.clone();
1734 self.detached.spawn(async move {
1735 if let Some(mgr) = crate::mcp::manager_ref::get() {
1738 mgr.stop_server(&name).await;
1739 }
1740 let _ = tx.send(Msg::McpServerStopped { name }).await;
1741 });
1742 },
1743 Cmd::PullOllamaModel { model } => {
1744 let tx = self.msg_tx.clone();
1745 self.detached.spawn(async move {
1746 dispatch_pull_ollama_model(tx, model).await;
1747 });
1748 },
1749 Cmd::OpenInSystem(path) => {
1750 self.detached.spawn(async move {
1751 let _ = tokio::task::spawn_blocking(move || {
1752 mermaid_model::utils::open_file(&path);
1753 })
1754 .await;
1755 });
1756 },
1757 Cmd::WriteImageToTemp {
1758 path,
1759 bytes,
1760 format: _,
1761 } => {
1762 self.detached.spawn(async move {
1763 if let Err(e) = tokio::fs::write(&path, &bytes).await {
1764 tracing::warn!(path = %path.display(), error = %e, "WriteImageToTemp failed");
1765 }
1766 });
1767 },
1768 Cmd::ReadClipboard => {
1769 let tx = self.msg_tx.clone();
1770 self.detached.spawn(async move {
1771 dispatch_read_clipboard(tx).await;
1772 });
1773 },
1774 Cmd::ProbeVision { model_id, warn } => {
1775 let tx = self.msg_tx.clone();
1776 let providers = self.providers.clone();
1777 self.detached.spawn(async move {
1778 dispatch_probe_vision(model_id, warn, providers, tx).await;
1779 });
1780 },
1781 Cmd::CopyToClipboard(text) => {
1782 let tx = self.msg_tx.clone();
1783 self.detached.spawn(async move {
1784 dispatch_copy_to_clipboard(text, tx).await;
1785 });
1786 },
1787 Cmd::Exit => {
1788 },
1793 Cmd::SetTerminalTitle(title) => {
1794 if !self.terminal_title_enabled {
1795 return;
1796 }
1797 self.detached.spawn_blocking(move || {
1803 use std::io::Write;
1804 let seq = format!("\x1b]2;{title}\x07");
1805 let mut stdout = std::io::stdout();
1806 let _ = stdout.write_all(seq.as_bytes());
1807 let _ = stdout.flush();
1808 });
1809 },
1810 Cmd::AlertUser => {
1811 if !self.terminal_title_enabled {
1812 return;
1813 }
1814 self.detached.spawn_blocking(|| {
1817 use std::io::Write;
1818 let mut stdout = std::io::stdout();
1819 let _ = stdout.write_all(b"\x07");
1820 let _ = stdout.flush();
1821 });
1822 },
1823 }
1824 }
1825
1826 fn queue_persistence(&mut self, job: PersistenceJob) {
1827 let previous = self.persistence_tail.take();
1828 let state = Arc::clone(&self.persistence_state);
1829 let tx = self.msg_tx.clone();
1830 self.persistence_tail = Some(tokio::spawn(async move {
1831 if let Some(previous) = previous
1832 && let Err(error) = previous.await
1833 {
1834 tracing::warn!(error = %error, "previous persistence job panicked");
1835 }
1836
1837 let result = tokio::task::spawn_blocking(move || {
1838 state
1839 .lock()
1840 .unwrap_or_else(|error| error.into_inner())
1841 .process(job)
1842 })
1843 .await;
1844
1845 match result {
1846 Ok((events, outcome)) => {
1847 if outcome.is_ok() || !events.is_empty() {
1851 let _ = tx.send(Msg::SessionSaved).await;
1852 }
1853 for event in events {
1854 fire_compaction_hook(&event).await;
1855 }
1856 if let Err(error) = outcome {
1857 tracing::warn!(
1858 error = %error,
1859 "persistence job failed; compaction barriers remain queued"
1860 );
1861 }
1862 },
1863 Err(error) => tracing::warn!(error = %error, "persistence job panicked"),
1864 }
1865 }));
1866 }
1867
1868 pub async fn shutdown(mut self) {
1872 for (id, scope) in self.scopes.iter() {
1873 tracing::debug!(turn = %id, "shutdown: cancelling scope");
1874 scope.cancel();
1875 }
1876
1877 if let Some(handle) = self.config_watch.take() {
1880 handle.abort();
1881 }
1882
1883 let shutdown_deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(5);
1885
1886 let owns_global_mcp = self.owns_global_mcp;
1887 let persistence_tail = self.persistence_tail.take();
1888 let persistence_state = Arc::clone(&self.persistence_state);
1889 let drain = async {
1890 if let Some(tail) = persistence_tail
1891 && let Err(error) = tail.await
1892 {
1893 tracing::warn!(error = %error, "shutdown: persistence chain panicked");
1894 }
1895 match tokio::task::spawn_blocking(move || {
1896 let mut state = persistence_state
1897 .lock()
1898 .unwrap_or_else(|error| error.into_inner());
1899 let drained = state.retry_all_blocked();
1900 if let Err(error) = state.flush_checkpoints() {
1904 tracing::warn!(
1905 error = %error,
1906 "shutdown: could not flush a session checkpoint; the log still has everything, so the next resume just folds further"
1907 );
1908 }
1909 drop(state);
1910 drained
1911 })
1912 .await
1913 {
1914 Ok((events, outcome)) => {
1915 for event in events {
1918 fire_compaction_hook(&event).await;
1919 }
1920 if let Err(error) = outcome {
1921 tracing::warn!(
1922 error = %error,
1923 "shutdown: compaction persistence barrier retry failed"
1924 );
1925 }
1926 },
1927 Err(error) => tracing::warn!(
1928 error = %error,
1929 "shutdown: compaction persistence barrier panicked"
1930 ),
1931 }
1932 if owns_global_mcp {
1936 let _ = tokio::time::timeout(
1942 std::time::Duration::from_secs(2),
1943 crate::mcp::manager_ref::wait_ready(),
1944 )
1945 .await;
1946 if let Some(mgr) = crate::mcp::manager_ref::get() {
1953 mgr.shutdown().await;
1954 }
1955 crate::searxng::shutdown().await;
1959 }
1960 for (id, mut scope) in self.scopes.drain() {
1966 if tokio::time::timeout(CANCEL_DRAIN_TIMEOUT, scope.drain())
1967 .await
1968 .is_err()
1969 {
1970 tracing::warn!(
1971 turn = %id,
1972 timeout_ms = CANCEL_DRAIN_TIMEOUT.as_millis(),
1973 "shutdown: scope drain timed out; aborting its remaining tasks"
1974 );
1975 }
1976 }
1977 while let Some(result) = self.detached.join_next().await {
1978 if let Err(e) = result
1979 && !e.is_cancelled()
1980 {
1981 tracing::warn!(error = %e, "shutdown: detached task panic");
1982 }
1983 }
1984 };
1985
1986 let _ = tokio::time::timeout_at(shutdown_deadline, drain).await;
1987 }
1988}
1989
1990fn note_stream_usage(
1998 tasks: &crate::providers::TaskBroker,
1999 usage: &Option<mermaid_model::models::TokenUsage>,
2000) {
2001 if let Some(usage) = usage {
2002 tasks.add_tokens(usage.completion_tokens as u64);
2003 }
2004}
2005
2006fn filter_suppressed(
2010 tools: Vec<mermaid_domain::ToolDefinition>,
2011 suppressed: &[&'static str],
2012) -> Vec<mermaid_domain::ToolDefinition> {
2013 if suppressed.is_empty() {
2014 return tools;
2015 }
2016 tools
2017 .into_iter()
2018 .filter(|t| !suppressed.contains(&t.name.as_str()))
2019 .collect()
2020}
2021
2022mod compaction;
2023mod memory;
2024mod model_call;
2025mod tool_call;
2026
2027use compaction::*;
2028use memory::*;
2029use model_call::*;
2030use tool_call::*;
2031
2032#[cfg(test)]
2033mod tests {
2034 #[test]
2038 fn compaction_row_maps_every_field_it_claims_to() {
2039 use mermaid_domain::{CompactionEvent, CompactionReviewStatus, CompactionTrigger};
2040 let record = CompactionEvent {
2041 id: "cmp-1".to_string(),
2042 trigger: CompactionTrigger::Manual,
2043 created_at: chrono::Local::now(),
2044 before_tokens: 9_000,
2045 after_tokens: 1_200,
2046 archived_message_count: 40,
2047 preserved_message_count: 6,
2048 preserved_turn_count: 3,
2049 summary_tokens: 450,
2050 duration_secs: 1.5,
2051 review_status: CompactionReviewStatus::Reviewed,
2052 review_error: None,
2053 focus: None,
2054 archive_path: None,
2055 };
2056 let row = compaction_row(
2057 &record,
2058 std::path::Path::new("/tmp/archive.json"),
2059 Some("task-7".to_string()),
2060 "sess-3".to_string(),
2061 );
2062 assert_eq!(row.id.as_deref(), Some("cmp-1"));
2063 assert_eq!(row.task_id.as_deref(), Some("task-7"));
2064 assert_eq!(row.session_id.as_deref(), Some("sess-3"));
2065 assert_eq!(row.source_token_estimate, Some(9_000));
2066 assert_eq!(row.summary_token_count, Some(450));
2067 assert_eq!(row.preserved_turns, Some(3));
2068 assert!(row.archive_path.is_some_and(|p| p.contains("archive.json")));
2069 assert_eq!(
2070 row.verification_status.as_deref(),
2071 Some(CompactionReviewStatus::Reviewed.as_str())
2072 );
2073 }
2074
2075 use super::*;
2076 use mermaid_domain::ToolCallId;
2077 use std::time::Duration;
2078
2079 fn runner() -> (EffectRunner, mpsc::Receiver<Msg>) {
2080 EffectRunner::pair(PathBuf::from("/tmp"))
2081 }
2082
2083 #[test]
2086 fn filter_suppressed_drops_only_the_named_tools() {
2087 let def = |name: &str| mermaid_domain::ToolDefinition {
2088 name: name.to_string(),
2089 description: String::new(),
2090 input_schema: serde_json::json!({}),
2091 };
2092 let tools = vec![def("task_create"), def("task_list"), def("task_update")];
2093 let kept = filter_suppressed(tools.clone(), &["task_create", "task_update"]);
2094 assert_eq!(
2095 kept.iter().map(|t| t.name.as_str()).collect::<Vec<_>>(),
2096 vec!["task_list"]
2097 );
2098 let kept = filter_suppressed(tools, &[]);
2099 assert_eq!(kept.len(), 3, "empty suppression list is a no-op");
2100 }
2101
2102 #[test]
2103 fn runtime_tool_payloads_are_redacted_before_serialization() {
2104 let payload = serde_json::json!({
2105 "url": "https://user:hunter2@example.test/page?X-Amz-Signature=opaque-signature#private",
2106 "authorization": "opaque-secret-value",
2107 "model_content": "Fetched page says OPENAI_API_KEY=sk-abcdefghijklmnop1234 and Authorization: Bearer abcdef123456ghijkl",
2108 });
2109 let serialized = redacted_json_string(&payload).expect("serialize redacted payload");
2110 assert!(
2111 !serialized.contains("hunter2"),
2112 "URL password leaked: {serialized}"
2113 );
2114 assert!(
2115 !serialized.contains("opaque-signature"),
2116 "signed URL leaked: {serialized}"
2117 );
2118 assert!(
2119 !serialized.contains("private"),
2120 "URL fragment leaked: {serialized}"
2121 );
2122 assert!(
2123 !serialized.contains("opaque-secret-value"),
2124 "credential-named field leaked: {serialized}"
2125 );
2126 assert!(
2127 !serialized.contains("abcdef123456ghijkl"),
2128 "bearer token leaked: {serialized}"
2129 );
2130 assert!(
2131 !serialized.contains("sk-abcdefghijklmnop1234"),
2132 "secret-shaped fetched content leaked: {serialized}"
2133 );
2134 assert!(serialized.contains("[REDACTED]"));
2135 }
2136
2137 #[test]
2138 fn project_walk_respects_gitignore_sorts_and_marks_dirs() {
2139 let root = std::env::temp_dir().join(format!(
2140 "mermaid-walk-{}-{:?}",
2141 std::process::id(),
2142 std::thread::current().id()
2143 ));
2144 let _ = std::fs::remove_dir_all(&root);
2145 std::fs::create_dir_all(root.join("src")).unwrap();
2146 std::fs::create_dir_all(root.join("target")).unwrap();
2147 std::fs::create_dir_all(root.join(".git")).unwrap();
2148 std::fs::write(root.join(".gitignore"), "target/\n").unwrap();
2149 std::fs::write(root.join("src/main.rs"), "fn main() {}").unwrap();
2150 std::fs::write(root.join("target/out.bin"), "ignored").unwrap();
2151 std::fs::write(root.join("README.md"), "readme").unwrap();
2152 std::fs::write(root.join(".hidden"), "hidden").unwrap();
2153
2154 let files = walk_project_files(&root);
2155 assert_eq!(
2156 files,
2157 vec![
2158 "README.md".to_string(),
2159 "src/".to_string(),
2160 "src/main.rs".to_string(),
2161 ],
2162 "sorted, dirs slash-marked, target/ ignored, dotfiles hidden"
2163 );
2164 let _ = std::fs::remove_dir_all(&root);
2165 }
2166
2167 #[test]
2168 fn new_child_suppresses_terminal_title() {
2169 let (tx, _rx) = mpsc::channel::<Msg>(MSG_CHANNEL_CAPACITY);
2173 let providers = Arc::new(ProviderFactory::new(mermaid_domain::Config::default()));
2174 let tools = Arc::new(ToolRegistry::new());
2175 let child = EffectRunner::new_child(tx, PathBuf::from("/tmp"), providers, tools);
2176 assert!(
2177 !child.terminal_title_enabled,
2178 "subagent child runner must suppress terminal-title escapes"
2179 );
2180 }
2181
2182 #[test]
2183 fn new_child_does_not_own_global_mcp_shutdown() {
2184 let (tx, _rx) = mpsc::channel::<Msg>(MSG_CHANNEL_CAPACITY);
2189 let providers = Arc::new(ProviderFactory::new(mermaid_domain::Config::default()));
2190 let tools = Arc::new(ToolRegistry::new());
2191 let child = EffectRunner::new_child(tx, PathBuf::from("/tmp"), providers, tools);
2192 assert!(
2193 !child.owns_global_mcp,
2194 "child runner must not reap the shared global MCP manager"
2195 );
2196 let (top, _rx2) = EffectRunner::pair(PathBuf::from("/tmp"));
2197 assert!(
2198 top.owns_global_mcp,
2199 "top-level runner still owns the global MCP reap"
2200 );
2201 }
2202
2203 #[test]
2204 fn parse_prune_plan_extracts_json_amid_prose() {
2205 let plan = parse_prune_plan(
2206 "Sure, here's the plan:\n```json\n{\"prune\": [\"a\", \"b\"], \"reason\": \"dupes\"}\n```\nDone.",
2207 )
2208 .expect("should parse");
2209 assert_eq!(plan.prune, vec!["a".to_string(), "b".to_string()]);
2210 assert_eq!(plan.reason, "dupes");
2211 }
2212
2213 #[test]
2214 fn parse_prune_plan_handles_empty_and_garbage() {
2215 let empty = parse_prune_plan("{\"prune\": [], \"reason\": \"all distinct\"}")
2216 .expect("empty plan parses");
2217 assert!(empty.prune.is_empty());
2218 assert!(parse_prune_plan("no json here").is_none());
2219 }
2220
2221 #[test]
2222 fn memory_title_from_text_is_short_and_nonempty() {
2223 assert_eq!(
2224 memory_title_from_text("prefer ripgrep over grep"),
2225 "prefer ripgrep over grep"
2226 );
2227 assert_eq!(memory_title_from_text(" "), "memory");
2228 let long = memory_title_from_text("one two three four five six seven eight nine ten");
2229 assert!(long.split_whitespace().count() <= 8);
2230 }
2231
2232 #[tokio::test]
2233 async fn dispatch_exit_is_noop_on_runner_state() {
2234 let (mut r, _rx) = runner();
2235 r.dispatch(Cmd::Exit);
2236 assert_eq!(r.scope_count(), 0);
2237 }
2238
2239 #[tokio::test]
2240 async fn dispatch_save_emits_session_saved() {
2241 let (mut r, mut rx) = runner();
2242 r.dispatch(Cmd::SaveConversation {
2243 snapshot: mermaid_domain::ConversationHistory::new(
2244 "/p".to_string(),
2245 "m".to_string(),
2246 chrono::Local::now(),
2247 ),
2248 events: Vec::new(),
2249 });
2250 let msg = tokio::time::timeout(Duration::from_millis(200), rx.recv())
2251 .await
2252 .expect("sender emits")
2253 .expect("channel alive");
2254 assert!(matches!(msg, Msg::SessionSaved));
2255 }
2256
2257 #[cfg(unix)]
2258 #[tokio::test]
2259 async fn init_mcp_servers_emits_incremental_errored_msgs() {
2260 let (tx, mut rx) = tokio::sync::mpsc::channel(8);
2264 let mut configs = std::collections::HashMap::new();
2265 for name in ["one", "two"] {
2266 configs.insert(
2267 name.to_string(),
2268 mermaid_domain::McpServerConfig {
2269 command: "/nonexistent/mermaid-test-mcp-binary".to_string(),
2270 ..Default::default()
2271 },
2272 );
2273 }
2274 dispatch_init_mcp_servers(configs, tx).await;
2275 let mut errored = Vec::new();
2276 while let Ok(msg) = rx.try_recv() {
2277 match msg {
2278 Msg::McpServerErrored { name, .. } => errored.push(name),
2279 other => panic!("unexpected msg: {other:?}"),
2280 }
2281 }
2282 errored.sort();
2283 assert_eq!(errored, vec!["one".to_string(), "two".to_string()]);
2284 assert!(crate::mcp::manager_ref::is_ready());
2285 }
2286
2287 #[tokio::test]
2288 async fn cancel_scope_emits_turn_cancelled_after_bounded_timeout() {
2289 let (mut r, mut rx) = runner();
2290 let turn = TurnId(77);
2291 {
2292 let scope = r.scope_mut(turn);
2293 scope.spawn(async {
2294 std::future::pending::<()>().await;
2295 });
2296 }
2297 assert_eq!(r.scope_count(), 1);
2298
2299 let start = std::time::Instant::now();
2300 r.dispatch(Cmd::CancelScope(turn));
2301 assert_eq!(r.scope_count(), 0);
2302 let msg = tokio::time::timeout(Duration::from_millis(500), rx.recv())
2303 .await
2304 .expect("bounded cancel should emit terminal message")
2305 .expect("channel alive");
2306 assert!(matches!(msg, Msg::TurnCancelled(t) if t == turn));
2307 assert!(
2308 start.elapsed() < Duration::from_millis(500),
2309 "cancel terminal message took {:?}",
2310 start.elapsed()
2311 );
2312 }
2313
2314 #[tokio::test]
2315 async fn cancel_scope_emits_turn_cancelled_even_after_reaping() {
2316 let (mut r, mut rx) = runner();
2322 let turn = TurnId(88);
2323 {
2324 let scope = r.scope_mut(turn);
2325 scope.spawn(async {}); }
2327 assert_eq!(r.scope_count(), 1);
2328
2329 tokio::time::sleep(Duration::from_millis(20)).await;
2331 r.dispatch(Cmd::Exit);
2332 assert_eq!(r.scope_count(), 0, "completed scope should be reaped");
2333
2334 r.dispatch(Cmd::CancelScope(turn));
2337 let msg = tokio::time::timeout(Duration::from_millis(500), rx.recv())
2338 .await
2339 .expect("cancel on a reaped scope must still emit a terminal message")
2340 .expect("channel alive");
2341 assert!(matches!(msg, Msg::TurnCancelled(t) if t == turn));
2342 }
2343
2344 #[tokio::test]
2345 async fn dispatch_call_model_creates_scope() {
2346 let (mut r, _rx) = runner();
2347 let turn = TurnId(7);
2348 let request = mermaid_domain::ChatRequest {
2349 model_id: "test/m".to_string(),
2350 messages: vec![],
2351 system_prompt: String::new(),
2352 instructions: None,
2353 reasoning: mermaid_model::models::ReasoningLevel::Medium,
2354 temperature: 0.7,
2355 max_tokens: 4096,
2356 tools: vec![],
2357
2358 ollama_num_ctx: None,
2359 ollama_allow_ram_offload: None,
2360 resolved_context_window: None,
2361 resolved_max_output: None,
2362 output_schema: None,
2363 suppress_auto_compact: false,
2364 suppressed_builtin_tools: Vec::new(),
2365 };
2366 r.dispatch(Cmd::CallModel { turn, request });
2367 assert_eq!(r.scope_count(), 1);
2368 }
2369
2370 #[tokio::test]
2374 async fn empty_scopes_are_reaped_on_next_dispatch() {
2375 let (mut r, mut rx) = runner();
2376 let turn = TurnId(42);
2377 let request = mermaid_domain::ChatRequest {
2378 model_id: "test/m".to_string(),
2379 messages: vec![],
2380 system_prompt: String::new(),
2381 instructions: None,
2382 reasoning: mermaid_model::models::ReasoningLevel::Medium,
2383 temperature: 0.7,
2384 max_tokens: 4096,
2385 tools: vec![],
2386
2387 ollama_num_ctx: None,
2388 ollama_allow_ram_offload: None,
2389 resolved_context_window: None,
2390 resolved_max_output: None,
2391 output_schema: None,
2392 suppress_auto_compact: false,
2393 suppressed_builtin_tools: Vec::new(),
2394 };
2395 r.dispatch(Cmd::CallModel { turn, request });
2396 assert_eq!(r.scope_count(), 1);
2397
2398 let msg = tokio::time::timeout(Duration::from_millis(200), rx.recv())
2403 .await
2404 .expect("upstream error arrived")
2405 .expect("channel alive");
2406 assert!(matches!(msg, Msg::UpstreamError { .. }));
2407
2408 tokio::task::yield_now().await;
2410
2411 r.dispatch(Cmd::SetTerminalTitle("x".to_string()));
2413 assert_eq!(
2414 r.scope_count(),
2415 0,
2416 "completed scope must be reaped on next dispatch"
2417 );
2418 }
2419
2420 #[tokio::test]
2421 async fn dispatch_execute_tool_under_turn_emits_tool_started() {
2422 let (mut r, mut rx) = runner();
2423 let turn = TurnId(7);
2424 let call_id = ToolCallId(1);
2425 let source = mermaid_model::models::tool_call::ToolCall {
2426 id: Some("c1".to_string()),
2427 function: mermaid_model::models::tool_call::FunctionCall {
2428 name: "read_file".to_string(),
2429 arguments: serde_json::json!({"path": "x"}),
2430 },
2431 };
2432 r.dispatch(Cmd::ExecuteTool {
2433 turn,
2434 call_id,
2435 source,
2436 dispatch: mermaid_domain::ToolDispatch {
2437 model_id: "ollama/test".to_string(),
2438 safety_mode: mermaid_runtime::SafetyMode::Ask,
2439 plan_file: None,
2440 plan_permissions: mermaid_domain::PlanPermissions::default(),
2441 context_percent: None,
2442 intent: None,
2443 session_id: "sess-test".to_string(),
2444 message_index: 0,
2445 scratchpad: None,
2446 },
2447 });
2448 let first = tokio::time::timeout(Duration::from_millis(200), rx.recv())
2449 .await
2450 .expect("some msg")
2451 .expect("channel alive");
2452 assert!(matches!(
2453 first,
2454 Msg::ToolStarted {
2455 turn: t,
2456 call_id: c,
2457 } if t == turn && c == call_id
2458 ));
2459 }
2460
2461 #[tokio::test]
2462 async fn cancel_scope_before_execute_tool_drops_pending_work() {
2463 let (mut r, _rx) = runner();
2464 let turn = TurnId(9);
2465 r.dispatch(Cmd::CallModel {
2466 turn,
2467 request: mermaid_domain::ChatRequest {
2468 model_id: "m".to_string(),
2469 messages: vec![],
2470 system_prompt: String::new(),
2471 instructions: None,
2472 reasoning: mermaid_model::models::ReasoningLevel::Medium,
2473 temperature: 0.7,
2474 max_tokens: 4096,
2475 tools: vec![],
2476
2477 ollama_num_ctx: None,
2478 ollama_allow_ram_offload: None,
2479 resolved_context_window: None,
2480 resolved_max_output: None,
2481 output_schema: None,
2482 suppress_auto_compact: false,
2483 suppressed_builtin_tools: Vec::new(),
2484 },
2485 });
2486 assert_eq!(r.scope_count(), 1);
2487
2488 r.dispatch(Cmd::CancelScope(turn));
2489 assert_eq!(r.scope_count(), 0);
2490 }
2491
2492 #[tokio::test]
2493 async fn tombstoned_turn_is_not_resurrected_by_late_scoped_cmd() {
2494 let (mut r, _rx) = runner();
2500 let req = || mermaid_domain::ChatRequest {
2501 model_id: "test/m".to_string(),
2502 messages: vec![],
2503 system_prompt: String::new(),
2504 instructions: None,
2505 reasoning: mermaid_model::models::ReasoningLevel::Medium,
2506 temperature: 0.7,
2507 max_tokens: 4096,
2508 tools: vec![],
2509 ollama_num_ctx: None,
2510 ollama_allow_ram_offload: None,
2511 resolved_context_window: None,
2512 resolved_max_output: None,
2513 output_schema: None,
2514 suppress_auto_compact: false,
2515 suppressed_builtin_tools: Vec::new(),
2516 };
2517 let turn = TurnId(123);
2518
2519 r.dispatch(Cmd::CallModel {
2520 turn,
2521 request: req(),
2522 });
2523 assert_eq!(r.scope_count(), 1);
2524
2525 r.dispatch(Cmd::CancelScope(turn));
2527 assert_eq!(r.scope_count(), 0);
2528
2529 r.dispatch(Cmd::CallModel {
2531 turn,
2532 request: req(),
2533 });
2534 assert_eq!(
2535 r.scope_count(),
2536 0,
2537 "a cancelled turn must not be resurrected by a late scoped Cmd"
2538 );
2539
2540 r.dispatch(Cmd::CallModel {
2542 turn: TurnId(124),
2543 request: req(),
2544 });
2545 assert_eq!(
2546 r.scope_count(),
2547 1,
2548 "a fresh turn must still create its scope normally"
2549 );
2550 }
2551
2552 #[tokio::test]
2553 async fn shutdown_drains_pending_saves() {
2554 let (mut r, _rx) = runner();
2555 for _ in 0..5 {
2556 r.dispatch(Cmd::SaveConversation {
2557 snapshot: mermaid_domain::ConversationHistory::new(
2558 "/p".to_string(),
2559 "m".to_string(),
2560 chrono::Local::now(),
2561 ),
2562 events: Vec::new(),
2563 });
2564 }
2565 let start = std::time::Instant::now();
2567 r.shutdown().await;
2568 assert!(start.elapsed() < Duration::from_secs(2));
2569 }
2570
2571 fn persistence_fixture(
2572 root: &std::path::Path,
2573 record_id: &str,
2574 ) -> (mermaid_domain::ConversationHistory, PendingCompactionSave) {
2575 let now = chrono::Local::now();
2576 let mut full = mermaid_domain::ConversationHistory::new(
2577 root.display().to_string(),
2578 "test/model".to_string(),
2579 now,
2580 );
2581 full.add_messages(
2582 &[mermaid_model::models::ChatMessage::user("raw history")],
2583 now,
2584 );
2585 let mut compacted = full.clone();
2586 compacted.replace_messages(
2587 vec![mermaid_model::models::ChatMessage::user(
2588 "compacted checkpoint",
2589 )],
2590 now,
2591 );
2592 let record = mermaid_domain::CompactionEvent {
2593 id: record_id.to_string(),
2594 trigger: mermaid_domain::CompactionTrigger::Manual,
2595 created_at: now,
2596 before_tokens: 100,
2597 after_tokens: 20,
2598 archived_message_count: 1,
2599 preserved_message_count: 1,
2600 preserved_turn_count: 1,
2601 summary_tokens: 10,
2602 duration_secs: 0.1,
2603 review_status: mermaid_domain::CompactionReviewStatus::Reviewed,
2604 review_error: None,
2605 focus: None,
2606 archive_path: None,
2607 };
2608 (
2609 full,
2610 PendingCompactionSave {
2611 record,
2612 conversation: compacted,
2613 events: vec![mermaid_domain::SessionEvent::Input {
2616 text: "compaction boundary".to_string(),
2617 }],
2618 events_appended: false,
2619 task_id: None,
2620 },
2621 )
2622 }
2623
2624 fn block_event_log(root: &std::path::Path, id: &str) {
2629 let dir = root.join(".mermaid").join("conversations");
2630 std::fs::create_dir_all(&dir).expect("conversations dir");
2631 std::fs::create_dir_all(dir.join(format!("{id}.jsonl"))).expect("plant a blocker");
2632 }
2633
2634 fn one_message_save(
2636 conversation: &mut mermaid_domain::ConversationHistory,
2637 text: &str,
2638 ) -> PersistenceJob {
2639 let message = mermaid_model::models::ChatMessage::user(text);
2640 conversation.add_messages(std::slice::from_ref(&message), chrono::Local::now());
2641 PersistenceJob::Conversation {
2642 snapshot: Box::new(conversation.clone()),
2643 events: vec![mermaid_domain::SessionEvent::Message { message }],
2644 }
2645 }
2646
2647 #[test]
2648 fn the_checkpoint_stops_being_written_on_every_save() {
2649 let root = std::env::temp_dir().join(format!(
2653 "mermaid-throttle-{}-{:?}",
2654 std::process::id(),
2655 std::thread::current().id()
2656 ));
2657 let _ = std::fs::remove_dir_all(&root);
2658 let manager = crate::session::ConversationManager::new(&root).unwrap();
2659 let mut conversation = mermaid_domain::ConversationHistory::new(
2660 root.display().to_string(),
2661 "test/model".to_string(),
2662 chrono::Local::now(),
2663 );
2664 let mut state = PersistenceState::new(root.clone());
2665
2666 state
2669 .process(one_message_save(&mut conversation, "first"))
2670 .1
2671 .unwrap();
2672 let checkpoint = manager
2673 .conversations_dir()
2674 .join(format!("{}.json", conversation.id));
2675 assert!(
2676 !checkpoint.exists(),
2677 "a single save must not rewrite the transcript"
2678 );
2679
2680 let resumed = manager.load_conversation(&conversation.id).unwrap();
2682 assert_eq!(resumed.messages().len(), 1);
2683 assert_eq!(resumed.messages()[0].content, "first");
2684
2685 for i in 0..CHECKPOINT_EVERY_EVENTS {
2687 state
2688 .process(one_message_save(&mut conversation, &format!("m{i}")))
2689 .1
2690 .unwrap();
2691 }
2692 assert!(
2693 checkpoint.exists(),
2694 "crossing {CHECKPOINT_EVERY_EVENTS} events must materialize a checkpoint"
2695 );
2696 let resumed = manager.load_conversation(&conversation.id).unwrap();
2697 assert_eq!(resumed.messages().len(), CHECKPOINT_EVERY_EVENTS + 1);
2698 let _ = std::fs::remove_dir_all(root);
2699 }
2700
2701 #[test]
2702 fn shutdown_flushes_the_checkpoint_it_was_holding() {
2703 let root = std::env::temp_dir().join(format!(
2704 "mermaid-flush-{}-{:?}",
2705 std::process::id(),
2706 std::thread::current().id()
2707 ));
2708 let _ = std::fs::remove_dir_all(&root);
2709 let manager = crate::session::ConversationManager::new(&root).unwrap();
2710 let mut conversation = mermaid_domain::ConversationHistory::new(
2711 root.display().to_string(),
2712 "test/model".to_string(),
2713 chrono::Local::now(),
2714 );
2715 let mut state = PersistenceState::new(root.clone());
2716 state
2717 .process(one_message_save(&mut conversation, "only message"))
2718 .1
2719 .unwrap();
2720
2721 let checkpoint = manager
2722 .conversations_dir()
2723 .join(format!("{}.json", conversation.id));
2724 assert!(!checkpoint.exists());
2725 state.flush_checkpoints().unwrap();
2726 assert!(
2727 checkpoint.exists(),
2728 "a clean exit must leave a current checkpoint"
2729 );
2730 let raw = std::fs::read_to_string(&checkpoint).unwrap();
2732 let value: serde_json::Value = serde_json::from_str(&raw).unwrap();
2733 assert!(
2734 value.get("checkpoint_seq").is_some(),
2735 "a flushed checkpoint must be placeable in its log: {raw}"
2736 );
2737 let _ = std::fs::remove_dir_all(root);
2738 }
2739
2740 #[test]
2741 fn a_failed_append_keeps_its_events_for_the_next_save() {
2742 let root = std::env::temp_dir().join(format!(
2745 "mermaid-unappended-{}-{:?}",
2746 std::process::id(),
2747 std::thread::current().id()
2748 ));
2749 let _ = std::fs::remove_dir_all(&root);
2750 let manager = crate::session::ConversationManager::new(&root).unwrap();
2751 let mut conversation = mermaid_domain::ConversationHistory::new(
2752 root.display().to_string(),
2753 "test/model".to_string(),
2754 chrono::Local::now(),
2755 );
2756 let mut state = PersistenceState::new(root.clone());
2757
2758 std::fs::create_dir_all(
2760 manager
2761 .conversations_dir()
2762 .join(format!("{}.jsonl", conversation.id)),
2763 )
2764 .expect("plant a blocker");
2765 let job = one_message_save(&mut conversation, "must survive");
2766 assert!(state.process(job).1.is_err(), "the append must fail");
2767 assert_eq!(
2768 state.unappended.get(&conversation.id).map(Vec::len),
2769 Some(1),
2770 "the batch must be held, not dropped"
2771 );
2772
2773 std::fs::remove_dir_all(
2775 manager
2776 .conversations_dir()
2777 .join(format!("{}.jsonl", conversation.id)),
2778 )
2779 .expect("unblock");
2780 state
2781 .process(one_message_save(&mut conversation, "and this one"))
2782 .1
2783 .unwrap();
2784 assert!(state.unappended.is_empty(), "the hold must clear");
2785 state.flush_checkpoints().unwrap();
2786
2787 let resumed = manager.load_conversation(&conversation.id).unwrap();
2788 let texts: Vec<&str> = resumed
2789 .messages()
2790 .iter()
2791 .map(|m| m.content.as_str())
2792 .collect();
2793 assert!(
2794 texts.contains(&"must survive"),
2795 "the event held over a failed append must reach the log: {texts:?}"
2796 );
2797 assert!(texts.contains(&"and this one"), "{texts:?}");
2798 let _ = std::fs::remove_dir_all(root);
2799 }
2800
2801 #[test]
2802 fn persistence_orders_compaction_before_newer_conversation_save() {
2803 let root = std::env::temp_dir().join(format!(
2804 "mermaid-persistence-order-{}-{:?}",
2805 std::process::id(),
2806 std::thread::current().id()
2807 ));
2808 let _ = std::fs::remove_dir_all(&root);
2809 let (full, compaction) = persistence_fixture(&root, "compact_ordered");
2810 let manager = crate::session::ConversationManager::new(&root).unwrap();
2811 manager.save_conversation(&full).unwrap();
2812
2813 let mut state = PersistenceState::new(root.clone());
2814 let (events, outcome) =
2815 state.process(PersistenceJob::Compaction(Box::new(compaction.clone())));
2816 outcome.unwrap();
2817 assert_eq!(events.len(), 1);
2818 let mut newer = compaction.conversation;
2819 let reply = mermaid_model::models::ChatMessage::assistant("new assistant reply");
2820 newer.add_messages(std::slice::from_ref(&reply), chrono::Local::now());
2821 let (_, outcome) = state.process(PersistenceJob::Conversation {
2826 snapshot: Box::new(newer),
2827 events: vec![mermaid_domain::SessionEvent::Message { message: reply }],
2828 });
2829 outcome.unwrap();
2830
2831 let loaded = crate::session::ConversationManager::new(&root)
2832 .unwrap()
2833 .load_conversation(&full.id)
2834 .unwrap();
2835 assert!(
2836 loaded
2837 .messages()
2838 .iter()
2839 .any(|message| message.content == "new assistant reply")
2840 );
2841 let _ = std::fs::remove_dir_all(root);
2842 }
2843
2844 #[test]
2845 fn failed_event_append_blocks_later_stripped_conversation_save() {
2846 let root = std::env::temp_dir().join(format!(
2847 "mermaid-persistence-barrier-{}-{:?}",
2848 std::process::id(),
2849 std::thread::current().id()
2850 ));
2851 let _ = std::fs::remove_dir_all(&root);
2852 let (full, compaction) = persistence_fixture(&root, "compact_blocked");
2853 let manager = crate::session::ConversationManager::new(&root).unwrap();
2854 manager.save_conversation(&full).unwrap();
2855 block_event_log(&root, &full.id);
2856
2857 let mut state = PersistenceState::new(root.clone());
2858 assert!(
2859 state
2860 .process(PersistenceJob::Compaction(Box::new(compaction.clone())))
2861 .1
2862 .is_err()
2863 );
2864 assert!(
2865 state
2866 .process(PersistenceJob::Conversation {
2867 snapshot: Box::new(compaction.conversation),
2868 events: Vec::new(),
2869 })
2870 .1
2871 .is_err()
2872 );
2873 assert_eq!(state.blocked.get(&full.id).map(VecDeque::len), Some(1));
2874
2875 let loaded = crate::session::ConversationManager::new(&root)
2876 .unwrap()
2877 .load_conversation(&full.id)
2878 .unwrap();
2879 assert_eq!(loaded.messages()[0].content, "raw history");
2880 let _ = std::fs::remove_dir_all(root);
2881 }
2882
2883 #[test]
2884 fn blocked_barrier_queues_a_new_compaction_instead_of_dropping_it() {
2885 let root = std::env::temp_dir().join(format!(
2886 "mermaid-persistence-queue-{}-{:?}",
2887 std::process::id(),
2888 std::thread::current().id()
2889 ));
2890 let _ = std::fs::remove_dir_all(&root);
2891 let (full, first) = persistence_fixture(&root, "compact_first");
2892 let mut second = first.clone();
2893 second.record.id = "compact_second".to_string();
2894 block_event_log(&root, &full.id);
2895
2896 let mut state = PersistenceState::new(root.clone());
2897 assert!(
2898 state
2899 .process(PersistenceJob::Compaction(Box::new(first)))
2900 .1
2901 .is_err()
2902 );
2903 assert!(
2906 state
2907 .process(PersistenceJob::Compaction(Box::new(second)))
2908 .1
2909 .is_err()
2910 );
2911 let queued = state.blocked.get(&full.id).expect("barrier queue");
2912 assert_eq!(queued.len(), 2);
2913 assert_eq!(queued[0].record.id, "compact_first");
2914 assert_eq!(queued[1].record.id, "compact_second");
2915 let _ = std::fs::remove_dir_all(root);
2916 }
2917
2918 #[test]
2919 fn retry_all_blocked_attempts_every_conversation() {
2920 let root = std::env::temp_dir().join(format!(
2921 "mermaid-persistence-drain-{}-{:?}",
2922 std::process::id(),
2923 std::thread::current().id()
2924 ));
2925 let _ = std::fs::remove_dir_all(&root);
2926 let (bad_full, bad) = persistence_fixture(&root, "compact_bad");
2927 let (mut good_full, mut good) = persistence_fixture(&root, "compact_good");
2928 block_event_log(&root, &bad_full.id);
2929 good_full.id = "20990101_000000_001".to_string();
2933 good.conversation.id = good_full.id.clone();
2934
2935 let mut state = PersistenceState::new(root.clone());
2936 state
2937 .blocked
2938 .entry(bad_full.id.clone())
2939 .or_default()
2940 .push_back(bad);
2941 state
2942 .blocked
2943 .entry(good_full.id.clone())
2944 .or_default()
2945 .push_back(good);
2946
2947 let (events, outcome) = state.retry_all_blocked();
2951 assert!(outcome.is_err());
2952 assert_eq!(events.len(), 1);
2953 assert_eq!(events[0].id, "compact_good");
2954 assert!(!state.blocked.contains_key(&good_full.id));
2955 assert_eq!(state.blocked.get(&bad_full.id).map(VecDeque::len), Some(1));
2956 let loaded = crate::session::ConversationManager::new(&root)
2957 .unwrap()
2958 .load_conversation(&good_full.id)
2959 .unwrap();
2960 assert_eq!(loaded.messages()[0].content, "compacted checkpoint");
2961 let _ = std::fs::remove_dir_all(root);
2962 }
2963
2964 #[test]
2965 fn partially_drained_barrier_reports_its_persisted_events() {
2966 let root = std::env::temp_dir().join(format!(
2967 "mermaid-persistence-partial-{}-{:?}",
2968 std::process::id(),
2969 std::thread::current().id()
2970 ));
2971 let _ = std::fs::remove_dir_all(&root);
2972 let (full, good) = persistence_fixture(&root, "compact_good");
2973 let mut bad = good.clone();
2974 bad.record.id = "compact_bad".to_string();
2975 bad.conversation.id = "../invalid".to_string();
2978
2979 let mut state = PersistenceState::new(root.clone());
2980 let queue = state.blocked.entry(full.id.clone()).or_default();
2981 queue.push_back(good);
2982 queue.push_back(bad);
2983
2984 let (events, outcome) = state.retry_blocked(&full.id);
2988 assert!(outcome.is_err());
2989 assert_eq!(events.len(), 1);
2990 assert_eq!(events[0].id, "compact_good");
2991 assert_eq!(state.blocked.get(&full.id).map(VecDeque::len), Some(1));
2992 let _ = std::fs::remove_dir_all(root);
2993 }
2994}