import { readFile } from "node:fs/promises"; import { db } from "@/db"; import { videoUploadJobs, videoUploads, type TiktokSettings, type YoutubeSettings, } from "@/db/schema"; 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"; import { getTiktokApiConfig, getYoutubeFollowerConfig, getYoutubePublishingConfig, } from "@/lib/config/server"; 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 { let apiKey: string; let channelId: string; try { ({ apiKey, channelId } = getYoutubeFollowerConfig()); } catch { 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; }; }>; }; 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() { let config: ReturnType; try { config = getTiktokApiConfig(); } catch { return null; } if (config.accessToken) return config.accessToken; const { refreshToken, clientKey, clientSecret } = config; if (!refreshToken || !clientKey || !clientSecret) return 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 { 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, clientId, clientSecret } = getYoutubePublishingConfig(); 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(/* turbopackIgnore: true */ 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(/* turbopackIgnore: true */ 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(/* turbopackIgnore: true */ 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/" }; } type DatabaseVideoUpload = typeof videoUploads.$inferSelect; type DatabaseVideoJob = typeof videoUploadJobs.$inferSelect; 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;