1use crate::error::{Result, StreamingError};
4use serde::{Deserialize, Serialize};
5use std::sync::Arc;
6use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
7use std::time::{Duration, Instant};
8use tokio::sync::RwLock;
9
10#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
12pub enum BackpressureStrategy {
13 Block,
15
16 DropOldest,
18
19 DropNewest,
21
22 Fail,
24
25 Adaptive,
27}
28
29#[derive(Debug, Clone, Serialize, Deserialize)]
31pub struct BackpressureConfig {
32 pub strategy: BackpressureStrategy,
34
35 pub high_watermark: f64,
37
38 pub low_watermark: f64,
40
41 pub max_latency: Duration,
43
44 pub sample_window: Duration,
46
47 pub adaptive: bool,
49}
50
51impl Default for BackpressureConfig {
52 fn default() -> Self {
53 Self {
54 strategy: BackpressureStrategy::Block,
55 high_watermark: 0.8,
56 low_watermark: 0.2,
57 max_latency: Duration::from_secs(1),
58 sample_window: Duration::from_secs(10),
59 adaptive: true,
60 }
61 }
62}
63
64#[derive(Debug, Clone)]
66pub struct LoadMetrics {
67 pub buffer_utilization: f64,
69
70 pub avg_latency: Duration,
72
73 pub peak_latency: Duration,
75
76 pub throughput: f64,
78
79 pub dropped_elements: u64,
81
82 pub backpressure_events: u64,
84}
85
86impl Default for LoadMetrics {
87 fn default() -> Self {
88 Self {
89 buffer_utilization: 0.0,
90 avg_latency: Duration::ZERO,
91 peak_latency: Duration::ZERO,
92 throughput: 0.0,
93 dropped_elements: 0,
94 backpressure_events: 0,
95 }
96 }
97}
98
99pub struct BackpressureManager {
101 config: BackpressureConfig,
102 metrics: Arc<RwLock<LoadMetrics>>,
103 buffer_capacity: AtomicUsize,
104 buffer_size: AtomicUsize,
105 elements_processed: AtomicU64,
106 elements_dropped: AtomicU64,
107 backpressure_events: AtomicU64,
108 last_sample: Arc<RwLock<Instant>>,
109 sample_start: Instant,
110}
111
112impl BackpressureManager {
113 pub fn new(config: BackpressureConfig, buffer_capacity: usize) -> Self {
115 Self {
116 config,
117 metrics: Arc::new(RwLock::new(LoadMetrics::default())),
118 buffer_capacity: AtomicUsize::new(buffer_capacity),
119 buffer_size: AtomicUsize::new(0),
120 elements_processed: AtomicU64::new(0),
121 elements_dropped: AtomicU64::new(0),
122 backpressure_events: AtomicU64::new(0),
123 last_sample: Arc::new(RwLock::new(Instant::now())),
124 sample_start: Instant::now(),
125 }
126 }
127
128 pub async fn should_apply_backpressure(&self) -> bool {
130 let utilization = self.buffer_utilization();
131 utilization >= self.config.high_watermark
132 }
133
134 pub async fn can_release_backpressure(&self) -> bool {
136 let utilization = self.buffer_utilization();
137 utilization <= self.config.low_watermark
138 }
139
140 pub async fn handle_element_arrival(&self) -> Result<bool> {
142 let current_size = self.buffer_size.load(Ordering::Relaxed);
143 let capacity = self.buffer_capacity.load(Ordering::Relaxed);
144
145 if current_size >= capacity {
146 self.backpressure_events.fetch_add(1, Ordering::Relaxed);
147
148 match self.config.strategy {
149 BackpressureStrategy::Block => {
150 return Ok(false);
151 }
152 BackpressureStrategy::DropOldest => {
153 self.elements_dropped.fetch_add(1, Ordering::Relaxed);
156 return Ok(true);
157 }
158 BackpressureStrategy::DropNewest => {
159 self.elements_dropped.fetch_add(1, Ordering::Relaxed);
162 return Ok(false);
163 }
164 BackpressureStrategy::Fail => {
165 return Err(StreamingError::BufferFull);
166 }
167 BackpressureStrategy::Adaptive => {
168 if self.should_apply_backpressure().await {
169 return Ok(false);
170 }
171 }
172 }
173 }
174
175 self.buffer_size.fetch_add(1, Ordering::Relaxed);
176 Ok(true)
177 }
178
179 pub async fn handle_element_processed(&self, latency: Duration) {
181 self.buffer_size.fetch_sub(1, Ordering::Relaxed);
182 self.elements_processed.fetch_add(1, Ordering::Relaxed);
183
184 self.update_metrics(latency).await;
186 }
187
188 async fn update_metrics(&self, latency: Duration) {
190 let now = Instant::now();
191 let last_sample = *self.last_sample.read().await;
192
193 if now.duration_since(last_sample) >= self.config.sample_window {
194 let mut metrics = self.metrics.write().await;
195 let mut last = self.last_sample.write().await;
196
197 metrics.buffer_utilization = self.buffer_utilization();
198 metrics.dropped_elements = self.elements_dropped.load(Ordering::Relaxed);
199 metrics.backpressure_events = self.backpressure_events.load(Ordering::Relaxed);
200
201 let elapsed = now.duration_since(self.sample_start).as_secs_f64();
202 let processed = self.elements_processed.load(Ordering::Relaxed);
203 metrics.throughput = processed as f64 / elapsed;
204
205 if latency > metrics.peak_latency {
206 metrics.peak_latency = latency;
207 }
208
209 let alpha = 0.1;
211 let new_latency_secs = latency.as_secs_f64();
212 let old_latency_secs = metrics.avg_latency.as_secs_f64();
213 let avg_latency_secs = alpha * new_latency_secs + (1.0 - alpha) * old_latency_secs;
214 metrics.avg_latency = Duration::from_secs_f64(avg_latency_secs);
215
216 *last = now;
217 }
218 }
219
220 fn buffer_utilization(&self) -> f64 {
222 let size = self.buffer_size.load(Ordering::Relaxed);
223 let capacity = self.buffer_capacity.load(Ordering::Relaxed);
224
225 if capacity == 0 {
226 0.0
227 } else {
228 size as f64 / capacity as f64
229 }
230 }
231
232 pub async fn metrics(&self) -> LoadMetrics {
234 self.metrics.read().await.clone()
235 }
236
237 pub fn set_capacity(&self, capacity: usize) {
239 self.buffer_capacity.store(capacity, Ordering::Relaxed);
240 }
241
242 pub fn capacity(&self) -> usize {
244 self.buffer_capacity.load(Ordering::Relaxed)
245 }
246
247 pub fn size(&self) -> usize {
249 self.buffer_size.load(Ordering::Relaxed)
250 }
251
252 pub async fn reset_metrics(&self) {
254 let mut metrics = self.metrics.write().await;
255 *metrics = LoadMetrics::default();
256
257 self.elements_processed.store(0, Ordering::Relaxed);
258 self.elements_dropped.store(0, Ordering::Relaxed);
259 self.backpressure_events.store(0, Ordering::Relaxed);
260 }
261
262 pub async fn adjust_capacity_adaptive(&self) {
264 let utilization = self.buffer_utilization();
266 let metrics = self.metrics().await;
267
268 if utilization > self.config.high_watermark && metrics.avg_latency < self.config.max_latency
269 {
270 let current = self.buffer_capacity.load(Ordering::Relaxed);
271 let new_capacity = (current as f64 * 1.2) as usize;
272 self.buffer_capacity.store(new_capacity, Ordering::Relaxed);
273 } else if utilization < self.config.low_watermark {
274 let current = self.buffer_capacity.load(Ordering::Relaxed);
275 let new_capacity = ((current as f64 * 0.8) as usize).max(64);
276 self.buffer_capacity.store(new_capacity, Ordering::Relaxed);
277 }
278 }
279}
280
281#[cfg(test)]
282mod tests {
283 use super::*;
284
285 #[tokio::test]
286 async fn test_backpressure_manager_creation() {
287 let config = BackpressureConfig::default();
288 let manager = BackpressureManager::new(config, 1000);
289
290 assert_eq!(manager.capacity(), 1000);
291 assert_eq!(manager.size(), 0);
292 }
293
294 #[tokio::test]
295 async fn test_buffer_utilization() {
296 let config = BackpressureConfig::default();
297 let manager = BackpressureManager::new(config, 100);
298
299 assert_eq!(manager.buffer_utilization(), 0.0);
300
301 for _ in 0..50 {
302 manager
303 .handle_element_arrival()
304 .await
305 .expect("backpressure element arrival should succeed");
306 }
307
308 assert!((manager.buffer_utilization() - 0.5).abs() < 0.01);
309 }
310
311 #[tokio::test]
312 async fn test_backpressure_application() {
313 let config = BackpressureConfig {
314 high_watermark: 0.5,
315 ..Default::default()
316 };
317 let manager = BackpressureManager::new(config, 100);
318
319 for _ in 0..55 {
320 manager
321 .handle_element_arrival()
322 .await
323 .expect("backpressure element arrival should succeed");
324 }
325
326 assert!(manager.should_apply_backpressure().await);
327 }
328
329 #[tokio::test]
330 async fn test_drop_oldest_admits_new_element() {
331 let config = BackpressureConfig {
332 strategy: BackpressureStrategy::DropOldest,
333 ..Default::default()
334 };
335 let manager = BackpressureManager::new(config, 2);
336
337 for _ in 0..2 {
339 assert!(
340 manager
341 .handle_element_arrival()
342 .await
343 .expect("arrival should succeed"),
344 );
345 }
346
347 let admitted = manager
349 .handle_element_arrival()
350 .await
351 .expect("arrival should succeed");
352 assert!(admitted, "DropOldest must admit the incoming element");
353 assert_eq!(manager.metrics().await.dropped_elements, 0); let dropped = manager.elements_dropped.load(Ordering::Relaxed);
356 assert_eq!(dropped, 1);
357 }
358
359 #[tokio::test]
360 async fn test_drop_newest_rejects_new_element() {
361 let config = BackpressureConfig {
362 strategy: BackpressureStrategy::DropNewest,
363 ..Default::default()
364 };
365 let manager = BackpressureManager::new(config, 2);
366
367 for _ in 0..2 {
369 assert!(
370 manager
371 .handle_element_arrival()
372 .await
373 .expect("arrival should succeed"),
374 );
375 }
376
377 let admitted = manager
379 .handle_element_arrival()
380 .await
381 .expect("arrival should succeed");
382 assert!(!admitted, "DropNewest must reject the incoming element");
383 let dropped = manager.elements_dropped.load(Ordering::Relaxed);
384 assert_eq!(dropped, 1);
385 assert_eq!(manager.size(), 2);
387 }
388
389 #[tokio::test]
390 async fn test_adaptive_capacity_adjustment() {
391 let config = BackpressureConfig::default();
392 let manager = BackpressureManager::new(config, 100);
393
394 let initial_capacity = manager.capacity();
395
396 for _ in 0..95 {
397 manager
398 .handle_element_arrival()
399 .await
400 .expect("backpressure element arrival should succeed");
401 }
402
403 manager.adjust_capacity_adaptive().await;
404 let new_capacity = manager.capacity();
405
406 assert!(new_capacity > initial_capacity);
407 }
408}