Skip to main content

agentsight_capture/view/
mod.rs

1// SPDX-License-Identifier: MIT
2// Copyright (c) 2026 eunomia-bpf org.
3
4mod 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(&current.source);
626    candidate_priority
627        .cmp(&current_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(&current.id))
636        .is_lt()
637}
638
639fn token_source_priority(source: &str) -> u8 {
640    match source {
641        // Network-observed response usage is the primary fact source. Native
642        // session logs enrich or backfill when no network call was captured.
643        "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}