1use 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#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize)]
34pub struct SessionIndexKey {
35 pub harness: String,
37 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#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
52#[serde(tag = "kind", rename_all = "snake_case")]
53pub enum SessionIndexChange {
54 Added {
56 descriptor: SessionDescriptor,
58 },
59 Updated {
61 descriptor: SessionDescriptor,
63 },
64 Removed {
66 key: SessionIndexKey,
68 },
69}
70
71#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
76pub struct SessionIndexDelta {
77 pub revision: u64,
79 pub changes: Vec<SessionIndexChange>,
81}
82
83pub(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 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
111pub(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 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 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(¬ifier);
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 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 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 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 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 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 root.parent()
561 .filter(|parent| parent.is_dir())
562 .map(Path::to_path_buf)
563}
564
565fn 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
577fn 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
590fn 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
604fn 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
612fn 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
690fn 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}
789
790#[cfg(test)]
791mod tests {
792 use super::*;
793
794 #[test]
795 fn preview_search_is_refused_before_opening_a_retained_index() {
796 let query = DiscoveryQuery {
797 search_previews: true,
798 query: Some("nebula".into()),
799 ..DiscoveryQuery::default()
800 };
801 let error = match SessionIndexSubscription::open(query, Arc::new(Notify::new())) {
802 Err(error) => error,
803 Ok(_) => panic!("preview search must not open live watchers"),
804 };
805 assert!(error.contains("use sessions.discover"), "{error}");
806 }
807
808 fn descriptor(id: &str, updated_at_ms: u64) -> SessionDescriptor {
809 SessionDescriptor {
810 locator: SessionLocator {
811 harness: HarnessId::new(HarnessId::CODEX),
812 session_id: id.into(),
813 storage: StorageLocator::File {
814 path: PathBuf::from(format!("/{id}.jsonl")),
815 },
816 },
817 cwd: None,
818 title: None,
819 preview_candidates: Vec::new(),
820 latest_message_candidates: Vec::new(),
821 updated_at_ms: Some(updated_at_ms),
822 message_count: None,
823 model: None,
824 parent_session_id: None,
825 child_session_count: 0,
826 nouns: Default::default(),
827 }
828 }
829
830 #[test]
831 fn resize_retains_index_watcher_and_cached_previews_until_commit() {
832 let root = std::env::temp_dir().join(format!(
833 "supercode-index-resize-{}-{}",
834 std::process::id(),
835 std::time::SystemTime::now()
836 .duration_since(UNIX_EPOCH)
837 .unwrap()
838 .as_nanos()
839 ));
840 fs::create_dir_all(&root).unwrap();
841 let root = root.canonicalize().unwrap();
842 let query = DiscoveryQuery {
843 harnesses: vec![HarnessId::new(HarnessId::CODEX)],
844 homes: crate::HarnessHomes {
845 codex: root.clone(),
846 ..crate::HarnessHomes::default()
847 },
848 limit: Some(1),
849 include_topic_candidates: true,
850 ..DiscoveryQuery::default()
851 };
852 let (mut index, _) =
853 SessionIndexSubscription::open(query, Arc::new(Notify::new())).unwrap();
854 for (id, updated) in [("newest", 3), ("middle", 2), ("oldest", 1)] {
855 let path = root.join(format!("{id}.jsonl"));
856 fs::write(&path, format!(
857 "{{\"type\":\"session_meta\",\"payload\":{{\"id\":\"{id}\",\"cwd\":\"/workspace\"}}}}\n{{\"type\":\"response_item\",\"payload\":{{\"type\":\"message\",\"role\":\"user\",\"content\":[{{\"type\":\"input_text\",\"text\":\"topic {id}\"}}]}}}}\n"
858 )).unwrap();
859 let mut row = descriptor(id, updated);
860 row.locator.storage = StorageLocator::File { path };
861 index
862 .raw
863 .insert(SessionIndexKey::from_locator(&row.locator), row);
864 }
865 index.paths = descriptor_path_map(&index.raw);
866 index.rebuild_current(&BTreeSet::new()).unwrap();
867 let original = index.current.clone();
868 assert!(!original
869 .values()
870 .next()
871 .unwrap()
872 .preview_candidates
873 .is_empty());
874 fs::write(root.join("newest.jsonl"), "").unwrap();
876 let queued = root.join("queued.jsonl");
878 fs::write(
879 &queued,
880 "{\"type\":\"session_meta\",\"payload\":{\"id\":\"queued\",\"cwd\":\"/workspace\"}}\n",
881 )
882 .unwrap();
883 let (sender, receiver) = mpsc::channel();
884 let _native_receiver = std::mem::replace(&mut index.receiver, receiver);
886 sender
887 .send(Ok(Event::new(notify::EventKind::Any).add_path(queued)))
888 .unwrap();
889 let watcher = &index._watcher as *const _;
890 let raw_row = index.raw.values().next().unwrap() as *const _;
891 let raw = index.raw.clone();
892 let reconcile = index.last_reconcile;
893 let prepared = index.prepare_resize(2).unwrap();
894 assert_eq!(prepared.revision, 2);
895 assert_eq!(prepared.page.receipt.total_matched, 3);
896 assert_eq!(prepared.page.sessions.len(), 2);
897 assert_eq!(
898 prepared.page.sessions[0],
899 *original.values().next().unwrap()
900 );
901 assert!(!prepared.page.sessions[1].preview_candidates.is_empty());
902 assert_eq!(index.current, original);
903 assert_eq!(index.revision, 1);
904 drop(prepared); assert!(index.prepare_resize(0).is_err());
906 assert!(index.prepare_resize(2049).is_err());
907 assert_eq!(index.current, original);
908 assert_eq!(index.revision, 1);
909 let prepared = index.prepare_resize(2).unwrap();
910 index.commit_resize(prepared);
911 assert_eq!(index.raw, raw);
912 assert_eq!(index.raw.values().next().unwrap() as *const _, raw_row);
913 assert_eq!(&index._watcher as *const _, watcher);
914 assert_eq!(index.last_reconcile, reconcile);
915 assert_eq!(index.prepare_resize(2).unwrap().revision, 2);
916 let delta = index.poll().unwrap().unwrap();
917 assert_eq!(delta.revision, 3);
918 let shrink = index.prepare_resize(1).unwrap();
919 assert_eq!(shrink.revision, 4);
920 index.commit_resize(shrink);
921 assert_eq!(index.current.len(), 1);
922 assert_eq!(index.prepare_resize(1).unwrap().revision, 4);
923 drop(index);
924 fs::remove_dir_all(root).unwrap();
925 }
926
927 #[test]
928 fn same_limit_receipt_counts_out_of_window_changes_without_visible_revision() {
929 let root = std::env::temp_dir().join(format!(
930 "supercode-index-total-{}-{}",
931 std::process::id(),
932 std::time::SystemTime::now()
933 .duration_since(UNIX_EPOCH)
934 .unwrap()
935 .as_nanos()
936 ));
937 fs::create_dir_all(&root).unwrap();
938 let root = root.canonicalize().unwrap();
939 let query = DiscoveryQuery {
940 harnesses: vec![HarnessId::new(HarnessId::CODEX)],
941 homes: crate::HarnessHomes {
942 codex: root.clone(),
943 ..crate::HarnessHomes::default()
944 },
945 limit: Some(1),
946 ..DiscoveryQuery::default()
947 };
948 let (mut index, _) =
949 SessionIndexSubscription::open(query, Arc::new(Notify::new())).unwrap();
950 let visible = descriptor("visible", u64::MAX);
951 index.raw = descriptor_map([visible]);
952 index.rebuild_current(&BTreeSet::new()).unwrap();
953 let (sender, receiver) = mpsc::channel();
954 let _native_receiver = std::mem::replace(&mut index.receiver, receiver);
955 let hidden = root.join("hidden.jsonl");
956 fs::write(
957 &hidden,
958 "{\"type\":\"session_meta\",\"payload\":{\"id\":\"hidden\",\"cwd\":\"/workspace\"}}\n",
959 )
960 .unwrap();
961 sender
962 .send(Ok(
963 Event::new(notify::EventKind::Any).add_path(hidden.clone())
964 ))
965 .unwrap();
966 assert!(index.poll().unwrap().is_none());
967 let added = index.prepare_resize(1).unwrap();
968 assert_eq!(added.revision, 1);
969 assert_eq!(added.page.receipt.total_matched, 2);
970 fs::remove_file(&hidden).unwrap();
971 sender
972 .send(Ok(Event::new(notify::EventKind::Any).add_path(hidden)))
973 .unwrap();
974 assert!(index.poll().unwrap().is_none());
975 let removed = index.prepare_resize(1).unwrap();
976 assert_eq!(removed.revision, 1);
977 assert_eq!(removed.page.receipt.total_matched, 1);
978 drop(index);
979 fs::remove_dir_all(root).unwrap();
980 }
981
982 #[test]
983 fn hermes_store_appends_surface_as_index_updates() {
984 let root = std::env::temp_dir().join(format!(
987 "supercode-index-hermes-{}-{}",
988 std::process::id(),
989 std::time::SystemTime::now()
990 .duration_since(UNIX_EPOCH)
991 .unwrap()
992 .as_nanos()
993 ));
994 fs::create_dir_all(&root).unwrap();
995 let db = root.join("state.db");
996 fs::copy(
997 PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures/hermes_home/state.db"),
998 &db,
999 )
1000 .unwrap();
1001 let query = DiscoveryQuery {
1002 harnesses: vec![HarnessId::new(HarnessId::HERMES)],
1003 homes: crate::HarnessHomes {
1004 hermes: db.clone(),
1005 claude_code: root.join("missing-claude"),
1006 codex: root.join("missing-codex"),
1007 ..crate::HarnessHomes::default()
1008 },
1009 ..DiscoveryQuery::default()
1010 };
1011 let (mut subscription, initial) =
1012 SessionIndexSubscription::open(query, Arc::new(Notify::new())).unwrap();
1013 assert!(initial.len() >= 2, "{initial:#?}");
1014 assert!(initial
1015 .iter()
1016 .all(|descriptor| descriptor.locator.harness.as_str() == HarnessId::HERMES));
1017 assert!(
1018 subscription.poll().unwrap().is_none(),
1019 "quiet store, quiet index"
1020 );
1021
1022 let target = initial[0].locator.session_id.clone();
1024 std::thread::sleep(Duration::from_millis(20));
1025 {
1026 let conn = rusqlite::Connection::open(&db).unwrap();
1027 conn.execute(
1028 "INSERT INTO messages (session_id, role, content, timestamp, active) VALUES (?1, 'assistant', 'index test append', ?2, 1)",
1029 rusqlite::params![target, 1_800_000_000.0_f64],
1030 )
1031 .unwrap();
1032 conn.execute(
1033 "UPDATE sessions SET message_count = message_count + 1, ended_at = ?2 WHERE id = ?1",
1034 rusqlite::params![target, 1_800_000_000.0_f64],
1035 )
1036 .unwrap();
1037 }
1038 let deadline = Instant::now() + Duration::from_secs(5);
1040 let delta = loop {
1041 if let Some(delta) = subscription.poll().unwrap() {
1042 break delta;
1043 }
1044 assert!(
1045 Instant::now() < deadline,
1046 "no index delta after the store append"
1047 );
1048 std::thread::sleep(Duration::from_millis(50));
1049 };
1050 assert_eq!(delta.changes.len(), 1, "{delta:#?}");
1051 match &delta.changes[0] {
1052 SessionIndexChange::Updated { descriptor } => {
1053 assert_eq!(descriptor.locator.session_id, target);
1054 assert_eq!(
1055 descriptor.message_count,
1056 initial[0].message_count.map(|count| count + 1)
1057 );
1058 }
1059 other => panic!("expected an update for {target}, got {other:?}"),
1060 }
1061 assert!(
1062 subscription.poll().unwrap().is_none(),
1063 "one append, one delta"
1064 );
1065 fs::remove_dir_all(&root).ok();
1066 }
1067
1068 #[test]
1069 fn index_delta_is_a_complete_deterministic_replacement_set() {
1070 let before = descriptor_map([descriptor("removed", 1), descriptor("updated", 2)]);
1071 let after = descriptor_map([descriptor("updated", 3), descriptor("added", 4)]);
1072 let changes = diff_descriptors(&before, &after);
1073 assert!(matches!(
1074 &changes[0],
1075 SessionIndexChange::Added { descriptor } if descriptor.locator.session_id == "added"
1076 ));
1077 assert!(matches!(
1078 &changes[1],
1079 SessionIndexChange::Updated { descriptor } if descriptor.locator.session_id == "updated"
1080 ));
1081 assert!(matches!(
1082 &changes[2],
1083 SessionIndexChange::Removed { key } if key.session_id == "removed"
1084 ));
1085 }
1086
1087 #[test]
1088 fn raw_index_projects_child_activity_into_one_root_row() {
1089 let root = descriptor("root", 10);
1090 let mut child = descriptor("child", 20);
1091 child.parent_session_id = Some("root".into());
1092 let query = DiscoveryQuery {
1093 harnesses: vec![HarnessId::new(HarnessId::CODEX)],
1094 limit: Some(100),
1095 ..DiscoveryQuery::default()
1096 };
1097
1098 let projected = HarnessCatalog::new()
1099 .project_index(&query, [root, child])
1100 .unwrap();
1101
1102 assert_eq!(projected.len(), 1);
1103 assert_eq!(projected[0].locator.session_id, "root");
1104 assert_eq!(projected[0].updated_at_ms, Some(20));
1105 assert_eq!(projected[0].child_session_count, 1);
1106 }
1107
1108 #[test]
1109 fn complete_raw_index_backfills_a_bounded_page_without_discovery() {
1110 let query = DiscoveryQuery {
1111 harnesses: vec![HarnessId::new(HarnessId::CODEX)],
1112 limit: Some(2),
1113 ..DiscoveryQuery::default()
1114 };
1115 let catalog = HarnessCatalog::new();
1116 let mut raw = descriptor_map([
1117 descriptor("oldest", 1),
1118 descriptor("middle", 2),
1119 descriptor("newest", 3),
1120 ]);
1121 let initial = catalog
1122 .project_index(&query, raw.values().cloned())
1123 .unwrap();
1124 assert_eq!(
1125 initial
1126 .iter()
1127 .map(|descriptor| descriptor.locator.session_id.as_str())
1128 .collect::<Vec<_>>(),
1129 ["newest", "middle"]
1130 );
1131
1132 raw.remove(&SessionIndexKey {
1133 harness: HarnessId::CODEX.into(),
1134 session_id: "newest".into(),
1135 });
1136 let after = catalog
1137 .project_index(&query, raw.values().cloned())
1138 .unwrap();
1139 assert_eq!(
1140 after
1141 .iter()
1142 .map(|descriptor| descriptor.locator.session_id.as_str())
1143 .collect::<Vec<_>>(),
1144 ["middle", "oldest"]
1145 );
1146 }
1147
1148 #[test]
1149 fn append_reuses_an_immutable_header_but_replacement_does_not() {
1150 let mut existing = descriptor("session", 1);
1151 existing.cwd = Some(PathBuf::from("/workspace"));
1152 existing.model = Some("model".into());
1153 let before = FileFingerprint {
1154 len: 100,
1155 modified_ns: 1,
1156 modified_ms: Some(1),
1157 identity: 7,
1158 };
1159 let append = FileFingerprint {
1160 len: 200,
1161 modified_ns: 2,
1162 modified_ms: Some(2),
1163 identity: 7,
1164 };
1165 let replacement = FileFingerprint {
1166 identity: 8,
1167 ..append
1168 };
1169
1170 assert!(can_reuse_header(&existing, before, append));
1171 assert!(!can_reuse_header(&existing, before, replacement));
1172 }
1173
1174 #[tokio::test]
1175 async fn filesystem_event_wakes_index_without_a_poll_timer() {
1176 let nonce = std::time::SystemTime::now()
1177 .duration_since(UNIX_EPOCH)
1178 .unwrap()
1179 .as_nanos();
1180 let root = std::env::temp_dir().join(format!(
1181 "supercode-session-index-{}-{nonce}",
1182 std::process::id()
1183 ));
1184 let codex = root.join("codex");
1185 fs::create_dir_all(&codex).unwrap();
1186 let query = DiscoveryQuery {
1187 harnesses: vec![HarnessId::new(HarnessId::CODEX)],
1188 homes: crate::HarnessHomes {
1189 codex: codex.clone(),
1190 ..crate::HarnessHomes::default()
1191 },
1192 limit: Some(10),
1193 ..DiscoveryQuery::default()
1194 };
1195 let notifier = Arc::new(Notify::new());
1196 let (mut index, initial) =
1197 SessionIndexSubscription::open(query, Arc::clone(¬ifier)).unwrap();
1198 assert!(initial.is_empty());
1199
1200 let session = codex.join("new.jsonl");
1201 fs::write(
1202 &session,
1203 concat!(
1204 "{\"type\":\"session_meta\",\"payload\":{\"id\":\"new\",\"cwd\":\"/workspace\"}}\n",
1205 "{\"type\":\"turn_context\",\"payload\":{\"cwd\":\"/workspace\",\"model\":\"gpt-test\"}}\n"
1206 ),
1207 )
1208 .unwrap();
1209
1210 tokio::time::timeout(Duration::from_secs(5), notifier.notified())
1211 .await
1212 .expect("filesystem invalidation should wake the index");
1213 let delta = index
1214 .poll()
1215 .unwrap()
1216 .expect("the filesystem event should produce a visible delta");
1217 assert!(matches!(
1218 &delta.changes[0],
1219 SessionIndexChange::Added { descriptor }
1220 if descriptor.locator.session_id == "new"
1221 ));
1222
1223 drop(index);
1224 fs::remove_dir_all(root).unwrap();
1225 }
1226}