websocket-server.ts 2.2 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071
  1. import { Buffer } from "node:buffer"
  2. import { Effect } from "effect"
  3. interface ConnectionData {
  4. readonly id: number
  5. }
  6. export interface WebSocketServerState {
  7. readonly headers: Array<Record<string, string>>
  8. readonly messages: string[]
  9. opens: number
  10. closes: number
  11. pongs: number
  12. }
  13. export interface WebSocketServerFixture {
  14. readonly url: string
  15. readonly state: WebSocketServerState
  16. }
  17. export interface WebSocketServerOptions {
  18. readonly upgrade?: (request: Request) => boolean
  19. readonly open?: (socket: Bun.ServerWebSocket<ConnectionData>) => void
  20. readonly message?: (socket: Bun.ServerWebSocket<ConnectionData>, message: string | Buffer) => void
  21. }
  22. export const makeWebSocketServer = (options: WebSocketServerOptions = {}) =>
  23. Effect.acquireRelease(
  24. Effect.sync(() => {
  25. const state: WebSocketServerState = { headers: [], messages: [], opens: 0, closes: 0, pongs: 0 }
  26. let connection = 0
  27. const server = Bun.serve<ConnectionData>({
  28. hostname: "127.0.0.1",
  29. port: 0,
  30. fetch(request, server) {
  31. state.headers.push(Object.fromEntries(request.headers.entries()))
  32. if ((options.upgrade?.(request) ?? true) && server.upgrade(request, { data: { id: connection++ } }))
  33. return undefined
  34. return new Response("WebSocket upgrade required", {
  35. status: 426,
  36. headers: { "x-upgrade-rejected": "true" },
  37. })
  38. },
  39. websocket: {
  40. open(socket) {
  41. state.opens++
  42. options.open?.(socket)
  43. },
  44. message(socket, message) {
  45. const text = typeof message === "string" ? message : message.toString()
  46. state.messages.push(text)
  47. options.message?.(socket, message)
  48. },
  49. close() {
  50. state.closes++
  51. },
  52. pong() {
  53. state.pongs++
  54. },
  55. },
  56. })
  57. return {
  58. server,
  59. fixture: {
  60. url: `${server.url.toString().replace(/^http/, "ws")}responses`,
  61. state,
  62. } satisfies WebSocketServerFixture,
  63. }
  64. }),
  65. ({ server }) => Effect.promise(() => server.stop(true)),
  66. ).pipe(Effect.map((item) => item.fixture))