package gateway import ( "context" "encoding/json" "errors" "fmt" "strings" "time" "github.com/anthropics/anthropic-sdk-go" "github.com/anthropics/anthropic-sdk-go/option" ) // Routing is how a tier becomes a model and an effort level. // // The model per tier is a deployment knob — a tenant on a different contract, // or a deployment pinning a version through an incident, changes it without a // spec edit. The *effort* per tier is not: "fast" and "deep" mean something // specific about how much work an answer is worth, and letting a deployment // redefine that would make the same spec behave differently in two places // while claiming the same tier. type Routing struct { Model string Effort anthropic.OutputConfigEffort } // Config is the gateway's whole configuration surface. // // Built once at startup from the environment and passed in frozen, per §10. // Nothing in this package reads the environment itself. type Config struct { APIKey string Fast Routing Balanced Routing Deep Routing // MaxOutputTokens applies when a request does not set its own. MaxOutputTokens int64 } // AnthropicGateway calls the Claude API. type AnthropicGateway struct { client anthropic.Client cfg Config } // Compile-time proof that this satisfies the boundary. var _ Gateway = (*AnthropicGateway)(nil) // NewAnthropic builds a gateway over the Claude API. // // A missing key is not an error here. The service has to boot without model // credentials — every endpoint that is not an agent run still works, and a // developer running migrations should not need a key to do it. The failure // surfaces at the first Complete, as a structured NotConfigured that the // runtime can end a run with, rather than as a panic at startup. func NewAnthropic(cfg Config) *AnthropicGateway { opts := []option.RequestOption{} if cfg.APIKey != "" { opts = append(opts, option.WithAPIKey(cfg.APIKey)) } return &AnthropicGateway{client: anthropic.NewClient(opts...), cfg: cfg} } // routing resolves a tier. An unknown tier has already been normalised by // ParseTier, so the default arm is reached only by a zero value. func (g *AnthropicGateway) routing(t Tier) Routing { switch t { case TierFast: return g.cfg.Fast case TierDeep: return g.cfg.Deep default: return g.cfg.Balanced } } // Complete sends one request and reports one result. // MaxAttempts is how many times a transient failure is retried. // // Three total, not three retries. Past that the problem is not transient and a // fourth call is just spending money on the same answer. const MaxAttempts = 3 // retryBackoff is the pause before each retry. // // Short, and deliberately so: this sits inside a run that already has a // wall-clock deadline, and a backoff long enough to be polite to the API is // long enough to spend the caller's whole budget waiting. A run that cannot // afford the wait dies on its deadline instead, which is the correct failure. var retryBackoff = []time.Duration{400 * time.Millisecond, 1200 * time.Millisecond} // Complete calls the model, retrying failures that are worth retrying. // // THE RETRY IS NOT DEFENSIVE POLISH. Error.Retryable() has existed since this // package was written and had ZERO callers — the classification was built and // never used, so a 529 "overloaded" killed a run that would have succeeded four // hundred milliseconds later. Found by a real overload during live testing, // where it presented as "the agent could not finish" with nothing to act on. // // Only genuinely transient failures qualify: rate limits, timeouts, and 5xx. // A 400 is a malformed request and will be malformed again; a 401 is a bad // credential and retrying it three times just gets refused three times. // // The run's context governs. A retry that would outlive the caller's deadline // does not happen — the deadline belongs to the run, not to this function, and // waiting past it would turn a bounded run into an unbounded one. func (g *AnthropicGateway) Complete(ctx context.Context, req Request) (*Response, error) { var last error for attempt := 0; attempt < MaxAttempts; attempt++ { if attempt > 0 { pause := retryBackoff[min(attempt-1, len(retryBackoff)-1)] select { case <-time.After(pause): case <-ctx.Done(): // Out of time. The ORIGINAL failure is returned rather than the // context error: "the model was overloaded" is what an operator // needs to see, and "context deadline exceeded" would hide it. return nil, last } } resp, err := g.complete(ctx, req) if err == nil { return resp, nil } last = err var gwErr *Error if !errors.As(err, &gwErr) || !gwErr.Retryable() { return resp, err } } return nil, last } // complete is one attempt. func (g *AnthropicGateway) complete(ctx context.Context, req Request) (*Response, error) { if g.cfg.APIKey == "" { return nil, &Error{ Code: CodeNotConfigured, Message: "no model credentials are configured for this deployment", } } if err := req.Validate(); err != nil { return nil, err } params, err := g.params(req) if err != nil { return nil, err } msg, err := g.client.Messages.New(ctx, params) if err != nil { return nil, translate(err) } return g.decode(msg, req) } // params builds the request both paths send. // // Extracted so the streaming and non-streaming calls cannot drift. They send // the same model, the same thinking config, the same cache breakpoint and the // same tools — an answer that differs 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 *AnthropicGateway) params(req Request) (anthropic.MessageNewParams, error) { route := g.routing(req.Tier) maxTokens := req.MaxOutputTokens if maxTokens <= 0 { maxTokens = g.cfg.MaxOutputTokens } messages, err := encodeMessages(req.Messages) if err != nil { return anthropic.MessageNewParams{}, err } params := anthropic.MessageNewParams{ Model: anthropic.Model(route.Model), MaxTokens: maxTokens, Messages: messages, // Adaptive thinking on every tier: the model decides how much to think, // and effort sets the ceiling on that. A fixed token budget for // reasoning is the deprecated shape and is rejected outright by the // current models. Thinking: anthropic.ThinkingConfigParamUnion{ OfAdaptive: &anthropic.ThinkingConfigAdaptiveParam{}, }, OutputConfig: anthropic.OutputConfigParam{Effort: route.Effort}, } if len(req.Tools) > 0 { params.Tools = encodeTools(req.Tools) } if s := strings.TrimSpace(req.System); s != "" { // One cached block. The system prompt is the stable prefix of every // turn in a run, and the render order is tools → system → messages, so // a breakpoint here is the one that survives the conversation growing. params.System = []anthropic.TextBlockParam{{ Text: s, CacheControl: anthropic.NewCacheControlEphemeralParam(), }} } return params, nil } // decode turns a finished message into a Response. // // Shared by both paths for the same reason params() is: a streamed message and // a non-streamed one are the same object by the time they get here, and reading // them differently would make streaming a second implementation of the answer. func (g *AnthropicGateway) decode(msg *anthropic.Message, req Request) (*Response, error) { route := g.routing(req.Tier) usage := Usage{ InputTokens: msg.Usage.InputTokens, OutputTokens: msg.Usage.OutputTokens, CacheReadTokens: msg.Usage.CacheReadInputTokens, CacheCreationTokens: msg.Usage.CacheCreationInputTokens, } // A refusal arrives as a successful HTTP response, so it has to be checked // before the content is read. It is still billed, and the usage is carried // on the error so the run's budget is charged for a turn that produced no // text — a refusal that cost nothing on the ledger is a refusal the loop // would happily repeat. if msg.StopReason == anthropic.StopReasonRefusal { return &Response{ StopReason: string(msg.StopReason), Usage: usage, Model: route.Model, Tier: req.Tier, }, &Error{ Code: CodeRefused, Message: "the model declined this request", Category: string(msg.StopDetails.Category), } } var ( text strings.Builder calls []ToolCall ) for _, block := range msg.Content { switch b := block.AsAny().(type) { case anthropic.TextBlock: text.WriteString(b.Text) case anthropic.ToolUseBlock: // The raw JSON, not a parsed value: current models vary their // string escaping inside tool inputs, so this is handed to the // handler's own decoder rather than matched on as a string here. calls = append(calls, ToolCall{ ID: b.ID, Name: b.Name, Input: json.RawMessage(b.JSON.Input.Raw()), }) } } return &Response{ Text: text.String(), ToolCalls: calls, StopReason: string(msg.StopReason), Usage: usage, Model: route.Model, Tier: req.Tier, }, nil } // encodeTools renders the tool definitions for the wire. func encodeTools(defs []ToolDef) []anthropic.ToolUnionParam { out := make([]anthropic.ToolUnionParam, 0, len(defs)) for _, d := range defs { schema := anthropic.ToolInputSchemaParam{} if props, ok := d.InputSchema["properties"].(map[string]any); ok { schema.Properties = props } if req, ok := d.InputSchema["required"].([]string); ok { schema.Required = req } tool := anthropic.ToolParam{ Name: d.Name, Description: anthropic.String(d.Description), InputSchema: schema, } out = append(out, anthropic.ToolUnionParam{OfTool: &tool}) } return out } // encodeMessages renders a conversation for the wire. // // Tool results are variadic within ONE user message. Splitting them across // several messages is accepted by the API and quietly teaches the model to stop // making parallel calls — a performance regression with no error to trace it // to, so the grouping is done here rather than left to callers. func encodeMessages(msgs []Message) ([]anthropic.MessageParam, error) { out := make([]anthropic.MessageParam, 0, len(msgs)) for i, m := range msgs { var blocks []anthropic.ContentBlockParamUnion if s := strings.TrimSpace(m.Text); s != "" { blocks = append(blocks, anthropic.NewTextBlock(m.Text)) } for _, c := range m.ToolCalls { var input any if len(c.Input) > 0 { if err := json.Unmarshal(c.Input, &input); err != nil { return nil, &Error{ Code: CodeInvalidRequest, Message: fmt.Sprintf("messages[%d]: tool call %s carries invalid JSON", i, c.Name), } } } blocks = append(blocks, anthropic.NewToolUseBlock(c.ID, input, c.Name)) } for _, r := range m.ToolResults { blocks = append(blocks, anthropic.NewToolResultBlock(r.CallID, r.Content, r.IsError)) } if len(blocks) == 0 { continue } if m.Role == RoleAssistant { out = append(out, anthropic.NewAssistantMessage(blocks...)) continue } out = append(out, anthropic.NewUserMessage(blocks...)) } return out, nil } // translate turns an SDK error into one the runtime can branch on. // // A single broad class would lose the distinction the loop actually needs: // whether sending the same request again could work. So the status is read and // mapped, and anything unrecognised stays CodeUpstream with its status intact // rather than being flattened into a generic failure. func translate(err error) error { if errors.Is(err, context.DeadlineExceeded) || errors.Is(err, context.Canceled) { return &Error{Code: CodeTimeout, Message: "the model call did not complete in time", Cause: err} } var apierr *anthropic.Error if !errors.As(err, &apierr) { return &Error{Code: CodeUpstream, Message: "the model call failed", Cause: err} } switch apierr.StatusCode { case 400: return &Error{Code: CodeInvalidRequest, Message: "the model rejected the request", Status: 400, Cause: err} case 401, 403: return &Error{Code: CodeUnauthorized, Message: "the model credentials were refused", Status: apierr.StatusCode, Cause: err} case 408: return &Error{Code: CodeTimeout, Message: "the model call timed out", Status: 408, Cause: err} case 429: return &Error{Code: CodeRateLimited, Message: "the model is rate limiting this deployment", Status: 429, Cause: err} case 529: // Anthropic's "overloaded" — the service is up and temporarily out of // capacity. Named separately from the 500s because it is the one that // actually happens, and because a run dying on it is a run that would // have succeeded a second later. return &Error{Code: CodeUpstream, Message: "the model is temporarily overloaded", Status: 529, Cause: err} default: // The status is IN the message, not only in the field. It cost an hour // of debugging to learn that "the model call failed" was a 529 rather // than a malformed tool schema, and the trajectory only records the // message. return &Error{ Code: CodeUpstream, Message: fmt.Sprintf("the model call failed (http %d)", apierr.StatusCode), Status: apierr.StatusCode, Cause: err, } } } /* ── Streaming ──────────────────────────────────────────────────────────── */ // 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 of that matter and they pull in opposite directions. // // TEXT IS STREAMED because a fifteen-second wait with nothing on screen reads // as broken. The reader wants the first sentence while the rest is still being // written, and that is the whole difference between a product and a spinner. // // TOOL CALLS ARE NOT. A tool call arrives as JSON assembled character by // character 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. Dispatching on a partial call would run a tool with arguments // the model had not finished choosing. So the accumulated message is 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 *AnthropicGateway) Stream(ctx context.Context, req Request, onDelta func(string)) (*Response, error) { if g.cfg.APIKey == "" { return nil, &Error{ Code: CodeNotConfigured, Message: "no model credentials are configured for this deployment", } } if err := req.Validate(); err != nil { return nil, err } params, err := g.params(req) if err != nil { return nil, err } stream := g.client.Messages.NewStreaming(ctx, params) defer stream.Close() var msg anthropic.Message for stream.Next() { event := stream.Current() if err := msg.Accumulate(event); err != nil { return nil, &Error{ Code: CodeUpstream, Message: "the streamed response could not be assembled", Cause: err, } } // Text only. A thinking delta is the model's private reasoning and is // not the answer; a tool-input delta is a fragment of JSON. Neither is // something to put in front of a reader. if event.Type == "content_block_delta" && event.Delta.Type == "text_delta" { if d := event.Delta.Text; d != "" && onDelta != nil { onDelta(d) } } } if err := stream.Err(); err != nil { return nil, translate(err) } return g.decode(&msg, req) } // StreamComplete runs a request through whichever path the gateway supports. // // A gateway that cannot stream is not a broken gateway — every fake in the test // suite is one, and so is any future provider without a streaming API. Falling // back to Complete and delivering the finished text as a single delta keeps the // caller's code identical either way, which is what stops streaming from // becoming a second code path through the loop. func StreamComplete(ctx context.Context, gw Gateway, req Request, onDelta func(string)) (*Response, error) { // Normalised once, here, so no implementation has to guard it. A caller // that does not want deltas passes nil — every eval and every test does — // and an implementation that took that literally would panic on the first // fragment. Making each Streamer remember the check is how one of them // eventually forgets. if onDelta == nil { onDelta = func(string) {} } if s, ok := gw.(Streamer); ok { return s.Stream(ctx, req, onDelta) } resp, err := gw.Complete(ctx, req) if err == nil && resp != nil && resp.Text != "" && onDelta != nil { onDelta(resp.Text) } return resp, err }