pub struct GrpcTaskDelivery;Expand description
Delivers by pushing a WorkerMessage::ActivityTask onto the worker’s
registered stream.
§Only two outcomes are reachable here, and that is not an oversight
mpsc::Sender::send().await fails only when the receiver is gone, i.e.
the worker’s stream is closed — which is exactly
Undeliverable::WorkerUnreachable.
A full channel does not error; it applies backpressure and the send waits.
So this transport never produces
DeliveryFailed.
I record that because my own design note previously claimed a full channel
would map to DeliveryFailed. It would — under try_send. Under the
send().await this path has always used, that outcome does not exist, and
switching to try_send to manufacture it would turn a worker that is merely
busy into a failed delivery. DeliveryFailed exists for the blocking
transport, which really can fail while its worker is alive.
§The intent is asked once, and the backpressure wait is not re-polled
The caller’s intent is checked immediately before the push. If the caller
loses its claim during the backpressure wait inside send().await, this
arm does not notice: there is no poll boundary to re-ask at, and reaching for
one would mean abandoning a send mid-flight.
That is safe, and by two independent mechanisms rather than one:
- The result cannot be recorded. The completion token is one per pass, minted by the caller and revoked once when the pass places nothing, so a task that lands after its pass withdrew carries an authorization the fences no longer honour.
- The external effect cannot be duplicated. The task’s
idempotency_keyis derived from(workflow_id, run_id, activity_id)alone — attempts and execution generations deliberately do not participate (aion-proto/src/worker.rs:143-146) — so a redelivery by whichever pass re-claimed the row carries the same key and is deduplicated at the worker.
The token covers the recording, the idempotency key covers the effect. If either ever stops holding, that is a finding about the fences or the key — not a reason to start polling a send.
Trait Implementations§
Source§impl Clone for GrpcTaskDelivery
impl Clone for GrpcTaskDelivery
Source§fn clone(&self) -> GrpcTaskDelivery
fn clone(&self) -> GrpcTaskDelivery
1.0.0 (const: unstable) · Source§fn clone_from(&mut self, source: &Self)
fn clone_from(&mut self, source: &Self)
source. Read moreimpl Copy for GrpcTaskDelivery
Source§impl Debug for GrpcTaskDelivery
impl Debug for GrpcTaskDelivery
Source§impl Default for GrpcTaskDelivery
impl Default for GrpcTaskDelivery
Source§fn default() -> GrpcTaskDelivery
fn default() -> GrpcTaskDelivery
Source§impl WorkerTaskDelivery for GrpcTaskDelivery
impl WorkerTaskDelivery for GrpcTaskDelivery
Source§fn deliver<'life0, 'life1, 'life2, 'life3, 'life4, 'async_trait>(
&'life0 self,
worker: &'life1 WorkerHandle,
task: &'life2 ProtoActivityTask,
intent: &'life3 SharedDeliveryIntent,
accepted: &'life4 dyn DeliveryAccepted,
) -> Pin<Box<dyn Future<Output = TaskDelivery> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
'life3: 'async_trait,
'life4: 'async_trait,
fn deliver<'life0, 'life1, 'life2, 'life3, 'life4, 'async_trait>(
&'life0 self,
worker: &'life1 WorkerHandle,
task: &'life2 ProtoActivityTask,
intent: &'life3 SharedDeliveryIntent,
accepted: &'life4 dyn DeliveryAccepted,
) -> Pin<Box<dyn Future<Output = TaskDelivery> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
'life3: 'async_trait,
'life4: 'async_trait,
Auto Trait Implementations§
impl Freeze for GrpcTaskDelivery
impl RefUnwindSafe for GrpcTaskDelivery
impl Send for GrpcTaskDelivery
impl Sync for GrpcTaskDelivery
impl Unpin for GrpcTaskDelivery
impl UnsafeUnpin for GrpcTaskDelivery
impl UnwindSafe for GrpcTaskDelivery
Blanket Implementations§
Source§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
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
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> IntoMaybeUndefined<T> for T
impl<T> IntoMaybeUndefined<T> for T
Source§fn into_maybe_undefined(self) -> MaybeUndefined<T>
fn into_maybe_undefined(self) -> MaybeUndefined<T>
Source§impl<T> IntoOption<T> for T
impl<T> IntoOption<T> for T
Source§fn into_option(self) -> Option<T>
fn into_option(self) -> Option<T>
Source§impl<T> IntoRequest<T> for T
impl<T> IntoRequest<T> for T
Source§fn into_request(self) -> Request<T>
fn into_request(self) -> Request<T>
T in a tonic::Request