Skip to main content

supercode_interchange/
catalog.rs

1//! Discovery and durable addressing for sessions written by external harnesses.
2//!
3//! [`HarnessCatalog`] is intentionally about persisted state. It does not
4//! claim that the process which wrote a session is still alive or attachable.
5
6use std::collections::{BTreeSet, HashMap, HashSet};
7use std::fs::{self, File};
8use std::io::{BufRead, BufReader, Read, Seek, SeekFrom};
9use std::path::{Component, Path, PathBuf};
10use std::time::UNIX_EPOCH;
11
12use rusqlite::Connection;
13use serde::{Deserialize, Serialize};
14use serde_json::Value;
15
16use crate::native_store::load_native_store_family;
17use crate::ontology::{Binding, OrchestratorBindingRow};
18use crate::session::{
19    hermes_capture_nouns, openclaw_agent_id_from_path, openclaw_capture_header_nouns,
20    percent_decode_path, OrchestrationNouns, SessionMeta, SessionSource,
21};
22use crate::{Error, Fidelity, Result, Session, SessionFollower};
23
24pub use crate::ontology::HarnessId;
25
26/// Durable storage address for a persisted session.
27#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
28#[serde(tag = "kind", rename_all = "snake_case")]
29pub enum StorageLocator {
30    /// One session stored in one file.
31    File {
32        /// Absolute or caller-resolvable path to the transcript.
33        path: PathBuf,
34    },
35    /// One logical session selected from a SQLite store.
36    Sqlite {
37        /// Path to the SQLite database.
38        path: PathBuf,
39        /// Harness-native stable selector, currently an OpenCode session id.
40        selector: String,
41    },
42}
43
44impl StorageLocator {
45    /// Return the underlying file or database path.
46    pub fn path(&self) -> &Path {
47        match self {
48            Self::File { path } | Self::Sqlite { path, .. } => path,
49        }
50    }
51}
52
53/// Stable identity for a persisted harness session.
54#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
55pub struct SessionLocator {
56    /// Harness which owns the storage format.
57    pub harness: HarnessId,
58    /// Harness-native session identity.
59    pub session_id: String,
60    /// Exact storage address needed to load the session again.
61    pub storage: StorageLocator,
62}
63
64/// Lightweight metadata returned by catalog discovery.
65#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
66pub struct SessionDescriptor {
67    /// Durable address accepted by [`HarnessCatalog::load`] and
68    /// [`HarnessCatalog::follow`].
69    pub locator: SessionLocator,
70    /// Working directory recorded by the harness.
71    pub cwd: Option<PathBuf>,
72    /// Harness-provided title, when cheaply available.
73    pub title: Option<String>,
74    /// Oldest-first bounded conversation messages for a fallback topic when
75    /// the harness does not publish a useful title. These are read only for
76    /// the returned page and interpreted by the same presentation projection
77    /// as an opened conversation.
78    #[serde(default, skip_serializing_if = "Vec::is_empty")]
79    pub preview_candidates: Vec<SessionPreviewCandidate>,
80    /// Newest-first bounded conversation messages for compact list previews.
81    /// These are read only for the returned page, never for the entire
82    /// catalog, and are interpreted by the same presentation projection as
83    /// an opened conversation.
84    #[serde(default, skip_serializing_if = "Vec::is_empty")]
85    pub latest_message_candidates: Vec<SessionPreviewCandidate>,
86    /// Time of the last recorded TURN as Unix epoch milliseconds.
87    ///
88    /// Not the transcript file's mtime: a harness rewrites its transcript for
89    /// reasons that are not conversation (Claude Code appends untimestamped
90    /// `bridge-session` records while a session merely sits open, and moves
91    /// mtime again on resume), so mtime ranks idle-but-open sessions above
92    /// genuinely active ones. Falls back to mtime only when no timestamped
93    /// record is readable.
94    pub updated_at_ms: Option<u64>,
95    /// Harness message-record count, when available without loading the session.
96    /// STORED message-record count in the native store (cheap line scan) —
97    /// for branched record-tree dialects this can exceed the active-path
98    /// message count a full load renders; `inspect` labels both (SUP-58).
99    pub message_count: Option<usize>,
100    /// Model recorded in lightweight session metadata.
101    pub model: Option<String>,
102    /// Direct parent session for a harness-native child rollout. Ordinary
103    /// conversation lists exclude these children, while tree/fidelity callers
104    /// can request them explicitly without losing the native relationship.
105    #[serde(default, skip_serializing_if = "Option::is_none")]
106    pub parent_session_id: Option<String>,
107    /// Number of proven harness-native descendants represented by this root.
108    /// Child identities remain behind the trusted catalog boundary until a
109    /// caller explicitly requests this session family.
110    #[serde(default, skip_serializing_if = "is_zero")]
111    pub child_session_count: usize,
112    /// ORCH-6: the ORCH-3 conversation nouns (`trigger`, `surface`,
113    /// `profile`, `recurrence`, `cross_surface`, `workspace`), flattened onto
114    /// the row so the wire stays additive. Filled by `finalize_nouns` for
115    /// every harness; only Hermes and OpenClaw publish more than the
116    /// `trigger`/`workspace` defaults today.
117    #[serde(flatten)]
118    pub nouns: OrchestrationNouns,
119}
120
121/// One bounded normalized conversation-message candidate for list projection.
122#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
123pub struct SessionPreviewCandidate {
124    /// Opaque identity of this native message boundary. Consumers may retain
125    /// it to reconcile bounded discovery windows without treating a growing
126    /// preview or a native-store heartbeat as a new conversation message.
127    #[serde(default, skip_serializing_if = "Option::is_none")]
128    pub cursor: Option<String>,
129    /// Canonical conversation role. Older clients may assume `user` when this
130    /// field is absent from an older server.
131    pub role: String,
132    /// Canonical text content.
133    pub content: String,
134    /// Canonical provenance used by the normal conversation visibility rules.
135    #[serde(default, skip_serializing_if = "HashMap::is_empty")]
136    pub metadata: HashMap<String, String>,
137}
138
139/// One stable newest-first discovery page.
140#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
141pub struct DiscoveryPage {
142    /// Sessions in this page.
143    pub sessions: Vec<SessionDescriptor>,
144    /// Opaque cursor for the next page, or `None` at the end.
145    pub next_cursor: Option<String>,
146    /// Coverage receipt (UNI-7 / SUP-54): what was asked and what came back,
147    /// so incomplete coverage can never read as complete.
148    #[serde(default)]
149    pub receipt: DiscoveryReceipt,
150}
151
152/// Coverage evidence for one discovery page.
153#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
154#[serde(default)]
155pub struct DiscoveryReceipt {
156    /// True only when bounded topic/latest message search was applied before
157    /// pagination. Omitted for ordinary metadata-only discovery.
158    #[serde(skip_serializing_if = "is_false")]
159    pub searched_previews: bool,
160    /// The requested lower time bound (epoch ms), when one was given.
161    #[serde(skip_serializing_if = "Option::is_none")]
162    pub requested_after_ms: Option<u64>,
163    /// The requested upper time bound (epoch ms), when one was given.
164    #[serde(skip_serializing_if = "Option::is_none")]
165    pub requested_before_ms: Option<u64>,
166    /// The requested page limit, when one was given.
167    #[serde(skip_serializing_if = "Option::is_none")]
168    pub requested_limit: Option<usize>,
169    /// Oldest `updated_at_ms` among the RETURNED sessions.
170    #[serde(skip_serializing_if = "Option::is_none")]
171    pub oldest_returned_ms: Option<u64>,
172    /// Newest `updated_at_ms` among the RETURNED sessions.
173    #[serde(skip_serializing_if = "Option::is_none")]
174    pub newest_returned_ms: Option<u64>,
175    /// Sessions in this page.
176    pub returned: usize,
177    /// Sessions matching the query across ALL pages.
178    pub total_matched: usize,
179    /// True when matches beyond this page exist (`next_cursor` is the resume
180    /// point). Explicit so a consumer cannot mistake a capped page for the
181    /// complete result set.
182    pub truncated: bool,
183}
184
185/// Append-aware index of the stable topic records in Codex `history.jsonl`.
186///
187/// Long-lived session-list clients can retain this index and refresh it after
188/// filesystem invalidations. Ordinary appends read and parse only the new
189/// bytes; truncation, replacement, and in-place rewrites rebuild the index so
190/// the result remains identical to a fresh catalog discovery.
191#[derive(Debug)]
192pub struct CodexHistoryTopicIndex {
193    path: PathBuf,
194    fingerprint: Option<CodexHistoryFingerprint>,
195    offset: u64,
196    trailing: Vec<u8>,
197    topics: HashMap<String, Vec<SessionPreviewCandidate>>,
198}
199
200#[derive(Debug, Clone, Copy, PartialEq, Eq)]
201struct CodexHistoryFingerprint {
202    len: u64,
203    modified_ns: u128,
204    identity: u128,
205}
206
207impl CodexHistoryTopicIndex {
208    /// Create an empty index for the Codex sessions root from
209    /// [`HarnessHomes::codex`]. Call [`Self::refresh`] before first use.
210    pub fn new(sessions_root: &Path) -> Self {
211        let root = sessions_root.parent().unwrap_or(sessions_root);
212        Self {
213            path: root.join("history.jsonl"),
214            fingerprint: None,
215            offset: 0,
216            trailing: Vec::new(),
217            topics: HashMap::new(),
218        }
219    }
220
221    /// Return the native history file watched by this index.
222    pub fn path(&self) -> &Path {
223        &self.path
224    }
225
226    /// Refresh from durable state and return session ids whose topic changed.
227    ///
228    /// A missing file is a valid empty history. Read errors leave the previous
229    /// successful index intact so a transient filesystem error cannot erase
230    /// topics from a live session list.
231    pub fn refresh(&mut self) -> Result<BTreeSet<String>> {
232        let metadata = match fs::metadata(&self.path) {
233            Ok(metadata) => metadata,
234            Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
235                let changed = self.topics.keys().cloned().collect();
236                self.fingerprint = None;
237                self.offset = 0;
238                self.trailing.clear();
239                self.topics.clear();
240                return Ok(changed);
241            }
242            Err(error) => return Err(error.into()),
243        };
244        let fingerprint = codex_history_fingerprint(&metadata)?;
245        if self.fingerprint == Some(fingerprint) {
246            return Ok(BTreeSet::new());
247        }
248
249        let is_append = self.fingerprint.is_some_and(|previous| {
250            previous.identity == fingerprint.identity
251                && previous.len < fingerprint.len
252                && self.offset <= previous.len
253        });
254        if is_append {
255            let mut file = File::open(&self.path)?;
256            file.seek(SeekFrom::Start(self.offset))?;
257            let mut bytes = Vec::with_capacity(
258                usize::try_from(fingerprint.len.saturating_sub(self.offset)).unwrap_or(0),
259            );
260            file.read_to_end(&mut bytes)?;
261            self.offset = file.stream_position()?;
262            let changed = self.ingest(bytes);
263            self.fingerprint = Some(CodexHistoryFingerprint {
264                len: self.offset,
265                ..fingerprint
266            });
267            return Ok(changed);
268        }
269
270        let previous = std::mem::take(&mut self.topics);
271        let mut file = File::open(&self.path)?;
272        let mut bytes = Vec::with_capacity(usize::try_from(fingerprint.len).unwrap_or(0));
273        file.read_to_end(&mut bytes)?;
274        self.offset = file.stream_position()?;
275        self.trailing.clear();
276        self.ingest(bytes);
277        self.fingerprint = Some(CodexHistoryFingerprint {
278            len: self.offset,
279            ..fingerprint
280        });
281        Ok(changed_topic_ids(&previous, &self.topics))
282    }
283
284    fn ingest(&mut self, bytes: Vec<u8>) -> BTreeSet<String> {
285        let mut input = std::mem::take(&mut self.trailing);
286        input.extend(bytes);
287        let complete_len = input
288            .iter()
289            .rposition(|byte| *byte == b'\n')
290            .map_or(0, |index| index + 1);
291        let mut changed = BTreeSet::new();
292        for line in input[..complete_len].split(|byte| *byte == b'\n') {
293            self.ingest_line(line, &mut changed);
294        }
295        self.trailing.extend_from_slice(&input[complete_len..]);
296        if !self.trailing.is_empty() {
297            let trailing = self.trailing.clone();
298            if self.ingest_line(&trailing, &mut changed) {
299                self.trailing.clear();
300            }
301        }
302        changed
303    }
304
305    fn ingest_line(&mut self, line: &[u8], changed: &mut BTreeSet<String>) -> bool {
306        let Ok(value) = serde_json::from_slice::<Value>(line) else {
307            return false;
308        };
309        let Some(session_id) = value.get("session_id").and_then(Value::as_str) else {
310            return true;
311        };
312        if self.topics.contains_key(session_id) {
313            return true;
314        }
315        let mut candidates = Vec::new();
316        push_message_candidate(&mut candidates, "user", value.get("text"), HashMap::new());
317        if !candidates.is_empty() {
318            self.topics.insert(session_id.to_string(), candidates);
319            changed.insert(session_id.to_string());
320        }
321        true
322    }
323}
324
325/// Configurable session roots for the built-in harnesses.
326#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
327#[serde(default)]
328pub struct HarnessHomes {
329    /// Directory containing Claude Code project session directories.
330    pub claude_code: PathBuf,
331    /// Directory containing Codex rollout sessions.
332    pub codex: PathBuf,
333    /// Directory containing Pi project session directories.
334    pub pi: PathBuf,
335    /// OpenCode data root, or an explicit `opencode*.db` path.
336    pub opencode: PathBuf,
337    /// Grok session root containing percent-encoded workspace directories.
338    pub grok: PathBuf,
339    /// Gemini CLI configuration root containing `projects.json` and `tmp/`.
340    pub gemini: PathBuf,
341    /// Goose `sessions.db`, or a directory containing it.
342    pub goose: PathBuf,
343    /// Supercode's native saved-session directory.
344    pub supercode: PathBuf,
345    /// OpenClaw home root (contains `agents/<agentId>/sessions/*.jsonl`,
346    /// openclaw >= 2026.7 — plain pi-v3 dialect files; see
347    /// `SessionSource::OpenClaw`). READ-ONLY discovery (UNI-16).
348    pub openclaw: PathBuf,
349    /// Hermes `state.db` SQLite store path (see `SessionSource::Hermes`).
350    /// READ-ONLY discovery (UNI-15).
351    pub hermes: PathBuf,
352    /// The orchestrator's home FOLDER (`SUPERCODE_ORCHESTRATOR_HOME`, default
353    /// `~/.supercode/orchestrator`). Unlike every other entry this addresses a
354    /// directory, because the orchestrator's serialization IS its folder
355    /// (`docs/ORCHESTRATOR-IR.md` §6): the root is the `default` profile and
356    /// `profiles/<name>/` are the named ones. READ-ONLY (ORC-7).
357    pub orchestrator: PathBuf,
358}
359
360/// Every profile folder under an orchestrator home, in listing order: the
361/// root (the implicit `default` profile) then each `profiles/<name>/`.
362///
363/// One helper, used by every reader that answers `--harness orchestrator`, so
364/// the folder layout of `docs/ORCHESTRATOR-IR.md` §6 is stated once.
365pub fn orchestrator_profile_dirs(root: &Path) -> Vec<(String, PathBuf)> {
366    if !root.is_dir() {
367        return Vec::new();
368    }
369    let mut dirs = vec![("default".to_string(), root.to_path_buf())];
370    if let Ok(entries) = fs::read_dir(root.join("profiles")) {
371        let mut named: Vec<(String, PathBuf)> = entries
372            .flatten()
373            .filter(|entry| entry.path().is_dir())
374            .filter_map(|entry| {
375                entry
376                    .file_name()
377                    .into_string()
378                    .ok()
379                    .map(|name| (name, entry.path()))
380            })
381            // A dot-prefixed folder is never a profile: `profiles/.trash/`
382            // holds the homes `profiles.delete` moved aside (ORC-13), and the
383            // package's own loader skips it for the same reason.
384            .filter(|(name, _)| !name.starts_with('.'))
385            .collect();
386        named.sort();
387        dirs.extend(named);
388    }
389    dirs
390}
391
392/// A Hermes home's session stores: its own `state.db` and each named profile's (Hermes gives every profile
393/// its own home and store under `profiles/<name>/`; a linked profile folder counts). The same folder walk as
394/// the orchestrator's profiles, since both lay profiles out as Hermes does.
395pub fn hermes_session_stores(db_path: &Path) -> Vec<PathBuf> {
396    let mut stores = vec![db_path.to_path_buf()];
397    let (Some(root), Some(name)) = (db_path.parent(), db_path.file_name()) else {
398        return stores;
399    };
400    for (_, dir) in orchestrator_profile_dirs(root).into_iter().skip(1) {
401        stores.push(dir.join(name));
402    }
403    stores
404}
405
406/// The running Claude Code session whose registry name is exactly the query, as `(name, id)`.
407/// Claude keeps one record per running session beside its projects (`sessions/<pid>.json`).
408fn claude_live_named(query: &DiscoveryQuery) -> Option<(String, String)> {
409    let wanted = query
410        .query
411        .as_deref()
412        .map(str::trim)
413        .filter(|q| !q.is_empty())?;
414    if !query.harnesses.is_empty()
415        && !query
416            .harnesses
417            .iter()
418            .any(|harness| harness.as_str() == HarnessId::CLAUDE_CODE)
419    {
420        return None;
421    }
422    let registry = query.homes.claude_code.parent()?.join("sessions");
423    std::fs::read_dir(registry)
424        .ok()?
425        .flatten()
426        .find_map(|entry| {
427            let record: Value = serde_json::from_slice(&std::fs::read(entry.path()).ok()?).ok()?;
428            let name = record.get("name")?.as_str()?;
429            let id = record.get("sessionId")?.as_str()?;
430            (name == wanted).then(|| (name.to_string(), id.to_string()))
431        })
432}
433
434impl Default for HarnessHomes {
435    fn default() -> Self {
436        let home = crate::user_home()
437            .map(std::path::PathBuf::into_os_string)
438            .map(PathBuf::from)
439            .unwrap_or_else(|| PathBuf::from("."));
440        let claude_root = std::env::var_os("CLAUDE_CONFIG_DIR")
441            .map(PathBuf::from)
442            .unwrap_or_else(|| home.join(".claude"));
443        let codex_root = std::env::var_os("CODEX_HOME")
444            .map(PathBuf::from)
445            .unwrap_or_else(|| home.join(".codex"));
446        let pi = std::env::var_os("PI_CODING_AGENT_SESSION_DIR")
447            .map(PathBuf::from)
448            .unwrap_or_else(|| {
449                std::env::var_os("PI_CODING_AGENT_DIR")
450                    .map(PathBuf::from)
451                    .unwrap_or_else(|| home.join(".pi/agent"))
452                    .join("sessions")
453            });
454        let opencode = std::env::var_os("OPENCODE_DB")
455            .map(PathBuf::from)
456            .unwrap_or_else(|| {
457                std::env::var_os("XDG_DATA_HOME")
458                    .map(PathBuf::from)
459                    .unwrap_or_else(|| home.join(".local/share"))
460                    .join("opencode")
461            });
462        let grok = std::env::var_os("GROK_HOME")
463            .map(PathBuf::from)
464            .unwrap_or_else(|| home.join(".grok"))
465            .join("sessions");
466        let gemini = std::env::var_os("GEMINI_CLI_HOME")
467            .map(PathBuf::from)
468            .unwrap_or_else(|| home.join(".gemini"));
469        // Blind-walk finding 2026-08-31: openclaw itself treats
470        // OPENCLAW_HOME as a HOME replacement (state lives at
471        // `$OPENCLAW_HOME/.openclaw`), and names the state dir directly with
472        // OPENCLAW_STATE_DIR. Mirror those semantics exactly so an inherited
473        // environment means the same thing to us and to any openclaw process
474        // we spawn.
475        let openclaw = std::env::var_os("OPENCLAW_STATE_DIR")
476            .map(PathBuf::from)
477            .or_else(|| {
478                std::env::var_os("OPENCLAW_HOME").map(|root| PathBuf::from(root).join(".openclaw"))
479            })
480            .unwrap_or_else(|| home.join(".openclaw"));
481        let orchestrator = std::env::var_os("SUPERCODE_ORCHESTRATOR_HOME")
482            .map(PathBuf::from)
483            .unwrap_or_else(|| home.join(".supercode/orchestrator"));
484        let hermes = std::env::var_os("HERMES_HOME")
485            .map(PathBuf::from)
486            .unwrap_or_else(|| home.join(".hermes"))
487            .join("state.db");
488        let goose = std::env::var_os("GOOSE_PATH_ROOT")
489            .map(PathBuf::from)
490            .map(|root| root.join("data/sessions/sessions.db"))
491            .unwrap_or_else(|| {
492                #[cfg(target_os = "macos")]
493                {
494                    home.join("Library/Application Support/Block/goose/sessions/sessions.db")
495                }
496                #[cfg(target_os = "windows")]
497                {
498                    std::env::var_os("APPDATA")
499                        .map(PathBuf::from)
500                        .unwrap_or_else(|| home.join("AppData/Roaming"))
501                        .join("Block/goose/sessions/sessions.db")
502                }
503                #[cfg(not(any(target_os = "macos", target_os = "windows")))]
504                {
505                    std::env::var_os("XDG_DATA_HOME")
506                        .map(PathBuf::from)
507                        .unwrap_or_else(|| home.join(".local/share"))
508                        .join("goose/sessions/sessions.db")
509                }
510            });
511        let supercode = std::env::var_os("SUPERCODE_HOME")
512            .map(PathBuf::from)
513            .unwrap_or_else(|| {
514                std::env::var_os("XDG_CONFIG_HOME")
515                    .map(PathBuf::from)
516                    .unwrap_or_else(|| home.join(".config"))
517                    .join("supercode")
518            })
519            .join("sessions");
520        Self {
521            claude_code: claude_root.join("projects"),
522            codex: codex_root.join("sessions"),
523            gemini,
524            goose,
525            supercode,
526            openclaw,
527            hermes,
528            orchestrator,
529            pi,
530            opencode,
531            grok,
532        }
533    }
534}
535
536fn codex_history_fingerprint(metadata: &fs::Metadata) -> Result<CodexHistoryFingerprint> {
537    let modified_ns = metadata
538        .modified()?
539        .duration_since(UNIX_EPOCH)
540        .map_err(|error| Error::Other(format!("history timestamp predates Unix epoch: {error}")))?
541        .as_nanos();
542    #[cfg(unix)]
543    let identity = {
544        use std::os::unix::fs::MetadataExt;
545        (u128::from(metadata.dev()) << 64) | u128::from(metadata.ino())
546    };
547    #[cfg(not(unix))]
548    let identity = 0;
549    Ok(CodexHistoryFingerprint {
550        len: metadata.len(),
551        modified_ns,
552        identity,
553    })
554}
555
556fn changed_topic_ids(
557    before: &HashMap<String, Vec<SessionPreviewCandidate>>,
558    after: &HashMap<String, Vec<SessionPreviewCandidate>>,
559) -> BTreeSet<String> {
560    before
561        .keys()
562        .chain(after.keys())
563        .filter(|session_id| before.get(*session_id) != after.get(*session_id))
564        .cloned()
565        .collect()
566}
567
568/// Filters and roots used for one catalog scan.
569#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
570#[serde(default)]
571pub struct DiscoveryQuery {
572    /// Only return sessions whose recorded working directory is this path.
573    pub workspace: Option<PathBuf>,
574    /// Widen `workspace` from that exact folder to the folder and every folder
575    /// under it (a Teams sync rule's root). Matched the way the exact filter
576    /// matches: only an absolute recorded cwd qualifies, both sides are
577    /// canonicalized (symlinks, macOS `/var` vs `/private/var`, Windows `\\?\`
578    /// prefixes), and containment is per path component, so `/work` does not
579    /// admit `/workshop`. Applied before any per-session enrichment. Ignored
580    /// without `workspace`; absent (old clients) means the exact match.
581    #[serde(skip_serializing_if = "std::ops::Not::not")]
582    pub workspace_subtree: bool,
583    /// Only return sessions whose recorded working directory belongs to the
584    /// same REPOSITORY FAMILY as this path (UNI-7 / SUP-54 decision): the
585    /// family identity is the realpath of the git common directory, so a main
586    /// checkout and all its worktrees join as one family; the origin URL is
587    /// the tie-breaker that also joins separate clones of the same remote. A
588    /// non-repository path degrades to exact-realpath matching. Combines with
589    /// `workspace` as AND when both are set (exact-cwd stays available).
590    pub workspace_family: Option<PathBuf>,
591    /// Only return sessions updated at or after this epoch-ms instant.
592    pub updated_after_ms: Option<u64>,
593    /// Only return sessions updated at or before this epoch-ms instant.
594    pub updated_before_ms: Option<u64>,
595    /// Harnesses to scan. Empty means all built-ins.
596    pub harnesses: Vec<HarnessId>,
597    /// Storage roots to scan.
598    pub homes: HarnessHomes,
599    /// Case-insensitive search over harness, id, title, workspace, and model.
600    pub query: Option<String>,
601    /// Also match bounded opening/latest message candidates. Explicitly opt-in:
602    /// this scans previews of eligible sessions before pagination, not just the
603    /// returned page. Requires a nonempty query; not a full-transcript search.
604    pub search_previews: bool,
605    /// Opaque cursor returned by a prior [`HarnessCatalog::discover_page`].
606    pub cursor: Option<String>,
607    /// Maximum number of results after newest-first sorting.
608    pub limit: Option<usize>,
609    /// Include oldest-first bounded topic candidates for harnesses whose
610    /// native store does not publish a useful title. Off by default because
611    /// topics are stable and list clients can retain them across refreshes.
612    pub include_topic_candidates: bool,
613    /// Include harness-native child rollouts such as Codex subagents. Off by
614    /// default because they are parts of a parent conversation, not chats the
615    /// user independently started. Translation/tree callers can opt in.
616    pub include_child_sessions: bool,
617    /// Restrict an explicit child-inclusive discovery to one root and every
618    /// descendant linked to it by native lineage. Applied before pagination.
619    pub root_session_id: Option<String>,
620    /// ORCH-6: only return sessions routed through this config home (Hermes
621    /// `profile_name` / gateway key namespace, OpenClaw agent id). Exact
622    /// match; a session with no profile never matches. `harnesses` is the
623    /// harness filter and needs no second spelling.
624    pub profile: Option<String>,
625}
626
627/// Repository-family identity (UNI-7 / SUP-54 decision): the realpath of the
628/// git common directory joins a main checkout with all its worktrees, and the
629/// origin URL is the tie-breaker that also joins separate clones of the same
630/// remote. Resolved from the filesystem alone (`.git` file/dir + config), no
631/// subprocess, so discovery stays deterministic and sandbox-friendly.
632#[derive(Debug, Clone, PartialEq, Eq)]
633struct RepoFamily {
634    /// Realpath of the git common directory, or the realpath of the queried
635    /// path itself when it is not inside a git repository.
636    identity: PathBuf,
637    /// True when `identity` is a git common directory (not a bare-path
638    /// degradation) — only then may origin tie-breaking apply.
639    is_repository: bool,
640    /// `[remote "origin"] url` from the common directory's config, if any.
641    origin_url: Option<String>,
642}
643
644impl RepoFamily {
645    fn of(path: &Path) -> Self {
646        let start = std::fs::canonicalize(path).unwrap_or_else(|_| path.to_path_buf());
647        let mut current = Some(start.as_path());
648        while let Some(dir) = current {
649            let dot_git = dir.join(".git");
650            if dot_git.is_dir() {
651                let identity = std::fs::canonicalize(&dot_git).unwrap_or_else(|_| dot_git.clone());
652                let origin_url = read_origin_url(&identity);
653                return Self {
654                    identity,
655                    is_repository: true,
656                    origin_url,
657                };
658            }
659            if dot_git.is_file() {
660                // Worktree checkout: `.git` is a one-line pointer file
661                // `gitdir: <main>/.git/worktrees/<name>`; the common directory
662                // is everything before `/worktrees/<name>`.
663                if let Ok(text) = std::fs::read_to_string(&dot_git) {
664                    if let Some(gitdir) = text
665                        .lines()
666                        .find_map(|line| line.trim().strip_prefix("gitdir:"))
667                    {
668                        let gitdir = PathBuf::from(gitdir.trim());
669                        let gitdir = if gitdir.is_absolute() {
670                            gitdir
671                        } else {
672                            dir.join(gitdir)
673                        };
674                        let common = gitdir
675                            .parent()
676                            .filter(|parent| parent.ends_with("worktrees"))
677                            .and_then(Path::parent)
678                            .map(Path::to_path_buf)
679                            .unwrap_or(gitdir);
680                        let identity = std::fs::canonicalize(&common).unwrap_or(common);
681                        let origin_url = read_origin_url(&identity);
682                        return Self {
683                            identity,
684                            is_repository: true,
685                            origin_url,
686                        };
687                    }
688                }
689            }
690            current = dir.parent();
691        }
692        Self {
693            identity: start,
694            is_repository: false,
695            origin_url: None,
696        }
697    }
698
699    /// Two paths join one family when their common directories match, or —
700    /// for genuine repositories only — when both declare the same origin URL
701    /// (the clone tie-breaker). Bare-path degradations never origin-match.
702    fn joins(&self, other: &Self) -> bool {
703        if self.identity == other.identity {
704            return true;
705        }
706        self.is_repository
707            && other.is_repository
708            && matches!((&self.origin_url, &other.origin_url), (Some(a), Some(b)) if a == b)
709    }
710}
711
712/// Minimal git-config scan for `[remote "origin"] url = ...`.
713fn read_origin_url(common_dir: &Path) -> Option<String> {
714    let text = std::fs::read_to_string(common_dir.join("config")).ok()?;
715    let mut in_origin = false;
716    for line in text.lines() {
717        let line = line.trim();
718        if line.starts_with('[') {
719            in_origin = line.starts_with("[remote \"origin\"]");
720            continue;
721        }
722        if in_origin {
723            if let Some(value) = line.strip_prefix("url") {
724                let value = value.trim_start();
725                if let Some(url) = value.strip_prefix('=') {
726                    let url = url.trim();
727                    if !url.is_empty() {
728                        return Some(url.to_string());
729                    }
730                }
731            }
732        }
733    }
734    None
735}
736
737/// Read-only entry point for discovering, loading, and following persisted
738/// harness sessions.
739#[derive(Debug, Default, Clone, Copy)]
740pub struct HarnessCatalog;
741
742impl HarnessCatalog {
743    /// Construct a catalog. It holds no cache or global mutable state.
744    pub fn new() -> Self {
745        Self
746    }
747
748    /// Discover sessions using lightweight headers/indexes rather than full
749    /// transcript normalization. Malformed or concurrently-created entries
750    /// are skipped without aborting the rest of the scan.
751    pub fn discover(&self, query: &DiscoveryQuery) -> Result<Vec<SessionDescriptor>> {
752        Ok(self.discover_page(query)?.sessions)
753    }
754
755    /// Scan file-backed session metadata without pagination or conversation
756    /// previews. Child sessions are retained so a long-lived caller can keep
757    /// an exact in-memory lineage index and project it without rescanning the
758    /// native stores.
759    pub fn discover_raw_index(&self, query: &DiscoveryQuery) -> Vec<SessionDescriptor> {
760        self.scan_descriptors(query, true)
761    }
762
763    /// Project a raw metadata index through the query's lineage, search, and
764    /// pagination rules without reading any transcript content.
765    pub fn project_index(
766        &self,
767        query: &DiscoveryQuery,
768        descriptors: impl IntoIterator<Item = SessionDescriptor>,
769    ) -> Result<Vec<SessionDescriptor>> {
770        Ok(self.project_index_page(query, descriptors)?.sessions)
771    }
772
773    /// Project a retained metadata index without reading transcripts, preserving
774    /// the full-match count and successor cursor of ordinary discovery.
775    pub fn project_index_page(
776        &self,
777        query: &DiscoveryQuery,
778        descriptors: impl IntoIterator<Item = SessionDescriptor>,
779    ) -> Result<DiscoveryPage> {
780        if query.search_previews {
781            return Err(Error::Other(
782                "preview search requires discover_page, not a metadata-only index projection"
783                    .into(),
784            ));
785        }
786        let mut found = descriptors.into_iter().collect::<Vec<_>>();
787        project_descriptors(query, &mut found);
788        let total_matched = found.len();
789        let (sessions, next_cursor) = paginate_descriptors(query, found)?;
790        let receipt = DiscoveryReceipt {
791            searched_previews: false,
792            requested_after_ms: query.updated_after_ms,
793            requested_before_ms: query.updated_before_ms,
794            requested_limit: query.limit,
795            oldest_returned_ms: sessions.iter().filter_map(|s| s.updated_at_ms).min(),
796            newest_returned_ms: sessions.iter().filter_map(|s| s.updated_at_ms).max(),
797            returned: sessions.len(),
798            total_matched,
799            truncated: next_cursor.is_some(),
800        };
801        Ok(DiscoveryPage {
802            sessions,
803            next_cursor,
804            receipt,
805        })
806    }
807
808    /// Add the bounded topic/latest-message previews used by list clients to
809    /// an already projected metadata page.
810    pub fn enrich_index_page(
811        &self,
812        query: &DiscoveryQuery,
813        mut sessions: Vec<SessionDescriptor>,
814    ) -> Result<Vec<SessionDescriptor>> {
815        enrich_descriptors(query, &mut sessions, None)?;
816        Ok(sessions)
817    }
818
819    /// Add list previews using a retained append-aware Codex history index.
820    ///
821    /// Results are identical to [`Self::enrich_index_page`], while a
822    /// long-lived caller avoids reparsing all of `history.jsonl` for every
823    /// active transcript write.
824    pub fn enrich_index_page_with_codex_history(
825        &self,
826        query: &DiscoveryQuery,
827        mut sessions: Vec<SessionDescriptor>,
828        codex_history: &CodexHistoryTopicIndex,
829    ) -> Result<Vec<SessionDescriptor>> {
830        enrich_descriptors(query, &mut sessions, Some(codex_history))?;
831        Ok(sessions)
832    }
833
834    /// Discover one stable page and return the cursor for its successor.
835    pub fn discover_page(&self, query: &DiscoveryQuery) -> Result<DiscoveryPage> {
836        if query.search_previews {
837            return self.discover_preview_page(query);
838        }
839        let mut page = self.project_index_page(
840            query,
841            self.scan_descriptors(query, query.include_child_sessions),
842        )?;
843        // A query that is exactly a running Claude session's name (`claude -n`, `/rename`, kept in
844        // Claude's session registry) names that session: it alone answers, under that name, not
845        // every session whose folder or text shares the words.
846        if query.cursor.is_none() {
847            if let Some((name, id)) = claude_live_named(query) {
848                if let Some(at) = page
849                    .sessions
850                    .iter()
851                    .position(|session| session.locator.session_id == id)
852                {
853                    let mut found = page.sessions.swap_remove(at);
854                    found.title.get_or_insert(name);
855                    page.sessions = vec![found];
856                    page.next_cursor = None;
857                }
858            }
859        }
860        // A query that is a session's whole id names that session: not the others whose folders
861        // or text carry the id (a session's scratch folders are named after it).
862        if let Some(id) = query.query.as_deref().map(str::trim) {
863            if page
864                .sessions
865                .iter()
866                .any(|session| session.locator.session_id == id)
867            {
868                page.sessions
869                    .retain(|session| session.locator.session_id == id);
870                page.next_cursor = None;
871            }
872        }
873        if page.sessions.is_empty() && query.cursor.is_none() {
874            if let Some(named) = self.claude_session_named(query)? {
875                page.sessions.push(named);
876            }
877        }
878        enrich_descriptors(query, &mut page.sessions, None)?;
879        Ok(page)
880    }
881
882    /// The Claude Code session a query names exactly, when nothing else matched
883    /// it. A name given on resume (`claude --resume <id> -n <name>`, `/rename`)
884    /// is a record appended wherever the conversation had got to, past the
885    /// header discovery reads; so, and only on a miss, the transcripts are
886    /// searched for the record itself, and the most recently written one that
887    /// carries it wins. The query's other filters still apply.
888    fn claude_session_named(&self, query: &DiscoveryQuery) -> Result<Option<SessionDescriptor>> {
889        let Some(name) = query
890            .query
891            .as_deref()
892            .map(str::trim)
893            .filter(|name| !name.is_empty())
894        else {
895            return Ok(None);
896        };
897        if !query.harnesses.is_empty()
898            && !query
899                .harnesses
900                .iter()
901                .any(|harness| harness.as_str() == HarnessId::CLAUDE_CODE)
902        {
903            return Ok(None);
904        }
905        let needle = format!("\"customTitle\":{}", Value::String(name.to_string()));
906        let mut best: Option<(std::time::SystemTime, String)> = None;
907        let Ok(projects) = std::fs::read_dir(&query.homes.claude_code) else {
908            return Ok(None);
909        };
910        for project in projects.flatten() {
911            let Ok(files) = std::fs::read_dir(project.path()) else {
912                continue;
913            };
914            for file in files.flatten() {
915                let path = file.path();
916                if path.extension().and_then(|value| value.to_str()) != Some("jsonl") {
917                    continue;
918                }
919                let Ok(bytes) = std::fs::read(&path) else {
920                    continue;
921                };
922                if !bytes
923                    .windows(needle.len())
924                    .any(|window| window == needle.as_bytes())
925                {
926                    continue;
927                }
928                let modified = file
929                    .metadata()
930                    .and_then(|meta| meta.modified())
931                    .unwrap_or(std::time::UNIX_EPOCH);
932                let Some(id) = path
933                    .file_stem()
934                    .map(|stem| stem.to_string_lossy().into_owned())
935                else {
936                    continue;
937                };
938                if best.as_ref().is_none_or(|(time, _)| modified > *time) {
939                    best = Some((modified, id));
940                }
941            }
942        }
943        let Some((_, id)) = best else {
944            return Ok(None);
945        };
946        let mut unnamed = query.clone();
947        unnamed.query = None;
948        unnamed.harnesses = vec![HarnessId::new(HarnessId::CLAUDE_CODE)];
949        unnamed.limit = None;
950        let found = self
951            .project_index_page(
952                &unnamed,
953                self.scan_descriptors(&unnamed, unnamed.include_child_sessions),
954            )?
955            .sessions
956            .into_iter()
957            .find(|descriptor| descriptor.locator.session_id == id);
958        Ok(found.map(|mut descriptor| {
959            descriptor.title = Some(name.to_string());
960            descriptor
961        }))
962    }
963
964    fn discover_preview_page(&self, query: &DiscoveryQuery) -> Result<DiscoveryPage> {
965        let search = query
966            .query
967            .as_deref()
968            .map(str::trim)
969            .filter(|text| !text.is_empty())
970            .ok_or_else(|| Error::Other("search_previews requires a nonempty query".into()))?
971            .to_lowercase();
972        if query.limit == Some(0) {
973            return Err(Error::Other("preview search limit must be positive".into()));
974        }
975        let cursor = query.cursor.as_deref().map(decode_cursor).transpose()?;
976        let mut eligible_query = query.clone();
977        eligible_query.query = None;
978        eligible_query.search_previews = false;
979        eligible_query.include_topic_candidates = true;
980        let mut eligible = self.scan_descriptors(query, query.include_child_sessions);
981        project_descriptors(&eligible_query, &mut eligible);
982        // Read Codex's first history topic once per scan, never once per row.
983        let topics = codex_history_topics(&query.homes.codex, &eligible).unwrap_or_default();
984        let mut sessions = Vec::new();
985        let mut total_matched = 0;
986        let mut cursor_seen = cursor.is_none();
987        let mut more = false;
988        for mut descriptor in eligible {
989            let metadata_match = descriptor_matches(&descriptor, &search);
990            if !metadata_match {
991                let topic = topics
992                    .get(&descriptor.locator.session_id)
993                    .filter(|_| descriptor.locator.harness.as_str() == HarnessId::CODEX);
994                enrich_descriptor(&eligible_query, &mut descriptor, topic);
995                if !descriptor
996                    .preview_candidates
997                    .iter()
998                    .chain(&descriptor.latest_message_candidates)
999                    .any(|candidate| candidate.content.to_lowercase().contains(&search))
1000                {
1001                    continue;
1002                }
1003            }
1004            total_matched += 1;
1005            if !cursor_seen {
1006                cursor_seen = cursor.as_ref() == Some(&descriptor_cursor_key(&descriptor));
1007                continue;
1008            }
1009            if sessions.len() < query.limit.unwrap_or(usize::MAX) {
1010                if metadata_match {
1011                    let topic = topics
1012                        .get(&descriptor.locator.session_id)
1013                        .filter(|_| descriptor.locator.harness.as_str() == HarnessId::CODEX);
1014                    enrich_descriptor(&eligible_query, &mut descriptor, topic);
1015                }
1016                sessions.push(descriptor);
1017            } else {
1018                more = true;
1019            }
1020        }
1021        if !cursor_seen {
1022            return Err(Error::Other("discovery cursor is stale or invalid".into()));
1023        }
1024        let next_cursor = more.then(|| sessions.last().map(encode_cursor)).flatten();
1025        let receipt = DiscoveryReceipt {
1026            searched_previews: true,
1027            requested_after_ms: query.updated_after_ms,
1028            requested_before_ms: query.updated_before_ms,
1029            requested_limit: query.limit,
1030            oldest_returned_ms: sessions.iter().filter_map(|s| s.updated_at_ms).min(),
1031            newest_returned_ms: sessions.iter().filter_map(|s| s.updated_at_ms).max(),
1032            returned: sessions.len(),
1033            total_matched,
1034            truncated: more,
1035        };
1036        Ok(DiscoveryPage {
1037            sessions,
1038            next_cursor,
1039            receipt,
1040        })
1041    }
1042
1043    fn scan_descriptors(
1044        &self,
1045        query: &DiscoveryQuery,
1046        include_child_sessions: bool,
1047    ) -> Vec<SessionDescriptor> {
1048        let selected: HashSet<&str> = if query.harnesses.is_empty() {
1049            [
1050                HarnessId::CLAUDE_CODE,
1051                HarnessId::CODEX,
1052                HarnessId::PI,
1053                HarnessId::OPENCODE,
1054                HarnessId::GROK,
1055                HarnessId::GEMINI,
1056                HarnessId::GOOSE,
1057                HarnessId::SUPERCODE,
1058                HarnessId::OPENCLAW,
1059                HarnessId::HERMES,
1060                HarnessId::ORCHESTRATOR,
1061            ]
1062            .into_iter()
1063            .collect()
1064        } else {
1065            query.harnesses.iter().map(HarnessId::as_str).collect()
1066        };
1067        // Built once per scan: a subtree scope canonicalizes its root once and
1068        // remembers each recorded cwd it has already judged.
1069        let scope = WorkspaceScope::of(query);
1070        let mut found = Vec::new();
1071        if selected.contains(HarnessId::CLAUDE_CODE) {
1072            discover_jsonl(
1073                &query.homes.claude_code,
1074                HarnessId::CLAUDE_CODE,
1075                scope.as_ref(),
1076                include_child_sessions,
1077                &mut found,
1078            );
1079        }
1080        if selected.contains(HarnessId::CODEX) {
1081            discover_jsonl(
1082                &query.homes.codex,
1083                HarnessId::CODEX,
1084                scope.as_ref(),
1085                include_child_sessions,
1086                &mut found,
1087            );
1088        }
1089        if selected.contains(HarnessId::PI) {
1090            discover_jsonl(
1091                &query.homes.pi,
1092                HarnessId::PI,
1093                scope.as_ref(),
1094                include_child_sessions,
1095                &mut found,
1096            );
1097        }
1098        if selected.contains(HarnessId::OPENCODE) {
1099            discover_opencode(&query.homes.opencode, scope.as_ref(), &mut found);
1100        }
1101        if selected.contains(HarnessId::GROK) {
1102            discover_grok(&query.homes.grok, scope.as_ref(), &mut found);
1103        }
1104        if selected.contains(HarnessId::GEMINI) {
1105            discover_gemini(&query.homes.gemini, scope.as_ref(), &mut found);
1106        }
1107        if selected.contains(HarnessId::GOOSE) {
1108            discover_goose(&query.homes.goose, scope.as_ref(), &mut found);
1109        }
1110        if selected.contains(HarnessId::OPENCLAW) {
1111            discover_openclaw(&query.homes.openclaw, scope.as_ref(), &mut found);
1112        }
1113        if selected.contains(HarnessId::HERMES) {
1114            for store in hermes_session_stores(&query.homes.hermes) {
1115                discover_hermes(&store, scope.as_ref(), &selected, &mut found);
1116            }
1117        }
1118        if selected.contains(HarnessId::ORCHESTRATOR) {
1119            discover_orchestrator(&query.homes.orchestrator, scope.as_ref(), &mut found);
1120        }
1121        if selected.contains(HarnessId::SUPERCODE) {
1122            discover_supercode(&query.homes.supercode, scope.as_ref(), &mut found);
1123        }
1124        for descriptor in &mut found {
1125            finalize_nouns(descriptor);
1126        }
1127        found
1128    }
1129
1130    /// Refresh one file-backed descriptor without rescanning its native store.
1131    ///
1132    /// This is the incremental counterpart to [`Self::discover_page`]: a
1133    /// filesystem notification is only an invalidation hint, so callers
1134    /// re-read the durable file and derive the complete current descriptor.
1135    /// `None` means the path disappeared or no longer contains a recognizable
1136    /// session. SQLite-backed harnesses retain their indexed discovery path.
1137    pub fn refresh_file_descriptor(
1138        &self,
1139        locator: &SessionLocator,
1140        workspace: Option<&Path>,
1141        include_topic_candidates: bool,
1142    ) -> Result<Option<SessionDescriptor>> {
1143        let Some(mut descriptor) = self.refresh_file_index_descriptor(locator, workspace)? else {
1144            return Ok(None);
1145        };
1146        if include_topic_candidates {
1147            descriptor.preview_candidates =
1148                topic_message_candidates(&descriptor.locator).unwrap_or_default();
1149        }
1150        descriptor.latest_message_candidates =
1151            latest_message_candidates(&descriptor.locator).unwrap_or_default();
1152        Ok(Some(descriptor))
1153    }
1154
1155    /// Refresh only stable list metadata for one file-backed descriptor.
1156    /// This avoids reading conversation preview windows for background index
1157    /// maintenance.
1158    pub fn refresh_file_index_descriptor(
1159        &self,
1160        locator: &SessionLocator,
1161        workspace: Option<&Path>,
1162    ) -> Result<Option<SessionDescriptor>> {
1163        let scope = workspace.map(WorkspaceScope::exact);
1164        self.refresh_file_index_descriptor_scoped(locator, scope.as_ref())
1165    }
1166
1167    /// [`Self::refresh_file_index_descriptor`] under a whole query's workspace
1168    /// filter, so a retained index honours `workspace_subtree` exactly as
1169    /// discovery does.
1170    pub fn refresh_file_index_descriptor_for(
1171        &self,
1172        locator: &SessionLocator,
1173        query: &DiscoveryQuery,
1174    ) -> Result<Option<SessionDescriptor>> {
1175        let scope = WorkspaceScope::of(query);
1176        self.refresh_file_index_descriptor_scoped(locator, scope.as_ref())
1177    }
1178
1179    fn refresh_file_index_descriptor_scoped(
1180        &self,
1181        locator: &SessionLocator,
1182        workspace: Option<&WorkspaceScope>,
1183    ) -> Result<Option<SessionDescriptor>> {
1184        let StorageLocator::File { path } = &locator.storage else {
1185            return Ok(None);
1186        };
1187        if !matches!(
1188            locator.harness.as_str(),
1189            HarnessId::CLAUDE_CODE | HarnessId::CODEX
1190        ) {
1191            return Ok(None);
1192        }
1193        if !path.is_file() {
1194            return Ok(None);
1195        }
1196        let Ok(meta) = read_header(path, locator.harness.as_str()) else {
1197            // Harnesses append the header and first turn non-atomically. A
1198            // transiently incomplete new file is not a service error; the
1199            // next native event or reconciliation pass will retry it.
1200            return Ok(None);
1201        };
1202        if workspace.is_some_and(|wanted| meta.cwd.as_deref().is_none_or(|cwd| !wanted.admits(cwd)))
1203        {
1204            return Ok(None);
1205        }
1206        let parent_session_id = meta.parent_session_id.or_else(|| {
1207            (locator.harness.as_str() == HarnessId::CLAUDE_CODE)
1208                .then(|| claude_subagent_parent_id(path))
1209                .flatten()
1210        });
1211        let tail = tail_facts(path, locator.harness.as_str());
1212        let descriptor = SessionDescriptor {
1213            locator: SessionLocator {
1214                harness: locator.harness.clone(),
1215                session_id: meta
1216                    .session_id
1217                    .unwrap_or_else(|| locator.session_id.clone()),
1218                storage: StorageLocator::File { path: path.clone() },
1219            },
1220            cwd: meta.cwd,
1221            title: meta.title,
1222            preview_candidates: Vec::new(),
1223            latest_message_candidates: Vec::new(),
1224            updated_at_ms: tail.last_turn_ms.or_else(|| modified_ms(path)),
1225            message_count: None,
1226            model: tail.model.or(meta.model),
1227            parent_session_id,
1228            child_session_count: 0,
1229            nouns: OrchestrationNouns::default(),
1230        };
1231        Ok(Some(descriptor))
1232    }
1233
1234    /// Load the complete normalized session named by a durable locator.
1235    pub fn load(&self, locator: &SessionLocator) -> Result<Session> {
1236        self.load_with_fidelity(locator, Fidelity::ByteLossless)
1237    }
1238
1239    /// [`Self::load`] at a declared fidelity.
1240    ///
1241    /// Read-only surfaces (a session mirror, `follow`) pass
1242    /// [`Fidelity::Semantic`] so a compacted transcript renders instead of
1243    /// erroring; every continuation/transfer/export caller keeps the strict
1244    /// default. See [`Session::load_with_fidelity`].
1245    pub fn load_with_fidelity(
1246        &self,
1247        locator: &SessionLocator,
1248        fidelity: Fidelity,
1249    ) -> Result<Session> {
1250        if let Some(session) = load_hermes_locator(locator) {
1251            return session;
1252        }
1253        match &locator.storage {
1254            StorageLocator::File { path } => {
1255                if let Some(session) = load_native_store_family(path)? {
1256                    Ok(session)
1257                } else {
1258                    Ok(Session::load_with_fidelity(path, fidelity)?)
1259                }
1260            }
1261            StorageLocator::Sqlite { path, selector } => {
1262                if locator.harness.as_str() == HarnessId::GOOSE {
1263                    Ok(Session::from_goose_sqlite(path, selector)?)
1264                } else {
1265                    Ok(Session::from_opencode_sqlite(path, Some(selector))?)
1266                }
1267            }
1268        }
1269    }
1270
1271    /// Load the selected parent transcript without recursively attaching
1272    /// Claude Code child sessions. This is the bounded frontend-view seam;
1273    /// lossless operations continue to use [`Self::load_with_fidelity`].
1274    #[doc(hidden)]
1275    pub fn load_parent_with_fidelity(
1276        &self,
1277        locator: &SessionLocator,
1278        fidelity: Fidelity,
1279    ) -> Result<Session> {
1280        if let Some(session) = load_hermes_locator(locator) {
1281            return session;
1282        }
1283        match &locator.storage {
1284            StorageLocator::File { path } => {
1285                if let Some(session) = load_native_store_family(path)? {
1286                    Ok(session)
1287                } else {
1288                    Ok(Session::load_parent_with_fidelity(path, fidelity)?)
1289                }
1290            }
1291            StorageLocator::Sqlite { path, selector } => {
1292                if locator.harness.as_str() == HarnessId::GOOSE {
1293                    Ok(Session::from_goose_sqlite(path, selector)?)
1294                } else {
1295                    Ok(Session::from_opencode_sqlite(path, Some(selector))?)
1296                }
1297            }
1298        }
1299    }
1300
1301    /// Load bounded parent-only human-visible history. Codex compaction
1302    /// changes resumable context but does not erase earlier visible turns.
1303    #[doc(hidden)]
1304    pub fn load_display_view(
1305        &self,
1306        locator: &SessionLocator,
1307        fidelity: Fidelity,
1308        message_limit: usize,
1309    ) -> Result<Session> {
1310        if let Some(session) = load_hermes_locator(locator) {
1311            let mut session = session?;
1312            if session.messages.len() > message_limit.max(1) {
1313                session
1314                    .messages
1315                    .drain(..session.messages.len() - message_limit.max(1));
1316            }
1317            return Ok(session);
1318        }
1319        match &locator.storage {
1320            StorageLocator::File { path } => {
1321                if let Some(mut session) = load_native_store_family(path)? {
1322                    if session.messages.len() > message_limit.max(1) {
1323                        session
1324                            .messages
1325                            .drain(..session.messages.len() - message_limit.max(1));
1326                    }
1327                    Ok(session)
1328                } else {
1329                    Ok(Session::load_display_view(path, fidelity, message_limit)?)
1330                }
1331            }
1332            StorageLocator::Sqlite { path, selector } => {
1333                let mut session = if locator.harness.as_str() == HarnessId::GOOSE {
1334                    Session::from_goose_sqlite_display(path, selector, message_limit)?
1335                } else {
1336                    Session::from_opencode_sqlite(path, Some(selector))?
1337                };
1338                if session.messages.len() > message_limit.max(1) {
1339                    session
1340                        .messages
1341                        .drain(..session.messages.len() - message_limit.max(1));
1342                }
1343                Ok(session)
1344            }
1345        }
1346    }
1347
1348    /// Open a passive change-triggered follower for a durable locator.
1349    pub fn follow(&self, locator: &SessionLocator) -> Result<SessionFollower> {
1350        self.follow_with_fidelity(locator, Fidelity::ByteLossless)
1351    }
1352
1353    /// [`Self::follow`] at a declared fidelity — see [`Self::load_with_fidelity`].
1354    pub fn follow_with_fidelity(
1355        &self,
1356        locator: &SessionLocator,
1357        fidelity: Fidelity,
1358    ) -> Result<SessionFollower> {
1359        SessionFollower::open_locator_with_fidelity(locator, fidelity)
1360    }
1361
1362    /// Follow a read-only view with explicit child-tree and history bounds.
1363    #[doc(hidden)]
1364    pub fn follow_read_view(
1365        &self,
1366        locator: &SessionLocator,
1367        fidelity: Fidelity,
1368        include_subagents: bool,
1369        message_limit: Option<usize>,
1370        max_message_chars: Option<usize>,
1371        display_history: bool,
1372    ) -> Result<SessionFollower> {
1373        SessionFollower::open_locator_with_view(
1374            locator,
1375            fidelity,
1376            include_subagents,
1377            message_limit,
1378            max_message_chars,
1379            display_history,
1380        )
1381    }
1382}
1383
1384/// ORCH-6: a Hermes row's durable address is the pair `{state.db, session
1385/// id}` — one store holds every session — so the generic file loader (which
1386/// opens the store's most recent session) would answer with the WRONG
1387/// conversation for every discovered locator but one, and its nouns with it.
1388/// Route by the locator's own session id instead.
1389fn load_hermes_locator(locator: &SessionLocator) -> Option<Result<Session>> {
1390    if locator.harness.as_str() != HarnessId::HERMES {
1391        return None;
1392    }
1393    let StorageLocator::File { path } = &locator.storage else {
1394        return None;
1395    };
1396    Some(Session::from_hermes_sqlite(path, Some(&locator.session_id)))
1397}
1398
1399/// ORCH-6: resolve every discovered row's `trigger` and `workspace` through
1400/// ORCH-3's own derivation. Harness scanners fill only what their native store
1401/// states (`surface`, `profile`, `recurrence`, `cross_surface`, and an explicit
1402/// `trigger` when the source says); this pass rebuilds a `SessionMeta` from
1403/// those facts plus `cwd` and re-reads the nouns off it, so a discovered row
1404/// and a loaded session can never disagree about the same session.
1405fn finalize_nouns(descriptor: &mut SessionDescriptor) {
1406    let mut meta = SessionMeta::new(SessionSource::Native);
1407    meta.cwd = descriptor.cwd.clone();
1408    meta.trigger = descriptor.nouns.trigger;
1409    meta.surface = descriptor.nouns.surface.clone();
1410    meta.profile = descriptor.nouns.profile.clone();
1411    meta.recurrence = descriptor.nouns.recurrence.clone();
1412    meta.cross_surface = descriptor.nouns.cross_surface.clone();
1413    descriptor.nouns = OrchestrationNouns::from_meta(&meta);
1414}
1415
1416fn project_descriptors(query: &DiscoveryQuery, found: &mut Vec<SessionDescriptor>) {
1417    roll_up_session_children(found, query.include_child_sessions);
1418    if let Some(root_session_id) = query.root_session_id.as_deref() {
1419        retain_session_family(found, root_session_id);
1420    }
1421    if let Some(family_path) = query.workspace_family.as_deref() {
1422        let family = RepoFamily::of(family_path);
1423        let mut cache: HashMap<PathBuf, bool> = HashMap::new();
1424        found.retain(|descriptor| {
1425            let Some(cwd) = descriptor.cwd.as_deref() else {
1426                return false;
1427            };
1428            *cache
1429                .entry(cwd.to_path_buf())
1430                .or_insert_with(|| RepoFamily::of(cwd).joins(&family))
1431        });
1432    }
1433    if let Some(profile) = query
1434        .profile
1435        .as_deref()
1436        .map(str::trim)
1437        .filter(|p| !p.is_empty())
1438    {
1439        found.retain(|descriptor| descriptor.nouns.profile.as_deref() == Some(profile));
1440    }
1441    if let Some(after) = query.updated_after_ms {
1442        found.retain(|descriptor| descriptor.updated_at_ms.is_some_and(|at| at >= after));
1443    }
1444    if let Some(before) = query.updated_before_ms {
1445        found.retain(|descriptor| descriptor.updated_at_ms.is_some_and(|at| at <= before));
1446    }
1447    found.sort_by(|a, b| {
1448        b.updated_at_ms
1449            .cmp(&a.updated_at_ms)
1450            .then_with(|| a.locator.harness.cmp(&b.locator.harness))
1451            .then_with(|| a.locator.session_id.cmp(&b.locator.session_id))
1452    });
1453    if let Some(search) = query
1454        .query
1455        .as_deref()
1456        .map(str::trim)
1457        .filter(|query| !query.is_empty())
1458    {
1459        let search = search.to_lowercase();
1460        found.retain(|descriptor| descriptor_matches(descriptor, &search));
1461    }
1462}
1463
1464fn paginate_descriptors(
1465    query: &DiscoveryQuery,
1466    found: Vec<SessionDescriptor>,
1467) -> Result<(Vec<SessionDescriptor>, Option<String>)> {
1468    let start = match query.cursor.as_deref() {
1469        Some(cursor) => {
1470            let key = decode_cursor(cursor)?;
1471            found
1472                .iter()
1473                .position(|descriptor| descriptor_cursor_key(descriptor) == key)
1474                .map(|index| index + 1)
1475                .ok_or_else(|| Error::Other("discovery cursor is stale or invalid".into()))?
1476        }
1477        None => 0,
1478    };
1479    let end = query
1480        .limit
1481        .map(|limit| start.saturating_add(limit).min(found.len()))
1482        .unwrap_or(found.len());
1483    let sessions = found[start.min(found.len())..end].to_vec();
1484    let next_cursor = (end < found.len())
1485        .then(|| sessions.last().map(encode_cursor))
1486        .flatten();
1487    Ok((sessions, next_cursor))
1488}
1489
1490fn enrich_descriptors(
1491    query: &DiscoveryQuery,
1492    sessions: &mut [SessionDescriptor],
1493    codex_history: Option<&CodexHistoryTopicIndex>,
1494) -> Result<()> {
1495    let codex_topics = if query.include_topic_candidates && codex_history.is_none() {
1496        codex_history_topics(&query.homes.codex, sessions).unwrap_or_default()
1497    } else {
1498        HashMap::new()
1499    };
1500    for descriptor in sessions {
1501        let topic = (descriptor.locator.harness.as_str() == HarnessId::CODEX)
1502            .then(|| {
1503                codex_history
1504                    .and_then(|history| history.topics.get(&descriptor.locator.session_id))
1505                    .or_else(|| codex_topics.get(&descriptor.locator.session_id))
1506            })
1507            .flatten();
1508        enrich_descriptor(query, descriptor, topic);
1509    }
1510    Ok(())
1511}
1512
1513fn enrich_descriptor(
1514    query: &DiscoveryQuery,
1515    descriptor: &mut SessionDescriptor,
1516    codex_topic: Option<&Vec<SessionPreviewCandidate>>,
1517) {
1518    if query.include_topic_candidates {
1519        descriptor.preview_candidates = codex_topic
1520            .cloned()
1521            .unwrap_or_else(|| topic_message_candidates(&descriptor.locator).unwrap_or_default());
1522    }
1523    descriptor.latest_message_candidates =
1524        latest_message_candidates(&descriptor.locator).unwrap_or_default();
1525}
1526
1527fn is_false(value: &bool) -> bool {
1528    !value
1529}
1530
1531fn descriptor_matches(descriptor: &SessionDescriptor, search: &str) -> bool {
1532    [
1533        Some(descriptor.locator.harness.as_str()),
1534        Some(descriptor.locator.session_id.as_str()),
1535        descriptor.title.as_deref(),
1536        descriptor.cwd.as_ref().and_then(|path| path.to_str()),
1537        descriptor.model.as_deref(),
1538    ]
1539    .into_iter()
1540    .flatten()
1541    .any(|value| value.to_lowercase().contains(search))
1542}
1543
1544fn descriptor_cursor_key(descriptor: &SessionDescriptor) -> (Option<u64>, String, String) {
1545    (
1546        descriptor.updated_at_ms,
1547        descriptor.locator.harness.as_str().to_string(),
1548        descriptor.locator.session_id.clone(),
1549    )
1550}
1551
1552fn encode_cursor(descriptor: &SessionDescriptor) -> String {
1553    let json = serde_json::to_vec(&descriptor_cursor_key(descriptor)).unwrap_or_default();
1554    let mut encoded = String::with_capacity(json.len() * 2);
1555    for byte in json {
1556        use std::fmt::Write;
1557        let _ = write!(&mut encoded, "{byte:02x}");
1558    }
1559    encoded
1560}
1561
1562fn decode_cursor(cursor: &str) -> Result<(Option<u64>, String, String)> {
1563    if cursor.len() % 2 != 0 {
1564        return Err(Error::Other("discovery cursor is invalid".into()));
1565    }
1566    let bytes = (0..cursor.len())
1567        .step_by(2)
1568        .map(|index| u8::from_str_radix(&cursor[index..index + 2], 16))
1569        .collect::<std::result::Result<Vec<_>, _>>()
1570        .map_err(|_| Error::Other("discovery cursor is invalid".into()))?;
1571    serde_json::from_slice(&bytes).map_err(|_| Error::Other("discovery cursor is invalid".into()))
1572}
1573
1574#[derive(Default)]
1575struct HeaderMeta {
1576    session_id: Option<String>,
1577    cwd: Option<PathBuf>,
1578    title: Option<String>,
1579    model: Option<String>,
1580    parent_session_id: Option<String>,
1581}
1582
1583fn discover_jsonl(
1584    root: &Path,
1585    harness: &str,
1586    workspace: Option<&WorkspaceScope>,
1587    include_child_sessions: bool,
1588    found: &mut Vec<SessionDescriptor>,
1589) {
1590    let mut files = Vec::new();
1591    collect_jsonl(root, harness, include_child_sessions, &mut files);
1592    for path in files {
1593        let Ok(meta) = read_header(&path, harness) else {
1594            continue;
1595        };
1596        if workspace.is_some_and(|wanted| meta.cwd.as_deref().is_none_or(|cwd| !wanted.admits(cwd)))
1597        {
1598            continue;
1599        }
1600        let session_id = meta.session_id.unwrap_or_else(|| {
1601            path.file_stem()
1602                .and_then(|value| value.to_str())
1603                .unwrap_or("unknown")
1604                .to_string()
1605        });
1606        let parent_session_id = meta.parent_session_id.or_else(|| {
1607            (harness == HarnessId::CLAUDE_CODE)
1608                .then(|| claude_subagent_parent_id(&path))
1609                .flatten()
1610        });
1611        let tail = tail_facts(&path, harness);
1612        found.push(SessionDescriptor {
1613            locator: SessionLocator {
1614                harness: HarnessId::new(harness),
1615                session_id,
1616                storage: StorageLocator::File { path: path.clone() },
1617            },
1618            cwd: meta.cwd,
1619            title: meta.title,
1620            preview_candidates: Vec::new(),
1621            latest_message_candidates: Vec::new(),
1622            updated_at_ms: tail.last_turn_ms.or_else(|| modified_ms(&path)),
1623            message_count: None,
1624            model: tail.model.or(meta.model),
1625            parent_session_id,
1626            child_session_count: if harness == HarnessId::CLAUDE_CODE && !include_child_sessions {
1627                count_claude_subagents(&path)
1628            } else {
1629                0
1630            },
1631            nouns: OrchestrationNouns::default(),
1632        });
1633    }
1634}
1635
1636fn collect_jsonl(root: &Path, harness: &str, include_child_sessions: bool, out: &mut Vec<PathBuf>) {
1637    let mut walked = HashSet::new();
1638    collect_jsonl_in(root, harness, include_child_sessions, out, &mut walked);
1639}
1640
1641/// A symlinked directory is walked like a real one: a harness home whose project directories
1642/// were moved to another volume and linked back (a box offloading its disk) still holds every
1643/// one of those sessions. `walked` holds each directory's canonical path, so a link cycle ends.
1644fn collect_jsonl_in(
1645    root: &Path,
1646    harness: &str,
1647    include_child_sessions: bool,
1648    out: &mut Vec<PathBuf>,
1649    walked: &mut HashSet<PathBuf>,
1650) {
1651    if !walked.insert(fs::canonicalize(root).unwrap_or_else(|_| root.to_path_buf())) {
1652        return;
1653    }
1654    let Ok(entries) = fs::read_dir(root) else {
1655        return;
1656    };
1657    for entry in entries.flatten() {
1658        let Ok(mut kind) = entry.file_type() else {
1659            continue;
1660        };
1661        let path = entry.path();
1662        if kind.is_symlink() {
1663            let Ok(target) = fs::metadata(&path) else {
1664                continue;
1665            };
1666            kind = target.file_type();
1667        }
1668        if kind.is_dir() {
1669            if harness == HarnessId::CLAUDE_CODE
1670                && path.file_name().and_then(|v| v.to_str()) == Some("subagents")
1671                && !include_child_sessions
1672            {
1673                continue;
1674            }
1675            collect_jsonl_in(&path, harness, include_child_sessions, out, walked);
1676        } else if kind.is_file()
1677            && (path.extension().and_then(|v| v.to_str()) == Some("jsonl")
1678                || (harness == HarnessId::CODEX
1679                    && path
1680                        .file_name()
1681                        .and_then(|v| v.to_str())
1682                        .is_some_and(|name| name.ends_with(".jsonl.zst"))))
1683        {
1684            out.push(path);
1685        }
1686    }
1687}
1688
1689fn read_header(path: &Path, harness: &str) -> Result<HeaderMeta> {
1690    let file = crate::session::open_session_reader(path)?;
1691    let mut result = HeaderMeta::default();
1692    let mut bytes = 0usize;
1693    for line in BufReader::new(file).lines().take(32) {
1694        let line = line?;
1695        bytes += line.len();
1696        if bytes > 256 * 1024 {
1697            break;
1698        }
1699        let Ok(value) = serde_json::from_str::<Value>(&line) else {
1700            continue;
1701        };
1702        update_header_meta(&mut result, &value, harness);
1703        if result.session_id.is_some() && result.cwd.is_some() && result.model.is_some() {
1704            break;
1705        }
1706    }
1707    if result.session_id.is_none() && result.cwd.is_none() {
1708        return Err(Error::Other(format!(
1709            "{} has no recognizable {harness} session header",
1710            path.display()
1711        )));
1712    }
1713    Ok(result)
1714}
1715
1716fn update_header_meta(result: &mut HeaderMeta, value: &Value, harness: &str) {
1717    match harness {
1718        HarnessId::CLAUDE_CODE => {
1719            fill_string(&mut result.session_id, value.get("sessionId"));
1720            fill_path(&mut result.cwd, value.get("cwd"));
1721            // The session's name (`claude -n`, `/rename`) is its title.
1722            fill_string(&mut result.title, value.get("customTitle"));
1723            fill_string(&mut result.title, value.get("agentName"));
1724            fill_string(
1725                &mut result.model,
1726                value.get("message").and_then(|v| v.get("model")),
1727            );
1728        }
1729        HarnessId::CODEX => {
1730            let payload = value.get("payload").unwrap_or(&Value::Null);
1731            if value.get("type").and_then(Value::as_str) == Some("session_meta") {
1732                fill_string(&mut result.session_id, payload.get("id"));
1733                fill_path(&mut result.cwd, payload.get("cwd"));
1734                fill_string(&mut result.title, payload.get("thread_name"));
1735                fill_string(&mut result.title, payload.get("title"));
1736                fill_string(
1737                    &mut result.parent_session_id,
1738                    payload.get("parent_thread_id"),
1739                );
1740                if let Some(parent) = payload
1741                    .pointer("/source/subagent/thread_spawn/parent_thread_id")
1742                    .and_then(Value::as_str)
1743                {
1744                    result.parent_session_id = Some(parent.to_string());
1745                }
1746                if result.title.is_none() {
1747                    result.title = payload
1748                        .pointer("/source/subagent/thread_spawn/agent_path")
1749                        .and_then(Value::as_str)
1750                        .and_then(|path| path.rsplit('/').find(|part| !part.is_empty()))
1751                        .map(humanize_topic);
1752                }
1753            }
1754            if value.get("type").and_then(Value::as_str) == Some("turn_context") {
1755                fill_path(&mut result.cwd, payload.get("cwd"));
1756                fill_string(&mut result.model, payload.get("model"));
1757            }
1758        }
1759        HarnessId::PI => {
1760            if value.get("type").and_then(Value::as_str) == Some("session") {
1761                fill_string(&mut result.session_id, value.get("id"));
1762                fill_path(&mut result.cwd, value.get("cwd"));
1763            }
1764            fill_string(
1765                &mut result.model,
1766                value.get("message").and_then(|v| v.get("model")),
1767            );
1768        }
1769        _ => {}
1770    }
1771}
1772
1773/// Collapse native child rollouts into their root conversation before sorting
1774/// and pagination. A child's write time contributes to the root so active
1775/// delegated work keeps the conversation visible without creating extra rows.
1776fn roll_up_session_children(found: &mut Vec<SessionDescriptor>, include_children: bool) {
1777    let by_id = found
1778        .iter()
1779        .enumerate()
1780        .map(|(index, descriptor)| {
1781            (
1782                (
1783                    descriptor.locator.harness.as_str().to_string(),
1784                    descriptor.locator.session_id.clone(),
1785                ),
1786                index,
1787            )
1788        })
1789        .collect::<HashMap<_, _>>();
1790    let mut root_updates = HashMap::<usize, u64>::new();
1791    let mut root_child_counts = HashMap::<usize, usize>::new();
1792
1793    for descriptor in found.iter() {
1794        let Some(mut parent_id) = descriptor.parent_session_id.as_deref() else {
1795            continue;
1796        };
1797        let harness = descriptor.locator.harness.as_str();
1798        let mut root = None;
1799        let mut visited = HashSet::new();
1800        while visited.insert(parent_id.to_string()) {
1801            let Some(&parent_index) = by_id.get(&(harness.to_string(), parent_id.to_string()))
1802            else {
1803                break;
1804            };
1805            root = Some(parent_index);
1806            let Some(next_parent) = found[parent_index].parent_session_id.as_deref() else {
1807                break;
1808            };
1809            parent_id = next_parent;
1810        }
1811        if let (Some(root), Some(updated_at_ms)) = (root, descriptor.updated_at_ms) {
1812            root_updates
1813                .entry(root)
1814                .and_modify(|current| *current = (*current).max(updated_at_ms))
1815                .or_insert(updated_at_ms);
1816        }
1817        if let Some(root) = root {
1818            *root_child_counts.entry(root).or_default() += 1;
1819        }
1820    }
1821
1822    for (root, child_updated_at_ms) in root_updates {
1823        found[root].updated_at_ms = Some(
1824            found[root]
1825                .updated_at_ms
1826                .unwrap_or_default()
1827                .max(child_updated_at_ms),
1828        );
1829    }
1830    for (root, child_count) in root_child_counts {
1831        found[root].child_session_count = child_count;
1832    }
1833    if !include_children {
1834        found.retain(|descriptor| descriptor.parent_session_id.is_none());
1835    }
1836}
1837
1838fn retain_session_family(found: &mut Vec<SessionDescriptor>, root_session_id: &str) {
1839    let parent_by_id = found
1840        .iter()
1841        .map(|descriptor| {
1842            (
1843                descriptor.locator.session_id.clone(),
1844                descriptor.parent_session_id.clone(),
1845            )
1846        })
1847        .collect::<HashMap<_, _>>();
1848    found.retain(|descriptor| {
1849        let mut current = descriptor.locator.session_id.clone();
1850        let mut visited = HashSet::new();
1851        while visited.insert(current.clone()) {
1852            if current == root_session_id {
1853                return true;
1854            }
1855            let Some(Some(parent)) = parent_by_id.get(&current) else {
1856                return false;
1857            };
1858            current = parent.clone();
1859        }
1860        false
1861    });
1862}
1863
1864fn claude_subagent_parent_id(path: &Path) -> Option<String> {
1865    let subagents = path.parent()?;
1866    if subagents.file_name()?.to_str()? != "subagents" {
1867        return None;
1868    }
1869    subagents
1870        .parent()?
1871        .file_name()?
1872        .to_str()
1873        .map(str::to_string)
1874}
1875
1876fn count_claude_subagents(parent_path: &Path) -> usize {
1877    let Some(parent) = parent_path.parent() else {
1878        return 0;
1879    };
1880    let Some(stem) = parent_path.file_stem() else {
1881        return 0;
1882    };
1883    let root = parent.join(stem).join("subagents");
1884    let mut files = Vec::new();
1885    collect_jsonl(&root, HarnessId::CLAUDE_CODE, true, &mut files);
1886    files.len()
1887}
1888
1889fn is_zero(value: &usize) -> bool {
1890    *value == 0
1891}
1892
1893fn humanize_topic(value: &str) -> String {
1894    let text = value.replace(['_', '-'], " ");
1895    let mut characters = text.chars();
1896    match characters.next() {
1897        Some(first) => first.to_uppercase().collect::<String>() + characters.as_str(),
1898        None => text,
1899    }
1900}
1901
1902fn discover_gemini(
1903    root: &Path,
1904    workspace: Option<&WorkspaceScope>,
1905    found: &mut Vec<SessionDescriptor>,
1906) {
1907    let slug_to_cwd = std::fs::read_to_string(root.join("projects.json"))
1908        .ok()
1909        .and_then(|text| serde_json::from_str::<Value>(&text).ok())
1910        .and_then(|value| value.get("projects").and_then(Value::as_object).cloned())
1911        .map(|projects| {
1912            projects
1913                .into_iter()
1914                .filter_map(|(cwd, slug)| Some((slug.as_str()?.to_string(), PathBuf::from(cwd))))
1915                .collect::<HashMap<_, _>>()
1916        })
1917        .unwrap_or_default();
1918    let mut files = Vec::new();
1919    collect_jsonl(&root.join("tmp"), HarnessId::GEMINI, false, &mut files);
1920    let worker_count = std::thread::available_parallelism()
1921        .map(usize::from)
1922        .unwrap_or(4)
1923        .clamp(1, 8)
1924        .min(files.len().max(1));
1925    let chunk_size = files.len().max(1).div_ceil(worker_count);
1926    let discovered = std::thread::scope(|scope| {
1927        files
1928            .chunks(chunk_size)
1929            .map(|paths| {
1930                scope.spawn(|| {
1931                    paths
1932                        .iter()
1933                        .filter_map(|path| gemini_descriptor(path, &slug_to_cwd, workspace))
1934                        .collect::<Vec<_>>()
1935                })
1936            })
1937            .collect::<Vec<_>>()
1938            .into_iter()
1939            .flat_map(|worker| {
1940                worker
1941                    .join()
1942                    .expect("Gemini discovery worker must not panic")
1943            })
1944            .collect::<Vec<_>>()
1945    });
1946    found.extend(discovered);
1947}
1948
1949fn gemini_descriptor(
1950    path: &Path,
1951    slug_to_cwd: &HashMap<String, PathBuf>,
1952    workspace: Option<&WorkspaceScope>,
1953) -> Option<SessionDescriptor> {
1954    if path
1955        .parent()
1956        .and_then(Path::file_name)
1957        .and_then(|name| name.to_str())
1958        != Some("chats")
1959    {
1960        return None;
1961    }
1962    let slug = path
1963        .parent()
1964        .and_then(Path::parent)
1965        .and_then(Path::file_name)
1966        .and_then(|name| name.to_str());
1967    let cwd = slug.and_then(|slug| slug_to_cwd.get(slug)).cloned();
1968    if workspace.is_some_and(|wanted| cwd.as_deref().is_none_or(|actual| !wanted.admits(actual))) {
1969        return None;
1970    }
1971
1972    // The native session id lives on line one. A small decoration budget keeps
1973    // the common title/model case without turning 1,800 sessions into a
1974    // sequential 60 MiB read before the list can render.
1975    let file = File::open(path).ok()?;
1976    let mut reader = BufReader::new(file.take(64 * 1024));
1977    let mut header = String::new();
1978    reader.read_line(&mut header).ok()?;
1979    let header = serde_json::from_str::<Value>(&header).ok()?;
1980    let session_id = header.get("sessionId")?.as_str()?.to_string();
1981    let mut model = None;
1982    for line in reader
1983        .take(4 * 1024)
1984        .lines()
1985        .map_while(std::result::Result::ok)
1986    {
1987        let Ok(value) = serde_json::from_str::<Value>(&line) else {
1988            continue;
1989        };
1990        let kind = value.get("type").and_then(Value::as_str);
1991        if kind != Some("user") && kind != Some("gemini") {
1992            continue;
1993        }
1994        if model.is_none() {
1995            model = value
1996                .get("model")
1997                .and_then(Value::as_str)
1998                .map(str::to_string);
1999        }
2000        if model.is_some() {
2001            break;
2002        }
2003    }
2004    Some(SessionDescriptor {
2005        locator: SessionLocator {
2006            harness: HarnessId::from(HarnessId::GEMINI),
2007            session_id,
2008            storage: StorageLocator::File {
2009                path: path.to_path_buf(),
2010            },
2011        },
2012        cwd,
2013        title: None,
2014        preview_candidates: Vec::new(),
2015        latest_message_candidates: Vec::new(),
2016        updated_at_ms: tail_facts(path, HarnessId::GEMINI)
2017            .last_turn_ms
2018            .or_else(|| modified_ms(path)),
2019        message_count: None,
2020        model,
2021        parent_session_id: None,
2022        child_session_count: 0,
2023        nouns: OrchestrationNouns::default(),
2024    })
2025}
2026
2027fn display_text(content: Option<&Value>) -> Option<String> {
2028    match content? {
2029        Value::String(text) => Some(text.clone()),
2030        Value::Array(parts) => Some(
2031            parts
2032                .iter()
2033                .filter_map(|part| part.get("text").and_then(Value::as_str))
2034                .collect::<Vec<_>>()
2035                .join(" ")
2036                .trim()
2037                .to_string(),
2038        ),
2039        _ => None,
2040    }
2041}
2042
2043/// Hermes (UNI-15): enumerate sessions from the single `state.db` SQLite
2044/// store, strictly read-only (the store is a live, shared, WAL,
2045/// single-writer database owned by a running Hermes install). A missing or
2046/// non-Hermes file is skipped silently, like every other absent home.
2047fn discover_hermes(
2048    db_path: &Path,
2049    workspace: Option<&WorkspaceScope>,
2050    selected: &HashSet<&str>,
2051    found: &mut Vec<SessionDescriptor>,
2052) {
2053    if !db_path.is_file() {
2054        return;
2055    }
2056    let Ok(conn) = Connection::open_with_flags(
2057        db_path,
2058        rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
2059    ) else {
2060        return;
2061    };
2062    let fingerprint_ok = ["sessions", "messages", "schema_version"].iter().all(|t| {
2063        conn.query_row(
2064            "SELECT 1 FROM sqlite_master WHERE type='table' AND name=?1",
2065            [t],
2066            |_| Ok(()),
2067        )
2068        .is_ok()
2069    });
2070    if !fingerprint_ok {
2071        return;
2072    }
2073    let Ok(mut statement) = conn.prepare(
2074        "SELECT id, cwd, title, model, message_count, started_at, ended_at, parent_session_id, \
2075         source, model_config FROM sessions ORDER BY started_at DESC",
2076    ) else {
2077        return;
2078    };
2079    let Ok(rows) = statement.query_map([], |row| {
2080        Ok((
2081            row.get::<_, String>(0)?,
2082            row.get::<_, Option<String>>(1)?,
2083            row.get::<_, Option<String>>(2)?,
2084            row.get::<_, Option<String>>(3)?,
2085            row.get::<_, Option<i64>>(4)?,
2086            row.get::<_, Option<f64>>(5)?,
2087            row.get::<_, Option<f64>>(6)?,
2088            row.get::<_, Option<String>>(7)?,
2089            row.get::<_, Option<String>>(8)?,
2090            row.get::<_, Option<String>>(9)?,
2091        ))
2092    }) else {
2093        return;
2094    };
2095    for row in rows.flatten() {
2096        let (
2097            id,
2098            cwd,
2099            title,
2100            model,
2101            message_count,
2102            started_at,
2103            ended_at,
2104            parent,
2105            source,
2106            model_config,
2107        ) = row;
2108        // a worker session the orchestrator mirrored into this store (hermes-compat row 11) is listed once:
2109        // as the worker's own session when this query also reads that worker's store, else here, as Hermes
2110        // lists it; one Hermes has carried past its mirror is Hermes's conversation too, and listed here
2111        let mirror = model_config
2112            .as_deref()
2113            .and_then(|c| serde_json::from_str::<serde_json::Value>(c).ok())
2114            .and_then(|c| c.get("_supercode_mirror").cloned());
2115        if let Some(mirror) = mirror {
2116            let worker_read = mirror
2117                .get("harness")
2118                .and_then(serde_json::Value::as_str)
2119                .is_some_and(|h| selected.contains(h));
2120            let mirrored = mirror.get("messages").and_then(serde_json::Value::as_i64);
2121            if worker_read && mirrored.is_none_or(|m| message_count.unwrap_or(0) <= m) {
2122                continue;
2123            }
2124            if mirror.get("continued_as").is_some()
2125                && mirrored.is_none_or(|m| message_count.unwrap_or(0) <= m)
2126            {
2127                continue;
2128            }
2129        }
2130        let cwd = cwd.map(PathBuf::from);
2131        if let Some(filter) = workspace {
2132            if !filter.admits_literal(cwd.as_deref()) {
2133                continue;
2134            }
2135        }
2136        let updated_at_ms = ended_at
2137            .or(started_at)
2138            .map(|seconds| (seconds * 1000.0) as u64);
2139        // ORCH-6: the same derivation `Session::from_hermes_sqlite` runs, over
2140        // the same row — discovery reads the gateway columns it already has
2141        // open instead of guessing a lighter-weight variant.
2142        let mut meta = SessionMeta::new(SessionSource::Hermes);
2143        meta.cwd = cwd.clone();
2144        if let Some(hermes_source) = source.filter(|value| !value.is_empty()) {
2145            meta.lineage
2146                .insert("hermes_source".to_string(), hermes_source);
2147        }
2148        if let Some(parent_id) = parent.as_deref() {
2149            meta.lineage.insert(
2150                "hermes_lineage_kind".to_string(),
2151                crate::session::hermes_lineage_kind(
2152                    &conn,
2153                    parent_id,
2154                    model_config.as_deref(),
2155                    started_at,
2156                )
2157                .to_string(),
2158            );
2159        }
2160        hermes_capture_nouns(&conn, &id, &mut meta);
2161        found.push(SessionDescriptor {
2162            locator: SessionLocator {
2163                harness: HarnessId::new(HarnessId::HERMES),
2164                session_id: id,
2165                storage: StorageLocator::File {
2166                    path: db_path.to_path_buf(),
2167                },
2168            },
2169            cwd,
2170            title: title.filter(|t| !t.is_empty()),
2171            preview_candidates: Vec::new(),
2172            latest_message_candidates: Vec::new(),
2173            updated_at_ms,
2174            message_count: message_count.map(|count| count.max(0) as usize),
2175            model,
2176            parent_session_id: parent,
2177            child_session_count: 0,
2178            nouns: OrchestrationNouns::from_meta(&meta),
2179        });
2180    }
2181}
2182
2183/// The orchestrator (ORC-7): every profile folder's `state.db` `bindings`
2184/// table is one row per conversation the orchestrator holds. A binding is not
2185/// a transcript — the transcript belongs to the WORKER harness it points at —
2186/// so the row carries the surface, trigger, profile and worker identity, and
2187/// its locator addresses the worker's own storage when the binding recorded
2188/// one. Strictly read-only, like every other store here.
2189fn discover_orchestrator(
2190    root: &Path,
2191    workspace: Option<&WorkspaceScope>,
2192    found: &mut Vec<SessionDescriptor>,
2193) {
2194    // A binding has no cwd of its own; a workspace filter can only exclude it.
2195    if workspace.is_some() {
2196        return;
2197    }
2198    for (profile, dir) in orchestrator_profile_dirs(root) {
2199        let db_path = dir.join("state.db");
2200        if !db_path.is_file() {
2201            continue;
2202        }
2203        let Ok(conn) = Connection::open_with_flags(
2204            &db_path,
2205            rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
2206        ) else {
2207            continue;
2208        };
2209        // the session entries the orchestrator wrote into Hermes's `gateway_routing` (hermes-compat row 10),
2210        // each binding whole in `metadata.supercode`; the rest are Hermes's own sessions, listed by Hermes's reader
2211        let sessions = dir.join("sessions");
2212        let scope = std::fs::canonicalize(&sessions)
2213            .unwrap_or(sessions)
2214            .display()
2215            .to_string();
2216        let Ok(mut statement) = conn.prepare(
2217            "SELECT json_extract(entry_json, '$.metadata.supercode.binding'), \
2218             CAST(strftime('%s', json_extract(entry_json, '$.metadata.supercode.binding.last_activity_at')) AS INTEGER) \
2219             FROM gateway_routing WHERE scope = ?1 AND json_extract(entry_json, '$.metadata.supercode') IS NOT NULL \
2220             ORDER BY json_extract(entry_json, '$.metadata.supercode.binding.last_activity_at') DESC",
2221        ) else {
2222            continue;
2223        };
2224        let Ok(rows) = statement.query_map([&scope], |row| {
2225            Ok((
2226                row.get::<_, Option<String>>(0)?,
2227                row.get::<_, Option<i64>>(1)?,
2228            ))
2229        }) else {
2230            continue;
2231        };
2232        let rows = rows.flatten().filter_map(|(json, epoch)| {
2233            let b: Binding = serde_json::from_str(&json?).ok()?;
2234            let text = |v: &Option<String>| v.clone().filter(|s| !s.is_empty());
2235            Some((
2236                OrchestratorBindingRow {
2237                    platform: b.key.platform.clone().unwrap_or_default(),
2238                    chat_type: b.key.kind.clone().unwrap_or_default(),
2239                    chat_id: text(&b.key.chat_id),
2240                    thread_id: text(&b.key.thread_id),
2241                    participant_id: text(&b.key.participant_id),
2242                    worker_harness: b.worker.harness.as_str().to_string(),
2243                    worker_session_id: text(&b.worker.session_id),
2244                    worker_locator: text(&b.worker.locator),
2245                    started_at: b.started_at.clone(),
2246                    last_activity_at: b.last_activity_at.clone(),
2247                    ended_at: b.ended_at.clone(),
2248                    end_reason: b.end_reason.map(|r| r.as_str().to_string()),
2249                    handoff_to: b.handoff.as_ref().and_then(|h| h.to.clone()),
2250                    handoff_state: b.handoff.as_ref().map(|h| h.state.clone()),
2251                    handoff_error: b.handoff.as_ref().and_then(|h| h.error.clone()),
2252                    recurrence_job_id: b.recurrence.as_ref().map(|r| r.job_id.clone()),
2253                },
2254                epoch,
2255            ))
2256        });
2257        for (row, last_activity_epoch) in rows {
2258            found.push(orchestrator_descriptor(
2259                &db_path,
2260                &profile,
2261                &row,
2262                last_activity_epoch,
2263            ));
2264        }
2265    }
2266}
2267
2268fn orchestrator_descriptor(
2269    db_path: &Path,
2270    profile: &str,
2271    row: &OrchestratorBindingRow,
2272    last_activity_epoch: Option<i64>,
2273) -> SessionDescriptor {
2274    let binding = Binding::from_orchestrator_row(profile, row);
2275    let nouns = binding.nouns();
2276    // The row's title names the worker the binding points at: that pair is
2277    // the only address from which the conversation itself can be read.
2278    let mut title = format!(
2279        "{} {}",
2280        row.worker_harness,
2281        row.worker_session_id
2282            .as_deref()
2283            .unwrap_or("(no worker session yet)")
2284    );
2285    if let Some(reason) = row.end_reason.as_deref().filter(|_| row.ended_at.is_some()) {
2286        title.push_str(&format!(" (ended: {reason})"));
2287    }
2288    SessionDescriptor {
2289        locator: SessionLocator {
2290            harness: HarnessId::new(HarnessId::ORCHESTRATOR),
2291            session_id: row.worker_session_id.clone().unwrap_or_default(),
2292            storage: StorageLocator::File {
2293                path: row
2294                    .worker_locator
2295                    .clone()
2296                    .map_or_else(|| db_path.to_path_buf(), PathBuf::from),
2297            },
2298        },
2299        cwd: None,
2300        title: Some(title),
2301        preview_candidates: Vec::new(),
2302        latest_message_candidates: Vec::new(),
2303        updated_at_ms: last_activity_epoch.map(|seconds| (seconds.max(0) as u64) * 1000),
2304        message_count: None,
2305        model: None,
2306        parent_session_id: None,
2307        child_session_count: 0,
2308        nouns,
2309    }
2310}
2311
2312/// OpenClaw >= 2026.7: `<home>/agents/<agentId>/sessions/<uuid>.jsonl` are
2313/// plain pi-v3 dialect session files. `.trajectory.jsonl` runtime traces and
2314/// `.trajectory-path.json` pointers live in the SAME directory and are
2315/// excluded by suffix plus a header check (their first line carries
2316/// `traceSchema`, never `type:"session"`). Read-only discovery (UNI-16).
2317fn discover_openclaw(
2318    root: &Path,
2319    workspace: Option<&WorkspaceScope>,
2320    found: &mut Vec<SessionDescriptor>,
2321) {
2322    let agents = root.join("agents");
2323    let Ok(agent_dirs) = std::fs::read_dir(&agents) else {
2324        return;
2325    };
2326    for agent_dir in agent_dirs.flatten() {
2327        let sessions = agent_dir.path().join("sessions");
2328        let Ok(files) = std::fs::read_dir(&sessions) else {
2329            continue;
2330        };
2331        for file in files.flatten() {
2332            let path = file.path();
2333            let name = file.file_name();
2334            let name = name.to_string_lossy();
2335            if !name.ends_with(".jsonl") || name.ends_with(".trajectory.jsonl") {
2336                continue;
2337            }
2338            let Ok(text) = std::fs::read_to_string(&path) else {
2339                continue;
2340            };
2341            let Some(header_line) = text.lines().find(|line| !line.trim().is_empty()) else {
2342                continue;
2343            };
2344            let Ok(header) = serde_json::from_str::<serde_json::Value>(header_line) else {
2345                continue;
2346            };
2347            if header.get("type").and_then(serde_json::Value::as_str) != Some("session") {
2348                continue;
2349            }
2350            let session_id = header
2351                .get("id")
2352                .and_then(serde_json::Value::as_str)
2353                .unwrap_or_else(|| name.trim_end_matches(".jsonl"))
2354                .to_string();
2355            let cwd = header
2356                .get("cwd")
2357                .and_then(serde_json::Value::as_str)
2358                .map(PathBuf::from);
2359            if let Some(filter) = workspace {
2360                if !filter.admits_literal(cwd.as_deref()) {
2361                    continue;
2362                }
2363            }
2364            let updated_at_ms = tail_facts(&path, HarnessId::OPENCLAW)
2365                .last_turn_ms
2366                .or_else(|| {
2367                    file.metadata()
2368                        .ok()
2369                        .and_then(|metadata| metadata.modified().ok())
2370                        .and_then(|modified| modified.duration_since(std::time::UNIX_EPOCH).ok())
2371                        .map(|elapsed| elapsed.as_millis() as u64)
2372                });
2373            let message_count = text
2374                .lines()
2375                .filter(|line| line.contains("\"type\":\"message\""))
2376                .count();
2377            // ORCH-6: the same two facts `Session::load` reads for an OpenClaw
2378            // file — the gateway `sessionKey` in the header, and the agent id
2379            // in `agents/<id>/`.
2380            let mut meta = SessionMeta::new(SessionSource::OpenClaw);
2381            meta.cwd = cwd.clone();
2382            openclaw_capture_header_nouns(&header, &mut meta);
2383            if meta.profile.is_none() {
2384                meta.profile = openclaw_agent_id_from_path(&path);
2385            }
2386            found.push(SessionDescriptor {
2387                locator: SessionLocator {
2388                    harness: HarnessId::new(HarnessId::OPENCLAW),
2389                    session_id,
2390                    storage: StorageLocator::File { path },
2391                },
2392                cwd,
2393                title: None,
2394                preview_candidates: Vec::new(),
2395                latest_message_candidates: Vec::new(),
2396                updated_at_ms,
2397                message_count: Some(message_count),
2398                model: None,
2399                parent_session_id: None,
2400                child_session_count: 0,
2401                nouns: OrchestrationNouns::from_meta(&meta),
2402            });
2403        }
2404    }
2405}
2406
2407fn discover_supercode(
2408    root: &Path,
2409    workspace: Option<&WorkspaceScope>,
2410    found: &mut Vec<SessionDescriptor>,
2411) {
2412    for info in list_native_store(root) {
2413        let path = if info.archived {
2414            root.join("archived").join(format!("{}.jsonl", info.name))
2415        } else {
2416            root.join(format!("{}.jsonl", info.name))
2417        };
2418        // The sidecar is the store's authoritative content whenever it exists
2419        // (see `read_native_store_header`), and a reduced session can outlive its
2420        // working transcript entirely. Address the file that IS there: a locator
2421        // naming a deleted `<name>.jsonl` is one discovery's own loader rejects.
2422        let sidecar = path.with_extension("sidecar.jsonl");
2423        let path = if path.is_file() {
2424            path
2425        } else {
2426            sidecar.clone()
2427        };
2428        let header = read_native_store_header(&path);
2429        if workspace.is_some_and(|wanted| {
2430            header
2431                .as_ref()
2432                .and_then(|meta| meta.cwd.as_deref())
2433                .is_none_or(|cwd| !wanted.admits(cwd))
2434        }) {
2435            continue;
2436        }
2437        let title = (!info.title.trim().is_empty()).then_some(info.title);
2438        // The native store's own records carry no clock; the sidecar beside it
2439        // stamps every turn. Prefer that, and keep mtime as the last resort.
2440        let updated_at_ms = tail_facts(&sidecar, HarnessId::SUPERCODE)
2441            .last_turn_ms
2442            .or_else(|| tail_facts(&path, HarnessId::SUPERCODE).last_turn_ms)
2443            .or_else(|| modified_ms(&path))
2444            .or_else(|| modified_ms(&sidecar));
2445        found.push(SessionDescriptor {
2446            locator: SessionLocator {
2447                harness: HarnessId::from(HarnessId::SUPERCODE),
2448                session_id: info.name,
2449                storage: StorageLocator::File { path: path.clone() },
2450            },
2451            cwd: header.as_ref().and_then(|meta| meta.cwd.clone()),
2452            title,
2453            preview_candidates: Vec::new(),
2454            latest_message_candidates: Vec::new(),
2455            updated_at_ms,
2456            message_count: None,
2457            model: header.and_then(|meta| meta.model),
2458            parent_session_id: None,
2459            child_session_count: 0,
2460            nouns: OrchestrationNouns::default(),
2461        });
2462    }
2463}
2464
2465/// Read only the bounded native envelope needed by discovery. Loading a
2466/// sidecar-backed session here used to deserialize the complete byte-lossless
2467/// transcript family, making a workspace list proportional to every saved
2468/// Supercode transcript on the machine.
2469fn read_native_store_header(path: &Path) -> Option<HeaderMeta> {
2470    let name = path.file_stem()?.to_str()?;
2471    let sidecar = path.with_file_name(format!("{name}.sidecar.jsonl"));
2472    let source_path = if sidecar.is_file() {
2473        sidecar
2474    } else {
2475        path.to_path_buf()
2476    };
2477    let file = File::open(source_path).ok()?;
2478    let mut result = HeaderMeta::default();
2479    let mut source = None;
2480    let mut bytes = 0usize;
2481    for line in BufReader::new(file).lines().take(32) {
2482        let line = line.ok()?;
2483        bytes += line.len();
2484        if bytes > 256 * 1024 {
2485            break;
2486        }
2487        let Ok(value) = serde_json::from_str::<Value>(&line) else {
2488            continue;
2489        };
2490        if source.is_none() {
2491            source = value.get("source").and_then(Value::as_str).map(|source| {
2492                if source == "claude_code" {
2493                    HarnessId::CLAUDE_CODE.to_string()
2494                } else {
2495                    source.to_string()
2496                }
2497            });
2498            fill_string(&mut result.session_id, value.get("session_id"));
2499        }
2500        if let Some(harness) = source.as_deref() {
2501            update_header_meta(&mut result, &value, harness);
2502        }
2503        if result.cwd.is_some() && result.model.is_some() {
2504            break;
2505        }
2506    }
2507    Some(result)
2508}
2509
2510#[derive(Deserialize)]
2511struct NativeStoreInfo {
2512    name: String,
2513    #[serde(default)]
2514    title: String,
2515    #[serde(skip)]
2516    archived: bool,
2517}
2518
2519fn list_native_store(root: &Path) -> Vec<NativeStoreInfo> {
2520    let mut sessions = Vec::new();
2521    for archived in [false, true] {
2522        let directory = if archived {
2523            root.join("archived")
2524        } else {
2525            root.to_path_buf()
2526        };
2527        let Ok(entries) = fs::read_dir(directory) else {
2528            continue;
2529        };
2530        for entry in entries.flatten() {
2531            let path = entry.path();
2532            if !path.to_string_lossy().ends_with(".meta.json") {
2533                continue;
2534            }
2535            let Ok(text) = fs::read_to_string(path) else {
2536                continue;
2537            };
2538            let Ok(mut info) = serde_json::from_str::<NativeStoreInfo>(&text) else {
2539                continue;
2540            };
2541            info.archived = archived;
2542            sessions.push(info);
2543        }
2544    }
2545    sessions.sort_by(|left, right| left.name.cmp(&right.name));
2546    sessions
2547}
2548
2549fn discover_grok(
2550    root: &Path,
2551    workspace: Option<&WorkspaceScope>,
2552    found: &mut Vec<SessionDescriptor>,
2553) {
2554    let Ok(workspaces) = fs::read_dir(root) else {
2555        return;
2556    };
2557    for workspace_entry in workspaces.flatten() {
2558        let encoded = workspace_entry.file_name();
2559        let Some(cwd) = encoded
2560            .to_str()
2561            .and_then(percent_decode_path)
2562            .map(PathBuf::from)
2563        else {
2564            continue;
2565        };
2566        if workspace.is_some_and(|wanted| !wanted.admits(&cwd)) {
2567            continue;
2568        }
2569        let Ok(sessions) = fs::read_dir(workspace_entry.path()) else {
2570            continue;
2571        };
2572        for session_entry in sessions.flatten() {
2573            let session_dir = session_entry.path();
2574            if !session_dir.is_dir() {
2575                continue;
2576            }
2577            let transcript = session_dir.join("chat_history.jsonl");
2578            if !transcript.is_file() {
2579                continue;
2580            }
2581            let Some(session_id) = session_dir
2582                .file_name()
2583                .and_then(|name| name.to_str())
2584                .map(str::to_string)
2585            else {
2586                continue;
2587            };
2588            let summary = fs::read_to_string(session_dir.join("summary.json"))
2589                .ok()
2590                .and_then(|text| serde_json::from_str::<Value>(&text).ok());
2591            let title = summary
2592                .as_ref()
2593                .and_then(|value| value.get("generated_title"))
2594                .and_then(Value::as_str)
2595                .filter(|title| !title.is_empty())
2596                .map(str::to_string);
2597            let model = summary
2598                .as_ref()
2599                .and_then(|value| value.get("current_model_id"))
2600                .and_then(Value::as_str)
2601                .map(str::to_string);
2602            let message_count = summary
2603                .as_ref()
2604                .and_then(|value| value.get("num_chat_messages"))
2605                .and_then(Value::as_u64)
2606                .and_then(|count| usize::try_from(count).ok());
2607            let updated_at_ms = summary
2608                .as_ref()
2609                .and_then(|value| value.get("updated_at"))
2610                .and_then(Value::as_str)
2611                .and_then(crate::sidecar::rfc3339_to_ms)
2612                .and_then(|millis| u64::try_from(millis).ok())
2613                .or_else(|| modified_ms(&transcript));
2614            found.push(SessionDescriptor {
2615                locator: SessionLocator {
2616                    harness: HarnessId::from(HarnessId::GROK),
2617                    session_id,
2618                    storage: StorageLocator::File { path: transcript },
2619                },
2620                cwd: Some(cwd.clone()),
2621                title,
2622                preview_candidates: Vec::new(),
2623                latest_message_candidates: Vec::new(),
2624                updated_at_ms,
2625                message_count,
2626                model,
2627                parent_session_id: None,
2628                child_session_count: 0,
2629                nouns: OrchestrationNouns::default(),
2630            });
2631        }
2632    }
2633}
2634
2635fn discover_opencode(
2636    root: &Path,
2637    workspace: Option<&WorkspaceScope>,
2638    found: &mut Vec<SessionDescriptor>,
2639) {
2640    let mut dbs = Vec::new();
2641    if root.is_file() {
2642        dbs.push(root.to_path_buf());
2643    } else if let Ok(entries) = fs::read_dir(root) {
2644        dbs.extend(entries.flatten().map(|entry| entry.path()).filter(|path| {
2645            path.file_name()
2646                .and_then(|v| v.to_str())
2647                .is_some_and(|name| name.starts_with("opencode") && name.ends_with(".db"))
2648        }));
2649    }
2650    dbs.sort();
2651    for db in dbs {
2652        let Ok(conn) = Connection::open_with_flags(
2653            &db,
2654            rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
2655        ) else {
2656            continue;
2657        };
2658        let has_model = conn.prepare("SELECT model FROM session LIMIT 0").is_ok();
2659        let model_column = if has_model { "s.model" } else { "NULL" };
2660        let query = format!(
2661            "SELECT s.id, s.directory, s.title, s.time_updated, {model_column}, COUNT(m.id) \
2662             FROM session s LEFT JOIN message m ON m.session_id = s.id \
2663             GROUP BY s.id ORDER BY s.time_updated DESC"
2664        );
2665        let Ok(mut stmt) = conn.prepare(&query) else {
2666            continue;
2667        };
2668        let Ok(rows) = stmt.query_map([], |row| {
2669            Ok((
2670                row.get::<_, String>(0)?,
2671                row.get::<_, String>(1)?,
2672                row.get::<_, String>(2)?,
2673                row.get::<_, i64>(3)?,
2674                row.get::<_, Option<String>>(4)?,
2675                row.get::<_, i64>(5)?,
2676            ))
2677        }) else {
2678            continue;
2679        };
2680        for row in rows.flatten() {
2681            let (id, cwd, title, updated, model, messages) = row;
2682            let cwd = PathBuf::from(cwd);
2683            if workspace.is_some_and(|wanted| !wanted.admits(&cwd)) {
2684                continue;
2685            }
2686            found.push(SessionDescriptor {
2687                locator: SessionLocator {
2688                    harness: HarnessId::from(HarnessId::OPENCODE),
2689                    session_id: id.clone(),
2690                    storage: StorageLocator::Sqlite {
2691                        path: db.clone(),
2692                        selector: id,
2693                    },
2694                },
2695                cwd: Some(cwd),
2696                title: (!title.is_empty()).then_some(title),
2697                preview_candidates: Vec::new(),
2698                latest_message_candidates: Vec::new(),
2699                updated_at_ms: u64::try_from(updated).ok(),
2700                message_count: usize::try_from(messages).ok(),
2701                model,
2702                parent_session_id: None,
2703                child_session_count: 0,
2704                nouns: OrchestrationNouns::default(),
2705            });
2706        }
2707    }
2708}
2709
2710fn discover_goose(
2711    root: &Path,
2712    workspace: Option<&WorkspaceScope>,
2713    found: &mut Vec<SessionDescriptor>,
2714) {
2715    let db = if root.is_file() {
2716        root.to_path_buf()
2717    } else if root.join("sessions.db").is_file() {
2718        root.join("sessions.db")
2719    } else {
2720        root.join("sessions/sessions.db")
2721    };
2722    let Ok(connection) = Connection::open_with_flags(
2723        &db,
2724        rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
2725    ) else {
2726        return;
2727    };
2728    let Ok(mut statement) = connection.prepare(
2729        "SELECT s.id, s.working_dir, s.name, s.updated_at, s.model_config_json, \
2730                COUNT(m.id) \
2731         FROM sessions s LEFT JOIN messages m ON m.session_id = s.id \
2732         WHERE s.archived_at IS NULL \
2733         GROUP BY s.id ORDER BY s.updated_at DESC",
2734    ) else {
2735        return;
2736    };
2737    let Ok(rows) = statement.query_map([], |row| {
2738        Ok((
2739            row.get::<_, String>(0)?,
2740            row.get::<_, String>(1)?,
2741            row.get::<_, String>(2)?,
2742            row.get::<_, String>(3)?,
2743            row.get::<_, Option<String>>(4)?,
2744            row.get::<_, i64>(5)?,
2745        ))
2746    }) else {
2747        return;
2748    };
2749    for row in rows.flatten() {
2750        let (id, cwd, title, updated_at, model_config, message_count) = row;
2751        let cwd = PathBuf::from(cwd);
2752        if workspace.is_some_and(|wanted| !wanted.admits(&cwd)) {
2753            continue;
2754        }
2755        let model = model_config
2756            .as_deref()
2757            .and_then(|value| serde_json::from_str::<Value>(value).ok())
2758            .and_then(|value| {
2759                value
2760                    .get("model_name")
2761                    .or_else(|| value.get("modelName"))
2762                    .and_then(Value::as_str)
2763                    .map(str::to_string)
2764            });
2765        let updated_at_ms = crate::sidecar::rfc3339_to_ms(&updated_at)
2766            .or_else(|| {
2767                // SQLite's CURRENT_TIMESTAMP uses `YYYY-MM-DD HH:MM:SS`.
2768                crate::sidecar::rfc3339_to_ms(&format!("{}Z", updated_at.replace(' ', "T")))
2769            })
2770            .and_then(|value| u64::try_from(value).ok());
2771        found.push(SessionDescriptor {
2772            locator: SessionLocator {
2773                harness: HarnessId::from(HarnessId::GOOSE),
2774                session_id: id.clone(),
2775                storage: StorageLocator::Sqlite {
2776                    path: db.clone(),
2777                    selector: id,
2778                },
2779            },
2780            cwd: Some(cwd),
2781            title: (!title.trim().is_empty()).then_some(title),
2782            preview_candidates: Vec::new(),
2783            latest_message_candidates: Vec::new(),
2784            updated_at_ms,
2785            message_count: usize::try_from(message_count).ok(),
2786            model,
2787            parent_session_id: None,
2788            child_session_count: 0,
2789            nouns: OrchestrationNouns::default(),
2790        });
2791    }
2792}
2793
2794const LATEST_PREVIEW_CANDIDATES: usize = 8;
2795const TOPIC_PREVIEW_HEAD_BYTES: u64 = 512 * 1024;
2796const LATEST_PREVIEW_TAIL_BYTES: u64 = 512 * 1024;
2797const LATEST_PREVIEW_MAX_BYTES: u64 = 4 * 1024 * 1024;
2798
2799fn topic_message_candidates(locator: &SessionLocator) -> Result<Vec<SessionPreviewCandidate>> {
2800    match &locator.storage {
2801        StorageLocator::File { path }
2802            if matches!(
2803                locator.harness.as_str(),
2804                HarnessId::CLAUDE_CODE | HarnessId::CODEX
2805            ) =>
2806        {
2807            topic_file_message_candidates(path, locator.harness.as_str())
2808        }
2809        _ => Ok(Vec::new()),
2810    }
2811}
2812
2813fn codex_history_topics(
2814    sessions_root: &Path,
2815    sessions: &[SessionDescriptor],
2816) -> Result<HashMap<String, Vec<SessionPreviewCandidate>>> {
2817    let wanted: HashSet<&str> = sessions
2818        .iter()
2819        .filter(|descriptor| descriptor.locator.harness.as_str() == HarnessId::CODEX)
2820        .map(|descriptor| descriptor.locator.session_id.as_str())
2821        .collect();
2822    if wanted.is_empty() {
2823        return Ok(HashMap::new());
2824    }
2825    let Some(root) = sessions_root.parent() else {
2826        return Ok(HashMap::new());
2827    };
2828    let file = match File::open(root.join("history.jsonl")) {
2829        Ok(file) => file,
2830        Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(HashMap::new()),
2831        Err(error) => return Err(error.into()),
2832    };
2833    let mut topics = HashMap::new();
2834    for line in BufReader::new(file).lines() {
2835        let Ok(value) = serde_json::from_str::<Value>(&line?) else {
2836            continue;
2837        };
2838        let Some(session_id) = value.get("session_id").and_then(Value::as_str) else {
2839            continue;
2840        };
2841        if !wanted.contains(session_id) || topics.contains_key(session_id) {
2842            continue;
2843        }
2844        let mut candidates = Vec::new();
2845        push_message_candidate(&mut candidates, "user", value.get("text"), HashMap::new());
2846        if !candidates.is_empty() {
2847            topics.insert(session_id.to_string(), candidates);
2848            if topics.len() == wanted.len() {
2849                break;
2850            }
2851        }
2852    }
2853    Ok(topics)
2854}
2855
2856fn latest_message_candidates(locator: &SessionLocator) -> Result<Vec<SessionPreviewCandidate>> {
2857    match &locator.storage {
2858        // A Hermes locator addresses one session inside the whole-store `state.db`; the file's
2859        // tail is SQLite pages, not transcript lines, so the preview comes from the store by id.
2860        StorageLocator::File { path } | StorageLocator::Sqlite { path, .. }
2861            if locator.harness.as_str() == HarnessId::HERMES =>
2862        {
2863            latest_hermes_message_candidates(path, &locator.session_id)
2864        }
2865        StorageLocator::File { path } => {
2866            latest_file_message_candidates(path, locator.harness.as_str())
2867        }
2868        StorageLocator::Sqlite { path, selector }
2869            if locator.harness.as_str() == HarnessId::OPENCODE =>
2870        {
2871            latest_opencode_message_candidates(path, selector)
2872        }
2873        StorageLocator::Sqlite { path, selector }
2874            if locator.harness.as_str() == HarnessId::GOOSE =>
2875        {
2876            latest_goose_message_candidates(path, selector)
2877        }
2878        StorageLocator::Sqlite { .. } => Ok(Vec::new()),
2879    }
2880}
2881
2882fn topic_file_message_candidates(
2883    path: &Path,
2884    harness: &str,
2885) -> Result<Vec<SessionPreviewCandidate>> {
2886    let mut file = File::open(path)?;
2887    let mut bytes = Vec::with_capacity(TOPIC_PREVIEW_HEAD_BYTES as usize);
2888    file.by_ref()
2889        .take(TOPIC_PREVIEW_HEAD_BYTES)
2890        .read_to_end(&mut bytes)?;
2891    if file.metadata()?.len() > TOPIC_PREVIEW_HEAD_BYTES {
2892        if let Some(newline) = bytes.iter().rposition(|byte| *byte == b'\n') {
2893            bytes.truncate(newline);
2894        }
2895    }
2896    let text = String::from_utf8(bytes).map_err(|_| {
2897        Error::Other(format!(
2898            "{} contains non-UTF-8 data in its topic-preview window",
2899            path.display()
2900        ))
2901    })?;
2902    if harness == HarnessId::CODEX {
2903        return Ok(codex_preview_candidates(text.lines(), false));
2904    }
2905    let mut candidates = Vec::new();
2906    for line in text.lines() {
2907        let Ok(value) = serde_json::from_str::<Value>(line) else {
2908            continue;
2909        };
2910        push_topic_message_candidate(&mut candidates, harness, &value);
2911        if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
2912            break;
2913        }
2914    }
2915    Ok(candidates)
2916}
2917
2918fn latest_file_message_candidates(
2919    path: &Path,
2920    harness: &str,
2921) -> Result<Vec<SessionPreviewCandidate>> {
2922    let mut candidates =
2923        latest_file_message_candidates_with_limit(path, harness, LATEST_PREVIEW_TAIL_BYTES)?;
2924    if candidates.is_empty() {
2925        candidates =
2926            latest_file_message_candidates_with_limit(path, harness, LATEST_PREVIEW_MAX_BYTES)?;
2927    }
2928    Ok(candidates)
2929}
2930
2931fn latest_file_message_candidates_with_limit(
2932    path: &Path,
2933    harness: &str,
2934    byte_limit: u64,
2935) -> Result<Vec<SessionPreviewCandidate>> {
2936    let mut file = File::open(path)?;
2937    let file_len = file.metadata()?.len();
2938    let start = file_len.saturating_sub(byte_limit);
2939    file.seek(SeekFrom::Start(start))?;
2940    let mut bytes = Vec::with_capacity((file_len - start) as usize);
2941    file.read_to_end(&mut bytes)?;
2942    if start > 0 {
2943        if let Some(newline) = bytes.iter().position(|byte| *byte == b'\n') {
2944            bytes.drain(..=newline);
2945        } else {
2946            return Ok(Vec::new());
2947        }
2948    }
2949    let text = String::from_utf8(bytes).map_err(|_| {
2950        Error::Other(format!(
2951            "{} contains non-UTF-8 data in its list-preview window",
2952            path.display()
2953        ))
2954    })?;
2955    if harness == HarnessId::CODEX {
2956        return Ok(codex_preview_candidates(text.lines().rev(), true));
2957    }
2958    let mut candidates = Vec::new();
2959    for line in text.lines().rev() {
2960        let Ok(value) = serde_json::from_str::<Value>(line) else {
2961            continue;
2962        };
2963        let (role, content, metadata) = match harness {
2964            HarnessId::CLAUDE_CODE => {
2965                let role = value.get("type").and_then(Value::as_str);
2966                if !matches!(role, Some("user" | "assistant")) {
2967                    continue;
2968                }
2969                let metadata = if role == Some("user") {
2970                    crate::session::claude_user_provenance(&value)
2971                        .into_iter()
2972                        .collect()
2973                } else {
2974                    HashMap::new()
2975                };
2976                (
2977                    role.unwrap_or_default(),
2978                    value
2979                        .get("message")
2980                        .and_then(|message| message.get("content")),
2981                    metadata,
2982                )
2983            }
2984            HarnessId::PI => {
2985                if value.get("type").and_then(Value::as_str) != Some("message") {
2986                    continue;
2987                }
2988                let message = value.get("message").unwrap_or(&Value::Null);
2989                let Some(role @ ("user" | "assistant")) =
2990                    message.get("role").and_then(Value::as_str)
2991                else {
2992                    continue;
2993                };
2994                (role, message.get("content"), HashMap::new())
2995            }
2996            HarnessId::GEMINI => {
2997                let Some(kind @ ("user" | "gemini")) = value.get("type").and_then(Value::as_str)
2998                else {
2999                    continue;
3000                };
3001                (
3002                    if kind == "gemini" {
3003                        "assistant"
3004                    } else {
3005                        "user"
3006                    },
3007                    value.get("content"),
3008                    HashMap::new(),
3009                )
3010            }
3011            HarnessId::GROK => {
3012                let Some(role @ ("user" | "assistant")) = value.get("type").and_then(Value::as_str)
3013                else {
3014                    continue;
3015                };
3016                (role, value.get("content"), HashMap::new())
3017            }
3018            HarnessId::SUPERCODE => {
3019                let Some(role @ ("user" | "assistant")) = value.get("role").and_then(Value::as_str)
3020                else {
3021                    continue;
3022                };
3023                (role, value.get("content"), HashMap::new())
3024            }
3025            _ => continue,
3026        };
3027        let mut metadata = metadata;
3028        if matches!(harness, HarnessId::CLAUDE_CODE | HarnessId::CODEX) {
3029            if let Some(timestamp) = value.get("timestamp").and_then(Value::as_str) {
3030                metadata.insert("timestamp".to_string(), timestamp.to_string());
3031            }
3032        }
3033        push_message_candidate_with_cursor(
3034            &mut candidates,
3035            role,
3036            content,
3037            metadata,
3038            Some(message_candidate_cursor(harness, &value)),
3039        );
3040        if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3041            break;
3042        }
3043    }
3044    Ok(candidates)
3045}
3046
3047struct CodexPreviewRecord {
3048    native: Value,
3049    role: String,
3050    text: String,
3051}
3052
3053// Pair adjacent conversational event/response mirrors one-to-one, within the
3054// bytes already read. This is intentionally narrower than the full codec's
3055// global assistant-text dedup: repeated same-kind records and separate pairs
3056// remain separate turns. Compare full display text BEFORE the 4096-character
3057// cap; prefer the canonical response's cursor/timestamp in either scan direction.
3058fn codex_preview_candidates<'a>(
3059    lines: impl Iterator<Item = &'a str>,
3060    latest: bool,
3061) -> Vec<SessionPreviewCandidate> {
3062    let mut candidates = Vec::new();
3063    let mut pending: Option<CodexPreviewRecord> = None;
3064    for line in lines {
3065        let Ok(native) = serde_json::from_str::<Value>(line) else {
3066            continue;
3067        };
3068        let Some((role, content)) = codex_preview_message(&native) else {
3069            continue;
3070        };
3071        let Some(text) = display_text(content).filter(|text| !text.trim().is_empty()) else {
3072            continue;
3073        };
3074        let current = CodexPreviewRecord {
3075            role: role.to_string(),
3076            text,
3077            native,
3078        };
3079        if let Some(previous) = pending.take() {
3080            if previous.role == current.role
3081                && previous.text == current.text
3082                && previous.native.get("type") != current.native.get("type")
3083            {
3084                let canonical = if previous.native.get("type").and_then(Value::as_str)
3085                    == Some("response_item")
3086                {
3087                    previous
3088                } else {
3089                    current
3090                };
3091                push_codex_preview_candidate(&mut candidates, canonical, latest);
3092            } else {
3093                push_codex_preview_candidate(&mut candidates, previous, latest);
3094                pending = Some(current);
3095            }
3096        } else {
3097            pending = Some(current);
3098        }
3099        if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3100            break;
3101        }
3102    }
3103    if let Some(last) = pending {
3104        push_codex_preview_candidate(&mut candidates, last, latest);
3105    }
3106    candidates
3107}
3108
3109fn push_codex_preview_candidate(
3110    candidates: &mut Vec<SessionPreviewCandidate>,
3111    record: CodexPreviewRecord,
3112    latest: bool,
3113) {
3114    let mut metadata = HashMap::new();
3115    if latest {
3116        if let Some(timestamp) = record.native.get("timestamp").and_then(Value::as_str) {
3117            metadata.insert("timestamp".to_string(), timestamp.to_string());
3118        }
3119    }
3120    let cursor = latest.then(|| message_candidate_cursor(HarnessId::CODEX, &record.native));
3121    push_message_candidate_with_cursor(
3122        candidates,
3123        &record.role,
3124        Some(&Value::String(record.text)),
3125        metadata,
3126        cursor,
3127    );
3128}
3129
3130// Codex collab rollouts can carry narration only as event_msg records.
3131fn codex_preview_message(value: &Value) -> Option<(&str, Option<&Value>)> {
3132    let payload = value.get("payload")?;
3133    match (
3134        value.get("type").and_then(Value::as_str)?,
3135        payload.get("type").and_then(Value::as_str)?,
3136    ) {
3137        ("response_item", "message") => {
3138            let role @ ("user" | "assistant") = payload.get("role").and_then(Value::as_str)? else {
3139                return None;
3140            };
3141            Some((role, payload.get("content")))
3142        }
3143        ("event_msg", "user_message") => Some(("user", payload.get("message"))),
3144        ("event_msg", "agent_message") => Some(("assistant", payload.get("message"))),
3145        _ => None,
3146    }
3147}
3148
3149fn push_topic_message_candidate(
3150    candidates: &mut Vec<SessionPreviewCandidate>,
3151    harness: &str,
3152    value: &Value,
3153) {
3154    let (role, content, metadata) = match harness {
3155        HarnessId::CLAUDE_CODE => {
3156            let role = value.get("type").and_then(Value::as_str);
3157            if !matches!(role, Some("user" | "assistant")) {
3158                return;
3159            }
3160            let metadata = if role == Some("user") {
3161                crate::session::claude_user_provenance(value)
3162                    .into_iter()
3163                    .collect()
3164            } else {
3165                HashMap::new()
3166            };
3167            (
3168                role.unwrap_or_default(),
3169                value
3170                    .get("message")
3171                    .and_then(|message| message.get("content")),
3172                metadata,
3173            )
3174        }
3175        _ => return,
3176    };
3177    push_message_candidate(candidates, role, content, metadata);
3178}
3179
3180fn latest_opencode_message_candidates(
3181    path: &Path,
3182    session_id: &str,
3183) -> Result<Vec<SessionPreviewCandidate>> {
3184    let connection = Connection::open_with_flags(
3185        path,
3186        rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
3187    )
3188    .map_err(|error| Error::Other(format!("failed to open list-preview store: {error}")))?;
3189    let mut statement = connection
3190        .prepare(
3191            "SELECT m.data, p.data FROM message m JOIN part p ON p.message_id = m.id \
3192         WHERE m.session_id = ?1 ORDER BY m.time_created DESC, p.time_created DESC LIMIT 32",
3193        )
3194        .map_err(|error| Error::Other(format!("failed to prepare list-preview query: {error}")))?;
3195    let rows = statement
3196        .query_map([session_id], |row| {
3197            Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
3198        })
3199        .map_err(|error| Error::Other(format!("failed to read list-preview rows: {error}")))?;
3200    let mut candidates = Vec::new();
3201    for row in rows.flatten() {
3202        let (Ok(message), Ok(part)) = (
3203            serde_json::from_str::<Value>(&row.0),
3204            serde_json::from_str::<Value>(&row.1),
3205        ) else {
3206            continue;
3207        };
3208        let Some(role @ ("user" | "assistant")) = message.get("role").and_then(Value::as_str)
3209        else {
3210            continue;
3211        };
3212        if part.get("type").and_then(Value::as_str) != Some("text") {
3213            continue;
3214        }
3215        push_message_candidate(&mut candidates, role, part.get("text"), HashMap::new());
3216        if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3217            break;
3218        }
3219    }
3220    Ok(candidates)
3221}
3222
3223fn latest_hermes_message_candidates(
3224    path: &Path,
3225    session_id: &str,
3226) -> Result<Vec<SessionPreviewCandidate>> {
3227    let connection = Connection::open_with_flags(
3228        path,
3229        rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
3230    )
3231    .map_err(|error| Error::Other(format!("failed to open list-preview store: {error}")))?;
3232    let mut statement = connection
3233        .prepare(
3234            "SELECT role, content FROM messages WHERE session_id = ?1 AND active = 1 \
3235             AND role IN ('user', 'assistant') AND content IS NOT NULL AND content != '' \
3236             ORDER BY timestamp DESC, id DESC LIMIT 32",
3237        )
3238        .map_err(|error| Error::Other(format!("failed to prepare list-preview query: {error}")))?;
3239    let rows = statement
3240        .query_map([session_id], |row| {
3241            Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
3242        })
3243        .map_err(|error| Error::Other(format!("failed to read list-preview rows: {error}")))?;
3244    let mut candidates = Vec::new();
3245    for (role, content) in rows.flatten() {
3246        push_message_candidate(
3247            &mut candidates,
3248            &role,
3249            Some(&Value::String(content)),
3250            HashMap::new(),
3251        );
3252        if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3253            break;
3254        }
3255    }
3256    Ok(candidates)
3257}
3258
3259fn latest_goose_message_candidates(
3260    path: &Path,
3261    session_id: &str,
3262) -> Result<Vec<SessionPreviewCandidate>> {
3263    let connection = Connection::open_with_flags(
3264        path,
3265        rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
3266    )
3267    .map_err(|error| Error::Other(format!("failed to open list-preview store: {error}")))?;
3268    let mut statement = connection
3269        .prepare(
3270            "SELECT role, content_json FROM messages WHERE session_id = ?1 \
3271         ORDER BY created_timestamp DESC, id DESC LIMIT 16",
3272        )
3273        .map_err(|error| Error::Other(format!("failed to prepare list-preview query: {error}")))?;
3274    let rows = statement
3275        .query_map([session_id], |row| {
3276            Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
3277        })
3278        .map_err(|error| Error::Other(format!("failed to read list-preview rows: {error}")))?;
3279    let mut candidates = Vec::new();
3280    for row in rows.flatten() {
3281        let (role, content) = row;
3282        if !matches!(role.as_str(), "user" | "assistant") {
3283            continue;
3284        }
3285        let Ok(content) = serde_json::from_str::<Value>(&content) else {
3286            continue;
3287        };
3288        push_message_candidate(&mut candidates, &role, Some(&content), HashMap::new());
3289        if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3290            break;
3291        }
3292    }
3293    Ok(candidates)
3294}
3295
3296fn push_message_candidate(
3297    candidates: &mut Vec<SessionPreviewCandidate>,
3298    role: &str,
3299    content: Option<&Value>,
3300    metadata: HashMap<String, String>,
3301) {
3302    push_message_candidate_with_cursor(candidates, role, content, metadata, None);
3303}
3304
3305fn push_message_candidate_with_cursor(
3306    candidates: &mut Vec<SessionPreviewCandidate>,
3307    role: &str,
3308    content: Option<&Value>,
3309    metadata: HashMap<String, String>,
3310    cursor: Option<String>,
3311) {
3312    if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3313        return;
3314    }
3315    let Some(text) = display_text(content).filter(|text| !text.trim().is_empty()) else {
3316        return;
3317    };
3318    const MAX_CHARS: usize = 4_096;
3319    candidates.push(SessionPreviewCandidate {
3320        cursor,
3321        role: role.to_string(),
3322        content: text.chars().take(MAX_CHARS).collect(),
3323        metadata,
3324    });
3325}
3326
3327fn message_candidate_cursor(harness: &str, value: &Value) -> String {
3328    let native_identity = value
3329        .get("uuid")
3330        .or_else(|| value.get("id"))
3331        .or_else(|| value.pointer("/message/id"))
3332        .or_else(|| value.pointer("/payload/id"))
3333        .and_then(Value::as_str)
3334        .or_else(|| value.get("timestamp").and_then(Value::as_str));
3335    let mut hasher = blake3::Hasher::new();
3336    hasher.update(b"supercode.session-preview-cursor.v1\0");
3337    hasher.update(harness.as_bytes());
3338    hasher.update(b"\0");
3339    if let Some(identity) = native_identity {
3340        hasher.update(identity.as_bytes());
3341    } else {
3342        // Some formats do not publish message ids. Hashing the complete native
3343        // record is still stable across discovery refreshes and reveals none
3344        // of the record itself to an untrusted presentation surface.
3345        hasher.update(value.to_string().as_bytes());
3346    }
3347    format!("v1:{}", &hasher.finalize().to_hex()[..24])
3348}
3349
3350fn fill_string(target: &mut Option<String>, value: Option<&Value>) {
3351    if target.is_none() {
3352        *target = value.and_then(Value::as_str).map(str::to_owned);
3353    }
3354}
3355
3356fn fill_path(target: &mut Option<PathBuf>, value: Option<&Value>) {
3357    if target.is_none() {
3358        *target = value.and_then(Value::as_str).map(PathBuf::from);
3359    }
3360}
3361
3362/// What the END of a JSONL transcript says about a session: when it last took
3363/// a turn, and which model took it.
3364#[derive(Default)]
3365struct TailFacts {
3366    /// Newest `timestamp` among the last records, as Unix epoch milliseconds.
3367    last_turn_ms: Option<u64>,
3368    /// Model on the LAST turn that named one.
3369    model: Option<String>,
3370}
3371
3372/// Read [`TailFacts`] from the last records of a JSONL transcript.
3373///
3374/// Both facts have to come from the conversation rather than from cheaper
3375/// stand-ins. Recency is not the file's mtime: harnesses touch a transcript
3376/// without saying anything (see [`SessionDescriptor::updated_at_ms`]). The model
3377/// is not the one in the header either — every dialect here records the model
3378/// per turn, so a mid-session switch (`/model`) leaves the opening record naming
3379/// a model the session has not used for hours.
3380///
3381/// Reading whole files is not an option — one catalog holds thousands of
3382/// sessions and a single transcript runs to tens of megabytes — so this seeks to
3383/// the end and walks backwards over a bounded window, which is where an
3384/// append-only log keeps its newest records. The window grows only while no
3385/// timestamp has been found, and gives up at [`TAIL_SCAN_LIMIT`]; a model the
3386/// window does not reach stays `None` and the caller keeps the header's.
3387fn tail_facts(path: &Path, harness: &str) -> TailFacts {
3388    let mut facts = TailFacts::default();
3389    let Ok(mut file) = File::open(path) else {
3390        return facts;
3391    };
3392    let Ok(len) = file.metadata().map(|meta| meta.len()) else {
3393        return facts;
3394    };
3395    let mut window = TAIL_SCAN_START.min(len);
3396    loop {
3397        if file.seek(SeekFrom::Start(len - window)).is_err() {
3398            return facts;
3399        }
3400        let Ok(size) = usize::try_from(window) else {
3401            return facts;
3402        };
3403        let mut buf = vec![0u8; size];
3404        if file.read_exact(&mut buf).is_err() {
3405            return facts;
3406        }
3407        // A transcript is append-only, so the newest record is the LAST one:
3408        // walk line boundaries backwards from the end and decode one record at a
3409        // time, stopping as soon as both facts are in hand. Decoding the whole
3410        // window instead would put a UTF-8 validation of every byte of every
3411        // transcript in the catalog on the path of one `discover`.
3412        //
3413        // `end` is the exclusive end of the line under inspection; the scan stops
3414        // at `floor`, because a window that starts mid-file almost certainly
3415        // starts mid-record and that partial first line belongs to the next,
3416        // wider window.
3417        let floor = if window < len {
3418            buf.iter().position(|byte| *byte == b'\n').map(|at| at + 1)
3419        } else {
3420            Some(0)
3421        };
3422        if let Some(floor) = floor {
3423            let mut end = buf.len();
3424            while end > floor && !(facts.last_turn_ms.is_some() && facts.model.is_some()) {
3425                let start = buf[floor..end]
3426                    .iter()
3427                    .rposition(|byte| *byte == b'\n')
3428                    .map_or(floor, |at| floor + at + 1);
3429                if let Ok(record) = std::str::from_utf8(&buf[start..end])
3430                    .map_err(|_| ())
3431                    .and_then(|line| serde_json::from_str::<Value>(line).map_err(|_| ()))
3432                {
3433                    if facts.last_turn_ms.is_none() {
3434                        facts.last_turn_ms = record_timestamp(&record, harness)
3435                            .and_then(crate::sidecar::rfc3339_to_ms)
3436                            .and_then(|millis| u64::try_from(millis).ok());
3437                    }
3438                    if facts.model.is_none() {
3439                        facts.model = record_model(&record, harness).map(str::to_owned);
3440                    }
3441                }
3442                end = start.saturating_sub(1);
3443            }
3444        }
3445        if facts.last_turn_ms.is_some() || window >= len || window >= TAIL_SCAN_LIMIT {
3446            return facts;
3447        }
3448        window = (window * 2).min(len).min(TAIL_SCAN_LIMIT);
3449    }
3450}
3451
3452/// When one transcript record was written, in the dialect that wrote it.
3453/// Every dialect discovery reads stamps its records with an RFC3339 UTC string;
3454/// only the key differs.
3455fn record_timestamp<'a>(record: &'a Value, harness: &str) -> Option<&'a str> {
3456    let key = match harness {
3457        HarnessId::SUPERCODE => "ts",
3458        _ => "timestamp",
3459    };
3460    record.get(key)?.as_str()
3461}
3462
3463/// The model one transcript record names, in the dialect that wrote it.
3464/// Mirrors the model arms of [`update_header_meta`], read newest-first.
3465fn record_model<'a>(record: &'a Value, harness: &str) -> Option<&'a str> {
3466    match harness {
3467        HarnessId::CLAUDE_CODE | HarnessId::PI => record.get("message")?.get("model")?.as_str(),
3468        HarnessId::CODEX => {
3469            if record.get("type")?.as_str()? != "turn_context" {
3470                return None;
3471            }
3472            record.get("payload")?.get("model")?.as_str()
3473        }
3474        _ => None,
3475    }
3476}
3477
3478/// Bytes read from a transcript's tail on the first attempt: comfortably more
3479/// than one record in every dialect discovery reads.
3480const TAIL_SCAN_START: u64 = 16 * 1024;
3481
3482/// Where the widening tail scan stops. A transcript whose last mebibyte holds
3483/// no timestamped record is not one whose recency this can honestly report.
3484const TAIL_SCAN_LIMIT: u64 = 1024 * 1024;
3485
3486fn modified_ms(path: &Path) -> Option<u64> {
3487    fs::metadata(path)
3488        .ok()?
3489        .modified()
3490        .ok()?
3491        .duration_since(UNIX_EPOCH)
3492        .ok()
3493        .and_then(|duration| u64::try_from(duration.as_millis()).ok())
3494}
3495
3496/// A discovery's workspace filter: the exact recorded folder (the default), or
3497/// the folder and every folder under it (`workspace_subtree`).
3498pub struct WorkspaceScope {
3499    wanted: PathBuf,
3500    /// Comparison keys of `wanted` (canonical and lexical) for subtree mode.
3501    roots: Vec<PathBuf>,
3502    subtree: bool,
3503    /// Subtree verdicts per recorded cwd: many sessions share one folder.
3504    judged: std::sync::Mutex<HashMap<PathBuf, bool>>,
3505}
3506
3507impl WorkspaceScope {
3508    /// The exact-folder filter every caller had before `workspace_subtree`.
3509    pub fn exact(wanted: &Path) -> Self {
3510        Self {
3511            wanted: wanted.to_path_buf(),
3512            roots: Vec::new(),
3513            subtree: false,
3514            judged: std::sync::Mutex::new(HashMap::new()),
3515        }
3516    }
3517
3518    /// The folder and every folder under it.
3519    pub fn subtree(root: &Path) -> Self {
3520        Self {
3521            wanted: root.to_path_buf(),
3522            roots: path_comparison_keys(root),
3523            subtree: true,
3524            judged: std::sync::Mutex::new(HashMap::new()),
3525        }
3526    }
3527
3528    /// The scope a query asks for, if it filters by workspace at all.
3529    pub fn of(query: &DiscoveryQuery) -> Option<Self> {
3530        query.workspace.as_deref().map(|wanted| {
3531            if query.workspace_subtree {
3532                Self::subtree(wanted)
3533            } else {
3534                Self::exact(wanted)
3535            }
3536        })
3537    }
3538
3539    /// Whether a session whose recorded working directory is `recorded`
3540    /// belongs to this scope. A relative recorded cwd never does, nor (on
3541    /// Windows) one without a drive or UNC prefix: neither says where the
3542    /// session ran.
3543    pub fn admits(&self, recorded: &Path) -> bool {
3544        if !recorded.is_absolute() {
3545            return false;
3546        }
3547        if !self.subtree {
3548            return same_path(recorded, &self.wanted);
3549        }
3550        let mut judged = self
3551            .judged
3552            .lock()
3553            .unwrap_or_else(std::sync::PoisonError::into_inner);
3554        *judged.entry(recorded.to_path_buf()).or_insert_with(|| {
3555            path_comparison_keys(recorded)
3556                .iter()
3557                .any(|key| self.roots.iter().any(|root| key.starts_with(root)))
3558        })
3559    }
3560
3561    /// Stores that historically compared the recorded cwd literally (Hermes,
3562    /// OpenClaw) keep that exact comparison; a subtree scope applies
3563    /// [`Self::admits`].
3564    fn admits_literal(&self, recorded: Option<&Path>) -> bool {
3565        if self.subtree {
3566            recorded.is_some_and(|cwd| self.admits(cwd))
3567        } else {
3568            recorded == Some(self.wanted.as_path())
3569        }
3570    }
3571}
3572
3573/// The forms of `path` a subtree comparison accepts: its canonical form when
3574/// it exists on this machine (resolving symlinks and aliases) and its lexical
3575/// absolute form (a deleted folder still names where a session ran). Windows
3576/// verbatim prefixes are dropped and case is folded there, as the filesystem
3577/// does.
3578fn path_comparison_keys(path: &Path) -> Vec<PathBuf> {
3579    let mut keys = Vec::with_capacity(2);
3580    if let Ok(canonical) = fs::canonicalize(path) {
3581        keys.push(comparison_key(&canonical));
3582    }
3583    let lexical = comparison_key(&normalize_path(path));
3584    if !keys.contains(&lexical) {
3585        keys.push(lexical);
3586    }
3587    keys
3588}
3589
3590#[cfg(windows)]
3591fn comparison_key(path: &Path) -> PathBuf {
3592    let text = path.to_string_lossy();
3593    let text = if let Some(unc) = text.strip_prefix(r"\\?\UNC\") {
3594        format!(r"\\{unc}")
3595    } else if let Some(local) = text.strip_prefix(r"\\?\") {
3596        local.to_string()
3597    } else {
3598        text.into_owned()
3599    };
3600    PathBuf::from(text.to_lowercase())
3601}
3602
3603#[cfg(not(windows))]
3604fn comparison_key(path: &Path) -> PathBuf {
3605    path.to_path_buf()
3606}
3607
3608/// A workspace filter is satisfiable only by a session whose RECORDED working
3609/// directory is absolute. A relative recorded cwd (OpenCode has shipped
3610/// literal `"."` session rows) carries no information about where the session
3611/// ran; resolving it against the discoverer's own current directory made such
3612/// a session match every workspace discovery happened to run from.
3613fn recorded_cwd_matches(recorded: &Path, wanted: &Path) -> bool {
3614    recorded.is_absolute() && same_path(recorded, wanted)
3615}
3616
3617fn same_path(left: &Path, right: &Path) -> bool {
3618    match (fs::canonicalize(left), fs::canonicalize(right)) {
3619        (Ok(left), Ok(right)) => left == right,
3620        _ => normalize_path(left) == normalize_path(right),
3621    }
3622}
3623
3624fn normalize_path(path: &Path) -> PathBuf {
3625    let absolute = if path.is_absolute() {
3626        path.to_path_buf()
3627    } else {
3628        std::env::current_dir()
3629            .unwrap_or_else(|_| PathBuf::from("."))
3630            .join(path)
3631    };
3632    let mut normalized = PathBuf::new();
3633    for component in absolute.components() {
3634        match component {
3635            Component::CurDir => {}
3636            Component::ParentDir => {
3637                normalized.pop();
3638            }
3639            other => normalized.push(other.as_os_str()),
3640        }
3641    }
3642    normalized
3643}
3644
3645#[cfg(test)]
3646mod tests {
3647    use super::*;
3648    use std::io::Write;
3649    use std::time::{SystemTime, UNIX_EPOCH};
3650
3651    fn temp_dir(label: &str) -> PathBuf {
3652        let nonce = SystemTime::now()
3653            .duration_since(UNIX_EPOCH)
3654            .unwrap()
3655            .as_nanos();
3656        let path = std::env::temp_dir().join(format!(
3657            "supercode-catalog-{label}-{}-{nonce}",
3658            std::process::id()
3659        ));
3660        fs::create_dir_all(&path).unwrap();
3661        path
3662    }
3663
3664    #[test]
3665    fn codex_history_index_reads_appends_and_repairs_replacements() {
3666        let root = temp_dir("codex-history-index");
3667        let sessions = root.join("sessions");
3668        fs::create_dir_all(&sessions).unwrap();
3669        let history = root.join("history.jsonl");
3670        fs::write(
3671            &history,
3672            "{\"session_id\":\"alpha\",\"text\":\"first topic\"}\n",
3673        )
3674        .unwrap();
3675
3676        let mut index = CodexHistoryTopicIndex::new(&sessions);
3677        assert_eq!(index.refresh().unwrap(), BTreeSet::from(["alpha".into()]));
3678        assert_eq!(index.topics["alpha"][0].content, "first topic");
3679        assert!(index.refresh().unwrap().is_empty());
3680
3681        let mut file = fs::OpenOptions::new().append(true).open(&history).unwrap();
3682        write!(
3683            file,
3684            "{{\"session_id\":\"alpha\",\"text\":\"later topic\"}}\n\
3685             {{\"session_id\":\"beta\",\"text\":\"second topic\"}}\n"
3686        )
3687        .unwrap();
3688        file.flush().unwrap();
3689        assert_eq!(index.refresh().unwrap(), BTreeSet::from(["beta".into()]));
3690        assert_eq!(index.topics["alpha"][0].content, "first topic");
3691        assert_eq!(index.topics["beta"][0].content, "second topic");
3692
3693        fs::write(
3694            &history,
3695            "{\"session_id\":\"gamma\",\"text\":\"replacement\"}\n",
3696        )
3697        .unwrap();
3698        assert_eq!(
3699            index.refresh().unwrap(),
3700            BTreeSet::from(["alpha".into(), "beta".into(), "gamma".into()])
3701        );
3702        assert!(!index.topics.contains_key("alpha"));
3703        assert_eq!(index.topics["gamma"][0].content, "replacement");
3704
3705        fs::remove_dir_all(root).ok();
3706    }
3707
3708    #[test]
3709    fn codex_history_index_retains_an_incomplete_appended_record() {
3710        let root = temp_dir("codex-history-partial");
3711        let sessions = root.join("sessions");
3712        fs::create_dir_all(&sessions).unwrap();
3713        let history = root.join("history.jsonl");
3714        fs::write(&history, "{\"session_id\":\"partial\",\"text\":\"hel").unwrap();
3715
3716        let mut index = CodexHistoryTopicIndex::new(&sessions);
3717        assert!(index.refresh().unwrap().is_empty());
3718        let mut file = fs::OpenOptions::new().append(true).open(&history).unwrap();
3719        writeln!(file, "lo\"}}").unwrap();
3720        file.flush().unwrap();
3721
3722        assert_eq!(index.refresh().unwrap(), BTreeSet::from(["partial".into()]));
3723        assert_eq!(index.topics["partial"][0].content, "hello");
3724        fs::remove_dir_all(root).ok();
3725    }
3726
3727    #[test]
3728    fn cached_codex_history_enrichment_matches_stateless_discovery() {
3729        let root = temp_dir("codex-history-parity");
3730        let sessions = root.join("sessions");
3731        let workspace = root.join("workspace");
3732        fs::create_dir_all(&sessions).unwrap();
3733        fs::create_dir_all(&workspace).unwrap();
3734        fs::write(
3735            sessions.join("rollout.jsonl"),
3736            format!(
3737                "{{\"type\":\"session_meta\",\"payload\":{{\"id\":\"alpha\",\"cwd\":{}}}}}\n{{\"type\":\"turn_context\",\"payload\":{{\"cwd\":{},\"model\":\"gpt-test\"}}}}\n{{\"type\":\"response_item\",\"payload\":{{\"type\":\"message\",\"role\":\"user\",\"content\":[{{\"type\":\"input_text\",\"text\":\"transcript fallback\"}}]}}}}\n",
3738                serde_json::to_string(&workspace.to_string_lossy()).unwrap(),
3739                serde_json::to_string(&workspace.to_string_lossy()).unwrap(),
3740            ),
3741        )
3742        .unwrap();
3743        fs::write(
3744            root.join("history.jsonl"),
3745            "{\"session_id\":\"alpha\",\"text\":\"history topic\"}\n",
3746        )
3747        .unwrap();
3748        let query = DiscoveryQuery {
3749            harnesses: vec![HarnessId::from(HarnessId::CODEX)],
3750            homes: HarnessHomes {
3751                codex: sessions.clone(),
3752                ..HarnessHomes::default()
3753            },
3754            include_topic_candidates: true,
3755            ..DiscoveryQuery::default()
3756        };
3757        let catalog = HarnessCatalog::new();
3758        let projected = catalog
3759            .project_index(&query, catalog.discover_raw_index(&query))
3760            .unwrap();
3761        let expected = catalog
3762            .enrich_index_page(&query, projected.clone())
3763            .unwrap();
3764        let mut history = CodexHistoryTopicIndex::new(&sessions);
3765        history.refresh().unwrap();
3766        let actual = catalog
3767            .enrich_index_page_with_codex_history(&query, projected, &history)
3768            .unwrap();
3769
3770        assert_eq!(actual, expected);
3771        assert_eq!(actual[0].preview_candidates[0].content, "history topic");
3772        fs::remove_dir_all(root).ok();
3773    }
3774
3775    #[test]
3776    fn locator_json_round_trip_preserves_sqlite_selector() {
3777        let locator = SessionLocator {
3778            harness: HarnessId::from(HarnessId::OPENCODE),
3779            session_id: "ses_123".into(),
3780            storage: StorageLocator::Sqlite {
3781                path: PathBuf::from("/tmp/opencode-dev.db"),
3782                selector: "ses_123".into(),
3783            },
3784        };
3785        let encoded = serde_json::to_string(&locator).unwrap();
3786        assert_eq!(
3787            serde_json::from_str::<SessionLocator>(&encoded).unwrap(),
3788            locator
3789        );
3790    }
3791
3792    #[test]
3793    fn discovers_filters_loads_and_follows_three_jsonl_harnesses() {
3794        let root = temp_dir("jsonl");
3795        let workspace = root.join("workspace");
3796        let other = root.join("other");
3797        fs::create_dir_all(&workspace).unwrap();
3798        fs::create_dir_all(&other).unwrap();
3799
3800        let claude = root.join("claude");
3801        let codex = root.join("codex");
3802        let pi = root.join("pi");
3803        fs::create_dir_all(&claude).unwrap();
3804        fs::create_dir_all(&codex).unwrap();
3805        fs::create_dir_all(&pi).unwrap();
3806        fs::write(
3807            claude.join("claude.jsonl"),
3808            format!(
3809                "{{\"type\":\"user\",\"sessionId\":\"cc-1\",\"cwd\":{},\"timestamp\":\"2026-01-01T00:00:01Z\",\"message\":{{\"role\":\"user\",\"content\":\"hi\"}}}}\n",
3810                serde_json::to_string(&workspace.to_string_lossy()).unwrap()
3811            ),
3812        )
3813        .unwrap();
3814        fs::write(
3815            codex.join("rollout.jsonl"),
3816            format!(
3817                "{{\"timestamp\":\"2026-01-01T00:00:00Z\",\"type\":\"session_meta\",\"payload\":{{\"id\":\"cx-1\",\"cwd\":{}}}}}\n{{\"timestamp\":\"2026-01-01T00:00:02Z\",\"type\":\"response_item\",\"payload\":{{\"type\":\"message\",\"role\":\"user\",\"content\":[{{\"type\":\"input_text\",\"text\":\"inspect codex\"}}]}}}}\n",
3818                serde_json::to_string(&workspace.to_string_lossy()).unwrap()
3819            ),
3820        )
3821        .unwrap();
3822        fs::write(
3823            pi.join("pi.jsonl"),
3824            format!(
3825                "{{\"type\":\"session\",\"version\":3,\"id\":\"pi-1\",\"timestamp\":\"2026-01-01T00:00:00Z\",\"cwd\":{}}}\n{{\"type\":\"message\",\"message\":{{\"role\":\"user\",\"content\":\"inspect pi\"}}}}\n",
3826                serde_json::to_string(&workspace.to_string_lossy()).unwrap()
3827            ),
3828        )
3829        .unwrap();
3830        fs::write(
3831            pi.join("unrelated.jsonl"),
3832            format!(
3833                "{{\"type\":\"session\",\"version\":3,\"id\":\"pi-2\",\"timestamp\":\"2026-01-01T00:00:00Z\",\"cwd\":{}}}\n",
3834                serde_json::to_string(&other.to_string_lossy()).unwrap()
3835            ),
3836        )
3837        .unwrap();
3838        fs::write(claude.join("partial.jsonl"), "{truncated").unwrap();
3839
3840        let query = DiscoveryQuery {
3841            workspace: Some(workspace),
3842            homes: HarnessHomes {
3843                claude_code: claude,
3844                codex,
3845                pi,
3846                opencode: root.join("missing-opencode"),
3847                grok: root.join("missing-grok"),
3848                gemini: root.join("missing-gemini"),
3849                goose: root.join("missing-goose"),
3850                supercode: root.join("missing-supercode"),
3851                openclaw: root.join("missing-openclaw"),
3852                hermes: root.join("missing-hermes"),
3853                orchestrator: root.join("missing-orchestrator"),
3854            },
3855            ..DiscoveryQuery::default()
3856        };
3857        let catalog = HarnessCatalog::new();
3858        let found = catalog.discover(&query).unwrap();
3859        assert_eq!(found.len(), 3);
3860        assert_eq!(
3861            found
3862                .iter()
3863                .map(|item| item.locator.harness.as_str())
3864                .collect::<HashSet<_>>(),
3865            HashSet::from([HarnessId::CLAUDE_CODE, HarnessId::CODEX, HarnessId::PI])
3866        );
3867        for descriptor in found {
3868            assert!(descriptor.preview_candidates.is_empty());
3869            assert_eq!(descriptor.latest_message_candidates.len(), 1);
3870            assert_eq!(descriptor.latest_message_candidates[0].role, "user");
3871            assert!(descriptor.latest_message_candidates[0].cursor.is_some());
3872            if descriptor.locator.harness.as_str() == HarnessId::CLAUDE_CODE {
3873                assert_eq!(
3874                    descriptor.latest_message_candidates[0]
3875                        .metadata
3876                        .get("timestamp")
3877                        .map(String::as_str),
3878                    Some("2026-01-01T00:00:01Z")
3879                );
3880            } else if descriptor.locator.harness.as_str() == HarnessId::CODEX {
3881                assert_eq!(
3882                    descriptor.latest_message_candidates[0]
3883                        .metadata
3884                        .get("timestamp")
3885                        .map(String::as_str),
3886                    Some("2026-01-01T00:00:02Z")
3887                );
3888            }
3889            let loaded = catalog.load(&descriptor.locator).unwrap();
3890            assert_eq!(
3891                loaded.meta.session_id.as_deref(),
3892                Some(descriptor.locator.session_id.as_str())
3893            );
3894            let mut follower = catalog.follow(&descriptor.locator).unwrap();
3895            assert!(matches!(
3896                follower.poll().unwrap(),
3897                Some(crate::SessionWatchEvent::SessionSnapshot { .. })
3898            ));
3899        }
3900        fs::remove_dir_all(root).ok();
3901    }
3902
3903    #[test]
3904    fn codex_event_messages_supply_bounded_native_order_previews() {
3905        // The list reader must handle the native event-only narration emitted by
3906        // collab sessions, without loading the transcript or widening its budgets.
3907        let root = temp_dir("codex-event-previews");
3908        let path = root.join("rollout.jsonl");
3909        let user = serde_json::json!({
3910            "timestamp": "2026-01-01T00:00:01Z", "type": "event_msg",
3911            "payload": {"type": "user_message", "message": "Investigate the worker"}
3912        });
3913        let answer = serde_json::json!({
3914            "timestamp": "2026-01-01T00:00:02Z", "type": "event_msg",
3915            "payload": {"type": "agent_message", "message": "Worker findings"}
3916        });
3917        let noise = serde_json::json!({
3918            "type": "event_msg", "payload": {"type": "token_count", "message": "not a message"}
3919        });
3920        fs::write(&path, format!("{user}\n{answer}\n{noise}\n{{partial")).unwrap();
3921        let latest = latest_file_message_candidates(&path, HarnessId::CODEX).unwrap();
3922        assert_eq!(latest.len(), 2);
3923        assert_eq!(latest[0].role, "assistant");
3924        assert_eq!(latest[0].content, "Worker findings");
3925        assert_eq!(latest[1].role, "user");
3926        assert_eq!(latest[1].content, "Investigate the worker");
3927        assert_eq!(
3928            latest[0].metadata.get("timestamp").map(String::as_str),
3929            Some("2026-01-01T00:00:02Z")
3930        );
3931        assert_eq!(
3932            latest[0].cursor.as_deref(),
3933            Some(message_candidate_cursor(HarnessId::CODEX, &answer).as_str())
3934        );
3935        let topics = topic_file_message_candidates(&path, HarnessId::CODEX).unwrap();
3936        assert_eq!(topics.len(), 2);
3937        assert_eq!(topics[0].content, "Investigate the worker");
3938
3939        let mut context_pairs = String::new();
3940        for index in 0..5 {
3941            let content = if index < 4 {
3942                format!("# AGENTS.md instructions for /work/{index}\n\n<INSTRUCTIONS>Context</INSTRUCTIONS>")
3943            } else {
3944                "The actual user request".to_string()
3945            };
3946            let event = serde_json::json!({
3947                "type": "event_msg", "payload": {"type": "user_message", "message": content}
3948            });
3949            let response = serde_json::json!({
3950                "type": "response_item", "payload": {"type": "message", "role": "user", "content": content}
3951            });
3952            context_pairs.push_str(&format!("{event}\n{response}\n"));
3953        }
3954        fs::write(&path, context_pairs).unwrap();
3955        let topics = topic_file_message_candidates(&path, HarnessId::CODEX).unwrap();
3956        assert!(topics
3957            .iter()
3958            .any(|candidate| candidate.content == "The actual user request"));
3959
3960        // Mirrors must not halve the previously visible response-item window.
3961        // Either physical order keeps the response's original native cursor.
3962        let mut paired = String::new();
3963        for index in 0..12 {
3964            let content = format!("answer {index}");
3965            let event = serde_json::json!({
3966                "type": "event_msg", "payload": {"type": "agent_message", "message": content}
3967            });
3968            let response = serde_json::json!({
3969                "id": format!("response-{index}"), "type": "response_item",
3970                "payload": {"type": "message", "role": "assistant", "content": content}
3971            });
3972            if index % 2 == 0 {
3973                paired.push_str(&format!("{event}\n{response}\n"));
3974            } else {
3975                paired.push_str(&format!("{response}\n{event}\n"));
3976            }
3977        }
3978        fs::write(&path, paired).unwrap();
3979        let latest = latest_file_message_candidates(&path, HarnessId::CODEX).unwrap();
3980        assert_eq!(latest.len(), LATEST_PREVIEW_CANDIDATES);
3981        assert_eq!(latest[0].content, "answer 11");
3982        assert_eq!(latest[1].content, "answer 10");
3983        assert_eq!(latest[7].content, "answer 4");
3984        for (offset, candidate) in latest.iter().enumerate() {
3985            let response = serde_json::json!({"id": format!("response-{}", 11 - offset)});
3986            assert_eq!(
3987                candidate.cursor.as_deref(),
3988                Some(message_candidate_cursor(HarnessId::CODEX, &response).as_str())
3989            );
3990        }
3991        let topics = topic_file_message_candidates(&path, HarnessId::CODEX).unwrap();
3992        assert_eq!(topics.len(), LATEST_PREVIEW_CANDIDATES);
3993        assert_eq!(topics[7].content, "answer 7");
3994
3995        let event = serde_json::json!({
3996            "type": "event_msg", "payload": {"type": "agent_message", "message": "again"}
3997        });
3998        let response = serde_json::json!({
3999            "type": "response_item", "payload": {"type": "message", "role": "assistant", "content": "again"}
4000        });
4001        for records in [
4002            format!("{event}\n{event}\n"),
4003            format!("{response}\n{response}\n"),
4004            format!("{event}\n{response}\n{event}\n{response}\n"),
4005            format!("{response}\n{event}\n{response}\n{event}\n"),
4006        ] {
4007            fs::write(&path, records).unwrap();
4008            assert_eq!(
4009                latest_file_message_candidates(&path, HarnessId::CODEX)
4010                    .unwrap()
4011                    .len(),
4012                2
4013            );
4014            assert_eq!(
4015                topic_file_message_candidates(&path, HarnessId::CODEX)
4016                    .unwrap()
4017                    .len(),
4018                2
4019            );
4020        }
4021
4022        // Equal truncated prefixes alone are not evidence of a mirrored turn.
4023        let prefix = "x".repeat(4096);
4024        let distinct_event = serde_json::json!({
4025            "type": "event_msg", "payload": {"type": "agent_message", "message": format!("{prefix}A")}
4026        });
4027        let distinct_response = serde_json::json!({
4028            "type": "response_item", "payload": {"type": "message", "role": "assistant", "content": format!("{prefix}B")}
4029        });
4030        fs::write(&path, format!("{distinct_event}\n{distinct_response}\n")).unwrap();
4031        assert_eq!(
4032            latest_file_message_candidates(&path, HarnessId::CODEX)
4033                .unwrap()
4034                .len(),
4035            2
4036        );
4037        assert_eq!(
4038            topic_file_message_candidates(&path, HarnessId::CODEX)
4039                .unwrap()
4040                .len(),
4041            2
4042        );
4043        let long = serde_json::json!({
4044            "type": "event_msg", "payload": {"type": "agent_message", "message": "x".repeat(5000)}
4045        });
4046        fs::write(&path, format!("{long}\n")).unwrap();
4047        let latest = latest_file_message_candidates(&path, HarnessId::CODEX).unwrap();
4048        assert_eq!(latest[0].content.len(), 4096);
4049        fs::remove_dir_all(root).ok();
4050    }
4051
4052    #[test]
4053    fn preview_cursor_tracks_native_boundary_not_growing_text() {
4054        let first = serde_json::json!({
4055            "timestamp": "2026-01-01T00:00:02Z",
4056            "type": "response_item",
4057            "payload": {"type": "message", "role": "assistant", "content": "partial"}
4058        });
4059        let grown = serde_json::json!({
4060            "timestamp": "2026-01-01T00:00:02Z",
4061            "type": "response_item",
4062            "payload": {"type": "message", "role": "assistant", "content": "partial and complete"}
4063        });
4064        let next = serde_json::json!({
4065            "timestamp": "2026-01-01T00:00:03Z",
4066            "type": "response_item",
4067            "payload": {"type": "message", "role": "assistant", "content": "next"}
4068        });
4069
4070        assert_eq!(
4071            message_candidate_cursor(HarnessId::CODEX, &first),
4072            message_candidate_cursor(HarnessId::CODEX, &grown)
4073        );
4074        assert_ne!(
4075            message_candidate_cursor(HarnessId::CODEX, &first),
4076            message_candidate_cursor(HarnessId::CODEX, &next)
4077        );
4078    }
4079
4080    #[test]
4081    fn codex_child_rollouts_roll_into_roots_before_pagination() {
4082        let root = temp_dir("codex-roots");
4083        let codex = root.join("codex");
4084        fs::create_dir_all(&codex).unwrap();
4085        // Recency is the newest record's `timestamp` (tail_facts), not the file's
4086        // mtime, so each rollout says when it last took a turn.
4087        let write_rollout = |name: &str, payload: Value, turn_seconds: u64| {
4088            fs::write(
4089                codex.join(format!("{name}.jsonl")),
4090                format!(
4091                    "{}\n",
4092                    serde_json::json!({
4093                        "timestamp": format!("1970-01-01T00:{:02}:{:02}Z", turn_seconds / 60, turn_seconds % 60),
4094                        "type": "session_meta",
4095                        "payload": payload,
4096                    })
4097                ),
4098            )
4099            .unwrap();
4100        };
4101        write_rollout(
4102            "parent",
4103            serde_json::json!({"id":"parent","cwd":"/project","source":"cli"}),
4104            100,
4105        );
4106        write_rollout(
4107            "other",
4108            serde_json::json!({"id":"other","cwd":"/project","source":"cli"}),
4109            200,
4110        );
4111        write_rollout(
4112            "child",
4113            serde_json::json!({
4114                "id": "child",
4115                "cwd": "/project",
4116                "parent_thread_id": "parent",
4117                "source": {"subagent":{"thread_spawn":{
4118                    "parent_thread_id":"parent",
4119                    "depth":1,
4120                    "agent_path":"/root/reviewer"
4121                }}}
4122            }),
4123            300,
4124        );
4125
4126        let catalog = HarnessCatalog::new();
4127        let query = DiscoveryQuery {
4128            harnesses: vec![HarnessId::from(HarnessId::CODEX)],
4129            homes: HarnessHomes {
4130                codex: codex.clone(),
4131                ..HarnessHomes::default()
4132            },
4133            limit: Some(1),
4134            ..DiscoveryQuery::default()
4135        };
4136        let roots = catalog.discover(&query).unwrap();
4137        assert_eq!(roots.len(), 1);
4138        assert_eq!(roots[0].locator.session_id, "parent");
4139        assert_eq!(roots[0].updated_at_ms, Some(300_000));
4140        assert_eq!(roots[0].parent_session_id, None);
4141        assert_eq!(roots[0].child_session_count, 1);
4142
4143        let tree = catalog
4144            .discover(&DiscoveryQuery {
4145                limit: None,
4146                include_child_sessions: true,
4147                root_session_id: Some("parent".into()),
4148                ..query
4149            })
4150            .unwrap();
4151        assert_eq!(tree.len(), 2);
4152        assert!(tree
4153            .iter()
4154            .all(|descriptor| descriptor.locator.session_id != "other"));
4155        let child = tree
4156            .iter()
4157            .find(|descriptor| descriptor.locator.session_id == "child")
4158            .unwrap();
4159        assert_eq!(child.parent_session_id.as_deref(), Some("parent"));
4160        fs::remove_dir_all(root).ok();
4161    }
4162
4163    #[test]
4164    fn discovers_loads_and_follows_opencode_sqlite() {
4165        let db = PathBuf::from(env!("CARGO_MANIFEST_DIR"))
4166            .join("../harness/tests/fixtures/opencode_fixture/opencode.db");
4167        let catalog = HarnessCatalog::new();
4168        let found = catalog
4169            .discover(&DiscoveryQuery {
4170                harnesses: vec![HarnessId::from(HarnessId::OPENCODE)],
4171                homes: HarnessHomes {
4172                    opencode: db,
4173                    ..HarnessHomes::default()
4174                },
4175                ..DiscoveryQuery::default()
4176            })
4177            .unwrap();
4178        assert!(!found.is_empty());
4179        for descriptor in found {
4180            assert_eq!(descriptor.locator.harness.as_str(), HarnessId::OPENCODE);
4181            assert_eq!(
4182                catalog.load(&descriptor.locator).unwrap().meta.session_id,
4183                Some(descriptor.locator.session_id.clone())
4184            );
4185            assert!(catalog.follow(&descriptor.locator).is_ok());
4186        }
4187    }
4188
4189    #[test]
4190    fn discovers_loads_and_follows_hermes_sqlite_by_session() {
4191        // A copy of the committed Hermes fixture store, so the test may append to it.
4192        let root = temp_dir("hermes-follow");
4193        fs::create_dir_all(&root).unwrap();
4194        let db = root.join("state.db");
4195        fs::copy(
4196            PathBuf::from(env!("CARGO_MANIFEST_DIR"))
4197                .join("../harness/tests/fixtures/hermes_home/state.db"),
4198            &db,
4199        )
4200        .unwrap();
4201        let catalog = HarnessCatalog::new();
4202        let query = DiscoveryQuery {
4203            harnesses: vec![HarnessId::from(HarnessId::HERMES)],
4204            homes: HarnessHomes {
4205                hermes: db.clone(),
4206                ..HarnessHomes::default()
4207            },
4208            ..DiscoveryQuery::default()
4209        };
4210        let found = catalog.discover(&query).unwrap();
4211        assert!(found.len() >= 2, "{found:#?}");
4212        // Every discovered locator loads AND follows as ITS OWN session (not the store's newest),
4213        // and its list preview is that session's own latest message rather than an empty read of
4214        // the store file's tail.
4215        for descriptor in &found {
4216            assert_eq!(descriptor.locator.harness.as_str(), HarnessId::HERMES);
4217            let loaded = catalog.load(&descriptor.locator).unwrap();
4218            let last_text = loaded.messages.iter().rev().find_map(|message| {
4219                (matches!(message.role, crate::Role::User | crate::Role::Assistant))
4220                    .then(|| message.content.clone())
4221                    .flatten()
4222            });
4223            // Lineage-only fixture rows have no messages and therefore no preview.
4224            assert_eq!(
4225                descriptor
4226                    .latest_message_candidates
4227                    .first()
4228                    .map(|c| c.content.as_str()),
4229                last_text.as_deref(),
4230                "{}",
4231                descriptor.locator.session_id
4232            );
4233            assert_eq!(
4234                catalog.load(&descriptor.locator).unwrap().meta.session_id,
4235                Some(descriptor.locator.session_id.clone())
4236            );
4237            let mut follower = catalog.follow(&descriptor.locator).unwrap();
4238            match follower.poll().unwrap() {
4239                Some(crate::watch::SessionWatchEvent::SessionSnapshot { session, .. }) => {
4240                    assert_eq!(
4241                        session.meta.session_id,
4242                        Some(descriptor.locator.session_id.clone())
4243                    );
4244                }
4245                other => panic!("expected an initial snapshot, got {other:?}"),
4246            }
4247        }
4248        // Append a message to ONE session: only that session's follower wakes, with exactly the new
4249        // message, while a sibling's follower stays quiet.
4250        let target = &found[0].locator;
4251        let sibling = &found[1].locator;
4252        let mut target_follower = catalog.follow(target).unwrap();
4253        let mut sibling_follower = catalog.follow(sibling).unwrap();
4254        target_follower.poll().unwrap();
4255        sibling_follower.poll().unwrap();
4256        std::thread::sleep(std::time::Duration::from_millis(20));
4257        {
4258            let conn = rusqlite::Connection::open(&db).unwrap();
4259            conn.execute(
4260                "INSERT INTO messages (session_id, role, content, timestamp, active) VALUES (?1, 'assistant', 'appended by the follow test', ?2, 1)",
4261                rusqlite::params![target.session_id, 1_800_000_000.0_f64],
4262            )
4263            .unwrap();
4264        }
4265        match target_follower.poll().unwrap() {
4266            Some(crate::watch::SessionWatchEvent::MessagesAppended {
4267                session_id,
4268                messages,
4269                ..
4270            }) => {
4271                assert_eq!(session_id, Some(target.session_id.clone()));
4272                assert_eq!(messages.len(), 1);
4273                assert_eq!(
4274                    messages[0].content.as_deref(),
4275                    Some("appended by the follow test")
4276                );
4277            }
4278            other => panic!("expected messages_appended for the target session, got {other:?}"),
4279        }
4280        assert!(
4281            sibling_follower.poll().unwrap().is_none(),
4282            "the sibling session must not wake"
4283        );
4284        fs::remove_dir_all(&root).ok();
4285    }
4286
4287    #[test]
4288    fn discovers_loads_and_follows_gemini_conversation_records() {
4289        let root = temp_dir("gemini");
4290        let workspace = root.join("workspace");
4291        let chats = root.join("gemini/tmp/demo/chats");
4292        fs::create_dir_all(&workspace).unwrap();
4293        fs::create_dir_all(&chats).unwrap();
4294        fs::write(
4295            root.join("gemini/projects.json"),
4296            serde_json::json!({
4297                "projects": {workspace.to_string_lossy(): "demo"}
4298            })
4299            .to_string(),
4300        )
4301        .unwrap();
4302        let transcript = chats.join("gemini-id.jsonl");
4303        fs::write(
4304            &transcript,
4305            include_str!("../../harness/tests/fixtures/gemini_session.jsonl"),
4306        )
4307        .unwrap();
4308
4309        let catalog = HarnessCatalog::new();
4310        let found = catalog
4311            .discover(&DiscoveryQuery {
4312                harnesses: vec![HarnessId::from(HarnessId::GEMINI)],
4313                homes: HarnessHomes {
4314                    gemini: root.join("gemini"),
4315                    ..HarnessHomes::default()
4316                },
4317                workspace: Some(workspace.clone()),
4318                ..DiscoveryQuery::default()
4319            })
4320            .unwrap();
4321
4322        assert_eq!(found.len(), 1);
4323        assert_eq!(found[0].cwd.as_deref(), Some(workspace.as_path()));
4324        assert_eq!(found[0].message_count, None);
4325        assert_eq!(found[0].model.as_deref(), Some("gemini-2.5-pro"));
4326        assert_eq!(found[0].title, None);
4327        assert!(found[0].preview_candidates.is_empty());
4328        assert_eq!(found[0].latest_message_candidates.len(), 3);
4329        assert_eq!(
4330            found[0].latest_message_candidates[0].content,
4331            "Fixture inspected."
4332        );
4333        let loaded = catalog.load(&found[0].locator).unwrap();
4334        assert_eq!(
4335            loaded.meta.session_id.as_deref(),
4336            Some("11111111-1111-4111-8111-111111111111")
4337        );
4338        assert_eq!(loaded.messages.len(), 4);
4339        assert!(matches!(
4340            catalog.follow(&found[0].locator).unwrap().poll().unwrap(),
4341            Some(crate::SessionWatchEvent::SessionSnapshot { .. })
4342        ));
4343        fs::remove_dir_all(root).ok();
4344    }
4345
4346    #[test]
4347    fn preview_search_filters_before_pagination_without_changing_metadata_search() {
4348        // Search/pagination must compose: filtering only the returned page loses
4349        // matches and gives a false total. Exercise the public catalog door.
4350        let root = temp_dir("preview-search");
4351        for (id, first, last) in [
4352            ("topic-hit", "NEBULA opening", "Finished"),
4353            ("latest-hit", "Ordinary opening", "Found the nebula"),
4354            ("no-hit", "Unrelated opening", "Finished"),
4355        ] {
4356            fs::write(root.join(format!("{id}.jsonl")), format!("{}\n{}\n",
4357                serde_json::json!({"sessionId": id, "cwd": "/work", "type": "user", "message": {"role": "user", "content": first}}),
4358                serde_json::json!({"sessionId": id, "type": "assistant", "message": {"role": "assistant", "content": last}}),
4359            )).unwrap();
4360        }
4361        let query: DiscoveryQuery = serde_json::from_value(serde_json::json!({
4362            "harnesses": ["claude-code"], "homes": {"claude_code": root},
4363            "query": "  nebula  ", "search_previews": true, "limit": 1
4364        }))
4365        .unwrap();
4366        let catalog = HarnessCatalog::new();
4367        let first = catalog.discover_page(&query).unwrap();
4368        assert!(first.receipt.searched_previews);
4369        assert_eq!(first.receipt.total_matched, 2);
4370        assert_eq!(first.sessions.len(), 1);
4371        assert!(first.receipt.truncated);
4372        let second = catalog
4373            .discover_page(&DiscoveryQuery {
4374                cursor: first.next_cursor.clone(),
4375                ..query.clone()
4376            })
4377            .unwrap();
4378        assert_eq!(second.receipt.total_matched, 2);
4379        assert_eq!(second.sessions.len(), 1);
4380        assert_ne!(first.sessions[0].locator, second.sessions[0].locator);
4381        assert!(!second.receipt.truncated);
4382        let mut metadata = serde_json::to_value(&query).unwrap();
4383        metadata["search_previews"] = false.into();
4384        let metadata_page = catalog
4385            .discover_page(&serde_json::from_value(metadata).unwrap())
4386            .unwrap();
4387        assert!(metadata_page.sessions.is_empty());
4388        assert!(serde_json::to_value(&metadata_page.receipt)
4389            .unwrap()
4390            .get("searched_previews")
4391            .is_none());
4392
4393        // Metadata remains part of the union; matching multiple candidates must
4394        // still yield one session, not one row per message.
4395        let all = catalog
4396            .discover_page(&DiscoveryQuery {
4397                query: Some("hit".into()),
4398                limit: None,
4399                ..query.clone()
4400            })
4401            .unwrap();
4402        assert_eq!(all.sessions.len(), 3);
4403        assert_eq!(all.receipt.total_matched, 3);
4404        let elsewhere = catalog
4405            .discover_page(&DiscoveryQuery {
4406                workspace: Some("/elsewhere".into()),
4407                ..query.clone()
4408            })
4409            .unwrap();
4410        assert_eq!(elsewhere.receipt.total_matched, 0);
4411        let excluded_by_time = catalog
4412            .discover_page(&DiscoveryQuery {
4413                updated_after_ms: Some(u64::MAX),
4414                ..query.clone()
4415            })
4416            .unwrap();
4417        assert_eq!(excluded_by_time.receipt.total_matched, 0);
4418        for invalid in [
4419            DiscoveryQuery {
4420                query: None,
4421                ..query.clone()
4422            },
4423            DiscoveryQuery {
4424                query: Some("  ".into()),
4425                ..query.clone()
4426            },
4427            DiscoveryQuery {
4428                limit: Some(0),
4429                ..query.clone()
4430            },
4431            DiscoveryQuery {
4432                cursor: Some("bad-cursor".into()),
4433                ..query.clone()
4434            },
4435            DiscoveryQuery {
4436                cursor: first.next_cursor,
4437                query: Some("absent".into()),
4438                ..query.clone()
4439            },
4440        ] {
4441            assert!(catalog.discover_page(&invalid).is_err());
4442        }
4443        assert!(catalog.project_index_page(&query, Vec::new()).is_err());
4444        fs::remove_dir_all(root).ok();
4445    }
4446
4447    #[test]
4448    fn preview_search_uses_codex_first_history_topic_and_bounded_candidates() {
4449        let root = temp_dir("preview-search-codex");
4450        let sessions = root.join("sessions");
4451        fs::create_dir_all(&sessions).unwrap();
4452        fs::write(
4453            root.join("history.jsonl"),
4454            format!(
4455                "{}\n{}\n",
4456                serde_json::json!({"session_id": "history-hit", "text": "Original nebula topic"}),
4457                serde_json::json!({"session_id": "history-hit", "text": "laterhistoryonly"}),
4458            ),
4459        )
4460        .unwrap();
4461        for id in ["history-hit", "latest-hit", "bounded"] {
4462            let mut content = format!(
4463                "{}\n",
4464                serde_json::json!({
4465                    "type": "session_meta", "payload": {"id": id, "cwd": "/work"}
4466                })
4467            );
4468            for index in 0..20 {
4469                let message = if id == "latest-hit" && index == 19 {
4470                    "Found NEBULA".to_string()
4471                } else if index == 10 {
4472                    "middlehistoryonly".to_string()
4473                } else {
4474                    format!("{}beyondtextcap", "x".repeat(4096))
4475                };
4476                content.push_str(&format!("{}\n", serde_json::json!({
4477                    "type": "event_msg", "payload": {"type": "agent_message", "message": message}
4478                })));
4479            }
4480            fs::write(sessions.join(format!("{id}.jsonl")), content).unwrap();
4481        }
4482        let catalog = HarnessCatalog::new();
4483        let query: DiscoveryQuery = serde_json::from_value(serde_json::json!({
4484            "harnesses": ["codex"], "homes": {"codex": sessions},
4485            "query": "nebula", "search_previews": true
4486        }))
4487        .unwrap();
4488        let page = catalog.discover_page(&query).unwrap();
4489        assert_eq!(page.receipt.total_matched, 2);
4490        for row in &page.sessions {
4491            assert!(row.preview_candidates.len() <= 8);
4492            assert!(row.latest_message_candidates.len() <= 8);
4493            assert!(row
4494                .preview_candidates
4495                .iter()
4496                .chain(&row.latest_message_candidates)
4497                .all(|candidate| candidate.content.chars().count() <= 4096));
4498        }
4499        for text in ["middlehistoryonly", "laterhistoryonly", "beyondtextcap"] {
4500            assert!(
4501                catalog
4502                    .discover_page(&DiscoveryQuery {
4503                        query: Some(text.into()),
4504                        ..query.clone()
4505                    })
4506                    .unwrap()
4507                    .sessions
4508                    .is_empty(),
4509                "not a full-history search: {text}"
4510            );
4511        }
4512        fs::remove_dir_all(root).unwrap();
4513    }
4514
4515    #[test]
4516    fn discovers_native_store_and_pages_search_results() {
4517        let root = temp_dir("supercode");
4518        let store_root = root.join("sessions");
4519        fs::create_dir_all(&store_root).unwrap();
4520        for (name, title) in [
4521            ("alpha", "Alpha planning"),
4522            ("beta", "Beta implementation"),
4523            ("gamma", "Gamma review"),
4524        ] {
4525            fs::write(
4526                store_root.join(format!("{name}.jsonl")),
4527                format!("{{\"role\":\"user\",\"content\":\"{title}\"}}\n"),
4528            )
4529            .unwrap();
4530            fs::write(
4531                store_root.join(format!("{name}.meta.json")),
4532                serde_json::json!({"name": name, "title": title}).to_string(),
4533            )
4534            .unwrap();
4535        }
4536        let catalog = HarnessCatalog::new();
4537        let base = DiscoveryQuery {
4538            harnesses: vec![HarnessId::from(HarnessId::SUPERCODE)],
4539            homes: HarnessHomes {
4540                supercode: store_root,
4541                ..HarnessHomes::default()
4542            },
4543            limit: Some(1),
4544            ..DiscoveryQuery::default()
4545        };
4546
4547        let first = catalog.discover_page(&base).unwrap();
4548        assert_eq!(first.sessions.len(), 1);
4549        assert!(first.next_cursor.is_some());
4550        let second = catalog
4551            .discover_page(&DiscoveryQuery {
4552                cursor: first.next_cursor,
4553                ..base.clone()
4554            })
4555            .unwrap();
4556        assert_eq!(second.sessions.len(), 1);
4557        assert_ne!(
4558            first.sessions[0].locator.session_id,
4559            second.sessions[0].locator.session_id
4560        );
4561        let search = catalog
4562            .discover_page(&DiscoveryQuery {
4563                limit: None,
4564                query: Some("implementation".into()),
4565                ..base
4566            })
4567            .unwrap();
4568        assert_eq!(search.sessions.len(), 1);
4569        assert_eq!(search.sessions[0].locator.session_id, "beta");
4570        assert_eq!(search.sessions[0].message_count, None);
4571        assert_eq!(
4572            catalog
4573                .load(&search.sessions[0].locator)
4574                .unwrap()
4575                .messages
4576                .len(),
4577            1
4578        );
4579        fs::remove_dir_all(root).ok();
4580    }
4581
4582    #[test]
4583    fn native_workspace_discovery_reads_bounded_sidecar_headers() {
4584        let root = temp_dir("supercode-bounded-header");
4585        let store_root = root.join("sessions");
4586        let workspace = root.join("project");
4587        fs::create_dir_all(&store_root).unwrap();
4588        fs::create_dir_all(&workspace).unwrap();
4589        let name = "bounded-native";
4590        fs::write(
4591            store_root.join(format!("{name}.meta.json")),
4592            serde_json::json!({"name": name, "title": "Bounded native"}).to_string(),
4593        )
4594        .unwrap();
4595        fs::write(
4596            store_root.join(format!("{name}.jsonl")),
4597            "{\"role\":\"user\",\"content\":\"projected view\"}\n",
4598        )
4599        .unwrap();
4600        let sidecar = [
4601            serde_json::json!({
4602                "supercode_native": 2,
4603                "source": "claude_code",
4604                "session_id": "native-session"
4605            })
4606            .to_string(),
4607            serde_json::json!({
4608                "type": "user",
4609                "sessionId": "native-session",
4610                "cwd": workspace,
4611                "message": {"role": "user", "content": "hello"}
4612            })
4613            .to_string(),
4614            serde_json::json!({
4615                "type": "assistant",
4616                "sessionId": "native-session",
4617                "cwd": workspace,
4618                "message": {"role": "assistant", "model": "claude-sonnet-5", "content": []}
4619            })
4620            .to_string(),
4621            // A full native-family parse rejects this trailing residue. Header
4622            // discovery must not touch it after it has enough metadata.
4623            "not-json".into(),
4624        ]
4625        .join("\n");
4626        fs::write(
4627            store_root.join(format!("{name}.sidecar.jsonl")),
4628            format!("{sidecar}\n"),
4629        )
4630        .unwrap();
4631
4632        let found = HarnessCatalog::new()
4633            .discover(&DiscoveryQuery {
4634                workspace: Some(workspace.clone()),
4635                harnesses: vec![HarnessId::from(HarnessId::SUPERCODE)],
4636                homes: HarnessHomes {
4637                    supercode: store_root,
4638                    ..HarnessHomes::default()
4639                },
4640                ..DiscoveryQuery::default()
4641            })
4642            .unwrap();
4643
4644        assert_eq!(found.len(), 1);
4645        assert_eq!(found[0].cwd.as_deref(), Some(workspace.as_path()));
4646        assert_eq!(found[0].model.as_deref(), Some("claude-sonnet-5"));
4647        assert_eq!(found[0].message_count, None);
4648        fs::remove_dir_all(root).ok();
4649    }
4650
4651    #[test]
4652    fn discovers_current_opencode_schema_without_a_session_model_column() {
4653        let root = temp_dir("opencode-current");
4654        let db = root.join("opencode.db");
4655        let conn = Connection::open(&db).unwrap();
4656        conn.execute_batch(
4657            "CREATE TABLE session (
4658                id TEXT PRIMARY KEY,
4659                directory TEXT NOT NULL,
4660                title TEXT NOT NULL,
4661                time_updated INTEGER NOT NULL
4662             );
4663             CREATE TABLE message (
4664                id TEXT PRIMARY KEY,
4665                session_id TEXT NOT NULL
4666             );
4667             INSERT INTO session VALUES ('ses_current', '/tmp/work', 'Current', 42);
4668             INSERT INTO message VALUES ('msg_current', 'ses_current');",
4669        )
4670        .unwrap();
4671        drop(conn);
4672
4673        let found = HarnessCatalog::new()
4674            .discover(&DiscoveryQuery {
4675                harnesses: vec![HarnessId::from(HarnessId::OPENCODE)],
4676                homes: HarnessHomes {
4677                    opencode: db,
4678                    ..HarnessHomes::default()
4679                },
4680                ..DiscoveryQuery::default()
4681            })
4682            .unwrap();
4683
4684        assert_eq!(found.len(), 1);
4685        assert_eq!(found[0].locator.session_id, "ses_current");
4686        assert_eq!(found[0].message_count, Some(1));
4687        assert_eq!(found[0].model, None);
4688        fs::remove_dir_all(root).ok();
4689    }
4690
4691    #[test]
4692    fn workspace_filter_never_matches_a_relative_recorded_cwd() {
4693        // OpenCode has shipped session rows whose `directory` is the literal
4694        // ".". Resolving that against the discoverer's own cwd made the
4695        // session match every workspace discovery ran from — the workspace
4696        // here IS the test process cwd, the exact aliasing that leaked.
4697        let root = temp_dir("opencode-relative-cwd");
4698        let db = root.join("opencode.db");
4699        let conn = Connection::open(&db).unwrap();
4700        let here = std::env::current_dir().unwrap();
4701        conn.execute_batch(&format!(
4702            "CREATE TABLE session (
4703                id TEXT PRIMARY KEY,
4704                directory TEXT NOT NULL,
4705                title TEXT NOT NULL,
4706                time_updated INTEGER NOT NULL
4707             );
4708             CREATE TABLE message (
4709                id TEXT PRIMARY KEY,
4710                session_id TEXT NOT NULL
4711             );
4712             INSERT INTO session VALUES ('ses_relative', '.', 'Ghost', 41);
4713             INSERT INTO session VALUES ('ses_here', '{}', 'Real', 42);",
4714            here.display()
4715        ))
4716        .unwrap();
4717        drop(conn);
4718
4719        let found = HarnessCatalog::new()
4720            .discover(&DiscoveryQuery {
4721                workspace: Some(here),
4722                harnesses: vec![HarnessId::from(HarnessId::OPENCODE)],
4723                homes: HarnessHomes {
4724                    opencode: db,
4725                    ..HarnessHomes::default()
4726                },
4727                ..DiscoveryQuery::default()
4728            })
4729            .unwrap();
4730
4731        assert_eq!(found.len(), 1);
4732        assert_eq!(found[0].locator.session_id, "ses_here");
4733        fs::remove_dir_all(root).ok();
4734    }
4735
4736    #[test]
4737    fn subtree_scope_admits_the_root_and_every_folder_under_it_only() {
4738        let base = temp_dir("subtree-scope");
4739        let work = base.join("work");
4740        let nested = work.join("repo").join("pkg");
4741        let sibling = base.join("workshop");
4742        fs::create_dir_all(&nested).unwrap();
4743        fs::create_dir_all(&sibling).unwrap();
4744        let scope = WorkspaceScope::subtree(&work);
4745        assert!(scope.admits(&work));
4746        assert!(scope.admits(&work.join("repo")));
4747        assert!(scope.admits(&nested));
4748        // A folder that no longer exists still names where its session ran.
4749        assert!(scope.admits(&work.join("deleted-since")));
4750        // `..` is resolved before containment is judged.
4751        assert!(!scope.admits(&work.join("..").join("workshop")));
4752        assert!(!scope.admits(&sibling), "/work must not admit /workshop");
4753        assert!(!scope.admits(&base));
4754        assert!(!scope.admits(Path::new(".")));
4755        assert!(!scope.admits(Path::new("repo/pkg")));
4756        let exact = WorkspaceScope::exact(&work);
4757        assert!(exact.admits(&work));
4758        assert!(!exact.admits(&nested), "the default stays the exact folder");
4759        fs::remove_dir_all(base).ok();
4760    }
4761
4762    #[cfg(windows)]
4763    #[test]
4764    fn subtree_scope_on_windows_ignores_verbatim_prefix_and_case_and_rejects_driveless_cwds() {
4765        let base = temp_dir("subtree-windows");
4766        let work = base.join("Work");
4767        fs::create_dir_all(work.join("repo")).unwrap();
4768        let scope = WorkspaceScope::subtree(&work);
4769        let verbatim = fs::canonicalize(work.join("repo")).unwrap();
4770        assert!(verbatim.to_string_lossy().starts_with(r"\\?\"));
4771        assert!(scope.admits(&verbatim));
4772        let folded = PathBuf::from(work.join("repo").to_string_lossy().to_uppercase());
4773        assert!(scope.admits(&folded));
4774        // `\Users\...` has no drive: Rust does not call it absolute, and it
4775        // says nothing about which volume the session ran on.
4776        let text = work.to_string_lossy().into_owned();
4777        let driveless = PathBuf::from(&text[2..]);
4778        assert!(driveless.has_root() && !driveless.is_absolute());
4779        assert!(!scope.admits(&driveless));
4780        assert!(!scope.admits(&driveless.join("repo")));
4781        fs::remove_dir_all(base).ok();
4782    }
4783
4784    #[test]
4785    fn subtree_scope_matches_through_a_symlinked_root() {
4786        let base = temp_dir("subtree-symlink");
4787        let real = base.join("real");
4788        fs::create_dir_all(real.join("repo")).unwrap();
4789        let alias = base.join("alias");
4790        #[cfg(unix)]
4791        let linked = std::os::unix::fs::symlink(&real, &alias).is_ok();
4792        #[cfg(windows)]
4793        let linked = std::os::windows::fs::symlink_dir(&real, &alias).is_ok()
4794            // Windows without Developer Mode refuses symlinks; a junction is
4795            // the unprivileged directory alias there.
4796            || std::process::Command::new("cmd")
4797                .arg("/C")
4798                .arg("mklink")
4799                .arg("/J")
4800                .arg(&alias)
4801                .arg(&real)
4802                .output()
4803                .is_ok_and(|output| output.status.success());
4804        if !linked {
4805            fs::remove_dir_all(base).ok();
4806            return;
4807        }
4808        // Root given through the alias, cwd recorded as the real path, and
4809        // the other way round (macOS records /private/var for /var).
4810        assert!(WorkspaceScope::subtree(&alias).admits(&real.join("repo")));
4811        assert!(WorkspaceScope::subtree(&real).admits(&alias.join("repo")));
4812        assert!(!WorkspaceScope::subtree(&alias).admits(&base));
4813        fs::remove_dir_all(base).ok();
4814    }
4815
4816    #[test]
4817    fn workspace_subtree_discovery_filters_before_enrichment_and_defaults_to_exact() {
4818        let base = temp_dir("subtree-discovery");
4819        let home = base.join("projects");
4820        let work = base.join("work");
4821        let nested = work.join("repo");
4822        let sibling = base.join("workshop");
4823        for dir in [&home, &nested, &sibling] {
4824            fs::create_dir_all(dir).unwrap();
4825        }
4826        let write = |id: &str, cwd: &Path| {
4827            let project = home.join(id);
4828            fs::create_dir_all(&project).unwrap();
4829            let line = serde_json::json!({
4830                "type": "user",
4831                "sessionId": id,
4832                "cwd": cwd,
4833                "message": {"role": "user", "content": format!("hello from {id}")}
4834            });
4835            fs::write(project.join(format!("{id}.jsonl")), format!("{line}\n")).unwrap();
4836        };
4837        write("in-root", &work);
4838        write("in-repo", &nested);
4839        write("sibling", &sibling);
4840        write("relative", Path::new("repo"));
4841        let query = |subtree: bool| DiscoveryQuery {
4842            workspace: Some(work.clone()),
4843            workspace_subtree: subtree,
4844            harnesses: vec![HarnessId::from(HarnessId::CLAUDE_CODE)],
4845            homes: HarnessHomes {
4846                claude_code: home.clone(),
4847                ..HarnessHomes::default()
4848            },
4849            ..DiscoveryQuery::default()
4850        };
4851        let ids = |subtree: bool| {
4852            let mut ids: Vec<String> = HarnessCatalog::new()
4853                .discover(&query(subtree))
4854                .unwrap()
4855                .into_iter()
4856                .map(|descriptor| descriptor.locator.session_id)
4857                .collect();
4858            ids.sort();
4859            ids
4860        };
4861        assert_eq!(ids(false), vec!["in-root"]);
4862        assert_eq!(ids(true), vec!["in-repo", "in-root"]);
4863
4864        // Wire: an old client never sends the field and gets the exact match;
4865        // the default is not serialized, so old servers see the same query.
4866        let old: DiscoveryQuery =
4867            serde_json::from_value(serde_json::json!({"workspace": work})).unwrap();
4868        assert!(!old.workspace_subtree);
4869        assert!(serde_json::to_value(&old)
4870            .unwrap()
4871            .get("workspace_subtree")
4872            .is_none());
4873        let new: DiscoveryQuery = serde_json::from_value(
4874            serde_json::json!({"workspace": work, "workspace_subtree": true}),
4875        )
4876        .unwrap();
4877        assert!(new.workspace_subtree);
4878        fs::remove_dir_all(base).ok();
4879    }
4880}