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(dirs::home_dir)
416}
417
418fn non_negative_i64_to_u64(value: i64) -> Option<u64> {
419    u64::try_from(value).ok()
420}
421
422fn system_time_from_ms(value: u64) -> SystemTime {
423    UNIX_EPOCH + Duration::from_millis(value)
424}
425
426fn iso_utc_from_ms(value: u64) -> String {
427    chrono::DateTime::<chrono::Utc>::from_timestamp_millis(value as i64)
428        .map(|dt| dt.to_rfc3339_opts(chrono::SecondsFormat::Millis, true))
429        .unwrap_or_default()
430}
431
432fn short_session_id(id: &str) -> String {
433    let compact = id.rsplit(['/', '\\']).next().unwrap_or(id).trim();
434    if compact.chars().count() <= 12 {
435        return compact.to_string();
436    }
437    let head = compact.chars().take(6).collect::<String>();
438    let tail = compact
439        .chars()
440        .rev()
441        .take(5)
442        .collect::<Vec<_>>()
443        .into_iter()
444        .rev()
445        .collect::<String>();
446    format!("{head}.{tail}")
447}
448
449fn clean_prompt_text(text: &str) -> Option<String> {
450    let text = text.split_whitespace().collect::<Vec<_>>().join(" ");
451    (!text.trim().is_empty()).then(|| text.trim().to_string())
452}
453
454fn view_id(session: &LocalSession) -> String {
455    format!("local:{}:{}", session.agent_type, session.display_id)
456}
457
458pub fn materialized_view(sessions: &[LocalSession]) -> MaterializedView {
459    let mut view = MaterializedView::new();
460    view.set_source(AGENT_NATIVE_SOURCE);
461    import_into_view(&mut view, sessions);
462    view
463}
464
465pub fn import_recent(view: &mut MaterializedView, limit: usize) {
466    let mut sessions = SessionCache::new().discover_cached(limit, Duration::ZERO);
467    // `discover_sessions` enriches on its way out, but this is a second entry
468    // point into the same data and every `report` subcommand without a --db
469    // goes through here, as does the snapshot the frontend renders.
470    enrich_cursor_sessions(&mut sessions);
471    import_into_view(view, &sessions);
472}
473
474pub fn import_into_view(view: &mut MaterializedView, sessions: &[LocalSession]) {
475    for session in sessions {
476        view.upsert_session(&session_row(session));
477        for row in llm_rows(session) {
478            view.apply_llm_call(&row);
479        }
480        for row in token_rows(session) {
481            view.apply_token_usage(&row);
482        }
483        for row in tool_rows(session) {
484            view.apply_tool_call(&row);
485        }
486    }
487}
488
489fn llm_rows(session: &LocalSession) -> Vec<LlmCallRow> {
490    let Some(prompt) = session.prompt_preview.as_ref() else {
491        return Vec::new();
492    };
493    let session_id = view_id(session);
494    let timestamp_ms = session
495        .events
496        .prompts
497        .first()
498        .and_then(|prompt| prompt.ts_ms)
499        .and_then(|ts| u64::try_from(ts).ok())
500        .or(session.start_timestamp_ms)
501        .unwrap_or_else(|| updated_ms(session));
502    let request = serde_json::json!({
503        "prompt": prompt,
504        "prompt_source": AGENT_NATIVE_SOURCE,
505        "session_id": session_id,
506        "agent_type": session.agent_type.as_str(),
507        "path": session.path.to_string_lossy(),
508    });
509
510    if session.model_usage.is_empty() {
511        let model = session
512            .model
513            .clone()
514            .unwrap_or_else(|| session.agent_type.clone());
515        return vec![llm_row_for_session(
516            &format!("{session_id}-{}", sanitize_id(&model)),
517            session,
518            &session_id,
519            timestamp_ms,
520            Some(model),
521            &session.usage,
522            request,
523        )];
524    }
525
526    session
527        .model_usage
528        .iter()
529        .map(|(model, usage)| {
530            llm_row_for_session(
531                &format!("{session_id}-{model}"),
532                session,
533                &session_id,
534                timestamp_ms,
535                Some(model.clone()),
536                usage,
537                request.clone(),
538            )
539        })
540        .collect()
541}
542
543fn llm_row_for_session(
544    id: &str,
545    session: &LocalSession,
546    session_id: &str,
547    timestamp_ms: u64,
548    model: Option<String>,
549    usage: &TokenUsage,
550    request: Value,
551) -> LlmCallRow {
552    LlmCallRow {
553        id: id.to_string(),
554        session_id: Some(session_id.to_string()),
555        conversation_id: session.conversation_id.clone(),
556        start_timestamp_ms: timestamp_ms,
557        end_timestamp_ms: session.end_timestamp_ms,
558        pid: None,
559        comm: Some(session.agent_type.clone()),
560        provider: None,
561        model,
562        call_kind: Some("agent_native_prompt".to_string()),
563        status: "observed".to_string(),
564        error_type: None,
565        finish_reason: None,
566        host: None,
567        path: Some(session.path.to_string_lossy().to_string()),
568        status_code: None,
569        input_tokens: usage.input_tokens,
570        output_tokens: usage.output_tokens,
571        total_tokens: usage.total_tokens,
572        request,
573        response: Value::Null,
574    }
575}
576
577pub fn observed_session_prompt_rows(audit_rows: &[AuditEventRow]) -> Vec<AuditEventRow> {
578    let mut rows = Vec::new();
579    let mut seen = HashSet::new();
580    let observed_exec_prompts = observed_codex_exec_prompts(audit_rows);
581    let mut seen_exec_prompts: Vec<ObservedCodexPrompt> = Vec::new();
582    for observed in observed_exec_prompts {
583        if seen_exec_prompts.iter().any(|seen| {
584            seen.prompt == observed.prompt
585                && timestamps_close(
586                    seen.timestamp_ms,
587                    observed.timestamp_ms,
588                    CODEX_EXEC_DEDUPE_WINDOW_MS,
589                )
590                && (!seen.native_exec || !observed.native_exec)
591        }) {
592            continue;
593        }
594        seen_exec_prompts.push(observed.clone());
595        rows.push(AuditEventRow {
596            id: format!(
597                "audit-codex-exec-prompt-{}-{}",
598                observed.timestamp_ms,
599                observed.pid.unwrap_or(0)
600            ),
601            timestamp_ms: observed.timestamp_ms,
602            audit_type: "llm".to_string(),
603            pid: observed.pid,
604            comm: observed.comm.or_else(|| Some("codex".to_string())),
605            subject: None,
606            action: Some("request".to_string()),
607            target: observed.target,
608            status: Some("observed".to_string()),
609            summary: Some(truncate_text(&observed.prompt, 160)),
610            details: serde_json::json!({
611                "text_content": observed.prompt,
612                "prompt_source": "local",
613            }),
614        });
615    }
616    for row in audit_rows {
617        if row.audit_type == "process" && row.action.as_deref() == Some("exec") {
618            continue;
619        }
620        if row.audit_type != "file" {
621            continue;
622        }
623        let Some(pid) = row.pid else {
624            continue;
625        };
626        let Some(path) = audit_session_path(row) else {
627            continue;
628        };
629        if !seen.insert((path.clone(), pid)) {
630            continue;
631        };
632        let Some(session) = agent_session::parse_session_path(&path) else {
633            continue;
634        };
635        let Some(prompt) = session.prompt_preview.as_ref() else {
636            continue;
637        };
638        rows.push(AuditEventRow {
639            id: format!(
640                "audit-agent-native-prompt-{}-{pid}",
641                sanitize_id(&session.display_id)
642            ),
643            timestamp_ms: row.timestamp_ms,
644            audit_type: "llm".to_string(),
645            pid: Some(pid),
646            comm: row
647                .comm
648                .clone()
649                .or_else(|| Some(session.agent_type.clone())),
650            subject: session.model.clone(),
651            action: Some("request".to_string()),
652            target: Some(path.to_string_lossy().to_string()),
653            status: Some("observed".to_string()),
654            summary: Some(truncate_text(prompt, 160)),
655            details: serde_json::json!({
656                "text_content": prompt,
657                "prompt_source": "local",
658                "session_id": view_id(&session),
659                "conversation_id": session.conversation_id.as_deref(),
660                "agent_type": session.agent_type,
661            }),
662        });
663    }
664    rows
665}
666
667pub fn observed_sessions_from_audit_rows(audit_rows: &[AuditEventRow]) -> Vec<LocalSession> {
668    let mut direct_paths = HashSet::new();
669    let mut codex_session_dirs = HashSet::new();
670    let observed_codex_prompts = observed_codex_exec_prompts(audit_rows);
671    let observed_codex_exec = observed_codex_exec_command(audit_rows);
672    let observed_window = observed_audit_window_ms(audit_rows);
673
674    for row in audit_rows.iter().filter(|row| row.audit_type == "file") {
675        for path in audit_file_paths(row) {
676            if let Some(session_path) =
677                agent_session::session_log_path_from_str(path.to_string_lossy().as_ref())
678            {
679                direct_paths.insert(session_path);
680            }
681            if let Some(dir) = observed_codex_sessions_dir(&path) {
682                codex_session_dirs.insert(dir);
683            }
684        }
685    }
686
687    let mut candidates = Vec::new();
688    for path in direct_paths {
689        if let Some(candidate) = agent_session::session_candidate_from_path(&path) {
690            candidates.push((candidate, false));
691        }
692    }
693    for dir in codex_session_dirs {
694        let dir_candidates =
695            agent_session::discover_session_files_in_dir(agent_session::AGENT_CODEX, &dir);
696        candidates.extend(
697            dir_candidates
698                .into_iter()
699                .map(|candidate| (candidate, true)),
700        );
701    }
702    candidates.sort_by_key(|(candidate, _)| Reverse(candidate.updated));
703
704    let mut seen_paths = HashSet::new();
705    let mut seen_sessions = HashSet::new();
706    let mut sessions = Vec::new();
707    for (candidate, is_codex_dir_fallback) in candidates.into_iter().take(75) {
708        if !seen_paths.insert(candidate.path.clone()) {
709            continue;
710        }
711        let Some(session) = agent_session::parse_session_file(&candidate) else {
712            continue;
713        };
714        if is_codex_dir_fallback && !observed_codex_exec {
715            continue;
716        }
717        if is_codex_dir_fallback
718            && !observed_codex_prompts.is_empty()
719            && !session_matches_observed_prompt(&session, &observed_codex_prompts)
720        {
721            continue;
722        }
723        if is_codex_dir_fallback && !session_is_in_observed_window(&session, observed_window) {
724            continue;
725        }
726        if seen_sessions.insert(session.display_id.clone()) {
727            sessions.push(session);
728        }
729    }
730    sessions
731}
732
733fn observed_codex_exec_prompts(audit_rows: &[AuditEventRow]) -> Vec<ObservedCodexPrompt> {
734    let prompts = audit_rows
735        .iter()
736        .filter(|row| row.audit_type == "process" && row.action.as_deref() == Some("exec"))
737        .filter_map(|row| {
738            let prompt = row
739                .details
740                .get("full_command")
741                .and_then(Value::as_str)
742                .and_then(codex_exec_prompt_from_command)?;
743            Some(ObservedCodexPrompt {
744                prompt,
745                timestamp_ms: row.timestamp_ms,
746                pid: row.pid,
747                native_exec: looks_like_native_codex_exec(row),
748                comm: row.comm.clone(),
749                target: row.target.clone(),
750            })
751        })
752        .collect::<Vec<_>>();
753    prompts
754        .iter()
755        .filter(|candidate| {
756            !prompts
757                .iter()
758                .any(|other| is_nearby_longer_prefix(candidate, other))
759        })
760        .cloned()
761        .collect()
762}
763
764fn observed_codex_exec_command(audit_rows: &[AuditEventRow]) -> bool {
765    audit_rows
766        .iter()
767        .filter(|row| row.audit_type == "process" && row.action.as_deref() == Some("exec"))
768        .filter_map(|row| row.details.get("full_command").and_then(Value::as_str))
769        .any(|command| codex_exec_command_tail(command).is_some())
770}
771
772fn codex_exec_prompt_from_command(command: &str) -> Option<String> {
773    agent_session::codex_exec_prompt(&codex_exec_command_tail(command)?)
774}
775
776fn codex_exec_command_tail(command: &str) -> Option<String> {
777    let tokens = command.split_whitespace().collect::<Vec<_>>();
778    let index = tokens.windows(2).enumerate().find_map(|(index, tokens)| {
779        (is_codex_executable_token(tokens[0], index == 0) && tokens[1] == "exec").then_some(index)
780    })?;
781    Some(tokens[index..].join(" "))
782}
783
784fn looks_like_native_codex_exec(row: &AuditEventRow) -> bool {
785    row.comm.as_deref() == Some("codex")
786        && row
787            .target
788            .as_deref()
789            .is_some_and(|target| is_codex_executable_token(target, true))
790}
791
792fn is_codex_executable_token(token: &str, allow_bare: bool) -> bool {
793    let token = token.trim_matches(|ch| matches!(ch, '"' | '\''));
794    (allow_bare && token == "codex")
795        || token.contains('/')
796            && Path::new(token).file_name().and_then(|name| name.to_str()) == Some("codex")
797}
798
799fn is_nearby_longer_prefix(candidate: &ObservedCodexPrompt, other: &ObservedCodexPrompt) -> bool {
800    other.prompt.len() > candidate.prompt.len()
801        && other.prompt.starts_with(candidate.prompt.as_str())
802        && timestamps_close(
803            candidate.timestamp_ms,
804            other.timestamp_ms,
805            CODEX_EXEC_DEDUPE_WINDOW_MS,
806        )
807        && (!candidate.native_exec || !other.native_exec)
808}
809
810fn timestamps_close(left: u64, right: u64, window_ms: u64) -> bool {
811    left.abs_diff(right) <= window_ms
812}
813
814fn session_matches_observed_prompt(
815    session: &LocalSession,
816    prompts: &[ObservedCodexPrompt],
817) -> bool {
818    let Some(preview) = session.prompt_preview.as_deref() else {
819        return false;
820    };
821    prompts.iter().any(|observed| {
822        let prompt = observed.prompt.as_str();
823        prompt_texts_overlap(prompt, preview)
824    })
825}
826
827fn prompt_texts_overlap(left: &str, right: &str) -> bool {
828    left == right || left.starts_with(right) || right.starts_with(left)
829}
830
831fn observed_audit_window_ms(audit_rows: &[AuditEventRow]) -> Option<(u64, u64)> {
832    let min = audit_rows.iter().map(|row| row.timestamp_ms).min()?;
833    let max = audit_rows.iter().map(|row| row.timestamp_ms).max()?;
834    Some((
835        min.saturating_sub(CODEX_FALLBACK_TIME_SLOP_MS),
836        max.saturating_add(CODEX_FALLBACK_TIME_SLOP_MS),
837    ))
838}
839
840fn session_is_in_observed_window(session: &LocalSession, window: Option<(u64, u64)>) -> bool {
841    let Some((min, max)) = window else {
842        return true;
843    };
844    let updated = updated_ms(session);
845    updated >= min && updated <= max
846}
847
848fn audit_session_path(row: &AuditEventRow) -> Option<PathBuf> {
849    row.target
850        .as_deref()
851        .and_then(agent_session::session_log_path_from_str)
852        .or_else(|| {
853            row.details
854                .get("filepath")
855                .and_then(Value::as_str)
856                .and_then(agent_session::session_log_path_from_str)
857        })
858        .or_else(|| {
859            row.details
860                .get("path")
861                .and_then(Value::as_str)
862                .and_then(agent_session::session_log_path_from_str)
863        })
864}
865
866fn audit_file_paths(row: &AuditEventRow) -> Vec<PathBuf> {
867    [
868        row.target.as_deref(),
869        row.details.get("filepath").and_then(Value::as_str),
870        row.details.get("path").and_then(Value::as_str),
871        row.details.get("fd_target").and_then(Value::as_str),
872    ]
873    .into_iter()
874    .flatten()
875    .filter_map(|raw| {
876        let path = PathBuf::from(raw.trim().trim_end_matches(" (deleted)"));
877        path.is_absolute().then_some(path)
878    })
879    .collect()
880}
881
882fn observed_codex_sessions_dir(path: &Path) -> Option<PathBuf> {
883    if !looks_like_codex_home_file(path) {
884        return None;
885    }
886    let home = path.parent()?;
887    let sessions = home.join("sessions");
888    sessions.is_dir().then_some(sessions)
889}
890
891fn looks_like_codex_home_file(path: &Path) -> bool {
892    let Some(name) = path.file_name().and_then(|name| name.to_str()) else {
893        return false;
894    };
895    (name.starts_with("state_")
896        || name.starts_with("logs_")
897        || matches!(name, "config.toml" | "auth.json" | "stat"))
898        && path
899            .parent()
900            .is_some_and(|parent| parent.join("sessions").is_dir())
901}
902
903fn session_row(session: &LocalSession) -> SessionRow {
904    let updated_ms = updated_ms(session);
905    SessionRow {
906        id: view_id(session),
907        agent_type: session.agent_type.clone(),
908        start_timestamp_ms: session
909            .start_timestamp_ms
910            .unwrap_or_else(|| updated_ms.saturating_sub(session.duration_ms)),
911        end_timestamp_ms: session.end_timestamp_ms.or(Some(updated_ms)),
912        status: "observed".to_string(),
913        model: session.model.clone(),
914        input_tokens: session.usage.input_tokens,
915        output_tokens: session.usage.output_tokens,
916        total_tokens: session.usage.total_tokens,
917        view_source: AGENT_NATIVE_SOURCE.to_string(),
918        confidence: Some(0.95),
919        attributes: serde_json::json!({
920            "session_id": session.session_id.clone(),
921            "conversation_id": session.conversation_id.as_deref(),
922            "path": session.path.to_string_lossy(),
923            "display_id": session.display_id,
924            "prompt_preview": session.prompt_preview.clone(),
925            "cwd": session.cwd.clone(),
926            "last_message_at": session.last_message_at.clone(),
927            "files": session.files,
928        }),
929    }
930}
931
932fn token_rows(session: &LocalSession) -> Vec<TokenUsageRow> {
933    let session_id = view_id(session);
934    session
935        .model_usage
936        .iter()
937        .filter(|(_, usage)| usage.total_tokens > 0)
938        .map(|(model, usage)| TokenUsageRow {
939            id: format!("token-{session_id}-{}", sanitize_id(model)),
940            llm_call_id: format!("{session_id}-{model}"),
941            timestamp_ms: updated_ms(session),
942            pid: None,
943            comm: Some(session.agent_type.clone()),
944            provider: None,
945            model: Some(model.clone()),
946            input_tokens: usage.input_tokens,
947            output_tokens: usage.output_tokens,
948            cache_creation_tokens: usage.cache_creation_tokens,
949            cache_read_tokens: usage.cache_read_tokens,
950            total_tokens: usage.total_tokens,
951            source: AGENT_NATIVE_SOURCE.to_string(),
952            view_source: AGENT_NATIVE_SOURCE.to_string(),
953            confidence: Some(0.95),
954        })
955        .collect()
956}
957
958fn tool_rows(session: &LocalSession) -> Vec<ToolCallRow> {
959    let session_id = view_id(session);
960    let timestamp_ms = updated_ms(session);
961    let mut rows = Vec::new();
962    for (tool, count) in &session.tools {
963        for index in 0..*count {
964            rows.push(ToolCallRow {
965                id: format!("tool-{session_id}-{}-{index}", sanitize_id(tool)),
966                session_id: Some(session_id.clone()),
967                conversation_id: session.conversation_id.clone(),
968                timestamp_ms,
969                tool_name: Some(tool.clone()),
970                tool_call_id: None,
971                start_timestamp_ms: Some(timestamp_ms),
972                end_timestamp_ms: Some(timestamp_ms),
973                duration_ms: None,
974                status: Some("observed".to_string()),
975                input: serde_json::json!({}),
976                output: serde_json::json!({}),
977                related_pid: None,
978                related_event_id: None,
979                view_source: AGENT_NATIVE_SOURCE.to_string(),
980                confidence: Some(0.95),
981            });
982        }
983    }
984    rows
985}
986
987fn updated_ms(session: &LocalSession) -> u64 {
988    session
989        .updated
990        .duration_since(UNIX_EPOCH)
991        .unwrap_or_default()
992        .as_millis() as u64
993}
994
995fn matches_filter(
996    session: &LocalSession,
997    pid_filter: Option<u32>,
998    text_filter: Option<&str>,
999) -> bool {
1000    if pid_filter.is_some() {
1001        return true;
1002    }
1003    let Some(filter) = text_filter else {
1004        return true;
1005    };
1006    let filter = filter.to_ascii_lowercase();
1007    session.agent_type.to_ascii_lowercase().contains(&filter)
1008        || session
1009            .prompt_preview
1010            .as_ref()
1011            .is_some_and(|prompt| prompt.to_ascii_lowercase().contains(&filter))
1012        || session
1013            .model
1014            .as_ref()
1015            .is_some_and(|model| model.to_ascii_lowercase().contains(&filter))
1016        || session
1017            .path
1018            .to_string_lossy()
1019            .to_ascii_lowercase()
1020            .contains(&filter)
1021}
1022
1023#[cfg(any(test, feature = "test-support"))]
1024pub fn create_temp_session_path(agent: &str) -> (tempfile::TempDir, PathBuf) {
1025    let temp = tempfile::tempdir().unwrap();
1026    let path = agent_session::fixture_session_path(agent, temp.path()).unwrap();
1027    fs::create_dir_all(path.parent().unwrap()).unwrap();
1028    fs::write(&path, "{}\n").unwrap();
1029    (temp, path)
1030}
1031
1032#[cfg(any(test, feature = "test-support"))]
1033pub fn parse_content_for_test(
1034    agent: &str,
1035    path: &std::path::Path,
1036    updated: std::time::SystemTime,
1037    content: &str,
1038) -> Option<LocalSession> {
1039    agent_session::parse_session_content(agent, path, updated, content)
1040}
1041
1042/// Fixture mirroring Cursor's `state.vscdb` on 3.15.6: one parent composer
1043/// with model, workspace, and legacy token bubbles, one subagent composer with
1044/// a NULL `lastUpdatedAt`, and one composer left on the "default" model.
1045#[cfg(any(test, feature = "test-support"))]
1046pub fn write_cursor_state_db_for_test(home: &Path) {
1047    let db_dir = home.join("Library/Application Support/Cursor/User/globalStorage");
1048    fs::create_dir_all(&db_dir).unwrap();
1049    let conn = rusqlite::Connection::open(db_dir.join("state.vscdb")).unwrap();
1050    // Real installs run WAL, which is what lets a read-only connection work
1051    // while Cursor holds a write transaction.
1052    conn.pragma_update(None, "journal_mode", "WAL").unwrap();
1053    conn.execute_batch(
1054        r#"CREATE TABLE composerHeaders (composerId TEXT PRIMARY KEY, workspaceId TEXT,
1055            createdAt INTEGER, lastUpdatedAt INTEGER, isArchived INTEGER,
1056            isSubagent INTEGER, recency INTEGER, checkpointAt INTEGER, value TEXT);
1057        CREATE TABLE cursorDiskKV (key TEXT UNIQUE ON CONFLICT REPLACE, value BLOB);
1058        INSERT INTO composerHeaders (composerId, workspaceId, createdAt, lastUpdatedAt, isSubagent)
1059        VALUES
1060        ('abc00000-0000-0000-0000-000000000abc', 'ws1', 1700000, 1900000, 0),
1061        ('def00000-0000-0000-0000-000000000def', 'ws1', 1750000, NULL, 1),
1062        ('aaa00000-0000-0000-0000-000000000aaa', 'ws2', 1600000, 1650000, 0);
1063        INSERT INTO cursorDiskKV (key, value) VALUES
1064        ('composerData:bbb00000-0000-0000-0000-000000000bbb',
1065         '{"composerId":"bbb00000-0000-0000-0000-000000000bbb","modelConfig":{"modelName":"claude-4.6-sonnet-medium-thinking","maxMode":false}}'),
1066        ('composerData:abc00000-0000-0000-0000-000000000abc',
1067         '{"composerId":"abc00000-0000-0000-0000-000000000abc","modelConfig":{"modelName":"claude-sonnet-4-6","maxMode":false},"workspaceIdentifier":{"uri":{"fsPath":"/work/repo"}}}'),
1068        ('composerData:aaa00000-0000-0000-0000-000000000aaa',
1069         '{"composerId":"aaa00000-0000-0000-0000-000000000aaa","modelConfig":{"modelName":"default","maxMode":false}}'),
1070        ('bubbleId:abc00000-0000-0000-0000-000000000abc:b1',
1071         '{"tokenCount":{"inputTokens":100,"outputTokens":40}}'),
1072        ('bubbleId:abc00000-0000-0000-0000-000000000abc:b2',
1073         '{"tokenCount":{"inputTokens":0,"outputTokens":0}}'),
1074        ('bubbleId:def00000-0000-0000-0000-000000000def:b1',
1075         '{"tokenCount":{"inputTokens":7,"outputTokens":3}}');"#,
1076    )
1077    .unwrap();
1078}
1079
1080#[cfg(any(test, feature = "test-support"))]
1081pub fn write_codex_state_db_for_test(home: &Path) {
1082    let codex_dir = home.join(".codex");
1083    fs::create_dir_all(&codex_dir).unwrap();
1084    let conn = rusqlite::Connection::open(codex_dir.join("state_5.sqlite")).unwrap();
1085    conn.execute_batch(
1086        "CREATE TABLE threads (
1087            id TEXT PRIMARY KEY,
1088            rollout_path TEXT,
1089            model TEXT,
1090            tokens_used INTEGER NOT NULL DEFAULT 0,
1091            preview TEXT,
1092            cwd TEXT,
1093            created_at_ms INTEGER,
1094            updated_at_ms INTEGER
1095        );
1096        INSERT INTO threads
1097        (id, rollout_path, model, tokens_used, preview, cwd, created_at_ms, updated_at_ms)
1098        VALUES
1099        ('019f49ca-54e7-7a91-82e7-a52b53cfd456', '/tmp/session.jsonl', 'gpt-web-ci', 33, 'web state prompt', '/work/repo', 1800000, 1900000);",
1100    )
1101    .unwrap();
1102}
1103
1104#[cfg(test)]
1105mod tests {
1106    use super::*;
1107
1108    const CODEX: &str = "/usr/bin/codex exec --skip-git-repo-check ";
1109
1110    #[test]
1111    fn agent_native_prompt_produces_llm_call_row() {
1112        let (_temp, path) = create_temp_session_path(agent_session::AGENT_CODEX);
1113        let session = parse_content_for_test(
1114            agent_session::AGENT_CODEX,
1115            &path,
1116            UNIX_EPOCH,
1117            "{\"type\":\"message\",\"content\":\"agentsight local codex prompt\"}\n",
1118        )
1119        .unwrap();
1120
1121        let view = materialized_view(&[session]);
1122        let rows = view.llm_call_rows(10);
1123
1124        assert_eq!(rows.len(), 1);
1125        assert_eq!(rows[0].comm.as_deref(), Some(agent_session::AGENT_CODEX));
1126        assert_eq!(
1127            rows[0].request.get("prompt").and_then(Value::as_str),
1128            Some("agentsight local codex prompt")
1129        );
1130    }
1131
1132    #[test]
1133    fn codex_state_db_produces_indexed_session_metadata() {
1134        let temp = tempfile::tempdir().unwrap();
1135        write_codex_state_db_for_test(temp.path());
1136
1137        let sessions = codex_state_sessions_in_home(temp.path(), 5);
1138
1139        assert_eq!(sessions.len(), 1);
1140        assert_eq!(sessions[0].display_id, "codex:019f49.fd456");
1141        assert_eq!(sessions[0].model.as_deref(), Some("gpt-web-ci"));
1142        assert_eq!(sessions[0].usage.total_tokens, 33);
1143        assert_eq!(
1144            sessions[0].prompt_preview.as_deref(),
1145            Some("web state prompt")
1146        );
1147        assert_eq!(sessions[0].cwd.as_deref(), Some("/work/repo"));
1148        assert_eq!(
1149            sessions[0].last_message_at.as_deref(),
1150            Some("1970-01-01T00:31:40.000Z")
1151        );
1152    }
1153
1154    #[test]
1155    fn cursor_state_db_path_checks_platform_layouts() {
1156        let temp = tempfile::tempdir().unwrap();
1157        assert!(cursor_state_db_path(temp.path()).is_none());
1158        write_cursor_state_db_for_test(temp.path());
1159        let path = cursor_state_db_path(temp.path()).unwrap();
1160        assert!(path.ends_with("Cursor/User/globalStorage/state.vscdb"));
1161    }
1162
1163    #[test]
1164    fn cursor_state_db_open_is_read_only_and_fails_closed() {
1165        let temp = tempfile::tempdir().unwrap();
1166        assert!(open_cursor_state_db(&temp.path().join("missing.vscdb")).is_none());
1167        write_cursor_state_db_for_test(temp.path());
1168        let conn = open_cursor_state_db(&cursor_state_db_path(temp.path()).unwrap()).unwrap();
1169        let denied = conn.execute(
1170            "INSERT INTO cursorDiskKV (key, value) VALUES ('x', 'y')",
1171            [],
1172        );
1173        assert!(denied.is_err());
1174    }
1175
1176    fn cursor_fixture_transcripts(home: &Path) -> PathBuf {
1177        let transcripts = home
1178            .join(".cursor/projects/repo/agent-transcripts/abc00000-0000-0000-0000-000000000abc");
1179        fs::create_dir_all(transcripts.join("subagents")).unwrap();
1180        let parent = transcripts.join("abc00000-0000-0000-0000-000000000abc.jsonl");
1181        fs::write(
1182            &parent,
1183            r#"{"role":"user","message":{"content":[{"type":"text","text":"check the build"}]}}"#,
1184        )
1185        .unwrap();
1186        fs::write(
1187            transcripts.join("subagents/def00000-0000-0000-0000-000000000def.jsonl"),
1188            "{}\n",
1189        )
1190        .unwrap();
1191        parent
1192    }
1193
1194    #[test]
1195    fn cursor_enrichment_fills_metadata_and_rolls_up_subagents() {
1196        let temp = tempfile::tempdir().unwrap();
1197        write_cursor_state_db_for_test(temp.path());
1198        let parent = cursor_fixture_transcripts(temp.path());
1199
1200        let mut sessions = vec![agent_session::parse_session_path(&parent).expect("parsed")];
1201        enrich_cursor_sessions_in_home(temp.path(), &mut sessions);
1202
1203        let session = &sessions[0];
1204        assert_eq!(session.model.as_deref(), Some("claude-sonnet-4-6"));
1205        assert_eq!(session.start_timestamp_ms, Some(1_700_000));
1206        assert_eq!(session.end_timestamp_ms, Some(1_900_000));
1207        assert_eq!(session.duration_ms, 200_000);
1208        assert_eq!(session.cwd.as_deref(), Some("/work/repo"));
1209        assert_eq!(
1210            session.last_message_at.as_deref(),
1211            Some("1970-01-01T00:31:40.000Z")
1212        );
1213        // 140 from the parent's bubbles plus 10 from the delegated run.
1214        assert_eq!(session.usage.total_tokens, 150);
1215        assert_eq!(
1216            session
1217                .model_usage
1218                .get("claude-sonnet-4-6")
1219                .map(|usage| usage.total_tokens),
1220            Some(150)
1221        );
1222    }
1223
1224    #[test]
1225    fn cursor_enrichment_reads_model_when_header_row_is_missing() {
1226        // Real installs have composerData with no composerHeaders row.
1227        let temp = tempfile::tempdir().unwrap();
1228        write_cursor_state_db_for_test(temp.path());
1229        let transcripts = temp
1230            .path()
1231            .join(".cursor/projects/repo/agent-transcripts/bbb00000-0000-0000-0000-000000000bbb");
1232        fs::create_dir_all(&transcripts).unwrap();
1233        let parent = transcripts.join("bbb00000-0000-0000-0000-000000000bbb.jsonl");
1234        fs::write(
1235            &parent,
1236            r#"{"role":"user","message":{"content":[{"type":"text","text":"check the build"}]}}"#,
1237        )
1238        .unwrap();
1239
1240        let parsed = agent_session::parse_session_path(&parent).expect("parsed");
1241        let mut sessions = vec![parsed.clone()];
1242        enrich_cursor_sessions_in_home(temp.path(), &mut sessions);
1243
1244        assert_eq!(
1245            sessions[0].model.as_deref(),
1246            Some("claude-4.6-sonnet-medium-thinking")
1247        );
1248        assert_eq!(sessions[0].start_timestamp_ms, parsed.start_timestamp_ms);
1249        assert_eq!(sessions[0].end_timestamp_ms, parsed.end_timestamp_ms);
1250        assert_eq!(sessions[0].last_message_at, parsed.last_message_at);
1251        assert_eq!(sessions[0].usage.total_tokens, parsed.usage.total_tokens);
1252    }
1253
1254    #[test]
1255    fn cursor_enrichment_missing_db_changes_nothing() {
1256        let temp = tempfile::tempdir().unwrap();
1257        let parent = cursor_fixture_transcripts(temp.path());
1258
1259        let parsed = agent_session::parse_session_path(&parent).expect("parsed");
1260        let mut sessions = vec![parsed.clone()];
1261        enrich_cursor_sessions_in_home(temp.path(), &mut sessions);
1262
1263        assert_eq!(sessions[0].model, parsed.model);
1264        assert_eq!(sessions[0].usage.total_tokens, parsed.usage.total_tokens);
1265        assert_eq!(sessions[0].cwd, parsed.cwd);
1266        assert_eq!(sessions[0].start_timestamp_ms, parsed.start_timestamp_ms);
1267    }
1268
1269    #[test]
1270    fn cursor_enrichment_reads_while_writer_holds_wal() {
1271        let temp = tempfile::tempdir().unwrap();
1272        write_cursor_state_db_for_test(temp.path());
1273        let parent = cursor_fixture_transcripts(temp.path());
1274
1275        let writer =
1276            rusqlite::Connection::open(cursor_state_db_path(temp.path()).unwrap()).unwrap();
1277        writer.execute_batch("BEGIN IMMEDIATE;").unwrap();
1278        writer
1279            .execute(
1280                "INSERT INTO cursorDiskKV (key, value) VALUES ('agentKv:x', 'held open')",
1281                [],
1282            )
1283            .unwrap();
1284
1285        let mut sessions = vec![agent_session::parse_session_path(&parent).expect("parsed")];
1286        enrich_cursor_sessions_in_home(temp.path(), &mut sessions);
1287        writer.execute_batch("ROLLBACK;").unwrap();
1288
1289        assert_eq!(sessions[0].model.as_deref(), Some("claude-sonnet-4-6"));
1290        assert_eq!(sessions[0].usage.total_tokens, 150);
1291    }
1292
1293    #[test]
1294    fn cursor_composer_data_reads_model_and_workspace() {
1295        let temp = tempfile::tempdir().unwrap();
1296        write_cursor_state_db_for_test(temp.path());
1297        let conn = open_cursor_state_db(&cursor_state_db_path(temp.path()).unwrap()).unwrap();
1298
1299        let pinned = cursor_composer_data(&conn, "abc00000-0000-0000-0000-000000000abc")
1300            .expect("pinned composer");
1301        assert_eq!(pinned.model.as_deref(), Some("claude-sonnet-4-6"));
1302        assert_eq!(pinned.workspace_path.as_deref(), Some("/work/repo"));
1303
1304        let unpinned = cursor_composer_data(&conn, "aaa00000-0000-0000-0000-000000000aaa")
1305            .expect("default composer");
1306        assert_eq!(unpinned.model, None);
1307        assert_eq!(unpinned.workspace_path, None);
1308
1309        assert!(cursor_composer_data(&conn, "not-a-composer").is_none());
1310    }
1311
1312    #[test]
1313    fn cursor_bubble_tokens_sums_by_bounded_range() {
1314        let temp = tempfile::tempdir().unwrap();
1315        write_cursor_state_db_for_test(temp.path());
1316        let conn = open_cursor_state_db(&cursor_state_db_path(temp.path()).unwrap()).unwrap();
1317
1318        let parent = cursor_bubble_tokens(&conn, "abc00000-0000-0000-0000-000000000abc");
1319        assert_eq!(parent.input_tokens, 100);
1320        assert_eq!(parent.output_tokens, 40);
1321        assert_eq!(parent.total_tokens, 140);
1322
1323        let child = cursor_bubble_tokens(&conn, "def00000-0000-0000-0000-000000000def");
1324        assert_eq!(child.total_tokens, 10);
1325
1326        let none = cursor_bubble_tokens(&conn, "aaa00000-0000-0000-0000-000000000aaa");
1327        assert_eq!(none.total_tokens, 0);
1328    }
1329
1330    #[test]
1331    fn cursor_subagent_ids_come_from_directory_layout() {
1332        let temp = tempfile::tempdir().unwrap();
1333        let transcripts = temp
1334            .path()
1335            .join(".cursor/projects/repo/agent-transcripts/abc");
1336        fs::create_dir_all(transcripts.join("subagents")).unwrap();
1337        let parent = transcripts.join("abc.jsonl");
1338        fs::write(&parent, "{}\n").unwrap();
1339        fs::write(transcripts.join("subagents/def.jsonl"), "{}\n").unwrap();
1340        fs::write(transcripts.join("subagents/aaa.jsonl"), "{}\n").unwrap();
1341        fs::write(transcripts.join("subagents/notes.txt"), "x").unwrap();
1342
1343        assert_eq!(cursor_subagent_ids(&parent), vec!["aaa", "def"]);
1344        let no_subagents = temp.path().join("elsewhere/abc.jsonl");
1345        assert!(cursor_subagent_ids(&no_subagents).is_empty());
1346    }
1347
1348    #[test]
1349    fn cursor_composer_header_reads_parent_and_subagent_rows() {
1350        let temp = tempfile::tempdir().unwrap();
1351        write_cursor_state_db_for_test(temp.path());
1352        let conn = open_cursor_state_db(&cursor_state_db_path(temp.path()).unwrap()).unwrap();
1353
1354        let parent = cursor_composer_header(&conn, "abc00000-0000-0000-0000-000000000abc")
1355            .expect("parent header");
1356        assert_eq!(parent.created_at_ms, Some(1_700_000));
1357        assert_eq!(parent.updated_at_ms, Some(1_900_000));
1358
1359        let child = cursor_composer_header(&conn, "def00000-0000-0000-0000-000000000def")
1360            .expect("subagent header");
1361        assert_eq!(child.created_at_ms, Some(1_750_000));
1362        assert_eq!(child.updated_at_ms, None);
1363
1364        assert!(cursor_composer_header(&conn, "not-a-composer").is_none());
1365    }
1366
1367    #[test]
1368    fn count_session_dirs_reports_cursor_root() {
1369        let temp = tempfile::tempdir().unwrap();
1370        assert!(agent_session::count_session_dirs_in_home(temp.path()).is_empty());
1371
1372        let transcripts = temp
1373            .path()
1374            .join(".cursor/projects/repo/agent-transcripts/abc");
1375        fs::create_dir_all(transcripts.join("subagents")).unwrap();
1376        fs::write(transcripts.join("abc.jsonl"), "{}\n").unwrap();
1377        fs::write(transcripts.join("subagents/def.jsonl"), "{}\n").unwrap();
1378
1379        let stats = agent_session::count_session_dirs_in_home(temp.path());
1380        assert_eq!(stats.len(), 1);
1381        assert_eq!(stats[0].agent, agent_session::AGENT_CURSOR);
1382        assert!(stats[0].dir.ends_with(".cursor/projects"));
1383        // Parents only: the subagent file folds into its parent session, so it
1384        // adds neither a session nor bytes.
1385        assert_eq!(stats[0].sessions, 1);
1386        assert_eq!(stats[0].bytes, 3);
1387    }
1388
1389    #[test]
1390    fn cursor_discovery_emits_parent_candidates_only() {
1391        let temp = tempfile::tempdir().unwrap();
1392        let project = temp.path().join(".cursor/projects/repo");
1393        let transcripts = project.join("agent-transcripts/abc");
1394        fs::create_dir_all(transcripts.join("subagents")).unwrap();
1395        fs::create_dir_all(project.join("canvases/node_modules/pkg")).unwrap();
1396        fs::write(transcripts.join("abc.jsonl"), "{}\n").unwrap();
1397        fs::write(transcripts.join("subagents/def.jsonl"), "{}\n").unwrap();
1398        fs::write(project.join("canvases/node_modules/pkg/data.jsonl"), "{}\n").unwrap();
1399
1400        let candidates: Vec<_> = agent_session::discover_session_files_in_home(temp.path())
1401            .into_iter()
1402            .filter(|candidate| candidate.agent == agent_session::AGENT_CURSOR)
1403            .collect();
1404
1405        assert_eq!(candidates.len(), 1);
1406        assert_eq!(candidates[0].path, transcripts.join("abc.jsonl"));
1407    }
1408
1409    #[test]
1410    fn cursor_candidate_updated_tracks_subagent_writes() {
1411        let temp = tempfile::tempdir().unwrap();
1412        let transcripts = temp
1413            .path()
1414            .join(".cursor/projects/repo/agent-transcripts/abc");
1415        fs::create_dir_all(transcripts.join("subagents")).unwrap();
1416        fs::write(transcripts.join("abc.jsonl"), "{}\n").unwrap();
1417        let child = transcripts.join("subagents/def.jsonl");
1418        fs::write(&child, "{}\n").unwrap();
1419
1420        // Bump only the child's mtime well past the parent's.
1421        let bumped = std::time::SystemTime::now() + std::time::Duration::from_secs(120);
1422        let handle = fs::File::options().write(true).open(&child).unwrap();
1423        handle.set_modified(bumped).unwrap();
1424
1425        let candidates: Vec<_> = agent_session::discover_session_files_in_home(temp.path())
1426            .into_iter()
1427            .filter(|candidate| candidate.agent == agent_session::AGENT_CURSOR)
1428            .collect();
1429
1430        assert_eq!(candidates.len(), 1);
1431        let parent_mtime = fs::metadata(transcripts.join("abc.jsonl"))
1432            .unwrap()
1433            .modified()
1434            .unwrap();
1435        assert!(candidates[0].updated > parent_mtime);
1436    }
1437
1438    #[test]
1439    fn cursor_duplicate_composer_prefers_real_workspace() {
1440        let temp = tempfile::tempdir().unwrap();
1441        let real = temp
1442            .path()
1443            .join(".cursor/projects/repo/agent-transcripts/abc");
1444        let stale = temp
1445            .path()
1446            .join(".cursor/projects/empty-window/agent-transcripts/abc");
1447        fs::create_dir_all(&real).unwrap();
1448        fs::create_dir_all(&stale).unwrap();
1449        fs::write(real.join("abc.jsonl"), "{}\n").unwrap();
1450        fs::write(stale.join("abc.jsonl"), "{}\n{}\n").unwrap();
1451
1452        let candidates: Vec<_> = agent_session::discover_session_files_in_home(temp.path())
1453            .into_iter()
1454            .filter(|candidate| candidate.agent == agent_session::AGENT_CURSOR)
1455            .collect();
1456
1457        assert_eq!(candidates.len(), 1);
1458        assert_eq!(candidates[0].path, real.join("abc.jsonl"));
1459    }
1460
1461    #[test]
1462    fn codex_state_db_uses_rollout_token_usage() {
1463        let temp = tempfile::tempdir().unwrap();
1464        write_codex_state_db_for_test(temp.path());
1465        let rollout = temp.path().join("session.jsonl");
1466        let mut content = "{}\n".repeat(CODEX_ROLLOUT_TAIL_BYTES as usize / 3 + 1);
1467        content.push_str(concat!(
1468            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}}}}"#,
1469            "\n"
1470        ));
1471        fs::write(&rollout, content).unwrap();
1472        let conn = rusqlite::Connection::open(temp.path().join(".codex/state_5.sqlite")).unwrap();
1473        conn.execute(
1474            "UPDATE threads SET rollout_path = ?1, tokens_used = 999999999",
1475            [rollout.to_string_lossy().as_ref()],
1476        )
1477        .unwrap();
1478
1479        let sessions = codex_state_sessions_in_home(temp.path(), 5);
1480
1481        assert_eq!(sessions[0].usage.input_tokens, 9_200);
1482        assert_eq!(sessions[0].usage.cache_read_tokens, 9_984);
1483        assert_eq!(sessions[0].usage.output_tokens, 11);
1484        assert_eq!(sessions[0].usage.total_tokens, 19_195);
1485    }
1486
1487    #[test]
1488    fn codex_state_db_errors_return_empty_for_jsonl_fallback() {
1489        let temp = tempfile::tempdir().unwrap();
1490        fs::create_dir_all(temp.path().join(".codex")).unwrap();
1491        fs::write(temp.path().join(".codex/state_5.sqlite"), "not sqlite").unwrap();
1492
1493        assert!(codex_state_sessions_in_home(temp.path(), 5).is_empty());
1494    }
1495
1496    #[test]
1497    fn crowded_codex_home_fallback_keeps_only_observed_exec_prompt() {
1498        let temp = tempfile::tempdir().unwrap();
1499        let state_path = write_codex_home(
1500            temp.path(),
1501            &[
1502                "unrelated historical prompt",
1503                "unrelated historical prompt",
1504                "unrelated historical prompt",
1505                "agentsight current run prompt",
1506                "agentsight current run historical prompt",
1507                "unrelated historical prompt",
1508                "unrelated historical prompt",
1509                "unrelated historical prompt",
1510            ],
1511        );
1512        let now = current_epoch_ms();
1513
1514        let rows = vec![
1515            exec_row(
1516                "audit-exec",
1517                now,
1518                "codex",
1519                &format!("{CODEX}agentsight current run prompt"),
1520            ),
1521            exec_row(
1522                "audit-exec-truncated",
1523                now + 1,
1524                "node",
1525                &format!("{CODEX}agentsight current run"),
1526            ),
1527            file_row("audit-file", now + 100, &state_path),
1528        ];
1529
1530        let sessions = observed_sessions_from_audit_rows(&rows);
1531
1532        assert_eq!(sessions.len(), 1);
1533        assert_eq!(
1534            sessions[0].prompt_preview.as_deref(),
1535            Some("agentsight current run prompt")
1536        );
1537    }
1538
1539    #[test]
1540    fn codex_home_fallback_accepts_time_window_when_exec_prompt_is_truncated() {
1541        let temp = tempfile::tempdir().unwrap();
1542        let state_path = write_codex_home(temp.path(), &["agentsight truncated command prompt"]);
1543        let now = current_epoch_ms();
1544
1545        let sessions = observed_sessions_from_audit_rows(&[
1546            exec_row(
1547                "audit-exec",
1548                now,
1549                "codex",
1550                &format!("{CODEX}-c model_provider=\"agentsight-mock"),
1551            ),
1552            file_row("audit-file", now + 100, &state_path),
1553        ]);
1554
1555        assert_eq!(sessions.len(), 1);
1556        assert_eq!(
1557            sessions[0].prompt_preview.as_deref(),
1558            Some("agentsight truncated command prompt")
1559        );
1560    }
1561
1562    #[test]
1563    fn codex_exec_prompt_rows_filter_only_wrapper_duplicates() {
1564        let rows = [
1565            (1_000, "node", "/usr/bin/node /opt/codex/bin/codex exec --skip-git-repo-check agentsight dedupe prompt"),
1566            (1_001, "codex", "/opt/codex/bin/codex exec --skip-git-repo-check agentsight dedupe prompt"),
1567            (2_000, "codex", "/usr/bin/codex exec --skip-git-repo-check agentsight short prompt"),
1568            (3_000, "codex", "/usr/bin/codex exec --skip-git-repo-check agentsight much longer unrelated prompt"),
1569            (10_000, "codex", "/usr/bin/codex exec --skip-git-repo-check agentsight repeated prompt"),
1570            (11_000, "codex", "/usr/bin/codex exec --skip-git-repo-check agentsight repeated prompt"),
1571            (20_000, "codex", "/usr/bin/codex exec --skip-git-repo-check agentsight repeated prompt"),
1572            (21_000, "docker", "docker exec codex exec agentsight should not parse"),
1573            (22_000, "docker", "docker exec container /usr/local/bin/codex exec agentsight should parse once"),
1574        ]
1575        .into_iter()
1576        .enumerate()
1577        .map(|(index, (ts, comm, command))| exec_row(&format!("audit-{index}"), ts, comm, command))
1578        .collect::<Vec<_>>();
1579
1580        let projected = observed_session_prompt_rows(&rows);
1581        assert_eq!(projected[0].comm.as_deref(), Some("node"));
1582        assert_eq!(projected[0].target.as_deref(), Some("/usr/bin/node"));
1583        let prompts = projected
1584            .into_iter()
1585            .map(|row| row.summary)
1586            .collect::<Vec<_>>();
1587
1588        assert_eq!(
1589            prompts,
1590            vec![
1591                Some("agentsight dedupe prompt".to_string()),
1592                Some("agentsight short prompt".to_string()),
1593                Some("agentsight much longer unrelated prompt".to_string()),
1594                Some("agentsight repeated prompt".to_string()),
1595                Some("agentsight repeated prompt".to_string()),
1596                Some("agentsight repeated prompt".to_string()),
1597                Some("agentsight should parse once".to_string()),
1598            ]
1599        );
1600    }
1601
1602    #[test]
1603    fn codex_fallback_time_window_rejects_stale_matching_session() {
1604        let (_temp, path) = create_temp_session_path(agent_session::AGENT_CODEX);
1605        let session = parse_content_for_test(
1606            agent_session::AGENT_CODEX,
1607            &path,
1608            UNIX_EPOCH,
1609            "{\"type\":\"message\",\"content\":\"agentsight repeated prompt\"}\n",
1610        )
1611        .unwrap();
1612
1613        assert!(session_matches_observed_prompt(
1614            &session,
1615            &[ObservedCodexPrompt {
1616                prompt: "agentsight repeated prompt".to_string(),
1617                timestamp_ms: current_epoch_ms(),
1618                pid: Some(42),
1619                native_exec: true,
1620                comm: Some("codex".to_string()),
1621                target: Some("/usr/bin/codex".to_string()),
1622            }]
1623        ));
1624        assert!(!session_is_in_observed_window(
1625            &session,
1626            Some((current_epoch_ms() - 1_000, current_epoch_ms() + 1_000))
1627        ));
1628    }
1629
1630    fn exec_row(id: &str, timestamp_ms: u64, comm: &str, full_command: &str) -> AuditEventRow {
1631        AuditEventRow {
1632            id: id.to_string(),
1633            timestamp_ms,
1634            audit_type: "process".to_string(),
1635            pid: Some(42),
1636            comm: Some(comm.to_string()),
1637            subject: None,
1638            action: Some("exec".to_string()),
1639            target: Some(format!("/usr/bin/{comm}")),
1640            status: Some("observed".to_string()),
1641            summary: None,
1642            details: serde_json::json!({ "full_command": full_command }),
1643        }
1644    }
1645
1646    fn file_row(id: &str, timestamp_ms: u64, path: &Path) -> AuditEventRow {
1647        AuditEventRow {
1648            id: id.to_string(),
1649            timestamp_ms,
1650            audit_type: "file".to_string(),
1651            pid: Some(42),
1652            comm: Some("codex".to_string()),
1653            subject: None,
1654            action: Some("write".to_string()),
1655            target: Some(path.to_string_lossy().to_string()),
1656            status: Some("observed".to_string()),
1657            summary: None,
1658            details: serde_json::json!({ "filepath": path.to_string_lossy() }),
1659        }
1660    }
1661
1662    fn current_epoch_ms() -> u64 {
1663        std::time::SystemTime::now()
1664            .duration_since(UNIX_EPOCH)
1665            .unwrap()
1666            .as_millis() as u64
1667    }
1668
1669    fn write_codex_home(root: &Path, prompts: &[&str]) -> PathBuf {
1670        let codex_home = root.join("codex-home");
1671        let sessions_dir = codex_home.join("sessions/2026/07/14");
1672        fs::create_dir_all(&sessions_dir).unwrap();
1673        for (index, prompt) in prompts.iter().enumerate() {
1674            fs::write(
1675                sessions_dir.join(format!(
1676                    "rollout-2026-07-14T00-00-{index:02}-session-{index}.jsonl"
1677                )),
1678                format!(
1679                    "{{\"timestamp\":\"2026-07-14T00:00:{index:02}.000Z\",\
1680                     \"type\":\"event_msg\",\
1681                     \"payload\":{{\"type\":\"user_message\",\"message\":\"{prompt}\"}}}}\n"
1682                ),
1683            )
1684            .unwrap();
1685        }
1686        let state_path = codex_home.join("stat");
1687        fs::write(&state_path, "").unwrap();
1688        state_path
1689    }
1690}