/** * BullMQ worker: generate lesson content for a blueprint. * * Idempotent: checks `job` ledger before generating. * On success: writes lesson to Redis cache and updates job status. * On failure: increments attempt count; sets status to 'failed' after processing. */ import { Worker } from 'bullmq'; import { eq } from 'drizzle-orm'; import { db } from '../src/lib/db'; import { jobs, blueprints, concepts, lessons } from '../src/lib/db/schema'; import { generateLesson } from '../src/lib/generation/generate-lesson'; import { verifyLesson } from '../src/lib/verification/verify-content'; import { getLessonResponse } from '../src/lib/db/queries'; import { setCachedLesson } from '../src/lib/cache/lesson'; import { enqueuePromoteBlueprint } from '../src/lib/jobs/queue'; import type { GenerateLessonJobData } from '../src/lib/jobs/queue'; const connection = { url: process.env.REDIS_URL ?? 'redis://localhost:6379', }; const worker = new Worker( 'generate-lesson', async (job) => { const { intentKey, blueprintId, idempotencyKey } = job.data; // 1. Idempotency check — find the job row by idempotency key const [jobRow] = await db .select({ id: jobs.id, status: jobs.status, attempts: jobs.attempts }) .from(jobs) .where(eq(jobs.id, idempotencyKey)) .limit(1); if (jobRow?.status === 'done') { // Already completed — warm cache if needed and return const cached = await getLessonResponse(intentKey); if (cached?.state === 'ready') { await setCachedLesson(intentKey, cached); } return; } // Mark running await db .update(jobs) .set({ status: 'running', attempts: (jobRow?.attempts ?? 0) + 1 }) .where(eq(jobs.id, idempotencyKey)); // 2. Check lesson doesn't already exist const existing = await getLessonResponse(intentKey); if (existing?.state === 'ready') { await setCachedLesson(intentKey, existing); await db.update(jobs).set({ status: 'done' }).where(eq(jobs.id, idempotencyKey)); return; } // 3. Fetch blueprint + concepts const [blueprint] = await db .select() .from(blueprints) .where(eq(blueprints.id, blueprintId)) .limit(1); if (!blueprint) throw new Error(`Blueprint not found: ${blueprintId}`); const conceptRows = await db .select({ id: concepts.id, name: concepts.name, ord: concepts.ord, retrievalCtxRef: concepts.retrievalCtxRef }) .from(concepts) .where(eq(concepts.blueprintId, blueprintId)); if (conceptRows.length === 0) throw new Error(`No concepts for blueprint: ${blueprintId}`); // 4. Generate lesson const { lessonId } = await generateLesson({ blueprintId, lessonOrd: 0, topicTitle: blueprint.title, conceptsToGenerate: conceptRows.map((c) => ({ id: c.id, name: c.name, retrievalCtxRef: c.retrievalCtxRef, })), }); // 5. T1 verification (non-blocking for serve — T2 + promotion handled by promote-blueprint job) await verifyLesson({ lessonId, runT2: false }); // 6. Write to Redis cache const lesson = await getLessonResponse(intentKey); if (lesson?.state === 'ready') { await setCachedLesson(intentKey, lesson); } // 7. Update job ledger await db.update(jobs).set({ status: 'done' }).where(eq(jobs.id, idempotencyKey)); // 8. Enqueue blueprint promotion (T2 verify + misconceptions) const promoteKey = `promote-blueprint:${blueprintId}:v1`; const [existingPromote] = await db .select({ id: jobs.id }) .from(jobs) .where(eq(jobs.idempotencyKey, promoteKey)) .limit(1); if (!existingPromote) { const [promoteJobRow] = await db .insert(jobs) .values({ type: 'promote-blueprint', idempotencyKey: promoteKey, status: 'pending', payloadJson: { blueprintId, intentKey, lessonId }, }) .returning({ id: jobs.id }); await enqueuePromoteBlueprint({ blueprintId, intentKey, idempotencyKey: promoteJobRow.id, }); } }, { connection, concurrency: 2 }, ); worker.on('failed', async (job, err) => { if (job?.data.idempotencyKey) { await db .update(jobs) .set({ status: 'failed' }) .where(eq(jobs.id, job.data.idempotencyKey)); } console.error('[generate-lesson] job failed:', err); }); export default worker;