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
27macro_rules! runtime_metrics {
48 (@munch
49 ctx { $ir:ident $cc:ident $gc:ident }
50 inner { $($inner:tt)* }
51 new { $($new:tt)* }
52 snap { $($snap:tt)* }
53 fields { $($fields:tt)* }
54 names { $($names:ident)* }
55 rest { sum $global:ident: core $core:ident, group $group:ident; $($rest:tt)* }
56 ) => {
57 runtime_metrics! {
58 @munch
59 ctx { $ir $cc $gc }
60 inner {
61 $($inner)*
62 pub(crate) $core: Vec<PaddedAtomicU64>,
63 pub(crate) $group: Vec<PaddedAtomicU64>,
64 }
65 new {
66 $($new)*
67 $core: zeroed_counters($cc),
68 $group: zeroed_counters($gc),
69 }
70 snap {
71 $($snap)*
72 let $core = load_counters(&$ir.$core);
73 let $global: u64 = $core.iter().sum();
74 let $group = load_counters(&$ir.$group);
75 }
76 fields {
77 $($fields)*
78 pub $global: u64,
79 pub $core: Vec<u64>,
80 pub $group: Vec<u64>,
81 }
82 names { $($names)* $global $core $group }
83 rest { $($rest)* }
84 }
85 };
86 (@munch
87 ctx { $ir:ident $cc:ident $gc:ident }
88 inner { $($inner:tt)* }
89 new { $($new:tt)* }
90 snap { $($snap:tt)* }
91 fields { $($fields:tt)* }
92 names { $($names:ident)* }
93 rest { sum $global:ident: core $core:ident; $($rest:tt)* }
94 ) => {
95 runtime_metrics! {
96 @munch
97 ctx { $ir $cc $gc }
98 inner {
99 $($inner)*
100 pub(crate) $core: Vec<PaddedAtomicU64>,
101 }
102 new {
103 $($new)*
104 $core: zeroed_counters($cc),
105 }
106 snap {
107 $($snap)*
108 let $core = load_counters(&$ir.$core);
109 let $global: u64 = $core.iter().sum();
110 }
111 fields {
112 $($fields)*
113 pub $global: u64,
114 pub $core: Vec<u64>,
115 }
116 names { $($names)* $global $core }
117 rest { $($rest)* }
118 }
119 };
120 (@munch
121 ctx { $ir:ident $cc:ident $gc:ident }
122 inner { $($inner:tt)* }
123 new { $($new:tt)* }
124 snap { $($snap:tt)* }
125 fields { $($fields:tt)* }
126 names { $($names:ident)* }
127 rest { sum $global:ident: group $group:ident; $($rest:tt)* }
128 ) => {
129 runtime_metrics! {
130 @munch
131 ctx { $ir $cc $gc }
132 inner {
133 $($inner)*
134 pub(crate) $group: Vec<PaddedAtomicU64>,
135 }
136 new {
137 $($new)*
138 $group: zeroed_counters($gc),
139 }
140 snap {
141 $($snap)*
142 let $group = load_counters(&$ir.$group);
143 let $global: u64 = $group.iter().sum();
144 }
145 fields {
146 $($fields)*
147 pub $global: u64,
148 pub $group: Vec<u64>,
149 }
150 names { $($names)* $global $group }
151 rest { $($rest)* }
152 }
153 };
154 (@munch
155 ctx { $ir:ident $cc:ident $gc:ident }
156 inner { $($inner:tt)* }
157 new { $($new:tt)* }
158 snap { $($snap:tt)* }
159 fields { $($fields:tt)* }
160 names { $($names:ident)* }
161 rest { max $global:ident: group $group:ident; $($rest:tt)* }
162 ) => {
163 runtime_metrics! {
164 @munch
165 ctx { $ir $cc $gc }
166 inner {
167 $($inner)*
168 pub(crate) $group: Vec<PaddedAtomicU64>,
169 }
170 new {
171 $($new)*
172 $group: zeroed_counters($gc),
173 }
174 snap {
175 $($snap)*
176 let $group = load_counters(&$ir.$group);
177 let $global = max_or_zero(&$group);
178 }
179 fields {
180 $($fields)*
181 pub $global: u64,
182 pub $group: Vec<u64>,
183 }
184 names { $($names)* $global $group }
185 rest { $($rest)* }
186 }
187 };
188 (@munch
189 ctx { $ir:ident $cc:ident $gc:ident }
190 inner { $($inner:tt)* }
191 new { $($new:tt)* }
192 snap { $($snap:tt)* }
193 fields { $($fields:tt)* }
194 names { $($names:ident)* }
195 rest { summax $sum:ident, $max:ident: group $group:ident; $($rest:tt)* }
196 ) => {
197 runtime_metrics! {
198 @munch
199 ctx { $ir $cc $gc }
200 inner {
201 $($inner)*
202 pub(crate) $group: Vec<PaddedAtomicU64>,
203 }
204 new {
205 $($new)*
206 $group: zeroed_counters($gc),
207 }
208 snap {
209 $($snap)*
210 let $group = load_counters(&$ir.$group);
211 let $sum: u64 = $group.iter().sum();
212 let $max = max_or_zero(&$group);
213 }
214 fields {
215 $($fields)*
216 pub $sum: u64,
217 pub $max: u64,
218 pub $group: Vec<u64>,
219 }
220 names { $($names)* $sum $max $group }
221 rest { $($rest)* }
222 }
223 };
224 (@munch
225 ctx { $ir:ident $cc:ident $gc:ident }
226 inner { $($inner:tt)* }
227 new { $($new:tt)* }
228 snap { $($snap:tt)* }
229 fields { $($fields:tt)* }
230 names { $($names:ident)* }
231 rest { counter $global:ident; $($rest:tt)* }
232 ) => {
233 runtime_metrics! {
234 @munch
235 ctx { $ir $cc $gc }
236 inner {
237 $($inner)*
238 pub(crate) $global: PaddedAtomicU64,
239 }
240 new {
241 $($new)*
242 $global: PaddedAtomicU64::new(0),
243 }
244 snap {
245 $($snap)*
246 let $global = $ir.$global.load_relaxed();
247 }
248 fields {
249 $($fields)*
250 pub $global: u64,
251 }
252 names { $($names)* $global }
253 rest { $($rest)* }
254 }
255 };
256 (@munch
257 ctx { $ir:ident $cc:ident $gc:ident }
258 inner { $($inner:tt)* }
259 new { $($new:tt)* }
260 snap { $($snap:tt)* }
261 fields { $($fields:tt)* }
262 names { $($names:ident)* }
263 rest { }
264 ) => {
265 #[derive(Debug)]
266 pub(crate) struct RuntimeMetricsInner {
267 $($inner)*
268 }
269
270 impl RuntimeMetricsInner {
271 pub(crate) fn new($cc: usize, $gc: usize) -> Self {
272 Self { $($new)* }
273 }
274 }
275
276 #[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
277 pub struct RuntimeMetricsSnapshot {
278 $($fields)*
279 }
280
281 impl RuntimeMetrics {
282 pub fn snapshot(&self) -> RuntimeMetricsSnapshot {
283 let $ir = &self.inner;
284 $($snap)*
285 RuntimeMetricsSnapshot { $($names),* }
286 }
287 }
288 };
289 ( $($manifest:tt)* ) => {
290 runtime_metrics! {
291 @munch
292 ctx { inner_counters core_count raft_group_count }
293 inner {}
294 new {}
295 snap {}
296 fields {}
297 names {}
298 rest { $($manifest)* }
299 }
300 };
301}
302
303fn zeroed_counters(len: usize) -> Vec<PaddedAtomicU64> {
304 (0..len).map(|_| PaddedAtomicU64::new(0)).collect()
305}
306
307fn load_counters(counters: &[PaddedAtomicU64]) -> Vec<u64> {
308 counters.iter().map(PaddedAtomicU64::load_relaxed).collect()
309}
310
311fn max_or_zero(values: &[u64]) -> u64 {
312 values.iter().copied().max().unwrap_or(0)
313}
314
315runtime_metrics! {
316 sum accepted_appends: core per_core_appends, group per_group_appends;
317 sum applied_mutations:
318 core per_core_applied_mutations, group per_group_applied_mutations;
319 sum mutation_apply_ns:
320 core per_core_mutation_apply_ns, group per_group_mutation_apply_ns;
321 sum group_lock_wait_ns:
322 core per_core_group_lock_wait_ns, group per_group_group_lock_wait_ns;
323 sum group_engine_exec_ns:
324 core per_core_group_engine_exec_ns, group per_group_group_engine_exec_ns;
325 sum group_mailbox_depth: group per_group_group_mailbox_depth;
326 max group_mailbox_max_depth: group per_group_group_mailbox_max_depth;
327 sum group_mailbox_full_events: group per_group_group_mailbox_full_events;
328 sum raft_write_many_batches:
329 core per_core_raft_write_many_batches, group per_group_raft_write_many_batches;
330 sum raft_write_many_commands:
331 core per_core_raft_write_many_commands, group per_group_raft_write_many_commands;
332 sum raft_write_many_logical_commands:
333 core per_core_raft_write_many_logical_commands,
334 group per_group_raft_write_many_logical_commands;
335 sum raft_write_many_responses:
336 core per_core_raft_write_many_responses, group per_group_raft_write_many_responses;
337 sum raft_write_many_submit_ns:
338 core per_core_raft_write_many_submit_ns, group per_group_raft_write_many_submit_ns;
339 sum raft_write_many_response_ns:
340 core per_core_raft_write_many_response_ns, group per_group_raft_write_many_response_ns;
341 sum raft_apply_entries: core per_core_raft_apply_entries, group per_group_raft_apply_entries;
342 sum raft_apply_ns: core per_core_raft_apply_ns, group per_group_raft_apply_ns;
343 sum raft_snapshot_builds: group per_group_raft_snapshot_builds;
344 sum raft_snapshot_build_ns: group per_group_raft_snapshot_build_ns;
345 summax raft_snapshot_body_bytes, raft_snapshot_body_bytes_max:
346 group per_group_raft_snapshot_body_bytes;
347 summax raft_snapshot_pointer_bytes, raft_snapshot_pointer_bytes_max:
348 group per_group_raft_snapshot_pointer_bytes;
349 summax raft_snapshot_streams, raft_snapshot_streams_max:
350 group per_group_raft_snapshot_streams;
351 sum raft_snapshot_external_uploads: group per_group_raft_snapshot_external_uploads;
352 sum raft_snapshot_inline_fallbacks: group per_group_raft_snapshot_inline_fallbacks;
353 sum live_read_waiters: core per_core_live_read_waiters;
354 sum live_read_backpressure_events: core per_core_live_read_backpressure_events;
355 sum routed_requests: core per_core_routed_requests;
356 sum mailbox_send_wait_ns: core per_core_mailbox_send_wait_ns;
357 sum mailbox_full_events: core per_core_mailbox_full_events;
358 sum wal_batches: core per_core_wal_batches, group per_group_wal_batches;
359 sum wal_records: core per_core_wal_records, group per_group_wal_records;
360 sum wal_write_ns: core per_core_wal_write_ns, group per_group_wal_write_ns;
361 sum wal_sync_ns: core per_core_wal_sync_ns, group per_group_wal_sync_ns;
362 counter cold_flush_uploads;
363 counter cold_flush_upload_bytes;
364 counter cold_flush_upload_ns;
365 counter cold_flush_publishes;
366 counter cold_flush_publish_bytes;
367 counter cold_flush_publish_ns;
368 counter cold_orphan_cleanup_attempts;
369 counter cold_orphan_cleanup_errors;
370 counter cold_orphan_bytes;
371 counter cold_gc_reclaimed;
372 counter cold_gc_errors;
373 counter cold_flush_write_errors;
374 sum cold_hot_bytes: group per_group_cold_hot_bytes;
375 max cold_hot_group_bytes_max: group per_group_cold_hot_bytes_max;
376 counter cold_hot_stream_bytes_max;
377 sum cold_backpressure_events:
378 core per_core_cold_backpressure_events, group per_group_cold_backpressure_events;
379 counter cold_backpressure_bytes;
380}
381
382#[derive(Debug, Clone, PartialEq, Eq)]
383pub struct RuntimeMailboxSnapshot {
384 pub depths: Vec<usize>,
385 pub capacities: Vec<usize>,
386}
387
388#[derive(Debug, Clone, Copy)]
389pub(crate) struct RaftWriteManySample {
390 pub(crate) command_count: u64,
391 pub(crate) logical_command_count: u64,
392 pub(crate) response_count: u64,
393 pub(crate) submit_ns: u64,
394 pub(crate) response_ns: u64,
395}
396
397#[derive(Debug, Clone, Copy)]
398pub(crate) struct RaftSnapshotBuildSample {
399 pub(crate) streams: u64,
400 pub(crate) body_bytes: u64,
401 pub(crate) pointer_bytes: u64,
402 pub(crate) build_ns: u64,
403 pub(crate) external_upload: bool,
404 pub(crate) inline_fallback: bool,
405}
406
407impl RuntimeMetricsInner {
408 pub(crate) fn record_routed_request(&self, core_id: CoreId, mailbox_send_wait_ns: u64) {
409 let index = usize::from(core_id.0);
410 self.per_core_routed_requests[index].fetch_add_relaxed(1);
411 self.per_core_mailbox_send_wait_ns[index].fetch_add_relaxed(mailbox_send_wait_ns);
412 }
413
414 pub(crate) fn record_mailbox_full(&self, core_id: CoreId) {
415 self.per_core_mailbox_full_events[usize::from(core_id.0)].fetch_add_relaxed(1);
416 }
417
418 pub(crate) fn record_append(&self, core_id: CoreId, group_id: RaftGroupId) {
419 self.record_append_batch(core_id, group_id, 1);
420 }
421
422 pub(crate) fn record_append_batch(&self, core_id: CoreId, group_id: RaftGroupId, count: u64) {
423 self.per_core_appends[usize::from(core_id.0)].fetch_add_relaxed(count);
424 self.per_group_appends[usize::try_from(group_id.0).expect("u32 fits usize")]
425 .fetch_add_relaxed(count);
426 }
427
428 pub(crate) fn record_applied_mutation(
429 &self,
430 core_id: CoreId,
431 group_id: RaftGroupId,
432 apply_ns: u64,
433 ) {
434 self.record_applied_mutation_batch(core_id, group_id, 1, apply_ns);
435 }
436
437 pub(crate) fn record_applied_mutation_batch(
438 &self,
439 core_id: CoreId,
440 group_id: RaftGroupId,
441 count: u64,
442 apply_ns: u64,
443 ) {
444 let core_index = usize::from(core_id.0);
445 let group_index = usize::try_from(group_id.0).expect("u32 fits usize");
446 self.per_core_applied_mutations[core_index].fetch_add_relaxed(count);
447 self.per_group_applied_mutations[group_index].fetch_add_relaxed(count);
448 self.per_core_mutation_apply_ns[core_index].fetch_add_relaxed(apply_ns);
449 self.per_group_mutation_apply_ns[group_index].fetch_add_relaxed(apply_ns);
450 }
451
452 pub(crate) fn record_group_engine_exec(
453 &self,
454 core_id: CoreId,
455 group_id: RaftGroupId,
456 exec_ns: u64,
457 ) {
458 let core_index = usize::from(core_id.0);
459 let group_index = usize::try_from(group_id.0).expect("u32 fits usize");
460 self.per_core_group_engine_exec_ns[core_index].fetch_add_relaxed(exec_ns);
461 self.per_group_group_engine_exec_ns[group_index].fetch_add_relaxed(exec_ns);
462 }
463
464 pub(crate) fn record_group_mailbox_enqueued(&self, group_id: RaftGroupId) {
465 let group_index = usize::try_from(group_id.0).expect("u32 fits usize");
466 let depth = self.per_group_group_mailbox_depth[group_index]
467 .fetch_add_relaxed(1)
468 .saturating_add(1);
469 self.per_group_group_mailbox_max_depth[group_index].fetch_max_relaxed(depth);
470 }
471
472 pub(crate) fn record_group_mailbox_dequeued(&self, group_id: RaftGroupId) {
473 let group_index = usize::try_from(group_id.0).expect("u32 fits usize");
474 self.per_group_group_mailbox_depth[group_index].fetch_sub_saturating_relaxed(1);
475 }
476
477 pub(crate) fn record_group_mailbox_full(&self, group_id: RaftGroupId) {
478 let group_index = usize::try_from(group_id.0).expect("u32 fits usize");
479 self.per_group_group_mailbox_full_events[group_index].fetch_add_relaxed(1);
480 }
481
482 pub(crate) fn record_raft_write_many(
483 &self,
484 core_id: CoreId,
485 group_id: RaftGroupId,
486 sample: RaftWriteManySample,
487 ) {
488 let core_index = usize::from(core_id.0);
489 let group_index = usize::try_from(group_id.0).expect("u32 fits usize");
490 self.per_core_raft_write_many_batches[core_index].fetch_add_relaxed(1);
491 self.per_group_raft_write_many_batches[group_index].fetch_add_relaxed(1);
492 self.per_core_raft_write_many_commands[core_index].fetch_add_relaxed(sample.command_count);
493 self.per_group_raft_write_many_commands[group_index]
494 .fetch_add_relaxed(sample.command_count);
495 self.per_core_raft_write_many_logical_commands[core_index]
496 .fetch_add_relaxed(sample.logical_command_count);
497 self.per_group_raft_write_many_logical_commands[group_index]
498 .fetch_add_relaxed(sample.logical_command_count);
499 self.per_core_raft_write_many_responses[core_index]
500 .fetch_add_relaxed(sample.response_count);
501 self.per_group_raft_write_many_responses[group_index]
502 .fetch_add_relaxed(sample.response_count);
503 self.per_core_raft_write_many_submit_ns[core_index].fetch_add_relaxed(sample.submit_ns);
504 self.per_group_raft_write_many_submit_ns[group_index].fetch_add_relaxed(sample.submit_ns);
505 self.per_core_raft_write_many_response_ns[core_index].fetch_add_relaxed(sample.response_ns);
506 self.per_group_raft_write_many_response_ns[group_index]
507 .fetch_add_relaxed(sample.response_ns);
508 }
509
510 pub(crate) fn record_raft_apply_batch(
511 &self,
512 core_id: CoreId,
513 group_id: RaftGroupId,
514 entry_count: u64,
515 apply_ns: u64,
516 ) {
517 let core_index = usize::from(core_id.0);
518 let group_index = usize::try_from(group_id.0).expect("u32 fits usize");
519 self.per_core_raft_apply_entries[core_index].fetch_add_relaxed(entry_count);
520 self.per_group_raft_apply_entries[group_index].fetch_add_relaxed(entry_count);
521 self.per_core_raft_apply_ns[core_index].fetch_add_relaxed(apply_ns);
522 self.per_group_raft_apply_ns[group_index].fetch_add_relaxed(apply_ns);
523 }
524
525 pub(crate) fn record_raft_snapshot_build(
526 &self,
527 group_id: RaftGroupId,
528 sample: RaftSnapshotBuildSample,
529 ) {
530 let group_index = usize::try_from(group_id.0).expect("u32 fits usize");
531 self.per_group_raft_snapshot_builds[group_index].fetch_add_relaxed(1);
532 self.per_group_raft_snapshot_build_ns[group_index].fetch_add_relaxed(sample.build_ns);
533 self.per_group_raft_snapshot_body_bytes[group_index].store_relaxed(sample.body_bytes);
534 self.per_group_raft_snapshot_pointer_bytes[group_index].store_relaxed(sample.pointer_bytes);
535 self.per_group_raft_snapshot_streams[group_index].store_relaxed(sample.streams);
536 if sample.external_upload {
537 self.per_group_raft_snapshot_external_uploads[group_index].fetch_add_relaxed(1);
538 }
539 if sample.inline_fallback {
540 self.per_group_raft_snapshot_inline_fallbacks[group_index].fetch_add_relaxed(1);
541 }
542 }
543
544 pub(crate) fn record_wal_batch(
545 &self,
546 core_id: CoreId,
547 group_id: RaftGroupId,
548 record_count: u64,
549 write_ns: u64,
550 sync_ns: u64,
551 ) {
552 let core_index = usize::from(core_id.0);
553 let group_index = usize::try_from(group_id.0).expect("u32 fits usize");
554 self.per_core_wal_batches[core_index].fetch_add_relaxed(1);
555 self.per_group_wal_batches[group_index].fetch_add_relaxed(1);
556 self.per_core_wal_records[core_index].fetch_add_relaxed(record_count);
557 self.per_group_wal_records[group_index].fetch_add_relaxed(record_count);
558 self.per_core_wal_write_ns[core_index].fetch_add_relaxed(write_ns);
559 self.per_group_wal_write_ns[group_index].fetch_add_relaxed(write_ns);
560 self.per_core_wal_sync_ns[core_index].fetch_add_relaxed(sync_ns);
561 self.per_group_wal_sync_ns[group_index].fetch_add_relaxed(sync_ns);
562 }
563
564 pub(crate) fn record_cold_upload(&self, bytes: u64, upload_ns: u64) {
565 self.cold_flush_uploads.fetch_add_relaxed(1);
566 self.cold_flush_upload_bytes.fetch_add_relaxed(bytes);
567 self.cold_flush_upload_ns.fetch_add_relaxed(upload_ns);
568 }
569
570 pub(crate) fn record_cold_publish(&self, bytes: u64, publish_ns: u64) {
571 self.cold_flush_publishes.fetch_add_relaxed(1);
572 self.cold_flush_publish_bytes.fetch_add_relaxed(bytes);
573 self.cold_flush_publish_ns.fetch_add_relaxed(publish_ns);
574 }
575
576 pub(crate) fn record_cold_gc_reclaimed(&self, entries: u64) {
577 self.cold_gc_reclaimed.fetch_add_relaxed(entries);
578 }
579
580 pub(crate) fn record_cold_flush_write_error(&self) {
581 self.cold_flush_write_errors.fetch_add_relaxed(1);
582 }
583
584 pub(crate) fn record_cold_gc_error(&self) {
585 self.cold_gc_errors.fetch_add_relaxed(1);
586 }
587
588 pub(crate) fn record_cold_hot_backlog(
589 &self,
590 group_id: RaftGroupId,
591 stream_hot_bytes: u64,
592 group_hot_bytes: u64,
593 ) {
594 let group_index = usize::try_from(group_id.0).expect("u32 fits usize");
595 self.per_group_cold_hot_bytes[group_index].store_relaxed(group_hot_bytes);
596 self.per_group_cold_hot_bytes_max[group_index].fetch_max_relaxed(group_hot_bytes);
597 self.cold_hot_stream_bytes_max
598 .fetch_max_relaxed(stream_hot_bytes);
599 }
600
601 pub(crate) fn record_cold_backpressure(
602 &self,
603 core_id: CoreId,
604 group_id: RaftGroupId,
605 incoming_bytes: u64,
606 _limit: u64,
607 ) {
608 let core_index = usize::from(core_id.0);
609 let group_index = usize::try_from(group_id.0).expect("u32 fits usize");
610 self.per_core_cold_backpressure_events[core_index].fetch_add_relaxed(1);
611 self.per_group_cold_backpressure_events[group_index].fetch_add_relaxed(1);
612 self.cold_backpressure_bytes
613 .fetch_add_relaxed(incoming_bytes);
614 }
615
616 pub(crate) fn record_read_watcher_added(&self, core_id: CoreId) {
617 self.record_read_watchers_added(core_id, 1);
618 }
619
620 pub(crate) fn record_read_watchers_added(&self, core_id: CoreId, count: usize) {
621 self.per_core_live_read_waiters[usize::from(core_id.0)]
622 .fetch_add_relaxed(u64::try_from(count).expect("watcher count fits u64"));
623 }
624
625 pub(crate) fn record_read_watchers_removed(&self, core_id: CoreId, count: usize) {
626 self.per_core_live_read_waiters[usize::from(core_id.0)]
627 .fetch_sub_relaxed(u64::try_from(count).expect("watcher count fits u64"));
628 }
629
630 pub(crate) fn record_live_read_backpressure(&self, core_id: CoreId) {
631 self.per_core_live_read_backpressure_events[usize::from(core_id.0)].fetch_add_relaxed(1);
632 }
633}
634
635pub(crate) fn elapsed_ns(started_at: Instant) -> u64 {
636 u64::try_from(started_at.elapsed().as_nanos()).unwrap_or(u64::MAX)
637}
638
639pub(crate) fn append_batch_payload_bytes(request: &AppendBatchRequest) -> u64 {
640 request
641 .payloads
642 .iter()
643 .map(|payload| u64::try_from(payload.len()).expect("payload len fits u64"))
644 .sum()
645}
646
647pub(crate) fn record_cold_backpressure_error(
648 metrics: &RuntimeMetricsInner,
649 placement: ShardPlacement,
650 incoming_bytes: u64,
651 admission: ColdWriteAdmission,
652 err: &GroupEngineError,
653) {
654 if !err.is_cold_backpressure() {
655 return;
656 }
657 metrics.record_cold_backpressure(
658 placement.core_id,
659 placement.raft_group_id,
660 incoming_bytes,
661 admission.max_hot_bytes_per_group.unwrap_or(0),
662 );
663}
664
665pub(crate) fn is_stale_cold_flush_candidate_error(err: &RuntimeError) -> bool {
666 match err.stream_error_code() {
667 Some(StreamErrorCode::StreamGone | StreamErrorCode::StreamNotFound) => true,
668 Some(StreamErrorCode::InvalidColdFlush) => err
669 .stream_error_context()
670 .iter()
671 .any(|context| matches!(context, StreamErrorContext::StaleColdFlushCandidate)),
672 _ => false,
673 }
674}
675
676pub(crate) async fn record_cold_hot_backlog(
677 group: &mut Box<dyn GroupEngine>,
678 metrics: &RuntimeMetricsInner,
679 stream_id: BucketStreamId,
680 placement: ShardPlacement,
681) {
682 if let Ok(backlog) = group.cold_hot_backlog(stream_id, placement).await {
683 metrics.record_cold_hot_backlog(
684 placement.raft_group_id,
685 backlog.stream_hot_bytes,
686 backlog.group_hot_bytes,
687 );
688 }
689}
690
691#[derive(Debug)]
692#[repr(align(128))]
693pub(crate) struct PaddedAtomicU64 {
694 value: AtomicU64,
695}
696
697impl PaddedAtomicU64 {
698 pub(crate) fn new(value: u64) -> Self {
699 Self {
700 value: AtomicU64::new(value),
701 }
702 }
703
704 pub(crate) fn load_relaxed(&self) -> u64 {
705 self.value.load(Ordering::Relaxed)
706 }
707
708 pub(crate) fn fetch_add_relaxed(&self, value: u64) -> u64 {
709 self.value.fetch_add(value, Ordering::Relaxed)
710 }
711
712 pub(crate) fn fetch_sub_relaxed(&self, value: u64) {
713 self.value.fetch_sub(value, Ordering::Relaxed);
714 }
715
716 pub(crate) fn fetch_sub_saturating_relaxed(&self, value: u64) {
717 let mut current = self.value.load(Ordering::Relaxed);
718 loop {
719 let next = current.saturating_sub(value);
720 match self.value.compare_exchange_weak(
721 current,
722 next,
723 Ordering::Relaxed,
724 Ordering::Relaxed,
725 ) {
726 Ok(_) => return,
727 Err(observed) => current = observed,
728 }
729 }
730 }
731
732 pub(crate) fn fetch_max_relaxed(&self, value: u64) {
733 self.value.fetch_max(value, Ordering::Relaxed);
734 }
735
736 pub(crate) fn store_relaxed(&self, value: u64) {
737 self.value.store(value, Ordering::Relaxed);
738 }
739}
740
741#[cfg(test)]
742mod metric_manifest_tests {
743 use std::sync::Arc;
744
745 use crate::metrics::RuntimeMetrics;
746 use crate::metrics::RuntimeMetricsInner;
747
748 const EXPECTED_SNAPSHOT_KEYS: [&str; 105] = [
752 "accepted_appends",
753 "per_core_appends",
754 "per_group_appends",
755 "applied_mutations",
756 "per_core_applied_mutations",
757 "per_group_applied_mutations",
758 "mutation_apply_ns",
759 "per_core_mutation_apply_ns",
760 "per_group_mutation_apply_ns",
761 "group_lock_wait_ns",
762 "per_core_group_lock_wait_ns",
763 "per_group_group_lock_wait_ns",
764 "group_engine_exec_ns",
765 "per_core_group_engine_exec_ns",
766 "per_group_group_engine_exec_ns",
767 "group_mailbox_depth",
768 "per_group_group_mailbox_depth",
769 "group_mailbox_max_depth",
770 "per_group_group_mailbox_max_depth",
771 "group_mailbox_full_events",
772 "per_group_group_mailbox_full_events",
773 "raft_write_many_batches",
774 "per_core_raft_write_many_batches",
775 "per_group_raft_write_many_batches",
776 "raft_write_many_commands",
777 "per_core_raft_write_many_commands",
778 "per_group_raft_write_many_commands",
779 "raft_write_many_logical_commands",
780 "per_core_raft_write_many_logical_commands",
781 "per_group_raft_write_many_logical_commands",
782 "raft_write_many_responses",
783 "per_core_raft_write_many_responses",
784 "per_group_raft_write_many_responses",
785 "raft_write_many_submit_ns",
786 "per_core_raft_write_many_submit_ns",
787 "per_group_raft_write_many_submit_ns",
788 "raft_write_many_response_ns",
789 "per_core_raft_write_many_response_ns",
790 "per_group_raft_write_many_response_ns",
791 "raft_apply_entries",
792 "per_core_raft_apply_entries",
793 "per_group_raft_apply_entries",
794 "raft_apply_ns",
795 "per_core_raft_apply_ns",
796 "per_group_raft_apply_ns",
797 "raft_snapshot_builds",
798 "per_group_raft_snapshot_builds",
799 "raft_snapshot_build_ns",
800 "per_group_raft_snapshot_build_ns",
801 "raft_snapshot_body_bytes",
802 "raft_snapshot_body_bytes_max",
803 "per_group_raft_snapshot_body_bytes",
804 "raft_snapshot_pointer_bytes",
805 "raft_snapshot_pointer_bytes_max",
806 "per_group_raft_snapshot_pointer_bytes",
807 "raft_snapshot_streams",
808 "raft_snapshot_streams_max",
809 "per_group_raft_snapshot_streams",
810 "raft_snapshot_external_uploads",
811 "per_group_raft_snapshot_external_uploads",
812 "raft_snapshot_inline_fallbacks",
813 "per_group_raft_snapshot_inline_fallbacks",
814 "live_read_waiters",
815 "per_core_live_read_waiters",
816 "live_read_backpressure_events",
817 "per_core_live_read_backpressure_events",
818 "routed_requests",
819 "per_core_routed_requests",
820 "mailbox_send_wait_ns",
821 "per_core_mailbox_send_wait_ns",
822 "mailbox_full_events",
823 "per_core_mailbox_full_events",
824 "wal_batches",
825 "per_core_wal_batches",
826 "per_group_wal_batches",
827 "wal_records",
828 "per_core_wal_records",
829 "per_group_wal_records",
830 "wal_write_ns",
831 "per_core_wal_write_ns",
832 "per_group_wal_write_ns",
833 "wal_sync_ns",
834 "per_core_wal_sync_ns",
835 "per_group_wal_sync_ns",
836 "cold_flush_uploads",
837 "cold_flush_upload_bytes",
838 "cold_flush_upload_ns",
839 "cold_flush_publishes",
840 "cold_flush_publish_bytes",
841 "cold_flush_publish_ns",
842 "cold_orphan_cleanup_attempts",
843 "cold_orphan_cleanup_errors",
844 "cold_orphan_bytes",
845 "cold_gc_reclaimed",
846 "cold_gc_errors",
847 "cold_flush_write_errors",
848 "cold_hot_bytes",
849 "per_group_cold_hot_bytes",
850 "cold_hot_group_bytes_max",
851 "per_group_cold_hot_bytes_max",
852 "cold_hot_stream_bytes_max",
853 "cold_backpressure_events",
854 "per_core_cold_backpressure_events",
855 "per_group_cold_backpressure_events",
856 "cold_backpressure_bytes",
857 ];
858
859 fn metrics_for_test() -> RuntimeMetrics {
860 RuntimeMetrics {
861 inner: Arc::new(RuntimeMetricsInner::new(2, 3)),
862 }
863 }
864
865 #[test]
866 fn snapshot_serializes_expected_field_names_in_order() {
867 let json =
868 serde_json::to_string(&metrics_for_test().snapshot()).expect("snapshot serializes");
869 assert_eq!(
872 json.matches("\":").count(),
873 EXPECTED_SNAPSHOT_KEYS.len(),
874 "unexpected number of serialized fields: {json}"
875 );
876 let mut last_position = None;
877 for name in EXPECTED_SNAPSHOT_KEYS {
878 let needle = format!("\"{name}\":");
879 let position = json
880 .find(&needle)
881 .unwrap_or_else(|| panic!("missing serialized key {name}"));
882 assert!(
883 last_position < Some(position),
884 "serialized key {name} out of declaration order"
885 );
886 last_position = Some(position);
887 }
888 }
889
890 #[test]
891 fn snapshot_vector_lengths_follow_metric_scope() {
892 let snapshot = metrics_for_test().snapshot();
893 assert_eq!(snapshot.per_core_appends.len(), 2);
894 assert_eq!(snapshot.per_group_appends.len(), 3);
895 assert_eq!(snapshot.per_core_routed_requests.len(), 2);
896 assert_eq!(snapshot.per_group_raft_snapshot_streams.len(), 3);
897 }
898
899 #[test]
900 fn snapshot_aggregates_sum_and_max_per_manifest() {
901 let metrics = metrics_for_test();
902 metrics.inner.per_core_appends[0].fetch_add_relaxed(3);
903 metrics.inner.per_core_appends[1].fetch_add_relaxed(4);
904 metrics.inner.per_group_group_mailbox_max_depth[1].fetch_max_relaxed(9);
905 metrics.inner.per_group_raft_snapshot_body_bytes[0].store_relaxed(5);
906 metrics.inner.per_group_raft_snapshot_body_bytes[2].store_relaxed(11);
907 let snapshot = metrics.snapshot();
908 assert_eq!(snapshot.accepted_appends, 7);
909 assert_eq!(snapshot.group_mailbox_max_depth, 9);
910 assert_eq!(snapshot.raft_snapshot_body_bytes, 16);
911 assert_eq!(snapshot.raft_snapshot_body_bytes_max, 11);
912 }
913}