Skip to main content

kcode_k1_order_store/
lib.rs

1use kcode_k1_transaction::GENESIS_PARENT;
2pub use kcode_k1_transaction::SubsystemId;
3pub use kcode_k1_transaction_store::TxId;
4use std::collections::HashMap;
5use std::fs::{self, File, OpenOptions};
6use std::io::{self, Read, Seek, SeekFrom, Write};
7use std::path::Path;
8
9const RECORD_BYTES: usize = 32;
10const RECORD_BYTES_U64: u64 = 32;
11
12type Entries = Vec<(TxId, SubsystemId)>;
13type Indexes = HashMap<TxId, usize>;
14
15pub struct OrderStore {
16    file: File,
17    entries: Entries,
18    indexes: Indexes,
19}
20
21impl OrderStore {
22    pub fn create(path: &Path) -> Result<Self, String> {
23        let file = match OpenOptions::new()
24            .read(true)
25            .write(true)
26            .create_new(true)
27            .open(path)
28        {
29            Ok(file) => file,
30            Err(error)
31                if matches!(
32                    error.kind(),
33                    io::ErrorKind::AlreadyExists
34                        | io::ErrorKind::NotFound
35                        | io::ErrorKind::InvalidInput
36                ) =>
37            {
38                return Err(format!("cannot create ordering file: {error}"));
39            }
40            Err(error) => fatal_io("create", error),
41        };
42
43        Ok(Self {
44            file,
45            entries: Vec::new(),
46            indexes: HashMap::new(),
47        })
48    }
49
50    pub fn open(path: &Path) -> Result<Self, String> {
51        let metadata = fs::symlink_metadata(path)
52            .map_err(|error| format!("cannot inspect ordering file: {error}"))?;
53
54        if !metadata.file_type().is_file() {
55            return Err("ordering path is not a regular file".to_owned());
56        }
57
58        let mut file = OpenOptions::new()
59            .read(true)
60            .write(true)
61            .open(path)
62            .map_err(|error| format!("cannot open ordering file: {error}"))?;
63        let (entries, indexes) = reconstruct(&mut file)?;
64
65        Ok(Self {
66            file,
67            entries,
68            indexes,
69        })
70    }
71
72    pub fn entries(&self) -> &[(TxId, SubsystemId)] {
73        &self.entries
74    }
75
76    pub fn index_of(&self, id: TxId) -> Option<usize> {
77        self.indexes.get(&id).copied()
78    }
79
80    pub fn commit(
81        &mut self,
82        retained_len: usize,
83        id: TxId,
84        subsystem: SubsystemId,
85    ) -> Result<(), String> {
86        if retained_len > self.entries.len() {
87            return Err("retained length exceeds the canonical order".to_owned());
88        }
89
90        if id == GENESIS_PARENT {
91            return Err("cannot commit the genesis sentinel".to_owned());
92        }
93
94        if self.indexes.contains_key(&id) {
95            return Err("transaction ID is already canonical".to_owned());
96        }
97
98        let expected_len = byte_len(self.entries.len())?;
99        let actual_len = self
100            .file
101            .metadata()
102            .unwrap_or_else(|error| fatal_io("inspect-before-commit", error))
103            .len();
104
105        if actual_len != expected_len {
106            fatal_io(
107                "verify-before-commit",
108                io::Error::other(format!(
109                    "ordering file length changed from {expected_len} to {actual_len}"
110                )),
111            );
112        }
113
114        if retained_len < self.entries.len() {
115            let retained_bytes = byte_len(retained_len)?;
116
117            self.file
118                .set_len(retained_bytes)
119                .unwrap_or_else(|error| fatal_io("truncate", error));
120            self.file
121                .sync_data()
122                .unwrap_or_else(|error| fatal_io("sync-truncation", error));
123            self.file
124                .seek(SeekFrom::Start(retained_bytes))
125                .unwrap_or_else(|error| fatal_io("seek-replacement", error));
126            write_record(&mut self.file, id, subsystem, "write-replacement");
127            self.file
128                .sync_data()
129                .unwrap_or_else(|error| fatal_io("sync-replacement", error));
130
131            for (removed_id, _) in self.entries.drain(retained_len..) {
132                self.indexes.remove(&removed_id);
133            }
134        } else {
135            self.file
136                .seek(SeekFrom::Start(expected_len))
137                .unwrap_or_else(|error| fatal_io("seek-append", error));
138            write_record(&mut self.file, id, subsystem, "write-append");
139            self.file
140                .sync_data()
141                .unwrap_or_else(|error| fatal_io("sync-append", error));
142        }
143
144        let index = self.entries.len();
145        self.indexes.insert(id, index);
146        self.entries.push((id, subsystem));
147        Ok(())
148    }
149}
150
151fn reconstruct(file: &mut File) -> Result<(Entries, Indexes), String> {
152    file.seek(SeekFrom::Start(0))
153        .map_err(|error| format!("cannot seek ordering file: {error}"))?;
154
155    let mut bytes = Vec::new();
156    file.read_to_end(&mut bytes)
157        .map_err(|error| format!("cannot read ordering file: {error}"))?;
158
159    if bytes.len() % RECORD_BYTES != 0 {
160        return Err("ordering file length is not a multiple of 32".to_owned());
161    }
162
163    let count = bytes.len() / RECORD_BYTES;
164    let mut entries = Vec::with_capacity(count);
165    let mut indexes = HashMap::with_capacity(count);
166
167    for record in bytes.chunks_exact(RECORD_BYTES) {
168        let id = TxId::from_bytes(
169            record[..12]
170                .try_into()
171                .expect("transaction ID range has fixed length"),
172        );
173
174        if id == GENESIS_PARENT {
175            return Err("ordering file contains the genesis sentinel".to_owned());
176        }
177
178        let subsystem = SubsystemId::from_bytes(
179            record[12..]
180                .try_into()
181                .expect("subsystem range has fixed length"),
182        )
183        .map_err(|error| format!("ordering file contains an invalid subsystem: {error}"))?;
184
185        let index = entries.len();
186
187        if indexes.insert(id, index).is_some() {
188            return Err("ordering file contains a duplicate transaction ID".to_owned());
189        }
190
191        entries.push((id, subsystem));
192    }
193
194    Ok((entries, indexes))
195}
196
197fn byte_len(records: usize) -> Result<u64, String> {
198    let records =
199        u64::try_from(records).map_err(|_| "ordering file length exceeds u64".to_owned())?;
200
201    records
202        .checked_mul(RECORD_BYTES_U64)
203        .ok_or_else(|| "ordering file length exceeds u64".to_owned())
204}
205
206fn encoded_record(id: TxId, subsystem: SubsystemId) -> [u8; RECORD_BYTES] {
207    let mut record = [0_u8; RECORD_BYTES];
208    record[..12].copy_from_slice(id.as_bytes());
209    record[12..].copy_from_slice(subsystem.as_bytes());
210    record
211}
212
213fn write_record(file: &mut File, id: TxId, subsystem: SubsystemId, operation: &'static str) {
214    let record = encoded_record(id, subsystem);
215
216    match file.write(&record) {
217        Ok(RECORD_BYTES) => {}
218        Ok(written) => fatal_io(
219            operation,
220            io::Error::new(
221                io::ErrorKind::WriteZero,
222                format!("short record write: wrote {written} of {RECORD_BYTES} bytes"),
223            ),
224        ),
225        Err(error) => fatal_io(operation, error),
226    }
227}
228
229fn fatal_io(operation: &str, error: io::Error) -> ! {
230    eprintln!("kcode-k1-order-store fatal {operation}: {error}");
231    std::process::abort()
232}
233
234#[cfg(test)]
235mod tests {
236    use super::*;
237    use std::path::{Path, PathBuf};
238    use std::process::Command;
239    use std::sync::atomic::{AtomicU64, Ordering};
240    use std::time::{Duration, Instant};
241
242    static NEXT_ROOT: AtomicU64 = AtomicU64::new(0);
243    const ABORT_CASE: &str = "KCODE_K1_ORDER_STORE_ABORT_CASE";
244    const ABORT_PATH: &str = "KCODE_K1_ORDER_STORE_ABORT_PATH";
245
246    struct TempRoot {
247        path: PathBuf,
248    }
249
250    impl TempRoot {
251        fn new(label: &str) -> Self {
252            let sequence = NEXT_ROOT.fetch_add(1, Ordering::Relaxed);
253            let path = std::env::temp_dir().join(format!(
254                "kcode-k1-order-store-{}-{}-{}",
255                std::process::id(),
256                sequence,
257                label
258            ));
259            let _ = fs::remove_dir_all(&path);
260            fs::create_dir_all(&path).unwrap();
261            Self { path }
262        }
263
264        fn file(&self) -> PathBuf {
265            self.path.join("ordering.dat")
266        }
267    }
268
269    impl Drop for TempRoot {
270        fn drop(&mut self) {
271            let _ = fs::remove_dir_all(&self.path);
272        }
273    }
274
275    fn id(value: u64) -> TxId {
276        let mut bytes = [0_u8; 12];
277        bytes[..8].copy_from_slice(&value.to_le_bytes());
278        TxId::from_bytes(bytes)
279    }
280
281    fn subsystem(value: u8) -> SubsystemId {
282        SubsystemId::from_bytes([value; 20]).unwrap()
283    }
284
285    fn write_records(path: &Path, records: &[(TxId, SubsystemId)]) {
286        let mut bytes = Vec::with_capacity(records.len() * RECORD_BYTES);
287
288        for (id, subsystem) in records {
289            bytes.extend_from_slice(&encoded_record(*id, *subsystem));
290        }
291
292        fs::write(path, bytes).unwrap();
293    }
294
295    fn run_abort_case(case: &str, path: &Path, test_name: &str) -> std::process::Output {
296        Command::new(std::env::current_exe().unwrap())
297            .arg("--exact")
298            .arg(test_name)
299            .arg("--nocapture")
300            .env(ABORT_CASE, case)
301            .env(ABORT_PATH, path)
302            .output()
303            .unwrap()
304    }
305
306    #[test]
307    fn creates_appends_replaces_and_reopens_exact_bytes() {
308        let root = TempRoot::new("lifecycle");
309        let path = root.file();
310        let first = (id(1), subsystem(b'a'));
311        let second = (id(2), subsystem(b'b'));
312        let replacement = (id(3), subsystem(b'c'));
313
314        let mut store = OrderStore::create(&path).unwrap();
315        assert!(path.is_file());
316        assert!(store.entries().is_empty());
317
318        store.commit(0, first.0, first.1).unwrap();
319        store.commit(1, second.0, second.1).unwrap();
320
321        let mut expected = Vec::new();
322        expected.extend_from_slice(&encoded_record(first.0, first.1));
323        expected.extend_from_slice(&encoded_record(second.0, second.1));
324        assert_eq!(fs::read(&path).unwrap(), expected);
325        assert_eq!(store.entries(), &[first, second]);
326        assert_eq!(store.index_of(first.0), Some(0));
327        assert_eq!(store.index_of(second.0), Some(1));
328
329        store.commit(1, replacement.0, replacement.1).unwrap();
330        assert_eq!(store.entries(), &[first, replacement]);
331        assert_eq!(store.index_of(second.0), None);
332        assert_eq!(store.index_of(replacement.0), Some(1));
333
334        drop(store);
335        let reopened = OrderStore::open(&path).unwrap();
336        assert_eq!(reopened.entries(), &[first, replacement]);
337        assert_eq!(reopened.index_of(first.0), Some(0));
338        assert_eq!(reopened.index_of(replacement.0), Some(1));
339    }
340
341    #[test]
342    fn rejects_invalid_commit_arguments_without_mutation() {
343        let root = TempRoot::new("arguments");
344        let path = root.file();
345        let mut store = OrderStore::create(&path).unwrap();
346        let first = id(1);
347        let entry_subsystem = subsystem(b'a');
348
349        store.commit(0, first, entry_subsystem).unwrap();
350        let before = fs::read(&path).unwrap();
351
352        assert!(store.commit(2, id(2), entry_subsystem).is_err());
353        assert!(store.commit(1, GENESIS_PARENT, entry_subsystem).is_err());
354        assert!(store.commit(1, first, entry_subsystem).is_err());
355        assert_eq!(fs::read(&path).unwrap(), before);
356        assert_eq!(store.entries().len(), 1);
357    }
358
359    #[test]
360    fn rejects_malformed_and_invalid_files() {
361        let malformed = TempRoot::new("malformed");
362        fs::write(malformed.file(), [0_u8; RECORD_BYTES - 1]).unwrap();
363        assert!(OrderStore::open(&malformed.file()).is_err());
364
365        let duplicate = TempRoot::new("duplicate");
366        write_records(
367            &duplicate.file(),
368            &[(id(1), subsystem(b'a')), (id(1), subsystem(b'b'))],
369        );
370        assert!(OrderStore::open(&duplicate.file()).is_err());
371
372        let genesis = TempRoot::new("genesis");
373        write_records(&genesis.file(), &[(GENESIS_PARENT, subsystem(b'a'))]);
374        assert!(OrderStore::open(&genesis.file()).is_err());
375
376        let invalid_subsystem = TempRoot::new("invalid-subsystem");
377        let mut record = [0_u8; RECORD_BYTES];
378        record[..12].copy_from_slice(id(2).as_bytes());
379        record[12..].fill(0xff);
380        fs::write(invalid_subsystem.file(), record).unwrap();
381        assert!(OrderStore::open(&invalid_subsystem.file()).is_err());
382
383        let wrong_type = TempRoot::new("wrong-type");
384        assert!(OrderStore::open(&wrong_type.path).is_err());
385    }
386
387    #[test]
388    fn create_requires_an_absent_path() {
389        let root = TempRoot::new("create");
390        let path = root.file();
391
392        drop(OrderStore::create(&path).unwrap());
393        assert!(OrderStore::create(&path).is_err());
394    }
395
396    #[test]
397    fn external_length_change_aborts() {
398        let root = TempRoot::new("external-abort");
399        let path = root.file();
400
401        if std::env::var(ABORT_CASE).as_deref() == Ok("external") {
402            let child_path = PathBuf::from(std::env::var(ABORT_PATH).unwrap());
403            let mut store = OrderStore::create(&child_path).unwrap();
404            OpenOptions::new()
405                .append(true)
406                .open(&child_path)
407                .unwrap()
408                .write_all(&[0])
409                .unwrap();
410            store.commit(0, id(1), subsystem(b'a')).unwrap();
411            unreachable!();
412        }
413
414        let output = run_abort_case("external", &path, "tests::external_length_change_aborts");
415        assert!(!output.status.success());
416
417        let stderr = String::from_utf8_lossy(&output.stderr);
418        assert_eq!(
419            stderr
420                .lines()
421                .filter(|line| line.contains("kcode-k1-order-store fatal"))
422                .count(),
423            1
424        );
425        assert!(stderr.contains("verify-before-commit"));
426    }
427
428    #[test]
429    fn fatal_helper_aborts() {
430        let root = TempRoot::new("direct-abort");
431
432        if std::env::var(ABORT_CASE).as_deref() == Ok("direct") {
433            fatal_io("test-fixture", io::Error::other("forced failure"));
434        }
435
436        let output = run_abort_case("direct", &root.file(), "tests::fatal_helper_aborts");
437        assert!(!output.status.success());
438
439        let stderr = String::from_utf8_lossy(&output.stderr);
440        assert_eq!(
441            stderr
442                .lines()
443                .filter(|line| line.contains("kcode-k1-order-store fatal"))
444                .count(),
445            1
446        );
447        assert!(stderr.contains("test-fixture"));
448        assert!(stderr.contains("forced failure"));
449    }
450
451    #[test]
452    fn reconstructs_million_entry_fixture() {
453        let root = TempRoot::new("million");
454        let path = root.file();
455        let count = 1_000_000_u64;
456        let entry_subsystem = subsystem(b'm');
457        let mut bytes = Vec::with_capacity(count as usize * RECORD_BYTES);
458
459        for value in 0..count {
460            bytes.extend_from_slice(id(value).as_bytes());
461            bytes.extend_from_slice(entry_subsystem.as_bytes());
462        }
463
464        fs::write(&path, bytes).unwrap();
465
466        let started = Instant::now();
467        let store = OrderStore::open(&path).unwrap();
468        assert!(started.elapsed() < Duration::from_secs(5));
469        assert_eq!(store.entries().len(), count as usize);
470        assert_eq!(store.entries()[0], (id(0), entry_subsystem));
471        assert_eq!(
472            store.entries()[count as usize - 1],
473            (id(count - 1), entry_subsystem)
474        );
475        assert_eq!(store.index_of(id(0)), Some(0));
476        assert_eq!(store.index_of(id(count - 1)), Some(count as usize - 1));
477    }
478}