From 715cf5614ec0d76b631da18afcaa7e18679aea73 Mon Sep 17 00:00:00 2001 From: beasty Date: Sat, 5 Sep 2026 14:44:47 +0200 Subject: [PATCH 1/9] fix: address canvas reliability and project review findings --- .forgejo/workflows/build-images.yml | 82 +- .github/workflows/verify.yml | 35 + .gitignore | 2 + app/app/rooms/[roomId]/page.tsx | 40 +- backend/stream-canvas/README.md | 71 +- .../drizzle/meta/0000_snapshot.json | 45 +- .../stream-canvas/drizzle/meta/_journal.json | 2 +- backend/stream-canvas/package.json | 5 +- backend/stream-canvas/src/auth.ts | 17 +- backend/stream-canvas/src/db.ts | 2 + backend/stream-canvas/src/deadline.ts | 20 + backend/stream-canvas/src/index.ts | 19 +- backend/stream-canvas/src/leader.ts | 81 +- backend/stream-canvas/src/lifecycle.test.ts | 322 +++++ backend/stream-canvas/src/obs-secret.ts | 8 + backend/stream-canvas/src/rate-limit.ts | 5 +- backend/stream-canvas/src/room-operations.ts | 17 + .../src/routes.integration.test.ts | 80 +- backend/stream-canvas/src/routes.ts | 103 +- backend/stream-canvas/src/security.test.ts | 2 +- .../stream-canvas/src/snapshot-writer.test.ts | 110 ++ backend/stream-canvas/src/snapshot-writer.ts | 87 ++ backend/stream-canvas/src/types.ts | 1 + backend/stream-canvas/src/ws-handler.ts | 315 +++-- backend/stream-canvas/tsconfig.test.json | 6 + biome.json | 4 +- components/__tests__/recovery.test.tsx | 103 ++ components/__tests__/room-settings.test.tsx | 84 ++ components/stream-canvas/CanvasEditor.tsx | 145 +- components/stream-canvas/CanvasMirror.tsx | 44 +- components/stream-canvas/TwitchPreview.tsx | 170 +++ components/stream-canvas/UserMultiSelect.tsx | 37 +- .../shapes/audio/AudioPlayerShape.tsx | 83 +- components/stream-canvas/shapes/shared.ts | 44 +- components/stream-canvas/use-media-url.ts | 55 + components/user-record-bootstrap.tsx | 52 +- convex/_generated/server.d.ts | 12 +- convex/_generated/server.js | 5 +- convex/schema.ts | 1 + convex/users.integration.test.ts | 73 + convex/users.ts | 4 +- e2e/canvas.spec.ts | 266 ++++ e2e/clerk.tsx | 13 + e2e/convex.ts | 6 + e2e/index.html | 2 + e2e/main.tsx | 72 + e2e/server.ts | 146 ++ lib/landing-content.ts | 6 +- .../__tests__/media-url.test.tsx | 80 ++ .../__tests__/upload-url-refresh.test.ts | 13 +- lib/stream-canvas/api.ts | 37 +- lib/stream-canvas/room-config.ts | 28 + package.json | 14 +- playwright.config.ts | 18 + pnpm-lock.yaml | 1173 ++++++++++++----- pnpm-workspace.yaml | 2 +- renovate.json | 4 +- scripts/delivery-changes.sh | 26 + scripts/delivery-changes.test.ts | 116 ++ scripts/test-runtime-config.ts | 121 ++ scripts/verify.sh | 19 + vitest.config.ts | 8 +- 62 files changed, 3714 insertions(+), 849 deletions(-) create mode 100644 .github/workflows/verify.yml create mode 100644 backend/stream-canvas/src/deadline.ts create mode 100644 backend/stream-canvas/src/lifecycle.test.ts create mode 100644 backend/stream-canvas/src/room-operations.ts create mode 100644 backend/stream-canvas/src/snapshot-writer.test.ts create mode 100644 backend/stream-canvas/src/snapshot-writer.ts create mode 100644 backend/stream-canvas/tsconfig.test.json create mode 100644 components/__tests__/recovery.test.tsx create mode 100644 components/__tests__/room-settings.test.tsx create mode 100644 components/stream-canvas/TwitchPreview.tsx create mode 100644 components/stream-canvas/use-media-url.ts create mode 100644 convex/users.integration.test.ts create mode 100644 e2e/canvas.spec.ts create mode 100644 e2e/clerk.tsx create mode 100644 e2e/convex.ts create mode 100644 e2e/index.html create mode 100644 e2e/main.tsx create mode 100644 e2e/server.ts create mode 100644 lib/stream-canvas/__tests__/media-url.test.tsx create mode 100644 lib/stream-canvas/room-config.ts create mode 100644 playwright.config.ts create mode 100644 scripts/delivery-changes.sh create mode 100644 scripts/delivery-changes.test.ts create mode 100644 scripts/test-runtime-config.ts create mode 100644 scripts/verify.sh diff --git a/.forgejo/workflows/build-images.yml b/.forgejo/workflows/build-images.yml index 4e20c38..a831126 100644 --- a/.forgejo/workflows/build-images.yml +++ b/.forgejo/workflows/build-images.yml @@ -1,6 +1,7 @@ name: Build and Push Container Images on: + pull_request: push: branches: - main @@ -36,6 +37,35 @@ env: INFISICAL_MACHINE_IDENTITY_ID: ${{ vars.INFISICAL_MACHINE_IDENTITY_ID }} jobs: + verify: + runs-on: personal + permissions: + contents: read + services: + postgres: + image: postgres:18.4-bookworm + env: + POSTGRES_DB: canvas_test + POSTGRES_USER: canvas_test + POSTGRES_HOST_AUTH_METHOD: trust + options: >- + --health-cmd "pg_isready -U canvas_test" + --health-interval 5s --health-timeout 5s --health-retries 10 + env: + CANVAS_TEST_DATABASE_URL: postgresql://canvas_test@postgres:5432/canvas_test + steps: + - uses: actions/checkout@v7 + - uses: pnpm/action-setup@v6 + with: + run_install: false + - uses: actions/setup-node@v7 + with: + node-version-file: .node-version + cache: pnpm + - run: pnpm install --frozen-lockfile --ignore-scripts + - run: pnpm exec playwright install --with-deps chromium + - run: bash scripts/verify.sh + changes: runs-on: personal if: github.ref == 'refs/heads/main' @@ -61,25 +91,7 @@ jobs: base="$(git rev-list --max-parents=0 "$head")" fi - changed="$(git diff --name-only "$base" "$head" || git diff --name-only HEAD~1 HEAD || true)" - - has_changed() { - printf '%s\n' "$changed" | grep -Eq "$1" - } - - set_output() { - name="$1" - pattern="$2" - if has_changed "$pattern"; then - echo "$name=true" >> "$GITHUB_OUTPUT" - else - echo "$name=false" >> "$GITHUB_OUTPUT" - fi - } - - set_output frontend '^(app|components|lib|public)/|^(\.dockerignore|package\.json|pnpm-lock\.yaml|pnpm-workspace\.yaml|next\.config\.mjs|postcss\.config\.mjs|tailwind\.config\.ts|tsconfig\.json|biome\.json|components\.json|Dockerfile|docker-entrypoint\.sh)$' - set_output convex '^convex/|^convex\.json$' - set_output stream_canvas '^backend/stream-canvas/|^\.dockerignore$' + git diff --name-only "$base" "$head" | bash scripts/delivery-changes.sh >> "$GITHUB_OUTPUT" tag-convex-changes: runs-on: personal @@ -97,31 +109,15 @@ jobs: run: | set -euo pipefail - previous_tag=$( - git tag --sort=-version:refname \ - | grep -Fxv "$CURRENT_TAG" \ - | head -n 1 - ) - - if [ -z "$previous_tag" ]; then - echo "No previous release tag found; deploying Convex by default." - echo "convex=true" >> "$GITHUB_OUTPUT" - exit 0 - fi - - if git diff --quiet "$previous_tag..$CURRENT_SHA" -- convex; then - echo "convex=false" >> "$GITHUB_OUTPUT" - else - echo "convex=true" >> "$GITHUB_OUTPUT" - fi + bash scripts/delivery-changes.sh tag >> "$GITHUB_OUTPUT" deploy-convex: - needs: [changes, tag-convex-changes] + needs: [verify, changes, tag-convex-changes] if: | - always() && + always() && needs.verify.result == 'success' && github.event_name != 'pull_request' && ((github.ref == 'refs/heads/main' && needs.changes.outputs.convex == 'true') || (startsWith(github.ref, 'refs/tags/v') && needs.tag-convex-changes.outputs.convex == 'true') || - github.event.inputs.deploy_convex == 'true') + (github.ref == 'refs/heads/main' && github.event.inputs.deploy_convex == 'true')) runs-on: personal steps: - uses: actions/checkout@v7 @@ -192,9 +188,9 @@ jobs: CLERK_JWT_ISSUER_DOMAIN: ${{ steps.ci-secrets.outputs.CLERK_JWT_ISSUER_DOMAIN }} build-frontend: - needs: changes + needs: [verify, changes] if: | - always() && + always() && needs.verify.result == 'success' && github.event_name != 'pull_request' && (startsWith(github.ref, 'refs/tags/v') || needs.changes.outputs.frontend == 'true' || github.event.inputs.build_frontend == 'true') @@ -290,9 +286,9 @@ jobs: NEXT_PUBLIC_TLDRAW_LICENSE_KEY=${{ steps.ci-secrets.outputs.NEXT_PUBLIC_TLDRAW_LICENSE_KEY }} build-stream-canvas: - needs: [changes, build-frontend] + needs: [verify, changes] if: | - always() && + always() && needs.verify.result == 'success' && github.event_name != 'pull_request' && (startsWith(github.ref, 'refs/tags/v') || needs.changes.outputs.stream_canvas == 'true' || github.event.inputs.build_stream_canvas == 'true') diff --git a/.github/workflows/verify.yml b/.github/workflows/verify.yml new file mode 100644 index 0000000..9b774a3 --- /dev/null +++ b/.github/workflows/verify.yml @@ -0,0 +1,35 @@ +name: Verify +on: + pull_request: + push: + branches: [main] +permissions: + contents: read +jobs: + verify: + runs-on: ubuntu-latest + services: + postgres: + image: postgres:18.4-bookworm + env: + POSTGRES_DB: canvas_test + POSTGRES_USER: canvas_test + POSTGRES_HOST_AUTH_METHOD: trust + ports: ["5432:5432"] + options: >- + --health-cmd "pg_isready -U canvas_test" + --health-interval 5s --health-timeout 5s --health-retries 10 + env: + CANVAS_TEST_DATABASE_URL: postgresql://canvas_test@127.0.0.1:5432/canvas_test + steps: + - uses: actions/checkout@v7 + - uses: pnpm/action-setup@v6 + with: + run_install: false + - uses: actions/setup-node@v7 + with: + node-version-file: .node-version + cache: pnpm + - run: pnpm install --frozen-lockfile --ignore-scripts + - run: pnpm exec playwright install --with-deps chromium + - run: bash scripts/verify.sh diff --git a/.gitignore b/.gitignore index 7738271..e71d9a9 100644 --- a/.gitignore +++ b/.gitignore @@ -14,6 +14,8 @@ backend/stream-canvas/node_modules backend/stream-canvas/data backend/stream-canvas/bun.lockb coverage +playwright-report +test-results dist build out diff --git a/app/app/rooms/[roomId]/page.tsx b/app/app/rooms/[roomId]/page.tsx index 8388ce8..962d2f5 100644 --- a/app/app/rooms/[roomId]/page.tsx +++ b/app/app/rooms/[roomId]/page.tsx @@ -13,31 +13,67 @@ export default function StreamCanvasRoomPage() { const [twitchChannel, setTwitchChannel] = useState(null); const [youtubePolicy, setYouTubePolicy] = useState("preview_only"); + const [loadedRoomId, setLoadedRoomId] = useState(null); + const [error, setError] = useState(null); + const [attempt, setAttempt] = useState(0); // Fetch this room's Twitch channel from the accessible rooms list + // biome-ignore lint/correctness/useExhaustiveDependencies: attempt explicitly retries the same request useEffect(() => { if (!isLoaded || !isSignedIn) return; let cancelled = false; + setError(null); getAccessibleRooms(getToken) .then((rooms) => { const room = rooms.find((r) => r.id === roomId); + if (!room) + throw new Error( + "This room is unavailable or you no longer have access.", + ); if (!cancelled && room) { setTwitchChannel(room.twitchChannel); setYouTubePolicy(room.youtubePolicy); + setLoadedRoomId(roomId); } }) .catch((err) => { if (!cancelled) - console.error("[stream-canvas] Failed to load room info:", err); + setError(err instanceof Error ? err.message : "Failed to load room"); }); return () => { cancelled = true; }; - }, [getToken, isLoaded, isSignedIn, roomId]); + }, [getToken, isLoaded, isSignedIn, roomId, attempt]); + + if (isLoaded && !isSignedIn) + return

Sign in to open this room.

; + if (error) + return ( +
+

{error}

+ +
+ ); + if (loadedRoomId !== roomId) + return ( +

+ Loading room… +

+ ); return (
{ max: config.databasePoolSize, connectionTimeoutMillis: 5_000, idleTimeoutMillis: 30_000, + statement_timeout: 10_000, + query_timeout: 12_000, application_name: "moddrop-stream-canvas", }); } diff --git a/backend/stream-canvas/src/deadline.ts b/backend/stream-canvas/src/deadline.ts new file mode 100644 index 0000000..a7b71ac --- /dev/null +++ b/backend/stream-canvas/src/deadline.ts @@ -0,0 +1,20 @@ +export async function withDeadline( + operation: Promise, + milliseconds: number, + label: string, +): Promise { + let timer: NodeJS.Timeout | undefined; + try { + return await Promise.race([ + operation, + new Promise((_resolve, reject) => { + timer = setTimeout( + () => reject(new Error(`${label} timed out`)), + milliseconds, + ); + }), + ]); + } finally { + clearTimeout(timer); + } +} diff --git a/backend/stream-canvas/src/index.ts b/backend/stream-canvas/src/index.ts index 91c7a49..9e2e3bf 100644 --- a/backend/stream-canvas/src/index.ts +++ b/backend/stream-canvas/src/index.ts @@ -4,6 +4,7 @@ import { cors } from "hono/cors"; import { sql } from "drizzle-orm"; import { WebSocketServer } from "ws"; import { config, validateServerConfig } from "./config.ts"; +import { withDeadline } from "./deadline.ts"; import { closeDatabase, db } from "./db.ts"; import { leaderState } from "./leader.ts"; import { objectStore } from "./object-store.ts"; @@ -25,7 +26,13 @@ app.get("/ready", async (c) => { return c.json({ status: "standby", role: "standby" }, 503); } try { - await Promise.all([db.execute(sql`SELECT 1`), objectStore.check()]); + await withDeadline( + Promise.all([db.execute(sql`SELECT 1`), objectStore.check()]), + 3_000, + "Readiness", + ); + if (!leaderState.isLeader) + return c.json({ status: "standby", role: "standby" }, 503); return c.json({ status: "ready", role: "leader" }); } catch (error) { console.error("[readiness] dependency check failed", error); @@ -53,9 +60,9 @@ const wss = new WebSocketServer({ perMessageDeflate: false, }); -leaderState.onDemote(async () => { +leaderState.onDemote(async (persist) => { for (const client of wss.clients) client.close(1012, "Service restarting"); - await closeAllRooms(); + await closeAllRooms(persist); }); leaderState.start(); @@ -111,7 +118,11 @@ async function performShutdown(): Promise { } function requestShutdown(): void { - void shutdown().then( + void withDeadline( + shutdown(), + 25_000, + "Shutdown; pending data may not be durable", + ).then( () => process.exit(0), (error) => { console.error("[stream-canvas] shutdown failed", error); diff --git a/backend/stream-canvas/src/leader.ts b/backend/stream-canvas/src/leader.ts index 0dda070..39448c0 100644 --- a/backend/stream-canvas/src/leader.ts +++ b/backend/stream-canvas/src/leader.ts @@ -1,5 +1,8 @@ import type { PoolClient } from "pg"; +import type { RoomSnapshot } from "@tldraw/sync-core"; +import { drizzle } from "drizzle-orm/node-postgres"; import { pool } from "./db.ts"; +import { canvasDocuments } from "./schema.ts"; const LEADER_LOCK_KEY = 1_296_315_460; @@ -8,7 +11,19 @@ class LeaderState { private timer: NodeJS.Timeout | null = null; private tickPromise: Promise | null = null; private stopping = false; - private demoteHandlers = new Set<() => void | Promise>(); + private demoteHandlers = new Set< + (persist: boolean) => void | Promise + >(); + private demoting: Promise | null = null; + private connectionFailed = false; + + private readonly onClientError = (error: Error) => { + console.error("[leader] retained PostgreSQL connection failed", error); + this.connectionFailed = true; + void this.demote(false).catch((failure) => + console.error("[leader] demotion failed", failure), + ); + }; isLeader = false; @@ -17,7 +32,7 @@ class LeaderState { this.runTick(); } - onDemote(handler: () => void | Promise): () => void { + onDemote(handler: (persist: boolean) => void | Promise): () => void { this.demoteHandlers.add(handler); return () => this.demoteHandlers.delete(handler); } @@ -27,7 +42,31 @@ class LeaderState { if (this.timer) clearTimeout(this.timer); this.timer = null; await this.tickPromise; - await this.demote(); + await this.demote(true); + } + + async persistSnapshot(roomId: string, snapshot: RoomSnapshot): Promise { + const client = this.client; + if (!client || this.connectionFailed) + throw new Error("Canvas ownership was lost"); + // Writes use the lock-owning connection. A successor cannot acquire the + // advisory lock until this connection and its queued writes have ended. + await drizzle(client) + .insert(canvasDocuments) + .values({ + roomId, + snapshot, + revision: snapshot.documentClock, + updatedAt: new Date(), + }) + .onConflictDoUpdate({ + target: canvasDocuments.roomId, + set: { + snapshot, + revision: snapshot.documentClock, + updatedAt: new Date(), + }, + }); } private runTick(): void { @@ -44,6 +83,7 @@ class LeaderState { private async tick(): Promise { try { + await this.demoting; if (this.isLeader && this.client) { await this.client.query("SELECT 1"); } else { @@ -51,7 +91,8 @@ class LeaderState { } } catch (error) { console.error("[leader] leadership connection failed", error); - await this.demote(); + this.connectionFailed = true; + await this.demote(false); } finally { if (!this.stopping) { this.timer = setTimeout( @@ -67,6 +108,8 @@ class LeaderState { private async tryAcquire(): Promise { const client = await pool.connect(); + this.connectionFailed = false; + client.on("error", this.onClientError); let retained = false; let destroy = false; try { @@ -75,7 +118,7 @@ class LeaderState { [LEADER_LOCK_KEY], ); if (!result.rows[0]?.acquired) return; - if (this.stopping) { + if (this.stopping || this.connectionFailed) { await client.query("SELECT pg_advisory_unlock($1)", [LEADER_LOCK_KEY]); return; } @@ -87,22 +130,34 @@ class LeaderState { destroy = true; throw error; } finally { - if (!retained) client.release(destroy); + if (!retained) { + client.release(destroy || this.connectionFailed); + client.removeListener("error", this.onClientError); + } } } - private async demote(): Promise { - const wasLeader = this.isLeader; + private demote(persist: boolean): Promise { this.isLeader = false; + if (this.demoting) return this.demoting; + this.demoting = this.performDemotion(persist).finally(() => { + this.demoting = null; + }); + return this.demoting; + } + + private async performDemotion(persist: boolean): Promise { + const failures: unknown[] = []; const client = this.client; - this.client = null; - if (wasLeader) { + if (!persist) this.client = null; + if (client) { console.warn("[leader] relinquishing stream-canvas leadership"); const results = await Promise.allSettled( - [...this.demoteHandlers].map((handler) => handler()), + [...this.demoteHandlers].map((handler) => handler(persist)), ); for (const result of results) { if (result.status === "rejected") { + failures.push(result.reason); console.error("[leader] demotion handler failed", result.reason); } } @@ -113,9 +168,13 @@ class LeaderState { } catch { // The connection may already be gone; PostgreSQL releases session locks. } finally { + this.client = null; client.release(true); + client.removeListener("error", this.onClientError); } } + if (failures.length) + throw new AggregateError(failures, "Leadership demotion failed"); } } diff --git a/backend/stream-canvas/src/lifecycle.test.ts b/backend/stream-canvas/src/lifecycle.test.ts new file mode 100644 index 0000000..7ba043c --- /dev/null +++ b/backend/stream-canvas/src/lifecycle.test.ts @@ -0,0 +1,322 @@ +import assert from "node:assert/strict"; +import { generateKeyPairSync, randomBytes, sign } from "node:crypto"; +import { once } from "node:events"; +import { mkdtemp, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import test from "node:test"; +import { setTimeout as delay } from "node:timers/promises"; + +test("PostgreSQL and WebSockets preserve canvas state, revoke sessions, and handle ownership loss", { + skip: + !process.env.CANVAS_TEST_DATABASE_URL && + "Set CANVAS_TEST_DATABASE_URL to an isolated PostgreSQL database", + timeout: 60_000, +}, async (t) => { + const directory = await mkdtemp(join(tmpdir(), "canvas-lifecycle-")); + const { publicKey, privateKey } = generateKeyPairSync("rsa", { + modulusLength: 2048, + }); + const origin = "http://localhost:4310"; + Object.assign(process.env, { + NODE_ENV: "test", + DATABASE_URL: process.env.CANVAS_TEST_DATABASE_URL, + UPLOADS_DIR: directory, + OBJECT_STORAGE_MODE: "filesystem", + CORS_ORIGINS: origin, + CLERK_JWT_KEY: publicKey.export({ type: "spki", format: "pem" }), + CLERK_JWT_ISSUER_DOMAIN: "https://identity.example", + OBS_TOKEN_SIGNING_SECRET: randomBytes(32).toString("hex"), + }); + const [ + { db, pool }, + { leaderState }, + { api }, + wsHandler, + { migrate }, + { serve }, + { WebSocket, WebSocketServer }, + sync, + tl, + ] = await Promise.all([ + import("./db.ts"), + import("./leader.ts"), + import("./routes.ts"), + import("./ws-handler.ts"), + import("drizzle-orm/node-postgres/migrator"), + import("@hono/node-server"), + import("ws"), + import("@tldraw/sync-core"), + import("tldraw"), + ]); + const { streamCanvasSchema } = await import("./tldraw-schema.ts"); + await migrate(db, { + migrationsFolder: new URL("../drizzle", import.meta.url).pathname, + }); + const server = serve({ fetch: api.fetch, port: 0, hostname: "127.0.0.1" }); + const wss = new WebSocketServer({ noServer: true }); + server.on("upgrade", (req, socket, head) => + wss.handleUpgrade(req, socket, head, (ws) => { + void wsHandler + .handleWebSocketUpgrade(ws, req) + .catch(() => ws.close(1011)); + }), + ); + if (!server.listening) await once(server, "listening"); + const address = server.address(); + assert.ok(address && typeof address !== "string"); + const base = `http://127.0.0.1:${address.port}`; + const sockets: InstanceType[] = []; + const createdRooms: string[] = []; + let demotions = 0; + const unsubscribe = leaderState.onDemote(async (persist) => { + demotions++; + await wsHandler.closeAllRooms(persist); + }); + t.after(async () => { + t.mock.restoreAll(); + for (const socket of sockets) socket.terminate(); + await leaderState.stop(); + unsubscribe(); + await new Promise((resolve) => wss.close(() => resolve())); + await new Promise((resolve) => server.close(() => resolve())); + for (const id of createdRooms) + await pool.query("DELETE FROM rooms WHERE id = $1", [id]); + await pool.end(); + await rm(directory, { recursive: true, force: true }); + }); + const owner = `user_${crypto.randomUUID().replaceAll("-", "")}`; + const member = `user_${crypto.randomUUID().replaceAll("-", "")}`; + function headers(user = owner) { + const now = Math.floor(Date.now() / 1000); + const head = Buffer.from( + JSON.stringify({ alg: "RS256", typ: "JWT", kid: "local-test" }), + ).toString("base64url"); + const body = Buffer.from( + JSON.stringify({ + sub: user, + iss: "https://identity.example", + azp: origin, + iat: now, + exp: now + 300, + }), + ).toString("base64url"); + const signature = sign( + "RSA-SHA256", + Buffer.from(`${head}.${body}`), + privateKey, + ).toString("base64url"); + return { + Origin: origin, + Authorization: `Bearer ${head}.${body}.${signature}`, + "Content-Type": "application/json", + }; + } + async function request( + path: string, + method = "GET", + body?: object, + user = owner, + ) { + const response = await fetch(base + path, { + method, + headers: headers(user), + body: body ? JSON.stringify(body) : undefined, + }); + assert.ok(response.ok, `${method} ${path}: ${response.status}`); + return response; + } + const created = await Promise.all( + Array.from({ length: 4 }, () => request("/api/rooms", "POST")), + ); + const rooms = await Promise.all( + created.map( + (response) => + response.json() as Promise<{ id: string; obsSetupSecret?: string }>, + ), + ); + const roomId = rooms[0]?.id; + assert.ok(roomId); + createdRooms.push(roomId); + assert.equal(new Set(rooms.map((room) => room.id)).size, 1); + assert.equal(created.filter((response) => response.status === 201).length, 1); + assert.equal(rooms.filter((room) => room.obsSetupSecret).length, 1); + const secret = rooms.find((room) => room.obsSetupSecret)?.obsSetupSecret; + assert.ok(secret); + await request(`/api/rooms/${roomId}`, "PATCH", { allowedUsers: [member] }); + leaderState.start(); + await until(() => leaderState.isLeader); + async function connect(token: string) { + const ws = new WebSocket( + `${base.replace("http:", "ws:")}/ws?roomId=${roomId}&token=${encodeURIComponent(token)}`, + { origin }, + ); + const messages: unknown[] = []; + let closed = false; + sockets.push(ws); + ws.on("close", () => { + closed = true; + }); + ws.on("message", (data) => { + messages.push(JSON.parse(data.toString())); + }); + await once(ws, "open"); + // Admission may perform database reads; ws buffers the protocol handshake. + ws.send( + JSON.stringify({ + type: "connect", + connectRequestId: crypto.randomUUID(), + lastServerClock: 0, + // The runtime exports this protocol helper but omits it from public declarations. + protocolVersion: ( + sync as typeof sync & { getTlsyncProtocolVersion(): number } + ).getTlsyncProtocolVersion(), + schema: streamCanvasSchema.serialize(), + }), + ); + await until( + () => messages.some((message) => isMessage(message, "connect")) || closed, + ); + return { + ws, + messages, + get closed() { + return closed; + }, + }; + } + async function editorToken(user = owner) { + return ( + (await ( + await request(`/api/rooms/${roomId}/ws-token`, "POST", undefined, user) + ).json()) as { token: string } + ).token; + } + const obsToken = ( + (await (await request("/obs/token", "POST", { secret })).json()) as { + token: string; + } + ).token; + const editor = await connect(await editorToken()); + const collaborator = await connect(await editorToken(member)); + const mirror = await connect(obsToken); + assert.ok( + mirror.messages.some( + (message) => isMessage(message, "connect") && message.isReadonly === true, + ), + ); + await request(`/api/rooms/${roomId}`, "PATCH", { + youtubePolicy: "disabled", + allowedUsers: [], + }); + await until(() => collaborator.closed); + assert.equal(editor.closed, false); + assert.equal(mirror.closed, false); + await until(() => + mirror.messages.some( + (message) => + isMessage(message, "custom") && + JSON.stringify(message).includes('"youtubePolicy":"disabled"'), + ), + ); + + let releaseWrite: () => void = () => {}; + let writeStarted = false; + const held = new Promise((resolve) => { + releaseWrite = resolve; + }); + const persist = leaderState.persistSnapshot.bind(leaderState); + t.mock.method( + leaderState, + "persistSnapshot", + async (...args: Parameters) => { + writeStarted = true; + await held; + return persist(...args); + }, + ); + const page = tl.PageRecordType.create({ + name: "Latest durable canvas", + index: "a2" as import("tldraw").IndexKey, + }); + editor.ws.send( + JSON.stringify({ + type: "push", + clientClock: 1, + diff: { [page.id]: ["put", page] }, + }), + ); + await until(() => writeStarted); + await until(() => + mirror.messages.some((message) => + JSON.stringify(message).includes("Latest durable canvas"), + ), + ); + await request(`/api/rooms/${roomId}/regenerate-secret`, "POST"); + await until(() => mirror.closed); + assert.equal(editor.closed, false); + const obsolete = await connect(obsToken); + assert.equal(obsolete.closed, true); + editor.ws.close(); + await until(() => editor.closed); + // tldraw removes a disconnected session after its reconnect grace period. + await delay(6_000); + let disposalFinished = false; + const disposal = wsHandler.disposeInactiveRoom(roomId).then(() => { + disposalFinished = true; + }); + await delay(20); + assert.equal(disposalFinished, false); + const reopening = connect(await editorToken()); + await delay(20); + releaseWrite(); + await disposal; + const reopened = await reopening; + assert.ok( + reopened.messages.some((message) => + JSON.stringify(message).includes("Latest durable canvas"), + ), + ); + t.mock.restoreAll(); + const persisted = await pool.query<{ snapshot: unknown }>( + "SELECT snapshot FROM canvas_documents WHERE room_id = $1", + [roomId], + ); + assert.match( + JSON.stringify(persisted.rows[0]?.snapshot), + /Latest durable canvas/, + ); + const lock = await pool.query<{ pid: number }>( + "SELECT pid FROM pg_locks WHERE locktype = 'advisory' AND objid = 1296315460 AND database = (SELECT oid FROM pg_database WHERE datname = current_database()) AND granted", + ); + assert.equal(lock.rows.length, 1); + await pool.query("SELECT pg_terminate_backend($1)", [lock.rows[0]?.pid]); + await until(() => demotions === 1 && reopened.closed); + await until(() => leaderState.isLeader); + assert.equal(demotions, 1); + const recovered = await connect(await editorToken()); + assert.ok( + recovered.messages.some((message) => + JSON.stringify(message).includes("Latest durable canvas"), + ), + ); +}); + +function isMessage( + value: unknown, + type: string, +): value is Record { + return ( + value !== null && + typeof value === "object" && + "type" in value && + value.type === type + ); +} +async function until(predicate: () => boolean) { + const end = Date.now() + 10_000; + while (!predicate()) { + if (Date.now() >= end) throw new Error("Lifecycle condition timed out"); + await delay(10); + } +} diff --git a/backend/stream-canvas/src/obs-secret.ts b/backend/stream-canvas/src/obs-secret.ts index ef8d1fc..04d5175 100644 --- a/backend/stream-canvas/src/obs-secret.ts +++ b/backend/stream-canvas/src/obs-secret.ts @@ -19,6 +19,14 @@ export function isHashedObsSecret(value: string): boolean { return value.startsWith(OBS_SECRET_PREFIX); } +/** Public ticket version, separate from the stored bootstrap credential. */ +export function obsCredentialVersion(storedSecret: string): string { + return createHash("sha256") + .update("obs-ticket-version\0") + .update(storedSecret) + .digest("base64url"); +} + export function verifyObsSecret(secret: string, storedSecret: string): boolean { const expected = isHashedObsSecret(storedSecret) ? storedSecret diff --git a/backend/stream-canvas/src/rate-limit.ts b/backend/stream-canvas/src/rate-limit.ts index 9539f15..f0a8003 100644 --- a/backend/stream-canvas/src/rate-limit.ts +++ b/backend/stream-canvas/src/rate-limit.ts @@ -63,7 +63,10 @@ export class FixedWindowRateLimit { } private cleanup(now: number): void { - if (this.entries.size < 1000 && this.entries.size <= this.options.maxEntries) { + if ( + this.entries.size < 1000 && + this.entries.size <= this.options.maxEntries + ) { return; } diff --git a/backend/stream-canvas/src/room-operations.ts b/backend/stream-canvas/src/room-operations.ts new file mode 100644 index 0000000..f672838 --- /dev/null +++ b/backend/stream-canvas/src/room-operations.ts @@ -0,0 +1,17 @@ +// Configuration commits and session admission share this queue. A connection +// cannot authenticate against a membership set that is being replaced. +const operations = new Map>(); + +export async function withRoomOperation( + roomId: string, + run: () => Promise, +): Promise { + const previous = operations.get(roomId); + const current = (previous ?? Promise.resolve()).catch(() => {}).then(run); + operations.set(roomId, current); + try { + return await current; + } finally { + if (operations.get(roomId) === current) operations.delete(roomId); + } +} diff --git a/backend/stream-canvas/src/routes.integration.test.ts b/backend/stream-canvas/src/routes.integration.test.ts index 03179e3..b9ab359 100644 --- a/backend/stream-canvas/src/routes.integration.test.ts +++ b/backend/stream-canvas/src/routes.integration.test.ts @@ -1,12 +1,14 @@ import assert from "node:assert/strict"; -import { mkdtemp } from "node:fs/promises"; +import { mkdtemp, rm } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; -import test from "node:test"; -import { generateKeyPairSync, sign } from "node:crypto"; +import test, { after } from "node:test"; +import { generateKeyPairSync, randomBytes, sign } from "node:crypto"; import type { IncomingMessage } from "node:http"; import { eq } from "drizzle-orm"; import { v4 as uuidv4 } from "uuid"; +import { Hono } from "hono"; +import { cors } from "hono/cors"; const origin = "https://moddrop.localhost:1355"; const issuer = "https://example.clerk.accounts.dev"; @@ -53,10 +55,13 @@ process.env.CLERK_JWT_KEY = publicKey.export({ }); process.env.CLERK_JWT_ISSUER_DOMAIN = issuer; process.env.CORS_ORIGINS = origin; -process.env.OBS_TOKEN_SIGNING_SECRET = - "test-signing-secret-with-at-least-32-bytes"; +process.env.OBS_TOKEN_SIGNING_SECRET = randomBytes(32).toString("hex"); const { db, pool } = await import("./db.ts"); +after(async () => { + await pool.end(); + await rm(tempRoot, { recursive: true, force: true }); +}); for (const statement of [ "CREATE TYPE youtube_policy AS ENUM ('disabled', 'preview_only', 'allow_on_air')", @@ -138,6 +143,71 @@ await db.insert(roomMembers).values({ clerkUserId: collaboratorUserId, }); +test("CORS permits the configured origin and bounded authorization preflight", async () => { + const app = new Hono().use("*", cors({ origin: [origin] })).route("/", api); + const response = await app.request("/api/rooms", { + method: "OPTIONS", + headers: { + Origin: origin, + "Access-Control-Request-Method": "POST", + "Access-Control-Request-Headers": "authorization,content-type", + }, + }); + assert.equal(response.status, 204); + assert.equal(response.headers.get("access-control-allow-origin"), origin); + assert.match( + response.headers.get("access-control-allow-headers") ?? "", + /authorization/i, + ); + const other = await app.request("/api/rooms", { + method: "OPTIONS", + headers: { + Origin: "https://other.example", + "Access-Control-Request-Method": "POST", + }, + }); + assert.equal(other.headers.get("access-control-allow-origin"), null); +}); + +test("object-write and metadata failures leave no successful upload record", async (t) => { + const { objectStore } = await import("./object-store.ts"); + const before = await db.select().from(uploads); + const upload = () => { + const body = new FormData(); + body.set( + "file", + new File([Uint8Array.from(minimalPng())], "fixture.png", { + type: "image/png", + }), + ); + return api.request(`/api/rooms/${roomId}/upload`, { + method: "POST", + headers: authHeaders(ownerUserId), + body, + }); + }; + t.mock.method(objectStore, "put", async () => { + throw new Error("test object store unavailable"); + }); + assert.equal((await upload()).status, 500); + assert.equal((await db.select().from(uploads)).length, before.length); + t.mock.restoreAll(); + const originalRemove = objectStore.remove.bind(objectStore); + const removed: string[] = []; + t.mock.method(objectStore, "remove", async (key: string) => { + removed.push(key); + await originalRemove(key); + }); + t.mock.method(db, "insert", () => { + throw new Error("test metadata insertion failed"); + }); + assert.equal((await upload()).status, 500); + assert.equal(removed.length, 1); + assert.equal(await objectStore.stat(removed[0] ?? ""), null); + assert.equal((await db.select().from(uploads)).length, before.length); + t.mock.restoreAll(); +}); + test("OBS token exchange accepts hashed room secrets", async () => { const requestBody = JSON.stringify({ secret: obsSecret }); const response = await api.request("/obs/token", { diff --git a/backend/stream-canvas/src/routes.ts b/backend/stream-canvas/src/routes.ts index 226fbdd..1716052 100644 --- a/backend/stream-canvas/src/routes.ts +++ b/backend/stream-canvas/src/routes.ts @@ -16,7 +16,13 @@ import { import { config } from "./config.ts"; import { db } from "./db.ts"; import { parseSingleByteRange } from "./http-range.ts"; -import { generateObsSecret, hashObsSecret } from "./obs-secret.ts"; +import { + generateObsSecret, + hashObsSecret, + obsCredentialVersion, +} from "./obs-secret.ts"; +import { withRoomOperation } from "./room-operations.ts"; +import { roomConfigChanged, revokeObsSessions } from "./ws-handler.ts"; import { objectStore } from "./object-store.ts"; import { FixedWindowRateLimit, rateLimitKeyFromHeaders } from "./rate-limit.ts"; import { @@ -180,7 +186,7 @@ api.post("/obs/token", async (c) => { return c.json({ error: "Invalid secret" }, 401); } - const token = mintObsToken(room.id); + const token = mintObsToken(room.id, obsCredentialVersion(room.obsSecret)); return c.json({ token, roomId: room.id, @@ -258,9 +264,16 @@ authed.post("/rooms", async (c) => { createdAt: now, updatedAt: now, }) + .onConflictDoNothing({ target: rooms.ownerClerkId }) .returning(); - if (!room) return c.json({ error: "Room creation failed" }, 500); + if (!room) { + const winner = await db.query.rooms.findFirst({ + where: eq(rooms.ownerClerkId, user.sub), + }); + if (!winner) return c.json({ error: "Room creation failed" }, 500); + return c.json(await ownerRoomResponse(winner)); + } return c.json(await ownerRoomResponse(room, obsSecret), 201); }); @@ -340,45 +353,44 @@ authed.patch("/rooms/:id", async (c) => { return c.json({ error: validation.error }, 400); } - await db.transaction(async (tx) => { - await tx - .update(rooms) - .set({ - ...(validation.value.twitchChannel !== undefined && { - twitchChannel: validation.value.twitchChannel, - }), - ...(validation.value.youtubePolicy !== undefined && { - youtubePolicy: validation.value.youtubePolicy, - youtubeRiskAcknowledgedAt: - validation.value.youtubePolicy === "allow_on_air" - ? new Date() - : null, - }), - updatedAt: new Date(), - }) - .where(eq(rooms.id, room.id)); - - if (validation.value.allowedUsers !== undefined) { - await tx.delete(roomMembers).where(eq(roomMembers.roomId, room.id)); - if (validation.value.allowedUsers.length > 0) { - await tx.insert(roomMembers).values( - validation.value.allowedUsers.map((clerkUserId) => ({ - roomId: room.id, - clerkUserId, - })), - ); + return withRoomOperation(room.id, async () => { + const updated = await db.transaction(async (tx) => { + const [updatedRoom] = await tx + .update(rooms) + .set({ + ...(validation.value.twitchChannel !== undefined && { + twitchChannel: validation.value.twitchChannel, + }), + ...(validation.value.youtubePolicy !== undefined && { + youtubePolicy: validation.value.youtubePolicy, + youtubeRiskAcknowledgedAt: + validation.value.youtubePolicy === "allow_on_air" + ? new Date() + : null, + }), + updatedAt: new Date(), + }) + .where(eq(rooms.id, room.id)) + .returning(); + if (!updatedRoom) throw new Error("Room disappeared during update"); + + if (validation.value.allowedUsers !== undefined) { + await tx.delete(roomMembers).where(eq(roomMembers.roomId, room.id)); + if (validation.value.allowedUsers.length > 0) { + await tx.insert(roomMembers).values( + validation.value.allowedUsers.map((clerkUserId) => ({ + roomId: room.id, + clerkUserId, + })), + ); + } } - } - }); + return updatedRoom; + }); - const updated = await db.query.rooms.findFirst({ - where: eq(rooms.id, room.id), + roomConfigChanged(updated, validation.value.allowedUsers); + return c.json(await ownerRoomResponse(updated)); }); - if (!updated) { - return c.json({ error: "Updated room could not be retrieved" }, 500); - } - - return c.json(await ownerRoomResponse(updated)); }); // Mint a short-lived editor WebSocket ticket. Raw Clerk JWTs never go in WS URLs. @@ -426,12 +438,15 @@ authed.post("/rooms/:id/regenerate-secret", async (c) => { if (!room) return c.json({ error: "Not found" }, 404); const newSecret = generateObsSecret(); - await db - .update(rooms) - .set({ obsSecret: hashObsSecret(newSecret), updatedAt: new Date() }) - .where(eq(rooms.id, room.id)); + return withRoomOperation(room.id, async () => { + await db + .update(rooms) + .set({ obsSecret: hashObsSecret(newSecret), updatedAt: new Date() }) + .where(eq(rooms.id, room.id)); - return c.json({ ok: true, obsSecret: newSecret }); + revokeObsSessions(room.id); + return c.json({ ok: true, obsSecret: newSecret }); + }); }); // OBS secrets are hashed at rest and cannot be revealed after creation. diff --git a/backend/stream-canvas/src/security.test.ts b/backend/stream-canvas/src/security.test.ts index 9c7043c..eefdee8 100644 --- a/backend/stream-canvas/src/security.test.ts +++ b/backend/stream-canvas/src/security.test.ts @@ -38,7 +38,7 @@ test("OBS secrets are hashed and verified without storing plaintext", () => { }); test("short-lived OBS and editor WebSocket tokens validate role, scope, and expiry", () => { - const obsToken = mintObsToken(roomId); + const obsToken = mintObsToken(roomId, "test-credential-version"); const obsClaims = verifyObsToken(obsToken); assert.equal(obsClaims?.roomId, roomId); assert.equal(obsClaims?.role, "obs"); diff --git a/backend/stream-canvas/src/snapshot-writer.test.ts b/backend/stream-canvas/src/snapshot-writer.test.ts new file mode 100644 index 0000000..fa14c20 --- /dev/null +++ b/backend/stream-canvas/src/snapshot-writer.test.ts @@ -0,0 +1,110 @@ +import assert from "node:assert/strict"; +import { setTimeout as delay } from "node:timers/promises"; +import test from "node:test"; +import type { RoomSnapshot } from "@tldraw/sync-core"; +import { RoomSnapshotWriter } from "./snapshot-writer.ts"; +import { withRoomOperation } from "./room-operations.ts"; + +const snapshot = (clock: number): RoomSnapshot => ({ + documentClock: clock, + documents: [], + tombstones: {}, + tombstoneHistoryStartsAtClock: 0, +}); + +test("snapshot creation is deferred and a burst writes only its latest state", async () => { + let clock = 0; + let snapshots = 0; + const written: number[] = []; + const writer = new RoomSnapshotWriter( + () => { + snapshots++; + return snapshot(clock); + }, + async (value) => { + written.push(value.documentClock ?? 0); + }, + "test", + ); + for (clock = 1; clock <= 100; clock++) writer.schedule(); + assert.equal(snapshots, 0); + await writer.flush(); + assert.deepEqual(written, [101]); + assert.equal(snapshots, 1); + writer.stop(); +}); + +test("flush drains changes arriving behind a pending write, without overlapping writers", async () => { + let clock = 1; + let release: () => void = () => {}; + const pending = new Promise((resolve) => { + release = resolve; + }); + const written: number[] = []; + const writer = new RoomSnapshotWriter( + () => snapshot(clock), + async (value) => { + if (!written.length) await pending; + written.push(value.documentClock ?? 0); + }, + "test", + ); + writer.schedule(); + const flush = writer.flush(); + clock = 2; + writer.schedule(); + release(); + await flush; + assert.deepEqual(written, [1, 2]); + writer.stop(); +}); + +test("failed flush is reported, keeps the latest state for retry, and respects a deadline", async () => { + let fail = true; + const written: number[] = []; + const writer = new RoomSnapshotWriter( + () => snapshot(3), + async (value) => { + if (fail) throw new Error("storage unavailable"); + written.push(value.documentClock ?? 0); + }, + "test", + ); + writer.schedule(); + await assert.rejects(writer.flush(), /storage unavailable/); + fail = false; + await writer.flush(); + assert.deepEqual(written, [3]); + writer.stop(); + const stalled = new RoomSnapshotWriter( + () => snapshot(1), + () => new Promise(() => {}), + "stalled", + ); + stalled.schedule(); + await assert.rejects(stalled.flush(10), /timed out/); + stalled.stop(); +}); + +test("room operations serialize a commit and admission, and recover from a failed operation", async () => { + const order: string[] = []; + await Promise.all([ + withRoomOperation("one", async () => { + await delay(5); + order.push("commit"); + }), + withRoomOperation("one", async () => { + order.push("admit"); + }), + ]); + assert.deepEqual(order, ["commit", "admit"]); + await assert.rejects( + withRoomOperation("one", async () => { + throw new Error("failed"); + }), + ); + assert.equal( + await withRoomOperation("one", async () => "recovered"), + "recovered", + ); +}); diff --git a/backend/stream-canvas/src/snapshot-writer.ts b/backend/stream-canvas/src/snapshot-writer.ts new file mode 100644 index 0000000..658e4ed --- /dev/null +++ b/backend/stream-canvas/src/snapshot-writer.ts @@ -0,0 +1,87 @@ +import type { RoomSnapshot } from "@tldraw/sync-core"; +import { withDeadline } from "./deadline.ts"; + +/** Coalesce changes before constructing the full snapshot; flush bypasses the delay. */ +export class RoomSnapshotWriter { + private dirty = false; + private stopped = false; + private timer: NodeJS.Timeout | undefined; + private draining: Promise | null = null; + private readonly snapshot: () => RoomSnapshot; + private readonly persist: (snapshot: RoomSnapshot) => Promise; + private readonly label: string; + private readonly debounceMs: number; + + constructor( + snapshot: () => RoomSnapshot, + persist: (snapshot: RoomSnapshot) => Promise, + label: string, + debounceMs = 100, + ) { + this.snapshot = snapshot; + this.persist = persist; + this.label = label; + this.debounceMs = debounceMs; + } + + schedule(): void { + if (this.stopped) return; + this.dirty = true; + this.arm(this.debounceMs); + } + + private arm(delay: number): void { + if (this.timer || this.draining || this.stopped) return; + this.timer = setTimeout(() => { + this.timer = undefined; + void this.drain().catch((error) => { + console.error( + `[canvas] ${this.label} persistence failed; retrying`, + error, + ); + this.arm(1_000); + }); + }, delay); + this.timer.unref(); + } + + private drain(): Promise { + if (this.draining) return this.draining; + this.draining = (async () => { + while (this.dirty && !this.stopped) { + this.dirty = false; + try { + await this.persist(this.snapshot()); + } catch (error) { + this.dirty = true; + throw error; + } + } + })().finally(() => { + this.draining = null; + if (this.dirty) this.arm(1_000); + }); + return this.draining; + } + + async flush(timeoutMs = 15_000): Promise { + clearTimeout(this.timer); + this.timer = undefined; + try { + await withDeadline( + this.drain(), + timeoutMs, + `${this.label} snapshot flush`, + ); + } catch (error) { + this.arm(1_000); + throw error; + } + } + + stop(): void { + this.stopped = true; + clearTimeout(this.timer); + this.timer = undefined; + } +} diff --git a/backend/stream-canvas/src/types.ts b/backend/stream-canvas/src/types.ts index f2f2426..9e89b50 100644 --- a/backend/stream-canvas/src/types.ts +++ b/backend/stream-canvas/src/types.ts @@ -13,6 +13,7 @@ export interface ClerkClaims { /** Decoded claims from a short-lived OBS token. */ export interface ObsTokenClaims { roomId: string; + credentialVersion: string; role: "obs"; scope: "stream-canvas-ws"; exp: number; diff --git a/backend/stream-canvas/src/ws-handler.ts b/backend/stream-canvas/src/ws-handler.ts index f551381..cbfa7d4 100644 --- a/backend/stream-canvas/src/ws-handler.ts +++ b/backend/stream-canvas/src/ws-handler.ts @@ -3,6 +3,7 @@ import { InMemorySyncStorage, type RoomSnapshot, TLSocketRoom, + TLSyncErrorCloseEventReason, } from "@tldraw/sync-core"; import { and, eq } from "drizzle-orm"; import type { TLRecord } from "tldraw"; @@ -20,11 +21,17 @@ import { isValidRoomId } from "./room-validation.ts"; import { canvasDocuments, roomMembers, rooms } from "./schema.ts"; import { streamCanvasSchema } from "./tldraw-schema.ts"; import type { ConnectionRole } from "./types.ts"; +import { obsCredentialVersion } from "./obs-secret.ts"; +import { withRoomOperation } from "./room-operations.ts"; +import { RoomSnapshotWriter } from "./snapshot-writer.ts"; interface ActiveRoom { room: TLSocketRoom; storage: InMemorySyncStorage; writer: RoomSnapshotWriter; + sessions: Map; + idleTimer?: NodeJS.Timeout; + disposing?: Promise; } const activeRooms = new Map(); @@ -35,71 +42,16 @@ const wsFailureLimiter = new FixedWindowRateLimit({ max: 30, }); -class RoomSnapshotWriter { - private pending: RoomSnapshot | null = null; - private drainPromise: Promise | null = null; - private readonly roomId: string; - - constructor(roomId: string) { - this.roomId = roomId; - } - - schedule(snapshot: RoomSnapshot): void { - this.pending = snapshot; - if (!this.drainPromise) { - this.drainPromise = this.drain() - .catch(async (error) => { - console.error( - `[canvas] snapshot persistence failed for room ${this.roomId}; retrying`, - error, - ); - await new Promise((resolve) => setTimeout(resolve, 1_000)); - }) - .finally(() => { - this.drainPromise = null; - if (this.pending) this.schedule(this.pending); - }); - } - } - - async flush(snapshot?: RoomSnapshot): Promise { - if (snapshot) this.schedule(snapshot); - while (this.drainPromise) await this.drainPromise; - } - - private async drain(): Promise { - while (this.pending) { - const snapshot = this.pending; - this.pending = null; - try { - await db - .insert(canvasDocuments) - .values({ - roomId: this.roomId, - snapshot, - revision: snapshot.documentClock, - updatedAt: new Date(), - }) - .onConflictDoUpdate({ - target: canvasDocuments.roomId, - set: { - snapshot, - revision: snapshot.documentClock, - updatedAt: new Date(), - }, - }); - } catch (error) { - // Keep a newer pending snapshot if one arrived while this write ran. - this.pending ??= snapshot; - throw error; - } - } - } -} - async function getOrCreateRoom(roomId: string): Promise { const existing = activeRooms.get(roomId); - if (existing) return existing; + if (existing?.disposing) { + await existing.disposing; + return getOrCreateRoom(roomId); + } + if (existing) { + clearTimeout(existing.idleTimer); + return existing; + } const loading = roomLoads.get(roomId); if (loading) return loading; @@ -119,45 +71,138 @@ async function loadRoom(roomId: string, epoch: number): Promise { if (epoch !== roomEpoch) { throw new Error("Room load cancelled during leadership handover"); } - const writer = new RoomSnapshotWriter(roomId); + const writer = new RoomSnapshotWriter( + () => storage.getSnapshot(), + (snapshot) => leaderState.persistSnapshot(roomId, snapshot), + `room ${roomId}`, + ); const storage = new InMemorySyncStorage({ ...(persisted ? { snapshot: persisted.snapshot as RoomSnapshot } : {}), onChange() { - writer.schedule(storage.getSnapshot()); + writer.schedule(); }, }); const room = new TLSocketRoom({ storage, schema: streamCanvasSchema, - onSessionRemoved(_room, { numSessionsRemaining }) { - if (numSessionsRemaining !== 0) return; - setTimeout(() => void disposeInactiveRoom(roomId), 30_000); + onSessionRemoved(_room, { sessionId, numSessionsRemaining }) { + active.sessions.delete(sessionId); + if (numSessionsRemaining !== 0 || activeRooms.get(roomId) !== active) + return; + scheduleIdleDisposal(roomId, active); }, }); - const active = { room, storage, writer }; + const active: ActiveRoom = { room, storage, writer, sessions: new Map() }; activeRooms.set(roomId, active); return active; } -async function disposeInactiveRoom(roomId: string): Promise { +function scheduleIdleDisposal(roomId: string, active: ActiveRoom): void { + if ( + activeRooms.get(roomId) !== active || + active.room.getNumActiveSessions() !== 0 + ) + return; + clearTimeout(active.idleTimer); + active.idleTimer = setTimeout(() => { + if (activeRooms.get(roomId) !== active) return; + void disposeInactiveRoom(roomId).catch((error) => { + console.error("[canvas] idle flush failed", error); + scheduleIdleDisposal(roomId, active); + }); + }, 30_000); + active.idleTimer.unref(); +} + +export async function disposeInactiveRoom(roomId: string): Promise { const active = activeRooms.get(roomId); if (active?.room.getNumActiveSessions() !== 0) return; - activeRooms.delete(roomId); - await active.writer.flush(active.storage.getSnapshot()); - active.room.close(); + if (active.disposing) return active.disposing; + clearTimeout(active.idleTimer); + active.disposing = (async () => { + await active.writer.flush(); + if (activeRooms.get(roomId) !== active) return; + activeRooms.delete(roomId); + active.writer.stop(); + active.room.close(); + })().finally(() => { + active.disposing = undefined; + }); + return active.disposing; } -export async function closeAllRooms(): Promise { +export async function closeAllRooms(persist = true): Promise { roomEpoch += 1; - await Promise.allSettled([...roomLoads.values()]); const entries = [...activeRooms.values()]; activeRooms.clear(); - await Promise.all( - entries.map(async ({ room, storage, writer }) => { - await writer.flush(storage.getSnapshot()); - room.close(); + // Stop accepting edits before taking the final snapshot, including sessions + // whose WebSocket close handshake has not completed. + for (const active of entries) { + clearTimeout(active.idleTimer); + active.room.close(); + if (!persist) active.writer.stop(); + } + await Promise.allSettled([...roomLoads.values()]); + const results = await Promise.allSettled( + entries.map(async (active) => { + try { + if (persist) await active.writer.flush(); + } finally { + active.writer.stop(); + } }), ); + const failures = results.filter((result) => result.status === "rejected"); + if (failures.length) + throw new AggregateError( + failures.map((result) => result.reason), + "Canvas flush failed", + ); +} + +export function roomConfigChanged( + room: typeof rooms.$inferSelect, + allowedUsers?: string[], +): void { + const active = activeRooms.get(room.id); + if (!active) return; + for (const [sessionId, auth] of active.sessions) { + if ( + auth.role === "editor" && + allowedUsers && + auth.userId !== room.ownerClerkId && + !allowedUsers.includes(auth.userId ?? "") + ) { + active.room.closeSession( + sessionId, + TLSyncErrorCloseEventReason.FORBIDDEN, + ); + } else { + active.room.sendCustomMessage(sessionId, roomConfigMessage(room)); + } + } +} + +export function revokeObsSessions(roomId: string): void { + const active = activeRooms.get(roomId); + if (!active) return; + for (const [sessionId, auth] of active.sessions) { + if (auth.role === "obs") + active.room.closeSession( + sessionId, + TLSyncErrorCloseEventReason.FORBIDDEN, + ); + } +} + +function roomConfigMessage( + room: Pick, +) { + return { + type: "room-config", + twitchChannel: room.twitchChannel, + youtubePolicy: room.youtubePolicy, + }; } interface AuthResult { @@ -204,7 +249,11 @@ export async function authenticateWebSocketUpgrade( const room = await db.query.rooms.findFirst({ where: eq(rooms.id, roomId), }); - if (!room) return null; + if ( + !room || + obsClaims.credentialVersion !== obsCredentialVersion(room.obsSecret) + ) + return null; return { role: "obs", roomId }; } @@ -214,6 +263,20 @@ export async function authenticateWebSocketUpgrade( export async function handleWebSocketUpgrade( ws: WebSocket, req: IncomingMessage, +): Promise { + // The client can send its connect frame immediately after HTTP upgrade. + // Keep it buffered until asynchronous admission installs room listeners. + ws.pause(); + try { + await admitWebSocket(ws, req); + } finally { + ws.resume(); + } +} + +async function admitWebSocket( + ws: WebSocket, + req: IncomingMessage, ): Promise { if (!leaderState.isLeader) { ws.close(1012, "Canvas leader is changing"); @@ -226,48 +289,76 @@ export async function handleWebSocketUpgrade( return; } - const auth = await authenticateWebSocketUpgrade(req); - if (!auth) { - wsFailureLimiter.consume(failureKey); + const roomId = new URL( + req.url ?? "", + "http://stream-canvas.local", + ).searchParams.get("roomId"); + if (!roomId || !isValidRoomId(roomId)) { ws.close(4001, "Unauthorized"); return; } - if (!leaderState.isLeader) { - ws.close(1012, "Canvas leader is changing"); - return; - } + await withRoomOperation(roomId, async () => { + const auth = await authenticateWebSocketUpgrade(req); + if (!auth) { + wsFailureLimiter.consume(failureKey); + ws.close(4001, "Unauthorized"); + return; + } + if (!leaderState.isLeader) { + ws.close(1012, "Canvas leader is changing"); + return; + } - if ( - !activeRooms.has(auth.roomId) && - activeRooms.size >= config.maxActiveRooms - ) { - ws.close(1013, "Too many active rooms"); - return; - } + if ( + !activeRooms.has(auth.roomId) && + activeRooms.size >= config.maxActiveRooms + ) { + ws.close(1013, "Too many active rooms"); + return; + } - let room: TLSocketRoom; - try { - ({ room } = await getOrCreateRoom(auth.roomId)); - } catch (error) { + let active: ActiveRoom; + try { + active = await getOrCreateRoom(auth.roomId); + } catch (error) { + if (!leaderState.isLeader) { + ws.close(1012, "Canvas leader is changing"); + return; + } + throw error; + } if (!leaderState.isLeader) { ws.close(1012, "Canvas leader is changing"); return; } - throw error; - } - if (!leaderState.isLeader) { - ws.close(1012, "Canvas leader is changing"); - return; - } - if (room.getNumActiveSessions() >= config.maxWsSessionsPerRoom) { - ws.close(1013, "Room session limit reached"); - return; - } + const { room } = active; + const metadata = await db.query.rooms.findFirst({ + where: eq(rooms.id, auth.roomId), + }); + if ( + !metadata || + !leaderState.isLeader || + activeRooms.get(auth.roomId) !== active || + room.isClosed() || + ws.readyState !== ws.OPEN + ) { + ws.close(1012, "Canvas leader is changing"); + return; + } + clearTimeout(active.idleTimer); + if (room.getNumActiveSessions() >= config.maxWsSessionsPerRoom) { + ws.close(1013, "Room session limit reached"); + return; + } - room.handleSocketConnect({ - sessionId: crypto.randomUUID(), - socket: ws, - isReadonly: auth.role === "obs", + const sessionId = crypto.randomUUID(); + active.sessions.set(sessionId, auth); + room.handleSocketConnect({ + sessionId, + socket: ws, + isReadonly: auth.role === "obs", + }); + room.sendCustomMessage(sessionId, roomConfigMessage(metadata)); }); } diff --git a/backend/stream-canvas/tsconfig.test.json b/backend/stream-canvas/tsconfig.test.json new file mode 100644 index 0000000..495c6f9 --- /dev/null +++ b/backend/stream-canvas/tsconfig.test.json @@ -0,0 +1,6 @@ +{ + "extends": "./tsconfig.json", + "compilerOptions": { "noEmit": true }, + "include": ["src/**/*.ts"], + "exclude": [] +} diff --git a/biome.json b/biome.json index 21a74eb..8a0423e 100644 --- a/biome.json +++ b/biome.json @@ -53,7 +53,9 @@ "**", "!.next", "!build", - "!dist", + "!**/dist", + "!test-results", + "!playwright-report", "!node_modules", "!convex/_generated", "!backend/stream-canvas/data", diff --git a/components/__tests__/recovery.test.tsx b/components/__tests__/recovery.test.tsx new file mode 100644 index 0000000..f5a7ebd --- /dev/null +++ b/components/__tests__/recovery.test.tsx @@ -0,0 +1,103 @@ +// @vitest-environment jsdom +import { + act, + cleanup, + fireEvent, + render, + screen, +} from "@testing-library/react"; +import { afterEach, beforeEach, expect, test, vi } from "vitest"; +import { UserRecordBootstrap } from "../user-record-bootstrap"; +import { UserMultiSelect } from "../stream-canvas/UserMultiSelect"; +import { TwitchPreview } from "../stream-canvas/TwitchPreview"; + +const state = vi.hoisted(() => ({ bootstrap: vi.fn(), authenticated: true })); +vi.mock("convex/react", () => ({ + useConvexAuth: () => ({ isAuthenticated: state.authenticated }), + useMutation: () => state.bootstrap, + useQuery: () => [ + { userId: "user_ava", username: "Ava" }, + { userId: "user_ben", username: "Ben" }, + ], +})); +vi.mock("sonner", () => ({ toast: { error: vi.fn() } })); +vi.mock("../stream-canvas/media-preferences", () => ({ + useMediaPreference: () => ({ enabled: false, volume: 1 }), +})); +beforeEach(() => { + vi.useFakeTimers(); + state.bootstrap.mockReset(); + state.authenticated = true; +}); +afterEach(() => { + cleanup(); + vi.useRealTimers(); + delete window.Twitch; + document.querySelectorAll("script").forEach((script) => { + script.remove(); + }); +}); + +test("profile bootstrap retries transient failures and cancels after sign out", async () => { + state.bootstrap + .mockRejectedValueOnce(new Error("temporary")) + .mockResolvedValue({}); + const { rerender } = render(); + await act(async () => {}); + await act(() => vi.advanceTimersByTimeAsync(1_000)); + expect(state.bootstrap).toHaveBeenCalledTimes(2); + state.authenticated = false; + rerender(); + await act(() => vi.advanceTimersByTimeAsync(10_000)); + expect(state.bootstrap).toHaveBeenCalledTimes(2); +}); + +test("collaborator selection follows arrow keys and Enter without submitting settings", () => { + const select = vi.fn(); + const submit = vi.fn((event: React.FormEvent) => event.preventDefault()); + render( +
+ + , + ); + const input = screen.getByRole("combobox"); + fireEvent.change(input, { target: { value: "a" } }); + fireEvent.keyDown(input, { key: "ArrowDown" }); + expect(input.getAttribute("aria-activedescendant")).toBe("members-results-1"); + fireEvent.keyDown(input, { key: "Enter" }); + expect(select).toHaveBeenCalledWith(["user_ben"]); + expect(submit).not.toHaveBeenCalled(); +}); + +test("failed Twitch script is removed and visible retry initializes a fresh preview", async () => { + render( + , + ); + const first = document.querySelector("script"); + expect(first).not.toBeNull(); + await act(async () => first?.dispatchEvent(new Event("error"))); + expect(first?.isConnected).toBe(false); + fireEvent.click(screen.getByRole("button", { name: "Retry Twitch preview" })); + const second = document.querySelector("script"); + expect(second).not.toBe(first); + const mounted = vi.fn(); + class Embed { + static VIDEO = "video"; + static VIDEO_READY = "ready"; + constructor() { + mounted(); + } + addEventListener() {} + getPlayer() { + return { setMuted() {}, setVolume() {} }; + } + } + window.Twitch = { Embed }; + await act(async () => second?.dispatchEvent(new Event("load"))); + expect(mounted).toHaveBeenCalledOnce(); + expect(screen.queryByRole("alert")).toBeNull(); +}); diff --git a/components/__tests__/room-settings.test.tsx b/components/__tests__/room-settings.test.tsx new file mode 100644 index 0000000..50cf92d --- /dev/null +++ b/components/__tests__/room-settings.test.tsx @@ -0,0 +1,84 @@ +// @vitest-environment jsdom +import { + cleanup, + fireEvent, + render, + screen, + waitFor, +} from "@testing-library/react"; +import { afterEach, expect, test, vi } from "vitest"; +import Settings from "@/app/app/settings/page"; +import RoomPage from "@/app/app/rooms/[roomId]/page"; + +const state = vi.hoisted(() => ({ + token: vi.fn(async () => "test-token"), + create: vi.fn(), + update: vi.fn(), + list: vi.fn(), + error: vi.fn(), + success: vi.fn(), +})); +vi.mock("@clerk/nextjs", () => ({ + useAuth: () => ({ isLoaded: true, isSignedIn: true, getToken: state.token }), + useClerk: () => ({}), +})); +vi.mock("next/navigation", () => ({ + useParams: () => ({ roomId: "room-fixture" }), +})); +vi.mock("sonner", () => ({ + toast: { error: state.error, success: state.success }, +})); +vi.mock("@/components/stream-canvas/UserMultiSelect", () => ({ + UserMultiSelect: () =>
Collaborators
, +})); +vi.mock("@/components/stream-canvas/CanvasEditor", () => ({ + CanvasEditor: () =>
Editor connected
, +})); +vi.mock("@/lib/stream-canvas/api", () => ({ + createRoom: state.create, + updateRoom: state.update, + getAccessibleRooms: state.list, + regenerateSecret: vi.fn(), +})); +afterEach(() => { + cleanup(); + vi.clearAllMocks(); +}); +const room = { + id: "room-fixture", + allowedUsers: [], + twitchChannel: null, + youtubePolicy: "preview_only", +}; + +test("room metadata failure is visible and retry opens the editor", async () => { + state.list + .mockRejectedValueOnce(new Error("Room info unavailable")) + .mockResolvedValue([room]); + render(); + expect((await screen.findByRole("alert")).textContent).toContain( + "Room info unavailable", + ); + fireEvent.click( + screen.getByRole("button", { name: "Retry room connection" }), + ); + expect(await screen.findByText("Editor connected")).toBeTruthy(); +}); + +test("failed settings submission stays retryable and reports success only after saving", async () => { + state.create.mockResolvedValue(room); + state.update + .mockRejectedValueOnce(new Error("Save unavailable")) + .mockResolvedValue(room); + render(); + const save = await screen.findByRole("button", { name: "Save changes" }); + fireEvent.click(save); + await waitFor(() => + expect(state.error).toHaveBeenCalledWith("Save unavailable"), + ); + expect(state.success).not.toHaveBeenCalled(); + fireEvent.click(save); + await waitFor(() => + expect(state.success).toHaveBeenCalledWith("Settings saved"), + ); +}); diff --git a/components/stream-canvas/CanvasEditor.tsx b/components/stream-canvas/CanvasEditor.tsx index 19c788f..f61dd9c 100644 --- a/components/stream-canvas/CanvasEditor.tsx +++ b/components/stream-canvas/CanvasEditor.tsx @@ -75,7 +75,9 @@ import { STREAM_ZONE, } from "@/lib/stream-canvas/stream-zone"; import type { YouTubePolicy } from "@/lib/stream-canvas/types"; +import { readRoomConfigMessage } from "@/lib/stream-canvas/room-config"; import { CanvasStylePanel } from "./MediaInspectorPanel"; +import { TwitchPreview } from "./TwitchPreview"; import { MediaPreferencesProvider, useMediaPreference, @@ -104,132 +106,7 @@ interface CanvasEditorProps { roomId: string; twitchChannel?: string | null; youtubePolicy: YouTubePolicy; -} - -interface TwitchPlayer { - setMuted(muted: boolean): void; - setVolume(volume: number): void; -} - -interface TwitchEmbedInstance { - addEventListener(event: string, callback: () => void): void; - getPlayer(): TwitchPlayer; -} - -interface TwitchEmbedConstructor { - new ( - element: HTMLElement, - options: Record, - ): TwitchEmbedInstance; - VIDEO: string; - VIDEO_READY: string; -} - -declare global { - interface Window { - Twitch?: { Embed: TwitchEmbedConstructor }; - } -} - -let twitchEmbedScriptPromise: Promise | null = null; - -function loadTwitchEmbed(): Promise { - if (window.Twitch?.Embed) return Promise.resolve(window.Twitch.Embed); - if (twitchEmbedScriptPromise) return twitchEmbedScriptPromise; - const promise = new Promise((resolve, reject) => { - const existing = document.querySelector( - 'script[src="https://embed.twitch.tv/embed/v1.js"]', - ); - const script = existing ?? document.createElement("script"); - const handleLoad = () => { - if (window.Twitch?.Embed) resolve(window.Twitch.Embed); - else reject(new Error("Twitch embed API did not initialize")); - }; - script.addEventListener("load", handleLoad, { once: true }); - script.addEventListener( - "error", - () => reject(new Error("Failed to load Twitch embed API")), - { once: true }, - ); - if (!existing) { - script.src = "https://embed.twitch.tv/embed/v1.js"; - script.async = true; - document.head.appendChild(script); - } - }).catch((error) => { - twitchEmbedScriptPromise = null; - throw error; - }); - twitchEmbedScriptPromise = promise; - return promise; -} - -function TwitchPreview({ - channel, - hostname, - interactive, -}: { - channel: string; - hostname: string; - interactive: boolean; -}) { - const hostRef = useRef(null); - const playerRef = useRef(null); - const preference = useMediaPreference("twitch-preview"); - const preferenceRef = useRef({ - enabled: preference.enabled, - volume: preference.volume, - }); - preferenceRef.current = { - enabled: preference.enabled, - volume: preference.volume, - }; - - useEffect(() => { - const host = hostRef.current; - if (!host) return; - let cancelled = false; - host.replaceChildren(); - void loadTwitchEmbed() - .then((Embed) => { - if (cancelled || !hostRef.current) return; - const embed = new Embed(hostRef.current, { - width: "100%", - height: "100%", - channel, - parent: [hostname], - layout: Embed.VIDEO, - autoplay: true, - muted: true, - }); - embed.addEventListener(Embed.VIDEO_READY, () => { - if (cancelled) return; - playerRef.current = embed.getPlayer(); - playerRef.current.setVolume(preferenceRef.current.volume); - playerRef.current.setMuted(!preferenceRef.current.enabled); - }); - }) - .catch((error) => console.error("[twitch] embed failed", error)); - return () => { - cancelled = true; - playerRef.current = null; - host.replaceChildren(); - }; - }, [channel, hostname]); - - useEffect(() => { - playerRef.current?.setVolume(preference.volume); - playerRef.current?.setMuted(!preference.enabled); - }, [preference.enabled, preference.volume]); - - return ( -
- ); + onMount?: import("tldraw").TLOnMountHandler; } /** @@ -805,9 +682,14 @@ function YouTubeInteractionController() { export function CanvasEditor({ roomId, - twitchChannel, - youtubePolicy, + twitchChannel: initialTwitchChannel, + youtubePolicy: initialYouTubePolicy, + onMount, }: CanvasEditorProps) { + const [{ twitchChannel, youtubePolicy }, setRoomConfig] = useState({ + twitchChannel: initialTwitchChannel ?? null, + youtubePolicy: initialYouTubePolicy, + }); const { getToken, userId } = useAuth(); const [interactiveShapeId, setInteractiveShapeId] = useState( null, @@ -861,6 +743,8 @@ export function CanvasEditor({ getToken, resolveUrl: (src: string, options?: { forceRefresh?: boolean }) => resolveEditorUploadUrl(roomId, src, getToken, options), + getRefreshDelayMs: (src: string) => + getEditorUploadUrlRefreshDelayMs(roomId, src), }), [roomId, getToken], ); @@ -894,6 +778,10 @@ export function CanvasEditor({ uri: getUri, assets, shapeUtils: syncShapeUtils, + onCustomMessageReceived(data: unknown) { + const roomConfig = readRoomConfigMessage(data); + if (roomConfig) setRoomConfig(roomConfig); + }, }); if (storeWithStatus.status === "loading") { @@ -920,6 +808,7 @@ export function CanvasEditor({ setAttempt((value) => value + 1)} + /> + ); +} + +function ConnectedCanvasMirror({ + obsSecret, + retry, +}: CanvasMirrorProps & { retry: () => void }) { const [error, setError] = useState(null); const [ready, setReady] = useState(false); const [roomId, setRoomId] = useState(null); @@ -143,6 +158,8 @@ export function CanvasMirror({ obsSecret }: CanvasMirrorProps) { getToken: async () => null, resolveUrl: (src: string, options?: { forceRefresh?: boolean }) => resolveObsUploadUrl(src, obsSecret, options), + getRefreshDelayMs: (src: string) => + getObsUploadUrlRefreshDelayMs(src), } : null, [obsSecret, roomId], @@ -158,12 +175,26 @@ export function CanvasMirror({ obsSecret }: CanvasMirrorProps) { uri: getUri, assets, shapeUtils: syncShapeUtils, + onCustomMessageReceived(data: unknown) { + const roomConfig = readRoomConfigMessage(data); + if (roomConfig) setYouTubePolicy(roomConfig.youtubePolicy); + }, }); if (error) { return ( -
+
{error} +
); } @@ -180,6 +211,13 @@ export function CanvasMirror({ obsSecret }: CanvasMirrorProps) { return (
OBS connection error. Refresh the source or regenerate the OBS URL. +
); } diff --git a/components/stream-canvas/TwitchPreview.tsx b/components/stream-canvas/TwitchPreview.tsx new file mode 100644 index 0000000..2700977 --- /dev/null +++ b/components/stream-canvas/TwitchPreview.tsx @@ -0,0 +1,170 @@ +"use client"; + +import { useEffect, useRef, useState } from "react"; +import { useMediaPreference } from "./media-preferences"; + +interface TwitchPlayer { + setMuted(muted: boolean): void; + setVolume(volume: number): void; +} + +interface TwitchEmbedInstance { + addEventListener(event: string, callback: () => void): void; + getPlayer(): TwitchPlayer; +} + +interface TwitchEmbedConstructor { + new ( + element: HTMLElement, + options: Record, + ): TwitchEmbedInstance; + VIDEO: string; + VIDEO_READY: string; +} + +declare global { + interface Window { + Twitch?: { Embed: TwitchEmbedConstructor }; + } +} + +let twitchEmbedScriptPromise: Promise | null = null; + +export function loadTwitchEmbed(): Promise { + if (window.Twitch?.Embed) return Promise.resolve(window.Twitch.Embed); + if (twitchEmbedScriptPromise) return twitchEmbedScriptPromise; + const promise = new Promise((resolve, reject) => { + const existing = document.querySelector( + 'script[src="https://embed.twitch.tv/embed/v1.js"]', + ); + const script = existing ?? document.createElement("script"); + const cleanup = () => { + clearTimeout(timeout); + script.removeEventListener("load", handleLoad); + script.removeEventListener("error", handleError); + }; + const fail = (message: string) => { + cleanup(); + script.remove(); + reject(new Error(message)); + }; + const handleLoad = () => { + if (window.Twitch?.Embed) { + cleanup(); + resolve(window.Twitch.Embed); + } else fail("Twitch embed API did not initialize"); + }; + const handleError = () => fail("Failed to load Twitch embed API"); + const timeout = setTimeout( + () => fail("Twitch embed API timed out"), + 10_000, + ); + script.addEventListener("load", handleLoad, { once: true }); + script.addEventListener("error", handleError, { once: true }); + if (!existing) { + script.src = "https://embed.twitch.tv/embed/v1.js"; + script.async = true; + document.head.appendChild(script); + } + }).catch((error) => { + twitchEmbedScriptPromise = null; + throw error; + }); + twitchEmbedScriptPromise = promise; + return promise; +} + +export function TwitchPreview({ + channel, + hostname, + interactive, +}: { + channel: string; + hostname: string; + interactive: boolean; +}) { + const [error, setError] = useState(null); + const [attempt, setAttempt] = useState(0); + const hostRef = useRef(null); + const playerRef = useRef(null); + const preference = useMediaPreference("twitch-preview"); + const preferenceRef = useRef({ + enabled: preference.enabled, + volume: preference.volume, + }); + preferenceRef.current = { + enabled: preference.enabled, + volume: preference.volume, + }; + + // biome-ignore lint/correctness/useExhaustiveDependencies: attempt explicitly retries script initialization + useEffect(() => { + const host = hostRef.current; + if (!host) return; + let cancelled = false; + host.replaceChildren(); + setError(null); + void loadTwitchEmbed() + .then((Embed) => { + if (cancelled || !hostRef.current) return; + const embed = new Embed(hostRef.current, { + width: "100%", + height: "100%", + channel, + parent: [hostname], + layout: Embed.VIDEO, + autoplay: true, + muted: true, + }); + embed.addEventListener(Embed.VIDEO_READY, () => { + if (cancelled) return; + playerRef.current = embed.getPlayer(); + playerRef.current.setVolume(preferenceRef.current.volume); + playerRef.current.setMuted(!preferenceRef.current.enabled); + }); + }) + .catch((error: unknown) => { + if (!cancelled) + setError( + error instanceof Error ? error.message : "Twitch preview failed", + ); + }); + return () => { + cancelled = true; + playerRef.current = null; + host.replaceChildren(); + }; + }, [channel, hostname, attempt]); + + useEffect(() => { + playerRef.current?.setVolume(preference.volume); + playerRef.current?.setMuted(!preference.enabled); + }, [preference.enabled, preference.volume]); + + return ( +
+
+ {error && ( +
+

{error}

+ +
+ )} +
+ ); +} diff --git a/components/stream-canvas/UserMultiSelect.tsx b/components/stream-canvas/UserMultiSelect.tsx index 9e8e86e..d165f7c 100644 --- a/components/stream-canvas/UserMultiSelect.tsx +++ b/components/stream-canvas/UserMultiSelect.tsx @@ -27,6 +27,7 @@ export function UserMultiSelect({ }: UserMultiSelectProps) { const [search, setSearch] = useState(""); const [open, setOpen] = useState(false); + const [activeIndex, setActiveIndex] = useState(0); const generatedId = useId(); const listboxId = `${inputId ?? generatedId}-results`; const inputRef = useRef(null); @@ -85,6 +86,10 @@ export function UserMultiSelect({ results?.filter((user) => !value.includes(user.userId)) ?? []; const hasSearch = search.trim().length > 0; const showResults = open && hasSearch; + const activeOptionIndex = Math.min( + activeIndex, + Math.max(0, filteredResults.length - 1), + ); return (
@@ -140,15 +145,33 @@ export function UserMultiSelect({ aria-autocomplete="list" aria-expanded={showResults} aria-controls={listboxId} + aria-activedescendant={ + showResults && filteredResults.length + ? `${listboxId}-${activeOptionIndex}` + : undefined + } value={search} onChange={(event) => { setSearch(event.target.value); + setActiveIndex(0); setOpen(true); }} onFocus={() => { if (hasSearch) setOpen(true); }} onKeyDown={(event) => { + if (event.key === "ArrowDown" || event.key === "ArrowUp") { + event.preventDefault(); + setOpen(true); + const count = filteredResults.length; + if (count) + setActiveIndex( + (activeOptionIndex + + (event.key === "ArrowDown" ? 1 : -1) + + count) % + count, + ); + } if (event.key === "Escape") { setOpen(false); } @@ -158,7 +181,7 @@ export function UserMultiSelect({ if (event.key === "Enter") { event.preventDefault(); if (showResults && filteredResults.length > 0) { - handleSelect(filteredResults[0]); + handleSelect(filteredResults[activeOptionIndex]); } } }} @@ -207,14 +230,20 @@ export function UserMultiSelect({
) : (
- {filteredResults.map((entry) => ( + {filteredResults.map((entry, index) => (
+ ) : null} {!isReadonly && !isInteractive ? ( + ); +} + +test.each(["read", "write"])( + "unavailable storage on %s keeps monitoring usable in memory", + (failure) => { + vi.spyOn(Storage.prototype, "getItem").mockImplementation(() => { + if (failure === "read") + throw new DOMException("Storage blocked", "SecurityError"); + return null; + }); + vi.spyOn(Storage.prototype, "setItem").mockImplementation(() => { + throw new DOMException("Storage full", "QuotaExceededError"); + }); + render( + + + , + ); + fireEvent.click(screen.getByRole("button", { name: "Monitoring off" })); + expect(screen.getByRole("button", { name: "Monitoring on" })).toBeTruthy(); + fireEvent.click(screen.getByRole("button", { name: "Monitoring on" })); + expect(screen.getByRole("button", { name: "Monitoring off" })).toBeTruthy(); + }, +); diff --git a/scripts/browser-fixture-cleanup.test.ts b/scripts/browser-fixture-cleanup.test.ts new file mode 100644 index 0000000..9a00941 --- /dev/null +++ b/scripts/browser-fixture-cleanup.test.ts @@ -0,0 +1,141 @@ +import { spawn } from "node:child_process"; +import { once } from "node:events"; +import { + mkdir, + mkdtemp, + readdir, + readFile, + rm, + writeFile, +} from "node:fs/promises"; +import { createServer } from "node:net"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { setTimeout as delay } from "node:timers/promises"; +import { expect, test } from "vitest"; + +test.each(["migration", "frontend", "cleanup", "signal"])( + "browser fixture releases partial startup resources after %s failure or interruption", + async (failure) => { + const root = await mkdtemp(join(tmpdir(), "browser-cleanup-test-")); + const storage = join(root, "temporary"); + const backendRoot = join(root, "backend/stream-canvas"); + const probe = createServer(); + probe.listen(0, "127.0.0.1"); + await once(probe, "listening"); + const address = probe.address(); + if (!address || typeof address === "string") + throw new Error("Missing port"); + const port = address.port; + await new Promise((resolve) => probe.close(() => resolve())); + let child: ReturnType | undefined; + try { + for (const directory of [ + "e2e", + "temporary", + "backend/stream-canvas/src", + "node_modules/vite", + "node_modules/pg", + ]) + await mkdir(join(root, directory), { recursive: true }); + await writeFile(join(root, "package.json"), '{"type":"module"}'); + // Run the real orchestrator with isolated local substitutes for its services. + await writeFile( + join(root, "e2e/server.ts"), + (await readFile("e2e/server.ts", "utf8")).replaceAll( + "4312", + String(port), + ), + ); + await writeFile( + join(backendRoot, "src/migrate.ts"), + failure === "signal" + ? 'import { writeFileSync } from "node:fs"; writeFileSync("migration-started", "yes"); setInterval(() => {}, 1000);' + : `process.exit(${failure === "migration" ? 1 : 0});`, + ); + await writeFile( + join(backendRoot, "src/index.ts"), + ` + import { createServer } from "node:http"; + import { writeFileSync } from "node:fs"; + const server = createServer((req, res) => { + res.setHeader("Content-Type", "application/json"); + res.end(JSON.stringify(req.url === "/ready" ? {} : { id: "fixture", obsSetupSecret: "synthetic" })); + }); + server.listen(Number(process.env.PORT), "127.0.0.1"); + process.on("SIGTERM", () => { writeFileSync("backend-stopped", "yes"); server.close(() => process.exit(0)); }); + `, + ); + await writeFile( + join(root, "node_modules/vite/package.json"), + '{"type":"module","exports":"./index.js"}', + ); + await writeFile( + join(root, "node_modules/vite/index.js"), + 'export async function createServer() { throw new Error("injected frontend startup failure"); }', + ); + await writeFile( + join(root, "node_modules/pg/package.json"), + '{"type":"module","exports":"./index.js"}', + ); + await writeFile( + join(root, "node_modules/pg/index.js"), + ` + import { writeFileSync } from "node:fs"; + export class Pool { + async query() { writeFileSync("database-cleanup-attempted", "yes"); ${failure === "cleanup" ? 'throw new Error("injected cleanup failure");' : ""} } + async end() {} + } + `, + ); + child = spawn(process.execPath, ["e2e/server.ts"], { + cwd: root, + env: { + NODE_ENV: "test", + PATH: process.env.PATH, + TMPDIR: storage, + TMP: storage, + TEMP: storage, + CANVAS_TEST_DATABASE_URL: "postgresql://unused.invalid/fixture", + }, + stdio: "ignore", + }); + const exited = once(child, "exit"); + if (failure === "signal") { + const deadline = Date.now() + 3_000; + while (!(await readdir(backendRoot)).includes("migration-started")) { + if (Date.now() > deadline) + throw new Error("Migration fixture did not start"); + await delay(10); + } + child.kill("SIGTERM"); + } + const [code, signal] = await exited; + expect(signal).toBeNull(); + expect(code).toBe(failure === "signal" ? 0 : 1); + expect(await readdir(storage)).toEqual([]); + if (failure === "frontend" || failure === "cleanup") { + expect( + await readFile(join(backendRoot, "backend-stopped"), "utf8"), + ).toBe("yes"); + expect( + await readFile( + join(backendRoot, "database-cleanup-attempted"), + "utf8", + ), + ).toBe("yes"); + } else { + expect(await readdir(backendRoot)).not.toContain( + "database-cleanup-attempted", + ); + } + } finally { + if (child && child.exitCode === null && child.signalCode === null) { + child.kill("SIGKILL"); + await once(child, "exit"); + } + await rm(root, { recursive: true, force: true }); + } + }, + 5_000, +); diff --git a/scripts/dev-lifecycle.test.ts b/scripts/dev-lifecycle.test.ts new file mode 100644 index 0000000..8340d8f --- /dev/null +++ b/scripts/dev-lifecycle.test.ts @@ -0,0 +1,60 @@ +import { execFileSync } from "node:child_process"; +import { + mkdirSync, + mkdtempSync, + readFileSync, + rmSync, + writeFileSync, +} from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { expect, test } from "vitest"; +import { parse } from "yaml"; + +test.each([false, true])( + "alternate database %s controls only its own lifecycle", + (alternate) => { + const directory = mkdtempSync(join(tmpdir(), "moddrop-dev-test-")); + try { + mkdirSync(join(directory, "scripts")); + mkdirSync(join(directory, "bin")); + writeFileSync( + join(directory, "scripts/dev.ts"), + readFileSync("scripts/dev.ts"), + ); + for (const command of ["docker", "pnpm"]) { + writeFileSync( + join(directory, "bin", command), + '#!/bin/sh\nprintf "%s\\n" "$0 $*" >> "$DEV_TEST_LOG"\n', + { mode: 0o755 }, + ); + } + const log = join(directory, "commands"); + execFileSync(process.execPath, [join(directory, "scripts/dev.ts")], { + env: { + NODE_ENV: "test", + PATH: `${join(directory, "bin")}:${process.env.PATH}`, + DEV_TEST_LOG: log, + NEXT_PUBLIC_CLERK_PUBLISHABLE_KEY: "pk_test_fixture", + CLERK_SECRET_KEY: ["sk", "test", "fixture"].join("_"), + CLERK_JWT_ISSUER_DOMAIN: "https://fixture.invalid", + ...(alternate + ? { MODDROP_DEV_DATABASE_URL: "postgresql://fixture.invalid/test" } + : {}), + }, + stdio: "pipe", + }); + const commands = readFileSync(log, "utf8"); + expect(commands.includes("docker")).toBe(!alternate); + expect(commands).toContain("db:migrate"); + expect(commands).toContain("dev:apps"); + } finally { + rmSync(directory, { recursive: true, force: true }); + } + }, +); + +test("development database is bound only to loopback", () => { + const compose = parse(readFileSync("compose.dev.yml", "utf8")); + expect(compose.services.postgres.ports).toEqual(["127.0.0.1:5432:5432"]); +}); diff --git a/scripts/dev.ts b/scripts/dev.ts index 26e1993..b363321 100644 --- a/scripts/dev.ts +++ b/scripts/dev.ts @@ -1,4 +1,4 @@ -import { spawn, type ChildProcess } from "node:child_process"; +import { type ChildProcess, spawn } from "node:child_process"; import { fileURLToPath } from "node:url"; const projectRoot = fileURLToPath(new URL("..", import.meta.url)); @@ -32,11 +32,13 @@ let exitCode = 0; try { assertRequiredEnvironment(); - console.log( - "[dev] Starting PostgreSQL and waiting for it to become healthy...", - ); - composeAttempted = true; - await run("docker", [...composeArgs, "up", "--wait", "postgres"]); + if (!process.env.MODDROP_DEV_DATABASE_URL) { + console.log( + "[dev] Starting PostgreSQL and waiting for it to become healthy...", + ); + composeAttempted = true; + await run("docker", [...composeArgs, "up", "--wait", "postgres"]); + } throwIfStopping(); console.log("[dev] Applying Drizzle migrations..."); From d9577df888d9bc0133489e79c31bba7e4e1f8545 Mon Sep 17 00:00:00 2001 From: beasty Date: Sun, 6 Sep 2026 13:10:36 +0200 Subject: [PATCH 8/9] fix: enforce YouTube policy for normalized provider URLs --- components/stream-canvas/shapes/shared.ts | 4 +- .../shapes/youtube/YouTubeEmbedShape.tsx | 54 ++++++++++++++----- e2e/media-regressions.spec.ts | 14 ++++- .../__tests__/extract-youtube-id.test.ts | 30 ++++++++++- 4 files changed, 83 insertions(+), 19 deletions(-) diff --git a/components/stream-canvas/shapes/shared.ts b/components/stream-canvas/shapes/shared.ts index df294cc..9acd000 100644 --- a/components/stream-canvas/shapes/shared.ts +++ b/components/stream-canvas/shapes/shared.ts @@ -23,7 +23,7 @@ import type { UploadUrlRefreshDelayMs } from "@/lib/stream-canvas/api"; import { AudioPlayerShapeUtil } from "./audio/AudioPlayerShape"; import { AudioPlayerTool } from "./audio/AudioPlayerTool"; import { - extractYouTubeId, + isYouTubeUrl, YouTubeEmbedShapeUtil, YouTubePolicyCtx, } from "./youtube/YouTubeEmbedShape"; @@ -44,7 +44,7 @@ const MAX_REFRESH_TIMEOUT_MS = 2 ** 31 - 1; class PolicyEmbedShapeUtil extends EmbedShapeUtil { static override type = "embed" as const; override component(shape: TLEmbedShape) { - if (!extractYouTubeId(shape.props.url)) return super.component(shape); + if (!isYouTubeUrl(shape.props.url)) return super.component(shape); return createElement(PolicyEmbed, {}, super.component(shape)); } } diff --git a/components/stream-canvas/shapes/youtube/YouTubeEmbedShape.tsx b/components/stream-canvas/shapes/youtube/YouTubeEmbedShape.tsx index de69d33..05f6380 100644 --- a/components/stream-canvas/shapes/youtube/YouTubeEmbedShape.tsx +++ b/components/stream-canvas/shapes/youtube/YouTubeEmbedShape.tsx @@ -92,22 +92,48 @@ export const YouTubePolicyCtx = createContext("preview_only"); * - youtube.com/live/ID */ export function extractYouTubeId(raw: string): string | null { - if (!raw) return null; - const trimmed = raw.trim(); - - // youtu.be short link - const shortMatch = trimmed.match( - /(?:https?:\/\/)?youtu\.be\/([a-zA-Z0-9_-]{11})/, - ); - if (shortMatch) return shortMatch[1]; + const url = parseYouTubeUrl(raw); + if (!url) return null; + let path: string[]; + try { + path = decodeURIComponent(url.pathname).split("/"); + } catch { + return null; + } + const id = + url.hostname === "youtu.be" + ? path[1] + : path[1] === "watch" + ? url.searchParams.get("v") + : ["embed", "shorts", "live"].includes(path[1]) + ? path[2] + : null; + return id && /^[a-zA-Z0-9_-]{11}$/.test(id) ? id : null; +} - // youtube.com variants - const longMatch = trimmed.match( - /(?:https?:\/\/)?(?:www\.|m\.)?youtube(?:-nocookie)?\.com\/(?:watch\?.*v=|embed\/|shorts\/|live\/)([a-zA-Z0-9_-]{11})/, - ); - if (longMatch) return longMatch[1]; +/** Owner policy applies to the provider, even without a recognized video ID. */ +export function isYouTubeUrl(raw: string): boolean { + return parseYouTubeUrl(raw) !== null; +} - return null; +function parseYouTubeUrl(raw: string): URL | null { + try { + const trimmed = raw.trim(); + const url = new URL( + trimmed.includes("://") ? trimmed : `https://${trimmed}`, + ); + if (url.protocol !== "https:" && url.protocol !== "http:") return null; + const host = url.hostname.replace(/\.$/, ""); + url.hostname = host; + return host === "youtu.be" || + ["youtube.com", "youtube-nocookie.com"].some( + (domain) => host === domain || host.endsWith(`.${domain}`), + ) + ? url + : null; + } catch { + return null; + } } const YOUTUBE_SYNC_THRESHOLD_PLAYING = 1.5; diff --git a/e2e/media-regressions.spec.ts b/e2e/media-regressions.spec.ts index 0c29cbd..c7cd72d 100644 --- a/e2e/media-regressions.spec.ts +++ b/e2e/media-regressions.spec.ts @@ -57,7 +57,7 @@ test("pasted and persisted YouTube embeds obey policy before OBS loads a player" if (!editor) throw new Error("Editor missing"); await editor.putExternalContent({ type: "url", - url: "https://www.youtube.com/watch?v=dQw4w9WgXcQ", + url: "https://WWW.YOUTUBE.COM/watch?v=%64Qw4w9WgXcQ", point: { x: 400, y: 300 }, }); editor.createShape({ @@ -70,6 +70,16 @@ test("pasted and persisted YouTube embeds obey policy before OBS loads a player" url: "https://www.youtube-nocookie.com/embed/dQw4w9WgXcQ", }, }); + editor.createShape({ + type: "embed", + x: 800, + y: 400, + props: { + w: 480, + h: 270, + url: "https://WWW.YOUTUBE.COM/watch?v=%64Qw4w9WgXcQ", + }, + }); }); const mirror = await context.newPage(); const youtubeRequests: string[] = []; @@ -78,7 +88,7 @@ test("pasted and persisted YouTube embeds obey policy before OBS loads a player" youtubeRequests.push(request.url()); }); await mirror.goto("/?view=mirror"); - await expect(mirror.locator(".tl-shape")).toHaveCount(2); + await expect(mirror.locator(".tl-shape")).toHaveCount(3); await expect(mirror.locator("iframe")).toHaveCount(0); expect(youtubeRequests).toEqual([]); await setPolicy(page, "disabled"); diff --git a/lib/stream-canvas/__tests__/extract-youtube-id.test.ts b/lib/stream-canvas/__tests__/extract-youtube-id.test.ts index 8abefe9..c1ca032 100644 --- a/lib/stream-canvas/__tests__/extract-youtube-id.test.ts +++ b/lib/stream-canvas/__tests__/extract-youtube-id.test.ts @@ -1,11 +1,22 @@ import { describe, expect, it } from "vitest"; -import { extractYouTubeId } from "@/components/stream-canvas/shapes/youtube/YouTubeEmbedShape"; +import { + extractYouTubeId, + isYouTubeUrl, +} from "@/components/stream-canvas/shapes/youtube/YouTubeEmbedShape"; const VIDEO_ID = "dQw4w9WgXcQ"; describe("extractYouTubeId", () => { it.each([ ["watch URL", `https://www.youtube.com/watch?v=${VIDEO_ID}`], + [ + "encoded ID and uppercase host", + "https://WWW.YOUTUBE.COM/watch?v=%64Qw4w9WgXcQ", + ], + [ + "encoded embed ID", + "https://www.youtube-nocookie.com/embed/%64Qw4w9WgXcQ", + ], ["watch URL without www", `https://youtube.com/watch?v=${VIDEO_ID}`], ["watch URL without protocol", `youtube.com/watch?v=${VIDEO_ID}`], [ @@ -34,7 +45,24 @@ describe("extractYouTubeId", () => { ["channel URL", "https://www.youtube.com/@somechannel"], ["watch URL without id", "https://www.youtube.com/watch"], ["id shorter than 11 chars", "https://youtu.be/short"], + ["lookalike host", `https://notyoutube.com/watch?v=${VIDEO_ID}`], + [ + "YouTube URL in unrelated query", + `https://example.com/?next=https://youtube.com/watch?v=${VIDEO_ID}`, + ], ])("returns null for a %s", (_label, url) => { expect(extractYouTubeId(url)).toBeNull(); }); }); + +it("classifies the provider independently of video ID syntax", () => { + expect( + isYouTubeUrl("https://WWW.YOUTUBE.COM/embed/videoseries?list=fixture"), + ).toBe(true); + expect(isYouTubeUrl("https://www.youtube-nocookie.com/embed/invalid")).toBe( + true, + ); + expect( + isYouTubeUrl("https://youtube.com.example.com/watch?v=dQw4w9WgXcQ"), + ).toBe(false); +}); From 74836a57d8526e3697604fc8119dc252d5d0dd2e Mon Sep 17 00:00:00 2001 From: beasty Date: Sun, 6 Sep 2026 14:15:36 +0200 Subject: [PATCH 9/9] release: restore gated GHCR delivery on the homelab runner --- .github/workflows/build-images.yml | 76 ++++++++++++++++++++++++++++++ .github/workflows/verify.yml | 12 +++-- backend/stream-canvas/README.md | 22 +++++---- package.json | 2 +- scripts/delivery-changes.test.ts | 30 +++++++++++- 5 files changed, 126 insertions(+), 16 deletions(-) create mode 100644 .github/workflows/build-images.yml diff --git a/.github/workflows/build-images.yml b/.github/workflows/build-images.yml new file mode 100644 index 0000000..ea71a45 --- /dev/null +++ b/.github/workflows/build-images.yml @@ -0,0 +1,76 @@ +name: Release container images + +on: + push: + tags: ["v*"] + workflow_dispatch: + +permissions: + contents: read + +jobs: + verify: + uses: ./.github/workflows/verify.yml + + deploy-convex: + needs: verify + if: startsWith(github.ref, 'refs/tags/v') || github.ref == 'refs/heads/main' + runs-on: arc-moddrop + steps: + - uses: actions/checkout@v7 + - uses: pnpm/action-setup@v6 + with: + run_install: false + - uses: actions/setup-node@v7 + with: + node-version-file: .node-version + cache: pnpm + - run: pnpm install --frozen-lockfile --ignore-scripts + - run: pnpm exec convex deploy + env: + CONVEX_DEPLOY_KEY: ${{ secrets.CONVEX_DEPLOY_KEY }} + CONVEX_DEPLOYMENT: ${{ secrets.CONVEX_DEPLOYMENT }} + CLERK_JWT_ISSUER_DOMAIN: ${{ secrets.CLERK_JWT_ISSUER_DOMAIN }} + + build: + needs: [verify, deploy-convex] + if: startsWith(github.ref, 'refs/tags/v') || github.ref == 'refs/heads/main' + runs-on: arc-moddrop + permissions: + contents: read + packages: write + strategy: + matrix: + include: + - image: moddrop-frontend + dockerfile: Dockerfile + - image: moddrop-stream-canvas + dockerfile: backend/stream-canvas/Dockerfile + steps: + - uses: actions/checkout@v7 + - uses: docker/setup-buildx-action@v4 + with: + driver: remote + endpoint: tcp://buildkitd.arc-runners.svc.cluster.local:1234 + - uses: docker/login-action@v4 + with: + registry: ghcr.io + username: ${{ github.actor }} + password: ${{ secrets.GITHUB_TOKEN }} + - uses: docker/metadata-action@v6 + id: meta + with: + images: ghcr.io/${{ github.repository_owner }}/${{ matrix.image }} + tags: | + type=semver,pattern={{version}} + type=sha,prefix= + - uses: docker/build-push-action@v7 + with: + context: . + file: ${{ matrix.dockerfile }} + push: true + tags: ${{ steps.meta.outputs.tags }} + labels: ${{ steps.meta.outputs.labels }} + build-args: | + CONVEX_DEPLOYMENT=${{ secrets.CONVEX_DEPLOYMENT }} + NEXT_PUBLIC_TLDRAW_LICENSE_KEY=${{ secrets.NEXT_PUBLIC_TLDRAW_LICENSE_KEY }} diff --git a/.github/workflows/verify.yml b/.github/workflows/verify.yml index 9b774a3..bf47db2 100644 --- a/.github/workflows/verify.yml +++ b/.github/workflows/verify.yml @@ -1,13 +1,16 @@ name: Verify on: pull_request: - push: - branches: [main] + workflow_call: + workflow_dispatch: permissions: contents: read jobs: verify: - runs-on: ubuntu-latest + runs-on: arc-moddrop + container: + image: node:26.5.0-bookworm + options: --user 0 services: postgres: image: postgres:18.4-bookworm @@ -15,12 +18,11 @@ jobs: POSTGRES_DB: canvas_test POSTGRES_USER: canvas_test POSTGRES_HOST_AUTH_METHOD: trust - ports: ["5432:5432"] options: >- --health-cmd "pg_isready -U canvas_test" --health-interval 5s --health-timeout 5s --health-retries 10 env: - CANVAS_TEST_DATABASE_URL: postgresql://canvas_test@127.0.0.1:5432/canvas_test + CANVAS_TEST_DATABASE_URL: postgresql://canvas_test@postgres:5432/canvas_test steps: - uses: actions/checkout@v7 - uses: pnpm/action-setup@v6 diff --git a/backend/stream-canvas/README.md b/backend/stream-canvas/README.md index 8cccf58..75778e8 100644 --- a/backend/stream-canvas/README.md +++ b/backend/stream-canvas/README.md @@ -109,15 +109,19 @@ deleted by this change. ## Verification -The Forgejo workflow is a legacy delivery target. The homelab `personal` runner -pool was retired on 2026-08-31; it cannot currently provide hosted verification. -If this workflow is used again, it requires a Docker-backed `personal` runner -with service-container networking. Its verification job explicitly uses the -pinned Node Bookworm container as root so Playwright can install Debian browser -dependencies without relying on host sudo configuration. Validate that job on -the replacement runner before enabling delivery. GitHub Actions must also be -enabled at repository level for `.github/workflows/verify.yml` to run; workflow -files alone do not turn it on. +GitHub verification runs on the homelab `arc-moddrop` Docker-backed runner. +The pinned Node Bookworm container runs as root so Playwright can install browser +dependencies; PostgreSQL is an isolated service container. Release tags run this +same gate before deploying Convex and publishing both images to GHCR. Manual +publication is restricted to `main`. The Forgejo workflow is retained as a legacy +copy; its `personal` runner pool was retired on 2026-08-31. + +After a release succeeds, update all three image references in the Homelab +Moddrop HelmRelease to the released tag and verified GHCR digests. Commit through +the Homelab GitOps workflow, reconcile Flux, and verify frontend and backend +readiness plus the public application. Cluster access is available through SSH +on `bunux`. Backend rollout uses Recreate so the old leader flushes and exits +before the new leader admits sessions. `pnpm run test` runs fast route and lifecycle unit tests and tears down its temporary storage and pools. Set `CANVAS_TEST_DATABASE_URL` to a disposable diff --git a/package.json b/package.json index b00ca15..f2396d7 100644 --- a/package.json +++ b/package.json @@ -1,7 +1,7 @@ { "name": "moddrop", "private": true, - "version": "0.6.1", + "version": "0.6.2", "type": "module", "packageManager": "pnpm@11.17.0+sha512.cca3cea332ad254bb84145f966d19f4879615210346fc92c79a047f23a0d7b3cca3c3792f0076ba1f1831d277efbcf0a9119b31a9a60eca7fb3d6231f331ef72", "scripts": { diff --git a/scripts/delivery-changes.test.ts b/scripts/delivery-changes.test.ts index f19d9d0..5573950 100644 --- a/scripts/delivery-changes.test.ts +++ b/scripts/delivery-changes.test.ts @@ -2,8 +2,36 @@ import { execFileSync } from "node:child_process"; import { mkdtempSync, readFileSync, rmSync, writeFileSync } from "node:fs"; import { tmpdir } from "node:os"; import { join, resolve } from "node:path"; -import { parse } from "yaml"; import { expect, test } from "vitest"; +import { parse } from "yaml"; + +test("GitHub releases use the PR verification gate before production actions", () => { + const workflow = parse( + readFileSync(".github/workflows/build-images.yml", "utf8"), + ); + const verification = parse( + readFileSync(".github/workflows/verify.yml", "utf8"), + ); + expect(workflow.jobs.verify.uses).toBe("./.github/workflows/verify.yml"); + expect(verification.on).toHaveProperty("pull_request"); + expect(verification.on).toHaveProperty("workflow_call"); + expect(verification.jobs.verify["runs-on"]).toBe("arc-moddrop"); + expect(workflow.jobs.build.needs).toContain("verify"); + expect(workflow.jobs.build.needs).toContain("deploy-convex"); + expect(workflow.jobs["deploy-convex"].needs).toBe("verify"); + for (const job of [workflow.jobs.build, workflow.jobs["deploy-convex"]]) { + const allowed = (ref: string) => + Function( + "github", + "startsWith", + `return (${job.if})`, + )({ ref }, (value: string, prefix: string) => value.startsWith(prefix)); + expect(allowed("refs/heads/feature")).toBe(false); + expect(allowed("refs/heads/main")).toBe(true); + expect(allowed("refs/tags/v0.6.2")).toBe(true); + expect(job.if).not.toContain("always()"); + } +}); test.each(["pnpm-lock.yaml", "pnpm-workspace.yaml", "package.json", ".npmrc"])( "shared install input %s rebuilds both images and Convex",