Skip to main content

agentsight_capture/view/
projection.rs

1// SPDX-License-Identifier: MIT
2// Copyright (c) 2026 eunomia-bpf org.
3
4use crate::event::Event;
5use crate::json::{i64_field as json_i64, parse_optional_value as parse_optional_json};
6use crate::model::{
7    AuditEventRow, LlmCallRow, NetworkTargetRow, ProcessNodeRow, ResourceSampleRow, TokenUsageRow,
8    ToolCallRow, ViewResult,
9};
10use crate::text::sanitize_ascii_identifier as sanitize_id;
11use crate::view::llm::TokenUsage;
12use crate::view::{
13    CanonicalEvent, EventKind, body_json, extract_model, extract_token_usage,
14    extract_token_usage_from_sse, normalize_event, provider_from_host,
15};
16use crate::view::{MaterializedView, PendingRequest};
17use serde_json::Value;
18
19const PENDING_REQUEST_TTL_MS: u64 = 5 * 60 * 1000;
20const MAX_PENDING_REQUESTS_PER_STREAM: usize = 16;
21
22impl MaterializedView {
23    pub fn ingest_event(&mut self, event: &Event) -> ViewResult<()> {
24        self.next_seq += 1;
25        let raw_id = format!(
26            "event-{}-{}-{}-{}",
27            event.timestamp,
28            sanitize_id(&event.source),
29            event.pid,
30            self.next_seq
31        );
32        let canonical = normalize_event(event, raw_id);
33        self.prune_pending(canonical.timestamp_ms);
34        if let Some(sample) = resource_sample_from_event(&canonical) {
35            self.emit_resource_sample(sample)?;
36        }
37        if let Some(target) = network_target_from_event(&canonical) {
38            self.emit_network_target(target)?;
39        }
40        self.ingest(&canonical)
41    }
42
43    fn ingest(&mut self, event: &CanonicalEvent) -> ViewResult<()> {
44        self.ingest_agent_specific_event(event)?;
45        match event.kind {
46            EventKind::LlmRequest => self.ingest_llm_request(event),
47            EventKind::LlmResponse | EventKind::LlmError => self.ingest_llm_response(event),
48            EventKind::HttpResponse if self.has_pending_llm_request(event) => {
49                self.ingest_llm_response(event)
50            }
51            EventKind::ProcessExec => self.ingest_process_audit(event, "exec"),
52            EventKind::ProcessExit => self.ingest_process_audit(event, "exit"),
53            EventKind::FsOpen if is_writable_open(event) => self.ingest_file_audit(event),
54            EventKind::FsWrite | EventKind::FsMutation => self.ingest_file_audit(event),
55            EventKind::Unknown if is_process_summary_write_event(event) => {
56                self.ingest_file_audit(event)
57            }
58            EventKind::Unknown if is_process_network_event(event) => {
59                self.ingest_network_audit(event)
60            }
61            _ => Ok(()),
62        }
63    }
64
65    fn ingest_llm_request(&mut self, event: &CanonicalEvent) -> ViewResult<()> {
66        let (Some(pid), Some(tid)) = (event.pid, event.tid) else {
67            return Ok(());
68        };
69        let req = PendingRequest {
70            event_id: event.event_id.clone(),
71            timestamp_ms: event.timestamp_ms,
72            pid,
73            comm: event.comm.clone().unwrap_or_default(),
74            provider: event.provider.clone(),
75            model: event.model.clone(),
76            host: event.host.clone(),
77            path: event.path.clone(),
78            request_id: event.request_id.clone(),
79            body_json: body_json(&event.attributes),
80        };
81        if req.body_json.is_none() && req.model.is_none() {
82            return Ok(());
83        }
84        self.insert_orphan_llm_request(&req)?;
85        let requests = self.pending.entry((pid, tid)).or_default();
86        requests.push_back(req);
87        while requests.len() > MAX_PENDING_REQUESTS_PER_STREAM {
88            requests.pop_front();
89        }
90        Ok(())
91    }
92
93    fn ingest_llm_response(&mut self, event: &CanonicalEvent) -> ViewResult<()> {
94        let Some(pid) = event.pid else {
95            return Ok(());
96        };
97        if let Some(tid) = event.tid
98            && let Some((req, confidence)) = self.take_matching_request(pid, tid, event)
99        {
100            return self.upsert_llm_pair(req, event, confidence);
101        }
102        self.insert_orphan_llm_response(event)
103    }
104
105    fn has_pending_llm_request(&self, event: &CanonicalEvent) -> bool {
106        let (Some(pid), Some(tid)) = (event.pid, event.tid) else {
107            return false;
108        };
109        self.pending
110            .get(&(pid, tid))
111            .is_some_and(|requests| !requests.is_empty())
112    }
113
114    fn take_matching_request(
115        &mut self,
116        pid: u32,
117        tid: u64,
118        resp: &CanonicalEvent,
119    ) -> Option<(PendingRequest, f32)> {
120        let requests = self.pending.get_mut(&(pid, tid))?;
121        let (req, confidence) = if let Some(resp_request_id) = resp.request_id.as_deref() {
122            let pos = requests
123                .iter()
124                .position(|req| req.request_id.as_deref() == Some(resp_request_id))?;
125            (requests.remove(pos)?, 0.95)
126        } else if requests.len() == 1 {
127            (requests.pop_front()?, 0.75)
128        } else {
129            let pos = {
130                let mut candidates = requests
131                    .iter()
132                    .enumerate()
133                    .filter(|(_, req)| req.body_json.is_some() || req.model.is_some())
134                    .map(|(idx, _)| idx);
135                let pos = candidates.next()?;
136                if candidates.next().is_some() {
137                    return None;
138                }
139                pos
140            };
141            (requests.remove(pos)?, 0.7)
142        };
143        if requests.is_empty() {
144            self.pending.remove(&(pid, tid));
145        }
146        Some((req, confidence))
147    }
148
149    fn prune_pending(&mut self, now_ms: u64) {
150        let cutoff = now_ms.saturating_sub(PENDING_REQUEST_TTL_MS);
151        self.pending.retain(|_, requests| {
152            while requests
153                .front()
154                .is_some_and(|req| req.timestamp_ms < cutoff)
155            {
156                requests.pop_front();
157            }
158            !requests.is_empty()
159        });
160    }
161
162    fn upsert_llm_pair(
163        &mut self,
164        req: PendingRequest,
165        resp: &CanonicalEvent,
166        confidence: f32,
167    ) -> ViewResult<()> {
168        let response_body = response_body_json(resp);
169        let model = req
170            .model
171            .clone()
172            .or_else(|| response_body.as_ref().and_then(extract_model))
173            .or_else(|| resp.model.clone())
174            .unwrap_or_else(|| "unknown".to_string());
175        let provider = req
176            .provider
177            .clone()
178            .or_else(|| req.host.as_deref().map(provider_from_host));
179        let llm_call_id = format!("llm-{}", req.event_id);
180        let status_code = resp.status_code;
181        let mut call_row = llm_call_row(
182            &llm_call_id,
183            req.timestamp_ms,
184            Some(resp.timestamp_ms),
185            req.pid,
186            &req.comm,
187            provider.as_deref(),
188            Some(&model),
189            req.host.as_deref(),
190            req.path.as_deref(),
191            status_code,
192            req.body_json.as_ref(),
193            response_body.as_ref(),
194        );
195        if let Some(usage) = self.ingest_response_usage_and_tools(
196            resp,
197            &llm_call_id,
198            req.pid,
199            &req.comm,
200            provider.as_deref(),
201            &model,
202            response_body.as_ref(),
203            confidence,
204        )? {
205            call_row.input_tokens = usage.input_tokens;
206            call_row.output_tokens = usage.output_tokens;
207            call_row.total_tokens = usage.total_tokens;
208        }
209        emit_llm_audit(
210            self,
211            &llm_call_id,
212            resp.timestamp_ms,
213            req.pid,
214            &req.comm,
215            Some(&model),
216            "call",
217            req.host.as_deref(),
218            if status_code.map(|c| c >= 400).unwrap_or(false) {
219                "failure"
220            } else {
221                "success"
222            },
223            "LLM call",
224            response_body.as_ref(),
225        )?;
226        self.emit_llm_call(call_row)
227    }
228
229    fn insert_orphan_llm_request(&mut self, req: &PendingRequest) -> ViewResult<()> {
230        let llm_call_id = format!("llm-{}", req.event_id);
231        let provider = req
232            .provider
233            .clone()
234            .or_else(|| req.host.as_deref().map(provider_from_host));
235        let call_row = llm_call_row(
236            &llm_call_id,
237            req.timestamp_ms,
238            None,
239            req.pid,
240            &req.comm,
241            provider.as_deref(),
242            req.model.as_deref(),
243            req.host.as_deref(),
244            req.path.as_deref(),
245            None,
246            req.body_json.as_ref(),
247            None,
248        );
249        emit_llm_audit(
250            self,
251            &llm_call_id,
252            req.timestamp_ms,
253            req.pid,
254            &req.comm,
255            req.model.as_deref(),
256            "request",
257            req.host.as_deref(),
258            "orphan_request",
259            "LLM request",
260            req.body_json.as_ref(),
261        )?;
262        self.emit_llm_call(call_row)
263    }
264
265    fn insert_orphan_llm_response(&mut self, resp: &CanonicalEvent) -> ViewResult<()> {
266        let response_body = response_body_json(resp);
267        let model = resp
268            .model
269            .clone()
270            .or_else(|| response_body.as_ref().and_then(extract_model))
271            .unwrap_or_else(|| "unknown".to_string());
272        let provider = resp
273            .provider
274            .clone()
275            .or_else(|| resp.host.as_deref().map(provider_from_host));
276        let pid = resp.pid.unwrap_or(0);
277        let comm = resp.comm.clone().unwrap_or_default();
278        let llm_call_id = format!("llm-orphan-{}", resp.event_id);
279        let mut call_row = llm_call_row(
280            &llm_call_id,
281            resp.timestamp_ms,
282            Some(resp.timestamp_ms),
283            pid,
284            &comm,
285            provider.as_deref(),
286            Some(&model),
287            resp.host.as_deref(),
288            resp.path.as_deref(),
289            resp.status_code,
290            None,
291            response_body.as_ref(),
292        );
293        if let Some(usage) = self.ingest_response_usage_and_tools(
294            resp,
295            &llm_call_id,
296            pid,
297            &comm,
298            provider.as_deref(),
299            &model,
300            response_body.as_ref(),
301            0.35,
302        )? {
303            call_row.input_tokens = usage.input_tokens;
304            call_row.output_tokens = usage.output_tokens;
305            call_row.total_tokens = usage.total_tokens;
306        }
307        emit_llm_audit(
308            self,
309            &llm_call_id,
310            resp.timestamp_ms,
311            pid,
312            &comm,
313            Some(&model),
314            "response",
315            resp.host.as_deref(),
316            "orphan_response",
317            "LLM response",
318            response_body.as_ref(),
319        )?;
320        self.emit_llm_call(call_row)
321    }
322
323    #[allow(clippy::too_many_arguments)]
324    fn ingest_response_usage_and_tools(
325        &mut self,
326        resp: &CanonicalEvent,
327        llm_call_id: &str,
328        pid: u32,
329        comm: &str,
330        provider: Option<&str>,
331        model: &str,
332        response_body: Option<&Value>,
333        confidence: f32,
334    ) -> ViewResult<Option<TokenUsageRow>> {
335        let usage = if resp.source == "sse_processor" {
336            extract_token_usage_from_sse(&resp.attributes)
337        } else {
338            response_body.map(extract_token_usage).unwrap_or_default()
339        };
340        let mut usage_row = None;
341        if !usage.is_empty() {
342            let token_id = format!("token-{llm_call_id}");
343            let row = token_usage_row(
344                &token_id,
345                llm_call_id,
346                resp.timestamp_ms,
347                pid,
348                Some(comm),
349                provider,
350                Some(model),
351                &usage,
352                "response_usage",
353                confidence,
354            );
355            self.emit_token_usage(row.clone())?;
356            usage_row = Some(row);
357        }
358        self.ingest_sse_tools(resp, llm_call_id, pid, confidence)?;
359        Ok(usage_row)
360    }
361
362    fn ingest_agent_specific_event(&mut self, event: &CanonicalEvent) -> ViewResult<()> {
363        self.ingest_claude_telemetry(event)?;
364        self.ingest_gemini_stdio_stats(event)
365    }
366
367    fn ingest_claude_telemetry(&mut self, event: &CanonicalEvent) -> ViewResult<()> {
368        let host = event.host.as_deref().unwrap_or_default();
369        if !host.contains("datadoghq.com") && event.source != "ssl" {
370            return Ok(());
371        }
372        let body = body_json(&event.attributes).or_else(|| {
373            event
374                .attributes
375                .get("data")
376                .and_then(|v| v.as_str())
377                .and_then(parse_json_str)
378        });
379        let Some(Value::Array(items)) = body else {
380            return Ok(());
381        };
382        let pid = event.pid.unwrap_or(0);
383        let comm = event.comm.as_deref().unwrap_or_default();
384        for (idx, item) in items.iter().enumerate() {
385            let message = item
386                .get("message")
387                .and_then(Value::as_str)
388                .unwrap_or_default();
389            if message == "tengu_api_success" {
390                let input = json_i64(item, "input_tokens");
391                let output = json_i64(item, "output_tokens");
392                let cache = json_i64(item, "cached_input_tokens");
393                let total = input + output + cache;
394                if total <= 0 {
395                    continue;
396                }
397                let model = item
398                    .get("model")
399                    .and_then(Value::as_str)
400                    .unwrap_or("unknown");
401                let llm_call_id = format!("claude-telemetry-{}-{idx}", event.event_id);
402                let usage = observed_token_usage(input, output, cache, total);
403                self.emit_token_usage(token_usage_row(
404                    &format!("token-{llm_call_id}"),
405                    &llm_call_id,
406                    event.timestamp_ms,
407                    pid,
408                    Some(comm),
409                    Some("anthropic"),
410                    Some(model),
411                    &usage,
412                    "claude_telemetry",
413                    0.80,
414                ))?;
415            } else if message == "tengu_tool_use_success" {
416                let tool_name = item.get("tool_name").and_then(Value::as_str).unwrap_or("?");
417                let duration_ms = item
418                    .get("duration_ms")
419                    .and_then(Value::as_i64)
420                    .map(|v| v as u64);
421                let request_id = item.get("request_id").and_then(Value::as_str);
422                self.emit_tool_call(ToolCallRow {
423                    id: format!("claude-tool-telemetry-{}-{idx}", event.event_id),
424                    session_id: None,
425                    conversation_id: None,
426                    timestamp_ms: event.timestamp_ms,
427                    tool_name: Some(tool_name.to_string()),
428                    tool_call_id: request_id.map(str::to_string),
429                    start_timestamp_ms: duration_ms.and_then(|d| event.timestamp_ms.checked_sub(d)),
430                    end_timestamp_ms: Some(event.timestamp_ms),
431                    duration_ms,
432                    status: Some("completed".to_string()),
433                    input: serde_json::json!({}),
434                    output: serde_json::json!({}),
435                    related_pid: Some(pid),
436                    related_event_id: Some(event.event_id.clone()),
437                    view_source: "view".to_string(),
438                    confidence: Some(0.75),
439                })?;
440            }
441        }
442        Ok(())
443    }
444
445    fn ingest_gemini_stdio_stats(&mut self, event: &CanonicalEvent) -> ViewResult<()> {
446        if !matches!(event.kind, EventKind::StdioMessage | EventKind::StdioRpc) {
447            return Ok(());
448        }
449        let Some(payload) = event.attributes.get("data").and_then(Value::as_str) else {
450            return Ok(());
451        };
452        let Some(obj) = parse_json_str(payload) else {
453            return Ok(());
454        };
455        let Some(models) = obj.pointer("/stats/models").and_then(Value::as_object) else {
456            return Ok(());
457        };
458        let pid = event.pid.unwrap_or(0);
459        let comm = event.comm.as_deref().unwrap_or("gemini");
460        for (model, stats) in models {
461            let tokens = stats.get("tokens").unwrap_or(stats);
462            let input = json_i64(tokens, "prompt").max(json_i64(tokens, "input"));
463            let output = json_i64(tokens, "candidates")
464                + json_i64(tokens, "thoughts")
465                + json_i64(tokens, "tool");
466            let cache = json_i64(tokens, "cached");
467            let total = json_i64(tokens, "total").max(input + output + cache);
468            if total <= 0 {
469                continue;
470            }
471            let llm_call_id = format!("gemini-stdout-{}-{}", event.event_id, sanitize_id(model));
472            let usage = observed_token_usage(input, output, cache, total);
473            self.emit_token_usage(token_usage_row(
474                &format!("token-{llm_call_id}"),
475                &llm_call_id,
476                event.timestamp_ms,
477                pid,
478                Some(comm),
479                Some("gcp.gen_ai"),
480                Some(model),
481                &usage,
482                "gemini_cli_stdout_stats",
483                0.85,
484            ))?;
485        }
486        Ok(())
487    }
488
489    fn ingest_sse_tools(
490        &mut self,
491        event: &CanonicalEvent,
492        llm_call_id: &str,
493        pid: u32,
494        confidence: f32,
495    ) -> ViewResult<()> {
496        if let Some(openai_json) = event
497            .attributes
498            .get("json_content")
499            .and_then(Value::as_str)
500            .and_then(parse_json_str)
501            && let Some(tool_calls) = openai_json.get("tool_calls").and_then(Value::as_array)
502        {
503            for (idx, tool_call) in tool_calls.iter().enumerate() {
504                let function = tool_call.get("function").unwrap_or(&Value::Null);
505                let name = function.get("name").and_then(Value::as_str).unwrap_or("?");
506                let tool_call_id = tool_call.get("id").and_then(Value::as_str);
507                let arguments = function.get("arguments").and_then(Value::as_str);
508                let tool_id = tool_call_id
509                    .map(str::to_string)
510                    .unwrap_or_else(|| format!("openai-tool-{idx}"));
511                self.emit_tool_call(ToolCallRow {
512                    id: format!("tool-{llm_call_id}-{tool_id}"),
513                    session_id: None,
514                    conversation_id: Some(format!("conv-{llm_call_id}")),
515                    timestamp_ms: event.timestamp_ms,
516                    tool_name: Some(name.to_string()),
517                    tool_call_id: tool_call_id.map(str::to_string),
518                    start_timestamp_ms: Some(event.timestamp_ms),
519                    end_timestamp_ms: None,
520                    duration_ms: None,
521                    status: Some("observed".to_string()),
522                    input: parse_optional_json(arguments),
523                    output: Value::Null,
524                    related_pid: Some(pid),
525                    related_event_id: Some(event.event_id.clone()),
526                    view_source: "view".to_string(),
527                    confidence: Some(confidence),
528                })?;
529            }
530        }
531
532        let Some(events) = event.attributes.get("sse_events").and_then(Value::as_array) else {
533            return Ok(());
534        };
535        for (idx, sse) in events.iter().enumerate() {
536            let Some(block) = sse.pointer("/parsed_data/content_block") else {
537                continue;
538            };
539            if block.get("type").and_then(Value::as_str) != Some("tool_use") {
540                continue;
541            }
542            let name = block.get("name").and_then(Value::as_str).unwrap_or("?");
543            let tool_call_id = block.get("id").and_then(Value::as_str);
544            let input_json = block.get("input").map(Value::to_string);
545            let tool_id = tool_call_id
546                .map(str::to_string)
547                .unwrap_or_else(|| format!("tool-{idx}"));
548            self.emit_tool_call(ToolCallRow {
549                id: format!("tool-{llm_call_id}-{tool_id}"),
550                session_id: None,
551                conversation_id: Some(format!("conv-{llm_call_id}")),
552                timestamp_ms: event.timestamp_ms,
553                tool_name: Some(name.to_string()),
554                tool_call_id: tool_call_id.map(str::to_string),
555                start_timestamp_ms: Some(event.timestamp_ms),
556                end_timestamp_ms: None,
557                duration_ms: None,
558                status: Some("observed".to_string()),
559                input: parse_optional_json(input_json.as_deref()),
560                output: Value::Null,
561                related_pid: Some(pid),
562                related_event_id: Some(event.event_id.clone()),
563                view_source: "view".to_string(),
564                confidence: Some(confidence),
565            })?;
566        }
567        Ok(())
568    }
569
570    fn ingest_process_audit(&mut self, event: &CanonicalEvent, action: &str) -> ViewResult<()> {
571        let target = event.attributes.get("filename").and_then(Value::as_str);
572        self.emit_audit_event(AuditEventRow {
573            id: format!("audit-{}", event.event_id),
574            timestamp_ms: event.timestamp_ms,
575            audit_type: "process".to_string(),
576            pid: event.pid,
577            comm: event.comm.clone(),
578            subject: event.comm.clone(),
579            action: Some(action.to_string()),
580            target: target.map(str::to_string),
581            status: Some(process_audit_status(action, &event.attributes).to_string()),
582            summary: event.summary.clone(),
583            details: event.attributes.clone(),
584        })?;
585        if let Some(row) = self
586            .process_node_id(event, action)
587            .and_then(|id| process_node_from_event(event, action, id))
588        {
589            self.emit_process_node(row)?;
590        }
591        Ok(())
592    }
593
594    fn process_node_id(&mut self, event: &CanonicalEvent, action: &str) -> Option<String> {
595        let pid = event.pid?;
596        match action {
597            "exec" => {
598                let id = self
599                    .active_processes
600                    .entry(pid)
601                    .or_insert_with(|| format!("process-{pid}-{}", event.timestamp_ms));
602                Some(id.clone())
603            }
604            "exit" => Some(
605                self.active_processes
606                    .remove(&pid)
607                    .unwrap_or_else(|| format!("process-{pid}-{}", event.timestamp_ms)),
608            ),
609            _ => None,
610        }
611    }
612
613    fn ingest_file_audit(&mut self, event: &CanonicalEvent) -> ViewResult<()> {
614        let target = event
615            .attributes
616            .get("path")
617            .or_else(|| event.attributes.get("filepath"))
618            .and_then(Value::as_str);
619        self.emit_audit_event(AuditEventRow {
620            id: format!("audit-{}", event.event_id),
621            timestamp_ms: event.timestamp_ms,
622            audit_type: "file".to_string(),
623            pid: event.pid,
624            comm: event.comm.clone(),
625            subject: event.comm.clone(),
626            action: Some("write".to_string()),
627            target: target.map(str::to_string),
628            status: Some("observed".to_string()),
629            summary: event.summary.clone(),
630            details: event.attributes.clone(),
631        })
632    }
633
634    fn ingest_network_audit(&mut self, event: &CanonicalEvent) -> ViewResult<()> {
635        let target = event
636            .attributes
637            .get("detail")
638            .or_else(|| event.attributes.get("host"))
639            .and_then(Value::as_str);
640        let action = process_network_action(&event.attributes).unwrap_or("network");
641        self.emit_audit_event(AuditEventRow {
642            id: format!("audit-{}", event.event_id),
643            timestamp_ms: event.timestamp_ms,
644            audit_type: "network".to_string(),
645            pid: event.pid,
646            comm: event.comm.clone(),
647            subject: event.comm.clone(),
648            action: Some(action.to_string()),
649            target: target.map(str::to_string),
650            status: Some("observed".to_string()),
651            summary: event.summary.clone(),
652            details: event.attributes.clone(),
653        })
654    }
655}
656
657#[allow(clippy::too_many_arguments)]
658fn emit_llm_audit(
659    view: &mut MaterializedView,
660    llm_call_id: &str,
661    timestamp_ms: u64,
662    pid: u32,
663    comm: &str,
664    subject: Option<&str>,
665    action: &str,
666    target: Option<&str>,
667    status: &str,
668    summary: &str,
669    details: Option<&Value>,
670) -> ViewResult<()> {
671    view.emit_audit_event(AuditEventRow {
672        id: format!("audit-{llm_call_id}-{action}"),
673        timestamp_ms,
674        audit_type: "llm".to_string(),
675        pid: Some(pid),
676        comm: Some(comm.to_string()),
677        subject: subject.map(str::to_string),
678        action: Some(action.to_string()),
679        target: target.map(str::to_string),
680        status: Some(status.to_string()),
681        summary: Some(summary.to_string()),
682        details: details.cloned().unwrap_or_else(|| serde_json::json!({})),
683    })
684}
685
686#[allow(clippy::too_many_arguments)]
687fn token_usage_row(
688    id: &str,
689    llm_call_id: &str,
690    timestamp_ms: u64,
691    pid: u32,
692    comm: Option<&str>,
693    provider: Option<&str>,
694    model: Option<&str>,
695    usage: &TokenUsage,
696    source: &str,
697    confidence: f32,
698) -> TokenUsageRow {
699    TokenUsageRow {
700        id: id.to_string(),
701        llm_call_id: llm_call_id.to_string(),
702        timestamp_ms,
703        pid: Some(pid),
704        comm: comm.map(str::to_string),
705        provider: provider.map(str::to_string),
706        model: model.map(str::to_string),
707        input_tokens: usage.input_tokens,
708        output_tokens: usage.output_tokens,
709        cache_creation_tokens: usage.cache_creation_tokens,
710        cache_read_tokens: usage.cache_read_tokens,
711        total_tokens: usage.total_tokens(),
712        source: source.to_string(),
713        view_source: "view".to_string(),
714        confidence: Some(confidence),
715    }
716}
717
718fn observed_token_usage(input: i64, output: i64, cache_read: i64, total: i64) -> TokenUsage {
719    TokenUsage {
720        input_tokens: input,
721        output_tokens: output,
722        cache_read_tokens: cache_read,
723        total_override: Some(total),
724        ..Default::default()
725    }
726}
727
728fn response_body_json(event: &CanonicalEvent) -> Option<Value> {
729    body_json(&event.attributes)
730        .or_else(|| (event.source == "sse_processor").then(|| event.attributes.clone()))
731}
732
733fn process_audit_status(action: &str, attributes: &Value) -> &'static str {
734    if action != "exit" {
735        return "observed";
736    }
737    match attributes.get("exit_code").and_then(Value::as_i64) {
738        Some(0) => "success",
739        Some(_) => "failure",
740        None => "observed",
741    }
742}
743
744fn process_node_from_event(
745    event: &CanonicalEvent,
746    action: &str,
747    id: String,
748) -> Option<ProcessNodeRow> {
749    let pid = event.pid?;
750    let status = process_audit_status(action, &event.attributes).to_string();
751    let argv = process_argv(&event.attributes);
752    Some(ProcessNodeRow {
753        id,
754        pid,
755        ppid: event.ppid,
756        root_pid: None,
757        start_timestamp_ms: (action == "exec").then_some(event.timestamp_ms),
758        end_timestamp_ms: (action == "exit").then_some(event.timestamp_ms),
759        comm: event.comm.clone(),
760        command: process_command(&event.attributes, &argv),
761        argv,
762        cwd: event
763            .attributes
764            .get("cwd")
765            .and_then(Value::as_str)
766            .map(str::to_string),
767        exit_code: (action == "exit")
768            .then(|| event.attributes.get("exit_code").and_then(Value::as_i64))
769            .flatten()
770            .map(|value| value as i32),
771        status: Some(status),
772        view_source: "view".to_string(),
773        confidence: event.confidence,
774    })
775}
776
777fn process_command(attributes: &Value, argv: &[String]) -> Option<String> {
778    attributes
779        .get("filename")
780        .and_then(Value::as_str)
781        .or_else(|| attributes.get("command").and_then(Value::as_str))
782        .map(str::to_string)
783        .or_else(|| argv.first().cloned())
784}
785
786fn process_argv(attributes: &Value) -> Vec<String> {
787    attributes
788        .get("argv")
789        .and_then(Value::as_array)
790        .map(|argv| {
791            argv.iter()
792                .filter_map(Value::as_str)
793                .map(str::to_string)
794                .collect()
795        })
796        .unwrap_or_default()
797}
798
799fn is_writable_open(event: &CanonicalEvent) -> bool {
800    let flags = event
801        .attributes
802        .get("flags")
803        .and_then(Value::as_i64)
804        .unwrap_or(0);
805    const O_ACCMODE: i64 = 0o3;
806    const O_CREAT: i64 = 0o100;
807    const O_TRUNC: i64 = 0o1000;
808    const O_APPEND: i64 = 0o2000;
809    (flags & O_ACCMODE) != 0 || (flags & (O_CREAT | O_TRUNC | O_APPEND)) != 0
810}
811
812fn is_process_summary_write_event(event: &CanonicalEvent) -> bool {
813    event.source == "process"
814        && process_event_name(&event.attributes) == Some("SUMMARY")
815        && event.attributes.get("type").and_then(Value::as_str) == Some("WRITE")
816}
817
818fn is_process_network_event(event: &CanonicalEvent) -> bool {
819    event.source == "process"
820        && process_network_action(&event.attributes).is_some_and(|name| name.starts_with("NET_"))
821}
822
823fn process_event_name(attributes: &Value) -> Option<&str> {
824    attributes.get("event").and_then(Value::as_str)
825}
826
827fn process_network_action(attributes: &Value) -> Option<&str> {
828    let event = process_event_name(attributes)?;
829    if event == "SUMMARY" {
830        attributes.get("type").and_then(Value::as_str)
831    } else {
832        Some(event)
833    }
834}
835
836fn parse_json_str(text: &str) -> Option<Value> {
837    serde_json::from_str(text).ok()
838}
839
840fn number_or_string(value: Option<&Value>) -> Option<f64> {
841    value.and_then(|v| {
842        v.as_f64()
843            .or_else(|| v.as_str().and_then(|s| s.parse::<f64>().ok()))
844    })
845}
846
847fn network_target_from_event(event: &CanonicalEvent) -> Option<NetworkTargetRow> {
848    let host = event.host.as_deref().filter(|host| !host.is_empty())?;
849    let path = event.path.as_deref().filter(|path| !path.is_empty());
850    let error_count = i64::from(
851        event.kind == EventKind::LlmError
852            || event.status_code.map(|code| code >= 400).unwrap_or(false),
853    );
854    Some(NetworkTargetRow {
855        pid: event.pid,
856        comm: event.comm.clone(),
857        host: host.to_string(),
858        path: path.map(str::to_string),
859        count: 1,
860        error_count,
861        first_timestamp_ms: Some(event.timestamp_ms),
862        last_timestamp_ms: Some(event.timestamp_ms),
863    })
864}
865
866fn resource_sample_from_event(event: &CanonicalEvent) -> Option<ResourceSampleRow> {
867    if event.kind != EventKind::ResourceSample {
868        return None;
869    }
870    let cpu = number_or_string(event.attributes.get("cpu").and_then(|v| v.get("percent")));
871    let rss_mb = number_or_string(event.attributes.get("memory").and_then(|v| v.get("rss_mb")));
872    Some(ResourceSampleRow {
873        timestamp_ms: event.timestamp_ms,
874        pid: event.pid,
875        comm: event.comm.clone(),
876        cpu_percent: cpu,
877        rss_mb: rss_mb.map(|v| v.max(0.0) as i64),
878    })
879}
880
881#[allow(clippy::too_many_arguments)]
882fn llm_call_row(
883    id: &str,
884    start_timestamp_ms: u64,
885    end_timestamp_ms: Option<u64>,
886    pid: u32,
887    comm: &str,
888    provider: Option<&str>,
889    model: Option<&str>,
890    host: Option<&str>,
891    path: Option<&str>,
892    status_code: Option<u16>,
893    request_body: Option<&Value>,
894    response_body: Option<&Value>,
895) -> LlmCallRow {
896    let status = llm_call_status(end_timestamp_ms, status_code);
897    let error_type = status_code
898        .filter(|code| *code >= 400)
899        .map(|code| format!("http_{code}"));
900    let session_id = request_body
901        .and_then(session_id_from_body)
902        .or_else(|| response_body.and_then(session_id_from_body));
903    let conversation_id = request_body
904        .and_then(conversation_id_from_body)
905        .or_else(|| response_body.and_then(conversation_id_from_body));
906    LlmCallRow {
907        id: id.to_string(),
908        session_id,
909        conversation_id,
910        start_timestamp_ms,
911        end_timestamp_ms,
912        pid: Some(pid),
913        comm: Some(comm.to_string()),
914        provider: provider.map(str::to_string),
915        model: model.map(str::to_string),
916        call_kind: path.and_then(call_kind_from_path).map(str::to_string),
917        status,
918        error_type,
919        finish_reason: response_body.and_then(finish_reason_from_body),
920        host: host.map(str::to_string),
921        path: path.map(str::to_string),
922        status_code,
923        input_tokens: 0,
924        output_tokens: 0,
925        total_tokens: 0,
926        request: request_body.cloned().unwrap_or(Value::Null),
927        response: response_body.cloned().unwrap_or(Value::Null),
928    }
929}
930
931fn llm_call_status(end_timestamp_ms: Option<u64>, status_code: Option<u16>) -> String {
932    if end_timestamp_ms.is_none() {
933        return "pending".to_string();
934    }
935    if status_code.map(|code| code >= 400).unwrap_or(false) {
936        "error".to_string()
937    } else {
938        "complete".to_string()
939    }
940}
941
942fn call_kind_from_path(path: &str) -> Option<&'static str> {
943    if path.contains("/v1/responses") || path.contains("/codex/responses") {
944        Some("responses")
945    } else if path.contains("/v1/messages") {
946        Some("messages")
947    } else if path.contains("/chat/completions") {
948        Some("chat")
949    } else if path.contains(":streamGenerateContent") {
950        Some("stream_generate_content")
951    } else if path.contains(":generateContent") {
952        Some("generate_content")
953    } else {
954        None
955    }
956}
957
958fn session_id_from_body(body: &Value) -> Option<String> {
959    string_at(body, &["session_id"])
960        .or_else(|| string_at(body, &["sessionId"]))
961        .or_else(|| string_at(body, &["metadata", "session_id"]))
962        .or_else(|| string_at(body, &["metadata", "sessionId"]))
963        .or_else(|| metadata_user_session_id(body))
964}
965
966fn conversation_id_from_body(body: &Value) -> Option<String> {
967    string_at(body, &["conversation_id"])
968        .or_else(|| string_at(body, &["conversationId"]))
969        .or_else(|| string_at(body, &["metadata", "conversation_id"]))
970        .or_else(|| string_at(body, &["metadata", "conversationId"]))
971        .or_else(|| string_at(body, &["response", "id"]))
972        .or_else(|| string_at(body, &["id"]))
973        .or_else(|| string_at(body, &["message_id"]))
974        .or_else(|| string_at(body, &["message", "id"]))
975}
976
977fn metadata_user_session_id(body: &Value) -> Option<String> {
978    let user_id = string_at(body, &["metadata", "user_id"])
979        .or_else(|| string_at(body, &["metadata", "userId"]))?;
980    if let Ok(json) = serde_json::from_str::<Value>(&user_id)
981        && let Some(session) = session_id_from_body(&json)
982    {
983        return Some(session);
984    }
985    user_id
986        .split("_session_")
987        .nth(1)
988        .map(str::to_string)
989        .filter(|value| !value.is_empty())
990}
991
992fn finish_reason_from_body(body: &Value) -> Option<String> {
993    string_at(body, &["finish_reason"])
994        .or_else(|| string_at(body, &["stop_reason"]))
995        .or_else(|| string_at(body, &["choices", "0", "finish_reason"]))
996        .or_else(|| string_at(body, &["candidates", "0", "finishReason"]))
997        .or_else(|| finish_reason_from_sse(body))
998}
999
1000fn finish_reason_from_sse(body: &Value) -> Option<String> {
1001    let events = body.get("sse_events")?.as_array()?;
1002    events.iter().rev().find_map(|event| {
1003        let parsed = event.get("parsed_data")?;
1004        finish_reason_from_body(parsed)
1005    })
1006}
1007
1008fn string_at(value: &Value, path: &[&str]) -> Option<String> {
1009    let mut current = value;
1010    for key in path {
1011        if let Ok(index) = key.parse::<usize>() {
1012            current = current.get(index)?;
1013        } else {
1014            current = current.get(*key)?;
1015        }
1016    }
1017    current
1018        .as_str()
1019        .map(str::to_string)
1020        .filter(|s| !s.is_empty())
1021}
1022
1023#[cfg(test)]
1024mod tests {
1025    use super::*;
1026    use serde_json::json;
1027
1028    fn process_node_id(
1029        view: &mut MaterializedView,
1030        timestamp: u64,
1031        event: &str,
1032        exit_code: Option<i32>,
1033    ) -> String {
1034        let mut data = json!({"event": event, "filename": format!("cmd-{timestamp}")});
1035        if let Some(code) = exit_code {
1036            data["exit_code"] = json!(code);
1037        }
1038        let event = Event::new_with_timestamp(
1039            timestamp,
1040            "process".to_string(),
1041            42,
1042            "cmd".to_string(),
1043            data,
1044        );
1045        view.ingest_event(&event).expect("ingest process event");
1046        view.export_snapshot(crate::model::SnapshotOptions { audit_limit: 100 })
1047            .process_nodes
1048            .into_iter()
1049            .find(|row| {
1050                row.command.as_deref() == Some(&format!("cmd-{timestamp}"))
1051                    || row.end_timestamp_ms == Some(timestamp)
1052            })
1053            .map(|row| row.id)
1054            .expect("process node update")
1055    }
1056
1057    #[test]
1058    fn process_node_id_survives_pid_reuse() {
1059        let mut view = MaterializedView::new();
1060        let first_exec = process_node_id(&mut view, 1_000, "EXEC", None);
1061        let second_execve = process_node_id(&mut view, 1_500, "EXEC", None);
1062        let first_exit = process_node_id(&mut view, 2_000, "EXIT", Some(0));
1063        let second_exec = process_node_id(&mut view, 3_000, "EXEC", None);
1064        let second_exit = process_node_id(&mut view, 4_000, "EXIT", Some(1));
1065
1066        assert_eq!(first_exec, second_execve);
1067        assert_eq!(first_exec, first_exit);
1068        assert_eq!(second_exec, second_exit);
1069        assert_ne!(first_exec, second_exec);
1070    }
1071
1072    #[test]
1073    fn llm_request_audit_survives_response_pairing() {
1074        let mut view = MaterializedView::new();
1075        let req = Event::new_with_timestamp(
1076            1_000,
1077            "http_parser".to_string(),
1078            42,
1079            "agent".to_string(),
1080            json!({
1081                "tid": 7,
1082                "message_type": "request",
1083                "method": "POST",
1084                "path": "/v1/messages",
1085                "headers": { "host": "api.anthropic.com" },
1086                "body": "{\"model\":\"claude-sonnet\"}"
1087            }),
1088        );
1089        let resp = Event::new_with_timestamp(
1090            2_000,
1091            "http_parser".to_string(),
1092            42,
1093            "agent".to_string(),
1094            json!({
1095                "tid": 7,
1096                "message_type": "response",
1097                "status_code": 200,
1098                "headers": { "host": "api.anthropic.com" },
1099                "body": "{\"usage\":{\"input_tokens\":1,\"output_tokens\":2}}"
1100            }),
1101        );
1102
1103        view.ingest_event(&req).expect("ingest request");
1104        view.ingest_event(&resp).expect("ingest response");
1105
1106        let snapshot = view.export_snapshot(crate::model::SnapshotOptions { audit_limit: 100 });
1107        let llm_actions = snapshot
1108            .audit_events
1109            .iter()
1110            .filter(|row| row.audit_type == "llm")
1111            .filter_map(|row| row.action.as_deref())
1112            .collect::<Vec<_>>();
1113        assert!(llm_actions.contains(&"request"));
1114        assert!(llm_actions.contains(&"call"));
1115    }
1116
1117    #[test]
1118    fn llm_call_health_promotes_pending_to_complete() {
1119        let mut view = MaterializedView::new();
1120        let req = Event::new_with_timestamp(
1121            1_000,
1122            "http_parser".to_string(),
1123            42,
1124            "agent".to_string(),
1125            json!({
1126                "tid": 7,
1127                "message_type": "request",
1128                "method": "POST",
1129                "path": "/v1/chat/completions",
1130                "headers": { "host": "api.openai.com" },
1131                "body": "{\"model\":\"gpt-test\",\"metadata\":{\"session_id\":\"sess-1\"}}"
1132            }),
1133        );
1134
1135        view.ingest_event(&req).expect("ingest request");
1136        let pending = view.llm_call_rows(10);
1137        assert_eq!(pending[0].status, "pending");
1138        assert_eq!(pending[0].session_id.as_deref(), Some("sess-1"));
1139        assert_eq!(pending[0].call_kind.as_deref(), Some("chat"));
1140
1141        let resp = Event::new_with_timestamp(
1142            2_000,
1143            "http_parser".to_string(),
1144            42,
1145            "agent".to_string(),
1146            json!({
1147                "tid": 7,
1148                "message_type": "response",
1149                "status_code": 200,
1150                "headers": { "host": "api.openai.com" },
1151                "body": "{\"choices\":[{\"finish_reason\":\"stop\"}],\"usage\":{\"prompt_tokens\":1,\"completion_tokens\":2}}"
1152            }),
1153        );
1154
1155        view.ingest_event(&resp).expect("ingest response");
1156        let complete = view.llm_call_rows(10);
1157        assert_eq!(complete[0].status, "complete");
1158        assert_eq!(complete[0].finish_reason.as_deref(), Some("stop"));
1159        assert_eq!(complete[0].total_tokens, 3);
1160    }
1161
1162    #[test]
1163    fn llm_call_health_marks_http_errors() {
1164        let mut view = MaterializedView::new();
1165        let req = Event::new_with_timestamp(
1166            1_000,
1167            "http_parser".to_string(),
1168            42,
1169            "agent".to_string(),
1170            json!({
1171                "tid": 7,
1172                "message_type": "request",
1173                "method": "POST",
1174                "path": "/v1/chat/completions",
1175                "headers": { "host": "api.openai.com" },
1176                "body": "{\"model\":\"gpt-test\"}"
1177            }),
1178        );
1179        let resp = Event::new_with_timestamp(
1180            2_000,
1181            "http_parser".to_string(),
1182            42,
1183            "agent".to_string(),
1184            json!({
1185                "tid": 7,
1186                "message_type": "response",
1187                "status_code": 429,
1188                "headers": { "host": "api.openai.com" },
1189                "body": "{\"error\":{\"message\":\"rate limited\"}}"
1190            }),
1191        );
1192
1193        view.ingest_event(&req).expect("ingest request");
1194        view.ingest_event(&resp).expect("ingest response");
1195        let calls = view.llm_call_rows(10);
1196        assert_eq!(calls[0].status, "error");
1197        assert_eq!(calls[0].error_type.as_deref(), Some("http_429"));
1198    }
1199}