Skip to main content

pjson_rs_domain/value_objects/
backpressure.rs

1use serde::{Deserialize, Serialize};
2
3/// Signal indicating client's receive buffer state for backpressure control
4///
5/// Clients send backpressure signals to inform the server about their processing
6/// capacity. The server uses these signals to throttle or pause frame transmission.
7#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
8#[non_exhaustive]
9pub enum BackpressureSignal {
10    /// Client is ready for more data, no throttling needed
11    #[default]
12    Ok,
13
14    /// Client's buffer is filling up, server should slow down transmission
15    SlowDown,
16
17    /// Client's buffer is full, server must pause transmission
18    Pause,
19}
20
21impl BackpressureSignal {
22    /// Returns true if this signal indicates the server should pause
23    pub fn should_pause(&self) -> bool {
24        matches!(self, BackpressureSignal::Pause)
25    }
26
27    /// Returns true if this signal indicates the server should slow down
28    pub fn should_throttle(&self) -> bool {
29        matches!(
30            self,
31            BackpressureSignal::SlowDown | BackpressureSignal::Pause
32        )
33    }
34
35    /// Get suggested delay in milliseconds based on backpressure signal
36    pub fn suggested_delay_ms(&self) -> u64 {
37        match self {
38            BackpressureSignal::Ok => 0,
39            BackpressureSignal::SlowDown => 100,
40            BackpressureSignal::Pause => u64::MAX, // Indefinite pause until resumed
41        }
42    }
43}
44
45impl std::fmt::Display for BackpressureSignal {
46    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
47        match self {
48            BackpressureSignal::Ok => write!(f, "OK"),
49            BackpressureSignal::SlowDown => write!(f, "SLOW_DOWN"),
50            BackpressureSignal::Pause => write!(f, "PAUSE"),
51        }
52    }
53}
54
55#[cfg(test)]
56mod tests {
57    use super::*;
58
59    #[test]
60    fn test_backpressure_signal_default() {
61        let signal = BackpressureSignal::default();
62        assert_eq!(signal, BackpressureSignal::Ok);
63        assert!(!signal.should_pause());
64        assert!(!signal.should_throttle());
65    }
66
67    #[test]
68    fn test_backpressure_signal_should_pause() {
69        assert!(BackpressureSignal::Pause.should_pause());
70        assert!(!BackpressureSignal::SlowDown.should_pause());
71        assert!(!BackpressureSignal::Ok.should_pause());
72    }
73
74    #[test]
75    fn test_backpressure_signal_should_throttle() {
76        assert!(BackpressureSignal::Pause.should_throttle());
77        assert!(BackpressureSignal::SlowDown.should_throttle());
78        assert!(!BackpressureSignal::Ok.should_throttle());
79    }
80
81    #[test]
82    fn test_backpressure_signal_suggested_delay() {
83        assert_eq!(BackpressureSignal::Ok.suggested_delay_ms(), 0);
84        assert_eq!(BackpressureSignal::SlowDown.suggested_delay_ms(), 100);
85        assert_eq!(BackpressureSignal::Pause.suggested_delay_ms(), u64::MAX);
86    }
87
88    #[test]
89    fn test_backpressure_signal_display() {
90        assert_eq!(BackpressureSignal::Ok.to_string(), "OK");
91        assert_eq!(BackpressureSignal::SlowDown.to_string(), "SLOW_DOWN");
92        assert_eq!(BackpressureSignal::Pause.to_string(), "PAUSE");
93    }
94
95    #[test]
96    fn test_backpressure_signal_serialization() {
97        let signal = BackpressureSignal::SlowDown;
98        let json = serde_json::to_string(&signal).unwrap();
99        let deserialized: BackpressureSignal = serde_json::from_str(&json).unwrap();
100        assert_eq!(signal, deserialized);
101    }
102}