返回 DeepSeek-Reasonix
responses.go
根目录 / internal / provider / responses / responses.go
1 // Package responses implements the OpenAI Responses API wire protocol.
2 // DeepSeek uses it statelessly and requires the complete input history on every
3 // request; compatible stateful endpoints may opt into previous_response_id.
4 package responses
5
6 import (
7 "bufio"
8 "bytes"
9 "context"
10 "crypto/sha256"
11 "encoding/hex"
12 "encoding/json"
13 "fmt"
14 "io"
15 "maps"
16 "net/http"
17 "strings"
18 "sync"
19 "sync/atomic"
20 "time"
21
22 "reasonix/internal/netclient"
23 "reasonix/internal/provider"
24 "reasonix/internal/provider/openai"
25 )
26
27 const (
28 defaultStreamIdleTimeout = 300 * time.Second
29 maxReplayableSearchItemBytes = 512 * 1024
30 )
31
32 func init() {
33 provider.RegisterReasoning("responses", ReasoningForConfig)
34 provider.RegisterReasoning("dashscope-responses", ReasoningForConfig)
35 provider.Register("responses", newFromConfig)
36 provider.Register("dashscope-responses", newFromConfig)
37 }
38
39 // Config holds Responses API provider settings.
40 type Config struct {
41 HTTPClient *http.Client
42 Name string
43 DisplayName string
44 Protocol string
45 APIKey string
46 BaseURL string
47 Model string
48 ModelInfo *provider.ModelInfo
49 Effort string
50 Mode string // stateful | stateless; empty uses vendor detection.
51 Stateful *bool // legacy form of Mode; nil preserves vendor detection.
52 WebSearch bool // expose the provider-executed web_search tool.
53 Proxy netclient.ProxySpec
54 KeyEnv string
55 KeySource string
56 RequestURL string // optional exact Responses request URL; empty derives from BaseURL
57 // MaxOutputTokens is the total provider output budget. Zero omits the field
58 // on official DeepSeek (server 384K ceiling) and unknown endpoints; MiMo
59 // still applies its 16K/32K ladder. Negative values omit it.
60 MaxOutputTokens int
61 // SessionCache controls DashScope's opt-in header. The header is never sent
62 // to non-DashScope endpoints even when this value is true.
63 SessionCache *bool
64 // Extra carries kind-specific options; "vision" (bool) enables embedding
65 // attached Images as input_image parts on user turns.
66 Extra map[string]any
67 }
68
69 func (c Config) mode() string {
70 mode := strings.ToLower(strings.TrimSpace(c.Mode))
71 if mode == "stateful" || mode == "stateless" {
72 return mode
73 }
74 if c.Stateful != nil {
75 if *c.Stateful {
76 return "stateful"
77 }
78 return "stateless"
79 }
80 if capabilitiesFor(DetectVendor(c.BaseURL)).stateless {
81 return "stateless"
82 }
83 return "stateful"
84 }
85
86 // DetectVendor lives in vendor.go (capabilities table): it covers dashscope/
87 // deepseek (incl. eu.deepseek.com) / mimo via exact-host matching.
88
89 type client struct {
90 identityHeaders http.Header
91 reasoning provider.ReasoningCapability
92 name string
93 identity provider.RequestIdentity
94 apiKey, keyEnv, keySource string
95 baseURL, requestURL, model, effort string
96 vendor, mode string
97 caps vendorCapabilities
98 sessionCache bool
99 search provider.SearchPolicy
100 maxOutputTokens int
101 vision bool // model accepts image input; embed Images as input_image parts
102 modelInfo provider.ModelInfo
103 http *http.Client
104 idleTimeout time.Duration
105 authed atomic.Bool
106
107 mu sync.Mutex
108 lastResponseID string
109 expectedPrefixDigest string
110 }
111
112 // New creates a Responses API provider.
113 func New(cfg Config) provider.Provider {
114 cfg.Extra = maps.Clone(cfg.Extra)
115 if cfg.Extra == nil {
116 cfg.Extra = map[string]any{}
117 }
118 if cfg.RequestURL != "" {
119 cfg.Extra["request_url"] = cfg.RequestURL
120 }
121 resolved := provider.ApplyOpenCodeGoContract("responses", provider.Config{BaseURL: cfg.BaseURL, Model: cfg.Model, Extra: cfg.Extra})
122 cfg.Extra = resolved.Extra
123 vendor := DetectVendor(cfg.BaseURL)
124 cap := capabilitiesFor(vendor)
125 // Explicit replay contracts apply to compatible gateways as well as exact
126 // vendor hosts. Do not inherit endpoint defaults, headers, or output limits.
127 if protocol, _ := cfg.Extra["reasoning_protocol"].(string); strings.EqualFold(strings.TrimSpace(protocol), "deepseek") || strings.EqualFold(strings.TrimSpace(protocol), "mimo") {
128 cap.toolCallReasoning = true
129 }
130 maxOutputTokens := cfg.MaxOutputTokens
131 // Official DeepSeek omits max_output_tokens (server 384K). MiMo still uses
132 // the 16K/32K effort ladder. Compact_ratio is independent.
133 if maxOutputTokens == 0 && vendor == "mimo" {
134 maxOutputTokens = responsesAutoOutputBudget(vendor, cfg.Effort)
135 } else if maxOutputTokens == 0 && vendor != "deepseek" && cap.defaultMaxOutputTokens > 0 {
136 maxOutputTokens = cap.defaultMaxOutputTokens
137 }
138 sessionCache := cap.sessionCacheHeader
139 if cfg.SessionCache != nil {
140 sessionCache = *cfg.SessionCache
141 }
142 vision, _ := cfg.Extra["vision"].(bool)
143 if cfg.ModelInfo != nil {
144 vision = cfg.ModelInfo.SupportsInput(provider.ModalityImage)
145 }
146 // Official DeepSeek image input is pinned to one SKU. Ignore metadata or
147 // Extra["vision"] for Flash/Pro.
148 vision = openai.DeepSeekImageInputAllowed(vendor == "deepseek", cfg.RequestURL, cfg.Model, cfg.ModelInfo != nil, vision)
149 httpClient := &http.Client{}
150 if built, err := netclient.NewHTTPClient(cfg.Proxy, netclient.TransportOptions{
151 DialTimeout: 30 * time.Second, KeepAlive: 30 * time.Second,
152 TLSHandshakeTimeout: 15 * time.Second, ResponseHeaderTimeout: 300 * time.Second,
153 }); err == nil {
154 httpClient = built
155 }
156 if cfg.HTTPClient != nil {
157 httpClient = cfg.HTTPClient
158 }
159 baseURL, requestURL := resolveEndpoints(cfg.BaseURL, cfg.RequestURL)
160 modelInfo := provider.ModelInfo{ID: cfg.Model, InputModalities: []provider.ModelModality{provider.ModalityText}}
161 if cfg.ModelInfo != nil {
162 modelInfo = *cfg.ModelInfo
163 modelInfo.ID = cfg.Model
164 }
165 clientWebSearch, _ := cfg.Extra["client_web_search"].(bool)
166 if reject, _ := cfg.Extra["reject_redirects"].(bool); reject {
167 httpClient.CheckRedirect = func(*http.Request, []*http.Request) error { return http.ErrUseLastResponse }
168 }
169 if vision {
170 modelInfo.InputModalities = []provider.ModelModality{provider.ModalityText, provider.ModalityImage}
171 } else if modelInfo.SupportsInput(provider.ModalityImage) {
172 modelInfo.InputModalities = []provider.ModelModality{provider.ModalityText}
173 }
174 return &client{
175 identityHeaders: provider.NewClientIdentityHeaders(),
176 name: cfg.Name,
177 identity: provider.RequestIdentity{Provider: cfg.Name, DisplayName: cfg.DisplayName, Protocol: cfg.Protocol},
178 apiKey: cfg.APIKey, keyEnv: cfg.KeyEnv, keySource: cfg.KeySource,
179 reasoning: ReasoningForConfig(provider.Config{BaseURL: cfg.BaseURL, Model: cfg.Model, Extra: cfg.Extra}),
180 baseURL: baseURL, requestURL: requestURL, model: cfg.Model, effort: cfg.Effort,
181 vendor: vendor, caps: cap, mode: cfg.mode(), sessionCache: sessionCache, search: provider.SearchPolicy{NativeEnabled: cfg.WebSearch, ClientEnabled: clientWebSearch}, maxOutputTokens: maxOutputTokens,
182 vision: vision,
183 modelInfo: modelInfo,
184 http: httpClient, idleTimeout: defaultStreamIdleTimeout,
185 }
186 }
187
188 func (c *client) ModelInfo() provider.ModelInfo {
189 if c == nil {
190 return provider.ModelInfo{}
191 }
192 info := c.modelInfo
193 info.InputModalities = append([]provider.ModelModality(nil), info.InputModalities...)
194 return info
195 }
196
197 func responsesReasoningDisabled(effort string) bool {
198 switch strings.ToLower(strings.TrimSpace(effort)) {
199 case "none", "disabled", "off":
200 return true
201 default:
202 return false
203 }
204 }
205
206 // responsesAutoOutputBudget is the MiMo (and similar) 16K/32K ladder.
207 // Official DeepSeek must not call this; it omits max_output_tokens instead.
208 func responsesAutoOutputBudget(vendor, effort string) int {
209 if responsesReasoningDisabled(effort) {
210 return provider.AutoOutputBudget(false, effort)
211 }
212 e := strings.ToLower(strings.TrimSpace(effort))
213 if vendor == "deepseek" && (e == "" || e == "auto") {
214 e = "high"
215 }
216 return provider.AutoOutputBudget(true, e)
217 }
218
219 func (c *client) Name() string { return c.name }
220
221 func (c *client) NativeToolSearchAvailable() bool {
222 return c != nil && provider.IsFirstPartyOpenAI(c.baseURL) && nativeToolSearchModel(c.model)
223 }
224
225 func nativeToolSearchModel(model string) bool {
226 model = strings.ToLower(strings.TrimSpace(model))
227 return strings.HasPrefix(model, "gpt-5.4") || strings.HasPrefix(model, "gpt-5.5") || strings.HasPrefix(model, "gpt-5.6")
228 }
229
230 func (c *client) sendOpts() provider.SendOptions {
231 return provider.SendOptions{Provider: c.name, ProviderDisplayName: c.identity.DisplayName, Protocol: c.identity.Protocol, KeyEnv: c.keyEnv, KeySource: c.keySource, KeyPresent: c.apiKey != "", RetryAuth: c.authed.Load()}
232 }
233
234 // ResetContext drops stateful continuation metadata. Full-input stateless mode
235 // is unaffected.
236 func (c *client) ResetContext() {
237 c.mu.Lock()
238 c.lastResponseID = ""
239 c.expectedPrefixDigest = ""
240 c.mu.Unlock()
241 }
242
243 func (c *client) Stream(ctx context.Context, req provider.Request) (<-chan provider.Chunk, error) {
244 if c.effort != "auto" && c.effort != "off" {
245 if err := c.reasoning.Validate(c.model, c.effort); err != nil {
246 return nil, err
247 }
248 }
249 if err := c.reasoning.Validate(c.model, req.EffortOverride); err != nil {
250 return nil, err
251 }
252 requestCtx := provider.WithRequestAttemptCounter(ctx)
253 body, usedPrevious, wireMessages := c.buildRequestBody(req)
254 resp, err := c.send(requestCtx, body)
255 if err != nil && usedPrevious && isStalePreviousResponseError(err) {
256 // A stateful response ID may expire server-side. Retrying once with full
257 // history is safe because no response body has started streaming.
258 c.ResetContext()
259 body, _, wireMessages = c.buildRequestBody(req)
260 resp, err = c.send(requestCtx, body)
261 }
262 if err != nil && isCommandCodeTransientResponsesError(c.requestURL, err) {
263 // Retry once only for Command Code's narrow anonymous-400 signature; the
264 // byte-identical body is safe to resend and named schema failures remain
265 // terminal.
266 resp, err = c.send(requestCtx, body)
267 }
268 if err != nil {
269 return nil, err
270 }
271 c.authed.Store(true)
272 out := make(chan provider.Chunk, 64)
273 go c.readStream(requestCtx, resp, out, wireMessages)
274 return out, nil
275 }
276
277 func (c *client) send(ctx context.Context, body map[string]any) (*http.Response, error) {
278 payload, err := json.Marshal(body)
279 if err != nil {
280 return nil, fmt.Errorf("responses: marshal request: %w", err)
281 }
282 newRequest := func(ctx context.Context) (*http.Request, error) {
283 req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.requestURL, bytes.NewReader(payload))
284 if err != nil {
285 return nil, err
286 }
287 req.Header.Set("Content-Type", "application/json")
288 req.Header.Set("Authorization", "Bearer "+c.apiKey)
289 provider.ApplyOpenCodeGoHeaders(req, c.baseURL, c.identityHeaders)
290 if c.caps.sessionCacheHeader && c.sessionCache {
291 req.Header.Set("x-dashscope-session-cache", "enable")
292 }
293 return req, nil
294 }
295 return provider.SendWithRetry(ctx, c.http, c.sendOpts(), newRequest)
296 }
297
298 func (c *client) buildRequestBody(req provider.Request) (map[string]any, bool, []provider.Message) {
299 messages := provider.SanitizeToolPairing(provider.ModelMessages(req.Messages))
300 body := map[string]any{"model": c.model, "stream": true}
301
302 effort := c.effort
303 if req.EffortOverride != "" {
304 effort = req.EffortOverride
305 }
306
307 switch effort {
308 case "auto":
309 effort = ""
310 case "disabled", "off":
311 effort = "none"
312 }
313 if effort != "" {
314 body["reasoning"] = map[string]any{"effort": effort}
315 }
316 maxOutputTokens := req.MaxTokens
317 if maxOutputTokens == 0 {
318 maxOutputTokens = c.maxOutputTokens
319 }
320 if maxOutputTokens == 0 && c.vendor == "mimo" {
321 maxOutputTokens = responsesAutoOutputBudget(c.vendor, c.effort)
322 } else if maxOutputTokens == 0 && c.vendor != "deepseek" && c.caps.defaultMaxOutputTokens > 0 {
323 maxOutputTokens = c.caps.defaultMaxOutputTokens
324 }
325 if maxOutputTokens > 0 {
326 body["max_output_tokens"] = maxOutputTokens
327 }
328 if req.ResponseFormat != nil && req.ResponseFormat.Type != "" {
329 // Structured output: Responses text.format. MiMo/DashScope/OpenAI
330 // all accept {"text":{"format":{"type":"json_object"}}}. The model
331 // only emits JSON when the instructions also demand it.
332 body["text"] = map[string]any{
333 "format": map[string]any{"type": req.ResponseFormat.Type},
334 }
335 }
336 if req.Temperature != nil && !c.caps.ignoresTemperature {
337 body["temperature"] = *req.Temperature
338 }
339 if c.search.NativeEnabled || len(req.Tools) > 0 {
340 body["tools"] = encodeResponsesTools(c, req)
341 }
342 instructions, rest := splitInstructions(messages)
343 if instructions != "" {
344 body["instructions"] = instructions
345 }
346
347 c.mu.Lock()
348 previousID, expectedDigest := c.lastResponseID, c.expectedPrefixDigest
349 c.mu.Unlock()
350 if c.canUseStatefulContinuation(messages, previousID, expectedDigest) {
351 body["input"] = messages[len(messages)-1].Content
352 body["previous_response_id"] = previousID
353 return body, true, messages
354 }
355
356 body["input"] = messagesToInput(rest, c.vision, c.search.NativeEnabled, c.caps.summaryRequired)
357 return body, false, messages
358 }
359
360 func inputImagePart(ref string) map[string]string {
361 switch provider.ClassifyImage(ref) {
362 case provider.ImageFileID:
363 return map[string]string{"type": "input_image", "file_id": ref}
364 case provider.ImageDataURL, provider.ImageHTTPURL:
365 return map[string]string{"type": "input_image", "image_url": ref}
366 default:
367 return nil
368 }
369 }
370
371 func splitInstructions(messages []provider.Message) (string, []provider.Message) {
372 if len(messages) == 0 || messages[0].Role != provider.RoleSystem {
373 return "", messages
374 }
375 return messages[0].Content, messages[1:]
376 }
377
378 func decodeReplayableWebSearchItem(raw json.RawMessage) (map[string]any, bool) {
379 if len(raw) == 0 || len(raw) > maxReplayableSearchItemBytes || !json.Valid(raw) {
380 return nil, false
381 }
382 var item map[string]any
383 if err := json.Unmarshal(raw, &item); err != nil || item["type"] != "web_search_call" {
384 return nil, false
385 }
386 id, _ := item["id"].(string)
387 status, _ := item["status"].(string)
388 if strings.TrimSpace(id) == "" || status != "completed" {
389 return nil, false
390 }
391 return item, true
392 }
393
394 func (c *client) conversationDigest(messages []provider.Message) string {
395 instructions, rest := splitInstructions(messages)
396 // Digest must mirror buildRequestBody's wire shape and vision/summary knobs;
397 // a mismatch skips previous_response_id and forces a full replay (cache-hit loss).
398 payload, _ := json.Marshal(struct {
399 Instructions string `json:"instructions,omitempty"`
400 Input []map[string]any `json:"input"`
401 }{Instructions: instructions, Input: messagesToInput(rest, c.vision, c.search.NativeEnabled, c.caps.summaryRequired)})
402 sum := sha256.Sum256(payload)
403 return hex.EncodeToString(sum[:])
404 }
405
406 type streamedCall struct {
407 id, name, arguments string
408 argChars int
409 completed bool
410 }
411
412 func (c *client) readStream(ctx context.Context, resp *http.Response, out chan<- provider.Chunk, requestMessages []provider.Message) {
413 defer resp.Body.Close()
414 defer close(out)
415
416 scanner := bufio.NewScanner(resp.Body)
417 scanner.Buffer(make([]byte, 64*1024), 4*1024*1024)
418 idle := c.idleTimeout
419 if idle <= 0 {
420 idle = defaultStreamIdleTimeout
421 }
422 watchDone := make(chan struct{})
423 activity := make(chan struct{}, 1)
424 var reasoningSnapshots responseReasoningSnapshots
425 var stalled atomic.Bool
426 go func() {
427 timer := time.NewTimer(idle)
428 defer timer.Stop()
429 for {
430 select {
431 case <-ctx.Done():
432 _ = resp.Body.Close()
433 return
434 case <-watchDone:
435 return
436 case <-activity:
437 if !timer.Stop() {
438 select {
439 case <-timer.C:
440 default:
441 }
442 }
443 timer.Reset(idle)
444 case <-timer.C:
445 stalled.Store(true)
446 _ = resp.Body.Close()
447 return
448 }
449 }
450 }()
451 defer close(watchDone)
452
453 calls := make(map[string]*streamedCall)
454 callOrder := make([]string, 0)
455 callForItem := func(itemID string) *streamedCall {
456 if call := calls[itemID]; call != nil {
457 return call
458 }
459 call := &streamedCall{id: itemID}
460 calls[itemID] = call
461 callOrder = append(callOrder, itemID)
462 return call
463 }
464 textDeltas := make(map[string]bool)
465 reasoningDeltas := make(map[string]bool)
466 seenSearchItems := make(map[string]struct{})
467 var responsesItems []json.RawMessage
468 var text, reasoning strings.Builder
469 reasoningID := ""
470 reasoningStatus := ""
471 terminal := false
472 failed := false
473 completedResponseID := ""
474
475 for scanner.Scan() {
476 select {
477 case activity <- struct{}{}:
478 default:
479 }
480 line := scanner.Text()
481 if !strings.HasPrefix(line, "data:") {
482 continue
483 }
484 data := strings.TrimSpace(strings.TrimPrefix(line, "data:"))
485 if data == "[DONE]" {
486 terminal = true
487 break
488 }
489 var event sseEvent
490 if json.Unmarshal([]byte(data), &event) != nil {
491 continue
492 }
493 key := fmt.Sprintf("%s:%d", event.ItemID, event.ContentIndex)
494 reasoningSnapshots.capture(event)
495 switch event.Type {
496 case "response.output_text.delta":
497 textDeltas[key] = true
498 text.WriteString(event.Delta)
499 if !sendChunk(ctx, out, provider.Chunk{Type: provider.ChunkText, Text: event.Delta}) {
500 return
501 }
502 case "response.output_text.done":
503 if event.Text != "" && !textDeltas[key] {
504 text.WriteString(event.Text)
505 if !sendChunk(ctx, out, provider.Chunk{Type: provider.ChunkText, Text: event.Text}) {
506 return
507 }
508 }
509 case "response.reasoning_text.delta", "response.reasoning_summary_text.delta":
510 reasoningDeltas[key] = true
511 reasoning.WriteString(event.Delta)
512 if !sendChunk(ctx, out, provider.Chunk{Type: provider.ChunkReasoning, Text: event.Delta}) {
513 return
514 }
515 case "response.reasoning_text.done", "response.reasoning_summary_text.done":
516 if event.Text != "" && !reasoningDeltas[key] {
517 reasoning.WriteString(event.Text)
518 if !sendChunk(ctx, out, provider.Chunk{Type: provider.ChunkReasoning, Text: event.Text}) {
519 return
520 }
521 }
522 case "response.output_item.added":
523 if event.Item != nil {
524 switch event.Item.Type {
525 case "function_call":
526 call := callForItem(event.Item.ID)
527 call.id = event.Item.CallID
528 call.name = event.Item.Name
529 if !sendChunk(ctx, out, provider.Chunk{Type: provider.ChunkToolCallStart, ToolCall: &provider.ToolCall{ID: call.id, Name: call.name}}) {
530 return
531 }
532 case "reasoning":
533 // Capture the provider-issued reasoning item id so the
534 // next turn's input reasoning item can carry it (the
535 // OpenAI Responses schema marks Reasoning.id required).
536 if event.Item.ID != "" {
537 // 多段推理(DeepSeek 长思考分多段)时末段 id 覆盖:round-trip
538 // 合并为一个 reasoning item 只带末段 id(服务端接受)。
539 reasoningID = event.Item.ID
540 }
541 }
542 }
543 case "response.function_call_arguments.delta":
544 call := callForItem(event.ItemID)
545 call.arguments += event.Delta
546 call.argChars += len(event.Delta)
547 if !sendChunk(ctx, out, provider.Chunk{Type: provider.ChunkToolCallArgsDelta, ToolCall: &provider.ToolCall{ID: call.id, Name: call.name}, ArgChars: call.argChars}) {
548 return
549 }
550 case "response.function_call_arguments.done":
551 call := callForItem(event.ItemID)
552 if event.Arguments != "" {
553 call.arguments = event.Arguments
554 }
555 if !completeFunctionCall(ctx, out, call) {
556 return
557 }
558 case "response.output_item.done":
559 if event.Item != nil && event.Item.Type == "web_search_call" && c.search.NativeEnabled {
560 if _, ok := decodeReplayableWebSearchItem(event.Item.Raw); ok {
561 key := event.Item.ID
562 if key == "" {
563 key = string(event.Item.Raw)
564 }
565 if _, seen := seenSearchItems[key]; !seen {
566 seenSearchItems[key] = struct{}{}
567 raw := append(json.RawMessage(nil), event.Item.Raw...)
568 responsesItems = append(responsesItems, raw)
569 if !emitSearchReplay(ctx, out, raw) {
570 return
571 }
572 }
573 }
574 }
575 if event.Item != nil {
576 switch event.Item.Type {
577 case "function_call":
578 if !finishFunctionCall(ctx, out, callForItem(event.Item.ID), event.Item) {
579 return
580 }
581 case "reasoning":
582 // Preserve the done event's final status when round-tripping
583 // the reasoning item so the input matches the wire schema.
584 if event.Item.Status != "" {
585 reasoningStatus = event.Item.Status
586 }
587 }
588 }
589 case "response.completed", "response.incomplete", "response.failed":
590 terminal = true
591 if event.Type != "response.failed" {
592 for _, item := range unclosedOutputCalls(event.Response, calls) {
593 if !finishFunctionCall(ctx, out, callForItem(item.ID), item) {
594 return
595 }
596 }
597 }
598 if event.Type == "response.incomplete" {
599 if !sendChunk(ctx, out, provider.Chunk{Type: provider.ChunkReasoning, ReasoningState: provider.ReasoningIncomplete}) {
600 return
601 }
602 }
603 completedResponseID = terminalResponseID(event)
604 if !emitTerminalResponseUsage(ctx, out, event) {
605 return
606 }
607 if event.Type == "response.failed" {
608 failed = true
609 err := fmt.Errorf("responses: response failed")
610 if event.Response != nil && event.Response.Error != nil {
611 if authErr := authErrorFromResponse(c, event.Response.Error); authErr != nil {
612 err = authErr
613 } else {
614 err = fmt.Errorf("responses: %s", event.Response.Error.Message)
615 }
616 }
617 if !sendChunk(ctx, out, provider.Chunk{Type: provider.ChunkError, Err: err}) {
618 return
619 }
620 }
621 }
622 if terminal {
623 break
624 }
625 }
626
627 if ctx.Err() != nil {
628 return
629 }
630 if err := scanner.Err(); err != nil {
631 var reason string
632 if stalled.Load() {
633 err = fmt.Errorf("responses: stream idle timeout after %s", idle)
634 reason = provider.StreamInterruptIdleTimeout
635 } else {
636 reason = provider.ClassifyStreamInterrupt(err)
637 }
638 _ = sendChunk(ctx, out, provider.Chunk{Type: provider.ChunkError, Err: provider.StreamInterrupt(err, reason)})
639 return
640 }
641 // Protocol-defined terminal response events are required. Connection close
642 // before a terminal event leaves the attempt uncommitted — including any
643 // complete tool calls already forwarded as speculative output.
644 if !terminal {
645 _ = sendChunk(ctx, out, provider.Chunk{Type: provider.ChunkError, Err: provider.StreamInterrupt(io.ErrUnexpectedEOF, provider.StreamInterruptPrematureEOF)})
646 return
647 }
648 if !reasoningSnapshots.emit(ctx, out) {
649 return
650 }
651 responsesItems = append(responsesItems, reasoningSnapshots.items...)
652 if len(reasoningSnapshots.items) > 0 {
653 reasoningID, reasoningStatus = reasoningSnapshots.metadata()
654 }
655 if completedResponseID != "" {
656 assistant := provider.Message{Role: provider.RoleAssistant, Content: text.String(), ReasoningContent: reasoning.String(), ReasoningID: reasoningID, ReasoningStatus: reasoningStatus, ResponsesItems: responsesItems}
657 for _, itemID := range callOrder {
658 call := calls[itemID]
659 if call.completed {
660 assistant.ToolCalls = append(assistant.ToolCalls, provider.ToolCall{ID: call.id, Name: call.name, Arguments: call.arguments})
661 }
662 }
663 expected := append(append([]provider.Message(nil), requestMessages...), assistant)
664 c.mu.Lock()
665 c.lastResponseID = completedResponseID
666 c.expectedPrefixDigest = c.conversationDigest(expected)
667 c.mu.Unlock()
668 } else {
669 c.ResetContext()
670 }
671 if !failed {
672 // 把 reasoning item 的 id/status 作为元数据 chunk 流给 Agent
673 // (空 Text,随 ChunkReasoning 语义)——Agent 持久化进 session,
674 // 下一轮 input reasoning item 回传 id/status(评审 #7234 第 1 点)。
675 if reasoningID != "" || reasoningStatus != "" {
676 if !sendChunk(ctx, out, provider.Chunk{Type: provider.ChunkReasoning, ReasoningID: reasoningID, ReasoningStatus: reasoningStatus}) {
677 return
678 }
679 }
680 _ = sendChunk(ctx, out, provider.Chunk{Type: provider.ChunkDone})
681 }
682 }
683
684 func sendChunk(ctx context.Context, out chan<- provider.Chunk, chunk provider.Chunk) bool {
685 select {
686 case out <- chunk:
687 return true
688 default:
689 }
690 notifySendChunkEnterBlocking()
691 select {
692 case out <- chunk:
693 return true
694 case <-ctx.Done():
695 return false
696 }
697 }
698
699 func usageFromResponse(response *sseResponse) *provider.Usage {
700 usage := &provider.Usage{}
701 if response == nil || response.Usage == nil {
702 return usage
703 }
704 u := response.Usage
705 cached, reasoning := 0, 0
706 if u.InputTokensDetails != nil {
707 cached = u.InputTokensDetails.CachedTokens
708 }
709 if u.OutputTokensDetails != nil {
710 reasoning = u.OutputTokensDetails.ReasoningTokens
711 }
712 miss := max(u.InputTokens-cached, 0)
713 total := u.TotalTokens
714 if total == 0 {
715 total = u.InputTokens + u.OutputTokens
716 }
717 return &provider.Usage{PromptTokens: u.InputTokens, CompletionTokens: u.OutputTokens, TotalTokens: total, CacheHitTokens: cached, CacheMissTokens: miss, ReasoningTokens: reasoning}
718 }
719
720 func authErrorFromResponse(c *client, responseError *sseError) error {
721 if responseError == nil {
722 return nil
723 }
724 value := strings.ToLower(responseError.Code + " " + responseError.Message)
725 if !strings.Contains(value, "auth") && !strings.Contains(value, "api key") && !strings.Contains(value, "unauthorized") && !strings.Contains(value, "forbidden") && !strings.Contains(value, "permission") {
726 return nil
727 }
728 status := http.StatusUnauthorized
729 if strings.Contains(value, "forbidden") || strings.Contains(value, "permission") {
730 status = http.StatusForbidden
731 }
732 return &provider.AuthError{Provider: c.name, ProviderDisplayName: c.identity.DisplayName, Protocol: c.identity.Protocol, KeyEnv: c.keyEnv, KeySource: c.keySource, Status: status, HasKey: c.apiKey != "", Body: responseError.Message}
733 }
734
735 type sseEvent struct {
736 Type string `json:"type"`
737 Delta string `json:"delta"`
738 Text string `json:"text"`
739 Arguments string `json:"arguments"`
740 ItemID string `json:"item_id"`
741 ContentIndex int `json:"content_index"`
742 Item *sseItem `json:"item"`
743 Response *sseResponse `json:"response"`
744 }
745
746 type sseItem struct {
747 ID, Type, CallID, Name, Arguments, Status string
748 Raw json.RawMessage
749 }
750
751 func (i *sseItem) UnmarshalJSON(data []byte) error {
752 var wire struct {
753 ID string `json:"id"`
754 Type string `json:"type"`
755 CallID string `json:"call_id"`
756 Name string `json:"name"`
757 Arguments string `json:"arguments"`
758 Status string `json:"status"`
759 }
760 if err := json.Unmarshal(data, &wire); err != nil {
761 return err
762 }
763 *i = sseItem{ID: wire.ID, Type: wire.Type, CallID: wire.CallID, Name: wire.Name, Arguments: wire.Arguments, Status: wire.Status, Raw: append(json.RawMessage(nil), data...)}
764 return nil
765 }
766
767 type sseResponse struct {
768 Output []json.RawMessage `json:"output"`
769 ID string `json:"id"`
770 Usage *sseUsage `json:"usage"`
771 Error *sseError `json:"error"`
772 IncompleteDetails incompleteDetails `json:"incomplete_details"`
773 }
774
775 type incompleteDetails struct {
776 Reason string `json:"reason"`
777 }
778 type sseError struct {
779 Message string `json:"message"`
780 Code string `json:"code"`
781 }
782 type sseUsage struct {
783 InputTokens int `json:"input_tokens"`
784 OutputTokens int `json:"output_tokens"`
785 TotalTokens int `json:"total_tokens"`
786 InputTokensDetails *struct {
787 CachedTokens int `json:"cached_tokens"`
788 } `json:"input_tokens_details"`
789 OutputTokensDetails *struct {
790 ReasoningTokens int `json:"reasoning_tokens"`
791 } `json:"output_tokens_details"`
792 }
793
793 lines GO