import type { Page } from "@playwright/test" export type SseConnectionRecord = { id: number url: string path: "/global/event" | "/event" | "/api/event" headers: Record openedAt: number endedAt?: number endedBy?: "close" | "disconnect" | "error" | "abort" error?: string } export type SseDeliveryAcknowledgement = { deliveryID: number connectionID: number bytes: number chunkCount: number deliveredAt: number eventID?: string } export type SseEventOptions = { id?: string event?: string retry?: number marker?: string } export type SseTransport = { server: string waitForConnection(options?: { after?: number; timeout?: number }): Promise send(payload: T, options?: SseEventOptions): Promise burst(payloads: readonly T[], options?: readonly SseEventOptions[]): Promise split(payload: T, cuts: readonly number[], options?: SseEventOptions): Promise heartbeat(options?: SseEventOptions): Promise writeRaw(value: string | Uint8Array, cuts?: readonly number[], marker?: string): Promise close(): Promise disconnect(message?: string): Promise error(message?: string): Promise connections(): Promise acknowledgements(): Promise } type BrowserCommand = | { type: "send"; deliveries: { payload: T; options?: SseEventOptions }[]; burst: boolean; cuts?: number[] } | { type: "raw"; bytes: number[]; cuts?: number[]; marker?: string } | { type: "end"; mode: "close" | "disconnect" | "error"; message?: string } | { type: "connections" } | { type: "acknowledgements" } type BrowserTransport = Window & { __testSseTransport?: { command: (command: BrowserCommand) => unknown } } export async function installSseTransport( page: Page, options: { server: string; retry?: number }, ): Promise> { const server = new URL(options.server).origin await page.addInitScript( ({ server, retry }) => { type Connection = SseConnectionRecord & { controller: ReadableStreamDefaultController } type ProbeWindow = Window & { __visualStabilityProbe?: { startedAt: number; markers: { at: number; label: string }[] } } const originalFetch = window.fetch.bind(window) const connections: Connection[] = [] const acknowledgements: SseDeliveryAcknowledgement[] = [] const encoder = new TextEncoder() let nextConnectionID = 0 let nextDeliveryID = 0 const current = () => connections.findLast((connection) => connection.endedAt === undefined) const chunks = (bytes: Uint8Array, cuts?: readonly number[]) => { const boundaries = [...new Set(cuts ?? [])] .filter((cut) => Number.isInteger(cut) && cut > 0 && cut < bytes.byteLength) .sort((a, b) => a - b) return [0, ...boundaries].map((start, index) => bytes.slice(start, boundaries[index] ?? bytes.byteLength)) } const marker = (label?: string) => { if (!label) return const probe = (window as ProbeWindow).__visualStabilityProbe if (!probe) return probe.markers.push({ at: performance.now() - probe.startedAt, label }) } const frame = (payload: unknown, eventOptions: SseEventOptions = {}) => [ eventOptions.event === undefined ? "" : `event: ${eventOptions.event}\n`, eventOptions.id === undefined ? "" : `id: ${eventOptions.id}\n`, eventOptions.retry === undefined ? "" : `retry: ${eventOptions.retry}\n`, `data: ${JSON.stringify(payload)}\n\n`, ].join("") const currentEvent = (input: unknown) => { if (!input || typeof input !== "object" || !("payload" in input)) return input const envelope = input as { directory?: string; payload?: unknown } if (!envelope.payload || typeof envelope.payload !== "object") return input const payload = envelope.payload as { id?: string; type?: string; properties?: unknown } if (!payload.type) return input return { id: payload.id ?? `evt_mock_${Date.now()}`, created: Date.now(), type: payload.type, data: payload.properties ?? {}, location: envelope.directory && envelope.directory !== "global" ? { directory: envelope.directory } : undefined, } } const acknowledge = ( connection: Connection, bytes: number, chunkCount: number, eventID?: string, ): SseDeliveryAcknowledgement => { const acknowledgement = { deliveryID: ++nextDeliveryID, connectionID: connection.id, bytes, chunkCount, deliveredAt: performance.now(), ...(eventID === undefined ? {} : { eventID }), } acknowledgements.push(acknowledgement) return acknowledgement } const end = (mode: "close" | "disconnect" | "error", message?: string) => { const connection = current() if (!connection) throw new Error("SSE transport has no active connection") connection.endedAt = performance.now() connection.endedBy = mode if (message) connection.error = message if (mode === "close") { connection.controller.close() return } const error = new DOMException( message ?? "SSE connection disconnected", mode === "error" ? "Error" : "NetworkError", ) connection.controller.error(error) } const command = (input: BrowserCommand) => { if (input.type === "connections") return connections.map(({ controller: _controller, ...connection }) => connection) if (input.type === "acknowledgements") return acknowledgements if (input.type === "end") return end(input.mode, input.message) const connection = current() if (!connection) throw new Error("SSE transport has no active connection") if (input.type === "raw") { marker(input.marker) const output = chunks(new Uint8Array(input.bytes), input.cuts) output.forEach((chunk) => connection.controller.enqueue(chunk)) return acknowledge(connection, input.bytes.length, output.length) } const encoded = input.deliveries.map((delivery) => { const payload = connection.path === "/api/event" ? currentEvent(delivery.payload) : delivery.payload return { delivery, payload, bytes: encoder.encode(frame(payload, delivery.options)) } }) encoded.forEach((item) => marker(item.delivery.options?.marker)) if (input.burst) { const bytes = encoder.encode(encoded.map((item) => frame(item.payload, item.delivery.options)).join("")) connection.controller.enqueue(bytes) return encoded.map((item) => acknowledge(connection, item.bytes.byteLength, 1, item.delivery.options?.id)) } const output = chunks(encoded[0]!.bytes, input.cuts) output.forEach((chunk) => connection.controller.enqueue(chunk)) return acknowledge(connection, encoded[0]!.bytes.byteLength, output.length, encoded[0]!.delivery.options?.id) } ;(window as BrowserTransport).__testSseTransport = { command } const fetch = (input: RequestInfo | URL, init?: RequestInit) => { const request = new Request(input, init) const url = new URL(request.url) if ( url.origin !== server || (url.pathname !== "/global/event" && url.pathname !== "/event" && url.pathname !== "/api/event") ) return originalFetch(request) const id = ++nextConnectionID const record = { id, url: url.href, path: url.pathname, headers: Object.fromEntries(request.headers.entries()), openedAt: performance.now(), } as Connection const stream = new ReadableStream({ start(controller) { record.controller = controller connections.push(record) if (retry !== undefined) controller.enqueue(encoder.encode(`retry: ${retry}\n\n`)) if (url.pathname === "/api/event") controller.enqueue( encoder.encode(frame({ id: `evt_mock_connected_${id}`, type: "server.connected", data: {} })), ) if (url.pathname === "/global/event") controller.enqueue( encoder.encode( frame({ payload: { id: `evt_mock_connected_${id}`, type: "server.connected", properties: {} }, }), ), ) request.signal.addEventListener( "abort", () => { if (record.endedAt !== undefined) return record.endedAt = performance.now() record.endedBy = "abort" controller.error(request.signal.reason ?? new DOMException("The operation was aborted", "AbortError")) }, { once: true }, ) }, cancel() { if (record.endedAt !== undefined) return record.endedAt = performance.now() record.endedBy = "disconnect" }, }) return Promise.resolve( new Response(stream, { status: 200, headers: { "cache-control": "no-cache", "content-type": "text/event-stream", }, }), ) } Object.defineProperty(window, "fetch", { configurable: true, writable: true, value: fetch }) }, { server, retry: options.retry }, ) const command = (input: BrowserCommand) => page.evaluate((input) => { const transport = (window as BrowserTransport).__testSseTransport if (!transport) throw new Error("SSE transport was not installed before page load") return transport.command(input as BrowserCommand) }, input) as Promise return { server, async waitForConnection(input = {}) { await page.waitForFunction( (after) => { const transport = (window as BrowserTransport).__testSseTransport const connections = transport?.command({ type: "connections" }) as SseConnectionRecord[] | undefined return connections?.some((connection) => connection.id > after) }, input.after ?? 0, { timeout: input.timeout }, ) return (await command({ type: "connections" })).findLast( (connection) => connection.id > (input.after ?? 0), )! }, send(payload, eventOptions) { return command({ type: "send", deliveries: [{ payload, options: eventOptions }], burst: false }) }, burst(payloads, eventOptions = []) { return command({ type: "send", deliveries: payloads.map((payload, index) => ({ payload, options: eventOptions[index] })), burst: true, }) }, split(payload, cuts, eventOptions) { return command({ type: "send", deliveries: [{ payload, options: eventOptions }], burst: false, cuts: [...cuts] }) }, heartbeat(eventOptions) { return command({ type: "send", deliveries: [ { payload: { directory: "global", payload: { type: "server.heartbeat", properties: {} } } as T, options: eventOptions, }, ], burst: false, }) }, writeRaw(value, cuts, marker) { return command({ type: "raw", bytes: Array.from(typeof value === "string" ? new TextEncoder().encode(value) : value), cuts: cuts ? [...cuts] : undefined, marker, }) }, close() { return command({ type: "end", mode: "close" }) }, disconnect(message) { return command({ type: "end", mode: "disconnect", message }) }, error(message) { return command({ type: "end", mode: "error", message }) }, connections() { return command({ type: "connections" }) }, acknowledgements() { return command({ type: "acknowledgements" }) }, } }