478 lines
16 KiB
Go
478 lines
16 KiB
Go
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
|
|
}
|