1use agent_session::{AgentSession, TokenUsage};
5use serde_json::Value;
6use std::cmp::Reverse;
7use std::collections::{BTreeMap, HashSet};
8use std::fs::File;
9use std::io::{Read, Seek, SeekFrom};
10use std::path::{Path, PathBuf};
11use std::time::{Duration, SystemTime, UNIX_EPOCH};
12
13use std::fs;
14
15use crate::model::{
16 AGENT_NATIVE_SOURCE, AuditEventRow, LlmCallRow, SessionRow, Snapshot, SnapshotOptions,
17 TokenUsageRow, ToolCallRow,
18};
19use crate::text::{sanitize_ascii_identifier as sanitize_id, truncate_text};
20use crate::view::MaterializedView;
21
22pub type LocalSession = AgentSession;
23pub type SessionCache = agent_session::SessionCache;
24const CODEX_EXEC_DEDUPE_WINDOW_MS: u64 = 2_000;
25const CODEX_FALLBACK_TIME_SLOP_MS: u64 = 30_000;
26const CODEX_ROLLOUT_TAIL_BYTES: u64 = 1024 * 1024;
27const CURSOR_AGENT_TYPE: &str = "cursor";
31
32#[derive(Clone, Debug)]
33struct ObservedCodexPrompt {
34 prompt: String,
35 timestamp_ms: u64,
36 pid: Option<u32>,
37 native_exec: bool,
38 comm: Option<String>,
39 target: Option<String>,
40}
41
42pub fn snapshot(
43 cache: &mut SessionCache,
44 pid_filter: Option<u32>,
45 text_filter: Option<&str>,
46 limit: usize,
47 max_age: Duration,
48) -> Snapshot {
49 let filtered = discover_sessions(cache, pid_filter, text_filter, limit, max_age);
50 materialized_view(&filtered).export_snapshot(SnapshotOptions { audit_limit: 0 })
51}
52
53pub fn discover_sessions(
54 cache: &mut SessionCache,
55 pid_filter: Option<u32>,
56 text_filter: Option<&str>,
57 limit: usize,
58 max_age: Duration,
59) -> Vec<LocalSession> {
60 let indexed_codex = codex_state_sessions(limit);
61 let mut sessions = if indexed_codex.is_empty() {
62 cache.discover_cached(limit, max_age)
63 } else {
64 let mut sessions = indexed_codex;
65 sessions.extend(cache.discover_cached_excluding(
66 limit,
67 max_age,
68 &[agent_session::AGENT_CODEX],
69 ));
70 sessions.sort_by_key(|session| Reverse(session.updated));
71 sessions.truncate(limit.clamp(1, 25));
72 sessions
73 };
74 let mut seen = HashSet::new();
75 sessions.retain(|session| seen.insert(session.display_id.clone()));
76 enrich_cursor_sessions(&mut sessions);
77 sessions
78 .into_iter()
79 .filter(|s| matches_filter(s, pid_filter, text_filter))
80 .collect()
81}
82
83fn codex_state_sessions(limit: usize) -> Vec<LocalSession> {
84 user_home_dir()
85 .as_deref()
86 .map(|home| codex_state_sessions_in_home(home, limit))
87 .unwrap_or_default()
88}
89
90fn codex_state_sessions_in_home(home: &Path, limit: usize) -> Vec<LocalSession> {
91 let db_path = home.join(".codex/state_5.sqlite");
92 let Ok(conn) =
93 rusqlite::Connection::open_with_flags(db_path, rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY)
94 else {
95 return Vec::new();
96 };
97 let Ok(mut stmt) = conn.prepare(
98 "SELECT id, rollout_path, model, tokens_used, preview, cwd, created_at_ms, updated_at_ms
99 FROM threads
100 ORDER BY updated_at_ms DESC
101 LIMIT ?1",
102 ) else {
103 return Vec::new();
104 };
105 let Ok(rows) = stmt.query_map([limit.clamp(1, 25) as i64], |row| {
106 let id: String = row.get(0)?;
107 let rollout_path: String = row.get(1)?;
108 let model: Option<String> = row.get(2)?;
109 let tokens_used: i64 = row.get(3)?;
110 let preview: Option<String> = row.get(4)?;
111 let cwd: Option<String> = row.get(5)?;
112 let created_at_ms: Option<i64> = row.get(6)?;
113 let updated_at_ms: Option<i64> = row.get(7)?;
114 Ok(codex_state_session(
115 id,
116 rollout_path,
117 model,
118 tokens_used,
119 preview,
120 cwd,
121 created_at_ms,
122 updated_at_ms,
123 ))
124 }) else {
125 return Vec::new();
126 };
127
128 rows.filter_map(Result::ok).collect()
129}
130
131fn codex_state_session(
132 id: String,
133 rollout_path: String,
134 model: Option<String>,
135 tokens_used: i64,
136 preview: Option<String>,
137 cwd: Option<String>,
138 created_at_ms: Option<i64>,
139 updated_at_ms: Option<i64>,
140) -> LocalSession {
141 let updated_ms = updated_at_ms.and_then(non_negative_i64_to_u64);
142 let created_ms = created_at_ms
143 .and_then(non_negative_i64_to_u64)
144 .or(updated_ms);
145 let updated = updated_ms.map(system_time_from_ms).unwrap_or(UNIX_EPOCH);
146 let path = PathBuf::from(rollout_path);
147 let usage = codex_rollout_usage(&path).unwrap_or(TokenUsage {
148 total_tokens: tokens_used.max(0),
149 ..Default::default()
150 });
151 let model = model.filter(|value| !value.is_empty());
152 let mut model_usage = BTreeMap::new();
153 if let Some(model) = model.as_deref() {
154 model_usage.insert(model.to_string(), usage.clone());
155 }
156 let prompt_preview = preview
157 .and_then(|text| clean_prompt_text(&text))
158 .map(|text| truncate_text(&text, 180));
159 let last_message_at = updated_ms.map(iso_utc_from_ms);
160
161 LocalSession {
162 agent_type: agent_session::AGENT_CODEX.to_string(),
163 session_id: id.clone(),
164 conversation_id: Some(id.clone()),
165 display_id: format!("{}:{}", agent_session::AGENT_CODEX, short_session_id(&id)),
166 path,
167 updated,
168 start_timestamp_ms: created_ms,
169 end_timestamp_ms: updated_ms,
170 model,
171 usage,
172 model_usage,
173 tools: BTreeMap::new(),
174 files: BTreeMap::new(),
175 prompt_preview,
176 duration_ms: created_ms
177 .zip(updated_ms)
178 .map(|(start, end)| end.saturating_sub(start))
179 .unwrap_or_default(),
180 cwd,
181 last_message_at,
182 events: Default::default(),
183 }
184}
185
186fn codex_rollout_usage(path: &Path) -> Option<TokenUsage> {
187 let mut file = File::open(path).ok()?;
188 let len = file.metadata().ok()?.len();
189 let window = len.min(CODEX_ROLLOUT_TAIL_BYTES);
190 file.seek(SeekFrom::Start(len - window)).ok()?;
191 let mut data = Vec::with_capacity(window as usize);
192 file.read_to_end(&mut data).ok()?;
193 agent_session::codex_total_token_usage(&String::from_utf8_lossy(&data))
194}
195
196const CURSOR_STATE_DB_CANDIDATES: [&str; 3] = [
197 "Library/Application Support/Cursor/User/globalStorage/state.vscdb",
198 ".config/Cursor/User/globalStorage/state.vscdb",
199 "AppData/Roaming/Cursor/User/globalStorage/state.vscdb",
200];
201
202fn cursor_state_db_path(home: &Path) -> Option<PathBuf> {
203 CURSOR_STATE_DB_CANDIDATES
204 .iter()
205 .map(|candidate| home.join(candidate))
206 .find(|path| path.is_file())
207}
208
209fn open_cursor_state_db(path: &Path) -> Option<rusqlite::Connection> {
210 rusqlite::Connection::open_with_flags(path, rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY).ok()
211}
212
213struct CursorComposerHeader {
214 created_at_ms: Option<u64>,
215 updated_at_ms: Option<u64>,
216}
217
218fn cursor_composer_header(
219 conn: &rusqlite::Connection,
220 composer_id: &str,
221) -> Option<CursorComposerHeader> {
222 conn.query_row(
223 "SELECT createdAt, lastUpdatedAt FROM composerHeaders WHERE composerId = ?1",
224 [composer_id],
225 |row| {
226 let created: Option<i64> = row.get(0)?;
227 let updated: Option<i64> = row.get(1)?;
228 Ok(CursorComposerHeader {
229 created_at_ms: created.and_then(non_negative_i64_to_u64),
230 updated_at_ms: updated.and_then(non_negative_i64_to_u64),
231 })
232 },
233 )
234 .ok()
235}
236
237fn cursor_kv_bytes(value: rusqlite::types::Value) -> Option<Vec<u8>> {
238 match value {
239 rusqlite::types::Value::Blob(bytes) => Some(bytes),
240 rusqlite::types::Value::Text(text) => Some(text.into_bytes()),
241 _ => None,
242 }
243}
244
245struct CursorComposerData {
246 model: Option<String>,
247 workspace_path: Option<String>,
248}
249
250fn cursor_composer_data(
251 conn: &rusqlite::Connection,
252 composer_id: &str,
253) -> Option<CursorComposerData> {
254 let raw: rusqlite::types::Value = conn
255 .query_row(
256 "SELECT value FROM cursorDiskKV WHERE key = ?1",
257 [format!("composerData:{composer_id}")],
258 |row| row.get(0),
259 )
260 .ok()?;
261 let value: Value = serde_json::from_slice(&cursor_kv_bytes(raw)?).ok()?;
262 let model = value
263 .pointer("/modelConfig/modelName")
264 .and_then(Value::as_str)
265 .filter(|name| !name.is_empty() && *name != "default")
266 .map(str::to_string);
267 let workspace_path = value
268 .pointer("/workspaceIdentifier/uri/fsPath")
269 .and_then(Value::as_str)
270 .filter(|path| !path.is_empty())
271 .map(str::to_string);
272 Some(CursorComposerData {
273 model,
274 workspace_path,
275 })
276}
277
278fn cursor_bubble_tokens(conn: &rusqlite::Connection, composer_id: &str) -> TokenUsage {
279 let mut usage = TokenUsage::default();
280 let Ok(mut stmt) = conn.prepare("SELECT value FROM cursorDiskKV WHERE key >= ?1 AND key < ?2")
281 else {
282 return usage;
283 };
284 let lower = format!("bubbleId:{composer_id}:");
285 let upper = format!("bubbleId:{composer_id};");
286 let Ok(rows) = stmt.query_map([lower, upper], |row| {
287 row.get::<_, rusqlite::types::Value>(0)
288 }) else {
289 return usage;
290 };
291 for raw in rows.filter_map(Result::ok).filter_map(cursor_kv_bytes) {
292 let Ok(value) = serde_json::from_slice::<Value>(&raw) else {
293 continue;
294 };
295 let input = value
296 .pointer("/tokenCount/inputTokens")
297 .and_then(Value::as_i64)
298 .unwrap_or_default()
299 .max(0);
300 let output = value
301 .pointer("/tokenCount/outputTokens")
302 .and_then(Value::as_i64)
303 .unwrap_or_default()
304 .max(0);
305 if input + output > 0 {
306 usage.input_tokens += input;
307 usage.output_tokens += output;
308 usage.total_tokens += input + output;
309 }
310 }
311 usage
312}
313
314fn cursor_subagent_ids(parent_transcript: &Path) -> Vec<String> {
315 let Some(subagents) = parent_transcript.parent().map(|dir| dir.join("subagents")) else {
316 return Vec::new();
317 };
318 let Ok(entries) = fs::read_dir(subagents) else {
319 return Vec::new();
320 };
321 let mut ids: Vec<String> = entries
322 .flatten()
323 .map(|entry| entry.path())
324 .filter(|path| path.extension().and_then(|ext| ext.to_str()) == Some("jsonl"))
325 .filter_map(|path| {
326 path.file_stem()
327 .and_then(|stem| stem.to_str())
328 .map(str::to_string)
329 })
330 .collect();
331 ids.sort();
332 ids
333}
334
335fn enrich_cursor_sessions(sessions: &mut [LocalSession]) {
336 if !sessions
337 .iter()
338 .any(|session| session.agent_type == CURSOR_AGENT_TYPE)
339 {
340 return;
341 }
342 let Some(home) = user_home_dir() else {
343 return;
344 };
345 enrich_cursor_sessions_in_home(&home, sessions);
346}
347
348fn enrich_cursor_sessions_in_home(home: &Path, sessions: &mut [LocalSession]) {
349 let Some(db_path) = cursor_state_db_path(home) else {
350 return;
351 };
352 let Some(conn) = open_cursor_state_db(&db_path) else {
353 return;
354 };
355 for session in sessions
356 .iter_mut()
357 .filter(|session| session.agent_type == CURSOR_AGENT_TYPE)
358 {
359 enrich_cursor_session(&conn, session);
360 }
361}
362
363fn enrich_cursor_session(conn: &rusqlite::Connection, session: &mut LocalSession) {
364 let composer_id = session.session_id.clone();
365 if let Some(header) = cursor_composer_header(conn, &composer_id) {
366 if header.created_at_ms.is_some() {
367 session.start_timestamp_ms = header.created_at_ms;
368 }
369 if let Some(updated_ms) = header.updated_at_ms {
370 session.end_timestamp_ms = Some(updated_ms);
371 session.last_message_at = Some(iso_utc_from_ms(updated_ms));
372 }
373 if let (Some(start), Some(end)) = (session.start_timestamp_ms, session.end_timestamp_ms) {
374 session.duration_ms = end.saturating_sub(start);
375 }
376 }
377 if let Some(data) = cursor_composer_data(conn, &composer_id) {
378 if data.model.is_some() {
379 session.model = data.model;
380 }
381 if data.workspace_path.is_some() {
382 session.cwd = data.workspace_path;
383 }
384 }
385 let mut usage = cursor_bubble_tokens(conn, &composer_id);
387 for child_id in cursor_subagent_ids(&session.path) {
388 let child = cursor_bubble_tokens(conn, &child_id);
389 usage.input_tokens += child.input_tokens;
390 usage.output_tokens += child.output_tokens;
391 usage.total_tokens += child.total_tokens;
392 }
393 if usage.total_tokens > 0 {
394 if let Some(model) = session.model.as_deref() {
395 session.model_usage.insert(model.to_string(), usage.clone());
396 }
397 session.usage = usage;
398 }
399}
400
401fn user_home_dir() -> Option<PathBuf> {
402 std::env::var("SUDO_USER")
403 .ok()
404 .and_then(|user| {
405 std::fs::read_to_string("/etc/passwd")
406 .ok()
407 .and_then(|passwd| {
408 passwd
409 .lines()
410 .find(|line| line.starts_with(&format!("{user}:")))
411 .and_then(|line| line.split(':').nth(5))
412 .map(PathBuf::from)
413 })
414 })
415 .or_else(|| {
416 std::env::var_os("HOME")
417 .map(PathBuf::from)
418 .filter(|home| home.is_absolute())
419 })
420 .or_else(dirs::home_dir)
421}
422
423fn non_negative_i64_to_u64(value: i64) -> Option<u64> {
424 u64::try_from(value).ok()
425}
426
427fn system_time_from_ms(value: u64) -> SystemTime {
428 UNIX_EPOCH + Duration::from_millis(value)
429}
430
431fn iso_utc_from_ms(value: u64) -> String {
432 chrono::DateTime::<chrono::Utc>::from_timestamp_millis(value as i64)
433 .map(|dt| dt.to_rfc3339_opts(chrono::SecondsFormat::Millis, true))
434 .unwrap_or_default()
435}
436
437fn short_session_id(id: &str) -> String {
438 let compact = id.rsplit(['/', '\\']).next().unwrap_or(id).trim();
439 if compact.chars().count() <= 12 {
440 return compact.to_string();
441 }
442 let head = compact.chars().take(6).collect::<String>();
443 let tail = compact
444 .chars()
445 .rev()
446 .take(5)
447 .collect::<Vec<_>>()
448 .into_iter()
449 .rev()
450 .collect::<String>();
451 format!("{head}.{tail}")
452}
453
454fn clean_prompt_text(text: &str) -> Option<String> {
455 let text = text.split_whitespace().collect::<Vec<_>>().join(" ");
456 (!text.trim().is_empty()).then(|| text.trim().to_string())
457}
458
459fn view_id(session: &LocalSession) -> String {
460 format!("local:{}:{}", session.agent_type, session.display_id)
461}
462
463pub fn materialized_view(sessions: &[LocalSession]) -> MaterializedView {
464 let mut view = MaterializedView::new();
465 view.set_source(AGENT_NATIVE_SOURCE);
466 import_into_view(&mut view, sessions);
467 view
468}
469
470pub fn import_recent(view: &mut MaterializedView, limit: usize) {
471 let mut sessions = SessionCache::new().discover_cached(limit, Duration::ZERO);
472 enrich_cursor_sessions(&mut sessions);
476 import_into_view(view, &sessions);
477}
478
479pub fn import_into_view(view: &mut MaterializedView, sessions: &[LocalSession]) {
480 for session in sessions {
481 view.upsert_session(&session_row(session));
482 for row in llm_rows(session) {
483 view.apply_llm_call(&row);
484 }
485 for row in token_rows(session) {
486 view.apply_token_usage(&row);
487 }
488 for row in tool_rows(session) {
489 view.apply_tool_call(&row);
490 }
491 }
492}
493
494fn llm_rows(session: &LocalSession) -> Vec<LlmCallRow> {
495 let Some(prompt) = session.prompt_preview.as_ref() else {
496 return Vec::new();
497 };
498 let session_id = view_id(session);
499 let timestamp_ms = session
500 .events
501 .prompts
502 .first()
503 .and_then(|prompt| prompt.ts_ms)
504 .and_then(|ts| u64::try_from(ts).ok())
505 .or(session.start_timestamp_ms)
506 .unwrap_or_else(|| updated_ms(session));
507 let request = serde_json::json!({
508 "prompt": prompt,
509 "prompt_source": AGENT_NATIVE_SOURCE,
510 "session_id": session_id,
511 "agent_type": session.agent_type.as_str(),
512 "path": session.path.to_string_lossy(),
513 });
514
515 if session.model_usage.is_empty() {
516 let model = session
517 .model
518 .clone()
519 .unwrap_or_else(|| session.agent_type.clone());
520 return vec![llm_row_for_session(
521 &format!("{session_id}-{}", sanitize_id(&model)),
522 session,
523 &session_id,
524 timestamp_ms,
525 Some(model),
526 &session.usage,
527 request,
528 )];
529 }
530
531 session
532 .model_usage
533 .iter()
534 .map(|(model, usage)| {
535 llm_row_for_session(
536 &format!("{session_id}-{model}"),
537 session,
538 &session_id,
539 timestamp_ms,
540 Some(model.clone()),
541 usage,
542 request.clone(),
543 )
544 })
545 .collect()
546}
547
548fn llm_row_for_session(
549 id: &str,
550 session: &LocalSession,
551 session_id: &str,
552 timestamp_ms: u64,
553 model: Option<String>,
554 usage: &TokenUsage,
555 request: Value,
556) -> LlmCallRow {
557 LlmCallRow {
558 id: id.to_string(),
559 session_id: Some(session_id.to_string()),
560 conversation_id: session.conversation_id.clone(),
561 start_timestamp_ms: timestamp_ms,
562 end_timestamp_ms: session.end_timestamp_ms,
563 pid: None,
564 comm: Some(session.agent_type.clone()),
565 provider: None,
566 model,
567 call_kind: Some("agent_native_prompt".to_string()),
568 status: "observed".to_string(),
569 error_type: None,
570 finish_reason: None,
571 host: None,
572 path: Some(session.path.to_string_lossy().to_string()),
573 status_code: None,
574 input_tokens: usage.input_tokens,
575 output_tokens: usage.output_tokens,
576 total_tokens: usage.total_tokens,
577 request,
578 response: Value::Null,
579 }
580}
581
582pub fn observed_session_prompt_rows(audit_rows: &[AuditEventRow]) -> Vec<AuditEventRow> {
583 let mut rows = Vec::new();
584 let mut seen = HashSet::new();
585 let observed_exec_prompts = observed_codex_exec_prompts(audit_rows);
586 let mut seen_exec_prompts: Vec<ObservedCodexPrompt> = Vec::new();
587 for observed in observed_exec_prompts {
588 if seen_exec_prompts.iter().any(|seen| {
589 seen.prompt == observed.prompt
590 && timestamps_close(
591 seen.timestamp_ms,
592 observed.timestamp_ms,
593 CODEX_EXEC_DEDUPE_WINDOW_MS,
594 )
595 && (!seen.native_exec || !observed.native_exec)
596 }) {
597 continue;
598 }
599 seen_exec_prompts.push(observed.clone());
600 rows.push(AuditEventRow {
601 id: format!(
602 "audit-codex-exec-prompt-{}-{}",
603 observed.timestamp_ms,
604 observed.pid.unwrap_or(0)
605 ),
606 timestamp_ms: observed.timestamp_ms,
607 audit_type: "llm".to_string(),
608 pid: observed.pid,
609 comm: observed.comm.or_else(|| Some("codex".to_string())),
610 subject: None,
611 action: Some("request".to_string()),
612 target: observed.target,
613 status: Some("observed".to_string()),
614 summary: Some(truncate_text(&observed.prompt, 160)),
615 details: serde_json::json!({
616 "text_content": observed.prompt,
617 "prompt_source": "local",
618 }),
619 });
620 }
621 for row in audit_rows {
622 if row.audit_type == "process" && row.action.as_deref() == Some("exec") {
623 continue;
624 }
625 if row.audit_type != "file" {
626 continue;
627 }
628 let Some(pid) = row.pid else {
629 continue;
630 };
631 let Some(path) = audit_session_path(row) else {
632 continue;
633 };
634 if !seen.insert((path.clone(), pid)) {
635 continue;
636 };
637 let Some(session) = agent_session::parse_session_path(&path) else {
638 continue;
639 };
640 let Some(prompt) = session.prompt_preview.as_ref() else {
641 continue;
642 };
643 rows.push(AuditEventRow {
644 id: format!(
645 "audit-agent-native-prompt-{}-{pid}",
646 sanitize_id(&session.display_id)
647 ),
648 timestamp_ms: row.timestamp_ms,
649 audit_type: "llm".to_string(),
650 pid: Some(pid),
651 comm: row
652 .comm
653 .clone()
654 .or_else(|| Some(session.agent_type.clone())),
655 subject: session.model.clone(),
656 action: Some("request".to_string()),
657 target: Some(path.to_string_lossy().to_string()),
658 status: Some("observed".to_string()),
659 summary: Some(truncate_text(prompt, 160)),
660 details: serde_json::json!({
661 "text_content": prompt,
662 "prompt_source": "local",
663 "session_id": view_id(&session),
664 "conversation_id": session.conversation_id.as_deref(),
665 "agent_type": session.agent_type,
666 }),
667 });
668 }
669 rows
670}
671
672pub fn observed_sessions_from_audit_rows(audit_rows: &[AuditEventRow]) -> Vec<LocalSession> {
673 let mut direct_paths = HashSet::new();
674 let mut codex_session_dirs = HashSet::new();
675 let observed_codex_prompts = observed_codex_exec_prompts(audit_rows);
676 let observed_codex_exec = observed_codex_exec_command(audit_rows);
677 let observed_window = observed_audit_window_ms(audit_rows);
678
679 for row in audit_rows.iter().filter(|row| row.audit_type == "file") {
680 for path in audit_file_paths(row) {
681 if let Some(session_path) =
682 agent_session::session_log_path_from_str(path.to_string_lossy().as_ref())
683 {
684 direct_paths.insert(session_path);
685 }
686 if let Some(dir) = observed_codex_sessions_dir(&path) {
687 codex_session_dirs.insert(dir);
688 }
689 }
690 }
691
692 let mut candidates = Vec::new();
693 for path in direct_paths {
694 if let Some(candidate) = agent_session::session_candidate_from_path(&path) {
695 candidates.push((candidate, false));
696 }
697 }
698 for dir in codex_session_dirs {
699 let dir_candidates =
700 agent_session::discover_session_files_in_dir(agent_session::AGENT_CODEX, &dir);
701 candidates.extend(
702 dir_candidates
703 .into_iter()
704 .map(|candidate| (candidate, true)),
705 );
706 }
707 candidates.sort_by_key(|(candidate, _)| Reverse(candidate.updated));
708
709 let mut seen_paths = HashSet::new();
710 let mut seen_sessions = HashSet::new();
711 let mut sessions = Vec::new();
712 for (candidate, is_codex_dir_fallback) in candidates.into_iter().take(75) {
713 if !seen_paths.insert(candidate.path.clone()) {
714 continue;
715 }
716 let Some(session) = agent_session::parse_session_file(&candidate) else {
717 continue;
718 };
719 if is_codex_dir_fallback && !observed_codex_exec {
720 continue;
721 }
722 if is_codex_dir_fallback
723 && !observed_codex_prompts.is_empty()
724 && !session_matches_observed_prompt(&session, &observed_codex_prompts)
725 {
726 continue;
727 }
728 if is_codex_dir_fallback && !session_is_in_observed_window(&session, observed_window) {
729 continue;
730 }
731 if seen_sessions.insert(session.display_id.clone()) {
732 sessions.push(session);
733 }
734 }
735 sessions
736}
737
738fn observed_codex_exec_prompts(audit_rows: &[AuditEventRow]) -> Vec<ObservedCodexPrompt> {
739 let prompts = audit_rows
740 .iter()
741 .filter(|row| row.audit_type == "process" && row.action.as_deref() == Some("exec"))
742 .filter_map(|row| {
743 let prompt = row
744 .details
745 .get("full_command")
746 .and_then(Value::as_str)
747 .and_then(codex_exec_prompt_from_command)?;
748 Some(ObservedCodexPrompt {
749 prompt,
750 timestamp_ms: row.timestamp_ms,
751 pid: row.pid,
752 native_exec: looks_like_native_codex_exec(row),
753 comm: row.comm.clone(),
754 target: row.target.clone(),
755 })
756 })
757 .collect::<Vec<_>>();
758 prompts
759 .iter()
760 .filter(|candidate| {
761 !prompts
762 .iter()
763 .any(|other| is_nearby_longer_prefix(candidate, other))
764 })
765 .cloned()
766 .collect()
767}
768
769fn observed_codex_exec_command(audit_rows: &[AuditEventRow]) -> bool {
770 audit_rows
771 .iter()
772 .filter(|row| row.audit_type == "process" && row.action.as_deref() == Some("exec"))
773 .filter_map(|row| row.details.get("full_command").and_then(Value::as_str))
774 .any(|command| codex_exec_command_tail(command).is_some())
775}
776
777fn codex_exec_prompt_from_command(command: &str) -> Option<String> {
778 agent_session::codex_exec_prompt(&codex_exec_command_tail(command)?)
779}
780
781fn codex_exec_command_tail(command: &str) -> Option<String> {
782 let tokens = command.split_whitespace().collect::<Vec<_>>();
783 let index = tokens.windows(2).enumerate().find_map(|(index, tokens)| {
784 (is_codex_executable_token(tokens[0], index == 0) && tokens[1] == "exec").then_some(index)
785 })?;
786 Some(tokens[index..].join(" "))
787}
788
789fn looks_like_native_codex_exec(row: &AuditEventRow) -> bool {
790 row.comm.as_deref() == Some("codex")
791 && row
792 .target
793 .as_deref()
794 .is_some_and(|target| is_codex_executable_token(target, true))
795}
796
797fn is_codex_executable_token(token: &str, allow_bare: bool) -> bool {
798 let token = token.trim_matches(|ch| matches!(ch, '"' | '\''));
799 (allow_bare && token == "codex")
800 || token.contains('/')
801 && Path::new(token).file_name().and_then(|name| name.to_str()) == Some("codex")
802}
803
804fn is_nearby_longer_prefix(candidate: &ObservedCodexPrompt, other: &ObservedCodexPrompt) -> bool {
805 other.prompt.len() > candidate.prompt.len()
806 && other.prompt.starts_with(candidate.prompt.as_str())
807 && timestamps_close(
808 candidate.timestamp_ms,
809 other.timestamp_ms,
810 CODEX_EXEC_DEDUPE_WINDOW_MS,
811 )
812 && (!candidate.native_exec || !other.native_exec)
813}
814
815fn timestamps_close(left: u64, right: u64, window_ms: u64) -> bool {
816 left.abs_diff(right) <= window_ms
817}
818
819fn session_matches_observed_prompt(
820 session: &LocalSession,
821 prompts: &[ObservedCodexPrompt],
822) -> bool {
823 let Some(preview) = session.prompt_preview.as_deref() else {
824 return false;
825 };
826 prompts.iter().any(|observed| {
827 let prompt = observed.prompt.as_str();
828 prompt_texts_overlap(prompt, preview)
829 })
830}
831
832fn prompt_texts_overlap(left: &str, right: &str) -> bool {
833 left == right || left.starts_with(right) || right.starts_with(left)
834}
835
836fn observed_audit_window_ms(audit_rows: &[AuditEventRow]) -> Option<(u64, u64)> {
837 let min = audit_rows.iter().map(|row| row.timestamp_ms).min()?;
838 let max = audit_rows.iter().map(|row| row.timestamp_ms).max()?;
839 Some((
840 min.saturating_sub(CODEX_FALLBACK_TIME_SLOP_MS),
841 max.saturating_add(CODEX_FALLBACK_TIME_SLOP_MS),
842 ))
843}
844
845fn session_is_in_observed_window(session: &LocalSession, window: Option<(u64, u64)>) -> bool {
846 let Some((min, max)) = window else {
847 return true;
848 };
849 let updated = updated_ms(session);
850 updated >= min && updated <= max
851}
852
853fn audit_session_path(row: &AuditEventRow) -> Option<PathBuf> {
854 row.target
855 .as_deref()
856 .and_then(agent_session::session_log_path_from_str)
857 .or_else(|| {
858 row.details
859 .get("filepath")
860 .and_then(Value::as_str)
861 .and_then(agent_session::session_log_path_from_str)
862 })
863 .or_else(|| {
864 row.details
865 .get("path")
866 .and_then(Value::as_str)
867 .and_then(agent_session::session_log_path_from_str)
868 })
869}
870
871fn audit_file_paths(row: &AuditEventRow) -> Vec<PathBuf> {
872 [
873 row.target.as_deref(),
874 row.details.get("filepath").and_then(Value::as_str),
875 row.details.get("path").and_then(Value::as_str),
876 row.details.get("fd_target").and_then(Value::as_str),
877 ]
878 .into_iter()
879 .flatten()
880 .filter_map(|raw| {
881 let path = PathBuf::from(raw.trim().trim_end_matches(" (deleted)"));
882 path.is_absolute().then_some(path)
883 })
884 .collect()
885}
886
887fn observed_codex_sessions_dir(path: &Path) -> Option<PathBuf> {
888 if !looks_like_codex_home_file(path) {
889 return None;
890 }
891 let home = path.parent()?;
892 let sessions = home.join("sessions");
893 sessions.is_dir().then_some(sessions)
894}
895
896fn looks_like_codex_home_file(path: &Path) -> bool {
897 let Some(name) = path.file_name().and_then(|name| name.to_str()) else {
898 return false;
899 };
900 (name.starts_with("state_")
901 || name.starts_with("logs_")
902 || matches!(name, "config.toml" | "auth.json" | "stat"))
903 && path
904 .parent()
905 .is_some_and(|parent| parent.join("sessions").is_dir())
906}
907
908fn session_row(session: &LocalSession) -> SessionRow {
909 let updated_ms = updated_ms(session);
910 SessionRow {
911 id: view_id(session),
912 agent_type: session.agent_type.clone(),
913 start_timestamp_ms: session
914 .start_timestamp_ms
915 .unwrap_or_else(|| updated_ms.saturating_sub(session.duration_ms)),
916 end_timestamp_ms: session.end_timestamp_ms.or(Some(updated_ms)),
917 status: "observed".to_string(),
918 model: session.model.clone(),
919 input_tokens: session.usage.input_tokens,
920 output_tokens: session.usage.output_tokens,
921 total_tokens: session.usage.total_tokens,
922 view_source: AGENT_NATIVE_SOURCE.to_string(),
923 confidence: Some(0.95),
924 attributes: serde_json::json!({
925 "session_id": session.session_id.clone(),
926 "conversation_id": session.conversation_id.as_deref(),
927 "path": session.path.to_string_lossy(),
928 "display_id": session.display_id,
929 "prompt_preview": session.prompt_preview.clone(),
930 "cwd": session.cwd.clone(),
931 "last_message_at": session.last_message_at.clone(),
932 "files": session.files,
933 }),
934 }
935}
936
937fn token_rows(session: &LocalSession) -> Vec<TokenUsageRow> {
938 let session_id = view_id(session);
939 session
940 .model_usage
941 .iter()
942 .filter(|(_, usage)| usage.total_tokens > 0)
943 .map(|(model, usage)| TokenUsageRow {
944 id: format!("token-{session_id}-{}", sanitize_id(model)),
945 llm_call_id: format!("{session_id}-{model}"),
946 timestamp_ms: updated_ms(session),
947 pid: None,
948 comm: Some(session.agent_type.clone()),
949 provider: None,
950 model: Some(model.clone()),
951 input_tokens: usage.input_tokens,
952 output_tokens: usage.output_tokens,
953 cache_creation_tokens: usage.cache_creation_tokens,
954 cache_read_tokens: usage.cache_read_tokens,
955 total_tokens: usage.total_tokens,
956 source: AGENT_NATIVE_SOURCE.to_string(),
957 view_source: AGENT_NATIVE_SOURCE.to_string(),
958 confidence: Some(0.95),
959 })
960 .collect()
961}
962
963fn tool_rows(session: &LocalSession) -> Vec<ToolCallRow> {
964 let session_id = view_id(session);
965 let timestamp_ms = updated_ms(session);
966 let mut rows = Vec::new();
967 for (tool, count) in &session.tools {
968 for index in 0..*count {
969 rows.push(ToolCallRow {
970 id: format!("tool-{session_id}-{}-{index}", sanitize_id(tool)),
971 session_id: Some(session_id.clone()),
972 conversation_id: session.conversation_id.clone(),
973 timestamp_ms,
974 tool_name: Some(tool.clone()),
975 tool_call_id: None,
976 start_timestamp_ms: Some(timestamp_ms),
977 end_timestamp_ms: Some(timestamp_ms),
978 duration_ms: None,
979 status: Some("observed".to_string()),
980 input: serde_json::json!({}),
981 output: serde_json::json!({}),
982 related_pid: None,
983 related_event_id: None,
984 view_source: AGENT_NATIVE_SOURCE.to_string(),
985 confidence: Some(0.95),
986 });
987 }
988 }
989 rows
990}
991
992fn updated_ms(session: &LocalSession) -> u64 {
993 session
994 .updated
995 .duration_since(UNIX_EPOCH)
996 .unwrap_or_default()
997 .as_millis() as u64
998}
999
1000fn matches_filter(
1001 session: &LocalSession,
1002 pid_filter: Option<u32>,
1003 text_filter: Option<&str>,
1004) -> bool {
1005 if pid_filter.is_some() {
1006 return true;
1007 }
1008 let Some(filter) = text_filter else {
1009 return true;
1010 };
1011 let filter = filter.to_ascii_lowercase();
1012 session.agent_type.to_ascii_lowercase().contains(&filter)
1013 || session
1014 .prompt_preview
1015 .as_ref()
1016 .is_some_and(|prompt| prompt.to_ascii_lowercase().contains(&filter))
1017 || session
1018 .model
1019 .as_ref()
1020 .is_some_and(|model| model.to_ascii_lowercase().contains(&filter))
1021 || session
1022 .path
1023 .to_string_lossy()
1024 .to_ascii_lowercase()
1025 .contains(&filter)
1026}
1027
1028#[cfg(any(test, feature = "test-support"))]
1029pub fn create_temp_session_path(agent: &str) -> (tempfile::TempDir, PathBuf) {
1030 let temp = tempfile::tempdir().unwrap();
1031 let path = agent_session::fixture_session_path(agent, temp.path()).unwrap();
1032 fs::create_dir_all(path.parent().unwrap()).unwrap();
1033 fs::write(&path, "{}\n").unwrap();
1034 (temp, path)
1035}
1036
1037#[cfg(any(test, feature = "test-support"))]
1038pub fn parse_content_for_test(
1039 agent: &str,
1040 path: &std::path::Path,
1041 updated: std::time::SystemTime,
1042 content: &str,
1043) -> Option<LocalSession> {
1044 agent_session::parse_session_content(agent, path, updated, content)
1045}
1046
1047#[cfg(any(test, feature = "test-support"))]
1051pub fn write_cursor_state_db_for_test(home: &Path) {
1052 let db_dir = home.join("Library/Application Support/Cursor/User/globalStorage");
1053 fs::create_dir_all(&db_dir).unwrap();
1054 let conn = rusqlite::Connection::open(db_dir.join("state.vscdb")).unwrap();
1055 conn.pragma_update(None, "journal_mode", "WAL").unwrap();
1058 conn.execute_batch(
1059 r#"CREATE TABLE composerHeaders (composerId TEXT PRIMARY KEY, workspaceId TEXT,
1060 createdAt INTEGER, lastUpdatedAt INTEGER, isArchived INTEGER,
1061 isSubagent INTEGER, recency INTEGER, checkpointAt INTEGER, value TEXT);
1062 CREATE TABLE cursorDiskKV (key TEXT UNIQUE ON CONFLICT REPLACE, value BLOB);
1063 INSERT INTO composerHeaders (composerId, workspaceId, createdAt, lastUpdatedAt, isSubagent)
1064 VALUES
1065 ('abc00000-0000-0000-0000-000000000abc', 'ws1', 1700000, 1900000, 0),
1066 ('def00000-0000-0000-0000-000000000def', 'ws1', 1750000, NULL, 1),
1067 ('aaa00000-0000-0000-0000-000000000aaa', 'ws2', 1600000, 1650000, 0);
1068 INSERT INTO cursorDiskKV (key, value) VALUES
1069 ('composerData:bbb00000-0000-0000-0000-000000000bbb',
1070 '{"composerId":"bbb00000-0000-0000-0000-000000000bbb","modelConfig":{"modelName":"claude-4.6-sonnet-medium-thinking","maxMode":false}}'),
1071 ('composerData:abc00000-0000-0000-0000-000000000abc',
1072 '{"composerId":"abc00000-0000-0000-0000-000000000abc","modelConfig":{"modelName":"claude-sonnet-4-6","maxMode":false},"workspaceIdentifier":{"uri":{"fsPath":"/work/repo"}}}'),
1073 ('composerData:aaa00000-0000-0000-0000-000000000aaa',
1074 '{"composerId":"aaa00000-0000-0000-0000-000000000aaa","modelConfig":{"modelName":"default","maxMode":false}}'),
1075 ('bubbleId:abc00000-0000-0000-0000-000000000abc:b1',
1076 '{"tokenCount":{"inputTokens":100,"outputTokens":40}}'),
1077 ('bubbleId:abc00000-0000-0000-0000-000000000abc:b2',
1078 '{"tokenCount":{"inputTokens":0,"outputTokens":0}}'),
1079 ('bubbleId:def00000-0000-0000-0000-000000000def:b1',
1080 '{"tokenCount":{"inputTokens":7,"outputTokens":3}}');"#,
1081 )
1082 .unwrap();
1083}
1084
1085#[cfg(any(test, feature = "test-support"))]
1086pub fn write_codex_state_db_for_test(home: &Path) {
1087 let codex_dir = home.join(".codex");
1088 fs::create_dir_all(&codex_dir).unwrap();
1089 let conn = rusqlite::Connection::open(codex_dir.join("state_5.sqlite")).unwrap();
1090 conn.execute_batch(
1091 "CREATE TABLE threads (
1092 id TEXT PRIMARY KEY,
1093 rollout_path TEXT,
1094 model TEXT,
1095 tokens_used INTEGER NOT NULL DEFAULT 0,
1096 preview TEXT,
1097 cwd TEXT,
1098 created_at_ms INTEGER,
1099 updated_at_ms INTEGER
1100 );
1101 INSERT INTO threads
1102 (id, rollout_path, model, tokens_used, preview, cwd, created_at_ms, updated_at_ms)
1103 VALUES
1104 ('019f49ca-54e7-7a91-82e7-a52b53cfd456', '/tmp/session.jsonl', 'gpt-web-ci', 33, 'web state prompt', '/work/repo', 1800000, 1900000);",
1105 )
1106 .unwrap();
1107}
1108
1109#[cfg(test)]
1110mod tests {
1111 use super::*;
1112
1113 const CODEX: &str = "/usr/bin/codex exec --skip-git-repo-check ";
1114
1115 #[test]
1116 fn agent_native_prompt_produces_llm_call_row() {
1117 let (_temp, path) = create_temp_session_path(agent_session::AGENT_CODEX);
1118 let session = parse_content_for_test(
1119 agent_session::AGENT_CODEX,
1120 &path,
1121 UNIX_EPOCH,
1122 "{\"type\":\"message\",\"content\":\"agentsight local codex prompt\"}\n",
1123 )
1124 .unwrap();
1125
1126 let view = materialized_view(&[session]);
1127 let rows = view.llm_call_rows(10);
1128
1129 assert_eq!(rows.len(), 1);
1130 assert_eq!(rows[0].comm.as_deref(), Some(agent_session::AGENT_CODEX));
1131 assert_eq!(
1132 rows[0].request.get("prompt").and_then(Value::as_str),
1133 Some("agentsight local codex prompt")
1134 );
1135 }
1136
1137 #[test]
1138 fn codex_state_db_produces_indexed_session_metadata() {
1139 let temp = tempfile::tempdir().unwrap();
1140 write_codex_state_db_for_test(temp.path());
1141
1142 let sessions = codex_state_sessions_in_home(temp.path(), 5);
1143
1144 assert_eq!(sessions.len(), 1);
1145 assert_eq!(sessions[0].display_id, "codex:019f49.fd456");
1146 assert_eq!(sessions[0].model.as_deref(), Some("gpt-web-ci"));
1147 assert_eq!(sessions[0].usage.total_tokens, 33);
1148 assert_eq!(
1149 sessions[0].prompt_preview.as_deref(),
1150 Some("web state prompt")
1151 );
1152 assert_eq!(sessions[0].cwd.as_deref(), Some("/work/repo"));
1153 assert_eq!(
1154 sessions[0].last_message_at.as_deref(),
1155 Some("1970-01-01T00:31:40.000Z")
1156 );
1157 }
1158
1159 #[test]
1160 fn cursor_state_db_path_checks_platform_layouts() {
1161 let temp = tempfile::tempdir().unwrap();
1162 assert!(cursor_state_db_path(temp.path()).is_none());
1163 write_cursor_state_db_for_test(temp.path());
1164 let path = cursor_state_db_path(temp.path()).unwrap();
1165 assert!(path.ends_with("Cursor/User/globalStorage/state.vscdb"));
1166 }
1167
1168 #[test]
1169 fn cursor_state_db_open_is_read_only_and_fails_closed() {
1170 let temp = tempfile::tempdir().unwrap();
1171 assert!(open_cursor_state_db(&temp.path().join("missing.vscdb")).is_none());
1172 write_cursor_state_db_for_test(temp.path());
1173 let conn = open_cursor_state_db(&cursor_state_db_path(temp.path()).unwrap()).unwrap();
1174 let denied = conn.execute(
1175 "INSERT INTO cursorDiskKV (key, value) VALUES ('x', 'y')",
1176 [],
1177 );
1178 assert!(denied.is_err());
1179 }
1180
1181 fn cursor_fixture_transcripts(home: &Path) -> PathBuf {
1182 let transcripts = home
1183 .join(".cursor/projects/repo/agent-transcripts/abc00000-0000-0000-0000-000000000abc");
1184 fs::create_dir_all(transcripts.join("subagents")).unwrap();
1185 let parent = transcripts.join("abc00000-0000-0000-0000-000000000abc.jsonl");
1186 fs::write(
1187 &parent,
1188 r#"{"role":"user","message":{"content":[{"type":"text","text":"check the build"}]}}"#,
1189 )
1190 .unwrap();
1191 fs::write(
1192 transcripts.join("subagents/def00000-0000-0000-0000-000000000def.jsonl"),
1193 "{}\n",
1194 )
1195 .unwrap();
1196 parent
1197 }
1198
1199 #[test]
1200 fn cursor_enrichment_fills_metadata_and_rolls_up_subagents() {
1201 let temp = tempfile::tempdir().unwrap();
1202 write_cursor_state_db_for_test(temp.path());
1203 let parent = cursor_fixture_transcripts(temp.path());
1204
1205 let mut sessions = vec![agent_session::parse_session_path(&parent).expect("parsed")];
1206 enrich_cursor_sessions_in_home(temp.path(), &mut sessions);
1207
1208 let session = &sessions[0];
1209 assert_eq!(session.model.as_deref(), Some("claude-sonnet-4-6"));
1210 assert_eq!(session.start_timestamp_ms, Some(1_700_000));
1211 assert_eq!(session.end_timestamp_ms, Some(1_900_000));
1212 assert_eq!(session.duration_ms, 200_000);
1213 assert_eq!(session.cwd.as_deref(), Some("/work/repo"));
1214 assert_eq!(
1215 session.last_message_at.as_deref(),
1216 Some("1970-01-01T00:31:40.000Z")
1217 );
1218 assert_eq!(session.usage.total_tokens, 150);
1220 assert_eq!(
1221 session
1222 .model_usage
1223 .get("claude-sonnet-4-6")
1224 .map(|usage| usage.total_tokens),
1225 Some(150)
1226 );
1227 }
1228
1229 #[test]
1230 fn cursor_enrichment_reads_model_when_header_row_is_missing() {
1231 let temp = tempfile::tempdir().unwrap();
1233 write_cursor_state_db_for_test(temp.path());
1234 let transcripts = temp
1235 .path()
1236 .join(".cursor/projects/repo/agent-transcripts/bbb00000-0000-0000-0000-000000000bbb");
1237 fs::create_dir_all(&transcripts).unwrap();
1238 let parent = transcripts.join("bbb00000-0000-0000-0000-000000000bbb.jsonl");
1239 fs::write(
1240 &parent,
1241 r#"{"role":"user","message":{"content":[{"type":"text","text":"check the build"}]}}"#,
1242 )
1243 .unwrap();
1244
1245 let parsed = agent_session::parse_session_path(&parent).expect("parsed");
1246 let mut sessions = vec![parsed.clone()];
1247 enrich_cursor_sessions_in_home(temp.path(), &mut sessions);
1248
1249 assert_eq!(
1250 sessions[0].model.as_deref(),
1251 Some("claude-4.6-sonnet-medium-thinking")
1252 );
1253 assert_eq!(sessions[0].start_timestamp_ms, parsed.start_timestamp_ms);
1254 assert_eq!(sessions[0].end_timestamp_ms, parsed.end_timestamp_ms);
1255 assert_eq!(sessions[0].last_message_at, parsed.last_message_at);
1256 assert_eq!(sessions[0].usage.total_tokens, parsed.usage.total_tokens);
1257 }
1258
1259 #[test]
1260 fn cursor_enrichment_missing_db_changes_nothing() {
1261 let temp = tempfile::tempdir().unwrap();
1262 let parent = cursor_fixture_transcripts(temp.path());
1263
1264 let parsed = agent_session::parse_session_path(&parent).expect("parsed");
1265 let mut sessions = vec![parsed.clone()];
1266 enrich_cursor_sessions_in_home(temp.path(), &mut sessions);
1267
1268 assert_eq!(sessions[0].model, parsed.model);
1269 assert_eq!(sessions[0].usage.total_tokens, parsed.usage.total_tokens);
1270 assert_eq!(sessions[0].cwd, parsed.cwd);
1271 assert_eq!(sessions[0].start_timestamp_ms, parsed.start_timestamp_ms);
1272 }
1273
1274 #[test]
1275 fn cursor_enrichment_reads_while_writer_holds_wal() {
1276 let temp = tempfile::tempdir().unwrap();
1277 write_cursor_state_db_for_test(temp.path());
1278 let parent = cursor_fixture_transcripts(temp.path());
1279
1280 let writer =
1281 rusqlite::Connection::open(cursor_state_db_path(temp.path()).unwrap()).unwrap();
1282 writer.execute_batch("BEGIN IMMEDIATE;").unwrap();
1283 writer
1284 .execute(
1285 "INSERT INTO cursorDiskKV (key, value) VALUES ('agentKv:x', 'held open')",
1286 [],
1287 )
1288 .unwrap();
1289
1290 let mut sessions = vec![agent_session::parse_session_path(&parent).expect("parsed")];
1291 enrich_cursor_sessions_in_home(temp.path(), &mut sessions);
1292 writer.execute_batch("ROLLBACK;").unwrap();
1293
1294 assert_eq!(sessions[0].model.as_deref(), Some("claude-sonnet-4-6"));
1295 assert_eq!(sessions[0].usage.total_tokens, 150);
1296 }
1297
1298 #[test]
1299 fn cursor_composer_data_reads_model_and_workspace() {
1300 let temp = tempfile::tempdir().unwrap();
1301 write_cursor_state_db_for_test(temp.path());
1302 let conn = open_cursor_state_db(&cursor_state_db_path(temp.path()).unwrap()).unwrap();
1303
1304 let pinned = cursor_composer_data(&conn, "abc00000-0000-0000-0000-000000000abc")
1305 .expect("pinned composer");
1306 assert_eq!(pinned.model.as_deref(), Some("claude-sonnet-4-6"));
1307 assert_eq!(pinned.workspace_path.as_deref(), Some("/work/repo"));
1308
1309 let unpinned = cursor_composer_data(&conn, "aaa00000-0000-0000-0000-000000000aaa")
1310 .expect("default composer");
1311 assert_eq!(unpinned.model, None);
1312 assert_eq!(unpinned.workspace_path, None);
1313
1314 assert!(cursor_composer_data(&conn, "not-a-composer").is_none());
1315 }
1316
1317 #[test]
1318 fn cursor_bubble_tokens_sums_by_bounded_range() {
1319 let temp = tempfile::tempdir().unwrap();
1320 write_cursor_state_db_for_test(temp.path());
1321 let conn = open_cursor_state_db(&cursor_state_db_path(temp.path()).unwrap()).unwrap();
1322
1323 let parent = cursor_bubble_tokens(&conn, "abc00000-0000-0000-0000-000000000abc");
1324 assert_eq!(parent.input_tokens, 100);
1325 assert_eq!(parent.output_tokens, 40);
1326 assert_eq!(parent.total_tokens, 140);
1327
1328 let child = cursor_bubble_tokens(&conn, "def00000-0000-0000-0000-000000000def");
1329 assert_eq!(child.total_tokens, 10);
1330
1331 let none = cursor_bubble_tokens(&conn, "aaa00000-0000-0000-0000-000000000aaa");
1332 assert_eq!(none.total_tokens, 0);
1333 }
1334
1335 #[test]
1336 fn cursor_subagent_ids_come_from_directory_layout() {
1337 let temp = tempfile::tempdir().unwrap();
1338 let transcripts = temp
1339 .path()
1340 .join(".cursor/projects/repo/agent-transcripts/abc");
1341 fs::create_dir_all(transcripts.join("subagents")).unwrap();
1342 let parent = transcripts.join("abc.jsonl");
1343 fs::write(&parent, "{}\n").unwrap();
1344 fs::write(transcripts.join("subagents/def.jsonl"), "{}\n").unwrap();
1345 fs::write(transcripts.join("subagents/aaa.jsonl"), "{}\n").unwrap();
1346 fs::write(transcripts.join("subagents/notes.txt"), "x").unwrap();
1347
1348 assert_eq!(cursor_subagent_ids(&parent), vec!["aaa", "def"]);
1349 let no_subagents = temp.path().join("elsewhere/abc.jsonl");
1350 assert!(cursor_subagent_ids(&no_subagents).is_empty());
1351 }
1352
1353 #[test]
1354 fn cursor_composer_header_reads_parent_and_subagent_rows() {
1355 let temp = tempfile::tempdir().unwrap();
1356 write_cursor_state_db_for_test(temp.path());
1357 let conn = open_cursor_state_db(&cursor_state_db_path(temp.path()).unwrap()).unwrap();
1358
1359 let parent = cursor_composer_header(&conn, "abc00000-0000-0000-0000-000000000abc")
1360 .expect("parent header");
1361 assert_eq!(parent.created_at_ms, Some(1_700_000));
1362 assert_eq!(parent.updated_at_ms, Some(1_900_000));
1363
1364 let child = cursor_composer_header(&conn, "def00000-0000-0000-0000-000000000def")
1365 .expect("subagent header");
1366 assert_eq!(child.created_at_ms, Some(1_750_000));
1367 assert_eq!(child.updated_at_ms, None);
1368
1369 assert!(cursor_composer_header(&conn, "not-a-composer").is_none());
1370 }
1371
1372 #[test]
1373 fn count_session_dirs_reports_cursor_root() {
1374 let temp = tempfile::tempdir().unwrap();
1375 assert!(agent_session::count_session_dirs_in_home(temp.path()).is_empty());
1376
1377 let transcripts = temp
1378 .path()
1379 .join(".cursor/projects/repo/agent-transcripts/abc");
1380 fs::create_dir_all(transcripts.join("subagents")).unwrap();
1381 fs::write(transcripts.join("abc.jsonl"), "{}\n").unwrap();
1382 fs::write(transcripts.join("subagents/def.jsonl"), "{}\n").unwrap();
1383
1384 let stats = agent_session::count_session_dirs_in_home(temp.path());
1385 assert_eq!(stats.len(), 1);
1386 assert_eq!(stats[0].agent, agent_session::AGENT_CURSOR);
1387 assert!(stats[0].dir.ends_with(".cursor/projects"));
1388 assert_eq!(stats[0].sessions, 1);
1391 assert_eq!(stats[0].bytes, 3);
1392 }
1393
1394 #[test]
1395 fn cursor_discovery_emits_parent_candidates_only() {
1396 let temp = tempfile::tempdir().unwrap();
1397 let project = temp.path().join(".cursor/projects/repo");
1398 let transcripts = project.join("agent-transcripts/abc");
1399 fs::create_dir_all(transcripts.join("subagents")).unwrap();
1400 fs::create_dir_all(project.join("canvases/node_modules/pkg")).unwrap();
1401 fs::write(transcripts.join("abc.jsonl"), "{}\n").unwrap();
1402 fs::write(transcripts.join("subagents/def.jsonl"), "{}\n").unwrap();
1403 fs::write(project.join("canvases/node_modules/pkg/data.jsonl"), "{}\n").unwrap();
1404
1405 let candidates: Vec<_> = agent_session::discover_session_files_in_home(temp.path())
1406 .into_iter()
1407 .filter(|candidate| candidate.agent == agent_session::AGENT_CURSOR)
1408 .collect();
1409
1410 assert_eq!(candidates.len(), 1);
1411 assert_eq!(candidates[0].path, transcripts.join("abc.jsonl"));
1412 }
1413
1414 #[test]
1415 fn cursor_candidate_updated_tracks_subagent_writes() {
1416 let temp = tempfile::tempdir().unwrap();
1417 let transcripts = temp
1418 .path()
1419 .join(".cursor/projects/repo/agent-transcripts/abc");
1420 fs::create_dir_all(transcripts.join("subagents")).unwrap();
1421 fs::write(transcripts.join("abc.jsonl"), "{}\n").unwrap();
1422 let child = transcripts.join("subagents/def.jsonl");
1423 fs::write(&child, "{}\n").unwrap();
1424
1425 let bumped = std::time::SystemTime::now() + std::time::Duration::from_secs(120);
1427 let handle = fs::File::options().write(true).open(&child).unwrap();
1428 handle.set_modified(bumped).unwrap();
1429
1430 let candidates: Vec<_> = agent_session::discover_session_files_in_home(temp.path())
1431 .into_iter()
1432 .filter(|candidate| candidate.agent == agent_session::AGENT_CURSOR)
1433 .collect();
1434
1435 assert_eq!(candidates.len(), 1);
1436 let parent_mtime = fs::metadata(transcripts.join("abc.jsonl"))
1437 .unwrap()
1438 .modified()
1439 .unwrap();
1440 assert!(candidates[0].updated > parent_mtime);
1441 }
1442
1443 #[test]
1444 fn cursor_duplicate_composer_prefers_real_workspace() {
1445 let temp = tempfile::tempdir().unwrap();
1446 let real = temp
1447 .path()
1448 .join(".cursor/projects/repo/agent-transcripts/abc");
1449 let stale = temp
1450 .path()
1451 .join(".cursor/projects/empty-window/agent-transcripts/abc");
1452 fs::create_dir_all(&real).unwrap();
1453 fs::create_dir_all(&stale).unwrap();
1454 fs::write(real.join("abc.jsonl"), "{}\n").unwrap();
1455 fs::write(stale.join("abc.jsonl"), "{}\n{}\n").unwrap();
1456
1457 let candidates: Vec<_> = agent_session::discover_session_files_in_home(temp.path())
1458 .into_iter()
1459 .filter(|candidate| candidate.agent == agent_session::AGENT_CURSOR)
1460 .collect();
1461
1462 assert_eq!(candidates.len(), 1);
1463 assert_eq!(candidates[0].path, real.join("abc.jsonl"));
1464 }
1465
1466 #[test]
1467 fn codex_state_db_uses_rollout_token_usage() {
1468 let temp = tempfile::tempdir().unwrap();
1469 write_codex_state_db_for_test(temp.path());
1470 let rollout = temp.path().join("session.jsonl");
1471 let mut content = "{}\n".repeat(CODEX_ROLLOUT_TAIL_BYTES as usize / 3 + 1);
1472 content.push_str(concat!(
1473 r#"{"type":"event_msg","payload":{"type":"token_count","info":{"total_token_usage":{"input_tokens":19184,"cached_input_tokens":9984,"output_tokens":11,"total_tokens":19195}}}}"#,
1474 "\n"
1475 ));
1476 fs::write(&rollout, content).unwrap();
1477 let conn = rusqlite::Connection::open(temp.path().join(".codex/state_5.sqlite")).unwrap();
1478 conn.execute(
1479 "UPDATE threads SET rollout_path = ?1, tokens_used = 999999999",
1480 [rollout.to_string_lossy().as_ref()],
1481 )
1482 .unwrap();
1483
1484 let sessions = codex_state_sessions_in_home(temp.path(), 5);
1485
1486 assert_eq!(sessions[0].usage.input_tokens, 9_200);
1487 assert_eq!(sessions[0].usage.cache_read_tokens, 9_984);
1488 assert_eq!(sessions[0].usage.output_tokens, 11);
1489 assert_eq!(sessions[0].usage.total_tokens, 19_195);
1490 }
1491
1492 #[test]
1493 fn codex_state_db_errors_return_empty_for_jsonl_fallback() {
1494 let temp = tempfile::tempdir().unwrap();
1495 fs::create_dir_all(temp.path().join(".codex")).unwrap();
1496 fs::write(temp.path().join(".codex/state_5.sqlite"), "not sqlite").unwrap();
1497
1498 assert!(codex_state_sessions_in_home(temp.path(), 5).is_empty());
1499 }
1500
1501 #[test]
1502 fn crowded_codex_home_fallback_keeps_only_observed_exec_prompt() {
1503 let temp = tempfile::tempdir().unwrap();
1504 let state_path = write_codex_home(
1505 temp.path(),
1506 &[
1507 "unrelated historical prompt",
1508 "unrelated historical prompt",
1509 "unrelated historical prompt",
1510 "agentsight current run prompt",
1511 "agentsight current run historical prompt",
1512 "unrelated historical prompt",
1513 "unrelated historical prompt",
1514 "unrelated historical prompt",
1515 ],
1516 );
1517 let now = current_epoch_ms();
1518
1519 let rows = vec![
1520 exec_row(
1521 "audit-exec",
1522 now,
1523 "codex",
1524 &format!("{CODEX}agentsight current run prompt"),
1525 ),
1526 exec_row(
1527 "audit-exec-truncated",
1528 now + 1,
1529 "node",
1530 &format!("{CODEX}agentsight current run"),
1531 ),
1532 file_row("audit-file", now + 100, &state_path),
1533 ];
1534
1535 let sessions = observed_sessions_from_audit_rows(&rows);
1536
1537 assert_eq!(sessions.len(), 1);
1538 assert_eq!(
1539 sessions[0].prompt_preview.as_deref(),
1540 Some("agentsight current run prompt")
1541 );
1542 }
1543
1544 #[test]
1545 fn codex_home_fallback_accepts_time_window_when_exec_prompt_is_truncated() {
1546 let temp = tempfile::tempdir().unwrap();
1547 let state_path = write_codex_home(temp.path(), &["agentsight truncated command prompt"]);
1548 let now = current_epoch_ms();
1549
1550 let sessions = observed_sessions_from_audit_rows(&[
1551 exec_row(
1552 "audit-exec",
1553 now,
1554 "codex",
1555 &format!("{CODEX}-c model_provider=\"agentsight-mock"),
1556 ),
1557 file_row("audit-file", now + 100, &state_path),
1558 ]);
1559
1560 assert_eq!(sessions.len(), 1);
1561 assert_eq!(
1562 sessions[0].prompt_preview.as_deref(),
1563 Some("agentsight truncated command prompt")
1564 );
1565 }
1566
1567 #[test]
1568 fn codex_exec_prompt_rows_filter_only_wrapper_duplicates() {
1569 let rows = [
1570 (1_000, "node", "/usr/bin/node /opt/codex/bin/codex exec --skip-git-repo-check agentsight dedupe prompt"),
1571 (1_001, "codex", "/opt/codex/bin/codex exec --skip-git-repo-check agentsight dedupe prompt"),
1572 (2_000, "codex", "/usr/bin/codex exec --skip-git-repo-check agentsight short prompt"),
1573 (3_000, "codex", "/usr/bin/codex exec --skip-git-repo-check agentsight much longer unrelated prompt"),
1574 (10_000, "codex", "/usr/bin/codex exec --skip-git-repo-check agentsight repeated prompt"),
1575 (11_000, "codex", "/usr/bin/codex exec --skip-git-repo-check agentsight repeated prompt"),
1576 (20_000, "codex", "/usr/bin/codex exec --skip-git-repo-check agentsight repeated prompt"),
1577 (21_000, "docker", "docker exec codex exec agentsight should not parse"),
1578 (22_000, "docker", "docker exec container /usr/local/bin/codex exec agentsight should parse once"),
1579 ]
1580 .into_iter()
1581 .enumerate()
1582 .map(|(index, (ts, comm, command))| exec_row(&format!("audit-{index}"), ts, comm, command))
1583 .collect::<Vec<_>>();
1584
1585 let projected = observed_session_prompt_rows(&rows);
1586 assert_eq!(projected[0].comm.as_deref(), Some("node"));
1587 assert_eq!(projected[0].target.as_deref(), Some("/usr/bin/node"));
1588 let prompts = projected
1589 .into_iter()
1590 .map(|row| row.summary)
1591 .collect::<Vec<_>>();
1592
1593 assert_eq!(
1594 prompts,
1595 vec![
1596 Some("agentsight dedupe prompt".to_string()),
1597 Some("agentsight short prompt".to_string()),
1598 Some("agentsight much longer unrelated prompt".to_string()),
1599 Some("agentsight repeated prompt".to_string()),
1600 Some("agentsight repeated prompt".to_string()),
1601 Some("agentsight repeated prompt".to_string()),
1602 Some("agentsight should parse once".to_string()),
1603 ]
1604 );
1605 }
1606
1607 #[test]
1608 fn codex_fallback_time_window_rejects_stale_matching_session() {
1609 let (_temp, path) = create_temp_session_path(agent_session::AGENT_CODEX);
1610 let session = parse_content_for_test(
1611 agent_session::AGENT_CODEX,
1612 &path,
1613 UNIX_EPOCH,
1614 "{\"type\":\"message\",\"content\":\"agentsight repeated prompt\"}\n",
1615 )
1616 .unwrap();
1617
1618 assert!(session_matches_observed_prompt(
1619 &session,
1620 &[ObservedCodexPrompt {
1621 prompt: "agentsight repeated prompt".to_string(),
1622 timestamp_ms: current_epoch_ms(),
1623 pid: Some(42),
1624 native_exec: true,
1625 comm: Some("codex".to_string()),
1626 target: Some("/usr/bin/codex".to_string()),
1627 }]
1628 ));
1629 assert!(!session_is_in_observed_window(
1630 &session,
1631 Some((current_epoch_ms() - 1_000, current_epoch_ms() + 1_000))
1632 ));
1633 }
1634
1635 fn exec_row(id: &str, timestamp_ms: u64, comm: &str, full_command: &str) -> AuditEventRow {
1636 AuditEventRow {
1637 id: id.to_string(),
1638 timestamp_ms,
1639 audit_type: "process".to_string(),
1640 pid: Some(42),
1641 comm: Some(comm.to_string()),
1642 subject: None,
1643 action: Some("exec".to_string()),
1644 target: Some(format!("/usr/bin/{comm}")),
1645 status: Some("observed".to_string()),
1646 summary: None,
1647 details: serde_json::json!({ "full_command": full_command }),
1648 }
1649 }
1650
1651 fn file_row(id: &str, timestamp_ms: u64, path: &Path) -> AuditEventRow {
1652 AuditEventRow {
1653 id: id.to_string(),
1654 timestamp_ms,
1655 audit_type: "file".to_string(),
1656 pid: Some(42),
1657 comm: Some("codex".to_string()),
1658 subject: None,
1659 action: Some("write".to_string()),
1660 target: Some(path.to_string_lossy().to_string()),
1661 status: Some("observed".to_string()),
1662 summary: None,
1663 details: serde_json::json!({ "filepath": path.to_string_lossy() }),
1664 }
1665 }
1666
1667 fn current_epoch_ms() -> u64 {
1668 std::time::SystemTime::now()
1669 .duration_since(UNIX_EPOCH)
1670 .unwrap()
1671 .as_millis() as u64
1672 }
1673
1674 fn write_codex_home(root: &Path, prompts: &[&str]) -> PathBuf {
1675 let codex_home = root.join("codex-home");
1676 let sessions_dir = codex_home.join("sessions/2026/07/14");
1677 fs::create_dir_all(&sessions_dir).unwrap();
1678 for (index, prompt) in prompts.iter().enumerate() {
1679 fs::write(
1680 sessions_dir.join(format!(
1681 "rollout-2026-07-14T00-00-{index:02}-session-{index}.jsonl"
1682 )),
1683 format!(
1684 "{{\"timestamp\":\"2026-07-14T00:00:{index:02}.000Z\",\
1685 \"type\":\"event_msg\",\
1686 \"payload\":{{\"type\":\"user_message\",\"message\":\"{prompt}\"}}}}\n"
1687 ),
1688 )
1689 .unwrap();
1690 }
1691 let state_path = codex_home.join("stat");
1692 fs::write(&state_path, "").unwrap();
1693 state_path
1694 }
1695}