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