httpapi-event.test.ts 3.5 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394
  1. import { afterEach, describe, expect } from "bun:test"
  2. import { Effect, Layer, Queue, Schema, Stream } from "effect"
  3. import { EventPaths } from "../../src/server/routes/instance/httpapi/groups/event"
  4. import { resetDatabase } from "../fixture/db"
  5. import { disposeAllInstances, TestInstance } from "../fixture/fixture"
  6. import { testEffect } from "../lib/effect"
  7. import { httpApiLayer, requestInDirectory } from "./httpapi-layer"
  8. const EventData = Schema.Struct({
  9. id: Schema.optional(Schema.String),
  10. type: Schema.String,
  11. properties: Schema.Record(Schema.String, Schema.Any),
  12. })
  13. const readEvent = (reader: Queue.Dequeue<Uint8Array>) =>
  14. Effect.gen(function* () {
  15. const value = yield* Queue.take(reader).pipe(
  16. Effect.timeoutOrElse({
  17. duration: "5 seconds",
  18. orElse: () => Effect.fail(new Error("timed out waiting for event")),
  19. }),
  20. )
  21. return Schema.decodeUnknownSync(EventData)(JSON.parse(new TextDecoder().decode(value).replace(/^data: /, "")))
  22. })
  23. const openEventStream = (directory: string) =>
  24. Effect.gen(function* () {
  25. const response = yield* requestInDirectory(EventPaths.event, directory)
  26. const reader = yield* Queue.unbounded<Uint8Array>()
  27. yield* response.stream.pipe(
  28. Stream.runForEach((value) => Queue.offer(reader, value)),
  29. Effect.forkScoped,
  30. )
  31. return { response, reader }
  32. })
  33. afterEach(async () => {
  34. await disposeAllInstances()
  35. await resetDatabase()
  36. })
  37. const it = testEffect(httpApiLayer)
  38. describe("event HttpApi", () => {
  39. it.instance(
  40. "serves event stream",
  41. () =>
  42. Effect.gen(function* () {
  43. const { directory } = yield* TestInstance
  44. const { response, reader } = yield* openEventStream(directory)
  45. expect(response.status).toBe(200)
  46. expect(response.headers["content-type"]).toContain("text/event-stream")
  47. expect(response.headers["cache-control"]).toBe("no-cache, no-transform")
  48. expect(response.headers["x-accel-buffering"]).toBe("no")
  49. expect(response.headers["x-content-type-options"]).toBe("nosniff")
  50. expect(yield* readEvent(reader)).toMatchObject({ type: "server.connected", properties: {} })
  51. }),
  52. { git: true, config: { formatter: false, lsp: false } },
  53. )
  54. it.instance(
  55. "keeps the event stream open after the initial event",
  56. () =>
  57. Effect.gen(function* () {
  58. const { directory } = yield* TestInstance
  59. const { reader } = yield* openEventStream(directory)
  60. expect(yield* readEvent(reader)).toMatchObject({ type: "server.connected", properties: {} })
  61. // If no second event arrives within 250ms, the stream is still open.
  62. const status = yield* Queue.take(reader).pipe(
  63. Effect.as("event" as const),
  64. Effect.timeoutOrElse({ duration: "250 millis", orElse: () => Effect.succeed("open" as const) }),
  65. )
  66. expect(status).toBe("open")
  67. }),
  68. { git: true, config: { formatter: false, lsp: false } },
  69. )
  70. it.instance(
  71. "delivers instance events after the initial event",
  72. () =>
  73. Effect.gen(function* () {
  74. const { directory } = yield* TestInstance
  75. const { reader } = yield* openEventStream(directory)
  76. expect(yield* readEvent(reader)).toMatchObject({ type: "server.connected", properties: {} })
  77. const created = yield* requestInDirectory("/session", directory, { method: "POST" })
  78. expect(created.status).toBe(200)
  79. expect(yield* readEvent(reader)).toMatchObject({ type: "session.created" })
  80. }),
  81. { git: true, config: { formatter: false, lsp: false } },
  82. )
  83. })