diff --git a/lib/video/video-platforms.ts b/lib/video/video-platforms.ts index a535ef2..f8d6dfd 100644 --- a/lib/video/video-platforms.ts +++ b/lib/video/video-platforms.ts @@ -6,8 +6,13 @@ import { type TiktokSettings, type YoutubeSettings, } from "@/db/schema"; -import { and, desc, eq } from "drizzle-orm"; +import { and, desc, eq, inArray, sql } from "drizzle-orm"; import { sse } from "@/lib/realtime/sse"; +import { + createVideoPublishingModule, + type VideoPublishJobRecord, + type VideoPublishUploadRecord, +} from "@/lib/video/video-publishing"; const VIDEO_TYPES = new Set(["video/mp4", "video/quicktime", "video/webm"]); @@ -330,105 +335,76 @@ async function publishTiktok(upload: typeof videoUploads.$inferSelect) { return { id: initData.data.publish_id, url: "https://www.tiktok.com/" }; } -export async function publishVideoJob(jobId: string) { - const job = await db.query.videoUploadJobs.findFirst({ - where: (table, { eq }) => eq(table.id, jobId), - }); - if (!job || job.status === "completed") return; - const upload = await db.query.videoUploads.findFirst({ - where: (table, { eq }) => eq(table.id, job.uploadId), - }); - if (!upload) throw new Error("Upload record not found."); - await db - .update(videoUploadJobs) - .set({ - status: "processing", - attempts: job.attempts + 1, - startedAt: new Date(), - error: null, - }) - .where(eq(videoUploadJobs.id, job.id)); - await sse.uploads.pub("update", { - uploadId: upload.id, - jobId: job.id, - platform: job.platform, - status: "processing", - }); - try { - const result = - job.platform === "youtube" - ? await publishYoutube(upload) - : await publishTiktok(upload); - await db - .update(videoUploadJobs) - .set({ - status: "completed", - externalId: result.id, - externalUrl: result.url, - completedAt: new Date(), - }) - .where(eq(videoUploadJobs.id, job.id)); - await sse.uploads.pub("update", { - uploadId: upload.id, - jobId: job.id, - platform: job.platform, - status: "completed", - }); - } catch (error) { - const message = - error instanceof Error ? error.message : "Platform upload failed."; - await db - .update(videoUploadJobs) - .set({ status: "failed", error: message }) - .where(eq(videoUploadJobs.id, job.id)); - await sse.uploads.pub("update", { - uploadId: upload.id, - jobId: job.id, - platform: job.platform, - status: "failed", - }); - throw error; - } -} +type DatabaseVideoUpload = typeof videoUploads.$inferSelect; +type DatabaseVideoJob = typeof videoUploadJobs.$inferSelect; -export async function processDueVideoUploads() { - const due = await db.query.videoUploads.findMany({ - where: (table, { eq, lte, or, isNull }) => - and( - or(eq(table.status, "scheduled"), eq(table.status, "partial")), - or(isNull(table.scheduledAt), lte(table.scheduledAt, new Date())), - ), - orderBy: [desc(videoUploads.createdAt)], - }); - for (const upload of due) { - const jobs = await db.query.videoUploadJobs.findMany({ - where: (table, { eq, and, ne }) => - and(eq(table.uploadId, upload.id), ne(table.status, "completed")), - }); - for (const job of jobs) { - await publishVideoJob(job.id).catch((error) => { - const message = - error instanceof Error ? error.message : "Platform upload failed."; - console.error(`[${job.platform}] ${job.id}: ${message}`); - }); - } - const remaining = await db.query.videoUploadJobs.findMany({ - where: (table, { eq, and, ne }) => - and(eq(table.uploadId, upload.id), ne(table.status, "completed")), - }); - const nextStatus = - remaining.length === 0 - ? "completed" - : remaining.some((job) => job.status === "completed") - ? "partial" - : "failed"; - await db - .update(videoUploads) - .set({ status: nextStatus, updatedAt: new Date() }) - .where(eq(videoUploads.id, upload.id)); - await sse.uploads.pub("update", { - uploadId: upload.id, - status: nextStatus, - }); - } -} +const videoPublishing = createVideoPublishingModule< + DatabaseVideoUpload & VideoPublishUploadRecord, + DatabaseVideoJob & VideoPublishJobRecord +>({ + store: { + claimJob: async (jobId, startedAt) => { + const [job] = await db + .update(videoUploadJobs) + .set({ + status: "processing", + attempts: sql`${videoUploadJobs.attempts} + 1`, + startedAt, + error: null, + }) + .where( + and( + eq(videoUploadJobs.id, jobId), + inArray(videoUploadJobs.status, ["pending", "failed"]), + ), + ) + .returning(); + return job ?? null; + }, + findUpload: async (uploadId) => (await db.query.videoUploads.findFirst({ + where: (table, { eq }) => eq(table.id, uploadId), + })) ?? null, + markCompleted: async (job, result, completedAt) => { + await db + .update(videoUploadJobs) + .set({ + status: "completed", + externalId: result.id, + externalUrl: result.url, + completedAt, + }) + .where(eq(videoUploadJobs.id, job.id)); + }, + markFailed: async (job, message) => { + await db + .update(videoUploadJobs) + .set({ status: "failed", error: message }) + .where(eq(videoUploadJobs.id, job.id)); + }, + findDueUploads: (now) => db.query.videoUploads.findMany({ + where: (table, { eq, lte, or, isNull }) => + and( + or(eq(table.status, "scheduled"), eq(table.status, "partial")), + or(isNull(table.scheduledAt), lte(table.scheduledAt, now)), + ), + orderBy: [desc(videoUploads.createdAt)], + }), + findJobs: (uploadId) => db.query.videoUploadJobs.findMany({ + where: (table, { eq }) => eq(table.uploadId, uploadId), + }), + updateUploadStatus: async (upload, status, updatedAt) => { + await db + .update(videoUploads) + .set({ status, updatedAt }) + .where(eq(videoUploads.id, upload.id)); + }, + }, + publishers: { + youtube: publishYoutube, + tiktok: publishTiktok, + }, + notify: (event) => sse.uploads.pub("update", event), +}); + +export const publishVideoJob = videoPublishing.publishJob; +export const processDueVideoUploads = videoPublishing.processDue; diff --git a/lib/video/video-publishing.test.ts b/lib/video/video-publishing.test.ts new file mode 100644 index 0000000..cc0bad1 --- /dev/null +++ b/lib/video/video-publishing.test.ts @@ -0,0 +1,180 @@ +import { describe, expect, test } from "bun:test"; +import { + createVideoPublishingModule, + type VideoPublishJobRecord, + type VideoPublishUploadRecord, + type VideoPublishingEvent, +} from "@/lib/video/video-publishing"; + +type TestUpload = VideoPublishUploadRecord & { title: string }; +type TestJob = VideoPublishJobRecord; + +function fixture() { + const upload: TestUpload = { + id: "upload-1", + title: "Test upload", + status: "processing", + scheduledAt: null, + createdAt: new Date("2026-08-09T00:00:00Z"), + }; + const job: TestJob = { + id: "job-1", + uploadId: upload.id, + platform: "youtube", + status: "pending", + attempts: 0, + }; + const events: VideoPublishingEvent[] = []; + const updates: string[] = []; + const jobs = [job]; + const publishing = createVideoPublishingModule({ + store: { + claimJob: async (_jobId, startedAt) => { + if (job.status === "completed" || job.status === "processing") return null; + job.status = "processing"; + job.attempts += 1; + void startedAt; + return job; + }, + findUpload: async () => upload, + markCompleted: async (nextJob) => { nextJob.status = "completed"; }, + markFailed: async (nextJob) => { nextJob.status = "failed"; }, + findDueUploads: async () => [upload], + findJobs: async () => jobs, + updateUploadStatus: async (_upload, status) => { + updates.push(status); + }, + }, + publishers: { + youtube: async (nextUpload) => ({ id: nextUpload.id, url: "https://youtube.test/video" }), + tiktok: async (nextUpload) => ({ id: nextUpload.id, url: "https://tiktok.test/video" }), + }, + notify: async (event) => { events.push(event); }, + now: () => new Date("2026-08-09T00:01:00Z"), + }); + + return { events, job, jobs, publishing, updates, upload }; +} + +describe("video publishing job module", () => { + test("keeps provider result and lifecycle notifications separate", async () => { + const { events, job, publishing } = fixture(); + + await publishing.publishJob(job.id); + + expect(job).toMatchObject({ status: "completed", attempts: 1 }); + expect(events.map((event) => event.status)).toEqual(["processing", "completed"]); + }); + + test("marks provider failures without hiding the original error", async () => { + const { events, job } = fixture(); + const publishing = createVideoPublishingModule({ + store: { + claimJob: async () => { + if (job.status === "processing" || job.status === "completed") return null; + job.status = "processing"; + return job; + }, + findUpload: async () => ({ + id: "upload-1", + title: "Test upload", + status: "processing", + scheduledAt: null, + createdAt: new Date(), + }), + markCompleted: async () => {}, + markFailed: async (nextJob) => { nextJob.status = "failed"; }, + findDueUploads: async () => [], + findJobs: async () => [], + updateUploadStatus: async () => {}, + }, + publishers: { + youtube: async () => { throw new Error("provider offline"); }, + tiktok: async () => ({ id: "id", url: "url" }), + }, + notify: async (event) => { events.push(event); }, + }); + + await expect(publishing.publishJob(job.id)).rejects.toThrow("provider offline"); + expect(job.status).toBe("failed"); + expect(events.at(-1)?.status).toBe("failed"); + }); + + test("allows only one caller to claim a job", async () => { + const { job, publishing } = fixture(); + + await Promise.all([publishing.publishJob(job.id), publishing.publishJob(job.id)]); + + expect(job).toMatchObject({ status: "completed", attempts: 1 }); + }); + + test("derives partial status from all platform jobs", async () => { + const { updates, upload } = fixture(); + const jobs: TestJob[] = [ + { + id: "job-1", + uploadId: upload.id, + platform: "youtube", + status: "completed", + attempts: 1, + }, + { + id: "job-2", + uploadId: upload.id, + platform: "tiktok", + status: "processing", + attempts: 1, + }, + ]; + const publishing = createVideoPublishingModule({ + store: { + claimJob: async () => null, + findUpload: async () => upload, + markCompleted: async () => {}, + markFailed: async () => {}, + findDueUploads: async () => [upload], + findJobs: async () => jobs, + updateUploadStatus: async (_upload, status) => { updates.push(status); }, + }, + publishers: { + youtube: async () => ({ id: "video", url: "url" }), + tiktok: async () => ({ id: "video", url: "url" }), + }, + notify: async () => {}, + }); + + await publishing.processDue(); + + expect(updates).toEqual(["partial"]); + }); + + test("does not fail publishing when lifecycle notification fails", async () => { + const { job, upload } = fixture(); + const notificationErrors: unknown[] = []; + const publishing = createVideoPublishingModule({ + store: { + claimJob: async () => { + job.status = "processing"; + return job; + }, + findUpload: async () => upload, + markCompleted: async (nextJob) => { nextJob.status = "completed"; }, + markFailed: async (nextJob) => { nextJob.status = "failed"; }, + findDueUploads: async () => [], + findJobs: async () => [job], + updateUploadStatus: async () => {}, + }, + publishers: { + youtube: async () => ({ id: "video", url: "url" }), + tiktok: async () => ({ id: "video", url: "url" }), + }, + notify: async () => { throw new Error("redis offline"); }, + reportNotificationFailure: (error) => notificationErrors.push(error), + }); + + await publishing.publishJob(job.id); + + expect(job.status).toBe("completed"); + expect(notificationErrors).toHaveLength(2); + }); +}); diff --git a/lib/video/video-publishing.ts b/lib/video/video-publishing.ts new file mode 100644 index 0000000..02937a0 --- /dev/null +++ b/lib/video/video-publishing.ts @@ -0,0 +1,142 @@ +export type VideoPlatform = "youtube" | "tiktok"; +export type VideoJobStatus = + | "pending" + | "processing" + | "completed" + | "failed"; +export type VideoUploadStatus = "draft" | "scheduled" | "processing" | "partial" | "completed" | "failed"; + +export type VideoPublishResult = { + id: string; + url: string; +}; + +export type VideoPublishJobRecord = { + id: string; + uploadId: string; + platform: VideoPlatform; + status: VideoJobStatus; + attempts: number; +}; + +export type VideoPublishUploadRecord = { + id: string; + status: VideoUploadStatus; + scheduledAt: Date | null; + createdAt: Date; +}; + +export type VideoPublishingEvent = { + uploadId: string; + jobId?: string; + platform?: VideoPlatform; + status: Exclude | VideoJobStatus; +}; + +export type VideoPublishingStore< + TUpload extends VideoPublishUploadRecord, + TJob extends VideoPublishJobRecord, +> = { + claimJob: (jobId: string, startedAt: Date) => Promise; + findUpload: (uploadId: string) => Promise; + markCompleted: (job: TJob, result: VideoPublishResult, completedAt: Date) => Promise; + markFailed: (job: TJob, message: string) => Promise; + findDueUploads: (now: Date) => Promise; + findJobs: (uploadId: string) => Promise; + updateUploadStatus: (upload: TUpload, status: VideoUploadStatus, updatedAt: Date) => Promise; +}; + +export type VideoPublishingModule = { + publishJob: (jobId: string) => Promise; + processDue: () => Promise; +}; + +export function createVideoPublishingModule< + TUpload extends VideoPublishUploadRecord, + TJob extends VideoPublishJobRecord, +>({ + store, + publishers, + notify, + reportNotificationFailure = (error) => { + console.error("Failed to publish video lifecycle notification", error); + }, + now = () => new Date(), +}: { + store: VideoPublishingStore; + publishers: Record Promise>; + notify: (event: VideoPublishingEvent) => Promise; + reportNotificationFailure?: (error: unknown) => void; + now?: () => Date; +}): VideoPublishingModule { + async function sendNotification(event: VideoPublishingEvent) { + try { + await notify(event); + } catch (error) { + reportNotificationFailure(error); + } + } + + async function publishJob(jobId: string) { + const job = await store.claimJob(jobId, now()); + if (!job) return; + + const upload = await store.findUpload(job.uploadId); + if (!upload) throw new Error("Upload record not found."); + + await sendNotification({ + uploadId: upload.id, + jobId: job.id, + platform: job.platform, + status: "processing", + }); + + try { + const result = await publishers[job.platform](upload); + await store.markCompleted(job, result, now()); + await sendNotification({ + uploadId: upload.id, + jobId: job.id, + platform: job.platform, + status: "completed", + }); + } catch (error) { + const message = error instanceof Error ? error.message : "Platform upload failed."; + await store.markFailed(job, message); + await sendNotification({ + uploadId: upload.id, + jobId: job.id, + platform: job.platform, + status: "failed", + }); + throw error; + } + } + + async function processDue() { + const dueUploads = await store.findDueUploads(now()); + + for (const upload of dueUploads) { + const jobs = await store.findJobs(upload.id); + for (const job of jobs.filter((item) => item.status !== "completed")) { + await publishJob(job.id).catch((error) => { + const message = error instanceof Error ? error.message : "Platform upload failed."; + console.error(`[${job.platform}] ${job.id}: ${message}`); + }); + } + + const currentJobs = await store.findJobs(upload.id); + const nextStatus: VideoUploadStatus = currentJobs.every( + (job) => job.status === "completed", + ) + ? "completed" + : currentJobs.some((job) => job.status === "completed") + ? "partial" + : "failed"; + await store.updateUploadStatus(upload, nextStatus, now()); + await sendNotification({ uploadId: upload.id, status: nextStatus }); + } + } + + return { publishJob, processDue }; +}