435 lines
14 KiB
TypeScript
435 lines
14 KiB
TypeScript
import { readFile } from "node:fs/promises";
|
|
import { db } from "@/db";
|
|
import {
|
|
videoUploadJobs,
|
|
videoUploads,
|
|
type TiktokSettings,
|
|
type YoutubeSettings,
|
|
} from "@/db/schema";
|
|
import { and, desc, eq } from "drizzle-orm";
|
|
import { sse } from "@/lib/realtime/sse";
|
|
|
|
const VIDEO_TYPES = new Set(["video/mp4", "video/quicktime", "video/webm"]);
|
|
|
|
export type LatestVideo = {
|
|
id: string;
|
|
title: string;
|
|
thumbnailUrl: string;
|
|
url: string;
|
|
createdAt: number;
|
|
platform: "youtube" | "tiktok";
|
|
};
|
|
|
|
export function validateVideoFile(file: File) {
|
|
if (!VIDEO_TYPES.has(file.type))
|
|
throw new Error("Use an MP4, MOV, or WebM video.");
|
|
if (file.size <= 0) throw new Error("Video cannot be empty.");
|
|
}
|
|
|
|
export function validateThumbnailFile(file: File) {
|
|
if (!file.type.startsWith("image/"))
|
|
throw new Error("Thumbnail must be an image.");
|
|
if (file.size <= 0 || file.size > 10 * 1024 * 1024)
|
|
throw new Error("Thumbnail must be smaller than 10 MB.");
|
|
}
|
|
|
|
export async function getLatestYoutubeVideos(): Promise<LatestVideo[]> {
|
|
const apiKey = process.env.YOUTUBE_API_KEY;
|
|
const channelId = process.env.YOUTUBE_CHANNEL_ID;
|
|
if (!apiKey || !channelId) return [];
|
|
const channels = await fetch(
|
|
`https://www.googleapis.com/youtube/v3/channels?part=contentDetails&id=${encodeURIComponent(channelId)}&key=${encodeURIComponent(apiKey)}`,
|
|
{ cache: "no-store" },
|
|
);
|
|
if (!channels.ok) throw new Error("YouTube channel lookup failed.");
|
|
const channelData = (await channels.json()) as {
|
|
items?: Array<{
|
|
contentDetails?: { relatedPlaylists?: { uploads?: string } };
|
|
}>;
|
|
};
|
|
const playlistId =
|
|
channelData.items?.[0]?.contentDetails?.relatedPlaylists?.uploads;
|
|
if (!playlistId) return [];
|
|
const response = await fetch(
|
|
`https://www.googleapis.com/youtube/v3/playlistItems?part=snippet,contentDetails&playlistId=${encodeURIComponent(playlistId)}&maxResults=10&key=${encodeURIComponent(apiKey)}`,
|
|
{ cache: "no-store" },
|
|
);
|
|
if (!response.ok) throw new Error("YouTube video list failed.");
|
|
const data = (await response.json()) as {
|
|
items?: Array<{
|
|
contentDetails?: { videoId?: string };
|
|
snippet?: {
|
|
title?: string;
|
|
publishedAt?: string;
|
|
thumbnails?: Record<string, { url: string }>;
|
|
};
|
|
}>;
|
|
};
|
|
return (data.items ?? []).flatMap((item) => {
|
|
const id = item.contentDetails?.videoId;
|
|
if (!id) return [];
|
|
return [
|
|
{
|
|
id,
|
|
title: item.snippet?.title ?? "Untitled",
|
|
thumbnailUrl:
|
|
item.snippet?.thumbnails?.high?.url ??
|
|
item.snippet?.thumbnails?.default?.url ??
|
|
"",
|
|
url: `https://www.youtube.com/watch?v=${id}`,
|
|
createdAt: Date.parse(item.snippet?.publishedAt ?? "") || 0,
|
|
platform: "youtube" as const,
|
|
},
|
|
];
|
|
});
|
|
}
|
|
|
|
async function getTiktokAccessToken() {
|
|
const refreshToken = process.env.TIKTOK_REFRESH_TOKEN;
|
|
const clientKey = process.env.TIKTOK_CLIENT_KEY;
|
|
const clientSecret = process.env.TIKTOK_CLIENT_SECRET;
|
|
if (!refreshToken || !clientKey || !clientSecret)
|
|
return process.env.TIKTOK_ACCESS_TOKEN ?? null;
|
|
const response = await fetch("https://open.tiktokapis.com/v2/oauth/token/", {
|
|
method: "POST",
|
|
headers: { "Content-Type": "application/x-www-form-urlencoded" },
|
|
body: new URLSearchParams({
|
|
client_key: clientKey,
|
|
client_secret: clientSecret,
|
|
grant_type: "refresh_token",
|
|
refresh_token: refreshToken,
|
|
}),
|
|
});
|
|
if (!response.ok) throw new Error("TikTok token refresh failed.");
|
|
const data = (await response.json()) as {
|
|
access_token?: string;
|
|
refresh_token?: string;
|
|
};
|
|
return data.access_token ?? null;
|
|
}
|
|
|
|
export async function getLatestTiktokVideos(): Promise<LatestVideo[]> {
|
|
const token = await getTiktokAccessToken();
|
|
if (!token) return [];
|
|
const response = await fetch(
|
|
"https://open.tiktokapis.com/v2/video/list/?fields=id,title,cover_image_url,share_url,create_time",
|
|
{
|
|
method: "POST",
|
|
headers: {
|
|
Authorization: `Bearer ${token}`,
|
|
"Content-Type": "application/json",
|
|
},
|
|
body: JSON.stringify({ max_count: 10 }),
|
|
cache: "no-store",
|
|
},
|
|
);
|
|
if (!response.ok) throw new Error("TikTok video list failed.");
|
|
const data = (await response.json()) as {
|
|
data?: {
|
|
videos?: Array<{
|
|
id: string;
|
|
title?: string;
|
|
cover_image_url?: string;
|
|
share_url?: string;
|
|
create_time?: number;
|
|
}>;
|
|
};
|
|
};
|
|
return (data.data?.videos ?? []).map((video) => ({
|
|
id: video.id,
|
|
title: video.title ?? "Untitled",
|
|
thumbnailUrl: video.cover_image_url ?? "",
|
|
url: video.share_url ?? `https://www.tiktok.com/`,
|
|
createdAt: video.create_time ?? 0,
|
|
platform: "tiktok" as const,
|
|
}));
|
|
}
|
|
|
|
async function youtubeAccessToken() {
|
|
const refreshToken = process.env.YOUTUBE_REFRESH_TOKEN;
|
|
const clientId = process.env.YOUTUBE_CLIENT_ID;
|
|
const clientSecret = process.env.YOUTUBE_CLIENT_SECRET;
|
|
if (!refreshToken || !clientId || !clientSecret)
|
|
throw new Error("YouTube OAuth is not configured.");
|
|
const response = await fetch("https://oauth2.googleapis.com/token", {
|
|
method: "POST",
|
|
headers: { "Content-Type": "application/x-www-form-urlencoded" },
|
|
body: new URLSearchParams({
|
|
client_id: clientId,
|
|
client_secret: clientSecret,
|
|
refresh_token: refreshToken,
|
|
grant_type: "refresh_token",
|
|
}),
|
|
});
|
|
if (!response.ok) throw new Error("YouTube token refresh failed.");
|
|
const data = (await response.json()) as { access_token?: string };
|
|
if (!data.access_token)
|
|
throw new Error("YouTube did not return an access token.");
|
|
return data.access_token;
|
|
}
|
|
|
|
async function publishYoutube(upload: typeof videoUploads.$inferSelect) {
|
|
const token = await youtubeAccessToken();
|
|
const settings = upload.youtubeSettings as YoutubeSettings;
|
|
const body = {
|
|
snippet: {
|
|
title: upload.title,
|
|
description: upload.description,
|
|
tags: settings.tags,
|
|
categoryId: settings.categoryId ?? "22",
|
|
},
|
|
status: {
|
|
privacyStatus: upload.scheduledAt
|
|
? "private"
|
|
: (settings.privacyStatus ?? "public"),
|
|
publishAt: upload.scheduledAt?.toISOString(),
|
|
license: settings.license ?? "youtube",
|
|
selfDeclaredMadeForKids: settings.madeForKids ?? false,
|
|
containsSyntheticMedia: settings.containsSyntheticMedia ?? false,
|
|
},
|
|
};
|
|
const init = await fetch(
|
|
"https://www.googleapis.com/upload/youtube/v3/videos?uploadType=resumable&part=snippet,status",
|
|
{
|
|
method: "POST",
|
|
headers: {
|
|
Authorization: `Bearer ${token}`,
|
|
"Content-Type": "application/json; charset=UTF-8",
|
|
"X-Upload-Content-Type": upload.mimeType,
|
|
"X-Upload-Content-Length": String(upload.size),
|
|
},
|
|
body: JSON.stringify(body),
|
|
},
|
|
);
|
|
if (!init.ok)
|
|
throw new Error(`YouTube upload initialization failed (${init.status}).`);
|
|
const location = init.headers.get("location");
|
|
if (!location) throw new Error("YouTube did not return an upload URL.");
|
|
const video = await readFile(upload.videoPath);
|
|
const result = await fetch(location, {
|
|
method: "PUT",
|
|
headers: {
|
|
Authorization: `Bearer ${token}`,
|
|
"Content-Type": upload.mimeType,
|
|
"Content-Length": String(video.byteLength),
|
|
},
|
|
body: video,
|
|
});
|
|
if (!result.ok) throw new Error(`YouTube upload failed (${result.status}).`);
|
|
const data = (await result.json()) as { id?: string };
|
|
if (!data.id) throw new Error("YouTube did not return a video ID.");
|
|
if (upload.thumbnailPath) {
|
|
const thumbnail = await readFile(upload.thumbnailPath);
|
|
const thumb = await fetch(
|
|
`https://www.googleapis.com/upload/youtube/v3/thumbnails/set?videoId=${encodeURIComponent(data.id)}`,
|
|
{
|
|
method: "POST",
|
|
headers: {
|
|
Authorization: `Bearer ${token}`,
|
|
"Content-Type": "image/jpeg",
|
|
"Content-Length": String(thumbnail.byteLength),
|
|
},
|
|
body: thumbnail,
|
|
},
|
|
);
|
|
if (!thumb.ok)
|
|
throw new Error(`YouTube thumbnail upload failed (${thumb.status}).`);
|
|
}
|
|
return { id: data.id, url: `https://www.youtube.com/watch?v=${data.id}` };
|
|
}
|
|
|
|
async function publishTiktok(upload: typeof videoUploads.$inferSelect) {
|
|
const token = await getTiktokAccessToken();
|
|
if (!token) throw new Error("TikTok OAuth is not configured.");
|
|
const settings = upload.tiktokSettings as TiktokSettings;
|
|
const creatorResponse = await fetch(
|
|
"https://open.tiktokapis.com/v2/post/publish/creator_info/query/",
|
|
{
|
|
method: "POST",
|
|
headers: {
|
|
Authorization: `Bearer ${token}`,
|
|
"Content-Type": "application/json; charset=UTF-8",
|
|
},
|
|
},
|
|
);
|
|
if (!creatorResponse.ok) {
|
|
const detail = await creatorResponse.text().catch(() => "");
|
|
throw new Error(
|
|
`TikTok creator information lookup failed (${creatorResponse.status})${detail ? `: ${detail}` : "."}`,
|
|
);
|
|
}
|
|
const creator = (await creatorResponse.json()) as {
|
|
data?: { privacy_level_options?: string[] };
|
|
};
|
|
const privacy =
|
|
settings.privacyLevel ??
|
|
creator.data?.privacy_level_options?.[0] ??
|
|
"SELF_ONLY";
|
|
const video = await readFile(upload.videoPath);
|
|
const chunkSize = Math.min(10 * 1024 * 1024, video.byteLength);
|
|
const totalChunkCount = Math.max(1, Math.floor(video.byteLength / chunkSize));
|
|
const init = await fetch(
|
|
"https://open.tiktokapis.com/v2/post/publish/video/init/",
|
|
{
|
|
method: "POST",
|
|
headers: {
|
|
Authorization: `Bearer ${token}`,
|
|
"Content-Type": "application/json; charset=UTF-8",
|
|
},
|
|
body: JSON.stringify({
|
|
post_info: {
|
|
title: upload.title,
|
|
privacy_level: privacy,
|
|
disable_comment: settings.disableComment ?? false,
|
|
disable_duet: settings.disableDuet ?? false,
|
|
disable_stitch: settings.disableStitch ?? false,
|
|
video_cover_timestamp_ms: settings.coverTimestampMs ?? 0,
|
|
brand_content_toggle: settings.brandContentToggle ?? false,
|
|
brand_organic_toggle: settings.brandOrganicToggle ?? false,
|
|
is_aigc: settings.isAigc ?? false,
|
|
},
|
|
source_info: {
|
|
source: "FILE_UPLOAD",
|
|
video_size: video.byteLength,
|
|
chunk_size: chunkSize,
|
|
total_chunk_count: totalChunkCount,
|
|
},
|
|
}),
|
|
},
|
|
);
|
|
if (!init.ok) {
|
|
const detail = await init.text().catch(() => "");
|
|
throw new Error(
|
|
`TikTok upload initialization failed (${init.status})${detail ? `: ${detail}` : "."}`,
|
|
);
|
|
}
|
|
const initData = (await init.json()) as {
|
|
data?: { publish_id?: string; upload_url?: string };
|
|
};
|
|
const uploadUrl = initData.data?.upload_url;
|
|
if (!uploadUrl || !initData.data?.publish_id)
|
|
throw new Error("TikTok did not return an upload ticket.");
|
|
for (let chunk = 0; chunk < totalChunkCount; chunk += 1) {
|
|
const start = chunk * chunkSize;
|
|
const end =
|
|
chunk === totalChunkCount - 1
|
|
? video.byteLength - 1
|
|
: start + chunkSize - 1;
|
|
const result = await fetch(uploadUrl, {
|
|
method: "PUT",
|
|
headers: {
|
|
"Content-Type": upload.mimeType,
|
|
"Content-Length": String(end - start + 1),
|
|
"Content-Range": `bytes ${start}-${end}/${video.byteLength}`,
|
|
},
|
|
body: video.subarray(start, end + 1),
|
|
});
|
|
if (!result.ok)
|
|
throw new Error(`TikTok video transfer failed (${result.status}).`);
|
|
}
|
|
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;
|
|
}
|
|
}
|
|
|
|
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,
|
|
});
|
|
}
|
|
}
|