Skip to main content

alopex_chirps/buffer/
message_buffer.rs

1use super::{BackpressureController, BackpressureLevel, PriorityQueue};
2use crate::MessageProfile;
3use thiserror::Error;
4
5#[derive(Debug, Error, PartialEq, Eq)]
6pub enum BufferError {
7    #[error("message buffer backpressure at {level:?}: requested {requested_bytes} bytes")]
8    Backpressure {
9        level: BackpressureLevel,
10        requested_bytes: usize,
11    },
12}
13
14#[derive(Debug, PartialEq, Eq)]
15pub struct BufferedMessage {
16    pub profile: MessageProfile,
17    pub payload: Vec<u8>,
18}
19
20/// Node-wide bounded receive buffer shared by all delivery profiles.
21#[derive(Debug)]
22pub struct MessageBuffer {
23    max_buffer_bytes: usize,
24    used_bytes: usize,
25    controller: BackpressureController,
26    queue: PriorityQueue<BufferedMessage>,
27    profile_bytes: [usize; 3],
28}
29
30impl MessageBuffer {
31    pub fn new(max_buffer_bytes: usize, warning_threshold: f32, limited_threshold: f32) -> Self {
32        Self {
33            max_buffer_bytes,
34            used_bytes: 0,
35            controller: BackpressureController::new(warning_threshold, limited_threshold),
36            queue: PriorityQueue::default(),
37            profile_bytes: [0; 3],
38        }
39    }
40
41    pub fn push(&mut self, profile: MessageProfile, payload: Vec<u8>) -> Result<(), BufferError> {
42        let requested_bytes = payload.len();
43        let projected = self.used_bytes.saturating_add(requested_bytes);
44        let level = if projected == self.max_buffer_bytes {
45            BackpressureLevel::Limited
46        } else {
47            self.controller.level(projected, self.max_buffer_bytes)
48        };
49        if !self.controller.allows(profile, level) {
50            return Err(BufferError::Backpressure {
51                level,
52                requested_bytes,
53            });
54        }
55
56        self.used_bytes = projected;
57        self.profile_bytes[profile_index(profile)] += requested_bytes;
58        self.queue
59            .push(profile, BufferedMessage { profile, payload });
60        Ok(())
61    }
62
63    pub fn pop(&mut self) -> Option<BufferedMessage> {
64        let (profile, message) = self.queue.pop()?;
65        let bytes = message.payload.len();
66        self.used_bytes = self.used_bytes.saturating_sub(bytes);
67        let profile_bytes = &mut self.profile_bytes[profile_index(profile)];
68        *profile_bytes = profile_bytes.saturating_sub(bytes);
69        Some(message)
70    }
71
72    pub fn max_buffer_bytes(&self) -> usize {
73        self.max_buffer_bytes
74    }
75
76    pub fn used_bytes(&self) -> usize {
77        self.used_bytes
78    }
79
80    pub fn backpressure_level(&self) -> BackpressureLevel {
81        self.controller
82            .level(self.used_bytes, self.max_buffer_bytes)
83    }
84
85    pub fn bytes_for(&self, profile: MessageProfile) -> usize {
86        self.profile_bytes[profile_index(profile)]
87    }
88
89    pub fn len(&self) -> usize {
90        self.queue.len()
91    }
92
93    pub fn is_empty(&self) -> bool {
94        self.queue.is_empty()
95    }
96}
97
98const fn profile_index(profile: MessageProfile) -> usize {
99    match profile {
100        MessageProfile::Control => 0,
101        MessageProfile::Durable => 1,
102        MessageProfile::Ephemeral => 2,
103    }
104}