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