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                    Some(descriptor)
436                } else {
437                    HarnessCatalog::new()
438                        .refresh_file_index_descriptor_for(&locator, &self.query)
439                        .map_err(|error| error.to_string())?
440                }
441            } else {
442                HarnessCatalog::new()
443                    .refresh_file_index_descriptor_for(&locator, &self.query)
444                    .map_err(|error| error.to_string())?
445            };
446        let Some(descriptor) = refreshed else {
447            return Ok(());
448        };
449        let key = SessionIndexKey::from_locator(&descriptor.locator);
450        if let Some(previous_key) = previous_key {
451            if previous_key != key {
452                self.raw.remove(&previous_key);
453                content_dirty.insert(previous_key);
454            }
455        }
456        self.paths.insert(event_path, key.clone());
457        self.raw.insert(key.clone(), descriptor);
458        content_dirty.insert(key);
459        Ok(())
460    }
461
462    fn rebuild_current(&mut self, content_dirty: &BTreeSet<SessionIndexKey>) -> Result<(), String> {
463        let page = self.project_current(&self.query, content_dirty)?;
464        self.current = descriptor_map(page.sessions);
465        Ok(())
466    }
467
468    fn project_current(
469        &self,
470        query: &DiscoveryQuery,
471        content_dirty: &BTreeSet<SessionIndexKey>,
472    ) -> Result<DiscoveryPage, String> {
473        let catalog = HarnessCatalog::new();
474        let mut page = catalog
475            .project_index_page(query, self.raw.values().cloned())
476            .map_err(|error| error.to_string())?;
477        let mut next = Vec::with_capacity(page.sessions.len());
478        for mut descriptor in page.sessions {
479            let key = SessionIndexKey::from_locator(&descriptor.locator);
480            if let Some(previous) = self.current.get(&key) {
481                descriptor.preview_candidates = previous.preview_candidates.clone();
482                descriptor.latest_message_candidates = previous.latest_message_candidates.clone();
483            }
484            if !self.current.contains_key(&key) || content_dirty.contains(&key) {
485                let enriched = match &self.codex_history {
486                    Some(history) => catalog.enrich_index_page_with_codex_history(
487                        query,
488                        vec![descriptor],
489                        history,
490                    ),
491                    None => catalog.enrich_index_page(query, vec![descriptor]),
492                };
493                descriptor = enriched
494                    .map_err(|error| error.to_string())?
495                    .pop()
496                    .expect("one descriptor remains one descriptor");
497            }
498            next.push(descriptor);
499        }
500        page.sessions = next;
501        Ok(page)
502    }
503}
504
505pub(crate) fn validate_query(query: &DiscoveryQuery) -> Result<(), String> {
506    if query.search_previews {
507        return Err(
508            "sessions.index.subscribe does not support preview search; use sessions.discover"
509                .into(),
510        );
511    }
512    if query.cursor.is_some() {
513        return Err("sessions.index.subscribe does not accept a cursor".into());
514    }
515    validate_limit(query.limit.unwrap_or(100))?;
516    if query.harnesses.is_empty()
517        || query.harnesses.iter().any(|harness| {
518            !matches!(
519                harness.as_str(),
520                HarnessId::CLAUDE_CODE | HarnessId::CODEX | HarnessId::HERMES
521            )
522        })
523    {
524        return Err(
525            "sessions.index.subscribe currently requires explicit claude-code, codex and/or hermes harnesses"
526                .into(),
527        );
528    }
529    Ok(())
530}
531
532pub(crate) fn validate_limit(limit: usize) -> Result<(), String> {
533    if limit == 0 || limit > MAX_SUBSCRIPTION_ROWS {
534        return Err(format!(
535            "session index limit must be between 1 and {MAX_SUBSCRIPTION_ROWS}"
536        ));
537    }
538    Ok(())
539}
540
541fn watch_roots(query: &DiscoveryQuery) -> BTreeSet<PathBuf> {
542    query
543        .harnesses
544        .iter()
545        .filter_map(|harness| match harness.as_str() {
546            HarnessId::CLAUDE_CODE => Some(query.homes.claude_code.clone()),
547            HarnessId::CODEX => Some(query.homes.codex.clone()),
548            _ => None,
549        })
550        .collect()
551}
552
553fn existing_watch_root(root: &Path) -> Option<PathBuf> {
554    if root.is_dir() {
555        return Some(root.to_path_buf());
556    }
557    // Watching an entire home directory because a harness has never created
558    // its store is disproportionate. One parent level catches the ordinary
559    // first-run mkdir; the recovery reconciliation handles rarer deeper gaps.
560    root.parent()
561        .filter(|parent| parent.is_dir())
562        .map(Path::to_path_buf)
563}
564
565/// Whole-store SQLite files named by the query (one file = every session of that harness).
566fn store_paths(query: &DiscoveryQuery) -> BTreeSet<PathBuf> {
567    query
568        .harnesses
569        .iter()
570        .flat_map(|harness| match harness.as_str() {
571            HarnessId::HERMES => hermes_session_stores(&query.homes.hermes),
572            _ => Vec::new(),
573        })
574        .collect()
575}
576
577/// A store's stamp set: the file itself and its WAL-mode siblings, which is where a live writer's
578/// appends land until a checkpoint. Absent files are kept as `None` so their creation is a change.
579fn scan_store_fingerprints(query: &DiscoveryQuery) -> BTreeMap<PathBuf, Option<FileFingerprint>> {
580    let mut stamps = BTreeMap::new();
581    for store in store_paths(query) {
582        for path in store_sibling_paths(&store) {
583            let stamp = file_fingerprint(&path);
584            stamps.insert(normalized_store_path(&path), stamp);
585        }
586    }
587    stamps
588}
589
590/// The files that hold a store's committed content: the database and its write-ahead log. The
591/// `-shm` sibling is SQLite's shared-memory index, which every reader writes (this index's own
592/// read-only enumeration included): stamping it made each enumeration trigger the next one.
593fn store_sibling_paths(store: &Path) -> [PathBuf; 2] {
594    let name = store
595        .file_name()
596        .and_then(|value| value.to_str())
597        .unwrap_or("state.db");
598    [
599        store.to_path_buf(),
600        store.with_file_name(format!("{name}-wal")),
601    ]
602}
603
604/// A watched store's `-shm` sibling: its changes are readers' bookkeeping, not content.
605fn is_store_shm(stores: &BTreeMap<PathBuf, Option<FileFingerprint>>, path: &Path) -> bool {
606    let path = normalized_store_path(path);
607    path.to_str()
608        .and_then(|value| value.strip_suffix("-shm"))
609        .is_some_and(|store| stores.contains_key(Path::new(store)))
610}
611
612/// Store siblings come and go, so canonicalize through the (stable) directory rather than the file.
613fn normalized_store_path(path: &Path) -> PathBuf {
614    match (path.parent(), path.file_name()) {
615        (Some(dir), Some(name)) => normalized_path(dir).join(name),
616        _ => path.to_path_buf(),
617    }
618}
619
620fn locator_for_path(query: &DiscoveryQuery, path: &Path) -> Option<SessionLocator> {
621    let claude_root = normalized_path(&query.homes.claude_code);
622    let codex_root = normalized_path(&query.homes.codex);
623    let harness = if query
624        .harnesses
625        .iter()
626        .any(|harness| harness.as_str() == HarnessId::CLAUDE_CODE)
627        && path.starts_with(&claude_root)
628    {
629        HarnessId::CLAUDE_CODE
630    } else if query
631        .harnesses
632        .iter()
633        .any(|harness| harness.as_str() == HarnessId::CODEX)
634        && path.starts_with(&codex_root)
635    {
636        HarnessId::CODEX
637    } else {
638        return None;
639    };
640    Some(SessionLocator {
641        harness: HarnessId::new(harness),
642        session_id: path
643            .file_stem()
644            .and_then(|value| value.to_str())
645            .unwrap_or("unknown")
646            .to_string(),
647        storage: StorageLocator::File {
648            path: path.to_path_buf(),
649        },
650    })
651}
652
653fn normalized_path(path: &Path) -> PathBuf {
654    fs::canonicalize(path).unwrap_or_else(|_| path.to_path_buf())
655}
656
657fn descriptor_path_map(
658    descriptors: &BTreeMap<SessionIndexKey, SessionDescriptor>,
659) -> BTreeMap<PathBuf, SessionIndexKey> {
660    descriptors
661        .iter()
662        .map(|(key, descriptor)| {
663            (
664                normalized_path(descriptor.locator.storage.path()),
665                key.clone(),
666            )
667        })
668        .collect()
669}
670
671fn scan_file_fingerprints(query: &DiscoveryQuery) -> BTreeMap<PathBuf, FileFingerprint> {
672    let mut paths = Vec::new();
673    for root in watch_roots(query) {
674        collect_jsonl_paths(&root, &mut paths);
675    }
676    paths
677        .into_iter()
678        .filter_map(|path| {
679            let path = normalized_path(&path);
680            file_fingerprint(&path).map(|fingerprint| (path, fingerprint))
681        })
682        .collect()
683}
684
685fn collect_jsonl_paths(root: &Path, paths: &mut Vec<PathBuf>) {
686    let mut walked = BTreeSet::new();
687    collect_jsonl_paths_in(root, paths, &mut walked);
688}
689
690/// Follows symlinked directories, as the catalog's walker does: a project directory moved to
691/// another volume and linked back is still this home's. `walked` ends a link cycle.
692fn collect_jsonl_paths_in(root: &Path, paths: &mut Vec<PathBuf>, walked: &mut BTreeSet<PathBuf>) {
693    if !walked.insert(fs::canonicalize(root).unwrap_or_else(|_| root.to_path_buf())) {
694        return;
695    }
696    let Ok(entries) = fs::read_dir(root) else {
697        return;
698    };
699    for entry in entries.flatten() {
700        let Ok(mut file_type) = entry.file_type() else {
701            continue;
702        };
703        let path = entry.path();
704        if file_type.is_symlink() {
705            let Ok(target) = fs::metadata(&path) else {
706                continue;
707            };
708            file_type = target.file_type();
709        }
710        if file_type.is_dir() {
711            collect_jsonl_paths_in(&path, paths, walked);
712        } else if file_type.is_file()
713            && path.extension().and_then(|value| value.to_str()) == Some("jsonl")
714        {
715            paths.push(path);
716        }
717    }
718}
719
720fn file_fingerprint(path: &Path) -> Option<FileFingerprint> {
721    let metadata = fs::metadata(path).ok()?;
722    let modified = metadata.modified().ok()?.duration_since(UNIX_EPOCH).ok()?;
723    #[cfg(unix)]
724    let identity = {
725        use std::os::unix::fs::MetadataExt;
726        (u128::from(metadata.dev()) << 64) | u128::from(metadata.ino())
727    };
728    #[cfg(not(unix))]
729    let identity = 0;
730    Some(FileFingerprint {
731        len: metadata.len(),
732        modified_ns: modified.as_nanos(),
733        modified_ms: u64::try_from(modified.as_millis()).ok(),
734        identity,
735    })
736}
737
738fn can_reuse_header(
739    descriptor: &SessionDescriptor,
740    previous: FileFingerprint,
741    current: FileFingerprint,
742) -> bool {
743    previous.identity == current.identity
744        && previous.len <= current.len
745        && descriptor.cwd.is_some()
746        && descriptor.model.is_some()
747        && !descriptor.locator.session_id.is_empty()
748}
749
750fn descriptor_map(
751    descriptors: impl IntoIterator<Item = SessionDescriptor>,
752) -> BTreeMap<SessionIndexKey, SessionDescriptor> {
753    descriptors
754        .into_iter()
755        .map(|descriptor| {
756            (
757                SessionIndexKey::from_locator(&descriptor.locator),
758                descriptor,
759            )
760        })
761        .collect()
762}
763
764fn diff_descriptors(
765    before: &BTreeMap<SessionIndexKey, SessionDescriptor>,
766    after: &BTreeMap<SessionIndexKey, SessionDescriptor>,
767) -> Vec<SessionIndexChange> {
768    let mut changes = Vec::new();
769    for (key, descriptor) in after {
770        match before.get(key) {
771            None => changes.push(SessionIndexChange::Added {
772                descriptor: descriptor.clone(),
773            }),
774            Some(previous) if previous != descriptor => {
775                changes.push(SessionIndexChange::Updated {
776                    descriptor: descriptor.clone(),
777                });
778            }
779            Some(_) => {}
780        }
781    }
782    for key in before.keys() {
783        if !after.contains_key(key) {
784            changes.push(SessionIndexChange::Removed { key: key.clone() });
785        }
786    }
787    changes
788}