/** * 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, and } 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 { appendStreamSegment, markStreamDone } from '../src/lib/cache/stream'; 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, locale = 'en', depth = 'standard', ageGroup = 'adult', difficultyLevel = 1, ord = 0 } = 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, undefined, locale, depth as import('../src/lib/generation/depth').LessonDepth, ageGroup as import('../src/schemas/preferences').AgeGroup); if (cached?.state === 'ready') { await setCachedLesson(intentKey, locale, cached, depth as import('../src/lib/generation/depth').LessonDepth, ageGroup as import('../src/schemas/preferences').AgeGroup); } return; } // Mark running await db .update(jobs) .set({ status: 'running', attempts: (jobRow?.attempts ?? 0) + 1 }) .where(eq(jobs.id, idempotencyKey)); // 2. Check lesson at this ord doesn't already exist const [existingLesson] = await db .select({ id: lessons.id }) .from(lessons) .where( and( eq(lessons.blueprintId, blueprintId), eq(lessons.locale, locale), eq(lessons.depth, depth), eq(lessons.ageGroup, ageGroup), eq(lessons.ord, ord), ), ) .limit(1); if (existingLesson) { const cached = await getLessonResponse(intentKey, undefined, locale, depth as import('../src/lib/generation/depth').LessonDepth, ageGroup as import('../src/schemas/preferences').AgeGroup); if (cached?.state === 'ready') { await setCachedLesson(intentKey, locale, cached, depth as import('../src/lib/generation/depth').LessonDepth, ageGroup as import('../src/schemas/preferences').AgeGroup); } await db.update(jobs).set({ status: 'done' }).where(eq(jobs.id, idempotencyKey)); return; } // 3. Fetch blueprint + concepts at this ord const [blueprint] = await db .select() .from(blueprints) .where(eq(blueprints.id, blueprintId)) .limit(1); if (!blueprint) throw new Error(`Blueprint not found: ${blueprintId}`); // For ord=0 use all concepts (backward compat); for ord>0 use concepts at that ord. const conceptRows = await db .select({ id: concepts.id, name: concepts.name, ord: concepts.ord, retrievalCtxRef: concepts.retrievalCtxRef }) .from(concepts) .where( ord === 0 ? eq(concepts.blueprintId, blueprintId) : and(eq(concepts.blueprintId, blueprintId), eq(concepts.ord, ord)), ); if (conceptRows.length === 0) throw new Error(`No concepts at ord=${ord} for blueprint: ${blueprintId}`); // 4. Generate lesson — onSegment pushes each segment to Redis as it's persisted, // so SSE clients receive progressive delivery without waiting for the full batch. const { lessonId } = await generateLesson({ blueprintId, lessonOrd: ord, topicTitle: blueprint.title, conceptsToGenerate: conceptRows.map((c) => ({ id: c.id, name: c.name, retrievalCtxRef: c.retrievalCtxRef, })), locale, depth: depth as import('../src/lib/generation/depth').LessonDepth, ageGroup: ageGroup as import('../src/schemas/preferences').AgeGroup, difficultyLevel: difficultyLevel as 1 | 2 | 3, onSegment: (seg) => appendStreamSegment( intentKey, locale, seg, depth as import('../src/lib/generation/depth').LessonDepth, ageGroup as import('../src/schemas/preferences').AgeGroup, ord, ), }); // Signal SSE clients that the stream is complete. await markStreamDone( intentKey, locale, depth as import('../src/lib/generation/depth').LessonDepth, ageGroup as import('../src/schemas/preferences').AgeGroup, ord, ); // 5. T1 verification (non-blocking for serve — T2 + promotion handled by promote-blueprint job) await verifyLesson({ lessonId, runT2: false }); // 6. Write to Redis cache (position-keyed) const lesson = await getLessonResponse(intentKey, undefined, locale, depth as import('../src/lib/generation/depth').LessonDepth, ageGroup as import('../src/schemas/preferences').AgeGroup); if (lesson?.state === 'ready') { await setCachedLesson(intentKey, locale, lesson, depth as import('../src/lib/generation/depth').LessonDepth, ageGroup as import('../src/schemas/preferences').AgeGroup, ord); } // 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;