diff --git a/service/src/job-cancellation.test.ts b/service/src/job-cancellation.test.ts index e23186e4..334fb8f7 100644 --- a/service/src/job-cancellation.test.ts +++ b/service/src/job-cancellation.test.ts @@ -1,5 +1,6 @@ import { expect, test } from 'bun:test'; import { EventEmitter } from 'node:events'; +import calculateSlot from 'cluster-key-slot'; import type IORedis from 'ioredis'; import type { Job, QueueEvents } from 'bullmq'; import { @@ -62,7 +63,7 @@ class FakeRedis { readonly existing = new Set(); readonly deleted: string[] = []; readonly transactions: FakeTransaction[] = []; - mgetFailures = 0; + getFailures = 0; cancellationFailures = 0; cancellationAttempts = 0; @@ -72,6 +73,10 @@ class FakeRedis { } async get(key: string): Promise { + if (this.getFailures > 0) { + this.getFailures -= 1; + throw new Error('command connection unavailable'); + } return this.existing.has(key) ? '1' : null; } @@ -93,14 +98,6 @@ class FakeRedis { return [1]; } - async mget(...keys: string[]): Promise> { - if (this.mgetFailures > 0) { - this.mgetFailures -= 1; - throw new Error('command connection unavailable'); - } - return keys.map(key => (this.existing.has(key) ? '1' : null)); - } - async del(key: string): Promise { this.deleted.push(key); this.existing.delete(key); @@ -263,7 +260,7 @@ test('subscriber reconnect retries durable-marker reconciliation', async () => { const controller = new AbortController(); await registry.register(target, controller); fake.existing.add(jobCancellationInternals.cancellationKey(target)); - fake.mgetFailures = 1; + fake.getFailures = 1; fake.subscriber.emit('ready'); await new Promise(resolve => setTimeout(resolve, 150)); @@ -597,3 +594,10 @@ test('queued removal frees waiting capacity and tolerates an activation race', a ).toBe(false); expect(removals).toBe(2); }); + +test('keys for one job share a Redis Cluster hash slot', () => { + const key = jobCancellationInternals.cancellationKey({ queueName: 'stateful-other-queue', jobId: '3q2X3EAPwWyx4u_oFIDG8' }); + const slot = calculateSlot(key); + expect(calculateSlot(`${key}:result`)).toBe(slot); + expect(calculateSlot(`${key}:execution`)).toBe(slot); +}); diff --git a/service/src/job-cancellation.ts b/service/src/job-cancellation.ts index 4da3253d..772d5915 100644 --- a/service/src/job-cancellation.ts +++ b/service/src/job-cancellation.ts @@ -1,5 +1,6 @@ import type IORedis from 'ioredis'; import type { Job, QueueEvents } from 'bullmq'; +import { hashTag } from './redis-connection'; const JOB_CANCELLATION_PREFIX = 'codeapi:job-cancellation:v1'; const JOB_CANCELLATION_CHANNEL = `${JOB_CANCELLATION_PREFIX}:events`; @@ -15,10 +16,12 @@ function targetKey(target: JobTarget): string { return `${target.queueName}:${target.jobId}`; } +/** Hash-tagged so each job's marker, result and execution keys share one + * Redis Cluster slot; the multi-key scripts below fail with CROSSSLOT otherwise. */ function cancellationKey(target: JobTarget): string { - return `${JOB_CANCELLATION_PREFIX}:${encodeURIComponent( - target.queueName, - )}:${encodeURIComponent(target.jobId)}`; + return `${JOB_CANCELLATION_PREFIX}:${hashTag( + `${encodeURIComponent(target.queueName)}:${encodeURIComponent(target.jobId)}`, + )}`; } function parseTarget(raw: string): JobTarget | undefined { @@ -137,8 +140,9 @@ export class JobCancellationRegistry { private async reconcile(): Promise { const entries = [...this.controllers.values()]; if (entries.length === 0) return; - const cancelled = await this.commands.mget( - ...entries.map(({ target }) => cancellationKey(target)), + // One GET per job: different jobs live in different cluster slots. + const cancelled = await Promise.all( + entries.map(({ target }) => this.commands.get(cancellationKey(target))), ); cancelled.forEach((value, index) => { if (value === '1') {