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