diff --git a/packages/api/src/cache/redisScript.ts b/packages/api/src/cache/redisScript.ts new file mode 100644 index 00000000000..2d9102c9f3e --- /dev/null +++ b/packages/api/src/cache/redisScript.ts @@ -0,0 +1,42 @@ +import { createHash } from 'node:crypto'; +import type { Redis, Cluster } from 'ioredis'; + +export type RedisScriptArg = string | number | Buffer; +export type RedisScriptClient = Pick; + +const scriptShas = new Map(); + +function scriptSha(script: string): string { + let sha = scriptShas.get(script); + if (sha == null) { + sha = createHash('sha1').update(script).digest('hex'); + scriptShas.set(script, sha); + } + return sha; +} + +export function isNoScriptError(error: unknown): boolean { + return error instanceof Error && error.message.includes('NOSCRIPT'); +} + +/** + * Runs a Lua script by its SHA1 (EVALSHA) so only the 40-byte digest crosses the wire on + * every call, and falls back to EVAL — which also loads the script into the server cache — + * when the server reports NOSCRIPT (first use, restart, SCRIPT FLUSH). Same semantics, + * atomicity, key slotting and return value as `client.eval(script, ...)`. + */ +export async function evalScript( + client: RedisScriptClient, + script: string, + numberOfKeys: number, + ...args: RedisScriptArg[] +): Promise { + try { + return (await client.evalsha(scriptSha(script), numberOfKeys, ...args)) as T; + } catch (error) { + if (!isNoScriptError(error)) { + throw error; + } + return (await client.eval(script, numberOfKeys, ...args)) as T; + } +} diff --git a/packages/api/src/stream/__tests__/RedisJobStore.spec.ts b/packages/api/src/stream/__tests__/RedisJobStore.spec.ts index 72b873051c8..3467158c730 100644 --- a/packages/api/src/stream/__tests__/RedisJobStore.spec.ts +++ b/packages/api/src/stream/__tests__/RedisJobStore.spec.ts @@ -7,6 +7,11 @@ jest.mock('~/cache/redisTelemetry', () => ({ instrumentIORedisClient: (client: unknown) => client, })); +/** Cold script cache: every EVALSHA reports NOSCRIPT so the store falls back to EVAL. */ +function evalshaNoScript(): jest.Mock { + return jest.fn().mockRejectedValue(new Error('NOSCRIPT No matching script. Please use EVAL.')); +} + type Deferred = { promise: Promise; resolve: (value: T) => void; @@ -48,6 +53,7 @@ describe('RedisJobStore', () => { const evalDrain = jest.fn().mockResolvedValue(1); const redis = { isCluster: true, + evalsha: evalshaNoScript(), eval: evalDrain, } as unknown as Cluster; const store = new RedisJobStore(redis); @@ -72,6 +78,7 @@ describe('RedisJobStore', () => { const evalBegin = jest.fn().mockResolvedValue(1); const redis = { isCluster: true, + evalsha: evalshaNoScript(), eval: evalBegin, } as unknown as Cluster; const store = new RedisJobStore(redis); @@ -97,6 +104,7 @@ describe('RedisJobStore', () => { const evalTransition = jest.fn().mockResolvedValue(0); const redis = { isCluster: true, + evalsha: evalshaNoScript(), eval: evalTransition, hgetall: jest.fn().mockResolvedValue({}), } as unknown as Cluster; @@ -174,6 +182,7 @@ describe('RedisJobStore', () => { const evalTransition = jest.fn().mockResolvedValue(0); const redis = { isCluster: true, + evalsha: evalshaNoScript(), eval: evalTransition, hgetall: jest.fn().mockResolvedValue({ streamId: 'stream-terminal-barrier', @@ -207,6 +216,7 @@ describe('RedisJobStore', () => { const evalTransition = jest.fn().mockResolvedValue(0); const redis = { isCluster: true, + evalsha: evalshaNoScript(), eval: evalTransition, } as unknown as Cluster; const store = new RedisJobStore(redis, { requiresActionTtl: 4321 }); @@ -231,6 +241,7 @@ describe('RedisJobStore', () => { const evalCommand = jest.fn().mockResolvedValue(1); const redis = { isCluster: true, + evalsha: evalshaNoScript(), eval: evalCommand, hgetall: jest .fn() @@ -285,6 +296,7 @@ describe('RedisJobStore', () => { const lrange = jest.fn(); const redis = { isCluster: true, + evalsha: evalshaNoScript(), eval: evalPeek, lrange, } as unknown as Cluster; @@ -310,6 +322,7 @@ describe('RedisJobStore', () => { .mockImplementation((...args: unknown[]) => ['', '', args[Number(args[1]) + 3]]); const redis = { isCluster: true, + evalsha: evalshaNoScript(), eval: evalJobCreation, hgetall: jest.fn(() => jobHashFromCreationCall(evalJobCreation.mock.calls[0])), sadd: jest.fn().mockResolvedValue(1), @@ -584,6 +597,7 @@ describe('RedisJobStore', () => { ]); const redis = { isCluster: true, + evalsha: evalshaNoScript(), eval: evalJobCreation, hgetall: jest.fn(() => jobHashFromCreationCall(evalJobCreation.mock.calls[0])), sadd: jest.fn().mockResolvedValue(1), @@ -639,6 +653,7 @@ describe('RedisJobStore', () => { const now = jest.spyOn(Date, 'now').mockReturnValue(100); const redis = { isCluster: true, + evalsha: evalshaNoScript(), eval: jest.fn().mockResolvedValue(['user-1', '', '100']), hgetall: jest.fn().mockResolvedValue({ streamId: 'stream-overlap', @@ -669,6 +684,7 @@ describe('RedisJobStore', () => { test('rejects creation when its durable epoch is already terminal', async () => { const redis = { isCluster: true, + evalsha: evalshaNoScript(), eval: jest.fn().mockResolvedValue(['', '', '100']), hgetall: jest.fn().mockResolvedValue({ streamId: 'stream-terminal-create', @@ -692,6 +708,7 @@ describe('RedisJobStore', () => { const evalRedis = jest.fn().mockResolvedValue(false); const redis = { isCluster: true, + evalsha: evalshaNoScript(), eval: evalRedis, } as unknown as Cluster; const store = new RedisJobStore(redis); @@ -770,6 +787,7 @@ describe('RedisJobStore', () => { }); const redis = { isCluster: true, + evalsha: evalshaNoScript(), eval: evalJobCreation, sadd: jest.fn((key: string) => { if (key === 'stream:running') { @@ -800,6 +818,7 @@ describe('RedisJobStore', () => { return job; }); + await waitFor(() => started.length === 1); expect(started).toEqual(['job']); evalResult.resolve(1); await waitFor(() => started.length === 6); @@ -879,6 +898,7 @@ describe('RedisJobStore', () => { }); const redis = { isCluster: true, + evalsha: evalshaNoScript(), eval: jest.fn(async (_script: string, keyCount: number, ...args: string[]) => { if (keyCount === 10) { durableHash = { ...durableHash, status: 'requires_action' }; @@ -922,6 +942,7 @@ describe('RedisJobStore', () => { const evalCommand = jest.fn().mockResolvedValue(0); const redis = { isCluster: true, + evalsha: evalshaNoScript(), eval: evalCommand, hgetall: jest.fn().mockResolvedValue({ streamId: 'stream-guarded', @@ -957,6 +978,7 @@ describe('RedisJobStore', () => { const evalCommand = jest.fn().mockResolvedValue(0); const redis = { isCluster: true, + evalsha: evalshaNoScript(), eval: evalCommand, } as unknown as Cluster; const store = new RedisJobStore(redis); @@ -995,6 +1017,7 @@ describe('RedisJobStore', () => { const evalCommand = jest.fn().mockResolvedValue(0); const redis = { isCluster: true, + evalsha: evalshaNoScript(), eval: evalCommand, } as unknown as Cluster; const store = new RedisJobStore(redis); @@ -1026,6 +1049,7 @@ describe('RedisJobStore', () => { const sadd = jest.fn().mockResolvedValue(1); const redis = { isCluster: true, + evalsha: evalshaNoScript(), eval: evalTransition, sadd, srem: jest.fn().mockResolvedValue(1), @@ -1083,6 +1107,7 @@ describe('RedisJobStore', () => { }; const redis = { isCluster: true, + evalsha: evalshaNoScript(), eval: evalTransition, sadd, srem: jest.fn().mockResolvedValue(1), @@ -1144,6 +1169,7 @@ describe('RedisJobStore', () => { const srem = jest.fn().mockResolvedValue(1); const redis = { isCluster: true, + evalsha: evalshaNoScript(), eval: jest.fn().mockResolvedValue(1), smembers: jest.fn().mockResolvedValue([member]), srem, @@ -1191,6 +1217,7 @@ describe('RedisJobStore', () => { const srem = jest.fn().mockResolvedValue(1); const redis = { isCluster: true, + evalsha: evalshaNoScript(), eval: evalClear, srem, hgetall: jest.fn().mockResolvedValue({}), @@ -1233,6 +1260,7 @@ describe('RedisJobStore', () => { const evalCommand = jest.fn().mockResolvedValue(0); const redis = { isCluster: true, + evalsha: evalshaNoScript(), eval: evalCommand, } as unknown as Cluster; const store = new RedisJobStore(redis); diff --git a/packages/api/src/stream/__tests__/RedisJobStore.stream_integration.spec.ts b/packages/api/src/stream/__tests__/RedisJobStore.stream_integration.spec.ts index df15322b643..9653d8195f0 100644 --- a/packages/api/src/stream/__tests__/RedisJobStore.stream_integration.spec.ts +++ b/packages/api/src/stream/__tests__/RedisJobStore.stream_integration.spec.ts @@ -8,7 +8,7 @@ import { PAUSE_PERSISTENCE_TIMEOUT_ERROR, STEER_ENQUEUE_RECEIPT_FULL, } from '../interfaces/IJobStore'; -import { clearRedisTestPrefix } from './helpers/redis'; +import { clearRedisTestPrefix, flushScriptCache } from './helpers/redis'; /** Suppress winston Console transport output (survives jest.resetModules) */ jest.spyOn(console, 'log').mockImplementation(); @@ -492,6 +492,7 @@ describe('RedisJobStore Integration Tests', () => { const streamId = `terminal-epoch-${Date.now()}`; const userId = 'terminal-epoch-user'; const now = jest.spyOn(Date, 'now').mockReturnValue(1000); + await flushScriptCache(ioredisClient); const originalEval = ioredisClient.eval.bind(ioredisClient) as ( script: string | Buffer, numberOfKeys: number, @@ -1632,7 +1633,7 @@ describe('RedisJobStore Integration Tests', () => { store.setCollectedUsage(streamId, [{ input_tokens: 1, output_tokens: 2 }]); const evalSpy = jest - .spyOn(ioredisClient, 'eval') + .spyOn(ioredisClient, 'evalsha') .mockRejectedValueOnce(new Error('replacement write failed')); try { await expect(store.createJob(streamId, 'user-1', streamId)).rejects.toThrow( @@ -4084,6 +4085,7 @@ describe('RedisJobStore Integration Tests', () => { String(Date.now() - 10_000), ); + await flushScriptCache(ioredisClient); const originalEval = ioredisClient.eval.bind(ioredisClient) as ( script: string | Buffer, numberOfKeys: number, diff --git a/packages/api/src/stream/__tests__/helpers/publisher.ts b/packages/api/src/stream/__tests__/helpers/publisher.ts index 9c2c65efabe..5811a83f817 100644 --- a/packages/api/src/stream/__tests__/helpers/publisher.ts +++ b/packages/api/src/stream/__tests__/helpers/publisher.ts @@ -6,6 +6,7 @@ export interface MockPublisher { get: jest.Mock; set: jest.Mock; del: jest.Mock; + evalsha: jest.Mock; eval: jest.Mock; } @@ -54,6 +55,9 @@ export function createMockPublisher(): MockPublisher { } return Promise.resolve(keys.length); }), + evalsha: jest + .fn() + .mockRejectedValue(new Error('NOSCRIPT No matching script. Please use EVAL.')), eval: jest.fn(), }; diff --git a/packages/api/src/stream/__tests__/helpers/redis.ts b/packages/api/src/stream/__tests__/helpers/redis.ts index 694bc98d72c..992d11673c2 100644 --- a/packages/api/src/stream/__tests__/helpers/redis.ts +++ b/packages/api/src/stream/__tests__/helpers/redis.ts @@ -34,6 +34,17 @@ export function createRedisTestClient(keyPrefix: string): RedisTestClient { return new IoRedis(primary.href, { ...redisOptions, lazyConnect: true }); } +/** + * Empties every node's Lua script cache so the next run of each script misses EVALSHA and + * falls back to EVAL — the only path where a test spying on `eval` can observe the script body. + */ +export async function flushScriptCache(redis: RedisTestClient): Promise { + const nodes = (redis as Cluster).isCluster + ? (redis as Cluster).nodes('master') + : [redis as Redis]; + await Promise.all(nodes.map((node) => node.script('FLUSH'))); +} + /** Delete only this suite's keys, including keys spread across cluster masters. */ export async function clearRedisTestPrefix( redis: RedisTestClient, diff --git a/packages/api/src/stream/__tests__/protocolRollout.stream_integration.spec.ts b/packages/api/src/stream/__tests__/protocolRollout.stream_integration.spec.ts index 94319d4dc4b..9763631e4f5 100644 --- a/packages/api/src/stream/__tests__/protocolRollout.stream_integration.spec.ts +++ b/packages/api/src/stream/__tests__/protocolRollout.stream_integration.spec.ts @@ -1,5 +1,10 @@ import type { SteerQueueItem } from '../interfaces/IJobStore'; -import { clearRedisTestPrefix, createRedisTestClient, type RedisTestClient } from './helpers/redis'; +import { + clearRedisTestPrefix, + createRedisTestClient, + flushScriptCache, + type RedisTestClient, +} from './helpers/redis'; import { InMemoryEventTransport } from '../implementations/InMemoryEventTransport'; import { RedisEventTransport } from '../implementations/RedisEventTransport'; import { GenerationJobManagerClass } from '../GenerationJobManager'; @@ -377,6 +382,7 @@ describe('Redis generation protocol rollout bridge', () => { racedPredecessor = await owner.createJob(streamId, userId, streamId, { initialMetadata: { generationProtocolVersion: 2 }, }); + await flushScriptCache(redis); injectLostCreateReply = true; } return observed; @@ -530,6 +536,7 @@ describe('Redis generation protocol rollout bridge', () => { initialMetadata: { generationProtocolVersion: 2 }, }); expect(await ownerStore.getJob(streamId)).toMatchObject({ providerAbortReady: true }); + await flushScriptCache(redis); injectLostCreateReply = true; await expect( diff --git a/packages/api/src/stream/implementations/RedisEventTransport.ts b/packages/api/src/stream/implementations/RedisEventTransport.ts index 274c089fd5c..6746f9ddcec 100644 --- a/packages/api/src/stream/implementations/RedisEventTransport.ts +++ b/packages/api/src/stream/implementations/RedisEventTransport.ts @@ -15,6 +15,7 @@ import { import { registerChunkPublicationCapability } from '~/stream/internal/chunkPublication'; import { GenerationPublicationFencedError } from '~/stream/interfaces/IJobStore'; import { instrumentIORedisClient, RedisUseCases } from '~/cache/redisTelemetry'; +import { evalScript } from '~/cache/redisScript'; /** * Redis key prefixes for pub/sub channels @@ -352,7 +353,8 @@ export class RedisEventTransport implements IEventTransport { * commits but its promise never settles, the marker timeout still releases attachment * admission instead of leaving every surviving local subscriber deferred forever. */ const operation = Promise.all([ - this.publisher.eval( + evalScript( + this.publisher, CAPTURE_SUBSCRIPTION_FRONTIER_LUA, 1, KEYS.sequence(streamId), @@ -471,7 +473,8 @@ export class RedisEventTransport implements IEventTransport { allowRetainedEpoch = false, requireActiveJob = false, ): Promise { - const seq = await this.publisher.eval( + const seq = await evalScript( + this.publisher, PUBLISH_SEQ_LUA, 3, KEYS.sequence(streamId), @@ -1372,7 +1375,8 @@ export class RedisEventTransport implements IEventTransport { data: event, generationId: replacedGenerationId, }); - const result = await this.publisher.eval( + const result = await evalScript( + this.publisher, PUBLISH_REPLACED_DONE_LUA, 2, KEYS.sequence(streamId), diff --git a/packages/api/src/stream/implementations/RedisJobStore.ts b/packages/api/src/stream/implementations/RedisJobStore.ts index 982462f8576..d8d26960f75 100644 --- a/packages/api/src/stream/implementations/RedisJobStore.ts +++ b/packages/api/src/stream/implementations/RedisJobStore.ts @@ -49,6 +49,7 @@ import { import { instrumentIORedisClient, RedisUseCases } from '~/cache/redisTelemetry'; import { RecoveredSteerPayloadMismatchError } from '~/stream/SteerRecovery'; import { createCheckpointNamespace } from '~/stream/checkpoints'; +import { evalScript } from '~/cache/redisScript'; const CLIENT_REQUEST_ID_PATTERN = /^[A-Za-z0-9:_-]{1,128}$/; @@ -2068,7 +2069,8 @@ export class RedisJobStore implements IJobStoreV2 { * cleanup or natural job expiry cannot make both keys disappear together. */ const createParkedTtl = recoveredSteerId != null ? this.runningStorageTtlSeconds() + parkedTtl : parkedTtl; - const previousOwner = await this.redis.eval( + const previousOwner = await evalScript( + this.redis, JOB_CREATE_LUA, 10, key, @@ -2281,7 +2283,8 @@ export class RedisJobStore implements IJobStoreV2 { if (creationAttemptId.length === 0 || replacedCreatedAts.length === 0) { return false; } - const acknowledged = await this.redis.eval( + const acknowledged = await evalScript( + this.redis, REPLACEMENT_RECEIPT_ACK_LUA, 1, KEYS.job(streamId), @@ -2334,7 +2337,8 @@ export class RedisJobStore implements IJobStoreV2 { ? Math.max(this.ttl.completed, TERMINAL_PERSISTENCE_RETENTION_TTL_S) : this.ttl.completed; const fields = Object.entries(serialized).flat(); - const updated = await this.redis.eval( + const updated = await evalScript( + this.redis, JOB_UPDATE_LUA, 4, key, @@ -2373,7 +2377,8 @@ export class RedisJobStore implements IJobStoreV2 { 'recoveryMethod' | 'recoveryOutcome' | 'recoveryCompletedAt' | 'recoveryFailureReason' >, ): Promise { - const settled = await this.redis.eval( + const settled = await evalScript( + this.redis, SETTLE_EARLY_BUFFER_RECOVERY_LUA, 1, KEYS.job(streamId), @@ -2390,7 +2395,8 @@ export class RedisJobStore implements IJobStoreV2 { overflowId: string, finalizedOverflow: EarlyBufferOverflowState, ): Promise { - const finalized = await this.redis.eval( + const finalized = await evalScript( + this.redis, FINALIZE_EARLY_BUFFER_OVERFLOW_LUA, 1, KEYS.job(streamId), @@ -2404,7 +2410,8 @@ export class RedisJobStore implements IJobStoreV2 { async hasSubscriberAttached(streamId: string, expectedCreatedAt: number): Promise { return ( Number( - await this.redis.eval( + await evalScript( + this.redis, HAS_SUBSCRIBER_ATTACHED_LUA, 1, KEYS.job(streamId), @@ -2423,7 +2430,8 @@ export class RedisJobStore implements IJobStoreV2 { ): Promise { return ( Number( - await this.redis.eval( + await evalScript( + this.redis, CLAIM_FIRST_SUBSCRIBER_LUA, 2, KEYS.job(streamId), @@ -2442,7 +2450,8 @@ export class RedisJobStore implements IJobStoreV2 { expectedCreatedAt: number, subscriberId: string, ): Promise { - await this.redis.eval( + await evalScript( + this.redis, DETACH_SUBSCRIBER_LUA, 2, KEYS.job(streamId), @@ -2459,7 +2468,8 @@ export class RedisJobStore implements IJobStoreV2 { ): Promise { return ( Number( - await this.redis.eval( + await evalScript( + this.redis, HAS_ACTIVE_SUBSCRIBER_LUA, 2, KEYS.job(streamId), @@ -2478,7 +2488,8 @@ export class RedisJobStore implements IJobStoreV2 { ): Promise { return ( Number( - await this.redis.eval( + await evalScript( + this.redis, PROVIDER_DRAIN_LUA, 1, KEYS.job(streamId), @@ -2496,7 +2507,8 @@ export class RedisJobStore implements IJobStoreV2 { ): Promise { return ( Number( - await this.redis.eval( + await evalScript( + this.redis, PROVIDER_BEGIN_LUA, 1, KEYS.job(streamId), @@ -2514,7 +2526,8 @@ export class RedisJobStore implements IJobStoreV2 { ): Promise { return ( Number( - await this.redis.eval( + await evalScript( + this.redis, TERMINAL_PERSISTENCE_FINALIZE_LUA, 1, KEYS.job(streamId), @@ -2562,7 +2575,8 @@ export class RedisJobStore implements IJobStoreV2 { } private async retainCleanupOwner(job: SerializableJobData): Promise { - await this.redis.eval( + await evalScript( + this.redis, OWNER_MEMBERSHIP_RECONCILE_LUA, 1, KEYS.userJobs(job.userId, job.tenantId), @@ -2652,7 +2666,8 @@ export class RedisJobStore implements IJobStoreV2 { } if (activeUserKey) { operations.push( - this.redis.eval( + evalScript( + this.redis, OWNER_MEMBERSHIP_RECONCILE_LUA, 1, activeUserKey, @@ -2913,7 +2928,8 @@ export class RedisJobStore implements IJobStoreV2 { // 1) Single-winner decision: an atomic CAS on the single-slot job hash. // Works identically on cluster and single-node, so two concurrent // resolves can never both win (and drive the run twice). - const result = await this.redis.eval( + const result = await evalScript( + this.redis, JOB_CAS_LUA, 10, key, @@ -2982,7 +2998,8 @@ export class RedisJobStore implements IJobStoreV2 { value: IdempotencyClaimValue, ttlSeconds: number, ): Promise { - const result = await this.redis.eval( + const result = await evalScript( + this.redis, IDEMPOTENCY_CLAIM_LUA, 1, KEYS.idempotency(key), @@ -3026,7 +3043,8 @@ export class RedisJobStore implements IJobStoreV2 { value: IdempotencyClaimValue, ttlSeconds: number, ): Promise { - const taken = await this.redis.eval( + const taken = await evalScript( + this.redis, IDEMPOTENCY_TAKEOVER_LUA, 1, KEYS.idempotency(key), @@ -3043,7 +3061,8 @@ export class RedisJobStore implements IJobStoreV2 { startedAt: number, ttlSeconds: number, ): Promise { - const marked = await this.redis.eval( + const marked = await evalScript( + this.redis, IDEMPOTENCY_MARK_STARTED_LUA, 1, KEYS.idempotency(key), @@ -3065,7 +3084,8 @@ export class RedisJobStore implements IJobStoreV2 { ttlSeconds: number, allowMissingClientRequestId = false, ): Promise { - const adopted = await this.redis.eval( + const adopted = await evalScript( + this.redis, IDEMPOTENCY_ADOPT_LIVE_JOB_LUA, 2, KEYS.idempotency(key), @@ -3082,7 +3102,8 @@ export class RedisJobStore implements IJobStoreV2 { } async releaseIdempotencyKey(key: string, expected?: IdempotencyClaimValue): Promise { - await this.redis.eval( + await evalScript( + this.redis, IDEMPOTENCY_RELEASE_LUA, 1, KEYS.idempotency(key), @@ -3094,7 +3115,8 @@ export class RedisJobStore implements IJobStoreV2 { const observedJob = await this.getJob(streamId); const targetCreatedAt = expectedCreatedAt ?? observedJob?.createdAt; const expectMissing = expectedCreatedAt == null && observedJob == null; - const deleted = await this.redis.eval( + const deleted = await evalScript( + this.redis, JOB_DELETE_LUA, 5, KEYS.job(streamId), @@ -3123,7 +3145,8 @@ export class RedisJobStore implements IJobStoreV2 { observedJob: SerializableJobData, now: number, ): Promise { - const deleted = await this.redis.eval( + const deleted = await evalScript( + this.redis, STALE_JOB_DELETE_LUA, 9, KEYS.job(streamId), @@ -3252,7 +3275,8 @@ export class RedisJobStore implements IJobStoreV2 { ) { const recovered = Number( - await this.redis.eval( + await evalScript( + this.redis, RECOVER_TERMINAL_PROVIDER_DRAIN_LUA, 1, KEYS.job(job.streamId), @@ -3331,7 +3355,8 @@ export class RedisJobStore implements IJobStoreV2 { // replacement at the same streamId is never cleared through its predecessor. The HDEL // and configured evidence-TTL reset happen atomically. The global retry // member includes this generation, so removing it cannot affect a successor. - const cleared = (await this.redis.eval( + const cleared = (await evalScript( + this.redis, 'if redis.call("HGET", KEYS[1], "createdAt") ~= ARGV[1] then return 0 end ' + 'local detachedStatus = redis.call("HGET", KEYS[1], "detachedAgentEventTerminalStatus") ' + 'if detachedStatus then redis.call("HSET", KEYS[1], "status", detachedStatus) end ' + @@ -4360,7 +4385,8 @@ export class RedisJobStore implements IJobStoreV2 { if (expectedCreatedAt == null) { return this.redis.get(KEYS.runSteps(streamId)); } - const data = await this.redis.eval( + const data = await evalScript( + this.redis, RUNSTEPS_READ_LUA, 2, KEYS.job(streamId), @@ -4391,7 +4417,8 @@ export class RedisJobStore implements IJobStoreV2 { streamId: string, expectedCreatedAt?: number, ): Promise { - await this.redis.eval( + await evalScript( + this.redis, CONTENT_CLEAR_LUA, 3, KEYS.chunks(streamId), @@ -4408,7 +4435,8 @@ export class RedisJobStore implements IJobStoreV2 { item: SteerQueueItem, expectedCreatedAt?: number, ): Promise { - const result = await this.redis.eval( + const result = await evalScript( + this.redis, STEER_ENQUEUE_LUA, 2, KEYS.job(streamId), @@ -4430,7 +4458,8 @@ export class RedisJobStore implements IJobStoreV2 { wantsPreempt: boolean, expectedCreatedAt?: number, ): Promise { - const result = await this.redis.eval( + const result = await evalScript( + this.redis, STEER_ENQUEUE_VERSIONED_LUA, 2, KEYS.job(streamId), @@ -4456,7 +4485,8 @@ export class RedisJobStore implements IJobStoreV2 { } async getSteerReceipt(streamId: string, clientSteerId: string): Promise { - const raw = await this.redis.eval( + const raw = await evalScript( + this.redis, STEER_RECEIPT_GET_LUA, 3, KEYS.steerReceipts(streamId), @@ -4483,7 +4513,8 @@ export class RedisJobStore implements IJobStoreV2 { wantsPreempt: boolean, expectedCreatedAt?: number, ): Promise { - const result = await this.redis.eval( + const result = await evalScript( + this.redis, STEER_ENQUEUE_RECEIPT_LUA, 4, KEYS.job(streamId), @@ -4515,7 +4546,8 @@ export class RedisJobStore implements IJobStoreV2 { } async drainSteers(streamId: string, expectedCreatedAt?: number): Promise { - const raw = await this.redis.eval( + const raw = await evalScript( + this.redis, STEER_DRAIN_LUA, 5, KEYS.job(streamId), @@ -4534,7 +4566,8 @@ export class RedisJobStore implements IJobStoreV2 { policy: TerminalSteerAdmissionPolicy, expectedCreatedAt?: number, ): Promise { - const raw = await this.redis.eval( + const raw = await evalScript( + this.redis, STEER_TERMINAL_ADMISSION_LUA, 5, KEYS.job(streamId), @@ -4566,7 +4599,8 @@ export class RedisJobStore implements IJobStoreV2 { if (items.length === 0) { return true; } - const restored = await this.redis.eval( + const restored = await evalScript( + this.redis, STEER_RESTORE_CLAIMED_LUA, 4, KEYS.job(streamId), @@ -4584,7 +4618,8 @@ export class RedisJobStore implements IJobStoreV2 { streamId: string, expectedCreatedAt?: number, ): Promise { - const raw = await this.redis.eval( + const raw = await evalScript( + this.redis, STEER_CLOSE_DRAIN_LUA, 6, KEYS.job(streamId), @@ -4606,7 +4641,8 @@ export class RedisJobStore implements IJobStoreV2 { const raw = expectedCreatedAt == null ? await this.redis.lrange(KEYS.steers(streamId), 0, -1) - : await this.redis.eval( + : await evalScript( + this.redis, STEER_PEEK_LUA, 2, KEYS.job(streamId), @@ -4617,7 +4653,8 @@ export class RedisJobStore implements IJobStoreV2 { } async peekClaimedSteers(streamId: string, expectedCreatedAt?: number): Promise { - const raw = await this.redis.eval( + const raw = await evalScript( + this.redis, STEER_PEEK_CLAIMED_LUA, 2, KEYS.job(streamId), @@ -4636,7 +4673,8 @@ export class RedisJobStore implements IJobStoreV2 { steerId: string, expectedCreatedAt?: number, ): Promise { - const removed = (await this.redis.eval( + const removed = (await evalScript( + this.redis, STEER_REMOVE_LUA, 3, KEYS.job(streamId), @@ -4661,7 +4699,8 @@ export class RedisJobStore implements IJobStoreV2 { steerId: string, expectedCreatedAt?: number, ): Promise { - const result = await this.redis.eval( + const result = await evalScript( + this.redis, STEER_ARM_LUA, 3, KEYS.job(streamId), @@ -4685,7 +4724,8 @@ export class RedisJobStore implements IJobStoreV2 { streamId: string, expectedCreatedAt?: number, ): Promise { - const changed = await this.redis.eval( + const changed = await evalScript( + this.redis, STEER_DOWNGRADE_PREEMPTS_LUA, 3, KEYS.job(streamId), @@ -4701,7 +4741,8 @@ export class RedisJobStore implements IJobStoreV2 { async parkSteers(streamId: string, payload: string, expectedCreatedAt?: number): Promise { const ttl = this.ttl.completed > 0 ? this.ttl.completed : PARKED_RECOVERY_TTL_S; - await this.redis.eval( + await evalScript( + this.redis, PARK_STEERS_LUA, 2, KEYS.job(streamId), @@ -4726,7 +4767,8 @@ export class RedisJobStore implements IJobStoreV2 { ownerTenantId?: string, requestedProtocolVersion: 1 | 2 = 1, ): Promise { - const claimed = await this.redis.eval( + const claimed = await evalScript( + this.redis, CLAIM_PARKED_LUA, 2, KEYS.parkedSteers(streamId), @@ -4756,7 +4798,8 @@ export class RedisJobStore implements IJobStoreV2 { ownerTenantId: string | undefined, expectedCreatedAt: number, ): Promise { - const consumed = await this.redis.eval( + const consumed = await evalScript( + this.redis, CONSUME_PARKED_STEER_LUA, 3, KEYS.job(streamId), @@ -4778,7 +4821,8 @@ export class RedisJobStore implements IJobStoreV2 { ownerTenantId?: string, expectedGenerationCreatedAt?: number, ): Promise { - const discarded = await this.redis.eval( + const discarded = await evalScript( + this.redis, DISCARD_STEER_LEFTOVER_LUA, 3, KEYS.steerReceipts(streamId), @@ -4844,7 +4888,8 @@ export class RedisJobStore implements IJobStoreV2 { // even when the pause's own EXPIRE no-op'd because this key didn't exist yet, while a // normally-running run still settles on the short running TTL. Both keys share the // {streamId} hash tag, so the multi-key eval stays on one slot under Redis Cluster. - const appended = await this.redis.eval( + const appended = await evalScript( + this.redis, CHUNK_APPEND_LUA, 8, key, @@ -4924,38 +4969,37 @@ export class RedisJobStore implements IJobStoreV2 { } const { expectedCreatedAt, events, settlers } = pending; - return this.redis - .eval( - CHUNK_APPEND_BATCH_LUA, - 8, - KEYS.chunks(streamId), - KEYS.job(streamId), - KEYS.steerReceipts(streamId), - KEYS.steerReceiptOrder(streamId), - KEYS.claimedSteers(streamId), - KEYS.steers(streamId), - KEYS.parkedSteers(streamId), - KEYS.generationEpoch(streamId), - String(this.runningStorageTtlSeconds()), - expectedCreatedAt != null ? String(expectedCreatedAt) : '', - String(Date.now()), - String(this.parkedRecoveryTtlSeconds()), - String(GENERATION_EPOCH_GRACE_TTL_S), - ...events, - ) - .then( - (appended) => { - const committed = appended === 1; - for (const settler of settlers) { - settler.resolve(committed); - } - }, - (err) => { - for (const settler of settlers) { - settler.reject(err); - } - }, - ); + return evalScript( + this.redis, + CHUNK_APPEND_BATCH_LUA, + 8, + KEYS.chunks(streamId), + KEYS.job(streamId), + KEYS.steerReceipts(streamId), + KEYS.steerReceiptOrder(streamId), + KEYS.claimedSteers(streamId), + KEYS.steers(streamId), + KEYS.parkedSteers(streamId), + KEYS.generationEpoch(streamId), + String(this.runningStorageTtlSeconds()), + expectedCreatedAt != null ? String(expectedCreatedAt) : '', + String(Date.now()), + String(this.parkedRecoveryTtlSeconds()), + String(GENERATION_EPOCH_GRACE_TTL_S), + ...events, + ).then( + (appended) => { + const committed = appended === 1; + for (const settler of settlers) { + settler.resolve(committed); + } + }, + (err) => { + for (const settler of settlers) { + settler.reject(err); + } + }, + ); } /** Persist a stream's pending coalesced appends now (pre-transition barrier). */ @@ -4983,7 +5027,8 @@ export class RedisJobStore implements IJobStoreV2 { let rawEntries: unknown; let rawDurableEventCount: unknown; if (includeDurableEventCount) { - const rawSnapshot = await this.redis.eval( + const rawSnapshot = await evalScript( + this.redis, CHUNKS_RECOVERY_READ_LUA, 2, KEYS.job(streamId), @@ -4995,7 +5040,8 @@ export class RedisJobStore implements IJobStoreV2 { rawEntries = expectedCreatedAt == null ? await this.redis.xrange(KEYS.chunks(streamId), '-', '+') - : await this.redis.eval( + : await evalScript( + this.redis, CHUNKS_READ_LUA, 2, KEYS.job(streamId), @@ -5041,7 +5087,8 @@ export class RedisJobStore implements IJobStoreV2 { runSteps: Agents.RunStep[], expectedCreatedAt?: number, ): Promise { - await this.redis.eval( + await evalScript( + this.redis, RUNSTEPS_SAVE_LUA, 2, KEYS.runSteps(streamId),