From 7f4a1a339c01748df56fe38e37bd42e8485ae566 Mon Sep 17 00:00:00 2001 From: gunshiz Date: Sun, 2 Aug 2026 18:43:45 +0700 Subject: [PATCH] feat(upload): add realtime SSE updates --- app/admin/upload/actions.ts | 7 +++++++ app/admin/upload/history/page.tsx | 4 ++++ app/admin/upload/video/page.tsx | 4 ++++ components/realtime-refresh.tsx | 15 +++++++++++++-- lib/realtime/realtime-refresh.ts | 5 +++++ lib/realtime/sse.test.ts | 29 +++++++++++++++++++++++++++++ lib/realtime/sse.ts | 18 ++++++++++++++++++ lib/video/video-platforms.ts | 23 +++++++++++++++++++++++ lib/video/video-upload-storage.ts | 9 +++++++++ 9 files changed, 112 insertions(+), 2 deletions(-) diff --git a/app/admin/upload/actions.ts b/app/admin/upload/actions.ts index 18eecb9..0b12f18 100644 --- a/app/admin/upload/actions.ts +++ b/app/admin/upload/actions.ts @@ -4,6 +4,7 @@ import { requireAdmin } from "@/lib/auth/auth"; import { db } from "@/db"; import { videoUploadJobs, videoUploads } from "@/db/schema"; import { eq } from "drizzle-orm"; +import { sse } from "@/lib/realtime/sse"; import { revalidatePath } from "next/cache"; export async function createVideoUpload(formData: FormData) { @@ -28,5 +29,11 @@ export async function retryVideoJob(jobId: string) { .update(videoUploads) .set({ status: "partial", updatedAt: new Date() }) .where(eq(videoUploads.id, job.uploadId)); + await sse.uploads.pub("update", { + uploadId: job.uploadId, + jobId: job.id, + platform: job.platform, + status: "pending", + }); revalidatePath("/admin/upload/history"); } diff --git a/app/admin/upload/history/page.tsx b/app/admin/upload/history/page.tsx index 72840b4..03db38c 100644 --- a/app/admin/upload/history/page.tsx +++ b/app/admin/upload/history/page.tsx @@ -1,5 +1,6 @@ import { Badge } from "@/components/ui/badge"; import { Button } from "@/components/ui/button"; +import { RealtimeRefresh } from "@/components/realtime-refresh"; import { Card, CardContent, @@ -44,6 +45,8 @@ export default async function AdminUploadHistoryPage() { ]); return ( + <> + Upload history @@ -92,5 +95,6 @@ export default async function AdminUploadHistoryPage() { )} + ); } diff --git a/app/admin/upload/video/page.tsx b/app/admin/upload/video/page.tsx index 31c1018..e5834ac 100644 --- a/app/admin/upload/video/page.tsx +++ b/app/admin/upload/video/page.tsx @@ -3,6 +3,7 @@ import { getLatestYoutubeVideos, } from "@/lib/video/video-platforms"; import { LatestCard } from "./latest-card"; +import { RealtimeRefresh } from "@/components/realtime-refresh"; export const dynamic = "force-dynamic"; @@ -13,9 +14,12 @@ export default async function AdminUploadVideoPage() { ]); return ( + <> +
+ ); } diff --git a/components/realtime-refresh.tsx b/components/realtime-refresh.tsx index b6cf965..bf43bb2 100644 --- a/components/realtime-refresh.tsx +++ b/components/realtime-refresh.tsx @@ -14,6 +14,11 @@ type FormRealtimeRefreshProps = { debounceMs?: number; }; +type UploadRealtimeRefreshProps = { + topic: "uploads"; + debounceMs?: number; +}; + type LeaderboardRealtimeRefreshProps = | { topic: "leaderboards"; @@ -30,12 +35,13 @@ type LeaderboardRealtimeRefreshProps = type RealtimeRefreshProps = | FormRealtimeRefreshProps + | UploadRealtimeRefreshProps | LeaderboardRealtimeRefreshProps; export function RealtimeRefresh(props: RealtimeRefreshProps) { const router = useRouter(); const topic = props.topic; - const formId = topic === "leaderboards" ? undefined : props.formId; + const formId = topic === "forms" || topic === "submissions" ? props.formId : undefined; const leaderboard = topic === "leaderboards" ? props.leaderboard : undefined; const channelId = topic === "leaderboards" ? props.channelId : undefined; const debounceMs = props.debounceMs ?? (topic === "leaderboards" ? 500 : 150); @@ -45,7 +51,9 @@ export function RealtimeRefresh(props: RealtimeRefreshProps) { ? leaderboard === "xp" ? { kind: "xp" as const } : { kind: "vc" as const, channelId } - : { kind: "form" as const, formId }; + : topic === "uploads" + ? { kind: "uploads" as const } + : { kind: "form" as const, formId }; return startRealtimeRefresh({ filter, @@ -60,6 +68,9 @@ export function RealtimeRefresh(props: RealtimeRefreshProps) { if (topic === "submissions") { return sse.submissions.sub("update", handlers.onForm, options); } + if (topic === "uploads") { + return sse.uploads.sub("update", handlers.onUpload, options); + } return sse.leaderboards.subMany({ xp: handlers.onXp, vc: handlers.onVc, diff --git a/lib/realtime/realtime-refresh.ts b/lib/realtime/realtime-refresh.ts index 5b15119..9dccfe3 100644 --- a/lib/realtime/realtime-refresh.ts +++ b/lib/realtime/realtime-refresh.ts @@ -1,11 +1,13 @@ export type RealtimeRefreshFilter = | { kind: "form"; formId?: string } + | { kind: "uploads" } | { kind: "xp" } | { kind: "vc"; channelId?: string }; export type RealtimeRefreshHandlers = { onOpen: () => void; onForm: (event: { formId: string }) => void; + onUpload: () => void; onXp: () => void; onVc: (event: { channelId: string }) => void; }; @@ -50,6 +52,9 @@ export function startRealtimeRefresh({ onXp() { if (filter.kind === "xp") scheduleRefresh(); }, + onUpload() { + if (filter.kind === "uploads") scheduleRefresh(); + }, onVc(event) { if ( filter.kind === "vc" diff --git a/lib/realtime/sse.test.ts b/lib/realtime/sse.test.ts index 14f8d3c..dd25f66 100644 --- a/lib/realtime/sse.test.ts +++ b/lib/realtime/sse.test.ts @@ -37,3 +37,32 @@ describe("leaderboard SSE contract", () => { )).toBeNull(); }); }); + +describe("upload SSE contract", () => { + test("registers an admin-only upload topic", () => { + expect(isSseTopic("uploads")).toBeTrue(); + expect(sse.uploads.adminOnly).toBeTrue(); + }); + + test("validates upload job status events", () => { + const message = JSON.stringify({ + event: "update", + data: { + uploadId: "upload-1", + jobId: "job-1", + platform: "youtube", + status: "completed", + }, + }); + + expect(sse.uploads.parseRedisMessage(message)).toEqual({ + event: "update", + data: { + uploadId: "upload-1", + jobId: "job-1", + platform: "youtube", + status: "completed", + }, + }); + }); +}); diff --git a/lib/realtime/sse.ts b/lib/realtime/sse.ts index c73a7ac..7299902 100644 --- a/lib/realtime/sse.ts +++ b/lib/realtime/sse.ts @@ -306,6 +306,24 @@ export const sse = createSseEndpoints({ }).strict(), }, }, + uploads: { + adminOnly: true, + events: { + update: z.object({ + uploadId: z.string().min(1), + jobId: z.string().min(1).optional(), + platform: z.enum(["youtube", "tiktok"]).optional(), + status: z.enum([ + "scheduled", + "processing", + "pending", + "completed", + "partial", + "failed", + ]), + }).strict(), + }, + }, leaderboards: { events: { xp: z.object({}).strict(), diff --git a/lib/video/video-platforms.ts b/lib/video/video-platforms.ts index 34c8f7d..3c6679a 100644 --- a/lib/video/video-platforms.ts +++ b/lib/video/video-platforms.ts @@ -7,6 +7,7 @@ import { 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"]); @@ -338,6 +339,12 @@ export async function publishVideoJob(jobId: string) { 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" @@ -352,6 +359,12 @@ export async function publishVideoJob(jobId: string) { 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."; @@ -359,6 +372,12 @@ export async function publishVideoJob(jobId: string) { .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; } } @@ -398,5 +417,9 @@ export async function processDueVideoUploads() { .update(videoUploads) .set({ status: nextStatus, updatedAt: new Date() }) .where(eq(videoUploads.id, upload.id)); + await sse.uploads.pub("update", { + uploadId: upload.id, + status: nextStatus, + }); } } diff --git a/lib/video/video-upload-storage.ts b/lib/video/video-upload-storage.ts index 40c2cfe..a77eb16 100644 --- a/lib/video/video-upload-storage.ts +++ b/lib/video/video-upload-storage.ts @@ -2,6 +2,7 @@ import "server-only"; import { mkdir, readdir, stat, writeFile } from "node:fs/promises"; import { randomUUID } from "node:crypto"; +import { sse } from "@/lib/realtime/sse"; import { db } from "@/db"; import { videoUploadJobs, @@ -107,6 +108,10 @@ export async function saveVideoUpload(formData: FormData, discordId: string) { { uploadId: upload.id, platform: "tiktok" }, ]) .returning({ id: videoUploadJobs.id }); + await sse.uploads.pub("update", { + uploadId: upload.id, + status: scheduledAt ? "scheduled" : "processing", + }); if (!scheduledAt) { await Promise.allSettled(jobs.map((job) => publishVideoJob(job.id))); const remaining = await db.query.videoUploadJobs.findMany({ @@ -120,6 +125,10 @@ export async function saveVideoUpload(formData: FormData, discordId: string) { updatedAt: new Date(), }) .where(eq(videoUploads.id, upload.id)); + await sse.uploads.pub("update", { + uploadId: upload.id, + status: remaining.length === 0 ? "completed" : "partial", + }); } return upload.id; }