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, Trigger};
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 trigger: Option<Trigger>,
1580}
1581
1582fn discover_jsonl(
1583 root: &Path,
1584 harness: &str,
1585 workspace: Option<&WorkspaceScope>,
1586 include_child_sessions: bool,
1587 found: &mut Vec<SessionDescriptor>,
1588) {
1589 let mut files = Vec::new();
1590 collect_jsonl(root, harness, include_child_sessions, &mut files);
1591 for path in files {
1592 let Ok(meta) = read_header(&path, harness) else {
1593 continue;
1594 };
1595 if workspace.is_some_and(|wanted| meta.cwd.as_deref().is_none_or(|cwd| !wanted.admits(cwd)))
1596 {
1597 continue;
1598 }
1599 let session_id = meta.session_id.unwrap_or_else(|| {
1600 path.file_stem()
1601 .and_then(|value| value.to_str())
1602 .unwrap_or("unknown")
1603 .to_string()
1604 });
1605 let parent_session_id = meta.parent_session_id.or_else(|| {
1606 (harness == HarnessId::CLAUDE_CODE)
1607 .then(|| claude_subagent_parent_id(&path))
1608 .flatten()
1609 });
1610 let tail = tail_facts(&path, harness);
1611 found.push(SessionDescriptor {
1612 locator: SessionLocator {
1613 harness: HarnessId::new(harness),
1614 session_id,
1615 storage: StorageLocator::File { path: path.clone() },
1616 },
1617 cwd: meta.cwd,
1618 title: meta.title,
1619 preview_candidates: Vec::new(),
1620 latest_message_candidates: Vec::new(),
1621 updated_at_ms: tail.last_turn_ms.or_else(|| modified_ms(&path)),
1622 message_count: None,
1623 model: tail.model.or(meta.model),
1624 parent_session_id,
1625 child_session_count: if harness == HarnessId::CLAUDE_CODE && !include_child_sessions {
1626 count_claude_subagents(&path)
1627 } else {
1628 0
1629 },
1630 nouns: OrchestrationNouns {
1631 trigger: meta.trigger,
1632 ..OrchestrationNouns::default()
1633 },
1634 });
1635 }
1636}
1637
1638fn collect_jsonl(root: &Path, harness: &str, include_child_sessions: bool, out: &mut Vec<PathBuf>) {
1639 let mut walked = HashSet::new();
1640 collect_jsonl_in(root, harness, include_child_sessions, out, &mut walked);
1641}
1642
1643fn collect_jsonl_in(
1647 root: &Path,
1648 harness: &str,
1649 include_child_sessions: bool,
1650 out: &mut Vec<PathBuf>,
1651 walked: &mut HashSet<PathBuf>,
1652) {
1653 if !walked.insert(fs::canonicalize(root).unwrap_or_else(|_| root.to_path_buf())) {
1654 return;
1655 }
1656 let Ok(entries) = fs::read_dir(root) else {
1657 return;
1658 };
1659 for entry in entries.flatten() {
1660 let Ok(mut kind) = entry.file_type() else {
1661 continue;
1662 };
1663 let path = entry.path();
1664 if kind.is_symlink() {
1665 let Ok(target) = fs::metadata(&path) else {
1666 continue;
1667 };
1668 kind = target.file_type();
1669 }
1670 if kind.is_dir() {
1671 if harness == HarnessId::CLAUDE_CODE
1672 && path.file_name().and_then(|v| v.to_str()) == Some("subagents")
1673 && !include_child_sessions
1674 {
1675 continue;
1676 }
1677 collect_jsonl_in(&path, harness, include_child_sessions, out, walked);
1678 } else if kind.is_file()
1679 && (path.extension().and_then(|v| v.to_str()) == Some("jsonl")
1680 || (harness == HarnessId::CODEX
1681 && path
1682 .file_name()
1683 .and_then(|v| v.to_str())
1684 .is_some_and(|name| name.ends_with(".jsonl.zst"))))
1685 {
1686 out.push(path);
1687 }
1688 }
1689}
1690
1691fn read_header(path: &Path, harness: &str) -> Result<HeaderMeta> {
1692 let file = crate::session::open_session_reader(path)?;
1693 let mut result = HeaderMeta::default();
1694 let mut bytes = 0usize;
1695 for line in BufReader::new(file).lines().take(32) {
1696 let line = line?;
1697 bytes += line.len();
1698 if bytes > 256 * 1024 {
1699 break;
1700 }
1701 let Ok(value) = serde_json::from_str::<Value>(&line) else {
1702 continue;
1703 };
1704 update_header_meta(&mut result, &value, harness);
1705 if result.session_id.is_some() && result.cwd.is_some() && result.model.is_some() {
1706 break;
1707 }
1708 }
1709 if result.session_id.is_none() && result.cwd.is_none() {
1710 return Err(Error::Other(format!(
1711 "{} has no recognizable {harness} session header",
1712 path.display()
1713 )));
1714 }
1715 Ok(result)
1716}
1717
1718fn update_header_meta(result: &mut HeaderMeta, value: &Value, harness: &str) {
1719 match harness {
1720 HarnessId::CLAUDE_CODE => {
1721 fill_string(&mut result.session_id, value.get("sessionId"));
1722 fill_path(&mut result.cwd, value.get("cwd"));
1723 fill_string(&mut result.title, value.get("customTitle"));
1725 fill_string(&mut result.title, value.get("agentName"));
1726 fill_string(
1727 &mut result.model,
1728 value.get("message").and_then(|v| v.get("model")),
1729 );
1730 }
1731 HarnessId::CODEX => {
1732 let payload = value.get("payload").unwrap_or(&Value::Null);
1733 if value.get("type").and_then(Value::as_str) == Some("session_meta") {
1734 fill_string(&mut result.session_id, payload.get("id"));
1735 if result.trigger.is_none() {
1736 result.trigger = crate::session::codex_start_trigger(payload);
1737 }
1738 fill_path(&mut result.cwd, payload.get("cwd"));
1739 fill_string(&mut result.title, payload.get("thread_name"));
1740 fill_string(&mut result.title, payload.get("title"));
1741 fill_string(
1742 &mut result.parent_session_id,
1743 payload.get("parent_thread_id"),
1744 );
1745 if let Some(parent) = payload
1746 .pointer("/source/subagent/thread_spawn/parent_thread_id")
1747 .and_then(Value::as_str)
1748 {
1749 result.parent_session_id = Some(parent.to_string());
1750 }
1751 if result.title.is_none() {
1752 result.title = payload
1753 .pointer("/source/subagent/thread_spawn/agent_path")
1754 .and_then(Value::as_str)
1755 .and_then(|path| path.rsplit('/').find(|part| !part.is_empty()))
1756 .map(humanize_topic);
1757 }
1758 }
1759 if value.get("type").and_then(Value::as_str) == Some("turn_context") {
1760 fill_path(&mut result.cwd, payload.get("cwd"));
1761 fill_string(&mut result.model, payload.get("model"));
1762 }
1763 }
1764 HarnessId::PI => {
1765 if value.get("type").and_then(Value::as_str) == Some("session") {
1766 fill_string(&mut result.session_id, value.get("id"));
1767 fill_path(&mut result.cwd, value.get("cwd"));
1768 }
1769 fill_string(
1770 &mut result.model,
1771 value.get("message").and_then(|v| v.get("model")),
1772 );
1773 }
1774 _ => {}
1775 }
1776}
1777
1778fn roll_up_session_children(found: &mut Vec<SessionDescriptor>, include_children: bool) {
1782 let by_id = found
1783 .iter()
1784 .enumerate()
1785 .map(|(index, descriptor)| {
1786 (
1787 (
1788 descriptor.locator.harness.as_str().to_string(),
1789 descriptor.locator.session_id.clone(),
1790 ),
1791 index,
1792 )
1793 })
1794 .collect::<HashMap<_, _>>();
1795 let mut root_updates = HashMap::<usize, u64>::new();
1796 let mut root_child_counts = HashMap::<usize, usize>::new();
1797
1798 for descriptor in found.iter() {
1799 let Some(mut parent_id) = descriptor.parent_session_id.as_deref() else {
1800 continue;
1801 };
1802 let harness = descriptor.locator.harness.as_str();
1803 let mut root = None;
1804 let mut visited = HashSet::new();
1805 while visited.insert(parent_id.to_string()) {
1806 let Some(&parent_index) = by_id.get(&(harness.to_string(), parent_id.to_string()))
1807 else {
1808 break;
1809 };
1810 root = Some(parent_index);
1811 let Some(next_parent) = found[parent_index].parent_session_id.as_deref() else {
1812 break;
1813 };
1814 parent_id = next_parent;
1815 }
1816 if let (Some(root), Some(updated_at_ms)) = (root, descriptor.updated_at_ms) {
1817 root_updates
1818 .entry(root)
1819 .and_modify(|current| *current = (*current).max(updated_at_ms))
1820 .or_insert(updated_at_ms);
1821 }
1822 if let Some(root) = root {
1823 *root_child_counts.entry(root).or_default() += 1;
1824 }
1825 }
1826
1827 for (root, child_updated_at_ms) in root_updates {
1828 found[root].updated_at_ms = Some(
1829 found[root]
1830 .updated_at_ms
1831 .unwrap_or_default()
1832 .max(child_updated_at_ms),
1833 );
1834 }
1835 for (root, child_count) in root_child_counts {
1836 found[root].child_session_count = child_count;
1837 }
1838 if !include_children {
1839 found.retain(|descriptor| descriptor.parent_session_id.is_none());
1840 }
1841}
1842
1843fn retain_session_family(found: &mut Vec<SessionDescriptor>, root_session_id: &str) {
1844 let parent_by_id = found
1845 .iter()
1846 .map(|descriptor| {
1847 (
1848 descriptor.locator.session_id.clone(),
1849 descriptor.parent_session_id.clone(),
1850 )
1851 })
1852 .collect::<HashMap<_, _>>();
1853 found.retain(|descriptor| {
1854 let mut current = descriptor.locator.session_id.clone();
1855 let mut visited = HashSet::new();
1856 while visited.insert(current.clone()) {
1857 if current == root_session_id {
1858 return true;
1859 }
1860 let Some(Some(parent)) = parent_by_id.get(¤t) else {
1861 return false;
1862 };
1863 current = parent.clone();
1864 }
1865 false
1866 });
1867}
1868
1869fn claude_subagent_parent_id(path: &Path) -> Option<String> {
1870 let subagents = path.parent()?;
1871 if subagents.file_name()?.to_str()? != "subagents" {
1872 return None;
1873 }
1874 subagents
1875 .parent()?
1876 .file_name()?
1877 .to_str()
1878 .map(str::to_string)
1879}
1880
1881fn count_claude_subagents(parent_path: &Path) -> usize {
1882 let Some(parent) = parent_path.parent() else {
1883 return 0;
1884 };
1885 let Some(stem) = parent_path.file_stem() else {
1886 return 0;
1887 };
1888 let root = parent.join(stem).join("subagents");
1889 let mut files = Vec::new();
1890 collect_jsonl(&root, HarnessId::CLAUDE_CODE, true, &mut files);
1891 files.len()
1892}
1893
1894fn is_zero(value: &usize) -> bool {
1895 *value == 0
1896}
1897
1898fn humanize_topic(value: &str) -> String {
1899 let text = value.replace(['_', '-'], " ");
1900 let mut characters = text.chars();
1901 match characters.next() {
1902 Some(first) => first.to_uppercase().collect::<String>() + characters.as_str(),
1903 None => text,
1904 }
1905}
1906
1907fn discover_gemini(
1908 root: &Path,
1909 workspace: Option<&WorkspaceScope>,
1910 found: &mut Vec<SessionDescriptor>,
1911) {
1912 let slug_to_cwd = std::fs::read_to_string(root.join("projects.json"))
1913 .ok()
1914 .and_then(|text| serde_json::from_str::<Value>(&text).ok())
1915 .and_then(|value| value.get("projects").and_then(Value::as_object).cloned())
1916 .map(|projects| {
1917 projects
1918 .into_iter()
1919 .filter_map(|(cwd, slug)| Some((slug.as_str()?.to_string(), PathBuf::from(cwd))))
1920 .collect::<HashMap<_, _>>()
1921 })
1922 .unwrap_or_default();
1923 let mut files = Vec::new();
1924 collect_jsonl(&root.join("tmp"), HarnessId::GEMINI, false, &mut files);
1925 let worker_count = std::thread::available_parallelism()
1926 .map(usize::from)
1927 .unwrap_or(4)
1928 .clamp(1, 8)
1929 .min(files.len().max(1));
1930 let chunk_size = files.len().max(1).div_ceil(worker_count);
1931 let discovered = std::thread::scope(|scope| {
1932 files
1933 .chunks(chunk_size)
1934 .map(|paths| {
1935 scope.spawn(|| {
1936 paths
1937 .iter()
1938 .filter_map(|path| gemini_descriptor(path, &slug_to_cwd, workspace))
1939 .collect::<Vec<_>>()
1940 })
1941 })
1942 .collect::<Vec<_>>()
1943 .into_iter()
1944 .flat_map(|worker| {
1945 worker
1946 .join()
1947 .expect("Gemini discovery worker must not panic")
1948 })
1949 .collect::<Vec<_>>()
1950 });
1951 found.extend(discovered);
1952}
1953
1954fn gemini_descriptor(
1955 path: &Path,
1956 slug_to_cwd: &HashMap<String, PathBuf>,
1957 workspace: Option<&WorkspaceScope>,
1958) -> Option<SessionDescriptor> {
1959 if path
1960 .parent()
1961 .and_then(Path::file_name)
1962 .and_then(|name| name.to_str())
1963 != Some("chats")
1964 {
1965 return None;
1966 }
1967 let slug = path
1968 .parent()
1969 .and_then(Path::parent)
1970 .and_then(Path::file_name)
1971 .and_then(|name| name.to_str());
1972 let cwd = slug.and_then(|slug| slug_to_cwd.get(slug)).cloned();
1973 if workspace.is_some_and(|wanted| cwd.as_deref().is_none_or(|actual| !wanted.admits(actual))) {
1974 return None;
1975 }
1976
1977 let file = File::open(path).ok()?;
1981 let mut reader = BufReader::new(file.take(64 * 1024));
1982 let mut header = String::new();
1983 reader.read_line(&mut header).ok()?;
1984 let header = serde_json::from_str::<Value>(&header).ok()?;
1985 let session_id = header.get("sessionId")?.as_str()?.to_string();
1986 let mut model = None;
1987 for line in reader
1988 .take(4 * 1024)
1989 .lines()
1990 .map_while(std::result::Result::ok)
1991 {
1992 let Ok(value) = serde_json::from_str::<Value>(&line) else {
1993 continue;
1994 };
1995 let kind = value.get("type").and_then(Value::as_str);
1996 if kind != Some("user") && kind != Some("gemini") {
1997 continue;
1998 }
1999 if model.is_none() {
2000 model = value
2001 .get("model")
2002 .and_then(Value::as_str)
2003 .map(str::to_string);
2004 }
2005 if model.is_some() {
2006 break;
2007 }
2008 }
2009 Some(SessionDescriptor {
2010 locator: SessionLocator {
2011 harness: HarnessId::from(HarnessId::GEMINI),
2012 session_id,
2013 storage: StorageLocator::File {
2014 path: path.to_path_buf(),
2015 },
2016 },
2017 cwd,
2018 title: None,
2019 preview_candidates: Vec::new(),
2020 latest_message_candidates: Vec::new(),
2021 updated_at_ms: tail_facts(path, HarnessId::GEMINI)
2022 .last_turn_ms
2023 .or_else(|| modified_ms(path)),
2024 message_count: None,
2025 model,
2026 parent_session_id: None,
2027 child_session_count: 0,
2028 nouns: OrchestrationNouns::default(),
2029 })
2030}
2031
2032fn display_text(content: Option<&Value>) -> Option<String> {
2033 match content? {
2034 Value::String(text) => Some(text.clone()),
2035 Value::Array(parts) => Some(
2036 parts
2037 .iter()
2038 .filter_map(|part| part.get("text").and_then(Value::as_str))
2039 .collect::<Vec<_>>()
2040 .join(" ")
2041 .trim()
2042 .to_string(),
2043 ),
2044 _ => None,
2045 }
2046}
2047
2048fn discover_hermes(
2053 db_path: &Path,
2054 workspace: Option<&WorkspaceScope>,
2055 selected: &HashSet<&str>,
2056 found: &mut Vec<SessionDescriptor>,
2057) {
2058 if !db_path.is_file() {
2059 return;
2060 }
2061 let Ok(conn) = Connection::open_with_flags(
2062 db_path,
2063 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
2064 ) else {
2065 return;
2066 };
2067 let fingerprint_ok = ["sessions", "messages", "schema_version"].iter().all(|t| {
2068 conn.query_row(
2069 "SELECT 1 FROM sqlite_master WHERE type='table' AND name=?1",
2070 [t],
2071 |_| Ok(()),
2072 )
2073 .is_ok()
2074 });
2075 if !fingerprint_ok {
2076 return;
2077 }
2078 let Ok(mut statement) = conn.prepare(
2079 "SELECT id, cwd, title, model, message_count, started_at, ended_at, parent_session_id, \
2080 source, model_config FROM sessions ORDER BY started_at DESC",
2081 ) else {
2082 return;
2083 };
2084 let Ok(rows) = statement.query_map([], |row| {
2085 Ok((
2086 row.get::<_, String>(0)?,
2087 row.get::<_, Option<String>>(1)?,
2088 row.get::<_, Option<String>>(2)?,
2089 row.get::<_, Option<String>>(3)?,
2090 row.get::<_, Option<i64>>(4)?,
2091 row.get::<_, Option<f64>>(5)?,
2092 row.get::<_, Option<f64>>(6)?,
2093 row.get::<_, Option<String>>(7)?,
2094 row.get::<_, Option<String>>(8)?,
2095 row.get::<_, Option<String>>(9)?,
2096 ))
2097 }) else {
2098 return;
2099 };
2100 for row in rows.flatten() {
2101 let (
2102 id,
2103 cwd,
2104 title,
2105 model,
2106 message_count,
2107 started_at,
2108 ended_at,
2109 parent,
2110 source,
2111 model_config,
2112 ) = row;
2113 let mirror = model_config
2117 .as_deref()
2118 .and_then(|c| serde_json::from_str::<serde_json::Value>(c).ok())
2119 .and_then(|c| c.get("_supercode_mirror").cloned());
2120 if let Some(mirror) = mirror {
2121 let worker_read = mirror
2122 .get("harness")
2123 .and_then(serde_json::Value::as_str)
2124 .is_some_and(|h| selected.contains(h));
2125 let mirrored = mirror.get("messages").and_then(serde_json::Value::as_i64);
2126 if worker_read && mirrored.is_none_or(|m| message_count.unwrap_or(0) <= m) {
2127 continue;
2128 }
2129 if mirror.get("continued_as").is_some()
2130 && mirrored.is_none_or(|m| message_count.unwrap_or(0) <= m)
2131 {
2132 continue;
2133 }
2134 }
2135 let cwd = cwd.map(PathBuf::from);
2136 if let Some(filter) = workspace {
2137 if !filter.admits_literal(cwd.as_deref()) {
2138 continue;
2139 }
2140 }
2141 let updated_at_ms = ended_at
2142 .or(started_at)
2143 .map(|seconds| (seconds * 1000.0) as u64);
2144 let mut meta = SessionMeta::new(SessionSource::Hermes);
2148 meta.cwd = cwd.clone();
2149 if let Some(hermes_source) = source.filter(|value| !value.is_empty()) {
2150 meta.lineage
2151 .insert("hermes_source".to_string(), hermes_source);
2152 }
2153 if let Some(parent_id) = parent.as_deref() {
2154 meta.lineage.insert(
2155 "hermes_lineage_kind".to_string(),
2156 crate::session::hermes_lineage_kind(
2157 &conn,
2158 parent_id,
2159 model_config.as_deref(),
2160 started_at,
2161 )
2162 .to_string(),
2163 );
2164 }
2165 hermes_capture_nouns(&conn, &id, &mut meta);
2166 found.push(SessionDescriptor {
2167 locator: SessionLocator {
2168 harness: HarnessId::new(HarnessId::HERMES),
2169 session_id: id,
2170 storage: StorageLocator::File {
2171 path: db_path.to_path_buf(),
2172 },
2173 },
2174 cwd,
2175 title: title.filter(|t| !t.is_empty()),
2176 preview_candidates: Vec::new(),
2177 latest_message_candidates: Vec::new(),
2178 updated_at_ms,
2179 message_count: message_count.map(|count| count.max(0) as usize),
2180 model,
2181 parent_session_id: parent,
2182 child_session_count: 0,
2183 nouns: OrchestrationNouns::from_meta(&meta),
2184 });
2185 }
2186}
2187
2188fn discover_orchestrator(
2195 root: &Path,
2196 workspace: Option<&WorkspaceScope>,
2197 found: &mut Vec<SessionDescriptor>,
2198) {
2199 if workspace.is_some() {
2201 return;
2202 }
2203 for (profile, dir) in orchestrator_profile_dirs(root) {
2204 let db_path = dir.join("state.db");
2205 if !db_path.is_file() {
2206 continue;
2207 }
2208 let Ok(conn) = Connection::open_with_flags(
2209 &db_path,
2210 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
2211 ) else {
2212 continue;
2213 };
2214 let sessions = dir.join("sessions");
2217 let scope = std::fs::canonicalize(&sessions)
2218 .unwrap_or(sessions)
2219 .display()
2220 .to_string();
2221 let Ok(mut statement) = conn.prepare(
2222 "SELECT json_extract(entry_json, '$.metadata.supercode.binding'), \
2223 CAST(strftime('%s', json_extract(entry_json, '$.metadata.supercode.binding.last_activity_at')) AS INTEGER) \
2224 FROM gateway_routing WHERE scope = ?1 AND json_extract(entry_json, '$.metadata.supercode') IS NOT NULL \
2225 ORDER BY json_extract(entry_json, '$.metadata.supercode.binding.last_activity_at') DESC",
2226 ) else {
2227 continue;
2228 };
2229 let Ok(rows) = statement.query_map([&scope], |row| {
2230 Ok((
2231 row.get::<_, Option<String>>(0)?,
2232 row.get::<_, Option<i64>>(1)?,
2233 ))
2234 }) else {
2235 continue;
2236 };
2237 let rows = rows.flatten().filter_map(|(json, epoch)| {
2238 let b: Binding = serde_json::from_str(&json?).ok()?;
2239 let text = |v: &Option<String>| v.clone().filter(|s| !s.is_empty());
2240 Some((
2241 OrchestratorBindingRow {
2242 platform: b.key.platform.clone().unwrap_or_default(),
2243 chat_type: b.key.kind.clone().unwrap_or_default(),
2244 chat_id: text(&b.key.chat_id),
2245 thread_id: text(&b.key.thread_id),
2246 participant_id: text(&b.key.participant_id),
2247 worker_harness: b.worker.harness.as_str().to_string(),
2248 worker_session_id: text(&b.worker.session_id),
2249 worker_locator: text(&b.worker.locator),
2250 started_at: b.started_at.clone(),
2251 last_activity_at: b.last_activity_at.clone(),
2252 ended_at: b.ended_at.clone(),
2253 end_reason: b.end_reason.map(|r| r.as_str().to_string()),
2254 handoff_to: b.handoff.as_ref().and_then(|h| h.to.clone()),
2255 handoff_state: b.handoff.as_ref().map(|h| h.state.clone()),
2256 handoff_error: b.handoff.as_ref().and_then(|h| h.error.clone()),
2257 recurrence_job_id: b.recurrence.as_ref().map(|r| r.job_id.clone()),
2258 },
2259 epoch,
2260 ))
2261 });
2262 for (row, last_activity_epoch) in rows {
2263 found.push(orchestrator_descriptor(
2264 &db_path,
2265 &profile,
2266 &row,
2267 last_activity_epoch,
2268 ));
2269 }
2270 }
2271}
2272
2273fn orchestrator_descriptor(
2274 db_path: &Path,
2275 profile: &str,
2276 row: &OrchestratorBindingRow,
2277 last_activity_epoch: Option<i64>,
2278) -> SessionDescriptor {
2279 let binding = Binding::from_orchestrator_row(profile, row);
2280 let nouns = binding.nouns();
2281 let mut title = format!(
2284 "{} {}",
2285 row.worker_harness,
2286 row.worker_session_id
2287 .as_deref()
2288 .unwrap_or("(no worker session yet)")
2289 );
2290 if let Some(reason) = row.end_reason.as_deref().filter(|_| row.ended_at.is_some()) {
2291 title.push_str(&format!(" (ended: {reason})"));
2292 }
2293 SessionDescriptor {
2294 locator: SessionLocator {
2295 harness: HarnessId::new(HarnessId::ORCHESTRATOR),
2296 session_id: row.worker_session_id.clone().unwrap_or_default(),
2297 storage: StorageLocator::File {
2298 path: row
2299 .worker_locator
2300 .clone()
2301 .map_or_else(|| db_path.to_path_buf(), PathBuf::from),
2302 },
2303 },
2304 cwd: None,
2305 title: Some(title),
2306 preview_candidates: Vec::new(),
2307 latest_message_candidates: Vec::new(),
2308 updated_at_ms: last_activity_epoch.map(|seconds| (seconds.max(0) as u64) * 1000),
2309 message_count: None,
2310 model: None,
2311 parent_session_id: None,
2312 child_session_count: 0,
2313 nouns,
2314 }
2315}
2316
2317fn discover_openclaw(
2323 root: &Path,
2324 workspace: Option<&WorkspaceScope>,
2325 found: &mut Vec<SessionDescriptor>,
2326) {
2327 let agents = root.join("agents");
2328 let Ok(agent_dirs) = std::fs::read_dir(&agents) else {
2329 return;
2330 };
2331 for agent_dir in agent_dirs.flatten() {
2332 let sessions = agent_dir.path().join("sessions");
2333 let Ok(files) = std::fs::read_dir(&sessions) else {
2334 continue;
2335 };
2336 for file in files.flatten() {
2337 let path = file.path();
2338 let name = file.file_name();
2339 let name = name.to_string_lossy();
2340 if !name.ends_with(".jsonl") || name.ends_with(".trajectory.jsonl") {
2341 continue;
2342 }
2343 let Ok(text) = std::fs::read_to_string(&path) else {
2344 continue;
2345 };
2346 let Some(header_line) = text.lines().find(|line| !line.trim().is_empty()) else {
2347 continue;
2348 };
2349 let Ok(header) = serde_json::from_str::<serde_json::Value>(header_line) else {
2350 continue;
2351 };
2352 if header.get("type").and_then(serde_json::Value::as_str) != Some("session") {
2353 continue;
2354 }
2355 let session_id = header
2356 .get("id")
2357 .and_then(serde_json::Value::as_str)
2358 .unwrap_or_else(|| name.trim_end_matches(".jsonl"))
2359 .to_string();
2360 let cwd = header
2361 .get("cwd")
2362 .and_then(serde_json::Value::as_str)
2363 .map(PathBuf::from);
2364 if let Some(filter) = workspace {
2365 if !filter.admits_literal(cwd.as_deref()) {
2366 continue;
2367 }
2368 }
2369 let updated_at_ms = tail_facts(&path, HarnessId::OPENCLAW)
2370 .last_turn_ms
2371 .or_else(|| {
2372 file.metadata()
2373 .ok()
2374 .and_then(|metadata| metadata.modified().ok())
2375 .and_then(|modified| modified.duration_since(std::time::UNIX_EPOCH).ok())
2376 .map(|elapsed| elapsed.as_millis() as u64)
2377 });
2378 let message_count = text
2379 .lines()
2380 .filter(|line| line.contains("\"type\":\"message\""))
2381 .count();
2382 let mut meta = SessionMeta::new(SessionSource::OpenClaw);
2386 meta.cwd = cwd.clone();
2387 openclaw_capture_header_nouns(&header, &mut meta);
2388 if meta.profile.is_none() {
2389 meta.profile = openclaw_agent_id_from_path(&path);
2390 }
2391 found.push(SessionDescriptor {
2392 locator: SessionLocator {
2393 harness: HarnessId::new(HarnessId::OPENCLAW),
2394 session_id,
2395 storage: StorageLocator::File { path },
2396 },
2397 cwd,
2398 title: None,
2399 preview_candidates: Vec::new(),
2400 latest_message_candidates: Vec::new(),
2401 updated_at_ms,
2402 message_count: Some(message_count),
2403 model: None,
2404 parent_session_id: None,
2405 child_session_count: 0,
2406 nouns: OrchestrationNouns::from_meta(&meta),
2407 });
2408 }
2409 }
2410}
2411
2412fn discover_supercode(
2413 root: &Path,
2414 workspace: Option<&WorkspaceScope>,
2415 found: &mut Vec<SessionDescriptor>,
2416) {
2417 for info in list_native_store(root) {
2418 let path = if info.archived {
2419 root.join("archived").join(format!("{}.jsonl", info.name))
2420 } else {
2421 root.join(format!("{}.jsonl", info.name))
2422 };
2423 let sidecar = path.with_extension("sidecar.jsonl");
2428 let path = if path.is_file() {
2429 path
2430 } else {
2431 sidecar.clone()
2432 };
2433 let header = read_native_store_header(&path);
2434 if workspace.is_some_and(|wanted| {
2435 header
2436 .as_ref()
2437 .and_then(|meta| meta.cwd.as_deref())
2438 .is_none_or(|cwd| !wanted.admits(cwd))
2439 }) {
2440 continue;
2441 }
2442 let title = (!info.title.trim().is_empty()).then_some(info.title);
2443 let updated_at_ms = tail_facts(&sidecar, HarnessId::SUPERCODE)
2446 .last_turn_ms
2447 .or_else(|| tail_facts(&path, HarnessId::SUPERCODE).last_turn_ms)
2448 .or_else(|| modified_ms(&path))
2449 .or_else(|| modified_ms(&sidecar));
2450 found.push(SessionDescriptor {
2451 locator: SessionLocator {
2452 harness: HarnessId::from(HarnessId::SUPERCODE),
2453 session_id: info.name,
2454 storage: StorageLocator::File { path: path.clone() },
2455 },
2456 cwd: header.as_ref().and_then(|meta| meta.cwd.clone()),
2457 title,
2458 preview_candidates: Vec::new(),
2459 latest_message_candidates: Vec::new(),
2460 updated_at_ms,
2461 message_count: None,
2462 model: header.and_then(|meta| meta.model),
2463 parent_session_id: None,
2464 child_session_count: 0,
2465 nouns: OrchestrationNouns::default(),
2466 });
2467 }
2468}
2469
2470fn read_native_store_header(path: &Path) -> Option<HeaderMeta> {
2475 let name = path.file_stem()?.to_str()?;
2476 let sidecar = path.with_file_name(format!("{name}.sidecar.jsonl"));
2477 let source_path = if sidecar.is_file() {
2478 sidecar
2479 } else {
2480 path.to_path_buf()
2481 };
2482 let file = File::open(source_path).ok()?;
2483 let mut result = HeaderMeta::default();
2484 let mut source = None;
2485 let mut bytes = 0usize;
2486 for line in BufReader::new(file).lines().take(32) {
2487 let line = line.ok()?;
2488 bytes += line.len();
2489 if bytes > 256 * 1024 {
2490 break;
2491 }
2492 let Ok(value) = serde_json::from_str::<Value>(&line) else {
2493 continue;
2494 };
2495 if source.is_none() {
2496 source = value.get("source").and_then(Value::as_str).map(|source| {
2497 if source == "claude_code" {
2498 HarnessId::CLAUDE_CODE.to_string()
2499 } else {
2500 source.to_string()
2501 }
2502 });
2503 fill_string(&mut result.session_id, value.get("session_id"));
2504 }
2505 if let Some(harness) = source.as_deref() {
2506 update_header_meta(&mut result, &value, harness);
2507 }
2508 if result.cwd.is_some() && result.model.is_some() {
2509 break;
2510 }
2511 }
2512 Some(result)
2513}
2514
2515#[derive(Deserialize)]
2516struct NativeStoreInfo {
2517 name: String,
2518 #[serde(default)]
2519 title: String,
2520 #[serde(skip)]
2521 archived: bool,
2522}
2523
2524fn list_native_store(root: &Path) -> Vec<NativeStoreInfo> {
2525 let mut sessions = Vec::new();
2526 for archived in [false, true] {
2527 let directory = if archived {
2528 root.join("archived")
2529 } else {
2530 root.to_path_buf()
2531 };
2532 let Ok(entries) = fs::read_dir(directory) else {
2533 continue;
2534 };
2535 for entry in entries.flatten() {
2536 let path = entry.path();
2537 if !path.to_string_lossy().ends_with(".meta.json") {
2538 continue;
2539 }
2540 let Ok(text) = fs::read_to_string(path) else {
2541 continue;
2542 };
2543 let Ok(mut info) = serde_json::from_str::<NativeStoreInfo>(&text) else {
2544 continue;
2545 };
2546 info.archived = archived;
2547 sessions.push(info);
2548 }
2549 }
2550 sessions.sort_by(|left, right| left.name.cmp(&right.name));
2551 sessions
2552}
2553
2554fn discover_grok(
2555 root: &Path,
2556 workspace: Option<&WorkspaceScope>,
2557 found: &mut Vec<SessionDescriptor>,
2558) {
2559 let Ok(workspaces) = fs::read_dir(root) else {
2560 return;
2561 };
2562 for workspace_entry in workspaces.flatten() {
2563 let encoded = workspace_entry.file_name();
2564 let Some(cwd) = encoded
2565 .to_str()
2566 .and_then(percent_decode_path)
2567 .map(PathBuf::from)
2568 else {
2569 continue;
2570 };
2571 if workspace.is_some_and(|wanted| !wanted.admits(&cwd)) {
2572 continue;
2573 }
2574 let Ok(sessions) = fs::read_dir(workspace_entry.path()) else {
2575 continue;
2576 };
2577 for session_entry in sessions.flatten() {
2578 let session_dir = session_entry.path();
2579 if !session_dir.is_dir() {
2580 continue;
2581 }
2582 let transcript = session_dir.join("chat_history.jsonl");
2583 if !transcript.is_file() {
2584 continue;
2585 }
2586 let Some(session_id) = session_dir
2587 .file_name()
2588 .and_then(|name| name.to_str())
2589 .map(str::to_string)
2590 else {
2591 continue;
2592 };
2593 let summary = fs::read_to_string(session_dir.join("summary.json"))
2594 .ok()
2595 .and_then(|text| serde_json::from_str::<Value>(&text).ok());
2596 let title = summary
2597 .as_ref()
2598 .and_then(|value| value.get("generated_title"))
2599 .and_then(Value::as_str)
2600 .filter(|title| !title.is_empty())
2601 .map(str::to_string);
2602 let model = summary
2603 .as_ref()
2604 .and_then(|value| value.get("current_model_id"))
2605 .and_then(Value::as_str)
2606 .map(str::to_string);
2607 let message_count = summary
2608 .as_ref()
2609 .and_then(|value| value.get("num_chat_messages"))
2610 .and_then(Value::as_u64)
2611 .and_then(|count| usize::try_from(count).ok());
2612 let updated_at_ms = summary
2613 .as_ref()
2614 .and_then(|value| value.get("updated_at"))
2615 .and_then(Value::as_str)
2616 .and_then(crate::sidecar::rfc3339_to_ms)
2617 .and_then(|millis| u64::try_from(millis).ok())
2618 .or_else(|| modified_ms(&transcript));
2619 found.push(SessionDescriptor {
2620 locator: SessionLocator {
2621 harness: HarnessId::from(HarnessId::GROK),
2622 session_id,
2623 storage: StorageLocator::File { path: transcript },
2624 },
2625 cwd: Some(cwd.clone()),
2626 title,
2627 preview_candidates: Vec::new(),
2628 latest_message_candidates: Vec::new(),
2629 updated_at_ms,
2630 message_count,
2631 model,
2632 parent_session_id: None,
2633 child_session_count: 0,
2634 nouns: OrchestrationNouns::default(),
2635 });
2636 }
2637 }
2638}
2639
2640fn discover_opencode(
2641 root: &Path,
2642 workspace: Option<&WorkspaceScope>,
2643 found: &mut Vec<SessionDescriptor>,
2644) {
2645 let mut dbs = Vec::new();
2646 if root.is_file() {
2647 dbs.push(root.to_path_buf());
2648 } else if let Ok(entries) = fs::read_dir(root) {
2649 dbs.extend(entries.flatten().map(|entry| entry.path()).filter(|path| {
2650 path.file_name()
2651 .and_then(|v| v.to_str())
2652 .is_some_and(|name| name.starts_with("opencode") && name.ends_with(".db"))
2653 }));
2654 }
2655 dbs.sort();
2656 for db in dbs {
2657 let Ok(conn) = Connection::open_with_flags(
2658 &db,
2659 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
2660 ) else {
2661 continue;
2662 };
2663 let has_model = conn.prepare("SELECT model FROM session LIMIT 0").is_ok();
2664 let model_column = if has_model { "s.model" } else { "NULL" };
2665 let query = format!(
2666 "SELECT s.id, s.directory, s.title, s.time_updated, {model_column}, COUNT(m.id) \
2667 FROM session s LEFT JOIN message m ON m.session_id = s.id \
2668 GROUP BY s.id ORDER BY s.time_updated DESC"
2669 );
2670 let Ok(mut stmt) = conn.prepare(&query) else {
2671 continue;
2672 };
2673 let Ok(rows) = stmt.query_map([], |row| {
2674 Ok((
2675 row.get::<_, String>(0)?,
2676 row.get::<_, String>(1)?,
2677 row.get::<_, String>(2)?,
2678 row.get::<_, i64>(3)?,
2679 row.get::<_, Option<String>>(4)?,
2680 row.get::<_, i64>(5)?,
2681 ))
2682 }) else {
2683 continue;
2684 };
2685 for row in rows.flatten() {
2686 let (id, cwd, title, updated, model, messages) = row;
2687 let cwd = PathBuf::from(cwd);
2688 if workspace.is_some_and(|wanted| !wanted.admits(&cwd)) {
2689 continue;
2690 }
2691 found.push(SessionDescriptor {
2692 locator: SessionLocator {
2693 harness: HarnessId::from(HarnessId::OPENCODE),
2694 session_id: id.clone(),
2695 storage: StorageLocator::Sqlite {
2696 path: db.clone(),
2697 selector: id,
2698 },
2699 },
2700 cwd: Some(cwd),
2701 title: (!title.is_empty()).then_some(title),
2702 preview_candidates: Vec::new(),
2703 latest_message_candidates: Vec::new(),
2704 updated_at_ms: u64::try_from(updated).ok(),
2705 message_count: usize::try_from(messages).ok(),
2706 model,
2707 parent_session_id: None,
2708 child_session_count: 0,
2709 nouns: OrchestrationNouns::default(),
2710 });
2711 }
2712 }
2713}
2714
2715fn discover_goose(
2716 root: &Path,
2717 workspace: Option<&WorkspaceScope>,
2718 found: &mut Vec<SessionDescriptor>,
2719) {
2720 let db = if root.is_file() {
2721 root.to_path_buf()
2722 } else if root.join("sessions.db").is_file() {
2723 root.join("sessions.db")
2724 } else {
2725 root.join("sessions/sessions.db")
2726 };
2727 let Ok(connection) = Connection::open_with_flags(
2728 &db,
2729 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
2730 ) else {
2731 return;
2732 };
2733 let Ok(mut statement) = connection.prepare(
2734 "SELECT s.id, s.working_dir, s.name, s.updated_at, s.model_config_json, \
2735 COUNT(m.id) \
2736 FROM sessions s LEFT JOIN messages m ON m.session_id = s.id \
2737 WHERE s.archived_at IS NULL \
2738 GROUP BY s.id ORDER BY s.updated_at DESC",
2739 ) else {
2740 return;
2741 };
2742 let Ok(rows) = statement.query_map([], |row| {
2743 Ok((
2744 row.get::<_, String>(0)?,
2745 row.get::<_, String>(1)?,
2746 row.get::<_, String>(2)?,
2747 row.get::<_, String>(3)?,
2748 row.get::<_, Option<String>>(4)?,
2749 row.get::<_, i64>(5)?,
2750 ))
2751 }) else {
2752 return;
2753 };
2754 for row in rows.flatten() {
2755 let (id, cwd, title, updated_at, model_config, message_count) = row;
2756 let cwd = PathBuf::from(cwd);
2757 if workspace.is_some_and(|wanted| !wanted.admits(&cwd)) {
2758 continue;
2759 }
2760 let model = model_config
2761 .as_deref()
2762 .and_then(|value| serde_json::from_str::<Value>(value).ok())
2763 .and_then(|value| {
2764 value
2765 .get("model_name")
2766 .or_else(|| value.get("modelName"))
2767 .and_then(Value::as_str)
2768 .map(str::to_string)
2769 });
2770 let updated_at_ms = crate::sidecar::rfc3339_to_ms(&updated_at)
2771 .or_else(|| {
2772 crate::sidecar::rfc3339_to_ms(&format!("{}Z", updated_at.replace(' ', "T")))
2774 })
2775 .and_then(|value| u64::try_from(value).ok());
2776 found.push(SessionDescriptor {
2777 locator: SessionLocator {
2778 harness: HarnessId::from(HarnessId::GOOSE),
2779 session_id: id.clone(),
2780 storage: StorageLocator::Sqlite {
2781 path: db.clone(),
2782 selector: id,
2783 },
2784 },
2785 cwd: Some(cwd),
2786 title: (!title.trim().is_empty()).then_some(title),
2787 preview_candidates: Vec::new(),
2788 latest_message_candidates: Vec::new(),
2789 updated_at_ms,
2790 message_count: usize::try_from(message_count).ok(),
2791 model,
2792 parent_session_id: None,
2793 child_session_count: 0,
2794 nouns: OrchestrationNouns::default(),
2795 });
2796 }
2797}
2798
2799const LATEST_PREVIEW_CANDIDATES: usize = 8;
2800const TOPIC_PREVIEW_HEAD_BYTES: u64 = 512 * 1024;
2801const LATEST_PREVIEW_TAIL_BYTES: u64 = 512 * 1024;
2802const LATEST_PREVIEW_MAX_BYTES: u64 = 4 * 1024 * 1024;
2803
2804fn topic_message_candidates(locator: &SessionLocator) -> Result<Vec<SessionPreviewCandidate>> {
2805 match &locator.storage {
2806 StorageLocator::File { path }
2807 if matches!(
2808 locator.harness.as_str(),
2809 HarnessId::CLAUDE_CODE | HarnessId::CODEX
2810 ) =>
2811 {
2812 topic_file_message_candidates(path, locator.harness.as_str())
2813 }
2814 _ => Ok(Vec::new()),
2815 }
2816}
2817
2818fn codex_history_topics(
2819 sessions_root: &Path,
2820 sessions: &[SessionDescriptor],
2821) -> Result<HashMap<String, Vec<SessionPreviewCandidate>>> {
2822 let wanted: HashSet<&str> = sessions
2823 .iter()
2824 .filter(|descriptor| descriptor.locator.harness.as_str() == HarnessId::CODEX)
2825 .map(|descriptor| descriptor.locator.session_id.as_str())
2826 .collect();
2827 if wanted.is_empty() {
2828 return Ok(HashMap::new());
2829 }
2830 let Some(root) = sessions_root.parent() else {
2831 return Ok(HashMap::new());
2832 };
2833 let file = match File::open(root.join("history.jsonl")) {
2834 Ok(file) => file,
2835 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(HashMap::new()),
2836 Err(error) => return Err(error.into()),
2837 };
2838 let mut topics = HashMap::new();
2839 for line in BufReader::new(file).lines() {
2840 let Ok(value) = serde_json::from_str::<Value>(&line?) else {
2841 continue;
2842 };
2843 let Some(session_id) = value.get("session_id").and_then(Value::as_str) else {
2844 continue;
2845 };
2846 if !wanted.contains(session_id) || topics.contains_key(session_id) {
2847 continue;
2848 }
2849 let mut candidates = Vec::new();
2850 push_message_candidate(&mut candidates, "user", value.get("text"), HashMap::new());
2851 if !candidates.is_empty() {
2852 topics.insert(session_id.to_string(), candidates);
2853 if topics.len() == wanted.len() {
2854 break;
2855 }
2856 }
2857 }
2858 Ok(topics)
2859}
2860
2861fn latest_message_candidates(locator: &SessionLocator) -> Result<Vec<SessionPreviewCandidate>> {
2862 match &locator.storage {
2863 StorageLocator::File { path } | StorageLocator::Sqlite { path, .. }
2866 if locator.harness.as_str() == HarnessId::HERMES =>
2867 {
2868 latest_hermes_message_candidates(path, &locator.session_id)
2869 }
2870 StorageLocator::File { path } => {
2871 latest_file_message_candidates(path, locator.harness.as_str())
2872 }
2873 StorageLocator::Sqlite { path, selector }
2874 if locator.harness.as_str() == HarnessId::OPENCODE =>
2875 {
2876 latest_opencode_message_candidates(path, selector)
2877 }
2878 StorageLocator::Sqlite { path, selector }
2879 if locator.harness.as_str() == HarnessId::GOOSE =>
2880 {
2881 latest_goose_message_candidates(path, selector)
2882 }
2883 StorageLocator::Sqlite { .. } => Ok(Vec::new()),
2884 }
2885}
2886
2887fn topic_file_message_candidates(
2888 path: &Path,
2889 harness: &str,
2890) -> Result<Vec<SessionPreviewCandidate>> {
2891 let mut file = File::open(path)?;
2892 let mut bytes = Vec::with_capacity(TOPIC_PREVIEW_HEAD_BYTES as usize);
2893 file.by_ref()
2894 .take(TOPIC_PREVIEW_HEAD_BYTES)
2895 .read_to_end(&mut bytes)?;
2896 if file.metadata()?.len() > TOPIC_PREVIEW_HEAD_BYTES {
2897 if let Some(newline) = bytes.iter().rposition(|byte| *byte == b'\n') {
2898 bytes.truncate(newline);
2899 }
2900 }
2901 let text = String::from_utf8(bytes).map_err(|_| {
2902 Error::Other(format!(
2903 "{} contains non-UTF-8 data in its topic-preview window",
2904 path.display()
2905 ))
2906 })?;
2907 if harness == HarnessId::CODEX {
2908 return Ok(codex_preview_candidates(text.lines(), false));
2909 }
2910 let mut candidates = Vec::new();
2911 for line in text.lines() {
2912 let Ok(value) = serde_json::from_str::<Value>(line) else {
2913 continue;
2914 };
2915 push_topic_message_candidate(&mut candidates, harness, &value);
2916 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
2917 break;
2918 }
2919 }
2920 Ok(candidates)
2921}
2922
2923fn latest_file_message_candidates(
2924 path: &Path,
2925 harness: &str,
2926) -> Result<Vec<SessionPreviewCandidate>> {
2927 let mut candidates =
2928 latest_file_message_candidates_with_limit(path, harness, LATEST_PREVIEW_TAIL_BYTES)?;
2929 if candidates.is_empty() {
2930 candidates =
2931 latest_file_message_candidates_with_limit(path, harness, LATEST_PREVIEW_MAX_BYTES)?;
2932 }
2933 Ok(candidates)
2934}
2935
2936fn latest_file_message_candidates_with_limit(
2937 path: &Path,
2938 harness: &str,
2939 byte_limit: u64,
2940) -> Result<Vec<SessionPreviewCandidate>> {
2941 let mut file = File::open(path)?;
2942 let file_len = file.metadata()?.len();
2943 let start = file_len.saturating_sub(byte_limit);
2944 file.seek(SeekFrom::Start(start))?;
2945 let mut bytes = Vec::with_capacity((file_len - start) as usize);
2946 file.read_to_end(&mut bytes)?;
2947 if start > 0 {
2948 if let Some(newline) = bytes.iter().position(|byte| *byte == b'\n') {
2949 bytes.drain(..=newline);
2950 } else {
2951 return Ok(Vec::new());
2952 }
2953 }
2954 let text = String::from_utf8(bytes).map_err(|_| {
2955 Error::Other(format!(
2956 "{} contains non-UTF-8 data in its list-preview window",
2957 path.display()
2958 ))
2959 })?;
2960 if harness == HarnessId::CODEX {
2961 return Ok(codex_preview_candidates(text.lines().rev(), true));
2962 }
2963 let mut candidates = Vec::new();
2964 for line in text.lines().rev() {
2965 let Ok(value) = serde_json::from_str::<Value>(line) else {
2966 continue;
2967 };
2968 let (role, content, metadata) = match harness {
2969 HarnessId::CLAUDE_CODE => {
2970 let role = value.get("type").and_then(Value::as_str);
2971 if !matches!(role, Some("user" | "assistant")) {
2972 continue;
2973 }
2974 let metadata = if role == Some("user") {
2975 crate::session::claude_user_provenance(&value)
2976 .into_iter()
2977 .collect()
2978 } else {
2979 HashMap::new()
2980 };
2981 (
2982 role.unwrap_or_default(),
2983 value
2984 .get("message")
2985 .and_then(|message| message.get("content")),
2986 metadata,
2987 )
2988 }
2989 HarnessId::PI => {
2990 if value.get("type").and_then(Value::as_str) != Some("message") {
2991 continue;
2992 }
2993 let message = value.get("message").unwrap_or(&Value::Null);
2994 let Some(role @ ("user" | "assistant")) =
2995 message.get("role").and_then(Value::as_str)
2996 else {
2997 continue;
2998 };
2999 (role, message.get("content"), HashMap::new())
3000 }
3001 HarnessId::GEMINI => {
3002 let Some(kind @ ("user" | "gemini")) = value.get("type").and_then(Value::as_str)
3003 else {
3004 continue;
3005 };
3006 (
3007 if kind == "gemini" {
3008 "assistant"
3009 } else {
3010 "user"
3011 },
3012 value.get("content"),
3013 HashMap::new(),
3014 )
3015 }
3016 HarnessId::GROK => {
3017 let Some(role @ ("user" | "assistant")) = value.get("type").and_then(Value::as_str)
3018 else {
3019 continue;
3020 };
3021 (role, value.get("content"), HashMap::new())
3022 }
3023 HarnessId::SUPERCODE => {
3024 let Some(role @ ("user" | "assistant")) = value.get("role").and_then(Value::as_str)
3025 else {
3026 continue;
3027 };
3028 (role, value.get("content"), HashMap::new())
3029 }
3030 _ => continue,
3031 };
3032 let mut metadata = metadata;
3033 if matches!(harness, HarnessId::CLAUDE_CODE | HarnessId::CODEX) {
3034 if let Some(timestamp) = value.get("timestamp").and_then(Value::as_str) {
3035 metadata.insert("timestamp".to_string(), timestamp.to_string());
3036 }
3037 }
3038 push_message_candidate_with_cursor(
3039 &mut candidates,
3040 role,
3041 content,
3042 metadata,
3043 Some(message_candidate_cursor(harness, &value)),
3044 );
3045 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3046 break;
3047 }
3048 }
3049 Ok(candidates)
3050}
3051
3052struct CodexPreviewRecord {
3053 native: Value,
3054 role: String,
3055 text: String,
3056}
3057
3058fn codex_preview_candidates<'a>(
3064 lines: impl Iterator<Item = &'a str>,
3065 latest: bool,
3066) -> Vec<SessionPreviewCandidate> {
3067 let mut candidates = Vec::new();
3068 let mut pending: Option<CodexPreviewRecord> = None;
3069 for line in lines {
3070 let Ok(native) = serde_json::from_str::<Value>(line) else {
3071 continue;
3072 };
3073 let Some((role, content)) = codex_preview_message(&native) else {
3074 continue;
3075 };
3076 let Some(text) = display_text(content).filter(|text| !text.trim().is_empty()) else {
3077 continue;
3078 };
3079 let current = CodexPreviewRecord {
3080 role: role.to_string(),
3081 text,
3082 native,
3083 };
3084 if let Some(previous) = pending.take() {
3085 if previous.role == current.role
3086 && previous.text == current.text
3087 && previous.native.get("type") != current.native.get("type")
3088 {
3089 let canonical = if previous.native.get("type").and_then(Value::as_str)
3090 == Some("response_item")
3091 {
3092 previous
3093 } else {
3094 current
3095 };
3096 push_codex_preview_candidate(&mut candidates, canonical, latest);
3097 } else {
3098 push_codex_preview_candidate(&mut candidates, previous, latest);
3099 pending = Some(current);
3100 }
3101 } else {
3102 pending = Some(current);
3103 }
3104 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3105 break;
3106 }
3107 }
3108 if let Some(last) = pending {
3109 push_codex_preview_candidate(&mut candidates, last, latest);
3110 }
3111 candidates
3112}
3113
3114fn push_codex_preview_candidate(
3115 candidates: &mut Vec<SessionPreviewCandidate>,
3116 record: CodexPreviewRecord,
3117 latest: bool,
3118) {
3119 let mut metadata = HashMap::new();
3120 if latest {
3121 if let Some(timestamp) = record.native.get("timestamp").and_then(Value::as_str) {
3122 metadata.insert("timestamp".to_string(), timestamp.to_string());
3123 }
3124 }
3125 let cursor = latest.then(|| message_candidate_cursor(HarnessId::CODEX, &record.native));
3126 push_message_candidate_with_cursor(
3127 candidates,
3128 &record.role,
3129 Some(&Value::String(record.text)),
3130 metadata,
3131 cursor,
3132 );
3133}
3134
3135fn codex_preview_message(value: &Value) -> Option<(&str, Option<&Value>)> {
3137 let payload = value.get("payload")?;
3138 match (
3139 value.get("type").and_then(Value::as_str)?,
3140 payload.get("type").and_then(Value::as_str)?,
3141 ) {
3142 ("response_item", "message") => {
3143 let role @ ("user" | "assistant") = payload.get("role").and_then(Value::as_str)? else {
3144 return None;
3145 };
3146 Some((role, payload.get("content")))
3147 }
3148 ("event_msg", "user_message") => Some(("user", payload.get("message"))),
3149 ("event_msg", "agent_message") => Some(("assistant", payload.get("message"))),
3150 _ => None,
3151 }
3152}
3153
3154fn push_topic_message_candidate(
3155 candidates: &mut Vec<SessionPreviewCandidate>,
3156 harness: &str,
3157 value: &Value,
3158) {
3159 let (role, content, metadata) = match harness {
3160 HarnessId::CLAUDE_CODE => {
3161 let role = value.get("type").and_then(Value::as_str);
3162 if !matches!(role, Some("user" | "assistant")) {
3163 return;
3164 }
3165 let metadata = if role == Some("user") {
3166 crate::session::claude_user_provenance(value)
3167 .into_iter()
3168 .collect()
3169 } else {
3170 HashMap::new()
3171 };
3172 (
3173 role.unwrap_or_default(),
3174 value
3175 .get("message")
3176 .and_then(|message| message.get("content")),
3177 metadata,
3178 )
3179 }
3180 _ => return,
3181 };
3182 push_message_candidate(candidates, role, content, metadata);
3183}
3184
3185fn latest_opencode_message_candidates(
3186 path: &Path,
3187 session_id: &str,
3188) -> Result<Vec<SessionPreviewCandidate>> {
3189 let connection = Connection::open_with_flags(
3190 path,
3191 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
3192 )
3193 .map_err(|error| Error::Other(format!("failed to open list-preview store: {error}")))?;
3194 let mut statement = connection
3195 .prepare(
3196 "SELECT m.data, p.data FROM message m JOIN part p ON p.message_id = m.id \
3197 WHERE m.session_id = ?1 ORDER BY m.time_created DESC, p.time_created DESC LIMIT 32",
3198 )
3199 .map_err(|error| Error::Other(format!("failed to prepare list-preview query: {error}")))?;
3200 let rows = statement
3201 .query_map([session_id], |row| {
3202 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
3203 })
3204 .map_err(|error| Error::Other(format!("failed to read list-preview rows: {error}")))?;
3205 let mut candidates = Vec::new();
3206 for row in rows.flatten() {
3207 let (Ok(message), Ok(part)) = (
3208 serde_json::from_str::<Value>(&row.0),
3209 serde_json::from_str::<Value>(&row.1),
3210 ) else {
3211 continue;
3212 };
3213 let Some(role @ ("user" | "assistant")) = message.get("role").and_then(Value::as_str)
3214 else {
3215 continue;
3216 };
3217 if part.get("type").and_then(Value::as_str) != Some("text") {
3218 continue;
3219 }
3220 push_message_candidate(&mut candidates, role, part.get("text"), HashMap::new());
3221 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3222 break;
3223 }
3224 }
3225 Ok(candidates)
3226}
3227
3228fn latest_hermes_message_candidates(
3229 path: &Path,
3230 session_id: &str,
3231) -> Result<Vec<SessionPreviewCandidate>> {
3232 let connection = Connection::open_with_flags(
3233 path,
3234 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
3235 )
3236 .map_err(|error| Error::Other(format!("failed to open list-preview store: {error}")))?;
3237 let mut statement = connection
3238 .prepare(
3239 "SELECT role, content FROM messages WHERE session_id = ?1 AND active = 1 \
3240 AND role IN ('user', 'assistant') AND content IS NOT NULL AND content != '' \
3241 ORDER BY timestamp DESC, id DESC LIMIT 32",
3242 )
3243 .map_err(|error| Error::Other(format!("failed to prepare list-preview query: {error}")))?;
3244 let rows = statement
3245 .query_map([session_id], |row| {
3246 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
3247 })
3248 .map_err(|error| Error::Other(format!("failed to read list-preview rows: {error}")))?;
3249 let mut candidates = Vec::new();
3250 for (role, content) in rows.flatten() {
3251 push_message_candidate(
3252 &mut candidates,
3253 &role,
3254 Some(&Value::String(content)),
3255 HashMap::new(),
3256 );
3257 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3258 break;
3259 }
3260 }
3261 Ok(candidates)
3262}
3263
3264fn latest_goose_message_candidates(
3265 path: &Path,
3266 session_id: &str,
3267) -> Result<Vec<SessionPreviewCandidate>> {
3268 let connection = Connection::open_with_flags(
3269 path,
3270 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
3271 )
3272 .map_err(|error| Error::Other(format!("failed to open list-preview store: {error}")))?;
3273 let mut statement = connection
3274 .prepare(
3275 "SELECT role, content_json FROM messages WHERE session_id = ?1 \
3276 ORDER BY created_timestamp DESC, id DESC LIMIT 16",
3277 )
3278 .map_err(|error| Error::Other(format!("failed to prepare list-preview query: {error}")))?;
3279 let rows = statement
3280 .query_map([session_id], |row| {
3281 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
3282 })
3283 .map_err(|error| Error::Other(format!("failed to read list-preview rows: {error}")))?;
3284 let mut candidates = Vec::new();
3285 for row in rows.flatten() {
3286 let (role, content) = row;
3287 if !matches!(role.as_str(), "user" | "assistant") {
3288 continue;
3289 }
3290 let Ok(content) = serde_json::from_str::<Value>(&content) else {
3291 continue;
3292 };
3293 push_message_candidate(&mut candidates, &role, Some(&content), HashMap::new());
3294 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3295 break;
3296 }
3297 }
3298 Ok(candidates)
3299}
3300
3301fn push_message_candidate(
3302 candidates: &mut Vec<SessionPreviewCandidate>,
3303 role: &str,
3304 content: Option<&Value>,
3305 metadata: HashMap<String, String>,
3306) {
3307 push_message_candidate_with_cursor(candidates, role, content, metadata, None);
3308}
3309
3310fn push_message_candidate_with_cursor(
3311 candidates: &mut Vec<SessionPreviewCandidate>,
3312 role: &str,
3313 content: Option<&Value>,
3314 metadata: HashMap<String, String>,
3315 cursor: Option<String>,
3316) {
3317 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3318 return;
3319 }
3320 let Some(text) = display_text(content).filter(|text| !text.trim().is_empty()) else {
3321 return;
3322 };
3323 const MAX_CHARS: usize = 4_096;
3324 candidates.push(SessionPreviewCandidate {
3325 cursor,
3326 role: role.to_string(),
3327 content: text.chars().take(MAX_CHARS).collect(),
3328 metadata,
3329 });
3330}
3331
3332fn message_candidate_cursor(harness: &str, value: &Value) -> String {
3333 let native_identity = value
3334 .get("uuid")
3335 .or_else(|| value.get("id"))
3336 .or_else(|| value.pointer("/message/id"))
3337 .or_else(|| value.pointer("/payload/id"))
3338 .and_then(Value::as_str)
3339 .or_else(|| value.get("timestamp").and_then(Value::as_str));
3340 let mut hasher = blake3::Hasher::new();
3341 hasher.update(b"supercode.session-preview-cursor.v1\0");
3342 hasher.update(harness.as_bytes());
3343 hasher.update(b"\0");
3344 if let Some(identity) = native_identity {
3345 hasher.update(identity.as_bytes());
3346 } else {
3347 hasher.update(value.to_string().as_bytes());
3351 }
3352 format!("v1:{}", &hasher.finalize().to_hex()[..24])
3353}
3354
3355fn fill_string(target: &mut Option<String>, value: Option<&Value>) {
3356 if target.is_none() {
3357 *target = value.and_then(Value::as_str).map(str::to_owned);
3358 }
3359}
3360
3361fn fill_path(target: &mut Option<PathBuf>, value: Option<&Value>) {
3362 if target.is_none() {
3363 *target = value.and_then(Value::as_str).map(PathBuf::from);
3364 }
3365}
3366
3367#[derive(Default)]
3370struct TailFacts {
3371 last_turn_ms: Option<u64>,
3373 model: Option<String>,
3375}
3376
3377fn tail_facts(path: &Path, harness: &str) -> TailFacts {
3393 let mut facts = TailFacts::default();
3394 let Ok(mut file) = File::open(path) else {
3395 return facts;
3396 };
3397 let Ok(len) = file.metadata().map(|meta| meta.len()) else {
3398 return facts;
3399 };
3400 let mut window = TAIL_SCAN_START.min(len);
3401 loop {
3402 if file.seek(SeekFrom::Start(len - window)).is_err() {
3403 return facts;
3404 }
3405 let Ok(size) = usize::try_from(window) else {
3406 return facts;
3407 };
3408 let mut buf = vec![0u8; size];
3409 if file.read_exact(&mut buf).is_err() {
3410 return facts;
3411 }
3412 let floor = if window < len {
3423 buf.iter().position(|byte| *byte == b'\n').map(|at| at + 1)
3424 } else {
3425 Some(0)
3426 };
3427 if let Some(floor) = floor {
3428 let mut end = buf.len();
3429 while end > floor && !(facts.last_turn_ms.is_some() && facts.model.is_some()) {
3430 let start = buf[floor..end]
3431 .iter()
3432 .rposition(|byte| *byte == b'\n')
3433 .map_or(floor, |at| floor + at + 1);
3434 if let Ok(record) = std::str::from_utf8(&buf[start..end])
3435 .map_err(|_| ())
3436 .and_then(|line| serde_json::from_str::<Value>(line).map_err(|_| ()))
3437 {
3438 if facts.last_turn_ms.is_none() {
3439 facts.last_turn_ms = record_timestamp(&record, harness)
3440 .and_then(crate::sidecar::rfc3339_to_ms)
3441 .and_then(|millis| u64::try_from(millis).ok());
3442 }
3443 if facts.model.is_none() {
3444 facts.model = record_model(&record, harness).map(str::to_owned);
3445 }
3446 }
3447 end = start.saturating_sub(1);
3448 }
3449 }
3450 if facts.last_turn_ms.is_some() || window >= len || window >= TAIL_SCAN_LIMIT {
3451 return facts;
3452 }
3453 window = (window * 2).min(len).min(TAIL_SCAN_LIMIT);
3454 }
3455}
3456
3457fn record_timestamp<'a>(record: &'a Value, harness: &str) -> Option<&'a str> {
3461 let key = match harness {
3462 HarnessId::SUPERCODE => "ts",
3463 _ => "timestamp",
3464 };
3465 record.get(key)?.as_str()
3466}
3467
3468fn record_model<'a>(record: &'a Value, harness: &str) -> Option<&'a str> {
3471 match harness {
3472 HarnessId::CLAUDE_CODE | HarnessId::PI => record.get("message")?.get("model")?.as_str(),
3473 HarnessId::CODEX => {
3474 if record.get("type")?.as_str()? != "turn_context" {
3475 return None;
3476 }
3477 record.get("payload")?.get("model")?.as_str()
3478 }
3479 _ => None,
3480 }
3481}
3482
3483const TAIL_SCAN_START: u64 = 16 * 1024;
3486
3487const TAIL_SCAN_LIMIT: u64 = 1024 * 1024;
3490
3491fn modified_ms(path: &Path) -> Option<u64> {
3492 fs::metadata(path)
3493 .ok()?
3494 .modified()
3495 .ok()?
3496 .duration_since(UNIX_EPOCH)
3497 .ok()
3498 .and_then(|duration| u64::try_from(duration.as_millis()).ok())
3499}
3500
3501pub struct WorkspaceScope {
3504 wanted: PathBuf,
3505 roots: Vec<PathBuf>,
3507 subtree: bool,
3508 judged: std::sync::Mutex<HashMap<PathBuf, bool>>,
3510}
3511
3512impl WorkspaceScope {
3513 pub fn exact(wanted: &Path) -> Self {
3515 Self {
3516 wanted: wanted.to_path_buf(),
3517 roots: Vec::new(),
3518 subtree: false,
3519 judged: std::sync::Mutex::new(HashMap::new()),
3520 }
3521 }
3522
3523 pub fn subtree(root: &Path) -> Self {
3530 let mut roots = path_comparison_keys(root);
3531 let own = roots.clone();
3532 for linked in linked_folders(root) {
3533 for key in path_comparison_keys(&linked) {
3534 let widens = own.iter().any(|root| root.starts_with(&key));
3535 let covered = roots.iter().any(|root| key.starts_with(root));
3536 if !widens && !covered {
3537 roots.push(key);
3538 }
3539 }
3540 }
3541 Self {
3542 wanted: root.to_path_buf(),
3543 roots,
3544 subtree: true,
3545 judged: std::sync::Mutex::new(HashMap::new()),
3546 }
3547 }
3548
3549 pub fn of(query: &DiscoveryQuery) -> Option<Self> {
3551 query.workspace.as_deref().map(|wanted| {
3552 if query.workspace_subtree {
3553 Self::subtree(wanted)
3554 } else {
3555 Self::exact(wanted)
3556 }
3557 })
3558 }
3559
3560 pub fn admits(&self, recorded: &Path) -> bool {
3565 if !recorded.is_absolute() {
3566 return false;
3567 }
3568 if !self.subtree {
3569 return same_path(recorded, &self.wanted);
3570 }
3571 let mut judged = self
3572 .judged
3573 .lock()
3574 .unwrap_or_else(std::sync::PoisonError::into_inner);
3575 *judged.entry(recorded.to_path_buf()).or_insert_with(|| {
3576 path_comparison_keys(recorded)
3577 .iter()
3578 .any(|key| self.roots.iter().any(|root| key.starts_with(root)))
3579 })
3580 }
3581
3582 fn admits_literal(&self, recorded: Option<&Path>) -> bool {
3586 if self.subtree {
3587 recorded.is_some_and(|cwd| self.admits(cwd))
3588 } else {
3589 recorded == Some(self.wanted.as_path())
3590 }
3591 }
3592}
3593
3594fn linked_folders(root: &Path) -> Vec<PathBuf> {
3596 let Ok(entries) = fs::read_dir(root) else {
3597 return Vec::new();
3598 };
3599 entries
3600 .filter_map(|entry| {
3601 let entry = entry.ok()?;
3602 let link = entry.file_type().ok()?.is_symlink();
3603 (link && fs::metadata(entry.path()).is_ok_and(|meta| meta.is_dir()))
3604 .then(|| entry.path())
3605 })
3606 .collect()
3607}
3608
3609fn path_comparison_keys(path: &Path) -> Vec<PathBuf> {
3615 let mut keys = Vec::with_capacity(2);
3616 if let Ok(canonical) = fs::canonicalize(path) {
3617 keys.push(comparison_key(&canonical));
3618 }
3619 let lexical = comparison_key(&normalize_path(path));
3620 if !keys.contains(&lexical) {
3621 keys.push(lexical);
3622 }
3623 keys
3624}
3625
3626#[cfg(windows)]
3627fn comparison_key(path: &Path) -> PathBuf {
3628 let text = path.to_string_lossy();
3629 let text = if let Some(unc) = text.strip_prefix(r"\\?\UNC\") {
3630 format!(r"\\{unc}")
3631 } else if let Some(local) = text.strip_prefix(r"\\?\") {
3632 local.to_string()
3633 } else {
3634 text.into_owned()
3635 };
3636 PathBuf::from(text.to_lowercase())
3637}
3638
3639#[cfg(not(windows))]
3640fn comparison_key(path: &Path) -> PathBuf {
3641 path.to_path_buf()
3642}
3643
3644fn recorded_cwd_matches(recorded: &Path, wanted: &Path) -> bool {
3650 recorded.is_absolute() && same_path(recorded, wanted)
3651}
3652
3653fn same_path(left: &Path, right: &Path) -> bool {
3654 match (fs::canonicalize(left), fs::canonicalize(right)) {
3655 (Ok(left), Ok(right)) => left == right,
3656 _ => normalize_path(left) == normalize_path(right),
3657 }
3658}
3659
3660fn normalize_path(path: &Path) -> PathBuf {
3661 let absolute = if path.is_absolute() {
3662 path.to_path_buf()
3663 } else {
3664 std::env::current_dir()
3665 .unwrap_or_else(|_| PathBuf::from("."))
3666 .join(path)
3667 };
3668 let mut normalized = PathBuf::new();
3669 for component in absolute.components() {
3670 match component {
3671 Component::CurDir => {}
3672 Component::ParentDir => {
3673 normalized.pop();
3674 }
3675 other => normalized.push(other.as_os_str()),
3676 }
3677 }
3678 normalized
3679}
3680
3681fn claude_custom_titles(path: &std::path::Path) -> Vec<String> {
3686 use std::io::{Read, Seek, SeekFrom};
3687 use std::sync::{Mutex, OnceLock};
3688 static CACHE: OnceLock<
3689 Mutex<std::collections::HashMap<std::path::PathBuf, (u64, Vec<String>)>>,
3690 > = OnceLock::new();
3691 const KEY: &[u8] = b"\"customTitle\":";
3692 let Ok(len) = std::fs::metadata(path).map(|meta| meta.len()) else {
3693 return Vec::new();
3694 };
3695 let cache = CACHE.get_or_init(Default::default);
3696 let known = cache.lock().ok().and_then(|cache| cache.get(path).cloned());
3697 let (from, mut titles) = match known {
3698 Some((seen, titles)) if seen == len => return titles,
3699 Some((seen, titles)) if seen < len => (seen.saturating_sub(4096), titles),
3701 _ => (0, Vec::new()),
3702 };
3703 let mut bytes = Vec::new();
3704 let read = std::fs::File::open(path).and_then(|mut file| {
3705 file.seek(SeekFrom::Start(from))?;
3706 file.read_to_end(&mut bytes)
3707 });
3708 if read.is_err() {
3709 return titles;
3710 }
3711 let mut at = 0;
3712 while let Some(offset) = bytes[at..]
3713 .windows(KEY.len())
3714 .position(|window| window == KEY)
3715 {
3716 let start = at + offset + KEY.len();
3717 if let Some(Ok(title)) = serde_json::Deserializer::from_slice(&bytes[start..])
3718 .into_iter::<String>()
3719 .next()
3720 {
3721 if !titles.contains(&title) {
3722 titles.push(title);
3723 }
3724 }
3725 at = start;
3726 }
3727 if let Ok(mut cache) = cache.lock() {
3728 cache.insert(path.to_path_buf(), (len, titles.clone()));
3729 }
3730 titles
3731}