Skip to main content

pamoja_sync/
memory.rs

1//! An in-memory store-and-forward queue.
2
3use std::collections::VecDeque;
4
5use pamoja_core::{Error, Result, Store};
6
7/// A fast in-memory first-in first-out queue.
8///
9/// Records live only for the lifetime of the process, which suits tests, the
10/// simulators, and the upper tier of a layered store. An optional capacity caps
11/// the number of buffered records and turns a full queue into an explicit
12/// backpressure signal rather than letting memory grow without bound.
13///
14/// # Examples
15///
16/// ```no_run
17/// use pamoja_core::Store;
18/// use pamoja_sync::MemoryStore;
19///
20/// # async fn run() -> pamoja_core::Result<()> {
21/// let mut store = MemoryStore::new();
22/// store.append(b"first").await?;
23/// store.append(b"second").await?;
24/// assert_eq!(store.len().await?, 2);
25/// assert_eq!(store.pop().await?, Some(b"first".to_vec()));
26/// # Ok(())
27/// # }
28/// ```
29#[derive(Clone, Debug, Default)]
30pub struct MemoryStore {
31    records: VecDeque<Vec<u8>>,
32    capacity: Option<usize>,
33}
34
35impl MemoryStore {
36    /// Creates an unbounded in-memory store.
37    ///
38    /// # Returns
39    ///
40    /// An empty store that grows to hold as many records as memory allows.
41    pub fn new() -> Self {
42        Self::default()
43    }
44
45    /// Creates a store that buffers at most `capacity` records.
46    ///
47    /// # Arguments
48    ///
49    /// * `capacity` - the maximum number of records to buffer; once reached,
50    ///   [`append`](Store::append) reports backpressure instead of growing.
51    ///
52    /// # Returns
53    ///
54    /// An empty capacity-bounded store.
55    pub fn with_capacity(capacity: usize) -> Self {
56        Self {
57            records: VecDeque::new(),
58            capacity: Some(capacity),
59        }
60    }
61}
62
63impl Store for MemoryStore {
64    async fn append(&mut self, record: &[u8]) -> Result<()> {
65        if let Some(capacity) = self.capacity {
66            if self.records.len() >= capacity {
67                return Err(Error::Io("store is at capacity".to_owned()));
68            }
69        }
70        self.records.push_back(record.to_vec());
71        Ok(())
72    }
73
74    async fn peek(&self) -> Result<Option<Vec<u8>>> {
75        Ok(self.records.front().cloned())
76    }
77
78    async fn pop(&mut self) -> Result<Option<Vec<u8>>> {
79        Ok(self.records.pop_front())
80    }
81
82    async fn len(&self) -> Result<usize> {
83        Ok(self.records.len())
84    }
85}
86
87#[cfg(test)]
88mod tests {
89    use super::*;
90
91    #[tokio::test]
92    async fn drains_in_first_in_first_out_order() {
93        let mut store = MemoryStore::new();
94        store.append(b"a").await.expect("append");
95        store.append(b"b").await.expect("append");
96        store.append(b"c").await.expect("append");
97
98        assert_eq!(store.len().await.expect("len"), 3);
99        assert_eq!(store.pop().await.expect("pop"), Some(b"a".to_vec()));
100        assert_eq!(store.pop().await.expect("pop"), Some(b"b".to_vec()));
101        assert_eq!(store.pop().await.expect("pop"), Some(b"c".to_vec()));
102        assert_eq!(store.pop().await.expect("pop"), None);
103    }
104
105    #[tokio::test]
106    async fn peek_returns_the_oldest_without_removing_it() {
107        let mut store = MemoryStore::new();
108        assert_eq!(store.peek().await.expect("peek"), None);
109
110        store.append(b"a").await.expect("append");
111        store.append(b"b").await.expect("append");
112        assert_eq!(store.peek().await.expect("peek"), Some(b"a".to_vec()));
113        assert_eq!(store.len().await.expect("len"), 2);
114        assert_eq!(store.pop().await.expect("pop"), Some(b"a".to_vec()));
115    }
116
117    #[tokio::test]
118    async fn empty_store_reports_empty() {
119        let store = MemoryStore::new();
120        assert_eq!(store.len().await.expect("len"), 0);
121        assert!(store.is_empty().await.expect("is_empty"));
122    }
123
124    #[tokio::test]
125    async fn bounded_store_rejects_when_full() {
126        let mut store = MemoryStore::with_capacity(1);
127        store.append(b"first").await.expect("append");
128
129        assert!(matches!(store.append(b"second").await, Err(Error::Io(_))));
130        assert_eq!(store.len().await.expect("len"), 1);
131
132        // Draining frees a slot so appends succeed again.
133        assert_eq!(store.pop().await.expect("pop"), Some(b"first".to_vec()));
134        store.append(b"second").await.expect("append after drain");
135        assert_eq!(store.pop().await.expect("pop"), Some(b"second".to_vec()));
136    }
137}