Skip to main content

StagedAtomic

Struct StagedAtomic 

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

A journal store that lends a (simulated) transaction.

Wraps any store and adds the capability. Every other method delegates, so a test gets a real journal with one extra thing it can do.

Implementations§

Source§

impl StagedAtomic

Source

pub fn wrap(inner: Arc<dyn JournalStore>) -> Arc<Self> ⓘ

Source

pub fn lose_commit_acknowledgement(&self)

From now on, every commit lands and its acknowledgement vanishes.

Models the nastier half of the in-doubt window: the transaction committed — statements applied, records appended — and the client was never told. This is the world in which a cheap abort is wrong: the writes are standing, permanent, with no reversal registered and none possible, and a group that settles Aborted over them has the journal claiming taken back whole about work nobody took back. The runtime’s only honest answer is quarantine, and this switch is how a test asks the question.

Source

pub fn applied(&self) -> Vec<Statement> ⓘ

Statements that were actually applied, in order.

Empty for a unit of work that refused: its statements were staged and discarded, which is what a rollback would have done to them.

Trait Implementations§

Source§

impl AtomicJournal for StagedAtomic

Source§

fn append_atomic<'life0, 'life1, 'async_trait>( &'life0 self, _run: RunId, epoch: Epoch, work: &'life1 dyn AtomicWork, ) -> Pin<Box<dyn Future<Output = Result<Vec<Record>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Run work and append what it returns, in one transaction, fenced by epoch exactly as an ordinary append is. Read more
Source§

impl Debug for StagedAtomic

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

impl JournalStore for StagedAtomic

Source§

fn is_shared(&self) -> bool

Whether more than one plane instance can write to this store. Read more
Source§

fn seals(&self) -> bool

Whether this store seals payloads on append, rewriting what it is handed. Read more
Source§

fn tenant(&self) -> &str

Whose rows this handle can reach. Read more
Source§

fn atomic(&self) -> Option<&dyn AtomicJournal>

This store’s own transaction, when a co-located resource can join it. Read more
Source§

fn append<'life0, 'async_trait>( &'life0 self, epoch: Epoch, batch: Vec<Append>, ) -> Pin<Box<dyn Future<Output = Result<Vec<Record>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Append a batch, sealing each record into the chain. Read more
Source§

fn read<'life0, 'async_trait>( &'life0 self, run: RunId, from: Seq, ) -> Pin<Box<dyn Future<Output = Result<Vec<Record>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Read a run’s records from from (inclusive, 1-based) onward. Read more
Source§

fn read_page<'life0, 'async_trait>( &'life0 self, run: RunId, from: Seq, limit: usize, ) -> Pin<Box<dyn Future<Output = Result<Vec<Record>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

At most limit of a run’s records from from (inclusive, 1-based). Read more
Source§

fn acquire<'life0, 'life1, 'async_trait>( &'life0 self, run: RunId, owner: &'life1 str, ttl: Duration, ) -> Pin<Box<dyn Future<Output = Result<Lease, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Take ownership of a run, returning the fencing epoch to write under. Read more
Source§

fn renew<'life0, 'life1, 'async_trait>( &'life0 self, run: RunId, owner: &'life1 str, epoch: Epoch, ttl: Duration, ) -> Pin<Box<dyn Future<Output = Result<Lease, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Extend a lease this caller still holds, without ever claiming one. Read more
Source§

fn release_lease<'life0, 'async_trait>( &'life0 self, run: RunId, epoch: Epoch, ) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Hand a lease back, so the next instance need not wait out the TTL. Read more
Source§

fn abandoned_runs<'life0, 'async_trait>( &'life0 self, limit: usize, ) -> Pin<Box<dyn Future<Output = Result<Vec<RunId>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Runs whose lease expired without being released — the runs an instance died holding. Read more
Source§

fn waiting_runs<'life0, 'async_trait>( &'life0 self, limit: usize, ) -> Pin<Box<dyn Future<Output = Result<Vec<WaitingRun>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

The runs that are waiting, and what each one waits for. Read more
Source§

fn admitted_as<'life0, 'life1, 'async_trait>( &'life0 self, key: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<Option<RunId>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

The run that holds this admission key, if one does. Read more
Source§

fn forget_admissions<'life0, 'async_trait>( &'life0 self, older_than: Timestamp, ) -> Pin<Box<dyn Future<Output = Result<usize, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Retire admission keys claimed before older_than. Returns how many. Read more
Source§

fn runs_by_outcome<'life0, 'life1, 'async_trait>( &'life0 self, outcome: &'life1 str, limit: usize, ) -> Pin<Box<dyn Future<Output = Result<Vec<RunId>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Concluded runs whose latest conclusion is outcome, newest first. Read more
Source§

fn count_by_outcome<'life0, 'life1, 'async_trait>( &'life0 self, outcome: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<u64, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

How many runs currently rest on this conclusion. Read more
Source§

fn runs_by_id<'life0, 'async_trait>( &'life0 self, after: Option<RunId>, limit: usize, ) -> Pin<Box<dyn Future<Output = Result<Vec<RunId>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Every run this tenant holds records for, by run id ascending, one bounded page strictly after after. Read more
Source§

fn recent_runs<'life0, 'async_trait>( &'life0 self, after: Option<(u64, RunId)>, limit: usize, ) -> Pin<Box<dyn Future<Output = Result<Vec<(RunId, u64)>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Runs ordered by last durable append, newest first, one bounded page. Read more
Source§

fn recent_runs_from<'life0, 'life1, 'async_trait>( &'life0 self, source: &'life1 str, after: Option<(u64, RunId)>, limit: usize, ) -> Pin<Box<dyn Future<Output = Result<Vec<(RunId, u64)>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

recent_runs, narrowed to the runs one producer admitted: those whose RunAdmitted carries an admission key whose source half (origin_source) is source. Read more
Source§

fn case_history<'life0, 'async_trait>( &'life0 self, case: CaseId, limit: usize, ) -> Pin<Box<dyn Future<Output = Result<Vec<Record>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Every record belonging to a case, oldest first. Read more
Source§

fn head<'life0, 'async_trait>( &'life0 self, run: RunId, ) -> Pin<Box<dyn Future<Output = Result<Head, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

The run’s current chain head.
Source§

fn seal<'life0, 'life1, 'async_trait>( &'life0 self, run: RunId, epoch: Epoch, outcome: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<Digest, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Close the chain and return its terminal hash — what a signature covers. Read more
Source§

fn checkpoint<'life0, 'async_trait>( &'life0 self, ) -> Pin<Box<dyn Future<Output = Result<Checkpoint, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

A commitment to the set of sealed runs. Read more
Source§

fn consistency_proof<'life0, 'async_trait>( &'life0 self, old_size: u64, ) -> Pin<Box<dyn Future<Output = Result<Vec<Digest>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Prove the log has only grown since a checkpoint of old_size. Read more
Source§

fn inclusion_proof<'life0, 'async_trait>( &'life0 self, run: RunId, ) -> Pin<Box<dyn Future<Output = Result<Option<Inclusion>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Prove a sealed run is in the log this checkpoint commits to. Read more
Source§

fn inclusion_proof_at<'life0, 'async_trait>( &'life0 self, run: RunId, size: u64, ) -> Pin<Box<dyn Future<Output = Result<Option<Inclusion>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Prove a sealed run is in the log as it stood at size leaves — the checkpoint an earlier checkpoint returned, not the live one. Read more
Source§

fn log_positions<'life0, 'life1, 'async_trait>( &'life0 self, runs: &'life1 [RunId], ) -> Pin<Box<dyn Future<Output = Result<Vec<Option<(u64, Digest)>>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Where each of runs sits in the log — its index and leaf — or None for a run that is not sealed, in the order asked. Read more
Source§

fn request_cancel<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, run: RunId, actor: &'life1 Operator, reason: &'life2 str, ) -> Pin<Box<dyn Future<Output = Result<bool, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait,

Ask a run to stop, durably. Read more
Source§

fn cancellation<'life0, 'async_trait>( &'life0 self, run: RunId, ) -> Pin<Box<dyn Future<Output = Result<Option<Cancellation>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

The pending stop request for a run, if one was made. Read more
Source§

fn verify<'life0, 'async_trait>( &'life0 self, run: RunId, ) -> Pin<Box<dyn Future<Output = Result<Digest, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Verify a run’s chain end to end.

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<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

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> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

impl<Unshared, Shared> IntoShared<Shared> for Unshared
where Shared: FromUnshared<Unshared>,

Source§

fn into_shared(self) -> Shared

Creates a shared type from an unshared type.
Source§

impl<T> MaybeSend for T
where T: Send,

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> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

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<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

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