alopex_chirps/buffer/
message_buffer.rs1use 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#[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}