refactor(llm): centralize InvalidRequestError, validate, and JSON POST
Phase A continuation of the ProviderShared dedupe pass. Three more
patterns lifted into ProviderShared so they're written once:
ProviderShared.invalidRequest(message) — replaces six identical
`const invalid = (message) => new InvalidRequestError({ message })`
one-liners across openai-chat, openai-responses, anthropic-messages,
gemini, openai-compatible-chat, and bedrock-converse. Each adapter
keeps a short `const invalid = ProviderShared.invalidRequest` alias
so the 27 callsite `yield* invalid("...")` patterns are unchanged.
Bedrock's SigV4 catch path and the openai-compatible-chat baseURL
guard both go through the helper now too.
ProviderShared.validateWith(decode) — replaces the identical
`(draft) => decode(draft).pipe(Effect.mapError((e) =>
invalid(e.message)))` lambda body in five adapters. Same line count
but shorter, names the pattern, and keeps the `decode → mapError →
InvalidRequestError` translation in one canonical spot.
ProviderShared.jsonPost({ url, body, headers }) — replaces the
five-adapter pattern of `HttpClientRequest.post(url).pipe(setHeaders,
bodyText)` for JSON-body POSTs. Sets `content-type: application/json`
last so caller headers can override everything except the
content-type. Bedrock uses it for both the bearer-auth and SigV4-
signed paths; SigV4 still signs against `baseHeaders` (which already
contained content-type) so the signature matches what the helper
ultimately sends.
Net change: -73 / +86 (+13 in shared.ts mostly JSDoc; -86 across the
six adapters). The `HttpClientRequest` and `InvalidRequestError`
imports are dropped from the five SSE adapters and from Bedrock since
they're no longer referenced directly.
Verified: `bun typecheck` clean, 106 pass / 0 fail / 0 skip
(unchanged).
This commit is contained in:
@@ -1,9 +1,8 @@
|
||||
import { Effect, Schema, Stream } from "effect"
|
||||
import { HttpClientRequest, type HttpClientResponse } from "effect/unstable/http"
|
||||
import type { HttpClientResponse } from "effect/unstable/http"
|
||||
import { Adapter } from "../adapter"
|
||||
import { capabilities, model as llmModel, type ModelInput } from "../llm"
|
||||
import {
|
||||
InvalidRequestError,
|
||||
Usage,
|
||||
type CacheHint,
|
||||
type FinishReason,
|
||||
@@ -205,7 +204,7 @@ const decodeChunk = (data: string) =>
|
||||
const encodeTarget = Schema.encodeSync(AnthropicTargetJson)
|
||||
const decodeTarget = Schema.decodeUnknownEffect(AnthropicMessagesDraft.pipe(Schema.decodeTo(AnthropicMessagesTarget)))
|
||||
|
||||
const invalid = (message: string) => new InvalidRequestError({ message })
|
||||
const invalid = ProviderShared.invalidRequest
|
||||
|
||||
const baseUrl = (request: LLMRequest) => (request.model.baseURL ?? "https://api.anthropic.com/v1").replace(/\/+$/, "")
|
||||
|
||||
@@ -348,14 +347,11 @@ const prepare = Effect.fn("AnthropicMessages.prepare")(function* (request: LLMRe
|
||||
|
||||
const toHttp = (target: AnthropicMessagesTarget, request: LLMRequest) =>
|
||||
Effect.succeed(
|
||||
HttpClientRequest.post(`${baseUrl(request)}/messages`).pipe(
|
||||
HttpClientRequest.setHeaders({
|
||||
"anthropic-version": "2023-06-01",
|
||||
...request.model.headers,
|
||||
"content-type": "application/json",
|
||||
}),
|
||||
HttpClientRequest.bodyText(encodeTarget(target), "application/json"),
|
||||
),
|
||||
ProviderShared.jsonPost({
|
||||
url: `${baseUrl(request)}/messages`,
|
||||
body: encodeTarget(target),
|
||||
headers: { "anthropic-version": "2023-06-01", ...request.model.headers },
|
||||
}),
|
||||
)
|
||||
|
||||
const mapFinishReason = (reason: string | null | undefined): FinishReason => {
|
||||
@@ -529,7 +525,7 @@ export const adapter = Adapter.define<AnthropicMessagesDraft, AnthropicMessagesT
|
||||
protocol: "anthropic-messages",
|
||||
redact: (target) => target,
|
||||
prepare,
|
||||
validate: (draft) => decodeTarget(draft).pipe(Effect.mapError((error) => invalid(error.message))),
|
||||
validate: ProviderShared.validateWith(decodeTarget),
|
||||
toHttp: (target, context) => toHttp(target, context.request),
|
||||
parse: events,
|
||||
})
|
||||
|
||||
@@ -2,11 +2,10 @@ import { EventStreamCodec } from "@smithy/eventstream-codec"
|
||||
import { fromUtf8, toUtf8 } from "@smithy/util-utf8"
|
||||
import { AwsV4Signer } from "aws4fetch"
|
||||
import { Effect, Option, Schema, Stream } from "effect"
|
||||
import { HttpClientRequest, HttpClientResponse } from "effect/unstable/http"
|
||||
import { HttpClientResponse } from "effect/unstable/http"
|
||||
import { Adapter } from "../adapter"
|
||||
import { capabilities, model as llmModel, type ModelInput } from "../llm"
|
||||
import {
|
||||
InvalidRequestError,
|
||||
Usage,
|
||||
type FinishReason,
|
||||
type LLMEvent,
|
||||
@@ -225,7 +224,7 @@ const decodeChunk = (data: unknown) =>
|
||||
const encodeTarget = Schema.encodeSync(Schema.fromJsonString(BedrockConverseTarget))
|
||||
const decodeTarget = Schema.decodeUnknownEffect(BedrockConverseDraft.pipe(Schema.decodeTo(BedrockConverseTarget)))
|
||||
|
||||
const invalid = (message: string) => new InvalidRequestError({ message })
|
||||
const invalid = ProviderShared.invalidRequest
|
||||
|
||||
const region = (request: LLMRequest) => {
|
||||
const fromNative = request.model.native?.aws_region
|
||||
@@ -401,9 +400,7 @@ const signRequest = (input: {
|
||||
return Object.fromEntries(signed.headers.entries())
|
||||
},
|
||||
catch: (error) =>
|
||||
new InvalidRequestError({
|
||||
message: `Bedrock Converse SigV4 signing failed: ${error instanceof Error ? error.message : String(error)}`,
|
||||
}),
|
||||
invalid(`Bedrock Converse SigV4 signing failed: ${error instanceof Error ? error.message : String(error)}`),
|
||||
})
|
||||
|
||||
const toHttp = Effect.fn("BedrockConverse.toHttp")(function* (target: BedrockConverseTarget, request: LLMRequest) {
|
||||
@@ -415,10 +412,7 @@ const toHttp = Effect.fn("BedrockConverse.toHttp")(function* (target: BedrockCon
|
||||
}
|
||||
|
||||
if (isBearerAuth(request.model.headers)) {
|
||||
return HttpClientRequest.post(url).pipe(
|
||||
HttpClientRequest.setHeaders(baseHeaders),
|
||||
HttpClientRequest.bodyText(body, "application/json"),
|
||||
)
|
||||
return ProviderShared.jsonPost({ url, body, headers: request.model.headers })
|
||||
}
|
||||
|
||||
const credentials = credentialsFromInput(request)
|
||||
@@ -427,11 +421,10 @@ const toHttp = Effect.fn("BedrockConverse.toHttp")(function* (target: BedrockCon
|
||||
"Bedrock Converse requires either a Bearer API key in headers or AWS credentials in model.native.aws_credentials",
|
||||
)
|
||||
}
|
||||
// SigV4 signs the request including content-type; keep `baseHeaders` so the
|
||||
// signed payload matches what `jsonPost` ultimately sends.
|
||||
const signed = yield* signRequest({ url, body, headers: baseHeaders, credentials })
|
||||
return HttpClientRequest.post(url).pipe(
|
||||
HttpClientRequest.setHeaders({ ...baseHeaders, ...signed }),
|
||||
HttpClientRequest.bodyText(body, "application/json"),
|
||||
)
|
||||
return ProviderShared.jsonPost({ url, body, headers: { ...baseHeaders, ...signed } })
|
||||
})
|
||||
|
||||
const mapFinishReason = (reason: string): FinishReason => {
|
||||
@@ -666,7 +659,7 @@ export const adapter = Adapter.define<BedrockConverseDraft, BedrockConverseTarge
|
||||
protocol: "bedrock-converse",
|
||||
redact: (target) => target,
|
||||
prepare,
|
||||
validate: (draft) => decodeTarget(draft).pipe(Effect.mapError((error) => invalid(error.message))),
|
||||
validate: ProviderShared.validateWith(decodeTarget),
|
||||
toHttp: (target, context) => toHttp(target, context.request),
|
||||
parse: parseStream,
|
||||
})
|
||||
|
||||
@@ -1,10 +1,9 @@
|
||||
import { Buffer } from "node:buffer"
|
||||
import { Effect, Schema, Stream } from "effect"
|
||||
import { HttpClientRequest, type HttpClientResponse } from "effect/unstable/http"
|
||||
import type { HttpClientResponse } from "effect/unstable/http"
|
||||
import { Adapter } from "../adapter"
|
||||
import { capabilities, model as llmModel, type ModelInput } from "../llm"
|
||||
import {
|
||||
InvalidRequestError,
|
||||
Usage,
|
||||
type FinishReason,
|
||||
type LLMEvent,
|
||||
@@ -151,7 +150,7 @@ const decodeChunk = (data: string) =>
|
||||
const encodeTarget = Schema.encodeSync(GeminiTargetJson)
|
||||
const decodeTarget = Schema.decodeUnknownEffect(GeminiDraft.pipe(Schema.decodeTo(GeminiTarget)))
|
||||
|
||||
const invalid = (message: string) => new InvalidRequestError({ message })
|
||||
const invalid = ProviderShared.invalidRequest
|
||||
|
||||
const baseUrl = (request: LLMRequest) =>
|
||||
(request.model.baseURL ?? "https://generativelanguage.googleapis.com/v1beta").replace(/\/+$/, "")
|
||||
@@ -315,13 +314,11 @@ const prepare = Effect.fn("Gemini.prepare")(function* (request: LLMRequest) {
|
||||
|
||||
const toHttp = (target: GeminiTarget, request: LLMRequest) =>
|
||||
Effect.succeed(
|
||||
HttpClientRequest.post(`${baseUrl(request)}/models/${request.model.id}:streamGenerateContent?alt=sse`).pipe(
|
||||
HttpClientRequest.setHeaders({
|
||||
...request.model.headers,
|
||||
"content-type": "application/json",
|
||||
}),
|
||||
HttpClientRequest.bodyText(encodeTarget(target), "application/json"),
|
||||
),
|
||||
ProviderShared.jsonPost({
|
||||
url: `${baseUrl(request)}/models/${request.model.id}:streamGenerateContent?alt=sse`,
|
||||
body: encodeTarget(target),
|
||||
headers: request.model.headers,
|
||||
}),
|
||||
)
|
||||
|
||||
const mapUsage = (usage: GeminiUsage | undefined) => {
|
||||
@@ -412,7 +409,7 @@ export const adapter = Adapter.define<GeminiDraft, GeminiTarget>({
|
||||
protocol: "gemini",
|
||||
redact: (target) => target,
|
||||
prepare,
|
||||
validate: (draft) => decodeTarget(draft).pipe(Effect.mapError((error) => invalid(error.message))),
|
||||
validate: ProviderShared.validateWith(decodeTarget),
|
||||
toHttp: (target, context) => toHttp(target, context.request),
|
||||
parse: events,
|
||||
})
|
||||
|
||||
@@ -1,9 +1,8 @@
|
||||
import { Effect, Schema, Stream } from "effect"
|
||||
import { HttpClientRequest, type HttpClientResponse } from "effect/unstable/http"
|
||||
import type { HttpClientResponse } from "effect/unstable/http"
|
||||
import { Adapter } from "../adapter"
|
||||
import { capabilities, model as llmModel, type ModelInput } from "../llm"
|
||||
import {
|
||||
InvalidRequestError,
|
||||
Usage,
|
||||
type FinishReason,
|
||||
type ContentPart,
|
||||
@@ -164,7 +163,7 @@ interface ParserState {
|
||||
|
||||
const decodeTarget = Schema.decodeUnknownEffect(OpenAIChatDraft.pipe(Schema.decodeTo(OpenAIChatTarget)))
|
||||
|
||||
const invalid = (message: string) => new InvalidRequestError({ message })
|
||||
const invalid = ProviderShared.invalidRequest
|
||||
|
||||
const baseUrl = (request: LLMRequest) => (request.model.baseURL ?? "https://api.openai.com/v1").replace(/\/+$/, "")
|
||||
|
||||
@@ -263,13 +262,11 @@ const prepare = Effect.fn("OpenAIChat.prepare")(function* (request: LLMRequest)
|
||||
|
||||
const toHttp = (target: OpenAIChatTarget, request: LLMRequest) =>
|
||||
Effect.succeed(
|
||||
HttpClientRequest.post(`${baseUrl(request)}/chat/completions`).pipe(
|
||||
HttpClientRequest.setHeaders({
|
||||
...request.model.headers,
|
||||
"content-type": "application/json",
|
||||
}),
|
||||
HttpClientRequest.bodyText(encodeTarget(target), "application/json"),
|
||||
),
|
||||
ProviderShared.jsonPost({
|
||||
url: `${baseUrl(request)}/chat/completions`,
|
||||
body: encodeTarget(target),
|
||||
headers: request.model.headers,
|
||||
}),
|
||||
)
|
||||
|
||||
const mapFinishReason = (reason: string | null | undefined): FinishReason => {
|
||||
@@ -371,7 +368,7 @@ export const adapter = Adapter.define<OpenAIChatDraft, OpenAIChatTarget>({
|
||||
protocol: "openai-chat",
|
||||
redact: (target) => target,
|
||||
prepare,
|
||||
validate: (draft) => decodeTarget(draft).pipe(Effect.mapError((error) => invalid(error.message))),
|
||||
validate: ProviderShared.validateWith(decodeTarget),
|
||||
toHttp: (target, context) => toHttp(target, context.request),
|
||||
parse: events,
|
||||
})
|
||||
|
||||
@@ -1,8 +1,7 @@
|
||||
import { Effect, Stream } from "effect"
|
||||
import { HttpClientRequest } from "effect/unstable/http"
|
||||
import { Adapter } from "../adapter"
|
||||
import { capabilities, model as llmModel, type ModelInput } from "../llm"
|
||||
import { InvalidRequestError, ProviderChunkError, type LLMError, type LLMRequest } from "../schema"
|
||||
import { ProviderChunkError, type LLMError, type LLMRequest } from "../schema"
|
||||
import { OpenAIChat, type OpenAIChatTarget } from "./openai-chat"
|
||||
import { families, type ProviderFamily } from "./openai-compatible-family"
|
||||
import { ProviderShared } from "./shared"
|
||||
@@ -20,7 +19,7 @@ export type ProviderFamilyModelInput = Omit<OpenAICompatibleChatModelInput, "pro
|
||||
readonly baseURL?: string
|
||||
}
|
||||
|
||||
const invalid = (message: string) => new InvalidRequestError({ message })
|
||||
const invalid = ProviderShared.invalidRequest
|
||||
|
||||
const isStringRecord = (value: unknown): value is Record<string, string> =>
|
||||
typeof value === "object" && value !== null && !Array.isArray(value) && Object.values(value).every((item) => typeof item === "string")
|
||||
@@ -42,14 +41,11 @@ const toHttp = (target: OpenAIChatTarget, request: LLMRequest) =>
|
||||
Effect.gen(function* () {
|
||||
const url = completionUrl(request)
|
||||
if (!url) return yield* invalid("OpenAI-compatible Chat requires a baseURL")
|
||||
|
||||
return HttpClientRequest.post(url).pipe(
|
||||
HttpClientRequest.setHeaders({
|
||||
...request.model.headers,
|
||||
"content-type": "application/json",
|
||||
}),
|
||||
HttpClientRequest.bodyText(ProviderShared.encodeJson(target), "application/json"),
|
||||
)
|
||||
return ProviderShared.jsonPost({
|
||||
url,
|
||||
body: ProviderShared.encodeJson(target),
|
||||
headers: request.model.headers,
|
||||
})
|
||||
})
|
||||
|
||||
const mapParseError = (error: LLMError) => {
|
||||
|
||||
@@ -1,9 +1,8 @@
|
||||
import { Effect, Schema, Stream } from "effect"
|
||||
import { HttpClientRequest, type HttpClientResponse } from "effect/unstable/http"
|
||||
import type { HttpClientResponse } from "effect/unstable/http"
|
||||
import { Adapter } from "../adapter"
|
||||
import { capabilities, model as llmModel, type ModelInput } from "../llm"
|
||||
import {
|
||||
InvalidRequestError,
|
||||
Usage,
|
||||
type FinishReason,
|
||||
type LLMEvent,
|
||||
@@ -149,7 +148,7 @@ interface ParserState {
|
||||
readonly tools: Record<string, ToolAccumulator>
|
||||
}
|
||||
|
||||
const invalid = (message: string) => new InvalidRequestError({ message })
|
||||
const invalid = ProviderShared.invalidRequest
|
||||
|
||||
const baseUrl = (request: LLMRequest) => (request.model.baseURL ?? "https://api.openai.com/v1").replace(/\/+$/, "")
|
||||
|
||||
@@ -239,13 +238,11 @@ const prepare = Effect.fn("OpenAIResponses.prepare")(function* (request: LLMRequ
|
||||
|
||||
const toHttp = (target: OpenAIResponsesTarget, request: LLMRequest) =>
|
||||
Effect.succeed(
|
||||
HttpClientRequest.post(`${baseUrl(request)}/responses`).pipe(
|
||||
HttpClientRequest.setHeaders({
|
||||
...request.model.headers,
|
||||
"content-type": "application/json",
|
||||
}),
|
||||
HttpClientRequest.bodyText(encodeTarget(target), "application/json"),
|
||||
),
|
||||
ProviderShared.jsonPost({
|
||||
url: `${baseUrl(request)}/responses`,
|
||||
body: encodeTarget(target),
|
||||
headers: request.model.headers,
|
||||
}),
|
||||
)
|
||||
|
||||
const mapUsage = (usage: OpenAIResponsesUsage | undefined) => {
|
||||
@@ -396,7 +393,7 @@ export const adapter = Adapter.define<OpenAIResponsesDraft, OpenAIResponsesTarge
|
||||
protocol: "openai-responses",
|
||||
redact: (target) => target,
|
||||
prepare,
|
||||
validate: (draft) => decodeTarget(draft).pipe(Effect.mapError((error) => invalid(error.message))),
|
||||
validate: ProviderShared.validateWith(decodeTarget),
|
||||
toHttp: (target, context) => toHttp(target, context.request),
|
||||
parse: events,
|
||||
})
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
import { Cause, Effect, Schema, Stream } from "effect"
|
||||
import * as Sse from "effect/unstable/encoding/Sse"
|
||||
import type { HttpClientResponse } from "effect/unstable/http"
|
||||
import { ProviderChunkError } from "../schema"
|
||||
import { HttpClientRequest, type HttpClientResponse } from "effect/unstable/http"
|
||||
import { InvalidRequestError, ProviderChunkError } from "../schema"
|
||||
|
||||
export const Json = Schema.fromJsonString(Schema.Unknown)
|
||||
export const decodeJson = Schema.decodeUnknownSync(Json)
|
||||
@@ -114,4 +114,41 @@ export const sse = <Chunk, State, Event>(input: {
|
||||
readonly onHalt?: (state: State) => ReadonlyArray<Event>
|
||||
}): Stream.Stream<Event, ProviderChunkError> => framed({ ...input, framing: sseFraming })
|
||||
|
||||
/**
|
||||
* Canonical `InvalidRequestError` constructor. Lift one-line `const invalid =
|
||||
* (message) => new InvalidRequestError({ message })` aliases out of every
|
||||
* adapter so the error constructor lives in one place. If we ever extend
|
||||
* `InvalidRequestError` with adapter context or trace metadata, the change
|
||||
* lands here.
|
||||
*/
|
||||
export const invalidRequest = (message: string) => new InvalidRequestError({ message })
|
||||
|
||||
/**
|
||||
* Build a `validate` step from a Schema decoder. Replaces the per-adapter
|
||||
* lambda body `(draft) => decode(draft).pipe(Effect.mapError((e) =>
|
||||
* invalid(e.message)))`. Any decode error is translated into
|
||||
* `InvalidRequestError` carrying the original parse-error message.
|
||||
*/
|
||||
export const validateWith =
|
||||
<A, I, E extends { readonly message: string }>(decode: (input: I) => Effect.Effect<A, E>) =>
|
||||
(draft: I) =>
|
||||
decode(draft).pipe(Effect.mapError((error) => invalidRequest(error.message)))
|
||||
|
||||
/**
|
||||
* Build an HTTP POST with a JSON body. Sets `content-type: application/json`
|
||||
* automatically (callers can't override it — every adapter today places it
|
||||
* last so caller headers win on everything else) and merges caller-supplied
|
||||
* headers. The body is passed pre-encoded so adapters can choose between
|
||||
* `Schema.encodeSync(target)` and `ProviderShared.encodeJson(target)`.
|
||||
*/
|
||||
export const jsonPost = (input: {
|
||||
readonly url: string
|
||||
readonly body: string
|
||||
readonly headers?: Record<string, string>
|
||||
}) =>
|
||||
HttpClientRequest.post(input.url).pipe(
|
||||
HttpClientRequest.setHeaders({ ...input.headers, "content-type": "application/json" }),
|
||||
HttpClientRequest.bodyText(input.body, "application/json"),
|
||||
)
|
||||
|
||||
export * as ProviderShared from "./shared"
|
||||
|
||||
Reference in New Issue
Block a user