Skip to main content

WorkerHandle

Struct WorkerHandle 

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

Cloneable handle for a registered worker stream.

A worker serves a SET of namespaces under a single task_queue, so it is indexed under one (namespace, task_queue, activity_type) key per namespace in its set. node is an OPTIONAL locality affinity (a locality, not a process — many handles may share a node id) used as a within-pool filter at selection time; None means the worker advertised no locality.

Implementations§

Source§

impl WorkerHandle

Source

pub const fn id(&self) -> WorkerId

Worker identifier assigned by this server process.

Source

pub const fn namespaces(&self) -> &BTreeSet<String>

Namespaces authorized for this worker stream. The worker is reachable for a dispatch only when its set includes the workflow’s namespace.

Source

pub fn identity(&self) -> &str

The identity string this worker registered under, exactly as the registration frame carried it (empty when the worker sent none).

Source

pub fn task_queue(&self) -> &str

Task queue (pool/flavour) this worker serves within each namespace.

Source

pub fn node(&self) -> Option<&str>

Optional locality affinity this worker advertised. None means the worker carries no node and is reachable only for unpinned dispatches.

Source

pub fn activity_types(&self) -> &BTreeSet<String>

Activity types advertised by this worker.

Source

pub const fn instance(&self) -> Option<&WorkerInstanceIdentity>

Optional durable deployment/instance association supplied at registration.

Source

pub const fn max_concurrency(&self) -> Option<NonZeroU32>

How many activities this worker runs AT ONCE, as it advertised, or None while it has registered but not yet said.

A dispatch PRECONDITION, read by [eligible_candidates_in_rotation]: a worker holding this many in-flight activities is passed over rather than pushed past its own admission. Pushing past it is what produced a worker that stopped reading its stream, answered no pings, and was then deregistered as lost while it was in fact working.

None is reachable only on the liminal transport, whose published registration frame has no capacity field, and only for the one round trip between that registration and the worker’s capacity announcement. It means UNKNOWN, never unlimited: selection treats an unknown capacity as full, because the alternative is to invent a number for a process that is about to state its own.

Source

pub fn is_connected(&self) -> bool

Whether the server still holds an OPEN push channel to this worker.

The explicit connected fact, as opposed to the absence of noise. Silence alone cannot tell a worker that is gone from one that is busy, and the expiry sweep used to read the second as the first; this is the third input that separates them.

gRPC answers from the stream sender itself — a closed receiver means the tonic task that owned the stream is gone. Liminal answers from its connection: the transport holds a live push connection or it does not.

Source

pub const fn delivery(&self) -> &WorkerDelivery

The transport this worker is delivered to through.

Source

pub const fn intervention_capabilities(&self) -> &InterventionCapabilities

The neutral intervention primitives this worker’s harness advertises (NOI-6).

The intervention router gates on this set and never routes an unadvertised primitive. Empty (the default for a plain activity worker) means the worker is observability-only: the router refuses every intervention command for it.

Source

pub fn sender(&self) -> Option<&WorkerTaskSender>

gRPC stream sender used by the gRPC dispatch path to push work, or None when this worker is delivered to over a non-gRPC transport (liminal).

The gRPC dispatch path registers every worker with a WorkerDelivery::Grpc delivery, so this is always Some for a gRPC-registered worker — the behaviour the path relied on before delivery became transport-agnostic.

Trait Implementations§

Source§

impl Clone for WorkerHandle

Source§

fn clone(&self) -> WorkerHandle

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 WorkerHandle

Source§

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

Formats the value using the given formatter. 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, !>

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