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
20#[derive(Default)]
21enum QueueLimiter {
22    #[default]
23    Unlimited,
24    Bounded(Arc<Semaphore>),
25}
26
27#[derive(Default)]
28pub enum QueueTicket {
29    #[default]
30    Unlimited,
31    Bounded(OwnedSemaphorePermit),
32}
33
34struct QueueState {
35    policy: QueuePolicy,
36    logs: QueueLimiter,
37    traces: QueueLimiter,
38    metrics: QueueLimiter,
39}
40
41impl Default for QueueState {
42    fn default() -> Self {
43        Self::from_policy(QueuePolicy::default())
44    }
45}
46
47impl QueueState {
48    fn from_policy(policy: QueuePolicy) -> Self {
49        Self {
50            logs: limiter(policy.logs_maxsize),
51            traces: limiter(policy.traces_maxsize),
52            metrics: limiter(policy.metrics_maxsize),
53            policy,
54        }
55    }
56}
57
58fn limiter(size: usize) -> QueueLimiter {
59    if size == 0 {
60        QueueLimiter::Unlimited
61    } else {
62        QueueLimiter::Bounded(Arc::new(Semaphore::new(size)))
63    }
64}
65
66static QUEUES: OnceLock<Mutex<QueueState>> = OnceLock::new();
67
68#[cfg_attr(test, mutants::skip)] // Equivalent mutants only swap in Mutex::default().
69fn default_queue_state_mutex() -> Mutex<QueueState> {
70    Mutex::new(QueueState::default())
71}
72
73fn queues() -> &'static Mutex<QueueState> {
74    QUEUES.get_or_init(default_queue_state_mutex)
75}
76
77pub fn set_queue_policy(policy: QueuePolicy) {
78    *crate::_lock::lock(queues()) = QueueState::from_policy(policy);
79}
80
81pub fn get_queue_policy() -> QueuePolicy {
82    crate::_lock::lock(queues()).policy.clone()
83}
84
85pub fn try_acquire(signal: Signal) -> Option<QueueTicket> {
86    let guard = crate::_lock::lock(queues());
87    let limiter = match signal {
88        Signal::Logs => &guard.logs,
89        Signal::Traces => &guard.traces,
90        Signal::Metrics => &guard.metrics,
91    };
92
93    match limiter {
94        QueueLimiter::Unlimited => Some(QueueTicket::Unlimited),
95        QueueLimiter::Bounded(semaphore) => semaphore
96            .clone()
97            .try_acquire_owned()
98            .map(QueueTicket::Bounded)
99            .ok()
100            .or_else(|| {
101                increment_dropped(signal, 1);
102                None
103            }),
104    }
105}
106
107#[cfg_attr(test, mutants::skip)] // Equivalent mutant: moving `ticket` into this function drops it at scope end.
108pub fn release(ticket: QueueTicket) {
109    match ticket {
110        QueueTicket::Unlimited => {}
111        QueueTicket::Bounded(permit) => drop(permit),
112    }
113}
114
115pub fn _reset_backpressure_for_tests() {
116    *crate::_lock::lock(queues()) = QueueState::default();
117}
118
119#[cfg(test)]
120mod tests {
121    use super::*;
122    use crate::testing::acquire_test_state_lock;
123
124    #[test]
125    fn backpressure_test_reset_helper_restores_default_policy() {
126        let _guard = acquire_test_state_lock();
127        set_queue_policy(QueuePolicy {
128            logs_maxsize: 3,
129            traces_maxsize: 4,
130            metrics_maxsize: 5,
131        });
132        assert_eq!(
133            get_queue_policy(),
134            QueuePolicy {
135                logs_maxsize: 3,
136                traces_maxsize: 4,
137                metrics_maxsize: 5,
138            }
139        );
140
141        _reset_backpressure_for_tests();
142
143        assert_eq!(get_queue_policy(), QueuePolicy::default());
144    }
145
146    #[test]
147    fn backpressure_test_queues_returns_shared_singleton() {
148        let _guard = acquire_test_state_lock();
149        let first = queues() as *const Mutex<QueueState>;
150        let second = queues() as *const Mutex<QueueState>;
151
152        assert_eq!(first, second);
153    }
154
155    #[test]
156    fn backpressure_test_acquire_obeys_policy_limits() {
157        let _guard = acquire_test_state_lock();
158        _reset_backpressure_for_tests();
159
160        set_queue_policy(QueuePolicy {
161            logs_maxsize: 1,
162            traces_maxsize: 0,
163            metrics_maxsize: 0,
164        });
165
166        let first = try_acquire(Signal::Logs).expect("bounded acquire should succeed once");
167        let second = try_acquire(Signal::Logs);
168
169        assert!(matches!(first, QueueTicket::Bounded(_)));
170        assert!(second.is_none());
171        assert_eq!(get_queue_policy().logs_maxsize, 1);
172        release(first);
173
174        _reset_backpressure_for_tests();
175        set_queue_policy(QueuePolicy::default());
176        let unlimited = try_acquire(Signal::Logs).expect("unlimited acquire should succeed");
177        assert!(matches!(unlimited, QueueTicket::Unlimited));
178    }
179}