1use std::sync::Arc;
2use std::sync::atomic::AtomicU64;
3use std::sync::atomic::Ordering;
4
5use ursula_shard::BucketStreamId;
6use ursula_shard::CoreId;
7use ursula_shard::RaftGroupId;
8use ursula_shard::ShardPlacement;
9use ursula_stream::StreamErrorCode;
10use ursula_stream::StreamErrorContext;
11
12use crate::engine::GroupEngine;
13use crate::engine::GroupEngineError;
14use crate::error::RuntimeError;
15use crate::request::AppendBatchRequest;
16use crate::request::ColdWriteAdmission;
17use crate::rt::time::Instant;
18
19pub(crate) const GROUP_ACTOR_MAX_WRITE_BATCH: usize = 64;
20pub(crate) const COLD_FLUSH_GROUP_BATCH_MAX_CHUNKS: usize = 64;
21
22#[derive(Debug, Clone)]
23pub struct RuntimeMetrics {
24 pub(crate) inner: Arc<RuntimeMetricsInner>,
25}
26
27impl RuntimeMetrics {
28 pub fn snapshot(&self) -> RuntimeMetricsSnapshot {
29 let per_core_appends = self
30 .inner
31 .per_core_appends
32 .iter()
33 .map(PaddedAtomicU64::load_relaxed)
34 .collect::<Vec<_>>();
35 let accepted_appends = per_core_appends.iter().sum();
36 let per_group_appends = self
37 .inner
38 .per_group_appends
39 .iter()
40 .map(PaddedAtomicU64::load_relaxed)
41 .collect();
42 let per_core_applied_mutations = self
43 .inner
44 .per_core_applied_mutations
45 .iter()
46 .map(PaddedAtomicU64::load_relaxed)
47 .collect::<Vec<_>>();
48 let applied_mutations = per_core_applied_mutations.iter().sum();
49 let per_group_applied_mutations = self
50 .inner
51 .per_group_applied_mutations
52 .iter()
53 .map(PaddedAtomicU64::load_relaxed)
54 .collect();
55 let per_core_mutation_apply_ns = self
56 .inner
57 .per_core_mutation_apply_ns
58 .iter()
59 .map(PaddedAtomicU64::load_relaxed)
60 .collect::<Vec<_>>();
61 let mutation_apply_ns = per_core_mutation_apply_ns.iter().sum();
62 let per_group_mutation_apply_ns = self
63 .inner
64 .per_group_mutation_apply_ns
65 .iter()
66 .map(PaddedAtomicU64::load_relaxed)
67 .collect();
68 let per_core_group_lock_wait_ns = self
69 .inner
70 .per_core_group_lock_wait_ns
71 .iter()
72 .map(PaddedAtomicU64::load_relaxed)
73 .collect::<Vec<_>>();
74 let group_lock_wait_ns = per_core_group_lock_wait_ns.iter().sum();
75 let per_group_group_lock_wait_ns = self
76 .inner
77 .per_group_group_lock_wait_ns
78 .iter()
79 .map(PaddedAtomicU64::load_relaxed)
80 .collect();
81 let per_core_group_engine_exec_ns = self
82 .inner
83 .per_core_group_engine_exec_ns
84 .iter()
85 .map(PaddedAtomicU64::load_relaxed)
86 .collect::<Vec<_>>();
87 let group_engine_exec_ns = per_core_group_engine_exec_ns.iter().sum();
88 let per_group_group_engine_exec_ns = self
89 .inner
90 .per_group_group_engine_exec_ns
91 .iter()
92 .map(PaddedAtomicU64::load_relaxed)
93 .collect();
94 let per_group_group_mailbox_depth = self
95 .inner
96 .per_group_group_mailbox_depth
97 .iter()
98 .map(PaddedAtomicU64::load_relaxed)
99 .collect::<Vec<_>>();
100 let group_mailbox_depth = per_group_group_mailbox_depth.iter().sum();
101 let per_group_group_mailbox_max_depth = self
102 .inner
103 .per_group_group_mailbox_max_depth
104 .iter()
105 .map(PaddedAtomicU64::load_relaxed)
106 .collect::<Vec<_>>();
107 let group_mailbox_max_depth = per_group_group_mailbox_max_depth
108 .iter()
109 .copied()
110 .max()
111 .unwrap_or(0);
112 let per_group_group_mailbox_full_events = self
113 .inner
114 .per_group_group_mailbox_full_events
115 .iter()
116 .map(PaddedAtomicU64::load_relaxed)
117 .collect::<Vec<_>>();
118 let group_mailbox_full_events = per_group_group_mailbox_full_events.iter().sum();
119 let per_core_raft_write_many_batches = self
120 .inner
121 .per_core_raft_write_many_batches
122 .iter()
123 .map(PaddedAtomicU64::load_relaxed)
124 .collect::<Vec<_>>();
125 let raft_write_many_batches = per_core_raft_write_many_batches.iter().sum();
126 let per_group_raft_write_many_batches = self
127 .inner
128 .per_group_raft_write_many_batches
129 .iter()
130 .map(PaddedAtomicU64::load_relaxed)
131 .collect();
132 let per_core_raft_write_many_commands = self
133 .inner
134 .per_core_raft_write_many_commands
135 .iter()
136 .map(PaddedAtomicU64::load_relaxed)
137 .collect::<Vec<_>>();
138 let raft_write_many_commands = per_core_raft_write_many_commands.iter().sum();
139 let per_group_raft_write_many_commands = self
140 .inner
141 .per_group_raft_write_many_commands
142 .iter()
143 .map(PaddedAtomicU64::load_relaxed)
144 .collect();
145 let per_core_raft_write_many_logical_commands = self
146 .inner
147 .per_core_raft_write_many_logical_commands
148 .iter()
149 .map(PaddedAtomicU64::load_relaxed)
150 .collect::<Vec<_>>();
151 let raft_write_many_logical_commands =
152 per_core_raft_write_many_logical_commands.iter().sum();
153 let per_group_raft_write_many_logical_commands = self
154 .inner
155 .per_group_raft_write_many_logical_commands
156 .iter()
157 .map(PaddedAtomicU64::load_relaxed)
158 .collect();
159 let per_core_raft_write_many_responses = self
160 .inner
161 .per_core_raft_write_many_responses
162 .iter()
163 .map(PaddedAtomicU64::load_relaxed)
164 .collect::<Vec<_>>();
165 let raft_write_many_responses = per_core_raft_write_many_responses.iter().sum();
166 let per_group_raft_write_many_responses = self
167 .inner
168 .per_group_raft_write_many_responses
169 .iter()
170 .map(PaddedAtomicU64::load_relaxed)
171 .collect();
172 let per_core_raft_write_many_submit_ns = self
173 .inner
174 .per_core_raft_write_many_submit_ns
175 .iter()
176 .map(PaddedAtomicU64::load_relaxed)
177 .collect::<Vec<_>>();
178 let raft_write_many_submit_ns = per_core_raft_write_many_submit_ns.iter().sum();
179 let per_group_raft_write_many_submit_ns = self
180 .inner
181 .per_group_raft_write_many_submit_ns
182 .iter()
183 .map(PaddedAtomicU64::load_relaxed)
184 .collect();
185 let per_core_raft_write_many_response_ns = self
186 .inner
187 .per_core_raft_write_many_response_ns
188 .iter()
189 .map(PaddedAtomicU64::load_relaxed)
190 .collect::<Vec<_>>();
191 let raft_write_many_response_ns = per_core_raft_write_many_response_ns.iter().sum();
192 let per_group_raft_write_many_response_ns = self
193 .inner
194 .per_group_raft_write_many_response_ns
195 .iter()
196 .map(PaddedAtomicU64::load_relaxed)
197 .collect();
198 let per_core_raft_apply_entries = self
199 .inner
200 .per_core_raft_apply_entries
201 .iter()
202 .map(PaddedAtomicU64::load_relaxed)
203 .collect::<Vec<_>>();
204 let raft_apply_entries = per_core_raft_apply_entries.iter().sum();
205 let per_group_raft_apply_entries = self
206 .inner
207 .per_group_raft_apply_entries
208 .iter()
209 .map(PaddedAtomicU64::load_relaxed)
210 .collect();
211 let per_core_raft_apply_ns = self
212 .inner
213 .per_core_raft_apply_ns
214 .iter()
215 .map(PaddedAtomicU64::load_relaxed)
216 .collect::<Vec<_>>();
217 let raft_apply_ns = per_core_raft_apply_ns.iter().sum();
218 let per_group_raft_apply_ns = self
219 .inner
220 .per_group_raft_apply_ns
221 .iter()
222 .map(PaddedAtomicU64::load_relaxed)
223 .collect();
224 let per_group_raft_snapshot_builds = self
225 .inner
226 .per_group_raft_snapshot_builds
227 .iter()
228 .map(PaddedAtomicU64::load_relaxed)
229 .collect::<Vec<_>>();
230 let raft_snapshot_builds = per_group_raft_snapshot_builds.iter().sum();
231 let per_group_raft_snapshot_build_ns = self
232 .inner
233 .per_group_raft_snapshot_build_ns
234 .iter()
235 .map(PaddedAtomicU64::load_relaxed)
236 .collect::<Vec<_>>();
237 let raft_snapshot_build_ns = per_group_raft_snapshot_build_ns.iter().sum();
238 let per_group_raft_snapshot_body_bytes = self
239 .inner
240 .per_group_raft_snapshot_body_bytes
241 .iter()
242 .map(PaddedAtomicU64::load_relaxed)
243 .collect::<Vec<_>>();
244 let raft_snapshot_body_bytes = per_group_raft_snapshot_body_bytes.iter().sum();
245 let raft_snapshot_body_bytes_max = per_group_raft_snapshot_body_bytes
246 .iter()
247 .copied()
248 .max()
249 .unwrap_or(0);
250 let per_group_raft_snapshot_pointer_bytes = self
251 .inner
252 .per_group_raft_snapshot_pointer_bytes
253 .iter()
254 .map(PaddedAtomicU64::load_relaxed)
255 .collect::<Vec<_>>();
256 let raft_snapshot_pointer_bytes = per_group_raft_snapshot_pointer_bytes.iter().sum();
257 let raft_snapshot_pointer_bytes_max = per_group_raft_snapshot_pointer_bytes
258 .iter()
259 .copied()
260 .max()
261 .unwrap_or(0);
262 let per_group_raft_snapshot_streams = self
263 .inner
264 .per_group_raft_snapshot_streams
265 .iter()
266 .map(PaddedAtomicU64::load_relaxed)
267 .collect::<Vec<_>>();
268 let raft_snapshot_streams = per_group_raft_snapshot_streams.iter().sum();
269 let raft_snapshot_streams_max = per_group_raft_snapshot_streams
270 .iter()
271 .copied()
272 .max()
273 .unwrap_or(0);
274 let per_group_raft_snapshot_external_uploads = self
275 .inner
276 .per_group_raft_snapshot_external_uploads
277 .iter()
278 .map(PaddedAtomicU64::load_relaxed)
279 .collect::<Vec<_>>();
280 let raft_snapshot_external_uploads = per_group_raft_snapshot_external_uploads.iter().sum();
281 let per_group_raft_snapshot_inline_fallbacks = self
282 .inner
283 .per_group_raft_snapshot_inline_fallbacks
284 .iter()
285 .map(PaddedAtomicU64::load_relaxed)
286 .collect::<Vec<_>>();
287 let raft_snapshot_inline_fallbacks = per_group_raft_snapshot_inline_fallbacks.iter().sum();
288 let per_core_live_read_waiters = self
289 .inner
290 .per_core_live_read_waiters
291 .iter()
292 .map(PaddedAtomicU64::load_relaxed)
293 .collect::<Vec<_>>();
294 let live_read_waiters = per_core_live_read_waiters.iter().sum();
295 let per_core_live_read_backpressure_events = self
296 .inner
297 .per_core_live_read_backpressure_events
298 .iter()
299 .map(PaddedAtomicU64::load_relaxed)
300 .collect::<Vec<_>>();
301 let live_read_backpressure_events = per_core_live_read_backpressure_events.iter().sum();
302 let per_core_routed_requests = self
303 .inner
304 .per_core_routed_requests
305 .iter()
306 .map(PaddedAtomicU64::load_relaxed)
307 .collect::<Vec<_>>();
308 let routed_requests = per_core_routed_requests.iter().sum();
309 let per_core_mailbox_send_wait_ns = self
310 .inner
311 .per_core_mailbox_send_wait_ns
312 .iter()
313 .map(PaddedAtomicU64::load_relaxed)
314 .collect::<Vec<_>>();
315 let mailbox_send_wait_ns = per_core_mailbox_send_wait_ns.iter().sum();
316 let per_core_mailbox_full_events = self
317 .inner
318 .per_core_mailbox_full_events
319 .iter()
320 .map(PaddedAtomicU64::load_relaxed)
321 .collect::<Vec<_>>();
322 let mailbox_full_events = per_core_mailbox_full_events.iter().sum();
323 let per_core_wal_batches = self
324 .inner
325 .per_core_wal_batches
326 .iter()
327 .map(PaddedAtomicU64::load_relaxed)
328 .collect::<Vec<_>>();
329 let wal_batches = per_core_wal_batches.iter().sum();
330 let per_group_wal_batches = self
331 .inner
332 .per_group_wal_batches
333 .iter()
334 .map(PaddedAtomicU64::load_relaxed)
335 .collect();
336 let per_core_wal_records = self
337 .inner
338 .per_core_wal_records
339 .iter()
340 .map(PaddedAtomicU64::load_relaxed)
341 .collect::<Vec<_>>();
342 let wal_records = per_core_wal_records.iter().sum();
343 let per_group_wal_records = self
344 .inner
345 .per_group_wal_records
346 .iter()
347 .map(PaddedAtomicU64::load_relaxed)
348 .collect();
349 let per_core_wal_write_ns = self
350 .inner
351 .per_core_wal_write_ns
352 .iter()
353 .map(PaddedAtomicU64::load_relaxed)
354 .collect::<Vec<_>>();
355 let wal_write_ns = per_core_wal_write_ns.iter().sum();
356 let per_group_wal_write_ns = self
357 .inner
358 .per_group_wal_write_ns
359 .iter()
360 .map(PaddedAtomicU64::load_relaxed)
361 .collect();
362 let per_core_wal_sync_ns = self
363 .inner
364 .per_core_wal_sync_ns
365 .iter()
366 .map(PaddedAtomicU64::load_relaxed)
367 .collect::<Vec<_>>();
368 let wal_sync_ns = per_core_wal_sync_ns.iter().sum();
369 let per_group_wal_sync_ns = self
370 .inner
371 .per_group_wal_sync_ns
372 .iter()
373 .map(PaddedAtomicU64::load_relaxed)
374 .collect();
375 let cold_flush_uploads = self.inner.cold_flush_uploads.load_relaxed();
376 let cold_flush_upload_bytes = self.inner.cold_flush_upload_bytes.load_relaxed();
377 let cold_flush_upload_ns = self.inner.cold_flush_upload_ns.load_relaxed();
378 let cold_flush_publishes = self.inner.cold_flush_publishes.load_relaxed();
379 let cold_flush_publish_bytes = self.inner.cold_flush_publish_bytes.load_relaxed();
380 let cold_flush_publish_ns = self.inner.cold_flush_publish_ns.load_relaxed();
381 let cold_orphan_cleanup_attempts = self.inner.cold_orphan_cleanup_attempts.load_relaxed();
382 let cold_orphan_cleanup_errors = self.inner.cold_orphan_cleanup_errors.load_relaxed();
383 let cold_orphan_bytes = self.inner.cold_orphan_bytes.load_relaxed();
384 let cold_gc_reclaimed = self.inner.cold_gc_reclaimed.load_relaxed();
385 let cold_gc_errors = self.inner.cold_gc_errors.load_relaxed();
386 let cold_flush_write_errors = self.inner.cold_flush_write_errors.load_relaxed();
387 let per_group_cold_hot_bytes = self
388 .inner
389 .per_group_cold_hot_bytes
390 .iter()
391 .map(PaddedAtomicU64::load_relaxed)
392 .collect::<Vec<_>>();
393 let cold_hot_bytes = per_group_cold_hot_bytes.iter().sum();
394 let per_group_cold_hot_bytes_max = self
395 .inner
396 .per_group_cold_hot_bytes_max
397 .iter()
398 .map(PaddedAtomicU64::load_relaxed)
399 .collect::<Vec<_>>();
400 let cold_hot_group_bytes_max = per_group_cold_hot_bytes_max
401 .iter()
402 .copied()
403 .max()
404 .unwrap_or(0);
405 let cold_hot_stream_bytes_max = self.inner.cold_hot_stream_bytes_max.load_relaxed();
406 let per_core_cold_backpressure_events = self
407 .inner
408 .per_core_cold_backpressure_events
409 .iter()
410 .map(PaddedAtomicU64::load_relaxed)
411 .collect::<Vec<_>>();
412 let cold_backpressure_events = per_core_cold_backpressure_events.iter().sum();
413 let per_group_cold_backpressure_events = self
414 .inner
415 .per_group_cold_backpressure_events
416 .iter()
417 .map(PaddedAtomicU64::load_relaxed)
418 .collect();
419 let cold_backpressure_bytes = self.inner.cold_backpressure_bytes.load_relaxed();
420
421 RuntimeMetricsSnapshot {
422 accepted_appends,
423 per_core_appends,
424 per_group_appends,
425 applied_mutations,
426 per_core_applied_mutations,
427 per_group_applied_mutations,
428 mutation_apply_ns,
429 per_core_mutation_apply_ns,
430 per_group_mutation_apply_ns,
431 group_lock_wait_ns,
432 per_core_group_lock_wait_ns,
433 per_group_group_lock_wait_ns,
434 group_engine_exec_ns,
435 per_core_group_engine_exec_ns,
436 per_group_group_engine_exec_ns,
437 group_mailbox_depth,
438 per_group_group_mailbox_depth,
439 group_mailbox_max_depth,
440 per_group_group_mailbox_max_depth,
441 group_mailbox_full_events,
442 per_group_group_mailbox_full_events,
443 raft_write_many_batches,
444 per_core_raft_write_many_batches,
445 per_group_raft_write_many_batches,
446 raft_write_many_commands,
447 per_core_raft_write_many_commands,
448 per_group_raft_write_many_commands,
449 raft_write_many_logical_commands,
450 per_core_raft_write_many_logical_commands,
451 per_group_raft_write_many_logical_commands,
452 raft_write_many_responses,
453 per_core_raft_write_many_responses,
454 per_group_raft_write_many_responses,
455 raft_write_many_submit_ns,
456 per_core_raft_write_many_submit_ns,
457 per_group_raft_write_many_submit_ns,
458 raft_write_many_response_ns,
459 per_core_raft_write_many_response_ns,
460 per_group_raft_write_many_response_ns,
461 raft_apply_entries,
462 per_core_raft_apply_entries,
463 per_group_raft_apply_entries,
464 raft_apply_ns,
465 per_core_raft_apply_ns,
466 per_group_raft_apply_ns,
467 raft_snapshot_builds,
468 per_group_raft_snapshot_builds,
469 raft_snapshot_build_ns,
470 per_group_raft_snapshot_build_ns,
471 raft_snapshot_body_bytes,
472 raft_snapshot_body_bytes_max,
473 per_group_raft_snapshot_body_bytes,
474 raft_snapshot_pointer_bytes,
475 raft_snapshot_pointer_bytes_max,
476 per_group_raft_snapshot_pointer_bytes,
477 raft_snapshot_streams,
478 raft_snapshot_streams_max,
479 per_group_raft_snapshot_streams,
480 raft_snapshot_external_uploads,
481 per_group_raft_snapshot_external_uploads,
482 raft_snapshot_inline_fallbacks,
483 per_group_raft_snapshot_inline_fallbacks,
484 live_read_waiters,
485 per_core_live_read_waiters,
486 live_read_backpressure_events,
487 per_core_live_read_backpressure_events,
488 routed_requests,
489 per_core_routed_requests,
490 mailbox_send_wait_ns,
491 per_core_mailbox_send_wait_ns,
492 mailbox_full_events,
493 per_core_mailbox_full_events,
494 wal_batches,
495 per_core_wal_batches,
496 per_group_wal_batches,
497 wal_records,
498 per_core_wal_records,
499 per_group_wal_records,
500 wal_write_ns,
501 per_core_wal_write_ns,
502 per_group_wal_write_ns,
503 wal_sync_ns,
504 per_core_wal_sync_ns,
505 per_group_wal_sync_ns,
506 cold_flush_uploads,
507 cold_flush_upload_bytes,
508 cold_flush_upload_ns,
509 cold_flush_publishes,
510 cold_flush_publish_bytes,
511 cold_flush_publish_ns,
512 cold_orphan_cleanup_attempts,
513 cold_orphan_cleanup_errors,
514 cold_gc_reclaimed,
515 cold_gc_errors,
516 cold_flush_write_errors,
517 cold_orphan_bytes,
518 cold_hot_bytes,
519 per_group_cold_hot_bytes,
520 cold_hot_group_bytes_max,
521 per_group_cold_hot_bytes_max,
522 cold_hot_stream_bytes_max,
523 cold_backpressure_events,
524 per_core_cold_backpressure_events,
525 per_group_cold_backpressure_events,
526 cold_backpressure_bytes,
527 }
528 }
529}
530
531#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
532pub struct RuntimeMetricsSnapshot {
533 pub accepted_appends: u64,
534 pub per_core_appends: Vec<u64>,
535 pub per_group_appends: Vec<u64>,
536 pub applied_mutations: u64,
537 pub per_core_applied_mutations: Vec<u64>,
538 pub per_group_applied_mutations: Vec<u64>,
539 pub mutation_apply_ns: u64,
540 pub per_core_mutation_apply_ns: Vec<u64>,
541 pub per_group_mutation_apply_ns: Vec<u64>,
542 pub group_lock_wait_ns: u64,
543 pub per_core_group_lock_wait_ns: Vec<u64>,
544 pub per_group_group_lock_wait_ns: Vec<u64>,
545 pub group_engine_exec_ns: u64,
546 pub per_core_group_engine_exec_ns: Vec<u64>,
547 pub per_group_group_engine_exec_ns: Vec<u64>,
548 pub group_mailbox_depth: u64,
549 pub per_group_group_mailbox_depth: Vec<u64>,
550 pub group_mailbox_max_depth: u64,
551 pub per_group_group_mailbox_max_depth: Vec<u64>,
552 pub group_mailbox_full_events: u64,
553 pub per_group_group_mailbox_full_events: Vec<u64>,
554 pub raft_write_many_batches: u64,
555 pub per_core_raft_write_many_batches: Vec<u64>,
556 pub per_group_raft_write_many_batches: Vec<u64>,
557 pub raft_write_many_commands: u64,
558 pub per_core_raft_write_many_commands: Vec<u64>,
559 pub per_group_raft_write_many_commands: Vec<u64>,
560 pub raft_write_many_logical_commands: u64,
561 pub per_core_raft_write_many_logical_commands: Vec<u64>,
562 pub per_group_raft_write_many_logical_commands: Vec<u64>,
563 pub raft_write_many_responses: u64,
564 pub per_core_raft_write_many_responses: Vec<u64>,
565 pub per_group_raft_write_many_responses: Vec<u64>,
566 pub raft_write_many_submit_ns: u64,
567 pub per_core_raft_write_many_submit_ns: Vec<u64>,
568 pub per_group_raft_write_many_submit_ns: Vec<u64>,
569 pub raft_write_many_response_ns: u64,
570 pub per_core_raft_write_many_response_ns: Vec<u64>,
571 pub per_group_raft_write_many_response_ns: Vec<u64>,
572 pub raft_apply_entries: u64,
573 pub per_core_raft_apply_entries: Vec<u64>,
574 pub per_group_raft_apply_entries: Vec<u64>,
575 pub raft_apply_ns: u64,
576 pub per_core_raft_apply_ns: Vec<u64>,
577 pub per_group_raft_apply_ns: Vec<u64>,
578 pub raft_snapshot_builds: u64,
579 pub per_group_raft_snapshot_builds: Vec<u64>,
580 pub raft_snapshot_build_ns: u64,
581 pub per_group_raft_snapshot_build_ns: Vec<u64>,
582 pub raft_snapshot_body_bytes: u64,
583 pub raft_snapshot_body_bytes_max: u64,
584 pub per_group_raft_snapshot_body_bytes: Vec<u64>,
585 pub raft_snapshot_pointer_bytes: u64,
586 pub raft_snapshot_pointer_bytes_max: u64,
587 pub per_group_raft_snapshot_pointer_bytes: Vec<u64>,
588 pub raft_snapshot_streams: u64,
589 pub raft_snapshot_streams_max: u64,
590 pub per_group_raft_snapshot_streams: Vec<u64>,
591 pub raft_snapshot_external_uploads: u64,
592 pub per_group_raft_snapshot_external_uploads: Vec<u64>,
593 pub raft_snapshot_inline_fallbacks: u64,
594 pub per_group_raft_snapshot_inline_fallbacks: Vec<u64>,
595 pub live_read_waiters: u64,
596 pub per_core_live_read_waiters: Vec<u64>,
597 pub live_read_backpressure_events: u64,
598 pub per_core_live_read_backpressure_events: Vec<u64>,
599 pub routed_requests: u64,
600 pub per_core_routed_requests: Vec<u64>,
601 pub mailbox_send_wait_ns: u64,
602 pub per_core_mailbox_send_wait_ns: Vec<u64>,
603 pub mailbox_full_events: u64,
604 pub per_core_mailbox_full_events: Vec<u64>,
605 pub wal_batches: u64,
606 pub per_core_wal_batches: Vec<u64>,
607 pub per_group_wal_batches: Vec<u64>,
608 pub wal_records: u64,
609 pub per_core_wal_records: Vec<u64>,
610 pub per_group_wal_records: Vec<u64>,
611 pub wal_write_ns: u64,
612 pub per_core_wal_write_ns: Vec<u64>,
613 pub per_group_wal_write_ns: Vec<u64>,
614 pub wal_sync_ns: u64,
615 pub per_core_wal_sync_ns: Vec<u64>,
616 pub per_group_wal_sync_ns: Vec<u64>,
617 pub cold_flush_uploads: u64,
618 pub cold_flush_upload_bytes: u64,
619 pub cold_flush_upload_ns: u64,
620 pub cold_flush_publishes: u64,
621 pub cold_flush_publish_bytes: u64,
622 pub cold_flush_publish_ns: u64,
623 pub cold_orphan_cleanup_attempts: u64,
624 pub cold_orphan_cleanup_errors: u64,
625 pub cold_orphan_bytes: u64,
626 pub cold_gc_reclaimed: u64,
627 pub cold_gc_errors: u64,
628 pub cold_flush_write_errors: u64,
629 pub cold_hot_bytes: u64,
630 pub per_group_cold_hot_bytes: Vec<u64>,
631 pub cold_hot_group_bytes_max: u64,
632 pub per_group_cold_hot_bytes_max: Vec<u64>,
633 pub cold_hot_stream_bytes_max: u64,
634 pub cold_backpressure_events: u64,
635 pub per_core_cold_backpressure_events: Vec<u64>,
636 pub per_group_cold_backpressure_events: Vec<u64>,
637 pub cold_backpressure_bytes: u64,
638}
639
640#[derive(Debug, Clone, PartialEq, Eq)]
641pub struct RuntimeMailboxSnapshot {
642 pub depths: Vec<usize>,
643 pub capacities: Vec<usize>,
644}
645
646#[derive(Debug)]
647pub(crate) struct RuntimeMetricsInner {
648 pub(crate) per_core_appends: Vec<PaddedAtomicU64>,
649 pub(crate) per_group_appends: Vec<PaddedAtomicU64>,
650 pub(crate) per_core_applied_mutations: Vec<PaddedAtomicU64>,
651 pub(crate) per_group_applied_mutations: Vec<PaddedAtomicU64>,
652 pub(crate) per_core_mutation_apply_ns: Vec<PaddedAtomicU64>,
653 pub(crate) per_group_mutation_apply_ns: Vec<PaddedAtomicU64>,
654 pub(crate) per_core_group_lock_wait_ns: Vec<PaddedAtomicU64>,
655 pub(crate) per_group_group_lock_wait_ns: Vec<PaddedAtomicU64>,
656 pub(crate) per_core_group_engine_exec_ns: Vec<PaddedAtomicU64>,
657 pub(crate) per_group_group_engine_exec_ns: Vec<PaddedAtomicU64>,
658 pub(crate) per_group_group_mailbox_depth: Vec<PaddedAtomicU64>,
659 pub(crate) per_group_group_mailbox_max_depth: Vec<PaddedAtomicU64>,
660 pub(crate) per_group_group_mailbox_full_events: Vec<PaddedAtomicU64>,
661 pub(crate) per_core_raft_write_many_batches: Vec<PaddedAtomicU64>,
662 pub(crate) per_group_raft_write_many_batches: Vec<PaddedAtomicU64>,
663 pub(crate) per_core_raft_write_many_commands: Vec<PaddedAtomicU64>,
664 pub(crate) per_group_raft_write_many_commands: Vec<PaddedAtomicU64>,
665 pub(crate) per_core_raft_write_many_logical_commands: Vec<PaddedAtomicU64>,
666 pub(crate) per_group_raft_write_many_logical_commands: Vec<PaddedAtomicU64>,
667 pub(crate) per_core_raft_write_many_responses: Vec<PaddedAtomicU64>,
668 pub(crate) per_group_raft_write_many_responses: Vec<PaddedAtomicU64>,
669 pub(crate) per_core_raft_write_many_submit_ns: Vec<PaddedAtomicU64>,
670 pub(crate) per_group_raft_write_many_submit_ns: Vec<PaddedAtomicU64>,
671 pub(crate) per_core_raft_write_many_response_ns: Vec<PaddedAtomicU64>,
672 pub(crate) per_group_raft_write_many_response_ns: Vec<PaddedAtomicU64>,
673 pub(crate) per_core_raft_apply_entries: Vec<PaddedAtomicU64>,
674 pub(crate) per_group_raft_apply_entries: Vec<PaddedAtomicU64>,
675 pub(crate) per_core_raft_apply_ns: Vec<PaddedAtomicU64>,
676 pub(crate) per_group_raft_apply_ns: Vec<PaddedAtomicU64>,
677 pub(crate) per_group_raft_snapshot_builds: Vec<PaddedAtomicU64>,
678 pub(crate) per_group_raft_snapshot_build_ns: Vec<PaddedAtomicU64>,
679 pub(crate) per_group_raft_snapshot_body_bytes: Vec<PaddedAtomicU64>,
680 pub(crate) per_group_raft_snapshot_pointer_bytes: Vec<PaddedAtomicU64>,
681 pub(crate) per_group_raft_snapshot_streams: Vec<PaddedAtomicU64>,
682 pub(crate) per_group_raft_snapshot_external_uploads: Vec<PaddedAtomicU64>,
683 pub(crate) per_group_raft_snapshot_inline_fallbacks: Vec<PaddedAtomicU64>,
684 pub(crate) per_core_live_read_waiters: Vec<PaddedAtomicU64>,
685 pub(crate) per_core_live_read_backpressure_events: Vec<PaddedAtomicU64>,
686 pub(crate) per_core_routed_requests: Vec<PaddedAtomicU64>,
687 pub(crate) per_core_mailbox_send_wait_ns: Vec<PaddedAtomicU64>,
688 pub(crate) per_core_mailbox_full_events: Vec<PaddedAtomicU64>,
689 pub(crate) per_core_wal_batches: Vec<PaddedAtomicU64>,
690 pub(crate) per_group_wal_batches: Vec<PaddedAtomicU64>,
691 pub(crate) per_core_wal_records: Vec<PaddedAtomicU64>,
692 pub(crate) per_group_wal_records: Vec<PaddedAtomicU64>,
693 pub(crate) per_core_wal_write_ns: Vec<PaddedAtomicU64>,
694 pub(crate) per_group_wal_write_ns: Vec<PaddedAtomicU64>,
695 pub(crate) per_core_wal_sync_ns: Vec<PaddedAtomicU64>,
696 pub(crate) per_group_wal_sync_ns: Vec<PaddedAtomicU64>,
697 pub(crate) cold_flush_uploads: PaddedAtomicU64,
698 pub(crate) cold_flush_upload_bytes: PaddedAtomicU64,
699 pub(crate) cold_flush_upload_ns: PaddedAtomicU64,
700 pub(crate) cold_flush_publishes: PaddedAtomicU64,
701 pub(crate) cold_flush_publish_bytes: PaddedAtomicU64,
702 pub(crate) cold_flush_publish_ns: PaddedAtomicU64,
703 pub(crate) cold_orphan_cleanup_attempts: PaddedAtomicU64,
704 pub(crate) cold_orphan_cleanup_errors: PaddedAtomicU64,
705 pub(crate) cold_gc_reclaimed: PaddedAtomicU64,
706 pub(crate) cold_gc_errors: PaddedAtomicU64,
707 pub(crate) cold_flush_write_errors: PaddedAtomicU64,
708 pub(crate) cold_orphan_bytes: PaddedAtomicU64,
709 pub(crate) per_group_cold_hot_bytes: Vec<PaddedAtomicU64>,
710 pub(crate) per_group_cold_hot_bytes_max: Vec<PaddedAtomicU64>,
711 pub(crate) cold_hot_stream_bytes_max: PaddedAtomicU64,
712 pub(crate) per_core_cold_backpressure_events: Vec<PaddedAtomicU64>,
713 pub(crate) per_group_cold_backpressure_events: Vec<PaddedAtomicU64>,
714 pub(crate) cold_backpressure_bytes: PaddedAtomicU64,
715}
716
717#[derive(Debug, Clone, Copy)]
718pub(crate) struct RaftWriteManySample {
719 pub(crate) command_count: u64,
720 pub(crate) logical_command_count: u64,
721 pub(crate) response_count: u64,
722 pub(crate) submit_ns: u64,
723 pub(crate) response_ns: u64,
724}
725
726#[derive(Debug, Clone, Copy)]
727pub(crate) struct RaftSnapshotBuildSample {
728 pub(crate) streams: u64,
729 pub(crate) body_bytes: u64,
730 pub(crate) pointer_bytes: u64,
731 pub(crate) build_ns: u64,
732 pub(crate) external_upload: bool,
733 pub(crate) inline_fallback: bool,
734}
735
736impl RuntimeMetricsInner {
737 pub(crate) fn new(core_count: usize, raft_group_count: usize) -> Self {
738 Self {
739 per_core_appends: (0..core_count).map(|_| PaddedAtomicU64::new(0)).collect(),
740 per_group_appends: (0..raft_group_count)
741 .map(|_| PaddedAtomicU64::new(0))
742 .collect(),
743 per_core_applied_mutations: (0..core_count).map(|_| PaddedAtomicU64::new(0)).collect(),
744 per_group_applied_mutations: (0..raft_group_count)
745 .map(|_| PaddedAtomicU64::new(0))
746 .collect(),
747 per_core_mutation_apply_ns: (0..core_count).map(|_| PaddedAtomicU64::new(0)).collect(),
748 per_group_mutation_apply_ns: (0..raft_group_count)
749 .map(|_| PaddedAtomicU64::new(0))
750 .collect(),
751 per_core_group_lock_wait_ns: (0..core_count).map(|_| PaddedAtomicU64::new(0)).collect(),
752 per_group_group_lock_wait_ns: (0..raft_group_count)
753 .map(|_| PaddedAtomicU64::new(0))
754 .collect(),
755 per_core_group_engine_exec_ns: (0..core_count)
756 .map(|_| PaddedAtomicU64::new(0))
757 .collect(),
758 per_group_group_engine_exec_ns: (0..raft_group_count)
759 .map(|_| PaddedAtomicU64::new(0))
760 .collect(),
761 per_group_group_mailbox_depth: (0..raft_group_count)
762 .map(|_| PaddedAtomicU64::new(0))
763 .collect(),
764 per_group_group_mailbox_max_depth: (0..raft_group_count)
765 .map(|_| PaddedAtomicU64::new(0))
766 .collect(),
767 per_group_group_mailbox_full_events: (0..raft_group_count)
768 .map(|_| PaddedAtomicU64::new(0))
769 .collect(),
770 per_core_raft_write_many_batches: (0..core_count)
771 .map(|_| PaddedAtomicU64::new(0))
772 .collect(),
773 per_group_raft_write_many_batches: (0..raft_group_count)
774 .map(|_| PaddedAtomicU64::new(0))
775 .collect(),
776 per_core_raft_write_many_commands: (0..core_count)
777 .map(|_| PaddedAtomicU64::new(0))
778 .collect(),
779 per_group_raft_write_many_commands: (0..raft_group_count)
780 .map(|_| PaddedAtomicU64::new(0))
781 .collect(),
782 per_core_raft_write_many_logical_commands: (0..core_count)
783 .map(|_| PaddedAtomicU64::new(0))
784 .collect(),
785 per_group_raft_write_many_logical_commands: (0..raft_group_count)
786 .map(|_| PaddedAtomicU64::new(0))
787 .collect(),
788 per_core_raft_write_many_responses: (0..core_count)
789 .map(|_| PaddedAtomicU64::new(0))
790 .collect(),
791 per_group_raft_write_many_responses: (0..raft_group_count)
792 .map(|_| PaddedAtomicU64::new(0))
793 .collect(),
794 per_core_raft_write_many_submit_ns: (0..core_count)
795 .map(|_| PaddedAtomicU64::new(0))
796 .collect(),
797 per_group_raft_write_many_submit_ns: (0..raft_group_count)
798 .map(|_| PaddedAtomicU64::new(0))
799 .collect(),
800 per_core_raft_write_many_response_ns: (0..core_count)
801 .map(|_| PaddedAtomicU64::new(0))
802 .collect(),
803 per_group_raft_write_many_response_ns: (0..raft_group_count)
804 .map(|_| PaddedAtomicU64::new(0))
805 .collect(),
806 per_core_raft_apply_entries: (0..core_count).map(|_| PaddedAtomicU64::new(0)).collect(),
807 per_group_raft_apply_entries: (0..raft_group_count)
808 .map(|_| PaddedAtomicU64::new(0))
809 .collect(),
810 per_core_raft_apply_ns: (0..core_count).map(|_| PaddedAtomicU64::new(0)).collect(),
811 per_group_raft_apply_ns: (0..raft_group_count)
812 .map(|_| PaddedAtomicU64::new(0))
813 .collect(),
814 per_group_raft_snapshot_builds: (0..raft_group_count)
815 .map(|_| PaddedAtomicU64::new(0))
816 .collect(),
817 per_group_raft_snapshot_build_ns: (0..raft_group_count)
818 .map(|_| PaddedAtomicU64::new(0))
819 .collect(),
820 per_group_raft_snapshot_body_bytes: (0..raft_group_count)
821 .map(|_| PaddedAtomicU64::new(0))
822 .collect(),
823 per_group_raft_snapshot_pointer_bytes: (0..raft_group_count)
824 .map(|_| PaddedAtomicU64::new(0))
825 .collect(),
826 per_group_raft_snapshot_streams: (0..raft_group_count)
827 .map(|_| PaddedAtomicU64::new(0))
828 .collect(),
829 per_group_raft_snapshot_external_uploads: (0..raft_group_count)
830 .map(|_| PaddedAtomicU64::new(0))
831 .collect(),
832 per_group_raft_snapshot_inline_fallbacks: (0..raft_group_count)
833 .map(|_| PaddedAtomicU64::new(0))
834 .collect(),
835 per_core_live_read_waiters: (0..core_count).map(|_| PaddedAtomicU64::new(0)).collect(),
836 per_core_live_read_backpressure_events: (0..core_count)
837 .map(|_| PaddedAtomicU64::new(0))
838 .collect(),
839 per_core_routed_requests: (0..core_count).map(|_| PaddedAtomicU64::new(0)).collect(),
840 per_core_mailbox_send_wait_ns: (0..core_count)
841 .map(|_| PaddedAtomicU64::new(0))
842 .collect(),
843 per_core_mailbox_full_events: (0..core_count)
844 .map(|_| PaddedAtomicU64::new(0))
845 .collect(),
846 per_core_wal_batches: (0..core_count).map(|_| PaddedAtomicU64::new(0)).collect(),
847 per_group_wal_batches: (0..raft_group_count)
848 .map(|_| PaddedAtomicU64::new(0))
849 .collect(),
850 per_core_wal_records: (0..core_count).map(|_| PaddedAtomicU64::new(0)).collect(),
851 per_group_wal_records: (0..raft_group_count)
852 .map(|_| PaddedAtomicU64::new(0))
853 .collect(),
854 per_core_wal_write_ns: (0..core_count).map(|_| PaddedAtomicU64::new(0)).collect(),
855 per_group_wal_write_ns: (0..raft_group_count)
856 .map(|_| PaddedAtomicU64::new(0))
857 .collect(),
858 per_core_wal_sync_ns: (0..core_count).map(|_| PaddedAtomicU64::new(0)).collect(),
859 per_group_wal_sync_ns: (0..raft_group_count)
860 .map(|_| PaddedAtomicU64::new(0))
861 .collect(),
862 cold_flush_uploads: PaddedAtomicU64::new(0),
863 cold_flush_upload_bytes: PaddedAtomicU64::new(0),
864 cold_flush_upload_ns: PaddedAtomicU64::new(0),
865 cold_flush_publishes: PaddedAtomicU64::new(0),
866 cold_flush_publish_bytes: PaddedAtomicU64::new(0),
867 cold_flush_publish_ns: PaddedAtomicU64::new(0),
868 cold_orphan_cleanup_attempts: PaddedAtomicU64::new(0),
869 cold_orphan_cleanup_errors: PaddedAtomicU64::new(0),
870 cold_gc_reclaimed: PaddedAtomicU64::new(0),
871 cold_gc_errors: PaddedAtomicU64::new(0),
872 cold_flush_write_errors: PaddedAtomicU64::new(0),
873 cold_orphan_bytes: PaddedAtomicU64::new(0),
874 per_group_cold_hot_bytes: (0..raft_group_count)
875 .map(|_| PaddedAtomicU64::new(0))
876 .collect(),
877 per_group_cold_hot_bytes_max: (0..raft_group_count)
878 .map(|_| PaddedAtomicU64::new(0))
879 .collect(),
880 cold_hot_stream_bytes_max: PaddedAtomicU64::new(0),
881 per_core_cold_backpressure_events: (0..core_count)
882 .map(|_| PaddedAtomicU64::new(0))
883 .collect(),
884 per_group_cold_backpressure_events: (0..raft_group_count)
885 .map(|_| PaddedAtomicU64::new(0))
886 .collect(),
887 cold_backpressure_bytes: PaddedAtomicU64::new(0),
888 }
889 }
890
891 pub(crate) fn record_routed_request(&self, core_id: CoreId, mailbox_send_wait_ns: u64) {
892 let index = usize::from(core_id.0);
893 self.per_core_routed_requests[index].fetch_add_relaxed(1);
894 self.per_core_mailbox_send_wait_ns[index].fetch_add_relaxed(mailbox_send_wait_ns);
895 }
896
897 pub(crate) fn record_mailbox_full(&self, core_id: CoreId) {
898 self.per_core_mailbox_full_events[usize::from(core_id.0)].fetch_add_relaxed(1);
899 }
900
901 pub(crate) fn record_append(&self, core_id: CoreId, group_id: RaftGroupId) {
902 self.record_append_batch(core_id, group_id, 1);
903 }
904
905 pub(crate) fn record_append_batch(&self, core_id: CoreId, group_id: RaftGroupId, count: u64) {
906 self.per_core_appends[usize::from(core_id.0)].fetch_add_relaxed(count);
907 self.per_group_appends[usize::try_from(group_id.0).expect("u32 fits usize")]
908 .fetch_add_relaxed(count);
909 }
910
911 pub(crate) fn record_applied_mutation(
912 &self,
913 core_id: CoreId,
914 group_id: RaftGroupId,
915 apply_ns: u64,
916 ) {
917 self.record_applied_mutation_batch(core_id, group_id, 1, apply_ns);
918 }
919
920 pub(crate) fn record_applied_mutation_batch(
921 &self,
922 core_id: CoreId,
923 group_id: RaftGroupId,
924 count: u64,
925 apply_ns: u64,
926 ) {
927 let core_index = usize::from(core_id.0);
928 let group_index = usize::try_from(group_id.0).expect("u32 fits usize");
929 self.per_core_applied_mutations[core_index].fetch_add_relaxed(count);
930 self.per_group_applied_mutations[group_index].fetch_add_relaxed(count);
931 self.per_core_mutation_apply_ns[core_index].fetch_add_relaxed(apply_ns);
932 self.per_group_mutation_apply_ns[group_index].fetch_add_relaxed(apply_ns);
933 }
934
935 pub(crate) fn record_group_engine_exec(
936 &self,
937 core_id: CoreId,
938 group_id: RaftGroupId,
939 exec_ns: u64,
940 ) {
941 let core_index = usize::from(core_id.0);
942 let group_index = usize::try_from(group_id.0).expect("u32 fits usize");
943 self.per_core_group_engine_exec_ns[core_index].fetch_add_relaxed(exec_ns);
944 self.per_group_group_engine_exec_ns[group_index].fetch_add_relaxed(exec_ns);
945 }
946
947 pub(crate) fn record_group_mailbox_enqueued(&self, group_id: RaftGroupId) {
948 let group_index = usize::try_from(group_id.0).expect("u32 fits usize");
949 let depth = self.per_group_group_mailbox_depth[group_index]
950 .fetch_add_relaxed(1)
951 .saturating_add(1);
952 self.per_group_group_mailbox_max_depth[group_index].fetch_max_relaxed(depth);
953 }
954
955 pub(crate) fn record_group_mailbox_dequeued(&self, group_id: RaftGroupId) {
956 let group_index = usize::try_from(group_id.0).expect("u32 fits usize");
957 self.per_group_group_mailbox_depth[group_index].fetch_sub_saturating_relaxed(1);
958 }
959
960 pub(crate) fn record_group_mailbox_full(&self, group_id: RaftGroupId) {
961 let group_index = usize::try_from(group_id.0).expect("u32 fits usize");
962 self.per_group_group_mailbox_full_events[group_index].fetch_add_relaxed(1);
963 }
964
965 pub(crate) fn record_raft_write_many(
966 &self,
967 core_id: CoreId,
968 group_id: RaftGroupId,
969 sample: RaftWriteManySample,
970 ) {
971 let core_index = usize::from(core_id.0);
972 let group_index = usize::try_from(group_id.0).expect("u32 fits usize");
973 self.per_core_raft_write_many_batches[core_index].fetch_add_relaxed(1);
974 self.per_group_raft_write_many_batches[group_index].fetch_add_relaxed(1);
975 self.per_core_raft_write_many_commands[core_index].fetch_add_relaxed(sample.command_count);
976 self.per_group_raft_write_many_commands[group_index]
977 .fetch_add_relaxed(sample.command_count);
978 self.per_core_raft_write_many_logical_commands[core_index]
979 .fetch_add_relaxed(sample.logical_command_count);
980 self.per_group_raft_write_many_logical_commands[group_index]
981 .fetch_add_relaxed(sample.logical_command_count);
982 self.per_core_raft_write_many_responses[core_index]
983 .fetch_add_relaxed(sample.response_count);
984 self.per_group_raft_write_many_responses[group_index]
985 .fetch_add_relaxed(sample.response_count);
986 self.per_core_raft_write_many_submit_ns[core_index].fetch_add_relaxed(sample.submit_ns);
987 self.per_group_raft_write_many_submit_ns[group_index].fetch_add_relaxed(sample.submit_ns);
988 self.per_core_raft_write_many_response_ns[core_index].fetch_add_relaxed(sample.response_ns);
989 self.per_group_raft_write_many_response_ns[group_index]
990 .fetch_add_relaxed(sample.response_ns);
991 }
992
993 pub(crate) fn record_raft_apply_batch(
994 &self,
995 core_id: CoreId,
996 group_id: RaftGroupId,
997 entry_count: u64,
998 apply_ns: u64,
999 ) {
1000 let core_index = usize::from(core_id.0);
1001 let group_index = usize::try_from(group_id.0).expect("u32 fits usize");
1002 self.per_core_raft_apply_entries[core_index].fetch_add_relaxed(entry_count);
1003 self.per_group_raft_apply_entries[group_index].fetch_add_relaxed(entry_count);
1004 self.per_core_raft_apply_ns[core_index].fetch_add_relaxed(apply_ns);
1005 self.per_group_raft_apply_ns[group_index].fetch_add_relaxed(apply_ns);
1006 }
1007
1008 pub(crate) fn record_raft_snapshot_build(
1009 &self,
1010 group_id: RaftGroupId,
1011 sample: RaftSnapshotBuildSample,
1012 ) {
1013 let group_index = usize::try_from(group_id.0).expect("u32 fits usize");
1014 self.per_group_raft_snapshot_builds[group_index].fetch_add_relaxed(1);
1015 self.per_group_raft_snapshot_build_ns[group_index].fetch_add_relaxed(sample.build_ns);
1016 self.per_group_raft_snapshot_body_bytes[group_index].store_relaxed(sample.body_bytes);
1017 self.per_group_raft_snapshot_pointer_bytes[group_index].store_relaxed(sample.pointer_bytes);
1018 self.per_group_raft_snapshot_streams[group_index].store_relaxed(sample.streams);
1019 if sample.external_upload {
1020 self.per_group_raft_snapshot_external_uploads[group_index].fetch_add_relaxed(1);
1021 }
1022 if sample.inline_fallback {
1023 self.per_group_raft_snapshot_inline_fallbacks[group_index].fetch_add_relaxed(1);
1024 }
1025 }
1026
1027 pub(crate) fn record_wal_batch(
1028 &self,
1029 core_id: CoreId,
1030 group_id: RaftGroupId,
1031 record_count: u64,
1032 write_ns: u64,
1033 sync_ns: u64,
1034 ) {
1035 let core_index = usize::from(core_id.0);
1036 let group_index = usize::try_from(group_id.0).expect("u32 fits usize");
1037 self.per_core_wal_batches[core_index].fetch_add_relaxed(1);
1038 self.per_group_wal_batches[group_index].fetch_add_relaxed(1);
1039 self.per_core_wal_records[core_index].fetch_add_relaxed(record_count);
1040 self.per_group_wal_records[group_index].fetch_add_relaxed(record_count);
1041 self.per_core_wal_write_ns[core_index].fetch_add_relaxed(write_ns);
1042 self.per_group_wal_write_ns[group_index].fetch_add_relaxed(write_ns);
1043 self.per_core_wal_sync_ns[core_index].fetch_add_relaxed(sync_ns);
1044 self.per_group_wal_sync_ns[group_index].fetch_add_relaxed(sync_ns);
1045 }
1046
1047 pub(crate) fn record_cold_upload(&self, bytes: u64, upload_ns: u64) {
1048 self.cold_flush_uploads.fetch_add_relaxed(1);
1049 self.cold_flush_upload_bytes.fetch_add_relaxed(bytes);
1050 self.cold_flush_upload_ns.fetch_add_relaxed(upload_ns);
1051 }
1052
1053 pub(crate) fn record_cold_publish(&self, bytes: u64, publish_ns: u64) {
1054 self.cold_flush_publishes.fetch_add_relaxed(1);
1055 self.cold_flush_publish_bytes.fetch_add_relaxed(bytes);
1056 self.cold_flush_publish_ns.fetch_add_relaxed(publish_ns);
1057 }
1058
1059 pub(crate) fn record_cold_gc_reclaimed(&self, entries: u64) {
1060 self.cold_gc_reclaimed.fetch_add_relaxed(entries);
1061 }
1062
1063 pub(crate) fn record_cold_flush_write_error(&self) {
1064 self.cold_flush_write_errors.fetch_add_relaxed(1);
1065 }
1066
1067 pub(crate) fn record_cold_gc_error(&self) {
1068 self.cold_gc_errors.fetch_add_relaxed(1);
1069 }
1070
1071 pub(crate) fn record_cold_hot_backlog(
1072 &self,
1073 group_id: RaftGroupId,
1074 stream_hot_bytes: u64,
1075 group_hot_bytes: u64,
1076 ) {
1077 let group_index = usize::try_from(group_id.0).expect("u32 fits usize");
1078 self.per_group_cold_hot_bytes[group_index].store_relaxed(group_hot_bytes);
1079 self.per_group_cold_hot_bytes_max[group_index].fetch_max_relaxed(group_hot_bytes);
1080 self.cold_hot_stream_bytes_max
1081 .fetch_max_relaxed(stream_hot_bytes);
1082 }
1083
1084 pub(crate) fn record_cold_backpressure(
1085 &self,
1086 core_id: CoreId,
1087 group_id: RaftGroupId,
1088 incoming_bytes: u64,
1089 _limit: u64,
1090 ) {
1091 let core_index = usize::from(core_id.0);
1092 let group_index = usize::try_from(group_id.0).expect("u32 fits usize");
1093 self.per_core_cold_backpressure_events[core_index].fetch_add_relaxed(1);
1094 self.per_group_cold_backpressure_events[group_index].fetch_add_relaxed(1);
1095 self.cold_backpressure_bytes
1096 .fetch_add_relaxed(incoming_bytes);
1097 }
1098
1099 pub(crate) fn record_read_watcher_added(&self, core_id: CoreId) {
1100 self.record_read_watchers_added(core_id, 1);
1101 }
1102
1103 pub(crate) fn record_read_watchers_added(&self, core_id: CoreId, count: usize) {
1104 self.per_core_live_read_waiters[usize::from(core_id.0)]
1105 .fetch_add_relaxed(u64::try_from(count).expect("watcher count fits u64"));
1106 }
1107
1108 pub(crate) fn record_read_watchers_removed(&self, core_id: CoreId, count: usize) {
1109 self.per_core_live_read_waiters[usize::from(core_id.0)]
1110 .fetch_sub_relaxed(u64::try_from(count).expect("watcher count fits u64"));
1111 }
1112
1113 pub(crate) fn record_live_read_backpressure(&self, core_id: CoreId) {
1114 self.per_core_live_read_backpressure_events[usize::from(core_id.0)].fetch_add_relaxed(1);
1115 }
1116}
1117
1118pub(crate) fn elapsed_ns(started_at: Instant) -> u64 {
1119 u64::try_from(started_at.elapsed().as_nanos()).unwrap_or(u64::MAX)
1120}
1121
1122pub(crate) fn append_batch_payload_bytes(request: &AppendBatchRequest) -> u64 {
1123 request
1124 .payloads
1125 .iter()
1126 .map(|payload| u64::try_from(payload.len()).expect("payload len fits u64"))
1127 .sum()
1128}
1129
1130pub(crate) fn record_cold_backpressure_error(
1131 metrics: &RuntimeMetricsInner,
1132 placement: ShardPlacement,
1133 incoming_bytes: u64,
1134 admission: ColdWriteAdmission,
1135 err: &GroupEngineError,
1136) {
1137 if !err.is_cold_backpressure() {
1138 return;
1139 }
1140 metrics.record_cold_backpressure(
1141 placement.core_id,
1142 placement.raft_group_id,
1143 incoming_bytes,
1144 admission.max_hot_bytes_per_group.unwrap_or(0),
1145 );
1146}
1147
1148pub(crate) fn is_stale_cold_flush_candidate_error(err: &RuntimeError) -> bool {
1149 match err.stream_error_code() {
1150 Some(StreamErrorCode::StreamGone | StreamErrorCode::StreamNotFound) => true,
1151 Some(StreamErrorCode::InvalidColdFlush) => err
1152 .stream_error_context()
1153 .iter()
1154 .any(|context| matches!(context, StreamErrorContext::StaleColdFlushCandidate)),
1155 _ => false,
1156 }
1157}
1158
1159pub(crate) async fn record_cold_hot_backlog(
1160 group: &mut Box<dyn GroupEngine>,
1161 metrics: &RuntimeMetricsInner,
1162 stream_id: BucketStreamId,
1163 placement: ShardPlacement,
1164) {
1165 if let Ok(backlog) = group.cold_hot_backlog(stream_id, placement).await {
1166 metrics.record_cold_hot_backlog(
1167 placement.raft_group_id,
1168 backlog.stream_hot_bytes,
1169 backlog.group_hot_bytes,
1170 );
1171 }
1172}
1173
1174#[derive(Debug)]
1175#[repr(align(128))]
1176pub(crate) struct PaddedAtomicU64 {
1177 value: AtomicU64,
1178}
1179
1180impl PaddedAtomicU64 {
1181 pub(crate) fn new(value: u64) -> Self {
1182 Self {
1183 value: AtomicU64::new(value),
1184 }
1185 }
1186
1187 pub(crate) fn load_relaxed(&self) -> u64 {
1188 self.value.load(Ordering::Relaxed)
1189 }
1190
1191 pub(crate) fn fetch_add_relaxed(&self, value: u64) -> u64 {
1192 self.value.fetch_add(value, Ordering::Relaxed)
1193 }
1194
1195 pub(crate) fn fetch_sub_relaxed(&self, value: u64) {
1196 self.value.fetch_sub(value, Ordering::Relaxed);
1197 }
1198
1199 pub(crate) fn fetch_sub_saturating_relaxed(&self, value: u64) {
1200 let mut current = self.value.load(Ordering::Relaxed);
1201 loop {
1202 let next = current.saturating_sub(value);
1203 match self.value.compare_exchange_weak(
1204 current,
1205 next,
1206 Ordering::Relaxed,
1207 Ordering::Relaxed,
1208 ) {
1209 Ok(_) => return,
1210 Err(observed) => current = observed,
1211 }
1212 }
1213 }
1214
1215 pub(crate) fn fetch_max_relaxed(&self, value: u64) {
1216 self.value.fetch_max(value, Ordering::Relaxed);
1217 }
1218
1219 pub(crate) fn store_relaxed(&self, value: u64) {
1220 self.value.store(value, Ordering::Relaxed);
1221 }
1222}