Skip to main content

Interface: PersistenceCore

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

Extended by​

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"; }>


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>


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>


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


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[]>


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>


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[]>


createLog()​

createLog(data): Promise<void>

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

Parameters​

data​

CreateLogInput

Returns​

Promise<void>


createRun()​

createRun(data): Promise<WorkflowRunRecord>

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

Parameters​

data​

CreateRunInput

Returns​

Promise<WorkflowRunRecord>


createStage()​

createStage(data): Promise<WorkflowStageRecord>

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

Parameters​

data​

CreateStageInput

Returns​

Promise<WorkflowStageRecord>


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>


deleteStage()​

deleteStage(id): Promise<void>

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

Parameters​

id​

string

Returns​

Promise<void>


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>


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>


getRun()​

getRun(id): Promise<WorkflowRunRecord | null>

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

Parameters​

id​

string

Returns​

Promise<WorkflowRunRecord | null>


getRunStatus()​

getRunStatus(id): Promise<Status | null>

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

Parameters​

id​

string

Returns​

Promise<Status | null>


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>


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[]>


getStuckRuns()​

getStuckRuns(stuckSince): Promise<WorkflowRunRecord[]>

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

Parameters​

stuckSince​

Date

Returns​

Promise<WorkflowRunRecord[]>


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[]>


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[]>


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>


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>


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[]>


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[]>


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>


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>


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>


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>


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>


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


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>


updateStage()​

updateStage(id, data): Promise<void>

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

Parameters​

id​

string

data​

UpdateStageInput

Returns​

Promise<void>


upsertStage()​

upsertStage(data): Promise<WorkflowStageRecord>

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

Parameters​

data​

UpsertStageInput

Returns​

Promise<WorkflowStageRecord>


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>