112 lines
3.1 KiB
TypeScript
112 lines
3.1 KiB
TypeScript
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 } from "./lib/ffmpeg/encode";
|
|
import { getRedisConnection } from "./lib/queue/client";
|
|
import { incrementQuota } from "./lib/quota";
|
|
import { getJobDir } from "./lib/storage";
|
|
import type { VideoJobData } from "./lib/types";
|
|
import { uploadToYouTube } from "./lib/youtube/upload";
|
|
|
|
async function updateJobStatus(jobId: string) {
|
|
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;
|
|
}
|
|
|
|
await prisma.job.update({
|
|
where: { id: jobId },
|
|
data: {
|
|
status,
|
|
completedAt: completed + failed === total ? new Date() : null,
|
|
},
|
|
});
|
|
|
|
}
|
|
|
|
async function processJobItem(data: VideoJobData) {
|
|
const item = await prisma.jobItem.findUniqueOrThrow({
|
|
where: { id: data.jobItemId },
|
|
include: { job: true },
|
|
});
|
|
|
|
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.job.imagePath,
|
|
audioPath: item.audioPath,
|
|
outputPath,
|
|
resolution: item.resolution,
|
|
includeWatermark: item.includeWatermark,
|
|
});
|
|
|
|
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,
|
|
},
|
|
});
|
|
|
|
await incrementQuota(data.userId, 1);
|
|
await cleanupFiles([outputPath]);
|
|
} catch (err) {
|
|
const message = err instanceof Error ? err.message : "Unknown error";
|
|
await prisma.jobItem.update({
|
|
where: { id: item.id },
|
|
data: { status: JobItemStatus.FAILED, error: message },
|
|
});
|
|
throw err;
|
|
} finally {
|
|
await updateJobStatus(data.jobId);
|
|
}
|
|
}
|
|
|
|
const worker = new Worker<VideoJobData>(
|
|
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("s2yt worker started, waiting for jobs...");
|