use std::collections::HashMap;
use std::fmt;
use std::fs::{File, OpenOptions};
use std::io::{Read, Seek, SeekFrom, Write};
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex, MutexGuard};
use lgwks_std::hash::{Digest, Hasher};
use crate::journal::frame::SaturatingFrom;
use lgwks_std::wire::{WireError, from_bytes, to_bytes};
use crate::effect::RunId;
use crate::journal::frame::{self, Cursor, HEAD_BYTES};
use crate::journal::owner::{self, Stage, StorageGate, StorageOwner, SubmitError};
use crate::script::run_store::{Appended, RunRecords, StagedRecord, StoredValue};
use crate::script::{FlowError, StepKey};
use super::definition::{DefinitionIdentity, Drift, STORE_FORMAT};
pub const MAX_RECORD_BYTES: usize = 256 * 1024;
pub const MAX_RECORDS_PER_RUN: u64 = 65_536;
pub const MAX_STORE_BYTES: u64 = 64 * 1024 * 1024;
const STORE_MAGIC: &[u8; 16] = b"lgwks-runstore\x00\x00";
const STORE_HEADER: [u8; 16] = {
let mut header = *STORE_MAGIC;
header[15] = STORE_FORMAT;
header
};
const VERSION_BYTE: usize = STORE_MAGIC.len() - 1;
fn check_format_version(header: [u8; STORE_HEADER.len()]) -> Result<(), StoreError> {
if header[..VERSION_BYTE] != STORE_MAGIC[..VERSION_BYTE] {
let refusal = Err(StoreError::NotAStore);
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "check_format_version: returning an error to the caller");
return refusal;
}
let found = header[VERSION_BYTE];
if found == STORE_FORMAT {
return Ok(());
}
Err(StoreError::FormatVersion {
found,
expected: STORE_FORMAT,
})
}
fn genesis_head() -> Digest {
let mut hasher = Hasher::new();
hasher.write_framed(b"lgwks-runstore/genesis");
hasher.finalize()
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
#[non_exhaustive]
pub enum StoreLimitKind {
RecordBytes,
Records,
StoreBytes,
}
#[derive(Debug)]
#[non_exhaustive]
pub enum StoreError {
Limit {
kind: StoreLimitKind,
requested: u64,
limit: u64,
},
Storage {
cause: std::io::Error,
},
NotAStore,
FormatVersion {
found: u8,
expected: u8,
},
Corrupt {
at: u64,
},
Encoding {
cause: WireError,
},
ForeignTenant {
owner: String,
asked: String,
},
UnknownRun {
run: String,
tenant: String,
},
Incompatible {
run: String,
drift: Drift,
},
}
impl StoreError {
pub(crate) fn storage(cause: std::io::Error) -> Self {
Self::Storage { cause }
}
}
impl fmt::Display for StoreError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match *self {
Self::Limit {
kind,
requested,
limit,
} => {
let what = match kind {
StoreLimitKind::RecordBytes => "one record's bytes",
StoreLimitKind::Records => "records in one run",
StoreLimitKind::StoreBytes => "bytes in one run's store",
};
write!(
formatter,
"{what}: {requested} exceeds the ceiling of {limit}"
)
}
Self::Storage { ref cause } => {
write!(formatter, "the run store's device refused: {cause}")
}
Self::NotAStore => {
formatter.write_str("this file is not a run store; refusing to read it as records")
}
Self::FormatVersion { found, expected } => write!(
formatter,
"this is a run store in format version {found}, and this build reads version \
{expected}; refusing to read records whose definition identity this version \
cannot reconstruct"
),
Self::Corrupt { at } => write!(
formatter,
"run store frame {at} does not follow from the records before it; \
the bytes are refused, not trimmed"
),
Self::Encoding { ref cause } => {
write!(formatter, "the step's value could not be archived: {cause}")
}
Self::ForeignTenant {
ref owner,
ref asked,
} => write!(
formatter,
"run belongs to tenant {owner:?}, not {asked:?}; refusing to read its records"
),
Self::UnknownRun {
ref run,
ref tenant,
} => write!(
formatter,
"no records for run {run} in tenant {tenant:?}'s store; refusing to resume a \
run this store cannot attribute to it"
),
Self::Incompatible { ref run, ref drift } => write!(
formatter,
"run {run} was recorded under a different definition: {}; refusing to replay \
it under this one",
drift
),
}
}
}
impl std::error::Error for StoreError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match *self {
Self::Storage { ref cause } => Some(cause),
Self::Encoding { ref cause } => Some(cause),
Self::Limit { .. }
| Self::NotAStore
| Self::FormatVersion { .. }
| Self::Corrupt { .. }
| Self::ForeignTenant { .. }
| Self::UnknownRun { .. }
| Self::Incompatible { .. } => None,
}
}
}
impl From<StoreError> for FlowError {
fn from(error: StoreError) -> Self {
match error {
StoreError::Incompatible { drift, .. } => Self::incompatible("", drift),
other => Self::Store {
at: Arc::from(""),
source: Box::new(other),
},
}
}
}
#[derive(Clone)]
pub struct RunStore {
inner: Arc<StoreInner>,
}
struct StoreInner {
owner: StorageOwner<Arc<Mutex<Index>>, Appended>,
path: PathBuf,
index: Arc<Mutex<Index>>,
unreadable: AtomicBool,
}
#[derive(Debug)]
struct Index {
runs: HashMap<RunId, RunIndex>,
committed: u64,
tail: Digest,
}
#[derive(Debug)]
struct RunIndex {
tenant: String,
definition: DefinitionIdentity,
steps: HashMap<String, Held>,
}
#[derive(Debug)]
struct Held {
path: String,
bytes: Vec<u8>,
}
impl fmt::Debug for RunStore {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("RunStore")
.field("path", &self.inner.path)
.field(
"runs",
&self.inner.index.lock().map_or(0, |index| index.runs.len()),
)
.field(
"committed",
&self.inner.index.lock().map_or(0, |index| index.committed),
)
.finish()
}
}
impl RunStore {
pub fn open(path: impl Into<PathBuf>) -> Result<Self, StoreError> {
Self::open_impl(path.into(), false)
}
pub fn open_with_stalled_device(path: impl Into<PathBuf>) -> Result<Self, StoreError> {
Self::open_impl(path.into(), true)
}
fn open_impl(path: PathBuf, stalled: bool) -> Result<Self, StoreError> {
let existed = path.exists();
let mut file = OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.open(&path)
.map_err(StoreError::storage)?;
let header = u64::saturating_from(STORE_HEADER.len());
if !existed || file.metadata().map_err(StoreError::storage)?.len() == 0 {
file.write_all(&STORE_HEADER)
.and_then(|()| file.sync_all())
.map_err(StoreError::storage)?;
}
let index = replay(&mut file)?;
if file.metadata().map_err(StoreError::storage)?.len() != index.committed {
file.set_len(index.committed).map_err(StoreError::storage)?;
file.sync_all().map_err(StoreError::storage)?;
}
debug_assert!(
index.committed >= header,
"the header is part of the committed length"
);
let index = Arc::new(Mutex::new(index));
let owner =
StorageOwner::spawn(file, Arc::clone(&index), stalled).map_err(StoreError::storage)?;
Ok(Self {
inner: Arc::new(StoreInner {
owner,
path,
index,
unreadable: AtomicBool::new(false),
}),
})
}
pub fn fail_next_index_read(&self) {
self.inner.unreadable.store(true, Ordering::SeqCst);
}
#[must_use]
pub fn storage_gate(&self) -> StorageGate {
self.inner.owner.gate()
}
pub fn release_device(&self) {
self.storage_gate().release();
}
#[cfg(test)]
pub(crate) fn fail_next_flush(&self) {
self.inner.owner.fail_next_flush();
}
pub fn open_in(dir: &Path, tenant: &str) -> Result<Self, StoreError> {
std::fs::create_dir_all(dir).map_err(StoreError::storage)?;
Self::open(dir.join(format!("{tenant}.runstore")))
}
#[must_use]
pub fn path(&self) -> &Path {
&self.inner.path
}
#[must_use]
pub fn record_count(&self, run: RunId) -> usize {
self.index()
.runs
.get(&run)
.map_or(0, |held| held.steps.len())
}
pub fn lookup(
&self,
tenant: &str,
run: RunId,
key: StepKey,
) -> Result<Option<Vec<u8>>, FlowError> {
Ok(RunRecords::lookup(self, tenant, run, key)?.map(|stored| stored.into_bytes()))
}
#[must_use]
pub fn committed_bytes(&self) -> u64 {
self.index().committed
}
#[must_use]
pub fn flush_counts(&self) -> (u64, u64) {
self.inner.owner.flush_counts()
}
#[must_use]
pub fn knows_run(&self, run: RunId) -> bool {
self.index().runs.contains_key(&run)
}
#[must_use]
pub fn tenant_of(&self, run: RunId) -> Option<String> {
let index = self.index();
index.runs.get(&run).map(|held| held.tenant.clone())
}
#[must_use]
pub fn definition_of(&self, run: RunId) -> Option<DefinitionIdentity> {
self.index()
.runs
.get(&run)
.map(|held| held.definition.clone())
}
pub fn check_definition(
&self,
run: RunId,
identity: &DefinitionIdentity,
) -> Result<(), StoreError> {
let Some(recorded) = self.definition_of(run) else {
return Ok(());
};
match identity.drift_from(&recorded) {
None => Ok(()),
Some(drift) => Err(StoreError::Incompatible {
run: run.id().to_hex(),
drift,
}),
}
}
fn owned(&self, record: StagedRecord<'_>) -> Staged {
Staged {
record: Stored {
run: record.run(),
key: record.key().as_bytes().to_vec(),
tenant: record.tenant().to_owned(),
path: record.path().to_owned(),
definition: StoredDefinition::of(record.definition()),
value: record.bytes().to_vec(),
},
key: record.key(),
}
}
fn index(&self) -> MutexGuard<'_, Index> {
owner::lock(&self.inner.index)
}
fn step_readable(&self) -> Result<(), StoreError> {
if self.inner.unreadable.swap(false, Ordering::SeqCst) {
let refusal = Err(StoreError::storage(std::io::Error::other(
"the run store's device refused a read of its index",
)));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "step_readable: returning an error to the caller");
return refusal;
}
Ok(())
}
}
impl RunRecords for RunStore {
fn lookup(
&self,
tenant: &str,
run: RunId,
key: StepKey,
) -> Result<Option<StoredValue>, FlowError> {
self.step_readable()?;
let index = self.index();
let Some(held) = index.runs.get(&run) else {
return Ok(None);
};
if held.tenant != tenant {
let refusal = Err(StoreError::ForeignTenant {
owner: held.tenant.clone(),
asked: tenant.to_owned(),
}
.into());
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "lookup: returning an error to the caller");
return refusal;
}
Ok(held
.steps
.get(&key.to_hex())
.map(|step| StoredValue::new(step.path.clone(), step.bytes.clone())))
}
fn append(
&self,
tenant: &str,
run: RunId,
key: StepKey,
path: &str,
definition: &DefinitionIdentity,
bytes: Vec<u8>,
) -> Result<Appended, FlowError> {
let staged = self.owned(RunRecords::stage(
self, tenant, run, key, path, definition, bytes,
));
self.commit(staged.record, staged.key)
.map_err(FlowError::from)
}
fn compatibility(
&self,
run: RunId,
definition: &DefinitionIdentity,
) -> Result<bool, FlowError> {
self.step_readable()?;
self.check_definition(run, definition)
.map(|()| true)
.map_err(FlowError::from)
}
fn drift(&self, run: RunId, definition: &DefinitionIdentity) -> Result<Drift, FlowError> {
let Some(recorded) = self.definition_of(run) else {
let refusal = Err(FlowError::failed(format!(
"this store holds no definition identity for run {}, so it cannot name a drift \
axis; refusing to invent one",
run.id().to_hex()
)));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "drift: returning an error to the caller");
return refusal;
};
match definition.drift_from(&recorded) {
Some(drift) => Ok(drift),
None => Err(FlowError::failed(format!(
"run {} was refused as incompatible and yet agrees with the declared \
definition; refusing to name an axis for a disagreement that is not there",
run.id().to_hex()
))),
}
}
fn append_async<'a>(
&'a self,
record: StagedRecord<'a>,
) -> crate::BoxFuture<'a, Result<Appended, FlowError>> {
let staged = self.owned(record);
Box::pin(async move {
self.commit_async(staged.record, staged.key)
.await
.map_err(FlowError::from)
})
}
}
impl From<SubmitError> for StoreError {
fn from(cause: SubmitError) -> Self {
Self::Storage {
cause: cause.into_io(),
}
}
}
struct Staged {
record: Stored,
key: StepKey,
}
impl RunStore {
fn commit(&self, stored: Stored, key: StepKey) -> Result<Appended, StoreError> {
self.inner
.owner
.submit(move |file, index| append_on_owner(file, index, &stored, key))
.map_err(StoreError::from)
}
fn commit_async<'a>(
&'a self,
stored: Stored,
key: StepKey,
) -> crate::BoxFuture<'a, Result<Appended, StoreError>> {
Box::pin(async move {
self.inner
.owner
.submit_async(move |file, index| append_on_owner(file, index, &stored, key))
.await
.map_err(StoreError::from)
})
}
}
fn append_on_owner(
file: &mut File,
shared: &Arc<Mutex<Index>>,
stored: &Stored,
key: StepKey,
) -> std::io::Result<Stage<Appended, Arc<Mutex<Index>>>> {
let mut index = owner::lock(shared);
decide_and_write(file, &mut index, stored, key)
}
fn decide_and_write(
file: &mut File,
index: &mut Index,
stored: &Stored,
key: StepKey,
) -> std::io::Result<Stage<Appended, Arc<Mutex<Index>>>> {
let owned = index.runs.get(&stored.run);
if let Some(owned) = owned {
if owned.tenant != stored.tenant {
let refusal = Err(refusal(StoreError::ForeignTenant {
owner: owned.tenant.clone(),
asked: stored.tenant.clone(),
}));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "decide_and_write: returning an error to the caller");
return refusal;
}
if let Some(step) = owned.steps.get(&key.to_hex()) {
return Ok(Stage::Settled(Ok(if step.bytes == stored.value {
Appended::AlreadyRecorded
} else {
Appended::Conflicting
})));
}
if u64::saturating_from(owned.steps.len()) >= MAX_RECORDS_PER_RUN {
let refusal = Err(refusal(StoreError::Limit {
kind: StoreLimitKind::Records,
requested: u64::saturating_from(owned.steps.len()).saturating_add(1),
limit: MAX_RECORDS_PER_RUN,
}));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "decide_and_write: returning an error to the caller");
return refusal;
}
}
let previous = tail_head(index);
let (framed, head) = frame(stored, &previous).map_err(refusal)?;
let staged = u64::saturating_from(framed.len());
let next = index
.committed
.checked_add(staged)
.ok_or_else(|| refusal(store_full(u64::MAX)))?;
if next > MAX_STORE_BYTES {
let refusal = Err(refusal(store_full(next)));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "decide_and_write: returning an error to the caller");
return refusal;
}
let on_disk = file
.metadata()
.map_err(StoreError::storage)
.map_err(refusal)?
.len();
if on_disk != index.committed {
let refusal = Err(refusal(StoreError::Corrupt {
at: u64::saturating_from(index.runs.get(&stored.run).map_or(0, |run| run.steps.len())),
}));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "decide_and_write: returning an error to the caller");
return refusal;
}
index.committed = next;
index.tail = head;
file.write_all(&framed)
.map_err(StoreError::storage)
.map_err(refusal)?;
let held = Held {
path: stored.path.clone(),
bytes: stored.value.clone(),
};
let key = hex_of(&stored.key);
let run = stored.run;
let tenant = stored.tenant.clone();
let definition = stored.definition.to_identity();
Ok(Stage::Unsynced {
answer: Appended::Recorded,
bytes: framed.len(),
settle: Box::new(move |index: &mut Arc<Mutex<Index>>| {
let mut index = owner::lock(index);
let entry = index.runs.entry(run).or_insert_with(|| RunIndex {
tenant: tenant.clone(),
definition: definition.clone(),
steps: HashMap::new(),
});
entry.steps.insert(key, held);
}),
})
}
fn refusal(cause: StoreError) -> std::io::Error {
std::io::Error::other(cause.to_string())
}
fn tail_head(index: &Index) -> Digest {
index.tail
}
#[derive(
Debug, Clone, lgwks_std::wire::Archive, lgwks_std::wire::Serialize, lgwks_std::wire::Deserialize,
)]
#[rkyv(
attr(non_exhaustive),
crate = lgwks_std::wire::rkyv,
compare(PartialEq),
derive(Debug)
)]
struct Stored {
#[rkyv(attr(doc = "The run this record belongs to."))]
run: RunId,
#[rkyv(attr(doc = "The step key's digest."))]
key: Vec<u8>,
#[rkyv(attr(doc = "The tenant that minted the run."))]
tenant: String,
#[rkyv(attr(doc = "The step's path, for attribution."))]
path: String,
#[rkyv(attr(doc = "The definition the record was written under, as its fields."))]
definition: StoredDefinition,
#[rkyv(attr(doc = "The archived value the step returned."))]
value: Vec<u8>,
}
#[derive(
Debug, Clone, lgwks_std::wire::Archive, lgwks_std::wire::Serialize, lgwks_std::wire::Deserialize,
)]
#[rkyv(
attr(non_exhaustive),
crate = lgwks_std::wire::rkyv,
compare(PartialEq),
derive(Debug)
)]
struct StoredDefinition {
#[rkyv(attr(doc = "The task name."))]
name: String,
#[rkyv(attr(doc = "The declared definition revision."))]
revision: u64,
#[rkyv(attr(doc = "The input digest."))]
input: Vec<u8>,
#[rkyv(attr(doc = "How many durable steps the definition declares."))]
steps: u64,
#[rkyv(attr(doc = "The declared durable-value schema id."))]
codec: String,
}
fn digest_from_record(bytes: &[u8]) -> Digest {
let mut digest = [0u8; 32];
let copied = bytes.len().min(digest.len());
digest[..copied].copy_from_slice(&bytes[..copied]);
Digest::from_bytes(digest)
}
impl StoredDefinition {
fn of(identity: &DefinitionIdentity) -> Self {
Self {
name: identity.name().to_owned(),
revision: identity.revision(),
input: identity.input().as_bytes().to_vec(),
steps: u64::saturating_from(identity.steps()),
codec: identity.codec().to_owned(),
}
}
fn to_identity(&self) -> DefinitionIdentity {
DefinitionIdentity::new(
&self.name,
self.revision,
digest_from_record(&self.input),
usize::saturating_from(self.steps),
)
.with_codec(&self.codec)
}
}
impl Stored {
fn head_from(&self, previous: &Digest, _archived: &[u8]) -> Digest {
let mut hasher = Hasher::new();
hasher.write_framed(previous.as_bytes());
hasher.write_framed(&self.run_key_bytes());
hasher.write_framed(&self.key);
hasher.write_framed(self.tenant.as_bytes());
hasher.write_framed(self.path.as_bytes());
hasher.write_framed(self.definition.name.as_bytes());
hasher.write_framed(&self.definition.revision.to_le_bytes());
hasher.write_framed(&self.definition.input);
hasher.write_framed(&self.definition.steps.to_le_bytes());
hasher.write_framed(self.definition.codec.as_bytes());
hasher.write_framed(&self.value);
hasher.finalize()
}
fn run_key_bytes(&self) -> Vec<u8> {
self.run.id().to_hex().into_bytes()
}
}
fn frame(stored: &Stored, previous: &Digest) -> Result<(Vec<u8>, Digest), StoreError> {
frame::frame_record(
stored,
previous,
MAX_RECORD_BYTES,
|stored| {
to_bytes::<WireError>(stored)
.map(|bytes| bytes.as_ref().to_vec())
.map_err(|cause| StoreError::Encoding { cause })
},
Stored::head_from,
record_too_large,
)
}
fn store_full(requested: u64) -> StoreError {
StoreError::Limit {
kind: StoreLimitKind::StoreBytes,
requested,
limit: MAX_STORE_BYTES,
}
}
fn record_too_large(len: usize) -> StoreError {
StoreError::Limit {
kind: StoreLimitKind::RecordBytes,
requested: u64::saturating_from(len),
limit: u64::saturating_from(MAX_RECORD_BYTES),
}
}
fn read_full(reader: &mut impl Read, buf: &mut [u8]) -> Result<bool, StoreError> {
let mut filled = 0usize;
while filled < buf.len() {
match reader.read(&mut buf[filled..]) {
Ok(0) => return Ok(false),
Ok(read) => filled = filled.saturating_add(read),
Err(ref error) if error.kind() == std::io::ErrorKind::Interrupted => {}
Err(cause) => {
let refusal = Err(StoreError::storage(cause));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "read_full: returning an error to the caller");
return refusal;
}
}
}
Ok(true)
}
struct Framed {
stored: Stored,
payload: Vec<u8>,
head: [u8; HEAD_BYTES],
declared: usize,
}
fn next_frame(file: &mut File, cursor: &Cursor<'_>) -> Result<Option<Framed>, StoreError> {
let at = cursor.at;
let corrupt = || StoreError::Corrupt { at };
let Some(raw) = frame::read_raw(
file,
cursor,
MAX_RECORD_BYTES,
StoreError::storage,
corrupt,
|previous, payload| {
let aligned = frame::decodable(payload);
from_bytes::<Stored, WireError>(aligned.as_slice())
.ok()
.map(|stored| stored.head_from(previous, payload))
},
)?
else {
return Ok(None);
};
let stored = from_bytes::<Stored, WireError>(&raw.payload).map_err(|error| {
lgwks_std::trace::debug!(?error, at, "next_frame: the payload did not decode");
StoreError::Corrupt { at }
})?;
Ok(Some(Framed {
stored,
declared: raw.payload.len(),
payload: raw.payload,
head: raw.head,
}))
}
fn replay(file: &mut File) -> Result<Index, StoreError> {
let total = file.metadata().map_err(StoreError::storage)?.len();
file.seek(SeekFrom::Start(0)).map_err(StoreError::storage)?;
let mut header = [0u8; STORE_HEADER.len()];
if !read_full(file, &mut header)? {
let refusal = Err(StoreError::NotAStore);
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "replay: returning an error to the caller");
return refusal;
}
check_format_version(header)?;
let mut index = Index {
runs: HashMap::new(),
committed: u64::saturating_from(STORE_HEADER.len()),
tail: genesis_head(),
};
let mut previous = genesis_head();
let mut at = 0u64;
while let Some(Framed {
stored,
payload,
head,
declared,
}) = next_frame(file, &Cursor::new(at, index.committed, &previous))?
{
if stored.head_from(&previous, &payload) != Digest::from_bytes(head) {
let refusal = Err(StoreError::Corrupt { at });
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "replay: returning an error to the caller");
return refusal;
}
let frame_len = frame::framed_len(declared);
if index
.committed
.checked_add(frame_len)
.is_none_or(|end| end > total)
{
break;
}
index.committed = index.committed.saturating_add(frame_len);
index.tail = stored.head_from(&previous, &payload);
previous = stored.head_from(&previous, &payload);
at = at.saturating_add(1);
if at > MAX_RECORDS_PER_RUN {
let refusal = Err(StoreError::Limit {
kind: StoreLimitKind::Records,
requested: at,
limit: MAX_RECORDS_PER_RUN,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "replay: returning an error to the caller");
return refusal;
}
let run = index.runs.entry(stored.run).or_insert_with(|| RunIndex {
tenant: stored.tenant.clone(),
definition: stored.definition.to_identity(),
steps: HashMap::new(),
});
run.steps.insert(
hex_of(&stored.key),
Held {
path: stored.path.clone(),
bytes: stored.value.clone(),
},
);
}
Ok(index)
}
fn hex_of(bytes: &[u8]) -> String {
const DIGITS: &[u8; 16] = b"0123456789abcdef";
let mut hex = Vec::with_capacity(bytes.len().saturating_mul(2));
for byte in bytes.iter().copied() {
hex.push(DIGITS[usize::from(byte >> 4) & 0x0f]);
hex.push(DIGITS[usize::from(byte) & 0x0f]);
}
hex.into_iter().map(char::from).collect()
}
#[cfg(test)]
mod tests {
use super::*;
use crate::journal::frame::probe::{Scratch, declared_at, frame_starts, with_prefix};
type TestResult = Result<(), Box<dyn std::error::Error>>;
fn written(path: &Path, count: u8) -> Result<Vec<u8>, Box<dyn std::error::Error>> {
let run = RunId::from_hex(&format!("a{}", "0".repeat(31)))?;
let identity =
DefinitionIdentity::new("tail", 1, Digest::from_bytes([1; 32]), usize::from(count));
let mut bytes = STORE_HEADER.to_vec();
let mut previous = genesis_head();
for step in 0..count {
let stored = Stored {
run,
key: vec![step; 32],
tenant: "acme".to_owned(),
path: format!("s{step}"),
definition: StoredDefinition::of(&identity),
value: vec![step; 40 + usize::from(step)],
};
let (framed, head) = frame(&stored, &previous)?;
bytes.extend_from_slice(&framed);
previous = head;
}
std::fs::write(path, &bytes)?;
Ok(bytes)
}
fn require_refused(path: &Path, before: &[u8], at: u64, why: &str) -> TestResult {
let outcome: TestResult = match RunStore::open(path) {
Err(StoreError::Corrupt { at: named }) => {
assert_eq!(named, at, "{why}");
Ok(())
}
Err(other) => Err(format!("{why}: expected Corrupt, got {other}").into()),
Ok(opened) => Err(format!(
"{why}: reopened at {} committed bytes, so an acknowledged frame was lost",
opened.committed_bytes()
)
.into()),
};
outcome?;
assert_eq!(std::fs::read(path)?, before, "{why}: refused bytes move");
Ok(())
}
#[test]
fn a_lengthened_acknowledged_final_record_is_refused_not_trimmed() -> TestResult {
let scratch = Scratch::new("store-lengthened")?;
let bytes = written(scratch.path(), 3)?;
let last = frame_starts(&bytes, STORE_HEADER.len())?[2];
let declared = declared_at(&bytes, last);
for extra in (1u32..=40).chain([100, 255, 1024, 4096]) {
let lied = with_prefix(&bytes, last, declared + extra);
std::fs::write(scratch.path(), &lied)?;
require_refused(scratch.path(), &lied, 2, &format!("final record L+{extra}"))?;
}
Ok(())
}
#[test]
fn an_inflated_record_with_records_behind_it_is_refused_untouched() -> TestResult {
let scratch = Scratch::new("store-inflated")?;
let bytes = written(scratch.path(), 3)?;
let middle = frame_starts(&bytes, STORE_HEADER.len())?[1];
let remaining = u32::try_from(bytes.len() - middle)?;
for declared in [remaining, remaining + 1, remaining + 31, remaining + 500] {
let lied = with_prefix(&bytes, middle, declared);
std::fs::write(scratch.path(), &lied)?;
require_refused(scratch.path(), &lied, 1, &format!("declared {declared}"))?;
}
Ok(())
}
#[test]
fn a_damaged_cut_record_with_an_acknowledged_one_behind_it_is_refused() -> TestResult {
let scratch = Scratch::new("store-damaged")?;
let bytes = written(scratch.path(), 3)?;
let middle = frame_starts(&bytes, STORE_HEADER.len())?[1];
let remaining = u32::try_from(bytes.len() - middle)?;
let mut lied = with_prefix(&bytes, middle, remaining + 7);
if let Some(byte) = lied.get_mut(middle + 12) {
*byte ^= 0x55;
}
std::fs::write(scratch.path(), &lied)?;
require_refused(scratch.path(), &lied, 1, "damaged middle, lengthened")
}
#[test]
fn an_append_cut_inside_the_final_record_is_repaired() -> TestResult {
let scratch = Scratch::new("store-cut")?;
let bytes = written(scratch.path(), 3)?;
let last = frame_starts(&bytes, STORE_HEADER.len())?[2];
let kept = u64::try_from(last)?;
let whole = bytes.len() - last;
for cut in [
1,
3,
4,
5,
whole >> 1,
whole - 33,
whole - 32,
whole - 31,
whole - 1,
] {
std::fs::write(
scratch.path(),
bytes.get(..last + cut).ok_or("cut past the end")?,
)?;
let store = RunStore::open(scratch.path())?;
assert_eq!(
store.committed_bytes(),
kept,
"cut {cut}: the two records survive"
);
drop(store);
assert_eq!(std::fs::metadata(scratch.path())?.len(), kept, "cut {cut}");
}
Ok(())
}
#[test]
fn a_ceiling_sized_noise_tail_is_trimmed_in_bounded_time() -> TestResult {
let scratch = Scratch::new("store-noise")?;
let mut bytes = written(scratch.path(), 1)?;
let kept = u64::try_from(bytes.len())?;
bytes.extend_from_slice(&u32::try_from(MAX_RECORD_BYTES)?.to_be_bytes());
let mut state = 0x9e37_79b9_7f4a_7c15u64;
for _ in 0..MAX_RECORD_BYTES + HEAD_BYTES - 1 {
state = state
.wrapping_mul(6_364_136_223_846_793_005)
.wrapping_add(1);
bytes.push(state.to_be_bytes()[0]);
}
std::fs::write(scratch.path(), &bytes)?;
let started = std::time::Instant::now();
let store = RunStore::open(scratch.path())?;
let elapsed = started.elapsed();
assert_eq!(
store.committed_bytes(),
kept,
"noise holds no acknowledged record"
);
assert!(elapsed.as_secs() < 5, "the search took {elapsed:?}");
Ok(())
}
}