Skip to main content

agentsight_capture/sources/
agent_native.rs

1// SPDX-License-Identifier: MIT
2// Copyright (c) 2026 eunomia-bpf org.
3
4use 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;
28// Local copy of agent_session::AGENT_CURSOR. cargo package verifies this crate
29// against the published agent-session, which lags behind the workspace copy, so
30// production code here cannot reference constants the registry version lacks.
31const 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
191/// Expand a lightweight indexed session only when a caller needs transcript events.
192/// Discovery remains bounded and cheap; the shared cache avoids reparsing unchanged files.
193pub 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
257/// Bound the authorized detail payload without losing the newest interaction.
258/// Previews remain available when older full text falls outside the budget.
259fn 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    // Roll up across delegated runs, or the sessions that delegated most under-report.
545    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    // `discover_sessions` enriches on its way out, but this is a second entry
632    // point into the same data and every `report` subcommand without a --db
633    // goes through here, as does the snapshot the frontend renders.
634    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/// Fixture mirroring Cursor's `state.vscdb` on 3.15.6: one parent composer
1208/// with model, workspace, and legacy token bubbles, one subagent composer with
1209/// a NULL `lastUpdatedAt`, and one composer left on the "default" model.
1210#[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    // Real installs run WAL, which is what lets a read-only connection work
1216    // while Cursor holds a write transaction.
1217    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        // 140 from the parent's bubbles plus 10 from the delegated run.
1454        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        // Real installs have composerData with no composerHeaders row.
1467        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        // Parents only: the subagent file folds into its parent session, so it
1624        // adds neither a session nor bytes.
1625        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        // Bump only the child's mtime well past the parent's.
1661        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}