Skip to main content

mj_controller/
sessionwiki.rs

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