1use std::collections::{BTreeSet, HashMap, HashSet};
7use std::fs::{self, File};
8use std::io::{BufRead, BufReader, Read, Seek, SeekFrom};
9use std::path::{Component, Path, PathBuf};
10use std::time::UNIX_EPOCH;
11
12use rusqlite::Connection;
13use serde::{Deserialize, Serialize};
14use serde_json::Value;
15
16use crate::native_store::load_native_store_family;
17use crate::ontology::{Binding, OrchestratorBindingRow};
18use crate::session::{
19 hermes_capture_nouns, openclaw_agent_id_from_path, openclaw_capture_header_nouns,
20 percent_decode_path, OrchestrationNouns, SessionMeta, SessionSource,
21};
22use crate::{Error, Fidelity, Result, Session, SessionFollower};
23
24pub use crate::ontology::HarnessId;
25
26#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
28#[serde(tag = "kind", rename_all = "snake_case")]
29pub enum StorageLocator {
30 File {
32 path: PathBuf,
34 },
35 Sqlite {
37 path: PathBuf,
39 selector: String,
41 },
42}
43
44impl StorageLocator {
45 pub fn path(&self) -> &Path {
47 match self {
48 Self::File { path } | Self::Sqlite { path, .. } => path,
49 }
50 }
51}
52
53#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
55pub struct SessionLocator {
56 pub harness: HarnessId,
58 pub session_id: String,
60 pub storage: StorageLocator,
62}
63
64#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
66pub struct SessionDescriptor {
67 pub locator: SessionLocator,
70 pub cwd: Option<PathBuf>,
72 pub title: Option<String>,
74 #[serde(default, skip_serializing_if = "Vec::is_empty")]
79 pub preview_candidates: Vec<SessionPreviewCandidate>,
80 #[serde(default, skip_serializing_if = "Vec::is_empty")]
85 pub latest_message_candidates: Vec<SessionPreviewCandidate>,
86 pub updated_at_ms: Option<u64>,
95 pub message_count: Option<usize>,
100 pub model: Option<String>,
102 #[serde(default, skip_serializing_if = "Option::is_none")]
106 pub parent_session_id: Option<String>,
107 #[serde(default, skip_serializing_if = "is_zero")]
111 pub child_session_count: usize,
112 #[serde(flatten)]
118 pub nouns: OrchestrationNouns,
119}
120
121#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
123pub struct SessionPreviewCandidate {
124 #[serde(default, skip_serializing_if = "Option::is_none")]
128 pub cursor: Option<String>,
129 pub role: String,
132 pub content: String,
134 #[serde(default, skip_serializing_if = "HashMap::is_empty")]
136 pub metadata: HashMap<String, String>,
137}
138
139#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
141pub struct DiscoveryPage {
142 pub sessions: Vec<SessionDescriptor>,
144 pub next_cursor: Option<String>,
146 #[serde(default)]
149 pub receipt: DiscoveryReceipt,
150}
151
152#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
154#[serde(default)]
155pub struct DiscoveryReceipt {
156 #[serde(skip_serializing_if = "is_false")]
159 pub searched_previews: bool,
160 #[serde(skip_serializing_if = "Option::is_none")]
162 pub requested_after_ms: Option<u64>,
163 #[serde(skip_serializing_if = "Option::is_none")]
165 pub requested_before_ms: Option<u64>,
166 #[serde(skip_serializing_if = "Option::is_none")]
168 pub requested_limit: Option<usize>,
169 #[serde(skip_serializing_if = "Option::is_none")]
171 pub oldest_returned_ms: Option<u64>,
172 #[serde(skip_serializing_if = "Option::is_none")]
174 pub newest_returned_ms: Option<u64>,
175 pub returned: usize,
177 pub total_matched: usize,
179 pub truncated: bool,
183}
184
185#[derive(Debug)]
192pub struct CodexHistoryTopicIndex {
193 path: PathBuf,
194 fingerprint: Option<CodexHistoryFingerprint>,
195 offset: u64,
196 trailing: Vec<u8>,
197 topics: HashMap<String, Vec<SessionPreviewCandidate>>,
198}
199
200#[derive(Debug, Clone, Copy, PartialEq, Eq)]
201struct CodexHistoryFingerprint {
202 len: u64,
203 modified_ns: u128,
204 identity: u128,
205}
206
207impl CodexHistoryTopicIndex {
208 pub fn new(sessions_root: &Path) -> Self {
211 let root = sessions_root.parent().unwrap_or(sessions_root);
212 Self {
213 path: root.join("history.jsonl"),
214 fingerprint: None,
215 offset: 0,
216 trailing: Vec::new(),
217 topics: HashMap::new(),
218 }
219 }
220
221 pub fn path(&self) -> &Path {
223 &self.path
224 }
225
226 pub fn refresh(&mut self) -> Result<BTreeSet<String>> {
232 let metadata = match fs::metadata(&self.path) {
233 Ok(metadata) => metadata,
234 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
235 let changed = self.topics.keys().cloned().collect();
236 self.fingerprint = None;
237 self.offset = 0;
238 self.trailing.clear();
239 self.topics.clear();
240 return Ok(changed);
241 }
242 Err(error) => return Err(error.into()),
243 };
244 let fingerprint = codex_history_fingerprint(&metadata)?;
245 if self.fingerprint == Some(fingerprint) {
246 return Ok(BTreeSet::new());
247 }
248
249 let is_append = self.fingerprint.is_some_and(|previous| {
250 previous.identity == fingerprint.identity
251 && previous.len < fingerprint.len
252 && self.offset <= previous.len
253 });
254 if is_append {
255 let mut file = File::open(&self.path)?;
256 file.seek(SeekFrom::Start(self.offset))?;
257 let mut bytes = Vec::with_capacity(
258 usize::try_from(fingerprint.len.saturating_sub(self.offset)).unwrap_or(0),
259 );
260 file.read_to_end(&mut bytes)?;
261 self.offset = file.stream_position()?;
262 let changed = self.ingest(bytes);
263 self.fingerprint = Some(CodexHistoryFingerprint {
264 len: self.offset,
265 ..fingerprint
266 });
267 return Ok(changed);
268 }
269
270 let previous = std::mem::take(&mut self.topics);
271 let mut file = File::open(&self.path)?;
272 let mut bytes = Vec::with_capacity(usize::try_from(fingerprint.len).unwrap_or(0));
273 file.read_to_end(&mut bytes)?;
274 self.offset = file.stream_position()?;
275 self.trailing.clear();
276 self.ingest(bytes);
277 self.fingerprint = Some(CodexHistoryFingerprint {
278 len: self.offset,
279 ..fingerprint
280 });
281 Ok(changed_topic_ids(&previous, &self.topics))
282 }
283
284 fn ingest(&mut self, bytes: Vec<u8>) -> BTreeSet<String> {
285 let mut input = std::mem::take(&mut self.trailing);
286 input.extend(bytes);
287 let complete_len = input
288 .iter()
289 .rposition(|byte| *byte == b'\n')
290 .map_or(0, |index| index + 1);
291 let mut changed = BTreeSet::new();
292 for line in input[..complete_len].split(|byte| *byte == b'\n') {
293 self.ingest_line(line, &mut changed);
294 }
295 self.trailing.extend_from_slice(&input[complete_len..]);
296 if !self.trailing.is_empty() {
297 let trailing = self.trailing.clone();
298 if self.ingest_line(&trailing, &mut changed) {
299 self.trailing.clear();
300 }
301 }
302 changed
303 }
304
305 fn ingest_line(&mut self, line: &[u8], changed: &mut BTreeSet<String>) -> bool {
306 let Ok(value) = serde_json::from_slice::<Value>(line) else {
307 return false;
308 };
309 let Some(session_id) = value.get("session_id").and_then(Value::as_str) else {
310 return true;
311 };
312 if self.topics.contains_key(session_id) {
313 return true;
314 }
315 let mut candidates = Vec::new();
316 push_message_candidate(&mut candidates, "user", value.get("text"), HashMap::new());
317 if !candidates.is_empty() {
318 self.topics.insert(session_id.to_string(), candidates);
319 changed.insert(session_id.to_string());
320 }
321 true
322 }
323}
324
325#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
327#[serde(default)]
328pub struct HarnessHomes {
329 pub claude_code: PathBuf,
331 pub codex: PathBuf,
333 pub pi: PathBuf,
335 pub opencode: PathBuf,
337 pub grok: PathBuf,
339 pub gemini: PathBuf,
341 pub goose: PathBuf,
343 pub supercode: PathBuf,
345 pub openclaw: PathBuf,
349 pub hermes: PathBuf,
352 pub orchestrator: PathBuf,
358}
359
360pub fn orchestrator_profile_dirs(root: &Path) -> Vec<(String, PathBuf)> {
366 if !root.is_dir() {
367 return Vec::new();
368 }
369 let mut dirs = vec![("default".to_string(), root.to_path_buf())];
370 if let Ok(entries) = fs::read_dir(root.join("profiles")) {
371 let mut named: Vec<(String, PathBuf)> = entries
372 .flatten()
373 .filter(|entry| entry.path().is_dir())
374 .filter_map(|entry| {
375 entry
376 .file_name()
377 .into_string()
378 .ok()
379 .map(|name| (name, entry.path()))
380 })
381 .filter(|(name, _)| !name.starts_with('.'))
385 .collect();
386 named.sort();
387 dirs.extend(named);
388 }
389 dirs
390}
391
392pub fn hermes_session_stores(db_path: &Path) -> Vec<PathBuf> {
396 let mut stores = vec![db_path.to_path_buf()];
397 let (Some(root), Some(name)) = (db_path.parent(), db_path.file_name()) else {
398 return stores;
399 };
400 for (_, dir) in orchestrator_profile_dirs(root).into_iter().skip(1) {
401 stores.push(dir.join(name));
402 }
403 stores
404}
405
406impl Default for HarnessHomes {
407 fn default() -> Self {
408 let home = std::env::var_os("HOME")
409 .map(PathBuf::from)
410 .unwrap_or_else(|| PathBuf::from("."));
411 let claude_root = std::env::var_os("CLAUDE_CONFIG_DIR")
412 .map(PathBuf::from)
413 .unwrap_or_else(|| home.join(".claude"));
414 let codex_root = std::env::var_os("CODEX_HOME")
415 .map(PathBuf::from)
416 .unwrap_or_else(|| home.join(".codex"));
417 let pi = std::env::var_os("PI_CODING_AGENT_SESSION_DIR")
418 .map(PathBuf::from)
419 .unwrap_or_else(|| {
420 std::env::var_os("PI_CODING_AGENT_DIR")
421 .map(PathBuf::from)
422 .unwrap_or_else(|| home.join(".pi/agent"))
423 .join("sessions")
424 });
425 let opencode = std::env::var_os("OPENCODE_DB")
426 .map(PathBuf::from)
427 .unwrap_or_else(|| {
428 std::env::var_os("XDG_DATA_HOME")
429 .map(PathBuf::from)
430 .unwrap_or_else(|| home.join(".local/share"))
431 .join("opencode")
432 });
433 let grok = std::env::var_os("GROK_HOME")
434 .map(PathBuf::from)
435 .unwrap_or_else(|| home.join(".grok"))
436 .join("sessions");
437 let gemini = std::env::var_os("GEMINI_CLI_HOME")
438 .map(PathBuf::from)
439 .unwrap_or_else(|| home.join(".gemini"));
440 let openclaw = std::env::var_os("OPENCLAW_STATE_DIR")
447 .map(PathBuf::from)
448 .or_else(|| {
449 std::env::var_os("OPENCLAW_HOME").map(|root| PathBuf::from(root).join(".openclaw"))
450 })
451 .unwrap_or_else(|| home.join(".openclaw"));
452 let orchestrator = std::env::var_os("SUPERCODE_ORCHESTRATOR_HOME")
453 .map(PathBuf::from)
454 .unwrap_or_else(|| home.join(".supercode/orchestrator"));
455 let hermes = std::env::var_os("HERMES_HOME")
456 .map(PathBuf::from)
457 .unwrap_or_else(|| home.join(".hermes"))
458 .join("state.db");
459 let goose = std::env::var_os("GOOSE_PATH_ROOT")
460 .map(PathBuf::from)
461 .map(|root| root.join("data/sessions/sessions.db"))
462 .unwrap_or_else(|| {
463 #[cfg(target_os = "macos")]
464 {
465 home.join("Library/Application Support/Block/goose/sessions/sessions.db")
466 }
467 #[cfg(target_os = "windows")]
468 {
469 std::env::var_os("APPDATA")
470 .map(PathBuf::from)
471 .unwrap_or_else(|| home.join("AppData/Roaming"))
472 .join("Block/goose/sessions/sessions.db")
473 }
474 #[cfg(not(any(target_os = "macos", target_os = "windows")))]
475 {
476 std::env::var_os("XDG_DATA_HOME")
477 .map(PathBuf::from)
478 .unwrap_or_else(|| home.join(".local/share"))
479 .join("goose/sessions/sessions.db")
480 }
481 });
482 let supercode = std::env::var_os("SUPERCODE_HOME")
483 .map(PathBuf::from)
484 .unwrap_or_else(|| {
485 std::env::var_os("XDG_CONFIG_HOME")
486 .map(PathBuf::from)
487 .unwrap_or_else(|| home.join(".config"))
488 .join("supercode")
489 })
490 .join("sessions");
491 Self {
492 claude_code: claude_root.join("projects"),
493 codex: codex_root.join("sessions"),
494 gemini,
495 goose,
496 supercode,
497 openclaw,
498 hermes,
499 orchestrator,
500 pi,
501 opencode,
502 grok,
503 }
504 }
505}
506
507fn codex_history_fingerprint(metadata: &fs::Metadata) -> Result<CodexHistoryFingerprint> {
508 let modified_ns = metadata
509 .modified()?
510 .duration_since(UNIX_EPOCH)
511 .map_err(|error| Error::Other(format!("history timestamp predates Unix epoch: {error}")))?
512 .as_nanos();
513 #[cfg(unix)]
514 let identity = {
515 use std::os::unix::fs::MetadataExt;
516 (u128::from(metadata.dev()) << 64) | u128::from(metadata.ino())
517 };
518 #[cfg(not(unix))]
519 let identity = 0;
520 Ok(CodexHistoryFingerprint {
521 len: metadata.len(),
522 modified_ns,
523 identity,
524 })
525}
526
527fn changed_topic_ids(
528 before: &HashMap<String, Vec<SessionPreviewCandidate>>,
529 after: &HashMap<String, Vec<SessionPreviewCandidate>>,
530) -> BTreeSet<String> {
531 before
532 .keys()
533 .chain(after.keys())
534 .filter(|session_id| before.get(*session_id) != after.get(*session_id))
535 .cloned()
536 .collect()
537}
538
539#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
541#[serde(default)]
542pub struct DiscoveryQuery {
543 pub workspace: Option<PathBuf>,
545 pub workspace_family: Option<PathBuf>,
553 pub updated_after_ms: Option<u64>,
555 pub updated_before_ms: Option<u64>,
557 pub harnesses: Vec<HarnessId>,
559 pub homes: HarnessHomes,
561 pub query: Option<String>,
563 pub search_previews: bool,
567 pub cursor: Option<String>,
569 pub limit: Option<usize>,
571 pub include_topic_candidates: bool,
575 pub include_child_sessions: bool,
579 pub root_session_id: Option<String>,
582 pub profile: Option<String>,
587}
588
589#[derive(Debug, Clone, PartialEq, Eq)]
595struct RepoFamily {
596 identity: PathBuf,
599 is_repository: bool,
602 origin_url: Option<String>,
604}
605
606impl RepoFamily {
607 fn of(path: &Path) -> Self {
608 let start = std::fs::canonicalize(path).unwrap_or_else(|_| path.to_path_buf());
609 let mut current = Some(start.as_path());
610 while let Some(dir) = current {
611 let dot_git = dir.join(".git");
612 if dot_git.is_dir() {
613 let identity = std::fs::canonicalize(&dot_git).unwrap_or_else(|_| dot_git.clone());
614 let origin_url = read_origin_url(&identity);
615 return Self {
616 identity,
617 is_repository: true,
618 origin_url,
619 };
620 }
621 if dot_git.is_file() {
622 if let Ok(text) = std::fs::read_to_string(&dot_git) {
626 if let Some(gitdir) = text
627 .lines()
628 .find_map(|line| line.trim().strip_prefix("gitdir:"))
629 {
630 let gitdir = PathBuf::from(gitdir.trim());
631 let gitdir = if gitdir.is_absolute() {
632 gitdir
633 } else {
634 dir.join(gitdir)
635 };
636 let common = gitdir
637 .parent()
638 .filter(|parent| parent.ends_with("worktrees"))
639 .and_then(Path::parent)
640 .map(Path::to_path_buf)
641 .unwrap_or(gitdir);
642 let identity = std::fs::canonicalize(&common).unwrap_or(common);
643 let origin_url = read_origin_url(&identity);
644 return Self {
645 identity,
646 is_repository: true,
647 origin_url,
648 };
649 }
650 }
651 }
652 current = dir.parent();
653 }
654 Self {
655 identity: start,
656 is_repository: false,
657 origin_url: None,
658 }
659 }
660
661 fn joins(&self, other: &Self) -> bool {
665 if self.identity == other.identity {
666 return true;
667 }
668 self.is_repository
669 && other.is_repository
670 && matches!((&self.origin_url, &other.origin_url), (Some(a), Some(b)) if a == b)
671 }
672}
673
674fn read_origin_url(common_dir: &Path) -> Option<String> {
676 let text = std::fs::read_to_string(common_dir.join("config")).ok()?;
677 let mut in_origin = false;
678 for line in text.lines() {
679 let line = line.trim();
680 if line.starts_with('[') {
681 in_origin = line.starts_with("[remote \"origin\"]");
682 continue;
683 }
684 if in_origin {
685 if let Some(value) = line.strip_prefix("url") {
686 let value = value.trim_start();
687 if let Some(url) = value.strip_prefix('=') {
688 let url = url.trim();
689 if !url.is_empty() {
690 return Some(url.to_string());
691 }
692 }
693 }
694 }
695 }
696 None
697}
698
699#[derive(Debug, Default, Clone, Copy)]
702pub struct HarnessCatalog;
703
704impl HarnessCatalog {
705 pub fn new() -> Self {
707 Self
708 }
709
710 pub fn discover(&self, query: &DiscoveryQuery) -> Result<Vec<SessionDescriptor>> {
714 Ok(self.discover_page(query)?.sessions)
715 }
716
717 pub fn discover_raw_index(&self, query: &DiscoveryQuery) -> Vec<SessionDescriptor> {
722 self.scan_descriptors(query, true)
723 }
724
725 pub fn project_index(
728 &self,
729 query: &DiscoveryQuery,
730 descriptors: impl IntoIterator<Item = SessionDescriptor>,
731 ) -> Result<Vec<SessionDescriptor>> {
732 Ok(self.project_index_page(query, descriptors)?.sessions)
733 }
734
735 pub fn project_index_page(
738 &self,
739 query: &DiscoveryQuery,
740 descriptors: impl IntoIterator<Item = SessionDescriptor>,
741 ) -> Result<DiscoveryPage> {
742 if query.search_previews {
743 return Err(Error::Other(
744 "preview search requires discover_page, not a metadata-only index projection"
745 .into(),
746 ));
747 }
748 let mut found = descriptors.into_iter().collect::<Vec<_>>();
749 project_descriptors(query, &mut found);
750 let total_matched = found.len();
751 let (sessions, next_cursor) = paginate_descriptors(query, found)?;
752 let receipt = DiscoveryReceipt {
753 searched_previews: false,
754 requested_after_ms: query.updated_after_ms,
755 requested_before_ms: query.updated_before_ms,
756 requested_limit: query.limit,
757 oldest_returned_ms: sessions.iter().filter_map(|s| s.updated_at_ms).min(),
758 newest_returned_ms: sessions.iter().filter_map(|s| s.updated_at_ms).max(),
759 returned: sessions.len(),
760 total_matched,
761 truncated: next_cursor.is_some(),
762 };
763 Ok(DiscoveryPage {
764 sessions,
765 next_cursor,
766 receipt,
767 })
768 }
769
770 pub fn enrich_index_page(
773 &self,
774 query: &DiscoveryQuery,
775 mut sessions: Vec<SessionDescriptor>,
776 ) -> Result<Vec<SessionDescriptor>> {
777 enrich_descriptors(query, &mut sessions, None)?;
778 Ok(sessions)
779 }
780
781 pub fn enrich_index_page_with_codex_history(
787 &self,
788 query: &DiscoveryQuery,
789 mut sessions: Vec<SessionDescriptor>,
790 codex_history: &CodexHistoryTopicIndex,
791 ) -> Result<Vec<SessionDescriptor>> {
792 enrich_descriptors(query, &mut sessions, Some(codex_history))?;
793 Ok(sessions)
794 }
795
796 pub fn discover_page(&self, query: &DiscoveryQuery) -> Result<DiscoveryPage> {
798 if query.search_previews {
799 return self.discover_preview_page(query);
800 }
801 let mut page = self.project_index_page(
802 query,
803 self.scan_descriptors(query, query.include_child_sessions),
804 )?;
805 if page.sessions.is_empty() && query.cursor.is_none() {
806 if let Some(named) = self.claude_session_named(query)? {
807 page.sessions.push(named);
808 }
809 }
810 enrich_descriptors(query, &mut page.sessions, None)?;
811 Ok(page)
812 }
813
814 fn claude_session_named(&self, query: &DiscoveryQuery) -> Result<Option<SessionDescriptor>> {
821 let Some(name) = query
822 .query
823 .as_deref()
824 .map(str::trim)
825 .filter(|name| !name.is_empty())
826 else {
827 return Ok(None);
828 };
829 if !query.harnesses.is_empty()
830 && !query
831 .harnesses
832 .iter()
833 .any(|harness| harness.as_str() == HarnessId::CLAUDE_CODE)
834 {
835 return Ok(None);
836 }
837 let needle = format!("\"customTitle\":{}", Value::String(name.to_string()));
838 let mut best: Option<(std::time::SystemTime, String)> = None;
839 let Ok(projects) = std::fs::read_dir(&query.homes.claude_code) else {
840 return Ok(None);
841 };
842 for project in projects.flatten() {
843 let Ok(files) = std::fs::read_dir(project.path()) else {
844 continue;
845 };
846 for file in files.flatten() {
847 let path = file.path();
848 if path.extension().and_then(|value| value.to_str()) != Some("jsonl") {
849 continue;
850 }
851 let Ok(bytes) = std::fs::read(&path) else {
852 continue;
853 };
854 if !bytes
855 .windows(needle.len())
856 .any(|window| window == needle.as_bytes())
857 {
858 continue;
859 }
860 let modified = file
861 .metadata()
862 .and_then(|meta| meta.modified())
863 .unwrap_or(std::time::UNIX_EPOCH);
864 let Some(id) = path
865 .file_stem()
866 .map(|stem| stem.to_string_lossy().into_owned())
867 else {
868 continue;
869 };
870 if best.as_ref().is_none_or(|(time, _)| modified > *time) {
871 best = Some((modified, id));
872 }
873 }
874 }
875 let Some((_, id)) = best else {
876 return Ok(None);
877 };
878 let mut unnamed = query.clone();
879 unnamed.query = None;
880 unnamed.harnesses = vec![HarnessId::new(HarnessId::CLAUDE_CODE)];
881 unnamed.limit = None;
882 let found = self
883 .project_index_page(
884 &unnamed,
885 self.scan_descriptors(&unnamed, unnamed.include_child_sessions),
886 )?
887 .sessions
888 .into_iter()
889 .find(|descriptor| descriptor.locator.session_id == id);
890 Ok(found.map(|mut descriptor| {
891 descriptor.title = Some(name.to_string());
892 descriptor
893 }))
894 }
895
896 fn discover_preview_page(&self, query: &DiscoveryQuery) -> Result<DiscoveryPage> {
897 let search = query
898 .query
899 .as_deref()
900 .map(str::trim)
901 .filter(|text| !text.is_empty())
902 .ok_or_else(|| Error::Other("search_previews requires a nonempty query".into()))?
903 .to_lowercase();
904 if query.limit == Some(0) {
905 return Err(Error::Other("preview search limit must be positive".into()));
906 }
907 let cursor = query.cursor.as_deref().map(decode_cursor).transpose()?;
908 let mut eligible_query = query.clone();
909 eligible_query.query = None;
910 eligible_query.search_previews = false;
911 eligible_query.include_topic_candidates = true;
912 let mut eligible = self.scan_descriptors(query, query.include_child_sessions);
913 project_descriptors(&eligible_query, &mut eligible);
914 let topics = codex_history_topics(&query.homes.codex, &eligible).unwrap_or_default();
916 let mut sessions = Vec::new();
917 let mut total_matched = 0;
918 let mut cursor_seen = cursor.is_none();
919 let mut more = false;
920 for mut descriptor in eligible {
921 let metadata_match = descriptor_matches(&descriptor, &search);
922 if !metadata_match {
923 let topic = topics
924 .get(&descriptor.locator.session_id)
925 .filter(|_| descriptor.locator.harness.as_str() == HarnessId::CODEX);
926 enrich_descriptor(&eligible_query, &mut descriptor, topic);
927 if !descriptor
928 .preview_candidates
929 .iter()
930 .chain(&descriptor.latest_message_candidates)
931 .any(|candidate| candidate.content.to_lowercase().contains(&search))
932 {
933 continue;
934 }
935 }
936 total_matched += 1;
937 if !cursor_seen {
938 cursor_seen = cursor.as_ref() == Some(&descriptor_cursor_key(&descriptor));
939 continue;
940 }
941 if sessions.len() < query.limit.unwrap_or(usize::MAX) {
942 if metadata_match {
943 let topic = topics
944 .get(&descriptor.locator.session_id)
945 .filter(|_| descriptor.locator.harness.as_str() == HarnessId::CODEX);
946 enrich_descriptor(&eligible_query, &mut descriptor, topic);
947 }
948 sessions.push(descriptor);
949 } else {
950 more = true;
951 }
952 }
953 if !cursor_seen {
954 return Err(Error::Other("discovery cursor is stale or invalid".into()));
955 }
956 let next_cursor = more.then(|| sessions.last().map(encode_cursor)).flatten();
957 let receipt = DiscoveryReceipt {
958 searched_previews: true,
959 requested_after_ms: query.updated_after_ms,
960 requested_before_ms: query.updated_before_ms,
961 requested_limit: query.limit,
962 oldest_returned_ms: sessions.iter().filter_map(|s| s.updated_at_ms).min(),
963 newest_returned_ms: sessions.iter().filter_map(|s| s.updated_at_ms).max(),
964 returned: sessions.len(),
965 total_matched,
966 truncated: more,
967 };
968 Ok(DiscoveryPage {
969 sessions,
970 next_cursor,
971 receipt,
972 })
973 }
974
975 fn scan_descriptors(
976 &self,
977 query: &DiscoveryQuery,
978 include_child_sessions: bool,
979 ) -> Vec<SessionDescriptor> {
980 let selected: HashSet<&str> = if query.harnesses.is_empty() {
981 [
982 HarnessId::CLAUDE_CODE,
983 HarnessId::CODEX,
984 HarnessId::PI,
985 HarnessId::OPENCODE,
986 HarnessId::GROK,
987 HarnessId::GEMINI,
988 HarnessId::GOOSE,
989 HarnessId::SUPERCODE,
990 HarnessId::OPENCLAW,
991 HarnessId::HERMES,
992 HarnessId::ORCHESTRATOR,
993 ]
994 .into_iter()
995 .collect()
996 } else {
997 query.harnesses.iter().map(HarnessId::as_str).collect()
998 };
999 let mut found = Vec::new();
1000 if selected.contains(HarnessId::CLAUDE_CODE) {
1001 discover_jsonl(
1002 &query.homes.claude_code,
1003 HarnessId::CLAUDE_CODE,
1004 query.workspace.as_deref(),
1005 include_child_sessions,
1006 &mut found,
1007 );
1008 }
1009 if selected.contains(HarnessId::CODEX) {
1010 discover_jsonl(
1011 &query.homes.codex,
1012 HarnessId::CODEX,
1013 query.workspace.as_deref(),
1014 include_child_sessions,
1015 &mut found,
1016 );
1017 }
1018 if selected.contains(HarnessId::PI) {
1019 discover_jsonl(
1020 &query.homes.pi,
1021 HarnessId::PI,
1022 query.workspace.as_deref(),
1023 include_child_sessions,
1024 &mut found,
1025 );
1026 }
1027 if selected.contains(HarnessId::OPENCODE) {
1028 discover_opencode(
1029 &query.homes.opencode,
1030 query.workspace.as_deref(),
1031 &mut found,
1032 );
1033 }
1034 if selected.contains(HarnessId::GROK) {
1035 discover_grok(&query.homes.grok, query.workspace.as_deref(), &mut found);
1036 }
1037 if selected.contains(HarnessId::GEMINI) {
1038 discover_gemini(&query.homes.gemini, query.workspace.as_deref(), &mut found);
1039 }
1040 if selected.contains(HarnessId::GOOSE) {
1041 discover_goose(&query.homes.goose, query.workspace.as_deref(), &mut found);
1042 }
1043 if selected.contains(HarnessId::OPENCLAW) {
1044 discover_openclaw(
1045 &query.homes.openclaw,
1046 query.workspace.as_deref(),
1047 &mut found,
1048 );
1049 }
1050 if selected.contains(HarnessId::HERMES) {
1051 for store in hermes_session_stores(&query.homes.hermes) {
1052 discover_hermes(&store, query.workspace.as_deref(), &selected, &mut found);
1053 }
1054 }
1055 if selected.contains(HarnessId::ORCHESTRATOR) {
1056 discover_orchestrator(
1057 &query.homes.orchestrator,
1058 query.workspace.as_deref(),
1059 &mut found,
1060 );
1061 }
1062 if selected.contains(HarnessId::SUPERCODE) {
1063 discover_supercode(
1064 &query.homes.supercode,
1065 query.workspace.as_deref(),
1066 &mut found,
1067 );
1068 }
1069 for descriptor in &mut found {
1070 finalize_nouns(descriptor);
1071 }
1072 found
1073 }
1074
1075 pub fn refresh_file_descriptor(
1083 &self,
1084 locator: &SessionLocator,
1085 workspace: Option<&Path>,
1086 include_topic_candidates: bool,
1087 ) -> Result<Option<SessionDescriptor>> {
1088 let Some(mut descriptor) = self.refresh_file_index_descriptor(locator, workspace)? else {
1089 return Ok(None);
1090 };
1091 if include_topic_candidates {
1092 descriptor.preview_candidates =
1093 topic_message_candidates(&descriptor.locator).unwrap_or_default();
1094 }
1095 descriptor.latest_message_candidates =
1096 latest_message_candidates(&descriptor.locator).unwrap_or_default();
1097 Ok(Some(descriptor))
1098 }
1099
1100 pub fn refresh_file_index_descriptor(
1104 &self,
1105 locator: &SessionLocator,
1106 workspace: Option<&Path>,
1107 ) -> Result<Option<SessionDescriptor>> {
1108 let StorageLocator::File { path } = &locator.storage else {
1109 return Ok(None);
1110 };
1111 if !matches!(
1112 locator.harness.as_str(),
1113 HarnessId::CLAUDE_CODE | HarnessId::CODEX
1114 ) {
1115 return Ok(None);
1116 }
1117 if !path.is_file() {
1118 return Ok(None);
1119 }
1120 let Ok(meta) = read_header(path, locator.harness.as_str()) else {
1121 return Ok(None);
1125 };
1126 if workspace.is_some_and(|wanted| {
1127 meta.cwd
1128 .as_deref()
1129 .is_none_or(|cwd| !recorded_cwd_matches(cwd, wanted))
1130 }) {
1131 return Ok(None);
1132 }
1133 let parent_session_id = meta.parent_session_id.or_else(|| {
1134 (locator.harness.as_str() == HarnessId::CLAUDE_CODE)
1135 .then(|| claude_subagent_parent_id(path))
1136 .flatten()
1137 });
1138 let tail = tail_facts(path, locator.harness.as_str());
1139 let descriptor = SessionDescriptor {
1140 locator: SessionLocator {
1141 harness: locator.harness.clone(),
1142 session_id: meta
1143 .session_id
1144 .unwrap_or_else(|| locator.session_id.clone()),
1145 storage: StorageLocator::File { path: path.clone() },
1146 },
1147 cwd: meta.cwd,
1148 title: meta.title,
1149 preview_candidates: Vec::new(),
1150 latest_message_candidates: Vec::new(),
1151 updated_at_ms: tail.last_turn_ms.or_else(|| modified_ms(path)),
1152 message_count: None,
1153 model: tail.model.or(meta.model),
1154 parent_session_id,
1155 child_session_count: 0,
1156 nouns: OrchestrationNouns::default(),
1157 };
1158 Ok(Some(descriptor))
1159 }
1160
1161 pub fn load(&self, locator: &SessionLocator) -> Result<Session> {
1163 self.load_with_fidelity(locator, Fidelity::ByteLossless)
1164 }
1165
1166 pub fn load_with_fidelity(
1173 &self,
1174 locator: &SessionLocator,
1175 fidelity: Fidelity,
1176 ) -> Result<Session> {
1177 if let Some(session) = load_hermes_locator(locator) {
1178 return session;
1179 }
1180 match &locator.storage {
1181 StorageLocator::File { path } => {
1182 if let Some(session) = load_native_store_family(path)? {
1183 Ok(session)
1184 } else {
1185 Ok(Session::load_with_fidelity(path, fidelity)?)
1186 }
1187 }
1188 StorageLocator::Sqlite { path, selector } => {
1189 if locator.harness.as_str() == HarnessId::GOOSE {
1190 Ok(Session::from_goose_sqlite(path, selector)?)
1191 } else {
1192 Ok(Session::from_opencode_sqlite(path, Some(selector))?)
1193 }
1194 }
1195 }
1196 }
1197
1198 #[doc(hidden)]
1202 pub fn load_parent_with_fidelity(
1203 &self,
1204 locator: &SessionLocator,
1205 fidelity: Fidelity,
1206 ) -> Result<Session> {
1207 if let Some(session) = load_hermes_locator(locator) {
1208 return session;
1209 }
1210 match &locator.storage {
1211 StorageLocator::File { path } => {
1212 if let Some(session) = load_native_store_family(path)? {
1213 Ok(session)
1214 } else {
1215 Ok(Session::load_parent_with_fidelity(path, fidelity)?)
1216 }
1217 }
1218 StorageLocator::Sqlite { path, selector } => {
1219 if locator.harness.as_str() == HarnessId::GOOSE {
1220 Ok(Session::from_goose_sqlite(path, selector)?)
1221 } else {
1222 Ok(Session::from_opencode_sqlite(path, Some(selector))?)
1223 }
1224 }
1225 }
1226 }
1227
1228 #[doc(hidden)]
1231 pub fn load_display_view(
1232 &self,
1233 locator: &SessionLocator,
1234 fidelity: Fidelity,
1235 message_limit: usize,
1236 ) -> Result<Session> {
1237 if let Some(session) = load_hermes_locator(locator) {
1238 let mut session = session?;
1239 if session.messages.len() > message_limit.max(1) {
1240 session
1241 .messages
1242 .drain(..session.messages.len() - message_limit.max(1));
1243 }
1244 return Ok(session);
1245 }
1246 match &locator.storage {
1247 StorageLocator::File { path } => {
1248 if let Some(mut session) = load_native_store_family(path)? {
1249 if session.messages.len() > message_limit.max(1) {
1250 session
1251 .messages
1252 .drain(..session.messages.len() - message_limit.max(1));
1253 }
1254 Ok(session)
1255 } else {
1256 Ok(Session::load_display_view(path, fidelity, message_limit)?)
1257 }
1258 }
1259 StorageLocator::Sqlite { path, selector } => {
1260 let mut session = if locator.harness.as_str() == HarnessId::GOOSE {
1261 Session::from_goose_sqlite_display(path, selector, message_limit)?
1262 } else {
1263 Session::from_opencode_sqlite(path, Some(selector))?
1264 };
1265 if session.messages.len() > message_limit.max(1) {
1266 session
1267 .messages
1268 .drain(..session.messages.len() - message_limit.max(1));
1269 }
1270 Ok(session)
1271 }
1272 }
1273 }
1274
1275 pub fn follow(&self, locator: &SessionLocator) -> Result<SessionFollower> {
1277 self.follow_with_fidelity(locator, Fidelity::ByteLossless)
1278 }
1279
1280 pub fn follow_with_fidelity(
1282 &self,
1283 locator: &SessionLocator,
1284 fidelity: Fidelity,
1285 ) -> Result<SessionFollower> {
1286 SessionFollower::open_locator_with_fidelity(locator, fidelity)
1287 }
1288
1289 #[doc(hidden)]
1291 pub fn follow_read_view(
1292 &self,
1293 locator: &SessionLocator,
1294 fidelity: Fidelity,
1295 include_subagents: bool,
1296 message_limit: Option<usize>,
1297 max_message_chars: Option<usize>,
1298 display_history: bool,
1299 ) -> Result<SessionFollower> {
1300 SessionFollower::open_locator_with_view(
1301 locator,
1302 fidelity,
1303 include_subagents,
1304 message_limit,
1305 max_message_chars,
1306 display_history,
1307 )
1308 }
1309}
1310
1311fn load_hermes_locator(locator: &SessionLocator) -> Option<Result<Session>> {
1317 if locator.harness.as_str() != HarnessId::HERMES {
1318 return None;
1319 }
1320 let StorageLocator::File { path } = &locator.storage else {
1321 return None;
1322 };
1323 Some(Session::from_hermes_sqlite(path, Some(&locator.session_id)))
1324}
1325
1326fn finalize_nouns(descriptor: &mut SessionDescriptor) {
1333 let mut meta = SessionMeta::new(SessionSource::Native);
1334 meta.cwd = descriptor.cwd.clone();
1335 meta.trigger = descriptor.nouns.trigger;
1336 meta.surface = descriptor.nouns.surface.clone();
1337 meta.profile = descriptor.nouns.profile.clone();
1338 meta.recurrence = descriptor.nouns.recurrence.clone();
1339 meta.cross_surface = descriptor.nouns.cross_surface.clone();
1340 descriptor.nouns = OrchestrationNouns::from_meta(&meta);
1341}
1342
1343fn project_descriptors(query: &DiscoveryQuery, found: &mut Vec<SessionDescriptor>) {
1344 roll_up_session_children(found, query.include_child_sessions);
1345 if let Some(root_session_id) = query.root_session_id.as_deref() {
1346 retain_session_family(found, root_session_id);
1347 }
1348 if let Some(family_path) = query.workspace_family.as_deref() {
1349 let family = RepoFamily::of(family_path);
1350 let mut cache: HashMap<PathBuf, bool> = HashMap::new();
1351 found.retain(|descriptor| {
1352 let Some(cwd) = descriptor.cwd.as_deref() else {
1353 return false;
1354 };
1355 *cache
1356 .entry(cwd.to_path_buf())
1357 .or_insert_with(|| RepoFamily::of(cwd).joins(&family))
1358 });
1359 }
1360 if let Some(profile) = query
1361 .profile
1362 .as_deref()
1363 .map(str::trim)
1364 .filter(|p| !p.is_empty())
1365 {
1366 found.retain(|descriptor| descriptor.nouns.profile.as_deref() == Some(profile));
1367 }
1368 if let Some(after) = query.updated_after_ms {
1369 found.retain(|descriptor| descriptor.updated_at_ms.is_some_and(|at| at >= after));
1370 }
1371 if let Some(before) = query.updated_before_ms {
1372 found.retain(|descriptor| descriptor.updated_at_ms.is_some_and(|at| at <= before));
1373 }
1374 found.sort_by(|a, b| {
1375 b.updated_at_ms
1376 .cmp(&a.updated_at_ms)
1377 .then_with(|| a.locator.harness.cmp(&b.locator.harness))
1378 .then_with(|| a.locator.session_id.cmp(&b.locator.session_id))
1379 });
1380 if let Some(search) = query
1381 .query
1382 .as_deref()
1383 .map(str::trim)
1384 .filter(|query| !query.is_empty())
1385 {
1386 let search = search.to_lowercase();
1387 found.retain(|descriptor| descriptor_matches(descriptor, &search));
1388 }
1389}
1390
1391fn paginate_descriptors(
1392 query: &DiscoveryQuery,
1393 found: Vec<SessionDescriptor>,
1394) -> Result<(Vec<SessionDescriptor>, Option<String>)> {
1395 let start = match query.cursor.as_deref() {
1396 Some(cursor) => {
1397 let key = decode_cursor(cursor)?;
1398 found
1399 .iter()
1400 .position(|descriptor| descriptor_cursor_key(descriptor) == key)
1401 .map(|index| index + 1)
1402 .ok_or_else(|| Error::Other("discovery cursor is stale or invalid".into()))?
1403 }
1404 None => 0,
1405 };
1406 let end = query
1407 .limit
1408 .map(|limit| start.saturating_add(limit).min(found.len()))
1409 .unwrap_or(found.len());
1410 let sessions = found[start.min(found.len())..end].to_vec();
1411 let next_cursor = (end < found.len())
1412 .then(|| sessions.last().map(encode_cursor))
1413 .flatten();
1414 Ok((sessions, next_cursor))
1415}
1416
1417fn enrich_descriptors(
1418 query: &DiscoveryQuery,
1419 sessions: &mut [SessionDescriptor],
1420 codex_history: Option<&CodexHistoryTopicIndex>,
1421) -> Result<()> {
1422 let codex_topics = if query.include_topic_candidates && codex_history.is_none() {
1423 codex_history_topics(&query.homes.codex, sessions).unwrap_or_default()
1424 } else {
1425 HashMap::new()
1426 };
1427 for descriptor in sessions {
1428 let topic = (descriptor.locator.harness.as_str() == HarnessId::CODEX)
1429 .then(|| {
1430 codex_history
1431 .and_then(|history| history.topics.get(&descriptor.locator.session_id))
1432 .or_else(|| codex_topics.get(&descriptor.locator.session_id))
1433 })
1434 .flatten();
1435 enrich_descriptor(query, descriptor, topic);
1436 }
1437 Ok(())
1438}
1439
1440fn enrich_descriptor(
1441 query: &DiscoveryQuery,
1442 descriptor: &mut SessionDescriptor,
1443 codex_topic: Option<&Vec<SessionPreviewCandidate>>,
1444) {
1445 if query.include_topic_candidates {
1446 descriptor.preview_candidates = codex_topic
1447 .cloned()
1448 .unwrap_or_else(|| topic_message_candidates(&descriptor.locator).unwrap_or_default());
1449 }
1450 descriptor.latest_message_candidates =
1451 latest_message_candidates(&descriptor.locator).unwrap_or_default();
1452}
1453
1454fn is_false(value: &bool) -> bool {
1455 !value
1456}
1457
1458fn descriptor_matches(descriptor: &SessionDescriptor, search: &str) -> bool {
1459 [
1460 Some(descriptor.locator.harness.as_str()),
1461 Some(descriptor.locator.session_id.as_str()),
1462 descriptor.title.as_deref(),
1463 descriptor.cwd.as_ref().and_then(|path| path.to_str()),
1464 descriptor.model.as_deref(),
1465 ]
1466 .into_iter()
1467 .flatten()
1468 .any(|value| value.to_lowercase().contains(search))
1469}
1470
1471fn descriptor_cursor_key(descriptor: &SessionDescriptor) -> (Option<u64>, String, String) {
1472 (
1473 descriptor.updated_at_ms,
1474 descriptor.locator.harness.as_str().to_string(),
1475 descriptor.locator.session_id.clone(),
1476 )
1477}
1478
1479fn encode_cursor(descriptor: &SessionDescriptor) -> String {
1480 let json = serde_json::to_vec(&descriptor_cursor_key(descriptor)).unwrap_or_default();
1481 let mut encoded = String::with_capacity(json.len() * 2);
1482 for byte in json {
1483 use std::fmt::Write;
1484 let _ = write!(&mut encoded, "{byte:02x}");
1485 }
1486 encoded
1487}
1488
1489fn decode_cursor(cursor: &str) -> Result<(Option<u64>, String, String)> {
1490 if cursor.len() % 2 != 0 {
1491 return Err(Error::Other("discovery cursor is invalid".into()));
1492 }
1493 let bytes = (0..cursor.len())
1494 .step_by(2)
1495 .map(|index| u8::from_str_radix(&cursor[index..index + 2], 16))
1496 .collect::<std::result::Result<Vec<_>, _>>()
1497 .map_err(|_| Error::Other("discovery cursor is invalid".into()))?;
1498 serde_json::from_slice(&bytes).map_err(|_| Error::Other("discovery cursor is invalid".into()))
1499}
1500
1501#[derive(Default)]
1502struct HeaderMeta {
1503 session_id: Option<String>,
1504 cwd: Option<PathBuf>,
1505 title: Option<String>,
1506 model: Option<String>,
1507 parent_session_id: Option<String>,
1508}
1509
1510fn discover_jsonl(
1511 root: &Path,
1512 harness: &str,
1513 workspace: Option<&Path>,
1514 include_child_sessions: bool,
1515 found: &mut Vec<SessionDescriptor>,
1516) {
1517 let mut files = Vec::new();
1518 collect_jsonl(root, harness, include_child_sessions, &mut files);
1519 for path in files {
1520 let Ok(meta) = read_header(&path, harness) else {
1521 continue;
1522 };
1523 if workspace.is_some_and(|wanted| {
1524 meta.cwd
1525 .as_deref()
1526 .is_none_or(|cwd| !recorded_cwd_matches(cwd, wanted))
1527 }) {
1528 continue;
1529 }
1530 let session_id = meta.session_id.unwrap_or_else(|| {
1531 path.file_stem()
1532 .and_then(|value| value.to_str())
1533 .unwrap_or("unknown")
1534 .to_string()
1535 });
1536 let parent_session_id = meta.parent_session_id.or_else(|| {
1537 (harness == HarnessId::CLAUDE_CODE)
1538 .then(|| claude_subagent_parent_id(&path))
1539 .flatten()
1540 });
1541 let tail = tail_facts(&path, harness);
1542 found.push(SessionDescriptor {
1543 locator: SessionLocator {
1544 harness: HarnessId::new(harness),
1545 session_id,
1546 storage: StorageLocator::File { path: path.clone() },
1547 },
1548 cwd: meta.cwd,
1549 title: meta.title,
1550 preview_candidates: Vec::new(),
1551 latest_message_candidates: Vec::new(),
1552 updated_at_ms: tail.last_turn_ms.or_else(|| modified_ms(&path)),
1553 message_count: None,
1554 model: tail.model.or(meta.model),
1555 parent_session_id,
1556 child_session_count: if harness == HarnessId::CLAUDE_CODE && !include_child_sessions {
1557 count_claude_subagents(&path)
1558 } else {
1559 0
1560 },
1561 nouns: OrchestrationNouns::default(),
1562 });
1563 }
1564}
1565
1566fn collect_jsonl(root: &Path, harness: &str, include_child_sessions: bool, out: &mut Vec<PathBuf>) {
1567 let mut walked = HashSet::new();
1568 collect_jsonl_in(root, harness, include_child_sessions, out, &mut walked);
1569}
1570
1571fn collect_jsonl_in(
1575 root: &Path,
1576 harness: &str,
1577 include_child_sessions: bool,
1578 out: &mut Vec<PathBuf>,
1579 walked: &mut HashSet<PathBuf>,
1580) {
1581 if !walked.insert(fs::canonicalize(root).unwrap_or_else(|_| root.to_path_buf())) {
1582 return;
1583 }
1584 let Ok(entries) = fs::read_dir(root) else {
1585 return;
1586 };
1587 for entry in entries.flatten() {
1588 let Ok(mut kind) = entry.file_type() else {
1589 continue;
1590 };
1591 let path = entry.path();
1592 if kind.is_symlink() {
1593 let Ok(target) = fs::metadata(&path) else {
1594 continue;
1595 };
1596 kind = target.file_type();
1597 }
1598 if kind.is_dir() {
1599 if harness == HarnessId::CLAUDE_CODE
1600 && path.file_name().and_then(|v| v.to_str()) == Some("subagents")
1601 && !include_child_sessions
1602 {
1603 continue;
1604 }
1605 collect_jsonl_in(&path, harness, include_child_sessions, out, walked);
1606 } else if kind.is_file()
1607 && (path.extension().and_then(|v| v.to_str()) == Some("jsonl")
1608 || (harness == HarnessId::CODEX
1609 && path
1610 .file_name()
1611 .and_then(|v| v.to_str())
1612 .is_some_and(|name| name.ends_with(".jsonl.zst"))))
1613 {
1614 out.push(path);
1615 }
1616 }
1617}
1618
1619fn read_header(path: &Path, harness: &str) -> Result<HeaderMeta> {
1620 let file = crate::session::open_session_reader(path)?;
1621 let mut result = HeaderMeta::default();
1622 let mut bytes = 0usize;
1623 for line in BufReader::new(file).lines().take(32) {
1624 let line = line?;
1625 bytes += line.len();
1626 if bytes > 256 * 1024 {
1627 break;
1628 }
1629 let Ok(value) = serde_json::from_str::<Value>(&line) else {
1630 continue;
1631 };
1632 update_header_meta(&mut result, &value, harness);
1633 if result.session_id.is_some() && result.cwd.is_some() && result.model.is_some() {
1634 break;
1635 }
1636 }
1637 if result.session_id.is_none() && result.cwd.is_none() {
1638 return Err(Error::Other(format!(
1639 "{} has no recognizable {harness} session header",
1640 path.display()
1641 )));
1642 }
1643 Ok(result)
1644}
1645
1646fn update_header_meta(result: &mut HeaderMeta, value: &Value, harness: &str) {
1647 match harness {
1648 HarnessId::CLAUDE_CODE => {
1649 fill_string(&mut result.session_id, value.get("sessionId"));
1650 fill_path(&mut result.cwd, value.get("cwd"));
1651 fill_string(&mut result.title, value.get("customTitle"));
1653 fill_string(&mut result.title, value.get("agentName"));
1654 fill_string(
1655 &mut result.model,
1656 value.get("message").and_then(|v| v.get("model")),
1657 );
1658 }
1659 HarnessId::CODEX => {
1660 let payload = value.get("payload").unwrap_or(&Value::Null);
1661 if value.get("type").and_then(Value::as_str) == Some("session_meta") {
1662 fill_string(&mut result.session_id, payload.get("id"));
1663 fill_path(&mut result.cwd, payload.get("cwd"));
1664 fill_string(&mut result.title, payload.get("thread_name"));
1665 fill_string(&mut result.title, payload.get("title"));
1666 fill_string(
1667 &mut result.parent_session_id,
1668 payload.get("parent_thread_id"),
1669 );
1670 if let Some(parent) = payload
1671 .pointer("/source/subagent/thread_spawn/parent_thread_id")
1672 .and_then(Value::as_str)
1673 {
1674 result.parent_session_id = Some(parent.to_string());
1675 }
1676 if result.title.is_none() {
1677 result.title = payload
1678 .pointer("/source/subagent/thread_spawn/agent_path")
1679 .and_then(Value::as_str)
1680 .and_then(|path| path.rsplit('/').find(|part| !part.is_empty()))
1681 .map(humanize_topic);
1682 }
1683 }
1684 if value.get("type").and_then(Value::as_str) == Some("turn_context") {
1685 fill_path(&mut result.cwd, payload.get("cwd"));
1686 fill_string(&mut result.model, payload.get("model"));
1687 }
1688 }
1689 HarnessId::PI => {
1690 if value.get("type").and_then(Value::as_str) == Some("session") {
1691 fill_string(&mut result.session_id, value.get("id"));
1692 fill_path(&mut result.cwd, value.get("cwd"));
1693 }
1694 fill_string(
1695 &mut result.model,
1696 value.get("message").and_then(|v| v.get("model")),
1697 );
1698 }
1699 _ => {}
1700 }
1701}
1702
1703fn roll_up_session_children(found: &mut Vec<SessionDescriptor>, include_children: bool) {
1707 let by_id = found
1708 .iter()
1709 .enumerate()
1710 .map(|(index, descriptor)| {
1711 (
1712 (
1713 descriptor.locator.harness.as_str().to_string(),
1714 descriptor.locator.session_id.clone(),
1715 ),
1716 index,
1717 )
1718 })
1719 .collect::<HashMap<_, _>>();
1720 let mut root_updates = HashMap::<usize, u64>::new();
1721 let mut root_child_counts = HashMap::<usize, usize>::new();
1722
1723 for descriptor in found.iter() {
1724 let Some(mut parent_id) = descriptor.parent_session_id.as_deref() else {
1725 continue;
1726 };
1727 let harness = descriptor.locator.harness.as_str();
1728 let mut root = None;
1729 let mut visited = HashSet::new();
1730 while visited.insert(parent_id.to_string()) {
1731 let Some(&parent_index) = by_id.get(&(harness.to_string(), parent_id.to_string()))
1732 else {
1733 break;
1734 };
1735 root = Some(parent_index);
1736 let Some(next_parent) = found[parent_index].parent_session_id.as_deref() else {
1737 break;
1738 };
1739 parent_id = next_parent;
1740 }
1741 if let (Some(root), Some(updated_at_ms)) = (root, descriptor.updated_at_ms) {
1742 root_updates
1743 .entry(root)
1744 .and_modify(|current| *current = (*current).max(updated_at_ms))
1745 .or_insert(updated_at_ms);
1746 }
1747 if let Some(root) = root {
1748 *root_child_counts.entry(root).or_default() += 1;
1749 }
1750 }
1751
1752 for (root, child_updated_at_ms) in root_updates {
1753 found[root].updated_at_ms = Some(
1754 found[root]
1755 .updated_at_ms
1756 .unwrap_or_default()
1757 .max(child_updated_at_ms),
1758 );
1759 }
1760 for (root, child_count) in root_child_counts {
1761 found[root].child_session_count = child_count;
1762 }
1763 if !include_children {
1764 found.retain(|descriptor| descriptor.parent_session_id.is_none());
1765 }
1766}
1767
1768fn retain_session_family(found: &mut Vec<SessionDescriptor>, root_session_id: &str) {
1769 let parent_by_id = found
1770 .iter()
1771 .map(|descriptor| {
1772 (
1773 descriptor.locator.session_id.clone(),
1774 descriptor.parent_session_id.clone(),
1775 )
1776 })
1777 .collect::<HashMap<_, _>>();
1778 found.retain(|descriptor| {
1779 let mut current = descriptor.locator.session_id.clone();
1780 let mut visited = HashSet::new();
1781 while visited.insert(current.clone()) {
1782 if current == root_session_id {
1783 return true;
1784 }
1785 let Some(Some(parent)) = parent_by_id.get(¤t) else {
1786 return false;
1787 };
1788 current = parent.clone();
1789 }
1790 false
1791 });
1792}
1793
1794fn claude_subagent_parent_id(path: &Path) -> Option<String> {
1795 let subagents = path.parent()?;
1796 if subagents.file_name()?.to_str()? != "subagents" {
1797 return None;
1798 }
1799 subagents
1800 .parent()?
1801 .file_name()?
1802 .to_str()
1803 .map(str::to_string)
1804}
1805
1806fn count_claude_subagents(parent_path: &Path) -> usize {
1807 let Some(parent) = parent_path.parent() else {
1808 return 0;
1809 };
1810 let Some(stem) = parent_path.file_stem() else {
1811 return 0;
1812 };
1813 let root = parent.join(stem).join("subagents");
1814 let mut files = Vec::new();
1815 collect_jsonl(&root, HarnessId::CLAUDE_CODE, true, &mut files);
1816 files.len()
1817}
1818
1819fn is_zero(value: &usize) -> bool {
1820 *value == 0
1821}
1822
1823fn humanize_topic(value: &str) -> String {
1824 let text = value.replace(['_', '-'], " ");
1825 let mut characters = text.chars();
1826 match characters.next() {
1827 Some(first) => first.to_uppercase().collect::<String>() + characters.as_str(),
1828 None => text,
1829 }
1830}
1831
1832fn discover_gemini(root: &Path, workspace: Option<&Path>, found: &mut Vec<SessionDescriptor>) {
1833 let slug_to_cwd = std::fs::read_to_string(root.join("projects.json"))
1834 .ok()
1835 .and_then(|text| serde_json::from_str::<Value>(&text).ok())
1836 .and_then(|value| value.get("projects").and_then(Value::as_object).cloned())
1837 .map(|projects| {
1838 projects
1839 .into_iter()
1840 .filter_map(|(cwd, slug)| Some((slug.as_str()?.to_string(), PathBuf::from(cwd))))
1841 .collect::<HashMap<_, _>>()
1842 })
1843 .unwrap_or_default();
1844 let mut files = Vec::new();
1845 collect_jsonl(&root.join("tmp"), HarnessId::GEMINI, false, &mut files);
1846 let worker_count = std::thread::available_parallelism()
1847 .map(usize::from)
1848 .unwrap_or(4)
1849 .clamp(1, 8)
1850 .min(files.len().max(1));
1851 let chunk_size = files.len().max(1).div_ceil(worker_count);
1852 let discovered = std::thread::scope(|scope| {
1853 files
1854 .chunks(chunk_size)
1855 .map(|paths| {
1856 scope.spawn(|| {
1857 paths
1858 .iter()
1859 .filter_map(|path| gemini_descriptor(path, &slug_to_cwd, workspace))
1860 .collect::<Vec<_>>()
1861 })
1862 })
1863 .collect::<Vec<_>>()
1864 .into_iter()
1865 .flat_map(|worker| {
1866 worker
1867 .join()
1868 .expect("Gemini discovery worker must not panic")
1869 })
1870 .collect::<Vec<_>>()
1871 });
1872 found.extend(discovered);
1873}
1874
1875fn gemini_descriptor(
1876 path: &Path,
1877 slug_to_cwd: &HashMap<String, PathBuf>,
1878 workspace: Option<&Path>,
1879) -> Option<SessionDescriptor> {
1880 if path
1881 .parent()
1882 .and_then(Path::file_name)
1883 .and_then(|name| name.to_str())
1884 != Some("chats")
1885 {
1886 return None;
1887 }
1888 let slug = path
1889 .parent()
1890 .and_then(Path::parent)
1891 .and_then(Path::file_name)
1892 .and_then(|name| name.to_str());
1893 let cwd = slug.and_then(|slug| slug_to_cwd.get(slug)).cloned();
1894 if workspace.is_some_and(|wanted| {
1895 cwd.as_deref()
1896 .is_none_or(|actual| !recorded_cwd_matches(actual, wanted))
1897 }) {
1898 return None;
1899 }
1900
1901 let file = File::open(path).ok()?;
1905 let mut reader = BufReader::new(file.take(64 * 1024));
1906 let mut header = String::new();
1907 reader.read_line(&mut header).ok()?;
1908 let header = serde_json::from_str::<Value>(&header).ok()?;
1909 let session_id = header.get("sessionId")?.as_str()?.to_string();
1910 let mut model = None;
1911 for line in reader
1912 .take(4 * 1024)
1913 .lines()
1914 .map_while(std::result::Result::ok)
1915 {
1916 let Ok(value) = serde_json::from_str::<Value>(&line) else {
1917 continue;
1918 };
1919 let kind = value.get("type").and_then(Value::as_str);
1920 if kind != Some("user") && kind != Some("gemini") {
1921 continue;
1922 }
1923 if model.is_none() {
1924 model = value
1925 .get("model")
1926 .and_then(Value::as_str)
1927 .map(str::to_string);
1928 }
1929 if model.is_some() {
1930 break;
1931 }
1932 }
1933 Some(SessionDescriptor {
1934 locator: SessionLocator {
1935 harness: HarnessId::from(HarnessId::GEMINI),
1936 session_id,
1937 storage: StorageLocator::File {
1938 path: path.to_path_buf(),
1939 },
1940 },
1941 cwd,
1942 title: None,
1943 preview_candidates: Vec::new(),
1944 latest_message_candidates: Vec::new(),
1945 updated_at_ms: tail_facts(path, HarnessId::GEMINI)
1946 .last_turn_ms
1947 .or_else(|| modified_ms(path)),
1948 message_count: None,
1949 model,
1950 parent_session_id: None,
1951 child_session_count: 0,
1952 nouns: OrchestrationNouns::default(),
1953 })
1954}
1955
1956fn display_text(content: Option<&Value>) -> Option<String> {
1957 match content? {
1958 Value::String(text) => Some(text.clone()),
1959 Value::Array(parts) => Some(
1960 parts
1961 .iter()
1962 .filter_map(|part| part.get("text").and_then(Value::as_str))
1963 .collect::<Vec<_>>()
1964 .join(" ")
1965 .trim()
1966 .to_string(),
1967 ),
1968 _ => None,
1969 }
1970}
1971
1972fn discover_hermes(
1977 db_path: &Path,
1978 workspace: Option<&Path>,
1979 selected: &HashSet<&str>,
1980 found: &mut Vec<SessionDescriptor>,
1981) {
1982 if !db_path.is_file() {
1983 return;
1984 }
1985 let Ok(conn) = Connection::open_with_flags(
1986 db_path,
1987 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
1988 ) else {
1989 return;
1990 };
1991 let fingerprint_ok = ["sessions", "messages", "schema_version"].iter().all(|t| {
1992 conn.query_row(
1993 "SELECT 1 FROM sqlite_master WHERE type='table' AND name=?1",
1994 [t],
1995 |_| Ok(()),
1996 )
1997 .is_ok()
1998 });
1999 if !fingerprint_ok {
2000 return;
2001 }
2002 let Ok(mut statement) = conn.prepare(
2003 "SELECT id, cwd, title, model, message_count, started_at, ended_at, parent_session_id, \
2004 source, model_config FROM sessions ORDER BY started_at DESC",
2005 ) else {
2006 return;
2007 };
2008 let Ok(rows) = statement.query_map([], |row| {
2009 Ok((
2010 row.get::<_, String>(0)?,
2011 row.get::<_, Option<String>>(1)?,
2012 row.get::<_, Option<String>>(2)?,
2013 row.get::<_, Option<String>>(3)?,
2014 row.get::<_, Option<i64>>(4)?,
2015 row.get::<_, Option<f64>>(5)?,
2016 row.get::<_, Option<f64>>(6)?,
2017 row.get::<_, Option<String>>(7)?,
2018 row.get::<_, Option<String>>(8)?,
2019 row.get::<_, Option<String>>(9)?,
2020 ))
2021 }) else {
2022 return;
2023 };
2024 for row in rows.flatten() {
2025 let (
2026 id,
2027 cwd,
2028 title,
2029 model,
2030 message_count,
2031 started_at,
2032 ended_at,
2033 parent,
2034 source,
2035 model_config,
2036 ) = row;
2037 let mirror = model_config
2041 .as_deref()
2042 .and_then(|c| serde_json::from_str::<serde_json::Value>(c).ok())
2043 .and_then(|c| c.get("_supercode_mirror").cloned());
2044 if let Some(mirror) = mirror {
2045 let worker_read = mirror
2046 .get("harness")
2047 .and_then(serde_json::Value::as_str)
2048 .is_some_and(|h| selected.contains(h));
2049 let mirrored = mirror.get("messages").and_then(serde_json::Value::as_i64);
2050 if worker_read && mirrored.is_none_or(|m| message_count.unwrap_or(0) <= m) {
2051 continue;
2052 }
2053 if mirror.get("continued_as").is_some()
2054 && mirrored.is_none_or(|m| message_count.unwrap_or(0) <= m)
2055 {
2056 continue;
2057 }
2058 }
2059 let cwd = cwd.map(PathBuf::from);
2060 if let Some(filter) = workspace {
2061 if cwd.as_deref() != Some(filter) {
2062 continue;
2063 }
2064 }
2065 let updated_at_ms = ended_at
2066 .or(started_at)
2067 .map(|seconds| (seconds * 1000.0) as u64);
2068 let mut meta = SessionMeta::new(SessionSource::Hermes);
2072 meta.cwd = cwd.clone();
2073 if let Some(hermes_source) = source.filter(|value| !value.is_empty()) {
2074 meta.lineage
2075 .insert("hermes_source".to_string(), hermes_source);
2076 }
2077 if let Some(parent_id) = parent.as_deref() {
2078 meta.lineage.insert(
2079 "hermes_lineage_kind".to_string(),
2080 crate::session::hermes_lineage_kind(
2081 &conn,
2082 parent_id,
2083 model_config.as_deref(),
2084 started_at,
2085 )
2086 .to_string(),
2087 );
2088 }
2089 hermes_capture_nouns(&conn, &id, &mut meta);
2090 found.push(SessionDescriptor {
2091 locator: SessionLocator {
2092 harness: HarnessId::new(HarnessId::HERMES),
2093 session_id: id,
2094 storage: StorageLocator::File {
2095 path: db_path.to_path_buf(),
2096 },
2097 },
2098 cwd,
2099 title: title.filter(|t| !t.is_empty()),
2100 preview_candidates: Vec::new(),
2101 latest_message_candidates: Vec::new(),
2102 updated_at_ms,
2103 message_count: message_count.map(|count| count.max(0) as usize),
2104 model,
2105 parent_session_id: parent,
2106 child_session_count: 0,
2107 nouns: OrchestrationNouns::from_meta(&meta),
2108 });
2109 }
2110}
2111
2112fn discover_orchestrator(
2119 root: &Path,
2120 workspace: Option<&Path>,
2121 found: &mut Vec<SessionDescriptor>,
2122) {
2123 if workspace.is_some() {
2125 return;
2126 }
2127 for (profile, dir) in orchestrator_profile_dirs(root) {
2128 let db_path = dir.join("state.db");
2129 if !db_path.is_file() {
2130 continue;
2131 }
2132 let Ok(conn) = Connection::open_with_flags(
2133 &db_path,
2134 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
2135 ) else {
2136 continue;
2137 };
2138 let sessions = dir.join("sessions");
2141 let scope = std::fs::canonicalize(&sessions)
2142 .unwrap_or(sessions)
2143 .display()
2144 .to_string();
2145 let Ok(mut statement) = conn.prepare(
2146 "SELECT json_extract(entry_json, '$.metadata.supercode.binding'), \
2147 CAST(strftime('%s', json_extract(entry_json, '$.metadata.supercode.binding.last_activity_at')) AS INTEGER) \
2148 FROM gateway_routing WHERE scope = ?1 AND json_extract(entry_json, '$.metadata.supercode') IS NOT NULL \
2149 ORDER BY json_extract(entry_json, '$.metadata.supercode.binding.last_activity_at') DESC",
2150 ) else {
2151 continue;
2152 };
2153 let Ok(rows) = statement.query_map([&scope], |row| {
2154 Ok((
2155 row.get::<_, Option<String>>(0)?,
2156 row.get::<_, Option<i64>>(1)?,
2157 ))
2158 }) else {
2159 continue;
2160 };
2161 let rows = rows.flatten().filter_map(|(json, epoch)| {
2162 let b: Binding = serde_json::from_str(&json?).ok()?;
2163 let text = |v: &Option<String>| v.clone().filter(|s| !s.is_empty());
2164 Some((
2165 OrchestratorBindingRow {
2166 platform: b.key.platform.clone().unwrap_or_default(),
2167 chat_type: b.key.kind.clone().unwrap_or_default(),
2168 chat_id: text(&b.key.chat_id),
2169 thread_id: text(&b.key.thread_id),
2170 participant_id: text(&b.key.participant_id),
2171 worker_harness: b.worker.harness.as_str().to_string(),
2172 worker_session_id: text(&b.worker.session_id),
2173 worker_locator: text(&b.worker.locator),
2174 started_at: b.started_at.clone(),
2175 last_activity_at: b.last_activity_at.clone(),
2176 ended_at: b.ended_at.clone(),
2177 end_reason: b.end_reason.map(|r| r.as_str().to_string()),
2178 handoff_to: b.handoff.as_ref().and_then(|h| h.to.clone()),
2179 handoff_state: b.handoff.as_ref().map(|h| h.state.clone()),
2180 handoff_error: b.handoff.as_ref().and_then(|h| h.error.clone()),
2181 recurrence_job_id: b.recurrence.as_ref().map(|r| r.job_id.clone()),
2182 },
2183 epoch,
2184 ))
2185 });
2186 for (row, last_activity_epoch) in rows {
2187 found.push(orchestrator_descriptor(
2188 &db_path,
2189 &profile,
2190 &row,
2191 last_activity_epoch,
2192 ));
2193 }
2194 }
2195}
2196
2197fn orchestrator_descriptor(
2198 db_path: &Path,
2199 profile: &str,
2200 row: &OrchestratorBindingRow,
2201 last_activity_epoch: Option<i64>,
2202) -> SessionDescriptor {
2203 let binding = Binding::from_orchestrator_row(profile, row);
2204 let nouns = binding.nouns();
2205 let mut title = format!(
2208 "{} {}",
2209 row.worker_harness,
2210 row.worker_session_id
2211 .as_deref()
2212 .unwrap_or("(no worker session yet)")
2213 );
2214 if let Some(reason) = row.end_reason.as_deref().filter(|_| row.ended_at.is_some()) {
2215 title.push_str(&format!(" (ended: {reason})"));
2216 }
2217 SessionDescriptor {
2218 locator: SessionLocator {
2219 harness: HarnessId::new(HarnessId::ORCHESTRATOR),
2220 session_id: row.worker_session_id.clone().unwrap_or_default(),
2221 storage: StorageLocator::File {
2222 path: row
2223 .worker_locator
2224 .clone()
2225 .map_or_else(|| db_path.to_path_buf(), PathBuf::from),
2226 },
2227 },
2228 cwd: None,
2229 title: Some(title),
2230 preview_candidates: Vec::new(),
2231 latest_message_candidates: Vec::new(),
2232 updated_at_ms: last_activity_epoch.map(|seconds| (seconds.max(0) as u64) * 1000),
2233 message_count: None,
2234 model: None,
2235 parent_session_id: None,
2236 child_session_count: 0,
2237 nouns,
2238 }
2239}
2240
2241fn discover_openclaw(root: &Path, workspace: Option<&Path>, found: &mut Vec<SessionDescriptor>) {
2247 let agents = root.join("agents");
2248 let Ok(agent_dirs) = std::fs::read_dir(&agents) else {
2249 return;
2250 };
2251 for agent_dir in agent_dirs.flatten() {
2252 let sessions = agent_dir.path().join("sessions");
2253 let Ok(files) = std::fs::read_dir(&sessions) else {
2254 continue;
2255 };
2256 for file in files.flatten() {
2257 let path = file.path();
2258 let name = file.file_name();
2259 let name = name.to_string_lossy();
2260 if !name.ends_with(".jsonl") || name.ends_with(".trajectory.jsonl") {
2261 continue;
2262 }
2263 let Ok(text) = std::fs::read_to_string(&path) else {
2264 continue;
2265 };
2266 let Some(header_line) = text.lines().find(|line| !line.trim().is_empty()) else {
2267 continue;
2268 };
2269 let Ok(header) = serde_json::from_str::<serde_json::Value>(header_line) else {
2270 continue;
2271 };
2272 if header.get("type").and_then(serde_json::Value::as_str) != Some("session") {
2273 continue;
2274 }
2275 let session_id = header
2276 .get("id")
2277 .and_then(serde_json::Value::as_str)
2278 .unwrap_or_else(|| name.trim_end_matches(".jsonl"))
2279 .to_string();
2280 let cwd = header
2281 .get("cwd")
2282 .and_then(serde_json::Value::as_str)
2283 .map(PathBuf::from);
2284 if let Some(filter) = workspace {
2285 if cwd.as_deref() != Some(filter) {
2286 continue;
2287 }
2288 }
2289 let updated_at_ms = tail_facts(&path, HarnessId::OPENCLAW)
2290 .last_turn_ms
2291 .or_else(|| {
2292 file.metadata()
2293 .ok()
2294 .and_then(|metadata| metadata.modified().ok())
2295 .and_then(|modified| modified.duration_since(std::time::UNIX_EPOCH).ok())
2296 .map(|elapsed| elapsed.as_millis() as u64)
2297 });
2298 let message_count = text
2299 .lines()
2300 .filter(|line| line.contains("\"type\":\"message\""))
2301 .count();
2302 let mut meta = SessionMeta::new(SessionSource::OpenClaw);
2306 meta.cwd = cwd.clone();
2307 openclaw_capture_header_nouns(&header, &mut meta);
2308 if meta.profile.is_none() {
2309 meta.profile = openclaw_agent_id_from_path(&path);
2310 }
2311 found.push(SessionDescriptor {
2312 locator: SessionLocator {
2313 harness: HarnessId::new(HarnessId::OPENCLAW),
2314 session_id,
2315 storage: StorageLocator::File { path },
2316 },
2317 cwd,
2318 title: None,
2319 preview_candidates: Vec::new(),
2320 latest_message_candidates: Vec::new(),
2321 updated_at_ms,
2322 message_count: Some(message_count),
2323 model: None,
2324 parent_session_id: None,
2325 child_session_count: 0,
2326 nouns: OrchestrationNouns::from_meta(&meta),
2327 });
2328 }
2329 }
2330}
2331
2332fn discover_supercode(root: &Path, workspace: Option<&Path>, found: &mut Vec<SessionDescriptor>) {
2333 for info in list_native_store(root) {
2334 let path = if info.archived {
2335 root.join("archived").join(format!("{}.jsonl", info.name))
2336 } else {
2337 root.join(format!("{}.jsonl", info.name))
2338 };
2339 let sidecar = path.with_extension("sidecar.jsonl");
2344 let path = if path.is_file() {
2345 path
2346 } else {
2347 sidecar.clone()
2348 };
2349 let header = read_native_store_header(&path);
2350 if workspace.is_some_and(|wanted| {
2351 header
2352 .as_ref()
2353 .and_then(|meta| meta.cwd.as_deref())
2354 .is_none_or(|cwd| !recorded_cwd_matches(cwd, wanted))
2355 }) {
2356 continue;
2357 }
2358 let title = (!info.title.trim().is_empty()).then_some(info.title);
2359 let updated_at_ms = tail_facts(&sidecar, HarnessId::SUPERCODE)
2362 .last_turn_ms
2363 .or_else(|| tail_facts(&path, HarnessId::SUPERCODE).last_turn_ms)
2364 .or_else(|| modified_ms(&path))
2365 .or_else(|| modified_ms(&sidecar));
2366 found.push(SessionDescriptor {
2367 locator: SessionLocator {
2368 harness: HarnessId::from(HarnessId::SUPERCODE),
2369 session_id: info.name,
2370 storage: StorageLocator::File { path: path.clone() },
2371 },
2372 cwd: header.as_ref().and_then(|meta| meta.cwd.clone()),
2373 title,
2374 preview_candidates: Vec::new(),
2375 latest_message_candidates: Vec::new(),
2376 updated_at_ms,
2377 message_count: None,
2378 model: header.and_then(|meta| meta.model),
2379 parent_session_id: None,
2380 child_session_count: 0,
2381 nouns: OrchestrationNouns::default(),
2382 });
2383 }
2384}
2385
2386fn read_native_store_header(path: &Path) -> Option<HeaderMeta> {
2391 let name = path.file_stem()?.to_str()?;
2392 let sidecar = path.with_file_name(format!("{name}.sidecar.jsonl"));
2393 let source_path = if sidecar.is_file() {
2394 sidecar
2395 } else {
2396 path.to_path_buf()
2397 };
2398 let file = File::open(source_path).ok()?;
2399 let mut result = HeaderMeta::default();
2400 let mut source = None;
2401 let mut bytes = 0usize;
2402 for line in BufReader::new(file).lines().take(32) {
2403 let line = line.ok()?;
2404 bytes += line.len();
2405 if bytes > 256 * 1024 {
2406 break;
2407 }
2408 let Ok(value) = serde_json::from_str::<Value>(&line) else {
2409 continue;
2410 };
2411 if source.is_none() {
2412 source = value.get("source").and_then(Value::as_str).map(|source| {
2413 if source == "claude_code" {
2414 HarnessId::CLAUDE_CODE.to_string()
2415 } else {
2416 source.to_string()
2417 }
2418 });
2419 fill_string(&mut result.session_id, value.get("session_id"));
2420 }
2421 if let Some(harness) = source.as_deref() {
2422 update_header_meta(&mut result, &value, harness);
2423 }
2424 if result.cwd.is_some() && result.model.is_some() {
2425 break;
2426 }
2427 }
2428 Some(result)
2429}
2430
2431#[derive(Deserialize)]
2432struct NativeStoreInfo {
2433 name: String,
2434 #[serde(default)]
2435 title: String,
2436 #[serde(skip)]
2437 archived: bool,
2438}
2439
2440fn list_native_store(root: &Path) -> Vec<NativeStoreInfo> {
2441 let mut sessions = Vec::new();
2442 for archived in [false, true] {
2443 let directory = if archived {
2444 root.join("archived")
2445 } else {
2446 root.to_path_buf()
2447 };
2448 let Ok(entries) = fs::read_dir(directory) else {
2449 continue;
2450 };
2451 for entry in entries.flatten() {
2452 let path = entry.path();
2453 if !path.to_string_lossy().ends_with(".meta.json") {
2454 continue;
2455 }
2456 let Ok(text) = fs::read_to_string(path) else {
2457 continue;
2458 };
2459 let Ok(mut info) = serde_json::from_str::<NativeStoreInfo>(&text) else {
2460 continue;
2461 };
2462 info.archived = archived;
2463 sessions.push(info);
2464 }
2465 }
2466 sessions.sort_by(|left, right| left.name.cmp(&right.name));
2467 sessions
2468}
2469
2470fn discover_grok(root: &Path, workspace: Option<&Path>, found: &mut Vec<SessionDescriptor>) {
2471 let Ok(workspaces) = fs::read_dir(root) else {
2472 return;
2473 };
2474 for workspace_entry in workspaces.flatten() {
2475 let encoded = workspace_entry.file_name();
2476 let Some(cwd) = encoded
2477 .to_str()
2478 .and_then(percent_decode_path)
2479 .map(PathBuf::from)
2480 else {
2481 continue;
2482 };
2483 if workspace.is_some_and(|wanted| !recorded_cwd_matches(&cwd, wanted)) {
2484 continue;
2485 }
2486 let Ok(sessions) = fs::read_dir(workspace_entry.path()) else {
2487 continue;
2488 };
2489 for session_entry in sessions.flatten() {
2490 let session_dir = session_entry.path();
2491 if !session_dir.is_dir() {
2492 continue;
2493 }
2494 let transcript = session_dir.join("chat_history.jsonl");
2495 if !transcript.is_file() {
2496 continue;
2497 }
2498 let Some(session_id) = session_dir
2499 .file_name()
2500 .and_then(|name| name.to_str())
2501 .map(str::to_string)
2502 else {
2503 continue;
2504 };
2505 let summary = fs::read_to_string(session_dir.join("summary.json"))
2506 .ok()
2507 .and_then(|text| serde_json::from_str::<Value>(&text).ok());
2508 let title = summary
2509 .as_ref()
2510 .and_then(|value| value.get("generated_title"))
2511 .and_then(Value::as_str)
2512 .filter(|title| !title.is_empty())
2513 .map(str::to_string);
2514 let model = summary
2515 .as_ref()
2516 .and_then(|value| value.get("current_model_id"))
2517 .and_then(Value::as_str)
2518 .map(str::to_string);
2519 let message_count = summary
2520 .as_ref()
2521 .and_then(|value| value.get("num_chat_messages"))
2522 .and_then(Value::as_u64)
2523 .and_then(|count| usize::try_from(count).ok());
2524 let updated_at_ms = summary
2525 .as_ref()
2526 .and_then(|value| value.get("updated_at"))
2527 .and_then(Value::as_str)
2528 .and_then(crate::sidecar::rfc3339_to_ms)
2529 .and_then(|millis| u64::try_from(millis).ok())
2530 .or_else(|| modified_ms(&transcript));
2531 found.push(SessionDescriptor {
2532 locator: SessionLocator {
2533 harness: HarnessId::from(HarnessId::GROK),
2534 session_id,
2535 storage: StorageLocator::File { path: transcript },
2536 },
2537 cwd: Some(cwd.clone()),
2538 title,
2539 preview_candidates: Vec::new(),
2540 latest_message_candidates: Vec::new(),
2541 updated_at_ms,
2542 message_count,
2543 model,
2544 parent_session_id: None,
2545 child_session_count: 0,
2546 nouns: OrchestrationNouns::default(),
2547 });
2548 }
2549 }
2550}
2551
2552fn discover_opencode(root: &Path, workspace: Option<&Path>, found: &mut Vec<SessionDescriptor>) {
2553 let mut dbs = Vec::new();
2554 if root.is_file() {
2555 dbs.push(root.to_path_buf());
2556 } else if let Ok(entries) = fs::read_dir(root) {
2557 dbs.extend(entries.flatten().map(|entry| entry.path()).filter(|path| {
2558 path.file_name()
2559 .and_then(|v| v.to_str())
2560 .is_some_and(|name| name.starts_with("opencode") && name.ends_with(".db"))
2561 }));
2562 }
2563 dbs.sort();
2564 for db in dbs {
2565 let Ok(conn) = Connection::open_with_flags(
2566 &db,
2567 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
2568 ) else {
2569 continue;
2570 };
2571 let has_model = conn.prepare("SELECT model FROM session LIMIT 0").is_ok();
2572 let model_column = if has_model { "s.model" } else { "NULL" };
2573 let query = format!(
2574 "SELECT s.id, s.directory, s.title, s.time_updated, {model_column}, COUNT(m.id) \
2575 FROM session s LEFT JOIN message m ON m.session_id = s.id \
2576 GROUP BY s.id ORDER BY s.time_updated DESC"
2577 );
2578 let Ok(mut stmt) = conn.prepare(&query) else {
2579 continue;
2580 };
2581 let Ok(rows) = stmt.query_map([], |row| {
2582 Ok((
2583 row.get::<_, String>(0)?,
2584 row.get::<_, String>(1)?,
2585 row.get::<_, String>(2)?,
2586 row.get::<_, i64>(3)?,
2587 row.get::<_, Option<String>>(4)?,
2588 row.get::<_, i64>(5)?,
2589 ))
2590 }) else {
2591 continue;
2592 };
2593 for row in rows.flatten() {
2594 let (id, cwd, title, updated, model, messages) = row;
2595 let cwd = PathBuf::from(cwd);
2596 if workspace.is_some_and(|wanted| !recorded_cwd_matches(&cwd, wanted)) {
2597 continue;
2598 }
2599 found.push(SessionDescriptor {
2600 locator: SessionLocator {
2601 harness: HarnessId::from(HarnessId::OPENCODE),
2602 session_id: id.clone(),
2603 storage: StorageLocator::Sqlite {
2604 path: db.clone(),
2605 selector: id,
2606 },
2607 },
2608 cwd: Some(cwd),
2609 title: (!title.is_empty()).then_some(title),
2610 preview_candidates: Vec::new(),
2611 latest_message_candidates: Vec::new(),
2612 updated_at_ms: u64::try_from(updated).ok(),
2613 message_count: usize::try_from(messages).ok(),
2614 model,
2615 parent_session_id: None,
2616 child_session_count: 0,
2617 nouns: OrchestrationNouns::default(),
2618 });
2619 }
2620 }
2621}
2622
2623fn discover_goose(root: &Path, workspace: Option<&Path>, found: &mut Vec<SessionDescriptor>) {
2624 let db = if root.is_file() {
2625 root.to_path_buf()
2626 } else if root.join("sessions.db").is_file() {
2627 root.join("sessions.db")
2628 } else {
2629 root.join("sessions/sessions.db")
2630 };
2631 let Ok(connection) = Connection::open_with_flags(
2632 &db,
2633 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
2634 ) else {
2635 return;
2636 };
2637 let Ok(mut statement) = connection.prepare(
2638 "SELECT s.id, s.working_dir, s.name, s.updated_at, s.model_config_json, \
2639 COUNT(m.id) \
2640 FROM sessions s LEFT JOIN messages m ON m.session_id = s.id \
2641 WHERE s.archived_at IS NULL \
2642 GROUP BY s.id ORDER BY s.updated_at DESC",
2643 ) else {
2644 return;
2645 };
2646 let Ok(rows) = statement.query_map([], |row| {
2647 Ok((
2648 row.get::<_, String>(0)?,
2649 row.get::<_, String>(1)?,
2650 row.get::<_, String>(2)?,
2651 row.get::<_, String>(3)?,
2652 row.get::<_, Option<String>>(4)?,
2653 row.get::<_, i64>(5)?,
2654 ))
2655 }) else {
2656 return;
2657 };
2658 for row in rows.flatten() {
2659 let (id, cwd, title, updated_at, model_config, message_count) = row;
2660 let cwd = PathBuf::from(cwd);
2661 if workspace.is_some_and(|wanted| !recorded_cwd_matches(&cwd, wanted)) {
2662 continue;
2663 }
2664 let model = model_config
2665 .as_deref()
2666 .and_then(|value| serde_json::from_str::<Value>(value).ok())
2667 .and_then(|value| {
2668 value
2669 .get("model_name")
2670 .or_else(|| value.get("modelName"))
2671 .and_then(Value::as_str)
2672 .map(str::to_string)
2673 });
2674 let updated_at_ms = crate::sidecar::rfc3339_to_ms(&updated_at)
2675 .or_else(|| {
2676 crate::sidecar::rfc3339_to_ms(&format!("{}Z", updated_at.replace(' ', "T")))
2678 })
2679 .and_then(|value| u64::try_from(value).ok());
2680 found.push(SessionDescriptor {
2681 locator: SessionLocator {
2682 harness: HarnessId::from(HarnessId::GOOSE),
2683 session_id: id.clone(),
2684 storage: StorageLocator::Sqlite {
2685 path: db.clone(),
2686 selector: id,
2687 },
2688 },
2689 cwd: Some(cwd),
2690 title: (!title.trim().is_empty()).then_some(title),
2691 preview_candidates: Vec::new(),
2692 latest_message_candidates: Vec::new(),
2693 updated_at_ms,
2694 message_count: usize::try_from(message_count).ok(),
2695 model,
2696 parent_session_id: None,
2697 child_session_count: 0,
2698 nouns: OrchestrationNouns::default(),
2699 });
2700 }
2701}
2702
2703const LATEST_PREVIEW_CANDIDATES: usize = 8;
2704const TOPIC_PREVIEW_HEAD_BYTES: u64 = 512 * 1024;
2705const LATEST_PREVIEW_TAIL_BYTES: u64 = 512 * 1024;
2706const LATEST_PREVIEW_MAX_BYTES: u64 = 4 * 1024 * 1024;
2707
2708fn topic_message_candidates(locator: &SessionLocator) -> Result<Vec<SessionPreviewCandidate>> {
2709 match &locator.storage {
2710 StorageLocator::File { path }
2711 if matches!(
2712 locator.harness.as_str(),
2713 HarnessId::CLAUDE_CODE | HarnessId::CODEX
2714 ) =>
2715 {
2716 topic_file_message_candidates(path, locator.harness.as_str())
2717 }
2718 _ => Ok(Vec::new()),
2719 }
2720}
2721
2722fn codex_history_topics(
2723 sessions_root: &Path,
2724 sessions: &[SessionDescriptor],
2725) -> Result<HashMap<String, Vec<SessionPreviewCandidate>>> {
2726 let wanted: HashSet<&str> = sessions
2727 .iter()
2728 .filter(|descriptor| descriptor.locator.harness.as_str() == HarnessId::CODEX)
2729 .map(|descriptor| descriptor.locator.session_id.as_str())
2730 .collect();
2731 if wanted.is_empty() {
2732 return Ok(HashMap::new());
2733 }
2734 let Some(root) = sessions_root.parent() else {
2735 return Ok(HashMap::new());
2736 };
2737 let file = match File::open(root.join("history.jsonl")) {
2738 Ok(file) => file,
2739 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(HashMap::new()),
2740 Err(error) => return Err(error.into()),
2741 };
2742 let mut topics = HashMap::new();
2743 for line in BufReader::new(file).lines() {
2744 let Ok(value) = serde_json::from_str::<Value>(&line?) else {
2745 continue;
2746 };
2747 let Some(session_id) = value.get("session_id").and_then(Value::as_str) else {
2748 continue;
2749 };
2750 if !wanted.contains(session_id) || topics.contains_key(session_id) {
2751 continue;
2752 }
2753 let mut candidates = Vec::new();
2754 push_message_candidate(&mut candidates, "user", value.get("text"), HashMap::new());
2755 if !candidates.is_empty() {
2756 topics.insert(session_id.to_string(), candidates);
2757 if topics.len() == wanted.len() {
2758 break;
2759 }
2760 }
2761 }
2762 Ok(topics)
2763}
2764
2765fn latest_message_candidates(locator: &SessionLocator) -> Result<Vec<SessionPreviewCandidate>> {
2766 match &locator.storage {
2767 StorageLocator::File { path } | StorageLocator::Sqlite { path, .. }
2770 if locator.harness.as_str() == HarnessId::HERMES =>
2771 {
2772 latest_hermes_message_candidates(path, &locator.session_id)
2773 }
2774 StorageLocator::File { path } => {
2775 latest_file_message_candidates(path, locator.harness.as_str())
2776 }
2777 StorageLocator::Sqlite { path, selector }
2778 if locator.harness.as_str() == HarnessId::OPENCODE =>
2779 {
2780 latest_opencode_message_candidates(path, selector)
2781 }
2782 StorageLocator::Sqlite { path, selector }
2783 if locator.harness.as_str() == HarnessId::GOOSE =>
2784 {
2785 latest_goose_message_candidates(path, selector)
2786 }
2787 StorageLocator::Sqlite { .. } => Ok(Vec::new()),
2788 }
2789}
2790
2791fn topic_file_message_candidates(
2792 path: &Path,
2793 harness: &str,
2794) -> Result<Vec<SessionPreviewCandidate>> {
2795 let mut file = File::open(path)?;
2796 let mut bytes = Vec::with_capacity(TOPIC_PREVIEW_HEAD_BYTES as usize);
2797 file.by_ref()
2798 .take(TOPIC_PREVIEW_HEAD_BYTES)
2799 .read_to_end(&mut bytes)?;
2800 if file.metadata()?.len() > TOPIC_PREVIEW_HEAD_BYTES {
2801 if let Some(newline) = bytes.iter().rposition(|byte| *byte == b'\n') {
2802 bytes.truncate(newline);
2803 }
2804 }
2805 let text = String::from_utf8(bytes).map_err(|_| {
2806 Error::Other(format!(
2807 "{} contains non-UTF-8 data in its topic-preview window",
2808 path.display()
2809 ))
2810 })?;
2811 if harness == HarnessId::CODEX {
2812 return Ok(codex_preview_candidates(text.lines(), false));
2813 }
2814 let mut candidates = Vec::new();
2815 for line in text.lines() {
2816 let Ok(value) = serde_json::from_str::<Value>(line) else {
2817 continue;
2818 };
2819 push_topic_message_candidate(&mut candidates, harness, &value);
2820 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
2821 break;
2822 }
2823 }
2824 Ok(candidates)
2825}
2826
2827fn latest_file_message_candidates(
2828 path: &Path,
2829 harness: &str,
2830) -> Result<Vec<SessionPreviewCandidate>> {
2831 let mut candidates =
2832 latest_file_message_candidates_with_limit(path, harness, LATEST_PREVIEW_TAIL_BYTES)?;
2833 if candidates.is_empty() {
2834 candidates =
2835 latest_file_message_candidates_with_limit(path, harness, LATEST_PREVIEW_MAX_BYTES)?;
2836 }
2837 Ok(candidates)
2838}
2839
2840fn latest_file_message_candidates_with_limit(
2841 path: &Path,
2842 harness: &str,
2843 byte_limit: u64,
2844) -> Result<Vec<SessionPreviewCandidate>> {
2845 let mut file = File::open(path)?;
2846 let file_len = file.metadata()?.len();
2847 let start = file_len.saturating_sub(byte_limit);
2848 file.seek(SeekFrom::Start(start))?;
2849 let mut bytes = Vec::with_capacity((file_len - start) as usize);
2850 file.read_to_end(&mut bytes)?;
2851 if start > 0 {
2852 if let Some(newline) = bytes.iter().position(|byte| *byte == b'\n') {
2853 bytes.drain(..=newline);
2854 } else {
2855 return Ok(Vec::new());
2856 }
2857 }
2858 let text = String::from_utf8(bytes).map_err(|_| {
2859 Error::Other(format!(
2860 "{} contains non-UTF-8 data in its list-preview window",
2861 path.display()
2862 ))
2863 })?;
2864 if harness == HarnessId::CODEX {
2865 return Ok(codex_preview_candidates(text.lines().rev(), true));
2866 }
2867 let mut candidates = Vec::new();
2868 for line in text.lines().rev() {
2869 let Ok(value) = serde_json::from_str::<Value>(line) else {
2870 continue;
2871 };
2872 let (role, content, metadata) = match harness {
2873 HarnessId::CLAUDE_CODE => {
2874 let role = value.get("type").and_then(Value::as_str);
2875 if !matches!(role, Some("user" | "assistant")) {
2876 continue;
2877 }
2878 let metadata = if role == Some("user") {
2879 crate::session::claude_user_provenance(&value)
2880 .into_iter()
2881 .collect()
2882 } else {
2883 HashMap::new()
2884 };
2885 (
2886 role.unwrap_or_default(),
2887 value
2888 .get("message")
2889 .and_then(|message| message.get("content")),
2890 metadata,
2891 )
2892 }
2893 HarnessId::PI => {
2894 if value.get("type").and_then(Value::as_str) != Some("message") {
2895 continue;
2896 }
2897 let message = value.get("message").unwrap_or(&Value::Null);
2898 let Some(role @ ("user" | "assistant")) =
2899 message.get("role").and_then(Value::as_str)
2900 else {
2901 continue;
2902 };
2903 (role, message.get("content"), HashMap::new())
2904 }
2905 HarnessId::GEMINI => {
2906 let Some(kind @ ("user" | "gemini")) = value.get("type").and_then(Value::as_str)
2907 else {
2908 continue;
2909 };
2910 (
2911 if kind == "gemini" {
2912 "assistant"
2913 } else {
2914 "user"
2915 },
2916 value.get("content"),
2917 HashMap::new(),
2918 )
2919 }
2920 HarnessId::GROK => {
2921 let Some(role @ ("user" | "assistant")) = value.get("type").and_then(Value::as_str)
2922 else {
2923 continue;
2924 };
2925 (role, value.get("content"), HashMap::new())
2926 }
2927 HarnessId::SUPERCODE => {
2928 let Some(role @ ("user" | "assistant")) = value.get("role").and_then(Value::as_str)
2929 else {
2930 continue;
2931 };
2932 (role, value.get("content"), HashMap::new())
2933 }
2934 _ => continue,
2935 };
2936 let mut metadata = metadata;
2937 if matches!(harness, HarnessId::CLAUDE_CODE | HarnessId::CODEX) {
2938 if let Some(timestamp) = value.get("timestamp").and_then(Value::as_str) {
2939 metadata.insert("timestamp".to_string(), timestamp.to_string());
2940 }
2941 }
2942 push_message_candidate_with_cursor(
2943 &mut candidates,
2944 role,
2945 content,
2946 metadata,
2947 Some(message_candidate_cursor(harness, &value)),
2948 );
2949 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
2950 break;
2951 }
2952 }
2953 Ok(candidates)
2954}
2955
2956struct CodexPreviewRecord {
2957 native: Value,
2958 role: String,
2959 text: String,
2960}
2961
2962fn codex_preview_candidates<'a>(
2968 lines: impl Iterator<Item = &'a str>,
2969 latest: bool,
2970) -> Vec<SessionPreviewCandidate> {
2971 let mut candidates = Vec::new();
2972 let mut pending: Option<CodexPreviewRecord> = None;
2973 for line in lines {
2974 let Ok(native) = serde_json::from_str::<Value>(line) else {
2975 continue;
2976 };
2977 let Some((role, content)) = codex_preview_message(&native) else {
2978 continue;
2979 };
2980 let Some(text) = display_text(content).filter(|text| !text.trim().is_empty()) else {
2981 continue;
2982 };
2983 let current = CodexPreviewRecord {
2984 role: role.to_string(),
2985 text,
2986 native,
2987 };
2988 if let Some(previous) = pending.take() {
2989 if previous.role == current.role
2990 && previous.text == current.text
2991 && previous.native.get("type") != current.native.get("type")
2992 {
2993 let canonical = if previous.native.get("type").and_then(Value::as_str)
2994 == Some("response_item")
2995 {
2996 previous
2997 } else {
2998 current
2999 };
3000 push_codex_preview_candidate(&mut candidates, canonical, latest);
3001 } else {
3002 push_codex_preview_candidate(&mut candidates, previous, latest);
3003 pending = Some(current);
3004 }
3005 } else {
3006 pending = Some(current);
3007 }
3008 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3009 break;
3010 }
3011 }
3012 if let Some(last) = pending {
3013 push_codex_preview_candidate(&mut candidates, last, latest);
3014 }
3015 candidates
3016}
3017
3018fn push_codex_preview_candidate(
3019 candidates: &mut Vec<SessionPreviewCandidate>,
3020 record: CodexPreviewRecord,
3021 latest: bool,
3022) {
3023 let mut metadata = HashMap::new();
3024 if latest {
3025 if let Some(timestamp) = record.native.get("timestamp").and_then(Value::as_str) {
3026 metadata.insert("timestamp".to_string(), timestamp.to_string());
3027 }
3028 }
3029 let cursor = latest.then(|| message_candidate_cursor(HarnessId::CODEX, &record.native));
3030 push_message_candidate_with_cursor(
3031 candidates,
3032 &record.role,
3033 Some(&Value::String(record.text)),
3034 metadata,
3035 cursor,
3036 );
3037}
3038
3039fn codex_preview_message(value: &Value) -> Option<(&str, Option<&Value>)> {
3041 let payload = value.get("payload")?;
3042 match (
3043 value.get("type").and_then(Value::as_str)?,
3044 payload.get("type").and_then(Value::as_str)?,
3045 ) {
3046 ("response_item", "message") => {
3047 let role @ ("user" | "assistant") = payload.get("role").and_then(Value::as_str)? else {
3048 return None;
3049 };
3050 Some((role, payload.get("content")))
3051 }
3052 ("event_msg", "user_message") => Some(("user", payload.get("message"))),
3053 ("event_msg", "agent_message") => Some(("assistant", payload.get("message"))),
3054 _ => None,
3055 }
3056}
3057
3058fn push_topic_message_candidate(
3059 candidates: &mut Vec<SessionPreviewCandidate>,
3060 harness: &str,
3061 value: &Value,
3062) {
3063 let (role, content, metadata) = match harness {
3064 HarnessId::CLAUDE_CODE => {
3065 let role = value.get("type").and_then(Value::as_str);
3066 if !matches!(role, Some("user" | "assistant")) {
3067 return;
3068 }
3069 let metadata = if role == Some("user") {
3070 crate::session::claude_user_provenance(value)
3071 .into_iter()
3072 .collect()
3073 } else {
3074 HashMap::new()
3075 };
3076 (
3077 role.unwrap_or_default(),
3078 value
3079 .get("message")
3080 .and_then(|message| message.get("content")),
3081 metadata,
3082 )
3083 }
3084 _ => return,
3085 };
3086 push_message_candidate(candidates, role, content, metadata);
3087}
3088
3089fn latest_opencode_message_candidates(
3090 path: &Path,
3091 session_id: &str,
3092) -> Result<Vec<SessionPreviewCandidate>> {
3093 let connection = Connection::open_with_flags(
3094 path,
3095 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
3096 )
3097 .map_err(|error| Error::Other(format!("failed to open list-preview store: {error}")))?;
3098 let mut statement = connection
3099 .prepare(
3100 "SELECT m.data, p.data FROM message m JOIN part p ON p.message_id = m.id \
3101 WHERE m.session_id = ?1 ORDER BY m.time_created DESC, p.time_created DESC LIMIT 32",
3102 )
3103 .map_err(|error| Error::Other(format!("failed to prepare list-preview query: {error}")))?;
3104 let rows = statement
3105 .query_map([session_id], |row| {
3106 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
3107 })
3108 .map_err(|error| Error::Other(format!("failed to read list-preview rows: {error}")))?;
3109 let mut candidates = Vec::new();
3110 for row in rows.flatten() {
3111 let (Ok(message), Ok(part)) = (
3112 serde_json::from_str::<Value>(&row.0),
3113 serde_json::from_str::<Value>(&row.1),
3114 ) else {
3115 continue;
3116 };
3117 let Some(role @ ("user" | "assistant")) = message.get("role").and_then(Value::as_str)
3118 else {
3119 continue;
3120 };
3121 if part.get("type").and_then(Value::as_str) != Some("text") {
3122 continue;
3123 }
3124 push_message_candidate(&mut candidates, role, part.get("text"), HashMap::new());
3125 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3126 break;
3127 }
3128 }
3129 Ok(candidates)
3130}
3131
3132fn latest_hermes_message_candidates(
3133 path: &Path,
3134 session_id: &str,
3135) -> Result<Vec<SessionPreviewCandidate>> {
3136 let connection = Connection::open_with_flags(
3137 path,
3138 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
3139 )
3140 .map_err(|error| Error::Other(format!("failed to open list-preview store: {error}")))?;
3141 let mut statement = connection
3142 .prepare(
3143 "SELECT role, content FROM messages WHERE session_id = ?1 AND active = 1 \
3144 AND role IN ('user', 'assistant') AND content IS NOT NULL AND content != '' \
3145 ORDER BY timestamp DESC, id DESC LIMIT 32",
3146 )
3147 .map_err(|error| Error::Other(format!("failed to prepare list-preview query: {error}")))?;
3148 let rows = statement
3149 .query_map([session_id], |row| {
3150 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
3151 })
3152 .map_err(|error| Error::Other(format!("failed to read list-preview rows: {error}")))?;
3153 let mut candidates = Vec::new();
3154 for (role, content) in rows.flatten() {
3155 push_message_candidate(
3156 &mut candidates,
3157 &role,
3158 Some(&Value::String(content)),
3159 HashMap::new(),
3160 );
3161 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3162 break;
3163 }
3164 }
3165 Ok(candidates)
3166}
3167
3168fn latest_goose_message_candidates(
3169 path: &Path,
3170 session_id: &str,
3171) -> Result<Vec<SessionPreviewCandidate>> {
3172 let connection = Connection::open_with_flags(
3173 path,
3174 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
3175 )
3176 .map_err(|error| Error::Other(format!("failed to open list-preview store: {error}")))?;
3177 let mut statement = connection
3178 .prepare(
3179 "SELECT role, content_json FROM messages WHERE session_id = ?1 \
3180 ORDER BY created_timestamp DESC, id DESC LIMIT 16",
3181 )
3182 .map_err(|error| Error::Other(format!("failed to prepare list-preview query: {error}")))?;
3183 let rows = statement
3184 .query_map([session_id], |row| {
3185 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
3186 })
3187 .map_err(|error| Error::Other(format!("failed to read list-preview rows: {error}")))?;
3188 let mut candidates = Vec::new();
3189 for row in rows.flatten() {
3190 let (role, content) = row;
3191 if !matches!(role.as_str(), "user" | "assistant") {
3192 continue;
3193 }
3194 let Ok(content) = serde_json::from_str::<Value>(&content) else {
3195 continue;
3196 };
3197 push_message_candidate(&mut candidates, &role, Some(&content), HashMap::new());
3198 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3199 break;
3200 }
3201 }
3202 Ok(candidates)
3203}
3204
3205fn push_message_candidate(
3206 candidates: &mut Vec<SessionPreviewCandidate>,
3207 role: &str,
3208 content: Option<&Value>,
3209 metadata: HashMap<String, String>,
3210) {
3211 push_message_candidate_with_cursor(candidates, role, content, metadata, None);
3212}
3213
3214fn push_message_candidate_with_cursor(
3215 candidates: &mut Vec<SessionPreviewCandidate>,
3216 role: &str,
3217 content: Option<&Value>,
3218 metadata: HashMap<String, String>,
3219 cursor: Option<String>,
3220) {
3221 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3222 return;
3223 }
3224 let Some(text) = display_text(content).filter(|text| !text.trim().is_empty()) else {
3225 return;
3226 };
3227 const MAX_CHARS: usize = 4_096;
3228 candidates.push(SessionPreviewCandidate {
3229 cursor,
3230 role: role.to_string(),
3231 content: text.chars().take(MAX_CHARS).collect(),
3232 metadata,
3233 });
3234}
3235
3236fn message_candidate_cursor(harness: &str, value: &Value) -> String {
3237 let native_identity = value
3238 .get("uuid")
3239 .or_else(|| value.get("id"))
3240 .or_else(|| value.pointer("/message/id"))
3241 .or_else(|| value.pointer("/payload/id"))
3242 .and_then(Value::as_str)
3243 .or_else(|| value.get("timestamp").and_then(Value::as_str));
3244 let mut hasher = blake3::Hasher::new();
3245 hasher.update(b"supercode.session-preview-cursor.v1\0");
3246 hasher.update(harness.as_bytes());
3247 hasher.update(b"\0");
3248 if let Some(identity) = native_identity {
3249 hasher.update(identity.as_bytes());
3250 } else {
3251 hasher.update(value.to_string().as_bytes());
3255 }
3256 format!("v1:{}", &hasher.finalize().to_hex()[..24])
3257}
3258
3259fn fill_string(target: &mut Option<String>, value: Option<&Value>) {
3260 if target.is_none() {
3261 *target = value.and_then(Value::as_str).map(str::to_owned);
3262 }
3263}
3264
3265fn fill_path(target: &mut Option<PathBuf>, value: Option<&Value>) {
3266 if target.is_none() {
3267 *target = value.and_then(Value::as_str).map(PathBuf::from);
3268 }
3269}
3270
3271#[derive(Default)]
3274struct TailFacts {
3275 last_turn_ms: Option<u64>,
3277 model: Option<String>,
3279}
3280
3281fn tail_facts(path: &Path, harness: &str) -> TailFacts {
3297 let mut facts = TailFacts::default();
3298 let Ok(mut file) = File::open(path) else {
3299 return facts;
3300 };
3301 let Ok(len) = file.metadata().map(|meta| meta.len()) else {
3302 return facts;
3303 };
3304 let mut window = TAIL_SCAN_START.min(len);
3305 loop {
3306 if file.seek(SeekFrom::Start(len - window)).is_err() {
3307 return facts;
3308 }
3309 let Ok(size) = usize::try_from(window) else {
3310 return facts;
3311 };
3312 let mut buf = vec![0u8; size];
3313 if file.read_exact(&mut buf).is_err() {
3314 return facts;
3315 }
3316 let floor = if window < len {
3327 buf.iter().position(|byte| *byte == b'\n').map(|at| at + 1)
3328 } else {
3329 Some(0)
3330 };
3331 if let Some(floor) = floor {
3332 let mut end = buf.len();
3333 while end > floor && !(facts.last_turn_ms.is_some() && facts.model.is_some()) {
3334 let start = buf[floor..end]
3335 .iter()
3336 .rposition(|byte| *byte == b'\n')
3337 .map_or(floor, |at| floor + at + 1);
3338 if let Ok(record) = std::str::from_utf8(&buf[start..end])
3339 .map_err(|_| ())
3340 .and_then(|line| serde_json::from_str::<Value>(line).map_err(|_| ()))
3341 {
3342 if facts.last_turn_ms.is_none() {
3343 facts.last_turn_ms = record_timestamp(&record, harness)
3344 .and_then(crate::sidecar::rfc3339_to_ms)
3345 .and_then(|millis| u64::try_from(millis).ok());
3346 }
3347 if facts.model.is_none() {
3348 facts.model = record_model(&record, harness).map(str::to_owned);
3349 }
3350 }
3351 end = start.saturating_sub(1);
3352 }
3353 }
3354 if facts.last_turn_ms.is_some() || window >= len || window >= TAIL_SCAN_LIMIT {
3355 return facts;
3356 }
3357 window = (window * 2).min(len).min(TAIL_SCAN_LIMIT);
3358 }
3359}
3360
3361fn record_timestamp<'a>(record: &'a Value, harness: &str) -> Option<&'a str> {
3365 let key = match harness {
3366 HarnessId::SUPERCODE => "ts",
3367 _ => "timestamp",
3368 };
3369 record.get(key)?.as_str()
3370}
3371
3372fn record_model<'a>(record: &'a Value, harness: &str) -> Option<&'a str> {
3375 match harness {
3376 HarnessId::CLAUDE_CODE | HarnessId::PI => record.get("message")?.get("model")?.as_str(),
3377 HarnessId::CODEX => {
3378 if record.get("type")?.as_str()? != "turn_context" {
3379 return None;
3380 }
3381 record.get("payload")?.get("model")?.as_str()
3382 }
3383 _ => None,
3384 }
3385}
3386
3387const TAIL_SCAN_START: u64 = 16 * 1024;
3390
3391const TAIL_SCAN_LIMIT: u64 = 1024 * 1024;
3394
3395fn modified_ms(path: &Path) -> Option<u64> {
3396 fs::metadata(path)
3397 .ok()?
3398 .modified()
3399 .ok()?
3400 .duration_since(UNIX_EPOCH)
3401 .ok()
3402 .and_then(|duration| u64::try_from(duration.as_millis()).ok())
3403}
3404
3405fn recorded_cwd_matches(recorded: &Path, wanted: &Path) -> bool {
3411 recorded.is_absolute() && same_path(recorded, wanted)
3412}
3413
3414fn same_path(left: &Path, right: &Path) -> bool {
3415 match (fs::canonicalize(left), fs::canonicalize(right)) {
3416 (Ok(left), Ok(right)) => left == right,
3417 _ => normalize_path(left) == normalize_path(right),
3418 }
3419}
3420
3421fn normalize_path(path: &Path) -> PathBuf {
3422 let absolute = if path.is_absolute() {
3423 path.to_path_buf()
3424 } else {
3425 std::env::current_dir()
3426 .unwrap_or_else(|_| PathBuf::from("."))
3427 .join(path)
3428 };
3429 let mut normalized = PathBuf::new();
3430 for component in absolute.components() {
3431 match component {
3432 Component::CurDir => {}
3433 Component::ParentDir => {
3434 normalized.pop();
3435 }
3436 other => normalized.push(other.as_os_str()),
3437 }
3438 }
3439 normalized
3440}
3441
3442#[cfg(test)]
3443mod tests {
3444 use super::*;
3445 use std::io::Write;
3446 use std::time::{SystemTime, UNIX_EPOCH};
3447
3448 fn temp_dir(label: &str) -> PathBuf {
3449 let nonce = SystemTime::now()
3450 .duration_since(UNIX_EPOCH)
3451 .unwrap()
3452 .as_nanos();
3453 let path = std::env::temp_dir().join(format!(
3454 "supercode-catalog-{label}-{}-{nonce}",
3455 std::process::id()
3456 ));
3457 fs::create_dir_all(&path).unwrap();
3458 path
3459 }
3460
3461 #[test]
3462 fn codex_history_index_reads_appends_and_repairs_replacements() {
3463 let root = temp_dir("codex-history-index");
3464 let sessions = root.join("sessions");
3465 fs::create_dir_all(&sessions).unwrap();
3466 let history = root.join("history.jsonl");
3467 fs::write(
3468 &history,
3469 "{\"session_id\":\"alpha\",\"text\":\"first topic\"}\n",
3470 )
3471 .unwrap();
3472
3473 let mut index = CodexHistoryTopicIndex::new(&sessions);
3474 assert_eq!(index.refresh().unwrap(), BTreeSet::from(["alpha".into()]));
3475 assert_eq!(index.topics["alpha"][0].content, "first topic");
3476 assert!(index.refresh().unwrap().is_empty());
3477
3478 let mut file = fs::OpenOptions::new().append(true).open(&history).unwrap();
3479 write!(
3480 file,
3481 "{{\"session_id\":\"alpha\",\"text\":\"later topic\"}}\n\
3482 {{\"session_id\":\"beta\",\"text\":\"second topic\"}}\n"
3483 )
3484 .unwrap();
3485 file.flush().unwrap();
3486 assert_eq!(index.refresh().unwrap(), BTreeSet::from(["beta".into()]));
3487 assert_eq!(index.topics["alpha"][0].content, "first topic");
3488 assert_eq!(index.topics["beta"][0].content, "second topic");
3489
3490 fs::write(
3491 &history,
3492 "{\"session_id\":\"gamma\",\"text\":\"replacement\"}\n",
3493 )
3494 .unwrap();
3495 assert_eq!(
3496 index.refresh().unwrap(),
3497 BTreeSet::from(["alpha".into(), "beta".into(), "gamma".into()])
3498 );
3499 assert!(!index.topics.contains_key("alpha"));
3500 assert_eq!(index.topics["gamma"][0].content, "replacement");
3501
3502 fs::remove_dir_all(root).ok();
3503 }
3504
3505 #[test]
3506 fn codex_history_index_retains_an_incomplete_appended_record() {
3507 let root = temp_dir("codex-history-partial");
3508 let sessions = root.join("sessions");
3509 fs::create_dir_all(&sessions).unwrap();
3510 let history = root.join("history.jsonl");
3511 fs::write(&history, "{\"session_id\":\"partial\",\"text\":\"hel").unwrap();
3512
3513 let mut index = CodexHistoryTopicIndex::new(&sessions);
3514 assert!(index.refresh().unwrap().is_empty());
3515 let mut file = fs::OpenOptions::new().append(true).open(&history).unwrap();
3516 writeln!(file, "lo\"}}").unwrap();
3517 file.flush().unwrap();
3518
3519 assert_eq!(index.refresh().unwrap(), BTreeSet::from(["partial".into()]));
3520 assert_eq!(index.topics["partial"][0].content, "hello");
3521 fs::remove_dir_all(root).ok();
3522 }
3523
3524 #[test]
3525 fn cached_codex_history_enrichment_matches_stateless_discovery() {
3526 let root = temp_dir("codex-history-parity");
3527 let sessions = root.join("sessions");
3528 let workspace = root.join("workspace");
3529 fs::create_dir_all(&sessions).unwrap();
3530 fs::create_dir_all(&workspace).unwrap();
3531 fs::write(
3532 sessions.join("rollout.jsonl"),
3533 format!(
3534 "{{\"type\":\"session_meta\",\"payload\":{{\"id\":\"alpha\",\"cwd\":{}}}}}\n{{\"type\":\"turn_context\",\"payload\":{{\"cwd\":{},\"model\":\"gpt-test\"}}}}\n{{\"type\":\"response_item\",\"payload\":{{\"type\":\"message\",\"role\":\"user\",\"content\":[{{\"type\":\"input_text\",\"text\":\"transcript fallback\"}}]}}}}\n",
3535 serde_json::to_string(&workspace.to_string_lossy()).unwrap(),
3536 serde_json::to_string(&workspace.to_string_lossy()).unwrap(),
3537 ),
3538 )
3539 .unwrap();
3540 fs::write(
3541 root.join("history.jsonl"),
3542 "{\"session_id\":\"alpha\",\"text\":\"history topic\"}\n",
3543 )
3544 .unwrap();
3545 let query = DiscoveryQuery {
3546 harnesses: vec![HarnessId::from(HarnessId::CODEX)],
3547 homes: HarnessHomes {
3548 codex: sessions.clone(),
3549 ..HarnessHomes::default()
3550 },
3551 include_topic_candidates: true,
3552 ..DiscoveryQuery::default()
3553 };
3554 let catalog = HarnessCatalog::new();
3555 let projected = catalog
3556 .project_index(&query, catalog.discover_raw_index(&query))
3557 .unwrap();
3558 let expected = catalog
3559 .enrich_index_page(&query, projected.clone())
3560 .unwrap();
3561 let mut history = CodexHistoryTopicIndex::new(&sessions);
3562 history.refresh().unwrap();
3563 let actual = catalog
3564 .enrich_index_page_with_codex_history(&query, projected, &history)
3565 .unwrap();
3566
3567 assert_eq!(actual, expected);
3568 assert_eq!(actual[0].preview_candidates[0].content, "history topic");
3569 fs::remove_dir_all(root).ok();
3570 }
3571
3572 #[test]
3573 fn locator_json_round_trip_preserves_sqlite_selector() {
3574 let locator = SessionLocator {
3575 harness: HarnessId::from(HarnessId::OPENCODE),
3576 session_id: "ses_123".into(),
3577 storage: StorageLocator::Sqlite {
3578 path: PathBuf::from("/tmp/opencode-dev.db"),
3579 selector: "ses_123".into(),
3580 },
3581 };
3582 let encoded = serde_json::to_string(&locator).unwrap();
3583 assert_eq!(
3584 serde_json::from_str::<SessionLocator>(&encoded).unwrap(),
3585 locator
3586 );
3587 }
3588
3589 #[test]
3590 fn discovers_filters_loads_and_follows_three_jsonl_harnesses() {
3591 let root = temp_dir("jsonl");
3592 let workspace = root.join("workspace");
3593 let other = root.join("other");
3594 fs::create_dir_all(&workspace).unwrap();
3595 fs::create_dir_all(&other).unwrap();
3596
3597 let claude = root.join("claude");
3598 let codex = root.join("codex");
3599 let pi = root.join("pi");
3600 fs::create_dir_all(&claude).unwrap();
3601 fs::create_dir_all(&codex).unwrap();
3602 fs::create_dir_all(&pi).unwrap();
3603 fs::write(
3604 claude.join("claude.jsonl"),
3605 format!(
3606 "{{\"type\":\"user\",\"sessionId\":\"cc-1\",\"cwd\":{},\"timestamp\":\"2026-01-01T00:00:01Z\",\"message\":{{\"role\":\"user\",\"content\":\"hi\"}}}}\n",
3607 serde_json::to_string(&workspace.to_string_lossy()).unwrap()
3608 ),
3609 )
3610 .unwrap();
3611 fs::write(
3612 codex.join("rollout.jsonl"),
3613 format!(
3614 "{{\"timestamp\":\"2026-01-01T00:00:00Z\",\"type\":\"session_meta\",\"payload\":{{\"id\":\"cx-1\",\"cwd\":{}}}}}\n{{\"timestamp\":\"2026-01-01T00:00:02Z\",\"type\":\"response_item\",\"payload\":{{\"type\":\"message\",\"role\":\"user\",\"content\":[{{\"type\":\"input_text\",\"text\":\"inspect codex\"}}]}}}}\n",
3615 serde_json::to_string(&workspace.to_string_lossy()).unwrap()
3616 ),
3617 )
3618 .unwrap();
3619 fs::write(
3620 pi.join("pi.jsonl"),
3621 format!(
3622 "{{\"type\":\"session\",\"version\":3,\"id\":\"pi-1\",\"timestamp\":\"2026-01-01T00:00:00Z\",\"cwd\":{}}}\n{{\"type\":\"message\",\"message\":{{\"role\":\"user\",\"content\":\"inspect pi\"}}}}\n",
3623 serde_json::to_string(&workspace.to_string_lossy()).unwrap()
3624 ),
3625 )
3626 .unwrap();
3627 fs::write(
3628 pi.join("unrelated.jsonl"),
3629 format!(
3630 "{{\"type\":\"session\",\"version\":3,\"id\":\"pi-2\",\"timestamp\":\"2026-01-01T00:00:00Z\",\"cwd\":{}}}\n",
3631 serde_json::to_string(&other.to_string_lossy()).unwrap()
3632 ),
3633 )
3634 .unwrap();
3635 fs::write(claude.join("partial.jsonl"), "{truncated").unwrap();
3636
3637 let query = DiscoveryQuery {
3638 workspace: Some(workspace),
3639 homes: HarnessHomes {
3640 claude_code: claude,
3641 codex,
3642 pi,
3643 opencode: root.join("missing-opencode"),
3644 grok: root.join("missing-grok"),
3645 gemini: root.join("missing-gemini"),
3646 goose: root.join("missing-goose"),
3647 supercode: root.join("missing-supercode"),
3648 openclaw: root.join("missing-openclaw"),
3649 hermes: root.join("missing-hermes"),
3650 orchestrator: root.join("missing-orchestrator"),
3651 },
3652 ..DiscoveryQuery::default()
3653 };
3654 let catalog = HarnessCatalog::new();
3655 let found = catalog.discover(&query).unwrap();
3656 assert_eq!(found.len(), 3);
3657 assert_eq!(
3658 found
3659 .iter()
3660 .map(|item| item.locator.harness.as_str())
3661 .collect::<HashSet<_>>(),
3662 HashSet::from([HarnessId::CLAUDE_CODE, HarnessId::CODEX, HarnessId::PI])
3663 );
3664 for descriptor in found {
3665 assert!(descriptor.preview_candidates.is_empty());
3666 assert_eq!(descriptor.latest_message_candidates.len(), 1);
3667 assert_eq!(descriptor.latest_message_candidates[0].role, "user");
3668 assert!(descriptor.latest_message_candidates[0].cursor.is_some());
3669 if descriptor.locator.harness.as_str() == HarnessId::CLAUDE_CODE {
3670 assert_eq!(
3671 descriptor.latest_message_candidates[0]
3672 .metadata
3673 .get("timestamp")
3674 .map(String::as_str),
3675 Some("2026-01-01T00:00:01Z")
3676 );
3677 } else if descriptor.locator.harness.as_str() == HarnessId::CODEX {
3678 assert_eq!(
3679 descriptor.latest_message_candidates[0]
3680 .metadata
3681 .get("timestamp")
3682 .map(String::as_str),
3683 Some("2026-01-01T00:00:02Z")
3684 );
3685 }
3686 let loaded = catalog.load(&descriptor.locator).unwrap();
3687 assert_eq!(
3688 loaded.meta.session_id.as_deref(),
3689 Some(descriptor.locator.session_id.as_str())
3690 );
3691 let mut follower = catalog.follow(&descriptor.locator).unwrap();
3692 assert!(matches!(
3693 follower.poll().unwrap(),
3694 Some(crate::SessionWatchEvent::SessionSnapshot { .. })
3695 ));
3696 }
3697 fs::remove_dir_all(root).ok();
3698 }
3699
3700 #[test]
3701 fn codex_event_messages_supply_bounded_native_order_previews() {
3702 let root = temp_dir("codex-event-previews");
3705 let path = root.join("rollout.jsonl");
3706 let user = serde_json::json!({
3707 "timestamp": "2026-01-01T00:00:01Z", "type": "event_msg",
3708 "payload": {"type": "user_message", "message": "Investigate the worker"}
3709 });
3710 let answer = serde_json::json!({
3711 "timestamp": "2026-01-01T00:00:02Z", "type": "event_msg",
3712 "payload": {"type": "agent_message", "message": "Worker findings"}
3713 });
3714 let noise = serde_json::json!({
3715 "type": "event_msg", "payload": {"type": "token_count", "message": "not a message"}
3716 });
3717 fs::write(&path, format!("{user}\n{answer}\n{noise}\n{{partial")).unwrap();
3718 let latest = latest_file_message_candidates(&path, HarnessId::CODEX).unwrap();
3719 assert_eq!(latest.len(), 2);
3720 assert_eq!(latest[0].role, "assistant");
3721 assert_eq!(latest[0].content, "Worker findings");
3722 assert_eq!(latest[1].role, "user");
3723 assert_eq!(latest[1].content, "Investigate the worker");
3724 assert_eq!(
3725 latest[0].metadata.get("timestamp").map(String::as_str),
3726 Some("2026-01-01T00:00:02Z")
3727 );
3728 assert_eq!(
3729 latest[0].cursor.as_deref(),
3730 Some(message_candidate_cursor(HarnessId::CODEX, &answer).as_str())
3731 );
3732 let topics = topic_file_message_candidates(&path, HarnessId::CODEX).unwrap();
3733 assert_eq!(topics.len(), 2);
3734 assert_eq!(topics[0].content, "Investigate the worker");
3735
3736 let mut context_pairs = String::new();
3737 for index in 0..5 {
3738 let content = if index < 4 {
3739 format!("# AGENTS.md instructions for /work/{index}\n\n<INSTRUCTIONS>Context</INSTRUCTIONS>")
3740 } else {
3741 "The actual user request".to_string()
3742 };
3743 let event = serde_json::json!({
3744 "type": "event_msg", "payload": {"type": "user_message", "message": content}
3745 });
3746 let response = serde_json::json!({
3747 "type": "response_item", "payload": {"type": "message", "role": "user", "content": content}
3748 });
3749 context_pairs.push_str(&format!("{event}\n{response}\n"));
3750 }
3751 fs::write(&path, context_pairs).unwrap();
3752 let topics = topic_file_message_candidates(&path, HarnessId::CODEX).unwrap();
3753 assert!(topics
3754 .iter()
3755 .any(|candidate| candidate.content == "The actual user request"));
3756
3757 let mut paired = String::new();
3760 for index in 0..12 {
3761 let content = format!("answer {index}");
3762 let event = serde_json::json!({
3763 "type": "event_msg", "payload": {"type": "agent_message", "message": content}
3764 });
3765 let response = serde_json::json!({
3766 "id": format!("response-{index}"), "type": "response_item",
3767 "payload": {"type": "message", "role": "assistant", "content": content}
3768 });
3769 if index % 2 == 0 {
3770 paired.push_str(&format!("{event}\n{response}\n"));
3771 } else {
3772 paired.push_str(&format!("{response}\n{event}\n"));
3773 }
3774 }
3775 fs::write(&path, paired).unwrap();
3776 let latest = latest_file_message_candidates(&path, HarnessId::CODEX).unwrap();
3777 assert_eq!(latest.len(), LATEST_PREVIEW_CANDIDATES);
3778 assert_eq!(latest[0].content, "answer 11");
3779 assert_eq!(latest[1].content, "answer 10");
3780 assert_eq!(latest[7].content, "answer 4");
3781 for (offset, candidate) in latest.iter().enumerate() {
3782 let response = serde_json::json!({"id": format!("response-{}", 11 - offset)});
3783 assert_eq!(
3784 candidate.cursor.as_deref(),
3785 Some(message_candidate_cursor(HarnessId::CODEX, &response).as_str())
3786 );
3787 }
3788 let topics = topic_file_message_candidates(&path, HarnessId::CODEX).unwrap();
3789 assert_eq!(topics.len(), LATEST_PREVIEW_CANDIDATES);
3790 assert_eq!(topics[7].content, "answer 7");
3791
3792 let event = serde_json::json!({
3793 "type": "event_msg", "payload": {"type": "agent_message", "message": "again"}
3794 });
3795 let response = serde_json::json!({
3796 "type": "response_item", "payload": {"type": "message", "role": "assistant", "content": "again"}
3797 });
3798 for records in [
3799 format!("{event}\n{event}\n"),
3800 format!("{response}\n{response}\n"),
3801 format!("{event}\n{response}\n{event}\n{response}\n"),
3802 format!("{response}\n{event}\n{response}\n{event}\n"),
3803 ] {
3804 fs::write(&path, records).unwrap();
3805 assert_eq!(
3806 latest_file_message_candidates(&path, HarnessId::CODEX)
3807 .unwrap()
3808 .len(),
3809 2
3810 );
3811 assert_eq!(
3812 topic_file_message_candidates(&path, HarnessId::CODEX)
3813 .unwrap()
3814 .len(),
3815 2
3816 );
3817 }
3818
3819 let prefix = "x".repeat(4096);
3821 let distinct_event = serde_json::json!({
3822 "type": "event_msg", "payload": {"type": "agent_message", "message": format!("{prefix}A")}
3823 });
3824 let distinct_response = serde_json::json!({
3825 "type": "response_item", "payload": {"type": "message", "role": "assistant", "content": format!("{prefix}B")}
3826 });
3827 fs::write(&path, format!("{distinct_event}\n{distinct_response}\n")).unwrap();
3828 assert_eq!(
3829 latest_file_message_candidates(&path, HarnessId::CODEX)
3830 .unwrap()
3831 .len(),
3832 2
3833 );
3834 assert_eq!(
3835 topic_file_message_candidates(&path, HarnessId::CODEX)
3836 .unwrap()
3837 .len(),
3838 2
3839 );
3840 let long = serde_json::json!({
3841 "type": "event_msg", "payload": {"type": "agent_message", "message": "x".repeat(5000)}
3842 });
3843 fs::write(&path, format!("{long}\n")).unwrap();
3844 let latest = latest_file_message_candidates(&path, HarnessId::CODEX).unwrap();
3845 assert_eq!(latest[0].content.len(), 4096);
3846 fs::remove_dir_all(root).ok();
3847 }
3848
3849 #[test]
3850 fn preview_cursor_tracks_native_boundary_not_growing_text() {
3851 let first = serde_json::json!({
3852 "timestamp": "2026-01-01T00:00:02Z",
3853 "type": "response_item",
3854 "payload": {"type": "message", "role": "assistant", "content": "partial"}
3855 });
3856 let grown = serde_json::json!({
3857 "timestamp": "2026-01-01T00:00:02Z",
3858 "type": "response_item",
3859 "payload": {"type": "message", "role": "assistant", "content": "partial and complete"}
3860 });
3861 let next = serde_json::json!({
3862 "timestamp": "2026-01-01T00:00:03Z",
3863 "type": "response_item",
3864 "payload": {"type": "message", "role": "assistant", "content": "next"}
3865 });
3866
3867 assert_eq!(
3868 message_candidate_cursor(HarnessId::CODEX, &first),
3869 message_candidate_cursor(HarnessId::CODEX, &grown)
3870 );
3871 assert_ne!(
3872 message_candidate_cursor(HarnessId::CODEX, &first),
3873 message_candidate_cursor(HarnessId::CODEX, &next)
3874 );
3875 }
3876
3877 #[test]
3878 fn codex_child_rollouts_roll_into_roots_before_pagination() {
3879 let root = temp_dir("codex-roots");
3880 let codex = root.join("codex");
3881 fs::create_dir_all(&codex).unwrap();
3882 let write_rollout =
3883 |name: &str, payload: Value, modified_seconds: u64| {
3884 let path = codex.join(format!("{name}.jsonl"));
3885 fs::write(
3886 &path,
3887 format!(
3888 "{}\n",
3889 serde_json::json!({
3890 "timestamp": "2026-01-01T00:00:00Z",
3891 "type": "session_meta",
3892 "payload": payload,
3893 })
3894 ),
3895 )
3896 .unwrap();
3897 File::open(&path)
3898 .unwrap()
3899 .set_times(fs::FileTimes::new().set_modified(
3900 UNIX_EPOCH + std::time::Duration::from_secs(modified_seconds),
3901 ))
3902 .unwrap();
3903 };
3904 write_rollout(
3905 "parent",
3906 serde_json::json!({"id":"parent","cwd":"/project","source":"cli"}),
3907 100,
3908 );
3909 write_rollout(
3910 "other",
3911 serde_json::json!({"id":"other","cwd":"/project","source":"cli"}),
3912 200,
3913 );
3914 write_rollout(
3915 "child",
3916 serde_json::json!({
3917 "id": "child",
3918 "cwd": "/project",
3919 "parent_thread_id": "parent",
3920 "source": {"subagent":{"thread_spawn":{
3921 "parent_thread_id":"parent",
3922 "depth":1,
3923 "agent_path":"/root/reviewer"
3924 }}}
3925 }),
3926 300,
3927 );
3928
3929 let catalog = HarnessCatalog::new();
3930 let query = DiscoveryQuery {
3931 harnesses: vec![HarnessId::from(HarnessId::CODEX)],
3932 homes: HarnessHomes {
3933 codex: codex.clone(),
3934 ..HarnessHomes::default()
3935 },
3936 limit: Some(1),
3937 ..DiscoveryQuery::default()
3938 };
3939 let roots = catalog.discover(&query).unwrap();
3940 assert_eq!(roots.len(), 1);
3941 assert_eq!(roots[0].locator.session_id, "parent");
3942 assert_eq!(roots[0].updated_at_ms, Some(300_000));
3943 assert_eq!(roots[0].parent_session_id, None);
3944 assert_eq!(roots[0].child_session_count, 1);
3945
3946 let tree = catalog
3947 .discover(&DiscoveryQuery {
3948 limit: None,
3949 include_child_sessions: true,
3950 root_session_id: Some("parent".into()),
3951 ..query
3952 })
3953 .unwrap();
3954 assert_eq!(tree.len(), 2);
3955 assert!(tree
3956 .iter()
3957 .all(|descriptor| descriptor.locator.session_id != "other"));
3958 let child = tree
3959 .iter()
3960 .find(|descriptor| descriptor.locator.session_id == "child")
3961 .unwrap();
3962 assert_eq!(child.parent_session_id.as_deref(), Some("parent"));
3963 fs::remove_dir_all(root).ok();
3964 }
3965
3966 #[test]
3967 fn discovers_loads_and_follows_opencode_sqlite() {
3968 let db = PathBuf::from(env!("CARGO_MANIFEST_DIR"))
3969 .join("../harness/tests/fixtures/opencode_fixture/opencode.db");
3970 let catalog = HarnessCatalog::new();
3971 let found = catalog
3972 .discover(&DiscoveryQuery {
3973 harnesses: vec![HarnessId::from(HarnessId::OPENCODE)],
3974 homes: HarnessHomes {
3975 opencode: db,
3976 ..HarnessHomes::default()
3977 },
3978 ..DiscoveryQuery::default()
3979 })
3980 .unwrap();
3981 assert!(!found.is_empty());
3982 for descriptor in found {
3983 assert_eq!(descriptor.locator.harness.as_str(), HarnessId::OPENCODE);
3984 assert_eq!(
3985 catalog.load(&descriptor.locator).unwrap().meta.session_id,
3986 Some(descriptor.locator.session_id.clone())
3987 );
3988 assert!(catalog.follow(&descriptor.locator).is_ok());
3989 }
3990 }
3991
3992 #[test]
3993 fn discovers_loads_and_follows_hermes_sqlite_by_session() {
3994 let root = temp_dir("hermes-follow");
3996 fs::create_dir_all(&root).unwrap();
3997 let db = root.join("state.db");
3998 fs::copy(
3999 PathBuf::from(env!("CARGO_MANIFEST_DIR"))
4000 .join("../harness/tests/fixtures/hermes_home/state.db"),
4001 &db,
4002 )
4003 .unwrap();
4004 let catalog = HarnessCatalog::new();
4005 let query = DiscoveryQuery {
4006 harnesses: vec![HarnessId::from(HarnessId::HERMES)],
4007 homes: HarnessHomes {
4008 hermes: db.clone(),
4009 ..HarnessHomes::default()
4010 },
4011 ..DiscoveryQuery::default()
4012 };
4013 let found = catalog.discover(&query).unwrap();
4014 assert!(found.len() >= 2, "{found:#?}");
4015 for descriptor in &found {
4019 assert_eq!(descriptor.locator.harness.as_str(), HarnessId::HERMES);
4020 let loaded = catalog.load(&descriptor.locator).unwrap();
4021 let last_text = loaded.messages.iter().rev().find_map(|message| {
4022 (matches!(message.role, crate::Role::User | crate::Role::Assistant))
4023 .then(|| message.content.clone())
4024 .flatten()
4025 });
4026 assert_eq!(
4028 descriptor
4029 .latest_message_candidates
4030 .first()
4031 .map(|c| c.content.as_str()),
4032 last_text.as_deref(),
4033 "{}",
4034 descriptor.locator.session_id
4035 );
4036 assert_eq!(
4037 catalog.load(&descriptor.locator).unwrap().meta.session_id,
4038 Some(descriptor.locator.session_id.clone())
4039 );
4040 let mut follower = catalog.follow(&descriptor.locator).unwrap();
4041 match follower.poll().unwrap() {
4042 Some(crate::watch::SessionWatchEvent::SessionSnapshot { session, .. }) => {
4043 assert_eq!(
4044 session.meta.session_id,
4045 Some(descriptor.locator.session_id.clone())
4046 );
4047 }
4048 other => panic!("expected an initial snapshot, got {other:?}"),
4049 }
4050 }
4051 let target = &found[0].locator;
4054 let sibling = &found[1].locator;
4055 let mut target_follower = catalog.follow(target).unwrap();
4056 let mut sibling_follower = catalog.follow(sibling).unwrap();
4057 target_follower.poll().unwrap();
4058 sibling_follower.poll().unwrap();
4059 std::thread::sleep(std::time::Duration::from_millis(20));
4060 {
4061 let conn = rusqlite::Connection::open(&db).unwrap();
4062 conn.execute(
4063 "INSERT INTO messages (session_id, role, content, timestamp, active) VALUES (?1, 'assistant', 'appended by the follow test', ?2, 1)",
4064 rusqlite::params![target.session_id, 1_800_000_000.0_f64],
4065 )
4066 .unwrap();
4067 }
4068 match target_follower.poll().unwrap() {
4069 Some(crate::watch::SessionWatchEvent::MessagesAppended {
4070 session_id,
4071 messages,
4072 ..
4073 }) => {
4074 assert_eq!(session_id, Some(target.session_id.clone()));
4075 assert_eq!(messages.len(), 1);
4076 assert_eq!(
4077 messages[0].content.as_deref(),
4078 Some("appended by the follow test")
4079 );
4080 }
4081 other => panic!("expected messages_appended for the target session, got {other:?}"),
4082 }
4083 assert!(
4084 sibling_follower.poll().unwrap().is_none(),
4085 "the sibling session must not wake"
4086 );
4087 fs::remove_dir_all(&root).ok();
4088 }
4089
4090 #[test]
4091 fn discovers_loads_and_follows_gemini_conversation_records() {
4092 let root = temp_dir("gemini");
4093 let workspace = root.join("workspace");
4094 let chats = root.join("gemini/tmp/demo/chats");
4095 fs::create_dir_all(&workspace).unwrap();
4096 fs::create_dir_all(&chats).unwrap();
4097 fs::write(
4098 root.join("gemini/projects.json"),
4099 serde_json::json!({
4100 "projects": {workspace.to_string_lossy(): "demo"}
4101 })
4102 .to_string(),
4103 )
4104 .unwrap();
4105 let transcript = chats.join("gemini-id.jsonl");
4106 fs::write(
4107 &transcript,
4108 include_str!("../../harness/tests/fixtures/gemini_session.jsonl"),
4109 )
4110 .unwrap();
4111
4112 let catalog = HarnessCatalog::new();
4113 let found = catalog
4114 .discover(&DiscoveryQuery {
4115 harnesses: vec![HarnessId::from(HarnessId::GEMINI)],
4116 homes: HarnessHomes {
4117 gemini: root.join("gemini"),
4118 ..HarnessHomes::default()
4119 },
4120 workspace: Some(workspace.clone()),
4121 ..DiscoveryQuery::default()
4122 })
4123 .unwrap();
4124
4125 assert_eq!(found.len(), 1);
4126 assert_eq!(found[0].cwd.as_deref(), Some(workspace.as_path()));
4127 assert_eq!(found[0].message_count, None);
4128 assert_eq!(found[0].model.as_deref(), Some("gemini-2.5-pro"));
4129 assert_eq!(found[0].title, None);
4130 assert!(found[0].preview_candidates.is_empty());
4131 assert_eq!(found[0].latest_message_candidates.len(), 3);
4132 assert_eq!(
4133 found[0].latest_message_candidates[0].content,
4134 "Fixture inspected."
4135 );
4136 let loaded = catalog.load(&found[0].locator).unwrap();
4137 assert_eq!(
4138 loaded.meta.session_id.as_deref(),
4139 Some("11111111-1111-4111-8111-111111111111")
4140 );
4141 assert_eq!(loaded.messages.len(), 4);
4142 assert!(matches!(
4143 catalog.follow(&found[0].locator).unwrap().poll().unwrap(),
4144 Some(crate::SessionWatchEvent::SessionSnapshot { .. })
4145 ));
4146 fs::remove_dir_all(root).ok();
4147 }
4148
4149 #[test]
4150 fn preview_search_filters_before_pagination_without_changing_metadata_search() {
4151 let root = temp_dir("preview-search");
4154 for (id, first, last) in [
4155 ("topic-hit", "NEBULA opening", "Finished"),
4156 ("latest-hit", "Ordinary opening", "Found the nebula"),
4157 ("no-hit", "Unrelated opening", "Finished"),
4158 ] {
4159 fs::write(root.join(format!("{id}.jsonl")), format!("{}\n{}\n",
4160 serde_json::json!({"sessionId": id, "cwd": "/work", "type": "user", "message": {"role": "user", "content": first}}),
4161 serde_json::json!({"sessionId": id, "type": "assistant", "message": {"role": "assistant", "content": last}}),
4162 )).unwrap();
4163 }
4164 let query: DiscoveryQuery = serde_json::from_value(serde_json::json!({
4165 "harnesses": ["claude-code"], "homes": {"claude_code": root},
4166 "query": " nebula ", "search_previews": true, "limit": 1
4167 }))
4168 .unwrap();
4169 let catalog = HarnessCatalog::new();
4170 let first = catalog.discover_page(&query).unwrap();
4171 assert!(first.receipt.searched_previews);
4172 assert_eq!(first.receipt.total_matched, 2);
4173 assert_eq!(first.sessions.len(), 1);
4174 assert!(first.receipt.truncated);
4175 let second = catalog
4176 .discover_page(&DiscoveryQuery {
4177 cursor: first.next_cursor.clone(),
4178 ..query.clone()
4179 })
4180 .unwrap();
4181 assert_eq!(second.receipt.total_matched, 2);
4182 assert_eq!(second.sessions.len(), 1);
4183 assert_ne!(first.sessions[0].locator, second.sessions[0].locator);
4184 assert!(!second.receipt.truncated);
4185 let mut metadata = serde_json::to_value(&query).unwrap();
4186 metadata["search_previews"] = false.into();
4187 let metadata_page = catalog
4188 .discover_page(&serde_json::from_value(metadata).unwrap())
4189 .unwrap();
4190 assert!(metadata_page.sessions.is_empty());
4191 assert!(serde_json::to_value(&metadata_page.receipt)
4192 .unwrap()
4193 .get("searched_previews")
4194 .is_none());
4195
4196 let all = catalog
4199 .discover_page(&DiscoveryQuery {
4200 query: Some("hit".into()),
4201 limit: None,
4202 ..query.clone()
4203 })
4204 .unwrap();
4205 assert_eq!(all.sessions.len(), 3);
4206 assert_eq!(all.receipt.total_matched, 3);
4207 let elsewhere = catalog
4208 .discover_page(&DiscoveryQuery {
4209 workspace: Some("/elsewhere".into()),
4210 ..query.clone()
4211 })
4212 .unwrap();
4213 assert_eq!(elsewhere.receipt.total_matched, 0);
4214 let excluded_by_time = catalog
4215 .discover_page(&DiscoveryQuery {
4216 updated_after_ms: Some(u64::MAX),
4217 ..query.clone()
4218 })
4219 .unwrap();
4220 assert_eq!(excluded_by_time.receipt.total_matched, 0);
4221 for invalid in [
4222 DiscoveryQuery {
4223 query: None,
4224 ..query.clone()
4225 },
4226 DiscoveryQuery {
4227 query: Some(" ".into()),
4228 ..query.clone()
4229 },
4230 DiscoveryQuery {
4231 limit: Some(0),
4232 ..query.clone()
4233 },
4234 DiscoveryQuery {
4235 cursor: Some("bad-cursor".into()),
4236 ..query.clone()
4237 },
4238 DiscoveryQuery {
4239 cursor: first.next_cursor,
4240 query: Some("absent".into()),
4241 ..query.clone()
4242 },
4243 ] {
4244 assert!(catalog.discover_page(&invalid).is_err());
4245 }
4246 assert!(catalog.project_index_page(&query, Vec::new()).is_err());
4247 fs::remove_dir_all(root).ok();
4248 }
4249
4250 #[test]
4251 fn preview_search_uses_codex_first_history_topic_and_bounded_candidates() {
4252 let root = temp_dir("preview-search-codex");
4253 let sessions = root.join("sessions");
4254 fs::create_dir_all(&sessions).unwrap();
4255 fs::write(
4256 root.join("history.jsonl"),
4257 format!(
4258 "{}\n{}\n",
4259 serde_json::json!({"session_id": "history-hit", "text": "Original nebula topic"}),
4260 serde_json::json!({"session_id": "history-hit", "text": "laterhistoryonly"}),
4261 ),
4262 )
4263 .unwrap();
4264 for id in ["history-hit", "latest-hit", "bounded"] {
4265 let mut content = format!(
4266 "{}\n",
4267 serde_json::json!({
4268 "type": "session_meta", "payload": {"id": id, "cwd": "/work"}
4269 })
4270 );
4271 for index in 0..20 {
4272 let message = if id == "latest-hit" && index == 19 {
4273 "Found NEBULA".to_string()
4274 } else if index == 10 {
4275 "middlehistoryonly".to_string()
4276 } else {
4277 format!("{}beyondtextcap", "x".repeat(4096))
4278 };
4279 content.push_str(&format!("{}\n", serde_json::json!({
4280 "type": "event_msg", "payload": {"type": "agent_message", "message": message}
4281 })));
4282 }
4283 fs::write(sessions.join(format!("{id}.jsonl")), content).unwrap();
4284 }
4285 let catalog = HarnessCatalog::new();
4286 let query: DiscoveryQuery = serde_json::from_value(serde_json::json!({
4287 "harnesses": ["codex"], "homes": {"codex": sessions},
4288 "query": "nebula", "search_previews": true
4289 }))
4290 .unwrap();
4291 let page = catalog.discover_page(&query).unwrap();
4292 assert_eq!(page.receipt.total_matched, 2);
4293 for row in &page.sessions {
4294 assert!(row.preview_candidates.len() <= 8);
4295 assert!(row.latest_message_candidates.len() <= 8);
4296 assert!(row
4297 .preview_candidates
4298 .iter()
4299 .chain(&row.latest_message_candidates)
4300 .all(|candidate| candidate.content.chars().count() <= 4096));
4301 }
4302 for text in ["middlehistoryonly", "laterhistoryonly", "beyondtextcap"] {
4303 assert!(
4304 catalog
4305 .discover_page(&DiscoveryQuery {
4306 query: Some(text.into()),
4307 ..query.clone()
4308 })
4309 .unwrap()
4310 .sessions
4311 .is_empty(),
4312 "not a full-history search: {text}"
4313 );
4314 }
4315 fs::remove_dir_all(root).unwrap();
4316 }
4317
4318 #[test]
4319 fn discovers_native_store_and_pages_search_results() {
4320 let root = temp_dir("supercode");
4321 let store_root = root.join("sessions");
4322 fs::create_dir_all(&store_root).unwrap();
4323 for (name, title) in [
4324 ("alpha", "Alpha planning"),
4325 ("beta", "Beta implementation"),
4326 ("gamma", "Gamma review"),
4327 ] {
4328 fs::write(
4329 store_root.join(format!("{name}.jsonl")),
4330 format!("{{\"role\":\"user\",\"content\":\"{title}\"}}\n"),
4331 )
4332 .unwrap();
4333 fs::write(
4334 store_root.join(format!("{name}.meta.json")),
4335 serde_json::json!({"name": name, "title": title}).to_string(),
4336 )
4337 .unwrap();
4338 }
4339 let catalog = HarnessCatalog::new();
4340 let base = DiscoveryQuery {
4341 harnesses: vec![HarnessId::from(HarnessId::SUPERCODE)],
4342 homes: HarnessHomes {
4343 supercode: store_root,
4344 ..HarnessHomes::default()
4345 },
4346 limit: Some(1),
4347 ..DiscoveryQuery::default()
4348 };
4349
4350 let first = catalog.discover_page(&base).unwrap();
4351 assert_eq!(first.sessions.len(), 1);
4352 assert!(first.next_cursor.is_some());
4353 let second = catalog
4354 .discover_page(&DiscoveryQuery {
4355 cursor: first.next_cursor,
4356 ..base.clone()
4357 })
4358 .unwrap();
4359 assert_eq!(second.sessions.len(), 1);
4360 assert_ne!(
4361 first.sessions[0].locator.session_id,
4362 second.sessions[0].locator.session_id
4363 );
4364 let search = catalog
4365 .discover_page(&DiscoveryQuery {
4366 limit: None,
4367 query: Some("implementation".into()),
4368 ..base
4369 })
4370 .unwrap();
4371 assert_eq!(search.sessions.len(), 1);
4372 assert_eq!(search.sessions[0].locator.session_id, "beta");
4373 assert_eq!(search.sessions[0].message_count, None);
4374 assert_eq!(
4375 catalog
4376 .load(&search.sessions[0].locator)
4377 .unwrap()
4378 .messages
4379 .len(),
4380 1
4381 );
4382 fs::remove_dir_all(root).ok();
4383 }
4384
4385 #[test]
4386 fn native_workspace_discovery_reads_bounded_sidecar_headers() {
4387 let root = temp_dir("supercode-bounded-header");
4388 let store_root = root.join("sessions");
4389 let workspace = root.join("project");
4390 fs::create_dir_all(&store_root).unwrap();
4391 fs::create_dir_all(&workspace).unwrap();
4392 let name = "bounded-native";
4393 fs::write(
4394 store_root.join(format!("{name}.meta.json")),
4395 serde_json::json!({"name": name, "title": "Bounded native"}).to_string(),
4396 )
4397 .unwrap();
4398 fs::write(
4399 store_root.join(format!("{name}.jsonl")),
4400 "{\"role\":\"user\",\"content\":\"projected view\"}\n",
4401 )
4402 .unwrap();
4403 let sidecar = [
4404 serde_json::json!({
4405 "supercode_native": 2,
4406 "source": "claude_code",
4407 "session_id": "native-session"
4408 })
4409 .to_string(),
4410 serde_json::json!({
4411 "type": "user",
4412 "sessionId": "native-session",
4413 "cwd": workspace,
4414 "message": {"role": "user", "content": "hello"}
4415 })
4416 .to_string(),
4417 serde_json::json!({
4418 "type": "assistant",
4419 "sessionId": "native-session",
4420 "cwd": workspace,
4421 "message": {"role": "assistant", "model": "claude-sonnet-5", "content": []}
4422 })
4423 .to_string(),
4424 "not-json".into(),
4427 ]
4428 .join("\n");
4429 fs::write(
4430 store_root.join(format!("{name}.sidecar.jsonl")),
4431 format!("{sidecar}\n"),
4432 )
4433 .unwrap();
4434
4435 let found = HarnessCatalog::new()
4436 .discover(&DiscoveryQuery {
4437 workspace: Some(workspace.clone()),
4438 harnesses: vec![HarnessId::from(HarnessId::SUPERCODE)],
4439 homes: HarnessHomes {
4440 supercode: store_root,
4441 ..HarnessHomes::default()
4442 },
4443 ..DiscoveryQuery::default()
4444 })
4445 .unwrap();
4446
4447 assert_eq!(found.len(), 1);
4448 assert_eq!(found[0].cwd.as_deref(), Some(workspace.as_path()));
4449 assert_eq!(found[0].model.as_deref(), Some("claude-sonnet-5"));
4450 assert_eq!(found[0].message_count, None);
4451 fs::remove_dir_all(root).ok();
4452 }
4453
4454 #[test]
4455 fn discovers_current_opencode_schema_without_a_session_model_column() {
4456 let root = temp_dir("opencode-current");
4457 let db = root.join("opencode.db");
4458 let conn = Connection::open(&db).unwrap();
4459 conn.execute_batch(
4460 "CREATE TABLE session (
4461 id TEXT PRIMARY KEY,
4462 directory TEXT NOT NULL,
4463 title TEXT NOT NULL,
4464 time_updated INTEGER NOT NULL
4465 );
4466 CREATE TABLE message (
4467 id TEXT PRIMARY KEY,
4468 session_id TEXT NOT NULL
4469 );
4470 INSERT INTO session VALUES ('ses_current', '/tmp/work', 'Current', 42);
4471 INSERT INTO message VALUES ('msg_current', 'ses_current');",
4472 )
4473 .unwrap();
4474 drop(conn);
4475
4476 let found = HarnessCatalog::new()
4477 .discover(&DiscoveryQuery {
4478 harnesses: vec![HarnessId::from(HarnessId::OPENCODE)],
4479 homes: HarnessHomes {
4480 opencode: db,
4481 ..HarnessHomes::default()
4482 },
4483 ..DiscoveryQuery::default()
4484 })
4485 .unwrap();
4486
4487 assert_eq!(found.len(), 1);
4488 assert_eq!(found[0].locator.session_id, "ses_current");
4489 assert_eq!(found[0].message_count, Some(1));
4490 assert_eq!(found[0].model, None);
4491 fs::remove_dir_all(root).ok();
4492 }
4493
4494 #[test]
4495 fn workspace_filter_never_matches_a_relative_recorded_cwd() {
4496 let root = temp_dir("opencode-relative-cwd");
4501 let db = root.join("opencode.db");
4502 let conn = Connection::open(&db).unwrap();
4503 let here = std::env::current_dir().unwrap();
4504 conn.execute_batch(&format!(
4505 "CREATE TABLE session (
4506 id TEXT PRIMARY KEY,
4507 directory TEXT NOT NULL,
4508 title TEXT NOT NULL,
4509 time_updated INTEGER NOT NULL
4510 );
4511 CREATE TABLE message (
4512 id TEXT PRIMARY KEY,
4513 session_id TEXT NOT NULL
4514 );
4515 INSERT INTO session VALUES ('ses_relative', '.', 'Ghost', 41);
4516 INSERT INTO session VALUES ('ses_here', '{}', 'Real', 42);",
4517 here.display()
4518 ))
4519 .unwrap();
4520 drop(conn);
4521
4522 let found = HarnessCatalog::new()
4523 .discover(&DiscoveryQuery {
4524 workspace: Some(here),
4525 harnesses: vec![HarnessId::from(HarnessId::OPENCODE)],
4526 homes: HarnessHomes {
4527 opencode: db,
4528 ..HarnessHomes::default()
4529 },
4530 ..DiscoveryQuery::default()
4531 })
4532 .unwrap();
4533
4534 assert_eq!(found.len(), 1);
4535 assert_eq!(found[0].locator.session_id, "ses_here");
4536 fs::remove_dir_all(root).ok();
4537 }
4538}