use std::collections::HashMap;
use std::fs::{File, OpenOptions};
use std::io::{BufReader, Read, Write};
use std::path::{Path, PathBuf};
use super::continuation::{
Continuation, ContinuationPolicy, ContinuationWatermark, SealPause, SettledAttempt,
UnresolvedAttempt, successor_path,
};
use super::frame::{HEAD_BYTES, LENGTH_BYTES, Piece, Prefix, SaturatingFrom, read_exact_or_eof};
use super::owner::{StorageGate, StorageOwner};
use super::{
AttemptStatus, ChainBreak, DurabilityPromise, DurableAck, EffectEvent, EffectEvidence,
EffectJournal, EventKind, JournalEntry, JournalError, JournalLimitKind, JournalPosition,
MAX_JOURNAL_BYTES, MAX_JOURNAL_EVENTS, Recovered, check_append_order, next_allowed_of,
recover_continued,
};
use lgwks_std::wire::{WireError, from_bytes};
pub(super) const MAX_FRAME_BYTES: usize = 64 * 1024;
pub const MAX_GENERATION_WALK: u64 = 4096;
type FramePiece = Piece;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum CorruptionKind {
Undecodable,
Framed,
Chain(ChainBreak),
Sealed,
}
impl CorruptionKind {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Undecodable => "the frame's bytes do not decode to an event",
Self::Framed => {
"the frame's length prefix cannot be true of any frame this journal writes"
}
Self::Chain(_) => "the stored head does not follow from the events before it",
Self::Sealed => "the file holds a second sealed checkpoint",
}
}
}
impl core::fmt::Display for CorruptionKind {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
f.write_str(self.as_str())
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Corruption {
at: u64,
kind: CorruptionKind,
}
impl Corruption {
#[must_use]
pub const fn new(at: u64, kind: CorruptionKind) -> Self {
Self { at, kind }
}
#[must_use]
pub const fn at(&self) -> u64 {
self.at
}
#[must_use]
pub const fn kind(&self) -> &CorruptionKind {
&self.kind
}
}
impl core::fmt::Display for Corruption {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
write!(f, "frame {} corrupt: {}", self.at, self.kind)
}
}
impl std::error::Error for Corruption {}
#[derive(Debug)]
enum ScanStop {
Complete(u64),
Torn(u64),
AmbiguousTail {
offset: u64,
},
}
fn read_exact_classified(
reader: &mut impl Read,
buf: &mut [u8],
) -> Result<FramePiece, JournalError> {
super::frame::read_piece(reader, buf).map_err(JournalError::Storage)
}
fn resolve_ambiguous_tail(
file: &mut File,
offset: u64,
position: JournalPosition,
index: u64,
) -> Result<u64, JournalError> {
let acknowledged = super::frame::cut_holds_acknowledged(
file,
offset,
&position.head(),
MAX_FRAME_BYTES,
JournalError::Storage,
|previous, payload| Some(super::chain_over_bytes(previous, payload)),
)?;
if acknowledged {
let refusal = Err(JournalError::Corrupt(Box::new(Corruption::new(
index,
CorruptionKind::Framed,
))));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "resolve_ambiguous_tail: the tail holds an acknowledged frame under a lying length");
return refusal;
}
Ok(offset)
}
struct Frame {
payload: Vec<u8>,
head: [u8; HEAD_BYTES],
payload_len: usize,
}
#[derive(Clone, Copy)]
enum Halt {
Complete,
Torn,
Ambiguous,
}
impl Halt {
const fn at(self, offset: u64) -> ScanStop {
match self {
Self::Complete => ScanStop::Complete(offset),
Self::Torn => ScanStop::Torn(offset),
Self::Ambiguous => ScanStop::AmbiguousTail { offset },
}
}
}
fn next_frame(
reader: &mut impl Read,
held: usize,
max_events: usize,
index: u64,
) -> Result<Result<Frame, Halt>, JournalError> {
let mut prefix = [0u8; LENGTH_BYTES];
match super::frame::read_prefix(reader, &mut prefix).map_err(JournalError::Storage)? {
Prefix::Eof => return Ok(Err(Halt::Complete)),
Prefix::Torn => return Ok(Err(Halt::Torn)),
Prefix::Full => {
let requested = u64::saturating_from(held).saturating_add(1);
if requested > u64::saturating_from(max_events) {
let refusal = Err(JournalError::CapacityExceeded {
resource: JournalLimitKind::Events,
limit: u64::saturating_from(max_events),
requested,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "next_frame: the file holds more events than the ceiling");
return refusal;
}
}
}
let payload_len = super::frame::declared_length(&prefix);
if !super::frame::is_possible_length(payload_len, MAX_FRAME_BYTES) {
let refusal = Err(JournalError::Corrupt(Box::new(Corruption::new(
index,
CorruptionKind::Framed,
))));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), payload_len, "next_frame: the declared length is one this journal never writes");
return refusal;
}
let mut payload = vec![0u8; payload_len];
if let FramePiece::Interrupted = read_exact_classified(reader, &mut payload)? {
return Ok(Err(Halt::Ambiguous));
}
let mut head = [0u8; HEAD_BYTES];
if let FramePiece::Interrupted = read_exact_classified(reader, &mut head)? {
return Ok(Err(Halt::Ambiguous));
}
Ok(Ok(Frame {
payload,
head,
payload_len,
}))
}
fn scan(reader: &mut impl Read, max_events: usize) -> Result<Scanned, JournalError> {
let mut entries = Vec::new();
let mut carried: Option<Continuation> = None;
let mut seal: Option<Continuation> = None;
let mut position = JournalPosition::genesis();
let mut chain_from = JournalPosition::genesis();
let mut offset = 0u64;
let mut index = 0u64;
loop {
let Frame {
payload,
head,
payload_len,
} = match next_frame(reader, entries.len(), max_events, index)? {
Ok(frame) => frame,
Err(halt) => {
return Ok(Scanned {
entries,
carried,
seal,
chain_from,
tail: position,
stop: halt.at(offset),
});
}
};
if let Some(decoded) = Continuation::from_payload(&payload) {
if seal.is_some() {
let refusal = Err(JournalError::Corrupt(Box::new(Corruption::new(
index,
CorruptionKind::Sealed,
))));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "scan: the file holds a second sealed checkpoint");
return refusal;
}
let checkpoint = match decoded {
Ok(checkpoint) => checkpoint,
Err(cause) => {
lgwks_std::trace::debug!(index, "scan: the sealed checkpoint did not decode");
return Err(cause);
}
};
position = verify_frame(&payload, &head, checkpoint.predecessor(), index)?;
if index == 0 {
carried = Some(checkpoint);
chain_from = position;
} else {
seal = Some(checkpoint);
}
} else {
if seal.is_some() {
let refusal = Err(JournalError::Corrupt(Box::new(Corruption::new(
index,
CorruptionKind::Sealed,
))));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "scan: an event follows this journal's own seal");
return refusal;
}
let event: EffectEvent = match from_bytes::<EffectEvent, WireError>(&payload) {
Ok(event) => event,
Err(error) => {
lgwks_std::trace::debug!(?error, index, "scan: the payload did not decode");
let refusal = Err(JournalError::Corrupt(Box::new(Corruption::new(
index,
CorruptionKind::Undecodable,
))));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "scan: returning an error to the caller");
return refusal;
}
};
position = verify_frame(&payload, &head, position, index)?;
entries.push(JournalEntry::new(position, event));
}
offset = offset.saturating_add(super::frame::framed_len(payload_len));
index = index.saturating_add(1);
}
}
fn verify_frame(
payload: &[u8],
head: &[u8; HEAD_BYTES],
position: JournalPosition,
index: u64,
) -> Result<JournalPosition, JournalError> {
let recomputed = JournalPosition {
sequence: position.sequence().saturating_add(1),
head: super::chain_over_bytes(&position.head(), payload),
};
let recorded = JournalPosition {
sequence: recomputed.sequence(),
head: lgwks_std::hash::Digest::from_bytes(*head),
};
if recorded != recomputed {
let refusal = Err(JournalError::Corrupt(Box::new(Corruption::new(
index,
CorruptionKind::Chain(ChainBreak::Disagreement {
at: recorded.sequence(),
recorded,
recomputed,
}),
))));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "verify_frame: returning an error to the caller");
return refusal;
}
Ok(recorded)
}
struct Scanned {
entries: Vec<JournalEntry>,
carried: Option<Continuation>,
seal: Option<Continuation>,
chain_from: JournalPosition,
tail: JournalPosition,
stop: ScanStop,
}
pub struct FileJournal {
path: PathBuf,
storage: StorageOwner<(), SealOutcome>,
view: FileView,
committed: Vec<JournalEntry>,
ladder: HashMap<crate::effect::EffectKey, EventKind>,
outcomes: HashMap<crate::effect::EffectKey, (JournalPosition, EffectEvidence)>,
disk_len: u64,
position: JournalPosition,
torn_tail_repaired: bool,
continuing: bool,
continuation: ContinuationPolicy,
carried: Option<Continuation>,
seal: Option<Continuation>,
chain_from: JournalPosition,
base: JournalPosition,
folded: Vec<SettledAttempt>,
sealed_by: Option<PathBuf>,
pause: Option<SealPause>,
}
impl core::fmt::Debug for FileJournal {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
f.debug_struct("FileJournal")
.field("path", &self.path)
.field("committed_events", &self.committed.len())
.field("disk_len", &self.disk_len)
.field("torn_tail_repaired", &self.torn_tail_repaired)
.field("continuing", &self.continuing)
.field("generation", &self.generation())
.field("sealed_by", &self.sealed_by)
.finish()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum SealOutcome {
Appended,
Sealed {
bytes: usize,
},
Paused(SealPause),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum OpenKind {
Plain,
Stalled,
Continuing(ContinuationPolicy),
}
impl OpenKind {
const fn stalled(self) -> bool {
matches!(self, Self::Stalled)
}
const fn continues(self) -> bool {
matches!(self, Self::Continuing(_))
}
const fn policy(self) -> ContinuationPolicy {
match self {
Self::Continuing(policy) => policy,
Self::Plain | Self::Stalled => ContinuationPolicy::declared(),
}
}
}
fn ladder_hold(
ladder: &mut HashMap<crate::effect::EffectKey, EventKind>,
key: crate::effect::EffectKey,
rung: EventKind,
) {
ladder.insert(key, rung);
}
fn successor_is_complete(successor: &Path, seal: &Continuation) -> bool {
let Ok(mut file) = OpenOptions::new().read(true).open(successor) else {
return false;
};
let mut prefix = [0u8; LENGTH_BYTES];
let Ok(Prefix::Full) = super::frame::read_prefix(&mut file, &mut prefix) else {
return false;
};
let declared = super::frame::declared_length(&prefix);
if !super::frame::is_possible_length(declared, MAX_FRAME_BYTES) {
return false;
}
let mut payload = vec![0u8; declared];
let mut head = [0u8; HEAD_BYTES];
if super::frame::read_piece(&mut file, &mut payload).is_err()
|| super::frame::read_piece(&mut file, &mut head).is_err()
{
return false;
}
match Continuation::from_payload(&payload) {
Some(Ok(carried)) => {
carried.predecessor() == seal.predecessor()
&& super::chain_over_bytes(&seal.predecessor().head(), &payload)
== lgwks_std::hash::Digest::from_bytes(head)
}
Some(Err(_)) | None => false,
}
}
#[cfg(unix)]
fn sync_directory(path: &Path, created_here: bool) -> std::io::Result<()> {
let Some(parent) = path.parent() else {
return Ok(());
};
File::open(parent)?.sync_all()?;
if created_here && let Some(grandparent) = parent.parent() {
File::open(grandparent)?.sync_all()?;
}
Ok(())
}
#[cfg(not(unix))]
fn sync_directory(_path: &Path, _created_here: bool) -> std::io::Result<()> {
Ok(())
}
fn write_successor(successor: &Path, frame: &[u8]) -> std::io::Result<bool> {
let created = successor.parent().is_some_and(|parent| !parent.exists());
if let Some(parent) = successor.parent() {
std::fs::create_dir_all(parent)?;
}
let written = match OpenOptions::new()
.write(true)
.create_new(true)
.open(successor)
{
Ok(mut file) => file.write_all(frame),
Err(ref error) if error.kind() == std::io::ErrorKind::AlreadyExists => {
let file = OpenOptions::new().read(true).write(true).open(successor)?;
file.set_len(0)?;
let mut handle = file;
handle.write_all(frame)
}
Err(error) => {
let refusal = Err(error);
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "write_successor: the successor could not be created");
return refusal;
}
};
written.map(|()| created)
}
fn seal_on_owner(
file: &mut File,
expected_len: u64,
frame: &[u8],
successor: &Path,
pause: Option<SealPause>,
) -> std::io::Result<super::owner::Stage<SealOutcome, ()>> {
fence_and_write(file, expected_len, frame)?;
let Some(paused) = paused_at(pause, SealPause::AfterSealWrite) else {
let created_directory = write_successor(successor, frame)?;
let Some(paused) = paused_at(pause, SealPause::AfterSuccessorWrite) else {
OpenOptions::new().read(true).open(successor)?.sync_all()?;
let Some(paused) = paused_at(pause, SealPause::AfterSuccessorSync) else {
sync_directory(successor, created_directory)?;
let Some(paused) = paused_at(pause, SealPause::AfterDirectorySync) else {
return Ok(super::owner::Stage::Unsynced {
answer: SealOutcome::Sealed { bytes: frame.len() },
bytes: frame.len(),
settle: Box::new(|_: &mut ()| {}),
});
};
return Ok(super::owner::Stage::Settled(Ok(paused)));
};
return Ok(super::owner::Stage::Settled(Ok(paused)));
};
return Ok(super::owner::Stage::Settled(Ok(paused)));
};
Ok(super::owner::Stage::Settled(Ok(paused)))
}
fn paused_at(pause: Option<SealPause>, boundary: SealPause) -> Option<SealOutcome> {
match pause {
Some(armed) if armed == boundary => Some(SealOutcome::Paused(boundary)),
_ => None,
}
}
struct FileView {
file: File,
}
impl FileView {
fn read_only(path: &Path) -> Result<Self, JournalError> {
Ok(Self {
file: OpenOptions::new()
.read(true)
.open(path)
.map_err(JournalError::Storage)?,
})
}
fn len(&self) -> Result<u64, JournalError> {
self.file
.metadata()
.map_err(JournalError::Storage)
.map(|meta| meta.len())
}
}
pub struct Replay {
reader: BufReader<File>,
position: JournalPosition,
yielded: u64,
frames: u64,
sealed: bool,
offset: u64,
done: bool,
}
impl Replay {
fn open(path: &Path) -> Result<Self, JournalError> {
let mut file = OpenOptions::new()
.read(true)
.open(path)
.map_err(JournalError::Storage)?;
file.seek_read_zero()?;
Ok(Self {
reader: BufReader::new(file),
position: JournalPosition::genesis(),
yielded: 0,
frames: 0,
sealed: false,
offset: 0,
done: false,
})
}
fn read_one(&mut self) -> Option<Result<EffectEvent, JournalError>> {
loop {
match self.read_frame()? {
Ok(Some(event)) => return Some(Ok(event)),
Ok(None) => {}
Err(error) => return Some(Err(error)),
}
}
}
fn read_frame(&mut self) -> Option<Result<Option<EffectEvent>, JournalError>> {
let frame = self.frames;
let mut prefix = [0u8; LENGTH_BYTES];
match read_exact_or_eof(&mut self.reader, &mut prefix) {
Err(error) => return Some(Err(error)),
Ok(None) => return None,
Ok(Some(read_len)) if read_len < LENGTH_BYTES => return None,
Ok(Some(_)) => {}
}
let limit = u64::saturating_from(MAX_JOURNAL_EVENTS);
if self.yielded >= limit {
let refusal = Err(JournalError::CapacityExceeded {
resource: JournalLimitKind::Events,
limit,
requested: self.yielded.saturating_add(1),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "Replay::read_frame: the stream reached the event ceiling");
return Some(refusal);
}
let payload_len = usize::saturating_from(u32::from_be_bytes(prefix));
if payload_len == 0 || payload_len > MAX_FRAME_BYTES {
let refusal = Err(JournalError::Corrupt(Box::new(Corruption::new(
frame,
CorruptionKind::Framed,
))));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "Replay::read_frame: a length no writer produces");
return Some(refusal);
}
let mut payload = vec![0u8; payload_len];
match read_exact_classified(&mut self.reader, &mut payload) {
Ok(FramePiece::Interrupted) => return self.end_of_acknowledged_prefix(),
Err(error) => return Some(Err(error)),
Ok(FramePiece::Filled) => {}
}
let mut head = [0u8; HEAD_BYTES];
match read_exact_classified(&mut self.reader, &mut head) {
Ok(FramePiece::Interrupted) => return self.end_of_acknowledged_prefix(),
Err(error) => return Some(Err(error)),
Ok(FramePiece::Filled) => {}
}
let admitted = self.admit(&payload, &head, frame);
self.frames = self.frames.saturating_add(1);
self.offset = self
.offset
.saturating_add(super::frame::framed_len(payload_len));
Some(admitted)
}
fn admit(
&mut self,
payload: &[u8],
head: &[u8; HEAD_BYTES],
frame: u64,
) -> Result<Option<EffectEvent>, JournalError> {
if self.sealed {
let refusal = Err(JournalError::Corrupt(Box::new(Corruption::new(
frame,
CorruptionKind::Sealed,
))));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "Replay::admit: a frame follows this journal's own seal");
return refusal;
}
if let Some(decoded) = Continuation::from_payload(payload) {
let checkpoint = decoded?;
let sealed_at = verify_frame(payload, head, checkpoint.predecessor(), frame)?;
self.position = sealed_at;
self.sealed = frame != 0;
return Ok(None);
}
let Ok(event) = from_bytes::<EffectEvent, WireError>(payload) else {
let refusal = Err(JournalError::Corrupt(Box::new(Corruption::new(
frame,
CorruptionKind::Undecodable,
))));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "Replay::admit: the payload did not decode");
return refusal;
};
self.position = verify_frame(payload, head, self.position, frame)?;
self.yielded = self.yielded.saturating_add(1);
Ok(Some(event))
}
fn end_of_acknowledged_prefix(&mut self) -> Option<Result<Option<EffectEvent>, JournalError>> {
match resolve_ambiguous_tail(
self.reader.get_mut(),
self.offset,
self.position,
self.yielded,
) {
Ok(_) => None,
Err(error) => Some(Err(error)),
}
}
}
impl Iterator for Replay {
type Item = Result<EffectEvent, JournalError>;
fn next(&mut self) -> Option<Self::Item> {
if self.done {
return None;
}
let item = self.read_one();
if item.as_ref().is_none_or(Result::is_err) {
self.done = true;
}
item
}
}
impl FileJournal {
pub fn open(path: impl AsRef<Path>) -> Result<Self, JournalError> {
Self::open_impl(path.as_ref(), OpenKind::Plain)
}
pub fn open_continuing(path: impl AsRef<Path>) -> Result<Self, JournalError> {
Self::open_impl(
path.as_ref(),
OpenKind::Continuing(ContinuationPolicy::declared()),
)
}
pub fn open_continuing_with(
path: impl AsRef<Path>,
policy: ContinuationPolicy,
) -> Result<Self, JournalError> {
Self::open_impl(path.as_ref(), OpenKind::Continuing(policy))
}
pub fn open_active(path: impl AsRef<Path>) -> Result<Self, JournalError> {
let mut here = path.as_ref().to_path_buf();
for _ in 0..MAX_GENERATION_WALK {
match Self::open_impl(&here, OpenKind::Continuing(ContinuationPolicy::declared())) {
Ok(journal) => return Ok(journal),
Err(JournalError::Superseded { path: next }) => here = PathBuf::from(next),
Err(other) => {
let refusal = Err(other);
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "FileJournal::open_active: a generation on the walk refused to open");
return refusal;
}
}
}
let next = successor_path(&here).display().to_string();
let refusal = Err(JournalError::Superseded { path: next });
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "open_active: the generation walk ran out before reaching a live journal");
refusal
}
pub fn arm_continuation_pause(&mut self, boundary: SealPause) {
self.pause = Some(boundary);
}
#[must_use]
pub fn generation(&self) -> u64 {
super::continuation::generation_of(&self.path)
}
#[must_use]
pub fn checkpoint(&self) -> Option<&Continuation> {
self.carried.as_ref()
}
#[must_use]
pub fn own_seal(&self) -> Option<&Continuation> {
self.seal.as_ref()
}
#[must_use]
pub const fn base(&self) -> JournalPosition {
self.base
}
#[must_use]
pub const fn chain_from(&self) -> JournalPosition {
self.chain_from
}
#[must_use]
pub fn sealed_by(&self) -> Option<&Path> {
self.sealed_by.as_deref()
}
#[must_use]
pub fn successor_path(&self) -> PathBuf {
successor_path(&self.path)
}
fn open_impl(path: &Path, kind: OpenKind) -> Result<Self, JournalError> {
let path = path.to_path_buf();
super::continuation::refuse_ambiguous_base(&path)?;
let mut file = OpenOptions::new()
.read(true)
.append(true)
.create(true)
.open(&path)
.map_err(JournalError::Storage)?;
file.try_lock().map_err(|error| match error {
std::fs::TryLockError::WouldBlock => JournalError::Locked {
path: path.display().to_string(),
},
std::fs::TryLockError::Error(io) => JournalError::Storage(io),
})?;
let file_len = file.metadata().map_err(JournalError::Storage)?.len();
if file_len > MAX_JOURNAL_BYTES {
let refusal = Err(JournalError::CapacityExceeded {
resource: JournalLimitKind::Bytes,
limit: MAX_JOURNAL_BYTES,
requested: file_len,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "open_impl: returning an error to the caller");
return refusal;
}
file.seek_read_zero()?;
let mut reader = BufReader::new(&mut file);
let scanned = scan(&mut reader, MAX_JOURNAL_EVENTS)?;
drop(reader);
let Scanned {
entries,
carried,
seal,
chain_from,
tail,
stop,
} = scanned;
if let Some(checkpoint) = seal.as_ref() {
let next = successor_path(&path);
if successor_is_complete(&next, checkpoint) {
let refusal = Err(JournalError::Superseded {
path: next.display().to_string(),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "open_impl: this journal is sealed and its successor is authoritative");
return refusal;
}
}
let (acked_len, torn_tail_repaired) = match stop {
ScanStop::Complete(len) => (len, false),
ScanStop::Torn(offset) => {
file.set_len(offset).map_err(JournalError::Storage)?;
file.sync_all().map_err(JournalError::Storage)?;
(offset, true)
}
ScanStop::AmbiguousTail { offset } => {
let index = u64::saturating_from(entries.len());
let offset = resolve_ambiguous_tail(&mut file, offset, tail, index)?;
file.set_len(offset).map_err(JournalError::Storage)?;
file.sync_all().map_err(JournalError::Storage)?;
(offset, true)
}
};
let mut ladder: HashMap<crate::effect::EffectKey, EventKind> = entries
.iter()
.map(|entry| (entry.event().key(), entry.event().kind()))
.collect();
let mut outcomes = HashMap::new();
for entry in &entries {
if let EffectEvent::OutcomeObserved { key, evidence } = *entry.event() {
outcomes.insert(key, (entry.position(), evidence));
}
}
let mut folded = Vec::new();
let mut base = JournalPosition::genesis();
if let Some(checkpoint) = carried.as_ref() {
base = checkpoint.predecessor();
folded = checkpoint.settled().to_vec();
for held in checkpoint.unresolved() {
if let Some(rung) = held.rung() {
ladder_hold(&mut ladder, held.key(), rung);
}
}
for held in checkpoint.settled() {
if let Some(rung) = held.rung() {
ladder_hold(&mut ladder, held.key(), rung);
}
}
}
let storage =
StorageOwner::spawn(file, (), kind.stalled()).map_err(JournalError::Storage)?;
let view = FileView::read_only(&path)?;
Ok(Self {
path,
storage,
view,
committed: entries,
position: tail,
ladder,
outcomes,
disk_len: acked_len,
torn_tail_repaired,
continuing: kind.continues(),
continuation: kind.policy(),
carried,
seal,
base,
chain_from,
folded,
sealed_by: None,
pause: None,
})
}
fn bound_events(&self, additional: u64) -> Result<(), JournalError> {
let requested = u64::saturating_from(self.committed.len()).saturating_add(additional);
let limit = u64::saturating_from(MAX_JOURNAL_EVENTS);
if requested > limit {
let refusal = Err(JournalError::CapacityExceeded {
resource: JournalLimitKind::Events,
limit,
requested,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "bound_events: returning an error to the caller");
return refusal;
}
Ok(())
}
#[must_use]
pub fn path(&self) -> &Path {
&self.path
}
#[must_use]
pub const fn torn_tail_repaired(&self) -> bool {
self.torn_tail_repaired
}
pub fn events(&self) -> impl Iterator<Item = &EffectEvent> {
self.committed.iter().map(JournalEntry::event)
}
pub fn replay(&self) -> Result<Replay, JournalError> {
Replay::open(&self.path)
}
#[must_use]
pub fn recover(&self) -> Recovered {
recover_continued(self.carried.as_ref(), self.events())
}
fn fence(&self) -> Result<(), JournalError> {
self.refuse_sealed()?;
if self.storage.poisoned() {
let refusal = Err(JournalError::Storage(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"a previous append failed while writing; the file may hold \
unacknowledged bytes, reopen to replay",
)));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "fence: returning an error to the caller");
return refusal;
}
if self.view.len()? != self.disk_len {
let refusal = Err(JournalError::Storage(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"the journal file moved under this controller; reopen before appending",
)));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "fence: returning an error to the caller");
return refusal;
}
Ok(())
}
fn refuse_sealed(&self) -> Result<(), JournalError> {
let Some(next) = self.sealed_by.as_ref() else {
return Ok(());
};
let refusal = Err(JournalError::Superseded {
path: next.display().to_string(),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "refuse_sealed: this handle sealed its journal and is read-only");
refusal
}
fn frame(
&self,
event: &EffectEvent,
from: JournalPosition,
) -> Result<(JournalPosition, Vec<u8>), JournalError> {
let sequence = from
.sequence()
.checked_add(1)
.ok_or(JournalError::Exhausted)?;
let (framed, head) = super::frame::frame_record(
event,
&from.head,
MAX_FRAME_BYTES,
|event| {
event
.to_bytes()
.map(|bytes| bytes.as_ref().to_vec())
.map_err(JournalError::Encoding)
},
|_event, previous, archived| super::chain_over_bytes(previous, archived),
|_len| {
JournalError::Storage(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"the event exceeds this journal's frame bound",
))
},
)?;
Ok((JournalPosition { sequence, head }, framed))
}
fn bound_bytes(&self, staged: usize) -> Result<(), JournalError> {
let requested = self.disk_len.saturating_add(u64::saturating_from(staged));
if requested > MAX_JOURNAL_BYTES {
let refusal = Err(JournalError::CapacityExceeded {
resource: JournalLimitKind::Bytes,
limit: MAX_JOURNAL_BYTES,
requested,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "bound_bytes: returning an error to the caller");
return refusal;
}
Ok(())
}
fn write_and_sync(&mut self, frames: &[u8]) -> Result<(), JournalError> {
self.write_and_sync_blocking(frames)
}
fn continuation_continues(&self) -> bool {
self.continuing
}
fn successor(&self) -> PathBuf {
successor_path(&self.path)
}
fn carried(&self) -> (Vec<SettledAttempt>, Vec<UnresolvedAttempt>) {
let recovered = self.recover();
let mut settled: Vec<SettledAttempt> = Vec::new();
let mut unresolved: Vec<UnresolvedAttempt> = Vec::new();
let mut folded_at: HashMap<crate::effect::ActionId, usize> = HashMap::new();
for attempt in recovered.attempts() {
match attempt.status() {
AttemptStatus::Prepared => {
unresolved.push(UnresolvedAttempt::new(
attempt.key(),
EventKind::IntentAdmitted,
));
}
AttemptStatus::OutcomeUnknown => {
unresolved.push(UnresolvedAttempt::new(
attempt.key(),
EventKind::DispatchPrepared,
));
}
status => {
let rung = if recovered
.history(attempt.key())
.last()
.is_some_and(|change| {
change.to() == AttemptStatus::Verified
|| change.to() == AttemptStatus::VerificationFailed
}) {
EventKind::Verified
} else {
EventKind::OutcomeObserved
};
let record =
SettledAttempt::new(attempt.key(), rung, status, attempt.verification());
match folded_at.get(&record.action()) {
Some(at) => {
if let Some(slot) = settled.get_mut(*at) {
*slot = record;
}
}
None => {
let _ = folded_at.insert(record.action(), settled.len());
settled.push(record);
}
}
}
}
}
(settled, unresolved)
}
fn write_and_sync_async<'a>(
&'a mut self,
frames: Vec<u8>,
) -> crate::BoxFuture<'a, Result<(), JournalError>> {
let expected = self.disk_len;
Box::pin(async move {
self.storage
.submit_async(move |file, _state| commit(file, expected, &frames))
.await
.map(|_| ())
.map_err(write_refusal)
})
}
fn write_and_sync_blocking(&self, frames: &[u8]) -> Result<(), JournalError> {
let staged = frames.to_vec();
let expected = self.disk_len;
self.storage
.submit(move |file, _state| commit(file, expected, &staged))
.map(|_| ())
.map_err(write_refusal)
}
fn accept(&mut self, event: &EffectEvent, position: JournalPosition, frame_len: usize) {
self.committed.push(JournalEntry::new(position, *event));
self.position = position;
self.ladder.insert(event.key(), event.kind());
if let EffectEvent::OutcomeObserved { key, evidence } = *event {
self.outcomes.insert(key, (position, evidence));
}
self.disk_len = self
.disk_len
.saturating_add(u64::saturating_from(frame_len));
}
fn prepare_append(
&self,
expected_tail: JournalPosition,
event: &EffectEvent,
) -> Result<(JournalPosition, Vec<u8>), JournalError> {
self.refuse_sealed()?;
let actual = self.tail();
check_append_order(
expected_tail,
actual,
event,
self.ladder.get(&event.key()).copied(),
)?;
self.refuse_walked(event)?;
self.fence()?;
self.bound_events(1)?;
let (position, frame) = self.frame(event, actual)?;
self.bound_bytes(frame.len())?;
Ok((position, frame))
}
pub fn continue_as_file(&mut self) -> Result<Option<Self>, JournalError> {
if !self.continuation_continues() {
return Ok(None);
}
if let Some(next) = self.sealed_by.clone() {
let refusal = Err(JournalError::Superseded {
path: next.display().to_string(),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "continue_as_new: this handle already sealed its journal");
return refusal;
}
let successor = self.successor();
let (settled, unresolved) = self.carried();
let checkpoint = Continuation::new(
self.generation().saturating_add(1),
self.tail(),
settled,
unresolved,
)?;
let payload = checkpoint.to_payload()?;
let (_position, framed) = self.frame_payload(payload, self.tail())?;
let expected = self.disk_len;
let pause = self.pause.take();
let sealed_at = successor.clone();
let answer = self
.storage
.submit(move |file, _state| seal_on_owner(file, expected, &framed, &sealed_at, pause))
.map_err(write_refusal)?;
match answer {
SealOutcome::Sealed { bytes } => {
self.disk_len = self.disk_len.saturating_add(u64::saturating_from(bytes));
self.sealed_by = Some(successor.clone());
Ok(Some(Self::open_impl(
&successor,
OpenKind::Continuing(self.continuation),
)?))
}
SealOutcome::Paused(boundary) => {
self.sealed_by = Some(successor);
let refusal = Err(JournalError::ContinuationPaused { boundary });
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "continue_as_new: an armed fault injector stopped the seal");
refusal
}
SealOutcome::Appended => {
let refusal = Err(JournalError::Storage(std::io::Error::other(
"the seal reported an append's answer; the journal's own steps disagree",
)));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "continue_as_new: returning an error to the caller");
refusal
}
}
}
fn refuse_walked(&self, event: &EffectEvent) -> Result<(), JournalError> {
let key = event.key();
if self.ladder.contains_key(&key) {
return Ok(());
}
let walked = self
.folded
.iter()
.find(|entry| entry.already_walked(key))
.map(|entry| entry.attempt());
match walked {
Some(latest) => {
let refusal = Err(JournalError::AttemptAlreadyWalked {
key: Box::new(key),
latest,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "refuse_walked: the sealed history already walked this attempt");
refusal
}
None => Ok(()),
}
}
fn frame_payload(
&self,
payload: Vec<u8>,
from: JournalPosition,
) -> Result<(JournalPosition, Vec<u8>), JournalError> {
self.fence()?;
let sequence = from
.sequence()
.checked_add(1)
.ok_or(JournalError::Exhausted)?;
let (framed, head) = super::frame::frame_record::<Vec<u8>, _, _, _, _>(
&payload,
&from.head,
MAX_FRAME_BYTES,
|archived: &Vec<u8>| Ok(archived.clone()),
|_payload, previous, archived| super::chain_over_bytes(previous, archived),
|_len| {
JournalError::Storage(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"the sealed checkpoint exceeds this journal's frame bound",
))
},
)?;
Ok((JournalPosition { sequence, head }, framed))
}
pub fn open_with_stalled_storage(path: impl AsRef<Path>) -> Result<Self, JournalError> {
Self::open_impl(path.as_ref(), OpenKind::Stalled)
}
#[must_use]
pub fn storage_gate(&self) -> StorageGate {
self.storage.gate()
}
pub fn release_storage(&self) {
self.storage_gate().release();
}
pub fn compare_and_append_all(
&mut self,
events: &[EffectEvent],
) -> Result<Vec<DurableAck>, JournalError> {
self.fence()?;
if events.is_empty() {
return Ok(Vec::new());
}
self.bound_events(u64::saturating_from(events.len()))?;
let mut position = self.tail();
let mut staged: HashMap<crate::effect::EffectKey, EventKind> = HashMap::new();
let mut frames = Vec::new();
let mut pending: Vec<(JournalPosition, usize)> = Vec::with_capacity(events.len());
for event in events {
let key = event.key();
let attempted = event.kind();
let expected =
next_allowed_of(staged.get(&key).or_else(|| self.ladder.get(&key)).copied());
if expected != Some(attempted) {
let refusal = Err(JournalError::OutOfOrder {
key: Box::new(key),
expected,
attempted,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "compare_and_append_all: returning an error to the caller");
return refusal;
}
staged.insert(key, attempted);
let (next, framed) = self.frame(event, position)?;
position = next;
self.bound_bytes(frames.len().saturating_add(framed.len()))?;
pending.push((position, framed.len()));
frames.extend_from_slice(&framed);
}
self.write_and_sync(&frames)?;
let mut acks = Vec::with_capacity(pending.len());
for (event, entry) in events.iter().zip(&pending) {
let (position, frame_len) = *entry;
self.accept(event, position, frame_len);
acks.push(DurableAck::new(position, self.durability()));
}
Ok(acks)
}
}
fn fence_and_write(file: &mut File, expected_len: u64, bytes: &[u8]) -> std::io::Result<()> {
let on_disk = file.metadata()?.len();
if on_disk != expected_len {
let refusal = Err(stale_file());
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "fence_and_write: returning an error to the caller");
return refusal;
}
file.write_all(bytes)
}
fn stale_file() -> std::io::Error {
std::io::Error::new(
std::io::ErrorKind::InvalidData,
"the journal file moved under this controller; reopen before appending",
)
}
fn commit(
file: &mut File,
expected_len: u64,
frames: &[u8],
) -> std::io::Result<super::owner::Stage<SealOutcome, ()>> {
fence_and_write(file, expected_len, frames)?;
Ok(super::owner::Stage::Unsynced {
answer: SealOutcome::Appended,
bytes: frames.len(),
settle: Box::new(|_: &mut ()| {}),
})
}
fn write_refusal(cause: super::owner::SubmitError) -> JournalError {
if cause.is_outcome_unknown() {
JournalError::OutcomeUnknown {
cause: cause.into_io(),
}
} else {
JournalError::Storage(cause.into_io())
}
}
impl EffectJournal for FileJournal {
fn durability(&self) -> DurabilityPromise {
DurabilityPromise::ProcessCrash
}
fn tail(&self) -> JournalPosition {
self.position
}
fn committed(&self) -> Result<Vec<EffectEvent>, JournalError> {
Ok(self.events().copied().collect())
}
fn committed_entries(&self) -> Result<Vec<JournalEntry>, JournalError> {
Ok(self.committed.clone())
}
fn committed_entry(
&self,
position: JournalPosition,
) -> Result<Option<JournalEntry>, JournalError> {
let Some(index) = position
.sequence()
.checked_sub(self.chain_from.sequence())
.and_then(|offset| offset.checked_sub(1))
.and_then(|n| usize::try_from(n).ok())
else {
return Ok(None);
};
Ok(self
.committed
.get(index)
.copied()
.filter(|entry| entry.position() == position))
}
fn outcome_at(
&self,
key: crate::effect::EffectKey,
) -> Result<Option<(JournalPosition, EffectEvidence)>, JournalError> {
Ok(self.outcomes.get(&key).copied())
}
fn reserve_handoff_capacity(&self, rungs: u64) -> Result<(), JournalError> {
self.bound_events(rungs)
}
fn continuation_watermark(&self) -> Result<ContinuationWatermark, JournalError> {
let events = u64::saturating_from(self.committed.len());
let limits = (u64::saturating_from(MAX_JOURNAL_EVENTS), MAX_JOURNAL_BYTES);
let watermark = ContinuationWatermark::measured(
events,
limits.0,
self.disk_len,
limits.1,
self.continuation,
);
Ok(match self.continuation_continues() {
true => watermark,
false => ContinuationWatermark::inert(events, limits.0, self.disk_len, limits.1),
})
}
fn continue_as_new(&mut self) -> Result<Option<Box<dyn EffectJournal>>, JournalError> {
let opened = self.continue_as_file()?;
Ok(opened.map(|journal| {
let boxed: Box<dyn EffectJournal> = Box::new(journal);
boxed
}))
}
fn compare_and_append(
&mut self,
expected_tail: JournalPosition,
event: &EffectEvent,
) -> Result<DurableAck, JournalError> {
let (position, frame) = self.prepare_append(expected_tail, event)?;
self.write_and_sync(&frame)?;
self.accept(event, position, frame.len());
Ok(DurableAck::new(position, self.durability()))
}
fn compare_and_append_async<'a>(
&'a mut self,
expected_tail: JournalPosition,
event: &'a EffectEvent,
) -> crate::BoxFuture<'a, Result<DurableAck, JournalError>> {
Box::pin(async move {
let (position, frame) = self.prepare_append(expected_tail, event)?;
self.write_and_sync_async(frame.clone()).await?;
self.accept(event, position, frame.len());
Ok(DurableAck::new(position, self.durability()))
})
}
fn confirm_outcome(
&mut self,
key: crate::effect::EffectKey,
evidence: EffectEvidence,
position: JournalPosition,
required: DurabilityPromise,
) -> Result<DurableAck, JournalError> {
if !self.durability().meets(required) {
let refusal = Err(JournalError::ReceiptUnavailable { required });
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "confirm_outcome: returning an error to the caller");
return refusal;
}
let settled = self
.committed
.iter()
.find(|entry| entry.position() == position)
.map(JournalEntry::event)
.copied();
match settled {
Some(EffectEvent::OutcomeObserved {
key: observed_key,
evidence: observed_evidence,
}) if observed_key == key && observed_evidence == evidence => {
Ok(DurableAck::new(position, self.durability()))
}
_ => Err(JournalError::ReceiptMismatch {
expected: position,
actual: self.tail(),
}),
}
}
}
trait ReplayCursor {
fn seek_read_zero(&mut self) -> Result<(), JournalError>;
}
impl ReplayCursor for File {
fn seek_read_zero(&mut self) -> Result<(), JournalError> {
use std::io::Seek;
use std::io::SeekFrom;
self.seek(SeekFrom::Start(0))
.map(|_: u64| ())
.map_err(JournalError::Storage)
}
}
#[cfg(test)]
pub(super) mod tests {
use super::*;
use crate::effect::{
ActionDigest, ActionId, AttemptId, EnvironmentEpoch, EnvironmentId, FlowRevision, RunId,
};
use crate::journal::chain;
use crate::journal::frame::probe::{declared_at, frame_starts, with_prefix};
use crate::journal::{AttemptStatus, EventKind};
type TestResult = Result<(), Box<dyn std::error::Error>>;
struct TempGuard(std::path::PathBuf);
impl Drop for TempGuard {
fn drop(&mut self) {
if self.0.is_dir() {
drop(std::fs::remove_dir_all(&self.0));
} else {
drop(std::fs::remove_file(&self.0));
}
}
}
const RUN: &str = "0102030405060708090a0b0c0d0e0f10";
const ACTION: &str = "1112131415161718191a1b1c1d1e1f20";
const ENV: &str = "2122232425262728292a2b2c2d2e2f30";
const FLOW_HEX: &str = "000102030405060708090a0b0c0d0e0f101112131415161718191a1b1c1d1e1f";
const DIGEST_HEX: &str = "f0f1f2f3f4f5f6f7f8f9fafbfcfdfeffe0e1e2e3e4e5e6e7e8e9eaebecedeeef";
pub(in crate::journal) fn scratch(name: &str) -> Result<PathBuf, std::io::Error> {
crate::journal::frame::probe::scratch_path("journal-file", name)
}
fn key_for_attempt(
attempt: &str,
) -> Result<crate::effect::EffectKey, Box<dyn std::error::Error>> {
let run = RunId::from_hex(RUN)?;
let action = ActionId::from_hex(ACTION)?;
let attempt = AttemptId::from_decimal(attempt)?;
let flow = FlowRevision::from_tagged("blake3_256", FLOW_HEX)?;
let digest = ActionDigest::from_tagged("blake3_256", DIGEST_HEX)?;
let environment = EnvironmentId::from_hex(ENV)?;
let epoch = EnvironmentEpoch::from_decimal("1")?;
Ok(crate::effect::EffectIdentity::new(run, environment, flow)
.key(action, attempt, digest, epoch))
}
fn key() -> Result<crate::effect::EffectKey, Box<dyn std::error::Error>> {
key_for_attempt("1")
}
fn subject(name: &str) -> Result<(PathBuf, TempGuard), Box<dyn std::error::Error>> {
let path = scratch(name)?;
let guard = TempGuard(path.clone());
Ok((path, guard))
}
#[test]
fn writing_a_successor_reports_the_directory_it_created() -> TestResult {
let (base, _guard) = subject("successor-dir")?;
std::fs::create_dir_all(&base)?;
let first = base.join("run.jrnl.cont").join("000001");
assert!(
write_successor(&first, b"frame")?,
"the first generation's directory is created by this write"
);
let second = base.join("run.jrnl.cont").join("000002");
assert!(
!write_successor(&second, b"frame")?,
"a later generation writes into a directory that already exists"
);
assert_eq!(std::fs::read(&second)?, b"frame");
Ok(())
}
#[test]
fn a_batched_ladder_is_four_acknowledgments_from_one_sync() -> TestResult {
let (path, _guard) = subject("batch")?;
let key = key()?;
let verdict = crate::journal::Verification::new(
crate::effect::Id128::from_hex(&"42".repeat(16))?,
1,
lgwks_std::hash::blake3(b"postcondition observed"),
crate::journal::VerificationResult::Satisfied,
);
let mut journal = FileJournal::open(&path)?;
let ladder = [
EffectEvent::IntentAdmitted { key },
EffectEvent::DispatchPrepared { key },
EffectEvent::OutcomeObserved {
key,
evidence: EffectEvidence::Applied,
},
EffectEvent::Verified {
key,
verification: verdict,
},
];
let acks = journal.compare_and_append_all(&ladder)?;
assert_eq!(acks.len(), 4, "one acknowledgment per rung");
assert!(
acks.iter()
.all(|ack| ack.promise() == DurabilityPromise::ProcessCrash),
"every batched acknowledgment carries the same earned promise"
);
drop(journal);
let reopened = FileJournal::open(&path)?;
assert_eq!(reopened.committed()?.len(), 4);
assert_eq!(
reopened.recover().status(key),
Some(AttemptStatus::Verified),
"the batch folded to the ladder's end state"
);
drop(reopened);
let mut journal = FileJournal::open(&path)?;
let before = journal.committed()?.len();
let refused = journal.compare_and_append_all(&[
EffectEvent::IntentAdmitted { key: key2()? },
EffectEvent::IntentAdmitted { key: key2()? },
]);
match refused {
Err(JournalError::OutOfOrder { .. }) => {}
Err(other) => {
return Err(format!("expected an out-of-order refusal, got {other}").into());
}
Ok(acks) => {
let _ = acks;
return Err("a double admission in one batch must be refused whole".into());
}
}
assert_eq!(
journal.committed()?.len(),
before,
"a refused batch commits nothing"
);
drop(journal);
let mut reopened = FileJournal::open(&path)?;
assert_eq!(
reopened.committed()?.len(),
before,
"and the file carries none of it"
);
{
use std::io::Write as _;
let mut other = std::fs::OpenOptions::new().append(true).open(&path)?;
other.write_all(&[0x00, 0x00, 0x00])?;
other.sync_all()?;
}
let empty = reopened.compare_and_append_all(&[]);
assert!(
matches!(empty, Err(JournalError::Storage(_))),
"a stale handle must not answer an empty batch with success"
);
Ok(())
}
fn fail_next_write_of(journal: &FileJournal) {
journal.storage.fail_next_commit();
}
fn attempt_key(n: u64) -> Result<crate::effect::EffectKey, Box<dyn std::error::Error>> {
key_for_attempt(&n.to_string())
}
fn key2() -> Result<crate::effect::EffectKey, Box<dyn std::error::Error>> {
key_for_attempt("7")
}
#[test]
fn a_batch_climbs_a_key_the_disk_already_knows() -> TestResult {
let (path, _guard) = subject("batch-committed")?;
let key = key()?;
let mut journal = FileJournal::open(&path)?;
journal.compare_and_append(journal.tail(), &EffectEvent::IntentAdmitted { key })?;
let acks = journal.compare_and_append_all(&[
EffectEvent::DispatchPrepared { key },
EffectEvent::OutcomeObserved {
key,
evidence: EffectEvidence::Applied,
},
])?;
assert_eq!(acks.len(), 2, "both rungs acknowledged");
assert_eq!(
journal.recover().status(key),
Some(AttemptStatus::Applied),
"the batch continued the committed ladder"
);
Ok(())
}
#[test]
fn a_replayed_journal_is_the_journal_that_was_written() -> TestResult {
let (path, _guard) = subject("replay")?;
let key = key()?;
{
let mut journal = FileJournal::open(&path)?;
journal.compare_and_append(journal.tail(), &EffectEvent::IntentAdmitted { key })?;
journal.compare_and_append(journal.tail(), &EffectEvent::DispatchPrepared { key })?;
}
let reopened = FileJournal::open(&path)?;
assert_eq!(reopened.committed()?.len(), 2);
assert!(!reopened.torn_tail_repaired());
assert_eq!(
reopened.recover().status(key),
Some(AttemptStatus::OutcomeUnknown)
);
assert_eq!(
reopened.durability(),
DurabilityPromise::ProcessCrash,
"the file-backed promise is the one a kill can check"
);
Ok(())
}
#[test]
fn confirm_outcome_attests_only_its_own_committed_outcome() -> TestResult {
let (path, _guard) = subject("confirm")?;
let key = key()?;
let mut journal = FileJournal::open(&path)?;
journal.compare_and_append(journal.tail(), &EffectEvent::IntentAdmitted { key })?;
journal.compare_and_append(journal.tail(), &EffectEvent::DispatchPrepared { key })?;
let ack = journal.compare_and_append(
journal.tail(),
&EffectEvent::OutcomeObserved {
key,
evidence: EffectEvidence::Applied,
},
)?;
journal.confirm_outcome(
key,
EffectEvidence::Applied,
ack.position(),
DurabilityPromise::ProcessCrash,
)?;
let wrong_evidence = journal.confirm_outcome(
key,
EffectEvidence::NotApplied,
ack.position(),
DurabilityPromise::ProcessCrash,
);
assert!(
matches!(wrong_evidence, Err(JournalError::ReceiptMismatch { .. })),
"a receipt for an outcome this journal did not commit must be refused"
);
let power_loss = journal.confirm_outcome(
key,
EffectEvidence::Applied,
ack.position(),
DurabilityPromise::PowerLoss,
);
assert!(
matches!(
power_loss,
Err(JournalError::ReceiptUnavailable {
required: DurabilityPromise::PowerLoss
})
),
"a file-backed journal must not claim power-loss durability"
);
Ok(())
}
#[test]
fn a_reopened_journal_answers_the_latest_outcome_by_key() -> TestResult {
let (path, _guard) = subject("outcome-index")?;
let key = key()?;
let mut journal = FileJournal::open(&path)?;
journal.compare_and_append(journal.tail(), &EffectEvent::IntentAdmitted { key })?;
journal.compare_and_append(journal.tail(), &EffectEvent::DispatchPrepared { key })?;
let outcome = EffectEvent::OutcomeObserved {
key,
evidence: EffectEvidence::Applied,
};
let ack = journal.compare_and_append(journal.tail(), &outcome)?;
assert_eq!(
journal.outcome_at(key)?,
Some((ack.position(), EffectEvidence::Applied)),
"the live index answers the committed outcome without scanning history"
);
drop(journal);
let reopened = FileJournal::open(&path)?;
assert_eq!(
reopened.outcome_at(key)?,
Some((ack.position(), EffectEvidence::Applied)),
"the replay rebuilds the outcome index, so a restart keeps it"
);
assert_eq!(
reopened
.committed_entry(ack.position())?
.map(|entry| *entry.event()),
Some(outcome),
"the acknowledged position still holds the exact outcome"
);
Ok(())
}
#[test]
fn a_rotted_length_prefix_is_refused_and_nothing_is_trimmed() -> TestResult {
let (path, _guard) = subject("length-rot")?;
{
let mut journal = FileJournal::open(&path)?;
for attempt in 1u64..=3 {
let attempt_key_n = attempt_key(attempt)?;
journal.compare_and_append(
journal.tail(),
&EffectEvent::IntentAdmitted { key: attempt_key_n },
)?;
}
}
let before = std::fs::metadata(&path)?.len();
assert!(before > 3, "three frames are on the disk");
let mut bytes = std::fs::read(&path)?;
let first_len = u32::from_be_bytes([bytes[0], bytes[1], bytes[2], bytes[3]]);
let len_bytes = u32::saturating_from(LENGTH_BYTES);
let head_bytes = u32::saturating_from(HEAD_BYTES);
let second = usize::saturating_from(
first_len
.saturating_add(len_bytes)
.saturating_add(head_bytes),
)
.min(bytes.len());
if second < bytes.len() {
bytes[second] ^= 0x40;
}
std::fs::write(&path, &bytes)?;
match FileJournal::open(&path) {
Err(JournalError::Corrupt(corruption)) => {
assert_eq!(corruption.at(), 1, "the rotted frame is named");
assert!(
matches!(corruption.kind(), CorruptionKind::Framed),
"the refusal names the framing, not the contents"
);
}
Err(other) => return Err(format!("expected a corruption refusal, got {other}").into()),
Ok(_) => return Err("a rotted length prefix must not reopen as a journal".into()),
}
assert_eq!(
std::fs::metadata(&path)?.len(),
before,
"refused bytes are never trimmed"
);
Ok(())
}
#[test]
fn an_undecodable_committed_frame_is_refused_not_trimmed() -> TestResult {
let (path, _guard) = subject("undecodable")?;
let key = key()?;
{
let mut journal = FileJournal::open(&path)?;
journal.compare_and_append(journal.tail(), &EffectEvent::IntentAdmitted { key })?;
}
let mut bytes = std::fs::read(&path)?;
let raw_len = u32::from_be_bytes([bytes[0], bytes[1], bytes[2], bytes[3]]);
let payload_len = usize::saturating_from(raw_len);
let mid = LENGTH_BYTES.saturating_add(payload_len >> 1);
bytes[mid] = bytes[mid].wrapping_add(0x55);
std::fs::write(&path, &bytes)?;
match FileJournal::open(&path) {
Err(JournalError::Corrupt(corruption)) => {
assert_eq!(corruption.at(), 0);
assert!(matches!(
corruption.kind(),
CorruptionKind::Undecodable | CorruptionKind::Chain(_)
));
}
Err(other) => return Err(format!("expected a corruption refusal, got {other}").into()),
Ok(_) => return Err("undecodable committed bytes must not reopen".into()),
}
Ok(())
}
#[test]
fn the_ladder_refuses_a_second_prepared_dispatch_from_a_replayed_view() -> TestResult {
let (path, _guard) = subject("ladder")?;
let key = key()?;
let mut journal = FileJournal::open(&path)?;
journal.compare_and_append(journal.tail(), &EffectEvent::IntentAdmitted { key })?;
journal.compare_and_append(journal.tail(), &EffectEvent::DispatchPrepared { key })?;
let refused =
journal.compare_and_append(journal.tail(), &EffectEvent::DispatchPrepared { key });
match refused {
Err(JournalError::OutOfOrder {
expected: Some(EventKind::OutcomeObserved),
..
}) => {}
Err(other) => {
return Err(format!("expected an out-of-order refusal, got {other}").into());
}
Ok(_) => return Err("a second dispatch of one attempt must be unrepresentable".into()),
}
Ok(())
}
fn frame_bytes(
position: JournalPosition,
event: &EffectEvent,
) -> Result<Vec<u8>, Box<dyn std::error::Error>> {
let payload = event.to_bytes()?;
let head = chain(position, event)?;
let length = u32::try_from(payload.len())?;
Ok(crate::journal::frame::encode(length, &payload, &head))
}
struct FaultyAfter {
inner: std::io::Cursor<Vec<u8>>,
serve: u64,
}
impl Read for FaultyAfter {
fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
if self.inner.position() >= self.serve {
let refusal = Err(std::io::Error::other("injected storage fault"));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "read: returning an error to the caller");
return refusal;
}
let remaining = self.serve.saturating_sub(self.inner.position());
let cap = usize::saturating_from(remaining).min(buf.len());
self.inner.read(&mut buf[..cap])
}
}
#[test]
fn a_storage_fault_mid_frame_is_storage_not_a_torn_tail() -> TestResult {
let event = EffectEvent::IntentAdmitted { key: key()? };
let frame = frame_bytes(JournalPosition::genesis(), &event)?;
let mid_head = frame.len().saturating_sub(HEAD_BYTES).saturating_add(1);
for serve in [1usize, LENGTH_BYTES.saturating_add(1), mid_head] {
let serve = u64::try_from(serve)?;
let mut faulty = FaultyAfter {
inner: std::io::Cursor::new(frame.clone()),
serve,
};
match scan(&mut faulty, MAX_JOURNAL_EVENTS) {
Err(JournalError::Storage(_)) => {}
Err(other) => {
return Err(format!(
"a fault {serve} bytes in was classified {other}, not storage"
)
.into());
}
Ok(scanned) => {
let stop = scanned.stop;
return Err(format!(
"a fault {serve} bytes in stopped the scan as {stop:?}, not storage"
)
.into());
}
}
}
Ok(())
}
#[test]
fn scanning_refuses_a_complete_event_beyond_the_limit() -> TestResult {
let first = EffectEvent::IntentAdmitted {
key: attempt_key(1)?,
};
let second = EffectEvent::IntentAdmitted {
key: attempt_key(2)?,
};
let genesis = JournalPosition::genesis();
let first_position = JournalPosition {
sequence: 1,
head: chain(genesis, &first)?,
};
let mut bytes = frame_bytes(genesis, &first)?;
bytes.extend_from_slice(&frame_bytes(first_position, &second)?);
let mut reader = std::io::Cursor::new(bytes);
match scan(&mut reader, 1) {
Err(JournalError::CapacityExceeded {
resource: JournalLimitKind::Events,
limit,
requested,
}) => {
assert_eq!(limit, 1);
assert_eq!(requested, 2);
}
Err(other) => return Err(format!("expected event-limit refusal, got {other}").into()),
Ok(scanned) => {
let entries = scanned.entries;
return Err(format!(
"two frames under a one-event limit must refuse, retained {}",
entries.len()
)
.into());
}
}
Ok(())
}
#[test]
fn open_refuses_an_over_limit_file_without_truncating_it() -> TestResult {
let (path, _guard) = subject("byte-limit")?;
let file = File::create(&path)?;
let requested = MAX_JOURNAL_BYTES.saturating_add(1);
file.set_len(requested)?;
drop(file);
match FileJournal::open(&path) {
Err(JournalError::CapacityExceeded {
resource: JournalLimitKind::Bytes,
limit,
requested: actual,
}) => {
assert_eq!(limit, MAX_JOURNAL_BYTES);
assert_eq!(actual, requested);
}
Err(other) => return Err(format!("expected byte-limit refusal, got {other}").into()),
Ok(_) => return Err("an over-limit journal must not be opened".into()),
}
assert_eq!(
std::fs::metadata(&path)?.len(),
requested,
"capacity refusal preserves the existing file byte-for-byte in length"
);
Ok(())
}
#[test]
fn batch_admission_refuses_history_over_the_event_limit_without_writing() -> TestResult {
let (path, _guard) = subject("event-limit")?;
let requested = MAX_JOURNAL_EVENTS.saturating_add(1);
let mut events = Vec::new();
for attempt in 1..=requested {
events.push(EffectEvent::IntentAdmitted {
key: attempt_key(u64::try_from(attempt)?)?,
});
}
let mut journal = FileJournal::open(&path)?;
match journal.compare_and_append_all(&events) {
Err(JournalError::CapacityExceeded {
resource: JournalLimitKind::Events,
limit,
requested: actual,
}) => {
assert_eq!(limit, u64::try_from(MAX_JOURNAL_EVENTS)?);
assert_eq!(actual, u64::try_from(requested)?);
}
Err(other) => return Err(format!("expected event-limit refusal, got {other}").into()),
Ok(acks) => {
return Err(format!(
"over-limit batch must refuse before write, returned {} acknowledgments",
acks.len()
)
.into());
}
}
assert!(
journal.committed()?.is_empty(),
"refused batch leaves the complete prior history unchanged"
);
assert_eq!(
std::fs::metadata(&path)?.len(),
0,
"refused batch writes no partial prefix"
);
Ok(())
}
fn three_frame_file(path: &Path) -> Result<(Vec<u8>, usize), Box<dyn std::error::Error>> {
{
let mut journal = FileJournal::open(path)?;
for attempt in 1u64..=3 {
journal.compare_and_append(
journal.tail(),
&EffectEvent::IntentAdmitted {
key: attempt_key(attempt)?,
},
)?;
}
}
let bytes = std::fs::read(path)?;
let first = u32::from_be_bytes([bytes[0], bytes[1], bytes[2], bytes[3]]);
let first_len = crate::journal::frame::framed_len(usize::try_from(first)?);
let next = first_len;
let mut second_prefix = [0u8; LENGTH_BYTES];
for (slot, offset) in second_prefix.iter_mut().zip(0u64..) {
let index = next.saturating_add(offset);
*slot = *bytes
.get(usize::try_from(index)?)
.ok_or("fixture file is shorter than its own first frame")?;
}
let second = u32::from_be_bytes(second_prefix);
let second_len = crate::journal::frame::framed_len(usize::try_from(second)?);
Ok((
bytes,
usize::try_from(first_len.saturating_add(second_len))?,
))
}
#[test]
fn an_inflated_length_over_a_complete_final_frame_is_refused_not_trimmed() -> TestResult {
let (path, _guard) = subject("inflated-length")?;
let (mut bytes, third_start) = three_frame_file(&path)?;
let before = bytes.len();
let inflated = u32::try_from(MAX_FRAME_BYTES)?;
let over_limit = inflated.saturating_add(1);
for (name, replacement) in [
("the maximum legal length", inflated),
("one byte over the limit", over_limit),
("a zero length", 0u32),
] {
for (slot, offset) in bytes[third_start..]
.iter_mut()
.take(LENGTH_BYTES)
.zip(0usize..)
{
*slot = replacement.to_be_bytes()[offset];
}
std::fs::write(&path, &bytes)?;
match FileJournal::open(&path) {
Err(JournalError::Corrupt(corruption)) => {
assert_eq!(corruption.at(), 2, "{name}: the lying frame is named");
assert!(
matches!(corruption.kind(), CorruptionKind::Framed),
"{name}: the refusal names the framing, got {:?}",
corruption.kind()
);
}
Err(other) => {
return Err(
format!("{name}: expected a corruption refusal, got {other}").into(),
);
}
Ok(_) => {
return Err(
format!("{name}: a lying length must not reopen as a journal").into(),
);
}
}
assert_eq!(
std::fs::metadata(&path)?.len(),
u64::try_from(before)?,
"{name}: refused bytes are never trimmed"
);
}
Ok(())
}
#[test]
fn a_short_final_frame_is_still_repaired_as_a_torn_tail() -> TestResult {
let (path, _guard) = subject("torn-tail")?;
let (bytes, third_start) = three_frame_file(&path)?;
let complete_len = u64::try_from(third_start)?;
let torn = &bytes[..third_start.saturating_add(10)];
std::fs::write(&path, torn)?;
let reopened = FileJournal::open(&path)?;
assert!(
reopened.torn_tail_repaired(),
"an interrupted append is repaired, not refused"
);
assert_eq!(
reopened.committed()?.len(),
2,
"the two kept frames survive"
);
assert_eq!(
std::fs::metadata(&path)?.len(),
complete_len,
"the repair restores the prefix every acknowledgment names"
);
Ok(())
}
fn require_refused_untouched(path: &Path, before: &[u8], at: u64, why: &str) -> TestResult {
let outcome: TestResult = match FileJournal::open(path) {
Err(JournalError::Corrupt(corruption)) => {
assert_eq!(corruption.at(), at, "{why}: the lying frame is named");
assert!(
matches!(corruption.kind(), CorruptionKind::Framed),
"{why}: the refusal names the framing, got {:?}",
corruption.kind()
);
Ok(())
}
Err(other) => Err(format!("{why}: expected a corruption refusal, got {other}").into()),
Ok(opened) => Err(format!(
"{why}: reopened as a journal of {} events, so an acknowledged frame was lost",
opened.events().count()
)
.into()),
};
outcome?;
assert_eq!(
lgwks_std::hash::blake3(&std::fs::read(path)?),
lgwks_std::hash::blake3(before),
"{why}: refused bytes are never touched"
);
Ok(())
}
#[test]
fn a_lengthened_acknowledged_final_frame_is_refused_not_trimmed() -> TestResult {
let (path, _guard) = subject("lengthened-final")?;
let (bytes, third) = three_frame_file(&path)?;
let declared = declared_at(&bytes, third);
for extra in 1u32..=1024 {
let lied = with_prefix(&bytes, third, declared + extra);
std::fs::write(&path, &lied)?;
require_refused_untouched(&path, &lied, 2, &format!("final frame L+{extra}"))?;
}
Ok(())
}
#[test]
fn a_lengthened_event_behind_a_carried_seal_is_refused_not_trimmed() -> TestResult {
let (dir, _guard) = subject("lengthened-successor")?;
std::fs::create_dir_all(&dir)?;
let successor = {
let mut journal = FileJournal::open_continuing_with(
dir.join("run.jrnl"),
crate::journal::ContinuationPolicy::declared(),
)?;
let first = EffectEvent::IntentAdmitted {
key: attempt_key(1)?,
};
journal.compare_and_append(journal.tail(), &first)?;
let mut next = journal
.continue_as_file()?
.ok_or("a continuing journal hands back its successor")?;
let second = EffectEvent::IntentAdmitted {
key: attempt_key(2)?,
};
next.compare_and_append(next.tail(), &second)?;
next.path().to_path_buf()
};
let bytes = std::fs::read(&successor)?;
let event = frame_starts(&bytes, 0)?[1];
let declared = declared_at(&bytes, event);
for extra in 1u32..=96 {
let lied = with_prefix(&bytes, event, declared + extra);
std::fs::write(&successor, &lied)?;
require_refused_untouched(&successor, &lied, 0, &format!("successor event L+{extra}"))?;
}
Ok(())
}
#[test]
fn an_inflated_non_final_length_is_refused_and_every_byte_survives() -> TestResult {
let (path, _guard) = subject("inflated-middle")?;
let (bytes, _) = three_frame_file(&path)?;
let middle = frame_starts(&bytes, 0)?[1];
let remaining = u32::try_from(bytes.len() - middle)?;
for declared in [
remaining,
remaining + 1,
remaining + 31,
remaining + 500,
u32::try_from(MAX_FRAME_BYTES)?,
] {
let lied = with_prefix(&bytes, middle, declared);
std::fs::write(&path, &lied)?;
require_refused_untouched(
&path,
&lied,
1,
&format!("middle frame declared {declared}"),
)?;
}
Ok(())
}
#[test]
fn an_append_cut_at_every_byte_of_the_final_frame_is_repaired() -> TestResult {
let (path, _guard) = subject("cut-every-byte")?;
let (bytes, third) = three_frame_file(&path)?;
for cut in third + 1..bytes.len() {
std::fs::write(&path, &bytes[..cut])?;
let reopened = FileJournal::open(&path).map_err(|error| {
format!(
"a final frame cut at byte {cut} of {} was refused: {error}",
bytes.len()
)
})?;
assert!(
reopened.torn_tail_repaired(),
"cut at {cut} was not reported as repaired"
);
assert_eq!(
reopened.committed()?.len(),
2,
"cut at {cut} lost an acknowledged frame"
);
drop(reopened);
assert_eq!(
std::fs::metadata(&path)?.len(),
u64::try_from(third)?,
"cut at {cut} was not trimmed to the acknowledged prefix"
);
}
Ok(())
}
#[test]
fn a_damaged_cut_frame_with_an_acknowledged_frame_behind_it_is_refused() -> TestResult {
let (path, _guard) = subject("damaged-middle")?;
let (bytes, _) = three_frame_file(&path)?;
let starts = frame_starts(&bytes, 0)?;
let middle = starts[1];
let remaining = u32::try_from(bytes.len() - middle)?;
let mut lied = with_prefix(&bytes, middle, remaining + 7);
lied[middle + LENGTH_BYTES + 3] ^= 0x55;
std::fs::write(&path, &lied)?;
require_refused_untouched(&path, &lied, 1, "damaged middle frame")
}
#[test]
fn a_final_frame_with_a_lying_length_and_a_damaged_head_is_the_stated_limit() -> TestResult {
let (path, _guard) = subject("two-faults")?;
let (bytes, third) = three_frame_file(&path)?;
let declared = declared_at(&bytes, third);
let mut lied = with_prefix(&bytes, third, declared + 5);
let last = lied.len() - 1;
lied[last] ^= 0xff;
std::fs::write(&path, &lied)?;
let reopened = FileJournal::open(&path)?;
assert!(reopened.torn_tail_repaired());
assert_eq!(reopened.committed()?.len(), 2);
Ok(())
}
#[test]
fn a_streaming_replay_refuses_a_lengthened_final_frame_and_then_ends() -> TestResult {
let (path, _guard) = subject("replay-lengthened")?;
let (bytes, third) = three_frame_file(&path)?;
let journal = FileJournal::open(&path)?;
let declared = declared_at(&bytes, third);
for extra in [1u32, 17, 32, 33, 400] {
std::fs::write(&path, with_prefix(&bytes, third, declared + extra))?;
let mut replay = journal.replay()?;
assert!(
matches!(replay.next(), Some(Ok(_))),
"L+{extra}: first event"
);
assert!(
matches!(replay.next(), Some(Ok(_))),
"L+{extra}: second event"
);
assert!(
matches!(replay.next(), Some(Err(JournalError::Corrupt(_)))),
"L+{extra}: the lengthened frame must be a refusal, not the end of the stream"
);
assert!(
replay.next().is_none(),
"L+{extra}: a refusal ends the stream"
);
}
std::fs::write(&path, &bytes[..third + 10])?;
let events = journal.replay()?.collect::<Result<Vec<_>, _>>()?;
assert_eq!(
events.len(),
2,
"a cut append still ends the stream quietly"
);
Ok(())
}
#[test]
fn a_failed_batch_write_latches_the_poison_like_a_single_append() -> TestResult {
let (path, _guard) = subject("batch-poison")?;
let event = |attempt: u64| -> Result<EffectEvent, Box<dyn std::error::Error>> {
Ok(EffectEvent::IntentAdmitted {
key: attempt_key(attempt)?,
})
};
let mut journal = FileJournal::open(&path)?;
journal.compare_and_append(journal.tail(), &event(200)?)?;
fail_next_write_of(&journal);
match journal.compare_and_append(journal.tail(), &event(201)?) {
Err(JournalError::OutcomeUnknown { .. }) => {}
Err(other) => return Err(format!("expected an unknown outcome, got {other}").into()),
Ok(_) => return Err("a failed write must not acknowledge".into()),
}
match journal.compare_and_append(journal.tail(), &event(201)?) {
Err(error) => assert!(
error.to_string().contains("previous append failed"),
"the poisoned handle refused with its own fence, not {error}"
),
Ok(_) => return Err("a poisoned handle must refuse further appends".into()),
}
drop(journal);
let mut journal = FileJournal::open(&path)?;
journal.compare_and_append(journal.tail(), &event(300)?)?;
fail_next_write_of(&journal);
match journal.compare_and_append_all(&[event(301)?, event(302)?]) {
Err(JournalError::OutcomeUnknown { .. }) => {}
Err(other) => {
return Err(format!(
"a failed batch write answered {other}, not the single append's unknown outcome"
)
.into());
}
Ok(acks) => {
let _ = acks;
return Err("a failed batch write must not acknowledge".into());
}
}
match journal.compare_and_append_all(&[event(303)?]) {
Err(error) => assert!(
error.to_string().contains("previous append failed"),
"the poisoned handle refused with its own fence, not {error}"
),
Ok(_) => return Err("a poisoned handle must refuse further batches".into()),
}
Ok(())
}
#[test]
#[cfg(feature = "rt")]
fn a_stalled_device_does_not_stop_the_task_waiting_on_it() -> TestResult {
use std::future::poll_fn;
use std::task::Poll;
use std::time::{Duration, Instant};
let (path, _guard) = subject("stalled-device")?;
let mut journal = FileJournal::open_with_stalled_storage(&path)?;
let gate = journal.storage_gate();
let event = EffectEvent::IntentAdmitted {
key: attempt_key(900)?,
};
let tail = journal.tail();
let runtime = crate::rt::runtime::Runtime::new()?;
let mut append = Box::pin(journal.compare_and_append_async(tail, &event));
let mut last = Instant::now();
let mut longest = Duration::ZERO;
let mut ticks = 0u64;
let ack = runtime.block_on(poll_fn(|cx| {
let now = Instant::now();
longest = longest.max(now.duration_since(last));
last = now;
ticks = ticks.saturating_add(1);
match append.as_mut().poll(cx) {
Poll::Ready(answer) => Poll::Ready(answer),
Poll::Pending => {
if ticks >= 64 {
gate.release();
}
cx.waker().wake_by_ref();
Poll::Pending
}
}
}))?;
drop(append);
assert!(
ticks >= 64,
"the runtime stopped after {ticks} wakeups before the parked flush \
was released: the durable write is running on the awaiting thread"
);
assert_eq!(
ack.promise(),
DurabilityPromise::ProcessCrash,
"the acknowledgment still has to be earned, and a released device earns it"
);
assert!(
longest < Duration::from_millis(500),
"the executor was unavailable for {longest:?} while a flush that \
could not finish was outstanding"
);
assert_eq!(
EffectJournal::committed(&journal)?.len(),
1,
"the frame was folded in"
);
drop(journal);
let reopened = FileJournal::open(&path)?;
assert_eq!(
EffectJournal::committed(&reopened)?.len(),
1,
"the released frame must be on the disk, not only in the handle"
);
Ok(())
}
#[test]
#[cfg(feature = "rt")]
fn a_dropped_waiter_poisons_the_handle_and_a_reopen_does_not_duplicate() -> TestResult {
use std::future::poll_fn;
use std::task::Poll;
let (path, _guard) = subject("dropped-waiter")?;
let mut journal = FileJournal::open_with_stalled_storage(&path)?;
let gate = journal.storage_gate();
let event = EffectEvent::IntentAdmitted {
key: attempt_key(901)?,
};
let tail = journal.tail();
let runtime = crate::rt::runtime::Runtime::new()?;
let mut append = Box::pin(journal.compare_and_append_async(tail, &event));
let started = runtime.block_on(poll_fn(|cx| {
let polled = append.as_mut().poll(cx);
if polled.is_pending() {
cx.waker().wake_by_ref();
}
Poll::Ready(polled)
}));
assert!(
started.is_pending(),
"the append completed before it could be abandoned"
);
drop(append);
gate.release();
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10);
while !journal.storage.poisoned() {
if std::time::Instant::now() >= deadline {
return Err("the storage owner never latched the poison".into());
}
std::thread::yield_now();
}
match journal.compare_and_append(journal.tail(), &event) {
Err(error) => assert!(
error.to_string().contains("previous append failed"),
"an abandoned append must refuse with its own fence, not {error}"
),
Ok(_) => return Err("a handle that lost its waiter must refuse".into()),
}
drop(journal);
let mut reopened = FileJournal::open(&path)?;
let landed = EffectJournal::committed(&reopened)?.len();
assert!(
landed <= 1,
"an abandoned append wrote {landed} frames for one attempt"
);
let same = EffectEvent::IntentAdmitted {
key: attempt_key(901)?,
};
assert!(
reopened.compare_and_append(reopened.tail(), &same).is_err(),
"the replayed ladder must refuse a rung the file already records"
);
Ok(())
}
}