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::{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
45const SESSION_PARSE_FORMAT_VERSION: i64 = 1;
48const CHANGE_TOKEN_SLOT_BITS: u32 = 10;
49
50fn session_change_token(token: i64) -> i64 {
51 token
52 .saturating_mul(1_i64 << (CHANGE_TOKEN_SLOT_BITS * 2))
53 .saturating_add(
54 i64::from(mj_transcript::summary::SUMMARY_VERSION)
55 .saturating_mul(1_i64 << CHANGE_TOKEN_SLOT_BITS),
56 )
57 .saturating_add(SESSION_PARSE_FORMAT_VERSION)
58}
59
60struct ArchiveFile {
62 path: PathBuf,
63 frontier: u64,
64 token: i64,
66}
67
68#[derive(Default)]
72struct Sessions {
73 records: mj_core::snapshot_map::SnapshotMap<String, SessionRecord>,
74 ownership: top_level::Snapshot,
75 project_directories: BTreeMap<String, std::result::Result<Option<PathBuf>, String>>,
79 live: BTreeMap<String, i64>,
81}
82
83impl Sessions {
84 fn of(state: &State, config: Option<&Config>) -> Self {
85 Self {
86 records: state.sessions.clone(),
87 ownership: top_level::Snapshot::from_state(state),
88 project_directories: project_directories_of(state, config),
89 live: live_tokens(state),
90 }
91 }
92}
93
94fn project_directories_of(
95 state: &State,
96 config: Option<&Config>,
97) -> BTreeMap<String, std::result::Result<Option<PathBuf>, String>> {
98 state
99 .sessions
100 .iter()
101 .map(|(session_id, record)| {
102 (
103 session_id.clone(),
104 indexed_project_directory(state, config, session_id, record)
105 .map_err(|error| format!("{error:#}")),
106 )
107 })
108 .collect()
109}
110
111fn indexed_project_directory(
114 state: &State,
115 config: Option<&Config>,
116 session_id: &str,
117 record: &SessionRecord,
118) -> Result<Option<PathBuf>> {
119 match state.checkout(session_id)?.effective() {
120 mj_core::state::Checkout::Attached { path } => Ok(Some(path.to_path_buf())),
121 mj_core::state::Checkout::ManagedWorktree {
122 project_directory, ..
123 } => Ok(project_directory.map(Path::to_path_buf)),
124 mj_core::state::Checkout::ManagedWorkspace => {
125 let Some(config) = config else {
126 return Ok(None);
127 };
128 let Some(bundle) = record.project_bundle(config) else {
129 return Ok(None);
130 };
131 let primary = bundle
132 .repositories
133 .iter()
134 .find(|repository| repository.id == bundle.primary_repo)
135 .with_context(|| {
136 format!(
137 "session bundle has no primary repository {:?}",
138 bundle.primary_repo
139 )
140 })?;
141 let workspace_root = match record.target.as_ref() {
142 Some(locator) => {
143 let backend = crate::controller::backend_locator(locator, record, config)
144 .context("resolve the session target for SessionWiki")?;
145 crate::controller::workspace_root(
146 &backend,
147 record.container_workspace.as_deref(),
148 )
149 }
150 None => mj_core::targets::container_workspace_root(
151 record.container_workspace.as_deref(),
152 ),
153 };
154 Ok(Some(PathBuf::from(
155 crate::server_runtime::api::agent_working_directory_at(
156 &workspace_root,
157 &primary.destination,
158 ),
159 )))
160 }
161 mj_core::state::Checkout::Borrowed { .. } => {
162 unreachable!("State::checkout resolves borrowed sessions")
163 }
164 }
165}
166
167fn live_tokens(state: &State) -> BTreeMap<String, i64> {
174 let activity = match crate::database::load_transcribed_session_activity() {
175 Ok(activity) => activity,
176 Err(error) => {
177 tracing::warn!(%error, "could not read session activity for SessionWiki");
178 return BTreeMap::new();
179 }
180 };
181 state
182 .sessions
183 .iter()
184 .filter(|(_, record)| record.state != mj_core::state::SessionState::Stopped)
185 .filter_map(|(session_id, _)| {
186 let watermark = activity.get(session_id)?;
187 Some((session_id.clone(), watermark.unwrap_or_default() / 1000))
188 })
189 .collect()
190}
191
192pub struct MjolnirAdapter {
194 sessions_dir: PathBuf,
195 sessions: std::sync::Mutex<Sessions>,
196 reload: bool,
198}
199
200impl MjolnirAdapter {
201 pub fn from_state(state: &State) -> Self {
204 Self {
205 sessions_dir: mj_core::config::sessions_dir(),
206 sessions: std::sync::Mutex::new(Sessions::of(state, None)),
207 reload: false,
208 }
209 }
210
211 pub(crate) fn from_state_with_config(state: &State, config: &Config) -> Self {
214 Self {
215 sessions_dir: mj_core::config::sessions_dir(),
216 sessions: std::sync::Mutex::new(Sessions::of(state, Some(config))),
217 reload: false,
218 }
219 }
220
221 pub fn reloading(state: &State) -> Self {
230 Self {
231 reload: true,
232 ..Self::from_state(state)
233 }
234 }
235
236 pub fn indexed_tags(&self) -> BTreeMap<String, tags::MjTags> {
245 let sessions = self
246 .sessions
247 .lock()
248 .unwrap_or_else(std::sync::PoisonError::into_inner);
249 sessions
250 .records
251 .iter()
252 .filter(|(id, _)| !sessions.ownership.owns_session(id))
253 .map(|(session_id, record)| {
254 (
255 session_id.clone(),
256 tags::MjTags {
257 target: Some(record.target_template_id.clone()).filter(|id| !id.is_empty()),
258 profile: Some(record.last_profile.clone()).filter(|id| !id.is_empty()),
259 harness: Some(record.harness_kind.id().to_owned()),
260 },
261 )
262 })
263 .collect()
264 }
265
266 fn reload(&self) {
267 if !self.reload {
268 return;
269 }
270 match Controller::load() {
271 Ok(controller) => {
272 *self
273 .sessions
274 .lock()
275 .unwrap_or_else(std::sync::PoisonError::into_inner) =
276 Sessions::of(&controller.state, Some(&controller.config))
277 }
278 Err(error) => {
279 tracing::warn!(%error, "could not refresh session records for SessionWiki")
280 }
281 }
282 }
283
284 fn checkpointed_transcript(&self, session_id: &str) -> Result<IndexedTranscript> {
287 let (newest, _) = self.newest_archives();
288 let archive = newest
289 .get(session_id)
290 .with_context(|| format!("no checkpoint archive for session {session_id}"))?;
291 let snapshot = mj_checkpoint::archive::read_archive_verified(&archive.path)
292 .with_context(|| format!("read checkpoint {}", archive.path.display()))?
293 .canonical_session()
294 .with_context(|| format!("read the transcript of session {session_id}"))?;
295 let mut evidence = provenance::Evidence::default();
296 for item in &snapshot.transcript {
297 if let mj_core::archive::CanonicalTranscriptBody::Tool { call, .. } = &item.body {
298 evidence.observe(call, item.created_at_ms);
299 }
300 }
301 let messages = summary_messages(mj_transcript::summary::TranscriptSummary::from_snapshot(
302 &snapshot,
303 ));
304 Ok(IndexedTranscript {
305 messages,
306 title: snapshot.session.session_title.clone(),
307 evidence,
308 })
309 }
310
311 fn projected_transcript(&self, session_id: &str) -> Result<IndexedTranscript> {
315 let projection = crate::database::load_materialized_session(session_id)
316 .with_context(|| format!("read the stored transcript of session {session_id}"))?
317 .with_context(|| format!("no stored transcript for session {session_id}"))?;
318 let mut evidence = provenance::Evidence::default();
319 for item in &projection.transcript {
320 if let mj_core::state::TranscriptBody::Tool { call, .. } = &item.body {
321 evidence.observe(call, item.created_at_ms);
322 }
323 }
324 Ok(IndexedTranscript {
325 messages: projected_messages(&projection),
326 title: projection.session_title.clone(),
327 evidence,
328 })
329 }
330
331 fn key_for(&self, session_id: &str) -> String {
334 format!("{}/{session_id}", self.sessions_dir.display())
335 }
336
337 fn newest_archives(&self) -> (BTreeMap<String, ArchiveFile>, bool) {
343 let mut newest: BTreeMap<String, ArchiveFile> = BTreeMap::new();
344 let mut had_error = false;
345 let entries = match std::fs::read_dir(&self.sessions_dir) {
346 Ok(entries) => entries,
347 Err(error) => {
348 if self.sessions_dir.exists() {
349 tracing::debug!(
350 directory = %self.sessions_dir.display(),
351 %error,
352 "could not list the checkpoint directory for SessionWiki"
353 );
354 had_error = true;
355 }
356 return (newest, had_error);
357 }
358 };
359 for entry in entries {
360 let Ok(entry) = entry else {
361 had_error = true;
362 continue;
363 };
364 let Some((session_id, frontier)) = checkpoint_archive_session(&entry.file_name())
365 else {
366 continue;
367 };
368 let token = entry
369 .metadata()
370 .ok()
371 .and_then(|metadata| metadata.modified().ok())
372 .and_then(|modified| modified.duration_since(std::time::UNIX_EPOCH).ok())
373 .map(|age| age.as_secs() as i64)
374 .unwrap_or(0);
375 let candidate = ArchiveFile {
376 path: entry.path(),
377 frontier,
378 token,
379 };
380 match newest.get(&session_id) {
381 Some(existing) if existing.frontier >= candidate.frontier => {}
382 _ => {
383 newest.insert(session_id, candidate);
384 }
385 }
386 }
387 (newest, had_error)
388 }
389}
390
391struct IndexedTranscript {
392 messages: Vec<Message>,
393 title: Option<String>,
394 evidence: provenance::Evidence,
395}
396
397fn checkpoint_archive_session(name: &std::ffi::OsStr) -> Option<(String, u64)> {
402 if let Some(parsed) = managed_checkpoint_archive_name(name) {
403 return Some((parsed.session_id, parsed.frontier));
404 }
405 let stem = name
406 .to_str()
407 .and_then(|name| name.strip_suffix(".hel.zip"))?;
408 mj_core::config::validate_id("session", stem)
409 .is_ok()
410 .then(|| (stem.to_owned(), 0))
411}
412
413fn projected_messages(projection: &mj_core::state::MaterializedSession) -> Vec<Message> {
415 summary_messages(mj_transcript::summary::TranscriptSummary::from_materialized(projection))
416}
417
418fn summary_messages(summary: mj_transcript::summary::TranscriptSummary) -> Vec<Message> {
419 use mj_transcript::summary::SummaryRole;
420 summary
421 .entries
422 .into_iter()
423 .filter_map(|entry| {
424 let role = match entry.role {
425 SummaryRole::User => Role::User,
426 SummaryRole::Assistant => Role::Assistant,
427 SummaryRole::Tool => Role::Tool,
428 SummaryRole::Plan => return None,
429 };
430 message(role, entry.body(), entry.created_at_ms)
431 })
432 .collect()
433}
434
435fn message(role: Role, text: String, created_at_ms: i64) -> Option<Message> {
437 let text = text.trim().to_owned();
438 (!text.is_empty()).then(|| Message {
439 role,
440 text,
441 ts: DateTime::from_timestamp_millis(created_at_ms),
442 })
443}
444
445fn parse_time(value: &str) -> Option<DateTime<Utc>> {
446 DateTime::parse_from_rfc3339(value)
447 .ok()
448 .map(|time| time.with_timezone(&Utc))
449}
450
451impl Adapter for MjolnirAdapter {
452 fn name(&self) -> &'static str {
453 TOOL
454 }
455
456 fn root(&self) -> Option<PathBuf> {
457 Some(self.sessions_dir.clone())
458 }
459
460 fn discover(&self) -> Discovered {
463 Discovered {
464 files: Vec::new(),
465 had_error: false,
466 }
467 }
468
469 fn parse(&self, _path: &Path) -> Result<Session> {
470 anyhow::bail!("Mjolnir sessions are parsed by key, not by file")
471 }
472
473 fn store(&self) -> Option<Store> {
474 self.reload();
475 let (newest, had_error) = self.newest_archives();
476 let mut files = Vec::with_capacity(newest.len());
477 let mut tokens: BTreeMap<String, i64> = BTreeMap::new();
478 let sessions = self
479 .sessions
480 .lock()
481 .unwrap_or_else(std::sync::PoisonError::into_inner);
482 for (session_id, archive) in newest {
483 if sessions.ownership.owns_session(&session_id) {
484 continue;
485 }
486 tokens.insert(session_id, archive.token);
487 files.push(archive.path);
488 }
489 tokens.extend(
494 sessions
495 .live
496 .iter()
497 .filter(|(id, _)| !sessions.ownership.owns_session(id))
498 .map(|(id, token)| (id.clone(), *token)),
499 );
500 for (session_id, token) in tokens.iter_mut() {
505 let updated = sessions
506 .records
507 .get(session_id)
508 .and_then(|record| parse_time(&record.updated_at))
509 .map(|updated| updated.timestamp());
510 if let Some(updated) = updated {
511 *token = (*token).max(updated);
512 }
513 }
514 let keys = tokens
515 .into_iter()
516 .map(|(session_id, token)| (self.key_for(&session_id), session_change_token(token)))
517 .collect();
518 Some(Store {
519 keys,
520 files,
521 had_error,
522 })
523 }
524
525 fn reconcile_scope(&self) -> Option<String> {
529 Some(format!("{}/", self.sessions_dir.display()))
530 }
531
532 fn parse_key(&self, key: &str) -> Result<Session> {
533 let session_id = key.rsplit('/').next().unwrap_or_default();
534 anyhow::ensure!(!session_id.is_empty(), "no session id in key {key:?}");
535 let sessions = self
536 .sessions
537 .lock()
538 .unwrap_or_else(std::sync::PoisonError::into_inner);
539 anyhow::ensure!(
540 !sessions.ownership.owns_session(session_id),
541 "sub-agent sessions are not indexed"
542 );
543 let IndexedTranscript {
544 messages,
545 title: snapshot_title,
546 evidence,
547 } = if sessions.live.contains_key(session_id) {
548 self.projected_transcript(session_id)?
549 } else {
550 self.checkpointed_transcript(session_id)?
551 };
552 let record = sessions.records.get(session_id);
553 let project = match sessions.project_directories.get(session_id) {
554 Some(Ok(Some(directory))) => directory.display().to_string(),
555 Some(Ok(None)) => String::new(),
556 Some(Err(error)) => {
557 anyhow::bail!("resolve the session project directory: {error}")
558 }
559 None => record
560 .and_then(|record| record.checkout().project_directory())
561 .map(|directory| directory.display().to_string())
562 .unwrap_or_default(),
563 };
564
565 let title = record
566 .and_then(|record| record.session_title_override.clone())
567 .or_else(|| record.and_then(|record| record.acp_session_title.clone()))
568 .or_else(|| snapshot_title.clone())
569 .unwrap_or_else(|| {
570 messages
571 .iter()
572 .find(|message| message.role == Role::User)
573 .map(|message| message.text.chars().take(80).collect())
574 .unwrap_or_default()
575 });
576
577 Ok(Session {
578 id: session_id.to_owned(),
579 tool: TOOL,
580 path: PathBuf::from(key),
581 project,
582 started: record.and_then(|record| parse_time(&record.created_at)),
583 ended: record.and_then(|record| parse_time(&record.updated_at)),
584 title,
585 subagent: sessions.ownership.owns_session(session_id),
586 messages,
587 touched: evidence.paths.into_iter().collect(),
588 edits: evidence.edits,
589 })
590 }
591}
592
593struct SharedMjolnirAdapter(Arc<MjolnirAdapter>);
601
602impl Adapter for SharedMjolnirAdapter {
603 fn name(&self) -> &'static str {
604 self.0.name()
605 }
606
607 fn root(&self) -> Option<PathBuf> {
608 self.0.root()
609 }
610
611 fn discover(&self) -> Discovered {
612 self.0.discover()
613 }
614
615 fn parse(&self, path: &Path) -> Result<Session> {
616 self.0.parse(path)
617 }
618
619 fn store(&self) -> Option<Store> {
620 self.0.store()
621 }
622
623 fn parse_key(&self, key: &str) -> Result<Session> {
624 self.0.parse_key(key)
625 }
626
627 fn reconcile_scope(&self) -> Option<String> {
628 self.0.reconcile_scope()
629 }
630}
631
632pub struct WikiIndexer {
639 inner: Arc<Indexer>,
640}
641
642#[derive(Default)]
643struct Indexer {
644 running: tokio::sync::Mutex<()>,
646 notify: tokio::sync::Notify,
647 requested: AtomicBool,
649 full_requested: AtomicBool,
651 in_flight: AtomicBool,
654 last_success: std::sync::Mutex<Option<Success>>,
655 native_scan_cache: crate::import::NativeScanCache,
656}
657
658#[derive(Clone, Copy)]
659struct Success {
660 at: Instant,
661 epoch_seconds: i64,
662}
663
664impl WikiIndexer {
665 pub fn spawn() -> Self {
668 let inner = Arc::new(Indexer {
669 native_scan_cache: crate::import::NativeScanCache::shared(),
670 ..Indexer::default()
671 });
672 if let Ok(handle) = tokio::runtime::Handle::try_current() {
673 let worker = Arc::clone(&inner);
674 handle.spawn(async move { worker.run().await });
675 }
676 Self { inner }
677 }
678
679 pub fn request_sync(&self, full: bool) {
681 if full {
682 self.inner.full_requested.store(true, Ordering::Release);
683 }
684 self.inner.requested.store(true, Ordering::Release);
685 self.inner.notify.notify_one();
686 }
687
688 #[cfg(test)]
691 pub(crate) fn inert() -> Self {
692 Self {
693 inner: Arc::new(Indexer::default()),
694 }
695 }
696
697 #[cfg(test)]
699 pub(crate) fn sync_requested(&self) -> bool {
700 self.inner.requested.load(Ordering::Acquire)
701 }
702
703 pub async fn sync_now(&self, full: bool) -> Result<()> {
705 self.inner.sync(full).await
706 }
707
708 pub fn status(&self) -> WikiStatus {
711 WikiStatus {
712 state: index_state(),
713 topping_up: self.inner.in_flight.load(Ordering::Acquire)
714 || self.inner.requested.load(Ordering::Acquire),
715 }
716 }
717
718 pub fn last_success(&self) -> Option<Instant> {
720 self.inner
721 .last_success
722 .lock()
723 .unwrap_or_else(std::sync::PoisonError::into_inner)
724 .map(|success| success.at)
725 }
726}
727
728impl Indexer {
729 async fn run(self: Arc<Self>) {
730 loop {
731 self.notify.notified().await;
732 while self.requested.swap(false, Ordering::AcqRel) {
733 let full = self.full_requested.swap(false, Ordering::AcqRel);
734 if let Err(error) = self.sync(full).await {
735 self.report(&error);
736 break;
741 }
742 }
743 }
744 }
745
746 fn report(&self, error: &anyhow::Error) {
750 if crate::database::is_busy_error(error) {
751 self.requested.store(true, Ordering::Release);
752 tracing::debug!(%error, "the SessionWiki index was busy; retrying on the next trigger");
753 } else {
754 tracing::warn!(%error, "could not sync sessions into SessionWiki");
755 }
756 }
757
758 async fn sync(&self, full: bool) -> Result<()> {
759 let _guard = self.running.lock().await;
760 let since = if full {
761 None
762 } else {
763 self.last_success
764 .lock()
765 .unwrap_or_else(std::sync::PoisonError::into_inner)
766 .map(|success| success.epoch_seconds - 60)
769 };
770 let started = Instant::now();
771 self.in_flight.store(true, Ordering::Release);
772 let cache = self.native_scan_cache.clone();
773 let ran = run_abandonable(move || sync_blocking(since, &cache)).await;
774 self.in_flight.store(false, Ordering::Release);
775 let ran = ran?;
776 if ran {
777 *self
778 .last_success
779 .lock()
780 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(Success {
781 at: started,
782 epoch_seconds: Utc::now().timestamp(),
783 });
784 }
785 Ok(())
786 }
787}
788
789async fn run_abandonable<T: Send + 'static>(
800 pass: impl FnOnce() -> Result<T> + Send + 'static,
801) -> Result<T> {
802 let (sender, receiver) = tokio::sync::oneshot::channel();
803 std::thread::Builder::new()
804 .name("sessionwiki-sync".to_owned())
805 .spawn(move || {
806 let _ = sender.send(pass());
809 })
810 .context("start the SessionWiki sync thread")?;
811 receiver
812 .await
813 .context("the SessionWiki sync thread stopped without an answer")?
814}
815
816fn sync_blocking(since: Option<i64>, cache: &crate::import::NativeScanCache) -> Result<bool> {
819 if !index_is_writable() {
820 return Ok(false);
821 }
822 mj_core::test_hooks::reach_test_hook("sessionwiki_sync_pass")?;
824 let controller =
825 Controller::load().context("load controller state for the SessionWiki sync")?;
826 let mjolnir = Arc::new(MjolnirAdapter::from_state_with_config(
831 &controller.state,
832 &controller.config,
833 ));
834 let owned: Vec<Box<dyn Adapter>> = vec![Box::new(SharedMjolnirAdapter(Arc::clone(&mjolnir)))];
835 let ownership = mjolnir
836 .sessions
837 .lock()
838 .unwrap_or_else(std::sync::PoisonError::into_inner)
839 .ownership
840 .clone();
841 let mut connection = sessionwiki::index::open().context("open the SessionWiki index")?;
842 top_level::prune(&mut connection, &BTreeSet::new(), &ownership)?;
843 sessionwiki::index::sync_with(&mut connection, &owned, since)
844 .context("sync Mjolnir sessions into SessionWiki")?;
845 let (native, excluded) = top_level::prepare(native_adapters(&controller.config), cache);
846 top_level::prune(&mut connection, &excluded, &ownership)?;
847 sessionwiki::index::sync_with(&mut connection, &native, since)
848 .context("sync native sessions into SessionWiki")?;
849 top_level::prune(&mut connection, &excluded, &ownership)?;
850 write_session_tags(&mut connection, &mjolnir.indexed_tags())
851 .context("store Mjolnir's session metadata in the SessionWiki index")?;
852 provenance::backfill(&mut connection, &mjolnir).context("backfill Mjolnir file provenance")?;
853 if since.is_none() {
854 record_first_build();
858 }
859 Ok(true)
860}
861
862fn write_session_tags(
872 connection: &mut rusqlite::Connection,
873 session_tags: &BTreeMap<String, tags::MjTags>,
874) -> Result<()> {
875 if session_tags.is_empty() {
876 return Ok(());
877 }
878 let transaction = connection
879 .transaction()
880 .context("open a transaction for the session metadata")?;
881 for (session_id, session) in session_tags {
882 if session.is_empty() {
883 continue;
884 }
885 tags::write(&transaction, session_id, session)?;
886 }
887 transaction
888 .commit()
889 .context("commit the session metadata")?;
890 Ok(())
891}
892
893fn native_adapters(config: &mj_core::config::Config) -> Vec<Box<dyn Adapter>> {
913 let mut seen: BTreeSet<(HarnessKind, &Path)> = BTreeSet::new();
916 let mut adapters: Vec<Box<dyn Adapter>> = Vec::new();
917 for (_, profile) in config.enabled_profiles() {
918 if !seen.insert((profile.kind, profile.home.as_path())) {
920 continue;
921 }
922 let adapter: Box<dyn Adapter> = match profile.kind {
923 HarnessKind::Codex => {
924 Box::new(sessionwiki::adapters::Codex::in_home(profile.home.clone()))
925 }
926 HarnessKind::Claude => Box::new(sessionwiki::adapters::ClaudeCode::in_home(
927 profile.home.clone(),
928 )),
929 kind => match HarnessAdapter::in_home(kind, profile.home.clone()) {
930 Some(adapter) => Box::new(adapter),
931 None => continue,
932 },
933 };
934 adapters.push(adapter);
935 }
936 adapters.extend(sessionwiki::adapters::all().into_iter().filter(|adapter| {
937 harness_adapters::harness_for_tool(adapter.name()) == Some(HarnessKind::OpenCode)
938 }));
939 adapters
940}
941
942fn index_is_isolated() -> bool {
955 static SAID: AtomicBool = AtomicBool::new(false);
956 if mj_core::config::session_index_is_resolved()
957 || std::env::var_os(mj_core::config::SESSION_INDEX_ENV).is_some()
958 {
959 return true;
960 }
961 if !SAID.swap(true, Ordering::AcqRel) {
962 tracing::debug!(
963 "this process did not resolve a session index location; SessionWiki is not used"
964 );
965 }
966 false
967}
968
969fn index_version_mismatch() -> bool {
978 static SAID: AtomicBool = AtomicBool::new(false);
979 let Ok(path) = sessionwiki::index::db_path() else {
980 return false;
981 };
982 if !path.exists() {
983 return false;
984 }
985 let version = rusqlite::Connection::open_with_flags(
986 &path,
987 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_URI,
988 )
989 .and_then(|connection| connection.pragma_query_value(None, "user_version", |row| row.get(0)));
990 let version: i64 = match version {
991 Ok(version) => version,
992 Err(error) => {
993 tracing::debug!(%error, "could not read the SessionWiki index schema version");
994 return false;
995 }
996 };
997 let mismatch = version != 0 && version != sessionwiki::index::SCHEMA_VERSION;
1000 if mismatch && !SAID.swap(true, Ordering::AcqRel) {
1001 tracing::warn!(
1002 found = version,
1003 expected = sessionwiki::index::SCHEMA_VERSION,
1004 path = %path.display(),
1005 "the SessionWiki index was written by another version; Mjolnir will not open it, because opening it would rebuild it. Install the matching sessionwiki command"
1006 );
1007 }
1008 mismatch
1009}
1010
1011fn index_is_writable() -> bool {
1012 index_is_isolated() && !index_version_mismatch()
1013}
1014
1015fn first_build_marker() -> PathBuf {
1018 mj_core::config::data_dir().join("sessionwiki-built")
1019}
1020
1021fn record_first_build() {
1022 let path = first_build_marker();
1023 let version = sessionwiki::index::SCHEMA_VERSION.to_string();
1024 if std::fs::read_to_string(&path).is_ok_and(|held| held.trim() == version) {
1025 return;
1026 }
1027 if let Err(error) = std::fs::write(&path, &version) {
1028 tracing::warn!(%error, path = %path.display(), "could not record the first SessionWiki build");
1029 }
1030}
1031
1032fn first_build_is_done() -> bool {
1034 std::fs::read_to_string(first_build_marker())
1035 .is_ok_and(|held| held.trim() == sessionwiki::index::SCHEMA_VERSION.to_string())
1036 && sessionwiki::index::db_path().is_ok_and(|path| path.exists())
1037}
1038
1039pub fn index_state() -> WikiIndexState {
1041 if !index_is_isolated() {
1042 return WikiIndexState::Indexing;
1043 }
1044 if index_version_mismatch() {
1045 return WikiIndexState::VersionMismatch;
1046 }
1047 if first_build_is_done() {
1048 WikiIndexState::Ready
1049 } else {
1050 WikiIndexState::Indexing
1051 }
1052}
1053
1054pub const DESTROY_SYNC_WAIT: Duration = Duration::from_secs(20);
1062
1063const DEFERRED_WRITE_RETRY: Duration = Duration::from_secs(5);
1067const DEFERRED_WRITE_LIMIT: Duration = Duration::from_secs(60 * 60);
1068
1069#[derive(Debug, Clone, PartialEq, Eq)]
1072pub enum IndexedBeforeDestroy {
1073 Current,
1076 Synced,
1078 WrittenDirectly,
1081 Deferred,
1084 Unavailable(&'static str),
1086 Failed(String),
1088}
1089
1090impl WikiIndexer {
1091 pub async fn index_before_destroy(
1101 &self,
1102 session_id: &str,
1103 wait: Duration,
1104 ) -> IndexedBeforeDestroy {
1105 if let Some(reason) = unwritable_reason() {
1106 return IndexedBeforeDestroy::Unavailable(reason);
1107 }
1108 let root = session_id.to_owned();
1109 let pending = match tokio::task::spawn_blocking(move || unindexed_session(&root)).await {
1110 Ok(Ok(pending)) => pending,
1111 Ok(Err(error)) => {
1112 return IndexedBeforeDestroy::Failed(format!(
1113 "could not tell whether the index holds the session: {error:#}"
1114 ));
1115 }
1116 Err(error) => {
1117 return IndexedBeforeDestroy::Failed(format!(
1118 "checking the index for the session stopped: {error}"
1119 ));
1120 }
1121 };
1122 if pending.is_empty() {
1123 return IndexedBeforeDestroy::Current;
1124 }
1125 let inner = Arc::clone(&self.inner);
1126 index_before_destroy_with(
1127 async move { inner.sync(false).await },
1128 wait,
1129 move || capture_sessions(&pending),
1130 DEFERRED_WRITE_RETRY,
1131 )
1132 .await
1133 }
1134}
1135
1136fn unwritable_reason() -> Option<&'static str> {
1138 if !index_is_isolated() {
1139 return Some("this process did not resolve a SessionWiki index of its own");
1140 }
1141 if index_version_mismatch() {
1142 return Some("the SessionWiki index was written by another SessionWiki version");
1143 }
1144 None
1145}
1146
1147async fn index_before_destroy_with<S, C>(
1150 sync: S,
1151 wait: Duration,
1152 capture: C,
1153 retry: Duration,
1154) -> IndexedBeforeDestroy
1155where
1156 S: std::future::Future<Output = Result<()>> + Send + 'static,
1157 C: FnOnce() -> Result<Vec<CapturedSession>> + Send + 'static,
1158{
1159 match tokio::time::timeout(wait, tokio::spawn(sync)).await {
1164 Ok(Ok(Ok(()))) => return IndexedBeforeDestroy::Synced,
1165 Ok(Ok(Err(error))) => tracing::warn!(
1166 error = %format!("{error:#}"),
1167 "the SessionWiki sync before a destroy failed; indexing the session on its own"
1168 ),
1169 Ok(Err(error)) => tracing::warn!(
1170 %error,
1171 "the SessionWiki sync before a destroy stopped; indexing the session on its own"
1172 ),
1173 Err(_) => tracing::info!(
1174 wait_seconds = wait.as_secs_f64(),
1175 "the SessionWiki sync did not finish in time; indexing the session on its own"
1176 ),
1177 }
1178 let captured = match tokio::task::spawn_blocking(capture).await {
1179 Ok(Ok(captured)) => Arc::new(captured),
1180 Ok(Err(error)) => {
1181 return IndexedBeforeDestroy::Failed(format!(
1182 "could not read the session to index it: {error:#}"
1183 ));
1184 }
1185 Err(error) => {
1186 return IndexedBeforeDestroy::Failed(format!(
1187 "reading the session to index it stopped: {error}"
1188 ));
1189 }
1190 };
1191 if captured.is_empty() {
1192 return IndexedBeforeDestroy::Current;
1193 }
1194 let attempt = Arc::clone(&captured);
1195 match tokio::task::spawn_blocking(move || write_captured(&attempt)).await {
1196 Ok(Ok(())) => IndexedBeforeDestroy::WrittenDirectly,
1197 Ok(Err(error)) if crate::database::is_busy_error(&error) => {
1198 write_captured_later(captured, retry);
1199 IndexedBeforeDestroy::Deferred
1200 }
1201 Ok(Err(error)) => IndexedBeforeDestroy::Failed(format!(
1202 "could not write the session into the index: {error:#}"
1203 )),
1204 Err(error) => IndexedBeforeDestroy::Failed(format!(
1205 "writing the session into the index stopped: {error}"
1206 )),
1207 }
1208}
1209
1210fn unindexed_session(root: &str) -> Result<Vec<String>> {
1212 let controller =
1213 Controller::load().context("load controller state to index a destroyed session")?;
1214 if !controller.state.sessions.contains_key(root) {
1215 return Ok(Vec::new());
1216 }
1217 unindexed(
1218 &MjolnirAdapter::from_state_with_config(&controller.state, &controller.config),
1219 &[root.to_owned()],
1220 )
1221}
1222
1223fn unindexed(adapter: &MjolnirAdapter, session_ids: &[String]) -> Result<Vec<String>> {
1227 let tokens: BTreeMap<String, i64> = adapter
1228 .store()
1229 .map(|store| store.keys.into_iter().collect())
1230 .unwrap_or_default();
1231 let connection = open_readonly().ok();
1233 let mut pending = Vec::new();
1234 for session_id in session_ids {
1235 let key = adapter.key_for(session_id);
1236 let Some(&token) = tokens.get(&key) else {
1238 continue;
1239 };
1240 let current = match &connection {
1241 Some(connection) => indexed_token(connection, &key)? == Some(token),
1242 None => false,
1243 };
1244 if !current {
1245 pending.push(session_id.clone());
1246 }
1247 }
1248 Ok(pending)
1249}
1250
1251fn indexed_token(connection: &rusqlite::Connection, key: &str) -> Result<Option<i64>> {
1254 use rusqlite::OptionalExtension;
1255 connection
1256 .query_row(
1257 "SELECT mtime FROM files WHERE path = ?1 AND archived_at IS NULL",
1258 [key],
1259 |row| row.get(0),
1260 )
1261 .optional()
1262 .context("read a session's change token from the SessionWiki index")
1263}
1264
1265struct CapturedSession {
1268 key: String,
1269 token: i64,
1270 session: Session,
1271 tags: tags::MjTags,
1272}
1273
1274fn capture_sessions(session_ids: &[String]) -> Result<Vec<CapturedSession>> {
1275 let controller =
1276 Controller::load().context("load controller state to index a destroyed session")?;
1277 capture_sessions_from(
1278 &MjolnirAdapter::from_state_with_config(&controller.state, &controller.config),
1279 session_ids,
1280 )
1281}
1282
1283fn capture_sessions_from(
1286 adapter: &MjolnirAdapter,
1287 session_ids: &[String],
1288) -> Result<Vec<CapturedSession>> {
1289 let tokens: BTreeMap<String, i64> = adapter
1290 .store()
1291 .map(|store| store.keys.into_iter().collect())
1292 .unwrap_or_default();
1293 let mut session_tags = adapter.indexed_tags();
1294 let mut captured = Vec::new();
1295 for session_id in session_ids {
1296 let key = adapter.key_for(session_id);
1297 let Some(&token) = tokens.get(&key) else {
1298 continue;
1299 };
1300 captured.push(CapturedSession {
1301 session: adapter.parse_key(&key)?,
1302 tags: session_tags.remove(session_id).unwrap_or_default(),
1303 key,
1304 token,
1305 });
1306 }
1307 Ok(captured)
1308}
1309
1310fn write_captured(captured: &Arc<Vec<CapturedSession>>) -> Result<()> {
1313 anyhow::ensure!(
1314 index_is_writable(),
1315 "this process may not write the SessionWiki index"
1316 );
1317 let mut connection = sessionwiki::index::open().context("open the SessionWiki index")?;
1318 for index in 0..captured.len() {
1319 let adapter: Box<dyn Adapter> = Box::new(CapturedAdapter {
1320 captured: Arc::clone(captured),
1321 index,
1322 });
1323 sessionwiki::index::sync_with(&mut connection, &[adapter], None)
1324 .context("index a session before it is destroyed")?;
1325 }
1326 let session_tags = captured
1327 .iter()
1328 .map(|captured| (captured.session.id.clone(), captured.tags.clone()))
1329 .collect();
1330 write_session_tags(&mut connection, &session_tags)
1331 .context("store Mjolnir's session metadata in the SessionWiki index")
1332}
1333
1334fn write_captured_later(captured: Arc<Vec<CapturedSession>>, retry: Duration) {
1339 let sessions = captured
1340 .iter()
1341 .map(|captured| captured.session.id.clone())
1342 .collect::<Vec<_>>();
1343 tracing::info!(
1344 ?sessions,
1345 "the SessionWiki index is busy; indexing the destroyed sessions once it is free"
1346 );
1347 tokio::spawn(async move {
1348 let started = Instant::now();
1349 loop {
1350 tokio::time::sleep(retry).await;
1351 let attempt = Arc::clone(&captured);
1352 let error = match tokio::task::spawn_blocking(move || write_captured(&attempt)).await {
1353 Ok(Ok(())) => {
1354 tracing::info!(?sessions, "indexed the destroyed sessions");
1355 return;
1356 }
1357 Ok(Err(error)) => error,
1358 Err(error) => anyhow::Error::new(error),
1359 };
1360 if !crate::database::is_busy_error(&error) || started.elapsed() >= DEFERRED_WRITE_LIMIT
1361 {
1362 tracing::warn!(
1363 ?sessions,
1364 error = %format!("{error:#}"),
1365 "gave up indexing destroyed sessions in SessionWiki"
1366 );
1367 return;
1368 }
1369 }
1370 });
1371}
1372
1373struct CapturedAdapter {
1376 captured: Arc<Vec<CapturedSession>>,
1377 index: usize,
1378}
1379
1380impl CapturedAdapter {
1381 fn captured(&self) -> &CapturedSession {
1382 &self.captured[self.index]
1383 }
1384}
1385
1386impl Adapter for CapturedAdapter {
1387 fn name(&self) -> &'static str {
1388 TOOL
1389 }
1390
1391 fn root(&self) -> Option<PathBuf> {
1392 Path::new(&self.captured().key)
1393 .parent()
1394 .map(Path::to_path_buf)
1395 }
1396
1397 fn discover(&self) -> Discovered {
1398 Discovered {
1399 files: Vec::new(),
1400 had_error: false,
1401 }
1402 }
1403
1404 fn parse(&self, _path: &Path) -> Result<Session> {
1405 anyhow::bail!("Mjolnir sessions are parsed by key, not by file")
1406 }
1407
1408 fn store(&self) -> Option<Store> {
1409 let captured = self.captured();
1410 Some(Store {
1411 keys: vec![(captured.key.clone(), captured.token)],
1412 files: Vec::new(),
1413 had_error: false,
1414 })
1415 }
1416
1417 fn parse_key(&self, key: &str) -> Result<Session> {
1418 let captured = self.captured();
1419 anyhow::ensure!(key == captured.key, "no captured session for key {key:?}");
1420 Ok(copy_session(&captured.session))
1421 }
1422
1423 fn reconcile_scope(&self) -> Option<String> {
1427 Some(format!("{}\0", self.captured().key))
1428 }
1429}
1430
1431fn copy_session(session: &Session) -> Session {
1434 Session {
1435 id: session.id.clone(),
1436 tool: session.tool,
1437 path: session.path.clone(),
1438 project: session.project.clone(),
1439 started: session.started,
1440 ended: session.ended,
1441 title: session.title.clone(),
1442 subagent: session.subagent,
1443 messages: session
1444 .messages
1445 .iter()
1446 .map(|message| Message {
1447 role: message.role,
1448 text: message.text.clone(),
1449 ts: message.ts,
1450 })
1451 .collect(),
1452 touched: session.touched.clone(),
1453 edits: session
1454 .edits
1455 .iter()
1456 .map(|edit| sessionwiki::model::EditEvent {
1457 path: edit.path.clone(),
1458 kind: edit.kind,
1459 snippet: edit.snippet.clone(),
1460 ts: edit.ts,
1461 })
1462 .collect(),
1463 }
1464}
1465
1466pub const MAX_WIKI_LIMIT: usize = 200;
1472pub const DEFAULT_WIKI_LIMIT: usize = 50;
1474const MIN_FULLTEXT_QUERY: usize = 3;
1477pub const SYNC_STALE_AFTER: std::time::Duration = std::time::Duration::from_secs(60);
1479
1480pub fn sync_is_stale(last_success: Option<Instant>) -> bool {
1482 last_success.is_none_or(|at| at.elapsed() >= SYNC_STALE_AFTER)
1483}
1484
1485pub fn query_rows(
1494 query: &str,
1495 limit: usize,
1496 live: &BTreeSet<String>,
1497 include_tool_matches: bool,
1498) -> Result<Vec<WikiRow>> {
1499 let limit = limit.clamp(1, MAX_WIKI_LIMIT);
1500 if !index_is_writable() {
1501 return Ok(Vec::new());
1505 }
1506 let connection = open_readonly()?;
1507 let ownership = top_level::current_snapshot()?;
1508 let cache = crate::import::NativeScanCache::shared();
1509 query_rows_from(
1510 &connection,
1511 query,
1512 limit,
1513 live,
1514 include_tool_matches,
1515 &ownership,
1516 &cache,
1517 )
1518}
1519
1520fn query_rows_from(
1521 connection: &rusqlite::Connection,
1522 query: &str,
1523 limit: usize,
1524 live: &BTreeSet<String>,
1525 include_tool_matches: bool,
1526 ownership: &top_level::Snapshot,
1527 cache: &crate::import::NativeScanCache,
1528) -> Result<Vec<WikiRow>> {
1529 let limit = limit.clamp(1, MAX_WIKI_LIMIT);
1530 let query = query.trim();
1531 if query.is_empty() {
1532 let rows = sessionwiki::index::recent(connection, MAX_WIKI_LIMIT, None, None, None, true)
1533 .context("list recent SessionWiki sessions")?;
1534 let mut visible = Vec::with_capacity(limit);
1535 for row in rows {
1536 if visible.len() >= limit {
1537 break;
1538 }
1539 if top_level::is_child(&row, ownership, cache)? {
1540 continue;
1541 }
1542 visible.push(wiki_row(row, None, live));
1543 }
1544 let mut rows = visible;
1545 fill_session_tags(connection, &mut rows)?;
1546 return Ok(rows);
1547 }
1548 let search_limit = MAX_WIKI_LIMIT;
1551 let hits = if query.chars().count() < MIN_FULLTEXT_QUERY {
1552 sessionwiki::index::search_like(connection, query, search_limit, None, None)
1553 } else {
1554 sessionwiki::index::search(connection, query, search_limit, None, None)
1555 }
1556 .context("search the SessionWiki index")?;
1557 let mut rows = Vec::with_capacity(hits.len().min(limit));
1559 for hit in hits {
1560 if rows.len() >= limit {
1561 break;
1562 }
1563 if top_level::is_child(&hit.row, ownership, cache)? {
1564 continue;
1565 }
1566 if !include_tool_matches
1573 && !matches!(hit.role.as_str(), "user" | "assistant")
1574 && !conversation_matches(connection, &hit.row, query)?
1575 {
1576 continue;
1577 }
1578 rows.push(wiki_row(hit.row, Some(hit.snippet), live));
1579 }
1580 if rows.len() < limit {
1584 let mut found: BTreeSet<String> = rows.iter().map(|row| row.id.clone()).collect();
1585 for row in named_like(connection, query)? {
1586 if rows.len() >= limit {
1587 break;
1588 }
1589 if found.contains(&row.session_id) || top_level::is_child(&row, ownership, cache)? {
1590 continue;
1591 }
1592 found.insert(row.session_id.clone());
1593 rows.push(wiki_row(row, None, live));
1594 }
1595 }
1596 fill_session_tags(connection, &mut rows)?;
1597 Ok(rows)
1598}
1599
1600const TEXT_SEARCH_MESSAGE_LIMIT: i64 = 20_000;
1604
1605pub fn session_text_matches(query: &str, live: &BTreeSet<String>) -> Result<Vec<SessionTextMatch>> {
1610 let query = query.trim();
1611 if query.is_empty() || !index_is_writable() {
1612 return Ok(Vec::new());
1613 }
1614 let connection = open_readonly()?;
1615 let ownership = top_level::current_snapshot()?;
1616 text_matches_in(&connection, query, live, &ownership)
1617}
1618
1619fn text_matches_in(
1620 connection: &rusqlite::Connection,
1621 query: &str,
1622 live: &BTreeSet<String>,
1623 ownership: &top_level::Snapshot,
1624) -> Result<Vec<SessionTextMatch>> {
1625 let mut statement;
1626 let rows = if query.chars().count() < MIN_FULLTEXT_QUERY {
1627 let pattern = format!(
1629 "%{}%",
1630 sessionwiki::util::nfc(query)
1631 .replace('\\', "\\\\")
1632 .replace('%', "\\%")
1633 .replace('_', "\\_")
1634 );
1635 statement = connection.prepare(
1636 "SELECT f.session_id, f.path, m.role, f.kind
1637 FROM messages m JOIN files f ON f.session_id = m.session_id
1638 WHERE f.tool = ?1 AND m.role IN ('user', 'assistant')
1639 AND m.text LIKE ?2 ESCAPE '\\'
1640 ORDER BY m.id DESC LIMIT ?3",
1641 )?;
1642 statement
1643 .query_map(
1644 rusqlite::params![TOOL, pattern, TEXT_SEARCH_MESSAGE_LIMIT],
1645 |row| {
1646 Ok((
1647 row.get::<_, String>(0)?,
1648 row.get::<_, String>(1)?,
1649 row.get::<_, String>(2)?,
1650 row.get::<_, String>(3)?,
1651 ))
1652 },
1653 )?
1654 .collect::<rusqlite::Result<Vec<_>>>()
1655 } else {
1656 let phrase = format!("\"{}\"", sessionwiki::util::nfc(query).replace('"', "\"\""));
1657 statement = connection.prepare(
1658 "SELECT f.session_id, f.path, m.role, f.kind
1659 FROM (SELECT rowid AS mid FROM msgs WHERE msgs MATCH ?2 LIMIT ?3) x
1660 JOIN messages m ON m.id = x.mid
1661 JOIN files f ON f.session_id = m.session_id
1662 WHERE f.tool = ?1 AND m.role IN ('user', 'assistant')",
1663 )?;
1664 statement
1665 .query_map(
1666 rusqlite::params![TOOL, phrase, TEXT_SEARCH_MESSAGE_LIMIT],
1667 |row| {
1668 Ok((
1669 row.get::<_, String>(0)?,
1670 row.get::<_, String>(1)?,
1671 row.get::<_, String>(2)?,
1672 row.get::<_, String>(3)?,
1673 ))
1674 },
1675 )?
1676 .collect::<rusqlite::Result<Vec<_>>>()
1677 }
1678 .context("search the indexed messages")?;
1679 let mut found = BTreeMap::<String, SessionTextMatchKind>::new();
1680 for (session_id, path, role, kind) in rows {
1681 if !live.contains(&session_id)
1682 || ownership.owns_indexed_child(&session_id, &path, TOOL, &kind)
1683 {
1684 continue;
1685 }
1686 let kind = if role == "user" {
1687 SessionTextMatchKind::User
1688 } else {
1689 SessionTextMatchKind::Agent
1690 };
1691 let entry = found.entry(session_id.clone()).or_insert(kind);
1692 *entry = (*entry).min(kind);
1693 }
1694 Ok(found
1695 .into_iter()
1696 .map(|(session_id, kind)| SessionTextMatch { session_id, kind })
1697 .collect())
1698}
1699
1700fn fill_session_tags(connection: &rusqlite::Connection, rows: &mut [WikiRow]) -> Result<()> {
1706 let ids: Vec<&str> = rows
1707 .iter()
1708 .filter(|row| row.tool == TOOL)
1709 .map(|row| row.id.as_str())
1710 .collect();
1711 let found = tags::read(connection, &ids).context("read the indexed session metadata")?;
1712 for row in rows.iter_mut().filter(|row| row.tool == TOOL) {
1713 let Some(session) = found.get(&row.id) else {
1714 continue;
1715 };
1716 row.target = session.target.clone();
1717 row.profile = session.profile.clone();
1718 row.harness = session.harness.clone();
1719 }
1720 Ok(())
1721}
1722
1723const NAME_SCAN_LIMIT: usize = 2_000;
1729
1730fn named_like(
1732 connection: &rusqlite::Connection,
1733 query: &str,
1734) -> Result<Vec<sessionwiki::index::SessionRow>> {
1735 let needle = query.to_lowercase();
1736 let sql = format!(
1737 "SELECT session_id, tool, path, project, title, started, msg_count, kind,
1738 archived_at IS NOT NULL
1739 FROM files ORDER BY started DESC LIMIT {NAME_SCAN_LIMIT}"
1740 );
1741 let mut statement = connection
1742 .prepare(&sql)
1743 .context("prepare recent SessionWiki metadata scan")?;
1744 let rows = statement
1745 .query_map([], |row| {
1746 Ok(sessionwiki::index::SessionRow {
1747 session_id: row.get(0)?,
1748 tool: row.get(1)?,
1749 path: row.get(2)?,
1750 project: row.get(3)?,
1751 title: row.get(4)?,
1752 started: row.get(5)?,
1753 msg_count: row.get(6)?,
1754 kind: row.get(7)?,
1755 preview: None,
1756 summary: None,
1757 tags: None,
1758 archived: row.get(8)?,
1759 account: None,
1760 })
1761 })
1762 .context("list recent SessionWiki metadata")?
1763 .collect::<rusqlite::Result<Vec<_>>>()
1764 .context("read recent SessionWiki metadata")?;
1765 Ok(rows
1766 .into_iter()
1767 .filter(|row| {
1768 row.title.to_lowercase().contains(&needle)
1769 || row.project.to_lowercase().contains(&needle)
1770 })
1771 .collect())
1772}
1773
1774fn conversation_matches(
1777 connection: &rusqlite::Connection,
1778 row: &sessionwiki::index::SessionRow,
1779 query: &str,
1780) -> Result<bool> {
1781 let session = sessionwiki::index::session_from_index(connection, row)
1782 .context("read an indexed session")?;
1783 Ok(!hit_transcript(&session, query, 0, 1).blocks.is_empty())
1784}
1785
1786pub fn brief(id: &str, max_chars: usize) -> Result<Option<String>> {
1788 if !index_is_writable() {
1789 return Ok(None);
1790 }
1791 let connection = open_readonly()?;
1792 let ownership = top_level::current_snapshot()?;
1793 let cache = crate::import::NativeScanCache::shared();
1794 let Some(row) = visible_row_by_id(&connection, id, &ownership, &cache)? else {
1795 return Ok(None);
1796 };
1797 let session = sessionwiki::index::session_from_index(&connection, &row)
1798 .context("read an indexed session")?;
1799 Ok(Some(sessionwiki::commands::brief_markdown(
1800 &session, max_chars, true,
1801 )))
1802}
1803
1804pub fn transcript_hits(
1812 id: &str,
1813 query: &str,
1814 context_messages: usize,
1815 per_message_chars: usize,
1816) -> Result<Option<WikiHitTranscript>> {
1817 if !index_is_writable() {
1818 return Ok(None);
1819 }
1820 let connection = open_readonly()?;
1821 let ownership = top_level::current_snapshot()?;
1822 let cache = crate::import::NativeScanCache::shared();
1823 let Some(row) = visible_row_by_id(&connection, id, &ownership, &cache)? else {
1824 return Ok(None);
1825 };
1826 let session = sessionwiki::index::session_from_index(&connection, &row)
1827 .context("read an indexed session")?;
1828 Ok(Some(hit_transcript(
1829 &session,
1830 query,
1831 context_messages,
1832 per_message_chars,
1833 )))
1834}
1835
1836fn hit_transcript(
1846 session: &Session,
1847 query: &str,
1848 context_messages: usize,
1849 per_message_chars: usize,
1850) -> WikiHitTranscript {
1851 let found = sessionwiki::grep::grep_session(
1852 session,
1853 query,
1854 &sessionwiki::grep::GrepOpts {
1855 context_messages,
1856 chars: per_message_chars,
1857 max_matches: None,
1858 anchor_roles: vec![Role::User, Role::Assistant],
1859 },
1860 );
1861 WikiHitTranscript {
1862 blocks: found
1863 .hits
1864 .into_iter()
1865 .map(|hit| WikiHitBlock {
1866 role: role_name(hit.role).to_owned(),
1867 text: hit.text,
1868 hits: hit.matches,
1869 omitted_before: hit.omitted_before,
1870 truncated: hit.truncated,
1871 })
1872 .collect(),
1873 omitted_after: found.omitted_after,
1874 }
1875}
1876
1877fn role_name(role: Role) -> &'static str {
1878 match role {
1879 Role::User => "user",
1880 Role::Assistant => "assistant",
1881 Role::Tool => "tool",
1882 }
1883}
1884
1885pub struct ArchivedSession {
1889 pub title: String,
1890 pub project_directory: Option<PathBuf>,
1893 pub snapshot: mj_core::archive::CanonicalSessionSnapshot,
1894}
1895
1896pub fn archived_session(id: &str) -> Result<Option<ArchivedSession>> {
1898 if !index_is_writable() {
1899 return Ok(None);
1900 }
1901 let connection = open_readonly()?;
1902 let ownership = top_level::current_snapshot()?;
1903 let cache = crate::import::NativeScanCache::shared();
1904 let Some(row) = visible_row_by_id(&connection, id, &ownership, &cache)? else {
1905 return Ok(None);
1906 };
1907 let session = sessionwiki::index::session_from_index(&connection, &row)
1908 .context("read an indexed session")?;
1909 let snapshot = snapshot_of(&session)?;
1910 Ok(Some(ArchivedSession {
1911 title: session.title.clone(),
1912 project_directory: project_directory_of(&session.project),
1913 snapshot,
1914 }))
1915}
1916
1917pub fn sessions_ready_to_archive(
1933 sessions: &mj_core::snapshot_map::SnapshotMap<String, SessionRecord>,
1934 subagents: &mj_core::snapshot_map::SnapshotMap<String, mj_core::subagent::SubagentRecord>,
1935 now: DateTime<Utc>,
1936 older_than_days: u32,
1937) -> Vec<String> {
1938 let state = mj_core::state::State {
1939 sessions: sessions.clone(),
1940 subagents: subagents.clone(),
1941 ..mj_core::state::State::default()
1942 };
1943 sessions_ready_to_archive_from_state(&state, now, older_than_days)
1944}
1945
1946pub(crate) fn sessions_ready_to_archive_from_state(
1947 state: &mj_core::state::State,
1948 now: DateTime<Utc>,
1949 older_than_days: u32,
1950) -> Vec<String> {
1951 let sessions = &state.sessions;
1952 let subagents = &state.subagents;
1953 let cutoff = now - chrono::Duration::days(i64::from(older_than_days));
1954 let aged = |session_id: &String| {
1955 sessions.get(session_id).is_some_and(|record| {
1956 let checkout = state.checkout(session_id).ok();
1957 let is_managed_worktree = checkout.as_ref().is_some_and(|checkout| {
1958 matches!(
1959 checkout.effective(),
1960 mj_core::state::Checkout::ManagedWorktree { worktree, .. }
1961 if worktree.kind == mj_core::state::ManagedCheckoutKind::Worktree
1962 )
1963 });
1964 let has_raw_path = checkout.as_ref().is_some_and(|checkout| {
1967 checkout.project_directory().is_some()
1968 && matches!(
1969 checkout,
1970 mj_core::state::Checkout::Attached { .. }
1971 | mj_core::state::Checkout::Borrowed { .. }
1972 )
1973 });
1974 record.state == mj_core::state::SessionState::Stopped
1975 && parse_time(&record.updated_at).is_some_and(|updated| updated <= cutoff)
1976 && (is_managed_worktree
1977 || has_raw_path
1978 || record
1979 .checkpoint
1980 .as_ref()
1981 .zip(record.publication.as_ref())
1982 .is_some_and(|(checkpoint, publication)| {
1983 publication.checkpoint_sha256 == checkpoint.sha256
1984 && publication.state == mj_core::state::PublicationState::Published
1985 && !publication.dirty
1986 && !publication.stashed
1987 }))
1988 })
1989 };
1990 let selected: BTreeSet<String> = sessions
1991 .keys()
1992 .filter(|session_id| aged(session_id))
1993 .filter(|session_id| {
1994 subagents
1995 .values()
1996 .filter(|child| &&child.parent_session_id == session_id)
1997 .filter(|child| sessions.contains_key(&child.child_session_id))
1999 .all(|child| aged(&child.child_session_id))
2000 })
2001 .cloned()
2002 .collect();
2003 let mut ordered: Vec<String> = selected.iter().cloned().collect();
2004 ordered.sort_by_key(|session_id| std::cmp::Reverse(ancestor_depth(session_id, subagents)));
2005 ordered
2006}
2007
2008fn ancestor_depth(
2012 session_id: &str,
2013 subagents: &mj_core::snapshot_map::SnapshotMap<String, mj_core::subagent::SubagentRecord>,
2014) -> usize {
2015 let mut depth = 0;
2016 let mut current = session_id;
2017 while let Some(parent) = subagents
2019 .get(current)
2020 .map(|child| child.parent_session_id.as_str())
2021 {
2022 depth += 1;
2023 if depth > subagents.len() {
2024 break;
2025 }
2026 current = parent;
2027 }
2028 depth
2029}
2030
2031pub use mj_core::state::ArchiveSpacePreview;
2037
2038pub fn archive_space_preview(older_than_days: Option<u32>) -> Result<ArchiveSpacePreview> {
2050 let controller =
2051 Controller::load().context("load the session records to size their storage")?;
2052 Ok(archive_space_over_state(
2053 &mj_core::config::sessions_dir(),
2054 &controller.state,
2055 Utc::now(),
2056 older_than_days,
2057 ))
2058}
2059
2060fn archive_space_over_state(
2062 sessions_root: &Path,
2063 state: &mj_core::state::State,
2064 now: DateTime<Utc>,
2065 older_than_days: Option<u32>,
2066) -> ArchiveSpacePreview {
2067 let sessions = &state.sessions;
2068 let mut preview = ArchiveSpacePreview {
2069 sessions: sessions.len(),
2070 bytes: sessions
2071 .iter()
2072 .map(|(session_id, record)| session_bytes(sessions_root, session_id, record))
2073 .sum(),
2074 reclaimable_sessions: 0,
2075 reclaimable_bytes: 0,
2076 };
2077 if let Some(days) = older_than_days {
2078 let aged = sessions_ready_to_archive_from_state(state, now, days);
2079 preview.reclaimable_sessions = aged.len();
2080 preview.reclaimable_bytes = aged
2081 .iter()
2082 .filter_map(|session_id| {
2083 sessions
2084 .get(session_id)
2085 .map(|record| session_bytes(sessions_root, session_id, record))
2086 })
2087 .sum();
2088 }
2089 preview
2090}
2091
2092#[cfg(test)]
2094fn archive_space_over(
2095 sessions_root: &Path,
2096 sessions: &mj_core::snapshot_map::SnapshotMap<String, SessionRecord>,
2097 subagents: &mj_core::snapshot_map::SnapshotMap<String, mj_core::subagent::SubagentRecord>,
2098 now: DateTime<Utc>,
2099 older_than_days: Option<u32>,
2100) -> ArchiveSpacePreview {
2101 let state = mj_core::state::State {
2102 sessions: sessions.clone(),
2103 subagents: subagents.clone(),
2104 ..mj_core::state::State::default()
2105 };
2106 archive_space_over_state(sessions_root, &state, now, older_than_days)
2107}
2108
2109fn session_bytes(sessions_root: &Path, session_id: &str, record: &SessionRecord) -> u64 {
2112 let checkpoint = record
2113 .checkpoint
2114 .as_ref()
2115 .and_then(|checkpoint| std::fs::metadata(&checkpoint.archive_path).ok())
2116 .filter(|metadata| metadata.is_file())
2117 .map(|metadata| metadata.len())
2118 .unwrap_or(0);
2119 let attachments = sessions_root
2120 .join(session_id)
2121 .join(mj_core::attachment::ATTACHMENT_DIR);
2122 let attachments = crate::import::claude::directory_size(&attachments).unwrap_or(0);
2123 checkpoint.saturating_add(attachments)
2124}
2125
2126pub fn indexed_with_messages(session_ids: &[String]) -> Result<BTreeSet<String>> {
2133 if !index_is_writable() {
2134 return Ok(BTreeSet::new());
2137 }
2138 let connection = open_readonly()?;
2139 let sessions_dir = mj_core::config::sessions_dir();
2140 let mut indexed = BTreeSet::new();
2141 for session_id in session_ids {
2142 let key = format!("{}/{session_id}", sessions_dir.display());
2143 let rows = sessionwiki::index::resolve(&connection, session_id)
2144 .context("look up a stopped session in the SessionWiki index")?;
2145 if rows
2146 .iter()
2147 .any(|row| row.tool == TOOL && row.path == key && row.msg_count > 0 && !row.archived)
2148 {
2149 indexed.insert(session_id.clone());
2150 }
2151 }
2152 Ok(indexed)
2153}
2154
2155fn open_readonly() -> Result<rusqlite::Connection> {
2156 sessionwiki::index::open_readonly().context("open the SessionWiki index")
2157}
2158
2159fn row_by_id(
2162 connection: &rusqlite::Connection,
2163 id: &str,
2164) -> Result<Option<sessionwiki::index::SessionRow>> {
2165 Ok(sessionwiki::index::resolve(connection, id)
2166 .context("look up an indexed session")?
2167 .into_iter()
2168 .find(|row| row.session_id == id))
2169}
2170
2171fn visible_row_by_id(
2174 connection: &rusqlite::Connection,
2175 id: &str,
2176 ownership: &top_level::Snapshot,
2177 cache: &crate::import::NativeScanCache,
2178) -> Result<Option<sessionwiki::index::SessionRow>> {
2179 let Some(row) = row_by_id(connection, id)? else {
2180 return Ok(None);
2181 };
2182 if top_level::is_child(&row, ownership, cache)? {
2183 Ok(None)
2184 } else {
2185 Ok(Some(row))
2186 }
2187}
2188
2189fn wiki_row(
2190 row: sessionwiki::index::SessionRow,
2191 snippet: Option<String>,
2192 live: &BTreeSet<String>,
2193) -> WikiRow {
2194 let hel_session_id = (row.tool == TOOL)
2197 .then(|| row.path.rsplit('/').next().unwrap_or_default().to_owned())
2198 .filter(|session_id| live.contains(session_id));
2199 let native_id = sessionwiki::index::native_id_of(&row.path);
2200 WikiRow {
2201 id: row.session_id,
2202 tool: row.tool,
2203 project: row.project,
2204 title: row.title,
2205 started: row.started,
2206 msgs: row.msg_count,
2207 preview: row.preview,
2208 archived: row.archived,
2209 native_id,
2210 snippet,
2211 hel_session_id,
2212 target: None,
2215 profile: None,
2216 harness: None,
2217 }
2218}
2219
2220fn project_directory_of(project: &str) -> Option<PathBuf> {
2228 if project.trim().is_empty() {
2229 return None;
2230 }
2231 let path = PathBuf::from(project);
2232 let repository = path
2233 .ancestors()
2234 .find(|ancestor| ancestor.file_name().is_some_and(|name| name == ".mj"))
2235 .and_then(std::path::Path::parent)
2236 .map(std::path::Path::to_path_buf)
2237 .unwrap_or(path);
2238 repository.is_dir().then_some(repository)
2239}
2240
2241fn snapshot_of(
2248 session: &sessionwiki::model::Session,
2249) -> Result<mj_core::archive::CanonicalSessionSnapshot> {
2250 use mj_core::archive::{
2251 CanonicalExecutionState, CanonicalSessionSnapshot, CanonicalSessionState,
2252 CanonicalTranscriptBody, CanonicalTranscriptItem,
2253 };
2254
2255 let started_ms = session
2256 .started
2257 .map(|time| time.timestamp_millis())
2258 .unwrap_or_default();
2259 let mut transcript: Vec<CanonicalTranscriptItem> = Vec::new();
2260 for message in &session.messages {
2261 let text = message.text.trim();
2262 if text.is_empty() {
2263 continue;
2264 }
2265 if transcript.is_empty() && message.role != Role::User {
2268 continue;
2269 }
2270 let position = transcript.len() as u64 + 1;
2271 let body = match message.role {
2272 Role::User => CanonicalTranscriptBody::User {
2273 content: vec![serde_json::json!({"type": "text", "text": text})],
2274 },
2275 Role::Assistant => CanonicalTranscriptBody::Agent {
2276 chunks: vec![serde_json::json!({
2277 "content": {"type": "text", "text": text}
2278 })],
2279 streaming: false,
2280 },
2281 Role::Tool => {
2283 let (call, terminal_outputs) = mj_transcript::summary::indexed_tool_call(
2284 text,
2285 &format!("wiki-tool-{position}"),
2286 );
2287 CanonicalTranscriptBody::Tool {
2288 call,
2289 terminal_outputs,
2290 terminal_refs: Vec::new(),
2291 presentation: None,
2292 }
2293 }
2294 };
2295 let created_at_ms = message
2296 .ts
2297 .map(|time| time.timestamp_millis())
2298 .unwrap_or(started_ms);
2299 transcript.push(CanonicalTranscriptItem {
2300 stable_id: format!("wiki-{position}"),
2301 position,
2302 latest_content_event_ordinal: matches!(body, CanonicalTranscriptBody::Agent { .. })
2305 .then_some(position),
2306 created_at_ms,
2307 last_changed_at_ms: created_at_ms,
2308 body,
2309 });
2310 }
2311 anyhow::ensure!(
2312 !transcript.is_empty(),
2313 "the archived session has no prompt to restore from"
2314 );
2315
2316 let event_frontier = transcript.len() as u64;
2317 let last_activity_at_ms = transcript.last().map(|item| item.last_changed_at_ms);
2318 Ok(CanonicalSessionSnapshot {
2319 command_ledger: None,
2320 assessment_state: None,
2321 event_frontier,
2322 event_frontier_digest: {
2326 use sha2::Digest;
2327 mj_core::hex::lower_hex(sha2::Sha256::digest(
2328 format!("sessionwiki:{}", session.id).as_bytes(),
2329 ))
2330 },
2331 session: CanonicalSessionState {
2332 execution: CanonicalExecutionState::Idle,
2333 last_activity_at_ms,
2334 session_title: Some(session.title.clone()).filter(|title| !title.trim().is_empty()),
2335 configuration: Default::default(),
2336 },
2337 transcript,
2338 queued_prompts: Vec::new(),
2339 })
2340}
2341
2342#[derive(Debug, Clone, PartialEq, Eq)]
2352pub enum WikiContinuation {
2353 Resume { session_id: String },
2355 Restore { wiki_id: String },
2358 Import {
2360 harness: HarnessKind,
2361 native_session_id: String,
2362 },
2363}
2364
2365pub fn wiki_continuation(
2371 wiki_id: &str,
2372 tool: &str,
2373 path: &Path,
2374 has_record: bool,
2375) -> Result<WikiContinuation> {
2376 if tool == TOOL {
2377 return Ok(match has_record {
2380 true => WikiContinuation::Resume {
2381 session_id: wiki_id.to_owned(),
2382 },
2383 false => WikiContinuation::Restore {
2384 wiki_id: wiki_id.to_owned(),
2385 },
2386 });
2387 }
2388 let harness = harness_adapters::harness_for_tool(tool)
2389 .with_context(|| format!("Mjolnir cannot continue a {tool} session"))?;
2390 let native_session_id = crate::import::native_session_id_from_path(harness, path)
2391 .with_context(|| {
2392 format!(
2393 "no {tool} session id in the indexed path {}",
2394 path.display()
2395 )
2396 })?;
2397 Ok(WikiContinuation::Import {
2398 harness,
2399 native_session_id,
2400 })
2401}
2402
2403pub fn wiki_session(
2408 wiki_id: &str,
2409 known_sessions: &BTreeSet<String>,
2410) -> Result<Option<WikiSessionInfo>> {
2411 if !index_is_writable() {
2412 return Ok(None);
2413 }
2414 let connection = open_readonly()?;
2415 let ownership = top_level::current_snapshot()?;
2416 let cache = crate::import::NativeScanCache::shared();
2417 wiki_session_from(&connection, wiki_id, known_sessions, &ownership, &cache)
2418}
2419
2420fn wiki_session_from(
2421 connection: &rusqlite::Connection,
2422 wiki_id: &str,
2423 known_sessions: &BTreeSet<String>,
2424 ownership: &top_level::Snapshot,
2425 cache: &crate::import::NativeScanCache,
2426) -> Result<Option<WikiSessionInfo>> {
2427 let Some(row) = visible_row_by_id(connection, wiki_id, ownership, cache)? else {
2428 return Ok(None);
2429 };
2430 let is_mjolnir = row.tool == TOOL;
2431 let mjolnir_session_id = is_mjolnir.then(|| row.session_id.clone());
2432 let has_record = mjolnir_session_id
2433 .as_deref()
2434 .is_some_and(|session_id| known_sessions.contains(session_id));
2435 let status = match (is_mjolnir, has_record) {
2436 (false, _) => WikiSessionStatus::Native,
2437 (true, true) => WikiSessionStatus::Mine,
2438 (true, false) => WikiSessionStatus::Archived,
2439 };
2440 let tags = match is_mjolnir {
2441 true => tags::read(connection, &[row.session_id.as_str()])
2442 .context("read the indexed session metadata")?
2443 .remove(&row.session_id)
2444 .unwrap_or_default(),
2445 false => tags::MjTags::default(),
2446 };
2447 let nothing_to_restore = status == WikiSessionStatus::Archived
2450 && !has_prompt(
2451 &sessionwiki::index::session_from_index(connection, &row)
2452 .context("read an indexed session")?,
2453 );
2454 let harness = tags
2455 .harness
2456 .as_deref()
2457 .and_then(|id| id.parse::<HarnessKind>().ok())
2458 .or_else(|| {
2459 (!is_mjolnir)
2460 .then(|| harness_adapters::harness_for_tool(&row.tool))
2461 .flatten()
2462 });
2463 Ok(Some(WikiSessionInfo {
2464 wiki_id: row.session_id,
2465 tool: row.tool,
2466 path: PathBuf::from(row.path),
2467 status,
2468 mjolnir_session_id,
2469 profile_id: tags.profile,
2470 target_template_id: tags.target,
2471 harness,
2472 title: row.title,
2473 project: row.project,
2474 nothing_to_restore,
2475 }))
2476}
2477
2478fn has_prompt(session: &sessionwiki::model::Session) -> bool {
2481 session
2482 .messages
2483 .iter()
2484 .any(|message| message.role == Role::User && !message.text.trim().is_empty())
2485}
2486
2487#[cfg(test)]
2488mod tests {
2489 use std::collections::BTreeMap;
2490 use std::path::Path;
2491
2492 mod continuation {
2495 use super::super::{WikiContinuation, wiki_continuation};
2496 use mj_core::config::HarnessKind;
2497 use std::path::Path;
2498
2499 #[test]
2500 fn a_claude_code_row_is_imported_with_the_uuid_from_its_path() {
2501 let path = Path::new(
2502 "/home/user/.claude/projects/-home-user-app/7f3a1c20-0b11-4a55-9e0d-2c8a5d6f1b44.jsonl",
2503 );
2504 assert_eq!(
2505 wiki_continuation("abc123", "claude-code", path, false).unwrap(),
2506 WikiContinuation::Import {
2507 harness: HarnessKind::Claude,
2508 native_session_id: "7f3a1c20-0b11-4a55-9e0d-2c8a5d6f1b44".to_owned(),
2509 }
2510 );
2511 }
2512
2513 #[test]
2516 fn a_codex_row_is_imported_with_the_uuid_from_its_rollout_name() {
2517 let path = Path::new(
2518 "/home/user/.codex/sessions/2026/09/18/rollout-2026-09-18T09-15-00-7f3a1c20-0b11-4a55-9e0d-2c8a5d6f1b44.jsonl",
2519 );
2520 assert_eq!(
2521 wiki_continuation("abc123", "codex", path, false).unwrap(),
2522 WikiContinuation::Import {
2523 harness: HarnessKind::Codex,
2524 native_session_id: "7f3a1c20-0b11-4a55-9e0d-2c8a5d6f1b44".to_owned(),
2525 }
2526 );
2527 }
2528 }
2529
2530 use mj_checkpoint::archive::{
2531 ArchiveInput, BundleManifest, CanonicalExecutionState, CanonicalSessionSnapshot,
2532 CanonicalSessionState, CanonicalTranscriptBody, CanonicalTranscriptItem, SessionManifest,
2533 TargetManifest, write_archive_atomic,
2534 };
2535
2536 use super::*;
2537
2538 fn wiki_session_for_test(
2539 wiki_id: &str,
2540 known_sessions: &BTreeSet<String>,
2541 ) -> Result<Option<WikiSessionInfo>> {
2542 let connection = open_readonly()?;
2543 wiki_session_from(
2544 &connection,
2545 wiki_id,
2546 known_sessions,
2547 &top_level::Snapshot::default(),
2548 &crate::import::NativeScanCache::new(),
2549 )
2550 }
2551
2552 fn query_rows_for_test(
2553 query: &str,
2554 limit: usize,
2555 include_tool_matches: bool,
2556 ) -> Result<Vec<WikiRow>> {
2557 query_rows_with_snapshot_for_test(
2558 query,
2559 limit,
2560 include_tool_matches,
2561 &top_level::Snapshot::default(),
2562 )
2563 }
2564
2565 fn query_rows_with_snapshot_for_test(
2566 query: &str,
2567 limit: usize,
2568 include_tool_matches: bool,
2569 ownership: &top_level::Snapshot,
2570 ) -> Result<Vec<WikiRow>> {
2571 let connection = open_readonly()?;
2572 query_rows_from(
2573 &connection,
2574 query,
2575 limit,
2576 &BTreeSet::new(),
2577 include_tool_matches,
2578 ownership,
2579 &crate::import::NativeScanCache::new(),
2580 )
2581 }
2582
2583 fn item(position: u64, body: CanonicalTranscriptBody) -> CanonicalTranscriptItem {
2584 let streamed = matches!(body, CanonicalTranscriptBody::Agent { .. });
2587 CanonicalTranscriptItem {
2588 stable_id: format!("item-{position}"),
2589 position,
2590 latest_content_event_ordinal: streamed.then_some(position),
2591 created_at_ms: 1_700_000_000_000 + i64::try_from(position).unwrap(),
2592 last_changed_at_ms: 1_700_000_000_000 + i64::try_from(position).unwrap(),
2593 body,
2594 }
2595 }
2596
2597 fn write_archive(directory: &Path, session_id: &str, frontier: u64) {
2600 let path = directory.join(format!(
2601 "{session_id}-{frontier}-archive-{}.hel.zip",
2602 "0".repeat(32)
2603 ));
2604 write_archive_atomic(
2605 &path,
2606 &ArchiveInput {
2607 session: SessionManifest {
2608 id: session_id.into(),
2609 title: "indexed session".into(),
2610 harness_kind: mj_core::config::HarnessKind::Codex,
2611 profile_id: "codex".into(),
2612 native_session_id: "native-session".into(),
2613 created_at: "2026-09-01T00:00:00Z".into(),
2614 checkpointed_at: "2026-09-01T01:00:00Z".into(),
2615 hel_version: "test".into(),
2616 relay_version: "test".into(),
2617 adapter_version: "test".into(),
2618 },
2619 target: TargetManifest {
2620 template_id: "local".into(),
2621 target_kind: "local-bare".into(),
2622 details: BTreeMap::new(),
2623 },
2624 bundle: BundleManifest {
2625 id: "project".into(),
2626 primary_repository: "project".into(),
2627 },
2628 canonical_session: CanonicalSessionSnapshot {
2629 command_ledger: None,
2630 assessment_state: None,
2631 event_frontier: 4,
2632 event_frontier_digest: "a".repeat(64),
2633 session: CanonicalSessionState {
2634 execution: CanonicalExecutionState::Idle,
2635 last_activity_at_ms: Some(1_700_000_000_004),
2636 session_title: Some("snapshot title".into()),
2637 configuration: Default::default(),
2638 },
2639 transcript: vec![
2640 item(
2641 1,
2642 CanonicalTranscriptBody::User {
2643 content: vec![serde_json::json!({
2644 "type": "text",
2645 "text": "index this session"
2646 })],
2647 },
2648 ),
2649 item(
2650 2,
2651 CanonicalTranscriptBody::Thought {
2652 chunks: vec![serde_json::json!({
2653 "content": {"type": "text", "text": "pondering"}
2654 })],
2655 streaming: false,
2656 },
2657 ),
2658 item(
2659 3,
2660 CanonicalTranscriptBody::Tool {
2661 call: serde_json::json!({
2662 "toolCallId": "call-1",
2663 "title": "Edit config.toml",
2664 "kind": "edit",
2665 "status": "completed",
2666 "locations": [{"path": "/old/container/config.toml"}]
2667 }),
2668 terminal_outputs: Vec::new(),
2669 terminal_refs: Vec::new(),
2670 presentation: None,
2671 },
2672 ),
2673 item(
2674 4,
2675 CanonicalTranscriptBody::Agent {
2676 chunks: vec![serde_json::json!({
2677 "content": {"type": "text", "text": "done"}
2678 })],
2679 streaming: false,
2680 },
2681 ),
2682 ],
2683 queued_prompts: Vec::new(),
2684 },
2685 native_artifacts: Vec::new(),
2686 repositories: Vec::new(),
2687 },
2688 )
2689 .unwrap();
2690 }
2691
2692 fn adapter(directory: &Path, session_id: &str) -> MjolnirAdapter {
2693 adapter_with_live(directory, session_id, BTreeMap::new())
2694 }
2695
2696 fn adapter_with_live(
2697 directory: &Path,
2698 session_id: &str,
2699 live: BTreeMap<String, i64>,
2700 ) -> MjolnirAdapter {
2701 let record = SessionRecord {
2702 project: None,
2703 id: session_id.into(),
2704 ..record_template()
2705 };
2706 MjolnirAdapter {
2707 sessions_dir: directory.to_path_buf(),
2708 sessions: std::sync::Mutex::new(Sessions {
2709 records: [(session_id.to_owned(), record)].into_iter().collect(),
2710 ownership: top_level::Snapshot::default(),
2711 project_directories: [(
2712 session_id.to_owned(),
2713 Ok(Some(PathBuf::from("/home/dev/project"))),
2714 )]
2715 .into_iter()
2716 .collect(),
2717 live,
2718 }),
2719 reload: false,
2720 }
2721 }
2722
2723 fn adapter_for_state(directory: &Path, state: &State, config: &Config) -> MjolnirAdapter {
2724 let sessions = Sessions {
2725 records: state.sessions.clone(),
2726 ownership: top_level::Snapshot::from_state(state),
2727 project_directories: project_directories_of(state, Some(config)),
2728 live: BTreeMap::new(),
2729 };
2730 MjolnirAdapter {
2731 sessions_dir: directory.to_path_buf(),
2732 sessions: std::sync::Mutex::new(sessions),
2733 reload: false,
2734 }
2735 }
2736
2737 fn sync_adapter(connection: &mut rusqlite::Connection, source: &Arc<MjolnirAdapter>) {
2738 let adapter: Box<dyn Adapter> = Box::new(SharedMjolnirAdapter(Arc::clone(source)));
2739 sessionwiki::index::sync_with(connection, &[adapter], None).unwrap();
2740 }
2741
2742 fn indexed_project(connection: &rusqlite::Connection, session_id: &str) -> String {
2743 connection
2744 .query_row(
2745 "SELECT project FROM files WHERE session_id = ?1",
2746 [session_id],
2747 |row| row.get(0),
2748 )
2749 .unwrap()
2750 }
2751
2752 fn indexed_mtime(connection: &rusqlite::Connection, session_id: &str) -> i64 {
2753 connection
2754 .query_row(
2755 "SELECT mtime FROM files WHERE session_id = ?1",
2756 [session_id],
2757 |row| row.get(0),
2758 )
2759 .unwrap()
2760 }
2761
2762 fn bundle_config() -> Config {
2763 let mut config = Config::default();
2764 config.bundles.insert(
2765 "project".into(),
2766 mj_core::config::ProjectBundle {
2767 primary_repo: "bifrost".into(),
2768 repositories: vec![mj_core::config::ProjectRepository {
2769 id: "bifrost".into(),
2770 github: Some("BrokkAi/bifrost".into()),
2771 local: None,
2772 destination: PathBuf::from("bifrost"),
2773 git_ref: None,
2774 }],
2775 },
2776 );
2777 config
2778 }
2779
2780 fn bundle_record(session_id: &str) -> SessionRecord {
2781 SessionRecord {
2782 id: session_id.into(),
2783 project_directory: None,
2784 container_workspace: Some(PathBuf::from(format!("/workspace/{session_id}"))),
2785 target: Some(mj_core::state::TargetLocator::LocalPodman {
2786 container_id: "test-container".into(),
2787 workspace_storage: Default::default(),
2788 borrowed_from: None,
2789 }),
2790 ..record_template()
2791 }
2792 }
2793
2794 fn state_with_record(session_id: &str, record: SessionRecord) -> State {
2795 let mut state = State::default();
2796 state.sessions.insert(session_id.to_owned(), record);
2797 state
2798 }
2799
2800 fn record_template() -> SessionRecord {
2801 SessionRecord {
2802 project: None,
2803 target_runtime: None,
2804 launch_base: None,
2805 launch_branch: None,
2806 checkout: None,
2807 publication: None,
2808 build_cache: None,
2809 container_workspace: None,
2810 subagents: None,
2811 create_managed_worktree: None,
2812 workspace_id: mj_core::workspace::DEFAULT_WORKSPACE_ID.to_owned(),
2813 archived: false,
2814 container_cpus: None,
2815 container_memory: None,
2816 id: "0123456789abcdef0123456789abcdef".into(),
2817 title: "indexed session".into(),
2818 harness_kind: mj_core::config::HarnessKind::Codex,
2819 last_profile: "codex".into(),
2820 bundle_id: "project".into(),
2821 project_directory: Some(PathBuf::from("/home/dev/project")),
2822 managed_worktree: None,
2823 review: None,
2824 target_template_id: "local-bare".into(),
2825 resource_allocation: None,
2826 additional_mounts: Vec::new(),
2827 state: mj_core::state::SessionState::Stopped,
2828 target: None,
2829 native_session_id: Some("native-session".into()),
2830 acp_session_title: Some("the harness title".into()),
2831 session_title_override: None,
2832 created_at: "2026-09-01T00:00:00Z".into(),
2833 updated_at: "2026-09-01T01:00:00Z".into(),
2834 viewed_through_event_ordinal: 0,
2835 draft_input: String::new(),
2836 last_error: None,
2837 last_checkpoint_error: None,
2838 checkpoint: None,
2839 }
2840 }
2841
2842 #[test]
2843 fn children_have_no_store_keys_metadata_or_pre_destroy_work() {
2844 let _held = tags::testing::lock();
2845 let (_index, _connection) = tags::testing::isolated_index();
2846 let directory = tempfile::tempdir().unwrap();
2847 let parent = "0123456789abcdef0123456789abcdef";
2848 let child = "fedcba9876543210fedcba9876543210";
2849 write_archive(directory.path(), parent, 1);
2850 write_archive(directory.path(), child, 1);
2851 let source = adapter_with_live(
2852 directory.path(),
2853 parent,
2854 BTreeMap::from([
2855 (parent.to_owned(), 1_900_000_000),
2856 (child.to_owned(), 1_900_000_001),
2857 ]),
2858 );
2859 {
2860 let mut sessions = source.sessions.lock().unwrap();
2861 sessions.ownership = top_level::Snapshot::for_test(BTreeSet::from([child.to_owned()]));
2862 sessions.records.insert(
2863 child.to_owned(),
2864 SessionRecord {
2865 id: child.into(),
2866 ..record_template()
2867 },
2868 );
2869 }
2870 let store = source.store().unwrap();
2871 assert_eq!(store.keys.len(), 1);
2872 assert_eq!(store.keys[0].0, source.key_for(parent));
2873 assert_eq!(store.files.len(), 1);
2874 assert_eq!(
2875 source.indexed_tags().keys().cloned().collect::<Vec<_>>(),
2876 [parent]
2877 );
2878 assert_eq!(
2879 unindexed(&source, &[parent.to_owned(), child.to_owned()]).unwrap(),
2880 [parent]
2881 );
2882 assert!(
2883 source
2884 .parse_key(&source.key_for(child))
2885 .unwrap_err()
2886 .to_string()
2887 .contains("sub-agent")
2888 );
2889 source.sessions.lock().unwrap().live.clear();
2891 assert_eq!(source.store().unwrap().keys.len(), 1);
2892 }
2893
2894 #[test]
2895 fn the_newest_checkpoint_of_each_session_is_one_indexed_key() {
2896 let directory = tempfile::tempdir().unwrap();
2897 let session_id = "0123456789abcdef0123456789abcdef";
2898 write_archive(directory.path(), session_id, 1);
2899 write_archive(directory.path(), session_id, 7);
2900 let adapter = adapter(directory.path(), session_id);
2901
2902 let store = adapter.store().expect("the adapter is a shared store");
2903 let key = format!("{}/{session_id}", directory.path().display());
2904 assert_eq!(
2905 store
2906 .keys
2907 .iter()
2908 .map(|(key, _)| key.as_str())
2909 .collect::<Vec<_>>(),
2910 vec![key.as_str()]
2911 );
2912 assert!(!store.had_error);
2913 assert_eq!(store.files.len(), 1);
2914 assert!(
2915 store.files[0]
2916 .file_name()
2917 .unwrap()
2918 .to_str()
2919 .unwrap()
2920 .contains("-7-archive-"),
2921 "the newest checkpoint is the one indexed: {:?}",
2922 store.files[0]
2923 );
2924 assert_eq!(
2925 adapter.reconcile_scope(),
2926 Some(format!("{}/", directory.path().display()))
2927 );
2928
2929 let session = adapter.parse_key(&key).unwrap();
2930 assert_eq!(session.id, session_id);
2931 assert_eq!(session.tool, "mjolnir");
2932 assert_eq!(session.path, PathBuf::from(&key));
2933 assert_eq!(session.project, "/home/dev/project");
2934 assert_eq!(session.title, "the harness title");
2935 assert!(!session.subagent);
2936 assert_eq!(
2937 session.messages.iter().map(|m| m.role).collect::<Vec<_>>(),
2938 vec![Role::User, Role::Tool, Role::Assistant]
2939 );
2940 assert_eq!(session.messages[0].text, "index this session");
2941 let tool: serde_json::Value = serde_json::from_str(&session.messages[1].text).unwrap();
2942 assert_eq!(tool["name"], "Edit");
2943 assert_eq!(tool["call"]["title"], "Edit config.toml");
2944 assert_eq!(session.messages[2].text, "done");
2945 assert_eq!(session.touched, vec!["/old/container/config.toml"]);
2946 }
2947
2948 #[test]
2949 fn a_bundle_session_is_indexed_and_searchable_by_its_primary_repository() {
2950 let _held = tags::testing::lock();
2951 let (_index_dir, mut connection) = tags::testing::isolated_index();
2952 let directory = tempfile::tempdir().unwrap();
2953 let session_id = "0123456789abcdef0123456789abcdef";
2954 write_archive(directory.path(), session_id, 1);
2955 let state = state_with_record(session_id, bundle_record(session_id));
2956 let source = Arc::new(adapter_for_state(
2957 directory.path(),
2958 &state,
2959 &bundle_config(),
2960 ));
2961
2962 sync_adapter(&mut connection, &source);
2963
2964 assert_eq!(
2965 indexed_project(&connection, session_id),
2966 format!("/workspace/{session_id}/bifrost")
2967 );
2968 assert!(
2969 query_rows_for_test("bifrost", 10, false)
2970 .unwrap()
2971 .iter()
2972 .any(|row| row.id == session_id),
2973 "the indexed bundle session is found by its repository name"
2974 );
2975 }
2976
2977 #[test]
2978 fn a_raw_local_session_keeps_its_directory_when_indexed() {
2979 let _held = tags::testing::lock();
2980 let (_index_dir, mut connection) = tags::testing::isolated_index();
2981 let directory = tempfile::tempdir().unwrap();
2982 let session_id = "fedcba9876543210fedcba9876543210";
2983 let project_directory = PathBuf::from("/home/jonathan/Projects/bifrost");
2984 write_archive(directory.path(), session_id, 1);
2985 let record = SessionRecord {
2986 id: session_id.into(),
2987 project_directory: Some(project_directory.clone()),
2988 ..record_template()
2989 };
2990 let state = state_with_record(session_id, record);
2991 let source = Arc::new(adapter_for_state(
2992 directory.path(),
2993 &state,
2994 &Config::default(),
2995 ));
2996
2997 sync_adapter(&mut connection, &source);
2998
2999 assert_eq!(
3000 indexed_project(&connection, session_id),
3001 project_directory.display().to_string()
3002 );
3003 }
3004
3005 #[test]
3006 fn a_row_from_the_previous_parse_format_is_reparsed_once() {
3007 let _held = tags::testing::lock();
3008 let (_index_dir, mut connection) = tags::testing::isolated_index();
3009 let directory = tempfile::tempdir().unwrap();
3010 let session_id = "0123456789abcdef0123456789abcdef";
3011 write_archive(directory.path(), session_id, 1);
3012 let mut record = bundle_record(session_id);
3013 record.updated_at = "2099-01-01T00:00:00Z".into();
3014 let previous_token = parse_time(&record.updated_at)
3015 .unwrap()
3016 .timestamp()
3017 .saturating_mul(1024)
3018 .saturating_add(i64::from(mj_transcript::summary::SUMMARY_VERSION));
3019 let state = state_with_record(session_id, record);
3020 let source = Arc::new(adapter_for_state(
3021 directory.path(),
3022 &state,
3023 &bundle_config(),
3024 ));
3025 let key = source.key_for(session_id);
3026 tags::testing::index_row(&connection, session_id, TOOL);
3027 connection
3028 .execute(
3029 "UPDATE files SET path = ?1, project = '', mtime = ?2 WHERE session_id = ?3",
3030 rusqlite::params![key, previous_token, session_id],
3031 )
3032 .unwrap();
3033
3034 sync_adapter(&mut connection, &source);
3035
3036 let parsed_token = indexed_mtime(&connection, session_id);
3037 assert_eq!(
3038 indexed_project(&connection, session_id),
3039 format!("/workspace/{session_id}/bifrost")
3040 );
3041 assert_eq!(
3042 parsed_token,
3043 source
3044 .store()
3045 .unwrap()
3046 .keys
3047 .into_iter()
3048 .find(|(path, _)| path == &key)
3049 .unwrap()
3050 .1
3051 );
3052
3053 source
3055 .sessions
3056 .lock()
3057 .unwrap()
3058 .project_directories
3059 .insert(session_id.into(), Err("unexpected second parse".into()));
3060 sync_adapter(&mut connection, &source);
3061
3062 assert_eq!(indexed_mtime(&connection, session_id), parsed_token);
3063 assert_eq!(
3064 indexed_project(&connection, session_id),
3065 format!("/workspace/{session_id}/bifrost")
3066 );
3067 }
3068
3069 #[test]
3070 fn provenance_backfill_repairs_an_unchanged_checkpoint_without_rebuilding_the_index() {
3071 let _held = tags::testing::lock();
3072 let (_index_dir, mut connection) = tags::testing::isolated_index();
3073 let directory = tempfile::tempdir().unwrap();
3074 write_archive(directory.path(), "old-session", 4);
3075 let source = adapter(directory.path(), "old-session");
3076 let key = source.key_for("old-session");
3077 tags::testing::index_row(&connection, "old-session", "mjolnir");
3078 connection
3079 .execute(
3080 "UPDATE files SET path = ?1 WHERE session_id = 'old-session'",
3081 [&key],
3082 )
3083 .unwrap();
3084 provenance::backfill(&mut connection, &source).unwrap();
3085 assert_eq!(
3086 sessionwiki::index::files_for(&connection, "old-session").unwrap(),
3087 vec!["/old/container/config.toml"]
3088 );
3089 provenance::backfill(&mut connection, &source).unwrap();
3090 assert_eq!(
3091 sessionwiki::index::sessions_for_file(&connection, "config.toml", 20)
3092 .unwrap()
3093 .len(),
3094 1
3095 );
3096 }
3097
3098 #[test]
3103 fn a_running_session_is_listed_with_its_own_change_token() {
3104 let directory = tempfile::tempdir().unwrap();
3105 let running = "0123456789abcdef0123456789abcdef";
3106 let never_checkpointed = "fedcba9876543210fedcba9876543210";
3107 write_archive(directory.path(), running, 3);
3108 let live = adapter_with_live(
3109 directory.path(),
3110 running,
3111 BTreeMap::from([
3112 (running.to_owned(), 1_900_000_000),
3113 (never_checkpointed.to_owned(), 1_900_000_001),
3114 ]),
3115 );
3116
3117 let store = live.store().expect("the adapter is a shared store");
3118 let key_of = |session_id: &str| format!("{}/{session_id}", directory.path().display());
3119 assert_eq!(
3120 store.keys,
3121 vec![
3122 (key_of(running), session_change_token(1_900_000_000)),
3123 (
3124 key_of(never_checkpointed),
3125 session_change_token(1_900_000_001)
3126 ),
3127 ],
3128 "a live session's own token replaces the checkpoint's"
3129 );
3130
3131 let stopped = adapter(directory.path(), running);
3134 let keys = stopped.store().expect("a shared store").keys;
3135 assert_eq!(keys.len(), 1);
3136 assert_eq!(keys[0].0, key_of(running));
3137 assert_ne!(keys[0].1, 1_900_000_000);
3138 assert_eq!(
3139 stopped.parse_key(&key_of(running)).unwrap().title,
3140 "the harness title",
3141 "a stopped session is parsed from its checkpoint"
3142 );
3143 }
3144
3145 #[test]
3148 fn a_rename_moves_a_session_change_token() {
3149 let directory = tempfile::tempdir().unwrap();
3150 let session_id = "0123456789abcdef0123456789abcdef";
3151 write_archive(directory.path(), session_id, 1);
3152 let adapter = adapter(directory.path(), session_id);
3153 let before = adapter.store().expect("a shared store").keys[0].1;
3154
3155 {
3156 let mut sessions = adapter.sessions.lock().unwrap();
3157 let record = sessions.records.get_mut(session_id).unwrap();
3158 record.session_title_override = Some("the new name".into());
3159 record.updated_at = "2099-01-01T00:00:00Z".into();
3160 }
3161 let after = adapter.store().expect("a shared store").keys[0].1;
3162 assert!(
3163 after > before,
3164 "a renamed session is re-indexed: {before} then {after}"
3165 );
3166 assert_eq!(
3167 adapter
3168 .parse_key(&format!("{}/{session_id}", directory.path().display()))
3169 .unwrap()
3170 .title,
3171 "the new name"
3172 );
3173 }
3174
3175 fn indexed(messages: Vec<(Role, &str)>) -> sessionwiki::model::Session {
3176 Session {
3177 id: "0123456789abcdef0123456789abcdef".into(),
3178 tool: "mjolnir",
3179 path: PathBuf::from("/sessions/0123456789abcdef0123456789abcdef"),
3180 project: "/home/dev/project".into(),
3181 started: DateTime::from_timestamp_millis(1_700_000_000_000),
3182 ended: None,
3183 title: "the archived session".into(),
3184 subagent: false,
3185 messages: messages
3186 .into_iter()
3187 .map(|(role, text)| Message {
3188 role,
3189 text: text.to_owned(),
3190 ts: None,
3191 })
3192 .collect(),
3193 touched: Vec::new(),
3194 edits: Vec::new(),
3195 }
3196 }
3197
3198 #[test]
3199 fn golden_wiki_resume_preview_search() {
3200 use std::fmt::Write as _;
3201
3202 fn response(out: &mut String, label: &str, transcript: &WikiHitTranscript) {
3203 writeln!(
3204 out,
3205 "=== {label} ({} preview blocks) ===",
3206 transcript.blocks.len()
3207 )
3208 .unwrap();
3209 writeln!(out, "{}", serde_json::to_string_pretty(transcript).unwrap()).unwrap();
3210 }
3211
3212 let mut out = String::new();
3213 let case_insensitive = indexed(vec![
3214 (Role::User, "Make the Tests green"),
3215 (Role::Assistant, "the tests are green now"),
3216 ]);
3217 response(
3218 &mut out,
3219 "case-insensitive query with user and assistant hits",
3220 &hit_transcript(&case_insensitive, "TESTS", 0, 4_000),
3221 );
3222
3223 let tool_only_and_assistant = indexed(vec![
3224 (Role::User, "make it build"),
3225 (Role::Tool, "cargo build --needle"),
3226 (Role::Assistant, "it builds"),
3227 ]);
3228 response(
3229 &mut out,
3230 "query found only in tool output",
3231 &hit_transcript(&tool_only_and_assistant, "needle", 1, 4_000),
3232 );
3233 response(
3234 &mut out,
3235 "tool context beside an assistant hit",
3236 &hit_transcript(&tool_only_and_assistant, "builds", 1, 4_000),
3237 );
3238
3239 mj_core::golden::assert_golden(
3240 env!("CARGO_MANIFEST_DIR"),
3241 "wiki-resume-preview-search",
3242 &out,
3243 );
3244 }
3245
3246 #[test]
3249 fn transcript_hits_keeps_context_and_marks_omissions() {
3250 let session = indexed(vec![
3251 (Role::User, "zero"),
3252 (Role::Assistant, "one needle one"),
3253 (Role::Tool, "two"),
3254 (Role::User, "three"),
3255 (Role::Assistant, "four"),
3256 (Role::Tool, "five"),
3257 (Role::User, "six needle six"),
3258 (Role::Assistant, "seven"),
3259 (Role::User, "eight"),
3260 ]);
3261
3262 let found = hit_transcript(&session, "needle", 1, 4_000);
3263
3264 let shown: Vec<(&str, &str, usize)> = found
3265 .blocks
3266 .iter()
3267 .map(|block| {
3268 (
3269 block.role.as_str(),
3270 block.text.as_str(),
3271 block.omitted_before,
3272 )
3273 })
3274 .collect();
3275 assert_eq!(
3276 shown,
3277 vec![
3278 ("user", "zero", 0),
3279 ("assistant", "one needle one", 0),
3280 ("tool", "two", 0),
3281 ("tool", "five", 2),
3282 ("user", "six needle six", 0),
3283 ("assistant", "seven", 0),
3284 ]
3285 );
3286 assert_eq!(found.omitted_after, 1, "the last message is not shown");
3287 assert!(found.blocks[0].hits.is_empty(), "context has no hits");
3288 }
3289
3290 #[test]
3293 fn transcript_hits_window_keeps_the_first_hit() {
3294 let filler = "x".repeat(4_000);
3295 let session = indexed(vec![(Role::User, &format!("{filler} needle {filler}"))]);
3296
3297 let found = hit_transcript(&session, "needle", 0, 100);
3298
3299 let block = &found.blocks[0];
3300 assert!(block.truncated);
3301 assert_eq!(block.text.chars().count(), 100);
3302 assert_eq!(block.hits.len(), 1, "the windowed text keeps its hit");
3303 let (start, end) = block.hits[0];
3304 assert_eq!(&block.text[start..end], "needle");
3305 assert!(
3306 start >= 20,
3307 "the window keeps lead-in before the hit, got {start}"
3308 );
3309 }
3310
3311 #[test]
3315 fn messages_before_the_first_prompt_are_dropped() {
3316 let snapshot = snapshot_of(&indexed(vec![
3317 (Role::Assistant, "still working"),
3318 (Role::User, "carry on"),
3319 ]))
3320 .unwrap();
3321 assert_eq!(snapshot.transcript.len(), 1);
3322 assert_eq!(snapshot.transcript[0].position, 1);
3323 snapshot.validate().unwrap();
3324
3325 let error = snapshot_of(&indexed(vec![(Role::Assistant, "nobody asked")])).unwrap_err();
3326 assert!(
3327 error.to_string().contains("no prompt"),
3328 "a session with no prompt cannot be restored: {error}"
3329 );
3330 assert!(!has_prompt(&indexed(vec![(
3332 Role::Assistant,
3333 "nobody asked"
3334 )])));
3335 assert!(has_prompt(&indexed(vec![(Role::User, "carry on")])));
3336 }
3337
3338 fn record(
3339 session_id: &str,
3340 state: mj_core::state::SessionState,
3341 updated_at: &str,
3342 ) -> SessionRecord {
3343 SessionRecord {
3344 project: None,
3345 id: session_id.into(),
3346 state,
3347 updated_at: updated_at.into(),
3348 ..record_template()
3349 }
3350 }
3351
3352 fn child(child_session_id: &str, parent_session_id: &str) -> mj_core::subagent::SubagentRecord {
3353 mj_core::subagent::SubagentRecord {
3354 child_session_id: child_session_id.into(),
3355 parent_session_id: parent_session_id.into(),
3356 task_name: "task".into(),
3357 profile_id: "codex".into(),
3358 model: None,
3359 effort: None,
3360 working_directory: PathBuf::new(),
3361 initial_prompt: "do the thing".into(),
3362 request_key: "key".into(),
3363 created_at: "2026-09-01T00:00:00Z".into(),
3364 noticed_turn: None,
3365 reported_finish: None,
3366 handback_tool: false,
3367 }
3368 }
3369
3370 fn ready(
3371 sessions: Vec<SessionRecord>,
3372 children: Vec<mj_core::subagent::SubagentRecord>,
3373 ) -> Vec<String> {
3374 let now = parse_time("2026-09-10T00:00:00Z").unwrap();
3375 sessions_ready_to_archive(
3376 &sessions
3377 .into_iter()
3378 .map(|record| (record.id.clone(), record))
3379 .collect(),
3380 &children
3381 .into_iter()
3382 .map(|child| (child.child_session_id.clone(), child))
3383 .collect(),
3384 now,
3385 3,
3386 )
3387 }
3388
3389 #[test]
3390 fn aged_clone_requires_clean_published_evidence_for_its_current_checkpoint() {
3391 let id = "0123456789abcdef0123456789abcdef";
3392 let root = PathBuf::from(format!("/srv/project/.mj/clones/{id}"));
3393 let mut session = record(
3394 id,
3395 mj_core::state::SessionState::Stopped,
3396 "2026-09-01T00:00:00Z",
3397 );
3398 session.project_directory = Some(root.clone());
3399 session.managed_worktree = Some(mj_core::state::ManagedWorktree {
3400 kind: mj_core::state::ManagedCheckoutKind::Clone,
3401 source_project_directory: "/srv/project".into(),
3402 source_repository: "/srv/project".into(),
3403 worktree_root: root,
3404 branch: "feature".into(),
3405 target: mj_core::state::ManagedWorktreeTarget::Local,
3406 base_commit: Some("1".repeat(40)),
3407 });
3408 session.checkpoint = Some(mj_core::state::CheckpointMetadata {
3409 archive_path: "sessions/checkpoint.hel.zip".into(),
3410 sha256: "a".repeat(64),
3411 created_at: "2026-09-01T00:00:00Z".into(),
3412 event_frontier: 0,
3413 });
3414 assert!(ready(vec![session.clone()], vec![]).is_empty());
3415 session.publication = Some(mj_core::state::PublicationAssessment {
3416 checkpoint_sha256: "a".repeat(64),
3417 state: mj_core::state::PublicationState::Published,
3418 dirty: false,
3419 stashed: false,
3420 saved_commits: vec!["2".repeat(40)],
3421 destinations: vec!["https://example.test/repository.git".into()],
3422 checked_at: "2026-09-01T01:00:00Z".into(),
3423 reason: Some("feature branch was pushed but not merged".into()),
3424 });
3425 assert_eq!(ready(vec![session.clone()], vec![]), vec![id]);
3426 session.publication.as_mut().unwrap().stashed = true;
3427 assert!(ready(vec![session.clone()], vec![]).is_empty());
3428 session.publication.as_mut().unwrap().stashed = false;
3429 session.publication.as_mut().unwrap().checkpoint_sha256 = "b".repeat(64);
3430 assert!(ready(vec![session], vec![]).is_empty());
3431 }
3432
3433 fn sized_session(
3435 root: &Path,
3436 session_id: &str,
3437 updated_at: &str,
3438 checkpoint_bytes: usize,
3439 attachment_bytes: &[usize],
3440 ) -> SessionRecord {
3441 let archive_path = root.join(format!("{session_id}.hel.zip"));
3442 std::fs::write(&archive_path, vec![b'c'; checkpoint_bytes]).unwrap();
3443 if !attachment_bytes.is_empty() {
3444 let attachments = root
3445 .join(session_id)
3446 .join(mj_core::attachment::ATTACHMENT_DIR);
3447 std::fs::create_dir_all(&attachments).unwrap();
3448 for (index, size) in attachment_bytes.iter().enumerate() {
3449 std::fs::write(attachments.join(format!("{index}.png")), vec![b'a'; *size])
3450 .unwrap();
3451 }
3452 }
3453 SessionRecord {
3454 project: None,
3455 checkpoint: Some(mj_core::state::CheckpointMetadata {
3456 archive_path,
3457 sha256: "0".repeat(64),
3458 created_at: updated_at.into(),
3459 event_frontier: 1,
3460 }),
3461 ..record(
3462 session_id,
3463 mj_core::state::SessionState::Stopped,
3464 updated_at,
3465 )
3466 }
3467 }
3468
3469 #[test]
3470 fn the_space_preview_sizes_every_session_and_only_the_aged_ones_as_reclaimable() {
3471 let directory = tempfile::tempdir().unwrap();
3472 let root = directory.path();
3473 let sessions: mj_core::snapshot_map::SnapshotMap<String, SessionRecord> = [
3474 sized_session(root, "old-stopped", "2026-09-01T00:00:00Z", 1000, &[10, 20]),
3475 sized_session(root, "just-stopped", "2026-09-09T00:00:00Z", 500, &[]),
3476 SessionRecord {
3479 project: None,
3480 checkpoint: Some(mj_core::state::CheckpointMetadata {
3481 archive_path: root.join("missing.hel.zip"),
3482 sha256: "0".repeat(64),
3483 created_at: "2026-09-01T00:00:00Z".into(),
3484 event_frontier: 1,
3485 }),
3486 ..record(
3487 "lost-checkpoint",
3488 mj_core::state::SessionState::Stopped,
3489 "2026-09-01T00:00:00Z",
3490 )
3491 },
3492 ]
3493 .into_iter()
3494 .map(|record| (record.id.clone(), record))
3495 .collect();
3496 let now = parse_time("2026-09-10T00:00:00Z").unwrap();
3497
3498 let all = archive_space_over(root, &sessions, &Default::default(), now, None);
3499 assert_eq!(all.sessions, 3);
3500 assert_eq!(all.bytes, 1530);
3501 assert_eq!(all.reclaimable_sessions, 0);
3502 assert_eq!(all.reclaimable_bytes, 0);
3503
3504 let aged = archive_space_over(root, &sessions, &Default::default(), now, Some(3));
3505 assert_eq!(aged.bytes, 1530);
3506 assert_eq!(
3507 (aged.reclaimable_sessions, aged.reclaimable_bytes),
3508 (2, 1030),
3509 "only the sessions the job would archive count, attachments included"
3510 );
3511 }
3512
3513 #[test]
3514 fn only_stopped_sessions_past_the_cut_off_are_archived() {
3515 use mj_core::state::SessionState;
3516 let selected = ready(
3517 vec![
3518 record("old-stopped", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3519 record(
3520 "just-stopped",
3521 SessionState::Stopped,
3522 "2026-09-09T00:00:00Z",
3523 ),
3524 record("old-running", SessionState::Running, "2026-09-01T00:00:00Z"),
3525 record("old-error", SessionState::Error, "2026-09-01T00:00:00Z"),
3526 record("unparsable", SessionState::Stopped, "not a time"),
3527 record("at-the-edge", SessionState::Stopped, "2026-09-07T00:00:00Z"),
3529 ],
3530 Vec::new(),
3531 );
3532 assert_eq!(selected, vec!["at-the-edge", "old-stopped"]);
3533 }
3534
3535 #[test]
3536 fn a_child_the_pass_is_not_archiving_holds_its_parent_back() {
3537 use mj_core::state::SessionState;
3538 let selected = ready(
3539 vec![
3540 record("parent", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3541 record(
3542 "running-child",
3543 SessionState::Running,
3544 "2026-09-01T00:00:00Z",
3545 ),
3546 ],
3547 vec![child("running-child", "parent")],
3548 );
3549 assert!(selected.is_empty(), "the parent must wait: {selected:?}");
3550
3551 let selected = ready(
3552 vec![
3553 record("parent", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3554 record("young-child", SessionState::Stopped, "2026-09-09T00:00:00Z"),
3555 ],
3556 vec![child("young-child", "parent")],
3557 );
3558 assert!(selected.is_empty(), "the parent must wait: {selected:?}");
3559
3560 let selected = ready(
3562 vec![record(
3563 "parent",
3564 SessionState::Stopped,
3565 "2026-09-01T00:00:00Z",
3566 )],
3567 vec![child("departed-child", "parent")],
3568 );
3569 assert_eq!(selected, vec!["parent"]);
3570 }
3571
3572 #[test]
3573 fn children_are_archived_before_their_parents() {
3574 use mj_core::state::SessionState;
3575 let selected = ready(
3576 vec![
3577 record("parent", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3578 record("child", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3579 record("grandchild", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3580 ],
3581 vec![child("child", "parent"), child("grandchild", "child")],
3582 );
3583 assert_eq!(selected, vec!["grandchild", "child", "parent"]);
3584 }
3585
3586 #[test]
3590 fn query_rows_returns_the_indexed_target_profile_and_harness() {
3591 let _held = tags::testing::lock();
3592 let (_directory, connection) = tags::testing::isolated_index();
3593 tags::testing::index_row(&connection, "mj-session", TOOL);
3594 tags::testing::index_row(&connection, "codex-session", "codex");
3595 tags::write(
3596 &connection,
3597 "mj-session",
3598 &tags::MjTags {
3599 target: Some("Prod-Box".into()),
3600 profile: Some("codex-Main".into()),
3601 harness: Some("codex".into()),
3602 },
3603 )
3604 .expect("write the session metadata");
3605
3606 let rows = query_rows_for_test("", 10, false).expect("query the index");
3607 let mjolnir = rows
3608 .iter()
3609 .find(|row| row.id == "mj-session")
3610 .expect("the Mjolnir row is returned");
3611 assert_eq!(mjolnir.target.as_deref(), Some("Prod-Box"));
3612 assert_eq!(mjolnir.profile.as_deref(), Some("codex-Main"));
3613 assert_eq!(mjolnir.harness.as_deref(), Some("codex"));
3614
3615 let codex = rows
3616 .iter()
3617 .find(|row| row.id == "codex-session")
3618 .expect("the Codex row is returned");
3619 assert_eq!(codex.target, None);
3620 assert_eq!(codex.profile, None);
3621 assert_eq!(codex.harness, None);
3622 }
3623
3624 #[test]
3625 fn standalone_indexed_children_are_hidden_from_every_history_listing() {
3626 let _held = tags::testing::lock();
3627 let (_directory, mut connection) = tags::testing::isolated_index();
3628 let home = tempfile::tempdir().unwrap();
3629 let project = home.path().join("projects/project");
3630 std::fs::create_dir_all(&project).unwrap();
3631 let paths = [
3632 project.join("00000000-0000-4000-8000-000000000001.jsonl"),
3633 project.join("00000000-0000-4000-8000-000000000002.jsonl"),
3634 project.join("00000000-0000-4000-8000-000000000003.jsonl"),
3635 project.join("00000000-0000-4000-8000-000000000004.jsonl"),
3636 ];
3637 let transcript = |timestamp: &str, content: &str, sidechain: bool| {
3638 format!(
3639 "{{\"type\":\"user\",\"cwd\":\"/src/project\",\"entrypoint\":\"cli\",\"timestamp\":\"{timestamp}\",\"isSidechain\":{sidechain},\"message\":{{\"role\":\"user\",\"content\":\"{content}\"}}}}\n"
3640 )
3641 };
3642 for (index, path) in paths.iter().enumerate() {
3643 let is_child = index >= 2;
3644 let content = if is_child {
3645 "quokka quokka quokka quokka quokka quokka"
3646 } else {
3647 "quokka parent conversation"
3648 };
3649 let timestamp = match index {
3650 0 => "2026-10-04T00:00:00Z",
3651 1 => "2026-10-03T00:00:00Z",
3652 _ => "2026-10-05T00:00:00Z",
3653 };
3654 std::fs::write(path, transcript(timestamp, content, index == 3)).unwrap();
3655 }
3656
3657 let adapter = Box::new(sessionwiki::adapters::ClaudeCode::in_home(
3660 home.path().to_owned(),
3661 ));
3662 sessionwiki::index::sync_with(&mut connection, &[adapter], None).unwrap();
3663 for path in &paths[2..] {
3664 let kind: String = connection
3665 .query_row(
3666 "SELECT kind FROM files WHERE path = ?1",
3667 [path.to_string_lossy()],
3668 |row| row.get(0),
3669 )
3670 .unwrap();
3671 assert_eq!(kind, "main", "the standalone sync misclassified {path:?}");
3672 }
3673
3674 let mj_child_id = "mj-child-session";
3675 let native_mj_child_id = "00000000-0000-4000-8000-000000000003";
3676 let mut state = mj_core::state::State::default();
3677 let mut child_record = crate::database::test_session(mj_child_id, "test-project");
3678 child_record.native_session_id = Some(native_mj_child_id.to_owned());
3679 state.sessions.insert(mj_child_id.to_owned(), child_record);
3680 state.subagents.insert(
3681 mj_child_id.to_owned(),
3682 mj_core::subagent::SubagentRecord {
3683 child_session_id: mj_child_id.to_owned(),
3684 parent_session_id: "mj-parent-session".to_owned(),
3685 task_name: "child".to_owned(),
3686 profile_id: "codex".to_owned(),
3687 model: None,
3688 effort: None,
3689 working_directory: PathBuf::from("/src/project"),
3690 initial_prompt: "child task".to_owned(),
3691 request_key: "test-child".to_owned(),
3692 created_at: "2026-10-05T00:00:00Z".to_owned(),
3693 noticed_turn: None,
3694 reported_finish: None,
3695 handback_tool: false,
3696 },
3697 );
3698 let ownership = top_level::Snapshot::from_state(&state);
3699 let ids_for_paths: Vec<String> = paths
3700 .iter()
3701 .map(|path| {
3702 connection
3703 .query_row(
3704 "SELECT session_id FROM files WHERE path = ?1",
3705 [path.to_string_lossy()],
3706 |row| row.get(0),
3707 )
3708 .unwrap()
3709 })
3710 .collect();
3711 for id in &ids_for_paths {
3712 connection
3713 .execute(
3714 "INSERT INTO touched(session_id, path) VALUES (?1, '/src/project/src/a.rs')",
3715 [id],
3716 )
3717 .unwrap();
3718 }
3719 connection
3720 .execute(
3721 "UPDATE files SET title = 'standalone-only-name' WHERE session_id = ?1",
3722 [&ids_for_paths[2]],
3723 )
3724 .unwrap();
3725
3726 let child_native_ids = BTreeSet::from([
3727 "00000000-0000-4000-8000-000000000003".to_owned(),
3728 "00000000-0000-4000-8000-000000000004".to_owned(),
3729 ]);
3730 let raw_hits = sessionwiki::index::search(&connection, "quokka", 10, None, None).unwrap();
3731 let first_search_ids: BTreeSet<_> = raw_hits
3732 .iter()
3733 .take(2)
3734 .filter_map(|hit| sessionwiki::index::native_id_of(&hit.row.path))
3735 .collect();
3736 assert_eq!(first_search_ids, child_native_ids);
3737 let raw_recent =
3738 sessionwiki::index::recent(&connection, 2, None, None, None, true).unwrap();
3739 let first_recent_ids: BTreeSet<_> = raw_recent
3740 .iter()
3741 .filter_map(|row| sessionwiki::index::native_id_of(&row.path))
3742 .collect();
3743 assert_eq!(first_recent_ids, child_native_ids);
3744 let resume = query_rows_with_snapshot_for_test("quokka", 2, false, &ownership).unwrap();
3745 let mut resume_ids: Vec<_> = resume.into_iter().map(|row| row.id).collect();
3746 resume_ids.sort();
3747 let mut parent_ids = ids_for_paths[..2].to_vec();
3748 parent_ids.sort();
3749 assert_eq!(resume_ids, parent_ids);
3750
3751 let recent = query_rows_with_snapshot_for_test("", 2, false, &ownership).unwrap();
3752 let mut recent_ids: Vec<_> = recent.into_iter().map(|row| row.id).collect();
3753 recent_ids.sort();
3754 assert_eq!(recent_ids, parent_ids);
3755 assert!(
3756 query_rows_with_snapshot_for_test("standalone-only-name", 10, false, &ownership,)
3757 .unwrap()
3758 .is_empty()
3759 );
3760
3761 let search_sessions = mj_core::history::HistoryRequest {
3762 request_id: "test-search".into(),
3763 query: mj_core::history::HistoryQuery::SearchSessions {
3764 query: "quokka".into(),
3765 limit: 2,
3766 },
3767 blame: None,
3768 };
3769 let history = history::query_in_with(
3770 &connection,
3771 &search_sessions,
3772 &ownership,
3773 &crate::import::NativeScanCache::new(),
3774 )
3775 .unwrap();
3776 let mut history_ids: Vec<String> = history["sessions"]
3777 .as_array()
3778 .unwrap()
3779 .iter()
3780 .map(|row| row["id"].as_str().unwrap().to_owned())
3781 .collect();
3782 history_ids.sort();
3783 assert_eq!(history_ids, parent_ids);
3784
3785 let trace = history::query_in_with(
3786 &connection,
3787 &mj_core::history::HistoryRequest {
3788 request_id: "test-trace".into(),
3789 query: mj_core::history::HistoryQuery::TraceFile {
3790 path: "src/a.rs".into(),
3791 limit: 2,
3792 },
3793 blame: None,
3794 },
3795 &ownership,
3796 &crate::import::NativeScanCache::new(),
3797 )
3798 .unwrap();
3799 let mut trace_ids: Vec<String> = trace["sessions"]
3800 .as_array()
3801 .unwrap()
3802 .iter()
3803 .map(|item| item["session"]["id"].as_str().unwrap().to_owned())
3804 .collect();
3805 trace_ids.sort();
3806 assert_eq!(trace_ids, parent_ids);
3807
3808 let blame_time = chrono::DateTime::parse_from_rfc3339("2026-10-05T00:00:00Z")
3809 .unwrap()
3810 .timestamp();
3811 let blame = history::query_in_with(
3812 &connection,
3813 &mj_core::history::HistoryRequest {
3814 request_id: "test-blame".into(),
3815 query: mj_core::history::HistoryQuery::BlameFile {
3816 path: PathBuf::from("src/a.rs"),
3817 start_line: 1,
3818 end_line: 1,
3819 },
3820 blame: Some(mj_core::history::BlameEvidence {
3821 repository: PathBuf::from("/src/project"),
3822 relative_path: PathBuf::from("src/a.rs"),
3823 porcelain: format!(
3824 "{} 1 1 1\nauthor-time {blame_time}\n\tline\n",
3825 "a".repeat(40),
3826 ),
3827 }),
3828 },
3829 &ownership,
3830 &crate::import::NativeScanCache::new(),
3831 )
3832 .unwrap();
3833 assert_eq!(blame["runs"][0]["status"], "confident");
3834 assert_eq!(
3835 blame["runs"][0]["sessions"][0]["session_id"],
3836 ids_for_paths[0]
3837 );
3838
3839 let child_brief = history::query_in_with(
3840 &connection,
3841 &mj_core::history::HistoryRequest {
3842 request_id: "test-child-brief".into(),
3843 query: mj_core::history::HistoryQuery::GetSessionBrief {
3844 session_id: ids_for_paths[2].clone(),
3845 max_chars: 200,
3846 },
3847 blame: None,
3848 },
3849 &ownership,
3850 &crate::import::NativeScanCache::new(),
3851 );
3852 assert!(child_brief.unwrap_err().to_string().contains("not found"));
3853 }
3854
3855 #[test]
3856 fn every_query_path_excludes_sub_agents_including_agent_history() {
3857 let _held = tags::testing::lock();
3858 let (_directory, connection) = tags::testing::isolated_index();
3859 for (session_id, kind) in [("main-session", "main"), ("sub-session", "sub")] {
3860 tags::testing::index_row(&connection, session_id, "claude");
3861 connection
3862 .execute(
3863 "UPDATE files SET kind = ?2 WHERE session_id = ?1",
3864 rusqlite::params![session_id, kind],
3865 )
3866 .expect("set the session kind");
3867 connection
3868 .execute(
3869 "INSERT INTO messages(session_id, role, text)
3870 VALUES (?1, 'user', 'fix the bridge derivation zq')",
3871 [session_id],
3872 )
3873 .expect("insert a message");
3874 connection
3875 .execute(
3876 "INSERT INTO msgs(rowid, text) VALUES (?1, 'fix the bridge derivation zq')",
3877 [connection.last_insert_rowid()],
3878 )
3879 .expect("index the message");
3880 }
3881 let ids = |query: &str, include_tool_matches: bool| {
3882 let mut ids: Vec<String> = query_rows_for_test(query, 10, include_tool_matches)
3883 .expect("query the index")
3884 .into_iter()
3885 .map(|row| row.id)
3886 .collect();
3887 ids.sort();
3888 ids
3889 };
3890
3891 for query in ["", "bridge derivation", "zq", "an indexed session"] {
3893 assert_eq!(ids(query, false), ["main-session"], "query {query:?}");
3894 assert_eq!(ids(query, true), ["main-session"], "query {query:?}");
3895 }
3896 }
3897
3898 #[test]
3906 fn a_phrase_only_in_a_sub_agents_transcript_does_not_match_its_parent() {
3907 let _held = tags::testing::lock();
3908 let (_directory, connection) = tags::testing::isolated_index();
3909 let message = |session_id: &str, role: &str, text: &str| {
3910 connection
3911 .execute(
3912 "INSERT INTO messages(session_id, role, text) VALUES (?1, ?2, ?3)",
3913 rusqlite::params![session_id, role, text],
3914 )
3915 .expect("insert a message");
3916 connection
3917 .execute(
3918 "INSERT INTO msgs(rowid, text) VALUES (?1, ?2)",
3919 rusqlite::params![connection.last_insert_rowid(), text],
3920 )
3921 .expect("index the message");
3922 };
3923 for (session_id, kind) in [("parent", "main"), ("child", "sub")] {
3924 tags::testing::index_row(&connection, session_id, "claude");
3925 connection
3926 .execute(
3927 "UPDATE files SET kind = ?2 WHERE session_id = ?1",
3928 rusqlite::params![session_id, kind],
3929 )
3930 .expect("set the session kind");
3931 }
3932 message("parent", "user", "look into the relay journal");
3933 message("parent", "tool", "Task {\"prompt\":\"read the journal\"}");
3934 message("parent", "tool", "the journal uses a quokka checksum");
3935 message(
3936 "parent",
3937 "assistant",
3938 "The journal is fine; the parent zebra ends here.",
3939 );
3940 message("child", "user", "read the journal");
3941 message("child", "assistant", "the journal uses a quokka checksum");
3942
3943 let ids = |query: &str, include_tool_matches: bool| {
3944 let mut ids: Vec<String> = query_rows_for_test(query, 10, include_tool_matches)
3945 .expect("query the index")
3946 .into_iter()
3947 .map(|row| row.id)
3948 .collect();
3949 ids.sort();
3950 ids
3951 };
3952 assert!(
3953 ids("quokka", false).is_empty(),
3954 "{:?}",
3955 ids("quokka", false)
3956 );
3957 assert_eq!(ids("quokka", true), ["parent"]);
3959 assert_eq!(ids("parent zebra", false), ["parent"]);
3960 }
3961
3962 #[test]
3964 fn short_query_scan_also_ignores_tool_only_matches() {
3965 let _held = tags::testing::lock();
3966 let (_directory, connection) = tags::testing::isolated_index();
3967 tags::testing::index_row(&connection, "parent", "claude");
3968 connection
3969 .execute(
3970 "INSERT INTO messages(session_id, role, text) VALUES ('parent', 'tool', 'qx')",
3971 [],
3972 )
3973 .expect("insert a message");
3974 assert!(
3975 query_rows_for_test("qx", 10, false)
3976 .expect("query the index")
3977 .is_empty()
3978 );
3979 }
3980
3981 #[test]
3982 fn full_text_search_fills_its_result_limit_after_skipping_sub_agents() {
3983 let _held = tags::testing::lock();
3984 let (_directory, connection) = tags::testing::isolated_index();
3985 for (id, kind, text) in [
3986 ("sub", "sub", "restic restic restic restic"),
3987 ("main", "main", "restic cleanup"),
3988 ] {
3989 tags::testing::index_row(&connection, id, "codex");
3990 connection
3991 .execute(
3992 "UPDATE files SET kind = ?2 WHERE session_id = ?1",
3993 rusqlite::params![id, kind],
3994 )
3995 .unwrap();
3996 connection
3997 .execute(
3998 "INSERT INTO messages(session_id, role, text) VALUES (?1, 'user', ?2)",
3999 rusqlite::params![id, text],
4000 )
4001 .unwrap();
4002 connection
4003 .execute(
4004 "INSERT INTO msgs(rowid, text) VALUES (?1, ?2)",
4005 rusqlite::params![connection.last_insert_rowid(), text],
4006 )
4007 .unwrap();
4008 }
4009
4010 let rows = query_rows_for_test("restic", 1, false).unwrap();
4011 assert_eq!(
4012 rows.iter().map(|row| row.id.as_str()).collect::<Vec<_>>(),
4013 ["main"]
4014 );
4015 }
4016
4017 fn block_on<F: std::future::Future>(future: F) -> F::Output {
4020 tokio::runtime::Builder::new_current_thread()
4021 .enable_all()
4022 .build()
4023 .unwrap()
4024 .block_on(future)
4025 }
4026
4027 #[test]
4033 fn a_destroy_indexes_the_session_itself_when_the_sync_outlasts_the_wait() {
4034 let _held = tags::testing::lock();
4035 let (_index_dir, _connection) = tags::testing::isolated_index();
4036 let directory = tempfile::tempdir().unwrap();
4037 let session_id = "0123456789abcdef0123456789abcdef";
4038 write_archive(directory.path(), session_id, 1);
4039 let source = adapter(directory.path(), session_id);
4040
4041 let started = Instant::now();
4042 let outcome = block_on(index_before_destroy_with(
4043 std::future::pending::<Result<()>>(),
4044 Duration::from_millis(200),
4045 move || capture_sessions_from(&source, &[session_id.to_owned()]),
4046 Duration::from_millis(50),
4047 ));
4048
4049 assert_eq!(outcome, IndexedBeforeDestroy::WrittenDirectly);
4050 assert!(
4051 started.elapsed() < Duration::from_secs(10),
4052 "the destroy must not wait for the pass: {:?}",
4053 started.elapsed()
4054 );
4055 let found = wiki_session_for_test(session_id, &BTreeSet::new())
4056 .unwrap()
4057 .expect("the session is found by its id");
4058 assert_eq!(found.status, WikiSessionStatus::Archived);
4059 assert_eq!(found.tool, TOOL);
4060 assert_eq!(
4061 found.path,
4062 PathBuf::from(format!("{}/{session_id}", directory.path().display()))
4063 );
4064 assert_eq!(found.title, "the harness title");
4065 assert_eq!(
4066 found.harness,
4067 Some(HarnessKind::Codex),
4068 "the session's metadata is written beside its row"
4069 );
4070 assert!(!found.nothing_to_restore);
4071 }
4072
4073 #[test]
4078 fn a_busy_index_takes_the_destroyed_session_once_it_is_free() {
4079 let _held = tags::testing::lock();
4080 let (_index_dir, writer) = tags::testing::isolated_index();
4081 let directory = tempfile::tempdir().unwrap();
4082 let session_id = "0123456789abcdef0123456789abcdef";
4083 write_archive(directory.path(), session_id, 1);
4084 let source = adapter(directory.path(), session_id);
4085
4086 block_on(async {
4087 writer.execute_batch("BEGIN IMMEDIATE").unwrap();
4088 let outcome = index_before_destroy_with(
4089 std::future::pending::<Result<()>>(),
4090 Duration::from_millis(50),
4091 move || capture_sessions_from(&source, &[session_id.to_owned()]),
4092 Duration::from_millis(50),
4093 )
4094 .await;
4095 assert_eq!(outcome, IndexedBeforeDestroy::Deferred);
4096 assert!(
4097 wiki_session_for_test(session_id, &BTreeSet::new())
4098 .unwrap()
4099 .is_none(),
4100 "nothing is written while the other writer holds the index"
4101 );
4102
4103 writer.execute_batch("COMMIT").unwrap();
4104 let deadline = Instant::now() + Duration::from_secs(30);
4105 while wiki_session_for_test(session_id, &BTreeSet::new())
4106 .unwrap()
4107 .is_none()
4108 {
4109 assert!(
4110 Instant::now() < deadline,
4111 "the deferred row never reached the index"
4112 );
4113 tokio::time::sleep(Duration::from_millis(50)).await;
4114 }
4115 });
4116 }
4117
4118 #[test]
4122 fn a_session_the_index_holds_as_it_is_now_needs_no_indexing() {
4123 let _held = tags::testing::lock();
4124 let (_index_dir, _connection) = tags::testing::isolated_index();
4125 let directory = tempfile::tempdir().unwrap();
4126 let session_id = "0123456789abcdef0123456789abcdef";
4127 let never_prompted = "fedcba9876543210fedcba9876543210";
4128 write_archive(directory.path(), session_id, 1);
4129 let source = adapter(directory.path(), session_id);
4130 let ids = [session_id.to_owned(), never_prompted.to_owned()];
4131
4132 assert_eq!(
4133 unindexed(&source, &ids).unwrap(),
4134 [session_id],
4135 "a session with no conversation has nothing to index"
4136 );
4137 let captured = Arc::new(capture_sessions_from(&source, &ids).unwrap());
4138 write_captured(&captured).unwrap();
4139 assert!(unindexed(&source, &ids).unwrap().is_empty());
4140
4141 source
4143 .sessions
4144 .lock()
4145 .unwrap()
4146 .records
4147 .get_mut(session_id)
4148 .unwrap()
4149 .updated_at = "2099-01-01T00:00:00Z".into();
4150 assert_eq!(unindexed(&source, &ids).unwrap(), [session_id]);
4151 }
4152
4153 #[test]
4155 fn adapters_cover_every_supported_harness_and_no_unsupported_tool() {
4156 use mj_core::config::HarnessProfile;
4157 let directory = tempfile::tempdir().unwrap();
4158 let mut config = Config::default();
4159 for (index, kind) in HarnessKind::ALL.into_iter().enumerate() {
4160 let home = directory.path().join(format!("home-{index}"));
4161 std::fs::create_dir_all(&home).unwrap();
4162 config.profiles.insert(
4163 format!("profile-{index}"),
4164 HarnessProfile {
4165 enabled: true,
4166 kind,
4167 home,
4168 environment: Default::default(),
4169 context_window_bytes: None,
4170 subagents: Default::default(),
4171 guardian_review_model: None,
4172 },
4173 );
4174 }
4175 let mut names: Vec<&str> = native_adapters(&config)
4176 .iter()
4177 .map(|adapter| adapter.name())
4178 .collect();
4179 names.sort_unstable();
4180 assert_eq!(
4181 names,
4182 [
4183 "claude-code",
4184 "codex",
4185 "grok-build",
4186 "kimi-code",
4187 "muse",
4188 "opencode"
4189 ],
4190 "one adapter per supported harness, none for tools Mjolnir cannot run"
4191 );
4192 }
4193
4194 mod text_search {
4196 use super::super::{SessionTextMatch, SessionTextMatchKind, text_matches_in};
4197 use std::collections::BTreeSet;
4198
4199 fn index(sessions: &[(&str, &[(&str, &str)])]) -> rusqlite::Connection {
4200 let connection = rusqlite::Connection::open_in_memory().unwrap();
4201 connection
4202 .execute_batch(
4203 "CREATE TABLE files(path TEXT PRIMARY KEY, session_id TEXT NOT NULL,
4204 tool TEXT NOT NULL, kind TEXT NOT NULL DEFAULT 'main');
4205 CREATE TABLE messages(id INTEGER PRIMARY KEY, session_id TEXT NOT NULL,
4206 role TEXT NOT NULL, text TEXT NOT NULL);
4207 CREATE VIRTUAL TABLE msgs USING fts5(
4208 text, content='messages', content_rowid='id', tokenize='trigram');",
4209 )
4210 .unwrap();
4211 for (id, messages) in sessions {
4212 connection
4213 .execute(
4214 "INSERT INTO files(path, session_id, tool) VALUES (?1, ?2, 'mjolnir')",
4215 rusqlite::params![format!("/checkpoints/{id}"), id],
4216 )
4217 .unwrap();
4218 for (role, text) in *messages {
4219 connection
4220 .execute(
4221 "INSERT INTO messages(session_id, role, text) VALUES (?1, ?2, ?3)",
4222 rusqlite::params![id, role, text],
4223 )
4224 .unwrap();
4225 let rowid = connection.last_insert_rowid();
4226 connection
4227 .execute(
4228 "INSERT INTO msgs(rowid, text) VALUES (?1, ?2)",
4229 rusqlite::params![rowid, text],
4230 )
4231 .unwrap();
4232 }
4233 }
4234 connection
4235 }
4236
4237 fn live(ids: &[&str]) -> BTreeSet<String> {
4238 ids.iter().map(|id| (*id).to_owned()).collect()
4239 }
4240
4241 fn matches(kinds: &[(&str, SessionTextMatchKind)]) -> Vec<SessionTextMatch> {
4242 kinds
4243 .iter()
4244 .map(|(id, kind)| SessionTextMatch {
4245 session_id: (*id).to_owned(),
4246 kind: *kind,
4247 })
4248 .collect()
4249 }
4250
4251 #[test]
4252 fn golden_sessions_filter_search() {
4253 let connection = index(&[
4254 ("said-by-user", &[("user", "please fix the Zebra crossing")]),
4255 ("said-by-agent", &[("assistant", "the zebra is fixed")]),
4256 (
4257 "only-in-tool",
4258 &[("tool", "zebra stack trace"), ("user", "hello")],
4259 ),
4260 (
4261 "both",
4262 &[("assistant", "a ZEBRA appears"), ("user", "a zebra please")],
4263 ),
4264 ("gone", &[("user", "zebra")]),
4265 ]);
4266 let live = live(&["said-by-user", "said-by-agent", "only-in-tool", "both"]);
4267 let found = text_matches_in(
4268 &connection,
4269 "zebra",
4270 &live,
4271 &super::super::top_level::Snapshot::default(),
4272 )
4273 .unwrap();
4274 assert_eq!(
4275 found,
4276 matches(&[
4277 ("both", SessionTextMatchKind::User),
4278 ("said-by-agent", SessionTextMatchKind::Agent),
4279 ("said-by-user", SessionTextMatchKind::User),
4280 ])
4281 );
4282 let rendered = format!(
4283 "=== controller text search response ({} matches) ===\n{}\n",
4284 found.len(),
4285 serde_json::to_string_pretty(&found).unwrap()
4286 );
4287 mj_core::golden::assert_golden(
4288 env!("CARGO_MANIFEST_DIR"),
4289 "sessions-filter-search",
4290 &rendered,
4291 );
4292 }
4293
4294 #[test]
4295 fn a_query_too_short_for_the_trigram_index_still_matches() {
4296 let connection = index(&[
4297 ("a", &[("user", "go to the zoo")]),
4298 ("b", &[("assistant", "zoo")]),
4299 ("c", &[("tool", "zoo")]),
4300 ]);
4301 assert_eq!(
4302 text_matches_in(
4303 &connection,
4304 "zo",
4305 &live(&["a", "b", "c"]),
4306 &super::super::top_level::Snapshot::default(),
4307 )
4308 .unwrap(),
4309 matches(&[
4310 ("a", SessionTextMatchKind::User),
4311 ("b", SessionTextMatchKind::Agent),
4312 ])
4313 );
4314 assert_eq!(
4316 text_matches_in(
4317 &connection,
4318 "%z",
4319 &live(&["a"]),
4320 &super::super::top_level::Snapshot::default(),
4321 )
4322 .unwrap(),
4323 Vec::new()
4324 );
4325 }
4326 }
4327}