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
readonlyfairnessGroupBy: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
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
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
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?
Returns
Promise<JobAckOutcome>
Implementation of
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?
Returns
Promise<JobAckOutcome>
Implementation of
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
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?
Returns
Promise<DequeueResult | null>
Implementation of
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
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
Returns
Promise<string[]>
Implementation of
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
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?
Returns
Promise<JobAckOutcome>
Implementation of
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
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
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?
Returns
Promise<JobAckOutcome>
Implementation of
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>