fix(video): make publishing transitions safe
This commit is contained in:
+78
-102
@@ -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;
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
@@ -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<VideoUploadStatus, "draft"> | VideoJobStatus;
|
||||
};
|
||||
|
||||
export type VideoPublishingStore<
|
||||
TUpload extends VideoPublishUploadRecord,
|
||||
TJob extends VideoPublishJobRecord,
|
||||
> = {
|
||||
claimJob: (jobId: string, startedAt: Date) => Promise<TJob | null>;
|
||||
findUpload: (uploadId: string) => Promise<TUpload | null>;
|
||||
markCompleted: (job: TJob, result: VideoPublishResult, completedAt: Date) => Promise<void>;
|
||||
markFailed: (job: TJob, message: string) => Promise<void>;
|
||||
findDueUploads: (now: Date) => Promise<TUpload[]>;
|
||||
findJobs: (uploadId: string) => Promise<TJob[]>;
|
||||
updateUploadStatus: (upload: TUpload, status: VideoUploadStatus, updatedAt: Date) => Promise<void>;
|
||||
};
|
||||
|
||||
export type VideoPublishingModule = {
|
||||
publishJob: (jobId: string) => Promise<void>;
|
||||
processDue: () => Promise<void>;
|
||||
};
|
||||
|
||||
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<TUpload, TJob>;
|
||||
publishers: Record<VideoPlatform, (upload: TUpload) => Promise<VideoPublishResult>>;
|
||||
notify: (event: VideoPublishingEvent) => Promise<unknown>;
|
||||
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 };
|
||||
}
|
||||
Reference in New Issue
Block a user