Skip to main content

Interface: WorkflowPersistence

Defined in: packages/workflow-engine/src/persistence/interface.ts:972

Full persistence contract: PersistenceCore (what the kernel calls) plus a transaction boundary that hands the callback this same surface. The pre-1.0 artifact methods (saveArtifact, loadArtifact, ... — replaced by the BlobStore port) and query helpers (getRunsByStatus, claimPendingRun, getStageById, ...) are no longer part of the contract; the built-in adapters still implement them as plain class methods.

Extends​

Methods​

acquireIdempotencyKey()​

acquireIdempotencyKey(key, commandType, options?): Promise<{ status: "acquired"; } | { result: unknown; status: "replay"; } | { status: "in_progress"; }>

Defined in: packages/workflow-engine/src/persistence/interface.ts:938

Atomically acquire an idempotency key for command execution.

If the key is currently in_progress (e.g. a previous dispatcher crashed between committing its transaction and calling completeIdempotencyKey), passing staleInProgressAfterMs allows the key to be reclaimed once it has been in progress for at least that long, measured against options.now (defaults to new Date()). Reclaiming is atomic: only one caller wins when multiple dispatchers race to reclaim the same stale key. When staleInProgressAfterMs is omitted, a stuck in_progress key is never reclaimed (matches prior behavior).

Parameters​

key​

string

commandType​

string

options?​
now?​

Date

staleInProgressAfterMs?​

number

Returns​

Promise<{ status: "acquired"; } | { result: unknown; status: "replay"; } | { status: "in_progress"; }>

Inherited from​

PersistenceCore.acquireIdempotencyKey


appendAnnotations()​

appendAnnotations(inputs): Promise<void>

Defined in: packages/workflow-engine/src/persistence/interface.ts:878

Append one or more annotations. Designed to be called both standalone (fire-and-forget from external attach) and inside an existing transaction (buffered during stage execution, flushed in the stage-completion transaction) -- called from run.create, external attach, and the stage-completion transactions in job-execute and stage-poll-suspended.

Rows with the same (workflowRunId, key, idempotencyKey) are deduped via the unique constraint; duplicates are silently skipped.

Parameters​

inputs​

CreateAnnotationInput[]

Returns​

Promise<void>

Inherited from​

PersistenceCore.appendAnnotations


appendOutboxEvents()​

appendOutboxEvents(events): Promise<void>

Defined in: packages/workflow-engine/src/persistence/interface.ts:901

Write events to the outbox. Sequences are auto-assigned per workflowRunId.

Parameters​

events​

CreateOutboxEventInput[]

Returns​

Promise<void>

Inherited from​

PersistenceCore.appendOutboxEvents


claimNextPendingRun()​

claimNextPendingRun(options?): Promise<WorkflowRunRecord | null>

Defined in: packages/workflow-engine/src/persistence/interface.ts:747

Atomically find and claim the next pending workflow run. Uses FOR UPDATE SKIP LOCKED pattern (in Postgres) to prevent race conditions when multiple workers try to claim workflows simultaneously.

Priority ordering: higher priority first, then oldest (FIFO within same priority).

Parameters​

options?​
now?​

Date

The kernel clock's time, written as startedAt/updatedAt.

serves?​

readonly ServedDefinition[]

The definitions the claiming host is built to serve. When supplied, a run is claimable only if it is pinned to one of these (workflowId, version) pairs, or is unpinned (definitionVersion is null, i.e. it predates the consumer's migration) and one of the pairs names its workflow. An empty array claims nothing: a host that serves no workflow has no work.

The workflow-id restriction on the unpinned arm matters in a fleet: without it every host claims the whole pre-migration population, and one whose registry lacks the workflow adopts a run only to fail it with WORKFLOW_NOT_FOUND. Deciding it in the query is what keeps a host from taking work it cannot do.

Omit it to claim any pending run, which is the pre-1.0 behaviour.

Returns​

Promise<WorkflowRunRecord | null>

The claimed workflow run (now with status RUNNING), or null if no pending runs

Inherited from​

PersistenceCore.claimNextPendingRun


claimUnpublishedOutboxEvents()​

claimUnpublishedOutboxEvents(limit?): Promise<OutboxRecord[]>

Defined in: packages/workflow-engine/src/persistence/interface.ts:916

Atomically claim up to limit unpublished events for this caller by setting publishedAt, ordered by (workflowRunId, sequence). Two processes flushing the same outbox concurrently must never both receive the same event: on Postgres this is one UPDATE ... FROM (SELECT ... FOR UPDATE SKIP LOCKED) RETURNING; other stores use a compare-and-set on publishedAt IS NULL per row. Events whose publication then fails are handed back with releaseOutboxEvents.

Parameters​

limit?​

number

Returns​

Promise<OutboxRecord[]>

Inherited from​

PersistenceCore.claimUnpublishedOutboxEvents


completeIdempotencyKey()​

completeIdempotencyKey(key, commandType, result): Promise<void>

Defined in: packages/workflow-engine/src/persistence/interface.ts:949

Mark an idempotency key as completed and cache the command result.

Parameters​

key​

string

commandType​

string

result​

unknown

Returns​

Promise<void>

Inherited from​

PersistenceCore.completeIdempotencyKey


countRunsByDefinitionVersion()​

countRunsByDefinitionVersion(filter?): Promise<DefinitionVersionCount[]>

Defined in: packages/workflow-engine/src/persistence/interface.ts:821

Run counts grouped by (workflowId, definitionVersion, status) — the query behind "has this version drained?". Returns an empty array when PersistenceCore.supportsDefinitionVersioning is false.

Parameters​

filter?​

DefinitionVersionCountFilter

Returns​

Promise<DefinitionVersionCount[]>

Inherited from​

PersistenceCore.countRunsByDefinitionVersion


createLog()​

createLog(data): Promise<void>

Defined in: packages/workflow-engine/src/persistence/interface.ts:864

Parameters​

data​

CreateLogInput

Returns​

Promise<void>

Inherited from​

PersistenceCore.createLog


createRun()​

createRun(data): Promise<WorkflowRunRecord>

Defined in: packages/workflow-engine/src/persistence/interface.ts:702

Parameters​

data​

CreateRunInput

Returns​

Promise<WorkflowRunRecord>

Inherited from​

PersistenceCore.createRun


createStage()​

createStage(data): Promise<WorkflowStageRecord>

Defined in: packages/workflow-engine/src/persistence/interface.ts:826

Parameters​

data​

CreateStageInput

Returns​

Promise<WorkflowStageRecord>

Inherited from​

PersistenceCore.createStage


deleteRun()​

deleteRun(id): Promise<void>

Defined in: packages/workflow-engine/src/persistence/interface.ts:736

Delete a run and everything the persistence owns under it: stage records, logs, artifacts, annotations, and (on the reference schema, through the cascade) workflow_steps. A missing id is a no-op. Never emits an outbox event; the kernel clears the step ledger and the blob store before calling it.

Parameters​

id​

string

Returns​

Promise<void>

Inherited from​

PersistenceCore.deleteRun


deleteStage()​

deleteStage(id): Promise<void>

Defined in: packages/workflow-engine/src/persistence/interface.ts:861

Parameters​

id​

string

Returns​

Promise<void>

Inherited from​

PersistenceCore.deleteStage


ensureDefinitionVersioningDetected()?​

optional ensureDefinitionVersioningDetected(): Promise<boolean>

Defined in: packages/workflow-engine/src/persistence/interface.ts:797

Optional. Confirms PersistenceCore.supportsDefinitionVersioning against the live database, once, and resolves to the answer every later call will give.

An adapter whose capability answer is derived from something other than the database — the Prisma adapter reads the generated client, which is regenerated before the migration is applied on every first migration and on any rolling deploy that ships code first — implements this so the disagreement is found before it becomes a 42703 on every claim. The kernel awaits it on the paths that would otherwise touch the versioning columns; an adapter that does not need it (the in-memory one, whose answer is its own configuration) omits it.

MUST be safe to call inside a caller's transaction, so it may not issue a statement that can fail — a failed statement aborts the whole transaction, including work the caller did before it.

Returns​

Promise<boolean>

Inherited from​

PersistenceCore.ensureDefinitionVersioningDetected


getDefinition()​

getDefinition(workflowId, version): Promise<WorkflowDefinitionRecord | null>

Defined in: packages/workflow-engine/src/persistence/interface.ts:811

Loads one stored definition snapshot, or null when absent.

Parameters​

workflowId​

string

version​

string

Returns​

Promise<WorkflowDefinitionRecord | null>

Inherited from​

PersistenceCore.getDefinition


getRun()​

getRun(id): Promise<WorkflowRunRecord | null>

Defined in: packages/workflow-engine/src/persistence/interface.ts:711

Parameters​

id​

string

Returns​

Promise<WorkflowRunRecord | null>

Inherited from​

PersistenceCore.getRun


getRunStatus()​

getRunStatus(id): Promise<Status | null>

Defined in: packages/workflow-engine/src/persistence/interface.ts:712

Parameters​

id​

string

Returns​

Promise<Status | null>

Inherited from​

PersistenceCore.getRunStatus


getStage()​

getStage(runId, stageId): Promise<WorkflowStageRecord | null>

Defined in: packages/workflow-engine/src/persistence/interface.ts:829

Parameters​

runId​

string

stageId​

string

Returns​

Promise<WorkflowStageRecord | null>

Inherited from​

PersistenceCore.getStage


getStagesByRun()​

getStagesByRun(runId, options?): Promise<WorkflowStageRecord[]>

Defined in: packages/workflow-engine/src/persistence/interface.ts:835

Ordered by executionGroup (actual execution/dependency order), with stageNumber (definition order) as a tiebreaker for stages that share an execution group (parallel stages).

Parameters​

runId​

string

options?​
orderBy?​

"asc" | "desc"

status?​

Status

Returns​

Promise<WorkflowStageRecord[]>

Inherited from​

PersistenceCore.getStagesByRun


getStuckRuns()​

getStuckRuns(stuckSince): Promise<WorkflowRunRecord[]>

Defined in: packages/workflow-engine/src/persistence/interface.ts:713

Parameters​

stuckSince​

Date

Returns​

Promise<WorkflowRunRecord[]>

Inherited from​

PersistenceCore.getStuckRuns


getSuspendedStages()​

getSuspendedStages(beforeDate, options?): Promise<WorkflowStageRecord[]>

Defined in: packages/workflow-engine/src/persistence/interface.ts:854

Suspended stages whose nextPollAt has come due, oldest deadline first, capped at limit.

Both narrowings belong here rather than in the caller. Without the ordering and the cap the adapter returns every ready row and the handler slices in JS, so a row that keeps coming back — a stage this build cannot serve, released for another host to take — permanently occupies a candidate slot and pushes servable stages out of the window. serves is the same filter claimNextPendingRun applies to runs, evaluated against the stage's run, so an unserving host does not even claim the poll lease of a stage that is not its work. An adapter whose schema has no definitionVersion column ignores serves; such a database has no pinned runs, so ignoring it is exact.

Parameters​

beforeDate​

Date

options?​
limit?​

number

serves?​

readonly ServedDefinition[]

Returns​

Promise<WorkflowStageRecord[]>

Inherited from​

PersistenceCore.getSuspendedStages


getUnpublishedOutboxEvents()​

getUnpublishedOutboxEvents(limit?): Promise<OutboxRecord[]>

Defined in: packages/workflow-engine/src/persistence/interface.ts:904

Read unpublished events ordered by (workflowRunId, sequence).

Parameters​

limit?​

number

Returns​

Promise<OutboxRecord[]>

Inherited from​

PersistenceCore.getUnpublishedOutboxEvents


incrementOutboxRetryCount()​

incrementOutboxRetryCount(id): Promise<number>

Defined in: packages/workflow-engine/src/persistence/interface.ts:891

Increment retry count for a failed outbox event. Returns new count.

Parameters​

id​

string

Returns​

Promise<number>

Inherited from​

PersistenceCore.incrementOutboxRetryCount


insertDefinitionIfAbsent()​

insertDefinitionIfAbsent(input): Promise<WorkflowDefinitionRecord | null>

Defined in: packages/workflow-engine/src/persistence/interface.ts:806

Stores a definition snapshot if (workflowId, version) is not already present, and returns the stored row either way — so a caller can tell whether an explicit version is being re-registered with a different structure. Returns null when PersistenceCore.supportsDefinitionVersioning is false.

Parameters​

input​

CreateDefinitionInput

Returns​

Promise<WorkflowDefinitionRecord | null>

Inherited from​

PersistenceCore.insertDefinitionIfAbsent


listAnnotations()​

listAnnotations(workflowRunId, filters?): Promise<WorkflowAnnotationRecord[]>

Defined in: packages/workflow-engine/src/persistence/interface.ts:884

List annotations for a run, optionally filtered. Returns rows ordered by createdAt ascending (so consumers get a timeline by default).

Parameters​

workflowRunId​

string

filters?​

AnnotationFilters

Returns​

Promise<WorkflowAnnotationRecord[]>

Inherited from​

PersistenceCore.listAnnotations


listRunsForPurge()​

listRunsForPurge(cutoff, statuses, limit): Promise<PurgeableRun[]>

Defined in: packages/workflow-engine/src/persistence/interface.ts:723

Terminal runs eligible for retention deletion, oldest first: runs whose status is one of statuses and which finished at or before cutoff (completedAt <= cutoff, or updatedAt <= cutoff for a terminal run with no completedAt). Returns at most limit runs, each with the ids of its stage records so the caller can clear a pluggable StepLedger before the row goes. Called by run.purge.

Parameters​

cutoff​

Date

statuses​

readonly PurgeableRunStatus[]

limit​

number

Returns​

Promise<PurgeableRun[]>

Inherited from​

PersistenceCore.listRunsForPurge


markOutboxEventsPublished()​

markOutboxEventsPublished(ids): Promise<void>

Defined in: packages/workflow-engine/src/persistence/interface.ts:922

Mark events as published.

Parameters​

ids​

string[]

Returns​

Promise<void>

Inherited from​

PersistenceCore.markOutboxEventsPublished


moveOutboxEventToDLQ()​

moveOutboxEventToDLQ(id): Promise<void>

Defined in: packages/workflow-engine/src/persistence/interface.ts:894

Move an outbox event to DLQ (sets dlqAt).

Parameters​

id​

string

Returns​

Promise<void>

Inherited from​

PersistenceCore.moveOutboxEventToDLQ


releaseIdempotencyKey()​

releaseIdempotencyKey(key, commandType): Promise<void>

Defined in: packages/workflow-engine/src/persistence/interface.ts:956

Release an in-progress idempotency key after command failure.

Parameters​

key​

string

commandType​

string

Returns​

Promise<void>

Inherited from​

PersistenceCore.releaseIdempotencyKey


releaseOutboxEvents()​

releaseOutboxEvents(ids): Promise<void>

Defined in: packages/workflow-engine/src/persistence/interface.ts:919

Un-claim events (clear publishedAt) so a later flush retries them.

Parameters​

ids​

string[]

Returns​

Promise<void>

Inherited from​

PersistenceCore.releaseOutboxEvents


replayDLQEvents()​

replayDLQEvents(maxEvents): Promise<number>

Defined in: packages/workflow-engine/src/persistence/interface.ts:897

Reset DLQ events so they can be reprocessed by outbox.flush. Returns count reset.

Parameters​

maxEvents​

number

Returns​

Promise<number>

Inherited from​

PersistenceCore.replayDLQEvents


supportsDefinitionVersioning()​

supportsDefinitionVersioning(): boolean

Defined in: packages/workflow-engine/src/persistence/interface.ts:777

Whether this adapter's schema carries the definition-versioning columns and table. false on a database that has not been migrated: the engine then behaves exactly as it did before versioning existed rather than failing to start.

Returns​

boolean

Inherited from​

PersistenceCore.supportsDefinitionVersioning


updateRun()​

updateRun(id, data): Promise<void>

Defined in: packages/workflow-engine/src/persistence/interface.ts:710

Updates a run's fields. version is incremented on every call, whether or not expectedVersion is supplied -- callers relying on optimistic concurrency (e.g. job.execute's claimed-run guard) can always detect a concurrent write, including unconditional writes like run.cancel.

Parameters​

id​

string

data​

UpdateRunInput

Returns​

Promise<void>

Inherited from​

PersistenceCore.updateRun


updateStage()​

updateStage(id, data): Promise<void>

Defined in: packages/workflow-engine/src/persistence/interface.ts:828

Parameters​

id​

string

data​

UpdateStageInput

Returns​

Promise<void>

Inherited from​

PersistenceCore.updateStage


upsertStage()​

upsertStage(data): Promise<WorkflowStageRecord>

Defined in: packages/workflow-engine/src/persistence/interface.ts:827

Parameters​

data​

UpsertStageInput

Returns​

Promise<WorkflowStageRecord>

Inherited from​

PersistenceCore.upsertStage


withTransaction()​

withTransaction<T>(fn): Promise<T>

Defined in: packages/workflow-engine/src/persistence/interface.ts:978

Execute operations within a transaction boundary. Redeclared (not merely inherited from PersistenceCore) so the callback receives the full WorkflowPersistence surface.

Type Parameters​

T​

T

Parameters​

fn​

(tx) => Promise<T>

Returns​

Promise<T>

Overrides​

PersistenceCore.withTransaction