Skip to main content

mj_controller/
sessionwiki.rs

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