Skip to main content

pamoja_sync/
file.rs

1//! A crash-safe on-disk store-and-forward queue.
2
3use std::collections::VecDeque;
4use std::fs;
5use std::io::Write;
6use std::path::{Path, PathBuf};
7
8use pamoja_core::{Error, Result, Store};
9
10/// The filename extension for a stored record.
11const RECORD_EXTENSION: &str = "rec";
12
13/// A durable first-in first-out queue backed by one file per record.
14///
15/// Each [`append`](Store::append) writes a record to its own sequence-numbered
16/// file, flushes it to disk, and atomically renames it into place, so a power
17/// loss mid-write leaves the queue consistent: a partially written record is
18/// never visible. [`pop`](Store::pop) reads and deletes the oldest record. The
19/// directory itself is the durable state, so a store reopened after a crash
20/// resumes with every record that was fully written.
21///
22/// Delivery is at-least-once: if the process stops between reading a record and
23/// deleting it, the next [`open`](FileStore::open) returns that record again, so
24/// consumers must tolerate the occasional redelivery.
25///
26/// # Examples
27///
28/// ```no_run
29/// use pamoja_core::Store;
30/// use pamoja_sync::FileStore;
31///
32/// # async fn run() -> pamoja_core::Result<()> {
33/// let mut store = FileStore::open("/var/lib/pamoja/outbox")?;
34/// store.append(b"reading").await?;
35/// if let Some(record) = store.pop().await? {
36///     // forward `record` over a transport, then it is gone from the queue
37///     let _ = record;
38/// }
39/// # Ok(())
40/// # }
41/// ```
42pub struct FileStore {
43    dir: PathBuf,
44    pending: VecDeque<u64>,
45    next: u64,
46}
47
48impl FileStore {
49    /// Opens a store rooted at `dir`, creating the directory if needed.
50    ///
51    /// Records left in the directory by a previous run are adopted in sequence
52    /// order, so the queue resumes where it left off.
53    ///
54    /// # Arguments
55    ///
56    /// * `dir` - the directory that holds the queue's record files.
57    ///
58    /// # Returns
59    ///
60    /// A store ready to append and drain records.
61    ///
62    /// # Errors
63    ///
64    /// Returns [`Error::Io`](pamoja_core::Error::Io) if the directory cannot be
65    /// created or scanned.
66    pub fn open(dir: impl AsRef<Path>) -> Result<Self> {
67        let dir = dir.as_ref().to_path_buf();
68        fs::create_dir_all(&dir).map_err(io)?;
69
70        let mut sequences = Vec::new();
71        for entry in fs::read_dir(&dir).map_err(io)? {
72            let path = entry.map_err(io)?.path();
73            if path.extension().and_then(|ext| ext.to_str()) != Some(RECORD_EXTENSION) {
74                continue;
75            }
76            if let Some(sequence) = path
77                .file_stem()
78                .and_then(|stem| stem.to_str())
79                .and_then(|stem| stem.parse::<u64>().ok())
80            {
81                sequences.push(sequence);
82            }
83        }
84        sequences.sort_unstable();
85        let next = sequences.last().map_or(0, |last| last + 1);
86
87        Ok(Self {
88            dir,
89            pending: sequences.into(),
90            next,
91        })
92    }
93
94    /// Returns the path of the record file for a sequence number.
95    fn record_path(&self, sequence: u64) -> PathBuf {
96        self.dir.join(format!("{sequence:020}.{RECORD_EXTENSION}"))
97    }
98}
99
100impl Store for FileStore {
101    async fn append(&mut self, record: &[u8]) -> Result<()> {
102        let sequence = self.next;
103        let final_path = self.record_path(sequence);
104        let temp_path = final_path.with_extension(format!("{RECORD_EXTENSION}.tmp"));
105
106        let mut file = fs::File::create(&temp_path).map_err(io)?;
107        file.write_all(record).map_err(io)?;
108        file.sync_all().map_err(io)?;
109        drop(file);
110        fs::rename(&temp_path, &final_path).map_err(io)?;
111
112        self.next += 1;
113        self.pending.push_back(sequence);
114        Ok(())
115    }
116
117    async fn peek(&self) -> Result<Option<Vec<u8>>> {
118        let Some(&sequence) = self.pending.front() else {
119            return Ok(None);
120        };
121        let record = fs::read(self.record_path(sequence)).map_err(io)?;
122        Ok(Some(record))
123    }
124
125    async fn pop(&mut self) -> Result<Option<Vec<u8>>> {
126        let Some(sequence) = self.pending.pop_front() else {
127            return Ok(None);
128        };
129        let path = self.record_path(sequence);
130        let record = fs::read(&path).map_err(io)?;
131        fs::remove_file(&path).map_err(io)?;
132        Ok(Some(record))
133    }
134
135    async fn len(&self) -> Result<usize> {
136        Ok(self.pending.len())
137    }
138}
139
140/// Maps a filesystem error onto the shared I/O error.
141fn io(error: std::io::Error) -> Error {
142    Error::Io(error.to_string())
143}
144
145#[cfg(test)]
146mod tests {
147    use super::*;
148
149    #[tokio::test]
150    async fn drains_in_first_in_first_out_order() {
151        let dir = tempfile::tempdir().expect("tempdir");
152        let mut store = FileStore::open(dir.path()).expect("open");
153
154        store.append(b"a").await.expect("append");
155        store.append(b"b").await.expect("append");
156        assert_eq!(store.len().await.expect("len"), 2);
157
158        assert_eq!(store.pop().await.expect("pop"), Some(b"a".to_vec()));
159        assert_eq!(store.pop().await.expect("pop"), Some(b"b".to_vec()));
160        assert_eq!(store.pop().await.expect("pop"), None);
161    }
162
163    #[tokio::test]
164    async fn records_survive_reopening() {
165        let dir = tempfile::tempdir().expect("tempdir");
166        {
167            let mut store = FileStore::open(dir.path()).expect("open");
168            store.append(b"durable").await.expect("append");
169            store.append(b"records").await.expect("append");
170        }
171
172        // A fresh store over the same directory resumes the queue in order.
173        let mut reopened = FileStore::open(dir.path()).expect("reopen");
174        assert_eq!(reopened.len().await.expect("len"), 2);
175        assert_eq!(
176            reopened.pop().await.expect("pop"),
177            Some(b"durable".to_vec())
178        );
179        assert_eq!(
180            reopened.pop().await.expect("pop"),
181            Some(b"records".to_vec())
182        );
183    }
184
185    #[tokio::test]
186    async fn peek_reads_the_oldest_without_removing_it() {
187        let dir = tempfile::tempdir().expect("tempdir");
188        let mut store = FileStore::open(dir.path()).expect("open");
189        store.append(b"oldest").await.expect("append");
190        store.append(b"newer").await.expect("append");
191
192        assert_eq!(store.peek().await.expect("peek"), Some(b"oldest".to_vec()));
193        assert_eq!(store.len().await.expect("len"), 2);
194        assert_eq!(store.pop().await.expect("pop"), Some(b"oldest".to_vec()));
195    }
196
197    #[tokio::test]
198    async fn popping_removes_the_record_file() {
199        let dir = tempfile::tempdir().expect("tempdir");
200        let mut store = FileStore::open(dir.path()).expect("open");
201        store.append(b"once").await.expect("append");
202        let _ = store.pop().await.expect("pop");
203
204        let reopened = FileStore::open(dir.path()).expect("reopen");
205        assert!(reopened.is_empty().await.expect("is_empty"));
206    }
207}