Skip to main content

agentsight_capture/sources/
sqlite.rs

1// SPDX-License-Identifier: MIT
2// Copyright (c) 2026 eunomia-bpf org.
3
4use 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(&current_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}