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