143 lines
4.2 KiB
TypeScript
143 lines
4.2 KiB
TypeScript
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 };
|
|
}
|