Skip to main content

oxigdal_streaming/core/
backpressure.rs

1//! Backpressure handling for stream processing.
2
3use 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/// Strategy for handling backpressure.
11#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
12pub enum BackpressureStrategy {
13    /// Block until capacity is available
14    Block,
15
16    /// Drop oldest elements
17    DropOldest,
18
19    /// Drop newest elements
20    DropNewest,
21
22    /// Fail the operation
23    Fail,
24
25    /// Adaptive strategy based on load
26    Adaptive,
27}
28
29/// Configuration for backpressure management.
30#[derive(Debug, Clone, Serialize, Deserialize)]
31pub struct BackpressureConfig {
32    /// Strategy to use
33    pub strategy: BackpressureStrategy,
34
35    /// High watermark (percentage)
36    pub high_watermark: f64,
37
38    /// Low watermark (percentage)
39    pub low_watermark: f64,
40
41    /// Maximum latency threshold
42    pub max_latency: Duration,
43
44    /// Sample window for metrics
45    pub sample_window: Duration,
46
47    /// Enable adaptive backpressure
48    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/// Metrics for load monitoring.
65#[derive(Debug, Clone)]
66pub struct LoadMetrics {
67    /// Current buffer utilization (0.0 to 1.0)
68    pub buffer_utilization: f64,
69
70    /// Average latency
71    pub avg_latency: Duration,
72
73    /// Peak latency
74    pub peak_latency: Duration,
75
76    /// Throughput (elements per second)
77    pub throughput: f64,
78
79    /// Number of dropped elements
80    pub dropped_elements: u64,
81
82    /// Number of backpressure events
83    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
99/// Manages backpressure for a stream.
100pub 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    /// Create a new backpressure manager.
114    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    /// Check if backpressure should be applied.
129    pub async fn should_apply_backpressure(&self) -> bool {
130        let utilization = self.buffer_utilization();
131        utilization >= self.config.high_watermark
132    }
133
134    /// Check if backpressure can be released.
135    pub async fn can_release_backpressure(&self) -> bool {
136        let utilization = self.buffer_utilization();
137        utilization <= self.config.low_watermark
138    }
139
140    /// Handle a new element arrival.
141    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                    // Evict an oldest buffered element elsewhere and admit the
154                    // newly-arrived one: the caller is told to accept it.
155                    self.elements_dropped.fetch_add(1, Ordering::Relaxed);
156                    return Ok(true);
157                }
158                BackpressureStrategy::DropNewest => {
159                    // Discard the just-arrived (newest) element: the caller is
160                    // told to reject it rather than admit it.
161                    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    /// Handle element processing completion.
180    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        // Update metrics
185        self.update_metrics(latency).await;
186    }
187
188    /// Update metrics based on current state.
189    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            // Simple moving average for latency
210            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    /// Get current buffer utilization.
221    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    /// Get current metrics.
233    pub async fn metrics(&self) -> LoadMetrics {
234        self.metrics.read().await.clone()
235    }
236
237    /// Set buffer capacity.
238    pub fn set_capacity(&self, capacity: usize) {
239        self.buffer_capacity.store(capacity, Ordering::Relaxed);
240    }
241
242    /// Get buffer capacity.
243    pub fn capacity(&self) -> usize {
244        self.buffer_capacity.load(Ordering::Relaxed)
245    }
246
247    /// Get current buffer size.
248    pub fn size(&self) -> usize {
249        self.buffer_size.load(Ordering::Relaxed)
250    }
251
252    /// Reset metrics.
253    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    /// Adaptive capacity adjustment based on load.
263    pub async fn adjust_capacity_adaptive(&self) {
264        // Use real-time buffer utilization instead of cached metrics
265        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        // Fill to capacity.
338        for _ in 0..2 {
339            assert!(
340                manager
341                    .handle_element_arrival()
342                    .await
343                    .expect("arrival should succeed"),
344            );
345        }
346
347        // At capacity, DropOldest signals: admit the new element.
348        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); // metrics not yet sampled
354        // The internal drop counter must have advanced.
355        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        // Fill to capacity.
368        for _ in 0..2 {
369            assert!(
370                manager
371                    .handle_element_arrival()
372                    .await
373                    .expect("arrival should succeed"),
374            );
375        }
376
377        // At capacity, DropNewest signals: reject the incoming element.
378        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        // Buffer size must not have grown past capacity.
386        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}