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