@superdurable/dex
v1.4.0
Published
TypeScript SDK for the Dex durable workflow engine
Readme
Dex SDK for TypeScript
Flow state I/O
Applications read and write Flow state through typed RPCs. The Client does not expose direct Attribute reads or writes, Channel publication, or pending-message mutation. This keeps each external state transition behind a Flow-owned method.
An RPC can stage channel.delete(context, messageId). Set the decorator option
isTransactional: true when a missing ID must abort all other RPC writes.
Attribute locks already select transactional execution; Channel deletion requires
the explicit option.
RPC state loading
RPCs receive ordinary Attributes and all Channel size metadata by default. AttributeMap entries and pending Channel messages are opt-in:
@rpc({
loadAttributeMapInstances: [items.load("tenant-a")],
loadChannels: [queued],
loadChannelMaps: [byTenant],
})
snapshot(context: Context): RPCResult<Snapshot> {
return { output: { messages: queued.pendingMessages(context) } };
}Put an AttributeMap or ChannelMap directly in its plural load option when every current instance is needed. Use the singular instance options for the less common exact-instance case. A selected empty queue returns an empty array; reading an unselected map entry throws AttributeMapNotLoadedError. An unselected pending-message snapshot throws ChannelMessagesNotLoadedError. Pending messages preserve FIFO order and include the server-assigned message ID. The snapshot does not change after the handler stages a publish or deletion.
Loading controls which data reaches the Worker. isTransactional controls atomic commit and Channel deletion validation. Attribute locks add isolation only among cooperating Steps and RPCs using the same lock. Write-only publish and delete operations do not require loading.
Invocation-time map instances
Decorator RPCOptions remains the registration-time contract for timeout,
transactionality, locks, and fixed state loads. Use invokeRPCWithOptions when
the caller chooses an exact map instance at runtime:
const result = await client.invokeRPCWithOptions(
directory.upsertCustomerProfile,
flowId,
profile,
{
lockAttributeMapInstances: [profiles.lock(partitionName)],
loadAttributeMapInstances: [profiles.load(partitionName)],
},
);The separate method avoids ambiguity with the positional runId. Selections are
additive, sorted, and deduplicated. Locking and loading are independent;
read-modify-write handlers need both for the same instance.
See the runnable sequential chunking and hash partitioning patterns for complete call sites.
Step and timeout-handler state loading
StepOptions provides the same five selections independently for waitFor and
execute: all AttributeMap instances, exact AttributeMap instances, Channels,
all ChannelMap instances, and exact ChannelMap instances. The methods receive
independent snapshots. Execute reads after the winning Wait consumes messages;
retries of one logical method call reuse its first snapshot.
FlowTimeoutHandlerOptions provides Execute-style timeouts, heartbeat detection,
retry, durability, Attribute locks, and the same state selections. Set it on
StartFlowOptions or SubFlowOptions only for a positive timeout using HANDLER.
Exhausted retries may proceed to a registered Step<void> using voidCodec;
read the final failure from Context.recoveryError.
class TimeoutRecoveryStep implements Step<void> {
public readonly inputCodec = voidCodec;
public getStepType(): string {
return "TimeoutRecoveryStep";
}
public execute(context: Context, _input: void): StepDecision {
const failure = context.recoveryError;
return forceFail(failure?.detail ?? "timeout recovery has no failure");
}
}This package targets Node.js 22 and 24. It provides strongly typed workflow
contracts and a Promise-based gRPC Client. The Client and Worker runtime use
@grpc/grpc-js. Blob caching uses the shared Rust DXBC implementation through
a Node-API addon.
Application values use Codec<T>. Flow, Step, RPC, Attribute, Channel, and Stream
definitions retain their input and output types. Client methods return Promise
because Node network I/O is asynchronous.
The Client uses @grpc/grpc-js directly. Rust is only the implementation
boundary for the shared BlobCache; TypeScript callbacks and network transport
stay in Node.
Omitted Step input codecs and RPC input/output codecs use JSON
(JSON.stringify / JSON.parse, no structural validation). Scalar wire kinds
still need stringCodec, booleanCodec, int64Codec, doubleCodec, or
bytesCodec. A @rpc() method with only Context and a void return stays a
procedure.
class ApproveOrder implements Step<{ orderId: string }> {
getStepType(): string {
return "ApproveOrder";
}
waitFor(_context: Context, _order: { orderId: string }): Wait {
return Wait.until(Timer.byDuration(1_000));
}
execute(_context: Context, order: { orderId: string }): StepDecision {
return gracefulComplete(order);
}
}
class Orders implements Flow<{ orderId: string }> {
readonly approve = new ApproveOrder();
getFlowType(): string {
return "Orders";
}
getSteps() {
return StepList.startStep(this.approve);
}
}
const orders = new Orders();
const registry = new Registry([orders]);Flows return all Steps once. Start with StepList.startStep(step) and append
heterogeneous Steps with .otherSteps(...). Use
StepList.withoutStartStep<void>(...) for RPC-triggered Steps, or
StepList.empty() when the Flow has no Steps.
Flow<StartInput> only types the starting Step and Client.startFlow() input;
StepList<StartInput> enforces that relationship during type checking.
Non-starting Steps may use unrelated input types. Flow defaults to void for
Flows without a start input.
Worker.start() synchronizes every registered Indexed Attribute with Dex
Server before binding its listener. Existing indexes return immediately;
failure or the default 120-second deadline aborts startup. An indexed
AttributeMap must provide one fixed indexKey.
StepOptions.waitForMethodTimeoutMs and executeMethodTimeoutMs bound the two
handler calls. Timer and channel conditions determine how long a Step waits.
Canceling Step executions
A successful Step can cancel queued or active executions while continuing with its normal decision:
return withCancelingSteps(
withCancelingSiblingSteps(
goTo(RecordQuote, quote),
QuoteCarrierA,
QuoteCarrierB,
),
GlobalQuoteTimeout,
);withCancelingSteps selects every current execution of each registered Step
type. withCancelingSiblingSteps selects only executions whose
Context.fromStepExecutionId matches the current execution. Both helpers return
a new decision; repeated calls form a union, and Flow-wide selection wins for
the same Step class. Unregistered selectors produce an invalid Step result.
Dex resolves one snapshot after the current execution succeeds. Completed, already-canceled, and absent targets are no-ops. Next Steps created by that decision are outside the snapshot. Dex immediately applies the next or close action; late decisions, writes, retries, and recovery Steps are discarded.
An RPC result may set cancelingSteps for Flow-wide selection. RPCs do not
support sibling selection because they have no Step execution lineage.
Soft Flow timeout
An optional Flow.handleTimeout makes a positive timeout use handler policy by
default. It may be synchronous or async and returns a normal StepDecision:
async handleTimeout(context: Context): Promise<StepDecision> {
await notifyExpiration(context);
return forceComplete("expired");
}
await client.startFlow(orders, "order-42", input, {
timeoutMs: 30 * 60_000,
timeoutPolicy: FlowTimeoutPolicy.HANDLER,
});FAIL produces FlowErrorType.FLOW_TIMEOUT and permits Flow retry; CANCEL
cancels without retry. Continue-as-new preserves the durable timer's deadline;
retry runs get a fresh budget. A zero or omitted timeout disables the
feature.
Opt an Attribute or AttributeMap into Attribute Store synchronization, then select Server-configured Stores for the Flow:
const email = new Attribute("customer-email", stringCodec)
.syncToAttributeStore();
const config: FlowConfig = { attributeStoreNames: ["profiles", "audit"] };Stores are asynchronous latest-state projections. Every enabled Attribute write
is sent to every selected Store. Deletion writes SQL NULL, and projection
failures do not roll back Flow Attributes. Omitting attributeStoreNames
preserves current targets; attributeStoreNames: [] disables future
synchronization while retaining protocol presence.
Waiting and map inspection
Wait.until, Wait.allOf, and Wait.anyOf use unnamed Conditions by default.
Do not add Condition IDs merely because a Condition is nested in one of these
waits. Every Condition in Wait.anyCombinationOf needs a non-empty user ID;
the same Condition object may be reused across combinations.
Client.waitForAttributeMatch overloads target the current run and return the
decoded matched value. Build a match with one of the six AttributeMatch
factories. String and boolean support equality operators; integer and double
support every operator. JSON, bytes, null, non-finite doubles, and invalid
ordering reject before transport. Every AttributeMap and ChannelMap
instance must be non-empty and must not contain /. AttributeMap.getMapSize
and getAllInstanceKeys include buffered sets and deletes. The matching
ChannelMap methods are RPC-only, include buffered publishes, and omit empty
instances. Keys are decoded and sorted. Use
forceCompleteIfChannelsEmpty(...) for conditional completion.
Request IDs are optional for both durable waits. When omitted, the server
derives a stable ID from the Step execution or Attribute condition. Reuse an
override only for the same wait. requestTimeoutMs is the total SDK call budget
across transparent transport reattachments. Omit it or use zero to wait
indefinitely. Temporal permits 10 in-flight Updates and 2,000 total Updates in
History per Workflow Execution. A caller timeout or cancellation does not
finish an accepted durable wait, so an abandoned handler can retain one
in-flight slot. Set internalHandlerTimeoutMs only when abandoned waits can
approach that limit. An active call transparently starts another generation,
which counts toward the 2,000-Update history limit. Prefer a value comfortably
longer than normal request timeouts and reconnect gaps. Zero disables timed
rollover. Request expiry throws RequestTimeoutError.
Async handlers
Step.execute, Step.waitFor, and RPC methods may be async and return a
Promise of their result. Because the Worker runs on Node's single event loop,
an async handler that awaits Client calls (startFlow, waitForFlow,
invokeRPC, ...) yields the loop, so the same Worker keeps serving other
WorkerService calls — including the child Flow or RPC you just started. One
Worker is enough; no sidecar Worker is needed.
async execute(context: Context, childId: string): Promise<StepDecision> {
await client.startFlow(childFlow, childId, input);
const result = await client.waitForFlow(childId, 5_000);
const output = result.singleOutput(stringCodec);
return gracefulComplete(output);
}Synchronous handlers stay valid — returning a plain StepDecision / Wait /
RPCResult needs no change. Two rules keep the loop healthy:
- Never block the event loop from a handler (
Atomics.wait,child_process.spawnSync). That freezes the Worker and defeats the point. - A long
await client.waitForFlow(...)holds the current WorkerService call until it resolves or times out. Prefer a short timeout plus Step retry (or a timer-backed re-check) over one unbounded poll.executeMethodTimeoutMs/waitForMethodTimeoutMsclock the full async handler duration.
Step progress and Streams
Annotate a Promise-returning Step handler with AsyncContext when it needs to
record progress. The heartbeat value argument is required. Passing undefined
explicitly sends a heartbeat without a Value and clears previously persisted
heartbeat details. Passing null sends a present JSON-null Value.
async execute(context: AsyncContext, input: ImportInput): Promise<StepDecision> {
const restored = context.hasLastHeartbeatValue()
? context.getLastHeartbeatValue(importCheckpointCodec)
: undefined;
for await (const page of remainingPages(restored)) {
importedRows.write(context, page.rows);
await context.recordHeartbeat(page.checkpoint, importCheckpointCodec);
}
return gracefulComplete();
}Omitting the heartbeat codec uses JSON. Read a restored value with the same
codec used to record it. hasLastHeartbeatValue() distinguishes an absent Value
from a present JSON-null Value; both decode to undefined.
Stream.write(context, value) is synchronous and fire-and-forget. It validates
and encodes locally, then writes a frame to the current WaitFor or Execute gRPC
response stream. A handler may write any number of messages to the same or
different Streams. Dex attempts each append in call order, but Stream Store
failures are not returned to the handler. RPC and Flow timeout Contexts reject
Step Stream writes.
For text deltas, create one writer and pass its bound write method directly
to the producer:
const progress = thinking.bufferedText(context, { flushIntervalMs: 500 });
await generateText(input, progress.write);The defaults are one second and a soft 16 KiB UTF-8 threshold. Invocation completion sends the tail before the final result or error. Chunks are concatenated exactly. Retry does not restore unsent text or deduplicate emitted batches.
External writes use client.writeStream(flowId, stream, source, value). Source
must be non-empty, may contain #, and may repeat; every write appends a new
message. Client.readStream returns it as StreamMessage.source. Step writes
use #<stepExecutionID> as source metadata.
Client.readStream moves forward one message at a time and can long-poll.
Client.listStreamMessages returns a non-blocking newest-first retained page.
Pass the registered typed Stream, a positive page size, and an empty token for
the first page. Pass StreamMessagesPage.nextPageToken unchanged to read older
messages until it is empty. The server caps page size at 1000 by default. A
trimmed anchor returns an empty page, and newer concurrent writes do not enter
an existing older-page chain.
Step durability defaults to the Flow configuration and then sync. Regular attempts default to two hours, heartbeat timeout defaults to one minute, and retry total duration defaults to four hours. The SDK accepts non-negative whole-second heartbeat timeout values, including two seconds, and leaves the deployment-configured minimum to Dex Server. Async durability first uses a seven-second local-activity window with at most three attempts; that phase ignores method timeouts and heartbeats but still writes Stream messages.
Every TypeScript Flow and Step must return an explicit durable name from
getFlowType() or getStepType(). Class names are never used as fallbacks
because bundlers and minifiers may rename them.
waitForFlow hydrates every output-bearing completion before resolving. For a
multi-output Flow, select by Step identity and decode each value with its codec:
const result = await client.waitForFlow(flowId);
const receipt = result.completions
.find((completion) => completion.stepType === "ChargeCard")
?.decode(receiptCodec);The completion array is read-only and keeps server collection order. Parallel
branch order is not deterministic. A no-output Flow returns an empty array;
singleOutput throws for zero or multiple completions. Every terminal status returns
a FlowResult; inspect status, errorType, and errorMessage for unsuccessful
completion.
SubFlows are normal, independently addressable Flows used as durable Conditions:
public waitFor(_context: Context, input: ChargeInput): Wait {
return Wait.until(SubFlow.run(this.chargeFlow, input));
}
public execute(context: Context, _input: ChargeInput): StepDecision {
const receipt = SubFlow.getConditionResults(context).singleOutput(receiptCodec);
return gracefulComplete(receipt);
}SubFlow.getFlowId(context, index) remains available for a running anyOf loser.
SubFlowOptions configures timing, timeout policy, retry, initial target Attributes,
Flow config, Condition ID, and reuse. Parent completion does not cancel an unfinished SubFlow.
Errors
Client calls reject with concrete DexServiceError subclasses. Existing-Flow
reads (describeFlow, waitForFlow, and timeTravel) use
FlowNotFoundError; operations that require a running Flow use
FlowNotActiveOrNotFoundError. Start conflicts, worker failures, RPC lock contention,
and long-poll timeouts use FlowAlreadyStartedError,
WorkerInvocationError, RpcLockConflictError, and LongPollTimeoutError.
RPCs also use this error when the target Flow cannot be found, including query-only
reads. A query-only RPC without locks, transactional execution, durable effects, or
Server-forced Update routing can read a retained terminal execution. At that confirmed
read boundary, this error means no readable target was found; return the application's
missing or unavailable result directly, without a lifecycle probe, a short timeout, a
retry, or historical Step-output decoding. Preserve other service and Worker failures.
For mutations and active-only RPCs, the error can mean either missing or closed and does
not prove that the requested action succeeded. RPC success does not prove that the Flow
is active.
Durable Step and Attribute waits reattach transparently instead of exposing
LongPollTimeoutError.
try {
await client.invokeRPC(orders.updateOrder, flowId, update);
} catch (error) {
if (error instanceof FlowNotActiveOrNotFoundError) {
// The Flow is missing or already closed.
} else {
throw error;
}
}Every service error retains code, subStatus, detail, operation,
flowId, and the original gRPC error as cause. WorkerInvocationError also
retains workerCode, workerErrorType, and workerErrorDetail. Registration,
serialization, and invalid handler returns use FlowDefinitionError,
ValueMappingError, and InvalidStepResultError instead of transport errors.
Source layout
Public contracts are grouped by domain under src/. The root src/index.ts
is a barrel that re-exports the supported package API; applications should
continue importing only from @superdurable/dex.
codec.ts: wire values and codecspersistence.ts: attributes, indexes, locks, and schemaswait.ts: channels, timers, conditions, and waitsstep.ts: Steps, movements, options, and decisionsrpc.ts: typed RPC contracts and decoratorsflow.ts: Flows, registration, and validationclient.ts: Promise-based FlowService Clientworker.ts: Worker gRPC service and lifecycleworker-dispatcher.ts: typed callback dispatch and response mappinginvocation-context.ts: invocation-scoped persistence and condition stateblob-cache.ts: DXBC BlobCache contract and Node-API binding loadergen/: checked-in protobuf and grpc-js bindings
Serialized object values use wire encoding json for JSON and raw for raw bytes.
Internal Blob references are opaque. Hydration sends the owning Flow ID with
each reference, and the cache is isolated by (flowId, blobRef).
Run npm run build:native once to stage the DXBC Node addon for the current
platform, then npm test for runtime contracts, npm run typecheck for strict
static contracts, and npm run docs:check for public API documentation. The
documentation check follows the actual exports of src/index.ts and requires
JSDoc for every public class, interface, type, overload, member, object-style
enum value, type parameter, input, and output; generated sources are excluded.
The comments appear in TypeScript language-service and IDE hovers.
Run ./run-integration-tests.sh for all 58 IWF
compatibility scenarios against an isolated dexcli dev environment. After
changing protos/dex.proto, run make generated-code from the repository root
and commit every server and SDK output.
Integration coverage
Run the complete integration suite with TypeScript source coverage:
npm run coverage:integrationThe terminal report lists coverage per SDK source file and every uncovered
line. Open coverage/index.html for annotated source, or inspect
coverage/coverage-summary.json programmatically. coverage/lcov.info is the
report uploaded by CI. Generated protobuf code under src/gen/ is excluded.
CI uploads the LCOV report to Codecov with GitHub OIDC, so no upload secret is
stored in this repository. The report uses the sdk-typescript-integration
flag and contributes to the TypeScript SDK component defined in the root
codecov.yml. After the first successful main upload, Codecov displays
project and patch coverage in its dashboard, GitHub checks, and PR comments.
The Actions run also publishes the complete HTML report as
sdk-typescript-integration-coverage.
The complete legacy IWF integration inventory lives under
test/integ. Its Flow fixtures retain the Java
suite's workflow behavior and its 58 assertions run against a real Dex server.
Releases
The npm package is published as @superdurable/dex.
After the initial bootstrap, publish a GitHub Release tagged
sdk-typescript/vX.Y.Z. The tag is the release version source of truth; CI
temporarily writes it to package.json and package-lock.json without
committing either file. The release workflow then runs type checks and tests,
inspects the tarball, and publishes through npm Trusted Publishing. Prerelease
versions use the next npm dist-tag; stable versions use latest.
Trusted Publishing can only be configured after the package exists. Bootstrap the first version from a maintainer workstation with 2FA:
cd sdk-typescript
npm ci
npm run typecheck
npm test
npm pack --dry-run
npm login
npm publish --access publicThen open the package settings on npmjs.com and add a GitHub Actions trusted
publisher with organization superdurable, repository dex, workflow
sdk-typescript-publish.yml, no environment, and npm publish permission.
Future releases use short-lived OIDC credentials and require no NPM_TOKEN.
After verifying the first OIDC release, configure npm publishing access to
require 2FA and disallow token-based publication.
Server protocol compatibility
Worker startup calls GetServerInfo, negotiates the highest common protocol,
synchronizes Attribute indexes, and then binds WorkerService. The initial
TypeScript SDK interval is [1,1]. Missing information, invalid or disjoint
intervals, and RPC failures stop startup before binding.
The build generates the diagnostic SDK version from package.json, so the
runtime value follows the version stamped by the release workflow. The artifact
version does not participate in compatibility decisions.
Search run inclusion
Dex Server excludes ContinuedAsNew runs by default. Application queries still
need FlowType when using shared index slots. An explicit inclusion option
removes only the default exclusion; it does not override query filters. Keep
the query, inclusion option, and page size consistent across pages.
Pass { includeContinuedAsNew: true } as the fourth argument of searchFlows(query, pageSize, nextPageToken, options) for execution-chain inspection. The default is false.
