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