Skip to main content

ursula_runtime/
metrics.rs

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    /// Create an empty metric set sized for a runtime topology.
30    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    /// Return the metric handle passed to one group engine.
37    pub fn group_engine_metrics(&self) -> GroupEngineMetrics {
38        GroupEngineMetrics {
39            inner: Arc::clone(&self.inner),
40        }
41    }
42
43    /// Records one aggregate Raft-log pressure snapshot pass.
44    pub fn record_raft_snapshot_pressure(&self, groups: u64) {
45        self.inner.record_raft_snapshot_pressure(groups);
46    }
47}
48
49/// Declares every runtime metric once and expands the four sections that were
50/// previously hand-replicated per metric: the `RuntimeMetricsInner` counter
51/// fields, `RuntimeMetricsInner::new`, the `RuntimeMetrics::snapshot`
52/// collection logic, and the public serialized `RuntimeMetricsSnapshot`
53/// struct.
54///
55/// Manifest grammar (entries must be listed in serialized snapshot-field
56/// order; every counter and snapshot field name is spelled explicitly so
57/// serialized names stay grep-able and byte-stable):
58///
59/// - `sum GLOBAL: core PER_CORE, group PER_GROUP;` — per-core and per-group
60///   counters; the global value is the sum across cores.
61/// - `sum GLOBAL: core PER_CORE;` — per-core counters summed into the global.
62/// - `sum GLOBAL: group PER_GROUP;` — per-group counters summed into the
63///   global.
64/// - `max GLOBAL: group PER_GROUP;` — per-group values; the global is the
65///   maximum.
66/// - `summax SUM, MAX: group PER_GROUP;` — per-group values exposed both as a
67///   sum and as a maximum.
68/// - `counter GLOBAL;` — a single global counter.
69macro_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    // Current largest per-group backlog. Cold-health consumes this gauge and
425    // must be able to recover after a flush. The separate per-group `*_max`
426    // series remains the lifetime high-water mark for diagnostics.
427    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    /// The serialized field names of [`RuntimeMetricsSnapshot`] in declaration
898    /// order, captured from the pre-macro hand-written struct. Metrics
899    /// endpoints and `ursulactl` depend on these names staying byte-identical.
900    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        // Every value is a number or an array of numbers, so each `":`
1065        // occurrence in the output belongs to exactly one field key.
1066        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}