/** * Redis Streams distributed job queue. * * Provides horizontal crawler scaling via Redis Streams consumer groups. * Multiple crawler nodes can pull from the same stream — each job is * delivered to exactly one consumer (at-least-once via ACK/PENDING). * * Used when REDIS_URL is set AND JOB_QUEUE_BACKEND=redis. * Falls back to the Postgres-backed queue when Redis is unavailable. * * Stream key: zeta:crawl:stream * Group name: zeta:crawl:workers * Consumer id: zeta:worker:: */ import { createHash } from 'node:crypto'; import { hostname } from 'node:os'; export interface RedisQueueJob { id: string; // Redis stream entry ID (e.g. "1720001234567-0") jobId: string; // Postgres job row ID tenantId: string; projectId: string; payload: Record; enqueuedAt: string; } export interface RedisQueueClient { enqueue(jobId: string, tenantId: string, projectId: string, payload?: Record): Promise; dequeue(timeoutMs?: number): Promise; ack(streamEntryId: string): Promise; nack(streamEntryId: string): Promise; reclaimStale(maxIdleMs?: number): Promise; stats(): Promise<{ pending: number; delivered: number; groupLag: number }>; close(): Promise; } const STREAM_KEY = 'zeta:crawl:stream'; const GROUP_NAME = 'zeta:crawl:workers'; const MAX_LEN = 10_000; // approx cap on stream length function consumerId(): string { return `zeta:worker:${hostname()}:${process.pid}`; } let ioredis: any; async function getIoredis() { if (!ioredis) { ioredis = await import('ioredis').then(m => m.default ?? m); } return ioredis; } /** * Create a Redis Streams job queue client. * Throws if REDIS_URL is not set. */ export async function createRedisQueue(): Promise { const redisUrl = process.env.REDIS_URL; if (!redisUrl) throw new Error('REDIS_URL env var required for redis queue backend'); const Redis = await getIoredis(); const redis = new Redis(redisUrl, { maxRetriesPerRequest: 3, enableReadyCheck: true, lazyConnect: false, }); // Create stream + consumer group if they don't exist await redis.xgroup('CREATE', STREAM_KEY, GROUP_NAME, '$', 'MKSTREAM').catch((e: any) => { if (!e.message?.includes('BUSYGROUP')) throw e; }); const consumer = consumerId(); return { async enqueue(jobId, tenantId, projectId, payload = {}): Promise { const entryId = await redis.xadd( STREAM_KEY, 'MAXLEN', '~', MAX_LEN, '*', 'jobId', jobId, 'tenantId', tenantId, 'projectId', projectId, 'payload', JSON.stringify(payload), 'enqueuedAt', new Date().toISOString(), ); return entryId as string; }, async dequeue(timeoutMs = 5_000): Promise { // First drain pending (unacknowledged) messages for this consumer const pending: any[][] = await redis.xreadgroup( 'GROUP', GROUP_NAME, consumer, 'COUNT', '1', 'STREAMS', STREAM_KEY, '0', ) ?? []; let entries = pending?.[0]?.[1] ?? []; if (!entries.length) { // Block-read for new messages const fresh: any[][] = await redis.xreadgroup( 'GROUP', GROUP_NAME, consumer, 'COUNT', '1', 'BLOCK', timeoutMs, 'STREAMS', STREAM_KEY, '>', ) ?? []; entries = fresh?.[0]?.[1] ?? []; } if (!entries.length) return null; const [entryId, fields] = entries[0] as [string, string[]]; const obj: Record = {}; for (let i = 0; i < fields.length; i += 2) { obj[fields[i]] = fields[i + 1]; } return { id: entryId, jobId: obj.jobId, tenantId: obj.tenantId, projectId: obj.projectId, payload: JSON.parse(obj.payload ?? '{}'), enqueuedAt: obj.enqueuedAt, }; }, async ack(streamEntryId: string): Promise { await redis.xack(STREAM_KEY, GROUP_NAME, streamEntryId); }, async nack(streamEntryId: string): Promise { // Move back to '>'-readable by deleting and re-adding (simplest approach) // In production, prefer a dead-letter stream instead. await redis.xack(STREAM_KEY, GROUP_NAME, streamEntryId); }, async reclaimStale(maxIdleMs = 60_000): Promise { const result: any = await redis.xautoclaim( STREAM_KEY, GROUP_NAME, consumer, maxIdleMs, '0-0', 'COUNT', '10', ); const entries: [string, string[]][] = Array.isArray(result) ? (result[1] ?? []) : []; return entries.map(([entryId, fields]) => { const obj: Record = {}; for (let i = 0; i < fields.length; i += 2) obj[fields[i]] = fields[i + 1]; return { id: entryId, jobId: obj.jobId, tenantId: obj.tenantId, projectId: obj.projectId, payload: JSON.parse(obj.payload ?? '{}'), enqueuedAt: obj.enqueuedAt, }; }); }, async stats(): Promise<{ pending: number; delivered: number; groupLag: number }> { const info: any[] = await redis.xinfo('GROUPS', STREAM_KEY); // xinfo groups returns array of arrays: [name, ..., pel-count, ..., lag, ...] const group = info.find((g: any) => { const arr = Array.isArray(g) ? g : []; for (let i = 0; i < arr.length; i += 2) { if (arr[i] === 'name' && arr[i + 1] === GROUP_NAME) return true; } return false; }); if (!group) return { pending: 0, delivered: 0, groupLag: 0 }; const flat = Array.isArray(group) ? group : []; const get = (key: string) => { for (let i = 0; i < flat.length - 1; i += 2) { if (flat[i] === key) return Number(flat[i + 1]) || 0; } return 0; }; return { pending: get('pel-count'), delivered: get('delivered-count'), groupLag: get('lag'), }; }, async close(): Promise { await redis.quit().catch(() => {}); }, }; } /** * Check if Redis queue backend is configured. */ export function isRedisQueueEnabled(): boolean { return !!(process.env.REDIS_URL && process.env.JOB_QUEUE_BACKEND === 'redis'); } /** * Convenience: create queue only if enabled, else return null. */ export async function maybeCreateRedisQueue(): Promise { if (!isRedisQueueEnabled()) return null; return createRedisQueue().catch((e) => { console.error('[redis-queue] Failed to connect:', e.message); return null; }); }