/** * llmindex.io — BullMQ workers: eval batches, duel scheduling, IRT refit triggers * Author: Simon-Pierre Boucher * Contact: contact@spboucher.ai * License: Proprietary — © Simon-Pierre Boucher, all rights reserved */ import { Worker, type Job } from 'bullmq'; import IORedis from 'ioredis'; import { z } from 'zod'; import { redisUrl } from './env'; import { runEvalBatch } from './eval-runner'; import { runDuelBatch } from './duel-runner'; import { runRefit } from './refit'; import { DOMAINS, DUEL_DOMAINS } from '@llmindex/scoring'; export const QUEUES = { eval: 'llmindex-eval', duel: 'llmindex-duel', refit: 'llmindex-refit', } as const; const evalJob = z.object({ modelSlug: z.string().min(1), domain: z.enum(DOMAINS), n: z.number().int().min(1).max(5000), kSamples: z.number().int().min(1).max(10).optional(), seed: z.string().optional(), }); const duelJob = z.object({ domain: z.enum(DUEL_DOMAINS as [string, ...string[]]), pairs: z.number().int().min(1).max(5000), seed: z.string().optional(), }); function connection() { return new IORedis(redisUrl(), { maxRetriesPerRequest: null }); } function main(): void { const opts = { connection: connection(), concurrency: 1 }; new Worker( QUEUES.eval, async (job: Job) => { const params = evalJob.parse(job.data); return runEvalBatch({ ...params, domain: params.domain }); }, opts, ); new Worker( QUEUES.duel, async (job: Job) => { const params = duelJob.parse(job.data); return runDuelBatch({ domain: params.domain as 'writing' | 'safety_refusal_quality', pairs: params.pairs, seed: params.seed, }); }, opts, ); new Worker(QUEUES.refit, async () => runRefit(), opts); console.log( `[worker] listening on queues ${Object.values(QUEUES).join(', ')} via ${redisUrl().replace(/\/\/.*@/, '//***@')}`, ); } main();