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::Instant;
22
23use anyhow::{Context, Result};
24use chrono::{DateTime, Utc};
25
26use mj_client::daemon::{
27 WikiHitBlock, WikiHitTranscript, WikiIndexState, WikiRow, WikiSessionInfo, WikiSessionStatus,
28 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: BTreeMap<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 pub async fn sync_now(&self, full: bool) -> Result<()> {
568 self.inner.sync(full).await
569 }
570
571 pub fn status(&self) -> WikiStatus {
574 WikiStatus {
575 state: index_state(),
576 topping_up: self.inner.in_flight.load(Ordering::Acquire)
577 || self.inner.requested.load(Ordering::Acquire),
578 }
579 }
580
581 pub fn last_success(&self) -> Option<Instant> {
583 self.inner
584 .last_success
585 .lock()
586 .unwrap_or_else(std::sync::PoisonError::into_inner)
587 .map(|success| success.at)
588 }
589}
590
591impl Indexer {
592 async fn run(self: Arc<Self>) {
593 loop {
594 self.notify.notified().await;
595 while self.requested.swap(false, Ordering::AcqRel) {
596 let full = self.full_requested.swap(false, Ordering::AcqRel);
597 if let Err(error) = self.sync(full).await {
598 self.report(&error);
599 break;
604 }
605 }
606 }
607 }
608
609 fn report(&self, error: &anyhow::Error) {
613 if crate::database::is_busy_error(error) {
614 self.requested.store(true, Ordering::Release);
615 tracing::debug!(%error, "the SessionWiki index was busy; retrying on the next trigger");
616 } else {
617 tracing::warn!(%error, "could not sync sessions into SessionWiki");
618 }
619 }
620
621 async fn sync(&self, full: bool) -> Result<()> {
622 let _guard = self.running.lock().await;
623 let since = if full {
624 None
625 } else {
626 self.last_success
627 .lock()
628 .unwrap_or_else(std::sync::PoisonError::into_inner)
629 .map(|success| success.epoch_seconds - 60)
632 };
633 let started = Instant::now();
634 let work = crate::upgrade::activity("SessionWiki sync")?;
635 self.in_flight.store(true, Ordering::Release);
636 let ran = tokio::task::spawn_blocking(move || {
637 let _work = work;
638 sync_blocking(since)
639 })
640 .await;
641 self.in_flight.store(false, Ordering::Release);
642 let ran = ran.context("run the SessionWiki sync")??;
643 if ran {
644 *self
645 .last_success
646 .lock()
647 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(Success {
648 at: started,
649 epoch_seconds: Utc::now().timestamp(),
650 });
651 }
652 Ok(())
653 }
654}
655
656fn sync_blocking(since: Option<i64>) -> Result<bool> {
659 if !index_is_writable() {
660 return Ok(false);
661 }
662 let controller =
663 Controller::load().context("load controller state for the SessionWiki sync")?;
664 let mjolnir = Arc::new(MjolnirAdapter::reloading(&controller.state));
667 let mut adapters: Vec<Box<dyn sessionwiki::adapters::Adapter>> =
668 vec![Box::new(SharedMjolnirAdapter(Arc::clone(&mjolnir)))];
669 adapters.extend(native_adapters(&controller.config));
670 let mut connection = sessionwiki::index::open().context("open the SessionWiki index")?;
671 sessionwiki::index::sync_with(&mut connection, &adapters, since)
672 .context("sync the SessionWiki index")?;
673 write_session_tags(&mut connection, &mjolnir.indexed_tags())
674 .context("store Mjolnir's session metadata in the SessionWiki index")?;
675 provenance::backfill(&mut connection, &mjolnir).context("backfill Mjolnir file provenance")?;
676 if since.is_none() {
677 record_first_build();
681 }
682 Ok(true)
683}
684
685fn write_session_tags(
695 connection: &mut rusqlite::Connection,
696 session_tags: &BTreeMap<String, tags::MjTags>,
697) -> Result<()> {
698 if session_tags.is_empty() {
699 return Ok(());
700 }
701 let transaction = connection
702 .transaction()
703 .context("open a transaction for the session metadata")?;
704 for (session_id, session) in session_tags {
705 if session.is_empty() {
706 continue;
707 }
708 tags::write(&transaction, session_id, session)?;
709 }
710 transaction
711 .commit()
712 .context("commit the session metadata")?;
713 Ok(())
714}
715
716fn native_adapters(config: &mj_core::config::Config) -> Vec<Box<dyn Adapter>> {
733 let mut seen: BTreeSet<(HarnessKind, &Path)> = BTreeSet::new();
736 let mut adapters: Vec<Box<dyn Adapter>> = Vec::new();
737 for (_, profile) in config.enabled_profiles() {
738 if !seen.insert((profile.kind, profile.home.as_path())) {
740 continue;
741 }
742 let adapter: Box<dyn Adapter> = match profile.kind {
743 HarnessKind::Codex => {
744 Box::new(sessionwiki::adapters::Codex::in_home(profile.home.clone()))
745 }
746 HarnessKind::Claude => Box::new(sessionwiki::adapters::ClaudeCode::in_home(
747 profile.home.clone(),
748 )),
749 kind => match HarnessAdapter::in_home(kind, profile.home.clone()) {
750 Some(adapter) => Box::new(adapter),
751 None => continue,
752 },
753 };
754 adapters.push(adapter);
755 }
756 adapters.extend(
757 sessionwiki::adapters::all()
758 .into_iter()
759 .filter(|adapter| !matches!(adapter.name(), "codex" | "claude-code")),
760 );
761 adapters
762}
763
764fn index_is_isolated() -> bool {
777 static SAID: AtomicBool = AtomicBool::new(false);
778 if mj_core::config::session_index_is_resolved()
779 || std::env::var_os(mj_core::config::SESSION_INDEX_ENV).is_some()
780 {
781 return true;
782 }
783 if !SAID.swap(true, Ordering::AcqRel) {
784 tracing::debug!(
785 "this process did not resolve a session index location; SessionWiki is not used"
786 );
787 }
788 false
789}
790
791fn index_version_mismatch() -> bool {
800 static SAID: AtomicBool = AtomicBool::new(false);
801 let Ok(path) = sessionwiki::index::db_path() else {
802 return false;
803 };
804 if !path.exists() {
805 return false;
806 }
807 let version = rusqlite::Connection::open_with_flags(
808 &path,
809 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_URI,
810 )
811 .and_then(|connection| connection.pragma_query_value(None, "user_version", |row| row.get(0)));
812 let version: i64 = match version {
813 Ok(version) => version,
814 Err(error) => {
815 tracing::debug!(%error, "could not read the SessionWiki index schema version");
816 return false;
817 }
818 };
819 let mismatch = version != 0 && version != sessionwiki::index::SCHEMA_VERSION;
822 if mismatch && !SAID.swap(true, Ordering::AcqRel) {
823 tracing::warn!(
824 found = version,
825 expected = sessionwiki::index::SCHEMA_VERSION,
826 path = %path.display(),
827 "the SessionWiki index was written by another version; Mjolnir will not open it, because opening it would rebuild it. Install the matching sessionwiki command"
828 );
829 }
830 mismatch
831}
832
833fn index_is_writable() -> bool {
834 index_is_isolated() && !index_version_mismatch()
835}
836
837fn first_build_marker() -> PathBuf {
840 mj_core::config::data_dir().join("sessionwiki-built")
841}
842
843fn record_first_build() {
844 let path = first_build_marker();
845 let version = sessionwiki::index::SCHEMA_VERSION.to_string();
846 if std::fs::read_to_string(&path).is_ok_and(|held| held.trim() == version) {
847 return;
848 }
849 if let Err(error) = std::fs::write(&path, &version) {
850 tracing::warn!(%error, path = %path.display(), "could not record the first SessionWiki build");
851 }
852}
853
854fn first_build_is_done() -> bool {
856 std::fs::read_to_string(first_build_marker())
857 .is_ok_and(|held| held.trim() == sessionwiki::index::SCHEMA_VERSION.to_string())
858 && sessionwiki::index::db_path().is_ok_and(|path| path.exists())
859}
860
861pub fn index_state() -> WikiIndexState {
863 if !index_is_isolated() {
864 return WikiIndexState::Indexing;
865 }
866 if index_version_mismatch() {
867 return WikiIndexState::VersionMismatch;
868 }
869 if first_build_is_done() {
870 WikiIndexState::Ready
871 } else {
872 WikiIndexState::Indexing
873 }
874}
875
876pub const MAX_WIKI_LIMIT: usize = 200;
882pub const DEFAULT_WIKI_LIMIT: usize = 50;
884const MIN_FULLTEXT_QUERY: usize = 3;
887pub const SYNC_STALE_AFTER: std::time::Duration = std::time::Duration::from_secs(60);
889
890pub fn sync_is_stale(last_success: Option<Instant>) -> bool {
892 last_success.is_none_or(|at| at.elapsed() >= SYNC_STALE_AFTER)
893}
894
895pub fn query_rows(
904 query: &str,
905 limit: usize,
906 live: &BTreeSet<String>,
907 include_subagents: bool,
908) -> Result<Vec<WikiRow>> {
909 let limit = limit.clamp(1, MAX_WIKI_LIMIT);
910 if !index_is_writable() {
911 return Ok(Vec::new());
915 }
916 let connection = open_readonly()?;
917 let query = query.trim();
918 if query.is_empty() {
919 let rows =
920 sessionwiki::index::recent(&connection, limit, None, None, None, include_subagents)
921 .context("list recent SessionWiki sessions")?;
922 let mut rows: Vec<WikiRow> = rows
923 .into_iter()
924 .map(|row| wiki_row(row, None, live))
925 .collect();
926 fill_session_tags(&connection, &mut rows)?;
927 return Ok(rows);
928 }
929 let hits = if query.chars().count() < MIN_FULLTEXT_QUERY {
930 sessionwiki::index::search_like(&connection, query, limit, None, None)
931 } else {
932 sessionwiki::index::search(&connection, query, limit, None, None)
933 }
934 .context("search the SessionWiki index")?;
935 let mut rows: Vec<WikiRow> = hits
937 .into_iter()
938 .filter(|hit| include_subagents || is_main_session(&hit.row))
939 .map(|hit| wiki_row(hit.row, Some(hit.snippet), live))
940 .collect();
941 let found: BTreeSet<String> = rows.iter().map(|row| row.id.clone()).collect();
945 for row in named_like(&connection, query, include_subagents)? {
946 if rows.len() >= limit {
947 break;
948 }
949 if found.contains(&row.session_id) {
950 continue;
951 }
952 rows.push(wiki_row(row, None, live));
953 }
954 fill_session_tags(&connection, &mut rows)?;
955 Ok(rows)
956}
957
958fn fill_session_tags(connection: &rusqlite::Connection, rows: &mut [WikiRow]) -> Result<()> {
964 let ids: Vec<&str> = rows
965 .iter()
966 .filter(|row| row.tool == TOOL)
967 .map(|row| row.id.as_str())
968 .collect();
969 let found = tags::read(connection, &ids).context("read the indexed session metadata")?;
970 for row in rows.iter_mut().filter(|row| row.tool == TOOL) {
971 let Some(session) = found.get(&row.id) else {
972 continue;
973 };
974 row.target = session.target.clone();
975 row.profile = session.profile.clone();
976 row.harness = session.harness.clone();
977 }
978 Ok(())
979}
980
981const NAME_SCAN_LIMIT: usize = 2_000;
985
986fn named_like(
988 connection: &rusqlite::Connection,
989 query: &str,
990 include_subagents: bool,
991) -> Result<Vec<sessionwiki::index::SessionRow>> {
992 let needle = query.to_lowercase();
993 let rows = sessionwiki::index::recent(
994 connection,
995 NAME_SCAN_LIMIT,
996 None,
997 None,
998 None,
999 include_subagents,
1000 )
1001 .context("list recent SessionWiki sessions")?;
1002 Ok(rows
1003 .into_iter()
1004 .filter(|row| {
1005 row.title.to_lowercase().contains(&needle)
1006 || row.project.to_lowercase().contains(&needle)
1007 })
1008 .collect())
1009}
1010
1011fn is_main_session(row: &sessionwiki::index::SessionRow) -> bool {
1014 row.kind == "main"
1015}
1016
1017pub fn brief(id: &str, max_chars: usize) -> Result<Option<String>> {
1019 if !index_is_writable() {
1020 return Ok(None);
1021 }
1022 let connection = open_readonly()?;
1023 let Some(row) = row_by_id(&connection, id)? else {
1024 return Ok(None);
1025 };
1026 let session = sessionwiki::index::session_from_index(&connection, &row)
1027 .context("read an indexed session")?;
1028 Ok(Some(sessionwiki::commands::brief_markdown(
1029 &session, max_chars, true,
1030 )))
1031}
1032
1033pub fn transcript_hits(
1041 id: &str,
1042 query: &str,
1043 context_messages: usize,
1044 per_message_chars: usize,
1045) -> Result<Option<WikiHitTranscript>> {
1046 if !index_is_writable() {
1047 return Ok(None);
1048 }
1049 let connection = open_readonly()?;
1050 let Some(row) = row_by_id(&connection, id)? else {
1051 return Ok(None);
1052 };
1053 let session = sessionwiki::index::session_from_index(&connection, &row)
1054 .context("read an indexed session")?;
1055 Ok(Some(hit_transcript(
1056 &session,
1057 query,
1058 context_messages,
1059 per_message_chars,
1060 )))
1061}
1062
1063fn hit_transcript(
1073 session: &Session,
1074 query: &str,
1075 context_messages: usize,
1076 per_message_chars: usize,
1077) -> WikiHitTranscript {
1078 let found = sessionwiki::grep::grep_session(
1079 session,
1080 query,
1081 &sessionwiki::grep::GrepOpts {
1082 context_messages,
1083 chars: per_message_chars,
1084 max_matches: None,
1085 anchor_roles: vec![Role::User, Role::Assistant],
1086 },
1087 );
1088 WikiHitTranscript {
1089 blocks: found
1090 .hits
1091 .into_iter()
1092 .map(|hit| WikiHitBlock {
1093 role: role_name(hit.role).to_owned(),
1094 text: hit.text,
1095 hits: hit.matches,
1096 omitted_before: hit.omitted_before,
1097 truncated: hit.truncated,
1098 })
1099 .collect(),
1100 omitted_after: found.omitted_after,
1101 }
1102}
1103
1104fn role_name(role: Role) -> &'static str {
1105 match role {
1106 Role::User => "user",
1107 Role::Assistant => "assistant",
1108 Role::Tool => "tool",
1109 }
1110}
1111
1112pub struct ArchivedSession {
1116 pub title: String,
1117 pub project_directory: Option<PathBuf>,
1120 pub snapshot: mj_core::archive::CanonicalSessionSnapshot,
1121}
1122
1123pub fn archived_session(id: &str) -> Result<Option<ArchivedSession>> {
1125 if !index_is_writable() {
1126 return Ok(None);
1127 }
1128 let connection = open_readonly()?;
1129 let Some(row) = row_by_id(&connection, id)? else {
1130 return Ok(None);
1131 };
1132 let session = sessionwiki::index::session_from_index(&connection, &row)
1133 .context("read an indexed session")?;
1134 let snapshot = snapshot_of(&session)?;
1135 Ok(Some(ArchivedSession {
1136 title: session.title.clone(),
1137 project_directory: project_directory_of(&session.project),
1138 snapshot,
1139 }))
1140}
1141
1142pub fn sessions_ready_to_archive(
1158 sessions: &BTreeMap<String, SessionRecord>,
1159 subagents: &BTreeMap<String, mj_core::subagent::SubagentRecord>,
1160 now: DateTime<Utc>,
1161 older_than_days: u32,
1162) -> Vec<String> {
1163 let cutoff = now - chrono::Duration::days(i64::from(older_than_days));
1164 let aged = |session_id: &String| {
1165 sessions.get(session_id).is_some_and(|record| {
1166 record.state == mj_core::state::SessionState::Stopped
1167 && parse_time(&record.updated_at).is_some_and(|updated| updated <= cutoff)
1168 })
1169 };
1170 let selected: BTreeSet<String> = sessions
1171 .keys()
1172 .filter(|session_id| aged(session_id))
1173 .filter(|session_id| {
1174 subagents
1175 .values()
1176 .filter(|child| &&child.parent_session_id == session_id)
1177 .filter(|child| sessions.contains_key(&child.child_session_id))
1179 .all(|child| aged(&child.child_session_id))
1180 })
1181 .cloned()
1182 .collect();
1183 let mut ordered: Vec<String> = selected.iter().cloned().collect();
1184 ordered.sort_by_key(|session_id| std::cmp::Reverse(ancestor_depth(session_id, subagents)));
1185 ordered
1186}
1187
1188fn ancestor_depth(
1192 session_id: &str,
1193 subagents: &BTreeMap<String, mj_core::subagent::SubagentRecord>,
1194) -> usize {
1195 let mut depth = 0;
1196 let mut current = session_id;
1197 while let Some(parent) = subagents
1199 .get(current)
1200 .map(|child| child.parent_session_id.as_str())
1201 {
1202 depth += 1;
1203 if depth > subagents.len() {
1204 break;
1205 }
1206 current = parent;
1207 }
1208 depth
1209}
1210
1211pub use mj_core::state::ArchiveSpacePreview;
1217
1218pub fn archive_space_preview(older_than_days: Option<u32>) -> Result<ArchiveSpacePreview> {
1230 let controller =
1231 Controller::load().context("load the session records to size their storage")?;
1232 Ok(archive_space_over(
1233 &mj_core::config::sessions_dir(),
1234 &controller.state.sessions,
1235 &controller.state.subagents,
1236 Utc::now(),
1237 older_than_days,
1238 ))
1239}
1240
1241fn archive_space_over(
1244 sessions_root: &Path,
1245 sessions: &BTreeMap<String, SessionRecord>,
1246 subagents: &BTreeMap<String, mj_core::subagent::SubagentRecord>,
1247 now: DateTime<Utc>,
1248 older_than_days: Option<u32>,
1249) -> ArchiveSpacePreview {
1250 let mut preview = ArchiveSpacePreview {
1251 sessions: sessions.len(),
1252 bytes: sessions
1253 .iter()
1254 .map(|(session_id, record)| session_bytes(sessions_root, session_id, record))
1255 .sum(),
1256 reclaimable_sessions: 0,
1257 reclaimable_bytes: 0,
1258 };
1259 if let Some(days) = older_than_days {
1260 let aged = sessions_ready_to_archive(sessions, subagents, now, days);
1261 preview.reclaimable_sessions = aged.len();
1262 preview.reclaimable_bytes = aged
1263 .iter()
1264 .filter_map(|session_id| {
1265 sessions
1266 .get(session_id)
1267 .map(|record| session_bytes(sessions_root, session_id, record))
1268 })
1269 .sum();
1270 }
1271 preview
1272}
1273
1274fn session_bytes(sessions_root: &Path, session_id: &str, record: &SessionRecord) -> u64 {
1277 let checkpoint = record
1278 .checkpoint
1279 .as_ref()
1280 .and_then(|checkpoint| std::fs::metadata(&checkpoint.archive_path).ok())
1281 .filter(|metadata| metadata.is_file())
1282 .map(|metadata| metadata.len())
1283 .unwrap_or(0);
1284 let attachments = sessions_root
1285 .join(session_id)
1286 .join(mj_core::attachment::ATTACHMENT_DIR);
1287 let attachments = crate::import::claude::directory_size(&attachments).unwrap_or(0);
1288 checkpoint.saturating_add(attachments)
1289}
1290
1291pub fn indexed_with_messages(session_ids: &[String]) -> Result<BTreeSet<String>> {
1298 if !index_is_writable() {
1299 return Ok(BTreeSet::new());
1302 }
1303 let connection = open_readonly()?;
1304 let sessions_dir = mj_core::config::sessions_dir();
1305 let mut indexed = BTreeSet::new();
1306 for session_id in session_ids {
1307 let key = format!("{}/{session_id}", sessions_dir.display());
1308 let rows = sessionwiki::index::resolve(&connection, session_id)
1309 .context("look up a stopped session in the SessionWiki index")?;
1310 if rows
1311 .iter()
1312 .any(|row| row.tool == TOOL && row.path == key && row.msg_count > 0 && !row.archived)
1313 {
1314 indexed.insert(session_id.clone());
1315 }
1316 }
1317 Ok(indexed)
1318}
1319
1320fn open_readonly() -> Result<rusqlite::Connection> {
1321 sessionwiki::index::open_readonly().context("open the SessionWiki index")
1322}
1323
1324fn row_by_id(
1327 connection: &rusqlite::Connection,
1328 id: &str,
1329) -> Result<Option<sessionwiki::index::SessionRow>> {
1330 Ok(sessionwiki::index::resolve(connection, id)
1331 .context("look up an indexed session")?
1332 .into_iter()
1333 .find(|row| row.session_id == id))
1334}
1335
1336fn wiki_row(
1337 row: sessionwiki::index::SessionRow,
1338 snippet: Option<String>,
1339 live: &BTreeSet<String>,
1340) -> WikiRow {
1341 let hel_session_id = (row.tool == TOOL)
1344 .then(|| row.path.rsplit('/').next().unwrap_or_default().to_owned())
1345 .filter(|session_id| live.contains(session_id));
1346 let native_id = sessionwiki::index::native_id_of(&row.path);
1347 WikiRow {
1348 id: row.session_id,
1349 tool: row.tool,
1350 project: row.project,
1351 title: row.title,
1352 started: row.started,
1353 msgs: row.msg_count,
1354 preview: row.preview,
1355 archived: row.archived,
1356 native_id,
1357 snippet,
1358 hel_session_id,
1359 target: None,
1362 profile: None,
1363 harness: None,
1364 }
1365}
1366
1367fn project_directory_of(project: &str) -> Option<PathBuf> {
1375 if project.trim().is_empty() {
1376 return None;
1377 }
1378 let path = PathBuf::from(project);
1379 let repository = path
1380 .ancestors()
1381 .find(|ancestor| ancestor.file_name().is_some_and(|name| name == ".mj"))
1382 .and_then(std::path::Path::parent)
1383 .map(std::path::Path::to_path_buf)
1384 .unwrap_or(path);
1385 repository.is_dir().then_some(repository)
1386}
1387
1388fn snapshot_of(
1395 session: &sessionwiki::model::Session,
1396) -> Result<mj_core::archive::CanonicalSessionSnapshot> {
1397 use mj_core::archive::{
1398 CanonicalExecutionState, CanonicalSessionSnapshot, CanonicalSessionState,
1399 CanonicalTranscriptBody, CanonicalTranscriptItem,
1400 };
1401
1402 let started_ms = session
1403 .started
1404 .map(|time| time.timestamp_millis())
1405 .unwrap_or_default();
1406 let mut transcript: Vec<CanonicalTranscriptItem> = Vec::new();
1407 for message in &session.messages {
1408 let text = message.text.trim();
1409 if text.is_empty() {
1410 continue;
1411 }
1412 if transcript.is_empty() && message.role != Role::User {
1415 continue;
1416 }
1417 let position = transcript.len() as u64 + 1;
1418 let body = match message.role {
1419 Role::User => CanonicalTranscriptBody::User {
1420 content: vec![serde_json::json!({"type": "text", "text": text})],
1421 },
1422 Role::Assistant => CanonicalTranscriptBody::Agent {
1423 chunks: vec![serde_json::json!({
1424 "content": {"type": "text", "text": text}
1425 })],
1426 streaming: false,
1427 },
1428 Role::Tool => {
1430 let (call, terminal_outputs) = mj_transcript::summary::indexed_tool_call(
1431 text,
1432 &format!("wiki-tool-{position}"),
1433 );
1434 CanonicalTranscriptBody::Tool {
1435 call,
1436 terminal_outputs,
1437 terminal_refs: Vec::new(),
1438 presentation: None,
1439 }
1440 }
1441 };
1442 let created_at_ms = message
1443 .ts
1444 .map(|time| time.timestamp_millis())
1445 .unwrap_or(started_ms);
1446 transcript.push(CanonicalTranscriptItem {
1447 stable_id: format!("wiki-{position}"),
1448 position,
1449 latest_content_event_ordinal: matches!(body, CanonicalTranscriptBody::Agent { .. })
1452 .then_some(position),
1453 created_at_ms,
1454 last_changed_at_ms: created_at_ms,
1455 body,
1456 });
1457 }
1458 anyhow::ensure!(
1459 !transcript.is_empty(),
1460 "the archived session has no prompt to restore from"
1461 );
1462
1463 let event_frontier = transcript.len() as u64;
1464 let last_activity_at_ms = transcript.last().map(|item| item.last_changed_at_ms);
1465 Ok(CanonicalSessionSnapshot {
1466 event_frontier,
1467 event_frontier_digest: {
1471 use sha2::Digest;
1472 mj_core::hex::lower_hex(sha2::Sha256::digest(
1473 format!("sessionwiki:{}", session.id).as_bytes(),
1474 ))
1475 },
1476 session: CanonicalSessionState {
1477 execution: CanonicalExecutionState::Idle,
1478 last_activity_at_ms,
1479 session_title: Some(session.title.clone()).filter(|title| !title.trim().is_empty()),
1480 configuration: BTreeMap::new(),
1481 },
1482 transcript,
1483 queued_prompts: Vec::new(),
1484 })
1485}
1486
1487#[derive(Debug, Clone, PartialEq, Eq)]
1497pub enum WikiContinuation {
1498 Resume { session_id: String },
1500 Restore { wiki_id: String },
1503 Import {
1505 harness: HarnessKind,
1506 native_session_id: String,
1507 },
1508}
1509
1510pub fn wiki_continuation(
1516 wiki_id: &str,
1517 tool: &str,
1518 path: &Path,
1519 has_record: bool,
1520) -> Result<WikiContinuation> {
1521 if tool == TOOL {
1522 return Ok(match has_record {
1525 true => WikiContinuation::Resume {
1526 session_id: wiki_id.to_owned(),
1527 },
1528 false => WikiContinuation::Restore {
1529 wiki_id: wiki_id.to_owned(),
1530 },
1531 });
1532 }
1533 let harness = harness_adapters::harness_for_tool(tool)
1534 .with_context(|| format!("Mjolnir cannot continue a {tool} session"))?;
1535 let native_session_id = crate::import::native_session_id_from_path(harness, path)
1536 .with_context(|| {
1537 format!(
1538 "no {tool} session id in the indexed path {}",
1539 path.display()
1540 )
1541 })?;
1542 Ok(WikiContinuation::Import {
1543 harness,
1544 native_session_id,
1545 })
1546}
1547
1548pub fn wiki_session(
1553 wiki_id: &str,
1554 known_sessions: &BTreeSet<String>,
1555) -> Result<Option<WikiSessionInfo>> {
1556 if !index_is_writable() {
1557 return Ok(None);
1558 }
1559 let connection = open_readonly()?;
1560 let Some(row) = row_by_id(&connection, wiki_id)? else {
1561 return Ok(None);
1562 };
1563 let is_mjolnir = row.tool == TOOL;
1564 let mjolnir_session_id = is_mjolnir.then(|| row.session_id.clone());
1565 let has_record = mjolnir_session_id
1566 .as_deref()
1567 .is_some_and(|session_id| known_sessions.contains(session_id));
1568 let status = match (is_mjolnir, has_record) {
1569 (false, _) => WikiSessionStatus::Native,
1570 (true, true) => WikiSessionStatus::Mine,
1571 (true, false) => WikiSessionStatus::Archived,
1572 };
1573 let tags = match is_mjolnir {
1574 true => tags::read(&connection, &[row.session_id.as_str()])
1575 .context("read the indexed session metadata")?
1576 .remove(&row.session_id)
1577 .unwrap_or_default(),
1578 false => tags::MjTags::default(),
1579 };
1580 let harness = tags
1581 .harness
1582 .as_deref()
1583 .and_then(|id| id.parse::<HarnessKind>().ok())
1584 .or_else(|| {
1585 (!is_mjolnir)
1586 .then(|| harness_adapters::harness_for_tool(&row.tool))
1587 .flatten()
1588 });
1589 Ok(Some(WikiSessionInfo {
1590 wiki_id: row.session_id,
1591 tool: row.tool,
1592 path: PathBuf::from(row.path),
1593 status,
1594 mjolnir_session_id,
1595 profile_id: tags.profile,
1596 target_template_id: tags.target,
1597 harness,
1598 title: row.title,
1599 project: row.project,
1600 }))
1601}
1602
1603#[cfg(test)]
1604mod tests {
1605 use std::collections::BTreeMap;
1606 use std::path::Path;
1607
1608 mod continuation {
1611 use super::super::{WikiContinuation, wiki_continuation};
1612 use mj_core::config::HarnessKind;
1613 use std::path::Path;
1614
1615 #[test]
1616 fn a_mjolnir_row_with_a_record_is_resumed_and_one_without_is_restored() {
1617 let path = Path::new("/home/user/.local/share/mj/sessions/session-7");
1618 assert_eq!(
1619 wiki_continuation("session-7", "mjolnir", path, true).unwrap(),
1620 WikiContinuation::Resume {
1621 session_id: "session-7".to_owned(),
1622 }
1623 );
1624 assert_eq!(
1625 wiki_continuation("session-7", "mjolnir", path, false).unwrap(),
1626 WikiContinuation::Restore {
1627 wiki_id: "session-7".to_owned(),
1628 }
1629 );
1630 }
1631
1632 #[test]
1633 fn a_claude_code_row_is_imported_with_the_uuid_from_its_path() {
1634 let path = Path::new(
1635 "/home/user/.claude/projects/-home-user-app/7f3a1c20-0b11-4a55-9e0d-2c8a5d6f1b44.jsonl",
1636 );
1637 assert_eq!(
1638 wiki_continuation("abc123", "claude-code", path, false).unwrap(),
1639 WikiContinuation::Import {
1640 harness: HarnessKind::Claude,
1641 native_session_id: "7f3a1c20-0b11-4a55-9e0d-2c8a5d6f1b44".to_owned(),
1642 }
1643 );
1644 }
1645
1646 #[test]
1649 fn a_codex_row_is_imported_with_the_uuid_from_its_rollout_name() {
1650 let path = Path::new(
1651 "/home/user/.codex/sessions/2026/09/18/rollout-2026-09-18T09-15-00-7f3a1c20-0b11-4a55-9e0d-2c8a5d6f1b44.jsonl",
1652 );
1653 assert_eq!(
1654 wiki_continuation("abc123", "codex", path, false).unwrap(),
1655 WikiContinuation::Import {
1656 harness: HarnessKind::Codex,
1657 native_session_id: "7f3a1c20-0b11-4a55-9e0d-2c8a5d6f1b44".to_owned(),
1658 }
1659 );
1660 }
1661
1662 #[test]
1663 fn an_unknown_tool_is_an_error_that_names_it() {
1664 let error = wiki_continuation("abc123", "opencode", Path::new("/tmp/s.jsonl"), false)
1665 .unwrap_err();
1666 assert!(
1667 format!("{error:#}").contains("opencode"),
1668 "the error has to name the tool: {error:#}"
1669 );
1670 }
1671 }
1672
1673 use mj_checkpoint::archive::{
1674 ArchiveInput, BundleManifest, CanonicalExecutionState, CanonicalSessionSnapshot,
1675 CanonicalSessionState, CanonicalTranscriptBody, CanonicalTranscriptItem, SessionManifest,
1676 TargetManifest, write_archive_atomic,
1677 };
1678
1679 use super::*;
1680
1681 fn item(position: u64, body: CanonicalTranscriptBody) -> CanonicalTranscriptItem {
1682 let streamed = matches!(body, CanonicalTranscriptBody::Agent { .. });
1685 CanonicalTranscriptItem {
1686 stable_id: format!("item-{position}"),
1687 position,
1688 latest_content_event_ordinal: streamed.then_some(position),
1689 created_at_ms: 1_700_000_000_000 + i64::try_from(position).unwrap(),
1690 last_changed_at_ms: 1_700_000_000_000 + i64::try_from(position).unwrap(),
1691 body,
1692 }
1693 }
1694
1695 fn write_archive(directory: &Path, session_id: &str, frontier: u64) {
1698 let path = directory.join(format!(
1699 "{session_id}-{frontier}-archive-{}.hel.zip",
1700 "0".repeat(32)
1701 ));
1702 write_archive_atomic(
1703 &path,
1704 &ArchiveInput {
1705 session: SessionManifest {
1706 id: session_id.into(),
1707 title: "indexed session".into(),
1708 harness_kind: mj_core::config::HarnessKind::Codex,
1709 profile_id: "codex".into(),
1710 native_session_id: "native-session".into(),
1711 created_at: "2026-09-01T00:00:00Z".into(),
1712 checkpointed_at: "2026-09-01T01:00:00Z".into(),
1713 hel_version: "test".into(),
1714 relay_version: "test".into(),
1715 adapter_version: "test".into(),
1716 },
1717 target: TargetManifest {
1718 template_id: "local".into(),
1719 target_kind: "local-bare".into(),
1720 details: BTreeMap::new(),
1721 },
1722 bundle: BundleManifest {
1723 id: "project".into(),
1724 primary_repository: "project".into(),
1725 },
1726 canonical_session: CanonicalSessionSnapshot {
1727 event_frontier: 4,
1728 event_frontier_digest: "a".repeat(64),
1729 session: CanonicalSessionState {
1730 execution: CanonicalExecutionState::Idle,
1731 last_activity_at_ms: Some(1_700_000_000_004),
1732 session_title: Some("snapshot title".into()),
1733 configuration: BTreeMap::new(),
1734 },
1735 transcript: vec![
1736 item(
1737 1,
1738 CanonicalTranscriptBody::User {
1739 content: vec![serde_json::json!({
1740 "type": "text",
1741 "text": "index this session"
1742 })],
1743 },
1744 ),
1745 item(
1746 2,
1747 CanonicalTranscriptBody::Thought {
1748 chunks: vec![serde_json::json!({
1749 "content": {"type": "text", "text": "pondering"}
1750 })],
1751 streaming: false,
1752 },
1753 ),
1754 item(
1755 3,
1756 CanonicalTranscriptBody::Tool {
1757 call: serde_json::json!({
1758 "toolCallId": "call-1",
1759 "title": "Edit config.toml",
1760 "kind": "edit",
1761 "status": "completed",
1762 "locations": [{"path": "/old/container/config.toml"}]
1763 }),
1764 terminal_outputs: Vec::new(),
1765 terminal_refs: Vec::new(),
1766 presentation: None,
1767 },
1768 ),
1769 item(
1770 4,
1771 CanonicalTranscriptBody::Agent {
1772 chunks: vec![serde_json::json!({
1773 "content": {"type": "text", "text": "done"}
1774 })],
1775 streaming: false,
1776 },
1777 ),
1778 ],
1779 queued_prompts: Vec::new(),
1780 },
1781 native_artifacts: Vec::new(),
1782 repositories: Vec::new(),
1783 },
1784 )
1785 .unwrap();
1786 }
1787
1788 fn adapter(directory: &Path, session_id: &str) -> MjolnirAdapter {
1789 adapter_with_live(directory, session_id, BTreeMap::new())
1790 }
1791
1792 fn adapter_with_live(
1793 directory: &Path,
1794 session_id: &str,
1795 live: BTreeMap<String, i64>,
1796 ) -> MjolnirAdapter {
1797 let record = SessionRecord {
1798 id: session_id.into(),
1799 ..record_template()
1800 };
1801 MjolnirAdapter {
1802 sessions_dir: directory.to_path_buf(),
1803 sessions: std::sync::Mutex::new(Sessions {
1804 records: BTreeMap::from([(session_id.to_owned(), record)]),
1805 subagent_ids: BTreeSet::new(),
1806 live,
1807 }),
1808 reload: false,
1809 }
1810 }
1811
1812 fn record_template() -> SessionRecord {
1813 SessionRecord {
1814 launch_base: None,
1815 build_cache: None,
1816 container_workspace: None,
1817 mjolnir_subagents: None,
1818 create_managed_worktree: None,
1819 workspace_id: mj_core::workspace::DEFAULT_WORKSPACE_ID.to_owned(),
1820 archived: false,
1821 container_cpus: None,
1822 container_memory: None,
1823 id: "0123456789abcdef0123456789abcdef".into(),
1824 title: "indexed session".into(),
1825 harness_kind: mj_core::config::HarnessKind::Codex,
1826 last_profile: "codex".into(),
1827 bundle_id: "project".into(),
1828 project_directory: Some(PathBuf::from("/home/dev/project")),
1829 managed_worktree: None,
1830 target_template_id: "local-bare".into(),
1831 resource_allocation: None,
1832 additional_mounts: Vec::new(),
1833 state: mj_core::state::SessionState::Stopped,
1834 target: None,
1835 native_session_id: Some("native-session".into()),
1836 acp_session_title: Some("the harness title".into()),
1837 session_title_override: None,
1838 created_at: "2026-09-01T00:00:00Z".into(),
1839 updated_at: "2026-09-01T01:00:00Z".into(),
1840 viewed_through_event_ordinal: 0,
1841 draft_input: String::new(),
1842 last_error: None,
1843 last_checkpoint_error: None,
1844 checkpoint: None,
1845 }
1846 }
1847
1848 #[test]
1849 fn the_newest_checkpoint_of_each_session_is_one_indexed_key() {
1850 let directory = tempfile::tempdir().unwrap();
1851 let session_id = "0123456789abcdef0123456789abcdef";
1852 write_archive(directory.path(), session_id, 1);
1853 write_archive(directory.path(), session_id, 7);
1854 let adapter = adapter(directory.path(), session_id);
1855
1856 let store = adapter.store().expect("the adapter is a shared store");
1857 let key = format!("{}/{session_id}", directory.path().display());
1858 assert_eq!(
1859 store
1860 .keys
1861 .iter()
1862 .map(|(key, _)| key.as_str())
1863 .collect::<Vec<_>>(),
1864 vec![key.as_str()]
1865 );
1866 assert!(!store.had_error);
1867 assert_eq!(store.files.len(), 1);
1868 assert!(
1869 store.files[0]
1870 .file_name()
1871 .unwrap()
1872 .to_str()
1873 .unwrap()
1874 .contains("-7-archive-"),
1875 "the newest checkpoint is the one indexed: {:?}",
1876 store.files[0]
1877 );
1878 assert_eq!(
1879 adapter.reconcile_scope(),
1880 Some(format!("{}/", directory.path().display()))
1881 );
1882
1883 let session = adapter.parse_key(&key).unwrap();
1884 assert_eq!(session.id, session_id);
1885 assert_eq!(session.tool, "mjolnir");
1886 assert_eq!(session.path, PathBuf::from(&key));
1887 assert_eq!(session.project, "/home/dev/project");
1888 assert_eq!(session.title, "the harness title");
1889 assert!(!session.subagent);
1890 assert_eq!(
1891 session.messages.iter().map(|m| m.role).collect::<Vec<_>>(),
1892 vec![Role::User, Role::Tool, Role::Assistant]
1893 );
1894 assert_eq!(session.messages[0].text, "index this session");
1895 let tool: serde_json::Value = serde_json::from_str(&session.messages[1].text).unwrap();
1896 assert_eq!(tool["name"], "Edit");
1897 assert_eq!(tool["call"]["title"], "Edit config.toml");
1898 assert_eq!(session.messages[2].text, "done");
1899 assert_eq!(session.touched, vec!["/old/container/config.toml"]);
1900 }
1901
1902 #[test]
1903 fn provenance_backfill_repairs_an_unchanged_checkpoint_without_rebuilding_the_index() {
1904 let _held = tags::testing::lock();
1905 let (_index_dir, mut connection) = tags::testing::isolated_index();
1906 let directory = tempfile::tempdir().unwrap();
1907 write_archive(directory.path(), "old-session", 4);
1908 let source = adapter(directory.path(), "old-session");
1909 let key = source.key_for("old-session");
1910 tags::testing::index_row(&connection, "old-session", "mjolnir");
1911 connection
1912 .execute(
1913 "UPDATE files SET path = ?1 WHERE session_id = 'old-session'",
1914 [&key],
1915 )
1916 .unwrap();
1917 provenance::backfill(&mut connection, &source).unwrap();
1918 assert_eq!(
1919 sessionwiki::index::files_for(&connection, "old-session").unwrap(),
1920 vec!["/old/container/config.toml"]
1921 );
1922 provenance::backfill(&mut connection, &source).unwrap();
1923 assert_eq!(
1924 sessionwiki::index::sessions_for_file(&connection, "config.toml", 20)
1925 .unwrap()
1926 .len(),
1927 1
1928 );
1929 }
1930
1931 fn projection(session_id: &str) -> mj_core::state::MaterializedSession {
1932 use mj_core::transcript::{TranscriptBody, TranscriptItem};
1933 let mut projected = mj_core::state::MaterializedSession::empty(session_id);
1934 let mut push = |position: u64, body: TranscriptBody| {
1935 let streamed = matches!(body, TranscriptBody::Agent { .. });
1936 projected
1937 .transcript
1938 .push(std::sync::Arc::new(TranscriptItem {
1939 stable_id: format!("item-{position}"),
1940 position,
1941 latest_content_event_ordinal: streamed.then_some(position),
1942 created_at_ms: 1_700_000_000_000 + i64::try_from(position).unwrap(),
1943 last_changed_at_ms: 1_700_000_000_000 + i64::try_from(position).unwrap(),
1944 body,
1945 }));
1946 };
1947 push(
1948 1,
1949 TranscriptBody::User {
1950 content: vec![serde_json::json!({"type": "text", "text": "still talking"})],
1951 },
1952 );
1953 push(
1954 2,
1955 TranscriptBody::Thought {
1956 chunks: vec![serde_json::json!({"content": {"type": "text", "text": "hmm"}})],
1957 streaming: false,
1958 },
1959 );
1960 push(
1961 3,
1962 TranscriptBody::Tool {
1963 call: serde_json::json!({"toolCallId": "c1", "title": "Read README.md"}),
1964 terminal_outputs: Vec::new(),
1965 terminal_refs: Vec::new(),
1966 presentation: None,
1967 },
1968 );
1969 push(
1970 4,
1971 TranscriptBody::Agent {
1972 chunks: vec![serde_json::json!({"content": {"type": "text", "text": "reading"}})],
1973 streaming: false,
1974 },
1975 );
1976 projected.session_title = Some("the live title".into());
1977 projected
1978 }
1979
1980 #[test]
1983 fn a_running_session_is_indexed_from_its_stored_transcript() {
1984 let session_id = "0123456789abcdef0123456789abcdef";
1985 let messages = projected_messages(&projection(session_id));
1986 assert_eq!(
1987 messages.iter().map(|m| m.role).collect::<Vec<_>>(),
1988 vec![Role::User, Role::Tool, Role::Assistant]
1989 );
1990 assert_eq!(messages[0].text, "still talking");
1991 assert_eq!(messages[2].text, "reading");
1992 let tool: serde_json::Value = serde_json::from_str(&messages[1].text).unwrap();
1993 assert_eq!(tool["name"], "Read");
1994 assert_eq!(tool["call"]["title"], "Read README.md");
1995 }
1996
1997 #[test]
2002 fn a_running_session_is_listed_with_its_own_change_token() {
2003 let directory = tempfile::tempdir().unwrap();
2004 let running = "0123456789abcdef0123456789abcdef";
2005 let never_checkpointed = "fedcba9876543210fedcba9876543210";
2006 write_archive(directory.path(), running, 3);
2007 let live = adapter_with_live(
2008 directory.path(),
2009 running,
2010 BTreeMap::from([
2011 (running.to_owned(), 1_900_000_000),
2012 (never_checkpointed.to_owned(), 1_900_000_001),
2013 ]),
2014 );
2015
2016 let store = live.store().expect("the adapter is a shared store");
2017 let key_of = |session_id: &str| format!("{}/{session_id}", directory.path().display());
2018 assert_eq!(
2019 store.keys,
2020 vec![
2021 (
2022 key_of(running),
2023 1_900_000_000 * 1024 + i64::from(mj_transcript::summary::SUMMARY_VERSION)
2024 ),
2025 (
2026 key_of(never_checkpointed),
2027 1_900_000_001 * 1024 + i64::from(mj_transcript::summary::SUMMARY_VERSION)
2028 ),
2029 ],
2030 "a live session's own token replaces the checkpoint's"
2031 );
2032
2033 let stopped = adapter(directory.path(), running);
2036 let keys = stopped.store().expect("a shared store").keys;
2037 assert_eq!(keys.len(), 1);
2038 assert_eq!(keys[0].0, key_of(running));
2039 assert_ne!(keys[0].1, 1_900_000_000);
2040 assert_eq!(
2041 stopped.parse_key(&key_of(running)).unwrap().title,
2042 "the harness title",
2043 "a stopped session is parsed from its checkpoint"
2044 );
2045 }
2046
2047 #[test]
2050 fn a_rename_moves_a_session_change_token() {
2051 let directory = tempfile::tempdir().unwrap();
2052 let session_id = "0123456789abcdef0123456789abcdef";
2053 write_archive(directory.path(), session_id, 1);
2054 let adapter = adapter(directory.path(), session_id);
2055 let before = adapter.store().expect("a shared store").keys[0].1;
2056
2057 {
2058 let mut sessions = adapter.sessions.lock().unwrap();
2059 let record = sessions.records.get_mut(session_id).unwrap();
2060 record.session_title_override = Some("the new name".into());
2061 record.updated_at = "2099-01-01T00:00:00Z".into();
2062 }
2063 let after = adapter.store().expect("a shared store").keys[0].1;
2064 assert!(
2065 after > before,
2066 "a renamed session is re-indexed: {before} then {after}"
2067 );
2068 assert_eq!(
2069 adapter
2070 .parse_key(&format!("{}/{session_id}", directory.path().display()))
2071 .unwrap()
2072 .title,
2073 "the new name"
2074 );
2075 }
2076
2077 fn indexed(messages: Vec<(Role, &str)>) -> sessionwiki::model::Session {
2078 Session {
2079 id: "0123456789abcdef0123456789abcdef".into(),
2080 tool: "mjolnir",
2081 path: PathBuf::from("/sessions/0123456789abcdef0123456789abcdef"),
2082 project: "/home/dev/project".into(),
2083 started: DateTime::from_timestamp_millis(1_700_000_000_000),
2084 ended: None,
2085 title: "the archived session".into(),
2086 subagent: false,
2087 messages: messages
2088 .into_iter()
2089 .map(|(role, text)| Message {
2090 role,
2091 text: text.to_owned(),
2092 ts: None,
2093 })
2094 .collect(),
2095 touched: Vec::new(),
2096 edits: Vec::new(),
2097 }
2098 }
2099
2100 #[test]
2103 fn transcript_hits_locates_case_insensitive_matches() {
2104 let session = indexed(vec![
2105 (Role::User, "Make the Tests green"),
2106 (Role::Assistant, "the tests are green now"),
2107 ]);
2108
2109 let found = hit_transcript(&session, "TESTS", 0, 4_000);
2110
2111 assert_eq!(found.blocks.len(), 2, "both messages contain the query");
2112 assert_eq!(found.blocks[0].role, "user");
2113 let (start, end) = found.blocks[0].hits[0];
2114 assert_eq!(&found.blocks[0].text[start..end], "Tests");
2115 let (start, end) = found.blocks[1].hits[0];
2116 assert_eq!(&found.blocks[1].text[start..end], "tests");
2117 assert!(!found.blocks[0].truncated);
2118 assert_eq!(found.omitted_after, 0);
2119 }
2120
2121 #[test]
2124 fn transcript_hits_keeps_context_and_marks_omissions() {
2125 let session = indexed(vec![
2126 (Role::User, "zero"),
2127 (Role::Assistant, "one needle one"),
2128 (Role::Tool, "two"),
2129 (Role::User, "three"),
2130 (Role::Assistant, "four"),
2131 (Role::Tool, "five"),
2132 (Role::User, "six needle six"),
2133 (Role::Assistant, "seven"),
2134 (Role::User, "eight"),
2135 ]);
2136
2137 let found = hit_transcript(&session, "needle", 1, 4_000);
2138
2139 let shown: Vec<(&str, &str, usize)> = found
2140 .blocks
2141 .iter()
2142 .map(|block| {
2143 (
2144 block.role.as_str(),
2145 block.text.as_str(),
2146 block.omitted_before,
2147 )
2148 })
2149 .collect();
2150 assert_eq!(
2151 shown,
2152 vec![
2153 ("user", "zero", 0),
2154 ("assistant", "one needle one", 0),
2155 ("tool", "two", 0),
2156 ("tool", "five", 2),
2157 ("user", "six needle six", 0),
2158 ("assistant", "seven", 0),
2159 ]
2160 );
2161 assert_eq!(found.omitted_after, 1, "the last message is not shown");
2162 assert!(found.blocks[0].hits.is_empty(), "context has no hits");
2163 }
2164
2165 #[test]
2171 fn transcript_hits_never_anchor_on_tool_output() {
2172 let session = indexed(vec![
2173 (Role::User, "make it build"),
2174 (Role::Tool, "cargo build --needle"),
2175 (Role::Assistant, "it builds"),
2176 ]);
2177
2178 let only_in_a_tool = hit_transcript(&session, "needle", 1, 4_000);
2179 assert!(
2180 only_in_a_tool.blocks.is_empty(),
2181 "tool output must not anchor a passage, got {:?}",
2182 only_in_a_tool.blocks
2183 );
2184
2185 let beside_a_match = hit_transcript(&session, "builds", 1, 4_000);
2186 let shown: Vec<(&str, bool)> = beside_a_match
2187 .blocks
2188 .iter()
2189 .map(|block| (block.role.as_str(), !block.hits.is_empty()))
2190 .collect();
2191 assert_eq!(
2192 shown,
2193 vec![("tool", false), ("assistant", true)],
2194 "a tool message is still context around a real match"
2195 );
2196 }
2197
2198 #[test]
2201 fn transcript_hits_window_keeps_the_first_hit() {
2202 let filler = "x".repeat(4_000);
2203 let session = indexed(vec![(Role::User, &format!("{filler} needle {filler}"))]);
2204
2205 let found = hit_transcript(&session, "needle", 0, 100);
2206
2207 let block = &found.blocks[0];
2208 assert!(block.truncated);
2209 assert_eq!(block.text.chars().count(), 100);
2210 assert_eq!(block.hits.len(), 1, "the windowed text keeps its hit");
2211 let (start, end) = block.hits[0];
2212 assert_eq!(&block.text[start..end], "needle");
2213 assert!(
2214 start >= 20,
2215 "the window keeps lead-in before the hit, got {start}"
2216 );
2217 }
2218
2219 #[test]
2223 fn a_restored_snapshot_is_a_valid_transcript_of_the_indexed_session() {
2224 let snapshot = snapshot_of(&indexed(vec![
2225 (Role::User, "make the tests green"),
2226 (Role::Tool, "Read src/lib.rs"),
2227 (Role::Assistant, "they are green now"),
2228 (Role::User, " "),
2229 ]))
2230 .unwrap();
2231
2232 snapshot.validate().expect("the snapshot is well formed");
2233 assert_eq!(snapshot.event_frontier, 3);
2234 assert_eq!(
2235 snapshot.session.session_title.as_deref(),
2236 Some("the archived session")
2237 );
2238 assert!(snapshot.session.last_activity_at_ms.is_some());
2239 let bodies = snapshot
2240 .transcript
2241 .iter()
2242 .map(|item| match &item.body {
2243 mj_core::archive::CanonicalTranscriptBody::User { content } => (
2244 "user",
2245 mj_core::transcript::materialized_content_text(content),
2246 ),
2247 mj_core::archive::CanonicalTranscriptBody::Agent { chunks, .. } => (
2248 "agent",
2249 mj_core::transcript::materialized_chunks_text(chunks),
2250 ),
2251 mj_core::archive::CanonicalTranscriptBody::Tool { call, .. } => (
2252 "tool",
2253 call["title"].as_str().unwrap_or_default().to_owned(),
2254 ),
2255 _ => ("other", String::new()),
2256 })
2257 .collect::<Vec<_>>();
2258 assert_eq!(
2259 bodies,
2260 vec![
2261 ("user", "make the tests green".to_owned()),
2262 ("tool", "Read src/lib.rs".to_owned()),
2263 ("agent", "they are green now".to_owned()),
2264 ],
2265 "the blank message is dropped and every other one keeps its role"
2266 );
2267 }
2268
2269 #[test]
2273 fn messages_before_the_first_prompt_are_dropped() {
2274 let snapshot = snapshot_of(&indexed(vec![
2275 (Role::Assistant, "still working"),
2276 (Role::User, "carry on"),
2277 ]))
2278 .unwrap();
2279 assert_eq!(snapshot.transcript.len(), 1);
2280 assert_eq!(snapshot.transcript[0].position, 1);
2281 snapshot.validate().unwrap();
2282
2283 let error = snapshot_of(&indexed(vec![(Role::Assistant, "nobody asked")])).unwrap_err();
2284 assert!(
2285 error.to_string().contains("no prompt"),
2286 "a session with no prompt cannot be restored: {error}"
2287 );
2288 }
2289
2290 fn record(
2291 session_id: &str,
2292 state: mj_core::state::SessionState,
2293 updated_at: &str,
2294 ) -> SessionRecord {
2295 SessionRecord {
2296 id: session_id.into(),
2297 state,
2298 updated_at: updated_at.into(),
2299 ..record_template()
2300 }
2301 }
2302
2303 fn child(child_session_id: &str, parent_session_id: &str) -> mj_core::subagent::SubagentRecord {
2304 mj_core::subagent::SubagentRecord {
2305 child_session_id: child_session_id.into(),
2306 parent_session_id: parent_session_id.into(),
2307 task_name: "task".into(),
2308 profile_id: "codex".into(),
2309 model: None,
2310 effort: None,
2311 working_directory: PathBuf::new(),
2312 initial_prompt: "do the thing".into(),
2313 request_key: "key".into(),
2314 created_at: "2026-09-01T00:00:00Z".into(),
2315 noticed_turn: None,
2316 }
2317 }
2318
2319 fn ready(
2320 sessions: Vec<SessionRecord>,
2321 children: Vec<mj_core::subagent::SubagentRecord>,
2322 ) -> Vec<String> {
2323 let now = parse_time("2026-09-10T00:00:00Z").unwrap();
2324 sessions_ready_to_archive(
2325 &sessions
2326 .into_iter()
2327 .map(|record| (record.id.clone(), record))
2328 .collect(),
2329 &children
2330 .into_iter()
2331 .map(|child| (child.child_session_id.clone(), child))
2332 .collect(),
2333 now,
2334 3,
2335 )
2336 }
2337
2338 fn sized_session(
2340 root: &Path,
2341 session_id: &str,
2342 updated_at: &str,
2343 checkpoint_bytes: usize,
2344 attachment_bytes: &[usize],
2345 ) -> SessionRecord {
2346 let archive_path = root.join(format!("{session_id}.hel.zip"));
2347 std::fs::write(&archive_path, vec![b'c'; checkpoint_bytes]).unwrap();
2348 if !attachment_bytes.is_empty() {
2349 let attachments = root
2350 .join(session_id)
2351 .join(mj_core::attachment::ATTACHMENT_DIR);
2352 std::fs::create_dir_all(&attachments).unwrap();
2353 for (index, size) in attachment_bytes.iter().enumerate() {
2354 std::fs::write(attachments.join(format!("{index}.png")), vec![b'a'; *size])
2355 .unwrap();
2356 }
2357 }
2358 SessionRecord {
2359 checkpoint: Some(mj_core::state::CheckpointMetadata {
2360 archive_path,
2361 sha256: "0".repeat(64),
2362 created_at: updated_at.into(),
2363 event_frontier: 1,
2364 }),
2365 ..record(
2366 session_id,
2367 mj_core::state::SessionState::Stopped,
2368 updated_at,
2369 )
2370 }
2371 }
2372
2373 #[test]
2374 fn the_space_preview_sizes_every_session_and_only_the_aged_ones_as_reclaimable() {
2375 let directory = tempfile::tempdir().unwrap();
2376 let root = directory.path();
2377 let sessions: BTreeMap<String, SessionRecord> = [
2378 sized_session(root, "old-stopped", "2026-09-01T00:00:00Z", 1000, &[10, 20]),
2379 sized_session(root, "just-stopped", "2026-09-09T00:00:00Z", 500, &[]),
2380 SessionRecord {
2383 checkpoint: Some(mj_core::state::CheckpointMetadata {
2384 archive_path: root.join("missing.hel.zip"),
2385 sha256: "0".repeat(64),
2386 created_at: "2026-09-01T00:00:00Z".into(),
2387 event_frontier: 1,
2388 }),
2389 ..record(
2390 "lost-checkpoint",
2391 mj_core::state::SessionState::Stopped,
2392 "2026-09-01T00:00:00Z",
2393 )
2394 },
2395 ]
2396 .into_iter()
2397 .map(|record| (record.id.clone(), record))
2398 .collect();
2399 let now = parse_time("2026-09-10T00:00:00Z").unwrap();
2400
2401 let all = archive_space_over(root, &sessions, &BTreeMap::new(), now, None);
2402 assert_eq!(all.sessions, 3);
2403 assert_eq!(all.bytes, 1530);
2404 assert_eq!(all.reclaimable_sessions, 0);
2405 assert_eq!(all.reclaimable_bytes, 0);
2406
2407 let aged = archive_space_over(root, &sessions, &BTreeMap::new(), now, Some(3));
2408 assert_eq!(aged.bytes, 1530);
2409 assert_eq!(
2410 (aged.reclaimable_sessions, aged.reclaimable_bytes),
2411 (2, 1030),
2412 "only the sessions the job would archive count, attachments included"
2413 );
2414 }
2415
2416 #[test]
2417 fn only_stopped_sessions_past_the_cut_off_are_archived() {
2418 use mj_core::state::SessionState;
2419 let selected = ready(
2420 vec![
2421 record("old-stopped", SessionState::Stopped, "2026-09-01T00:00:00Z"),
2422 record(
2423 "just-stopped",
2424 SessionState::Stopped,
2425 "2026-09-09T00:00:00Z",
2426 ),
2427 record("old-running", SessionState::Running, "2026-09-01T00:00:00Z"),
2428 record("old-error", SessionState::Error, "2026-09-01T00:00:00Z"),
2429 record("unparsable", SessionState::Stopped, "not a time"),
2430 record("at-the-edge", SessionState::Stopped, "2026-09-07T00:00:00Z"),
2432 ],
2433 Vec::new(),
2434 );
2435 assert_eq!(selected, vec!["at-the-edge", "old-stopped"]);
2436 }
2437
2438 #[test]
2439 fn a_child_the_pass_is_not_archiving_holds_its_parent_back() {
2440 use mj_core::state::SessionState;
2441 let selected = ready(
2442 vec![
2443 record("parent", SessionState::Stopped, "2026-09-01T00:00:00Z"),
2444 record(
2445 "running-child",
2446 SessionState::Running,
2447 "2026-09-01T00:00:00Z",
2448 ),
2449 ],
2450 vec![child("running-child", "parent")],
2451 );
2452 assert!(selected.is_empty(), "the parent must wait: {selected:?}");
2453
2454 let selected = ready(
2455 vec![
2456 record("parent", SessionState::Stopped, "2026-09-01T00:00:00Z"),
2457 record("young-child", SessionState::Stopped, "2026-09-09T00:00:00Z"),
2458 ],
2459 vec![child("young-child", "parent")],
2460 );
2461 assert!(selected.is_empty(), "the parent must wait: {selected:?}");
2462
2463 let selected = ready(
2465 vec![record(
2466 "parent",
2467 SessionState::Stopped,
2468 "2026-09-01T00:00:00Z",
2469 )],
2470 vec![child("departed-child", "parent")],
2471 );
2472 assert_eq!(selected, vec!["parent"]);
2473 }
2474
2475 #[test]
2476 fn children_are_archived_before_their_parents() {
2477 use mj_core::state::SessionState;
2478 let selected = ready(
2479 vec![
2480 record("parent", SessionState::Stopped, "2026-09-01T00:00:00Z"),
2481 record("child", SessionState::Stopped, "2026-09-01T00:00:00Z"),
2482 record("grandchild", SessionState::Stopped, "2026-09-01T00:00:00Z"),
2483 ],
2484 vec![child("child", "parent"), child("grandchild", "child")],
2485 );
2486 assert_eq!(selected, vec!["grandchild", "child", "parent"]);
2487 }
2488
2489 #[test]
2490 fn native_adapters_cover_every_enabled_profile_home() {
2491 use mj_core::config::{Config, HarnessKind, HarnessProfile};
2492
2493 fn profile(kind: HarnessKind, home: &str, enabled: bool) -> HarnessProfile {
2494 HarnessProfile {
2495 enabled,
2496 kind,
2497 home: PathBuf::from(home),
2498 environment: BTreeMap::new(),
2499 context_window_bytes: None,
2500 guardian_review_model: None,
2501 }
2502 }
2503
2504 let mut config = Config::default();
2505 for (id, built) in [
2506 (
2507 "codex",
2508 profile(HarnessKind::Codex, "/home/dev/.codex3", true),
2509 ),
2510 (
2511 "codex-ds",
2512 profile(HarnessKind::Codex, "/home/dev/.codex-ds", true),
2513 ),
2514 (
2516 "codex-alt",
2517 profile(HarnessKind::Codex, "/home/dev/.codex3", true),
2518 ),
2519 (
2520 "codex-off",
2521 profile(HarnessKind::Codex, "/home/dev/.codex-off", false),
2522 ),
2523 (
2524 "claude",
2525 profile(HarnessKind::Claude, "/home/dev/.claude4", true),
2526 ),
2527 ("kimi", profile(HarnessKind::Kimi, "/home/dev/.kimi", true)),
2528 ("grok", profile(HarnessKind::Grok, "/home/dev/.grok", true)),
2529 ("muse", profile(HarnessKind::Muse, "/home/dev/muse", true)),
2530 (
2531 "muse-off",
2532 profile(HarnessKind::Muse, "/home/dev/muse-off", false),
2533 ),
2534 ] {
2535 config.profiles.insert(id.into(), built);
2536 }
2537
2538 let adapters = native_adapters(&config);
2539 let roots: Vec<(&str, Option<PathBuf>)> = adapters
2540 .iter()
2541 .map(|adapter| (adapter.name(), adapter.root()))
2542 .collect();
2543
2544 let codex: Vec<&Option<PathBuf>> = roots
2545 .iter()
2546 .filter(|(name, _)| *name == "codex")
2547 .map(|(_, root)| root)
2548 .collect();
2549 assert_eq!(
2550 codex,
2551 vec![
2552 &Some(PathBuf::from("/home/dev/.codex3/sessions")),
2553 &Some(PathBuf::from("/home/dev/.codex-ds/sessions")),
2554 ],
2555 "one adapter per enabled Codex home, deduplicated: {roots:?}"
2556 );
2557
2558 let claude: Vec<&Option<PathBuf>> = roots
2559 .iter()
2560 .filter(|(name, _)| *name == "claude-code")
2561 .map(|(_, root)| root)
2562 .collect();
2563 assert_eq!(
2564 claude,
2565 vec![&Some(PathBuf::from("/home/dev/.claude4/projects"))],
2566 "one adapter for the enabled Claude home: {roots:?}"
2567 );
2568
2569 for (_, root) in &roots {
2570 let Some(root) = root else { continue };
2571 let text = root.to_string_lossy();
2572 assert!(
2573 !text.contains(".codex-off"),
2574 "a disabled profile must not be indexed: {roots:?}"
2575 );
2576 assert!(
2577 !text.ends_with("/.codex/sessions") && !text.ends_with("/.claude/projects"),
2578 "the stock homes are not indexed unless a profile names them: {roots:?}"
2579 );
2580 }
2581
2582 for (name, root) in [
2585 ("kimi-code", PathBuf::from("/home/dev/.kimi/sessions")),
2586 ("grok-build", PathBuf::from("/home/dev/.grok/sessions")),
2587 (
2588 "muse",
2589 mj_checkpoint::native::muse_sessions_root(Path::new("/home/dev/muse")).unwrap(),
2590 ),
2591 ] {
2592 let found: Vec<&Option<PathBuf>> = roots
2593 .iter()
2594 .filter(|(found, _)| *found == name)
2595 .map(|(_, root)| root)
2596 .collect();
2597 assert_eq!(found, vec![&Some(root)], "one {name} adapter: {roots:?}");
2598 }
2599
2600 for (_, root) in &roots {
2601 let Some(root) = root else { continue };
2602 assert!(
2603 !root.to_string_lossy().contains("muse-off"),
2604 "a disabled profile must not be indexed: {roots:?}"
2605 );
2606 }
2607
2608 assert!(
2609 roots.iter().any(|(name, _)| *name == "gemini"),
2610 "the other built-in adapters are kept: {roots:?}"
2611 );
2612 }
2613
2614 #[test]
2618 fn query_rows_returns_the_indexed_target_profile_and_harness() {
2619 let _held = tags::testing::lock();
2620 let (_directory, connection) = tags::testing::isolated_index();
2621 tags::testing::index_row(&connection, "mj-session", TOOL);
2622 tags::testing::index_row(&connection, "codex-session", "codex");
2623 tags::write(
2624 &connection,
2625 "mj-session",
2626 &tags::MjTags {
2627 target: Some("Prod-Box".into()),
2628 profile: Some("codex-Main".into()),
2629 harness: Some("codex".into()),
2630 },
2631 )
2632 .expect("write the session metadata");
2633
2634 let rows = query_rows("", 10, &BTreeSet::new(), false).expect("query the index");
2635 let mjolnir = rows
2636 .iter()
2637 .find(|row| row.id == "mj-session")
2638 .expect("the Mjolnir row is returned");
2639 assert_eq!(mjolnir.target.as_deref(), Some("Prod-Box"));
2640 assert_eq!(mjolnir.profile.as_deref(), Some("codex-Main"));
2641 assert_eq!(mjolnir.harness.as_deref(), Some("codex"));
2642
2643 let codex = rows
2644 .iter()
2645 .find(|row| row.id == "codex-session")
2646 .expect("the Codex row is returned");
2647 assert_eq!(codex.target, None);
2648 assert_eq!(codex.profile, None);
2649 assert_eq!(codex.harness, None);
2650 }
2651
2652 #[test]
2653 fn one_flag_keeps_sub_agents_out_of_every_query_path() {
2654 let _held = tags::testing::lock();
2655 let (_directory, connection) = tags::testing::isolated_index();
2656 for (session_id, kind) in [("main-session", "main"), ("sub-session", "sub")] {
2657 tags::testing::index_row(&connection, session_id, "claude");
2658 connection
2659 .execute(
2660 "UPDATE files SET kind = ?2 WHERE session_id = ?1",
2661 rusqlite::params![session_id, kind],
2662 )
2663 .expect("set the session kind");
2664 connection
2665 .execute(
2666 "INSERT INTO messages(session_id, role, text)
2667 VALUES (?1, 'user', 'fix the bridge derivation zq')",
2668 [session_id],
2669 )
2670 .expect("insert a message");
2671 connection
2672 .execute(
2673 "INSERT INTO msgs(rowid, text) VALUES (?1, 'fix the bridge derivation zq')",
2674 [connection.last_insert_rowid()],
2675 )
2676 .expect("index the message");
2677 }
2678 let ids = |query: &str, include_subagents: bool| {
2679 let mut ids: Vec<String> = query_rows(query, 10, &BTreeSet::new(), include_subagents)
2680 .expect("query the index")
2681 .into_iter()
2682 .map(|row| row.id)
2683 .collect();
2684 ids.sort();
2685 ids
2686 };
2687
2688 for query in ["", "bridge derivation", "zq", "an indexed session"] {
2690 assert_eq!(ids(query, false), ["main-session"], "query {query:?}");
2691 assert_eq!(
2692 ids(query, true),
2693 ["main-session", "sub-session"],
2694 "query {query:?}"
2695 );
2696 }
2697 }
2698}