Skip to main content

Class: InMemoryWorkflowPersistence

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:84

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 InMemoryWorkflowPersistence(opts?): InMemoryWorkflowPersistence

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:109

Parameters​

opts?​

InMemoryPersistenceOptions = {}

Returns​

InMemoryWorkflowPersistence

Methods​

acquireIdempotencyKey()​

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

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:741

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/testing/in-memory-persistence.ts:850

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/testing/in-memory-persistence.ts:645

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

Parameters​

events​

CreateOutboxEventInput[]

Returns​

Promise<void>

Implementation of​

WorkflowPersistence.appendOutboxEvents


claimNextPendingRun()​

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

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:294

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

attempt?​

number = 0

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/testing/in-memory-persistence.ts:276

Parameters​

id​

string

Returns​

Promise<boolean>


claimUnpublishedOutboxEvents()​

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

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:680

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


clear()​

clear(): void

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:1047

Clear all data - useful between tests

Returns​

void


completeIdempotencyKey()​

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

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:774

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/testing/in-memory-persistence.ts:998

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/testing/in-memory-persistence.ts:628

Parameters​

data​

CreateLogInput

Returns​

Promise<void>

Implementation of​

WorkflowPersistence.createLog


createRun()​

createRun(data): Promise<WorkflowRunRecord>

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:137

Parameters​

data​

CreateRunInput

Returns​

Promise<WorkflowRunRecord>

Implementation of​

WorkflowPersistence.createRun


createStage()​

createStage(data): Promise<WorkflowStageRecord>

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:365

Parameters​

data​

CreateStageInput

Returns​

Promise<WorkflowStageRecord>

Implementation of​

WorkflowPersistence.createStage


deleteArtifact()​

deleteArtifact(runId, key): Promise<void>

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

Parameters​

runId​

string

key​

string

Returns​

Promise<void>


deleteRun()​

deleteRun(id): Promise<void>

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:228

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/testing/in-memory-persistence.ts:607

Parameters​

id​

string

Returns​

Promise<void>

Implementation of​

WorkflowPersistence.deleteStage


getAllAnnotations()​

getAllAnnotations(): WorkflowAnnotationRecord[]

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:1096

Get all annotations for inspection

Returns​

WorkflowAnnotationRecord[]


getAllArtifacts()​

getAllArtifacts(): WorkflowArtifactRecord[]

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:1089

Get all artifacts for inspection

Returns​

WorkflowArtifactRecord[]


getAllLogs()​

getAllLogs(): WorkflowLogRecord[]

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:1082

Get all logs for inspection

Returns​

WorkflowLogRecord[]


getAllRuns()​

getAllRuns(): WorkflowRunRecord[]

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:1065

Get all runs for inspection

Returns​

WorkflowRunRecord[]


getAllStages()​

getAllStages(): WorkflowStageRecord[]

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:1072

Get all stages for inspection

Returns​

WorkflowStageRecord[]


getDefinition()​

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

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:990

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/testing/in-memory-persistence.ts:578

Parameters​

runId​

string

Returns​

Promise<WorkflowStageRecord | null>


getFirstSuspendedStageReadyToResume()​

getFirstSuspendedStageReadyToResume(runId): Promise<WorkflowStageRecord | null>

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:567

Parameters​

runId​

string

Returns​

Promise<WorkflowStageRecord | null>


getLastCompletedStage()​

getLastCompletedStage(runId): Promise<WorkflowStageRecord | null>

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:585

Parameters​

runId​

string

Returns​

Promise<WorkflowStageRecord | null>


getLastCompletedStageBefore()​

getLastCompletedStageBefore(runId, executionGroup): Promise<WorkflowStageRecord | null>

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:595

Parameters​

runId​

string

executionGroup​

number

Returns​

Promise<WorkflowStageRecord | null>


getRun()​

getRun(id): Promise<WorkflowRunRecord | null>

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:192

Parameters​

id​

string

Returns​

Promise<WorkflowRunRecord | null>

Implementation of​

WorkflowPersistence.getRun


getRunsByStatus()​

getRunsByStatus(status): Promise<WorkflowRunRecord[]>

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:202

Parameters​

status​

Status

Returns​

Promise<WorkflowRunRecord[]>


getRunStatus()​

getRunStatus(id): Promise<Status | null>

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:197

Parameters​

id​

string

Returns​

Promise<Status | null>

Implementation of​

WorkflowPersistence.getRunStatus


getStage()​

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

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:481

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/testing/in-memory-persistence.ts:490

Parameters​

id​

string

Returns​

Promise<WorkflowStageRecord | null>


getStageIdForArtifact()​

getStageIdForArtifact(runId, stageId): Promise<string | null>

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:838

Parameters​

runId​

string

stageId​

string

Returns​

Promise<string | null>


getStagesByRun()​

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

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:495

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/testing/in-memory-persistence.ts:255

Parameters​

stuckSince​

Date

Returns​

Promise<WorkflowRunRecord[]>

Implementation of​

WorkflowPersistence.getStuckRuns


getSuspendedStages()​

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

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:524

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/testing/in-memory-persistence.ts:667

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/testing/in-memory-persistence.ts:824

Parameters​

runId​

string

key​

string

Returns​

Promise<boolean>


incrementOutboxRetryCount()​

incrementOutboxRetryCount(id): Promise<number>

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

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/testing/in-memory-persistence.ts:973

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/testing/in-memory-persistence.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[]>

Implementation of​

WorkflowPersistence.listAnnotations


listArtifacts()​

listArtifacts(runId): Promise<WorkflowArtifactRecord[]>

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:832

Parameters​

runId​

string

Returns​

Promise<WorkflowArtifactRecord[]>


listRunsForPurge()​

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

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:208

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/testing/in-memory-persistence.ts:819

Parameters​

runId​

string

key​

string

Returns​

Promise<unknown>


markOutboxEventsPublished()​

markOutboxEventsPublished(ids): Promise<void>

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:704

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/testing/in-memory-persistence.ts:720

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/testing/in-memory-persistence.ts:789

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/testing/in-memory-persistence.ts:697

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/testing/in-memory-persistence.ts:726

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/testing/in-memory-persistence.ts:799

Parameters​

data​

SaveArtifactInput

Returns​

Promise<void>


saveStageOutput()​

saveStageOutput(runId, workflowType, stageId, output): Promise<string>

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:944

Parameters​

runId​

string

workflowType​

string

stageId​

string

output​

unknown

Returns​

Promise<string>


supportsDefinitionVersioning()​

supportsDefinitionVersioning(): boolean

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:969

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/testing/in-memory-persistence.ts:165

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/testing/in-memory-persistence.ts:420

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/testing/in-memory-persistence.ts:448

Parameters​

workflowRunId​

string

stageId​

string

data​

UpdateStageInput

Returns​

Promise<void>


upsertStage()​

upsertStage(data): Promise<WorkflowStageRecord>

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:402

Parameters​

data​

UpsertStageInput

Returns​

Promise<WorkflowStageRecord>

Implementation of​

WorkflowPersistence.upsertStage


withTransaction()​

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

Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:127

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