1use std::collections::VecDeque;
4use std::fs;
5use std::io::Write;
6use std::path::{Path, PathBuf};
7
8use pamoja_core::{Error, Result, Store};
9
10const RECORD_EXTENSION: &str = "rec";
12
13pub struct FileStore {
43 dir: PathBuf,
44 pending: VecDeque<u64>,
45 next: u64,
46}
47
48impl FileStore {
49 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 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
140fn 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 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}