diff --git a/package.json b/package.json index cbac5e5..147c01d 100644 --- a/package.json +++ b/package.json @@ -60,7 +60,7 @@ "c8": "^11.0.0", "del-cli": "^7.0.0", "hot-hook": "^1.0.0", - "ioredis": "^5.11.1", + "ioredis": "^6.0.0", "knex": "3.1.0", "kysely": "^0.29.2", "mysql2": "^3.0.0", @@ -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", "knex": "^3.0.0", "kysely": "^0.29.0" }, diff --git a/tests/_utils/register_driver_test_suite.ts b/tests/_utils/register_driver_test_suite.ts index 8bfdb4b..730f22a 100644 --- a/tests/_utils/register_driver_test_suite.ts +++ b/tests/_utils/register_driver_test_suite.ts @@ -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') @@ -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') }) @@ -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') }) @@ -2972,17 +2981,17 @@ 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') @@ -2990,16 +2999,16 @@ export function registerDriverTestSuite(options: DriverTestSuiteOptions) { 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') @@ -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') @@ -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') diff --git a/tests/_utils/register_worker_concurrency_suite.ts b/tests/_utils/register_worker_concurrency_suite.ts index 56d92ff..9cd099f 100644 --- a/tests/_utils/register_worker_concurrency_suite.ts +++ b/tests/_utils/register_worker_concurrency_suite.ts @@ -200,12 +200,15 @@ export function registerWorkerConcurrencyTestSuite(options: WorkerConcurrencyTes }) => { const jobStartTimes: Map = new Map() const jobEndTimes: Map = new Map() + const jobCount = 4 + const { promise: allJobsStarted, resolve: releaseJobs } = Promise.withResolvers() + 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()) } } @@ -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 } })() @@ -279,12 +282,11 @@ 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')! @@ -292,28 +294,20 @@ export function registerWorkerConcurrencyTestSuite(options: WorkerConcurrencyTes // 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` ) }) } diff --git a/tests/adapter.spec.ts b/tests/adapter.spec.ts index 8b27877..3b42b28 100644 --- a/tests/adapter.spec.ts +++ b/tests/adapter.spec.ts @@ -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) diff --git a/yarn.lock b/yarn.lock index 587646a..5e6ffc1 100644 --- a/yarn.lock +++ b/yarn.lock @@ -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" @@ -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: @@ -547,13 +547,6 @@ __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" @@ -561,6 +554,13 @@ __metadata: 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" @@ -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 @@ -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: