Skip to main content

mj_controller/
sessionwiki.rs

1//! Publishing Mjolnir's own checkpointed sessions into the user's SessionWiki
2//! index, and the daemon-side job that keeps that index current.
3//!
4//! SessionWiki keeps one searchable index of AI coding sessions across every
5//! tool a user runs. Mjolnir links it as a library and registers
6//! [`MjolnirAdapter`] beside SessionWiki's built-in adapters, so a Mjolnir
7//! session is searchable next to a Claude Code or Codex one. The adapter is a
8//! "shared store" adapter: checkpoints are not one-file-per-session in a shape
9//! SessionWiki can parse, so the indexer enumerates sessions by key and asks
10//! this adapter to parse the ones whose checkpoint changed.
11
12mod harness_adapters;
13pub(crate) mod history;
14mod provenance;
15pub mod tags;
16
17use std::collections::{BTreeMap, BTreeSet};
18use std::path::{Path, PathBuf};
19use std::sync::Arc;
20use std::sync::atomic::{AtomicBool, Ordering};
21use std::time::Instant;
22
23use anyhow::{Context, Result};
24use chrono::{DateTime, Utc};
25
26use mj_client::daemon::{
27    WikiHitBlock, WikiHitTranscript, WikiIndexState, WikiRow, WikiSessionInfo, WikiSessionStatus,
28    WikiStatus,
29};
30use mj_core::config::HarnessKind;
31use mj_core::state::{SessionRecord, State};
32use sessionwiki::adapters::{Adapter, Discovered, Store};
33use sessionwiki::model::{Message, Role, Session};
34
35use crate::controller::Controller;
36use crate::controller::checkpoint::managed_checkpoint_archive_name;
37use harness_adapters::HarnessAdapter;
38
39/// The tool name every Mjolnir instance publishes under. One name means one
40/// search partition; reconciliation is scoped per instance instead (see
41/// [`Adapter::reconcile_scope`]).
42const TOOL: &str = "mjolnir";
43
44/// One checkpoint archive on disk, reduced to what indexing needs.
45struct ArchiveFile {
46    path: PathBuf,
47    frontier: u64,
48    /// Modification time in epoch seconds, SessionWiki's change token.
49    token: i64,
50}
51
52/// What indexing needs from controller state: which sessions exist, which of
53/// them are sub-agent children, and which are still running with their
54/// conversation in the daemon's own database rather than in a checkpoint.
55#[derive(Default)]
56struct Sessions {
57    records: BTreeMap<String, SessionRecord>,
58    subagent_ids: BTreeSet<String>,
59    /// Session id to change token, for sessions indexed from the projection.
60    live: BTreeMap<String, i64>,
61}
62
63impl Sessions {
64    fn of(state: &State) -> Self {
65        Self {
66            records: state.sessions.clone(),
67            subagent_ids: state.subagents.keys().cloned().collect(),
68            live: live_tokens(state),
69        }
70    }
71}
72
73/// The change token of every session whose transcript is still only in the
74/// daemon's database: its activity watermark in whole seconds.
75///
76/// A stopped session keeps being indexed from its checkpoint, which never
77/// changes again. Everything else is indexed from the projection, so a running
78/// session is findable before it has ever been closed.
79fn live_tokens(state: &State) -> BTreeMap<String, i64> {
80    let activity = match crate::database::load_transcribed_session_activity() {
81        Ok(activity) => activity,
82        Err(error) => {
83            tracing::warn!(%error, "could not read session activity for SessionWiki");
84            return BTreeMap::new();
85        }
86    };
87    state
88        .sessions
89        .iter()
90        .filter(|(_, record)| record.state != mj_core::state::SessionState::Stopped)
91        .filter_map(|(session_id, _)| {
92            let watermark = activity.get(session_id)?;
93            Some((session_id.clone(), watermark.unwrap_or_default() / 1000))
94        })
95        .collect()
96}
97
98/// Mjolnir's sessions, as SessionWiki sees them.
99pub struct MjolnirAdapter {
100    sessions_dir: PathBuf,
101    sessions: std::sync::Mutex<Sessions>,
102    /// Re-read controller state when the indexer reaches this adapter.
103    reload: bool,
104}
105
106impl MjolnirAdapter {
107    /// A fixed view of the given state, which is what a caller with a state in
108    /// hand wants.
109    pub fn from_state(state: &State) -> Self {
110        Self {
111            sessions_dir: mj_core::config::sessions_dir(),
112            sessions: std::sync::Mutex::new(Sessions::of(state)),
113            reload: false,
114        }
115    }
116
117    /// The same, but re-reading controller state when the indexer reaches this
118    /// adapter.
119    ///
120    /// One sync pass walks every other tool's store first, which can take
121    /// minutes on a large corpus. Without the reload, sessions that closed
122    /// during that walk would be indexed with no record: no project, no start
123    /// time, and the title guessed from the first prompt. Their checkpoints do
124    /// not change afterwards, so nothing would ever correct them.
125    pub fn reloading(state: &State) -> Self {
126        Self {
127            reload: true,
128            ..Self::from_state(state)
129        }
130    }
131
132    /// Mjolnir's own metadata for every session this adapter knows about, to
133    /// be stored in the index beside the transcripts.
134    ///
135    /// Read from the adapter's own snapshot rather than from the controller
136    /// state the sync loaded, because [`MjolnirAdapter::reloading`] replaces
137    /// that snapshot when the indexer reaches this adapter. A session that
138    /// closed during a long first pass is indexed from the reloaded state, so
139    /// its metadata has to come from the same state that produced its row.
140    pub fn indexed_tags(&self) -> BTreeMap<String, tags::MjTags> {
141        let sessions = self
142            .sessions
143            .lock()
144            .unwrap_or_else(std::sync::PoisonError::into_inner);
145        sessions
146            .records
147            .iter()
148            .map(|(session_id, record)| {
149                (
150                    session_id.clone(),
151                    tags::MjTags {
152                        target: Some(record.target_template_id.clone()).filter(|id| !id.is_empty()),
153                        profile: Some(record.last_profile.clone()).filter(|id| !id.is_empty()),
154                        harness: Some(record.harness_kind.id().to_owned()),
155                    },
156                )
157            })
158            .collect()
159    }
160
161    fn reload(&self) {
162        if !self.reload {
163            return;
164        }
165        match Controller::load() {
166            Ok(controller) => {
167                *self
168                    .sessions
169                    .lock()
170                    .unwrap_or_else(std::sync::PoisonError::into_inner) =
171                    Sessions::of(&controller.state)
172            }
173            Err(error) => {
174                tracing::warn!(%error, "could not refresh session records for SessionWiki")
175            }
176        }
177    }
178
179    /// The conversation of a stopped session, read from its newest checkpoint,
180    /// with the title the checkpoint recorded.
181    fn checkpointed_transcript(&self, session_id: &str) -> Result<IndexedTranscript> {
182        let (newest, _) = self.newest_archives();
183        let archive = newest
184            .get(session_id)
185            .with_context(|| format!("no checkpoint archive for session {session_id}"))?;
186        let snapshot = mj_checkpoint::archive::read_archive_verified(&archive.path)
187            .with_context(|| format!("read checkpoint {}", archive.path.display()))?
188            .canonical_session()
189            .with_context(|| format!("read the transcript of session {session_id}"))?;
190        let mut evidence = provenance::Evidence::default();
191        for item in &snapshot.transcript {
192            if let mj_core::archive::CanonicalTranscriptBody::Tool { call, .. } = &item.body {
193                evidence.observe(call, item.created_at_ms);
194            }
195        }
196        let messages = summary_messages(mj_transcript::summary::TranscriptSummary::from_snapshot(
197            &snapshot,
198        ));
199        Ok(IndexedTranscript {
200            messages,
201            title: snapshot.session.session_title.clone(),
202            evidence,
203        })
204    }
205
206    /// The conversation of a session that has not stopped, read from the
207    /// daemon's own projection. It is the same conversation the checkpoint
208    /// would hold, minus whatever has not happened yet.
209    fn projected_transcript(&self, session_id: &str) -> Result<IndexedTranscript> {
210        let projection = crate::database::load_materialized_session(session_id)
211            .with_context(|| format!("read the stored transcript of session {session_id}"))?
212            .with_context(|| format!("no stored transcript for session {session_id}"))?;
213        let mut evidence = provenance::Evidence::default();
214        for item in &projection.transcript {
215            if let mj_core::state::TranscriptBody::Tool { call, .. } = &item.body {
216                evidence.observe(call, item.created_at_ms);
217            }
218        }
219        Ok(IndexedTranscript {
220            messages: projected_messages(&projection),
221            title: projection.session_title.clone(),
222            evidence,
223        })
224    }
225
226    /// The stable key for one session: its checkpoint directory and id. The
227    /// directory is per instance, which is what scopes reconciliation.
228    fn key_for(&self, session_id: &str) -> String {
229        format!("{}/{session_id}", self.sessions_dir.display())
230    }
231
232    /// The newest checkpoint of every session in the directory, by session id.
233    ///
234    /// `had_error` is true when the directory exists but could not be read in
235    /// full; the indexer then skips deletion reconciliation rather than
236    /// archiving every Mjolnir session off a partial listing.
237    fn newest_archives(&self) -> (BTreeMap<String, ArchiveFile>, bool) {
238        let mut newest: BTreeMap<String, ArchiveFile> = BTreeMap::new();
239        let mut had_error = false;
240        let entries = match std::fs::read_dir(&self.sessions_dir) {
241            Ok(entries) => entries,
242            Err(error) => {
243                if self.sessions_dir.exists() {
244                    tracing::debug!(
245                        directory = %self.sessions_dir.display(),
246                        %error,
247                        "could not list the checkpoint directory for SessionWiki"
248                    );
249                    had_error = true;
250                }
251                return (newest, had_error);
252            }
253        };
254        for entry in entries {
255            let Ok(entry) = entry else {
256                had_error = true;
257                continue;
258            };
259            let Some((session_id, frontier)) = checkpoint_archive_session(&entry.file_name())
260            else {
261                continue;
262            };
263            let token = entry
264                .metadata()
265                .ok()
266                .and_then(|metadata| metadata.modified().ok())
267                .and_then(|modified| modified.duration_since(std::time::UNIX_EPOCH).ok())
268                .map(|age| age.as_secs() as i64)
269                .unwrap_or(0);
270            let candidate = ArchiveFile {
271                path: entry.path(),
272                frontier,
273                token,
274            };
275            match newest.get(&session_id) {
276                Some(existing) if existing.frontier >= candidate.frontier => {}
277                _ => {
278                    newest.insert(session_id, candidate);
279                }
280            }
281        }
282        (newest, had_error)
283    }
284}
285
286struct IndexedTranscript {
287    messages: Vec<Message>,
288    title: Option<String>,
289    evidence: provenance::Evidence,
290}
291
292/// The session a checkpoint file name belongs to, with its generation.
293///
294/// Managed checkpoints carry a frontier and a nonce; an imported archive is
295/// named for its session alone and counts as generation zero.
296fn checkpoint_archive_session(name: &std::ffi::OsStr) -> Option<(String, u64)> {
297    if let Some(parsed) = managed_checkpoint_archive_name(name) {
298        return Some((parsed.session_id, parsed.frontier));
299    }
300    let stem = name
301        .to_str()
302        .and_then(|name| name.strip_suffix(".hel.zip"))?;
303    mj_core::config::validate_id("session", stem)
304        .is_ok()
305        .then(|| (stem.to_owned(), 0))
306}
307
308/// A running session's conversation, as SessionWiki stores it.
309fn projected_messages(projection: &mj_core::state::MaterializedSession) -> Vec<Message> {
310    summary_messages(mj_transcript::summary::TranscriptSummary::from_materialized(projection))
311}
312
313fn summary_messages(summary: mj_transcript::summary::TranscriptSummary) -> Vec<Message> {
314    use mj_transcript::summary::SummaryRole;
315    summary
316        .entries
317        .into_iter()
318        .filter_map(|entry| {
319            let role = match entry.role {
320                SummaryRole::User => Role::User,
321                SummaryRole::Assistant => Role::Assistant,
322                SummaryRole::Tool => Role::Tool,
323                SummaryRole::Plan => return None,
324            };
325            message(role, entry.body(), entry.created_at_ms)
326        })
327        .collect()
328}
329
330/// One indexed message, or nothing when the item carried no text.
331fn message(role: Role, text: String, created_at_ms: i64) -> Option<Message> {
332    let text = text.trim().to_owned();
333    (!text.is_empty()).then(|| Message {
334        role,
335        text,
336        ts: DateTime::from_timestamp_millis(created_at_ms),
337    })
338}
339
340fn parse_time(value: &str) -> Option<DateTime<Utc>> {
341    DateTime::parse_from_rfc3339(value)
342        .ok()
343        .map(|time| time.with_timezone(&Utc))
344}
345
346impl Adapter for MjolnirAdapter {
347    fn name(&self) -> &'static str {
348        TOOL
349    }
350
351    fn root(&self) -> Option<PathBuf> {
352        Some(self.sessions_dir.clone())
353    }
354
355    /// Unused: this is a shared-store adapter, so the indexer enumerates
356    /// sessions through [`Adapter::store`] instead of walking files.
357    fn discover(&self) -> Discovered {
358        Discovered {
359            files: Vec::new(),
360            had_error: false,
361        }
362    }
363
364    fn parse(&self, _path: &Path) -> Result<Session> {
365        anyhow::bail!("Mjolnir sessions are parsed by key, not by file")
366    }
367
368    fn store(&self) -> Option<Store> {
369        self.reload();
370        let (newest, had_error) = self.newest_archives();
371        let mut files = Vec::with_capacity(newest.len());
372        let mut tokens: BTreeMap<String, i64> = BTreeMap::new();
373        for (session_id, archive) in newest {
374            tokens.insert(session_id, archive.token);
375            files.push(archive.path);
376        }
377        // A session that is still running is indexed from the projection, and
378        // its own token replaces any checkpoint token it has: the conversation
379        // has moved on since that checkpoint was written. Listing it also
380        // keeps reconciliation from archiving a running session.
381        let sessions = self
382            .sessions
383            .lock()
384            .unwrap_or_else(std::sync::PoisonError::into_inner);
385        let live = sessions.live.clone();
386        tokens.extend(live);
387        // A rename changes the record and not the conversation, so the
388        // record's own last update is part of the change token. Without it a
389        // renamed session would keep its old title in the index for as long as
390        // its transcript stood still.
391        for (session_id, token) in tokens.iter_mut() {
392            let updated = sessions
393                .records
394                .get(session_id)
395                .and_then(|record| parse_time(&record.updated_at))
396                .map(|updated| updated.timestamp());
397            if let Some(updated) = updated {
398                *token = (*token).max(updated);
399            }
400        }
401        let keys = tokens
402            .into_iter()
403            .map(|(session_id, token)| {
404                (
405                    self.key_for(&session_id),
406                    token
407                        .saturating_mul(1024)
408                        .saturating_add(i64::from(mj_transcript::summary::SUMMARY_VERSION)),
409                )
410            })
411            .collect();
412        Some(Store {
413            keys,
414            files,
415            had_error,
416        })
417    }
418
419    /// Every Mjolnir instance publishes under one tool name, so this instance
420    /// speaks only for keys under its own checkpoint directory. Without the
421    /// scope, two instances would archive each other's rows on every sync.
422    fn reconcile_scope(&self) -> Option<String> {
423        Some(format!("{}/", self.sessions_dir.display()))
424    }
425
426    fn parse_key(&self, key: &str) -> Result<Session> {
427        let session_id = key.rsplit('/').next().unwrap_or_default();
428        anyhow::ensure!(!session_id.is_empty(), "no session id in key {key:?}");
429        let sessions = self
430            .sessions
431            .lock()
432            .unwrap_or_else(std::sync::PoisonError::into_inner);
433        let IndexedTranscript {
434            messages,
435            title: snapshot_title,
436            evidence,
437        } = if sessions.live.contains_key(session_id) {
438            self.projected_transcript(session_id)?
439        } else {
440            self.checkpointed_transcript(session_id)?
441        };
442        let record = sessions.records.get(session_id);
443
444        let title = record
445            .and_then(|record| record.session_title_override.clone())
446            .or_else(|| record.and_then(|record| record.acp_session_title.clone()))
447            .or_else(|| snapshot_title.clone())
448            .unwrap_or_else(|| {
449                messages
450                    .iter()
451                    .find(|message| message.role == Role::User)
452                    .map(|message| message.text.chars().take(80).collect())
453                    .unwrap_or_default()
454            });
455
456        Ok(Session {
457            id: session_id.to_owned(),
458            tool: TOOL,
459            path: PathBuf::from(key),
460            project: record
461                .and_then(|record| record.project_directory.as_ref())
462                .map(|directory| directory.display().to_string())
463                .unwrap_or_default(),
464            started: record.and_then(|record| parse_time(&record.created_at)),
465            ended: record.and_then(|record| parse_time(&record.updated_at)),
466            title,
467            subagent: sessions.subagent_ids.contains(session_id),
468            messages,
469            touched: evidence.paths.into_iter().collect(),
470            edits: evidence.edits,
471        })
472    }
473}
474
475/// The Mjolnir adapter handed to the indexer while the sync keeps its own
476/// handle on it.
477///
478/// The indexer takes `Box<dyn Adapter>` and consumes the list, but the sync has
479/// to ask the same adapter for its final session snapshot once the walk is over
480/// (see [`MjolnirAdapter::indexed_tags`]). Sharing the adapter is the only way
481/// both can hold it.
482struct SharedMjolnirAdapter(Arc<MjolnirAdapter>);
483
484impl Adapter for SharedMjolnirAdapter {
485    fn name(&self) -> &'static str {
486        self.0.name()
487    }
488
489    fn root(&self) -> Option<PathBuf> {
490        self.0.root()
491    }
492
493    fn discover(&self) -> Discovered {
494        self.0.discover()
495    }
496
497    fn parse(&self, path: &Path) -> Result<Session> {
498        self.0.parse(path)
499    }
500
501    fn store(&self) -> Option<Store> {
502        self.0.store()
503    }
504
505    fn parse_key(&self, key: &str) -> Result<Session> {
506        self.0.parse_key(key)
507    }
508
509    fn reconcile_scope(&self) -> Option<String> {
510        self.0.reconcile_scope()
511    }
512}
513
514/// The daemon's SessionWiki sync job.
515///
516/// Triggers coalesce: a request while a sync is running marks a rerun instead
517/// of queueing a second one, so a burst of closing sessions costs one extra
518/// pass. Syncs are single-flight because SessionWiki holds a write transaction
519/// per adapter batch, and two writers only produce a busy error.
520pub struct WikiIndexer {
521    inner: Arc<Indexer>,
522}
523
524#[derive(Default)]
525struct Indexer {
526    /// Held for the whole of one run: this is what makes syncs single-flight.
527    running: tokio::sync::Mutex<()>,
528    notify: tokio::sync::Notify,
529    /// A trigger arrived; the worker has not consumed it yet.
530    requested: AtomicBool,
531    /// At least one waiting trigger asked for a full sync.
532    full_requested: AtomicBool,
533    /// A sync pass is running now. A surface shows this as "topping up", so a
534    /// user knows more results may arrive.
535    in_flight: AtomicBool,
536    last_success: std::sync::Mutex<Option<Success>>,
537}
538
539#[derive(Clone, Copy)]
540struct Success {
541    at: Instant,
542    epoch_seconds: i64,
543}
544
545impl WikiIndexer {
546    /// Start the background sync worker. Without a Tokio runtime (some tests
547    /// build a runtime state without one) the indexer stays inert.
548    pub fn spawn() -> Self {
549        let inner = Arc::new(Indexer::default());
550        if let Ok(handle) = tokio::runtime::Handle::try_current() {
551            let worker = Arc::clone(&inner);
552            handle.spawn(async move { worker.run().await });
553        }
554        Self { inner }
555    }
556
557    /// Ask for a sync. Returns immediately; the work happens in the background.
558    pub fn request_sync(&self, full: bool) {
559        if full {
560            self.inner.full_requested.store(true, Ordering::Release);
561        }
562        self.inner.requested.store(true, Ordering::Release);
563        self.inner.notify.notify_one();
564    }
565
566    /// Run a sync and wait for it, joining a sync already in flight.
567    pub async fn sync_now(&self, full: bool) -> Result<()> {
568        self.inner.sync(full).await
569    }
570
571    /// The state of the index and whether a sync is running, for the surfaces
572    /// that say so while the first build is under way.
573    pub fn status(&self) -> WikiStatus {
574        WikiStatus {
575            state: index_state(),
576            topping_up: self.inner.in_flight.load(Ordering::Acquire)
577                || self.inner.requested.load(Ordering::Acquire),
578        }
579    }
580
581    /// When the last sync succeeded, for callers that trigger on staleness.
582    pub fn last_success(&self) -> Option<Instant> {
583        self.inner
584            .last_success
585            .lock()
586            .unwrap_or_else(std::sync::PoisonError::into_inner)
587            .map(|success| success.at)
588    }
589}
590
591impl Indexer {
592    async fn run(self: Arc<Self>) {
593        loop {
594            self.notify.notified().await;
595            while self.requested.swap(false, Ordering::AcqRel) {
596                let full = self.full_requested.swap(false, Ordering::AcqRel);
597                if let Err(error) = self.sync(full).await {
598                    self.report(&error);
599                    // A failure waits for the next trigger rather than
600                    // retrying straight away: a busy index stays busy for as
601                    // long as the other writer holds it, and a spin would only
602                    // add to the contention.
603                    break;
604                }
605            }
606        }
607    }
608
609    /// Log a failed sync at the level its cause deserves. A busy index is an
610    /// expected collision with another writer, not a fault: mark a rerun and
611    /// say so only in debug output.
612    fn report(&self, error: &anyhow::Error) {
613        if crate::database::is_busy_error(error) {
614            self.requested.store(true, Ordering::Release);
615            tracing::debug!(%error, "the SessionWiki index was busy; retrying on the next trigger");
616        } else {
617            tracing::warn!(%error, "could not sync sessions into SessionWiki");
618        }
619    }
620
621    async fn sync(&self, full: bool) -> Result<()> {
622        let _guard = self.running.lock().await;
623        let since = if full {
624            None
625        } else {
626            self.last_success
627                .lock()
628                .unwrap_or_else(std::sync::PoisonError::into_inner)
629                // A minute of overlap covers checkpoints written while the
630                // previous run was reading the directory.
631                .map(|success| success.epoch_seconds - 60)
632        };
633        let started = Instant::now();
634        let work = crate::upgrade::activity("SessionWiki sync")?;
635        self.in_flight.store(true, Ordering::Release);
636        let ran = tokio::task::spawn_blocking(move || {
637            let _work = work;
638            sync_blocking(since)
639        })
640        .await;
641        self.in_flight.store(false, Ordering::Release);
642        let ran = ran.context("run the SessionWiki sync")??;
643        if ran {
644            *self
645                .last_success
646                .lock()
647                .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(Success {
648                at: started,
649                epoch_seconds: Utc::now().timestamp(),
650            });
651        }
652        Ok(())
653    }
654}
655
656/// One synchronous sync pass. Returns false when this process must not touch
657/// the index, so a refused run never records a success it did not have.
658fn sync_blocking(since: Option<i64>) -> Result<bool> {
659    if !index_is_writable() {
660        return Ok(false);
661    }
662    let controller =
663        Controller::load().context("load controller state for the SessionWiki sync")?;
664    // Mjolnir's own sessions go first: a cold index walks every other tool's
665    // store for many minutes, and a just-closed session should not wait on it.
666    let mjolnir = Arc::new(MjolnirAdapter::reloading(&controller.state));
667    let mut adapters: Vec<Box<dyn sessionwiki::adapters::Adapter>> =
668        vec![Box::new(SharedMjolnirAdapter(Arc::clone(&mjolnir)))];
669    adapters.extend(native_adapters(&controller.config));
670    let mut connection = sessionwiki::index::open().context("open the SessionWiki index")?;
671    sessionwiki::index::sync_with(&mut connection, &adapters, since)
672        .context("sync the SessionWiki index")?;
673    write_session_tags(&mut connection, &mjolnir.indexed_tags())
674        .context("store Mjolnir's session metadata in the SessionWiki index")?;
675    provenance::backfill(&mut connection, &mjolnir).context("backfill Mjolnir file provenance")?;
676    if since.is_none() {
677        // A full pass has walked every store, so the index is complete enough
678        // for a search to be trusted. The marker is what a later daemon reads
679        // instead of walking the corpus again to find out.
680        record_first_build();
681    }
682    Ok(true)
683}
684
685/// Store each session's target, profile and harness in the index, in one
686/// transaction.
687///
688/// Every session is written on every sync rather than only the changed ones:
689/// the write is a delete and three inserts, which is nothing beside the
690/// transcript indexing in the same pass, and it is what makes a Move or a
691/// profile switch show up without tracking which records changed. It is also
692/// what gives sessions indexed before this existed their metadata, with no
693/// migration and no re-index.
694fn write_session_tags(
695    connection: &mut rusqlite::Connection,
696    session_tags: &BTreeMap<String, tags::MjTags>,
697) -> Result<()> {
698    if session_tags.is_empty() {
699        return Ok(());
700    }
701    let transaction = connection
702        .transaction()
703        .context("open a transaction for the session metadata")?;
704    for (session_id, session) in session_tags {
705        if session.is_empty() {
706            continue;
707        }
708        tags::write(&transaction, session_id, session)?;
709    }
710    transaction
711        .commit()
712        .context("commit the session metadata")?;
713    Ok(())
714}
715
716/// The non-Mjolnir adapters this install indexes.
717///
718/// Mjolnir's configured harness profiles decide which harness homes are
719/// indexed, not the stock `~/.codex` and `~/.claude` locations. A user who
720/// runs several profile homes expects every session Mjolnir can start to be
721/// searchable, and a home no profile names is not Mjolnir's to walk. So the
722/// stock Codex and Claude adapters are dropped and one adapter per enabled
723/// profile home takes their place; every other built-in adapter is kept as is.
724///
725/// Kimi Code, Grok Build and Muse have no SessionWiki adapter at all, so
726/// Mjolnir supplies one per enabled profile home of its own (see
727/// [`harness_adapters`]). Without them those sessions would never appear in
728/// the Resume dialog's search.
729///
730/// Each per-home adapter reports a reconcile scope covering only its own root,
731/// so a sync of one install never archives the rows of another.
732fn native_adapters(config: &mj_core::config::Config) -> Vec<Box<dyn Adapter>> {
733    // Two profiles may share one home, and two harnesses may share one home
734    // path without sharing sessions, so the kind is part of the identity.
735    let mut seen: BTreeSet<(HarnessKind, &Path)> = BTreeSet::new();
736    let mut adapters: Vec<Box<dyn Adapter>> = Vec::new();
737    for (_, profile) in config.enabled_profiles() {
738        // A second adapter for the same home would only walk it twice.
739        if !seen.insert((profile.kind, profile.home.as_path())) {
740            continue;
741        }
742        let adapter: Box<dyn Adapter> = match profile.kind {
743            HarnessKind::Codex => {
744                Box::new(sessionwiki::adapters::Codex::in_home(profile.home.clone()))
745            }
746            HarnessKind::Claude => Box::new(sessionwiki::adapters::ClaudeCode::in_home(
747                profile.home.clone(),
748            )),
749            kind => match HarnessAdapter::in_home(kind, profile.home.clone()) {
750                Some(adapter) => Box::new(adapter),
751                None => continue,
752            },
753        };
754        adapters.push(adapter);
755    }
756    adapters.extend(
757        sessionwiki::adapters::all()
758            .into_iter()
759            .filter(|adapter| !matches!(adapter.name(), "codex" | "claude-code")),
760    );
761    adapters
762}
763
764// ---------------------------------------------------------------------------
765// Which index, and whether it may be touched
766// ---------------------------------------------------------------------------
767
768/// Whether this process may open the index at all.
769///
770/// Indexing is always on, so a process that never resolved where its index
771/// belongs must not reach for one: it would walk the user's real session
772/// stores and write the user's real index. Only Mjolnir's own startup resolves
773/// it (see `mj_core::config::apply_instance_flag`), so this refuses every unit
774/// test that builds a daemon runtime directly and every other embedder, unless
775/// it names an index of its own with `SESSIONWIKI_DATA`.
776fn index_is_isolated() -> bool {
777    static SAID: AtomicBool = AtomicBool::new(false);
778    if mj_core::config::session_index_is_resolved()
779        || std::env::var_os(mj_core::config::SESSION_INDEX_ENV).is_some()
780    {
781        return true;
782    }
783    if !SAID.swap(true, Ordering::AcqRel) {
784        tracing::debug!(
785            "this process did not resolve a session index location; SessionWiki is not used"
786        );
787    }
788    false
789}
790
791/// Whether the index on disk was written by a SessionWiki at another schema
792/// version.
793///
794/// SessionWiki's own `open` drops and rebuilds its whole cache when the file's
795/// `user_version` differs from the version it was built with, which on a large
796/// corpus costs tens of minutes. Mjolnir will not do that to a user who also
797/// runs the `sessionwiki` command: it reads the version without SessionWiki and
798/// stands aside.
799fn index_version_mismatch() -> bool {
800    static SAID: AtomicBool = AtomicBool::new(false);
801    let Ok(path) = sessionwiki::index::db_path() else {
802        return false;
803    };
804    if !path.exists() {
805        return false;
806    }
807    let version = rusqlite::Connection::open_with_flags(
808        &path,
809        rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_URI,
810    )
811    .and_then(|connection| connection.pragma_query_value(None, "user_version", |row| row.get(0)));
812    let version: i64 = match version {
813        Ok(version) => version,
814        Err(error) => {
815            tracing::debug!(%error, "could not read the SessionWiki index schema version");
816            return false;
817        }
818    };
819    // Zero is an index SessionWiki has not finished creating; it is not a
820    // different version.
821    let mismatch = version != 0 && version != sessionwiki::index::SCHEMA_VERSION;
822    if mismatch && !SAID.swap(true, Ordering::AcqRel) {
823        tracing::warn!(
824            found = version,
825            expected = sessionwiki::index::SCHEMA_VERSION,
826            path = %path.display(),
827            "the SessionWiki index was written by another version;              Mjolnir will not open it, because opening it would rebuild it.              Install the matching sessionwiki command"
828        );
829    }
830    mismatch
831}
832
833fn index_is_writable() -> bool {
834    index_is_isolated() && !index_version_mismatch()
835}
836
837/// The file recording that one full sync has completed, holding the schema
838/// version it completed at.
839fn first_build_marker() -> PathBuf {
840    mj_core::config::data_dir().join("sessionwiki-built")
841}
842
843fn record_first_build() {
844    let path = first_build_marker();
845    let version = sessionwiki::index::SCHEMA_VERSION.to_string();
846    if std::fs::read_to_string(&path).is_ok_and(|held| held.trim() == version) {
847        return;
848    }
849    if let Err(error) = std::fs::write(&path, &version) {
850        tracing::warn!(%error, path = %path.display(), "could not record the first SessionWiki build");
851    }
852}
853
854/// Whether this index has completed a full build at this schema version.
855fn first_build_is_done() -> bool {
856    std::fs::read_to_string(first_build_marker())
857        .is_ok_and(|held| held.trim() == sessionwiki::index::SCHEMA_VERSION.to_string())
858        && sessionwiki::index::db_path().is_ok_and(|path| path.exists())
859}
860
861/// What a surface should say about this index right now.
862pub fn index_state() -> WikiIndexState {
863    if !index_is_isolated() {
864        return WikiIndexState::Indexing;
865    }
866    if index_version_mismatch() {
867        return WikiIndexState::VersionMismatch;
868    }
869    if first_build_is_done() {
870        WikiIndexState::Ready
871    } else {
872        WikiIndexState::Indexing
873    }
874}
875
876// ---------------------------------------------------------------------------
877// Queries and restore
878// ---------------------------------------------------------------------------
879
880/// The largest page a caller may ask a wiki query for.
881pub const MAX_WIKI_LIMIT: usize = 200;
882/// The page size a caller that names none gets.
883pub const DEFAULT_WIKI_LIMIT: usize = 50;
884/// SessionWiki's full-text index needs three characters; shorter queries fall
885/// back to a substring scan.
886const MIN_FULLTEXT_QUERY: usize = 3;
887/// How stale the index may be before a query triggers a background sync.
888pub const SYNC_STALE_AFTER: std::time::Duration = std::time::Duration::from_secs(60);
889
890/// Whether a query should trigger a bounded background sync before it answers.
891pub fn sync_is_stale(last_success: Option<Instant>) -> bool {
892    last_success.is_none_or(|at| at.elapsed() >= SYNC_STALE_AFTER)
893}
894
895/// One page of the index, newest first or best match first.
896///
897/// `live` is the set of session ids this daemon still holds, which is what
898/// decides whether a Mjolnir row names a session the user can simply resume.
899/// `include_subagents` decides, for every path below, whether sub-agent
900/// sessions are answered at all; a resume list never wants them.
901/// Runs SQLite work, so callers on the async runtime wrap it in
902/// `spawn_blocking`.
903pub fn query_rows(
904    query: &str,
905    limit: usize,
906    live: &BTreeSet<String>,
907    include_subagents: bool,
908) -> Result<Vec<WikiRow>> {
909    let limit = limit.clamp(1, MAX_WIKI_LIMIT);
910    if !index_is_writable() {
911        // Nothing to answer from: either this process has no index of its own
912        // or the one on disk is at another version. The status beside the rows
913        // says which.
914        return Ok(Vec::new());
915    }
916    let connection = open_readonly()?;
917    let query = query.trim();
918    if query.is_empty() {
919        let rows =
920            sessionwiki::index::recent(&connection, limit, None, None, None, include_subagents)
921                .context("list recent SessionWiki sessions")?;
922        let mut rows: Vec<WikiRow> = rows
923            .into_iter()
924            .map(|row| wiki_row(row, None, live))
925            .collect();
926        fill_session_tags(&connection, &mut rows)?;
927        return Ok(rows);
928    }
929    let hits = if query.chars().count() < MIN_FULLTEXT_QUERY {
930        sessionwiki::index::search_like(&connection, query, limit, None, None)
931    } else {
932        sessionwiki::index::search(&connection, query, limit, None, None)
933    }
934    .context("search the SessionWiki index")?;
935    // SessionWiki's full-text search has no sub-agent filter of its own.
936    let mut rows: Vec<WikiRow> = hits
937        .into_iter()
938        .filter(|hit| include_subagents || is_main_session(&hit.row))
939        .map(|hit| wiki_row(hit.row, Some(hit.snippet), live))
940        .collect();
941    // SessionWiki searches message text alone, so a session known by a title
942    // or a project that is never said out loud would be unfindable. Those
943    // matches follow the full-text ones rather than displacing them.
944    let found: BTreeSet<String> = rows.iter().map(|row| row.id.clone()).collect();
945    for row in named_like(&connection, query, include_subagents)? {
946        if rows.len() >= limit {
947            break;
948        }
949        if found.contains(&row.session_id) {
950            continue;
951        }
952        rows.push(wiki_row(row, None, live));
953    }
954    fill_session_tags(&connection, &mut rows)?;
955    Ok(rows)
956}
957
958/// Fill in the target, profile and harness of every Mjolnir row on this page
959/// from the index's own tags, in one query.
960///
961/// Only Mjolnir writes those tags, so a row from another tool keeps `None` and
962/// is not even asked about.
963fn fill_session_tags(connection: &rusqlite::Connection, rows: &mut [WikiRow]) -> Result<()> {
964    let ids: Vec<&str> = rows
965        .iter()
966        .filter(|row| row.tool == TOOL)
967        .map(|row| row.id.as_str())
968        .collect();
969    let found = tags::read(connection, &ids).context("read the indexed session metadata")?;
970    for row in rows.iter_mut().filter(|row| row.tool == TOOL) {
971        let Some(session) = found.get(&row.id) else {
972            continue;
973        };
974        row.target = session.target.clone();
975        row.profile = session.profile.clone();
976        row.harness = session.harness.clone();
977    }
978    Ok(())
979}
980
981/// How far back a title or project match looks. Those columns have no index of
982/// their own, so this is a scan of the most recent sessions rather than of the
983/// whole corpus.
984const NAME_SCAN_LIMIT: usize = 2_000;
985
986/// Indexed sessions whose title or project contains the query, ignoring case.
987fn named_like(
988    connection: &rusqlite::Connection,
989    query: &str,
990    include_subagents: bool,
991) -> Result<Vec<sessionwiki::index::SessionRow>> {
992    let needle = query.to_lowercase();
993    let rows = sessionwiki::index::recent(
994        connection,
995        NAME_SCAN_LIMIT,
996        None,
997        None,
998        None,
999        include_subagents,
1000    )
1001    .context("list recent SessionWiki sessions")?;
1002    Ok(rows
1003        .into_iter()
1004        .filter(|row| {
1005            row.title.to_lowercase().contains(&needle)
1006                || row.project.to_lowercase().contains(&needle)
1007        })
1008        .collect())
1009}
1010
1011/// Whether an indexed session is one a person started rather than a
1012/// sub-agent. SessionWiki's own `recent` filter tests the same `kind`.
1013fn is_main_session(row: &sessionwiki::index::SessionRow) -> bool {
1014    row.kind == "main"
1015}
1016
1017/// The briefing for one indexed session, or `None` when the id names none.
1018pub fn brief(id: &str, max_chars: usize) -> Result<Option<String>> {
1019    if !index_is_writable() {
1020        return Ok(None);
1021    }
1022    let connection = open_readonly()?;
1023    let Some(row) = row_by_id(&connection, id)? else {
1024        return Ok(None);
1025    };
1026    let session = sessionwiki::index::session_from_index(&connection, &row)
1027        .context("read an indexed session")?;
1028    Ok(Some(sessionwiki::commands::brief_markdown(
1029        &session, max_chars, true,
1030    )))
1031}
1032
1033/// The passages of one indexed session that match `query`, or `None` when the
1034/// id names no indexed session.
1035///
1036/// Every matching message is returned with `context_messages` neighbours on
1037/// each side; overlapping groups are merged and each group's first block says
1038/// how many messages were skipped before it. Each block's text is capped at
1039/// `per_message_chars` characters, keeping the window around its first match.
1040pub fn transcript_hits(
1041    id: &str,
1042    query: &str,
1043    context_messages: usize,
1044    per_message_chars: usize,
1045) -> Result<Option<WikiHitTranscript>> {
1046    if !index_is_writable() {
1047        return Ok(None);
1048    }
1049    let connection = open_readonly()?;
1050    let Some(row) = row_by_id(&connection, id)? else {
1051        return Ok(None);
1052    };
1053    let session = sessionwiki::index::session_from_index(&connection, &row)
1054        .context("read an indexed session")?;
1055    Ok(Some(hit_transcript(
1056        &session,
1057        query,
1058        context_messages,
1059        per_message_chars,
1060    )))
1061}
1062
1063/// The matching passages of one loaded session, converted from SessionWiki's
1064/// own grep. Pure, so the conversion can be tested without an index on disk.
1065///
1066/// Matching, redaction and the excerpt window are `sessionwiki::grep`'s, so the
1067/// `sessionwiki grep` CLI and this preview report the same hits. Tool output
1068/// never anchors a passage: it is machine chatter the reader did not write,
1069/// a hit buried in it would open the preview on a wall of command output, and
1070/// the preview collapses tool runs anyway. Tool messages still appear as
1071/// context around a real match.
1072fn hit_transcript(
1073    session: &Session,
1074    query: &str,
1075    context_messages: usize,
1076    per_message_chars: usize,
1077) -> WikiHitTranscript {
1078    let found = sessionwiki::grep::grep_session(
1079        session,
1080        query,
1081        &sessionwiki::grep::GrepOpts {
1082            context_messages,
1083            chars: per_message_chars,
1084            max_matches: None,
1085            anchor_roles: vec![Role::User, Role::Assistant],
1086        },
1087    );
1088    WikiHitTranscript {
1089        blocks: found
1090            .hits
1091            .into_iter()
1092            .map(|hit| WikiHitBlock {
1093                role: role_name(hit.role).to_owned(),
1094                text: hit.text,
1095                hits: hit.matches,
1096                omitted_before: hit.omitted_before,
1097                truncated: hit.truncated,
1098            })
1099            .collect(),
1100        omitted_after: found.omitted_after,
1101    }
1102}
1103
1104fn role_name(role: Role) -> &'static str {
1105    match role {
1106        Role::User => "user",
1107        Role::Assistant => "assistant",
1108        Role::Tool => "tool",
1109    }
1110}
1111
1112/// What a restore needs from the index: the transcript as a snapshot the
1113/// compaction pipeline accepts, plus the title and project of the session it
1114/// came from.
1115pub struct ArchivedSession {
1116    pub title: String,
1117    /// The project directory the session ran in, when the row names one that
1118    /// still exists.
1119    pub project_directory: Option<PathBuf>,
1120    pub snapshot: mj_core::archive::CanonicalSessionSnapshot,
1121}
1122
1123/// Load one indexed session for restore, or `None` when the id names none.
1124pub fn archived_session(id: &str) -> Result<Option<ArchivedSession>> {
1125    if !index_is_writable() {
1126        return Ok(None);
1127    }
1128    let connection = open_readonly()?;
1129    let Some(row) = row_by_id(&connection, id)? else {
1130        return Ok(None);
1131    };
1132    let session = sessionwiki::index::session_from_index(&connection, &row)
1133        .context("read an indexed session")?;
1134    let snapshot = snapshot_of(&session)?;
1135    Ok(Some(ArchivedSession {
1136        title: session.title.clone(),
1137        project_directory: project_directory_of(&session.project),
1138        snapshot,
1139    }))
1140}
1141
1142// ---------------------------------------------------------------------------
1143// The archive job
1144// ---------------------------------------------------------------------------
1145
1146/// The stopped sessions that `archive_after_days = older_than_days` has caught,
1147/// children before their parents.
1148///
1149/// A session qualifies when its record is `Stopped`, its last update is at
1150/// least that many days old, and every sub-agent child it still has is being
1151/// archived in the same pass. The child rule is what keeps the pass from
1152/// destroying a session it did not choose: archiving a parent tears its
1153/// children down with it, so a child that is still running, or stopped but not
1154/// yet old enough, holds its parent back until the next pass.
1155///
1156/// Pure over controller state, so the rule can be tested without a daemon.
1157pub fn sessions_ready_to_archive(
1158    sessions: &BTreeMap<String, SessionRecord>,
1159    subagents: &BTreeMap<String, mj_core::subagent::SubagentRecord>,
1160    now: DateTime<Utc>,
1161    older_than_days: u32,
1162) -> Vec<String> {
1163    let cutoff = now - chrono::Duration::days(i64::from(older_than_days));
1164    let aged = |session_id: &String| {
1165        sessions.get(session_id).is_some_and(|record| {
1166            record.state == mj_core::state::SessionState::Stopped
1167                && parse_time(&record.updated_at).is_some_and(|updated| updated <= cutoff)
1168        })
1169    };
1170    let selected: BTreeSet<String> = sessions
1171        .keys()
1172        .filter(|session_id| aged(session_id))
1173        .filter(|session_id| {
1174            subagents
1175                .values()
1176                .filter(|child| &&child.parent_session_id == session_id)
1177                // A child whose record is already gone holds nothing open.
1178                .filter(|child| sessions.contains_key(&child.child_session_id))
1179                .all(|child| aged(&child.child_session_id))
1180        })
1181        .cloned()
1182        .collect();
1183    let mut ordered: Vec<String> = selected.iter().cloned().collect();
1184    ordered.sort_by_key(|session_id| std::cmp::Reverse(ancestor_depth(session_id, subagents)));
1185    ordered
1186}
1187
1188/// How many sub-agent parents a session has above it. Deeper sessions are
1189/// archived first so a parent never tears down a child the pass still has to
1190/// visit.
1191fn ancestor_depth(
1192    session_id: &str,
1193    subagents: &BTreeMap<String, mj_core::subagent::SubagentRecord>,
1194) -> usize {
1195    let mut depth = 0;
1196    let mut current = session_id;
1197    // Bounded by the map: a cycle cannot outlive one pass over every entry.
1198    while let Some(parent) = subagents
1199        .get(current)
1200        .map(|child| child.parent_session_id.as_str())
1201    {
1202        depth += 1;
1203        if depth > subagents.len() {
1204            break;
1205        }
1206        current = parent;
1207    }
1208    depth
1209}
1210
1211/// How much disk Mjolnir's own copies of sessions use, and how much an
1212/// `archive_after_days` value would free. "Mjolnir's own copy" is the
1213/// checkpoint archive plus the session's image attachments; the conversation
1214/// itself lives in the SessionWiki index and is not counted, because archiving
1215/// keeps it. The type lives in `mj-core` so the terminal UI can name it too.
1216pub use mj_core::state::ArchiveSpacePreview;
1217
1218/// The space every session uses now and, when `older_than_days` is set, the
1219/// space archiving after that many days would reclaim.
1220///
1221/// The reclaim figure uses the archive job's own selection rule but not its
1222/// "is it indexed yet" gate: that gate depends on how far the hourly index
1223/// sync has got, so applying it would make the estimate swing between zero and
1224/// the true value while the first index builds. This answers what the policy
1225/// would reclaim, not what the next tick happens to reclaim.
1226///
1227/// Walks the filesystem, so callers on the async runtime must run it in a
1228/// blocking task.
1229pub fn archive_space_preview(older_than_days: Option<u32>) -> Result<ArchiveSpacePreview> {
1230    let controller =
1231        Controller::load().context("load the session records to size their storage")?;
1232    Ok(archive_space_over(
1233        &mj_core::config::sessions_dir(),
1234        &controller.state.sessions,
1235        &controller.state.subagents,
1236        Utc::now(),
1237        older_than_days,
1238    ))
1239}
1240
1241/// The sizing itself, over given records and a given sessions directory, so it
1242/// can be tested without the live data directory.
1243fn archive_space_over(
1244    sessions_root: &Path,
1245    sessions: &BTreeMap<String, SessionRecord>,
1246    subagents: &BTreeMap<String, mj_core::subagent::SubagentRecord>,
1247    now: DateTime<Utc>,
1248    older_than_days: Option<u32>,
1249) -> ArchiveSpacePreview {
1250    let mut preview = ArchiveSpacePreview {
1251        sessions: sessions.len(),
1252        bytes: sessions
1253            .iter()
1254            .map(|(session_id, record)| session_bytes(sessions_root, session_id, record))
1255            .sum(),
1256        reclaimable_sessions: 0,
1257        reclaimable_bytes: 0,
1258    };
1259    if let Some(days) = older_than_days {
1260        let aged = sessions_ready_to_archive(sessions, subagents, now, days);
1261        preview.reclaimable_sessions = aged.len();
1262        preview.reclaimable_bytes = aged
1263            .iter()
1264            .filter_map(|session_id| {
1265                sessions
1266                    .get(session_id)
1267                    .map(|record| session_bytes(sessions_root, session_id, record))
1268            })
1269            .sum();
1270    }
1271    preview
1272}
1273
1274/// What archiving one session would free: its checkpoint archive and its
1275/// attachments. Anything already missing counts as zero.
1276fn session_bytes(sessions_root: &Path, session_id: &str, record: &SessionRecord) -> u64 {
1277    let checkpoint = record
1278        .checkpoint
1279        .as_ref()
1280        .and_then(|checkpoint| std::fs::metadata(&checkpoint.archive_path).ok())
1281        .filter(|metadata| metadata.is_file())
1282        .map(|metadata| metadata.len())
1283        .unwrap_or(0);
1284    let attachments = sessions_root
1285        .join(session_id)
1286        .join(mj_core::attachment::ATTACHMENT_DIR);
1287    let attachments = crate::import::claude::directory_size(&attachments).unwrap_or(0);
1288    checkpoint.saturating_add(attachments)
1289}
1290
1291/// Which of `session_ids` the index holds under this instance's own key, with
1292/// at least one message and not already archived.
1293///
1294/// This is the gate the archive job will not cross: Mjolnir only deletes its
1295/// own copy of a conversation SessionWiki has actually stored. Runs SQLite
1296/// work, so callers on the async runtime wrap it in `spawn_blocking`.
1297pub fn indexed_with_messages(session_ids: &[String]) -> Result<BTreeSet<String>> {
1298    if !index_is_writable() {
1299        // An index this daemon will not open holds nothing it may act on, and
1300        // the archive job deletes data, so it must find nothing here.
1301        return Ok(BTreeSet::new());
1302    }
1303    let connection = open_readonly()?;
1304    let sessions_dir = mj_core::config::sessions_dir();
1305    let mut indexed = BTreeSet::new();
1306    for session_id in session_ids {
1307        let key = format!("{}/{session_id}", sessions_dir.display());
1308        let rows = sessionwiki::index::resolve(&connection, session_id)
1309            .context("look up a stopped session in the SessionWiki index")?;
1310        if rows
1311            .iter()
1312            .any(|row| row.tool == TOOL && row.path == key && row.msg_count > 0 && !row.archived)
1313        {
1314            indexed.insert(session_id.clone());
1315        }
1316    }
1317    Ok(indexed)
1318}
1319
1320fn open_readonly() -> Result<rusqlite::Connection> {
1321    sessionwiki::index::open_readonly().context("open the SessionWiki index")
1322}
1323
1324/// The one row an id names exactly. `resolve` matches prefixes, which is right
1325/// for a person typing and wrong for a client passing an id back.
1326fn row_by_id(
1327    connection: &rusqlite::Connection,
1328    id: &str,
1329) -> Result<Option<sessionwiki::index::SessionRow>> {
1330    Ok(sessionwiki::index::resolve(connection, id)
1331        .context("look up an indexed session")?
1332        .into_iter()
1333        .find(|row| row.session_id == id))
1334}
1335
1336fn wiki_row(
1337    row: sessionwiki::index::SessionRow,
1338    snippet: Option<String>,
1339    live: &BTreeSet<String>,
1340) -> WikiRow {
1341    // Only this daemon's own sessions can be live here, and only under the key
1342    // shape the adapter writes: the checkpoint directory and the session id.
1343    let hel_session_id = (row.tool == TOOL)
1344        .then(|| row.path.rsplit('/').next().unwrap_or_default().to_owned())
1345        .filter(|session_id| live.contains(session_id));
1346    let native_id = sessionwiki::index::native_id_of(&row.path);
1347    WikiRow {
1348        id: row.session_id,
1349        tool: row.tool,
1350        project: row.project,
1351        title: row.title,
1352        started: row.started,
1353        msgs: row.msg_count,
1354        preview: row.preview,
1355        archived: row.archived,
1356        native_id,
1357        snippet,
1358        hel_session_id,
1359        // Filled in by `fill_session_tags` from the index's own tags; the row
1360        // itself does not carry them.
1361        target: None,
1362        profile: None,
1363        harness: None,
1364    }
1365}
1366
1367/// The project a restored session should open.
1368///
1369/// A Mjolnir session runs in a managed worktree under the repository it was
1370/// started from, and that worktree is gone once the session is archived. The
1371/// repository above it is what the user still has, so a worktree path is
1372/// reduced to it. Any other path is used as it stands, and a path that no
1373/// longer exists is left for the caller to replace.
1374fn project_directory_of(project: &str) -> Option<PathBuf> {
1375    if project.trim().is_empty() {
1376        return None;
1377    }
1378    let path = PathBuf::from(project);
1379    let repository = path
1380        .ancestors()
1381        .find(|ancestor| ancestor.file_name().is_some_and(|name| name == ".mj"))
1382        .and_then(std::path::Path::parent)
1383        .map(std::path::Path::to_path_buf)
1384        .unwrap_or(path);
1385    repository.is_dir().then_some(repository)
1386}
1387
1388/// Rebuild an indexed transcript as a canonical snapshot.
1389///
1390/// The snapshot is only ever read by the compaction pipeline, which wants
1391/// turns: a user message opens a turn and assistant and tool items attach to
1392/// it. Messages before the first user message therefore have nowhere to go and
1393/// are dropped, and a session with no user message at all cannot be restored.
1394fn snapshot_of(
1395    session: &sessionwiki::model::Session,
1396) -> Result<mj_core::archive::CanonicalSessionSnapshot> {
1397    use mj_core::archive::{
1398        CanonicalExecutionState, CanonicalSessionSnapshot, CanonicalSessionState,
1399        CanonicalTranscriptBody, CanonicalTranscriptItem,
1400    };
1401
1402    let started_ms = session
1403        .started
1404        .map(|time| time.timestamp_millis())
1405        .unwrap_or_default();
1406    let mut transcript: Vec<CanonicalTranscriptItem> = Vec::new();
1407    for message in &session.messages {
1408        let text = message.text.trim();
1409        if text.is_empty() {
1410            continue;
1411        }
1412        // Compaction attaches assistant and tool items to the open turn, so an
1413        // item before the first user message would be dropped anyway.
1414        if transcript.is_empty() && message.role != Role::User {
1415            continue;
1416        }
1417        let position = transcript.len() as u64 + 1;
1418        let body = match message.role {
1419            Role::User => CanonicalTranscriptBody::User {
1420                content: vec![serde_json::json!({"type": "text", "text": text})],
1421            },
1422            Role::Assistant => CanonicalTranscriptBody::Agent {
1423                chunks: vec![serde_json::json!({
1424                    "content": {"type": "text", "text": text}
1425                })],
1426                streaming: false,
1427            },
1428            // New indexes retain the shared projection; legacy rows contain only a title.
1429            Role::Tool => {
1430                let (call, terminal_outputs) = mj_transcript::summary::indexed_tool_call(
1431                    text,
1432                    &format!("wiki-tool-{position}"),
1433                );
1434                CanonicalTranscriptBody::Tool {
1435                    call,
1436                    terminal_outputs,
1437                    terminal_refs: Vec::new(),
1438                    presentation: None,
1439                }
1440            }
1441        };
1442        let created_at_ms = message
1443            .ts
1444            .map(|time| time.timestamp_millis())
1445            .unwrap_or(started_ms);
1446        transcript.push(CanonicalTranscriptItem {
1447            stable_id: format!("wiki-{position}"),
1448            position,
1449            // The validator wants an ordinal on agent messages and on nothing
1450            // else; one event per item makes the item's own position right.
1451            latest_content_event_ordinal: matches!(body, CanonicalTranscriptBody::Agent { .. })
1452                .then_some(position),
1453            created_at_ms,
1454            last_changed_at_ms: created_at_ms,
1455            body,
1456        });
1457    }
1458    anyhow::ensure!(
1459        !transcript.is_empty(),
1460        "the archived session has no prompt to restore from"
1461    );
1462
1463    let event_frontier = transcript.len() as u64;
1464    let last_activity_at_ms = transcript.last().map(|item| item.last_changed_at_ms);
1465    Ok(CanonicalSessionSnapshot {
1466        event_frontier,
1467        // Not a relay frontier, so there is no recorded digest to carry. It has
1468        // to be a well-formed non-genesis digest, and deriving it from the
1469        // session makes two restores of one session agree.
1470        event_frontier_digest: {
1471            use sha2::Digest;
1472            mj_core::hex::lower_hex(sha2::Sha256::digest(
1473                format!("sessionwiki:{}", session.id).as_bytes(),
1474            ))
1475        },
1476        session: CanonicalSessionState {
1477            execution: CanonicalExecutionState::Idle,
1478            last_activity_at_ms,
1479            session_title: Some(session.title.clone()).filter(|title| !title.trim().is_empty()),
1480            configuration: BTreeMap::new(),
1481        },
1482        transcript,
1483        queued_prompts: Vec::new(),
1484    })
1485}
1486
1487// ---------------------------------------------------------------------------
1488// Continuing an indexed session
1489// ---------------------------------------------------------------------------
1490
1491/// What continuing one indexed session means.
1492///
1493/// An agent that found a session with SessionWiki should not have to know
1494/// whose session it was, so the branch lives here and `mj resume --wiki` takes
1495/// it on the agent's behalf.
1496#[derive(Debug, Clone, PartialEq, Eq)]
1497pub enum WikiContinuation {
1498    /// A Mjolnir session this daemon still has a record of: resume it.
1499    Resume { session_id: String },
1500    /// A Mjolnir session whose record the archive job destroyed: start a new
1501    /// session seeded with a compacted hand-off.
1502    Restore { wiki_id: String },
1503    /// Another tool's session: import it, then resume what the import made.
1504    Import {
1505        harness: HarnessKind,
1506        native_session_id: String,
1507    },
1508}
1509
1510/// How to continue the indexed session a row describes.
1511///
1512/// Pure over the row so the branch can be tested without an index:
1513/// `path` is the row's stored path and `has_record` says whether controller
1514/// state still holds a session with this id.
1515pub fn wiki_continuation(
1516    wiki_id: &str,
1517    tool: &str,
1518    path: &Path,
1519    has_record: bool,
1520) -> Result<WikiContinuation> {
1521    if tool == TOOL {
1522        // `MjolnirAdapter::parse_key` names the session by its own Mjolnir id,
1523        // so a Mjolnir row's SessionWiki id is the session id.
1524        return Ok(match has_record {
1525            true => WikiContinuation::Resume {
1526                session_id: wiki_id.to_owned(),
1527            },
1528            false => WikiContinuation::Restore {
1529                wiki_id: wiki_id.to_owned(),
1530            },
1531        });
1532    }
1533    let harness = harness_adapters::harness_for_tool(tool)
1534        .with_context(|| format!("Mjolnir cannot continue a {tool} session"))?;
1535    let native_session_id = crate::import::native_session_id_from_path(harness, path)
1536        .with_context(|| {
1537            format!(
1538                "no {tool} session id in the indexed path {}",
1539                path.display()
1540            )
1541        })?;
1542    Ok(WikiContinuation::Import {
1543        harness,
1544        native_session_id,
1545    })
1546}
1547
1548/// What one indexed session is, as far as continuing it is concerned.
1549///
1550/// Read through [`wiki_session`]; the daemon serves it for `mj resume --wiki`
1551/// and for `mj sessions --session` when the id names no Mjolnir session.
1552pub fn wiki_session(
1553    wiki_id: &str,
1554    known_sessions: &BTreeSet<String>,
1555) -> Result<Option<WikiSessionInfo>> {
1556    if !index_is_writable() {
1557        return Ok(None);
1558    }
1559    let connection = open_readonly()?;
1560    let Some(row) = row_by_id(&connection, wiki_id)? else {
1561        return Ok(None);
1562    };
1563    let is_mjolnir = row.tool == TOOL;
1564    let mjolnir_session_id = is_mjolnir.then(|| row.session_id.clone());
1565    let has_record = mjolnir_session_id
1566        .as_deref()
1567        .is_some_and(|session_id| known_sessions.contains(session_id));
1568    let status = match (is_mjolnir, has_record) {
1569        (false, _) => WikiSessionStatus::Native,
1570        (true, true) => WikiSessionStatus::Mine,
1571        (true, false) => WikiSessionStatus::Archived,
1572    };
1573    let tags = match is_mjolnir {
1574        true => tags::read(&connection, &[row.session_id.as_str()])
1575            .context("read the indexed session metadata")?
1576            .remove(&row.session_id)
1577            .unwrap_or_default(),
1578        false => tags::MjTags::default(),
1579    };
1580    let harness = tags
1581        .harness
1582        .as_deref()
1583        .and_then(|id| id.parse::<HarnessKind>().ok())
1584        .or_else(|| {
1585            (!is_mjolnir)
1586                .then(|| harness_adapters::harness_for_tool(&row.tool))
1587                .flatten()
1588        });
1589    Ok(Some(WikiSessionInfo {
1590        wiki_id: row.session_id,
1591        tool: row.tool,
1592        path: PathBuf::from(row.path),
1593        status,
1594        mjolnir_session_id,
1595        profile_id: tags.profile,
1596        target_template_id: tags.target,
1597        harness,
1598        title: row.title,
1599        project: row.project,
1600    }))
1601}
1602
1603#[cfg(test)]
1604mod tests {
1605    use std::collections::BTreeMap;
1606    use std::path::Path;
1607
1608    /// The dispatch `mj resume --wiki` takes, over rows built by hand: the
1609    /// branch has to be right without an index behind it.
1610    mod continuation {
1611        use super::super::{WikiContinuation, wiki_continuation};
1612        use mj_core::config::HarnessKind;
1613        use std::path::Path;
1614
1615        #[test]
1616        fn a_mjolnir_row_with_a_record_is_resumed_and_one_without_is_restored() {
1617            let path = Path::new("/home/user/.local/share/mj/sessions/session-7");
1618            assert_eq!(
1619                wiki_continuation("session-7", "mjolnir", path, true).unwrap(),
1620                WikiContinuation::Resume {
1621                    session_id: "session-7".to_owned(),
1622                }
1623            );
1624            assert_eq!(
1625                wiki_continuation("session-7", "mjolnir", path, false).unwrap(),
1626                WikiContinuation::Restore {
1627                    wiki_id: "session-7".to_owned(),
1628                }
1629            );
1630        }
1631
1632        #[test]
1633        fn a_claude_code_row_is_imported_with_the_uuid_from_its_path() {
1634            let path = Path::new(
1635                "/home/user/.claude/projects/-home-user-app/7f3a1c20-0b11-4a55-9e0d-2c8a5d6f1b44.jsonl",
1636            );
1637            assert_eq!(
1638                wiki_continuation("abc123", "claude-code", path, false).unwrap(),
1639                WikiContinuation::Import {
1640                    harness: HarnessKind::Claude,
1641                    native_session_id: "7f3a1c20-0b11-4a55-9e0d-2c8a5d6f1b44".to_owned(),
1642                }
1643            );
1644        }
1645
1646        /// A Codex rollout's file name is a timestamp and the thread UUID, so
1647        /// the stem alone is not the id `mj import codex --session` takes.
1648        #[test]
1649        fn a_codex_row_is_imported_with_the_uuid_from_its_rollout_name() {
1650            let path = Path::new(
1651                "/home/user/.codex/sessions/2026/09/18/rollout-2026-09-18T09-15-00-7f3a1c20-0b11-4a55-9e0d-2c8a5d6f1b44.jsonl",
1652            );
1653            assert_eq!(
1654                wiki_continuation("abc123", "codex", path, false).unwrap(),
1655                WikiContinuation::Import {
1656                    harness: HarnessKind::Codex,
1657                    native_session_id: "7f3a1c20-0b11-4a55-9e0d-2c8a5d6f1b44".to_owned(),
1658                }
1659            );
1660        }
1661
1662        #[test]
1663        fn an_unknown_tool_is_an_error_that_names_it() {
1664            let error = wiki_continuation("abc123", "opencode", Path::new("/tmp/s.jsonl"), false)
1665                .unwrap_err();
1666            assert!(
1667                format!("{error:#}").contains("opencode"),
1668                "the error has to name the tool: {error:#}"
1669            );
1670        }
1671    }
1672
1673    use mj_checkpoint::archive::{
1674        ArchiveInput, BundleManifest, CanonicalExecutionState, CanonicalSessionSnapshot,
1675        CanonicalSessionState, CanonicalTranscriptBody, CanonicalTranscriptItem, SessionManifest,
1676        TargetManifest, write_archive_atomic,
1677    };
1678
1679    use super::*;
1680
1681    fn item(position: u64, body: CanonicalTranscriptBody) -> CanonicalTranscriptItem {
1682        // Only an agent message carries a content ordinal; the snapshot
1683        // validator rejects one on any other item and demands one here.
1684        let streamed = matches!(body, CanonicalTranscriptBody::Agent { .. });
1685        CanonicalTranscriptItem {
1686            stable_id: format!("item-{position}"),
1687            position,
1688            latest_content_event_ordinal: streamed.then_some(position),
1689            created_at_ms: 1_700_000_000_000 + i64::try_from(position).unwrap(),
1690            last_changed_at_ms: 1_700_000_000_000 + i64::try_from(position).unwrap(),
1691            body,
1692        }
1693    }
1694
1695    /// A managed checkpoint with one prompt, one reply, one tool call, and one
1696    /// thought, which is every transcript shape the adapter decides about.
1697    fn write_archive(directory: &Path, session_id: &str, frontier: u64) {
1698        let path = directory.join(format!(
1699            "{session_id}-{frontier}-archive-{}.hel.zip",
1700            "0".repeat(32)
1701        ));
1702        write_archive_atomic(
1703            &path,
1704            &ArchiveInput {
1705                session: SessionManifest {
1706                    id: session_id.into(),
1707                    title: "indexed session".into(),
1708                    harness_kind: mj_core::config::HarnessKind::Codex,
1709                    profile_id: "codex".into(),
1710                    native_session_id: "native-session".into(),
1711                    created_at: "2026-09-01T00:00:00Z".into(),
1712                    checkpointed_at: "2026-09-01T01:00:00Z".into(),
1713                    hel_version: "test".into(),
1714                    relay_version: "test".into(),
1715                    adapter_version: "test".into(),
1716                },
1717                target: TargetManifest {
1718                    template_id: "local".into(),
1719                    target_kind: "local-bare".into(),
1720                    details: BTreeMap::new(),
1721                },
1722                bundle: BundleManifest {
1723                    id: "project".into(),
1724                    primary_repository: "project".into(),
1725                },
1726                canonical_session: CanonicalSessionSnapshot {
1727                    event_frontier: 4,
1728                    event_frontier_digest: "a".repeat(64),
1729                    session: CanonicalSessionState {
1730                        execution: CanonicalExecutionState::Idle,
1731                        last_activity_at_ms: Some(1_700_000_000_004),
1732                        session_title: Some("snapshot title".into()),
1733                        configuration: BTreeMap::new(),
1734                    },
1735                    transcript: vec![
1736                        item(
1737                            1,
1738                            CanonicalTranscriptBody::User {
1739                                content: vec![serde_json::json!({
1740                                    "type": "text",
1741                                    "text": "index this session"
1742                                })],
1743                            },
1744                        ),
1745                        item(
1746                            2,
1747                            CanonicalTranscriptBody::Thought {
1748                                chunks: vec![serde_json::json!({
1749                                    "content": {"type": "text", "text": "pondering"}
1750                                })],
1751                                streaming: false,
1752                            },
1753                        ),
1754                        item(
1755                            3,
1756                            CanonicalTranscriptBody::Tool {
1757                                call: serde_json::json!({
1758                                    "toolCallId": "call-1",
1759                                    "title": "Edit config.toml",
1760                                    "kind": "edit",
1761                                    "status": "completed",
1762                                    "locations": [{"path": "/old/container/config.toml"}]
1763                                }),
1764                                terminal_outputs: Vec::new(),
1765                                terminal_refs: Vec::new(),
1766                                presentation: None,
1767                            },
1768                        ),
1769                        item(
1770                            4,
1771                            CanonicalTranscriptBody::Agent {
1772                                chunks: vec![serde_json::json!({
1773                                    "content": {"type": "text", "text": "done"}
1774                                })],
1775                                streaming: false,
1776                            },
1777                        ),
1778                    ],
1779                    queued_prompts: Vec::new(),
1780                },
1781                native_artifacts: Vec::new(),
1782                repositories: Vec::new(),
1783            },
1784        )
1785        .unwrap();
1786    }
1787
1788    fn adapter(directory: &Path, session_id: &str) -> MjolnirAdapter {
1789        adapter_with_live(directory, session_id, BTreeMap::new())
1790    }
1791
1792    fn adapter_with_live(
1793        directory: &Path,
1794        session_id: &str,
1795        live: BTreeMap<String, i64>,
1796    ) -> MjolnirAdapter {
1797        let record = SessionRecord {
1798            id: session_id.into(),
1799            ..record_template()
1800        };
1801        MjolnirAdapter {
1802            sessions_dir: directory.to_path_buf(),
1803            sessions: std::sync::Mutex::new(Sessions {
1804                records: BTreeMap::from([(session_id.to_owned(), record)]),
1805                subagent_ids: BTreeSet::new(),
1806                live,
1807            }),
1808            reload: false,
1809        }
1810    }
1811
1812    fn record_template() -> SessionRecord {
1813        SessionRecord {
1814            launch_base: None,
1815            build_cache: None,
1816            container_workspace: None,
1817            mjolnir_subagents: None,
1818            create_managed_worktree: None,
1819            workspace_id: mj_core::workspace::DEFAULT_WORKSPACE_ID.to_owned(),
1820            archived: false,
1821            container_cpus: None,
1822            container_memory: None,
1823            id: "0123456789abcdef0123456789abcdef".into(),
1824            title: "indexed session".into(),
1825            harness_kind: mj_core::config::HarnessKind::Codex,
1826            last_profile: "codex".into(),
1827            bundle_id: "project".into(),
1828            project_directory: Some(PathBuf::from("/home/dev/project")),
1829            managed_worktree: None,
1830            target_template_id: "local-bare".into(),
1831            resource_allocation: None,
1832            additional_mounts: Vec::new(),
1833            state: mj_core::state::SessionState::Stopped,
1834            target: None,
1835            native_session_id: Some("native-session".into()),
1836            acp_session_title: Some("the harness title".into()),
1837            session_title_override: None,
1838            created_at: "2026-09-01T00:00:00Z".into(),
1839            updated_at: "2026-09-01T01:00:00Z".into(),
1840            viewed_through_event_ordinal: 0,
1841            draft_input: String::new(),
1842            last_error: None,
1843            last_checkpoint_error: None,
1844            checkpoint: None,
1845        }
1846    }
1847
1848    #[test]
1849    fn the_newest_checkpoint_of_each_session_is_one_indexed_key() {
1850        let directory = tempfile::tempdir().unwrap();
1851        let session_id = "0123456789abcdef0123456789abcdef";
1852        write_archive(directory.path(), session_id, 1);
1853        write_archive(directory.path(), session_id, 7);
1854        let adapter = adapter(directory.path(), session_id);
1855
1856        let store = adapter.store().expect("the adapter is a shared store");
1857        let key = format!("{}/{session_id}", directory.path().display());
1858        assert_eq!(
1859            store
1860                .keys
1861                .iter()
1862                .map(|(key, _)| key.as_str())
1863                .collect::<Vec<_>>(),
1864            vec![key.as_str()]
1865        );
1866        assert!(!store.had_error);
1867        assert_eq!(store.files.len(), 1);
1868        assert!(
1869            store.files[0]
1870                .file_name()
1871                .unwrap()
1872                .to_str()
1873                .unwrap()
1874                .contains("-7-archive-"),
1875            "the newest checkpoint is the one indexed: {:?}",
1876            store.files[0]
1877        );
1878        assert_eq!(
1879            adapter.reconcile_scope(),
1880            Some(format!("{}/", directory.path().display()))
1881        );
1882
1883        let session = adapter.parse_key(&key).unwrap();
1884        assert_eq!(session.id, session_id);
1885        assert_eq!(session.tool, "mjolnir");
1886        assert_eq!(session.path, PathBuf::from(&key));
1887        assert_eq!(session.project, "/home/dev/project");
1888        assert_eq!(session.title, "the harness title");
1889        assert!(!session.subagent);
1890        assert_eq!(
1891            session.messages.iter().map(|m| m.role).collect::<Vec<_>>(),
1892            vec![Role::User, Role::Tool, Role::Assistant]
1893        );
1894        assert_eq!(session.messages[0].text, "index this session");
1895        let tool: serde_json::Value = serde_json::from_str(&session.messages[1].text).unwrap();
1896        assert_eq!(tool["name"], "Edit");
1897        assert_eq!(tool["call"]["title"], "Edit config.toml");
1898        assert_eq!(session.messages[2].text, "done");
1899        assert_eq!(session.touched, vec!["/old/container/config.toml"]);
1900    }
1901
1902    #[test]
1903    fn provenance_backfill_repairs_an_unchanged_checkpoint_without_rebuilding_the_index() {
1904        let _held = tags::testing::lock();
1905        let (_index_dir, mut connection) = tags::testing::isolated_index();
1906        let directory = tempfile::tempdir().unwrap();
1907        write_archive(directory.path(), "old-session", 4);
1908        let source = adapter(directory.path(), "old-session");
1909        let key = source.key_for("old-session");
1910        tags::testing::index_row(&connection, "old-session", "mjolnir");
1911        connection
1912            .execute(
1913                "UPDATE files SET path = ?1 WHERE session_id = 'old-session'",
1914                [&key],
1915            )
1916            .unwrap();
1917        provenance::backfill(&mut connection, &source).unwrap();
1918        assert_eq!(
1919            sessionwiki::index::files_for(&connection, "old-session").unwrap(),
1920            vec!["/old/container/config.toml"]
1921        );
1922        provenance::backfill(&mut connection, &source).unwrap();
1923        assert_eq!(
1924            sessionwiki::index::sessions_for_file(&connection, "config.toml", 20)
1925                .unwrap()
1926                .len(),
1927            1
1928        );
1929    }
1930
1931    fn projection(session_id: &str) -> mj_core::state::MaterializedSession {
1932        use mj_core::transcript::{TranscriptBody, TranscriptItem};
1933        let mut projected = mj_core::state::MaterializedSession::empty(session_id);
1934        let mut push = |position: u64, body: TranscriptBody| {
1935            let streamed = matches!(body, TranscriptBody::Agent { .. });
1936            projected
1937                .transcript
1938                .push(std::sync::Arc::new(TranscriptItem {
1939                    stable_id: format!("item-{position}"),
1940                    position,
1941                    latest_content_event_ordinal: streamed.then_some(position),
1942                    created_at_ms: 1_700_000_000_000 + i64::try_from(position).unwrap(),
1943                    last_changed_at_ms: 1_700_000_000_000 + i64::try_from(position).unwrap(),
1944                    body,
1945                }));
1946        };
1947        push(
1948            1,
1949            TranscriptBody::User {
1950                content: vec![serde_json::json!({"type": "text", "text": "still talking"})],
1951            },
1952        );
1953        push(
1954            2,
1955            TranscriptBody::Thought {
1956                chunks: vec![serde_json::json!({"content": {"type": "text", "text": "hmm"}})],
1957                streaming: false,
1958            },
1959        );
1960        push(
1961            3,
1962            TranscriptBody::Tool {
1963                call: serde_json::json!({"toolCallId": "c1", "title": "Read README.md"}),
1964                terminal_outputs: Vec::new(),
1965                terminal_refs: Vec::new(),
1966                presentation: None,
1967            },
1968        );
1969        push(
1970            4,
1971            TranscriptBody::Agent {
1972                chunks: vec![serde_json::json!({"content": {"type": "text", "text": "reading"}})],
1973                streaming: false,
1974            },
1975        );
1976        projected.session_title = Some("the live title".into());
1977        projected
1978    }
1979
1980    /// A session that has never been checkpointed is indexed from the
1981    /// daemon's own projection, with the same roles a checkpoint would give.
1982    #[test]
1983    fn a_running_session_is_indexed_from_its_stored_transcript() {
1984        let session_id = "0123456789abcdef0123456789abcdef";
1985        let messages = projected_messages(&projection(session_id));
1986        assert_eq!(
1987            messages.iter().map(|m| m.role).collect::<Vec<_>>(),
1988            vec![Role::User, Role::Tool, Role::Assistant]
1989        );
1990        assert_eq!(messages[0].text, "still talking");
1991        assert_eq!(messages[2].text, "reading");
1992        let tool: serde_json::Value = serde_json::from_str(&messages[1].text).unwrap();
1993        assert_eq!(tool["name"], "Read");
1994        assert_eq!(tool["call"]["title"], "Read README.md");
1995    }
1996
1997    /// A running session is listed under the same key as a stopped one, with
1998    /// its own change token, so it is searchable before it is ever closed and
1999    /// reconciliation never archives it. When it stops, the key stays and the
2000    /// checkpoint becomes its source.
2001    #[test]
2002    fn a_running_session_is_listed_with_its_own_change_token() {
2003        let directory = tempfile::tempdir().unwrap();
2004        let running = "0123456789abcdef0123456789abcdef";
2005        let never_checkpointed = "fedcba9876543210fedcba9876543210";
2006        write_archive(directory.path(), running, 3);
2007        let live = adapter_with_live(
2008            directory.path(),
2009            running,
2010            BTreeMap::from([
2011                (running.to_owned(), 1_900_000_000),
2012                (never_checkpointed.to_owned(), 1_900_000_001),
2013            ]),
2014        );
2015
2016        let store = live.store().expect("the adapter is a shared store");
2017        let key_of = |session_id: &str| format!("{}/{session_id}", directory.path().display());
2018        assert_eq!(
2019            store.keys,
2020            vec![
2021                (
2022                    key_of(running),
2023                    1_900_000_000 * 1024 + i64::from(mj_transcript::summary::SUMMARY_VERSION)
2024                ),
2025                (
2026                    key_of(never_checkpointed),
2027                    1_900_000_001 * 1024 + i64::from(mj_transcript::summary::SUMMARY_VERSION)
2028                ),
2029            ],
2030            "a live session's own token replaces the checkpoint's"
2031        );
2032
2033        // Once it stops it leaves the live set, and the checkpoint's own
2034        // modification time is the token again.
2035        let stopped = adapter(directory.path(), running);
2036        let keys = stopped.store().expect("a shared store").keys;
2037        assert_eq!(keys.len(), 1);
2038        assert_eq!(keys[0].0, key_of(running));
2039        assert_ne!(keys[0].1, 1_900_000_000);
2040        assert_eq!(
2041            stopped.parse_key(&key_of(running)).unwrap().title,
2042            "the harness title",
2043            "a stopped session is parsed from its checkpoint"
2044        );
2045    }
2046
2047    /// Renaming a session leaves its conversation untouched, so only the
2048    /// record's own last update can tell the index the title moved.
2049    #[test]
2050    fn a_rename_moves_a_session_change_token() {
2051        let directory = tempfile::tempdir().unwrap();
2052        let session_id = "0123456789abcdef0123456789abcdef";
2053        write_archive(directory.path(), session_id, 1);
2054        let adapter = adapter(directory.path(), session_id);
2055        let before = adapter.store().expect("a shared store").keys[0].1;
2056
2057        {
2058            let mut sessions = adapter.sessions.lock().unwrap();
2059            let record = sessions.records.get_mut(session_id).unwrap();
2060            record.session_title_override = Some("the new name".into());
2061            record.updated_at = "2099-01-01T00:00:00Z".into();
2062        }
2063        let after = adapter.store().expect("a shared store").keys[0].1;
2064        assert!(
2065            after > before,
2066            "a renamed session is re-indexed: {before} then {after}"
2067        );
2068        assert_eq!(
2069            adapter
2070                .parse_key(&format!("{}/{session_id}", directory.path().display()))
2071                .unwrap()
2072                .title,
2073            "the new name"
2074        );
2075    }
2076
2077    fn indexed(messages: Vec<(Role, &str)>) -> sessionwiki::model::Session {
2078        Session {
2079            id: "0123456789abcdef0123456789abcdef".into(),
2080            tool: "mjolnir",
2081            path: PathBuf::from("/sessions/0123456789abcdef0123456789abcdef"),
2082            project: "/home/dev/project".into(),
2083            started: DateTime::from_timestamp_millis(1_700_000_000_000),
2084            ended: None,
2085            title: "the archived session".into(),
2086            subagent: false,
2087            messages: messages
2088                .into_iter()
2089                .map(|(role, text)| Message {
2090                    role,
2091                    text: text.to_owned(),
2092                    ts: None,
2093                })
2094                .collect(),
2095            touched: Vec::new(),
2096            edits: Vec::new(),
2097        }
2098    }
2099
2100    /// A hit is found whatever the case of the query or of the transcript, and
2101    /// the reported range covers the matched text in the returned block.
2102    #[test]
2103    fn transcript_hits_locates_case_insensitive_matches() {
2104        let session = indexed(vec![
2105            (Role::User, "Make the Tests green"),
2106            (Role::Assistant, "the tests are green now"),
2107        ]);
2108
2109        let found = hit_transcript(&session, "TESTS", 0, 4_000);
2110
2111        assert_eq!(found.blocks.len(), 2, "both messages contain the query");
2112        assert_eq!(found.blocks[0].role, "user");
2113        let (start, end) = found.blocks[0].hits[0];
2114        assert_eq!(&found.blocks[0].text[start..end], "Tests");
2115        let (start, end) = found.blocks[1].hits[0];
2116        assert_eq!(&found.blocks[1].text[start..end], "tests");
2117        assert!(!found.blocks[0].truncated);
2118        assert_eq!(found.omitted_after, 0);
2119    }
2120
2121    /// Context messages come back around each hit, with the gap between two
2122    /// groups counted rather than silently closed.
2123    #[test]
2124    fn transcript_hits_keeps_context_and_marks_omissions() {
2125        let session = indexed(vec![
2126            (Role::User, "zero"),
2127            (Role::Assistant, "one needle one"),
2128            (Role::Tool, "two"),
2129            (Role::User, "three"),
2130            (Role::Assistant, "four"),
2131            (Role::Tool, "five"),
2132            (Role::User, "six needle six"),
2133            (Role::Assistant, "seven"),
2134            (Role::User, "eight"),
2135        ]);
2136
2137        let found = hit_transcript(&session, "needle", 1, 4_000);
2138
2139        let shown: Vec<(&str, &str, usize)> = found
2140            .blocks
2141            .iter()
2142            .map(|block| {
2143                (
2144                    block.role.as_str(),
2145                    block.text.as_str(),
2146                    block.omitted_before,
2147                )
2148            })
2149            .collect();
2150        assert_eq!(
2151            shown,
2152            vec![
2153                ("user", "zero", 0),
2154                ("assistant", "one needle one", 0),
2155                ("tool", "two", 0),
2156                ("tool", "five", 2),
2157                ("user", "six needle six", 0),
2158                ("assistant", "seven", 0),
2159            ]
2160        );
2161        assert_eq!(found.omitted_after, 1, "the last message is not shown");
2162        assert!(found.blocks[0].hits.is_empty(), "context has no hits");
2163    }
2164
2165    /// A query that only occurs in tool output finds nothing, and a tool
2166    /// message beside a real match still comes back as context. Tool text is
2167    /// machine chatter: anchoring a passage on it opens the preview on command
2168    /// output the reader never wrote, and the preview collapses tool runs, so
2169    /// the match could not be shown even if it were returned.
2170    #[test]
2171    fn transcript_hits_never_anchor_on_tool_output() {
2172        let session = indexed(vec![
2173            (Role::User, "make it build"),
2174            (Role::Tool, "cargo build --needle"),
2175            (Role::Assistant, "it builds"),
2176        ]);
2177
2178        let only_in_a_tool = hit_transcript(&session, "needle", 1, 4_000);
2179        assert!(
2180            only_in_a_tool.blocks.is_empty(),
2181            "tool output must not anchor a passage, got {:?}",
2182            only_in_a_tool.blocks
2183        );
2184
2185        let beside_a_match = hit_transcript(&session, "builds", 1, 4_000);
2186        let shown: Vec<(&str, bool)> = beside_a_match
2187            .blocks
2188            .iter()
2189            .map(|block| (block.role.as_str(), !block.hits.is_empty()))
2190            .collect();
2191        assert_eq!(
2192            shown,
2193            vec![("tool", false), ("assistant", true)],
2194            "a tool message is still context around a real match"
2195        );
2196    }
2197
2198    /// A long message is cut down to the caller's budget around its first hit,
2199    /// not from the start, so the match is always in what comes back.
2200    #[test]
2201    fn transcript_hits_window_keeps_the_first_hit() {
2202        let filler = "x".repeat(4_000);
2203        let session = indexed(vec![(Role::User, &format!("{filler} needle {filler}"))]);
2204
2205        let found = hit_transcript(&session, "needle", 0, 100);
2206
2207        let block = &found.blocks[0];
2208        assert!(block.truncated);
2209        assert_eq!(block.text.chars().count(), 100);
2210        assert_eq!(block.hits.len(), 1, "the windowed text keeps its hit");
2211        let (start, end) = block.hits[0];
2212        assert_eq!(&block.text[start..end], "needle");
2213        assert!(
2214            start >= 20,
2215            "the window keeps lead-in before the hit, got {start}"
2216        );
2217    }
2218
2219    /// The snapshot a restore hands to compaction has to satisfy the same
2220    /// validator a real checkpoint does, and has to carry every message in
2221    /// order.
2222    #[test]
2223    fn a_restored_snapshot_is_a_valid_transcript_of_the_indexed_session() {
2224        let snapshot = snapshot_of(&indexed(vec![
2225            (Role::User, "make the tests green"),
2226            (Role::Tool, "Read src/lib.rs"),
2227            (Role::Assistant, "they are green now"),
2228            (Role::User, "  "),
2229        ]))
2230        .unwrap();
2231
2232        snapshot.validate().expect("the snapshot is well formed");
2233        assert_eq!(snapshot.event_frontier, 3);
2234        assert_eq!(
2235            snapshot.session.session_title.as_deref(),
2236            Some("the archived session")
2237        );
2238        assert!(snapshot.session.last_activity_at_ms.is_some());
2239        let bodies = snapshot
2240            .transcript
2241            .iter()
2242            .map(|item| match &item.body {
2243                mj_core::archive::CanonicalTranscriptBody::User { content } => (
2244                    "user",
2245                    mj_core::transcript::materialized_content_text(content),
2246                ),
2247                mj_core::archive::CanonicalTranscriptBody::Agent { chunks, .. } => (
2248                    "agent",
2249                    mj_core::transcript::materialized_chunks_text(chunks),
2250                ),
2251                mj_core::archive::CanonicalTranscriptBody::Tool { call, .. } => (
2252                    "tool",
2253                    call["title"].as_str().unwrap_or_default().to_owned(),
2254                ),
2255                _ => ("other", String::new()),
2256            })
2257            .collect::<Vec<_>>();
2258        assert_eq!(
2259            bodies,
2260            vec![
2261                ("user", "make the tests green".to_owned()),
2262                ("tool", "Read src/lib.rs".to_owned()),
2263                ("agent", "they are green now".to_owned()),
2264            ],
2265            "the blank message is dropped and every other one keeps its role"
2266        );
2267    }
2268
2269    /// Compaction attaches assistant and tool items to the open turn, so an
2270    /// index that starts mid-conversation must not produce a snapshot whose
2271    /// first item has no turn to join.
2272    #[test]
2273    fn messages_before_the_first_prompt_are_dropped() {
2274        let snapshot = snapshot_of(&indexed(vec![
2275            (Role::Assistant, "still working"),
2276            (Role::User, "carry on"),
2277        ]))
2278        .unwrap();
2279        assert_eq!(snapshot.transcript.len(), 1);
2280        assert_eq!(snapshot.transcript[0].position, 1);
2281        snapshot.validate().unwrap();
2282
2283        let error = snapshot_of(&indexed(vec![(Role::Assistant, "nobody asked")])).unwrap_err();
2284        assert!(
2285            error.to_string().contains("no prompt"),
2286            "a session with no prompt cannot be restored: {error}"
2287        );
2288    }
2289
2290    fn record(
2291        session_id: &str,
2292        state: mj_core::state::SessionState,
2293        updated_at: &str,
2294    ) -> SessionRecord {
2295        SessionRecord {
2296            id: session_id.into(),
2297            state,
2298            updated_at: updated_at.into(),
2299            ..record_template()
2300        }
2301    }
2302
2303    fn child(child_session_id: &str, parent_session_id: &str) -> mj_core::subagent::SubagentRecord {
2304        mj_core::subagent::SubagentRecord {
2305            child_session_id: child_session_id.into(),
2306            parent_session_id: parent_session_id.into(),
2307            task_name: "task".into(),
2308            profile_id: "codex".into(),
2309            model: None,
2310            effort: None,
2311            working_directory: PathBuf::new(),
2312            initial_prompt: "do the thing".into(),
2313            request_key: "key".into(),
2314            created_at: "2026-09-01T00:00:00Z".into(),
2315            noticed_turn: None,
2316        }
2317    }
2318
2319    fn ready(
2320        sessions: Vec<SessionRecord>,
2321        children: Vec<mj_core::subagent::SubagentRecord>,
2322    ) -> Vec<String> {
2323        let now = parse_time("2026-09-10T00:00:00Z").unwrap();
2324        sessions_ready_to_archive(
2325            &sessions
2326                .into_iter()
2327                .map(|record| (record.id.clone(), record))
2328                .collect(),
2329            &children
2330                .into_iter()
2331                .map(|child| (child.child_session_id.clone(), child))
2332                .collect(),
2333            now,
2334            3,
2335        )
2336    }
2337
2338    /// A session whose checkpoint archive and attachments sit under `root`.
2339    fn sized_session(
2340        root: &Path,
2341        session_id: &str,
2342        updated_at: &str,
2343        checkpoint_bytes: usize,
2344        attachment_bytes: &[usize],
2345    ) -> SessionRecord {
2346        let archive_path = root.join(format!("{session_id}.hel.zip"));
2347        std::fs::write(&archive_path, vec![b'c'; checkpoint_bytes]).unwrap();
2348        if !attachment_bytes.is_empty() {
2349            let attachments = root
2350                .join(session_id)
2351                .join(mj_core::attachment::ATTACHMENT_DIR);
2352            std::fs::create_dir_all(&attachments).unwrap();
2353            for (index, size) in attachment_bytes.iter().enumerate() {
2354                std::fs::write(attachments.join(format!("{index}.png")), vec![b'a'; *size])
2355                    .unwrap();
2356            }
2357        }
2358        SessionRecord {
2359            checkpoint: Some(mj_core::state::CheckpointMetadata {
2360                archive_path,
2361                sha256: "0".repeat(64),
2362                created_at: updated_at.into(),
2363                event_frontier: 1,
2364            }),
2365            ..record(
2366                session_id,
2367                mj_core::state::SessionState::Stopped,
2368                updated_at,
2369            )
2370        }
2371    }
2372
2373    #[test]
2374    fn the_space_preview_sizes_every_session_and_only_the_aged_ones_as_reclaimable() {
2375        let directory = tempfile::tempdir().unwrap();
2376        let root = directory.path();
2377        let sessions: BTreeMap<String, SessionRecord> = [
2378            sized_session(root, "old-stopped", "2026-09-01T00:00:00Z", 1000, &[10, 20]),
2379            sized_session(root, "just-stopped", "2026-09-09T00:00:00Z", 500, &[]),
2380            // A record whose checkpoint file is already gone counts as zero
2381            // rather than failing the whole estimate.
2382            SessionRecord {
2383                checkpoint: Some(mj_core::state::CheckpointMetadata {
2384                    archive_path: root.join("missing.hel.zip"),
2385                    sha256: "0".repeat(64),
2386                    created_at: "2026-09-01T00:00:00Z".into(),
2387                    event_frontier: 1,
2388                }),
2389                ..record(
2390                    "lost-checkpoint",
2391                    mj_core::state::SessionState::Stopped,
2392                    "2026-09-01T00:00:00Z",
2393                )
2394            },
2395        ]
2396        .into_iter()
2397        .map(|record| (record.id.clone(), record))
2398        .collect();
2399        let now = parse_time("2026-09-10T00:00:00Z").unwrap();
2400
2401        let all = archive_space_over(root, &sessions, &BTreeMap::new(), now, None);
2402        assert_eq!(all.sessions, 3);
2403        assert_eq!(all.bytes, 1530);
2404        assert_eq!(all.reclaimable_sessions, 0);
2405        assert_eq!(all.reclaimable_bytes, 0);
2406
2407        let aged = archive_space_over(root, &sessions, &BTreeMap::new(), now, Some(3));
2408        assert_eq!(aged.bytes, 1530);
2409        assert_eq!(
2410            (aged.reclaimable_sessions, aged.reclaimable_bytes),
2411            (2, 1030),
2412            "only the sessions the job would archive count, attachments included"
2413        );
2414    }
2415
2416    #[test]
2417    fn only_stopped_sessions_past_the_cut_off_are_archived() {
2418        use mj_core::state::SessionState;
2419        let selected = ready(
2420            vec![
2421                record("old-stopped", SessionState::Stopped, "2026-09-01T00:00:00Z"),
2422                record(
2423                    "just-stopped",
2424                    SessionState::Stopped,
2425                    "2026-09-09T00:00:00Z",
2426                ),
2427                record("old-running", SessionState::Running, "2026-09-01T00:00:00Z"),
2428                record("old-error", SessionState::Error, "2026-09-01T00:00:00Z"),
2429                record("unparsable", SessionState::Stopped, "not a time"),
2430                // Exactly the cut-off counts as old enough.
2431                record("at-the-edge", SessionState::Stopped, "2026-09-07T00:00:00Z"),
2432            ],
2433            Vec::new(),
2434        );
2435        assert_eq!(selected, vec!["at-the-edge", "old-stopped"]);
2436    }
2437
2438    #[test]
2439    fn a_child_the_pass_is_not_archiving_holds_its_parent_back() {
2440        use mj_core::state::SessionState;
2441        let selected = ready(
2442            vec![
2443                record("parent", SessionState::Stopped, "2026-09-01T00:00:00Z"),
2444                record(
2445                    "running-child",
2446                    SessionState::Running,
2447                    "2026-09-01T00:00:00Z",
2448                ),
2449            ],
2450            vec![child("running-child", "parent")],
2451        );
2452        assert!(selected.is_empty(), "the parent must wait: {selected:?}");
2453
2454        let selected = ready(
2455            vec![
2456                record("parent", SessionState::Stopped, "2026-09-01T00:00:00Z"),
2457                record("young-child", SessionState::Stopped, "2026-09-09T00:00:00Z"),
2458            ],
2459            vec![child("young-child", "parent")],
2460        );
2461        assert!(selected.is_empty(), "the parent must wait: {selected:?}");
2462
2463        // A child whose record is already gone holds nothing open.
2464        let selected = ready(
2465            vec![record(
2466                "parent",
2467                SessionState::Stopped,
2468                "2026-09-01T00:00:00Z",
2469            )],
2470            vec![child("departed-child", "parent")],
2471        );
2472        assert_eq!(selected, vec!["parent"]);
2473    }
2474
2475    #[test]
2476    fn children_are_archived_before_their_parents() {
2477        use mj_core::state::SessionState;
2478        let selected = ready(
2479            vec![
2480                record("parent", SessionState::Stopped, "2026-09-01T00:00:00Z"),
2481                record("child", SessionState::Stopped, "2026-09-01T00:00:00Z"),
2482                record("grandchild", SessionState::Stopped, "2026-09-01T00:00:00Z"),
2483            ],
2484            vec![child("child", "parent"), child("grandchild", "child")],
2485        );
2486        assert_eq!(selected, vec!["grandchild", "child", "parent"]);
2487    }
2488
2489    #[test]
2490    fn native_adapters_cover_every_enabled_profile_home() {
2491        use mj_core::config::{Config, HarnessKind, HarnessProfile};
2492
2493        fn profile(kind: HarnessKind, home: &str, enabled: bool) -> HarnessProfile {
2494            HarnessProfile {
2495                enabled,
2496                kind,
2497                home: PathBuf::from(home),
2498                environment: BTreeMap::new(),
2499                context_window_bytes: None,
2500                guardian_review_model: None,
2501            }
2502        }
2503
2504        let mut config = Config::default();
2505        for (id, built) in [
2506            (
2507                "codex",
2508                profile(HarnessKind::Codex, "/home/dev/.codex3", true),
2509            ),
2510            (
2511                "codex-ds",
2512                profile(HarnessKind::Codex, "/home/dev/.codex-ds", true),
2513            ),
2514            // A second profile on one home must not add a second adapter.
2515            (
2516                "codex-alt",
2517                profile(HarnessKind::Codex, "/home/dev/.codex3", true),
2518            ),
2519            (
2520                "codex-off",
2521                profile(HarnessKind::Codex, "/home/dev/.codex-off", false),
2522            ),
2523            (
2524                "claude",
2525                profile(HarnessKind::Claude, "/home/dev/.claude4", true),
2526            ),
2527            ("kimi", profile(HarnessKind::Kimi, "/home/dev/.kimi", true)),
2528            ("grok", profile(HarnessKind::Grok, "/home/dev/.grok", true)),
2529            ("muse", profile(HarnessKind::Muse, "/home/dev/muse", true)),
2530            (
2531                "muse-off",
2532                profile(HarnessKind::Muse, "/home/dev/muse-off", false),
2533            ),
2534        ] {
2535            config.profiles.insert(id.into(), built);
2536        }
2537
2538        let adapters = native_adapters(&config);
2539        let roots: Vec<(&str, Option<PathBuf>)> = adapters
2540            .iter()
2541            .map(|adapter| (adapter.name(), adapter.root()))
2542            .collect();
2543
2544        let codex: Vec<&Option<PathBuf>> = roots
2545            .iter()
2546            .filter(|(name, _)| *name == "codex")
2547            .map(|(_, root)| root)
2548            .collect();
2549        assert_eq!(
2550            codex,
2551            vec![
2552                &Some(PathBuf::from("/home/dev/.codex3/sessions")),
2553                &Some(PathBuf::from("/home/dev/.codex-ds/sessions")),
2554            ],
2555            "one adapter per enabled Codex home, deduplicated: {roots:?}"
2556        );
2557
2558        let claude: Vec<&Option<PathBuf>> = roots
2559            .iter()
2560            .filter(|(name, _)| *name == "claude-code")
2561            .map(|(_, root)| root)
2562            .collect();
2563        assert_eq!(
2564            claude,
2565            vec![&Some(PathBuf::from("/home/dev/.claude4/projects"))],
2566            "one adapter for the enabled Claude home: {roots:?}"
2567        );
2568
2569        for (_, root) in &roots {
2570            let Some(root) = root else { continue };
2571            let text = root.to_string_lossy();
2572            assert!(
2573                !text.contains(".codex-off"),
2574                "a disabled profile must not be indexed: {roots:?}"
2575            );
2576            assert!(
2577                !text.ends_with("/.codex/sessions") && !text.ends_with("/.claude/projects"),
2578                "the stock homes are not indexed unless a profile names them: {roots:?}"
2579            );
2580        }
2581
2582        // SessionWiki has no adapter for these three, so Mjolnir supplies one
2583        // per enabled profile home under its own tool name.
2584        for (name, root) in [
2585            ("kimi-code", PathBuf::from("/home/dev/.kimi/sessions")),
2586            ("grok-build", PathBuf::from("/home/dev/.grok/sessions")),
2587            (
2588                "muse",
2589                mj_checkpoint::native::muse_sessions_root(Path::new("/home/dev/muse")).unwrap(),
2590            ),
2591        ] {
2592            let found: Vec<&Option<PathBuf>> = roots
2593                .iter()
2594                .filter(|(found, _)| *found == name)
2595                .map(|(_, root)| root)
2596                .collect();
2597            assert_eq!(found, vec![&Some(root)], "one {name} adapter: {roots:?}");
2598        }
2599
2600        for (_, root) in &roots {
2601            let Some(root) = root else { continue };
2602            assert!(
2603                !root.to_string_lossy().contains("muse-off"),
2604                "a disabled profile must not be indexed: {roots:?}"
2605            );
2606        }
2607
2608        assert!(
2609            roots.iter().any(|(name, _)| *name == "gemini"),
2610            "the other built-in adapters are kept: {roots:?}"
2611        );
2612    }
2613
2614    /// A Mjolnir row carries the target, profile and harness the sync stored
2615    /// in the index; a row from another tool carries none, because only
2616    /// Mjolnir writes those tags.
2617    #[test]
2618    fn query_rows_returns_the_indexed_target_profile_and_harness() {
2619        let _held = tags::testing::lock();
2620        let (_directory, connection) = tags::testing::isolated_index();
2621        tags::testing::index_row(&connection, "mj-session", TOOL);
2622        tags::testing::index_row(&connection, "codex-session", "codex");
2623        tags::write(
2624            &connection,
2625            "mj-session",
2626            &tags::MjTags {
2627                target: Some("Prod-Box".into()),
2628                profile: Some("codex-Main".into()),
2629                harness: Some("codex".into()),
2630            },
2631        )
2632        .expect("write the session metadata");
2633
2634        let rows = query_rows("", 10, &BTreeSet::new(), false).expect("query the index");
2635        let mjolnir = rows
2636            .iter()
2637            .find(|row| row.id == "mj-session")
2638            .expect("the Mjolnir row is returned");
2639        assert_eq!(mjolnir.target.as_deref(), Some("Prod-Box"));
2640        assert_eq!(mjolnir.profile.as_deref(), Some("codex-Main"));
2641        assert_eq!(mjolnir.harness.as_deref(), Some("codex"));
2642
2643        let codex = rows
2644            .iter()
2645            .find(|row| row.id == "codex-session")
2646            .expect("the Codex row is returned");
2647        assert_eq!(codex.target, None);
2648        assert_eq!(codex.profile, None);
2649        assert_eq!(codex.harness, None);
2650    }
2651
2652    #[test]
2653    fn one_flag_keeps_sub_agents_out_of_every_query_path() {
2654        let _held = tags::testing::lock();
2655        let (_directory, connection) = tags::testing::isolated_index();
2656        for (session_id, kind) in [("main-session", "main"), ("sub-session", "sub")] {
2657            tags::testing::index_row(&connection, session_id, "claude");
2658            connection
2659                .execute(
2660                    "UPDATE files SET kind = ?2 WHERE session_id = ?1",
2661                    rusqlite::params![session_id, kind],
2662                )
2663                .expect("set the session kind");
2664            connection
2665                .execute(
2666                    "INSERT INTO messages(session_id, role, text)
2667                     VALUES (?1, 'user', 'fix the bridge derivation zq')",
2668                    [session_id],
2669                )
2670                .expect("insert a message");
2671            connection
2672                .execute(
2673                    "INSERT INTO msgs(rowid, text) VALUES (?1, 'fix the bridge derivation zq')",
2674                    [connection.last_insert_rowid()],
2675                )
2676                .expect("index the message");
2677        }
2678        let ids = |query: &str, include_subagents: bool| {
2679            let mut ids: Vec<String> = query_rows(query, 10, &BTreeSet::new(), include_subagents)
2680                .expect("query the index")
2681                .into_iter()
2682                .map(|row| row.id)
2683                .collect();
2684            ids.sort();
2685            ids
2686        };
2687
2688        // The recent list, full-text search, short-query scan, and title match.
2689        for query in ["", "bridge derivation", "zq", "an indexed session"] {
2690            assert_eq!(ids(query, false), ["main-session"], "query {query:?}");
2691            assert_eq!(
2692                ids(query, true),
2693                ["main-session", "sub-session"],
2694                "query {query:?}"
2695            );
2696        }
2697    }
2698}