Skip to main content

Class: PrismaJobQueue

Defined in: packages/workflow-engine/src/persistence/prisma/job-queue.ts:114

Implements​

Constructors​

Constructor​

new PrismaJobQueue(prisma, options?): PrismaJobQueue

Defined in: packages/workflow-engine/src/persistence/prisma/job-queue.ts:137

Parameters​

prisma​

EnginePrismaClient

options?​

PrismaJobQueueOptions = {}

Returns​

PrismaJobQueue

Properties​

fairnessGroupBy​

readonly fairnessGroupBy: string | null

Defined in: packages/workflow-engine/src/persistence/prisma/job-queue.ts:133

The dotted groupBy path this queue's fairness cap reads, or null when fairness is off. Read by createSpillingJobTransport so a spilled payload still carries its group key.

Implementation of​

JobQueue.fairnessGroupBy

Methods​

adoptWorkerId()​

adoptWorkerId(workerId): string

Defined in: packages/workflow-engine/src/persistence/prisma/job-queue.ts:207

Take the host's worker id unless this queue was constructed with one of its own; returns the id it will stamp on job_queue.workerId.

Parameters​

workerId​

string

Returns​

string

Implementation of​

JobQueue.adoptWorkerId


cancelByRun()​

cancelByRun(workflowRunId): Promise<number>

Defined in: packages/workflow-engine/src/persistence/prisma/job-queue.ts:866

Cancel all pending/suspended jobs for a workflow run.

Parameters​

workflowRunId​

string

Returns​

Promise<number>

Implementation of​

JobQueue.cancelByRun


complete()​

complete(jobId, fence?): Promise<JobAckOutcome>

Defined in: packages/workflow-engine/src/persistence/prisma/job-queue.ts:703

Mark job as completed

Parameters​

jobId​

string

fence?​

JobAckFence

Returns​

Promise<JobAckOutcome>

Implementation of​

JobQueue.complete


defer()​

defer(jobId, nextPollAt, reason, fence?): Promise<JobAckOutcome>

Defined in: packages/workflow-engine/src/persistence/prisma/job-queue.ts:764

Return a claimed job to PENDING with a later nextPollAt, without counting the claim as an attempt: the dequeue incremented attempt, so this decrements it back. See JobQueue.defer for why declining is not failing.

Parameters​

jobId​

string

nextPollAt​

Date

reason​

string

fence?​

JobAckFence

Returns​

Promise<JobAckOutcome>

Implementation of​

JobQueue.defer


deleteByRunAndStages()​

deleteByRunAndStages(workflowRunId, stageIds): Promise<number>

Defined in: packages/workflow-engine/src/persistence/prisma/job-queue.ts:266

Remove every job row for the given stages of a run, whatever their status.

Parameters​

workflowRunId​

string

stageIds​

string[]

Returns​

Promise<number>

Implementation of​

JobQueue.deleteByRunAndStages


dequeue()​

dequeue(options?): Promise<DequeueResult | null>

Defined in: packages/workflow-engine/src/persistence/prisma/job-queue.ts:288

Atomically dequeue the next available job Uses FOR UPDATE SKIP LOCKED (PostgreSQL) or optimistic locking (SQLite)

Parameters​

options?​

DequeueOptions

Returns​

Promise<DequeueResult | null>

Implementation of​

JobQueue.dequeue


enqueue()​

enqueue(options): Promise<string>

Defined in: packages/workflow-engine/src/persistence/prisma/job-queue.ts:222

Add a job to the queue, replacing any row already queued for the same (workflowRunId, stageId) — see the JobQueue.enqueueParallel contract.

Parameters​

options​

EnqueueJobInput

Returns​

Promise<string>


enqueueParallel()​

enqueueParallel(jobs): Promise<string[]>

Defined in: packages/workflow-engine/src/persistence/prisma/job-queue.ts:236

Enqueue multiple stages in parallel (same execution group).

Idempotent on (workflowRunId, stageId): rows already queued for those pairs are removed in the same transaction as the insert, so a run.rerunFrom re-enqueue or a run.reapStuck recovery sweep over a stage that still carries its previous (terminal) job row leaves exactly one PENDING row with attempt back at 0.

Parameters​

jobs​

EnqueueJobInput[]

Returns​

Promise<string[]>

Implementation of​

JobQueue.enqueueParallel


expireRunawayJobs()​

expireRunawayJobs(absoluteTimeoutMs): Promise<number>

Defined in: packages/workflow-engine/src/persistence/prisma/job-queue.ts:1024

Fail every RUNNING job whose claim has exceeded the coarse absolute cap.

The coarse tier of a two-tier expiry: releaseStaleJobs handles the fine-grained heartbeat signal and requeues for another worker when a worker dies, but is defeated by a worker that is alive but wedged (e.g. infinite loop or hung network call) because it keeps heartbeating.

startedAt is the per-claim stamp that the heartbeat never touches (a suspended job that resumes is re-claimed and gets a fresh startedAt, so this measures one attempt's wall time, not the run's). A job past this cap is failed terminally with LEASE_ABSOLUTE_CAP rather than requeued, as a job that hung for the full cap will hang again.

On PostgreSQL, the deadline is derived in-database for the same single-clock reason the heartbeat sweep is. The run itself is resolved afterwards by run.reapStuck on its next pass.

Parameters​

absoluteTimeoutMs​

number

Returns​

Promise<number>

Implementation of​

JobQueue.expireRunawayJobs


fail()​

fail(jobId, error, shouldRetry?, fence?): Promise<JobAckOutcome>

Defined in: packages/workflow-engine/src/persistence/prisma/job-queue.ts:795

Mark job as failed

Parameters​

jobId​

string

error​

string

shouldRetry?​

boolean = false

fence?​

JobAckFence

Returns​

Promise<JobAckOutcome>

Implementation of​

JobQueue.fail


getJobsByWorkflowRun()​

getJobsByWorkflowRun(workflowRunId): Promise<JobRecord[]>

Defined in: packages/workflow-engine/src/persistence/prisma/job-queue.ts:885

Get all job rows for a workflow run (any status).

Parameters​

workflowRunId​

string

Returns​

Promise<JobRecord[]>

Implementation of​

JobQueue.getJobsByWorkflowRun


getWorkerId()​

getWorkerId(): string

Defined in: packages/workflow-engine/src/persistence/prisma/job-queue.ts:213

The id this queue stamps on the jobs it claims.

Returns​

string


releaseStaleJobs()​

releaseStaleJobs(staleThresholdMs?): Promise<number>

Defined in: packages/workflow-engine/src/persistence/prisma/job-queue.ts:952

Release stale locks (for crashed workers). Stamps lastError with the LEASE_HEARTBEAT_LOST prefix so an operator can tell a reclaimed lease from a stage-level failure. The fine-grained tier of a two-tier expiry whose coarse tier is expireRunawayJobs.

On PostgreSQL, the deadline is derived in-database from the same now() the claim stamped, so the sweep does not depend on the sweeping host's system clock. "updatedAt" is set explicitly because a raw UPDATE does not trigger Prisma's @updatedAt.

SQLite has no now() AT TIME ZONE, so it keeps the application-clock comparison (single-process by nature, where the two clocks are the same clock anyway).

Parameters​

staleThresholdMs?​

number = 300000

Returns​

Promise<number>

Implementation of​

JobQueue.releaseStaleJobs


suspend()​

suspend(jobId, nextPollAt, fence?): Promise<JobAckOutcome>

Defined in: packages/workflow-engine/src/persistence/prisma/job-queue.ts:728

Mark job as suspended (for async-batch)

Parameters​

jobId​

string

nextPollAt​

Date

fence?​

JobAckFence

Returns​

Promise<JobAckOutcome>

Implementation of​

JobQueue.suspend


touchJob()​

touchJob(jobId): Promise<void>

Defined in: packages/workflow-engine/src/persistence/prisma/job-queue.ts:918

Refresh a running job's lease without changing status.

Parameters​

jobId​

string

Returns​

Promise<void>

Implementation of​

JobQueue.touchJob