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