File
Blob: src/worker/durable-objects/inbox-ws.ts
| 1 | import { DurableObject } from "cloudflare:workers"; |
| 2 | import type { NewEmailNotification } from "@/shared/contracts"; |
| 3 | |
| 4 | export class InboxWebSocket extends DurableObject { |
| 5 | constructor(ctx: DurableObjectState, env: Env) { |
| 6 | super(ctx, env); |
| 7 | this.ctx.setWebSocketAutoResponse(new WebSocketRequestResponsePair("ping", "pong")); |
| 8 | } |
| 9 | |
| 10 | async notifyNewEmail(payload: NewEmailNotification): Promise<void> { |
| 11 | this.broadcast(JSON.stringify({ type: "new_email", ...payload })); |
| 12 | } |
| 13 | |
| 14 | async fetch(request: Request) { |
| 15 | if (request.headers.get("upgrade") === "websocket") { |
| 16 | const pair = new WebSocketPair(); |
| 17 | const client = pair[0]; |
| 18 | const server = pair[1]; |
| 19 | const url = new URL(request.url); |
| 20 | const address = url.searchParams.get("address") ?? "unknown"; |
| 21 | |
| 22 | this.ctx.acceptWebSocket(server, [address]); |
| 23 | return new Response(null, { |
| 24 | status: 101, |
| 25 | webSocket: client, |
| 26 | }); |
| 27 | } |
| 28 | |
| 29 | return new Response("Expected WebSocket upgrade", { status: 400 }); |
| 30 | } |
| 31 | |
| 32 | private broadcast(message: string) { |
| 33 | for (const socket of this.ctx.getWebSockets()) { |
| 34 | socket.send(message); |
| 35 | } |
| 36 | } |
| 37 | |
| 38 | async webSocketMessage(_ws: WebSocket, message: string | ArrayBuffer) { |
| 39 | if (typeof message === "string" && message === "ping") { |
| 40 | return; |
| 41 | } |
| 42 | } |
| 43 | |
| 44 | async webSocketClose(ws: WebSocket, code: number, reason: string, _wasClean: boolean) { |
| 45 | ws.close(code, reason); |
| 46 | } |
| 47 | |
| 48 | async webSocketError(ws: WebSocket) { |
| 49 | ws.close(1011, "WebSocket error"); |
| 50 | } |
| 51 | } |