Skip to main content

telltale_vm/
buffer.rs

1//! Bounded buffers with backpressure.
2//!
3//! Matches the Lean `BoundedBuffer` from `runtime.md ยง6`.
4//! Ring buffer with configurable mode and backpressure policy.
5
6use std::collections::BTreeMap;
7
8use serde::{Deserialize, Serialize};
9
10use crate::coroutine::Value;
11use crate::session::Edge;
12
13/// Buffer delivery mode.
14#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
15pub enum BufferMode {
16    /// First-in, first-out. All messages delivered in order.
17    Fifo,
18    /// Only the latest value is retained. Overwrites on enqueue.
19    LatestValue,
20}
21
22/// Policy when a buffer is full.
23#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
24pub enum BackpressurePolicy {
25    /// Block the sender until space is available.
26    Block,
27    /// Drop the message silently.
28    Drop,
29    /// Return an error to the sender.
30    Error,
31    /// Resize the buffer up to a maximum capacity.
32    Resize {
33        /// Upper bound on buffer capacity.
34        max_capacity: usize,
35    },
36}
37
38/// Configuration for a buffer.
39#[derive(Debug, Clone, Serialize, Deserialize)]
40pub struct BufferConfig {
41    /// Delivery mode.
42    pub mode: BufferMode,
43    /// Initial capacity.
44    pub initial_capacity: usize,
45    /// Backpressure policy.
46    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/// Bounded ring buffer for inter-role message passing.
60#[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/// Result of attempting to enqueue a value.
73#[derive(Debug)]
74pub enum EnqueueResult {
75    /// Value was enqueued successfully.
76    Ok,
77    /// Buffer is full; sender should block.
78    WouldBlock,
79    /// Value was dropped per policy.
80    Dropped,
81    /// Buffer full and policy is Error.
82    Full,
83}
84
85/// Signed value payload used by authenticated transport layers.
86#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
87pub struct SignedValue<V> {
88    /// The application payload.
89    pub payload: Value,
90    /// The signature/proof attached to the payload.
91    pub signature: V,
92}
93
94/// Signed FIFO for a single edge.
95pub type SignedBuffer<V> = BoundedBuffer<SignedValue<V>>;
96
97/// Signed buffers indexed by sid-qualified edge.
98pub type SignedBuffers<V> = BTreeMap<Edge, SignedBuffer<V>>;
99
100/// Signed dequeue failure.
101#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
102pub enum SignedDequeueError {
103    /// Signature verification failed for the dequeued payload.
104    VerificationFailed,
105}
106
107/// Result of signed enqueue attempts.
108#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
109pub enum SignedEnqueueResult {
110    /// Value was enqueued successfully.
111    Ok,
112    /// Buffer is full; sender should block.
113    Blocked,
114    /// Value was dropped per policy.
115    Dropped,
116    /// Error path for failed enqueue.
117    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
131/// Enqueue one signed payload into per-edge signed buffers.
132pub 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
144/// Dequeue and verify one signed payload.
145///
146/// The verifier is provided by the caller so this buffer module stays
147/// independent from any specific verification backend.
148///
149/// # Errors
150///
151/// Returns [`SignedDequeueError::VerificationFailed`] if the signature does not verify.
152pub 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    /// Create a new buffer with the given configuration.
175    #[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    /// Try to enqueue a value.
193    pub fn enqueue(&mut self, v: T) -> EnqueueResult {
194        match self.mode {
195            BufferMode::LatestValue => {
196                // Overwrite the single slot.
197                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    /// Dequeue a value, if available.
226    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    /// Whether the buffer is empty.
249    #[must_use]
250    pub fn is_empty(&self) -> bool {
251        self.count == 0
252    }
253
254    /// Whether the buffer is full (FIFO mode only).
255    #[must_use]
256    pub fn is_full(&self) -> bool {
257        self.count >= self.capacity
258    }
259
260    /// Current number of buffered values.
261    #[must_use]
262    pub fn len(&self) -> usize {
263        self.count
264    }
265
266    /// Buffer capacity.
267    #[must_use]
268    pub fn capacity(&self) -> usize {
269        self.capacity
270    }
271
272    /// Current epoch.
273    #[must_use]
274    pub fn epoch(&self) -> usize {
275        self.epoch
276    }
277
278    /// Advance the epoch (used during session draining).
279    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        // Copy existing elements in order.
295        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(); // remove 1
333        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}