Skip to main content

cherenkov/control/
activity.rs

1//! Paged views of the worker's copied expert counters; no GPU access.
2
3use super::stats::{Expert, Layer, Page, Rates, Summary};
4use crate::qwen4_exp::gpu::{ExpertActivity, ExpertCounters, LayerStats};
5use anyhow::{Result, ensure};
6use serde::Serialize;
7use std::ops::Range;
8
9fn range(total: usize, offset: usize, limit: usize) -> Result<Range<usize>> {
10    ensure!(
11        (1..=128).contains(&limit),
12        "limit must be between 1 and 128"
13    );
14    ensure!(offset <= total, "offset exceeds available entries");
15
16    Ok(offset..offset.saturating_add(limit).min(total))
17}
18
19fn page<T: Serialize>(total: usize, range: Range<usize>, data: Vec<T>) -> Result<Page<T>> {
20    // Leave room for framing and observation metadata under the 64 KiB limit.
21    let mut bytes = 0;
22    let mut bounded = Vec::new();
23
24    for entry in data {
25        bytes += serde_json::to_vec(&entry)?.len() + 1;
26
27        if bytes > 60 * crate::units::BYTES_PER_KIB {
28            break;
29        }
30
31        bounded.push(entry);
32    }
33
34    ensure!(
35        range.is_empty() || !bounded.is_empty(),
36        "one statistics entry exceeds the control frame"
37    );
38
39    let end = range.start + bounded.len();
40
41    Ok(Page {
42        offset: range.start,
43        next_offset: (end < total).then_some(end),
44        total,
45        data: bounded,
46    })
47}
48
49pub(super) fn layers(
50    activity: &ExpertActivity,
51    offset: usize,
52    limit: usize,
53) -> Result<Page<Layer>> {
54    let total = activity.layers.len();
55    let range = range(total, offset, limit)?;
56    let data = range
57        .clone()
58        .map(|layer| {
59            let start = layer * activity.experts_per_layer;
60            let mut counters = ExpertCounters::default();
61
62            for &record in &activity.records[start..start + activity.experts_per_layer] {
63                counters.add(record);
64            }
65
66            Layer {
67                layer,
68                name: activity.layer_prefixes[layer].clone(),
69                experts: activity.experts_per_layer,
70                counters,
71                streaming: activity.layers[layer].clone(),
72            }
73        })
74        .collect();
75
76    page(total, range, data)
77}
78
79pub(super) fn summary(activity: &ExpertActivity) -> Summary {
80    let mut total = LayerStats::default();
81
82    for layer in &activity.layers {
83        total.add(layer);
84    }
85
86    let reads = total.reads();
87    let phase = total.phases;
88    let gaps = phase.router_to_resident_seconds + phase.resident_to_fetched_seconds;
89    let stages = phase.resident_seconds + phase.fetched_stage_seconds;
90
91    Summary {
92        prefill: activity.prefill,
93        streaming: total,
94        reads,
95        rates: Rates {
96            bytes_per_second: ratio(reads.completed_bytes as f64, activity.elapsed_seconds),
97            reads_per_second: ratio(reads.completed_reads as f64, activity.elapsed_seconds),
98            gpu_handoff_gap_fraction: ratio(gaps, gaps + stages),
99            gpu_stage_fraction: ratio(stages, gaps + stages),
100        },
101    }
102}
103
104fn ratio(numerator: f64, denominator: f64) -> Option<f64> {
105    (denominator > 0.0).then(|| numerator / denominator)
106}
107
108pub(super) fn experts(
109    activity: &ExpertActivity,
110    layer: usize,
111    offset: usize,
112    limit: usize,
113) -> Result<Page<Expert>> {
114    ensure!(layer < activity.layers.len(), "unknown expert layer");
115
116    let total = activity.experts_per_layer;
117    let range = range(total, offset, limit)?;
118    let data = range
119        .clone()
120        .map(|expert| Expert {
121            layer,
122            expert,
123            counters: activity.records[layer * total + expert],
124        })
125        .collect();
126
127    page(total, range, data)
128}