Skip to main content

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​

CreateAnnotationInput[]

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​

CreateOutboxEventInput[]

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?​

DefinitionVersionCountFilter

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​

CreateLogInput

Returns​

Promise<void>

Implementation of​

WorkflowPersistence.createLog


createRun()​

createRun(data): Promise<WorkflowRunRecord>

Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:306

Parameters​

data​

CreateRunInput

Returns​

Promise<WorkflowRunRecord>

Implementation of​

WorkflowPersistence.createRun


createStage()​

createStage(data): Promise<WorkflowStageRecord>

Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:801

Parameters​

data​

CreateStageInput

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​

WorkflowPersistence.deleteRun


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​

WorkflowPersistence.getRun


getRunsByStatus()​

getRunsByStatus(status): Promise<WorkflowRunRecord[]>

Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:374

Parameters​

status​

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​

WorkflowPersistence.getStage


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?​

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​

CreateDefinitionInput

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​

SaveArtifactInput

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​

UpdateRunInput

Returns​

Promise<void>

Implementation of​

WorkflowPersistence.updateRun


updateStage()​

updateStage(id, data): Promise<void>

Defined in: packages/workflow-engine/src/persistence/prisma/persistence.ts:851

Parameters​

id​

string

data​

UpdateStageInput

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​

UpdateStageInput

Returns​

Promise<void>


upsertStage()​

upsertStage(data): Promise<WorkflowStageRecord>

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

Parameters​

data​

UpsertStageInput

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>

Implementation of​

WorkflowPersistence.withTransaction