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