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