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:
| exit | meaning | cursor |
|---|---|---|
commit | processing of this event is done | advances past the event |
suspend | this stage is done, the chain continues | unmoved |
pause_until | this stage is done, processing pauses until at | unmoved |
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>
impl<'inv> StagedOp<'inv>
Sourcepub async fn suspend(
self,
) -> Result<Suspended<'inv>, Box<dyn Error + Send + Sync>>
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).
Sourcepub fn pause_until(self, at: DateTime<Utc>) -> Handled<'inv>
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.
Sourcepub fn commit(self) -> Handled<'inv>
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.
Trait Implementations§
Source§impl AtomicOperation for StagedOp<'_>
impl AtomicOperation for StagedOp<'_>
Source§fn maybe_now(&self) -> Option<DateTime<Utc>>
fn maybe_now(&self) -> Option<DateTime<Utc>>
Source§fn clock(&self) -> &ClockHandle
fn clock(&self) -> &ClockHandle
Source§fn connection(&mut self) -> &mut Connection
fn connection(&mut self) -> &mut Connection
Source§fn as_executor(&mut self) -> OneTimeExecutor<'_, &mut Connection>
fn as_executor(&mut self) -> OneTimeExecutor<'_, &mut Connection>
sqlx::Executor implementation that statements should be
executed through. Read moreSource§fn add_commit_hook_dyn(
&mut self,
type_id: TypeId,
hook: Box<dyn DynHook>,
) -> Result<(), Box<dyn DynHook>>
fn add_commit_hook_dyn( &mut self, type_id: TypeId, hook: Box<dyn DynHook>, ) -> Result<(), Box<dyn DynHook>>
add_commit_hook.Source§fn commit_hook_dyn(&self, type_id: TypeId) -> Option<&dyn DynHook>
fn commit_hook_dyn(&self, type_id: TypeId) -> Option<&dyn DynHook>
commit_hook.Source§fn supports_hooks(&self) -> bool
fn supports_hooks(&self) -> bool
Source§fn savepoint_parts(&mut self) -> (&mut Connection, HookSlot<'_>)
fn savepoint_parts(&mut self) -> (&mut Connection, HookSlot<'_>)
SAVEPOINT folds into when released. Implementing this is the
only thing an operation must do to get the whole of
SavepointOperation — with_savepoint, begin_savepoint, and
arbitrary-depth nesting — for free. Read moreSource§fn add_commit_hook<H>(&mut self, hook: H) -> Result<(), H>where
H: CommitHook,
Self: Sized,
fn add_commit_hook<H>(&mut self, hook: H) -> Result<(), H>where
H: CommitHook,
Self: Sized,
Source§fn commit_hook<H>(&self) -> Option<&H>where
H: CommitHook,
Self: Sized,
fn commit_hook<H>(&self) -> Option<&H>where
H: CommitHook,
Self: Sized,
H,
if this operation supports commit hooks and one is registered.
Returns the hook a subsequent add_commit_hook::<H> call would merge into.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> BatchIsolation for Twhere
T: SavepointOperation + ?Sized,
impl<T> BatchIsolation for Twhere
T: SavepointOperation + ?Sized,
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>> + 'awhere
T: 'a,
V: 'a,
E: 'a,
F: AsyncFnOnce(&mut SavepointOp<'_>, &T) -> Result<V, E> + Clone + Sync + 'a,
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>> + 'awhere
T: 'a,
V: 'a,
E: 'a,
F: AsyncFnOnce(&mut SavepointOp<'_>, &T) -> Result<V, E> + Clone + Sync + 'a,
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>> + 'awhere
T: 'a,
E: Error + 'static,
F: AsyncFnOnce(&mut SavepointOp<'_>, &[T]) -> Result<(), E> + Clone + Sync + 'a,
fn run_bisected<'a, T, E, F>(
&'a mut self,
items: &'a [T],
budget: BisectBudget,
f: F,
) -> impl Future<Output = Result<BisectOutcomes<E>, Error>> + 'awhere
T: 'a,
E: Error + 'static,
F: AsyncFnOnce(&mut SavepointOp<'_>, &[T]) -> Result<(), E> + Clone + Sync + 'a,
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
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
run_bisected with a caller-supplied notion of
which failures are transient. Read moreSource§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Source§impl<T> FutureExt for T
impl<T> FutureExt for T
Source§fn with_context(self, otel_cx: Context) -> WithContext<Self> ⓘ
fn with_context(self, otel_cx: Context) -> WithContext<Self> ⓘ
Source§fn with_current_context(self) -> WithContext<Self> ⓘ
fn with_current_context(self) -> WithContext<Self> ⓘ
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
Source§fn in_current_span(self) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
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 moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
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 moreSource§impl<T> SavepointOperation for Twhere
T: AtomicOperation + ?Sized,
impl<T> SavepointOperation for Twhere
T: AtomicOperation + ?Sized,
Source§fn with_savepoint<T, E, F>(
&mut self,
f: F,
) -> impl Future<Output = Result<Result<T, E>, Error>>
fn with_savepoint<T, E, F>( &mut self, f: F, ) -> impl Future<Output = Result<Result<T, E>, Error>>
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 moreSource§fn begin_savepoint(
&mut self,
) -> impl Future<Output = Result<SavepointOp<'_>, Error>> + Send
fn begin_savepoint( &mut self, ) -> impl Future<Output = Result<SavepointOp<'_>, Error>> + Send
SAVEPOINT scope explicitly — see
DbOp::begin_savepoint. Must be finished
with release or
rollback; dropping it rolls back. Read more