Skip to content

Commit f3d54ae

Browse files
buslaclaude
andcommitted
feat(redis): cluster-safe keys for the api, worker and file-server paths
ElastiCache Serverless (Valkey) is always cluster mode, so every multi-key script, MULTI/EXEC and MGET must stay in one hash slot. Follows the key shapes of upstream PR LibreChat-AI#17 so this can be dropped when that lands: - redis-connection.ts: isClusterMode (USE_REDIS_CLUSTER), bullmqPrefix ('{codeapi}' in cluster mode, BullMQ's default otherwise), hashTag. - runtime-session registry: rtsx:{sess,lock,gen,ckptseq}:{<id>} so the record+lock scripts are single-slot (not covered by PR LibreChat-AI#17). - replay state: exec_state/exec_lock/exec_result/tool_history:{<id>}; the stale-execution sweep reads with pipelined GETs instead of a cross-slot MGET. - BullMQ prefix on every Queue/Worker/QueueEvents, including hosted apps. Bridge, egress ledger and tool-call server are not covered (not used by this deployment). Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Wzy53QgK2dCkSKvEFTySD1
1 parent 84be06a commit f3d54ae

9 files changed

Lines changed: 85 additions & 45 deletions

File tree

‎service/src/hosted-app/queue.ts‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ import { hostedAppOperationTimeoutMs } from '../config';
66
import logger from '../logger';
77
import { withHostedAppQueueDeadline } from './queue-deadline';
88
import type { HostedAppJobData, HostedAppJobName, HostedAppJobResult } from './jobs';
9+
import { bullmqPrefix } from '../redis-connection';
910

1011
/** Hosted apps are stateful-only. A fixed isolated queue prevents an ordinary
1112
* stateless worker from ever receiving an AWS lifecycle job. */
@@ -24,9 +25,9 @@ function resources(): { queue: HostedAppQueue; events: QueueEvents } {
2425
if (!hostedAppQueue || !hostedAppQueueEvents) {
2526
hostedAppQueue = new Queue<HostedAppJobData, HostedAppJobResult, HostedAppJobName>(
2627
HOSTED_APP_QUEUE_NAME,
27-
{ connection },
28+
{ prefix: bullmqPrefix(), connection },
2829
);
29-
hostedAppQueueEvents = new QueueEvents(HOSTED_APP_QUEUE_NAME, { connection });
30+
hostedAppQueueEvents = new QueueEvents(HOSTED_APP_QUEUE_NAME, { prefix: bullmqPrefix(), connection });
3031
setMaxListeners(0, hostedAppQueue, hostedAppQueueEvents);
3132
/* These resources are created after lifecycle startup, on first use, so the
3233
* ordinary queue listener registration never sees them. An unhandled

‎service/src/hosted-app/worker.ts‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ import type { HostedAppJob, HostedAppJobData, HostedAppJobName, HostedAppJobResu
88
import { HOSTED_APP_QUEUE_NAME } from './queue';
99
import logger from '../logger';
1010
import { workerRunning } from '../metrics';
11+
import { bullmqPrefix } from '../redis-connection';
1112

1213
export function serializedHostedAppFailure(error: unknown): Error {
1314
if (error instanceof HostedAppControlPlaneError) {
@@ -101,7 +102,7 @@ export const hostedAppWorker: Worker<
101102
HostedAppJobName
102103
> | undefined = env.HOSTED_APPS_ENABLED
103104
? new Worker(HOSTED_APP_QUEUE_NAME, processHostedAppJob, {
104-
connection,
105+
prefix: bullmqPrefix(), connection,
105106
/* Lifecycle transitions are serialized again by their per-app Redis lock.
106107
* This modest concurrency allows unrelated apps to launch in parallel while
107108
* the fleet-wide AWS throttle remains authoritative. */

‎service/src/queue.ts‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ import type {
1919
import logger from './logger';
2020
import { redisKeepAliveOptions } from './redis-options';
2121
import { bullmqQueueJobs, registerBullmqQueueMetricsCollector } from './metrics';
22+
import { bullmqPrefix } from './redis-connection';
2223

2324
const MAX_RECONNECT_ATTEMPTS = 5;
2425
const RECONNECT_DELAY = 2000;
@@ -87,8 +88,8 @@ function getQueueResources(
8788
const existing = queueResources.get(name);
8889
if (existing != null) return existing;
8990

90-
const queue = new Queue<t.JobData, t.JobResult, Jobs.execute>(name, { connection });
91-
const events = new QueueEvents(name, { connection });
91+
const queue = new Queue<t.JobData, t.JobResult, Jobs.execute>(name, { prefix: bullmqPrefix(), connection });
92+
const events = new QueueEvents(name, { prefix: bullmqPrefix(), connection });
9293
setMaxListeners(0, queue, events);
9394
const resources = { queue, events };
9495
queueResources.set(name, resources);

‎service/src/redis-connection.ts‎

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,31 @@
1+
/** Minimal Redis Cluster key helpers. Same names and key shapes as upstream
2+
* PR #17 (LibreChat-AI/code-interpreter, "Add Redis Cluster Support") so this
3+
* file can be dropped when that lands. ElastiCache Serverless is always
4+
* cluster mode: every multi-key script, MULTI/EXEC or MGET must stay in one
5+
* hash slot.
6+
*
7+
* ponytail: covers what the api, worker and file-server use (session
8+
* registry, replay state, BullMQ). bridge/*, egress-ledger and
9+
* tool-call-server still do cross-slot work; tag them before running those
10+
* components against a cluster. */
11+
12+
export function isClusterMode(): boolean {
13+
return process.env.USE_REDIS_CLUSTER === 'true' || (process.env.REDIS_HOST ?? '').includes(',');
14+
}
15+
16+
/** One hash tag keeps every BullMQ key in one slot on a cluster. Standalone
17+
* keeps BullMQ's default so existing queues are untouched. */
18+
export function bullmqPrefix(): string {
19+
return isClusterMode() ? '{codeapi}' : 'bull';
20+
}
21+
22+
/** Wrap an id in a hash tag so every key built from it lands in one slot.
23+
* Redis Cluster hashes only the substring between the first `{` and `}`. */
24+
export function hashTag(id: string): string {
25+
return `{${id}}`;
26+
}
27+
28+
/** Reverse of `hashTag`. */
29+
export function stripHashTag(raw: string): string {
30+
return raw.startsWith('{') && raw.endsWith('}') ? raw.slice(1, -1) : raw;
31+
}

‎service/src/runtime-session/registry.test.ts‎

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -181,10 +181,10 @@ describe('runtime session lock', () => {
181181
releaseDelayedSet();
182182
await lateSetComplete;
183183
for (let attempt = 0; attempt < 20; attempt += 1) {
184-
if (await mock.get('rtsx:lock:rt_late_lock') == null) break;
184+
if (await mock.get('rtsx:lock:{rt_late_lock}') == null) break;
185185
await new Promise(resolve => setTimeout(resolve, 10));
186186
}
187-
expect(await mock.get('rtsx:lock:rt_late_lock')).toBeNull();
187+
expect(await mock.get('rtsx:lock:{rt_late_lock}')).toBeNull();
188188
});
189189

190190
test('an applied SET whose replay resolves null cannot leave a ghost lock', async () => {
@@ -228,10 +228,10 @@ describe('runtime session lock', () => {
228228
await expect(acquire).rejects.toThrow('job deadline');
229229
finishReplay();
230230
for (let attempt = 0; attempt < 20; attempt += 1) {
231-
if (await mock.get('rtsx:lock:rt_ambiguous_lock') == null) break;
231+
if (await mock.get('rtsx:lock:{rt_ambiguous_lock}') == null) break;
232232
await new Promise(resolve => setTimeout(resolve, 10));
233233
}
234-
expect(await mock.get('rtsx:lock:rt_ambiguous_lock')).toBeNull();
234+
expect(await mock.get('rtsx:lock:{rt_ambiguous_lock}')).toBeNull();
235235
});
236236

237237
test('an applied SET whose replay quickly resolves null is still released', async () => {
@@ -254,10 +254,10 @@ describe('runtime session lock', () => {
254254

255255
expect(await acquireRuntimeSessionLock('rt_fast_null_lock', 60_000)).toBeNull();
256256
for (let attempt = 0; attempt < 20; attempt += 1) {
257-
if (await mock.get('rtsx:lock:rt_fast_null_lock') == null) break;
257+
if (await mock.get('rtsx:lock:{rt_fast_null_lock}') == null) break;
258258
await new Promise(resolve => setTimeout(resolve, 10));
259259
}
260-
expect(await mock.get('rtsx:lock:rt_fast_null_lock')).toBeNull();
260+
expect(await mock.get('rtsx:lock:{rt_fast_null_lock}')).toBeNull();
261261
});
262262

263263
test('an applied SET whose response quickly rejects is still released', async () => {
@@ -282,10 +282,10 @@ describe('runtime session lock', () => {
282282
acquireRuntimeSessionLock('rt_fast_reject_lock', 60_000),
283283
).rejects.toThrow('connection closed after apply');
284284
for (let attempt = 0; attempt < 20; attempt += 1) {
285-
if (await mock.get('rtsx:lock:rt_fast_reject_lock') == null) break;
285+
if (await mock.get('rtsx:lock:{rt_fast_reject_lock}') == null) break;
286286
await new Promise(resolve => setTimeout(resolve, 10));
287287
}
288-
expect(await mock.get('rtsx:lock:rt_fast_reject_lock')).toBeNull();
288+
expect(await mock.get('rtsx:lock:{rt_fast_reject_lock}')).toBeNull();
289289
});
290290

291291
test('lock polling catches abort between its precheck and listener registration', async () => {
@@ -400,7 +400,7 @@ describe('fenced record writes', () => {
400400
});
401401

402402
test('reads a corrupt record as missing instead of throwing', async () => {
403-
await mock.set('rtsx:sess:rt_bad', '{not valid json');
403+
await mock.set('rtsx:sess:{rt_bad}', '{not valid json');
404404
expect(await readRuntimeSessionRecord('rt_bad')).toBeNull();
405405
});
406406

‎service/src/runtime-session/registry.ts‎

Lines changed: 11 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@ import {
77
} from '../config';
88
import logger from '../logger';
99
import type { HostedAppRecordDetails } from '../hosted-app/record';
10+
import { hashTag } from '../redis-connection';
1011

1112
export { RUNTIME_SESSION_REDIS_COMMAND_TIMEOUT_MS } from '../config';
1213

@@ -294,7 +295,7 @@ export async function acquireRuntimeSessionLock(
294295
try {
295296
result = await runRegistryCommand(
296297
'Runtime session lock acquire',
297-
() => redis.set(`${LOCK_PREFIX}${runtimeSessionId}`, token, 'PX', ttlMs, 'NX'),
298+
() => redis.set(`${LOCK_PREFIX}${hashTag(runtimeSessionId)}`, token, 'PX', ttlMs, 'NX'),
298299
options,
299300
/* A caller-side deadline cannot cancel ioredis. The original SET may
300301
* have succeeded even if an automatic replay eventually resolves null,
@@ -404,7 +405,7 @@ export async function releaseRuntimeSessionLock(
404405
try {
405406
await runRegistryCommand(
406407
'Runtime session lock release',
407-
() => redis.releaseRuntimeSessionLockScript(`${LOCK_PREFIX}${runtimeSessionId}`, token),
408+
() => redis.releaseRuntimeSessionLockScript(`${LOCK_PREFIX}${hashTag(runtimeSessionId)}`, token),
408409
{ signal: options.signal, timeoutMs: remainingMs },
409410
);
410411
return;
@@ -453,7 +454,7 @@ export async function renewRuntimeSessionLock(
453454
const result = await runRegistryCommand(
454455
'Runtime session lock renewal',
455456
() => redis.renewRuntimeSessionLockScript(
456-
`${LOCK_PREFIX}${runtimeSessionId}`,
457+
`${LOCK_PREFIX}${hashTag(runtimeSessionId)}`,
457458
token,
458459
String(ttlMs),
459460
),
@@ -475,7 +476,7 @@ export async function readRuntimeSessionRecord(
475476
): Promise<RuntimeSessionRecord | null> {
476477
const data = await runRegistryCommand(
477478
'Runtime session record read',
478-
() => redis.get(`${SESS_PREFIX}${runtimeSessionId}`),
479+
() => redis.get(`${SESS_PREFIX}${hashTag(runtimeSessionId)}`),
479480
options,
480481
);
481482
if (data == null) return null;
@@ -502,8 +503,8 @@ export async function writeRuntimeSessionRecord(
502503
const result = await runRegistryCommand(
503504
'Runtime session record write',
504505
() => redis.writeRuntimeSessionRecordScript(
505-
`${SESS_PREFIX}${record.runtime_session_id}`,
506-
`${LOCK_PREFIX}${record.runtime_session_id}`,
506+
`${SESS_PREFIX}${hashTag(record.runtime_session_id)}`,
507+
`${LOCK_PREFIX}${hashTag(record.runtime_session_id)}`,
507508
lockToken,
508509
JSON.stringify(record),
509510
String(ttlSeconds),
@@ -533,7 +534,7 @@ export async function allocateRuntimeSessionGeneration(
533534
if (!Number.isSafeInteger(initialGeneration) || initialGeneration < 1) {
534535
throw new Error('Runtime session generation must be a positive safe integer');
535536
}
536-
const key = `${GEN_PREFIX}${runtimeSessionId}`;
537+
const key = `${GEN_PREFIX}${hashTag(runtimeSessionId)}`;
537538
const rawGeneration = await runRegistryCommand(
538539
'Runtime session generation allocation',
539540
() => redis.allocateRuntimeSessionGenerationScript(
@@ -561,7 +562,7 @@ export async function allocateCheckpointSequence(
561562
retainedMax = 0,
562563
options: RuntimeSessionRedisCommandOptions = {},
563564
): Promise<number> {
564-
const key = `${CKPT_SEQ_PREFIX}${runtimeSessionId}`;
565+
const key = `${CKPT_SEQ_PREFIX}${hashTag(runtimeSessionId)}`;
565566
return runRegistryCommand(
566567
'Runtime session checkpoint sequence allocation',
567568
() => redis.allocateCheckpointSequenceScript(
@@ -583,8 +584,8 @@ export async function removeRuntimeSession(
583584
const result = await runRegistryCommand(
584585
'Runtime session record removal',
585586
() => redis.removeRuntimeSessionScript(
586-
`${SESS_PREFIX}${runtimeSessionId}`,
587-
`${LOCK_PREFIX}${runtimeSessionId}`,
587+
`${SESS_PREFIX}${hashTag(runtimeSessionId)}`,
588+
`${LOCK_PREFIX}${hashTag(runtimeSessionId)}`,
588589
lockToken,
589590
),
590591
options,

‎service/src/sandbox-backend/lambda-microvm.test.ts‎

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -129,7 +129,7 @@ beforeAll(() => {
129129
await new Promise((resolve) => setTimeout(resolve, executeDelayMs));
130130
}
131131
if (stealSessionLockOnExecute) {
132-
await mock.set('rtsx:lock:rt_session_1', 'stolen');
132+
await mock.set('rtsx:lock:{rt_session_1}', 'stolen');
133133
}
134134
return new Response(JSON.stringify(executeResponseBody), {
135135
status: executeStatus,
@@ -728,7 +728,7 @@ describe('LambdaMicrovmSandboxBackend session execution', () => {
728728
});
729729
/* Model a worker that persisted its intent and whose RunMicrovm reached AWS,
730730
* but died before it could record the returned MicroVM id. */
731-
await mock.set('rtsx:gen:rt_session_1', '7');
731+
await mock.set('rtsx:gen:{rt_session_1}', '7');
732732
const lock = await acquireRuntimeSessionLock('rt_session_1', 60_000);
733733
expect(lock).not.toBeNull();
734734
await writeRuntimeSessionRecord({
@@ -861,7 +861,7 @@ describe('LambdaMicrovmSandboxBackend session execution', () => {
861861
* volatile registry (including its generation counter) is lost, while the
862862
* provider retains its client-token idempotency history. */
863863
await fake.terminateMicrovm(firstVmId);
864-
await mock.del('rtsx:sess:rt_session_1', 'rtsx:gen:rt_session_1');
864+
await mock.del('rtsx:sess:{rt_session_1}', 'rtsx:gen:{rt_session_1}');
865865

866866
await expect(
867867
makeBackend(fake, { imageVersion: '4' }).execute(request(), sessionContext()),
@@ -1439,7 +1439,7 @@ describe('LambdaMicrovmSandboxBackend session execution', () => {
14391439
const fake = fakeClient();
14401440
const backend = makeBackend(fake);
14411441
onExecute = async () => {
1442-
await mock.del('rtsx:sess:rt_session_1');
1442+
await mock.del('rtsx:sess:{rt_session_1}');
14431443
};
14441444

14451445
try {
@@ -1680,7 +1680,7 @@ describe('LambdaMicrovmSandboxBackend auto-checkpoint', () => {
16801680
const originalGet = redisWithGet.get.bind(mock);
16811681
let failSessionRead = false;
16821682
redisWithGet.get = async (key: string): Promise<string | null> => {
1683-
if (failSessionRead && key === 'rtsx:sess:rt_ckpt_1') {
1683+
if (failSessionRead && key === 'rtsx:sess:{rt_ckpt_1}') {
16841684
throw new Error('registry read unavailable');
16851685
}
16861686
return originalGet(key);

0 commit comments

Comments
 (0)