use std::fmt;
use std::io;
use std::path::Path;
use std::time::SystemTime;
use lgwks_std::hash::{Digest, blake3};
pub const MAX_STABILITY_READS: u32 = 8;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
#[non_exhaustive]
pub enum Drift {
Length,
Modified,
Contents,
Truncated,
}
impl Drift {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Length => "length",
Self::Modified => "modified",
Self::Contents => "contents",
Self::Truncated => "truncated",
}
}
}
impl fmt::Display for Drift {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str(self.as_str())
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Reading {
bytes: Vec<u8>,
reported: Option<u64>,
modified: Option<SystemTime>,
digest: Digest,
}
impl Reading {
#[must_use]
pub fn of_bytes(bytes: impl Into<Vec<u8>>) -> Self {
let bytes = bytes.into();
let digest = blake3(&bytes);
Self {
bytes,
reported: None,
modified: None,
digest,
}
}
#[must_use]
pub fn of_subject(bytes: impl Into<Vec<u8>>, reported: u64, modified: SystemTime) -> Self {
let mut reading = Self::of_bytes(bytes);
reading.reported = Some(reported);
reading.modified = Some(modified);
reading
}
#[must_use]
pub fn bytes(&self) -> &[u8] {
&self.bytes
}
#[must_use]
pub fn reported(&self) -> Option<u64> {
self.reported
}
#[must_use]
pub fn modified(&self) -> Option<SystemTime> {
self.modified
}
#[must_use]
pub const fn digest(&self) -> Digest {
self.digest
}
#[must_use]
pub fn is_whole(&self) -> bool {
self.reported.is_none_or(|reported| {
u64::try_from(self.bytes.len()).is_ok_and(|read| read == reported)
})
}
#[must_use]
pub fn agrees_with(&self, other: &Reading) -> Option<Drift> {
if !self.is_whole() || !other.is_whole() {
return Some(Drift::Truncated);
}
if self.reported != other.reported {
return Some(Drift::Length);
}
if self.modified != other.modified {
return Some(Drift::Modified);
}
if self.digest != other.digest {
return Some(Drift::Contents);
}
None
}
}
#[derive(Debug)]
#[non_exhaustive]
pub enum ReadFailure {
Unstable(Unstable),
Unreadable {
cause: Box<dyn std::error::Error + Send + Sync>,
},
}
impl ReadFailure {
#[must_use]
pub const fn unstable(&self) -> Option<&Unstable> {
match *self {
Self::Unstable(ref unstable) => Some(unstable),
Self::Unreadable { .. } => None,
}
}
}
impl fmt::Display for ReadFailure {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match *self {
Self::Unstable(unstable) => write!(formatter, "{unstable}"),
Self::Unreadable { ref cause } => {
write!(formatter, "the subject could not be read: {cause}")
}
}
}
}
impl std::error::Error for ReadFailure {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match *self {
Self::Unstable(_) => None,
Self::Unreadable { ref cause } => Some(cause.as_ref()),
}
}
}
impl From<io::Error> for ReadFailure {
fn from(cause: io::Error) -> Self {
Self::Unreadable {
cause: Box::new(cause),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Unstable {
reads: u32,
drift: Drift,
held: Digest,
observed: Digest,
}
impl Unstable {
#[must_use]
pub const fn reads(&self) -> u32 {
self.reads
}
#[must_use]
pub const fn drift(&self) -> Drift {
self.drift
}
#[must_use]
pub const fn held(&self) -> Digest {
self.held
}
#[must_use]
pub const fn observed(&self) -> Digest {
self.observed
}
}
impl fmt::Display for Unstable {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(
formatter,
"the subject did not settle: {} moved across {} read(s)",
self.drift, self.reads
)
}
}
pub fn read_stable<R>(mut read: R) -> Result<Reading, ReadFailure>
where
R: FnMut() -> Result<Reading, ReadFailure>,
{
let mut held = read()?;
let mut reads: u32 = 1;
let mut last_move: Option<(Drift, Digest, Digest)> = None;
while reads < MAX_STABILITY_READS {
let observed = read()?;
reads = reads.saturating_add(1);
match held.agrees_with(&observed) {
None => return Ok(observed),
Some(axis) => {
last_move = Some((axis, held.digest(), observed.digest()));
held = observed;
}
}
}
let refusal = match last_move {
Some((axis, before, after)) => Err(ReadFailure::Unstable(Unstable {
reads,
drift: axis,
held: before,
observed: after,
})),
None => Err(ReadFailure::Unstable(Unstable {
reads,
drift: Drift::Contents,
held: held.digest(),
observed: held.digest(),
})),
};
lgwks_std::trace::debug!(
error = ?refusal.as_ref().err(),
"read_stable: returning an error to the caller"
);
refusal
}
pub async fn read_stable_file(path: &Path) -> Result<Reading, ReadFailure> {
let owned = path.to_path_buf();
lgwks_std::task::spawn_blocking(move || settle_file(&owned)).await
}
fn read_file(path: &Path) -> Result<Reading, ReadFailure> {
let bytes = std::fs::read(path)?;
let metadata = std::fs::metadata(path)?;
let mut reading = Reading::of_bytes(bytes);
reading.reported = Some(metadata.len());
reading.modified = metadata.modified().ok();
Ok(reading)
}
fn settle_file(path: &Path) -> Result<Reading, ReadFailure> {
read_stable(|| read_file(path))
}
#[cfg(test)]
mod tests {
use super::{Drift, MAX_STABILITY_READS, ReadFailure, Reading, read_stable, read_stable_file};
use std::error::Error;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::SystemTime;
type TestResult = Result<(), Box<dyn Error>>;
fn failed(cause: impl Into<String>) -> Box<dyn Error> {
Box::new(std::io::Error::other(cause.into()))
}
fn epoch_nanos() -> u128 {
match std::time::SystemTime::now().duration_since(SystemTime::UNIX_EPOCH) {
Ok(elapsed) => elapsed.as_nanos(),
Err(_) => 0,
}
}
fn stamp() -> SystemTime {
SystemTime::UNIX_EPOCH
}
struct Script {
readings: Vec<Reading>,
cursor: AtomicUsize,
}
impl Script {
fn new(readings: Vec<Reading>) -> Self {
Self {
readings,
cursor: AtomicUsize::new(0),
}
}
fn next(&self) -> Result<Reading, ReadFailure> {
let index = self.cursor.fetch_add(1, Ordering::Relaxed);
match self.readings.get(index) {
Some(reading) => Ok(reading.clone()),
None => Err(ReadFailure::Unreadable {
cause: Box::new(std::io::Error::other("the script has no more readings")),
}),
}
}
}
#[test]
fn two_agreeing_readings_settle() -> TestResult {
let script = Script::new(vec![
Reading::of_subject(b"{\"a\":1}".to_vec(), 7, stamp()),
Reading::of_subject(b"{\"a\":1}".to_vec(), 7, stamp()),
Reading::of_subject(b"{\"a\":2}".to_vec(), 7, stamp()),
]);
let settled = read_stable(|| script.next())?;
assert_eq!(
settled.bytes(),
b"{\"a\":1}",
"the settled value must be the one two reads agreed on, not a later one"
);
Ok(())
}
#[test]
fn every_drift_axis_is_reported_by_name() -> TestResult {
let base = Reading::of_subject(b"12345".to_vec(), 5, stamp());
let cases: [(Reading, Drift, &str); 4] = [
(
Reading::of_subject(b"123456".to_vec(), 6, stamp()),
Drift::Length,
"a subject that grew",
),
(
Reading::of_subject(
b"12345".to_vec(),
5,
stamp()
.checked_add(std::time::Duration::from_secs(1))
.ok_or_else(|| failed("one second after the epoch is representable"))?,
),
Drift::Modified,
"a subject that was touched",
),
(
Reading::of_subject(b"12435".to_vec(), 5, stamp()),
Drift::Contents,
"a subject rewritten in place to the same length",
),
(
Reading::of_subject(b"123".to_vec(), 5, stamp()),
Drift::Truncated,
"a read cut short by a concurrent writer",
),
];
for (other, drift, why) in cases {
assert_eq!(
base.agrees_with(&other),
Some(drift),
"{why}: the axis that moved is what a caller triages on"
);
}
assert_eq!(Drift::Contents.as_str(), "contents");
assert_eq!(
base.agrees_with(&base),
None,
"a reading agrees with itself"
);
Ok(())
}
#[test]
fn a_subject_that_never_settles_is_refused_at_the_bound() -> TestResult {
let width = usize::try_from(MAX_STABILITY_READS)?;
let mut drawn = Vec::new();
for index in 0..=MAX_STABILITY_READS {
drawn.push(Reading::of_subject(
format!("{index:0>width$}").into_bytes(),
u64::from(MAX_STABILITY_READS),
stamp(),
));
}
let script = Script::new(drawn);
let outcome = read_stable(|| script.next());
let Err(ReadFailure::Unstable(unstable)) = outcome else {
return Err(failed(format!(
"a subject whose bytes change on every read must be refused: {outcome:?}"
)));
};
assert_eq!(
unstable.reads(),
MAX_STABILITY_READS,
"the attempt must stop at the declared bound rather than loop"
);
assert_eq!(
unstable.drift(),
Drift::Contents,
"same length, same timestamp, different bytes: that is the contents axis"
);
assert_ne!(
unstable.held(),
unstable.observed(),
"the refusal carries both digests, so a report can name what it saw"
);
Ok(())
}
#[test]
fn an_unreadable_subject_is_its_own_arm() -> TestResult {
let outcome = read_stable(|| {
Err::<Reading, _>(ReadFailure::Unreadable {
cause: Box::new(std::io::Error::from(std::io::ErrorKind::PermissionDenied)),
})
});
let Some(failure) = outcome.as_ref().err() else {
return Err(failed("an unreadable subject must be refused"));
};
assert!(
failure.unstable().is_none(),
"a subject that could not be read has no reading to report: {failure}"
);
assert!(
failure.source().is_some(),
"the platform's own error is retained as the cause: {failure}"
);
Ok(())
}
fn owned_dir(tag: &str) -> Result<std::path::PathBuf, std::io::Error> {
let base = std::env::temp_dir();
for attempt in 0..32u32 {
let candidate = base.join(format!("lgwks-stability-{tag}-{}-{attempt}", epoch_nanos()));
match std::fs::create_dir(&candidate) {
Ok(()) => return Ok(candidate),
Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {}
Err(error) => {
let refusal = Err::<std::path::PathBuf, std::io::Error>(error);
lgwks_std::trace::debug!(
error = ?refusal.as_ref().err(),
"owned_dir: returning an error to the caller"
);
return refusal;
}
}
}
let refusal = Err::<std::path::PathBuf, std::io::Error>(std::io::Error::other(
"no unused scratch directory name after 32 attempts",
));
lgwks_std::trace::debug!(
error = ?refusal.as_ref().err(),
"owned_dir: returning an error to the caller"
);
refusal
}
#[test]
fn a_real_file_is_read_stably_and_a_missing_one_is_unreadable() -> TestResult {
let directory = owned_dir("file")?;
let path = directory.join("subject.json");
std::fs::write(&path, b"{\"a\":1}\n")?;
let settled = lgwks_std::task::block_on(read_stable_file(&path))?;
assert_eq!(
settled.bytes(),
b"{\"a\":1}\n",
"a file nobody is writing reads the bytes it holds"
);
assert!(
settled.is_whole(),
"the reading must carry the subject's own length, so a cut read is visible"
);
let absent = path.with_extension("absent");
let outcome = lgwks_std::task::block_on(read_stable_file(&absent));
assert!(
matches!(outcome, Err(ReadFailure::Unreadable { .. })),
"a file that does not exist is unreadable, never unstable: {outcome:?}"
);
drop(std::fs::remove_dir_all(&directory));
Ok(())
}
}