The platform could only talk to one vendor. Moving off Claude — for cost, or
because a client asks for Gemini — meant a rewrite behind an interface that
already had exactly the right shape and one implementation.
`openai` is not only OpenAI. Groq, Gemini's compatibility endpoint, OpenRouter,
Together, vLLM and a local Ollama all serve the chat-completions shape, so one
implementation reaches all of them and the difference between them is a base
URL and three model ids. That is why this is one file and not a package per
vendor.
`routing.go` had the vendor baked into the routing table every provider has to
read: effort was `anthropic.OutputConfigEffort`. Nothing was wrong with that
while there was one implementation; it became wrong the moment there were two,
because the OpenAI path would have had to import the Anthropic SDK to learn how
hard to think. Effort is now the platform's own three-value vocabulary and each
implementation maps it onto whatever its API calls the same idea.
THE ACCOUNTING DIFFERS BETWEEN THE TWO WIRES, and getting it wrong would have
been invisible. OpenAI reports prompt_tokens INCLUSIVE of the cached prefix;
Anthropic reports input tokens EXCLUSIVE of it and carries the cache
separately. Usage.Total() adds all four fields, so copying both numbers across
verbatim bills the cached prefix twice — worst on long conversations, which is
exactly where I3's budget matters most. The run would still answer; it would
just hit BudgetExceeded early, for no visible reason. normalise() subtracts,
and there is a test named after it.
Streamed tool calls are keyed by their wire index, not appended in arrival
order. Providers interleave the fragments of parallel calls, so appending
splices one call's arguments onto another's — and the result is usually two
calls that are each valid JSON and both wrong, which means the tools run with
inputs the model never chose and nothing errors. Mutation-checked: ignoring the
index produces `{"day"{"week":"friday"}:"next"}` and the test catches it.
Three configuration mistakes are refused at startup rather than at runtime:
- MODEL_BASE_URL without MODEL_PROVIDER=openai. The anthropic path has one
endpoint and ignores the field, so this is a deployment that believes it
switched providers and did not — every run still goes to Anthropic and is
still billed there, with nothing in the logs to say so. Cost is the whole
reason this change exists, and that is the one mistake that silently
defeats it.
- An unrecognised MODEL_PROVIDER, once at boot instead of once per run.
- A production deployment with no credential — except against localhost,
which needs none, and demanding one would make the free local path
impossible to configure.
reasoning_effort is opt-in via MODEL_REASONING_EFFORT. Reasoning models accept
it; most others reject the entire request with a 400 rather than ignoring an
unknown key, so every deployment would have had to opt out instead.
`make eval-live` now reads the same environment the service does and logs which
provider answered, because a suite that cannot say which model produced a
result is a suite whose result cannot be compared with another run's. That is
the point of this change: §12 leaves model hosting open, and this makes the
decision cheap to reverse and possible to settle on evidence. Weigh the I7 case
heaviest — a cheaper model that follows the planted injection is a security
regression, not a saving.
Default behaviour is unchanged: MODEL_PROVIDER unset means anthropic, and
ANTHROPIC_API_KEY still works, so no existing deployment needs an edit.
NOT verified against a live provider — no credential was available on this
machine. Tested against a fake endpoint covering both paths, and the three
guarantees above are mutation-checked.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01PJvibeSc1JYXjatankqM1g
751 lines
25 KiB
Go
751 lines
25 KiB
Go
package gateway
|
|
|
|
import (
|
|
"bufio"
|
|
"bytes"
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"net/http"
|
|
"strings"
|
|
"time"
|
|
)
|
|
|
|
// OpenAIGateway calls any service that speaks the OpenAI chat-completions API.
|
|
//
|
|
// ONE IMPLEMENTATION, MANY PROVIDERS. Groq, Gemini (through its compatibility
|
|
// endpoint), OpenRouter, Together, vLLM and a local Ollama all serve this same
|
|
// shape, so the difference between them is a base URL and a model id — not a
|
|
// package each. That is the whole reason this file exists: the platform needed
|
|
// a way off a single vendor's pricing without a rewrite per alternative.
|
|
//
|
|
// Hand-rolled over net/http rather than an SDK, per §10. The surface actually
|
|
// used here is one endpoint and one event stream; a dependency for that buys a
|
|
// version to keep current and a second opinion about retries, and this package
|
|
// already has its own.
|
|
type OpenAIGateway struct {
|
|
cfg Config
|
|
http *http.Client
|
|
}
|
|
|
|
// Compile-time proof that this satisfies the boundary and can stream.
|
|
var (
|
|
_ Gateway = (*OpenAIGateway)(nil)
|
|
_ Streamer = (*OpenAIGateway)(nil)
|
|
)
|
|
|
|
// DefaultOpenAIBaseURL is where an unconfigured deployment points.
|
|
const DefaultOpenAIBaseURL = "https://api.openai.com/v1"
|
|
|
|
// openAIHTTPTimeout bounds a single call at the transport.
|
|
//
|
|
// Above the deepest tier's deadline on purpose. The run's own context is what
|
|
// should end a slow call — that failure is a Deadline the runtime can report
|
|
// against a budget — and a transport timeout firing first would present the
|
|
// same event as an unexplained upstream error instead.
|
|
const openAIHTTPTimeout = 10 * time.Minute
|
|
|
|
// NewOpenAI builds a gateway over an OpenAI-compatible service.
|
|
//
|
|
// A missing key is not an error here, for the same reason it is not one for
|
|
// Anthropic: the service has to boot without model credentials, and the
|
|
// failure belongs at the first Complete as a structured NotConfigured a run
|
|
// can end with. A local Ollama legitimately needs no key at all, which is why
|
|
// the check is deferred rather than dropped — see complete().
|
|
func NewOpenAI(cfg Config) *OpenAIGateway {
|
|
return &OpenAIGateway{cfg: cfg, http: &http.Client{Timeout: openAIHTTPTimeout}}
|
|
}
|
|
|
|
// endpoint is the chat-completions URL for this deployment.
|
|
func (g *OpenAIGateway) endpoint() string {
|
|
base := strings.TrimRight(strings.TrimSpace(g.cfg.BaseURL), "/")
|
|
if base == "" {
|
|
base = DefaultOpenAIBaseURL
|
|
}
|
|
return base + "/chat/completions"
|
|
}
|
|
|
|
// routing resolves a tier against this gateway's table.
|
|
func (g *OpenAIGateway) routing(t Tier) Routing { return g.cfg.routingFor(t) }
|
|
|
|
// needsCredential reports whether this deployment must present a key.
|
|
//
|
|
// A hosted provider does; a local Ollama does not, and demanding one would
|
|
// make the zero-cost development path impossible to configure. The base URL is
|
|
// the only signal available — a deployment that has pointed this at its own
|
|
// machine has already said the call is not leaving it.
|
|
func (g *OpenAIGateway) needsCredential() bool {
|
|
base := strings.TrimSpace(g.cfg.BaseURL)
|
|
if base == "" {
|
|
return true
|
|
}
|
|
return !strings.Contains(base, "localhost") && !strings.Contains(base, "127.0.0.1")
|
|
}
|
|
|
|
// Complete calls the model, retrying failures that are worth retrying.
|
|
//
|
|
// Same policy as every other provider — see withRetry, which is shared
|
|
// precisely so the two cannot drift.
|
|
func (g *OpenAIGateway) Complete(ctx context.Context, req Request) (*Response, error) {
|
|
return withRetry(ctx, func() (*Response, error) { return g.complete(ctx, req) })
|
|
}
|
|
|
|
// complete is one attempt.
|
|
func (g *OpenAIGateway) complete(ctx context.Context, req Request) (*Response, error) {
|
|
body, err := g.params(req, false)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
httpResp, err := g.post(ctx, body)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer httpResp.Body.Close()
|
|
|
|
raw, err := io.ReadAll(httpResp.Body)
|
|
if err != nil {
|
|
return nil, &Error{Code: CodeUpstream, Message: "the model response could not be read", Cause: err}
|
|
}
|
|
if httpResp.StatusCode >= 400 {
|
|
return nil, translateOpenAI(httpResp.StatusCode, raw)
|
|
}
|
|
|
|
var decoded oaiResponse
|
|
if err := json.Unmarshal(raw, &decoded); err != nil {
|
|
return nil, &Error{
|
|
Code: CodeUpstream,
|
|
Message: "the model returned a response this gateway could not parse",
|
|
Cause: err,
|
|
}
|
|
}
|
|
if len(decoded.Choices) == 0 {
|
|
return nil, &Error{Code: CodeUpstream, Message: "the model returned no choices"}
|
|
}
|
|
|
|
choice := decoded.Choices[0]
|
|
return g.decode(req, decoded.Model, choice.FinishReason, choice.Message, decoded.Usage)
|
|
}
|
|
|
|
// params builds the request body both paths send.
|
|
//
|
|
// Extracted for the same reason the Anthropic path extracts its own: an answer
|
|
// that differed depending on whether it was streamed would be the worst kind of
|
|
// bug to chase, because the transport is the last place anybody looks.
|
|
func (g *OpenAIGateway) params(req Request, stream bool) (*oaiRequest, error) {
|
|
if err := req.Validate(); err != nil {
|
|
return nil, err
|
|
}
|
|
route := g.routing(req.Tier)
|
|
|
|
maxTokens := req.MaxOutputTokens
|
|
if maxTokens <= 0 {
|
|
maxTokens = g.cfg.MaxOutputTokens
|
|
}
|
|
|
|
body := &oaiRequest{
|
|
Model: route.Model,
|
|
Messages: encodeOpenAIMessages(req.System, req.Messages),
|
|
MaxTokens: maxTokens,
|
|
Tools: encodeOpenAITools(req.Tools),
|
|
}
|
|
if g.cfg.SendReasoningEffort {
|
|
body.ReasoningEffort = openAIEffort(route.Effort)
|
|
}
|
|
if stream {
|
|
body.Stream = true
|
|
// Usage is omitted from a stream unless it is asked for, and a call
|
|
// whose cost is unknown is a call the run's budget cannot be charged
|
|
// for. I3 needs every call measured, so this is not optional.
|
|
body.StreamOptions = &oaiStreamOptions{IncludeUsage: true}
|
|
}
|
|
return body, nil
|
|
}
|
|
|
|
// post sends the request body.
|
|
func (g *OpenAIGateway) post(ctx context.Context, body *oaiRequest) (*http.Response, error) {
|
|
if g.cfg.APIKey == "" && g.needsCredential() {
|
|
return nil, &Error{
|
|
Code: CodeNotConfigured,
|
|
Message: "no model credentials are configured for this deployment",
|
|
}
|
|
}
|
|
|
|
encoded, err := json.Marshal(body)
|
|
if err != nil {
|
|
return nil, &Error{Code: CodeInvalidRequest, Message: "the request could not be encoded", Cause: err}
|
|
}
|
|
|
|
httpReq, err := http.NewRequestWithContext(ctx, http.MethodPost, g.endpoint(), bytes.NewReader(encoded))
|
|
if err != nil {
|
|
return nil, &Error{Code: CodeInvalidRequest, Message: "the request could not be built", Cause: err}
|
|
}
|
|
httpReq.Header.Set("Content-Type", "application/json")
|
|
if g.cfg.APIKey != "" {
|
|
httpReq.Header.Set("Authorization", "Bearer "+g.cfg.APIKey)
|
|
}
|
|
|
|
resp, err := g.http.Do(httpReq)
|
|
if err != nil {
|
|
if errors.Is(err, context.DeadlineExceeded) || errors.Is(err, context.Canceled) {
|
|
return nil, &Error{Code: CodeTimeout, Message: "the model call did not complete in time", Cause: err}
|
|
}
|
|
return nil, &Error{Code: CodeUpstream, Message: "the model call failed", Cause: err}
|
|
}
|
|
return resp, nil
|
|
}
|
|
|
|
// decode turns a finished choice into a Response.
|
|
//
|
|
// Shared by both paths, so a streamed answer and a non-streamed one are read
|
|
// by the same code rather than by two implementations of the same reading.
|
|
func (g *OpenAIGateway) decode(
|
|
req Request, model, finish string, msg oaiMessage, usage oaiUsage,
|
|
) (*Response, error) {
|
|
route := g.routing(req.Tier)
|
|
if model == "" {
|
|
model = route.Model
|
|
}
|
|
|
|
counted := usage.normalise()
|
|
|
|
// A refusal arrives as a successful HTTP response, so it is checked before
|
|
// the content is read. It is still billed, and the usage rides on the
|
|
// Response rather than being dropped — a refusal that cost nothing on the
|
|
// ledger is a refusal the loop would happily repeat.
|
|
if refusal := strings.TrimSpace(msg.Refusal); refusal != "" || finish == "content_filter" {
|
|
category := finish
|
|
if refusal != "" {
|
|
category = "refusal"
|
|
}
|
|
return &Response{
|
|
StopReason: openAIStopReason(finish),
|
|
Usage: counted,
|
|
Model: model,
|
|
Tier: req.Tier,
|
|
}, &Error{
|
|
Code: CodeRefused,
|
|
Message: "the model declined this request",
|
|
Category: category,
|
|
}
|
|
}
|
|
|
|
var calls []ToolCall
|
|
for _, c := range msg.ToolCalls {
|
|
args := strings.TrimSpace(c.Function.Arguments)
|
|
if args == "" {
|
|
// An argumentless call is legitimate; an empty string is not valid
|
|
// JSON, and the handler's decoder would reject it for a reason that
|
|
// has nothing to do with the caller's request.
|
|
args = "{}"
|
|
}
|
|
calls = append(calls, ToolCall{
|
|
ID: c.ID,
|
|
Name: c.Function.Name,
|
|
// The raw JSON, not a parsed value — handed to the handler's own
|
|
// decoder rather than matched on as a string here.
|
|
Input: json.RawMessage(args),
|
|
})
|
|
}
|
|
|
|
return &Response{
|
|
Text: msg.Content,
|
|
ToolCalls: calls,
|
|
StopReason: openAIStopReason(finish),
|
|
Usage: counted,
|
|
Model: model,
|
|
Tier: req.Tier,
|
|
}, nil
|
|
}
|
|
|
|
/* ── Wire types ─────────────────────────────────────────────────────────── */
|
|
|
|
type oaiRequest struct {
|
|
Model string `json:"model"`
|
|
Messages []oaiMessage `json:"messages"`
|
|
Tools []oaiTool `json:"tools,omitempty"`
|
|
MaxTokens int64 `json:"max_tokens,omitempty"`
|
|
Stream bool `json:"stream,omitempty"`
|
|
StreamOptions *oaiStreamOptions `json:"stream_options,omitempty"`
|
|
|
|
// ReasoningEffort is omitted unless a deployment opted in. Most non-
|
|
// reasoning models reject the whole request rather than ignoring the key.
|
|
ReasoningEffort string `json:"reasoning_effort,omitempty"`
|
|
}
|
|
|
|
type oaiStreamOptions struct {
|
|
IncludeUsage bool `json:"include_usage"`
|
|
}
|
|
|
|
// oaiMessage is one wire message. It doubles as a streamed delta, because the
|
|
// two carry the same fields and differ only in how much of each is present.
|
|
type oaiMessage struct {
|
|
Role string `json:"role,omitempty"`
|
|
Content string `json:"content,omitempty"`
|
|
Refusal string `json:"refusal,omitempty"`
|
|
ToolCalls []oaiToolCall `json:"tool_calls,omitempty"`
|
|
|
|
// ToolCallID is set only on a role:"tool" message, correlating a result
|
|
// with the call that asked for it.
|
|
ToolCallID string `json:"tool_call_id,omitempty"`
|
|
}
|
|
|
|
type oaiToolCall struct {
|
|
// Index orders a call within a streamed response. Absent when complete,
|
|
// which is why it is a pointer: index 0 and "no index" are different
|
|
// things, and reading a missing field as 0 merges every streamed call
|
|
// into the first one.
|
|
Index *int `json:"index,omitempty"`
|
|
ID string `json:"id,omitempty"`
|
|
Type string `json:"type,omitempty"`
|
|
Function oaiFunctionRef `json:"function"`
|
|
}
|
|
|
|
type oaiFunctionRef struct {
|
|
Name string `json:"name,omitempty"`
|
|
Arguments string `json:"arguments,omitempty"`
|
|
}
|
|
|
|
type oaiTool struct {
|
|
Type string `json:"type"`
|
|
Function oaiFunctionDef `json:"function"`
|
|
}
|
|
|
|
type oaiFunctionDef struct {
|
|
Name string `json:"name"`
|
|
Description string `json:"description,omitempty"`
|
|
Parameters map[string]any `json:"parameters,omitempty"`
|
|
}
|
|
|
|
type oaiResponse struct {
|
|
Model string `json:"model"`
|
|
Choices []oaiChoice `json:"choices"`
|
|
Usage oaiUsage `json:"usage"`
|
|
}
|
|
|
|
type oaiChoice struct {
|
|
Message oaiMessage `json:"message"`
|
|
Delta oaiMessage `json:"delta"`
|
|
FinishReason string `json:"finish_reason"`
|
|
}
|
|
|
|
type oaiUsage struct {
|
|
PromptTokens int64 `json:"prompt_tokens"`
|
|
CompletionTokens int64 `json:"completion_tokens"`
|
|
PromptTokensDetails struct {
|
|
CachedTokens int64 `json:"cached_tokens"`
|
|
} `json:"prompt_tokens_details"`
|
|
}
|
|
|
|
// normalise converts OpenAI's accounting into this platform's.
|
|
//
|
|
// THE SUBTRACTION IS THE WHOLE FUNCTION, and getting it wrong would corrupt
|
|
// every budget quietly. OpenAI reports `prompt_tokens` INCLUSIVE of the cached
|
|
// prefix; Anthropic reports input tokens EXCLUSIVE of it, and carries the cache
|
|
// separately. Usage.Total() adds all four fields, so copying both numbers
|
|
// across verbatim would bill the cached prefix twice — and it would do it
|
|
// worst on long conversations, which is exactly where a budget matters most.
|
|
//
|
|
// Clamped at zero rather than trusted: a provider that reports more cached
|
|
// tokens than prompt tokens is wrong, but a negative charge would be a bug
|
|
// that hands a run free budget rather than one that shows up as a wrong number.
|
|
func (u oaiUsage) normalise() Usage {
|
|
cached := u.PromptTokensDetails.CachedTokens
|
|
fresh := u.PromptTokens - cached
|
|
if fresh < 0 {
|
|
fresh = 0
|
|
}
|
|
return Usage{
|
|
InputTokens: fresh,
|
|
OutputTokens: u.CompletionTokens,
|
|
CacheReadTokens: cached,
|
|
// No creation figure on this wire. Left at zero rather than guessed:
|
|
// an invented number is worse than an absent one, because it looks
|
|
// like a measurement.
|
|
CacheCreationTokens: 0,
|
|
}
|
|
}
|
|
|
|
/* ── Encoding ───────────────────────────────────────────────────────────── */
|
|
|
|
// openAIEffort maps the platform's effort vocabulary onto OpenAI's.
|
|
//
|
|
// Three of ours onto three of theirs, preserving the ordering rather than the
|
|
// spelling: their scale runs minimal/low/medium/high, so "high" here is their
|
|
// "medium" and "xhigh" is their "high". Matching the words instead of the
|
|
// positions would have made `fast` and `balanced` nearly indistinguishable.
|
|
func openAIEffort(e Effort) string {
|
|
switch e {
|
|
case EffortLow:
|
|
return "low"
|
|
case EffortXhigh:
|
|
return "high"
|
|
default:
|
|
return "medium"
|
|
}
|
|
}
|
|
|
|
// openAIStopReason maps a finish_reason onto the vocabulary the trajectories
|
|
// already use.
|
|
//
|
|
// Translated rather than passed through, so a trajectory reads the same
|
|
// whichever provider answered. An eval comparing two providers is comparing
|
|
// the run, and it should not have to know that one says "tool_calls" where the
|
|
// other says "tool_use".
|
|
func openAIStopReason(finish string) string {
|
|
switch finish {
|
|
case "tool_calls", "function_call":
|
|
return "tool_use"
|
|
case "stop":
|
|
return "end_turn"
|
|
case "length":
|
|
return "max_tokens"
|
|
case "content_filter":
|
|
return "refusal"
|
|
default:
|
|
return finish
|
|
}
|
|
}
|
|
|
|
// encodeOpenAITools renders the tool definitions for the wire.
|
|
//
|
|
// The whole input schema is passed through, not just its properties: this API
|
|
// validates arguments against what it is given, so dropping `type`, `enum` or
|
|
// a nested object's own required list would let the model send arguments the
|
|
// handler then has to reject.
|
|
func encodeOpenAITools(defs []ToolDef) []oaiTool {
|
|
if len(defs) == 0 {
|
|
return nil
|
|
}
|
|
out := make([]oaiTool, 0, len(defs))
|
|
for _, d := range defs {
|
|
params := d.InputSchema
|
|
if params == nil {
|
|
params = map[string]any{"type": "object", "properties": map[string]any{}}
|
|
} else if _, ok := params["type"]; !ok {
|
|
// A schema without a type is rejected by some providers and
|
|
// silently accepted by others. Copied rather than mutated: the
|
|
// caller's map is shared across every call in a run.
|
|
cloned := make(map[string]any, len(params)+1)
|
|
for k, v := range params {
|
|
cloned[k] = v
|
|
}
|
|
cloned["type"] = "object"
|
|
params = cloned
|
|
}
|
|
out = append(out, oaiTool{
|
|
Type: "function",
|
|
Function: oaiFunctionDef{
|
|
Name: d.Name,
|
|
Description: d.Description,
|
|
Parameters: params,
|
|
},
|
|
})
|
|
}
|
|
return out
|
|
}
|
|
|
|
// encodeOpenAIMessages renders a conversation for the wire.
|
|
//
|
|
// TWO SHAPE DIFFERENCES from the Anthropic path, and both are load-bearing:
|
|
//
|
|
// - The system prompt is a MESSAGE here, not a top-level field, and it must
|
|
// come first.
|
|
// - A tool result is its OWN message with role "tool", one per result —
|
|
// where Anthropic carries them as blocks inside a single user turn. So the
|
|
// grouping the other encoder is careful to preserve has to be undone here,
|
|
// in the same order, or a result arrives detached from its call.
|
|
//
|
|
// Ordering within a turn matters: results are emitted before any text in the
|
|
// same message, because they answer the assistant turn that preceded them.
|
|
func encodeOpenAIMessages(system string, msgs []Message) []oaiMessage {
|
|
out := make([]oaiMessage, 0, len(msgs)+1)
|
|
|
|
if s := strings.TrimSpace(system); s != "" {
|
|
out = append(out, oaiMessage{Role: "system", Content: s})
|
|
}
|
|
|
|
for _, m := range msgs {
|
|
for _, r := range m.ToolResults {
|
|
// IsError has no home on this wire — there is no error flag on a
|
|
// tool message. The handler's own error payload is already in the
|
|
// content, per §4, so the model still sees what went wrong; what
|
|
// is lost is the structured marker, and inventing a prefix for it
|
|
// would put prose in a channel that carries data.
|
|
out = append(out, oaiMessage{
|
|
Role: "tool",
|
|
ToolCallID: r.CallID,
|
|
Content: r.Content,
|
|
})
|
|
}
|
|
|
|
hasText := strings.TrimSpace(m.Text) != ""
|
|
if !hasText && len(m.ToolCalls) == 0 {
|
|
continue
|
|
}
|
|
|
|
msg := oaiMessage{Role: string(m.Role), Content: m.Text}
|
|
for _, c := range m.ToolCalls {
|
|
args := strings.TrimSpace(string(c.Input))
|
|
if args == "" {
|
|
args = "{}"
|
|
}
|
|
msg.ToolCalls = append(msg.ToolCalls, oaiToolCall{
|
|
ID: c.ID,
|
|
Type: "function",
|
|
Function: oaiFunctionRef{Name: c.Name, Arguments: args},
|
|
})
|
|
}
|
|
out = append(out, msg)
|
|
}
|
|
return out
|
|
}
|
|
|
|
/* ── Errors ─────────────────────────────────────────────────────────────── */
|
|
|
|
// translateOpenAI turns an HTTP failure into one the runtime can branch on.
|
|
//
|
|
// Mapped by status, mirroring the Anthropic path, because the distinction the
|
|
// loop needs is the same one either way: whether sending this request again
|
|
// could work. The upstream message is carried through when there is one — a
|
|
// 400 that says which tool schema is malformed is worth more than "the model
|
|
// rejected the request", and the trajectory only records the message.
|
|
func translateOpenAI(status int, body []byte) error {
|
|
detail := openAIErrorMessage(body)
|
|
|
|
withDetail := func(base string) string {
|
|
if detail == "" {
|
|
return base
|
|
}
|
|
return base + ": " + detail
|
|
}
|
|
|
|
switch {
|
|
case status == 400 || status == 404 || status == 422:
|
|
// 404 belongs here, not with the 5xx: on these providers it almost
|
|
// always means the model id does not exist on this endpoint, which is
|
|
// a configuration mistake and will fail identically next time.
|
|
return &Error{Code: CodeInvalidRequest, Message: withDetail("the model rejected the request"), Status: status}
|
|
case status == 401 || status == 403:
|
|
return &Error{Code: CodeUnauthorized, Message: withDetail("the model credentials were refused"), Status: status}
|
|
case status == 408:
|
|
return &Error{Code: CodeTimeout, Message: withDetail("the model call timed out"), Status: status}
|
|
case status == 429:
|
|
return &Error{Code: CodeRateLimited, Message: withDetail("the model is rate limiting this deployment"), Status: status}
|
|
default:
|
|
return &Error{
|
|
Code: CodeUpstream,
|
|
Message: withDetail(fmt.Sprintf("the model call failed (http %d)", status)),
|
|
Status: status,
|
|
}
|
|
}
|
|
}
|
|
|
|
// openAIErrorMessage digs the human-readable reason out of an error body.
|
|
//
|
|
// Best-effort by design: providers agree on the envelope often enough to be
|
|
// worth reading and not often enough to depend on, so an unparseable body
|
|
// yields nothing rather than failing a failure.
|
|
func openAIErrorMessage(body []byte) string {
|
|
var envelope struct {
|
|
Error struct {
|
|
Message string `json:"message"`
|
|
} `json:"error"`
|
|
Message string `json:"message"`
|
|
}
|
|
if err := json.Unmarshal(body, &envelope); err != nil {
|
|
return ""
|
|
}
|
|
if m := strings.TrimSpace(envelope.Error.Message); m != "" {
|
|
return m
|
|
}
|
|
return strings.TrimSpace(envelope.Message)
|
|
}
|
|
|
|
/* ── Streaming ──────────────────────────────────────────────────────────── */
|
|
|
|
// maxSSELine caps a single server-sent-event line.
|
|
//
|
|
// One event carries one delta, but a tool call's arguments arrive as a single
|
|
// field that can be large, and the default scanner limit of 64KB is low enough
|
|
// to be hit by a real request. A cap is still wanted: an unbounded line from a
|
|
// misbehaving upstream would be read straight into memory.
|
|
const maxSSELine = 1 << 20
|
|
|
|
// Stream is Complete, with the assistant's text delivered as it arrives.
|
|
//
|
|
// §6: "Stream partial assistant text as it arrives; buffer tool calls until
|
|
// complete." Both halves matter and they pull in opposite directions.
|
|
//
|
|
// TEXT IS STREAMED because a fifteen-second wait with nothing on screen reads
|
|
// as broken.
|
|
//
|
|
// TOOL CALLS ARE NOT. On this wire a call's arguments arrive as a JSON string
|
|
// assembled across many events, and a half-built argument object is not a
|
|
// smaller version of the finished one — it is a different object, usually an
|
|
// invalid one. So the fragments are accumulated by index and decoded only once
|
|
// the stream closes, by exactly the same code the non-streaming path uses.
|
|
//
|
|
// onDelta is called from this goroutine, in order, and must not block for long
|
|
// — it is on the path between the model and the reader.
|
|
func (g *OpenAIGateway) Stream(ctx context.Context, req Request, onDelta func(string)) (*Response, error) {
|
|
body, err := g.params(req, true)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
httpResp, err := g.post(ctx, body)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer httpResp.Body.Close()
|
|
|
|
if httpResp.StatusCode >= 400 {
|
|
raw, _ := io.ReadAll(httpResp.Body)
|
|
return nil, translateOpenAI(httpResp.StatusCode, raw)
|
|
}
|
|
|
|
acc, err := accumulateSSE(httpResp.Body, onDelta)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return g.decode(req, acc.model, acc.finishReason, acc.message(), acc.usage)
|
|
}
|
|
|
|
// streamAccumulator assembles a streamed response.
|
|
//
|
|
// Tool calls are keyed by their wire index rather than appended in arrival
|
|
// order: providers interleave the fragments of parallel calls, so arrival
|
|
// order is not call order, and appending would splice one call's arguments
|
|
// onto another's.
|
|
type streamAccumulator struct {
|
|
text strings.Builder
|
|
refusal strings.Builder
|
|
model string
|
|
finishReason string
|
|
usage oaiUsage
|
|
|
|
calls map[int]*oaiToolCall
|
|
order []int
|
|
}
|
|
|
|
// message renders the accumulated stream as the finished message the shared
|
|
// decoder reads.
|
|
func (a *streamAccumulator) message() oaiMessage {
|
|
msg := oaiMessage{
|
|
Role: "assistant",
|
|
Content: a.text.String(),
|
|
Refusal: a.refusal.String(),
|
|
}
|
|
for _, idx := range a.order {
|
|
msg.ToolCalls = append(msg.ToolCalls, *a.calls[idx])
|
|
}
|
|
return msg
|
|
}
|
|
|
|
// accumulateSSE reads the event stream to its end.
|
|
func accumulateSSE(r io.Reader, onDelta func(string)) (*streamAccumulator, error) {
|
|
acc := &streamAccumulator{calls: map[int]*oaiToolCall{}}
|
|
|
|
scanner := bufio.NewScanner(r)
|
|
scanner.Buffer(make([]byte, 0, 64*1024), maxSSELine)
|
|
|
|
for scanner.Scan() {
|
|
line := strings.TrimSpace(scanner.Text())
|
|
if line == "" {
|
|
continue
|
|
}
|
|
// Some providers emit "data: {...}", others "data:{...}". Comment
|
|
// lines beginning ":" are keep-alives and carry nothing.
|
|
if !strings.HasPrefix(line, "data:") {
|
|
continue
|
|
}
|
|
payload := strings.TrimSpace(strings.TrimPrefix(line, "data:"))
|
|
if payload == "" || payload == "[DONE]" {
|
|
continue
|
|
}
|
|
|
|
var chunk oaiResponse
|
|
if err := json.Unmarshal([]byte(payload), &chunk); err != nil {
|
|
// One malformed event is not a failed response. Skipping it keeps
|
|
// a keep-alive or a provider-specific event from ending a stream
|
|
// that is otherwise fine.
|
|
continue
|
|
}
|
|
|
|
if chunk.Model != "" {
|
|
acc.model = chunk.Model
|
|
}
|
|
// The usage chunk arrives last and carries no choices. Guarded rather
|
|
// than assumed: a zero usage overwriting a real one would silently
|
|
// hand the run a free turn.
|
|
if chunk.Usage.PromptTokens > 0 || chunk.Usage.CompletionTokens > 0 {
|
|
acc.usage = chunk.Usage
|
|
}
|
|
if len(chunk.Choices) == 0 {
|
|
continue
|
|
}
|
|
|
|
choice := chunk.Choices[0]
|
|
if choice.FinishReason != "" {
|
|
acc.finishReason = choice.FinishReason
|
|
}
|
|
if d := choice.Delta.Content; d != "" {
|
|
acc.text.WriteString(d)
|
|
if onDelta != nil {
|
|
onDelta(d)
|
|
}
|
|
}
|
|
// A refusal is accumulated but never streamed to the reader: it is not
|
|
// the answer, and putting it on screen would show a declined request
|
|
// as though it were one.
|
|
if d := choice.Delta.Refusal; d != "" {
|
|
acc.refusal.WriteString(d)
|
|
}
|
|
acc.addToolCallDeltas(choice.Delta.ToolCalls)
|
|
}
|
|
|
|
if err := scanner.Err(); err != nil {
|
|
return nil, &Error{
|
|
Code: CodeUpstream,
|
|
Message: "the streamed response could not be assembled",
|
|
Cause: err,
|
|
}
|
|
}
|
|
return acc, nil
|
|
}
|
|
|
|
// addToolCallDeltas folds one event's tool-call fragments into the accumulator.
|
|
func (a *streamAccumulator) addToolCallDeltas(deltas []oaiToolCall) {
|
|
for _, d := range deltas {
|
|
idx := 0
|
|
if d.Index != nil {
|
|
idx = *d.Index
|
|
}
|
|
|
|
call, seen := a.calls[idx]
|
|
if !seen {
|
|
call = &oaiToolCall{Type: "function"}
|
|
a.calls[idx] = call
|
|
a.order = append(a.order, idx)
|
|
}
|
|
|
|
// The id and name arrive once, on the opening fragment. Assigned only
|
|
// when non-empty so a later fragment carrying empty strings — which is
|
|
// the common shape — does not erase them.
|
|
if d.ID != "" {
|
|
call.ID = d.ID
|
|
}
|
|
if d.Type != "" {
|
|
call.Type = d.Type
|
|
}
|
|
if d.Function.Name != "" {
|
|
call.Function.Name = d.Function.Name
|
|
}
|
|
// Arguments are the fragmented field: concatenated, never replaced.
|
|
call.Function.Arguments += d.Function.Arguments
|
|
}
|
|
}
|