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, 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    /// Copy the materialized rows without copying output sinks. Callers can
85    /// merge read-only supplemental sources for one response without mutating
86    /// the live capture view or publishing duplicate rows.
87    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(&current.source);
653    candidate_priority
654        .cmp(&current_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(&current.id))
663        .is_lt()
664}
665
666fn token_source_priority(source: &str) -> u8 {
667    match source {
668        // Network-observed response usage is the primary fact source. Native
669        // session logs enrich or backfill when no network call was captured.
670        "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}