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