control-server.test.ts 3.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101
  1. import { expect, test } from "bun:test"
  2. import { Effect, Queue, Schema } from "effect"
  3. import { SimulationControlServer } from "../src/control-server"
  4. import { availableEndpoint, connect } from "./fixture/websocket"
  5. const Request = Schema.Struct({ id: Schema.optional(Schema.Number) })
  6. test("awaits accepted socket cleanup before the server scope closes", async () => {
  7. const endpoint = availableEndpoint()
  8. let cleaned = false
  9. await Effect.runPromise(
  10. Effect.scoped(
  11. Effect.gen(function* () {
  12. yield* SimulationControlServer.start({
  13. endpoint,
  14. label: "control server test",
  15. data: () => ({}),
  16. decode: Schema.decodeUnknownEffect(Schema.fromJsonString(Request)),
  17. handle: () => Effect.succeed({ ok: true }),
  18. close: () =>
  19. Effect.promise(async () => {
  20. await Bun.sleep(25)
  21. cleaned = true
  22. }),
  23. })
  24. yield* connect(endpoint)
  25. }),
  26. ),
  27. )
  28. expect(cleaned).toBe(true)
  29. const url = new URL(endpoint)
  30. const rebound = Bun.serve({ hostname: url.hostname, port: Number(url.port), fetch: () => new Response() })
  31. await rebound.stop(true)
  32. })
  33. test("continues serving after a response targets a closed socket", async () => {
  34. const endpoint = availableEndpoint()
  35. await Effect.runPromise(
  36. Effect.scoped(
  37. Effect.gen(function* () {
  38. yield* SimulationControlServer.start({
  39. endpoint,
  40. label: "control server test",
  41. data: () => ({}),
  42. decode: Schema.decodeUnknownEffect(Schema.fromJsonString(Request)),
  43. handle: () => Effect.sleep(25).pipe(Effect.as({ ok: true })),
  44. })
  45. const closed = yield* connect(endpoint)
  46. closed.send(JSON.stringify({ id: 1 }))
  47. closed.close()
  48. yield* Effect.sleep(50)
  49. const socket = yield* connect(endpoint)
  50. const messages = yield* Queue.unbounded<unknown>()
  51. socket.addEventListener("message", (event) => Queue.offerUnsafe(messages, JSON.parse(String(event.data))))
  52. socket.send(JSON.stringify({ id: 2 }))
  53. expect(yield* Queue.take(messages)).toMatchObject({ id: 2, result: { ok: true } })
  54. }),
  55. ),
  56. )
  57. })
  58. test("disconnects and cleans up when an outbound message exceeds the queue bound", async () => {
  59. const endpoint = availableEndpoint()
  60. let cleaned = false
  61. let delivered = false
  62. await Effect.runPromise(
  63. Effect.scoped(
  64. Effect.gen(function* () {
  65. yield* SimulationControlServer.start({
  66. endpoint,
  67. label: "control server test",
  68. data: () => ({}),
  69. decode: Schema.decodeUnknownEffect(Schema.fromJsonString(Request)),
  70. handle: (socket) =>
  71. Effect.gen(function* () {
  72. yield* socket.send("x".repeat(64 * 1024 * 1024 + 1))
  73. delivered = true
  74. return { ok: true }
  75. }),
  76. close: () => Effect.sync(() => void (cleaned = true)),
  77. })
  78. const socket = yield* connect(endpoint)
  79. const closed = yield* Queue.unbounded<void>()
  80. const messages: string[] = []
  81. socket.addEventListener("close", () => Queue.offerUnsafe(closed, undefined))
  82. socket.addEventListener("message", (event) => messages.push(String(event.data)))
  83. socket.send(JSON.stringify({ id: 1 }))
  84. yield* Queue.take(closed).pipe(Effect.timeout("5 seconds"))
  85. expect(cleaned).toBe(true)
  86. expect(delivered).toBe(false)
  87. expect(messages).toEqual([])
  88. }),
  89. ),
  90. )
  91. })