websocket.test.ts 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529
  1. import { describe, expect, test } from "bun:test"
  2. import { Deferred, Effect, Exit, Fiber, Layer } from "effect"
  3. import { Socket } from "effect/unstable/socket"
  4. import { existsSync } from "node:fs"
  5. import { HttpRecorder } from "../src"
  6. import { layerSocketWithMode } from "../src/websocket/recorder"
  7. import { failureText, readCassette, seedCassetteDirectory, tempDirectory, withEnvironment } from "./support"
  8. const unavailableSocket = Socket.make({
  9. runRaw: () => Effect.die(new Error("unexpected live WebSocket run")),
  10. writer: Effect.succeed(() => Effect.die(new Error("unexpected live WebSocket write"))),
  11. })
  12. class EchoWebSocket extends EventTarget {
  13. readonly protocol = ""
  14. readonly extensions = ""
  15. bufferedAmount = 0
  16. binaryType: BinaryType = "blob"
  17. readyState = 0
  18. constructor(readonly url: string) {
  19. super()
  20. queueMicrotask(() => {
  21. this.readyState = 1
  22. this.dispatchEvent(new Event("open"))
  23. })
  24. }
  25. send(data: string | ArrayBufferLike | Blob | ArrayBufferView) {
  26. queueMicrotask(() => this.dispatchEvent(new MessageEvent("message", { data })))
  27. }
  28. close(code = 1000, reason = "") {
  29. if (this.readyState === 3) return
  30. this.readyState = 3
  31. this.dispatchEvent(new CloseEvent("close", { code, reason, wasClean: code === 1000 }))
  32. }
  33. }
  34. describe("WebSocket", () => {
  35. test("constructor recording is complete when the recorder layer closes", async () => {
  36. using directory = tempDirectory("http-recorder-websocket-constructor-")
  37. const recorder = HttpRecorder.layerWebSocketConstructor("websocket/constructor-record", {
  38. directory: directory.path,
  39. }).pipe(
  40. Layer.provide(
  41. Layer.succeed(Socket.WebSocketConstructor, (url) => new EchoWebSocket(url) as unknown as globalThis.WebSocket),
  42. ),
  43. )
  44. await withEnvironment("CI", undefined, () =>
  45. Effect.runPromise(
  46. Effect.gen(function* () {
  47. const socket = yield* Socket.makeWebSocket("wss://echo.example.test/one", {
  48. protocols: ["echo.v1"],
  49. closeCodeIsError: () => false,
  50. })
  51. const write = yield* socket.writer
  52. yield* socket.runString(() => write(new Socket.CloseEvent(1000, "complete")).pipe(Effect.orDie), {
  53. onOpen: write("hello").pipe(Effect.orDie),
  54. })
  55. }).pipe(Effect.scoped, Effect.provide(recorder)),
  56. ),
  57. )
  58. expect(readCassette(`${directory.path}/websocket/constructor-record.json`).interactions).toEqual([
  59. {
  60. transport: "websocket",
  61. connection: {
  62. sequence: 0,
  63. url: "wss://echo.example.test/one",
  64. protocols: ["echo.v1"],
  65. close: { code: 1000, reason: "complete" },
  66. },
  67. events: [
  68. { direction: "client", kind: "text", body: "hello" },
  69. { direction: "server", kind: "text", body: "hello" },
  70. ],
  71. },
  72. ])
  73. })
  74. test("constructor replay validates dynamic URLs and protocols without opening a live socket", async () => {
  75. using directory = tempDirectory("http-recorder-websocket-constructor-")
  76. await seedCassetteDirectory(directory.path, "websocket/constructor", [
  77. {
  78. transport: "websocket",
  79. connection: {
  80. sequence: 0,
  81. url: "wss://events.example.test/workspaces/one",
  82. protocols: ["events.v1"],
  83. close: { code: 1000, reason: "complete" },
  84. },
  85. events: [
  86. { direction: "client", kind: "text", body: '{"type":"subscribe"}' },
  87. { direction: "server", kind: "text", body: '{"type":"ready"}' },
  88. ],
  89. },
  90. ])
  91. const unavailableConstructor = () => {
  92. throw new Error("unexpected live WebSocket construction")
  93. }
  94. const recorder = HttpRecorder.layerWebSocketConstructor("websocket/constructor", {
  95. directory: directory.path,
  96. }).pipe(Layer.provide(Layer.succeed(Socket.WebSocketConstructor, unavailableConstructor)))
  97. const received = await Effect.runPromise(
  98. Effect.gen(function* () {
  99. const socket = yield* Socket.makeWebSocket("wss://events.example.test/workspaces/one", {
  100. protocols: ["events.v1"],
  101. closeCodeIsError: () => false,
  102. })
  103. const write = yield* socket.writer
  104. const received: string[] = []
  105. yield* socket.runString(
  106. (message) => {
  107. received.push(message)
  108. },
  109. {
  110. onOpen: write('{"type":"subscribe"}').pipe(Effect.orDie),
  111. },
  112. )
  113. return received
  114. }).pipe(Effect.scoped, Effect.provide(recorder)),
  115. )
  116. expect(received).toEqual(['{"type":"ready"}'])
  117. })
  118. test("constructor replay rejects a different dynamic URL", async () => {
  119. using directory = tempDirectory("http-recorder-websocket-constructor-")
  120. await seedCassetteDirectory(directory.path, "websocket/constructor-mismatch", [
  121. {
  122. transport: "websocket",
  123. connection: {
  124. sequence: 0,
  125. url: "wss://events.example.test/workspaces/one",
  126. protocols: [],
  127. close: { code: 1000, reason: "complete" },
  128. },
  129. events: [],
  130. },
  131. ])
  132. const recorder = HttpRecorder.layerWebSocketConstructor("websocket/constructor-mismatch", {
  133. directory: directory.path,
  134. }).pipe(
  135. Layer.provide(
  136. Layer.succeed(Socket.WebSocketConstructor, () => {
  137. throw new Error("unexpected live WebSocket construction")
  138. }),
  139. ),
  140. )
  141. const exit = await Effect.runPromise(
  142. Effect.gen(function* () {
  143. const socket = yield* Socket.makeWebSocket("wss://events.example.test/workspaces/two")
  144. yield* socket.runString(() => {})
  145. }).pipe(Effect.scoped, Effect.exit, Effect.provide(recorder)),
  146. )
  147. expect(Exit.isFailure(exit)).toBe(true)
  148. })
  149. test("records WebSocket frames in observed client/server order", async () => {
  150. using directory = tempDirectory("http-recorder-websocket-")
  151. const response = JSON.stringify({
  152. type: "response.completed",
  153. token: "server-secret",
  154. })
  155. const upstream = Socket.make({
  156. runRaw: (handler, options) =>
  157. Effect.gen(function* () {
  158. if (options?.onOpen) yield* options.onOpen
  159. const result = handler(response)
  160. if (Effect.isEffect(result)) yield* result
  161. }),
  162. writer: Effect.succeed(() => Effect.void),
  163. })
  164. await Effect.runPromise(
  165. Effect.gen(function* () {
  166. const socket = yield* Socket.Socket
  167. const write = yield* socket.writer
  168. yield* socket.runRaw(() => {}, {
  169. onOpen: write(JSON.stringify({ type: "response.create", token: "client-secret" })).pipe(Effect.orDie),
  170. })
  171. }).pipe(
  172. Effect.scoped,
  173. Effect.provide(
  174. layerSocketWithMode("websocket/record", {
  175. directory: directory.path,
  176. metadata: { provider: "test" },
  177. mode: "record",
  178. }).pipe(Layer.provide(Layer.succeed(Socket.Socket, upstream))),
  179. ),
  180. ),
  181. )
  182. expect(readCassette(`${directory.path}/websocket/record.json`)).toMatchObject({
  183. interactions: [
  184. {
  185. transport: "websocket",
  186. events: [
  187. {
  188. direction: "client",
  189. kind: "text",
  190. body: '{"type":"response.create","token":"[REDACTED]"}',
  191. },
  192. {
  193. direction: "server",
  194. kind: "text",
  195. body: '{"type":"response.completed","token":"[REDACTED]"}',
  196. },
  197. ],
  198. },
  199. ],
  200. })
  201. })
  202. test("WebSocket replay preserves causal frame ordering", async () => {
  203. using directory = tempDirectory("http-recorder-websocket-")
  204. await seedCassetteDirectory(directory.path, "websocket/replay", [
  205. {
  206. transport: "websocket",
  207. events: [
  208. {
  209. direction: "server",
  210. kind: "text",
  211. body: '{"type":"session.created"}',
  212. },
  213. {
  214. direction: "client",
  215. kind: "text",
  216. body: '{"type":"response.create","prompt":"hello"}',
  217. },
  218. {
  219. direction: "server",
  220. kind: "text",
  221. body: '{"type":"response.completed"}',
  222. },
  223. ],
  224. },
  225. ])
  226. const received: string[] = []
  227. await Effect.runPromise(
  228. Effect.gen(function* () {
  229. const socket = yield* Socket.Socket
  230. const write = yield* socket.writer
  231. yield* socket.runRaw((message) =>
  232. Effect.gen(function* () {
  233. if (typeof message !== "string") return
  234. received.push(message)
  235. const event: unknown = JSON.parse(message)
  236. if (typeof event !== "object" || event === null || !("type" in event)) return
  237. if (event.type === "session.created") yield* write('{"prompt":"hello","type":"response.create"}')
  238. }),
  239. )
  240. }).pipe(
  241. Effect.scoped,
  242. Effect.provide(
  243. layerSocketWithMode("websocket/replay", {
  244. directory: directory.path,
  245. compareClientMessagesAsJson: true,
  246. mode: "replay",
  247. }).pipe(Layer.provide(Layer.succeed(Socket.Socket, unavailableSocket))),
  248. ),
  249. ),
  250. )
  251. expect(received).toEqual(['{"type":"session.created"}', '{"type":"response.completed"}'])
  252. })
  253. test("the public socket decorator replays a causal provider conversation", async () => {
  254. using directory = tempDirectory("http-recorder-websocket-")
  255. await seedCassetteDirectory(directory.path, "websocket/public-layer", [
  256. {
  257. transport: "websocket",
  258. events: [
  259. {
  260. direction: "server",
  261. kind: "text",
  262. body: '{"type":"session.created"}',
  263. },
  264. {
  265. direction: "client",
  266. kind: "text",
  267. body: '{"type":"response.create","prompt":"first"}',
  268. },
  269. {
  270. direction: "server",
  271. kind: "text",
  272. body: '{"type":"response.completed","id":"first"}',
  273. },
  274. {
  275. direction: "client",
  276. kind: "text",
  277. body: '{"type":"response.create","prompt":"second"}',
  278. },
  279. {
  280. direction: "server",
  281. kind: "text",
  282. body: '{"type":"response.completed","id":"second"}',
  283. },
  284. ],
  285. },
  286. ])
  287. const received: string[] = []
  288. await Effect.runPromise(
  289. Effect.gen(function* () {
  290. const socket = yield* Socket.Socket
  291. const write = yield* socket.writer
  292. yield* socket.runString((message) =>
  293. Effect.gen(function* () {
  294. received.push(message)
  295. const event: unknown = JSON.parse(message)
  296. if (typeof event !== "object" || event === null) return
  297. if ("type" in event && event.type === "session.created") {
  298. yield* write('{"prompt":"first","type":"response.create"}')
  299. return
  300. }
  301. if ("id" in event && event.id === "first") {
  302. yield* write('{"prompt":"second","type":"response.create"}')
  303. return
  304. }
  305. yield* write(new Socket.CloseEvent(1000, "done"))
  306. }),
  307. )
  308. }).pipe(
  309. Effect.scoped,
  310. Effect.provide(
  311. HttpRecorder.layerSocket("websocket/public-layer", { directory: directory.path }).pipe(
  312. Layer.provide(Layer.succeed(Socket.Socket, unavailableSocket)),
  313. ),
  314. ),
  315. ),
  316. )
  317. expect(received).toEqual([
  318. '{"type":"session.created"}',
  319. '{"type":"response.completed","id":"first"}',
  320. '{"type":"response.completed","id":"second"}',
  321. ])
  322. })
  323. test("WebSocket replay runs message handlers concurrently", async () => {
  324. using directory = tempDirectory("http-recorder-websocket-")
  325. await seedCassetteDirectory(directory.path, "websocket/concurrent-handlers", [
  326. {
  327. transport: "websocket",
  328. events: [
  329. { direction: "server", kind: "text", body: "first" },
  330. { direction: "server", kind: "text", body: "second" },
  331. ],
  332. },
  333. ])
  334. await Effect.runPromise(
  335. Effect.gen(function* () {
  336. const socket = yield* Socket.Socket
  337. const second = yield* Deferred.make<void>()
  338. yield* socket.runString((message) =>
  339. message === "first" ? Deferred.await(second) : Deferred.succeed(second, undefined),
  340. )
  341. }).pipe(
  342. Effect.scoped,
  343. Effect.provide(
  344. layerSocketWithMode("websocket/concurrent-handlers", { directory: directory.path, mode: "replay" }).pipe(
  345. Layer.provide(Layer.succeed(Socket.Socket, unavailableSocket)),
  346. ),
  347. ),
  348. ),
  349. )
  350. })
  351. test("rejected concurrent replay does not consume the next interaction", async () => {
  352. using directory = tempDirectory("http-recorder-websocket-")
  353. await seedCassetteDirectory(directory.path, "websocket/concurrent-runs", [
  354. { transport: "websocket", events: [{ direction: "server", kind: "text", body: "first" }] },
  355. { transport: "websocket", events: [{ direction: "server", kind: "text", body: "second" }] },
  356. ])
  357. const received: string[] = []
  358. await Effect.runPromise(
  359. Effect.gen(function* () {
  360. const socket = yield* Socket.Socket
  361. const started = yield* Deferred.make<void>()
  362. const release = yield* Deferred.make<void>()
  363. const first = yield* socket
  364. .runString((message) =>
  365. Effect.gen(function* () {
  366. received.push(message)
  367. yield* Deferred.succeed(started, undefined)
  368. yield* Deferred.await(release)
  369. }),
  370. )
  371. .pipe(Effect.forkChild)
  372. yield* Deferred.await(started)
  373. const concurrent = yield* Effect.exit(socket.runString(() => Effect.void))
  374. expect(failureText(concurrent)).toContain("Concurrent runs")
  375. yield* Deferred.succeed(release, undefined)
  376. yield* Fiber.join(first)
  377. yield* socket.runString((message) => Effect.sync(() => received.push(message)))
  378. }).pipe(
  379. Effect.scoped,
  380. Effect.provide(
  381. layerSocketWithMode("websocket/concurrent-runs", { directory: directory.path, mode: "replay" }).pipe(
  382. Layer.provide(Layer.succeed(Socket.Socket, unavailableSocket)),
  383. ),
  384. ),
  385. ),
  386. )
  387. expect(received).toEqual(["first", "second"])
  388. })
  389. test("WebSocket replay rejects close with unconsumed events", async () => {
  390. using directory = tempDirectory("http-recorder-websocket-")
  391. await seedCassetteDirectory(directory.path, "websocket/early-close", [
  392. {
  393. transport: "websocket",
  394. events: [{ direction: "client", kind: "text", body: "expected" }],
  395. },
  396. ])
  397. const exit = await Effect.runPromise(
  398. Effect.gen(function* () {
  399. const socket = yield* Socket.Socket
  400. const write = yield* socket.writer
  401. return yield* Effect.exit(
  402. socket.runRaw(() => {}, {
  403. onOpen: write(new Socket.CloseEvent(1000)).pipe(Effect.orDie),
  404. }),
  405. )
  406. }).pipe(
  407. Effect.scoped,
  408. Effect.provide(
  409. layerSocketWithMode("websocket/early-close", { directory: directory.path, mode: "replay" }).pipe(
  410. Layer.provide(Layer.succeed(Socket.Socket, unavailableSocket)),
  411. ),
  412. ),
  413. ),
  414. )
  415. expect(failureText(exit)).toContain("closed with unconsumed events")
  416. })
  417. test("failed WebSocket runs do not write complete cassettes", async () => {
  418. using directory = tempDirectory("http-recorder-websocket-")
  419. const exit = await Effect.runPromise(
  420. Effect.gen(function* () {
  421. const socket = yield* Socket.Socket
  422. return yield* Effect.exit(socket.runRaw(() => {}))
  423. }).pipe(
  424. Effect.scoped,
  425. Effect.provide(
  426. layerSocketWithMode("websocket/failed-run", { directory: directory.path, mode: "record" }).pipe(
  427. Layer.provide(
  428. Layer.succeed(
  429. Socket.Socket,
  430. Socket.make({
  431. runRaw: () => Effect.die(new Error("connection failed")),
  432. writer: Effect.succeed(() => Effect.void),
  433. }),
  434. ),
  435. ),
  436. ),
  437. ),
  438. ),
  439. )
  440. expect(Exit.isFailure(exit)).toBe(true)
  441. expect(existsSync(`${directory.path}/websocket/failed-run.json`)).toBe(false)
  442. })
  443. test("WebSocket replay preserves binary frame kinds across reconnects", async () => {
  444. using directory = tempDirectory("http-recorder-websocket-")
  445. const interaction = {
  446. transport: "websocket" as const,
  447. events: [
  448. {
  449. direction: "client" as const,
  450. kind: "binary" as const,
  451. body: Buffer.from([1, 2]).toString("base64"),
  452. bodyEncoding: "base64" as const,
  453. },
  454. {
  455. direction: "server" as const,
  456. kind: "binary" as const,
  457. body: Buffer.from([3, 4]).toString("base64"),
  458. bodyEncoding: "base64" as const,
  459. },
  460. ],
  461. }
  462. await seedCassetteDirectory(directory.path, "websocket/binary", [interaction, interaction])
  463. const received: number[][] = []
  464. await Effect.runPromise(
  465. Effect.gen(function* () {
  466. const socket = yield* Socket.Socket
  467. const write = yield* socket.writer
  468. const run = socket.runRaw(
  469. (message) => {
  470. if (typeof message === "string") throw new Error("Expected a binary WebSocket frame")
  471. received.push([...message])
  472. },
  473. { onOpen: write(new Uint8Array([1, 2])).pipe(Effect.orDie) },
  474. )
  475. yield* run
  476. yield* run
  477. }).pipe(
  478. Effect.scoped,
  479. Effect.provide(
  480. layerSocketWithMode("websocket/binary", { directory: directory.path, mode: "replay" }).pipe(
  481. Layer.provide(Layer.succeed(Socket.Socket, unavailableSocket)),
  482. ),
  483. ),
  484. ),
  485. )
  486. expect(received).toEqual([
  487. [3, 4],
  488. [3, 4],
  489. ])
  490. })
  491. })