1use crate::model::{AuditEventRow, LlmCallRow, ProcessNodeRow, ViewResult};
5use crate::sinks::sqlite::SqliteStore;
6use crate::sources::agent_native;
7use crate::text::{clean_prompt_text, extract_prompt_text, truncate_text};
8use crate::view::MaterializedView;
9use serde_json::{Value, json};
10use std::collections::BTreeSet;
11use std::path::Path;
12
13const PROMPT_DEDUP_WINDOW_MS: u64 = 10_000;
14#[cfg_attr(not(test), allow(dead_code))]
15pub fn load_view(path: impl AsRef<Path>) -> ViewResult<MaterializedView> {
16 load_view_inner(path, false)
17}
18
19pub fn load_view_with_observed_session_prompts(
20 path: impl AsRef<Path>,
21) -> ViewResult<MaterializedView> {
22 load_view_inner(path, true)
23}
24
25fn load_view_inner(
26 path: impl AsRef<Path>,
27 include_observed_session_prompts: bool,
28) -> ViewResult<MaterializedView> {
29 let store = SqliteStore::open_readonly(path)?;
30 let mut view = MaterializedView::new();
31 view.set_source("sqlite");
32
33 let mut llm_rows = Vec::new();
34 if let Ok(rows) = store.all_llm_call_rows() {
35 for row in &rows {
36 view.apply_llm_call(row);
37 }
38 llm_rows = rows;
39 }
40 if let Ok(rows) = store.token_usage_rows() {
41 for row in rows {
42 view.apply_token_usage(&row);
43 }
44 }
45 let mut audit_rows = Vec::new();
46 if let Ok(rows) = store.all_audit_event_rows() {
47 for row in &rows {
48 if include_observed_session_prompts && is_reprojected_llm_request(row) {
49 continue;
50 }
51 view.apply_audit_event(row);
52 }
53 audit_rows = rows;
54 }
55 let mut process_pids = BTreeSet::new();
56 if let Ok(rows) = store.process_node_rows() {
57 for row in &rows {
58 process_pids.insert(row.pid);
59 view.upsert_process_node(row);
60 }
61 }
62 if let Ok(rows) = store.tool_call_rows() {
63 for row in rows {
64 view.apply_tool_call(&row);
65 }
66 }
67 if let Ok(rows) = store.network_target_rows() {
68 for row in rows {
69 view.upsert_network_target(&row);
70 }
71 }
72 if let Ok(rows) = store.resource_sample_rows() {
73 for row in rows {
74 view.apply_resource_sample(&row);
75 }
76 }
77 if include_observed_session_prompts {
78 import_observed_process_nodes(&mut view, &llm_rows, &process_pids);
79 let observed_sessions = agent_native::observed_sessions_from_audit_rows(&audit_rows);
80 agent_native::import_into_view(&mut view, &observed_sessions);
81 let current_llm_rows = view.llm_call_rows(usize::MAX);
82 let mut prompt_rows = llm_call_prompt_rows(¤t_llm_rows);
83 let mut local_prompt_rows = agent_native::observed_session_prompt_rows(&audit_rows);
84 local_prompt_rows.sort_by_key(|row| {
85 row.details
86 .get("session_id")
87 .and_then(Value::as_str)
88 .is_none()
89 });
90 append_deduped_local_session_prompt_rows(&mut prompt_rows, local_prompt_rows);
91 for row in local_prompt_llm_call_rows(&prompt_rows) {
92 view.apply_llm_call(&row);
93 }
94 for row in prompt_rows {
95 view.apply_audit_event(&row);
96 }
97 }
98
99 Ok(view)
100}
101
102fn import_observed_process_nodes(
103 view: &mut MaterializedView,
104 llm_rows: &[LlmCallRow],
105 existing_pids: &BTreeSet<u32>,
106) {
107 for row in llm_rows {
108 let Some(pid) = row.pid else {
109 continue;
110 };
111 if existing_pids.contains(&pid) {
112 continue;
113 }
114 let comm = row.comm.clone();
115 let command = comm.clone().unwrap_or_else(|| format!("pid {}", pid));
116 view.upsert_process_node(&ProcessNodeRow {
117 id: format!("process-{}-observed", pid),
118 pid,
119 ppid: None,
120 root_pid: Some(pid),
121 start_timestamp_ms: Some(row.start_timestamp_ms),
122 end_timestamp_ms: None,
123 comm,
124 command: Some(command),
125 argv: Vec::new(),
126 cwd: None,
127 exit_code: None,
128 status: Some("observed".to_string()),
129 view_source: "sqlite".to_string(),
130 confidence: Some(0.5),
131 });
132 }
133}
134
135fn is_reprojected_llm_request(row: &AuditEventRow) -> bool {
136 row.audit_type == "llm" && row.action.as_deref() == Some("request")
137}
138
139fn llm_call_prompt_rows(rows: &[LlmCallRow]) -> Vec<AuditEventRow> {
140 let mut prompts = Vec::new();
141 for row in rows {
142 if row.request.is_null() || row.request.as_object().is_some_and(|obj| obj.is_empty()) {
143 continue;
144 }
145 let Some(text) = extract_prompt_text(&row.request) else {
146 continue;
147 };
148 let is_agent_native = row.request.get("prompt_source").and_then(Value::as_str)
149 == Some(crate::model::AGENT_NATIVE_SOURCE);
150 let prompt_source = if is_agent_native { "local" } else { "ssl" };
151 prompts.push(AuditEventRow {
152 id: format!("audit-{}-request", row.id),
153 timestamp_ms: row.start_timestamp_ms,
154 audit_type: "llm".to_string(),
155 pid: row.pid,
156 comm: row.comm.clone(),
157 subject: row.model.clone(),
158 action: Some("request".to_string()),
159 target: row.host.clone(),
160 status: Some("observed".to_string()),
161 summary: Some(truncate_text(&text, 160)),
162 details: json!({
163 "text_content": text,
164 "prompt_source": prompt_source,
165 "session_id": row.request.get("session_id").and_then(Value::as_str),
166 "request": row.request,
167 "provider": row.provider,
168 "path": row.path,
169 }),
170 });
171 }
172 prompts
173}
174
175fn append_deduped_local_session_prompt_rows(
176 ssl_rows: &mut Vec<AuditEventRow>,
177 local_rows: Vec<AuditEventRow>,
178) {
179 for local in local_rows {
180 let Some(local_text) = prompt_text_from_details(&local.details) else {
181 ssl_rows.push(local);
182 continue;
183 };
184 let duplicate = ssl_rows.iter().any(|ssl| {
185 let source = ssl.details.get("prompt_source").and_then(Value::as_str);
186 let is_session_bound_local = source == Some("local")
187 && ssl
188 .details
189 .get("session_id")
190 .and_then(Value::as_str)
191 .is_some();
192 if source != Some("ssl") && !is_session_bound_local {
193 return false;
194 }
195 if let (Some(local_pid), Some(ssl_pid)) = (local.pid, ssl.pid)
196 && local_pid != ssl_pid
197 {
198 return false;
199 }
200 if !is_session_bound_local
201 && local.timestamp_ms.abs_diff(ssl.timestamp_ms) > PROMPT_DEDUP_WINDOW_MS
202 {
203 return false;
204 }
205 if let (Some(local_model), Some(ssl_model)) =
206 (local.subject.as_deref(), ssl.subject.as_deref())
207 && local_model != ssl_model
208 {
209 return false;
210 }
211 let Some(ssl_text) = prompt_text_from_details(&ssl.details) else {
212 return false;
213 };
214 prompt_texts_match_or_truncated(&local_text, &ssl_text)
215 });
216 if !duplicate {
217 ssl_rows.push(local);
218 }
219 }
220}
221
222fn prompt_texts_match_or_truncated(left: &str, right: &str) -> bool {
223 if left.eq_ignore_ascii_case(right) {
224 return true;
225 }
226 let min_len = left.len().min(right.len());
227 if min_len < 48 {
228 return false;
229 }
230 let left_lower = left.to_ascii_lowercase();
231 let right_lower = right.to_ascii_lowercase();
232 left_lower.starts_with(&right_lower) || right_lower.starts_with(&left_lower)
233}
234
235fn local_prompt_llm_call_rows(prompt_rows: &[AuditEventRow]) -> Vec<LlmCallRow> {
236 prompt_rows
237 .iter()
238 .filter(|row| {
239 row.details.get("prompt_source").and_then(Value::as_str) == Some("local")
240 && row
241 .details
242 .get("session_id")
243 .and_then(Value::as_str)
244 .is_none()
245 && row.audit_type == "llm"
246 && row.action.as_deref() == Some("request")
247 })
248 .filter_map(local_prompt_llm_call_row)
249 .collect()
250}
251
252fn local_prompt_llm_call_row(row: &AuditEventRow) -> Option<LlmCallRow> {
253 let text = prompt_text_from_details(&row.details)?;
254 let session_id = row
255 .details
256 .get("session_id")
257 .and_then(Value::as_str)
258 .map(ToString::to_string);
259 let conversation_id = row
260 .details
261 .get("conversation_id")
262 .and_then(Value::as_str)
263 .map(ToString::to_string);
264 Some(LlmCallRow {
265 id: format!("llm-{}", row.id),
266 session_id,
267 conversation_id,
268 start_timestamp_ms: row.timestamp_ms,
269 end_timestamp_ms: None,
270 pid: row.pid,
271 comm: row.comm.clone(),
272 provider: None,
273 model: row.subject.clone().or_else(|| row.comm.clone()),
274 call_kind: Some("agent_native_prompt".to_string()),
275 status: row.status.clone().unwrap_or_else(|| "observed".to_string()),
276 error_type: None,
277 finish_reason: None,
278 host: None,
279 path: row.target.clone(),
280 status_code: None,
281 input_tokens: 0,
282 output_tokens: 0,
283 total_tokens: 0,
284 request: json!({
285 "prompt": text,
286 "prompt_source": "local",
287 "target": row.target.as_deref(),
288 }),
289 response: Value::Null,
290 })
291}
292
293fn prompt_text_from_details(details: &Value) -> Option<String> {
294 details
295 .get("text_content")
296 .and_then(Value::as_str)
297 .or_else(|| details.get("prompt").and_then(Value::as_str))
298 .and_then(clean_prompt_text)
299}
300
301#[cfg(test)]
302mod tests {
303 use super::*;
304 use crate::model::ViewSink;
305 use serde_json::json;
306
307 #[test]
308 fn dedupes_local_prompt_only_when_ssl_matches_model_and_text() {
309 for (name, local_model, local_details, expected_rows) in [
310 (
311 "same model and text",
312 Some("claude-opus-4-6"),
313 json!({"text_content": "Run the command.", "prompt_source": "local"}),
314 1,
315 ),
316 (
317 "legacy prompt field",
318 Some("claude-opus-4-6"),
319 json!({"prompt": "Run the command.", "prompt_source": "local"}),
320 1,
321 ),
322 (
323 "different model",
324 Some("claude-haiku-4-5"),
325 json!({"text_content": "Run the command.", "prompt_source": "local"}),
326 2,
327 ),
328 (
329 "missing model",
330 None,
331 json!({"text_content": "Run the command.", "prompt_source": "local"}),
332 1,
333 ),
334 ] {
335 let ssl_rows = [ssl_call_row("claude-opus-4-6", "Run the command.")];
336 let mut prompt_rows = llm_call_prompt_rows(&ssl_rows);
337 let mut local =
338 local_prompt_row("local-prompt", 1_500, local_model, "Run the command.");
339 local.details = local_details;
340
341 append_deduped_local_session_prompt_rows(&mut prompt_rows, vec![local]);
342
343 assert_eq!(prompt_rows.len(), expected_rows, "{name}");
344 }
345 }
346
347 #[test]
348 fn prompt_text_dedupe_accepts_truncated_prefix() {
349 assert!(prompt_texts_match_or_truncated(
350 "Reply with exactly: agentsight-codex-ignore-user-c",
351 "Reply with exactly: agentsight-codex-ignore-user-config",
352 ));
353 assert!(!prompt_texts_match_or_truncated(
354 "Reply with exactly: agentsight-codex-",
355 "Reply with exactly: agentsight-codex-real-smoke",
356 ));
357 }
358
359 #[test]
360 fn observed_codex_exec_prompt_reprojects_as_llm_call() {
361 let temp = tempfile::tempdir().unwrap();
362 let db = temp.path().join("codex.db");
363 let store = SqliteStore::open(&db).unwrap();
364 store
365 .connection()
366 .execute(
367 "INSERT INTO audit_events (
368 id, timestamp_ms, audit_type, pid, comm, action, target, status, details_json
369 ) VALUES (
370 'audit-1', 1000, 'process', 42, 'codex', 'exec', '/tmp/tools/bin/codex',
371 'observed',
372 '{\"full_command\":\"/tmp/tools/bin/codex exec --skip-git-repo-check -c model=gpt agentsight local codex prompt\"}'
373 )",
374 [],
375 )
376 .unwrap();
377
378 let view = load_view_with_observed_session_prompts(&db).unwrap();
379 let rows = view.llm_call_rows(10);
380
381 assert_eq!(rows.len(), 1);
382 assert_eq!(rows[0].comm.as_deref(), Some("codex"));
383 assert_eq!(
384 rows[0].request.get("prompt").and_then(Value::as_str),
385 Some("agentsight local codex prompt")
386 );
387 }
388
389 #[test]
390 fn codex_exec_prompt_dedupes_against_ssl_row_without_local_model() {
391 let temp = tempfile::tempdir().unwrap();
392 let db = temp.path().join("codex.db");
393 let mut store = SqliteStore::open(&db).unwrap();
394 store
395 .llm_call(&ssl_call_row(
396 "gpt-agentsight-mock",
397 "agentsight local codex prompt",
398 ))
399 .unwrap();
400 store
401 .audit_event(&AuditEventRow {
402 id: "audit-1".to_string(),
403 timestamp_ms: 1_500,
404 audit_type: "process".to_string(),
405 pid: Some(42),
406 comm: Some("codex".to_string()),
407 subject: None,
408 action: Some("exec".to_string()),
409 target: Some("/tmp/tools/bin/codex".to_string()),
410 status: Some("observed".to_string()),
411 summary: None,
412 details: json!({
413 "full_command": concat!(
414 "/tmp/tools/bin/codex exec --skip-git-repo-check ",
415 "-c model=gpt agentsight local codex prompt"
416 ),
417 }),
418 })
419 .unwrap();
420 drop(store);
421
422 let view = load_view_with_observed_session_prompts(&db).unwrap();
423 let rows = view.llm_call_rows(10);
424 let prompts = view.audit_rows(Some("llm"), 10);
425
426 assert_eq!(rows.len(), 1);
427 assert_eq!(rows[0].id, "ssl-call");
428 assert_eq!(prompts.len(), 1);
429 assert_eq!(
430 prompts[0]
431 .details
432 .get("prompt_source")
433 .and_then(Value::as_str),
434 Some("ssl")
435 );
436 }
437
438 #[test]
439 fn observed_codex_home_reprojects_session_prompt_as_llm_call() {
440 let temp = tempfile::tempdir().unwrap();
441 let state_path = write_codex_home(temp.path(), "agentsight inferred codex home prompt");
442 let db = temp.path().join("codex.db");
443 let store = SqliteStore::open(&db).unwrap();
444 insert_exec_event(
445 &store,
446 current_epoch_ms(),
447 "/usr/bin/codex exec --skip-git-repo-check agentsight inferred codex home prompt",
448 );
449 insert_file_event(&store, current_epoch_ms() + 100, &state_path);
450
451 let view = load_view_with_observed_session_prompts(&db).unwrap();
452 let rows = view.llm_call_rows(10);
453 let snapshot = view.export_snapshot(crate::model::SnapshotOptions { audit_limit: 10 });
454
455 assert_eq!(rows.len(), 1);
456 assert_eq!(rows[0].comm.as_deref(), Some("codex"));
457 assert_eq!(
458 rows[0].request.get("prompt").and_then(Value::as_str),
459 Some("agentsight inferred codex home prompt")
460 );
461 assert_eq!(snapshot.sessions.len(), 1);
462 }
463
464 #[test]
465 fn observed_codex_home_without_exec_prompt_does_not_import_sessions() {
466 let temp = tempfile::tempdir().unwrap();
467 let state_path = write_codex_home(temp.path(), "agentsight unrelated codex home prompt");
468 let db = temp.path().join("codex.db");
469 let store = SqliteStore::open(&db).unwrap();
470 insert_file_event(&store, current_epoch_ms(), &state_path);
471
472 let view = load_view_with_observed_session_prompts(&db).unwrap();
473
474 assert!(view.llm_call_rows(10).is_empty());
475 }
476
477 fn ssl_call_row(model: &str, text: &str) -> LlmCallRow {
478 LlmCallRow {
479 id: "ssl-call".to_string(),
480 session_id: None,
481 conversation_id: None,
482 start_timestamp_ms: 1_000,
483 end_timestamp_ms: None,
484 pid: Some(42),
485 comm: Some("HTTP Client".to_string()),
486 provider: Some("anthropic".to_string()),
487 model: Some(model.to_string()),
488 call_kind: Some("messages".to_string()),
489 status: "pending".to_string(),
490 error_type: None,
491 finish_reason: None,
492 host: Some("api.anthropic.com".to_string()),
493 path: Some("/v1/messages".to_string()),
494 status_code: None,
495 input_tokens: 0,
496 output_tokens: 0,
497 total_tokens: 0,
498 request: json!({
499 "model": model,
500 "messages": [
501 {
502 "role": "user",
503 "content": [
504 {
505 "type": "text",
506 "text": text
507 }
508 ]
509 }
510 ]
511 }),
512 response: Value::Null,
513 }
514 }
515
516 fn local_prompt_row(
517 id: &str,
518 timestamp_ms: u64,
519 model: Option<&str>,
520 text: &str,
521 ) -> AuditEventRow {
522 AuditEventRow {
523 id: id.to_string(),
524 timestamp_ms,
525 audit_type: "llm".to_string(),
526 pid: Some(42),
527 comm: Some("claude".to_string()),
528 subject: model.map(ToString::to_string),
529 action: Some("request".to_string()),
530 target: agent_session::fixture_session_path(
531 agent_session::AGENT_CLAUDE,
532 std::path::Path::new("/home/user"),
533 )
534 .map(|path| path.to_string_lossy().to_string()),
535 status: Some("observed".to_string()),
536 summary: Some(text.to_string()),
537 details: json!({
538 "text_content": text,
539 "prompt_source": "local"
540 }),
541 }
542 }
543
544 fn current_epoch_ms() -> u64 {
545 std::time::SystemTime::now()
546 .duration_since(std::time::UNIX_EPOCH)
547 .unwrap()
548 .as_millis() as u64
549 }
550
551 fn write_codex_home(root: &std::path::Path, prompt: &str) -> std::path::PathBuf {
552 let codex_home = root.join("codex-home");
553 let session_dir = codex_home.join("sessions/2026/07/11");
554 std::fs::create_dir_all(&session_dir).unwrap();
555 std::fs::write(
556 session_dir.join("rollout-2026-07-11T00-00-00.jsonl"),
557 format!(
558 "{{\"timestamp\":\"2026-07-11T00:00:00.000Z\",\
559 \"type\":\"event_msg\",\
560 \"payload\":{{\"type\":\"user_message\",\"message\":\"{prompt}\"}}}}\n"
561 ),
562 )
563 .unwrap();
564 let state_path = codex_home.join("stat");
565 std::fs::write(&state_path, "").unwrap();
566 state_path
567 }
568
569 fn insert_exec_event(store: &SqliteStore, timestamp_ms: u64, full_command: &str) {
570 store
571 .connection()
572 .execute(
573 "INSERT INTO audit_events (
574 id, timestamp_ms, audit_type, pid, comm, action, target, status, details_json
575 ) VALUES (
576 'audit-exec-1', ?1, 'process', 42, 'codex', 'exec', '/usr/bin/codex',
577 'observed', ?2
578 )",
579 rusqlite::params![
580 timestamp_ms,
581 json!({"full_command": full_command}).to_string()
582 ],
583 )
584 .unwrap();
585 }
586
587 fn insert_file_event(store: &SqliteStore, timestamp_ms: u64, path: &std::path::Path) {
588 store
589 .connection()
590 .execute(
591 "INSERT INTO audit_events (
592 id, timestamp_ms, audit_type, pid, comm, action, target, status, details_json
593 ) VALUES (
594 'audit-file-1', ?1, 'file', 42, 'codex', 'write', ?2, 'observed', ?3
595 )",
596 rusqlite::params![
597 timestamp_ms,
598 path.to_string_lossy().as_ref(),
599 json!({"filepath": path.to_string_lossy()}).to_string(),
600 ],
601 )
602 .unwrap();
603 }
604}