Temporalex

Workflow orchestration for Elixir, built on the Temporal Core SDK (Rust) over Rustler NIFs.

Temporalex workflows read top-to-bottom as sequential code. Concurrency is explicit and structured — there is no implicit event loop. The runtime uses a deterministic cooperative scheduler that owns thread ordering, so command sequences are reproducible from the same activation transcript regardless of BEAM scheduling or mailbox timing.

Status: 0.5.x. This line is an architectural rewrite around a deterministic core, a Temporalex.Backend boundary that isolates Temporal Core / Rust details, and structured concurrency primitives phase and parallel. The 0.x line is not backwards-compatible with 0.3.0. See CHANGELOG.md for the migration notes.

** NEXT and FINAL STAGE Before ALPHA TESTING ** Fix the usability of the SDK to make it easier and simpler than the vibe coded version Retest and verify and finish the liveview demo app

Core design and scheduler authored by @hansihe.


Install

# mix.exs
defp deps do
[{:temporalex, "~> 0.4.0"}]
end

Requirements: Elixir ~> 1.17, Rust toolchain (the NIF crate compiles on first build against temporalio/sdk-rust v0.4.0).

Run a Temporal dev server

brew install temporal
temporal server start-dev

The Web UI lands at http://localhost:8233; the gRPC endpoint at localhost:7233.


Define and call a workflow

A workflow module declares what it is — behaviour, identity, address — and use generates the call-side surface:

defmodule Greet do
use Temporalex.Workflow, queue: "greetings"
@impl true
def id(name), do: "greet-\#{name}"
@impl true
def run(name), do: {:ok, "Hello, \#{name}!"}
end
greeting = Greet.execute!("Fresha") # start, wait, get the answer
handle = Greet.start!("Fresha") # start and move on
greeting = Temporalex.await!(handle) # collect later
booking_id # policy, when a call carries it
|> Booking.new()
|> Temporalex.retry(max_attempts: 3)
|> Temporalex.fairness(salon_id)
|> Temporalex.execute!()

id/1 derives the workflow id — Temporal's idempotency key — so a duplicate start attaches to the running execution and a webhook can signal it knowing only the business key: Booking.signal!(booking_id, "confirmed"). Design and rationale: docs/rfcs/0002-client-surface.md.

Define an activity

defmodule MyApp.Activities.Payment do
use Temporalex.Activity
defactivity charge(amount), start_to_close_timeout: 30_000 do
{:ok, "charge-#{amount}"}
end
# Local activity: runs in-process on the same worker, durable via a
# history marker. Use for short, deterministic work where the network
# round-trip to schedule a regular activity isn't worth it.
defactivity stamp(prefix), local: true, start_to_close_timeout: 5_000 do
{:ok, "#{prefix}-#{System.unique_integer([:positive])}"}
end
end

Structured errors

Activities can raise Temporalex.ApplicationError (or return {:error, reason}); the workflow sees a typed exception with the cause preserved:

defactivity charge(amount) do
if amount > 10_000 do
raise %Temporalex.ApplicationError{
message: "amount exceeds limit",
type: "AmountTooLarge",
non_retryable: true
}
else
{:ok, amount}
end
end
# In the workflow:
case Activities.charge(amount) do
{:ok, charge} -> ...
{:error, %Temporalex.ActivityFailure{cause: %{type: "AmountTooLarge"}}} -> ...
end

Define a workflow

defmodule MyApp.Workflows.Checkout do
use Temporalex.Workflow
alias Temporalex.Workflow.API
def handle_query("status", _args, state), do: {:reply, state}
def run(args) do
API.publish_state(:charging)
{:ok, charge} = MyApp.Activities.Payment.charge(args["amount"])
API.publish_state(:awaiting_confirmation)
confirmed =
API.phase(false,
signal: %{
"confirm" => fn _args, _state -> {:stop, true} end,
"cancel" => fn _args, _state -> {:stop, false} end
},
timeout: :timer.minutes(5)
)
case confirmed do
{:timeout, _state} -> {:error, :timed_out}
true -> {:ok, %{charge: charge, confirmed: true}}
false -> {:error, :user_cancelled}
end
end
end

Clients and workers

The worker is not the connection — that's the client. They are different processes with different jobs:

ClientWorker
Isthe gRPC connection to the servera poller + executor bound to one task queue
Knowstarget, namespace, codecwhich workflows/activities it can run
Needed byanyone who starts/signals/queriesonly nodes that run workflow code

One client can back many workers. And a node that only starts workflows — a Phoenix deployment creating bookings, say — needs just a client and no worker at all: it never polls.

The worker entry in your supervision tree is the deployment topology written down: "this node serves these workflows." That is a property of the deployable, not of the workflow code, which is why it cannot be derived and must be stated.

Task queues

A task queue is a rendezvous string — nothing more. Starting a workflow writes its tasks under a name; workers long-poll a name; whoever polls the name you wrote to gets the work. There is no server-side registry of which worker runs which workflow type, so every start names its queue.

The queue is the unit of…
decouplingcallers name a queue, never a host, pod, or process — whoever polls that name picks the work up
scalingmore capacity = more workers polling the same name
deploymentone queue ≈ one deployable, with its own release cadence and blast radius
fairness & versioningfairness keys and worker deployment versions attach per queue

Queues are not provisioned; a queue springs into existence the first time anyone uses its name. Which is what makes a typo dangerous: it does not error, it creates a new empty queue — and a workflow started there sits "Running" forever, because nothing polls it.

Start a client and a worker

children = [
{Temporalex.Client,
name: MyApp.Temporal,
backend: Temporalex.Backend.TemporalCore,
target: "http://127.0.0.1:7233",
namespace: "default",
# Starts that pass no :task_queue inherit this one. Without it they go to
# "default" — which nothing here polls (see "Task queues" above).
task_queue: "checkout",
# :etf (default) preserves full Elixir term fidelity.
# :json makes payloads renderable by `temporal` CLI and non-Elixir
# clients, at the cost of lossy term encoding (atoms → strings,
# tuples → unsupported).
payload_codec: :etf},
{Temporalex.Worker,
name: MyApp.Worker,
client: MyApp.Temporal,
task_queue: "checkout",
workflows: [MyApp.Workflows.Checkout],
activities: [MyApp.Activities.Payment]}
]
Supervisor.start_link(children, strategy: :one_for_one)

Metrics

Temporal's core records worker-side metrics — poller counts, slot usage, and schedule_to_start latency, which is the signal you autoscale workers on. They are off by default. The exporter belongs to the runtime, so it is configured on the client:

{Temporalex.Client,
name: MyApp.Temporal,
backend: Temporalex.Backend.TemporalCore,
target: "http://127.0.0.1:7233",
namespace: "default",
telemetry: [
prometheus: [bind_address: "0.0.0.0:9464"],
global_tags: %{"service" => "checkout-worker", "env" => "prod"}
]}

GET /metrics on that address then serves, among others:

temporal_num_pollers
temporal_worker_task_slots_available
temporal_worker_task_slots_used
temporal_workflow_task_schedule_to_start_latency
temporal_workflow_task_execution_latency
temporal_workflow_endtoend_latency
temporal_workflow_completed

For OTLP instead — to an OpenTelemetry Collector, or any agent with an OTLP intake:

telemetry: [
otlp: [
url: "http://localhost:4317",
protocol: :grpc, # or :http
metric_temporality: :cumulative, # or :delta
metric_periodicity_ms: 1_000,
headers: %{"authorization" => "Bearer ..."}
]
]

:prometheus and :otlp are mutually exclusive on one runtime. Other keys: metric_prefix (default "temporal_"), attach_service_name, and durations_as_seconds.

Build id

Workers report a build id, which lands on every WorkflowTaskCompleted event and so answers "which release executed this task?" from history or the Web UI:

{Temporalex.Worker,
name: MyApp.Worker,
client: MyApp.Temporal,
task_queue: "checkout",
build_id: System.get_env("BUILD_ID", "dev"),
workflows: [MyApp.Workflows.Checkout]}

Without a :versioning option the strategy is None, so the build id identifies but does not affect task routing. Defaults to "temporalex-<version>". Pass versioning: [deployment_name: ..., use_versioning: true, default_behavior: :pinned | :auto_upgrade] to opt a worker into deployment-based routing.

Drive workflows from a client

{:ok, handle} =
Temporalex.Client.start_workflow(
MyApp.Temporal,
MyApp.Workflows.Checkout,
%{"amount" => 100},
workflow_id: "checkout-#{order_id}"
)
:ok = Temporalex.Client.signal_workflow(handle, "confirm")
{:ok, status} = Temporalex.Client.query_workflow(handle, "status")
{:ok, result} = Temporalex.Client.get_result(handle)

The full client surface: start_workflow, get_result, signal_workflow, query_workflow, update_workflow, cancel_workflow, terminate_workflow, describe_workflow.


Programming model

Workflows are a single run/1 function. Concurrency enters only through phase and parallel, which act as structured concurrency scopes — every async handler spawned within a scope must complete before the scope returns.

PrimitivePurpose
Activities.Module.fun(args)Execute an activity. Blocks until resolved.
API.sleep(ms)Durable timer.
API.wait_for_signal(name)Pop one signal from the buffer.
API.publish_state(state)Update the snapshot that queries see.
API.now/0API.random/0API.uuid4/0Deterministic time/random.
API.patched?(id)Workflow versioning, replay-safe.
API.phase(state, opts)Message-processing scope with signal/update handlers and an optional :timeout.
API.parallel(fns)Cooperatively scheduled fan-out. Results in input order.
API.update_state(fn)Atomically transform the enclosing phase's state from inside an {:async, fn, _} handler.
API.execute_child_workflow(mod, input, opts)Start a child workflow, block until it completes.
API.start_child_workflow(mod, input, opts)Start a child non-blocking; returns a ChildHandle.
API.await_child_workflow(handle)Block until a started child completes.
API.signal_child_workflow(handle_or_id, name, args)Send a durable signal to a child workflow.
API.cancel_child_workflow(handle_or_id)Request cancellation of a child workflow.

Full details, return-value contracts, and the determinism rationale:


Testing

Temporalex.Backend.Test is an in-memory backend that lets you drive a worker with core activation structs directly — no Temporal server required. The same Temporalex.Server and Temporalex.Core.Executor that handle real traffic also handle the test backend, so workflow code under test runs the production codepath.

start_supervised!(
{Temporalex.Client, name: MyApp.TestClient, backend: Temporalex.Backend.Test}
)
start_supervised!(
{Temporalex.Worker,
name: MyApp.TestWorker,
client: MyApp.TestClient,
workflows: [MyApp.Workflows.Checkout],
activities: [MyApp.Activities.Payment]}
)

See test/temporalex/server_integration_test.exs for full activation and activity-task transcripts.


Project layout

lib/temporalex/
workflow.ex use Temporalex.Workflow
workflow/api.ex sequential primitives, phase, parallel
activity.ex defactivity macro
activity/context.ex heartbeat, cancelled? for activity bodies
client.ex start / get_result / signal / query / update / cancel / terminate / describe
worker.ex Supervisor — what users add to their tree
server.ex Worker server: backend state, executor registry, activation routing
core/executor.ex deterministic workflow executor (scheduler + replay)
core/structs.ex internal protocol: Activation, Job, Command, Completion, Op
core/test_harness.ex in-process harness for testing the core directly
backend.ex Backend behaviour
backend/test.ex in-memory backend for tests
backend/temporal_core.ex Rustler-backed Temporal Core backend
native.ex Rustler NIF surface (do not call directly)
native/temporalex_nif/
src/ Rust NIF crate

Contributing

The architecture is documented in docs/. Start with docs/sdk_overview.md. The docs/implementation_principles.md admission rule applies to any new workflow API: a primitive only enters the public surface if it has a precise replay contract and can be tested without the real Temporal backend.

Run the quality gates before committing — CI enforces all four:

mix format
mix test
mix credo --strict
mix dialyzer

Credo is configured in .credo.exs. The few relaxations (Temporal-style exception names, NIF stub arities) are commented in the config; individually exempted functions carry credo:disable-for-next-line comments at the site explaining why.

License

MIT — see LICENSE.