use std::io::{Read, Seek, SeekFrom};
use lgwks_std::hash::Digest;
use super::JournalError;
pub(crate) trait SaturatingFrom<T> {
fn saturating_from(value: T) -> Self;
}
impl SaturatingFrom<usize> for u64 {
fn saturating_from(value: usize) -> Self {
match u64::try_from(value) {
Ok(narrowed) => narrowed,
Err(_wider_than_a_file_offset) => Self::MAX,
}
}
}
impl SaturatingFrom<u64> for usize {
fn saturating_from(value: u64) -> Self {
match usize::try_from(value) {
Ok(narrowed) => narrowed,
Err(_wider_than_an_index) => Self::MAX,
}
}
}
impl SaturatingFrom<u32> for usize {
fn saturating_from(value: u32) -> Self {
match usize::try_from(value) {
Ok(narrowed) => narrowed,
Err(_wider_than_an_index) => Self::MAX,
}
}
}
impl SaturatingFrom<usize> for u32 {
fn saturating_from(value: usize) -> Self {
match u32::try_from(value) {
Ok(narrowed) => narrowed,
Err(_wider_than_a_length_prefix) => Self::MAX,
}
}
}
impl SaturatingFrom<u64> for u32 {
fn saturating_from(value: u64) -> Self {
match u32::try_from(value) {
Ok(narrowed) => narrowed,
Err(_wider_than_a_report_count) => Self::MAX,
}
}
}
impl SaturatingFrom<u128> for u64 {
fn saturating_from(value: u128) -> Self {
match u64::try_from(value) {
Ok(narrowed) => narrowed,
Err(_wider_than_a_nanosecond_count) => Self::MAX,
}
}
}
impl SaturatingFrom<usize> for u8 {
fn saturating_from(value: usize) -> Self {
match u8::try_from(value) {
Ok(narrowed) => narrowed,
Err(_wider_than_a_byte) => Self::MAX,
}
}
}
pub(crate) fn read_exact_or_eof(
reader: &mut impl Read,
buf: &mut [u8],
) -> Result<Option<usize>, JournalError> {
let mut filled = 0usize;
while filled < buf.len() {
let read = reader
.read(&mut buf[filled..])
.map_err(JournalError::Storage)?;
if read == 0 {
break;
}
filled = filled.saturating_add(read);
}
if filled == 0 {
Ok(None)
} else {
Ok(Some(filled))
}
}
pub(crate) const LENGTH_BYTES: usize = 4;
pub(crate) const HEAD_BYTES: usize = 32;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum Piece {
Filled,
Interrupted,
}
pub(crate) fn read_piece(reader: &mut impl Read, buf: &mut [u8]) -> Result<Piece, std::io::Error> {
match reader.read_exact(buf) {
Ok(()) => Ok(Piece::Filled),
Err(error) if error.kind() == std::io::ErrorKind::UnexpectedEof => Ok(Piece::Interrupted),
Err(error) => Err(error),
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum Prefix {
Eof,
Torn,
Full,
}
pub(crate) fn read_prefix(
reader: &mut impl Read,
buf: &mut [u8; LENGTH_BYTES],
) -> Result<Prefix, std::io::Error> {
let mut filled = 0usize;
while filled < buf.len() {
match reader.read(&mut buf[filled..]) {
Ok(0) => break,
Ok(read) => filled = filled.saturating_add(read),
Err(ref error) if error.kind() == std::io::ErrorKind::Interrupted => {}
Err(error) => {
let refusal = Err(error);
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "read_prefix: returning an error to the caller");
return refusal;
}
}
}
Ok(match filled {
0 => Prefix::Eof,
n if n < buf.len() => Prefix::Torn,
_ => Prefix::Full,
})
}
pub(crate) fn declared_length(prefix: &[u8; LENGTH_BYTES]) -> usize {
usize::saturating_from(u32::from_be_bytes(*prefix))
}
pub(crate) fn is_possible_length(payload_len: usize, max_frame_bytes: usize) -> bool {
payload_len != 0 && payload_len <= max_frame_bytes
}
pub(crate) fn writable_length(payload_len: usize, max_frame_bytes: usize) -> Option<u32> {
u32::try_from(payload_len)
.ok()
.filter(|_| is_possible_length(payload_len, max_frame_bytes))
}
pub(crate) fn framed_len(payload_len: usize) -> u64 {
u64::saturating_from(
LENGTH_BYTES
.saturating_add(payload_len)
.saturating_add(HEAD_BYTES),
)
}
pub(crate) fn encode(payload_len: u32, payload: &[u8], head: &Digest) -> Vec<u8> {
let claimed = usize::saturating_from(payload_len);
debug_assert_eq!(
claimed,
payload.len(),
"a caller framed a payload under a length that is not its own"
);
let mut frame = Vec::with_capacity(usize::saturating_from(framed_len(claimed)));
frame.extend_from_slice(&payload_len.to_be_bytes());
frame.extend_from_slice(payload);
frame.extend_from_slice(head.as_bytes());
frame
}
pub(crate) fn frame_record<R, A, H, C, E>(
record: &R,
previous: &Digest,
max_frame_bytes: usize,
archive: A,
head: H,
refuse: C,
) -> Result<(Vec<u8>, Digest), E>
where
A: FnOnce(&R) -> Result<Vec<u8>, E>,
H: FnOnce(&R, &Digest, &[u8]) -> Digest,
C: FnOnce(usize) -> E,
{
let payload = archive(record)?;
let Some(payload_len) = writable_length(payload.len(), max_frame_bytes) else {
{
lgwks_std::trace::debug!(
payload_len = payload.len(),
max_frame_bytes,
"frame_record: the archived record exceeds the frame ceiling"
);
return Err(refuse(payload.len()));
};
};
let head = head(record, previous, &payload);
Ok((encode(payload_len, &payload, &head), head))
}
pub(crate) fn holds_acknowledged_frame<H>(
suffix: &[u8],
previous: &Digest,
max_frame_bytes: usize,
head_of: H,
) -> bool
where
H: Fn(&Digest, &[u8]) -> Option<Digest>,
{
let authenticates = |previous: &Digest, payload: &[u8], stored: &[u8]| {
head_of(previous, payload).is_some_and(|head| head.as_bytes().as_slice() == stored)
};
let room = suffix.len().saturating_sub(HEAD_BYTES);
for len in 1..=room.min(max_frame_bytes) {
let (Some(payload), Some(stored)) = (
suffix.get(..len),
suffix.get(len..len.saturating_add(HEAD_BYTES)),
) else {
continue;
};
if authenticates(previous, payload, stored) {
return true;
}
if let Some(implied) = head_of(previous, payload)
&& frame_chains_from(
suffix,
len.saturating_add(HEAD_BYTES),
&implied,
max_frame_bytes,
&authenticates,
)
{
return true;
}
}
let mut at = HEAD_BYTES.saturating_add(1);
while at
.saturating_add(LENGTH_BYTES)
.saturating_add(1)
.saturating_add(HEAD_BYTES)
<= suffix.len()
{
if let Some(before) = suffix
.get(at.saturating_sub(HEAD_BYTES)..at)
.and_then(|bytes| <[u8; HEAD_BYTES]>::try_from(bytes).ok())
&& frame_chains_from(
suffix,
at,
&Digest::from_bytes(before),
max_frame_bytes,
&authenticates,
)
{
return true;
}
at = at.saturating_add(1);
}
false
}
fn frame_chains_from(
suffix: &[u8],
at: usize,
previous: &Digest,
max_frame_bytes: usize,
authenticates: &impl Fn(&Digest, &[u8], &[u8]) -> bool,
) -> bool {
let payload_at = at.saturating_add(LENGTH_BYTES);
let Some(prefix) = suffix
.get(at..payload_at)
.and_then(|bytes| <[u8; LENGTH_BYTES]>::try_from(bytes).ok())
else {
return false;
};
let declared = declared_length(&prefix);
if !is_possible_length(declared, max_frame_bytes) {
return false;
}
let head_at = payload_at.saturating_add(declared);
let (Some(payload), Some(stored)) = (
suffix.get(payload_at..head_at),
suffix.get(head_at..head_at.saturating_add(HEAD_BYTES)),
) else {
return false;
};
authenticates(previous, payload, stored)
}
pub(crate) fn cut_holds_acknowledged<R, E, H>(
file: &mut R,
offset: u64,
previous: &Digest,
max_frame_bytes: usize,
storage: fn(std::io::Error) -> E,
head_of: H,
) -> Result<bool, E>
where
R: Read + Seek,
H: Fn(&Digest, &[u8]) -> Option<Digest>,
{
let start = offset.saturating_add(u64::saturating_from(LENGTH_BYTES));
let len = file.seek(SeekFrom::End(0)).map_err(storage)?;
let behind = len.saturating_sub(start);
let ceiling = u64::saturating_from(max_frame_bytes.saturating_add(HEAD_BYTES));
if behind > ceiling {
let grew = Err(storage(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"the file grew while its tail was being resolved; reopen it",
)));
lgwks_std::trace::debug!(
behind,
ceiling,
"cut_holds_acknowledged: more behind the prefix than a cut frame leaves"
);
return grew;
}
file.seek(SeekFrom::Start(start)).map_err(storage)?;
let Ok(width) = usize::try_from(behind) else {
let wider = Err(storage(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"the file named more bytes than this host can hold a frame in",
)));
lgwks_std::trace::debug!(
behind,
"cut_holds_acknowledged: the tail is wider than this host's address space"
);
return wider;
};
let mut suffix = vec![0u8; width];
file.read_exact(&mut suffix).map_err(storage)?;
Ok(holds_acknowledged_frame(
&suffix,
previous,
max_frame_bytes,
head_of,
))
}
#[cfg(feature = "script")]
pub(crate) struct Raw {
pub(crate) payload: Vec<u8>,
pub(crate) head: [u8; HEAD_BYTES],
}
#[cfg(feature = "script")]
pub(crate) enum Decodable<'a> {
Borrowed(&'a [u8]),
Copied(lgwks_std::wire::AlignedVec),
}
#[cfg(feature = "script")]
impl<'a> Decodable<'a> {
pub(crate) fn as_slice(&self) -> &[u8] {
match *self {
Self::Borrowed(bytes) => bytes,
Self::Copied(ref copy) => copy.as_slice(),
}
}
}
#[cfg(feature = "script")]
pub(crate) fn decodable(payload: &[u8]) -> Decodable<'_> {
const ARCHIVE_ALIGN: usize = 16;
if payload.as_ptr().align_offset(ARCHIVE_ALIGN) == 0 {
return Decodable::Borrowed(payload);
}
let mut copy = lgwks_std::wire::AlignedVec::with_capacity(payload.len());
copy.extend_from_slice(payload);
Decodable::Copied(copy)
}
#[cfg(feature = "script")]
pub(crate) struct Cursor<'chain> {
pub(crate) at: u64,
pub(crate) offset: u64,
pub(crate) previous: &'chain Digest,
}
#[cfg(feature = "script")]
impl<'chain> Cursor<'chain> {
pub(crate) fn new(at: u64, offset: u64, previous: &'chain Digest) -> Self {
Self {
at,
offset,
previous,
}
}
}
#[cfg(feature = "script")]
pub(crate) fn read_raw<R, E, H>(
file: &mut R,
at: &Cursor<'_>,
max_frame_bytes: usize,
storage: fn(std::io::Error) -> E,
corrupt: impl FnOnce() -> E,
head_of: H,
) -> Result<Option<Raw>, E>
where
R: Read + Seek,
H: Fn(&Digest, &[u8]) -> Option<Digest>,
{
let mut prefix = [0u8; LENGTH_BYTES];
let declared = match read_prefix(file, &mut prefix).map_err(storage)? {
Prefix::Eof | Prefix::Torn => return Ok(None),
Prefix::Full => declared_length(&prefix),
};
if !is_possible_length(declared, max_frame_bytes) {
lgwks_std::trace::debug!(
declared,
"read_raw: the declared length is one the writer never produces"
);
return Err(corrupt());
}
let mut payload = vec![0u8; declared];
let mut head = [0u8; HEAD_BYTES];
for piece in [&mut payload[..], &mut head[..]] {
if let Piece::Interrupted = read_piece(file, piece).map_err(storage)? {
let lies = cut_holds_acknowledged(
file,
at.offset,
at.previous,
max_frame_bytes,
storage,
head_of,
)?;
return if lies { Err(corrupt()) } else { Ok(None) };
}
}
Ok(Some(Raw { payload, head }))
}
#[cfg(test)]
pub(crate) mod probe {
use super::{LENGTH_BYTES, framed_len};
#[cfg(feature = "script")]
use std::path::Path;
use std::path::PathBuf;
use std::sync::atomic::{AtomicU64, Ordering};
pub(crate) fn scratch_path(prefix: &str, name: &str) -> Result<PathBuf, std::io::Error> {
let unique = SCRATCH_COUNTER.fetch_add(1, Ordering::Relaxed);
let Ok(since_epoch) = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH)
else {
let refusal = Err(std::io::Error::other(
"the host clock is before the Unix epoch, so no scratch path can be named uniquely",
));
lgwks_std::trace::warn!(
%prefix,
%name,
"journal test scratch: the clock cannot name a unique path"
);
return refusal;
};
Ok(std::env::temp_dir().join(format!(
"lgwks-{prefix}-{name}-{}-{unique}",
since_epoch.as_nanos()
)))
}
static SCRATCH_COUNTER: AtomicU64 = AtomicU64::new(0);
#[cfg(feature = "script")]
pub(crate) struct Scratch(PathBuf);
#[cfg(feature = "script")]
impl Scratch {
#[cfg(feature = "script")]
pub(crate) fn new(name: &str) -> Result<Self, std::io::Error> {
scratch_path("frame", name).map(Self)
}
pub(crate) fn path(&self) -> &Path {
&self.0
}
}
#[cfg(feature = "script")]
impl Drop for Scratch {
fn drop(&mut self) {
drop(std::fs::remove_file(&self.0));
}
}
pub(crate) fn frame_starts(
bytes: &[u8],
header: usize,
) -> Result<Vec<usize>, Box<dyn std::error::Error>> {
let mut starts = Vec::new();
let mut at = header;
while at < bytes.len() {
starts.push(at);
at = at.saturating_add(usize::try_from(framed_len(usize::try_from(declared_at(
bytes, at,
))?))?);
}
Ok(starts)
}
pub(crate) fn with_prefix(bytes: &[u8], at: usize, declared: u32) -> Vec<u8> {
let mut out = bytes.to_vec();
for (slot, byte) in out.iter_mut().skip(at).zip(declared.to_be_bytes()) {
*slot = byte;
}
out
}
pub(crate) fn declared_at(bytes: &[u8], at: usize) -> u32 {
let mut prefix = [0u8; LENGTH_BYTES];
for (slot, byte) in prefix.iter_mut().zip(bytes.iter().skip(at)) {
*slot = *byte;
}
u32::from_be_bytes(prefix)
}
}
#[cfg(test)]
mod tests {
use super::{
Digest, HEAD_BYTES, LENGTH_BYTES, Piece, Prefix, SaturatingFrom, declared_length, encode,
framed_len, is_possible_length, read_piece, read_prefix, writable_length,
};
use std::io::Cursor;
type TestResult = Result<(), Box<dyn std::error::Error>>;
fn head() -> Digest {
Digest::from_bytes([7u8; 32])
}
fn framed(payload: &[u8]) -> Result<Vec<u8>, Box<dyn std::error::Error>> {
let length =
writable_length(payload.len(), usize::MAX).ok_or("the grammar framed this payload")?;
let out = encode(length, payload, &head());
assert_eq!(
usize::saturating_from(framed_len(payload.len())),
out.len(),
"a frame's bytes are exactly the length the grammar says a frame is"
);
Ok(out)
}
#[test]
fn a_frame_round_trips_through_the_shared_grammar() -> TestResult {
let payload: Vec<u8> = (0..200u16)
.map(u8::try_from)
.collect::<Result<Vec<u8>, _>>()?;
let bytes = framed(&payload)?;
assert_eq!(&bytes[..LENGTH_BYTES], &200u32.to_be_bytes());
assert_eq!(
&bytes[LENGTH_BYTES..LENGTH_BYTES + payload.len()],
&payload[..]
);
assert_eq!(&bytes[bytes.len() - HEAD_BYTES..], head().as_bytes());
Ok(())
}
#[test]
fn a_prefix_classifies_eof_full_and_torn() -> TestResult {
let payload: Vec<u8> = vec![1u8, 2, 3];
let bytes = framed(&payload)?;
let mut clean = Cursor::new(Vec::new());
let mut buffer = [0u8; LENGTH_BYTES];
assert_eq!(read_prefix(&mut clean, &mut buffer)?, Prefix::Eof);
let mut whole = Cursor::new(bytes.clone());
assert_eq!(read_prefix(&mut whole, &mut buffer)?, Prefix::Full);
assert_eq!(declared_length(&buffer), 3);
for cut in 1..LENGTH_BYTES {
let zeroed: Vec<u8> = bytes[..cut].iter().map(|_| 0u8).collect();
let mut torn = Cursor::new(zeroed);
let mut scratch = [0u8; LENGTH_BYTES];
assert_eq!(read_prefix(&mut torn, &mut scratch)?, Prefix::Torn);
}
Ok(())
}
#[test]
fn a_frame_cut_inside_a_piece_is_interrupted() -> TestResult {
let payload: Vec<u8> = vec![9u8; 16];
let bytes = framed(&payload)?;
for cut in [LENGTH_BYTES + 1, LENGTH_BYTES + 8, bytes.len() - HEAD_BYTES] {
let mut source = Cursor::new(bytes[..cut].to_vec());
let wanted = bytes.len().saturating_sub(cut).saturating_add(1);
let mut piece = vec![0u8; wanted];
assert_eq!(read_piece(&mut source, &mut piece)?, Piece::Interrupted);
}
Ok(())
}
#[test]
fn only_a_possible_length_is_possible() {
assert!(is_possible_length(1, 64 * 1024));
assert!(is_possible_length(64 * 1024, 64 * 1024));
assert!(!is_possible_length(0, 64 * 1024));
assert!(!is_possible_length(64 * 1024 + 1, 64 * 1024));
}
#[test]
fn a_frames_length_accounts_for_every_piece() {
assert_eq!(framed_len(0), 36);
assert_eq!(framed_len(100), 136);
assert_eq!(framed_len(usize::MAX), u64::MAX);
}
#[test]
fn a_frames_bytes_are_prefix_payload_and_head() -> TestResult {
for size in [1usize, 64, 4096] {
let payload = vec![7u8; size];
let bytes = framed(&payload)?;
assert_eq!(bytes.len(), LENGTH_BYTES + size + HEAD_BYTES);
assert_eq!(&bytes[LENGTH_BYTES..LENGTH_BYTES + size], &payload[..]);
assert_eq!(
&bytes[bytes.len() - HEAD_BYTES..],
head().as_bytes(),
"the stored head is the frame's last {HEAD_BYTES} bytes"
);
}
Ok(())
}
fn chained(previous: &Digest, payload: &[u8]) -> Option<Digest> {
let mut hasher = lgwks_std::hash::Hasher::new();
hasher.update(previous.as_bytes());
hasher.update(payload);
Some(hasher.finalize())
}
fn chain_of(payloads: &[&[u8]]) -> Result<(Vec<u8>, Vec<Digest>), Box<dyn std::error::Error>> {
let mut bytes = Vec::new();
let mut previous = Digest::from_bytes([0u8; HEAD_BYTES]);
let mut chained_from = Vec::new();
for payload in payloads {
let head = chained(&previous, payload).ok_or("chained refused to hash a head")?;
let length =
writable_length(payload.len(), CEILING).ok_or("the grammar framed this payload")?;
bytes.extend_from_slice(&encode(length, payload, &head));
chained_from.push(previous);
previous = head;
}
Ok((bytes, chained_from))
}
const CEILING: usize = 64 * 1024;
#[test]
fn the_cut_frame_authenticates_only_when_it_is_whole_behind_a_lying_prefix() -> TestResult {
let payload: Vec<u8> = (0u8..100).collect();
let (bytes, from) = chain_of(&[&payload])?;
let previous = from[0];
let behind_prefix = &bytes[LENGTH_BYTES..];
assert!(super::holds_acknowledged_frame(
behind_prefix,
&previous,
CEILING,
chained
));
for cut in 0..behind_prefix.len() {
assert!(
!super::holds_acknowledged_frame(
&behind_prefix[..cut],
&previous,
CEILING,
chained
),
"a frame cut {cut} bytes in, short of its head, is an interrupted append"
);
}
Ok(())
}
#[test]
fn a_whole_payload_with_a_short_head_authenticates_nothing() -> TestResult {
let payload = vec![9u8; 64];
let (bytes, from) = chain_of(&[&payload])?;
let behind_prefix = &bytes[LENGTH_BYTES..];
for missing in 1..=HEAD_BYTES {
let cut = behind_prefix.len() - missing;
assert!(!super::holds_acknowledged_frame(
&behind_prefix[..cut],
&from[0],
CEILING,
chained
));
}
Ok(())
}
#[test]
fn a_later_frame_authenticates_when_the_cut_frame_cannot() -> TestResult {
let (bytes, from) = chain_of(&[&[1u8; 40], &[2u8; 50], &[3u8; 60]])?;
let mut behind_prefix = bytes[LENGTH_BYTES..].to_vec();
behind_prefix[5] ^= 0x80;
assert!(super::holds_acknowledged_frame(
&behind_prefix,
&from[0],
CEILING,
chained
));
let first_only = &behind_prefix[..40 + HEAD_BYTES];
assert!(!super::holds_acknowledged_frame(
first_only, &from[0], CEILING, chained
));
Ok(())
}
#[test]
fn a_frame_behind_a_damaged_head_chains_from_the_head_its_payload_implies() -> TestResult {
let (bytes, from) = chain_of(&[&[1u8; 40], &[2u8; 50]])?;
let mut behind_prefix = bytes[LENGTH_BYTES..].to_vec();
behind_prefix[40 + 7] ^= 0x01;
assert!(super::holds_acknowledged_frame(
&behind_prefix,
&from[0],
CEILING,
chained
));
Ok(())
}
#[test]
fn noise_and_the_longest_possible_tail_hold_nothing() -> TestResult {
let previous = Digest::from_bytes([5u8; HEAD_BYTES]);
assert!(!super::holds_acknowledged_frame(
&[],
&previous,
CEILING,
chained
));
let longest: Vec<u8> = (0..CEILING + HEAD_BYTES)
.map(|index| u8::saturating_from(index % 251))
.collect();
assert!(!super::holds_acknowledged_frame(
&longest, &previous, CEILING, chained
));
let (bytes, from) = chain_of(&[&[7u8; 30]])?;
assert!(!super::holds_acknowledged_frame(
&bytes[LENGTH_BYTES..],
&from[0],
CEILING,
|_previous, _payload| None
));
Ok(())
}
}