import assert from 'node:assert' import type { IndexBloatOptions, JobMatchStrategy, ManagedIndex, ManagedFunction } from './types.ts' import { indexKeysRaw, indexIncludeRaw, indexPredicateRaw, displayIndexDefinition, extractFunctionBody, normalizeFunctionBody } from './drifter.ts' import type { ExpectedColumns, ExpectedConstraints } from './drifter.ts' import schemaManifest from './schema.json' with { type: 'json' } import { normalizeSchemaName, resolveSchemaName } from './tools.ts' export interface SqlQuery { text: string values: unknown[] } export const PG_ERROR = { divisionByZero: '22012', lockNotAvailable: '55P03', foreignKeyViolation: '23503', undefinedTable: '42P01', checkViolation: '23514' } export const DEFAULT_SCHEMA = 'pgboss' export const MIGRATE_RACE_MESSAGE = 'division by zero' export const CREATE_RACE_MESSAGE = 'already exists' export const SINGLE_QUOTE_REGEX = /'/g const FIFTEEN_MINUTES = 60 * 15 const FORTEEN_DAYS = 60 * 60 * 24 * 14 const SEVEN_DAYS = 60 * 60 * 24 * 7 // A bam row stuck at 'in_progress' past this means the process that claimed it died or was // stopped mid-command (bam.ts #processCommands returns without marking the row on #stopped or a // crash), and getNextBamCommand needs to reclaim it or every future async migration wedges forever. // This is the *fallback* backstop only, deliberately long (24h) because false-reclaim is the worse // failure: reclaiming a still-running CREATE INDEX CONCURRENTLY runs two builds on the same index at // once, and an interrupted CONCURRENTLY leaves an INVALID index needing manual cleanup. The intended // *primary* trigger is a liveness check against pg_stat_progress_create_index (reclaim as soon as no // backend is actually building), which recovers in minutes instead of a day; the 24h timeout only // covers cases liveness can't classify (non-index commands, or a build that never registered). const BAM_STALE_SECONDS = 60 * 60 * 24 // Liveness grace window: on the native-Postgres path we don't trust a "no build running" reading // until a claimed command has had time to actually start and register in pg_stat_progress_create_index // (pool latency between claiming the row and issuing the build). Below this age, an in_progress row is // only reclaimed via the 24h fallback, never via liveness, so a genuinely-running build is never // yanked out from under itself. const BAM_LIVENESS_GRACE_SECONDS = 60 * 5 export const JOB_STATES = Object.freeze({ created: 'created', retry: 'retry', active: 'active', completed: 'completed', cancelled: 'cancelled', failed: 'failed' }) export const QUEUE_POLICIES = Object.freeze({ standard: 'standard', short: 'short', singleton: 'singleton', stately: 'stately', exclusive: 'exclusive', key_strict_fifo: 'key_strict_fifo' }) /** * How the expression in a schedule row is read: a cron expression, or an RFC 5545 recurrence rule. * * Stored on the row rather than derived from the expression on every pass, so the format is decided * once, by whoever writes the schedule, and every reader agrees with that decision. A row written * straight into the table with SQL has to name its own kind; the column defaults to cron, which is * what every schedule written before rules existed is. */ export const SCHEDULE_KINDS = Object.freeze({ cron: 'cron', rrule: 'rrule' } as const) /** * What a schedule does about occurrences that came due while no cron pass ran. * * A pass sends what the last minute holds, so nothing outside that window is sent at all: a * deployment that was down, or between deploys, for an hour never sends the occurrences that hour * held. `skip` is that behavior, and stays the default, since a job appearing on the first pass * after a deploy for work whose moment has passed is a surprise nobody asked for. * * `once` sends a single job however many occurrences were missed, which is what a schedule whose * job reads the current state of the world wants: a nightly report that missed three nights is one * report, not three. One job needs no occurrence identity to be worth running, which is what a * policy sending a job per missed occurrence would need and has no way to carry: the forwarded job * gets the schedule's `data` and nothing else, so a handler could not tell which hour of an outage * each of twelve identical jobs was for. * * Part of the options blob rather than a column of its own: the pass reads every schedule row * anyway, and nothing queries the table by policy. */ export const SCHEDULE_MISSED_POLICIES = Object.freeze({ skip: 'skip', once: 'once' } as const) const QUEUE_DEFAULTS = { expire_seconds: FIFTEEN_MINUTES, retention_seconds: FORTEEN_DAYS, deletion_seconds: SEVEN_DAYS, retry_limit: 2, retry_delay: 0, warning_queued: 0, retry_backoff: false, partition: false } // The root of the job table hierarchy, under both install shapes: the partitioned parent that // COMMON_JOB_TABLE and every per-queue partition attach to, and the plain physical table that // carries the rows and indexes itself where noTablePartitioning is set. Also the template every // partition is cloned from (CREATE TABLE ... LIKE). export const BASE_JOB_TABLE = 'job' export const COMMON_JOB_TABLE = 'job_common' interface CreateOptions { createSchema?: boolean noTablePartitioning?: boolean noDeferrableConstraints?: boolean noAdvisoryLocks?: boolean noCoveringIndexes?: boolean inlineTableIndexes?: boolean } // Folds the statements that follow a CREATE TABLE (ADD PRIMARY KEY, ADD CONSTRAINT, CREATE // [UNIQUE] INDEX on that table) into the CREATE TABLE itself, in CockroachDB's inline syntax. See // inlineTableIndexes. It rewrites the same statements the standalone path runs rather than keeping // a second copy of every definition, so the two paths cannot drift apart. Anything it does not // recognise throws: silently dropping an index is the failure this must never have. export const SERVER_VERSION = 'SELECT version() AS version' export function inlineIntoCreateTable (createTable: string, statements: string[]): string { const clauses = statements .flatMap(statement => statement.split(';')) .map(statement => statement.trim().replace(/\s+/g, ' ')) .filter(Boolean) .map(statement => { let m = statement.match(/^ALTER TABLE \S+ ADD PRIMARY KEY (\(.+\))$/i) if (m) return `PRIMARY KEY ${m[1]}` m = statement.match(/^ALTER TABLE \S+ ADD CONSTRAINT (\S+) (.+)$/i) if (m) return `CONSTRAINT ${m[1]} ${m[2]}` m = statement.match(/^CREATE (UNIQUE )?INDEX (?:IF NOT EXISTS )?(\S+) ON \S+ (\(.+)$/i) if (m) { assert(!/\bINCLUDE\s*\(/i.test(m[3]), `inlineIntoCreateTable: covering index ${m[2]} has no inline form`) return `${m[1] ?? ''}INDEX ${m[2]} ${m[3]}` } throw new Error(`inlineIntoCreateTable: cannot inline: ${statement}`) }) const close = createTable.lastIndexOf(')') assert(close > 0, 'inlineIntoCreateTable: no column list to extend') return `${createTable.slice(0, close).trimEnd()},\n ${clauses.join(',\n ')}\n ${createTable.slice(close)}` } export function create (schema: string, version: number, options?: CreateOptions) { const noPartitioning = options?.noTablePartitioning ?? false const noDeferrable = options?.noDeferrableConstraints ?? false const noLocks = options?.noAdvisoryLocks ?? false const noCovering = options?.noCoveringIndexes ?? false // Only meaningful without partitioning: a partitioned job table's indexes live on the partitions. const inline = (options?.inlineTableIndexes ?? false) && noPartitioning if (inline) { return locked(schema, createInline(schema, version, { createSchema: options?.createSchema, noDeferrable, noCovering }), undefined, noLocks) } const commands = [ options?.createSchema ? createSchema(schema) : '', createEnumJobState(schema), createClockFunction(schema), createTableVersion(schema), createTableQueue(schema), createTableSchedule(schema), createTableSubscription(schema), createTableBam(schema), // Partition-helper functions are only used by the partitioned architecture. // They are unused when partitioning is disabled, and job_table_format's // IMMUTABLE + format() body is rejected at create time by databases like // CockroachDB, so skip them entirely in noTablePartitioning mode. noPartitioning ? '' : jobTableFormatFunction(schema), noPartitioning ? '' : jobTableRunFunction(schema), noPartitioning ? '' : jobTableRunAsyncFunction(schema), createTableJob(schema, noPartitioning), createPrimaryKeyJob(schema), noPartitioning ? createTableJobIndexes(schema, noDeferrable, noCovering) : createTableJobCommon(schema), createTableWarning(schema), createIndexWarning(schema), createTableQueueStats(schema, noPartitioning), createIndexQueueStats(schema, noCovering), noPartitioning ? '' : ensureQueueStatsPartitions(schema), createTableJobDependency(schema), createIndexJobDependencyParent(schema), createTableInstance(schema), createQueueFunction(schema, noPartitioning), deleteQueueFunction(schema, noPartitioning), insertVersion(schema, version) ] return locked(schema, commands, undefined, noLocks) } // The install for inlineTableIndexes (CockroachDB): the same objects as create(), non-partitioned, // with every table's indexes and constraints declared inside its CREATE TABLE. function createInline (schema: string, version: number, options: { createSchema?: boolean, noDeferrable: boolean, noCovering: boolean }) { return [ options.createSchema ? createSchema(schema) : '', createEnumJobState(schema), createClockFunction(schema), createTableVersion(schema), createTableQueue(schema), createTableSchedule(schema), createTableSubscription(schema), createTableBam(schema), inlineIntoCreateTable(createTableJob(schema, true), [ createPrimaryKeyJob(schema), createTableJobIndexes(schema, options.noDeferrable, options.noCovering) ]), inlineIntoCreateTable(createTableWarning(schema), [createIndexWarning(schema)]), inlineIntoCreateTable(createTableQueueStats(schema, true), [createIndexQueueStats(schema, options.noCovering)]), inlineIntoCreateTable(createTableJobDependency(schema), [createIndexJobDependencyParent(schema)]), createTableInstance(schema), createQueueFunction(schema, true), deleteQueueFunction(schema, true), insertVersion(schema, version) ] } function createSchema (schema: string) { return `CREATE SCHEMA IF NOT EXISTS ${schema}` } function createEnumJobState (schema: string) { // ENUM definition order is important // base type is numeric and first values are less than last values return ` CREATE TYPE ${schema}.job_state AS ENUM ( '${JOB_STATES.created}', '${JOB_STATES.retry}', '${JOB_STATES.active}', '${JOB_STATES.completed}', '${JOB_STATES.cancelled}', '${JOB_STATES.failed}' ) ` } // The one place pg-boss SQL reads the clock. A single-statement LANGUAGE sql function with no SET // clause, so the planner inlines it and plans are identical to calling pg_catalog.now() directly; // STABLE matches the body and is what CockroachDB requires to inline. A TestClock swaps the body // (pg-boss #689). // // Column defaults stay on pg_catalog.now(). Every pg-boss write names its timestamps, so a default // that read this function would only ever serve rows pg-boss did not write - and it would tie the // function to every table that carries the default (create_queue() copies them into each partition // with LIKE job INCLUDING DEFAULTS), which CockroachDB records as a dependency and refuses to drop. // Only create_queue()'s own body names it, which is why v42's uninstall restores that body before // dropping the function. export const CLOCK_FUNCTION_BODY = 'SELECT pg_catalog.now();' // Sessions opt into the fake clock through this setting, so an instance without a TestClock on // the same schema, or one started after a killed run left the override behind, stays on real time. export const CLOCK_OVERRIDE_SETTING = 'pgboss.test_clock' // Where a TestClock keeps the fake time. Deliberately a name no user would choose for their own // table: attach() takes over whatever is sitting under it, so it has to be unmistakably ours. export function clockTable (schema: string) { return `${schema}.__pgboss_test_clock` } // The body a TestClock installs: the single row of the clock table for a session that opted in, else // the real clock. The subquery defeats inlining, which is fine in tests and never happens in production. export function clockOverrideBody (schema: string) { return `SELECT COALESCE(CASE WHEN current_setting('${CLOCK_OVERRIDE_SETTING}', true) = 'on' THEN (SELECT c.now FROM ${clockTable(schema)} c LIMIT 1) END, pg_catalog.now());` } export function enableClockOverride () { return `SET ${CLOCK_OVERRIDE_SETTING} = 'on'` } export function disableClockOverride () { return `RESET ${CLOCK_OVERRIDE_SETTING}` } // The stored source of the clock function, for spotting an override a killed test run left behind. // prosrc rather than pg_get_functiondef: the latter is unsupported on CockroachDB, and the caller // only needs to know whether the body reads CLOCK_OVERRIDE_SETTING, not to diff it. CockroachDB // rewrites what it stores, so this must never be compared against CLOCK_FUNCTION_BODY for equality. export function getClockFunctionSource (schema: string) { return ` SELECT p.prosrc AS source FROM pg_proc p JOIN pg_namespace n ON n.oid = p.pronamespace WHERE n.nspname = '${resolveSchemaName(schema).replace(SINGLE_QUOTE_REGEX, "''")}' AND p.proname = 'job_now' ` } // Whether a stored body is a TestClock's rather than the shipped one. A positive test for the // setting, not a diff against CLOCK_FUNCTION_BODY: CockroachDB stores its own rewriting of the // canonical body ('SELECT now():::TIMESTAMPTZ;'), so equality would call every CockroachDB install // overridden. Any spelling of the override still names the setting. export function clockFunctionIsOverridden (source: string | null | undefined) { return !!source?.includes(CLOCK_OVERRIDE_SETTING) } // Puts the clock function back the way a fresh install leaves it and clears the table the override // body read from. Used to undo a TestClock that was never released, e.g. a killed test run. export function restoreClockFunction (schema: string) { return ` ${createClockFunction(schema, { replace: true })} DROP TABLE IF EXISTS ${clockTable(schema)}; ` } export function createClockFunction (schema: string, options: { replace?: boolean, body?: string } = {}) { const { replace = false, body = CLOCK_FUNCTION_BODY } = options return ` CREATE ${replace ? 'OR REPLACE ' : ''}FUNCTION ${schema}.job_now() RETURNS timestamp with time zone AS $$ ${body} $$ LANGUAGE sql STABLE; ` } // The *_on columns are single-row interval claims, one per background pass, so that one instance // per interval does the work rather than every instance racing to: // // reindex_on the index-bloat maintenance pass, against every instance issuing the same // REINDEX. // monitor_backoff_on the odd one out - a deadline rather than a last-run stamp. Before it, no // instance may start another queue-stats aggregate, so autovacuum gets a // horizon-free window it cannot miss. Null on a database that has never // needed one. function createTableVersion (schema: string) { return ` CREATE TABLE ${schema}.version ( version int primary key, cron_on timestamp with time zone, bam_on timestamp with time zone, flow_on timestamp with time zone, reindex_on timestamp with time zone, monitor_backoff_on timestamp with time zone ) ` } // Two stamps, not one. monitor_claim_on is the interval claim that decides which instance runs a // monitor pass; monitor_on is when this queue's counts were actually written, and is stamped only by // the aggregate that wrote them. Splitting them is what lets a pass be claimed and then skip the // aggregate - because the vacuum backoff is in force, or because another instance holds the stats // try-lock - without capturedOn claiming a freshness the counts do not have. // created_delta / completed_delta / failed_delta are not gauges like the counts beside them: they // are how many jobs went through between the previous monitor pass and the latest one. delta_on and // delta_seconds are the window those three cover, null until a pass counts. wait_bins and run_bins // are the wait and run times of the jobs that finished in that window, as histograms of LATENCY_SLOTS // counts (see LATENCY_BINS), and ready_oldest_seconds is how long the oldest job ready to run had // waited at the pass. Like the deltas, all of them are null until a pass counts, and a pass that // counts writes a value: every slot null when nothing finished, 0 when nothing was waiting. /* eslint-disable no-restricted-syntax -- column defaults stay on the real clock: every pg-boss write names its timestamps through job_now() */ function createTableQueue (schema: string) { return ` CREATE TABLE ${schema}.queue ( name text NOT NULL, policy text NOT NULL, retry_limit int NOT NULL, retry_delay int NOT NULL, retry_backoff bool NOT NULL, retry_delay_max int, expire_seconds int NOT NULL, retention_seconds int NOT NULL, deletion_seconds int NOT NULL, dead_letter text REFERENCES ${schema}.queue (name) CHECK (dead_letter IS DISTINCT FROM name), partition bool NOT NULL, table_name text NOT NULL, deferred_count int NOT NULL default 0, blocked_count int NOT NULL default 0, queued_count int NOT NULL default 0, ready_count int NOT NULL default 0, warning_queued int NOT NULL default 0, active_count int NOT NULL default 0, failed_count int NOT NULL default 0, total_count int NOT NULL default 0, created_delta int NOT NULL default 0, completed_delta int NOT NULL default 0, failed_delta int NOT NULL default 0, delta_on timestamp with time zone, delta_seconds int, wait_bins int[], run_bins int[], ready_oldest_seconds int, ready_history int[] NOT NULL default '{}', heartbeat_seconds int, notify bool NOT NULL DEFAULT false, singletons_active text[], monitor_claim_on timestamp with time zone, monitor_on timestamp with time zone, maintain_on timestamp with time zone, created_on timestamp with time zone not null default now(), updated_on timestamp with time zone not null default now(), PRIMARY KEY (name) ) ` } // `cron` holds the expression whatever its format, and `kind` says which format that is: the column // predates rules and renaming it would break every consumer reading the table, from the dashboard to // a hand-written query. // // `timezone` defaults to UTC rather than to null, so a row written straight into the table with SQL // gets the zone schedule() would have given it. Nullable still, because an instance on an older // release can write a null during a rolling upgrade, which is what the read-side COALESCE covers. /* eslint-enable no-restricted-syntax */ /* eslint-disable no-restricted-syntax -- column defaults stay on the real clock: every pg-boss write names its timestamps through job_now() */ function createTableSchedule (schema: string) { return ` CREATE TABLE ${schema}.schedule ( name text REFERENCES ${schema}.queue ON DELETE CASCADE, key text not null DEFAULT '', kind text not null DEFAULT '${SCHEDULE_KINDS.cron}' CHECK (kind IN ('${SCHEDULE_KINDS.cron}', '${SCHEDULE_KINDS.rrule}')), cron text not null, timezone text DEFAULT 'UTC', data jsonb, options jsonb, created_on timestamp with time zone not null default now(), updated_on timestamp with time zone not null default now(), last_job_id uuid, PRIMARY KEY (name, key) ) ` } /* eslint-enable no-restricted-syntax */ /* eslint-disable no-restricted-syntax -- column defaults stay on the real clock: every pg-boss write names its timestamps through job_now() */ function createTableSubscription (schema: string) { return ` CREATE TABLE ${schema}.subscription ( event text not null, name text not null REFERENCES ${schema}.queue ON DELETE CASCADE, created_on timestamp with time zone not null default now(), updated_on timestamp with time zone not null default now(), PRIMARY KEY(event, name) ) ` } // created_on defaults to clock_timestamp(), not ${schema}.job_now(), so multiple job_table_run_async() enqueues // within a single migration transaction keep their insertion order, BAM applies queued commands in // created_on order, and some migrations enqueue an ordered drop-then-rebuild pair (see v33). /* eslint-enable no-restricted-syntax */ /* eslint-disable no-restricted-syntax -- bam.created_on orders several enqueues within one transaction; the schema clock would tie them */ function createTableBam (schema: string) { return ` CREATE TABLE ${schema}.bam ( id uuid PRIMARY KEY default gen_random_uuid(), name text NOT NULL, version int NOT NULL, status text NOT NULL DEFAULT 'pending', queue text, table_name text NOT NULL, command text NOT NULL, error text, created_on timestamp with time zone NOT NULL DEFAULT clock_timestamp(), started_on timestamp with time zone, completed_on timestamp with time zone ) ` } /* eslint-enable no-restricted-syntax */ /* eslint-disable no-restricted-syntax -- column defaults stay on the real clock: every pg-boss write names its timestamps through job_now() */ export function createTableWarning (schema: string) { return ` CREATE TABLE ${schema}.warning ( id uuid PRIMARY KEY default gen_random_uuid(), type text NOT NULL, message text NOT NULL, data jsonb, created_on timestamp with time zone NOT NULL DEFAULT now() ) ` } export function createIndexWarning (schema: string) { return `CREATE INDEX warning_i1 ON ${schema}.warning (created_on DESC)` } /* eslint-enable no-restricted-syntax */ export function createTableJobDependency (schema: string) { return ` CREATE TABLE ${schema}.job_dependency ( child_name text NOT NULL, child_id uuid NOT NULL, parent_name text NOT NULL, parent_id uuid NOT NULL, PRIMARY KEY (child_name, child_id, parent_name, parent_id) ) ` } // An instance is live until it misses this many heartbeats, and its row is deleted once its heartbeat // has not moved for this many days. export const INSTANCE_QUIET_BEATS = 3 export const INSTANCE_RETENTION_DAYS = 7 // A crash never sets stopped_on, so a crash loop leaves a quiet row per restart, and a deploy whose // processes never call stop() leaves one per replica. Registering keeps the newest of these dead rows // for its own name (for its host, when unnamed), and maintenance keeps the newest overall, so neither // can grow the table for the whole retention. Keyed on the name rather than the host because a // Kubernetes pod that is recreated comes back under a new hostname. export const INSTANCE_DEAD_KEPT_PER_NAME = 20 export const INSTANCE_DEAD_KEPT = 1000 // Stopped, or quiet: the complement of getInstances' live. function instanceDead (schema: string) { return `(stopped_on IS NOT NULL OR heartbeat_on < ${schema}.job_now() - heartbeat_seconds * ${INSTANCE_QUIET_BEATS} * interval '1 second')` } // Deletes the dead rows ranked past `keep`, newest start first, among those `where` selects. Ranked // before locking and locked with SKIP LOCKED, so a registration and maintenance pruning at once // cannot deadlock, and neither deletes a row the other's ranking kept. Used on every backend, unlike // the fetch's SKIP LOCKED: where it can pass over an unlocked row (CockroachDB), that row is only // left for the next prune. function pruneDeadInstances (schema: string, where: string, keep: number) { return ` DELETE FROM ${schema}.instance WHERE id IN ( SELECT id FROM ${schema}.instance WHERE id IN ( SELECT id FROM ( SELECT id, row_number() OVER (ORDER BY started_on DESC, id DESC) as n FROM ${schema}.instance WHERE ${where} AND ${instanceDead(schema)} ) ranked WHERE n > ${keep} ) FOR UPDATE SKIP LOCKED ) ` } // Run by an instance after it registers: $1 its id, $2 its name, $3 its host. export function pruneInstanceLives (schema: string, keep: number) { return pruneDeadInstances(schema, 'id <> $1 AND (CASE WHEN $2::text IS NULL THEN name IS NULL AND host = $3 ELSE name = $2 END)', keep) } export function trimDeadInstances (schema: string, keep: number) { return pruneDeadInstances(schema, 'true', keep) } // Crash restarts: how many lives in a row on this name and host ended without stop() before this one // started. Worked out from the rows rather than carried from one life to the next, because at // registration a row that still reads live is either a sibling process (pm2, a second PgBoss) or a // predecessor that crashed seconds ago, and only its next missed heartbeats tell which. So a crashed // life counts toward an instance only once it is dead, and only if its last heartbeat came before the // instance started, which a live sibling's keeps moving past. // // $1 this instance's id, $2 its name, $3 its host, $4 its pid, $5 when its process started: a row // with this pid whose heartbeat predates that is an earlier process that reused the pid, dead at once // (a container restart, pid 1 each time), and a row with this pid started since is another PgBoss // object in this same process, never a predecessor. function crashSlot (schema: string) { return ` me AS (SELECT id, started_on FROM ${schema}.instance WHERE id = $1), slot AS ( SELECT i.* FROM ${schema}.instance i, me WHERE i.id <> me.id AND i.host = $3 AND (i.name = $2 OR (i.name IS NULL AND $2::text IS NULL)) AND i.started_on < me.started_on AND i.heartbeat_on <= me.started_on AND NOT (i.pid = $4 AND i.started_on >= $5::timestamptz) ), boundary AS (SELECT max(started_on) as at FROM slot WHERE stopped_on IS NOT NULL), streak AS ( SELECT s.* FROM slot s, boundary b WHERE s.stopped_on IS NULL AND (b.at IS NULL OR s.started_on > b.at) )` } function crashDead (schema: string) { return `(heartbeat_on < ${schema}.job_now() - heartbeat_seconds * ${INSTANCE_QUIET_BEATS} * interval '1 second' OR (pid = $4 AND heartbeat_on < $5::timestamptz))` } // Numbers the dead lives of the current streak from the earliest one still kept, whose count stands // for any the pruning has taken, and writes each later life's count and this instance's. Rewriting // the streak is what keeps the counts exact: a life that crashed before its own recheck is corrected // by its successor, so the earliest kept row is always right when the pruning moves past it. export function countCrashRestarts (schema: string) { return ` WITH ${crashSlot(schema)}, crashed AS ( SELECT id, crash_restarts, crash_restarts_since, heartbeat_on, row_number() OVER (ORDER BY started_on, id) as n, count(*) OVER () as k FROM streak WHERE ${crashDead(schema)} ), head AS ( SELECT crash_restarts as base, COALESCE(crash_restarts_since, heartbeat_on) as since, k FROM crashed WHERE n = 1 ), counts AS ( SELECT c.id, h.base + c.n - 1 as restarts, h.since FROM crashed c, head h WHERE c.n > 1 UNION ALL SELECT me.id, COALESCE(h.base + h.k, 0), h.since FROM me LEFT JOIN head h ON true ) UPDATE ${schema}.instance i SET crash_restarts = counts.restarts::int, crash_restarts_since = counts.since FROM counts WHERE i.id = counts.id ` } // When the rows this count could not yet judge would go quiet, or null when there are none: the // registrar counts again then. export function crashRecountAt (schema: string) { return ` WITH ${crashSlot(schema)} SELECT max(heartbeat_on + heartbeat_seconds * ${INSTANCE_QUIET_BEATS} * interval '1 second') as "recountAt" FROM streak WHERE NOT ${crashDead(schema)} ` } // Instance registry statements. $1 is the instance id throughout. Registering replaces the row, so a // PgBoss object restarted after stop() reads as live again from a new started_on. A heartbeat is an // upsert too: a row pruned while its process was paused (a laptop asleep past the retention) comes // back on the next beat, keeping the started_on the instance registered with ($20). const INSTANCE_COLUMNS = `id, name, host, pid, version, node_version, heartbeat_seconds, supervise, schedule, migrate, persist_queue_stats, persist_warnings, pool_max, pool_total, pool_idle, pool_waiting, workers, metrics, config, application_name, started_on, heartbeat_on` const INSTANCE_VALUES = `$1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17::text::jsonb, $18::text::jsonb, $19::text::jsonb, current_setting('application_name')` const INSTANCE_BEAT = `pool_max = EXCLUDED.pool_max, pool_total = EXCLUDED.pool_total, pool_idle = EXCLUDED.pool_idle, pool_waiting = EXCLUDED.pool_waiting, workers = EXCLUDED.workers, metrics = EXCLUDED.metrics, heartbeat_on = EXCLUDED.heartbeat_on` export function registerInstance (schema: string) { return ` INSERT INTO ${schema}.instance (${INSTANCE_COLUMNS}) VALUES (${INSTANCE_VALUES}, ${schema}.job_now(), ${schema}.job_now()) ON CONFLICT (id) DO UPDATE SET name = EXCLUDED.name, host = EXCLUDED.host, pid = EXCLUDED.pid, version = EXCLUDED.version, node_version = EXCLUDED.node_version, heartbeat_seconds = EXCLUDED.heartbeat_seconds, supervise = EXCLUDED.supervise, schedule = EXCLUDED.schedule, migrate = EXCLUDED.migrate, persist_queue_stats = EXCLUDED.persist_queue_stats, persist_warnings = EXCLUDED.persist_warnings, application_name = EXCLUDED.application_name, config = EXCLUDED.config, started_on = EXCLUDED.started_on, stopped_on = NULL, crash_restarts = 0, crash_restarts_since = NULL, ${INSTANCE_BEAT} RETURNING started_on as "startedOn" ` } export function heartbeatInstance (schema: string) { return ` INSERT INTO ${schema}.instance (${INSTANCE_COLUMNS}) VALUES (${INSTANCE_VALUES}, $20::timestamptz, ${schema}.job_now()) ON CONFLICT (id) DO UPDATE SET ${INSTANCE_BEAT} ` } export function stopInstance (schema: string) { return `UPDATE ${schema}.instance SET stopped_on = ${schema}.job_now() WHERE id = $1` } // Rows whose heartbeat has not moved for the retention, stopped or quiet alike. A live instance's // row never qualifies, since its heartbeat keeps moving. export function deleteOldInstances (schema: string, days: number) { return ` DELETE FROM ${schema}.instance WHERE heartbeat_on < ${schema}.job_now() - interval '${days} days' ` } export function getInstances (schema: string) { return ` SELECT id, name, host, pid, version, node_version as "nodeVersion", application_name as "applicationName", heartbeat_seconds as "heartbeatSeconds", supervise, schedule, migrate, persist_queue_stats as "persistQueueStats", persist_warnings as "persistWarnings", pool_max as "poolMax", pool_total as "poolTotal", pool_idle as "poolIdle", pool_waiting as "poolWaiting", workers, metrics, config, crash_restarts as "crashRestarts", crash_restarts_since as "crashRestartsSince", started_on as "startedOn", heartbeat_on as "heartbeatOn", stopped_on as "stoppedOn", stopped_on IS NULL AND heartbeat_on >= ${schema}.job_now() - heartbeat_seconds * ${INSTANCE_QUIET_BEATS} * interval '1 second' as live FROM ${schema}.instance ORDER BY started_on, id ` } export function createIndexJobDependencyParent (schema: string) { return `CREATE INDEX IF NOT EXISTS job_dep_parent_idx ON ${schema}.job_dependency (parent_name, parent_id)` } /* eslint-disable no-restricted-syntax -- column defaults stay on the real clock: every pg-boss write names its timestamps through job_now() */ // One row per PgBoss object, written by the object itself at start() and on each heartbeat, so a // database can say which instances share it and what each is doing. A crashed process never sets // stopped_on; it goes quiet instead, when heartbeat_on stops moving, and heartbeat_seconds says how // long that takes for this row, since the interval is per instance. workers is one entry per work() // call. The pool columns are node-postgres's counts at the last heartbeat, null for a pool pg-boss // was handed and cannot read. metrics is the process's CPU, memory and event loop at the last // heartbeat against its container's limits (see nurse.ts), null until the first sample lands. // config is the options it runs with, as the registrar's allowlist picks them, for comparing instances. // application_name is the one the registering session carried, so pg_stat_activity joins to the row // wherever it is unique. export function createTableInstance (schema: string) { return ` CREATE TABLE ${schema}.instance ( id uuid PRIMARY KEY, name text, host text NOT NULL, pid int NOT NULL, version text NOT NULL, node_version text NOT NULL, application_name text, heartbeat_seconds int NOT NULL, supervise bool NOT NULL, schedule bool NOT NULL, migrate bool NOT NULL, persist_queue_stats bool NOT NULL, persist_warnings bool NOT NULL, pool_max int, pool_total int, pool_idle int, pool_waiting int, workers jsonb NOT NULL DEFAULT '[]'::jsonb, metrics jsonb, config jsonb NOT NULL DEFAULT '{}'::jsonb, crash_restarts int NOT NULL DEFAULT 0, crash_restarts_since timestamptz, started_on timestamptz NOT NULL DEFAULT now(), heartbeat_on timestamptz NOT NULL DEFAULT now(), stopped_on timestamptz ) ` } /* eslint-enable no-restricted-syntax */ // Anchored so a schema name that itself contains these substrings (e.g. `job_intake`) isn't // mangled: `\.job\y` matches only the base table reference (`schema.job`, not `schema.job_i5` whose // `job` is followed by `_`, nor `.job_dependency`), and `\yjob_i(\d+)` matches only the bare // index-name tokens (job_i1..9), never the `job_i` inside a schema name. Mirrors formatJobTable() // in migrationStore.ts; the migration that fixed this (v37) carries its own frozen copy. export function jobTableFormatFunction (schema: string) { return ` CREATE FUNCTION ${schema}.job_table_format(command text, table_name text) RETURNS text AS $$ SELECT format( regexp_replace( regexp_replace(command, '\\.job\\y', '.%1$I', 'g'), '\\yjob_i(\\d+)', '%1$s_i\\1', 'g' ), table_name ); $$ LANGUAGE sql IMMUTABLE; ` } function jobTableRunFunction (schema: string) { return ` CREATE FUNCTION ${schema}.job_table_run(command text, tbl_name text DEFAULT NULL, queue_name text DEFAULT NULL) RETURNS VOID AS $$ DECLARE tbl RECORD; BEGIN IF queue_name IS NOT NULL THEN SELECT table_name INTO tbl_name FROM ${schema}.queue WHERE name = queue_name; END IF; IF tbl_name IS NOT NULL THEN EXECUTE ${schema}.job_table_format(command, tbl_name); RETURN; END IF; EXECUTE ${schema}.job_table_format(command, '${COMMON_JOB_TABLE}'); FOR tbl IN SELECT table_name FROM ${schema}.queue WHERE partition = true LOOP EXECUTE ${schema}.job_table_format(command, tbl.table_name); END LOOP; END; $$ LANGUAGE plpgsql; ` } function jobTableRunAsyncFunction (schema: string) { return ` CREATE FUNCTION ${schema}.job_table_run_async(command_name text, version int, command text, tbl_name text DEFAULT NULL, queue_name text DEFAULT NULL) RETURNS VOID AS $$ BEGIN IF queue_name IS NOT NULL THEN SELECT table_name INTO tbl_name FROM ${schema}.queue WHERE name = queue_name; END IF; IF tbl_name IS NOT NULL THEN INSERT INTO ${schema}.bam (name, version, status, queue, table_name, command) VALUES ( command_name, version, 'pending', queue_name, tbl_name, ${schema}.job_table_format(command, tbl_name) ); RETURN; END IF; INSERT INTO ${schema}.bam (name, version, status, queue, table_name, command) SELECT command_name, version, 'pending', NULL, '${COMMON_JOB_TABLE}', ${schema}.job_table_format(command, '${COMMON_JOB_TABLE}') UNION ALL SELECT command_name, version, 'pending', queue.name, queue.table_name, ${schema}.job_table_format(command, queue.table_name) FROM ${schema}.queue WHERE partition = true; END; $$ LANGUAGE plpgsql; ` } /* eslint-disable no-restricted-syntax -- column defaults stay on the real clock: every pg-boss write names its timestamps through job_now() */ function createTableJob (schema: string, noPartitioning = false) { // source_name / source_id / source_created_on / source_retry_count / source_output are // dead-letter provenance: where a job in a dead-letter queue came from, stamped at the transfer // so the original queue, id, enqueue time, retry count and final output survive the move. The // original's output is provenance rather than the copy's own `output`: the copy is a new job, // and its `output` is whatever its own run produces. // // source_root_id is the one lineage column that is not about the last hop. source_id names the job // that just failed, which after a redrive is the redriven copy, not the job send() returned; each // round trip through a dead letter queue would add a hop, and retention deletes the failed rows the // hops go through. The root is the id of the first job in the chain, copied onto every dead letter // row and every redriven job after it, so one lookup finds them all however many trips it took. const partitionClause = noPartitioning ? '' : 'PARTITION BY LIST (name)' return ` CREATE TABLE ${schema}.job ( id uuid not null default gen_random_uuid(), name text not null, priority integer not null default(0), data jsonb, state ${schema}.job_state not null default '${JOB_STATES.created}', retry_limit integer not null default ${QUEUE_DEFAULTS.retry_limit}, retry_count integer not null default 0, retry_delay integer not null default ${QUEUE_DEFAULTS.retry_delay}, retry_backoff boolean not null default ${QUEUE_DEFAULTS.retry_backoff}, retry_delay_max integer, expire_seconds int not null default ${QUEUE_DEFAULTS.expire_seconds}, deletion_seconds int not null default ${QUEUE_DEFAULTS.deletion_seconds}, singleton_key text, singleton_on timestamp without time zone, group_id text, group_tier text, start_after timestamp with time zone not null default now(), created_on timestamp with time zone not null default now(), started_on timestamp with time zone, completed_on timestamp with time zone, keep_until timestamp with time zone NOT NULL default now() + interval '${QUEUE_DEFAULTS.retention_seconds}', output jsonb, dead_letter text, policy text, heartbeat_on timestamp with time zone, heartbeat_seconds int, blocked boolean not null default false, blocking boolean not null default false, pending_dependencies int not null default 0, source_name text, source_id uuid, source_created_on timestamp with time zone, source_retry_count int, source_output jsonb, source_root_id uuid, trace_context jsonb, upsert_by_key bool ) ${partitionClause} ` } // retry_count is in the minimal set because a worker settles against it. const JOB_COLUMNS_MIN = 'id, name, data, retry_count as "retryCount", expire_seconds as "expireInSeconds", heartbeat_seconds as "heartbeatSeconds", group_id as "groupId", group_tier as "groupTier"' const JOB_COLUMNS_ALL = `${JOB_COLUMNS_MIN}, policy, state, priority, retry_limit as "retryLimit", retry_delay as "retryDelay", retry_backoff as "retryBackoff", retry_delay_max as "retryDelayMax", start_after as "startAfter", started_on as "startedOn", singleton_key as "singletonKey", singleton_on as "singletonOn", deletion_seconds as "deleteAfterSeconds", heartbeat_on as "heartbeatOn", created_on as "createdOn", completed_on as "completedOn", keep_until as "keepUntil", dead_letter as "deadLetter", blocked, blocking, pending_dependencies as "pendingDependencies", output, source_name as "sourceName", source_id as "sourceId", source_created_on as "sourceCreatedOn", source_retry_count as "sourceRetryCount", source_output as "sourceOutput", source_root_id as "sourceRootId" ` /* eslint-enable no-restricted-syntax */ function createTableJobCommon (schema: string) { return ` CREATE TABLE ${schema}.${COMMON_JOB_TABLE} (LIKE ${schema}.job INCLUDING GENERATED INCLUDING DEFAULTS); SELECT ${schema}.job_table_run($cmd$${createPrimaryKeyJob(schema)}$cmd$, '${COMMON_JOB_TABLE}'); SELECT ${schema}.job_table_run($cmd$${createQueueForeignKeyJob(schema)}$cmd$, '${COMMON_JOB_TABLE}'); SELECT ${schema}.job_table_run($cmd$${createQueueForeignKeyJobDeadLetter(schema)}$cmd$, '${COMMON_JOB_TABLE}'); SELECT ${schema}.job_table_run($cmd$${createIndexJobPolicyShort(schema)}$cmd$, '${COMMON_JOB_TABLE}'); SELECT ${schema}.job_table_run($cmd$${createIndexJobPolicySingleton(schema)}$cmd$, '${COMMON_JOB_TABLE}'); SELECT ${schema}.job_table_run($cmd$${createIndexJobPolicyStately(schema)}$cmd$, '${COMMON_JOB_TABLE}'); SELECT ${schema}.job_table_run($cmd$${createIndexJobPolicyExclusive(schema)}$cmd$, '${COMMON_JOB_TABLE}'); SELECT ${schema}.job_table_run($cmd$${createIndexJobPolicyKeyStrictFifo(schema)}$cmd$, '${COMMON_JOB_TABLE}'); SELECT ${schema}.job_table_run($cmd$${createIndexJobPolicyKeyStrictFifoHeads(schema)}$cmd$, '${COMMON_JOB_TABLE}'); SELECT ${schema}.job_table_run($cmd$${createCheckConstraintKeyStrictFifo(schema)}$cmd$, '${COMMON_JOB_TABLE}'); SELECT ${schema}.job_table_run($cmd$${createIndexJobThrottle(schema)}$cmd$, '${COMMON_JOB_TABLE}'); SELECT ${schema}.job_table_run($cmd$${createIndexJobFetch(schema)}$cmd$, '${COMMON_JOB_TABLE}'); SELECT ${schema}.job_table_run($cmd$${createIndexJobGroupConcurrency(schema)}$cmd$, '${COMMON_JOB_TABLE}'); SELECT ${schema}.job_table_run($cmd$${createIndexJobBlocking(schema)}$cmd$, '${COMMON_JOB_TABLE}'); SELECT ${schema}.job_table_run($cmd$${createIndexJobSourceRoot(schema)}$cmd$, '${COMMON_JOB_TABLE}'); SELECT ${schema}.job_table_run($cmd$${createIndexJobUpsert(schema)}$cmd$, '${COMMON_JOB_TABLE}'); ALTER TABLE ${schema}.job ATTACH PARTITION ${schema}.${COMMON_JOB_TABLE} DEFAULT; ` } // Creates indexes directly on job table when partitioning is disabled function createTableJobIndexes (schema: string, noDeferrableConstraints = false, noCoveringIndex = false) { return ` ${createQueueForeignKeyJob(schema, noDeferrableConstraints)}; ${createQueueForeignKeyJobDeadLetter(schema, noDeferrableConstraints)}; ${createIndexJobPolicyShort(schema)}; ${createIndexJobPolicySingleton(schema)}; ${createIndexJobPolicyStately(schema)}; ${createIndexJobPolicyExclusive(schema)}; ${createIndexJobPolicyKeyStrictFifo(schema)}; ${createIndexJobPolicyKeyStrictFifoHeads(schema, noCoveringIndex)}; ${createCheckConstraintKeyStrictFifo(schema)}; ${createIndexJobThrottle(schema)}; ${createIndexJobFetch(schema, noCoveringIndex)}; ${createIndexJobGroupConcurrency(schema)}; ${createIndexJobBlocking(schema)}; ${createIndexJobSourceRoot(schema)}; ${createIndexJobUpsert(schema)}; ` } function createQueueFunction (schema: string, noPartitioning = false) { if (noPartitioning) { // Simplified version without table partitioning support return ` CREATE FUNCTION ${schema}.create_queue(queue_name text, options jsonb) RETURNS VOID AS $$ BEGIN INSERT INTO ${schema}.queue ( name, policy, retry_limit, retry_delay, retry_backoff, retry_delay_max, expire_seconds, retention_seconds, deletion_seconds, warning_queued, dead_letter, partition, table_name, heartbeat_seconds, created_on, updated_on ) VALUES ( queue_name, options->>'policy', COALESCE((options->>'retryLimit')::int, ${QUEUE_DEFAULTS.retry_limit}), COALESCE((options->>'retryDelay')::int, ${QUEUE_DEFAULTS.retry_delay}), COALESCE((options->>'retryBackoff')::bool, ${QUEUE_DEFAULTS.retry_backoff}), (options->>'retryDelayMax')::int, COALESCE((options->>'expireInSeconds')::int, ${QUEUE_DEFAULTS.expire_seconds}), COALESCE((options->>'retentionSeconds')::int, ${QUEUE_DEFAULTS.retention_seconds}), COALESCE((options->>'deleteAfterSeconds')::int, ${QUEUE_DEFAULTS.deletion_seconds}), COALESCE((options->>'warningQueueSize')::int, ${QUEUE_DEFAULTS.warning_queued}), options->>'deadLetter', false, '${BASE_JOB_TABLE}', (options->>'heartbeatSeconds')::int, ${schema}.job_now(), ${schema}.job_now() ) ON CONFLICT DO NOTHING; END; $$ LANGUAGE plpgsql; ` } return ` CREATE FUNCTION ${schema}.create_queue(queue_name text, options jsonb) RETURNS VOID AS $$ DECLARE tablename varchar := CASE WHEN options->>'partition' = 'true' THEN 'j' || encode(sha224(queue_name::bytea), 'hex') ELSE '${COMMON_JOB_TABLE}' END; queue_created_on timestamptz; BEGIN WITH q as ( INSERT INTO ${schema}.queue ( name, policy, retry_limit, retry_delay, retry_backoff, retry_delay_max, expire_seconds, retention_seconds, deletion_seconds, warning_queued, dead_letter, partition, table_name, heartbeat_seconds, notify, created_on, updated_on ) VALUES ( queue_name, options->>'policy', COALESCE((options->>'retryLimit')::int, ${QUEUE_DEFAULTS.retry_limit}), COALESCE((options->>'retryDelay')::int, ${QUEUE_DEFAULTS.retry_delay}), COALESCE((options->>'retryBackoff')::bool, ${QUEUE_DEFAULTS.retry_backoff}), (options->>'retryDelayMax')::int, COALESCE((options->>'expireInSeconds')::int, ${QUEUE_DEFAULTS.expire_seconds}), COALESCE((options->>'retentionSeconds')::int, ${QUEUE_DEFAULTS.retention_seconds}), COALESCE((options->>'deleteAfterSeconds')::int, ${QUEUE_DEFAULTS.deletion_seconds}), COALESCE((options->>'warningQueueSize')::int, ${QUEUE_DEFAULTS.warning_queued}), options->>'deadLetter', COALESCE((options->>'partition')::bool, ${QUEUE_DEFAULTS.partition}), tablename, (options->>'heartbeatSeconds')::int, COALESCE((options->>'notify')::bool, false), ${schema}.job_now(), ${schema}.job_now() ) ON CONFLICT DO NOTHING RETURNING created_on ) SELECT created_on into queue_created_on from q; IF queue_created_on IS NULL OR options->>'partition' IS DISTINCT FROM 'true' THEN RETURN; END IF; EXECUTE format('CREATE TABLE ${schema}.%I (LIKE ${schema}.job INCLUDING DEFAULTS)', tablename); EXECUTE ${schema}.job_table_format($cmd$${createPrimaryKeyJob(schema)}$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$${createQueueForeignKeyJob(schema)}$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$${createQueueForeignKeyJobDeadLetter(schema)}$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$${createIndexJobFetch(schema)}$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$${createIndexJobThrottle(schema)}$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$${createIndexJobGroupConcurrency(schema)}$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$${createIndexJobBlocking(schema)}$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$${createIndexJobSourceRoot(schema)}$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$${createIndexJobUpsert(schema)}$cmd$, tablename); IF options->>'policy' = 'short' THEN EXECUTE ${schema}.job_table_format($cmd$${createIndexJobPolicyShort(schema)}$cmd$, tablename); ELSIF options->>'policy' = 'singleton' THEN EXECUTE ${schema}.job_table_format($cmd$${createIndexJobPolicySingleton(schema)}$cmd$, tablename); ELSIF options->>'policy' = 'stately' THEN EXECUTE ${schema}.job_table_format($cmd$${createIndexJobPolicyStately(schema)}$cmd$, tablename); ELSIF options->>'policy' = 'exclusive' THEN EXECUTE ${schema}.job_table_format($cmd$${createIndexJobPolicyExclusive(schema)}$cmd$, tablename); ELSIF options->>'policy' = '${QUEUE_POLICIES.key_strict_fifo}' THEN EXECUTE ${schema}.job_table_format($cmd$${createIndexJobPolicyKeyStrictFifo(schema)}$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$${createIndexJobPolicyKeyStrictFifoHeads(schema)}$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$${createCheckConstraintKeyStrictFifo(schema)}$cmd$, tablename); END IF; EXECUTE format('ALTER TABLE ${schema}.%I ADD CONSTRAINT cjc CHECK (name=%L)', tablename, queue_name); EXECUTE format('ALTER TABLE ${schema}.job ATTACH PARTITION ${schema}.%I FOR VALUES IN (%L)', tablename, queue_name); END; $$ LANGUAGE plpgsql; ` } function deleteQueueFunction (schema: string, noPartitioning = false) { const deleteJobsSql = noPartitioning ? `DELETE FROM ${schema}.job WHERE name = queue_name;` : ` SELECT table_name, partition FROM ${schema}.queue WHERE name = queue_name INTO v_table, v_partition; IF v_partition THEN EXECUTE format('DROP TABLE IF EXISTS ${schema}.%I', v_table); ELSE EXECUTE format('DELETE FROM ${schema}.%I WHERE name = %L', v_table, queue_name); END IF; ` const declareBlock = noPartitioning ? '' : ` DECLARE v_table varchar; v_partition bool;` return ` CREATE FUNCTION ${schema}.delete_queue(queue_name text) RETURNS VOID AS $$${declareBlock} BEGIN ${deleteJobsSql} DELETE FROM ${schema}.queue WHERE name = queue_name; END; $$ LANGUAGE plpgsql; ` } export function createQueue (schema: string, name: string, options: unknown, noAdvisoryLocks?: boolean) { const sql = `SELECT ${schema}.create_queue('${name}', '${JSON.stringify(options)}'::jsonb)` return locked(schema, sql, 'create-queue', noAdvisoryLocks) } // LISTEN/NOTIFY channels share a single database-global namespace and are limited to // NAMEDATALEN (63 bytes), unlike the rest of pg-boss which is schema-bound. Derive a // stable, collision-resistant channel from the schema so separate pg-boss instances // (and other services) on the same database never clash. Payload carries the queue name. // // Returns a SQL scalar expression (not a value) hashed in-database with sha224, matching // the convention used by advisoryLock() and partition table naming. Both the producer // (inlined into the insert) and the listener (resolved once at startup) derive the channel // from this single expression, so they always agree. The 'pgboss_' prefix keeps the // channel human-recognizable in pg_stat_activity; 24 hex chars leaves ample headroom under // the 63-byte identifier limit. Channels are already scoped to a single database, so unlike // advisoryLock there is no need to mix in current_database(). // // normalizeSchemaName, not resolveSchemaName: the channel is never matched against the catalog, so // it only has to agree across instances on the same schema. See the note on the helper. export function notifyChannelSql (schema: string): string { return `('pgboss_' || left(encode(sha224('${normalizeSchemaName(schema)}'::bytea), 'hex'), 24))` } // Parameter-less statement that wakes workers on a notify-enabled queue. Embedded into // flow batches so it commits in the same transaction as the inserts. export function notifyQueue (schema: string, name: string): string { return `SELECT pg_notify(${notifyChannelSql(schema)}, '${name}')` } // Dropping a queue's own table (partition: true) needs ACCESS EXCLUSIVE on it, on job and on job_common, // and an insert into job_common or a query through job takes those in other orders. So the locks are // taken first, all with NOWAIT: the transaction never waits holding one, so it can never be part of a // deadlock, and the caller tries again when one is busy. The table is read in the same transaction, so // a stale cache cannot name the wrong one, and a queue already gone is left alone. export function deleteQueue (schema: string, name: string, noAdvisoryLocks?: boolean, partitioned = false) { const sql = partitioned ? ` DO $$ DECLARE v_table text; v_partition bool; BEGIN SELECT table_name, partition FROM ${schema}.queue WHERE name = '${name}' INTO v_table, v_partition; IF NOT FOUND THEN RETURN; END IF; IF v_partition AND to_regclass(format('${schema}.%I', v_table)) IS NOT NULL THEN EXECUTE format('LOCK TABLE ${schema}.${COMMON_JOB_TABLE}, ${schema}.%I, ${schema}.${BASE_JOB_TABLE} IN ACCESS EXCLUSIVE MODE NOWAIT', v_table); END IF; PERFORM ${schema}.delete_queue('${name}'); END $$ ` : `SELECT ${schema}.delete_queue('${name}')` return locked(schema, sql, 'delete-queue', noAdvisoryLocks) } function createPrimaryKeyJob (schema: string) { return `ALTER TABLE ${schema}.job ADD PRIMARY KEY (name, id)` } function createQueueForeignKeyJob (schema: string, noPartitioning = false) { const deferrable = noPartitioning ? '' : ' DEFERRABLE INITIALLY DEFERRED' return `ALTER TABLE ${schema}.job ADD CONSTRAINT q_fkey FOREIGN KEY (name) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT${deferrable}` } function createQueueForeignKeyJobDeadLetter (schema: string, noPartitioning = false) { const deferrable = noPartitioning ? '' : ' DEFERRABLE INITIALLY DEFERRED' return `ALTER TABLE ${schema}.job ADD CONSTRAINT dlq_fkey FOREIGN KEY (dead_letter) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT${deferrable}` } function createIndexJobPolicyShort (schema: string) { return `CREATE UNIQUE INDEX job_i1 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = '${JOB_STATES.created}' AND policy = '${QUEUE_POLICIES.short}'` } function createIndexJobPolicySingleton (schema: string) { return `CREATE UNIQUE INDEX job_i2 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = '${JOB_STATES.active}' AND policy = '${QUEUE_POLICIES.singleton}'` } function createIndexJobPolicyStately (schema: string) { return `CREATE UNIQUE INDEX job_i3 ON ${schema}.job (name, state, COALESCE(singleton_key, '')) WHERE state <= '${JOB_STATES.active}' AND policy = '${QUEUE_POLICIES.stately}'` } function createIndexJobThrottle (schema: string) { return `CREATE UNIQUE INDEX job_i4 ON ${schema}.job (name, singleton_on, COALESCE(singleton_key, '')) WHERE state <> '${JOB_STATES.cancelled}' AND singleton_on IS NOT NULL` } function createIndexJobFetch (schema: string, noCoveringIndex = false) { // Ordered to match the fetch's ORDER BY (priority desc, created_on) // // Two details are load-bearing, and dropping either costs an order of magnitude: // // - `start_after` is a trailing KEY column, not INCLUDE and not absent. A non-leading key column // is still evaluated as an Index Cond, so not-yet-due rows are filtered inside the index. With // it absent the predicate becomes a Filter needing a heap fetch per candidate, and an idle // queue sitting on a scheduled backlog goes 0.04 ms -> 18 ms at 50k deferred, growing // linearly. INCLUDE cannot help: FOR UPDATE forces an Index Scan, never an Index Only Scan. // // - `id` is deliberately NOT a key column. The obvious shape, (name, priority DESC, created_on, // id), also satisfies the ordering, but id is a random uuid, so it defeats btree // deduplication and dominates size: 24 MB against 19 MB for this shape at 500k rows, on the // hot insert path and on every vacuum's bulkdelete. It buys nothing, because fetchNextJob // dropped id from its ORDER BY for the same reason (a random uuid was never creation order); // with the ordering ending at created_on, this index satisfies it outright and the plan // carries no sort node at all. // // Also measured: neutral for key_strict_fifo (that fetch never reaches this index. The // strict_fifo_heads CTE on job_i10 feeds the outer CTE), and a ~7x win on CockroachDB, which // takes the noSkipLocked path and so does not depend on Incremental Sort at all. // // noCoveringIndex is accepted for signature symmetry with the other index builders; this shape // has no covering payload to strip. return `CREATE INDEX job_i11 ON ${schema}.job (name, priority DESC, created_on, start_after) WHERE state < '${JOB_STATES.active}' AND NOT blocked` } function createIndexJobPolicyExclusive (schema: string) { return `CREATE UNIQUE INDEX job_i6 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state <= '${JOB_STATES.active}' AND policy = '${QUEUE_POLICIES.exclusive}'` } function createIndexJobPolicyKeyStrictFifo (schema: string) { return `CREATE UNIQUE INDEX job_i8 ON ${schema}.job (name, singleton_key) WHERE state IN ('${JOB_STATES.active}', '${JOB_STATES.retry}', '${JOB_STATES.failed}') AND policy = '${QUEUE_POLICIES.key_strict_fifo}'` } function createIndexJobPolicyKeyStrictFifoHeads (schema: string, noCoveringIndex = false) { // Column order mirrors the strict_fifo_heads ORDER BY (state DESC puts a retry job ahead // of created siblings) so the DISTINCT ON head-per-key scan stays an ordered index scan. // // INCLUDE (start_after) is load-bearing, not a covering-index habit: `policy` and `blocked` // are satisfied by the partial predicate, so start_after is the only remaining filter in the // heads CTE. Without it every index entry needs a heap fetch and the scan degrades from an // Index Only Scan to an Index Scan (measured at 200k queued rows / 2,000 keys: 47ms -> 111ms, // ~400k extra buffer hits). It costs no HOT updates that job_i11 doesn't already cost, since // job_i11 indexes start_after too, and it only widens a partial index that key_strict_fifo // rows enter. Backends without covering indexes (the CockroachDB profile) fall back to the // narrow form and pay the heap fetches; they take a different fetch path anyway (noSkipLocked). const include = noCoveringIndex ? '' : ' INCLUDE (start_after)' return `CREATE INDEX job_i10 ON ${schema}.job (name, singleton_key, state DESC, created_on, id)${include} WHERE state < '${JOB_STATES.active}' AND NOT blocked AND policy = '${QUEUE_POLICIES.key_strict_fifo}'` } function createCheckConstraintKeyStrictFifo (schema: string) { return `ALTER TABLE ${schema}.job ADD CONSTRAINT job_key_strict_fifo_singleton_key_check CHECK (NOT (policy = '${QUEUE_POLICIES.key_strict_fifo}' AND singleton_key IS NULL))` } function createIndexJobGroupConcurrency (schema: string) { return `CREATE INDEX job_i7 ON ${schema}.job (name, group_id) WHERE state = '${JOB_STATES.active}' AND group_id IS NOT NULL` } // Partial index supporting the background flow resolver (Navigator): lets it find completed // blocking parents with an index scan instead of a partition-wide scan. The `state = completed` // predicate keeps still-running and permanently-failed blocking parents out of the index, so // non-flow queues (and high-partition-count deployments) carry an empty index that costs nothing. function createIndexJobBlocking (schema: string) { return `CREATE INDEX job_i9 ON ${schema}.job (name, id) WHERE blocking AND state = '${JOB_STATES.completed}'` } // Finds every job in a dead letter chain from its root, for a lineage lookup. Keyed on the root // alone, with no name: a chain crosses queues, from the source queue to its dead letter queue and // back. Partial on the column being set, so only jobs that have // been through a dead letter queue are in it and a queue that never dead-letters carries it empty. function createIndexJobSourceRoot (schema: string) { return `CREATE INDEX job_i12 ON ${schema}.job (source_root_id) WHERE source_root_id IS NOT NULL` } // At most one waiting job per key among the jobs upsert() inserted by singletonKey, which is what makes // concurrent upserts of one key agree: the second insert waits on the first's index entry, conflicts, // and edits that job instead. upsert_by_key is set only on those inserts, so send() keeps its own policy's // rules and a queue that never upserts carries the index empty. Only created, not retry, so a failing // job never collides with a newer upsert on its way back to retry. function createIndexJobUpsert (schema: string) { return `CREATE UNIQUE INDEX job_i13 ON ${schema}.job (name, singleton_key) WHERE state = '${JOB_STATES.created}' AND upsert_by_key` } // The interval claim for a monitor pass. It stamps monitor_claim_on, never monitor_on, which only // the stats aggregate writes, so a pass that skips the aggregate leaves capturedOn aging. // // The stats aggregate pins the vacuum horizon, so after a slow one setMonitorBackoff() writes // monitor_backoff_on and no instance starts another until it passes. The backoff comes back as // refreshStats instead of blocking the claim: the pass still runs failJobsByTimeout and // failJobsByHeartbeat, and only the aggregate waits. // // A queue never claimed falls back to monitor_on, so a queue that crossed the v40 upgrade without // its seed (CockroachDB, see noAddColumnBackfill) is not due all at once. A new queue has neither // stamp and is due immediately. export function trySetQueueMonitorTime (schema: string, queues: string[], seconds: number, skipLocked?: boolean): SqlQuery { return { text: ` WITH due AS ( SELECT name FROM ${schema}.queue WHERE name = ANY($1::text[]) AND EXTRACT( EPOCH FROM (${schema}.job_now() - COALESCE(monitor_claim_on, monitor_on, ${schema}.job_now() - interval '1 week') ) ) >= ${seconds} ${queueRowLock(skipLocked)} ) UPDATE ${schema}.queue SET monitor_claim_on = ${schema}.job_now() FROM due WHERE queue.name = due.name RETURNING queue.name, NOT EXISTS (SELECT 1 FROM ${schema}.version WHERE monitor_backoff_on > ${schema}.job_now()) as "refreshStats" `, values: [queues] } } // How much of autovacuum_naptime the stats aggregate may spend pinning the horizon before the // backoff engages. At the default 60s naptime this is 6s: every measured aggregate below that size // (0.19s at 1M rows, 2.1-3.6s at 10M rows / 3GB) leaves a free window of at least 54s out of every // 60s, which no phase alignment of the autovacuum launcher can miss for long. Above it, the free // window stops being reliably longer than a naptime and has to be enforced rather than assumed. const MONITOR_PIN_BUDGET_RATIO = 0.1 // Ceiling on the duration-proportional half of the backoff, so one pathological table cannot stall // the counts indefinitely. Ten minutes already dwarfs any plausible naptime. const MONITOR_BACKOFF_MAX_SECONDS = 600 // Absolute ceiling, for a server whose naptime has itself been tuned into the hours. const MONITOR_BACKOFF_CAP_SECONDS = 3600 // Defer the next stats aggregate long enough for autovacuum to get a clean run at the job table. // // The length follows from when a vacuum decides what it may remove: OldestXmin is computed once, // when the worker starts on the relation, and never revisited. So the requirement is not that a // vacuum *finish* between two aggregates - it is that the instant a worker computes its cutoff // falls in a window with no aggregate running. The launcher retries every autovacuum_naptime, so a // free window strictly longer than one naptime cannot be missed by any phase alignment, while // anything shorter is a probability that two 60s timers can conspire to defeat for hours. // // 2 x naptime the launcher gets two full ticks inside the free window, because its tick is not // the moment the cutoff is taken: the worker still has to connect, build its table // list, and may vacuum other tables first, with autovacuum_max_workers saturated. // >= elapsed a pass that pinned for T seconds is followed by at least T seconds free, capping // the aggregate's duty cycle at 50%. This is PlanetScale's "one concurrent // analytics worker" limit expressed in time rather than in workers. // // naptime is read from the server rather than assumed, so an operator who has tuned it up gets a // proportionally longer backoff - which is the correct answer, not a coincidence. // // Returns a row only when the backoff actually engaged, so the caller warns exactly when it does. // // GREATEST, not assignment: the deadline may only ever move outward. Instances measure their own // passes and write independently, so a marginal pass elsewhere must not cut short a long backoff // already in force - instance A measuring 600s writes a 600s deadline, and instance B measuring 7s // would otherwise overwrite it with 120s a moment later and re-open the window A had closed. // `backoffSeconds` is therefore reported as the time actually remaining on the winning deadline // rather than as the length this pass computed, so the warning never claims a backoff shorter or // longer than the one in force. export function setMonitorBackoff (schema: string, elapsedSeconds: number): SqlQuery { return { text: ` WITH budget AS ( SELECT naptime, LEAST( GREATEST(2 * naptime, LEAST($1::float8, ${MONITOR_BACKOFF_MAX_SECONDS})), ${MONITOR_BACKOFF_CAP_SECONDS} ) AS backoff FROM ( SELECT EXTRACT(EPOCH FROM current_setting('autovacuum_naptime')::interval)::float8 AS naptime ) n ) UPDATE ${schema}.version v SET monitor_backoff_on = GREATEST( COALESCE(v.monitor_backoff_on, ${schema}.job_now()), ${schema}.job_now() + make_interval(secs => b.backoff) ) FROM budget b WHERE $1::float8 > ${MONITOR_PIN_BUDGET_RATIO} * b.naptime RETURNING b.naptime as "naptimeSeconds", EXTRACT(EPOCH FROM (v.monitor_backoff_on - ${schema}.job_now()))::float8 as "backoffSeconds", v.monitor_backoff_on as "backoffUntil" `, values: [elapsedSeconds] } } export function trySetQueueDeletionTime (schema: string, queues: string[], seconds: number, skipLocked?: boolean): SqlQuery { return trySetQueueTimestamp(schema, queues, 'maintain_on', seconds, skipLocked) } // The cron claim, which also answers with the timestamp it replaced and how old that timestamp was. // The timestamp is when an instance last ran a pass, so the pass reads it to find out how long // scheduling was off, which is the window a schedule's `missed` policy catches up over. Null on a // database no pass has ever run against, where there is no gap to catch up on. // // Both answers have to come from a CTE of their own: RETURNING sees the row as the UPDATE leaves // it, while a WITH sub-statement reads the snapshot the whole statement was planned against, so it // sees the value the UPDATE is replacing. // // This answers on both outcomes rather than on a claim alone, which is why `claimed` is a column // and not the row count: an instance that was refused needs to know when the row it lost to comes // due, so its next attempt can be measured from the claim that beat it rather than from its own // failure. See Timekeeper.onCron(). Reading `claim` through EXISTS does not make the UPDATE // conditional - a data-modifying CTE always runs, exactly once, whether or not the outer query // reads it. export function trySetCronTime (schema: string, seconds: number) { return ` WITH prior AS ( SELECT cron_on, EXTRACT( EPOCH FROM (${schema}.job_now() - cron_on) ) as elapsed FROM ${schema}.version ), claim AS ( ${trySetTimestamp(schema, 'cron_on', seconds)} ) SELECT prior.cron_on as "priorCronOn", prior.elapsed as "elapsed", EXISTS (SELECT 1 FROM claim) as "claimed" FROM prior ` } export function trySetBamTime (schema: string, seconds: number) { return trySetTimestamp(schema, 'bam_on', seconds) } export function trySetFlowTime (schema: string, seconds: number) { return trySetTimestamp(schema, 'flow_on', seconds) } export function trySetReindexTime (schema: string, seconds: number) { return trySetTimestamp(schema, 'reindex_on', seconds) } // The claim that decides which instance in a deployment runs an interval pass: whoever moves the // timestamp owns the interval, and everyone else's UPDATE matches nothing. The COALESCE lets the // first claim through while the column is still null. // // The comparison includes the interval itself. On a real clock that is one instant out of a // microsecond-resolution range and changes nothing, but a fake clock has no jitter to hide behind: // TestClock fires each timer exactly one period after the last and moves job_now() with it, so the // elapsed time here is exactly `seconds` on every tick - the one value a strict `>` refuses every // time, which ran a pass on every other tick. Nothing is accepted early either way, so the looser // comparison costs nothing on a real clock. // // It is not what keeps a real deployment's passes on schedule. That is the timer: see ClaimTimer, // which anchors the next attempt to the moment this statement stamps the row rather than to a grid // fixed before it ran. function trySetTimestamp (schema: string, column: string, seconds: number) { return ` UPDATE ${schema}.version SET ${column} = ${schema}.job_now() WHERE EXTRACT( EPOCH FROM (${schema}.job_now() - COALESCE(${column}, ${schema}.job_now() - interval '1 week') ) ) >= ${seconds} RETURNING true ` } function trySetQueueTimestamp (schema: string, queues: string[], column: string, seconds: number, skipLocked?: boolean): SqlQuery { return { text: ` WITH due AS ( SELECT name FROM ${schema}.queue WHERE name = ANY($1::text[]) AND EXTRACT( EPOCH FROM (${schema}.job_now() - COALESCE(${column}, ${schema}.job_now() - interval '1 week') ) ) >= ${seconds} ${queueRowLock(skipLocked)} ) UPDATE ${schema}.queue SET ${column} = ${schema}.job_now() FROM due WHERE queue.name = due.name RETURNING queue.name `, values: [queues] } } // Every statement that writes more than one queue row locks them in name order first. A single // UPDATE locks rows in whatever order its plan visits them, and every write moves a row, so two // multi-row writers - two instances claiming the same queues, or a claim and cacheQueueStats - could // otherwise take the same rows in opposite orders and deadlock. With skipLocked a claim also skips // rows another session holds: a locked row is in another instance's pass, and the next interval // claims it again. Without it the claim waits, which the name order alone keeps deadlock-free. // NO KEY UPDATE is the lock a plain UPDATE of these columns takes; FOR UPDATE would also block the // KEY SHARE lock that inserting a job takes on its queue row through the foreign key. function queueRowLock (skipLocked?: boolean) { return `ORDER BY name FOR NO KEY UPDATE${skipLocked ? ' SKIP LOCKED' : ''}` } // Whether the queue claims may skip locked rows. Only on stock Postgres: the name order already // prevents the deadlock everywhere, and skipping is an optimization the other backends are not // tested with. PGlite has a single connection, so it has nothing to skip. export function queueClaimSkipLocked (backend: string | undefined, noSkipLocked?: boolean): boolean { return backend === 'postgres' && !noSkipLocked } export function updateQueue (schema: string) { return ` WITH options as (SELECT $2::jsonb as data) UPDATE ${schema}.queue SET retry_limit = COALESCE((o.data->>'retryLimit')::int, retry_limit), retry_delay = COALESCE((o.data->>'retryDelay')::int, retry_delay), retry_backoff = COALESCE((o.data->>'retryBackoff')::bool, retry_backoff), retry_delay_max = CASE WHEN (o.data -> 'retryDelayMax') IS NOT NULL THEN (o.data->>'retryDelayMax')::int ELSE retry_delay_max END, expire_seconds = COALESCE((o.data->>'expireInSeconds')::int, expire_seconds), retention_seconds = COALESCE((o.data->>'retentionSeconds')::int, retention_seconds), deletion_seconds = COALESCE((o.data->>'deleteAfterSeconds')::int, deletion_seconds), warning_queued = COALESCE((o.data->>'warningQueueSize')::int, warning_queued), heartbeat_seconds = CASE WHEN (o.data -> 'heartbeatSeconds') IS NOT NULL THEN (o.data->>'heartbeatSeconds')::int ELSE heartbeat_seconds END, notify = COALESCE((o.data->>'notify')::bool, notify), dead_letter = CASE WHEN (o.data -> 'deadLetter') IS NOT NULL THEN o.data->>'deadLetter' ELSE dead_letter END, updated_on = ${schema}.job_now() FROM options o WHERE name = $1 ` } export function currentDatabase () { return 'SELECT current_database() AS name' } export function getQueues (schema: string, names?: string[]): SqlQuery { const hasNames = names && names.length > 0 return { text: ` SELECT q.name, q.policy, q.retry_limit as "retryLimit", q.retry_delay as "retryDelay", q.retry_backoff as "retryBackoff", q.retry_delay_max as "retryDelayMax", q.expire_seconds as "expireInSeconds", q.retention_seconds as "retentionSeconds", q.deletion_seconds as "deleteAfterSeconds", q.partition, q.heartbeat_seconds as "heartbeatSeconds", q.notify, q.dead_letter as "deadLetter", q.deferred_count as "deferredCount", q.blocked_count as "blockedCount", q.warning_queued as "warningQueueSize", q.queued_count as "queuedCount", q.ready_count as "readyCount", q.active_count as "activeCount", q.failed_count as "failedCount", q.total_count as "totalCount", q.created_delta as "createdDelta", q.completed_delta as "completedDelta", q.failed_delta as "failedDelta", q.delta_seconds as "deltaSeconds", q.delta_on as "deltaOn", q.wait_bins as "waitBins", q.run_bins as "runBins", q.ready_oldest_seconds as "readyOldestSeconds", q.singletons_active as "singletonsActive", q.table_name as "table", q.created_on as "createdOn", q.updated_on as "updatedOn" FROM ${schema}.queue q ${hasNames ? 'WHERE q.name = ANY($1::text[])' : ''} `, values: hasNames ? [names] : [] } } export function deleteJobsById (schema: string, table: string, fenced?: boolean) { return ` WITH results as ( DELETE FROM ${schema}.${table} WHERE name = $1 AND id = ANY($2::uuid[]) ${fenced ? attemptFence(3) : ''} RETURNING id ) ${settledCountAndIds()} ` } export function deleteQueuedJobs (schema: string, table: string) { return `DELETE from ${schema}.${table} WHERE name = $1 and state < '${JOB_STATES.active}'` } export function deleteStoredJobs (schema: string, table: string) { return `DELETE from ${schema}.${table} WHERE name = $1 and state > '${JOB_STATES.active}'` } export function truncateTable (schema: string, table: string) { return `TRUNCATE ${schema}.${table}` } // The cached counts of a queue whose table was just truncated, written without counting: they are // zero. Only the truncate paths use it. A DELETE leaves concurrent sends and the other states in // place, so its counts can only come from a recount, whose cost the delete cannot bound. The // throughput counters are the monitor's and are left alone. With `one`, $1 is the queue's name. export function zeroQueueStats (schema: string, one?: boolean) { return ` UPDATE ${schema}.queue SET deferred_count = 0, blocked_count = 0, queued_count = 0, ready_count = 0, active_count = 0, failed_count = 0, total_count = 0, singletons_active = NULL, monitor_on = ${schema}.job_now() FROM ( SELECT name FROM ${schema}.queue${one ? '\n WHERE name = $1' : ''} ${queueRowLock()} ) q WHERE queue.name = q.name ` } export function deleteAllJobs (schema: string, table: string) { return `DELETE from ${schema}.${table} WHERE name = $1` } // Named and aliased rather than SELECT *, so the timestamp columns and last_job_id reach callers // in the same camelCase shape every other read in the API uses. // // The zone is coalesced because the column is nullable and Schedule.timezone is not: schedule() // stores 'UTC' for a zone it was not given, and the v41 migration retires the nulls already in the // table, but an instance on an older release can still write one during a rolling upgrade, and a // row written with SQL can leave it out. Null read back as null would evaluate in the host's local // zone, which no caller ever asked for and which differs between instances. const SCHEDULE_COLUMNS = ` name, key, kind, cron, COALESCE(timezone, 'UTC') as timezone, data, options, created_on as "createdOn", updated_on as "updatedOn", last_job_id as "lastJobId" ` export function getSchedules (schema: string) { return `SELECT ${SCHEDULE_COLUMNS} FROM ${schema}.schedule ORDER BY name, key` } export function getSchedulesByQueue (schema: string) { return `SELECT ${SCHEDULE_COLUMNS} FROM ${schema}.schedule WHERE name = $1 ORDER BY key` } export function getSchedulesByQueueAndKey (schema: string) { return `SELECT ${SCHEDULE_COLUMNS} FROM ${schema}.schedule WHERE name = $1 AND COALESCE(key, '') = $2` } // Records the job each schedule most recently produced. Written from the send-it handler after the // job exists, one statement per batch rather than one per schedule, with the (name, key, job id) // triples carried in a JSON recordset. // // The caller must supply at most one record per (name, key): postgres leaves it unspecified which // source row an UPDATE ... FROM uses when several join the same target, so duplicates would make // the resulting last_job_id arbitrary rather than latest. // // updated_on is deliberately left alone: it tracks edits to the definition, and a firing schedule // has not been edited. export function setScheduleLastJobIds (schema: string) { return ` UPDATE ${schema}.schedule s SET last_job_id = x."jobId" FROM json_to_recordset($1::text::json) AS x (name text, key text, "jobId" uuid) WHERE s.name = x.name AND COALESCE(s.key, '') = x.key ` } /** * Relabels the kind of one or more schedules, as the cron pass does when a row's stored kind * disagrees with the expression beside it. * * `updated_on` is deliberately untouched: it tracks edits to the definition, and a relabel is the * pass agreeing with what the row already said, not a change to what the schedule does. * * Matched on the expression as well as the primary key, so a schedule() upsert landing between the * pass's read and this write keeps the kind it was stored with instead of being stamped with the * previous expression's. A row that misses the update because it changed underneath is relabelled * by the next pass if it still needs it. */ export function setScheduleKinds (schema: string) { return ` UPDATE ${schema}.schedule s SET kind = k.kind FROM json_to_recordset($1::text::json) as k (name text, key text, kind text, cron text) WHERE s.name = k.name AND COALESCE(s.key, '') = k.key AND s.cron = k.cron ` } export function schedule (schema: string) { return ` INSERT INTO ${schema}.schedule (name, key, kind, cron, timezone, data, options, created_on, updated_on) VALUES ($1, $2, $3, $4, $5, $6, $7, ${schema}.job_now(), ${schema}.job_now()) ON CONFLICT (name, key) DO UPDATE SET kind = EXCLUDED.kind, cron = EXCLUDED.cron, timezone = EXCLUDED.timezone, data = EXCLUDED.data, options = EXCLUDED.options, updated_on = ${schema}.job_now() ` } export function unschedule (schema: string) { return ` DELETE FROM ${schema}.schedule WHERE name = $1 AND COALESCE(key, '') = $2 ` } export function subscribe (schema: string) { return ` INSERT INTO ${schema}.subscription (event, name, created_on, updated_on) VALUES ($1, $2, ${schema}.job_now(), ${schema}.job_now()) ON CONFLICT (event, name) DO UPDATE SET event = EXCLUDED.event, name = EXCLUDED.name, updated_on = ${schema}.job_now() ` } export function unsubscribe (schema: string) { return ` DELETE FROM ${schema}.subscription WHERE event = $1 and name = $2 ` } export function getQueuesForEvent (schema: string) { return ` SELECT name FROM ${schema}.subscription WHERE event = $1 ` } export function getTime (schema: string) { return `SELECT round(date_part('epoch', ${schema}.job_now()) * 1000) as time` } export function insertWarning (schema: string) { return ` INSERT INTO ${schema}.warning (type, message, data, created_on) VALUES ($1, $2, $3::text::jsonb, ${schema}.job_now()) ` } export function getWarnings (schema: string): string { return ` SELECT id, type, message, data, created_on as "createdOn" FROM ${schema}.warning WHERE ($1::text IS NULL OR type = $1) ORDER BY created_on DESC LIMIT $2 OFFSET $3 ` } export function getWarningsCount (schema: string): string { return ` SELECT COUNT(*)::int as count FROM ${schema}.warning WHERE ($1::text IS NULL OR type = $1) ` } export function deleteOldWarnings (schema: string, days: number): string { return ` DELETE FROM ${schema}.warning WHERE created_on < ${schema}.job_now() - interval '${days} days' ` } // The delta columns are nullable, unlike the gauges: a snapshot captured before they were counted // has no value, and zero would chart as an idle queue. Passes are not evenly spaced, so a rate is // sum(delta) / sum(delta_seconds), never delta / bucket width. delta_on trails captured_on by // DELTA_LAG, since the window a pass counts ends that far behind it. /* eslint-disable no-restricted-syntax -- column defaults stay on the real clock: every pg-boss write names its timestamps through job_now() */ export function createTableQueueStats (schema: string, noPartitioning = false): string { return ` CREATE TABLE ${schema}.queue_stats ( id uuid NOT NULL DEFAULT gen_random_uuid(), name text NOT NULL, deferred_count int NOT NULL DEFAULT 0, queued_count int NOT NULL DEFAULT 0, ready_count int NOT NULL DEFAULT 0, active_count int NOT NULL DEFAULT 0, failed_count int NOT NULL DEFAULT 0, total_count int NOT NULL DEFAULT 0, created_delta int, completed_delta int, failed_delta int, delta_seconds int, delta_on timestamptz, wait_bins int[], run_bins int[], ready_oldest_seconds int, captured_on timestamptz NOT NULL DEFAULT now(), ${noPartitioning ? 'PRIMARY KEY (id)' : 'PRIMARY KEY (id, captured_on)'} ) ${noPartitioning ? '' : 'PARTITION BY RANGE (captured_on)'} ` } /* eslint-enable no-restricted-syntax */ export function createIndexQueueStats (schema: string, noCoveringIndex = false): string { // The three delta columns are deliberately *not* included. Adding them would // make a fresh install's index differ from a migrated one unless the migration // rebuilt it, and rebuilding the covering index on a partitioned table that is // large on exactly the installations that care about throughput is a heavy // price for making three int columns index-only. The lookup still uses this // index; it just visits the heap for those three values. const include = noCoveringIndex ? '' : 'INCLUDE (deferred_count, queued_count, ready_count, active_count, failed_count, total_count)' return `CREATE INDEX queue_stats_i1 ON ${schema}.queue_stats (name, captured_on DESC) ${include}` } // Idempotently create the daily partitions for today and tomorrow (UTC). Both the day suffix and // the range bounds are derived in SQL from the UTC calendar date, and the bounds are emitted as // explicit `+00` timestamptz literals. This keeps partitioning correct regardless of the database // session TimeZone (a bare date literal like '2026-06-25' would otherwise be cast to timestamptz in // the session TZ, so rows written near UTC midnight could fall outside every existing partition). // Computing the date in SQL (rather than interpolating new Date()) also keeps emitted DDL, including // the v35 migration and exported create plans, deterministic and apply-time accurate. export function ensureQueueStatsPartitions (schema: string): string { return ` DO $$ DECLARE d date; i int; part_name text; BEGIN FOR i IN 0..1 LOOP d := (${schema}.job_now() AT TIME ZONE 'UTC')::date + i; part_name := 'queue_stats_' || to_char(d, 'YYYYMMDD'); IF NOT EXISTS ( SELECT 1 FROM pg_class c JOIN pg_namespace n ON n.oid = c.relnamespace WHERE n.nspname = '${resolveSchemaName(schema)}' AND c.relname = part_name ) THEN EXECUTE format( 'CREATE TABLE ${schema}.%I PARTITION OF ${schema}.queue_stats FOR VALUES FROM (%L) TO (%L)', part_name, to_char(d, 'YYYY-MM-DD') || ' 00:00:00+00', to_char(d + 1, 'YYYY-MM-DD') || ' 00:00:00+00' ); END IF; END LOOP; END; $$ ` } export function dropOldQueueStatsPartitions (schema: string, days: number): string { return ` DO $$ DECLARE r record; cutoff date := (${schema}.job_now() AT TIME ZONE 'UTC')::date - ${days}; suffix text; part_date date; BEGIN FOR r IN SELECT c.relname FROM pg_inherits i JOIN pg_class p ON p.oid = i.inhparent JOIN pg_class c ON c.oid = i.inhrelid JOIN pg_namespace n ON n.oid = p.relnamespace WHERE n.nspname = '${resolveSchemaName(schema)}' AND p.relname = 'queue_stats' LOOP suffix := substring(r.relname FROM 'queue_stats_(.*)$'); IF suffix ~ '^[0-9]{8}$' THEN part_date := to_date(suffix, 'YYYYMMDD'); IF part_date < cutoff THEN EXECUTE 'DROP TABLE IF EXISTS ${schema}.' || quote_ident(r.relname); END IF; END IF; END LOOP; END; $$ ` } export function deleteOldQueueStats (schema: string, days: number): string { return ` DELETE FROM ${schema}.queue_stats WHERE captured_on < ${schema}.job_now() - interval '${days} days' ` } export function insertQueueStats (schema: string, queues: string[], noAdvisoryLocks?: boolean): string { const sql = ` INSERT INTO ${schema}.queue_stats (name, deferred_count, queued_count, ready_count, active_count, failed_count, total_count, created_delta, completed_delta, failed_delta, delta_seconds, delta_on, wait_bins, run_bins, ready_oldest_seconds, captured_on) SELECT name, deferred_count, queued_count, ready_count, active_count, failed_count, total_count, created_delta, completed_delta, failed_delta, delta_seconds, delta_on, wait_bins, run_bins, ready_oldest_seconds, ${schema}.job_now() FROM ${schema}.queue WHERE name = ANY(${serializeArrayParam(queues)}) ` return locked(schema, sql, 'queue-stats-insert', noAdvisoryLocks) } // Cheap single-row read of the cached counts the monitor maintains on the queue table. capturedOn // is monitor_on. The moment those counts were last refreshed, or NULL if the queue has never been // monitored (so the caller knows to recompute rather than trust default-zero counts). // // cacheAgeMs is computed here rather than from capturedOn on the client: pg-boss shares the // caller's pool, and a global pg-types parser (Temporal.Instant, say) returns timestamptz values // `new Date()` cannot coerce. // // monitorBackoff rides along because the caller must serve this cache even when it is stale while // the vacuum-safety backoff is in force: a forced refresh runs the same whole-table aggregate the // backoff exists to space out, and a dashboard polling { force: true } would otherwise walk // straight back into the pin the monitor just backed away from. export function getQueueStatsCache (schema: string): string { return ` SELECT name, deferred_count as "deferredCount", queued_count as "queuedCount", ready_count as "readyCount", active_count as "activeCount", failed_count as "failedCount", total_count as "totalCount", created_delta as "createdDelta", completed_delta as "completedDelta", failed_delta as "failedDelta", delta_seconds as "deltaSeconds", delta_on as "deltaOn", wait_bins as "waitBins", run_bins as "runBins", ready_oldest_seconds as "readyOldestSeconds", table_name as "table", monitor_on as "capturedOn", (extract(epoch from (${schema}.job_now() - monitor_on)) * 1000)::float8 as "cacheAgeMs", (SELECT monitor_backoff_on > ${schema}.job_now() FROM ${schema}.version) as "monitorBackoff" FROM ${schema}.queue WHERE name = $1 ` } export function getQueueStatsHistory (schema: string): string { return ` SELECT name, deferred_count as "deferredCount", queued_count as "queuedCount", ready_count as "readyCount", active_count as "activeCount", failed_count as "failedCount", total_count as "totalCount", created_delta as "createdDelta", completed_delta as "completedDelta", failed_delta as "failedDelta", delta_seconds as "deltaSeconds", delta_on as "deltaOn", wait_bins as "waitBins", run_bins as "runBins", ready_oldest_seconds as "readyOldestSeconds", captured_on as "capturedOn" FROM ${schema}.queue_stats WHERE name = $1 AND ($2::timestamptz IS NULL OR captured_on >= $2) AND ($3::timestamptz IS NULL OR captured_on <= $3) ORDER BY captured_on DESC LIMIT $4 ` } // Per-bucket aggregate over a count column. The function name can't be a bind parameter, so it's // interpolated, safe because the manager validates `aggregate` against this whitelist first. Every // result is cast back to int: it honors the int count contract (avg rounds) and keeps Postgres // returning the value as a JS number rather than a numeric string. const STATS_AGG = { max: (c: string) => `max(${c})::int`, min: (c: string) => `min(${c})::int`, avg: (c: string) => `round(avg(${c}))::int` } as const // Downsampled history: group the recorded series into fixed-width time buckets and collapse each // bucket's counts with `aggregate`, so a wide window returns a manageable, representative sample // instead of just the newest `limit` raw rows. // // mode 'bucket', $5 is the bucket width in seconds (explicit resolution). // mode 'auto', $5 is maxDataPoints; the width is derived so the series fits in $5 points. // from/to sets the range, but they cannot exceed the data's own min/max values. // // The three delta columns are summed rather than passed through `aggregate`, and // that is not an oversight. Every other column here is a gauge, where the // question a wider bucket asks is "how high did it get" or "what was it // typically"; these three are counters, where the only meaningful answer is "how // many in total". Averaging them would report a rate per capture interval // labelled as a count, which reads plausible and is wrong by whatever the // bucket width happens to be. // // They are also bucketed by a different time. A gauge belongs to the moment it was read, // captured_on. The counters in that same row describe an interval that ended DELTA_LAG earlier, // delta_on, so keyed on captured_on they would land a bucket or two late and sit beside the wrong // gauges. Each side is bucketed by its own time and the two are joined per bucket, which is why the // newest bucket carries null counters: the pass that counts its interval has not run yet, and a // later read fills it in. // // The join is not symmetric, though. A counter's bucket can hold no capture at all: passes drift, // so delta_on (a pass time minus the lag) lands a bucket away from the previous capture; a bucket // narrower than the monitor interval has more empty buckets than full ones. Emitting that bucket on its own // would chart the queue as empty, since a gauge with no capture reads as zero. So a counter bucket // with no gauges is folded into the newest gauge bucket at or before it (the capture at the end of // the counted interval reflects the state it left), or into the first gauge bucket in range when // none precedes it. Every bucket returned has a capture behind its gauges. // // The bucket key avoids date_bin() (PG14+): pg-boss supports PostgreSQL 13+ and CockroachDB/ // YugabyteDB, none of which can rely on it. to_timestamp / extract(epoch) / floor exist on all of // them (extract returns double on PG13, numeric on PG14+; floor/division handle both identically), // and buckets align to the Unix epoch so their boundaries are stable across calls. // // Wait and run histograms are added up per returned bucket in SQL (passes, slots, histograms), so a // bucket comes back as one histogram however many passes it covers. Each pass lands in the bucket // its counters were placed in, then its counts are summed per bucket and slot, a slot with no job in // any pass as 0. A bucket no pass measured has no histogram row and comes back null. export function getQueueStatsHistoryBucketed (schema: string, aggregate: 'max' | 'min' | 'avg', mode: 'bucket' | 'auto'): string { const agg = STATS_AGG[aggregate] const widthCte = mode === 'auto' ? `WITH extent AS ( SELECT min(captured_on) AS lo, max(captured_on) AS hi FROM ${schema}.queue_stats WHERE name = $1 ), bounds AS ( SELECT greatest(coalesce($2::timestamptz, lo), lo) AS lo, least(coalesce($3::timestamptz, hi), hi) AS hi FROM extent ), w AS ( SELECT greatest(1, ceil(extract(epoch from (hi - lo))::float8 / greatest($5, 1)::float8)::bigint) AS secs FROM bounds )` : 'WITH w AS (SELECT greatest($5, 1)::bigint AS secs)' // Hard-cap auto-mode at maxDataPoints. Epoch-aligned bucketing can straddle a boundary and emit // one bucket more than the target, so cap the row count at the smaller of the user's limit and // maxDataPoints. ORDER BY DESC means the cap drops the oldest (straddle) bucket and keeps the // newest N. Explicit bucketSeconds has no target to overshoot, so it keeps the raw limit. const limit = mode === 'auto' ? 'least($4, $5)' : '$4' // float8 on both sides: CockroachDB has no float / int operator, and extract() is a float there. const bucket = (column: string) => `to_timestamp(floor(extract(epoch from ${column})::float8 / w.secs::float8) * w.secs::float8)` return ` ${widthCte}, gauges AS ( SELECT ${bucket('captured_on')} as bucket, ${agg('deferred_count')} as "deferredCount", ${agg('queued_count')} as "queuedCount", ${agg('ready_count')} as "readyCount", ${agg('active_count')} as "activeCount", ${agg('failed_count')} as "failedCount", ${agg('total_count')} as "totalCount" FROM ${schema}.queue_stats, w WHERE name = $1 AND ($2::timestamptz IS NULL OR captured_on >= $2) AND ($3::timestamptz IS NULL OR captured_on <= $3) GROUP BY 1 ), counters AS ( SELECT ${bucket('delta_on')} as bucket, sum(created_delta)::int as "createdDelta", sum(completed_delta)::int as "completedDelta", sum(failed_delta)::int as "failedDelta", sum(delta_seconds)::int as "deltaSeconds", max(delta_on) as "deltaOn", max(ready_oldest_seconds) as "readyOldestSeconds" FROM ${schema}.queue_stats, w WHERE name = $1 AND delta_on IS NOT NULL AND ($2::timestamptz IS NULL OR delta_on >= $2) AND ($3::timestamptz IS NULL OR delta_on <= $3) GROUP BY 1 ), placed AS ( SELECT COALESCE( max(g.bucket) OVER (ORDER BY COALESCE(g.bucket, c.bucket) ROWS UNBOUNDED PRECEDING), min(g.bucket) OVER () ) as bucket, g."deferredCount", g."queuedCount", g."readyCount", g."activeCount", g."failedCount", g."totalCount", c."createdDelta", c."completedDelta", c."failedDelta", c."deltaSeconds", c."deltaOn", c."readyOldestSeconds", c.bucket as "counterBucket" FROM gauges g FULL JOIN counters c ON c.bucket = g.bucket ), passes AS ( SELECT ${bucket('delta_on')} as "counterBucket", wait_bins, run_bins FROM ${schema}.queue_stats, w WHERE name = $1 AND delta_on IS NOT NULL AND wait_bins IS NOT NULL AND ($2::timestamptz IS NULL OR delta_on >= $2) AND ($3::timestamptz IS NULL OR delta_on <= $3) ), slots AS ( SELECT p.bucket, u.slot, coalesce(sum(u.w), 0)::int AS w, coalesce(sum(u.r), 0)::int AS r FROM (SELECT DISTINCT bucket, "counterBucket" FROM placed) p JOIN passes ps ON ps."counterBucket" = p."counterBucket", unnest(ps.wait_bins, ps.run_bins) WITH ORDINALITY AS u(w, r, slot) GROUP BY 1, 2 ), histograms AS ( SELECT bucket, array_agg(w ORDER BY slot) as "waitBins", array_agg(r ORDER BY slot) as "runBins" FROM slots GROUP BY 1 ), folded AS ( SELECT bucket as "capturedOn", max("deferredCount")::int as "deferredCount", max("queuedCount")::int as "queuedCount", max("readyCount")::int as "readyCount", max("activeCount")::int as "activeCount", max("failedCount")::int as "failedCount", max("totalCount")::int as "totalCount", sum("createdDelta")::int as "createdDelta", sum("completedDelta")::int as "completedDelta", sum("failedDelta")::int as "failedDelta", sum("deltaSeconds")::int as "deltaSeconds", max("deltaOn") as "deltaOn", max("readyOldestSeconds")::int as "readyOldestSeconds" FROM placed WHERE bucket IS NOT NULL GROUP BY 1 ) SELECT f.*, h."waitBins", h."runBins" FROM folded f LEFT JOIN histograms h ON h.bucket = f."capturedOn" ORDER BY f."capturedOn" DESC LIMIT ${limit} ` } export function getVersion (schema: string) { return `SELECT version from ${schema}.version` } export function setVersion (schema: string, version: number) { return `UPDATE ${schema}.version SET version = '${version}'` } export function versionTableExists (schema: string) { return `SELECT to_regclass('${schema}.version') as name` } // Installed pg-boss schemas whose name differs from the configured one by case alone. Postgres // folds a bare name and stores a quoted one verbatim, so `MySchema` and `"MySchema"` are two // schemas that look nearly identical in config. Used on the install path to tell a caller who // mis-spelled the quoting that their data is next door, rather than silently installing a second, // empty schema beside the populated one. export function getSchemaCaseVariants (schema: string): string { const resolved = resolveSchemaName(schema).replace(SINGLE_QUOTE_REGEX, "''") return ` SELECT n.nspname as name FROM pg_namespace n JOIN pg_class c ON c.relnamespace = n.oid AND c.relname = 'version' AND c.relkind IN ('r', 'p') WHERE lower(n.nspname) = lower('${resolved}') AND n.nspname <> '${resolved}' ORDER BY n.nspname ` } export function getPartitionedQueueTables (schema: string) { return `SELECT table_name, policy FROM ${schema}.queue WHERE partition = true` } export function insertVersion (schema: string, version: number) { return `INSERT INTO ${schema}.version(version) VALUES ('${version}')` } interface GroupConcurrencyConfig { default: number tiers?: Record } interface FetchJobOptions { schema: string table: string name: string policy: string | undefined limit: number includeMetadata?: boolean // Returns the job's stored trace context as "__traceContext", which Manager strips off the job. includeTraceContext?: boolean priority?: boolean orderByCreatedOn?: boolean ignoreStartAfter?: boolean ignoreSingletons: string[] | null ignoreGroups?: string[] | null groupConcurrency?: number | GroupConcurrencyConfig minPriority?: number maxPriority?: number } interface FetchQueryParams { values: unknown[] ignoreSingletonsParam: string ignoreGroupsParam: string defaultGroupLimitParam: string tiersParam: string minPriorityParam: string maxPriorityParam: string } function buildFetchParams (options: FetchJobOptions): FetchQueryParams { const { ignoreSingletons, ignoreGroups, groupConcurrency, minPriority, maxPriority } = options const hasIgnoreSingletons = ignoreSingletons != null && ignoreSingletons.length > 0 const hasIgnoreGroups = ignoreGroups != null && ignoreGroups.length > 0 const hasGroupConcurrency = groupConcurrency != null const hasMinPriority = minPriority != null const hasMaxPriority = maxPriority != null const groupConcurrencyConfig = hasGroupConcurrency ? (typeof groupConcurrency === 'number' ? { default: groupConcurrency } : groupConcurrency) : null const hasTiers = groupConcurrencyConfig?.tiers && Object.keys(groupConcurrencyConfig.tiers).length > 0 const values: unknown[] = [] let paramIndex = 0 let ignoreSingletonsParam = '' let ignoreGroupsParam = '' let defaultGroupLimitParam = '' let tiersParam = '' let minPriorityParam = '' let maxPriorityParam = '' if (hasIgnoreSingletons) { paramIndex++ ignoreSingletonsParam = `$${paramIndex}::text[]` // job_i2/job_i3 key singleton/stately jobs on the empty key (COALESCE(singleton_key, '')) // as one slot, so a keyless active job must block keyless pending jobs the same way a keyed // one blocks its key. Map null -> '' here so the WHERE clause's COALESCE comparison (below) // never has to compare against a literal NULL array element, which would make `<> ALL(...)` // evaluate to NULL (excluding every row) instead of the intended per-key filter. values.push(ignoreSingletons.map(key => key ?? '')) } if (hasIgnoreGroups) { paramIndex++ ignoreGroupsParam = `$${paramIndex}::text[]` values.push(ignoreGroups) } if (hasGroupConcurrency && groupConcurrencyConfig) { paramIndex++ defaultGroupLimitParam = `$${paramIndex}::int` values.push(groupConcurrencyConfig.default) if (hasTiers) { paramIndex++ tiersParam = `$${paramIndex}::text::jsonb` values.push(JSON.stringify(groupConcurrencyConfig.tiers)) } } if (hasMinPriority) { paramIndex++ minPriorityParam = `$${paramIndex}::int` values.push(minPriority) } if (hasMaxPriority) { paramIndex++ maxPriorityParam = `$${paramIndex}::int` values.push(maxPriority) } return { values, ignoreSingletonsParam, ignoreGroupsParam, defaultGroupLimitParam, tiersParam, minPriorityParam, maxPriorityParam } } /** * Builds the fetch query for claiming jobs from the queue. * * With SKIP LOCKED (noSkipLocked=false, the default), uses SELECT FOR UPDATE SKIP * LOCKED, which lets multiple workers efficiently fetch different jobs simultaneously. * * With noSkipLocked=true, omits FOR UPDATE SKIP LOCKED and adds an additional state * check in the WHERE clause. This pattern works better with distributed databases like * CockroachDB where SKIP LOCKED has performance issues and can unexpectedly skip * unlocked rows. * * Trade-off when noSkipLocked is set: under high contention, workers may receive fewer * jobs per fetch as concurrent updates to the same rows will result in some workers * getting empty results. This is acceptable for job queues where processing time * exceeds fetch time. */ export function fetchNextJob (options: FetchJobOptions, noSkipLocked = false): SqlQuery { const { schema, table, name, policy, limit, includeMetadata, includeTraceContext = false, ignoreStartAfter = false, groupConcurrency, minPriority, maxPriority } = options const keyStrictFifo = policy === QUEUE_POLICIES.key_strict_fifo const singletonFetch = limit > 1 && (policy === QUEUE_POLICIES.singleton || policy === QUEUE_POLICIES.stately) const hasIgnoreSingletons = options.ignoreSingletons != null && options.ignoreSingletons.length > 0 const hasIgnoreGroups = options.ignoreGroups != null && options.ignoreGroups.length > 0 const hasGroupConcurrency = groupConcurrency != null const hasMinPriority = minPriority != null const hasMaxPriority = maxPriority != null const groupConcurrencyConfig = hasGroupConcurrency ? (typeof groupConcurrency === 'number' ? { default: groupConcurrency } : groupConcurrency) : null const hasTiers = hasGroupConcurrency && groupConcurrencyConfig?.tiers && Object.keys(groupConcurrencyConfig.tiers).length > 0 const hasSingleGroupConcurrency = hasGroupConcurrency && !hasTiers && groupConcurrencyConfig?.default === 1 const hasActiveGroupCounts = hasGroupConcurrency && !hasSingleGroupConcurrency const params = buildFetchParams(options) const groupLimit = hasTiers ? `COALESCE((${params.tiersParam} ->> group_tier)::int, ${params.defaultGroupLimitParam})` : params.defaultGroupLimitParam const activeGroupCountExpression = hasActiveGroupCounts ? 'COALESCE(((SELECT counts FROM active_group_count_map) ->> j.group_id)::int, 0)' : '' const selectCols = [ 'j.id', singletonFetch ? 'j.singleton_key' : '', hasGroupConcurrency ? 'j.group_id, j.group_tier' : '', hasActiveGroupCounts ? `${activeGroupCountExpression} as active_cnt` : '' ].filter(Boolean).join(', ') // For limits above 1, aggregate active counts into a single JSONB value. Each // candidate uses a keyed lookup through an uncorrelated InitPlan, so the planner // cannot turn stale active-group estimates into a per-candidate relation scan. const activeGroupCountMapCte = hasActiveGroupCounts ? `active_group_count_map AS MATERIALIZED ( SELECT COALESCE(jsonb_object_agg(group_id, active_cnt), '{}'::jsonb) as counts FROM ( SELECT group_id, COUNT(*)::int as active_cnt FROM ${schema}.${table} WHERE name = '${name}' AND state = '${JOB_STATES.active}' AND group_id IS NOT NULL GROUP BY group_id ) active_groups ), ` : '' // With noSkipLocked, omit FOR UPDATE SKIP LOCKED as it performs poorly // in distributed databases like CockroachDB const lockClause = noSkipLocked ? '' : 'FOR UPDATE OF j SKIP LOCKED' // Column references are qualified with j. throughout so both the base case and // the groupConcurrency branches share one set of expressions. const groupConcurrencyFilter = hasGroupConcurrency ? hasSingleGroupConcurrency ? `(j.group_id IS NULL OR NOT EXISTS ( SELECT 1 FROM ${schema}.${table} active_group_probe WHERE active_group_probe.name = '${name}' AND active_group_probe.state = '${JOB_STATES.active}' AND active_group_probe.group_id IS NOT NULL AND active_group_probe.group_id = j.group_id ))` : `(j.group_id IS NULL OR ${activeGroupCountExpression} < ${groupLimit})` : '' // state DESC (retry > created in the enum) makes a retry job its key's head. job_i8 // guarantees at most one job per key in active/retry/failed, and the NOT EXISTS blocker // below rejects every sibling of a retry job, so if created_on picked the head, an older // deferred job whose start_after has since arrived would claim the head slot while being // unfetchable, and the retry job (fetchable but not the head) would deadlock the key. const strictFifoHeadsCte = keyStrictFifo ? `strict_fifo_heads AS MATERIALIZED ( SELECT DISTINCT ON (h.singleton_key) h.id FROM ${schema}.${table} h WHERE h.name = '${name}' AND h.state < '${JOB_STATES.active}' AND NOT h.blocked AND h.policy = '${QUEUE_POLICIES.key_strict_fifo}' ${!ignoreStartAfter ? `AND h.start_after <= ${schema}.job_now()` : ''} ORDER BY h.singleton_key, h.state DESC, h.created_on, h.id ), ` : '' const whereConditions = [ `j.name = '${name}'`, `j.state < '${JOB_STATES.active}'`, 'NOT j.blocked', // `<=` (not `<`) so a job inserted with start_after = ${schema}.job_now() is immediately // fetchable in the next statement. `${schema}.job_now()` is transaction-scoped; on backends with coarse // clock resolution (notably PGlite) consecutive autocommit statements often share the same // timestamp, so `<` would leave freshly-inserted jobs invisible until the clock ticks. // NOTIFY gating already uses `start_after <= ${schema}.job_now()` for the same reason. !ignoreStartAfter ? `j.start_after <= ${schema}.job_now()` : '', keyStrictFifo ? 'j.id IN (SELECT id FROM strict_fifo_heads)' : '', keyStrictFifo ? `NOT EXISTS ( SELECT 1 FROM ${schema}.${table} b WHERE b.name = j.name AND b.singleton_key = j.singleton_key AND b.state IN ('${JOB_STATES.active}', '${JOB_STATES.retry}', '${JOB_STATES.failed}') AND b.policy = '${QUEUE_POLICIES.key_strict_fifo}' AND b.id <> j.id )` : '', hasIgnoreSingletons ? `COALESCE(j.singleton_key, '') <> ALL(${params.ignoreSingletonsParam})` : '', hasIgnoreGroups ? `(j.group_id IS NULL OR j.group_id <> ALL(${params.ignoreGroupsParam}))` : '', hasMinPriority ? `j.priority >= ${params.minPriorityParam}` : '', hasMaxPriority ? `j.priority <= ${params.maxPriorityParam}` : '', groupConcurrencyFilter ].filter(Boolean).join('\n AND ') // Unconditional, and with no `id` tiebreak. // // The two conditionals this replaces existed only to serve the `priority` and `orderByCreatedOn` // fetch options, which were requested as performance escapes from this sort. job_i11 now satisfies // the ordering directly, so the escapes have nothing to escape. See the deprecation in // manager.fetch(). Dropping them also deletes two branches from the hottest SQL builder here. // // `id` goes too. It is a random uuid, so it never provided creation order: a batch insert shares // one ${schema}.job_now() and therefore ties on created_on, and those ties resolve today in random uuid order, // not insertion order. Without it the index satisfies the ordering outright rather than through // an Incremental Sort (0.097 -> 0.031 ms and 50 -> 5 buffers at limit=1) and ties fall back to // index order, which tracks insertion order better than a uuid does. // // key_strict_fifo keeps its own id tiebreak in strict_fifo_heads, where DISTINCT ON does need a // total order to pick a deterministic head per key. const nextCte = ` next AS ( SELECT ${selectCols} FROM ${schema}.${table} j WHERE ${whereConditions} ORDER BY j.priority desc, j.created_on LIMIT ${limit} ${lockClause} )` const singletonCte = singletonFetch ? `, singleton_ranking AS ( SELECT id, ${hasGroupConcurrency ? 'group_id, group_tier, ' : ''}${hasActiveGroupCounts ? 'active_cnt, ' : ''} row_number() OVER (PARTITION BY singleton_key) as singleton_rn FROM next )` : '' const groupConcurrencyCtes = hasGroupConcurrency ? `, group_ranking AS ( SELECT t.id , t.group_id , t.group_tier ${singletonFetch ? ', singleton_rn' : ''} , ROW_NUMBER() OVER (PARTITION BY t.group_id ORDER BY t.id) as group_rn , ${hasActiveGroupCounts ? 't.active_cnt' : '0'} as active_cnt FROM ${singletonFetch ? 'singleton_ranking' : 'next'} t ${singletonFetch ? 'WHERE singleton_rn = 1' : ''} ), group_filtered AS ( SELECT id FROM group_ranking WHERE group_id IS NULL OR (active_cnt + group_rn) <= ${groupLimit} )` : '' const finalCte = hasGroupConcurrency ? 'group_filtered' : (singletonFetch) ? 'singleton_ranking' : 'next' // An uncorrelated array InitPlan makes the selected ids a one-time input to the // UPDATE. Without it, stale estimates can make Postgres put the inlined ranking // query on the inner side of a nested loop and execute it once per job table row. const updateSource = hasGroupConcurrency ? '' : `FROM ${finalCte}` const updateMatch = hasGroupConcurrency ? `j.id = ANY (ARRAY(SELECT id FROM ${finalCte}))` : `j.id = ${finalCte}.id` // Without SKIP LOCKED, add a state check to prevent duplicate claims // when multiple workers try to claim the same jobs concurrently const distributedStateCheck = noSkipLocked ? `AND j.state < '${JOB_STATES.active}'` : '' return { text: ` WITH ${strictFifoHeadsCte} ${activeGroupCountMapCte} ${nextCte} ${singletonCte} ${groupConcurrencyCtes} UPDATE ${schema}.${table} j SET state = '${JOB_STATES.active}', started_on = ${schema}.job_now(), heartbeat_on = ${schema}.job_now(), retry_count = CASE WHEN started_on IS NOT NULL THEN retry_count + 1 ELSE retry_count END ${updateSource} WHERE name = '${name}' AND ${updateMatch} ${singletonFetch && !hasGroupConcurrency ? 'AND singleton_rn = 1' : ''} ${distributedStateCheck} RETURNING j.${includeMetadata ? JOB_COLUMNS_ALL : JOB_COLUMNS_MIN}${includeTraceContext ? ', j.trace_context as "__traceContext"' : ''} `, values: params.values } } // The attempt fence (#925). A worker whose claim lapsed (a heartbeat the database stopped seeing, an // expiry) keeps running its handler, and by the time it settles, the job can be `active` again under // another worker's claim. Matching on id and state alone lets the stale settle land on that newer // attempt: its complete() overwrites the newer attempt's outcome, and its fail() deletes the live // attempt and re-queues it. Every claim increments retry_count (fetchNextJob), so the value a worker // fetched names its attempt. `param` is a text[] of 'id:retry_count' pairs (attemptPairs), and a // job whose retry_count has moved on is left alone. Always paired with state = 'active': a job failed // back to `retry` still carries the stale attempt's retry_count until it is claimed again. // // The pair is matched with = ANY because Postgres hashes that lookup against a parameter array, so a // settle costs the same fenced as unfenced. Measured at 5000 jobs: 58ms against 54ms unfenced. Looking // each id's position up with array_position scans the array once per row (149ms), and joining an // unnest() of the two arrays is planned as a nested loop over every pair (1063ms). // // restore() clears started_on, so the claim after it does not increment. The fence cannot tell a // restored job's next claim from the one before it. pg-boss itself never restores a job whose handler // is still running (the localGroupConcurrency excess is restored before its handler starts), but a // caller of the public restore() can. function attemptFence (param: number): string { return `AND state = '${JOB_STATES.active}' AND (id::text || ':' || retry_count::text) = ANY($${param}::text[])` } // The attemptFence parameter for jobs `ids` fetched at `attempts` (parallel arrays). export function attemptPairs (ids: string[], attempts: number[]): string[] { return ids.map((id, i) => `${canonicalUuid(id)}:${attempts[i]}`) } const CANONICAL_UUID = /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/ // The fence compares id::text, which Postgres and CockroachDB always render lower case and hyphenated, // while the id lookup beside it takes a uuid parameter and so accepts any spelling the uuid type does // (upper case, braces, missing hyphens). A caller building { id, retryCount } from an id stored in // another spelling would pass the lookup and silently miss the fence, so the pair is built from the // canonical form. Done here rather than in SQL, so the fence stays a hashed = ANY over a parameter. // Ids from fetch() and work() are already canonical and take the regex test alone. Anything that is // not 32 hex digits is left as it is: it is not a uuid, and the lookup rejects it. function canonicalUuid (id: string): string { if (CANONICAL_UUID.test(id)) return id const hex = id.replace(/[{}-]/g, '').toLowerCase() if (!/^[0-9a-f]{32}$/.test(hex)) return id return `${hex.slice(0, 8)}-${hex.slice(8, 12)}-${hex.slice(12, 16)}-${hex.slice(16, 20)}-${hex.slice(20)}` } // The fence for the per-job-output statements, which carry each job's retry_count in the recordset. function recordsetAttemptFence (alias: string, input: string): string { return `AND ${alias}state = '${JOB_STATES.active}' AND ${alias}retry_count = ${input}.retry_count` } // Shared SET/WHERE body for marking jobs completed (no RETURNING). Used by the // single-statement completeJobs() and the distributed completeJobsDistributed(). function completeJobsUpdate (schema: string, table: string, includeQueued?: boolean, fenced?: boolean): string { return `UPDATE ${schema}.${table} SET completed_on = ${schema}.job_now(), state = '${JOB_STATES.completed}', output = $3::jsonb, blocked = ${includeQueued ? 'false' : 'blocked'}, pending_dependencies = ${includeQueued ? '0' : 'pending_dependencies'} WHERE name = $1 AND id = ANY($2::uuid[]) AND ${includeQueued ? `state < '${JOB_STATES.completed}'` : `state = '${JOB_STATES.active}'` } ${fenced ? attemptFence(4) : ''}` } // Shared dependency-unblocking fragments. Both consume a `decremented` CTE // (child_name, child_id, n) that the caller defines, and are reused by the standard // completeJobs() and the distributed decrementDependents(). function lockedChildrenCte (schema: string): string { return `locked_children AS ( SELECT j.name, j.id, d.n FROM ${schema}.job j JOIN decremented d ON d.child_name = j.name AND d.child_id = j.id WHERE j.blocked ORDER BY j.name, j.id FOR UPDATE OF j )` } // A child released by its last parent has its start_after moved up to the release, so its wait (in // the monitor's histograms and ready_oldest_seconds) counts from when it could first run rather than // from when the flow was sent. A start_after still in the future is kept. function unblockChildrenUpdate (schema: string): string { return `UPDATE ${schema}.job j SET pending_dependencies = GREATEST(j.pending_dependencies - lc.n, 0), blocked = GREATEST(j.pending_dependencies - lc.n, 0) > 0, start_after = CASE WHEN GREATEST(j.pending_dependencies - lc.n, 0) = 0 THEN GREATEST(j.start_after, ${schema}.job_now()) ELSE j.start_after END FROM locked_children lc WHERE j.name = lc.name AND j.id = lc.id` } // Dependency unblocking is intentionally NOT done here. Completion is the hot path; chasing // dependents inline (joining job_dependency and the partitioned job table) made completion // scale with partition count (see issue #824). The background resolver (Navigator) handles // unblocking out of band, driven by the job_i9 partial index. export function completeJobs (schema: string, table: string, includeQueued?: boolean, fenced?: boolean) { return ` WITH results AS ( ${completeJobsUpdate(schema, table, includeQueued, fenced)} RETURNING id ) ${settledCountAndIds()} ` } // Per-job-output completion: each job's output is supplied via a JSON recordset ($2) and applied by // id, so a batch can be completed with distinct outputs in a single statement. Mirrors completeJobs // (only active jobs; same dependency-unblocking), but sources output from the input join. export function completeJobsWithOutputs (schema: string, table: string) { return ` WITH input AS ( SELECT * FROM json_to_recordset($2::text::json) AS x (id uuid, retry_count int, output jsonb) ), results AS ( UPDATE ${schema}.${table} j SET completed_on = ${schema}.job_now(), state = '${JOB_STATES.completed}', output = i.output FROM input i WHERE j.name = $1 AND j.id = i.id ${recordsetAttemptFence('j.', 'i')} RETURNING j.id ) ${settledCountAndIds()} ` } // Distributed equivalent of completeJobsWithOutputs: a single mutation that returns the completed // ids. Dependency unblocking is handled out of band by the background resolver (Navigator), so // completion does no dependency work on any backend. export function completeJobsWithOutputsDistributed (schema: string, table: string) { return ` WITH input AS ( SELECT * FROM json_to_recordset($2::text::json) AS x (id uuid, retry_count int, output jsonb) ) UPDATE ${schema}.${table} j SET completed_on = ${schema}.job_now(), state = '${JOB_STATES.completed}', output = i.output FROM input i WHERE j.name = $1 AND j.id = i.id ${recordsetAttemptFence('j.', 'i')} RETURNING j.id ` } export function cancelJobs (schema: string, table: string, fenced?: boolean) { return ` WITH results as ( UPDATE ${schema}.${table} SET completed_on = ${schema}.job_now(), state = '${JOB_STATES.cancelled}' WHERE name = $1 AND id = ANY($2::uuid[]) AND state < '${JOB_STATES.completed}' ${fenced ? attemptFence(3) : ''} RETURNING id ) ${settledCountAndIds()} ` } // A resumed job's start_after moves up to now, as a released flow child's does, so its wait (in the // monitor's histograms and ready_oldest_seconds) counts from when it could run again rather than // from when it was first sent. A start_after still in the future is kept. It also loses upsert_by_key, // here and in restoreJobs, so it never collides in job_i13 with a newer job upserted by the same key, // and queues beside it. export function resumeJobs (schema: string, table: string) { return ` WITH results as ( UPDATE ${schema}.${table} SET completed_on = NULL, state = '${JOB_STATES.created}', upsert_by_key = NULL, start_after = GREATEST(start_after, ${schema}.job_now()) WHERE name = $1 AND id = ANY($2::uuid[]) AND state = '${JOB_STATES.cancelled}' RETURNING 1 ) SELECT COUNT(*) from results ` } export function restoreJobs (schema: string, table: string) { return ` UPDATE ${schema}.${table} SET state = '${JOB_STATES.created}', started_on = NULL, heartbeat_on = NULL, upsert_by_key = NULL WHERE name = $1 AND id = ANY($2::uuid[]) ` } // A `startAfter` string is either an absolute date time or a delay expressed as a Postgres // interval. A trailing 'Z' has always marked a date time; a leading ISO 8601 calendar date // (YYYY-MM-DD) marks one as well, which is what lets the other 8601 zone designators through // ('+00:00', '+05:30', '-08:00') along with zone-less and date-only strings, all of which // used to reach the interval cast and fail as 'invalid input syntax for type interval'. // // Recognition is only ever widened, so anything that resolves as an interval today still // does: bare seconds ('0', '300'), phrases ('5 minutes'), ISO 8601 durations ('PT1H') and // the year-month form ('2027-01', which Postgres reads as 2027 years 1 mon) carry neither // mark. A date time with an explicit offset resolves to that exact instant. // // A zone-less date time would otherwise be cast in the database session's TimeZone, so the // public entry points pin it to UTC before it gets here (Attorney.pinZonelessDateTime). This // cast is still session-TZ dependent for anything that reaches it unpinned, which is why the // pin lives at the boundary rather than in this expression: a caller may legitimately pass a // form Postgres resolves itself ('2027-01-01 08:00:00 America/New_York'), and rewriting those // in SQL would mean re-implementing timestamp parsing in a regex. function isDateTimeString (expression: string) { return `(right(${expression}, 1) = 'Z' OR ${expression} ~ '^[0-9]{4}-[0-9]{2}-[0-9]{2}')` } interface InsertJobsOptions { table: string name: string returnId?: boolean notify?: boolean // Whether a job may name the throttle slot it is filed in. Only the cron pass asks for it, so the // statement a public insert() builds does not declare the column and a caller naming it sets // nothing. slots?: boolean // Marks the jobs as inserted by upsert() by singletonKey, for job_i13. Set by the statement, never // read from the recordset, so insert() cannot mark a job. upsertByKey?: boolean } // The queue is LEFT JOINed and every NOT NULL column it supplies falls back to the schema default, so // a job for a queue deleted after the caller cached it still reaches q_fkey and fails there (23503), // rather than producing no row, which would read the same as a singleton or throttle refusal. The // fallbacks never apply to a queue that exists, since its own columns are NOT NULL. export function insertJobs (schema: string, { table, name, returnId = true, notify = false, slots = false, upsertByKey = false }: InsertJobsOptions) { // When notify is enabled we always RETURN start_after so the wrapper below can gate // the NOTIFY on immediate availability, regardless of whether the caller wants ids. const returning = notify ? 'RETURNING id, start_after' : returnId ? 'RETURNING id' : '' // A caller that knows the slot names it outright: the cron pass files an occurrence in the slot // the occurrence falls in, and an offset off now() cannot pin that, since now() here is // insert time. Only in the statement the pass asks for, because insert() stringifies caller // objects straight into the recordset below, so a column declared for everyone would be a live, // undeclared and unvalidated option on the public path, where a bad value surfaces as a raw // postgres error. Prefixed as well, and not called singletonOn, which is a column fetching a job // hands back, so a job read from one queue and inserted into another cannot fill it in by accident. const slotClause = slots ? 'WHEN "__singletonSlot" IS NOT NULL THEN CAST("__singletonSlot" as timestamp)' : '' const insert = ` INSERT INTO ${schema}.${table} ( id, name, data, priority, start_after, created_on, singleton_key, singleton_on, group_id, group_tier, expire_seconds, deletion_seconds, keep_until, retry_limit, retry_delay, retry_backoff, retry_delay_max, policy, dead_letter, heartbeat_seconds, blocked, blocking, pending_dependencies, trace_context${upsertByKey ? ', upsert_by_key' : ''} ) SELECT COALESCE(id, gen_random_uuid()) as id, '${name}' as name, data, COALESCE(priority, 0) as priority, j.start_after, ${schema}.job_now() as created_on, "singletonKey", CASE ${slotClause} WHEN "singletonSeconds" IS NOT NULL THEN 'epoch'::timestamp + '1s'::interval * ("singletonSeconds"::float8 * floor(( date_part('epoch', ${schema}.job_now()) + COALESCE("singletonOffset",0)::float8) / "singletonSeconds"::float8 )) ELSE NULL END as singleton_on, "groupId" as group_id, "groupTier" as group_tier, COALESCE("expireInSeconds", q.expire_seconds, ${QUEUE_DEFAULTS.expire_seconds}) as expire_seconds, COALESCE("deleteAfterSeconds", q.deletion_seconds, ${QUEUE_DEFAULTS.deletion_seconds}) as deletion_seconds, j.start_after + (COALESCE("retentionSeconds", q.retention_seconds, ${QUEUE_DEFAULTS.retention_seconds}) * interval '1s') as keep_until, COALESCE("retryLimit", q.retry_limit, ${QUEUE_DEFAULTS.retry_limit}) as retry_limit, COALESCE("retryDelay", q.retry_delay, ${QUEUE_DEFAULTS.retry_delay}) as retry_delay, COALESCE("retryBackoff", q.retry_backoff, ${QUEUE_DEFAULTS.retry_backoff}) as retry_backoff, COALESCE("retryDelayMax", q.retry_delay_max) as retry_delay_max, q.policy, COALESCE("deadLetter", q.dead_letter) as dead_letter, COALESCE("heartbeatSeconds", q.heartbeat_seconds) as heartbeat_seconds, COALESCE(blocked, false) as blocked, COALESCE(blocking, false) as blocking, COALESCE("pendingDependencies", 0) as pending_dependencies, "__traceContext" as trace_context${upsertByKey ? ', true as upsert_by_key' : ''} FROM ( SELECT *, CASE WHEN ${isDateTimeString('"startAfter"')} THEN CAST("startAfter" as timestamp with time zone) ELSE ${schema}.job_now() + CAST(COALESCE("startAfter",'0') as interval) END as start_after FROM json_to_recordset($1::text::json) as x ( id uuid, priority integer, data jsonb, "startAfter" text, "retryLimit" integer, "retryDelay" integer, "retryDelayMax" integer, "retryBackoff" boolean, "singletonKey" text, "singletonSeconds" integer, "singletonOffset" integer, ${slots ? '"__singletonSlot" text,' : ''} "groupId" text, "groupTier" text, "expireInSeconds" integer, "deleteAfterSeconds" integer, "retentionSeconds" integer, "deadLetter" text, "heartbeatSeconds" integer, blocked boolean, blocking boolean, "pendingDependencies" integer, "__traceContext" jsonb ) ) j LEFT JOIN ${schema}.queue q ON q.name = '${name}' ON CONFLICT DO NOTHING ${returning} ` if (!notify) { return insert } // Fire a single transactional NOTIFY (committed atomically with the insert) only when // at least one inserted row is immediately runnable. Future-dated/throttled jobs are // left to the polling floor. The `notified` CTE is referenced from the final WHERE so // Postgres actually evaluates it; pg_notify runs at most once thanks to LIMIT 1. The // comparator shapes the output rows to honor returnId without changing notify behavior. const comparator = returnId ? '>= 0' : '< 0' return ` WITH ins AS ( ${insert} ), notified AS ( SELECT pg_notify(${notifyChannelSql(schema)}, '${name}') FROM ins WHERE start_after <= ${schema}.job_now() LIMIT 1 ) SELECT id FROM ins WHERE (SELECT count(*) FROM notified) ${comparator} ` } // Self-contained (parameter-less) insert for one queue's slice of a flow batch. The JSON // payload is embedded directly so the whole flow can be sent as a single multi-statement // round-trip regardless of db adapter. Guarded so a skipped row (ON CONFLICT) raises // 'division by zero', aborting the surrounding transaction. The divisor references the // row count so it isn't constant-folded at plan time. export function insertFlowJobs (schema: string, { table, name }: { table: string, name: string }, jobs: unknown[]): string { const insert = insertJobs(schema, { table, name, returnId: true }) .replace('$1', () => serializeJsonParam(jobs)) return ` WITH ins AS ( ${insert} ) SELECT 1 / (CASE WHEN (SELECT count(*) FROM ins) = ${jobs.length} THEN 1 ELSE 0 END) ` } export function failJobsById (schema: string, table: string, fenced?: boolean) { const where = `name = $1 AND id = ANY($2::uuid[]) AND state < '${JOB_STATES.completed}' ${fenced ? attemptFence(4) : ''}` const output = '$3::jsonb' return failJobs(schema, table, where, output, true) } export function failJobsByTimeout (schema: string, table: string, queues: string[], noAdvisoryLocks?: boolean): string { const where = `state = '${JOB_STATES.active}' AND (started_on + expire_seconds * interval '1s') < ${schema}.job_now() AND name = ANY(${serializeArrayParam(queues)})` const output = '\'{ "value": { "message": "job timed out" } }\'::jsonb' return locked(schema, failJobs(schema, table, where, output), table + 'failJobsByTimeout', noAdvisoryLocks) } export function failJobsByHeartbeat (schema: string, table: string, queues: string[], noAdvisoryLocks?: boolean): string { const where = `state = '${JOB_STATES.active}' AND heartbeat_seconds IS NOT NULL AND (heartbeat_on + heartbeat_seconds * interval '1s') < ${schema}.job_now() AND name = ANY(${serializeArrayParam(queues)})` const output = '\'{ "value": { "message": "job heartbeat timeout" } }\'::jsonb' return locked(schema, failJobs(schema, table, where, output), table + 'failJobsByHeartbeat', noAdvisoryLocks) } export function touchJobs (schema: string, table: string, fenced?: boolean) { return ` WITH results AS ( UPDATE ${schema}.${table} SET heartbeat_on = ${schema}.job_now() WHERE name = $1 AND id = ANY($2::uuid[]) AND state = '${JOB_STATES.active}' ${fenced ? attemptFence(3) : ''} RETURNING id ) ${settledCountAndIds()} ` } // `returnIds` is off for the supervisor's bulk maintenance (failJobsByTimeout/failJobsByHeartbeat), // which only ever reads the count and would otherwise ship a whole sweep's worth of ids back. function failJobs (schema: string, table: string, where: string, output: string, returnIds = false) { return ` WITH ${failJobsBody(schema, table, where, output)} ${returnIds ? settledCountAndIds() : 'SELECT COUNT(*) FROM results'} ` } // The tail of every settle statement: how many rows the mutation touched, and which. The count is // what callers report as `affected`; the ids are what tells a transactional handler's own settle // apart from a claim that went away, since a short count alone cannot say which of its jobs landed. // Aggregated from the same `results` CTE the count comes from, so it costs no extra round trip. function settledCountAndIds () { return "SELECT COUNT(*), COALESCE(array_agg(id), '{}'::uuid[]) AS ids FROM results" } // The CTE chain shared by failJobs() and failJobsByIdWithOutputs(): delete the matched jobs and // re-insert them as retry (when retries remain) or failed (+ dead letter). `where` selects the rows // to fail and `output` is the SQL expression stored on each re-inserted job. Returned without the // leading `WITH` or trailing `SELECT` so callers can prepend extra CTEs (e.g. an output map). // When `forceTerminal` is set, every re-inserted job goes straight to the terminal `failed` state // regardless of remaining retries, so the dlq_jobs CTE routes it to the dead letter queue (if any) // immediately. This backs the perJobResults `deadletter` disposition. // // Both re-inserts carry the source_* provenance columns. A job in a dead letter queue that its // own worker fails is deleted and re-inserted here like any other, and without them it would // forget which queue it came from and become unroutable for redrive. // // The dead letter copy takes the failed job's output as source_output, not as its own output. The // copy is a new job that has not run yet, so its output starts empty like any other. function failJobsBody (schema: string, table: string, where: string, output: string, forceTerminal = false) { const state = forceTerminal ? `'${JOB_STATES.failed}'::${schema}.job_state` : `CASE WHEN retry_count < retry_limit THEN '${JOB_STATES.retry}'::${schema}.job_state ELSE '${JOB_STATES.failed}'::${schema}.job_state END` const completedOn = forceTerminal ? `${schema}.job_now()` : `CASE WHEN retry_count < retry_limit THEN NULL ELSE ${schema}.job_now() END` return `deleted_jobs AS ( DELETE FROM ${schema}.${table} WHERE ${where} RETURNING * ), retried_jobs AS ( INSERT INTO ${schema}.${table} ( id, name, priority, data, state, retry_limit, retry_count, retry_delay, retry_backoff, retry_delay_max, start_after, started_on, singleton_key, singleton_on, group_id, group_tier, expire_seconds, deletion_seconds, created_on, completed_on, keep_until, policy, output, dead_letter, heartbeat_on, heartbeat_seconds, blocked, blocking, pending_dependencies, source_name, source_id, source_created_on, source_retry_count, source_output, source_root_id, trace_context ) SELECT id, name, priority, data, ${state} as state, retry_limit, retry_count, retry_delay, retry_backoff, retry_delay_max, CASE WHEN retry_count = retry_limit THEN start_after WHEN NOT retry_backoff THEN ${schema}.job_now() + retry_delay * interval '1' ELSE ${schema}.job_now() + LEAST( retry_delay_max, GREATEST(retry_delay, 1) * ( 2 ^ LEAST(16, retry_count + 1) / 2 + 2 ^ LEAST(16, retry_count + 1) / 2 * random() ) ) * interval '1s' END as start_after, started_on, singleton_key, singleton_on, group_id, group_tier, expire_seconds, deletion_seconds, created_on, ${completedOn} as completed_on, keep_until, policy, ${output}, dead_letter, NULL as heartbeat_on, heartbeat_seconds, blocked, blocking, pending_dependencies, source_name, source_id, source_created_on, source_retry_count, source_output, source_root_id, trace_context FROM deleted_jobs ON CONFLICT DO NOTHING RETURNING * ), failed_jobs as ( INSERT INTO ${schema}.${table} ( id, name, priority, data, state, retry_limit, retry_count, retry_delay, retry_backoff, retry_delay_max, start_after, started_on, singleton_key, singleton_on, group_id, group_tier, expire_seconds, deletion_seconds, created_on, completed_on, keep_until, policy, output, dead_letter, heartbeat_on, heartbeat_seconds, blocked, blocking, pending_dependencies, source_name, source_id, source_created_on, source_retry_count, source_output, source_root_id, trace_context ) SELECT id, name, priority, data, '${JOB_STATES.failed}'::${schema}.job_state as state, retry_limit, retry_count, retry_delay, retry_backoff, retry_delay_max, start_after, started_on, singleton_key, singleton_on, group_id, group_tier, expire_seconds, deletion_seconds, created_on, ${schema}.job_now() as completed_on, keep_until, policy, ${output}, dead_letter, NULL as heartbeat_on, heartbeat_seconds, blocked, blocking, pending_dependencies, source_name, source_id, source_created_on, source_retry_count, source_output, source_root_id, trace_context FROM deleted_jobs WHERE id NOT IN (SELECT id from retried_jobs) RETURNING * ), results as ( SELECT * FROM retried_jobs UNION ALL SELECT * FROM failed_jobs ), dlq_jobs as ( INSERT INTO ${schema}.job (name, priority, data, retry_limit, retry_backoff, retry_delay, start_after, created_on, keep_until, deletion_seconds, expire_seconds, singleton_key, group_id, group_tier, heartbeat_seconds, source_name, source_id, source_created_on, source_retry_count, source_output, source_root_id, trace_context) SELECT r.dead_letter, r.priority, r.data, q.retry_limit, q.retry_backoff, q.retry_delay, ${schema}.job_now(), ${schema}.job_now(), ${schema}.job_now() + q.retention_seconds * interval '1s', q.deletion_seconds, q.expire_seconds, r.singleton_key, r.group_id, r.group_tier, q.heartbeat_seconds, r.name, r.id, r.created_on, r.retry_count, r.output, COALESCE(r.source_root_id, r.id), r.trace_context FROM results r JOIN ${schema}.queue q ON q.name = r.dead_letter WHERE state = '${JOB_STATES.failed}' )` } export function failJobsByIdWithOutputs (schema: string, table: string) { // Output is supplied per job via a JSON recordset ($2). `where` and the output expression both // reference the output_map CTE so each re-inserted job keeps its own output. Constant number of // statements regardless of batch size. const where = `name = $1 AND (id, retry_count) IN (SELECT id, retry_count FROM output_map) AND state = '${JOB_STATES.active}'` const output = '(SELECT om.output FROM output_map om WHERE om.id = deleted_jobs.id)' return ` WITH output_map AS ( SELECT * FROM json_to_recordset($2::text::json) AS x (id uuid, retry_count int, output jsonb) ), ${failJobsBody(schema, table, where, output)} ${settledCountAndIds()} ` } // Like failJobsByIdWithOutputs, but fails every job terminally (forceTerminal) so it routes straight // to the dead letter queue, bypassing remaining retries. Backs the perJobResults `deadletter` status. export function deadLetterJobsByIdWithOutputs (schema: string, table: string) { const where = `name = $1 AND (id, retry_count) IN (SELECT id, retry_count FROM output_map) AND state = '${JOB_STATES.active}'` const output = '(SELECT om.output FROM output_map om WHERE om.id = deleted_jobs.id)' return ` WITH output_map AS ( SELECT * FROM json_to_recordset($2::text::json) AS x (id uuid, retry_count int, output jsonb) ), ${failJobsBody(schema, table, where, output, true)} ${settledCountAndIds()} ` } // Distributed mode: separate queries to avoid CockroachDB's multi-mutation CTE limitation // The five timestamp columns a re-insert binds back are carried alongside the row as text. These rows // come straight from a `SELECT *`, so every column arrives through whatever parser the application // installed on the pool, and pg-boss shares that pool: a parser returning anything but a Date (a // Temporal type, a Luxon DateTime) leaves an object that `pg` encodes as JSON and Postgres rejects // with `invalid input syntax for type timestamp with time zone`. The text rendering carries the // offset and full precision, so it round-trips unchanged, which is what releaseBamCommand does with // a claim's started_on for the same reason. const REBOUND_TIMESTAMPS_AS_TEXT = `started_on::text as started_on_text, singleton_on::text as singleton_on_text, created_on::text as created_on_text, keep_until::text as keep_until_text, start_after::text as start_after_text, source_created_on::text as source_created_on_text` export function selectJobsToFailById (schema: string, table: string, fenced?: boolean): SqlQuery { return { text: `SELECT *, ${REBOUND_TIMESTAMPS_AS_TEXT} FROM ${schema}.${table} WHERE name = $1 AND id = ANY($2::uuid[]) AND state < '${JOB_STATES.completed}' ${fenced ? attemptFence(3) : ''}`, values: [] } } export function deleteJobsToFail (schema: string, table: string): SqlQuery { return { text: `DELETE FROM ${schema}.${table} WHERE name = $1 AND id = ANY($2::uuid[])`, values: [] } } // Distributed mode: the predicate-based maintenance expiry equivalents of selectJobsToFailById. // The supervisor's failJobsByTimeout/failJobsByHeartbeat use the multi-mutation failJobs() CTE, // which CockroachDB rejects, so in distributed mode we select the timed-out jobs here and re-insert // them separately (delete via deleteJobsByIds, re-insert via insertRetryJob), all in one transaction. export function selectJobsToFailByTimeout (schema: string, table: string, queues: string[]): SqlQuery { return { text: `SELECT *, ${REBOUND_TIMESTAMPS_AS_TEXT} FROM ${schema}.${table} WHERE state = '${JOB_STATES.active}' AND (started_on + expire_seconds * interval '1s') < ${schema}.job_now() AND name = ANY(${serializeArrayParam(queues)})`, values: [] } } export function selectJobsToFailByHeartbeat (schema: string, table: string, queues: string[]): SqlQuery { return { text: `SELECT *, ${REBOUND_TIMESTAMPS_AS_TEXT} FROM ${schema}.${table} WHERE state = '${JOB_STATES.active}' AND heartbeat_seconds IS NOT NULL AND (heartbeat_on + heartbeat_seconds * interval '1s') < ${schema}.job_now() AND name = ANY(${serializeArrayParam(queues)})`, values: [] } } export function deleteJobsByIds (schema: string, table: string): SqlQuery { return { text: `DELETE FROM ${schema}.${table} WHERE id = ANY($1::uuid[])`, values: [] } } // Distributed mode: complete jobs as a single-table mutation. Dependency unblocking is handled // out of band by the background resolver (Navigator), so completion does no dependency work. export function completeJobsDistributed (schema: string, table: string, includeQueued?: boolean, fenced?: boolean): string { return ` ${completeJobsUpdate(schema, table, includeQueued, fenced)} RETURNING id ` } // Decrement pending_dependencies for children of the given completed parent jobs, unblocking // any that reach zero. Only the final UPDATE mutates job, so this is a single mutation acceptable // to CockroachDB. Used by the distributed flow resolver path. $1 is the parent queue name, $2 the // list of resolved parent ids for that queue. export function decrementDependents (schema: string): string { return ` WITH decremented AS ( SELECT d.child_name, d.child_id, COUNT(*)::int AS n FROM ${schema}.job_dependency d WHERE d.parent_name = $1 AND d.parent_id = ANY($2::uuid[]) GROUP BY d.child_name, d.child_id ), ${lockedChildrenCte(schema)} ${unblockChildrenUpdate(schema)} ` } // Background flow resolver (Navigator) batch size: the max number of completed blocking parents // locked per audit statement. The resolver loops until a batch drains, so this only bounds the // lock footprint and per-statement cost. export const FLOW_BATCH_SIZE = 1000 // Standard (multi-mutation CTE) flow audit. Locks a batch of completed blocking parents in the // given partition table, decrements their children's pending_dependencies (reusing the shared // unblock fragments, which reach across partitions via the parent job table), unblocks children // that reach zero, and clears `blocking` on the resolved parents so they leave the job_i9 index // and are never reprocessed. $1 is the chunk of queue names (for partition pruning). Returns the // number of parents resolved so the caller can loop until a batch drains. export function resolveFlowJobs (schema: string, table: string, names: string[]): SqlQuery { return { text: ` WITH locked_parents AS ( SELECT j.name, j.id FROM ${schema}.${table} j WHERE j.blocking AND j.state = '${JOB_STATES.completed}' AND j.name = ANY($1::text[]) ORDER BY j.name, j.id FOR UPDATE OF j SKIP LOCKED LIMIT ${FLOW_BATCH_SIZE} ), decremented AS ( SELECT d.child_name, d.child_id, COUNT(*)::int AS n FROM ${schema}.job_dependency d JOIN locked_parents p ON d.parent_name = p.name AND d.parent_id = p.id GROUP BY d.child_name, d.child_id ), ${lockedChildrenCte(schema)}, unblocked AS ( ${unblockChildrenUpdate(schema)} RETURNING 1 ), cleared AS ( UPDATE ${schema}.${table} j SET blocking = false FROM locked_parents p WHERE j.name = p.name AND j.id = p.id RETURNING 1 ) SELECT COUNT(*)::int AS resolved FROM cleared `, values: [names] } } // Distributed flow audit (CockroachDB / noMultiMutationCte). Locks a batch of completed blocking // parents without mutating, so the caller can run the single-mutation decrementDependents() and // clearBlocking() separately within one transaction. $1 is the chunk of queue names; SKIP LOCKED // is omitted under noSkipLocked. export function selectBlockingParents (schema: string, table: string, names: string[], noSkipLocked?: boolean): SqlQuery { return { text: ` SELECT name, id FROM ${schema}.${table} WHERE blocking AND state = '${JOB_STATES.completed}' AND name = ANY($1::text[]) ORDER BY name, id FOR UPDATE${noSkipLocked ? '' : ' SKIP LOCKED'} LIMIT ${FLOW_BATCH_SIZE} `, values: [names] } } // Distributed flow audit: clear `blocking` on resolved parents (single mutation). $1 is the parent // queue name, $2 the list of resolved parent ids for that queue. export function clearBlocking (schema: string): string { return ` UPDATE ${schema}.job SET blocking = false WHERE name = $1 AND id = ANY($2::uuid[]) ` } export function insertRetryJob (schema: string, table: string): string { return ` INSERT INTO ${schema}.${table} ( id, name, priority, data, state, retry_limit, retry_count, retry_delay, retry_backoff, retry_delay_max, start_after, started_on, singleton_key, singleton_on, group_id, group_tier, expire_seconds, deletion_seconds, created_on, completed_on, keep_until, policy, output, dead_letter, heartbeat_on, heartbeat_seconds, blocked, blocking, pending_dependencies, source_name, source_id, source_created_on, source_retry_count, source_output, source_root_id, trace_context ) VALUES ( $1, $2, $3, $4::text::jsonb, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22, $23::text::jsonb, $24, $25, $26, $27, $28, $29, $30, $31, $32, $33, $34::text::jsonb, $35, $36::text::jsonb ) ON CONFLICT DO NOTHING RETURNING id ` } export function insertDeadLetterJob (schema: string): string { return ` INSERT INTO ${schema}.job (name, data, priority, retry_limit, retry_backoff, retry_delay, start_after, created_on, keep_until, deletion_seconds, expire_seconds, singleton_key, group_id, group_tier, heartbeat_seconds, source_name, source_id, source_created_on, source_retry_count, source_output, source_root_id, trace_context) SELECT $1, $2::text::jsonb, $9, q.retry_limit, q.retry_backoff, q.retry_delay, ${schema}.job_now(), ${schema}.job_now(), ${schema}.job_now() + q.retention_seconds * interval '1s', q.deletion_seconds, q.expire_seconds, $8, $10, $11, q.heartbeat_seconds, $4, $5, $6, $7, $3::text::jsonb, COALESCE($12::uuid, $5::uuid), $13::text::jsonb FROM ${schema}.queue q WHERE q.name = $1 ` } // The candidate predicate shared by redriveJobs and previewRedrive, so a preview can never count // a job the redrive would skip or skip one it would move. Parameters: $1 dead letter queue, // $2 destination override, $3 sourceName, $4 data (jsonb containment), $5 createdBefore, // $6 ids. Each filter is off when its parameter is null. Only jobs not yet active are // candidates: a job the dead letter queue's own workers failed stays where it is. // // A key_strict_fifo job whose key is held by an active, retrying or failed job waits behind it, as // it would for fetch. job_i8 allows one holder per key, so a collision could not fail it in place, // and left queued it would be a candidate on every call and a draining loop would never end. function redriveWhere (schema: string, table: string): string { return `j.name = $1 AND j.state < '${JOB_STATES.active}' AND NOT EXISTS ( SELECT 1 FROM ${schema}.${table} k WHERE j.policy = '${QUEUE_POLICIES.key_strict_fifo}' AND k.name = j.name AND k.singleton_key = j.singleton_key AND k.policy = '${QUEUE_POLICIES.key_strict_fifo}' AND k.state IN ('${JOB_STATES.active}', '${JOB_STATES.retry}', '${JOB_STATES.failed}') ) AND ($3::text IS NULL OR j.source_name = $3) AND ($4::text::jsonb IS NULL OR j.data @> $4::text::jsonb) AND ($5::timestamptz IS NULL OR j.created_on < $5) AND ($6::uuid[] IS NULL OR j.id = ANY($6::uuid[]))` } // The columns a redriven job is created with, shared by both redrive paths. `m` is the dead letter // row and `q` the destination queue. Re-created jobs get a new id, `created` state, retry_count 0, // cleared output, NULL source_* except source_root_id, which carries the chain's first job on (see // createTableJob), and every queue-config column (retry/retention/policy/expiry/ // heartbeat/dead_letter) from the destination queue as it is configured now, per-job overrides // from the original send() are not preserved, since the DLQ copy never stored them. `dead_letter` // is the same value send() falls back to, so a second terminal failure re-enters the DLQ. // Job-identity columns (priority, singleton_key, group_id, group_tier) are carried over instead. const REDRIVE_INSERT_COLUMNS = `(id, name, data, priority, retry_limit, retry_backoff, retry_delay, retry_delay_max, expire_seconds, start_after, created_on, keep_until, deletion_seconds, policy, singleton_key, group_id, group_tier, heartbeat_seconds, dead_letter, source_root_id, trace_context)` function redriveInsertValues (schema: string, newId: string, destination: string): string { return `${newId}, COALESCE(${destination}, m.source_name), m.data, m.priority, q.retry_limit, q.retry_backoff, q.retry_delay, q.retry_delay_max, q.expire_seconds, ${schema}.job_now(), ${schema}.job_now(), ${schema}.job_now() + q.retention_seconds * interval '1s', q.deletion_seconds, q.policy, m.singleton_key, m.group_id, m.group_tier, q.heartbeat_seconds, q.dead_letter, COALESCE(m.source_root_id, m.source_id), m.trace_context` } // What a job that could not be re-created becomes: failed, in place, in the dead letter queue, with // the reason as its output. Never deleted: before 12.35 it was, and the job was simply gone. As a // failed job it is removed by the dead letter queue's own deleteAfterSeconds like any other, so // nothing accumulates, and until then it can be seen and retried. Failed rather than left waiting because a // waiting job is a candidate again on the next call, and redrive takes candidates oldest first, so // a handful of them would fill every later batch and stall the drain. Failed is also the state // redrive already leaves alone, and retrying the job makes it a candidate again once whatever it // collided with has finished. // // The output it had is kept as source_output when that is empty, which only happens on a job // dead-lettered before source_output existed: those copied the original's output into their own. // `j` is the dead letter row and `q` the destination queue. // // Failing two key_strict_fifo jobs with the same key would violate job_i8, so of those only the // oldest is failed. The rest stay queued behind it, and redriveWhere skips them until it is resolved. // `conflicts` is a predicate on `m.id` selecting the dead letter rows that were not re-created. function redriveConflictFailable (schema: string, table: string, conflicts: string): string { return `(j.policy IS DISTINCT FROM '${QUEUE_POLICIES.key_strict_fifo}' OR j.singleton_key IS NULL OR j.id IN ( SELECT DISTINCT ON (m.singleton_key) m.id FROM ${schema}.${table} m WHERE ${conflicts} AND m.policy = '${QUEUE_POLICIES.key_strict_fifo}' ORDER BY m.singleton_key, m.created_on, m.id ))` } function redriveConflictSet (schema: string): string { return `state = '${JOB_STATES.failed}', completed_on = ${schema}.job_now(), source_output = COALESCE(j.source_output, j.output), output = jsonb_build_object( 'message', 'Not redriven: queue ' || q.name || CASE WHEN j.singleton_key IS NULL THEN ' already has a job this one conflicts with' ELSE ' already has a job with singletonKey ' || j.singleton_key END || ' under its ' || COALESCE(q.policy, 'standard') || ' policy', 'reason', 'redrive_conflict', 'destination', q.name, 'policy', q.policy, 'singletonKey', j.singleton_key )` } // Dead-letter redrive. Moves un-started jobs out of a dead-letter queue and // re-creates them as fresh jobs on their original source queue (or $2 destination override), // oldest-first, capped at $7 and narrowed by the filters in redriveWhere. The JOIN in // `candidates` only matches jobs whose destination queue exists, so legacy/orphaned jobs // (NULL source_name, no override) are never deleted. They stay in the DLQ rather than being lost. // // The insert's ON CONFLICT DO NOTHING is load-bearing: a destination queue's short/stately/exclusive // policy can collide on (name, singleton_key), with a job already there or with another job in the // same batch, and skipping just that row is preferable to aborting the whole batch. The dead letter // row of a job that was not re-created is failed in place with the reason (redriveConflictSet). // // Every candidate is given its new id up front, so `ins` RETURNING says exactly which ones were // re-created: a data-modifying CTE cannot see the others' writes, so there is no other way to tell // within one statement. `candidates` is referenced more than once and is FOR UPDATE, so it is // materialized, and each row keeps the one id it was given. export function redriveJobs (schema: string, table: string): string { return ` WITH candidates AS ( SELECT j.id, gen_random_uuid() AS new_id FROM ${schema}.${table} j JOIN ${schema}.queue q ON q.name = COALESCE($2, j.source_name) WHERE ${redriveWhere(schema, table)} ORDER BY j.created_on LIMIT $7 FOR UPDATE OF j SKIP LOCKED ), ins AS ( INSERT INTO ${schema}.job ${REDRIVE_INSERT_COLUMNS} SELECT ${redriveInsertValues(schema, 'c.new_id', '$2')} FROM candidates c JOIN ${schema}.${table} m ON m.id = c.id JOIN ${schema}.queue q ON q.name = COALESCE($2, m.source_name) ORDER BY m.created_on ON CONFLICT DO NOTHING RETURNING id ), settled AS ( SELECT c.id, (i.id IS NOT NULL) AS inserted FROM candidates c LEFT JOIN ins i ON i.id = c.new_id ), removed AS ( DELETE FROM ${schema}.${table} WHERE id IN (SELECT id FROM settled WHERE inserted) ), conflicted AS ( UPDATE ${schema}.${table} j SET ${redriveConflictSet(schema)} FROM settled s, ${schema}.queue q WHERE j.id = s.id AND NOT s.inserted AND q.name = COALESCE($2, j.source_name) AND ${redriveConflictFailable(schema, table, 'm.id IN (SELECT id FROM settled WHERE NOT inserted)')} ) SELECT count(*)::int AS moved FROM ins ` } // Distributed redrive (noMultiMutationCte). CockroachDB refuses redriveJobs' DELETE and INSERT on // one table in one statement, so the manager runs the steps below in a transaction instead: lock // the candidates and give each its new id, re-create them, delete the ones that were re-created, // and fail the rest in place. Same predicate, same order, same limit, so the two // paths move the same jobs. Both inserts run oldest-first, so on Postgres the older of two // colliding jobs is the one kept, on either path. CockroachDB does not honor that ORDER BY when it // resolves ON CONFLICT, so which one it keeps is arbitrary; only the counts are the same there. export function selectRedriveCandidates (schema: string, table: string): string { return ` SELECT j.id, gen_random_uuid() AS new_id FROM ${schema}.${table} j WHERE ${redriveWhere(schema, table)} AND EXISTS (SELECT 1 FROM ${schema}.queue q WHERE q.name = COALESCE($2, j.source_name)) ORDER BY j.created_on LIMIT $7 FOR UPDATE ` } // $1 candidate ids, $2 destination override, $3 the new id for each candidate, in the same order. // Returns the new ids of the jobs that were re-created. export function insertRedrivenJobs (schema: string, table: string): string { return ` INSERT INTO ${schema}.job ${REDRIVE_INSERT_COLUMNS} SELECT ${redriveInsertValues(schema, 'p.new_id', '$2')} FROM unnest($1::uuid[], $3::uuid[]) AS p (id, new_id) JOIN ${schema}.${table} m ON m.id = p.id JOIN ${schema}.queue q ON q.name = COALESCE($2, m.source_name) ORDER BY m.created_on ON CONFLICT DO NOTHING RETURNING id ` } // $1 the dead letter ids that were not re-created, $2 destination override. export function failRedriveConflicts (schema: string, table: string): string { return ` UPDATE ${schema}.${table} j SET ${redriveConflictSet(schema)} FROM ${schema}.queue q WHERE j.id = ANY($1::uuid[]) AND q.name = COALESCE($2, j.source_name) AND ${redriveConflictFailable(schema, table, 'm.id = ANY($1::uuid[])')} ` } // What a redrive with the same filter would do, without doing it: matching jobs grouped by the // queue each would land in. The LEFT JOIN keeps jobs the redrive would leave behind (no source // and no override, or a source queue since deleted) so they can be reported rather than vanish. export function previewRedrive (schema: string, table: string): string { return ` SELECT COALESCE($2, j.source_name) AS destination, (q.name IS NOT NULL) AS routable, count(*)::int AS count FROM ${schema}.${table} j LEFT JOIN ${schema}.queue q ON q.name = COALESCE($2, j.source_name) WHERE ${redriveWhere(schema, table)} GROUP BY 1, 2 ` } export function deletion (schema: string, table: string, queues: string[], noAdvisoryLocks?: boolean): string { const sql = ` DELETE FROM ${schema}.${table} WHERE name = ANY(${serializeArrayParam(queues)}) AND ( (deletion_seconds > 0 AND completed_on + deletion_seconds * interval '1s' < ${schema}.job_now()) OR (state < '${JOB_STATES.active}' AND keep_until < ${schema}.job_now()) ) ` return locked(schema, sql, table + 'deletion', noAdvisoryLocks) } // start_after moves up to now, as in resumeJobs. export function retryJobs (schema: string, table: string) { return ` WITH results as ( UPDATE ${schema}.job SET state = '${JOB_STATES.retry}', retry_limit = retry_limit + 1, completed_on = NULL, start_after = GREATEST(start_after, ${schema}.job_now()) WHERE name = $1 AND id = ANY($2::uuid[]) AND state = '${JOB_STATES.failed}' RETURNING 1 ) SELECT COUNT(*) from results ` } // Partial in-place edit of not-yet-active jobs, preserving id/state/singleton identity. // The payload ($1) is a jsonb object of ONLY the fields the caller supplied; each column is // left untouched unless its key is present (`(o.data -> 'key') IS NOT NULL`, which is true for a // key carrying JSON null and false only for an absent key, the same answer as jsonb_exists(); that // function does not exist on CockroachDB), so an update that carries just // `data` never clobbers an existing start_after/priority/etc. Targeting is by id or // singleton_key; when by key, `match` picks which of several pre-active matches to edit // (newest/oldest = one row via ORDER BY + LIMIT; all = every match). When `notify` is set the // edit emits a single pg_notify iff a touched row ends up runnable (start_after <= ${schema}.job_now()), // closing the wake-up gap for jobs pulled forward. Callers needing insert-on-miss compose this // with insertJobs (see Manager.upsert). // // Two guards inside are easy to read as incidental. When only startAfter moves, keep_until slides by // the same original retention window (keep_until - start_after), so pulling a job forward or back // never leaves keep_until in the past, where the deletion sweep would treat it as expired and remove // a still-pending job. And the UPDATE re-checks `state < active` on the locked row, not just in the // unlocked target CTE: under READ COMMITTED a concurrent fetchNextJob can activate a candidate // between target selection and the UPDATE, and EvalPlanQual re-evaluates this predicate on the // freshly-locked row, so the repeat prevents mutating a job a worker has already started. export function updateJob (schema: string, table: string, name: string, by: 'id' | 'singletonKey', match: JobMatchStrategy, notify = false) { const targetPredicate = by === 'id' ? "job.id = (o.data->>'id')::uuid" : "job.singleton_key = o.data->>'singletonKey'" const ordering = (by === 'singletonKey' && match !== 'all') ? `ORDER BY job.created_on ${match === 'oldest' ? 'ASC' : 'DESC'} LIMIT 1` : '' // Resolve the incoming startAfter the same way insertJobs does (absolute date time vs. // relative interval), falling back to the row's current start_after when not supplied. const resolvedStartAfter = ` CASE WHEN (o.data -> 'startAfter') IS NOT NULL THEN CASE WHEN ${isDateTimeString("o.data->>'startAfter'")} THEN (o.data->>'startAfter')::timestamptz ELSE ${schema}.job_now() + CAST(o.data->>'startAfter' AS interval) END ELSE job.start_after END` const tail = notify ? `, notified AS ( SELECT pg_notify(${notifyChannelSql(schema)}, '${name}') FROM upd WHERE start_after <= ${schema}.job_now() LIMIT 1 ) SELECT id FROM upd WHERE (SELECT count(*) FROM notified) >= 0` : ` SELECT id FROM upd` return ` WITH o AS (SELECT $1::text::jsonb AS data), target AS ( SELECT job.id FROM ${schema}.${table} job, o WHERE job.name = '${name}' AND job.state < '${JOB_STATES.active}' AND ${targetPredicate} ${ordering} ), upd AS ( UPDATE ${schema}.${table} job SET data = CASE WHEN (o.data -> 'data') IS NOT NULL THEN o.data->'data' ELSE job.data END, priority = COALESCE((o.data->>'priority')::int, job.priority), start_after = ${resolvedStartAfter}, keep_until = CASE WHEN (o.data -> 'retentionSeconds') IS NOT NULL THEN (${resolvedStartAfter}) + ((o.data->>'retentionSeconds')::int * interval '1s') WHEN (o.data -> 'startAfter') IS NOT NULL THEN (${resolvedStartAfter}) + (job.keep_until - job.start_after) ELSE job.keep_until END, expire_seconds = COALESCE((o.data->>'expireInSeconds')::int, job.expire_seconds), deletion_seconds = COALESCE((o.data->>'deleteAfterSeconds')::int, job.deletion_seconds), retry_limit = COALESCE((o.data->>'retryLimit')::int, job.retry_limit), retry_delay = COALESCE((o.data->>'retryDelay')::int, job.retry_delay), retry_backoff = COALESCE((o.data->>'retryBackoff')::bool, job.retry_backoff), retry_delay_max = CASE WHEN (o.data -> 'retryDelayMax') IS NOT NULL THEN (o.data->>'retryDelayMax')::int ELSE job.retry_delay_max END, dead_letter = CASE WHEN (o.data -> 'deadLetter') IS NOT NULL THEN o.data->>'deadLetter' ELSE job.dead_letter END, heartbeat_seconds = CASE WHEN (o.data -> 'heartbeatSeconds') IS NOT NULL THEN (o.data->>'heartbeatSeconds')::int ELSE job.heartbeat_seconds END, group_id = CASE WHEN (o.data -> 'groupId') IS NOT NULL THEN o.data->>'groupId' ELSE job.group_id END, group_tier = CASE WHEN (o.data -> 'groupTier') IS NOT NULL THEN o.data->>'groupTier' ELSE job.group_tier END FROM o WHERE job.id IN (SELECT id FROM target) AND job.state < '${JOB_STATES.active}' RETURNING job.id, job.start_after )${tail} ` } // The throughput window. It starts when this queue was last *counted*, which is not monitor_on: // a pass that doesn't count (persistQueueStats off on that instance, or a forced getQueueStats // refresh) stamps monitor_on too, and windowing on it made every such pass swallow the jobs that // finished before it. Only a statement that counts moves delta_on, so an uncounted pass in between // leaves the window open and the next counted pass picks its jobs up. // // It ends DELTA_LAG behind the pass, not at now(). created_on and completed_on come from job_now(), // the start of the transaction that wrote them, so a job sent inside an application's transaction, // or completed by a transactional worker, carries a stamp from before a pass it commits after. A // window ending at now() stepped past that stamp while the row was still invisible. With the lag, // a transaction that commits within the lag of starting lands ahead of the end, whatever the // backend and whatever role it runs as, and is counted in the window its stamp belongs to. The // counters are DELTA_LAG behind the gauges in the same row, which is what delta_on is for. // // A transaction that runs longer than the lag lands behind a window that has already been counted // and recorded. The pass never waits for it (an open report or backup holding every queue's // counters back is the thing this replaced); the snapshots are trued up after the fact instead. // So the latest hour of history is provisional: a snapshot's counters can rise after it is written, // never fall. // // A window whose start is more than DELTA_RESET_MAX behind its end starts afresh instead: counting // was off for a while, and landing hours of work on one snapshot would chart as a spike at the // moment it was switched back on. Null in, null out, so those comparisons count nothing, the same // as a queue never counted. An instance whose passes are further apart than half of that uses two // of its intervals instead (deltaResetMax), or every window would start afresh and nothing would // ever be counted. Instances that disagree on it can only disagree about whether a gap resets, so // it is safe to derive per instance. // // The lag and the true-up horizon are fixed rather than configured. delta_on is shared by every instance, and a // window that one instance measured with a different lag than the last would count an interval // twice or skip it. The lag only decides how often the true-up runs, not what is counted: measured // under a mixed load whose longest routine transactions held 5s, a lag of 0 or 1s set off a true-up // on 90-97% of passes, and 5s or more on none. Ten seconds leaves headroom over that and puts the // counters within a pass of the gauges. See research/delta-true-up.md on planning. const DELTA_LAG = "interval '10 seconds'" const DELTA_RESET_MAX_SECONDS = 2 * 60 * 60 const DELTA_RESET_MAX = `interval '${DELTA_RESET_MAX_SECONDS} seconds'` const DELTA_TRUE_UP_MAX = "interval '1 hour'" // Test seams for the intervals above. Nothing in production passes them: a test cannot wait a // minute for the lag or an hour for the true-up horizon. export interface DeltaWindowOptions { lag?: string, resetMax?: string, trueUpMax?: string } // The reset threshold for an instance whose counting passes run every intervalSeconds. export function deltaResetMax (intervalSeconds: number): string { return intervalSeconds * 2 > DELTA_RESET_MAX_SECONDS ? `interval '${intervalSeconds * 2} seconds'` : DELTA_RESET_MAX } function deltaWindowEnd (schema: string, lag = DELTA_LAG): string { return `(${schema}.job_now() - ${lag})` } function deltaWindowStart (alias: string, end: string, resetMax = DELTA_RESET_MAX): string { return `(CASE WHEN ${alias}.delta_on > ${end} - ${resetMax} THEN ${alias}.delta_on END)` } // The SET clause for a statement that counts. The right-hand side reads the row as it was before // the update, so delta_seconds is measured from the same window the counts used, and an idle queue // (no job rows, so no stats row) still records zero over a real number of seconds. The window never // moves backwards, even if the end lands before the last one. function throughputAssignments (end: string, resetMax?: string): string { const start = deltaWindowStart('queue', end, resetMax) // CASE rather than GREATEST(0, ...): GREATEST skips NULLs, and a pass with no window has to // record null seconds, not zero. const seconds = `round(extract(epoch from (${end} - ${start})))::int` return ` created_delta = COALESCE(stats."createdDelta", 0), completed_delta = COALESCE(stats."completedDelta", 0), failed_delta = COALESCE(stats."failedDelta", 0), delta_seconds = CASE WHEN ${seconds} < 0 THEN 0 ELSE ${seconds} END, delta_on = GREATEST(queue.delta_on, ${end}), wait_bins = COALESCE(stats."waitBins", ${EMPTY_BINS}), run_bins = COALESCE(stats."runBins", ${EMPTY_BINS}), ready_oldest_seconds = COALESCE(stats."readyOldestSeconds", 0),` } // The windows a true-up may still revise, per queue: every recorded snapshot whose window ends // after the anchor `h`, up to the newest one, `top`, with the counters they hold between them. // // A snapshot's window is not stored; it is rebuilt. Each counted pass starts where the last one // ended and records its end as delta_on, so a snapshot's window runs from the previous snapshot's // delta_on to its own, exactly, with no rounding (delta_seconds is rounded, so start = delta_on - // delta_seconds would not tile). The anchor is the newest snapshot at least DELTA_TRUE_UP_MAX // behind, so only windows younger than the horizon are revised. A snapshot that started a window // afresh (first counted pass, or a reset after a gap; both record null seconds) counted nothing // before its end on purpose, so it anchors too: truing up across it would land hours of work on it. // // The recount only sees rows that still exist, so the anchor also stays inside the queue's // retention. Retention deletes a finished job once completed_on + deletion_seconds has passed, and // a queued one once keep_until (start_after + retention_seconds) has. Rows it removed from a window // that was already counted make the recount fall short of the snapshots, and on a queue that // deletes within the hour that shortfall hides every late commit. The anchor is therefore never // earlier than the first window end past job_now() minus the shorter of the two, where nothing has // been deleted yet; a window straddling that point is not revised. // // When a pass counted but its snapshot was never written (the insert is a separate statement and // can fail), the next snapshot's rebuilt window covers both passes. Its stored counters then fall // short by that pass's jobs, and the true-up restores them there, which is the right place to the // resolution the history has. // The anchor rule and the reach, shared so trueUpWindows (the fix) and trueUpSettled (the monitor's // check) always agree on which windows are revised. `s` is a queue_stats row. function trueUpAnchor (schema: string, trueUpMax: string): string { return `s.delta_on <= ${schema}.job_now() - ${trueUpMax} OR s.delta_seconds IS NULL` } // `q` is the queue row. Per-job overrides of either retention are not seen here. function trueUpRetentionFloor (schema: string, q: string): string { return `s.delta_on >= ${schema}.job_now() - (CASE WHEN ${q}.deletion_seconds > 0 AND ${q}.deletion_seconds < ${q}.retention_seconds THEN ${q}.deletion_seconds ELSE ${q}.retention_seconds END) * interval '1s'` } // The anchor from the horizon rule `b` and the retention floor `f`, whichever is later. A null floor // means no window ends inside the retention, so there is nothing to revise and the anchor is null. function trueUpAnchorAt (b: string, f: string): string { return `CASE WHEN ${f} IS NULL OR ${f} > ${b} THEN ${f} ELSE ${b} END` } function trueUpReach (schema: string, trueUpMax: string): string { return `s.captured_on >= ${schema}.job_now() - 2 * ${trueUpMax} AND s.delta_on IS NOT NULL` } // Every recorded snapshot in reach, each carrying its queue's anchor `h`; the windows are the rows // with delta_on > h. One read of queue_stats: the anchor is a window aggregate over the same rows. // CASE rather than FILTER, which CockroachDB doesn't take on a window function. function trueUpWindows (schema: string, queues: string, trueUpMax = DELTA_TRUE_UP_MAX): string { return ` SELECT x.id, x.captured_on, x.name, x.delta_on, x.created_delta, x.completed_delta, x.failed_delta, ${trueUpAnchorAt('x.b', 'x.f')} AS h FROM ( SELECT s.id, s.captured_on, s.name, s.delta_on, s.created_delta, s.completed_delta, s.failed_delta, COALESCE( max(CASE WHEN ${trueUpAnchor(schema, trueUpMax)} THEN s.delta_on END) OVER (PARTITION BY s.name), min(s.delta_on) OVER (PARTITION BY s.name) ) AS b, min(CASE WHEN ${trueUpRetentionFloor(schema, 'q')} THEN s.delta_on END) OVER (PARTITION BY s.name) AS f FROM ${schema}.queue_stats s JOIN ${schema}.queue q ON q.name = s.name WHERE s.name = ANY(${queues}) AND ${trueUpReach(schema, trueUpMax)} ) x` } // The same anchor and windows, reduced to what the monitor's check needs: per queue, the anchor, // the newest window end `top`, and what the windows between them hold. Two index lookups per queue // on queue_stats (name, captured_on), laterally, rather than trueUpWindows' window aggregate: that // sorts every snapshot in reach by queue name, which on a thousand queues cost more than the check. function trueUpSettled (schema: string, alias: string, trueUpMax = DELTA_TRUE_UP_MAX): string { return ` LEFT JOIN LATERAL ( SELECT ${trueUpAnchorAt('x.b', 'x.f')} AS h FROM ( SELECT COALESCE( max(s.delta_on) FILTER (WHERE ${trueUpAnchor(schema, trueUpMax)}), min(s.delta_on) ) AS b, min(s.delta_on) FILTER (WHERE ${trueUpRetentionFloor(schema, alias)}) AS f FROM ${schema}.queue_stats s WHERE s.name = ${alias}.name AND ${trueUpReach(schema, trueUpMax)} ) x ) a ON true LEFT JOIN LATERAL ( SELECT max(s.delta_on) AS top, sum(s.created_delta + s.completed_delta + s.failed_delta) AS settled FROM ${schema}.queue_stats s WHERE s.name = ${alias}.name AND ${trueUpReach(schema, trueUpMax)} AND s.delta_on > a.h ) t ON true` } // Revise the recent snapshots' counters for jobs that committed after the pass that counted their // window. Run only for the queues the monitor's aggregate flagged: // that check is three filtered counts riding a scan the monitor makes anyway, and this is a second // scan of the queue's rows, paid once per late commit rather than on every pass. // // Each job row is placed in its window by width_bucket over the queue's window ends, a binary // search rather than a join against every window, then counted per window and compared with what // the snapshot holds. The thresholds are [h, end1, end2, ...], so width_bucket returns i for // h <= stamp < end_i, the window of the snapshot ordered i, and past the last end a bucket no // window claims. Counters only rise: retention and deleteJob() remove rows a window already // counted, and a recount that went down would unwrite them. Two instances truing up the same // snapshot write the same GREATEST, so nothing here needs to be serialized for correctness; the // stats try-lock is taken so they don't both scan. One row when it scanned, none when another // instance held the lock, the same contract as the aggregate's pinSeconds. export function trueUpQueueStats (schema: string, table: string, queues: string[], noAdvisoryLocks?: boolean, window: DeltaWindowOptions = {}): string { const names = serializeArrayParam(queues) const trueUpMax = window.trueUpMax ?? DELTA_TRUE_UP_MAX const lock = tryAdvisoryLock(schema, 'queue-stats', noAdvisoryLocks) const sql = ` WITH ${lock.cte}win AS ( SELECT w.*, row_number() OVER (PARTITION BY w.name ORDER BY w.delta_on, w.captured_on) AS i FROM (${trueUpWindows(schema, names, trueUpMax)}) w WHERE w.delta_on > w.h${lock.guard} ), th AS ( SELECT w.name, w.h, array_prepend(w.h, array_agg(w.delta_on ORDER BY w.i)) AS t FROM win w GROUP BY w.name, w.h ), binned AS ( SELECT j.name, CASE WHEN j.created_on >= th.h THEN width_bucket(j.created_on, th.t) END AS cb, CASE WHEN j.state IN ('${JOB_STATES.completed}', '${JOB_STATES.failed}') AND j.completed_on >= th.h THEN width_bucket(j.completed_on, th.t) END AS fb, j.state FROM ${schema}.${table} j JOIN th ON th.name = j.name WHERE j.name = ANY(${names}) ), grouped AS ( SELECT name, cb, fb, state, count(*)::int AS n FROM binned WHERE cb IS NOT NULL OR fb IS NOT NULL GROUP BY 1, 2, 3, 4 ), counted AS ( SELECT name, cb AS i, sum(n)::int AS created, 0 AS completed, 0 AS failed FROM grouped WHERE cb IS NOT NULL GROUP BY 1, 2 UNION ALL SELECT name, fb AS i, 0, COALESCE(sum(n) FILTER (WHERE state = '${JOB_STATES.completed}'), 0)::int, COALESCE(sum(n) FILTER (WHERE state = '${JOB_STATES.failed}'), 0)::int FROM grouped WHERE fb IS NOT NULL GROUP BY 1, 2 ), recount AS ( SELECT w.id, w.captured_on, sum(c.created)::int AS created, sum(c.completed)::int AS completed, sum(c.failed)::int AS failed FROM win w JOIN counted c ON c.name = w.name AND c.i = w.i GROUP BY w.id, w.captured_on ), raised AS ( UPDATE ${schema}.queue_stats s SET created_delta = GREATEST(s.created_delta, r.created), completed_delta = GREATEST(s.completed_delta, r.completed), failed_delta = GREATEST(s.failed_delta, r.failed) FROM recount r WHERE s.id = r.id AND s.captured_on = r.captured_on${lock.guard} AND (r.created > s.created_delta OR r.completed > s.completed_delta OR r.failed > s.failed_delta) RETURNING s.id ) SELECT (SELECT count(*) FROM raised)::int AS raised, ${PIN_SECONDS_SQL} as "pinSeconds" WHERE true${lock.guard} ` return transaction(sql) } // Wait and run times, as the monitor records them: a histogram per counted pass, of the jobs that // finished in its window. Slot 0 holds times under 10 ms, slots 1 to LATENCY_BINS bins that each // grow by √2 (slot k runs from 10 ms · √2^(k-1) to 10 ms · √2^k), and the last slot everything past // about 23 hours. Log-spaced because the times span six orders of magnitude, and a percentile read // from them lies in the same bin as the exact one, a factor of √2 at most. Histograms rather than // percentiles, because histograms add: a reader sums them across passes, buckets or queues and // reads any percentile from the sum. // Stored as 48 slots in slot order, null where no job landed: Postgres keeps a null as one bit // rather than four bytes, and most slots are empty. getQueueStats() hands them out with nulls as 0. // Readers that add slots in SQL coalesce, since a null plus a count is null. export const LATENCY_BINS = 46 export const LATENCY_SLOTS = LATENCY_BINS + 2 export const LATENCY_MIN_SECONDS = 0.01 // A measured histogram in which no job finished: every slot null, not a null array. const EMPTY_BINS = `'{${new Array(LATENCY_SLOTS).fill('NULL').join(',')}}'::int[]` // Both bins travel packed in one integer (wait * LATENCY_PACK + run), so the aggregate needs one // array and one filter. Two arrays, each with its own filter evaluated on every row of the table, // cost twice as much: measured on 2.5M rows, +55 ms on a ~560 ms pass packed, +80 to 120 ms apart. const LATENCY_PACK = 64 // date_part rather than extract: extract returns numeric on PostgreSQL 14 and later, and the numeric // arithmetic on every finished job was a large share of the histograms' cost. date_part is float8 // on every supported backend, as the bin math needs anyway. const waitSeconds = "date_part('epoch', (j.started_on - GREATEST(j.created_on, j.start_after)))" const runSeconds = "date_part('epoch', (j.completed_on - j.started_on))" // The slot is width_bucket's, written out: CockroachDB's width_bucket takes decimals, not float8. // ln(t / 10 ms) over ln(√2), plus one, clamped to slot 0 below 10 ms and the last slot past the end. // A duration of zero or less (stamps from two clocks) lands in slot 0. function latencyBin (seconds: string): string { const lo = Math.log(LATENCY_MIN_SECONDS) const width = Math.log(2) / 2 const raw = `floor((ln(GREATEST((${seconds})::float8, 0.001::float8)) - ${lo}::float8) / ${width}::float8)::int + 1` return `LEAST(GREATEST(${raw}, 0), ${LATENCY_SLOTS - 1})` } // The aggregate's array is unnested once per queue, in a lateral join after the aggregate, and // counted by packed value: at most LATENCY_SLOTS² rows (ps, ns) however many jobs finished. Both // histograms are read from those, so the per-job array is walked once rather than once each. function latencyCounts (packed: string): string { return `LEFT JOIN LATERAL ( SELECT array_agg(g.p) AS ps, array_agg(g.n) AS ns FROM (SELECT p, count(*)::int AS n FROM unnest(${packed}) AS u(p) GROUP BY p) g ) latency ON true` } // One histogram from latencyCounts: every slot in slot order, null where no job landed. A pass in // which nothing finished still writes the 48 slots (unnest of a null array is no rows, and the left // join keeps every slot), so a pass that counted says so, as its deltas do. floor() of a float // division rather than integer division, which CockroachDB answers in decimal. function latencySlotOf (which: 'wait' | 'run'): string { return which === 'wait' ? `floor(p / ${LATENCY_PACK}.0)::int` : `(p % ${LATENCY_PACK})::int` } function latencyHistogram (which: 'wait' | 'run'): string { return `(SELECT array_agg(c.n ORDER BY s.slot) FROM generate_series(0, ${LATENCY_SLOTS - 1}) AS s(slot) LEFT JOIN (SELECT ${latencySlotOf(which)} AS slot, sum(n)::int AS n FROM unnest(latency.ps, latency.ns) AS u(p, n) GROUP BY 1) c ON c.slot = s.slot)` } // Every count the monitor keeps, from one pass over the queue's table. // // Six of them are gauges — what the queue looks like right now. Three are not: // createdDelta counts the jobs created since the last pass, and completedDelta // and failedDelta the jobs that *finished*, which is the only way to answer // "how many jobs did this queue get through" from a table of current state. // Five hundred arriving and five hundred leaving looks identical to a still // queue in every gauge here. // // The window is the queue's own `delta_on`, joined in // rather than passed as a fixed interval. That watermark is what makes the three // counters exact across a skipped or backed-off pass: nothing is counted twice, // because the window starts where the last one ended, and nothing is missed, // because a late pass simply covers a longer window, and says so in // delta_seconds. A queue that has never been counted has a null watermark and // counts zero, which is the honest answer for a first pass that has nothing to // compare against. // // The join is against `queue`, which holds one row per queue — Postgres hashes // it once and probes per row. The alternative, a second pass over the job table // filtered on completed_on, would be a whole extra scan of the largest table in // the schema, and there is no index on that column to make it cheaper. export function getQueueStats (schema: string, table: string, queues: string[], throughput = false, window: DeltaWindowOptions = {}): SqlQuery { // A queued job is exactly one of blocked (waiting on a flow parent), deferred (start_after still // ahead) or ready, so the three add up to queuedCount. Blocked wins over deferred: a job whose // parent has not finished cannot run when its start_after comes round. const queued = `j.state < '${JOB_STATES.active}'` const blocked = `${queued} AND j.blocked` const deferred = `${queued} AND NOT j.blocked AND j.start_after > ${schema}.job_now()` const ready = `${queued} AND NOT j.blocked AND j.start_after <= ${schema}.job_now()` // Counted only with persistQueueStats. Otherwise the aggregate does what it did before throughput // existed: no join, no extra counts, no cost. The measured price is in the `persistQueueStats` docs. const end = deltaWindowEnd(schema, window.lag) const inWindow = (column: string) => `j.${column} >= ${deltaWindowStart('q', end, window.resetMax)} AND j.${column} < ${end}` // The true-up check: the same jobs recounted across the windows already recorded (see // trueUpSettled). More than the snapshots hold means something committed after its window was // counted, and the monitor follows up with trueUpQueueStats for this queue. const settled = (column: string) => `j.${column} >= q.h AND j.${column} < q.top` const counters = throughput ? { select: ` "createdDelta", "completedDelta", "failedDelta", ${latencyHistogram('wait')} as "waitBins", ${latencyHistogram('run')} as "runBins", "readyOldestSeconds", COALESCE("recount" > "settled", false) as "trueUp",`, counts: ` (count(*) FILTER (WHERE ${inWindow('created_on')}))::int as "createdDelta", (count(*) FILTER (WHERE j.state = '${JOB_STATES.completed}' AND ${inWindow('completed_on')}))::int as "completedDelta", (count(*) FILTER (WHERE j.state = '${JOB_STATES.failed}' AND ${inWindow('completed_on')}))::int as "failedDelta", array_agg(${latencyBin(waitSeconds)} * ${LATENCY_PACK} + ${latencyBin(runSeconds)}) FILTER (WHERE j.state IN ('${JOB_STATES.completed}', '${JOB_STATES.failed}') AND j.started_on IS NOT NULL AND ${inWindow('completed_on')}) as "latencyBins", round(extract(epoch from (${schema}.job_now() - min(GREATEST(j.created_on, j.start_after)) FILTER (WHERE ${ready}))))::int as "readyOldestSeconds", sum( CASE WHEN ${settled('created_on')} THEN 1 ELSE 0 END + CASE WHEN ${settled('completed_on')} AND j.state IN ('${JOB_STATES.completed}', '${JOB_STATES.failed}') THEN 1 ELSE 0 END )::int as "recount", max(q.settled) as "settled",`, join: `JOIN ( SELECT q.name, q.delta_on, a.h, t.top, t.settled FROM ${schema}.queue q${trueUpSettled(schema, 'q', window.trueUpMax)} WHERE q.name = ANY($1::text[]) ) q ON q.name = j.name`, lateral: latencyCounts('stats."latencyBins"') } : { select: '', counts: '', join: '', lateral: '' } return { text: ` SELECT name, "deferredCount", "queuedCount", "readyCount", "blockedCount", "activeCount", "failedCount", "totalCount",${counters.select} "singletonsActive" FROM ( SELECT j.name, (count(*) FILTER (WHERE ${deferred}))::int as "deferredCount", (count(*) FILTER (WHERE ${blocked}))::int as "blockedCount", (count(*) FILTER (WHERE ${ready}))::int as "readyCount", (count(*) FILTER (WHERE ${queued}))::int as "queuedCount", (count(*) FILTER (WHERE j.state = '${JOB_STATES.active}'))::int as "activeCount", (count(*) FILTER (WHERE j.state = '${JOB_STATES.failed}'))::int as "failedCount", count(*)::int as "totalCount",${counters.counts} array_agg(j.singleton_key) FILTER (WHERE j.policy IN ('${QUEUE_POLICIES.singleton}','${QUEUE_POLICIES.stately}') AND j.state = '${JOB_STATES.active}') as "singletonsActive" FROM ${schema}.${table} j ${counters.join} WHERE j.name = ANY($1::text[]) GROUP BY 1 ) stats ${counters.lateral} `, values: [queues] } } // Length of the recent-ready-count sliding window kept on queue.ready_history for the dashboard // sparkline. One sample is appended per monitor cycle (default 60s), so this is roughly the last // READY_HISTORY_SIZE minutes of trend. Sized to comfortably render the sparkline (the widest is the // ~160px detail card) without over-collecting, more points than pixels add nothing visible. export const READY_HISTORY_SIZE = 60 /* eslint-disable no-restricted-syntax -- how long this transaction has held its snapshot: real elapsed time, not job time */ const PIN_SECONDS_SQL = 'EXTRACT(EPOCH FROM (clock_timestamp() - transaction_timestamp()))::float8' /* eslint-enable no-restricted-syntax */ export function cacheQueueStats (schema: string, table: string, queues: string[], noAdvisoryLocks?: boolean, throughput?: boolean, window: DeltaWindowOptions = {}): string { const statsQuery = getQueueStats(schema, table, queues, throughput, window) // The aggregate only produces these when counting, so the assignment has to // disappear with them rather than reference a column that is not there. const throughputSet = throughput ? throughputAssignments(deltaWindowEnd(schema, window.lag), window.resetMax) : '' // Serialize the $1 parameter for use in the multi-statement transaction below const statsText = statsQuery.text.replaceAll('$1::text[]', serializeArrayParam(queues)) const lock = tryAdvisoryLock(schema, 'queue-stats', noAdvisoryLocks) // Two columns in here are not counts and are easy to mistake for incidental: // // monitor_on is stamped by this statement and by refreshQueueStats - the two that write counts - // and never by the claim in trySetQueueMonitorTime, which writes monitor_claim_on instead. So capturedOn means exactly one thing, "when these counts were // written", and a claimed pass that skipped the aggregate (backed off, or beaten to the lock // below) leaves it aging rather than advancing it over unchanged counts. // // ready_history is an always-on sliding window of recent ready counts for the dashboard sparkline, // maintained here independently of persistQueueStats: prepend the newest sample, keep the newest // READY_HISTORY_SIZE, stored newest-first. Built with unnest + array_agg rather than array slicing, // which CockroachDB lacks. // // pinSeconds is how long this transaction has held the MVCC horizon, measured by the server // rather than by a stopwatch around the call. The two are not the same number: executeSql waits // for a pool connection inside the client's timing window, and a saturated pool inflates it // without pinning anything (measured: a 1 ms transaction timed at 1.97 s behind one busy // connection). Backing off on that number would defer monitoring for pool contention and then // blame the job table for it. transaction_timestamp() is the BEGIN; clock_timestamp() is // volatile, so this is evaluated as rows are returned - after the aggregate CTE has run, which is // the part that pins. It undercounts the tail (the remaining rows and the COMMIT), the safe // direction to be wrong in. // // The queue rows are locked in name order (see queueRowLock) by the subquery that feeds the update. // Its count over stats is there only to finish the aggregate before the first row is locked, so no // queue row stays locked through the job table scan. const sql = ` WITH ${lock.cte}stats AS (SELECT * FROM (${statsText}) agg WHERE true${lock.guard}) UPDATE ${schema}.queue SET deferred_count = COALESCE(stats."deferredCount", 0), blocked_count = COALESCE(stats."blockedCount", 0), queued_count = COALESCE(stats."queuedCount", 0), ready_count = COALESCE(stats."readyCount", 0), active_count = COALESCE(stats."activeCount", 0), failed_count = COALESCE(stats."failedCount", 0), total_count = COALESCE(stats."totalCount", 0),${throughputSet} singletons_active = stats."singletonsActive", monitor_on = ${schema}.job_now(), ready_history = ( SELECT COALESCE(array_agg(v ORDER BY ord), '{}'::int[]) FROM ( SELECT v, ord FROM ( SELECT COALESCE(stats."readyCount", 0)::int AS v, 0::bigint AS ord UNION ALL SELECT h.v, h.ord FROM unnest(COALESCE(queue.ready_history, '{}'::int[])) WITH ORDINALITY AS h(v, ord) ) merged ORDER BY ord LIMIT ${READY_HISTORY_SIZE} ) capped ) FROM ( SELECT name FROM ${schema}.queue WHERE name = ANY(${serializeArrayParam(queues)})${lock.guard} AND (SELECT count(*) FROM stats) >= 0 ORDER BY name FOR NO KEY UPDATE ) q LEFT JOIN stats ON stats.name = q.name WHERE queue.name = q.name${lock.guard} RETURNING queue.name, queue.queued_count as "queuedCount", queue.warning_queued as "warningQueueSize",${throughput ? '\n COALESCE(stats."trueUp", false) as "trueUp",' : ''} ${PIN_SECONDS_SQL} as "pinSeconds" ` // transaction(), not locked(): the lock is taken inside the statement with try rather than by a // preceding blocking one. The wrapper is still wanted for its SET LOCAL timeouts. return transaction(sql) } // Recompute one queue's counts from the job table and write them back to the queue-table cache // (including monitor_on, so subsequent reads are served from cache), returning the fresh counts. // Backs getQueueStats(name, { force: true }) and the first read of a never-monitored queue. // // Shares the monitor's try-lock, and needs it more than the monitor does. Correctness never // required one - concurrent forced refreshes are idempotent, each a valid point-in-time snapshot // with last write wins - but the horizon does: nothing else gates this path. Two instances reading // the same stale cache both decide to refresh and both scan the whole table, and two *different* // queue names on the same table do not even contend for the cache entry, so the concurrency is // bounded only by how many instances are running. Sharing the monitor's key also keeps a forced // read from overlapping a supervise pass already in flight. // // Returning no rows is the documented outcome of losing the race: the caller serves the cached // counts it is already holding, which is what a concurrent refresh would have converged on anyway. // Single statement, so the lock is scoped to its implicit transaction and released with it. // // firstCapture is the one case that must not lose. A queue with no capture yet has no cached counts // to fall back on - only the columns' default zeros - so skipping the scan would answer "this queue // is empty" for a queue holding thousands of jobs, which is a fabricated number rather than a stale // one. And the lock key is global rather than per-queue, so a first read collides with any supervise // aggregate anywhere in the schema, which on exactly the slow-aggregate deployments this subsystem // targets is not a rare race. Skipping the lock is bounded: it can happen at most once per queue, // because the scan it runs is what populates the cache that gates every later read. // // It never counts throughput. Only the monitor does, so there is one writer of the counting window; // this leaves delta_on alone and the next monitor pass counts across it. export function refreshQueueStats (schema: string, table: string, name: string, options: { noAdvisoryLocks?: boolean, firstCapture?: boolean } = {}): string { const statsQuery = getQueueStats(schema, table, [name]) const statsText = statsQuery.text.replace('$1::text[]', serializeArrayParam([name])) const lock = tryAdvisoryLock(schema, 'queue-stats', options.noAdvisoryLocks || options.firstCapture) return ` WITH ${lock.cte}stats AS (SELECT * FROM (${statsText}) agg WHERE true${lock.guard}) UPDATE ${schema}.queue SET deferred_count = COALESCE(stats."deferredCount", 0), blocked_count = COALESCE(stats."blockedCount", 0), queued_count = COALESCE(stats."queuedCount", 0), ready_count = COALESCE(stats."readyCount", 0), active_count = COALESCE(stats."activeCount", 0), failed_count = COALESCE(stats."failedCount", 0), total_count = COALESCE(stats."totalCount", 0), singletons_active = stats."singletonsActive", monitor_on = ${schema}.job_now() FROM ( SELECT q.name FROM unnest(${serializeArrayParam([name])}) AS q(name) ) q LEFT JOIN stats ON stats.name = q.name WHERE queue.name = q.name${lock.guard} RETURNING queue.name, queue.deferred_count as "deferredCount", queue.queued_count as "queuedCount", queue.ready_count as "readyCount", queue.active_count as "activeCount", queue.failed_count as "failedCount", queue.total_count as "totalCount", queue.created_delta as "createdDelta", queue.completed_delta as "completedDelta", queue.failed_delta as "failedDelta", queue.delta_seconds as "deltaSeconds", queue.delta_on as "deltaOn", queue.monitor_on as "capturedOn" ` } // Serialize a string array for embedding directly in SQL as PostgreSQL array literal export function serializeArrayParam (values: string[]): string { const escaped = values.map(v => `'${v.replace(SINGLE_QUOTE_REGEX, "''")}'`) return `ARRAY[${escaped.join(',')}]::text[]` } // Serialize a JSON-serializable value for embedding directly in SQL as a quoted literal export function serializeJsonParam (value: unknown): string { return `'${JSON.stringify(value).replace(SINGLE_QUOTE_REGEX, "''")}'` } export function transaction (query: string | string[]): string { const sql = Array.isArray(query) ? query.join(';\n') : query return ` BEGIN; SET LOCAL lock_timeout = 30000; SET LOCAL idle_in_transaction_session_timeout = 30000; ${sql}; COMMIT; ` } export function locked (schema: string, query: string | string[], key?: string, noAdvisoryLocks?: boolean): string { const statements = Array.isArray(query) ? query : [query] return transaction(noAdvisoryLocks ? statements : [advisoryLock(schema, key), ...statements]) } // normalizeSchemaName, not resolveSchemaName: the key is opaque to postgres and never compared // against the catalog, so it only has to agree across instances on the same schema. See the note // on the helper. export function advisoryLockKey (schema: string, key?: string) { return `('x' || encode(sha224((current_database() || '.pgboss.${normalizeSchemaName(schema)}${key || ''}')::bytea), 'hex'))::bit(64)::bigint` } function advisoryLock (schema: string, key?: string) { return `SELECT pg_advisory_xact_lock(${advisoryLockKey(schema, key)})` } // The stats aggregate takes its lock with try, not wait, and abandons the pass if another instance // already holds it. // // Waiting is worse than skipping here, and not by a little. A backend blocked on a lock has already // sent BEGIN and taken its snapshot, so it advertises backend_xmin for the whole wait - confirmed // on pg_stat_activity: wait_event_type 'Lock' with backend_xmin set. Blocking to run a second // whole-table scan would therefore hold the horizon for the first instance's aggregate *and* its // own, turning one 3s pin into a contiguous 6s one, and three instances into 9s. That is the // overlapping-analytics regime this whole subsystem exists to prevent, manufactured by the lock // meant to prevent duplicate work. // // So the loser does nothing. trySetQueueMonitorTime stamped monitor_claim_on a statement earlier, // when it claimed the interval and before this lock was attempted, so those queues sit out one // interval rather than retrying at once - their counts end up at most one monitorIntervalSeconds // behind, and monitor_on keeps aging so capturedOn says so. // Against a staleness budget measured in hours, losing one sample costs a gap in the sparkline; // winning the race to duplicate a scan nobody asked for costs the job table. // // The guard is a one-time filter, so a lost race skips the scan rather than running and discarding // it, verified on the plan: `Parallel Seq Scan on job (never executed)`. It also has to gate the // UPDATE, not just the aggregate: an empty stats CTE still LEFT JOINs, and the COALESCE(…, 0) would // write every count to zero. function tryAdvisoryLock (schema: string, key: string, noAdvisoryLocks?: boolean) { if (noAdvisoryLocks) { return { cte: '', guard: '' } } return { cte: `lock AS MATERIALIZED (SELECT pg_try_advisory_xact_lock(${advisoryLockKey(schema, key)}) AS got),\n `, guard: ' AND (SELECT got FROM lock)' } } export function assertMigration (schema: string, version: number) { // raises 'division by zero' if already on desired schema version return `SELECT version::int/(version::int-${version}) from ${schema}.version` } export function findJobs (schema: string, table: string, options: { queued: boolean, byKey: boolean, byData: boolean, byId: boolean }) { const { queued, byKey, byData, byId } = options let paramIndex = 1 const whereConditions = [] if (byId) { ++paramIndex whereConditions.push(`AND id = $${paramIndex}`) } if (byKey) { ++paramIndex whereConditions.push(`AND singleton_key = $${paramIndex}`) } if (byData) { ++paramIndex whereConditions.push(`AND data @> $${paramIndex}::text::jsonb`) } if (queued) { whereConditions.push(`AND state < '${JOB_STATES.active}'`) } return ` SELECT ${JOB_COLUMNS_ALL} FROM ${schema}.${table} WHERE name = $1 ${whereConditions.join('\n ')} ` } export function getJobById (schema: string, table: string) { return ` SELECT ${JOB_COLUMNS_ALL} FROM ${schema}.${table} WHERE name = $1 AND id = $2 ` } // Pass `deps` to embed the payload as a literal (parameter-less) so the statement can be // concatenated into a flow batch; omit it to get the parameterized ($1) form. export function insertDependencies (schema: string, deps?: unknown[]) { const sql = ` INSERT INTO ${schema}.job_dependency (child_name, child_id, parent_name, parent_id) SELECT child_name, child_id, parent_name, parent_id FROM json_to_recordset($1::text::json) AS x ( child_name text, child_id uuid, parent_name text, parent_id uuid ) ON CONFLICT DO NOTHING ` return deps ? sql.replace('$1', () => serializeJsonParam(deps)) : sql } export function getDependencies (schema: string) { return ` SELECT parent_name as "parentName", parent_id as "parentId" FROM ${schema}.job_dependency WHERE child_name = $1 AND child_id = $2 ` } export function getDependents (schema: string) { return ` SELECT child_name as "childName", child_id as "childId" FROM ${schema}.job_dependency WHERE parent_name = $1 AND parent_id = $2 ` } export function cleanupDependencies (schema: string, table: string, queues: string[], noAdvisoryLocks?: boolean): string { const sql = ` DELETE FROM ${schema}.job_dependency WHERE (child_name = ANY(${serializeArrayParam(queues)}) AND NOT EXISTS ( SELECT 1 FROM ${schema}.${table} j WHERE j.name = child_name AND j.id = child_id )) OR (parent_name = ANY(${serializeArrayParam(queues)}) AND NOT EXISTS ( SELECT 1 FROM ${schema}.${table} j WHERE j.name = parent_name AND j.id = parent_id )) ` return locked(schema, sql, table + 'cleanupDependencies', noAdvisoryLocks) } export function getBlockedKeys (schema: string, table: string) { return ` SELECT DISTINCT singleton_key as "singletonKey" FROM ${schema}.${table} WHERE name = $1 AND state = '${JOB_STATES.failed}' AND policy = '${QUEUE_POLICIES.key_strict_fifo}' ` } export function getNextBamCommand (schema: string, { useLiveness = false }: { useLiveness?: boolean } = {}) { // Head-of-line note (shared by both variants): process all 'pending' commands (oldest first) // before retrying any 'failed' or stale 'in_progress' one, so a permanently-failing (or // crashed-mid-flight) command can't sit at the head of the queue and starve everything behind it. // Within a status, created_on preserves enqueue order. if (!useLiveness) { // Timeout-only path for engines without pg_stat_progress_create_index (CockroachDB/YugabyteDB): // a stuck in_progress row is reclaimed purely on the 24h fallback. No CONCURRENTLY healing, so no // reclaimed flag is emitted (bam.ts skips healing when noIndexProgressView is set anyway). return ` WITH candidate AS ( SELECT c.id, c.status AS prior_status, c.started_on AS prior_started_on FROM ${schema}.bam c WHERE ( c.status IN ('pending', 'failed') OR (c.status = 'in_progress' AND c.started_on < ${schema}.job_now() - interval '${BAM_STALE_SECONDS} seconds') ) AND NOT EXISTS ( SELECT 1 FROM ${schema}.bam g WHERE g.status = 'in_progress' AND g.started_on >= ${schema}.job_now() - interval '${BAM_STALE_SECONDS} seconds' ) ORDER BY (c.status != 'pending'), c.created_on LIMIT 1 ) UPDATE ${schema}.bam b SET status = 'in_progress', started_on = ${schema}.job_now() FROM candidate WHERE b.id = candidate.id AND b.status = candidate.prior_status AND b.started_on IS NOT DISTINCT FROM candidate.prior_started_on RETURNING b.id, b.name, b.version, b.status, b.queue, b.table_name as "table", b.command, b.error, b.created_on as "createdOn", b.started_on as "startedOn", b.completed_on as "completedOn", (candidate.prior_status <> 'pending') as reattempt, candidate.prior_status as "priorStatus", candidate.prior_started_on::text as "priorStartedOn", b.started_on::text as "claimedStartedOn" ` } // Native-Postgres liveness path. An in_progress row counts as "stale" (reclaimable) when it is past // the grace window AND no backend is actually building its index right now. The same predicate, // negated, defines a genuinely-live command that must still block the queue, so a running build is // NEVER reclaimed (no matter how long it runs), and a dead one recovers within the grace window. // There is deliberately no 24h absolute cap here: liveBuild=true always means a build is in flight, so // capping on elapsed time would reclaim a genuinely-running build and start a second // CREATE INDEX CONCURRENTLY on the same index. The exact double-build this path exists to prevent. // (The timeout-only path's BAM_STALE_SECONDS fallback covers engines with no way to detect liveness.) // // liveBuild(tableCol): is a CREATE INDEX CONCURRENTLY actively building this table's index right now? // Detected via pg_locks, NOT pg_stat_progress_create_index, and that choice is load-bearing for // multi-instance safety. pg_stat_progress_* is filtered to the querying role's OWN backends (only a // superuser or a member of pg_read_all_stats sees another role's builds), so a progress-view check // silently reads a peer's live build as "dead" whenever pg-boss instances connect under different DB // roles, and the heal step (bamHealProbe/bamHealDrop in bam.ts) would then DROP INDEX CONCURRENTLY a // live index mid-build, racing the builder into a double CREATE. pg_locks, by contrast, is cluster-wide // and visible to every role (verified empirically). CREATE INDEX CONCURRENTLY holds a // ShareUpdateExclusiveLock on the target table for the ENTIRE build and releases it the instant the // statement finishes or the backend dies, so a granted SUExclusive lock on the row's table is a // crash-safe, role-agnostic "build in flight" signal, instances may run under different roles with no // loss of safety. Ordinary queue DML never takes SUExclusive (it uses AccessShare/RowShare/ // RowExclusive), so it can't false-trigger; a concurrent autovacuum/ANALYZE on the same table DOES take // SUExclusive and reads as "live", but that false positive is in the SAFE direction. It only briefly // DEFERS a reclaim, never drops a live index. Scoped to the current database (l.database) because // pg_locks is cluster-wide while relation OIDs are only unique per database. Correlating on the table // is enough because the queue runs only one in_progress command at a time. const liveBuild = (tableCol: string) => `EXISTS ( SELECT 1 FROM pg_locks l WHERE l.locktype = 'relation' AND l.granted AND l.mode = 'ShareUpdateExclusiveLock' AND l.database = (SELECT oid FROM pg_database WHERE datname = current_database()) AND l.relation = to_regclass(quote_ident('${resolveSchemaName(schema)}') || '.' || quote_ident(${tableCol})) )` const stale = (startedCol: string, tableCol: string) => `( ${startedCol} < ${schema}.job_now() - interval '${BAM_LIVENESS_GRACE_SECONDS} seconds' AND NOT ${liveBuild(tableCol)} )` // The candidate CTE's FOR UPDATE OF c SKIP LOCKED is defense-in-depth against a double-claim. The // upstream trySetBamTime throttle serializes healers to one per interval, but a build running // longer than bamIntervalSeconds lets a second instance pass the throttle while the first is still // working. SKIP LOCKED makes the claim itself mutually exclusive: whichever instance locks the head // row wins, and the other skips it (and, with LIMIT 1, claims nothing) rather than re-running the // same UPDATE and re-driving the command with a stale prior_status. OF c scopes the lock to the // candidate row only, so the NOT EXISTS probe over bam g is never itself locked. Postgres-only: // SKIP LOCKED is deliberately absent from the timeout-only variant, which also serves CockroachDB // and YugabyteDB, where it performs poorly and can skip unexpectedly. // // The returned `reattempt` asks whether this command was already tried once (a stale-reclaimed // in_progress, or a prior 'failed'). Either way an interrupted or failed CREATE INDEX CONCURRENTLY // may have left an INVALID index, so the runner heals (drop-then-rebuild) before re-running. Fresh // 'pending' rows have nothing to heal. return ` WITH candidate AS ( SELECT c.id, c.status AS prior_status, c.started_on AS prior_started_on FROM ${schema}.bam c WHERE ( c.status IN ('pending', 'failed') OR (c.status = 'in_progress' AND ${stale('c.started_on', 'c.table_name')}) ) AND NOT EXISTS ( SELECT 1 FROM ${schema}.bam g WHERE g.status = 'in_progress' AND NOT ${stale('g.started_on', 'g.table_name')} ) ORDER BY (c.status != 'pending'), c.created_on LIMIT 1 FOR UPDATE OF c SKIP LOCKED ) UPDATE ${schema}.bam b SET status = 'in_progress', started_on = ${schema}.job_now() FROM candidate WHERE b.id = candidate.id AND b.status = candidate.prior_status AND b.started_on IS NOT DISTINCT FROM candidate.prior_started_on RETURNING b.id, b.name, b.version, b.status, b.queue, b.table_name as "table", b.command, b.error, b.created_on as "createdOn", b.started_on as "startedOn", b.completed_on as "completedOn", (candidate.prior_status <> 'pending') as reattempt, candidate.prior_status as "priorStatus", candidate.prior_started_on::text as "priorStartedOn", b.started_on::text as "claimedStartedOn" ` } // Derives the drop-then-rebuild heal step for a re-attempted BAM index build: an interrupted // CREATE INDEX CONCURRENTLY leaves an INVALID index in the catalog, and the command's own // IF NOT EXISTS would then skip forever (marking the row done while the index stays invalid). // Dropping it first lets the re-run rebuild cleanly. Returns null for non-index commands (which // have nothing to heal) so the caller just re-runs them. DROP ... CONCURRENTLY runs outside a // transaction, matching the CREATE, and IF EXISTS makes it a no-op when there's nothing to drop. export function bamHealDrop (schema: string, command: string): string | null { const match = command.match(/CREATE\s+(?:UNIQUE\s+)?INDEX\s+CONCURRENTLY\s+(?:IF\s+NOT\s+EXISTS\s+)?("?[\w$]+"?)/i) if (!match) return null return `DROP INDEX CONCURRENTLY IF EXISTS ${schema}.${match[1]}` } // Probe run before bamHealDrop: returns `invalid = true` only when the index the re-attempted command // would build already exists AND is INVALID (an interrupted CREATE INDEX CONCURRENTLY left a stub). // A build that actually finished but whose BAM row was never marked completed, e.g. a graceful stop // landed between the CREATE succeeding and markCompleted, leaves a VALID index; dropping that would // tear down a live production index for the whole rebuild window, so the caller must NOT heal it (the // re-run's own IF NOT EXISTS then no-ops and just marks the row done). Returns null for non-index // commands and for an absent index (no row), where there is nothing to drop. Mirrors bamHealDrop's // CONCURRENTLY-only recognition so it is non-null exactly when bamHealDrop is. export function bamHealProbe (schema: string, command: string): string | null { const match = command.match(/CREATE\s+(?:UNIQUE\s+)?INDEX\s+CONCURRENTLY\s+(?:IF\s+NOT\s+EXISTS\s+)?"?([\w$]+)"?/i) if (!match) return null return ` SELECT NOT i.indisvalid AS invalid FROM pg_class c JOIN pg_index i ON i.indexrelid = c.oid JOIN pg_namespace n ON n.oid = c.relnamespace WHERE n.nspname = '${resolveSchemaName(schema).replace(SINGLE_QUOTE_REGEX, "''")}' AND c.relname = '${match[1].replace(SINGLE_QUOTE_REGEX, "''")}' ` } // Undoes a claim that never ran: hands the row back exactly as getNextBamCommand found it, restoring // both the status and the started_on the claim overwrote. Used when a stop lands between the claim // UPDATE returning and the runner picking the command up. Leaving the row in_progress instead would // block every other BAM command until it went stale - BAM_LIVENESS_GRACE_SECONDS on native Postgres, // BAM_STALE_SECONDS (24 hours) on the timeout-only backends - and restoring started_on is what makes a // released stale-in_progress row immediately reclaimable again rather than starting that clock over. export function releaseBamCommand (schema: string, id: string, priorStatus: string, priorStartedOn: string | null, claimedStartedOn: string) { // priorStartedOn arrives as the text rendering of the original timestamptz (the claim casts it in // SQL rather than round-tripping a JS Date), so it carries its offset and full precision back. const startedOn = priorStartedOn == null ? 'NULL' : `'${priorStartedOn.replace(SINGLE_QUOTE_REGEX, "''")}'::timestamptz` // Compare-and-swap on the claim this runner actually took, not on the id alone. The timeout-only // claim has no SKIP LOCKED, so two overlapping claims can both return the // same row - verified: both get rowCount 1 and the same prior_status. Releasing on the id alone // would then reset a row a peer is actively building back to 'pending', and the next poll would // start a second CREATE INDEX CONCURRENTLY on the same index. Matching started_on makes the release // a no-op unless the claim is still ours. return ` UPDATE ${schema}.bam SET status = '${priorStatus.replace(SINGLE_QUOTE_REGEX, "''")}', started_on = ${startedOn} WHERE id = '${id}' AND status = 'in_progress' AND started_on = '${claimedStartedOn.replace(SINGLE_QUOTE_REGEX, "''")}'::timestamptz ` } export function setBamCompleted (schema: string, id: string) { // error is cleared: a row reaching 'completed' is commonly a retry of a prior 'failed' attempt, and // getBamEntries() returns error for every status - a stale message would keep reading as a failure // on a command that succeeded. return ` UPDATE ${schema}.bam SET status = 'completed', completed_on = ${schema}.job_now(), error = NULL WHERE id = '${id}' ` } export function setBamFailed (schema: string, id: string, error: string) { const escapedError = error.replace(/'/g, "''") return ` UPDATE ${schema}.bam SET status = 'failed', error = '${escapedError}', completed_on = ${schema}.job_now() WHERE id = '${id}' ` } export function getBamStatus (schema: string) { return ` SELECT status, count(*)::int as count, max(created_on) as "lastCreatedOn" FROM ${schema}.bam GROUP BY status ` } export function getBamEntries (schema: string) { return ` SELECT id, name, version, status, queue, table_name as "table", command, error, created_on as "createdOn", started_on as "startedOn", completed_on as "completedOn" FROM ${schema}.bam ORDER BY version, created_on ` } // --- drift detection: live-catalog probes --- // // pg-boss-specific probes that read the live database for drift detection: partitioning topology and // in-flight BAM builds. Their results feed the generic engine in drifter.ts, which knows nothing about // pg-boss. // Probe for the default partition; its presence means the partitioned architecture is in use, so // drift detection reads it from the live database rather than trusting a boss config flag. export function jobCommonExists (schema: string) { return `SELECT to_regclass('${schema}.${COMMON_JOB_TABLE}') as name` } // Per-queue partition tables and their policy, so the expected index set can be computed per table. export function getManagedQueuePartitions (schema: string) { return `SELECT table_name as "table", policy FROM ${schema}.queue WHERE partition = true` } // The CREATE INDEX command text of every BAM row not yet completed, used to tell a genuinely-missing // index apart from one an async build is still working on (or retrying after a failure). export function getIncompleteBamCommands (schema: string) { return `SELECT command FROM ${schema}.bam WHERE status <> 'completed'` } // Extracts the index name a BAM command builds, so an incomplete build can be matched to a missing // index. Mirrors the CREATE INDEX shapes bamHealDrop recognises. export function bamCommandIndexName (command: string): string | null { const match = command.match(/CREATE\s+(?:UNIQUE\s+)?INDEX\s+(?:CONCURRENTLY\s+)?(?:IF\s+NOT\s+EXISTS\s+)?"?([\w$]+)"?/i) return match ? match[1] : null } // --- expected schema shape (all derived from the generated manifest) --- // // The expected tables/columns/constraints/functions/indexes/enum pg-boss compares the live catalog // against. Every one is read from schema.json (regenerate with `npm run gen:manifest` after a // DDL change) with the placeholder schema substituted back, so nothing here duplicates the DDL as // hand-written literals. The only hand-maintained pg-boss knowledge is the per-queue partition index // distribution rule just below. These produce the `expected` inputs to drifter.ts's generic diff. interface QueuePartition { table: string policy?: string | null } // job_iN partial indexes that gate on a queue policy: a per-queue partition table (partition:true) // only receives the index for its own policy. The shared // job_common table and a non-partitioned job table carry all of them at once. Keep in sync with the // createIndexJobPolicy* builders and the ELSIF ladder in createQueueFunction. const POLICY_JOB_INDEXES: Record = { 1: QUEUE_POLICIES.short, 2: QUEUE_POLICIES.singleton, 3: QUEUE_POLICIES.stately, 6: QUEUE_POLICIES.exclusive, 8: QUEUE_POLICIES.key_strict_fifo, 10: QUEUE_POLICIES.key_strict_fifo } // job_iN indexes with no policy gate, created on every job table regardless of policy // (throttle i4, fetch i11, group-concurrency i7, blocking i9, source root i12). 5 is absent, not // missing: the fetch index was replaced in v40 and the retired number is not reused. const BASE_JOB_INDEXES = [4, 7, 9, 11, 12] // The fixed (non-job) managed tables; job/job_common/partitions are handled separately. const FIXED_MANAGED_TABLES = ['version', 'queue', 'schedule', 'subscription', 'bam', 'warning', 'queue_stats', 'job_dependency', 'instance'] // Selects the manifest section for the live architecture, and substitutes the real schema name back in // for the placeholder the manifest stores. function manifestSection (partitioned: boolean) { return partitioned ? schemaManifest.partitioned : schemaManifest.nonPartitioned } function applyManifestSchema (text: string, schema: string): string { return text.split(schemaManifest.schemaToken).join(schema) } // The argument types in a manifest function's definition, without names or defaults, as DROP FUNCTION // takes them: "create_queue(queue_name text, options jsonb)" gives "text, jsonb". function functionArgTypes (def: string): string { const args = def.slice(def.indexOf('(') + 1, def.indexOf(')')) return args.split(',') .map(arg => arg.trim().replace(/\s+DEFAULT\s+.*$/i, '')) .filter(Boolean) .map(arg => arg.split(/\s+/).slice(1).join(' ')) .join(', ') } // Removes everything create() installs, for a schema pg-boss shares with other objects; a schema of // its own is simply dropped. Read from the manifest, so it covers whatever this version installs. // Dropping job takes every queue's own table with it: a partitioned queue's table is a partition of // job, and without partitioning every queue's jobs are in job itself. queue_stats' daily partitions // go with queue_stats the same way. export function uninstall (schema: string, partitioned = true): string { const section = manifestSection(partitioned) const tables = section.tables.map(table => `${schema}.${table}`).join(', ') const functions = section.functions.map(fn => `${schema}.${fn.name}(${functionArgTypes(fn.def)})`).join(', ') // Functions first: CockroachDB records a function as depending on the tables its body names. return [ `DROP FUNCTION IF EXISTS ${functions};`, `DROP TABLE IF EXISTS ${tables};`, `DROP TYPE IF EXISTS ${schema}.job_state;` ].join('\n') } // The job_state enum values in declaration order, from the manifest (both sections carry the same enum). // Order is significant. The numeric base type makes created < retry < … < failed load-bearing. export const EXPECTED_JOB_STATES: readonly string[] = schemaManifest.partitioned.enum // The tables pg-boss expects to exist. The manifest lists the fixed tables plus job (and job_common in // partitioned mode); per-queue partition tables are dynamic, so they are appended from the live set. export function expectedManagedTables (schema: string, partitioned: boolean, partitions: QueuePartition[] = []): string[] { const tables = [...manifestSection(partitioned).tables] if (partitioned) for (const p of partitions) tables.push(p.table) return tables } // The columns pg-boss expects on each managed table, taken from the manifest (introspected column type, // nullability, and default). Fixed tables carry the full defaults + types maps so column defaults, data // types, and nullability are diffed; the job table (shared by job_common and every per-queue partition // in partitioned mode) is name-only. Its FKs are profile-dependent and keep_until's interval default is // not worth pinning per-partition. A table with no live columns is skipped by the diff, so listing one // that does not exist yet is harmless. export function expectedManagedColumns (schema: string, partitioned: boolean, partitions: QueuePartition[] = []): ExpectedColumns[] { const fixed = new Set(FIXED_MANAGED_TABLES) const byTable = new Map() const jobColumns: string[] = [] for (const c of manifestSection(partitioned).columns) { if (c.table === BASE_JOB_TABLE) { jobColumns.push(c.column); continue } if (!fixed.has(c.table)) continue // job_common / anything else derives from the job template below let entry = byTable.get(c.table) if (!entry) byTable.set(c.table, entry = { table: c.table, columns: [], defaults: {}, types: {} }) entry.columns.push(c.column) entry.types![c.column] = { type: applyManifestSchema(c.type, schema), notNull: c.notNull } if (c.default != null) entry.defaults![c.column] = applyManifestSchema(c.default, schema) } const out = [...byTable.values()] out.push({ table: BASE_JOB_TABLE, columns: jobColumns }) if (partitioned) { out.push({ table: COMMON_JOB_TABLE, columns: jobColumns }) for (const p of partitions) out.push({ table: p.table, columns: jobColumns }) } return out } // The constraints pg-boss expects on each FIXED managed table, taken verbatim from the manifest's // pg_get_constraintdef capture (job/job_common/partitions are excluded: their FKs are profile-dependent // DEFERRABLE). Compared as a normalized set, so order is irrelevant. The fresh-install integration test // verifies the manifest matches live. export function expectedManagedConstraints (schema: string, partitioned: boolean): ExpectedConstraints[] { const fixed = new Set(FIXED_MANAGED_TABLES) const byTable = new Map() for (const { table, def } of manifestSection(partitioned).constraints) { if (!fixed.has(table)) continue const list = byTable.get(table) ?? byTable.set(table, []).get(table)! list.push(applyManifestSchema(def, schema)) } return [...byTable].map(([table, constraints]) => ({ table, constraints })) } // The functions pg-boss expects to exist in `schema`, from the manifest's pg_get_functiondef capture // (the partitioned architecture has the three job_table_* helpers plus partition-aware create_queue/ // delete_queue; non-partitioned mode has neither the helpers nor the partition branches). Each entry // carries the whitespace-normalised body used for the diff and the full statement for remediation. // Postgres stores a function body verbatim, so the manifest body compares equal to the live one. export function expectedManagedFunctions (schema: string, partitioned: boolean, options: { clockOverride?: boolean } = {}): ManagedFunction[] { return manifestSection(partitioned).functions.map(fn => { let def = applyManifestSchema(fn.def, schema) if (fn.name === 'job_now' && options.clockOverride) { def = def.replace(extractFunctionBody(def), ` ${clockOverrideBody(schema)} `) } return { name: fn.name, expectedBody: normalizeFunctionBody(extractFunctionBody(def)), definition: def.replace(/\s+/g, ' ').trim() } }) } // Computes the set of indexes pg-boss expects to exist in `schema`, sourced from the manifest. // `partitioned` reflects the live database (job_common present ⇒ partitioned architecture), so this // needs no boss config. The manifest holds every fixed index (the static ones plus the job_iN set on // job/job_common); each per-queue partition is dynamic, so its indexes are templated here from the // job_common set. The base indexes always, plus the one for the queue's policy. keys/predicate/ // definition are derived from the catalog-canonical pg_get_indexdef the manifest stores. export function expectedManagedIndexes (schema: string, partitioned: boolean, partitions: QueuePartition[] = []): ManagedIndex[] { const managed = (name: string, table: string, indexdef: string): ManagedIndex => { const def = applyManifestSchema(indexdef, schema) return { name, table, keys: indexKeysRaw(def), include: indexIncludeRaw(def), predicate: indexPredicateRaw(def), definition: displayIndexDefinition(def) } } const jobTable = partitioned ? COMMON_JOB_TABLE : BASE_JOB_TABLE const out: ManagedIndex[] = [] const jobIndexes: Array<{ n: number, def: string }> = [] // The manifest holds only the managed indexes (constraint-backing *_pkey indexes are excluded at // generation), so every row is an expectation. for (const idx of manifestSection(partitioned).indexes) { out.push(managed(idx.name, idx.table, idx.def)) const n = idx.table === jobTable ? idx.name.match(/_i(\d+)$/) : null if (n) jobIndexes.push({ n: Number(n[1]), def: idx.def }) } // Per-queue partition tables are dynamic, so template each applicable job_common index onto them. if (partitioned) { for (const p of partitions) { for (const { n, def } of jobIndexes) { if (!BASE_JOB_INDEXES.includes(n) && POLICY_JOB_INDEXES[n] !== p.policy) continue out.push(managed(`${p.table}_i${n}`, p.table, def.split(COMMON_JOB_TABLE).join(p.table))) } } } return out } // --- Index bloat detection and reindexing (issue #876) --------------------------------------- // // Autovacuum reclaims heap space but never returns btree pages to the OS, so a job index sizes // itself to the largest backlog its queue has ever held and stays there. A drained 1M-job peak // leaves job_common_i11 at 38 MB and job_common_pkey at 51 MB holding 1k live rows, and every // subsequent vacuum with dead tuples to clean does a full physical scan of every one of those // pages (measured: 17,969 buffer hits + 3,892 reads vs 750 hits / 0 reads after a rebuild), which // is what drains IO burst credits on managed Postgres. // // The bloat is a high-water mark rather than an unbounded leak, empty pages ARE recycled by later // inserts, so the trigger is density, not elapsed time. export const REINDEX_DEFAULTS = { // 1 MB. Below this the absolute waste is irrelevant and reltuples/relpages is noisy. minPages: 128, // A freshly rebuilt job index packs 140-170 entries per page; a bloated one holds ~0.2. Five is // three orders of magnitude clear of healthy on the indexes that carry fixed-width keys. maxEntriesPerPage: 5, // Density alone is not enough: singleton_key is unbounded text and is indexed by i1-i4, i6, i8 and // i10, so a healthy index over ~1.3 kB keys legitimately holds fewer than five entries per page // (measured: 150k distinct 1,286-byte keys, freshly built, 4.0 entries/page) and would be rebuilt // every interval forever. What separates that from bloat is the index's size against the size its // own live entries need, which pg_stats gives for free (measured on the two fixtures): // // reported bloat (drained 1M peak) 0 live entries, 4,861 pages 4861x // healthy wide-key index 150k x 1,291 bytes, 37,502 pages 1.6x // // Four keeps a wide margin on both sides, and unlike a comparison against the heap it does not // quietly stop working when the heap itself is bloated or cannot be truncated. minSizeRatio: 4, // 2 GB. A large index that still passes the density gate would carry a genuinely expensive // rebuild, so leave it to an operator who has chosen the moment. maxIndexBytes: 2 * 1024 * 1024 * 1024 } // float8, not bigint: pg_relation_size() and reltuples::bigint are int8, which node-postgres hands // back as a *string* (it registers no int8 parser), so the exported IndexBloat type would be lying // about `bytes` and `entries` and a consumer comparing them numerically would be comparing text. // int4 is not an option - the index that exceeds maxIndexBytes is exactly the one that overflows it. // float8 is exact to 2^53 bytes, and reltuples is a float4 estimate to begin with. // The pages the index's *live* entries actually need, from the average width of what it indexes. // pg_stats carries that per column, and. The part that matters here, since i1-i4 and i6 key on // `COALESCE(singleton_key, '')`, ANALYZE also collects it for an index's expression columns, keyed // by the index relation. Plain attnums read the table's stats, attnum 0 (an expression) falls back // to the index's own. 20 bytes covers the btree tuple header and line pointer. const INDEX_WIDTH_JOIN = `LEFT JOIN LATERAL ( SELECT COALESCE(sum(COALESCE(ts.avg_width, xs.avg_width, 0)), 0) AS key_width FROM unnest(x.indkey::int2[]) WITH ORDINALITY AS k(attnum, pos) LEFT JOIN pg_attribute a ON a.attrelid = x.indrelid AND a.attnum = k.attnum LEFT JOIN pg_stats ts ON ts.schemaname = n.nspname AND ts.tablename = t.relname AND ts.attname = a.attname LEFT JOIN pg_attribute ia ON ia.attrelid = i.oid AND ia.attnum = k.pos::int2 LEFT JOIN pg_stats xs ON xs.schemaname = n.nspname AND xs.tablename = i.relname AND xs.attname = ia.attname ) w ON true` const EXPECTED_PAGES = 'GREATEST(ceil((i.reltuples * (w.key_width + 20)) / 8192.0), 1)' const INDEX_STAT_COLUMNS = ` floor(i.reltuples)::float8 as entries, pg_relation_size(i.oid)::float8 as bytes,` function numericThreshold (value: unknown, fallback: number, name: string): number { if (value === undefined || value === null) return fallback assert(typeof value === 'number' && Number.isFinite(value) && value >= 0, `configuration assert: reindex.${name} must be a finite number >= 0`) return value } // The job tables in scope. Every queue row names the physical table its jobs live on, so the // subquery covers the whole installation on its own: `job_common` for an unpartitioned queue, // `j` for a partitioned one, and `job` itself where noTablePartitioning makes `job` a // plain table rather than a partitioned parent. // // The two literals cover the one thing it cannot. The shared table is created by the install with // all of its indexes and outlives every queue row, so it has to stay in scope before the first queue // exists and after the last one is dropped - which is precisely when a drained backlog has left // bloat behind and nothing points at the table any more. Whichever of the two names the install // shape doesn't use simply matches nothing, and on a partitioned install `job` is the parent, whose // indexes are partitioned (relkind 'I') and already excluded by the relkind = 'i' filter. function jobTableScope (schema: string, tables?: string[]) { return tables?.length ? `t.relname = ANY(${serializeArrayParam(tables)})` : `t.relname IN (SELECT table_name FROM ${schema}.queue UNION SELECT '${BASE_JOB_TABLE}' UNION SELECT '${COMMON_JOB_TABLE}')` } /** * Bloated job indexes, by the pg_class density heuristic. Deliberately avoids pgstattuple: it is an * extension that may not be installed, and relpages/reltuples separates a bloated index from a * healthy one by three orders of magnitude on their own. * * PostgreSQL-only (including PGlite, where both the bloat and REINDEX CONCURRENTLY behave normally). * Callers must gate on noReindex. See the note in boss.ts #reindex. * * `owned` is returned rather than filtered on, so an installation whose role cannot reindex still * gets told what is bloated. Only leaf indexes (relkind 'i') are candidates. pg-boss creates * i1..i10 directly on each partition, and a partitioned index (relkind 'I') cannot be reindexed * except through its leaves anyway. * * Caveat: relpages/reltuples are refreshed by VACUUM/ANALYZE, so they go stale where autovacuum is * disabled for the table. That is acceptable. The entire failure mode assumes autovacuum runs. */ export function getBloatedIndexes (schema: string, tables?: string[], options?: IndexBloatOptions): string { // These are interpolated into the statement, and supervise()/getReindexCommands() accept them // per call, where the constructor's validation never runs. Reject anything that isn't a finite // number here rather than trusting the caller. const minPages = numericThreshold(options?.minPages, REINDEX_DEFAULTS.minPages, 'minPages') const maxEntriesPerPage = numericThreshold(options?.maxEntriesPerPage, REINDEX_DEFAULTS.maxEntriesPerPage, 'maxEntriesPerPage') const minSizeRatio = numericThreshold(options?.minSizeRatio, REINDEX_DEFAULTS.minSizeRatio, 'minSizeRatio') // The `reltuples >= 0` gate is not redundant. PG 14+ writes -1 rather than 0 for a relation that // has never been vacuumed or analyzed (commit 3d351d916b2), and -1 would pass both thresholds // below - a negative count is under any density, and a negative expected size is under any page // count - on statistics that mean "unknown" rather than "empty". return ` SELECT i.relname as name, t.relname as table, i.relpages as pages, ${INDEX_STAT_COLUMNS} pg_has_role(current_user, i.relowner, 'USAGE') as owned FROM pg_class i JOIN pg_index x ON x.indexrelid = i.oid JOIN pg_class t ON t.oid = x.indrelid JOIN pg_namespace n ON n.oid = i.relnamespace ${INDEX_WIDTH_JOIN} WHERE n.nspname = '${resolveSchemaName(schema).replace(SINGLE_QUOTE_REGEX, "''")}' AND i.relkind = 'i' AND x.indisvalid AND ${jobTableScope(schema, tables)} AND i.relpages > ${minPages} AND i.reltuples >= 0 AND i.reltuples / GREATEST(i.relpages, 1) < ${maxEntriesPerPage} AND i.relpages > ${minSizeRatio} * ${EXPECTED_PAGES} ORDER BY pg_relation_size(i.oid) DESC ` } /** * Every owned leaf index on the job tables, ignoring the density gate. The target list for * `{ force: true }`. */ export function getJobIndexes (schema: string, tables?: string[]): string { return ` SELECT i.relname as name, t.relname as table, i.relpages as pages, ${INDEX_STAT_COLUMNS} pg_has_role(current_user, i.relowner, 'USAGE') as owned FROM pg_class i JOIN pg_index x ON x.indexrelid = i.oid JOIN pg_class t ON t.oid = x.indrelid JOIN pg_namespace n ON n.oid = i.relnamespace WHERE n.nspname = '${resolveSchemaName(schema).replace(SINGLE_QUOTE_REGEX, "''")}' AND i.relkind = 'i' AND x.indisvalid AND ${jobTableScope(schema, tables)} ORDER BY pg_relation_size(i.oid) DESC ` } /** * Invalid `*_ccnew` stubs left behind by an interrupted REINDEX CONCURRENTLY. Postgres appends * `_ccnew`, then `_ccnew1`, `_ccnew2`, ... on collision. These must be dropped before a retry or * they accumulate, and they are dead weight that vacuum still walks. Mirrors bamHealDrop's role for * interrupted CREATE INDEX CONCURRENTLY. */ export function getReindexLeftovers (schema: string, tables?: string[], noIndexProgressView = false, anyOwner = false): string { // A rebuild that is still running looks exactly like a leftover from one that died: its transient // index is invalid until the swap. Nothing in the catalog distinguishes them, so ask whether a // build is actually in flight on that table - the same liveness source bam.ts uses to decide a // claimed CREATE INDEX CONCURRENTLY is dead, and pg_stat_progress_create_index covers REINDEX // CONCURRENTLY too. Without it, a forced pass (which skips the interval claim) could drop the // stub another instance is in the middle of building. // // Best effort by nature: an unprivileged role sees only its own backends in the progress view, and // a backend that skips the view has no liveness to read at all. const liveBuild = noIndexProgressView ? '' : 'AND NOT EXISTS (SELECT 1 FROM pg_stat_progress_create_index p WHERE p.relid = x.indrelid)' // The background pass can only drop what it owns, so filtering there saves an error. The command // list is the opposite case - it is handed to an operator who may run it as the owning role, and // it already emits REINDEX for indexes this connection does not own, so withholding the DROP that // has to precede one would hand over a list that fails on the first statement. const owner = anyOwner ? '' : "AND pg_has_role(current_user, i.relowner, 'USAGE')" return ` SELECT i.relname as name FROM pg_class i JOIN pg_index x ON x.indexrelid = i.oid JOIN pg_class t ON t.oid = x.indrelid JOIN pg_namespace n ON n.oid = i.relnamespace WHERE n.nspname = '${resolveSchemaName(schema).replace(SINGLE_QUOTE_REGEX, "''")}' AND i.relkind = 'i' AND NOT x.indisvalid AND i.relname ~ '_ccnew[0-9]*$' AND ${jobTableScope(schema, tables)} ${owner} ${liveBuild} ` } // Index names come from the catalog, so they are raw identifiers that may need quoting (a queue // name can contain anything, and the derived partition/index names inherit it). The schema is // already validated and carries its own quoting - see assertPostgresObjectName. function quoteIdentifier (name: string) { return `"${name.replace(/"/g, '""')}"` } /** * The sources that can hold the MVCC horizon back, and the catalog each is read from. Split out so * an installation whose role cannot read one of them still gets an answer from the rest, and so the * warning can name which sources it could not check rather than implying a clean bill of health. * * `pg_stat_activity.backend_xmin` is readable by an unprivileged role (it is not one of the columns * restricted to the owning user), as are the other three views, so the degraded path is defensive * rather than expected. A managed provider may still revoke them. */ /* eslint-disable no-restricted-syntax -- these measure real backend and vacuum age against pg_stat_* timestamps Postgres wrote; a fake clock would compare two different clocks */ export const XMIN_HORIZON_SOURCES = { // Restricted to backends whose transaction was already open when the failed vacuum ran ($1). // Every backend executing a query advertises a backend_xmin, including the one asking this // question, so presence alone is not evidence. The ordering against a vacuum that reclaimed // nothing is. Compared as timestamps inside one statement rather than as two ages read from two // queries a moment apart, which would let a connection that started after the vacuum qualify. // // A NULL xact_start passes rather than being dropped, and that is the difference between seeing // the most common cause of a pinned horizon and being blind to it. pg_stat_activity hands an // ordinary role only pid, application_name, usename and backend_xmin for a backend owned by a // *different* role; xact_start reads NULL and query reads '' (measured on // PG18). Since `NULL < $1` is NULL, the strict form silently discarded exactly the holder an // operator most needs named (another application's analytical query) leaving every source NULL, // no holder attributable, and the warning withheld. Admitting them costs a possible false // positive from a young foreign backend, which in practice loses the max(age()) to any real // holder, and opaqueBackends below reports how many could not be timed. backends: 'SELECT max(age(backend_xmin)) FROM pg_catalog.pg_stat_activity WHERE backend_xmin IS NOT NULL AND (xact_start IS NULL OR xact_start < $1)', slots: 'SELECT max(age(xmin)) FROM pg_catalog.pg_replication_slots WHERE xmin IS NOT NULL', slotsCatalog: 'SELECT max(age(catalog_xmin)) FROM pg_catalog.pg_replication_slots WHERE catalog_xmin IS NOT NULL', standbys: 'SELECT max(age(backend_xmin)) FROM pg_catalog.pg_stat_replication WHERE backend_xmin IS NOT NULL', prepared: 'SELECT max(age(transaction)) FROM pg_catalog.pg_prepared_xacts' } as const export type XminHorizonSource = keyof typeof XMIN_HORIZON_SOURCES export const XMIN_HORIZON_QUERY_SOURCES = Object.keys(XMIN_HORIZON_SOURCES) as readonly XminHorizonSource[] /** * How far behind the MVCC horizon is, per holder class, in **transactions**. * * `age()` measures transactions elapsed since the pinned xid, not wall time, which is the metric * that matters: dead tuples accumulate per transaction, so a horizon pinned across an idle night * costs nothing while the same age on a busy queue is a backlog of rows autovacuum cannot reclaim. * Measured directly. A REPEATABLE READ holder open on an idle database reports 0, and 2,000 after * 2,000 transactions with the same holder still open. * * Also returns the oldest in-transaction backend's wall-clock age, purely so the warning can say * how long as well as how far; nothing keys off it. */ export function getXminHorizon (lastVacuum: Date, sources: readonly XminHorizonSource[] = XMIN_HORIZON_QUERY_SOURCES): SqlQuery { // There is no statement with no columns, and an empty list here would silently build `SELECT ,`. // Callers that have narrowed to nothing must stop asking, not ask for nothing. assert(sources.length, 'getXminHorizon requires at least one source') const columns = sources.map(name => `(${XMIN_HORIZON_SOURCES[name]}) as "${name}"`) // Identity for the oldest holding backend, so the warning can name what is interfering instead of // only its class. Every column here is readable by an ordinary role for another role's backend // except state, which reads NULL. Self is excluded by pid: the backend asking this question is // always holding a snapshot of its own. // // The query text is deliberately NOT collected, and truncating it would not make it collectable. // pg_stat_activity.query is the raw statement, so an un-parameterized one carries its literals - // emails, tokens, keys - and they sit in the WHERE clause, at the front, where a prefix cut keeps // them. Nor is size the issue: track_activity_query_size already caps it at 1 kB by default. The // objection is where it would end up. Reading that column takes pg_read_all_stats, superuser or // the same role, and it vanishes with the backend; warning.data needs only SELECT on this schema, // survives for warningRetentionDays, is served by getWarnings and the dashboard, and is emitted as // an event most apps forward to their logs. Copying one into the other widens who can read another // application's SQL, makes it durable, and sends it off-box. // // state is kept because it carries the part that changes what the operator does: 'idle in // transaction' is an abandoned transaction (fix the client, or set // idle_in_transaction_session_timeout) while 'active' is a genuinely long-running query. Which // query it is follows from pid, application_name and role, looked up live where the catalog's own // privilege rules still apply. // // selfApplicationName is what makes "ours or theirs" answerable. pg-boss names the pool it owns // 'pgboss', or 'pgboss:' for a registered instance, so a holder matching this connection's own // value is pg-boss doing it to itself - the monitor's own aggregate, most likely - which has a // completely different fix from an external reporting tool holding a transaction open. It is // compared rather than hardcoded because an adapter-supplied pool sets whatever the host app // chose, and claiming that is definitely pg-boss would be a guess. const backendIdentity = sources.includes('backends') ? `, (SELECT to_jsonb(h) FROM ( SELECT pid, application_name as "applicationName", usename::text as "userName", state, age(backend_xmin) as age, extract(epoch from (now() - xact_start))::int as "xactSeconds" FROM pg_catalog.pg_stat_activity WHERE backend_xmin IS NOT NULL AND (xact_start IS NULL OR xact_start < $1) AND pid <> pg_backend_pid() ORDER BY age(backend_xmin) DESC, xact_start NULLS LAST LIMIT 1 ) h) as "backendHolder", (SELECT count(*)::int FROM pg_catalog.pg_stat_activity WHERE backend_xmin IS NOT NULL AND xact_start IS NULL AND pid <> pg_backend_pid()) as "opaqueBackends", current_setting('application_name') as "selfApplicationName"` : '' const text = ` SELECT ${columns.join(',\n ')}, (SELECT max(extract(epoch from (now() - xact_start)))::int FROM pg_catalog.pg_stat_activity WHERE backend_xmin IS NOT NULL AND xact_start IS NOT NULL) as "oldestTransactionSeconds"${backendIdentity} ` // Only backends carries the placeholder, so the value goes along only when it is actually // referenced - binding one to a statement that has none is an error, not a no-op. return { text, values: sources.includes('backends') ? [lastVacuum] : [] } } /** * Dead-tuple accounting for pg-boss's own job tables, plus the autovacuum settings each table is * actually measured against. * * This is the direct measurement the `xmin_horizon` check is built on. `age(backend_xmin)` is only * a proxy for it, and a poor one: it counts cluster transactions, so a producer batching 500 jobs * per `insert()` moves it ~9x slower than one calling `send()` per job for identical garbage, and * any unrelated application sharing the database moves it for reasons that have nothing to do with * pg-boss at all. `n_dead_tup` is the harm itself. * * The budget comes from Postgres rather than from a pg-boss constant: * `autovacuum_vacuum_threshold + autovacuum_vacuum_scale_factor * n_live_tup` is the point the * server itself decides a table needs vacuuming. Per-table reloptions win over the cluster * settings, so an operator who has tuned a queue's table is measured against their own numbers. * * `lastVacuum` is the later of the autovacuum and manual timestamps: what matters is that a vacuum * ran, not who asked for it. GREATEST ignores NULLs, so it is NULL only when neither has ever run. * Returned as epoch milliseconds rather than timestamptz: pg-boss shares the caller's pool, and a * global pg-types parser (Temporal.Instant, say) returns values `new Date()` cannot coerce. * Its age is computed server-side rather than against the client clock, so the check does not * inherit the skew that `clock_skew` exists to report. * * PostgreSQL-only. Callers must gate on noMonitorVacuum - see the note in boss.ts #checkXminHorizon. */ export function getJobTableGarbage (schema: string, tables?: string[]): string { return ` SELECT t.relname as name, s.n_live_tup as "liveTuples", s.n_dead_tup as "deadTuples", (extract(epoch from GREATEST(s.last_autovacuum, s.last_vacuum)) * 1000)::float8 as "lastVacuum", extract(epoch from (now() - GREATEST(s.last_autovacuum, s.last_vacuum)))::int as "vacuumAgeSeconds", coalesce((SELECT option_value::int FROM pg_options_to_table(t.reloptions) WHERE option_name = 'autovacuum_vacuum_threshold'), current_setting('autovacuum_vacuum_threshold')::int) as "threshold", coalesce((SELECT option_value::float8 FROM pg_options_to_table(t.reloptions) WHERE option_name = 'autovacuum_vacuum_scale_factor'), current_setting('autovacuum_vacuum_scale_factor')::float8) as "scaleFactor", coalesce((SELECT option_value::boolean FROM pg_options_to_table(t.reloptions) WHERE option_name = 'autovacuum_enabled'), true) as "autovacuumEnabled" FROM pg_catalog.pg_class t JOIN pg_catalog.pg_namespace n ON n.oid = t.relnamespace JOIN pg_catalog.pg_stat_user_tables s ON s.relid = t.oid WHERE n.nspname = '${resolveSchemaName(schema).replace(SINGLE_QUOTE_REGEX, "''")}' AND t.relkind = 'r' AND ${jobTableScope(schema, tables)} ` } /* eslint-enable no-restricted-syntax */ export function reindexIndex (schema: string, name: string): string { return `REINDEX INDEX CONCURRENTLY ${schema}.${quoteIdentifier(name)}` } export function dropIndexConcurrently (schema: string, name: string): string { return `DROP INDEX CONCURRENTLY IF EXISTS ${schema}.${quoteIdentifier(name)}` } /** * The statement list an operator runs by hand: every stale stub dropped first, then the rebuilds. * * Stubs are emitted whole rather than paired to the index they came from. Postgres does not append * `_ccnew`. It truncates the *base* so the result fits in 63 bytes, so `j_i11` (61 chars) * leaves behind `j__ccnew` with the `i11` gone, and a prefix match never fires on a * partitioned queue (every partition table name is `'j' || sha224(queue_name)`, 57 chars). Dropping * them unconditionally is also what the background pass does: an invalid stub is dead weight that * vacuum still walks, whichever index it came from. */ export function buildReindexCommands (schema: string, targets: { name: string }[], leftovers: { name: string }[]): string[] { if (!targets.length) return [] return [ ...leftovers.map(({ name }) => dropIndexConcurrently(schema, name)), ...targets.map(({ name }) => reindexIndex(schema, name)) ] }