bab2fbc7f6
Cleanup of the Bedrock adapter (ba1705d) following parallel review
passes for code reuse, code quality, and efficiency.
- Drop dead `text` join helper and unused `TextPart` import.
- Schema-validate `model.native.aws_credentials` instead of seven
manual `typeof` guards in `credentialsFromInput`. Removes the
unsafe `as Record<string, unknown>` cast and fixes the dead
`native?.region` fallback (the `model()` constructor only writes
`aws_region`).
- Skip the JSON.parse → JSON.stringify → Schema.fromJsonString triple
round-trip in the frame consumer. The eventstream codec already
hands us a UTF-8 payload; parse once and feed the wrapped object
directly to `Schema.decodeUnknownSync(BedrockChunk)`.
- Replace O(n²) buffer concat in `consumeFrames` with a cursor-based
state `{ buffer, offset }`. Compaction happens once per network
chunk via `appendChunk` instead of per frame; frame slicing is
zero-copy via `subarray`. Bounded buffer growth regardless of
stream length.
- Rename `ParserState.finishReason` → `pendingStopReason` (raw
string) and defer the `mapFinishReason` call to the single emit
site, plus the `onHalt` fallback. Tightens the helper's signature
to `(reason: string)` so the chunk-typed `messageStop.stopReason`
flows through without the optional widening.
- Restructure `signRequest` to take an object parameter (was four
positional args), and replace the manual `forEach`-into-record with
`Object.fromEntries(signed.headers.entries())`.
- Inline single-use `status` and `useTools` variables.
- Widen `fixedResponse` to accept `ConstructorParameters<Response>[0]`
so binary fixtures (`Uint8Array`, streams) flow without casts. The
Bedrock test's `fixedBytes` helper now wraps it cleanly.
- Tidy `captureResponseBody` into a ternary returning the union shape
directly so the call site spreads the captured object without
reaching for `bodyEncoding` explicitly.
Verified: `bun typecheck` clean, 106 pass / 0 fail / 0 skip
(unchanged from before the refactor).
87 lines
3.2 KiB
TypeScript
87 lines
3.2 KiB
TypeScript
import { Effect, Layer, Ref } from "effect"
|
|
import { HttpClient, HttpClientRequest, HttpClientResponse } from "effect/unstable/http"
|
|
import { RequestExecutor } from "../../src/executor"
|
|
|
|
export type HandlerInput = {
|
|
readonly request: HttpClientRequest.HttpClientRequest
|
|
readonly text: string
|
|
readonly respond: (body: ConstructorParameters<typeof Response>[0], init?: ResponseInit) => HttpClientResponse.HttpClientResponse
|
|
}
|
|
|
|
export type Handler = (input: HandlerInput) => Effect.Effect<HttpClientResponse.HttpClientResponse>
|
|
|
|
const handlerLayer = (handler: Handler): Layer.Layer<HttpClient.HttpClient> =>
|
|
Layer.succeed(
|
|
HttpClient.HttpClient,
|
|
HttpClient.make((request) =>
|
|
Effect.gen(function* () {
|
|
const web = yield* HttpClientRequest.toWeb(request).pipe(Effect.orDie)
|
|
const text = yield* Effect.promise(() => web.text())
|
|
return yield* handler({
|
|
request,
|
|
text,
|
|
respond: (body, init) => HttpClientResponse.fromWeb(request, new Response(body, init)),
|
|
})
|
|
}),
|
|
),
|
|
)
|
|
|
|
const executorWith = (layer: Layer.Layer<HttpClient.HttpClient>) =>
|
|
RequestExecutor.layer.pipe(Layer.provide(layer))
|
|
|
|
const SSE_HEADERS = { "content-type": "text/event-stream" } as const
|
|
|
|
/**
|
|
* Layer that returns a single fixed response body. Use for stream-parser
|
|
* fixture tests where the request shape is irrelevant. The body type widens
|
|
* to whatever `Response` accepts so binary fixtures (`Uint8Array`,
|
|
* `ReadableStream`, etc.) flow through without casts.
|
|
*/
|
|
export const fixedResponse = (
|
|
body: ConstructorParameters<typeof Response>[0],
|
|
init: ResponseInit = { headers: SSE_HEADERS },
|
|
) => executorWith(handlerLayer((input) => Effect.succeed(input.respond(body, init))))
|
|
|
|
/**
|
|
* Layer that builds a response per request. Useful for echo servers.
|
|
*/
|
|
export const dynamicResponse = (handler: Handler) => executorWith(handlerLayer(handler))
|
|
|
|
/**
|
|
* Layer that emits the supplied SSE chunks and then aborts mid-stream. Used to
|
|
* exercise transport errors that surface during parsing.
|
|
*/
|
|
export const truncatedStream = (chunks: ReadonlyArray<string>) =>
|
|
dynamicResponse((input) =>
|
|
Effect.sync(() => {
|
|
const encoder = new TextEncoder()
|
|
const stream = new ReadableStream({
|
|
start(controller) {
|
|
for (const chunk of chunks) controller.enqueue(encoder.encode(chunk))
|
|
controller.error(new Error("connection reset"))
|
|
},
|
|
})
|
|
return input.respond(stream, { headers: SSE_HEADERS })
|
|
}),
|
|
)
|
|
|
|
/**
|
|
* Layer that returns successive bodies on each request. Useful for scripting
|
|
* multi-step model exchanges (e.g. tool-call loops). The last body in the
|
|
* array is reused if the test makes more requests than scripted.
|
|
*/
|
|
export const scriptedResponses = (bodies: ReadonlyArray<string>, init: ResponseInit = { headers: SSE_HEADERS }) => {
|
|
if (bodies.length === 0) throw new Error("scriptedResponses requires at least one body")
|
|
return Layer.unwrap(
|
|
Effect.gen(function* () {
|
|
const cursor = yield* Ref.make(0)
|
|
return dynamicResponse((input) =>
|
|
Effect.gen(function* () {
|
|
const index = yield* Ref.getAndUpdate(cursor, (n) => n + 1)
|
|
return input.respond(bodies[index] ?? bodies[bodies.length - 1], init)
|
|
}),
|
|
)
|
|
}),
|
|
)
|
|
}
|