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), Extra: c.ExtraContent, }) } 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"` // ExtraContent is the provider's own metadata on the call, round-tripped // as raw JSON. See ToolCall.Extra for why it is not optional. ExtraContent json.RawMessage `json:"extra_content,omitempty"` } 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}, ExtraContent: c.Extra, }) } 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 { // Gemini wraps its error in a one-element ARRAY — `[{"error":{...}}]` — // where OpenAI, Groq and the rest send the object bare. Unwrapped here // rather than tolerated as "no detail", because the detail is the whole // value of the field: for two weeks the trajectory said only "the model // rejected the request" when the body said "Function call is missing a // thought_signature", and the difference was a day of diagnosis. body = bytes.TrimSpace(body) if bytes.HasPrefix(body, []byte("[")) { var many []json.RawMessage if err := json.Unmarshal(body, &many); err != nil || len(many) == 0 { return "" } body = many[0] } 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 } // Provider metadata arrives whole on one fragment, like the id. Kept // when non-empty so a later empty fragment does not erase it. if len(d.ExtraContent) > 0 { call.ExtraContent = d.ExtraContent } // Arguments are the fragmented field: concatenated, never replaced. call.Function.Arguments += d.Function.Arguments } }