1use std::collections::BTreeMap;
23use std::time::{Duration, Instant};
24
25use crate::{Error, Result};
26
27use crate::bus::query::{Answer, RepeatingQuery, declare_repeating};
28use crate::bus::write::CallTarget;
29use crate::model::registry::SliceSet;
30use crate::report::{BenchReport, OriginLatency};
31
32pub struct BenchSpec<'a> {
34 pub target: &'a CallTarget,
35 pub producer: &'a str,
36 pub procedure: &'a str,
37 pub count: usize,
39 pub concurrency: usize,
41 pub timeout: Duration,
42 pub force: bool,
45}
46
47const FRAMEWORK_READS: [&str; 2] = ["introspect", "describe"];
54
55fn check_idempotent(slices: Option<&SliceSet>, producer: &str, procedure: &str) -> Result<()> {
62 if FRAMEWORK_READS.contains(&procedure) {
63 return Ok(());
64 }
65 let Some(slices) = slices else {
66 return Err(Error::unaskable(
67 format!("{producer}/{procedure}"),
68 "no registry is loaded, so its idempotence is unknown — a benchmark \
69 repeats a call N times, and \"not asked\" is not \"safe to repeat\" \
70 (RFC 09 §5.1 O4). Load a registry, or pass --i-know.",
71 ));
72 };
73 let decl = slices
74 .get(producer)
75 .and_then(|s| s.procedures.iter().find(|p| p.path == procedure));
76 match decl {
77 Some(d) if d.idempotent == Some(true) => Ok(()),
78 Some(d) => Err(Error::unaskable(
79 format!("{producer}/{procedure}"),
80 format!(
81 "declares kind = {:?}, idempotent = {} — repeating it is a write \
82 into a live fleet, not a measurement. Pass --i-know to mean it.",
83 d.kind,
84 match d.idempotent {
85 Some(false) => "false",
86 _ => "(undeclared)",
87 }
88 ),
89 )),
90 None => Err(Error::unaskable(
91 format!("{producer}/{procedure}"),
92 "the loaded registry does not declare it, so nothing says it is safe \
93 to repeat. Pass --i-know to bench it anyway.",
94 )),
95 }
96}
97
98#[derive(Debug, Default, PartialEq, Eq)]
109struct Tally {
110 completed: usize,
111 errors: usize,
112 silent: usize,
113 panicked: usize,
114}
115
116impl Tally {
117 fn record(
120 &mut self,
121 joined: std::result::Result<
122 Result<Vec<(crate::bus::query::FleetAnswer, Duration)>>,
123 tokio::task::JoinError,
124 >,
125 per_origin: &mut BTreeMap<String, Vec<Duration>>,
126 ) {
127 let Ok(result) = joined else {
131 self.panicked += 1;
132 return;
133 };
134 let Ok(answers) = result else {
135 self.errors += 1;
136 return;
137 };
138 self.completed += 1;
139 if answers.is_empty() {
140 self.silent += 1;
143 return;
144 }
145 for (answer, at) in answers {
146 match answer.answer {
147 Answer::Value(_) => per_origin.entry(answer.origin).or_default().push(at),
148 Answer::Error { .. } => self.errors += 1,
149 }
150 }
151 }
152}
153
154fn percentile(sorted: &[Duration], p: f64) -> f64 {
156 if sorted.is_empty() {
157 return 0.0;
158 }
159 let rank = ((p / 100.0) * sorted.len() as f64).ceil() as usize;
160 let idx = rank.saturating_sub(1).min(sorted.len() - 1);
161 sorted[idx].as_secs_f64() * 1000.0
162}
163
164pub async fn run_bench(
166 fleet: &crate::Fleet<'_>,
167 spec: BenchSpec<'_>,
168 slices: Option<&SliceSet>,
169) -> Result<BenchReport> {
170 if !spec.force {
171 check_idempotent(slices, spec.producer, spec.procedure)?;
172 }
173 if spec.count == 0 {
174 return Err(Error::unaskable("--calls 0", "measures nothing"));
175 }
176
177 let segments: Vec<&str> = spec.procedure.split('/').collect();
178 let relative = match spec.target {
179 CallTarget::Host(id) => {
180 let origin = zenkey::origin::RemoteOrigin::from_host(id.clone());
181 zenkey::selector::rpc_at(&origin, spec.producer, &segments).to_string()
182 }
183 CallTarget::Fleet => zenkey::selector::fleet_rpc(spec.producer, &segments).to_string(),
184 CallTarget::Service(origin) => zenkey::selector::service_rpc(origin, &segments).to_string(),
185 };
186 let key = fleet.wire(relative);
187
188 let querier = std::sync::Arc::new(
191 declare_repeating(fleet, &key, spec.timeout)
192 .await
193 .map_err(|e| Error::bus("declare querier", key.clone(), e))?,
194 );
195
196 let concurrency = spec.concurrency.max(1).min(spec.count);
197 let started = Instant::now();
198 let mut per_origin: BTreeMap<String, Vec<Duration>> = BTreeMap::new();
199 let mut tally = Tally::default();
200
201 let mut issued = 0usize;
202 while issued < spec.count {
203 let batch = concurrency.min(spec.count - issued);
204 let mut set = Vec::with_capacity(batch);
205 for _ in 0..batch {
206 let q: std::sync::Arc<RepeatingQuery> = querier.clone();
207 set.push(tokio::spawn(async move { q.fetch_timed().await }));
208 }
209 issued += batch;
210 for handle in set {
211 tally.record(handle.await, &mut per_origin);
212 }
213 }
214 let Tally {
215 completed,
216 errors,
217 silent,
218 panicked,
219 } = tally;
220 let elapsed = started.elapsed();
221 std::sync::Arc::try_unwrap(querier)
222 .map_err(|_| Error::Internal("bench tasks outlived the run".into()))?
223 .undeclare()
224 .await?;
225
226 let origins = per_origin
227 .into_iter()
228 .map(|(origin, mut samples)| {
229 samples.sort_unstable();
230 OriginLatency {
231 origin,
232 replies: samples.len(),
233 min_ms: samples[0].as_secs_f64() * 1000.0,
234 p50_ms: percentile(&samples, 50.0),
235 p95_ms: percentile(&samples, 95.0),
236 p99_ms: percentile(&samples, 99.0),
237 max_ms: samples[samples.len() - 1].as_secs_f64() * 1000.0,
238 }
239 })
240 .collect();
241
242 Ok(BenchReport {
243 key,
244 requested: spec.count,
245 completed,
246 concurrency,
247 errors,
248 silent,
249 panicked,
250 elapsed_s: elapsed.as_secs_f64(),
251 calls_per_s: if elapsed.as_secs_f64() > 0.0 {
252 completed as f64 / elapsed.as_secs_f64()
253 } else {
254 0.0
255 },
256 origins,
257 })
258}
259
260#[cfg(test)]
261mod tests {
262 use super::*;
263 use zenkey::slice::{ProcedureDecl, RegistrySlice};
264
265 fn slices(kind: &str, idempotent: Option<bool>) -> SliceSet {
266 let mut trigger = ProcedureDecl::new("capture/trigger");
267 trigger.kind = Some(zenkey::Declared::parse(kind));
268 trigger.reply = Some("Ack".into());
269 trigger.idempotent = idempotent;
270 let mut slice = RegistrySlice::new("1.0", "t", "netring");
271 slice.procedures = vec![trigger];
272 SliceSet::from_slices(vec![slice])
273 }
274
275 #[test]
279 fn only_a_declared_idempotent_procedure_benches_by_default() {
280 let ok = slices("read", Some(true));
281 assert!(check_idempotent(Some(&ok), "netring", "capture/trigger").is_ok());
282
283 for (kind, idem) in [("write", Some(false)), ("read", None)] {
284 let s = slices(kind, idem);
285 let err = check_idempotent(Some(&s), "netring", "capture/trigger")
286 .unwrap_err()
287 .to_string();
288 assert!(err.contains("--i-know"), "{err}");
289 }
290
291 let s = slices("read", Some(true));
293 assert!(check_idempotent(Some(&s), "netring", "other").is_err());
294 let err = check_idempotent(None, "netring", "capture/trigger")
295 .unwrap_err()
296 .to_string();
297 assert!(err.contains("O4"), "{err}");
298 }
299
300 #[test]
306 fn the_conventions_own_reads_need_no_registry_permission() {
307 for p in ["introspect", "describe"] {
308 assert!(check_idempotent(None, "anything", p).is_ok(), "{p}");
309 }
310 assert!(check_idempotent(None, "anything", "introspect/all").is_err());
312 }
313
314 #[tokio::test]
319 async fn a_panicked_call_is_its_own_population_and_reaches_a_ledger() {
320 let mut per_origin: BTreeMap<String, Vec<Duration>> = BTreeMap::new();
321 let mut tally = Tally::default();
322
323 let join_error = tokio::spawn(async { panic!("a call fell over") })
324 .await
325 .expect_err("the task panicked");
326 tally.record(Err(join_error), &mut per_origin);
327 assert_eq!(
328 tally,
329 Tally {
330 completed: 0,
331 errors: 0,
332 silent: 0,
333 panicked: 1,
334 },
335 "the panic reaches its own ledger and no other"
336 );
337
338 tally.record(
340 Ok(Err(Error::bus("get", "", "the GET failed"))),
341 &mut per_origin,
342 );
343 tally.record(Ok(Ok(vec![])), &mut per_origin);
344 tally.record(
345 Ok(Ok(vec![(
346 crate::bus::query::FleetAnswer {
347 origin: "h-3fa9c2d41b7e".into(),
348 key: "v1/h-3fa9c2d41b7e/@rpc/netring/capture/trigger".into(),
349 encoding: None,
350 attachment: None,
351 answer: Answer::Value(zenoh::bytes::ZBytes::from(b"{}".to_vec())),
352 },
353 Duration::from_millis(3),
354 )])),
355 &mut per_origin,
356 );
357 assert_eq!(
358 tally,
359 Tally {
360 completed: 2,
361 errors: 1,
362 silent: 1,
363 panicked: 1,
364 }
365 );
366 assert_eq!(per_origin["h-3fa9c2d41b7e"], vec![Duration::from_millis(3)]);
367 }
368
369 #[test]
370 fn percentiles_are_nearest_rank_and_survive_one_sample() {
371 let d = |ms: u64| Duration::from_millis(ms);
372 let one = [d(7)];
373 assert_eq!(percentile(&one, 50.0), 7.0);
374 assert_eq!(percentile(&one, 99.0), 7.0);
375
376 let ten: Vec<Duration> = (1..=10).map(d).collect();
377 assert_eq!(percentile(&ten, 50.0), 5.0);
378 assert_eq!(percentile(&ten, 95.0), 10.0);
379 assert_eq!(percentile(&ten, 100.0), 10.0);
380 assert_eq!(percentile(&[], 50.0), 0.0);
382 }
383}