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;
}