1use std::collections::{BTreeSet, HashMap, HashSet};
7use std::fs::{self, File};
8use std::io::{BufRead, BufReader, Read, Seek, SeekFrom};
9use std::path::{Component, Path, PathBuf};
10use std::time::UNIX_EPOCH;
11
12use rusqlite::Connection;
13use serde::{Deserialize, Serialize};
14use serde_json::Value;
15
16use crate::native_store::load_native_store_family;
17use crate::ontology::{Binding, OrchestratorBindingRow};
18use crate::session::{
19 hermes_capture_nouns, openclaw_agent_id_from_path, openclaw_capture_header_nouns,
20 percent_decode_path, OrchestrationNouns, SessionMeta, SessionSource,
21};
22use crate::{Error, Fidelity, Result, Session, SessionFollower};
23
24pub use crate::ontology::HarnessId;
25
26#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
28#[serde(tag = "kind", rename_all = "snake_case")]
29pub enum StorageLocator {
30 File {
32 path: PathBuf,
34 },
35 Sqlite {
37 path: PathBuf,
39 selector: String,
41 },
42}
43
44impl StorageLocator {
45 pub fn path(&self) -> &Path {
47 match self {
48 Self::File { path } | Self::Sqlite { path, .. } => path,
49 }
50 }
51}
52
53#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
55pub struct SessionLocator {
56 pub harness: HarnessId,
58 pub session_id: String,
60 pub storage: StorageLocator,
62}
63
64#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
66pub struct SessionDescriptor {
67 pub locator: SessionLocator,
70 pub cwd: Option<PathBuf>,
72 pub title: Option<String>,
74 #[serde(default, skip_serializing_if = "Vec::is_empty")]
79 pub preview_candidates: Vec<SessionPreviewCandidate>,
80 #[serde(default, skip_serializing_if = "Vec::is_empty")]
85 pub latest_message_candidates: Vec<SessionPreviewCandidate>,
86 pub updated_at_ms: Option<u64>,
95 pub message_count: Option<usize>,
100 pub model: Option<String>,
102 #[serde(default, skip_serializing_if = "Option::is_none")]
106 pub parent_session_id: Option<String>,
107 #[serde(default, skip_serializing_if = "is_zero")]
111 pub child_session_count: usize,
112 #[serde(flatten)]
118 pub nouns: OrchestrationNouns,
119}
120
121#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
123pub struct SessionPreviewCandidate {
124 #[serde(default, skip_serializing_if = "Option::is_none")]
128 pub cursor: Option<String>,
129 pub role: String,
132 pub content: String,
134 #[serde(default, skip_serializing_if = "HashMap::is_empty")]
136 pub metadata: HashMap<String, String>,
137}
138
139#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
141pub struct DiscoveryPage {
142 pub sessions: Vec<SessionDescriptor>,
144 pub next_cursor: Option<String>,
146 #[serde(default)]
149 pub receipt: DiscoveryReceipt,
150}
151
152#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
154#[serde(default)]
155pub struct DiscoveryReceipt {
156 #[serde(skip_serializing_if = "is_false")]
159 pub searched_previews: bool,
160 #[serde(skip_serializing_if = "Option::is_none")]
162 pub requested_after_ms: Option<u64>,
163 #[serde(skip_serializing_if = "Option::is_none")]
165 pub requested_before_ms: Option<u64>,
166 #[serde(skip_serializing_if = "Option::is_none")]
168 pub requested_limit: Option<usize>,
169 #[serde(skip_serializing_if = "Option::is_none")]
171 pub oldest_returned_ms: Option<u64>,
172 #[serde(skip_serializing_if = "Option::is_none")]
174 pub newest_returned_ms: Option<u64>,
175 pub returned: usize,
177 pub total_matched: usize,
179 pub truncated: bool,
183}
184
185#[derive(Debug)]
192pub struct CodexHistoryTopicIndex {
193 path: PathBuf,
194 fingerprint: Option<CodexHistoryFingerprint>,
195 offset: u64,
196 trailing: Vec<u8>,
197 topics: HashMap<String, Vec<SessionPreviewCandidate>>,
198}
199
200#[derive(Debug, Clone, Copy, PartialEq, Eq)]
201struct CodexHistoryFingerprint {
202 len: u64,
203 modified_ns: u128,
204 identity: u128,
205}
206
207impl CodexHistoryTopicIndex {
208 pub fn new(sessions_root: &Path) -> Self {
211 let root = sessions_root.parent().unwrap_or(sessions_root);
212 Self {
213 path: root.join("history.jsonl"),
214 fingerprint: None,
215 offset: 0,
216 trailing: Vec::new(),
217 topics: HashMap::new(),
218 }
219 }
220
221 pub fn path(&self) -> &Path {
223 &self.path
224 }
225
226 pub fn refresh(&mut self) -> Result<BTreeSet<String>> {
232 let metadata = match fs::metadata(&self.path) {
233 Ok(metadata) => metadata,
234 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
235 let changed = self.topics.keys().cloned().collect();
236 self.fingerprint = None;
237 self.offset = 0;
238 self.trailing.clear();
239 self.topics.clear();
240 return Ok(changed);
241 }
242 Err(error) => return Err(error.into()),
243 };
244 let fingerprint = codex_history_fingerprint(&metadata)?;
245 if self.fingerprint == Some(fingerprint) {
246 return Ok(BTreeSet::new());
247 }
248
249 let is_append = self.fingerprint.is_some_and(|previous| {
250 previous.identity == fingerprint.identity
251 && previous.len < fingerprint.len
252 && self.offset <= previous.len
253 });
254 if is_append {
255 let mut file = File::open(&self.path)?;
256 file.seek(SeekFrom::Start(self.offset))?;
257 let mut bytes = Vec::with_capacity(
258 usize::try_from(fingerprint.len.saturating_sub(self.offset)).unwrap_or(0),
259 );
260 file.read_to_end(&mut bytes)?;
261 self.offset = file.stream_position()?;
262 let changed = self.ingest(bytes);
263 self.fingerprint = Some(CodexHistoryFingerprint {
264 len: self.offset,
265 ..fingerprint
266 });
267 return Ok(changed);
268 }
269
270 let previous = std::mem::take(&mut self.topics);
271 let mut file = File::open(&self.path)?;
272 let mut bytes = Vec::with_capacity(usize::try_from(fingerprint.len).unwrap_or(0));
273 file.read_to_end(&mut bytes)?;
274 self.offset = file.stream_position()?;
275 self.trailing.clear();
276 self.ingest(bytes);
277 self.fingerprint = Some(CodexHistoryFingerprint {
278 len: self.offset,
279 ..fingerprint
280 });
281 Ok(changed_topic_ids(&previous, &self.topics))
282 }
283
284 fn ingest(&mut self, bytes: Vec<u8>) -> BTreeSet<String> {
285 let mut input = std::mem::take(&mut self.trailing);
286 input.extend(bytes);
287 let complete_len = input
288 .iter()
289 .rposition(|byte| *byte == b'\n')
290 .map_or(0, |index| index + 1);
291 let mut changed = BTreeSet::new();
292 for line in input[..complete_len].split(|byte| *byte == b'\n') {
293 self.ingest_line(line, &mut changed);
294 }
295 self.trailing.extend_from_slice(&input[complete_len..]);
296 if !self.trailing.is_empty() {
297 let trailing = self.trailing.clone();
298 if self.ingest_line(&trailing, &mut changed) {
299 self.trailing.clear();
300 }
301 }
302 changed
303 }
304
305 fn ingest_line(&mut self, line: &[u8], changed: &mut BTreeSet<String>) -> bool {
306 let Ok(value) = serde_json::from_slice::<Value>(line) else {
307 return false;
308 };
309 let Some(session_id) = value.get("session_id").and_then(Value::as_str) else {
310 return true;
311 };
312 if self.topics.contains_key(session_id) {
313 return true;
314 }
315 let mut candidates = Vec::new();
316 push_message_candidate(&mut candidates, "user", value.get("text"), HashMap::new());
317 if !candidates.is_empty() {
318 self.topics.insert(session_id.to_string(), candidates);
319 changed.insert(session_id.to_string());
320 }
321 true
322 }
323}
324
325#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
327#[serde(default)]
328pub struct HarnessHomes {
329 pub claude_code: PathBuf,
331 pub codex: PathBuf,
333 pub pi: PathBuf,
335 pub opencode: PathBuf,
337 pub grok: PathBuf,
339 pub gemini: PathBuf,
341 pub goose: PathBuf,
343 pub supercode: PathBuf,
345 pub openclaw: PathBuf,
349 pub hermes: PathBuf,
352 pub orchestrator: PathBuf,
358}
359
360pub fn orchestrator_profile_dirs(root: &Path) -> Vec<(String, PathBuf)> {
366 if !root.is_dir() {
367 return Vec::new();
368 }
369 let mut dirs = vec![("default".to_string(), root.to_path_buf())];
370 if let Ok(entries) = fs::read_dir(root.join("profiles")) {
371 let mut named: Vec<(String, PathBuf)> = entries
372 .flatten()
373 .filter(|entry| entry.path().is_dir())
374 .filter_map(|entry| {
375 entry
376 .file_name()
377 .into_string()
378 .ok()
379 .map(|name| (name, entry.path()))
380 })
381 .filter(|(name, _)| !name.starts_with('.'))
385 .collect();
386 named.sort();
387 dirs.extend(named);
388 }
389 dirs
390}
391
392pub fn hermes_session_stores(db_path: &Path) -> Vec<PathBuf> {
396 let mut stores = vec![db_path.to_path_buf()];
397 let (Some(root), Some(name)) = (db_path.parent(), db_path.file_name()) else {
398 return stores;
399 };
400 for (_, dir) in orchestrator_profile_dirs(root).into_iter().skip(1) {
401 stores.push(dir.join(name));
402 }
403 stores
404}
405
406fn claude_live_named(query: &DiscoveryQuery) -> Option<(String, String)> {
409 let wanted = query
410 .query
411 .as_deref()
412 .map(str::trim)
413 .filter(|q| !q.is_empty())?;
414 if !query.harnesses.is_empty()
415 && !query
416 .harnesses
417 .iter()
418 .any(|harness| harness.as_str() == HarnessId::CLAUDE_CODE)
419 {
420 return None;
421 }
422 let registry = query.homes.claude_code.parent()?.join("sessions");
423 std::fs::read_dir(registry)
424 .ok()?
425 .flatten()
426 .find_map(|entry| {
427 let record: Value = serde_json::from_slice(&std::fs::read(entry.path()).ok()?).ok()?;
428 let name = record.get("name")?.as_str()?;
429 let id = record.get("sessionId")?.as_str()?;
430 (name == wanted).then(|| (name.to_string(), id.to_string()))
431 })
432}
433
434impl Default for HarnessHomes {
435 fn default() -> Self {
436 let home = crate::user_home()
437 .map(std::path::PathBuf::into_os_string)
438 .map(PathBuf::from)
439 .unwrap_or_else(|| PathBuf::from("."));
440 let claude_root = std::env::var_os("CLAUDE_CONFIG_DIR")
441 .map(PathBuf::from)
442 .unwrap_or_else(|| home.join(".claude"));
443 let codex_root = std::env::var_os("CODEX_HOME")
444 .map(PathBuf::from)
445 .unwrap_or_else(|| home.join(".codex"));
446 let pi = std::env::var_os("PI_CODING_AGENT_SESSION_DIR")
447 .map(PathBuf::from)
448 .unwrap_or_else(|| {
449 std::env::var_os("PI_CODING_AGENT_DIR")
450 .map(PathBuf::from)
451 .unwrap_or_else(|| home.join(".pi/agent"))
452 .join("sessions")
453 });
454 let opencode = std::env::var_os("OPENCODE_DB")
455 .map(PathBuf::from)
456 .unwrap_or_else(|| {
457 std::env::var_os("XDG_DATA_HOME")
458 .map(PathBuf::from)
459 .unwrap_or_else(|| home.join(".local/share"))
460 .join("opencode")
461 });
462 let grok = std::env::var_os("GROK_HOME")
463 .map(PathBuf::from)
464 .unwrap_or_else(|| home.join(".grok"))
465 .join("sessions");
466 let gemini = std::env::var_os("GEMINI_CLI_HOME")
467 .map(PathBuf::from)
468 .unwrap_or_else(|| home.join(".gemini"));
469 let openclaw = std::env::var_os("OPENCLAW_STATE_DIR")
476 .map(PathBuf::from)
477 .or_else(|| {
478 std::env::var_os("OPENCLAW_HOME").map(|root| PathBuf::from(root).join(".openclaw"))
479 })
480 .unwrap_or_else(|| home.join(".openclaw"));
481 let orchestrator = std::env::var_os("SUPERCODE_ORCHESTRATOR_HOME")
482 .map(PathBuf::from)
483 .unwrap_or_else(|| home.join(".supercode/orchestrator"));
484 let hermes = std::env::var_os("HERMES_HOME")
485 .map(PathBuf::from)
486 .unwrap_or_else(|| home.join(".hermes"))
487 .join("state.db");
488 let goose = std::env::var_os("GOOSE_PATH_ROOT")
489 .map(PathBuf::from)
490 .map(|root| root.join("data/sessions/sessions.db"))
491 .unwrap_or_else(|| {
492 #[cfg(target_os = "macos")]
493 {
494 home.join("Library/Application Support/Block/goose/sessions/sessions.db")
495 }
496 #[cfg(target_os = "windows")]
497 {
498 std::env::var_os("APPDATA")
499 .map(PathBuf::from)
500 .unwrap_or_else(|| home.join("AppData/Roaming"))
501 .join("Block/goose/sessions/sessions.db")
502 }
503 #[cfg(not(any(target_os = "macos", target_os = "windows")))]
504 {
505 std::env::var_os("XDG_DATA_HOME")
506 .map(PathBuf::from)
507 .unwrap_or_else(|| home.join(".local/share"))
508 .join("goose/sessions/sessions.db")
509 }
510 });
511 let supercode = std::env::var_os("SUPERCODE_HOME")
512 .map(PathBuf::from)
513 .unwrap_or_else(|| {
514 std::env::var_os("XDG_CONFIG_HOME")
515 .map(PathBuf::from)
516 .unwrap_or_else(|| home.join(".config"))
517 .join("supercode")
518 })
519 .join("sessions");
520 Self {
521 claude_code: claude_root.join("projects"),
522 codex: codex_root.join("sessions"),
523 gemini,
524 goose,
525 supercode,
526 openclaw,
527 hermes,
528 orchestrator,
529 pi,
530 opencode,
531 grok,
532 }
533 }
534}
535
536fn codex_history_fingerprint(metadata: &fs::Metadata) -> Result<CodexHistoryFingerprint> {
537 let modified_ns = metadata
538 .modified()?
539 .duration_since(UNIX_EPOCH)
540 .map_err(|error| Error::Other(format!("history timestamp predates Unix epoch: {error}")))?
541 .as_nanos();
542 #[cfg(unix)]
543 let identity = {
544 use std::os::unix::fs::MetadataExt;
545 (u128::from(metadata.dev()) << 64) | u128::from(metadata.ino())
546 };
547 #[cfg(not(unix))]
548 let identity = 0;
549 Ok(CodexHistoryFingerprint {
550 len: metadata.len(),
551 modified_ns,
552 identity,
553 })
554}
555
556fn changed_topic_ids(
557 before: &HashMap<String, Vec<SessionPreviewCandidate>>,
558 after: &HashMap<String, Vec<SessionPreviewCandidate>>,
559) -> BTreeSet<String> {
560 before
561 .keys()
562 .chain(after.keys())
563 .filter(|session_id| before.get(*session_id) != after.get(*session_id))
564 .cloned()
565 .collect()
566}
567
568#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
570#[serde(default)]
571pub struct DiscoveryQuery {
572 pub workspace: Option<PathBuf>,
574 #[serde(skip_serializing_if = "std::ops::Not::not")]
582 pub workspace_subtree: bool,
583 pub workspace_family: Option<PathBuf>,
591 pub updated_after_ms: Option<u64>,
593 pub updated_before_ms: Option<u64>,
595 pub harnesses: Vec<HarnessId>,
597 pub homes: HarnessHomes,
599 pub query: Option<String>,
601 pub search_previews: bool,
605 pub cursor: Option<String>,
607 pub limit: Option<usize>,
609 pub include_topic_candidates: bool,
613 pub include_child_sessions: bool,
617 pub root_session_id: Option<String>,
620 pub profile: Option<String>,
625}
626
627#[derive(Debug, Clone, PartialEq, Eq)]
633struct RepoFamily {
634 identity: PathBuf,
637 is_repository: bool,
640 origin_url: Option<String>,
642}
643
644impl RepoFamily {
645 fn of(path: &Path) -> Self {
646 let start = std::fs::canonicalize(path).unwrap_or_else(|_| path.to_path_buf());
647 let mut current = Some(start.as_path());
648 while let Some(dir) = current {
649 let dot_git = dir.join(".git");
650 if dot_git.is_dir() {
651 let identity = std::fs::canonicalize(&dot_git).unwrap_or_else(|_| dot_git.clone());
652 let origin_url = read_origin_url(&identity);
653 return Self {
654 identity,
655 is_repository: true,
656 origin_url,
657 };
658 }
659 if dot_git.is_file() {
660 if let Ok(text) = std::fs::read_to_string(&dot_git) {
664 if let Some(gitdir) = text
665 .lines()
666 .find_map(|line| line.trim().strip_prefix("gitdir:"))
667 {
668 let gitdir = PathBuf::from(gitdir.trim());
669 let gitdir = if gitdir.is_absolute() {
670 gitdir
671 } else {
672 dir.join(gitdir)
673 };
674 let common = gitdir
675 .parent()
676 .filter(|parent| parent.ends_with("worktrees"))
677 .and_then(Path::parent)
678 .map(Path::to_path_buf)
679 .unwrap_or(gitdir);
680 let identity = std::fs::canonicalize(&common).unwrap_or(common);
681 let origin_url = read_origin_url(&identity);
682 return Self {
683 identity,
684 is_repository: true,
685 origin_url,
686 };
687 }
688 }
689 }
690 current = dir.parent();
691 }
692 Self {
693 identity: start,
694 is_repository: false,
695 origin_url: None,
696 }
697 }
698
699 fn joins(&self, other: &Self) -> bool {
703 if self.identity == other.identity {
704 return true;
705 }
706 self.is_repository
707 && other.is_repository
708 && matches!((&self.origin_url, &other.origin_url), (Some(a), Some(b)) if a == b)
709 }
710}
711
712fn read_origin_url(common_dir: &Path) -> Option<String> {
714 let text = std::fs::read_to_string(common_dir.join("config")).ok()?;
715 let mut in_origin = false;
716 for line in text.lines() {
717 let line = line.trim();
718 if line.starts_with('[') {
719 in_origin = line.starts_with("[remote \"origin\"]");
720 continue;
721 }
722 if in_origin {
723 if let Some(value) = line.strip_prefix("url") {
724 let value = value.trim_start();
725 if let Some(url) = value.strip_prefix('=') {
726 let url = url.trim();
727 if !url.is_empty() {
728 return Some(url.to_string());
729 }
730 }
731 }
732 }
733 }
734 None
735}
736
737#[derive(Debug, Default, Clone, Copy)]
740pub struct HarnessCatalog;
741
742impl HarnessCatalog {
743 pub fn new() -> Self {
745 Self
746 }
747
748 pub fn discover(&self, query: &DiscoveryQuery) -> Result<Vec<SessionDescriptor>> {
752 Ok(self.discover_page(query)?.sessions)
753 }
754
755 pub fn discover_raw_index(&self, query: &DiscoveryQuery) -> Vec<SessionDescriptor> {
760 self.scan_descriptors(query, true)
761 }
762
763 pub fn project_index(
766 &self,
767 query: &DiscoveryQuery,
768 descriptors: impl IntoIterator<Item = SessionDescriptor>,
769 ) -> Result<Vec<SessionDescriptor>> {
770 Ok(self.project_index_page(query, descriptors)?.sessions)
771 }
772
773 pub fn project_index_page(
776 &self,
777 query: &DiscoveryQuery,
778 descriptors: impl IntoIterator<Item = SessionDescriptor>,
779 ) -> Result<DiscoveryPage> {
780 if query.search_previews {
781 return Err(Error::Other(
782 "preview search requires discover_page, not a metadata-only index projection"
783 .into(),
784 ));
785 }
786 let mut found = descriptors.into_iter().collect::<Vec<_>>();
787 project_descriptors(query, &mut found);
788 let total_matched = found.len();
789 let (sessions, next_cursor) = paginate_descriptors(query, found)?;
790 let receipt = DiscoveryReceipt {
791 searched_previews: false,
792 requested_after_ms: query.updated_after_ms,
793 requested_before_ms: query.updated_before_ms,
794 requested_limit: query.limit,
795 oldest_returned_ms: sessions.iter().filter_map(|s| s.updated_at_ms).min(),
796 newest_returned_ms: sessions.iter().filter_map(|s| s.updated_at_ms).max(),
797 returned: sessions.len(),
798 total_matched,
799 truncated: next_cursor.is_some(),
800 };
801 Ok(DiscoveryPage {
802 sessions,
803 next_cursor,
804 receipt,
805 })
806 }
807
808 pub fn enrich_index_page(
811 &self,
812 query: &DiscoveryQuery,
813 mut sessions: Vec<SessionDescriptor>,
814 ) -> Result<Vec<SessionDescriptor>> {
815 enrich_descriptors(query, &mut sessions, None)?;
816 Ok(sessions)
817 }
818
819 pub fn enrich_index_page_with_codex_history(
825 &self,
826 query: &DiscoveryQuery,
827 mut sessions: Vec<SessionDescriptor>,
828 codex_history: &CodexHistoryTopicIndex,
829 ) -> Result<Vec<SessionDescriptor>> {
830 enrich_descriptors(query, &mut sessions, Some(codex_history))?;
831 Ok(sessions)
832 }
833
834 pub fn discover_page(&self, query: &DiscoveryQuery) -> Result<DiscoveryPage> {
836 if query.search_previews {
837 return self.discover_preview_page(query);
838 }
839 let mut page = self.project_index_page(
840 query,
841 self.scan_descriptors(query, query.include_child_sessions),
842 )?;
843 if query.cursor.is_none() {
847 if let Some((name, id)) = claude_live_named(query) {
848 if let Some(at) = page
849 .sessions
850 .iter()
851 .position(|session| session.locator.session_id == id)
852 {
853 let mut found = page.sessions.swap_remove(at);
854 found.title.get_or_insert(name);
855 page.sessions = vec![found];
856 page.next_cursor = None;
857 }
858 }
859 }
860 if let Some(id) = query.query.as_deref().map(str::trim) {
863 if page
864 .sessions
865 .iter()
866 .any(|session| session.locator.session_id == id)
867 {
868 page.sessions
869 .retain(|session| session.locator.session_id == id);
870 page.next_cursor = None;
871 }
872 }
873 if page.sessions.is_empty() && query.cursor.is_none() {
874 if let Some(named) = self.claude_session_named(query)? {
875 page.sessions.push(named);
876 }
877 }
878 enrich_descriptors(query, &mut page.sessions, None)?;
879 Ok(page)
880 }
881
882 fn claude_session_named(&self, query: &DiscoveryQuery) -> Result<Option<SessionDescriptor>> {
889 let Some(name) = query
890 .query
891 .as_deref()
892 .map(str::trim)
893 .filter(|name| !name.is_empty())
894 else {
895 return Ok(None);
896 };
897 if !query.harnesses.is_empty()
898 && !query
899 .harnesses
900 .iter()
901 .any(|harness| harness.as_str() == HarnessId::CLAUDE_CODE)
902 {
903 return Ok(None);
904 }
905 let mut best: Option<(std::time::SystemTime, String)> = None;
906 let Ok(projects) = std::fs::read_dir(&query.homes.claude_code) else {
907 return Ok(None);
908 };
909 for project in projects.flatten() {
910 let Ok(files) = std::fs::read_dir(project.path()) else {
911 continue;
912 };
913 for file in files.flatten() {
914 let path = file.path();
915 if path.extension().and_then(|value| value.to_str()) != Some("jsonl") {
916 continue;
917 }
918 if !claude_custom_titles(&path)
919 .iter()
920 .any(|title| title == name)
921 {
922 continue;
923 }
924 let modified = file
925 .metadata()
926 .and_then(|meta| meta.modified())
927 .unwrap_or(std::time::UNIX_EPOCH);
928 let Some(id) = path
929 .file_stem()
930 .map(|stem| stem.to_string_lossy().into_owned())
931 else {
932 continue;
933 };
934 if best.as_ref().is_none_or(|(time, _)| modified > *time) {
935 best = Some((modified, id));
936 }
937 }
938 }
939 let Some((_, id)) = best else {
940 return Ok(None);
941 };
942 let mut unnamed = query.clone();
943 unnamed.query = None;
944 unnamed.harnesses = vec![HarnessId::new(HarnessId::CLAUDE_CODE)];
945 unnamed.limit = None;
946 let found = self
947 .project_index_page(
948 &unnamed,
949 self.scan_descriptors(&unnamed, unnamed.include_child_sessions),
950 )?
951 .sessions
952 .into_iter()
953 .find(|descriptor| descriptor.locator.session_id == id);
954 Ok(found.map(|mut descriptor| {
955 descriptor.title = Some(name.to_string());
956 descriptor
957 }))
958 }
959
960 fn discover_preview_page(&self, query: &DiscoveryQuery) -> Result<DiscoveryPage> {
961 let search = query
962 .query
963 .as_deref()
964 .map(str::trim)
965 .filter(|text| !text.is_empty())
966 .ok_or_else(|| Error::Other("search_previews requires a nonempty query".into()))?
967 .to_lowercase();
968 if query.limit == Some(0) {
969 return Err(Error::Other("preview search limit must be positive".into()));
970 }
971 let cursor = query.cursor.as_deref().map(decode_cursor).transpose()?;
972 let mut eligible_query = query.clone();
973 eligible_query.query = None;
974 eligible_query.search_previews = false;
975 eligible_query.include_topic_candidates = true;
976 let mut eligible = self.scan_descriptors(query, query.include_child_sessions);
977 project_descriptors(&eligible_query, &mut eligible);
978 let topics = codex_history_topics(&query.homes.codex, &eligible).unwrap_or_default();
980 let mut sessions = Vec::new();
981 let mut total_matched = 0;
982 let mut cursor_seen = cursor.is_none();
983 let mut more = false;
984 for mut descriptor in eligible {
985 let metadata_match = descriptor_matches(&descriptor, &search);
986 if !metadata_match {
987 let topic = topics
988 .get(&descriptor.locator.session_id)
989 .filter(|_| descriptor.locator.harness.as_str() == HarnessId::CODEX);
990 enrich_descriptor(&eligible_query, &mut descriptor, topic);
991 if !descriptor
992 .preview_candidates
993 .iter()
994 .chain(&descriptor.latest_message_candidates)
995 .any(|candidate| candidate.content.to_lowercase().contains(&search))
996 {
997 continue;
998 }
999 }
1000 total_matched += 1;
1001 if !cursor_seen {
1002 cursor_seen = cursor.as_ref() == Some(&descriptor_cursor_key(&descriptor));
1003 continue;
1004 }
1005 if sessions.len() < query.limit.unwrap_or(usize::MAX) {
1006 if metadata_match {
1007 let topic = topics
1008 .get(&descriptor.locator.session_id)
1009 .filter(|_| descriptor.locator.harness.as_str() == HarnessId::CODEX);
1010 enrich_descriptor(&eligible_query, &mut descriptor, topic);
1011 }
1012 sessions.push(descriptor);
1013 } else {
1014 more = true;
1015 }
1016 }
1017 if !cursor_seen {
1018 return Err(Error::Other("discovery cursor is stale or invalid".into()));
1019 }
1020 let next_cursor = more.then(|| sessions.last().map(encode_cursor)).flatten();
1021 let receipt = DiscoveryReceipt {
1022 searched_previews: true,
1023 requested_after_ms: query.updated_after_ms,
1024 requested_before_ms: query.updated_before_ms,
1025 requested_limit: query.limit,
1026 oldest_returned_ms: sessions.iter().filter_map(|s| s.updated_at_ms).min(),
1027 newest_returned_ms: sessions.iter().filter_map(|s| s.updated_at_ms).max(),
1028 returned: sessions.len(),
1029 total_matched,
1030 truncated: more,
1031 };
1032 Ok(DiscoveryPage {
1033 sessions,
1034 next_cursor,
1035 receipt,
1036 })
1037 }
1038
1039 fn scan_descriptors(
1040 &self,
1041 query: &DiscoveryQuery,
1042 include_child_sessions: bool,
1043 ) -> Vec<SessionDescriptor> {
1044 let selected: HashSet<&str> = if query.harnesses.is_empty() {
1045 [
1046 HarnessId::CLAUDE_CODE,
1047 HarnessId::CODEX,
1048 HarnessId::PI,
1049 HarnessId::OPENCODE,
1050 HarnessId::GROK,
1051 HarnessId::GEMINI,
1052 HarnessId::GOOSE,
1053 HarnessId::SUPERCODE,
1054 HarnessId::OPENCLAW,
1055 HarnessId::HERMES,
1056 HarnessId::ORCHESTRATOR,
1057 ]
1058 .into_iter()
1059 .collect()
1060 } else {
1061 query.harnesses.iter().map(HarnessId::as_str).collect()
1062 };
1063 let scope = WorkspaceScope::of(query);
1066 let mut found = Vec::new();
1067 if selected.contains(HarnessId::CLAUDE_CODE) {
1068 discover_jsonl(
1069 &query.homes.claude_code,
1070 HarnessId::CLAUDE_CODE,
1071 scope.as_ref(),
1072 include_child_sessions,
1073 &mut found,
1074 );
1075 }
1076 if selected.contains(HarnessId::CODEX) {
1077 discover_jsonl(
1078 &query.homes.codex,
1079 HarnessId::CODEX,
1080 scope.as_ref(),
1081 include_child_sessions,
1082 &mut found,
1083 );
1084 }
1085 if selected.contains(HarnessId::PI) {
1086 discover_jsonl(
1087 &query.homes.pi,
1088 HarnessId::PI,
1089 scope.as_ref(),
1090 include_child_sessions,
1091 &mut found,
1092 );
1093 }
1094 if selected.contains(HarnessId::OPENCODE) {
1095 discover_opencode(&query.homes.opencode, scope.as_ref(), &mut found);
1096 }
1097 if selected.contains(HarnessId::GROK) {
1098 discover_grok(&query.homes.grok, scope.as_ref(), &mut found);
1099 }
1100 if selected.contains(HarnessId::GEMINI) {
1101 discover_gemini(&query.homes.gemini, scope.as_ref(), &mut found);
1102 }
1103 if selected.contains(HarnessId::GOOSE) {
1104 discover_goose(&query.homes.goose, scope.as_ref(), &mut found);
1105 }
1106 if selected.contains(HarnessId::OPENCLAW) {
1107 discover_openclaw(&query.homes.openclaw, scope.as_ref(), &mut found);
1108 }
1109 if selected.contains(HarnessId::HERMES) {
1110 for store in hermes_session_stores(&query.homes.hermes) {
1111 discover_hermes(&store, scope.as_ref(), &selected, &mut found);
1112 }
1113 }
1114 if selected.contains(HarnessId::ORCHESTRATOR) {
1115 discover_orchestrator(&query.homes.orchestrator, scope.as_ref(), &mut found);
1116 }
1117 if selected.contains(HarnessId::SUPERCODE) {
1118 discover_supercode(&query.homes.supercode, scope.as_ref(), &mut found);
1119 }
1120 for descriptor in &mut found {
1121 finalize_nouns(descriptor);
1122 }
1123 found
1124 }
1125
1126 pub fn refresh_file_descriptor(
1134 &self,
1135 locator: &SessionLocator,
1136 workspace: Option<&Path>,
1137 include_topic_candidates: bool,
1138 ) -> Result<Option<SessionDescriptor>> {
1139 let Some(mut descriptor) = self.refresh_file_index_descriptor(locator, workspace)? else {
1140 return Ok(None);
1141 };
1142 if include_topic_candidates {
1143 descriptor.preview_candidates =
1144 topic_message_candidates(&descriptor.locator).unwrap_or_default();
1145 }
1146 descriptor.latest_message_candidates =
1147 latest_message_candidates(&descriptor.locator).unwrap_or_default();
1148 Ok(Some(descriptor))
1149 }
1150
1151 pub fn refresh_file_index_descriptor(
1155 &self,
1156 locator: &SessionLocator,
1157 workspace: Option<&Path>,
1158 ) -> Result<Option<SessionDescriptor>> {
1159 let scope = workspace.map(WorkspaceScope::exact);
1160 self.refresh_file_index_descriptor_scoped(locator, scope.as_ref())
1161 }
1162
1163 pub fn refresh_file_index_descriptor_for(
1167 &self,
1168 locator: &SessionLocator,
1169 query: &DiscoveryQuery,
1170 ) -> Result<Option<SessionDescriptor>> {
1171 let scope = WorkspaceScope::of(query);
1172 self.refresh_file_index_descriptor_scoped(locator, scope.as_ref())
1173 }
1174
1175 fn refresh_file_index_descriptor_scoped(
1176 &self,
1177 locator: &SessionLocator,
1178 workspace: Option<&WorkspaceScope>,
1179 ) -> Result<Option<SessionDescriptor>> {
1180 let StorageLocator::File { path } = &locator.storage else {
1181 return Ok(None);
1182 };
1183 if !matches!(
1184 locator.harness.as_str(),
1185 HarnessId::CLAUDE_CODE | HarnessId::CODEX
1186 ) {
1187 return Ok(None);
1188 }
1189 if !path.is_file() {
1190 return Ok(None);
1191 }
1192 let Ok(meta) = read_header(path, locator.harness.as_str()) else {
1193 return Ok(None);
1197 };
1198 if workspace.is_some_and(|wanted| meta.cwd.as_deref().is_none_or(|cwd| !wanted.admits(cwd)))
1199 {
1200 return Ok(None);
1201 }
1202 let parent_session_id = meta.parent_session_id.or_else(|| {
1203 (locator.harness.as_str() == HarnessId::CLAUDE_CODE)
1204 .then(|| claude_subagent_parent_id(path))
1205 .flatten()
1206 });
1207 let tail = tail_facts(path, locator.harness.as_str());
1208 let descriptor = SessionDescriptor {
1209 locator: SessionLocator {
1210 harness: locator.harness.clone(),
1211 session_id: meta
1212 .session_id
1213 .unwrap_or_else(|| locator.session_id.clone()),
1214 storage: StorageLocator::File { path: path.clone() },
1215 },
1216 cwd: meta.cwd,
1217 title: meta.title,
1218 preview_candidates: Vec::new(),
1219 latest_message_candidates: Vec::new(),
1220 updated_at_ms: tail.last_turn_ms.or_else(|| modified_ms(path)),
1221 message_count: None,
1222 model: tail.model.or(meta.model),
1223 parent_session_id,
1224 child_session_count: 0,
1225 nouns: OrchestrationNouns::default(),
1226 };
1227 Ok(Some(descriptor))
1228 }
1229
1230 pub fn load(&self, locator: &SessionLocator) -> Result<Session> {
1232 self.load_with_fidelity(locator, Fidelity::ByteLossless)
1233 }
1234
1235 pub fn load_with_fidelity(
1242 &self,
1243 locator: &SessionLocator,
1244 fidelity: Fidelity,
1245 ) -> Result<Session> {
1246 if let Some(session) = load_hermes_locator(locator) {
1247 return session;
1248 }
1249 match &locator.storage {
1250 StorageLocator::File { path } => {
1251 if let Some(session) = load_native_store_family(path)? {
1252 Ok(session)
1253 } else {
1254 Ok(Session::load_with_fidelity(path, fidelity)?)
1255 }
1256 }
1257 StorageLocator::Sqlite { path, selector } => {
1258 if locator.harness.as_str() == HarnessId::GOOSE {
1259 Ok(Session::from_goose_sqlite(path, selector)?)
1260 } else {
1261 Ok(Session::from_opencode_sqlite(path, Some(selector))?)
1262 }
1263 }
1264 }
1265 }
1266
1267 #[doc(hidden)]
1271 pub fn load_parent_with_fidelity(
1272 &self,
1273 locator: &SessionLocator,
1274 fidelity: Fidelity,
1275 ) -> Result<Session> {
1276 if let Some(session) = load_hermes_locator(locator) {
1277 return session;
1278 }
1279 match &locator.storage {
1280 StorageLocator::File { path } => {
1281 if let Some(session) = load_native_store_family(path)? {
1282 Ok(session)
1283 } else {
1284 Ok(Session::load_parent_with_fidelity(path, fidelity)?)
1285 }
1286 }
1287 StorageLocator::Sqlite { path, selector } => {
1288 if locator.harness.as_str() == HarnessId::GOOSE {
1289 Ok(Session::from_goose_sqlite(path, selector)?)
1290 } else {
1291 Ok(Session::from_opencode_sqlite(path, Some(selector))?)
1292 }
1293 }
1294 }
1295 }
1296
1297 #[doc(hidden)]
1300 pub fn load_display_view(
1301 &self,
1302 locator: &SessionLocator,
1303 fidelity: Fidelity,
1304 message_limit: usize,
1305 ) -> Result<Session> {
1306 if let Some(session) = load_hermes_locator(locator) {
1307 let mut session = session?;
1308 if session.messages.len() > message_limit.max(1) {
1309 session
1310 .messages
1311 .drain(..session.messages.len() - message_limit.max(1));
1312 }
1313 return Ok(session);
1314 }
1315 match &locator.storage {
1316 StorageLocator::File { path } => {
1317 if let Some(mut session) = load_native_store_family(path)? {
1318 if session.messages.len() > message_limit.max(1) {
1319 session
1320 .messages
1321 .drain(..session.messages.len() - message_limit.max(1));
1322 }
1323 Ok(session)
1324 } else {
1325 Ok(Session::load_display_view(path, fidelity, message_limit)?)
1326 }
1327 }
1328 StorageLocator::Sqlite { path, selector } => {
1329 let mut session = if locator.harness.as_str() == HarnessId::GOOSE {
1330 Session::from_goose_sqlite_display(path, selector, message_limit)?
1331 } else {
1332 Session::from_opencode_sqlite(path, Some(selector))?
1333 };
1334 if session.messages.len() > message_limit.max(1) {
1335 session
1336 .messages
1337 .drain(..session.messages.len() - message_limit.max(1));
1338 }
1339 Ok(session)
1340 }
1341 }
1342 }
1343
1344 pub fn follow(&self, locator: &SessionLocator) -> Result<SessionFollower> {
1346 self.follow_with_fidelity(locator, Fidelity::ByteLossless)
1347 }
1348
1349 pub fn follow_with_fidelity(
1351 &self,
1352 locator: &SessionLocator,
1353 fidelity: Fidelity,
1354 ) -> Result<SessionFollower> {
1355 SessionFollower::open_locator_with_fidelity(locator, fidelity)
1356 }
1357
1358 #[doc(hidden)]
1360 pub fn follow_read_view(
1361 &self,
1362 locator: &SessionLocator,
1363 fidelity: Fidelity,
1364 include_subagents: bool,
1365 message_limit: Option<usize>,
1366 max_message_chars: Option<usize>,
1367 display_history: bool,
1368 ) -> Result<SessionFollower> {
1369 SessionFollower::open_locator_with_view(
1370 locator,
1371 fidelity,
1372 include_subagents,
1373 message_limit,
1374 max_message_chars,
1375 display_history,
1376 )
1377 }
1378}
1379
1380fn load_hermes_locator(locator: &SessionLocator) -> Option<Result<Session>> {
1386 if locator.harness.as_str() != HarnessId::HERMES {
1387 return None;
1388 }
1389 let StorageLocator::File { path } = &locator.storage else {
1390 return None;
1391 };
1392 Some(Session::from_hermes_sqlite(path, Some(&locator.session_id)))
1393}
1394
1395fn finalize_nouns(descriptor: &mut SessionDescriptor) {
1402 let mut meta = SessionMeta::new(SessionSource::Native);
1403 meta.cwd = descriptor.cwd.clone();
1404 meta.trigger = descriptor.nouns.trigger;
1405 meta.surface = descriptor.nouns.surface.clone();
1406 meta.profile = descriptor.nouns.profile.clone();
1407 meta.recurrence = descriptor.nouns.recurrence.clone();
1408 meta.cross_surface = descriptor.nouns.cross_surface.clone();
1409 descriptor.nouns = OrchestrationNouns::from_meta(&meta);
1410}
1411
1412fn project_descriptors(query: &DiscoveryQuery, found: &mut Vec<SessionDescriptor>) {
1413 roll_up_session_children(found, query.include_child_sessions);
1414 if let Some(root_session_id) = query.root_session_id.as_deref() {
1415 retain_session_family(found, root_session_id);
1416 }
1417 if let Some(family_path) = query.workspace_family.as_deref() {
1418 let family = RepoFamily::of(family_path);
1419 let mut cache: HashMap<PathBuf, bool> = HashMap::new();
1420 found.retain(|descriptor| {
1421 let Some(cwd) = descriptor.cwd.as_deref() else {
1422 return false;
1423 };
1424 *cache
1425 .entry(cwd.to_path_buf())
1426 .or_insert_with(|| RepoFamily::of(cwd).joins(&family))
1427 });
1428 }
1429 if let Some(profile) = query
1430 .profile
1431 .as_deref()
1432 .map(str::trim)
1433 .filter(|p| !p.is_empty())
1434 {
1435 found.retain(|descriptor| descriptor.nouns.profile.as_deref() == Some(profile));
1436 }
1437 if let Some(after) = query.updated_after_ms {
1438 found.retain(|descriptor| descriptor.updated_at_ms.is_some_and(|at| at >= after));
1439 }
1440 if let Some(before) = query.updated_before_ms {
1441 found.retain(|descriptor| descriptor.updated_at_ms.is_some_and(|at| at <= before));
1442 }
1443 found.sort_by(|a, b| {
1444 b.updated_at_ms
1445 .cmp(&a.updated_at_ms)
1446 .then_with(|| a.locator.harness.cmp(&b.locator.harness))
1447 .then_with(|| a.locator.session_id.cmp(&b.locator.session_id))
1448 });
1449 if let Some(search) = query
1450 .query
1451 .as_deref()
1452 .map(str::trim)
1453 .filter(|query| !query.is_empty())
1454 {
1455 let search = search.to_lowercase();
1456 found.retain(|descriptor| descriptor_matches(descriptor, &search));
1457 }
1458}
1459
1460fn paginate_descriptors(
1461 query: &DiscoveryQuery,
1462 found: Vec<SessionDescriptor>,
1463) -> Result<(Vec<SessionDescriptor>, Option<String>)> {
1464 let start = match query.cursor.as_deref() {
1465 Some(cursor) => {
1466 let key = decode_cursor(cursor)?;
1467 found
1468 .iter()
1469 .position(|descriptor| descriptor_cursor_key(descriptor) == key)
1470 .map(|index| index + 1)
1471 .ok_or_else(|| Error::Other("discovery cursor is stale or invalid".into()))?
1472 }
1473 None => 0,
1474 };
1475 let end = query
1476 .limit
1477 .map(|limit| start.saturating_add(limit).min(found.len()))
1478 .unwrap_or(found.len());
1479 let sessions = found[start.min(found.len())..end].to_vec();
1480 let next_cursor = (end < found.len())
1481 .then(|| sessions.last().map(encode_cursor))
1482 .flatten();
1483 Ok((sessions, next_cursor))
1484}
1485
1486fn enrich_descriptors(
1487 query: &DiscoveryQuery,
1488 sessions: &mut [SessionDescriptor],
1489 codex_history: Option<&CodexHistoryTopicIndex>,
1490) -> Result<()> {
1491 let codex_topics = if query.include_topic_candidates && codex_history.is_none() {
1492 codex_history_topics(&query.homes.codex, sessions).unwrap_or_default()
1493 } else {
1494 HashMap::new()
1495 };
1496 for descriptor in sessions {
1497 let topic = (descriptor.locator.harness.as_str() == HarnessId::CODEX)
1498 .then(|| {
1499 codex_history
1500 .and_then(|history| history.topics.get(&descriptor.locator.session_id))
1501 .or_else(|| codex_topics.get(&descriptor.locator.session_id))
1502 })
1503 .flatten();
1504 enrich_descriptor(query, descriptor, topic);
1505 }
1506 Ok(())
1507}
1508
1509fn enrich_descriptor(
1510 query: &DiscoveryQuery,
1511 descriptor: &mut SessionDescriptor,
1512 codex_topic: Option<&Vec<SessionPreviewCandidate>>,
1513) {
1514 if query.include_topic_candidates {
1515 descriptor.preview_candidates = codex_topic
1516 .cloned()
1517 .unwrap_or_else(|| topic_message_candidates(&descriptor.locator).unwrap_or_default());
1518 }
1519 descriptor.latest_message_candidates =
1520 latest_message_candidates(&descriptor.locator).unwrap_or_default();
1521}
1522
1523fn is_false(value: &bool) -> bool {
1524 !value
1525}
1526
1527fn descriptor_matches(descriptor: &SessionDescriptor, search: &str) -> bool {
1528 [
1529 Some(descriptor.locator.harness.as_str()),
1530 Some(descriptor.locator.session_id.as_str()),
1531 descriptor.title.as_deref(),
1532 descriptor.cwd.as_ref().and_then(|path| path.to_str()),
1533 descriptor.model.as_deref(),
1534 ]
1535 .into_iter()
1536 .flatten()
1537 .any(|value| value.to_lowercase().contains(search))
1538}
1539
1540fn descriptor_cursor_key(descriptor: &SessionDescriptor) -> (Option<u64>, String, String) {
1541 (
1542 descriptor.updated_at_ms,
1543 descriptor.locator.harness.as_str().to_string(),
1544 descriptor.locator.session_id.clone(),
1545 )
1546}
1547
1548fn encode_cursor(descriptor: &SessionDescriptor) -> String {
1549 let json = serde_json::to_vec(&descriptor_cursor_key(descriptor)).unwrap_or_default();
1550 let mut encoded = String::with_capacity(json.len() * 2);
1551 for byte in json {
1552 use std::fmt::Write;
1553 let _ = write!(&mut encoded, "{byte:02x}");
1554 }
1555 encoded
1556}
1557
1558fn decode_cursor(cursor: &str) -> Result<(Option<u64>, String, String)> {
1559 if cursor.len() % 2 != 0 {
1560 return Err(Error::Other("discovery cursor is invalid".into()));
1561 }
1562 let bytes = (0..cursor.len())
1563 .step_by(2)
1564 .map(|index| u8::from_str_radix(&cursor[index..index + 2], 16))
1565 .collect::<std::result::Result<Vec<_>, _>>()
1566 .map_err(|_| Error::Other("discovery cursor is invalid".into()))?;
1567 serde_json::from_slice(&bytes).map_err(|_| Error::Other("discovery cursor is invalid".into()))
1568}
1569
1570#[derive(Default)]
1571struct HeaderMeta {
1572 session_id: Option<String>,
1573 cwd: Option<PathBuf>,
1574 title: Option<String>,
1575 model: Option<String>,
1576 parent_session_id: Option<String>,
1577}
1578
1579fn discover_jsonl(
1580 root: &Path,
1581 harness: &str,
1582 workspace: Option<&WorkspaceScope>,
1583 include_child_sessions: bool,
1584 found: &mut Vec<SessionDescriptor>,
1585) {
1586 let mut files = Vec::new();
1587 collect_jsonl(root, harness, include_child_sessions, &mut files);
1588 for path in files {
1589 let Ok(meta) = read_header(&path, harness) else {
1590 continue;
1591 };
1592 if workspace.is_some_and(|wanted| meta.cwd.as_deref().is_none_or(|cwd| !wanted.admits(cwd)))
1593 {
1594 continue;
1595 }
1596 let session_id = meta.session_id.unwrap_or_else(|| {
1597 path.file_stem()
1598 .and_then(|value| value.to_str())
1599 .unwrap_or("unknown")
1600 .to_string()
1601 });
1602 let parent_session_id = meta.parent_session_id.or_else(|| {
1603 (harness == HarnessId::CLAUDE_CODE)
1604 .then(|| claude_subagent_parent_id(&path))
1605 .flatten()
1606 });
1607 let tail = tail_facts(&path, harness);
1608 found.push(SessionDescriptor {
1609 locator: SessionLocator {
1610 harness: HarnessId::new(harness),
1611 session_id,
1612 storage: StorageLocator::File { path: path.clone() },
1613 },
1614 cwd: meta.cwd,
1615 title: meta.title,
1616 preview_candidates: Vec::new(),
1617 latest_message_candidates: Vec::new(),
1618 updated_at_ms: tail.last_turn_ms.or_else(|| modified_ms(&path)),
1619 message_count: None,
1620 model: tail.model.or(meta.model),
1621 parent_session_id,
1622 child_session_count: if harness == HarnessId::CLAUDE_CODE && !include_child_sessions {
1623 count_claude_subagents(&path)
1624 } else {
1625 0
1626 },
1627 nouns: OrchestrationNouns::default(),
1628 });
1629 }
1630}
1631
1632fn collect_jsonl(root: &Path, harness: &str, include_child_sessions: bool, out: &mut Vec<PathBuf>) {
1633 let mut walked = HashSet::new();
1634 collect_jsonl_in(root, harness, include_child_sessions, out, &mut walked);
1635}
1636
1637fn collect_jsonl_in(
1641 root: &Path,
1642 harness: &str,
1643 include_child_sessions: bool,
1644 out: &mut Vec<PathBuf>,
1645 walked: &mut HashSet<PathBuf>,
1646) {
1647 if !walked.insert(fs::canonicalize(root).unwrap_or_else(|_| root.to_path_buf())) {
1648 return;
1649 }
1650 let Ok(entries) = fs::read_dir(root) else {
1651 return;
1652 };
1653 for entry in entries.flatten() {
1654 let Ok(mut kind) = entry.file_type() else {
1655 continue;
1656 };
1657 let path = entry.path();
1658 if kind.is_symlink() {
1659 let Ok(target) = fs::metadata(&path) else {
1660 continue;
1661 };
1662 kind = target.file_type();
1663 }
1664 if kind.is_dir() {
1665 if harness == HarnessId::CLAUDE_CODE
1666 && path.file_name().and_then(|v| v.to_str()) == Some("subagents")
1667 && !include_child_sessions
1668 {
1669 continue;
1670 }
1671 collect_jsonl_in(&path, harness, include_child_sessions, out, walked);
1672 } else if kind.is_file()
1673 && (path.extension().and_then(|v| v.to_str()) == Some("jsonl")
1674 || (harness == HarnessId::CODEX
1675 && path
1676 .file_name()
1677 .and_then(|v| v.to_str())
1678 .is_some_and(|name| name.ends_with(".jsonl.zst"))))
1679 {
1680 out.push(path);
1681 }
1682 }
1683}
1684
1685fn read_header(path: &Path, harness: &str) -> Result<HeaderMeta> {
1686 let file = crate::session::open_session_reader(path)?;
1687 let mut result = HeaderMeta::default();
1688 let mut bytes = 0usize;
1689 for line in BufReader::new(file).lines().take(32) {
1690 let line = line?;
1691 bytes += line.len();
1692 if bytes > 256 * 1024 {
1693 break;
1694 }
1695 let Ok(value) = serde_json::from_str::<Value>(&line) else {
1696 continue;
1697 };
1698 update_header_meta(&mut result, &value, harness);
1699 if result.session_id.is_some() && result.cwd.is_some() && result.model.is_some() {
1700 break;
1701 }
1702 }
1703 if result.session_id.is_none() && result.cwd.is_none() {
1704 return Err(Error::Other(format!(
1705 "{} has no recognizable {harness} session header",
1706 path.display()
1707 )));
1708 }
1709 Ok(result)
1710}
1711
1712fn update_header_meta(result: &mut HeaderMeta, value: &Value, harness: &str) {
1713 match harness {
1714 HarnessId::CLAUDE_CODE => {
1715 fill_string(&mut result.session_id, value.get("sessionId"));
1716 fill_path(&mut result.cwd, value.get("cwd"));
1717 fill_string(&mut result.title, value.get("customTitle"));
1719 fill_string(&mut result.title, value.get("agentName"));
1720 fill_string(
1721 &mut result.model,
1722 value.get("message").and_then(|v| v.get("model")),
1723 );
1724 }
1725 HarnessId::CODEX => {
1726 let payload = value.get("payload").unwrap_or(&Value::Null);
1727 if value.get("type").and_then(Value::as_str) == Some("session_meta") {
1728 fill_string(&mut result.session_id, payload.get("id"));
1729 fill_path(&mut result.cwd, payload.get("cwd"));
1730 fill_string(&mut result.title, payload.get("thread_name"));
1731 fill_string(&mut result.title, payload.get("title"));
1732 fill_string(
1733 &mut result.parent_session_id,
1734 payload.get("parent_thread_id"),
1735 );
1736 if let Some(parent) = payload
1737 .pointer("/source/subagent/thread_spawn/parent_thread_id")
1738 .and_then(Value::as_str)
1739 {
1740 result.parent_session_id = Some(parent.to_string());
1741 }
1742 if result.title.is_none() {
1743 result.title = payload
1744 .pointer("/source/subagent/thread_spawn/agent_path")
1745 .and_then(Value::as_str)
1746 .and_then(|path| path.rsplit('/').find(|part| !part.is_empty()))
1747 .map(humanize_topic);
1748 }
1749 }
1750 if value.get("type").and_then(Value::as_str) == Some("turn_context") {
1751 fill_path(&mut result.cwd, payload.get("cwd"));
1752 fill_string(&mut result.model, payload.get("model"));
1753 }
1754 }
1755 HarnessId::PI => {
1756 if value.get("type").and_then(Value::as_str) == Some("session") {
1757 fill_string(&mut result.session_id, value.get("id"));
1758 fill_path(&mut result.cwd, value.get("cwd"));
1759 }
1760 fill_string(
1761 &mut result.model,
1762 value.get("message").and_then(|v| v.get("model")),
1763 );
1764 }
1765 _ => {}
1766 }
1767}
1768
1769fn roll_up_session_children(found: &mut Vec<SessionDescriptor>, include_children: bool) {
1773 let by_id = found
1774 .iter()
1775 .enumerate()
1776 .map(|(index, descriptor)| {
1777 (
1778 (
1779 descriptor.locator.harness.as_str().to_string(),
1780 descriptor.locator.session_id.clone(),
1781 ),
1782 index,
1783 )
1784 })
1785 .collect::<HashMap<_, _>>();
1786 let mut root_updates = HashMap::<usize, u64>::new();
1787 let mut root_child_counts = HashMap::<usize, usize>::new();
1788
1789 for descriptor in found.iter() {
1790 let Some(mut parent_id) = descriptor.parent_session_id.as_deref() else {
1791 continue;
1792 };
1793 let harness = descriptor.locator.harness.as_str();
1794 let mut root = None;
1795 let mut visited = HashSet::new();
1796 while visited.insert(parent_id.to_string()) {
1797 let Some(&parent_index) = by_id.get(&(harness.to_string(), parent_id.to_string()))
1798 else {
1799 break;
1800 };
1801 root = Some(parent_index);
1802 let Some(next_parent) = found[parent_index].parent_session_id.as_deref() else {
1803 break;
1804 };
1805 parent_id = next_parent;
1806 }
1807 if let (Some(root), Some(updated_at_ms)) = (root, descriptor.updated_at_ms) {
1808 root_updates
1809 .entry(root)
1810 .and_modify(|current| *current = (*current).max(updated_at_ms))
1811 .or_insert(updated_at_ms);
1812 }
1813 if let Some(root) = root {
1814 *root_child_counts.entry(root).or_default() += 1;
1815 }
1816 }
1817
1818 for (root, child_updated_at_ms) in root_updates {
1819 found[root].updated_at_ms = Some(
1820 found[root]
1821 .updated_at_ms
1822 .unwrap_or_default()
1823 .max(child_updated_at_ms),
1824 );
1825 }
1826 for (root, child_count) in root_child_counts {
1827 found[root].child_session_count = child_count;
1828 }
1829 if !include_children {
1830 found.retain(|descriptor| descriptor.parent_session_id.is_none());
1831 }
1832}
1833
1834fn retain_session_family(found: &mut Vec<SessionDescriptor>, root_session_id: &str) {
1835 let parent_by_id = found
1836 .iter()
1837 .map(|descriptor| {
1838 (
1839 descriptor.locator.session_id.clone(),
1840 descriptor.parent_session_id.clone(),
1841 )
1842 })
1843 .collect::<HashMap<_, _>>();
1844 found.retain(|descriptor| {
1845 let mut current = descriptor.locator.session_id.clone();
1846 let mut visited = HashSet::new();
1847 while visited.insert(current.clone()) {
1848 if current == root_session_id {
1849 return true;
1850 }
1851 let Some(Some(parent)) = parent_by_id.get(¤t) else {
1852 return false;
1853 };
1854 current = parent.clone();
1855 }
1856 false
1857 });
1858}
1859
1860fn claude_subagent_parent_id(path: &Path) -> Option<String> {
1861 let subagents = path.parent()?;
1862 if subagents.file_name()?.to_str()? != "subagents" {
1863 return None;
1864 }
1865 subagents
1866 .parent()?
1867 .file_name()?
1868 .to_str()
1869 .map(str::to_string)
1870}
1871
1872fn count_claude_subagents(parent_path: &Path) -> usize {
1873 let Some(parent) = parent_path.parent() else {
1874 return 0;
1875 };
1876 let Some(stem) = parent_path.file_stem() else {
1877 return 0;
1878 };
1879 let root = parent.join(stem).join("subagents");
1880 let mut files = Vec::new();
1881 collect_jsonl(&root, HarnessId::CLAUDE_CODE, true, &mut files);
1882 files.len()
1883}
1884
1885fn is_zero(value: &usize) -> bool {
1886 *value == 0
1887}
1888
1889fn humanize_topic(value: &str) -> String {
1890 let text = value.replace(['_', '-'], " ");
1891 let mut characters = text.chars();
1892 match characters.next() {
1893 Some(first) => first.to_uppercase().collect::<String>() + characters.as_str(),
1894 None => text,
1895 }
1896}
1897
1898fn discover_gemini(
1899 root: &Path,
1900 workspace: Option<&WorkspaceScope>,
1901 found: &mut Vec<SessionDescriptor>,
1902) {
1903 let slug_to_cwd = std::fs::read_to_string(root.join("projects.json"))
1904 .ok()
1905 .and_then(|text| serde_json::from_str::<Value>(&text).ok())
1906 .and_then(|value| value.get("projects").and_then(Value::as_object).cloned())
1907 .map(|projects| {
1908 projects
1909 .into_iter()
1910 .filter_map(|(cwd, slug)| Some((slug.as_str()?.to_string(), PathBuf::from(cwd))))
1911 .collect::<HashMap<_, _>>()
1912 })
1913 .unwrap_or_default();
1914 let mut files = Vec::new();
1915 collect_jsonl(&root.join("tmp"), HarnessId::GEMINI, false, &mut files);
1916 let worker_count = std::thread::available_parallelism()
1917 .map(usize::from)
1918 .unwrap_or(4)
1919 .clamp(1, 8)
1920 .min(files.len().max(1));
1921 let chunk_size = files.len().max(1).div_ceil(worker_count);
1922 let discovered = std::thread::scope(|scope| {
1923 files
1924 .chunks(chunk_size)
1925 .map(|paths| {
1926 scope.spawn(|| {
1927 paths
1928 .iter()
1929 .filter_map(|path| gemini_descriptor(path, &slug_to_cwd, workspace))
1930 .collect::<Vec<_>>()
1931 })
1932 })
1933 .collect::<Vec<_>>()
1934 .into_iter()
1935 .flat_map(|worker| {
1936 worker
1937 .join()
1938 .expect("Gemini discovery worker must not panic")
1939 })
1940 .collect::<Vec<_>>()
1941 });
1942 found.extend(discovered);
1943}
1944
1945fn gemini_descriptor(
1946 path: &Path,
1947 slug_to_cwd: &HashMap<String, PathBuf>,
1948 workspace: Option<&WorkspaceScope>,
1949) -> Option<SessionDescriptor> {
1950 if path
1951 .parent()
1952 .and_then(Path::file_name)
1953 .and_then(|name| name.to_str())
1954 != Some("chats")
1955 {
1956 return None;
1957 }
1958 let slug = path
1959 .parent()
1960 .and_then(Path::parent)
1961 .and_then(Path::file_name)
1962 .and_then(|name| name.to_str());
1963 let cwd = slug.and_then(|slug| slug_to_cwd.get(slug)).cloned();
1964 if workspace.is_some_and(|wanted| cwd.as_deref().is_none_or(|actual| !wanted.admits(actual))) {
1965 return None;
1966 }
1967
1968 let file = File::open(path).ok()?;
1972 let mut reader = BufReader::new(file.take(64 * 1024));
1973 let mut header = String::new();
1974 reader.read_line(&mut header).ok()?;
1975 let header = serde_json::from_str::<Value>(&header).ok()?;
1976 let session_id = header.get("sessionId")?.as_str()?.to_string();
1977 let mut model = None;
1978 for line in reader
1979 .take(4 * 1024)
1980 .lines()
1981 .map_while(std::result::Result::ok)
1982 {
1983 let Ok(value) = serde_json::from_str::<Value>(&line) else {
1984 continue;
1985 };
1986 let kind = value.get("type").and_then(Value::as_str);
1987 if kind != Some("user") && kind != Some("gemini") {
1988 continue;
1989 }
1990 if model.is_none() {
1991 model = value
1992 .get("model")
1993 .and_then(Value::as_str)
1994 .map(str::to_string);
1995 }
1996 if model.is_some() {
1997 break;
1998 }
1999 }
2000 Some(SessionDescriptor {
2001 locator: SessionLocator {
2002 harness: HarnessId::from(HarnessId::GEMINI),
2003 session_id,
2004 storage: StorageLocator::File {
2005 path: path.to_path_buf(),
2006 },
2007 },
2008 cwd,
2009 title: None,
2010 preview_candidates: Vec::new(),
2011 latest_message_candidates: Vec::new(),
2012 updated_at_ms: tail_facts(path, HarnessId::GEMINI)
2013 .last_turn_ms
2014 .or_else(|| modified_ms(path)),
2015 message_count: None,
2016 model,
2017 parent_session_id: None,
2018 child_session_count: 0,
2019 nouns: OrchestrationNouns::default(),
2020 })
2021}
2022
2023fn display_text(content: Option<&Value>) -> Option<String> {
2024 match content? {
2025 Value::String(text) => Some(text.clone()),
2026 Value::Array(parts) => Some(
2027 parts
2028 .iter()
2029 .filter_map(|part| part.get("text").and_then(Value::as_str))
2030 .collect::<Vec<_>>()
2031 .join(" ")
2032 .trim()
2033 .to_string(),
2034 ),
2035 _ => None,
2036 }
2037}
2038
2039fn discover_hermes(
2044 db_path: &Path,
2045 workspace: Option<&WorkspaceScope>,
2046 selected: &HashSet<&str>,
2047 found: &mut Vec<SessionDescriptor>,
2048) {
2049 if !db_path.is_file() {
2050 return;
2051 }
2052 let Ok(conn) = Connection::open_with_flags(
2053 db_path,
2054 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
2055 ) else {
2056 return;
2057 };
2058 let fingerprint_ok = ["sessions", "messages", "schema_version"].iter().all(|t| {
2059 conn.query_row(
2060 "SELECT 1 FROM sqlite_master WHERE type='table' AND name=?1",
2061 [t],
2062 |_| Ok(()),
2063 )
2064 .is_ok()
2065 });
2066 if !fingerprint_ok {
2067 return;
2068 }
2069 let Ok(mut statement) = conn.prepare(
2070 "SELECT id, cwd, title, model, message_count, started_at, ended_at, parent_session_id, \
2071 source, model_config FROM sessions ORDER BY started_at DESC",
2072 ) else {
2073 return;
2074 };
2075 let Ok(rows) = statement.query_map([], |row| {
2076 Ok((
2077 row.get::<_, String>(0)?,
2078 row.get::<_, Option<String>>(1)?,
2079 row.get::<_, Option<String>>(2)?,
2080 row.get::<_, Option<String>>(3)?,
2081 row.get::<_, Option<i64>>(4)?,
2082 row.get::<_, Option<f64>>(5)?,
2083 row.get::<_, Option<f64>>(6)?,
2084 row.get::<_, Option<String>>(7)?,
2085 row.get::<_, Option<String>>(8)?,
2086 row.get::<_, Option<String>>(9)?,
2087 ))
2088 }) else {
2089 return;
2090 };
2091 for row in rows.flatten() {
2092 let (
2093 id,
2094 cwd,
2095 title,
2096 model,
2097 message_count,
2098 started_at,
2099 ended_at,
2100 parent,
2101 source,
2102 model_config,
2103 ) = row;
2104 let mirror = model_config
2108 .as_deref()
2109 .and_then(|c| serde_json::from_str::<serde_json::Value>(c).ok())
2110 .and_then(|c| c.get("_supercode_mirror").cloned());
2111 if let Some(mirror) = mirror {
2112 let worker_read = mirror
2113 .get("harness")
2114 .and_then(serde_json::Value::as_str)
2115 .is_some_and(|h| selected.contains(h));
2116 let mirrored = mirror.get("messages").and_then(serde_json::Value::as_i64);
2117 if worker_read && mirrored.is_none_or(|m| message_count.unwrap_or(0) <= m) {
2118 continue;
2119 }
2120 if mirror.get("continued_as").is_some()
2121 && mirrored.is_none_or(|m| message_count.unwrap_or(0) <= m)
2122 {
2123 continue;
2124 }
2125 }
2126 let cwd = cwd.map(PathBuf::from);
2127 if let Some(filter) = workspace {
2128 if !filter.admits_literal(cwd.as_deref()) {
2129 continue;
2130 }
2131 }
2132 let updated_at_ms = ended_at
2133 .or(started_at)
2134 .map(|seconds| (seconds * 1000.0) as u64);
2135 let mut meta = SessionMeta::new(SessionSource::Hermes);
2139 meta.cwd = cwd.clone();
2140 if let Some(hermes_source) = source.filter(|value| !value.is_empty()) {
2141 meta.lineage
2142 .insert("hermes_source".to_string(), hermes_source);
2143 }
2144 if let Some(parent_id) = parent.as_deref() {
2145 meta.lineage.insert(
2146 "hermes_lineage_kind".to_string(),
2147 crate::session::hermes_lineage_kind(
2148 &conn,
2149 parent_id,
2150 model_config.as_deref(),
2151 started_at,
2152 )
2153 .to_string(),
2154 );
2155 }
2156 hermes_capture_nouns(&conn, &id, &mut meta);
2157 found.push(SessionDescriptor {
2158 locator: SessionLocator {
2159 harness: HarnessId::new(HarnessId::HERMES),
2160 session_id: id,
2161 storage: StorageLocator::File {
2162 path: db_path.to_path_buf(),
2163 },
2164 },
2165 cwd,
2166 title: title.filter(|t| !t.is_empty()),
2167 preview_candidates: Vec::new(),
2168 latest_message_candidates: Vec::new(),
2169 updated_at_ms,
2170 message_count: message_count.map(|count| count.max(0) as usize),
2171 model,
2172 parent_session_id: parent,
2173 child_session_count: 0,
2174 nouns: OrchestrationNouns::from_meta(&meta),
2175 });
2176 }
2177}
2178
2179fn discover_orchestrator(
2186 root: &Path,
2187 workspace: Option<&WorkspaceScope>,
2188 found: &mut Vec<SessionDescriptor>,
2189) {
2190 if workspace.is_some() {
2192 return;
2193 }
2194 for (profile, dir) in orchestrator_profile_dirs(root) {
2195 let db_path = dir.join("state.db");
2196 if !db_path.is_file() {
2197 continue;
2198 }
2199 let Ok(conn) = Connection::open_with_flags(
2200 &db_path,
2201 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
2202 ) else {
2203 continue;
2204 };
2205 let sessions = dir.join("sessions");
2208 let scope = std::fs::canonicalize(&sessions)
2209 .unwrap_or(sessions)
2210 .display()
2211 .to_string();
2212 let Ok(mut statement) = conn.prepare(
2213 "SELECT json_extract(entry_json, '$.metadata.supercode.binding'), \
2214 CAST(strftime('%s', json_extract(entry_json, '$.metadata.supercode.binding.last_activity_at')) AS INTEGER) \
2215 FROM gateway_routing WHERE scope = ?1 AND json_extract(entry_json, '$.metadata.supercode') IS NOT NULL \
2216 ORDER BY json_extract(entry_json, '$.metadata.supercode.binding.last_activity_at') DESC",
2217 ) else {
2218 continue;
2219 };
2220 let Ok(rows) = statement.query_map([&scope], |row| {
2221 Ok((
2222 row.get::<_, Option<String>>(0)?,
2223 row.get::<_, Option<i64>>(1)?,
2224 ))
2225 }) else {
2226 continue;
2227 };
2228 let rows = rows.flatten().filter_map(|(json, epoch)| {
2229 let b: Binding = serde_json::from_str(&json?).ok()?;
2230 let text = |v: &Option<String>| v.clone().filter(|s| !s.is_empty());
2231 Some((
2232 OrchestratorBindingRow {
2233 platform: b.key.platform.clone().unwrap_or_default(),
2234 chat_type: b.key.kind.clone().unwrap_or_default(),
2235 chat_id: text(&b.key.chat_id),
2236 thread_id: text(&b.key.thread_id),
2237 participant_id: text(&b.key.participant_id),
2238 worker_harness: b.worker.harness.as_str().to_string(),
2239 worker_session_id: text(&b.worker.session_id),
2240 worker_locator: text(&b.worker.locator),
2241 started_at: b.started_at.clone(),
2242 last_activity_at: b.last_activity_at.clone(),
2243 ended_at: b.ended_at.clone(),
2244 end_reason: b.end_reason.map(|r| r.as_str().to_string()),
2245 handoff_to: b.handoff.as_ref().and_then(|h| h.to.clone()),
2246 handoff_state: b.handoff.as_ref().map(|h| h.state.clone()),
2247 handoff_error: b.handoff.as_ref().and_then(|h| h.error.clone()),
2248 recurrence_job_id: b.recurrence.as_ref().map(|r| r.job_id.clone()),
2249 },
2250 epoch,
2251 ))
2252 });
2253 for (row, last_activity_epoch) in rows {
2254 found.push(orchestrator_descriptor(
2255 &db_path,
2256 &profile,
2257 &row,
2258 last_activity_epoch,
2259 ));
2260 }
2261 }
2262}
2263
2264fn orchestrator_descriptor(
2265 db_path: &Path,
2266 profile: &str,
2267 row: &OrchestratorBindingRow,
2268 last_activity_epoch: Option<i64>,
2269) -> SessionDescriptor {
2270 let binding = Binding::from_orchestrator_row(profile, row);
2271 let nouns = binding.nouns();
2272 let mut title = format!(
2275 "{} {}",
2276 row.worker_harness,
2277 row.worker_session_id
2278 .as_deref()
2279 .unwrap_or("(no worker session yet)")
2280 );
2281 if let Some(reason) = row.end_reason.as_deref().filter(|_| row.ended_at.is_some()) {
2282 title.push_str(&format!(" (ended: {reason})"));
2283 }
2284 SessionDescriptor {
2285 locator: SessionLocator {
2286 harness: HarnessId::new(HarnessId::ORCHESTRATOR),
2287 session_id: row.worker_session_id.clone().unwrap_or_default(),
2288 storage: StorageLocator::File {
2289 path: row
2290 .worker_locator
2291 .clone()
2292 .map_or_else(|| db_path.to_path_buf(), PathBuf::from),
2293 },
2294 },
2295 cwd: None,
2296 title: Some(title),
2297 preview_candidates: Vec::new(),
2298 latest_message_candidates: Vec::new(),
2299 updated_at_ms: last_activity_epoch.map(|seconds| (seconds.max(0) as u64) * 1000),
2300 message_count: None,
2301 model: None,
2302 parent_session_id: None,
2303 child_session_count: 0,
2304 nouns,
2305 }
2306}
2307
2308fn discover_openclaw(
2314 root: &Path,
2315 workspace: Option<&WorkspaceScope>,
2316 found: &mut Vec<SessionDescriptor>,
2317) {
2318 let agents = root.join("agents");
2319 let Ok(agent_dirs) = std::fs::read_dir(&agents) else {
2320 return;
2321 };
2322 for agent_dir in agent_dirs.flatten() {
2323 let sessions = agent_dir.path().join("sessions");
2324 let Ok(files) = std::fs::read_dir(&sessions) else {
2325 continue;
2326 };
2327 for file in files.flatten() {
2328 let path = file.path();
2329 let name = file.file_name();
2330 let name = name.to_string_lossy();
2331 if !name.ends_with(".jsonl") || name.ends_with(".trajectory.jsonl") {
2332 continue;
2333 }
2334 let Ok(text) = std::fs::read_to_string(&path) else {
2335 continue;
2336 };
2337 let Some(header_line) = text.lines().find(|line| !line.trim().is_empty()) else {
2338 continue;
2339 };
2340 let Ok(header) = serde_json::from_str::<serde_json::Value>(header_line) else {
2341 continue;
2342 };
2343 if header.get("type").and_then(serde_json::Value::as_str) != Some("session") {
2344 continue;
2345 }
2346 let session_id = header
2347 .get("id")
2348 .and_then(serde_json::Value::as_str)
2349 .unwrap_or_else(|| name.trim_end_matches(".jsonl"))
2350 .to_string();
2351 let cwd = header
2352 .get("cwd")
2353 .and_then(serde_json::Value::as_str)
2354 .map(PathBuf::from);
2355 if let Some(filter) = workspace {
2356 if !filter.admits_literal(cwd.as_deref()) {
2357 continue;
2358 }
2359 }
2360 let updated_at_ms = tail_facts(&path, HarnessId::OPENCLAW)
2361 .last_turn_ms
2362 .or_else(|| {
2363 file.metadata()
2364 .ok()
2365 .and_then(|metadata| metadata.modified().ok())
2366 .and_then(|modified| modified.duration_since(std::time::UNIX_EPOCH).ok())
2367 .map(|elapsed| elapsed.as_millis() as u64)
2368 });
2369 let message_count = text
2370 .lines()
2371 .filter(|line| line.contains("\"type\":\"message\""))
2372 .count();
2373 let mut meta = SessionMeta::new(SessionSource::OpenClaw);
2377 meta.cwd = cwd.clone();
2378 openclaw_capture_header_nouns(&header, &mut meta);
2379 if meta.profile.is_none() {
2380 meta.profile = openclaw_agent_id_from_path(&path);
2381 }
2382 found.push(SessionDescriptor {
2383 locator: SessionLocator {
2384 harness: HarnessId::new(HarnessId::OPENCLAW),
2385 session_id,
2386 storage: StorageLocator::File { path },
2387 },
2388 cwd,
2389 title: None,
2390 preview_candidates: Vec::new(),
2391 latest_message_candidates: Vec::new(),
2392 updated_at_ms,
2393 message_count: Some(message_count),
2394 model: None,
2395 parent_session_id: None,
2396 child_session_count: 0,
2397 nouns: OrchestrationNouns::from_meta(&meta),
2398 });
2399 }
2400 }
2401}
2402
2403fn discover_supercode(
2404 root: &Path,
2405 workspace: Option<&WorkspaceScope>,
2406 found: &mut Vec<SessionDescriptor>,
2407) {
2408 for info in list_native_store(root) {
2409 let path = if info.archived {
2410 root.join("archived").join(format!("{}.jsonl", info.name))
2411 } else {
2412 root.join(format!("{}.jsonl", info.name))
2413 };
2414 let sidecar = path.with_extension("sidecar.jsonl");
2419 let path = if path.is_file() {
2420 path
2421 } else {
2422 sidecar.clone()
2423 };
2424 let header = read_native_store_header(&path);
2425 if workspace.is_some_and(|wanted| {
2426 header
2427 .as_ref()
2428 .and_then(|meta| meta.cwd.as_deref())
2429 .is_none_or(|cwd| !wanted.admits(cwd))
2430 }) {
2431 continue;
2432 }
2433 let title = (!info.title.trim().is_empty()).then_some(info.title);
2434 let updated_at_ms = tail_facts(&sidecar, HarnessId::SUPERCODE)
2437 .last_turn_ms
2438 .or_else(|| tail_facts(&path, HarnessId::SUPERCODE).last_turn_ms)
2439 .or_else(|| modified_ms(&path))
2440 .or_else(|| modified_ms(&sidecar));
2441 found.push(SessionDescriptor {
2442 locator: SessionLocator {
2443 harness: HarnessId::from(HarnessId::SUPERCODE),
2444 session_id: info.name,
2445 storage: StorageLocator::File { path: path.clone() },
2446 },
2447 cwd: header.as_ref().and_then(|meta| meta.cwd.clone()),
2448 title,
2449 preview_candidates: Vec::new(),
2450 latest_message_candidates: Vec::new(),
2451 updated_at_ms,
2452 message_count: None,
2453 model: header.and_then(|meta| meta.model),
2454 parent_session_id: None,
2455 child_session_count: 0,
2456 nouns: OrchestrationNouns::default(),
2457 });
2458 }
2459}
2460
2461fn read_native_store_header(path: &Path) -> Option<HeaderMeta> {
2466 let name = path.file_stem()?.to_str()?;
2467 let sidecar = path.with_file_name(format!("{name}.sidecar.jsonl"));
2468 let source_path = if sidecar.is_file() {
2469 sidecar
2470 } else {
2471 path.to_path_buf()
2472 };
2473 let file = File::open(source_path).ok()?;
2474 let mut result = HeaderMeta::default();
2475 let mut source = None;
2476 let mut bytes = 0usize;
2477 for line in BufReader::new(file).lines().take(32) {
2478 let line = line.ok()?;
2479 bytes += line.len();
2480 if bytes > 256 * 1024 {
2481 break;
2482 }
2483 let Ok(value) = serde_json::from_str::<Value>(&line) else {
2484 continue;
2485 };
2486 if source.is_none() {
2487 source = value.get("source").and_then(Value::as_str).map(|source| {
2488 if source == "claude_code" {
2489 HarnessId::CLAUDE_CODE.to_string()
2490 } else {
2491 source.to_string()
2492 }
2493 });
2494 fill_string(&mut result.session_id, value.get("session_id"));
2495 }
2496 if let Some(harness) = source.as_deref() {
2497 update_header_meta(&mut result, &value, harness);
2498 }
2499 if result.cwd.is_some() && result.model.is_some() {
2500 break;
2501 }
2502 }
2503 Some(result)
2504}
2505
2506#[derive(Deserialize)]
2507struct NativeStoreInfo {
2508 name: String,
2509 #[serde(default)]
2510 title: String,
2511 #[serde(skip)]
2512 archived: bool,
2513}
2514
2515fn list_native_store(root: &Path) -> Vec<NativeStoreInfo> {
2516 let mut sessions = Vec::new();
2517 for archived in [false, true] {
2518 let directory = if archived {
2519 root.join("archived")
2520 } else {
2521 root.to_path_buf()
2522 };
2523 let Ok(entries) = fs::read_dir(directory) else {
2524 continue;
2525 };
2526 for entry in entries.flatten() {
2527 let path = entry.path();
2528 if !path.to_string_lossy().ends_with(".meta.json") {
2529 continue;
2530 }
2531 let Ok(text) = fs::read_to_string(path) else {
2532 continue;
2533 };
2534 let Ok(mut info) = serde_json::from_str::<NativeStoreInfo>(&text) else {
2535 continue;
2536 };
2537 info.archived = archived;
2538 sessions.push(info);
2539 }
2540 }
2541 sessions.sort_by(|left, right| left.name.cmp(&right.name));
2542 sessions
2543}
2544
2545fn discover_grok(
2546 root: &Path,
2547 workspace: Option<&WorkspaceScope>,
2548 found: &mut Vec<SessionDescriptor>,
2549) {
2550 let Ok(workspaces) = fs::read_dir(root) else {
2551 return;
2552 };
2553 for workspace_entry in workspaces.flatten() {
2554 let encoded = workspace_entry.file_name();
2555 let Some(cwd) = encoded
2556 .to_str()
2557 .and_then(percent_decode_path)
2558 .map(PathBuf::from)
2559 else {
2560 continue;
2561 };
2562 if workspace.is_some_and(|wanted| !wanted.admits(&cwd)) {
2563 continue;
2564 }
2565 let Ok(sessions) = fs::read_dir(workspace_entry.path()) else {
2566 continue;
2567 };
2568 for session_entry in sessions.flatten() {
2569 let session_dir = session_entry.path();
2570 if !session_dir.is_dir() {
2571 continue;
2572 }
2573 let transcript = session_dir.join("chat_history.jsonl");
2574 if !transcript.is_file() {
2575 continue;
2576 }
2577 let Some(session_id) = session_dir
2578 .file_name()
2579 .and_then(|name| name.to_str())
2580 .map(str::to_string)
2581 else {
2582 continue;
2583 };
2584 let summary = fs::read_to_string(session_dir.join("summary.json"))
2585 .ok()
2586 .and_then(|text| serde_json::from_str::<Value>(&text).ok());
2587 let title = summary
2588 .as_ref()
2589 .and_then(|value| value.get("generated_title"))
2590 .and_then(Value::as_str)
2591 .filter(|title| !title.is_empty())
2592 .map(str::to_string);
2593 let model = summary
2594 .as_ref()
2595 .and_then(|value| value.get("current_model_id"))
2596 .and_then(Value::as_str)
2597 .map(str::to_string);
2598 let message_count = summary
2599 .as_ref()
2600 .and_then(|value| value.get("num_chat_messages"))
2601 .and_then(Value::as_u64)
2602 .and_then(|count| usize::try_from(count).ok());
2603 let updated_at_ms = summary
2604 .as_ref()
2605 .and_then(|value| value.get("updated_at"))
2606 .and_then(Value::as_str)
2607 .and_then(crate::sidecar::rfc3339_to_ms)
2608 .and_then(|millis| u64::try_from(millis).ok())
2609 .or_else(|| modified_ms(&transcript));
2610 found.push(SessionDescriptor {
2611 locator: SessionLocator {
2612 harness: HarnessId::from(HarnessId::GROK),
2613 session_id,
2614 storage: StorageLocator::File { path: transcript },
2615 },
2616 cwd: Some(cwd.clone()),
2617 title,
2618 preview_candidates: Vec::new(),
2619 latest_message_candidates: Vec::new(),
2620 updated_at_ms,
2621 message_count,
2622 model,
2623 parent_session_id: None,
2624 child_session_count: 0,
2625 nouns: OrchestrationNouns::default(),
2626 });
2627 }
2628 }
2629}
2630
2631fn discover_opencode(
2632 root: &Path,
2633 workspace: Option<&WorkspaceScope>,
2634 found: &mut Vec<SessionDescriptor>,
2635) {
2636 let mut dbs = Vec::new();
2637 if root.is_file() {
2638 dbs.push(root.to_path_buf());
2639 } else if let Ok(entries) = fs::read_dir(root) {
2640 dbs.extend(entries.flatten().map(|entry| entry.path()).filter(|path| {
2641 path.file_name()
2642 .and_then(|v| v.to_str())
2643 .is_some_and(|name| name.starts_with("opencode") && name.ends_with(".db"))
2644 }));
2645 }
2646 dbs.sort();
2647 for db in dbs {
2648 let Ok(conn) = Connection::open_with_flags(
2649 &db,
2650 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
2651 ) else {
2652 continue;
2653 };
2654 let has_model = conn.prepare("SELECT model FROM session LIMIT 0").is_ok();
2655 let model_column = if has_model { "s.model" } else { "NULL" };
2656 let query = format!(
2657 "SELECT s.id, s.directory, s.title, s.time_updated, {model_column}, COUNT(m.id) \
2658 FROM session s LEFT JOIN message m ON m.session_id = s.id \
2659 GROUP BY s.id ORDER BY s.time_updated DESC"
2660 );
2661 let Ok(mut stmt) = conn.prepare(&query) else {
2662 continue;
2663 };
2664 let Ok(rows) = stmt.query_map([], |row| {
2665 Ok((
2666 row.get::<_, String>(0)?,
2667 row.get::<_, String>(1)?,
2668 row.get::<_, String>(2)?,
2669 row.get::<_, i64>(3)?,
2670 row.get::<_, Option<String>>(4)?,
2671 row.get::<_, i64>(5)?,
2672 ))
2673 }) else {
2674 continue;
2675 };
2676 for row in rows.flatten() {
2677 let (id, cwd, title, updated, model, messages) = row;
2678 let cwd = PathBuf::from(cwd);
2679 if workspace.is_some_and(|wanted| !wanted.admits(&cwd)) {
2680 continue;
2681 }
2682 found.push(SessionDescriptor {
2683 locator: SessionLocator {
2684 harness: HarnessId::from(HarnessId::OPENCODE),
2685 session_id: id.clone(),
2686 storage: StorageLocator::Sqlite {
2687 path: db.clone(),
2688 selector: id,
2689 },
2690 },
2691 cwd: Some(cwd),
2692 title: (!title.is_empty()).then_some(title),
2693 preview_candidates: Vec::new(),
2694 latest_message_candidates: Vec::new(),
2695 updated_at_ms: u64::try_from(updated).ok(),
2696 message_count: usize::try_from(messages).ok(),
2697 model,
2698 parent_session_id: None,
2699 child_session_count: 0,
2700 nouns: OrchestrationNouns::default(),
2701 });
2702 }
2703 }
2704}
2705
2706fn discover_goose(
2707 root: &Path,
2708 workspace: Option<&WorkspaceScope>,
2709 found: &mut Vec<SessionDescriptor>,
2710) {
2711 let db = if root.is_file() {
2712 root.to_path_buf()
2713 } else if root.join("sessions.db").is_file() {
2714 root.join("sessions.db")
2715 } else {
2716 root.join("sessions/sessions.db")
2717 };
2718 let Ok(connection) = Connection::open_with_flags(
2719 &db,
2720 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
2721 ) else {
2722 return;
2723 };
2724 let Ok(mut statement) = connection.prepare(
2725 "SELECT s.id, s.working_dir, s.name, s.updated_at, s.model_config_json, \
2726 COUNT(m.id) \
2727 FROM sessions s LEFT JOIN messages m ON m.session_id = s.id \
2728 WHERE s.archived_at IS NULL \
2729 GROUP BY s.id ORDER BY s.updated_at DESC",
2730 ) else {
2731 return;
2732 };
2733 let Ok(rows) = statement.query_map([], |row| {
2734 Ok((
2735 row.get::<_, String>(0)?,
2736 row.get::<_, String>(1)?,
2737 row.get::<_, String>(2)?,
2738 row.get::<_, String>(3)?,
2739 row.get::<_, Option<String>>(4)?,
2740 row.get::<_, i64>(5)?,
2741 ))
2742 }) else {
2743 return;
2744 };
2745 for row in rows.flatten() {
2746 let (id, cwd, title, updated_at, model_config, message_count) = row;
2747 let cwd = PathBuf::from(cwd);
2748 if workspace.is_some_and(|wanted| !wanted.admits(&cwd)) {
2749 continue;
2750 }
2751 let model = model_config
2752 .as_deref()
2753 .and_then(|value| serde_json::from_str::<Value>(value).ok())
2754 .and_then(|value| {
2755 value
2756 .get("model_name")
2757 .or_else(|| value.get("modelName"))
2758 .and_then(Value::as_str)
2759 .map(str::to_string)
2760 });
2761 let updated_at_ms = crate::sidecar::rfc3339_to_ms(&updated_at)
2762 .or_else(|| {
2763 crate::sidecar::rfc3339_to_ms(&format!("{}Z", updated_at.replace(' ', "T")))
2765 })
2766 .and_then(|value| u64::try_from(value).ok());
2767 found.push(SessionDescriptor {
2768 locator: SessionLocator {
2769 harness: HarnessId::from(HarnessId::GOOSE),
2770 session_id: id.clone(),
2771 storage: StorageLocator::Sqlite {
2772 path: db.clone(),
2773 selector: id,
2774 },
2775 },
2776 cwd: Some(cwd),
2777 title: (!title.trim().is_empty()).then_some(title),
2778 preview_candidates: Vec::new(),
2779 latest_message_candidates: Vec::new(),
2780 updated_at_ms,
2781 message_count: usize::try_from(message_count).ok(),
2782 model,
2783 parent_session_id: None,
2784 child_session_count: 0,
2785 nouns: OrchestrationNouns::default(),
2786 });
2787 }
2788}
2789
2790const LATEST_PREVIEW_CANDIDATES: usize = 8;
2791const TOPIC_PREVIEW_HEAD_BYTES: u64 = 512 * 1024;
2792const LATEST_PREVIEW_TAIL_BYTES: u64 = 512 * 1024;
2793const LATEST_PREVIEW_MAX_BYTES: u64 = 4 * 1024 * 1024;
2794
2795fn topic_message_candidates(locator: &SessionLocator) -> Result<Vec<SessionPreviewCandidate>> {
2796 match &locator.storage {
2797 StorageLocator::File { path }
2798 if matches!(
2799 locator.harness.as_str(),
2800 HarnessId::CLAUDE_CODE | HarnessId::CODEX
2801 ) =>
2802 {
2803 topic_file_message_candidates(path, locator.harness.as_str())
2804 }
2805 _ => Ok(Vec::new()),
2806 }
2807}
2808
2809fn codex_history_topics(
2810 sessions_root: &Path,
2811 sessions: &[SessionDescriptor],
2812) -> Result<HashMap<String, Vec<SessionPreviewCandidate>>> {
2813 let wanted: HashSet<&str> = sessions
2814 .iter()
2815 .filter(|descriptor| descriptor.locator.harness.as_str() == HarnessId::CODEX)
2816 .map(|descriptor| descriptor.locator.session_id.as_str())
2817 .collect();
2818 if wanted.is_empty() {
2819 return Ok(HashMap::new());
2820 }
2821 let Some(root) = sessions_root.parent() else {
2822 return Ok(HashMap::new());
2823 };
2824 let file = match File::open(root.join("history.jsonl")) {
2825 Ok(file) => file,
2826 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(HashMap::new()),
2827 Err(error) => return Err(error.into()),
2828 };
2829 let mut topics = HashMap::new();
2830 for line in BufReader::new(file).lines() {
2831 let Ok(value) = serde_json::from_str::<Value>(&line?) else {
2832 continue;
2833 };
2834 let Some(session_id) = value.get("session_id").and_then(Value::as_str) else {
2835 continue;
2836 };
2837 if !wanted.contains(session_id) || topics.contains_key(session_id) {
2838 continue;
2839 }
2840 let mut candidates = Vec::new();
2841 push_message_candidate(&mut candidates, "user", value.get("text"), HashMap::new());
2842 if !candidates.is_empty() {
2843 topics.insert(session_id.to_string(), candidates);
2844 if topics.len() == wanted.len() {
2845 break;
2846 }
2847 }
2848 }
2849 Ok(topics)
2850}
2851
2852fn latest_message_candidates(locator: &SessionLocator) -> Result<Vec<SessionPreviewCandidate>> {
2853 match &locator.storage {
2854 StorageLocator::File { path } | StorageLocator::Sqlite { path, .. }
2857 if locator.harness.as_str() == HarnessId::HERMES =>
2858 {
2859 latest_hermes_message_candidates(path, &locator.session_id)
2860 }
2861 StorageLocator::File { path } => {
2862 latest_file_message_candidates(path, locator.harness.as_str())
2863 }
2864 StorageLocator::Sqlite { path, selector }
2865 if locator.harness.as_str() == HarnessId::OPENCODE =>
2866 {
2867 latest_opencode_message_candidates(path, selector)
2868 }
2869 StorageLocator::Sqlite { path, selector }
2870 if locator.harness.as_str() == HarnessId::GOOSE =>
2871 {
2872 latest_goose_message_candidates(path, selector)
2873 }
2874 StorageLocator::Sqlite { .. } => Ok(Vec::new()),
2875 }
2876}
2877
2878fn topic_file_message_candidates(
2879 path: &Path,
2880 harness: &str,
2881) -> Result<Vec<SessionPreviewCandidate>> {
2882 let mut file = File::open(path)?;
2883 let mut bytes = Vec::with_capacity(TOPIC_PREVIEW_HEAD_BYTES as usize);
2884 file.by_ref()
2885 .take(TOPIC_PREVIEW_HEAD_BYTES)
2886 .read_to_end(&mut bytes)?;
2887 if file.metadata()?.len() > TOPIC_PREVIEW_HEAD_BYTES {
2888 if let Some(newline) = bytes.iter().rposition(|byte| *byte == b'\n') {
2889 bytes.truncate(newline);
2890 }
2891 }
2892 let text = String::from_utf8(bytes).map_err(|_| {
2893 Error::Other(format!(
2894 "{} contains non-UTF-8 data in its topic-preview window",
2895 path.display()
2896 ))
2897 })?;
2898 if harness == HarnessId::CODEX {
2899 return Ok(codex_preview_candidates(text.lines(), false));
2900 }
2901 let mut candidates = Vec::new();
2902 for line in text.lines() {
2903 let Ok(value) = serde_json::from_str::<Value>(line) else {
2904 continue;
2905 };
2906 push_topic_message_candidate(&mut candidates, harness, &value);
2907 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
2908 break;
2909 }
2910 }
2911 Ok(candidates)
2912}
2913
2914fn latest_file_message_candidates(
2915 path: &Path,
2916 harness: &str,
2917) -> Result<Vec<SessionPreviewCandidate>> {
2918 let mut candidates =
2919 latest_file_message_candidates_with_limit(path, harness, LATEST_PREVIEW_TAIL_BYTES)?;
2920 if candidates.is_empty() {
2921 candidates =
2922 latest_file_message_candidates_with_limit(path, harness, LATEST_PREVIEW_MAX_BYTES)?;
2923 }
2924 Ok(candidates)
2925}
2926
2927fn latest_file_message_candidates_with_limit(
2928 path: &Path,
2929 harness: &str,
2930 byte_limit: u64,
2931) -> Result<Vec<SessionPreviewCandidate>> {
2932 let mut file = File::open(path)?;
2933 let file_len = file.metadata()?.len();
2934 let start = file_len.saturating_sub(byte_limit);
2935 file.seek(SeekFrom::Start(start))?;
2936 let mut bytes = Vec::with_capacity((file_len - start) as usize);
2937 file.read_to_end(&mut bytes)?;
2938 if start > 0 {
2939 if let Some(newline) = bytes.iter().position(|byte| *byte == b'\n') {
2940 bytes.drain(..=newline);
2941 } else {
2942 return Ok(Vec::new());
2943 }
2944 }
2945 let text = String::from_utf8(bytes).map_err(|_| {
2946 Error::Other(format!(
2947 "{} contains non-UTF-8 data in its list-preview window",
2948 path.display()
2949 ))
2950 })?;
2951 if harness == HarnessId::CODEX {
2952 return Ok(codex_preview_candidates(text.lines().rev(), true));
2953 }
2954 let mut candidates = Vec::new();
2955 for line in text.lines().rev() {
2956 let Ok(value) = serde_json::from_str::<Value>(line) else {
2957 continue;
2958 };
2959 let (role, content, metadata) = match harness {
2960 HarnessId::CLAUDE_CODE => {
2961 let role = value.get("type").and_then(Value::as_str);
2962 if !matches!(role, Some("user" | "assistant")) {
2963 continue;
2964 }
2965 let metadata = if role == Some("user") {
2966 crate::session::claude_user_provenance(&value)
2967 .into_iter()
2968 .collect()
2969 } else {
2970 HashMap::new()
2971 };
2972 (
2973 role.unwrap_or_default(),
2974 value
2975 .get("message")
2976 .and_then(|message| message.get("content")),
2977 metadata,
2978 )
2979 }
2980 HarnessId::PI => {
2981 if value.get("type").and_then(Value::as_str) != Some("message") {
2982 continue;
2983 }
2984 let message = value.get("message").unwrap_or(&Value::Null);
2985 let Some(role @ ("user" | "assistant")) =
2986 message.get("role").and_then(Value::as_str)
2987 else {
2988 continue;
2989 };
2990 (role, message.get("content"), HashMap::new())
2991 }
2992 HarnessId::GEMINI => {
2993 let Some(kind @ ("user" | "gemini")) = value.get("type").and_then(Value::as_str)
2994 else {
2995 continue;
2996 };
2997 (
2998 if kind == "gemini" {
2999 "assistant"
3000 } else {
3001 "user"
3002 },
3003 value.get("content"),
3004 HashMap::new(),
3005 )
3006 }
3007 HarnessId::GROK => {
3008 let Some(role @ ("user" | "assistant")) = value.get("type").and_then(Value::as_str)
3009 else {
3010 continue;
3011 };
3012 (role, value.get("content"), HashMap::new())
3013 }
3014 HarnessId::SUPERCODE => {
3015 let Some(role @ ("user" | "assistant")) = value.get("role").and_then(Value::as_str)
3016 else {
3017 continue;
3018 };
3019 (role, value.get("content"), HashMap::new())
3020 }
3021 _ => continue,
3022 };
3023 let mut metadata = metadata;
3024 if matches!(harness, HarnessId::CLAUDE_CODE | HarnessId::CODEX) {
3025 if let Some(timestamp) = value.get("timestamp").and_then(Value::as_str) {
3026 metadata.insert("timestamp".to_string(), timestamp.to_string());
3027 }
3028 }
3029 push_message_candidate_with_cursor(
3030 &mut candidates,
3031 role,
3032 content,
3033 metadata,
3034 Some(message_candidate_cursor(harness, &value)),
3035 );
3036 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3037 break;
3038 }
3039 }
3040 Ok(candidates)
3041}
3042
3043struct CodexPreviewRecord {
3044 native: Value,
3045 role: String,
3046 text: String,
3047}
3048
3049fn codex_preview_candidates<'a>(
3055 lines: impl Iterator<Item = &'a str>,
3056 latest: bool,
3057) -> Vec<SessionPreviewCandidate> {
3058 let mut candidates = Vec::new();
3059 let mut pending: Option<CodexPreviewRecord> = None;
3060 for line in lines {
3061 let Ok(native) = serde_json::from_str::<Value>(line) else {
3062 continue;
3063 };
3064 let Some((role, content)) = codex_preview_message(&native) else {
3065 continue;
3066 };
3067 let Some(text) = display_text(content).filter(|text| !text.trim().is_empty()) else {
3068 continue;
3069 };
3070 let current = CodexPreviewRecord {
3071 role: role.to_string(),
3072 text,
3073 native,
3074 };
3075 if let Some(previous) = pending.take() {
3076 if previous.role == current.role
3077 && previous.text == current.text
3078 && previous.native.get("type") != current.native.get("type")
3079 {
3080 let canonical = if previous.native.get("type").and_then(Value::as_str)
3081 == Some("response_item")
3082 {
3083 previous
3084 } else {
3085 current
3086 };
3087 push_codex_preview_candidate(&mut candidates, canonical, latest);
3088 } else {
3089 push_codex_preview_candidate(&mut candidates, previous, latest);
3090 pending = Some(current);
3091 }
3092 } else {
3093 pending = Some(current);
3094 }
3095 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3096 break;
3097 }
3098 }
3099 if let Some(last) = pending {
3100 push_codex_preview_candidate(&mut candidates, last, latest);
3101 }
3102 candidates
3103}
3104
3105fn push_codex_preview_candidate(
3106 candidates: &mut Vec<SessionPreviewCandidate>,
3107 record: CodexPreviewRecord,
3108 latest: bool,
3109) {
3110 let mut metadata = HashMap::new();
3111 if latest {
3112 if let Some(timestamp) = record.native.get("timestamp").and_then(Value::as_str) {
3113 metadata.insert("timestamp".to_string(), timestamp.to_string());
3114 }
3115 }
3116 let cursor = latest.then(|| message_candidate_cursor(HarnessId::CODEX, &record.native));
3117 push_message_candidate_with_cursor(
3118 candidates,
3119 &record.role,
3120 Some(&Value::String(record.text)),
3121 metadata,
3122 cursor,
3123 );
3124}
3125
3126fn codex_preview_message(value: &Value) -> Option<(&str, Option<&Value>)> {
3128 let payload = value.get("payload")?;
3129 match (
3130 value.get("type").and_then(Value::as_str)?,
3131 payload.get("type").and_then(Value::as_str)?,
3132 ) {
3133 ("response_item", "message") => {
3134 let role @ ("user" | "assistant") = payload.get("role").and_then(Value::as_str)? else {
3135 return None;
3136 };
3137 Some((role, payload.get("content")))
3138 }
3139 ("event_msg", "user_message") => Some(("user", payload.get("message"))),
3140 ("event_msg", "agent_message") => Some(("assistant", payload.get("message"))),
3141 _ => None,
3142 }
3143}
3144
3145fn push_topic_message_candidate(
3146 candidates: &mut Vec<SessionPreviewCandidate>,
3147 harness: &str,
3148 value: &Value,
3149) {
3150 let (role, content, metadata) = match harness {
3151 HarnessId::CLAUDE_CODE => {
3152 let role = value.get("type").and_then(Value::as_str);
3153 if !matches!(role, Some("user" | "assistant")) {
3154 return;
3155 }
3156 let metadata = if role == Some("user") {
3157 crate::session::claude_user_provenance(value)
3158 .into_iter()
3159 .collect()
3160 } else {
3161 HashMap::new()
3162 };
3163 (
3164 role.unwrap_or_default(),
3165 value
3166 .get("message")
3167 .and_then(|message| message.get("content")),
3168 metadata,
3169 )
3170 }
3171 _ => return,
3172 };
3173 push_message_candidate(candidates, role, content, metadata);
3174}
3175
3176fn latest_opencode_message_candidates(
3177 path: &Path,
3178 session_id: &str,
3179) -> Result<Vec<SessionPreviewCandidate>> {
3180 let connection = Connection::open_with_flags(
3181 path,
3182 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
3183 )
3184 .map_err(|error| Error::Other(format!("failed to open list-preview store: {error}")))?;
3185 let mut statement = connection
3186 .prepare(
3187 "SELECT m.data, p.data FROM message m JOIN part p ON p.message_id = m.id \
3188 WHERE m.session_id = ?1 ORDER BY m.time_created DESC, p.time_created DESC LIMIT 32",
3189 )
3190 .map_err(|error| Error::Other(format!("failed to prepare list-preview query: {error}")))?;
3191 let rows = statement
3192 .query_map([session_id], |row| {
3193 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
3194 })
3195 .map_err(|error| Error::Other(format!("failed to read list-preview rows: {error}")))?;
3196 let mut candidates = Vec::new();
3197 for row in rows.flatten() {
3198 let (Ok(message), Ok(part)) = (
3199 serde_json::from_str::<Value>(&row.0),
3200 serde_json::from_str::<Value>(&row.1),
3201 ) else {
3202 continue;
3203 };
3204 let Some(role @ ("user" | "assistant")) = message.get("role").and_then(Value::as_str)
3205 else {
3206 continue;
3207 };
3208 if part.get("type").and_then(Value::as_str) != Some("text") {
3209 continue;
3210 }
3211 push_message_candidate(&mut candidates, role, part.get("text"), HashMap::new());
3212 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3213 break;
3214 }
3215 }
3216 Ok(candidates)
3217}
3218
3219fn latest_hermes_message_candidates(
3220 path: &Path,
3221 session_id: &str,
3222) -> Result<Vec<SessionPreviewCandidate>> {
3223 let connection = Connection::open_with_flags(
3224 path,
3225 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
3226 )
3227 .map_err(|error| Error::Other(format!("failed to open list-preview store: {error}")))?;
3228 let mut statement = connection
3229 .prepare(
3230 "SELECT role, content FROM messages WHERE session_id = ?1 AND active = 1 \
3231 AND role IN ('user', 'assistant') AND content IS NOT NULL AND content != '' \
3232 ORDER BY timestamp DESC, id DESC LIMIT 32",
3233 )
3234 .map_err(|error| Error::Other(format!("failed to prepare list-preview query: {error}")))?;
3235 let rows = statement
3236 .query_map([session_id], |row| {
3237 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
3238 })
3239 .map_err(|error| Error::Other(format!("failed to read list-preview rows: {error}")))?;
3240 let mut candidates = Vec::new();
3241 for (role, content) in rows.flatten() {
3242 push_message_candidate(
3243 &mut candidates,
3244 &role,
3245 Some(&Value::String(content)),
3246 HashMap::new(),
3247 );
3248 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3249 break;
3250 }
3251 }
3252 Ok(candidates)
3253}
3254
3255fn latest_goose_message_candidates(
3256 path: &Path,
3257 session_id: &str,
3258) -> Result<Vec<SessionPreviewCandidate>> {
3259 let connection = Connection::open_with_flags(
3260 path,
3261 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
3262 )
3263 .map_err(|error| Error::Other(format!("failed to open list-preview store: {error}")))?;
3264 let mut statement = connection
3265 .prepare(
3266 "SELECT role, content_json FROM messages WHERE session_id = ?1 \
3267 ORDER BY created_timestamp DESC, id DESC LIMIT 16",
3268 )
3269 .map_err(|error| Error::Other(format!("failed to prepare list-preview query: {error}")))?;
3270 let rows = statement
3271 .query_map([session_id], |row| {
3272 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
3273 })
3274 .map_err(|error| Error::Other(format!("failed to read list-preview rows: {error}")))?;
3275 let mut candidates = Vec::new();
3276 for row in rows.flatten() {
3277 let (role, content) = row;
3278 if !matches!(role.as_str(), "user" | "assistant") {
3279 continue;
3280 }
3281 let Ok(content) = serde_json::from_str::<Value>(&content) else {
3282 continue;
3283 };
3284 push_message_candidate(&mut candidates, &role, Some(&content), HashMap::new());
3285 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3286 break;
3287 }
3288 }
3289 Ok(candidates)
3290}
3291
3292fn push_message_candidate(
3293 candidates: &mut Vec<SessionPreviewCandidate>,
3294 role: &str,
3295 content: Option<&Value>,
3296 metadata: HashMap<String, String>,
3297) {
3298 push_message_candidate_with_cursor(candidates, role, content, metadata, None);
3299}
3300
3301fn push_message_candidate_with_cursor(
3302 candidates: &mut Vec<SessionPreviewCandidate>,
3303 role: &str,
3304 content: Option<&Value>,
3305 metadata: HashMap<String, String>,
3306 cursor: Option<String>,
3307) {
3308 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3309 return;
3310 }
3311 let Some(text) = display_text(content).filter(|text| !text.trim().is_empty()) else {
3312 return;
3313 };
3314 const MAX_CHARS: usize = 4_096;
3315 candidates.push(SessionPreviewCandidate {
3316 cursor,
3317 role: role.to_string(),
3318 content: text.chars().take(MAX_CHARS).collect(),
3319 metadata,
3320 });
3321}
3322
3323fn message_candidate_cursor(harness: &str, value: &Value) -> String {
3324 let native_identity = value
3325 .get("uuid")
3326 .or_else(|| value.get("id"))
3327 .or_else(|| value.pointer("/message/id"))
3328 .or_else(|| value.pointer("/payload/id"))
3329 .and_then(Value::as_str)
3330 .or_else(|| value.get("timestamp").and_then(Value::as_str));
3331 let mut hasher = blake3::Hasher::new();
3332 hasher.update(b"supercode.session-preview-cursor.v1\0");
3333 hasher.update(harness.as_bytes());
3334 hasher.update(b"\0");
3335 if let Some(identity) = native_identity {
3336 hasher.update(identity.as_bytes());
3337 } else {
3338 hasher.update(value.to_string().as_bytes());
3342 }
3343 format!("v1:{}", &hasher.finalize().to_hex()[..24])
3344}
3345
3346fn fill_string(target: &mut Option<String>, value: Option<&Value>) {
3347 if target.is_none() {
3348 *target = value.and_then(Value::as_str).map(str::to_owned);
3349 }
3350}
3351
3352fn fill_path(target: &mut Option<PathBuf>, value: Option<&Value>) {
3353 if target.is_none() {
3354 *target = value.and_then(Value::as_str).map(PathBuf::from);
3355 }
3356}
3357
3358#[derive(Default)]
3361struct TailFacts {
3362 last_turn_ms: Option<u64>,
3364 model: Option<String>,
3366}
3367
3368fn tail_facts(path: &Path, harness: &str) -> TailFacts {
3384 let mut facts = TailFacts::default();
3385 let Ok(mut file) = File::open(path) else {
3386 return facts;
3387 };
3388 let Ok(len) = file.metadata().map(|meta| meta.len()) else {
3389 return facts;
3390 };
3391 let mut window = TAIL_SCAN_START.min(len);
3392 loop {
3393 if file.seek(SeekFrom::Start(len - window)).is_err() {
3394 return facts;
3395 }
3396 let Ok(size) = usize::try_from(window) else {
3397 return facts;
3398 };
3399 let mut buf = vec![0u8; size];
3400 if file.read_exact(&mut buf).is_err() {
3401 return facts;
3402 }
3403 let floor = if window < len {
3414 buf.iter().position(|byte| *byte == b'\n').map(|at| at + 1)
3415 } else {
3416 Some(0)
3417 };
3418 if let Some(floor) = floor {
3419 let mut end = buf.len();
3420 while end > floor && !(facts.last_turn_ms.is_some() && facts.model.is_some()) {
3421 let start = buf[floor..end]
3422 .iter()
3423 .rposition(|byte| *byte == b'\n')
3424 .map_or(floor, |at| floor + at + 1);
3425 if let Ok(record) = std::str::from_utf8(&buf[start..end])
3426 .map_err(|_| ())
3427 .and_then(|line| serde_json::from_str::<Value>(line).map_err(|_| ()))
3428 {
3429 if facts.last_turn_ms.is_none() {
3430 facts.last_turn_ms = record_timestamp(&record, harness)
3431 .and_then(crate::sidecar::rfc3339_to_ms)
3432 .and_then(|millis| u64::try_from(millis).ok());
3433 }
3434 if facts.model.is_none() {
3435 facts.model = record_model(&record, harness).map(str::to_owned);
3436 }
3437 }
3438 end = start.saturating_sub(1);
3439 }
3440 }
3441 if facts.last_turn_ms.is_some() || window >= len || window >= TAIL_SCAN_LIMIT {
3442 return facts;
3443 }
3444 window = (window * 2).min(len).min(TAIL_SCAN_LIMIT);
3445 }
3446}
3447
3448fn record_timestamp<'a>(record: &'a Value, harness: &str) -> Option<&'a str> {
3452 let key = match harness {
3453 HarnessId::SUPERCODE => "ts",
3454 _ => "timestamp",
3455 };
3456 record.get(key)?.as_str()
3457}
3458
3459fn record_model<'a>(record: &'a Value, harness: &str) -> Option<&'a str> {
3462 match harness {
3463 HarnessId::CLAUDE_CODE | HarnessId::PI => record.get("message")?.get("model")?.as_str(),
3464 HarnessId::CODEX => {
3465 if record.get("type")?.as_str()? != "turn_context" {
3466 return None;
3467 }
3468 record.get("payload")?.get("model")?.as_str()
3469 }
3470 _ => None,
3471 }
3472}
3473
3474const TAIL_SCAN_START: u64 = 16 * 1024;
3477
3478const TAIL_SCAN_LIMIT: u64 = 1024 * 1024;
3481
3482fn modified_ms(path: &Path) -> Option<u64> {
3483 fs::metadata(path)
3484 .ok()?
3485 .modified()
3486 .ok()?
3487 .duration_since(UNIX_EPOCH)
3488 .ok()
3489 .and_then(|duration| u64::try_from(duration.as_millis()).ok())
3490}
3491
3492pub struct WorkspaceScope {
3495 wanted: PathBuf,
3496 roots: Vec<PathBuf>,
3498 subtree: bool,
3499 judged: std::sync::Mutex<HashMap<PathBuf, bool>>,
3501}
3502
3503impl WorkspaceScope {
3504 pub fn exact(wanted: &Path) -> Self {
3506 Self {
3507 wanted: wanted.to_path_buf(),
3508 roots: Vec::new(),
3509 subtree: false,
3510 judged: std::sync::Mutex::new(HashMap::new()),
3511 }
3512 }
3513
3514 pub fn subtree(root: &Path) -> Self {
3516 Self {
3517 wanted: root.to_path_buf(),
3518 roots: path_comparison_keys(root),
3519 subtree: true,
3520 judged: std::sync::Mutex::new(HashMap::new()),
3521 }
3522 }
3523
3524 pub fn of(query: &DiscoveryQuery) -> Option<Self> {
3526 query.workspace.as_deref().map(|wanted| {
3527 if query.workspace_subtree {
3528 Self::subtree(wanted)
3529 } else {
3530 Self::exact(wanted)
3531 }
3532 })
3533 }
3534
3535 pub fn admits(&self, recorded: &Path) -> bool {
3540 if !recorded.is_absolute() {
3541 return false;
3542 }
3543 if !self.subtree {
3544 return same_path(recorded, &self.wanted);
3545 }
3546 let mut judged = self
3547 .judged
3548 .lock()
3549 .unwrap_or_else(std::sync::PoisonError::into_inner);
3550 *judged.entry(recorded.to_path_buf()).or_insert_with(|| {
3551 path_comparison_keys(recorded)
3552 .iter()
3553 .any(|key| self.roots.iter().any(|root| key.starts_with(root)))
3554 })
3555 }
3556
3557 fn admits_literal(&self, recorded: Option<&Path>) -> bool {
3561 if self.subtree {
3562 recorded.is_some_and(|cwd| self.admits(cwd))
3563 } else {
3564 recorded == Some(self.wanted.as_path())
3565 }
3566 }
3567}
3568
3569fn path_comparison_keys(path: &Path) -> Vec<PathBuf> {
3575 let mut keys = Vec::with_capacity(2);
3576 if let Ok(canonical) = fs::canonicalize(path) {
3577 keys.push(comparison_key(&canonical));
3578 }
3579 let lexical = comparison_key(&normalize_path(path));
3580 if !keys.contains(&lexical) {
3581 keys.push(lexical);
3582 }
3583 keys
3584}
3585
3586#[cfg(windows)]
3587fn comparison_key(path: &Path) -> PathBuf {
3588 let text = path.to_string_lossy();
3589 let text = if let Some(unc) = text.strip_prefix(r"\\?\UNC\") {
3590 format!(r"\\{unc}")
3591 } else if let Some(local) = text.strip_prefix(r"\\?\") {
3592 local.to_string()
3593 } else {
3594 text.into_owned()
3595 };
3596 PathBuf::from(text.to_lowercase())
3597}
3598
3599#[cfg(not(windows))]
3600fn comparison_key(path: &Path) -> PathBuf {
3601 path.to_path_buf()
3602}
3603
3604fn recorded_cwd_matches(recorded: &Path, wanted: &Path) -> bool {
3610 recorded.is_absolute() && same_path(recorded, wanted)
3611}
3612
3613fn same_path(left: &Path, right: &Path) -> bool {
3614 match (fs::canonicalize(left), fs::canonicalize(right)) {
3615 (Ok(left), Ok(right)) => left == right,
3616 _ => normalize_path(left) == normalize_path(right),
3617 }
3618}
3619
3620fn normalize_path(path: &Path) -> PathBuf {
3621 let absolute = if path.is_absolute() {
3622 path.to_path_buf()
3623 } else {
3624 std::env::current_dir()
3625 .unwrap_or_else(|_| PathBuf::from("."))
3626 .join(path)
3627 };
3628 let mut normalized = PathBuf::new();
3629 for component in absolute.components() {
3630 match component {
3631 Component::CurDir => {}
3632 Component::ParentDir => {
3633 normalized.pop();
3634 }
3635 other => normalized.push(other.as_os_str()),
3636 }
3637 }
3638 normalized
3639}
3640
3641fn claude_custom_titles(path: &std::path::Path) -> Vec<String> {
3646 use std::io::{Read, Seek, SeekFrom};
3647 use std::sync::{Mutex, OnceLock};
3648 static CACHE: OnceLock<
3649 Mutex<std::collections::HashMap<std::path::PathBuf, (u64, Vec<String>)>>,
3650 > = OnceLock::new();
3651 const KEY: &[u8] = b"\"customTitle\":";
3652 let Ok(len) = std::fs::metadata(path).map(|meta| meta.len()) else {
3653 return Vec::new();
3654 };
3655 let cache = CACHE.get_or_init(Default::default);
3656 let known = cache.lock().ok().and_then(|cache| cache.get(path).cloned());
3657 let (from, mut titles) = match known {
3658 Some((seen, titles)) if seen == len => return titles,
3659 Some((seen, titles)) if seen < len => (seen.saturating_sub(4096), titles),
3661 _ => (0, Vec::new()),
3662 };
3663 let mut bytes = Vec::new();
3664 let read = std::fs::File::open(path).and_then(|mut file| {
3665 file.seek(SeekFrom::Start(from))?;
3666 file.read_to_end(&mut bytes)
3667 });
3668 if read.is_err() {
3669 return titles;
3670 }
3671 let mut at = 0;
3672 while let Some(offset) = bytes[at..]
3673 .windows(KEY.len())
3674 .position(|window| window == KEY)
3675 {
3676 let start = at + offset + KEY.len();
3677 if let Some(Ok(title)) = serde_json::Deserializer::from_slice(&bytes[start..])
3678 .into_iter::<String>()
3679 .next()
3680 {
3681 if !titles.contains(&title) {
3682 titles.push(title);
3683 }
3684 }
3685 at = start;
3686 }
3687 if let Ok(mut cache) = cache.lock() {
3688 cache.insert(path.to_path_buf(), (len, titles.clone()));
3689 }
3690 titles
3691}