1mod canonical;
5pub(crate) mod llm;
6pub mod process_select;
7mod projection;
8pub mod session_process_match;
9
10pub(crate) use canonical::{CanonicalEvent, EventKind, normalize_event};
11pub(crate) use llm::{
12 body_json, extract_model, extract_token_usage, extract_token_usage_from_sse, provider_from_host,
13};
14
15use crate::model::{
16 AGENT_NATIVE_SOURCE, AuditEventRow, LlmCallRow, NetworkTargetRow, ProcessNodeRow,
17 ResourceSampleRow, SessionRow, Snapshot, SnapshotOptions, SnapshotSummary, TokenSummary,
18 TokenUsageRow, ToolCallRow, ViewResult, ViewSink,
19};
20use chrono::{SecondsFormat, Utc};
21use serde_json::Value;
22use std::collections::{BTreeMap, BTreeSet, HashMap, VecDeque};
23use std::sync::{Arc, Mutex};
24
25pub type SharedMaterializedView = Arc<Mutex<MaterializedView>>;
26
27const MAX_AUDIT_EVENTS_IN_MEMORY: usize = 20_000;
28const MAX_RESOURCE_SAMPLES_IN_MEMORY: usize = 10_000;
29
30#[derive(Default)]
31pub struct MaterializedView {
32 source: String,
33 llm_calls: BTreeMap<String, LlmCallRow>,
34 token_usage: BTreeMap<String, TokenUsageRow>,
35 audit_events: BTreeMap<String, AuditEventRow>,
36 process_nodes: BTreeMap<String, ProcessNodeRow>,
37 tool_calls: BTreeMap<String, ToolCallRow>,
38 sessions: BTreeMap<String, SessionRow>,
39 network_targets: BTreeMap<String, NetworkTargetRow>,
40 resource_samples: Vec<ResourceSampleRow>,
41 audit_order: VecDeque<String>,
42 sinks: Vec<Box<dyn ViewSink>>,
43 pending: HashMap<(u32, u64), VecDeque<PendingRequest>>,
44 active_processes: HashMap<u32, String>,
45 counts: ViewCounts,
46 start_timestamp_ms: Option<u64>,
47 end_timestamp_ms: Option<u64>,
48 max_audit_events: Option<usize>,
49 max_resource_samples: Option<usize>,
50 next_seq: u64,
51}
52
53#[derive(Default)]
54struct ViewCounts {
55 llm_calls: i64,
56 token_usage: i64,
57 audit_events: i64,
58 process_nodes: i64,
59 tool_calls: i64,
60 sessions: i64,
61 network_targets: i64,
62 resource_samples: i64,
63}
64
65#[derive(Debug, Clone)]
66struct PendingRequest {
67 event_id: String,
68 timestamp_ms: u64,
69 pid: u32,
70 comm: String,
71 provider: Option<String>,
72 model: Option<String>,
73 host: Option<String>,
74 path: Option<String>,
75 request_id: Option<String>,
76 body_json: Option<Value>,
77}
78
79impl MaterializedView {
80 pub fn new() -> Self {
81 Self::default()
82 }
83
84 pub fn bounded() -> Self {
85 let mut view = Self::new();
86 view.max_audit_events = Some(MAX_AUDIT_EVENTS_IN_MEMORY);
87 view.max_resource_samples = Some(MAX_RESOURCE_SAMPLES_IN_MEMORY);
88 view
89 }
90
91 pub fn shared_bounded() -> SharedMaterializedView {
92 Arc::new(Mutex::new(Self::bounded()))
93 }
94
95 pub fn add_sink(&mut self, sink: Box<dyn ViewSink>) {
96 self.sinks.push(sink);
97 }
98
99 pub fn set_source(&mut self, source: impl Into<String>) {
100 self.source = source.into();
101 }
102
103 pub fn emit_llm_call(&mut self, row: LlmCallRow) -> ViewResult<()> {
104 self.apply_llm_call(&row);
105 self.publish(|sink| sink.llm_call(&row))
106 }
107
108 pub fn emit_token_usage(&mut self, row: TokenUsageRow) -> ViewResult<()> {
109 self.apply_token_usage(&row);
110 self.publish(|sink| sink.token_usage(&row))
111 }
112
113 pub fn emit_audit_event(&mut self, row: AuditEventRow) -> ViewResult<()> {
114 self.apply_audit_event(&row);
115 self.publish(|sink| sink.audit_event(&row))
116 }
117
118 pub fn emit_process_node(&mut self, row: ProcessNodeRow) -> ViewResult<()> {
119 self.upsert_process_node(&row);
120 self.publish(|sink| sink.process_node(&row))
121 }
122
123 pub fn emit_tool_call(&mut self, row: ToolCallRow) -> ViewResult<()> {
124 self.apply_tool_call(&row);
125 self.publish(|sink| sink.tool_call(&row))
126 }
127
128 pub fn emit_network_target(&mut self, row: NetworkTargetRow) -> ViewResult<()> {
129 self.upsert_network_target(&row);
130 self.publish(|sink| sink.network_target(&row))
131 }
132
133 pub fn emit_resource_sample(&mut self, row: ResourceSampleRow) -> ViewResult<()> {
134 self.apply_resource_sample(&row);
135 self.publish(|sink| sink.resource_sample(&row))
136 }
137
138 fn publish<F>(&mut self, mut publish: F) -> ViewResult<()>
139 where
140 F: FnMut(&mut dyn ViewSink) -> ViewResult<()>,
141 {
142 let mut first_error = None;
143 for sink in &mut self.sinks {
144 if let Err(error) = publish(sink.as_mut()) {
145 log::warn!("MaterializedView: failed to publish view row: {}", error);
146 first_error.get_or_insert_with(|| error.to_string());
147 }
148 }
149 if let Some(error) = first_error {
150 return Err(std::io::Error::other(error).into());
151 }
152 Ok(())
153 }
154}
155
156impl MaterializedView {
157 pub fn apply_llm_call(&mut self, row: &LlmCallRow) {
158 if !self.llm_calls.contains_key(&row.id) {
159 self.counts.llm_calls += 1;
160 }
161 self.observe(Some(row.start_timestamp_ms));
162 self.observe(row.end_timestamp_ms);
163 self.llm_calls.insert(row.id.clone(), row.clone());
164 }
165
166 pub fn apply_token_usage(&mut self, row: &TokenUsageRow) {
167 if !self.token_usage.contains_key(&row.id) {
168 self.counts.token_usage += 1;
169 }
170 self.observe(Some(row.timestamp_ms));
171 self.token_usage.insert(row.id.clone(), row.clone());
172 }
173
174 pub fn apply_audit_event(&mut self, row: &AuditEventRow) {
175 if !self.audit_events.contains_key(&row.id) {
176 self.counts.audit_events += 1;
177 if self.max_audit_events.is_some() {
178 self.audit_order.push_back(row.id.clone());
179 }
180 }
181 self.observe(Some(row.timestamp_ms));
182 self.audit_events.insert(row.id.clone(), row.clone());
183 if let Some(max) = self.max_audit_events {
184 while self.audit_events.len() > max {
185 let Some(id) = self.audit_order.pop_front() else {
186 break;
187 };
188 self.audit_events.remove(&id);
189 }
190 }
191 }
192
193 pub fn apply_tool_call(&mut self, row: &ToolCallRow) {
194 if !self.tool_calls.contains_key(&row.id) {
195 self.counts.tool_calls += 1;
196 }
197 self.observe(Some(row.timestamp_ms));
198 self.tool_calls.insert(row.id.clone(), row.clone());
199 }
200
201 pub fn apply_resource_sample(&mut self, row: &ResourceSampleRow) {
202 self.counts.resource_samples += 1;
203 self.observe(Some(row.timestamp_ms));
204 self.resource_samples.push(row.clone());
205 if let Some(max) = self.max_resource_samples {
206 let overflow = self.resource_samples.len().saturating_sub(max);
207 if overflow > 0 {
208 self.resource_samples.drain(0..overflow);
209 }
210 }
211 }
212
213 pub fn upsert_session(&mut self, row: &SessionRow) {
214 self.observe(Some(row.start_timestamp_ms));
215 self.observe(row.end_timestamp_ms);
216 let Some(existing) = self.sessions.get_mut(&row.id) else {
217 self.counts.sessions += 1;
218 self.sessions.insert(row.id.clone(), row.clone());
219 return;
220 };
221
222 existing.start_timestamp_ms = existing.start_timestamp_ms.min(row.start_timestamp_ms);
223 existing.end_timestamp_ms = max_optional(existing.end_timestamp_ms, row.end_timestamp_ms);
224 if row.model.as_deref().is_some_and(|model| model != "unknown") || existing.model.is_none()
225 {
226 existing.model = row.model.clone();
227 }
228 existing.input_tokens = existing.input_tokens.max(row.input_tokens);
229 existing.output_tokens = existing.output_tokens.max(row.output_tokens);
230 existing.total_tokens = existing.total_tokens.max(row.total_tokens);
231 existing.confidence = max_optional(existing.confidence, row.confidence);
232 }
233
234 pub fn upsert_network_target(&mut self, row: &NetworkTargetRow) {
235 self.observe(row.first_timestamp_ms);
236 self.observe(row.last_timestamp_ms);
237 let key = network_target_key(row);
238 let Some(existing) = self.network_targets.get_mut(&key) else {
239 self.counts.network_targets += 1;
240 self.network_targets.insert(key, row.clone());
241 return;
242 };
243
244 existing.count += row.count;
245 existing.error_count += row.error_count;
246 existing.first_timestamp_ms =
247 min_optional(existing.first_timestamp_ms, row.first_timestamp_ms);
248 existing.last_timestamp_ms =
249 max_optional(existing.last_timestamp_ms, row.last_timestamp_ms);
250 }
251
252 pub fn upsert_process_node(&mut self, row: &ProcessNodeRow) {
253 self.observe(row.start_timestamp_ms);
254 self.observe(row.end_timestamp_ms);
255 let Some(existing) = self.process_nodes.get_mut(&row.id) else {
256 self.counts.process_nodes += 1;
257 self.process_nodes.insert(row.id.clone(), row.clone());
258 return;
259 };
260
261 existing.start_timestamp_ms =
262 min_optional(existing.start_timestamp_ms, row.start_timestamp_ms);
263 existing.end_timestamp_ms = max_optional(existing.end_timestamp_ms, row.end_timestamp_ms);
264 if row.ppid.is_some() {
265 existing.ppid = row.ppid;
266 }
267 if row.root_pid.is_some() {
268 existing.root_pid = row.root_pid;
269 }
270 if row.comm.is_some() {
271 existing.comm = row.comm.clone();
272 }
273 if row.command.is_some() {
274 existing.command = row.command.clone();
275 }
276 if existing.argv.is_empty() && !row.argv.is_empty() {
277 existing.argv = row.argv.clone();
278 }
279 if row.cwd.is_some() {
280 existing.cwd = row.cwd.clone();
281 }
282 if row.exit_code.is_some() {
283 existing.exit_code = row.exit_code;
284 }
285 if row.status.is_some() {
286 existing.status = row.status.clone();
287 }
288 existing.confidence = max_optional(existing.confidence, row.confidence);
289 }
290
291 pub fn export_snapshot(&self, options: SnapshotOptions) -> Snapshot {
292 Snapshot {
293 schema_version: 1,
294 generated_at: Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true),
295 summary: self.snapshot_summary(options),
296 token_summary: self.token_summary("model"),
297 network_targets: self.network_targets(),
298 process_nodes: self.process_nodes(),
299 audit_events: self.audit_events(options.audit_limit),
300 resource_samples: self.resource_sample_rows(),
301 sessions: self.sessions(),
302 tool_calls: self.tool_calls.values().cloned().collect(),
303 }
304 }
305
306 fn snapshot_summary(&self, options: SnapshotOptions) -> SnapshotSummary {
307 let (input_tokens, output_tokens, total_tokens) =
308 self.effective_tokens()
309 .into_iter()
310 .fold((0, 0, 0), |acc, token| {
311 (
312 acc.0 + token.input_tokens,
313 acc.1 + token.output_tokens,
314 acc.2 + token.total_tokens,
315 )
316 });
317
318 SnapshotSummary {
319 source: if self.source.is_empty() {
320 "materialized_view".to_string()
321 } else {
322 self.source.clone()
323 },
324 view_events: self.view_events(),
325 llm_calls: self.counts.llm_calls,
326 token_usage_rows: self.counts.token_usage,
327 audit_events: self.counts.audit_events,
328 sessions: self.counts.sessions,
329 input_tokens,
330 output_tokens,
331 total_tokens,
332 start_timestamp_ms: self.start_timestamp_ms,
333 end_timestamp_ms: self.end_timestamp_ms,
334 audit_limit: options.audit_limit,
335 }
336 }
337
338 pub fn token_summary(&self, group_by: &str) -> Vec<TokenSummary> {
339 let mut groups: BTreeMap<String, TokenSummaryGroup> = BTreeMap::new();
340 for token in self.effective_tokens() {
341 let group = self.token_group(token, group_by);
342 let entry = groups
343 .entry(group.clone())
344 .or_insert_with(|| TokenSummaryGroup::new(group));
345 entry.row.input_tokens += token.input_tokens;
346 entry.row.output_tokens += token.output_tokens;
347 entry.row.cache_creation_tokens += token.cache_creation_tokens;
348 entry.row.cache_read_tokens += token.cache_read_tokens;
349 entry.row.total_tokens += token.total_tokens;
350 entry.row.calls += 1;
351 if let Some(session_key) = self.token_session_key(token) {
352 entry.sessions.insert(session_key);
353 }
354 }
355 let mut rows = groups
356 .into_values()
357 .map(|mut group| {
358 group.row.sessions = group.sessions.len() as i64;
359 group.row
360 })
361 .collect::<Vec<_>>();
362 sort_token_summary(&mut rows);
363 rows
364 }
365
366 pub fn audit_rows(&self, audit_type: Option<&str>, limit: usize) -> Vec<AuditEventRow> {
367 let mut rows = self
368 .audit_events
369 .values()
370 .filter(|row| audit_type.is_none_or(|audit_type| row.audit_type == audit_type))
371 .cloned()
372 .collect::<Vec<_>>();
373 rows.sort_by_key(|b| std::cmp::Reverse(b.timestamp_ms));
374 rows.truncate(limit.clamp(1, 10_000));
375 rows
376 }
377
378 pub fn llm_call_rows(&self, limit: usize) -> Vec<LlmCallRow> {
379 let token_totals = self.effective_token_totals_by_call();
380 let mut rows = self
381 .llm_calls
382 .values()
383 .cloned()
384 .map(|mut row| {
385 if let Some((input, output, total)) = token_totals.get(&row.id) {
386 row.input_tokens = *input;
387 row.output_tokens = *output;
388 row.total_tokens = *total;
389 }
390 row
391 })
392 .collect::<Vec<_>>();
393 rows.sort_by_key(|b| std::cmp::Reverse(b.start_timestamp_ms));
394 rows.truncate(limit.clamp(1, 10_000));
395 rows
396 }
397
398 fn resource_sample_rows(&self) -> Vec<ResourceSampleRow> {
399 let mut rows = self.resource_samples.clone();
400 rows.sort_by(|a, b| {
401 a.timestamp_ms
402 .cmp(&b.timestamp_ms)
403 .then_with(|| a.pid.cmp(&b.pid))
404 .then_with(|| a.comm.cmp(&b.comm))
405 });
406 rows
407 }
408
409 fn network_targets(&self) -> Vec<NetworkTargetRow> {
410 let mut rows = self.network_targets.values().cloned().collect::<Vec<_>>();
411 rows.sort_by(|a, b| {
412 b.count
413 .cmp(&a.count)
414 .then_with(|| a.host.cmp(&b.host))
415 .then_with(|| a.path.cmp(&b.path))
416 });
417 rows
418 }
419
420 fn audit_events(&self, limit: usize) -> Vec<AuditEventRow> {
421 let mut rows = self.audit_events.values().cloned().collect::<Vec<_>>();
422 rows.sort_by(|a, b| {
423 a.timestamp_ms
424 .cmp(&b.timestamp_ms)
425 .then_with(|| a.id.cmp(&b.id))
426 });
427 let limit = limit.min(100_000);
428 if rows.len() > limit {
429 rows.drain(0..rows.len() - limit);
430 }
431 rows
432 }
433
434 fn process_nodes(&self) -> Vec<ProcessNodeRow> {
435 let mut rows = self.process_nodes.values().cloned().collect::<Vec<_>>();
436 rows.sort_by(|a, b| {
437 a.start_timestamp_ms
438 .cmp(&b.start_timestamp_ms)
439 .then_with(|| a.pid.cmp(&b.pid))
440 .then_with(|| a.id.cmp(&b.id))
441 });
442 rows
443 }
444
445 fn sessions(&self) -> Vec<SessionRow> {
446 let mut rows = self.sessions.values().cloned().collect::<Vec<_>>();
447 rows.sort_by(|a, b| {
448 a.start_timestamp_ms
449 .cmp(&b.start_timestamp_ms)
450 .then_with(|| a.id.cmp(&b.id))
451 });
452 rows
453 }
454
455 fn view_events(&self) -> i64 {
456 self.counts.llm_calls
457 + self.counts.token_usage
458 + self.counts.audit_events
459 + self.counts.process_nodes
460 + self.counts.tool_calls
461 + self.counts.sessions
462 + self.counts.network_targets
463 + self.counts.resource_samples
464 }
465
466 fn effective_tokens(&self) -> Vec<&TokenUsageRow> {
467 let mut selected: BTreeMap<String, &TokenUsageRow> = BTreeMap::new();
468 let mut gemini_totals = BTreeMap::new();
469 for token in self.token_usage.values() {
470 let Some(key) = token.pid.zip(token.model.as_deref()) else {
471 continue;
472 };
473 let totals = gemini_totals.entry(key).or_insert((0, 0));
474 match token.source.as_str() {
475 "response_usage" | "orphan_response_usage" => totals.0 += token.total_tokens,
476 "gemini_cli_stdout_stats" => totals.1 = totals.1.max(token.total_tokens),
477 _ => {}
478 }
479 }
480 for token in self.token_usage.values() {
481 if let Some((network, stdout)) = token
482 .pid
483 .zip(token.model.as_deref())
484 .and_then(|key| gemini_totals.get(&key))
485 {
486 let network_source = matches!(
487 token.source.as_str(),
488 "response_usage" | "orphan_response_usage"
489 );
490 if (token.source == "gemini_cli_stdout_stats" && network >= stdout)
491 || (network_source && stdout > network)
492 || (token.source == "gemini_cli_stdout_stats" && token.total_tokens < *stdout)
493 {
494 continue;
495 }
496 }
497 let key = if token.source == "gemini_cli_stdout_stats" {
498 token
499 .pid
500 .zip(token.model.as_deref())
501 .map(|(pid, model)| format!("gemini-stdout\0{pid}\0{model}"))
502 .unwrap_or_else(|| token.id.clone())
503 } else if token.llm_call_id.is_empty() {
504 token.id.clone()
505 } else {
506 token.llm_call_id.clone()
507 };
508 match selected.get(&key) {
509 Some(current) if !token_has_higher_priority(token, current) => {}
510 _ => {
511 selected.insert(key, token);
512 }
513 }
514 }
515 selected.into_values().collect()
516 }
517
518 fn effective_token_totals_by_call(&self) -> BTreeMap<String, (i64, i64, i64)> {
519 let mut totals = BTreeMap::new();
520 for token in self.effective_tokens() {
521 totals.insert(
522 token.llm_call_id.clone(),
523 (token.input_tokens, token.output_tokens, token.total_tokens),
524 );
525 }
526 totals
527 }
528
529 fn observe(&mut self, timestamp: Option<u64>) {
530 observe_timestamp(
531 &mut self.start_timestamp_ms,
532 &mut self.end_timestamp_ms,
533 timestamp,
534 );
535 }
536
537 fn token_group(&self, token: &TokenUsageRow, group_by: &str) -> String {
538 match group_by {
539 "provider" => token.provider.clone(),
540 "comm" => token.comm.clone(),
541 "pid" => token.pid.map(|pid| pid.to_string()),
542 "dir" | "cwd" | "directory" => self.token_working_dir(token),
543 _ => token.model.clone(),
544 }
545 .filter(|value| !value.is_empty())
546 .unwrap_or_else(|| "unknown".to_string())
547 }
548
549 fn token_working_dir(&self, token: &TokenUsageRow) -> Option<String> {
550 self.token_session_key(token)
551 .and_then(|session_id| self.sessions.get(&session_id).and_then(session_cwd))
552 .or_else(|| self.token_process_cwd(token))
553 }
554
555 fn token_session_key(&self, token: &TokenUsageRow) -> Option<String> {
556 if let Some(session_id) = self
557 .llm_calls
558 .get(&token.llm_call_id)
559 .and_then(|row| row.session_id.as_ref())
560 .filter(|session_id| !session_id.is_empty())
561 {
562 return Some(session_id.clone());
563 }
564
565 self.sessions
566 .keys()
567 .find(|session_id| {
568 let session_id = session_id.as_str();
569 token.llm_call_id == session_id
570 || token
571 .llm_call_id
572 .strip_prefix(session_id)
573 .is_some_and(|suffix| suffix.starts_with('-'))
574 })
575 .cloned()
576 }
577
578 fn token_process_cwd(&self, token: &TokenUsageRow) -> Option<String> {
579 let pid = token.pid.or_else(|| {
580 self.llm_calls
581 .get(&token.llm_call_id)
582 .and_then(|row| row.pid)
583 })?;
584 self.process_nodes
585 .values()
586 .find(|row| row.pid == pid && row.cwd.as_deref().is_some_and(|cwd| !cwd.is_empty()))
587 .and_then(|row| row.cwd.clone())
588 }
589}
590
591struct TokenSummaryGroup {
592 row: TokenSummary,
593 sessions: BTreeSet<String>,
594}
595
596impl TokenSummaryGroup {
597 fn new(group: String) -> Self {
598 Self {
599 row: TokenSummary {
600 group,
601 input_tokens: 0,
602 output_tokens: 0,
603 cache_creation_tokens: 0,
604 cache_read_tokens: 0,
605 total_tokens: 0,
606 calls: 0,
607 sessions: 0,
608 },
609 sessions: BTreeSet::new(),
610 }
611 }
612}
613
614fn session_cwd(session: &SessionRow) -> Option<String> {
615 session
616 .attributes
617 .get("cwd")
618 .and_then(Value::as_str)
619 .filter(|cwd| !cwd.is_empty())
620 .map(str::to_string)
621}
622
623fn token_has_higher_priority(candidate: &TokenUsageRow, current: &TokenUsageRow) -> bool {
624 let candidate_priority = token_source_priority(&candidate.source);
625 let current_priority = token_source_priority(¤t.source);
626 candidate_priority
627 .cmp(¤t_priority)
628 .then_with(|| {
629 current
630 .confidence
631 .unwrap_or_default()
632 .partial_cmp(&candidate.confidence.unwrap_or_default())
633 .unwrap_or(std::cmp::Ordering::Equal)
634 })
635 .then_with(|| candidate.id.cmp(¤t.id))
636 .is_lt()
637}
638
639fn token_source_priority(source: &str) -> u8 {
640 match source {
641 "response_usage" => 0,
644 "orphan_response_usage" => 1,
645 "gemini_cli_stdout_stats" => 2,
646 "claude_telemetry" => 3,
647 AGENT_NATIVE_SOURCE => 4,
648 _ => 5,
649 }
650}
651
652fn network_target_key(row: &NetworkTargetRow) -> String {
653 format!(
654 "{}\0{}\0{}",
655 row.pid.unwrap_or_default(),
656 row.host,
657 row.path.as_deref().unwrap_or_default()
658 )
659}
660
661fn observe_timestamp(start: &mut Option<u64>, end: &mut Option<u64>, timestamp: Option<u64>) {
662 let Some(timestamp) = timestamp else {
663 return;
664 };
665 *start = Some(start.map_or(timestamp, |current| current.min(timestamp)));
666 *end = Some(end.map_or(timestamp, |current| current.max(timestamp)));
667}
668
669fn sort_token_summary(rows: &mut [TokenSummary]) {
670 rows.sort_by(|a, b| {
671 b.total_tokens
672 .cmp(&a.total_tokens)
673 .then_with(|| a.group.cmp(&b.group))
674 });
675}
676
677fn min_optional<T: PartialOrd>(left: Option<T>, right: Option<T>) -> Option<T> {
678 match (left, right) {
679 (Some(left), Some(right)) => Some(if left <= right { left } else { right }),
680 (Some(value), None) | (None, Some(value)) => Some(value),
681 (None, None) => None,
682 }
683}
684
685fn max_optional<T: PartialOrd>(left: Option<T>, right: Option<T>) -> Option<T> {
686 match (left, right) {
687 (Some(left), Some(right)) => Some(if left >= right { left } else { right }),
688 (Some(value), None) | (None, Some(value)) => Some(value),
689 (None, None) => None,
690 }
691}
692
693#[cfg(test)]
694mod tests {
695 use super::*;
696 use serde_json::json;
697
698 fn audit_row(timestamp_ms: u64) -> AuditEventRow {
699 AuditEventRow {
700 id: format!("audit-{timestamp_ms}"),
701 timestamp_ms,
702 audit_type: "file".to_string(),
703 pid: Some(1),
704 comm: Some("test".to_string()),
705 subject: None,
706 action: Some("write".to_string()),
707 target: Some(format!("/tmp/{timestamp_ms}")),
708 status: Some("observed".to_string()),
709 summary: None,
710 details: json!({}),
711 }
712 }
713
714 fn session_row(id: &str, cwd: &str) -> SessionRow {
715 SessionRow {
716 id: id.to_string(),
717 agent_type: "codex".to_string(),
718 start_timestamp_ms: 1_000,
719 end_timestamp_ms: Some(2_000),
720 status: "observed".to_string(),
721 model: Some("gpt-5".to_string()),
722 input_tokens: 0,
723 output_tokens: 0,
724 total_tokens: 0,
725 view_source: AGENT_NATIVE_SOURCE.to_string(),
726 confidence: Some(0.95),
727 attributes: json!({ "cwd": cwd }),
728 }
729 }
730
731 fn token_row(
732 id: &str,
733 llm_call_id: &str,
734 model: &str,
735 input_tokens: i64,
736 output_tokens: i64,
737 cache_creation_tokens: i64,
738 cache_read_tokens: i64,
739 total_tokens: i64,
740 ) -> TokenUsageRow {
741 TokenUsageRow {
742 id: id.to_string(),
743 llm_call_id: llm_call_id.to_string(),
744 timestamp_ms: 1_500,
745 pid: None,
746 comm: Some("codex".to_string()),
747 provider: None,
748 model: Some(model.to_string()),
749 input_tokens,
750 output_tokens,
751 cache_creation_tokens,
752 cache_read_tokens,
753 total_tokens,
754 source: AGENT_NATIVE_SOURCE.to_string(),
755 view_source: AGENT_NATIVE_SOURCE.to_string(),
756 confidence: Some(0.95),
757 }
758 }
759
760 fn process_node(pid: u32, cwd: &str) -> ProcessNodeRow {
761 ProcessNodeRow {
762 id: format!("process-{pid}"),
763 pid,
764 ppid: None,
765 root_pid: Some(pid),
766 start_timestamp_ms: Some(1_000),
767 end_timestamp_ms: None,
768 comm: Some("agent".to_string()),
769 command: Some("agent".to_string()),
770 argv: Vec::new(),
771 cwd: Some(cwd.to_string()),
772 exit_code: None,
773 status: Some("observed".to_string()),
774 view_source: "process".to_string(),
775 confidence: Some(0.8),
776 }
777 }
778
779 fn llm_call_row(id: &str, pid: u32, session_id: Option<&str>) -> LlmCallRow {
780 LlmCallRow {
781 id: id.to_string(),
782 session_id: session_id.map(str::to_string),
783 conversation_id: None,
784 start_timestamp_ms: 1_100,
785 end_timestamp_ms: Some(1_400),
786 pid: Some(pid),
787 comm: Some("agent".to_string()),
788 provider: Some("anthropic".to_string()),
789 model: Some("claude-sonnet-4".to_string()),
790 call_kind: Some("messages".to_string()),
791 status: "ok".to_string(),
792 error_type: None,
793 finish_reason: None,
794 host: Some("api.anthropic.com".to_string()),
795 path: Some("/v1/messages".to_string()),
796 status_code: Some(200),
797 input_tokens: 0,
798 output_tokens: 0,
799 total_tokens: 0,
800 request: json!({}),
801 response: json!({}),
802 }
803 }
804
805 #[test]
806 fn audit_retention_keeps_counters_and_recent_rows() {
807 let mut view = MaterializedView::bounded();
808 for timestamp_ms in 0..(MAX_AUDIT_EVENTS_IN_MEMORY as u64 + 5) {
809 view.apply_audit_event(&audit_row(timestamp_ms));
810 }
811
812 let snapshot = view.export_snapshot(SnapshotOptions {
813 audit_limit: MAX_AUDIT_EVENTS_IN_MEMORY + 10,
814 });
815 assert_eq!(
816 snapshot.summary.audit_events,
817 MAX_AUDIT_EVENTS_IN_MEMORY as i64 + 5
818 );
819 assert_eq!(snapshot.audit_events.len(), MAX_AUDIT_EVENTS_IN_MEMORY);
820 assert_eq!(snapshot.audit_events[0].timestamp_ms, 5);
821 assert_eq!(snapshot.summary.start_timestamp_ms, Some(0));
822 assert_eq!(
823 snapshot.summary.end_timestamp_ms,
824 Some(MAX_AUDIT_EVENTS_IN_MEMORY as u64 + 4)
825 );
826 }
827
828 #[test]
829 fn resource_retention_keeps_counters() {
830 let mut view = MaterializedView::bounded();
831 for timestamp_ms in 0..(MAX_RESOURCE_SAMPLES_IN_MEMORY as u64 + 5) {
832 view.apply_resource_sample(&ResourceSampleRow {
833 timestamp_ms,
834 pid: Some(1),
835 comm: Some("test".to_string()),
836 cpu_percent: Some(1.0),
837 rss_mb: Some(2),
838 });
839 }
840
841 let snapshot = view.export_snapshot(SnapshotOptions { audit_limit: 0 });
842 assert_eq!(
843 snapshot.summary.view_events,
844 MAX_RESOURCE_SAMPLES_IN_MEMORY as i64 + 5
845 );
846 assert_eq!(
847 snapshot.resource_samples.len(),
848 MAX_RESOURCE_SAMPLES_IN_MEMORY
849 );
850 assert_eq!(snapshot.resource_samples[0].timestamp_ms, 5);
851 }
852
853 #[test]
854 fn token_summary_groups_agent_native_tokens_by_session_dir() {
855 let mut view = MaterializedView::new();
856 view.upsert_session(&session_row("local:codex:session-1", "/repo/one"));
857 view.apply_token_usage(&token_row(
858 "token-1",
859 "local:codex:session-1-gpt-5",
860 "gpt-5",
861 10,
862 5,
863 2,
864 3,
865 20,
866 ));
867 view.apply_token_usage(&token_row(
868 "token-2",
869 "local:codex:session-1-gpt-4",
870 "gpt-4",
871 7,
872 2,
873 0,
874 1,
875 10,
876 ));
877
878 let rows = view.token_summary("dir");
879
880 assert_eq!(rows.len(), 1);
881 assert_eq!(rows[0].group, "/repo/one");
882 assert_eq!(rows[0].input_tokens, 17);
883 assert_eq!(rows[0].output_tokens, 7);
884 assert_eq!(rows[0].cache_creation_tokens, 2);
885 assert_eq!(rows[0].cache_read_tokens, 4);
886 assert_eq!(rows[0].total_tokens, 30);
887 assert_eq!(rows[0].calls, 2);
888 assert_eq!(rows[0].sessions, 1);
889 }
890
891 #[test]
892 fn token_summary_groups_saved_tokens_by_process_dir() {
893 let mut view = MaterializedView::new();
894 view.upsert_process_node(&process_node(42, "/repo/saved"));
895 view.apply_llm_call(&llm_call_row("llm-1", 42, Some("session-1")));
896 view.apply_token_usage(&TokenUsageRow {
897 source: "response_usage".to_string(),
898 view_source: "sqlite".to_string(),
899 ..token_row("token-1", "llm-1", "claude-sonnet-4", 11, 13, 0, 0, 24)
900 });
901
902 let rows = view.token_summary("dir");
903
904 assert_eq!(rows.len(), 1);
905 assert_eq!(rows[0].group, "/repo/saved");
906 assert_eq!(rows[0].total_tokens, 24);
907 assert_eq!(rows[0].sessions, 1);
908 }
909
910 #[test]
911 fn gemini_stdout_tokens_are_fallback_for_network_usage() {
912 let mut view = MaterializedView::new();
913 view.apply_token_usage(&TokenUsageRow {
914 pid: Some(42),
915 comm: Some("node".to_string()),
916 source: "response_usage".to_string(),
917 ..token_row("token-network", "llm-network", "gemini", 11, 4, 0, 0, 15)
918 });
919 view.apply_token_usage(&TokenUsageRow {
920 pid: Some(42),
921 comm: Some("node".to_string()),
922 source: "gemini_cli_stdout_stats".to_string(),
923 ..token_row("token-stdout", "llm-stdout", "gemini", 11, 4, 0, 0, 15)
924 });
925
926 let snapshot = view.export_snapshot(SnapshotOptions { audit_limit: 0 });
927 assert_eq!(snapshot.summary.total_tokens, 15);
928 }
929
930 #[test]
931 fn gemini_stdout_tokens_cover_partial_network_capture() {
932 let mut view = MaterializedView::new();
933 for (id, total) in [("one", 7), ("two", 8)] {
934 view.apply_token_usage(&TokenUsageRow {
935 pid: Some(42),
936 source: "response_usage".to_string(),
937 ..token_row(id, id, "gemini", total, 0, 0, 0, total)
938 });
939 }
940 view.apply_token_usage(&TokenUsageRow {
941 pid: Some(42),
942 source: "gemini_cli_stdout_stats".to_string(),
943 ..token_row("old-stdout", "old-stdout", "gemini", 15, 0, 0, 0, 15)
944 });
945 view.apply_token_usage(&TokenUsageRow {
946 pid: Some(42),
947 source: "gemini_cli_stdout_stats".to_string(),
948 ..token_row("stdout", "stdout", "gemini", 30, 0, 0, 0, 30)
949 });
950
951 let snapshot = view.export_snapshot(SnapshotOptions { audit_limit: 0 });
952 assert_eq!(snapshot.summary.total_tokens, 30);
953 }
954
955 #[test]
956 fn token_summary_groups_missing_dir_as_unknown() {
957 let mut view = MaterializedView::new();
958 view.apply_token_usage(&token_row("token-1", "llm-1", "gpt-5", 1, 2, 3, 4, 10));
959
960 let rows = view.token_summary("dir");
961
962 assert_eq!(rows.len(), 1);
963 assert_eq!(rows[0].group, "unknown");
964 assert_eq!(rows[0].total_tokens, 10);
965 assert_eq!(rows[0].sessions, 0);
966 }
967
968 #[test]
969 fn process_node_preserves_first_non_empty_argv() {
970 let mut view = MaterializedView::new();
971 let first = ProcessNodeRow {
972 id: "pid:42:start:100".to_string(),
973 pid: 42,
974 ppid: Some(1),
975 root_pid: Some(42),
976 start_timestamp_ms: Some(100),
977 end_timestamp_ms: None,
978 comm: Some("agent".to_string()),
979 command: Some("agent".to_string()),
980 argv: vec![
981 "agent".to_string(),
982 "--model".to_string(),
983 "gpt-test".to_string(),
984 ],
985 cwd: Some("/tmp".to_string()),
986 exit_code: None,
987 status: Some("running".to_string()),
988 view_source: "process".to_string(),
989 confidence: Some(1.0),
990 };
991 let mut later = first.clone();
992 later.end_timestamp_ms = Some(200);
993 later.command = Some("agent-exit".to_string());
994 later.argv = vec!["agent-exit".to_string()];
995 later.exit_code = Some(0);
996 later.status = Some("success".to_string());
997
998 view.upsert_process_node(&first);
999 view.upsert_process_node(&later);
1000
1001 let snapshot = view.export_snapshot(SnapshotOptions { audit_limit: 0 });
1002 let row = snapshot.process_nodes.first().expect("process node");
1003 assert_eq!(row.argv, first.argv);
1004 assert_eq!(row.end_timestamp_ms, Some(200));
1005 assert_eq!(row.exit_code, Some(0));
1006 }
1007}