import { Worker } from "bullmq"; import path from "path"; import { JobItemStatus, JobStatus } from "@prisma/client"; import { QUEUE_NAME } from "../lib/constants"; import { prisma } from "../lib/db"; import { cleanupFiles, encodeVideo, getFfmpegPath } from "../lib/ffmpeg/encode"; import { getRedisConnection } from "../lib/queue/client"; import { getJobDir } from "../lib/storage"; import type { VideoJobData } from "../lib/types"; import { formatYouTubeErrorForUser } from "../lib/youtube/errors"; import { uploadToYouTube } from "../lib/youtube/upload"; import { deliverJobWebhook } from "../lib/webhooks"; async function updateJobStatus(jobId: string) { const job = await prisma.job.findUniqueOrThrow({ where: { id: jobId }, select: { webhookUrl: true }, }); const items = await prisma.jobItem.findMany({ where: { jobId } }); const completed = items.filter((i) => i.status === JobItemStatus.COMPLETED).length; const failed = items.filter((i) => i.status === JobItemStatus.FAILED).length; const total = items.length; let status: JobStatus = JobStatus.PROCESSING; if (completed + failed === total) { if (completed === total) status = JobStatus.COMPLETED; else if (failed === total) status = JobStatus.FAILED; else status = JobStatus.PARTIAL; } const completedAt = completed + failed === total ? new Date() : null; await prisma.job.update({ where: { id: jobId }, data: { status, completedAt, }, }); if (completedAt) { const event = status === JobStatus.COMPLETED ? "job.completed" : status === JobStatus.FAILED ? "job.failed" : "job.partial"; await deliverJobWebhook(job.webhookUrl, { event, jobId, status, completedAt: completedAt.toISOString(), }); } } async function processJobItem(data: VideoJobData) { const item = await prisma.jobItem.findUniqueOrThrow({ where: { id: data.jobItemId }, include: { job: true }, }); if (item.job.userId !== data.userId || item.jobId !== data.jobId) { const message = "Job ownership mismatch"; console.error(`[security] ${message} for item ${item.id}`); await prisma.jobItem.update({ where: { id: item.id }, data: { status: JobItemStatus.FAILED, error: message }, }); await updateJobStatus(item.jobId); throw new Error(message); } const jobDir = getJobDir(data.userId, data.jobId); const outputPath = path.join(jobDir, `${item.id}.mp4`); try { await prisma.jobItem.update({ where: { id: item.id }, data: { status: JobItemStatus.ENCODING }, }); await prisma.job.update({ where: { id: data.jobId }, data: { status: JobStatus.PROCESSING }, }); await encodeVideo({ imagePath: item.itemImagePath || item.job.imagePath, audioPath: item.audioPath, outputPath, resolution: item.resolution, includeWatermark: item.includeWatermark, songTitle: item.songTitle || item.title, artist: item.artist, layout: { template: (item.layoutTemplate as | "COVER_LEFT_TEXT_RIGHT" | "COVER_TOP_TEXT_BOTTOM" | "COVER_RIGHT_TEXT_LEFT" | "CENTERED_COMPACT" | "LOWER_LEFT_COVER_TEXT" | "LOWER_RIGHT_COVER_TEXT" | null) ?? null, blurFill: item.blurFill ?? false, blurAmount: item.blurAmount ?? 55, blurOpacity: item.blurOpacity ?? 100, textPadding: item.textPadding ?? 48, titleArtistGap: item.titleArtistGap ?? 10, titleBold: item.titleBold ?? true, textOffsetX: item.textOffsetX ?? 0, textOffsetY: item.textOffsetY ?? 0, }, backgroundImagePath: item.backgroundImagePath, watermark: { mode: (item.watermarkMode as "none" | "default" | "text" | "logo") || "default", text: item.watermarkText, logoPath: item.watermarkLogoPath, fontKey: (item.watermarkFontKey as | "system" | "custom" | "inter" | "montserrat" | "roboto" | "oswald" | "playfair") || "system", fontPath: item.watermarkFontPath, position: (item.watermarkPosition as | "top-left" | "top-right" | "bottom-left" | "bottom-right" | "center") || "bottom-right", offsetX: item.watermarkOffsetX ?? 20, offsetY: item.watermarkOffsetY ?? 20, }, }); await prisma.jobItem.update({ where: { id: item.id }, data: { status: JobItemStatus.UPLOADING, outputPath }, }); const youtubeVideoId = await uploadToYouTube(data.userId, outputPath, item); await prisma.jobItem.update({ where: { id: item.id }, data: { status: JobItemStatus.COMPLETED, youtubeVideoId, error: null, }, }); // Track successful renders for instance statistics. await prisma.user.update({ where: { id: data.userId }, data: { createdVideoCount: { increment: 1 } }, }); await deliverJobWebhook(item.job.webhookUrl, { event: "job.item.completed", jobId: data.jobId, status: "COMPLETED", itemId: item.id, youtubeVideoId, }); await cleanupFiles([outputPath]); } catch (err) { const message = formatYouTubeErrorForUser(err); console.error(`[youtube] Job item ${item.id} ("${item.title}") failed: ${message}`); await prisma.jobItem.update({ where: { id: item.id }, data: { status: JobItemStatus.FAILED, error: message }, }); await deliverJobWebhook(item.job.webhookUrl, { event: "job.item.failed", jobId: data.jobId, status: "FAILED", itemId: item.id, error: message, }); throw new Error(message); } finally { await updateJobStatus(data.jobId); } } const worker = new Worker( QUEUE_NAME, async (job) => { await processJobItem(job.data); }, { connection: getRedisConnection(), concurrency: 2, }, ); worker.on("completed", (job) => { console.log(`Job item ${job.data.jobItemId} completed`); }); worker.on("failed", (job, err) => { console.error(`Job item ${job?.data.jobItemId} failed:`, err.message); }); console.log(`Songs2VID worker started (ffmpeg: ${getFfmpegPath()}), waiting for jobs...`);