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(dirs::home_dir)
416}
417
418fn non_negative_i64_to_u64(value: i64) -> Option<u64> {
419 u64::try_from(value).ok()
420}
421
422fn system_time_from_ms(value: u64) -> SystemTime {
423 UNIX_EPOCH + Duration::from_millis(value)
424}
425
426fn iso_utc_from_ms(value: u64) -> String {
427 chrono::DateTime::<chrono::Utc>::from_timestamp_millis(value as i64)
428 .map(|dt| dt.to_rfc3339_opts(chrono::SecondsFormat::Millis, true))
429 .unwrap_or_default()
430}
431
432fn short_session_id(id: &str) -> String {
433 let compact = id.rsplit(['/', '\\']).next().unwrap_or(id).trim();
434 if compact.chars().count() <= 12 {
435 return compact.to_string();
436 }
437 let head = compact.chars().take(6).collect::<String>();
438 let tail = compact
439 .chars()
440 .rev()
441 .take(5)
442 .collect::<Vec<_>>()
443 .into_iter()
444 .rev()
445 .collect::<String>();
446 format!("{head}.{tail}")
447}
448
449fn clean_prompt_text(text: &str) -> Option<String> {
450 let text = text.split_whitespace().collect::<Vec<_>>().join(" ");
451 (!text.trim().is_empty()).then(|| text.trim().to_string())
452}
453
454fn view_id(session: &LocalSession) -> String {
455 format!("local:{}:{}", session.agent_type, session.display_id)
456}
457
458pub fn materialized_view(sessions: &[LocalSession]) -> MaterializedView {
459 let mut view = MaterializedView::new();
460 view.set_source(AGENT_NATIVE_SOURCE);
461 import_into_view(&mut view, sessions);
462 view
463}
464
465pub fn import_recent(view: &mut MaterializedView, limit: usize) {
466 let mut sessions = SessionCache::new().discover_cached(limit, Duration::ZERO);
467 enrich_cursor_sessions(&mut sessions);
471 import_into_view(view, &sessions);
472}
473
474pub fn import_into_view(view: &mut MaterializedView, sessions: &[LocalSession]) {
475 for session in sessions {
476 view.upsert_session(&session_row(session));
477 for row in llm_rows(session) {
478 view.apply_llm_call(&row);
479 }
480 for row in token_rows(session) {
481 view.apply_token_usage(&row);
482 }
483 for row in tool_rows(session) {
484 view.apply_tool_call(&row);
485 }
486 }
487}
488
489fn llm_rows(session: &LocalSession) -> Vec<LlmCallRow> {
490 let Some(prompt) = session.prompt_preview.as_ref() else {
491 return Vec::new();
492 };
493 let session_id = view_id(session);
494 let timestamp_ms = session
495 .events
496 .prompts
497 .first()
498 .and_then(|prompt| prompt.ts_ms)
499 .and_then(|ts| u64::try_from(ts).ok())
500 .or(session.start_timestamp_ms)
501 .unwrap_or_else(|| updated_ms(session));
502 let request = serde_json::json!({
503 "prompt": prompt,
504 "prompt_source": AGENT_NATIVE_SOURCE,
505 "session_id": session_id,
506 "agent_type": session.agent_type.as_str(),
507 "path": session.path.to_string_lossy(),
508 });
509
510 if session.model_usage.is_empty() {
511 let model = session
512 .model
513 .clone()
514 .unwrap_or_else(|| session.agent_type.clone());
515 return vec![llm_row_for_session(
516 &format!("{session_id}-{}", sanitize_id(&model)),
517 session,
518 &session_id,
519 timestamp_ms,
520 Some(model),
521 &session.usage,
522 request,
523 )];
524 }
525
526 session
527 .model_usage
528 .iter()
529 .map(|(model, usage)| {
530 llm_row_for_session(
531 &format!("{session_id}-{model}"),
532 session,
533 &session_id,
534 timestamp_ms,
535 Some(model.clone()),
536 usage,
537 request.clone(),
538 )
539 })
540 .collect()
541}
542
543fn llm_row_for_session(
544 id: &str,
545 session: &LocalSession,
546 session_id: &str,
547 timestamp_ms: u64,
548 model: Option<String>,
549 usage: &TokenUsage,
550 request: Value,
551) -> LlmCallRow {
552 LlmCallRow {
553 id: id.to_string(),
554 session_id: Some(session_id.to_string()),
555 conversation_id: session.conversation_id.clone(),
556 start_timestamp_ms: timestamp_ms,
557 end_timestamp_ms: session.end_timestamp_ms,
558 pid: None,
559 comm: Some(session.agent_type.clone()),
560 provider: None,
561 model,
562 call_kind: Some("agent_native_prompt".to_string()),
563 status: "observed".to_string(),
564 error_type: None,
565 finish_reason: None,
566 host: None,
567 path: Some(session.path.to_string_lossy().to_string()),
568 status_code: None,
569 input_tokens: usage.input_tokens,
570 output_tokens: usage.output_tokens,
571 total_tokens: usage.total_tokens,
572 request,
573 response: Value::Null,
574 }
575}
576
577pub fn observed_session_prompt_rows(audit_rows: &[AuditEventRow]) -> Vec<AuditEventRow> {
578 let mut rows = Vec::new();
579 let mut seen = HashSet::new();
580 let observed_exec_prompts = observed_codex_exec_prompts(audit_rows);
581 let mut seen_exec_prompts: Vec<ObservedCodexPrompt> = Vec::new();
582 for observed in observed_exec_prompts {
583 if seen_exec_prompts.iter().any(|seen| {
584 seen.prompt == observed.prompt
585 && timestamps_close(
586 seen.timestamp_ms,
587 observed.timestamp_ms,
588 CODEX_EXEC_DEDUPE_WINDOW_MS,
589 )
590 && (!seen.native_exec || !observed.native_exec)
591 }) {
592 continue;
593 }
594 seen_exec_prompts.push(observed.clone());
595 rows.push(AuditEventRow {
596 id: format!(
597 "audit-codex-exec-prompt-{}-{}",
598 observed.timestamp_ms,
599 observed.pid.unwrap_or(0)
600 ),
601 timestamp_ms: observed.timestamp_ms,
602 audit_type: "llm".to_string(),
603 pid: observed.pid,
604 comm: observed.comm.or_else(|| Some("codex".to_string())),
605 subject: None,
606 action: Some("request".to_string()),
607 target: observed.target,
608 status: Some("observed".to_string()),
609 summary: Some(truncate_text(&observed.prompt, 160)),
610 details: serde_json::json!({
611 "text_content": observed.prompt,
612 "prompt_source": "local",
613 }),
614 });
615 }
616 for row in audit_rows {
617 if row.audit_type == "process" && row.action.as_deref() == Some("exec") {
618 continue;
619 }
620 if row.audit_type != "file" {
621 continue;
622 }
623 let Some(pid) = row.pid else {
624 continue;
625 };
626 let Some(path) = audit_session_path(row) else {
627 continue;
628 };
629 if !seen.insert((path.clone(), pid)) {
630 continue;
631 };
632 let Some(session) = agent_session::parse_session_path(&path) else {
633 continue;
634 };
635 let Some(prompt) = session.prompt_preview.as_ref() else {
636 continue;
637 };
638 rows.push(AuditEventRow {
639 id: format!(
640 "audit-agent-native-prompt-{}-{pid}",
641 sanitize_id(&session.display_id)
642 ),
643 timestamp_ms: row.timestamp_ms,
644 audit_type: "llm".to_string(),
645 pid: Some(pid),
646 comm: row
647 .comm
648 .clone()
649 .or_else(|| Some(session.agent_type.clone())),
650 subject: session.model.clone(),
651 action: Some("request".to_string()),
652 target: Some(path.to_string_lossy().to_string()),
653 status: Some("observed".to_string()),
654 summary: Some(truncate_text(prompt, 160)),
655 details: serde_json::json!({
656 "text_content": prompt,
657 "prompt_source": "local",
658 "session_id": view_id(&session),
659 "conversation_id": session.conversation_id.as_deref(),
660 "agent_type": session.agent_type,
661 }),
662 });
663 }
664 rows
665}
666
667pub fn observed_sessions_from_audit_rows(audit_rows: &[AuditEventRow]) -> Vec<LocalSession> {
668 let mut direct_paths = HashSet::new();
669 let mut codex_session_dirs = HashSet::new();
670 let observed_codex_prompts = observed_codex_exec_prompts(audit_rows);
671 let observed_codex_exec = observed_codex_exec_command(audit_rows);
672 let observed_window = observed_audit_window_ms(audit_rows);
673
674 for row in audit_rows.iter().filter(|row| row.audit_type == "file") {
675 for path in audit_file_paths(row) {
676 if let Some(session_path) =
677 agent_session::session_log_path_from_str(path.to_string_lossy().as_ref())
678 {
679 direct_paths.insert(session_path);
680 }
681 if let Some(dir) = observed_codex_sessions_dir(&path) {
682 codex_session_dirs.insert(dir);
683 }
684 }
685 }
686
687 let mut candidates = Vec::new();
688 for path in direct_paths {
689 if let Some(candidate) = agent_session::session_candidate_from_path(&path) {
690 candidates.push((candidate, false));
691 }
692 }
693 for dir in codex_session_dirs {
694 let dir_candidates =
695 agent_session::discover_session_files_in_dir(agent_session::AGENT_CODEX, &dir);
696 candidates.extend(
697 dir_candidates
698 .into_iter()
699 .map(|candidate| (candidate, true)),
700 );
701 }
702 candidates.sort_by_key(|(candidate, _)| Reverse(candidate.updated));
703
704 let mut seen_paths = HashSet::new();
705 let mut seen_sessions = HashSet::new();
706 let mut sessions = Vec::new();
707 for (candidate, is_codex_dir_fallback) in candidates.into_iter().take(75) {
708 if !seen_paths.insert(candidate.path.clone()) {
709 continue;
710 }
711 let Some(session) = agent_session::parse_session_file(&candidate) else {
712 continue;
713 };
714 if is_codex_dir_fallback && !observed_codex_exec {
715 continue;
716 }
717 if is_codex_dir_fallback
718 && !observed_codex_prompts.is_empty()
719 && !session_matches_observed_prompt(&session, &observed_codex_prompts)
720 {
721 continue;
722 }
723 if is_codex_dir_fallback && !session_is_in_observed_window(&session, observed_window) {
724 continue;
725 }
726 if seen_sessions.insert(session.display_id.clone()) {
727 sessions.push(session);
728 }
729 }
730 sessions
731}
732
733fn observed_codex_exec_prompts(audit_rows: &[AuditEventRow]) -> Vec<ObservedCodexPrompt> {
734 let prompts = audit_rows
735 .iter()
736 .filter(|row| row.audit_type == "process" && row.action.as_deref() == Some("exec"))
737 .filter_map(|row| {
738 let prompt = row
739 .details
740 .get("full_command")
741 .and_then(Value::as_str)
742 .and_then(codex_exec_prompt_from_command)?;
743 Some(ObservedCodexPrompt {
744 prompt,
745 timestamp_ms: row.timestamp_ms,
746 pid: row.pid,
747 native_exec: looks_like_native_codex_exec(row),
748 comm: row.comm.clone(),
749 target: row.target.clone(),
750 })
751 })
752 .collect::<Vec<_>>();
753 prompts
754 .iter()
755 .filter(|candidate| {
756 !prompts
757 .iter()
758 .any(|other| is_nearby_longer_prefix(candidate, other))
759 })
760 .cloned()
761 .collect()
762}
763
764fn observed_codex_exec_command(audit_rows: &[AuditEventRow]) -> bool {
765 audit_rows
766 .iter()
767 .filter(|row| row.audit_type == "process" && row.action.as_deref() == Some("exec"))
768 .filter_map(|row| row.details.get("full_command").and_then(Value::as_str))
769 .any(|command| codex_exec_command_tail(command).is_some())
770}
771
772fn codex_exec_prompt_from_command(command: &str) -> Option<String> {
773 agent_session::codex_exec_prompt(&codex_exec_command_tail(command)?)
774}
775
776fn codex_exec_command_tail(command: &str) -> Option<String> {
777 let tokens = command.split_whitespace().collect::<Vec<_>>();
778 let index = tokens.windows(2).enumerate().find_map(|(index, tokens)| {
779 (is_codex_executable_token(tokens[0], index == 0) && tokens[1] == "exec").then_some(index)
780 })?;
781 Some(tokens[index..].join(" "))
782}
783
784fn looks_like_native_codex_exec(row: &AuditEventRow) -> bool {
785 row.comm.as_deref() == Some("codex")
786 && row
787 .target
788 .as_deref()
789 .is_some_and(|target| is_codex_executable_token(target, true))
790}
791
792fn is_codex_executable_token(token: &str, allow_bare: bool) -> bool {
793 let token = token.trim_matches(|ch| matches!(ch, '"' | '\''));
794 (allow_bare && token == "codex")
795 || token.contains('/')
796 && Path::new(token).file_name().and_then(|name| name.to_str()) == Some("codex")
797}
798
799fn is_nearby_longer_prefix(candidate: &ObservedCodexPrompt, other: &ObservedCodexPrompt) -> bool {
800 other.prompt.len() > candidate.prompt.len()
801 && other.prompt.starts_with(candidate.prompt.as_str())
802 && timestamps_close(
803 candidate.timestamp_ms,
804 other.timestamp_ms,
805 CODEX_EXEC_DEDUPE_WINDOW_MS,
806 )
807 && (!candidate.native_exec || !other.native_exec)
808}
809
810fn timestamps_close(left: u64, right: u64, window_ms: u64) -> bool {
811 left.abs_diff(right) <= window_ms
812}
813
814fn session_matches_observed_prompt(
815 session: &LocalSession,
816 prompts: &[ObservedCodexPrompt],
817) -> bool {
818 let Some(preview) = session.prompt_preview.as_deref() else {
819 return false;
820 };
821 prompts.iter().any(|observed| {
822 let prompt = observed.prompt.as_str();
823 prompt_texts_overlap(prompt, preview)
824 })
825}
826
827fn prompt_texts_overlap(left: &str, right: &str) -> bool {
828 left == right || left.starts_with(right) || right.starts_with(left)
829}
830
831fn observed_audit_window_ms(audit_rows: &[AuditEventRow]) -> Option<(u64, u64)> {
832 let min = audit_rows.iter().map(|row| row.timestamp_ms).min()?;
833 let max = audit_rows.iter().map(|row| row.timestamp_ms).max()?;
834 Some((
835 min.saturating_sub(CODEX_FALLBACK_TIME_SLOP_MS),
836 max.saturating_add(CODEX_FALLBACK_TIME_SLOP_MS),
837 ))
838}
839
840fn session_is_in_observed_window(session: &LocalSession, window: Option<(u64, u64)>) -> bool {
841 let Some((min, max)) = window else {
842 return true;
843 };
844 let updated = updated_ms(session);
845 updated >= min && updated <= max
846}
847
848fn audit_session_path(row: &AuditEventRow) -> Option<PathBuf> {
849 row.target
850 .as_deref()
851 .and_then(agent_session::session_log_path_from_str)
852 .or_else(|| {
853 row.details
854 .get("filepath")
855 .and_then(Value::as_str)
856 .and_then(agent_session::session_log_path_from_str)
857 })
858 .or_else(|| {
859 row.details
860 .get("path")
861 .and_then(Value::as_str)
862 .and_then(agent_session::session_log_path_from_str)
863 })
864}
865
866fn audit_file_paths(row: &AuditEventRow) -> Vec<PathBuf> {
867 [
868 row.target.as_deref(),
869 row.details.get("filepath").and_then(Value::as_str),
870 row.details.get("path").and_then(Value::as_str),
871 row.details.get("fd_target").and_then(Value::as_str),
872 ]
873 .into_iter()
874 .flatten()
875 .filter_map(|raw| {
876 let path = PathBuf::from(raw.trim().trim_end_matches(" (deleted)"));
877 path.is_absolute().then_some(path)
878 })
879 .collect()
880}
881
882fn observed_codex_sessions_dir(path: &Path) -> Option<PathBuf> {
883 if !looks_like_codex_home_file(path) {
884 return None;
885 }
886 let home = path.parent()?;
887 let sessions = home.join("sessions");
888 sessions.is_dir().then_some(sessions)
889}
890
891fn looks_like_codex_home_file(path: &Path) -> bool {
892 let Some(name) = path.file_name().and_then(|name| name.to_str()) else {
893 return false;
894 };
895 (name.starts_with("state_")
896 || name.starts_with("logs_")
897 || matches!(name, "config.toml" | "auth.json" | "stat"))
898 && path
899 .parent()
900 .is_some_and(|parent| parent.join("sessions").is_dir())
901}
902
903fn session_row(session: &LocalSession) -> SessionRow {
904 let updated_ms = updated_ms(session);
905 SessionRow {
906 id: view_id(session),
907 agent_type: session.agent_type.clone(),
908 start_timestamp_ms: session
909 .start_timestamp_ms
910 .unwrap_or_else(|| updated_ms.saturating_sub(session.duration_ms)),
911 end_timestamp_ms: session.end_timestamp_ms.or(Some(updated_ms)),
912 status: "observed".to_string(),
913 model: session.model.clone(),
914 input_tokens: session.usage.input_tokens,
915 output_tokens: session.usage.output_tokens,
916 total_tokens: session.usage.total_tokens,
917 view_source: AGENT_NATIVE_SOURCE.to_string(),
918 confidence: Some(0.95),
919 attributes: serde_json::json!({
920 "session_id": session.session_id.clone(),
921 "conversation_id": session.conversation_id.as_deref(),
922 "path": session.path.to_string_lossy(),
923 "display_id": session.display_id,
924 "prompt_preview": session.prompt_preview.clone(),
925 "cwd": session.cwd.clone(),
926 "last_message_at": session.last_message_at.clone(),
927 "files": session.files,
928 }),
929 }
930}
931
932fn token_rows(session: &LocalSession) -> Vec<TokenUsageRow> {
933 let session_id = view_id(session);
934 session
935 .model_usage
936 .iter()
937 .filter(|(_, usage)| usage.total_tokens > 0)
938 .map(|(model, usage)| TokenUsageRow {
939 id: format!("token-{session_id}-{}", sanitize_id(model)),
940 llm_call_id: format!("{session_id}-{model}"),
941 timestamp_ms: updated_ms(session),
942 pid: None,
943 comm: Some(session.agent_type.clone()),
944 provider: None,
945 model: Some(model.clone()),
946 input_tokens: usage.input_tokens,
947 output_tokens: usage.output_tokens,
948 cache_creation_tokens: usage.cache_creation_tokens,
949 cache_read_tokens: usage.cache_read_tokens,
950 total_tokens: usage.total_tokens,
951 source: AGENT_NATIVE_SOURCE.to_string(),
952 view_source: AGENT_NATIVE_SOURCE.to_string(),
953 confidence: Some(0.95),
954 })
955 .collect()
956}
957
958fn tool_rows(session: &LocalSession) -> Vec<ToolCallRow> {
959 let session_id = view_id(session);
960 let timestamp_ms = updated_ms(session);
961 let mut rows = Vec::new();
962 for (tool, count) in &session.tools {
963 for index in 0..*count {
964 rows.push(ToolCallRow {
965 id: format!("tool-{session_id}-{}-{index}", sanitize_id(tool)),
966 session_id: Some(session_id.clone()),
967 conversation_id: session.conversation_id.clone(),
968 timestamp_ms,
969 tool_name: Some(tool.clone()),
970 tool_call_id: None,
971 start_timestamp_ms: Some(timestamp_ms),
972 end_timestamp_ms: Some(timestamp_ms),
973 duration_ms: None,
974 status: Some("observed".to_string()),
975 input: serde_json::json!({}),
976 output: serde_json::json!({}),
977 related_pid: None,
978 related_event_id: None,
979 view_source: AGENT_NATIVE_SOURCE.to_string(),
980 confidence: Some(0.95),
981 });
982 }
983 }
984 rows
985}
986
987fn updated_ms(session: &LocalSession) -> u64 {
988 session
989 .updated
990 .duration_since(UNIX_EPOCH)
991 .unwrap_or_default()
992 .as_millis() as u64
993}
994
995fn matches_filter(
996 session: &LocalSession,
997 pid_filter: Option<u32>,
998 text_filter: Option<&str>,
999) -> bool {
1000 if pid_filter.is_some() {
1001 return true;
1002 }
1003 let Some(filter) = text_filter else {
1004 return true;
1005 };
1006 let filter = filter.to_ascii_lowercase();
1007 session.agent_type.to_ascii_lowercase().contains(&filter)
1008 || session
1009 .prompt_preview
1010 .as_ref()
1011 .is_some_and(|prompt| prompt.to_ascii_lowercase().contains(&filter))
1012 || session
1013 .model
1014 .as_ref()
1015 .is_some_and(|model| model.to_ascii_lowercase().contains(&filter))
1016 || session
1017 .path
1018 .to_string_lossy()
1019 .to_ascii_lowercase()
1020 .contains(&filter)
1021}
1022
1023#[cfg(any(test, feature = "test-support"))]
1024pub fn create_temp_session_path(agent: &str) -> (tempfile::TempDir, PathBuf) {
1025 let temp = tempfile::tempdir().unwrap();
1026 let path = agent_session::fixture_session_path(agent, temp.path()).unwrap();
1027 fs::create_dir_all(path.parent().unwrap()).unwrap();
1028 fs::write(&path, "{}\n").unwrap();
1029 (temp, path)
1030}
1031
1032#[cfg(any(test, feature = "test-support"))]
1033pub fn parse_content_for_test(
1034 agent: &str,
1035 path: &std::path::Path,
1036 updated: std::time::SystemTime,
1037 content: &str,
1038) -> Option<LocalSession> {
1039 agent_session::parse_session_content(agent, path, updated, content)
1040}
1041
1042#[cfg(any(test, feature = "test-support"))]
1046pub fn write_cursor_state_db_for_test(home: &Path) {
1047 let db_dir = home.join("Library/Application Support/Cursor/User/globalStorage");
1048 fs::create_dir_all(&db_dir).unwrap();
1049 let conn = rusqlite::Connection::open(db_dir.join("state.vscdb")).unwrap();
1050 conn.pragma_update(None, "journal_mode", "WAL").unwrap();
1053 conn.execute_batch(
1054 r#"CREATE TABLE composerHeaders (composerId TEXT PRIMARY KEY, workspaceId TEXT,
1055 createdAt INTEGER, lastUpdatedAt INTEGER, isArchived INTEGER,
1056 isSubagent INTEGER, recency INTEGER, checkpointAt INTEGER, value TEXT);
1057 CREATE TABLE cursorDiskKV (key TEXT UNIQUE ON CONFLICT REPLACE, value BLOB);
1058 INSERT INTO composerHeaders (composerId, workspaceId, createdAt, lastUpdatedAt, isSubagent)
1059 VALUES
1060 ('abc00000-0000-0000-0000-000000000abc', 'ws1', 1700000, 1900000, 0),
1061 ('def00000-0000-0000-0000-000000000def', 'ws1', 1750000, NULL, 1),
1062 ('aaa00000-0000-0000-0000-000000000aaa', 'ws2', 1600000, 1650000, 0);
1063 INSERT INTO cursorDiskKV (key, value) VALUES
1064 ('composerData:bbb00000-0000-0000-0000-000000000bbb',
1065 '{"composerId":"bbb00000-0000-0000-0000-000000000bbb","modelConfig":{"modelName":"claude-4.6-sonnet-medium-thinking","maxMode":false}}'),
1066 ('composerData:abc00000-0000-0000-0000-000000000abc',
1067 '{"composerId":"abc00000-0000-0000-0000-000000000abc","modelConfig":{"modelName":"claude-sonnet-4-6","maxMode":false},"workspaceIdentifier":{"uri":{"fsPath":"/work/repo"}}}'),
1068 ('composerData:aaa00000-0000-0000-0000-000000000aaa',
1069 '{"composerId":"aaa00000-0000-0000-0000-000000000aaa","modelConfig":{"modelName":"default","maxMode":false}}'),
1070 ('bubbleId:abc00000-0000-0000-0000-000000000abc:b1',
1071 '{"tokenCount":{"inputTokens":100,"outputTokens":40}}'),
1072 ('bubbleId:abc00000-0000-0000-0000-000000000abc:b2',
1073 '{"tokenCount":{"inputTokens":0,"outputTokens":0}}'),
1074 ('bubbleId:def00000-0000-0000-0000-000000000def:b1',
1075 '{"tokenCount":{"inputTokens":7,"outputTokens":3}}');"#,
1076 )
1077 .unwrap();
1078}
1079
1080#[cfg(any(test, feature = "test-support"))]
1081pub fn write_codex_state_db_for_test(home: &Path) {
1082 let codex_dir = home.join(".codex");
1083 fs::create_dir_all(&codex_dir).unwrap();
1084 let conn = rusqlite::Connection::open(codex_dir.join("state_5.sqlite")).unwrap();
1085 conn.execute_batch(
1086 "CREATE TABLE threads (
1087 id TEXT PRIMARY KEY,
1088 rollout_path TEXT,
1089 model TEXT,
1090 tokens_used INTEGER NOT NULL DEFAULT 0,
1091 preview TEXT,
1092 cwd TEXT,
1093 created_at_ms INTEGER,
1094 updated_at_ms INTEGER
1095 );
1096 INSERT INTO threads
1097 (id, rollout_path, model, tokens_used, preview, cwd, created_at_ms, updated_at_ms)
1098 VALUES
1099 ('019f49ca-54e7-7a91-82e7-a52b53cfd456', '/tmp/session.jsonl', 'gpt-web-ci', 33, 'web state prompt', '/work/repo', 1800000, 1900000);",
1100 )
1101 .unwrap();
1102}
1103
1104#[cfg(test)]
1105mod tests {
1106 use super::*;
1107
1108 const CODEX: &str = "/usr/bin/codex exec --skip-git-repo-check ";
1109
1110 #[test]
1111 fn agent_native_prompt_produces_llm_call_row() {
1112 let (_temp, path) = create_temp_session_path(agent_session::AGENT_CODEX);
1113 let session = parse_content_for_test(
1114 agent_session::AGENT_CODEX,
1115 &path,
1116 UNIX_EPOCH,
1117 "{\"type\":\"message\",\"content\":\"agentsight local codex prompt\"}\n",
1118 )
1119 .unwrap();
1120
1121 let view = materialized_view(&[session]);
1122 let rows = view.llm_call_rows(10);
1123
1124 assert_eq!(rows.len(), 1);
1125 assert_eq!(rows[0].comm.as_deref(), Some(agent_session::AGENT_CODEX));
1126 assert_eq!(
1127 rows[0].request.get("prompt").and_then(Value::as_str),
1128 Some("agentsight local codex prompt")
1129 );
1130 }
1131
1132 #[test]
1133 fn codex_state_db_produces_indexed_session_metadata() {
1134 let temp = tempfile::tempdir().unwrap();
1135 write_codex_state_db_for_test(temp.path());
1136
1137 let sessions = codex_state_sessions_in_home(temp.path(), 5);
1138
1139 assert_eq!(sessions.len(), 1);
1140 assert_eq!(sessions[0].display_id, "codex:019f49.fd456");
1141 assert_eq!(sessions[0].model.as_deref(), Some("gpt-web-ci"));
1142 assert_eq!(sessions[0].usage.total_tokens, 33);
1143 assert_eq!(
1144 sessions[0].prompt_preview.as_deref(),
1145 Some("web state prompt")
1146 );
1147 assert_eq!(sessions[0].cwd.as_deref(), Some("/work/repo"));
1148 assert_eq!(
1149 sessions[0].last_message_at.as_deref(),
1150 Some("1970-01-01T00:31:40.000Z")
1151 );
1152 }
1153
1154 #[test]
1155 fn cursor_state_db_path_checks_platform_layouts() {
1156 let temp = tempfile::tempdir().unwrap();
1157 assert!(cursor_state_db_path(temp.path()).is_none());
1158 write_cursor_state_db_for_test(temp.path());
1159 let path = cursor_state_db_path(temp.path()).unwrap();
1160 assert!(path.ends_with("Cursor/User/globalStorage/state.vscdb"));
1161 }
1162
1163 #[test]
1164 fn cursor_state_db_open_is_read_only_and_fails_closed() {
1165 let temp = tempfile::tempdir().unwrap();
1166 assert!(open_cursor_state_db(&temp.path().join("missing.vscdb")).is_none());
1167 write_cursor_state_db_for_test(temp.path());
1168 let conn = open_cursor_state_db(&cursor_state_db_path(temp.path()).unwrap()).unwrap();
1169 let denied = conn.execute(
1170 "INSERT INTO cursorDiskKV (key, value) VALUES ('x', 'y')",
1171 [],
1172 );
1173 assert!(denied.is_err());
1174 }
1175
1176 fn cursor_fixture_transcripts(home: &Path) -> PathBuf {
1177 let transcripts = home
1178 .join(".cursor/projects/repo/agent-transcripts/abc00000-0000-0000-0000-000000000abc");
1179 fs::create_dir_all(transcripts.join("subagents")).unwrap();
1180 let parent = transcripts.join("abc00000-0000-0000-0000-000000000abc.jsonl");
1181 fs::write(
1182 &parent,
1183 r#"{"role":"user","message":{"content":[{"type":"text","text":"check the build"}]}}"#,
1184 )
1185 .unwrap();
1186 fs::write(
1187 transcripts.join("subagents/def00000-0000-0000-0000-000000000def.jsonl"),
1188 "{}\n",
1189 )
1190 .unwrap();
1191 parent
1192 }
1193
1194 #[test]
1195 fn cursor_enrichment_fills_metadata_and_rolls_up_subagents() {
1196 let temp = tempfile::tempdir().unwrap();
1197 write_cursor_state_db_for_test(temp.path());
1198 let parent = cursor_fixture_transcripts(temp.path());
1199
1200 let mut sessions = vec![agent_session::parse_session_path(&parent).expect("parsed")];
1201 enrich_cursor_sessions_in_home(temp.path(), &mut sessions);
1202
1203 let session = &sessions[0];
1204 assert_eq!(session.model.as_deref(), Some("claude-sonnet-4-6"));
1205 assert_eq!(session.start_timestamp_ms, Some(1_700_000));
1206 assert_eq!(session.end_timestamp_ms, Some(1_900_000));
1207 assert_eq!(session.duration_ms, 200_000);
1208 assert_eq!(session.cwd.as_deref(), Some("/work/repo"));
1209 assert_eq!(
1210 session.last_message_at.as_deref(),
1211 Some("1970-01-01T00:31:40.000Z")
1212 );
1213 assert_eq!(session.usage.total_tokens, 150);
1215 assert_eq!(
1216 session
1217 .model_usage
1218 .get("claude-sonnet-4-6")
1219 .map(|usage| usage.total_tokens),
1220 Some(150)
1221 );
1222 }
1223
1224 #[test]
1225 fn cursor_enrichment_reads_model_when_header_row_is_missing() {
1226 let temp = tempfile::tempdir().unwrap();
1228 write_cursor_state_db_for_test(temp.path());
1229 let transcripts = temp
1230 .path()
1231 .join(".cursor/projects/repo/agent-transcripts/bbb00000-0000-0000-0000-000000000bbb");
1232 fs::create_dir_all(&transcripts).unwrap();
1233 let parent = transcripts.join("bbb00000-0000-0000-0000-000000000bbb.jsonl");
1234 fs::write(
1235 &parent,
1236 r#"{"role":"user","message":{"content":[{"type":"text","text":"check the build"}]}}"#,
1237 )
1238 .unwrap();
1239
1240 let parsed = agent_session::parse_session_path(&parent).expect("parsed");
1241 let mut sessions = vec![parsed.clone()];
1242 enrich_cursor_sessions_in_home(temp.path(), &mut sessions);
1243
1244 assert_eq!(
1245 sessions[0].model.as_deref(),
1246 Some("claude-4.6-sonnet-medium-thinking")
1247 );
1248 assert_eq!(sessions[0].start_timestamp_ms, parsed.start_timestamp_ms);
1249 assert_eq!(sessions[0].end_timestamp_ms, parsed.end_timestamp_ms);
1250 assert_eq!(sessions[0].last_message_at, parsed.last_message_at);
1251 assert_eq!(sessions[0].usage.total_tokens, parsed.usage.total_tokens);
1252 }
1253
1254 #[test]
1255 fn cursor_enrichment_missing_db_changes_nothing() {
1256 let temp = tempfile::tempdir().unwrap();
1257 let parent = cursor_fixture_transcripts(temp.path());
1258
1259 let parsed = agent_session::parse_session_path(&parent).expect("parsed");
1260 let mut sessions = vec![parsed.clone()];
1261 enrich_cursor_sessions_in_home(temp.path(), &mut sessions);
1262
1263 assert_eq!(sessions[0].model, parsed.model);
1264 assert_eq!(sessions[0].usage.total_tokens, parsed.usage.total_tokens);
1265 assert_eq!(sessions[0].cwd, parsed.cwd);
1266 assert_eq!(sessions[0].start_timestamp_ms, parsed.start_timestamp_ms);
1267 }
1268
1269 #[test]
1270 fn cursor_enrichment_reads_while_writer_holds_wal() {
1271 let temp = tempfile::tempdir().unwrap();
1272 write_cursor_state_db_for_test(temp.path());
1273 let parent = cursor_fixture_transcripts(temp.path());
1274
1275 let writer =
1276 rusqlite::Connection::open(cursor_state_db_path(temp.path()).unwrap()).unwrap();
1277 writer.execute_batch("BEGIN IMMEDIATE;").unwrap();
1278 writer
1279 .execute(
1280 "INSERT INTO cursorDiskKV (key, value) VALUES ('agentKv:x', 'held open')",
1281 [],
1282 )
1283 .unwrap();
1284
1285 let mut sessions = vec![agent_session::parse_session_path(&parent).expect("parsed")];
1286 enrich_cursor_sessions_in_home(temp.path(), &mut sessions);
1287 writer.execute_batch("ROLLBACK;").unwrap();
1288
1289 assert_eq!(sessions[0].model.as_deref(), Some("claude-sonnet-4-6"));
1290 assert_eq!(sessions[0].usage.total_tokens, 150);
1291 }
1292
1293 #[test]
1294 fn cursor_composer_data_reads_model_and_workspace() {
1295 let temp = tempfile::tempdir().unwrap();
1296 write_cursor_state_db_for_test(temp.path());
1297 let conn = open_cursor_state_db(&cursor_state_db_path(temp.path()).unwrap()).unwrap();
1298
1299 let pinned = cursor_composer_data(&conn, "abc00000-0000-0000-0000-000000000abc")
1300 .expect("pinned composer");
1301 assert_eq!(pinned.model.as_deref(), Some("claude-sonnet-4-6"));
1302 assert_eq!(pinned.workspace_path.as_deref(), Some("/work/repo"));
1303
1304 let unpinned = cursor_composer_data(&conn, "aaa00000-0000-0000-0000-000000000aaa")
1305 .expect("default composer");
1306 assert_eq!(unpinned.model, None);
1307 assert_eq!(unpinned.workspace_path, None);
1308
1309 assert!(cursor_composer_data(&conn, "not-a-composer").is_none());
1310 }
1311
1312 #[test]
1313 fn cursor_bubble_tokens_sums_by_bounded_range() {
1314 let temp = tempfile::tempdir().unwrap();
1315 write_cursor_state_db_for_test(temp.path());
1316 let conn = open_cursor_state_db(&cursor_state_db_path(temp.path()).unwrap()).unwrap();
1317
1318 let parent = cursor_bubble_tokens(&conn, "abc00000-0000-0000-0000-000000000abc");
1319 assert_eq!(parent.input_tokens, 100);
1320 assert_eq!(parent.output_tokens, 40);
1321 assert_eq!(parent.total_tokens, 140);
1322
1323 let child = cursor_bubble_tokens(&conn, "def00000-0000-0000-0000-000000000def");
1324 assert_eq!(child.total_tokens, 10);
1325
1326 let none = cursor_bubble_tokens(&conn, "aaa00000-0000-0000-0000-000000000aaa");
1327 assert_eq!(none.total_tokens, 0);
1328 }
1329
1330 #[test]
1331 fn cursor_subagent_ids_come_from_directory_layout() {
1332 let temp = tempfile::tempdir().unwrap();
1333 let transcripts = temp
1334 .path()
1335 .join(".cursor/projects/repo/agent-transcripts/abc");
1336 fs::create_dir_all(transcripts.join("subagents")).unwrap();
1337 let parent = transcripts.join("abc.jsonl");
1338 fs::write(&parent, "{}\n").unwrap();
1339 fs::write(transcripts.join("subagents/def.jsonl"), "{}\n").unwrap();
1340 fs::write(transcripts.join("subagents/aaa.jsonl"), "{}\n").unwrap();
1341 fs::write(transcripts.join("subagents/notes.txt"), "x").unwrap();
1342
1343 assert_eq!(cursor_subagent_ids(&parent), vec!["aaa", "def"]);
1344 let no_subagents = temp.path().join("elsewhere/abc.jsonl");
1345 assert!(cursor_subagent_ids(&no_subagents).is_empty());
1346 }
1347
1348 #[test]
1349 fn cursor_composer_header_reads_parent_and_subagent_rows() {
1350 let temp = tempfile::tempdir().unwrap();
1351 write_cursor_state_db_for_test(temp.path());
1352 let conn = open_cursor_state_db(&cursor_state_db_path(temp.path()).unwrap()).unwrap();
1353
1354 let parent = cursor_composer_header(&conn, "abc00000-0000-0000-0000-000000000abc")
1355 .expect("parent header");
1356 assert_eq!(parent.created_at_ms, Some(1_700_000));
1357 assert_eq!(parent.updated_at_ms, Some(1_900_000));
1358
1359 let child = cursor_composer_header(&conn, "def00000-0000-0000-0000-000000000def")
1360 .expect("subagent header");
1361 assert_eq!(child.created_at_ms, Some(1_750_000));
1362 assert_eq!(child.updated_at_ms, None);
1363
1364 assert!(cursor_composer_header(&conn, "not-a-composer").is_none());
1365 }
1366
1367 #[test]
1368 fn count_session_dirs_reports_cursor_root() {
1369 let temp = tempfile::tempdir().unwrap();
1370 assert!(agent_session::count_session_dirs_in_home(temp.path()).is_empty());
1371
1372 let transcripts = temp
1373 .path()
1374 .join(".cursor/projects/repo/agent-transcripts/abc");
1375 fs::create_dir_all(transcripts.join("subagents")).unwrap();
1376 fs::write(transcripts.join("abc.jsonl"), "{}\n").unwrap();
1377 fs::write(transcripts.join("subagents/def.jsonl"), "{}\n").unwrap();
1378
1379 let stats = agent_session::count_session_dirs_in_home(temp.path());
1380 assert_eq!(stats.len(), 1);
1381 assert_eq!(stats[0].agent, agent_session::AGENT_CURSOR);
1382 assert!(stats[0].dir.ends_with(".cursor/projects"));
1383 assert_eq!(stats[0].sessions, 1);
1386 assert_eq!(stats[0].bytes, 3);
1387 }
1388
1389 #[test]
1390 fn cursor_discovery_emits_parent_candidates_only() {
1391 let temp = tempfile::tempdir().unwrap();
1392 let project = temp.path().join(".cursor/projects/repo");
1393 let transcripts = project.join("agent-transcripts/abc");
1394 fs::create_dir_all(transcripts.join("subagents")).unwrap();
1395 fs::create_dir_all(project.join("canvases/node_modules/pkg")).unwrap();
1396 fs::write(transcripts.join("abc.jsonl"), "{}\n").unwrap();
1397 fs::write(transcripts.join("subagents/def.jsonl"), "{}\n").unwrap();
1398 fs::write(project.join("canvases/node_modules/pkg/data.jsonl"), "{}\n").unwrap();
1399
1400 let candidates: Vec<_> = agent_session::discover_session_files_in_home(temp.path())
1401 .into_iter()
1402 .filter(|candidate| candidate.agent == agent_session::AGENT_CURSOR)
1403 .collect();
1404
1405 assert_eq!(candidates.len(), 1);
1406 assert_eq!(candidates[0].path, transcripts.join("abc.jsonl"));
1407 }
1408
1409 #[test]
1410 fn cursor_candidate_updated_tracks_subagent_writes() {
1411 let temp = tempfile::tempdir().unwrap();
1412 let transcripts = temp
1413 .path()
1414 .join(".cursor/projects/repo/agent-transcripts/abc");
1415 fs::create_dir_all(transcripts.join("subagents")).unwrap();
1416 fs::write(transcripts.join("abc.jsonl"), "{}\n").unwrap();
1417 let child = transcripts.join("subagents/def.jsonl");
1418 fs::write(&child, "{}\n").unwrap();
1419
1420 let bumped = std::time::SystemTime::now() + std::time::Duration::from_secs(120);
1422 let handle = fs::File::options().write(true).open(&child).unwrap();
1423 handle.set_modified(bumped).unwrap();
1424
1425 let candidates: Vec<_> = agent_session::discover_session_files_in_home(temp.path())
1426 .into_iter()
1427 .filter(|candidate| candidate.agent == agent_session::AGENT_CURSOR)
1428 .collect();
1429
1430 assert_eq!(candidates.len(), 1);
1431 let parent_mtime = fs::metadata(transcripts.join("abc.jsonl"))
1432 .unwrap()
1433 .modified()
1434 .unwrap();
1435 assert!(candidates[0].updated > parent_mtime);
1436 }
1437
1438 #[test]
1439 fn cursor_duplicate_composer_prefers_real_workspace() {
1440 let temp = tempfile::tempdir().unwrap();
1441 let real = temp
1442 .path()
1443 .join(".cursor/projects/repo/agent-transcripts/abc");
1444 let stale = temp
1445 .path()
1446 .join(".cursor/projects/empty-window/agent-transcripts/abc");
1447 fs::create_dir_all(&real).unwrap();
1448 fs::create_dir_all(&stale).unwrap();
1449 fs::write(real.join("abc.jsonl"), "{}\n").unwrap();
1450 fs::write(stale.join("abc.jsonl"), "{}\n{}\n").unwrap();
1451
1452 let candidates: Vec<_> = agent_session::discover_session_files_in_home(temp.path())
1453 .into_iter()
1454 .filter(|candidate| candidate.agent == agent_session::AGENT_CURSOR)
1455 .collect();
1456
1457 assert_eq!(candidates.len(), 1);
1458 assert_eq!(candidates[0].path, real.join("abc.jsonl"));
1459 }
1460
1461 #[test]
1462 fn codex_state_db_uses_rollout_token_usage() {
1463 let temp = tempfile::tempdir().unwrap();
1464 write_codex_state_db_for_test(temp.path());
1465 let rollout = temp.path().join("session.jsonl");
1466 let mut content = "{}\n".repeat(CODEX_ROLLOUT_TAIL_BYTES as usize / 3 + 1);
1467 content.push_str(concat!(
1468 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}}}}"#,
1469 "\n"
1470 ));
1471 fs::write(&rollout, content).unwrap();
1472 let conn = rusqlite::Connection::open(temp.path().join(".codex/state_5.sqlite")).unwrap();
1473 conn.execute(
1474 "UPDATE threads SET rollout_path = ?1, tokens_used = 999999999",
1475 [rollout.to_string_lossy().as_ref()],
1476 )
1477 .unwrap();
1478
1479 let sessions = codex_state_sessions_in_home(temp.path(), 5);
1480
1481 assert_eq!(sessions[0].usage.input_tokens, 9_200);
1482 assert_eq!(sessions[0].usage.cache_read_tokens, 9_984);
1483 assert_eq!(sessions[0].usage.output_tokens, 11);
1484 assert_eq!(sessions[0].usage.total_tokens, 19_195);
1485 }
1486
1487 #[test]
1488 fn codex_state_db_errors_return_empty_for_jsonl_fallback() {
1489 let temp = tempfile::tempdir().unwrap();
1490 fs::create_dir_all(temp.path().join(".codex")).unwrap();
1491 fs::write(temp.path().join(".codex/state_5.sqlite"), "not sqlite").unwrap();
1492
1493 assert!(codex_state_sessions_in_home(temp.path(), 5).is_empty());
1494 }
1495
1496 #[test]
1497 fn crowded_codex_home_fallback_keeps_only_observed_exec_prompt() {
1498 let temp = tempfile::tempdir().unwrap();
1499 let state_path = write_codex_home(
1500 temp.path(),
1501 &[
1502 "unrelated historical prompt",
1503 "unrelated historical prompt",
1504 "unrelated historical prompt",
1505 "agentsight current run prompt",
1506 "agentsight current run historical prompt",
1507 "unrelated historical prompt",
1508 "unrelated historical prompt",
1509 "unrelated historical prompt",
1510 ],
1511 );
1512 let now = current_epoch_ms();
1513
1514 let rows = vec![
1515 exec_row(
1516 "audit-exec",
1517 now,
1518 "codex",
1519 &format!("{CODEX}agentsight current run prompt"),
1520 ),
1521 exec_row(
1522 "audit-exec-truncated",
1523 now + 1,
1524 "node",
1525 &format!("{CODEX}agentsight current run"),
1526 ),
1527 file_row("audit-file", now + 100, &state_path),
1528 ];
1529
1530 let sessions = observed_sessions_from_audit_rows(&rows);
1531
1532 assert_eq!(sessions.len(), 1);
1533 assert_eq!(
1534 sessions[0].prompt_preview.as_deref(),
1535 Some("agentsight current run prompt")
1536 );
1537 }
1538
1539 #[test]
1540 fn codex_home_fallback_accepts_time_window_when_exec_prompt_is_truncated() {
1541 let temp = tempfile::tempdir().unwrap();
1542 let state_path = write_codex_home(temp.path(), &["agentsight truncated command prompt"]);
1543 let now = current_epoch_ms();
1544
1545 let sessions = observed_sessions_from_audit_rows(&[
1546 exec_row(
1547 "audit-exec",
1548 now,
1549 "codex",
1550 &format!("{CODEX}-c model_provider=\"agentsight-mock"),
1551 ),
1552 file_row("audit-file", now + 100, &state_path),
1553 ]);
1554
1555 assert_eq!(sessions.len(), 1);
1556 assert_eq!(
1557 sessions[0].prompt_preview.as_deref(),
1558 Some("agentsight truncated command prompt")
1559 );
1560 }
1561
1562 #[test]
1563 fn codex_exec_prompt_rows_filter_only_wrapper_duplicates() {
1564 let rows = [
1565 (1_000, "node", "/usr/bin/node /opt/codex/bin/codex exec --skip-git-repo-check agentsight dedupe prompt"),
1566 (1_001, "codex", "/opt/codex/bin/codex exec --skip-git-repo-check agentsight dedupe prompt"),
1567 (2_000, "codex", "/usr/bin/codex exec --skip-git-repo-check agentsight short prompt"),
1568 (3_000, "codex", "/usr/bin/codex exec --skip-git-repo-check agentsight much longer unrelated prompt"),
1569 (10_000, "codex", "/usr/bin/codex exec --skip-git-repo-check agentsight repeated prompt"),
1570 (11_000, "codex", "/usr/bin/codex exec --skip-git-repo-check agentsight repeated prompt"),
1571 (20_000, "codex", "/usr/bin/codex exec --skip-git-repo-check agentsight repeated prompt"),
1572 (21_000, "docker", "docker exec codex exec agentsight should not parse"),
1573 (22_000, "docker", "docker exec container /usr/local/bin/codex exec agentsight should parse once"),
1574 ]
1575 .into_iter()
1576 .enumerate()
1577 .map(|(index, (ts, comm, command))| exec_row(&format!("audit-{index}"), ts, comm, command))
1578 .collect::<Vec<_>>();
1579
1580 let projected = observed_session_prompt_rows(&rows);
1581 assert_eq!(projected[0].comm.as_deref(), Some("node"));
1582 assert_eq!(projected[0].target.as_deref(), Some("/usr/bin/node"));
1583 let prompts = projected
1584 .into_iter()
1585 .map(|row| row.summary)
1586 .collect::<Vec<_>>();
1587
1588 assert_eq!(
1589 prompts,
1590 vec![
1591 Some("agentsight dedupe prompt".to_string()),
1592 Some("agentsight short prompt".to_string()),
1593 Some("agentsight much longer unrelated prompt".to_string()),
1594 Some("agentsight repeated prompt".to_string()),
1595 Some("agentsight repeated prompt".to_string()),
1596 Some("agentsight repeated prompt".to_string()),
1597 Some("agentsight should parse once".to_string()),
1598 ]
1599 );
1600 }
1601
1602 #[test]
1603 fn codex_fallback_time_window_rejects_stale_matching_session() {
1604 let (_temp, path) = create_temp_session_path(agent_session::AGENT_CODEX);
1605 let session = parse_content_for_test(
1606 agent_session::AGENT_CODEX,
1607 &path,
1608 UNIX_EPOCH,
1609 "{\"type\":\"message\",\"content\":\"agentsight repeated prompt\"}\n",
1610 )
1611 .unwrap();
1612
1613 assert!(session_matches_observed_prompt(
1614 &session,
1615 &[ObservedCodexPrompt {
1616 prompt: "agentsight repeated prompt".to_string(),
1617 timestamp_ms: current_epoch_ms(),
1618 pid: Some(42),
1619 native_exec: true,
1620 comm: Some("codex".to_string()),
1621 target: Some("/usr/bin/codex".to_string()),
1622 }]
1623 ));
1624 assert!(!session_is_in_observed_window(
1625 &session,
1626 Some((current_epoch_ms() - 1_000, current_epoch_ms() + 1_000))
1627 ));
1628 }
1629
1630 fn exec_row(id: &str, timestamp_ms: u64, comm: &str, full_command: &str) -> AuditEventRow {
1631 AuditEventRow {
1632 id: id.to_string(),
1633 timestamp_ms,
1634 audit_type: "process".to_string(),
1635 pid: Some(42),
1636 comm: Some(comm.to_string()),
1637 subject: None,
1638 action: Some("exec".to_string()),
1639 target: Some(format!("/usr/bin/{comm}")),
1640 status: Some("observed".to_string()),
1641 summary: None,
1642 details: serde_json::json!({ "full_command": full_command }),
1643 }
1644 }
1645
1646 fn file_row(id: &str, timestamp_ms: u64, path: &Path) -> AuditEventRow {
1647 AuditEventRow {
1648 id: id.to_string(),
1649 timestamp_ms,
1650 audit_type: "file".to_string(),
1651 pid: Some(42),
1652 comm: Some("codex".to_string()),
1653 subject: None,
1654 action: Some("write".to_string()),
1655 target: Some(path.to_string_lossy().to_string()),
1656 status: Some("observed".to_string()),
1657 summary: None,
1658 details: serde_json::json!({ "filepath": path.to_string_lossy() }),
1659 }
1660 }
1661
1662 fn current_epoch_ms() -> u64 {
1663 std::time::SystemTime::now()
1664 .duration_since(UNIX_EPOCH)
1665 .unwrap()
1666 .as_millis() as u64
1667 }
1668
1669 fn write_codex_home(root: &Path, prompts: &[&str]) -> PathBuf {
1670 let codex_home = root.join("codex-home");
1671 let sessions_dir = codex_home.join("sessions/2026/07/14");
1672 fs::create_dir_all(&sessions_dir).unwrap();
1673 for (index, prompt) in prompts.iter().enumerate() {
1674 fs::write(
1675 sessions_dir.join(format!(
1676 "rollout-2026-07-14T00-00-{index:02}-session-{index}.jsonl"
1677 )),
1678 format!(
1679 "{{\"timestamp\":\"2026-07-14T00:00:{index:02}.000Z\",\
1680 \"type\":\"event_msg\",\
1681 \"payload\":{{\"type\":\"user_message\",\"message\":\"{prompt}\"}}}}\n"
1682 ),
1683 )
1684 .unwrap();
1685 }
1686 let state_path = codex_home.join("stat");
1687 fs::write(&state_path, "").unwrap();
1688 state_path
1689 }
1690}