1use std::hint::spin_loop;
12use std::sync::Arc;
13use std::sync::atomic::AtomicU64;
14use std::sync::atomic::Ordering;
15use std::thread;
16
17use qubit_fast_cas::CasCell;
18#[cfg(feature = "serde")]
19use serde::Deserialize;
20#[cfg(feature = "serde")]
21use serde::Deserializer;
22#[cfg(feature = "serde")]
23use serde::Serialize;
24#[cfg(feature = "serde")]
25use serde::de::Error;
26
27use crate::MetricError;
28use crate::internal::OperationState;
29#[cfg(feature = "serde")]
30use crate::validation::validate_metrics;
31#[cfg(feature = "serde")]
32use crate::validation::validate_snapshot_counts;
33
34#[derive(Clone, Debug, Eq, PartialEq)]
36pub struct Metric {
37 pub(crate) id: Arc<str>,
39 pub(crate) name: Arc<str>,
41 pub(crate) total: Option<u64>,
43}
44
45impl Metric {
46 #[must_use]
51 pub fn new(id: &str, name: &str) -> Self {
52 Self {
53 id: Arc::from(id),
54 name: Arc::from(name),
55 total: None,
56 }
57 }
58
59 #[must_use]
64 pub const fn total(mut self, total: u64) -> Self {
65 self.total = Some(total);
66 self
67 }
68
69 #[must_use]
71 pub fn id(&self) -> &str {
72 &self.id
73 }
74
75 #[must_use]
77 pub fn name(&self) -> &str {
78 &self.name
79 }
80
81 #[must_use]
83 pub const fn configured_total(&self) -> Option<u64> {
84 self.total
85 }
86}
87
88#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
90pub struct MetricDelta {
91 started: u64,
93 unclassified: u64,
95 succeeded: u64,
97 failed: u64,
99 cancelled: u64,
101}
102
103impl MetricDelta {
104 #[must_use]
106 pub const fn new() -> Self {
107 Self {
108 started: 0,
109 unclassified: 0,
110 succeeded: 0,
111 failed: 0,
112 cancelled: 0,
113 }
114 }
115
116 #[must_use]
118 pub const fn started(mut self, count: u64) -> Self {
119 self.started = count;
120 self
121 }
122
123 #[must_use]
125 pub const fn unclassified(mut self, count: u64) -> Self {
126 self.unclassified = count;
127 self
128 }
129
130 #[must_use]
132 pub const fn succeeded(mut self, count: u64) -> Self {
133 self.succeeded = count;
134 self
135 }
136
137 #[must_use]
139 pub const fn failed(mut self, count: u64) -> Self {
140 self.failed = count;
141 self
142 }
143
144 #[must_use]
146 pub const fn cancelled(mut self, count: u64) -> Self {
147 self.cancelled = count;
148 self
149 }
150}
151
152#[derive(Clone)]
158pub struct MetricHandle {
159 inner: Arc<MetricInner>,
161 operation_state: Arc<OperationState>,
163}
164
165impl MetricHandle {
166 pub(crate) fn new(
168 metric: Metric,
169 operation_state: Arc<OperationState>,
170 ) -> Self {
171 Self {
172 inner: Arc::new(MetricInner::new(metric)),
173 operation_state,
174 }
175 }
176
177 #[must_use]
179 pub fn id(&self) -> &str {
180 self.inner.metric.id()
181 }
182
183 #[must_use]
185 pub fn name(&self) -> &str {
186 self.inner.metric.name()
187 }
188
189 pub fn start(&self, count: u64) -> Result<(), MetricError> {
196 self.apply_delta(MetricDelta::new().started(count))
197 }
198
199 pub fn complete(&self, count: u64) -> Result<(), MetricError> {
206 self.apply_delta(MetricDelta::new().unclassified(count))
207 }
208
209 pub fn succeed(&self, count: u64) -> Result<(), MetricError> {
216 self.apply_delta(MetricDelta::new().succeeded(count))
217 }
218
219 pub fn fail(&self, count: u64) -> Result<(), MetricError> {
226 self.apply_delta(MetricDelta::new().failed(count))
227 }
228
229 pub fn cancel(&self, count: u64) -> Result<(), MetricError> {
236 self.apply_delta(MetricDelta::new().cancelled(count))
237 }
238
239 pub fn apply_delta(&self, delta: MetricDelta) -> Result<(), MetricError> {
247 let metric_id = self.id();
248 let total = self.inner.metric.configured_total();
249 let _update_guard = self.operation_state.enter_update(metric_id)?;
250
251 self.inner.with_update(|counts| {
252 let mut next = *counts;
253 apply_delta_to_counts(&mut next, delta, metric_id)?;
254 next.validate(metric_id, total)?;
255 *counts = next;
256 Ok(())
257 })
258 }
259
260 #[must_use]
264 pub fn snapshot(&self) -> MetricSnapshot {
265 let counts = self.inner.snapshot_counts();
266 MetricSnapshot::from_counts(&self.inner.metric, counts)
267 }
268}
269
270struct MetricInner {
272 metric: Metric,
274 gate: CasCell,
276 active: AtomicU64,
278 completed_unclassified: AtomicU64,
280 succeeded: AtomicU64,
282 failed: AtomicU64,
284 cancelled: AtomicU64,
286}
287
288impl MetricInner {
289 fn new(metric: Metric) -> Self {
291 Self {
292 metric,
293 gate: CasCell::new(0),
294 active: AtomicU64::new(0),
295 completed_unclassified: AtomicU64::new(0),
296 succeeded: AtomicU64::new(0),
297 failed: AtomicU64::new(0),
298 cancelled: AtomicU64::new(0),
299 }
300 }
301
302 fn with_update<R, F>(&self, mut update: F) -> Result<R, MetricError>
304 where
305 F: FnMut(&mut MetricCounts) -> Result<R, MetricError>,
306 {
307 let mut attempts = 0;
308 loop {
309 let version = self.gate.load();
310 if version & 1 != 0 {
311 wait_for_contention(attempts);
312 attempts += 1;
313 continue;
314 }
315
316 match self.gate.compare_set(version, version.wrapping_add(1)) {
317 Ok(()) => {
318 let _guard = MetricGateGuard::new(
319 &self.gate,
320 version.wrapping_add(2),
321 );
322 let mut counts = self.read_counts();
323 let result = update(&mut counts);
324 if result.is_ok() {
325 self.write_counts(&counts);
326 }
327 return result;
328 }
329 Err(_) => {
330 wait_for_contention(attempts);
331 attempts += 1;
332 }
333 }
334 }
335 }
336
337 fn read_counts(&self) -> MetricCounts {
339 MetricCounts {
340 active: self.active.load(Ordering::Acquire),
341 completed_unclassified: self
342 .completed_unclassified
343 .load(Ordering::Acquire),
344 succeeded: self.succeeded.load(Ordering::Acquire),
345 failed: self.failed.load(Ordering::Acquire),
346 cancelled: self.cancelled.load(Ordering::Acquire),
347 }
348 }
349
350 fn write_counts(&self, counts: &MetricCounts) {
352 self.active.store(counts.active, Ordering::Release);
353 self.completed_unclassified
354 .store(counts.completed_unclassified, Ordering::Release);
355 self.succeeded.store(counts.succeeded, Ordering::Release);
356 self.failed.store(counts.failed, Ordering::Release);
357 self.cancelled.store(counts.cancelled, Ordering::Release);
358 }
359
360 fn snapshot_counts(&self) -> MetricCounts {
362 let mut attempts = 0;
363 loop {
364 let start = self.gate.load();
365 if start & 1 != 0 {
366 wait_for_contention(attempts);
367 attempts += 1;
368 continue;
369 }
370
371 let counts = self.read_counts();
372 if start == self.gate.load() {
373 return counts;
374 }
375
376 wait_for_contention(attempts);
377 attempts += 1;
378 }
379 }
380}
381
382#[derive(Clone, Copy)]
384struct MetricCounts {
385 active: u64,
387 completed_unclassified: u64,
389 succeeded: u64,
391 failed: u64,
393 cancelled: u64,
395}
396
397impl MetricCounts {
398 fn completed(self) -> Option<u64> {
402 self.completed_unclassified
403 .checked_add(self.succeeded)?
404 .checked_add(self.failed)?
405 .checked_add(self.cancelled)
406 }
407
408 fn occupied(self) -> Option<u64> {
410 self.completed()?.checked_add(self.active)
411 }
412
413 fn validate(
415 self,
416 metric_id: &str,
417 total: Option<u64>,
418 ) -> Result<(), MetricError> {
419 let occupied =
420 self.occupied().ok_or_else(|| MetricError::CountOverflow {
421 metric_id: metric_id.into(),
422 })?;
423 if let Some(total) = total
424 && occupied > total
425 {
426 return Err(MetricError::TotalExceeded {
427 metric_id: metric_id.into(),
428 total,
429 attempted: occupied,
430 });
431 }
432 Ok(())
433 }
434}
435
436fn apply_delta_to_counts(
438 counts: &mut MetricCounts,
439 delta: MetricDelta,
440 metric_id: &str,
441) -> Result<(), MetricError> {
442 let terminal_delta = delta
443 .unclassified
444 .checked_add(delta.succeeded)
445 .and_then(|value| value.checked_add(delta.failed))
446 .and_then(|value| value.checked_add(delta.cancelled))
447 .ok_or_else(|| MetricError::CountOverflow {
448 metric_id: metric_id.into(),
449 })?;
450 let available_active = counts
451 .active
452 .checked_add(delta.started)
453 .ok_or_else(|| MetricError::CountOverflow {
454 metric_id: metric_id.into(),
455 })?;
456 if terminal_delta > available_active {
457 return Err(MetricError::InsufficientActive {
458 metric_id: metric_id.into(),
459 requested: terminal_delta,
460 available: available_active,
461 });
462 }
463
464 counts.active = available_active - terminal_delta;
465 counts.completed_unclassified = counts
466 .completed_unclassified
467 .checked_add(delta.unclassified)
468 .ok_or_else(|| MetricError::CountOverflow {
469 metric_id: metric_id.into(),
470 })?;
471 counts.succeeded = counts
472 .succeeded
473 .checked_add(delta.succeeded)
474 .ok_or_else(|| MetricError::CountOverflow {
475 metric_id: metric_id.into(),
476 })?;
477 counts.failed =
478 counts.failed.checked_add(delta.failed).ok_or_else(|| {
479 MetricError::CountOverflow {
480 metric_id: metric_id.into(),
481 }
482 })?;
483 counts.cancelled = counts
484 .cancelled
485 .checked_add(delta.cancelled)
486 .ok_or_else(|| MetricError::CountOverflow {
487 metric_id: metric_id.into(),
488 })?;
489 Ok(())
490}
491
492struct MetricGateGuard<'gate> {
494 gate: &'gate CasCell,
495 next_version: u64,
496}
497
498impl<'gate> MetricGateGuard<'gate> {
499 fn new(gate: &'gate CasCell, next_version: u64) -> Self {
500 Self { gate, next_version }
501 }
502}
503
504impl Drop for MetricGateGuard<'_> {
505 fn drop(&mut self) {
506 self.gate.store(self.next_version);
507 }
508}
509
510#[cfg_attr(feature = "serde", derive(Serialize))]
512#[derive(Clone, Debug, Eq, PartialEq)]
513pub struct MetricSnapshot {
514 id: Arc<str>,
516 name: Arc<str>,
518 total: Option<u64>,
520 completed: u64,
522 active: u64,
524 succeeded: u64,
526 failed: u64,
528 cancelled: u64,
530}
531
532impl MetricSnapshot {
533 fn from_counts(metric: &Metric, counts: MetricCounts) -> Self {
535 Self {
536 id: Arc::clone(&metric.id),
537 name: Arc::clone(&metric.name),
538 total: metric.total,
539 completed: counts
540 .completed()
541 .expect("validated metric counts must fit in u64"),
542 active: counts.active,
543 succeeded: counts.succeeded,
544 failed: counts.failed,
545 cancelled: counts.cancelled,
546 }
547 }
548 #[must_use]
550 pub fn id(&self) -> &str {
551 &self.id
552 }
553 #[must_use]
555 pub fn name(&self) -> &str {
556 &self.name
557 }
558 #[must_use]
560 pub const fn total(&self) -> Option<u64> {
561 self.total
562 }
563 #[must_use]
565 pub const fn completed(&self) -> u64 {
566 self.completed
567 }
568 #[must_use]
570 pub const fn unclassified(&self) -> u64 {
571 let classified = self
572 .succeeded
573 .saturating_add(self.failed)
574 .saturating_add(self.cancelled);
575 self.completed.saturating_sub(classified)
576 }
577 #[must_use]
579 pub const fn active(&self) -> u64 {
580 self.active
581 }
582 #[must_use]
584 pub const fn succeeded(&self) -> u64 {
585 self.succeeded
586 }
587 #[must_use]
589 pub const fn failed(&self) -> u64 {
590 self.failed
591 }
592 #[must_use]
594 pub const fn cancelled(&self) -> u64 {
595 self.cancelled
596 }
597 #[must_use]
599 pub fn completion_fraction(&self) -> Option<f64> {
600 self.total
601 .filter(|total| *total > 0)
602 .map(|total| self.completed as f64 / total as f64)
603 }
604}
605
606#[cfg(feature = "serde")]
608#[derive(Deserialize)]
609struct MetricSnapshotWire {
610 id: Arc<str>,
612 name: Arc<str>,
614 total: Option<u64>,
616 completed: u64,
618 active: u64,
620 succeeded: u64,
622 failed: u64,
624 cancelled: u64,
626}
627
628#[cfg(feature = "serde")]
629impl<'de> Deserialize<'de> for MetricSnapshot {
630 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
632 where
633 D: Deserializer<'de>,
634 {
635 let wire = MetricSnapshotWire::deserialize(deserializer)?;
636 let snapshot = Self {
637 id: wire.id,
638 name: wire.name,
639 total: wire.total,
640 completed: wire.completed,
641 active: wire.active,
642 succeeded: wire.succeeded,
643 failed: wire.failed,
644 cancelled: wire.cancelled,
645 };
646 let metric = Metric {
647 id: Arc::clone(&snapshot.id),
648 name: Arc::clone(&snapshot.name),
649 total: snapshot.total,
650 };
651 validate_metrics(std::slice::from_ref(&metric))
652 .map_err(Error::custom)?;
653 validate_snapshot_counts(&snapshot).map_err(Error::custom)?;
654 Ok(snapshot)
655 }
656}
657
658#[inline]
660fn wait_for_contention(attempts: usize) {
661 if attempts > 0 && attempts.is_multiple_of(16) {
662 thread::yield_now();
663 } else {
664 spin_loop();
665 }
666}