effect.ts 8.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211
  1. import { NodeFileSystem } from "@effect/platform-node"
  2. import { Effect, Layer, Option, Ref } from "effect"
  3. import {
  4. FetchHttpClient,
  5. HttpClient,
  6. HttpClientError,
  7. HttpClientRequest,
  8. HttpClientResponse,
  9. } from "effect/unstable/http"
  10. import { redactedErrorRequest, mismatchDetail, requestDiff } from "./diff"
  11. import { defaultMatcher, decodeJson, type RequestMatcher } from "./matching"
  12. import { redactHeaders, redactUrl, type SecretFinding } from "./redaction"
  13. import {
  14. httpInteractions,
  15. type Cassette,
  16. type CassetteMetadata,
  17. type HttpInteraction,
  18. type ResponseSnapshot,
  19. } from "./schema"
  20. import * as CassetteService from "./cassette"
  21. export const DEFAULT_REQUEST_HEADERS: ReadonlyArray<string> = ["content-type", "accept", "openai-beta"]
  22. const DEFAULT_RESPONSE_HEADERS: ReadonlyArray<string> = ["content-type"]
  23. export type RecordReplayMode = "record" | "replay" | "passthrough"
  24. export interface RecordReplayOptions {
  25. readonly mode?: RecordReplayMode
  26. readonly directory?: string
  27. readonly metadata?: CassetteMetadata
  28. readonly redact?: {
  29. readonly headers?: ReadonlyArray<string>
  30. readonly query?: ReadonlyArray<string>
  31. readonly url?: (url: string) => string
  32. }
  33. readonly requestHeaders?: ReadonlyArray<string>
  34. readonly responseHeaders?: ReadonlyArray<string>
  35. readonly redactBody?: (body: unknown) => unknown
  36. readonly dispatch?: "match" | "sequential"
  37. readonly match?: RequestMatcher
  38. }
  39. const responseHeaders = (
  40. response: HttpClientResponse.HttpClientResponse,
  41. allow: ReadonlyArray<string>,
  42. redact: ReadonlyArray<string> | undefined,
  43. ) => {
  44. const merged = redactHeaders(response.headers as Record<string, string>, allow, redact)
  45. if (!merged["content-type"]) merged["content-type"] = "text/event-stream"
  46. return merged
  47. }
  48. const BINARY_CONTENT_TYPES: ReadonlyArray<string> = ["vnd.amazon.eventstream", "octet-stream"]
  49. const isBinaryContentType = (contentType: string | undefined) => {
  50. if (!contentType) return false
  51. const lower = contentType.toLowerCase()
  52. return BINARY_CONTENT_TYPES.some((token) => lower.includes(token))
  53. }
  54. const captureResponseBody = (response: HttpClientResponse.HttpClientResponse, contentType: string | undefined) =>
  55. isBinaryContentType(contentType)
  56. ? response.arrayBuffer.pipe(
  57. Effect.map((bytes) => ({ body: Buffer.from(bytes).toString("base64"), bodyEncoding: "base64" as const })),
  58. )
  59. : response.text.pipe(Effect.map((body) => ({ body })))
  60. const decodeResponseBody = (snapshot: ResponseSnapshot) =>
  61. snapshot.bodyEncoding === "base64" ? Buffer.from(snapshot.body, "base64") : snapshot.body
  62. const fixtureMissing = (request: HttpClientRequest.HttpClientRequest, name: string) =>
  63. new HttpClientError.HttpClientError({
  64. reason: new HttpClientError.TransportError({
  65. request: redactedErrorRequest(request),
  66. description: `Fixture "${name}" not found. Run with RECORD=true to create it.`,
  67. }),
  68. })
  69. const fixtureMismatch = (request: HttpClientRequest.HttpClientRequest, name: string, detail: string) =>
  70. new HttpClientError.HttpClientError({
  71. reason: new HttpClientError.TransportError({
  72. request: redactedErrorRequest(request),
  73. description: `Fixture "${name}" does not match the current request: ${detail}. Run with RECORD=true to update it.`,
  74. }),
  75. })
  76. const unsafeCassette = (
  77. request: HttpClientRequest.HttpClientRequest,
  78. name: string,
  79. findings: ReadonlyArray<SecretFinding>,
  80. ) =>
  81. new HttpClientError.HttpClientError({
  82. reason: new HttpClientError.TransportError({
  83. request: redactedErrorRequest(request),
  84. description: `Refusing to write cassette "${name}" because it contains possible secrets: ${findings
  85. .map((item) => `${item.path} (${item.reason})`)
  86. .join(", ")}`,
  87. }),
  88. })
  89. export const recordingLayer = (
  90. name: string,
  91. options: Omit<RecordReplayOptions, "directory"> = {},
  92. ): Layer.Layer<HttpClient.HttpClient, never, HttpClient.HttpClient | CassetteService.Service> =>
  93. Layer.effect(
  94. HttpClient.HttpClient,
  95. Effect.gen(function* () {
  96. const upstream = yield* HttpClient.HttpClient
  97. const cassetteService = yield* CassetteService.Service
  98. const requestHeadersAllow = options.requestHeaders ?? DEFAULT_REQUEST_HEADERS
  99. const responseHeadersAllow = options.responseHeaders ?? DEFAULT_RESPONSE_HEADERS
  100. const match = options.match ?? defaultMatcher
  101. const mode = options.mode ?? "replay"
  102. const sequential = options.dispatch === "sequential"
  103. const replay = yield* Ref.make<Cassette | undefined>(undefined)
  104. const cursor = yield* Ref.make(0)
  105. const snapshotRequest = (request: HttpClientRequest.HttpClientRequest) =>
  106. Effect.gen(function* () {
  107. const web = yield* HttpClientRequest.toWeb(request).pipe(Effect.orDie)
  108. const raw = yield* Effect.promise(() => web.text())
  109. const body = options.redactBody
  110. ? Option.match(decodeJson(raw), {
  111. onNone: () => raw,
  112. onSome: (parsed) => JSON.stringify(options.redactBody?.(parsed)),
  113. })
  114. : raw
  115. return {
  116. method: web.method,
  117. url: redactUrl(web.url, options.redact?.query, options.redact?.url),
  118. headers: redactHeaders(
  119. Object.fromEntries(web.headers.entries()),
  120. requestHeadersAllow,
  121. options.redact?.headers,
  122. ),
  123. body,
  124. }
  125. })
  126. const selectInteraction = (cassette: Cassette, incoming: HttpInteraction["request"]) =>
  127. Effect.gen(function* () {
  128. const interactions = httpInteractions(cassette)
  129. if (sequential) {
  130. const index = yield* Ref.get(cursor)
  131. const interaction = interactions[index]
  132. if (!interaction)
  133. return { interaction, detail: `interaction ${index + 1} of ${interactions.length} not recorded` }
  134. if (!match(incoming, interaction.request)) {
  135. return { interaction: undefined, detail: requestDiff(interaction.request, incoming).join("\n") }
  136. }
  137. yield* Ref.update(cursor, (n) => n + 1)
  138. return { interaction, detail: "" }
  139. }
  140. const interaction = interactions.find((candidate) => match(incoming, candidate.request))
  141. return { interaction, detail: interaction ? "" : mismatchDetail(cassette, incoming) }
  142. })
  143. const loadReplay = (request: HttpClientRequest.HttpClientRequest) =>
  144. Effect.gen(function* () {
  145. const cached = yield* Ref.get(replay)
  146. if (cached) return cached
  147. const cassette = yield* cassetteService.read(name).pipe(Effect.mapError(() => fixtureMissing(request, name)))
  148. yield* Ref.set(replay, cassette)
  149. return cassette
  150. })
  151. return HttpClient.make((request) => {
  152. if (mode === "passthrough") return upstream.execute(request)
  153. if (mode === "record") {
  154. return Effect.gen(function* () {
  155. const currentRequest = yield* snapshotRequest(request)
  156. const response = yield* upstream.execute(request)
  157. const headers = responseHeaders(response, responseHeadersAllow, options.redact?.headers)
  158. const captured = yield* captureResponseBody(response, headers["content-type"])
  159. const interaction: HttpInteraction = {
  160. transport: "http",
  161. request: currentRequest,
  162. response: { status: response.status, headers, ...captured },
  163. }
  164. const result = yield* cassetteService.append(name, interaction, options.metadata).pipe(Effect.orDie)
  165. const findings = result.findings
  166. if (findings.length > 0) return yield* unsafeCassette(request, name, findings)
  167. return HttpClientResponse.fromWeb(
  168. request,
  169. new Response(decodeResponseBody(interaction.response), interaction.response),
  170. )
  171. })
  172. }
  173. return Effect.gen(function* () {
  174. const cassette = yield* loadReplay(request)
  175. const incoming = yield* snapshotRequest(request)
  176. const { interaction, detail } = yield* selectInteraction(cassette, incoming)
  177. if (!interaction) return yield* fixtureMismatch(request, name, detail)
  178. return HttpClientResponse.fromWeb(
  179. request,
  180. new Response(decodeResponseBody(interaction.response), interaction.response),
  181. )
  182. })
  183. })
  184. }),
  185. )
  186. export const cassetteLayer = (name: string, options: RecordReplayOptions = {}): Layer.Layer<HttpClient.HttpClient> =>
  187. recordingLayer(name, options).pipe(
  188. Layer.provide(CassetteService.layer({ directory: options.directory })),
  189. Layer.provide(FetchHttpClient.layer),
  190. Layer.provide(NodeFileSystem.layer),
  191. )