Add admin viewers SSE so device and lock changes appear live without polling.
This commit is contained in:
parent
f6d484aa56
commit
55695aa47c
5 changed files with 195 additions and 3 deletions
|
|
@ -63,13 +63,15 @@ async function proxy(request: NextRequest, pathSegments: string[]) {
|
||||||
const search = url.search;
|
const search = url.search;
|
||||||
const isLogin = request.method === "POST" && targetPath === "v1/auth/login";
|
const isLogin = request.method === "POST" && targetPath === "v1/auth/login";
|
||||||
const isLogout = request.method === "POST" && targetPath === "v1/auth/logout";
|
const isLogout = request.method === "POST" && targetPath === "v1/auth/logout";
|
||||||
|
const isSse =
|
||||||
|
request.method === "GET" && targetPath === "v1/admin/viewers/events";
|
||||||
|
|
||||||
const headers = new Headers();
|
const headers = new Headers();
|
||||||
const contentType = request.headers.get("content-type");
|
const contentType = request.headers.get("content-type");
|
||||||
if (contentType) headers.set("content-type", contentType);
|
if (contentType) headers.set("content-type", contentType);
|
||||||
const cookie = request.headers.get("cookie");
|
const cookie = request.headers.get("cookie");
|
||||||
if (cookie) headers.set("cookie", cookie);
|
if (cookie) headers.set("cookie", cookie);
|
||||||
headers.set("accept", "application/json");
|
headers.set("accept", isSse ? "text/event-stream" : "application/json");
|
||||||
|
|
||||||
const body =
|
const body =
|
||||||
request.method !== "GET" && request.method !== "HEAD"
|
request.method !== "GET" && request.method !== "HEAD"
|
||||||
|
|
@ -86,6 +88,10 @@ async function proxy(request: NextRequest, pathSegments: string[]) {
|
||||||
headers,
|
headers,
|
||||||
body,
|
body,
|
||||||
redirect: "manual",
|
redirect: "manual",
|
||||||
|
// SSE must not be aborted when the proxy function "returns" the stream
|
||||||
|
...(isSse
|
||||||
|
? { cache: "no-store" as RequestCache, signal: request.signal }
|
||||||
|
: {}),
|
||||||
});
|
});
|
||||||
|
|
||||||
// Special-case auth so the session cookie is owned by admin.vonas.nl
|
// Special-case auth so the session cookie is owned by admin.vonas.nl
|
||||||
|
|
@ -130,6 +136,34 @@ async function proxy(request: NextRequest, pathSegments: string[]) {
|
||||||
return response;
|
return response;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if (isSse) {
|
||||||
|
if (!upstream.ok) {
|
||||||
|
const errBody = await upstream.arrayBuffer();
|
||||||
|
return new NextResponse(errBody, {
|
||||||
|
status: upstream.status,
|
||||||
|
headers: {
|
||||||
|
"content-type":
|
||||||
|
upstream.headers.get("content-type") ?? "application/json",
|
||||||
|
},
|
||||||
|
});
|
||||||
|
}
|
||||||
|
if (!upstream.body) {
|
||||||
|
return NextResponse.json(
|
||||||
|
{ error: { code: "UPSTREAM_ERROR", message: "SSE stream missing body" } },
|
||||||
|
{ status: 502 }
|
||||||
|
);
|
||||||
|
}
|
||||||
|
return new NextResponse(upstream.body, {
|
||||||
|
status: 200,
|
||||||
|
headers: {
|
||||||
|
"Content-Type": "text/event-stream; charset=utf-8",
|
||||||
|
"Cache-Control": "no-cache, no-transform",
|
||||||
|
Connection: "keep-alive",
|
||||||
|
"X-Accel-Buffering": "no",
|
||||||
|
},
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
const responseBody = await upstream.arrayBuffer();
|
const responseBody = await upstream.arrayBuffer();
|
||||||
const response = new NextResponse(responseBody, { status: upstream.status });
|
const response = new NextResponse(responseBody, { status: upstream.status });
|
||||||
response.headers.set(
|
response.headers.set(
|
||||||
|
|
|
||||||
|
|
@ -194,6 +194,48 @@ export default function ViewersPage() {
|
||||||
.catch(() => undefined);
|
.catch(() => undefined);
|
||||||
}, [load]);
|
}, [load]);
|
||||||
|
|
||||||
|
// Live updates via SSE (device link, login lock, CRUD). Geen interval-poll.
|
||||||
|
// Overslaan tijdens typen; formuliervelden zitten in aparte state en blijven staan.
|
||||||
|
useEffect(() => {
|
||||||
|
let debounce: ReturnType<typeof setTimeout> | null = null;
|
||||||
|
const isTyping = () => {
|
||||||
|
const el = document.activeElement;
|
||||||
|
if (!el || !(el instanceof HTMLElement)) return false;
|
||||||
|
const tag = el.tagName;
|
||||||
|
return tag === "INPUT" || tag === "TEXTAREA" || tag === "SELECT" || el.isContentEditable;
|
||||||
|
};
|
||||||
|
const scheduleRefresh = () => {
|
||||||
|
if (isTyping()) return;
|
||||||
|
if (debounce) clearTimeout(debounce);
|
||||||
|
debounce = setTimeout(() => {
|
||||||
|
if (!isTyping()) load();
|
||||||
|
}, 200);
|
||||||
|
};
|
||||||
|
|
||||||
|
const es = new EventSource("/api/v1/admin/viewers/events");
|
||||||
|
es.addEventListener("viewers", (ev) => {
|
||||||
|
try {
|
||||||
|
const data = JSON.parse((ev as MessageEvent).data) as { type?: string };
|
||||||
|
if (data.type === "hello") return;
|
||||||
|
} catch {
|
||||||
|
// still refresh on malformed payloads
|
||||||
|
}
|
||||||
|
scheduleRefresh();
|
||||||
|
});
|
||||||
|
|
||||||
|
// Fallback als je terugkomt op het tabblad (bijv. na korte netwerk-dip)
|
||||||
|
const onVisibility = () => {
|
||||||
|
if (document.visibilityState === "visible") scheduleRefresh();
|
||||||
|
};
|
||||||
|
document.addEventListener("visibilitychange", onVisibility);
|
||||||
|
|
||||||
|
return () => {
|
||||||
|
if (debounce) clearTimeout(debounce);
|
||||||
|
es.close();
|
||||||
|
document.removeEventListener("visibilitychange", onVisibility);
|
||||||
|
};
|
||||||
|
}, [load]);
|
||||||
|
|
||||||
const stats = useMemo(() => {
|
const stats = useMemo(() => {
|
||||||
const active = viewers.filter((v) => v.enabled).length;
|
const active = viewers.filter((v) => v.enabled).length;
|
||||||
const appOn = viewers.filter((v) => v.appAccess).length;
|
const appOn = viewers.filter((v) => v.appAccess).length;
|
||||||
|
|
@ -809,6 +851,9 @@ export default function ViewersPage() {
|
||||||
/>
|
/>
|
||||||
</label>
|
</label>
|
||||||
<div className="sidebar-tools">
|
<div className="sidebar-tools">
|
||||||
|
<button type="button" onClick={() => load()} title="Lijst opnieuw laden">
|
||||||
|
Vernieuwen
|
||||||
|
</button>
|
||||||
{playConfig?.configured ? (
|
{playConfig?.configured ? (
|
||||||
<>
|
<>
|
||||||
<button type="button" disabled={bulkBusy} onClick={() => void bulkAction("sync-all")}>
|
<button type="button" disabled={bulkBusy} onClick={() => void bulkAction("sync-all")}>
|
||||||
|
|
|
||||||
40
apps/master-api/src/viewer/admin-events.ts
Normal file
40
apps/master-api/src/viewer/admin-events.ts
Normal file
|
|
@ -0,0 +1,40 @@
|
||||||
|
/**
|
||||||
|
* In-process fan-out for admin UI SSE (/api/v1/admin/viewers/events).
|
||||||
|
* Single master-api instance — fine for this deploy. No Redis needed.
|
||||||
|
*/
|
||||||
|
|
||||||
|
export type ViewerAdminEvent = {
|
||||||
|
type: string;
|
||||||
|
viewerId?: string;
|
||||||
|
at: number;
|
||||||
|
};
|
||||||
|
|
||||||
|
type Listener = (event: ViewerAdminEvent) => void;
|
||||||
|
|
||||||
|
const listeners = new Set<Listener>();
|
||||||
|
|
||||||
|
export function notifyViewersChanged(
|
||||||
|
type: string,
|
||||||
|
opts?: { viewerId?: string }
|
||||||
|
): void {
|
||||||
|
if (listeners.size === 0) return;
|
||||||
|
const event: ViewerAdminEvent = {
|
||||||
|
type,
|
||||||
|
viewerId: opts?.viewerId,
|
||||||
|
at: Date.now(),
|
||||||
|
};
|
||||||
|
for (const listener of listeners) {
|
||||||
|
try {
|
||||||
|
listener(event);
|
||||||
|
} catch {
|
||||||
|
// never break the mutation path for a broken SSE client
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
export function subscribeViewersChanged(listener: Listener): () => void {
|
||||||
|
listeners.add(listener);
|
||||||
|
return () => {
|
||||||
|
listeners.delete(listener);
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
@ -8,6 +8,7 @@ import { ViewerProfileService } from "./profiles";
|
||||||
import type { DownloadService } from "../downloads/service";
|
import type { DownloadService } from "../downloads/service";
|
||||||
import type { SubtitleSource } from "../opensubtitles/client";
|
import type { SubtitleSource } from "../opensubtitles/client";
|
||||||
import { GooglePlayAccessService } from "../google-play/service";
|
import { GooglePlayAccessService } from "../google-play/service";
|
||||||
|
import { subscribeViewersChanged } from "./admin-events";
|
||||||
|
|
||||||
export function registerViewerRoutes(
|
export function registerViewerRoutes(
|
||||||
app: FastifyInstance,
|
app: FastifyInstance,
|
||||||
|
|
@ -24,6 +25,50 @@ export function registerViewerRoutes(
|
||||||
return { viewers: await viewers.listViewers() };
|
return { viewers: await viewers.listViewers() };
|
||||||
});
|
});
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Live updates for admin UI (device link, login lock, CRUD).
|
||||||
|
* One long-lived connection per open viewers page; heartbeats keep proxies happy.
|
||||||
|
*/
|
||||||
|
app.get(
|
||||||
|
"/api/v1/admin/viewers/events",
|
||||||
|
{
|
||||||
|
preHandler: requireAdmin,
|
||||||
|
config: { rateLimit: false },
|
||||||
|
},
|
||||||
|
(request, reply) => {
|
||||||
|
reply.hijack();
|
||||||
|
const res = reply.raw;
|
||||||
|
res.writeHead(200, {
|
||||||
|
"Content-Type": "text/event-stream; charset=utf-8",
|
||||||
|
"Cache-Control": "no-cache, no-transform",
|
||||||
|
Connection: "keep-alive",
|
||||||
|
"X-Accel-Buffering": "no",
|
||||||
|
});
|
||||||
|
|
||||||
|
const write = (chunk: string) => {
|
||||||
|
if (!res.writableEnded) res.write(chunk);
|
||||||
|
};
|
||||||
|
|
||||||
|
write(": connected\n\n");
|
||||||
|
write(`event: viewers\ndata: ${JSON.stringify({ type: "hello", at: Date.now() })}\n\n`);
|
||||||
|
|
||||||
|
const unsubscribe = subscribeViewersChanged((event) => {
|
||||||
|
write(`event: viewers\ndata: ${JSON.stringify(event)}\n\n`);
|
||||||
|
});
|
||||||
|
|
||||||
|
const ping = setInterval(() => {
|
||||||
|
write(`: ping ${Date.now()}\n\n`);
|
||||||
|
}, 20_000);
|
||||||
|
|
||||||
|
const cleanup = () => {
|
||||||
|
clearInterval(ping);
|
||||||
|
unsubscribe();
|
||||||
|
};
|
||||||
|
|
||||||
|
request.raw.on("close", cleanup);
|
||||||
|
request.raw.on("error", cleanup);
|
||||||
|
}
|
||||||
|
);
|
||||||
app.post("/api/v1/admin/viewers", { preHandler: requireAdmin }, async (request) => {
|
app.post("/api/v1/admin/viewers", { preHandler: requireAdmin }, async (request) => {
|
||||||
const body = request.body as {
|
const body = request.body as {
|
||||||
email?: string;
|
email?: string;
|
||||||
|
|
|
||||||
|
|
@ -22,6 +22,7 @@ import {
|
||||||
kidsSeriesSql,
|
kidsSeriesSql,
|
||||||
isKidsBlockedTitle,
|
isKidsBlockedTitle,
|
||||||
} from "./kids-filter";
|
} from "./kids-filter";
|
||||||
|
import { notifyViewersChanged } from "./admin-events";
|
||||||
|
|
||||||
const CODE_TTL_MS = 10 * 60 * 1000;
|
const CODE_TTL_MS = 10 * 60 * 1000;
|
||||||
const CODE_ALPHABET = "ABCDEFGHJKLMNPQRSTUVWXYZ23456789";
|
const CODE_ALPHABET = "ABCDEFGHJKLMNPQRSTUVWXYZ23456789";
|
||||||
|
|
@ -116,6 +117,7 @@ export class ViewerService {
|
||||||
},
|
},
|
||||||
});
|
});
|
||||||
await ensureOwnerProfile(viewer.id, viewer.name);
|
await ensureOwnerProfile(viewer.id, viewer.name);
|
||||||
|
notifyViewersChanged("viewer.created", { viewerId: viewer.id });
|
||||||
return viewer;
|
return viewer;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -217,6 +219,7 @@ export class ViewerService {
|
||||||
if (!existing) throw new AppError("NOT_FOUND", "Gebruiker niet gevonden", 404);
|
if (!existing) throw new AppError("NOT_FOUND", "Gebruiker niet gevonden", 404);
|
||||||
// Cascade: profiles, devices, prefs, favorites, iptv-lijn; addon-tokens blijven (viewerUserId → null)
|
// Cascade: profiles, devices, prefs, favorites, iptv-lijn; addon-tokens blijven (viewerUserId → null)
|
||||||
await prisma.viewerUser.delete({ where: { id } });
|
await prisma.viewerUser.delete({ where: { id } });
|
||||||
|
notifyViewersChanged("viewer.deleted", { viewerId: id });
|
||||||
return { ok: true };
|
return { ok: true };
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -291,14 +294,24 @@ export class ViewerService {
|
||||||
} else if (patch.passwordHash) {
|
} else if (patch.passwordHash) {
|
||||||
await this.clearLoginLock(updated.email.toLowerCase());
|
await this.clearLoginLock(updated.email.toLowerCase());
|
||||||
}
|
}
|
||||||
|
if (Object.keys(patch).length > 0) {
|
||||||
|
notifyViewersChanged("viewer.updated", { viewerId: updated.id });
|
||||||
|
}
|
||||||
return updated;
|
return updated;
|
||||||
}
|
}
|
||||||
|
|
||||||
async revokeDevice(deviceId: string) {
|
async revokeDevice(deviceId: string) {
|
||||||
|
const device = await prisma.viewerDevice.findUnique({
|
||||||
|
where: { id: deviceId },
|
||||||
|
select: { id: true, viewerUserId: true },
|
||||||
|
});
|
||||||
await prisma.viewerDevice.update({
|
await prisma.viewerDevice.update({
|
||||||
where: { id: deviceId },
|
where: { id: deviceId },
|
||||||
data: { revoked: true, refreshTokenHash: hashToken(`revoked:${deviceId}:${Date.now()}`) },
|
data: { revoked: true, refreshTokenHash: hashToken(`revoked:${deviceId}:${Date.now()}`) },
|
||||||
});
|
});
|
||||||
|
if (device) {
|
||||||
|
notifyViewersChanged("device.revoked", { viewerId: device.viewerUserId });
|
||||||
|
}
|
||||||
return { ok: true };
|
return { ok: true };
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -448,6 +461,7 @@ export class ViewerService {
|
||||||
|
|
||||||
if (!device) return { status: "already_linked" as const };
|
if (!device) return { status: "already_linked" as const };
|
||||||
|
|
||||||
|
notifyViewersChanged("device.linked", { viewerId: viewer.id });
|
||||||
return {
|
return {
|
||||||
status: "linked" as const,
|
status: "linked" as const,
|
||||||
refreshToken,
|
refreshToken,
|
||||||
|
|
@ -553,6 +567,7 @@ export class ViewerService {
|
||||||
await this.clearLoginLock(email);
|
await this.clearLoginLock(email);
|
||||||
await ensureOwnerProfile(viewer.id, viewer.name);
|
await ensureOwnerProfile(viewer.id, viewer.name);
|
||||||
|
|
||||||
|
notifyViewersChanged("device.login", { viewerId: viewer.id });
|
||||||
return {
|
return {
|
||||||
refreshToken,
|
refreshToken,
|
||||||
device: { id: device.id, name: device.name, platform: device.platform },
|
device: { id: device.id, name: device.name, platform: device.platform },
|
||||||
|
|
@ -590,15 +605,28 @@ export class ViewerService {
|
||||||
const failedCount = baseCount + 1;
|
const failedCount = baseCount + 1;
|
||||||
const lockedUntil =
|
const lockedUntil =
|
||||||
failedCount >= LOGIN_MAX_FAILURES ? new Date(now + LOGIN_LOCK_MS) : null;
|
failedCount >= LOGIN_MAX_FAILURES ? new Date(now + LOGIN_LOCK_MS) : null;
|
||||||
return prisma.viewerLoginLock.upsert({
|
const lock = await prisma.viewerLoginLock.upsert({
|
||||||
where: { email },
|
where: { email },
|
||||||
create: { email, failedCount, lockedUntil },
|
create: { email, failedCount, lockedUntil },
|
||||||
update: { failedCount, lockedUntil },
|
update: { failedCount, lockedUntil },
|
||||||
});
|
});
|
||||||
|
const viewer = await prisma.viewerUser.findUnique({
|
||||||
|
where: { email },
|
||||||
|
select: { id: true },
|
||||||
|
});
|
||||||
|
notifyViewersChanged("login.lock", { viewerId: viewer?.id });
|
||||||
|
return lock;
|
||||||
}
|
}
|
||||||
|
|
||||||
private async clearLoginLock(email: string) {
|
private async clearLoginLock(email: string) {
|
||||||
await prisma.viewerLoginLock.deleteMany({ where: { email } });
|
const result = await prisma.viewerLoginLock.deleteMany({ where: { email } });
|
||||||
|
if (result.count > 0) {
|
||||||
|
const viewer = await prisma.viewerUser.findUnique({
|
||||||
|
where: { email },
|
||||||
|
select: { id: true },
|
||||||
|
});
|
||||||
|
notifyViewersChanged("login.unlock", { viewerId: viewer?.id });
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async authFromBearer(authHeader?: string): Promise<AuthedViewer> {
|
async authFromBearer(authHeader?: string): Promise<AuthedViewer> {
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue