Skip to main content

FileKernelJournal

Struct FileKernelJournal 

Source
pub struct FileKernelJournal { /* private fields */ }
Expand description

Cross-process atomic reference implementation of KernelJournal (spec Task 8b).

The atomicity primitive is POSIX link(2) (std::fs::hard_link, reached through tokio::fs::hard_link so the blocking call stays off the reactor): content is written to a private temp file and fsynced, then hard-linked into its final name. link fails with AlreadyExists if the name is taken, and it publishes already-complete content — so a crash can never leave a half-written .rec, only an orphan temp file that the naming rule ignores. The journal root must therefore live on a filesystem that supports hard links; that requirement is the price of real CAS.

Why the record filename is <step_seq>.rec and contains no digest. The filename is the collision domain. Two writers racing on the same head both compute the same next step_seq, so they contend for one name and exactly one wins. Folding a per-writer value (the new record’s digest) into the name would give the racers different names — both links would succeed and the chain would fork. Only the predecessor-determined part of the identity may appear in the name (spec Task 8b correction (a)).

The pre-link head check is not a TOCTOU hole: it can only reject an append that link would have accepted (a stale expected_head whose step_seq slot happens to be free), never accept one link would have rejected. Every acceptance is still decided by the atomic link — the publish half is factored into FileKernelJournal::publish_record / FileKernelJournal::publish_checkpoint so that claim is directly testable without the pre-check in front of it.

Checkpoint installs use the same primitive on a separate ordinal space (<ordinal>.ckpt), so two processes installing on the same predecessor also contend for one name.

Implementations§

Source§

impl FileKernelJournal

Source

pub fn new(root: impl AsRef<Path>) -> Self

Trait Implementations§

Source§

impl KernelJournal for FileKernelJournal

Source§

fn stage_outbound_envelope<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, operation_id: &'life1 str, envelope_json: &'life2 str, ) -> Pin<Box<dyn Future<Output = JournalResult<()>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait,

Persist the byte-identical canonical input envelope before attempting its record append. Read more
Source§

fn read_outbound_envelope<'life0, 'life1, 'async_trait>( &'life0 self, operation_id: &'life1 str, ) -> Pin<Box<dyn Future<Output = JournalResult<Option<String>>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Return the staged outbound envelope, if an append-before crash left one behind.
Source§

fn clear_outbound_envelope<'life0, 'life1, 'async_trait>( &'life0 self, operation_id: &'life1 str, ) -> Pin<Box<dyn Future<Output = JournalResult<()>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Clear a staged outbound envelope after its input is durably owned or rejected.
Source§

fn compare_and_append<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, operation_id: &'life1 str, expected_head: Option<&'life2 str>, record: JournalRecordInput, ) -> Pin<Box<dyn Future<Output = JournalResult<JournalAppendReceipt>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait,

Atomically append record iff the operation’s head is exactly expected_head. Read more
Source§

fn head<'life0, 'life1, 'async_trait>( &'life0 self, operation_id: &'life1 str, ) -> Pin<Box<dyn Future<Output = JournalResult<Option<JournalHead>>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

The current head, or None when the operation has no records and no pruned anchor.
Source§

fn read_from<'life0, 'life1, 'async_trait>( &'life0 self, operation_id: &'life1 str, from_step_seq: u64, ) -> Pin<Box<dyn Future<Output = JournalResult<Vec<JournalEntry>>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Records with step_seq >= from_step_seq, in chain order. 0 returns the whole retained chain (Rust has no default arguments; node’s optional parameter defaults to the same).
Source§

fn records_after<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, operation_id: &'life1 str, after_head: Option<&'life2 str>, ) -> Pin<Box<dyn Future<Output = JournalResult<Vec<JournalEntry>>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait,

Records strictly after the record whose digest is after_head — the digest-anchored cursor §9.1 names records_after(operation_id, checkpoint_head). None returns everything retained. Read more
Source§

fn compare_and_install_checkpoint<'life0, 'life1, 'life2, 'life3, 'async_trait>( &'life0 self, operation_id: &'life1 str, previous_checkpoint_id: Option<&'life2 str>, covered_head: &'life3 str, checkpoint: CheckpointCandidate, ) -> Pin<Box<dyn Future<Output = JournalResult<InstalledCheckpoint>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait, 'life3: 'async_trait,

Atomically install checkpoint iff the operation’s checkpoint pointer is exactly previous_checkpoint_id, and covered_head is the record digest at checkpoint.through_step_seq. Read more
Source§

fn latest_checkpoint<'life0, 'life1, 'async_trait>( &'life0 self, operation_id: &'life1 str, ) -> Pin<Box<dyn Future<Output = JournalResult<Option<InstalledCheckpoint>>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

The highest-ordinal installed checkpoint, acknowledged or not.
Source§

fn ack_checkpoint<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, operation_id: &'life1 str, checkpoint_id: &'life2 str, ) -> Pin<Box<dyn Future<Output = JournalResult<InstalledCheckpoint>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait,

Record the durable acknowledgement that opens the prefix-reclamation boundary. Idempotent. Read more
Source§

fn prune_acked_prefix<'life0, 'life1, 'async_trait>( &'life0 self, operation_id: &'life1 str, ) -> Pin<Box<dyn Future<Output = JournalResult<JournalPruneReceipt>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Reclaim the record prefix covered by the latest acknowledged checkpoint. A no-op while no checkpoint is acknowledged. The pruned boundary is retained as an anchor so Self::head and the next CAS still resolve on a fully-pruned chain (Task 8b correction (d)).

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> PolicyExt for T
where T: ?Sized,

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. Read more
Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more