Struct Worker

Source
pub struct Worker { /* private fields */ }

Implementations§

Source§

impl Worker

Source

pub fn sticky_cache(self, options: StickyCacheOptions) -> Result<Self>

Explicitly enable a bounded durable-history cache. Disabled by default.

The encoded-byte limit excludes transient JSON decoding, replay, keys and cursor metadata. This never retains live workflow instances or session resources.

Source

pub fn build_id(self, build_id: impl Into<String>) -> Self

Select the deployment build identity used for registration, routing and cache keys.

An omitted build ID uses the SDK identity for cache claims, preserving the Server’s unversioned registration default. An explicit ID must be 1..=255 bytes.

Source

pub fn sticky_cache_metrics(&self) -> Result<StickyCacheMetrics>

Source§

impl Worker

Source

pub fn worker_sessions(self, enabled: bool) -> Self

Opt in to remote activity sessions. Sticky execution stays unsupported.

Source

pub fn max_concurrent_worker_sessions(self, count: usize) -> Self

Source

pub fn capabilities<I, S>(self, capabilities: I) -> Self
where I: IntoIterator<Item = S>, S: Into<String>,

Additional resource requirements this worker can actually satisfy.

Source

pub fn worker_session( &self, options: WorkerSessionOptions, ) -> Result<WorkerSession>

Obtain a shared handle for this registered worker and immutable session options.

Source§

impl Worker

Source

pub fn new(client: Client, task_queue: impl Into<String>) -> Self

Source

pub fn worker_id(self, worker_id: impl Into<String>) -> Self

Source

pub fn poll_timeout(self, timeout: Duration) -> Self

Source

pub fn heartbeat_interval(self, interval: Duration) -> Self

Source

pub fn local_activities(self, enabled: bool) -> Self

Explicitly enable inline local activities in this workflow worker.

The worker advertises the capability to Server only when enabled. These callbacks must yield to Tokio and have idempotent side effects. Cooperative prepared local supervision is a separate capability.

Source

pub fn cooperative_cancellation(self, enabled: bool) -> Self

Opt in to separate cooperative workflow cancellation requests.

Registration requires the Server to acknowledge this exact worker, namespace, queue and capabilities with compatible protocol 1.20. Task processing uses canonical delivery and activity ownership fences. Existing terminal cancel and terminate operations remain terminal. The default is disabled, preserving ordinary protocol 1.19 workers. Call Worker::register before Worker::run_once. Worker::run and Worker::run_until register automatically.

Source

pub fn retry_policy(self, policy: WorkerRetryPolicy) -> Self

Configure bounded retries for task-poll acquisition and worker heartbeats.

Source

pub fn recover_transient_outages(self, enabled: bool) -> Self

Keep run and run_until alive during retryable poll and heartbeat outages.

Retries preserve the poll request identity and use capped exponential backoff from WorkerRetryPolicy. Shutdown interrupts a retry wait. An in-flight poll is allowed to finish so leased work is not discarded. Authentication, protocol, codec, handler and task settlement failures remain errors. Registration and deregistration are not retried here.

The default is disabled, preserving the bounded retry contract. run_once always remains bounded. max_retries = 0 disables retries, even when this option is enabled.

Source

pub fn on_worker_heartbeat<F>(self, observer: F) -> Self
where F: Fn(&WorkerHeartbeatObservation) + Send + Sync + 'static,

Source

pub fn max_concurrent_workflow_tasks(self, count: usize) -> Self

Source

pub fn max_concurrent_activity_tasks(self, count: usize) -> Self

Source

pub fn validate_registration(&self) -> Result<(), DuplicateRegistrationError>

Validate local handler admission without making a Server request.

Registration methods retain their existing signatures. A duplicate keeps the original handler and makes this Worker invalid, including its clones. register, run, run_until and run_once enforce this check before any network request. Build a new Worker to correct an invalid configuration.

Source

pub fn register_workflow<F, Fut>( &mut self, workflow_type: impl Into<String>, handler: F, )
where F: Fn(WorkflowContext, Value) -> Fut + Send + Sync + 'static, Fut: Future<Output = Result<Value>> + Send + 'static,

Register a workflow handler.

An uncaught Error returned by the handler fails the workflow run and is reported to clients as Error::WorkflowFailed. Errors that occur while acquiring or decoding a worker task remain worker-operation failures and do not get converted into workflow outcomes.

Source

pub fn register_typed_workflow<I, O, F, Fut>( &mut self, workflow_type: impl Into<String>, handler: F, )
where I: DeserializeOwned + Send + 'static, O: Serialize + Send + 'static, F: Fn(WorkflowContext, I) -> Fut + Send + Sync + 'static, Fut: Future<Output = Result<O>> + Send + 'static,

Register a workflow with one Serde request value and a Serde result.

This is an ergonomic adapter over the same fixed Avro Value protocol as Worker::register_workflow_avro_value. It does not create or publish a workflow-specific schema. A task must contain zero arguments for a unit request or exactly one argument for every other request type.

See the runnable hello_world example for typed workflow and activity contracts with retry and timeout policy.

Source

pub fn register_workflow_avro_value<F, Fut>( &mut self, workflow_type: impl Into<String>, handler: F, )
where F: Fn(WorkflowContext, AvroValue) -> Fut + Send + Sync + 'static, Fut: Future<Output = Result<AvroValue>> + Send + 'static,

Register a workflow on the lossless fixed Avro Value surface.

Source

pub fn register_replayed_workflow<S, Factory, F, Fut>( &mut self, workflow_type: impl Into<String>, state_factory: Factory, handler: F, )
where S: Clone + Send + Sync + 'static, Factory: Fn() -> S + Send + Sync + 'static, F: Fn(WorkflowContext, Value, WorkflowInstance<S>) -> Fut + Send + Sync + 'static, Fut: Future<Output = Result<Value>> + Send + 'static,

Register a workflow whose typed instance state can be reconstructed for queries.

state_factory creates a fresh instance for every normal workflow task and query replay. The workflow handler is the single source of truth for state transitions: it updates WorkflowInstance after activities and signals resolve. Query replay runs this same handler over committed history and discards any commands it would emit.

Source

pub fn register_typed_replayed_workflow<I, O, S, Factory, F, Fut>( &mut self, workflow_type: impl Into<String>, state_factory: Factory, handler: F, )
where I: DeserializeOwned + Send + 'static, O: Serialize + Send + 'static, S: Clone + Send + Sync + 'static, Factory: Fn() -> S + Send + Sync + 'static, F: Fn(WorkflowContext, I, WorkflowInstance<S>) -> Fut + Send + Sync + 'static, Fut: Future<Output = Result<O>> + Send + 'static,

Register a replayable workflow with one Serde request value and result.

Normal task execution and instance-state query replay both decode and encode through the fixed Avro Value codec. The state factory and handler otherwise follow Worker::register_replayed_workflow.

Source

pub fn declare_workflow_signals( &mut self, workflow_type: &str, signal_names: &[&str], ) -> Result<()>

Declare the signal names consumed by a registered workflow.

Call this after registering the workflow and before starting the worker. Names read through wait_signal or signals need a declaration for Server command admission. Arguments remain arbitrary positional values, including lossless Avro values. This replaces the previous declaration; an empty slice declares no signals. Names are sorted and deduplicated. Existing runs retain the declarations captured when they started.

Source

pub fn set_workflow_definition_sources( &mut self, workflow_type: &str, sources: &[&str], ) -> Result<()>

Bind a registered workflow to compile-time embedded source for safe redrive.

Supply include_str! values for the workflow body and every helper whose behavior can affect replay. Register the workflow once, then set its identity before starting the worker. Duplicate registrations are invalid. The server rejects reusing a worker ID when its prior fingerprint disappears. Workflows without source identity remain runnable but cannot be safely redriven. For example, pass &[include_str!("workflows.rs")] when the handler and replay-sensitive helpers are in a sibling workflows.rs file.

Source

pub fn register_replayed_workflow_avro_value<S, Factory, F, Fut>( &mut self, workflow_type: impl Into<String>, state_factory: Factory, handler: F, )
where S: Clone + Send + Sync + 'static, Factory: Fn() -> S + Send + Sync + 'static, F: Fn(WorkflowContext, AvroValue, WorkflowInstance<S>) -> Fut + Send + Sync + 'static, Fut: Future<Output = Result<AvroValue>> + Send + 'static,

Register a replayable workflow on the lossless fixed Avro Value surface.

Source

pub fn register_activity<F, Fut>( &mut self, activity_type: impl Into<String>, handler: F, )
where F: Fn(ActivityContext, Value) -> Fut + Send + Sync + 'static, Fut: Future<Output = Result<Value>> + Send + 'static,

Source

pub fn register_typed_activity<I, O, F, Fut>( &mut self, activity_type: impl Into<String>, handler: F, )
where I: DeserializeOwned + Send + 'static, O: Serialize + Send + 'static, F: Fn(ActivityContext, I) -> Fut + Send + Sync + 'static, Fut: Future<Output = Result<O>> + Send + 'static,

Register an activity with one Serde request value and a Serde result.

Inputs and results use the platform’s fixed Avro Value schema. Shape mismatches and unsupported Serde values return Error::HandlerType with the activity name and Rust type.

Source

pub fn register_activity_avro_value<F, Fut>( &mut self, activity_type: impl Into<String>, handler: F, )
where F: Fn(ActivityContext, AvroValue) -> Fut + Send + Sync + 'static, Fut: Future<Output = Result<AvroValue>> + Send + 'static,

Register an activity on the lossless fixed Avro Value surface.

Source

pub fn register_query<F, Fut>( &mut self, workflow_type: impl Into<String>, query_name: impl Into<String>, handler: F, )
where F: Fn(QueryContext, Value) -> Fut + Send + Sync + 'static, Fut: Future<Output = Result<Value>> + Send + 'static,

Register a named, read-only query handler for a workflow type.

The workflow type must also be registered with Worker::register_workflow before the worker runs. The handler receives only an immutable committed state snapshot and normalized query arguments.

Source

pub fn register_query_avro_value<F, Fut>( &mut self, workflow_type: impl Into<String>, query_name: impl Into<String>, handler: F, )
where F: Fn(QueryContext, AvroValue) -> Fut + Send + Sync + 'static, Fut: Future<Output = Result<AvroValue>> + Send + 'static,

Register a query handler on the lossless fixed Avro Value surface.

Source

pub fn register_replayed_query<S, F, Fut>( &mut self, workflow_type: impl Into<String>, query_name: impl Into<String>, handler: F, )
where S: Clone + Send + Sync + 'static, F: Fn(QueryContext, Arc<S>, Value) -> Fut + Send + Sync + 'static, Fut: Future<Output = Result<Value>> + Send + 'static,

Register a named query against deterministically replayed instance state.

The workflow type must use Worker::register_replayed_workflow with the same state type S. The handler receives an immutable, detached state clone, so successful and failed queries cannot affect workflow execution or the state reconstructed by a later query.

Source

pub fn register_replayed_query_avro_value<S, F, Fut>( &mut self, workflow_type: impl Into<String>, query_name: impl Into<String>, handler: F, )
where S: Clone + Send + Sync + 'static, F: Fn(QueryContext, Arc<S>, AvroValue) -> Fut + Send + Sync + 'static, Fut: Future<Output = Result<AvroValue>> + Send + 'static,

Register a replayed-state query on the lossless fixed Avro Value surface.

Source

pub fn register_update<F, Fut>( &mut self, workflow_type: impl Into<String>, update_name: impl Into<String>, handler: F, )
where F: Fn(QueryContext, Value) -> Fut + Send + Sync + 'static, Fut: Future<Output = Result<Value>> + Send + 'static,

Register a synchronous workflow update handler.

Source

pub fn register_update_avro_value<F, Fut>( &mut self, workflow_type: impl Into<String>, update_name: impl Into<String>, handler: F, )
where F: Fn(QueryContext, AvroValue) -> Fut + Send + Sync + 'static, Fut: Future<Output = Result<AvroValue>> + Send + 'static,

Register an update handler on the lossless fixed Avro Value surface.

Source

pub async fn register(&self) -> Result<RegisterWorkerResponse>

Source

pub async fn run(&self) -> Result<()>

Run until shutdown or a terminal worker error occurs.

Empty long-poll expirations do not stop the worker. Retryable poll and heartbeat failures use WorkerRetryPolicy independently, while authentication, protocol, and other non-retryable failures are returned. Enable Worker::recover_transient_outages to keep retrying transient poll and heartbeat outages beyond the ordinary retry budget.

Source

pub async fn run_until<F>(&self, shutdown: F) -> Result<()>
where F: Future<Output = ()>,

Run until shutdown resolves or a terminal worker error occurs.

This has the same liveness and terminal-error contract as Worker::run.

Source

pub async fn run_once(&self) -> Result<usize>

Poll and settle at most one task from each enabled task family.

A workflow may reach its server-enforced run deadline while this worker holds a task. When the completion endpoint authoritatively rejects that selected task and run with recorded=false, reason=run_timed_out, and terminal run_status=failed, the workflow tick is considered settled: the late command was not recorded and cannot replace the terminal run. Every other completion rejection remains an error. This worker-level race handling is distinct from WorkflowResultOptions::timeout, which only bounds how long a client waits for a result.

Direct callers of Client::complete_workflow_task continue to receive the original Error::Http status and response body.

Trait Implementations§

Source§

impl Clone for Worker

Source§

fn clone(&self) -> Worker

Returns a copy of the value. Read more
1.0.0 · Source§

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

Performs copy-assignment from source. Read more

Auto Trait Implementations§

§

impl Freeze for Worker

§

impl !RefUnwindSafe for Worker

§

impl Send for Worker

§

impl Sync for Worker

§

impl Unpin for Worker

§

impl !UnwindSafe for Worker

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

Source§

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

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

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

§

impl<T> Instrument for T

§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided [Span], returning an Instrumented wrapper. Read more
§

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.

§

impl<T> PolicyExt for T
where T: ?Sized,

§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns [Action::Follow] only if self and other return Action::Follow. Read more
§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns [Action::Follow] if either self or other returns Action::Follow. Read more
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 = Infallible

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.
§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

§

fn vzip(self) -> V

§

impl<T> WithSubscriber for T

§

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
§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a [WithDispatch] wrapper. Read more
§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,