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