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