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