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