1use std::collections::BTreeMap;
23use std::time::{Duration, Instant};
24
25use anyhow::{Result, anyhow, bail};
26use zenoh::Session;
27
28use crate::query::{Answer, RepeatingQuery, declare_repeating};
29use crate::registry::SliceSet;
30use crate::report::{BenchReport, OriginLatency};
31use crate::write::CallTarget;
32
33pub struct BenchSpec<'a> {
35 pub target: &'a CallTarget,
36 pub producer: &'a str,
37 pub procedure: &'a str,
38 pub count: usize,
40 pub concurrency: usize,
42 pub timeout: Duration,
43 pub force: bool,
46}
47
48const FRAMEWORK_READS: [&str; 2] = ["introspect", "describe"];
55
56fn check_idempotent(slices: Option<&SliceSet>, producer: &str, procedure: &str) -> Result<()> {
63 if FRAMEWORK_READS.contains(&procedure) {
64 return Ok(());
65 }
66 let Some(slices) = slices else {
67 bail!(
68 "no registry loaded, so {producer}/{procedure}'s idempotence is unknown — a \
69 benchmark 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) => bail!(
79 "{producer}/{procedure} declares kind = {:?}, idempotent = {} — repeating it is a \
80 write into a live fleet, not a measurement. Pass --i-know to mean it.",
81 d.kind,
82 match d.idempotent {
83 Some(false) => "false",
84 _ => "(undeclared)",
85 }
86 ),
87 None => bail!(
88 "the loaded registry does not declare {producer}/{procedure}, so nothing says it \
89 is safe to repeat. Pass --i-know to bench it anyway."
90 ),
91 }
92}
93
94fn percentile(sorted: &[Duration], p: f64) -> f64 {
96 if sorted.is_empty() {
97 return 0.0;
98 }
99 let rank = ((p / 100.0) * sorted.len() as f64).ceil() as usize;
100 let idx = rank.saturating_sub(1).min(sorted.len() - 1);
101 sorted[idx].as_secs_f64() * 1000.0
102}
103
104pub async fn bench_rpc(
106 session: &Session,
107 base: &str,
108 spec: BenchSpec<'_>,
109 slices: Option<&SliceSet>,
110) -> Result<BenchReport> {
111 if !spec.force {
112 check_idempotent(slices, spec.producer, spec.procedure)?;
113 }
114 if spec.count == 0 {
115 bail!("--count 0 measures nothing");
116 }
117
118 let segments: Vec<&str> = spec.procedure.split('/').collect();
119 let relative = match spec.target {
120 CallTarget::Host(id) => {
121 let origin = zenkey::origin::RemoteOrigin::from_host(id.clone());
122 zenkey::selector::rpc_at(&origin, spec.producer, &segments).to_string()
123 }
124 CallTarget::Fleet => zenkey::selector::fleet_rpc(spec.producer, &segments).to_string(),
125 CallTarget::Service(origin) => zenkey::selector::service_rpc(origin, &segments).to_string(),
126 };
127 let key = zenkey::grammar::with_base(base, relative);
128
129 let querier = std::sync::Arc::new(
132 declare_repeating(session, base, &key, spec.timeout)
133 .await
134 .map_err(|e| anyhow!("declare querier {key}: {e}"))?,
135 );
136
137 let concurrency = spec.concurrency.max(1).min(spec.count);
138 let started = Instant::now();
139 let mut per_origin: BTreeMap<String, Vec<Duration>> = BTreeMap::new();
140 let mut errors = 0usize;
141 let mut silent = 0usize;
142 let mut completed = 0usize;
143
144 let mut issued = 0usize;
145 while issued < spec.count {
146 let batch = concurrency.min(spec.count - issued);
147 let mut set = Vec::with_capacity(batch);
148 for _ in 0..batch {
149 let q: std::sync::Arc<RepeatingQuery> = querier.clone();
150 set.push(tokio::spawn(async move { q.fetch_timed().await }));
151 }
152 issued += batch;
153 for handle in set {
154 let Ok(result) = handle.await else { continue };
155 let answers = match result {
156 Ok(a) => a,
157 Err(_) => {
158 errors += 1;
159 continue;
160 }
161 };
162 completed += 1;
163 if answers.is_empty() {
164 silent += 1;
167 continue;
168 }
169 for (answer, at) in answers {
170 match answer.answer {
171 Answer::Value(_) => per_origin.entry(answer.origin).or_default().push(at),
172 Answer::Error { .. } => errors += 1,
173 }
174 }
175 }
176 }
177 let elapsed = started.elapsed();
178 std::sync::Arc::try_unwrap(querier)
179 .map_err(|_| anyhow!("bench tasks outlived the run"))?
180 .undeclare()
181 .await?;
182
183 let origins = per_origin
184 .into_iter()
185 .map(|(origin, mut samples)| {
186 samples.sort_unstable();
187 OriginLatency {
188 origin,
189 replies: samples.len(),
190 min_ms: samples[0].as_secs_f64() * 1000.0,
191 p50_ms: percentile(&samples, 50.0),
192 p95_ms: percentile(&samples, 95.0),
193 p99_ms: percentile(&samples, 99.0),
194 max_ms: samples[samples.len() - 1].as_secs_f64() * 1000.0,
195 }
196 })
197 .collect();
198
199 Ok(BenchReport {
200 key,
201 requested: spec.count,
202 completed,
203 concurrency,
204 errors,
205 silent,
206 elapsed_s: elapsed.as_secs_f64(),
207 calls_per_s: if elapsed.as_secs_f64() > 0.0 {
208 completed as f64 / elapsed.as_secs_f64()
209 } else {
210 0.0
211 },
212 origins,
213 })
214}
215
216#[cfg(test)]
217mod tests {
218 use super::*;
219 use zenkey::slice::{ProcedureDecl, RegistrySlice};
220
221 fn slices(kind: &str, idempotent: Option<bool>) -> SliceSet {
222 SliceSet::from_slices(vec![RegistrySlice {
223 version: "1.0".into(),
224 app: "t".into(),
225 convention: 1,
226 name: "netring".into(),
227 service_origin: None,
228 description: None,
229 subjects: vec![],
230 procedures: vec![ProcedureDecl {
231 path: "capture/trigger".into(),
232 kind: kind.into(),
233 reply: Some("Ack".into()),
234 request: None,
235 encoding: None,
236 fanout: None,
237 idempotent,
238 since: None,
239 description: None,
240 }],
241 blob: vec![],
242 media: vec![],
243 deprecated: vec![],
244 }])
245 }
246
247 #[test]
251 fn only_a_declared_idempotent_procedure_benches_by_default() {
252 let ok = slices("read", Some(true));
253 assert!(check_idempotent(Some(&ok), "netring", "capture/trigger").is_ok());
254
255 for (kind, idem) in [("write", Some(false)), ("read", None)] {
256 let s = slices(kind, idem);
257 let err = check_idempotent(Some(&s), "netring", "capture/trigger")
258 .unwrap_err()
259 .to_string();
260 assert!(err.contains("--i-know"), "{err}");
261 }
262
263 let s = slices("read", Some(true));
265 assert!(check_idempotent(Some(&s), "netring", "other").is_err());
266 let err = check_idempotent(None, "netring", "capture/trigger")
267 .unwrap_err()
268 .to_string();
269 assert!(err.contains("O4"), "{err}");
270 }
271
272 #[test]
278 fn the_conventions_own_reads_need_no_registry_permission() {
279 for p in ["introspect", "describe"] {
280 assert!(check_idempotent(None, "anything", p).is_ok(), "{p}");
281 }
282 assert!(check_idempotent(None, "anything", "introspect/all").is_err());
284 }
285
286 #[test]
287 fn percentiles_are_nearest_rank_and_survive_one_sample() {
288 let d = |ms: u64| Duration::from_millis(ms);
289 let one = [d(7)];
290 assert_eq!(percentile(&one, 50.0), 7.0);
291 assert_eq!(percentile(&one, 99.0), 7.0);
292
293 let ten: Vec<Duration> = (1..=10).map(d).collect();
294 assert_eq!(percentile(&ten, 50.0), 5.0);
295 assert_eq!(percentile(&ten, 95.0), 10.0);
296 assert_eq!(percentile(&ten, 100.0), 10.0);
297 assert_eq!(percentile(&[], 50.0), 0.0);
299 }
300}