Expand description
§Durable Workflow Rust SDK
durable-workflow is the first-party Rust client and worker SDK for
Durable Workflow. Rust applications can start
and inspect workflows, run workflow and activity handlers, and exchange typed
Avro values with PHP, Python, and Rust workers through one durable runtime.
Use the SDK with a self-hosted Durable Workflow Server or a managed Durable Workflow Cloud namespace. Your workers remain ordinary Rust processes and scale independently from the runtime.
§Install
Rust 1.86 or newer is required.
cargo add durable-workflowApplications using Tokio entry points should also enable the Tokio features they need:
cargo add tokio --features macros,rt-multi-thread§Run the example
examples/hello_world.rs is a complete typed
workflow and activity example. It starts a workflow, runs two activities,
demonstrates retry and failure handling, waits for completion, and prints the
result.
Start a bootstrapped self-hosted Server, then run the example from this repository:
DURABLE_WORKFLOW_SERVER_URL=http://127.0.0.1:8080 \
DURABLE_WORKFLOW_TOKEN=dev-token \
cargo run --example hello_worldPass the Server origin without a trailing /api. TASK_QUEUE changes the
default rust-workers queue, and GREETING_NAME changes the example input.
For a provisioned Cloud namespace, use the exact runtime URL and the separate client and worker credentials shown in Cloud:
DURABLE_WORKFLOW_RUNTIME_URL=https://cloud.durable-workflow.com/api/runtime/v1/namespaces/your-namespace-id \
DURABLE_WORKFLOW_RUNTIME_NAMESPACE=your.namespace \
DURABLE_WORKFLOW_CLIENT_TOKEN=your-client-token \
DURABLE_WORKFLOW_WORKER_TOKEN=your-worker-token \
cargo run --example hello_worldThe namespace runtime URL is already complete. Do not append another /api.
§Core API
Clientstarts, signals, queries, updates, cancels, terminates, describes, and awaits workflow executions.Workerregisters workflow, activity, query and update handlers, declares workflow signals, and long-polls task queues.WorkflowContextprovides durable activities, timers, conditions, child workflows, side effects, version markers, parallel operations, selection, sagas, message streams, memo, search attributes, and continue-as-new.- Typed registration and result helpers preserve Serde request and result types over the fixed Avro Value protocol.
- Activity options cover retries, start-to-close, schedule-to-start, schedule-to-close, heartbeat timeouts, cancellation, and heartbeats.
§Handler registration
Register each workflow and activity wire name once per Worker. Named query and update handlers are also scoped to their workflow type. Dynamic, typed, Avro and replayed adapters share these namespaces. Repeating a registration, including the same callback, is an error. Different Workers or handler kinds may reuse a name, and queries or updates on different workflow types may reuse a name.
Registration methods keep their existing signatures. A duplicate preserves the
first handler and invalidates that Worker. Call worker.validate_registration()?
after configuring it for a local check. register, run, run_until and
run_once also reject it before contacting Server with Error::WorkerLoop
and a duplicate_registration diagnostic. Local validation returns the typed
DuplicateRegistrationError, exposing the handler kind, wire name, optional
workflow scope, and both source locations and authoring adapters. It converts
into the existing worker-error variant when used with ?, preserving exhaustive
matches on the SDK’s existing Error enum.
Construct a new Worker to correct an invalid configuration. Choose the desired
handler before registration instead of relying on replacement order.
After registering a workflow, declare every signal name it reads through
wait_signal or signals before starting the worker:
worker.declare_workflow_signals("orders", &["finish", "changed"])?;Server records those names and positional argument contracts when starting a run. Changing worker declarations does not change existing runs.
The SDK writes Avro payloads only. The fixed recursive Value schema preserves nulls, booleans, signed 64-bit integers, finite doubles, bytes, UTF-8 strings, lists, and string-keyed maps across official SDKs without customer-managed schemas or a registry.
Large payloads are uploaded automatically when the namespace advertises runtime
storage. The SDK follows its inline threshold and upload limit, including batches
that exceed the ordinary request limit. Uploads and downloads use the same
runtime URL, namespace, and credential role, with size and SHA-256 verification.
Client::builder(...).max_external_payload_bytes(...) limits unique downloaded
bytes per response (64 MiB by default). Provider credentials are not needed;
payload fetches never follow redirects.
§Examples
§Local activities
Enable inline local execution with Worker::new(client, queue).local_activities(true).
Register the callback with the ordinary activity registration methods, then call
ctx.local_activity(...), local_activity_typed(...) or
local_activity_avro_value(...) from workflow code. LocalActivityOptions sets
bounded retries and start-to-close, schedule-to-close and heartbeat timeouts.
Local activities run in the workflow worker and bypass the activity queue. Server records their attempts, heartbeat details and terminal result. Committed results replay without running the callback. A worker lost before completion is acknowledged can execute the callback again, so external effects must be idempotent.
Callbacks must yield to Tokio. Timeout or lost workflow lease drops the async callback without requiring application heartbeats. Blocking work needs separate process supervision. Inline local execution and cooperative prepared local supervision use separate worker profiles. Registration rejects combining them. Database-generated execution and failure IDs become available after Server commits history. Use the failure kind, timeout kind and attempt number to branch during fresh execution, rather than testing whether an ID has been assigned.
Available since SDK 3.1.0. Default workers keep local execution disabled. Opted-in workers negotiate their local capability during registration.
§Worker sessions
Enable Worker::worker_sessions(true) and declare the resource requirements that
the worker actually satisfies with capabilities(...). Route a remote call with
ctx.activity(...).in_worker_session(options) and use .typed::<Output>() or
.avro_value() for its result. A parallel group’s in_worker_session(...) applies
the same routing to each activity leaf.
Server creates a session on the first admitted activity. After registration,
worker.worker_session(options) also provides an explicit shared handle for
create(), renew() and close(reason). ActivityContext::worker_session() exposes
the current handle. Heartbeats renew the holder lease without extending the
absolute TTL. Graceful worker shutdown drains activities and closes held sessions
before deregistering.
lease_seconds(...) bounds holder authority and ttl_seconds(...) bounds the
session’s total lifetime. Reacquisition retains the original TTL deadline. Renew
the handle explicitly while keeping an idle resource alive. Activity heartbeats
also renew the lease. An async callback is dropped when its locally observed lease
or TTL ends, even without application heartbeats. Callbacks must yield to Tokio.
Set max_concurrent_worker_sessions(...) to bound the worker’s session registry
and max_concurrent_activities(...) on the options to bound each session’s activity
concurrency. Uncreated or failed handles release their local slot when dropped.
An admitted session keeps its slot while its holder lease is active.
Session memory is process-local. A replacement holder must rebuild its resources, and an interrupted activity can execute again. Use idempotency and attempt fencing for external side effects. Committed results replay from history without rebuilding resources or rerunning callbacks. Changing recorded session options fails replay. Local activities cannot use session routing.
The runnable session example uses a real process-local cache and prints its resource generation. Use Server 2.5.1 / Native 2.4.1 for session history and original TTL preservation. Session support starts with SDK 3.2.0.
§Bounded sticky execution
Enable a durable-history cache explicitly on a Worker:
use durable_workflow::{Client, Result, StickyCacheOptions, Worker};
use std::time::Duration;
fn configure(client: Client) -> Result<Worker> {
Worker::new(client, "rust-workers")
.build_id("orders-v1")
.sticky_cache(StickyCacheOptions::new(32)
.max_history_bytes(8 * 1024 * 1024)
.ttl(Duration::from_secs(60)))
}Caching is disabled by default. A positive capacity opts in, and zero disables it. The defaults are 16 MiB of retained encoded history and a 300-second TTL. Set whole-second TTLs between 1 and 3600 seconds. Registration must confirm the worker, queue, namespace, build and sticky capability before it can poll with caching. Use Server 2.5.10 or newer. SDK support starts with 3.4.0.
The cache retains immutable wire histories, never live workflow instances or session resources. Replays decode a fresh snapshot. Entry and encoded-byte bounds use LRU eviction, and reads do not extend expiry. Decoding, replay, keys and cursor metadata use additional memory. Size the cache within your worker’s memory budget.
A warm replay validates the inline prefix and fetches the retained tail cursor under the current task’s lease and attempt. Missing, expired, changed or invalid entries fall back to complete durable history. Worker replacement needs no cache transfer. Committed side effects replay, and cancellation delivery remains governed by canonical Server history. Terminal completion and worker shutdown clear entries.
worker.sticky_cache_metrics() reports hits, misses, evictions, forced cold
replays, retained entries and encoded history bytes. This reduces history fetching
for warm runs. Measure CPU, memory and completion rate with your own workloads.
An explicit build_id(...) pins new runs to that deployment build. A different
build cannot take an existing pinned run. Omit it to retain unversioned routing.
The runnable example below prints its result and cache counters.
§Deploying workflow patches
Use a stable change ID with ctx.patched(...) or ctx.get_version(...) to
preserve deployment decisions through replay. An old unmarked history selects
the legacy version -1 if the next durable operation was already recorded.
New executions record one decision per change ID, even when code asks again.
Some older workers recorded a marker for every repeated call. Rust replays
consistent markers at their original positions, retaining the recorded history
and timestamps. Keep those call boundaries, or replace them with
ctx.deprecate_patch(...), while such histories remain active. Replay consumes
one matching historical marker per call and never skips an intervening activity,
timer or side effect. Conflicting recorded versions and unsupported current
ranges remain replay errors.
§Recovering from Server outages
For a long-running service worker, enable .recover_transient_outages(true)
on Worker. Retryable poll and worker-heartbeat failures then keep retrying
with capped exponential backoff. A retried poll keeps its original request ID.
run_until interrupts retry waits when shutdown arrives and settles any poll
response already in flight.
The default preserves bounded retries. run_once remains bounded in either
mode, and WorkerRetryPolicy.max_retries = 0 disables ordinary retries.
Authentication, protocol, codec, handler and task settlement failures still
return an error. Registration and deregistration are outside this recovery
option. Keep a process supervisor for startup failures and process crashes.
§Runnable examples
| Example | Demonstrates |
|---|---|
hello_world.rs | Typed worker, workflow, activities, retries, and completion |
activity_options.rs | Activity retry and timeout policies |
local_activities.rs | Explicit local execution, retries, durable timers and remote work |
worker_sessions.rs | Typed activity routing, holder-local resources and graceful session close |
sticky_execution.rs | Opt-in bounded history, signal replay and cache counters |
condition_search_attributes.rs | Durable conditions and typed search attributes |
continue_as_new.rs | Bounded histories and continue-as-new |
parallel_saga.rs | Deterministic parallel work and saga compensation |
§Documentation
- Rust SDK landing page
- Generated API reference
- Rust SDK guide
- Self-hosted Server guide
- Cloud early access
§Compatibility
The SDK supports cooperative requests and bounded,
replayable cleanup. Enable Worker::cooperative_cancellation(true) against a
Server that advertises protocol 1.20 and the required capabilities. This profile
supervises async Activity futures independently of application heartbeats.
Independently cancellable scopes remain disabled. See the
cancellation guide
and v3 migration guide.
The crate publishes its supported Server and worker-protocol ranges in
[package.metadata.durable-workflow] in Cargo.toml. Runtime
capability manifests, not matching package version strings, determine protocol
compatibility. Stable releases follow semantic versioning.
§Development
cargo fmt --all --check
cargo test --all-targets --all-features
cargo doc --all-features --no-deps
cargo packageReplay and codec defects require a minimal regression fixture. See
CONTRIBUTING.md for the corpus rules.
§License
Durable Workflow Rust SDK is released under the MIT License.
Macros§
- json
- Construct a
serde_json::Valuefrom a JSON literal. - wait_
condition - Create a durable condition wait whose predicate definition is fingerprinted from its inline Rust tokens.
Structs§
- Activity
Call - Activity
Context - Activity
Failure - A stable, machine-readable terminal activity failure.
- Activity
Heartbeat Response - Activity
Options - Options recorded atomically on one deterministic
schedule_activitycommand. - Activity
Options Error - A machine-readable activity-options validation failure.
- Activity
Retry Policy - Durable server-side retry policy for one activity execution.
- Activity
Task - Activity
Task Rejection - A worker-side activity settlement or heartbeat rejected by durable state.
- Cancel
Durable Operation Call - Future returned by
DurableOperationHandle::cancel. - Cancellation
Context - Original cancellation metadata restored from canonical request history.
- Cancellation
Delivery - One committed delivery marker, bound to the original cancellation request.
- Cancellation
Delivery Receipt - A delivery acknowledgment bound to one workflow task and durable run.
- Cancellation
History - Validated canonical cancellation state for one workflow run.
- Cancellation
Lineage - One immutable local request in the ordered cancellation lineage.
- Cancellation
Request - The Server’s original request identity and immutable cleanup deadline.
- Cancellation
Scope Open Receipt - A scope opening proved against complete history on its original claim. This proof does not advertise scoped workflow execution support.
- Cancellation
Shield - A workflow-local cleanup scope. Dropping it restores cancellation checks.
- Child
Workflow Avro Result - Lossless successful child result for fixed Avro Value workflows.
- Child
Workflow Call - Future returned by
WorkflowContext::start_child_workflow. - Child
Workflow Failure - A stable, machine-readable child workflow failure delivered to its parent.
- Child
Workflow Options - Options recorded with a child-workflow command.
- Child
Workflow Result - A successful child result together with its durable parent-child identity.
- Child
Workflow Retry Policy - Durable retry policy for one child workflow invocation.
- Client
- Client
Builder - Condition
Wait Call - Future returned by
WorkflowContext::wait_condition. - Condition
Wait Options - Stable identity and optional durable timeout for a condition wait.
- Continue
AsNew Options - Optional routing overrides for a continue-as-new transition.
- Continue
AsNew Options Error - A stable validation error raised before a continue-as-new command is emitted.
- Cooperative
Cancellation Options - Optional reason and runtime-owned cleanup limit for a cooperative request.
- Cooperative
Cancellation Requested - Cancellation delivered at its committed authored call, with original identity.
- Cooperative
Workflow Task - One actual protocol 1.20 claim and its original pending observation.
- Cooperative
Workflow Task Poll - An explicit cooperative poll, including ordinary idle and stop outcomes.
- Duplicate
Registration Error - Ambiguous local configuration, rejected before contacting Server.
- Durable
Operation Await Call - Future returned by
DurableOperationHandle::await_result. - Durable
Operation Cancelled - Typed result of explicitly awaiting a cancelled non-winning operation.
- Durable
Operation Handle - Stable reference to one member of a durable selection group.
- Handler
Registration - The authoring adapter and source location of a handler registration.
- History
Event - Local
Activity Options - Retry and timeout settings for an activity executed by the workflow worker.
- Message
Stream - Message
Stream Message - Parallel
Call - Future returned by
WorkflowContext::parallel. - Parallel
Completion - One successful leaf retained when another parallel member failed.
- Parallel
Failure - A deterministic join failed after some siblings had already completed.
- Parallel
Group Error - Stable validation error returned before an invalid group emits commands.
- Parallel
Group Metadata - Stable identity for one enclosing deterministic parallel group.
- Payload
Envelope - Poll
Activity Task Response - Poll
Query Task Response - Poll
Workflow Task Response - Protocol
Failure - A stable failure returned when a server rejects an SDK protocol version.
- Query
Context - Immutable state supplied to a registered query handler.
- Query
Failure - A stable, machine-readable workflow query or query-task settlement failure.
- Query
Signal - One decoded signal in the committed workflow-history snapshot.
- Query
Task - An ephemeral server-routed query task.
- Register
Worker Response - Replay
Failure - A stable, machine-readable failure raised when workflow code no longer reconstructs the durable command stream recorded in history.
- Saga
- Workflow-local deterministic saga compensation helper.
- Saga
Compensation Failure - A forward saga failure followed by a terminal compensation failure.
- Scoped
Cancellation Context - Original scope ancestry carried by a candidate cooperative child request. Reading this immutable metadata does not authorize entering a scope body.
- Scoped
Cancellation Lineage - One immutable scope address and its bounded cleanup deadline.
- Search
Attribute Update - Validated typed workflow-side search-attribute mutation.
- Select
Call - Future returned by
WorkflowContext::select. - Selection
Result - The one winner committed for a durable selection group.
- Signal
Call - Sticky
Cache Metrics - Replay counters and retained encoded history. Decoding/replay memory is additional.
- Sticky
Cache Options - Opt-in bounds for a worker’s durable-history cache.
- Timer
Call - Future returned by
WorkflowContext::sleep. - Uuid
- A Universally Unique Identifier (UUID).
- Worker
- Worker
Deregistration Envelope - Result of gracefully removing a worker-plane registration.
- Worker
Heartbeat Observation - Worker
Retry Policy - Bounded retry policy for worker poll acquisition and worker heartbeats.
- Worker
Session - One worker’s acknowledged session lease.
- Worker
Session Options - Routing and lifetime options for one worker-held session.
- Workflow
Cancellation Request - Acknowledgment of a request targeting the current durable run.
- Workflow
Cancellation Requested - Cooperative workflow cancellation observed at an author-controlled point.
- Workflow
Command Options - Optional structured fields for a cancellation or termination request.
- Workflow
Command Rejection - A stable rejection returned by instance- or selected-run lifecycle commands.
- Workflow
Command Result - The accepted, machine-readable result of a lifecycle command.
- Workflow
Context - Workflow
Description - Workflow
Handle - Workflow
History Budget - Public history-budget information attached to the current workflow task.
- Workflow
Identity - The identity of one durable workflow execution.
- Workflow
Instance - Typed local state owned by one deterministic workflow invocation.
- Workflow
Redrive Result - A new run continuing a failed run from its recorded activity boundary.
- Workflow
Result Options - Workflow
Start Options - Server-enforced timeout policy for a workflow start.
- Workflow
Stream Append Item - One item for direct or replay-safe append.
- Workflow
Stream Append Result - Durable acceptance and deduplication outcome for an append request.
- Workflow
Stream Description - Lifecycle and backlog metadata for one run-scoped Workflow Stream.
- Workflow
Stream Item - One durable item at its stable zero-based offset.
- Workflow
Stream Page - One bounded at-least-once subscription page.
- Workflow
Task - Workflow
Task Heartbeat - Successful renewal of the exact workflow-task claim with optional observation.
- Workflow
Terminal Outcome - A typed terminal workflow outcome with durable identity and failure metadata.
Enums§
- Activity
Backoff - Backoff intervals for one durable activity retry policy.
- Activity
Failure Kind - Stable terminal categories returned when an awaited activity does not succeed.
- Activity
Options Error Kind - Stable validation categories for
ActivityOptions. - Avro
Value - Native adapter for the fixed language-neutral Avro Value schema.
- Cancellation
Call Kind - The authored durable call where canonical cancellation is delivered.
- Cancellation
Delivery Reply - Delivery either commits or parks the parent until canonical child cleanup.
- Cancellation
Policy - Cancellation at an awaiting operation. Activities default to Try, children to Abandon.
- Child
Workflow Failure Kind - Stable terminal categories returned when an awaited child does not succeed.
- Condition
Wait Options Error - Validation failure for a durable condition-wait definition.
- Condition
Wait Result - Unambiguous terminal result of a durable condition wait.
- Error
- Handler
Kind - The registered handler family reported by
Error::HandlerType. - Handler
Value Kind - Whether a typed handler failed to adapt its input or result.
- Parallel
Avro Result - Lossless fixed-Avro counterpart to
ParallelResult. - Parallel
Operation - A deferred durable leaf or nested group for
WorkflowContext::parallel. - Parallel
Result - One input-ordered result returned by
WorkflowContext::parallel. - Parent
Close Policy - Server behavior when a parent closes while its child is still open.
- Registration
Kind - Handler namespaces checked independently within a Worker.
- Search
Attribute Update Error - Validation failure for a typed workflow search-attribute update.
- Search
Attribute Value - One public typed search-attribute value.
- Selection
Key - Stable user-facing identity for one member of a durable selection group.
- Value
- Represents any valid JSON value.
- Worker
Poll Outcome - Stable classification for worker poll responses.
- Workflow
Command Kind - The lifecycle command sent to a workflow execution.
- Workflow
Terminal Kind - Stable terminal categories returned by
WorkflowHandle::result.
Constants§
- AVRO_
VALUE_ SCHEMA_ FINGERPRINT - AVRO_
VALUE_ SCHEMA_ FINGERPRINT_ HEX - AVRO_
VALUE_ SCHEMA_ JSON - Canonical Avro Value schema packaged with the crate and parsed by the runtime.
- CONDITION_
WAIT_ MINIMUM_ WORKER_ PROTOCOL_ VERSION - First additive worker protocol that defines external durable condition waits.
- CONDITION_
WAIT_ OCCURRENCE_ IDENTITY_ CAPABILITY - Worker-registration capability for authored condition-wait occurrence identity.
- CONDITION_
WAIT_ OCCURRENCE_ IDENTITY_ MINIMUM_ WORKER_ PROTOCOL_ VERSION - First additive worker protocol that preserves authored condition-wait occurrences.
- CONTROL_
PLANE_ VERSION - DEFAULT_
CODEC - DURABLE_
SELECTION_ CAPABILITY - Worker-registration capability for persisted first-completion selection.
- DURABLE_
SELECTION_ MINIMUM_ WORKER_ PROTOCOL_ VERSION - First additive worker protocol that defines durable selection groups.
- MEMO_
UPSERTS_ CAPABILITY - Worker-registration capability for portable memo upserts.
- MEMO_
UPSERT_ MINIMUM_ WORKER_ PROTOCOL_ VERSION - First additive worker protocol that defines portable memo upserts.
- MESSAGE_
STREAMS_ CAPABILITY - Worker-registration capability for durable named input streams.
- MESSAGE_
STREAMS_ MINIMUM_ WORKER_ PROTOCOL_ VERSION - MESSAGE_
STREAM_ CURSOR_ SCHEMA - MESSAGE_
STREAM_ MAX_ BATCH - MESSAGE_
STREAM_ SCHEMA - MESSAGE_
STREAM_ SIGNAL - PORTABLE_
WORKER_ AFFINITY_ MINIMUM_ PROTOCOL_ VERSION - First additive worker protocol that defines portable worker-affinity features.
- QUERY_
TASKS_ CAPABILITY - Worker-registration capability for server-routed read-only queries.
- QUERY_
TASK_ MINIMUM_ WORKER_ PROTOCOL_ VERSION - First additive worker protocol that defines query-task transport.
- SDK_
VERSION - SEARCH_
ATTRIBUTE_ UPDATE_ MINIMUM_ WORKER_ PROTOCOL_ VERSION - First additive worker protocol that defines typed search-attribute upserts.
- TYPED_
SEARCH_ ATTRIBUTES_ CAPABILITY - Worker-registration capability for canonical typed search attributes.
- TYPED_
SEARCH_ ATTRIBUTES_ MINIMUM_ WORKER_ PROTOCOL_ VERSION - First additive worker protocol that preserves declared search-attribute types.
- WORKER_
PROTOCOL_ VERSION - WORKFLOW_
UPDATES_ CAPABILITY - Worker-registration capability for synchronous workflow updates.
Functions§
- decode_
avro_ value - decode_
payload - encode_
avro_ value - encode_
payload - portable_
worker_ affinity_ capability_ manifest - Manifest for the default service-worker profile.
- worker_
protocol_ supports_ message_ streams