/* eslint-disable no-restricted-syntax -- everything above the v42 entry keeps its original SQL verbatim: the pre-v42 migrations and the versioned DDL snapshot maps they use. The clock guard applies from v42 on; a new snapshot builder for v42 or later must use the schema clock even though it sits inside this region. */ import assert from 'node:assert' import * as plans from './plans.ts' import { resolveSchemaName } from './tools.ts' import * as types from './types.ts' // Options for rendering an async (BAM) migration as inline, self-contained DDL instead // of a job_table_run_async() enqueue call. Used by CLI / exported migrations, which run // in a context with no BAM worker to process the queued commands. See issue #766. interface MigrateOptions { inlineAsync?: boolean // Partitioned queue metadata used to expand inlined index builds, in addition // to job_common. Supplied by callers that hold a live connection (the CLI); empty for a // purely static export, which can only target job_common. Each entry carries the queue's // policy because policy-scoped builds (v38's job_i10) must skip partitions of other policies. partitionTables?: types.MigrationPartition[] } // 12.28.0 replaced the plain `string[]` partitionTables shape with { tableName, policy } records. // TypeScript callers get a compile error; this keeps JavaScript callers from silently exporting a // migration whose target list is wrong (a bare string has no policy, so every policy-scoped build // would quietly drop it). function assertPartitionMetadata (partitionTables: types.MigrationPartition[]) { for (const partition of partitionTables) { assert( typeof partition === 'object' && partition !== null && typeof partition.tableName === 'string' && typeof partition.policy === 'string', `partitionTables entries must be { tableName, policy } records, received: ${JSON.stringify(partition)}. The string[] form was removed in 12.28.0.` ) } } // Mirrors the SQL job_table_format() function (src/plans.ts): rewrites a command targeting // the base `job` table to target a specific partition table. function formatJobTable (command: string, table: string) { // Anchor both rewrites so a schema name that itself contains these substrings (e.g. `job_intake`) // isn't mangled: `.job\b` only matches the base table reference (`schema.job`, not `schema.job_i5` // whose `job` is followed by `_`), and `job_iN` only matches the bare index-name tokens (job_i1..9), // never the `job_i` inside an arbitrary schema name. return command .replace(/\.job\b/g, `.${table}`) .replace(/\bjob_i(\d+)/g, `${table}_i$1`) } // Derives the direct index DDL that a job_table_run_async() command would eventually run // via BAM, one statement per target table, each prefixed with a provenance comment. The // CONCURRENTLY keyword is preserved (these are emitted after COMMIT), and IF NOT EXISTS is // added so the script is safe to re-run. function inlineAsyncCommand (schema: string, asyncMigration: string | types.AsyncMigrationCommand, version: number, partitionTables: types.MigrationPartition[]) { // Two spellings: the legacy string form embeds the whole job_table_run_async() call as SQL, so // name/body/table-pin have to be parsed back out of it; the record form carries them as fields. const isLegacy = typeof asyncMigration === 'string' const asyncCommand = isLegacy ? asyncMigration : asyncMigration.command const nameMatch = isLegacy ? asyncCommand.match(/job_table_run_async\(\s*'([^']+)'/) : null const bodyMatch = isLegacy ? asyncCommand.match(/\$\$([\s\S]*?)\$\$/) : null // An explicit table arg after the $$ body pins the command to a single table (e.g. i8 → // job_common); without it the command fans out across job_common + every partition. const tableMatch = isLegacy ? asyncCommand.match(/\$\$\s*,\s*'([^']+)'/) : null assert(!isLegacy || (nameMatch && bodyMatch), `Unable to inline async migration command: ${asyncCommand}`) const commandName = isLegacy ? nameMatch![1] : asyncMigration.name const body = (isLegacy ? bodyMatch![1] : asyncMigration.command).trim() // A policy-scoped build only belongs on partitions of that policy; unscoped builds fan out to all. const partitionPolicy = isLegacy ? undefined : asyncMigration.partitionPolicy const eligiblePartitions = partitionTables .filter(partition => !partitionPolicy || partition.policy === partitionPolicy) .map(partition => partition.tableName) const targetTables = tableMatch ? [tableMatch[1]] : ['job_common', ...eligiblePartitions] return targetTables.map(table => { // Add IF NOT EXISTS so the exported script is re-runnable. The negative lookahead keeps it // idempotent when the async command already spells out IF NOT EXISTS (the live BAM path needs it // there for its own idempotency, e.g. migration v36's job_i9 build), avoiding a double insert. const ddl = formatJobTable(body, table).replace( /(CREATE (?:UNIQUE )?INDEX CONCURRENTLY)(?! IF NOT EXISTS) /, '$1 IF NOT EXISTS ' ) const comment = `-- inlined from ${schema}.job_table_run_async (migration v${version}, command: ${commandName})` return `${comment}\n${ddl}` }) } function renderAsyncCommand (schema: string, asyncMigration: string | types.AsyncMigrationCommand, version: number) { if (typeof asyncMigration === 'string') { return asyncMigration.replace(/\$VERSION\$/g, String(version)) } const name = asyncMigration.name.replaceAll("'", "''") const policy = asyncMigration.partitionPolicy?.replaceAll("'", "''") if (!policy) { return `SELECT ${schema}.job_table_run_async('${name}', ${version}, $$${asyncMigration.command}$$)` } // queue_name is passed alongside table_name so the enqueued rows record which queue each build // belongs to, matching what the unscoped fan-out inside job_table_run_async() writes. For a // partition row the function re-derives table_name from queue_name and lands on the same value; // job_common has no queue of its own, so its row keeps queue NULL like the fan-out does. return `SELECT ${schema}.job_table_run_async( '${name}', ${version}, $$${asyncMigration.command}$$, targets.table_name, targets.queue_name ) FROM ( SELECT 'job_common'::text AS table_name, NULL::text AS queue_name UNION ALL SELECT table_name, name FROM ${schema}.queue WHERE partition = true AND policy = '${policy}' ) targets` } // A migration split into its transactional block and the post-COMMIT CONCURRENTLY index // builds. The concurrent statements cannot run inside a transaction, so a programmatic // caller (the CLI apply path) must execute each one on its own; a printed script can simply // concatenate them (psql runs top-level statements individually). See migrate(). interface MigrationCommands { sql: string concurrent: string[] } function flatten (schema: string, commands: string[], version: number, noAdvisoryLocks?: boolean) { commands.unshift(plans.assertMigration(schema, version)) commands.push(plans.setVersion(schema, version)) return plans.locked(schema, commands, undefined, noAdvisoryLocks) } function rollback (schema: string, version: number, migrations?: types.Migration[], noAdvisoryLocks?: boolean) { migrations = migrations || getAll(schema) const result = migrations.find(i => i.version === version) assert(result, `Version ${version} not found.`) // Async (BAM) index builds are enqueued as bam rows, and the BAM runner does not filter by schema // version, so a row left unfinished by a rollback would rebuild, one version later, the very index // the rollback just dropped. Clear this version's unfinished rows as part of the same transaction. // Only migrations that enqueue async commands can have any, and the migration that introduced BAM // is also the first one with async commands, so bam is guaranteed to exist here; ordering the delete // ahead of the migration's own steps keeps it before that migration's DROP TABLE bam. // // No guard against a build that is genuinely running: the uninstall steps drop the index the build // holds a ShareUpdateExclusiveLock on, so a live build already makes the whole rollback fail on the // transaction's lock_timeout, taking this delete with it. const clearAsync = result.async?.length ? [`DELETE FROM ${schema}.bam WHERE version = ${result.version} AND status <> 'completed'`] : [] return flatten(schema, [...clearAsync, ...(result.uninstall || [])], result.previous, noAdvisoryLocks) } function next (schema: string, version: number, migrations?: types.Migration[], noAdvisoryLocks?: boolean) { migrations = migrations || getAll(schema) const result = migrations.find(i => i.previous === version) assert(result, `Version ${version} not found.`) return flatten(schema, result.install, result.version, noAdvisoryLocks) } // Builds the migration as separate pieces: the transactional block plus any inlined // CONCURRENTLY index builds (when options.inlineAsync). Callers that execute SQL // programmatically must run `concurrent` statements individually, outside a transaction. function migrateCommands (schema: string, version: number, migrations?: types.Migration[], noAdvisoryLocks?: boolean, options: MigrateOptions = {}): MigrationCommands { migrations = migrations || getAll(schema) // Refuse to migrate from a real DB version older than the oldest migration can start from. // Without this floor, `filter(i => i.previous >= version)` happily selects the whole chain for any // version below the minimum `previous`, applying migrations over missing intermediate steps. A // cryptic mid-transaction failure, or worse a "success" that stamps the latest version onto an // incomplete schema. Version 0 is the sentinel for a full "from scratch" export (getMigrationPlans) // and is intentionally exempt. // Only floor a valid numeric version; a non-numeric/garbage version falls through to the // "Version X not found" assert below. Version 0 is the full-export sentinel and is exempt. if (Number.isInteger(version) && version !== 0) { const minPrevious = Math.min(...migrations.map(i => i.previous)) assert(version >= minPrevious, `Cannot migrate pg-boss schema from version ${version}: the oldest supported starting version is ${minPrevious}. ` + 'Upgrade to a schema at or above that version using an older pg-boss release first.') } assertPartitionMetadata(options.partitionTables || []) const concurrent: string[] = [] const result = migrations .filter(i => i.previous >= version!) .sort((a, b) => a.version - b.version) .reduce((acc, migration) => { acc.install = acc.install.concat(migration.install) if (migration.async) { if (options.inlineAsync) { // Bypass BAM: emit the real index DDL (run after COMMIT) instead of enqueuing it. for (const cmd of migration.async) { concurrent.push(...inlineAsyncCommand(schema, cmd, migration.version, options.partitionTables || [])) } } else { const bamCommands = migration.async.map(cmd => renderAsyncCommand(schema, cmd, migration.version) ) acc.install = acc.install.concat(bamCommands) } } acc.version = migration.version return acc }, { install: [] as string[], version }) assert(result.install.length > 0, `Version ${version} not found.`) return { sql: flatten(schema, result.install, result.version!, noAdvisoryLocks), concurrent } } // Renders a migration as a single SQL script. The inlined CONCURRENTLY builds are appended // after COMMIT; this form is meant for printing/export (psql runs them individually). To // apply programmatically, use migrateCommands() and execute `concurrent` separately. function migrate (schema: string, version: number, migrations?: types.Migration[], noAdvisoryLocks?: boolean, options: MigrateOptions = {}) { const { sql, concurrent } = migrateCommands(schema, version, migrations, noAdvisoryLocks, options) return concurrent.length ? `${sql}\n${concurrent.join(';\n')};` : sql } const createQueueFn: Record string> = { 26: (schema) => ` CREATE OR REPLACE 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 'job_common' 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 ) VALUES ( queue_name, options->>'policy', COALESCE((options->>'retryLimit')::int, 2), COALESCE((options->>'retryDelay')::int, 0), COALESCE((options->>'retryBackoff')::bool, false), (options->>'retryDelayMax')::int, COALESCE((options->>'expireInSeconds')::int, 900), COALESCE((options->>'retentionSeconds')::int, 1209600), COALESCE((options->>'deleteAfterSeconds')::int, 604800), COALESCE((options->>'warningQueueSize')::int, 0), options->>'deadLetter', COALESCE((options->>'partition')::bool, false), tablename ) 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 format('ALTER TABLE ${schema}.%1$I ADD PRIMARY KEY (name, id)', tablename); EXECUTE format('ALTER TABLE ${schema}.%1$I ADD CONSTRAINT q_fkey FOREIGN KEY (name) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED', tablename); EXECUTE format('ALTER TABLE ${schema}.%1$I ADD CONSTRAINT dlq_fkey FOREIGN KEY (dead_letter) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED', tablename); EXECUTE format('CREATE INDEX %1$s_i5 ON ${schema}.%1$I (name, start_after) INCLUDE (priority, created_on, id) WHERE state < ''active''', tablename); EXECUTE format('CREATE UNIQUE INDEX %1$s_i4 ON ${schema}.%1$I (name, singleton_on, COALESCE(singleton_key, '''')) WHERE state <> ''cancelled'' AND singleton_on IS NOT NULL', tablename); IF options->>'policy' = 'short' THEN EXECUTE format('CREATE UNIQUE INDEX %1$s_i1 ON ${schema}.%1$I (name, COALESCE(singleton_key, '''')) WHERE state = ''created'' AND policy = ''short''', tablename); ELSIF options->>'policy' = 'singleton' THEN EXECUTE format('CREATE UNIQUE INDEX %1$s_i2 ON ${schema}.%1$I (name, COALESCE(singleton_key, '''')) WHERE state = ''active'' AND policy = ''singleton''', tablename); ELSIF options->>'policy' = 'stately' THEN EXECUTE format('CREATE UNIQUE INDEX %1$s_i3 ON ${schema}.%1$I (name, state, COALESCE(singleton_key, '''')) WHERE state <= ''active'' AND policy = ''stately''', tablename); ELSIF options->>'policy' = 'exclusive' THEN EXECUTE format('CREATE UNIQUE INDEX %1$s_i6 ON ${schema}.%1$I (name, COALESCE(singleton_key, '''')) WHERE state <= ''active'' AND policy = ''exclusive''', 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; `, 27: (schema) => ` CREATE OR REPLACE 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 'job_common' 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 ) VALUES ( queue_name, options->>'policy', COALESCE((options->>'retryLimit')::int, 2), COALESCE((options->>'retryDelay')::int, 0), COALESCE((options->>'retryBackoff')::bool, false), (options->>'retryDelayMax')::int, COALESCE((options->>'expireInSeconds')::int, 900), COALESCE((options->>'retentionSeconds')::int, 1209600), COALESCE((options->>'deleteAfterSeconds')::int, 604800), COALESCE((options->>'warningQueueSize')::int, 0), options->>'deadLetter', COALESCE((options->>'partition')::bool, false), tablename ) 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$ALTER TABLE ${schema}.job ADD PRIMARY KEY (name, id)$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT q_fkey FOREIGN KEY (name) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT dlq_fkey FOREIGN KEY (dead_letter) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i5 ON ${schema}.job (name, start_after) INCLUDE (priority, created_on, id) WHERE state < 'active'$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i4 ON ${schema}.job (name, singleton_on, COALESCE(singleton_key, '')) WHERE state <> 'cancelled' AND singleton_on IS NOT NULL$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i7 ON ${schema}.job (name, group_id) WHERE state = 'active' AND group_id IS NOT NULL$cmd$, tablename); IF options->>'policy' = 'short' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i1 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = 'created' AND policy = 'short'$cmd$, tablename); ELSIF options->>'policy' = 'singleton' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i2 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = 'active' AND policy = 'singleton'$cmd$, tablename); ELSIF options->>'policy' = 'stately' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i3 ON ${schema}.job (name, state, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'stately'$cmd$, tablename); ELSIF options->>'policy' = 'exclusive' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i6 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'exclusive'$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; `, 28: (schema) => ` CREATE OR REPLACE 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 'job_common' 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 ) VALUES ( queue_name, options->>'policy', COALESCE((options->>'retryLimit')::int, 2), COALESCE((options->>'retryDelay')::int, 0), COALESCE((options->>'retryBackoff')::bool, false), (options->>'retryDelayMax')::int, COALESCE((options->>'expireInSeconds')::int, 900), COALESCE((options->>'retentionSeconds')::int, 1209600), COALESCE((options->>'deleteAfterSeconds')::int, 604800), COALESCE((options->>'warningQueueSize')::int, 0), options->>'deadLetter', COALESCE((options->>'partition')::bool, false), tablename ) 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$ALTER TABLE ${schema}.job ADD PRIMARY KEY (name, id)$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT q_fkey FOREIGN KEY (name) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT dlq_fkey FOREIGN KEY (dead_letter) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i5 ON ${schema}.job (name, start_after) INCLUDE (priority, created_on, id) WHERE state < 'active'$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i4 ON ${schema}.job (name, singleton_on, COALESCE(singleton_key, '')) WHERE state <> 'cancelled' AND singleton_on IS NOT NULL$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i7 ON ${schema}.job (name, group_id) WHERE state = 'active' AND group_id IS NOT NULL$cmd$, tablename); IF options->>'policy' = 'short' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i1 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = 'created' AND policy = 'short'$cmd$, tablename); ELSIF options->>'policy' = 'singleton' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i2 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = 'active' AND policy = 'singleton'$cmd$, tablename); ELSIF options->>'policy' = 'stately' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i3 ON ${schema}.job (name, state, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'stately'$cmd$, tablename); ELSIF options->>'policy' = 'exclusive' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i6 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'exclusive'$cmd$, tablename); ELSIF options->>'policy' = 'key_strict_fifo' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i8 ON ${schema}.job (name, singleton_key) WHERE state IN ('active', 'retry', 'failed') AND policy = 'key_strict_fifo'$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT job_key_strict_fifo_singleton_key_check CHECK (NOT (policy = 'key_strict_fifo' AND singleton_key IS NULL))$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; `, 30: (schema) => ` CREATE OR REPLACE 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 'job_common' 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 ) VALUES ( queue_name, options->>'policy', COALESCE((options->>'retryLimit')::int, 2), COALESCE((options->>'retryDelay')::int, 0), COALESCE((options->>'retryBackoff')::bool, false), (options->>'retryDelayMax')::int, COALESCE((options->>'expireInSeconds')::int, 900), COALESCE((options->>'retentionSeconds')::int, 1209600), COALESCE((options->>'deleteAfterSeconds')::int, 604800), COALESCE((options->>'warningQueueSize')::int, 0), options->>'deadLetter', COALESCE((options->>'partition')::bool, false), tablename, (options->>'heartbeatSeconds')::int ) 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$ALTER TABLE ${schema}.job ADD PRIMARY KEY (name, id)$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT q_fkey FOREIGN KEY (name) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT dlq_fkey FOREIGN KEY (dead_letter) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i5 ON ${schema}.job (name, start_after) INCLUDE (priority, created_on, id) WHERE state < 'active'$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i4 ON ${schema}.job (name, singleton_on, COALESCE(singleton_key, '')) WHERE state <> 'cancelled' AND singleton_on IS NOT NULL$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i7 ON ${schema}.job (name, group_id) WHERE state = 'active' AND group_id IS NOT NULL$cmd$, tablename); IF options->>'policy' = 'short' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i1 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = 'created' AND policy = 'short'$cmd$, tablename); ELSIF options->>'policy' = 'singleton' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i2 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = 'active' AND policy = 'singleton'$cmd$, tablename); ELSIF options->>'policy' = 'stately' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i3 ON ${schema}.job (name, state, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'stately'$cmd$, tablename); ELSIF options->>'policy' = 'exclusive' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i6 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'exclusive'$cmd$, tablename); ELSIF options->>'policy' = 'key_strict_fifo' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i8 ON ${schema}.job (name, singleton_key) WHERE state IN ('active', 'retry', 'failed') AND policy = 'key_strict_fifo'$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT job_key_strict_fifo_singleton_key_check CHECK (NOT (policy = 'key_strict_fifo' AND singleton_key IS NULL))$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; `, 31: (schema) => ` CREATE OR REPLACE 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 'job_common' 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 ) VALUES ( queue_name, options->>'policy', COALESCE((options->>'retryLimit')::int, 2), COALESCE((options->>'retryDelay')::int, 0), COALESCE((options->>'retryBackoff')::bool, false), (options->>'retryDelayMax')::int, COALESCE((options->>'expireInSeconds')::int, 900), COALESCE((options->>'retentionSeconds')::int, 1209600), COALESCE((options->>'deleteAfterSeconds')::int, 604800), COALESCE((options->>'warningQueueSize')::int, 0), options->>'deadLetter', COALESCE((options->>'partition')::bool, false), tablename, (options->>'heartbeatSeconds')::int ) 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$ALTER TABLE ${schema}.job ADD PRIMARY KEY (name, id)$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT q_fkey FOREIGN KEY (name) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT dlq_fkey FOREIGN KEY (dead_letter) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i5 ON ${schema}.job (name, start_after) INCLUDE (priority, created_on, id) WHERE state < 'active' AND NOT blocked$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i4 ON ${schema}.job (name, singleton_on, COALESCE(singleton_key, '')) WHERE state <> 'cancelled' AND singleton_on IS NOT NULL$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i7 ON ${schema}.job (name, group_id) WHERE state = 'active' AND group_id IS NOT NULL$cmd$, tablename); IF options->>'policy' = 'short' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i1 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = 'created' AND policy = 'short'$cmd$, tablename); ELSIF options->>'policy' = 'singleton' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i2 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = 'active' AND policy = 'singleton'$cmd$, tablename); ELSIF options->>'policy' = 'stately' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i3 ON ${schema}.job (name, state, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'stately'$cmd$, tablename); ELSIF options->>'policy' = 'exclusive' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i6 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'exclusive'$cmd$, tablename); ELSIF options->>'policy' = 'key_strict_fifo' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i8 ON ${schema}.job (name, singleton_key) WHERE state IN ('active', 'retry', 'failed') AND policy = 'key_strict_fifo'$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT job_key_strict_fifo_singleton_key_check CHECK (NOT (policy = 'key_strict_fifo' AND singleton_key IS NULL))$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; `, 32: (schema) => ` CREATE OR REPLACE 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 'job_common' 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 ) VALUES ( queue_name, options->>'policy', COALESCE((options->>'retryLimit')::int, 2), COALESCE((options->>'retryDelay')::int, 0), COALESCE((options->>'retryBackoff')::bool, false), (options->>'retryDelayMax')::int, COALESCE((options->>'expireInSeconds')::int, 900), COALESCE((options->>'retentionSeconds')::int, 1209600), COALESCE((options->>'deleteAfterSeconds')::int, 604800), COALESCE((options->>'warningQueueSize')::int, 0), options->>'deadLetter', COALESCE((options->>'partition')::bool, false), tablename, (options->>'heartbeatSeconds')::int, COALESCE((options->>'notify')::bool, false) ) 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$ALTER TABLE ${schema}.job ADD PRIMARY KEY (name, id)$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT q_fkey FOREIGN KEY (name) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT dlq_fkey FOREIGN KEY (dead_letter) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i5 ON ${schema}.job (name, start_after) INCLUDE (priority, created_on, id) WHERE state < 'active' AND NOT blocked$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i4 ON ${schema}.job (name, singleton_on, COALESCE(singleton_key, '')) WHERE state <> 'cancelled' AND singleton_on IS NOT NULL$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i7 ON ${schema}.job (name, group_id) WHERE state = 'active' AND group_id IS NOT NULL$cmd$, tablename); IF options->>'policy' = 'short' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i1 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = 'created' AND policy = 'short'$cmd$, tablename); ELSIF options->>'policy' = 'singleton' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i2 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = 'active' AND policy = 'singleton'$cmd$, tablename); ELSIF options->>'policy' = 'stately' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i3 ON ${schema}.job (name, state, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'stately'$cmd$, tablename); ELSIF options->>'policy' = 'exclusive' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i6 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'exclusive'$cmd$, tablename); ELSIF options->>'policy' = 'key_strict_fifo' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i8 ON ${schema}.job (name, singleton_key) WHERE state IN ('active', 'retry', 'failed') AND policy = 'key_strict_fifo'$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT job_key_strict_fifo_singleton_key_check CHECK (NOT (policy = 'key_strict_fifo' AND singleton_key IS NULL))$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; `, 33: (schema) => ` CREATE OR REPLACE 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 'job_common' 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 ) VALUES ( queue_name, options->>'policy', COALESCE((options->>'retryLimit')::int, 2), COALESCE((options->>'retryDelay')::int, 0), COALESCE((options->>'retryBackoff')::bool, false), (options->>'retryDelayMax')::int, COALESCE((options->>'expireInSeconds')::int, 900), COALESCE((options->>'retentionSeconds')::int, 1209600), COALESCE((options->>'deleteAfterSeconds')::int, 604800), COALESCE((options->>'warningQueueSize')::int, 0), options->>'deadLetter', COALESCE((options->>'partition')::bool, false), tablename, (options->>'heartbeatSeconds')::int, COALESCE((options->>'notify')::bool, false) ) 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$ALTER TABLE ${schema}.job ADD PRIMARY KEY (name, id)$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT q_fkey FOREIGN KEY (name) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT dlq_fkey FOREIGN KEY (dead_letter) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i5 ON ${schema}.job (name, start_after) WHERE state < 'active' AND NOT blocked$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i4 ON ${schema}.job (name, singleton_on, COALESCE(singleton_key, '')) WHERE state <> 'cancelled' AND singleton_on IS NOT NULL$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i7 ON ${schema}.job (name, group_id) WHERE state = 'active' AND group_id IS NOT NULL$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i9 ON ${schema}.job (name, id) WHERE blocking AND state = 'completed'$cmd$, tablename); IF options->>'policy' = 'short' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i1 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = 'created' AND policy = 'short'$cmd$, tablename); ELSIF options->>'policy' = 'singleton' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i2 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = 'active' AND policy = 'singleton'$cmd$, tablename); ELSIF options->>'policy' = 'stately' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i3 ON ${schema}.job (name, state, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'stately'$cmd$, tablename); ELSIF options->>'policy' = 'exclusive' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i6 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'exclusive'$cmd$, tablename); ELSIF options->>'policy' = 'key_strict_fifo' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i8 ON ${schema}.job (name, singleton_key) WHERE state IN ('active', 'retry', 'failed') AND policy = 'key_strict_fifo'$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT job_key_strict_fifo_singleton_key_check CHECK (NOT (policy = 'key_strict_fifo' AND singleton_key IS NULL))$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; `, 38: (schema) => ` CREATE OR REPLACE 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 'job_common' 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 ) VALUES ( queue_name, options->>'policy', COALESCE((options->>'retryLimit')::int, 2), COALESCE((options->>'retryDelay')::int, 0), COALESCE((options->>'retryBackoff')::bool, false), (options->>'retryDelayMax')::int, COALESCE((options->>'expireInSeconds')::int, 900), COALESCE((options->>'retentionSeconds')::int, 1209600), COALESCE((options->>'deleteAfterSeconds')::int, 604800), COALESCE((options->>'warningQueueSize')::int, 0), options->>'deadLetter', COALESCE((options->>'partition')::bool, false), tablename, (options->>'heartbeatSeconds')::int, COALESCE((options->>'notify')::bool, false) ) 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$ALTER TABLE ${schema}.job ADD PRIMARY KEY (name, id)$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT q_fkey FOREIGN KEY (name) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT dlq_fkey FOREIGN KEY (dead_letter) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i5 ON ${schema}.job (name, start_after) WHERE state < 'active' AND NOT blocked$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i4 ON ${schema}.job (name, singleton_on, COALESCE(singleton_key, '')) WHERE state <> 'cancelled' AND singleton_on IS NOT NULL$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i7 ON ${schema}.job (name, group_id) WHERE state = 'active' AND group_id IS NOT NULL$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i9 ON ${schema}.job (name, id) WHERE blocking AND state = 'completed'$cmd$, tablename); IF options->>'policy' = 'short' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i1 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = 'created' AND policy = 'short'$cmd$, tablename); ELSIF options->>'policy' = 'singleton' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i2 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = 'active' AND policy = 'singleton'$cmd$, tablename); ELSIF options->>'policy' = 'stately' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i3 ON ${schema}.job (name, state, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'stately'$cmd$, tablename); ELSIF options->>'policy' = 'exclusive' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i6 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'exclusive'$cmd$, tablename); ELSIF options->>'policy' = 'key_strict_fifo' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i8 ON ${schema}.job (name, singleton_key) WHERE state IN ('active', 'retry', 'failed') AND policy = 'key_strict_fifo'$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i10 ON ${schema}.job (name, singleton_key, state DESC, created_on, id) INCLUDE (start_after) WHERE state < 'active' AND NOT blocked AND policy = 'key_strict_fifo'$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT job_key_strict_fifo_singleton_key_check CHECK (NOT (policy = 'key_strict_fifo' AND singleton_key IS NULL))$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; `, 40: (schema) => ` CREATE OR REPLACE 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 'job_common' 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 ) VALUES ( queue_name, options->>'policy', COALESCE((options->>'retryLimit')::int, 2), COALESCE((options->>'retryDelay')::int, 0), COALESCE((options->>'retryBackoff')::bool, false), (options->>'retryDelayMax')::int, COALESCE((options->>'expireInSeconds')::int, 900), COALESCE((options->>'retentionSeconds')::int, 1209600), COALESCE((options->>'deleteAfterSeconds')::int, 604800), COALESCE((options->>'warningQueueSize')::int, 0), options->>'deadLetter', COALESCE((options->>'partition')::bool, false), tablename, (options->>'heartbeatSeconds')::int, COALESCE((options->>'notify')::bool, false) ) 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$ALTER TABLE ${schema}.job ADD PRIMARY KEY (name, id)$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT q_fkey FOREIGN KEY (name) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT dlq_fkey FOREIGN KEY (dead_letter) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i11 ON ${schema}.job (name, priority DESC, created_on, start_after) WHERE state < 'active' AND NOT blocked$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i4 ON ${schema}.job (name, singleton_on, COALESCE(singleton_key, '')) WHERE state <> 'cancelled' AND singleton_on IS NOT NULL$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i7 ON ${schema}.job (name, group_id) WHERE state = 'active' AND group_id IS NOT NULL$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i9 ON ${schema}.job (name, id) WHERE blocking AND state = 'completed'$cmd$, tablename); IF options->>'policy' = 'short' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i1 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = 'created' AND policy = 'short'$cmd$, tablename); ELSIF options->>'policy' = 'singleton' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i2 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = 'active' AND policy = 'singleton'$cmd$, tablename); ELSIF options->>'policy' = 'stately' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i3 ON ${schema}.job (name, state, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'stately'$cmd$, tablename); ELSIF options->>'policy' = 'exclusive' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i6 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'exclusive'$cmd$, tablename); ELSIF options->>'policy' = 'key_strict_fifo' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i8 ON ${schema}.job (name, singleton_key) WHERE state IN ('active', 'retry', 'failed') AND policy = 'key_strict_fifo'$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i10 ON ${schema}.job (name, singleton_key, state DESC, created_on, id) INCLUDE (start_after) WHERE state < 'active' AND NOT blocked AND policy = 'key_strict_fifo'$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT job_key_strict_fifo_singleton_key_check CHECK (NOT (policy = 'key_strict_fifo' AND singleton_key IS NULL))$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; `, 42: (schema) => ` CREATE OR REPLACE 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 'job_common' 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, 2), COALESCE((options->>'retryDelay')::int, 0), COALESCE((options->>'retryBackoff')::bool, false), (options->>'retryDelayMax')::int, COALESCE((options->>'expireInSeconds')::int, 900), COALESCE((options->>'retentionSeconds')::int, 1209600), COALESCE((options->>'deleteAfterSeconds')::int, 604800), COALESCE((options->>'warningQueueSize')::int, 0), options->>'deadLetter', COALESCE((options->>'partition')::bool, false), 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$ALTER TABLE ${schema}.job ADD PRIMARY KEY (name, id)$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT q_fkey FOREIGN KEY (name) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT dlq_fkey FOREIGN KEY (dead_letter) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i11 ON ${schema}.job (name, priority DESC, created_on, start_after) WHERE state < 'active' AND NOT blocked$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i4 ON ${schema}.job (name, singleton_on, COALESCE(singleton_key, '')) WHERE state <> 'cancelled' AND singleton_on IS NOT NULL$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i7 ON ${schema}.job (name, group_id) WHERE state = 'active' AND group_id IS NOT NULL$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i9 ON ${schema}.job (name, id) WHERE blocking AND state = 'completed'$cmd$, tablename); IF options->>'policy' = 'short' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i1 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = 'created' AND policy = 'short'$cmd$, tablename); ELSIF options->>'policy' = 'singleton' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i2 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = 'active' AND policy = 'singleton'$cmd$, tablename); ELSIF options->>'policy' = 'stately' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i3 ON ${schema}.job (name, state, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'stately'$cmd$, tablename); ELSIF options->>'policy' = 'exclusive' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i6 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'exclusive'$cmd$, tablename); ELSIF options->>'policy' = 'key_strict_fifo' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i8 ON ${schema}.job (name, singleton_key) WHERE state IN ('active', 'retry', 'failed') AND policy = 'key_strict_fifo'$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i10 ON ${schema}.job (name, singleton_key, state DESC, created_on, id) INCLUDE (start_after) WHERE state < 'active' AND NOT blocked AND policy = 'key_strict_fifo'$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT job_key_strict_fifo_singleton_key_check CHECK (NOT (policy = 'key_strict_fifo' AND singleton_key IS NULL))$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; `, 43: (schema) => ` CREATE OR REPLACE 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 'job_common' 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, 2), COALESCE((options->>'retryDelay')::int, 0), COALESCE((options->>'retryBackoff')::bool, false), (options->>'retryDelayMax')::int, COALESCE((options->>'expireInSeconds')::int, 900), COALESCE((options->>'retentionSeconds')::int, 1209600), COALESCE((options->>'deleteAfterSeconds')::int, 604800), COALESCE((options->>'warningQueueSize')::int, 0), options->>'deadLetter', COALESCE((options->>'partition')::bool, false), 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$ALTER TABLE ${schema}.job ADD PRIMARY KEY (name, id)$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT q_fkey FOREIGN KEY (name) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT dlq_fkey FOREIGN KEY (dead_letter) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i11 ON ${schema}.job (name, priority DESC, created_on, start_after) WHERE state < 'active' AND NOT blocked$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i4 ON ${schema}.job (name, singleton_on, COALESCE(singleton_key, '')) WHERE state <> 'cancelled' AND singleton_on IS NOT NULL$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i7 ON ${schema}.job (name, group_id) WHERE state = 'active' AND group_id IS NOT NULL$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i9 ON ${schema}.job (name, id) WHERE blocking AND state = 'completed'$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i12 ON ${schema}.job (source_root_id) WHERE source_root_id IS NOT NULL$cmd$, tablename); IF options->>'policy' = 'short' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i1 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = 'created' AND policy = 'short'$cmd$, tablename); ELSIF options->>'policy' = 'singleton' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i2 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = 'active' AND policy = 'singleton'$cmd$, tablename); ELSIF options->>'policy' = 'stately' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i3 ON ${schema}.job (name, state, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'stately'$cmd$, tablename); ELSIF options->>'policy' = 'exclusive' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i6 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'exclusive'$cmd$, tablename); ELSIF options->>'policy' = 'key_strict_fifo' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i8 ON ${schema}.job (name, singleton_key) WHERE state IN ('active', 'retry', 'failed') AND policy = 'key_strict_fifo'$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i10 ON ${schema}.job (name, singleton_key, state DESC, created_on, id) INCLUDE (start_after) WHERE state < 'active' AND NOT blocked AND policy = 'key_strict_fifo'$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT job_key_strict_fifo_singleton_key_check CHECK (NOT (policy = 'key_strict_fifo' AND singleton_key IS NULL))$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; `, 45: (schema) => ` CREATE OR REPLACE 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 'job_common' 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, 2), COALESCE((options->>'retryDelay')::int, 0), COALESCE((options->>'retryBackoff')::bool, false), (options->>'retryDelayMax')::int, COALESCE((options->>'expireInSeconds')::int, 900), COALESCE((options->>'retentionSeconds')::int, 1209600), COALESCE((options->>'deleteAfterSeconds')::int, 604800), COALESCE((options->>'warningQueueSize')::int, 0), options->>'deadLetter', COALESCE((options->>'partition')::bool, false), 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$ALTER TABLE ${schema}.job ADD PRIMARY KEY (name, id)$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT q_fkey FOREIGN KEY (name) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT dlq_fkey FOREIGN KEY (dead_letter) REFERENCES ${schema}.queue (name) ON DELETE RESTRICT DEFERRABLE INITIALLY DEFERRED$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i11 ON ${schema}.job (name, priority DESC, created_on, start_after) WHERE state < 'active' AND NOT blocked$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i4 ON ${schema}.job (name, singleton_on, COALESCE(singleton_key, '')) WHERE state <> 'cancelled' AND singleton_on IS NOT NULL$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i7 ON ${schema}.job (name, group_id) WHERE state = 'active' AND group_id IS NOT NULL$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i9 ON ${schema}.job (name, id) WHERE blocking AND state = 'completed'$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i12 ON ${schema}.job (source_root_id) WHERE source_root_id IS NOT NULL$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i13 ON ${schema}.job (name, singleton_key) WHERE state = 'created' AND upsert_by_key$cmd$, tablename); IF options->>'policy' = 'short' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i1 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = 'created' AND policy = 'short'$cmd$, tablename); ELSIF options->>'policy' = 'singleton' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i2 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state = 'active' AND policy = 'singleton'$cmd$, tablename); ELSIF options->>'policy' = 'stately' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i3 ON ${schema}.job (name, state, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'stately'$cmd$, tablename); ELSIF options->>'policy' = 'exclusive' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i6 ON ${schema}.job (name, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'exclusive'$cmd$, tablename); ELSIF options->>'policy' = 'key_strict_fifo' THEN EXECUTE ${schema}.job_table_format($cmd$CREATE UNIQUE INDEX job_i8 ON ${schema}.job (name, singleton_key) WHERE state IN ('active', 'retry', 'failed') AND policy = 'key_strict_fifo'$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$CREATE INDEX job_i10 ON ${schema}.job (name, singleton_key, state DESC, created_on, id) INCLUDE (start_after) WHERE state < 'active' AND NOT blocked AND policy = 'key_strict_fifo'$cmd$, tablename); EXECUTE ${schema}.job_table_format($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT job_key_strict_fifo_singleton_key_check CHECK (NOT (policy = 'key_strict_fifo' AND singleton_key IS NULL))$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; ` } // The noTablePartitioning create_queue, frozen per version like createQueueFn above. It needed no // snapshot until v42: every earlier migration that replaced create_queue did so to change index DDL, // and only the partitioned body carries any. const createQueueNoPartitionFn: Record string> = { 41: (schema) => ` CREATE OR REPLACE 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 ) VALUES ( queue_name, options->>'policy', COALESCE((options->>'retryLimit')::int, 2), COALESCE((options->>'retryDelay')::int, 0), COALESCE((options->>'retryBackoff')::bool, false), (options->>'retryDelayMax')::int, COALESCE((options->>'expireInSeconds')::int, 900), COALESCE((options->>'retentionSeconds')::int, 1209600), COALESCE((options->>'deleteAfterSeconds')::int, 604800), COALESCE((options->>'warningQueueSize')::int, 0), options->>'deadLetter', false, 'job', (options->>'heartbeatSeconds')::int ) ON CONFLICT DO NOTHING; END; $$ LANGUAGE plpgsql; `, 42: (schema) => ` CREATE OR REPLACE 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, 2), COALESCE((options->>'retryDelay')::int, 0), COALESCE((options->>'retryBackoff')::bool, false), (options->>'retryDelayMax')::int, COALESCE((options->>'expireInSeconds')::int, 900), COALESCE((options->>'retentionSeconds')::int, 1209600), COALESCE((options->>'deleteAfterSeconds')::int, 604800), COALESCE((options->>'warningQueueSize')::int, 0), options->>'deadLetter', false, 'job', (options->>'heartbeatSeconds')::int, ${schema}.job_now(), ${schema}.job_now() ) ON CONFLICT DO NOTHING; END; $$ LANGUAGE plpgsql; ` } // Frozen per-version snapshots of the queue_stats DDL, version-keyed like createQueueFn above. A // migration must always emit the DDL as it was authored for that schema version; the plans.* builders // track the *current* schema and will drift as it evolves, so the migration copies the DDL here // rather than importing it. When a later version changes queue_stats, add a new keyed entry and leave // the older ones untouched. const createTableQueueStatsFn: Record string> = { 35: (schema, noPartitioning) => noPartitioning ? ` 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, captured_on timestamptz NOT NULL DEFAULT now(), PRIMARY KEY (id) ) ` : ` 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, captured_on timestamptz NOT NULL DEFAULT now(), PRIMARY KEY (id, captured_on) ) PARTITION BY RANGE (captured_on) ` } const createIndexQueueStatsFn: Record string> = { 35: (schema, noCovering) => { const cols = '(name, captured_on DESC)' const include = 'INCLUDE (deferred_count, queued_count, ready_count, active_count, failed_count, total_count)' return noCovering ? `CREATE INDEX queue_stats_i1 ON ${schema}.queue_stats ${cols}` : `CREATE INDEX queue_stats_i1 ON ${schema}.queue_stats ${cols} ${include}` } } // The nspname comparison needs the resolved catalog name, same as plans.ensureQueueStatsPartitions. // This does change the SQL a frozen migration emits, but only for names the check could never have // matched anyway (quoted, or bare mixed-case), and only from "silently creates a partition it // already believes is missing" to "checks correctly". const ensureQueueStatsPartitionsFn: Record string> = { 35: (schema) => ` DO $$ DECLARE d date; i int; part_name text; BEGIN FOR i IN 0..1 LOOP d := (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; $$ ` } // Frozen job_table_format() bodies, one per schema version, so the v37 migration is an immutable // snapshot even if plans.jobTableFormatFunction drifts later. 37 is the anchored regexp_replace // fix (matches only the base table reference and bare job_iN tokens); 36 is the prior naive // replace(), kept solely so v37's rollback restores the exact previous definition. const jobTableFormatFn: Record string> = { 36: (schema) => ` CREATE OR REPLACE FUNCTION ${schema}.job_table_format(command text, table_name text) RETURNS text AS $$ SELECT format( replace( replace(command, '.job', '.%1$I'), 'job_i', '%1$s_i' ), table_name ); $$ LANGUAGE sql IMMUTABLE; `, 37: (schema) => ` CREATE OR REPLACE 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; ` } // Lowest schema version a migration can start from (the smallest `previous` in the set). Below // this there is no chain to apply; callers use it as an honest floor / offline fallback. function getMinVersion (schema: string): number { return Math.min(...getAll(schema).map(i => i.previous)) } // The migration set shaped for a resolved config, so the flag order lives in one place rather than // at every positional call site (contractor, the CLI, tests that inject a list). function getAllForConfig (config: Pick): types.Migration[] { return getAll(config.schema, config.noTablePartitioning, config.noCoveringIndexes, config.noAddColumnBackfill) } // Authoring note on seeding a column you just added: that write is rejected by a backend where // ADD COLUMN is a schema-change job (CockroachDB, see noAddColumnBackfill), so gate it on that flag // and give the runtime a way to read the same answer without it - v40 falls back to monitor_on, v41 // relabels on the first cron pass. When the seeded value cannot be recomputed at read time, gating // it out loses data: enqueue it as a bam row instead, so it runs after the migration commits. function getAll (schema: string, noPartitioning = false, noCovering = false, noAddColumnBackfill = false): types.Migration[] { return [ { release: '11.1.0', version: 26, previous: 25, install: [ createQueueFn[26](schema), `CREATE UNIQUE INDEX job_i6 ON ${schema}.job_common (name, COALESCE(singleton_key, '')) WHERE state <= 'active' AND policy = 'exclusive'` ], uninstall: [ `DROP INDEX ${schema}.job_i6` ] }, { release: '12.6.0', version: 27, previous: 26, install: [ `ALTER TABLE ${schema}.version ADD COLUMN IF NOT EXISTS bam_on timestamp with time zone`, ` CREATE TABLE IF NOT EXISTS ${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 now(), started_on timestamp with time zone, completed_on timestamp with time zone ) `, `CREATE FUNCTION ${schema}.job_table_format(command text, table_name text) RETURNS text AS $$ SELECT format( replace( replace(command, '.job', '.%1$I'), 'job_i', '%1$s_i' ), table_name ); $$ LANGUAGE sql IMMUTABLE; `, ` CREATE OR REPLACE 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, 'job_common', ${schema}.job_table_format(command, 'job_common') 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; `, ` CREATE OR REPLACE 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, 'job_common'); 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; `, `ALTER TABLE ${schema}.job ADD COLUMN IF NOT EXISTS group_id text`, `ALTER TABLE ${schema}.job ADD COLUMN IF NOT EXISTS group_tier text`, createQueueFn[27](schema), `ALTER INDEX IF EXISTS ${schema}.job_i1 RENAME TO job_common_i1`, `ALTER INDEX IF EXISTS ${schema}.job_i2 RENAME TO job_common_i2`, `ALTER INDEX IF EXISTS ${schema}.job_i3 RENAME TO job_common_i3`, `ALTER INDEX IF EXISTS ${schema}.job_i4 RENAME TO job_common_i4`, `ALTER INDEX IF EXISTS ${schema}.job_i5 RENAME TO job_common_i5`, `ALTER INDEX IF EXISTS ${schema}.job_i6 RENAME TO job_common_i6`, `ALTER INDEX IF EXISTS ${schema}.job_i7 RENAME TO job_common_i7` ], async: [ `SELECT ${schema}.job_table_run_async( 'group_concurency_index', $VERSION$, $$ CREATE INDEX CONCURRENTLY IF NOT EXISTS job_i7 ON ${schema}.job (name, group_id) WHERE state = 'active' AND group_id IS NOT NULL $$ )` ], uninstall: [ `ALTER INDEX ${schema}.job_common_i6 RENAME TO job_i6`, `ALTER INDEX ${schema}.job_common_i5 RENAME TO job_i5`, `ALTER INDEX ${schema}.job_common_i4 RENAME TO job_i4`, `ALTER INDEX ${schema}.job_common_i3 RENAME TO job_i3`, `ALTER INDEX ${schema}.job_common_i2 RENAME TO job_i2`, `ALTER INDEX ${schema}.job_common_i1 RENAME TO job_i1`, `SELECT ${schema}.job_table_run('DROP INDEX IF EXISTS ${schema}.job_i7')`, createQueueFn[26](schema), `DROP FUNCTION ${schema}.job_table_run(text, text, text)`, `DROP FUNCTION ${schema}.job_table_run_async(text, int, text, text, text)`, `DROP FUNCTION ${schema}.job_table_format(text, text)`, `DROP TABLE ${schema}.bam`, `ALTER TABLE ${schema}.version DROP COLUMN bam_on`, `ALTER TABLE ${schema}.job DROP COLUMN group_tier`, `ALTER TABLE ${schema}.job DROP COLUMN group_id` ] }, { release: '12.10.0', version: 28, previous: 27, install: [ `SELECT ${schema}.job_table_run($cmd$ALTER TABLE ${schema}.job ADD CONSTRAINT job_key_strict_fifo_singleton_key_check CHECK (NOT (policy = 'key_strict_fifo' AND singleton_key IS NULL))$cmd$, 'job_common')`, createQueueFn[28](schema) ], async: [ `SELECT ${schema}.job_table_run_async( 'key_strict_fifo_index', $VERSION$, $$ CREATE UNIQUE INDEX CONCURRENTLY IF NOT EXISTS job_i8 ON ${schema}.job (name, singleton_key) WHERE state IN ('active', 'retry', 'failed') AND policy = 'key_strict_fifo' $$ , 'job_common')` ], uninstall: [ `SELECT ${schema}.job_table_run('DROP INDEX IF EXISTS ${schema}.job_i8')`, `SELECT ${schema}.job_table_run('ALTER TABLE ${schema}.job DROP CONSTRAINT IF EXISTS job_key_strict_fifo_singleton_key_check')`, createQueueFn[27](schema) ] }, { release: '12.11.0', version: 29, previous: 28, install: [ `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() )`, `CREATE INDEX warning_i1 ON ${schema}.warning (created_on DESC)` ], uninstall: [ `DROP INDEX ${schema}.warning_i1`, `DROP TABLE ${schema}.warning` ] }, { release: '12.12.0', version: 30, previous: 29, install: [ `ALTER TABLE ${schema}.job ADD COLUMN heartbeat_on timestamp with time zone`, `ALTER TABLE ${schema}.job ADD COLUMN heartbeat_seconds int`, `ALTER TABLE ${schema}.queue ADD COLUMN heartbeat_seconds int`, createQueueFn[30](schema) ], uninstall: [ createQueueFn[28](schema), `ALTER TABLE ${schema}.queue DROP COLUMN heartbeat_seconds`, `ALTER TABLE ${schema}.job DROP COLUMN heartbeat_seconds`, `ALTER TABLE ${schema}.job DROP COLUMN heartbeat_on` ] }, { release: '12.19.0', version: 31, previous: 30, install: [ `ALTER TABLE ${schema}.job ADD COLUMN blocked boolean NOT NULL DEFAULT false`, `ALTER TABLE ${schema}.job ADD COLUMN blocking boolean NOT NULL DEFAULT false`, `ALTER TABLE ${schema}.job ADD COLUMN pending_dependencies int NOT NULL DEFAULT 0`, ` CREATE TABLE IF NOT EXISTS ${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) ) `, `CREATE INDEX IF NOT EXISTS job_dep_parent_idx ON ${schema}.job_dependency (parent_name, parent_id)`, // NOTE: the v31 job_i5 rebuild (adding `AND NOT blocked`) is intentionally omitted, v33 // drops and rebuilds job_i5 again (slimming off the covering INCLUDE), so on a multi-version // upgrade (<= v31 -> >= v33, applied as one migration transaction) this only built a covering // index that v33 immediately throws away. The whole migration runs in a single transaction, // so no worker observes the pre-v33 shape; the old job_i5 simply persists untouched until v33 // replaces it. Anyone who already migrated to exactly v31/v32 keeps the index they built then, // so removing the build here does not affect them. New partitions created while on v31 still // get the correct shape from createQueueFn[31] below. createQueueFn[31](schema) ], uninstall: [ `DROP INDEX IF EXISTS ${schema}.job_dep_parent_idx`, `DROP TABLE IF EXISTS ${schema}.job_dependency`, createQueueFn[30](schema), `SELECT ${schema}.job_table_run($cmd$DROP INDEX IF EXISTS ${schema}.job_i5$cmd$)`, `SELECT ${schema}.job_table_run($cmd$CREATE INDEX job_i5 ON ${schema}.job (name, start_after) INCLUDE (priority, created_on, id) WHERE state < 'active'$cmd$)`, `ALTER TABLE ${schema}.job DROP COLUMN pending_dependencies`, `ALTER TABLE ${schema}.job DROP COLUMN blocking`, `ALTER TABLE ${schema}.job DROP COLUMN blocked` ] }, { release: '12.21.0', version: 32, previous: 31, install: [ `ALTER TABLE ${schema}.queue ADD COLUMN notify boolean NOT NULL DEFAULT false`, createQueueFn[32](schema), `ALTER TABLE ${schema}.queue ADD COLUMN failed_count int NOT NULL DEFAULT 0`, `ALTER TABLE ${schema}.queue ADD COLUMN ready_count int NOT NULL DEFAULT 0` ], uninstall: [ `ALTER TABLE ${schema}.queue DROP COLUMN ready_count`, `ALTER TABLE ${schema}.queue DROP COLUMN failed_count`, createQueueFn[31](schema), `ALTER TABLE ${schema}.queue DROP COLUMN notify` ] }, { release: '12.22.0', version: 33, previous: 32, install: [ `ALTER TABLE ${schema}.version ADD COLUMN IF NOT EXISTS flow_on timestamp with time zone`, // The job_i9 build and job_i5 reshape are run OFF the migration transaction, via BAM as // CONCURRENTLY DDL, see the `async` block below. The original v33 ran them synchronously here // via job_table_run(), taking SHARE/ACCESS EXCLUSIVE locks on job_common + every partition // inside the migration transaction, which deadlocked live workers polling job_common during a // rolling deploy (issue #832). // // Only databases that have NOT yet passed v33 execute this install, and they always have the // covering job_i5 (the slim form is introduced here), so the reshape never needlessly // rebuilds an already-slim index. Databases already past v33 keep what they built then; they // pick up only the bam default change, carried separately by migration v36. // // Must run BEFORE the enqueues below and therefore in v33, not v36: migrations apply in // version order, and the ordered pair enqueued here needs the default already changed. See // createTableBam for why the default is clock_timestamp(). `ALTER TABLE ${schema}.bam ALTER COLUMN created_on SET DEFAULT clock_timestamp()`, createQueueFn[33](schema) ], async: [ `SELECT ${schema}.job_table_run_async( 'flow_resolver_index', $VERSION$, $$ CREATE INDEX CONCURRENTLY IF NOT EXISTS job_i9 ON ${schema}.job (name, id) WHERE blocking AND state = 'completed' $$ )`, // CONCURRENTLY cannot reshape in place, so this is an ordered drop-then-rebuild pair; BAM // applies them in created_on order (see the clock_timestamp() default set above). `SELECT ${schema}.job_table_run_async( 'fetch_index_drop', $VERSION$, $$ DROP INDEX CONCURRENTLY IF EXISTS ${schema}.job_i5 $$ )`, `SELECT ${schema}.job_table_run_async( 'fetch_index', $VERSION$, $$ CREATE INDEX CONCURRENTLY IF NOT EXISTS job_i5 ON ${schema}.job (name, start_after) WHERE state < 'active' AND NOT blocked $$ )` ], uninstall: [ createQueueFn[32](schema), `SELECT ${schema}.job_table_run($cmd$DROP INDEX IF EXISTS ${schema}.job_i5$cmd$)`, `SELECT ${schema}.job_table_run($cmd$CREATE INDEX job_i5 ON ${schema}.job (name, start_after) INCLUDE (priority, created_on, id) WHERE state < 'active' AND NOT blocked$cmd$)`, `SELECT ${schema}.job_table_run($cmd$DROP INDEX IF EXISTS ${schema}.job_i9$cmd$)`, `ALTER TABLE ${schema}.version DROP COLUMN flow_on` ] }, { release: '12.23.0', version: 34, previous: 33, // Plain columns on the partitioned parent cascade to job_common (DEFAULT partition) and every // existing/future partition, so no job_table_run fan-out or createQueueFn bump is needed. install: [ `ALTER TABLE ${schema}.job ADD COLUMN IF NOT EXISTS source_name text`, `ALTER TABLE ${schema}.job ADD COLUMN IF NOT EXISTS source_id uuid`, `ALTER TABLE ${schema}.job ADD COLUMN IF NOT EXISTS source_created_on timestamp with time zone`, `ALTER TABLE ${schema}.job ADD COLUMN IF NOT EXISTS source_retry_count int` ], uninstall: [ `ALTER TABLE ${schema}.job DROP COLUMN source_name`, `ALTER TABLE ${schema}.job DROP COLUMN source_id`, `ALTER TABLE ${schema}.job DROP COLUMN source_created_on`, `ALTER TABLE ${schema}.job DROP COLUMN source_retry_count` ] }, { release: '12.24.0', version: 35, previous: 34, // Mirror plans.create(): honor noTablePartitioning so upgrades on non-partitioning // deployments (e.g. CockroachDB, which rejects declarative RANGE partitioning) get a plain // queue_stats table instead of a partitioned one they could never maintain. noCovering is a // separate axis (CockroachDB sets it, YugabyteDB doesn't) gating the index's covering INCLUDE. // queue.ready_history is NOT NULL DEFAULT '{}' so existing rows backfill with an empty window // that fills in over the next monitor cycles. install: [ ...(noPartitioning ? [ createTableQueueStatsFn[35](schema, true), createIndexQueueStatsFn[35](schema, noCovering) ] : [ createTableQueueStatsFn[35](schema, false), createIndexQueueStatsFn[35](schema, noCovering), ensureQueueStatsPartitionsFn[35](schema) ]), `ALTER TABLE ${schema}.queue ADD COLUMN ready_history int[] NOT NULL DEFAULT '{}'` ], uninstall: [ `ALTER TABLE ${schema}.queue DROP COLUMN ready_history`, `DROP TABLE IF EXISTS ${schema}.queue_stats` ] }, { release: '12.24.1', version: 36, previous: 35, // Carry only the bam.created_on default change (now() -> clock_timestamp()) to databases that // already ran v33 and so won't re-run its install. This keeps a fully-migrated database's schema // identical to a fresh install (plans.create builds the bam table with this default), and lets a // future async migration enqueue an ordered drop-then-rebuild without the now() tie. The ALTER // is idempotent, so databases that just ran v33's copy of it (a multi-version upgrade) are // unaffected. No index work here. That lives in v33, which the deadlock-affected (pre-v33) // databases run; databases already past v33 keep the indexes they built and skip the churn. install: [ `ALTER TABLE ${schema}.bam ALTER COLUMN created_on SET DEFAULT clock_timestamp()` ], // The default change is forward-compatible and harmless to keep, so rollback leaves it in place. uninstall: [] }, { release: '12.26.0', version: 37, previous: 36, // Only installed where partitioning is enabled, plans.create() creates job_table_format() // solely in that case, so a noPartitioning database has none to replace. install: noPartitioning ? [] : [jobTableFormatFn[37](schema)], // Restore the prior (naive) definition on rollback so the schema matches v36 exactly. uninstall: noPartitioning ? [] : [jobTableFormatFn[36](schema)] }, { release: '12.28.0', version: 38, previous: 37, install: noPartitioning // noCovering backends (CockroachDB) have no INCLUDE, so they get the narrow index. The // partitioned path is PostgreSQL-only and always covering. ? [`CREATE INDEX job_i10 ON ${schema}.job (name, singleton_key, state DESC, created_on, id)${noCovering ? '' : ' INCLUDE (start_after)'} WHERE state < 'active' AND NOT blocked AND policy = 'key_strict_fifo'`] : [createQueueFn[38](schema)], async: noPartitioning ? [] : [ { name: 'key_strict_fifo_head_index', partitionPolicy: 'key_strict_fifo', command: `CREATE INDEX CONCURRENTLY IF NOT EXISTS job_i10 ON ${schema}.job (name, singleton_key, state DESC, created_on, id) INCLUDE (start_after) WHERE state < 'active' AND NOT blocked AND policy = 'key_strict_fifo'` } ], uninstall: noPartitioning ? [`DROP INDEX IF EXISTS ${schema}.job_i10`] : [ createQueueFn[33](schema), `SELECT ${schema}.job_table_run($cmd$DROP INDEX IF EXISTS ${schema}.job_i10$cmd$)` ] }, { release: '12.29.0', version: 39, previous: 38, install: [ `ALTER TABLE ${schema}.version ADD COLUMN IF NOT EXISTS reindex_on timestamp with time zone` ], uninstall: [ `ALTER TABLE ${schema}.version DROP COLUMN reindex_on` ] }, { release: '12.30.0', version: 40, previous: 39, install: [ `ALTER TABLE ${schema}.version ADD COLUMN IF NOT EXISTS monitor_backoff_on timestamp with time zone`, `ALTER TABLE ${schema}.queue ADD COLUMN IF NOT EXISTS monitor_claim_on timestamp with time zone`, // Seed the claim from the existing pace so an upgrade does not make every queue immediately // eligible and stampede one aggregate per queue on the first supervise pass after deploy. // // Left out on a backend that cannot write a column in the transaction that added it // (CockroachDB, see noAddColumnBackfill), where it fails the whole migration with // "column is being backfilled" and is what kept a CockroachDB deployment on schema 39 from // reaching 40 at all. Nothing is lost by omitting it: trySetQueueMonitorTime falls back to // monitor_on when the claim is NULL, which is the value this statement would have written. ...(noAddColumnBackfill ? [] : [`UPDATE ${schema}.queue SET monitor_claim_on = monitor_on WHERE monitor_claim_on IS NULL`]), ...(noPartitioning // Single transaction, so both statements commit together and no window exists to protect. ? [ `CREATE INDEX job_i11 ON ${schema}.job (name, priority DESC, created_on, start_after) WHERE state < 'active' AND NOT blocked`, `DROP INDEX IF EXISTS ${schema}.job_i5` ] : [createQueueFn[40](schema)]) ], // Build before retire. BAM applies these in created_on order, so job_i11 is complete on a // table before job_i5 is dropped from it. async: noPartitioning ? [] : [ { name: 'fetch_index_priority_build', command: `CREATE INDEX CONCURRENTLY IF NOT EXISTS job_i11 ON ${schema}.job (name, priority DESC, created_on, start_after) WHERE state < 'active' AND NOT blocked` }, { name: 'fetch_index_retire', command: `DROP INDEX CONCURRENTLY IF EXISTS ${schema}.job_i5` } ], // Restore before retire, the mirror of the install. Rollback can land at any point in the async // sequence (before the build, between build and retire, or after both) and IF NOT EXISTS is // what makes all three the same statement: v40 never reshapes job_i5, it only drops it, so a // job_i5 still present is already the v39 shape and is left exactly where it is rather than // being dropped and rebuilt for nothing. uninstall: [ ...(noPartitioning ? [ `CREATE INDEX IF NOT EXISTS job_i5 ON ${schema}.job (name, start_after) WHERE state < 'active' AND NOT blocked`, `DROP INDEX IF EXISTS ${schema}.job_i11` ] : [ createQueueFn[38](schema), `SELECT ${schema}.job_table_run($cmd$CREATE INDEX IF NOT EXISTS job_i5 ON ${schema}.job (name, start_after) WHERE state < 'active' AND NOT blocked$cmd$)`, `SELECT ${schema}.job_table_run($cmd$DROP INDEX IF EXISTS ${schema}.job_i11$cmd$)` ]), `ALTER TABLE ${schema}.queue DROP COLUMN monitor_claim_on`, `ALTER TABLE ${schema}.version DROP COLUMN monitor_backoff_on` ] }, { release: '12.31.0', version: 41, previous: 40, // `kind` says which format the expression in `cron` is in. The default labels every row cron, // which is what a table this migration has never seen holds: cron was the only format a // schedule could be written in. The UPDATE is for the table it has seen before, since // `uninstall` drops the column rather than remembering it, so a rollback to v40 and a // re-upgrade would otherwise relabel every rule as cron from the default and leave a row that // reads fine and never fires. Reading the expression puts the label back. The two patterns // are the detection isRrule() performs, in the terms both postgres and CockroachDB's regexp // engine share: `^` anchors the whole string in one and not the other, so a property on a // line below the first is matched on the whitespace before it instead. No cron expression // matches either, and cannot, since no cron field contains `=`, `:` or `;`. // // On a backend that cannot write a column in the transaction that added it (CockroachDB, see // noAddColumnBackfill) the label is left to the pass instead: the whole migration is one // transaction, and the UPDATE fails there with "column is being backfilled". Nothing is lost. // A row whose kind disagrees with its expression is read the way it is written and relabelled // by setScheduleKinds on the first pass that reaches it, which is the same fallback a rolling // upgrade relies on, and that statement CockroachDB accepts. What it costs is the window // before that pass: getSchedules() reports `cron` for a rule row that has not been read yet. // // `last_job_id` is the job each schedule most recently produced, so a schedule can be joined // to its last run. Nullable and unconstrained on purpose: the referenced job is subject to // retention and will eventually be deleted, so a foreign key would either block retention or // null the column back out. // // The zone backfill is on the same pass because the same rows are already being read. A null // timezone is how `schedule({ tz: null })` landed on a release whose destructuring default // only covered `undefined`, and how a row written with SQL leaves it out. Nothing ever chose // it: schedule()'s default is UTC and the docs say UTC, but cron-parser reads a non-string // zone as unset, so those rows have been firing in the local zone of whichever instance took // the pass. UTC is what their authors asked for. The default on the column keeps the next // hand-written insert from making another one, and matches what a fresh install now builds. install: [ `ALTER TABLE ${schema}.schedule ADD COLUMN IF NOT EXISTS kind text NOT NULL DEFAULT 'cron' CHECK (kind IN ('cron', 'rrule'))`, ...(noAddColumnBackfill ? [] : [`UPDATE ${schema}.schedule SET kind = 'rrule' WHERE kind = 'cron' AND (cron ~* '(^|[[:space:]]|;)FREQ=' OR cron ~* '(^|[[:space:]])(DTSTART|RRULE|RDATE|EXDATE)[;:]')`]), `UPDATE ${schema}.schedule SET timezone = 'UTC' WHERE timezone IS NULL`, `ALTER TABLE ${schema}.schedule ALTER COLUMN timezone SET DEFAULT 'UTC'`, `ALTER TABLE ${schema}.schedule ADD COLUMN IF NOT EXISTS last_job_id uuid` ], // Dropping `kind` drops its CHECK with it, since the constraint belongs to the column. The // zone default is left in place: a v40 instance names the column on every write, so the // default it would fall back on never applies, and keeping it costs a rollback nothing. uninstall: [ `ALTER TABLE ${schema}.schedule DROP COLUMN kind`, `ALTER TABLE ${schema}.schedule DROP COLUMN last_job_id` ] }, /* eslint-enable no-restricted-syntax */ { release: '12.33.0', version: 42, previous: 41, install: [ `CREATE OR REPLACE FUNCTION ${schema}.job_now() RETURNS timestamp with time zone LANGUAGE sql STABLE AS $function$ SELECT pg_catalog.now(); $function$ `, noPartitioning ? createQueueNoPartitionFn[42](schema) : createQueueFn[42](schema) ], uninstall: [ noPartitioning ? createQueueNoPartitionFn[41](schema) : createQueueFn[40](schema), `DROP FUNCTION ${schema}.job_now()` ] }, { release: '12.35.0', version: 43, previous: 42, // No statement rewrites a table: every added column is nullable or has a constant default, // which matters on queue_stats, partitioned and large on a busy installation. Existing rows // are not backfilled. Snapshots already in queue_stats keep null deltas, jobs already in a dead // letter queue keep the output they were copied with, and a row dead-lettered before // source_root_id existed has none, so its first redrive falls back to source_id. // // job_i12 is built the way v40 built job_i11: inline without partitioning, otherwise through // create_queue for new partitions and BAM, concurrently, for the tables that already exist. // // One ALTER per table, not one per column. On CockroachDB every ALTER is a schema change job // inside the migration's transaction, and seven of them held it open long enough to lose // RETRY_SERIALIZABLE against DDL on another schema in the same cluster. install: [ `ALTER TABLE ${schema}.queue ADD COLUMN created_delta int NOT NULL DEFAULT 0, ADD COLUMN completed_delta int NOT NULL DEFAULT 0, ADD COLUMN failed_delta int NOT NULL DEFAULT 0, ADD COLUMN delta_on timestamp with time zone, ADD COLUMN delta_seconds int`, `ALTER TABLE ${schema}.queue_stats ADD COLUMN created_delta int, ADD COLUMN completed_delta int, ADD COLUMN failed_delta int, ADD COLUMN delta_seconds int, ADD COLUMN delta_on timestamp with time zone`, `ALTER TABLE ${schema}.job ADD COLUMN IF NOT EXISTS source_output jsonb, ADD COLUMN IF NOT EXISTS source_root_id uuid`, noPartitioning ? `CREATE INDEX job_i12 ON ${schema}.job (source_root_id) WHERE source_root_id IS NOT NULL` : createQueueFn[43](schema) ], async: noPartitioning ? [] : [ { name: 'source_root_index_build', command: `CREATE INDEX CONCURRENTLY IF NOT EXISTS job_i12 ON ${schema}.job (source_root_id) WHERE source_root_id IS NOT NULL` } ], // The index goes before its column. Dropping the column would take the index with it, but // partitioned, create_queue has to be restored first, or a queue created mid-rollback would // build an index on a column that is about to be gone. uninstall: [ ...(noPartitioning ? [`DROP INDEX IF EXISTS ${schema}.job_i12`] : [ createQueueFn[42](schema), `SELECT ${schema}.job_table_run($cmd$DROP INDEX IF EXISTS ${schema}.job_i12$cmd$)` ]), `ALTER TABLE ${schema}.queue DROP COLUMN created_delta, DROP COLUMN completed_delta, DROP COLUMN failed_delta, DROP COLUMN delta_on, DROP COLUMN delta_seconds`, `ALTER TABLE ${schema}.queue_stats DROP COLUMN created_delta, DROP COLUMN completed_delta, DROP COLUMN failed_delta, DROP COLUMN delta_seconds, DROP COLUMN delta_on`, `ALTER TABLE ${schema}.job DROP COLUMN source_output, DROP COLUMN source_root_id` ] }, { release: '12.36.0', version: 44, previous: 43, // The instance table starts empty: each instance registers on its next start(). // blocked_count's constant default adds it without a table rewrite, and existing queues read 0 // until the next monitor pass counts them. It goes on the queue row only, not queue_stats, so // the covering index on queue_stats needs no rebuild. // The histogram and ready_oldest_seconds columns are nullable with no default, like the v43 // deltas, so no statement rewrites a table, and a snapshot captured before them reads as not // counted rather than as a queue with no waits. // trace_context is nullable with no default too, so adding it rewrites no rows. install: [ /* eslint-disable no-restricted-syntax -- column defaults stay on the real clock: every pg-boss write names its timestamps through job_now() */ `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 */ `ALTER TABLE ${schema}.queue ADD COLUMN blocked_count int NOT NULL DEFAULT 0, ADD COLUMN wait_bins int[], ADD COLUMN run_bins int[], ADD COLUMN ready_oldest_seconds int`, `ALTER TABLE ${schema}.queue_stats ADD COLUMN wait_bins int[], ADD COLUMN run_bins int[], ADD COLUMN ready_oldest_seconds int`, `ALTER TABLE ${schema}.job ADD COLUMN IF NOT EXISTS trace_context jsonb` ], uninstall: [ `DROP TABLE ${schema}.instance`, `ALTER TABLE ${schema}.queue DROP COLUMN blocked_count, DROP COLUMN wait_bins, DROP COLUMN run_bins, DROP COLUMN ready_oldest_seconds`, `ALTER TABLE ${schema}.queue_stats DROP COLUMN wait_bins, DROP COLUMN run_bins, DROP COLUMN ready_oldest_seconds`, `ALTER TABLE ${schema}.job DROP COLUMN trace_context` ] }, { release: '12.37.0', version: 45, previous: 44, // upsert_by_key is nullable with no default, so adding it rewrites no rows, and no existing job // carries it. job_i13 therefore starts empty and cannot fail on existing duplicates. // It is built the way v43 built job_i12: inline without partitioning, otherwise through // create_queue for new partitions and BAM, concurrently, for the tables that already exist. install: [ `ALTER TABLE ${schema}.job ADD COLUMN IF NOT EXISTS upsert_by_key bool`, noPartitioning ? `CREATE UNIQUE INDEX job_i13 ON ${schema}.job (name, singleton_key) WHERE state = 'created' AND upsert_by_key` : createQueueFn[45](schema) ], async: noPartitioning ? [] : [ { name: 'upsert_index_build', command: `CREATE UNIQUE INDEX CONCURRENTLY IF NOT EXISTS job_i13 ON ${schema}.job (name, singleton_key) WHERE state = 'created' AND upsert_by_key` } ], // The index goes before its column, and create_queue is restored first, as in v43. CASCADE because // CockroachDB will not drop a unique index without it. uninstall: [ ...(noPartitioning ? [`DROP INDEX IF EXISTS ${schema}.job_i13 CASCADE`] : [ createQueueFn[43](schema), `SELECT ${schema}.job_table_run($cmd$DROP INDEX IF EXISTS ${schema}.job_i13$cmd$)` ]), `ALTER TABLE ${schema}.job DROP COLUMN upsert_by_key` ] } ] } export { rollback, next, migrate, migrateCommands, getAll, getAllForConfig, getMinVersion, }