1use 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}