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
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)] fn 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)] pub 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}