Files
triggerdotdev--trigger.dev/apps/webapp/app/services/taskIdentifierCache.server.ts
T
Eric Allam 567e2a2c32 feat(webapp,redis): handle READONLY / LOADING during ElastiCache failover (#3548)
## Summary

During an ElastiCache role swap (failover) or node-type change (vertical
scale), the ioredis TCP/TLS connection stays open but the server starts
answering with `READONLY` (the client is talking to a node that became a
replica) or `LOADING` (node still loading data from disk). Without an
explicit hook, those errors surface to caller code as `ReplyError`
instances — every write op on the affected connection fails until the
cluster fully cuts over.

This PR adds `reconnectOnError` to every prod ioredis client so the
disconnect + reconnect + retry cycle absorbs these errors and caller
code never sees them.

## Fix

```ts
export function defaultReconnectOnError(err: Error): boolean | 1 | 2 {
  const msg = err.message ?? "";
  if (msg.startsWith("READONLY") || msg.startsWith("LOADING")) return 2;
  return false;
}
```

Returning `2` tells ioredis to disconnect, reconnect, and re-issue the
failed command. After reconnect, DNS / SG state routes the new socket to
a writable node.

The helper lives in `@internal/redis` and is wired into both the shared
`createRedisClient` (which covers RunQueue, schedule-engine,
redis-worker, and every other internal-package consumer) and the direct
`new Redis(...)` call sites in the webapp.

V1-only marqs files are intentionally not migrated.

## Test plan

- [x] `pnpm run typecheck --filter webapp`
- [x] `pnpm run typecheck --filter @internal/run-engine`
- [x] Verified end-to-end against a live ElastiCache vertical-scale
event — caller-surfaced errors went from tens of thousands during the
cutover window down to a handful per ioredis client
- [ ] Confirm steady-state behavior unchanged after deploy
2026-05-11 07:17:07 +01:00

118 lines
2.8 KiB
TypeScript

import { Redis } from "ioredis";
import { defaultReconnectOnError } from "@internal/redis";
import type { TaskTriggerSource } from "@trigger.dev/database";
import { env } from "~/env.server";
import { singleton } from "~/utils/singleton";
import { logger } from "./logger.server";
const KEY_PREFIX = "tids:";
type CachedTaskIdentifier = {
s: string;
ts: TaskTriggerSource;
live: boolean;
};
export type TaskIdentifierEntry = {
slug: string;
triggerSource: TaskTriggerSource;
isInLatestDeployment: boolean;
};
function buildKey(environmentId: string): string {
return `${KEY_PREFIX}${environmentId}`;
}
function encode(entry: TaskIdentifierEntry): string {
return JSON.stringify({
s: entry.slug,
ts: entry.triggerSource,
live: entry.isInLatestDeployment,
} satisfies CachedTaskIdentifier);
}
function decode(raw: string): TaskIdentifierEntry {
const parsed = JSON.parse(raw) as CachedTaskIdentifier;
return {
slug: parsed.s,
triggerSource: parsed.ts,
isInLatestDeployment: parsed.live,
};
}
function initializeRedis(): Redis | undefined {
const host = env.CACHE_REDIS_HOST;
if (!host) {
return undefined;
}
return new Redis({
connectionName: "taskIdentifierCache",
host,
port: env.CACHE_REDIS_PORT,
username: env.CACHE_REDIS_USERNAME,
password: env.CACHE_REDIS_PASSWORD,
keyPrefix: "tr:",
enableAutoPipelining: true,
reconnectOnError: defaultReconnectOnError,
...(env.CACHE_REDIS_TLS_DISABLED === "true" ? {} : { tls: {} }),
});
}
const redis = singleton("taskIdentifierCache", initializeRedis);
export async function populateTaskIdentifierCache(
environmentId: string,
identifiers: TaskIdentifierEntry[]
): Promise<void> {
if (!redis) return;
try {
const key = buildKey(environmentId);
const pipeline = redis.pipeline();
pipeline.del(key);
if (identifiers.length > 0) {
pipeline.sadd(key, ...identifiers.map(encode));
}
await pipeline.exec();
} catch (error) {
logger.error("Failed to populate task identifier cache", {
environmentId,
error,
});
}
}
export async function invalidateTaskIdentifierCache(environmentId: string): Promise<void> {
if (!redis) return;
try {
const key = buildKey(environmentId);
await redis.del(key);
} catch (error) {
logger.error("Failed to invalidate task identifier cache", {
environmentId,
error,
});
}
}
export async function getTaskIdentifiersFromCache(
environmentId: string
): Promise<TaskIdentifierEntry[] | null> {
if (!redis) return null;
try {
const key = buildKey(environmentId);
const members = await redis.smembers(key);
if (members.length === 0) return null;
return members.map(decode);
} catch (error) {
logger.error("Failed to get task identifiers from cache", {
environmentId,
error,
});
return null;
}
}