Skip to main content

kcode_pending_object_store/
lib.rs

1//! Durable, checksummed pending-object sidecar storage.
2//!
3//! Each object is installed through a synchronized same-directory temporary
4//! file and is read only after its format, metadata, length, and checksum have
5//! been verified.
6
7use std::{
8    collections::HashSet,
9    fs::{File, OpenOptions},
10    io::{Read, Seek, Write},
11    path::{Path, PathBuf},
12};
13
14use anyhow::{Context as _, ensure};
15use sha2::{Digest, Sha256};
16
17const MAGIC: &[u8] = b"KSPENDING01\n";
18const FIXED_HEADER_BYTES: usize = 2 + 4 + 4 + 8 + 32;
19
20#[derive(Clone, Debug, Eq, PartialEq)]
21pub struct StoredPendingObject {
22    pub file_name: String,
23    pub media_type: String,
24    pub bytes: Vec<u8>,
25}
26
27#[derive(Clone, Debug)]
28pub struct PendingObjectStore {
29    directory: PathBuf,
30    namespace: String,
31    format_id: String,
32}
33
34impl PendingObjectStore {
35    pub fn new(
36        directory: impl Into<PathBuf>,
37        namespace: impl Into<String>,
38        format_id: impl Into<String>,
39    ) -> anyhow::Result<Self> {
40        let namespace = namespace.into();
41        let format_id = format_id.into();
42        validate_namespace(&namespace)?;
43        ensure!(!format_id.is_empty(), "format ID cannot be empty");
44        u16::try_from(format_id.len()).context("format ID is too long")?;
45        Ok(Self {
46            directory: directory.into(),
47            namespace,
48            format_id,
49        })
50    }
51
52    pub fn install(
53        &self,
54        position: u64,
55        file_name: &str,
56        media_type: &str,
57        bytes: &[u8],
58    ) -> anyhow::Result<()> {
59        ensure!(
60            !file_name.trim().is_empty(),
61            "object filename cannot be empty"
62        );
63        ensure!(
64            !media_type.trim().is_empty(),
65            "object media type cannot be empty"
66        );
67
68        let format_id = self.format_id.as_bytes();
69        let format_id_len = u16::try_from(format_id.len()).context("format ID is too long")?;
70        let file_name_len =
71            u32::try_from(file_name.len()).context("object filename exceeds 4 GiB")?;
72        let media_type_len =
73            u32::try_from(media_type.len()).context("object media type exceeds 4 GiB")?;
74        let object_len = u64::try_from(bytes.len()).context("object exceeds addressable size")?;
75        let final_path = self.final_path(position);
76        let temp_path = self.temp_path(position);
77
78        ensure!(
79            !final_path.exists(),
80            "pending object file {} already exists",
81            final_path.display()
82        );
83        if temp_path.exists() {
84            std::fs::remove_file(&temp_path)
85                .with_context(|| format!("removing stale temporary {}", temp_path.display()))?;
86        }
87
88        let mut file = OpenOptions::new()
89            .create_new(true)
90            .write(true)
91            .open(&temp_path)
92            .with_context(|| format!("creating temporary {}", temp_path.display()))?;
93        file.write_all(MAGIC)?;
94        file.write_all(&format_id_len.to_le_bytes())?;
95        file.write_all(&file_name_len.to_le_bytes())?;
96        file.write_all(&media_type_len.to_le_bytes())?;
97        file.write_all(&object_len.to_le_bytes())?;
98        file.write_all(&Sha256::digest(bytes))?;
99        file.write_all(format_id)?;
100        file.write_all(file_name.as_bytes())?;
101        file.write_all(media_type.as_bytes())?;
102        file.write_all(bytes)?;
103        file.sync_all()
104            .with_context(|| format!("synchronizing temporary {}", temp_path.display()))?;
105        std::fs::rename(&temp_path, &final_path).with_context(|| {
106            format!(
107                "renaming {} to {}",
108                temp_path.display(),
109                final_path.display()
110            )
111        })?;
112        sync_directory(&self.directory)
113    }
114
115    pub fn read(&self, position: u64) -> anyhow::Result<StoredPendingObject> {
116        let path = self.final_path(position);
117        let mut file = File::open(&path)
118            .with_context(|| format!("opening pending object {}", path.display()))?;
119
120        let mut magic = vec![0_u8; MAGIC.len()];
121        file.read_exact(&mut magic)?;
122        ensure!(
123            magic == MAGIC,
124            "{} is not a pending-object file",
125            path.display()
126        );
127
128        let mut fixed = [0_u8; FIXED_HEADER_BYTES];
129        file.read_exact(&mut fixed)?;
130        let format_id_len = u16::from_le_bytes(fixed[0..2].try_into().unwrap()) as usize;
131        let file_name_len = u32::from_le_bytes(fixed[2..6].try_into().unwrap()) as usize;
132        let media_type_len = u32::from_le_bytes(fixed[6..10].try_into().unwrap()) as usize;
133        let object_len = u64::from_le_bytes(fixed[10..18].try_into().unwrap());
134        let checksum = &fixed[18..50];
135
136        let variable_len = format_id_len
137            .checked_add(file_name_len)
138            .and_then(|value| value.checked_add(media_type_len))
139            .context("pending-object header length overflow")?;
140        let expected_file_len = u64::try_from(MAGIC.len() + FIXED_HEADER_BYTES)
141            .context("pending-object fixed header does not fit u64")?
142            .checked_add(
143                u64::try_from(variable_len)
144                    .context("pending-object variable header does not fit u64")?,
145            )
146            .and_then(|value| value.checked_add(object_len))
147            .context("pending-object declared length overflow")?;
148        ensure!(
149            file.metadata()?.len() == expected_file_len,
150            "pending-object declared length differs from file length"
151        );
152
153        let mut variable = vec![0_u8; variable_len];
154        file.read_exact(&mut variable)?;
155        let format_id = std::str::from_utf8(&variable[..format_id_len])?;
156        ensure!(
157            format_id == self.format_id,
158            "unsupported pending-object format {format_id}"
159        );
160
161        let file_name_end = format_id_len + file_name_len;
162        let file_name = std::str::from_utf8(&variable[format_id_len..file_name_end])?.to_owned();
163        let media_type =
164            std::str::from_utf8(&variable[file_name_end..file_name_end + media_type_len])?
165                .to_owned();
166        ensure!(!file_name.trim().is_empty(), "object filename is empty");
167        ensure!(!media_type.trim().is_empty(), "object media type is empty");
168
169        let object_len =
170            usize::try_from(object_len).context("pending object does not fit memory")?;
171        let mut bytes = vec![0_u8; object_len];
172        file.read_exact(&mut bytes)?;
173        ensure!(
174            file.stream_position()? == expected_file_len,
175            "pending-object file has trailing bytes"
176        );
177        let actual_checksum = Sha256::digest(&bytes);
178        ensure!(
179            &actual_checksum[..] == checksum,
180            "pending-object checksum mismatch"
181        );
182
183        Ok(StoredPendingObject {
184            file_name,
185            media_type,
186            bytes,
187        })
188    }
189
190    pub fn verify_all(&self, referenced_positions: &[u64]) -> anyhow::Result<()> {
191        for &position in referenced_positions {
192            self.read(position)?;
193        }
194        Ok(())
195    }
196
197    pub fn reconcile(&self, referenced_positions: &[u64]) -> anyhow::Result<()> {
198        let referenced = referenced_positions.iter().copied().collect::<HashSet<_>>();
199        let mut removed = false;
200
201        if self.directory.exists() {
202            for entry in std::fs::read_dir(&self.directory)
203                .with_context(|| format!("reading directory {}", self.directory.display()))?
204            {
205                let entry = entry?;
206                let Some(name) = entry.file_name().to_str().map(str::to_owned) else {
207                    continue;
208                };
209                let Some((position, temporary)) = self.parse_file_name(&name) else {
210                    continue;
211                };
212                if temporary || !referenced.contains(&position) {
213                    std::fs::remove_file(entry.path())
214                        .with_context(|| format!("removing {}", entry.path().display()))?;
215                    removed = true;
216                }
217            }
218        }
219
220        if removed {
221            sync_directory(&self.directory)?;
222        }
223        self.verify_all(referenced_positions)
224    }
225
226    pub fn delete_all(&self) -> anyhow::Result<()> {
227        if self.directory.exists() {
228            for entry in std::fs::read_dir(&self.directory)
229                .with_context(|| format!("reading directory {}", self.directory.display()))?
230            {
231                let entry = entry?;
232                let Some(name) = entry.file_name().to_str().map(str::to_owned) else {
233                    continue;
234                };
235                if self.parse_file_name(&name).is_some() {
236                    std::fs::remove_file(entry.path())
237                        .with_context(|| format!("removing {}", entry.path().display()))?;
238                }
239            }
240        }
241        sync_directory(&self.directory)
242    }
243
244    fn final_path(&self, position: u64) -> PathBuf {
245        self.directory
246            .join(format!("{}-{position}.pending-object", self.namespace))
247    }
248
249    fn temp_path(&self, position: u64) -> PathBuf {
250        self.directory
251            .join(format!("{}-{position}.pending-object.tmp", self.namespace))
252    }
253
254    fn parse_file_name(&self, name: &str) -> Option<(u64, bool)> {
255        let tail = name.strip_prefix(&format!("{}-", self.namespace))?;
256        let (number, temporary) = if let Some(number) = tail.strip_suffix(".pending-object.tmp") {
257            (number, true)
258        } else {
259            (tail.strip_suffix(".pending-object")?, false)
260        };
261        if number.is_empty() || number.starts_with('0') && number != "0" {
262            return None;
263        }
264        let position = number.parse::<u64>().ok()?;
265        (position.to_string() == number).then_some((position, temporary))
266    }
267}
268
269fn validate_namespace(namespace: &str) -> anyhow::Result<()> {
270    ensure!(!namespace.is_empty(), "namespace cannot be empty");
271    ensure!(namespace.len() <= 255, "namespace exceeds 255 characters");
272    ensure!(
273        namespace
274            .bytes()
275            .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_')),
276        "namespace contains characters that are unsafe in filenames"
277    );
278    Ok(())
279}
280
281fn sync_directory(directory: &Path) -> anyhow::Result<()> {
282    File::open(directory)
283        .with_context(|| format!("opening directory {}", directory.display()))?
284        .sync_all()
285        .with_context(|| format!("synchronizing directory {}", directory.display()))
286}
287
288#[cfg(test)]
289mod tests {
290    use std::time::{SystemTime, UNIX_EPOCH};
291
292    use super::*;
293
294    struct TestDirectory(PathBuf);
295
296    impl TestDirectory {
297        fn new(label: &str) -> Self {
298            Self(std::env::temp_dir().join(format!(
299                "pending-object-store-{label}-{}-{}",
300                std::process::id(),
301                SystemTime::now()
302                    .duration_since(UNIX_EPOCH)
303                    .unwrap()
304                    .as_nanos()
305            )))
306        }
307    }
308
309    impl Drop for TestDirectory {
310        fn drop(&mut self) {
311            if self.0.exists() {
312                std::fs::remove_dir_all(&self.0).unwrap();
313            }
314        }
315    }
316
317    #[test]
318    fn construction_validates_configuration_without_touching_filesystem() {
319        let directory = TestDirectory::new("new");
320        let store = PendingObjectStore::new(&directory.0, "session_1", "0.2.1").unwrap();
321        assert!(!store.directory.exists());
322        assert!(PendingObjectStore::new(&directory.0, "", "0.2.1").is_err());
323        assert!(PendingObjectStore::new(&directory.0, "bad/name", "0.2.1").is_err());
324        assert!(PendingObjectStore::new(&directory.0, "a".repeat(256), "0.2.1").is_err());
325        assert!(PendingObjectStore::new(&directory.0, "valid", "").is_err());
326        assert!(PendingObjectStore::new(&directory.0, "valid", "x".repeat(65_536)).is_err());
327        assert!(!directory.0.exists());
328    }
329
330    #[test]
331    fn install_matches_the_session_log_format_and_read_verifies_it() {
332        let directory = TestDirectory::new("format");
333        std::fs::create_dir(&directory.0).unwrap();
334        let store = PendingObjectStore::new(&directory.0, "session-7", "0.2.1").unwrap();
335        store
336            .install(3, "notes.txt", "text/plain", b"durable bytes")
337            .unwrap();
338
339        let mut expected = Vec::new();
340        expected.extend_from_slice(MAGIC);
341        expected.extend_from_slice(&5_u16.to_le_bytes());
342        expected.extend_from_slice(&9_u32.to_le_bytes());
343        expected.extend_from_slice(&10_u32.to_le_bytes());
344        expected.extend_from_slice(&13_u64.to_le_bytes());
345        expected.extend_from_slice(&Sha256::digest(b"durable bytes"));
346        expected.extend_from_slice(b"0.2.1notes.txttext/plaindurable bytes");
347        assert_eq!(
348            std::fs::read(directory.0.join("session-7-3.pending-object")).unwrap(),
349            expected
350        );
351        assert!(!directory.0.join("session-7-3.pending-object.tmp").exists());
352        assert_eq!(
353            store.read(3).unwrap(),
354            StoredPendingObject {
355                file_name: "notes.txt".into(),
356                media_type: "text/plain".into(),
357                bytes: b"durable bytes".to_vec(),
358            }
359        );
360
361        assert!(
362            store
363                .install(3, "replacement", "text/plain", b"new")
364                .is_err()
365        );
366        assert_eq!(
367            std::fs::read(directory.0.join("session-7-3.pending-object")).unwrap(),
368            expected
369        );
370    }
371
372    #[test]
373    fn install_removes_only_its_exact_stale_temporary_file() {
374        let directory = TestDirectory::new("temporary");
375        std::fs::create_dir(&directory.0).unwrap();
376        let store = PendingObjectStore::new(&directory.0, "space", "format").unwrap();
377        let stale = directory.0.join("space-4.pending-object.tmp");
378        let other = directory.0.join("space-04.pending-object.tmp");
379        std::fs::write(&stale, b"stale").unwrap();
380        std::fs::write(&other, b"keep").unwrap();
381
382        store.install(4, "a", "b", b"c").unwrap();
383
384        assert!(!stale.exists());
385        assert!(other.exists());
386        assert!(
387            PendingObjectStore::new(directory.0.join("missing"), "space", "format")
388                .unwrap()
389                .install(0, "a", "b", b"c")
390                .is_err()
391        );
392    }
393
394    #[test]
395    fn reads_reject_wrong_format_length_checksum_and_empty_metadata() {
396        let directory = TestDirectory::new("invalid");
397        std::fs::create_dir(&directory.0).unwrap();
398        let store = PendingObjectStore::new(&directory.0, "object", "format-a").unwrap();
399
400        store.install(0, "file", "type", b"bytes").unwrap();
401        assert!(
402            PendingObjectStore::new(&directory.0, "object", "format-b")
403                .unwrap()
404                .read(0)
405                .is_err()
406        );
407
408        let path = directory.0.join("object-0.pending-object");
409        let valid = std::fs::read(&path).unwrap();
410        let mut corrupted = valid.clone();
411        *corrupted.last_mut().unwrap() ^= 0xff;
412        std::fs::write(&path, &corrupted).unwrap();
413        assert!(store.read(0).is_err());
414
415        let mut trailing = valid.clone();
416        trailing.push(0);
417        std::fs::write(&path, trailing).unwrap();
418        assert!(store.read(0).is_err());
419
420        let mut empty_name = valid;
421        let fixed_start = MAGIC.len();
422        empty_name[fixed_start + 2..fixed_start + 6].copy_from_slice(&0_u32.to_le_bytes());
423        std::fs::write(&path, empty_name).unwrap();
424        assert!(store.read(0).is_err());
425    }
426
427    #[test]
428    fn verification_does_not_clean_and_reconciliation_is_exact() {
429        let directory = TestDirectory::new("reconcile");
430        std::fs::create_dir(&directory.0).unwrap();
431        let store = PendingObjectStore::new(&directory.0, "abc", "format").unwrap();
432        store.install(1, "one", "type", b"one").unwrap();
433        store.install(2, "two", "type", b"two").unwrap();
434
435        let temporary = directory.0.join("abc-8.pending-object.tmp");
436        let leading_zero = directory.0.join("abc-08.pending-object.tmp");
437        let similar_prefix = directory.0.join("abc-other-9.pending-object");
438        let malformed = directory.0.join("abc-x.pending-object");
439        std::fs::write(&temporary, b"remove").unwrap();
440        std::fs::write(&leading_zero, b"keep").unwrap();
441        std::fs::write(&similar_prefix, b"keep").unwrap();
442        std::fs::write(&malformed, b"keep").unwrap();
443
444        store.verify_all(&[1]).unwrap();
445        assert!(temporary.exists());
446
447        store.reconcile(&[1]).unwrap();
448        assert!(directory.0.join("abc-1.pending-object").exists());
449        assert!(!directory.0.join("abc-2.pending-object").exists());
450        assert!(!temporary.exists());
451        assert!(leading_zero.exists());
452        assert!(similar_prefix.exists());
453        assert!(malformed.exists());
454        assert!(store.reconcile(&[99]).is_err());
455    }
456
457    #[test]
458    fn delete_all_removes_only_recognized_names() {
459        let directory = TestDirectory::new("delete");
460        std::fs::create_dir(&directory.0).unwrap();
461        let store = PendingObjectStore::new(&directory.0, "abc", "format").unwrap();
462        store.install(0, "zero", "type", b"zero").unwrap();
463
464        let temporary = directory.0.join("abc-1.pending-object.tmp");
465        let leading_zero = directory.0.join("abc-01.pending-object");
466        let similar_prefix = directory.0.join("abcd-2.pending-object");
467        std::fs::write(&temporary, b"remove").unwrap();
468        std::fs::write(&leading_zero, b"keep").unwrap();
469        std::fs::write(&similar_prefix, b"keep").unwrap();
470
471        store.delete_all().unwrap();
472
473        assert!(!directory.0.join("abc-0.pending-object").exists());
474        assert!(!temporary.exists());
475        assert!(leading_zero.exists());
476        assert!(similar_prefix.exists());
477    }
478}