返回 DeepSeek-Reasonix
live_official_recovery_matrix_test.go
根目录 / internal / agent / live_official_recovery_matrix_test.go
1 //go:build live
2
3 package agent
4
5 import (
6 "bytes"
7 "context"
8 "encoding/json"
9 "errors"
10 "io"
11 "net/http"
12 "net/http/httptest"
13 "os"
14 "strings"
15 "sync"
16 "sync/atomic"
17 "testing"
18 "time"
19
20 "reasonix/internal/event"
21 "reasonix/internal/provider"
22 "reasonix/internal/provider/anthropic"
23 "reasonix/internal/provider/openai"
24 "reasonix/internal/provider/responses"
25 "reasonix/internal/tool"
26 )
27
28 // Only synthetic echo is exposed: no shell, filesystem reader, MCP, or access
29 // to credentials. Faults affect real completed upstream streams locally, not
30 // the public service. Raw requests/responses remain in memory.
31 func TestLiveOfficialRecoveryMatrix(t *testing.T) {
32 key := os.Getenv("DEEPSEEK_API_KEY")
33 if key == "" {
34 t.Skip("DEEPSEEK_API_KEY not set")
35 }
36 for _, model := range []string{"deepseek-v4-flash", "deepseek-v4-pro"} {
37 for _, protocol := range []string{"chat", "responses", "anthropic"} {
38 for _, scenario := range []string{"low", "max", "disabled", "cut_once", "missing_once", "missing_persistent", "followup_503", "server_replay_rejection", "cancel_before_commit"} {
39 t.Run(model+"/"+protocol+"/"+scenario, func(t *testing.T) {
40 ctx, cancel := context.WithTimeout(context.Background(), 120*time.Second)
41 defer cancel()
42 proxy := &officialRecoveryProxy{protocol: protocol, scenario: scenario, cancel: cancel}
43 srv := httptest.NewServer(proxy)
44 defer srv.Close()
45 effort := "high"
46 if scenario == "low" || scenario == "max" || scenario == "disabled" {
47 effort = scenario
48 }
49 p := officialMatrixProvider(t, key, model, protocol, effort, srv.URL)
50 var executions atomic.Int32
51 reg := tool.NewRegistry()
52 reg.Add(liveRecoveryEchoTool{executions: &executions})
53 sink := &recordSink{}
54 session := NewSession("You are a concise tool-using assistant. Call echo exactly once when asked, then report the marker. Do not repeat a completed tool.")
55 a := New(p, reg, session, Options{MaxSteps: 4, MaxOutputTokens: 2048, MissingReasoningWarnStateDir: t.TempDir()}, sink)
56 start := time.Now()
57 err := a.Run(ctx, "Call echo exactly once, then report its result.")
58 wantRequests, wantExecutions := 2, int32(1)
59 switch scenario {
60 case "cut_once":
61 wantRequests, wantExecutions = 1, 0
62 if !provider.IsStreamInterrupted(err) {
63 t.Fatalf("expected stream failure: %v", err)
64 }
65 case "followup_503":
66 var api *provider.APIError
67 if !errors.As(err, &api) || api.Status != 503 {
68 t.Fatalf("expected upstream failure: %v", err)
69 }
70 case "server_replay_rejection":
71 wantRequests = 3
72 case "missing_once":
73 if protocol == "anthropic" {
74 wantRequests = 3
75 }
76 case "missing_persistent":
77 if protocol == "anthropic" {
78 wantExecutions = 0
79 if err == nil {
80 t.Fatal("strict missing proof accepted")
81 }
82 }
83 case "cancel_before_commit":
84 wantRequests, wantExecutions = 1, 0
85 if !errors.Is(err, context.Canceled) {
86 t.Fatalf("cancellation error=%v", err)
87 }
88 }
89 stopped := scenario == "cut_once" || scenario == "followup_503" || scenario == "cancel_before_commit" || scenario == "missing_persistent" && protocol == "anthropic"
90 if !stopped && err != nil {
91 t.Fatalf("live run: %v", err)
92 }
93 if scenario == "server_replay_rejection" && err == nil {
94 t.Logf("before_continuation_tool_executions=%d", executions.Load())
95 if err := a.Run(ctx, "Without calling tools, report the already known marker from the completed request."); err != nil {
96 t.Fatalf("post-repair continuation: %v", err)
97 }
98 wantRequests = 4
99 }
100 proxy.mu.Lock()
101 requests, upstream, mutations := proxy.requests, proxy.upstream, proxy.mutations
102 proxy.mu.Unlock()
103 if requests != wantRequests || executions.Load() != wantExecutions {
104 t.Fatalf("requests=%d executions=%d want=%d/%d err=%v", requests, executions.Load(), wantRequests, wantExecutions, err)
105 }
106 if strings.HasPrefix(scenario, "missing") || scenario == "cut_once" {
107 if mutations == 0 {
108 t.Fatal("requested fault was not injected")
109 }
110 }
111 if (scenario == "cut_once" || scenario == "followup_503") && len(sink.kinds(event.Retrying)) != 0 {
112 t.Fatal("transport failure retried automatically")
113 }
114 if !stopped {
115 msgs := session.Snapshot()
116 if len(msgs) == 0 || strings.TrimSpace(msgs[len(msgs)-1].Content) == "" {
117 t.Fatal("no final text")
118 }
119 }
120 prompt, completion, accounted, unknown := 0, 0, 0, false
121 for _, e := range sink.kinds(event.Usage) {
122 if u := e.Usage; u != nil {
123 prompt += u.PromptTokens
124 completion += u.CompletionTokens
125 accounted += u.RequestCount
126 unknown = unknown || u.Unknown
127 }
128 }
129 if (scenario == "cut_once" || scenario == "cancel_before_commit" || scenario == "followup_503" || scenario == "server_replay_rejection") && !unknown {
130 t.Fatal("missing terminal usage must stay unknown")
131 }
132 if accounted != requests {
133 t.Fatalf("accounted=%d HTTP=%d", accounted, requests)
134 }
135 t.Logf("protocol=%s scenario=%s http_attempts=%d upstream_requests=%d mutations=%d tool_executions=%d retry_events=%d prompt=%d completion=%d accounted_requests=%d usage_unknown=%t elapsed_ms=%d", protocol, scenario, requests, upstream, mutations, executions.Load(), len(sink.kinds(event.Retrying)), prompt, completion, accounted, unknown, time.Since(start).Milliseconds())
136 })
137 }
138 }
139 }
140 }
141
142 func officialMatrixProvider(t *testing.T, key, model, protocol, effort, url string) provider.Provider {
143 t.Helper()
144 extra := map[string]any{"api_key_env": "DEEPSEEK_API_KEY", "request_url": url, "reasoning_protocol": "deepseek", "thinking": "enabled", "effort": effort, "reject_redirects": true}
145 if effort == "disabled" {
146 extra["thinking"] = "disabled"
147 }
148 var p provider.Provider
149 var err error
150 switch protocol {
151 case "chat":
152 p, err = openai.New(provider.Config{Name: "official-live-chat", BaseURL: "https://api.deepseek.com", Model: model, APIKey: key, Extra: extra})
153 case "anthropic":
154 p, err = anthropic.New(provider.Config{Name: "official-live-anthropic", BaseURL: "https://api.deepseek.com/anthropic", Model: model, APIKey: key, Extra: extra})
155 case "responses":
156 p = responses.New(responses.Config{Name: "official-live-responses", BaseURL: "https://api.deepseek.com", RequestURL: url, Model: model, APIKey: key, KeyEnv: "DEEPSEEK_API_KEY", Effort: effort, Mode: "stateless", MaxOutputTokens: 2048})
157 default:
158 t.Fatal("unsupported test protocol")
159 }
160 if err != nil {
161 t.Fatal(err)
162 }
163 if c, ok := p.(interface{ CloseIdleConnections() }); ok {
164 t.Cleanup(c.CloseIdleConnections)
165 }
166 return p
167 }
168
169 type officialRecoveryProxy struct {
170 mu sync.Mutex
171 protocol, scenario string
172 upstreamURL, sessionHeader string
173 statuses []int
174 wireTools, wireStops []string
175 requests, upstream, mutations int
176 bodies [][]byte
177 searchResponses [][]byte
178 cancel context.CancelFunc
179 }
180
181 func (p *officialRecoveryProxy) ServeHTTP(w http.ResponseWriter, r *http.Request) {
182 body, err := io.ReadAll(r.Body)
183 if err != nil {
184 http.Error(w, "read test request", 400)
185 return
186 }
187 p.mu.Lock()
188 p.requests++
189 n := p.requests
190 p.bodies = append(p.bodies, body)
191 p.mu.Unlock()
192 if n > 16 || n > 5 && p.scenario != "continuity" {
193 http.Error(w, "live test request limit", 401)
194 return
195 }
196 if p.scenario == "followup_503" && n == 2 {
197 w.WriteHeader(503)
198 _, _ = io.WriteString(w, `{"error":{"message":"injected temporary service failure"}}`)
199 return
200 }
201 // Drop proof and replace the call identity only in the outbound request.
202 // This forces the official service to validate unavailable replay content,
203 // while the canonical local transcript still retains the completed tool.
204 if p.scenario == "server_replay_rejection" && n == 2 {
205 body = breakOfficialReplay(p.protocol, body)
206 }
207 path := map[string]string{"chat": "/chat/completions", "responses": "/responses", "anthropic": "/anthropic/v1/messages"}[p.protocol]
208 target := p.upstreamURL
209 if target == "" {
210 target = "https://api.deepseek.com" + path
211 }
212 req, err := http.NewRequestWithContext(r.Context(), http.MethodPost, target, bytes.NewReader(body))
213 if err != nil {
214 http.Error(w, "create upstream", 500)
215 return
216 }
217 for _, name := range []string{"Authorization", "x-api-key", "anthropic-version", "Content-Type", "User-Agent", "x-opencode-session"} {
218 if v := r.Header.Get(name); v != "" {
219 req.Header.Set(name, v)
220 }
221 }
222 if p.sessionHeader != "" {
223 req.Header.Set("User-Agent", "Reasonix/live-validation")
224 req.Header.Set("x-opencode-session", p.sessionHeader)
225 }
226 client := &http.Client{Timeout: 90 * time.Second, CheckRedirect: func(*http.Request, []*http.Request) error { return http.ErrUseLastResponse }}
227 p.mu.Lock()
228 p.upstream++
229 p.mu.Unlock()
230 resp, err := client.Do(req)
231 if err != nil {
232 http.Error(w, "upstream failed", 502)
233 return
234 }
235 defer resp.Body.Close()
236 p.mu.Lock()
237 p.statuses = append(p.statuses, resp.StatusCode)
238 p.mu.Unlock()
239 data, err := io.ReadAll(resp.Body)
240 if err != nil {
241 http.Error(w, "upstream read failed", 502)
242 return
243 }
244 if p.scenario == "search" {
245 p.mu.Lock()
246 p.searchResponses = append(p.searchResponses, bytes.Clone(data))
247 p.mu.Unlock()
248 }
249 if p.protocol == "chat" && resp.StatusCode == 200 {
250 var names, stops []string
251 for _, line := range bytes.Split(data, []byte("\n")) {
252 if !bytes.HasPrefix(line, []byte("data:")) {
253 continue
254 }
255 var frame struct {
256 Choices []struct {
257 FinishReason string `json:"finish_reason"`
258 Delta struct {
259 ToolCalls []struct {
260 Function struct {
261 Name string `json:"name"`
262 } `json:"function"`
263 } `json:"tool_calls"`
264 } `json:"delta"`
265 } `json:"choices"`
266 }
267 if json.Unmarshal(bytes.TrimSpace(bytes.TrimPrefix(line, []byte("data:"))), &frame) != nil {
268 continue
269 }
270 for _, choice := range frame.Choices {
271 if choice.FinishReason != "" {
272 stops = append(stops, choice.FinishReason)
273 }
274 for _, call := range choice.Delta.ToolCalls {
275 if call.Function.Name != "" {
276 names = append(names, call.Function.Name)
277 }
278 }
279 }
280 }
281 p.mu.Lock()
282 p.wireTools = append(p.wireTools, names...)
283 p.wireStops = append(p.wireStops, stops...)
284 p.mu.Unlock()
285 }
286 if resp.StatusCode == 200 {
287 if p.scenario == "cancel_before_commit" && n == 1 {
288 p.cancel()
289 }
290 if p.scenario == "cut_once" && n == 1 {
291 data = data[:len(data)/2]
292 p.mu.Lock()
293 p.mutations++
294 p.mu.Unlock()
295 }
296 if strings.HasPrefix(p.scenario, "missing") && (n == 1 || p.scenario == "missing_persistent") {
297 switch p.protocol {
298 case "chat":
299 s := &liveReasoningStripProxy{}
300 data = s.stripReasoning(data)
301 p.mu.Lock()
302 p.mutations += int(s.strippedFields.Load())
303 p.mu.Unlock()
304 case "responses":
305 var count int
306 data, count = stripResponsesReasoningEvents(data)
307 p.mu.Lock()
308 p.mutations += count
309 p.mu.Unlock()
310 case "anthropic":
311 var count int
312 data, count = stripOfficialThinking(data)
313 p.mu.Lock()
314 p.mutations += count
315 p.mu.Unlock()
316 }
317 }
318 }
319 w.Header().Set("Content-Type", resp.Header.Get("Content-Type"))
320 w.WriteHeader(resp.StatusCode)
321 _, _ = w.Write(data)
322 }
323 func stripOfficialThinking(data []byte) ([]byte, int) {
324 var out bytes.Buffer
325 count := 0
326 thinking := map[int]bool{}
327 for _, frame := range bytes.Split(data, []byte("\n\n")) {
328 skip := false
329 for _, line := range bytes.Split(frame, []byte("\n")) {
330 if !bytes.HasPrefix(line, []byte("data:")) {
331 continue
332 }
333 var e struct {
334 Type string `json:"type"`
335 Index int `json:"index"`
336 Block struct {
337 Type string `json:"type"`
338 } `json:"content_block"`
339 }
340 if json.Unmarshal(bytes.TrimSpace(bytes.TrimPrefix(line, []byte("data:"))), &e) != nil {
341 continue
342 }
343 if e.Type == "content_block_start" && e.Block.Type == "thinking" {
344 thinking[e.Index] = true
345 }
346 if strings.HasPrefix(e.Type, "content_block_") && thinking[e.Index] {
347 skip = true
348 }
349 }
350 if skip {
351 count++
352 continue
353 }
354 out.Write(frame)
355 out.WriteString("\n\n")
356 }
357 return out.Bytes(), count
358 }
359
360 func breakOfficialReplay(protocol string, body []byte) []byte {
361 var request map[string]any
362 if json.Unmarshal(body, &request) != nil {
363 return body
364 }
365 const replacement = "call_live_missing_proof"
366 if protocol == "responses" {
367 var input []any
368 for _, v := range request["input"].([]any) {
369 item := v.(map[string]any)
370 if item["type"] == "reasoning" {
371 continue
372 }
373 if item["type"] == "function_call" {
374 item["id"] = "fc_live_missing_proof"
375 item["call_id"] = replacement
376 }
377 if item["type"] == "function_call_output" {
378 item["call_id"] = replacement
379 }
380 input = append(input, item)
381 }
382 request["input"] = input
383 } else {
384 for _, v := range request["messages"].([]any) {
385 msg := v.(map[string]any)
386 if protocol == "chat" {
387 delete(msg, "reasoning_content")
388 if calls, ok := msg["tool_calls"].([]any); ok {
389 for _, c := range calls {
390 c.(map[string]any)["id"] = replacement
391 }
392 }
393 if msg["role"] == "tool" {
394 msg["tool_call_id"] = replacement
395 }
396 } else if blocks, ok := msg["content"].([]any); ok {
397 var content []any
398 for _, v := range blocks {
399 b := v.(map[string]any)
400 if b["type"] == "thinking" {
401 continue
402 }
403 if b["type"] == "tool_use" {
404 b["id"] = replacement
405 }
406 if b["type"] == "tool_result" {
407 b["tool_use_id"] = replacement
408 }
409 content = append(content, b)
410 }
411 msg["content"] = content
412 }
413 }
414 }
415 encoded, err := json.Marshal(request)
416 if err != nil {
417 return body
418 }
419 return encoded
420 }
421
421 lines GO