Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions .devcontainer/devcontainer-lock.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
{
"features": {
"ghcr.io/devcontainers-extra/features/bun:1": {
"version": "1.1.0",
"resolved": "ghcr.io/devcontainers-extra/features/bun@sha256:0624284ecaead9dd4c6654616a7f939cfa4ebcbc60593700a74e35b1767befa5",
"integrity": "sha256:0624284ecaead9dd4c6654616a7f939cfa4ebcbc60593700a74e35b1767befa5"
}
}
}
8 changes: 8 additions & 0 deletions .devcontainer/devcontainer.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
{
"image": "mcr.microsoft.com/devcontainers/base:debian",
"features": {
"ghcr.io/devcontainers-extra/features/bun:1": {
"version": "1.3.14"
}
}
}
7 changes: 7 additions & 0 deletions .github/workflows/test.yml
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,13 @@ jobs:
turbo-${{ runner.os }}-${{ hashFiles('turbo.json', '**/package.json') }}-
turbo-${{ runner.os }}-

- name: Lint http-recorder changes
run: bun lint packages/http-recorder/src/socket.ts packages/http-recorder/test/record-replay.test.ts

- name: Run http-recorder tests with coverage
working-directory: packages/http-recorder
run: bun test --coverage --coverage-dir=./coverage

- name: Run unit tests
timeout-minutes: 10
# Default turbo concurrency (10) oversubscribes ubuntu-latest's 4
Expand Down
34 changes: 17 additions & 17 deletions packages/http-recorder/src/socket.ts
Original file line number Diff line number Diff line change
Expand Up @@ -116,14 +116,14 @@ const openSnapshot = (request: WebSocketRequest, redactor: Redactor) => {
return { url: snapshot.url, headers: snapshot.headers }
}

const makeRecordingSocket = (
upstream: Socket.Socket,
cassette: CassetteService.Interface,
name: string,
request: WebSocketRequest,
options: WebSocketRecorderOptions,
redactor: Redactor,
) =>
const makeRecordingSocket = (input: {
upstream: Socket.Socket
cassette: CassetteService.Interface
name: string
request: WebSocketRequest
options: WebSocketRecorderOptions
redactor: Redactor
}) =>
Effect.gen(function* () {
const active = yield* Ref.make<ActiveRecording | undefined>(undefined)
const writeLock = yield* Semaphore.make(1)
Expand All @@ -140,11 +140,11 @@ const makeRecordingSocket = (
}
const occupied = yield* Ref.modify(active, (current) => [current !== undefined, current ?? state])
if (occupied) return yield* Effect.die("Concurrent runs of a recorded WebSocket are not supported")
yield* upstream
yield* input.upstream
.runRaw(
(message) => {
if (!Ref.getUnsafe(state.accepting)) throw new Error("WebSocket received a frame after closing")
state.events.push(redactEvent(encodeEvent("server", message), redactor))
state.events.push(redactEvent(encodeEvent("server", message), input.redactor))
return handler(message)
},
{
Expand All @@ -163,15 +163,15 @@ const makeRecordingSocket = (
yield* Ref.set(state.accepting, false)
yield* Ref.set(active, undefined)
if (!Exit.isSuccess(exit) || !state.opened || !state.valid) return
yield* cassette
yield* input.cassette
.append(
name,
input.name,
{
transport: "websocket",
open: openSnapshot(request, redactor),
open: openSnapshot(input.request, input.redactor),
events: [...state.events],
},
options.metadata,
input.options.metadata,
)
.pipe(Effect.orDie)
}),
Expand All @@ -180,7 +180,7 @@ const makeRecordingSocket = (
),
)
}),
writer: upstream.writer.pipe(
writer: input.upstream.writer.pipe(
Effect.map(
(write) => (message) =>
writeLock.withPermit(
Expand All @@ -189,7 +189,7 @@ const makeRecordingSocket = (
const state = yield* Ref.get(active)
if (!state || !(yield* Ref.get(state.accepting)))
return yield* Effect.die("WebSocket writer used without an active socket run")
const event = redactEvent(encodeEvent("client", message), redactor)
const event = redactEvent(encodeEvent("client", message), input.redactor)
yield* state.eventLock.withPermit(Effect.sync(() => state.events.push(event)))
return yield* write(message).pipe(Effect.onError(() => Effect.sync(() => (state.valid = false))))
}),
Expand Down Expand Up @@ -293,7 +293,7 @@ const recordingLayer = (
const cassette = yield* CassetteService.Service
const redactor = make(options.redact)
if ((forcedMode ?? (yield* resolveAutoMode(cassette, name))) === "record")
return yield* makeRecordingSocket(upstream, cassette, name, request, options, redactor)
return yield* makeRecordingSocket({ upstream, cassette, name, request, options, redactor })
return yield* makeReplaySocket(cassette, name, request, options, redactor)
}),
)
Expand Down
38 changes: 30 additions & 8 deletions packages/http-recorder/test/record-replay.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import { NodeFileSystem } from "@effect/platform-node"
import { describe, expect, test } from "bun:test"
import { Cause, Deferred, Effect, Exit, Layer, Scope, Stream } from "effect"

Check warning on line 3 in packages/http-recorder/test/record-replay.test.ts

View workflow job for this annotation

GitHub Actions / unit

eslint(no-unused-vars)

Identifier 'Stream' is imported but never used.
import { Headers, HttpBody, HttpClient, HttpClientRequest } from "effect/unstable/http"

Check warning on line 4 in packages/http-recorder/test/record-replay.test.ts

View workflow job for this annotation

GitHub Actions / unit

eslint(no-unused-vars)

Identifier 'Headers' is imported but never used.
import { Socket } from "effect/unstable/socket"
import * as fs from "node:fs"
Expand Down Expand Up @@ -42,7 +42,7 @@
effect: Effect.Effect<A, E, HttpClient.HttpClient>,
) => Effect.runPromise(effect.pipe(Effect.provide(HttpRecorder.http(name, options))))

const runRecorder = <A, E>(effect: Effect.Effect<A, E, HttpRecorderInternal.Cassette.Service | Scope.Scope>) =>

Check warning on line 45 in packages/http-recorder/test/record-replay.test.ts

View workflow job for this annotation

GitHub Actions / unit

eslint(no-unused-vars)

Variable 'runRecorder' is declared but never used. Unused variables should start with a '_'.
Effect.runPromise(
Effect.scoped(
effect.pipe(
Expand Down Expand Up @@ -207,9 +207,14 @@
})
})

test("records WebSocket frames in observed client/server order", async () => {
test("records WebSocket frames in order with metadata and redaction", async () => {
const directory = fs.mkdtempSync(path.join(os.tmpdir(), "http-recorder-websocket-"))
const response = JSON.stringify({ type: "response.completed", token: "server-secret" })
const response = JSON.stringify({
type: "response.completed",
token: "server-secret",
privateValue: "server-value",
})
const received: Array<string | Uint8Array> = []
let receive: ((message: string | Uint8Array) => Effect.Effect<unknown, unknown, unknown> | void) | undefined
const upstream = Socket.make({
runRaw: (handler, options) =>
Expand All @@ -230,29 +235,46 @@
Effect.gen(function* () {
const socket = yield* Socket.Socket
const write = yield* socket.writer
yield* socket.runRaw(() => {}, {
onOpen: write(JSON.stringify({ type: "response.create", token: "client-secret" })),
})
yield* socket.runRaw(
(message) => {
received.push(message)
},
{
onOpen: write(
JSON.stringify({ type: "response.create", token: "client-secret", privateValue: "client-value" }),
),
},
)
}).pipe(
Effect.scoped,
Effect.provide(
HttpRecorderInternal.socketLayer(
"websocket/record",
{ url: "wss://example.test/realtime", headers: { "content-type": "application/json" } },
{ directory, metadata: { provider: "test" }, mode: "record" },
{ directory, metadata: { provider: "test" }, redact: { jsonFields: ["privateValue"] }, mode: "record" },
).pipe(Layer.provide(Layer.succeed(Socket.Socket, upstream))),
),
),
)

expect(received).toEqual([response])
expect(JSON.parse(fs.readFileSync(path.join(directory, "websocket/record.json"), "utf8"))).toMatchObject({
metadata: { name: "websocket/record", provider: "test" },
interactions: [
{
transport: "websocket",
open: { url: "wss://example.test/realtime", headers: { "content-type": "application/json" } },
events: [
{ direction: "client", kind: "text", body: '{"type":"response.create","token":"[REDACTED]"}' },
{ direction: "server", kind: "text", body: '{"type":"response.completed","token":"[REDACTED]"}' },
{
direction: "client",
kind: "text",
body: '{"type":"response.create","token":"[REDACTED]","privateValue":"[REDACTED]"}',
},
{
direction: "server",
kind: "text",
body: '{"type":"response.completed","token":"[REDACTED]","privateValue":"[REDACTED]"}',
},
],
},
],
Expand Down
Loading