refactor(realtime): centralize refresh policy
This commit is contained in:
@@ -3,10 +3,7 @@
|
||||
import { useEffect } from "react";
|
||||
import { useRouter } from "next/navigation";
|
||||
import { sse } from "@/lib/realtime/sse";
|
||||
import {
|
||||
startRealtimeRefresh,
|
||||
type RealtimeRefreshHandlers,
|
||||
} from "@/lib/realtime/realtime-refresh";
|
||||
import { startRealtimeRefresh } from "@/lib/realtime/realtime-refresh";
|
||||
|
||||
type FormRealtimeRefreshProps = {
|
||||
topic: "forms" | "submissions";
|
||||
@@ -67,21 +64,29 @@ export function RealtimeRefresh(props: RealtimeRefreshProps) {
|
||||
filter,
|
||||
debounceMs,
|
||||
refresh: () => router.refresh(),
|
||||
subscribe(handlers: RealtimeRefreshHandlers) {
|
||||
const options = { onopen: handlers.onOpen };
|
||||
subscribe(emit) {
|
||||
const options = { onopen: () => emit({ kind: "open" }) };
|
||||
|
||||
if (topic === "forms") {
|
||||
return sse.forms.sub("update", handlers.onForm, options);
|
||||
return sse.forms.sub(
|
||||
"update",
|
||||
(event) => emit({ kind: "form", formId: event.formId }),
|
||||
options,
|
||||
);
|
||||
}
|
||||
if (topic === "submissions") {
|
||||
return sse.submissions.sub("update", handlers.onForm, options);
|
||||
return sse.submissions.sub(
|
||||
"update",
|
||||
(event) => emit({ kind: "form", formId: event.formId }),
|
||||
options,
|
||||
);
|
||||
}
|
||||
if (topic === "uploads") {
|
||||
return sse.uploads.sub("update", handlers.onUpload, options);
|
||||
return sse.uploads.sub("update", () => emit({ kind: "upload" }), options);
|
||||
}
|
||||
return sse.leaderboards.subMany({
|
||||
xp: handlers.onXp,
|
||||
vc: handlers.onVc,
|
||||
xp: () => emit({ kind: "xp" }),
|
||||
vc: (event) => emit({ kind: "vc", channelId: event.channelId }),
|
||||
}, options);
|
||||
},
|
||||
});
|
||||
|
||||
@@ -1,58 +1,68 @@
|
||||
import { describe, expect, test } from "bun:test";
|
||||
import {
|
||||
matchesRealtimeRefresh,
|
||||
startRealtimeRefresh,
|
||||
type RealtimeRefreshHandlers,
|
||||
type RealtimeRefreshEvent,
|
||||
} from "@/lib/realtime/realtime-refresh";
|
||||
|
||||
function waitForRefresh() {
|
||||
return new Promise((resolve) => setTimeout(resolve, 5));
|
||||
}
|
||||
|
||||
describe("profile realtime refresh", () => {
|
||||
test("refreshes for both XP and voice events", async () => {
|
||||
let handlers: RealtimeRefreshHandlers | undefined;
|
||||
describe("realtime refresh policy", () => {
|
||||
test("matches page intent against one event stream", () => {
|
||||
expect(matchesRealtimeRefresh({ kind: "profile" }, { kind: "xp" })).toBeTrue();
|
||||
expect(matchesRealtimeRefresh({ kind: "profile" }, { kind: "vc", channelId: "1" })).toBeTrue();
|
||||
expect(matchesRealtimeRefresh({ kind: "vc", channelId: "1" }, { kind: "vc", channelId: "2" })).toBeFalse();
|
||||
expect(matchesRealtimeRefresh({ kind: "form", formId: "form-1" }, { kind: "form", formId: "form-2" })).toBeFalse();
|
||||
expect(matchesRealtimeRefresh({ kind: "uploads" }, { kind: "upload" })).toBeTrue();
|
||||
});
|
||||
|
||||
test("refreshes for matching events and reconnects", async () => {
|
||||
let emit: ((event: RealtimeRefreshEvent) => void) | undefined;
|
||||
let refreshes = 0;
|
||||
const clean = startRealtimeRefresh({
|
||||
filter: { kind: "profile" },
|
||||
debounceMs: 0,
|
||||
refresh: () => { refreshes += 1; },
|
||||
subscribe: (nextHandlers) => {
|
||||
handlers = nextHandlers;
|
||||
subscribe: (nextEmit) => {
|
||||
emit = nextEmit;
|
||||
return { clean() {} };
|
||||
},
|
||||
});
|
||||
|
||||
handlers?.onXp();
|
||||
emit?.({ kind: "open" });
|
||||
emit?.({ kind: "xp" });
|
||||
await waitForRefresh();
|
||||
handlers?.onVc({ channelId: "1234567890" });
|
||||
emit?.({ kind: "vc", channelId: "1234567890" });
|
||||
await waitForRefresh();
|
||||
emit?.({ kind: "open" });
|
||||
await waitForRefresh();
|
||||
|
||||
expect(refreshes).toBe(2);
|
||||
expect(refreshes).toBe(3);
|
||||
clean();
|
||||
});
|
||||
|
||||
test("debounces profile event bursts and refreshes after reconnect", async () => {
|
||||
let handlers: RealtimeRefreshHandlers | undefined;
|
||||
test("debounces event bursts and cancels on cleanup", async () => {
|
||||
let emit: ((event: RealtimeRefreshEvent) => void) | undefined;
|
||||
let refreshes = 0;
|
||||
const clean = startRealtimeRefresh({
|
||||
filter: { kind: "profile" },
|
||||
filter: { kind: "xp" },
|
||||
debounceMs: 1,
|
||||
refresh: () => { refreshes += 1; },
|
||||
subscribe: (nextHandlers) => {
|
||||
handlers = nextHandlers;
|
||||
subscribe: (nextEmit) => {
|
||||
emit = nextEmit;
|
||||
return { clean() {} };
|
||||
},
|
||||
});
|
||||
|
||||
handlers?.onXp();
|
||||
handlers?.onVc({ channelId: "1234567890" });
|
||||
emit?.({ kind: "xp" });
|
||||
emit?.({ kind: "xp" });
|
||||
await waitForRefresh();
|
||||
expect(refreshes).toBe(1);
|
||||
|
||||
handlers?.onOpen();
|
||||
handlers?.onOpen();
|
||||
await waitForRefresh();
|
||||
expect(refreshes).toBe(2);
|
||||
emit?.({ kind: "xp" });
|
||||
clean();
|
||||
await waitForRefresh();
|
||||
expect(refreshes).toBe(1);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -5,27 +5,39 @@ export type RealtimeRefreshFilter =
|
||||
| { 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;
|
||||
};
|
||||
export type RealtimeRefreshEvent =
|
||||
| { kind: "open" }
|
||||
| { kind: "form"; formId: string }
|
||||
| { kind: "upload" }
|
||||
| { kind: "xp" }
|
||||
| { kind: "vc"; channelId: string };
|
||||
|
||||
type StartRealtimeRefreshOptions = {
|
||||
filter: RealtimeRefreshFilter;
|
||||
debounceMs: number;
|
||||
refresh: () => void;
|
||||
subscribe: (handlers: RealtimeRefreshHandlers) => { clean: () => void };
|
||||
};
|
||||
export type RealtimeRefreshSubscription = { clean: () => void };
|
||||
|
||||
export function matchesRealtimeRefresh(
|
||||
filter: RealtimeRefreshFilter,
|
||||
event: Exclude<RealtimeRefreshEvent, { kind: "open" }>,
|
||||
) {
|
||||
if (filter.kind === "form") {
|
||||
return event.kind === "form" && (!filter.formId || event.formId === filter.formId);
|
||||
}
|
||||
if (filter.kind === "uploads") return event.kind === "upload";
|
||||
if (filter.kind === "xp") return event.kind === "xp";
|
||||
if (filter.kind === "profile") return event.kind === "xp" || event.kind === "vc";
|
||||
return event.kind === "vc" && (!filter.channelId || event.channelId === filter.channelId);
|
||||
}
|
||||
|
||||
export function startRealtimeRefresh({
|
||||
filter,
|
||||
debounceMs,
|
||||
refresh,
|
||||
subscribe,
|
||||
}: StartRealtimeRefreshOptions) {
|
||||
}: {
|
||||
filter: RealtimeRefreshFilter;
|
||||
debounceMs: number;
|
||||
refresh: () => void;
|
||||
subscribe: (emit: (event: RealtimeRefreshEvent) => void) => RealtimeRefreshSubscription;
|
||||
}) {
|
||||
let connectedOnce = false;
|
||||
let refreshTimer: ReturnType<typeof setTimeout> | undefined;
|
||||
|
||||
@@ -34,39 +46,14 @@ export function startRealtimeRefresh({
|
||||
refreshTimer = setTimeout(refresh, debounceMs);
|
||||
};
|
||||
|
||||
const subscription = subscribe({
|
||||
onOpen() {
|
||||
if (connectedOnce) {
|
||||
scheduleRefresh();
|
||||
} else {
|
||||
connectedOnce = true;
|
||||
}
|
||||
},
|
||||
onForm(event) {
|
||||
if (
|
||||
filter.kind === "form"
|
||||
&& (!filter.formId || event.formId === filter.formId)
|
||||
) {
|
||||
scheduleRefresh();
|
||||
}
|
||||
},
|
||||
onXp() {
|
||||
if (filter.kind === "xp" || filter.kind === "profile") scheduleRefresh();
|
||||
},
|
||||
onUpload() {
|
||||
if (filter.kind === "uploads") scheduleRefresh();
|
||||
},
|
||||
onVc(event) {
|
||||
if (
|
||||
filter.kind === "profile"
|
||||
|| (
|
||||
filter.kind === "vc"
|
||||
&& (!filter.channelId || event.channelId === filter.channelId)
|
||||
)
|
||||
) {
|
||||
scheduleRefresh();
|
||||
}
|
||||
},
|
||||
const subscription = subscribe((event) => {
|
||||
if (event.kind === "open") {
|
||||
if (connectedOnce) scheduleRefresh();
|
||||
connectedOnce = true;
|
||||
return;
|
||||
}
|
||||
|
||||
if (matchesRealtimeRefresh(filter, event)) scheduleRefresh();
|
||||
});
|
||||
|
||||
return () => {
|
||||
|
||||
Reference in New Issue
Block a user