Skip to main content

mj_controller/
sessionwiki.rs

1//! Publishing Mjolnir's own checkpointed sessions into the user's SessionWiki
2//! index, and the daemon-side job that keeps that index current.
3//!
4//! SessionWiki keeps one searchable index of AI coding sessions across every
5//! tool a user runs. Mjolnir links it as a library and registers
6//! [`MjolnirAdapter`] beside SessionWiki's built-in adapters, so a Mjolnir
7//! session is searchable next to a Claude Code or Codex one. The adapter is a
8//! "shared store" adapter: checkpoints are not one-file-per-session in a shape
9//! SessionWiki can parse, so the indexer enumerates sessions by key and asks
10//! this adapter to parse the ones whose checkpoint changed.
11
12mod 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    SessionTextMatch, SessionTextMatchKind, WikiHitBlock, WikiHitTranscript, WikiIndexState,
28    WikiRow, WikiSessionInfo, WikiSessionStatus, 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
39/// The tool name every Mjolnir instance publishes under. One name means one
40/// search partition; reconciliation is scoped per instance instead (see
41/// [`Adapter::reconcile_scope`]).
42const TOOL: &str = "mjolnir";
43
44/// One checkpoint archive on disk, reduced to what indexing needs.
45struct ArchiveFile {
46    path: PathBuf,
47    frontier: u64,
48    /// Modification time in epoch seconds, SessionWiki's change token.
49    token: i64,
50}
51
52/// What indexing needs from controller state: which sessions exist, which of
53/// them are sub-agent children, and which are still running with their
54/// conversation in the daemon's own database rather than in a checkpoint.
55#[derive(Default)]
56struct Sessions {
57    records: mj_core::snapshot_map::SnapshotMap<String, SessionRecord>,
58    subagent_ids: BTreeSet<String>,
59    /// Session id to change token, for sessions indexed from the projection.
60    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
73/// The change token of every session whose transcript is still only in the
74/// daemon's database: its activity watermark in whole seconds.
75///
76/// A stopped session keeps being indexed from its checkpoint, which never
77/// changes again. Everything else is indexed from the projection, so a running
78/// session is findable before it has ever been closed.
79fn 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
98/// Mjolnir's sessions, as SessionWiki sees them.
99pub struct MjolnirAdapter {
100    sessions_dir: PathBuf,
101    sessions: std::sync::Mutex<Sessions>,
102    /// Re-read controller state when the indexer reaches this adapter.
103    reload: bool,
104}
105
106impl MjolnirAdapter {
107    /// A fixed view of the given state, which is what a caller with a state in
108    /// hand wants.
109    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    /// The same, but re-reading controller state when the indexer reaches this
118    /// adapter.
119    ///
120    /// One sync pass walks every other tool's store first, which can take
121    /// minutes on a large corpus. Without the reload, sessions that closed
122    /// during that walk would be indexed with no record: no project, no start
123    /// time, and the title guessed from the first prompt. Their checkpoints do
124    /// not change afterwards, so nothing would ever correct them.
125    pub fn reloading(state: &State) -> Self {
126        Self {
127            reload: true,
128            ..Self::from_state(state)
129        }
130    }
131
132    /// Mjolnir's own metadata for every session this adapter knows about, to
133    /// be stored in the index beside the transcripts.
134    ///
135    /// Read from the adapter's own snapshot rather than from the controller
136    /// state the sync loaded, because [`MjolnirAdapter::reloading`] replaces
137    /// that snapshot when the indexer reaches this adapter. A session that
138    /// closed during a long first pass is indexed from the reloaded state, so
139    /// its metadata has to come from the same state that produced its row.
140    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    /// The conversation of a stopped session, read from its newest checkpoint,
180    /// with the title the checkpoint recorded.
181    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    /// The conversation of a session that has not stopped, read from the
207    /// daemon's own projection. It is the same conversation the checkpoint
208    /// would hold, minus whatever has not happened yet.
209    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    /// The stable key for one session: its checkpoint directory and id. The
227    /// directory is per instance, which is what scopes reconciliation.
228    fn key_for(&self, session_id: &str) -> String {
229        format!("{}/{session_id}", self.sessions_dir.display())
230    }
231
232    /// The newest checkpoint of every session in the directory, by session id.
233    ///
234    /// `had_error` is true when the directory exists but could not be read in
235    /// full; the indexer then skips deletion reconciliation rather than
236    /// archiving every Mjolnir session off a partial listing.
237    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
292/// The session a checkpoint file name belongs to, with its generation.
293///
294/// Managed checkpoints carry a frontier and a nonce; an imported archive is
295/// named for its session alone and counts as generation zero.
296fn 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
308/// A running session's conversation, as SessionWiki stores it.
309fn 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
330/// One indexed message, or nothing when the item carried no text.
331fn 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    /// Unused: this is a shared-store adapter, so the indexer enumerates
356    /// sessions through [`Adapter::store`] instead of walking files.
357    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        // A session that is still running is indexed from the projection, and
378        // its own token replaces any checkpoint token it has: the conversation
379        // has moved on since that checkpoint was written. Listing it also
380        // keeps reconciliation from archiving a running session.
381        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        // A rename changes the record and not the conversation, so the
388        // record's own last update is part of the change token. Without it a
389        // renamed session would keep its old title in the index for as long as
390        // its transcript stood still.
391        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    /// Every Mjolnir instance publishes under one tool name, so this instance
420    /// speaks only for keys under its own checkpoint directory. Without the
421    /// scope, two instances would archive each other's rows on every sync.
422    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
475/// The Mjolnir adapter handed to the indexer while the sync keeps its own
476/// handle on it.
477///
478/// The indexer takes `Box<dyn Adapter>` and consumes the list, but the sync has
479/// to ask the same adapter for its final session snapshot once the walk is over
480/// (see [`MjolnirAdapter::indexed_tags`]). Sharing the adapter is the only way
481/// both can hold it.
482struct 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
514/// The daemon's SessionWiki sync job.
515///
516/// Triggers coalesce: a request while a sync is running marks a rerun instead
517/// of queueing a second one, so a burst of closing sessions costs one extra
518/// pass. Syncs are single-flight because SessionWiki holds a write transaction
519/// per adapter batch, and two writers only produce a busy error.
520pub struct WikiIndexer {
521    inner: Arc<Indexer>,
522}
523
524#[derive(Default)]
525struct Indexer {
526    /// Held for the whole of one run: this is what makes syncs single-flight.
527    running: tokio::sync::Mutex<()>,
528    notify: tokio::sync::Notify,
529    /// A trigger arrived; the worker has not consumed it yet.
530    requested: AtomicBool,
531    /// At least one waiting trigger asked for a full sync.
532    full_requested: AtomicBool,
533    /// A sync pass is running now. A surface shows this as "topping up", so a
534    /// user knows more results may arrive.
535    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    /// Start the background sync worker. Without a Tokio runtime (some tests
547    /// build a runtime state without one) the indexer stays inert.
548    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    /// Ask for a sync. Returns immediately; the work happens in the background.
558    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    /// An indexer with no worker, so a test can see a request that nothing
567    /// takes and no test runs a sync against the real index.
568    #[cfg(test)]
569    pub(crate) fn inert() -> Self {
570        Self {
571            inner: Arc::new(Indexer::default()),
572        }
573    }
574
575    /// Whether a sync has been requested and not yet taken by the worker.
576    #[cfg(test)]
577    pub(crate) fn sync_requested(&self) -> bool {
578        self.inner.requested.load(Ordering::Acquire)
579    }
580
581    /// Run a sync and wait for it, joining a sync already in flight.
582    pub async fn sync_now(&self, full: bool) -> Result<()> {
583        self.inner.sync(full).await
584    }
585
586    /// The state of the index and whether a sync is running, for the surfaces
587    /// that say so while the first build is under way.
588    pub fn status(&self) -> WikiStatus {
589        WikiStatus {
590            state: index_state(),
591            topping_up: self.inner.in_flight.load(Ordering::Acquire)
592                || self.inner.requested.load(Ordering::Acquire),
593        }
594    }
595
596    /// When the last sync succeeded, for callers that trigger on staleness.
597    pub fn last_success(&self) -> Option<Instant> {
598        self.inner
599            .last_success
600            .lock()
601            .unwrap_or_else(std::sync::PoisonError::into_inner)
602            .map(|success| success.at)
603    }
604}
605
606impl Indexer {
607    async fn run(self: Arc<Self>) {
608        loop {
609            self.notify.notified().await;
610            while self.requested.swap(false, Ordering::AcqRel) {
611                let full = self.full_requested.swap(false, Ordering::AcqRel);
612                if let Err(error) = self.sync(full).await {
613                    self.report(&error);
614                    // A failure waits for the next trigger rather than
615                    // retrying straight away: a busy index stays busy for as
616                    // long as the other writer holds it, and a spin would only
617                    // add to the contention.
618                    break;
619                }
620            }
621        }
622    }
623
624    /// Log a failed sync at the level its cause deserves. A busy index is an
625    /// expected collision with another writer, not a fault: mark a rerun and
626    /// say so only in debug output.
627    fn report(&self, error: &anyhow::Error) {
628        if crate::database::is_busy_error(error) {
629            self.requested.store(true, Ordering::Release);
630            tracing::debug!(%error, "the SessionWiki index was busy; retrying on the next trigger");
631        } else {
632            tracing::warn!(%error, "could not sync sessions into SessionWiki");
633        }
634    }
635
636    async fn sync(&self, full: bool) -> Result<()> {
637        let _guard = self.running.lock().await;
638        let since = if full {
639            None
640        } else {
641            self.last_success
642                .lock()
643                .unwrap_or_else(std::sync::PoisonError::into_inner)
644                // A minute of overlap covers checkpoints written while the
645                // previous run was reading the directory.
646                .map(|success| success.epoch_seconds - 60)
647        };
648        let started = Instant::now();
649        let work = crate::upgrade::activity("SessionWiki sync")?;
650        self.in_flight.store(true, Ordering::Release);
651        let ran = tokio::task::spawn_blocking(move || {
652            let _work = work;
653            sync_blocking(since)
654        })
655        .await;
656        self.in_flight.store(false, Ordering::Release);
657        let ran = ran.context("run the SessionWiki sync")??;
658        if ran {
659            *self
660                .last_success
661                .lock()
662                .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(Success {
663                at: started,
664                epoch_seconds: Utc::now().timestamp(),
665            });
666        }
667        Ok(())
668    }
669}
670
671/// One synchronous sync pass. Returns false when this process must not touch
672/// the index, so a refused run never records a success it did not have.
673fn sync_blocking(since: Option<i64>) -> Result<bool> {
674    if !index_is_writable() {
675        return Ok(false);
676    }
677    let controller =
678        Controller::load().context("load controller state for the SessionWiki sync")?;
679    // Mjolnir's own sessions go first: a cold index walks every other tool's
680    // store for many minutes, and a just-closed session should not wait on it.
681    let mjolnir = Arc::new(MjolnirAdapter::reloading(&controller.state));
682    let mut adapters: Vec<Box<dyn sessionwiki::adapters::Adapter>> =
683        vec![Box::new(SharedMjolnirAdapter(Arc::clone(&mjolnir)))];
684    adapters.extend(native_adapters(&controller.config));
685    let mut connection = sessionwiki::index::open().context("open the SessionWiki index")?;
686    sessionwiki::index::sync_with(&mut connection, &adapters, since)
687        .context("sync the SessionWiki index")?;
688    write_session_tags(&mut connection, &mjolnir.indexed_tags())
689        .context("store Mjolnir's session metadata in the SessionWiki index")?;
690    provenance::backfill(&mut connection, &mjolnir).context("backfill Mjolnir file provenance")?;
691    if since.is_none() {
692        // A full pass has walked every store, so the index is complete enough
693        // for a search to be trusted. The marker is what a later daemon reads
694        // instead of walking the corpus again to find out.
695        record_first_build();
696    }
697    Ok(true)
698}
699
700/// Store each session's target, profile and harness in the index, in one
701/// transaction.
702///
703/// Every session is written on every sync rather than only the changed ones:
704/// the write is a delete and three inserts, which is nothing beside the
705/// transcript indexing in the same pass, and it is what makes a Move or a
706/// profile switch show up without tracking which records changed. It is also
707/// what gives sessions indexed before this existed their metadata, with no
708/// migration and no re-index.
709fn write_session_tags(
710    connection: &mut rusqlite::Connection,
711    session_tags: &BTreeMap<String, tags::MjTags>,
712) -> Result<()> {
713    if session_tags.is_empty() {
714        return Ok(());
715    }
716    let transaction = connection
717        .transaction()
718        .context("open a transaction for the session metadata")?;
719    for (session_id, session) in session_tags {
720        if session.is_empty() {
721            continue;
722        }
723        tags::write(&transaction, session_id, session)?;
724    }
725    transaction
726        .commit()
727        .context("commit the session metadata")?;
728    Ok(())
729}
730
731/// The non-Mjolnir adapters this install indexes.
732///
733/// Mjolnir's configured harness profiles decide which harness homes are
734/// indexed, not the stock `~/.codex` and `~/.claude` locations. A user who
735/// runs several profile homes expects every session Mjolnir can start to be
736/// searchable, and a home no profile names is not Mjolnir's to walk. So the
737/// stock Codex and Claude adapters are dropped and one adapter per enabled
738/// profile home takes their place; every other built-in adapter is kept as is.
739///
740/// Kimi Code, Grok Build and Muse have no SessionWiki adapter at all, so
741/// Mjolnir supplies one per enabled profile home of its own (see
742/// [`harness_adapters`]). Without them those sessions would never appear in
743/// the Resume dialog's search.
744///
745/// Each per-home adapter reports a reconcile scope covering only its own root,
746/// so a sync of one install never archives the rows of another.
747fn native_adapters(config: &mj_core::config::Config) -> Vec<Box<dyn Adapter>> {
748    // Two profiles may share one home, and two harnesses may share one home
749    // path without sharing sessions, so the kind is part of the identity.
750    let mut seen: BTreeSet<(HarnessKind, &Path)> = BTreeSet::new();
751    let mut adapters: Vec<Box<dyn Adapter>> = Vec::new();
752    for (_, profile) in config.enabled_profiles() {
753        // A second adapter for the same home would only walk it twice.
754        if !seen.insert((profile.kind, profile.home.as_path())) {
755            continue;
756        }
757        let adapter: Box<dyn Adapter> = match profile.kind {
758            HarnessKind::Codex => {
759                Box::new(sessionwiki::adapters::Codex::in_home(profile.home.clone()))
760            }
761            HarnessKind::Claude => Box::new(sessionwiki::adapters::ClaudeCode::in_home(
762                profile.home.clone(),
763            )),
764            kind => match HarnessAdapter::in_home(kind, profile.home.clone()) {
765                Some(adapter) => Box::new(adapter),
766                None => continue,
767            },
768        };
769        adapters.push(adapter);
770    }
771    adapters.extend(
772        sessionwiki::adapters::all()
773            .into_iter()
774            .filter(|adapter| !matches!(adapter.name(), "codex" | "claude-code")),
775    );
776    adapters
777}
778
779// ---------------------------------------------------------------------------
780// Which index, and whether it may be touched
781// ---------------------------------------------------------------------------
782
783/// Whether this process may open the index at all.
784///
785/// Indexing is always on, so a process that never resolved where its index
786/// belongs must not reach for one: it would walk the user's real session
787/// stores and write the user's real index. Only Mjolnir's own startup resolves
788/// it (see `mj_core::config::apply_instance_flag`), so this refuses every unit
789/// test that builds a daemon runtime directly and every other embedder, unless
790/// it names an index of its own with `SESSIONWIKI_DATA`.
791fn index_is_isolated() -> bool {
792    static SAID: AtomicBool = AtomicBool::new(false);
793    if mj_core::config::session_index_is_resolved()
794        || std::env::var_os(mj_core::config::SESSION_INDEX_ENV).is_some()
795    {
796        return true;
797    }
798    if !SAID.swap(true, Ordering::AcqRel) {
799        tracing::debug!(
800            "this process did not resolve a session index location; SessionWiki is not used"
801        );
802    }
803    false
804}
805
806/// Whether the index on disk was written by a SessionWiki at another schema
807/// version.
808///
809/// SessionWiki's own `open` drops and rebuilds its whole cache when the file's
810/// `user_version` differs from the version it was built with, which on a large
811/// corpus costs tens of minutes. Mjolnir will not do that to a user who also
812/// runs the `sessionwiki` command: it reads the version without SessionWiki and
813/// stands aside.
814fn index_version_mismatch() -> bool {
815    static SAID: AtomicBool = AtomicBool::new(false);
816    let Ok(path) = sessionwiki::index::db_path() else {
817        return false;
818    };
819    if !path.exists() {
820        return false;
821    }
822    let version = rusqlite::Connection::open_with_flags(
823        &path,
824        rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_URI,
825    )
826    .and_then(|connection| connection.pragma_query_value(None, "user_version", |row| row.get(0)));
827    let version: i64 = match version {
828        Ok(version) => version,
829        Err(error) => {
830            tracing::debug!(%error, "could not read the SessionWiki index schema version");
831            return false;
832        }
833    };
834    // Zero is an index SessionWiki has not finished creating; it is not a
835    // different version.
836    let mismatch = version != 0 && version != sessionwiki::index::SCHEMA_VERSION;
837    if mismatch && !SAID.swap(true, Ordering::AcqRel) {
838        tracing::warn!(
839            found = version,
840            expected = sessionwiki::index::SCHEMA_VERSION,
841            path = %path.display(),
842            "the SessionWiki index was written by another version;              Mjolnir will not open it, because opening it would rebuild it.              Install the matching sessionwiki command"
843        );
844    }
845    mismatch
846}
847
848fn index_is_writable() -> bool {
849    index_is_isolated() && !index_version_mismatch()
850}
851
852/// The file recording that one full sync has completed, holding the schema
853/// version it completed at.
854fn first_build_marker() -> PathBuf {
855    mj_core::config::data_dir().join("sessionwiki-built")
856}
857
858fn record_first_build() {
859    let path = first_build_marker();
860    let version = sessionwiki::index::SCHEMA_VERSION.to_string();
861    if std::fs::read_to_string(&path).is_ok_and(|held| held.trim() == version) {
862        return;
863    }
864    if let Err(error) = std::fs::write(&path, &version) {
865        tracing::warn!(%error, path = %path.display(), "could not record the first SessionWiki build");
866    }
867}
868
869/// Whether this index has completed a full build at this schema version.
870fn first_build_is_done() -> bool {
871    std::fs::read_to_string(first_build_marker())
872        .is_ok_and(|held| held.trim() == sessionwiki::index::SCHEMA_VERSION.to_string())
873        && sessionwiki::index::db_path().is_ok_and(|path| path.exists())
874}
875
876/// What a surface should say about this index right now.
877pub fn index_state() -> WikiIndexState {
878    if !index_is_isolated() {
879        return WikiIndexState::Indexing;
880    }
881    if index_version_mismatch() {
882        return WikiIndexState::VersionMismatch;
883    }
884    if first_build_is_done() {
885        WikiIndexState::Ready
886    } else {
887        WikiIndexState::Indexing
888    }
889}
890
891// ---------------------------------------------------------------------------
892// Indexing a session before it is destroyed
893// ---------------------------------------------------------------------------
894
895/// How long a destroy waits for a sync pass before it indexes the session on
896/// its own. A first build can run for many minutes, and a destroy must not
897/// wait for it.
898pub const DESTROY_SYNC_WAIT: Duration = Duration::from_secs(20);
899
900/// How often rows the index was too busy to take are offered again, and for
901/// how long. A first build holds the index while it parses one tool's
902/// sessions and frees it between tools.
903const DEFERRED_WRITE_RETRY: Duration = Duration::from_secs(5);
904const DEFERRED_WRITE_LIMIT: Duration = Duration::from_secs(60 * 60);
905
906/// How a session about to be destroyed, with its sub-agents, reached the
907/// index.
908#[derive(Debug, Clone, PartialEq, Eq)]
909pub enum IndexedBeforeDestroy {
910    /// The index already held every one of them as they are now, or none of
911    /// them had a conversation to index.
912    Current,
913    /// A sync pass that began after the request finished in time.
914    Synced,
915    /// The sync pass did not finish in time, so their rows were written on
916    /// their own.
917    WrittenDirectly,
918    /// The index was busy. Their rows were read before the destroy and are
919    /// written in the background once the index is free.
920    Deferred,
921    /// This process may not write the index, and why.
922    Unavailable(&'static str),
923    /// Indexing failed, and why. The destroy goes ahead.
924    Failed(String),
925}
926
927impl WikiIndexer {
928    /// Put a session and its sub-agents into the index while their records
929    /// and stored conversations still exist.
930    ///
931    /// Sessions enter the index only on a sync pass, and destroy deletes the
932    /// record and the conversation, so a session created and destroyed
933    /// between two passes was never findable (R2-11). This asks for an
934    /// incremental pass and waits `wait` for it. A pass that does not finish
935    /// in time, such as a first build, is left running, and the sessions are
936    /// indexed on their own from the same rows the pass would write.
937    pub async fn index_before_destroy(
938        &self,
939        session_id: &str,
940        wait: Duration,
941    ) -> IndexedBeforeDestroy {
942        if let Some(reason) = unwritable_reason() {
943            return IndexedBeforeDestroy::Unavailable(reason);
944        }
945        let root = session_id.to_owned();
946        let pending = match tokio::task::spawn_blocking(move || unindexed_session_tree(&root)).await
947        {
948            Ok(Ok(pending)) => pending,
949            Ok(Err(error)) => {
950                return IndexedBeforeDestroy::Failed(format!(
951                    "could not tell whether the index holds the session: {error:#}"
952                ));
953            }
954            Err(error) => {
955                return IndexedBeforeDestroy::Failed(format!(
956                    "checking the index for the session stopped: {error}"
957                ));
958            }
959        };
960        if pending.is_empty() {
961            return IndexedBeforeDestroy::Current;
962        }
963        let inner = Arc::clone(&self.inner);
964        index_before_destroy_with(
965            async move { inner.sync(false).await },
966            wait,
967            move || capture_sessions(&pending),
968            DEFERRED_WRITE_RETRY,
969        )
970        .await
971    }
972}
973
974/// Why this process may not write the index, if it may not.
975fn unwritable_reason() -> Option<&'static str> {
976    if !index_is_isolated() {
977        return Some("this process did not resolve a SessionWiki index of its own");
978    }
979    if index_version_mismatch() {
980        return Some("the SessionWiki index was written by another SessionWiki version");
981    }
982    None
983}
984
985/// The bounded wait and its fallback, with the sync pass and the reading of
986/// the sessions passed in so a test can stand in for either.
987async fn index_before_destroy_with<S, C>(
988    sync: S,
989    wait: Duration,
990    capture: C,
991    retry: Duration,
992) -> IndexedBeforeDestroy
993where
994    S: std::future::Future<Output = Result<()>> + Send + 'static,
995    C: FnOnce() -> Result<Vec<CapturedSession>> + Send + 'static,
996{
997    // Spawned rather than awaited here, so a wait that runs out drops only
998    // the handle and the pass still finishes. Dropping `Indexer::sync`
999    // part-way would release the single-flight lock while its blocking half
1000    // was still writing.
1001    match tokio::time::timeout(wait, tokio::spawn(sync)).await {
1002        Ok(Ok(Ok(()))) => return IndexedBeforeDestroy::Synced,
1003        Ok(Ok(Err(error))) => tracing::warn!(
1004            error = %format!("{error:#}"),
1005            "the SessionWiki sync before a destroy failed; indexing the session on its own"
1006        ),
1007        Ok(Err(error)) => tracing::warn!(
1008            %error,
1009            "the SessionWiki sync before a destroy stopped; indexing the session on its own"
1010        ),
1011        Err(_) => tracing::info!(
1012            wait_seconds = wait.as_secs_f64(),
1013            "the SessionWiki sync did not finish in time; indexing the session on its own"
1014        ),
1015    }
1016    let captured = match tokio::task::spawn_blocking(capture).await {
1017        Ok(Ok(captured)) => Arc::new(captured),
1018        Ok(Err(error)) => {
1019            return IndexedBeforeDestroy::Failed(format!(
1020                "could not read the session to index it: {error:#}"
1021            ));
1022        }
1023        Err(error) => {
1024            return IndexedBeforeDestroy::Failed(format!(
1025                "reading the session to index it stopped: {error}"
1026            ));
1027        }
1028    };
1029    if captured.is_empty() {
1030        return IndexedBeforeDestroy::Current;
1031    }
1032    let attempt = Arc::clone(&captured);
1033    match tokio::task::spawn_blocking(move || write_captured(&attempt)).await {
1034        Ok(Ok(())) => IndexedBeforeDestroy::WrittenDirectly,
1035        Ok(Err(error)) if crate::database::is_busy_error(&error) => {
1036            write_captured_later(captured, retry);
1037            IndexedBeforeDestroy::Deferred
1038        }
1039        Ok(Err(error)) => IndexedBeforeDestroy::Failed(format!(
1040            "could not write the session into the index: {error:#}"
1041        )),
1042        Err(error) => IndexedBeforeDestroy::Failed(format!(
1043            "writing the session into the index stopped: {error}"
1044        )),
1045    }
1046}
1047
1048/// A session and every sub-agent below it that has a conversation the index
1049/// does not hold as it is now. Empty when the session has no record.
1050fn unindexed_session_tree(root: &str) -> Result<Vec<String>> {
1051    let controller =
1052        Controller::load().context("load controller state to index a destroyed session")?;
1053    unindexed(
1054        &MjolnirAdapter::from_state(&controller.state),
1055        &session_tree(&controller.state, root),
1056    )
1057}
1058
1059/// Those of `session_ids` that have a conversation the index does not hold
1060/// as it is now: no row, an archived row, or a row with an older change
1061/// token than the adapter lists.
1062fn unindexed(adapter: &MjolnirAdapter, session_ids: &[String]) -> Result<Vec<String>> {
1063    let tokens: BTreeMap<String, i64> = adapter
1064        .store()
1065        .map(|store| store.keys.into_iter().collect())
1066        .unwrap_or_default();
1067    // No index yet holds nothing.
1068    let connection = open_readonly().ok();
1069    let mut pending = Vec::new();
1070    for session_id in session_ids {
1071        let key = adapter.key_for(session_id);
1072        // Only a session with a stored or checkpointed conversation is listed.
1073        let Some(&token) = tokens.get(&key) else {
1074            continue;
1075        };
1076        let current = match &connection {
1077            Some(connection) => indexed_token(connection, &key)? == Some(token),
1078            None => false,
1079        };
1080        if !current {
1081            pending.push(session_id.clone());
1082        }
1083    }
1084    Ok(pending)
1085}
1086
1087/// A session and every sub-agent below it, parents first.
1088fn session_tree(state: &State, root: &str) -> Vec<String> {
1089    if !state.sessions.contains_key(root) {
1090        return Vec::new();
1091    }
1092    let mut tree = vec![root.to_owned()];
1093    let mut seen = BTreeSet::from([root.to_owned()]);
1094    let mut next = 0;
1095    while let Some(parent) = tree.get(next).cloned() {
1096        next += 1;
1097        for child in state.subagents.values() {
1098            if child.parent_session_id == parent
1099                && state.sessions.contains_key(&child.child_session_id)
1100                && seen.insert(child.child_session_id.clone())
1101            {
1102                tree.push(child.child_session_id.clone());
1103            }
1104        }
1105    }
1106    tree
1107}
1108
1109/// The change token the index holds for a live row. SessionWiki stores a
1110/// shared-store token in the `mtime` column.
1111fn indexed_token(connection: &rusqlite::Connection, key: &str) -> Result<Option<i64>> {
1112    use rusqlite::OptionalExtension;
1113    connection
1114        .query_row(
1115            "SELECT mtime FROM files WHERE path = ?1 AND archived_at IS NULL",
1116            [key],
1117            |row| row.get(0),
1118        )
1119        .optional()
1120        .context("read a session's change token from the SessionWiki index")
1121}
1122
1123/// One session's index row and metadata, read while its record and
1124/// conversation still exist, so they can be written after both are gone.
1125struct CapturedSession {
1126    key: String,
1127    token: i64,
1128    session: Session,
1129    tags: tags::MjTags,
1130}
1131
1132fn capture_sessions(session_ids: &[String]) -> Result<Vec<CapturedSession>> {
1133    let controller =
1134        Controller::load().context("load controller state to index a destroyed session")?;
1135    capture_sessions_from(&MjolnirAdapter::from_state(&controller.state), session_ids)
1136}
1137
1138/// The rows a sync pass would write for these sessions, built by the same
1139/// adapter. A session with no conversation to index is left out.
1140fn capture_sessions_from(
1141    adapter: &MjolnirAdapter,
1142    session_ids: &[String],
1143) -> Result<Vec<CapturedSession>> {
1144    let tokens: BTreeMap<String, i64> = adapter
1145        .store()
1146        .map(|store| store.keys.into_iter().collect())
1147        .unwrap_or_default();
1148    let mut session_tags = adapter.indexed_tags();
1149    let mut captured = Vec::new();
1150    for session_id in session_ids {
1151        let key = adapter.key_for(session_id);
1152        let Some(&token) = tokens.get(&key) else {
1153            continue;
1154        };
1155        captured.push(CapturedSession {
1156            session: adapter.parse_key(&key)?,
1157            tags: session_tags.remove(session_id).unwrap_or_default(),
1158            key,
1159            token,
1160        });
1161    }
1162    Ok(captured)
1163}
1164
1165/// Write captured rows through SessionWiki's own indexing, one session at a
1166/// time, and their metadata beside them.
1167fn write_captured(captured: &Arc<Vec<CapturedSession>>) -> Result<()> {
1168    anyhow::ensure!(
1169        index_is_writable(),
1170        "this process may not write the SessionWiki index"
1171    );
1172    let mut connection = sessionwiki::index::open().context("open the SessionWiki index")?;
1173    for index in 0..captured.len() {
1174        let adapter: Box<dyn Adapter> = Box::new(CapturedAdapter {
1175            captured: Arc::clone(captured),
1176            index,
1177        });
1178        sessionwiki::index::sync_with(&mut connection, &[adapter], None)
1179            .context("index a session before it is destroyed")?;
1180    }
1181    let session_tags = captured
1182        .iter()
1183        .map(|captured| (captured.session.id.clone(), captured.tags.clone()))
1184        .collect();
1185    write_session_tags(&mut connection, &session_tags)
1186        .context("store Mjolnir's session metadata in the SessionWiki index")
1187}
1188
1189/// Offer rows the index was too busy to take until it takes them, or until
1190/// [`DEFERRED_WRITE_LIMIT`] passes. The rows live only in this task: a
1191/// daemon that stops before the index is free loses them, and the sessions
1192/// logged here are then not found by id.
1193fn write_captured_later(captured: Arc<Vec<CapturedSession>>, retry: Duration) {
1194    let sessions = captured
1195        .iter()
1196        .map(|captured| captured.session.id.clone())
1197        .collect::<Vec<_>>();
1198    tracing::info!(
1199        ?sessions,
1200        "the SessionWiki index is busy; indexing the destroyed sessions once it is free"
1201    );
1202    tokio::spawn(async move {
1203        let started = Instant::now();
1204        loop {
1205            tokio::time::sleep(retry).await;
1206            let attempt = Arc::clone(&captured);
1207            let error = match tokio::task::spawn_blocking(move || write_captured(&attempt)).await {
1208                Ok(Ok(())) => {
1209                    tracing::info!(?sessions, "indexed the destroyed sessions");
1210                    return;
1211                }
1212                Ok(Err(error)) => error,
1213                Err(error) => anyhow::Error::new(error),
1214            };
1215            if !crate::database::is_busy_error(&error) || started.elapsed() >= DEFERRED_WRITE_LIMIT
1216            {
1217                tracing::warn!(
1218                    ?sessions,
1219                    error = %format!("{error:#}"),
1220                    "gave up indexing destroyed sessions in SessionWiki"
1221                );
1222                return;
1223            }
1224        }
1225    });
1226}
1227
1228/// One captured session, offered to SessionWiki as a shared store that
1229/// lists only it.
1230struct CapturedAdapter {
1231    captured: Arc<Vec<CapturedSession>>,
1232    index: usize,
1233}
1234
1235impl CapturedAdapter {
1236    fn captured(&self) -> &CapturedSession {
1237        &self.captured[self.index]
1238    }
1239}
1240
1241impl Adapter for CapturedAdapter {
1242    fn name(&self) -> &'static str {
1243        TOOL
1244    }
1245
1246    fn root(&self) -> Option<PathBuf> {
1247        Path::new(&self.captured().key)
1248            .parent()
1249            .map(Path::to_path_buf)
1250    }
1251
1252    fn discover(&self) -> Discovered {
1253        Discovered {
1254            files: Vec::new(),
1255            had_error: false,
1256        }
1257    }
1258
1259    fn parse(&self, _path: &Path) -> Result<Session> {
1260        anyhow::bail!("Mjolnir sessions are parsed by key, not by file")
1261    }
1262
1263    fn store(&self) -> Option<Store> {
1264        let captured = self.captured();
1265        Some(Store {
1266            keys: vec![(captured.key.clone(), captured.token)],
1267            files: Vec::new(),
1268            had_error: false,
1269        })
1270    }
1271
1272    fn parse_key(&self, key: &str) -> Result<Session> {
1273        let captured = self.captured();
1274        anyhow::ensure!(key == captured.key, "no captured session for key {key:?}");
1275        Ok(copy_session(&captured.session))
1276    }
1277
1278    /// A prefix no key starts with, since keys hold no NUL. This store lists
1279    /// one session, not every session of this instance, so reconciliation
1280    /// must not archive the rows it does not list.
1281    fn reconcile_scope(&self) -> Option<String> {
1282        Some(format!("{}\0", self.captured().key))
1283    }
1284}
1285
1286/// A copy of an indexed session, for a write that may be retried.
1287/// SessionWiki's model does not implement `Clone`.
1288fn copy_session(session: &Session) -> Session {
1289    Session {
1290        id: session.id.clone(),
1291        tool: session.tool,
1292        path: session.path.clone(),
1293        project: session.project.clone(),
1294        started: session.started,
1295        ended: session.ended,
1296        title: session.title.clone(),
1297        subagent: session.subagent,
1298        messages: session
1299            .messages
1300            .iter()
1301            .map(|message| Message {
1302                role: message.role,
1303                text: message.text.clone(),
1304                ts: message.ts,
1305            })
1306            .collect(),
1307        touched: session.touched.clone(),
1308        edits: session
1309            .edits
1310            .iter()
1311            .map(|edit| sessionwiki::model::EditEvent {
1312                path: edit.path.clone(),
1313                kind: edit.kind,
1314                snippet: edit.snippet.clone(),
1315                ts: edit.ts,
1316            })
1317            .collect(),
1318    }
1319}
1320
1321// ---------------------------------------------------------------------------
1322// Queries and restore
1323// ---------------------------------------------------------------------------
1324
1325/// The largest page a caller may ask a wiki query for.
1326pub const MAX_WIKI_LIMIT: usize = 200;
1327/// The page size a caller that names none gets.
1328pub const DEFAULT_WIKI_LIMIT: usize = 50;
1329/// SessionWiki's full-text index needs three characters; shorter queries fall
1330/// back to a substring scan.
1331const MIN_FULLTEXT_QUERY: usize = 3;
1332/// How stale the index may be before a query triggers a background sync.
1333pub const SYNC_STALE_AFTER: std::time::Duration = std::time::Duration::from_secs(60);
1334
1335/// Whether a query should trigger a bounded background sync before it answers.
1336pub fn sync_is_stale(last_success: Option<Instant>) -> bool {
1337    last_success.is_none_or(|at| at.elapsed() >= SYNC_STALE_AFTER)
1338}
1339
1340/// One page of the index, newest first or best match first.
1341///
1342/// `live` is the set of session ids this daemon still holds, which is what
1343/// decides whether a Mjolnir row names a session the user can simply resume.
1344/// `include_subagents` decides, for every path below, whether sub-agent
1345/// sessions are answered at all; a resume list never wants them.
1346/// Runs SQLite work, so callers on the async runtime wrap it in
1347/// `spawn_blocking`.
1348pub fn query_rows(
1349    query: &str,
1350    limit: usize,
1351    live: &BTreeSet<String>,
1352    include_subagents: bool,
1353) -> Result<Vec<WikiRow>> {
1354    let limit = limit.clamp(1, MAX_WIKI_LIMIT);
1355    if !index_is_writable() {
1356        // Nothing to answer from: either this process has no index of its own
1357        // or the one on disk is at another version. The status beside the rows
1358        // says which.
1359        return Ok(Vec::new());
1360    }
1361    let connection = open_readonly()?;
1362    let query = query.trim();
1363    if query.is_empty() {
1364        let rows =
1365            sessionwiki::index::recent(&connection, limit, None, None, None, include_subagents)
1366                .context("list recent SessionWiki sessions")?;
1367        let mut rows: Vec<WikiRow> = rows
1368            .into_iter()
1369            .map(|row| wiki_row(row, None, live))
1370            .collect();
1371        fill_session_tags(&connection, &mut rows)?;
1372        return Ok(rows);
1373    }
1374    // The library search ranks sub-agents too. Ask for enough candidates to
1375    // fill the caller's result limit after those unresumable rows are removed.
1376    let search_limit = if include_subagents {
1377        limit
1378    } else {
1379        MAX_WIKI_LIMIT
1380    };
1381    let hits = if query.chars().count() < MIN_FULLTEXT_QUERY {
1382        sessionwiki::index::search_like(&connection, query, search_limit, None, None)
1383    } else {
1384        sessionwiki::index::search(&connection, query, search_limit, None, None)
1385    }
1386    .context("search the SessionWiki index")?;
1387    // SessionWiki's full-text search has no sub-agent filter of its own.
1388    let mut rows = Vec::with_capacity(hits.len().min(limit));
1389    for hit in hits {
1390        if rows.len() >= limit {
1391            break;
1392        }
1393        if !include_subagents && !is_main_session(&hit.row) {
1394            continue;
1395        }
1396        // A match only in tool text is not one the preview can show: it
1397        // never anchors on tool output. It is also how a sub-agent's words
1398        // reach its parent, as the Task prompt and result Claude Code records
1399        // in the parent's transcript. Keep such a hit only when the
1400        // conversation itself matches too. The agents' history search, which
1401        // asks for sub-agents, keeps tool matches.
1402        if !include_subagents
1403            && !matches!(hit.role.as_str(), "user" | "assistant")
1404            && !conversation_matches(&connection, &hit.row, query)?
1405        {
1406            continue;
1407        }
1408        rows.push(wiki_row(hit.row, Some(hit.snippet), live));
1409    }
1410    // SessionWiki searches message text alone, so a session known by a title
1411    // or a project that is never said out loud would be unfindable. Those
1412    // matches follow the full-text ones rather than displacing them.
1413    if rows.len() < limit {
1414        let found: BTreeSet<String> = rows.iter().map(|row| row.id.clone()).collect();
1415        for row in named_like(&connection, query, include_subagents)? {
1416            if rows.len() >= limit {
1417                break;
1418            }
1419            if found.contains(&row.session_id) {
1420                continue;
1421            }
1422            rows.push(wiki_row(row, None, live));
1423        }
1424    }
1425    fill_session_tags(&connection, &mut rows)?;
1426    Ok(rows)
1427}
1428
1429/// How many message rows one text search reads before it stops. A term common
1430/// enough to pass this can miss sessions whose only match ranks past it; the
1431/// person narrows the query, as with the resume dialog's own search.
1432const TEXT_SEARCH_MESSAGE_LIMIT: i64 = 20_000;
1433
1434/// The live Mjolnir sessions whose user or agent messages contain `query`,
1435/// ignoring case, and where: a session with a user match reports the user
1436/// match. Tool messages never match. Runs SQLite work, so callers on the async
1437/// runtime wrap it in `spawn_blocking`.
1438pub fn session_text_matches(query: &str, live: &BTreeSet<String>) -> Result<Vec<SessionTextMatch>> {
1439    let query = query.trim();
1440    if query.is_empty() || !index_is_writable() {
1441        return Ok(Vec::new());
1442    }
1443    let connection = open_readonly()?;
1444    text_matches_in(&connection, query, live)
1445}
1446
1447fn text_matches_in(
1448    connection: &rusqlite::Connection,
1449    query: &str,
1450    live: &BTreeSet<String>,
1451) -> Result<Vec<SessionTextMatch>> {
1452    let mut statement;
1453    let rows = if query.chars().count() < MIN_FULLTEXT_QUERY {
1454        // Too short for the trigram index; scan the newest messages instead.
1455        let pattern = format!(
1456            "%{}%",
1457            sessionwiki::util::nfc(query)
1458                .replace('\\', "\\\\")
1459                .replace('%', "\\%")
1460                .replace('_', "\\_")
1461        );
1462        statement = connection.prepare(
1463            "SELECT f.path, m.role
1464             FROM messages m JOIN files f ON f.session_id = m.session_id
1465             WHERE f.tool = ?1 AND m.role IN ('user', 'assistant')
1466               AND m.text LIKE ?2 ESCAPE '\\'
1467             ORDER BY m.id DESC LIMIT ?3",
1468        )?;
1469        statement
1470            .query_map(
1471                rusqlite::params![TOOL, pattern, TEXT_SEARCH_MESSAGE_LIMIT],
1472                |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
1473            )?
1474            .collect::<rusqlite::Result<Vec<_>>>()
1475    } else {
1476        let phrase = format!("\"{}\"", sessionwiki::util::nfc(query).replace('"', "\"\""));
1477        statement = connection.prepare(
1478            "SELECT f.path, m.role
1479             FROM (SELECT rowid AS mid FROM msgs WHERE msgs MATCH ?2 LIMIT ?3) x
1480             JOIN messages m ON m.id = x.mid
1481             JOIN files f ON f.session_id = m.session_id
1482             WHERE f.tool = ?1 AND m.role IN ('user', 'assistant')",
1483        )?;
1484        statement
1485            .query_map(
1486                rusqlite::params![TOOL, phrase, TEXT_SEARCH_MESSAGE_LIMIT],
1487                |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
1488            )?
1489            .collect::<rusqlite::Result<Vec<_>>>()
1490    }
1491    .context("search the indexed messages")?;
1492    let mut found = BTreeMap::<String, SessionTextMatchKind>::new();
1493    for (path, role) in rows {
1494        // The adapter writes the session id as the last path segment.
1495        let session_id = path.rsplit('/').next().unwrap_or_default();
1496        if !live.contains(session_id) {
1497            continue;
1498        }
1499        let kind = if role == "user" {
1500            SessionTextMatchKind::User
1501        } else {
1502            SessionTextMatchKind::Agent
1503        };
1504        let entry = found.entry(session_id.to_owned()).or_insert(kind);
1505        *entry = (*entry).min(kind);
1506    }
1507    Ok(found
1508        .into_iter()
1509        .map(|(session_id, kind)| SessionTextMatch { session_id, kind })
1510        .collect())
1511}
1512
1513/// Fill in the target, profile and harness of every Mjolnir row on this page
1514/// from the index's own tags, in one query.
1515///
1516/// Only Mjolnir writes those tags, so a row from another tool keeps `None` and
1517/// is not even asked about.
1518fn fill_session_tags(connection: &rusqlite::Connection, rows: &mut [WikiRow]) -> Result<()> {
1519    let ids: Vec<&str> = rows
1520        .iter()
1521        .filter(|row| row.tool == TOOL)
1522        .map(|row| row.id.as_str())
1523        .collect();
1524    let found = tags::read(connection, &ids).context("read the indexed session metadata")?;
1525    for row in rows.iter_mut().filter(|row| row.tool == TOOL) {
1526        let Some(session) = found.get(&row.id) else {
1527            continue;
1528        };
1529        row.target = session.target.clone();
1530        row.profile = session.profile.clone();
1531        row.harness = session.harness.clone();
1532    }
1533    Ok(())
1534}
1535
1536/// How far back a title or project match looks. Those columns have no index of
1537/// their own, so this is a scan of the most recent sessions rather than of the
1538/// whole corpus. Fetch only metadata here: the regular recent-session query
1539/// also loads a preview, summary and tags for every row, none of which title
1540/// matching needs.
1541const NAME_SCAN_LIMIT: usize = 2_000;
1542
1543/// Indexed sessions whose title or project contains the query, ignoring case.
1544fn named_like(
1545    connection: &rusqlite::Connection,
1546    query: &str,
1547    include_subagents: bool,
1548) -> Result<Vec<sessionwiki::index::SessionRow>> {
1549    let needle = query.to_lowercase();
1550    let sql = format!(
1551        "SELECT session_id, tool, path, project, title, started, msg_count, kind,
1552                archived_at IS NOT NULL
1553         FROM files WHERE {} ORDER BY started DESC LIMIT {NAME_SCAN_LIMIT}",
1554        if include_subagents {
1555            "1=1"
1556        } else {
1557            "kind = 'main'"
1558        }
1559    );
1560    let mut statement = connection
1561        .prepare(&sql)
1562        .context("prepare recent SessionWiki metadata scan")?;
1563    let rows = statement
1564        .query_map([], |row| {
1565            Ok(sessionwiki::index::SessionRow {
1566                session_id: row.get(0)?,
1567                tool: row.get(1)?,
1568                path: row.get(2)?,
1569                project: row.get(3)?,
1570                title: row.get(4)?,
1571                started: row.get(5)?,
1572                msg_count: row.get(6)?,
1573                kind: row.get(7)?,
1574                preview: None,
1575                summary: None,
1576                tags: None,
1577                archived: row.get(8)?,
1578                account: None,
1579            })
1580        })
1581        .context("list recent SessionWiki metadata")?
1582        .collect::<rusqlite::Result<Vec<_>>>()
1583        .context("read recent SessionWiki metadata")?;
1584    Ok(rows
1585        .into_iter()
1586        .filter(|row| {
1587            row.title.to_lowercase().contains(&needle)
1588                || row.project.to_lowercase().contains(&needle)
1589        })
1590        .collect())
1591}
1592
1593/// Whether a user or assistant message of an indexed session matches the
1594/// query, by the same rule the preview's passages use.
1595fn conversation_matches(
1596    connection: &rusqlite::Connection,
1597    row: &sessionwiki::index::SessionRow,
1598    query: &str,
1599) -> Result<bool> {
1600    let session = sessionwiki::index::session_from_index(connection, row)
1601        .context("read an indexed session")?;
1602    Ok(!hit_transcript(&session, query, 0, 1).blocks.is_empty())
1603}
1604
1605/// Whether an indexed session is one a person started rather than a
1606/// sub-agent. SessionWiki's own `recent` filter tests the same `kind`.
1607fn is_main_session(row: &sessionwiki::index::SessionRow) -> bool {
1608    row.kind == "main"
1609}
1610
1611/// The briefing for one indexed session, or `None` when the id names none.
1612pub fn brief(id: &str, max_chars: usize) -> Result<Option<String>> {
1613    if !index_is_writable() {
1614        return Ok(None);
1615    }
1616    let connection = open_readonly()?;
1617    let Some(row) = row_by_id(&connection, id)? else {
1618        return Ok(None);
1619    };
1620    let session = sessionwiki::index::session_from_index(&connection, &row)
1621        .context("read an indexed session")?;
1622    Ok(Some(sessionwiki::commands::brief_markdown(
1623        &session, max_chars, true,
1624    )))
1625}
1626
1627/// The passages of one indexed session that match `query`, or `None` when the
1628/// id names no indexed session.
1629///
1630/// Every matching message is returned with `context_messages` neighbours on
1631/// each side; overlapping groups are merged and each group's first block says
1632/// how many messages were skipped before it. Each block's text is capped at
1633/// `per_message_chars` characters, keeping the window around its first match.
1634pub fn transcript_hits(
1635    id: &str,
1636    query: &str,
1637    context_messages: usize,
1638    per_message_chars: usize,
1639) -> Result<Option<WikiHitTranscript>> {
1640    if !index_is_writable() {
1641        return Ok(None);
1642    }
1643    let connection = open_readonly()?;
1644    let Some(row) = row_by_id(&connection, id)? else {
1645        return Ok(None);
1646    };
1647    let session = sessionwiki::index::session_from_index(&connection, &row)
1648        .context("read an indexed session")?;
1649    Ok(Some(hit_transcript(
1650        &session,
1651        query,
1652        context_messages,
1653        per_message_chars,
1654    )))
1655}
1656
1657/// The matching passages of one loaded session, converted from SessionWiki's
1658/// own grep. Pure, so the conversion can be tested without an index on disk.
1659///
1660/// Matching, redaction and the excerpt window are `sessionwiki::grep`'s, so the
1661/// `sessionwiki grep` CLI and this preview report the same hits. Tool output
1662/// never anchors a passage: it is machine chatter the reader did not write,
1663/// a hit buried in it would open the preview on a wall of command output, and
1664/// the preview collapses tool runs anyway. Tool messages still appear as
1665/// context around a real match.
1666fn hit_transcript(
1667    session: &Session,
1668    query: &str,
1669    context_messages: usize,
1670    per_message_chars: usize,
1671) -> WikiHitTranscript {
1672    let found = sessionwiki::grep::grep_session(
1673        session,
1674        query,
1675        &sessionwiki::grep::GrepOpts {
1676            context_messages,
1677            chars: per_message_chars,
1678            max_matches: None,
1679            anchor_roles: vec![Role::User, Role::Assistant],
1680        },
1681    );
1682    WikiHitTranscript {
1683        blocks: found
1684            .hits
1685            .into_iter()
1686            .map(|hit| WikiHitBlock {
1687                role: role_name(hit.role).to_owned(),
1688                text: hit.text,
1689                hits: hit.matches,
1690                omitted_before: hit.omitted_before,
1691                truncated: hit.truncated,
1692            })
1693            .collect(),
1694        omitted_after: found.omitted_after,
1695    }
1696}
1697
1698fn role_name(role: Role) -> &'static str {
1699    match role {
1700        Role::User => "user",
1701        Role::Assistant => "assistant",
1702        Role::Tool => "tool",
1703    }
1704}
1705
1706/// What a restore needs from the index: the transcript as a snapshot the
1707/// compaction pipeline accepts, plus the title and project of the session it
1708/// came from.
1709pub struct ArchivedSession {
1710    pub title: String,
1711    /// The project directory the session ran in, when the row names one that
1712    /// still exists.
1713    pub project_directory: Option<PathBuf>,
1714    pub snapshot: mj_core::archive::CanonicalSessionSnapshot,
1715}
1716
1717/// Load one indexed session for restore, or `None` when the id names none.
1718pub fn archived_session(id: &str) -> Result<Option<ArchivedSession>> {
1719    if !index_is_writable() {
1720        return Ok(None);
1721    }
1722    let connection = open_readonly()?;
1723    let Some(row) = row_by_id(&connection, id)? else {
1724        return Ok(None);
1725    };
1726    let session = sessionwiki::index::session_from_index(&connection, &row)
1727        .context("read an indexed session")?;
1728    let snapshot = snapshot_of(&session)?;
1729    Ok(Some(ArchivedSession {
1730        title: session.title.clone(),
1731        project_directory: project_directory_of(&session.project),
1732        snapshot,
1733    }))
1734}
1735
1736// ---------------------------------------------------------------------------
1737// The archive job
1738// ---------------------------------------------------------------------------
1739
1740/// The stopped sessions that `archive_after_days = older_than_days` has caught,
1741/// children before their parents.
1742///
1743/// A session qualifies when its record is `Stopped`, its last update is at
1744/// least that many days old, and every sub-agent child it still has is being
1745/// archived in the same pass. The child rule is what keeps the pass from
1746/// destroying a session it did not choose: archiving a parent tears its
1747/// children down with it, so a child that is still running, or stopped but not
1748/// yet old enough, holds its parent back until the next pass.
1749///
1750/// Pure over controller state, so the rule can be tested without a daemon.
1751pub fn sessions_ready_to_archive(
1752    sessions: &mj_core::snapshot_map::SnapshotMap<String, SessionRecord>,
1753    subagents: &mj_core::snapshot_map::SnapshotMap<String, mj_core::subagent::SubagentRecord>,
1754    now: DateTime<Utc>,
1755    older_than_days: u32,
1756) -> Vec<String> {
1757    let cutoff = now - chrono::Duration::days(i64::from(older_than_days));
1758    let aged = |session_id: &String| {
1759        sessions.get(session_id).is_some_and(|record| {
1760            record.state == mj_core::state::SessionState::Stopped
1761                && parse_time(&record.updated_at).is_some_and(|updated| updated <= cutoff)
1762                && (record.managed_worktree.as_ref().is_some_and(|checkout| {
1763                    checkout.kind == mj_core::state::ManagedCheckoutKind::Worktree
1764                }) || (record.managed_worktree.is_none() && record.project_directory.is_some())
1765                    || record
1766                        .checkpoint
1767                        .as_ref()
1768                        .zip(record.publication.as_ref())
1769                        .is_some_and(|(checkpoint, publication)| {
1770                            publication.checkpoint_sha256 == checkpoint.sha256
1771                                && publication.state == mj_core::state::PublicationState::Published
1772                                && !publication.dirty
1773                                && !publication.stashed
1774                        }))
1775        })
1776    };
1777    let selected: BTreeSet<String> = sessions
1778        .keys()
1779        .filter(|session_id| aged(session_id))
1780        .filter(|session_id| {
1781            subagents
1782                .values()
1783                .filter(|child| &&child.parent_session_id == session_id)
1784                // A child whose record is already gone holds nothing open.
1785                .filter(|child| sessions.contains_key(&child.child_session_id))
1786                .all(|child| aged(&child.child_session_id))
1787        })
1788        .cloned()
1789        .collect();
1790    let mut ordered: Vec<String> = selected.iter().cloned().collect();
1791    ordered.sort_by_key(|session_id| std::cmp::Reverse(ancestor_depth(session_id, subagents)));
1792    ordered
1793}
1794
1795/// How many sub-agent parents a session has above it. Deeper sessions are
1796/// archived first so a parent never tears down a child the pass still has to
1797/// visit.
1798fn ancestor_depth(
1799    session_id: &str,
1800    subagents: &mj_core::snapshot_map::SnapshotMap<String, mj_core::subagent::SubagentRecord>,
1801) -> usize {
1802    let mut depth = 0;
1803    let mut current = session_id;
1804    // Bounded by the map: a cycle cannot outlive one pass over every entry.
1805    while let Some(parent) = subagents
1806        .get(current)
1807        .map(|child| child.parent_session_id.as_str())
1808    {
1809        depth += 1;
1810        if depth > subagents.len() {
1811            break;
1812        }
1813        current = parent;
1814    }
1815    depth
1816}
1817
1818/// How much disk Mjolnir's own copies of sessions use, and how much an
1819/// `archive_after_days` value would free. "Mjolnir's own copy" is the
1820/// checkpoint archive plus the session's image attachments; the conversation
1821/// itself lives in the SessionWiki index and is not counted, because archiving
1822/// keeps it. The type lives in `mj-core` so the terminal UI can name it too.
1823pub use mj_core::state::ArchiveSpacePreview;
1824
1825/// The space every session uses now and, when `older_than_days` is set, the
1826/// space archiving after that many days would reclaim.
1827///
1828/// The reclaim figure uses the archive job's own selection rule but not its
1829/// "is it indexed yet" gate: that gate depends on how far the hourly index
1830/// sync has got, so applying it would make the estimate swing between zero and
1831/// the true value while the first index builds. This answers what the policy
1832/// would reclaim, not what the next tick happens to reclaim.
1833///
1834/// Walks the filesystem, so callers on the async runtime must run it in a
1835/// blocking task.
1836pub fn archive_space_preview(older_than_days: Option<u32>) -> Result<ArchiveSpacePreview> {
1837    let controller =
1838        Controller::load().context("load the session records to size their storage")?;
1839    Ok(archive_space_over(
1840        &mj_core::config::sessions_dir(),
1841        &controller.state.sessions,
1842        &controller.state.subagents,
1843        Utc::now(),
1844        older_than_days,
1845    ))
1846}
1847
1848/// The sizing itself, over given records and a given sessions directory, so it
1849/// can be tested without the live data directory.
1850fn archive_space_over(
1851    sessions_root: &Path,
1852    sessions: &mj_core::snapshot_map::SnapshotMap<String, SessionRecord>,
1853    subagents: &mj_core::snapshot_map::SnapshotMap<String, mj_core::subagent::SubagentRecord>,
1854    now: DateTime<Utc>,
1855    older_than_days: Option<u32>,
1856) -> ArchiveSpacePreview {
1857    let mut preview = ArchiveSpacePreview {
1858        sessions: sessions.len(),
1859        bytes: sessions
1860            .iter()
1861            .map(|(session_id, record)| session_bytes(sessions_root, session_id, record))
1862            .sum(),
1863        reclaimable_sessions: 0,
1864        reclaimable_bytes: 0,
1865    };
1866    if let Some(days) = older_than_days {
1867        let aged = sessions_ready_to_archive(sessions, subagents, now, days);
1868        preview.reclaimable_sessions = aged.len();
1869        preview.reclaimable_bytes = aged
1870            .iter()
1871            .filter_map(|session_id| {
1872                sessions
1873                    .get(session_id)
1874                    .map(|record| session_bytes(sessions_root, session_id, record))
1875            })
1876            .sum();
1877    }
1878    preview
1879}
1880
1881/// What archiving one session would free: its checkpoint archive and its
1882/// attachments. Anything already missing counts as zero.
1883fn session_bytes(sessions_root: &Path, session_id: &str, record: &SessionRecord) -> u64 {
1884    let checkpoint = record
1885        .checkpoint
1886        .as_ref()
1887        .and_then(|checkpoint| std::fs::metadata(&checkpoint.archive_path).ok())
1888        .filter(|metadata| metadata.is_file())
1889        .map(|metadata| metadata.len())
1890        .unwrap_or(0);
1891    let attachments = sessions_root
1892        .join(session_id)
1893        .join(mj_core::attachment::ATTACHMENT_DIR);
1894    let attachments = crate::import::claude::directory_size(&attachments).unwrap_or(0);
1895    checkpoint.saturating_add(attachments)
1896}
1897
1898/// Which of `session_ids` the index holds under this instance's own key, with
1899/// at least one message and not already archived.
1900///
1901/// This is the gate the archive job will not cross: Mjolnir only deletes its
1902/// own copy of a conversation SessionWiki has actually stored. Runs SQLite
1903/// work, so callers on the async runtime wrap it in `spawn_blocking`.
1904pub fn indexed_with_messages(session_ids: &[String]) -> Result<BTreeSet<String>> {
1905    if !index_is_writable() {
1906        // An index this daemon will not open holds nothing it may act on, and
1907        // the archive job deletes data, so it must find nothing here.
1908        return Ok(BTreeSet::new());
1909    }
1910    let connection = open_readonly()?;
1911    let sessions_dir = mj_core::config::sessions_dir();
1912    let mut indexed = BTreeSet::new();
1913    for session_id in session_ids {
1914        let key = format!("{}/{session_id}", sessions_dir.display());
1915        let rows = sessionwiki::index::resolve(&connection, session_id)
1916            .context("look up a stopped session in the SessionWiki index")?;
1917        if rows
1918            .iter()
1919            .any(|row| row.tool == TOOL && row.path == key && row.msg_count > 0 && !row.archived)
1920        {
1921            indexed.insert(session_id.clone());
1922        }
1923    }
1924    Ok(indexed)
1925}
1926
1927fn open_readonly() -> Result<rusqlite::Connection> {
1928    sessionwiki::index::open_readonly().context("open the SessionWiki index")
1929}
1930
1931/// The one row an id names exactly. `resolve` matches prefixes, which is right
1932/// for a person typing and wrong for a client passing an id back.
1933fn row_by_id(
1934    connection: &rusqlite::Connection,
1935    id: &str,
1936) -> Result<Option<sessionwiki::index::SessionRow>> {
1937    Ok(sessionwiki::index::resolve(connection, id)
1938        .context("look up an indexed session")?
1939        .into_iter()
1940        .find(|row| row.session_id == id))
1941}
1942
1943fn wiki_row(
1944    row: sessionwiki::index::SessionRow,
1945    snippet: Option<String>,
1946    live: &BTreeSet<String>,
1947) -> WikiRow {
1948    // Only this daemon's own sessions can be live here, and only under the key
1949    // shape the adapter writes: the checkpoint directory and the session id.
1950    let hel_session_id = (row.tool == TOOL)
1951        .then(|| row.path.rsplit('/').next().unwrap_or_default().to_owned())
1952        .filter(|session_id| live.contains(session_id));
1953    let native_id = sessionwiki::index::native_id_of(&row.path);
1954    WikiRow {
1955        id: row.session_id,
1956        tool: row.tool,
1957        project: row.project,
1958        title: row.title,
1959        started: row.started,
1960        msgs: row.msg_count,
1961        preview: row.preview,
1962        archived: row.archived,
1963        native_id,
1964        snippet,
1965        hel_session_id,
1966        // Filled in by `fill_session_tags` from the index's own tags; the row
1967        // itself does not carry them.
1968        target: None,
1969        profile: None,
1970        harness: None,
1971    }
1972}
1973
1974/// The project a restored session should open.
1975///
1976/// A Mjolnir session runs in a managed worktree under the repository it was
1977/// started from, and that worktree is gone once the session is archived. The
1978/// repository above it is what the user still has, so a worktree path is
1979/// reduced to it. Any other path is used as it stands, and a path that no
1980/// longer exists is left for the caller to replace.
1981fn project_directory_of(project: &str) -> Option<PathBuf> {
1982    if project.trim().is_empty() {
1983        return None;
1984    }
1985    let path = PathBuf::from(project);
1986    let repository = path
1987        .ancestors()
1988        .find(|ancestor| ancestor.file_name().is_some_and(|name| name == ".mj"))
1989        .and_then(std::path::Path::parent)
1990        .map(std::path::Path::to_path_buf)
1991        .unwrap_or(path);
1992    repository.is_dir().then_some(repository)
1993}
1994
1995/// Rebuild an indexed transcript as a canonical snapshot.
1996///
1997/// The snapshot is only ever read by the compaction pipeline, which wants
1998/// turns: a user message opens a turn and assistant and tool items attach to
1999/// it. Messages before the first user message therefore have nowhere to go and
2000/// are dropped, and a session with no user message at all cannot be restored.
2001fn snapshot_of(
2002    session: &sessionwiki::model::Session,
2003) -> Result<mj_core::archive::CanonicalSessionSnapshot> {
2004    use mj_core::archive::{
2005        CanonicalExecutionState, CanonicalSessionSnapshot, CanonicalSessionState,
2006        CanonicalTranscriptBody, CanonicalTranscriptItem,
2007    };
2008
2009    let started_ms = session
2010        .started
2011        .map(|time| time.timestamp_millis())
2012        .unwrap_or_default();
2013    let mut transcript: Vec<CanonicalTranscriptItem> = Vec::new();
2014    for message in &session.messages {
2015        let text = message.text.trim();
2016        if text.is_empty() {
2017            continue;
2018        }
2019        // Compaction attaches assistant and tool items to the open turn, so an
2020        // item before the first user message would be dropped anyway.
2021        if transcript.is_empty() && message.role != Role::User {
2022            continue;
2023        }
2024        let position = transcript.len() as u64 + 1;
2025        let body = match message.role {
2026            Role::User => CanonicalTranscriptBody::User {
2027                content: vec![serde_json::json!({"type": "text", "text": text})],
2028            },
2029            Role::Assistant => CanonicalTranscriptBody::Agent {
2030                chunks: vec![serde_json::json!({
2031                    "content": {"type": "text", "text": text}
2032                })],
2033                streaming: false,
2034            },
2035            // New indexes retain the shared projection; legacy rows contain only a title.
2036            Role::Tool => {
2037                let (call, terminal_outputs) = mj_transcript::summary::indexed_tool_call(
2038                    text,
2039                    &format!("wiki-tool-{position}"),
2040                );
2041                CanonicalTranscriptBody::Tool {
2042                    call,
2043                    terminal_outputs,
2044                    terminal_refs: Vec::new(),
2045                    presentation: None,
2046                }
2047            }
2048        };
2049        let created_at_ms = message
2050            .ts
2051            .map(|time| time.timestamp_millis())
2052            .unwrap_or(started_ms);
2053        transcript.push(CanonicalTranscriptItem {
2054            stable_id: format!("wiki-{position}"),
2055            position,
2056            // The validator wants an ordinal on agent messages and on nothing
2057            // else; one event per item makes the item's own position right.
2058            latest_content_event_ordinal: matches!(body, CanonicalTranscriptBody::Agent { .. })
2059                .then_some(position),
2060            created_at_ms,
2061            last_changed_at_ms: created_at_ms,
2062            body,
2063        });
2064    }
2065    anyhow::ensure!(
2066        !transcript.is_empty(),
2067        "the archived session has no prompt to restore from"
2068    );
2069
2070    let event_frontier = transcript.len() as u64;
2071    let last_activity_at_ms = transcript.last().map(|item| item.last_changed_at_ms);
2072    Ok(CanonicalSessionSnapshot {
2073        command_ledger: None,
2074        assessment_state: None,
2075        event_frontier,
2076        // Not a relay frontier, so there is no recorded digest to carry. It has
2077        // to be a well-formed non-genesis digest, and deriving it from the
2078        // session makes two restores of one session agree.
2079        event_frontier_digest: {
2080            use sha2::Digest;
2081            mj_core::hex::lower_hex(sha2::Sha256::digest(
2082                format!("sessionwiki:{}", session.id).as_bytes(),
2083            ))
2084        },
2085        session: CanonicalSessionState {
2086            execution: CanonicalExecutionState::Idle,
2087            last_activity_at_ms,
2088            session_title: Some(session.title.clone()).filter(|title| !title.trim().is_empty()),
2089            configuration: Default::default(),
2090        },
2091        transcript,
2092        queued_prompts: Vec::new(),
2093    })
2094}
2095
2096// ---------------------------------------------------------------------------
2097// Continuing an indexed session
2098// ---------------------------------------------------------------------------
2099
2100/// What continuing one indexed session means.
2101///
2102/// An agent that found a session with SessionWiki should not have to know
2103/// whose session it was, so the branch lives here and `mj resume --wiki` takes
2104/// it on the agent's behalf.
2105#[derive(Debug, Clone, PartialEq, Eq)]
2106pub enum WikiContinuation {
2107    /// A Mjolnir session this daemon still has a record of: resume it.
2108    Resume { session_id: String },
2109    /// A Mjolnir session whose record the archive job destroyed: start a new
2110    /// session seeded with a compacted hand-off.
2111    Restore { wiki_id: String },
2112    /// Another tool's session: import it, then resume what the import made.
2113    Import {
2114        harness: HarnessKind,
2115        native_session_id: String,
2116    },
2117}
2118
2119/// How to continue the indexed session a row describes.
2120///
2121/// Pure over the row so the branch can be tested without an index:
2122/// `path` is the row's stored path and `has_record` says whether controller
2123/// state still holds a session with this id.
2124pub fn wiki_continuation(
2125    wiki_id: &str,
2126    tool: &str,
2127    path: &Path,
2128    has_record: bool,
2129) -> Result<WikiContinuation> {
2130    if tool == TOOL {
2131        // `MjolnirAdapter::parse_key` names the session by its own Mjolnir id,
2132        // so a Mjolnir row's SessionWiki id is the session id.
2133        return Ok(match has_record {
2134            true => WikiContinuation::Resume {
2135                session_id: wiki_id.to_owned(),
2136            },
2137            false => WikiContinuation::Restore {
2138                wiki_id: wiki_id.to_owned(),
2139            },
2140        });
2141    }
2142    let harness = harness_adapters::harness_for_tool(tool)
2143        .with_context(|| format!("Mjolnir cannot continue a {tool} session"))?;
2144    let native_session_id = crate::import::native_session_id_from_path(harness, path)
2145        .with_context(|| {
2146            format!(
2147                "no {tool} session id in the indexed path {}",
2148                path.display()
2149            )
2150        })?;
2151    Ok(WikiContinuation::Import {
2152        harness,
2153        native_session_id,
2154    })
2155}
2156
2157/// What one indexed session is, as far as continuing it is concerned.
2158///
2159/// Read through [`wiki_session`]; the daemon serves it for `mj resume --wiki`
2160/// and for `mj sessions --session` when the id names no Mjolnir session.
2161pub fn wiki_session(
2162    wiki_id: &str,
2163    known_sessions: &BTreeSet<String>,
2164) -> Result<Option<WikiSessionInfo>> {
2165    if !index_is_writable() {
2166        return Ok(None);
2167    }
2168    let connection = open_readonly()?;
2169    let Some(row) = row_by_id(&connection, wiki_id)? else {
2170        return Ok(None);
2171    };
2172    let is_mjolnir = row.tool == TOOL;
2173    let mjolnir_session_id = is_mjolnir.then(|| row.session_id.clone());
2174    let has_record = mjolnir_session_id
2175        .as_deref()
2176        .is_some_and(|session_id| known_sessions.contains(session_id));
2177    let status = match (is_mjolnir, has_record) {
2178        (false, _) => WikiSessionStatus::Native,
2179        (true, true) => WikiSessionStatus::Mine,
2180        (true, false) => WikiSessionStatus::Archived,
2181    };
2182    let tags = match is_mjolnir {
2183        true => tags::read(&connection, &[row.session_id.as_str()])
2184            .context("read the indexed session metadata")?
2185            .remove(&row.session_id)
2186            .unwrap_or_default(),
2187        false => tags::MjTags::default(),
2188    };
2189    // Only an archived row is continued by restoring its transcript, and a
2190    // restore needs a prompt to open the first turn.
2191    let nothing_to_restore = status == WikiSessionStatus::Archived
2192        && !has_prompt(
2193            &sessionwiki::index::session_from_index(&connection, &row)
2194                .context("read an indexed session")?,
2195        );
2196    let harness = tags
2197        .harness
2198        .as_deref()
2199        .and_then(|id| id.parse::<HarnessKind>().ok())
2200        .or_else(|| {
2201            (!is_mjolnir)
2202                .then(|| harness_adapters::harness_for_tool(&row.tool))
2203                .flatten()
2204        });
2205    Ok(Some(WikiSessionInfo {
2206        wiki_id: row.session_id,
2207        tool: row.tool,
2208        path: PathBuf::from(row.path),
2209        status,
2210        mjolnir_session_id,
2211        profile_id: tags.profile,
2212        target_template_id: tags.target,
2213        harness,
2214        title: row.title,
2215        project: row.project,
2216        nothing_to_restore,
2217    }))
2218}
2219
2220/// Whether an indexed transcript holds a prompt, which is what
2221/// [`snapshot_of`] needs to open a turn.
2222fn has_prompt(session: &sessionwiki::model::Session) -> bool {
2223    session
2224        .messages
2225        .iter()
2226        .any(|message| message.role == Role::User && !message.text.trim().is_empty())
2227}
2228
2229#[cfg(test)]
2230mod tests {
2231    use std::collections::BTreeMap;
2232    use std::path::Path;
2233
2234    /// The dispatch `mj resume --wiki` takes, over rows built by hand: the
2235    /// branch has to be right without an index behind it.
2236    mod continuation {
2237        use super::super::{WikiContinuation, wiki_continuation};
2238        use mj_core::config::HarnessKind;
2239        use std::path::Path;
2240
2241        #[test]
2242        fn a_mjolnir_row_with_a_record_is_resumed_and_one_without_is_restored() {
2243            let path = Path::new("/home/user/.local/share/mj/sessions/session-7");
2244            assert_eq!(
2245                wiki_continuation("session-7", "mjolnir", path, true).unwrap(),
2246                WikiContinuation::Resume {
2247                    session_id: "session-7".to_owned(),
2248                }
2249            );
2250            assert_eq!(
2251                wiki_continuation("session-7", "mjolnir", path, false).unwrap(),
2252                WikiContinuation::Restore {
2253                    wiki_id: "session-7".to_owned(),
2254                }
2255            );
2256        }
2257
2258        #[test]
2259        fn a_claude_code_row_is_imported_with_the_uuid_from_its_path() {
2260            let path = Path::new(
2261                "/home/user/.claude/projects/-home-user-app/7f3a1c20-0b11-4a55-9e0d-2c8a5d6f1b44.jsonl",
2262            );
2263            assert_eq!(
2264                wiki_continuation("abc123", "claude-code", path, false).unwrap(),
2265                WikiContinuation::Import {
2266                    harness: HarnessKind::Claude,
2267                    native_session_id: "7f3a1c20-0b11-4a55-9e0d-2c8a5d6f1b44".to_owned(),
2268                }
2269            );
2270        }
2271
2272        /// A Codex rollout's file name is a timestamp and the thread UUID, so
2273        /// the stem alone is not the id `mj import codex --session` takes.
2274        #[test]
2275        fn a_codex_row_is_imported_with_the_uuid_from_its_rollout_name() {
2276            let path = Path::new(
2277                "/home/user/.codex/sessions/2026/09/18/rollout-2026-09-18T09-15-00-7f3a1c20-0b11-4a55-9e0d-2c8a5d6f1b44.jsonl",
2278            );
2279            assert_eq!(
2280                wiki_continuation("abc123", "codex", path, false).unwrap(),
2281                WikiContinuation::Import {
2282                    harness: HarnessKind::Codex,
2283                    native_session_id: "7f3a1c20-0b11-4a55-9e0d-2c8a5d6f1b44".to_owned(),
2284                }
2285            );
2286        }
2287
2288        #[test]
2289        fn an_unknown_tool_is_an_error_that_names_it() {
2290            let error = wiki_continuation("abc123", "opencode", Path::new("/tmp/s.jsonl"), false)
2291                .unwrap_err();
2292            assert!(
2293                format!("{error:#}").contains("opencode"),
2294                "the error has to name the tool: {error:#}"
2295            );
2296        }
2297    }
2298
2299    use mj_checkpoint::archive::{
2300        ArchiveInput, BundleManifest, CanonicalExecutionState, CanonicalSessionSnapshot,
2301        CanonicalSessionState, CanonicalTranscriptBody, CanonicalTranscriptItem, SessionManifest,
2302        TargetManifest, write_archive_atomic,
2303    };
2304
2305    use super::*;
2306
2307    fn item(position: u64, body: CanonicalTranscriptBody) -> CanonicalTranscriptItem {
2308        // Only an agent message carries a content ordinal; the snapshot
2309        // validator rejects one on any other item and demands one here.
2310        let streamed = matches!(body, CanonicalTranscriptBody::Agent { .. });
2311        CanonicalTranscriptItem {
2312            stable_id: format!("item-{position}"),
2313            position,
2314            latest_content_event_ordinal: streamed.then_some(position),
2315            created_at_ms: 1_700_000_000_000 + i64::try_from(position).unwrap(),
2316            last_changed_at_ms: 1_700_000_000_000 + i64::try_from(position).unwrap(),
2317            body,
2318        }
2319    }
2320
2321    /// A managed checkpoint with one prompt, one reply, one tool call, and one
2322    /// thought, which is every transcript shape the adapter decides about.
2323    fn write_archive(directory: &Path, session_id: &str, frontier: u64) {
2324        let path = directory.join(format!(
2325            "{session_id}-{frontier}-archive-{}.hel.zip",
2326            "0".repeat(32)
2327        ));
2328        write_archive_atomic(
2329            &path,
2330            &ArchiveInput {
2331                session: SessionManifest {
2332                    id: session_id.into(),
2333                    title: "indexed session".into(),
2334                    harness_kind: mj_core::config::HarnessKind::Codex,
2335                    profile_id: "codex".into(),
2336                    native_session_id: "native-session".into(),
2337                    created_at: "2026-09-01T00:00:00Z".into(),
2338                    checkpointed_at: "2026-09-01T01:00:00Z".into(),
2339                    hel_version: "test".into(),
2340                    relay_version: "test".into(),
2341                    adapter_version: "test".into(),
2342                },
2343                target: TargetManifest {
2344                    template_id: "local".into(),
2345                    target_kind: "local-bare".into(),
2346                    details: BTreeMap::new(),
2347                },
2348                bundle: BundleManifest {
2349                    id: "project".into(),
2350                    primary_repository: "project".into(),
2351                },
2352                canonical_session: CanonicalSessionSnapshot {
2353                    command_ledger: None,
2354                    assessment_state: None,
2355                    event_frontier: 4,
2356                    event_frontier_digest: "a".repeat(64),
2357                    session: CanonicalSessionState {
2358                        execution: CanonicalExecutionState::Idle,
2359                        last_activity_at_ms: Some(1_700_000_000_004),
2360                        session_title: Some("snapshot title".into()),
2361                        configuration: Default::default(),
2362                    },
2363                    transcript: vec![
2364                        item(
2365                            1,
2366                            CanonicalTranscriptBody::User {
2367                                content: vec![serde_json::json!({
2368                                    "type": "text",
2369                                    "text": "index this session"
2370                                })],
2371                            },
2372                        ),
2373                        item(
2374                            2,
2375                            CanonicalTranscriptBody::Thought {
2376                                chunks: vec![serde_json::json!({
2377                                    "content": {"type": "text", "text": "pondering"}
2378                                })],
2379                                streaming: false,
2380                            },
2381                        ),
2382                        item(
2383                            3,
2384                            CanonicalTranscriptBody::Tool {
2385                                call: serde_json::json!({
2386                                    "toolCallId": "call-1",
2387                                    "title": "Edit config.toml",
2388                                    "kind": "edit",
2389                                    "status": "completed",
2390                                    "locations": [{"path": "/old/container/config.toml"}]
2391                                }),
2392                                terminal_outputs: Vec::new(),
2393                                terminal_refs: Vec::new(),
2394                                presentation: None,
2395                            },
2396                        ),
2397                        item(
2398                            4,
2399                            CanonicalTranscriptBody::Agent {
2400                                chunks: vec![serde_json::json!({
2401                                    "content": {"type": "text", "text": "done"}
2402                                })],
2403                                streaming: false,
2404                            },
2405                        ),
2406                    ],
2407                    queued_prompts: Vec::new(),
2408                },
2409                native_artifacts: Vec::new(),
2410                repositories: Vec::new(),
2411            },
2412        )
2413        .unwrap();
2414    }
2415
2416    fn adapter(directory: &Path, session_id: &str) -> MjolnirAdapter {
2417        adapter_with_live(directory, session_id, BTreeMap::new())
2418    }
2419
2420    fn adapter_with_live(
2421        directory: &Path,
2422        session_id: &str,
2423        live: BTreeMap<String, i64>,
2424    ) -> MjolnirAdapter {
2425        let record = SessionRecord {
2426            project: None,
2427            id: session_id.into(),
2428            ..record_template()
2429        };
2430        MjolnirAdapter {
2431            sessions_dir: directory.to_path_buf(),
2432            sessions: std::sync::Mutex::new(Sessions {
2433                records: [(session_id.to_owned(), record)].into_iter().collect(),
2434                subagent_ids: BTreeSet::new(),
2435                live,
2436            }),
2437            reload: false,
2438        }
2439    }
2440
2441    fn record_template() -> SessionRecord {
2442        SessionRecord {
2443            project: None,
2444            target_runtime: None,
2445            launch_base: None,
2446            launch_branch: None,
2447            checkout: None,
2448            publication: None,
2449            build_cache: None,
2450            container_workspace: None,
2451            subagents: None,
2452            create_managed_worktree: None,
2453            workspace_id: mj_core::workspace::DEFAULT_WORKSPACE_ID.to_owned(),
2454            archived: false,
2455            container_cpus: None,
2456            container_memory: None,
2457            id: "0123456789abcdef0123456789abcdef".into(),
2458            title: "indexed session".into(),
2459            harness_kind: mj_core::config::HarnessKind::Codex,
2460            last_profile: "codex".into(),
2461            bundle_id: "project".into(),
2462            project_directory: Some(PathBuf::from("/home/dev/project")),
2463            managed_worktree: None,
2464            target_template_id: "local-bare".into(),
2465            resource_allocation: None,
2466            additional_mounts: Vec::new(),
2467            state: mj_core::state::SessionState::Stopped,
2468            target: None,
2469            native_session_id: Some("native-session".into()),
2470            acp_session_title: Some("the harness title".into()),
2471            session_title_override: None,
2472            created_at: "2026-09-01T00:00:00Z".into(),
2473            updated_at: "2026-09-01T01:00:00Z".into(),
2474            viewed_through_event_ordinal: 0,
2475            draft_input: String::new(),
2476            last_error: None,
2477            last_checkpoint_error: None,
2478            checkpoint: None,
2479        }
2480    }
2481
2482    #[test]
2483    fn the_newest_checkpoint_of_each_session_is_one_indexed_key() {
2484        let directory = tempfile::tempdir().unwrap();
2485        let session_id = "0123456789abcdef0123456789abcdef";
2486        write_archive(directory.path(), session_id, 1);
2487        write_archive(directory.path(), session_id, 7);
2488        let adapter = adapter(directory.path(), session_id);
2489
2490        let store = adapter.store().expect("the adapter is a shared store");
2491        let key = format!("{}/{session_id}", directory.path().display());
2492        assert_eq!(
2493            store
2494                .keys
2495                .iter()
2496                .map(|(key, _)| key.as_str())
2497                .collect::<Vec<_>>(),
2498            vec![key.as_str()]
2499        );
2500        assert!(!store.had_error);
2501        assert_eq!(store.files.len(), 1);
2502        assert!(
2503            store.files[0]
2504                .file_name()
2505                .unwrap()
2506                .to_str()
2507                .unwrap()
2508                .contains("-7-archive-"),
2509            "the newest checkpoint is the one indexed: {:?}",
2510            store.files[0]
2511        );
2512        assert_eq!(
2513            adapter.reconcile_scope(),
2514            Some(format!("{}/", directory.path().display()))
2515        );
2516
2517        let session = adapter.parse_key(&key).unwrap();
2518        assert_eq!(session.id, session_id);
2519        assert_eq!(session.tool, "mjolnir");
2520        assert_eq!(session.path, PathBuf::from(&key));
2521        assert_eq!(session.project, "/home/dev/project");
2522        assert_eq!(session.title, "the harness title");
2523        assert!(!session.subagent);
2524        assert_eq!(
2525            session.messages.iter().map(|m| m.role).collect::<Vec<_>>(),
2526            vec![Role::User, Role::Tool, Role::Assistant]
2527        );
2528        assert_eq!(session.messages[0].text, "index this session");
2529        let tool: serde_json::Value = serde_json::from_str(&session.messages[1].text).unwrap();
2530        assert_eq!(tool["name"], "Edit");
2531        assert_eq!(tool["call"]["title"], "Edit config.toml");
2532        assert_eq!(session.messages[2].text, "done");
2533        assert_eq!(session.touched, vec!["/old/container/config.toml"]);
2534    }
2535
2536    #[test]
2537    fn provenance_backfill_repairs_an_unchanged_checkpoint_without_rebuilding_the_index() {
2538        let _held = tags::testing::lock();
2539        let (_index_dir, mut connection) = tags::testing::isolated_index();
2540        let directory = tempfile::tempdir().unwrap();
2541        write_archive(directory.path(), "old-session", 4);
2542        let source = adapter(directory.path(), "old-session");
2543        let key = source.key_for("old-session");
2544        tags::testing::index_row(&connection, "old-session", "mjolnir");
2545        connection
2546            .execute(
2547                "UPDATE files SET path = ?1 WHERE session_id = 'old-session'",
2548                [&key],
2549            )
2550            .unwrap();
2551        provenance::backfill(&mut connection, &source).unwrap();
2552        assert_eq!(
2553            sessionwiki::index::files_for(&connection, "old-session").unwrap(),
2554            vec!["/old/container/config.toml"]
2555        );
2556        provenance::backfill(&mut connection, &source).unwrap();
2557        assert_eq!(
2558            sessionwiki::index::sessions_for_file(&connection, "config.toml", 20)
2559                .unwrap()
2560                .len(),
2561            1
2562        );
2563    }
2564
2565    fn projection(session_id: &str) -> mj_core::state::MaterializedSession {
2566        use mj_core::transcript::{TranscriptBody, TranscriptItem};
2567        let mut projected = mj_core::state::MaterializedSession::empty(session_id);
2568        let mut push = |position: u64, body: TranscriptBody| {
2569            let streamed = matches!(body, TranscriptBody::Agent { .. });
2570            projected
2571                .transcript
2572                .push(std::sync::Arc::new(TranscriptItem {
2573                    stable_id: format!("item-{position}"),
2574                    position,
2575                    latest_content_event_ordinal: streamed.then_some(position),
2576                    created_at_ms: 1_700_000_000_000 + i64::try_from(position).unwrap(),
2577                    last_changed_at_ms: 1_700_000_000_000 + i64::try_from(position).unwrap(),
2578                    body,
2579                }));
2580        };
2581        push(
2582            1,
2583            TranscriptBody::User {
2584                content: vec![serde_json::json!({"type": "text", "text": "still talking"})],
2585            },
2586        );
2587        push(
2588            2,
2589            TranscriptBody::Thought {
2590                chunks: vec![serde_json::json!({"content": {"type": "text", "text": "hmm"}})],
2591                streaming: false,
2592            },
2593        );
2594        push(
2595            3,
2596            TranscriptBody::Tool {
2597                call: serde_json::json!({"toolCallId": "c1", "title": "Read README.md"}),
2598                terminal_outputs: Vec::new(),
2599                terminal_refs: Vec::new(),
2600                presentation: None,
2601            },
2602        );
2603        push(
2604            4,
2605            TranscriptBody::Agent {
2606                chunks: vec![serde_json::json!({"content": {"type": "text", "text": "reading"}})],
2607                streaming: false,
2608            },
2609        );
2610        projected.session_title = Some("the live title".into());
2611        projected
2612    }
2613
2614    /// A session that has never been checkpointed is indexed from the
2615    /// daemon's own projection, with the same roles a checkpoint would give.
2616    #[test]
2617    fn a_running_session_is_indexed_from_its_stored_transcript() {
2618        let session_id = "0123456789abcdef0123456789abcdef";
2619        let messages = projected_messages(&projection(session_id));
2620        assert_eq!(
2621            messages.iter().map(|m| m.role).collect::<Vec<_>>(),
2622            vec![Role::User, Role::Tool, Role::Assistant]
2623        );
2624        assert_eq!(messages[0].text, "still talking");
2625        assert_eq!(messages[2].text, "reading");
2626        let tool: serde_json::Value = serde_json::from_str(&messages[1].text).unwrap();
2627        assert_eq!(tool["name"], "Read");
2628        assert_eq!(tool["call"]["title"], "Read README.md");
2629    }
2630
2631    /// A running session is listed under the same key as a stopped one, with
2632    /// its own change token, so it is searchable before it is ever closed and
2633    /// reconciliation never archives it. When it stops, the key stays and the
2634    /// checkpoint becomes its source.
2635    #[test]
2636    fn a_running_session_is_listed_with_its_own_change_token() {
2637        let directory = tempfile::tempdir().unwrap();
2638        let running = "0123456789abcdef0123456789abcdef";
2639        let never_checkpointed = "fedcba9876543210fedcba9876543210";
2640        write_archive(directory.path(), running, 3);
2641        let live = adapter_with_live(
2642            directory.path(),
2643            running,
2644            BTreeMap::from([
2645                (running.to_owned(), 1_900_000_000),
2646                (never_checkpointed.to_owned(), 1_900_000_001),
2647            ]),
2648        );
2649
2650        let store = live.store().expect("the adapter is a shared store");
2651        let key_of = |session_id: &str| format!("{}/{session_id}", directory.path().display());
2652        assert_eq!(
2653            store.keys,
2654            vec![
2655                (
2656                    key_of(running),
2657                    1_900_000_000 * 1024 + i64::from(mj_transcript::summary::SUMMARY_VERSION)
2658                ),
2659                (
2660                    key_of(never_checkpointed),
2661                    1_900_000_001 * 1024 + i64::from(mj_transcript::summary::SUMMARY_VERSION)
2662                ),
2663            ],
2664            "a live session's own token replaces the checkpoint's"
2665        );
2666
2667        // Once it stops it leaves the live set, and the checkpoint's own
2668        // modification time is the token again.
2669        let stopped = adapter(directory.path(), running);
2670        let keys = stopped.store().expect("a shared store").keys;
2671        assert_eq!(keys.len(), 1);
2672        assert_eq!(keys[0].0, key_of(running));
2673        assert_ne!(keys[0].1, 1_900_000_000);
2674        assert_eq!(
2675            stopped.parse_key(&key_of(running)).unwrap().title,
2676            "the harness title",
2677            "a stopped session is parsed from its checkpoint"
2678        );
2679    }
2680
2681    /// Renaming a session leaves its conversation untouched, so only the
2682    /// record's own last update can tell the index the title moved.
2683    #[test]
2684    fn a_rename_moves_a_session_change_token() {
2685        let directory = tempfile::tempdir().unwrap();
2686        let session_id = "0123456789abcdef0123456789abcdef";
2687        write_archive(directory.path(), session_id, 1);
2688        let adapter = adapter(directory.path(), session_id);
2689        let before = adapter.store().expect("a shared store").keys[0].1;
2690
2691        {
2692            let mut sessions = adapter.sessions.lock().unwrap();
2693            let record = sessions.records.get_mut(session_id).unwrap();
2694            record.session_title_override = Some("the new name".into());
2695            record.updated_at = "2099-01-01T00:00:00Z".into();
2696        }
2697        let after = adapter.store().expect("a shared store").keys[0].1;
2698        assert!(
2699            after > before,
2700            "a renamed session is re-indexed: {before} then {after}"
2701        );
2702        assert_eq!(
2703            adapter
2704                .parse_key(&format!("{}/{session_id}", directory.path().display()))
2705                .unwrap()
2706                .title,
2707            "the new name"
2708        );
2709    }
2710
2711    fn indexed(messages: Vec<(Role, &str)>) -> sessionwiki::model::Session {
2712        Session {
2713            id: "0123456789abcdef0123456789abcdef".into(),
2714            tool: "mjolnir",
2715            path: PathBuf::from("/sessions/0123456789abcdef0123456789abcdef"),
2716            project: "/home/dev/project".into(),
2717            started: DateTime::from_timestamp_millis(1_700_000_000_000),
2718            ended: None,
2719            title: "the archived session".into(),
2720            subagent: false,
2721            messages: messages
2722                .into_iter()
2723                .map(|(role, text)| Message {
2724                    role,
2725                    text: text.to_owned(),
2726                    ts: None,
2727                })
2728                .collect(),
2729            touched: Vec::new(),
2730            edits: Vec::new(),
2731        }
2732    }
2733
2734    /// A hit is found whatever the case of the query or of the transcript, and
2735    /// the reported range covers the matched text in the returned block.
2736    #[test]
2737    fn transcript_hits_locates_case_insensitive_matches() {
2738        let session = indexed(vec![
2739            (Role::User, "Make the Tests green"),
2740            (Role::Assistant, "the tests are green now"),
2741        ]);
2742
2743        let found = hit_transcript(&session, "TESTS", 0, 4_000);
2744
2745        assert_eq!(found.blocks.len(), 2, "both messages contain the query");
2746        assert_eq!(found.blocks[0].role, "user");
2747        let (start, end) = found.blocks[0].hits[0];
2748        assert_eq!(&found.blocks[0].text[start..end], "Tests");
2749        let (start, end) = found.blocks[1].hits[0];
2750        assert_eq!(&found.blocks[1].text[start..end], "tests");
2751        assert!(!found.blocks[0].truncated);
2752        assert_eq!(found.omitted_after, 0);
2753    }
2754
2755    /// Context messages come back around each hit, with the gap between two
2756    /// groups counted rather than silently closed.
2757    #[test]
2758    fn transcript_hits_keeps_context_and_marks_omissions() {
2759        let session = indexed(vec![
2760            (Role::User, "zero"),
2761            (Role::Assistant, "one needle one"),
2762            (Role::Tool, "two"),
2763            (Role::User, "three"),
2764            (Role::Assistant, "four"),
2765            (Role::Tool, "five"),
2766            (Role::User, "six needle six"),
2767            (Role::Assistant, "seven"),
2768            (Role::User, "eight"),
2769        ]);
2770
2771        let found = hit_transcript(&session, "needle", 1, 4_000);
2772
2773        let shown: Vec<(&str, &str, usize)> = found
2774            .blocks
2775            .iter()
2776            .map(|block| {
2777                (
2778                    block.role.as_str(),
2779                    block.text.as_str(),
2780                    block.omitted_before,
2781                )
2782            })
2783            .collect();
2784        assert_eq!(
2785            shown,
2786            vec![
2787                ("user", "zero", 0),
2788                ("assistant", "one needle one", 0),
2789                ("tool", "two", 0),
2790                ("tool", "five", 2),
2791                ("user", "six needle six", 0),
2792                ("assistant", "seven", 0),
2793            ]
2794        );
2795        assert_eq!(found.omitted_after, 1, "the last message is not shown");
2796        assert!(found.blocks[0].hits.is_empty(), "context has no hits");
2797    }
2798
2799    /// A query that only occurs in tool output finds nothing, and a tool
2800    /// message beside a real match still comes back as context. Tool text is
2801    /// machine chatter: anchoring a passage on it opens the preview on command
2802    /// output the reader never wrote, and the preview collapses tool runs, so
2803    /// the match could not be shown even if it were returned.
2804    #[test]
2805    fn transcript_hits_never_anchor_on_tool_output() {
2806        let session = indexed(vec![
2807            (Role::User, "make it build"),
2808            (Role::Tool, "cargo build --needle"),
2809            (Role::Assistant, "it builds"),
2810        ]);
2811
2812        let only_in_a_tool = hit_transcript(&session, "needle", 1, 4_000);
2813        assert!(
2814            only_in_a_tool.blocks.is_empty(),
2815            "tool output must not anchor a passage, got {:?}",
2816            only_in_a_tool.blocks
2817        );
2818
2819        let beside_a_match = hit_transcript(&session, "builds", 1, 4_000);
2820        let shown: Vec<(&str, bool)> = beside_a_match
2821            .blocks
2822            .iter()
2823            .map(|block| (block.role.as_str(), !block.hits.is_empty()))
2824            .collect();
2825        assert_eq!(
2826            shown,
2827            vec![("tool", false), ("assistant", true)],
2828            "a tool message is still context around a real match"
2829        );
2830    }
2831
2832    /// A long message is cut down to the caller's budget around its first hit,
2833    /// not from the start, so the match is always in what comes back.
2834    #[test]
2835    fn transcript_hits_window_keeps_the_first_hit() {
2836        let filler = "x".repeat(4_000);
2837        let session = indexed(vec![(Role::User, &format!("{filler} needle {filler}"))]);
2838
2839        let found = hit_transcript(&session, "needle", 0, 100);
2840
2841        let block = &found.blocks[0];
2842        assert!(block.truncated);
2843        assert_eq!(block.text.chars().count(), 100);
2844        assert_eq!(block.hits.len(), 1, "the windowed text keeps its hit");
2845        let (start, end) = block.hits[0];
2846        assert_eq!(&block.text[start..end], "needle");
2847        assert!(
2848            start >= 20,
2849            "the window keeps lead-in before the hit, got {start}"
2850        );
2851    }
2852
2853    /// The snapshot a restore hands to compaction has to satisfy the same
2854    /// validator a real checkpoint does, and has to carry every message in
2855    /// order.
2856    #[test]
2857    fn a_restored_snapshot_is_a_valid_transcript_of_the_indexed_session() {
2858        let snapshot = snapshot_of(&indexed(vec![
2859            (Role::User, "make the tests green"),
2860            (Role::Tool, "Read src/lib.rs"),
2861            (Role::Assistant, "they are green now"),
2862            (Role::User, "  "),
2863        ]))
2864        .unwrap();
2865
2866        snapshot.validate().expect("the snapshot is well formed");
2867        assert_eq!(snapshot.event_frontier, 3);
2868        assert_eq!(
2869            snapshot.session.session_title.as_deref(),
2870            Some("the archived session")
2871        );
2872        assert!(snapshot.session.last_activity_at_ms.is_some());
2873        let bodies = snapshot
2874            .transcript
2875            .iter()
2876            .map(|item| match &item.body {
2877                mj_core::archive::CanonicalTranscriptBody::User { content } => (
2878                    "user",
2879                    mj_core::transcript::materialized_content_text(content),
2880                ),
2881                mj_core::archive::CanonicalTranscriptBody::Agent { chunks, .. } => (
2882                    "agent",
2883                    mj_core::transcript::materialized_chunks_text(chunks),
2884                ),
2885                mj_core::archive::CanonicalTranscriptBody::Tool { call, .. } => (
2886                    "tool",
2887                    call["title"].as_str().unwrap_or_default().to_owned(),
2888                ),
2889                _ => ("other", String::new()),
2890            })
2891            .collect::<Vec<_>>();
2892        assert_eq!(
2893            bodies,
2894            vec![
2895                ("user", "make the tests green".to_owned()),
2896                ("tool", "Read src/lib.rs".to_owned()),
2897                ("agent", "they are green now".to_owned()),
2898            ],
2899            "the blank message is dropped and every other one keeps its role"
2900        );
2901    }
2902
2903    /// Compaction attaches assistant and tool items to the open turn, so an
2904    /// index that starts mid-conversation must not produce a snapshot whose
2905    /// first item has no turn to join.
2906    #[test]
2907    fn messages_before_the_first_prompt_are_dropped() {
2908        let snapshot = snapshot_of(&indexed(vec![
2909            (Role::Assistant, "still working"),
2910            (Role::User, "carry on"),
2911        ]))
2912        .unwrap();
2913        assert_eq!(snapshot.transcript.len(), 1);
2914        assert_eq!(snapshot.transcript[0].position, 1);
2915        snapshot.validate().unwrap();
2916
2917        let error = snapshot_of(&indexed(vec![(Role::Assistant, "nobody asked")])).unwrap_err();
2918        assert!(
2919            error.to_string().contains("no prompt"),
2920            "a session with no prompt cannot be restored: {error}"
2921        );
2922        // `mj sessions --session` offers a restore by the same rule.
2923        assert!(!has_prompt(&indexed(vec![(
2924            Role::Assistant,
2925            "nobody asked"
2926        )])));
2927        assert!(has_prompt(&indexed(vec![(Role::User, "carry on")])));
2928    }
2929
2930    fn record(
2931        session_id: &str,
2932        state: mj_core::state::SessionState,
2933        updated_at: &str,
2934    ) -> SessionRecord {
2935        SessionRecord {
2936            project: None,
2937            id: session_id.into(),
2938            state,
2939            updated_at: updated_at.into(),
2940            ..record_template()
2941        }
2942    }
2943
2944    fn child(child_session_id: &str, parent_session_id: &str) -> mj_core::subagent::SubagentRecord {
2945        mj_core::subagent::SubagentRecord {
2946            child_session_id: child_session_id.into(),
2947            parent_session_id: parent_session_id.into(),
2948            task_name: "task".into(),
2949            profile_id: "codex".into(),
2950            model: None,
2951            effort: None,
2952            working_directory: PathBuf::new(),
2953            initial_prompt: "do the thing".into(),
2954            request_key: "key".into(),
2955            created_at: "2026-09-01T00:00:00Z".into(),
2956            noticed_turn: None,
2957            handback_tool: false,
2958        }
2959    }
2960
2961    fn ready(
2962        sessions: Vec<SessionRecord>,
2963        children: Vec<mj_core::subagent::SubagentRecord>,
2964    ) -> Vec<String> {
2965        let now = parse_time("2026-09-10T00:00:00Z").unwrap();
2966        sessions_ready_to_archive(
2967            &sessions
2968                .into_iter()
2969                .map(|record| (record.id.clone(), record))
2970                .collect(),
2971            &children
2972                .into_iter()
2973                .map(|child| (child.child_session_id.clone(), child))
2974                .collect(),
2975            now,
2976            3,
2977        )
2978    }
2979
2980    #[test]
2981    fn aged_clone_requires_clean_published_evidence_for_its_current_checkpoint() {
2982        let id = "0123456789abcdef0123456789abcdef";
2983        let root = PathBuf::from(format!("/srv/project/.mj/clones/{id}"));
2984        let mut session = record(
2985            id,
2986            mj_core::state::SessionState::Stopped,
2987            "2026-09-01T00:00:00Z",
2988        );
2989        session.project_directory = Some(root.clone());
2990        session.managed_worktree = Some(mj_core::state::ManagedWorktree {
2991            kind: mj_core::state::ManagedCheckoutKind::Clone,
2992            source_project_directory: "/srv/project".into(),
2993            source_repository: "/srv/project".into(),
2994            worktree_root: root,
2995            branch: "feature".into(),
2996            target: mj_core::state::ManagedWorktreeTarget::Local,
2997            base_commit: Some("1".repeat(40)),
2998        });
2999        session.checkpoint = Some(mj_core::state::CheckpointMetadata {
3000            archive_path: "sessions/checkpoint.hel.zip".into(),
3001            sha256: "a".repeat(64),
3002            created_at: "2026-09-01T00:00:00Z".into(),
3003            event_frontier: 0,
3004        });
3005        assert!(ready(vec![session.clone()], vec![]).is_empty());
3006        session.publication = Some(mj_core::state::PublicationAssessment {
3007            checkpoint_sha256: "a".repeat(64),
3008            state: mj_core::state::PublicationState::Published,
3009            dirty: false,
3010            stashed: false,
3011            saved_commits: vec!["2".repeat(40)],
3012            destinations: vec!["https://example.test/repository.git".into()],
3013            checked_at: "2026-09-01T01:00:00Z".into(),
3014            reason: Some("feature branch was pushed but not merged".into()),
3015        });
3016        assert_eq!(ready(vec![session.clone()], vec![]), vec![id]);
3017        session.publication.as_mut().unwrap().stashed = true;
3018        assert!(ready(vec![session.clone()], vec![]).is_empty());
3019        session.publication.as_mut().unwrap().stashed = false;
3020        session.publication.as_mut().unwrap().checkpoint_sha256 = "b".repeat(64);
3021        assert!(ready(vec![session], vec![]).is_empty());
3022    }
3023
3024    /// A session whose checkpoint archive and attachments sit under `root`.
3025    fn sized_session(
3026        root: &Path,
3027        session_id: &str,
3028        updated_at: &str,
3029        checkpoint_bytes: usize,
3030        attachment_bytes: &[usize],
3031    ) -> SessionRecord {
3032        let archive_path = root.join(format!("{session_id}.hel.zip"));
3033        std::fs::write(&archive_path, vec![b'c'; checkpoint_bytes]).unwrap();
3034        if !attachment_bytes.is_empty() {
3035            let attachments = root
3036                .join(session_id)
3037                .join(mj_core::attachment::ATTACHMENT_DIR);
3038            std::fs::create_dir_all(&attachments).unwrap();
3039            for (index, size) in attachment_bytes.iter().enumerate() {
3040                std::fs::write(attachments.join(format!("{index}.png")), vec![b'a'; *size])
3041                    .unwrap();
3042            }
3043        }
3044        SessionRecord {
3045            project: None,
3046            checkpoint: Some(mj_core::state::CheckpointMetadata {
3047                archive_path,
3048                sha256: "0".repeat(64),
3049                created_at: updated_at.into(),
3050                event_frontier: 1,
3051            }),
3052            ..record(
3053                session_id,
3054                mj_core::state::SessionState::Stopped,
3055                updated_at,
3056            )
3057        }
3058    }
3059
3060    #[test]
3061    fn the_space_preview_sizes_every_session_and_only_the_aged_ones_as_reclaimable() {
3062        let directory = tempfile::tempdir().unwrap();
3063        let root = directory.path();
3064        let sessions: mj_core::snapshot_map::SnapshotMap<String, SessionRecord> = [
3065            sized_session(root, "old-stopped", "2026-09-01T00:00:00Z", 1000, &[10, 20]),
3066            sized_session(root, "just-stopped", "2026-09-09T00:00:00Z", 500, &[]),
3067            // A record whose checkpoint file is already gone counts as zero
3068            // rather than failing the whole estimate.
3069            SessionRecord {
3070                project: None,
3071                checkpoint: Some(mj_core::state::CheckpointMetadata {
3072                    archive_path: root.join("missing.hel.zip"),
3073                    sha256: "0".repeat(64),
3074                    created_at: "2026-09-01T00:00:00Z".into(),
3075                    event_frontier: 1,
3076                }),
3077                ..record(
3078                    "lost-checkpoint",
3079                    mj_core::state::SessionState::Stopped,
3080                    "2026-09-01T00:00:00Z",
3081                )
3082            },
3083        ]
3084        .into_iter()
3085        .map(|record| (record.id.clone(), record))
3086        .collect();
3087        let now = parse_time("2026-09-10T00:00:00Z").unwrap();
3088
3089        let all = archive_space_over(root, &sessions, &Default::default(), now, None);
3090        assert_eq!(all.sessions, 3);
3091        assert_eq!(all.bytes, 1530);
3092        assert_eq!(all.reclaimable_sessions, 0);
3093        assert_eq!(all.reclaimable_bytes, 0);
3094
3095        let aged = archive_space_over(root, &sessions, &Default::default(), now, Some(3));
3096        assert_eq!(aged.bytes, 1530);
3097        assert_eq!(
3098            (aged.reclaimable_sessions, aged.reclaimable_bytes),
3099            (2, 1030),
3100            "only the sessions the job would archive count, attachments included"
3101        );
3102    }
3103
3104    #[test]
3105    fn only_stopped_sessions_past_the_cut_off_are_archived() {
3106        use mj_core::state::SessionState;
3107        let selected = ready(
3108            vec![
3109                record("old-stopped", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3110                record(
3111                    "just-stopped",
3112                    SessionState::Stopped,
3113                    "2026-09-09T00:00:00Z",
3114                ),
3115                record("old-running", SessionState::Running, "2026-09-01T00:00:00Z"),
3116                record("old-error", SessionState::Error, "2026-09-01T00:00:00Z"),
3117                record("unparsable", SessionState::Stopped, "not a time"),
3118                // Exactly the cut-off counts as old enough.
3119                record("at-the-edge", SessionState::Stopped, "2026-09-07T00:00:00Z"),
3120            ],
3121            Vec::new(),
3122        );
3123        assert_eq!(selected, vec!["at-the-edge", "old-stopped"]);
3124    }
3125
3126    #[test]
3127    fn a_child_the_pass_is_not_archiving_holds_its_parent_back() {
3128        use mj_core::state::SessionState;
3129        let selected = ready(
3130            vec![
3131                record("parent", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3132                record(
3133                    "running-child",
3134                    SessionState::Running,
3135                    "2026-09-01T00:00:00Z",
3136                ),
3137            ],
3138            vec![child("running-child", "parent")],
3139        );
3140        assert!(selected.is_empty(), "the parent must wait: {selected:?}");
3141
3142        let selected = ready(
3143            vec![
3144                record("parent", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3145                record("young-child", SessionState::Stopped, "2026-09-09T00:00:00Z"),
3146            ],
3147            vec![child("young-child", "parent")],
3148        );
3149        assert!(selected.is_empty(), "the parent must wait: {selected:?}");
3150
3151        // A child whose record is already gone holds nothing open.
3152        let selected = ready(
3153            vec![record(
3154                "parent",
3155                SessionState::Stopped,
3156                "2026-09-01T00:00:00Z",
3157            )],
3158            vec![child("departed-child", "parent")],
3159        );
3160        assert_eq!(selected, vec!["parent"]);
3161    }
3162
3163    #[test]
3164    fn children_are_archived_before_their_parents() {
3165        use mj_core::state::SessionState;
3166        let selected = ready(
3167            vec![
3168                record("parent", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3169                record("child", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3170                record("grandchild", SessionState::Stopped, "2026-09-01T00:00:00Z"),
3171            ],
3172            vec![child("child", "parent"), child("grandchild", "child")],
3173        );
3174        assert_eq!(selected, vec!["grandchild", "child", "parent"]);
3175    }
3176
3177    #[test]
3178    fn native_adapters_cover_every_enabled_profile_home() {
3179        use mj_core::config::{Config, HarnessKind, HarnessProfile};
3180
3181        fn profile(kind: HarnessKind, home: &str, enabled: bool) -> HarnessProfile {
3182            HarnessProfile {
3183                enabled,
3184                kind,
3185                home: PathBuf::from(home),
3186                environment: Default::default(),
3187                context_window_bytes: None,
3188                subagents: Default::default(),
3189                guardian_review_model: None,
3190            }
3191        }
3192
3193        let mut config = Config::default();
3194        for (id, built) in [
3195            (
3196                "codex",
3197                profile(HarnessKind::Codex, "/home/dev/.codex3", true),
3198            ),
3199            (
3200                "codex-ds",
3201                profile(HarnessKind::Codex, "/home/dev/.codex-ds", true),
3202            ),
3203            // A second profile on one home must not add a second adapter.
3204            (
3205                "codex-alt",
3206                profile(HarnessKind::Codex, "/home/dev/.codex3", true),
3207            ),
3208            (
3209                "codex-off",
3210                profile(HarnessKind::Codex, "/home/dev/.codex-off", false),
3211            ),
3212            (
3213                "claude",
3214                profile(HarnessKind::Claude, "/home/dev/.claude4", true),
3215            ),
3216            ("kimi", profile(HarnessKind::Kimi, "/home/dev/.kimi", true)),
3217            ("grok", profile(HarnessKind::Grok, "/home/dev/.grok", true)),
3218            ("muse", profile(HarnessKind::Muse, "/home/dev/muse", true)),
3219            (
3220                "muse-off",
3221                profile(HarnessKind::Muse, "/home/dev/muse-off", false),
3222            ),
3223        ] {
3224            config.profiles.insert(id.into(), built);
3225        }
3226
3227        let adapters = native_adapters(&config);
3228        let roots: Vec<(&str, Option<PathBuf>)> = adapters
3229            .iter()
3230            .map(|adapter| (adapter.name(), adapter.root()))
3231            .collect();
3232
3233        let codex: Vec<&Option<PathBuf>> = roots
3234            .iter()
3235            .filter(|(name, _)| *name == "codex")
3236            .map(|(_, root)| root)
3237            .collect();
3238        assert_eq!(
3239            codex,
3240            vec![
3241                &Some(PathBuf::from("/home/dev/.codex3/sessions")),
3242                &Some(PathBuf::from("/home/dev/.codex-ds/sessions")),
3243            ],
3244            "one adapter per enabled Codex home, deduplicated: {roots:?}"
3245        );
3246
3247        let claude: Vec<&Option<PathBuf>> = roots
3248            .iter()
3249            .filter(|(name, _)| *name == "claude-code")
3250            .map(|(_, root)| root)
3251            .collect();
3252        assert_eq!(
3253            claude,
3254            vec![&Some(PathBuf::from("/home/dev/.claude4/projects"))],
3255            "one adapter for the enabled Claude home: {roots:?}"
3256        );
3257
3258        for (_, root) in &roots {
3259            let Some(root) = root else { continue };
3260            let text = root.to_string_lossy();
3261            assert!(
3262                !text.contains(".codex-off"),
3263                "a disabled profile must not be indexed: {roots:?}"
3264            );
3265            assert!(
3266                !text.ends_with("/.codex/sessions") && !text.ends_with("/.claude/projects"),
3267                "the stock homes are not indexed unless a profile names them: {roots:?}"
3268            );
3269        }
3270
3271        // SessionWiki has no adapter for these three, so Mjolnir supplies one
3272        // per enabled profile home under its own tool name.
3273        for (name, root) in [
3274            ("kimi-code", PathBuf::from("/home/dev/.kimi/sessions")),
3275            ("grok-build", PathBuf::from("/home/dev/.grok/sessions")),
3276            (
3277                "muse",
3278                mj_checkpoint::native::muse_sessions_root(Path::new("/home/dev/muse")).unwrap(),
3279            ),
3280        ] {
3281            let found: Vec<&Option<PathBuf>> = roots
3282                .iter()
3283                .filter(|(found, _)| *found == name)
3284                .map(|(_, root)| root)
3285                .collect();
3286            assert_eq!(found, vec![&Some(root)], "one {name} adapter: {roots:?}");
3287        }
3288
3289        for (_, root) in &roots {
3290            let Some(root) = root else { continue };
3291            assert!(
3292                !root.to_string_lossy().contains("muse-off"),
3293                "a disabled profile must not be indexed: {roots:?}"
3294            );
3295        }
3296
3297        assert!(
3298            roots.iter().any(|(name, _)| *name == "gemini"),
3299            "the other built-in adapters are kept: {roots:?}"
3300        );
3301    }
3302
3303    /// A Mjolnir row carries the target, profile and harness the sync stored
3304    /// in the index; a row from another tool carries none, because only
3305    /// Mjolnir writes those tags.
3306    #[test]
3307    fn query_rows_returns_the_indexed_target_profile_and_harness() {
3308        let _held = tags::testing::lock();
3309        let (_directory, connection) = tags::testing::isolated_index();
3310        tags::testing::index_row(&connection, "mj-session", TOOL);
3311        tags::testing::index_row(&connection, "codex-session", "codex");
3312        tags::write(
3313            &connection,
3314            "mj-session",
3315            &tags::MjTags {
3316                target: Some("Prod-Box".into()),
3317                profile: Some("codex-Main".into()),
3318                harness: Some("codex".into()),
3319            },
3320        )
3321        .expect("write the session metadata");
3322
3323        let rows = query_rows("", 10, &BTreeSet::new(), false).expect("query the index");
3324        let mjolnir = rows
3325            .iter()
3326            .find(|row| row.id == "mj-session")
3327            .expect("the Mjolnir row is returned");
3328        assert_eq!(mjolnir.target.as_deref(), Some("Prod-Box"));
3329        assert_eq!(mjolnir.profile.as_deref(), Some("codex-Main"));
3330        assert_eq!(mjolnir.harness.as_deref(), Some("codex"));
3331
3332        let codex = rows
3333            .iter()
3334            .find(|row| row.id == "codex-session")
3335            .expect("the Codex row is returned");
3336        assert_eq!(codex.target, None);
3337        assert_eq!(codex.profile, None);
3338        assert_eq!(codex.harness, None);
3339    }
3340
3341    #[test]
3342    fn one_flag_keeps_sub_agents_out_of_every_query_path() {
3343        let _held = tags::testing::lock();
3344        let (_directory, connection) = tags::testing::isolated_index();
3345        for (session_id, kind) in [("main-session", "main"), ("sub-session", "sub")] {
3346            tags::testing::index_row(&connection, session_id, "claude");
3347            connection
3348                .execute(
3349                    "UPDATE files SET kind = ?2 WHERE session_id = ?1",
3350                    rusqlite::params![session_id, kind],
3351                )
3352                .expect("set the session kind");
3353            connection
3354                .execute(
3355                    "INSERT INTO messages(session_id, role, text)
3356                     VALUES (?1, 'user', 'fix the bridge derivation zq')",
3357                    [session_id],
3358                )
3359                .expect("insert a message");
3360            connection
3361                .execute(
3362                    "INSERT INTO msgs(rowid, text) VALUES (?1, 'fix the bridge derivation zq')",
3363                    [connection.last_insert_rowid()],
3364                )
3365                .expect("index the message");
3366        }
3367        let ids = |query: &str, include_subagents: bool| {
3368            let mut ids: Vec<String> = query_rows(query, 10, &BTreeSet::new(), include_subagents)
3369                .expect("query the index")
3370                .into_iter()
3371                .map(|row| row.id)
3372                .collect();
3373            ids.sort();
3374            ids
3375        };
3376
3377        // The recent list, full-text search, short-query scan, and title match.
3378        for query in ["", "bridge derivation", "zq", "an indexed session"] {
3379            assert_eq!(ids(query, false), ["main-session"], "query {query:?}");
3380            assert_eq!(
3381                ids(query, true),
3382                ["main-session", "sub-session"],
3383                "query {query:?}"
3384            );
3385        }
3386    }
3387
3388    /// I1-5: a phrase only a sub-agent wrote reaches its parent's index as
3389    /// tool text (Claude Code records the Task prompt and the sub-agent's
3390    /// answer as the parent's tool call and tool result). The resume search
3391    /// matched the parent on it while the preview, which never anchors on
3392    /// tool output, said "no hits". A match counts only where the preview can
3393    /// show it.
3394    #[test]
3395    fn a_phrase_only_in_a_sub_agents_transcript_does_not_match_its_parent() {
3396        let _held = tags::testing::lock();
3397        let (_directory, connection) = tags::testing::isolated_index();
3398        let message = |session_id: &str, role: &str, text: &str| {
3399            connection
3400                .execute(
3401                    "INSERT INTO messages(session_id, role, text) VALUES (?1, ?2, ?3)",
3402                    rusqlite::params![session_id, role, text],
3403                )
3404                .expect("insert a message");
3405            connection
3406                .execute(
3407                    "INSERT INTO msgs(rowid, text) VALUES (?1, ?2)",
3408                    rusqlite::params![connection.last_insert_rowid(), text],
3409                )
3410                .expect("index the message");
3411        };
3412        for (session_id, kind) in [("parent", "main"), ("child", "sub")] {
3413            tags::testing::index_row(&connection, session_id, "claude");
3414            connection
3415                .execute(
3416                    "UPDATE files SET kind = ?2 WHERE session_id = ?1",
3417                    rusqlite::params![session_id, kind],
3418                )
3419                .expect("set the session kind");
3420        }
3421        message("parent", "user", "look into the relay journal");
3422        message("parent", "tool", "Task {\"prompt\":\"read the journal\"}");
3423        message("parent", "tool", "the journal uses a quokka checksum");
3424        message(
3425            "parent",
3426            "assistant",
3427            "The journal is fine; the parent zebra ends here.",
3428        );
3429        message("child", "user", "read the journal");
3430        message("child", "assistant", "the journal uses a quokka checksum");
3431
3432        let ids = |query: &str, include_subagents: bool| {
3433            let mut ids: Vec<String> = query_rows(query, 10, &BTreeSet::new(), include_subagents)
3434                .expect("query the index")
3435                .into_iter()
3436                .map(|row| row.id)
3437                .collect();
3438            ids.sort();
3439            ids
3440        };
3441        assert!(
3442            ids("quokka", false).is_empty(),
3443            "{:?}",
3444            ids("quokka", false)
3445        );
3446        // The agents' history search asks for sub-agents and keeps tool text.
3447        assert_eq!(ids("quokka", true), ["child", "parent"]);
3448        assert_eq!(ids("parent zebra", false), ["parent"]);
3449    }
3450
3451    #[test]
3452    fn short_query_scan_also_ignores_tool_only_matches() {
3453        let _held = tags::testing::lock();
3454        let (_directory, connection) = tags::testing::isolated_index();
3455        tags::testing::index_row(&connection, "parent", "claude");
3456        connection
3457            .execute(
3458                "INSERT INTO messages(session_id, role, text) VALUES ('parent', 'tool', 'qx')",
3459                [],
3460            )
3461            .expect("insert a message");
3462        assert!(
3463            query_rows("qx", 10, &BTreeSet::new(), false)
3464                .expect("query the index")
3465                .is_empty()
3466        );
3467    }
3468
3469    #[test]
3470    fn full_text_search_fills_its_result_limit_after_skipping_sub_agents() {
3471        let _held = tags::testing::lock();
3472        let (_directory, connection) = tags::testing::isolated_index();
3473        for (id, kind, text) in [
3474            ("sub", "sub", "restic restic restic restic"),
3475            ("main", "main", "restic cleanup"),
3476        ] {
3477            tags::testing::index_row(&connection, id, "codex");
3478            connection
3479                .execute(
3480                    "UPDATE files SET kind = ?2 WHERE session_id = ?1",
3481                    rusqlite::params![id, kind],
3482                )
3483                .unwrap();
3484            connection
3485                .execute(
3486                    "INSERT INTO messages(session_id, role, text) VALUES (?1, 'user', ?2)",
3487                    rusqlite::params![id, text],
3488                )
3489                .unwrap();
3490            connection
3491                .execute(
3492                    "INSERT INTO msgs(rowid, text) VALUES (?1, ?2)",
3493                    rusqlite::params![connection.last_insert_rowid(), text],
3494                )
3495                .unwrap();
3496        }
3497
3498        let rows = query_rows("restic", 1, &BTreeSet::new(), false).unwrap();
3499        assert_eq!(
3500            rows.iter().map(|row| row.id.as_str()).collect::<Vec<_>>(),
3501            ["main"]
3502        );
3503    }
3504
3505    /// A runtime of the test's own, so a test can hold the index lock without
3506    /// holding it across an await.
3507    fn block_on<F: std::future::Future>(future: F) -> F::Output {
3508        tokio::runtime::Builder::new_current_thread()
3509            .enable_all()
3510            .build()
3511            .unwrap()
3512            .block_on(future)
3513    }
3514
3515    /// R2-11's fallback. A destroy that outwaits a running sync pass, such as
3516    /// a first build, indexes the session on its own from the rows the pass
3517    /// would write, and does not wait for the pass. The session is then found
3518    /// by its id with no record left, as after the destroy.
3519    #[test]
3520    fn a_destroy_indexes_the_session_itself_when_the_sync_outlasts_the_wait() {
3521        let _held = tags::testing::lock();
3522        let (_index_dir, _connection) = tags::testing::isolated_index();
3523        let directory = tempfile::tempdir().unwrap();
3524        let session_id = "0123456789abcdef0123456789abcdef";
3525        write_archive(directory.path(), session_id, 1);
3526        let source = adapter(directory.path(), session_id);
3527
3528        let started = Instant::now();
3529        let outcome = block_on(index_before_destroy_with(
3530            std::future::pending::<Result<()>>(),
3531            Duration::from_millis(200),
3532            move || capture_sessions_from(&source, &[session_id.to_owned()]),
3533            Duration::from_millis(50),
3534        ));
3535
3536        assert_eq!(outcome, IndexedBeforeDestroy::WrittenDirectly);
3537        assert!(
3538            started.elapsed() < Duration::from_secs(10),
3539            "the destroy must not wait for the pass: {:?}",
3540            started.elapsed()
3541        );
3542        let found = wiki_session(session_id, &BTreeSet::new())
3543            .unwrap()
3544            .expect("the session is found by its id");
3545        assert_eq!(found.status, WikiSessionStatus::Archived);
3546        assert_eq!(found.tool, TOOL);
3547        assert_eq!(
3548            found.path,
3549            PathBuf::from(format!("{}/{session_id}", directory.path().display()))
3550        );
3551        assert_eq!(found.title, "the harness title");
3552        assert_eq!(
3553            found.harness,
3554            Some(HarnessKind::Codex),
3555            "the session's metadata is written beside its row"
3556        );
3557        assert!(!found.nothing_to_restore);
3558    }
3559
3560    /// A sync pass that finishes within the wait has indexed the session, so
3561    /// nothing is read or written on its own.
3562    #[test]
3563    fn a_sync_that_finishes_in_time_is_all_a_destroy_waits_for() {
3564        let outcome = block_on(index_before_destroy_with(
3565            async { Ok(()) },
3566            DESTROY_SYNC_WAIT,
3567            || -> Result<Vec<CapturedSession>> {
3568                panic!("a finished pass leaves nothing to index on its own")
3569            },
3570            Duration::from_millis(50),
3571        ));
3572        assert_eq!(outcome, IndexedBeforeDestroy::Synced);
3573    }
3574
3575    /// While another writer holds the index, as a first build does while it
3576    /// parses one tool's sessions, the destroy goes ahead and the rows it read
3577    /// are written once the index is free.
3578    #[test]
3579    fn a_busy_index_takes_the_destroyed_session_once_it_is_free() {
3580        let _held = tags::testing::lock();
3581        let (_index_dir, writer) = tags::testing::isolated_index();
3582        let directory = tempfile::tempdir().unwrap();
3583        let session_id = "0123456789abcdef0123456789abcdef";
3584        write_archive(directory.path(), session_id, 1);
3585        let source = adapter(directory.path(), session_id);
3586
3587        block_on(async {
3588            writer.execute_batch("BEGIN IMMEDIATE").unwrap();
3589            let outcome = index_before_destroy_with(
3590                std::future::pending::<Result<()>>(),
3591                Duration::from_millis(50),
3592                move || capture_sessions_from(&source, &[session_id.to_owned()]),
3593                Duration::from_millis(50),
3594            )
3595            .await;
3596            assert_eq!(outcome, IndexedBeforeDestroy::Deferred);
3597            assert!(
3598                wiki_session(session_id, &BTreeSet::new())
3599                    .unwrap()
3600                    .is_none(),
3601                "nothing is written while the other writer holds the index"
3602            );
3603
3604            writer.execute_batch("COMMIT").unwrap();
3605            let deadline = Instant::now() + Duration::from_secs(30);
3606            while wiki_session(session_id, &BTreeSet::new())
3607                .unwrap()
3608                .is_none()
3609            {
3610                assert!(
3611                    Instant::now() < deadline,
3612                    "the deferred row never reached the index"
3613                );
3614                tokio::time::sleep(Duration::from_millis(50)).await;
3615            }
3616        });
3617    }
3618
3619    /// A destroy of a session the index already holds as it is now runs no
3620    /// sync pass, so destroying a workspace or a parent with sub-agents costs
3621    /// one pass at most. The token compared is the one SessionWiki stored.
3622    #[test]
3623    fn a_session_the_index_holds_as_it_is_now_needs_no_indexing() {
3624        let _held = tags::testing::lock();
3625        let (_index_dir, _connection) = tags::testing::isolated_index();
3626        let directory = tempfile::tempdir().unwrap();
3627        let session_id = "0123456789abcdef0123456789abcdef";
3628        let never_prompted = "fedcba9876543210fedcba9876543210";
3629        write_archive(directory.path(), session_id, 1);
3630        let source = adapter(directory.path(), session_id);
3631        let ids = [session_id.to_owned(), never_prompted.to_owned()];
3632
3633        assert_eq!(
3634            unindexed(&source, &ids).unwrap(),
3635            [session_id],
3636            "a session with no conversation has nothing to index"
3637        );
3638        let captured = Arc::new(capture_sessions_from(&source, &ids).unwrap());
3639        write_captured(&captured).unwrap();
3640        assert!(unindexed(&source, &ids).unwrap().is_empty());
3641
3642        // A rename moves the change token, so the row is stale again.
3643        source
3644            .sessions
3645            .lock()
3646            .unwrap()
3647            .records
3648            .get_mut(session_id)
3649            .unwrap()
3650            .updated_at = "2099-01-01T00:00:00Z".into();
3651        assert_eq!(unindexed(&source, &ids).unwrap(), [session_id]);
3652    }
3653
3654    #[test]
3655    fn a_session_tree_holds_the_sub_agents_below_it_and_nothing_else() {
3656        let record = |id: &str| {
3657            (
3658                id.to_owned(),
3659                SessionRecord {
3660                    project: None,
3661                    id: id.into(),
3662                    ..record_template()
3663                },
3664            )
3665        };
3666        let state = State {
3667            sessions: [
3668                record("parent"),
3669                record("child"),
3670                record("grandchild"),
3671                record("sibling"),
3672            ]
3673            .into_iter()
3674            .collect(),
3675            subagents: [
3676                ("child".to_owned(), child("child", "parent")),
3677                ("grandchild".to_owned(), child("grandchild", "child")),
3678                ("sibling".to_owned(), child("sibling", "other-parent")),
3679            ]
3680            .into_iter()
3681            .collect(),
3682            ..State::default()
3683        };
3684        assert_eq!(
3685            session_tree(&state, "parent"),
3686            ["parent", "child", "grandchild"]
3687        );
3688        assert!(session_tree(&state, "unknown").is_empty());
3689    }
3690
3691    /// The text search over an in-memory index shaped like SessionWiki's.
3692    mod text_search {
3693        use super::super::{SessionTextMatch, SessionTextMatchKind, text_matches_in};
3694        use std::collections::BTreeSet;
3695
3696        fn index(sessions: &[(&str, &[(&str, &str)])]) -> rusqlite::Connection {
3697            let connection = rusqlite::Connection::open_in_memory().unwrap();
3698            connection
3699                .execute_batch(
3700                    "CREATE TABLE files(path TEXT PRIMARY KEY, session_id TEXT NOT NULL,
3701                                        tool TEXT NOT NULL);
3702                     CREATE TABLE messages(id INTEGER PRIMARY KEY, session_id TEXT NOT NULL,
3703                                           role TEXT NOT NULL, text TEXT NOT NULL);
3704                     CREATE VIRTUAL TABLE msgs USING fts5(
3705                         text, content='messages', content_rowid='id', tokenize='trigram');",
3706                )
3707                .unwrap();
3708            for (id, messages) in sessions {
3709                connection
3710                    .execute(
3711                        "INSERT INTO files VALUES (?1, ?2, 'mjolnir')",
3712                        rusqlite::params![format!("/checkpoints/{id}"), id],
3713                    )
3714                    .unwrap();
3715                for (role, text) in *messages {
3716                    connection
3717                        .execute(
3718                            "INSERT INTO messages(session_id, role, text) VALUES (?1, ?2, ?3)",
3719                            rusqlite::params![id, role, text],
3720                        )
3721                        .unwrap();
3722                    let rowid = connection.last_insert_rowid();
3723                    connection
3724                        .execute(
3725                            "INSERT INTO msgs(rowid, text) VALUES (?1, ?2)",
3726                            rusqlite::params![rowid, text],
3727                        )
3728                        .unwrap();
3729                }
3730            }
3731            connection
3732        }
3733
3734        fn live(ids: &[&str]) -> BTreeSet<String> {
3735            ids.iter().map(|id| (*id).to_owned()).collect()
3736        }
3737
3738        fn matches(kinds: &[(&str, SessionTextMatchKind)]) -> Vec<SessionTextMatch> {
3739            kinds
3740                .iter()
3741                .map(|(id, kind)| SessionTextMatch {
3742                    session_id: (*id).to_owned(),
3743                    kind: *kind,
3744                })
3745                .collect()
3746        }
3747
3748        #[test]
3749        fn user_and_agent_messages_match_and_tool_output_does_not() {
3750            let connection = index(&[
3751                ("said-by-user", &[("user", "please fix the Zebra crossing")]),
3752                ("said-by-agent", &[("assistant", "the zebra is fixed")]),
3753                (
3754                    "only-in-tool",
3755                    &[("tool", "zebra stack trace"), ("user", "hello")],
3756                ),
3757                (
3758                    "both",
3759                    &[("assistant", "a ZEBRA appears"), ("user", "a zebra please")],
3760                ),
3761                ("gone", &[("user", "zebra")]),
3762            ]);
3763            let live = live(&["said-by-user", "said-by-agent", "only-in-tool", "both"]);
3764            assert_eq!(
3765                text_matches_in(&connection, "zebra", &live).unwrap(),
3766                matches(&[
3767                    ("both", SessionTextMatchKind::User),
3768                    ("said-by-agent", SessionTextMatchKind::Agent),
3769                    ("said-by-user", SessionTextMatchKind::User),
3770                ])
3771            );
3772        }
3773
3774        #[test]
3775        fn a_query_too_short_for_the_trigram_index_still_matches() {
3776            let connection = index(&[
3777                ("a", &[("user", "go to the zoo")]),
3778                ("b", &[("assistant", "zoo")]),
3779                ("c", &[("tool", "zoo")]),
3780            ]);
3781            assert_eq!(
3782                text_matches_in(&connection, "zo", &live(&["a", "b", "c"])).unwrap(),
3783                matches(&[
3784                    ("a", SessionTextMatchKind::User),
3785                    ("b", SessionTextMatchKind::Agent),
3786                ])
3787            );
3788            // A percent sign is text, not a wildcard.
3789            assert_eq!(
3790                text_matches_in(&connection, "%z", &live(&["a"])).unwrap(),
3791                Vec::new()
3792            );
3793        }
3794    }
3795}