1use std::collections::BTreeMap;
7
8use serde::{Deserialize, Serialize};
9
10use crate::coroutine::Value;
11use crate::session::Edge;
12
13#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
15pub enum BufferMode {
16 Fifo,
18 LatestValue,
20}
21
22#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
24pub enum BackpressurePolicy {
25 Block,
27 Drop,
29 Error,
31 Resize {
33 max_capacity: usize,
35 },
36}
37
38#[derive(Debug, Clone, Serialize, Deserialize)]
40pub struct BufferConfig {
41 pub mode: BufferMode,
43 pub initial_capacity: usize,
45 pub policy: BackpressurePolicy,
47}
48
49impl Default for BufferConfig {
50 fn default() -> Self {
51 Self {
52 mode: BufferMode::Fifo,
53 initial_capacity: 64,
54 policy: BackpressurePolicy::Block,
55 }
56 }
57}
58
59#[derive(Debug, Clone, Serialize, Deserialize)]
61pub struct BoundedBuffer<T = Value> {
62 data: Vec<Option<T>>,
63 head: usize,
64 tail: usize,
65 capacity: usize,
66 count: usize,
67 epoch: usize,
68 mode: BufferMode,
69 policy: BackpressurePolicy,
70}
71
72#[derive(Debug)]
74pub enum EnqueueResult {
75 Ok,
77 WouldBlock,
79 Dropped,
81 Full,
83}
84
85#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
87pub struct SignedValue<V> {
88 pub payload: Value,
90 pub signature: V,
92}
93
94pub type SignedBuffer<V> = BoundedBuffer<SignedValue<V>>;
96
97pub type SignedBuffers<V> = BTreeMap<Edge, SignedBuffer<V>>;
99
100#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
102pub enum SignedDequeueError {
103 VerificationFailed,
105}
106
107#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
109pub enum SignedEnqueueResult {
110 Ok,
112 Blocked,
114 Dropped,
116 Error(String),
118}
119
120impl From<EnqueueResult> for SignedEnqueueResult {
121 fn from(value: EnqueueResult) -> Self {
122 match value {
123 EnqueueResult::Ok => Self::Ok,
124 EnqueueResult::WouldBlock => Self::Blocked,
125 EnqueueResult::Dropped => Self::Dropped,
126 EnqueueResult::Full => Self::Error("buffer full".to_string()),
127 }
128 }
129}
130
131pub fn signed_enqueue<V>(
133 buffers: &mut SignedBuffers<V>,
134 edge: Edge,
135 payload: Value,
136 signature: V,
137) -> SignedEnqueueResult {
138 let queue = buffers
139 .entry(edge)
140 .or_insert_with(|| BoundedBuffer::new(&BufferConfig::default()));
141 queue.enqueue(SignedValue { payload, signature }).into()
142}
143
144pub fn signed_dequeue_verified<V, F>(
153 buffers: &mut SignedBuffers<V>,
154 edge: &Edge,
155 verifier: F,
156) -> Result<Option<Value>, SignedDequeueError>
157where
158 F: Fn(&Value, &V) -> bool,
159{
160 let Some(queue) = buffers.get_mut(edge) else {
161 return Ok(None);
162 };
163 let Some(signed) = queue.dequeue() else {
164 return Ok(None);
165 };
166 if verifier(&signed.payload, &signed.signature) {
167 Ok(Some(signed.payload))
168 } else {
169 Err(SignedDequeueError::VerificationFailed)
170 }
171}
172
173impl<T> BoundedBuffer<T> {
174 #[must_use]
176 pub fn new(config: &BufferConfig) -> Self {
177 let capacity = config.initial_capacity.max(1);
178 let mut data = Vec::with_capacity(capacity);
179 data.resize_with(capacity, || None);
180 Self {
181 data,
182 head: 0,
183 tail: 0,
184 capacity,
185 count: 0,
186 epoch: 0,
187 mode: config.mode,
188 policy: config.policy,
189 }
190 }
191
192 pub fn enqueue(&mut self, v: T) -> EnqueueResult {
194 match self.mode {
195 BufferMode::LatestValue => {
196 self.data[0] = Some(v);
198 self.count = 1;
199 EnqueueResult::Ok
200 }
201 BufferMode::Fifo => {
202 if self.count >= self.capacity {
203 match self.policy {
204 BackpressurePolicy::Block => EnqueueResult::WouldBlock,
205 BackpressurePolicy::Drop => EnqueueResult::Dropped,
206 BackpressurePolicy::Error => EnqueueResult::Full,
207 BackpressurePolicy::Resize { max_capacity } => {
208 if self.capacity < max_capacity {
209 self.grow();
210 self.enqueue_fifo(v);
211 EnqueueResult::Ok
212 } else {
213 EnqueueResult::Full
214 }
215 }
216 }
217 } else {
218 self.enqueue_fifo(v);
219 EnqueueResult::Ok
220 }
221 }
222 }
223 }
224
225 pub fn dequeue(&mut self) -> Option<T> {
227 match self.mode {
228 BufferMode::LatestValue => {
229 if self.count > 0 {
230 self.count = 0;
231 self.data[0].take()
232 } else {
233 None
234 }
235 }
236 BufferMode::Fifo => {
237 if self.count == 0 {
238 return None;
239 }
240 let val = self.data[self.head].take();
241 self.head = (self.head + 1) % self.capacity;
242 self.count -= 1;
243 val
244 }
245 }
246 }
247
248 #[must_use]
250 pub fn is_empty(&self) -> bool {
251 self.count == 0
252 }
253
254 #[must_use]
256 pub fn is_full(&self) -> bool {
257 self.count >= self.capacity
258 }
259
260 #[must_use]
262 pub fn len(&self) -> usize {
263 self.count
264 }
265
266 #[must_use]
268 pub fn capacity(&self) -> usize {
269 self.capacity
270 }
271
272 #[must_use]
274 pub fn epoch(&self) -> usize {
275 self.epoch
276 }
277
278 pub fn advance_epoch(&mut self) {
280 self.epoch += 1;
281 }
282
283 fn enqueue_fifo(&mut self, v: T) {
284 self.data[self.tail] = Some(v);
285 self.tail = (self.tail + 1) % self.capacity;
286 self.count += 1;
287 }
288
289 fn grow(&mut self) {
290 let new_cap = self.capacity * 2;
291 let mut new_data = Vec::with_capacity(new_cap);
292 new_data.resize_with(new_cap, || None);
293
294 for (i, slot) in new_data.iter_mut().enumerate().take(self.count) {
296 let idx = (self.head + i) % self.capacity;
297 *slot = self.data[idx].take();
298 }
299
300 self.data = new_data;
301 self.head = 0;
302 self.tail = self.count;
303 self.capacity = new_cap;
304 }
305}
306
307#[cfg(test)]
308mod tests {
309 use super::*;
310
311 #[test]
312 fn test_fifo_basic() {
313 let mut buf = BoundedBuffer::new(&BufferConfig::default());
314 buf.enqueue(Value::Nat(1));
315 buf.enqueue(Value::Nat(2));
316 assert_eq!(buf.len(), 2);
317 assert_eq!(buf.dequeue(), Some(Value::Nat(1)));
318 assert_eq!(buf.dequeue(), Some(Value::Nat(2)));
319 assert!(buf.is_empty());
320 }
321
322 #[test]
323 fn test_fifo_wraparound() {
324 let config = BufferConfig {
325 initial_capacity: 3,
326 ..Default::default()
327 };
328 let mut buf = BoundedBuffer::new(&config);
329
330 buf.enqueue(Value::Nat(1));
331 buf.enqueue(Value::Nat(2));
332 buf.dequeue(); buf.enqueue(Value::Nat(3));
334 buf.enqueue(Value::Nat(4));
335
336 assert_eq!(buf.dequeue(), Some(Value::Nat(2)));
337 assert_eq!(buf.dequeue(), Some(Value::Nat(3)));
338 assert_eq!(buf.dequeue(), Some(Value::Nat(4)));
339 assert!(buf.is_empty());
340 }
341
342 #[test]
343 fn test_latest_value_overwrites() {
344 let config = BufferConfig {
345 mode: BufferMode::LatestValue,
346 initial_capacity: 1,
347 policy: BackpressurePolicy::Block,
348 };
349 let mut buf = BoundedBuffer::new(&config);
350
351 buf.enqueue(Value::Nat(1));
352 buf.enqueue(Value::Nat(2));
353 buf.enqueue(Value::Nat(3));
354
355 assert_eq!(buf.dequeue(), Some(Value::Nat(3)));
356 assert!(buf.is_empty());
357 }
358
359 #[test]
360 fn test_backpressure_block() {
361 let config = BufferConfig {
362 initial_capacity: 2,
363 policy: BackpressurePolicy::Block,
364 ..Default::default()
365 };
366 let mut buf = BoundedBuffer::new(&config);
367 buf.enqueue(Value::Nat(1));
368 buf.enqueue(Value::Nat(2));
369 assert!(matches!(
370 buf.enqueue(Value::Nat(3)),
371 EnqueueResult::WouldBlock
372 ));
373 }
374
375 #[test]
376 fn test_backpressure_resize() {
377 let config = BufferConfig {
378 initial_capacity: 2,
379 policy: BackpressurePolicy::Resize { max_capacity: 8 },
380 ..Default::default()
381 };
382 let mut buf = BoundedBuffer::new(&config);
383 buf.enqueue(Value::Nat(1));
384 buf.enqueue(Value::Nat(2));
385 assert!(matches!(buf.enqueue(Value::Nat(3)), EnqueueResult::Ok));
386 assert_eq!(buf.len(), 3);
387 }
388
389 #[test]
390 fn test_signed_buffer_alias_and_enqueue_result_mapping() {
391 let edge = Edge::new(7usize, "A", "B");
392 let signed = SignedValue {
393 payload: Value::Nat(9),
394 signature: vec![0_u8, 1_u8],
395 };
396 let mut buffers: SignedBuffers<Vec<u8>> = BTreeMap::new();
397 assert_eq!(
398 signed_enqueue(
399 &mut buffers,
400 edge.clone(),
401 signed.payload.clone(),
402 signed.signature.clone(),
403 ),
404 SignedEnqueueResult::Ok
405 );
406 assert_eq!(buffers.get(&edge).map(BoundedBuffer::len), Some(1));
407 assert_eq!(
408 buffers.get_mut(&edge).and_then(BoundedBuffer::dequeue),
409 Some(signed)
410 );
411
412 assert_eq!(
413 SignedEnqueueResult::from(EnqueueResult::Ok),
414 SignedEnqueueResult::Ok
415 );
416 assert_eq!(
417 SignedEnqueueResult::from(EnqueueResult::WouldBlock),
418 SignedEnqueueResult::Blocked
419 );
420 assert_eq!(
421 SignedEnqueueResult::from(EnqueueResult::Dropped),
422 SignedEnqueueResult::Dropped
423 );
424 assert!(matches!(
425 SignedEnqueueResult::from(EnqueueResult::Full),
426 SignedEnqueueResult::Error(_)
427 ));
428 }
429
430 #[test]
431 fn test_signed_dequeue_verified_success() {
432 let edge = Edge::new(11usize, "A", "B");
433 let mut buffers: SignedBuffers<Vec<u8>> = BTreeMap::new();
434 let _ = signed_enqueue(&mut buffers, edge.clone(), Value::Nat(5), vec![1, 2, 3]);
435 let out = signed_dequeue_verified(&mut buffers, &edge, |payload, signature| {
436 *payload == Value::Nat(5) && signature == &vec![1, 2, 3]
437 })
438 .expect("signature must verify");
439 assert_eq!(out, Some(Value::Nat(5)));
440 }
441
442 #[test]
443 fn test_signed_dequeue_verified_failure() {
444 let edge = Edge::new(12usize, "A", "B");
445 let mut buffers: SignedBuffers<Vec<u8>> = BTreeMap::new();
446 let _ = signed_enqueue(&mut buffers, edge.clone(), Value::Nat(5), vec![1, 2, 3]);
447 let result = signed_dequeue_verified(&mut buffers, &edge, |_payload, _signature| false);
448 assert_eq!(result, Err(SignedDequeueError::VerificationFailed));
449 }
450}