use core::fmt;
use std::collections::HashMap;
use std::io;
use lgwks_std::hash::{Digest, Hasher, blake3};
use lgwks_std::wire::{AlignedVec, WireError};
use crate::effect::{EffectKey, Id128};
use frame::SaturatingFrom;
pub use crate::ecs::EffectEvidence;
use crate::BoxFuture;
mod wire_form;
pub use wire_form::{
ArchivedEffectEvent, ArchivedVerificationResult, EffectEvent, VerificationResult,
};
mod file;
pub(crate) mod frame;
pub(crate) mod owner;
pub mod continuation;
pub use continuation::{
CONTINUATION_WATERMARK_DENOMINATOR, CONTINUATION_WATERMARK_NUMERATOR, Continuation,
ContinuationPolicy, ContinuationWatermark, MAX_CHECKPOINT_SETTLED, MAX_CHECKPOINT_UNRESOLVED,
SealPause, SettledAttempt, UnresolvedAttempt, is_generation_path, successor_path,
};
pub use file::{Corruption, CorruptionKind, FileJournal, Replay};
pub use owner::StorageGate;
const GENESIS_DOMAIN: &[u8] = b"lgwks.journal.v1.genesis";
pub const MAX_JOURNAL_EVENTS: usize = 100_000;
pub const MAX_JOURNAL_BYTES: u64 = 64 * 1024 * 1024;
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
#[non_exhaustive]
pub enum DurabilityPromise {
Ephemeral,
ProcessCrash,
PowerLoss,
}
impl DurabilityPromise {
#[must_use]
pub const fn survives_process_crash(self) -> bool {
matches!(self, Self::ProcessCrash | Self::PowerLoss)
}
#[must_use]
pub const fn meets(self, required: Self) -> bool {
match (self, required) {
(_, Self::Ephemeral) => true,
(Self::ProcessCrash | Self::PowerLoss, Self::ProcessCrash) => true,
(Self::PowerLoss, Self::PowerLoss) => true,
(Self::Ephemeral, _) => false,
(Self::ProcessCrash, Self::PowerLoss) => false,
}
}
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Ephemeral => "ephemeral",
Self::ProcessCrash => "process_crash",
Self::PowerLoss => "power_loss",
}
}
}
impl fmt::Display for DurabilityPromise {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(self.as_str())
}
}
#[derive(
Debug,
Clone,
Copy,
PartialEq,
Eq,
Hash,
lgwks_std::wire::Archive,
lgwks_std::wire::Serialize,
lgwks_std::wire::Deserialize,
)]
#[rkyv(crate = lgwks_std::wire::rkyv, compare(PartialEq), derive(Debug))]
pub struct JournalPosition {
sequence: u64,
head: Digest,
}
impl JournalPosition {
#[must_use]
pub fn genesis() -> Self {
Self {
sequence: 0,
head: blake3(GENESIS_DOMAIN),
}
}
#[must_use]
pub const fn sequence(self) -> u64 {
self.sequence
}
#[must_use]
pub const fn head(self) -> Digest {
self.head
}
}
impl fmt::Display for JournalPosition {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "{} @ {}", self.sequence, self.head)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
#[non_exhaustive]
pub enum EventKind {
IntentAdmitted,
DispatchPrepared,
OutcomeObserved,
Verified,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
#[non_exhaustive]
pub enum JournalLimitKind {
Events,
Bytes,
ContinuationWatermark,
Generation,
CheckpointActions,
CheckpointUnresolved,
}
impl EventKind {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::IntentAdmitted => "intent_admitted",
Self::DispatchPrepared => "dispatch_prepared",
Self::OutcomeObserved => "outcome_observed",
Self::Verified => "verified",
}
}
#[must_use]
pub const fn next(self) -> Option<Self> {
match self {
Self::IntentAdmitted => Some(Self::DispatchPrepared),
Self::DispatchPrepared => Some(Self::OutcomeObserved),
Self::OutcomeObserved | Self::Verified => Some(Self::Verified),
}
}
}
impl fmt::Display for EventKind {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(self.as_str())
}
}
impl VerificationResult {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Satisfied => "satisfied",
Self::NotSatisfied => "not_satisfied",
}
}
}
impl fmt::Display for VerificationResult {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(self.as_str())
}
}
#[derive(
Debug,
Clone,
Copy,
PartialEq,
Eq,
Hash,
lgwks_std::wire::Archive,
lgwks_std::wire::Serialize,
lgwks_std::wire::Deserialize,
)]
#[rkyv(crate = lgwks_std::wire::rkyv, compare(PartialEq), derive(Debug))]
pub struct Verification {
predicate: Id128,
predicate_version: u64,
observations: Digest,
result: VerificationResult,
}
impl Verification {
#[must_use]
pub const fn new(
predicate: Id128,
predicate_version: u64,
observations: Digest,
result: VerificationResult,
) -> Self {
Self {
predicate,
predicate_version,
observations,
result,
}
}
#[must_use]
pub const fn predicate(self) -> Id128 {
self.predicate
}
#[must_use]
pub const fn predicate_version(self) -> u64 {
self.predicate_version
}
#[must_use]
pub const fn observations(self) -> Digest {
self.observations
}
#[must_use]
pub const fn result(self) -> VerificationResult {
self.result
}
}
impl EffectEvent {
#[must_use]
pub const fn key(self) -> EffectKey {
match self {
Self::IntentAdmitted { key }
| Self::DispatchPrepared { key }
| Self::OutcomeObserved { key, .. }
| Self::Verified { key, .. } => key,
}
}
#[must_use]
pub const fn kind(self) -> EventKind {
match self {
Self::IntentAdmitted { .. } => EventKind::IntentAdmitted,
Self::DispatchPrepared { .. } => EventKind::DispatchPrepared,
Self::OutcomeObserved { .. } => EventKind::OutcomeObserved,
Self::Verified { .. } => EventKind::Verified,
}
}
pub fn to_bytes(self) -> Result<AlignedVec, WireError> {
lgwks_std::wire::to_bytes::<WireError>(&self)
}
}
impl fmt::Display for EffectEvent {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "{} {}", self.kind(), self.key())
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub struct JournalEntry {
position: JournalPosition,
event: EffectEvent,
}
impl JournalEntry {
#[must_use]
pub const fn new(position: JournalPosition, event: EffectEvent) -> Self {
Self { position, event }
}
#[must_use]
pub const fn position(&self) -> JournalPosition {
self.position
}
#[must_use]
pub const fn event(&self) -> &EffectEvent {
&self.event
}
}
#[derive(Debug)]
#[non_exhaustive]
pub enum JournalError {
TailMismatch {
expected: JournalPosition,
actual: JournalPosition,
},
SnapshotStale {
recovered_events: u64,
committed_events: u64,
},
PromiseUnmet {
required: DurabilityPromise,
offered: DurabilityPromise,
},
ReceiptUnavailable {
required: DurabilityPromise,
},
ReceiptMismatch {
expected: JournalPosition,
actual: JournalPosition,
},
EntryMismatch {
position: JournalPosition,
expected: Box<EffectEvent>,
actual: Box<EffectEvent>,
},
OutOfOrder {
key: Box<EffectKey>,
expected: Option<EventKind>,
attempted: EventKind,
},
Exhausted,
Superseded {
path: String,
},
AttemptAlreadyWalked {
key: Box<EffectKey>,
latest: crate::effect::AttemptId,
},
ContinuationPaused {
boundary: SealPause,
},
CapacityExceeded {
resource: JournalLimitKind,
limit: u64,
requested: u64,
},
Locked {
path: String,
},
Storage(io::Error),
OutcomeUnknown {
cause: io::Error,
},
Encoding(WireError),
Corrupt(Box<Corruption>),
}
impl fmt::Display for JournalError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match *self {
Self::TailMismatch { expected, actual } => write!(
f,
"journal tail moved: expected {expected}, committed {actual}"
),
Self::SnapshotStale {
recovered_events,
committed_events,
} => write!(
f,
"journal advanced during recovery: folded {recovered_events} events, committed {committed_events}"
),
Self::PromiseUnmet { required, offered } => write!(
f,
"journal promises {offered}, which is below the {required} this needs"
),
Self::ReceiptUnavailable { required } => write!(
f,
"journal cannot attest that the recorded outcome meets {required}"
),
Self::ReceiptMismatch { expected, actual } => write!(
f,
"journal receipt names {actual}, not the outcome committed at {expected}"
),
Self::EntryMismatch {
position,
ref expected,
ref actual,
} => write!(
f,
"journal entry at {position} is {actual}, not the acknowledged {expected}"
),
Self::OutOfOrder {
ref key,
expected,
attempted,
} => match expected {
Some(next) => write!(
f,
"{attempted} cannot follow what is committed for {key}; \
the next recorded fact is {next}"
),
None => write!(
f,
"{attempted} cannot be appended for {key}; \
its recorded facts are complete"
),
},
Self::Exhausted => f.write_str("journal position exhausted"),
Self::Superseded { ref path } => write!(
f,
"this journal is sealed and its successor is authoritative; open {path}"
),
Self::AttemptAlreadyWalked { ref key, latest } => write!(
f,
"{key} is an attempt at or below the latest this journal walked \
for its action, which is {latest}; the sealed history already \
records it"
),
Self::ContinuationPaused { boundary } => write!(
f,
"the continuation was armed to stop at {boundary}, so no successor \
was established"
),
Self::CapacityExceeded {
resource,
limit,
requested,
} => write!(
f,
"journal {resource:?} limit exceeded: limit {limit}, requested {requested}"
),
Self::Locked { ref path } => write!(
f,
"journal at {path} is held by another writer; \
wait for the holder to finish, then reopen"
),
Self::Storage(ref cause) => {
write!(f, "journal storage refused the append: {cause}")
}
Self::OutcomeUnknown { ref cause } => write!(
f,
"journal append may have committed, acknowledgment lost: {cause}; \
reconcile by readback rather than re-appending"
),
Self::Encoding(ref cause) => {
write!(
f,
"journal could not encode the event for the chain: {cause}"
)
}
Self::Corrupt(ref corruption) => {
write!(f, "journal refused its own committed bytes: {corruption}")
}
}
}
}
impl std::error::Error for JournalError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match *self {
Self::Storage(ref cause) => Some(cause),
Self::OutcomeUnknown { ref cause } => Some(cause),
Self::Encoding(ref cause) => Some(cause),
Self::Corrupt(ref cause) => Some(cause),
_ => None,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub struct DurableAck {
position: JournalPosition,
promise: DurabilityPromise,
}
impl DurableAck {
#[must_use]
pub const fn new(position: JournalPosition, promise: DurabilityPromise) -> Self {
Self { position, promise }
}
#[must_use]
pub const fn position(self) -> JournalPosition {
self.position
}
#[must_use]
pub const fn promise(self) -> DurabilityPromise {
self.promise
}
}
pub trait EffectJournal {
fn durability(&self) -> DurabilityPromise;
fn tail(&self) -> JournalPosition;
fn committed(&self) -> Result<Vec<EffectEvent>, JournalError>;
fn committed_entries(&self) -> Result<Vec<JournalEntry>, JournalError> {
Err(JournalError::ReceiptUnavailable {
required: self.durability(),
})
}
fn committed_entry(
&self,
position: JournalPosition,
) -> Result<Option<JournalEntry>, JournalError> {
Ok(self
.committed_entries()?
.into_iter()
.find(|entry| entry.position() == position))
}
fn outcome_at(
&self,
key: EffectKey,
) -> Result<Option<(JournalPosition, EffectEvidence)>, JournalError> {
Ok(self
.committed_entries()?
.into_iter()
.rev()
.find_map(|entry| match *entry.event() {
EffectEvent::OutcomeObserved {
key: held,
evidence,
} if held == key => Some((entry.position(), evidence)),
_ => None,
}))
}
fn compare_and_append(
&mut self,
expected_tail: JournalPosition,
event: &EffectEvent,
) -> Result<DurableAck, JournalError>;
fn compare_and_append_async<'a>(
&'a mut self,
expected_tail: JournalPosition,
event: &'a EffectEvent,
) -> BoxFuture<'a, Result<DurableAck, JournalError>> {
Box::pin(async move { self.compare_and_append(expected_tail, event) })
}
fn confirm_outcome(
&mut self,
_key: EffectKey,
_evidence: EffectEvidence,
_position: JournalPosition,
required: DurabilityPromise,
) -> Result<DurableAck, JournalError> {
Err(JournalError::ReceiptUnavailable { required })
}
fn admit_external_handoff(&self) -> Result<DurabilityPromise, JournalError> {
let offered = self.durability();
if offered.survives_process_crash() {
return Ok(offered);
}
Err(JournalError::PromiseUnmet {
required: DurabilityPromise::ProcessCrash,
offered,
})
}
fn reserve_handoff_capacity(&self, _rungs: u64) -> Result<(), JournalError> {
Ok(())
}
fn continuation_watermark(&self) -> Result<ContinuationWatermark, JournalError> {
Ok(ContinuationWatermark::inert(
0,
u64::saturating_from(MAX_JOURNAL_EVENTS),
0,
MAX_JOURNAL_BYTES,
))
}
fn continue_as_new(&mut self) -> Result<Option<Box<dyn EffectJournal>>, JournalError> {
Ok(None)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
#[non_exhaustive]
pub enum AttemptStatus {
Prepared,
OutcomeUnknown,
Applied,
NotApplied,
Verified,
VerificationFailed,
}
impl AttemptStatus {
#[must_use]
pub const fn is_uncertain(self) -> bool {
matches!(self, Self::OutcomeUnknown)
}
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Prepared => "prepared",
Self::OutcomeUnknown => "outcome_unknown",
Self::Applied => "applied",
Self::NotApplied => "not_applied",
Self::Verified => "verified",
Self::VerificationFailed => "verification_failed",
}
}
}
impl fmt::Display for AttemptStatus {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(self.as_str())
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub struct Attempt {
key: EffectKey,
status: AttemptStatus,
verification: Option<Verification>,
}
impl Attempt {
#[must_use]
pub const fn new(key: EffectKey, status: AttemptStatus) -> Self {
Self {
key,
status,
verification: None,
}
}
#[must_use]
pub const fn verification(&self) -> Option<Verification> {
self.verification
}
#[must_use]
pub const fn key(&self) -> EffectKey {
self.key
}
#[must_use]
pub const fn status(&self) -> AttemptStatus {
self.status
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub struct Transition {
from: Option<AttemptStatus>,
to: AttemptStatus,
position: usize,
}
impl Transition {
#[must_use]
pub const fn from(&self) -> Option<AttemptStatus> {
self.from
}
#[must_use]
pub const fn to(&self) -> AttemptStatus {
self.to
}
#[must_use]
pub const fn position(&self) -> usize {
self.position
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct Recovered {
attempts: Vec<Attempt>,
index: HashMap<EffectKey, usize>,
transitions: Vec<Vec<Transition>>,
}
impl Recovered {
#[must_use]
pub fn status(&self, key: EffectKey) -> Option<AttemptStatus> {
self.index
.get(&key)
.and_then(|index| self.attempts.get(*index))
.map(|attempt| attempt.status)
}
#[must_use]
pub fn verification(&self, key: EffectKey) -> Option<Verification> {
self.index
.get(&key)
.and_then(|index| self.attempts.get(*index))
.and_then(Attempt::verification)
}
#[must_use]
pub fn history(&self, key: EffectKey) -> &[Transition] {
self.index
.get(&key)
.and_then(|index| self.transitions.get(*index))
.map_or(&[], Vec::as_slice)
}
#[must_use]
pub fn len(&self) -> usize {
self.attempts.len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.attempts.is_empty()
}
#[must_use]
pub fn uncertain(&self) -> Vec<EffectKey> {
self.attempts
.iter()
.filter(|attempt| attempt.status.is_uncertain())
.map(|attempt| attempt.key)
.collect()
}
#[must_use]
pub fn attempts(&self) -> &[Attempt] {
&self.attempts
}
}
#[must_use]
pub fn recover<'a>(events: impl IntoIterator<Item = &'a EffectEvent>) -> Recovered {
recover_continued(None, events)
}
#[must_use]
pub fn recover_continued<'a>(
base: Option<&Continuation>,
events: impl IntoIterator<Item = &'a EffectEvent>,
) -> Recovered {
let mut recovered = Recovered::default();
if let Some(checkpoint) = base {
for carried in checkpoint.settled() {
let (Some(status), Some(_rung)) = (carried.status(), carried.rung()) else {
continue;
};
seed(
&mut recovered,
carried.key(),
status,
carried.verification(),
);
}
for carried in checkpoint.unresolved() {
let Some(status) = carried.recovered_status() else {
continue;
};
seed(&mut recovered, carried.key(), status, None);
}
}
fold_events(&mut recovered, events);
recovered
}
fn seed(
recovered: &mut Recovered,
key: EffectKey,
status: AttemptStatus,
verification: Option<Verification>,
) {
let index = recovered.attempts.len();
let mut attempt = Attempt::new(key, status);
attempt.verification = verification;
recovered.attempts.push(attempt);
recovered.transitions.push(vec![Transition {
from: None,
to: status,
position: index,
}]);
recovered.index.insert(key, index);
}
fn fold_events<'a>(recovered: &mut Recovered, events: impl IntoIterator<Item = &'a EffectEvent>) {
for (position, event) in events.into_iter().enumerate() {
let status = match *event {
EffectEvent::IntentAdmitted { .. } => AttemptStatus::Prepared,
EffectEvent::DispatchPrepared { .. } => AttemptStatus::OutcomeUnknown,
EffectEvent::OutcomeObserved { evidence, .. } => match evidence {
EffectEvidence::Applied => AttemptStatus::Applied,
EffectEvidence::NotApplied => AttemptStatus::NotApplied,
},
EffectEvent::Verified { verification, .. } => match verification.result() {
VerificationResult::Satisfied => AttemptStatus::Verified,
VerificationResult::NotSatisfied => AttemptStatus::VerificationFailed,
},
};
let verification = match *event {
EffectEvent::Verified { verification, .. } => Some(verification),
_ => None,
};
let key = event.key();
match recovered.index.get(&key).copied() {
Some(index) => {
if let Some(attempt) = recovered.attempts.get_mut(index) {
if let Some(trail) = recovered.transitions.get_mut(index) {
trail.push(Transition {
from: Some(attempt.status),
to: status,
position,
});
}
attempt.status = status;
if verification.is_some() {
attempt.verification = verification;
}
}
}
None => {
let index = recovered.attempts.len();
let mut attempt = Attempt::new(key, status);
attempt.verification = verification;
recovered.attempts.push(attempt);
recovered.transitions.push(vec![Transition {
from: None,
to: status,
position,
}]);
recovered.index.insert(key, index);
}
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum ChainBreak {
Disagreement {
at: u64,
recorded: JournalPosition,
recomputed: JournalPosition,
},
Unencodable {
at: u64,
},
}
impl ChainBreak {
#[must_use]
pub const fn at(self) -> u64 {
match self {
Self::Disagreement { at, .. } | Self::Unencodable { at } => at,
}
}
#[must_use]
pub const fn recorded(self) -> Option<JournalPosition> {
match self {
Self::Disagreement { recorded, .. } => Some(recorded),
Self::Unencodable { .. } => None,
}
}
#[must_use]
pub const fn recomputed(self) -> Option<JournalPosition> {
match self {
Self::Disagreement { recomputed, .. } => Some(recomputed),
Self::Unencodable { .. } => None,
}
}
}
impl fmt::Display for ChainBreak {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match *self {
Self::Disagreement {
at,
recorded,
recomputed,
} => write!(
f,
"journal chain breaks at {at}: recorded {recorded}, recomputed {recomputed}"
),
Self::Unencodable { at } => write!(
f,
"journal entry {at} could not be re-encoded, so the chain could not be checked"
),
}
}
}
impl std::error::Error for ChainBreak {}
fn chain(previous: JournalPosition, event: &EffectEvent) -> Result<Digest, JournalError> {
let mut hasher = Hasher::new();
hasher.update(previous.head().as_bytes());
hasher.update(&event.to_bytes().map_err(JournalError::Encoding)?);
Ok(hasher.finalize())
}
pub(crate) fn chain_over_bytes(previous: &Digest, payload: &[u8]) -> Digest {
let mut hasher = Hasher::new();
hasher.update(previous.as_bytes());
hasher.update(payload);
hasher.finalize()
}
fn next_allowed_of(last: Option<EventKind>) -> Option<EventKind> {
last.map_or(Some(EventKind::IntentAdmitted), EventKind::next)
}
fn check_append_order(
expected_tail: JournalPosition,
actual: JournalPosition,
event: &EffectEvent,
last: Option<EventKind>,
) -> Result<(), JournalError> {
if expected_tail != actual {
let refusal = Err(JournalError::TailMismatch {
expected: expected_tail,
actual,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "check_append_order: the caller's view of the tail is stale");
return refusal;
}
let attempted = event.kind();
let expected = next_allowed_of(last);
if expected != Some(attempted) {
let refusal = Err(JournalError::OutOfOrder {
key: Box::new(event.key()),
expected,
attempted,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "check_append_order: the event does not follow the ladder");
return refusal;
}
Ok(())
}
pub fn verify_chain(entries: &[JournalEntry]) -> Result<JournalPosition, ChainBreak> {
verify_chain_from(JournalPosition::genesis(), entries)
}
pub fn verify_chain_from(
base: JournalPosition,
entries: &[JournalEntry],
) -> Result<JournalPosition, ChainBreak> {
let mut position = base;
for entry in entries {
let sequence = position.sequence().saturating_add(1);
let head = match chain(position, entry.event()) {
Ok(head) => head,
Err(_) => {
let refusal = Err(ChainBreak::Unencodable { at: sequence });
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "verify_chain: returning an error to the caller");
return refusal;
}
};
let recomputed = JournalPosition { sequence, head };
let recorded = entry.position();
if recorded.sequence() != sequence || !recorded.head().ct_eq(&head) {
let refusal = Err(ChainBreak::Disagreement {
at: sequence,
recorded,
recomputed,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "verify_chain: returning an error to the caller");
return refusal;
}
position = recomputed;
}
Ok(position)
}
#[derive(Debug, Clone)]
pub struct MemoryJournal {
committed: Vec<JournalEntry>,
ladder: std::collections::HashMap<EffectKey, EventKind>,
outcomes: std::collections::HashMap<EffectKey, (JournalPosition, EffectEvidence)>,
committed_bytes: u64,
}
impl Default for MemoryJournal {
fn default() -> Self {
Self::new()
}
}
impl MemoryJournal {
#[must_use]
pub fn new() -> Self {
Self {
committed: Vec::new(),
ladder: std::collections::HashMap::new(),
outcomes: std::collections::HashMap::new(),
committed_bytes: 0,
}
}
#[must_use]
pub fn committed(&self) -> &[JournalEntry] {
&self.committed
}
pub fn events(&self) -> impl Iterator<Item = &EffectEvent> {
self.committed.iter().map(JournalEntry::event)
}
pub fn verify(&self) -> Result<JournalPosition, ChainBreak> {
verify_chain(&self.committed)
}
#[must_use]
pub fn recover(&self) -> Recovered {
recover(self.events())
}
}
impl EffectJournal for MemoryJournal {
fn durability(&self) -> DurabilityPromise {
DurabilityPromise::Ephemeral
}
fn tail(&self) -> JournalPosition {
match self.committed.last() {
Some(entry) => entry.position(),
None => JournalPosition::genesis(),
}
}
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(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: EffectKey,
) -> Result<Option<(JournalPosition, EffectEvidence)>, JournalError> {
Ok(self.outcomes.get(&key).copied())
}
fn continuation_watermark(&self) -> Result<ContinuationWatermark, JournalError> {
let events = u64::saturating_from(self.committed.len());
let limit = u64::saturating_from(MAX_JOURNAL_EVENTS);
Ok(ContinuationWatermark::inert(
events,
limit,
self.committed_bytes,
MAX_JOURNAL_BYTES,
))
}
fn compare_and_append(
&mut self,
expected_tail: JournalPosition,
event: &EffectEvent,
) -> Result<DurableAck, JournalError> {
let actual = self.tail();
let key = event.key();
check_append_order(expected_tail, actual, event, self.ladder.get(&key).copied())?;
let requested_events = u64::saturating_from(self.committed.len()).saturating_add(1);
let event_limit = u64::saturating_from(MAX_JOURNAL_EVENTS);
if requested_events > event_limit {
let refusal = Err(JournalError::CapacityExceeded {
resource: JournalLimitKind::Events,
limit: event_limit,
requested: requested_events,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "compare_and_append: returning an error to the caller");
return refusal;
}
let payload_len =
u64::saturating_from(event.to_bytes().map_err(JournalError::Encoding)?.len());
let requested_bytes = self
.committed_bytes
.saturating_add(payload_len)
.saturating_add(36);
if requested_bytes > MAX_JOURNAL_BYTES {
let refusal = Err(JournalError::CapacityExceeded {
resource: JournalLimitKind::Bytes,
limit: MAX_JOURNAL_BYTES,
requested: requested_bytes,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "compare_and_append: returning an error to the caller");
return refusal;
}
let sequence = actual
.sequence()
.checked_add(1)
.ok_or(JournalError::Exhausted)?;
let position = JournalPosition {
sequence,
head: chain(actual, event)?,
};
self.committed.push(JournalEntry::new(position, *event));
self.committed_bytes = requested_bytes;
self.ladder.insert(key, event.kind());
if let EffectEvent::OutcomeObserved { evidence, .. } = *event {
self.outcomes.insert(key, (position, evidence));
}
Ok(DurableAck::new(position, self.durability()))
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::effect::{
ActionDigest, ActionId, AttemptId, EnvironmentEpoch, EnvironmentId, FlowRevision, RunId,
};
const RUN: &str = "0102030405060708090a0b0c0d0e0f10";
const ACTION: &str = "1112131415161718191a1b1c1d1e1f20";
const OTHER_ACTION: &str = "3132333435363738393a3b3c3d3e3f40";
const ENV: &str = "2122232425262728292a2b2c2d2e2f30";
const FLOW_HEX: &str = "000102030405060708090a0b0c0d0e0f101112131415161718191a1b1c1d1e1f";
const DIGEST_HEX: &str = "f0f1f2f3f4f5f6f7f8f9fafbfcfdfeffe0e1e2e3e4e5e6e7e8e9eaebecedeeef";
const PREDICATE: &str = "4142434445464748494a4b4c4d4e4f50";
type TestResult = Result<(), Box<dyn std::error::Error>>;
fn key(attempt: &str, epoch: &str) -> Result<EffectKey, Box<dyn std::error::Error>> {
build(ACTION, attempt, epoch)
}
fn other_key(attempt: &str, epoch: &str) -> Result<EffectKey, Box<dyn std::error::Error>> {
build(OTHER_ACTION, attempt, epoch)
}
fn build(
action: &str,
attempt: &str,
epoch: &str,
) -> Result<EffectKey, Box<dyn std::error::Error>> {
let key_run = RunId::from_hex(RUN)?;
let key_action = ActionId::from_hex(action)?;
let key_attempt = AttemptId::from_decimal(attempt)?;
let key_flow_revision = FlowRevision::from_tagged("blake3_256", FLOW_HEX)?;
let key_digest = ActionDigest::from_tagged("blake3_256", DIGEST_HEX)?;
let key_environment = EnvironmentId::from_hex(ENV)?;
let key_epoch = EnvironmentEpoch::from_decimal(epoch)?;
Ok(
crate::effect::EffectIdentity::new(key_run, key_environment, key_flow_revision).key(
key_action,
key_attempt,
key_digest,
key_epoch,
),
)
}
fn satisfied() -> Result<Verification, Box<dyn std::error::Error>> {
Ok(Verification::new(
Id128::from_hex(PREDICATE)?,
1,
blake3(b"observations"),
VerificationResult::Satisfied,
))
}
fn append(
journal: &mut MemoryJournal,
event: EffectEvent,
) -> Result<DurableAck, Box<dyn std::error::Error>> {
let tail = journal.tail();
Ok(journal.compare_and_append(tail, &event)?)
}
#[test]
fn memory_journal_refuses_history_beyond_its_declared_limit() -> TestResult {
let mut journal = MemoryJournal::new();
for attempt in 1_u64..=100_000 {
let key = key(&attempt.to_string(), "1")?;
journal.compare_and_append(journal.tail(), &EffectEvent::IntentAdmitted { key })?;
}
let extra_key = key("100001", "1")?;
let extra = journal.compare_and_append(
journal.tail(),
&EffectEvent::IntentAdmitted { key: extra_key },
);
assert!(
matches!(
&extra,
Err(JournalError::CapacityExceeded {
resource: JournalLimitKind::Events,
limit: 100_000,
requested: 100_001,
})
),
"the refusal names both the hard limit and requested size: {extra:?}"
);
assert_eq!(
journal.committed().len(),
100_000,
"history must remain complete at the declared limit; extra append was {extra:?}"
);
Ok(())
}
fn admit_and_prepare(
journal: &mut MemoryJournal,
key: EffectKey,
) -> Result<(), Box<dyn std::error::Error>> {
append(journal, EffectEvent::IntentAdmitted { key })?;
append(journal, EffectEvent::DispatchPrepared { key })?;
Ok(())
}
fn applied_journal(key: EffectKey) -> Result<MemoryJournal, Box<dyn std::error::Error>> {
let mut journal = MemoryJournal::new();
admit_and_prepare(&mut journal, key)?;
append(
&mut journal,
EffectEvent::OutcomeObserved {
key,
evidence: EffectEvidence::Applied,
},
)?;
Ok(journal)
}
#[test]
fn genesis_is_not_a_zero_head() -> TestResult {
let genesis = JournalPosition::genesis();
assert_eq!(genesis.sequence(), 0);
assert_ne!(genesis.head(), blake3(&[]));
assert_eq!(genesis, JournalPosition::genesis());
Ok(())
}
#[test]
fn an_append_advances_the_sequence_and_moves_the_head() -> TestResult {
let key = key("1", "1")?;
let mut journal = MemoryJournal::new();
let genesis = journal.tail();
let ack = append(&mut journal, EffectEvent::IntentAdmitted { key })?;
assert_eq!(ack.position().sequence(), 1);
assert_ne!(ack.position().head(), genesis.head());
assert_eq!(journal.tail(), ack.position());
Ok(())
}
#[test]
fn the_acknowledgment_names_the_promise_that_was_claimed() -> TestResult {
let key = key("1", "1")?;
let mut journal = MemoryJournal::new();
let ack = append(&mut journal, EffectEvent::IntentAdmitted { key })?;
assert_eq!(ack.promise(), DurabilityPromise::Ephemeral);
Ok(())
}
#[test]
fn positioned_readback_resolves_one_exact_sequence_and_head() -> TestResult {
let key = key("1", "1")?;
let mut journal = MemoryJournal::new();
let ack = append(&mut journal, EffectEvent::IntentAdmitted { key })?;
let entries = EffectJournal::committed_entries(&journal)?;
let expected = entries.first().copied().ok_or("committed entry missing")?;
assert_eq!(
journal.committed_entry(ack.position())?,
Some(expected),
"the acknowledged position returns its exact committed entry"
);
let wrong_head = JournalPosition {
sequence: ack.position().sequence(),
head: blake3(b"not the committed head"),
};
assert_eq!(
journal.committed_entry(wrong_head)?,
None,
"a matching sequence cannot authorize a different history head"
);
Ok(())
}
#[test]
fn two_events_in_another_order_are_a_different_head() -> TestResult {
let first = key("1", "1")?;
let second = other_key("1", "1")?;
let mut forwards = MemoryJournal::new();
append(&mut forwards, EffectEvent::IntentAdmitted { key: first })?;
append(&mut forwards, EffectEvent::IntentAdmitted { key: second })?;
let mut backwards = MemoryJournal::new();
append(&mut backwards, EffectEvent::IntentAdmitted { key: second })?;
append(&mut backwards, EffectEvent::IntentAdmitted { key: first })?;
assert_eq!(forwards.tail().sequence(), backwards.tail().sequence());
assert_ne!(forwards.tail().head(), backwards.tail().head());
Ok(())
}
#[test]
fn a_stale_tail_is_refused_and_changes_nothing() -> TestResult {
let key = key("1", "1")?;
let mut journal = MemoryJournal::new();
let genesis = journal.tail();
let committed = append(&mut journal, EffectEvent::IntentAdmitted { key })?;
let refused = journal.compare_and_append(genesis, &EffectEvent::DispatchPrepared { key });
match refused {
Err(JournalError::TailMismatch { expected, actual }) => {
assert_eq!(expected, genesis);
assert_eq!(actual, committed.position());
}
other => return Err(format!("expected a tail mismatch, got {other:?}").into()),
}
assert_eq!(journal.committed().len(), 1);
assert_eq!(journal.tail(), committed.position());
Ok(())
}
#[test]
fn a_second_controller_reading_the_same_tail_is_fenced() -> TestResult {
let key = key("1", "1")?;
let mut journal = MemoryJournal::new();
let both_saw = journal.tail();
let winner = journal.compare_and_append(both_saw, &EffectEvent::IntentAdmitted { key })?;
let loser = journal.compare_and_append(both_saw, &EffectEvent::DispatchPrepared { key });
assert!(
matches!(loser, Err(JournalError::TailMismatch { .. })),
"the second writer from a stale reading must be refused"
);
assert_eq!(journal.committed().len(), 1);
assert_eq!(journal.tail(), winner.position());
Ok(())
}
#[test]
fn a_dispatch_before_its_intent_is_refused() -> TestResult {
let key = key("1", "1")?;
let mut journal = MemoryJournal::new();
let refused =
journal.compare_and_append(journal.tail(), &EffectEvent::DispatchPrepared { key });
match refused {
Err(JournalError::OutOfOrder {
key: refused_key,
expected,
attempted,
}) => {
assert_eq!(*refused_key, key);
assert_eq!(expected, Some(EventKind::IntentAdmitted));
assert_eq!(attempted, EventKind::DispatchPrepared);
}
other => return Err(format!("expected an ordering refusal, got {other:?}").into()),
}
assert!(journal.committed().is_empty());
Ok(())
}
#[test]
fn the_ladder_is_walked_once_per_key() -> TestResult {
let key = key("1", "1")?;
let mut journal = MemoryJournal::new();
append(&mut journal, EffectEvent::IntentAdmitted { key })?;
let refused =
journal.compare_and_append(journal.tail(), &EffectEvent::IntentAdmitted { key });
assert!(
matches!(
refused,
Err(JournalError::OutOfOrder {
expected: Some(EventKind::DispatchPrepared),
attempted: EventKind::IntentAdmitted,
..
})
),
"a second intent for one key is how a blind resend becomes representable"
);
assert_eq!(journal.committed().len(), 1);
Ok(())
}
#[test]
fn a_verified_attempt_refuses_everything_but_a_later_verification() -> TestResult {
let key = key("1", "1")?;
let mut journal = applied_journal(key)?;
append(
&mut journal,
EffectEvent::Verified {
key,
verification: satisfied()?,
},
)?;
let refused = journal.compare_and_append(
journal.tail(),
&EffectEvent::OutcomeObserved {
key,
evidence: EffectEvidence::NotApplied,
},
);
assert!(
matches!(
refused,
Err(JournalError::OutOfOrder {
expected: Some(EventKind::Verified),
attempted: EventKind::OutcomeObserved,
..
})
),
"only a later verification follows a verification; the effect is not re-settled"
);
Ok(())
}
#[test]
fn an_ephemeral_journal_reports_itself_and_refuses_the_boundary() -> TestResult {
let journal = MemoryJournal::new();
assert_eq!(journal.durability(), DurabilityPromise::Ephemeral);
match journal.admit_external_handoff() {
Err(JournalError::PromiseUnmet { required, offered }) => {
assert_eq!(required, DurabilityPromise::ProcessCrash);
assert_eq!(offered, DurabilityPromise::Ephemeral);
}
other => {
return Err(format!(
"an in-memory journal must not host an external handoff, got {other:?}"
)
.into());
}
}
Ok(())
}
#[test]
fn only_a_crash_surviving_promise_clears_the_boundary() {
assert!(!DurabilityPromise::Ephemeral.survives_process_crash());
assert!(DurabilityPromise::ProcessCrash.survives_process_crash());
assert!(DurabilityPromise::PowerLoss.survives_process_crash());
assert!(DurabilityPromise::PowerLoss > DurabilityPromise::ProcessCrash);
}
#[test]
fn a_prepared_dispatch_with_no_result_recovers_unknown() -> TestResult {
let key = key("1", "1")?;
let mut journal = MemoryJournal::new();
admit_and_prepare(&mut journal, key)?;
let recovered = journal.recover();
assert_eq!(recovered.len(), 1);
assert_eq!(recovered.status(key), Some(AttemptStatus::OutcomeUnknown));
assert_eq!(recovered.uncertain(), vec![key]);
Ok(())
}
#[test]
fn an_observed_not_applied_settles_the_unknown() -> TestResult {
let key = key("1", "1")?;
let mut journal = MemoryJournal::new();
admit_and_prepare(&mut journal, key)?;
append(
&mut journal,
EffectEvent::OutcomeObserved {
key,
evidence: EffectEvidence::NotApplied,
},
)?;
let recovered = journal.recover();
assert_eq!(recovered.status(key), Some(AttemptStatus::NotApplied));
assert!(
recovered.uncertain().is_empty(),
"evidence that the effect did not land is what ends the uncertainty"
);
Ok(())
}
#[test]
fn an_admitted_intent_alone_is_not_uncertain() -> TestResult {
let key = key("1", "1")?;
let mut journal = MemoryJournal::new();
append(&mut journal, EffectEvent::IntentAdmitted { key })?;
assert_eq!(journal.recover().status(key), Some(AttemptStatus::Prepared));
assert!(journal.recover().uncertain().is_empty());
Ok(())
}
#[test]
fn a_satisfied_predicate_is_the_only_verified_state() -> TestResult {
let key = key("1", "1")?;
let mut journal = applied_journal(key)?;
append(
&mut journal,
EffectEvent::Verified {
key,
verification: satisfied()?,
},
)?;
assert_eq!(journal.recover().status(key), Some(AttemptStatus::Verified));
Ok(())
}
#[test]
fn a_predicate_that_did_not_hold_is_not_a_verification() -> TestResult {
let key = key("1", "1")?;
let mut journal = applied_journal(key)?;
append(
&mut journal,
EffectEvent::Verified {
key,
verification: Verification::new(
Id128::from_hex(PREDICATE)?,
1,
blake3(b"observations"),
VerificationResult::NotSatisfied,
),
},
)?;
let status = journal.recover().status(key);
assert_eq!(status, Some(AttemptStatus::VerificationFailed));
assert_ne!(status, Some(AttemptStatus::Applied));
assert_ne!(status, Some(AttemptStatus::Verified));
assert!(
!status.is_some_and(AttemptStatus::is_uncertain),
"the effect is known to have landed, so a failed predicate is not uncertainty"
);
Ok(())
}
fn verdict(
version: u64,
result: VerificationResult,
) -> Result<Verification, Box<dyn std::error::Error>> {
Ok(Verification::new(
Id128::from_hex(PREDICATE)?,
version,
blake3(format!("observations at version {version}").as_bytes()),
result,
))
}
#[test]
fn a_verification_can_be_revised_in_both_directions() -> TestResult {
let key = key("1", "1")?;
let mut journal = applied_journal(key)?;
for (version, result) in [
(1, VerificationResult::Satisfied),
(2, VerificationResult::NotSatisfied),
(3, VerificationResult::Satisfied),
] {
append(
&mut journal,
EffectEvent::Verified {
key,
verification: verdict(version, result)?,
},
)?;
}
let recovered = journal.recover();
assert_eq!(recovered.status(key), Some(AttemptStatus::Verified));
assert_eq!(
recovered
.verification(key)
.map(|found| found.predicate_version()),
Some(3),
"the status is qualified by the version it was decided at"
);
let trail: Vec<_> = recovered
.history(key)
.iter()
.map(|change| (change.from(), change.to(), change.position()))
.collect();
assert_eq!(
trail,
vec![
(None, AttemptStatus::Prepared, 0),
(
Some(AttemptStatus::Prepared),
AttemptStatus::OutcomeUnknown,
1
),
(
Some(AttemptStatus::OutcomeUnknown),
AttemptStatus::Applied,
2
),
(Some(AttemptStatus::Applied), AttemptStatus::Verified, 3),
(
Some(AttemptStatus::Verified),
AttemptStatus::VerificationFailed,
4
),
(
Some(AttemptStatus::VerificationFailed),
AttemptStatus::Verified,
5
),
],
"what changed, what it superseded, and where in the journal"
);
Ok(())
}
#[test]
fn recovery_lists_attempts_in_admission_order() -> TestResult {
let first = key("1", "1")?;
let second = other_key("2", "1")?;
let mut journal = MemoryJournal::new();
admit_and_prepare(&mut journal, second)?;
admit_and_prepare(&mut journal, first)?;
let attempts = journal.recover().attempts().to_vec();
assert_eq!(attempts.len(), 2);
assert_eq!(attempts[0].key(), second);
assert_eq!(attempts[1].key(), first);
Ok(())
}
#[test]
fn verify_accepts_an_untampered_chain() -> TestResult {
let key = key("1", "1")?;
let mut journal = MemoryJournal::new();
admit_and_prepare(&mut journal, key)?;
assert_eq!(journal.verify()?, journal.tail());
Ok(())
}
#[test]
fn verify_rejects_a_rewritten_position() -> TestResult {
let key = key("1", "1")?;
let mut journal = MemoryJournal::new();
admit_and_prepare(&mut journal, key)?;
let mut tampered = journal.committed().to_vec();
let last = tampered
.pop()
.ok_or("the journal should hold two entries after admit and prepare")?;
tampered.push(JournalEntry::new(JournalPosition::genesis(), *last.event()));
match verify_chain(&tampered) {
Err(broken) => assert_eq!(broken.at(), 2),
Ok(position) => {
return Err(
format!("a rewritten position must be detected, got {position}").into(),
);
}
}
Ok(())
}
#[test]
fn verify_rejects_a_reordered_chain() -> TestResult {
let key = key("1", "1")?;
let mut journal = MemoryJournal::new();
admit_and_prepare(&mut journal, key)?;
let mut reordered = journal.committed().to_vec();
reordered.reverse();
match verify_chain(&reordered) {
Err(broken) => assert_eq!(broken.at(), 1),
Ok(position) => {
return Err(format!("a reordered chain must be detected, got {position}").into());
}
}
Ok(())
}
#[test]
fn a_key_encodes_deterministically() -> TestResult {
let key = key("1", "1")?;
let bytes = key.to_bytes()?;
assert_eq!(
bytes.as_slice(),
key.to_bytes()?.as_slice(),
"one key encodes the same way twice"
);
assert!(
!bytes.is_empty(),
"and the encoding carries something: {bytes:?}"
);
Ok(())
}
#[test]
fn keys_differing_in_one_field_differ_in_bytes() -> TestResult {
let base = key("1", "1")?;
let later_attempt = key("2", "1")?;
let later_epoch = key("1", "2")?;
let other = other_key("1", "1")?;
assert_ne!(
base.to_bytes()?.as_slice(),
later_attempt.to_bytes()?.as_slice()
);
assert_ne!(
base.to_bytes()?.as_slice(),
later_epoch.to_bytes()?.as_slice()
);
assert_ne!(base.to_bytes()?.as_slice(), other.to_bytes()?.as_slice());
Ok(())
}
#[test]
fn an_event_encoding_carries_its_kind_and_its_key() -> TestResult {
let key = key("1", "1")?;
let admitted = EffectEvent::IntentAdmitted { key }.to_bytes()?;
let prepared = EffectEvent::DispatchPrepared { key }.to_bytes()?;
let applied = EffectEvent::OutcomeObserved {
key,
evidence: EffectEvidence::Applied,
}
.to_bytes()?;
let not_applied = EffectEvent::OutcomeObserved {
key,
evidence: EffectEvidence::NotApplied,
}
.to_bytes()?;
assert_ne!(admitted.as_slice(), prepared.as_slice());
assert_ne!(applied.as_slice(), not_applied.as_slice());
for (bytes, event) in [
(admitted, EffectEvent::IntentAdmitted { key }),
(prepared, EffectEvent::DispatchPrepared { key }),
(
applied,
EffectEvent::OutcomeObserved {
key,
evidence: EffectEvidence::Applied,
},
),
(
not_applied,
EffectEvent::OutcomeObserved {
key,
evidence: EffectEvidence::NotApplied,
},
),
] {
assert_eq!(
lgwks_std::wire::from_bytes::<EffectEvent, WireError>(&bytes)?,
event
);
}
Ok(())
}
#[test]
fn a_verification_encoding_names_its_predicate_and_its_version() -> TestResult {
let key = key("1", "1")?;
let first = EffectEvent::Verified {
key,
verification: satisfied()?,
}
.to_bytes()?;
let next_version = EffectEvent::Verified {
key,
verification: Verification::new(
Id128::from_hex(PREDICATE)?,
2,
blake3(b"observations"),
VerificationResult::Satisfied,
),
}
.to_bytes()?;
assert_ne!(first.as_slice(), next_version.as_slice());
assert_eq!(
lgwks_std::wire::from_bytes::<EffectEvent, WireError>(&first)?,
EffectEvent::Verified {
key,
verification: satisfied()?,
}
);
assert_eq!(
lgwks_std::wire::from_bytes::<EffectEvent, WireError>(&next_version)?,
EffectEvent::Verified {
key,
verification: Verification::new(
Id128::from_hex(PREDICATE)?,
2,
blake3(b"observations"),
VerificationResult::Satisfied,
),
}
);
Ok(())
}
}