internal-effect.ts 7.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189
  1. import { NodeFileSystem } from "@effect/platform-node"
  2. import { Deferred, Effect, Layer, Option, Ref } 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.js"
  14. import { defaultMatcher, selectSequential } from "./matching.js"
  15. import { makeReplayState, resolveAutoMode } from "./recorder.js"
  16. import { make, type Redactor } from "./redactor.js"
  17. import { redactUrl } from "./redaction.js"
  18. import { httpInteractions } from "./schema.js"
  19. import type { CassetteMetadata, HttpInteraction, RequestMatcher, ResponseSnapshot } from "./types.js"
  20. export { defaultMatcher }
  21. export type RecordReplayMode = "auto" | "record" | "replay" | "passthrough"
  22. export interface RecordReplayOptions {
  23. readonly mode?: RecordReplayMode
  24. readonly directory?: string
  25. readonly metadata?: CassetteMetadata
  26. readonly redactor?: Redactor
  27. readonly match?: RequestMatcher
  28. }
  29. const TEXT_CONTENT_TYPES = new Set([
  30. "application/graphql",
  31. "application/javascript",
  32. "application/json",
  33. "application/sql",
  34. "application/x-www-form-urlencoded",
  35. "application/xml",
  36. "application/yaml",
  37. "image/svg+xml",
  38. ])
  39. const isTextContentType = (contentType: string | undefined) => {
  40. const mediaType = contentType?.split(";", 1)[0]?.trim().toLowerCase()
  41. if (!mediaType) return false
  42. return (
  43. mediaType.startsWith("text/") ||
  44. mediaType.endsWith("+json") ||
  45. mediaType.endsWith("+xml") ||
  46. TEXT_CONTENT_TYPES.has(mediaType)
  47. )
  48. }
  49. const captureResponseBody = (response: HttpClientResponse.HttpClientResponse, contentType: string | undefined) =>
  50. response.arrayBuffer.pipe(
  51. Effect.map((bytes) =>
  52. isTextContentType(contentType)
  53. ? { body: new TextDecoder().decode(bytes) }
  54. : { body: Buffer.from(bytes).toString("base64"), bodyEncoding: "base64" as const },
  55. ),
  56. )
  57. const decodeResponseBody = (snapshot: ResponseSnapshot) =>
  58. snapshot.bodyEncoding === "base64" ? Buffer.from(snapshot.body, "base64") : snapshot.body
  59. const responseFromSnapshot = (request: HttpClientRequest.HttpClientRequest, snapshot: ResponseSnapshot) =>
  60. HttpClientResponse.fromWeb(
  61. request,
  62. new Response(
  63. request.method === "HEAD" || snapshot.status === 204 || snapshot.status === 205 || snapshot.status === 304
  64. ? null
  65. : decodeResponseBody(snapshot),
  66. snapshot,
  67. ),
  68. )
  69. export const redactedErrorRequest = (request: HttpClientRequest.HttpClientRequest) =>
  70. HttpClientRequest.makeWith(
  71. request.method,
  72. redactUrl(request.url),
  73. UrlParams.empty,
  74. Option.none(),
  75. Headers.empty,
  76. HttpBody.empty,
  77. )
  78. const transportError = (request: HttpClientRequest.HttpClientRequest, description: string) =>
  79. new HttpClientError.HttpClientError({
  80. reason: new HttpClientError.TransportError({ request: redactedErrorRequest(request), description }),
  81. })
  82. export const recordingLayer = (
  83. name: string,
  84. options: Omit<RecordReplayOptions, "directory"> = {},
  85. ): Layer.Layer<HttpClient.HttpClient, never, HttpClient.HttpClient | CassetteService.Service> =>
  86. Layer.effect(
  87. HttpClient.HttpClient,
  88. Effect.gen(function* () {
  89. const upstream = yield* HttpClient.HttpClient
  90. const cassetteService = yield* CassetteService.Service
  91. const redactor = options.redactor ?? make()
  92. const match = options.match ?? defaultMatcher
  93. const requested = options.mode ?? "auto"
  94. const mode = requested === "auto" ? yield* resolveAutoMode(cassetteService, name) : requested
  95. const snapshotRequest = (request: HttpClientRequest.HttpClientRequest) =>
  96. Effect.gen(function* () {
  97. const web = yield* HttpClientRequest.toWeb(request).pipe(Effect.orDie)
  98. return redactor.request({
  99. method: web.method,
  100. url: web.url,
  101. headers: Object.fromEntries(web.headers.entries()),
  102. body: yield* Effect.promise(() => web.text()),
  103. })
  104. })
  105. if (mode === "passthrough") return upstream
  106. if (mode === "record") {
  107. const initial = yield* Deferred.make<void>()
  108. yield* Deferred.succeed(initial, undefined)
  109. const tail = yield* Ref.make(initial)
  110. return HttpClient.make((request) =>
  111. Effect.gen(function* () {
  112. const completed = yield* Deferred.make<void>()
  113. const previous = yield* Ref.modify(tail, (current) => [current, completed])
  114. return yield* Effect.gen(function* () {
  115. const incoming = yield* snapshotRequest(request)
  116. const response = yield* upstream.execute(request)
  117. const captured = yield* captureResponseBody(response, response.headers["content-type"])
  118. const responseSnapshot: ResponseSnapshot = {
  119. status: response.status,
  120. headers: response.headers as Record<string, string>,
  121. ...captured,
  122. }
  123. const interaction: HttpInteraction = {
  124. transport: "http",
  125. request: incoming,
  126. response: redactor.response(responseSnapshot),
  127. }
  128. yield* Deferred.await(previous)
  129. yield* cassetteService
  130. .append(name, interaction, options.metadata)
  131. .pipe(
  132. Effect.catchTag("UnsafeCassetteError", (error) =>
  133. Effect.fail(transportError(request, error.message)),
  134. ),
  135. )
  136. return responseFromSnapshot(request, responseSnapshot)
  137. }).pipe(Effect.ensuring(Deferred.succeed(completed, undefined)))
  138. }),
  139. )
  140. }
  141. const replay = yield* makeReplayState(cassetteService, name, httpInteractions)
  142. return HttpClient.make((request) =>
  143. Effect.gen(function* () {
  144. const incoming = yield* snapshotRequest(request)
  145. const claimed = yield* replay
  146. .claim((interaction, index, interactions) => {
  147. const result = selectSequential(interactions, incoming, match, index)
  148. if (result.interaction) return Effect.void
  149. return Effect.fail(
  150. transportError(request, `Fixture "${name}" does not match the current request: ${result.detail}.`),
  151. )
  152. })
  153. .pipe(
  154. Effect.mapError((error) =>
  155. error._tag === "CassetteNotFoundError"
  156. ? transportError(
  157. request,
  158. `Fixture "${name}" not found. Run locally to record it (CI=true forces replay).`,
  159. )
  160. : error,
  161. ),
  162. )
  163. return responseFromSnapshot(request, claimed.interaction.response)
  164. }),
  165. )
  166. }),
  167. )
  168. export const cassetteLayer = (name: string, options: RecordReplayOptions = {}): Layer.Layer<HttpClient.HttpClient> =>
  169. recordingLayer(name, options).pipe(
  170. Layer.provide(CassetteService.fileSystem({ directory: options.directory })),
  171. Layer.provide(FetchHttpClient.layer),
  172. Layer.provide(NodeFileSystem.layer),
  173. )