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