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