Skip to main content

supercode_harness/
session_index.rs

1//! Revisioned session-list subscriptions for latency-sensitive frontends.
2//!
3//! Native filesystem events are treated as invalidation hints, never as the
4//! session record itself. Each hint causes a bounded re-read of the affected
5//! Claude Code or Codex transcript; a slow periodic catalog reconciliation
6//! repairs dropped/coalesced platform events and fills a page after removals.
7
8use std::collections::{BTreeMap, BTreeSet};
9use std::fs;
10use std::path::{Path, PathBuf};
11use std::sync::atomic::{AtomicBool, Ordering};
12use std::sync::{mpsc, Arc};
13use std::time::{Duration, Instant, UNIX_EPOCH};
14
15use notify::event::{AccessKind, AccessMode};
16use notify::{Event, EventKind, RecommendedWatcher, RecursiveMode, Watcher};
17use serde::Serialize;
18use tokio::sync::Notify;
19
20use supercode_interchange::catalog::{hermes_session_stores, CodexHistoryTopicIndex};
21
22use crate::{
23    DiscoveryPage, DiscoveryQuery, HarnessCatalog, HarnessId, SessionDescriptor, SessionLocator,
24    StorageLocator,
25};
26
27const RECONCILE_INTERVAL: Duration = Duration::from_secs(60);
28const MAX_SUBSCRIPTION_ROWS: usize = 2_048;
29const INVALIDATION_QUEUE_CAPACITY: usize = 1_024;
30
31/// Stable public identity for a session-index change. Persistence paths remain
32/// inside the trusted host and are sent only as part of complete descriptors.
33#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize)]
34pub struct SessionIndexKey {
35    /// Owning harness id.
36    pub harness: String,
37    /// Harness-native durable session id.
38    pub session_id: String,
39}
40
41impl SessionIndexKey {
42    fn from_locator(locator: &SessionLocator) -> Self {
43        Self {
44            harness: locator.harness.as_str().to_string(),
45            session_id: locator.session_id.clone(),
46        }
47    }
48}
49
50/// One complete replacement in a revisioned index delta.
51#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
52#[serde(tag = "kind", rename_all = "snake_case")]
53pub enum SessionIndexChange {
54    /// A session entered the bounded result page.
55    Added {
56        /// Complete current descriptor.
57        descriptor: SessionDescriptor,
58    },
59    /// A visible session's descriptor changed.
60    Updated {
61        /// Complete replacement descriptor.
62        descriptor: SessionDescriptor,
63    },
64    /// A session disappeared from the bounded result page.
65    Removed {
66        /// Stable identity of the removed descriptor.
67        key: SessionIndexKey,
68    },
69}
70
71/// One subscription poll result. Revisions start at one for the initial
72/// snapshot and increase by exactly one for each non-empty delta batch or
73/// committed window resize. They describe the visible window, not changes to
74/// out-of-window inventory totals returned by a same-limit read.
75#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
76pub struct SessionIndexDelta {
77    /// Monotonic subscription-local revision.
78    pub revision: u64,
79    /// Complete replacement changes in deterministic identity order.
80    pub changes: Vec<SessionIndexChange>,
81}
82
83/// Filesystem-backed index subscription. Dropping it drops the platform
84/// watcher and callback channel, so unsubscribe has deterministic cleanup.
85pub(crate) struct SessionIndexSubscription {
86    query: DiscoveryQuery,
87    raw: BTreeMap<SessionIndexKey, SessionDescriptor>,
88    paths: BTreeMap<PathBuf, SessionIndexKey>,
89    current: BTreeMap<SessionIndexKey, SessionDescriptor>,
90    fingerprints: BTreeMap<PathBuf, FileFingerprint>,
91    /// Whole-store SQLite files (Hermes `state.db` + its `-wal`/`-shm`), keyed by path. A store
92    /// holds every session in one file, so a stamp change means "re-enumerate this store", not
93    /// "this one path is one session". `None` = the file is absent.
94    store_fingerprints: BTreeMap<PathBuf, Option<FileFingerprint>>,
95    codex_history: Option<CodexHistoryTopicIndex>,
96    revision: u64,
97    receiver: mpsc::Receiver<notify::Result<Event>>,
98    overflowed: Arc<AtomicBool>,
99    _watcher: RecommendedWatcher,
100    last_reconcile: Instant,
101}
102
103#[derive(Debug, Clone, Copy, PartialEq, Eq)]
104struct FileFingerprint {
105    len: u64,
106    modified_ns: u128,
107    modified_ms: Option<u64>,
108    identity: u128,
109}
110
111/// A complete replacement prepared without changing the subscription. The
112/// transport can serialize it before committing, so a failed response leaves
113/// the old window and revision usable.
114pub(crate) struct PreparedIndexResize {
115    limit: usize,
116    pub(crate) revision: u64,
117    pub(crate) page: DiscoveryPage,
118}
119
120impl SessionIndexSubscription {
121    pub(crate) fn homes(&self) -> &crate::HarnessHomes {
122        &self.query.homes
123    }
124
125    pub(crate) fn prepare_resize(&self, limit: usize) -> Result<PreparedIndexResize, String> {
126        let mut query = self.query.clone();
127        query.limit = Some(limit);
128        validate_query(&query)?;
129        // This is a snapshot of the known index, not a filesystem barrier.
130        // Pending invalidations stay queued and produce the next normal delta.
131        let page = self.project_current(&query, &BTreeSet::new())?;
132        Ok(PreparedIndexResize {
133            limit,
134            revision: if self.query.limit == Some(limit) {
135                self.revision
136            } else {
137                self.revision.saturating_add(1)
138            },
139            page,
140        })
141    }
142
143    pub(crate) fn commit_resize(&mut self, prepared: PreparedIndexResize) {
144        self.query.limit = Some(prepared.limit);
145        self.revision = prepared.revision;
146        self.current = descriptor_map(prepared.page.sessions);
147    }
148
149    pub(crate) fn open(
150        mut query: DiscoveryQuery,
151        notifier: Arc<Notify>,
152    ) -> Result<(Self, Vec<SessionDescriptor>), String> {
153        validate_query(&query)?;
154        query.cursor = None;
155        query.limit = Some(query.limit.unwrap_or(100));
156
157        let catalog = HarnessCatalog::new();
158        let raw = descriptor_map(catalog.discover_raw_index(&query));
159        let projected = catalog
160            .project_index(&query, raw.values().cloned())
161            .map_err(|error| error.to_string())?;
162        let mut codex_history = (query.include_topic_candidates
163            && query
164                .harnesses
165                .iter()
166                .any(|harness| harness.as_str() == HarnessId::CODEX))
167        .then(|| CodexHistoryTopicIndex::new(&query.homes.codex));
168        if let Some(history) = &mut codex_history {
169            // Discovery has always treated an unavailable history file as a
170            // soft fallback to transcript topics. Preserve that behavior.
171            let _ = history.refresh();
172        }
173        let initial = match &codex_history {
174            Some(history) => {
175                catalog.enrich_index_page_with_codex_history(&query, projected, history)
176            }
177            None => catalog.enrich_index_page(&query, projected),
178        }
179        .map_err(|error| error.to_string())?;
180        let paths = descriptor_path_map(&raw);
181        let current = descriptor_map(initial.iter().cloned());
182        let fingerprints = scan_file_fingerprints(&query);
183        let store_fingerprints = scan_store_fingerprints(&query);
184        let (sender, receiver) = mpsc::sync_channel(INVALIDATION_QUEUE_CAPACITY);
185        let overflowed = Arc::new(AtomicBool::new(false));
186        let callback_overflowed = Arc::clone(&overflowed);
187        let callback_notifier = Arc::clone(&notifier);
188        let mut watcher = notify::recommended_watcher(move |event| {
189            if sender.try_send(event).is_err() {
190                callback_overflowed.store(true, Ordering::Release);
191            }
192            callback_notifier.notify_one();
193        })
194        .map_err(|error| error.to_string())?;
195        for root in watch_roots(&query) {
196            if let Some(watched) = existing_watch_root(&root) {
197                watcher
198                    .watch(&watched, RecursiveMode::Recursive)
199                    .map_err(|error| format!("cannot watch {}: {error}", watched.display()))?;
200            }
201        }
202        for store in store_paths(&query) {
203            // The store's directory, not the store file: a WAL-mode writer creates and removes the
204            // `-wal`/`-shm` siblings, and a first run creates the store itself.
205            let Some(dir) = store.parent() else { continue };
206            if let Some(watched) = existing_watch_root(dir) {
207                watcher
208                    .watch(&watched, RecursiveMode::NonRecursive)
209                    .map_err(|error| format!("cannot watch {}: {error}", watched.display()))?;
210            }
211        }
212        if let Some(history) = &codex_history {
213            let target = if history.path().is_file() {
214                history.path()
215            } else {
216                history.path().parent().unwrap_or(history.path())
217            };
218            if target.exists() {
219                watcher
220                    .watch(target, RecursiveMode::NonRecursive)
221                    .map_err(|error| format!("cannot watch {}: {error}", target.display()))?;
222            }
223        }
224
225        Ok((
226            Self {
227                query,
228                raw,
229                paths,
230                current,
231                fingerprints,
232                store_fingerprints,
233                codex_history,
234                revision: 1,
235                receiver,
236                overflowed,
237                _watcher: watcher,
238                last_reconcile: Instant::now(),
239            },
240            initial,
241        ))
242    }
243
244    /// Drain and coalesce native invalidations once. No events means no I/O
245    /// until the minute-scale metadata-only recovery sweep becomes due.
246    pub(crate) fn poll(&mut self) -> Result<Option<SessionIndexDelta>, String> {
247        let mut paths = BTreeSet::new();
248        let mut sweep = self.overflowed.swap(false, Ordering::AcqRel);
249        let mut stores = false;
250        while let Ok(event) = self.receiver.try_recv() {
251            match event {
252                // an open, a read or a close without writing changes nothing: on Linux, notify reports every
253                // open (IN_OPEN) in a watched directory, so a Hermes gateway opening its store per request drove
254                // a full re-index per open (a core at 100% in Open Autonomy's reporter); a close after writing stays
255                Ok(event) if matches!(event.kind, EventKind::Access(access) if access != AccessKind::Close(AccessMode::Write)) =>
256                    {}
257                Ok(event) => {
258                    if event.paths.is_empty() {
259                        sweep = true;
260                    }
261                    for path in event.paths {
262                        if is_store_shm(&self.store_fingerprints, &path) {
263                            continue;
264                        }
265                        if path.extension().and_then(|value| value.to_str()) == Some("jsonl") {
266                            paths.insert(path);
267                        } else if self
268                            .store_fingerprints
269                            .contains_key(&normalized_store_path(&path))
270                        {
271                            stores = true;
272                        } else {
273                            sweep = true;
274                        }
275                    }
276                }
277                Err(_) => sweep = true,
278            }
279        }
280        if self.last_reconcile.elapsed() >= RECONCILE_INTERVAL {
281            sweep = true;
282        }
283        if paths.is_empty() && !sweep && !stores {
284            return Ok(None);
285        }
286
287        let mut content_dirty = BTreeSet::new();
288        let history_path = self
289            .codex_history
290            .as_ref()
291            .map(|history| normalized_path(history.path()));
292        if let Some(history) = &mut self.codex_history {
293            if let Ok(changed) = history.refresh() {
294                content_dirty.extend(changed.into_iter().map(|session_id| SessionIndexKey {
295                    harness: HarnessId::CODEX.to_string(),
296                    session_id,
297                }));
298            }
299        }
300        if sweep {
301            self.reconcile_filesystem(&mut content_dirty)?;
302        }
303        if sweep || stores {
304            self.reconcile_stores(&mut content_dirty)?;
305        }
306        for path in paths {
307            if history_path
308                .as_ref()
309                .is_some_and(|history_path| normalized_path(&path) == *history_path)
310            {
311                continue;
312            }
313            self.refresh_path(&path, &mut content_dirty)?;
314        }
315        // Nothing that feeds the index moved: whatever woke this poll (a reader's -wal close or attribute change, an
316        // unrelated file beside a store) changed no session, and the rebuild would reproduce the index it has. On a
317        // Hermes store its gateway's every read woke it, and rebuilding the whole index each time held a core.
318        if content_dirty.is_empty() {
319            return Ok(None);
320        }
321        let before = self.current.clone();
322        self.rebuild_current(&content_dirty)?;
323        let changes = diff_descriptors(&before, &self.current);
324        if changes.is_empty() {
325            return Ok(None);
326        }
327        self.revision = self.revision.saturating_add(1);
328        Ok(Some(SessionIndexDelta {
329            revision: self.revision,
330            changes,
331        }))
332    }
333
334    fn reconcile_filesystem(
335        &mut self,
336        content_dirty: &mut BTreeSet<SessionIndexKey>,
337    ) -> Result<(), String> {
338        self.last_reconcile = Instant::now();
339        let next = scan_file_fingerprints(&self.query);
340        let changed = self
341            .fingerprints
342            .keys()
343            .chain(next.keys())
344            .filter(|path| self.fingerprints.get(*path) != next.get(*path))
345            .cloned()
346            .collect::<BTreeSet<_>>();
347        for path in changed {
348            self.refresh_path(&path, content_dirty)?;
349        }
350        self.fingerprints = next;
351        Ok(())
352    }
353
354    /// Re-enumerate every whole-store harness whose store stamps moved. Rows are diffed by value:
355    /// a store keeps its sessions' `message_count`/`ended_at` current on every append, so a
356    /// descriptor that compares equal is unchanged and one that differs is content-dirty.
357    fn reconcile_stores(
358        &mut self,
359        content_dirty: &mut BTreeSet<SessionIndexKey>,
360    ) -> Result<(), String> {
361        let next = scan_store_fingerprints(&self.query);
362        if next == self.store_fingerprints {
363            return Ok(());
364        }
365        self.store_fingerprints = next;
366        let mut query = self.query.clone();
367        query
368            .harnesses
369            .retain(|harness| harness.as_str() == HarnessId::HERMES);
370        if query.harnesses.is_empty() {
371            return Ok(());
372        }
373        let fresh = descriptor_map(HarnessCatalog::new().discover_raw_index(&query));
374        let stale = self
375            .raw
376            .keys()
377            .filter(|key| key.harness == HarnessId::HERMES)
378            .cloned()
379            .collect::<Vec<_>>();
380        for key in stale {
381            if !fresh.contains_key(&key) {
382                self.raw.remove(&key);
383                content_dirty.insert(key);
384            }
385        }
386        for (key, descriptor) in fresh {
387            if self.raw.get(&key) != Some(&descriptor) {
388                self.raw.insert(key.clone(), descriptor);
389                content_dirty.insert(key);
390            }
391        }
392        Ok(())
393    }
394
395    fn refresh_path(
396        &mut self,
397        path: &Path,
398        content_dirty: &mut BTreeSet<SessionIndexKey>,
399    ) -> Result<(), String> {
400        if path.extension().and_then(|value| value.to_str()) != Some("jsonl") {
401            return Ok(());
402        }
403        let event_path = normalized_path(path);
404        let previous_key = self.paths.get(&event_path).cloned();
405        let previous = previous_key
406            .as_ref()
407            .and_then(|key| self.raw.get(key))
408            .cloned();
409        let previous_fingerprint = self.fingerprints.get(&event_path).copied();
410        let fingerprint = file_fingerprint(&event_path);
411
412        let Some(fingerprint) = fingerprint else {
413            self.fingerprints.remove(&event_path);
414            if let Some(key) = previous_key {
415                self.paths.remove(&event_path);
416                self.raw.remove(&key);
417                content_dirty.insert(key);
418            }
419            return Ok(());
420        };
421        self.fingerprints.insert(event_path.clone(), fingerprint);
422
423        let locator = previous
424            .as_ref()
425            .map(|descriptor| descriptor.locator.clone())
426            .or_else(|| locator_for_path(&self.query, &event_path));
427        let Some(locator) = locator else {
428            return Ok(());
429        };
430        let refreshed =
431            if let (Some(descriptor), Some(old)) = (previous.as_ref(), previous_fingerprint) {
432                if can_reuse_header(descriptor, old, fingerprint) {
433                    let mut descriptor = descriptor.clone();
434                    descriptor.updated_at_ms = fingerprint.modified_ms;
435                    // A conversation's name can arrive after its header was read: Codex names (and renames) a thread
436                    // in its own index, read on every write since its header never changes; Claude Code writes
437                    // `ai-title` within the header's bounds, read while an untitled transcript is still within them.
438                    if let StorageLocator::File { path } = &descriptor.locator.storage {
439                        let harness = descriptor.locator.harness.as_str();
440                        if harness == HarnessId::CODEX
441                            || (descriptor.title.is_none() && old.len < TOPIC_HEADER_BYTES)
442                        {
443                            if let Some(name) = supercode_interchange::catalog::native_topic(
444                                harness,
445                                path,
446                                &descriptor.locator.session_id,
447                            ) {
448                                descriptor.title = Some(name);
449                            }
450                        }
451                    }
452                    Some(descriptor)
453                } else {
454                    HarnessCatalog::new()
455                        .refresh_file_index_descriptor_for(&locator, &self.query)
456                        .map_err(|error| error.to_string())?
457                }
458            } else {
459                HarnessCatalog::new()
460                    .refresh_file_index_descriptor_for(&locator, &self.query)
461                    .map_err(|error| error.to_string())?
462            };
463        let Some(descriptor) = refreshed else {
464            return Ok(());
465        };
466        let key = SessionIndexKey::from_locator(&descriptor.locator);
467        if let Some(previous_key) = previous_key {
468            if previous_key != key {
469                self.raw.remove(&previous_key);
470                content_dirty.insert(previous_key);
471            }
472        }
473        self.paths.insert(event_path, key.clone());
474        self.raw.insert(key.clone(), descriptor);
475        content_dirty.insert(key);
476        Ok(())
477    }
478
479    fn rebuild_current(&mut self, content_dirty: &BTreeSet<SessionIndexKey>) -> Result<(), String> {
480        let page = self.project_current(&self.query, content_dirty)?;
481        self.current = descriptor_map(page.sessions);
482        Ok(())
483    }
484
485    fn project_current(
486        &self,
487        query: &DiscoveryQuery,
488        content_dirty: &BTreeSet<SessionIndexKey>,
489    ) -> Result<DiscoveryPage, String> {
490        let catalog = HarnessCatalog::new();
491        let mut page = catalog
492            .project_index_page(query, self.raw.values().cloned())
493            .map_err(|error| error.to_string())?;
494        let mut next = Vec::with_capacity(page.sessions.len());
495        for mut descriptor in page.sessions {
496            let key = SessionIndexKey::from_locator(&descriptor.locator);
497            if let Some(previous) = self.current.get(&key) {
498                descriptor.preview_candidates = previous.preview_candidates.clone();
499                descriptor.latest_message_candidates = previous.latest_message_candidates.clone();
500            }
501            if !self.current.contains_key(&key) || content_dirty.contains(&key) {
502                let enriched = match &self.codex_history {
503                    Some(history) => catalog.enrich_index_page_with_codex_history(
504                        query,
505                        vec![descriptor],
506                        history,
507                    ),
508                    None => catalog.enrich_index_page(query, vec![descriptor]),
509                };
510                descriptor = enriched
511                    .map_err(|error| error.to_string())?
512                    .pop()
513                    .expect("one descriptor remains one descriptor");
514            }
515            next.push(descriptor);
516        }
517        page.sessions = next;
518        Ok(page)
519    }
520}
521
522pub(crate) fn validate_query(query: &DiscoveryQuery) -> Result<(), String> {
523    if query.search_previews {
524        return Err(
525            "sessions.index.subscribe does not support preview search; use sessions.discover"
526                .into(),
527        );
528    }
529    if query.cursor.is_some() {
530        return Err("sessions.index.subscribe does not accept a cursor".into());
531    }
532    validate_limit(query.limit.unwrap_or(100))?;
533    if query.harnesses.is_empty()
534        || query.harnesses.iter().any(|harness| {
535            !matches!(
536                harness.as_str(),
537                HarnessId::CLAUDE_CODE | HarnessId::CODEX | HarnessId::HERMES
538            )
539        })
540    {
541        return Err(
542            "sessions.index.subscribe currently requires explicit claude-code, codex and/or hermes harnesses"
543                .into(),
544        );
545    }
546    Ok(())
547}
548
549pub(crate) fn validate_limit(limit: usize) -> Result<(), String> {
550    if limit == 0 || limit > MAX_SUBSCRIPTION_ROWS {
551        return Err(format!(
552            "session index limit must be between 1 and {MAX_SUBSCRIPTION_ROWS}"
553        ));
554    }
555    Ok(())
556}
557
558fn watch_roots(query: &DiscoveryQuery) -> BTreeSet<PathBuf> {
559    query
560        .harnesses
561        .iter()
562        .filter_map(|harness| match harness.as_str() {
563            HarnessId::CLAUDE_CODE => Some(query.homes.claude_code.clone()),
564            HarnessId::CODEX => Some(query.homes.codex.clone()),
565            _ => None,
566        })
567        .collect()
568}
569
570fn existing_watch_root(root: &Path) -> Option<PathBuf> {
571    if root.is_dir() {
572        return Some(root.to_path_buf());
573    }
574    // Watching an entire home directory because a harness has never created
575    // its store is disproportionate. One parent level catches the ordinary
576    // first-run mkdir; the recovery reconciliation handles rarer deeper gaps.
577    root.parent()
578        .filter(|parent| parent.is_dir())
579        .map(Path::to_path_buf)
580}
581
582/// Whole-store SQLite files named by the query (one file = every session of that harness).
583fn store_paths(query: &DiscoveryQuery) -> BTreeSet<PathBuf> {
584    query
585        .harnesses
586        .iter()
587        .flat_map(|harness| match harness.as_str() {
588            HarnessId::HERMES => hermes_session_stores(&query.homes.hermes),
589            _ => Vec::new(),
590        })
591        .collect()
592}
593
594/// A store's stamp set: the file itself and its WAL-mode siblings, which is where a live writer's
595/// appends land until a checkpoint. Absent files are kept as `None` so their creation is a change.
596fn scan_store_fingerprints(query: &DiscoveryQuery) -> BTreeMap<PathBuf, Option<FileFingerprint>> {
597    let mut stamps = BTreeMap::new();
598    for store in store_paths(query) {
599        for path in store_sibling_paths(&store) {
600            let stamp = file_fingerprint(&path);
601            stamps.insert(normalized_store_path(&path), stamp);
602        }
603    }
604    stamps
605}
606
607/// The files that hold a store's committed content: the database and its write-ahead log. The
608/// `-shm` sibling is SQLite's shared-memory index, which every reader writes (this index's own
609/// read-only enumeration included): stamping it made each enumeration trigger the next one.
610fn store_sibling_paths(store: &Path) -> [PathBuf; 2] {
611    let name = store
612        .file_name()
613        .and_then(|value| value.to_str())
614        .unwrap_or("state.db");
615    [
616        store.to_path_buf(),
617        store.with_file_name(format!("{name}-wal")),
618    ]
619}
620
621/// A watched store's `-shm` sibling: its changes are readers' bookkeeping, not content.
622fn is_store_shm(stores: &BTreeMap<PathBuf, Option<FileFingerprint>>, path: &Path) -> bool {
623    let path = normalized_store_path(path);
624    path.to_str()
625        .and_then(|value| value.strip_suffix("-shm"))
626        .is_some_and(|store| stores.contains_key(Path::new(store)))
627}
628
629/// Store siblings come and go, so canonicalize through the (stable) directory rather than the file.
630fn normalized_store_path(path: &Path) -> PathBuf {
631    match (path.parent(), path.file_name()) {
632        (Some(dir), Some(name)) => normalized_path(dir).join(name),
633        _ => path.to_path_buf(),
634    }
635}
636
637fn locator_for_path(query: &DiscoveryQuery, path: &Path) -> Option<SessionLocator> {
638    let claude_root = normalized_path(&query.homes.claude_code);
639    let codex_root = normalized_path(&query.homes.codex);
640    let harness = if query
641        .harnesses
642        .iter()
643        .any(|harness| harness.as_str() == HarnessId::CLAUDE_CODE)
644        && path.starts_with(&claude_root)
645    {
646        HarnessId::CLAUDE_CODE
647    } else if query
648        .harnesses
649        .iter()
650        .any(|harness| harness.as_str() == HarnessId::CODEX)
651        && path.starts_with(&codex_root)
652    {
653        HarnessId::CODEX
654    } else {
655        return None;
656    };
657    Some(SessionLocator {
658        harness: HarnessId::new(harness),
659        session_id: path
660            .file_stem()
661            .and_then(|value| value.to_str())
662            .unwrap_or("unknown")
663            .to_string(),
664        storage: StorageLocator::File {
665            path: path.to_path_buf(),
666        },
667    })
668}
669
670fn normalized_path(path: &Path) -> PathBuf {
671    fs::canonicalize(path).unwrap_or_else(|_| path.to_path_buf())
672}
673
674fn descriptor_path_map(
675    descriptors: &BTreeMap<SessionIndexKey, SessionDescriptor>,
676) -> BTreeMap<PathBuf, SessionIndexKey> {
677    descriptors
678        .iter()
679        .map(|(key, descriptor)| {
680            (
681                normalized_path(descriptor.locator.storage.path()),
682                key.clone(),
683            )
684        })
685        .collect()
686}
687
688fn scan_file_fingerprints(query: &DiscoveryQuery) -> BTreeMap<PathBuf, FileFingerprint> {
689    let mut paths = Vec::new();
690    for root in watch_roots(query) {
691        collect_jsonl_paths(&root, &mut paths);
692    }
693    paths
694        .into_iter()
695        .filter_map(|path| {
696            let path = normalized_path(&path);
697            file_fingerprint(&path).map(|fingerprint| (path, fingerprint))
698        })
699        .collect()
700}
701
702fn collect_jsonl_paths(root: &Path, paths: &mut Vec<PathBuf>) {
703    let mut walked = BTreeSet::new();
704    collect_jsonl_paths_in(root, paths, &mut walked);
705}
706
707/// Follows symlinked directories, as the catalog's walker does: a project directory moved to
708/// another volume and linked back is still this home's. `walked` ends a link cycle.
709fn collect_jsonl_paths_in(root: &Path, paths: &mut Vec<PathBuf>, walked: &mut BTreeSet<PathBuf>) {
710    if !walked.insert(fs::canonicalize(root).unwrap_or_else(|_| root.to_path_buf())) {
711        return;
712    }
713    let Ok(entries) = fs::read_dir(root) else {
714        return;
715    };
716    for entry in entries.flatten() {
717        let Ok(mut file_type) = entry.file_type() else {
718            continue;
719        };
720        let path = entry.path();
721        if file_type.is_symlink() {
722            let Ok(target) = fs::metadata(&path) else {
723                continue;
724            };
725            file_type = target.file_type();
726        }
727        if file_type.is_dir() {
728            collect_jsonl_paths_in(&path, paths, walked);
729        } else if file_type.is_file()
730            && path.extension().and_then(|value| value.to_str()) == Some("jsonl")
731        {
732            paths.push(path);
733        }
734    }
735}
736
737fn file_fingerprint(path: &Path) -> Option<FileFingerprint> {
738    let metadata = fs::metadata(path).ok()?;
739    let modified = metadata.modified().ok()?.duration_since(UNIX_EPOCH).ok()?;
740    #[cfg(unix)]
741    let identity = {
742        use std::os::unix::fs::MetadataExt;
743        (u128::from(metadata.dev()) << 64) | u128::from(metadata.ino())
744    };
745    #[cfg(not(unix))]
746    let identity = 0;
747    Some(FileFingerprint {
748        len: metadata.len(),
749        modified_ns: modified.as_nanos(),
750        modified_ms: u64::try_from(modified.as_millis()).ok(),
751        identity,
752    })
753}
754
755/// The bytes a header read covers (`read_header`'s bound): a Claude Code transcript past them without a topic will not
756/// gain one there.
757const TOPIC_HEADER_BYTES: u64 = 256 * 1024;
758
759fn can_reuse_header(
760    descriptor: &SessionDescriptor,
761    previous: FileFingerprint,
762    current: FileFingerprint,
763) -> bool {
764    previous.identity == current.identity
765        && previous.len <= current.len
766        && descriptor.cwd.is_some()
767        && descriptor.model.is_some()
768        && !descriptor.locator.session_id.is_empty()
769}
770
771fn descriptor_map(
772    descriptors: impl IntoIterator<Item = SessionDescriptor>,
773) -> BTreeMap<SessionIndexKey, SessionDescriptor> {
774    descriptors
775        .into_iter()
776        .map(|descriptor| {
777            (
778                SessionIndexKey::from_locator(&descriptor.locator),
779                descriptor,
780            )
781        })
782        .collect()
783}
784
785fn diff_descriptors(
786    before: &BTreeMap<SessionIndexKey, SessionDescriptor>,
787    after: &BTreeMap<SessionIndexKey, SessionDescriptor>,
788) -> Vec<SessionIndexChange> {
789    let mut changes = Vec::new();
790    for (key, descriptor) in after {
791        match before.get(key) {
792            None => changes.push(SessionIndexChange::Added {
793                descriptor: descriptor.clone(),
794            }),
795            Some(previous) if previous != descriptor => {
796                changes.push(SessionIndexChange::Updated {
797                    descriptor: descriptor.clone(),
798                });
799            }
800            Some(_) => {}
801        }
802    }
803    for key in before.keys() {
804        if !after.contains_key(key) {
805            changes.push(SessionIndexChange::Removed { key: key.clone() });
806        }
807    }
808    changes
809}