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