effect.ts 5.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142
  1. import { NodeFileSystem } from "@effect/platform-node"
  2. import { Effect, Layer, Option } from "effect"
  3. import {
  4. FetchHttpClient,
  5. Headers,
  6. HttpBody,
  7. HttpClient,
  8. HttpClientError,
  9. HttpClientRequest,
  10. HttpClientResponse,
  11. UrlParams,
  12. } from "effect/unstable/http"
  13. import * as CassetteService from "./cassette"
  14. import { defaultMatcher, selectSequential, type RequestMatcher } from "./matching"
  15. import { makeReplayState, resolveAutoMode } from "./recorder"
  16. import { defaults, type Redactor } from "./redactor"
  17. import { redactUrl } from "./redaction"
  18. import { httpInteractions, type CassetteMetadata, type HttpInteraction, type ResponseSnapshot } from "./schema"
  19. export type RecordReplayMode = "auto" | "record" | "replay" | "passthrough"
  20. export interface RecordReplayOptions {
  21. readonly mode?: RecordReplayMode
  22. readonly directory?: string
  23. readonly metadata?: CassetteMetadata
  24. readonly redactor?: Redactor
  25. readonly match?: RequestMatcher
  26. }
  27. const BINARY_CONTENT_TYPES: ReadonlyArray<string> = ["vnd.amazon.eventstream", "octet-stream"]
  28. const isBinaryContentType = (contentType: string | undefined) =>
  29. contentType !== undefined && BINARY_CONTENT_TYPES.some((token) => contentType.toLowerCase().includes(token))
  30. const captureResponseBody = (response: HttpClientResponse.HttpClientResponse, contentType: string | undefined) =>
  31. isBinaryContentType(contentType)
  32. ? response.arrayBuffer.pipe(
  33. Effect.map((bytes) => ({ body: Buffer.from(bytes).toString("base64"), bodyEncoding: "base64" as const })),
  34. )
  35. : response.text.pipe(Effect.map((body) => ({ body })))
  36. const decodeResponseBody = (snapshot: ResponseSnapshot) =>
  37. snapshot.bodyEncoding === "base64" ? Buffer.from(snapshot.body, "base64") : snapshot.body
  38. export const redactedErrorRequest = (request: HttpClientRequest.HttpClientRequest) =>
  39. HttpClientRequest.makeWith(
  40. request.method,
  41. redactUrl(request.url),
  42. UrlParams.empty,
  43. Option.none(),
  44. Headers.empty,
  45. HttpBody.empty,
  46. )
  47. const transportError = (request: HttpClientRequest.HttpClientRequest, description: string) =>
  48. new HttpClientError.HttpClientError({
  49. reason: new HttpClientError.TransportError({ request: redactedErrorRequest(request), description }),
  50. })
  51. export const recordingLayer = (
  52. name: string,
  53. options: Omit<RecordReplayOptions, "directory"> = {},
  54. ): Layer.Layer<HttpClient.HttpClient, never, HttpClient.HttpClient | CassetteService.Service> =>
  55. Layer.effect(
  56. HttpClient.HttpClient,
  57. Effect.gen(function* () {
  58. const upstream = yield* HttpClient.HttpClient
  59. const cassetteService = yield* CassetteService.Service
  60. const redactor = options.redactor ?? defaults()
  61. const match = options.match ?? defaultMatcher
  62. const requested = options.mode ?? "auto"
  63. const mode = requested === "auto" ? yield* resolveAutoMode(cassetteService, name) : requested
  64. const replay = yield* makeReplayState(cassetteService, name, httpInteractions)
  65. const snapshotRequest = (request: HttpClientRequest.HttpClientRequest) =>
  66. Effect.gen(function* () {
  67. const web = yield* HttpClientRequest.toWeb(request).pipe(Effect.orDie)
  68. return redactor.request({
  69. method: web.method,
  70. url: web.url,
  71. headers: Object.fromEntries(web.headers.entries()),
  72. body: yield* Effect.promise(() => web.text()),
  73. })
  74. })
  75. return HttpClient.make((request) => {
  76. if (mode === "passthrough") return upstream.execute(request)
  77. if (mode === "record") {
  78. return Effect.gen(function* () {
  79. const incoming = yield* snapshotRequest(request)
  80. const response = yield* upstream.execute(request)
  81. const captured = yield* captureResponseBody(response, response.headers["content-type"])
  82. const interaction: HttpInteraction = {
  83. transport: "http",
  84. request: incoming,
  85. response: redactor.response({
  86. status: response.status,
  87. headers: response.headers as Record<string, string>,
  88. ...captured,
  89. }),
  90. }
  91. yield* cassetteService
  92. .append(name, interaction, options.metadata)
  93. .pipe(
  94. Effect.catchTag("UnsafeCassetteError", (error) => Effect.fail(transportError(request, error.message))),
  95. )
  96. return HttpClientResponse.fromWeb(
  97. request,
  98. new Response(decodeResponseBody(interaction.response), interaction.response),
  99. )
  100. })
  101. }
  102. return Effect.gen(function* () {
  103. const incoming = yield* snapshotRequest(request)
  104. const interactions = yield* replay.load.pipe(
  105. Effect.mapError(() =>
  106. transportError(request, `Fixture "${name}" not found. Run locally to record it (CI=true forces replay).`),
  107. ),
  108. )
  109. const result = selectSequential(interactions, incoming, match, yield* replay.cursor)
  110. if (!result.interaction)
  111. return yield* Effect.fail(
  112. transportError(request, `Fixture "${name}" does not match the current request: ${result.detail}.`),
  113. )
  114. yield* replay.advance
  115. return HttpClientResponse.fromWeb(
  116. request,
  117. new Response(decodeResponseBody(result.interaction.response), result.interaction.response),
  118. )
  119. })
  120. })
  121. }),
  122. )
  123. export const cassetteLayer = (name: string, options: RecordReplayOptions = {}): Layer.Layer<HttpClient.HttpClient> =>
  124. recordingLayer(name, options).pipe(
  125. Layer.provide(CassetteService.fileSystem({ directory: options.directory })),
  126. Layer.provide(FetchHttpClient.layer),
  127. Layer.provide(NodeFileSystem.layer),
  128. )