pub struct Worker { /* private fields */ }Implementations§
Source§impl Worker
impl Worker
Sourcepub fn sticky_cache(self, options: StickyCacheOptions) -> Result<Self>
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.
Sourcepub fn build_id(self, build_id: impl Into<String>) -> Self
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.
pub fn sticky_cache_metrics(&self) -> Result<StickyCacheMetrics>
Source§impl Worker
impl Worker
Sourcepub fn worker_sessions(self, enabled: bool) -> Self
pub fn worker_sessions(self, enabled: bool) -> Self
Opt in to remote activity sessions. Sticky execution stays unsupported.
pub fn max_concurrent_worker_sessions(self, count: usize) -> Self
Sourcepub fn capabilities<I, S>(self, capabilities: I) -> Self
pub fn capabilities<I, S>(self, capabilities: I) -> Self
Additional resource requirements this worker can actually satisfy.
Sourcepub fn worker_session(
&self,
options: WorkerSessionOptions,
) -> Result<WorkerSession>
pub fn worker_session( &self, options: WorkerSessionOptions, ) -> Result<WorkerSession>
Obtain a shared handle for this registered worker and immutable session options.
Source§impl Worker
impl Worker
pub fn new(client: Client, task_queue: impl Into<String>) -> Self
pub fn worker_id(self, worker_id: impl Into<String>) -> Self
pub fn poll_timeout(self, timeout: Duration) -> Self
pub fn heartbeat_interval(self, interval: Duration) -> Self
Sourcepub fn local_activities(self, enabled: bool) -> Self
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.
Sourcepub fn cooperative_cancellation(self, enabled: bool) -> Self
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.
Sourcepub fn retry_policy(self, policy: WorkerRetryPolicy) -> Self
pub fn retry_policy(self, policy: WorkerRetryPolicy) -> Self
Configure bounded retries for task-poll acquisition and worker heartbeats.
Sourcepub fn recover_transient_outages(self, enabled: bool) -> Self
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.
pub fn on_worker_heartbeat<F>(self, observer: F) -> Self
pub fn max_concurrent_workflow_tasks(self, count: usize) -> Self
pub fn max_concurrent_activity_tasks(self, count: usize) -> Self
Sourcepub fn validate_registration(&self) -> Result<(), DuplicateRegistrationError>
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.
Sourcepub fn register_workflow<F, Fut>(
&mut self,
workflow_type: impl Into<String>,
handler: F,
)
pub fn register_workflow<F, Fut>( &mut self, workflow_type: impl Into<String>, handler: F, )
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.
Sourcepub fn register_typed_workflow<I, O, F, Fut>(
&mut self,
workflow_type: impl Into<String>,
handler: F,
)
pub fn register_typed_workflow<I, O, F, Fut>( &mut self, workflow_type: impl Into<String>, handler: F, )
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.
Sourcepub fn register_workflow_avro_value<F, Fut>(
&mut self,
workflow_type: impl Into<String>,
handler: F,
)
pub fn register_workflow_avro_value<F, Fut>( &mut self, workflow_type: impl Into<String>, handler: F, )
Register a workflow on the lossless fixed Avro Value surface.
Sourcepub fn register_replayed_workflow<S, Factory, F, Fut>(
&mut self,
workflow_type: impl Into<String>,
state_factory: Factory,
handler: F,
)
pub fn register_replayed_workflow<S, Factory, F, Fut>( &mut self, workflow_type: impl Into<String>, state_factory: Factory, handler: F, )
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.
Sourcepub fn register_typed_replayed_workflow<I, O, S, Factory, F, Fut>(
&mut self,
workflow_type: impl Into<String>,
state_factory: Factory,
handler: F,
)
pub fn register_typed_replayed_workflow<I, O, S, Factory, F, Fut>( &mut self, workflow_type: impl Into<String>, state_factory: Factory, handler: F, )
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.
Sourcepub fn declare_workflow_signals(
&mut self,
workflow_type: &str,
signal_names: &[&str],
) -> Result<()>
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.
Sourcepub fn set_workflow_definition_sources(
&mut self,
workflow_type: &str,
sources: &[&str],
) -> Result<()>
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.
Sourcepub fn register_replayed_workflow_avro_value<S, Factory, F, Fut>(
&mut self,
workflow_type: impl Into<String>,
state_factory: Factory,
handler: F,
)
pub fn register_replayed_workflow_avro_value<S, Factory, F, Fut>( &mut self, workflow_type: impl Into<String>, state_factory: Factory, handler: F, )
Register a replayable workflow on the lossless fixed Avro Value surface.
pub fn register_activity<F, Fut>( &mut self, activity_type: impl Into<String>, handler: F, )
Sourcepub fn register_typed_activity<I, O, F, Fut>(
&mut self,
activity_type: impl Into<String>,
handler: F,
)
pub fn register_typed_activity<I, O, F, Fut>( &mut self, activity_type: impl Into<String>, handler: F, )
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.
Sourcepub fn register_activity_avro_value<F, Fut>(
&mut self,
activity_type: impl Into<String>,
handler: F,
)
pub fn register_activity_avro_value<F, Fut>( &mut self, activity_type: impl Into<String>, handler: F, )
Register an activity on the lossless fixed Avro Value surface.
Sourcepub fn register_query<F, Fut>(
&mut self,
workflow_type: impl Into<String>,
query_name: impl Into<String>,
handler: F,
)
pub fn register_query<F, Fut>( &mut self, workflow_type: impl Into<String>, query_name: impl Into<String>, handler: F, )
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.
Sourcepub fn register_query_avro_value<F, Fut>(
&mut self,
workflow_type: impl Into<String>,
query_name: impl Into<String>,
handler: F,
)
pub fn register_query_avro_value<F, Fut>( &mut self, workflow_type: impl Into<String>, query_name: impl Into<String>, handler: F, )
Register a query handler on the lossless fixed Avro Value surface.
Sourcepub fn register_replayed_query<S, F, Fut>(
&mut self,
workflow_type: impl Into<String>,
query_name: impl Into<String>,
handler: F,
)
pub fn register_replayed_query<S, F, Fut>( &mut self, workflow_type: impl Into<String>, query_name: impl Into<String>, handler: F, )
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.
Sourcepub fn register_replayed_query_avro_value<S, F, Fut>(
&mut self,
workflow_type: impl Into<String>,
query_name: impl Into<String>,
handler: F,
)
pub fn register_replayed_query_avro_value<S, F, Fut>( &mut self, workflow_type: impl Into<String>, query_name: impl Into<String>, handler: F, )
Register a replayed-state query on the lossless fixed Avro Value surface.
Sourcepub fn register_update<F, Fut>(
&mut self,
workflow_type: impl Into<String>,
update_name: impl Into<String>,
handler: F,
)
pub fn register_update<F, Fut>( &mut self, workflow_type: impl Into<String>, update_name: impl Into<String>, handler: F, )
Register a synchronous workflow update handler.
Sourcepub fn register_update_avro_value<F, Fut>(
&mut self,
workflow_type: impl Into<String>,
update_name: impl Into<String>,
handler: F,
)
pub fn register_update_avro_value<F, Fut>( &mut self, workflow_type: impl Into<String>, update_name: impl Into<String>, handler: F, )
Register an update handler on the lossless fixed Avro Value surface.
pub async fn register(&self) -> Result<RegisterWorkerResponse>
Sourcepub async fn run(&self) -> Result<()>
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.
Sourcepub async fn run_until<F>(&self, shutdown: F) -> Result<()>
pub async fn run_until<F>(&self, shutdown: F) -> Result<()>
Run until shutdown resolves or a terminal worker error occurs.
This has the same liveness and terminal-error contract as Worker::run.
Sourcepub async fn run_once(&self) -> Result<usize>
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.