provide_telemetry/
backpressure.rs1use 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)] pub 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}