Crate durable_workflow

Source
Expand description

§Durable Workflow Rust SDK

CI Crates.io API docs Rust 1.86+ License: MIT

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-workflow

Applications 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_world

Pass 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_world

The namespace runtime URL is already complete. Do not append another /api.

§Core API

  • Client starts, signals, queries, updates, cancels, terminates, describes, and awaits workflow executions.
  • Worker registers workflow, activity, query and update handlers, declares workflow signals, and long-polls task queues.
  • WorkflowContext provides 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

ExampleDemonstrates
hello_world.rsTyped worker, workflow, activities, retries, and completion
activity_options.rsActivity retry and timeout policies
local_activities.rsExplicit local execution, retries, durable timers and remote work
worker_sessions.rsTyped activity routing, holder-local resources and graceful session close
sticky_execution.rsOpt-in bounded history, signal replay and cache counters
condition_search_attributes.rsDurable conditions and typed search attributes
continue_as_new.rsBounded histories and continue-as-new
parallel_saga.rsDeterministic parallel work and saga compensation

§Documentation

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

Replay 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::Value from a JSON literal.
wait_condition
Create a durable condition wait whose predicate definition is fingerprinted from its inline Rust tokens.

Structs§

ActivityCall
ActivityContext
ActivityFailure
A stable, machine-readable terminal activity failure.
ActivityHeartbeatResponse
ActivityOptions
Options recorded atomically on one deterministic schedule_activity command.
ActivityOptionsError
A machine-readable activity-options validation failure.
ActivityRetryPolicy
Durable server-side retry policy for one activity execution.
ActivityTask
ActivityTaskRejection
A worker-side activity settlement or heartbeat rejected by durable state.
CancelDurableOperationCall
Future returned by DurableOperationHandle::cancel.
CancellationContext
Original cancellation metadata restored from canonical request history.
CancellationDelivery
One committed delivery marker, bound to the original cancellation request.
CancellationDeliveryReceipt
A delivery acknowledgment bound to one workflow task and durable run.
CancellationHistory
Validated canonical cancellation state for one workflow run.
CancellationLineage
One immutable local request in the ordered cancellation lineage.
CancellationRequest
The Server’s original request identity and immutable cleanup deadline.
CancellationScopeOpenReceipt
A scope opening proved against complete history on its original claim. This proof does not advertise scoped workflow execution support.
CancellationShield
A workflow-local cleanup scope. Dropping it restores cancellation checks.
ChildWorkflowAvroResult
Lossless successful child result for fixed Avro Value workflows.
ChildWorkflowCall
Future returned by WorkflowContext::start_child_workflow.
ChildWorkflowFailure
A stable, machine-readable child workflow failure delivered to its parent.
ChildWorkflowOptions
Options recorded with a child-workflow command.
ChildWorkflowResult
A successful child result together with its durable parent-child identity.
ChildWorkflowRetryPolicy
Durable retry policy for one child workflow invocation.
Client
ClientBuilder
ConditionWaitCall
Future returned by WorkflowContext::wait_condition.
ConditionWaitOptions
Stable identity and optional durable timeout for a condition wait.
ContinueAsNewOptions
Optional routing overrides for a continue-as-new transition.
ContinueAsNewOptionsError
A stable validation error raised before a continue-as-new command is emitted.
CooperativeCancellationOptions
Optional reason and runtime-owned cleanup limit for a cooperative request.
CooperativeCancellationRequested
Cancellation delivered at its committed authored call, with original identity.
CooperativeWorkflowTask
One actual protocol 1.20 claim and its original pending observation.
CooperativeWorkflowTaskPoll
An explicit cooperative poll, including ordinary idle and stop outcomes.
DuplicateRegistrationError
Ambiguous local configuration, rejected before contacting Server.
DurableOperationAwaitCall
Future returned by DurableOperationHandle::await_result.
DurableOperationCancelled
Typed result of explicitly awaiting a cancelled non-winning operation.
DurableOperationHandle
Stable reference to one member of a durable selection group.
HandlerRegistration
The authoring adapter and source location of a handler registration.
HistoryEvent
LocalActivityOptions
Retry and timeout settings for an activity executed by the workflow worker.
MessageStream
MessageStreamMessage
ParallelCall
Future returned by WorkflowContext::parallel.
ParallelCompletion
One successful leaf retained when another parallel member failed.
ParallelFailure
A deterministic join failed after some siblings had already completed.
ParallelGroupError
Stable validation error returned before an invalid group emits commands.
ParallelGroupMetadata
Stable identity for one enclosing deterministic parallel group.
PayloadEnvelope
PollActivityTaskResponse
PollQueryTaskResponse
PollWorkflowTaskResponse
ProtocolFailure
A stable failure returned when a server rejects an SDK protocol version.
QueryContext
Immutable state supplied to a registered query handler.
QueryFailure
A stable, machine-readable workflow query or query-task settlement failure.
QuerySignal
One decoded signal in the committed workflow-history snapshot.
QueryTask
An ephemeral server-routed query task.
RegisterWorkerResponse
ReplayFailure
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.
SagaCompensationFailure
A forward saga failure followed by a terminal compensation failure.
ScopedCancellationContext
Original scope ancestry carried by a candidate cooperative child request. Reading this immutable metadata does not authorize entering a scope body.
ScopedCancellationLineage
One immutable scope address and its bounded cleanup deadline.
SearchAttributeUpdate
Validated typed workflow-side search-attribute mutation.
SelectCall
Future returned by WorkflowContext::select.
SelectionResult
The one winner committed for a durable selection group.
SignalCall
StickyCacheMetrics
Replay counters and retained encoded history. Decoding/replay memory is additional.
StickyCacheOptions
Opt-in bounds for a worker’s durable-history cache.
TimerCall
Future returned by WorkflowContext::sleep.
Uuid
A Universally Unique Identifier (UUID).
Worker
WorkerDeregistrationEnvelope
Result of gracefully removing a worker-plane registration.
WorkerHeartbeatObservation
WorkerRetryPolicy
Bounded retry policy for worker poll acquisition and worker heartbeats.
WorkerSession
One worker’s acknowledged session lease.
WorkerSessionOptions
Routing and lifetime options for one worker-held session.
WorkflowCancellationRequest
Acknowledgment of a request targeting the current durable run.
WorkflowCancellationRequested
Cooperative workflow cancellation observed at an author-controlled point.
WorkflowCommandOptions
Optional structured fields for a cancellation or termination request.
WorkflowCommandRejection
A stable rejection returned by instance- or selected-run lifecycle commands.
WorkflowCommandResult
The accepted, machine-readable result of a lifecycle command.
WorkflowContext
WorkflowDescription
WorkflowHandle
WorkflowHistoryBudget
Public history-budget information attached to the current workflow task.
WorkflowIdentity
The identity of one durable workflow execution.
WorkflowInstance
Typed local state owned by one deterministic workflow invocation.
WorkflowRedriveResult
A new run continuing a failed run from its recorded activity boundary.
WorkflowResultOptions
WorkflowStartOptions
Server-enforced timeout policy for a workflow start.
WorkflowStreamAppendItem
One item for direct or replay-safe append.
WorkflowStreamAppendResult
Durable acceptance and deduplication outcome for an append request.
WorkflowStreamDescription
Lifecycle and backlog metadata for one run-scoped Workflow Stream.
WorkflowStreamItem
One durable item at its stable zero-based offset.
WorkflowStreamPage
One bounded at-least-once subscription page.
WorkflowTask
WorkflowTaskHeartbeat
Successful renewal of the exact workflow-task claim with optional observation.
WorkflowTerminalOutcome
A typed terminal workflow outcome with durable identity and failure metadata.

Enums§

ActivityBackoff
Backoff intervals for one durable activity retry policy.
ActivityFailureKind
Stable terminal categories returned when an awaited activity does not succeed.
ActivityOptionsErrorKind
Stable validation categories for ActivityOptions.
AvroValue
Native adapter for the fixed language-neutral Avro Value schema.
CancellationCallKind
The authored durable call where canonical cancellation is delivered.
CancellationDeliveryReply
Delivery either commits or parks the parent until canonical child cleanup.
CancellationPolicy
Cancellation at an awaiting operation. Activities default to Try, children to Abandon.
ChildWorkflowFailureKind
Stable terminal categories returned when an awaited child does not succeed.
ConditionWaitOptionsError
Validation failure for a durable condition-wait definition.
ConditionWaitResult
Unambiguous terminal result of a durable condition wait.
Error
HandlerKind
The registered handler family reported by Error::HandlerType.
HandlerValueKind
Whether a typed handler failed to adapt its input or result.
ParallelAvroResult
Lossless fixed-Avro counterpart to ParallelResult.
ParallelOperation
A deferred durable leaf or nested group for WorkflowContext::parallel.
ParallelResult
One input-ordered result returned by WorkflowContext::parallel.
ParentClosePolicy
Server behavior when a parent closes while its child is still open.
RegistrationKind
Handler namespaces checked independently within a Worker.
SearchAttributeUpdateError
Validation failure for a typed workflow search-attribute update.
SearchAttributeValue
One public typed search-attribute value.
SelectionKey
Stable user-facing identity for one member of a durable selection group.
Value
Represents any valid JSON value.
WorkerPollOutcome
Stable classification for worker poll responses.
WorkflowCommandKind
The lifecycle command sent to a workflow execution.
WorkflowTerminalKind
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

Type Aliases§

Result