Skip to content

@persistent-ai/fireflow-executor / server / PostgresExecutionStore

Class: PostgresExecutionStore ​

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:105

Implements ​

Constructors ​

Constructor ​

new PostgresExecutionStore(db): PostgresExecutionStore

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:106

Parameters ​

db ​

NodePgDatabase<Record<string, unknown>> & object

Returns ​

PostgresExecutionStore

Methods ​

batchGetDirectChildren() ​

batchGetDirectChildren(parentIds): Promise<ChildFeedItem[]>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:630

Batch-fetch direct children of multiple parents in one query. Used by the chain auto-loader in TraceView: fetches single children of chain nodes.

Parameters ​

parentIds ​

string[]

Returns ​

Promise<ChildFeedItem[]>

Implementation of ​

IExecutionStore.batchGetDirectChildren


claimExecution() ​

claimExecution(executionId, workerId, timeoutMs): Promise<boolean>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:660

Claim an execution for a worker atomically Uses PostgreSQL row-level locking to ensure only one worker can claim an execution Also updates the execution's processing fields in the same transaction

Parameters ​

executionId ​

string

workerId ​

string

timeoutMs ​

number

Returns ​

Promise<boolean>

Implementation of ​

IExecutionStore.claimExecution


create() ​

create(row): Promise<void>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:110

Parameters ​

row ​
commitHash ​

string

completedAt ​

Date | null

createdAt ​

Date

errorMessage ​

string | null

errorNodeId ​

string | null

executionDepth ​

number

externalEvents ​

object[] | null

failureCount ​

number

flowId ​

string

id ​

string

integration ​

object & object | null

lastFailureAt ​

Date | null

lastFailureReason ​

string | null

options ​

{ breakpoints?: string[]; debug?: boolean; execution?: { flowTimeoutMs?: number; maxConcurrency?: number; nodeTimeoutMs?: number; }; } | null

ownerId ​

string

parentExecutionId ​

string | null

path ​

string

processingStartedAt ​

Date | null

processingWorkerId ​

string | null

ref ​

string

rootExecutionId ​

string | null

startedAt ​

Date | null

status ​

ExecutionStatus

updatedAt ​

Date

workspaceId ​

string

Returns ​

Promise<void>

Implementation of ​

IExecutionStore.create


delete() ​

delete(id): Promise<boolean>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:140

Parameters ​

id ​

string

Returns ​

Promise<boolean>

Implementation of ​

IExecutionStore.delete


expireOldClaims() ​

expireOldClaims(): Promise<number>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:863

Expire old claims that have passed their expiration time

Returns ​

Promise<number>

Implementation of ​

IExecutionStore.expireOldClaims


extendClaim() ​

extendClaim(executionId, workerId, timeoutMs): Promise<boolean>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:795

Extend an execution claim (heartbeat)

Parameters ​

executionId ​

string

workerId ​

string

timeoutMs ​

number

Returns ​

Promise<boolean>

Implementation of ​

IExecutionStore.extendClaim


get() ​

get(id): Promise<{ commitHash: string; completedAt: Date | null; createdAt: Date; errorMessage: string | null; errorNodeId: string | null; executionDepth: number; externalEvents: object[] | null; failureCount: number; flowId: string; id: string; integration: object & object | null; lastFailureAt: Date | null; lastFailureReason: string | null; options: { breakpoints?: string[]; debug?: boolean; execution?: { flowTimeoutMs?: number; maxConcurrency?: number; nodeTimeoutMs?: number; }; } | null; ownerId: string; parentExecutionId: string | null; path: string; processingStartedAt: Date | null; processingWorkerId: string | null; ref: string; rootExecutionId: string | null; startedAt: Date | null; status: ExecutionStatus; updatedAt: Date; workspaceId: string; } | null>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:126

Parameters ​

id ​

string

Returns ​

Promise<{ commitHash: string; completedAt: Date | null; createdAt: Date; errorMessage: string | null; errorNodeId: string | null; executionDepth: number; externalEvents: object[] | null; failureCount: number; flowId: string; id: string; integration: object & object | null; lastFailureAt: Date | null; lastFailureReason: string | null; options: { breakpoints?: string[]; debug?: boolean; execution?: { flowTimeoutMs?: number; maxConcurrency?: number; nodeTimeoutMs?: number; }; } | null; ownerId: string; parentExecutionId: string | null; path: string; processingStartedAt: Date | null; processingWorkerId: string | null; ref: string; rootExecutionId: string | null; startedAt: Date | null; status: ExecutionStatus; updatedAt: Date; workspaceId: string; } | null>

Implementation of ​

IExecutionStore.get


getActiveClaims() ​

getActiveClaims(): Promise<ExecutionClaim[]>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:819

Get all active claims

Returns ​

Promise<ExecutionClaim[]>

Implementation of ​

IExecutionStore.getActiveClaims


getChildExecutions() ​

getChildExecutions(parentId): Promise<object[]>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:333

Get all child executions of a parent

Parameters ​

parentId ​

string

Returns ​

Promise<object[]>

Implementation of ​

IExecutionStore.getChildExecutions


getChildFeed() ​

getChildFeed(rootExecutionId, limit, cursor?): Promise<{ items: ChildFeedItem[]; nextCursor: string | null; }>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:425

Get a cursor-paginated feed of child executions for a root execution. Returns newest-first. Used by the Execution Feed UI.

Parameters ​

rootExecutionId ​

string

limit ​

number

cursor? ​

Date

Returns ​

Promise<{ items: ChildFeedItem[]; nextCursor: string | null; }>

Implementation of ​

IExecutionStore.getChildFeed


getClaimForExecution() ​

getClaimForExecution(executionId): Promise<ExecutionClaim | null>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:838

Get claim for a specific execution

Parameters ​

executionId ​

string

Returns ​

Promise<ExecutionClaim | null>

Implementation of ​

IExecutionStore.getClaimForExecution


getDirectChildFeed() ​

getDirectChildFeed(parentExecutionId, limit, cursor?): Promise<{ items: ChildFeedItem[]; nextCursor: string | null; }>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:591

Get a cursor-paginated feed of DIRECT children of a specific execution. Used by the TraceView children panel (works at any depth, not just root executions).

Parameters ​

parentExecutionId ​

string

limit ​

number

cursor? ​

Date

Returns ​

Promise<{ items: ChildFeedItem[]; nextCursor: string | null; }>

Implementation of ​

IExecutionStore.getDirectChildFeed


getExecutionMessages() ​

getExecutionMessages(executionId): Promise<DbosExecutionMessage[]>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:1090

Parameters ​

executionId ​

string

Returns ​

Promise<DbosExecutionMessage[]>

Implementation of ​

IExecutionStore.getExecutionMessages


getExecutionsNeedingRecovery() ​

getExecutionsNeedingRecovery(limit?): Promise<object[]>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:909

Get executions that need recovery Finds executions that are:

  1. Status = 'created' but no active claim
  2. Status = 'running' but claim expired
  3. Status = 'created' with failureCount < 5 and lastFailureAt > 1 minute ago

Parameters ​

limit? ​

number = 1000

Returns ​

Promise<object[]>

Implementation of ​

IExecutionStore.getExecutionsNeedingRecovery


getExecutionStats() ​

getExecutionStats(rootExecutionId): Promise<ExecutionStats>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:499

Get aggregate stats for a root execution's entire subtree. Used by the Execution Feed UI to show throughput and status counts.

Parameters ​

rootExecutionId ​

string

Returns ​

Promise<ExecutionStats>

Implementation of ​

IExecutionStore.getExecutionStats


getExecutionSteps() ​

getExecutionSteps(executionId): Promise<DbosExecutionStep[]>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:1052

Parameters ​

executionId ​

string

Returns ​

Promise<DbosExecutionStep[]>

Implementation of ​

IExecutionStore.getExecutionSteps


getExecutionSystemMeta() ​

getExecutionSystemMeta(executionId): Promise<DbosExecutionSystemMeta | null>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:1119

Parameters ​

executionId ​

string

Returns ​

Promise<DbosExecutionSystemMeta | null>

Implementation of ​

IExecutionStore.getExecutionSystemMeta


getExecutionTree() ​

getExecutionTree(executionId): Promise<ExecutionTreeNode[]>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:346

Get the execution subtree rooted at executionId. Works for both root executions (depth=0) and non-root executions (depth>0). For non-root executions, fetches the full tree by rootExecutionId first, then traverses downward from executionId.

Parameters ​

executionId ​

string

Returns ​

Promise<ExecutionTreeNode[]>

Implementation of ​

IExecutionStore.getExecutionTree


getRecoveryHistory() ​

getRecoveryHistory(executionId): Promise<object[]>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:977

Get recovery history for an execution

Parameters ​

executionId ​

string

Returns ​

Promise<object[]>

Implementation of ​

IExecutionStore.getRecoveryHistory


getRootExecutions() ​

getRootExecutions(coordinate, limit?, after?, filters?): Promise<RootExecution[]>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:242

Parameters ​

coordinate ​
path ​

string

ref? ​

string

workspaceId ​

string

limit? ​

number = 50

after? ​

Date | null

filters? ​
createdFrom? ​

Date

createdTo? ​

Date

status? ​

ExecutionStatus

Returns ​

Promise<RootExecution[]>

Implementation of ​

IExecutionStore.getRootExecutions


listExecutionStreams() ​

listExecutionStreams(executionId): Promise<Omit<DbosStreamSummary, "label">[]>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:1023

Parameters ​

executionId ​

string

Returns ​

Promise<Omit<DbosStreamSummary, "label">[]>

Implementation of ​

IExecutionStore.listExecutionStreams


recordRecovery() ​

recordRecovery(executionId, workerId, reason, previousStatus?, previousWorkerId?): Promise<void>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:957

Record a recovery action for an execution

Parameters ​

executionId ​

string

workerId ​

string

reason ​

string

previousStatus? ​

string

previousWorkerId? ​

string

Returns ​

Promise<void>

Implementation of ​

IExecutionStore.recordRecovery


releaseExecution() ​

releaseExecution(executionId, workerId): Promise<void>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:761

Release an execution claim Also clears the processing fields in the execution row

Parameters ​

executionId ​

string

workerId ​

string

Returns ​

Promise<void>

Implementation of ​

IExecutionStore.releaseExecution


releaseRecoveryLock() ​

releaseRecoveryLock(lockId): Promise<void>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:1005

Release a PostgreSQL advisory lock

Parameters ​

lockId ​

number

Returns ​

Promise<void>

Implementation of ​

IExecutionStore.releaseRecoveryLock


tryAcquireRecoveryLock() ​

tryAcquireRecoveryLock(lockId): Promise<boolean>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:989

Try to acquire a PostgreSQL advisory lock Returns true if lock was acquired, false if another session holds it

Parameters ​

lockId ​

number

Returns ​

Promise<boolean>

Implementation of ​

IExecutionStore.tryAcquireRecoveryLock


updateExecutionStatus() ​

updateExecutionStatus(params): Promise<boolean>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:208

Atomically update execution status and related fields This method avoids the need to fetch the execution first

Parameters ​

params ​

UpdateExecutionStatusParams

Returns ​

Promise<boolean>

Implementation of ​

IExecutionStore.updateExecutionStatus


updateExecutionStatusIfNotTerminal() ​

updateExecutionStatusIfNotTerminal(params): Promise<boolean>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:225

Like updateExecutionStatus, but guarded: only updates the row when it is NOT already in a terminal state. The WHERE clause makes this a no-op (returns false) for rows already completed/failed/stopped, so it is safe to call unconditionally from the workflow finally block — on the success/normal-error paths the row is already terminal (written via a DBOS step), and only the cancel path (where those steps throw) leaves a non-terminal row to flip.

Parameters ​

params ​

UpdateExecutionStatusParams

Returns ​

Promise<boolean>

Implementation of ​

IExecutionStore.updateExecutionStatusIfNotTerminal


updateFailureInfo() ​

updateFailureInfo(executionId, failureCount, reason): Promise<boolean>

Defined in: packages/fireflow-executor/server/stores/postgres/postgres-execution-store.ts:883

Update failure count and reason for an execution

Parameters ​

executionId ​

string

failureCount ​

number

reason ​

string

Returns ​

Promise<boolean>

Implementation of ​

IExecutionStore.updateFailureInfo

Licensed under BUSL-1.1