cherenkov/qwen4_exp/gpu/
activity.rs1pub(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}