返回 DeepSeek-Reasonix
anthropic.go
根目录 / internal / provider / anthropic / anthropic.go
1 // Package anthropic implements the Anthropic Messages API provider (POST
2 // /v1/messages, SSE streaming) with a hand-written net/http client — no SDK. It
3 // self-registers under the "anthropic" kind, so any Claude model is a config
4 // instance rather than code.
5 //
6 // Two notes, both rooted in the transport-agnostic provider.Message abstraction:
7 //
8 // - Extended thinking is opt-in (provider config thinking="adaptive"). Anthropic
9 // requires the *signed* thinking block be replayed on the next turn when a tool
10 // call followed thinking, so Message carries ReasoningSignature alongside
11 // ReasoningContent and this provider replays the signed block on the next
12 // request. DeepSeek's Anthropic endpoint instead uses unsigned thinking blocks,
13 // thinking.type enabled|disabled, and output_config.effort; requests carrying
14 // tools must replay all provider reasoning. Some other compatible gateways such
15 // as LongCat use the binary toggle without output_config. (redacted_thinking
16 // blocks are not yet captured/replayed.)
17 // - Native Anthropic requests omit temperature/top_p. Current Claude models
18 // (Opus 4.8/4.7) reject sampling parameters with a 400; Anthropic steers
19 // behavior via prompting instead. DeepSeek's compatible endpoint accepts the
20 // caller's temperature, so that field is preserved only for DeepSeek.
21 package anthropic
22
23 import (
24 "bytes"
25 "context"
26 "encoding/json"
27 "fmt"
28 "net/http"
29 "strings"
30 "sync"
31 "sync/atomic"
32 "time"
33
34 "reasonix/internal/netclient"
35 "reasonix/internal/provider"
36 "reasonix/internal/provider/openai"
37 )
38
39 // defaultStreamIdleTimeout caps how long a started SSE stream may go silent before
40 // it's treated as a dropped connection — a half-open TCP connection (proxy switched
41 // mid-stream) sends no RST, so scanner.Scan() would block forever. Generous on
42 // purpose; live streams emit far more often. Stored per-client (client.idleTimeout)
43 // so a test can shorten it without a shared global that races other watchdogs.
44 const defaultStreamIdleTimeout = 300 * time.Second
45
46 const (
47 // anthropicVersion is the required API version header value.
48 anthropicVersion = "2023-06-01"
49 // defaultBaseURL is the first-party endpoint; config may override it (e.g. a
50 // gateway). Bedrock/Vertex use a different request shape and are out of scope.
51 defaultBaseURL = "https://api.anthropic.com"
52 // defaultMaxTokens is the mandatory Anthropic fallback when neither config
53 // nor request supplies max_tokens. Native Anthropic stays at 16K; official
54 // DeepSeek sends the documented 384K ceiling because budget_tokens is ignored.
55 defaultMaxTokens = provider.DefaultOrdinaryOutputTokens
56 )
57
58 func init() {
59 provider.RegisterReasoning("anthropic", ReasoningForConfig)
60 provider.Register("anthropic", New)
61 }
62
63 // New builds an Anthropic provider from a resolved config.
64 func New(cfg provider.Config) (provider.Provider, error) {
65 cfg = provider.ApplyOpenCodeGoContract("anthropic", cfg)
66 if cfg.Model == "" {
67 return nil, fmt.Errorf("anthropic: model is required for provider %q", cfg.Name)
68 }
69 name := cfg.Name
70 if name == "" {
71 name = "anthropic"
72 }
73 baseURL := cfg.BaseURL
74 if baseURL == "" {
75 baseURL = defaultBaseURL
76 }
77 root, requestURL := resolveEndpoints(baseURL, cfg.Extra)
78 officialDeepSeek := openai.IsDeepSeek(root)
79 reasoningProtocol, _ := cfg.Extra["reasoning_protocol"].(string)
80 reasoningProtocol = strings.ToLower(strings.TrimSpace(reasoningProtocol))
81 deepSeekReplay := officialDeepSeek
82 switch reasoningProtocol {
83 case "deepseek":
84 deepSeekReplay = true
85 case "none":
86 deepSeekReplay = false
87 }
88 keyEnv, _ := cfg.Extra["api_key_env"].(string) // for actionable auth errors
89 keySource, _ := cfg.Extra["api_key_source"].(string)
90 thinking, _ := cfg.Extra["thinking"].(string)
91 thinking = strings.ToLower(strings.TrimSpace(thinking))
92 effort, err := configuredEffort(cfg)
93 if err != nil {
94 return nil, err
95 }
96 vision, _ := cfg.Extra["vision"].(bool)
97 modelInfo := provider.ModelInfo{ID: cfg.Model, InputModalities: []provider.ModelModality{provider.ModalityText}}
98 if cfg.ModelInfo != nil {
99 modelInfo = *cfg.ModelInfo
100 modelInfo.ID = cfg.Model
101 }
102 if cfg.ModelInfo != nil {
103 vision = modelInfo.SupportsInput(provider.ModalityImage)
104 }
105 // Official DeepSeek image input is pinned to one SKU even when metadata
106 // claims otherwise.
107 vision = openai.DeepSeekImageInputAllowed(officialDeepSeek, requestURL, cfg.Model, cfg.ModelInfo != nil, vision)
108 if vision {
109 modelInfo.InputModalities = []provider.ModelModality{provider.ModalityText, provider.ModalityImage}
110 } else if modelInfo.SupportsInput(provider.ModalityImage) {
111 modelInfo.InputModalities = []provider.ModelModality{provider.ModalityText}
112 }
113 webSearch, _ := cfg.Extra["web_search"].(bool)
114 clientWebSearch, _ := cfg.Extra["client_web_search"].(bool)
115 headers, _ := cfg.Extra["headers"].(map[string]string)
116 authHeader, _ := cfg.Extra["auth_header"].(bool)
117 maxOutputTokens, _ := cfg.Extra["max_output_tokens"].(int)
118 if maxOutputTokens <= 0 {
119 // Messages requires max_tokens. 0 = automatic; negative also falls back
120 // because the wire field is mandatory.
121 if officialDeepSeek {
122 maxOutputTokens = provider.DeepSeekMaxOutputTokens
123 } else {
124 // Native Anthropic and unknown gateways: conservative ordinary default.
125 maxOutputTokens = defaultMaxTokens
126 if strings.EqualFold(thinking, "adaptive") || strings.EqualFold(thinking, "enabled") {
127 maxOutputTokens = provider.AutoOutputBudget(true, effort)
128 }
129 }
130 }
131 httpClient, err := newHTTPClient(cfg)
132 if err != nil {
133 return nil, fmt.Errorf("anthropic: network: %w", err)
134 }
135 if reject, _ := cfg.Extra["reject_redirects"].(bool); reject {
136 httpClient.CheckRedirect = func(*http.Request, []*http.Request) error { return http.ErrUseLastResponse }
137 }
138 return &client{
139 identityHeaders: provider.NewClientIdentityHeaders(),
140 reasoning: ReasoningForConfig(cfg),
141 name: name,
142 identity: provider.RequestIdentity{Provider: name, DisplayName: cfg.DisplayName, Protocol: cfg.Protocol},
143 apiKey: cfg.APIKey,
144 keyEnv: keyEnv,
145 keySource: keySource,
146 baseURL: root,
147 requestURL: requestURL,
148 model: cfg.Model,
149 nativeAnthropic: strings.EqualFold(root, defaultBaseURL),
150 deepseek: deepSeekReplay,
151 thinking: thinking,
152 effort: effort,
153 vision: vision,
154 modelInfo: modelInfo,
155 mimo: provider.IsMiMoEndpoint(root),
156 search: provider.SearchPolicy{NativeEnabled: webSearch, ClientEnabled: clientWebSearch},
157 headers: cleanCustomHeaders(headers),
158 authHeader: authHeader,
159 defaultMaxTokens: maxOutputTokens,
160 http: httpClient, // no overall timeout; lifecycle is ctx-driven
161 idleTimeout: defaultStreamIdleTimeout,
162 }, nil
163 }
164
165 func newHTTPClient(cfg provider.Config) (*http.Client, error) {
166 if cfg.HTTPClient != nil {
167 return cfg.HTTPClient, nil
168 }
169 spec, _ := cfg.Extra["proxy_spec"].(netclient.ProxySpec)
170 return netclient.NewHTTPClient(spec, netclient.TransportOptions{
171 DialTimeout: 30 * time.Second,
172 KeepAlive: 30 * time.Second,
173 TLSHandshakeTimeout: 15 * time.Second,
174 ResponseHeaderTimeout: 300 * time.Second,
175 })
176 }
177
178 type client struct {
179 identityHeaders http.Header
180 reasoning provider.ReasoningCapability
181 name string
182 identity provider.RequestIdentity
183 apiKey string
184 keyEnv string // api_key_env name, surfaced in auth errors
185 keySource string // source of keyEnv, surfaced in auth errors
186 baseURL string
187 requestURL string
188 model string
189 nativeAnthropic bool // first-party endpoint: documented default-5m cache-write pricing applies
190 deepseek bool // official DeepSeek Anthropic endpoint: unsigned reasoning replay + automatic cache
191 thinking string // "adaptive" enables extended thinking; "" = off (config-driven)
192 effort string // output_config.effort: low|medium|high|xhigh|max; "" = provider default
193 vision bool // model accepts image input — embed attached images as base64 image blocks
194 modelInfo provider.ModelInfo
195 mimo bool // true for MiMo — upgrades legacy tuple schemas to Draft 2020-12
196 search provider.SearchPolicy
197 headers map[string]string
198 authHeader bool // send Authorization: Bearer instead of Anthropic's x-api-key header
199 defaultMaxTokens int
200 http *http.Client
201 idleTimeout time.Duration // SSE stall watchdog window; defaultStreamIdleTimeout unless a test overrides
202 authed atomic.Bool // a request has succeeded — gate transient-401 retry
203 }
204
205 func (c *client) Name() string { return c.name }
206
207 func (c *client) ModelInfo() provider.ModelInfo {
208 if c == nil {
209 return provider.ModelInfo{}
210 }
211 info := c.modelInfo
212 info.InputModalities = append([]provider.ModelModality(nil), info.InputModalities...)
213 return info
214 }
215
216 func (c *client) deepSeekThinkingEnabled() bool {
217 return c != nil && c.deepseek && c.thinking != "disabled" && c.effort != "disabled"
218 }
219
220 func (c *client) RequiresAssistantReasoningReplay(m provider.Message) bool {
221 if c == nil {
222 return false
223 }
224 if !c.deepseek {
225 return c.requiresReceivedReasoning(m)
226 }
227 // A turn carrying provider-issued reasoning must replay it, tools or not:
228 // DeepSeek's thinking mode 400s when a stored thinking block is not passed
229 // back. Without stored reasoning there is nothing to replay — plain text
230 // turns stay out of projection so healthy histories keep their backing.
231 if strings.TrimSpace(m.ReasoningContent) != "" {
232 return true
233 }
234 activity := len(m.ToolCalls) > 0 || len(m.ServerSearch) > 0
235 return activity && c.deepSeekThinkingEnabled()
236 }
237
238 func (c *client) AllowsEmptyReasoningFallback() bool { return false }
239
240 func (c *client) MissingToolCallReasoningWarningIdentity() string {
241 if c == nil {
242 return ""
243 }
244 protocol := "anthropic"
245 if c.deepseek {
246 protocol = "deepseek-anthropic"
247 }
248 return strings.Join([]string{
249 "anthropic", strings.TrimSpace(c.name), strings.TrimSpace(c.requestURL),
250 strings.TrimSpace(c.model), protocol, strings.TrimSpace(c.thinking), strings.TrimSpace(c.effort),
251 }, "\x00")
252 }
253
254 func (c *client) sendOpts() provider.SendOptions {
255 return provider.SendOptions{
256 Provider: c.name,
257 ProviderDisplayName: c.identity.DisplayName,
258 Protocol: c.identity.Protocol,
259 KeyEnv: c.keyEnv,
260 KeySource: c.keySource,
261 KeyPresent: c.apiKey != "",
262 RetryAuth: c.authed.Load(),
263 }
264 }
265
266 func cleanCustomHeaders(in map[string]string) map[string]string {
267 if len(in) == 0 {
268 return nil
269 }
270 out := make(map[string]string, len(in))
271 for name, value := range in {
272 name = strings.TrimSpace(name)
273 if name == "" || reservedCustomHeader(name) {
274 continue
275 }
276 out[name] = strings.TrimSpace(value)
277 }
278 if len(out) == 0 {
279 return nil
280 }
281 return out
282 }
283
284 func reservedCustomHeader(name string) bool {
285 switch strings.ToLower(strings.TrimSpace(name)) {
286 case "content-type", "accept", "x-api-key", "authorization", "anthropic-version", "anthropic-beta":
287 return true
288 default:
289 return false
290 }
291 }
292
293 func applyCustomHeaders(h http.Header, headers map[string]string) {
294 for name, value := range cleanCustomHeaders(headers) {
295 h.Set(name, value)
296 }
297 }
298
299 // bufPool reuses byte buffers for JSON-marshalled request bodies, reducing GC
300 // churn from repeated alloc/free of ~10-100KB buffers per turn.
301 var bufPool = sync.Pool{
302 New: func() any { return new(bytes.Buffer) },
303 }
304
305 func (c *client) Stream(ctx context.Context, req provider.Request) (<-chan provider.Chunk, error) {
306 if err := c.reasoning.Validate(c.model, req.EffortOverride); err != nil {
307 return nil, err
308 }
309 requestCtx := provider.WithRequestAttemptCounter(ctx)
310 buf := bufPool.Get().(*bytes.Buffer)
311 buf.Reset()
312 if err := json.NewEncoder(buf).Encode(c.buildRequest(requestCtx, req)); err != nil {
313 bufPool.Put(buf)
314 return nil, fmt.Errorf("%s: marshal request: %w", c.name, err)
315 }
316 body := make([]byte, buf.Len())
317 copy(body, buf.Bytes())
318 bufPool.Put(buf)
319
320 newReq := func(ctx context.Context) (*http.Request, error) {
321 httpReq, err := http.NewRequestWithContext(ctx, http.MethodPost, c.requestURL, bytes.NewReader(body))
322 if err != nil {
323 return nil, err
324 }
325 httpReq.Header.Set("Content-Type", "application/json")
326 httpReq.Header.Set("Accept", "text/event-stream")
327 if c.authHeader {
328 httpReq.Header.Set("Authorization", "Bearer "+c.apiKey)
329 } else {
330 httpReq.Header.Set("x-api-key", c.apiKey)
331 }
332 httpReq.Header.Set("anthropic-version", anthropicVersion)
333 if visionRequestUsesFileID(req) {
334 httpReq.Header.Set("anthropic-beta", "files-api-2025-04-14")
335 }
336 applyCustomHeaders(httpReq.Header, c.headers)
337 provider.ApplyOpenCodeGoHeaders(httpReq, c.baseURL, c.identityHeaders)
338 return httpReq, nil
339 }
340 resp, err := provider.SendWithRetry(requestCtx, c.http, c.sendOpts(), newReq)
341 if err != nil {
342 return nil, provider.AnnotateToolSchemaError(err, req.Tools)
343 }
344 c.authed.Store(true)
345
346 out := make(chan provider.Chunk)
347 go c.readStream(requestCtx, resp, out)
348 return out, nil
349 }
350
351 // buildRequest converts the transport-agnostic Request into the Messages API shape:
352 // RoleSystem messages lift to the top-level `system` field; assistant tool calls
353 // become `tool_use` blocks; RoleTool results become `tool_result` blocks in a user
354 // turn. Consecutive same-role messages are coalesced because the API requires
355 // alternating user/assistant turns (tool results are user turns).
356 func (c *client) buildRequest(ctx context.Context, req provider.Request) anthRequest {
357 system, msgs := c.buildMessages(req.Messages)
358
359 tools := encodeAnthTools(c, req)
360 if !c.deepseek {
361 markPromptCacheBreakpoints(system, tools, msgs)
362 }
363
364 maxTokens := req.MaxTokens
365 if maxTokens <= 0 {
366 maxTokens = c.defaultMaxTokens
367 if maxTokens <= 0 {
368 maxTokens = defaultMaxTokens
369 }
370 }
371 r := anthRequest{
372 Model: c.model,
373 MaxTokens: maxTokens,
374 System: system,
375 Messages: msgs,
376 Tools: tools,
377 Stream: true,
378 }
379 effort := c.effort
380 if req.EffortOverride != "" {
381 effort = req.EffortOverride
382 }
383 // Extended thinking is provider-specific. DeepSeek defaults to enabled and
384 // accepts output_config.effort alongside its binary toggle. Anthropic proper
385 // uses type=adaptive plus display/output_config. LongCat-style compatible
386 // gateways use the simpler enabled|disabled knob and reject output_config.
387 if c.deepseek {
388 c.applyDeepSeekThinking(&r, req)
389 } else {
390 thinking := c.thinking
391 if effort != "" && thinking == "" {
392 thinking = "adaptive"
393 }
394 switch thinking {
395 case "adaptive":
396 r.Thinking = &thinkingConfig{Type: "adaptive", Display: "summarized"}
397 if effort != "" {
398 r.OutputConfig = &outputConfig{Effort: effort}
399 }
400 case "enabled", "disabled":
401 t := c.thinking
402 if effort == "enabled" || effort == "disabled" {
403 t = effort
404 }
405 r.Thinking = &thinkingConfig{Type: t}
406 }
407 }
408 return r
409 }
410
411 // readStream parses the Messages API SSE stream into Chunks. Text deltas emit live;
412 // each tool_use content block emits a ChunkToolCallStart when its id+name are known
413 // and a complete ChunkToolCall when the block closes; usage is assembled from
414 // message_start/message_delta usage (compatible gateways may put every counter
415 // in the final delta) and emitted once before ChunkDone.
416 func (c *client) readStream(ctx context.Context, resp *http.Response, out chan<- provider.Chunk) {
417 defer resp.Body.Close()
418 defer close(out)
419
420 // Close the body if the stream stalls past c.idleTimeout so scanner.Scan()
421 // unblocks instead of hanging on a half-open connection. The watchdog owns the
422 // timer; the read loop only pings the buffered activity channel (no Timer.Reset
423 // race). A context cancel already unblocks the scan via the transport.
424 idleTimeout := c.idleTimeout
425 if idleTimeout <= 0 { // zero-value client (constructed without New)
426 idleTimeout = defaultStreamIdleTimeout
427 }
428 done := make(chan struct{})
429 defer close(done)
430 activity := make(chan struct{}, 1)
431 var stalled atomic.Bool
432 go func() {
433 idle := time.NewTimer(idleTimeout)
434 defer idle.Stop()
435 for {
436 select {
437 case <-ctx.Done():
438 resp.Body.Close()
439 return
440 case <-idle.C:
441 stalled.Store(true)
442 resp.Body.Close()
443 return
444 case <-activity:
445 if !idle.Stop() {
446 select {
447 case <-idle.C:
448 default:
449 }
450 }
451 idle.Reset(idleTimeout)
452 case <-done:
453 return
454 }
455 }
456 }()
457
458 send := func(chunk provider.Chunk) bool {
459 return sendChunk(ctx, out, chunk)
460 }
461
462 tools := map[int]*provider.ToolCall{} // tool_use blocks, keyed by content index
463 searches := newSearchStream()
464 thinking := map[int]*provider.ThinkingBlock{}
465 argBuckets := map[int]int{} // last emitted 2KB progress bucket per block
466 var inTok, outTok, cacheCreate, cacheRead int
467 var stopReason string
468 haveUsage := false
469 mergeUsage := func(usage *wireUsage) {
470 if usage == nil {
471 return
472 }
473 // The native Anthropic stream reports input/cache counters in
474 // message_start and output_tokens in message_delta. Compatible gateways
475 // such as LongCat report all counters in message_delta instead. Counters
476 // are cumulative and non-negative, so retaining the largest value also
477 // tolerates gateways that repeat partial usage in both events.
478 inTok = max(inTok, usage.InputTokens)
479 outTok = max(outTok, usage.OutputTokens)
480 cacheCreate = max(cacheCreate, usage.CacheCreationInputTokens)
481 cacheRead = max(cacheRead, usage.CacheReadInputTokens)
482 haveUsage = true
483 }
484
485 scanner := provider.NewStreamScanner(resp.Body, 1024*1024)
486
487 for scanner.Scan() {
488 select { // ping the idle watchdog; non-blocking so a full buffer is fine
489 case activity <- struct{}{}:
490 default:
491 }
492 line := strings.TrimSpace(scanner.Text())
493 // SSE carries `event:` and `data:` lines; the data JSON's own `type` field
494 // is authoritative, so we only need the data payloads.
495 if !strings.HasPrefix(line, "data:") {
496 continue
497 }
498 data := strings.TrimSpace(strings.TrimPrefix(line, "data:"))
499 if data == "" {
500 continue
501 }
502
503 var ev streamEvent
504 if err := json.Unmarshal([]byte(data), &ev); err != nil {
505 send(provider.Chunk{Type: provider.ChunkError, Err: scanner.DecodeError(c.name, data, err)})
506 return
507 }
508
509 if !updateThinkingStream(thinking, ev, send) {
510 return
511 }
512 switch ev.Type {
513 case "message_start":
514 if ev.Message != nil && ev.Message.Usage != nil {
515 mergeUsage(ev.Message.Usage)
516 }
517 case "content_block_start":
518 if chunk := beginContentBlock(ev.Index, ev.ContentBlock, tools, searches); chunk != nil {
519 if !send(*chunk) {
520 return
521 }
522 }
523 case "content_block_delta":
524 if ev.Delta == nil {
525 continue
526 }
527 switch ev.Delta.Type {
528 case "text_delta":
529 if ev.Delta.Text != "" {
530 if !send(provider.Chunk{Type: provider.ChunkText, Text: ev.Delta.Text}) {
531 return
532 }
533 }
534 case "input_json_delta":
535 if tc := tools[ev.Index]; tc != nil {
536 tc.Arguments += ev.Delta.PartialJSON
537 // Progress ticks for large streaming argument payloads, one
538 // per 2KB bucket (see the openai provider for rationale).
539 if bucket := len(tc.Arguments) / 2048; bucket > argBuckets[ev.Index] {
540 argBuckets[ev.Index] = bucket
541 if !send(provider.Chunk{Type: provider.ChunkToolCallArgsDelta, ToolCall: &provider.ToolCall{ID: tc.ID, Name: tc.Name}, ArgChars: len(tc.Arguments)}) {
542 return
543 }
544 }
545 }
546 if next := searches.argsDelta(ev.Index, ev.Delta.PartialJSON); next != nil {
547 if !send(provider.Chunk{Type: provider.ChunkServerSearch, ServerSearch: next}) {
548 return
549 }
550 }
551 case "web_search_tool_result_delta":
552 // Some DeepSeek-compatible streams deliver the result array in a
553 // delta instead of the block-start content; without this the card
554 // stays empty and the model-written source list is all that shows.
555 if next := searches.resultsDelta(ev.Index, ev.Delta.WebSearchResults); next != nil {
556 if !send(provider.Chunk{Type: provider.ChunkServerSearch, ServerSearch: next}) {
557 return
558 }
559 }
560 }
561 case "content_block_stop":
562 if tc := tools[ev.Index]; tc != nil {
563 if !send(provider.Chunk{Type: provider.ChunkToolCall, ToolCall: tc}) {
564 return
565 }
566 delete(tools, ev.Index)
567 }
568 case "message_delta":
569 if ev.Delta != nil && ev.Delta.StopReason != "" {
570 stopReason = ev.Delta.StopReason
571 }
572 mergeUsage(ev.Usage)
573 case "message_stop":
574 // Anthropic's terminal event. Tool blocks may already have closed;
575 // without this, the attempt stays speculative and is not committed.
576 // Stop reading immediately so a post-terminal connection reset cannot
577 // reclassify a complete response as interrupted.
578 goto finalize
579 case "error":
580 msg := "stream error"
581 if ev.Error != nil && ev.Error.Message != "" {
582 msg = ev.Error.Message
583 }
584 send(provider.Chunk{Type: provider.ChunkError, Err: fmt.Errorf("%s: %s", c.name, msg)})
585 return
586 }
587 }
588
589 if ctx.Err() != nil {
590 return
591 }
592 if err := streamScanEndError(c.name, idleTimeout, stalled.Load(), scanner.Err(), stopReason); err != nil {
593 send(provider.Chunk{Type: provider.ChunkError, Err: err})
594 return
595 }
596 goto finalize
597
598 finalize:
599 if len(thinking) > 0 {
600 send(provider.Chunk{Type: provider.ChunkReasoning, ReasoningState: provider.ReasoningIncomplete})
601 }
602 if haveUsage {
603 cacheWriteBilledTokens := 0.0
604 if cacheCreate > 0 && c.nativeAnthropic {
605 cacheWriteBilledTokens = float64(cacheCreate) * cacheWrite5MinuteInputMultiplier
606 }
607 usage := &provider.Usage{
608 PromptTokens: inTok + cacheCreate + cacheRead,
609 CompletionTokens: outTok,
610 TotalTokens: inTok + cacheCreate + cacheRead + outTok,
611 CacheHitTokens: cacheRead,
612 CacheMissTokens: inTok + cacheCreate,
613 CacheWriteTokens: cacheCreate,
614 CacheWriteBilledTokens: cacheWriteBilledTokens,
615 FinishReason: mapStopReason(stopReason),
616 }
617 provider.ApplyRequestAttemptCount(ctx, usage)
618 if !send(provider.Chunk{Type: provider.ChunkUsage, Usage: usage}) {
619 return
620 }
621 }
622 send(provider.Chunk{Type: provider.ChunkDone})
623 }
624
625 func sendChunk(ctx context.Context, out chan<- provider.Chunk, chunk provider.Chunk) bool {
626 select {
627 case out <- chunk:
628 return true
629 default:
630 }
631 select {
632 case <-ctx.Done():
633 return false
634 case out <- chunk:
635 return true
636 }
637 }
638
639 // mapStopReason translates Anthropic stop reasons to the OpenAI-style finish
640 // reasons the agent already recognises (it surfaces abnormal ones like "length").
641 func mapStopReason(s string) string {
642 switch s {
643 case "end_turn", "stop_sequence":
644 return "stop"
645 case "tool_use":
646 return "tool_calls"
647 case "max_tokens":
648 return "length"
649 default:
650 return s // "refusal", "pause_turn", "" — pass through
651 }
652 }
653
654 // Messages API wire protocol
655
656 const cacheWrite5MinuteInputMultiplier = 1.25
657
658 func ephemeral() *cacheControl { return &cacheControl{Type: "ephemeral"} }
659
660 type cacheControl struct {
661 Type string `json:"type"`
662 }
663
664 type anthRequest struct {
665 Model string `json:"model"`
666 MaxTokens int `json:"max_tokens"`
667 System []textBlock `json:"system,omitempty"`
668 Messages []anthMessage `json:"messages"`
669 Tools []anthTool `json:"tools,omitempty"`
670 Temperature *float64 `json:"temperature,omitempty"`
671 Thinking *thinkingConfig `json:"thinking,omitempty"`
672 OutputConfig *outputConfig `json:"output_config,omitempty"`
673 Stream bool `json:"stream"`
674 }
675
676 type thinkingConfig struct {
677 Type string `json:"type"` // "adaptive"
678 Display string `json:"display,omitempty"` // "summarized" to stream the reasoning text
679 }
680
681 type outputConfig struct {
682 Effort string `json:"effort,omitempty"` // low | high | max
683 }
684
685 type textBlock struct {
686 Type string `json:"type"`
687 Text string `json:"text"`
688 CacheControl *cacheControl `json:"cache_control,omitempty"`
689 }
690
691 type anthMessage struct {
692 Role string `json:"role"`
693 Content []contentBlock `json:"content"`
694 }
695
696 // contentBlock is the union of the block kinds we emit in a request: text,
697 // tool_use (echoing a prior assistant call), and tool_result. Unused fields are
698 // omitted so each block serialises to its canonical shape.
699 type contentBlock struct {
700 Data string `json:"data,omitempty"`
701 Type string `json:"type"`
702 Text string `json:"text,omitempty"` // text
703 Thinking string `json:"thinking,omitempty"` // thinking
704 Signature string `json:"signature,omitempty"` // thinking
705 ID string `json:"id,omitempty"` // tool_use
706 Name string `json:"name,omitempty"` // tool_use
707 Input json.RawMessage `json:"input,omitempty"` // tool_use
708 ToolUseID string `json:"tool_use_id,omitempty"` // tool_result
709 Content any `json:"content,omitempty"` // tool_result: string, or []contentBlock when the result carries images
710 Source *imageSource `json:"source,omitempty"` // image
711 CacheControl *cacheControl `json:"cache_control,omitempty"`
712 }
713
714 type imageSource struct {
715 Type string `json:"type"` // base64 | url | file
716 MediaType string `json:"media_type,omitempty"`
717 Data string `json:"data,omitempty"`
718 URL string `json:"url,omitempty"`
719 FileID string `json:"file_id,omitempty"`
720 }
721
722 // toolResultBlocks builds array content for a tool_result whose message carries
723 // images: the text first, then one image block per parseable data URL. It
724 // returns nil when nothing parses, so text-only results keep plain string
725 // content — byte-identical serialization to previous releases.
726 func toolResultBlocks(text string, images []string) []contentBlock {
727 var imgs []contentBlock
728 for _, url := range images {
729 if mt, data, ok := provider.ParseImageDataURL(url); ok {
730 imgs = append(imgs, contentBlock{Type: "image", Source: &imageSource{Type: "base64", MediaType: mt, Data: data}})
731 }
732 }
733 if imgs == nil {
734 return nil
735 }
736 return append([]contentBlock{{Type: "text", Text: text}}, imgs...)
737 }
738
739 type anthTool struct {
740 Type string `json:"type,omitempty"` // "web_search" for server-side search; empty for named tools
741 Name string `json:"name,omitempty"`
742 Description string `json:"description,omitempty"`
743 InputSchema json.RawMessage `json:"input_schema,omitempty"`
744 Strict bool `json:"strict,omitempty"`
745 DeferLoading bool `json:"defer_loading,omitempty"`
746 CacheControl *cacheControl `json:"cache_control,omitempty"`
747 }
748
749 // streamEvent is the discriminated SSE event; read the fields matching Type.
750 type streamEvent struct {
751 Type string `json:"type"`
752 Index int `json:"index"`
753 Message *struct {
754 Usage *wireUsage `json:"usage"`
755 } `json:"message"`
756 ContentBlock *streamContentBlock `json:"content_block"`
757 Delta *struct {
758 Type string `json:"type"` // text_delta | thinking_delta | signature_delta | input_json_delta | web_search_tool_result_delta
759 Text string `json:"text"` // text_delta
760 Thinking string `json:"thinking"` // thinking_delta
761 Signature string `json:"signature"` // signature_delta
762 PartialJSON string `json:"partial_json"` // input_json_delta
763 StopReason string `json:"stop_reason"` // message_delta
764 WebSearchResults json.RawMessage `json:"results"` // web_search_tool_result_delta
765 } `json:"delta"`
766 Usage *wireUsage `json:"usage"` // message_delta (cumulative output_tokens)
767 Error *struct {
768 Type string `json:"type"`
769 Message string `json:"message"`
770 } `json:"error"`
771 }
772
773 type wireUsage struct {
774 InputTokens int `json:"input_tokens"`
775 OutputTokens int `json:"output_tokens"`
776 CacheCreationInputTokens int `json:"cache_creation_input_tokens"`
777 CacheReadInputTokens int `json:"cache_read_input_tokens"`
778 }
779
779 lines GO