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
4 changes: 2 additions & 2 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,7 @@
"c8": "^11.0.0",
"del-cli": "^7.0.0",
"hot-hook": "^1.0.0",
"ioredis": "^5.11.1",

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I don't know if you want your test suite to test both versions or not (I don't think it's worth it), but in a "ideal" world/test coverage scenario, if we want to support both versions, we might want to have e2e tests that cover both.

"ioredis": "^6.0.0",
"knex": "3.1.0",
"kysely": "^0.29.2",
"mysql2": "^3.0.0",
Expand All @@ -76,7 +76,7 @@
"@opentelemetry/api": "^1.9.0",
"@opentelemetry/core": "^1.30.0 || ^2.0.0",
"@opentelemetry/instrumentation": ">=0.200.0 <0.300.0",
"ioredis": "^5.0.0",
"ioredis": "^5.0.0 || ^6.0.0",

@theoludwig theoludwig Aug 17, 2026 •

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should we support both or the breaking change is fine?

Suggested change
"ioredis": "^5.0.0 || ^6.0.0",
"ioredis": "^6.0.0",

"knex": "^3.0.0",
"kysely": "^0.29.0"
},
Expand Down
69 changes: 39 additions & 30 deletions tests/_utils/register_driver_test_suite.ts
Original file line number Diff line number Diff line change
Expand Up @@ -831,7 +831,19 @@ export function registerDriverTestSuite(options: DriverTestSuiteOptions) {

test('renewJobs should keep an active job from being recovered as stalled', async ({
assert,
cleanup,
}) => {
const stalledThreshold = 60_000
const realNow = Date.now
let clockOffset = 0
Date.now = () => realNow() + clockOffset
cleanup(() => {
Date.now = realNow
})
const advancePastStalledThreshold = () => {
clockOffset += stalledThreshold + 1
}

const adapter = await options.createAdapter()
adapter.setWorkerId('worker-1')

Expand All @@ -844,25 +856,22 @@ export function registerDriverTestSuite(options: DriverTestSuiteOptions) {
const job = acquired.find((candidate) => candidate?.id === 'long-running')
assert.isDefined(job)

// Both jobs run longer than the 200ms stalled threshold; only one is renewed.
// A slow run only makes the control job more stalled, and the renewed job
// has a 200ms margin between each renewal and the next recovery.
await new Promise((resolve) => setTimeout(resolve, 300))
advancePastStalledThreshold()
assert.equal(await adapter.renewJobs('test-queue', [job!]), 1)

const first = await adapter.recoverStalledJobs('test-queue', 200, 1, 100)
const first = await adapter.recoverStalledJobs('test-queue', stalledThreshold, 1, 100)

// The control job proves the threshold had passed; the renewed job stays active.
assert.equal(first.recovered, 1)
assert.equal((await adapter.getJob('long-running', 'test-queue'))!.status, 'active')
assert.equal((await adapter.getJob('not-renewed', 'test-queue'))!.status, 'pending')

// 300ms after the first renewal, the job would be stalled again: only the
// second renewal, with the same lease, keeps it active.
await new Promise((resolve) => setTimeout(resolve, 300))
// Past the threshold since the first renewal, only the second renewal,
// with the same lease, keeps the job active.
advancePastStalledThreshold()
assert.equal(await adapter.renewJobs('test-queue', [job!]), 1)

const second = await adapter.recoverStalledJobs('test-queue', 200, 1, 100)
const second = await adapter.recoverStalledJobs('test-queue', stalledThreshold, 1, 100)
assert.equal(second.recovered, 0)
assert.equal((await adapter.getJob('long-running', 'test-queue'))!.status, 'active')
})
Expand Down Expand Up @@ -2690,30 +2699,30 @@ export function registerDriverTestSuite(options: DriverTestSuiteOptions) {
name: 'TestJob',
payload: { n: 1 },
attempts: 0,
dedup: { id: 'TestJob::ext-1', ttl: 400, extend: true },
dedup: { id: 'TestJob::ext-1', ttl: 1000, extend: true },
})

await new Promise((r) => setTimeout(r, 250))
await new Promise((r) => setTimeout(r, 550))

const second = await adapter.pushOn('ext-queue', {
id: 'ext-uuid-2',
name: 'TestJob',
payload: { n: 2 },
attempts: 0,
dedup: { id: 'TestJob::ext-1', ttl: 400, extend: true },
dedup: { id: 'TestJob::ext-1', ttl: 1000, extend: true },
})
assert.equal(second && typeof second === 'object' && second.outcome, 'extended')

await new Promise((r) => setTimeout(r, 250))
await new Promise((r) => setTimeout(r, 550))

// Without the extend at T=250, the window would have expired at T=400. With it,
// only 250ms of the new 400ms window have passed.
// Without the extend at T=550, the window would have expired at T=1000. With it,
// only 550ms of the new 1000ms window have passed.
const third = await adapter.pushOn('ext-queue', {
id: 'ext-uuid-3',
name: 'TestJob',
payload: { n: 3 },
attempts: 0,
dedup: { id: 'TestJob::ext-1', ttl: 400, extend: true },
dedup: { id: 'TestJob::ext-1', ttl: 1000, extend: true },
})
assert.equal(third && typeof third === 'object' && third.outcome, 'extended')
})
Expand Down Expand Up @@ -2972,34 +2981,34 @@ export function registerDriverTestSuite(options: DriverTestSuiteOptions) {
name: 'TestJob',
payload: { version: 1 },
attempts: 0,
dedup: { id: 'TestJob::debounce-1', ttl: 400, extend: true, replace: true },
dedup: { id: 'TestJob::debounce-1', ttl: 1000, extend: true, replace: true },
})

await new Promise((r) => setTimeout(r, 250))
await new Promise((r) => setTimeout(r, 550))

const second = await adapter.pushOn('debounce-queue', {
id: 'debounce-uuid-2',
name: 'TestJob',
payload: { version: 2 },
attempts: 0,
dedup: { id: 'TestJob::debounce-1', ttl: 400, extend: true, replace: true },
dedup: { id: 'TestJob::debounce-1', ttl: 1000, extend: true, replace: true },
})
assert.equal(second && typeof second === 'object' && second.outcome, 'replaced')
assert.equal(second && typeof second === 'object' && second.jobId, 'debounce-uuid-1')

const midRecord = await adapter.getJob('debounce-uuid-1', 'debounce-queue')
assert.deepEqual(midRecord!.data.payload, { version: 2 })

// 500ms total elapsed > original 400ms TTL, but the second dispatch reset
// the window at T=250. Only 250ms into the new window → still alive.
await new Promise((r) => setTimeout(r, 250))
// 1100ms total elapsed > original 1000ms TTL, but the second dispatch reset
// the window at T=550. Only 550ms into the new window → still alive.
await new Promise((r) => setTimeout(r, 550))

const third = await adapter.pushOn('debounce-queue', {
id: 'debounce-uuid-3',
name: 'TestJob',
payload: { version: 3 },
attempts: 0,
dedup: { id: 'TestJob::debounce-1', ttl: 400, extend: true, replace: true },
dedup: { id: 'TestJob::debounce-1', ttl: 1000, extend: true, replace: true },
})
assert.equal(third && typeof third === 'object' && third.outcome, 'replaced')
assert.equal(third && typeof third === 'object' && third.jobId, 'debounce-uuid-1')
Expand Down Expand Up @@ -3052,11 +3061,11 @@ export function registerDriverTestSuite(options: DriverTestSuiteOptions) {
name: 'TestJob',
payload: { n: 1 },
attempts: 0,
dedup: { id: 'TestJob::active-ext-1', ttl: 400, extend: true },
dedup: { id: 'TestJob::active-ext-1', ttl: 1000, extend: true },
})

// Move to active mid-window.
await new Promise((r) => setTimeout(r, 250))
await new Promise((r) => setTimeout(r, 550))
const popped = await adapter.popFrom('active-ext-queue')
assert.equal(popped!.id, 'active-ext-uuid-1')

Expand All @@ -3067,22 +3076,22 @@ export function registerDriverTestSuite(options: DriverTestSuiteOptions) {
name: 'TestJob',
payload: { n: 2 },
attempts: 0,
dedup: { id: 'TestJob::active-ext-1', ttl: 400, extend: true },
dedup: { id: 'TestJob::active-ext-1', ttl: 1000, extend: true },
})
assert.equal(second && typeof second === 'object' && second.outcome, 'extended')
assert.equal(second && typeof second === 'object' && second.jobId, 'active-ext-uuid-1')

// Without the extend, the slot would have expired by now (250 + 250 > 400).
// With the extend at T=250, the window restarted; at T=500 only 250ms into
// Without the extend, the slot would have expired by now (550 + 550 > 1000).
// With the extend at T=550, the window restarted; at T=1100 only 550ms into
// new window → still blocking.
await new Promise((r) => setTimeout(r, 250))
await new Promise((r) => setTimeout(r, 550))

const third = await adapter.pushOn('active-ext-queue', {
id: 'active-ext-uuid-3',
name: 'TestJob',
payload: { n: 3 },
attempts: 0,
dedup: { id: 'TestJob::active-ext-1', ttl: 400, extend: true },
dedup: { id: 'TestJob::active-ext-1', ttl: 1000, extend: true },
})
assert.equal(third && typeof third === 'object' && third.outcome, 'extended')
assert.equal(third && typeof third === 'object' && third.jobId, 'active-ext-uuid-1')
Expand Down
36 changes: 15 additions & 21 deletions tests/_utils/register_worker_concurrency_suite.ts
Original file line number Diff line number Diff line change
Expand Up @@ -200,12 +200,15 @@ export function registerWorkerConcurrencyTestSuite(options: WorkerConcurrencyTes
}) => {
const jobStartTimes: Map<string, number> = new Map()
const jobEndTimes: Map<string, number> = new Map()
const jobCount = 4
const { promise: allJobsStarted, resolve: releaseJobs } = Promise.withResolvers<void>()
const sequentialExecutionTimeout = setTimeout(2000, undefined, { ref: false })

class SlowJob extends Job<{ jobId: string }> {
async execute() {
jobStartTimes.set(this.payload.jobId, Date.now())
// Simulate a slow job (300ms)
await setTimeout(300)
if (jobStartTimes.size === jobCount) releaseJobs()
await Promise.race([allJobsStarted, sequentialExecutionTimeout])
jobEndTimes.set(this.payload.jobId, Date.now())
}
}
Expand Down Expand Up @@ -237,7 +240,7 @@ export function registerWorkerConcurrencyTestSuite(options: WorkerConcurrencyTes
while (cycles < maxCycles) {
const cycle = await worker.processCycle(['default'])
cycles++
if (cycle?.type === 'idle' && jobEndTimes.size === 4) break
if (cycle?.type === 'idle' && jobEndTimes.size === jobCount) break
}
})()

Expand Down Expand Up @@ -279,41 +282,32 @@ export function registerWorkerConcurrencyTestSuite(options: WorkerConcurrencyTes
await processingPromise

// All 4 jobs should have been executed
assert.equal(jobStartTimes.size, 4, 'All 4 jobs should have started')
assert.equal(jobEndTimes.size, 4, 'All 4 jobs should have completed')
assert.equal(jobStartTimes.size, jobCount, 'All 4 jobs should have started')
assert.equal(jobEndTimes.size, jobCount, 'All 4 jobs should have completed')

// Verify concurrent execution: jobs 1, 2, 3 should start BEFORE job 0 ends
// If they ran sequentially, job 1 would start after job 0's 300ms execution
const job0Start = jobStartTimes.get('job-0')!
// If they ran sequentially, job 0 would only end after the sequential execution timeout
const job0End = jobEndTimes.get('job-0')!
const job1Start = jobStartTimes.get('job-1')!
const job2Start = jobStartTimes.get('job-2')!
const job3Start = jobStartTimes.get('job-3')!

// Job 1 should start before job 0 ends (proving concurrency)
assert.isTrue(
job1Start < job0End,
`Job 1 should start (${job1Start}) before job 0 ends (${job0End}) - concurrent execution`
job1Start <= job0End,
`Job 1 should start (${job1Start}) no later than job 0 ends (${job0End}) - concurrent execution`
)

// Job 2 should start before job 0 ends
assert.isTrue(
job2Start < job0End,
`Job 2 should start (${job2Start}) before job 0 ends (${job0End}) - concurrent execution`
job2Start <= job0End,
`Job 2 should start (${job2Start}) no later than job 0 ends (${job0End}) - concurrent execution`
)

// Job 3 should start before job 0 ends
assert.isTrue(
job3Start < job0End,
`Job 3 should start (${job3Start}) before job 0 ends (${job0End}) - concurrent execution`
)

// All jobs should start within a reasonable time window (not sequentially)
const maxStartDiff = Math.max(job1Start, job2Start, job3Start) - job0Start
assert.isBelow(
maxStartDiff,
250,
`All jobs should start within 250ms of each other (actual: ${maxStartDiff}ms)`
job3Start <= job0End,
`Job 3 should start (${job3Start}) no later than job 0 ends (${job0End}) - concurrent execution`
)
})
}
4 changes: 2 additions & 2 deletions tests/adapter.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -976,10 +976,10 @@ test.group('Adapter | Redis', (group) => {
.exec()

await adapter.migrate()
const firstMembers = await connection.zrange('schedules::due', 0, -1, 'WITHSCORES')
const firstMembers = await connection.zrange('schedules::due', 0, '-1', 'WITHSCORES')

await adapter.migrate()
const secondMembers = await connection.zrange('schedules::due', 0, -1, 'WITHSCORES')
const secondMembers = await connection.zrange('schedules::due', 0, '-1', 'WITHSCORES')

assert.deepEqual(firstMembers, ['idempotent-schedule', nextRunAt.toString()])
assert.deepEqual(secondMembers, firstMembers)
Expand Down
31 changes: 15 additions & 16 deletions yarn.lock
Original file line number Diff line number Diff line change
Expand Up @@ -64,7 +64,7 @@ __metadata:
cron-parser: "npm:^5.6.1"
del-cli: "npm:^7.0.0"
hot-hook: "npm:^1.0.0"
ioredis: "npm:^5.11.1"
ioredis: "npm:^6.0.0"
knex: "npm:3.1.0"
kysely: "npm:^0.29.2"
mysql2: "npm:^3.0.0"
Expand All @@ -79,7 +79,7 @@ __metadata:
"@opentelemetry/api": ^1.9.0
"@opentelemetry/core": ^1.30.0 || ^2.0.0
"@opentelemetry/instrumentation": ">=0.200.0 <0.300.0"
ioredis: ^5.0.0
ioredis: ^5.0.0 || ^6.0.0
knex: ^3.0.0
kysely: ^0.29.0
peerDependenciesMeta:
Expand Down Expand Up @@ -547,20 +547,20 @@ __metadata:
languageName: node
linkType: hard

"@ioredis/commands@npm:1.10.0":
version: 1.10.0
resolution: "@ioredis/commands@npm:1.10.0"
checksum: 10c0/baf91e62d0e64ef2b5f7ca4413dc2456fe250e87483beac4a1c8ef1fe5ad0d2fcdeb9b89d4556d8ef6c7455c64a964359d729601fdb06b2f4c76c35dd59afa99
languageName: node
linkType: hard

"@ioredis/commands@npm:1.5.1":
version: 1.5.1
resolution: "@ioredis/commands@npm:1.5.1"
checksum: 10c0/cb8f6d13cff0753e3e7ef001fb895491985d9a623248192538f13bc2fd9bfdfde3c18cf2ba6f20ec8ceaa681b0771070d3a09b82eed044c798bcfef5e3ae54b3
languageName: node
linkType: hard

"@ioredis/commands@npm:2.0.0":
version: 2.0.0
resolution: "@ioredis/commands@npm:2.0.0"
checksum: 10c0/2fb5edb9782790c24a375cc78bd69d45ac6c268d7d7ed8e534f2d4ecc28707bde9882a04eecd68be369ae784801bf93fa267b9a76f86418b252a08f341934e35
languageName: node
linkType: hard

"@isaacs/fs-minipass@npm:^4.0.0":
version: 4.0.1
resolution: "@isaacs/fs-minipass@npm:4.0.1"
Expand Down Expand Up @@ -3414,18 +3414,17 @@ __metadata:
languageName: node
linkType: hard

"ioredis@npm:^5.11.1":
version: 5.11.1
resolution: "ioredis@npm:5.11.1"
"ioredis@npm:^6.0.0":
version: 6.0.0
resolution: "ioredis@npm:6.0.0"
dependencies:
"@ioredis/commands": "npm:1.10.0"
"@ioredis/commands": "npm:2.0.0"
cluster-key-slot: "npm:1.1.1"
debug: "npm:4.4.3"
denque: "npm:2.1.0"
redis-errors: "npm:1.2.0"
redis-parser: "npm:3.0.0"
standard-as-callback: "npm:2.1.0"
checksum: 10c0/a8b27043cf2c045dfc93f40a32ce24cf9f8b57799a37f4234c4b925c365ccf131629590f94a512f546fda2ba8ed034009c94c4933ecd44c50bc166636d929fd6
checksum: 10c0/6f74e1d298440c0d814e6bf9b3acdb0991d4c1862d3fa64c3828668bcc306823f7845710abdc5eab7986a35d86ec8c86c4bb9c9c7d727636d385ecb2cd40abd7
languageName: node
linkType: hard

Expand Down Expand Up @@ -5021,7 +5020,7 @@ __metadata:
languageName: node
linkType: hard

"redis-parser@npm:3.0.0, redis-parser@npm:^3.0.0":
"redis-parser@npm:^3.0.0":
version: 3.0.0
resolution: "redis-parser@npm:3.0.0"
dependencies:
Expand Down
Loading