use kcode_k1_transaction::GENESIS_PARENT;
pub use kcode_k1_transaction::SubsystemId;
pub use kcode_k1_transaction_store::TxId;
use std::collections::HashMap;
use std::fs::{self, File, OpenOptions};
use std::io::{self, Read, Seek, SeekFrom, Write};
use std::path::Path;
const RECORD_BYTES: usize = 32;
const RECORD_BYTES_U64: u64 = 32;
type Entries = Vec<(TxId, SubsystemId)>;
type Indexes = HashMap<TxId, usize>;
pub struct OrderStore {
file: File,
entries: Entries,
indexes: Indexes,
}
impl OrderStore {
pub fn create(path: &Path) -> Result<Self, String> {
let file = match OpenOptions::new()
.read(true)
.write(true)
.create_new(true)
.open(path)
{
Ok(file) => file,
Err(error)
if matches!(
error.kind(),
io::ErrorKind::AlreadyExists
| io::ErrorKind::NotFound
| io::ErrorKind::InvalidInput
) =>
{
return Err(format!("cannot create ordering file: {error}"));
}
Err(error) => fatal_io("create", error),
};
Ok(Self {
file,
entries: Vec::new(),
indexes: HashMap::new(),
})
}
pub fn open(path: &Path) -> Result<Self, String> {
let metadata = fs::symlink_metadata(path)
.map_err(|error| format!("cannot inspect ordering file: {error}"))?;
if !metadata.file_type().is_file() {
return Err("ordering path is not a regular file".to_owned());
}
let mut file = OpenOptions::new()
.read(true)
.write(true)
.open(path)
.map_err(|error| format!("cannot open ordering file: {error}"))?;
let (entries, indexes) = reconstruct(&mut file)?;
Ok(Self {
file,
entries,
indexes,
})
}
pub fn entries(&self) -> &[(TxId, SubsystemId)] {
&self.entries
}
pub fn index_of(&self, id: TxId) -> Option<usize> {
self.indexes.get(&id).copied()
}
pub fn commit(
&mut self,
retained_len: usize,
id: TxId,
subsystem: SubsystemId,
) -> Result<(), String> {
if retained_len > self.entries.len() {
return Err("retained length exceeds the canonical order".to_owned());
}
if id == GENESIS_PARENT {
return Err("cannot commit the genesis sentinel".to_owned());
}
if self.indexes.contains_key(&id) {
return Err("transaction ID is already canonical".to_owned());
}
let expected_len = byte_len(self.entries.len())?;
let actual_len = self
.file
.metadata()
.unwrap_or_else(|error| fatal_io("inspect-before-commit", error))
.len();
if actual_len != expected_len {
fatal_io(
"verify-before-commit",
io::Error::other(format!(
"ordering file length changed from {expected_len} to {actual_len}"
)),
);
}
if retained_len < self.entries.len() {
let retained_bytes = byte_len(retained_len)?;
self.file
.set_len(retained_bytes)
.unwrap_or_else(|error| fatal_io("truncate", error));
self.file
.sync_data()
.unwrap_or_else(|error| fatal_io("sync-truncation", error));
self.file
.seek(SeekFrom::Start(retained_bytes))
.unwrap_or_else(|error| fatal_io("seek-replacement", error));
write_record(&mut self.file, id, subsystem, "write-replacement");
self.file
.sync_data()
.unwrap_or_else(|error| fatal_io("sync-replacement", error));
for (removed_id, _) in self.entries.drain(retained_len..) {
self.indexes.remove(&removed_id);
}
} else {
self.file
.seek(SeekFrom::Start(expected_len))
.unwrap_or_else(|error| fatal_io("seek-append", error));
write_record(&mut self.file, id, subsystem, "write-append");
self.file
.sync_data()
.unwrap_or_else(|error| fatal_io("sync-append", error));
}
let index = self.entries.len();
self.indexes.insert(id, index);
self.entries.push((id, subsystem));
Ok(())
}
}
fn reconstruct(file: &mut File) -> Result<(Entries, Indexes), String> {
file.seek(SeekFrom::Start(0))
.map_err(|error| format!("cannot seek ordering file: {error}"))?;
let mut bytes = Vec::new();
file.read_to_end(&mut bytes)
.map_err(|error| format!("cannot read ordering file: {error}"))?;
if bytes.len() % RECORD_BYTES != 0 {
return Err("ordering file length is not a multiple of 32".to_owned());
}
let count = bytes.len() / RECORD_BYTES;
let mut entries = Vec::with_capacity(count);
let mut indexes = HashMap::with_capacity(count);
for record in bytes.chunks_exact(RECORD_BYTES) {
let id = TxId::from_bytes(
record[..12]
.try_into()
.expect("transaction ID range has fixed length"),
);
if id == GENESIS_PARENT {
return Err("ordering file contains the genesis sentinel".to_owned());
}
let subsystem = SubsystemId::from_bytes(
record[12..]
.try_into()
.expect("subsystem range has fixed length"),
)
.map_err(|error| format!("ordering file contains an invalid subsystem: {error}"))?;
let index = entries.len();
if indexes.insert(id, index).is_some() {
return Err("ordering file contains a duplicate transaction ID".to_owned());
}
entries.push((id, subsystem));
}
Ok((entries, indexes))
}
fn byte_len(records: usize) -> Result<u64, String> {
let records =
u64::try_from(records).map_err(|_| "ordering file length exceeds u64".to_owned())?;
records
.checked_mul(RECORD_BYTES_U64)
.ok_or_else(|| "ordering file length exceeds u64".to_owned())
}
fn encoded_record(id: TxId, subsystem: SubsystemId) -> [u8; RECORD_BYTES] {
let mut record = [0_u8; RECORD_BYTES];
record[..12].copy_from_slice(id.as_bytes());
record[12..].copy_from_slice(subsystem.as_bytes());
record
}
fn write_record(file: &mut File, id: TxId, subsystem: SubsystemId, operation: &'static str) {
let record = encoded_record(id, subsystem);
match file.write(&record) {
Ok(RECORD_BYTES) => {}
Ok(written) => fatal_io(
operation,
io::Error::new(
io::ErrorKind::WriteZero,
format!("short record write: wrote {written} of {RECORD_BYTES} bytes"),
),
),
Err(error) => fatal_io(operation, error),
}
}
fn fatal_io(operation: &str, error: io::Error) -> ! {
eprintln!("kcode-k1-order-store fatal {operation}: {error}");
std::process::abort()
}
#[cfg(test)]
mod tests {
use super::*;
use std::path::{Path, PathBuf};
use std::process::Command;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{Duration, Instant};
static NEXT_ROOT: AtomicU64 = AtomicU64::new(0);
const ABORT_CASE: &str = "KCODE_K1_ORDER_STORE_ABORT_CASE";
const ABORT_PATH: &str = "KCODE_K1_ORDER_STORE_ABORT_PATH";
struct TempRoot {
path: PathBuf,
}
impl TempRoot {
fn new(label: &str) -> Self {
let sequence = NEXT_ROOT.fetch_add(1, Ordering::Relaxed);
let path = std::env::temp_dir().join(format!(
"kcode-k1-order-store-{}-{}-{}",
std::process::id(),
sequence,
label
));
let _ = fs::remove_dir_all(&path);
fs::create_dir_all(&path).unwrap();
Self { path }
}
fn file(&self) -> PathBuf {
self.path.join("ordering.dat")
}
}
impl Drop for TempRoot {
fn drop(&mut self) {
let _ = fs::remove_dir_all(&self.path);
}
}
fn id(value: u64) -> TxId {
let mut bytes = [0_u8; 12];
bytes[..8].copy_from_slice(&value.to_le_bytes());
TxId::from_bytes(bytes)
}
fn subsystem(value: u8) -> SubsystemId {
SubsystemId::from_bytes([value; 20]).unwrap()
}
fn write_records(path: &Path, records: &[(TxId, SubsystemId)]) {
let mut bytes = Vec::with_capacity(records.len() * RECORD_BYTES);
for (id, subsystem) in records {
bytes.extend_from_slice(&encoded_record(*id, *subsystem));
}
fs::write(path, bytes).unwrap();
}
fn run_abort_case(case: &str, path: &Path, test_name: &str) -> std::process::Output {
Command::new(std::env::current_exe().unwrap())
.arg("--exact")
.arg(test_name)
.arg("--nocapture")
.env(ABORT_CASE, case)
.env(ABORT_PATH, path)
.output()
.unwrap()
}
#[test]
fn creates_appends_replaces_and_reopens_exact_bytes() {
let root = TempRoot::new("lifecycle");
let path = root.file();
let first = (id(1), subsystem(b'a'));
let second = (id(2), subsystem(b'b'));
let replacement = (id(3), subsystem(b'c'));
let mut store = OrderStore::create(&path).unwrap();
assert!(path.is_file());
assert!(store.entries().is_empty());
store.commit(0, first.0, first.1).unwrap();
store.commit(1, second.0, second.1).unwrap();
let mut expected = Vec::new();
expected.extend_from_slice(&encoded_record(first.0, first.1));
expected.extend_from_slice(&encoded_record(second.0, second.1));
assert_eq!(fs::read(&path).unwrap(), expected);
assert_eq!(store.entries(), &[first, second]);
assert_eq!(store.index_of(first.0), Some(0));
assert_eq!(store.index_of(second.0), Some(1));
store.commit(1, replacement.0, replacement.1).unwrap();
assert_eq!(store.entries(), &[first, replacement]);
assert_eq!(store.index_of(second.0), None);
assert_eq!(store.index_of(replacement.0), Some(1));
drop(store);
let reopened = OrderStore::open(&path).unwrap();
assert_eq!(reopened.entries(), &[first, replacement]);
assert_eq!(reopened.index_of(first.0), Some(0));
assert_eq!(reopened.index_of(replacement.0), Some(1));
}
#[test]
fn rejects_invalid_commit_arguments_without_mutation() {
let root = TempRoot::new("arguments");
let path = root.file();
let mut store = OrderStore::create(&path).unwrap();
let first = id(1);
let entry_subsystem = subsystem(b'a');
store.commit(0, first, entry_subsystem).unwrap();
let before = fs::read(&path).unwrap();
assert!(store.commit(2, id(2), entry_subsystem).is_err());
assert!(store.commit(1, GENESIS_PARENT, entry_subsystem).is_err());
assert!(store.commit(1, first, entry_subsystem).is_err());
assert_eq!(fs::read(&path).unwrap(), before);
assert_eq!(store.entries().len(), 1);
}
#[test]
fn rejects_malformed_and_invalid_files() {
let malformed = TempRoot::new("malformed");
fs::write(malformed.file(), [0_u8; RECORD_BYTES - 1]).unwrap();
assert!(OrderStore::open(&malformed.file()).is_err());
let duplicate = TempRoot::new("duplicate");
write_records(
&duplicate.file(),
&[(id(1), subsystem(b'a')), (id(1), subsystem(b'b'))],
);
assert!(OrderStore::open(&duplicate.file()).is_err());
let genesis = TempRoot::new("genesis");
write_records(&genesis.file(), &[(GENESIS_PARENT, subsystem(b'a'))]);
assert!(OrderStore::open(&genesis.file()).is_err());
let invalid_subsystem = TempRoot::new("invalid-subsystem");
let mut record = [0_u8; RECORD_BYTES];
record[..12].copy_from_slice(id(2).as_bytes());
record[12..].fill(0xff);
fs::write(invalid_subsystem.file(), record).unwrap();
assert!(OrderStore::open(&invalid_subsystem.file()).is_err());
let wrong_type = TempRoot::new("wrong-type");
assert!(OrderStore::open(&wrong_type.path).is_err());
}
#[test]
fn create_requires_an_absent_path() {
let root = TempRoot::new("create");
let path = root.file();
drop(OrderStore::create(&path).unwrap());
assert!(OrderStore::create(&path).is_err());
}
#[test]
fn external_length_change_aborts() {
let root = TempRoot::new("external-abort");
let path = root.file();
if std::env::var(ABORT_CASE).as_deref() == Ok("external") {
let child_path = PathBuf::from(std::env::var(ABORT_PATH).unwrap());
let mut store = OrderStore::create(&child_path).unwrap();
OpenOptions::new()
.append(true)
.open(&child_path)
.unwrap()
.write_all(&[0])
.unwrap();
store.commit(0, id(1), subsystem(b'a')).unwrap();
unreachable!();
}
let output = run_abort_case("external", &path, "tests::external_length_change_aborts");
assert!(!output.status.success());
let stderr = String::from_utf8_lossy(&output.stderr);
assert_eq!(
stderr
.lines()
.filter(|line| line.contains("kcode-k1-order-store fatal"))
.count(),
1
);
assert!(stderr.contains("verify-before-commit"));
}
#[test]
fn fatal_helper_aborts() {
let root = TempRoot::new("direct-abort");
if std::env::var(ABORT_CASE).as_deref() == Ok("direct") {
fatal_io("test-fixture", io::Error::other("forced failure"));
}
let output = run_abort_case("direct", &root.file(), "tests::fatal_helper_aborts");
assert!(!output.status.success());
let stderr = String::from_utf8_lossy(&output.stderr);
assert_eq!(
stderr
.lines()
.filter(|line| line.contains("kcode-k1-order-store fatal"))
.count(),
1
);
assert!(stderr.contains("test-fixture"));
assert!(stderr.contains("forced failure"));
}
#[test]
fn reconstructs_million_entry_fixture() {
let root = TempRoot::new("million");
let path = root.file();
let count = 1_000_000_u64;
let entry_subsystem = subsystem(b'm');
let mut bytes = Vec::with_capacity(count as usize * RECORD_BYTES);
for value in 0..count {
bytes.extend_from_slice(id(value).as_bytes());
bytes.extend_from_slice(entry_subsystem.as_bytes());
}
fs::write(&path, bytes).unwrap();
let started = Instant::now();
let store = OrderStore::open(&path).unwrap();
assert!(started.elapsed() < Duration::from_secs(5));
assert_eq!(store.entries().len(), count as usize);
assert_eq!(store.entries()[0], (id(0), entry_subsystem));
assert_eq!(
store.entries()[count as usize - 1],
(id(count - 1), entry_subsystem)
);
assert_eq!(store.index_of(id(0)), Some(0));
assert_eq!(store.index_of(id(count - 1)), Some(count as usize - 1));
}
}