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::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
27/// Declares every runtime metric once and expands the four sections that were
28/// previously hand-replicated per metric: the `RuntimeMetricsInner` counter
29/// fields, `RuntimeMetricsInner::new`, the `RuntimeMetrics::snapshot`
30/// collection logic, and the public serialized `RuntimeMetricsSnapshot`
31/// struct.
32///
33/// Manifest grammar (entries must be listed in serialized snapshot-field
34/// order; every counter and snapshot field name is spelled explicitly so
35/// serialized names stay grep-able and byte-stable):
36///
37/// - `sum GLOBAL: core PER_CORE, group PER_GROUP;` — per-core and per-group
38///   counters; the global value is the sum across cores.
39/// - `sum GLOBAL: core PER_CORE;` — per-core counters summed into the global.
40/// - `sum GLOBAL: group PER_GROUP;` — per-group counters summed into the
41///   global.
42/// - `max GLOBAL: group PER_GROUP;` — per-group values; the global is the
43///   maximum.
44/// - `summax SUM, MAX: group PER_GROUP;` — per-group values exposed both as a
45///   sum and as a maximum.
46/// - `counter GLOBAL;` — a single global counter.
47macro_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    /// The serialized field names of [`RuntimeMetricsSnapshot`] in declaration
749    /// order, captured from the pre-macro hand-written struct. Metrics
750    /// endpoints and `ursulactl` depend on these names staying byte-identical.
751    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        // Every value is a number or an array of numbers, so each `":`
870        // occurrence in the output belongs to exactly one field key.
871        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}