Interface: Persistence
Defined in: packages/workflow-engine/src/kernel/ports.ts:246
Metadata storage port: run, stage, log, annotation, outbox, and
idempotency-key operations. Artifact payloads are handled by the
BlobStore port instead (see kernel/helpers/create-storage-shim.ts).
Derives from PersistenceCore (persistence/interface.ts) -- the ~26
methods the kernel's handlers/helpers actually call -- instead of
hand-duplicating signatures here. Keeping this as a distinct name (not
just an alias) lets kernel code read Persistence while the rest of
the package reasons about PersistenceCore directly; the two are
structurally identical. WorkflowPersistence (the full 41-method
contract implemented by PrismaWorkflowPersistence /
InMemoryWorkflowPersistence) is a structural superset of
PersistenceCore, so both satisfy this port with no adapter needed --
narrowing the port's declared surface to what the kernel actually calls
only loosens what createKernel demands, it doesn't break anything
that already provided the full contract.
Annotation append call sites (for context): run.create, external
kernel.annotations.attach, and the stage-completion transactions in
job-execute and stage-poll-suspended.
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
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
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?
Returns
Promise<DefinitionVersionCount[]>
Inherited from
PersistenceCore.countRunsByDefinitionVersion
createLog()
createLog(
data):Promise<void>
Defined in: packages/workflow-engine/src/persistence/interface.ts:864
Parameters
data
Returns
Promise<void>
Inherited from
createRun()
createRun(
data):Promise<WorkflowRunRecord>
Defined in: packages/workflow-engine/src/persistence/interface.ts:702
Parameters
data
Returns
Promise<WorkflowRunRecord>
Inherited from
createStage()
createStage(
data):Promise<WorkflowStageRecord>
Defined in: packages/workflow-engine/src/persistence/interface.ts:826
Parameters
data
Returns
Promise<WorkflowStageRecord>
Inherited from
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
deleteStage()
deleteStage(
id):Promise<void>
Defined in: packages/workflow-engine/src/persistence/interface.ts:861
Parameters
id
string
Returns
Promise<void>
Inherited from
ensureDefinitionVersioningDetected()?
optionalensureDefinitionVersioningDetected():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
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
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
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
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?
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
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
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?
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
Returns
Promise<void>
Inherited from
updateStage()
updateStage(
id,data):Promise<void>
Defined in: packages/workflow-engine/src/persistence/interface.ts:828
Parameters
id
string
data
Returns
Promise<void>
Inherited from
upsertStage()
upsertStage(
data):Promise<WorkflowStageRecord>
Defined in: packages/workflow-engine/src/persistence/interface.ts:827
Parameters
data
Returns
Promise<WorkflowStageRecord>
Inherited from
withTransaction()
withTransaction<
T>(fn):Promise<T>
Defined in: packages/workflow-engine/src/persistence/interface.ts:699
Execute operations within a transaction boundary. The callback
receives a PersistenceCore-scoped handle, which is all the kernel
ever needs inside a transaction. (WorkflowPersistence redeclares
this method with a WorkflowPersistence-scoped tx so existing
callers that use artifact methods inside a transaction are
unaffected -- see below.)
Type Parameters
T
T
Parameters
fn
(tx) => Promise<T>
Returns
Promise<T>