diff --git a/README.md b/README.md index 71eafd1..856408f 100644 --- a/README.md +++ b/README.md @@ -267,7 +267,7 @@ const { jobId, deduped } = await SaveDraftJob.dispatch({ content: '...' }) - `retryJob` does not touch the dedup entry — a retried job continues to occupy the dedup slot. TTL runs on wall-clock time, so long-running retries may outlive the TTL window. Use a generous TTL or no TTL if retries must stay deduped. - Atomicity: - **Redis**: a single Lua script per dispatch performs the dedup-key lookup, state check (pending/delayed ZSCORE), payload swap, and TTL refresh atomically. - - **Knex/Kysely**: transactional `SELECT ... FOR UPDATE` + insert/update inside a transaction. On PostgreSQL and SQLite, a partial unique index makes concurrent first dispatches race-free: a savepoint catches the unique-constraint violation and returns `{ deduped: 'skipped' }` pointing at the winner. MySQL has no partial unique index, see the caveat below. + - **Knex/Kysely**: a unique index on `(queue, dedup_id)` lets a single job own a dedup id, whatever its status. Each dispatch runs in a transaction: an existing owner is locked with `SELECT ... FOR UPDATE`, then skipped, extended or replaced under that lock. When there is no owner to lock, concurrent first dispatches race on the unique index: one inserts the job, the others run again, lock the winner and return `{ deduped: 'skipped' }` pointing at it. Works the same on PostgreSQL, MySQL and SQLite. - **SyncAdapter**: executes inline, no dedup support. ### Caveats @@ -278,7 +278,6 @@ const { jobId, deduped } = await SaveDraftJob.dispatch({ content: '...' }) - Scheduled jobs (`.schedule()`) do not support dedup — each cron/interval fire is an independent dispatch. - With no `ttl`, dedup persists until the job is removed (completed/failed without retention). When retention keeps the record, re-dispatch stays blocked until the record is pruned. - With `ttl`, dedup expires after the window — a new job (new UUID) is created. The old job still runs. -- Knex/Kysely MySQL concurrent race: MySQL does not support partial unique indexes, so two `pushOn` calls with the same dedup id firing at the exact same instant can both succeed. Serialize at the app layer if strict guarantees are required, or use Postgres / SQLite / Redis (all of which serialize correctly via the partial unique index or Lua atomicity). ## Job History & Retention @@ -473,6 +472,23 @@ await schema.dropJobsTable() +#### Migrating SQL deduplication after an upgrade + +Jobs tables created by earlier versions have no unique index on `(queue, dedup_id)`, and dispatches racing +with a worker could enqueue a job twice. Run `addDedupColumns()` once from a migration to create it: + +```typescript +// Knex +await new KnexQueueSchemaService(connection).addDedupColumns('queue_jobs') + +// Kysely +await new KyselyQueueSchemaService(db, { dialect: 'postgres' }).addDedupColumns('queue_jobs') +``` + +The migration is idempotent. When several jobs share a dedup id, the latest one keeps it and the +others only lose their dedup id: no job is removed. The previous dedup indexes are dropped. Stop +dispatching jobs with `.dedup()` while it runs. + #### Migrating SQL schedules after an upgrade The schedules table stores its dates as epoch milliseconds (`bigint`), so they do not depend on the diff --git a/src/drivers/knex_adapter.ts b/src/drivers/knex_adapter.ts index 6b59e08..c3bdb95 100644 --- a/src/drivers/knex_adapter.ts +++ b/src/drivers/knex_adapter.ts @@ -503,19 +503,39 @@ export class KnexAdapter implements Adapter { ) } - #pushWithDedupTxn( + /** + * A single job owns a dedup slot, whatever its status: the unique + * (queue, dedup_id) index rejects a second owner. An existing owner is + * locked for the transaction. When there is none to lock and a concurrent + * dispatch inserts first, the transaction runs again and locks that owner. + */ + async #pushWithDedupTxn( queue: string, jobData: JobData, insertRow: Record, dedup: NonNullable ): Promise { - return this.#connection.transaction(async (trx) => { - const now = Date.now() - const existingResult = await this.#resolveExistingDedup(trx, queue, jobData, dedup, now) - if (existingResult) return existingResult - - return this.#insertDedup(trx, queue, jobData.id, insertRow, dedup, now) - }) + while (true) { + try { + return await this.#connection.transaction(async (trx) => { + const now = Date.now() + const existingResult = await this.#resolveExistingDedup(trx, queue, jobData, dedup, now) + if (existingResult) return existingResult + + await trx(this.#jobsTable).insert({ + ...insertRow, + dedup_id: dedup.id, + dedup_at: now, + dedup_ttl: dedup.ttl ?? null, + }) + return { outcome: 'added' as DedupOutcome, jobId: jobData.id } + }) + } catch (err) { + if (!this.#isUniqueViolation(err) && !this.#isDeadlock(err)) throw err + // The job id itself is the duplicate + if (await this.getJob(jobData.id, queue)) throw err + } + } } async #resolveExistingDedup( @@ -525,14 +545,20 @@ export class KnexAdapter implements Adapter { dedup: NonNullable, now: number ): Promise { - const existing = await trx(this.#jobsTable) + const owner = await trx(this.#jobsTable) .where('queue', queue) .where('dedup_id', dedup.id) .orderBy('dedup_at', 'desc') - .forUpdate() - .first() + .first('id') - if (!existing) return null + if (!owner) return null + + // The owner is locked through its primary key, as workers lock a job. + // Locking it through the dedup index deadlocks with them on MySQL. + const existing = await trx(this.#jobsTable).where({ id: owner.id, queue }).forUpdate().first() + + // Removed or released in the meantime: the insert tells whether the slot is free + if (existing?.dedup_id !== dedup.id) return null const dedupAt = existing.dedup_at != null ? Number(existing.dedup_at) : null const dedupTtl = existing.dedup_ttl != null ? Number(existing.dedup_ttl) : null @@ -563,67 +589,31 @@ export class KnexAdapter implements Adapter { } // Release the expired dedup slot. The old job keeps running to completion. - const status = existing.status as JobStatus - if (status === 'pending' || status === 'delayed' || status === 'active') { - await trx(this.#jobsTable) - .where({ id: existing.id, queue }) - .update({ dedup_id: null, dedup_at: null, dedup_ttl: null }) - } + await trx(this.#jobsTable) + .where({ id: existing.id, queue }) + .update({ dedup_id: null, dedup_at: null, dedup_ttl: null }) return null } - async #insertDedup( - trx: Knex.Transaction, - queue: string, - jobId: string, - insertRow: Record, - dedup: NonNullable, - now: number - ): Promise { - let raceLost = false - try { - await trx.transaction(async (sp) => { - await sp(this.#jobsTable).insert({ - ...insertRow, - dedup_id: dedup.id, - dedup_at: now, - dedup_ttl: dedup.ttl ?? null, - }) - }) - } catch (err) { - if (this.#isUniqueViolation(err)) { - raceLost = true - } else { - throw err - } - } - - if (raceLost) { - const winner = await trx(this.#jobsTable) - .where('queue', queue) - .where('dedup_id', dedup.id) - .whereIn('status', ['pending', 'delayed']) - .orderBy('dedup_at', 'desc') - .first() - if (winner) { - return { outcome: 'skipped' as DedupOutcome, jobId: winner.id as string } - } - } - - return { outcome: 'added' as DedupOutcome, jobId } - } - #isUniqueViolation(err: unknown): boolean { if (!err || typeof err !== 'object') return false const e = err as { code?: string; message?: string } return ( e.code === '23505' || e.code === 'SQLITE_CONSTRAINT_UNIQUE' || + e.code === 'ER_DUP_ENTRY' || /UNIQUE constraint/i.test(e.message ?? '') ) } + /** + * InnoDB picks a victim among dispatches inserting the same dedup id. + */ + #isDeadlock(err: unknown): boolean { + return (err as { code?: string } | null)?.code === 'ER_LOCK_DEADLOCK' + } + async pushMany(jobs: JobData[]): Promise { return this.pushManyOn('default', jobs) } diff --git a/src/drivers/kysely_adapter.ts b/src/drivers/kysely_adapter.ts index e808c93..80d4d92 100644 --- a/src/drivers/kysely_adapter.ts +++ b/src/drivers/kysely_adapter.ts @@ -465,29 +465,49 @@ export class KyselyAdapter implements Adapter { await this.#jobs(this.#connection).insertInto(this.#jobsTable).values(row).execute() } + /** + * A single job owns a dedup slot, whatever its status: the unique + * (queue, dedup_id) index rejects a second owner. An existing owner is + * locked for the transaction. When there is none to lock and a concurrent + * dispatch inserts first, the transaction runs again and locks that owner. + */ async #pushWithDedup( connection: Kysely, queue: string, jobData: JobData, insertRow: Partial & Pick ): Promise { - try { - return await this.#withTransaction(connection, async (trx) => { - const now = Date.now() - const dedup = jobData.dedup! - const existingResult = await this.#resolveExistingDedup(trx, queue, jobData, dedup, now) - if (existingResult) return existingResult - - return this.#insertDedup(trx, queue, jobData.id, insertRow, dedup, now) - }) - } catch (error) { - if (this.#isMissingDedupColumn(error)) { - throw new Error( - `Dedup columns missing on "${this.#jobsTable}". Run KyselyQueueSchemaService.addDedupColumns() before dispatching jobs with .dedup().`, - { cause: error } - ) + const dedup = jobData.dedup! + + while (true) { + try { + return await this.#withTransaction(connection, async (trx) => { + const now = Date.now() + const existingResult = await this.#resolveExistingDedup(trx, queue, jobData, dedup, now) + if (existingResult) return existingResult + + await this.#jobs(trx) + .insertInto(this.#jobsTable) + .values({ + ...insertRow, + dedup_id: dedup.id, + dedup_at: now, + dedup_ttl: dedup.ttl ?? null, + }) + .execute() + return { outcome: 'added', jobId: jobData.id } + }) + } catch (error) { + if (this.#isMissingDedupColumn(error)) { + throw new Error( + `Dedup columns missing on "${this.#jobsTable}". Run KyselyQueueSchemaService.addDedupColumns() before dispatching jobs with .dedup().`, + { cause: error } + ) + } + if (!this.#isUniqueViolation(error) && !this.#isDeadlock(error)) throw error + // The job id itself is the duplicate + if (await this.getJob(jobData.id, queue)) throw error } - throw error } } @@ -498,18 +518,31 @@ export class KyselyAdapter implements Adapter { dedup: NonNullable, now: number ): Promise { - let existingQuery = this.#jobs(trx) + const owner = await this.#jobs(trx) .selectFrom(this.#jobsTable) - .selectAll() + .select('id') .where('queue', '=', queue) .where('dedup_id', '=', dedup.id) .orderBy('dedup_at', 'desc') .limit(1) + .executeTakeFirst() + + if (!owner) return null + + // The owner is locked through its primary key, as workers lock a job. + // Locking it through the dedup index deadlocks with them on MySQL. + let existingQuery = this.#jobs(trx) + .selectFrom(this.#jobsTable) + .selectAll() + .where('id', '=', owner.id) + .where('queue', '=', queue) if (this.#supportsSkipLocked()) existingQuery = existingQuery.forUpdate() const existing = await existingQuery.executeTakeFirst() - if (!existing) return null + + // Removed or released in the meantime: the insert tells whether the slot is free + if (existing?.dedup_id !== dedup.id) return null const dedupAt = existing.dedup_at == null ? null : Number(existing.dedup_at) const dedupTtl = existing.dedup_ttl == null ? null : Number(existing.dedup_ttl) @@ -547,62 +580,15 @@ export class KyselyAdapter implements Adapter { return { outcome: 'skipped', jobId: existing.id } } - if ( - existing.status === 'pending' || - existing.status === 'delayed' || - existing.status === 'active' - ) { - await this.#jobs(trx) - .updateTable(this.#jobsTable) - .set({ dedup_id: null, dedup_at: null, dedup_ttl: null }) - .where('id', '=', existing.id) - .where('queue', '=', queue) - .execute() - } - - return null - } - - async #insertDedup( - trx: Transaction, - queue: string, - jobId: string, - insertRow: Partial & Pick, - dedup: NonNullable, - now: number - ): Promise { - const savepoint = `queue_dedup_${randomUUID().replaceAll('-', '')}` - await sql`savepoint ${sql.id(savepoint)}`.execute(trx) - - try { - await this.#jobs(trx) - .insertInto(this.#jobsTable) - .values({ - ...insertRow, - dedup_id: dedup.id, - dedup_at: now, - dedup_ttl: dedup.ttl ?? null, - }) - .execute() - await sql`release savepoint ${sql.id(savepoint)}`.execute(trx) - return { outcome: 'added', jobId } - } catch (error) { - await sql`rollback to savepoint ${sql.id(savepoint)}`.execute(trx) - await sql`release savepoint ${sql.id(savepoint)}`.execute(trx) - if (!this.#isUniqueViolation(error)) throw error - } - - const winner = await this.#jobs(trx) - .selectFrom(this.#jobsTable) - .select('id') + // Release the expired dedup slot. The old job keeps running to completion. + await this.#jobs(trx) + .updateTable(this.#jobsTable) + .set({ dedup_id: null, dedup_at: null, dedup_ttl: null }) + .where('id', '=', existing.id) .where('queue', '=', queue) - .where('dedup_id', '=', dedup.id) - .where('status', 'in', ['pending', 'delayed']) - .orderBy('dedup_at', 'desc') - .executeTakeFirst() + .execute() - if (!winner) throw new Error(`Unable to resolve concurrent dedup dispatch for "${dedup.id}"`) - return { outcome: 'skipped', jobId: winner.id } + return null } async pushMany(jobs: JobData[]): Promise { @@ -962,6 +948,13 @@ export class KyselyAdapter implements Adapter { ) } + /** + * InnoDB picks a victim among dispatches inserting the same dedup id. + */ + #isDeadlock(error: unknown): boolean { + return (error as { code?: string } | null)?.code === 'ER_LOCK_DEADLOCK' + } + #isMissingDedupColumn(error: unknown): boolean { if (!error || typeof error !== 'object') return false const message = (error as { message?: string }).message ?? '' diff --git a/src/services/knex_queue_schema.ts b/src/services/knex_queue_schema.ts index 5a69c33..00e0fd2 100644 --- a/src/services/knex_queue_schema.ts +++ b/src/services/knex_queue_schema.ts @@ -51,20 +51,22 @@ export class KnexQueueSchemaService { table.index(['queue', 'status', 'score']) table.index(['queue', 'status', 'execute_at']) table.index(['queue', 'status', 'finished_at']) - table.index(['queue', 'dedup_id']) extend?.(table) }) - await this.#createDedupActiveUniqueIndex(tableName) + await this.#createDedupUniqueIndex(tableName) } /** * Idempotent migration: adds dedup columns (dedup_id, dedup_at, dedup_ttl) - * and a (queue, dedup_id) index to an existing jobs table. + * and a unique (queue, dedup_id) index to an existing jobs table. * - * Safe to run multiple times. Uses hasColumn checks so it won't fail on re-runs. - * For large Postgres tables, consider pausing workers during the run. + * Safe to run multiple times. Run it again when upgrading from a version + * that created a non-unique index: jobs sharing a dedup id keep running, and + * only the latest one keeps the dedup slot. + * + * Stop dispatching jobs with .dedup() during the run. */ async addDedupColumns(tableName: string = 'queue_jobs'): Promise { const hasDedupId = await this.#connection.schema.hasColumn(tableName, 'dedup_id') @@ -79,30 +81,74 @@ export class KnexQueueSchemaService { }) } - if (!hasDedupId) { - await this.#connection.schema.alterTable(tableName, (table) => { - table.index(['queue', 'dedup_id']) - }) + // The derived table is what allows MySQL to read the table it updates. + await this.#connection.raw( + `UPDATE ?? SET dedup_id = NULL, dedup_at = NULL, dedup_ttl = NULL + WHERE (id, queue) IN ( + SELECT id, queue FROM ( + SELECT DISTINCT older.id, older.queue + FROM ?? AS older + JOIN ?? AS newer ON newer.queue = older.queue AND newer.dedup_id = older.dedup_id + WHERE newer.dedup_at > older.dedup_at + OR (newer.dedup_at = older.dedup_at AND newer.id > older.id) + ) AS stale + )`, + [tableName, tableName, tableName] + ) + + await this.#createDedupUniqueIndex(tableName) + await this.#dropLegacyDedupIndexes(tableName) + } + + /** + * Unique index on (queue, dedup_id): a single job owns a dedup slot, whatever + * its status. Jobs without dedup hold NULL, which is never a duplicate. + */ + async #createDedupUniqueIndex(tableName: string): Promise { + const indexName = `${tableName}_dedup_uidx` + + if (this.#dialect() !== 'mysql') { + await this.#connection.raw('CREATE UNIQUE INDEX IF NOT EXISTS ?? ON ?? (??, ??)', [ + indexName, + tableName, + 'queue', + 'dedup_id', + ]) + return } - await this.#createDedupActiveUniqueIndex(tableName) + // MySQL has no CREATE INDEX IF NOT EXISTS + try { + await this.#connection.raw('CREATE UNIQUE INDEX ?? ON ?? (??, ??)', [ + indexName, + tableName, + 'queue', + 'dedup_id', + ]) + } catch (error) { + if ((error as { code?: string }).code !== 'ER_DUP_KEYNAME') throw error + } } /** - * Partial unique index on (queue, dedup_id) for active dedup slots. - * Prevents two concurrent inserts with the same dedup_id from both succeeding. - * Only PG and SQLite support partial unique indexes; MySQL is skipped. + * Drops the indexes the unique index replaces: the non-unique one, and the + * partial unique one PostgreSQL and SQLite had on pending and delayed jobs. */ - async #createDedupActiveUniqueIndex(tableName: string): Promise { - const client = this.#connection.client.config.client - if (client !== 'pg' && client !== 'better-sqlite3' && client !== 'sqlite3') return + async #dropLegacyDedupIndexes(tableName: string): Promise { + if (this.#dialect() !== 'mysql') { + await this.#connection.raw('DROP INDEX IF EXISTS ??', [`${tableName}_queue_dedup_id_index`]) + await this.#connection.raw('DROP INDEX IF EXISTS ??', [`${tableName}_dedup_active_uidx`]) + return + } - const indexName = `${tableName}_dedup_active_uidx` - await this.#connection.raw( - `CREATE UNIQUE INDEX IF NOT EXISTS ?? ON ?? ("queue", "dedup_id") ` + - `WHERE "dedup_id" IS NOT NULL AND "status" IN ('pending', 'delayed')`, - [indexName, tableName] - ) + try { + await this.#connection.raw('DROP INDEX ?? ON ??', [ + `${tableName}_queue_dedup_id_index`, + tableName, + ]) + } catch (error) { + if ((error as { code?: string }).code !== 'ER_CANT_DROP_FIELD_OR_KEY') throw error + } } /** diff --git a/src/services/kysely_queue_schema.ts b/src/services/kysely_queue_schema.ts index d93b952..75e41b2 100644 --- a/src/services/kysely_queue_schema.ts +++ b/src/services/kysely_queue_schema.ts @@ -57,7 +57,7 @@ export class KyselyQueueSchemaService { .execute() await this.#createJobsIndexes(tableName) - await this.#createDedupActiveUniqueIndex(tableName) + await this.#createDedupUniqueIndex(tableName) } async addDedupColumns(tableName: string = 'queue_jobs'): Promise { @@ -84,21 +84,23 @@ export class KyselyQueueSchemaService { await this.#connection.schema.alterTable(tableName).addColumn('dedup_ttl', 'bigint').execute() } - const index = this.#connection.schema - .createIndex(`${tableName}_queue_dedup_idx`) - .on(tableName) - .columns(['queue', 'dedup_id']) + // The derived table is what allows MySQL to read the table it updates. + await sql` + update ${sql.table(tableName)} set dedup_id = null, dedup_at = null, dedup_ttl = null + where (id, queue) in ( + select id, queue from ( + select distinct older.id, older.queue + from ${sql.table(tableName)} as older + join ${sql.table(tableName)} as newer + on newer.queue = older.queue and newer.dedup_id = older.dedup_id + where newer.dedup_at > older.dedup_at + or (newer.dedup_at = older.dedup_at and newer.id > older.id) + ) as stale + ) + `.execute(this.#connection) - if (this.#dialect === 'mysql') { - try { - await index.execute() - } catch (error) { - if (!this.#isDuplicateIndexError(error)) throw error - } - } else { - await index.ifNotExists().execute() - } - await this.#createDedupActiveUniqueIndex(tableName) + await this.#createDedupUniqueIndex(tableName) + await this.#dropLegacyDedupIndexes(tableName) } async createSchedulesTable(tableName: string = 'queue_schedules'): Promise { @@ -506,7 +508,6 @@ export class KyselyQueueSchemaService { ['status_score', ['queue', 'status', 'score']], ['status_execute', ['queue', 'status', 'execute_at']], ['status_finished', ['queue', 'status', 'finished_at']], - ['queue_dedup', ['queue', 'dedup_id']], ] for (const [suffix, columns] of indexes) { @@ -518,17 +519,47 @@ export class KyselyQueueSchemaService { } } - async #createDedupActiveUniqueIndex(tableName: string): Promise { - if (this.#dialect === 'mysql') return - - await this.#connection.schema - .createIndex(`${tableName}_dedup_active_uidx`) - .ifNotExists() + /** + * Unique index on (queue, dedup_id): a single job owns a dedup slot, whatever + * its status. Jobs without dedup hold NULL, which is never a duplicate. + */ + async #createDedupUniqueIndex(tableName: string): Promise { + const index = this.#connection.schema + .createIndex(`${tableName}_dedup_uidx`) .unique() .on(tableName) .columns(['queue', 'dedup_id']) - .where(sql`dedup_id is not null and status in ('pending', 'delayed')`) - .execute() + + // MySQL has no CREATE INDEX IF NOT EXISTS + if (this.#dialect === 'mysql') { + try { + await index.execute() + } catch (error) { + if (!this.#isDuplicateIndexError(error)) throw error + } + } else { + await index.ifNotExists().execute() + } + } + + /** + * Drops the indexes the unique index replaces: the non-unique one, and the + * partial unique one PostgreSQL and SQLite had on pending and delayed jobs. + */ + async #dropLegacyDedupIndexes(tableName: string): Promise { + const legacyIndex = this.#connection.schema.dropIndex(`${tableName}_queue_dedup_idx`) + + if (this.#dialect === 'mysql') { + try { + await legacyIndex.on(tableName).execute() + } catch (error) { + if ((error as { code?: string }).code !== 'ER_CANT_DROP_FIELD_OR_KEY') throw error + } + return + } + + await legacyIndex.ifExists().execute() + await this.#connection.schema.dropIndex(`${tableName}_dedup_active_uidx`).ifExists().execute() } /** LONGTEXT on MySQL, whose TEXT stops at 64 KB; TEXT elsewhere. */ diff --git a/tests/_utils/register_driver_test_suite.ts b/tests/_utils/register_driver_test_suite.ts index 8bfdb4b..a803271 100644 --- a/tests/_utils/register_driver_test_suite.ts +++ b/tests/_utils/register_driver_test_suite.ts @@ -1,5 +1,5 @@ import { test as JapaTest } from '@japa/runner' -import type { AcquiredJob, Adapter, JobLease } from '../../src/contracts/adapter.js' +import type { AcquiredJob, Adapter, JobLease, PushResult } from '../../src/contracts/adapter.js' interface DriverTestSuiteOptions { test: typeof JapaTest @@ -12,7 +12,7 @@ interface DriverTestSuiteOptions { supportsConcurrency?: boolean /** * Whether concurrent dispatches have an atomic dedup constraint. - * MySQL has no partial unique indexes, matching the documented Knex limitation. + * Some in-memory adapters do not serialize concurrent deduplicated dispatches. * @default true */ supportsAtomicDedup?: boolean @@ -3089,6 +3089,86 @@ export function registerDriverTestSuite(options: DriverTestSuiteOptions) { }) if (options.supportsAtomicDedup !== false) { + test('dedup: concurrent dispatches stay deduplicated while jobs are acquired', async ({ + assert, + }) => { + const adapter = await options.createAdapter() + adapter.setWorkerId('worker-1') + + let dispatching = true + const acquiring = (async () => { + while (dispatching) await adapter.popFrom('concurrent-dedup-queue') + })() + + const results = await Promise.all( + Array.from({ length: 4 }, async (_, dispatcher) => { + const dispatches = [] + for (let index = 0; index < 500; index++) { + dispatches.push( + await adapter.pushOn('concurrent-dedup-queue', { + id: `${dispatcher}-${index}`, + name: 'TestJob', + payload: {}, + attempts: 0, + dedup: { id: String(index % 200), ttl: 60_000 }, + }) + ) + } + return dispatches + }) + ).finally(() => (dispatching = false)) + await acquiring + + const dispatches = results.flat().map((result) => result as PushResult) + assert.lengthOf( + dispatches.filter(({ outcome }) => outcome === 'added'), + 200 + ) + + // Every dispatch points at a job that was stored + for (const jobId of new Set(dispatches.map((result) => result.jobId))) { + assert.isNotNull(await adapter.getJob(jobId, 'concurrent-dedup-queue')) + } + }).timeout(20_000) + + test('dedup: concurrent dispatches do not fail while jobs are completed', async ({ + assert, + }) => { + const adapter = await options.createAdapter() + adapter.setWorkerId('worker-1') + + let dispatching = true + let completed = 0 + const workers = Array.from({ length: 2 }, async () => { + while (dispatching || (await adapter.sizeOf('concurrent-dedup-queue')) > 0) { + const job = await adapter.popFrom('concurrent-dedup-queue') + if (job && (await adapter.completeJob(job, 'concurrent-dedup-queue'))) completed++ + } + }) + + const results = await Promise.all( + Array.from({ length: 4 }, async (_, dispatcher) => { + const dispatches = [] + for (let index = 0; index < 200; index++) { + dispatches.push( + await adapter.pushOn('concurrent-dedup-queue', { + id: `${dispatcher}-${index}`, + name: 'TestJob', + payload: {}, + attempts: 0, + dedup: { id: String(index % 3) }, + }) + ) + } + return dispatches + }) + ).finally(() => (dispatching = false)) + await Promise.all(workers) + + const added = results.flat().filter((result) => (result as PushResult).outcome === 'added') + assert.equal(completed, added.length) + }).timeout(20_000) + test('dedup: concurrent pushOn with same id - only one wins, rest skipped', async ({ assert, }) => { diff --git a/tests/adapter.spec.ts b/tests/adapter.spec.ts index 8b27877..e72f37d 100644 --- a/tests/adapter.spec.ts +++ b/tests/adapter.spec.ts @@ -3,6 +3,7 @@ import { tmpdir } from 'node:os' import { join } from 'node:path' import Knex from 'knex' import { test } from '@japa/runner' +import type { Assert } from '@japa/assert' import { Redis } from 'ioredis' import { CronExpressionParser } from 'cron-parser' import { MemoryAdapter } from './_mocks/memory_adapter.js' @@ -1970,6 +1971,45 @@ test.group('Adapter | Redis', (group) => { }) }) +/** + * A jobs table as the previous versions left it: no unique index, and expired + * owners that still hold their dedup id. + */ +async function assertDedupMigration( + assert: Assert, + connection: ReturnType, + tableName: string +) { + const schemaService = new KnexQueueSchemaService(connection) + const job = (id: string, dedupId: string, dedupAt: number) => ({ + id, + queue: 'default', + status: 'completed', + data: '{}', + dedup_id: dedupId, + dedup_at: dedupAt, + dedup_ttl: 10, + }) + + await connection.raw('DROP INDEX ??', [`${tableName}_dedup_uidx`]) + await connection.schema.alterTable(tableName, (table) => table.index(['queue', 'dedup_id'])) + await connection(tableName).insert([ + job('older', 'shared', 1_000), + job('newer', 'shared', 2_000), + job('alone', 'alone', 1_000), + ]) + + await schemaService.addDedupColumns(tableName) + await schemaService.addDedupColumns(tableName) + + assert.deepEqual(await connection(tableName).select('id', 'dedup_id').orderBy('id'), [ + { id: 'alone', dedup_id: 'alone' }, + { id: 'newer', dedup_id: 'shared' }, + { id: 'older', dedup_id: null }, + ]) + await assert.rejects(() => connection(tableName).insert(job('duplicate', 'shared', 3_000))) +} + test.group('Adapter | Knex (SQLite)', (group) => { let connection: ReturnType let adapter: KnexAdapter @@ -2003,6 +2043,10 @@ test.group('Adapter | Knex (SQLite)', (group) => { }, }) + test('addDedupColumns should leave a dedup slot to its latest job', async ({ assert }) => { + await assertDedupMigration(assert, connection, 'queue_jobs') + }) + test('listSchedules should execute a single SQL query in Knex adapter', async ({ assert }) => { const knexAdapter = new KnexAdapter({ connection }) @@ -2110,6 +2154,10 @@ test.group('Adapter | Knex (PostgreSQL)', (group) => { }, }) + test('addDedupColumns should leave a dedup slot to its latest job', async ({ assert }) => { + await assertDedupMigration(assert, connection, tableName) + }) + test('listSchedules should execute a single SQL query in Knex PostgreSQL adapter', async ({ assert, }) => { diff --git a/tests/kysely_adapter.spec.ts b/tests/kysely_adapter.spec.ts index 51f827b..df1674f 100644 --- a/tests/kysely_adapter.spec.ts +++ b/tests/kysely_adapter.spec.ts @@ -10,6 +10,7 @@ import { KyselyAdapter, KyselyQueueSchemaService, type QueueDatabase, + type QueueJobTable, } from '../src/drivers/kysely_adapter.js' import { registerDriverTestSuite } from './_utils/register_driver_test_suite.js' @@ -174,7 +175,6 @@ test.group('Adapter | Kysely (MySQL)', (group) => { registerDriverTestSuite({ test, - supportsAtomicDedup: false, createAdapter: () => { adapter = new KyselyAdapter({ connection, @@ -190,6 +190,49 @@ test.group('Adapter | Kysely (MySQL)', (group) => { await schema.addDedupColumns(tableName) await schema.addDedupColumns(tableName) }) + + test('addDedupColumns should leave a dedup slot to its latest job', async ({ assert }) => { + const jobs = connection.withTables>() + const job = (id: string, dedupId: string, dedupAt: number) => ({ + id, + queue: 'default', + status: 'completed' as const, + data: '{}', + dedup_id: dedupId, + dedup_at: dedupAt, + dedup_ttl: 10, + }) + + // A jobs table as the previous versions left it: no unique index, and + // expired owners that still hold their dedup id. + await connection.schema.dropIndex(`${tableName}_dedup_uidx`).on(tableName).execute() + await connection.schema + .createIndex(`${tableName}_queue_dedup_idx`) + .on(tableName) + .columns(['queue', 'dedup_id']) + .execute() + await jobs + .insertInto(tableName) + .values([job('older', 'shared', 1_000), job('newer', 'shared', 2_000)]) + .execute() + + await schema.addDedupColumns(tableName) + await schema.addDedupColumns(tableName) + + assert.deepEqual( + await jobs.selectFrom(tableName).select(['id', 'dedup_id']).orderBy('id').execute(), + [ + { id: 'newer', dedup_id: 'shared' }, + { id: 'older', dedup_id: null }, + ] + ) + await assert.rejects(() => + jobs + .insertInto(tableName) + .values(job('duplicate', 'shared', 3_000)) + .execute() + ) + }) }) test.group('KyselyAdapter | SQLite file, two connections', () => {