Workflows
A workflow.lua entry defines a durable Temporal workflow that orchestrates activities and maintains state across failures and restarts.
This page is an API reference with partial recipes. Entry declarations, worker registration, activity implementations, security policies, and surrounding application data are shown only where they are relevant to a specific contract.
Definition
- name: order_workflow
kind: workflow.lua
source: file://order_workflow.lua
method: main
modules:
- funcs
- time
- workflow
meta:
temporal:
workflow:
worker: app:worker
Metadata Fields
| Field | Required | Description |
|---|---|---|
worker |
Yes | Reference to temporal.worker entry |
name |
No | Custom workflow type name (defaults to entry ID) |
Basic Implementation
local funcs = require("funcs")
local time = require("time")
local function main(order)
local payment, err = funcs.call("app:charge_payment", {
amount = order.total,
customer = order.customer_id
})
if err then
return {status = "failed", error = tostring(err)}
end
time.sleep("1h")
local shipment, err = funcs.call("app:ship_order", {
order_id = order.id,
address = order.shipping_address
})
if err then
local _, refund_err = funcs.call("app:refund_payment", payment.id)
if refund_err then
return {
status = "failed",
error = tostring(err),
compensation_error = tostring(refund_err)
}
end
return {status = "failed", error = tostring(err)}
end
return {
status = "completed",
payment_id = payment.id,
tracking = shipment.tracking_number
}
end
return { main = main }
Workflow Module
The workflow module provides workflow-specific operations.
workflow.info()
Get workflow execution information:
local workflow = require("workflow")
local info, info_err = workflow.info()
if info_err then return nil, info_err end
print(info.workflow_id) -- Workflow execution ID
print(info.run_id) -- Current run ID
print(info.workflow_type) -- Workflow type name
print(info.task_queue) -- Task queue name
print(info.namespace) -- Temporal namespace
print(info.attempt) -- Current attempt number
print(info.history_length) -- Number of history events
print(info.history_size) -- History size in bytes
workflow.exec()
Execute a child workflow synchronously and wait for its result:
local result, err = workflow.exec("app:child_workflow", input_data)
if err then
return nil, err
end
Use this form when the parent must wait for the child result inline.
workflow.version()
Handle code changes with deterministic versioning:
local version, err = workflow.version("payment-v2", 1, 2)
if err then
return nil, err
end
if version == 1 then
return funcs.call("app:old_payment", input)
else
return funcs.call("app:new_payment", input)
end
Parameters:
change_id- Unique identifier for this changemin_supported- Minimum supported versionmax_supported- Maximum (current) version
The version number is deterministic per workflow execution. Existing in-flight workflows continue using their recorded version, while new workflows use max_supported.
workflow.attrs()
Update search attributes and memo:
local updated, err = workflow.attrs({
search = {
status = "processing",
customer_id = order.customer_id,
order_total = order.total
},
memo = {
notes = "Priority customer",
source = "web"
}
})
if err then
return nil, err
end
Search attributes are indexed and queryable via Temporal visibility APIs. Memo is arbitrary non-indexed data attached to the workflow.
workflow.history_length() / workflow.history_size()
Monitor workflow history growth:
local length, length_err = workflow.history_length()
if length_err then return nil, length_err end
local size, size_err = workflow.history_size()
if size_err then return nil, size_err end
if length > 10000 then
-- Consider continue-as-new to reset history
end
Starting Workflows
Basic Spawn
Start a workflow from any code using process.spawn():
local pid, err = process.spawn(
"app:order_workflow", -- workflow entry
"app:worker", -- temporal worker
{order_id = "123"} -- input
)
if err then
return nil, err
end
The host parameter is the temporal worker (not a process host). The workflow runs durably on Temporal infrastructure.
Spawn with Monitoring
Monitor workflows to receive EXIT events when they complete:
local pid, err = process.spawn_monitored(
"app:order_workflow",
"app:worker",
{order_id = "123"}
)
if err then
return nil, err
end
local events = process.events()
local event, open = events:receive()
if not open then
return nil, errors.new({kind = errors.INTERNAL, message = "process event channel closed"})
end
if event.kind == process.event.EXIT then
local result = event.result.value
local error = event.result.error
end
Spawn with Name
Assign a name to a workflow for idempotent starts:
local spawner = process
.with_options({})
:with_name("order-" .. order.id)
local pid, err = spawner:spawn_monitored(
"app:order_workflow",
"app:worker",
{order_id = order.id}
)
if err then
return nil, err
end
When a name is provided, Temporal uses it to deduplicate workflow starts. Spawning with the same name while a workflow is running returns the existing workflow's PID by default.
Spawn with Explicit Workflow ID
Set a specific Temporal workflow ID:
local spawner = process
.with_options({
["workflow.id"] = "order-" .. order.id,
})
local pid, err = spawner:spawn_monitored(
"app:order_workflow",
"app:worker",
order
)
if err then
return nil, err
end
ID Conflict Policies
Control behavior when spawning a workflow with an ID that already exists:
-- Fail if workflow already exists
local spawner = process
.with_options({
["workflow.id"] = "order-123",
["workflow.id_conflict_policy"] = "fail",
})
local pid, err = spawner:spawn("app:order_workflow", "app:worker", order)
if err then
-- Workflow already running with this ID
end
-- Error when already started (alternative approach)
local spawner = process
.with_options({
["workflow.id"] = "order-123",
["workflow.execution_error_when_already_started"] = true,
})
local pid, err = spawner:spawn("app:order_workflow", "app:worker", order)
if err then return nil, err end
-- Reuse existing (default behavior with explicit ID)
local spawner = process
.with_options({
["workflow.id"] = "order-123",
})
local pid, err = spawner:spawn("app:order_workflow", "app:worker", order)
if err then return nil, err end
-- Returns existing workflow PID if already running
| Policy | Behavior |
|---|---|
"use_existing" |
Return existing workflow PID (default with explicit ID) |
"fail" |
Return error if workflow exists |
"terminate_existing" |
Terminate existing and start new |
Workflow Start Options
Pass Temporal workflow options via with_options():
local spawner = process.with_options({
["workflow.id"] = "order-123",
["workflow.execution_timeout"] = "24h",
["workflow.run_timeout"] = "1h",
["workflow.task_timeout"] = "30s",
["workflow.id_conflict_policy"] = "fail",
["workflow.retry_policy"] = {
initial_interval = 1000,
backoff_coefficient = 2.0,
maximum_interval = 300000,
maximum_attempts = 3,
},
["workflow.cron_schedule"] = "0 */6 * * *",
["workflow.search_attributes"] = {
customer_id = "cust-123"
},
["workflow.memo"] = {
source = "api"
},
["workflow.start_delay"] = "5m",
["workflow.parent_close_policy"] = "terminate",
})
Options Reference
| Option | Type | Description |
|---|---|---|
temporal.workflow.id |
string | Explicit workflow execution ID |
temporal.workflow.task_queue |
string | Override task queue |
temporal.workflow.execution_timeout |
duration | Total workflow execution timeout |
temporal.workflow.run_timeout |
duration | Single run timeout |
temporal.workflow.task_timeout |
duration | Workflow task processing timeout |
temporal.workflow.id_conflict_policy |
string | use_existing, fail, terminate_existing |
temporal.workflow.id_reuse_policy |
string | allow_duplicate, allow_duplicate_failed_only, reject_duplicate |
temporal.workflow.execution_error_when_already_started |
boolean | Error if workflow already running |
temporal.workflow.retry_policy |
table | Retry policy (see below) |
temporal.workflow.cron_schedule |
string | Cron expression for recurring workflows |
temporal.workflow.memo |
table | Non-indexed workflow metadata |
temporal.workflow.search_attributes |
table | Indexed queryable attributes |
temporal.workflow.enable_eager_start |
boolean | Start execution immediately |
temporal.workflow.start_delay |
duration | Delay before workflow starts |
temporal.workflow.parent_close_policy |
string | Child behavior on parent close |
temporal.workflow.wait_for_cancellation |
boolean | Wait for cancellation to finish |
temporal.workflow.namespace |
string | Temporal namespace override |
temporal.workflow.name |
string | Workflow type name to start, when it differs from the registry ID |
temporal.workflow.versioning_intent |
string | compatible (inherit the build ID) or default (use assignment rules) |
temporal.workflow.priority |
table | Task priority: priority_key (number), fairness_key (string), fairness_weight (number) |
workflow.summary |
string | Human-readable summary shown in the Temporal UI |
workflow.details |
string | Human-readable details shown in the Temporal UI |
workflow.versioning_override |
table | Worker versioning override: mode is auto_upgrade, or pinned with deployment_name and build_id |
Every option is also accepted under its short key (workflow.id, workflow.task_queue, ...); the temporal.workflow. prefix is a legacy alias. summary and details have no temporal.workflow. alias.
Duration values accept strings ("5s", "10m", "1h") or milliseconds as numbers.
Legacy temporal.workflow.* aliases remain accepted for compatibility. New code should use the canonical workflow.* names shown above.
A pinned version override requires both the mode and deployment version:
["workflow.versioning_override"] = {
mode = "pinned",
version = {
deployment_name = "orders",
build_id = "orders-v2",
},
}
Use the string "auto_upgrade" for an auto-upgrade override.
Parent Close Policy
Controls what happens to child workflows when the parent closes:
| Policy | Behavior |
|---|---|
"terminate" |
Terminate child workflow |
"abandon" |
Let child continue independently |
"request_cancel" |
Send cancellation request to child |
Startup Messages
Queue signals with a workflow start. The first non-empty startup message is sent atomically with the start. Remaining startup messages are sent sequentially in builder order after the workflow starts, but they can interleave with signals sent concurrently by other callers:
local spawner = process
.with_options({})
:with_name("counter-workflow")
:with_message("increment", {amount = 2})
:with_message("increment", {amount = 1})
:with_message("increment", {amount = 4})
local pid, err = spawner:spawn_monitored(
"app:counter_workflow",
"app:worker",
{initial = 0}
)
if err then return nil, err end
With the use_existing conflict policy, startup messages are still delivered when a second spawn resolves to an existing workflow:
-- First spawn starts the workflow with initial messages
local first = process
.with_options({})
:with_name("my-counter")
:with_message("increment", {amount = 3})
local pid, first_err = first:spawn("app:counter_workflow", "app:worker", {initial = 0})
if first_err then return nil, first_err end
-- Second spawn reuses existing workflow and delivers new messages
local second = process
.with_options({})
:with_name("my-counter")
:with_message("increment", {amount = 2})
local pid2, second_err = second:spawn("app:counter_workflow", "app:worker", {initial = 999})
if second_err then return nil, second_err end
-- pid2 == pid (same workflow), input {initial = 999} is ignored
-- But the increment message with amount=2 is delivered
Context Propagation
Pass context values that are accessible inside the workflow and its activities:
local spawner = process.with_context({
user_id = "user-1",
tenant = "tenant-1",
request_id = "req-abc",
})
local pid, err = spawner:spawn_monitored(
"app:order_workflow",
"app:worker",
order
)
if err then return nil, err end
Inside the workflow (or any activity it calls), read context via the ctx module:
local ctx = require("ctx")
local user_id, user_err = ctx.get("user_id") -- "user-1"
if user_err then return nil, user_err end
local tenant, tenant_err = ctx.get("tenant") -- "tenant-1"
if tenant_err then return nil, tenant_err end
local all, err = ctx.all() -- {user_id="user-1", tenant="tenant-1", request_id="req-abc"}
if err then
return nil, err
end
Security Context
The actor and scope of the caller travel with the workflow, separately from ctx values and under stronger rules. They are carried in two Temporal headers:
| Header | Content |
|---|---|
wippy-security |
JSON envelope: actor ID, actor metadata, policy IDs, and the audience |
wippy-security-signature |
HMAC-SHA256 over that envelope, keyed by the client's security_hmac_key |
The audience is the ID of the execution the header was minted for — the workflow ID for a start or a signal, the activity ID for an activity. A header replayed against a different execution fails the audience check, so a captured header cannot be reused elsewhere.
Verification happens before the workflow body runs. The signature must match one of the client's keys, the audience must equal this execution's ID, and every policy named in the envelope must resolve in the local security registry. Any of those failing fails the workflow execution — it is not a warning and the workflow does not run with a reduced context. The same is true of an envelope that is internally inconsistent, such as an actor without a scope or policies without an actor.
Configure the keys on the temporal.client entry. Starting a workflow from a context that has an actor or a scope requires a signing key; without one the start fails rather than proceeding unsigned.
Secured workflows reject unsigned signals
A workflow running under a security context requires every incoming signal to carry a signed relay ticket — headers wippy-relay-signal and wippy-relay-signal-signature — bound to that workflow ID and that signal name. An unsigned or mis-addressed signal is rejected instead of delivered. Signals sent by Wippy processes through process.send are signed automatically. Signals injected from outside Wippy — the Temporal CLI, tctl, or another SDK — carry no ticket and therefore fail against a secured workflow. Drive a secured workflow only from Wippy.
Deterministic child and activity IDs
Under a security context, a child workflow or activity started without an explicit ID gets a derived one instead of a random one, because the ID is the audience the header is signed for and must be reproducible on replay:
| Started from a secured workflow | Generated ID |
|---|---|
| Child workflow | <parentWorkflowID>-<parentRunID>-child-<N> |
| Activity | <parentWorkflowID>-<parentRunID>-activity-<N> |
N counts within the workflow execution. An explicitly supplied temporal.workflow.id or activity ID is used as-is and becomes the audience. Without a security context, IDs are left to Temporal as before.
From HTTP Handlers
local function handler()
local req, req_err = http.request()
if req_err then
return nil, req_err
end
local body, body_err = req:body()
if body_err then
return nil, body_err
end
local order, decode_err = json.decode(body)
if decode_err then
return nil, decode_err
end
local request_id, header_err = req:header("X-Request-ID")
if header_err then
return nil, header_err
end
local spawner = process
.with_context({request_id = request_id})
:with_options({
["workflow.id"] = "order-" .. order.id,
["workflow.id_conflict_policy"] = "fail",
})
local pid, err = spawner:spawn(
"app:order_workflow",
"app:worker",
order
)
local res, res_err = http.response()
if res_err then
return nil, res_err
end
if err then
local status_err = res:set_status(409)
if status_err then
return nil, status_err
end
local write_err = res:write_json({error = tostring(err)})
if write_err then return nil, write_err end
return true
end
local status_err = res:set_status(202)
if status_err then
return nil, status_err
end
local write_err = res:write_json({
workflow_id = tostring(pid),
status = "started"
})
if write_err then return nil, write_err end
return true
end
Signals
Workflows receive signals via the process messaging system. Signals are durable - they survive workflow replays.
Inbox Pattern
Receive all messages through the process inbox:
local function main(order)
local inbox = process.inbox()
while true do
local msg, open = inbox:receive()
if not open then
return nil, errors.new({kind = errors.INTERNAL, message = "workflow inbox closed"})
end
local topic = msg:topic()
if topic == "approve" then
break
elseif topic == "cancel" then
local payload = msg:payload()
local data
if payload then
local payload_err
data, payload_err = payload:data()
if payload_err then return nil, payload_err end
end
local reason = type(data) == "table" and data.reason or nil
return {status = "cancelled", reason = reason}
end
end
return process_order(order)
end
Topic-Based Subscription
Subscribe to specific topics using process.listen():
local function main(input)
local results = {}
local job_ch, job_err = process.listen("add_job")
if job_err then return nil, job_err end
local exit_ch, exit_err = process.listen("exit")
if exit_err then return nil, exit_err end
while true do
local result = channel.select{
job_ch:case_receive(),
exit_ch:case_receive()
}
if result.channel == exit_ch then
break
elseif result.channel == job_ch then
if not result.ok then
break
end
local job_data = result.value
local activity_result, err = funcs.call(
"app:echo_activity",
{job_id = job_data.id, data = job_data}
)
if err then
return nil, err
end
table.insert(results, {
job_id = job_data.id,
result = activity_result
})
end
end
return {total_jobs = #results, results = results}
end
By default, process.listen() returns raw payload data. Use {message = true} to receive Message objects with sender information:
local ch, err = process.listen("request", {message = true})
if err then return nil, err end
local msg, open = ch:receive()
if not open then
return nil, errors.new({kind = errors.INTERNAL, message = "request channel closed"})
end
local sender = msg:from()
local payload = msg:payload()
local data
if payload then
local payload_err
data, payload_err = payload:data()
if payload_err then return nil, payload_err end
end
Serialized Signal Handling
Use one channel.select() loop when signals mutate shared workflow state. This preserves deterministic mutation order and lets the finish branch return without leaving blocked handler coroutines:
local function main(input)
local counter = input.initial or 0
local function send_reply(pid, topic, payload)
local sent, err = process.send(pid, topic, payload)
if err then error(err) end
return sent
end
local function message_data(msg)
local payload = msg:payload()
if not payload then return nil end
return payload:data()
end
local increment_ch, increment_err = process.listen("increment", {message = true})
if increment_err then return nil, increment_err end
local decrement_ch, decrement_err = process.listen("decrement", {message = true})
if decrement_err then return nil, decrement_err end
local finish_ch, finish_err = process.listen("finish", {message = true})
if finish_err then return nil, finish_err end
while true do
local result = channel.select{
increment_ch:case_receive(),
decrement_ch:case_receive(),
finish_ch:case_receive()
}
if not result.ok then
return nil, errors.new({kind = errors.INTERNAL, message = "signal channel closed"})
end
local msg = result.value
local reply_to = msg:from()
if result.channel == finish_ch then
send_reply(reply_to, "ack")
send_reply(reply_to, "ok", {message = "finishing", value = counter})
return {final_counter = counter}
end
local data, payload_err = message_data(msg)
if payload_err then return nil, payload_err end
if type(data) ~= "table" or type(data.amount) ~= "number" then
send_reply(reply_to, "nak", "amount must be a number")
elseif result.channel == decrement_ch and counter - data.amount < 0 then
send_reply(reply_to, "nak", "would result in negative value")
else
send_reply(reply_to, "ack")
if result.channel == increment_ch then
counter = counter + data.amount
else
counter = counter - data.amount
end
send_reply(reply_to, "ok", {value = counter})
end
end
end
Signal Acknowledgment
Implement request-reply patterns by sending responses back to the sender:
-- Workflow side
local ch, err = process.listen("get_status", {message = true})
if err then return nil, err end
local msg, open = ch:receive()
if not open then return nil, errors.new({kind = errors.INTERNAL, message = "status channel closed"}) end
local sent, send_err = process.send(msg:from(), "status_response", {status = "processing", progress = 75})
if send_err then return nil, send_err end
-- Caller side
local response_ch, listen_err = process.listen("status_response")
if listen_err then return nil, listen_err end
local sent, send_err = process.send(workflow_pid, "get_status", {})
if send_err then return nil, send_err end
local timeout, timeout_err = time.after("5s")
if timeout_err then return nil, timeout_err end
local result = channel.select{
response_ch:case_receive(),
timeout:case_receive()
}
if result.channel == response_ch then
if not result.ok then
return nil, errors.new({kind = errors.INTERNAL, message = "status response channel closed"})
end
return result.value
end
if not result.ok then
return nil, errors.new({kind = errors.INTERNAL, message = "status timeout channel closed"})
end
return nil, errors.new({kind = errors.TIMEOUT, message = "status request timed out", retryable = true})
Cross-Workflow Signaling
Workflows can send signals to other workflows using their PID:
-- Sender workflow
local function main(input)
local target_pid = input.target
local response_ch, listen_err = process.listen("cross_host_pong")
if listen_err then return nil, listen_err end
local ok, err = process.send(target_pid, "cross_host_ping", {data = "hello"})
if err then
return {ok = false, error = tostring(err)}
end
local response, open = response_ch:receive()
if not open then
return {ok = false, error = "cross_host_pong channel closed"}
end
return {ok = true, received = response}
end
Child Workflows
Synchronous Child (workflow.exec)
Execute a child workflow and wait for the result:
local result, err = workflow.exec("app:child_workflow", input_data)
if err then
return nil, err
end
Asynchronous Child (process.spawn)
Spawn a child workflow without blocking, then wait for its completion via events:
local events_ch = process.events()
local child_pid, err = process.spawn(
"app:child_workflow",
"app:worker",
{message = "hello from parent"}
)
if err then
return {status = "spawn_failed", error = tostring(err)}
end
-- Wait for child EXIT event
local event, open = events_ch:receive()
if not open then
return nil, errors.new({kind = errors.INTERNAL, message = "process event channel closed"})
end
if event.kind == process.event.EXIT then
local child_result = event.result.value
local child_error = event.result.error
end
Error Propagation from Children
When a child workflow returns an error, it appears in the EXIT event:
local events_ch = process.events()
local child_pid, err = process.spawn(
"app:error_child_workflow",
"app:worker"
)
if err then
return nil, err
end
local event, open = events_ch:receive()
if not open then
return nil, errors.new({kind = errors.INTERNAL, message = "process event channel closed"})
end
if event.result.error then
local child_err = event.result.error
-- Error objects have kind(), retryable(), message() methods
print(child_err:kind()) -- e.g. "NotFound"
print(child_err:retryable()) -- false
print(child_err:message()) -- error message text
end
Executing Workflows Synchronously (process.exec)
Run a workflow and wait for its result in one call:
local result, err = process.exec(
"app:hello_workflow",
"app:worker",
{name = "world"}
)
if err then
return nil, err
end
-- result contains the workflow return value
Monitoring and Linking
Post-Start Monitoring
Monitor a workflow after it has already started:
local pid, err = process.spawn(
"app:long_workflow",
"app:worker",
{iterations = 100}
)
if err then
return nil, err
end
-- Monitor later
local ok, monitor_err = process.monitor(pid)
if monitor_err then
return nil, monitor_err
end
local events_ch = process.events()
local event, open = events_ch:receive() -- EXIT when workflow completes
if not open then
return nil, errors.new({kind = errors.INTERNAL, message = "process event channel closed"})
end
Post-Start Linking
Link to a running workflow to receive LINK_DOWN on abnormal termination:
local ok, err = process.set_options({trap_links = true})
if err then
return nil, err
end
local pid, err = process.spawn(
"app:long_workflow",
"app:worker",
{iterations = 100}
)
if err then
return nil, err
end
-- Link after workflow has started
time.sleep("200ms")
local linked, link_err = process.link(pid)
if link_err then return nil, link_err end
-- If workflow is terminated, receive LINK_DOWN
local terminated, terminate_err = process.terminate(pid)
if terminate_err then return nil, terminate_err end
local events_ch = process.events()
local event, open = events_ch:receive()
if not open then
return nil, errors.new({kind = errors.INTERNAL, message = "process event channel closed"})
end
-- event.kind == process.event.LINK_DOWN
LINK_DOWN events require trap_links = true in process options. Without it, a linked process termination propagates the failure.
Unmonitor / Unlink
Remove monitoring or linking:
local unmonitored, unmonitor_err = process.unmonitor(pid)
if unmonitor_err then return nil, unmonitor_err end
local unlinked, unlink_err = process.unlink(pid)
if unlink_err then return nil, unlink_err end
After unmonitoring or unlinking, events for that process are no longer delivered.
Termination and Cancellation
Terminate
Force-terminate a running workflow:
local ok, err = process.terminate(workflow_pid)
Monitored callers receive an EXIT event with an error.
Cancel
Request graceful cancellation with an optional reason:
local ok, err = process.cancel(workflow_pid, "cancelled by operator")
Concurrent Work
Use coroutine.spawn() and channels for parallel work inside workflows:
local function main(input)
local worker_count = input.workers or 3
local job_count = input.jobs or 6
local work_queue = channel.new(10)
local results = channel.new(10)
for w = 1, worker_count do
coroutine.spawn(function()
while true do
local job, ok = work_queue:receive()
if not ok then break end
time.sleep(10 * time.MILLISECOND)
results:send({worker = w, job = job, result = job * 2})
end
end)
end
for j = 1, job_count do
work_queue:send(j)
end
work_queue:close()
local total = 0
local processed = {}
for _ = 1, job_count do
local r, open = results:receive()
if not open then
return nil, errors.new({kind = errors.INTERNAL, message = "results channel closed"})
end
total = total + r.result
table.insert(processed, r)
end
return {total = total, processed = processed}
end
All channel operations and sleeps inside coroutines are replay-safe.
Timers
Durable timers survive restarts:
local time = require("time")
time.sleep("24h")
time.sleep("5m")
time.sleep("30s")
time.sleep(100 * time.MILLISECOND)
Track elapsed time:
local start = time.now()
time.sleep("1s")
local elapsed = time.now():sub(start):milliseconds()
Determinism
Workflow code must be deterministic. The same inputs must produce the same sequence of commands.
Replay-Safe Operations
These operations are automatically intercepted and their results recorded. On replay, recorded values are returned:
-- Activity calls
local data = funcs.call("app:fetch_data", id)
-- Durable sleep
time.sleep("1h")
-- Current time
local now = time.now()
-- UUID generation
local id = uuid.v4()
-- Crypto operations
local bytes = crypto.random.bytes(32)
-- Child workflows
local result = workflow.exec("app:child", input)
-- Versioning
local v = workflow.version("change-1", 1, 2)
Non-Deterministic (Avoid)
-- Don't use wall clock time
local now = os.time() -- non-deterministic
-- Don't use random directly
local r = math.random() -- non-deterministic
-- Don't do I/O in workflow code
local file = io.open("data.txt") -- non-deterministic
-- Don't use global mutable state
counter = counter + 1 -- non-deterministic across replays
Error Handling
Activity Errors
Activity errors carry structured metadata:
local result, err = funcs.call("app:risky_activity", order)
if err then
print(err:kind()) -- error classification (e.g. "NotFound", "Internal")
print(err:retryable()) -- whether the error is retryable
print(err:message()) -- human-readable error message
end
Activity Failure Modes
Configure retry behavior for activity calls:
local executor = funcs.new():with_options({
["activity.retry_policy"] = {
maximum_attempts = 1,
}
})
local result, err = executor:call("app:unreliable_activity", input)
if err then
local kind = err:kind() -- "Internal" for runtime errors
local retryable = err:retryable()
end
Child Workflow Errors
Errors from child workflows (via process.exec or EXIT events) carry the same metadata:
local result, err = process.exec("app:error_workflow", "app:worker")
if err then
print(err:kind()) -- e.g. "NotFound"
print(err:retryable()) -- false
print(err:message()) -- error details
end
Compensation Pattern (Saga)
local function run_compensations(compensations)
local first_err
for _, comp in ipairs(compensations) do
local _, err = funcs.call(comp.action, comp.args)
if err and not first_err then
first_err = err
end
end
if first_err then return nil, first_err end
return true
end
local function main(order)
local compensations = {}
local reservation, err = funcs.call("app:reserve_inventory", order.items)
if err then
return {status = "failed", step = "inventory", error = tostring(err)}
end
table.insert(compensations, 1, {
action = "app:release_inventory",
args = reservation.id
})
local payment, err = funcs.call("app:charge_payment", order.payment)
if err then
local _, compensation_err = run_compensations(compensations)
if compensation_err then
return {status = "failed", step = "payment", error = tostring(err), compensation_error = tostring(compensation_err)}
end
return {status = "failed", step = "payment", error = tostring(err)}
end
table.insert(compensations, 1, {
action = "app:refund_payment",
args = payment.id
})
local shipment, err = funcs.call("app:ship_order", order.shipping)
if err then
local _, compensation_err = run_compensations(compensations)
if compensation_err then
return {status = "failed", step = "shipping", error = tostring(err), compensation_error = tostring(compensation_err)}
end
return {status = "failed", step = "shipping", error = tostring(err)}
end
return {status = "completed", tracking = shipment.tracking}
end
Compensations run in reverse registration order. If more than one compensation
fails, the workflow still attempts the remaining actions and reports the first
failure through compensation_error.
See Also
- Overview - Client and worker configuration
- Activities - Activity definitions and options
- Process - Process management API
- Functions - Function invocation
- Channels - Channel operations