Skip to main content

StagedOp

Struct StagedOp 

Source
pub struct StagedOp<'inv> { /* private fields */ }
Expand description

One stage of a keyed subscriber’s processing of one event: an open op, whose work lands whichever exit is taken.

Implements AtomicOperation — use it like any atomic operation, then take exactly one exit:

exitmeaningcursor
commitprocessing of this event is doneadvances past the event
suspendthis stage is done, the chain continuesunmoved
pause_untilthis stage is done, processing pauses until atunmoved

The verbs are about the event, not the transaction: every exit lands this stage’s writes, and what differs is what happens to the cursor. commit means what IsolatedOp::commit means — work and checkpoint land together — so a keyed subscriber that never stages reads exactly like a singleton one: consume, work, commit. suspend and pause_until land the work and leave the checkpoint where it was.

As with IsolatedOp, there is no mutable access to the raw es_entity::DbOp — the op can only land through one of the exits, so no stage can commit without the runner knowing what it meant.

Sealed to its invocation by the same brand Handled carries, so a staged op cannot be stashed and resumed from a later event:

use obix::StagedOp;

struct Evil {
    stash: std::sync::Mutex<Option<StagedOp<'static>>>,
}

async fn stash_it(op: StagedOp<'_>, evil: &Evil) {
    // error[E0521]: borrowed data escapes outside of function
    *evil.stash.lock().unwrap() = Some(op);
}

Implementations§

Source§

impl<'inv> StagedOp<'inv>

Source

pub async fn suspend( self, ) -> Result<Suspended<'inv>, Box<dyn Error + Send + Sync>>

This stage is done, the event is not: the returned Suspended holds no open transaction, so the subscriber can await external I/O before resume opens the next stage’s op.

The cursor does not move — this event is still being processed, and a crash here replays it (with everything this stage committed already durable).

Source

pub fn pause_until(self, at: DateTime<Utc>) -> Handled<'inv>

This stage is done and processing pauses, with the cursor still parked before this event, until at.

The op’s work lands, but the checkpoint the runner folds in is still pre-this-event: the next run re-reads the event and re-evaluates. The resume time is domain knowledge (a retry schedule owned by the consumer’s entities) — the one fact obix cannot derive.

Source

pub fn commit(self) -> Handled<'inv>

Processing of this event is done: the runner folds the checkpoint at this event’s sequence into the same transaction, so work and cursor advance together — what IsolatedOp::commit does for a singleton.

Methods from Deref<Target = DbOp<'static>>§

Source

pub fn maybe_now(&self) -> Option<DateTime<Utc>>

Returns the optionally cached chrono::DateTime

Trait Implementations§

Source§

impl AtomicOperation for StagedOp<'_>

Source§

fn maybe_now(&self) -> Option<DateTime<Utc>>

Function for querying when the operation is taking place - if it is cached.
Source§

fn clock(&self) -> &ClockHandle

Returns the clock handle for time operations. Read more
Source§

fn connection(&mut self) -> &mut Connection

Returns the raw underlying connection. The desired way to represent this would actually be as a GAT: Read more
Source§

fn as_executor(&mut self) -> OneTimeExecutor<'_, &mut Connection>

Returns the sqlx::Executor implementation that statements should be executed through. Read more
Source§

fn add_commit_hook_dyn( &mut self, type_id: TypeId, hook: Box<dyn DynHook>, ) -> Result<(), Box<dyn DynHook>>

Object-safe, type-erased form of add_commit_hook.
Source§

fn commit_hook_dyn(&self, type_id: TypeId) -> Option<&dyn DynHook>

Object-safe, type-erased form of commit_hook.
Source§

fn supports_hooks(&self) -> bool

Whether this operation supports commit hooks. Read more
Source§

fn savepoint_parts(&mut self) -> (&mut Connection, HookSlot<'_>)

Simultaneous access to the connection and the commit-hook buffer a nested SAVEPOINT folds into when released. Implementing this is the only thing an operation must do to get the whole of SavepointOperationwith_savepoint, begin_savepoint, and arbitrary-depth nesting — for free. Read more
Source§

fn add_commit_hook<H>(&mut self, hook: H) -> Result<(), H>
where H: CommitHook, Self: Sized,

Registers a commit hook that will run pre_commit before and post_commit after the transaction commits. Returns Ok(()) if the hook was registered, Err(hook) if hooks are not supported.
Source§

fn commit_hook<H>(&self) -> Option<&H>
where H: CommitHook, Self: Sized,

Typed shared access to the currently-accumulating commit hook of type H, if this operation supports commit hooks and one is registered. Returns the hook a subsequent add_commit_hook::<H> call would merge into.
Source§

impl Deref for StagedOp<'_>

Source§

type Target = DbOp<'static>

The resulting type after dereferencing.
Source§

fn deref(&self) -> &Self::Target

Dereferences the value.

Auto Trait Implementations§

§

impl<'inv> !RefUnwindSafe for StagedOp<'inv>

§

impl<'inv> !Sync for StagedOp<'inv>

§

impl<'inv> !UnwindSafe for StagedOp<'inv>

§

impl<'inv> Freeze for StagedOp<'inv>

§

impl<'inv> Send for StagedOp<'inv>

§

impl<'inv> Unpin for StagedOp<'inv>

§

impl<'inv> UnsafeUnpin for StagedOp<'inv>

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

Source§

fn run_isolated<'a, T, V, E, F>( &'a mut self, items: &'a [T], f: F, ) -> impl Future<Output = Result<Vec<Result<V, E>>, Error>> + 'a
where T: 'a, V: 'a, E: 'a, F: AsyncFnOnce(&mut SavepointOp<'_>, &T) -> Result<V, E> + Clone + Sync + 'a,

Runs f once per item, each inside its own SAVEPOINT, in item order. Read more
Source§

fn run_bisected<'a, T, E, F>( &'a mut self, items: &'a [T], budget: BisectBudget, f: F, ) -> impl Future<Output = Result<BisectOutcomes<E>, Error>> + 'a
where T: 'a, E: Error + 'static, F: AsyncFnOnce(&mut SavepointOp<'_>, &[T]) -> Result<(), E> + Clone + Sync + 'a,

Probes the whole slice at once, bisecting only on failure. Read more
Source§

fn run_bisected_with<'a, T, E, F, P>( &'a mut self, items: &'a [T], budget: BisectBudget, policy: TransientPolicy<P>, f: F, ) -> impl Future<Output = Result<BisectOutcomes<E>, Error>> + 'a
where T: 'a, E: Display + 'a, P: Fn(&E) -> bool + 'a, F: AsyncFnOnce(&mut SavepointOp<'_>, &[T]) -> Result<(), E> + Clone + Sync + 'a,

run_bisected with a caller-supplied notion of which failures are transient. 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> FutureExt for T

Source§

fn with_context(self, otel_cx: Context) -> WithContext<Self>

Attaches the provided Context to this type, returning a WithContext wrapper. Read more
Source§

fn with_current_context(self) -> WithContext<Self>

Attaches the current Context to this type, returning a WithContext wrapper. Read more
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<P, T> Receiver for P
where P: Deref<Target = T> + ?Sized, T: ?Sized,

Source§

type Target = T

🔬This is a nightly-only experimental API. (arbitrary_self_types)
The target type on which the method may be called.
Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T> SavepointOperation for T
where T: AtomicOperation + ?Sized,

Source§

fn with_savepoint<T, E, F>( &mut self, f: F, ) -> impl Future<Output = Result<Result<T, E>, Error>>
where F: AsyncFnOnce(&mut SavepointOp<'_>) -> Result<T, E>,

Runs f inside a SAVEPOINT, keeping its work on Ok and undoing it on Err — see DbOp::with_savepoint for the full contract, including the two layers of Result. Read more
Source§

fn begin_savepoint( &mut self, ) -> impl Future<Output = Result<SavepointOp<'_>, Error>> + Send

Begins a SAVEPOINT scope explicitly — see DbOp::begin_savepoint. Must be finished with release or rollback; dropping it rolls back. Read more
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