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