Skip to main content

CompletionFences

Struct CompletionFences 

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

Process-incarnation registry of the only generations allowed to complete.

Recovery deliberately starts with an empty registry. Re-dispatch issues a fresh token, so a result carrying any pre-recovery token is refused whether it arrives before or after the new dispatch.

Implementations§

Source§

impl CompletionFences

Source

pub fn issue( &self, workflow_id: &WorkflowId, run_id: &RunId, activity_id: &ActivityId, attempt: u32, ) -> Result<CompletionToken, ServerError>

Authorize one more delivery of one attempt and return its token.

One generation is held per execution site, and what a call is decides what happens to it:

  • a different idempotency_key is a different EXECUTION GENERATION (reset / continue-as-new): the prior generation is superseded whole;
  • a higher attempt within the same generation is a genuine RETRY: the prior attempt’s tokens are superseded;
  • the SAME attempt within the same generation is a REDELIVERY of work a worker may still be executing, so the new token JOINS the outstanding set rather than replacing it. Both workers hold a token that can be accepted, and the first accepted completion consumes both.

A re-dispatch of an attempt that has already been superseded registers nothing: it is logged with both attempt numbers, and the token it returns is refused as a stale generation when presented, which is the truthful answer for a worker whose attempt no longer exists.

§Errors

Returns ServerError::LockPoisoned when fence state cannot be trusted.

Source

pub fn accept( &self, workflow_id: &WorkflowId, activity_id: &ActivityId, submitted: &CompletionToken, ) -> Result<AcceptedGeneration, ServerError>

Consume the current generation when submitted is one of the tokens it still has outstanding.

Comparison and consumption share one mutex critical section, so two concurrent submissions cannot both become truth: the whole generation is removed by the first, and every later presentation of any of its tokens finds no generation at all.

§Errors

Returns a typed rejection for a missing generation or a token that is not outstanding in it, or ServerError::LockPoisoned when fence state cannot be trusted.

Source

pub fn restore_if_absent( &self, workflow_id: &WorkflowId, activity_id: &ActivityId, accepted: &AcceptedGeneration, ) -> Result<(), ServerError>

Restore a consumed generation after accepted-path settlement fails.

A concurrently issued newer generation always wins; the accepted generation is restored only while the execution site has no current one.

§Errors

Returns ServerError::LockPoisoned when fence state cannot be trusted.

Source

pub fn revoke( &self, workflow_id: &WorkflowId, activity_id: &ActivityId, token: &CompletionToken, ) -> Result<(), ServerError>

Revoke exactly the token this caller issued, and nothing else.

A dispatch that could not place its task withdraws its OWN authorization. It never withdraws a sibling token still held by a live worker executing the same attempt, and it never disturbs a newer retry — a newer attempt already replaced the outstanding set, so the old token is simply absent and the removal is a no-op. The generation is dropped once its last outstanding token is gone.

§Errors

Returns ServerError::LockPoisoned when fence state cannot be trusted.

Source

pub fn revoke_current( &self, workflow_id: &WorkflowId, activity_id: &ActivityId, ) -> Result<(), ServerError>

Revoke whichever generation is current while parking for recovery.

A park retires the execution site deliberately, so it takes every token of the current generation with it — including a redelivery’s sibling.

§Errors

Returns ServerError::LockPoisoned when fence state cannot be trusted.

Source

pub fn arm_lease(&self, token: &CompletionToken) -> Result<(), ServerError>

Mark token’s lease append as in flight: a completion for it waits at Self::lease_settled until Self::settle_lease runs.

§Errors

Returns ServerError::LockPoisoned if the gate lock is poisoned.

Source

pub fn settle_lease(&self, token: &CompletionToken) -> Result<(), ServerError>

The lease append for token has landed, failed, or been abandoned — release every completion waiting on it.

§Errors

Returns ServerError::LockPoisoned if the gate lock is poisoned.

Source

pub fn lease_pending( &self, token: &CompletionToken, ) -> Result<bool, ServerError>

Whether token’s lease append is still in flight.

§Errors

Returns ServerError::LockPoisoned if the gate lock is poisoned.

Source

pub async fn lease_settled( &self, token: &CompletionToken, ) -> Result<(), ServerError>

Wait until token’s lease append is no longer in flight. Returns at once for a token that was never armed or is already settled.

The completion entry points await this before handing a result to the fence; Self::accept also waits when it can host the wait itself.

§Errors

Returns ServerError::LockPoisoned if the gate lock is poisoned.

Trait Implementations§

Source§

impl Clone for CompletionFences

Source§

fn clone(&self) -> CompletionFences

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more
Source§

impl Debug for CompletionFences

Source§

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

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

impl Default for CompletionFences

Source§

fn default() -> CompletionFences

Returns the “default value” for a type. Read more

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> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> DynClone for T
where T: Clone,

Source§

fn __clone_box(&self, _: Private) -> *mut ()

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> FromRef<T> for T
where T: Clone,

Source§

fn from_ref(input: &T) -> T

Converts to this type from a reference to the input type.
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> IntoMaybeUndefined<T> for T

Source§

fn into_maybe_undefined(self) -> MaybeUndefined<T>

Converts this value into a three-state builder argument.
Source§

impl<T> IntoOption<T> for T

Source§

fn into_option(self) -> Option<T>

Converts this value into an optional builder argument.
Source§

impl<T> IntoRequest<T> for T

Source§

fn into_request(self) -> Request<T>

Wrap the input message T in a tonic::Request
Source§

impl<L> LayerExt<L> for L

Source§

fn named_layer<S>(&self, service: S) -> Layered<<L as Layer<S>>::Service, S>
where L: Layer<S>,

Applies the layer to a service and wraps it in Layered.
Source§

impl<T> Pointable for T

Source§

const ALIGN: usize

The alignment of pointer.
Source§

type Init = T

The type for initializers.
Source§

unsafe fn init(init: <T as Pointable>::Init) -> usize

Initializes a with the given initializer. Read more
Source§

unsafe fn deref<'a>(ptr: usize) -> &'a T

Dereferences the given pointer. Read more
Source§

unsafe fn deref_mut<'a>(ptr: usize) -> &'a mut T

Mutably dereferences the given pointer. Read more
Source§

unsafe fn drop(ptr: usize)

Drops the object pointed to by the given pointer. 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> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. 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, <T as TryFrom<U>>::Error>

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