Skip to main content

provide_telemetry/
backpressure.rs

1// SPDX-FileCopyrightText: Copyright (C) 2026 provide.io llc
2// SPDX-License-Identifier: Apache-2.0
3// SPDX-Comment: Part of provide-telemetry.
4//
5
6use std::sync::{Arc, Mutex, OnceLock};
7
8use tokio::sync::{OwnedSemaphorePermit, Semaphore};
9
10use crate::health::increment_dropped;
11use crate::sampling::Signal;
12
13#[derive(Clone, Debug, Default, PartialEq, Eq)]
14pub struct QueuePolicy {
15    pub logs_maxsize: usize,
16    pub traces_maxsize: usize,
17    pub metrics_maxsize: usize,
18}
19
20enum QueueLimiter {
21    Unlimited,
22    Bounded(Arc<Semaphore>),
23}
24
25pub enum QueueTicket {
26    Unlimited,
27    Bounded(OwnedSemaphorePermit),
28}
29
30struct QueueState {
31    policy: QueuePolicy,
32    logs: QueueLimiter,
33    traces: QueueLimiter,
34    metrics: QueueLimiter,
35}
36
37impl Default for QueueState {
38    fn default() -> Self {
39        Self::from_policy(QueuePolicy::default())
40    }
41}
42
43impl QueueState {
44    fn from_policy(policy: QueuePolicy) -> Self {
45        Self {
46            logs: limiter(policy.logs_maxsize),
47            traces: limiter(policy.traces_maxsize),
48            metrics: limiter(policy.metrics_maxsize),
49            policy,
50        }
51    }
52}
53
54fn limiter(size: usize) -> QueueLimiter {
55    if size == 0 {
56        QueueLimiter::Unlimited
57    } else {
58        QueueLimiter::Bounded(Arc::new(Semaphore::new(size)))
59    }
60}
61
62static QUEUES: OnceLock<Mutex<QueueState>> = OnceLock::new();
63
64fn queues() -> &'static Mutex<QueueState> {
65    QUEUES.get_or_init(|| Mutex::new(QueueState::default()))
66}
67
68pub fn set_queue_policy(policy: QueuePolicy) {
69    *queues().lock().expect("queue lock poisoned") = QueueState::from_policy(policy);
70}
71
72pub fn get_queue_policy() -> QueuePolicy {
73    queues().lock().expect("queue lock poisoned").policy.clone()
74}
75
76pub fn try_acquire(signal: Signal) -> Option<QueueTicket> {
77    let guard = queues().lock().expect("queue lock poisoned");
78    let limiter = match signal {
79        Signal::Logs => &guard.logs,
80        Signal::Traces => &guard.traces,
81        Signal::Metrics => &guard.metrics,
82    };
83
84    match limiter {
85        QueueLimiter::Unlimited => Some(QueueTicket::Unlimited),
86        QueueLimiter::Bounded(semaphore) => semaphore
87            .clone()
88            .try_acquire_owned()
89            .map(QueueTicket::Bounded)
90            .ok()
91            .or_else(|| {
92                increment_dropped(signal, 1);
93                None
94            }),
95    }
96}
97
98#[cfg_attr(test, mutants::skip)] // Equivalent mutant: moving `ticket` into this function drops it at scope end.
99pub fn release(ticket: QueueTicket) {
100    match ticket {
101        QueueTicket::Unlimited => {}
102        QueueTicket::Bounded(permit) => drop(permit),
103    }
104}
105
106pub fn _reset_backpressure_for_tests() {
107    *queues().lock().expect("queue lock poisoned") = QueueState::default();
108}
109
110#[cfg(test)]
111mod tests {
112    use super::*;
113    use crate::testing::acquire_test_state_lock;
114
115    #[test]
116    fn backpressure_test_reset_helper_restores_default_policy() {
117        let _guard = acquire_test_state_lock();
118        set_queue_policy(QueuePolicy {
119            logs_maxsize: 3,
120            traces_maxsize: 4,
121            metrics_maxsize: 5,
122        });
123        assert_eq!(
124            get_queue_policy(),
125            QueuePolicy {
126                logs_maxsize: 3,
127                traces_maxsize: 4,
128                metrics_maxsize: 5,
129            }
130        );
131
132        _reset_backpressure_for_tests();
133
134        assert_eq!(get_queue_policy(), QueuePolicy::default());
135    }
136
137    #[test]
138    fn backpressure_test_queues_returns_shared_singleton() {
139        let _guard = acquire_test_state_lock();
140        let first = queues() as *const Mutex<QueueState>;
141        let second = queues() as *const Mutex<QueueState>;
142
143        assert_eq!(first, second);
144    }
145}