Workflows & activities
Workflows orchestrate; activities do. The workflow method describes the durable control flow, while every side effect (HTTP calls, emails, database writes) belongs in an activity so it can be independently leased and its outcome recorded for replay.
Defining a workflow
Section titled “Defining a workflow”[Workflow("greeting")]public sealed class GreetingWorkflow{ [WorkflowRun] public async ValueTask<string> RunAsync( IWorkflowContext ctx, string name, CancellationToken cancellationToken ) { // ... }}[Workflow("greeting")] sets the workflow’s durable identity: the name stored with
every run in Postgres. It defaults to the class name, but pin it explicitly for anything
long-lived, because renaming a class without a pinned name orphans its in-flight runs.
[WorkflowRun] marks the single public entry-point method. Its signature is
(IWorkflowContext, input, CancellationToken): exactly one context, exactly one input
parameter, and an optional trailing cancellation token. The method can return void,
Task, ValueTask, or their generic forms.
Inputs and outputs are stored as JSON in Postgres, so any JSON-serializable type works. Records are a great fit:
public sealed record SignupInput(string Company, string Email);Defining activities
Section titled “Defining activities”public sealed class EmailActivities{ [Activity("send-welcome")] public string SendWelcome(string email) { // side effects go here }}An activity class is a normal class resolved from your DI container, so constructor
inject whatever the side effect needs, like an HttpClient or a DbContext. Methods
can be sync, Task, or ValueTask, take multiple parameters, and accept an optional
trailing CancellationToken that is cancelled if the activity loses its lease. Like
workflows, [Activity] names default to the method name; pin them for anything
long-lived.
Calling activities from a workflow
Section titled “Calling activities from a workflow”var receipt = await ctx.Activity( (PaymentActivities a) => a.Charge(input.UserName, input.Amount), cancellationToken);The lambda must be a direct method call on the activity class; that expression is how PgWorkflows knows which activity to enqueue and with which arguments. The arguments are evaluated and serialized at the call site, the activity runs on whichever worker leases the job, and the workflow parks until the result lands.
ctx.Activity awaits one activity. To run several concurrently, create pending
activities with ctx.CallActivity and await them together with ctx.WhenAll; see
fan-in fan-out.
Execution and idempotency
Section titled “Execution and idempotency”A worker executes an activity under an expiring lease and records its outcome in Postgres. Once recorded, workflow replay returns that stored outcome. Activity execution is nevertheless at least once, not exactly once: a worker can complete an external side effect and die before it records success, after which another worker reclaims and executes the job.
Design activity side effects to be idempotent, normally by including a stable business key in their
input. Delegate activities registered through the low-level RegisterActivity API can accept an
ActivityExecutionContext; its JobId is stable across attempts and is suitable when the external
system supports idempotency keys. Cancellation is cooperative, so a handler that loses its lease
may continue running; its stale database outcome is rejected, but PgWorkflows cannot undo or fence
an external call.
Deterministic workflow code
Section titled “Deterministic workflow code”Workflow methods replay from the beginning. Conditions and loops are safe when they depend on the
workflow input or recorded activity results. Current time, random values, mutable process state,
configuration changes, and direct I/O must not change the order of ctx.* calls. Put
nondeterministic decisions behind activities so their results are persisted.
Durable calls are matched by sequence within their own kind (activities, timers, signal waits, and failure hooks each have a counter). Incompatible changes that add, remove, or reorder calls require a new durable workflow name. Keep the previous workflow and activity contracts registered until its in-flight runs finish.
Registration
Section titled “Registration”builder.Services.AddPgWorkflows(pg => pg.UsePostgres(connectionString) .AddWorkflow<GreetingWorkflow>() .AddActivities<EmailActivities>());AddWorkflow<T>() registers one workflow class; AddActivities<T>() registers every
[Activity] method on a class. Workflow and activity names must be unique across the
registration, and duplicates fail at startup.