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
13#[cfg(any(test, feature = "test-support"))]
14use std::fs;
15
16use crate::model::{
17    AGENT_NATIVE_SOURCE, AuditEventRow, LlmCallRow, SessionRow, Snapshot, SnapshotOptions,
18    TokenUsageRow, ToolCallRow,
19};
20use crate::text::{sanitize_ascii_identifier as sanitize_id, truncate_text};
21use crate::view::MaterializedView;
22
23pub type LocalSession = AgentSession;
24pub type SessionCache = agent_session::SessionCache;
25const CODEX_EXEC_DEDUPE_WINDOW_MS: u64 = 2_000;
26const CODEX_FALLBACK_TIME_SLOP_MS: u64 = 30_000;
27const CODEX_ROLLOUT_TAIL_BYTES: u64 = 1024 * 1024;
28
29#[derive(Clone, Debug)]
30struct ObservedCodexPrompt {
31    prompt: String,
32    timestamp_ms: u64,
33    pid: Option<u32>,
34    native_exec: bool,
35    comm: Option<String>,
36    target: Option<String>,
37}
38
39pub fn snapshot(
40    cache: &mut SessionCache,
41    pid_filter: Option<u32>,
42    text_filter: Option<&str>,
43    limit: usize,
44    max_age: Duration,
45) -> Snapshot {
46    let filtered = discover_sessions(cache, pid_filter, text_filter, limit, max_age);
47    materialized_view(&filtered).export_snapshot(SnapshotOptions { audit_limit: 0 })
48}
49
50pub fn discover_sessions(
51    cache: &mut SessionCache,
52    pid_filter: Option<u32>,
53    text_filter: Option<&str>,
54    limit: usize,
55    max_age: Duration,
56) -> Vec<LocalSession> {
57    let indexed_codex = codex_state_sessions(limit);
58    let mut sessions = if indexed_codex.is_empty() {
59        cache.discover_cached(limit, max_age)
60    } else {
61        let mut sessions = indexed_codex;
62        sessions.extend(cache.discover_cached_excluding(
63            limit,
64            max_age,
65            &[agent_session::AGENT_CODEX],
66        ));
67        sessions.sort_by_key(|session| Reverse(session.updated));
68        sessions.truncate(limit.clamp(1, 25));
69        sessions
70    };
71    let mut seen = HashSet::new();
72    sessions.retain(|session| seen.insert(session.display_id.clone()));
73    sessions
74        .into_iter()
75        .filter(|s| matches_filter(s, pid_filter, text_filter))
76        .collect()
77}
78
79fn codex_state_sessions(limit: usize) -> Vec<LocalSession> {
80    user_home_dir()
81        .as_deref()
82        .map(|home| codex_state_sessions_in_home(home, limit))
83        .unwrap_or_default()
84}
85
86fn codex_state_sessions_in_home(home: &Path, limit: usize) -> Vec<LocalSession> {
87    let db_path = home.join(".codex/state_5.sqlite");
88    let Ok(conn) =
89        rusqlite::Connection::open_with_flags(db_path, rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY)
90    else {
91        return Vec::new();
92    };
93    let Ok(mut stmt) = conn.prepare(
94        "SELECT id, rollout_path, model, tokens_used, preview, cwd, created_at_ms, updated_at_ms
95         FROM threads
96         ORDER BY updated_at_ms DESC
97         LIMIT ?1",
98    ) else {
99        return Vec::new();
100    };
101    let Ok(rows) = stmt.query_map([limit.clamp(1, 25) as i64], |row| {
102        let id: String = row.get(0)?;
103        let rollout_path: String = row.get(1)?;
104        let model: Option<String> = row.get(2)?;
105        let tokens_used: i64 = row.get(3)?;
106        let preview: Option<String> = row.get(4)?;
107        let cwd: Option<String> = row.get(5)?;
108        let created_at_ms: Option<i64> = row.get(6)?;
109        let updated_at_ms: Option<i64> = row.get(7)?;
110        Ok(codex_state_session(
111            id,
112            rollout_path,
113            model,
114            tokens_used,
115            preview,
116            cwd,
117            created_at_ms,
118            updated_at_ms,
119        ))
120    }) else {
121        return Vec::new();
122    };
123
124    rows.filter_map(Result::ok).collect()
125}
126
127fn codex_state_session(
128    id: String,
129    rollout_path: String,
130    model: Option<String>,
131    tokens_used: i64,
132    preview: Option<String>,
133    cwd: Option<String>,
134    created_at_ms: Option<i64>,
135    updated_at_ms: Option<i64>,
136) -> LocalSession {
137    let updated_ms = updated_at_ms.and_then(non_negative_i64_to_u64);
138    let created_ms = created_at_ms
139        .and_then(non_negative_i64_to_u64)
140        .or(updated_ms);
141    let updated = updated_ms.map(system_time_from_ms).unwrap_or(UNIX_EPOCH);
142    let path = PathBuf::from(rollout_path);
143    let usage = codex_rollout_usage(&path).unwrap_or(TokenUsage {
144        total_tokens: tokens_used.max(0),
145        ..Default::default()
146    });
147    let model = model.filter(|value| !value.is_empty());
148    let mut model_usage = BTreeMap::new();
149    if let Some(model) = model.as_deref() {
150        model_usage.insert(model.to_string(), usage.clone());
151    }
152    let prompt_preview = preview
153        .and_then(|text| clean_prompt_text(&text))
154        .map(|text| truncate_text(&text, 180));
155    let last_message_at = updated_ms.map(iso_utc_from_ms);
156
157    LocalSession {
158        agent_type: agent_session::AGENT_CODEX.to_string(),
159        session_id: id.clone(),
160        conversation_id: Some(id.clone()),
161        display_id: format!("{}:{}", agent_session::AGENT_CODEX, short_session_id(&id)),
162        path,
163        updated,
164        start_timestamp_ms: created_ms,
165        end_timestamp_ms: updated_ms,
166        model,
167        usage,
168        model_usage,
169        tools: BTreeMap::new(),
170        files: BTreeMap::new(),
171        prompt_preview,
172        duration_ms: created_ms
173            .zip(updated_ms)
174            .map(|(start, end)| end.saturating_sub(start))
175            .unwrap_or_default(),
176        cwd,
177        last_message_at,
178        events: Default::default(),
179    }
180}
181
182fn codex_rollout_usage(path: &Path) -> Option<TokenUsage> {
183    let mut file = File::open(path).ok()?;
184    let len = file.metadata().ok()?.len();
185    let window = len.min(CODEX_ROLLOUT_TAIL_BYTES);
186    file.seek(SeekFrom::Start(len - window)).ok()?;
187    let mut data = Vec::with_capacity(window as usize);
188    file.read_to_end(&mut data).ok()?;
189    agent_session::codex_total_token_usage(&String::from_utf8_lossy(&data))
190}
191
192fn user_home_dir() -> Option<PathBuf> {
193    std::env::var("SUDO_USER")
194        .ok()
195        .and_then(|user| {
196            std::fs::read_to_string("/etc/passwd")
197                .ok()
198                .and_then(|passwd| {
199                    passwd
200                        .lines()
201                        .find(|line| line.starts_with(&format!("{user}:")))
202                        .and_then(|line| line.split(':').nth(5))
203                        .map(PathBuf::from)
204                })
205        })
206        .or_else(dirs::home_dir)
207}
208
209fn non_negative_i64_to_u64(value: i64) -> Option<u64> {
210    u64::try_from(value).ok()
211}
212
213fn system_time_from_ms(value: u64) -> SystemTime {
214    UNIX_EPOCH + Duration::from_millis(value)
215}
216
217fn iso_utc_from_ms(value: u64) -> String {
218    chrono::DateTime::<chrono::Utc>::from_timestamp_millis(value as i64)
219        .map(|dt| dt.to_rfc3339_opts(chrono::SecondsFormat::Millis, true))
220        .unwrap_or_default()
221}
222
223fn short_session_id(id: &str) -> String {
224    let compact = id.rsplit(['/', '\\']).next().unwrap_or(id).trim();
225    if compact.chars().count() <= 12 {
226        return compact.to_string();
227    }
228    let head = compact.chars().take(6).collect::<String>();
229    let tail = compact
230        .chars()
231        .rev()
232        .take(5)
233        .collect::<Vec<_>>()
234        .into_iter()
235        .rev()
236        .collect::<String>();
237    format!("{head}.{tail}")
238}
239
240fn clean_prompt_text(text: &str) -> Option<String> {
241    let text = text.split_whitespace().collect::<Vec<_>>().join(" ");
242    (!text.trim().is_empty()).then(|| text.trim().to_string())
243}
244
245fn view_id(session: &LocalSession) -> String {
246    format!("local:{}:{}", session.agent_type, session.display_id)
247}
248
249pub fn materialized_view(sessions: &[LocalSession]) -> MaterializedView {
250    let mut view = MaterializedView::new();
251    view.set_source(AGENT_NATIVE_SOURCE);
252    import_into_view(&mut view, sessions);
253    view
254}
255
256pub fn import_recent(view: &mut MaterializedView, limit: usize) {
257    let sessions = SessionCache::new().discover_cached(limit, Duration::ZERO);
258    import_into_view(view, &sessions);
259}
260
261pub fn import_into_view(view: &mut MaterializedView, sessions: &[LocalSession]) {
262    for session in sessions {
263        view.upsert_session(&session_row(session));
264        for row in llm_rows(session) {
265            view.apply_llm_call(&row);
266        }
267        for row in token_rows(session) {
268            view.apply_token_usage(&row);
269        }
270        for row in tool_rows(session) {
271            view.apply_tool_call(&row);
272        }
273    }
274}
275
276fn llm_rows(session: &LocalSession) -> Vec<LlmCallRow> {
277    let Some(prompt) = session.prompt_preview.as_ref() else {
278        return Vec::new();
279    };
280    let session_id = view_id(session);
281    let timestamp_ms = session
282        .events
283        .prompts
284        .first()
285        .and_then(|prompt| prompt.ts_ms)
286        .and_then(|ts| u64::try_from(ts).ok())
287        .or(session.start_timestamp_ms)
288        .unwrap_or_else(|| updated_ms(session));
289    let request = serde_json::json!({
290        "prompt": prompt,
291        "prompt_source": AGENT_NATIVE_SOURCE,
292        "session_id": session_id,
293        "agent_type": session.agent_type.as_str(),
294        "path": session.path.to_string_lossy(),
295    });
296
297    if session.model_usage.is_empty() {
298        let model = session
299            .model
300            .clone()
301            .unwrap_or_else(|| session.agent_type.clone());
302        return vec![llm_row_for_session(
303            &format!("{session_id}-{}", sanitize_id(&model)),
304            session,
305            &session_id,
306            timestamp_ms,
307            Some(model),
308            &session.usage,
309            request,
310        )];
311    }
312
313    session
314        .model_usage
315        .iter()
316        .map(|(model, usage)| {
317            llm_row_for_session(
318                &format!("{session_id}-{model}"),
319                session,
320                &session_id,
321                timestamp_ms,
322                Some(model.clone()),
323                usage,
324                request.clone(),
325            )
326        })
327        .collect()
328}
329
330fn llm_row_for_session(
331    id: &str,
332    session: &LocalSession,
333    session_id: &str,
334    timestamp_ms: u64,
335    model: Option<String>,
336    usage: &TokenUsage,
337    request: Value,
338) -> LlmCallRow {
339    LlmCallRow {
340        id: id.to_string(),
341        session_id: Some(session_id.to_string()),
342        conversation_id: session.conversation_id.clone(),
343        start_timestamp_ms: timestamp_ms,
344        end_timestamp_ms: session.end_timestamp_ms,
345        pid: None,
346        comm: Some(session.agent_type.clone()),
347        provider: None,
348        model,
349        call_kind: Some("agent_native_prompt".to_string()),
350        status: "observed".to_string(),
351        error_type: None,
352        finish_reason: None,
353        host: None,
354        path: Some(session.path.to_string_lossy().to_string()),
355        status_code: None,
356        input_tokens: usage.input_tokens,
357        output_tokens: usage.output_tokens,
358        total_tokens: usage.total_tokens,
359        request,
360        response: Value::Null,
361    }
362}
363
364pub fn observed_session_prompt_rows(audit_rows: &[AuditEventRow]) -> Vec<AuditEventRow> {
365    let mut rows = Vec::new();
366    let mut seen = HashSet::new();
367    let observed_exec_prompts = observed_codex_exec_prompts(audit_rows);
368    let mut seen_exec_prompts: Vec<ObservedCodexPrompt> = Vec::new();
369    for observed in observed_exec_prompts {
370        if seen_exec_prompts.iter().any(|seen| {
371            seen.prompt == observed.prompt
372                && timestamps_close(
373                    seen.timestamp_ms,
374                    observed.timestamp_ms,
375                    CODEX_EXEC_DEDUPE_WINDOW_MS,
376                )
377                && (!seen.native_exec || !observed.native_exec)
378        }) {
379            continue;
380        }
381        seen_exec_prompts.push(observed.clone());
382        rows.push(AuditEventRow {
383            id: format!(
384                "audit-codex-exec-prompt-{}-{}",
385                observed.timestamp_ms,
386                observed.pid.unwrap_or(0)
387            ),
388            timestamp_ms: observed.timestamp_ms,
389            audit_type: "llm".to_string(),
390            pid: observed.pid,
391            comm: observed.comm.or_else(|| Some("codex".to_string())),
392            subject: None,
393            action: Some("request".to_string()),
394            target: observed.target,
395            status: Some("observed".to_string()),
396            summary: Some(truncate_text(&observed.prompt, 160)),
397            details: serde_json::json!({
398                "text_content": observed.prompt,
399                "prompt_source": "local",
400            }),
401        });
402    }
403    for row in audit_rows {
404        if row.audit_type == "process" && row.action.as_deref() == Some("exec") {
405            continue;
406        }
407        if row.audit_type != "file" {
408            continue;
409        }
410        let Some(pid) = row.pid else {
411            continue;
412        };
413        let Some(path) = audit_session_path(row) else {
414            continue;
415        };
416        if !seen.insert((path.clone(), pid)) {
417            continue;
418        };
419        let Some(session) = agent_session::parse_session_path(&path) else {
420            continue;
421        };
422        let Some(prompt) = session.prompt_preview.as_ref() else {
423            continue;
424        };
425        rows.push(AuditEventRow {
426            id: format!(
427                "audit-agent-native-prompt-{}-{pid}",
428                sanitize_id(&session.display_id)
429            ),
430            timestamp_ms: row.timestamp_ms,
431            audit_type: "llm".to_string(),
432            pid: Some(pid),
433            comm: row
434                .comm
435                .clone()
436                .or_else(|| Some(session.agent_type.clone())),
437            subject: session.model.clone(),
438            action: Some("request".to_string()),
439            target: Some(path.to_string_lossy().to_string()),
440            status: Some("observed".to_string()),
441            summary: Some(truncate_text(prompt, 160)),
442            details: serde_json::json!({
443                "text_content": prompt,
444                "prompt_source": "local",
445                "session_id": view_id(&session),
446                "conversation_id": session.conversation_id.as_deref(),
447                "agent_type": session.agent_type,
448            }),
449        });
450    }
451    rows
452}
453
454pub fn observed_sessions_from_audit_rows(audit_rows: &[AuditEventRow]) -> Vec<LocalSession> {
455    let mut direct_paths = HashSet::new();
456    let mut codex_session_dirs = HashSet::new();
457    let observed_codex_prompts = observed_codex_exec_prompts(audit_rows);
458    let observed_codex_exec = observed_codex_exec_command(audit_rows);
459    let observed_window = observed_audit_window_ms(audit_rows);
460
461    for row in audit_rows.iter().filter(|row| row.audit_type == "file") {
462        for path in audit_file_paths(row) {
463            if let Some(session_path) =
464                agent_session::session_log_path_from_str(path.to_string_lossy().as_ref())
465            {
466                direct_paths.insert(session_path);
467            }
468            if let Some(dir) = observed_codex_sessions_dir(&path) {
469                codex_session_dirs.insert(dir);
470            }
471        }
472    }
473
474    let mut candidates = Vec::new();
475    for path in direct_paths {
476        if let Some(candidate) = agent_session::session_candidate_from_path(&path) {
477            candidates.push((candidate, false));
478        }
479    }
480    for dir in codex_session_dirs {
481        let dir_candidates =
482            agent_session::discover_session_files_in_dir(agent_session::AGENT_CODEX, &dir);
483        candidates.extend(
484            dir_candidates
485                .into_iter()
486                .map(|candidate| (candidate, true)),
487        );
488    }
489    candidates.sort_by_key(|(candidate, _)| Reverse(candidate.updated));
490
491    let mut seen_paths = HashSet::new();
492    let mut seen_sessions = HashSet::new();
493    let mut sessions = Vec::new();
494    for (candidate, is_codex_dir_fallback) in candidates.into_iter().take(75) {
495        if !seen_paths.insert(candidate.path.clone()) {
496            continue;
497        }
498        let Some(session) = agent_session::parse_session_file(&candidate) else {
499            continue;
500        };
501        if is_codex_dir_fallback && !observed_codex_exec {
502            continue;
503        }
504        if is_codex_dir_fallback
505            && !observed_codex_prompts.is_empty()
506            && !session_matches_observed_prompt(&session, &observed_codex_prompts)
507        {
508            continue;
509        }
510        if is_codex_dir_fallback && !session_is_in_observed_window(&session, observed_window) {
511            continue;
512        }
513        if seen_sessions.insert(session.display_id.clone()) {
514            sessions.push(session);
515        }
516    }
517    sessions
518}
519
520fn observed_codex_exec_prompts(audit_rows: &[AuditEventRow]) -> Vec<ObservedCodexPrompt> {
521    let prompts = audit_rows
522        .iter()
523        .filter(|row| row.audit_type == "process" && row.action.as_deref() == Some("exec"))
524        .filter_map(|row| {
525            let prompt = row
526                .details
527                .get("full_command")
528                .and_then(Value::as_str)
529                .and_then(codex_exec_prompt_from_command)?;
530            Some(ObservedCodexPrompt {
531                prompt,
532                timestamp_ms: row.timestamp_ms,
533                pid: row.pid,
534                native_exec: looks_like_native_codex_exec(row),
535                comm: row.comm.clone(),
536                target: row.target.clone(),
537            })
538        })
539        .collect::<Vec<_>>();
540    prompts
541        .iter()
542        .filter(|candidate| {
543            !prompts
544                .iter()
545                .any(|other| is_nearby_longer_prefix(candidate, other))
546        })
547        .cloned()
548        .collect()
549}
550
551fn observed_codex_exec_command(audit_rows: &[AuditEventRow]) -> bool {
552    audit_rows
553        .iter()
554        .filter(|row| row.audit_type == "process" && row.action.as_deref() == Some("exec"))
555        .filter_map(|row| row.details.get("full_command").and_then(Value::as_str))
556        .any(|command| codex_exec_command_tail(command).is_some())
557}
558
559fn codex_exec_prompt_from_command(command: &str) -> Option<String> {
560    agent_session::codex_exec_prompt(&codex_exec_command_tail(command)?)
561}
562
563fn codex_exec_command_tail(command: &str) -> Option<String> {
564    let tokens = command.split_whitespace().collect::<Vec<_>>();
565    let index = tokens.windows(2).enumerate().find_map(|(index, tokens)| {
566        (is_codex_executable_token(tokens[0], index == 0) && tokens[1] == "exec").then_some(index)
567    })?;
568    Some(tokens[index..].join(" "))
569}
570
571fn looks_like_native_codex_exec(row: &AuditEventRow) -> bool {
572    row.comm.as_deref() == Some("codex")
573        && row
574            .target
575            .as_deref()
576            .is_some_and(|target| is_codex_executable_token(target, true))
577}
578
579fn is_codex_executable_token(token: &str, allow_bare: bool) -> bool {
580    let token = token.trim_matches(|ch| matches!(ch, '"' | '\''));
581    (allow_bare && token == "codex")
582        || token.contains('/')
583            && Path::new(token).file_name().and_then(|name| name.to_str()) == Some("codex")
584}
585
586fn is_nearby_longer_prefix(candidate: &ObservedCodexPrompt, other: &ObservedCodexPrompt) -> bool {
587    other.prompt.len() > candidate.prompt.len()
588        && other.prompt.starts_with(candidate.prompt.as_str())
589        && timestamps_close(
590            candidate.timestamp_ms,
591            other.timestamp_ms,
592            CODEX_EXEC_DEDUPE_WINDOW_MS,
593        )
594        && (!candidate.native_exec || !other.native_exec)
595}
596
597fn timestamps_close(left: u64, right: u64, window_ms: u64) -> bool {
598    left.abs_diff(right) <= window_ms
599}
600
601fn session_matches_observed_prompt(
602    session: &LocalSession,
603    prompts: &[ObservedCodexPrompt],
604) -> bool {
605    let Some(preview) = session.prompt_preview.as_deref() else {
606        return false;
607    };
608    prompts.iter().any(|observed| {
609        let prompt = observed.prompt.as_str();
610        prompt_texts_overlap(prompt, preview)
611    })
612}
613
614fn prompt_texts_overlap(left: &str, right: &str) -> bool {
615    left == right || left.starts_with(right) || right.starts_with(left)
616}
617
618fn observed_audit_window_ms(audit_rows: &[AuditEventRow]) -> Option<(u64, u64)> {
619    let min = audit_rows.iter().map(|row| row.timestamp_ms).min()?;
620    let max = audit_rows.iter().map(|row| row.timestamp_ms).max()?;
621    Some((
622        min.saturating_sub(CODEX_FALLBACK_TIME_SLOP_MS),
623        max.saturating_add(CODEX_FALLBACK_TIME_SLOP_MS),
624    ))
625}
626
627fn session_is_in_observed_window(session: &LocalSession, window: Option<(u64, u64)>) -> bool {
628    let Some((min, max)) = window else {
629        return true;
630    };
631    let updated = updated_ms(session);
632    updated >= min && updated <= max
633}
634
635fn audit_session_path(row: &AuditEventRow) -> Option<PathBuf> {
636    row.target
637        .as_deref()
638        .and_then(agent_session::session_log_path_from_str)
639        .or_else(|| {
640            row.details
641                .get("filepath")
642                .and_then(Value::as_str)
643                .and_then(agent_session::session_log_path_from_str)
644        })
645        .or_else(|| {
646            row.details
647                .get("path")
648                .and_then(Value::as_str)
649                .and_then(agent_session::session_log_path_from_str)
650        })
651}
652
653fn audit_file_paths(row: &AuditEventRow) -> Vec<PathBuf> {
654    [
655        row.target.as_deref(),
656        row.details.get("filepath").and_then(Value::as_str),
657        row.details.get("path").and_then(Value::as_str),
658        row.details.get("fd_target").and_then(Value::as_str),
659    ]
660    .into_iter()
661    .flatten()
662    .filter_map(|raw| {
663        let path = PathBuf::from(raw.trim().trim_end_matches(" (deleted)"));
664        path.is_absolute().then_some(path)
665    })
666    .collect()
667}
668
669fn observed_codex_sessions_dir(path: &Path) -> Option<PathBuf> {
670    if !looks_like_codex_home_file(path) {
671        return None;
672    }
673    let home = path.parent()?;
674    let sessions = home.join("sessions");
675    sessions.is_dir().then_some(sessions)
676}
677
678fn looks_like_codex_home_file(path: &Path) -> bool {
679    let Some(name) = path.file_name().and_then(|name| name.to_str()) else {
680        return false;
681    };
682    (name.starts_with("state_")
683        || name.starts_with("logs_")
684        || matches!(name, "config.toml" | "auth.json" | "stat"))
685        && path
686            .parent()
687            .is_some_and(|parent| parent.join("sessions").is_dir())
688}
689
690fn session_row(session: &LocalSession) -> SessionRow {
691    let updated_ms = updated_ms(session);
692    SessionRow {
693        id: view_id(session),
694        agent_type: session.agent_type.clone(),
695        start_timestamp_ms: session
696            .start_timestamp_ms
697            .unwrap_or_else(|| updated_ms.saturating_sub(session.duration_ms)),
698        end_timestamp_ms: session.end_timestamp_ms.or(Some(updated_ms)),
699        status: "observed".to_string(),
700        model: session.model.clone(),
701        input_tokens: session.usage.input_tokens,
702        output_tokens: session.usage.output_tokens,
703        total_tokens: session.usage.total_tokens,
704        view_source: AGENT_NATIVE_SOURCE.to_string(),
705        confidence: Some(0.95),
706        attributes: serde_json::json!({
707            "session_id": session.session_id.clone(),
708            "conversation_id": session.conversation_id.as_deref(),
709            "path": session.path.to_string_lossy(),
710            "display_id": session.display_id,
711            "prompt_preview": session.prompt_preview.clone(),
712            "cwd": session.cwd.clone(),
713            "last_message_at": session.last_message_at.clone(),
714            "files": session.files,
715        }),
716    }
717}
718
719fn token_rows(session: &LocalSession) -> Vec<TokenUsageRow> {
720    let session_id = view_id(session);
721    session
722        .model_usage
723        .iter()
724        .filter(|(_, usage)| usage.total_tokens > 0)
725        .map(|(model, usage)| TokenUsageRow {
726            id: format!("token-{session_id}-{}", sanitize_id(model)),
727            llm_call_id: format!("{session_id}-{model}"),
728            timestamp_ms: updated_ms(session),
729            pid: None,
730            comm: Some(session.agent_type.clone()),
731            provider: None,
732            model: Some(model.clone()),
733            input_tokens: usage.input_tokens,
734            output_tokens: usage.output_tokens,
735            cache_creation_tokens: usage.cache_creation_tokens,
736            cache_read_tokens: usage.cache_read_tokens,
737            total_tokens: usage.total_tokens,
738            source: AGENT_NATIVE_SOURCE.to_string(),
739            view_source: AGENT_NATIVE_SOURCE.to_string(),
740            confidence: Some(0.95),
741        })
742        .collect()
743}
744
745fn tool_rows(session: &LocalSession) -> Vec<ToolCallRow> {
746    let session_id = view_id(session);
747    let timestamp_ms = updated_ms(session);
748    let mut rows = Vec::new();
749    for (tool, count) in &session.tools {
750        for index in 0..*count {
751            rows.push(ToolCallRow {
752                id: format!("tool-{session_id}-{}-{index}", sanitize_id(tool)),
753                session_id: Some(session_id.clone()),
754                conversation_id: session.conversation_id.clone(),
755                timestamp_ms,
756                tool_name: Some(tool.clone()),
757                tool_call_id: None,
758                start_timestamp_ms: Some(timestamp_ms),
759                end_timestamp_ms: Some(timestamp_ms),
760                duration_ms: None,
761                status: Some("observed".to_string()),
762                input: serde_json::json!({}),
763                output: serde_json::json!({}),
764                related_pid: None,
765                related_event_id: None,
766                view_source: AGENT_NATIVE_SOURCE.to_string(),
767                confidence: Some(0.95),
768            });
769        }
770    }
771    rows
772}
773
774fn updated_ms(session: &LocalSession) -> u64 {
775    session
776        .updated
777        .duration_since(UNIX_EPOCH)
778        .unwrap_or_default()
779        .as_millis() as u64
780}
781
782fn matches_filter(
783    session: &LocalSession,
784    pid_filter: Option<u32>,
785    text_filter: Option<&str>,
786) -> bool {
787    if pid_filter.is_some() {
788        return true;
789    }
790    let Some(filter) = text_filter else {
791        return true;
792    };
793    let filter = filter.to_ascii_lowercase();
794    session.agent_type.to_ascii_lowercase().contains(&filter)
795        || session
796            .prompt_preview
797            .as_ref()
798            .is_some_and(|prompt| prompt.to_ascii_lowercase().contains(&filter))
799        || session
800            .model
801            .as_ref()
802            .is_some_and(|model| model.to_ascii_lowercase().contains(&filter))
803        || session
804            .path
805            .to_string_lossy()
806            .to_ascii_lowercase()
807            .contains(&filter)
808}
809
810#[cfg(any(test, feature = "test-support"))]
811pub fn create_temp_session_path(agent: &str) -> (tempfile::TempDir, PathBuf) {
812    let temp = tempfile::tempdir().unwrap();
813    let path = agent_session::fixture_session_path(agent, temp.path()).unwrap();
814    fs::create_dir_all(path.parent().unwrap()).unwrap();
815    fs::write(&path, "{}\n").unwrap();
816    (temp, path)
817}
818
819#[cfg(any(test, feature = "test-support"))]
820pub fn parse_content_for_test(
821    agent: &str,
822    path: &std::path::Path,
823    updated: std::time::SystemTime,
824    content: &str,
825) -> Option<LocalSession> {
826    agent_session::parse_session_content(agent, path, updated, content)
827}
828
829#[cfg(any(test, feature = "test-support"))]
830pub fn write_codex_state_db_for_test(home: &Path) {
831    let codex_dir = home.join(".codex");
832    fs::create_dir_all(&codex_dir).unwrap();
833    let conn = rusqlite::Connection::open(codex_dir.join("state_5.sqlite")).unwrap();
834    conn.execute_batch(
835        "CREATE TABLE threads (
836            id TEXT PRIMARY KEY,
837            rollout_path TEXT,
838            model TEXT,
839            tokens_used INTEGER NOT NULL DEFAULT 0,
840            preview TEXT,
841            cwd TEXT,
842            created_at_ms INTEGER,
843            updated_at_ms INTEGER
844        );
845        INSERT INTO threads
846        (id, rollout_path, model, tokens_used, preview, cwd, created_at_ms, updated_at_ms)
847        VALUES
848        ('019f49ca-54e7-7a91-82e7-a52b53cfd456', '/tmp/session.jsonl', 'gpt-web-ci', 33, 'web state prompt', '/work/repo', 1800000, 1900000);",
849    )
850    .unwrap();
851}
852
853#[cfg(test)]
854mod tests {
855    use super::*;
856
857    const CODEX: &str = "/usr/bin/codex exec --skip-git-repo-check ";
858
859    #[test]
860    fn agent_native_prompt_produces_llm_call_row() {
861        let (_temp, path) = create_temp_session_path(agent_session::AGENT_CODEX);
862        let session = parse_content_for_test(
863            agent_session::AGENT_CODEX,
864            &path,
865            UNIX_EPOCH,
866            "{\"type\":\"message\",\"content\":\"agentsight local codex prompt\"}\n",
867        )
868        .unwrap();
869
870        let view = materialized_view(&[session]);
871        let rows = view.llm_call_rows(10);
872
873        assert_eq!(rows.len(), 1);
874        assert_eq!(rows[0].comm.as_deref(), Some(agent_session::AGENT_CODEX));
875        assert_eq!(
876            rows[0].request.get("prompt").and_then(Value::as_str),
877            Some("agentsight local codex prompt")
878        );
879    }
880
881    #[test]
882    fn codex_state_db_produces_indexed_session_metadata() {
883        let temp = tempfile::tempdir().unwrap();
884        write_codex_state_db_for_test(temp.path());
885
886        let sessions = codex_state_sessions_in_home(temp.path(), 5);
887
888        assert_eq!(sessions.len(), 1);
889        assert_eq!(sessions[0].display_id, "codex:019f49.fd456");
890        assert_eq!(sessions[0].model.as_deref(), Some("gpt-web-ci"));
891        assert_eq!(sessions[0].usage.total_tokens, 33);
892        assert_eq!(
893            sessions[0].prompt_preview.as_deref(),
894            Some("web state prompt")
895        );
896        assert_eq!(sessions[0].cwd.as_deref(), Some("/work/repo"));
897        assert_eq!(
898            sessions[0].last_message_at.as_deref(),
899            Some("1970-01-01T00:31:40.000Z")
900        );
901    }
902
903    #[test]
904    fn codex_state_db_uses_rollout_token_usage() {
905        let temp = tempfile::tempdir().unwrap();
906        write_codex_state_db_for_test(temp.path());
907        let rollout = temp.path().join("session.jsonl");
908        let mut content = "{}\n".repeat(CODEX_ROLLOUT_TAIL_BYTES as usize / 3 + 1);
909        content.push_str(concat!(
910            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}}}}"#,
911            "\n"
912        ));
913        fs::write(&rollout, content).unwrap();
914        let conn = rusqlite::Connection::open(temp.path().join(".codex/state_5.sqlite")).unwrap();
915        conn.execute(
916            "UPDATE threads SET rollout_path = ?1, tokens_used = 999999999",
917            [rollout.to_string_lossy().as_ref()],
918        )
919        .unwrap();
920
921        let sessions = codex_state_sessions_in_home(temp.path(), 5);
922
923        assert_eq!(sessions[0].usage.input_tokens, 9_200);
924        assert_eq!(sessions[0].usage.cache_read_tokens, 9_984);
925        assert_eq!(sessions[0].usage.output_tokens, 11);
926        assert_eq!(sessions[0].usage.total_tokens, 19_195);
927    }
928
929    #[test]
930    fn codex_state_db_errors_return_empty_for_jsonl_fallback() {
931        let temp = tempfile::tempdir().unwrap();
932        fs::create_dir_all(temp.path().join(".codex")).unwrap();
933        fs::write(temp.path().join(".codex/state_5.sqlite"), "not sqlite").unwrap();
934
935        assert!(codex_state_sessions_in_home(temp.path(), 5).is_empty());
936    }
937
938    #[test]
939    fn crowded_codex_home_fallback_keeps_only_observed_exec_prompt() {
940        let temp = tempfile::tempdir().unwrap();
941        let state_path = write_codex_home(
942            temp.path(),
943            &[
944                "unrelated historical prompt",
945                "unrelated historical prompt",
946                "unrelated historical prompt",
947                "agentsight current run prompt",
948                "agentsight current run historical prompt",
949                "unrelated historical prompt",
950                "unrelated historical prompt",
951                "unrelated historical prompt",
952            ],
953        );
954        let now = current_epoch_ms();
955
956        let rows = vec![
957            exec_row(
958                "audit-exec",
959                now,
960                "codex",
961                &format!("{CODEX}agentsight current run prompt"),
962            ),
963            exec_row(
964                "audit-exec-truncated",
965                now + 1,
966                "node",
967                &format!("{CODEX}agentsight current run"),
968            ),
969            file_row("audit-file", now + 100, &state_path),
970        ];
971
972        let sessions = observed_sessions_from_audit_rows(&rows);
973
974        assert_eq!(sessions.len(), 1);
975        assert_eq!(
976            sessions[0].prompt_preview.as_deref(),
977            Some("agentsight current run prompt")
978        );
979    }
980
981    #[test]
982    fn codex_home_fallback_accepts_time_window_when_exec_prompt_is_truncated() {
983        let temp = tempfile::tempdir().unwrap();
984        let state_path = write_codex_home(temp.path(), &["agentsight truncated command prompt"]);
985        let now = current_epoch_ms();
986
987        let sessions = observed_sessions_from_audit_rows(&[
988            exec_row(
989                "audit-exec",
990                now,
991                "codex",
992                &format!("{CODEX}-c model_provider=\"agentsight-mock"),
993            ),
994            file_row("audit-file", now + 100, &state_path),
995        ]);
996
997        assert_eq!(sessions.len(), 1);
998        assert_eq!(
999            sessions[0].prompt_preview.as_deref(),
1000            Some("agentsight truncated command prompt")
1001        );
1002    }
1003
1004    #[test]
1005    fn codex_exec_prompt_rows_filter_only_wrapper_duplicates() {
1006        let rows = [
1007            (1_000, "node", "/usr/bin/node /opt/codex/bin/codex exec --skip-git-repo-check agentsight dedupe prompt"),
1008            (1_001, "codex", "/opt/codex/bin/codex exec --skip-git-repo-check agentsight dedupe prompt"),
1009            (2_000, "codex", "/usr/bin/codex exec --skip-git-repo-check agentsight short prompt"),
1010            (3_000, "codex", "/usr/bin/codex exec --skip-git-repo-check agentsight much longer unrelated prompt"),
1011            (10_000, "codex", "/usr/bin/codex exec --skip-git-repo-check agentsight repeated prompt"),
1012            (11_000, "codex", "/usr/bin/codex exec --skip-git-repo-check agentsight repeated prompt"),
1013            (20_000, "codex", "/usr/bin/codex exec --skip-git-repo-check agentsight repeated prompt"),
1014            (21_000, "docker", "docker exec codex exec agentsight should not parse"),
1015            (22_000, "docker", "docker exec container /usr/local/bin/codex exec agentsight should parse once"),
1016        ]
1017        .into_iter()
1018        .enumerate()
1019        .map(|(index, (ts, comm, command))| exec_row(&format!("audit-{index}"), ts, comm, command))
1020        .collect::<Vec<_>>();
1021
1022        let projected = observed_session_prompt_rows(&rows);
1023        assert_eq!(projected[0].comm.as_deref(), Some("node"));
1024        assert_eq!(projected[0].target.as_deref(), Some("/usr/bin/node"));
1025        let prompts = projected
1026            .into_iter()
1027            .map(|row| row.summary)
1028            .collect::<Vec<_>>();
1029
1030        assert_eq!(
1031            prompts,
1032            vec![
1033                Some("agentsight dedupe prompt".to_string()),
1034                Some("agentsight short prompt".to_string()),
1035                Some("agentsight much longer unrelated prompt".to_string()),
1036                Some("agentsight repeated prompt".to_string()),
1037                Some("agentsight repeated prompt".to_string()),
1038                Some("agentsight repeated prompt".to_string()),
1039                Some("agentsight should parse once".to_string()),
1040            ]
1041        );
1042    }
1043
1044    #[test]
1045    fn codex_fallback_time_window_rejects_stale_matching_session() {
1046        let (_temp, path) = create_temp_session_path(agent_session::AGENT_CODEX);
1047        let session = parse_content_for_test(
1048            agent_session::AGENT_CODEX,
1049            &path,
1050            UNIX_EPOCH,
1051            "{\"type\":\"message\",\"content\":\"agentsight repeated prompt\"}\n",
1052        )
1053        .unwrap();
1054
1055        assert!(session_matches_observed_prompt(
1056            &session,
1057            &[ObservedCodexPrompt {
1058                prompt: "agentsight repeated prompt".to_string(),
1059                timestamp_ms: current_epoch_ms(),
1060                pid: Some(42),
1061                native_exec: true,
1062                comm: Some("codex".to_string()),
1063                target: Some("/usr/bin/codex".to_string()),
1064            }]
1065        ));
1066        assert!(!session_is_in_observed_window(
1067            &session,
1068            Some((current_epoch_ms() - 1_000, current_epoch_ms() + 1_000))
1069        ));
1070    }
1071
1072    fn exec_row(id: &str, timestamp_ms: u64, comm: &str, full_command: &str) -> AuditEventRow {
1073        AuditEventRow {
1074            id: id.to_string(),
1075            timestamp_ms,
1076            audit_type: "process".to_string(),
1077            pid: Some(42),
1078            comm: Some(comm.to_string()),
1079            subject: None,
1080            action: Some("exec".to_string()),
1081            target: Some(format!("/usr/bin/{comm}")),
1082            status: Some("observed".to_string()),
1083            summary: None,
1084            details: serde_json::json!({ "full_command": full_command }),
1085        }
1086    }
1087
1088    fn file_row(id: &str, timestamp_ms: u64, path: &Path) -> AuditEventRow {
1089        AuditEventRow {
1090            id: id.to_string(),
1091            timestamp_ms,
1092            audit_type: "file".to_string(),
1093            pid: Some(42),
1094            comm: Some("codex".to_string()),
1095            subject: None,
1096            action: Some("write".to_string()),
1097            target: Some(path.to_string_lossy().to_string()),
1098            status: Some("observed".to_string()),
1099            summary: None,
1100            details: serde_json::json!({ "filepath": path.to_string_lossy() }),
1101        }
1102    }
1103
1104    fn current_epoch_ms() -> u64 {
1105        std::time::SystemTime::now()
1106            .duration_since(UNIX_EPOCH)
1107            .unwrap()
1108            .as_millis() as u64
1109    }
1110
1111    fn write_codex_home(root: &Path, prompts: &[&str]) -> PathBuf {
1112        let codex_home = root.join("codex-home");
1113        let sessions_dir = codex_home.join("sessions/2026/07/14");
1114        fs::create_dir_all(&sessions_dir).unwrap();
1115        for (index, prompt) in prompts.iter().enumerate() {
1116            fs::write(
1117                sessions_dir.join(format!(
1118                    "rollout-2026-07-14T00-00-{index:02}-session-{index}.jsonl"
1119                )),
1120                format!(
1121                    "{{\"timestamp\":\"2026-07-14T00:00:{index:02}.000Z\",\
1122                     \"type\":\"event_msg\",\
1123                     \"payload\":{{\"type\":\"user_message\",\"message\":\"{prompt}\"}}}}\n"
1124                ),
1125            )
1126            .unwrap();
1127        }
1128        let state_path = codex_home.join("stat");
1129        fs::write(&state_path, "").unwrap();
1130        state_path
1131    }
1132}