diff --git a/apps/web-next/package.json b/apps/web-next/package.json index 002d542..90f30d0 100644 --- a/apps/web-next/package.json +++ b/apps/web-next/package.json @@ -40,6 +40,7 @@ }, "devDependencies": { "@playwright/test": "^1.50.0", + "socket.io-client": "^4.7.0", "@repo/core-eslint": "workspace:*", "@repo/core-testing": "workspace:*", "@repo/core-typescript": "workspace:*", diff --git a/apps/web-next/src/__tests__/realtime-ping.test.ts b/apps/web-next/src/__tests__/realtime-ping.test.ts new file mode 100644 index 0000000..a008e5e --- /dev/null +++ b/apps/web-next/src/__tests__/realtime-ping.test.ts @@ -0,0 +1,97 @@ +// e2e proof-of-life: connect → subscribe → emit ping → receive pong via the +// production-shaped binder + Socket.IO server, with cookie-session auth +// against a seeded test session. +import "reflect-metadata"; +import { describe, it, expect, beforeEach, afterEach } from "vitest"; +import { createServer, type Server as HttpServer } from "node:http"; +import { Server as IOServer } from "socket.io"; +import { io as ioClient } from "socket.io-client"; +import type { AddressInfo } from "node:net"; +import { + RealtimeHandlerRegistry, + SocketIORealtimeBroadcaster, + SocketIORealtimeServer, + realtimePingInboundDescriptor, + realtimePongChannel, + type IRealtimeAuthenticator, +} from "@repo/core-realtime"; + +describe("e2e: realtime-ping exercises all four checkpoints", () => { + let httpServer: HttpServer; + let realtimeServer: SocketIORealtimeServer; + let port: number; + + beforeEach(async () => { + httpServer = createServer(); + const io = new IOServer(httpServer); + const broadcaster = new SocketIORealtimeBroadcaster(io); + const registry = new RealtimeHandlerRegistry(); + registry.register(realtimePingInboundDescriptor(broadcaster)); + registry.registerChannel(realtimePongChannel); + + const authenticator: IRealtimeAuthenticator = { + authenticate: async ({ cookies }) => + cookies.session === "valid-session" + ? { userId: "user_test", roles: [] } + : null, + }; + + realtimeServer = new SocketIORealtimeServer({ + httpServer, + io, + authenticator, + registry, + }); + await realtimeServer.start(); + + await new Promise((r) => httpServer.listen(0, r)); + port = (httpServer.address() as AddressInfo).port; + }); + + afterEach(async () => { + await realtimeServer.stop(); + await new Promise((r) => httpServer.close(() => r())); + }); + + it("authenticated client gets pong after ping", async () => { + const client = ioClient(`http://localhost:${port}`, { + extraHeaders: { Cookie: "session=valid-session" }, + }); + await new Promise((r) => client.on("connect", () => r())); + + const subAck = await new Promise<{ ok: boolean }>((r) => + client.emit("subscribe", "realtime.pong", r), + ); + expect(subAck.ok).toBe(true); + + const pongs: { at: string; echo: string }[] = []; + client.on("realtime.pong", (p) => pongs.push(p)); + + const sentAt = "2026-05-08T12:00:00.000Z"; + const pingAck = await new Promise<{ ok: boolean }>((r) => + client.emit("realtime.ping", { at: sentAt }, r), + ); + expect(pingAck.ok).toBe(true); + + // Pong arrives synchronously after handler runs. + await new Promise((r) => setImmediate(r)); + + expect(pongs).toHaveLength(1); + expect(pongs[0]).toEqual({ at: sentAt, echo: "user_test" }); + + client.disconnect(); + }); + + it("anonymous client cannot subscribe to pong (authenticated scope)", async () => { + const client = ioClient(`http://localhost:${port}`); + await new Promise((r) => client.on("connect", () => r())); + + const subAck = await new Promise<{ ok: boolean; error?: string }>((r) => + client.emit("subscribe", "realtime.pong", r), + ); + + expect(subAck.ok).toBe(false); + expect(subAck.error).toBe("forbidden"); + client.disconnect(); + }); +}); diff --git a/packages/core-realtime/src/realtime-handler-registry.ts b/packages/core-realtime/src/realtime-handler-registry.ts index b0db24f..78e5839 100644 --- a/packages/core-realtime/src/realtime-handler-registry.ts +++ b/packages/core-realtime/src/realtime-handler-registry.ts @@ -1,17 +1,24 @@ import type { z } from "zod"; +import type { RealtimeChannelDescriptor } from "./realtime-channel"; import type { IInboundDescriptor } from "./realtime-handler.interface"; export interface IRealtimeHandlerRegistry { register(entry: IInboundDescriptor>): void; getInboundDescriptor(channelName: string): IInboundDescriptor | null; list(): IInboundDescriptor[]; + /** Register an outbound-only channel so Gate 2 can authorize subscriptions to it. */ + registerChannel(descriptor: RealtimeChannelDescriptor): void; + listChannels(): RealtimeChannelDescriptor[]; } export class RealtimeHandlerRegistry implements IRealtimeHandlerRegistry { private readonly entries = new Map>(); + private readonly channels = new Map>(); register(entry: IInboundDescriptor>): void { this.entries.set(entry.descriptor.name, entry as IInboundDescriptor); + // Also add the descriptor to the channel map so Gate 2 can authorize subscriptions. + this.channels.set(entry.descriptor.name, entry.descriptor); } getInboundDescriptor(channelName: string): IInboundDescriptor | null { @@ -21,4 +28,12 @@ export class RealtimeHandlerRegistry implements IRealtimeHandlerRegistry { list(): IInboundDescriptor[] { return Array.from(this.entries.values()); } + + registerChannel(descriptor: RealtimeChannelDescriptor): void { + this.channels.set(descriptor.name, descriptor); + } + + listChannels(): RealtimeChannelDescriptor[] { + return Array.from(this.channels.values()); + } } diff --git a/packages/core-realtime/src/socket-io-realtime-server.ts b/packages/core-realtime/src/socket-io-realtime-server.ts index 649efec..33c91be 100644 --- a/packages/core-realtime/src/socket-io-realtime-server.ts +++ b/packages/core-realtime/src/socket-io-realtime-server.ts @@ -55,11 +55,12 @@ export class SocketIORealtimeServer implements IRealtimeServer { // Gate 2: subscribe. socket.on("subscribe", async (requestedName: string, ack?: (r: unknown) => void) => { // Find a registered descriptor whose name (or template) matches requestedName. + // listChannels() covers both inbound descriptors and outbound-only channels. let matched: { descriptor: { name: string; scope: unknown }; params: Record } | null = null; - for (const entry of registry.list()) { - const m = matchChannelTemplate(entry.descriptor.name, requestedName); + for (const descriptor of registry.listChannels()) { + const m = matchChannelTemplate(descriptor.name, requestedName); if (m) { - matched = { descriptor: entry.descriptor, params: m.params }; + matched = { descriptor, params: m.params }; break; } } diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index b512d6b..eec3a8d 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -245,6 +245,9 @@ importers: jsdom: specifier: ^25.0.0 version: 25.0.1 + socket.io-client: + specifier: ^4.7.0 + version: 4.8.3 tsx: specifier: ^4.0.0 version: 4.21.0