Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 18 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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

Expand Down Expand Up @@ -473,6 +472,23 @@ await schema.dropJobsTable()

</details>

#### 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
Expand Down
108 changes: 49 additions & 59 deletions src/drivers/knex_adapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, unknown>,
dedup: NonNullable<JobData['dedup']>
): Promise<PushResult> {
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(
Expand All @@ -525,14 +545,20 @@ export class KnexAdapter implements Adapter {
dedup: NonNullable<JobData['dedup']>,
now: number
): Promise<PushResult | null> {
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
Expand Down Expand Up @@ -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<string, unknown>,
dedup: NonNullable<JobData['dedup']>,
now: number
): Promise<PushResult> {
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<void> {
return this.pushManyOn('default', jobs)
}
Expand Down
139 changes: 66 additions & 73 deletions src/drivers/kysely_adapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -465,29 +465,49 @@ export class KyselyAdapter<DB = QueueDatabase> 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<DB>,
queue: string,
jobData: JobData,
insertRow: Partial<JobRow> & Pick<JobRow, 'id' | 'queue' | 'status' | 'data'>
): Promise<PushResult> {
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
}
}

Expand All @@ -498,18 +518,31 @@ export class KyselyAdapter<DB = QueueDatabase> implements Adapter {
dedup: NonNullable<JobData['dedup']>,
now: number
): Promise<PushResult | null> {
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)
Expand Down Expand Up @@ -547,62 +580,15 @@ export class KyselyAdapter<DB = QueueDatabase> 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<DB>,
queue: string,
jobId: string,
insertRow: Partial<JobRow> & Pick<JobRow, 'id' | 'queue' | 'status' | 'data'>,
dedup: NonNullable<JobData['dedup']>,
now: number
): Promise<PushResult> {
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<void> {
Expand Down Expand Up @@ -962,6 +948,13 @@ export class KyselyAdapter<DB = QueueDatabase> 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 ?? ''
Expand Down
Loading
Loading