Class: PrismaWorkflowPersistence
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:199
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.
Implements
Constructors
Constructor
new PrismaWorkflowPersistence(
prisma,options?):PrismaWorkflowPersistence
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:214
Parameters
prisma
EnginePrismaClient
options?
PrismaWorkflowPersistenceOptions = {}
Returns
PrismaWorkflowPersistence
Methods
acquireIdempotencyKey()
acquireIdempotencyKey(
key,commandType,options?):Promise<{status:"acquired"; } | {result:unknown;status:"replay"; } | {status:"in_progress"; }>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1545
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"; }>
Implementation of
WorkflowPersistence.acquireIdempotencyKey
appendAnnotations()
appendAnnotations(
inputs):Promise<void>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1193
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>
Implementation of
WorkflowPersistence.appendAnnotations
appendOutboxEvents()
appendOutboxEvents(
events):Promise<void>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1354
Write events to the outbox. Sequences are auto-assigned per workflowRunId.
Parameters
events
Returns
Promise<void>
Implementation of
WorkflowPersistence.appendOutboxEvents
claimNextPendingRun()
claimNextPendingRun(
options?):Promise<WorkflowRunRecord|null>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:460
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
serves?
readonly ServedDefinition[]
Returns
Promise<WorkflowRunRecord | null>
The claimed workflow run (now with status RUNNING), or null if no pending runs
Implementation of
WorkflowPersistence.claimNextPendingRun
claimPendingRun()
claimPendingRun(
id):Promise<boolean>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:441
Parameters
id
string
Returns
Promise<boolean>
claimUnpublishedOutboxEvents()
claimUnpublishedOutboxEvents(
limit?):Promise<OutboxRecord[]>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1441
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[]>
Implementation of
WorkflowPersistence.claimUnpublishedOutboxEvents
completeIdempotencyKey()
completeIdempotencyKey(
key,commandType,result):Promise<void>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1642
Mark an idempotency key as completed and cache the command result.
Parameters
key
string
commandType
string
result
unknown
Returns
Promise<void>
Implementation of
WorkflowPersistence.completeIdempotencyKey
countRunsByDefinitionVersion()
countRunsByDefinitionVersion(
filter?):Promise<DefinitionVersionCount[]>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:752
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[]>
Implementation of
WorkflowPersistence.countRunsByDefinitionVersion
createLog()
createLog(
data):Promise<void>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1099
Parameters
data
Returns
Promise<void>
Implementation of
createRun()
createRun(
data):Promise<WorkflowRunRecord>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:306
Parameters
data
Returns
Promise<WorkflowRunRecord>
Implementation of
createStage()
createStage(
data):Promise<WorkflowStageRecord>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:801
Parameters
data
Returns
Promise<WorkflowStageRecord>
Implementation of
WorkflowPersistence.createStage
deleteArtifact()
deleteArtifact(
runId,key):Promise<void>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1159
Parameters
runId
string
key
string
Returns
Promise<void>
deleteRun()
deleteRun(
id):Promise<void>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:415
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>
Implementation of
deleteStage()
deleteStage(
id):Promise<void>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1091
Parameters
id
string
Returns
Promise<void>
Implementation of
WorkflowPersistence.deleteStage
ensureDefinitionVersioningDetected()
ensureDefinitionVersioningDetected():
Promise<boolean>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:266
Confirms the synchronous capability guess against the database, once
per client, and returns the answer every later
supportsDefinitionVersioning() will give.
The guess is structural: it reads the generated client, so it is
true the moment prisma generate runs, which is before migrate deploy on every first migration and on any rolling deploy that ships
code ahead of its migration. Left unconfirmed, run.create fails
P2022, insertDefinitionIfAbsent fails P2021 and every claim dies
on a raw 42703 — the process cannot serve a single run.
The confirmation only ever turns versioning off: a client with no
workflowDefinition delegate cannot use the tables however migrated
the database is. It is a catalogue read (see
probeDefinitionVersioningSchema), so it cannot abort a caller's
transaction, and it is skipped entirely when definitionVersioning
was configured explicitly.
Called by the kernel from the paths that would otherwise touch the
columns — run.create's snapshot recording, run.claimPending and
run.listVersions — so consumers do not have to call it themselves.
Returns
Promise<boolean>
Implementation of
WorkflowPersistence.ensureDefinitionVersioningDetected
getDefinition()
getDefinition(
workflowId,version):Promise<WorkflowDefinitionRecord|null>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:740
Loads one stored definition snapshot, or null when absent.
Parameters
workflowId
string
version
string
Returns
Promise<WorkflowDefinitionRecord | null>
Implementation of
WorkflowPersistence.getDefinition
getFirstFailedStage()
getFirstFailedStage(
runId):Promise<WorkflowStageRecord|null>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1043
Parameters
runId
string
Returns
Promise<WorkflowStageRecord | null>
getFirstSuspendedStageReadyToResume()
getFirstSuspendedStageReadyToResume(
runId):Promise<WorkflowStageRecord|null>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1027
Parameters
runId
string
Returns
Promise<WorkflowStageRecord | null>
getLastCompletedStage()
getLastCompletedStage(
runId):Promise<WorkflowStageRecord|null>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1059
Parameters
runId
string
Returns
Promise<WorkflowStageRecord | null>
getLastCompletedStageBefore()
getLastCompletedStageBefore(
runId,executionGroup):Promise<WorkflowStageRecord|null>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1074
Parameters
runId
string
executionGroup
number
Returns
Promise<WorkflowStageRecord | null>
getRun()
getRun(
id):Promise<WorkflowRunRecord|null>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:361
Parameters
id
string
Returns
Promise<WorkflowRunRecord | null>
Implementation of
getRunsByStatus()
getRunsByStatus(
status):Promise<WorkflowRunRecord[]>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:374
Parameters
status
Returns
Promise<WorkflowRunRecord[]>
getRunStatus()
getRunStatus(
id):Promise<Status|null>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:366
Parameters
id
string
Returns
Promise<Status | null>
Implementation of
WorkflowPersistence.getRunStatus
getStage()
getStage(
runId,stageId):Promise<WorkflowStageRecord|null>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:967
Parameters
runId
string
stageId
string
Returns
Promise<WorkflowStageRecord | null>
Implementation of
getStageById()
getStageById(
id):Promise<WorkflowStageRecord|null>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:979
Parameters
id
string
Returns
Promise<WorkflowStageRecord | null>
getStageIdForArtifact()
getStageIdForArtifact(
runId,stageId):Promise<string|null>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1176
Parameters
runId
string
stageId
string
Returns
Promise<string | null>
getStagesByRun()
getStagesByRun(
runId,options?):Promise<WorkflowStageRecord[]>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:984
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[]>
Implementation of
WorkflowPersistence.getStagesByRun
getStuckRuns()
getStuckRuns(
stuckSince):Promise<WorkflowRunRecord[]>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:422
Parameters
stuckSince
Date
Returns
Promise<WorkflowRunRecord[]>
Implementation of
WorkflowPersistence.getStuckRuns
getSuspendedStages()
getSuspendedStages(
beforeDate,options?):Promise<WorkflowStageRecord[]>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1001
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[]>
Implementation of
WorkflowPersistence.getSuspendedStages
getUnpublishedOutboxEvents()
getUnpublishedOutboxEvents(
limit?):Promise<OutboxRecord[]>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1431
Read unpublished events ordered by (workflowRunId, sequence).
Parameters
limit?
number
Returns
Promise<OutboxRecord[]>
Implementation of
WorkflowPersistence.getUnpublishedOutboxEvents
hasArtifact()
hasArtifact(
runId,key):Promise<boolean>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1149
Parameters
runId
string
key
string
Returns
Promise<boolean>
incrementOutboxRetryCount()
incrementOutboxRetryCount(
id):Promise<number>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1509
Increment retry count for a failed outbox event. Returns new count.
Parameters
id
string
Returns
Promise<number>
Implementation of
WorkflowPersistence.incrementOutboxRetryCount
insertDefinitionIfAbsent()
insertDefinitionIfAbsent(
input):Promise<WorkflowDefinitionRecord|null>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:695
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>
Implementation of
WorkflowPersistence.insertDefinitionIfAbsent
listAnnotations()
listAnnotations(
workflowRunId,filters?):Promise<WorkflowAnnotationRecord[]>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1248
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[]>
Implementation of
WorkflowPersistence.listAnnotations
listArtifacts()
listArtifacts(
runId):Promise<WorkflowArtifactRecord[]>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1167
Parameters
runId
string
Returns
Promise<WorkflowArtifactRecord[]>
listRunsForPurge()
listRunsForPurge(
cutoff,statuses,limit):Promise<PurgeableRun[]>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:382
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[]>
Implementation of
WorkflowPersistence.listRunsForPurge
loadArtifact()
loadArtifact(
runId,key):Promise<unknown>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1140
Parameters
runId
string
key
string
Returns
Promise<unknown>
markOutboxEventsPublished()
markOutboxEventsPublished(
ids):Promise<void>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1497
Mark events as published.
Parameters
ids
string[]
Returns
Promise<void>
Implementation of
WorkflowPersistence.markOutboxEventsPublished
moveOutboxEventToDLQ()
moveOutboxEventToDLQ(
id):Promise<void>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1518
Move an outbox event to DLQ (sets dlqAt).
Parameters
id
string
Returns
Promise<void>
Implementation of
WorkflowPersistence.moveOutboxEventToDLQ
releaseIdempotencyKey()
releaseIdempotencyKey(
key,commandType):Promise<void>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1653
Release an in-progress idempotency key after command failure.
Parameters
key
string
commandType
string
Returns
Promise<void>
Implementation of
WorkflowPersistence.releaseIdempotencyKey
releaseOutboxEvents()
releaseOutboxEvents(
ids):Promise<void>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1489
Un-claim events (clear publishedAt) so a later flush retries them.
Parameters
ids
string[]
Returns
Promise<void>
Implementation of
WorkflowPersistence.releaseOutboxEvents
replayDLQEvents()
replayDLQEvents(
maxEvents):Promise<number>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1525
Reset DLQ events so they can be reprocessed by outbox.flush. Returns count reset.
Parameters
maxEvents
number
Returns
Promise<number>
Implementation of
WorkflowPersistence.replayDLQEvents
saveArtifact()
saveArtifact(
data):Promise<void>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1115
Parameters
data
Returns
Promise<void>
saveStageOutput()
saveStageOutput(
runId,workflowType,stageId,output):Promise<string>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:1313
Parameters
runId
string
workflowType
string
stageId
string
output
unknown
Returns
Promise<string>
supportsDefinitionVersioning()
supportsDefinitionVersioning():
boolean
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:239
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
Implementation of
WorkflowPersistence.supportsDefinitionVersioning
updateRun()
updateRun(
id,data):Promise<void>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:328
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>
Implementation of
updateStage()
updateStage(
id,data):Promise<void>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:851
Parameters
id
string
data
Returns
Promise<void>
Implementation of
WorkflowPersistence.updateStage
updateStageByRunAndStageId()
updateStageByRunAndStageId(
workflowRunId,stageId,data):Promise<void>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:884
Parameters
workflowRunId
string
stageId
string
data
Returns
Promise<void>
upsertStage()
upsertStage(
data):Promise<WorkflowStageRecord>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:821
Parameters
data
Returns
Promise<WorkflowStageRecord>
Implementation of
WorkflowPersistence.upsertStage
withTransaction()
withTransaction<
T>(fn):Promise<T>
Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:279
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>