1mod harness_adapters;
13pub(crate) mod history;
14mod provenance;
15pub mod tags;
16
17use std::collections::{BTreeMap, BTreeSet};
18use std::path::{Path, PathBuf};
19use std::sync::Arc;
20use std::sync::atomic::{AtomicBool, Ordering};
21use std::time::{Duration, Instant};
22
23use anyhow::{Context, Result};
24use chrono::{DateTime, Utc};
25
26use mj_client::daemon::{
27 SessionTextMatch, SessionTextMatchKind, WikiHitBlock, WikiHitTranscript, WikiIndexState,
28 WikiRow, WikiSessionInfo, WikiSessionStatus, WikiStatus,
29};
30use mj_core::config::HarnessKind;
31use mj_core::state::{SessionRecord, State};
32use sessionwiki::adapters::{Adapter, Discovered, Store};
33use sessionwiki::model::{Message, Role, Session};
34
35use crate::controller::Controller;
36use crate::controller::checkpoint::managed_checkpoint_archive_name;
37use harness_adapters::HarnessAdapter;
38
39const TOOL: &str = "mjolnir";
43
44struct ArchiveFile {
46 path: PathBuf,
47 frontier: u64,
48 token: i64,
50}
51
52#[derive(Default)]
56struct Sessions {
57 records: mj_core::snapshot_map::SnapshotMap<String, SessionRecord>,
58 subagent_ids: BTreeSet<String>,
59 live: BTreeMap<String, i64>,
61}
62
63impl Sessions {
64 fn of(state: &State) -> Self {
65 Self {
66 records: state.sessions.clone(),
67 subagent_ids: state.subagents.keys().cloned().collect(),
68 live: live_tokens(state),
69 }
70 }
71}
72
73fn live_tokens(state: &State) -> BTreeMap<String, i64> {
80 let activity = match crate::database::load_transcribed_session_activity() {
81 Ok(activity) => activity,
82 Err(error) => {
83 tracing::warn!(%error, "could not read session activity for SessionWiki");
84 return BTreeMap::new();
85 }
86 };
87 state
88 .sessions
89 .iter()
90 .filter(|(_, record)| record.state != mj_core::state::SessionState::Stopped)
91 .filter_map(|(session_id, _)| {
92 let watermark = activity.get(session_id)?;
93 Some((session_id.clone(), watermark.unwrap_or_default() / 1000))
94 })
95 .collect()
96}
97
98pub struct MjolnirAdapter {
100 sessions_dir: PathBuf,
101 sessions: std::sync::Mutex<Sessions>,
102 reload: bool,
104}
105
106impl MjolnirAdapter {
107 pub fn from_state(state: &State) -> Self {
110 Self {
111 sessions_dir: mj_core::config::sessions_dir(),
112 sessions: std::sync::Mutex::new(Sessions::of(state)),
113 reload: false,
114 }
115 }
116
117 pub fn reloading(state: &State) -> Self {
126 Self {
127 reload: true,
128 ..Self::from_state(state)
129 }
130 }
131
132 pub fn indexed_tags(&self) -> BTreeMap<String, tags::MjTags> {
141 let sessions = self
142 .sessions
143 .lock()
144 .unwrap_or_else(std::sync::PoisonError::into_inner);
145 sessions
146 .records
147 .iter()
148 .map(|(session_id, record)| {
149 (
150 session_id.clone(),
151 tags::MjTags {
152 target: Some(record.target_template_id.clone()).filter(|id| !id.is_empty()),
153 profile: Some(record.last_profile.clone()).filter(|id| !id.is_empty()),
154 harness: Some(record.harness_kind.id().to_owned()),
155 },
156 )
157 })
158 .collect()
159 }
160
161 fn reload(&self) {
162 if !self.reload {
163 return;
164 }
165 match Controller::load() {
166 Ok(controller) => {
167 *self
168 .sessions
169 .lock()
170 .unwrap_or_else(std::sync::PoisonError::into_inner) =
171 Sessions::of(&controller.state)
172 }
173 Err(error) => {
174 tracing::warn!(%error, "could not refresh session records for SessionWiki")
175 }
176 }
177 }
178
179 fn checkpointed_transcript(&self, session_id: &str) -> Result<IndexedTranscript> {
182 let (newest, _) = self.newest_archives();
183 let archive = newest
184 .get(session_id)
185 .with_context(|| format!("no checkpoint archive for session {session_id}"))?;
186 let snapshot = mj_checkpoint::archive::read_archive_verified(&archive.path)
187 .with_context(|| format!("read checkpoint {}", archive.path.display()))?
188 .canonical_session()
189 .with_context(|| format!("read the transcript of session {session_id}"))?;
190 let mut evidence = provenance::Evidence::default();
191 for item in &snapshot.transcript {
192 if let mj_core::archive::CanonicalTranscriptBody::Tool { call, .. } = &item.body {
193 evidence.observe(call, item.created_at_ms);
194 }
195 }
196 let messages = summary_messages(mj_transcript::summary::TranscriptSummary::from_snapshot(
197 &snapshot,
198 ));
199 Ok(IndexedTranscript {
200 messages,
201 title: snapshot.session.session_title.clone(),
202 evidence,
203 })
204 }
205
206 fn projected_transcript(&self, session_id: &str) -> Result<IndexedTranscript> {
210 let projection = crate::database::load_materialized_session(session_id)
211 .with_context(|| format!("read the stored transcript of session {session_id}"))?
212 .with_context(|| format!("no stored transcript for session {session_id}"))?;
213 let mut evidence = provenance::Evidence::default();
214 for item in &projection.transcript {
215 if let mj_core::state::TranscriptBody::Tool { call, .. } = &item.body {
216 evidence.observe(call, item.created_at_ms);
217 }
218 }
219 Ok(IndexedTranscript {
220 messages: projected_messages(&projection),
221 title: projection.session_title.clone(),
222 evidence,
223 })
224 }
225
226 fn key_for(&self, session_id: &str) -> String {
229 format!("{}/{session_id}", self.sessions_dir.display())
230 }
231
232 fn newest_archives(&self) -> (BTreeMap<String, ArchiveFile>, bool) {
238 let mut newest: BTreeMap<String, ArchiveFile> = BTreeMap::new();
239 let mut had_error = false;
240 let entries = match std::fs::read_dir(&self.sessions_dir) {
241 Ok(entries) => entries,
242 Err(error) => {
243 if self.sessions_dir.exists() {
244 tracing::debug!(
245 directory = %self.sessions_dir.display(),
246 %error,
247 "could not list the checkpoint directory for SessionWiki"
248 );
249 had_error = true;
250 }
251 return (newest, had_error);
252 }
253 };
254 for entry in entries {
255 let Ok(entry) = entry else {
256 had_error = true;
257 continue;
258 };
259 let Some((session_id, frontier)) = checkpoint_archive_session(&entry.file_name())
260 else {
261 continue;
262 };
263 let token = entry
264 .metadata()
265 .ok()
266 .and_then(|metadata| metadata.modified().ok())
267 .and_then(|modified| modified.duration_since(std::time::UNIX_EPOCH).ok())
268 .map(|age| age.as_secs() as i64)
269 .unwrap_or(0);
270 let candidate = ArchiveFile {
271 path: entry.path(),
272 frontier,
273 token,
274 };
275 match newest.get(&session_id) {
276 Some(existing) if existing.frontier >= candidate.frontier => {}
277 _ => {
278 newest.insert(session_id, candidate);
279 }
280 }
281 }
282 (newest, had_error)
283 }
284}
285
286struct IndexedTranscript {
287 messages: Vec<Message>,
288 title: Option<String>,
289 evidence: provenance::Evidence,
290}
291
292fn checkpoint_archive_session(name: &std::ffi::OsStr) -> Option<(String, u64)> {
297 if let Some(parsed) = managed_checkpoint_archive_name(name) {
298 return Some((parsed.session_id, parsed.frontier));
299 }
300 let stem = name
301 .to_str()
302 .and_then(|name| name.strip_suffix(".hel.zip"))?;
303 mj_core::config::validate_id("session", stem)
304 .is_ok()
305 .then(|| (stem.to_owned(), 0))
306}
307
308fn projected_messages(projection: &mj_core::state::MaterializedSession) -> Vec<Message> {
310 summary_messages(mj_transcript::summary::TranscriptSummary::from_materialized(projection))
311}
312
313fn summary_messages(summary: mj_transcript::summary::TranscriptSummary) -> Vec<Message> {
314 use mj_transcript::summary::SummaryRole;
315 summary
316 .entries
317 .into_iter()
318 .filter_map(|entry| {
319 let role = match entry.role {
320 SummaryRole::User => Role::User,
321 SummaryRole::Assistant => Role::Assistant,
322 SummaryRole::Tool => Role::Tool,
323 SummaryRole::Plan => return None,
324 };
325 message(role, entry.body(), entry.created_at_ms)
326 })
327 .collect()
328}
329
330fn message(role: Role, text: String, created_at_ms: i64) -> Option<Message> {
332 let text = text.trim().to_owned();
333 (!text.is_empty()).then(|| Message {
334 role,
335 text,
336 ts: DateTime::from_timestamp_millis(created_at_ms),
337 })
338}
339
340fn parse_time(value: &str) -> Option<DateTime<Utc>> {
341 DateTime::parse_from_rfc3339(value)
342 .ok()
343 .map(|time| time.with_timezone(&Utc))
344}
345
346impl Adapter for MjolnirAdapter {
347 fn name(&self) -> &'static str {
348 TOOL
349 }
350
351 fn root(&self) -> Option<PathBuf> {
352 Some(self.sessions_dir.clone())
353 }
354
355 fn discover(&self) -> Discovered {
358 Discovered {
359 files: Vec::new(),
360 had_error: false,
361 }
362 }
363
364 fn parse(&self, _path: &Path) -> Result<Session> {
365 anyhow::bail!("Mjolnir sessions are parsed by key, not by file")
366 }
367
368 fn store(&self) -> Option<Store> {
369 self.reload();
370 let (newest, had_error) = self.newest_archives();
371 let mut files = Vec::with_capacity(newest.len());
372 let mut tokens: BTreeMap<String, i64> = BTreeMap::new();
373 for (session_id, archive) in newest {
374 tokens.insert(session_id, archive.token);
375 files.push(archive.path);
376 }
377 let sessions = self
382 .sessions
383 .lock()
384 .unwrap_or_else(std::sync::PoisonError::into_inner);
385 let live = sessions.live.clone();
386 tokens.extend(live);
387 for (session_id, token) in tokens.iter_mut() {
392 let updated = sessions
393 .records
394 .get(session_id)
395 .and_then(|record| parse_time(&record.updated_at))
396 .map(|updated| updated.timestamp());
397 if let Some(updated) = updated {
398 *token = (*token).max(updated);
399 }
400 }
401 let keys = tokens
402 .into_iter()
403 .map(|(session_id, token)| {
404 (
405 self.key_for(&session_id),
406 token
407 .saturating_mul(1024)
408 .saturating_add(i64::from(mj_transcript::summary::SUMMARY_VERSION)),
409 )
410 })
411 .collect();
412 Some(Store {
413 keys,
414 files,
415 had_error,
416 })
417 }
418
419 fn reconcile_scope(&self) -> Option<String> {
423 Some(format!("{}/", self.sessions_dir.display()))
424 }
425
426 fn parse_key(&self, key: &str) -> Result<Session> {
427 let session_id = key.rsplit('/').next().unwrap_or_default();
428 anyhow::ensure!(!session_id.is_empty(), "no session id in key {key:?}");
429 let sessions = self
430 .sessions
431 .lock()
432 .unwrap_or_else(std::sync::PoisonError::into_inner);
433 let IndexedTranscript {
434 messages,
435 title: snapshot_title,
436 evidence,
437 } = if sessions.live.contains_key(session_id) {
438 self.projected_transcript(session_id)?
439 } else {
440 self.checkpointed_transcript(session_id)?
441 };
442 let record = sessions.records.get(session_id);
443
444 let title = record
445 .and_then(|record| record.session_title_override.clone())
446 .or_else(|| record.and_then(|record| record.acp_session_title.clone()))
447 .or_else(|| snapshot_title.clone())
448 .unwrap_or_else(|| {
449 messages
450 .iter()
451 .find(|message| message.role == Role::User)
452 .map(|message| message.text.chars().take(80).collect())
453 .unwrap_or_default()
454 });
455
456 Ok(Session {
457 id: session_id.to_owned(),
458 tool: TOOL,
459 path: PathBuf::from(key),
460 project: record
461 .and_then(|record| record.project_directory.as_ref())
462 .map(|directory| directory.display().to_string())
463 .unwrap_or_default(),
464 started: record.and_then(|record| parse_time(&record.created_at)),
465 ended: record.and_then(|record| parse_time(&record.updated_at)),
466 title,
467 subagent: sessions.subagent_ids.contains(session_id),
468 messages,
469 touched: evidence.paths.into_iter().collect(),
470 edits: evidence.edits,
471 })
472 }
473}
474
475struct SharedMjolnirAdapter(Arc<MjolnirAdapter>);
483
484impl Adapter for SharedMjolnirAdapter {
485 fn name(&self) -> &'static str {
486 self.0.name()
487 }
488
489 fn root(&self) -> Option<PathBuf> {
490 self.0.root()
491 }
492
493 fn discover(&self) -> Discovered {
494 self.0.discover()
495 }
496
497 fn parse(&self, path: &Path) -> Result<Session> {
498 self.0.parse(path)
499 }
500
501 fn store(&self) -> Option<Store> {
502 self.0.store()
503 }
504
505 fn parse_key(&self, key: &str) -> Result<Session> {
506 self.0.parse_key(key)
507 }
508
509 fn reconcile_scope(&self) -> Option<String> {
510 self.0.reconcile_scope()
511 }
512}
513
514pub struct WikiIndexer {
521 inner: Arc<Indexer>,
522}
523
524#[derive(Default)]
525struct Indexer {
526 running: tokio::sync::Mutex<()>,
528 notify: tokio::sync::Notify,
529 requested: AtomicBool,
531 full_requested: AtomicBool,
533 in_flight: AtomicBool,
536 last_success: std::sync::Mutex<Option<Success>>,
537}
538
539#[derive(Clone, Copy)]
540struct Success {
541 at: Instant,
542 epoch_seconds: i64,
543}
544
545impl WikiIndexer {
546 pub fn spawn() -> Self {
549 let inner = Arc::new(Indexer::default());
550 if let Ok(handle) = tokio::runtime::Handle::try_current() {
551 let worker = Arc::clone(&inner);
552 handle.spawn(async move { worker.run().await });
553 }
554 Self { inner }
555 }
556
557 pub fn request_sync(&self, full: bool) {
559 if full {
560 self.inner.full_requested.store(true, Ordering::Release);
561 }
562 self.inner.requested.store(true, Ordering::Release);
563 self.inner.notify.notify_one();
564 }
565
566 #[cfg(test)]
569 pub(crate) fn inert() -> Self {
570 Self {
571 inner: Arc::new(Indexer::default()),
572 }
573 }
574
575 #[cfg(test)]
577 pub(crate) fn sync_requested(&self) -> bool {
578 self.inner.requested.load(Ordering::Acquire)
579 }
580
581 pub async fn sync_now(&self, full: bool) -> Result<()> {
583 self.inner.sync(full).await
584 }
585
586 pub fn status(&self) -> WikiStatus {
589 WikiStatus {
590 state: index_state(),
591 topping_up: self.inner.in_flight.load(Ordering::Acquire)
592 || self.inner.requested.load(Ordering::Acquire),
593 }
594 }
595
596 pub fn last_success(&self) -> Option<Instant> {
598 self.inner
599 .last_success
600 .lock()
601 .unwrap_or_else(std::sync::PoisonError::into_inner)
602 .map(|success| success.at)
603 }
604}
605
606impl Indexer {
607 async fn run(self: Arc<Self>) {
608 loop {
609 self.notify.notified().await;
610 while self.requested.swap(false, Ordering::AcqRel) {
611 let full = self.full_requested.swap(false, Ordering::AcqRel);
612 if let Err(error) = self.sync(full).await {
613 self.report(&error);
614 break;
619 }
620 }
621 }
622 }
623
624 fn report(&self, error: &anyhow::Error) {
628 if crate::database::is_busy_error(error) {
629 self.requested.store(true, Ordering::Release);
630 tracing::debug!(%error, "the SessionWiki index was busy; retrying on the next trigger");
631 } else {
632 tracing::warn!(%error, "could not sync sessions into SessionWiki");
633 }
634 }
635
636 async fn sync(&self, full: bool) -> Result<()> {
637 let _guard = self.running.lock().await;
638 let since = if full {
639 None
640 } else {
641 self.last_success
642 .lock()
643 .unwrap_or_else(std::sync::PoisonError::into_inner)
644 .map(|success| success.epoch_seconds - 60)
647 };
648 let started = Instant::now();
649 let work = crate::upgrade::activity("SessionWiki sync")?;
650 self.in_flight.store(true, Ordering::Release);
651 let ran = tokio::task::spawn_blocking(move || {
652 let _work = work;
653 sync_blocking(since)
654 })
655 .await;
656 self.in_flight.store(false, Ordering::Release);
657 let ran = ran.context("run the SessionWiki sync")??;
658 if ran {
659 *self
660 .last_success
661 .lock()
662 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(Success {
663 at: started,
664 epoch_seconds: Utc::now().timestamp(),
665 });
666 }
667 Ok(())
668 }
669}
670
671fn sync_blocking(since: Option<i64>) -> Result<bool> {
674 if !index_is_writable() {
675 return Ok(false);
676 }
677 let controller =
678 Controller::load().context("load controller state for the SessionWiki sync")?;
679 let mjolnir = Arc::new(MjolnirAdapter::reloading(&controller.state));
682 let mut adapters: Vec<Box<dyn sessionwiki::adapters::Adapter>> =
683 vec![Box::new(SharedMjolnirAdapter(Arc::clone(&mjolnir)))];
684 adapters.extend(native_adapters(&controller.config));
685 let mut connection = sessionwiki::index::open().context("open the SessionWiki index")?;
686 sessionwiki::index::sync_with(&mut connection, &adapters, since)
687 .context("sync the SessionWiki index")?;
688 write_session_tags(&mut connection, &mjolnir.indexed_tags())
689 .context("store Mjolnir's session metadata in the SessionWiki index")?;
690 provenance::backfill(&mut connection, &mjolnir).context("backfill Mjolnir file provenance")?;
691 if since.is_none() {
692 record_first_build();
696 }
697 Ok(true)
698}
699
700fn write_session_tags(
710 connection: &mut rusqlite::Connection,
711 session_tags: &BTreeMap<String, tags::MjTags>,
712) -> Result<()> {
713 if session_tags.is_empty() {
714 return Ok(());
715 }
716 let transaction = connection
717 .transaction()
718 .context("open a transaction for the session metadata")?;
719 for (session_id, session) in session_tags {
720 if session.is_empty() {
721 continue;
722 }
723 tags::write(&transaction, session_id, session)?;
724 }
725 transaction
726 .commit()
727 .context("commit the session metadata")?;
728 Ok(())
729}
730
731fn native_adapters(config: &mj_core::config::Config) -> Vec<Box<dyn Adapter>> {
748 let mut seen: BTreeSet<(HarnessKind, &Path)> = BTreeSet::new();
751 let mut adapters: Vec<Box<dyn Adapter>> = Vec::new();
752 for (_, profile) in config.enabled_profiles() {
753 if !seen.insert((profile.kind, profile.home.as_path())) {
755 continue;
756 }
757 let adapter: Box<dyn Adapter> = match profile.kind {
758 HarnessKind::Codex => {
759 Box::new(sessionwiki::adapters::Codex::in_home(profile.home.clone()))
760 }
761 HarnessKind::Claude => Box::new(sessionwiki::adapters::ClaudeCode::in_home(
762 profile.home.clone(),
763 )),
764 kind => match HarnessAdapter::in_home(kind, profile.home.clone()) {
765 Some(adapter) => Box::new(adapter),
766 None => continue,
767 },
768 };
769 adapters.push(adapter);
770 }
771 adapters.extend(
772 sessionwiki::adapters::all()
773 .into_iter()
774 .filter(|adapter| !matches!(adapter.name(), "codex" | "claude-code")),
775 );
776 adapters
777}
778
779fn index_is_isolated() -> bool {
792 static SAID: AtomicBool = AtomicBool::new(false);
793 if mj_core::config::session_index_is_resolved()
794 || std::env::var_os(mj_core::config::SESSION_INDEX_ENV).is_some()
795 {
796 return true;
797 }
798 if !SAID.swap(true, Ordering::AcqRel) {
799 tracing::debug!(
800 "this process did not resolve a session index location; SessionWiki is not used"
801 );
802 }
803 false
804}
805
806fn index_version_mismatch() -> bool {
815 static SAID: AtomicBool = AtomicBool::new(false);
816 let Ok(path) = sessionwiki::index::db_path() else {
817 return false;
818 };
819 if !path.exists() {
820 return false;
821 }
822 let version = rusqlite::Connection::open_with_flags(
823 &path,
824 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_URI,
825 )
826 .and_then(|connection| connection.pragma_query_value(None, "user_version", |row| row.get(0)));
827 let version: i64 = match version {
828 Ok(version) => version,
829 Err(error) => {
830 tracing::debug!(%error, "could not read the SessionWiki index schema version");
831 return false;
832 }
833 };
834 let mismatch = version != 0 && version != sessionwiki::index::SCHEMA_VERSION;
837 if mismatch && !SAID.swap(true, Ordering::AcqRel) {
838 tracing::warn!(
839 found = version,
840 expected = sessionwiki::index::SCHEMA_VERSION,
841 path = %path.display(),
842 "the SessionWiki index was written by another version; Mjolnir will not open it, because opening it would rebuild it. Install the matching sessionwiki command"
843 );
844 }
845 mismatch
846}
847
848fn index_is_writable() -> bool {
849 index_is_isolated() && !index_version_mismatch()
850}
851
852fn first_build_marker() -> PathBuf {
855 mj_core::config::data_dir().join("sessionwiki-built")
856}
857
858fn record_first_build() {
859 let path = first_build_marker();
860 let version = sessionwiki::index::SCHEMA_VERSION.to_string();
861 if std::fs::read_to_string(&path).is_ok_and(|held| held.trim() == version) {
862 return;
863 }
864 if let Err(error) = std::fs::write(&path, &version) {
865 tracing::warn!(%error, path = %path.display(), "could not record the first SessionWiki build");
866 }
867}
868
869fn first_build_is_done() -> bool {
871 std::fs::read_to_string(first_build_marker())
872 .is_ok_and(|held| held.trim() == sessionwiki::index::SCHEMA_VERSION.to_string())
873 && sessionwiki::index::db_path().is_ok_and(|path| path.exists())
874}
875
876pub fn index_state() -> WikiIndexState {
878 if !index_is_isolated() {
879 return WikiIndexState::Indexing;
880 }
881 if index_version_mismatch() {
882 return WikiIndexState::VersionMismatch;
883 }
884 if first_build_is_done() {
885 WikiIndexState::Ready
886 } else {
887 WikiIndexState::Indexing
888 }
889}
890
891pub const DESTROY_SYNC_WAIT: Duration = Duration::from_secs(20);
899
900const DEFERRED_WRITE_RETRY: Duration = Duration::from_secs(5);
904const DEFERRED_WRITE_LIMIT: Duration = Duration::from_secs(60 * 60);
905
906#[derive(Debug, Clone, PartialEq, Eq)]
909pub enum IndexedBeforeDestroy {
910 Current,
913 Synced,
915 WrittenDirectly,
918 Deferred,
921 Unavailable(&'static str),
923 Failed(String),
925}
926
927impl WikiIndexer {
928 pub async fn index_before_destroy(
938 &self,
939 session_id: &str,
940 wait: Duration,
941 ) -> IndexedBeforeDestroy {
942 if let Some(reason) = unwritable_reason() {
943 return IndexedBeforeDestroy::Unavailable(reason);
944 }
945 let root = session_id.to_owned();
946 let pending = match tokio::task::spawn_blocking(move || unindexed_session_tree(&root)).await
947 {
948 Ok(Ok(pending)) => pending,
949 Ok(Err(error)) => {
950 return IndexedBeforeDestroy::Failed(format!(
951 "could not tell whether the index holds the session: {error:#}"
952 ));
953 }
954 Err(error) => {
955 return IndexedBeforeDestroy::Failed(format!(
956 "checking the index for the session stopped: {error}"
957 ));
958 }
959 };
960 if pending.is_empty() {
961 return IndexedBeforeDestroy::Current;
962 }
963 let inner = Arc::clone(&self.inner);
964 index_before_destroy_with(
965 async move { inner.sync(false).await },
966 wait,
967 move || capture_sessions(&pending),
968 DEFERRED_WRITE_RETRY,
969 )
970 .await
971 }
972}
973
974fn unwritable_reason() -> Option<&'static str> {
976 if !index_is_isolated() {
977 return Some("this process did not resolve a SessionWiki index of its own");
978 }
979 if index_version_mismatch() {
980 return Some("the SessionWiki index was written by another SessionWiki version");
981 }
982 None
983}
984
985async fn index_before_destroy_with<S, C>(
988 sync: S,
989 wait: Duration,
990 capture: C,
991 retry: Duration,
992) -> IndexedBeforeDestroy
993where
994 S: std::future::Future<Output = Result<()>> + Send + 'static,
995 C: FnOnce() -> Result<Vec<CapturedSession>> + Send + 'static,
996{
997 match tokio::time::timeout(wait, tokio::spawn(sync)).await {
1002 Ok(Ok(Ok(()))) => return IndexedBeforeDestroy::Synced,
1003 Ok(Ok(Err(error))) => tracing::warn!(
1004 error = %format!("{error:#}"),
1005 "the SessionWiki sync before a destroy failed; indexing the session on its own"
1006 ),
1007 Ok(Err(error)) => tracing::warn!(
1008 %error,
1009 "the SessionWiki sync before a destroy stopped; indexing the session on its own"
1010 ),
1011 Err(_) => tracing::info!(
1012 wait_seconds = wait.as_secs_f64(),
1013 "the SessionWiki sync did not finish in time; indexing the session on its own"
1014 ),
1015 }
1016 let captured = match tokio::task::spawn_blocking(capture).await {
1017 Ok(Ok(captured)) => Arc::new(captured),
1018 Ok(Err(error)) => {
1019 return IndexedBeforeDestroy::Failed(format!(
1020 "could not read the session to index it: {error:#}"
1021 ));
1022 }
1023 Err(error) => {
1024 return IndexedBeforeDestroy::Failed(format!(
1025 "reading the session to index it stopped: {error}"
1026 ));
1027 }
1028 };
1029 if captured.is_empty() {
1030 return IndexedBeforeDestroy::Current;
1031 }
1032 let attempt = Arc::clone(&captured);
1033 match tokio::task::spawn_blocking(move || write_captured(&attempt)).await {
1034 Ok(Ok(())) => IndexedBeforeDestroy::WrittenDirectly,
1035 Ok(Err(error)) if crate::database::is_busy_error(&error) => {
1036 write_captured_later(captured, retry);
1037 IndexedBeforeDestroy::Deferred
1038 }
1039 Ok(Err(error)) => IndexedBeforeDestroy::Failed(format!(
1040 "could not write the session into the index: {error:#}"
1041 )),
1042 Err(error) => IndexedBeforeDestroy::Failed(format!(
1043 "writing the session into the index stopped: {error}"
1044 )),
1045 }
1046}
1047
1048fn unindexed_session_tree(root: &str) -> Result<Vec<String>> {
1051 let controller =
1052 Controller::load().context("load controller state to index a destroyed session")?;
1053 unindexed(
1054 &MjolnirAdapter::from_state(&controller.state),
1055 &session_tree(&controller.state, root),
1056 )
1057}
1058
1059fn unindexed(adapter: &MjolnirAdapter, session_ids: &[String]) -> Result<Vec<String>> {
1063 let tokens: BTreeMap<String, i64> = adapter
1064 .store()
1065 .map(|store| store.keys.into_iter().collect())
1066 .unwrap_or_default();
1067 let connection = open_readonly().ok();
1069 let mut pending = Vec::new();
1070 for session_id in session_ids {
1071 let key = adapter.key_for(session_id);
1072 let Some(&token) = tokens.get(&key) else {
1074 continue;
1075 };
1076 let current = match &connection {
1077 Some(connection) => indexed_token(connection, &key)? == Some(token),
1078 None => false,
1079 };
1080 if !current {
1081 pending.push(session_id.clone());
1082 }
1083 }
1084 Ok(pending)
1085}
1086
1087fn session_tree(state: &State, root: &str) -> Vec<String> {
1089 if !state.sessions.contains_key(root) {
1090 return Vec::new();
1091 }
1092 let mut tree = vec![root.to_owned()];
1093 let mut seen = BTreeSet::from([root.to_owned()]);
1094 let mut next = 0;
1095 while let Some(parent) = tree.get(next).cloned() {
1096 next += 1;
1097 for child in state.subagents.values() {
1098 if child.parent_session_id == parent
1099 && state.sessions.contains_key(&child.child_session_id)
1100 && seen.insert(child.child_session_id.clone())
1101 {
1102 tree.push(child.child_session_id.clone());
1103 }
1104 }
1105 }
1106 tree
1107}
1108
1109fn indexed_token(connection: &rusqlite::Connection, key: &str) -> Result<Option<i64>> {
1112 use rusqlite::OptionalExtension;
1113 connection
1114 .query_row(
1115 "SELECT mtime FROM files WHERE path = ?1 AND archived_at IS NULL",
1116 [key],
1117 |row| row.get(0),
1118 )
1119 .optional()
1120 .context("read a session's change token from the SessionWiki index")
1121}
1122
1123struct CapturedSession {
1126 key: String,
1127 token: i64,
1128 session: Session,
1129 tags: tags::MjTags,
1130}
1131
1132fn capture_sessions(session_ids: &[String]) -> Result<Vec<CapturedSession>> {
1133 let controller =
1134 Controller::load().context("load controller state to index a destroyed session")?;
1135 capture_sessions_from(&MjolnirAdapter::from_state(&controller.state), session_ids)
1136}
1137
1138fn capture_sessions_from(
1141 adapter: &MjolnirAdapter,
1142 session_ids: &[String],
1143) -> Result<Vec<CapturedSession>> {
1144 let tokens: BTreeMap<String, i64> = adapter
1145 .store()
1146 .map(|store| store.keys.into_iter().collect())
1147 .unwrap_or_default();
1148 let mut session_tags = adapter.indexed_tags();
1149 let mut captured = Vec::new();
1150 for session_id in session_ids {
1151 let key = adapter.key_for(session_id);
1152 let Some(&token) = tokens.get(&key) else {
1153 continue;
1154 };
1155 captured.push(CapturedSession {
1156 session: adapter.parse_key(&key)?,
1157 tags: session_tags.remove(session_id).unwrap_or_default(),
1158 key,
1159 token,
1160 });
1161 }
1162 Ok(captured)
1163}
1164
1165fn write_captured(captured: &Arc<Vec<CapturedSession>>) -> Result<()> {
1168 anyhow::ensure!(
1169 index_is_writable(),
1170 "this process may not write the SessionWiki index"
1171 );
1172 let mut connection = sessionwiki::index::open().context("open the SessionWiki index")?;
1173 for index in 0..captured.len() {
1174 let adapter: Box<dyn Adapter> = Box::new(CapturedAdapter {
1175 captured: Arc::clone(captured),
1176 index,
1177 });
1178 sessionwiki::index::sync_with(&mut connection, &[adapter], None)
1179 .context("index a session before it is destroyed")?;
1180 }
1181 let session_tags = captured
1182 .iter()
1183 .map(|captured| (captured.session.id.clone(), captured.tags.clone()))
1184 .collect();
1185 write_session_tags(&mut connection, &session_tags)
1186 .context("store Mjolnir's session metadata in the SessionWiki index")
1187}
1188
1189fn write_captured_later(captured: Arc<Vec<CapturedSession>>, retry: Duration) {
1194 let sessions = captured
1195 .iter()
1196 .map(|captured| captured.session.id.clone())
1197 .collect::<Vec<_>>();
1198 tracing::info!(
1199 ?sessions,
1200 "the SessionWiki index is busy; indexing the destroyed sessions once it is free"
1201 );
1202 tokio::spawn(async move {
1203 let started = Instant::now();
1204 loop {
1205 tokio::time::sleep(retry).await;
1206 let attempt = Arc::clone(&captured);
1207 let error = match tokio::task::spawn_blocking(move || write_captured(&attempt)).await {
1208 Ok(Ok(())) => {
1209 tracing::info!(?sessions, "indexed the destroyed sessions");
1210 return;
1211 }
1212 Ok(Err(error)) => error,
1213 Err(error) => anyhow::Error::new(error),
1214 };
1215 if !crate::database::is_busy_error(&error) || started.elapsed() >= DEFERRED_WRITE_LIMIT
1216 {
1217 tracing::warn!(
1218 ?sessions,
1219 error = %format!("{error:#}"),
1220 "gave up indexing destroyed sessions in SessionWiki"
1221 );
1222 return;
1223 }
1224 }
1225 });
1226}
1227
1228struct CapturedAdapter {
1231 captured: Arc<Vec<CapturedSession>>,
1232 index: usize,
1233}
1234
1235impl CapturedAdapter {
1236 fn captured(&self) -> &CapturedSession {
1237 &self.captured[self.index]
1238 }
1239}
1240
1241impl Adapter for CapturedAdapter {
1242 fn name(&self) -> &'static str {
1243 TOOL
1244 }
1245
1246 fn root(&self) -> Option<PathBuf> {
1247 Path::new(&self.captured().key)
1248 .parent()
1249 .map(Path::to_path_buf)
1250 }
1251
1252 fn discover(&self) -> Discovered {
1253 Discovered {
1254 files: Vec::new(),
1255 had_error: false,
1256 }
1257 }
1258
1259 fn parse(&self, _path: &Path) -> Result<Session> {
1260 anyhow::bail!("Mjolnir sessions are parsed by key, not by file")
1261 }
1262
1263 fn store(&self) -> Option<Store> {
1264 let captured = self.captured();
1265 Some(Store {
1266 keys: vec![(captured.key.clone(), captured.token)],
1267 files: Vec::new(),
1268 had_error: false,
1269 })
1270 }
1271
1272 fn parse_key(&self, key: &str) -> Result<Session> {
1273 let captured = self.captured();
1274 anyhow::ensure!(key == captured.key, "no captured session for key {key:?}");
1275 Ok(copy_session(&captured.session))
1276 }
1277
1278 fn reconcile_scope(&self) -> Option<String> {
1282 Some(format!("{}\0", self.captured().key))
1283 }
1284}
1285
1286fn copy_session(session: &Session) -> Session {
1289 Session {
1290 id: session.id.clone(),
1291 tool: session.tool,
1292 path: session.path.clone(),
1293 project: session.project.clone(),
1294 started: session.started,
1295 ended: session.ended,
1296 title: session.title.clone(),
1297 subagent: session.subagent,
1298 messages: session
1299 .messages
1300 .iter()
1301 .map(|message| Message {
1302 role: message.role,
1303 text: message.text.clone(),
1304 ts: message.ts,
1305 })
1306 .collect(),
1307 touched: session.touched.clone(),
1308 edits: session
1309 .edits
1310 .iter()
1311 .map(|edit| sessionwiki::model::EditEvent {
1312 path: edit.path.clone(),
1313 kind: edit.kind,
1314 snippet: edit.snippet.clone(),
1315 ts: edit.ts,
1316 })
1317 .collect(),
1318 }
1319}
1320
1321pub const MAX_WIKI_LIMIT: usize = 200;
1327pub const DEFAULT_WIKI_LIMIT: usize = 50;
1329const MIN_FULLTEXT_QUERY: usize = 3;
1332pub const SYNC_STALE_AFTER: std::time::Duration = std::time::Duration::from_secs(60);
1334
1335pub fn sync_is_stale(last_success: Option<Instant>) -> bool {
1337 last_success.is_none_or(|at| at.elapsed() >= SYNC_STALE_AFTER)
1338}
1339
1340pub fn query_rows(
1349 query: &str,
1350 limit: usize,
1351 live: &BTreeSet<String>,
1352 include_subagents: bool,
1353) -> Result<Vec<WikiRow>> {
1354 let limit = limit.clamp(1, MAX_WIKI_LIMIT);
1355 if !index_is_writable() {
1356 return Ok(Vec::new());
1360 }
1361 let connection = open_readonly()?;
1362 let query = query.trim();
1363 if query.is_empty() {
1364 let rows =
1365 sessionwiki::index::recent(&connection, limit, None, None, None, include_subagents)
1366 .context("list recent SessionWiki sessions")?;
1367 let mut rows: Vec<WikiRow> = rows
1368 .into_iter()
1369 .map(|row| wiki_row(row, None, live))
1370 .collect();
1371 fill_session_tags(&connection, &mut rows)?;
1372 return Ok(rows);
1373 }
1374 let search_limit = if include_subagents {
1377 limit
1378 } else {
1379 MAX_WIKI_LIMIT
1380 };
1381 let hits = if query.chars().count() < MIN_FULLTEXT_QUERY {
1382 sessionwiki::index::search_like(&connection, query, search_limit, None, None)
1383 } else {
1384 sessionwiki::index::search(&connection, query, search_limit, None, None)
1385 }
1386 .context("search the SessionWiki index")?;
1387 let mut rows = Vec::with_capacity(hits.len().min(limit));
1389 for hit in hits {
1390 if rows.len() >= limit {
1391 break;
1392 }
1393 if !include_subagents && !is_main_session(&hit.row) {
1394 continue;
1395 }
1396 if !include_subagents
1403 && !matches!(hit.role.as_str(), "user" | "assistant")
1404 && !conversation_matches(&connection, &hit.row, query)?
1405 {
1406 continue;
1407 }
1408 rows.push(wiki_row(hit.row, Some(hit.snippet), live));
1409 }
1410 if rows.len() < limit {
1414 let found: BTreeSet<String> = rows.iter().map(|row| row.id.clone()).collect();
1415 for row in named_like(&connection, query, include_subagents)? {
1416 if rows.len() >= limit {
1417 break;
1418 }
1419 if found.contains(&row.session_id) {
1420 continue;
1421 }
1422 rows.push(wiki_row(row, None, live));
1423 }
1424 }
1425 fill_session_tags(&connection, &mut rows)?;
1426 Ok(rows)
1427}
1428
1429const TEXT_SEARCH_MESSAGE_LIMIT: i64 = 20_000;
1433
1434pub fn session_text_matches(query: &str, live: &BTreeSet<String>) -> Result<Vec<SessionTextMatch>> {
1439 let query = query.trim();
1440 if query.is_empty() || !index_is_writable() {
1441 return Ok(Vec::new());
1442 }
1443 let connection = open_readonly()?;
1444 text_matches_in(&connection, query, live)
1445}
1446
1447fn text_matches_in(
1448 connection: &rusqlite::Connection,
1449 query: &str,
1450 live: &BTreeSet<String>,
1451) -> Result<Vec<SessionTextMatch>> {
1452 let mut statement;
1453 let rows = if query.chars().count() < MIN_FULLTEXT_QUERY {
1454 let pattern = format!(
1456 "%{}%",
1457 sessionwiki::util::nfc(query)
1458 .replace('\\', "\\\\")
1459 .replace('%', "\\%")
1460 .replace('_', "\\_")
1461 );
1462 statement = connection.prepare(
1463 "SELECT f.path, m.role
1464 FROM messages m JOIN files f ON f.session_id = m.session_id
1465 WHERE f.tool = ?1 AND m.role IN ('user', 'assistant')
1466 AND m.text LIKE ?2 ESCAPE '\\'
1467 ORDER BY m.id DESC LIMIT ?3",
1468 )?;
1469 statement
1470 .query_map(
1471 rusqlite::params![TOOL, pattern, TEXT_SEARCH_MESSAGE_LIMIT],
1472 |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
1473 )?
1474 .collect::<rusqlite::Result<Vec<_>>>()
1475 } else {
1476 let phrase = format!("\"{}\"", sessionwiki::util::nfc(query).replace('"', "\"\""));
1477 statement = connection.prepare(
1478 "SELECT f.path, m.role
1479 FROM (SELECT rowid AS mid FROM msgs WHERE msgs MATCH ?2 LIMIT ?3) x
1480 JOIN messages m ON m.id = x.mid
1481 JOIN files f ON f.session_id = m.session_id
1482 WHERE f.tool = ?1 AND m.role IN ('user', 'assistant')",
1483 )?;
1484 statement
1485 .query_map(
1486 rusqlite::params![TOOL, phrase, TEXT_SEARCH_MESSAGE_LIMIT],
1487 |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
1488 )?
1489 .collect::<rusqlite::Result<Vec<_>>>()
1490 }
1491 .context("search the indexed messages")?;
1492 let mut found = BTreeMap::<String, SessionTextMatchKind>::new();
1493 for (path, role) in rows {
1494 let session_id = path.rsplit('/').next().unwrap_or_default();
1496 if !live.contains(session_id) {
1497 continue;
1498 }
1499 let kind = if role == "user" {
1500 SessionTextMatchKind::User
1501 } else {
1502 SessionTextMatchKind::Agent
1503 };
1504 let entry = found.entry(session_id.to_owned()).or_insert(kind);
1505 *entry = (*entry).min(kind);
1506 }
1507 Ok(found
1508 .into_iter()
1509 .map(|(session_id, kind)| SessionTextMatch { session_id, kind })
1510 .collect())
1511}
1512
1513fn fill_session_tags(connection: &rusqlite::Connection, rows: &mut [WikiRow]) -> Result<()> {
1519 let ids: Vec<&str> = rows
1520 .iter()
1521 .filter(|row| row.tool == TOOL)
1522 .map(|row| row.id.as_str())
1523 .collect();
1524 let found = tags::read(connection, &ids).context("read the indexed session metadata")?;
1525 for row in rows.iter_mut().filter(|row| row.tool == TOOL) {
1526 let Some(session) = found.get(&row.id) else {
1527 continue;
1528 };
1529 row.target = session.target.clone();
1530 row.profile = session.profile.clone();
1531 row.harness = session.harness.clone();
1532 }
1533 Ok(())
1534}
1535
1536const NAME_SCAN_LIMIT: usize = 2_000;
1542
1543fn named_like(
1545 connection: &rusqlite::Connection,
1546 query: &str,
1547 include_subagents: bool,
1548) -> Result<Vec<sessionwiki::index::SessionRow>> {
1549 let needle = query.to_lowercase();
1550 let sql = format!(
1551 "SELECT session_id, tool, path, project, title, started, msg_count, kind,
1552 archived_at IS NOT NULL
1553 FROM files WHERE {} ORDER BY started DESC LIMIT {NAME_SCAN_LIMIT}",
1554 if include_subagents {
1555 "1=1"
1556 } else {
1557 "kind = 'main'"
1558 }
1559 );
1560 let mut statement = connection
1561 .prepare(&sql)
1562 .context("prepare recent SessionWiki metadata scan")?;
1563 let rows = statement
1564 .query_map([], |row| {
1565 Ok(sessionwiki::index::SessionRow {
1566 session_id: row.get(0)?,
1567 tool: row.get(1)?,
1568 path: row.get(2)?,
1569 project: row.get(3)?,
1570 title: row.get(4)?,
1571 started: row.get(5)?,
1572 msg_count: row.get(6)?,
1573 kind: row.get(7)?,
1574 preview: None,
1575 summary: None,
1576 tags: None,
1577 archived: row.get(8)?,
1578 account: None,
1579 })
1580 })
1581 .context("list recent SessionWiki metadata")?
1582 .collect::<rusqlite::Result<Vec<_>>>()
1583 .context("read recent SessionWiki metadata")?;
1584 Ok(rows
1585 .into_iter()
1586 .filter(|row| {
1587 row.title.to_lowercase().contains(&needle)
1588 || row.project.to_lowercase().contains(&needle)
1589 })
1590 .collect())
1591}
1592
1593fn conversation_matches(
1596 connection: &rusqlite::Connection,
1597 row: &sessionwiki::index::SessionRow,
1598 query: &str,
1599) -> Result<bool> {
1600 let session = sessionwiki::index::session_from_index(connection, row)
1601 .context("read an indexed session")?;
1602 Ok(!hit_transcript(&session, query, 0, 1).blocks.is_empty())
1603}
1604
1605fn is_main_session(row: &sessionwiki::index::SessionRow) -> bool {
1608 row.kind == "main"
1609}
1610
1611pub fn brief(id: &str, max_chars: usize) -> Result<Option<String>> {
1613 if !index_is_writable() {
1614 return Ok(None);
1615 }
1616 let connection = open_readonly()?;
1617 let Some(row) = row_by_id(&connection, id)? else {
1618 return Ok(None);
1619 };
1620 let session = sessionwiki::index::session_from_index(&connection, &row)
1621 .context("read an indexed session")?;
1622 Ok(Some(sessionwiki::commands::brief_markdown(
1623 &session, max_chars, true,
1624 )))
1625}
1626
1627pub fn transcript_hits(
1635 id: &str,
1636 query: &str,
1637 context_messages: usize,
1638 per_message_chars: usize,
1639) -> Result<Option<WikiHitTranscript>> {
1640 if !index_is_writable() {
1641 return Ok(None);
1642 }
1643 let connection = open_readonly()?;
1644 let Some(row) = row_by_id(&connection, id)? else {
1645 return Ok(None);
1646 };
1647 let session = sessionwiki::index::session_from_index(&connection, &row)
1648 .context("read an indexed session")?;
1649 Ok(Some(hit_transcript(
1650 &session,
1651 query,
1652 context_messages,
1653 per_message_chars,
1654 )))
1655}
1656
1657fn hit_transcript(
1667 session: &Session,
1668 query: &str,
1669 context_messages: usize,
1670 per_message_chars: usize,
1671) -> WikiHitTranscript {
1672 let found = sessionwiki::grep::grep_session(
1673 session,
1674 query,
1675 &sessionwiki::grep::GrepOpts {
1676 context_messages,
1677 chars: per_message_chars,
1678 max_matches: None,
1679 anchor_roles: vec![Role::User, Role::Assistant],
1680 },
1681 );
1682 WikiHitTranscript {
1683 blocks: found
1684 .hits
1685 .into_iter()
1686 .map(|hit| WikiHitBlock {
1687 role: role_name(hit.role).to_owned(),
1688 text: hit.text,
1689 hits: hit.matches,
1690 omitted_before: hit.omitted_before,
1691 truncated: hit.truncated,
1692 })
1693 .collect(),
1694 omitted_after: found.omitted_after,
1695 }
1696}
1697
1698fn role_name(role: Role) -> &'static str {
1699 match role {
1700 Role::User => "user",
1701 Role::Assistant => "assistant",
1702 Role::Tool => "tool",
1703 }
1704}
1705
1706pub struct ArchivedSession {
1710 pub title: String,
1711 pub project_directory: Option<PathBuf>,
1714 pub snapshot: mj_core::archive::CanonicalSessionSnapshot,
1715}
1716
1717pub fn archived_session(id: &str) -> Result<Option<ArchivedSession>> {
1719 if !index_is_writable() {
1720 return Ok(None);
1721 }
1722 let connection = open_readonly()?;
1723 let Some(row) = row_by_id(&connection, id)? else {
1724 return Ok(None);
1725 };
1726 let session = sessionwiki::index::session_from_index(&connection, &row)
1727 .context("read an indexed session")?;
1728 let snapshot = snapshot_of(&session)?;
1729 Ok(Some(ArchivedSession {
1730 title: session.title.clone(),
1731 project_directory: project_directory_of(&session.project),
1732 snapshot,
1733 }))
1734}
1735
1736pub fn sessions_ready_to_archive(
1752 sessions: &mj_core::snapshot_map::SnapshotMap<String, SessionRecord>,
1753 subagents: &mj_core::snapshot_map::SnapshotMap<String, mj_core::subagent::SubagentRecord>,
1754 now: DateTime<Utc>,
1755 older_than_days: u32,
1756) -> Vec<String> {
1757 let cutoff = now - chrono::Duration::days(i64::from(older_than_days));
1758 let aged = |session_id: &String| {
1759 sessions.get(session_id).is_some_and(|record| {
1760 record.state == mj_core::state::SessionState::Stopped
1761 && parse_time(&record.updated_at).is_some_and(|updated| updated <= cutoff)
1762 && (record.managed_worktree.as_ref().is_some_and(|checkout| {
1763 checkout.kind == mj_core::state::ManagedCheckoutKind::Worktree
1764 }) || (record.managed_worktree.is_none() && record.project_directory.is_some())
1765 || record
1766 .checkpoint
1767 .as_ref()
1768 .zip(record.publication.as_ref())
1769 .is_some_and(|(checkpoint, publication)| {
1770 publication.checkpoint_sha256 == checkpoint.sha256
1771 && publication.state == mj_core::state::PublicationState::Published
1772 && !publication.dirty
1773 && !publication.stashed
1774 }))
1775 })
1776 };
1777 let selected: BTreeSet<String> = sessions
1778 .keys()
1779 .filter(|session_id| aged(session_id))
1780 .filter(|session_id| {
1781 subagents
1782 .values()
1783 .filter(|child| &&child.parent_session_id == session_id)
1784 .filter(|child| sessions.contains_key(&child.child_session_id))
1786 .all(|child| aged(&child.child_session_id))
1787 })
1788 .cloned()
1789 .collect();
1790 let mut ordered: Vec<String> = selected.iter().cloned().collect();
1791 ordered.sort_by_key(|session_id| std::cmp::Reverse(ancestor_depth(session_id, subagents)));
1792 ordered
1793}
1794
1795fn ancestor_depth(
1799 session_id: &str,
1800 subagents: &mj_core::snapshot_map::SnapshotMap<String, mj_core::subagent::SubagentRecord>,
1801) -> usize {
1802 let mut depth = 0;
1803 let mut current = session_id;
1804 while let Some(parent) = subagents
1806 .get(current)
1807 .map(|child| child.parent_session_id.as_str())
1808 {
1809 depth += 1;
1810 if depth > subagents.len() {
1811 break;
1812 }
1813 current = parent;
1814 }
1815 depth
1816}
1817
1818pub use mj_core::state::ArchiveSpacePreview;
1824
1825pub fn archive_space_preview(older_than_days: Option<u32>) -> Result<ArchiveSpacePreview> {
1837 let controller =
1838 Controller::load().context("load the session records to size their storage")?;
1839 Ok(archive_space_over(
1840 &mj_core::config::sessions_dir(),
1841 &controller.state.sessions,
1842 &controller.state.subagents,
1843 Utc::now(),
1844 older_than_days,
1845 ))
1846}
1847
1848fn archive_space_over(
1851 sessions_root: &Path,
1852 sessions: &mj_core::snapshot_map::SnapshotMap<String, SessionRecord>,
1853 subagents: &mj_core::snapshot_map::SnapshotMap<String, mj_core::subagent::SubagentRecord>,
1854 now: DateTime<Utc>,
1855 older_than_days: Option<u32>,
1856) -> ArchiveSpacePreview {
1857 let mut preview = ArchiveSpacePreview {
1858 sessions: sessions.len(),
1859 bytes: sessions
1860 .iter()
1861 .map(|(session_id, record)| session_bytes(sessions_root, session_id, record))
1862 .sum(),
1863 reclaimable_sessions: 0,
1864 reclaimable_bytes: 0,
1865 };
1866 if let Some(days) = older_than_days {
1867 let aged = sessions_ready_to_archive(sessions, subagents, now, days);
1868 preview.reclaimable_sessions = aged.len();
1869 preview.reclaimable_bytes = aged
1870 .iter()
1871 .filter_map(|session_id| {
1872 sessions
1873 .get(session_id)
1874 .map(|record| session_bytes(sessions_root, session_id, record))
1875 })
1876 .sum();
1877 }
1878 preview
1879}
1880
1881fn session_bytes(sessions_root: &Path, session_id: &str, record: &SessionRecord) -> u64 {
1884 let checkpoint = record
1885 .checkpoint
1886 .as_ref()
1887 .and_then(|checkpoint| std::fs::metadata(&checkpoint.archive_path).ok())
1888 .filter(|metadata| metadata.is_file())
1889 .map(|metadata| metadata.len())
1890 .unwrap_or(0);
1891 let attachments = sessions_root
1892 .join(session_id)
1893 .join(mj_core::attachment::ATTACHMENT_DIR);
1894 let attachments = crate::import::claude::directory_size(&attachments).unwrap_or(0);
1895 checkpoint.saturating_add(attachments)
1896}
1897
1898pub fn indexed_with_messages(session_ids: &[String]) -> Result<BTreeSet<String>> {
1905 if !index_is_writable() {
1906 return Ok(BTreeSet::new());
1909 }
1910 let connection = open_readonly()?;
1911 let sessions_dir = mj_core::config::sessions_dir();
1912 let mut indexed = BTreeSet::new();
1913 for session_id in session_ids {
1914 let key = format!("{}/{session_id}", sessions_dir.display());
1915 let rows = sessionwiki::index::resolve(&connection, session_id)
1916 .context("look up a stopped session in the SessionWiki index")?;
1917 if rows
1918 .iter()
1919 .any(|row| row.tool == TOOL && row.path == key && row.msg_count > 0 && !row.archived)
1920 {
1921 indexed.insert(session_id.clone());
1922 }
1923 }
1924 Ok(indexed)
1925}
1926
1927fn open_readonly() -> Result<rusqlite::Connection> {
1928 sessionwiki::index::open_readonly().context("open the SessionWiki index")
1929}
1930
1931fn row_by_id(
1934 connection: &rusqlite::Connection,
1935 id: &str,
1936) -> Result<Option<sessionwiki::index::SessionRow>> {
1937 Ok(sessionwiki::index::resolve(connection, id)
1938 .context("look up an indexed session")?
1939 .into_iter()
1940 .find(|row| row.session_id == id))
1941}
1942
1943fn wiki_row(
1944 row: sessionwiki::index::SessionRow,
1945 snippet: Option<String>,
1946 live: &BTreeSet<String>,
1947) -> WikiRow {
1948 let hel_session_id = (row.tool == TOOL)
1951 .then(|| row.path.rsplit('/').next().unwrap_or_default().to_owned())
1952 .filter(|session_id| live.contains(session_id));
1953 let native_id = sessionwiki::index::native_id_of(&row.path);
1954 WikiRow {
1955 id: row.session_id,
1956 tool: row.tool,
1957 project: row.project,
1958 title: row.title,
1959 started: row.started,
1960 msgs: row.msg_count,
1961 preview: row.preview,
1962 archived: row.archived,
1963 native_id,
1964 snippet,
1965 hel_session_id,
1966 target: None,
1969 profile: None,
1970 harness: None,
1971 }
1972}
1973
1974fn project_directory_of(project: &str) -> Option<PathBuf> {
1982 if project.trim().is_empty() {
1983 return None;
1984 }
1985 let path = PathBuf::from(project);
1986 let repository = path
1987 .ancestors()
1988 .find(|ancestor| ancestor.file_name().is_some_and(|name| name == ".mj"))
1989 .and_then(std::path::Path::parent)
1990 .map(std::path::Path::to_path_buf)
1991 .unwrap_or(path);
1992 repository.is_dir().then_some(repository)
1993}
1994
1995fn snapshot_of(
2002 session: &sessionwiki::model::Session,
2003) -> Result<mj_core::archive::CanonicalSessionSnapshot> {
2004 use mj_core::archive::{
2005 CanonicalExecutionState, CanonicalSessionSnapshot, CanonicalSessionState,
2006 CanonicalTranscriptBody, CanonicalTranscriptItem,
2007 };
2008
2009 let started_ms = session
2010 .started
2011 .map(|time| time.timestamp_millis())
2012 .unwrap_or_default();
2013 let mut transcript: Vec<CanonicalTranscriptItem> = Vec::new();
2014 for message in &session.messages {
2015 let text = message.text.trim();
2016 if text.is_empty() {
2017 continue;
2018 }
2019 if transcript.is_empty() && message.role != Role::User {
2022 continue;
2023 }
2024 let position = transcript.len() as u64 + 1;
2025 let body = match message.role {
2026 Role::User => CanonicalTranscriptBody::User {
2027 content: vec![serde_json::json!({"type": "text", "text": text})],
2028 },
2029 Role::Assistant => CanonicalTranscriptBody::Agent {
2030 chunks: vec![serde_json::json!({
2031 "content": {"type": "text", "text": text}
2032 })],
2033 streaming: false,
2034 },
2035 Role::Tool => {
2037 let (call, terminal_outputs) = mj_transcript::summary::indexed_tool_call(
2038 text,
2039 &format!("wiki-tool-{position}"),
2040 );
2041 CanonicalTranscriptBody::Tool {
2042 call,
2043 terminal_outputs,
2044 terminal_refs: Vec::new(),
2045 presentation: None,
2046 }
2047 }
2048 };
2049 let created_at_ms = message
2050 .ts
2051 .map(|time| time.timestamp_millis())
2052 .unwrap_or(started_ms);
2053 transcript.push(CanonicalTranscriptItem {
2054 stable_id: format!("wiki-{position}"),
2055 position,
2056 latest_content_event_ordinal: matches!(body, CanonicalTranscriptBody::Agent { .. })
2059 .then_some(position),
2060 created_at_ms,
2061 last_changed_at_ms: created_at_ms,
2062 body,
2063 });
2064 }
2065 anyhow::ensure!(
2066 !transcript.is_empty(),
2067 "the archived session has no prompt to restore from"
2068 );
2069
2070 let event_frontier = transcript.len() as u64;
2071 let last_activity_at_ms = transcript.last().map(|item| item.last_changed_at_ms);
2072 Ok(CanonicalSessionSnapshot {
2073 command_ledger: None,
2074 assessment_state: None,
2075 event_frontier,
2076 event_frontier_digest: {
2080 use sha2::Digest;
2081 mj_core::hex::lower_hex(sha2::Sha256::digest(
2082 format!("sessionwiki:{}", session.id).as_bytes(),
2083 ))
2084 },
2085 session: CanonicalSessionState {
2086 execution: CanonicalExecutionState::Idle,
2087 last_activity_at_ms,
2088 session_title: Some(session.title.clone()).filter(|title| !title.trim().is_empty()),
2089 configuration: Default::default(),
2090 },
2091 transcript,
2092 queued_prompts: Vec::new(),
2093 })
2094}
2095
2096#[derive(Debug, Clone, PartialEq, Eq)]
2106pub enum WikiContinuation {
2107 Resume { session_id: String },
2109 Restore { wiki_id: String },
2112 Import {
2114 harness: HarnessKind,
2115 native_session_id: String,
2116 },
2117}
2118
2119pub fn wiki_continuation(
2125 wiki_id: &str,
2126 tool: &str,
2127 path: &Path,
2128 has_record: bool,
2129) -> Result<WikiContinuation> {
2130 if tool == TOOL {
2131 return Ok(match has_record {
2134 true => WikiContinuation::Resume {
2135 session_id: wiki_id.to_owned(),
2136 },
2137 false => WikiContinuation::Restore {
2138 wiki_id: wiki_id.to_owned(),
2139 },
2140 });
2141 }
2142 let harness = harness_adapters::harness_for_tool(tool)
2143 .with_context(|| format!("Mjolnir cannot continue a {tool} session"))?;
2144 let native_session_id = crate::import::native_session_id_from_path(harness, path)
2145 .with_context(|| {
2146 format!(
2147 "no {tool} session id in the indexed path {}",
2148 path.display()
2149 )
2150 })?;
2151 Ok(WikiContinuation::Import {
2152 harness,
2153 native_session_id,
2154 })
2155}
2156
2157pub fn wiki_session(
2162 wiki_id: &str,
2163 known_sessions: &BTreeSet<String>,
2164) -> Result<Option<WikiSessionInfo>> {
2165 if !index_is_writable() {
2166 return Ok(None);
2167 }
2168 let connection = open_readonly()?;
2169 let Some(row) = row_by_id(&connection, wiki_id)? else {
2170 return Ok(None);
2171 };
2172 let is_mjolnir = row.tool == TOOL;
2173 let mjolnir_session_id = is_mjolnir.then(|| row.session_id.clone());
2174 let has_record = mjolnir_session_id
2175 .as_deref()
2176 .is_some_and(|session_id| known_sessions.contains(session_id));
2177 let status = match (is_mjolnir, has_record) {
2178 (false, _) => WikiSessionStatus::Native,
2179 (true, true) => WikiSessionStatus::Mine,
2180 (true, false) => WikiSessionStatus::Archived,
2181 };
2182 let tags = match is_mjolnir {
2183 true => tags::read(&connection, &[row.session_id.as_str()])
2184 .context("read the indexed session metadata")?
2185 .remove(&row.session_id)
2186 .unwrap_or_default(),
2187 false => tags::MjTags::default(),
2188 };
2189 let nothing_to_restore = status == WikiSessionStatus::Archived
2192 && !has_prompt(
2193 &sessionwiki::index::session_from_index(&connection, &row)
2194 .context("read an indexed session")?,
2195 );
2196 let harness = tags
2197 .harness
2198 .as_deref()
2199 .and_then(|id| id.parse::<HarnessKind>().ok())
2200 .or_else(|| {
2201 (!is_mjolnir)
2202 .then(|| harness_adapters::harness_for_tool(&row.tool))
2203 .flatten()
2204 });
2205 Ok(Some(WikiSessionInfo {
2206 wiki_id: row.session_id,
2207 tool: row.tool,
2208 path: PathBuf::from(row.path),
2209 status,
2210 mjolnir_session_id,
2211 profile_id: tags.profile,
2212 target_template_id: tags.target,
2213 harness,
2214 title: row.title,
2215 project: row.project,
2216 nothing_to_restore,
2217 }))
2218}
2219
2220fn has_prompt(session: &sessionwiki::model::Session) -> bool {
2223 session
2224 .messages
2225 .iter()
2226 .any(|message| message.role == Role::User && !message.text.trim().is_empty())
2227}
2228
2229#[cfg(test)]
2230mod tests {
2231 use std::collections::BTreeMap;
2232 use std::path::Path;
2233
2234 mod continuation {
2237 use super::super::{WikiContinuation, wiki_continuation};
2238 use mj_core::config::HarnessKind;
2239 use std::path::Path;
2240
2241 #[test]
2242 fn a_mjolnir_row_with_a_record_is_resumed_and_one_without_is_restored() {
2243 let path = Path::new("/home/user/.local/share/mj/sessions/session-7");
2244 assert_eq!(
2245 wiki_continuation("session-7", "mjolnir", path, true).unwrap(),
2246 WikiContinuation::Resume {
2247 session_id: "session-7".to_owned(),
2248 }
2249 );
2250 assert_eq!(
2251 wiki_continuation("session-7", "mjolnir", path, false).unwrap(),
2252 WikiContinuation::Restore {
2253 wiki_id: "session-7".to_owned(),
2254 }
2255 );
2256 }
2257
2258 #[test]
2259 fn a_claude_code_row_is_imported_with_the_uuid_from_its_path() {
2260 let path = Path::new(
2261 "/home/user/.claude/projects/-home-user-app/7f3a1c20-0b11-4a55-9e0d-2c8a5d6f1b44.jsonl",
2262 );
2263 assert_eq!(
2264 wiki_continuation("abc123", "claude-code", path, false).unwrap(),
2265 WikiContinuation::Import {
2266 harness: HarnessKind::Claude,
2267 native_session_id: "7f3a1c20-0b11-4a55-9e0d-2c8a5d6f1b44".to_owned(),
2268 }
2269 );
2270 }
2271
2272 #[test]
2275 fn a_codex_row_is_imported_with_the_uuid_from_its_rollout_name() {
2276 let path = Path::new(
2277 "/home/user/.codex/sessions/2026/09/18/rollout-2026-09-18T09-15-00-7f3a1c20-0b11-4a55-9e0d-2c8a5d6f1b44.jsonl",
2278 );
2279 assert_eq!(
2280 wiki_continuation("abc123", "codex", path, false).unwrap(),
2281 WikiContinuation::Import {
2282 harness: HarnessKind::Codex,
2283 native_session_id: "7f3a1c20-0b11-4a55-9e0d-2c8a5d6f1b44".to_owned(),
2284 }
2285 );
2286 }
2287
2288 #[test]
2289 fn an_unknown_tool_is_an_error_that_names_it() {
2290 let error = wiki_continuation("abc123", "opencode", Path::new("/tmp/s.jsonl"), false)
2291 .unwrap_err();
2292 assert!(
2293 format!("{error:#}").contains("opencode"),
2294 "the error has to name the tool: {error:#}"
2295 );
2296 }
2297 }
2298
2299 use mj_checkpoint::archive::{
2300 ArchiveInput, BundleManifest, CanonicalExecutionState, CanonicalSessionSnapshot,
2301 CanonicalSessionState, CanonicalTranscriptBody, CanonicalTranscriptItem, SessionManifest,
2302 TargetManifest, write_archive_atomic,
2303 };
2304
2305 use super::*;
2306
2307 fn item(position: u64, body: CanonicalTranscriptBody) -> CanonicalTranscriptItem {
2308 let streamed = matches!(body, CanonicalTranscriptBody::Agent { .. });
2311 CanonicalTranscriptItem {
2312 stable_id: format!("item-{position}"),
2313 position,
2314 latest_content_event_ordinal: streamed.then_some(position),
2315 created_at_ms: 1_700_000_000_000 + i64::try_from(position).unwrap(),
2316 last_changed_at_ms: 1_700_000_000_000 + i64::try_from(position).unwrap(),
2317 body,
2318 }
2319 }
2320
2321 fn write_archive(directory: &Path, session_id: &str, frontier: u64) {
2324 let path = directory.join(format!(
2325 "{session_id}-{frontier}-archive-{}.hel.zip",
2326 "0".repeat(32)
2327 ));
2328 write_archive_atomic(
2329 &path,
2330 &ArchiveInput {
2331 session: SessionManifest {
2332 id: session_id.into(),
2333 title: "indexed session".into(),
2334 harness_kind: mj_core::config::HarnessKind::Codex,
2335 profile_id: "codex".into(),
2336 native_session_id: "native-session".into(),
2337 created_at: "2026-09-01T00:00:00Z".into(),
2338 checkpointed_at: "2026-09-01T01:00:00Z".into(),
2339 hel_version: "test".into(),
2340 relay_version: "test".into(),
2341 adapter_version: "test".into(),
2342 },
2343 target: TargetManifest {
2344 template_id: "local".into(),
2345 target_kind: "local-bare".into(),
2346 details: BTreeMap::new(),
2347 },
2348 bundle: BundleManifest {
2349 id: "project".into(),
2350 primary_repository: "project".into(),
2351 },
2352 canonical_session: CanonicalSessionSnapshot {
2353 command_ledger: None,
2354 assessment_state: None,
2355 event_frontier: 4,
2356 event_frontier_digest: "a".repeat(64),
2357 session: CanonicalSessionState {
2358 execution: CanonicalExecutionState::Idle,
2359 last_activity_at_ms: Some(1_700_000_000_004),
2360 session_title: Some("snapshot title".into()),
2361 configuration: Default::default(),
2362 },
2363 transcript: vec![
2364 item(
2365 1,
2366 CanonicalTranscriptBody::User {
2367 content: vec![serde_json::json!({
2368 "type": "text",
2369 "text": "index this session"
2370 })],
2371 },
2372 ),
2373 item(
2374 2,
2375 CanonicalTranscriptBody::Thought {
2376 chunks: vec![serde_json::json!({
2377 "content": {"type": "text", "text": "pondering"}
2378 })],
2379 streaming: false,
2380 },
2381 ),
2382 item(
2383 3,
2384 CanonicalTranscriptBody::Tool {
2385 call: serde_json::json!({
2386 "toolCallId": "call-1",
2387 "title": "Edit config.toml",
2388 "kind": "edit",
2389 "status": "completed",
2390 "locations": [{"path": "/old/container/config.toml"}]
2391 }),
2392 terminal_outputs: Vec::new(),
2393 terminal_refs: Vec::new(),
2394 presentation: None,
2395 },
2396 ),
2397 item(
2398 4,
2399 CanonicalTranscriptBody::Agent {
2400 chunks: vec![serde_json::json!({
2401 "content": {"type": "text", "text": "done"}
2402 })],
2403 streaming: false,
2404 },
2405 ),
2406 ],
2407 queued_prompts: Vec::new(),
2408 },
2409 native_artifacts: Vec::new(),
2410 repositories: Vec::new(),
2411 },
2412 )
2413 .unwrap();
2414 }
2415
2416 fn adapter(directory: &Path, session_id: &str) -> MjolnirAdapter {
2417 adapter_with_live(directory, session_id, BTreeMap::new())
2418 }
2419
2420 fn adapter_with_live(
2421 directory: &Path,
2422 session_id: &str,
2423 live: BTreeMap<String, i64>,
2424 ) -> MjolnirAdapter {
2425 let record = SessionRecord {
2426 project: None,
2427 id: session_id.into(),
2428 ..record_template()
2429 };
2430 MjolnirAdapter {
2431 sessions_dir: directory.to_path_buf(),
2432 sessions: std::sync::Mutex::new(Sessions {
2433 records: [(session_id.to_owned(), record)].into_iter().collect(),
2434 subagent_ids: BTreeSet::new(),
2435 live,
2436 }),
2437 reload: false,
2438 }
2439 }
2440
2441 fn record_template() -> SessionRecord {
2442 SessionRecord {
2443 project: None,
2444 target_runtime: None,
2445 launch_base: None,
2446 launch_branch: None,
2447 checkout: None,
2448 publication: None,
2449 build_cache: None,
2450 container_workspace: None,
2451 subagents: None,
2452 create_managed_worktree: None,
2453 workspace_id: mj_core::workspace::DEFAULT_WORKSPACE_ID.to_owned(),
2454 archived: false,
2455 container_cpus: None,
2456 container_memory: None,
2457 id: "0123456789abcdef0123456789abcdef".into(),
2458 title: "indexed session".into(),
2459 harness_kind: mj_core::config::HarnessKind::Codex,
2460 last_profile: "codex".into(),
2461 bundle_id: "project".into(),
2462 project_directory: Some(PathBuf::from("/home/dev/project")),
2463 managed_worktree: None,
2464 target_template_id: "local-bare".into(),
2465 resource_allocation: None,
2466 additional_mounts: Vec::new(),
2467 state: mj_core::state::SessionState::Stopped,
2468 target: None,
2469 native_session_id: Some("native-session".into()),
2470 acp_session_title: Some("the harness title".into()),
2471 session_title_override: None,
2472 created_at: "2026-09-01T00:00:00Z".into(),
2473 updated_at: "2026-09-01T01:00:00Z".into(),
2474 viewed_through_event_ordinal: 0,
2475 draft_input: String::new(),
2476 last_error: None,
2477 last_checkpoint_error: None,
2478 checkpoint: None,
2479 }
2480 }
2481
2482 #[test]
2483 fn the_newest_checkpoint_of_each_session_is_one_indexed_key() {
2484 let directory = tempfile::tempdir().unwrap();
2485 let session_id = "0123456789abcdef0123456789abcdef";
2486 write_archive(directory.path(), session_id, 1);
2487 write_archive(directory.path(), session_id, 7);
2488 let adapter = adapter(directory.path(), session_id);
2489
2490 let store = adapter.store().expect("the adapter is a shared store");
2491 let key = format!("{}/{session_id}", directory.path().display());
2492 assert_eq!(
2493 store
2494 .keys
2495 .iter()
2496 .map(|(key, _)| key.as_str())
2497 .collect::<Vec<_>>(),
2498 vec![key.as_str()]
2499 );
2500 assert!(!store.had_error);
2501 assert_eq!(store.files.len(), 1);
2502 assert!(
2503 store.files[0]
2504 .file_name()
2505 .unwrap()
2506 .to_str()
2507 .unwrap()
2508 .contains("-7-archive-"),
2509 "the newest checkpoint is the one indexed: {:?}",
2510 store.files[0]
2511 );
2512 assert_eq!(
2513 adapter.reconcile_scope(),
2514 Some(format!("{}/", directory.path().display()))
2515 );
2516
2517 let session = adapter.parse_key(&key).unwrap();
2518 assert_eq!(session.id, session_id);
2519 assert_eq!(session.tool, "mjolnir");
2520 assert_eq!(session.path, PathBuf::from(&key));
2521 assert_eq!(session.project, "/home/dev/project");
2522 assert_eq!(session.title, "the harness title");
2523 assert!(!session.subagent);
2524 assert_eq!(
2525 session.messages.iter().map(|m| m.role).collect::<Vec<_>>(),
2526 vec![Role::User, Role::Tool, Role::Assistant]
2527 );
2528 assert_eq!(session.messages[0].text, "index this session");
2529 let tool: serde_json::Value = serde_json::from_str(&session.messages[1].text).unwrap();
2530 assert_eq!(tool["name"], "Edit");
2531 assert_eq!(tool["call"]["title"], "Edit config.toml");
2532 assert_eq!(session.messages[2].text, "done");
2533 assert_eq!(session.touched, vec!["/old/container/config.toml"]);
2534 }
2535
2536 #[test]
2537 fn provenance_backfill_repairs_an_unchanged_checkpoint_without_rebuilding_the_index() {
2538 let _held = tags::testing::lock();
2539 let (_index_dir, mut connection) = tags::testing::isolated_index();
2540 let directory = tempfile::tempdir().unwrap();
2541 write_archive(directory.path(), "old-session", 4);
2542 let source = adapter(directory.path(), "old-session");
2543 let key = source.key_for("old-session");
2544 tags::testing::index_row(&connection, "old-session", "mjolnir");
2545 connection
2546 .execute(
2547 "UPDATE files SET path = ?1 WHERE session_id = 'old-session'",
2548 [&key],
2549 )
2550 .unwrap();
2551 provenance::backfill(&mut connection, &source).unwrap();
2552 assert_eq!(
2553 sessionwiki::index::files_for(&connection, "old-session").unwrap(),
2554 vec!["/old/container/config.toml"]
2555 );
2556 provenance::backfill(&mut connection, &source).unwrap();
2557 assert_eq!(
2558 sessionwiki::index::sessions_for_file(&connection, "config.toml", 20)
2559 .unwrap()
2560 .len(),
2561 1
2562 );
2563 }
2564
2565 fn projection(session_id: &str) -> mj_core::state::MaterializedSession {
2566 use mj_core::transcript::{TranscriptBody, TranscriptItem};
2567 let mut projected = mj_core::state::MaterializedSession::empty(session_id);
2568 let mut push = |position: u64, body: TranscriptBody| {
2569 let streamed = matches!(body, TranscriptBody::Agent { .. });
2570 projected
2571 .transcript
2572 .push(std::sync::Arc::new(TranscriptItem {
2573 stable_id: format!("item-{position}"),
2574 position,
2575 latest_content_event_ordinal: streamed.then_some(position),
2576 created_at_ms: 1_700_000_000_000 + i64::try_from(position).unwrap(),
2577 last_changed_at_ms: 1_700_000_000_000 + i64::try_from(position).unwrap(),
2578 body,
2579 }));
2580 };
2581 push(
2582 1,
2583 TranscriptBody::User {
2584 content: vec![serde_json::json!({"type": "text", "text": "still talking"})],
2585 },
2586 );
2587 push(
2588 2,
2589 TranscriptBody::Thought {
2590 chunks: vec![serde_json::json!({"content": {"type": "text", "text": "hmm"}})],
2591 streaming: false,
2592 },
2593 );
2594 push(
2595 3,
2596 TranscriptBody::Tool {
2597 call: serde_json::json!({"toolCallId": "c1", "title": "Read README.md"}),
2598 terminal_outputs: Vec::new(),
2599 terminal_refs: Vec::new(),
2600 presentation: None,
2601 },
2602 );
2603 push(
2604 4,
2605 TranscriptBody::Agent {
2606 chunks: vec![serde_json::json!({"content": {"type": "text", "text": "reading"}})],
2607 streaming: false,
2608 },
2609 );
2610 projected.session_title = Some("the live title".into());
2611 projected
2612 }
2613
2614 #[test]
2617 fn a_running_session_is_indexed_from_its_stored_transcript() {
2618 let session_id = "0123456789abcdef0123456789abcdef";
2619 let messages = projected_messages(&projection(session_id));
2620 assert_eq!(
2621 messages.iter().map(|m| m.role).collect::<Vec<_>>(),
2622 vec![Role::User, Role::Tool, Role::Assistant]
2623 );
2624 assert_eq!(messages[0].text, "still talking");
2625 assert_eq!(messages[2].text, "reading");
2626 let tool: serde_json::Value = serde_json::from_str(&messages[1].text).unwrap();
2627 assert_eq!(tool["name"], "Read");
2628 assert_eq!(tool["call"]["title"], "Read README.md");
2629 }
2630
2631 #[test]
2636 fn a_running_session_is_listed_with_its_own_change_token() {
2637 let directory = tempfile::tempdir().unwrap();
2638 let running = "0123456789abcdef0123456789abcdef";
2639 let never_checkpointed = "fedcba9876543210fedcba9876543210";
2640 write_archive(directory.path(), running, 3);
2641 let live = adapter_with_live(
2642 directory.path(),
2643 running,
2644 BTreeMap::from([
2645 (running.to_owned(), 1_900_000_000),
2646 (never_checkpointed.to_owned(), 1_900_000_001),
2647 ]),
2648 );
2649
2650 let store = live.store().expect("the adapter is a shared store");
2651 let key_of = |session_id: &str| format!("{}/{session_id}", directory.path().display());
2652 assert_eq!(
2653 store.keys,
2654 vec![
2655 (
2656 key_of(running),
2657 1_900_000_000 * 1024 + i64::from(mj_transcript::summary::SUMMARY_VERSION)
2658 ),
2659 (
2660 key_of(never_checkpointed),
2661 1_900_000_001 * 1024 + i64::from(mj_transcript::summary::SUMMARY_VERSION)
2662 ),
2663 ],
2664 "a live session's own token replaces the checkpoint's"
2665 );
2666
2667 let stopped = adapter(directory.path(), running);
2670 let keys = stopped.store().expect("a shared store").keys;
2671 assert_eq!(keys.len(), 1);
2672 assert_eq!(keys[0].0, key_of(running));
2673 assert_ne!(keys[0].1, 1_900_000_000);
2674 assert_eq!(
2675 stopped.parse_key(&key_of(running)).unwrap().title,
2676 "the harness title",
2677 "a stopped session is parsed from its checkpoint"
2678 );
2679 }
2680
2681 #[test]
2684 fn a_rename_moves_a_session_change_token() {
2685 let directory = tempfile::tempdir().unwrap();
2686 let session_id = "0123456789abcdef0123456789abcdef";
2687 write_archive(directory.path(), session_id, 1);
2688 let adapter = adapter(directory.path(), session_id);
2689 let before = adapter.store().expect("a shared store").keys[0].1;
2690
2691 {
2692 let mut sessions = adapter.sessions.lock().unwrap();
2693 let record = sessions.records.get_mut(session_id).unwrap();
2694 record.session_title_override = Some("the new name".into());
2695 record.updated_at = "2099-01-01T00:00:00Z".into();
2696 }
2697 let after = adapter.store().expect("a shared store").keys[0].1;
2698 assert!(
2699 after > before,
2700 "a renamed session is re-indexed: {before} then {after}"
2701 );
2702 assert_eq!(
2703 adapter
2704 .parse_key(&format!("{}/{session_id}", directory.path().display()))
2705 .unwrap()
2706 .title,
2707 "the new name"
2708 );
2709 }
2710
2711 fn indexed(messages: Vec<(Role, &str)>) -> sessionwiki::model::Session {
2712 Session {
2713 id: "0123456789abcdef0123456789abcdef".into(),
2714 tool: "mjolnir",
2715 path: PathBuf::from("/sessions/0123456789abcdef0123456789abcdef"),
2716 project: "/home/dev/project".into(),
2717 started: DateTime::from_timestamp_millis(1_700_000_000_000),
2718 ended: None,
2719 title: "the archived session".into(),
2720 subagent: false,
2721 messages: messages
2722 .into_iter()
2723 .map(|(role, text)| Message {
2724 role,
2725 text: text.to_owned(),
2726 ts: None,
2727 })
2728 .collect(),
2729 touched: Vec::new(),
2730 edits: Vec::new(),
2731 }
2732 }
2733
2734 #[test]
2737 fn transcript_hits_locates_case_insensitive_matches() {
2738 let session = indexed(vec![
2739 (Role::User, "Make the Tests green"),
2740 (Role::Assistant, "the tests are green now"),
2741 ]);
2742
2743 let found = hit_transcript(&session, "TESTS", 0, 4_000);
2744
2745 assert_eq!(found.blocks.len(), 2, "both messages contain the query");
2746 assert_eq!(found.blocks[0].role, "user");
2747 let (start, end) = found.blocks[0].hits[0];
2748 assert_eq!(&found.blocks[0].text[start..end], "Tests");
2749 let (start, end) = found.blocks[1].hits[0];
2750 assert_eq!(&found.blocks[1].text[start..end], "tests");
2751 assert!(!found.blocks[0].truncated);
2752 assert_eq!(found.omitted_after, 0);
2753 }
2754
2755 #[test]
2758 fn transcript_hits_keeps_context_and_marks_omissions() {
2759 let session = indexed(vec![
2760 (Role::User, "zero"),
2761 (Role::Assistant, "one needle one"),
2762 (Role::Tool, "two"),
2763 (Role::User, "three"),
2764 (Role::Assistant, "four"),
2765 (Role::Tool, "five"),
2766 (Role::User, "six needle six"),
2767 (Role::Assistant, "seven"),
2768 (Role::User, "eight"),
2769 ]);
2770
2771 let found = hit_transcript(&session, "needle", 1, 4_000);
2772
2773 let shown: Vec<(&str, &str, usize)> = found
2774 .blocks
2775 .iter()
2776 .map(|block| {
2777 (
2778 block.role.as_str(),
2779 block.text.as_str(),
2780 block.omitted_before,
2781 )
2782 })
2783 .collect();
2784 assert_eq!(
2785 shown,
2786 vec![
2787 ("user", "zero", 0),
2788 ("assistant", "one needle one", 0),
2789 ("tool", "two", 0),
2790 ("tool", "five", 2),
2791 ("user", "six needle six", 0),
2792 ("assistant", "seven", 0),
2793 ]
2794 );
2795 assert_eq!(found.omitted_after, 1, "the last message is not shown");
2796 assert!(found.blocks[0].hits.is_empty(), "context has no hits");
2797 }
2798
2799 #[test]
2805 fn transcript_hits_never_anchor_on_tool_output() {
2806 let session = indexed(vec![
2807 (Role::User, "make it build"),
2808 (Role::Tool, "cargo build --needle"),
2809 (Role::Assistant, "it builds"),
2810 ]);
2811
2812 let only_in_a_tool = hit_transcript(&session, "needle", 1, 4_000);
2813 assert!(
2814 only_in_a_tool.blocks.is_empty(),
2815 "tool output must not anchor a passage, got {:?}",
2816 only_in_a_tool.blocks
2817 );
2818
2819 let beside_a_match = hit_transcript(&session, "builds", 1, 4_000);
2820 let shown: Vec<(&str, bool)> = beside_a_match
2821 .blocks
2822 .iter()
2823 .map(|block| (block.role.as_str(), !block.hits.is_empty()))
2824 .collect();
2825 assert_eq!(
2826 shown,
2827 vec![("tool", false), ("assistant", true)],
2828 "a tool message is still context around a real match"
2829 );
2830 }
2831
2832 #[test]
2835 fn transcript_hits_window_keeps_the_first_hit() {
2836 let filler = "x".repeat(4_000);
2837 let session = indexed(vec![(Role::User, &format!("{filler} needle {filler}"))]);
2838
2839 let found = hit_transcript(&session, "needle", 0, 100);
2840
2841 let block = &found.blocks[0];
2842 assert!(block.truncated);
2843 assert_eq!(block.text.chars().count(), 100);
2844 assert_eq!(block.hits.len(), 1, "the windowed text keeps its hit");
2845 let (start, end) = block.hits[0];
2846 assert_eq!(&block.text[start..end], "needle");
2847 assert!(
2848 start >= 20,
2849 "the window keeps lead-in before the hit, got {start}"
2850 );
2851 }
2852
2853 #[test]
2857 fn a_restored_snapshot_is_a_valid_transcript_of_the_indexed_session() {
2858 let snapshot = snapshot_of(&indexed(vec![
2859 (Role::User, "make the tests green"),
2860 (Role::Tool, "Read src/lib.rs"),
2861 (Role::Assistant, "they are green now"),
2862 (Role::User, " "),
2863 ]))
2864 .unwrap();
2865
2866 snapshot.validate().expect("the snapshot is well formed");
2867 assert_eq!(snapshot.event_frontier, 3);
2868 assert_eq!(
2869 snapshot.session.session_title.as_deref(),
2870 Some("the archived session")
2871 );
2872 assert!(snapshot.session.last_activity_at_ms.is_some());
2873 let bodies = snapshot
2874 .transcript
2875 .iter()
2876 .map(|item| match &item.body {
2877 mj_core::archive::CanonicalTranscriptBody::User { content } => (
2878 "user",
2879 mj_core::transcript::materialized_content_text(content),
2880 ),
2881 mj_core::archive::CanonicalTranscriptBody::Agent { chunks, .. } => (
2882 "agent",
2883 mj_core::transcript::materialized_chunks_text(chunks),
2884 ),
2885 mj_core::archive::CanonicalTranscriptBody::Tool { call, .. } => (
2886 "tool",
2887 call["title"].as_str().unwrap_or_default().to_owned(),
2888 ),
2889 _ => ("other", String::new()),
2890 })
2891 .collect::<Vec<_>>();
2892 assert_eq!(
2893 bodies,
2894 vec![
2895 ("user", "make the tests green".to_owned()),
2896 ("tool", "Read src/lib.rs".to_owned()),
2897 ("agent", "they are green now".to_owned()),
2898 ],
2899 "the blank message is dropped and every other one keeps its role"
2900 );
2901 }
2902
2903 #[test]
2907 fn messages_before_the_first_prompt_are_dropped() {
2908 let snapshot = snapshot_of(&indexed(vec![
2909 (Role::Assistant, "still working"),
2910 (Role::User, "carry on"),
2911 ]))
2912 .unwrap();
2913 assert_eq!(snapshot.transcript.len(), 1);
2914 assert_eq!(snapshot.transcript[0].position, 1);
2915 snapshot.validate().unwrap();
2916
2917 let error = snapshot_of(&indexed(vec![(Role::Assistant, "nobody asked")])).unwrap_err();
2918 assert!(
2919 error.to_string().contains("no prompt"),
2920 "a session with no prompt cannot be restored: {error}"
2921 );
2922 assert!(!has_prompt(&indexed(vec![(
2924 Role::Assistant,
2925 "nobody asked"
2926 )])));
2927 assert!(has_prompt(&indexed(vec![(Role::User, "carry on")])));
2928 }
2929
2930 fn record(
2931 session_id: &str,
2932 state: mj_core::state::SessionState,
2933 updated_at: &str,
2934 ) -> SessionRecord {
2935 SessionRecord {
2936 project: None,
2937 id: session_id.into(),
2938 state,
2939 updated_at: updated_at.into(),
2940 ..record_template()
2941 }
2942 }
2943
2944 fn child(child_session_id: &str, parent_session_id: &str) -> mj_core::subagent::SubagentRecord {
2945 mj_core::subagent::SubagentRecord {
2946 child_session_id: child_session_id.into(),
2947 parent_session_id: parent_session_id.into(),
2948 task_name: "task".into(),
2949 profile_id: "codex".into(),
2950 model: None,
2951 effort: None,
2952 working_directory: PathBuf::new(),
2953 initial_prompt: "do the thing".into(),
2954 request_key: "key".into(),
2955 created_at: "2026-09-01T00:00:00Z".into(),
2956 noticed_turn: None,
2957 handback_tool: false,
2958 }
2959 }
2960
2961 fn ready(
2962 sessions: Vec<SessionRecord>,
2963 children: Vec<mj_core::subagent::SubagentRecord>,
2964 ) -> Vec<String> {
2965 let now = parse_time("2026-09-10T00:00:00Z").unwrap();
2966 sessions_ready_to_archive(
2967 &sessions
2968 .into_iter()
2969 .map(|record| (record.id.clone(), record))
2970 .collect(),
2971 &children
2972 .into_iter()
2973 .map(|child| (child.child_session_id.clone(), child))
2974 .collect(),
2975 now,
2976 3,
2977 )
2978 }
2979
2980 #[test]
2981 fn aged_clone_requires_clean_published_evidence_for_its_current_checkpoint() {
2982 let id = "0123456789abcdef0123456789abcdef";
2983 let root = PathBuf::from(format!("/srv/project/.mj/clones/{id}"));
2984 let mut session = record(
2985 id,
2986 mj_core::state::SessionState::Stopped,
2987 "2026-09-01T00:00:00Z",
2988 );
2989 session.project_directory = Some(root.clone());
2990 session.managed_worktree = Some(mj_core::state::ManagedWorktree {
2991 kind: mj_core::state::ManagedCheckoutKind::Clone,
2992 source_project_directory: "/srv/project".into(),
2993 source_repository: "/srv/project".into(),
2994 worktree_root: root,
2995 branch: "feature".into(),
2996 target: mj_core::state::ManagedWorktreeTarget::Local,
2997 base_commit: Some("1".repeat(40)),
2998 });
2999 session.checkpoint = Some(mj_core::state::CheckpointMetadata {
3000 archive_path: "sessions/checkpoint.hel.zip".into(),
3001 sha256: "a".repeat(64),
3002 created_at: "2026-09-01T00:00:00Z".into(),
3003 event_frontier: 0,
3004 });
3005 assert!(ready(vec![session.clone()], vec![]).is_empty());
3006 session.publication = Some(mj_core::state::PublicationAssessment {
3007 checkpoint_sha256: "a".repeat(64),
3008 state: mj_core::state::PublicationState::Published,
3009 dirty: false,
3010 stashed: false,
3011 saved_commits: vec!["2".repeat(40)],
3012 destinations: vec!["https://example.test/repository.git".into()],
3013 checked_at: "2026-09-01T01:00:00Z".into(),
3014 reason: Some("feature branch was pushed but not merged".into()),
3015 });
3016 assert_eq!(ready(vec![session.clone()], vec![]), vec![id]);
3017 session.publication.as_mut().unwrap().stashed = true;
3018 assert!(ready(vec![session.clone()], vec![]).is_empty());
3019 session.publication.as_mut().unwrap().stashed = false;
3020 session.publication.as_mut().unwrap().checkpoint_sha256 = "b".repeat(64);
3021 assert!(ready(vec![session], vec![]).is_empty());
3022 }
3023
3024 fn sized_session(
3026 root: &Path,
3027 session_id: &str,
3028 updated_at: &str,
3029 checkpoint_bytes: usize,
3030 attachment_bytes: &[usize],
3031 ) -> SessionRecord {
3032 let archive_path = root.join(format!("{session_id}.hel.zip"));
3033 std::fs::write(&archive_path, vec![b'c'; checkpoint_bytes]).unwrap();
3034 if !attachment_bytes.is_empty() {
3035 let attachments = root
3036 .join(session_id)
3037 .join(mj_core::attachment::ATTACHMENT_DIR);
3038 std::fs::create_dir_all(&attachments).unwrap();
3039 for (index, size) in attachment_bytes.iter().enumerate() {
3040 std::fs::write(attachments.join(format!("{index}.png")), vec![b'a'; *size])
3041 .unwrap();
3042 }
3043 }
3044 SessionRecord {
3045 project: None,
3046 checkpoint: Some(mj_core::state::CheckpointMetadata {
3047 archive_path,
3048 sha256: "0".repeat(64),
3049 created_at: updated_at.into(),
3050 event_frontier: 1,
3051 }),
3052 ..record(
3053 session_id,
3054 mj_core::state::SessionState::Stopped,
3055 updated_at,
3056 )
3057 }
3058 }
3059
3060 #[test]
3061 fn the_space_preview_sizes_every_session_and_only_the_aged_ones_as_reclaimable() {
3062 let directory = tempfile::tempdir().unwrap();
3063 let root = directory.path();
3064 let sessions: mj_core::snapshot_map::SnapshotMap<String, SessionRecord> = [
3065 sized_session(root, "old-stopped", "2026-09-01T00:00:00Z", 1000, &[10, 20]),
3066 sized_session(root, "just-stopped", "2026-09-09T00:00:00Z", 500, &[]),
3067 SessionRecord {
3070 project: None,
3071 checkpoint: Some(mj_core::state::CheckpointMetadata {
3072 archive_path: root.join("missing.hel.zip"),
3073 sha256: "0".repeat(64),
3074 created_at: "2026-09-01T00:00:00Z".into(),
3075 event_frontier: 1,
3076 }),
3077 ..record(
3078 "lost-checkpoint",
3079 mj_core::state::SessionState::Stopped,
3080 "2026-09-01T00:00:00Z",
3081 )
3082 },
3083 ]
3084 .into_iter()
3085 .map(|record| (record.id.clone(), record))
3086 .collect();
3087 let now = parse_time("2026-09-10T00:00:00Z").unwrap();
3088
3089 let all = archive_space_over(root, &sessions, &Default::default(), now, None);
3090 assert_eq!(all.sessions, 3);
3091 assert_eq!(all.bytes, 1530);
3092 assert_eq!(all.reclaimable_sessions, 0);
3093 assert_eq!(all.reclaimable_bytes, 0);
3094
3095 let aged = archive_space_over(root, &sessions, &Default::default(), now, Some(3));
3096 assert_eq!(aged.bytes, 1530);
3097 assert_eq!(
3098 (aged.reclaimable_sessions, aged.reclaimable_bytes),
3099 (2, 1030),
3100 "only the sessions the job would archive count, attachments included"
3101 );
3102 }
3103
3104 #[test]
3105 fn only_stopped_sessions_past_the_cut_off_are_archived() {
3106 use mj_core::state::SessionState;
3107 let selected = ready(
3108 vec![
3109 record("old-stopped", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3110 record(
3111 "just-stopped",
3112 SessionState::Stopped,
3113 "2026-09-09T00:00:00Z",
3114 ),
3115 record("old-running", SessionState::Running, "2026-09-01T00:00:00Z"),
3116 record("old-error", SessionState::Error, "2026-09-01T00:00:00Z"),
3117 record("unparsable", SessionState::Stopped, "not a time"),
3118 record("at-the-edge", SessionState::Stopped, "2026-09-07T00:00:00Z"),
3120 ],
3121 Vec::new(),
3122 );
3123 assert_eq!(selected, vec!["at-the-edge", "old-stopped"]);
3124 }
3125
3126 #[test]
3127 fn a_child_the_pass_is_not_archiving_holds_its_parent_back() {
3128 use mj_core::state::SessionState;
3129 let selected = ready(
3130 vec![
3131 record("parent", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3132 record(
3133 "running-child",
3134 SessionState::Running,
3135 "2026-09-01T00:00:00Z",
3136 ),
3137 ],
3138 vec![child("running-child", "parent")],
3139 );
3140 assert!(selected.is_empty(), "the parent must wait: {selected:?}");
3141
3142 let selected = ready(
3143 vec![
3144 record("parent", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3145 record("young-child", SessionState::Stopped, "2026-09-09T00:00:00Z"),
3146 ],
3147 vec![child("young-child", "parent")],
3148 );
3149 assert!(selected.is_empty(), "the parent must wait: {selected:?}");
3150
3151 let selected = ready(
3153 vec![record(
3154 "parent",
3155 SessionState::Stopped,
3156 "2026-09-01T00:00:00Z",
3157 )],
3158 vec![child("departed-child", "parent")],
3159 );
3160 assert_eq!(selected, vec!["parent"]);
3161 }
3162
3163 #[test]
3164 fn children_are_archived_before_their_parents() {
3165 use mj_core::state::SessionState;
3166 let selected = ready(
3167 vec![
3168 record("parent", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3169 record("child", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3170 record("grandchild", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3171 ],
3172 vec![child("child", "parent"), child("grandchild", "child")],
3173 );
3174 assert_eq!(selected, vec!["grandchild", "child", "parent"]);
3175 }
3176
3177 #[test]
3178 fn native_adapters_cover_every_enabled_profile_home() {
3179 use mj_core::config::{Config, HarnessKind, HarnessProfile};
3180
3181 fn profile(kind: HarnessKind, home: &str, enabled: bool) -> HarnessProfile {
3182 HarnessProfile {
3183 enabled,
3184 kind,
3185 home: PathBuf::from(home),
3186 environment: Default::default(),
3187 context_window_bytes: None,
3188 subagents: Default::default(),
3189 guardian_review_model: None,
3190 }
3191 }
3192
3193 let mut config = Config::default();
3194 for (id, built) in [
3195 (
3196 "codex",
3197 profile(HarnessKind::Codex, "/home/dev/.codex3", true),
3198 ),
3199 (
3200 "codex-ds",
3201 profile(HarnessKind::Codex, "/home/dev/.codex-ds", true),
3202 ),
3203 (
3205 "codex-alt",
3206 profile(HarnessKind::Codex, "/home/dev/.codex3", true),
3207 ),
3208 (
3209 "codex-off",
3210 profile(HarnessKind::Codex, "/home/dev/.codex-off", false),
3211 ),
3212 (
3213 "claude",
3214 profile(HarnessKind::Claude, "/home/dev/.claude4", true),
3215 ),
3216 ("kimi", profile(HarnessKind::Kimi, "/home/dev/.kimi", true)),
3217 ("grok", profile(HarnessKind::Grok, "/home/dev/.grok", true)),
3218 ("muse", profile(HarnessKind::Muse, "/home/dev/muse", true)),
3219 (
3220 "muse-off",
3221 profile(HarnessKind::Muse, "/home/dev/muse-off", false),
3222 ),
3223 ] {
3224 config.profiles.insert(id.into(), built);
3225 }
3226
3227 let adapters = native_adapters(&config);
3228 let roots: Vec<(&str, Option<PathBuf>)> = adapters
3229 .iter()
3230 .map(|adapter| (adapter.name(), adapter.root()))
3231 .collect();
3232
3233 let codex: Vec<&Option<PathBuf>> = roots
3234 .iter()
3235 .filter(|(name, _)| *name == "codex")
3236 .map(|(_, root)| root)
3237 .collect();
3238 assert_eq!(
3239 codex,
3240 vec![
3241 &Some(PathBuf::from("/home/dev/.codex3/sessions")),
3242 &Some(PathBuf::from("/home/dev/.codex-ds/sessions")),
3243 ],
3244 "one adapter per enabled Codex home, deduplicated: {roots:?}"
3245 );
3246
3247 let claude: Vec<&Option<PathBuf>> = roots
3248 .iter()
3249 .filter(|(name, _)| *name == "claude-code")
3250 .map(|(_, root)| root)
3251 .collect();
3252 assert_eq!(
3253 claude,
3254 vec![&Some(PathBuf::from("/home/dev/.claude4/projects"))],
3255 "one adapter for the enabled Claude home: {roots:?}"
3256 );
3257
3258 for (_, root) in &roots {
3259 let Some(root) = root else { continue };
3260 let text = root.to_string_lossy();
3261 assert!(
3262 !text.contains(".codex-off"),
3263 "a disabled profile must not be indexed: {roots:?}"
3264 );
3265 assert!(
3266 !text.ends_with("/.codex/sessions") && !text.ends_with("/.claude/projects"),
3267 "the stock homes are not indexed unless a profile names them: {roots:?}"
3268 );
3269 }
3270
3271 for (name, root) in [
3274 ("kimi-code", PathBuf::from("/home/dev/.kimi/sessions")),
3275 ("grok-build", PathBuf::from("/home/dev/.grok/sessions")),
3276 (
3277 "muse",
3278 mj_checkpoint::native::muse_sessions_root(Path::new("/home/dev/muse")).unwrap(),
3279 ),
3280 ] {
3281 let found: Vec<&Option<PathBuf>> = roots
3282 .iter()
3283 .filter(|(found, _)| *found == name)
3284 .map(|(_, root)| root)
3285 .collect();
3286 assert_eq!(found, vec![&Some(root)], "one {name} adapter: {roots:?}");
3287 }
3288
3289 for (_, root) in &roots {
3290 let Some(root) = root else { continue };
3291 assert!(
3292 !root.to_string_lossy().contains("muse-off"),
3293 "a disabled profile must not be indexed: {roots:?}"
3294 );
3295 }
3296
3297 assert!(
3298 roots.iter().any(|(name, _)| *name == "gemini"),
3299 "the other built-in adapters are kept: {roots:?}"
3300 );
3301 }
3302
3303 #[test]
3307 fn query_rows_returns_the_indexed_target_profile_and_harness() {
3308 let _held = tags::testing::lock();
3309 let (_directory, connection) = tags::testing::isolated_index();
3310 tags::testing::index_row(&connection, "mj-session", TOOL);
3311 tags::testing::index_row(&connection, "codex-session", "codex");
3312 tags::write(
3313 &connection,
3314 "mj-session",
3315 &tags::MjTags {
3316 target: Some("Prod-Box".into()),
3317 profile: Some("codex-Main".into()),
3318 harness: Some("codex".into()),
3319 },
3320 )
3321 .expect("write the session metadata");
3322
3323 let rows = query_rows("", 10, &BTreeSet::new(), false).expect("query the index");
3324 let mjolnir = rows
3325 .iter()
3326 .find(|row| row.id == "mj-session")
3327 .expect("the Mjolnir row is returned");
3328 assert_eq!(mjolnir.target.as_deref(), Some("Prod-Box"));
3329 assert_eq!(mjolnir.profile.as_deref(), Some("codex-Main"));
3330 assert_eq!(mjolnir.harness.as_deref(), Some("codex"));
3331
3332 let codex = rows
3333 .iter()
3334 .find(|row| row.id == "codex-session")
3335 .expect("the Codex row is returned");
3336 assert_eq!(codex.target, None);
3337 assert_eq!(codex.profile, None);
3338 assert_eq!(codex.harness, None);
3339 }
3340
3341 #[test]
3342 fn one_flag_keeps_sub_agents_out_of_every_query_path() {
3343 let _held = tags::testing::lock();
3344 let (_directory, connection) = tags::testing::isolated_index();
3345 for (session_id, kind) in [("main-session", "main"), ("sub-session", "sub")] {
3346 tags::testing::index_row(&connection, session_id, "claude");
3347 connection
3348 .execute(
3349 "UPDATE files SET kind = ?2 WHERE session_id = ?1",
3350 rusqlite::params![session_id, kind],
3351 )
3352 .expect("set the session kind");
3353 connection
3354 .execute(
3355 "INSERT INTO messages(session_id, role, text)
3356 VALUES (?1, 'user', 'fix the bridge derivation zq')",
3357 [session_id],
3358 )
3359 .expect("insert a message");
3360 connection
3361 .execute(
3362 "INSERT INTO msgs(rowid, text) VALUES (?1, 'fix the bridge derivation zq')",
3363 [connection.last_insert_rowid()],
3364 )
3365 .expect("index the message");
3366 }
3367 let ids = |query: &str, include_subagents: bool| {
3368 let mut ids: Vec<String> = query_rows(query, 10, &BTreeSet::new(), include_subagents)
3369 .expect("query the index")
3370 .into_iter()
3371 .map(|row| row.id)
3372 .collect();
3373 ids.sort();
3374 ids
3375 };
3376
3377 for query in ["", "bridge derivation", "zq", "an indexed session"] {
3379 assert_eq!(ids(query, false), ["main-session"], "query {query:?}");
3380 assert_eq!(
3381 ids(query, true),
3382 ["main-session", "sub-session"],
3383 "query {query:?}"
3384 );
3385 }
3386 }
3387
3388 #[test]
3395 fn a_phrase_only_in_a_sub_agents_transcript_does_not_match_its_parent() {
3396 let _held = tags::testing::lock();
3397 let (_directory, connection) = tags::testing::isolated_index();
3398 let message = |session_id: &str, role: &str, text: &str| {
3399 connection
3400 .execute(
3401 "INSERT INTO messages(session_id, role, text) VALUES (?1, ?2, ?3)",
3402 rusqlite::params![session_id, role, text],
3403 )
3404 .expect("insert a message");
3405 connection
3406 .execute(
3407 "INSERT INTO msgs(rowid, text) VALUES (?1, ?2)",
3408 rusqlite::params![connection.last_insert_rowid(), text],
3409 )
3410 .expect("index the message");
3411 };
3412 for (session_id, kind) in [("parent", "main"), ("child", "sub")] {
3413 tags::testing::index_row(&connection, session_id, "claude");
3414 connection
3415 .execute(
3416 "UPDATE files SET kind = ?2 WHERE session_id = ?1",
3417 rusqlite::params![session_id, kind],
3418 )
3419 .expect("set the session kind");
3420 }
3421 message("parent", "user", "look into the relay journal");
3422 message("parent", "tool", "Task {\"prompt\":\"read the journal\"}");
3423 message("parent", "tool", "the journal uses a quokka checksum");
3424 message(
3425 "parent",
3426 "assistant",
3427 "The journal is fine; the parent zebra ends here.",
3428 );
3429 message("child", "user", "read the journal");
3430 message("child", "assistant", "the journal uses a quokka checksum");
3431
3432 let ids = |query: &str, include_subagents: bool| {
3433 let mut ids: Vec<String> = query_rows(query, 10, &BTreeSet::new(), include_subagents)
3434 .expect("query the index")
3435 .into_iter()
3436 .map(|row| row.id)
3437 .collect();
3438 ids.sort();
3439 ids
3440 };
3441 assert!(
3442 ids("quokka", false).is_empty(),
3443 "{:?}",
3444 ids("quokka", false)
3445 );
3446 assert_eq!(ids("quokka", true), ["child", "parent"]);
3448 assert_eq!(ids("parent zebra", false), ["parent"]);
3449 }
3450
3451 #[test]
3452 fn short_query_scan_also_ignores_tool_only_matches() {
3453 let _held = tags::testing::lock();
3454 let (_directory, connection) = tags::testing::isolated_index();
3455 tags::testing::index_row(&connection, "parent", "claude");
3456 connection
3457 .execute(
3458 "INSERT INTO messages(session_id, role, text) VALUES ('parent', 'tool', 'qx')",
3459 [],
3460 )
3461 .expect("insert a message");
3462 assert!(
3463 query_rows("qx", 10, &BTreeSet::new(), false)
3464 .expect("query the index")
3465 .is_empty()
3466 );
3467 }
3468
3469 #[test]
3470 fn full_text_search_fills_its_result_limit_after_skipping_sub_agents() {
3471 let _held = tags::testing::lock();
3472 let (_directory, connection) = tags::testing::isolated_index();
3473 for (id, kind, text) in [
3474 ("sub", "sub", "restic restic restic restic"),
3475 ("main", "main", "restic cleanup"),
3476 ] {
3477 tags::testing::index_row(&connection, id, "codex");
3478 connection
3479 .execute(
3480 "UPDATE files SET kind = ?2 WHERE session_id = ?1",
3481 rusqlite::params![id, kind],
3482 )
3483 .unwrap();
3484 connection
3485 .execute(
3486 "INSERT INTO messages(session_id, role, text) VALUES (?1, 'user', ?2)",
3487 rusqlite::params![id, text],
3488 )
3489 .unwrap();
3490 connection
3491 .execute(
3492 "INSERT INTO msgs(rowid, text) VALUES (?1, ?2)",
3493 rusqlite::params![connection.last_insert_rowid(), text],
3494 )
3495 .unwrap();
3496 }
3497
3498 let rows = query_rows("restic", 1, &BTreeSet::new(), false).unwrap();
3499 assert_eq!(
3500 rows.iter().map(|row| row.id.as_str()).collect::<Vec<_>>(),
3501 ["main"]
3502 );
3503 }
3504
3505 fn block_on<F: std::future::Future>(future: F) -> F::Output {
3508 tokio::runtime::Builder::new_current_thread()
3509 .enable_all()
3510 .build()
3511 .unwrap()
3512 .block_on(future)
3513 }
3514
3515 #[test]
3520 fn a_destroy_indexes_the_session_itself_when_the_sync_outlasts_the_wait() {
3521 let _held = tags::testing::lock();
3522 let (_index_dir, _connection) = tags::testing::isolated_index();
3523 let directory = tempfile::tempdir().unwrap();
3524 let session_id = "0123456789abcdef0123456789abcdef";
3525 write_archive(directory.path(), session_id, 1);
3526 let source = adapter(directory.path(), session_id);
3527
3528 let started = Instant::now();
3529 let outcome = block_on(index_before_destroy_with(
3530 std::future::pending::<Result<()>>(),
3531 Duration::from_millis(200),
3532 move || capture_sessions_from(&source, &[session_id.to_owned()]),
3533 Duration::from_millis(50),
3534 ));
3535
3536 assert_eq!(outcome, IndexedBeforeDestroy::WrittenDirectly);
3537 assert!(
3538 started.elapsed() < Duration::from_secs(10),
3539 "the destroy must not wait for the pass: {:?}",
3540 started.elapsed()
3541 );
3542 let found = wiki_session(session_id, &BTreeSet::new())
3543 .unwrap()
3544 .expect("the session is found by its id");
3545 assert_eq!(found.status, WikiSessionStatus::Archived);
3546 assert_eq!(found.tool, TOOL);
3547 assert_eq!(
3548 found.path,
3549 PathBuf::from(format!("{}/{session_id}", directory.path().display()))
3550 );
3551 assert_eq!(found.title, "the harness title");
3552 assert_eq!(
3553 found.harness,
3554 Some(HarnessKind::Codex),
3555 "the session's metadata is written beside its row"
3556 );
3557 assert!(!found.nothing_to_restore);
3558 }
3559
3560 #[test]
3563 fn a_sync_that_finishes_in_time_is_all_a_destroy_waits_for() {
3564 let outcome = block_on(index_before_destroy_with(
3565 async { Ok(()) },
3566 DESTROY_SYNC_WAIT,
3567 || -> Result<Vec<CapturedSession>> {
3568 panic!("a finished pass leaves nothing to index on its own")
3569 },
3570 Duration::from_millis(50),
3571 ));
3572 assert_eq!(outcome, IndexedBeforeDestroy::Synced);
3573 }
3574
3575 #[test]
3579 fn a_busy_index_takes_the_destroyed_session_once_it_is_free() {
3580 let _held = tags::testing::lock();
3581 let (_index_dir, writer) = tags::testing::isolated_index();
3582 let directory = tempfile::tempdir().unwrap();
3583 let session_id = "0123456789abcdef0123456789abcdef";
3584 write_archive(directory.path(), session_id, 1);
3585 let source = adapter(directory.path(), session_id);
3586
3587 block_on(async {
3588 writer.execute_batch("BEGIN IMMEDIATE").unwrap();
3589 let outcome = index_before_destroy_with(
3590 std::future::pending::<Result<()>>(),
3591 Duration::from_millis(50),
3592 move || capture_sessions_from(&source, &[session_id.to_owned()]),
3593 Duration::from_millis(50),
3594 )
3595 .await;
3596 assert_eq!(outcome, IndexedBeforeDestroy::Deferred);
3597 assert!(
3598 wiki_session(session_id, &BTreeSet::new())
3599 .unwrap()
3600 .is_none(),
3601 "nothing is written while the other writer holds the index"
3602 );
3603
3604 writer.execute_batch("COMMIT").unwrap();
3605 let deadline = Instant::now() + Duration::from_secs(30);
3606 while wiki_session(session_id, &BTreeSet::new())
3607 .unwrap()
3608 .is_none()
3609 {
3610 assert!(
3611 Instant::now() < deadline,
3612 "the deferred row never reached the index"
3613 );
3614 tokio::time::sleep(Duration::from_millis(50)).await;
3615 }
3616 });
3617 }
3618
3619 #[test]
3623 fn a_session_the_index_holds_as_it_is_now_needs_no_indexing() {
3624 let _held = tags::testing::lock();
3625 let (_index_dir, _connection) = tags::testing::isolated_index();
3626 let directory = tempfile::tempdir().unwrap();
3627 let session_id = "0123456789abcdef0123456789abcdef";
3628 let never_prompted = "fedcba9876543210fedcba9876543210";
3629 write_archive(directory.path(), session_id, 1);
3630 let source = adapter(directory.path(), session_id);
3631 let ids = [session_id.to_owned(), never_prompted.to_owned()];
3632
3633 assert_eq!(
3634 unindexed(&source, &ids).unwrap(),
3635 [session_id],
3636 "a session with no conversation has nothing to index"
3637 );
3638 let captured = Arc::new(capture_sessions_from(&source, &ids).unwrap());
3639 write_captured(&captured).unwrap();
3640 assert!(unindexed(&source, &ids).unwrap().is_empty());
3641
3642 source
3644 .sessions
3645 .lock()
3646 .unwrap()
3647 .records
3648 .get_mut(session_id)
3649 .unwrap()
3650 .updated_at = "2099-01-01T00:00:00Z".into();
3651 assert_eq!(unindexed(&source, &ids).unwrap(), [session_id]);
3652 }
3653
3654 #[test]
3655 fn a_session_tree_holds_the_sub_agents_below_it_and_nothing_else() {
3656 let record = |id: &str| {
3657 (
3658 id.to_owned(),
3659 SessionRecord {
3660 project: None,
3661 id: id.into(),
3662 ..record_template()
3663 },
3664 )
3665 };
3666 let state = State {
3667 sessions: [
3668 record("parent"),
3669 record("child"),
3670 record("grandchild"),
3671 record("sibling"),
3672 ]
3673 .into_iter()
3674 .collect(),
3675 subagents: [
3676 ("child".to_owned(), child("child", "parent")),
3677 ("grandchild".to_owned(), child("grandchild", "child")),
3678 ("sibling".to_owned(), child("sibling", "other-parent")),
3679 ]
3680 .into_iter()
3681 .collect(),
3682 ..State::default()
3683 };
3684 assert_eq!(
3685 session_tree(&state, "parent"),
3686 ["parent", "child", "grandchild"]
3687 );
3688 assert!(session_tree(&state, "unknown").is_empty());
3689 }
3690
3691 mod text_search {
3693 use super::super::{SessionTextMatch, SessionTextMatchKind, text_matches_in};
3694 use std::collections::BTreeSet;
3695
3696 fn index(sessions: &[(&str, &[(&str, &str)])]) -> rusqlite::Connection {
3697 let connection = rusqlite::Connection::open_in_memory().unwrap();
3698 connection
3699 .execute_batch(
3700 "CREATE TABLE files(path TEXT PRIMARY KEY, session_id TEXT NOT NULL,
3701 tool TEXT NOT NULL);
3702 CREATE TABLE messages(id INTEGER PRIMARY KEY, session_id TEXT NOT NULL,
3703 role TEXT NOT NULL, text TEXT NOT NULL);
3704 CREATE VIRTUAL TABLE msgs USING fts5(
3705 text, content='messages', content_rowid='id', tokenize='trigram');",
3706 )
3707 .unwrap();
3708 for (id, messages) in sessions {
3709 connection
3710 .execute(
3711 "INSERT INTO files VALUES (?1, ?2, 'mjolnir')",
3712 rusqlite::params![format!("/checkpoints/{id}"), id],
3713 )
3714 .unwrap();
3715 for (role, text) in *messages {
3716 connection
3717 .execute(
3718 "INSERT INTO messages(session_id, role, text) VALUES (?1, ?2, ?3)",
3719 rusqlite::params![id, role, text],
3720 )
3721 .unwrap();
3722 let rowid = connection.last_insert_rowid();
3723 connection
3724 .execute(
3725 "INSERT INTO msgs(rowid, text) VALUES (?1, ?2)",
3726 rusqlite::params![rowid, text],
3727 )
3728 .unwrap();
3729 }
3730 }
3731 connection
3732 }
3733
3734 fn live(ids: &[&str]) -> BTreeSet<String> {
3735 ids.iter().map(|id| (*id).to_owned()).collect()
3736 }
3737
3738 fn matches(kinds: &[(&str, SessionTextMatchKind)]) -> Vec<SessionTextMatch> {
3739 kinds
3740 .iter()
3741 .map(|(id, kind)| SessionTextMatch {
3742 session_id: (*id).to_owned(),
3743 kind: *kind,
3744 })
3745 .collect()
3746 }
3747
3748 #[test]
3749 fn user_and_agent_messages_match_and_tool_output_does_not() {
3750 let connection = index(&[
3751 ("said-by-user", &[("user", "please fix the Zebra crossing")]),
3752 ("said-by-agent", &[("assistant", "the zebra is fixed")]),
3753 (
3754 "only-in-tool",
3755 &[("tool", "zebra stack trace"), ("user", "hello")],
3756 ),
3757 (
3758 "both",
3759 &[("assistant", "a ZEBRA appears"), ("user", "a zebra please")],
3760 ),
3761 ("gone", &[("user", "zebra")]),
3762 ]);
3763 let live = live(&["said-by-user", "said-by-agent", "only-in-tool", "both"]);
3764 assert_eq!(
3765 text_matches_in(&connection, "zebra", &live).unwrap(),
3766 matches(&[
3767 ("both", SessionTextMatchKind::User),
3768 ("said-by-agent", SessionTextMatchKind::Agent),
3769 ("said-by-user", SessionTextMatchKind::User),
3770 ])
3771 );
3772 }
3773
3774 #[test]
3775 fn a_query_too_short_for_the_trigram_index_still_matches() {
3776 let connection = index(&[
3777 ("a", &[("user", "go to the zoo")]),
3778 ("b", &[("assistant", "zoo")]),
3779 ("c", &[("tool", "zoo")]),
3780 ]);
3781 assert_eq!(
3782 text_matches_in(&connection, "zo", &live(&["a", "b", "c"])).unwrap(),
3783 matches(&[
3784 ("a", SessionTextMatchKind::User),
3785 ("b", SessionTextMatchKind::Agent),
3786 ])
3787 );
3788 assert_eq!(
3790 text_matches_in(&connection, "%z", &live(&["a"])).unwrap(),
3791 Vec::new()
3792 );
3793 }
3794 }
3795}