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(mut meta) = read_header(path, locator.harness.as_str()) else {
1193 return Ok(None);
1197 };
1198 untitled_by_topic(&mut meta, path, locator.harness.as_str());
1199 if workspace.is_some_and(|wanted| meta.cwd.as_deref().is_none_or(|cwd| !wanted.admits(cwd)))
1200 {
1201 return Ok(None);
1202 }
1203 let parent_session_id = meta.parent_session_id.or_else(|| {
1204 (locator.harness.as_str() == HarnessId::CLAUDE_CODE)
1205 .then(|| claude_subagent_parent_id(path))
1206 .flatten()
1207 });
1208 let tail = tail_facts(path, locator.harness.as_str());
1209 let descriptor = SessionDescriptor {
1210 locator: SessionLocator {
1211 harness: locator.harness.clone(),
1212 session_id: meta
1213 .session_id
1214 .unwrap_or_else(|| locator.session_id.clone()),
1215 storage: StorageLocator::File { path: path.clone() },
1216 },
1217 cwd: meta.cwd,
1218 title: meta.title,
1219 preview_candidates: Vec::new(),
1220 latest_message_candidates: Vec::new(),
1221 updated_at_ms: tail.last_turn_ms.or_else(|| modified_ms(path)),
1222 message_count: None,
1223 model: tail.model.or(meta.model),
1224 parent_session_id,
1225 child_session_count: 0,
1226 nouns: OrchestrationNouns::default(),
1227 };
1228 Ok(Some(descriptor))
1229 }
1230
1231 pub fn load(&self, locator: &SessionLocator) -> Result<Session> {
1233 self.load_with_fidelity(locator, Fidelity::ByteLossless)
1234 }
1235
1236 pub fn load_with_fidelity(
1243 &self,
1244 locator: &SessionLocator,
1245 fidelity: Fidelity,
1246 ) -> Result<Session> {
1247 if let Some(session) = load_hermes_locator(locator) {
1248 return session;
1249 }
1250 match &locator.storage {
1251 StorageLocator::File { path } => {
1252 if let Some(session) = load_native_store_family(path)? {
1253 Ok(session)
1254 } else {
1255 Ok(Session::load_with_fidelity(path, fidelity)?)
1256 }
1257 }
1258 StorageLocator::Sqlite { path, selector } => {
1259 if locator.harness.as_str() == HarnessId::GOOSE {
1260 Ok(Session::from_goose_sqlite(path, selector)?)
1261 } else {
1262 Ok(Session::from_opencode_sqlite(path, Some(selector))?)
1263 }
1264 }
1265 }
1266 }
1267
1268 #[doc(hidden)]
1272 pub fn load_parent_with_fidelity(
1273 &self,
1274 locator: &SessionLocator,
1275 fidelity: Fidelity,
1276 ) -> Result<Session> {
1277 if let Some(session) = load_hermes_locator(locator) {
1278 return session;
1279 }
1280 match &locator.storage {
1281 StorageLocator::File { path } => {
1282 if let Some(session) = load_native_store_family(path)? {
1283 Ok(session)
1284 } else {
1285 Ok(Session::load_parent_with_fidelity(path, fidelity)?)
1286 }
1287 }
1288 StorageLocator::Sqlite { path, selector } => {
1289 if locator.harness.as_str() == HarnessId::GOOSE {
1290 Ok(Session::from_goose_sqlite(path, selector)?)
1291 } else {
1292 Ok(Session::from_opencode_sqlite(path, Some(selector))?)
1293 }
1294 }
1295 }
1296 }
1297
1298 #[doc(hidden)]
1301 pub fn load_display_view(
1302 &self,
1303 locator: &SessionLocator,
1304 fidelity: Fidelity,
1305 message_limit: usize,
1306 ) -> Result<Session> {
1307 if let Some(session) = load_hermes_locator(locator) {
1308 let mut session = session?;
1309 if session.messages.len() > message_limit.max(1) {
1310 session
1311 .messages
1312 .drain(..session.messages.len() - message_limit.max(1));
1313 }
1314 return Ok(session);
1315 }
1316 match &locator.storage {
1317 StorageLocator::File { path } => {
1318 if let Some(mut session) = load_native_store_family(path)? {
1319 if session.messages.len() > message_limit.max(1) {
1320 session
1321 .messages
1322 .drain(..session.messages.len() - message_limit.max(1));
1323 }
1324 Ok(session)
1325 } else {
1326 Ok(Session::load_display_view(path, fidelity, message_limit)?)
1327 }
1328 }
1329 StorageLocator::Sqlite { path, selector } => {
1330 let mut session = if locator.harness.as_str() == HarnessId::GOOSE {
1331 Session::from_goose_sqlite_display(path, selector, message_limit)?
1332 } else {
1333 Session::from_opencode_sqlite(path, Some(selector))?
1334 };
1335 if session.messages.len() > message_limit.max(1) {
1336 session
1337 .messages
1338 .drain(..session.messages.len() - message_limit.max(1));
1339 }
1340 Ok(session)
1341 }
1342 }
1343 }
1344
1345 pub fn follow(&self, locator: &SessionLocator) -> Result<SessionFollower> {
1347 self.follow_with_fidelity(locator, Fidelity::ByteLossless)
1348 }
1349
1350 pub fn follow_with_fidelity(
1352 &self,
1353 locator: &SessionLocator,
1354 fidelity: Fidelity,
1355 ) -> Result<SessionFollower> {
1356 SessionFollower::open_locator_with_fidelity(locator, fidelity)
1357 }
1358
1359 #[doc(hidden)]
1361 pub fn follow_read_view(
1362 &self,
1363 locator: &SessionLocator,
1364 fidelity: Fidelity,
1365 include_subagents: bool,
1366 message_limit: Option<usize>,
1367 max_message_chars: Option<usize>,
1368 display_history: bool,
1369 ) -> Result<SessionFollower> {
1370 SessionFollower::open_locator_with_view(
1371 locator,
1372 fidelity,
1373 include_subagents,
1374 message_limit,
1375 max_message_chars,
1376 display_history,
1377 )
1378 }
1379}
1380
1381fn load_hermes_locator(locator: &SessionLocator) -> Option<Result<Session>> {
1387 if locator.harness.as_str() != HarnessId::HERMES {
1388 return None;
1389 }
1390 let StorageLocator::File { path } = &locator.storage else {
1391 return None;
1392 };
1393 Some(Session::from_hermes_sqlite(path, Some(&locator.session_id)))
1394}
1395
1396fn finalize_nouns(descriptor: &mut SessionDescriptor) {
1403 let mut meta = SessionMeta::new(SessionSource::Native);
1404 meta.cwd = descriptor.cwd.clone();
1405 meta.trigger = descriptor.nouns.trigger;
1406 meta.surface = descriptor.nouns.surface.clone();
1407 meta.profile = descriptor.nouns.profile.clone();
1408 meta.recurrence = descriptor.nouns.recurrence.clone();
1409 meta.cross_surface = descriptor.nouns.cross_surface.clone();
1410 descriptor.nouns = OrchestrationNouns::from_meta(&meta);
1411}
1412
1413fn project_descriptors(query: &DiscoveryQuery, found: &mut Vec<SessionDescriptor>) {
1414 roll_up_session_children(found, query.include_child_sessions);
1415 if let Some(root_session_id) = query.root_session_id.as_deref() {
1416 retain_session_family(found, root_session_id);
1417 }
1418 if let Some(family_path) = query.workspace_family.as_deref() {
1419 let family = RepoFamily::of(family_path);
1420 let mut cache: HashMap<PathBuf, bool> = HashMap::new();
1421 found.retain(|descriptor| {
1422 let Some(cwd) = descriptor.cwd.as_deref() else {
1423 return false;
1424 };
1425 *cache
1426 .entry(cwd.to_path_buf())
1427 .or_insert_with(|| RepoFamily::of(cwd).joins(&family))
1428 });
1429 }
1430 if let Some(profile) = query
1431 .profile
1432 .as_deref()
1433 .map(str::trim)
1434 .filter(|p| !p.is_empty())
1435 {
1436 found.retain(|descriptor| descriptor.nouns.profile.as_deref() == Some(profile));
1437 }
1438 if let Some(after) = query.updated_after_ms {
1439 found.retain(|descriptor| descriptor.updated_at_ms.is_some_and(|at| at >= after));
1440 }
1441 if let Some(before) = query.updated_before_ms {
1442 found.retain(|descriptor| descriptor.updated_at_ms.is_some_and(|at| at <= before));
1443 }
1444 found.sort_by(|a, b| {
1445 b.updated_at_ms
1446 .cmp(&a.updated_at_ms)
1447 .then_with(|| a.locator.harness.cmp(&b.locator.harness))
1448 .then_with(|| a.locator.session_id.cmp(&b.locator.session_id))
1449 });
1450 if let Some(search) = query
1451 .query
1452 .as_deref()
1453 .map(str::trim)
1454 .filter(|query| !query.is_empty())
1455 {
1456 let search = search.to_lowercase();
1457 found.retain(|descriptor| descriptor_matches(descriptor, &search));
1458 }
1459}
1460
1461fn paginate_descriptors(
1462 query: &DiscoveryQuery,
1463 found: Vec<SessionDescriptor>,
1464) -> Result<(Vec<SessionDescriptor>, Option<String>)> {
1465 let start = match query.cursor.as_deref() {
1466 Some(cursor) => {
1467 let key = decode_cursor(cursor)?;
1468 found
1469 .iter()
1470 .position(|descriptor| descriptor_cursor_key(descriptor) == key)
1471 .map(|index| index + 1)
1472 .ok_or_else(|| Error::Other("discovery cursor is stale or invalid".into()))?
1473 }
1474 None => 0,
1475 };
1476 let end = query
1477 .limit
1478 .map(|limit| start.saturating_add(limit).min(found.len()))
1479 .unwrap_or(found.len());
1480 let sessions = found[start.min(found.len())..end].to_vec();
1481 let next_cursor = (end < found.len())
1482 .then(|| sessions.last().map(encode_cursor))
1483 .flatten();
1484 Ok((sessions, next_cursor))
1485}
1486
1487fn enrich_descriptors(
1488 query: &DiscoveryQuery,
1489 sessions: &mut [SessionDescriptor],
1490 codex_history: Option<&CodexHistoryTopicIndex>,
1491) -> Result<()> {
1492 let codex_topics = if query.include_topic_candidates && codex_history.is_none() {
1493 codex_history_topics(&query.homes.codex, sessions).unwrap_or_default()
1494 } else {
1495 HashMap::new()
1496 };
1497 for descriptor in sessions {
1498 let topic = (descriptor.locator.harness.as_str() == HarnessId::CODEX)
1499 .then(|| {
1500 codex_history
1501 .and_then(|history| history.topics.get(&descriptor.locator.session_id))
1502 .or_else(|| codex_topics.get(&descriptor.locator.session_id))
1503 })
1504 .flatten();
1505 enrich_descriptor(query, descriptor, topic);
1506 }
1507 Ok(())
1508}
1509
1510fn enrich_descriptor(
1511 query: &DiscoveryQuery,
1512 descriptor: &mut SessionDescriptor,
1513 codex_topic: Option<&Vec<SessionPreviewCandidate>>,
1514) {
1515 if query.include_topic_candidates {
1516 descriptor.preview_candidates = codex_topic
1517 .cloned()
1518 .unwrap_or_else(|| topic_message_candidates(&descriptor.locator).unwrap_or_default());
1519 }
1520 descriptor.latest_message_candidates =
1521 latest_message_candidates(&descriptor.locator).unwrap_or_default();
1522}
1523
1524fn is_false(value: &bool) -> bool {
1525 !value
1526}
1527
1528fn descriptor_matches(descriptor: &SessionDescriptor, search: &str) -> bool {
1529 [
1530 Some(descriptor.locator.harness.as_str()),
1531 Some(descriptor.locator.session_id.as_str()),
1532 descriptor.title.as_deref(),
1533 descriptor.cwd.as_ref().and_then(|path| path.to_str()),
1534 descriptor.model.as_deref(),
1535 ]
1536 .into_iter()
1537 .flatten()
1538 .any(|value| value.to_lowercase().contains(search))
1539}
1540
1541fn descriptor_cursor_key(descriptor: &SessionDescriptor) -> (Option<u64>, String, String) {
1542 (
1543 descriptor.updated_at_ms,
1544 descriptor.locator.harness.as_str().to_string(),
1545 descriptor.locator.session_id.clone(),
1546 )
1547}
1548
1549fn encode_cursor(descriptor: &SessionDescriptor) -> String {
1550 let json = serde_json::to_vec(&descriptor_cursor_key(descriptor)).unwrap_or_default();
1551 let mut encoded = String::with_capacity(json.len() * 2);
1552 for byte in json {
1553 use std::fmt::Write;
1554 let _ = write!(&mut encoded, "{byte:02x}");
1555 }
1556 encoded
1557}
1558
1559fn decode_cursor(cursor: &str) -> Result<(Option<u64>, String, String)> {
1560 if cursor.len() % 2 != 0 {
1561 return Err(Error::Other("discovery cursor is invalid".into()));
1562 }
1563 let bytes = (0..cursor.len())
1564 .step_by(2)
1565 .map(|index| u8::from_str_radix(&cursor[index..index + 2], 16))
1566 .collect::<std::result::Result<Vec<_>, _>>()
1567 .map_err(|_| Error::Other("discovery cursor is invalid".into()))?;
1568 serde_json::from_slice(&bytes).map_err(|_| Error::Other("discovery cursor is invalid".into()))
1569}
1570
1571#[derive(Default)]
1572struct HeaderMeta {
1573 session_id: Option<String>,
1574 cwd: Option<PathBuf>,
1575 title: Option<String>,
1576 model: Option<String>,
1577 parent_session_id: Option<String>,
1578 trigger: Option<Trigger>,
1581 topic: Option<String>,
1584}
1585
1586fn discover_jsonl(
1587 root: &Path,
1588 harness: &str,
1589 workspace: Option<&WorkspaceScope>,
1590 include_child_sessions: bool,
1591 found: &mut Vec<SessionDescriptor>,
1592) {
1593 let mut files = Vec::new();
1594 collect_jsonl(root, harness, include_child_sessions, &mut files);
1595 for path in files {
1596 let Ok(mut meta) = read_header(&path, harness) else {
1597 continue;
1598 };
1599 untitled_by_topic(&mut meta, &path, harness);
1600 if workspace.is_some_and(|wanted| meta.cwd.as_deref().is_none_or(|cwd| !wanted.admits(cwd)))
1601 {
1602 continue;
1603 }
1604 let session_id = meta.session_id.unwrap_or_else(|| {
1605 path.file_stem()
1606 .and_then(|value| value.to_str())
1607 .unwrap_or("unknown")
1608 .to_string()
1609 });
1610 let parent_session_id = meta.parent_session_id.or_else(|| {
1611 (harness == HarnessId::CLAUDE_CODE)
1612 .then(|| claude_subagent_parent_id(&path))
1613 .flatten()
1614 });
1615 let tail = tail_facts(&path, harness);
1616 found.push(SessionDescriptor {
1617 locator: SessionLocator {
1618 harness: HarnessId::new(harness),
1619 session_id,
1620 storage: StorageLocator::File { path: path.clone() },
1621 },
1622 cwd: meta.cwd,
1623 title: meta.title,
1624 preview_candidates: Vec::new(),
1625 latest_message_candidates: Vec::new(),
1626 updated_at_ms: tail.last_turn_ms.or_else(|| modified_ms(&path)),
1627 message_count: None,
1628 model: tail.model.or(meta.model),
1629 parent_session_id,
1630 child_session_count: if harness == HarnessId::CLAUDE_CODE && !include_child_sessions {
1631 count_claude_subagents(&path)
1632 } else {
1633 0
1634 },
1635 nouns: OrchestrationNouns {
1636 trigger: meta.trigger,
1637 ..OrchestrationNouns::default()
1638 },
1639 });
1640 }
1641}
1642
1643fn collect_jsonl(root: &Path, harness: &str, include_child_sessions: bool, out: &mut Vec<PathBuf>) {
1644 let mut walked = HashSet::new();
1645 collect_jsonl_in(root, harness, include_child_sessions, out, &mut walked);
1646}
1647
1648fn collect_jsonl_in(
1652 root: &Path,
1653 harness: &str,
1654 include_child_sessions: bool,
1655 out: &mut Vec<PathBuf>,
1656 walked: &mut HashSet<PathBuf>,
1657) {
1658 if !walked.insert(fs::canonicalize(root).unwrap_or_else(|_| root.to_path_buf())) {
1659 return;
1660 }
1661 let Ok(entries) = fs::read_dir(root) else {
1662 return;
1663 };
1664 for entry in entries.flatten() {
1665 let Ok(mut kind) = entry.file_type() else {
1666 continue;
1667 };
1668 let path = entry.path();
1669 if kind.is_symlink() {
1670 let Ok(target) = fs::metadata(&path) else {
1671 continue;
1672 };
1673 kind = target.file_type();
1674 }
1675 if kind.is_dir() {
1676 if harness == HarnessId::CLAUDE_CODE
1677 && path.file_name().and_then(|v| v.to_str()) == Some("subagents")
1678 && !include_child_sessions
1679 {
1680 continue;
1681 }
1682 collect_jsonl_in(&path, harness, include_child_sessions, out, walked);
1683 } else if kind.is_file()
1684 && (path.extension().and_then(|v| v.to_str()) == Some("jsonl")
1685 || (harness == HarnessId::CODEX
1686 && path
1687 .file_name()
1688 .and_then(|v| v.to_str())
1689 .is_some_and(|name| name.ends_with(".jsonl.zst"))))
1690 {
1691 out.push(path);
1692 }
1693 }
1694}
1695
1696fn read_header(path: &Path, harness: &str) -> Result<HeaderMeta> {
1697 let file = crate::session::open_session_reader(path)?;
1698 let mut result = HeaderMeta::default();
1699 let mut bytes = 0usize;
1700 for line in BufReader::new(file).lines().take(32) {
1701 let line = line?;
1702 bytes += line.len();
1703 if bytes > 256 * 1024 {
1704 break;
1705 }
1706 let Ok(value) = serde_json::from_str::<Value>(&line) else {
1707 continue;
1708 };
1709 update_header_meta(&mut result, &value, harness);
1710 if result.session_id.is_some()
1712 && result.cwd.is_some()
1713 && result.model.is_some()
1714 && (harness != HarnessId::CLAUDE_CODE
1715 || result.title.is_some()
1716 || result.topic.is_some())
1717 {
1718 break;
1719 }
1720 }
1721 if result.session_id.is_none() && result.cwd.is_none() {
1722 return Err(Error::Other(format!(
1723 "{} has no recognizable {harness} session header",
1724 path.display()
1725 )));
1726 }
1727 Ok(result)
1728}
1729
1730fn untitled_by_topic(meta: &mut HeaderMeta, path: &Path, harness: &str) {
1733 if meta.title.is_some() {
1734 return;
1735 }
1736 meta.title = meta.topic.take().or_else(|| {
1737 (harness == HarnessId::CODEX)
1738 .then(|| {
1739 meta.session_id
1740 .as_deref()
1741 .and_then(|id| codex_thread_name(path, id))
1742 })
1743 .flatten()
1744 });
1745}
1746
1747pub fn native_topic(harness: &str, path: &Path, session_id: &str) -> Option<String> {
1751 match harness {
1752 HarnessId::CODEX => codex_thread_name(path, session_id),
1753 HarnessId::CLAUDE_CODE => read_header(path, harness).ok().and_then(|meta| meta.topic),
1754 _ => None,
1755 }
1756}
1757
1758fn codex_thread_name(rollout: &Path, id: &str) -> Option<String> {
1761 use std::sync::{Mutex, OnceLock};
1762 type Names = (u64, Option<std::time::SystemTime>, HashMap<String, String>);
1763 static CACHE: OnceLock<Mutex<HashMap<PathBuf, Names>>> = OnceLock::new();
1764 let index = rollout
1765 .ancestors()
1766 .find(|dir| dir.file_name().is_some_and(|name| name == "sessions"))?
1767 .parent()?
1768 .join("session_index.jsonl");
1769 let stat = std::fs::metadata(&index).ok()?;
1770 let stamp = (stat.len(), stat.modified().ok());
1771 let mut cache = CACHE.get_or_init(Default::default).lock().ok()?;
1772 let fresh = cache
1773 .get(&index)
1774 .is_some_and(|(len, modified, _)| (*len, *modified) == stamp);
1775 if !fresh {
1776 let mut names = HashMap::new();
1777 if let Ok(file) = File::open(&index) {
1778 for line in BufReader::new(file)
1779 .lines()
1780 .map_while(std::result::Result::ok)
1781 {
1782 let Ok(value) = serde_json::from_str::<Value>(&line) else {
1783 continue;
1784 };
1785 if let (Some(id), Some(name)) = (
1786 value.get("id").and_then(Value::as_str),
1787 value.get("thread_name").and_then(Value::as_str),
1788 ) {
1789 if !name.trim().is_empty() {
1790 names.insert(id.to_string(), name.trim().to_string());
1791 }
1792 }
1793 }
1794 }
1795 cache.insert(index.clone(), (stamp.0, stamp.1, names));
1796 }
1797 cache.get(&index)?.2.get(id).cloned()
1798}
1799
1800fn update_header_meta(result: &mut HeaderMeta, value: &Value, harness: &str) {
1801 match harness {
1802 HarnessId::CLAUDE_CODE => {
1803 fill_string(&mut result.session_id, value.get("sessionId"));
1804 fill_path(&mut result.cwd, value.get("cwd"));
1805 fill_string(&mut result.title, value.get("customTitle"));
1807 fill_string(&mut result.title, value.get("agentName"));
1808 if value.get("type").and_then(Value::as_str) == Some("ai-title") {
1809 fill_string(&mut result.topic, value.get("aiTitle"));
1810 }
1811 fill_string(
1812 &mut result.model,
1813 value.get("message").and_then(|v| v.get("model")),
1814 );
1815 }
1816 HarnessId::CODEX => {
1817 let payload = value.get("payload").unwrap_or(&Value::Null);
1818 if value.get("type").and_then(Value::as_str) == Some("session_meta") {
1819 fill_string(&mut result.session_id, payload.get("id"));
1820 if result.trigger.is_none() {
1821 result.trigger = crate::session::codex_start_trigger(payload);
1822 }
1823 fill_path(&mut result.cwd, payload.get("cwd"));
1824 fill_string(&mut result.title, payload.get("thread_name"));
1825 fill_string(&mut result.title, payload.get("title"));
1826 fill_string(
1827 &mut result.parent_session_id,
1828 payload.get("parent_thread_id"),
1829 );
1830 if let Some(parent) = payload
1831 .pointer("/source/subagent/thread_spawn/parent_thread_id")
1832 .and_then(Value::as_str)
1833 {
1834 result.parent_session_id = Some(parent.to_string());
1835 }
1836 if result.title.is_none() {
1837 result.title = payload
1838 .pointer("/source/subagent/thread_spawn/agent_path")
1839 .and_then(Value::as_str)
1840 .and_then(|path| path.rsplit('/').find(|part| !part.is_empty()))
1841 .map(humanize_topic);
1842 }
1843 }
1844 if value.get("type").and_then(Value::as_str) == Some("turn_context") {
1845 fill_path(&mut result.cwd, payload.get("cwd"));
1846 fill_string(&mut result.model, payload.get("model"));
1847 }
1848 }
1849 HarnessId::PI => {
1850 if value.get("type").and_then(Value::as_str) == Some("session") {
1851 fill_string(&mut result.session_id, value.get("id"));
1852 fill_path(&mut result.cwd, value.get("cwd"));
1853 }
1854 fill_string(
1855 &mut result.model,
1856 value.get("message").and_then(|v| v.get("model")),
1857 );
1858 }
1859 _ => {}
1860 }
1861}
1862
1863fn roll_up_session_children(found: &mut Vec<SessionDescriptor>, include_children: bool) {
1867 let by_id = found
1868 .iter()
1869 .enumerate()
1870 .map(|(index, descriptor)| {
1871 (
1872 (
1873 descriptor.locator.harness.as_str().to_string(),
1874 descriptor.locator.session_id.clone(),
1875 ),
1876 index,
1877 )
1878 })
1879 .collect::<HashMap<_, _>>();
1880 let mut root_updates = HashMap::<usize, u64>::new();
1881 let mut root_child_counts = HashMap::<usize, usize>::new();
1882
1883 for descriptor in found.iter() {
1884 let Some(mut parent_id) = descriptor.parent_session_id.as_deref() else {
1885 continue;
1886 };
1887 let harness = descriptor.locator.harness.as_str();
1888 let mut root = None;
1889 let mut visited = HashSet::new();
1890 while visited.insert(parent_id.to_string()) {
1891 let Some(&parent_index) = by_id.get(&(harness.to_string(), parent_id.to_string()))
1892 else {
1893 break;
1894 };
1895 root = Some(parent_index);
1896 let Some(next_parent) = found[parent_index].parent_session_id.as_deref() else {
1897 break;
1898 };
1899 parent_id = next_parent;
1900 }
1901 if let (Some(root), Some(updated_at_ms)) = (root, descriptor.updated_at_ms) {
1902 root_updates
1903 .entry(root)
1904 .and_modify(|current| *current = (*current).max(updated_at_ms))
1905 .or_insert(updated_at_ms);
1906 }
1907 if let Some(root) = root {
1908 *root_child_counts.entry(root).or_default() += 1;
1909 }
1910 }
1911
1912 for (root, child_updated_at_ms) in root_updates {
1913 found[root].updated_at_ms = Some(
1914 found[root]
1915 .updated_at_ms
1916 .unwrap_or_default()
1917 .max(child_updated_at_ms),
1918 );
1919 }
1920 for (root, child_count) in root_child_counts {
1921 found[root].child_session_count = child_count;
1922 }
1923 if !include_children {
1924 found.retain(|descriptor| descriptor.parent_session_id.is_none());
1925 }
1926}
1927
1928fn retain_session_family(found: &mut Vec<SessionDescriptor>, root_session_id: &str) {
1929 let parent_by_id = found
1930 .iter()
1931 .map(|descriptor| {
1932 (
1933 descriptor.locator.session_id.clone(),
1934 descriptor.parent_session_id.clone(),
1935 )
1936 })
1937 .collect::<HashMap<_, _>>();
1938 found.retain(|descriptor| {
1939 let mut current = descriptor.locator.session_id.clone();
1940 let mut visited = HashSet::new();
1941 while visited.insert(current.clone()) {
1942 if current == root_session_id {
1943 return true;
1944 }
1945 let Some(Some(parent)) = parent_by_id.get(¤t) else {
1946 return false;
1947 };
1948 current = parent.clone();
1949 }
1950 false
1951 });
1952}
1953
1954fn claude_subagent_parent_id(path: &Path) -> Option<String> {
1955 let subagents = path.parent()?;
1956 if subagents.file_name()?.to_str()? != "subagents" {
1957 return None;
1958 }
1959 subagents
1960 .parent()?
1961 .file_name()?
1962 .to_str()
1963 .map(str::to_string)
1964}
1965
1966fn count_claude_subagents(parent_path: &Path) -> usize {
1967 let Some(parent) = parent_path.parent() else {
1968 return 0;
1969 };
1970 let Some(stem) = parent_path.file_stem() else {
1971 return 0;
1972 };
1973 let root = parent.join(stem).join("subagents");
1974 let mut files = Vec::new();
1975 collect_jsonl(&root, HarnessId::CLAUDE_CODE, true, &mut files);
1976 files.len()
1977}
1978
1979fn is_zero(value: &usize) -> bool {
1980 *value == 0
1981}
1982
1983fn humanize_topic(value: &str) -> String {
1984 let text = value.replace(['_', '-'], " ");
1985 let mut characters = text.chars();
1986 match characters.next() {
1987 Some(first) => first.to_uppercase().collect::<String>() + characters.as_str(),
1988 None => text,
1989 }
1990}
1991
1992fn discover_gemini(
1993 root: &Path,
1994 workspace: Option<&WorkspaceScope>,
1995 found: &mut Vec<SessionDescriptor>,
1996) {
1997 let slug_to_cwd = std::fs::read_to_string(root.join("projects.json"))
1998 .ok()
1999 .and_then(|text| serde_json::from_str::<Value>(&text).ok())
2000 .and_then(|value| value.get("projects").and_then(Value::as_object).cloned())
2001 .map(|projects| {
2002 projects
2003 .into_iter()
2004 .filter_map(|(cwd, slug)| Some((slug.as_str()?.to_string(), PathBuf::from(cwd))))
2005 .collect::<HashMap<_, _>>()
2006 })
2007 .unwrap_or_default();
2008 let mut files = Vec::new();
2009 collect_jsonl(&root.join("tmp"), HarnessId::GEMINI, false, &mut files);
2010 let worker_count = std::thread::available_parallelism()
2011 .map(usize::from)
2012 .unwrap_or(4)
2013 .clamp(1, 8)
2014 .min(files.len().max(1));
2015 let chunk_size = files.len().max(1).div_ceil(worker_count);
2016 let discovered = std::thread::scope(|scope| {
2017 files
2018 .chunks(chunk_size)
2019 .map(|paths| {
2020 scope.spawn(|| {
2021 paths
2022 .iter()
2023 .filter_map(|path| gemini_descriptor(path, &slug_to_cwd, workspace))
2024 .collect::<Vec<_>>()
2025 })
2026 })
2027 .collect::<Vec<_>>()
2028 .into_iter()
2029 .flat_map(|worker| {
2030 worker
2031 .join()
2032 .expect("Gemini discovery worker must not panic")
2033 })
2034 .collect::<Vec<_>>()
2035 });
2036 found.extend(discovered);
2037}
2038
2039fn gemini_descriptor(
2040 path: &Path,
2041 slug_to_cwd: &HashMap<String, PathBuf>,
2042 workspace: Option<&WorkspaceScope>,
2043) -> Option<SessionDescriptor> {
2044 if path
2045 .parent()
2046 .and_then(Path::file_name)
2047 .and_then(|name| name.to_str())
2048 != Some("chats")
2049 {
2050 return None;
2051 }
2052 let slug = path
2053 .parent()
2054 .and_then(Path::parent)
2055 .and_then(Path::file_name)
2056 .and_then(|name| name.to_str());
2057 let cwd = slug.and_then(|slug| slug_to_cwd.get(slug)).cloned();
2058 if workspace.is_some_and(|wanted| cwd.as_deref().is_none_or(|actual| !wanted.admits(actual))) {
2059 return None;
2060 }
2061
2062 let file = File::open(path).ok()?;
2066 let mut reader = BufReader::new(file.take(64 * 1024));
2067 let mut header = String::new();
2068 reader.read_line(&mut header).ok()?;
2069 let header = serde_json::from_str::<Value>(&header).ok()?;
2070 let session_id = header.get("sessionId")?.as_str()?.to_string();
2071 let mut model = None;
2072 for line in reader
2073 .take(4 * 1024)
2074 .lines()
2075 .map_while(std::result::Result::ok)
2076 {
2077 let Ok(value) = serde_json::from_str::<Value>(&line) else {
2078 continue;
2079 };
2080 let kind = value.get("type").and_then(Value::as_str);
2081 if kind != Some("user") && kind != Some("gemini") {
2082 continue;
2083 }
2084 if model.is_none() {
2085 model = value
2086 .get("model")
2087 .and_then(Value::as_str)
2088 .map(str::to_string);
2089 }
2090 if model.is_some() {
2091 break;
2092 }
2093 }
2094 Some(SessionDescriptor {
2095 locator: SessionLocator {
2096 harness: HarnessId::from(HarnessId::GEMINI),
2097 session_id,
2098 storage: StorageLocator::File {
2099 path: path.to_path_buf(),
2100 },
2101 },
2102 cwd,
2103 title: None,
2104 preview_candidates: Vec::new(),
2105 latest_message_candidates: Vec::new(),
2106 updated_at_ms: tail_facts(path, HarnessId::GEMINI)
2107 .last_turn_ms
2108 .or_else(|| modified_ms(path)),
2109 message_count: None,
2110 model,
2111 parent_session_id: None,
2112 child_session_count: 0,
2113 nouns: OrchestrationNouns::default(),
2114 })
2115}
2116
2117fn display_text(content: Option<&Value>) -> Option<String> {
2118 match content? {
2119 Value::String(text) => Some(text.clone()),
2120 Value::Array(parts) => Some(
2121 parts
2122 .iter()
2123 .filter_map(|part| part.get("text").and_then(Value::as_str))
2124 .collect::<Vec<_>>()
2125 .join(" ")
2126 .trim()
2127 .to_string(),
2128 ),
2129 _ => None,
2130 }
2131}
2132
2133fn discover_hermes(
2138 db_path: &Path,
2139 workspace: Option<&WorkspaceScope>,
2140 selected: &HashSet<&str>,
2141 found: &mut Vec<SessionDescriptor>,
2142) {
2143 if !db_path.is_file() {
2144 return;
2145 }
2146 let Ok(conn) = Connection::open_with_flags(
2147 db_path,
2148 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
2149 ) else {
2150 return;
2151 };
2152 let fingerprint_ok = ["sessions", "messages", "schema_version"].iter().all(|t| {
2153 conn.query_row(
2154 "SELECT 1 FROM sqlite_master WHERE type='table' AND name=?1",
2155 [t],
2156 |_| Ok(()),
2157 )
2158 .is_ok()
2159 });
2160 if !fingerprint_ok {
2161 return;
2162 }
2163 let Ok(mut statement) = conn.prepare(
2164 "SELECT id, cwd, title, model, message_count, started_at, ended_at, parent_session_id, \
2165 source, model_config FROM sessions ORDER BY started_at DESC",
2166 ) else {
2167 return;
2168 };
2169 let Ok(rows) = statement.query_map([], |row| {
2170 Ok((
2171 row.get::<_, String>(0)?,
2172 row.get::<_, Option<String>>(1)?,
2173 row.get::<_, Option<String>>(2)?,
2174 row.get::<_, Option<String>>(3)?,
2175 row.get::<_, Option<i64>>(4)?,
2176 row.get::<_, Option<f64>>(5)?,
2177 row.get::<_, Option<f64>>(6)?,
2178 row.get::<_, Option<String>>(7)?,
2179 row.get::<_, Option<String>>(8)?,
2180 row.get::<_, Option<String>>(9)?,
2181 ))
2182 }) else {
2183 return;
2184 };
2185 for row in rows.flatten() {
2186 let (
2187 id,
2188 cwd,
2189 title,
2190 model,
2191 message_count,
2192 started_at,
2193 ended_at,
2194 parent,
2195 source,
2196 model_config,
2197 ) = row;
2198 let mirror = model_config
2202 .as_deref()
2203 .and_then(|c| serde_json::from_str::<serde_json::Value>(c).ok())
2204 .and_then(|c| c.get("_supercode_mirror").cloned());
2205 if let Some(mirror) = mirror {
2206 let worker_read = mirror
2207 .get("harness")
2208 .and_then(serde_json::Value::as_str)
2209 .is_some_and(|h| selected.contains(h));
2210 let mirrored = mirror.get("messages").and_then(serde_json::Value::as_i64);
2211 if worker_read && mirrored.is_none_or(|m| message_count.unwrap_or(0) <= m) {
2212 continue;
2213 }
2214 if mirror.get("continued_as").is_some()
2215 && mirrored.is_none_or(|m| message_count.unwrap_or(0) <= m)
2216 {
2217 continue;
2218 }
2219 }
2220 let cwd = cwd.map(PathBuf::from);
2221 if let Some(filter) = workspace {
2222 if !filter.admits_literal(cwd.as_deref()) {
2223 continue;
2224 }
2225 }
2226 let updated_at_ms = ended_at
2227 .or(started_at)
2228 .map(|seconds| (seconds * 1000.0) as u64);
2229 let mut meta = SessionMeta::new(SessionSource::Hermes);
2233 meta.cwd = cwd.clone();
2234 if let Some(hermes_source) = source.filter(|value| !value.is_empty()) {
2235 meta.lineage
2236 .insert("hermes_source".to_string(), hermes_source);
2237 }
2238 if let Some(parent_id) = parent.as_deref() {
2239 meta.lineage.insert(
2240 "hermes_lineage_kind".to_string(),
2241 crate::session::hermes_lineage_kind(
2242 &conn,
2243 parent_id,
2244 model_config.as_deref(),
2245 started_at,
2246 )
2247 .to_string(),
2248 );
2249 }
2250 hermes_capture_nouns(&conn, &id, &mut meta);
2251 found.push(SessionDescriptor {
2252 locator: SessionLocator {
2253 harness: HarnessId::new(HarnessId::HERMES),
2254 session_id: id,
2255 storage: StorageLocator::File {
2256 path: db_path.to_path_buf(),
2257 },
2258 },
2259 cwd,
2260 title: title.filter(|t| !t.is_empty()),
2261 preview_candidates: Vec::new(),
2262 latest_message_candidates: Vec::new(),
2263 updated_at_ms,
2264 message_count: message_count.map(|count| count.max(0) as usize),
2265 model,
2266 parent_session_id: parent,
2267 child_session_count: 0,
2268 nouns: OrchestrationNouns::from_meta(&meta),
2269 });
2270 }
2271}
2272
2273fn discover_orchestrator(
2280 root: &Path,
2281 workspace: Option<&WorkspaceScope>,
2282 found: &mut Vec<SessionDescriptor>,
2283) {
2284 if workspace.is_some() {
2286 return;
2287 }
2288 for (profile, dir) in orchestrator_profile_dirs(root) {
2289 let db_path = dir.join("state.db");
2290 if !db_path.is_file() {
2291 continue;
2292 }
2293 let Ok(conn) = Connection::open_with_flags(
2294 &db_path,
2295 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
2296 ) else {
2297 continue;
2298 };
2299 let sessions = dir.join("sessions");
2302 let scope = std::fs::canonicalize(&sessions)
2303 .unwrap_or(sessions)
2304 .display()
2305 .to_string();
2306 let Ok(mut statement) = conn.prepare(
2307 "SELECT json_extract(entry_json, '$.metadata.supercode.binding'), \
2308 CAST(strftime('%s', json_extract(entry_json, '$.metadata.supercode.binding.last_activity_at')) AS INTEGER) \
2309 FROM gateway_routing WHERE scope = ?1 AND json_extract(entry_json, '$.metadata.supercode') IS NOT NULL \
2310 ORDER BY json_extract(entry_json, '$.metadata.supercode.binding.last_activity_at') DESC",
2311 ) else {
2312 continue;
2313 };
2314 let Ok(rows) = statement.query_map([&scope], |row| {
2315 Ok((
2316 row.get::<_, Option<String>>(0)?,
2317 row.get::<_, Option<i64>>(1)?,
2318 ))
2319 }) else {
2320 continue;
2321 };
2322 let rows = rows.flatten().filter_map(|(json, epoch)| {
2323 let b: Binding = serde_json::from_str(&json?).ok()?;
2324 let text = |v: &Option<String>| v.clone().filter(|s| !s.is_empty());
2325 Some((
2326 OrchestratorBindingRow {
2327 platform: b.key.platform.clone().unwrap_or_default(),
2328 chat_type: b.key.kind.clone().unwrap_or_default(),
2329 chat_id: text(&b.key.chat_id),
2330 thread_id: text(&b.key.thread_id),
2331 participant_id: text(&b.key.participant_id),
2332 worker_harness: b.worker.harness.as_str().to_string(),
2333 worker_session_id: text(&b.worker.session_id),
2334 worker_locator: text(&b.worker.locator),
2335 started_at: b.started_at.clone(),
2336 last_activity_at: b.last_activity_at.clone(),
2337 ended_at: b.ended_at.clone(),
2338 end_reason: b.end_reason.map(|r| r.as_str().to_string()),
2339 handoff_to: b.handoff.as_ref().and_then(|h| h.to.clone()),
2340 handoff_state: b.handoff.as_ref().map(|h| h.state.clone()),
2341 handoff_error: b.handoff.as_ref().and_then(|h| h.error.clone()),
2342 recurrence_job_id: b.recurrence.as_ref().map(|r| r.job_id.clone()),
2343 },
2344 epoch,
2345 ))
2346 });
2347 for (row, last_activity_epoch) in rows {
2348 found.push(orchestrator_descriptor(
2349 &db_path,
2350 &profile,
2351 &row,
2352 last_activity_epoch,
2353 ));
2354 }
2355 }
2356}
2357
2358fn orchestrator_descriptor(
2359 db_path: &Path,
2360 profile: &str,
2361 row: &OrchestratorBindingRow,
2362 last_activity_epoch: Option<i64>,
2363) -> SessionDescriptor {
2364 let binding = Binding::from_orchestrator_row(profile, row);
2365 let nouns = binding.nouns();
2366 let mut title = format!(
2369 "{} {}",
2370 row.worker_harness,
2371 row.worker_session_id
2372 .as_deref()
2373 .unwrap_or("(no worker session yet)")
2374 );
2375 if let Some(reason) = row.end_reason.as_deref().filter(|_| row.ended_at.is_some()) {
2376 title.push_str(&format!(" (ended: {reason})"));
2377 }
2378 SessionDescriptor {
2379 locator: SessionLocator {
2380 harness: HarnessId::new(HarnessId::ORCHESTRATOR),
2381 session_id: row.worker_session_id.clone().unwrap_or_default(),
2382 storage: StorageLocator::File {
2383 path: row
2384 .worker_locator
2385 .clone()
2386 .map_or_else(|| db_path.to_path_buf(), PathBuf::from),
2387 },
2388 },
2389 cwd: None,
2390 title: Some(title),
2391 preview_candidates: Vec::new(),
2392 latest_message_candidates: Vec::new(),
2393 updated_at_ms: last_activity_epoch.map(|seconds| (seconds.max(0) as u64) * 1000),
2394 message_count: None,
2395 model: None,
2396 parent_session_id: None,
2397 child_session_count: 0,
2398 nouns,
2399 }
2400}
2401
2402fn discover_openclaw(
2408 root: &Path,
2409 workspace: Option<&WorkspaceScope>,
2410 found: &mut Vec<SessionDescriptor>,
2411) {
2412 let agents = root.join("agents");
2413 let Ok(agent_dirs) = std::fs::read_dir(&agents) else {
2414 return;
2415 };
2416 for agent_dir in agent_dirs.flatten() {
2417 let sessions = agent_dir.path().join("sessions");
2418 let Ok(files) = std::fs::read_dir(&sessions) else {
2419 continue;
2420 };
2421 for file in files.flatten() {
2422 let path = file.path();
2423 let name = file.file_name();
2424 let name = name.to_string_lossy();
2425 if !name.ends_with(".jsonl") || name.ends_with(".trajectory.jsonl") {
2426 continue;
2427 }
2428 let Ok(text) = std::fs::read_to_string(&path) else {
2429 continue;
2430 };
2431 let Some(header_line) = text.lines().find(|line| !line.trim().is_empty()) else {
2432 continue;
2433 };
2434 let Ok(header) = serde_json::from_str::<serde_json::Value>(header_line) else {
2435 continue;
2436 };
2437 if header.get("type").and_then(serde_json::Value::as_str) != Some("session") {
2438 continue;
2439 }
2440 let session_id = header
2441 .get("id")
2442 .and_then(serde_json::Value::as_str)
2443 .unwrap_or_else(|| name.trim_end_matches(".jsonl"))
2444 .to_string();
2445 let cwd = header
2446 .get("cwd")
2447 .and_then(serde_json::Value::as_str)
2448 .map(PathBuf::from);
2449 if let Some(filter) = workspace {
2450 if !filter.admits_literal(cwd.as_deref()) {
2451 continue;
2452 }
2453 }
2454 let updated_at_ms = tail_facts(&path, HarnessId::OPENCLAW)
2455 .last_turn_ms
2456 .or_else(|| {
2457 file.metadata()
2458 .ok()
2459 .and_then(|metadata| metadata.modified().ok())
2460 .and_then(|modified| modified.duration_since(std::time::UNIX_EPOCH).ok())
2461 .map(|elapsed| elapsed.as_millis() as u64)
2462 });
2463 let message_count = text
2464 .lines()
2465 .filter(|line| line.contains("\"type\":\"message\""))
2466 .count();
2467 let mut meta = SessionMeta::new(SessionSource::OpenClaw);
2471 meta.cwd = cwd.clone();
2472 openclaw_capture_header_nouns(&header, &mut meta);
2473 if meta.profile.is_none() {
2474 meta.profile = openclaw_agent_id_from_path(&path);
2475 }
2476 found.push(SessionDescriptor {
2477 locator: SessionLocator {
2478 harness: HarnessId::new(HarnessId::OPENCLAW),
2479 session_id,
2480 storage: StorageLocator::File { path },
2481 },
2482 cwd,
2483 title: None,
2484 preview_candidates: Vec::new(),
2485 latest_message_candidates: Vec::new(),
2486 updated_at_ms,
2487 message_count: Some(message_count),
2488 model: None,
2489 parent_session_id: None,
2490 child_session_count: 0,
2491 nouns: OrchestrationNouns::from_meta(&meta),
2492 });
2493 }
2494 }
2495}
2496
2497fn discover_supercode(
2498 root: &Path,
2499 workspace: Option<&WorkspaceScope>,
2500 found: &mut Vec<SessionDescriptor>,
2501) {
2502 for info in list_native_store(root) {
2503 let path = if info.archived {
2504 root.join("archived").join(format!("{}.jsonl", info.name))
2505 } else {
2506 root.join(format!("{}.jsonl", info.name))
2507 };
2508 let sidecar = path.with_extension("sidecar.jsonl");
2513 let path = if path.is_file() {
2514 path
2515 } else {
2516 sidecar.clone()
2517 };
2518 let header = read_native_store_header(&path);
2519 if workspace.is_some_and(|wanted| {
2520 header
2521 .as_ref()
2522 .and_then(|meta| meta.cwd.as_deref())
2523 .is_none_or(|cwd| !wanted.admits(cwd))
2524 }) {
2525 continue;
2526 }
2527 let title = (!info.title.trim().is_empty()).then_some(info.title);
2528 let updated_at_ms = tail_facts(&sidecar, HarnessId::SUPERCODE)
2531 .last_turn_ms
2532 .or_else(|| tail_facts(&path, HarnessId::SUPERCODE).last_turn_ms)
2533 .or_else(|| modified_ms(&path))
2534 .or_else(|| modified_ms(&sidecar));
2535 found.push(SessionDescriptor {
2536 locator: SessionLocator {
2537 harness: HarnessId::from(HarnessId::SUPERCODE),
2538 session_id: info.name,
2539 storage: StorageLocator::File { path: path.clone() },
2540 },
2541 cwd: header.as_ref().and_then(|meta| meta.cwd.clone()),
2542 title,
2543 preview_candidates: Vec::new(),
2544 latest_message_candidates: Vec::new(),
2545 updated_at_ms,
2546 message_count: None,
2547 model: header.and_then(|meta| meta.model),
2548 parent_session_id: None,
2549 child_session_count: 0,
2550 nouns: OrchestrationNouns::default(),
2551 });
2552 }
2553}
2554
2555fn read_native_store_header(path: &Path) -> Option<HeaderMeta> {
2560 let name = path.file_stem()?.to_str()?;
2561 let sidecar = path.with_file_name(format!("{name}.sidecar.jsonl"));
2562 let source_path = if sidecar.is_file() {
2563 sidecar
2564 } else {
2565 path.to_path_buf()
2566 };
2567 let file = File::open(source_path).ok()?;
2568 let mut result = HeaderMeta::default();
2569 let mut source = None;
2570 let mut bytes = 0usize;
2571 for line in BufReader::new(file).lines().take(32) {
2572 let line = line.ok()?;
2573 bytes += line.len();
2574 if bytes > 256 * 1024 {
2575 break;
2576 }
2577 let Ok(value) = serde_json::from_str::<Value>(&line) else {
2578 continue;
2579 };
2580 if source.is_none() {
2581 source = value.get("source").and_then(Value::as_str).map(|source| {
2582 if source == "claude_code" {
2583 HarnessId::CLAUDE_CODE.to_string()
2584 } else {
2585 source.to_string()
2586 }
2587 });
2588 fill_string(&mut result.session_id, value.get("session_id"));
2589 }
2590 if let Some(harness) = source.as_deref() {
2591 update_header_meta(&mut result, &value, harness);
2592 }
2593 if result.cwd.is_some() && result.model.is_some() {
2594 break;
2595 }
2596 }
2597 Some(result)
2598}
2599
2600#[derive(Deserialize)]
2601struct NativeStoreInfo {
2602 name: String,
2603 #[serde(default)]
2604 title: String,
2605 #[serde(skip)]
2606 archived: bool,
2607}
2608
2609fn list_native_store(root: &Path) -> Vec<NativeStoreInfo> {
2610 let mut sessions = Vec::new();
2611 for archived in [false, true] {
2612 let directory = if archived {
2613 root.join("archived")
2614 } else {
2615 root.to_path_buf()
2616 };
2617 let Ok(entries) = fs::read_dir(directory) else {
2618 continue;
2619 };
2620 for entry in entries.flatten() {
2621 let path = entry.path();
2622 if !path.to_string_lossy().ends_with(".meta.json") {
2623 continue;
2624 }
2625 let Ok(text) = fs::read_to_string(path) else {
2626 continue;
2627 };
2628 let Ok(mut info) = serde_json::from_str::<NativeStoreInfo>(&text) else {
2629 continue;
2630 };
2631 info.archived = archived;
2632 sessions.push(info);
2633 }
2634 }
2635 sessions.sort_by(|left, right| left.name.cmp(&right.name));
2636 sessions
2637}
2638
2639fn discover_grok(
2640 root: &Path,
2641 workspace: Option<&WorkspaceScope>,
2642 found: &mut Vec<SessionDescriptor>,
2643) {
2644 let Ok(workspaces) = fs::read_dir(root) else {
2645 return;
2646 };
2647 for workspace_entry in workspaces.flatten() {
2648 let encoded = workspace_entry.file_name();
2649 let Some(cwd) = encoded
2650 .to_str()
2651 .and_then(percent_decode_path)
2652 .map(PathBuf::from)
2653 else {
2654 continue;
2655 };
2656 if workspace.is_some_and(|wanted| !wanted.admits(&cwd)) {
2657 continue;
2658 }
2659 let Ok(sessions) = fs::read_dir(workspace_entry.path()) else {
2660 continue;
2661 };
2662 for session_entry in sessions.flatten() {
2663 let session_dir = session_entry.path();
2664 if !session_dir.is_dir() {
2665 continue;
2666 }
2667 let transcript = session_dir.join("chat_history.jsonl");
2668 if !transcript.is_file() {
2669 continue;
2670 }
2671 let Some(session_id) = session_dir
2672 .file_name()
2673 .and_then(|name| name.to_str())
2674 .map(str::to_string)
2675 else {
2676 continue;
2677 };
2678 let summary = fs::read_to_string(session_dir.join("summary.json"))
2679 .ok()
2680 .and_then(|text| serde_json::from_str::<Value>(&text).ok());
2681 let title = summary
2682 .as_ref()
2683 .and_then(|value| value.get("generated_title"))
2684 .and_then(Value::as_str)
2685 .filter(|title| !title.is_empty())
2686 .map(str::to_string);
2687 let model = summary
2688 .as_ref()
2689 .and_then(|value| value.get("current_model_id"))
2690 .and_then(Value::as_str)
2691 .map(str::to_string);
2692 let message_count = summary
2693 .as_ref()
2694 .and_then(|value| value.get("num_chat_messages"))
2695 .and_then(Value::as_u64)
2696 .and_then(|count| usize::try_from(count).ok());
2697 let updated_at_ms = summary
2698 .as_ref()
2699 .and_then(|value| value.get("updated_at"))
2700 .and_then(Value::as_str)
2701 .and_then(crate::sidecar::rfc3339_to_ms)
2702 .and_then(|millis| u64::try_from(millis).ok())
2703 .or_else(|| modified_ms(&transcript));
2704 found.push(SessionDescriptor {
2705 locator: SessionLocator {
2706 harness: HarnessId::from(HarnessId::GROK),
2707 session_id,
2708 storage: StorageLocator::File { path: transcript },
2709 },
2710 cwd: Some(cwd.clone()),
2711 title,
2712 preview_candidates: Vec::new(),
2713 latest_message_candidates: Vec::new(),
2714 updated_at_ms,
2715 message_count,
2716 model,
2717 parent_session_id: None,
2718 child_session_count: 0,
2719 nouns: OrchestrationNouns::default(),
2720 });
2721 }
2722 }
2723}
2724
2725fn discover_opencode(
2726 root: &Path,
2727 workspace: Option<&WorkspaceScope>,
2728 found: &mut Vec<SessionDescriptor>,
2729) {
2730 let mut dbs = Vec::new();
2731 if root.is_file() {
2732 dbs.push(root.to_path_buf());
2733 } else if let Ok(entries) = fs::read_dir(root) {
2734 dbs.extend(entries.flatten().map(|entry| entry.path()).filter(|path| {
2735 path.file_name()
2736 .and_then(|v| v.to_str())
2737 .is_some_and(|name| name.starts_with("opencode") && name.ends_with(".db"))
2738 }));
2739 }
2740 dbs.sort();
2741 for db in dbs {
2742 let Ok(conn) = Connection::open_with_flags(
2743 &db,
2744 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
2745 ) else {
2746 continue;
2747 };
2748 let has_model = conn.prepare("SELECT model FROM session LIMIT 0").is_ok();
2749 let model_column = if has_model { "s.model" } else { "NULL" };
2750 let query = format!(
2751 "SELECT s.id, s.directory, s.title, s.time_updated, {model_column}, COUNT(m.id) \
2752 FROM session s LEFT JOIN message m ON m.session_id = s.id \
2753 GROUP BY s.id ORDER BY s.time_updated DESC"
2754 );
2755 let Ok(mut stmt) = conn.prepare(&query) else {
2756 continue;
2757 };
2758 let Ok(rows) = stmt.query_map([], |row| {
2759 Ok((
2760 row.get::<_, String>(0)?,
2761 row.get::<_, String>(1)?,
2762 row.get::<_, String>(2)?,
2763 row.get::<_, i64>(3)?,
2764 row.get::<_, Option<String>>(4)?,
2765 row.get::<_, i64>(5)?,
2766 ))
2767 }) else {
2768 continue;
2769 };
2770 for row in rows.flatten() {
2771 let (id, cwd, title, updated, model, messages) = row;
2772 let cwd = PathBuf::from(cwd);
2773 if workspace.is_some_and(|wanted| !wanted.admits(&cwd)) {
2774 continue;
2775 }
2776 found.push(SessionDescriptor {
2777 locator: SessionLocator {
2778 harness: HarnessId::from(HarnessId::OPENCODE),
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.is_empty()).then_some(title),
2787 preview_candidates: Vec::new(),
2788 latest_message_candidates: Vec::new(),
2789 updated_at_ms: u64::try_from(updated).ok(),
2790 message_count: usize::try_from(messages).ok(),
2791 model,
2792 parent_session_id: None,
2793 child_session_count: 0,
2794 nouns: OrchestrationNouns::default(),
2795 });
2796 }
2797 }
2798}
2799
2800fn discover_goose(
2801 root: &Path,
2802 workspace: Option<&WorkspaceScope>,
2803 found: &mut Vec<SessionDescriptor>,
2804) {
2805 let db = if root.is_file() {
2806 root.to_path_buf()
2807 } else if root.join("sessions.db").is_file() {
2808 root.join("sessions.db")
2809 } else {
2810 root.join("sessions/sessions.db")
2811 };
2812 let Ok(connection) = Connection::open_with_flags(
2813 &db,
2814 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
2815 ) else {
2816 return;
2817 };
2818 let Ok(mut statement) = connection.prepare(
2819 "SELECT s.id, s.working_dir, s.name, s.updated_at, s.model_config_json, \
2820 COUNT(m.id) \
2821 FROM sessions s LEFT JOIN messages m ON m.session_id = s.id \
2822 WHERE s.archived_at IS NULL \
2823 GROUP BY s.id ORDER BY s.updated_at DESC",
2824 ) else {
2825 return;
2826 };
2827 let Ok(rows) = statement.query_map([], |row| {
2828 Ok((
2829 row.get::<_, String>(0)?,
2830 row.get::<_, String>(1)?,
2831 row.get::<_, String>(2)?,
2832 row.get::<_, String>(3)?,
2833 row.get::<_, Option<String>>(4)?,
2834 row.get::<_, i64>(5)?,
2835 ))
2836 }) else {
2837 return;
2838 };
2839 for row in rows.flatten() {
2840 let (id, cwd, title, updated_at, model_config, message_count) = row;
2841 let cwd = PathBuf::from(cwd);
2842 if workspace.is_some_and(|wanted| !wanted.admits(&cwd)) {
2843 continue;
2844 }
2845 let model = model_config
2846 .as_deref()
2847 .and_then(|value| serde_json::from_str::<Value>(value).ok())
2848 .and_then(|value| {
2849 value
2850 .get("model_name")
2851 .or_else(|| value.get("modelName"))
2852 .and_then(Value::as_str)
2853 .map(str::to_string)
2854 });
2855 let updated_at_ms = crate::sidecar::rfc3339_to_ms(&updated_at)
2856 .or_else(|| {
2857 crate::sidecar::rfc3339_to_ms(&format!("{}Z", updated_at.replace(' ', "T")))
2859 })
2860 .and_then(|value| u64::try_from(value).ok());
2861 found.push(SessionDescriptor {
2862 locator: SessionLocator {
2863 harness: HarnessId::from(HarnessId::GOOSE),
2864 session_id: id.clone(),
2865 storage: StorageLocator::Sqlite {
2866 path: db.clone(),
2867 selector: id,
2868 },
2869 },
2870 cwd: Some(cwd),
2871 title: (!title.trim().is_empty()).then_some(title),
2872 preview_candidates: Vec::new(),
2873 latest_message_candidates: Vec::new(),
2874 updated_at_ms,
2875 message_count: usize::try_from(message_count).ok(),
2876 model,
2877 parent_session_id: None,
2878 child_session_count: 0,
2879 nouns: OrchestrationNouns::default(),
2880 });
2881 }
2882}
2883
2884const LATEST_PREVIEW_CANDIDATES: usize = 8;
2885const TOPIC_PREVIEW_HEAD_BYTES: u64 = 512 * 1024;
2886const LATEST_PREVIEW_TAIL_BYTES: u64 = 512 * 1024;
2887const LATEST_PREVIEW_MAX_BYTES: u64 = 4 * 1024 * 1024;
2888
2889fn topic_message_candidates(locator: &SessionLocator) -> Result<Vec<SessionPreviewCandidate>> {
2890 match &locator.storage {
2891 StorageLocator::File { path }
2892 if matches!(
2893 locator.harness.as_str(),
2894 HarnessId::CLAUDE_CODE | HarnessId::CODEX
2895 ) =>
2896 {
2897 topic_file_message_candidates(path, locator.harness.as_str())
2898 }
2899 _ => Ok(Vec::new()),
2900 }
2901}
2902
2903fn codex_history_topics(
2904 sessions_root: &Path,
2905 sessions: &[SessionDescriptor],
2906) -> Result<HashMap<String, Vec<SessionPreviewCandidate>>> {
2907 let wanted: HashSet<&str> = sessions
2908 .iter()
2909 .filter(|descriptor| descriptor.locator.harness.as_str() == HarnessId::CODEX)
2910 .map(|descriptor| descriptor.locator.session_id.as_str())
2911 .collect();
2912 if wanted.is_empty() {
2913 return Ok(HashMap::new());
2914 }
2915 let Some(root) = sessions_root.parent() else {
2916 return Ok(HashMap::new());
2917 };
2918 let file = match File::open(root.join("history.jsonl")) {
2919 Ok(file) => file,
2920 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(HashMap::new()),
2921 Err(error) => return Err(error.into()),
2922 };
2923 let mut topics = HashMap::new();
2924 for line in BufReader::new(file).lines() {
2925 let Ok(value) = serde_json::from_str::<Value>(&line?) else {
2926 continue;
2927 };
2928 let Some(session_id) = value.get("session_id").and_then(Value::as_str) else {
2929 continue;
2930 };
2931 if !wanted.contains(session_id) || topics.contains_key(session_id) {
2932 continue;
2933 }
2934 let mut candidates = Vec::new();
2935 push_message_candidate(&mut candidates, "user", value.get("text"), HashMap::new());
2936 if !candidates.is_empty() {
2937 topics.insert(session_id.to_string(), candidates);
2938 if topics.len() == wanted.len() {
2939 break;
2940 }
2941 }
2942 }
2943 Ok(topics)
2944}
2945
2946fn latest_message_candidates(locator: &SessionLocator) -> Result<Vec<SessionPreviewCandidate>> {
2947 match &locator.storage {
2948 StorageLocator::File { path } | StorageLocator::Sqlite { path, .. }
2951 if locator.harness.as_str() == HarnessId::HERMES =>
2952 {
2953 latest_hermes_message_candidates(path, &locator.session_id)
2954 }
2955 StorageLocator::File { path } => {
2956 latest_file_message_candidates(path, locator.harness.as_str())
2957 }
2958 StorageLocator::Sqlite { path, selector }
2959 if locator.harness.as_str() == HarnessId::OPENCODE =>
2960 {
2961 latest_opencode_message_candidates(path, selector)
2962 }
2963 StorageLocator::Sqlite { path, selector }
2964 if locator.harness.as_str() == HarnessId::GOOSE =>
2965 {
2966 latest_goose_message_candidates(path, selector)
2967 }
2968 StorageLocator::Sqlite { .. } => Ok(Vec::new()),
2969 }
2970}
2971
2972fn topic_file_message_candidates(
2973 path: &Path,
2974 harness: &str,
2975) -> Result<Vec<SessionPreviewCandidate>> {
2976 let mut file = File::open(path)?;
2977 let mut bytes = Vec::with_capacity(TOPIC_PREVIEW_HEAD_BYTES as usize);
2978 file.by_ref()
2979 .take(TOPIC_PREVIEW_HEAD_BYTES)
2980 .read_to_end(&mut bytes)?;
2981 if file.metadata()?.len() > TOPIC_PREVIEW_HEAD_BYTES {
2982 if let Some(newline) = bytes.iter().rposition(|byte| *byte == b'\n') {
2983 bytes.truncate(newline);
2984 }
2985 }
2986 let text = String::from_utf8(bytes).map_err(|_| {
2987 Error::Other(format!(
2988 "{} contains non-UTF-8 data in its topic-preview window",
2989 path.display()
2990 ))
2991 })?;
2992 if harness == HarnessId::CODEX {
2993 return Ok(codex_preview_candidates(text.lines(), false));
2994 }
2995 let mut candidates = Vec::new();
2996 for line in text.lines() {
2997 let Ok(value) = serde_json::from_str::<Value>(line) else {
2998 continue;
2999 };
3000 push_topic_message_candidate(&mut candidates, harness, &value);
3001 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3002 break;
3003 }
3004 }
3005 Ok(candidates)
3006}
3007
3008fn latest_file_message_candidates(
3009 path: &Path,
3010 harness: &str,
3011) -> Result<Vec<SessionPreviewCandidate>> {
3012 let mut candidates =
3013 latest_file_message_candidates_with_limit(path, harness, LATEST_PREVIEW_TAIL_BYTES)?;
3014 if candidates.is_empty() {
3015 candidates =
3016 latest_file_message_candidates_with_limit(path, harness, LATEST_PREVIEW_MAX_BYTES)?;
3017 }
3018 Ok(candidates)
3019}
3020
3021fn latest_file_message_candidates_with_limit(
3022 path: &Path,
3023 harness: &str,
3024 byte_limit: u64,
3025) -> Result<Vec<SessionPreviewCandidate>> {
3026 let mut file = File::open(path)?;
3027 let file_len = file.metadata()?.len();
3028 let start = file_len.saturating_sub(byte_limit);
3029 file.seek(SeekFrom::Start(start))?;
3030 let mut bytes = Vec::with_capacity((file_len - start) as usize);
3031 file.read_to_end(&mut bytes)?;
3032 if start > 0 {
3033 if let Some(newline) = bytes.iter().position(|byte| *byte == b'\n') {
3034 bytes.drain(..=newline);
3035 } else {
3036 return Ok(Vec::new());
3037 }
3038 }
3039 let text = String::from_utf8(bytes).map_err(|_| {
3040 Error::Other(format!(
3041 "{} contains non-UTF-8 data in its list-preview window",
3042 path.display()
3043 ))
3044 })?;
3045 if harness == HarnessId::CODEX {
3046 return Ok(codex_preview_candidates(text.lines().rev(), true));
3047 }
3048 let mut candidates = Vec::new();
3049 for line in text.lines().rev() {
3050 let Ok(value) = serde_json::from_str::<Value>(line) else {
3051 continue;
3052 };
3053 let (role, content, metadata) = match harness {
3054 HarnessId::CLAUDE_CODE => {
3055 let role = value.get("type").and_then(Value::as_str);
3056 if !matches!(role, Some("user" | "assistant")) {
3057 continue;
3058 }
3059 let metadata = if role == Some("user") {
3060 crate::session::claude_user_provenance(&value)
3061 .into_iter()
3062 .collect()
3063 } else {
3064 HashMap::new()
3065 };
3066 (
3067 role.unwrap_or_default(),
3068 value
3069 .get("message")
3070 .and_then(|message| message.get("content")),
3071 metadata,
3072 )
3073 }
3074 HarnessId::PI => {
3075 if value.get("type").and_then(Value::as_str) != Some("message") {
3076 continue;
3077 }
3078 let message = value.get("message").unwrap_or(&Value::Null);
3079 let Some(role @ ("user" | "assistant")) =
3080 message.get("role").and_then(Value::as_str)
3081 else {
3082 continue;
3083 };
3084 (role, message.get("content"), HashMap::new())
3085 }
3086 HarnessId::GEMINI => {
3087 let Some(kind @ ("user" | "gemini")) = value.get("type").and_then(Value::as_str)
3088 else {
3089 continue;
3090 };
3091 (
3092 if kind == "gemini" {
3093 "assistant"
3094 } else {
3095 "user"
3096 },
3097 value.get("content"),
3098 HashMap::new(),
3099 )
3100 }
3101 HarnessId::GROK => {
3102 let Some(role @ ("user" | "assistant")) = value.get("type").and_then(Value::as_str)
3103 else {
3104 continue;
3105 };
3106 (role, value.get("content"), HashMap::new())
3107 }
3108 HarnessId::SUPERCODE => {
3109 let Some(role @ ("user" | "assistant")) = value.get("role").and_then(Value::as_str)
3110 else {
3111 continue;
3112 };
3113 (role, value.get("content"), HashMap::new())
3114 }
3115 _ => continue,
3116 };
3117 let mut metadata = metadata;
3118 if matches!(harness, HarnessId::CLAUDE_CODE | HarnessId::CODEX) {
3119 if let Some(timestamp) = value.get("timestamp").and_then(Value::as_str) {
3120 metadata.insert("timestamp".to_string(), timestamp.to_string());
3121 }
3122 }
3123 push_message_candidate_with_cursor(
3124 &mut candidates,
3125 role,
3126 content,
3127 metadata,
3128 Some(message_candidate_cursor(harness, &value)),
3129 );
3130 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3131 break;
3132 }
3133 }
3134 Ok(candidates)
3135}
3136
3137struct CodexPreviewRecord {
3138 native: Value,
3139 role: String,
3140 text: String,
3141}
3142
3143fn codex_preview_candidates<'a>(
3149 lines: impl Iterator<Item = &'a str>,
3150 latest: bool,
3151) -> Vec<SessionPreviewCandidate> {
3152 let mut candidates = Vec::new();
3153 let mut pending: Option<CodexPreviewRecord> = None;
3154 for line in lines {
3155 let Ok(native) = serde_json::from_str::<Value>(line) else {
3156 continue;
3157 };
3158 let Some((role, content)) = codex_preview_message(&native) else {
3159 continue;
3160 };
3161 let Some(text) = display_text(content).filter(|text| !text.trim().is_empty()) else {
3162 continue;
3163 };
3164 let current = CodexPreviewRecord {
3165 role: role.to_string(),
3166 text,
3167 native,
3168 };
3169 if let Some(previous) = pending.take() {
3170 if previous.role == current.role
3171 && previous.text == current.text
3172 && previous.native.get("type") != current.native.get("type")
3173 {
3174 let canonical = if previous.native.get("type").and_then(Value::as_str)
3175 == Some("response_item")
3176 {
3177 previous
3178 } else {
3179 current
3180 };
3181 push_codex_preview_candidate(&mut candidates, canonical, latest);
3182 } else {
3183 push_codex_preview_candidate(&mut candidates, previous, latest);
3184 pending = Some(current);
3185 }
3186 } else {
3187 pending = Some(current);
3188 }
3189 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3190 break;
3191 }
3192 }
3193 if let Some(last) = pending {
3194 push_codex_preview_candidate(&mut candidates, last, latest);
3195 }
3196 candidates
3197}
3198
3199fn push_codex_preview_candidate(
3200 candidates: &mut Vec<SessionPreviewCandidate>,
3201 record: CodexPreviewRecord,
3202 latest: bool,
3203) {
3204 let mut metadata = HashMap::new();
3205 if latest {
3206 if let Some(timestamp) = record.native.get("timestamp").and_then(Value::as_str) {
3207 metadata.insert("timestamp".to_string(), timestamp.to_string());
3208 }
3209 }
3210 let cursor = latest.then(|| message_candidate_cursor(HarnessId::CODEX, &record.native));
3211 push_message_candidate_with_cursor(
3212 candidates,
3213 &record.role,
3214 Some(&Value::String(record.text)),
3215 metadata,
3216 cursor,
3217 );
3218}
3219
3220fn codex_preview_message(value: &Value) -> Option<(&str, Option<&Value>)> {
3222 let payload = value.get("payload")?;
3223 match (
3224 value.get("type").and_then(Value::as_str)?,
3225 payload.get("type").and_then(Value::as_str)?,
3226 ) {
3227 ("response_item", "message") => {
3228 let role @ ("user" | "assistant") = payload.get("role").and_then(Value::as_str)? else {
3229 return None;
3230 };
3231 Some((role, payload.get("content")))
3232 }
3233 ("event_msg", "user_message") => Some(("user", payload.get("message"))),
3234 ("event_msg", "agent_message") => Some(("assistant", payload.get("message"))),
3235 _ => None,
3236 }
3237}
3238
3239fn push_topic_message_candidate(
3240 candidates: &mut Vec<SessionPreviewCandidate>,
3241 harness: &str,
3242 value: &Value,
3243) {
3244 let (role, content, metadata) = match harness {
3245 HarnessId::CLAUDE_CODE => {
3246 let role = value.get("type").and_then(Value::as_str);
3247 if !matches!(role, Some("user" | "assistant")) {
3248 return;
3249 }
3250 let metadata = if role == Some("user") {
3251 crate::session::claude_user_provenance(value)
3252 .into_iter()
3253 .collect()
3254 } else {
3255 HashMap::new()
3256 };
3257 (
3258 role.unwrap_or_default(),
3259 value
3260 .get("message")
3261 .and_then(|message| message.get("content")),
3262 metadata,
3263 )
3264 }
3265 _ => return,
3266 };
3267 push_message_candidate(candidates, role, content, metadata);
3268}
3269
3270fn latest_opencode_message_candidates(
3271 path: &Path,
3272 session_id: &str,
3273) -> Result<Vec<SessionPreviewCandidate>> {
3274 let connection = Connection::open_with_flags(
3275 path,
3276 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
3277 )
3278 .map_err(|error| Error::Other(format!("failed to open list-preview store: {error}")))?;
3279 let mut statement = connection
3280 .prepare(
3281 "SELECT m.data, p.data FROM message m JOIN part p ON p.message_id = m.id \
3282 WHERE m.session_id = ?1 ORDER BY m.time_created DESC, p.time_created DESC LIMIT 32",
3283 )
3284 .map_err(|error| Error::Other(format!("failed to prepare list-preview query: {error}")))?;
3285 let rows = statement
3286 .query_map([session_id], |row| {
3287 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
3288 })
3289 .map_err(|error| Error::Other(format!("failed to read list-preview rows: {error}")))?;
3290 let mut candidates = Vec::new();
3291 for row in rows.flatten() {
3292 let (Ok(message), Ok(part)) = (
3293 serde_json::from_str::<Value>(&row.0),
3294 serde_json::from_str::<Value>(&row.1),
3295 ) else {
3296 continue;
3297 };
3298 let Some(role @ ("user" | "assistant")) = message.get("role").and_then(Value::as_str)
3299 else {
3300 continue;
3301 };
3302 if part.get("type").and_then(Value::as_str) != Some("text") {
3303 continue;
3304 }
3305 push_message_candidate(&mut candidates, role, part.get("text"), HashMap::new());
3306 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3307 break;
3308 }
3309 }
3310 Ok(candidates)
3311}
3312
3313fn latest_hermes_message_candidates(
3314 path: &Path,
3315 session_id: &str,
3316) -> Result<Vec<SessionPreviewCandidate>> {
3317 let connection = Connection::open_with_flags(
3318 path,
3319 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
3320 )
3321 .map_err(|error| Error::Other(format!("failed to open list-preview store: {error}")))?;
3322 let mut statement = connection
3323 .prepare(
3324 "SELECT role, content FROM messages WHERE session_id = ?1 AND active = 1 \
3325 AND role IN ('user', 'assistant') AND content IS NOT NULL AND content != '' \
3326 ORDER BY timestamp DESC, id DESC LIMIT 32",
3327 )
3328 .map_err(|error| Error::Other(format!("failed to prepare list-preview query: {error}")))?;
3329 let rows = statement
3330 .query_map([session_id], |row| {
3331 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
3332 })
3333 .map_err(|error| Error::Other(format!("failed to read list-preview rows: {error}")))?;
3334 let mut candidates = Vec::new();
3335 for (role, content) in rows.flatten() {
3336 push_message_candidate(
3337 &mut candidates,
3338 &role,
3339 Some(&Value::String(content)),
3340 HashMap::new(),
3341 );
3342 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3343 break;
3344 }
3345 }
3346 Ok(candidates)
3347}
3348
3349fn latest_goose_message_candidates(
3350 path: &Path,
3351 session_id: &str,
3352) -> Result<Vec<SessionPreviewCandidate>> {
3353 let connection = Connection::open_with_flags(
3354 path,
3355 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX,
3356 )
3357 .map_err(|error| Error::Other(format!("failed to open list-preview store: {error}")))?;
3358 let mut statement = connection
3359 .prepare(
3360 "SELECT role, content_json FROM messages WHERE session_id = ?1 \
3361 ORDER BY created_timestamp DESC, id DESC LIMIT 16",
3362 )
3363 .map_err(|error| Error::Other(format!("failed to prepare list-preview query: {error}")))?;
3364 let rows = statement
3365 .query_map([session_id], |row| {
3366 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
3367 })
3368 .map_err(|error| Error::Other(format!("failed to read list-preview rows: {error}")))?;
3369 let mut candidates = Vec::new();
3370 for row in rows.flatten() {
3371 let (role, content) = row;
3372 if !matches!(role.as_str(), "user" | "assistant") {
3373 continue;
3374 }
3375 let Ok(content) = serde_json::from_str::<Value>(&content) else {
3376 continue;
3377 };
3378 push_message_candidate(&mut candidates, &role, Some(&content), HashMap::new());
3379 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3380 break;
3381 }
3382 }
3383 Ok(candidates)
3384}
3385
3386fn push_message_candidate(
3387 candidates: &mut Vec<SessionPreviewCandidate>,
3388 role: &str,
3389 content: Option<&Value>,
3390 metadata: HashMap<String, String>,
3391) {
3392 push_message_candidate_with_cursor(candidates, role, content, metadata, None);
3393}
3394
3395fn push_message_candidate_with_cursor(
3396 candidates: &mut Vec<SessionPreviewCandidate>,
3397 role: &str,
3398 content: Option<&Value>,
3399 metadata: HashMap<String, String>,
3400 cursor: Option<String>,
3401) {
3402 if candidates.len() >= LATEST_PREVIEW_CANDIDATES {
3403 return;
3404 }
3405 let Some(text) = display_text(content).filter(|text| !text.trim().is_empty()) else {
3406 return;
3407 };
3408 const MAX_CHARS: usize = 4_096;
3409 candidates.push(SessionPreviewCandidate {
3410 cursor,
3411 role: role.to_string(),
3412 content: text.chars().take(MAX_CHARS).collect(),
3413 metadata,
3414 });
3415}
3416
3417fn message_candidate_cursor(harness: &str, value: &Value) -> String {
3418 let native_identity = value
3419 .get("uuid")
3420 .or_else(|| value.get("id"))
3421 .or_else(|| value.pointer("/message/id"))
3422 .or_else(|| value.pointer("/payload/id"))
3423 .and_then(Value::as_str)
3424 .or_else(|| value.get("timestamp").and_then(Value::as_str));
3425 let mut hasher = blake3::Hasher::new();
3426 hasher.update(b"supercode.session-preview-cursor.v1\0");
3427 hasher.update(harness.as_bytes());
3428 hasher.update(b"\0");
3429 if let Some(identity) = native_identity {
3430 hasher.update(identity.as_bytes());
3431 } else {
3432 hasher.update(value.to_string().as_bytes());
3436 }
3437 format!("v1:{}", &hasher.finalize().to_hex()[..24])
3438}
3439
3440fn fill_string(target: &mut Option<String>, value: Option<&Value>) {
3441 if target.is_none() {
3442 *target = value.and_then(Value::as_str).map(str::to_owned);
3443 }
3444}
3445
3446fn fill_path(target: &mut Option<PathBuf>, value: Option<&Value>) {
3447 if target.is_none() {
3448 *target = value.and_then(Value::as_str).map(PathBuf::from);
3449 }
3450}
3451
3452#[derive(Default)]
3455struct TailFacts {
3456 last_turn_ms: Option<u64>,
3458 model: Option<String>,
3460}
3461
3462fn tail_facts(path: &Path, harness: &str) -> TailFacts {
3478 let mut facts = TailFacts::default();
3479 let Ok(mut file) = File::open(path) else {
3480 return facts;
3481 };
3482 let Ok(len) = file.metadata().map(|meta| meta.len()) else {
3483 return facts;
3484 };
3485 let mut window = TAIL_SCAN_START.min(len);
3486 loop {
3487 if file.seek(SeekFrom::Start(len - window)).is_err() {
3488 return facts;
3489 }
3490 let Ok(size) = usize::try_from(window) else {
3491 return facts;
3492 };
3493 let mut buf = vec![0u8; size];
3494 if file.read_exact(&mut buf).is_err() {
3495 return facts;
3496 }
3497 let floor = if window < len {
3508 buf.iter().position(|byte| *byte == b'\n').map(|at| at + 1)
3509 } else {
3510 Some(0)
3511 };
3512 if let Some(floor) = floor {
3513 let mut end = buf.len();
3514 while end > floor && !(facts.last_turn_ms.is_some() && facts.model.is_some()) {
3515 let start = buf[floor..end]
3516 .iter()
3517 .rposition(|byte| *byte == b'\n')
3518 .map_or(floor, |at| floor + at + 1);
3519 if let Ok(record) = std::str::from_utf8(&buf[start..end])
3520 .map_err(|_| ())
3521 .and_then(|line| serde_json::from_str::<Value>(line).map_err(|_| ()))
3522 {
3523 if facts.last_turn_ms.is_none() {
3524 facts.last_turn_ms = record_timestamp(&record, harness)
3525 .and_then(crate::sidecar::rfc3339_to_ms)
3526 .and_then(|millis| u64::try_from(millis).ok());
3527 }
3528 if facts.model.is_none() {
3529 facts.model = record_model(&record, harness).map(str::to_owned);
3530 }
3531 }
3532 end = start.saturating_sub(1);
3533 }
3534 }
3535 if facts.last_turn_ms.is_some() || window >= len || window >= TAIL_SCAN_LIMIT {
3536 return facts;
3537 }
3538 window = (window * 2).min(len).min(TAIL_SCAN_LIMIT);
3539 }
3540}
3541
3542fn record_timestamp<'a>(record: &'a Value, harness: &str) -> Option<&'a str> {
3546 let key = match harness {
3547 HarnessId::SUPERCODE => "ts",
3548 _ => "timestamp",
3549 };
3550 record.get(key)?.as_str()
3551}
3552
3553fn record_model<'a>(record: &'a Value, harness: &str) -> Option<&'a str> {
3556 match harness {
3557 HarnessId::CLAUDE_CODE | HarnessId::PI => record.get("message")?.get("model")?.as_str(),
3558 HarnessId::CODEX => {
3559 if record.get("type")?.as_str()? != "turn_context" {
3560 return None;
3561 }
3562 record.get("payload")?.get("model")?.as_str()
3563 }
3564 _ => None,
3565 }
3566}
3567
3568const TAIL_SCAN_START: u64 = 16 * 1024;
3571
3572const TAIL_SCAN_LIMIT: u64 = 1024 * 1024;
3575
3576fn modified_ms(path: &Path) -> Option<u64> {
3577 fs::metadata(path)
3578 .ok()?
3579 .modified()
3580 .ok()?
3581 .duration_since(UNIX_EPOCH)
3582 .ok()
3583 .and_then(|duration| u64::try_from(duration.as_millis()).ok())
3584}
3585
3586pub struct WorkspaceScope {
3589 wanted: PathBuf,
3590 roots: Vec<PathBuf>,
3592 subtree: bool,
3593 judged: std::sync::Mutex<HashMap<PathBuf, bool>>,
3595}
3596
3597impl WorkspaceScope {
3598 pub fn exact(wanted: &Path) -> Self {
3600 Self {
3601 wanted: wanted.to_path_buf(),
3602 roots: Vec::new(),
3603 subtree: false,
3604 judged: std::sync::Mutex::new(HashMap::new()),
3605 }
3606 }
3607
3608 pub fn subtree(root: &Path) -> Self {
3615 let mut roots = path_comparison_keys(root);
3616 let own = roots.clone();
3617 for linked in linked_folders(root) {
3618 for key in path_comparison_keys(&linked) {
3619 let widens = own.iter().any(|root| root.starts_with(&key));
3620 let covered = roots.iter().any(|root| key.starts_with(root));
3621 if !widens && !covered {
3622 roots.push(key);
3623 }
3624 }
3625 }
3626 Self {
3627 wanted: root.to_path_buf(),
3628 roots,
3629 subtree: true,
3630 judged: std::sync::Mutex::new(HashMap::new()),
3631 }
3632 }
3633
3634 pub fn of(query: &DiscoveryQuery) -> Option<Self> {
3636 query.workspace.as_deref().map(|wanted| {
3637 if query.workspace_subtree {
3638 Self::subtree(wanted)
3639 } else {
3640 Self::exact(wanted)
3641 }
3642 })
3643 }
3644
3645 pub fn admits(&self, recorded: &Path) -> bool {
3650 if !recorded.is_absolute() {
3651 return false;
3652 }
3653 if !self.subtree {
3654 return same_path(recorded, &self.wanted);
3655 }
3656 let mut judged = self
3657 .judged
3658 .lock()
3659 .unwrap_or_else(std::sync::PoisonError::into_inner);
3660 *judged.entry(recorded.to_path_buf()).or_insert_with(|| {
3661 path_comparison_keys(recorded)
3662 .iter()
3663 .any(|key| self.roots.iter().any(|root| key.starts_with(root)))
3664 })
3665 }
3666
3667 fn admits_literal(&self, recorded: Option<&Path>) -> bool {
3671 if self.subtree {
3672 recorded.is_some_and(|cwd| self.admits(cwd))
3673 } else {
3674 recorded == Some(self.wanted.as_path())
3675 }
3676 }
3677}
3678
3679fn linked_folders(root: &Path) -> Vec<PathBuf> {
3681 let Ok(entries) = fs::read_dir(root) else {
3682 return Vec::new();
3683 };
3684 entries
3685 .filter_map(|entry| {
3686 let entry = entry.ok()?;
3687 let link = entry.file_type().ok()?.is_symlink();
3688 (link && fs::metadata(entry.path()).is_ok_and(|meta| meta.is_dir()))
3689 .then(|| entry.path())
3690 })
3691 .collect()
3692}
3693
3694fn path_comparison_keys(path: &Path) -> Vec<PathBuf> {
3700 let mut keys = Vec::with_capacity(2);
3701 if let Ok(canonical) = fs::canonicalize(path) {
3702 keys.push(comparison_key(&canonical));
3703 }
3704 let lexical = comparison_key(&normalize_path(path));
3705 if !keys.contains(&lexical) {
3706 keys.push(lexical);
3707 }
3708 keys
3709}
3710
3711#[cfg(windows)]
3712fn comparison_key(path: &Path) -> PathBuf {
3713 let text = path.to_string_lossy();
3714 let text = if let Some(unc) = text.strip_prefix(r"\\?\UNC\") {
3715 format!(r"\\{unc}")
3716 } else if let Some(local) = text.strip_prefix(r"\\?\") {
3717 local.to_string()
3718 } else {
3719 text.into_owned()
3720 };
3721 PathBuf::from(text.to_lowercase())
3722}
3723
3724#[cfg(not(windows))]
3725fn comparison_key(path: &Path) -> PathBuf {
3726 path.to_path_buf()
3727}
3728
3729fn recorded_cwd_matches(recorded: &Path, wanted: &Path) -> bool {
3735 recorded.is_absolute() && same_path(recorded, wanted)
3736}
3737
3738fn same_path(left: &Path, right: &Path) -> bool {
3739 match (fs::canonicalize(left), fs::canonicalize(right)) {
3740 (Ok(left), Ok(right)) => left == right,
3741 _ => normalize_path(left) == normalize_path(right),
3742 }
3743}
3744
3745fn normalize_path(path: &Path) -> PathBuf {
3746 let absolute = if path.is_absolute() {
3747 path.to_path_buf()
3748 } else {
3749 std::env::current_dir()
3750 .unwrap_or_else(|_| PathBuf::from("."))
3751 .join(path)
3752 };
3753 let mut normalized = PathBuf::new();
3754 for component in absolute.components() {
3755 match component {
3756 Component::CurDir => {}
3757 Component::ParentDir => {
3758 normalized.pop();
3759 }
3760 other => normalized.push(other.as_os_str()),
3761 }
3762 }
3763 normalized
3764}
3765
3766fn claude_custom_titles(path: &std::path::Path) -> Vec<String> {
3771 use std::io::{Read, Seek, SeekFrom};
3772 use std::sync::{Mutex, OnceLock};
3773 static CACHE: OnceLock<
3774 Mutex<std::collections::HashMap<std::path::PathBuf, (u64, Vec<String>)>>,
3775 > = OnceLock::new();
3776 const KEY: &[u8] = b"\"customTitle\":";
3777 let Ok(len) = std::fs::metadata(path).map(|meta| meta.len()) else {
3778 return Vec::new();
3779 };
3780 let cache = CACHE.get_or_init(Default::default);
3781 let known = cache.lock().ok().and_then(|cache| cache.get(path).cloned());
3782 let (from, mut titles) = match known {
3783 Some((seen, titles)) if seen == len => return titles,
3784 Some((seen, titles)) if seen < len => (seen.saturating_sub(4096), titles),
3786 _ => (0, Vec::new()),
3787 };
3788 let mut bytes = Vec::new();
3789 let read = std::fs::File::open(path).and_then(|mut file| {
3790 file.seek(SeekFrom::Start(from))?;
3791 file.read_to_end(&mut bytes)
3792 });
3793 if read.is_err() {
3794 return titles;
3795 }
3796 let mut at = 0;
3797 while let Some(offset) = bytes[at..]
3798 .windows(KEY.len())
3799 .position(|window| window == KEY)
3800 {
3801 let start = at + offset + KEY.len();
3802 if let Some(Ok(title)) = serde_json::Deserializer::from_slice(&bytes[start..])
3803 .into_iter::<String>()
3804 .next()
3805 {
3806 if !titles.contains(&title) {
3807 titles.push(title);
3808 }
3809 }
3810 at = start;
3811 }
3812 if let Ok(mut cache) = cache.lock() {
3813 cache.insert(path.to_path_buf(), (len, titles.clone()));
3814 }
3815 titles
3816}