Add live video: WHEP relay route and WebRTC player

- POST/DELETE /api/cameras/[id]/live/whep check the session, resolve the
  camera from the registry, and relay WebRTC signaling to MediaMTX on
  localhost. The browser never sees MediaMTX's address, its error text or
  a camera login; video flows browser↔MediaMTX, not through Next.
- LivePlayer negotiates with the browser's own RTCPeerConnection and
  shows connecting, reconnecting and failed states with Retry.
- The camera page gains a Live section; the pop-out plays live video and
  falls back to snapshot polling if it can't. Both offer Main/Sub.
- Stop a leftover MediaMTX from a server that didn't exit cleanly: record
  its pid and, on startup, stop it only if that pid is still our binary
  with our config.
- Export vrek log: both cameras watched live on main and sub (4 of the
  goal's 6 checks).

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Michael Mainguy 2026-09-19 16:40:21 -05:00
parent 5824db8ea9
commit 41021a62ad
20 changed files with 1193 additions and 34 deletions

View File

@ -801,3 +801,23 @@
{"id":"evt-r8ad4s4jhgpw","type":"edge.added","subject":"iss-4atmz13","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"kind":"parent_of","from":"iss-4atmz13","to":"iss-yd2sq2q"},"at":"2026-09-19T20:33:23.229Z","parents":["evt-t3hv9z68czv1"],"hash":"ef606c14c0ef29a9538119652ef1482d7ed8b9ec4af4966afc80cdf614ea0173"}
{"id":"evt-88zs5h5gg17f","type":"edge.added","subject":"iss-yd2sq2q","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"kind":"tagged","from":"iss-yd2sq2q","to":"area:video"},"at":"2026-09-19T20:33:23.230Z","parents":["evt-r8ad4s4jhgpw"],"hash":"f04eaa848b5e077d22a68283a76b3d1be03ac60ff6fd3d0a720b234d509cf572"}
{"id":"evt-69m4zgpr4mzg","type":"edge.added","subject":"iss-yd2sq2q","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"kind":"discovered_from","from":"iss-yd2sq2q","to":"iss-nk6zrzv"},"at":"2026-09-19T20:33:23.231Z","parents":["evt-88zs5h5gg17f"],"hash":"c4b77e1ae675f302e51d418c5433781a0e4300741a9dd346ddf78bec9b826e40"}
{"id":"evt-qehb93vm0ycn","type":"node.created","subject":"ver-mmjvk5b","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"kind":"verification","title":"On startup the bridge stops a leftover MediaMTX recorded in .data/mediamtx/mediamtx.pid (SIGTERM, then SIGKILL after about 3 s), but only when that pid still runs our binary with our config; the pid file is written on spawn and removed on exit","body":"","status":"pending","owner":"prn-q80g8mz","attrs":{}},"at":"2026-09-19T20:35:44.472Z","parents":["evt-69m4zgpr4mzg"],"hash":"cb195ffff80a349f25e0976b83c0269220a9643eee71ad0711b7d1cc52b357cd"}
{"id":"evt-k1jkcdnpvpfj","type":"edge.added","subject":"ver-mmjvk5b","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"kind":"evidence_for","from":"ver-mmjvk5b","to":"iss-yd2sq2q"},"at":"2026-09-19T20:35:44.473Z","parents":["evt-qehb93vm0ycn"],"hash":"a4325f9e3429ed46171791ff9fedafb98cf249f188410985b79cd4ffa04d930e"}
{"id":"evt-egns6xcvsxz0","type":"verification.recorded","subject":"ver-mmjvk5b","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"result":"pass","evidence":"src/lib/mediamtx-leftover.test.ts (fake process table; never signals a reused pid), mediamtx-supervisor.test.ts (onSpawn/onExit), video.test.ts (\"leftover MediaMTX and pid file\"). Full suite 738 tests pass, 99.86% lines; tsc and eslint clean; next build compiles."},"at":"2026-09-19T20:35:44.474Z","parents":["evt-k1jkcdnpvpfj"],"hash":"3879a5cf742d0ef6bb36b7ccb4deae76e9972d197f1b71e9080cdad86e99c372"}
{"id":"evt-4rqs4x7276j6","type":"node.created","subject":"ver-r06y7ck","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"kind":"verification","title":"Real check by the user: after force-killing the server (kill -9 on next-server), the next start logs \"[video] stopped a MediaMTX left running by an earlier server\" and MediaMTX starts normally","body":"","status":"pending","owner":"prn-q80g8mz","attrs":{}},"at":"2026-09-19T20:35:46.362Z","parents":["evt-egns6xcvsxz0"],"hash":"2162360435f0db9b49fe200ef289849460a1d1a92317974e6930f77f933fec7c"}
{"id":"evt-4ycqh2c3rdqp","type":"edge.added","subject":"ver-r06y7ck","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"kind":"evidence_for","from":"ver-r06y7ck","to":"iss-yd2sq2q"},"at":"2026-09-19T20:35:46.364Z","parents":["evt-4rqs4x7276j6"],"hash":"1bbff09fcdc57d01100e0ea486995230a0b7b97090c29133d69ecaeca5b9f6f6"}
{"id":"evt-e0y77cepccqd","type":"verification.recorded","subject":"ver-r06y7ck","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"result":"pending","evidence":"User test with the new build"},"at":"2026-09-19T20:35:46.365Z","parents":["evt-4ycqh2c3rdqp"],"hash":"18780a1e5853ebbe44fdbe416b1ef1dc602361687279dfd17154df7c517479bf"}
{"id":"evt-zjaa2vr6m15t","type":"node.status_changed","subject":"iss-yd2sq2q","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"from":"open","to":"in_progress"},"at":"2026-09-19T20:35:48.277Z","parents":["evt-e0y77cepccqd"],"hash":"d1fc6f039f2c5c2bbbc8ef4b38676bbbe2919e0a18d74b50abc7b99d1f901a35"}
{"id":"evt-zc7rkzaets70","type":"node.created","subject":"ver-qsqa18n","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"kind":"verification","title":"POST/DELETE /api/cameras/[id]/live/whep check the session, validate the id and stream, relay signaling to MediaMTX on localhost with the bridge login, and never return MediaMTX's address, its error text or a camera credential; the player negotiates WebRTC in the browser, shows connecting/reconnecting/failed with Retry, and the pop-out falls back to snapshots; the camera page and pop-out both offer Main/Sub","body":"","status":"pending","owner":"prn-q80g8mz","attrs":{}},"at":"2026-09-19T21:35:04.892Z","parents":["evt-zjaa2vr6m15t"],"hash":"0a12104cdbccd44c7db496b321ed67e47f5038c3fe1851573b914cea58a33aaa"}
{"id":"evt-dqw4sny9j4j0","type":"edge.added","subject":"ver-qsqa18n","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"kind":"evidence_for","from":"ver-qsqa18n","to":"iss-cbx21zy"},"at":"2026-09-19T21:35:04.894Z","parents":["evt-zc7rkzaets70"],"hash":"53827af283759ae1b215bad43a3dd797e5f3b07e57d286f4626fb70ac0a9bf42"}
{"id":"evt-r3afr6svxe26","type":"verification.recorded","subject":"ver-qsqa18n","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"result":"pass","evidence":"src/app/api/cameras/[id]/live/whep/route.test.ts (18 tests), live-player.test.tsx (fake RTCPeerConnection), live-view.test.tsx, live-panel.test.tsx, page.test.tsx. Full suite 769 tests pass, 99.81% lines; tsc and eslint clean; next build lists ƒ /api/cameras/[id]/live/whep. WHEP behaviour checked against MediaMTX v1.21.0 internal/servers/webrtc/http_server.go."},"at":"2026-09-19T21:35:04.895Z","parents":["evt-dqw4sny9j4j0"],"hash":"4fcd50e634ab3957602b8fd5e4399e045338f26d7acee865c08db969fef007ec"}
{"id":"evt-t8adbjc357z1","type":"node.created","subject":"ver-pgf465d","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"kind":"verification","title":"Real check by the user: both cameras play live in the pop-out and on the camera page, on main and sub, with the session required and no camera credentials or MediaMTX address visible in the browser","body":"","status":"pending","owner":"prn-q80g8mz","attrs":{}},"at":"2026-09-19T21:35:07.527Z","parents":["evt-r3afr6svxe26"],"hash":"5e823af787b849d1b4b3715c9cd639ff971bde25046f099dd2543d8e8437ef19"}
{"id":"evt-vcksn6nsjyhr","type":"edge.added","subject":"ver-pgf465d","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"kind":"evidence_for","from":"ver-pgf465d","to":"iss-cbx21zy"},"at":"2026-09-19T21:35:07.530Z","parents":["evt-t8adbjc357z1"],"hash":"e68a14bd5ee9d09fde7ce89631455a5d647fa53c39ff167934438ab2884c8703"}
{"id":"evt-r1qj2e2yyrwe","type":"verification.recorded","subject":"ver-pgf465d","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"result":"pending","evidence":"User test after npm run build && npm start (or npm run dev): watch each camera, switch Main/Sub, and check the browser's network tab"},"at":"2026-09-19T21:35:07.531Z","parents":["evt-vcksn6nsjyhr"],"hash":"204c81adcc35d0725e4ecf38bc659526e87a2e20297238b039c03225573a455d"}
{"id":"evt-kmm00xbmnw9k","type":"node.status_changed","subject":"iss-cbx21zy","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"from":"open","to":"in_progress"},"at":"2026-09-19T21:35:09.464Z","parents":["evt-r1qj2e2yyrwe"],"hash":"5b8d3b666a7b0807cafaae7f631bb591c504a29213af8ac479153072fa81ef28"}
{"id":"evt-j1n23nmpk1yd","type":"node.created","subject":"mea-95yr6gy","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"kind":"measurement","title":"User watched both cameras live, on main and sub, in the browser","body":"Four of the goal's six checks, observed by the user on 2026-09-19 in their own browser after the live player shipped: cameras …af2e and …af54, each on the main (4096×1860) and sub (1200×536) streams, played live over WebRTC through the app's session-checked WHEP route (\"looks good\"). No firewall prompt appeared, which fits watching from the same Mac, where the loopback candidate is used. The two remaining checks are motion recording, one per camera, which isn't designed yet (iss-g456j72).","status":"recorded","owner":null,"attrs":{"value":4,"applies_at":"2026-09-19"}},"at":"2026-09-19T21:39:58.273Z","parents":["evt-kmm00xbmnw9k"],"hash":"7b614246944c81a570fc5469993b0b4117da1c1d02aa0defe736024e8aff700b"}
{"id":"evt-3tvvjzpx0fpx","type":"edge.added","subject":"mea-95yr6gy","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"kind":"measures","from":"mea-95yr6gy","to":"gol-sxakryh"},"at":"2026-09-19T21:39:58.275Z","parents":["evt-j1n23nmpk1yd"],"hash":"85b9d4aabec09896e90b5696e1ba0f3ae5d23251bcef503d259a7823241692fe"}
{"id":"evt-f58yed61x6g9","type":"node.created","subject":"ver-jra5mms","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"kind":"verification","title":"Real check by the user: both cameras play live in the pop-out and on the camera page, on main and sub","body":"","status":"pending","owner":"prn-q80g8mz","attrs":{}},"at":"2026-09-19T21:40:00.026Z","parents":["evt-3tvvjzpx0fpx"],"hash":"b8b09e4b2ac8ebac9d16b7909fb23ccc5d561362c2f740aca515bb9fe2ee601f"}
{"id":"evt-2mf7gbv44fmb","type":"edge.added","subject":"ver-jra5mms","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"kind":"evidence_for","from":"ver-jra5mms","to":"iss-cbx21zy"},"at":"2026-09-19T21:40:00.027Z","parents":["evt-f58yed61x6g9"],"hash":"d162456933fd05a2e6e019e5106a569ef7635cea9018d32972b1c3b143df207f"}
{"id":"evt-w404t95c5wxk","type":"verification.recorded","subject":"ver-jra5mms","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"result":"pass","evidence":"User on 2026-09-19: \"looks good\", both cameras, main and sub (measurement mea for gol-sxakryh, value 4 of 6)."},"at":"2026-09-19T21:40:00.028Z","parents":["evt-2mf7gbv44fmb"],"hash":"f54aa94f712f9045480d6c6a0e29d17c3c2791405805ecb3aafab5b8bd36f9d3"}
{"id":"evt-ew0n4dq70n0y","type":"node.status_changed","subject":"iss-cbx21zy","actor":"prn-q80g8mz","actor_kind":"agent","session":null,"payload":{"from":"in_progress","to":"done"},"at":"2026-09-19T21:40:01.764Z","parents":["evt-w404t95c5wxk"],"hash":"44defac6820b215f010e288c5dd6d875b0d6a0608439728ce7231e10999f0a6f"}

View File

@ -0,0 +1,174 @@
import { beforeEach, describe, expect, it, vi } from "vitest";
import { ctx, ID, record, target, url } from "../../../../../../../test/camera-routes";
const getCameraRecord = vi.fn();
const ensureStreamPath = vi.fn();
// Access control is tested in src/app/api/access.test.ts; here requests are allowed.
vi.mock("@/lib/access", () => ({ apiAccessDenied: async () => null }));
vi.mock("@/lib/camera-registry", async (importOriginal) => ({
...(await importOriginal<typeof import("@/lib/camera-registry")>()),
getCameraRecord,
}));
vi.mock("@/lib/video", async (importOriginal) => ({
...(await importOriginal<typeof import("@/lib/video")>()),
ensureStreamPath,
mediamtxAuthHeader: () => "Basic YXBwOnNlY3JldA==",
}));
const { POST, DELETE } = await import("./route");
const { VideoBridgeError } = await import("@/lib/video");
const OFFER = "v=0\r\no=- 1 1 IN IP4 0.0.0.0\r\nm=video 9 UDP/TLS/RTP/SAVPF 96\r\n";
const ANSWER = "v=0\r\no=- 2 2 IN IP4 127.0.0.1\r\n";
const SESSION = "123e4567-e89b-12d3-a456-426614174000";
const PATH = `cam-${ID}-main`;
const post = (query = "?stream=main", body = OFFER, type = "application/sdp", id?: string) =>
POST(new Request(url(`live/whep${query}`), { method: "POST", headers: { "Content-Type": type }, body }), ctx(id));
const del = (query: string, id?: string) =>
DELETE(new Request(url(`live/whep${query}`), { method: "DELETE" }), ctx(id));
/** A stand-in MediaMTX on localhost. */
function stubMediamtx(handler: (url: string, init: RequestInit) => Response) {
const mock = vi.fn(async (u: string | URL, init?: RequestInit) => handler(String(u), init ?? {}));
vi.stubGlobal("fetch", mock);
return mock;
}
const created = () =>
new Response(ANSWER, { status: 201, headers: { Location: `/${PATH}/whep/${SESSION}` } });
beforeEach(() => {
getCameraRecord.mockReset().mockResolvedValue(record);
ensureStreamPath.mockReset().mockResolvedValue(PATH);
});
describe("POST /api/cameras/[id]/live/whep", () => {
it("relays the offer to MediaMTX with the bridge login and returns the answer", async () => {
const fetchMock = stubMediamtx(() => created());
const res = await post();
expect(res.status).toBe(201);
expect(res.headers.get("Content-Type")).toBe("application/sdp");
expect(res.headers.get("Cache-Control")).toBe("no-store");
expect(res.headers.get("X-Whep-Session")).toBe(SESSION);
expect(await res.text()).toBe(ANSWER);
expect(ensureStreamPath).toHaveBeenCalledWith(target, "main");
const [calledUrl, init] = fetchMock.mock.calls[0];
expect(calledUrl).toBe(`http://127.0.0.1:8889/${PATH}/whep`);
expect(init!.body).toBe(OFFER);
expect((init!.headers as Record<string, string>).Authorization).toBe("Basic YXBwOnNlY3JldA==");
});
it("asks for the sub stream when requested, and defaults to main", async () => {
stubMediamtx(() => created());
await post("?stream=sub");
expect(ensureStreamPath).toHaveBeenLastCalledWith(target, "sub");
await post("?stream=nonsense");
expect(ensureStreamPath).toHaveBeenLastCalledWith(target, "main");
await post("");
expect(ensureStreamPath).toHaveBeenLastCalledWith(target, "main");
});
it("never returns a MediaMTX address or its error text", async () => {
const fetchMock = stubMediamtx(() => new Response(`path 'x' not found: rtsp://camera:hunter2@1.2.3.4`, { status: 404 }));
vi.spyOn(console, "warn").mockImplementation(() => {});
const res = await post();
expect(res.status).toBe(503);
const body = await res.text();
expect(body).toBe(JSON.stringify({ error: "The camera's live stream could not be started", code: "video" }));
expect(body).not.toContain("hunter2");
expect(body).not.toContain("127.0.0.1");
expect(fetchMock).toHaveBeenCalled();
});
it("says live video is unavailable when the bridge isn't running, without calling MediaMTX", async () => {
const fetchMock = stubMediamtx(() => created());
ensureStreamPath.mockRejectedValue(new VideoBridgeError("Video bridge is not running"));
vi.spyOn(console, "warn").mockImplementation(() => {});
const res = await post();
expect(res.status).toBe(503);
expect(await res.json()).toEqual({ error: "Live video is not available right now", code: "video" });
expect(fetchMock).not.toHaveBeenCalled();
});
it("reports camera problems with the shared error codes", async () => {
const { CameraAuthError } = await import("@/lib/camera");
ensureStreamPath.mockRejectedValue(new CameraAuthError("rejected"));
vi.spyOn(console, "warn").mockImplementation(() => {});
const res = await post();
expect(res.status).toBe(401);
expect(await res.json()).toMatchObject({ code: "auth" });
});
it("survives MediaMTX being unreachable", async () => {
vi.stubGlobal("fetch", vi.fn().mockRejectedValue(new Error("ECONNREFUSED")));
vi.spyOn(console, "warn").mockImplementation(() => {});
expect((await post()).status).toBe(503);
});
it("leaves the session header empty when MediaMTX sends no usable Location", async () => {
stubMediamtx(() => new Response(ANSWER, { status: 201 }));
expect((await post()).headers.get("X-Whep-Session")).toBe("");
});
it.each([
["a non-SDP content type", "?stream=main", OFFER, "application/json", 415],
["an empty offer", "?stream=main", "", "application/sdp", 400],
["an oversized offer", "?stream=main", "v".repeat(65_537), "application/sdp", 400],
])("rejects %s", async (_l, query, body, type, status) => {
const fetchMock = stubMediamtx(() => created());
expect((await post(query, body, type)).status).toBe(status);
expect(fetchMock).not.toHaveBeenCalled();
});
it("rejects a malformed id and an unknown camera before touching the bridge", async () => {
expect((await post("?stream=main", OFFER, "application/sdp", "../x")).status).toBe(400);
getCameraRecord.mockResolvedValue(null);
expect((await post()).status).toBe(404);
getCameraRecord.mockResolvedValue({ ...record, host: "8.8.8.8" });
expect((await post()).status).toBe(400);
expect(ensureStreamPath).not.toHaveBeenCalled();
});
});
describe("DELETE /api/cameras/[id]/live/whep", () => {
it("ends the session on MediaMTX using the path built from the registry", async () => {
const fetchMock = stubMediamtx(() => new Response(null, { status: 200 }));
const res = await del(`?stream=main&session=${SESSION}`);
expect(res.status).toBe(204);
const [calledUrl, init] = fetchMock.mock.calls[0];
expect(calledUrl).toBe(`http://127.0.0.1:8889/${PATH}/whep/${SESSION}`);
expect(init!.method).toBe("DELETE");
});
it.each(["", "?stream=main", "?stream=main&session=not-a-uuid", "?session=../../etc"])(
"rejects %j without calling MediaMTX",
async (query) => {
const fetchMock = stubMediamtx(() => new Response(null, { status: 200 }));
expect((await del(query)).status).toBe(400);
expect(fetchMock).not.toHaveBeenCalled();
},
);
it("still answers when the bridge is stopped or MediaMTX is unreachable", async () => {
const video = await import("@/lib/video");
const auth = vi.spyOn(video, "mediamtxAuthHeader").mockImplementation(() => {
throw new video.VideoBridgeError("Video bridge not started");
});
expect((await del(`?stream=sub&session=${SESSION}`)).status).toBe(503);
auth.mockRestore();
vi.stubGlobal("fetch", vi.fn().mockRejectedValue(new Error("ECONNREFUSED")));
vi.spyOn(console, "warn").mockImplementation(() => {});
expect((await del(`?stream=sub&session=${SESSION}`)).status).toBe(204);
});
it("checks the camera id like every other route", async () => {
getCameraRecord.mockResolvedValue(null);
expect((await del(`?stream=main&session=${SESSION}`)).status).toBe(404);
});
});

View File

@ -0,0 +1,105 @@
import { z } from "zod";
import { apiAccessDenied } from "@/lib/access";
import { cameraErrorResponse, cameraTarget } from "@/lib/camera-route";
import { mediamtxWebrtcUrl, pathName } from "@/lib/mediamtx-config";
import { ensureStreamPath, mediamtxAuthHeader, VideoBridgeError } from "@/lib/video";
/**
* WebRTC signaling (WHEP) for one camera stream, relayed to MediaMTX on localhost (vrek
* iss-cbx21zy). The browser never learns MediaMTX's address or a camera login: it posts
* its SDP offer here, and gets the answer back. Like every entry point, this checks the
* session and resolves the camera from the registry (pri-m1csgrm).
*/
const streamSchema = z.enum(["main", "sub"]).catch("main");
/** MediaMTX's session id: the uuid at the end of the Location it returns. */
const sessionSchema = z.string().uuid();
/** An SDP offer is a few KB; anything much larger isn't one. */
const MAX_SDP_BYTES = 64 * 1024;
const SDP = "application/sdp";
function unavailable(detail: string) {
return Response.json({ error: detail, code: "video" }, { status: 503 });
}
export async function POST(request: Request, ctx: RouteContext<"/api/cameras/[id]/live/whep">) {
const denied = await apiAccessDenied();
if (denied) return denied;
const target = await cameraTarget(ctx.params);
if (target instanceof Response) return target;
if (!request.headers.get("content-type")?.startsWith(SDP)) {
return Response.json({ error: "Expected an SDP offer" }, { status: 415 });
}
const offer = await request.text();
if (!offer || offer.length > MAX_SDP_BYTES) {
return Response.json({ error: "Invalid SDP offer" }, { status: 400 });
}
const stream = streamSchema.parse(new URL(request.url).searchParams.get("stream"));
let name: string;
try {
name = await ensureStreamPath(target, stream);
} catch (err) {
if (err instanceof VideoBridgeError) {
console.warn("[video]", err.message);
return unavailable("Live video is not available right now");
}
return cameraErrorResponse(err);
}
const res = await fetch(mediamtxWebrtcUrl(`/${name}/whep`), {
method: "POST",
headers: { "Content-Type": SDP, Authorization: mediamtxAuthHeader() },
body: offer,
signal: AbortSignal.timeout(15_000),
}).catch((err: unknown) => {
console.warn("[video] WHEP relay failed:", err);
return null;
});
if (!res || res.status !== 201) {
// MediaMTX's own error text can name paths and sources; keep it in the log.
if (res) console.warn(`[video] MediaMTX refused the WHEP offer (HTTP ${res.status})`);
return unavailable("The camera's live stream could not be started");
}
const session = res.headers.get("location")?.split("/").pop() ?? "";
return new Response(await res.text(), {
status: 201,
headers: {
"Content-Type": SDP,
"Cache-Control": "no-store",
// The browser sends this back to hang up; it only controls its own session.
"X-Whep-Session": sessionSchema.safeParse(session).success ? session : "",
},
});
}
/** Ends a session the browser started, so MediaMTX stops pulling the camera right away. */
export async function DELETE(request: Request, ctx: RouteContext<"/api/cameras/[id]/live/whep">) {
const denied = await apiAccessDenied();
if (denied) return denied;
const target = await cameraTarget(ctx.params);
if (target instanceof Response) return target;
const params = new URL(request.url).searchParams;
const session = sessionSchema.safeParse(params.get("session"));
if (!session.success) return Response.json({ error: "Invalid session" }, { status: 400 });
const stream = streamSchema.parse(params.get("stream"));
let auth: string;
try {
auth = mediamtxAuthHeader();
} catch {
return unavailable("Live video is not available right now");
}
// The path is built from the camera in the registry, never from the request.
await fetch(mediamtxWebrtcUrl(`/${pathName(target.id, stream)}/whep/${session.data}`), {
method: "DELETE",
headers: { Authorization: auth },
signal: AbortSignal.timeout(5_000),
}).catch((err: unknown) => console.warn("[video] WHEP teardown failed:", err));
return new Response(null, { status: 204 });
}

View File

@ -0,0 +1,23 @@
// @vitest-environment jsdom
import { fireEvent, render, screen } from "@testing-library/react";
import { describe, expect, it, vi } from "vitest";
import LivePanel from "./live-panel";
vi.mock("./live/live-player", () => ({
default: ({ cameraId, name, stream }: { cameraId: string; name: string; stream: string }) => (
<div data-testid="player">{`${name} ${cameraId} ${stream}`}</div>
),
}));
const ID = "11111111-2222-3333-4444-555555555555";
describe("LivePanel", () => {
it("plays the main stream for this camera, and can switch to the sub stream", () => {
render(<LivePanel cameraId={ID} name="Porch" />);
expect(screen.getByRole("heading", { name: "Live" })).toBeTruthy();
expect(screen.getByTestId("player").textContent).toBe(`Porch ${ID} main`);
fireEvent.change(screen.getByLabelText(/Stream/), { target: { value: "sub" } });
expect(screen.getByTestId("player").textContent).toBe(`Porch ${ID} sub`);
});
});

View File

@ -0,0 +1,23 @@
"use client";
import { useState } from "react";
import LivePlayer from "./live/live-player";
import StreamSelect from "./live/stream-select";
import type { StreamKind } from "./live/use-whep";
/** Live video on the camera's own page (vrek iss-cbx21zy). */
export default function LivePanel({ cameraId, name }: { cameraId: string; name: string }) {
const [stream, setStream] = useState<StreamKind>("main");
return (
<section>
<div className="flex flex-wrap items-center justify-between gap-2">
<h2 className="text-lg font-medium">Live</h2>
<StreamSelect value={stream} onChange={setStream} />
</div>
<div className="relative mt-2 aspect-video overflow-hidden rounded bg-black text-white">
<LivePlayer cameraId={cameraId} name={name} stream={stream} />
</div>
</section>
);
}

View File

@ -0,0 +1,198 @@
// @vitest-environment jsdom
// The browser's RTCPeerConnection is replaced by a scripted fake and fetch is stubbed, so
// no real peer connection or network is used (vrek pri-e14bahk).
import { act, fireEvent, screen, waitFor } from "@testing-library/react";
import { beforeEach, describe, expect, it, vi } from "vitest";
import { json, renderWithQuery, stubFetch } from "../../../../../test/dom";
import LivePlayer from "./live-player";
const m = vi.hoisted(() => ({ createPeerConnection: vi.fn() }));
vi.mock("./webrtc", () => ({ createPeerConnection: m.createPeerConnection }));
const createPeerConnection = m.createPeerConnection;
const ID = "11111111-2222-3333-4444-555555555555";
const SESSION = "123e4567-e89b-12d3-a456-426614174000";
const whep = () => `POST /api/cameras/${ID}/live/whep`;
class FakePeer extends EventTarget {
static last: FakePeer;
transceivers: string[] = [];
localDescription: { type: string; sdp: string } | null = null;
remoteDescription: RTCSessionDescriptionInit | null = null;
iceGatheringState: RTCIceGatheringState = "complete";
connectionState: RTCPeerConnectionState = "new";
closed = false;
constructor() {
super();
FakePeer.last = this;
}
addTransceiver(kind: string) {
this.transceivers.push(kind);
}
async createOffer() {
return { type: "offer", sdp: "offer-sdp" };
}
async setLocalDescription(desc: { type: string; sdp: string }) {
this.localDescription = desc;
}
async setRemoteDescription(desc: RTCSessionDescriptionInit) {
this.remoteDescription = desc;
}
close() {
this.closed = true;
}
/** Test helpers. */
connect(state: RTCPeerConnectionState) {
this.connectionState = state;
this.dispatchEvent(new Event("connectionstatechange"));
}
emitTrack(stream: unknown) {
const event = new Event("track") as Event & { streams: MediaStream[]; track: unknown };
Object.assign(event, { streams: [stream], track: {} });
this.dispatchEvent(event);
}
}
const answer = (headers: Record<string, string> = { "X-Whep-Session": SESSION }) =>
new Response("answer-sdp", { status: 201, headers });
beforeEach(() => {
createPeerConnection.mockReset().mockImplementation(() => new FakePeer());
});
const render = (props: Partial<Parameters<typeof LivePlayer>[0]> = {}) =>
renderWithQuery(<LivePlayer cameraId={ID} name="Porch" stream="main" {...props} />);
describe("LivePlayer", () => {
it("offers to our own route, plays the answer, and shows the picture", async () => {
const fetchMock = stubFetch({ [whep()]: () => answer() });
render();
expect(screen.getByRole("status").textContent).toContain("Connecting…");
await waitFor(() => expect(FakePeer.last.remoteDescription).toEqual({ type: "answer", sdp: "answer-sdp" }));
expect(FakePeer.last.transceivers).toEqual(["video", "audio"]);
const [url, init] = fetchMock.mock.calls[0];
expect(String(url)).toBe(`/api/cameras/${ID}/live/whep?stream=main`);
expect((init!.headers as Record<string, string>)["Content-Type"]).toBe("application/sdp");
expect(init!.body).toBe("offer-sdp");
// jsdom has no MediaStream; the player only passes it to <video>.
const media = { id: "fake-stream" } as unknown as MediaStream;
await act(async () => {
FakePeer.last.emitTrack(media);
FakePeer.last.connect("connected");
});
const video = screen.getByLabelText("Porch live") as HTMLVideoElement;
expect(video.srcObject).toBe(media);
expect(screen.queryByRole("status")).toBeNull();
});
it("asks for the stream it was given", async () => {
const fetchMock = stubFetch({ [whep()]: () => answer() });
render({ stream: "sub" });
await waitFor(() => expect(fetchMock).toHaveBeenCalled());
expect(String(fetchMock.mock.calls[0][0])).toBe(`/api/cameras/${ID}/live/whep?stream=sub`);
});
it("waits for ICE candidates before offering, but not forever", async () => {
vi.useFakeTimers({ shouldAdvanceTime: true });
try {
const fetchMock = stubFetch({ [whep()]: () => answer() });
createPeerConnection.mockImplementation(() => {
const peer = new FakePeer();
peer.iceGatheringState = "gathering";
return peer;
});
render();
await act(async () => {});
expect(fetchMock).not.toHaveBeenCalled();
await act(() => vi.advanceTimersByTimeAsync(1600));
expect(fetchMock).toHaveBeenCalledTimes(1);
} finally {
vi.useRealTimers();
}
});
it("sends the offer as soon as gathering finishes", async () => {
const fetchMock = stubFetch({ [whep()]: () => answer() });
createPeerConnection.mockImplementation(() => {
const peer = new FakePeer();
peer.iceGatheringState = "gathering";
return peer;
});
render();
await act(async () => {});
await act(async () => {
FakePeer.last.iceGatheringState = "complete";
FakePeer.last.dispatchEvent(new Event("icegatheringstatechange"));
});
await waitFor(() => expect(fetchMock).toHaveBeenCalled());
});
it("shows reconnecting, then a failure with Retry that starts a new session", async () => {
stubFetch({ [whep()]: () => answer() });
const onFailed = vi.fn();
render({ onFailed });
await waitFor(() => expect(FakePeer.last.remoteDescription).not.toBeNull());
await act(async () => FakePeer.last.connect("disconnected"));
expect(screen.getByRole("status").textContent).toContain("Reconnecting…");
await act(async () => FakePeer.last.connect("failed"));
expect(screen.getByRole("alert").textContent).toContain("The live connection dropped.");
expect(onFailed).toHaveBeenLastCalledWith(true);
const first = FakePeer.last;
fireEvent.click(screen.getByRole("button", { name: "Try live again" }));
await waitFor(() => expect(FakePeer.last).not.toBe(first));
expect(first.closed).toBe(true);
expect(screen.queryByRole("alert")).toBeNull();
});
it("explains a refused stream with the server's own message", async () => {
stubFetch({ [whep()]: () => json({ error: "Live video is not available right now", code: "video" }, 503) });
render();
expect((await screen.findByRole("alert")).textContent).toContain("Live video is not available right now");
});
it("tells a signed-out user to sign in", async () => {
stubFetch({ [whep()]: () => json({ error: "Sign in required" }, 401) });
render();
expect((await screen.findByRole("alert")).textContent).toContain("Sign in again to watch live video");
});
it("falls back to fixed text when the server sends no message", async () => {
stubFetch({ [whep()]: () => new Response("nope", { status: 500 }) });
render();
expect((await screen.findByRole("alert")).textContent).toContain("Live video could not be started.");
});
it("closes the connection and hangs up the session when it goes away", async () => {
const fetchMock = stubFetch({
[whep()]: () => answer(),
[`DELETE /api/cameras/${ID}/live/whep`]: () => new Response(null, { status: 204 }),
});
const { unmount } = render();
await waitFor(() => expect(FakePeer.last.remoteDescription).not.toBeNull());
const peer = FakePeer.last;
unmount();
expect(peer.closed).toBe(true);
await waitFor(() => {
const hangup = fetchMock.mock.calls.find(([, init]) => init?.method === "DELETE");
expect(String(hangup?.[0])).toBe(`/api/cameras/${ID}/live/whep?stream=main&session=${SESSION}`);
expect(hangup?.[1]?.keepalive).toBe(true);
});
});
it("doesn't hang up when it never got a session", async () => {
const fetchMock = stubFetch({ [whep()]: () => answer({}) });
const { unmount } = render();
await waitFor(() => expect(FakePeer.last.remoteDescription).not.toBeNull());
unmount();
await new Promise((r) => setTimeout(r, 20));
expect(fetchMock.mock.calls.some(([, init]) => init?.method === "DELETE")).toBe(false);
});
});

View File

@ -0,0 +1,85 @@
"use client";
import { CircleAlert, LoaderCircle, RefreshCw } from "lucide-react";
import { useCallback, useEffect, useRef, useState } from "react";
import { useWhep, type StreamKind } from "./use-whep";
/**
* The live picture for one camera stream, with its connection state (vrek iss-cbx21zy).
* `onFailed` lets the page fall back to snapshot polling.
*/
export default function LivePlayer(props: {
cameraId: string;
name: string;
stream: StreamKind;
onFailed?: (failed: boolean) => void;
}) {
// Retrying remounts the session, so its state starts fresh.
const [attempt, setAttempt] = useState(0);
const retry = useCallback(() => setAttempt((n) => n + 1), []);
return <LiveSession key={`${props.stream}-${attempt}`} {...props} retry={retry} />;
}
function LiveSession({
cameraId,
name,
stream,
onFailed,
retry,
}: {
cameraId: string;
name: string;
stream: StreamKind;
onFailed?: (failed: boolean) => void;
retry: () => void;
}) {
const playback = useWhep(cameraId, stream);
const videoRef = useRef<HTMLVideoElement>(null);
useEffect(() => {
const video = videoRef.current;
if (video) video.srcObject = playback.stream;
}, [playback.stream]);
const failed = playback.state === "failed";
useEffect(() => onFailed?.(failed), [failed, onFailed]);
return (
<div className="absolute inset-0 flex items-center justify-center">
<video
ref={videoRef}
autoPlay
playsInline
muted
aria-label={`${name} live`}
className={`max-h-full max-w-full ${playback.state === "playing" ? "" : "opacity-40"}`}
/>
{playback.state !== "playing" && (
<div
role={failed ? "alert" : "status"}
className="absolute inset-x-0 bottom-0 flex flex-wrap items-center justify-between gap-2 bg-black/70 px-4 py-2 text-sm"
>
<span className={`inline-flex items-center gap-1.5 ${failed ? "text-red-300" : ""}`}>
{failed ? (
<CircleAlert aria-hidden className="size-4 shrink-0" />
) : (
<LoaderCircle aria-hidden className="size-4 shrink-0 animate-spin motion-reduce:animate-none" />
)}
{failed
? (playback.message ?? "Live video could not be started.")
: playback.state === "reconnecting"
? "Reconnecting…"
: "Connecting…"}
</span>
{failed && (
<button onClick={retry} className="inline-flex items-center gap-1 underline">
<RefreshCw aria-hidden className="size-4" />
Try live again
</button>
)}
</div>
)}
</div>
);
}

View File

@ -9,21 +9,56 @@ vi.mock("next/navigation", () => ({
usePathname: () => "/cameras/x/live",
}));
// The player has its own tests; here it is a marker that can report a failure.
const m = vi.hoisted(() => ({ onFailed: undefined as ((failed: boolean) => void) | undefined }));
vi.mock("./live-player", () => ({
default: ({ stream, onFailed }: { stream: string; onFailed?: (failed: boolean) => void }) => {
m.onFailed = onFailed;
return <div data-testid="player">live {stream}</div>;
},
}));
const ID = "11111111-2222-3333-4444-555555555555";
const snap = `GET /api/cameras/${ID}/snapshot`;
const frame = () => new Response(new Blob(["jpeg"], { type: "image/jpeg" }));
const overlay = () => screen.getByText("Porch").closest("div")!;
const view = (intervalMs = 1000) => renderWithQuery(<LiveView id={ID} name="Porch" intervalMs={intervalMs} />);
/** Makes the player report that live video failed, so the pop-out falls back to snapshots. */
const failLive = () => act(async () => m.onFailed!(true));
beforeEach(() => {
m.onFailed = undefined;
vi.spyOn(URL, "createObjectURL").mockReturnValue("blob:frame");
vi.spyOn(URL, "revokeObjectURL").mockImplementation(() => {});
});
describe("LiveView", () => {
it("shows the live frame, the camera name, the frame time and a rate selector", async () => {
it("plays the main stream live, with the camera name and a stream switch", async () => {
const fetchMock = stubFetch({ [snap]: frame });
view();
expect(screen.getByTestId("player").textContent).toBe("live main");
expect(screen.getByText("Porch")).toBeTruthy();
expect((screen.getByRole("combobox") as HTMLSelectElement).value).toBe("main");
// No snapshot polling while live video is playing.
await new Promise((r) => setTimeout(r, 30));
expect(calls(fetchMock, snap)).toHaveLength(0);
});
it("switches between the main and sub streams", () => {
stubFetch({ [snap]: frame });
renderWithQuery(<LiveView id={ID} name="Porch" intervalMs={1000} />);
expect(screen.getByText("Loading…")).toBeTruthy();
view();
fireEvent.change(screen.getByRole("combobox"), { target: { value: "sub" } });
expect(screen.getByTestId("player").textContent).toBe("live sub");
});
it("falls back to snapshots when live video fails, and says so", async () => {
stubFetch({ [snap]: frame });
view();
await failLive();
expect(screen.queryByTestId("player")).toBeNull();
expect(screen.getByText("Snapshots (live video unavailable)")).toBeTruthy();
await waitFor(() => expect(screen.getByRole("img", { name: "Porch" }).getAttribute("src")).toBe("blob:frame"));
// Never stretched beyond the frame's own size (fnd-dw1pqcw).
const img = screen.getByRole("img", { name: "Porch" });
@ -38,7 +73,7 @@ describe("LiveView", () => {
vi.useFakeTimers({ shouldAdvanceTime: true });
try {
stubFetch({ [snap]: frame });
renderWithQuery(<LiveView id={ID} name="Porch" intervalMs={5000} />);
view(5000);
expect(overlay().className).toContain("opacity-100");
await act(() => vi.advanceTimersByTimeAsync(3000));
expect(overlay().className).toContain("opacity-0");
@ -54,7 +89,9 @@ describe("LiveView", () => {
it("stops on a failed frame, keeps the overlay up, and retries on request", async () => {
let fail = true;
const fetchMock = stubFetch({ [snap]: () => (fail ? json({ error: "Camera timed out" }, 502) : frame()) });
renderWithQuery(<LiveView id={ID} name="Porch" intervalMs={20} />);
view(20);
await failLive();
expect((await screen.findByRole("alert")).textContent).toContain("Stopped: Camera timed out");
const afterFailure = calls(fetchMock, snap).length;
await new Promise((r) => setTimeout(r, 100));
@ -68,7 +105,8 @@ describe("LiveView", () => {
it("sends a signed-out user to sign in and back to this window", async () => {
stubFetch({ [snap]: () => json({ error: "Sign in required" }, 401) });
renderWithQuery(<LiveView id={ID} name="Porch" intervalMs={1000} />);
view();
await failLive();
expect((await screen.findByRole("alert")).textContent).toContain("Your session has ended.");
expect(screen.getByRole("link", { name: "Sign in again" }).getAttribute("href")).toBe(
`/login?next=${encodeURIComponent(`/cameras/${ID}/live`)}`,
@ -77,7 +115,8 @@ describe("LiveView", () => {
it("points camera login or setup problems back to the dashboard", async () => {
stubFetch({ [snap]: () => json({ error: "bad login", code: "auth" }, 401) });
renderWithQuery(<LiveView id={ID} name="Porch" intervalMs={1000} />);
view();
await failLive();
expect((await screen.findByRole("alert")).textContent).toContain("Fix it from the dashboard");
});
@ -86,7 +125,7 @@ describe("LiveView", () => {
const seen: unknown[] = [];
const card = new BroadcastChannel("camera-popouts");
card.onmessage = (e) => seen.push(e.data);
const { unmount } = renderWithQuery(<LiveView id={ID} name="Porch" intervalMs={1000} />);
const { unmount } = view();
await waitFor(() => expect(seen).toContainEqual({ type: "open", id: ID }));
unmount();
await waitFor(() => expect(seen).toContainEqual({ type: "closed", id: ID }));

View File

@ -6,6 +6,9 @@ import { useCameraSnapshot, useResetCamera } from "../../../camera-queries";
import FrameImage from "../../../frame-image";
import { announcePopout } from "../../../popout";
import RefreshRateSelect from "../../../refresh-rate-select";
import LivePlayer from "./live-player";
import StreamSelect from "./stream-select";
import type { StreamKind } from "./use-whep";
/** How long the overlay stays after the mouse stops moving. */
const OVERLAY_MS = 2500;
@ -26,16 +29,29 @@ function useRecentActivity() {
return { active, poke };
}
export default function LiveView({ id, name, intervalMs }: { id: string; name: string; intervalMs: number }) {
/** Snapshot polling: what the pop-out falls back to when live video can't play. */
function SnapshotView({
id,
name,
intervalMs,
onUpdated,
}: {
id: string;
name: string;
intervalMs: number;
onUpdated: (at: number) => void;
}) {
const snapshot = useCameraSnapshot(id, intervalMs, true);
const reset = useResetCamera(id);
const { active, poke } = useRecentActivity();
const problem = snapshot.error;
const updatedAt = snapshot.dataUpdatedAt;
useEffect(() => announcePopout(id), [id]);
useEffect(() => {
if (updatedAt > 0) onUpdated(updatedAt);
}, [updatedAt, onUpdated]);
return (
<div className="absolute inset-0" onMouseMove={poke}>
<>
{snapshot.data && (
// Shown at its own size at most: stretching a small camera snapshot only blurs it
// (vrek fnd-dw1pqcw). Larger windows get a border instead.
@ -51,24 +67,6 @@ export default function LiveView({ id, name, intervalMs }: { id: string; name: s
<div className="absolute inset-0 flex items-center justify-center text-sm text-zinc-400">Loading…</div>
)}
<div
className={`absolute inset-x-0 top-0 flex flex-wrap items-center justify-between gap-3 bg-black/60 px-4 py-2 text-sm transition-opacity motion-reduce:transition-none ${
active || problem ? "opacity-100" : "opacity-0"
} focus-within:opacity-100`}
>
<span className="font-medium">
{name}
{snapshot.dataUpdatedAt > 0 && (
<span className="ml-2 font-normal text-zinc-300">
<time dateTime={new Date(snapshot.dataUpdatedAt).toISOString()}>
{new Date(snapshot.dataUpdatedAt).toLocaleTimeString()}
</time>
</span>
)}
</span>
<RefreshRateSelect value={intervalMs} />
</div>
{problem && (
<div
role="alert"
@ -83,7 +81,10 @@ export default function LiveView({ id, name, intervalMs }: { id: string; name: s
: "This camera needs attention. Fix it from the dashboard, then retry."}
</span>
{problem.kind === "signed-out" ? (
<a href={`/login?next=${encodeURIComponent(`/cameras/${id}/live`)}`} className="inline-flex items-center gap-1 underline">
<a
href={`/login?next=${encodeURIComponent(`/cameras/${id}/live`)}`}
className="inline-flex items-center gap-1 underline"
>
<LogIn aria-hidden className="size-4" />
Sign in again
</a>
@ -95,6 +96,48 @@ export default function LiveView({ id, name, intervalMs }: { id: string; name: s
)}
</div>
)}
</>
);
}
export default function LiveView({ id, name, intervalMs }: { id: string; name: string; intervalMs: number }) {
const { active, poke } = useRecentActivity();
const [stream, setStream] = useState<StreamKind>("main");
const [liveFailed, setLiveFailed] = useState(false);
const [frameAt, setFrameAt] = useState(0);
useEffect(() => announcePopout(id), [id]);
return (
<div className="absolute inset-0" onMouseMove={poke}>
{liveFailed ? (
<SnapshotView id={id} name={name} intervalMs={intervalMs} onUpdated={setFrameAt} />
) : (
<LivePlayer cameraId={id} name={name} stream={stream} onFailed={setLiveFailed} />
)}
<div
className={`absolute inset-x-0 top-0 flex flex-wrap items-center justify-between gap-3 bg-black/60 px-4 py-2 text-sm transition-opacity motion-reduce:transition-none ${
active || liveFailed ? "opacity-100" : "opacity-0"
} focus-within:opacity-100`}
>
<span className="font-medium">
{name}
{liveFailed && frameAt > 0 && (
<span className="ml-2 font-normal text-zinc-300">
<time dateTime={new Date(frameAt).toISOString()}>{new Date(frameAt).toLocaleTimeString()}</time>
</span>
)}
</span>
{liveFailed ? (
<span className="flex items-center gap-3">
<span className="text-zinc-300">Snapshots (live video unavailable)</span>
<RefreshRateSelect value={intervalMs} />
</span>
) : (
<StreamSelect value={stream} onChange={setStream} />
)}
</div>
</div>
);
}

View File

@ -0,0 +1,28 @@
"use client";
import { Video } from "lucide-react";
import type { StreamKind } from "./use-whep";
/** Main is the camera's full-resolution stream; Sub is its smaller, lighter one. */
export default function StreamSelect({
value,
onChange,
}: {
value: StreamKind;
onChange: (stream: StreamKind) => void;
}) {
return (
<label className="flex items-center gap-2">
<Video aria-hidden className="size-4" />
Stream
<select
value={value}
onChange={(e) => onChange(e.target.value as StreamKind)}
className="rounded border border-zinc-600 bg-black/40 px-2 py-1"
>
<option value="main">Main (full size)</option>
<option value="sub">Sub (smaller)</option>
</select>
</label>
);
}

View File

@ -0,0 +1,116 @@
"use client";
import { useEffect, useRef, useState } from "react";
import { createPeerConnection } from "./webrtc";
/**
* Plays one camera stream over WebRTC by posting an SDP offer to our own WHEP route
* (vrek iss-cbx21zy). The route relays it to MediaMTX; the browser never sees MediaMTX
* or a camera login.
*/
export type PlayerState = "connecting" | "playing" | "reconnecting" | "failed";
export type StreamKind = "main" | "sub";
/** How long to wait for ICE candidates before sending the offer anyway. */
const GATHER_MS = 1_500;
function whepUrl(cameraId: string, stream: StreamKind, session?: string) {
const query = new URLSearchParams({ stream, ...(session ? { session } : {}) });
return `/api/cameras/${cameraId}/live/whep?${query}`;
}
/** Resolves once ICE gathering finishes, or after GATHER_MS with what it has. */
function iceGathered(pc: RTCPeerConnection): Promise<void> {
if (pc.iceGatheringState === "complete") return Promise.resolve();
return new Promise((resolve) => {
const done = () => {
clearTimeout(timer);
pc.removeEventListener("icegatheringstatechange", check);
resolve();
};
const check = () => pc.iceGatheringState === "complete" && done();
const timer = setTimeout(done, GATHER_MS);
pc.addEventListener("icegatheringstatechange", check);
});
}
export interface WhepPlayback {
state: PlayerState;
stream: MediaStream | null;
/** Message for the failed state; fixed text, never an upstream error. */
message: string | null;
}
/**
* One playback attempt. Retrying means remounting the component that calls this, which
* starts the state over without touching state from inside an effect.
*/
export function useWhep(cameraId: string, stream: StreamKind): WhepPlayback {
const [state, setState] = useState<PlayerState>("connecting");
const [media, setMedia] = useState<MediaStream | null>(null);
const [message, setMessage] = useState<string | null>(null);
const sessionRef = useRef<string | null>(null);
useEffect(() => {
let cancelled = false;
const pc = createPeerConnection();
const abort = new AbortController();
pc.addEventListener("track", (e: RTCTrackEvent) => {
if (!cancelled) setMedia(e.streams[0] ?? new MediaStream([e.track]));
});
pc.addEventListener("connectionstatechange", () => {
if (cancelled) return;
if (pc.connectionState === "connected") setState("playing");
else if (pc.connectionState === "disconnected") setState("reconnecting");
else if (pc.connectionState === "failed") {
setState("failed");
setMessage("The live connection dropped.");
}
});
(async () => {
pc.addTransceiver("video", { direction: "recvonly" });
pc.addTransceiver("audio", { direction: "recvonly" });
await pc.setLocalDescription(await pc.createOffer());
await iceGathered(pc);
if (cancelled) return;
const res = await fetch(whepUrl(cameraId, stream), {
method: "POST",
headers: { "Content-Type": "application/sdp" },
body: pc.localDescription?.sdp ?? "",
signal: abort.signal,
});
if (!res.ok) {
const body = await res.json().catch(() => ({}));
throw new Error(
res.status === 401
? "Your session has ended. Sign in again to watch live video."
: (body.error ?? "Live video could not be started."),
);
}
sessionRef.current = res.headers.get("X-Whep-Session") || null;
await pc.setRemoteDescription({ type: "answer", sdp: await res.text() });
})().catch((err: Error) => {
if (cancelled) return;
setState("failed");
setMessage(err.message);
});
return () => {
cancelled = true;
abort.abort();
const session = sessionRef.current;
sessionRef.current = null;
pc.close();
if (session) {
// keepalive so the hang-up still goes out while the page is closing.
void fetch(whepUrl(cameraId, stream, session), { method: "DELETE", keepalive: true }).catch(() => {});
}
};
}, [cameraId, stream]);
return { state, stream: media, message };
}

View File

@ -0,0 +1,8 @@
/**
* The browser's WebRTC entry point, in its own module so tests can replace it (vrek
* pri-e14bahk). No STUN or TURN servers: MediaMTX is on the LAN, so its own candidates
* are enough and nothing about this network is sent to a third party.
*/
export function createPeerConnection(): RTCPeerConnection {
return new RTCPeerConnection({ iceServers: [] });
}

View File

@ -11,6 +11,10 @@ vi.mock("@/lib/camera-registry", async (importOriginal) => ({
getCameraRecord,
}));
vi.mock("./live-panel", () => ({
default: (props: object) => <div data-testid="live">{JSON.stringify(props)}</div>,
}));
// The async sections are tested on their own; here they are a marker carrying the target.
vi.mock("./settings-sections", () => ({
CameraSettingsSections: ({ target }: { target: object }) => (
@ -63,6 +67,7 @@ describe("/cameras/[id]", () => {
expect(html).toContain('href="/"');
expect(html).toContain("All cameras");
expect(html).toContain(">Refresh</button>");
expect(html.replace(/&quot;/g, '"')).toContain(JSON.stringify({ cameraId: ID, name: "I91ET" }));
expect(html).toContain(JSON.stringify({ id: ID, host: "192.168.17.129", port: 80 }).replace(/"/g, "&quot;"));
for (const svg of html.match(/<svg[^>]*>/g)!) expect(svg).toContain('aria-hidden="true"');
});

View File

@ -6,6 +6,7 @@ import { requirePageAccess } from "@/lib/access";
import { isAllowedHost } from "@/lib/camera";
import { getCameraRecord, isValidCameraId } from "@/lib/camera-registry";
import { RefreshButton } from "./client-controls";
import LivePanel from "./live-panel";
import { CameraSettingsSections } from "./settings-sections";
/**
@ -41,6 +42,10 @@ export default async function CameraPage({ params }: PageProps<"/cameras/[id]">)
<RefreshButton />
</div>
<div className="mt-6">
<LivePanel cameraId={record.id} name={record.name ?? "Unnamed camera"} />
</div>
<div className="mt-6">
<Suspense
fallback={

View File

@ -0,0 +1,126 @@
// Process lookups and signals are fakes; no real process is inspected or signalled except
// this test process's own, read-only (vrek pri-e14bahk).
import { mkdtemp, readFile, rm, stat, writeFile } from "node:fs/promises";
import os from "node:os";
import path from "node:path";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import {
defaultProcessDeps,
forgetPid,
isOurMediamtx,
recordPid,
stopLeftoverMediamtx,
type ProcessDeps,
} from "./mediamtx-leftover";
const BIN = "/app/bin/mediamtx";
const CONF = "/app/.data/mediamtx/mediamtx.yml";
const OURS = `${BIN} ${CONF}`;
let dir: string;
beforeEach(async () => {
dir = await mkdtemp(path.join(os.tmpdir(), "leftover-"));
});
afterEach(() => rm(dir, { recursive: true, force: true }));
const pidFile = () => path.join(dir, "mediamtx.pid");
const gone = async () => expect(stat(pidFile())).rejects.toThrow();
/** A fake process table: `commands` answers commandOf in order, repeating the last. */
function processes(...commands: (string | null)[]) {
let i = 0;
return {
commandOf: vi.fn<ProcessDeps["commandOf"]>(async () => commands[Math.min(i++, commands.length - 1)]),
kill: vi.fn<ProcessDeps["kill"]>(),
sleep: vi.fn<ProcessDeps["sleep"]>(async () => {}),
};
}
describe("recordPid / forgetPid", () => {
it("writes the pid owner-only and removes it", async () => {
await recordPid(dir, 4242);
expect(await readFile(pidFile(), "utf8")).toBe("4242\n");
expect((await stat(pidFile())).mode & 0o777).toBe(0o600);
await forgetPid(dir);
await gone();
await forgetPid(dir); // already gone: fine
});
});
describe("isOurMediamtx", () => {
it.each([
[OURS, true],
[`${BIN} /other/mediamtx.yml`, false],
[`/usr/local/bin/mediamtx ${CONF}`, false],
["/usr/sbin/sshd -D", false],
[`${BIN}x ${CONF}`, false],
])("%s → %s", (command, ours) => {
expect(isOurMediamtx(command, BIN, CONF)).toBe(ours);
});
});
describe("stopLeftoverMediamtx", () => {
it("does nothing without a pid file", async () => {
const deps = processes(OURS);
expect(await stopLeftoverMediamtx(dir, BIN, CONF, deps)).toBe("none");
expect(deps.commandOf).not.toHaveBeenCalled();
});
it.each(["", "garbage", "1", "-5", "3.5"])("ignores and removes a bad pid file %j", async (text) => {
await writeFile(pidFile(), text);
const deps = processes(OURS);
expect(await stopLeftoverMediamtx(dir, BIN, CONF, deps)).toBe("none");
expect(deps.kill).not.toHaveBeenCalled();
await gone();
});
it("forgets a pid that's no longer running", async () => {
await recordPid(dir, 4242);
const deps = processes(null);
expect(await stopLeftoverMediamtx(dir, BIN, CONF, deps)).toBe("none");
expect(deps.kill).not.toHaveBeenCalled();
await gone();
});
it("never signals a process that isn't our MediaMTX (a reused pid)", async () => {
await recordPid(dir, 4242);
const deps = processes("/usr/sbin/sshd -D");
expect(await stopLeftoverMediamtx(dir, BIN, CONF, deps)).toBe("not-ours");
expect(deps.kill).not.toHaveBeenCalled();
await gone();
});
it("stops our leftover MediaMTX with SIGTERM", async () => {
await recordPid(dir, 4242);
const deps = processes(OURS, OURS, null);
expect(await stopLeftoverMediamtx(dir, BIN, CONF, deps)).toBe("stopped");
expect(deps.kill.mock.calls).toEqual([[4242, "SIGTERM"]]);
await gone();
});
it("kills it if it ignores SIGTERM for about 3 s", async () => {
await recordPid(dir, 4242);
const deps = processes(OURS);
expect(await stopLeftoverMediamtx(dir, BIN, CONF, deps)).toBe("killed");
expect(deps.kill.mock.calls).toEqual([
[4242, "SIGTERM"],
[4242, "SIGKILL"],
]);
expect(deps.sleep).toHaveBeenCalledTimes(30);
await gone();
});
});
describe("defaultProcessDeps", () => {
it("reads a running process's command line, and null for a missing one", async () => {
expect(await defaultProcessDeps.commandOf(process.pid)).toContain("node");
expect(await defaultProcessDeps.commandOf(2 ** 22 + 12345)).toBeNull();
});
it("signals through process.kill and sleeps with a timer", async () => {
const kill = vi.spyOn(process, "kill").mockImplementation(() => true);
defaultProcessDeps.kill(4242, "SIGTERM");
expect(kill).toHaveBeenCalledWith(4242, "SIGTERM");
await defaultProcessDeps.sleep(1);
});
});

View File

@ -0,0 +1,87 @@
import "server-only";
import { execFile } from "node:child_process";
import { readFile, rm, writeFile } from "node:fs/promises";
import path from "node:path";
import { promisify } from "node:util";
/**
* Finds and stops a MediaMTX left running by an earlier server that died without a clean
* exit (vrek iss-yd2sq2q), so the new one can bind its ports. The running MediaMTX's pid
* is recorded in a file; a leftover is stopped only if that pid is still running our
* binary with our config, so an unrelated process that reused the pid is never touched.
*/
export const PID_FILE = "mediamtx.pid";
export interface ProcessDeps {
/** The process's full command line, or null if it isn't running. */
commandOf(pid: number): Promise<string | null>;
kill(pid: number, signal: NodeJS.Signals): void;
sleep(ms: number): Promise<void>;
}
export const defaultProcessDeps: ProcessDeps = {
async commandOf(pid) {
try {
const { stdout } = await promisify(execFile)("ps", ["-p", String(pid), "-o", "command="]);
return stdout.trim() || null;
} catch {
return null; // ps exits non-zero when there's no such process
}
},
kill: (pid, signal) => process.kill(pid, signal),
sleep: (ms) => new Promise((r) => setTimeout(r, ms)),
};
export async function recordPid(dir: string, pid: number): Promise<void> {
await writeFile(path.join(dir, PID_FILE), `${pid}\n`, { mode: 0o600 });
}
export async function forgetPid(dir: string): Promise<void> {
await rm(path.join(dir, PID_FILE), { force: true });
}
/** Whether a command line is our MediaMTX: our binary, run with our config file. */
export function isOurMediamtx(command: string, binary: string, configPath: string): boolean {
return command.startsWith(`${binary} `) && command.endsWith(` ${configPath}`);
}
/**
* Stops the MediaMTX recorded in `dir` if it's still ours and still running: SIGTERM,
* then SIGKILL if it hasn't exited after about 3 s. Returns what it did.
*/
export async function stopLeftoverMediamtx(
dir: string,
binary: string,
configPath: string,
deps: ProcessDeps = defaultProcessDeps,
): Promise<"none" | "not-ours" | "stopped" | "killed"> {
const text = await readFile(path.join(dir, PID_FILE), "utf8").catch(() => null);
const pid = Number(text?.trim());
if (!Number.isInteger(pid) || pid <= 1) {
await forgetPid(dir);
return "none";
}
const command = await deps.commandOf(pid);
if (command === null) {
await forgetPid(dir);
return "none";
}
if (!isOurMediamtx(command, binary, configPath)) {
await forgetPid(dir);
return "not-ours";
}
deps.kill(pid, "SIGTERM");
for (let i = 0; i < 30; i++) {
await deps.sleep(100);
if ((await deps.commandOf(pid)) === null) {
await forgetPid(dir);
return "stopped";
}
}
deps.kill(pid, "SIGKILL");
await forgetPid(dir);
return "killed";
}

View File

@ -6,6 +6,7 @@ import { beforeEach, describe, expect, it, vi } from "vitest";
import { MediamtxSupervisor, type SupervisorOptions } from "./mediamtx-supervisor";
class FakeChild extends EventEmitter {
pid = 4000 + children.length;
stdout = new PassThrough();
stderr = new PassThrough();
kill = vi.fn();
@ -147,6 +148,18 @@ describe("MediamtxSupervisor", () => {
expect(timers).toHaveLength(1);
});
it("reports each process's pid when it starts and when it ends", () => {
const onSpawn = vi.fn();
const onExit = vi.fn();
const { sup } = supervisor({ onSpawn, onExit });
sup.start();
expect(onSpawn).toHaveBeenCalledWith(4000);
children[0].emit("exit", 1, null);
expect(onExit).toHaveBeenCalledTimes(1);
timers[0].fn();
expect(onSpawn).toHaveBeenLastCalledWith(4001);
});
it("stops MediaMTX and doesn't restart it", () => {
const { sup } = supervisor();
sup.start();

View File

@ -27,6 +27,9 @@ export interface SupervisorOptions {
clearTimer?: (handle: unknown) => void;
now?: () => number;
log?: (line: string) => void;
/** Called with each MediaMTX process's pid when it starts, and when it ends. */
onSpawn?: (pid: number) => void;
onExit?: () => void;
}
const MIN_BACKOFF_MS = 1_000;
@ -87,6 +90,7 @@ export class MediamtxSupervisor {
stdio: ["ignore", "pipe", "pipe"],
});
this.child = child;
if (child.pid !== undefined) this.opts.onSpawn?.(child.pid);
const forward = (chunk: Buffer) => {
for (const line of chunk.toString("utf8").split("\n")) if (line.trim()) this.log(`[mediamtx] ${line.trimEnd()}`);
@ -107,6 +111,7 @@ export class MediamtxSupervisor {
});
child.once("exit", (code, signal) => {
this.opts.onExit?.();
if (this.child !== child) return;
this.child = null;
if (this.stopping) return;

View File

@ -102,6 +102,44 @@ describe("startVideoBridge", () => {
});
});
describe("leftover MediaMTX and pid file (iss-yd2sq2q)", () => {
it("stops a leftover before starting, and says so", async () => {
const stopLeftover = vi.fn(async () => "stopped" as const);
const log = vi.spyOn(console, "log").mockImplementation(() => {});
await video.startVideoBridge(factory, stopLeftover);
const bridgeDir = path.join(dir, ".data", "mediamtx");
expect(stopLeftover).toHaveBeenCalledWith(bridgeDir, path.join(dir, "bin", "mediamtx"), path.join(bridgeDir, "mediamtx.yml"));
expect(stopLeftover.mock.invocationCallOrder[0]).toBeLessThan(fake.start.mock.invocationCallOrder[0]);
expect(log).toHaveBeenCalledWith("[video] stopped a MediaMTX left running by an earlier server (stopped)");
});
it("stays quiet when there was nothing to stop", async () => {
const log = vi.spyOn(console, "log").mockImplementation(() => {});
await video.startVideoBridge(factory, async () => "not-ours");
expect(log).not.toHaveBeenCalled();
});
it("records the running pid, and removes it when MediaMTX or the server exits", async () => {
const once = vi.spyOn(process, "once");
await video.startVideoBridge(factory, async () => "none");
const pidFile = path.join(dir, ".data", "mediamtx", "mediamtx.pid");
const settle = () => new Promise((r) => setTimeout(r, 20));
options.onSpawn!(4242);
await settle();
expect(await readFile(pidFile, "utf8")).toBe("4242\n");
options.onExit!();
await settle();
await expect(stat(pidFile)).rejects.toThrow();
options.onSpawn!(4243);
await settle();
const onExit = once.mock.calls.find(([event]) => event === "exit")![1] as () => void;
onExit();
await expect(stat(pidFile)).rejects.toThrow();
});
});
describe("waitForApi", () => {
it("gives up after its attempts", async () => {
stubApi(() => Promise.reject(new Error("ECONNREFUSED")));

View File

@ -1,5 +1,6 @@
import "server-only";
import { randomBytes } from "node:crypto";
import { rmSync } from "node:fs";
import { chmod, mkdir, rename, writeFile } from "node:fs/promises";
import path from "node:path";
import { rtspSourceWithLogin, type CameraTarget } from "./camera";
@ -11,6 +12,7 @@ import {
type ApiLogin,
type StreamKind,
} from "./mediamtx-config";
import { forgetPid, PID_FILE, recordPid, stopLeftoverMediamtx } from "./mediamtx-leftover";
import { MediamtxSupervisor, type BridgeStatus } from "./mediamtx-supervisor";
/**
@ -91,19 +93,35 @@ async function writeConfig(dir: string, text: string): Promise<string> {
export async function startVideoBridge(
makeSupervisor: (opts: ConstructorParameters<typeof MediamtxSupervisor>[0]) => MediamtxSupervisor = (o) =>
new MediamtxSupervisor(o),
stopLeftover: typeof stopLeftoverMediamtx = stopLeftoverMediamtx,
): Promise<BridgeStatus> {
if (g[KEY]) return g[KEY].supervisor.status();
const login: ApiLogin = { user: "app", pass: randomBytes(24).toString("hex") };
const dir = bridgeDir();
const configPath = await writeConfig(dir, mediamtxConfig(login));
const binary = mediamtxBinary();
const configPath = path.join(dir, "mediamtx.yml");
// A MediaMTX from a server that didn't exit cleanly would hold the ports (iss-yd2sq2q).
const leftover = await stopLeftover(dir, binary, configPath);
if (leftover === "stopped" || leftover === "killed") {
console.log(`[video] stopped a MediaMTX left running by an earlier server (${leftover})`);
}
await writeConfig(dir, mediamtxConfig(login));
const supervisor = makeSupervisor({
binary: mediamtxBinary(),
binary,
configPath,
cwd: dir,
waitReady: () => waitForApi(login),
onSpawn: (pid) => void recordPid(dir, pid).catch(() => {}),
onExit: () => void forgetPid(dir).catch(() => {}),
});
g[KEY] = { supervisor, login, paths: new Map(), pathsEpoch: 0 };
process.once("exit", () => supervisor.stop());
process.once("exit", () => {
supervisor.stop();
// Synchronous: nothing async runs once the process is exiting.
rmSync(path.join(dir, PID_FILE), { force: true });
});
supervisor.start();
return supervisor.status();
}