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