Skip to main content

cherenkov/qwen4_exp/gpu/
activity.rs

1//! Cumulative routed-expert activity. Verification and draft rows count even
2//! when later rolled back. IO counters describe requested reads, not completion.
3
4pub(super) mod layers;
5pub(super) mod reads;
6use super::*;
7pub use layers::LayerStats;
8use reads::ReadSource;
9use serde::{Deserialize, Serialize};
10
11#[derive(Debug, Default, Clone, Copy, Serialize, Deserialize)]
12pub struct ExpertCounters {
13    pub selected_rows: u64,
14    pub cache_hits: u64,
15    pub cache_misses: u64,
16    pub prefetch_requests: u64,
17    pub read_requests: u64,
18    pub read_bytes_requested: u64,
19}
20
21impl ExpertCounters {
22    pub fn add(&mut self, other: Self) {
23        self.selected_rows += other.selected_rows;
24        self.cache_hits += other.cache_hits;
25        self.cache_misses += other.cache_misses;
26        self.prefetch_requests += other.prefetch_requests;
27        self.read_requests += other.read_requests;
28        self.read_bytes_requested += other.read_bytes_requested;
29    }
30
31    fn lookup(&mut self, rows: usize, resident: bool) {
32        self.selected_rows += rows as u64;
33        self.cache_hits += u64::from(resident);
34        self.cache_misses += u64::from(!resident);
35    }
36}
37
38#[derive(Debug, Default, Clone)]
39pub struct ExpertActivity {
40    pub prefill: prefill::PrefillStats,
41    pub layer_prefixes: Vec<String>,
42    pub experts_per_layer: usize,
43    pub records: Vec<ExpertCounters>,
44    pub layers: Vec<LayerStats>,
45    pub elapsed_seconds: f64,
46    pub gpu_timestamps_available: bool,
47    pub gpu_timing: GpuTiming,
48}
49
50impl ExpertActivity {
51    pub(super) fn layer_for_record(&self, record: usize) -> usize {
52        record / self.experts_per_layer
53    }
54
55    pub(crate) fn validate_dimensions(&self) -> Result<()> {
56        anyhow::ensure!(
57            self.layer_prefixes.len() == self.layers.len(),
58            "expert activity layer names and stats differ in length"
59        );
60        anyhow::ensure!(
61            self.layers.len().checked_mul(self.experts_per_layer) == Some(self.records.len()),
62            "expert activity record count does not match layer dimensions"
63        );
64
65        Ok(())
66    }
67
68    pub(super) fn new(layout: &crate::qwen4_exp::ExpertLayout) -> Self {
69        Self {
70            prefill: Default::default(),
71            layers: vec![LayerStats::default(); layout.layers],
72            elapsed_seconds: 0.0,
73            gpu_timestamps_available: false,
74            gpu_timing: GpuTiming::NotInitialized,
75            layer_prefixes: layout.layer_prefixes.clone(),
76            experts_per_layer: layout.experts,
77            records: vec![ExpertCounters::default(); layout.layers * layout.experts],
78        }
79    }
80
81    pub(super) fn lookup(&mut self, record: usize, rows: usize, resident: bool) {
82        self.records[record].lookup(rows, resident);
83    }
84
85    pub(super) fn read(&mut self, record: usize, bytes: usize) {
86        self.records[record].read_requests += 1;
87        self.records[record].read_bytes_requested += bytes as u64;
88    }
89}
90
91impl Gpu<'_> {
92    pub(super) fn record_cut_eligible(
93        &mut self,
94        layer: usize,
95        records: &[usize],
96        missing: &[bool],
97        weights: &[f32],
98    ) {
99        if self.cut_w <= 0.0 || self.fake_experts {
100            return;
101        }
102
103        for ((&record, &missing), &weight) in records.iter().zip(missing).zip(weights) {
104            if missing && weight < self.cut_w {
105                let kind = self.res.kind(record) as usize;
106                self.activity.layers[layer].quant[kind].eligible_weak_misses += 1;
107            }
108        }
109    }
110    pub fn expert_activity(&self) -> &ExpertActivity {
111        &self.activity
112    }
113
114    pub(super) fn record_routed_experts(&mut self, layer: usize, ids: &[u32], experts: &[u32]) {
115        for &expert in experts {
116            let record = self.record_id(layer, expert);
117            let rows = ids.iter().filter(|&&id| id == expert).count();
118
119            self.activity
120                .lookup(record, rows, self.res.is_member(record));
121        }
122    }
123
124    pub(super) fn record_read_plan(
125        &mut self,
126        plan: &mut residency::ReadPlan,
127        need: &[usize],
128        source: ReadSource,
129    ) {
130        if self.fake_experts {
131            return;
132        }
133
134        if let Some(&record) = need.first() {
135            plan.observe(
136                &self.read_tracker,
137                self.activity.layer_for_record(record),
138                source,
139            );
140        }
141
142        for (index, bytes) in plan.read_requests() {
143            self.activity.read(need[index], bytes);
144        }
145    }
146
147    pub fn copy_expert_activity(&self, snapshot: &mut ExpertActivity) {
148        snapshot
149            .layer_prefixes
150            .clone_from(&self.activity.layer_prefixes);
151        snapshot.records.clone_from(&self.activity.records);
152        snapshot.layers.clone_from(&self.activity.layers);
153
154        snapshot.prefill = self.activity.prefill;
155        snapshot.experts_per_layer = self.activity.experts_per_layer;
156        snapshot.elapsed_seconds = self.activity_started.elapsed().as_secs_f64();
157        snapshot.gpu_timestamps_available = self.phase_timer.is_some();
158
159        snapshot.gpu_timing.clone_from(&self.activity.gpu_timing);
160
161        for (layer, reads) in snapshot.layers.iter_mut().zip(self.read_tracker.snapshot()) {
162            for (quant, reads) in layer.quant.iter_mut().zip(reads) {
163                quant.reads = reads;
164            }
165        }
166    }
167
168    pub(super) fn record_quant_selection(&mut self, layer: usize, ids: &[u32], experts: &[u32]) {
169        for &expert in experts {
170            let kind = self.res.kind(self.record_id(layer, expert));
171            let stats = &mut self.activity.layers[layer].quant[kind as usize];
172            stats.selected_experts += 1;
173            stats.selected_rows += ids.iter().filter(|&&id| id == expert).count() as u64;
174        }
175    }
176
177    pub(super) fn record_prediction(&mut self, layer: usize, predicted: &[u32], selected: &[u32]) {
178        let stats = &mut self.activity.layers[layer].prediction;
179        stats.target_batches += 1;
180        stats.selected_unpredicted +=
181            selected.iter().filter(|id| !predicted.contains(id)).count() as u64;
182
183        for expert in predicted {
184            if !selected.contains(expert) {
185                stats.predicted_unused += 1;
186
187                continue;
188            }
189
190            stats.predicted_selected += 1;
191            let record = layer * self.activity.experts_per_layer + *expert as usize;
192            let ticket = self
193                .pending
194                .as_ref()
195                .and_then(|pending| pending.tickets.iter().find(|(r, _)| *r == record));
196
197            if let Some((_, ticket)) = ticket {
198                let ready = ticket.done();
199                stats.needed_prefetch_ready += u64::from(ready);
200                stats.needed_prefetch_late += u64::from(!ready);
201            }
202        }
203    }
204}