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 publication: None,
2343 build_cache: None,
2344 container_workspace: None,
2345 mjolnir_subagents: None,
2346 create_managed_worktree: None,
2347 workspace_id: mj_core::workspace::DEFAULT_WORKSPACE_ID.to_owned(),
2348 archived: false,
2349 container_cpus: None,
2350 container_memory: None,
2351 id: "0123456789abcdef0123456789abcdef".into(),
2352 title: "indexed session".into(),
2353 harness_kind: mj_core::config::HarnessKind::Codex,
2354 last_profile: "codex".into(),
2355 bundle_id: "project".into(),
2356 project_directory: Some(PathBuf::from("/home/dev/project")),
2357 managed_worktree: None,
2358 target_template_id: "local-bare".into(),
2359 resource_allocation: None,
2360 additional_mounts: Vec::new(),
2361 state: mj_core::state::SessionState::Stopped,
2362 target: None,
2363 native_session_id: Some("native-session".into()),
2364 acp_session_title: Some("the harness title".into()),
2365 session_title_override: None,
2366 created_at: "2026-09-01T00:00:00Z".into(),
2367 updated_at: "2026-09-01T01:00:00Z".into(),
2368 viewed_through_event_ordinal: 0,
2369 draft_input: String::new(),
2370 last_error: None,
2371 last_checkpoint_error: None,
2372 checkpoint: None,
2373 }
2374 }
2375
2376 #[test]
2377 fn the_newest_checkpoint_of_each_session_is_one_indexed_key() {
2378 let directory = tempfile::tempdir().unwrap();
2379 let session_id = "0123456789abcdef0123456789abcdef";
2380 write_archive(directory.path(), session_id, 1);
2381 write_archive(directory.path(), session_id, 7);
2382 let adapter = adapter(directory.path(), session_id);
2383
2384 let store = adapter.store().expect("the adapter is a shared store");
2385 let key = format!("{}/{session_id}", directory.path().display());
2386 assert_eq!(
2387 store
2388 .keys
2389 .iter()
2390 .map(|(key, _)| key.as_str())
2391 .collect::<Vec<_>>(),
2392 vec![key.as_str()]
2393 );
2394 assert!(!store.had_error);
2395 assert_eq!(store.files.len(), 1);
2396 assert!(
2397 store.files[0]
2398 .file_name()
2399 .unwrap()
2400 .to_str()
2401 .unwrap()
2402 .contains("-7-archive-"),
2403 "the newest checkpoint is the one indexed: {:?}",
2404 store.files[0]
2405 );
2406 assert_eq!(
2407 adapter.reconcile_scope(),
2408 Some(format!("{}/", directory.path().display()))
2409 );
2410
2411 let session = adapter.parse_key(&key).unwrap();
2412 assert_eq!(session.id, session_id);
2413 assert_eq!(session.tool, "mjolnir");
2414 assert_eq!(session.path, PathBuf::from(&key));
2415 assert_eq!(session.project, "/home/dev/project");
2416 assert_eq!(session.title, "the harness title");
2417 assert!(!session.subagent);
2418 assert_eq!(
2419 session.messages.iter().map(|m| m.role).collect::<Vec<_>>(),
2420 vec![Role::User, Role::Tool, Role::Assistant]
2421 );
2422 assert_eq!(session.messages[0].text, "index this session");
2423 let tool: serde_json::Value = serde_json::from_str(&session.messages[1].text).unwrap();
2424 assert_eq!(tool["name"], "Edit");
2425 assert_eq!(tool["call"]["title"], "Edit config.toml");
2426 assert_eq!(session.messages[2].text, "done");
2427 assert_eq!(session.touched, vec!["/old/container/config.toml"]);
2428 }
2429
2430 #[test]
2431 fn provenance_backfill_repairs_an_unchanged_checkpoint_without_rebuilding_the_index() {
2432 let _held = tags::testing::lock();
2433 let (_index_dir, mut connection) = tags::testing::isolated_index();
2434 let directory = tempfile::tempdir().unwrap();
2435 write_archive(directory.path(), "old-session", 4);
2436 let source = adapter(directory.path(), "old-session");
2437 let key = source.key_for("old-session");
2438 tags::testing::index_row(&connection, "old-session", "mjolnir");
2439 connection
2440 .execute(
2441 "UPDATE files SET path = ?1 WHERE session_id = 'old-session'",
2442 [&key],
2443 )
2444 .unwrap();
2445 provenance::backfill(&mut connection, &source).unwrap();
2446 assert_eq!(
2447 sessionwiki::index::files_for(&connection, "old-session").unwrap(),
2448 vec!["/old/container/config.toml"]
2449 );
2450 provenance::backfill(&mut connection, &source).unwrap();
2451 assert_eq!(
2452 sessionwiki::index::sessions_for_file(&connection, "config.toml", 20)
2453 .unwrap()
2454 .len(),
2455 1
2456 );
2457 }
2458
2459 fn projection(session_id: &str) -> mj_core::state::MaterializedSession {
2460 use mj_core::transcript::{TranscriptBody, TranscriptItem};
2461 let mut projected = mj_core::state::MaterializedSession::empty(session_id);
2462 let mut push = |position: u64, body: TranscriptBody| {
2463 let streamed = matches!(body, TranscriptBody::Agent { .. });
2464 projected
2465 .transcript
2466 .push(std::sync::Arc::new(TranscriptItem {
2467 stable_id: format!("item-{position}"),
2468 position,
2469 latest_content_event_ordinal: streamed.then_some(position),
2470 created_at_ms: 1_700_000_000_000 + i64::try_from(position).unwrap(),
2471 last_changed_at_ms: 1_700_000_000_000 + i64::try_from(position).unwrap(),
2472 body,
2473 }));
2474 };
2475 push(
2476 1,
2477 TranscriptBody::User {
2478 content: vec![serde_json::json!({"type": "text", "text": "still talking"})],
2479 },
2480 );
2481 push(
2482 2,
2483 TranscriptBody::Thought {
2484 chunks: vec![serde_json::json!({"content": {"type": "text", "text": "hmm"}})],
2485 streaming: false,
2486 },
2487 );
2488 push(
2489 3,
2490 TranscriptBody::Tool {
2491 call: serde_json::json!({"toolCallId": "c1", "title": "Read README.md"}),
2492 terminal_outputs: Vec::new(),
2493 terminal_refs: Vec::new(),
2494 presentation: None,
2495 },
2496 );
2497 push(
2498 4,
2499 TranscriptBody::Agent {
2500 chunks: vec![serde_json::json!({"content": {"type": "text", "text": "reading"}})],
2501 streaming: false,
2502 },
2503 );
2504 projected.session_title = Some("the live title".into());
2505 projected
2506 }
2507
2508 #[test]
2511 fn a_running_session_is_indexed_from_its_stored_transcript() {
2512 let session_id = "0123456789abcdef0123456789abcdef";
2513 let messages = projected_messages(&projection(session_id));
2514 assert_eq!(
2515 messages.iter().map(|m| m.role).collect::<Vec<_>>(),
2516 vec![Role::User, Role::Tool, Role::Assistant]
2517 );
2518 assert_eq!(messages[0].text, "still talking");
2519 assert_eq!(messages[2].text, "reading");
2520 let tool: serde_json::Value = serde_json::from_str(&messages[1].text).unwrap();
2521 assert_eq!(tool["name"], "Read");
2522 assert_eq!(tool["call"]["title"], "Read README.md");
2523 }
2524
2525 #[test]
2530 fn a_running_session_is_listed_with_its_own_change_token() {
2531 let directory = tempfile::tempdir().unwrap();
2532 let running = "0123456789abcdef0123456789abcdef";
2533 let never_checkpointed = "fedcba9876543210fedcba9876543210";
2534 write_archive(directory.path(), running, 3);
2535 let live = adapter_with_live(
2536 directory.path(),
2537 running,
2538 BTreeMap::from([
2539 (running.to_owned(), 1_900_000_000),
2540 (never_checkpointed.to_owned(), 1_900_000_001),
2541 ]),
2542 );
2543
2544 let store = live.store().expect("the adapter is a shared store");
2545 let key_of = |session_id: &str| format!("{}/{session_id}", directory.path().display());
2546 assert_eq!(
2547 store.keys,
2548 vec![
2549 (
2550 key_of(running),
2551 1_900_000_000 * 1024 + i64::from(mj_transcript::summary::SUMMARY_VERSION)
2552 ),
2553 (
2554 key_of(never_checkpointed),
2555 1_900_000_001 * 1024 + i64::from(mj_transcript::summary::SUMMARY_VERSION)
2556 ),
2557 ],
2558 "a live session's own token replaces the checkpoint's"
2559 );
2560
2561 let stopped = adapter(directory.path(), running);
2564 let keys = stopped.store().expect("a shared store").keys;
2565 assert_eq!(keys.len(), 1);
2566 assert_eq!(keys[0].0, key_of(running));
2567 assert_ne!(keys[0].1, 1_900_000_000);
2568 assert_eq!(
2569 stopped.parse_key(&key_of(running)).unwrap().title,
2570 "the harness title",
2571 "a stopped session is parsed from its checkpoint"
2572 );
2573 }
2574
2575 #[test]
2578 fn a_rename_moves_a_session_change_token() {
2579 let directory = tempfile::tempdir().unwrap();
2580 let session_id = "0123456789abcdef0123456789abcdef";
2581 write_archive(directory.path(), session_id, 1);
2582 let adapter = adapter(directory.path(), session_id);
2583 let before = adapter.store().expect("a shared store").keys[0].1;
2584
2585 {
2586 let mut sessions = adapter.sessions.lock().unwrap();
2587 let record = sessions.records.get_mut(session_id).unwrap();
2588 record.session_title_override = Some("the new name".into());
2589 record.updated_at = "2099-01-01T00:00:00Z".into();
2590 }
2591 let after = adapter.store().expect("a shared store").keys[0].1;
2592 assert!(
2593 after > before,
2594 "a renamed session is re-indexed: {before} then {after}"
2595 );
2596 assert_eq!(
2597 adapter
2598 .parse_key(&format!("{}/{session_id}", directory.path().display()))
2599 .unwrap()
2600 .title,
2601 "the new name"
2602 );
2603 }
2604
2605 fn indexed(messages: Vec<(Role, &str)>) -> sessionwiki::model::Session {
2606 Session {
2607 id: "0123456789abcdef0123456789abcdef".into(),
2608 tool: "mjolnir",
2609 path: PathBuf::from("/sessions/0123456789abcdef0123456789abcdef"),
2610 project: "/home/dev/project".into(),
2611 started: DateTime::from_timestamp_millis(1_700_000_000_000),
2612 ended: None,
2613 title: "the archived session".into(),
2614 subagent: false,
2615 messages: messages
2616 .into_iter()
2617 .map(|(role, text)| Message {
2618 role,
2619 text: text.to_owned(),
2620 ts: None,
2621 })
2622 .collect(),
2623 touched: Vec::new(),
2624 edits: Vec::new(),
2625 }
2626 }
2627
2628 #[test]
2631 fn transcript_hits_locates_case_insensitive_matches() {
2632 let session = indexed(vec![
2633 (Role::User, "Make the Tests green"),
2634 (Role::Assistant, "the tests are green now"),
2635 ]);
2636
2637 let found = hit_transcript(&session, "TESTS", 0, 4_000);
2638
2639 assert_eq!(found.blocks.len(), 2, "both messages contain the query");
2640 assert_eq!(found.blocks[0].role, "user");
2641 let (start, end) = found.blocks[0].hits[0];
2642 assert_eq!(&found.blocks[0].text[start..end], "Tests");
2643 let (start, end) = found.blocks[1].hits[0];
2644 assert_eq!(&found.blocks[1].text[start..end], "tests");
2645 assert!(!found.blocks[0].truncated);
2646 assert_eq!(found.omitted_after, 0);
2647 }
2648
2649 #[test]
2652 fn transcript_hits_keeps_context_and_marks_omissions() {
2653 let session = indexed(vec![
2654 (Role::User, "zero"),
2655 (Role::Assistant, "one needle one"),
2656 (Role::Tool, "two"),
2657 (Role::User, "three"),
2658 (Role::Assistant, "four"),
2659 (Role::Tool, "five"),
2660 (Role::User, "six needle six"),
2661 (Role::Assistant, "seven"),
2662 (Role::User, "eight"),
2663 ]);
2664
2665 let found = hit_transcript(&session, "needle", 1, 4_000);
2666
2667 let shown: Vec<(&str, &str, usize)> = found
2668 .blocks
2669 .iter()
2670 .map(|block| {
2671 (
2672 block.role.as_str(),
2673 block.text.as_str(),
2674 block.omitted_before,
2675 )
2676 })
2677 .collect();
2678 assert_eq!(
2679 shown,
2680 vec![
2681 ("user", "zero", 0),
2682 ("assistant", "one needle one", 0),
2683 ("tool", "two", 0),
2684 ("tool", "five", 2),
2685 ("user", "six needle six", 0),
2686 ("assistant", "seven", 0),
2687 ]
2688 );
2689 assert_eq!(found.omitted_after, 1, "the last message is not shown");
2690 assert!(found.blocks[0].hits.is_empty(), "context has no hits");
2691 }
2692
2693 #[test]
2699 fn transcript_hits_never_anchor_on_tool_output() {
2700 let session = indexed(vec![
2701 (Role::User, "make it build"),
2702 (Role::Tool, "cargo build --needle"),
2703 (Role::Assistant, "it builds"),
2704 ]);
2705
2706 let only_in_a_tool = hit_transcript(&session, "needle", 1, 4_000);
2707 assert!(
2708 only_in_a_tool.blocks.is_empty(),
2709 "tool output must not anchor a passage, got {:?}",
2710 only_in_a_tool.blocks
2711 );
2712
2713 let beside_a_match = hit_transcript(&session, "builds", 1, 4_000);
2714 let shown: Vec<(&str, bool)> = beside_a_match
2715 .blocks
2716 .iter()
2717 .map(|block| (block.role.as_str(), !block.hits.is_empty()))
2718 .collect();
2719 assert_eq!(
2720 shown,
2721 vec![("tool", false), ("assistant", true)],
2722 "a tool message is still context around a real match"
2723 );
2724 }
2725
2726 #[test]
2729 fn transcript_hits_window_keeps_the_first_hit() {
2730 let filler = "x".repeat(4_000);
2731 let session = indexed(vec![(Role::User, &format!("{filler} needle {filler}"))]);
2732
2733 let found = hit_transcript(&session, "needle", 0, 100);
2734
2735 let block = &found.blocks[0];
2736 assert!(block.truncated);
2737 assert_eq!(block.text.chars().count(), 100);
2738 assert_eq!(block.hits.len(), 1, "the windowed text keeps its hit");
2739 let (start, end) = block.hits[0];
2740 assert_eq!(&block.text[start..end], "needle");
2741 assert!(
2742 start >= 20,
2743 "the window keeps lead-in before the hit, got {start}"
2744 );
2745 }
2746
2747 #[test]
2751 fn a_restored_snapshot_is_a_valid_transcript_of_the_indexed_session() {
2752 let snapshot = snapshot_of(&indexed(vec![
2753 (Role::User, "make the tests green"),
2754 (Role::Tool, "Read src/lib.rs"),
2755 (Role::Assistant, "they are green now"),
2756 (Role::User, " "),
2757 ]))
2758 .unwrap();
2759
2760 snapshot.validate().expect("the snapshot is well formed");
2761 assert_eq!(snapshot.event_frontier, 3);
2762 assert_eq!(
2763 snapshot.session.session_title.as_deref(),
2764 Some("the archived session")
2765 );
2766 assert!(snapshot.session.last_activity_at_ms.is_some());
2767 let bodies = snapshot
2768 .transcript
2769 .iter()
2770 .map(|item| match &item.body {
2771 mj_core::archive::CanonicalTranscriptBody::User { content } => (
2772 "user",
2773 mj_core::transcript::materialized_content_text(content),
2774 ),
2775 mj_core::archive::CanonicalTranscriptBody::Agent { chunks, .. } => (
2776 "agent",
2777 mj_core::transcript::materialized_chunks_text(chunks),
2778 ),
2779 mj_core::archive::CanonicalTranscriptBody::Tool { call, .. } => (
2780 "tool",
2781 call["title"].as_str().unwrap_or_default().to_owned(),
2782 ),
2783 _ => ("other", String::new()),
2784 })
2785 .collect::<Vec<_>>();
2786 assert_eq!(
2787 bodies,
2788 vec![
2789 ("user", "make the tests green".to_owned()),
2790 ("tool", "Read src/lib.rs".to_owned()),
2791 ("agent", "they are green now".to_owned()),
2792 ],
2793 "the blank message is dropped and every other one keeps its role"
2794 );
2795 }
2796
2797 #[test]
2801 fn messages_before_the_first_prompt_are_dropped() {
2802 let snapshot = snapshot_of(&indexed(vec![
2803 (Role::Assistant, "still working"),
2804 (Role::User, "carry on"),
2805 ]))
2806 .unwrap();
2807 assert_eq!(snapshot.transcript.len(), 1);
2808 assert_eq!(snapshot.transcript[0].position, 1);
2809 snapshot.validate().unwrap();
2810
2811 let error = snapshot_of(&indexed(vec![(Role::Assistant, "nobody asked")])).unwrap_err();
2812 assert!(
2813 error.to_string().contains("no prompt"),
2814 "a session with no prompt cannot be restored: {error}"
2815 );
2816 assert!(!has_prompt(&indexed(vec![(
2818 Role::Assistant,
2819 "nobody asked"
2820 )])));
2821 assert!(has_prompt(&indexed(vec![(Role::User, "carry on")])));
2822 }
2823
2824 fn record(
2825 session_id: &str,
2826 state: mj_core::state::SessionState,
2827 updated_at: &str,
2828 ) -> SessionRecord {
2829 SessionRecord {
2830 id: session_id.into(),
2831 state,
2832 updated_at: updated_at.into(),
2833 ..record_template()
2834 }
2835 }
2836
2837 fn child(child_session_id: &str, parent_session_id: &str) -> mj_core::subagent::SubagentRecord {
2838 mj_core::subagent::SubagentRecord {
2839 child_session_id: child_session_id.into(),
2840 parent_session_id: parent_session_id.into(),
2841 task_name: "task".into(),
2842 profile_id: "codex".into(),
2843 model: None,
2844 effort: None,
2845 working_directory: PathBuf::new(),
2846 initial_prompt: "do the thing".into(),
2847 request_key: "key".into(),
2848 created_at: "2026-09-01T00:00:00Z".into(),
2849 noticed_turn: None,
2850 handback_tool: false,
2851 }
2852 }
2853
2854 fn ready(
2855 sessions: Vec<SessionRecord>,
2856 children: Vec<mj_core::subagent::SubagentRecord>,
2857 ) -> Vec<String> {
2858 let now = parse_time("2026-09-10T00:00:00Z").unwrap();
2859 sessions_ready_to_archive(
2860 &sessions
2861 .into_iter()
2862 .map(|record| (record.id.clone(), record))
2863 .collect(),
2864 &children
2865 .into_iter()
2866 .map(|child| (child.child_session_id.clone(), child))
2867 .collect(),
2868 now,
2869 3,
2870 )
2871 }
2872
2873 #[test]
2874 fn aged_clone_requires_clean_published_evidence_for_its_current_checkpoint() {
2875 let id = "0123456789abcdef0123456789abcdef";
2876 let root = PathBuf::from(format!("/srv/project/.mj/clones/{id}"));
2877 let mut session = record(
2878 id,
2879 mj_core::state::SessionState::Stopped,
2880 "2026-09-01T00:00:00Z",
2881 );
2882 session.project_directory = Some(root.clone());
2883 session.managed_worktree = Some(mj_core::state::ManagedWorktree {
2884 kind: mj_core::state::ManagedCheckoutKind::Clone,
2885 source_project_directory: "/srv/project".into(),
2886 source_repository: "/srv/project".into(),
2887 worktree_root: root,
2888 branch: "feature".into(),
2889 target: mj_core::state::ManagedWorktreeTarget::Local,
2890 base_commit: Some("1".repeat(40)),
2891 });
2892 session.checkpoint = Some(mj_core::state::CheckpointMetadata {
2893 archive_path: "sessions/checkpoint.hel.zip".into(),
2894 sha256: "a".repeat(64),
2895 created_at: "2026-09-01T00:00:00Z".into(),
2896 event_frontier: 0,
2897 });
2898 assert!(ready(vec![session.clone()], vec![]).is_empty());
2899 session.publication = Some(mj_core::state::PublicationAssessment {
2900 checkpoint_sha256: "a".repeat(64),
2901 state: mj_core::state::PublicationState::Published,
2902 dirty: false,
2903 stashed: false,
2904 saved_commits: vec!["2".repeat(40)],
2905 destinations: vec!["https://example.test/repository.git".into()],
2906 checked_at: "2026-09-01T01:00:00Z".into(),
2907 reason: Some("feature branch was pushed but not merged".into()),
2908 });
2909 assert_eq!(ready(vec![session.clone()], vec![]), vec![id]);
2910 session.publication.as_mut().unwrap().stashed = true;
2911 assert!(ready(vec![session.clone()], vec![]).is_empty());
2912 session.publication.as_mut().unwrap().stashed = false;
2913 session.publication.as_mut().unwrap().checkpoint_sha256 = "b".repeat(64);
2914 assert!(ready(vec![session], vec![]).is_empty());
2915 }
2916
2917 fn sized_session(
2919 root: &Path,
2920 session_id: &str,
2921 updated_at: &str,
2922 checkpoint_bytes: usize,
2923 attachment_bytes: &[usize],
2924 ) -> SessionRecord {
2925 let archive_path = root.join(format!("{session_id}.hel.zip"));
2926 std::fs::write(&archive_path, vec![b'c'; checkpoint_bytes]).unwrap();
2927 if !attachment_bytes.is_empty() {
2928 let attachments = root
2929 .join(session_id)
2930 .join(mj_core::attachment::ATTACHMENT_DIR);
2931 std::fs::create_dir_all(&attachments).unwrap();
2932 for (index, size) in attachment_bytes.iter().enumerate() {
2933 std::fs::write(attachments.join(format!("{index}.png")), vec![b'a'; *size])
2934 .unwrap();
2935 }
2936 }
2937 SessionRecord {
2938 checkpoint: Some(mj_core::state::CheckpointMetadata {
2939 archive_path,
2940 sha256: "0".repeat(64),
2941 created_at: updated_at.into(),
2942 event_frontier: 1,
2943 }),
2944 ..record(
2945 session_id,
2946 mj_core::state::SessionState::Stopped,
2947 updated_at,
2948 )
2949 }
2950 }
2951
2952 #[test]
2953 fn the_space_preview_sizes_every_session_and_only_the_aged_ones_as_reclaimable() {
2954 let directory = tempfile::tempdir().unwrap();
2955 let root = directory.path();
2956 let sessions: BTreeMap<String, SessionRecord> = [
2957 sized_session(root, "old-stopped", "2026-09-01T00:00:00Z", 1000, &[10, 20]),
2958 sized_session(root, "just-stopped", "2026-09-09T00:00:00Z", 500, &[]),
2959 SessionRecord {
2962 checkpoint: Some(mj_core::state::CheckpointMetadata {
2963 archive_path: root.join("missing.hel.zip"),
2964 sha256: "0".repeat(64),
2965 created_at: "2026-09-01T00:00:00Z".into(),
2966 event_frontier: 1,
2967 }),
2968 ..record(
2969 "lost-checkpoint",
2970 mj_core::state::SessionState::Stopped,
2971 "2026-09-01T00:00:00Z",
2972 )
2973 },
2974 ]
2975 .into_iter()
2976 .map(|record| (record.id.clone(), record))
2977 .collect();
2978 let now = parse_time("2026-09-10T00:00:00Z").unwrap();
2979
2980 let all = archive_space_over(root, &sessions, &BTreeMap::new(), now, None);
2981 assert_eq!(all.sessions, 3);
2982 assert_eq!(all.bytes, 1530);
2983 assert_eq!(all.reclaimable_sessions, 0);
2984 assert_eq!(all.reclaimable_bytes, 0);
2985
2986 let aged = archive_space_over(root, &sessions, &BTreeMap::new(), now, Some(3));
2987 assert_eq!(aged.bytes, 1530);
2988 assert_eq!(
2989 (aged.reclaimable_sessions, aged.reclaimable_bytes),
2990 (2, 1030),
2991 "only the sessions the job would archive count, attachments included"
2992 );
2993 }
2994
2995 #[test]
2996 fn only_stopped_sessions_past_the_cut_off_are_archived() {
2997 use mj_core::state::SessionState;
2998 let selected = ready(
2999 vec![
3000 record("old-stopped", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3001 record(
3002 "just-stopped",
3003 SessionState::Stopped,
3004 "2026-09-09T00:00:00Z",
3005 ),
3006 record("old-running", SessionState::Running, "2026-09-01T00:00:00Z"),
3007 record("old-error", SessionState::Error, "2026-09-01T00:00:00Z"),
3008 record("unparsable", SessionState::Stopped, "not a time"),
3009 record("at-the-edge", SessionState::Stopped, "2026-09-07T00:00:00Z"),
3011 ],
3012 Vec::new(),
3013 );
3014 assert_eq!(selected, vec!["at-the-edge", "old-stopped"]);
3015 }
3016
3017 #[test]
3018 fn a_child_the_pass_is_not_archiving_holds_its_parent_back() {
3019 use mj_core::state::SessionState;
3020 let selected = ready(
3021 vec![
3022 record("parent", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3023 record(
3024 "running-child",
3025 SessionState::Running,
3026 "2026-09-01T00:00:00Z",
3027 ),
3028 ],
3029 vec![child("running-child", "parent")],
3030 );
3031 assert!(selected.is_empty(), "the parent must wait: {selected:?}");
3032
3033 let selected = ready(
3034 vec![
3035 record("parent", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3036 record("young-child", SessionState::Stopped, "2026-09-09T00:00:00Z"),
3037 ],
3038 vec![child("young-child", "parent")],
3039 );
3040 assert!(selected.is_empty(), "the parent must wait: {selected:?}");
3041
3042 let selected = ready(
3044 vec![record(
3045 "parent",
3046 SessionState::Stopped,
3047 "2026-09-01T00:00:00Z",
3048 )],
3049 vec![child("departed-child", "parent")],
3050 );
3051 assert_eq!(selected, vec!["parent"]);
3052 }
3053
3054 #[test]
3055 fn children_are_archived_before_their_parents() {
3056 use mj_core::state::SessionState;
3057 let selected = ready(
3058 vec![
3059 record("parent", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3060 record("child", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3061 record("grandchild", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3062 ],
3063 vec![child("child", "parent"), child("grandchild", "child")],
3064 );
3065 assert_eq!(selected, vec!["grandchild", "child", "parent"]);
3066 }
3067
3068 #[test]
3069 fn native_adapters_cover_every_enabled_profile_home() {
3070 use mj_core::config::{Config, HarnessKind, HarnessProfile};
3071
3072 fn profile(kind: HarnessKind, home: &str, enabled: bool) -> HarnessProfile {
3073 HarnessProfile {
3074 enabled,
3075 kind,
3076 home: PathBuf::from(home),
3077 environment: BTreeMap::new(),
3078 context_window_bytes: None,
3079 guardian_review_model: None,
3080 }
3081 }
3082
3083 let mut config = Config::default();
3084 for (id, built) in [
3085 (
3086 "codex",
3087 profile(HarnessKind::Codex, "/home/dev/.codex3", true),
3088 ),
3089 (
3090 "codex-ds",
3091 profile(HarnessKind::Codex, "/home/dev/.codex-ds", true),
3092 ),
3093 (
3095 "codex-alt",
3096 profile(HarnessKind::Codex, "/home/dev/.codex3", true),
3097 ),
3098 (
3099 "codex-off",
3100 profile(HarnessKind::Codex, "/home/dev/.codex-off", false),
3101 ),
3102 (
3103 "claude",
3104 profile(HarnessKind::Claude, "/home/dev/.claude4", true),
3105 ),
3106 ("kimi", profile(HarnessKind::Kimi, "/home/dev/.kimi", true)),
3107 ("grok", profile(HarnessKind::Grok, "/home/dev/.grok", true)),
3108 ("muse", profile(HarnessKind::Muse, "/home/dev/muse", true)),
3109 (
3110 "muse-off",
3111 profile(HarnessKind::Muse, "/home/dev/muse-off", false),
3112 ),
3113 ] {
3114 config.profiles.insert(id.into(), built);
3115 }
3116
3117 let adapters = native_adapters(&config);
3118 let roots: Vec<(&str, Option<PathBuf>)> = adapters
3119 .iter()
3120 .map(|adapter| (adapter.name(), adapter.root()))
3121 .collect();
3122
3123 let codex: Vec<&Option<PathBuf>> = roots
3124 .iter()
3125 .filter(|(name, _)| *name == "codex")
3126 .map(|(_, root)| root)
3127 .collect();
3128 assert_eq!(
3129 codex,
3130 vec![
3131 &Some(PathBuf::from("/home/dev/.codex3/sessions")),
3132 &Some(PathBuf::from("/home/dev/.codex-ds/sessions")),
3133 ],
3134 "one adapter per enabled Codex home, deduplicated: {roots:?}"
3135 );
3136
3137 let claude: Vec<&Option<PathBuf>> = roots
3138 .iter()
3139 .filter(|(name, _)| *name == "claude-code")
3140 .map(|(_, root)| root)
3141 .collect();
3142 assert_eq!(
3143 claude,
3144 vec![&Some(PathBuf::from("/home/dev/.claude4/projects"))],
3145 "one adapter for the enabled Claude home: {roots:?}"
3146 );
3147
3148 for (_, root) in &roots {
3149 let Some(root) = root else { continue };
3150 let text = root.to_string_lossy();
3151 assert!(
3152 !text.contains(".codex-off"),
3153 "a disabled profile must not be indexed: {roots:?}"
3154 );
3155 assert!(
3156 !text.ends_with("/.codex/sessions") && !text.ends_with("/.claude/projects"),
3157 "the stock homes are not indexed unless a profile names them: {roots:?}"
3158 );
3159 }
3160
3161 for (name, root) in [
3164 ("kimi-code", PathBuf::from("/home/dev/.kimi/sessions")),
3165 ("grok-build", PathBuf::from("/home/dev/.grok/sessions")),
3166 (
3167 "muse",
3168 mj_checkpoint::native::muse_sessions_root(Path::new("/home/dev/muse")).unwrap(),
3169 ),
3170 ] {
3171 let found: Vec<&Option<PathBuf>> = roots
3172 .iter()
3173 .filter(|(found, _)| *found == name)
3174 .map(|(_, root)| root)
3175 .collect();
3176 assert_eq!(found, vec![&Some(root)], "one {name} adapter: {roots:?}");
3177 }
3178
3179 for (_, root) in &roots {
3180 let Some(root) = root else { continue };
3181 assert!(
3182 !root.to_string_lossy().contains("muse-off"),
3183 "a disabled profile must not be indexed: {roots:?}"
3184 );
3185 }
3186
3187 assert!(
3188 roots.iter().any(|(name, _)| *name == "gemini"),
3189 "the other built-in adapters are kept: {roots:?}"
3190 );
3191 }
3192
3193 #[test]
3197 fn query_rows_returns_the_indexed_target_profile_and_harness() {
3198 let _held = tags::testing::lock();
3199 let (_directory, connection) = tags::testing::isolated_index();
3200 tags::testing::index_row(&connection, "mj-session", TOOL);
3201 tags::testing::index_row(&connection, "codex-session", "codex");
3202 tags::write(
3203 &connection,
3204 "mj-session",
3205 &tags::MjTags {
3206 target: Some("Prod-Box".into()),
3207 profile: Some("codex-Main".into()),
3208 harness: Some("codex".into()),
3209 },
3210 )
3211 .expect("write the session metadata");
3212
3213 let rows = query_rows("", 10, &BTreeSet::new(), false).expect("query the index");
3214 let mjolnir = rows
3215 .iter()
3216 .find(|row| row.id == "mj-session")
3217 .expect("the Mjolnir row is returned");
3218 assert_eq!(mjolnir.target.as_deref(), Some("Prod-Box"));
3219 assert_eq!(mjolnir.profile.as_deref(), Some("codex-Main"));
3220 assert_eq!(mjolnir.harness.as_deref(), Some("codex"));
3221
3222 let codex = rows
3223 .iter()
3224 .find(|row| row.id == "codex-session")
3225 .expect("the Codex row is returned");
3226 assert_eq!(codex.target, None);
3227 assert_eq!(codex.profile, None);
3228 assert_eq!(codex.harness, None);
3229 }
3230
3231 #[test]
3232 fn one_flag_keeps_sub_agents_out_of_every_query_path() {
3233 let _held = tags::testing::lock();
3234 let (_directory, connection) = tags::testing::isolated_index();
3235 for (session_id, kind) in [("main-session", "main"), ("sub-session", "sub")] {
3236 tags::testing::index_row(&connection, session_id, "claude");
3237 connection
3238 .execute(
3239 "UPDATE files SET kind = ?2 WHERE session_id = ?1",
3240 rusqlite::params![session_id, kind],
3241 )
3242 .expect("set the session kind");
3243 connection
3244 .execute(
3245 "INSERT INTO messages(session_id, role, text)
3246 VALUES (?1, 'user', 'fix the bridge derivation zq')",
3247 [session_id],
3248 )
3249 .expect("insert a message");
3250 connection
3251 .execute(
3252 "INSERT INTO msgs(rowid, text) VALUES (?1, 'fix the bridge derivation zq')",
3253 [connection.last_insert_rowid()],
3254 )
3255 .expect("index the message");
3256 }
3257 let ids = |query: &str, include_subagents: bool| {
3258 let mut ids: Vec<String> = query_rows(query, 10, &BTreeSet::new(), include_subagents)
3259 .expect("query the index")
3260 .into_iter()
3261 .map(|row| row.id)
3262 .collect();
3263 ids.sort();
3264 ids
3265 };
3266
3267 for query in ["", "bridge derivation", "zq", "an indexed session"] {
3269 assert_eq!(ids(query, false), ["main-session"], "query {query:?}");
3270 assert_eq!(
3271 ids(query, true),
3272 ["main-session", "sub-session"],
3273 "query {query:?}"
3274 );
3275 }
3276 }
3277
3278 #[test]
3285 fn a_phrase_only_in_a_sub_agents_transcript_does_not_match_its_parent() {
3286 let _held = tags::testing::lock();
3287 let (_directory, connection) = tags::testing::isolated_index();
3288 let message = |session_id: &str, role: &str, text: &str| {
3289 connection
3290 .execute(
3291 "INSERT INTO messages(session_id, role, text) VALUES (?1, ?2, ?3)",
3292 rusqlite::params![session_id, role, text],
3293 )
3294 .expect("insert a message");
3295 connection
3296 .execute(
3297 "INSERT INTO msgs(rowid, text) VALUES (?1, ?2)",
3298 rusqlite::params![connection.last_insert_rowid(), text],
3299 )
3300 .expect("index the message");
3301 };
3302 for (session_id, kind) in [("parent", "main"), ("child", "sub")] {
3303 tags::testing::index_row(&connection, session_id, "claude");
3304 connection
3305 .execute(
3306 "UPDATE files SET kind = ?2 WHERE session_id = ?1",
3307 rusqlite::params![session_id, kind],
3308 )
3309 .expect("set the session kind");
3310 }
3311 message("parent", "user", "look into the relay journal");
3312 message("parent", "tool", "Task {\"prompt\":\"read the journal\"}");
3313 message("parent", "tool", "the journal uses a quokka checksum");
3314 message(
3315 "parent",
3316 "assistant",
3317 "The journal is fine; the parent zebra ends here.",
3318 );
3319 message("child", "user", "read the journal");
3320 message("child", "assistant", "the journal uses a quokka checksum");
3321
3322 let ids = |query: &str, include_subagents: bool| {
3323 let mut ids: Vec<String> = query_rows(query, 10, &BTreeSet::new(), include_subagents)
3324 .expect("query the index")
3325 .into_iter()
3326 .map(|row| row.id)
3327 .collect();
3328 ids.sort();
3329 ids
3330 };
3331 assert!(
3332 ids("quokka", false).is_empty(),
3333 "{:?}",
3334 ids("quokka", false)
3335 );
3336 assert_eq!(ids("quokka", true), ["child", "parent"]);
3338 assert_eq!(ids("parent zebra", false), ["parent"]);
3339 }
3340
3341 #[test]
3342 fn short_query_scan_also_ignores_tool_only_matches() {
3343 let _held = tags::testing::lock();
3344 let (_directory, connection) = tags::testing::isolated_index();
3345 tags::testing::index_row(&connection, "parent", "claude");
3346 connection
3347 .execute(
3348 "INSERT INTO messages(session_id, role, text) VALUES ('parent', 'tool', 'qx')",
3349 [],
3350 )
3351 .expect("insert a message");
3352 assert!(
3353 query_rows("qx", 10, &BTreeSet::new(), false)
3354 .expect("query the index")
3355 .is_empty()
3356 );
3357 }
3358
3359 #[test]
3360 fn full_text_search_fills_its_result_limit_after_skipping_sub_agents() {
3361 let _held = tags::testing::lock();
3362 let (_directory, connection) = tags::testing::isolated_index();
3363 for (id, kind, text) in [
3364 ("sub", "sub", "restic restic restic restic"),
3365 ("main", "main", "restic cleanup"),
3366 ] {
3367 tags::testing::index_row(&connection, id, "codex");
3368 connection
3369 .execute(
3370 "UPDATE files SET kind = ?2 WHERE session_id = ?1",
3371 rusqlite::params![id, kind],
3372 )
3373 .unwrap();
3374 connection
3375 .execute(
3376 "INSERT INTO messages(session_id, role, text) VALUES (?1, 'user', ?2)",
3377 rusqlite::params![id, text],
3378 )
3379 .unwrap();
3380 connection
3381 .execute(
3382 "INSERT INTO msgs(rowid, text) VALUES (?1, ?2)",
3383 rusqlite::params![connection.last_insert_rowid(), text],
3384 )
3385 .unwrap();
3386 }
3387
3388 let rows = query_rows("restic", 1, &BTreeSet::new(), false).unwrap();
3389 assert_eq!(
3390 rows.iter().map(|row| row.id.as_str()).collect::<Vec<_>>(),
3391 ["main"]
3392 );
3393 }
3394
3395 fn block_on<F: std::future::Future>(future: F) -> F::Output {
3398 tokio::runtime::Builder::new_current_thread()
3399 .enable_all()
3400 .build()
3401 .unwrap()
3402 .block_on(future)
3403 }
3404
3405 #[test]
3410 fn a_destroy_indexes_the_session_itself_when_the_sync_outlasts_the_wait() {
3411 let _held = tags::testing::lock();
3412 let (_index_dir, _connection) = tags::testing::isolated_index();
3413 let directory = tempfile::tempdir().unwrap();
3414 let session_id = "0123456789abcdef0123456789abcdef";
3415 write_archive(directory.path(), session_id, 1);
3416 let source = adapter(directory.path(), session_id);
3417
3418 let started = Instant::now();
3419 let outcome = block_on(index_before_destroy_with(
3420 std::future::pending::<Result<()>>(),
3421 Duration::from_millis(200),
3422 move || capture_sessions_from(&source, &[session_id.to_owned()]),
3423 Duration::from_millis(50),
3424 ));
3425
3426 assert_eq!(outcome, IndexedBeforeDestroy::WrittenDirectly);
3427 assert!(
3428 started.elapsed() < Duration::from_secs(10),
3429 "the destroy must not wait for the pass: {:?}",
3430 started.elapsed()
3431 );
3432 let found = wiki_session(session_id, &BTreeSet::new())
3433 .unwrap()
3434 .expect("the session is found by its id");
3435 assert_eq!(found.status, WikiSessionStatus::Archived);
3436 assert_eq!(found.tool, TOOL);
3437 assert_eq!(
3438 found.path,
3439 PathBuf::from(format!("{}/{session_id}", directory.path().display()))
3440 );
3441 assert_eq!(found.title, "the harness title");
3442 assert_eq!(
3443 found.harness,
3444 Some(HarnessKind::Codex),
3445 "the session's metadata is written beside its row"
3446 );
3447 assert!(!found.nothing_to_restore);
3448 }
3449
3450 #[test]
3453 fn a_sync_that_finishes_in_time_is_all_a_destroy_waits_for() {
3454 let outcome = block_on(index_before_destroy_with(
3455 async { Ok(()) },
3456 DESTROY_SYNC_WAIT,
3457 || -> Result<Vec<CapturedSession>> {
3458 panic!("a finished pass leaves nothing to index on its own")
3459 },
3460 Duration::from_millis(50),
3461 ));
3462 assert_eq!(outcome, IndexedBeforeDestroy::Synced);
3463 }
3464
3465 #[test]
3469 fn a_busy_index_takes_the_destroyed_session_once_it_is_free() {
3470 let _held = tags::testing::lock();
3471 let (_index_dir, writer) = tags::testing::isolated_index();
3472 let directory = tempfile::tempdir().unwrap();
3473 let session_id = "0123456789abcdef0123456789abcdef";
3474 write_archive(directory.path(), session_id, 1);
3475 let source = adapter(directory.path(), session_id);
3476
3477 block_on(async {
3478 writer.execute_batch("BEGIN IMMEDIATE").unwrap();
3479 let outcome = index_before_destroy_with(
3480 std::future::pending::<Result<()>>(),
3481 Duration::from_millis(50),
3482 move || capture_sessions_from(&source, &[session_id.to_owned()]),
3483 Duration::from_millis(50),
3484 )
3485 .await;
3486 assert_eq!(outcome, IndexedBeforeDestroy::Deferred);
3487 assert!(
3488 wiki_session(session_id, &BTreeSet::new())
3489 .unwrap()
3490 .is_none(),
3491 "nothing is written while the other writer holds the index"
3492 );
3493
3494 writer.execute_batch("COMMIT").unwrap();
3495 let deadline = Instant::now() + Duration::from_secs(30);
3496 while wiki_session(session_id, &BTreeSet::new())
3497 .unwrap()
3498 .is_none()
3499 {
3500 assert!(
3501 Instant::now() < deadline,
3502 "the deferred row never reached the index"
3503 );
3504 tokio::time::sleep(Duration::from_millis(50)).await;
3505 }
3506 });
3507 }
3508
3509 #[test]
3513 fn a_session_the_index_holds_as_it_is_now_needs_no_indexing() {
3514 let _held = tags::testing::lock();
3515 let (_index_dir, _connection) = tags::testing::isolated_index();
3516 let directory = tempfile::tempdir().unwrap();
3517 let session_id = "0123456789abcdef0123456789abcdef";
3518 let never_prompted = "fedcba9876543210fedcba9876543210";
3519 write_archive(directory.path(), session_id, 1);
3520 let source = adapter(directory.path(), session_id);
3521 let ids = [session_id.to_owned(), never_prompted.to_owned()];
3522
3523 assert_eq!(
3524 unindexed(&source, &ids).unwrap(),
3525 [session_id],
3526 "a session with no conversation has nothing to index"
3527 );
3528 let captured = Arc::new(capture_sessions_from(&source, &ids).unwrap());
3529 write_captured(&captured).unwrap();
3530 assert!(unindexed(&source, &ids).unwrap().is_empty());
3531
3532 source
3534 .sessions
3535 .lock()
3536 .unwrap()
3537 .records
3538 .get_mut(session_id)
3539 .unwrap()
3540 .updated_at = "2099-01-01T00:00:00Z".into();
3541 assert_eq!(unindexed(&source, &ids).unwrap(), [session_id]);
3542 }
3543
3544 #[test]
3545 fn a_session_tree_holds_the_sub_agents_below_it_and_nothing_else() {
3546 let record = |id: &str| {
3547 (
3548 id.to_owned(),
3549 SessionRecord {
3550 id: id.into(),
3551 ..record_template()
3552 },
3553 )
3554 };
3555 let state = State {
3556 sessions: BTreeMap::from([
3557 record("parent"),
3558 record("child"),
3559 record("grandchild"),
3560 record("sibling"),
3561 ]),
3562 subagents: BTreeMap::from([
3563 ("child".to_owned(), child("child", "parent")),
3564 ("grandchild".to_owned(), child("grandchild", "child")),
3565 ("sibling".to_owned(), child("sibling", "other-parent")),
3566 ]),
3567 ..State::default()
3568 };
3569 assert_eq!(
3570 session_tree(&state, "parent"),
3571 ["parent", "child", "grandchild"]
3572 );
3573 assert!(session_tree(&state, "unknown").is_empty());
3574 }
3575}