Files
krow_backend/go-api/internal/gateway/anthropic.go
Aravind cf99866e12
Some checks failed
CI / test (push) Failing after 4m38s
CI / fixture (push) Failing after 9s
create employee table
2026-09-05 10:44:47 +05:30

511 lines
18 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"
)
// 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 against this gateway's table.
func (g *AnthropicGateway) routing(t Tier) Routing { return g.cfg.routingFor(t) }
// sdkEffort maps the platform's effort vocabulary onto Anthropic's.
//
// A one-to-one mapping today, which is exactly why the neutral type is worth
// having: the platform's three levels are a statement about how much a turn is
// worth, and this function is where that statement meets one vendor's spelling
// of it. `max` is not reachable — see FromConfig.
func sdkEffort(e Effort) anthropic.OutputConfigEffort {
switch e {
case EffortLow:
return anthropic.OutputConfigEffortLow
case EffortXhigh:
return anthropic.OutputConfigEffortXhigh
default:
return anthropic.OutputConfigEffortHigh
}
}
// 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) {
return withRetry(ctx, func() (*Response, error) { return g.complete(ctx, req) })
}
// withRetry runs one attempt until it succeeds, fails terminally, or runs out
// of attempts.
//
// SHARED BY EVERY PROVIDER, and it has to be. The retry policy is a property of
// this platform's runs — bounded attempts, short backoff, the caller's deadline
// winning — not of any one vendor's API. Left as a method, the second provider
// would have grown its own copy, and the two would have drifted the first time
// either was tuned.
func withRetry(ctx context.Context, once func() (*Response, error)) (*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 := once()
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: sdkEffort(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}
}
// The upstream reason, carried through when there is one — the same rule the
// OpenAI path already follows, and for the same reason it gives: a 400 that
// says which model id or tool schema was rejected is worth more than "the
// model rejected the request", and the trajectory records only the message.
//
// This was not a hypothetical. Every run on a deployment failed with a bare
// `gateway.invalid_request`, and neither the run detail nor anything a
// client could read said why — the one fact needed to fix it was discarded
// here, three lines from where it arrived. `Cause` keeps the full error for
// a Go caller; nothing reads it by the time a run is persisted.
detail := anthropicErrorMessage(apierr.RawJSON())
withDetail := func(base string) string {
if detail == "" {
return base
}
return base + ": " + detail
}
switch apierr.StatusCode {
case 400:
return &Error{Code: CodeInvalidRequest, Message: withDetail("the model rejected the request"), Status: 400, Cause: err}
case 401, 403:
return &Error{Code: CodeUnauthorized, Message: withDetail("the model credentials were refused"), Status: apierr.StatusCode, Cause: err}
case 408:
return &Error{Code: CodeTimeout, Message: withDetail("the model call timed out"), Status: 408, Cause: err}
case 429:
return &Error{Code: CodeRateLimited, Message: withDetail("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: withDetail(fmt.Sprintf("the model call failed (http %d)", apierr.StatusCode)),
Status: apierr.StatusCode, Cause: err,
}
}
}
// anthropicErrorMessage digs the human-readable reason out of an error body.
//
// Best-effort, exactly like its OpenAI counterpart: the envelope is documented
// and stable enough to be worth reading, and an unparseable body yields nothing
// rather than failing a failure.
func anthropicErrorMessage(raw string) string {
if raw == "" {
return ""
}
var envelope struct {
Error struct {
Message string `json:"message"`
Type string `json:"type"`
} `json:"error"`
}
if err := json.Unmarshal([]byte(raw), &envelope); err != nil {
return ""
}
if envelope.Error.Message != "" {
return envelope.Error.Message
}
return envelope.Error.Type
}
/* ── 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
}