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
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
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?
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
Returns
Promise<void>
Implementation of
createRun()
createRun(
data):Promise<WorkflowRunRecord>
Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:137
Parameters
data
Returns
Promise<WorkflowRunRecord>
Implementation of
createStage()
createStage(
data):Promise<WorkflowStageRecord>
Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:365
Parameters
data
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
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
getAllArtifacts()
getAllArtifacts():
WorkflowArtifactRecord[]
Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:1089
Get all artifacts for inspection
Returns
getAllLogs()
getAllLogs():
WorkflowLogRecord[]
Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:1082
Get all logs for inspection
Returns
getAllRuns()
getAllRuns():
WorkflowRunRecord[]
Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:1065
Get all runs for inspection
Returns
getAllStages()
getAllStages():
WorkflowStageRecord[]
Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:1072
Get all stages for inspection
Returns
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
getRunsByStatus()
getRunsByStatus(
status):Promise<WorkflowRunRecord[]>
Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:202
Parameters
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
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?
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
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
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
Returns
Promise<void>
Implementation of
updateStage()
updateStage(
id,data):Promise<void>
Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:420
Parameters
id
string
data
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
Returns
Promise<void>
upsertStage()
upsertStage(
data):Promise<WorkflowStageRecord>
Defined in: packages/workflow-engine/src/testing/in-memory-persistence.ts:402
Parameters
data
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>