Buckets:
| import { NodeFileSystem } from "@effect/platform-node" | |
| import { Deferred, Effect, Layer, Option, Ref } from "effect" | |
| import { | |
| FetchHttpClient, | |
| Headers, | |
| HttpBody, | |
| HttpClient, | |
| HttpClientError, | |
| HttpClientRequest, | |
| HttpClientResponse, | |
| UrlParams, | |
| } from "effect/unstable/http" | |
| import * as CassetteService from "./cassette.js" | |
| import { defaultMatcher, selectSequential } from "./matching.js" | |
| import { makeReplayState, resolveAutoMode } from "./recorder.js" | |
| import { make, type Redactor } from "./redactor.js" | |
| import { redactUrl } from "./redaction.js" | |
| import { httpInteractions } from "./schema.js" | |
| import type { CassetteMetadata, HttpInteraction, RequestMatcher, ResponseSnapshot } from "./types.js" | |
| export { defaultMatcher } | |
| export type RecordReplayMode = "auto" | "record" | "replay" | "passthrough" | |
| export interface RecordReplayOptions { | |
| readonly mode?: RecordReplayMode | |
| readonly directory?: string | |
| readonly metadata?: CassetteMetadata | |
| readonly redactor?: Redactor | |
| readonly match?: RequestMatcher | |
| } | |
| const TEXT_CONTENT_TYPES = new Set([ | |
| "application/graphql", | |
| "application/javascript", | |
| "application/json", | |
| "application/sql", | |
| "application/x-www-form-urlencoded", | |
| "application/xml", | |
| "application/yaml", | |
| "image/svg+xml", | |
| ]) | |
| const isTextContentType = (contentType: string | undefined) => { | |
| const mediaType = contentType?.split(";", 1)[0]?.trim().toLowerCase() | |
| if (!mediaType) return false | |
| return ( | |
| mediaType.startsWith("text/") || | |
| mediaType.endsWith("+json") || | |
| mediaType.endsWith("+xml") || | |
| TEXT_CONTENT_TYPES.has(mediaType) | |
| ) | |
| } | |
| const captureResponseBody = (response: HttpClientResponse.HttpClientResponse, contentType: string | undefined) => | |
| response.arrayBuffer.pipe( | |
| Effect.map((bytes) => | |
| isTextContentType(contentType) | |
| ? { body: new TextDecoder().decode(bytes) } | |
| : { body: Buffer.from(bytes).toString("base64"), bodyEncoding: "base64" as const }, | |
| ), | |
| ) | |
| const decodeResponseBody = (snapshot: ResponseSnapshot) => | |
| snapshot.bodyEncoding === "base64" ? Buffer.from(snapshot.body, "base64") : snapshot.body | |
| const responseFromSnapshot = (request: HttpClientRequest.HttpClientRequest, snapshot: ResponseSnapshot) => | |
| HttpClientResponse.fromWeb( | |
| request, | |
| new Response( | |
| request.method === "HEAD" || snapshot.status === 204 || snapshot.status === 205 || snapshot.status === 304 | |
| ? null | |
| : decodeResponseBody(snapshot), | |
| snapshot, | |
| ), | |
| ) | |
| export const redactedErrorRequest = (request: HttpClientRequest.HttpClientRequest) => | |
| HttpClientRequest.makeWith( | |
| request.method, | |
| redactUrl(request.url), | |
| UrlParams.empty, | |
| Option.none(), | |
| Headers.empty, | |
| HttpBody.empty, | |
| ) | |
| const transportError = (request: HttpClientRequest.HttpClientRequest, description: string) => | |
| new HttpClientError.HttpClientError({ | |
| reason: new HttpClientError.TransportError({ request: redactedErrorRequest(request), description }), | |
| }) | |
| export const recordingLayer = ( | |
| name: string, | |
| options: Omit<RecordReplayOptions, "directory"> = {}, | |
| ): Layer.Layer<HttpClient.HttpClient, never, HttpClient.HttpClient | CassetteService.Service> => | |
| Layer.effect( | |
| HttpClient.HttpClient, | |
| Effect.gen(function* () { | |
| const upstream = yield* HttpClient.HttpClient | |
| const cassetteService = yield* CassetteService.Service | |
| const redactor = options.redactor ?? make() | |
| const match = options.match ?? defaultMatcher | |
| const requested = options.mode ?? "auto" | |
| const mode = requested === "auto" ? yield* resolveAutoMode(cassetteService, name) : requested | |
| const snapshotRequest = (request: HttpClientRequest.HttpClientRequest) => | |
| Effect.gen(function* () { | |
| const web = yield* HttpClientRequest.toWeb(request).pipe(Effect.orDie) | |
| return redactor.request({ | |
| method: web.method, | |
| url: web.url, | |
| headers: Object.fromEntries(web.headers.entries()), | |
| body: yield* Effect.promise(() => web.text()), | |
| }) | |
| }) | |
| if (mode === "passthrough") return upstream | |
| if (mode === "record") { | |
| const initial = yield* Deferred.make<void>() | |
| yield* Deferred.succeed(initial, undefined) | |
| const tail = yield* Ref.make(initial) | |
| return HttpClient.make((request) => | |
| Effect.gen(function* () { | |
| const completed = yield* Deferred.make<void>() | |
| const previous = yield* Ref.modify(tail, (current) => [current, completed]) | |
| return yield* Effect.gen(function* () { | |
| const incoming = yield* snapshotRequest(request) | |
| const response = yield* upstream.execute(request) | |
| const captured = yield* captureResponseBody(response, response.headers["content-type"]) | |
| const responseSnapshot: ResponseSnapshot = { | |
| status: response.status, | |
| headers: response.headers as Record<string, string>, | |
| ...captured, | |
| } | |
| const interaction: HttpInteraction = { | |
| transport: "http", | |
| request: incoming, | |
| response: redactor.response(responseSnapshot), | |
| } | |
| yield* Deferred.await(previous) | |
| yield* cassetteService | |
| .append(name, interaction, options.metadata) | |
| .pipe( | |
| Effect.catchTag("UnsafeCassetteError", (error) => | |
| Effect.fail(transportError(request, error.message)), | |
| ), | |
| ) | |
| return responseFromSnapshot(request, responseSnapshot) | |
| }).pipe(Effect.ensuring(Deferred.succeed(completed, undefined))) | |
| }), | |
| ) | |
| } | |
| const replay = yield* makeReplayState(cassetteService, name, httpInteractions) | |
| return HttpClient.make((request) => | |
| Effect.gen(function* () { | |
| const incoming = yield* snapshotRequest(request) | |
| const claimed = yield* replay | |
| .claim((interaction, index, interactions) => { | |
| const result = selectSequential(interactions, incoming, match, index) | |
| if (result.interaction) return Effect.void | |
| return Effect.fail( | |
| transportError(request, `Fixture "${name}" does not match the current request: ${result.detail}.`), | |
| ) | |
| }) | |
| .pipe( | |
| Effect.mapError((error) => | |
| error._tag === "CassetteNotFoundError" | |
| ? transportError( | |
| request, | |
| `Fixture "${name}" not found. Run locally to record it (CI=true forces replay).`, | |
| ) | |
| : error, | |
| ), | |
| ) | |
| return responseFromSnapshot(request, claimed.interaction.response) | |
| }), | |
| ) | |
| }), | |
| ) | |
| export const cassetteLayer = (name: string, options: RecordReplayOptions = {}): Layer.Layer<HttpClient.HttpClient> => | |
| recordingLayer(name, options).pipe( | |
| Layer.provide(CassetteService.fileSystem({ directory: options.directory })), | |
| Layer.provide(FetchHttpClient.layer), | |
| Layer.provide(NodeFileSystem.layer), | |
| ) | |
Xet Storage Details
- Size:
- 7.16 kB
- Xet hash:
- 9645ec8ffad976fc4cec8d14c907c84415c3d3a0f78b15ef74064cc313924dce
·
Xet efficiently stores files, intelligently splitting them into unique chunks and accelerating uploads and downloads. More info.