feat(upload): add realtime SSE updates
This commit is contained in:
@@ -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");
|
||||
}
|
||||
|
||||
@@ -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 (
|
||||
<>
|
||||
<RealtimeRefresh topic="uploads" />
|
||||
<Card>
|
||||
<CardHeader>
|
||||
<CardTitle>Upload history</CardTitle>
|
||||
@@ -92,5 +95,6 @@ export default async function AdminUploadHistoryPage() {
|
||||
)}
|
||||
</CardContent>
|
||||
</Card>
|
||||
</>
|
||||
);
|
||||
}
|
||||
|
||||
@@ -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 (
|
||||
<>
|
||||
<RealtimeRefresh topic="uploads" debounceMs={500} />
|
||||
<section className="grid gap-4 xl:grid-cols-2">
|
||||
<LatestCard title="Latest YouTube videos" videos={youtube} />
|
||||
<LatestCard title="Latest TikTok videos" videos={tiktok} />
|
||||
</section>
|
||||
</>
|
||||
);
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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",
|
||||
},
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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(),
|
||||
|
||||
@@ -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,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user