返回 DeepSeek-Reasonix
retry_e2e_test.go
根目录 / internal / agent / retry_e2e_test.go
1 package agent
2
3 import (
4 "bytes"
5 "context"
6 "encoding/json"
7 "errors"
8 "io"
9 "net/http"
10 "net/http/httptest"
11 "reflect"
12 "strings"
13 "sync"
14 "sync/atomic"
15 "testing"
16 "time"
17
18 "reasonix/internal/event"
19 "reasonix/internal/provider"
20 "reasonix/internal/provider/anthropic"
21 "reasonix/internal/provider/openai"
22 "reasonix/internal/tool"
23 )
24
25 type recordSink struct {
26 mu sync.Mutex
27 evs []event.Event
28 recovery []event.ProtocolRecoveryAudit
29 }
30
31 type textSignalSink struct {
32 *recordSink
33 textSeen chan struct{}
34 once sync.Once
35 }
36
37 func (s *textSignalSink) Emit(e event.Event) {
38 s.recordSink.Emit(e)
39 if e.Kind == event.Text {
40 s.once.Do(func() { close(s.textSeen) })
41 }
42 }
43
44 func (s *recordSink) Emit(e event.Event) {
45 s.mu.Lock()
46 defer s.mu.Unlock()
47 s.evs = append(s.evs, e)
48 }
49
50 func (s *recordSink) kinds(k event.Kind) []event.Event {
51 s.mu.Lock()
52 defer s.mu.Unlock()
53 var out []event.Event
54 for _, e := range s.evs {
55 if e.Kind == k {
56 out = append(out, e)
57 }
58 }
59 return out
60 }
61
62 func (s *recordSink) RecordProtocolRecovery(a event.ProtocolRecoveryAudit) {
63 s.mu.Lock()
64 defer s.mu.Unlock()
65 s.recovery = append(s.recovery, a)
66 }
67
68 func (s *recordSink) recoveryCount(kind event.ProtocolRecoveryKind) int {
69 s.mu.Lock()
70 defer s.mu.Unlock()
71 var count int
72 for _, audit := range s.recovery {
73 if audit.Kind == kind {
74 count++
75 }
76 }
77 return count
78 }
79
80 // TestAgentReturnsHTTPFailureWithoutRetry exercises the real OpenAI adapter.
81 func TestAgentReturnsHTTPFailureWithoutRetry(t *testing.T) {
82 for _, status := range []int{429, 503} {
83 t.Run(http.StatusText(status), func(t *testing.T) {
84 var calls atomic.Int32
85 srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
86 calls.Add(1)
87 w.Header().Set("Retry-After", "120")
88 w.WriteHeader(status)
89 _, _ = io.WriteString(w, `{"error":{"message":"upstream unavailable"}}`)
90 }))
91 defer srv.Close()
92 p, err := openai.New(provider.Config{Name: "deepseek", BaseURL: srv.URL, Model: "deepseek-v4-flash", APIKey: "k", HTTPClient: srv.Client()})
93 if err != nil {
94 t.Fatal(err)
95 }
96 sink := &recordSink{}
97 a := New(p, tool.NewRegistry(), NewSession(""), Options{}, sink)
98 err = a.Run(withNoClosedLoop(t.Context()), "hi")
99 var api *provider.APIError
100 if !errors.As(err, &api) || api.Status != status || !strings.Contains(api.Body, "upstream unavailable") {
101 t.Fatalf("err=%v", err)
102 }
103 if calls.Load() != 1 || len(sink.kinds(event.Retrying)) != 0 {
104 t.Fatalf("calls=%d retries=%+v", calls.Load(), sink.kinds(event.Retrying))
105 }
106 })
107 }
108 }
109
110 // TestDeepSeekFlashMissingReasoningRecoveryWithRealSSE exercises the actual
111 // OpenAI-compatible decoder shape used by the official Flash endpoint. The
112 // first response emits a tool call without reasoning_content; the second exact
113 // request includes it; only the adopted call reaches the session and UI.
114 func TestDeepSeekFlashMissingReasoningRecoveryWithRealSSE(t *testing.T) {
115 var mu sync.Mutex
116 var bodies [][]byte
117 srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
118 body, _ := io.ReadAll(r.Body)
119 mu.Lock()
120 bodies = append(bodies, append([]byte(nil), body...))
121 requestNo := len(bodies)
122 mu.Unlock()
123
124 w.Header().Set("Content-Type", "text/event-stream")
125
126 if requestNo == 1 {
127 _, _ = io.WriteString(w, `data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"c1","type":"function","function":{"name":"echo","arguments":"{\"text\":\"hi\"}"}}]},"finish_reason":"tool_calls"}]}`+"\n\n")
128 _, _ = io.WriteString(w, `data: {"choices":[],"usage":{"prompt_tokens":10,"completion_tokens":2,"total_tokens":12,"prompt_cache_hit_tokens":10,"prompt_cache_miss_tokens":0}}`+"\n\n")
129 _, _ = io.WriteString(w, "data: [DONE]\n\n")
130 return
131 }
132 _, _ = io.WriteString(w, `data: {"choices":[{"delta":{"content":"done"},"finish_reason":"stop"}]}`+"\n\n")
133 _, _ = io.WriteString(w, "data: [DONE]\n\n")
134 }))
135 defer srv.Close()
136
137 prov, err := openai.New(provider.Config{
138 Name: "deepseek", BaseURL: srv.URL, Model: "deepseek-v4-flash", APIKey: "k",
139 Extra: map[string]any{"reasoning_protocol": "deepseek", "thinking": "enabled"},
140 })
141 if err != nil {
142 t.Fatalf("New provider: %v", err)
143 }
144 sink := &recordSink{}
145 a := New(prov, echoRegistry(), NewSession(""), Options{}, sink)
146 if err := a.Run(context.Background(), "go"); err != nil {
147 t.Fatalf("Run: %v", err)
148 }
149
150 mu.Lock()
151 gotBodies := append([][]byte(nil), bodies...)
152 mu.Unlock()
153 if len(gotBodies) != 2 {
154 t.Fatalf("HTTP requests = %d, want malformed + recovery + final", len(gotBodies))
155 }
156 if !bytes.Contains(gotBodies[1], []byte(`"role":"tool"`)) {
157 t.Fatalf("recovery request changed bytes:\nfirst=%s\nretry=%s", gotBodies[0], gotBodies[1])
158 }
159 var toolTurns int
160 for _, message := range a.Session().Messages {
161 if message.Role == provider.RoleAssistant && len(message.ToolCalls) > 0 {
162 toolTurns++
163 if message.ReasoningContent != "" {
164 t.Fatalf("adopted reasoning = %q", message.ReasoningContent)
165 }
166 }
167 }
168 if toolTurns != 1 {
169 t.Fatalf("saved tool turns = %d, want 1", toolTurns)
170 }
171 // One partial dispatch from the adopted SSE plus one full execution
172 // dispatch. The discarded malformed stream must not add a third card.
173 if got := len(sink.kinds(event.ToolDispatch)); got != 2 {
174 t.Fatalf("tool dispatch events = %d, want adopted partial + full", got)
175 }
176 for _, notice := range sink.kinds(event.Notice) {
177 if strings.Contains(notice.Text, "reasoning_content") || strings.Contains(notice.Detail, "reasoning_content") {
178 t.Fatalf("protocol warning leaked to UI: %+v", notice)
179 }
180 }
181 }
182
183 // TestDeepSeekOpenAIReasoningReplay400RepairsOldHistory drives the OpenAI
184 // adapter through the shared stale-history recovery path. The first request
185 // replays an old assistant reasoning turn and is rejected; the repair retry
186 // empties provider-visible reasoning while preserving canonical history. The
187 // key itself stays: DeepSeek thinking mode rejects an assistant turn without it.
188 func TestDeepSeekOpenAIReasoningReplay400RepairsOldHistory(t *testing.T) {
189 var mu sync.Mutex
190 var bodies [][]byte
191 srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
192 body, _ := io.ReadAll(r.Body)
193 mu.Lock()
194 bodies = append(bodies, append([]byte(nil), body...))
195 requestNo := len(bodies)
196 mu.Unlock()
197
198 if requestNo == 1 {
199 w.WriteHeader(http.StatusBadRequest)
200 _, _ = io.WriteString(w, `{"error":{"message":"The reasoning_content in the thinking mode must be passed back to the API"}}`)
201 return
202 }
203 w.Header().Set("Content-Type", "text/event-stream")
204 _, _ = io.WriteString(w, `data: {"choices":[{"delta":{"content":"recovered"},"finish_reason":"stop"}]}`+"\n\n")
205 _, _ = io.WriteString(w, "data: [DONE]\n\n")
206 }))
207 defer srv.Close()
208
209 prov, err := openai.New(provider.Config{
210 Name: "deepseek-openai", BaseURL: srv.URL, Model: "deepseek-v4-pro", APIKey: "k",
211 Extra: map[string]any{"reasoning_protocol": "deepseek", "thinking": "enabled"},
212 })
213 if err != nil {
214 t.Fatalf("New provider: %v", err)
215 }
216 session := NewSession("system")
217 session.Add(provider.Message{Role: provider.RoleUser, Content: "earlier"})
218 session.Add(provider.Message{Role: provider.RoleAssistant, Content: "old answer", ReasoningContent: "stale thinking"})
219 sink := &recordSink{}
220 a := New(prov, echoRegistry(), session, Options{}, sink)
221
222 if err := a.Run(withNoClosedLoop(context.Background()), "next"); err != nil {
223 t.Fatalf("Run: %v", err)
224 }
225
226 mu.Lock()
227 gotBodies := append([][]byte(nil), bodies...)
228 mu.Unlock()
229 if len(gotBodies) != 2 {
230 t.Fatalf("HTTP requests = %d, want rejected attempt plus one repair retry", len(gotBodies))
231 }
232 if !bytes.Contains(gotBodies[0], []byte(`"reasoning_content":"stale thinking"`)) {
233 t.Fatalf("first request did not replay old reasoning: %s", gotBodies[0])
234 }
235 if bytes.Contains(gotBodies[1], []byte("stale thinking")) || !bytes.Contains(gotBodies[1], []byte(`"content":"old answer","reasoning_content":""`)) {
236 t.Fatalf("repair retry still carries old reasoning: %s", gotBodies[1])
237 }
238 if !bytes.Contains(gotBodies[1], []byte("old answer")) {
239 t.Fatalf("repair retry lost visible assistant text: %s", gotBodies[1])
240 }
241 var first, second map[string]json.RawMessage
242 if err := json.Unmarshal(gotBodies[0], &first); err != nil {
243 t.Fatalf("decode first request: %v", err)
244 }
245 if err := json.Unmarshal(gotBodies[1], &second); err != nil {
246 t.Fatalf("decode repair request: %v", err)
247 }
248 delete(first, "messages")
249 delete(second, "messages")
250 if !reflect.DeepEqual(first, second) {
251 t.Fatalf("repair retry changed non-message fields:\nfirst=%s\nretry=%s", gotBodies[0], gotBodies[1])
252 }
253 for _, message := range session.Snapshot() {
254 if message.Role == provider.RoleAssistant && message.Content == "old answer" && message.ReasoningContent != "stale thinking" {
255 t.Fatalf("canonical history lost old reasoning: %+v", message)
256 }
257 }
258 if got := sink.recoveryCount(event.ProtocolRecoveryReasoningReplay400Detected); got != 1 {
259 t.Fatalf("reasoning_replay_400_detected audits = %d, want 1", got)
260 }
261 if got := sink.recoveryCount(event.ProtocolRecoveryReasoningReplay400Recovered); got != 1 {
262 t.Fatalf("reasoning_replay_400_recovered audits = %d, want 1", got)
263 }
264 }
265
266 func TestDeepSeekOpenAIReasoningReplay400StripsOldToolHistory(t *testing.T) {
267 var mu sync.Mutex
268 var bodies [][]byte
269 srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
270 body, _ := io.ReadAll(r.Body)
271 mu.Lock()
272 bodies = append(bodies, append([]byte(nil), body...))
273 requestNo := len(bodies)
274 mu.Unlock()
275 if requestNo == 1 {
276 w.WriteHeader(http.StatusBadRequest)
277 _, _ = io.WriteString(w, `{"error":{"message":"The reasoning_content in the thinking mode must be passed back to the API"}}`)
278 return
279 }
280 w.Header().Set("Content-Type", "text/event-stream")
281 _, _ = io.WriteString(w, `data: {"choices":[{"delta":{"content":"recovered"},"finish_reason":"stop"}]}`+"\n\n")
282 _, _ = io.WriteString(w, "data: [DONE]\n\n")
283 }))
284 defer srv.Close()
285
286 prov, err := openai.New(provider.Config{
287 Name: "deepseek-openai", BaseURL: srv.URL, Model: "deepseek-v4-pro", APIKey: "k",
288 Extra: map[string]any{"reasoning_protocol": "deepseek", "thinking": "enabled"},
289 })
290 if err != nil {
291 t.Fatalf("New provider: %v", err)
292 }
293 session := NewSession("system")
294 session.Add(provider.Message{Role: provider.RoleUser, Content: "earlier"})
295 session.Add(provider.Message{
296 Role: provider.RoleAssistant, Content: "I will inspect the file", ReasoningContent: "stale thinking",
297 ToolCalls: []provider.ToolCall{{ID: "old-call", Name: "read_file", Arguments: `{"path":"old.go"}`}},
298 })
299 session.Add(provider.Message{Role: provider.RoleTool, ToolCallID: "old-call", Name: "read_file", Content: "old result"})
300 a := New(prov, echoRegistry(), session, Options{}, &recordSink{})
301 if err := a.Run(withNoClosedLoop(context.Background()), "next"); err != nil {
302 t.Fatalf("Run: %v", err)
303 }
304
305 mu.Lock()
306 gotBodies := append([][]byte(nil), bodies...)
307 mu.Unlock()
308 if len(gotBodies) != 2 {
309 t.Fatalf("HTTP requests = %d, want rejected attempt plus one repair retry", len(gotBodies))
310 }
311 if !bytes.Contains(gotBodies[0], []byte("old-call")) || !bytes.Contains(gotBodies[0], []byte("stale thinking")) {
312 t.Fatalf("first request did not contain old tool history: %s", gotBodies[0])
313 }
314 if bytes.Contains(gotBodies[1], []byte(`"tool_calls"`)) || bytes.Contains(gotBodies[1], []byte(`"role":"tool"`)) || bytes.Contains(gotBodies[1], []byte("stale thinking")) {
315 t.Fatalf("repair retry retained stale tool history: %s", gotBodies[1])
316 }
317 if !bytes.Contains(gotBodies[1], []byte("old result")) {
318 t.Fatal("repair lost the completed tool output")
319 }
320 if !bytes.Contains(gotBodies[1], []byte("I will inspect the file")) {
321 t.Fatalf("repair retry lost visible old assistant text: %s", gotBodies[1])
322 }
323 }
324
325 func TestGLMToolTurnWithoutReasoningContinuesWithoutRecovery(t *testing.T) {
326 var mu sync.Mutex
327 var bodies [][]byte
328 srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
329 body, _ := io.ReadAll(r.Body)
330 mu.Lock()
331 bodies = append(bodies, append([]byte(nil), body...))
332 requestNo := len(bodies)
333 mu.Unlock()
334
335 w.Header().Set("Content-Type", "text/event-stream")
336 if requestNo == 1 {
337 _, _ = io.WriteString(w, `data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"c1","type":"function","function":{"name":"echo","arguments":"{\"text\":\"hi\"}"}}]},"finish_reason":"tool_calls"}]}`+"\n\n")
338 } else {
339 _, _ = io.WriteString(w, `data: {"choices":[{"delta":{"content":"done"},"finish_reason":"stop"}]}`+"\n\n")
340 }
341 _, _ = io.WriteString(w, "data: [DONE]\n\n")
342 }))
343 defer srv.Close()
344
345 prov, err := openai.New(provider.Config{
346 Name: "glm", BaseURL: srv.URL, Model: "glm-5.2", APIKey: "k",
347 Extra: map[string]any{"reasoning_protocol": "glm"},
348 })
349 if err != nil {
350 t.Fatalf("New provider: %v", err)
351 }
352 sink := &recordSink{}
353 a := New(prov, echoRegistry(), NewSession(""), Options{}, sink)
354 if err := a.Run(context.Background(), "go"); err != nil {
355 t.Fatalf("Run: %v", err)
356 }
357
358 mu.Lock()
359 requestBodies := append([][]byte(nil), bodies...)
360 mu.Unlock()
361 if len(requestBodies) != 2 {
362 t.Fatalf("HTTP requests = %d, want tool turn and final turn without recovery", len(requestBodies))
363 }
364 if !bytes.Contains(requestBodies[1], []byte(`"reasoning_content":""`)) {
365 t.Fatal("GLM replay did not preserve the empty reasoning_content field required for tool history")
366 }
367 if got := sink.recoveryCount(event.ProtocolRecoveryMissingReasoningRetryAttempted); got != 0 {
368 t.Fatalf("missing-reasoning retries = %d, want 0", got)
369 }
370 if got := len(sink.kinds(event.ToolResult)); got != 1 {
371 t.Fatalf("tool results = %d, want 1", got)
372 }
373 }
374
375 func TestGLMTextWithoutReasoningStreamsBeforeResponseCompletes(t *testing.T) {
376 responseStarted := make(chan struct{})
377 releaseResponse := make(chan struct{})
378 var releaseOnce sync.Once
379
380 srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
381 w.Header().Set("Content-Type", "text/event-stream")
382 _, _ = io.WriteString(w, `data: {"choices":[{"delta":{"content":"streamed"}}]}`+"\n\n")
383 w.(http.Flusher).Flush()
384 close(responseStarted)
385 <-releaseResponse
386 _, _ = io.WriteString(w, "data: [DONE]\n\n")
387 }))
388 t.Cleanup(srv.Close)
389 t.Cleanup(func() { releaseOnce.Do(func() { close(releaseResponse) }) })
390
391 prov, err := openai.New(provider.Config{
392 Name: "glm", BaseURL: srv.URL, Model: "glm-5.2", APIKey: "k",
393 Extra: map[string]any{"reasoning_protocol": "glm"},
394 })
395 if err != nil {
396 t.Fatalf("New provider: %v", err)
397 }
398 sink := &textSignalSink{recordSink: &recordSink{}, textSeen: make(chan struct{})}
399 a := New(prov, tool.NewRegistry(), NewSession(""), Options{}, sink)
400 done := make(chan error, 1)
401 go func() { done <- a.Run(withNoClosedLoop(context.Background()), "reply with streamed") }()
402
403 <-responseStarted
404 select {
405 case <-sink.textSeen:
406 case err := <-done:
407 t.Fatalf("Run completed before the held response was released: %v", err)
408 case <-time.After(2 * time.Second):
409 t.Fatal("GLM text stayed buffered until the response completed")
410 }
411 releaseOnce.Do(func() { close(releaseResponse) })
412 if err := <-done; err != nil {
413 t.Fatalf("Run: %v", err)
414 }
415 }
416
417 func TestGLMReasoningOverflowFailsBeforeToolExecution(t *testing.T) {
418 var requests atomic.Int32
419 srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
420 requestNo := requests.Add(1)
421 w.Header().Set("Content-Type", "text/event-stream")
422 if requestNo == 1 {
423 _, _ = io.WriteString(w, `data: {"choices":[{"delta":{"reasoning_content":"`+strings.Repeat("reason", 16)+`"}}]}`+"\n\n")
424 _, _ = io.WriteString(w, `data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"c1","type":"function","function":{"name":"echo","arguments":"{\"text\":\"must not run\"}"}}]},"finish_reason":"tool_calls"}]}`+"\n\n")
425 } else {
426 _, _ = io.WriteString(w, `data: {"choices":[{"delta":{"content":"unexpected continuation"},"finish_reason":"stop"}]}`+"\n\n")
427 }
428 _, _ = io.WriteString(w, "data: [DONE]\n\n")
429 }))
430 defer srv.Close()
431
432 prov, err := openai.New(provider.Config{
433 Name: "glm", BaseURL: srv.URL, Model: "glm-5.2", APIKey: "k",
434 Extra: map[string]any{"reasoning_protocol": "glm"},
435 })
436 if err != nil {
437 t.Fatalf("New provider: %v", err)
438 }
439 sink := &recordSink{}
440 a := New(prov, echoRegistry(), NewSession(""), Options{ReasoningByteLimit: 16}, sink)
441 var replayErr *ReasoningReplayError
442 if err := a.Run(withNoClosedLoop(context.Background()), "go"); !errors.As(err, &replayErr) || replayErr.Kind != ReasoningReplayOverflow {
443 t.Fatalf("Run error = %v, want ReasoningReplayOverflow", err)
444 }
445 if got := len(sink.kinds(event.ToolResult)); got != 0 {
446 t.Fatalf("tool results = %d, want 0 after incomplete GLM reasoning", got)
447 }
448 if got := requests.Load(); got != 1 {
449 t.Fatalf("HTTP requests = %d, want no continuation after incomplete GLM reasoning", got)
450 }
451 }
452
453 // TestDeepSeekAnthropicThinking400CatchAndRepair drives the full self-heal
454 // against the real Anthropic-adapter wire shape: the first request replays the
455 // stored thinking block and the server rejects it with DeepSeek's documented
456 // 400; the agent must repair the projection once (stripping all reasoning),
457 // retry, and keep the strong projection for the following run.
458 func TestDeepSeekAnthropicThinking400CatchAndRepair(t *testing.T) {
459 var mu sync.Mutex
460 var bodies [][]byte
461 srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
462 body, _ := io.ReadAll(r.Body)
463 mu.Lock()
464 bodies = append(bodies, append([]byte(nil), body...))
465 requestNo := len(bodies)
466 mu.Unlock()
467
468 if requestNo == 1 {
469 w.WriteHeader(http.StatusBadRequest)
470 _, _ = w.Write([]byte(`{"error":{"message":"The ` + "`content[].thinking`" + ` in the thinking mode must be passed back to the API","type":"invalid_request_error"}}`))
471 return
472 }
473 w.Header().Set("Content-Type", "text/event-stream")
474 _, _ = io.WriteString(w, "data: {\"type\":\"message_start\",\"message\":{\"usage\":{\"input_tokens\":20,\"output_tokens\":1}}}\n\n")
475 if requestNo == 2 {
476 _, _ = io.WriteString(w, "data: {\"type\":\"content_block_start\",\"index\":0,\"content_block\":{\"type\":\"thinking\",\"thinking\":\"\"}}\n\n")
477 _, _ = io.WriteString(w, "data: {\"type\":\"content_block_delta\",\"index\":0,\"delta\":{\"type\":\"thinking_delta\",\"thinking\":\"fresh thinking\"}}\n\n")
478 _, _ = io.WriteString(w, "data: {\"type\":\"content_block_stop\",\"index\":0}\n\n")
479 }
480 _, _ = io.WriteString(w, "data: {\"type\":\"content_block_start\",\"index\":1,\"content_block\":{\"type\":\"text\",\"text\":\"\"}}\n\n")
481 _, _ = io.WriteString(w, "data: {\"type\":\"content_block_delta\",\"index\":1,\"delta\":{\"type\":\"text_delta\",\"text\":\"repaired answer\"}}\n\n")
482 _, _ = io.WriteString(w, "data: {\"type\":\"content_block_stop\",\"index\":1}\n\n")
483 _, _ = io.WriteString(w, "data: {\"type\":\"message_delta\",\"delta\":{\"stop_reason\":\"end_turn\"},\"usage\":{\"output_tokens\":5}}\n\n")
484 _, _ = io.WriteString(w, "data: {\"type\":\"message_stop\"}\n\n")
485 }))
486 defer srv.Close()
487
488 prov, err := anthropic.New(provider.Config{
489 Name: "deepseek-anthropic", BaseURL: srv.URL, Model: "deepseek-v4-flash", APIKey: "k",
490 Extra: map[string]any{"reasoning_protocol": "deepseek", "thinking": "enabled"},
491 })
492 if err != nil {
493 t.Fatalf("New provider: %v", err)
494 }
495 session := NewSession("system")
496 session.Add(provider.Message{Role: provider.RoleUser, Content: "earlier"})
497 session.Add(provider.Message{Role: provider.RoleAssistant, Content: "old answer", ReasoningContent: "stale thinking"})
498 sink := &recordSink{}
499 a := New(prov, echoRegistry(), session, Options{}, sink)
500
501 if err := a.Run(withNoClosedLoop(context.Background()), "next"); err != nil {
502 t.Fatalf("Run: %v", err)
503 }
504
505 mu.Lock()
506 gotBodies := append([][]byte(nil), bodies...)
507 mu.Unlock()
508 if len(gotBodies) != 2 {
509 t.Fatalf("HTTP requests = %d, want rejected attempt plus one repair retry", len(gotBodies))
510 }
511 if !bytes.Contains(gotBodies[0], []byte(`"type":"thinking"`)) || !bytes.Contains(gotBodies[0], []byte("stale thinking")) {
512 t.Fatalf("first request did not replay the stored thinking block: %s", gotBodies[0])
513 }
514 if bytes.Contains(gotBodies[1], []byte(`"type":"thinking"`)) || bytes.Contains(gotBodies[1], []byte("stale thinking")) {
515 t.Fatalf("repair retry still carries thinking blocks: %s", gotBodies[1])
516 }
517 if !bytes.Contains(gotBodies[1], []byte("old answer")) {
518 t.Fatalf("repair retry lost the visible assistant text: %s", gotBodies[1])
519 }
520 // Only Messages may change between the rejected request and its repair.
521 var first, second map[string]json.RawMessage
522 if err := json.Unmarshal(gotBodies[0], &first); err != nil || json.Unmarshal(gotBodies[1], &second) != nil {
523 t.Fatalf("decode request bodies: %v", err)
524 }
525 delete(first, "messages")
526 delete(second, "messages")
527 if !reflect.DeepEqual(first, second) {
528 t.Fatalf("repair retry changed non-message fields:\nfirst=%s\nretry=%s", gotBodies[0], gotBodies[1])
529 }
530 if got := sink.recoveryCount(event.ProtocolRecoveryReasoningReplay400Detected); got != 1 {
531 t.Fatalf("reasoning_replay_400_detected audits = %d, want 1", got)
532 }
533 if got := sink.recoveryCount(event.ProtocolRecoveryReasoningReplay400Recovered); got != 1 {
534 t.Fatalf("reasoning_replay_400_recovered audits = %d, want 1", got)
535 }
536 var repairNotices int
537 for _, e := range sink.kinds(event.Notice) {
538 if e.Code == event.NoticeCodeReasoningReplayRepair && e.Level == event.LevelWarn {
539 repairNotices++
540 }
541 }
542 if repairNotices != 1 {
543 t.Fatalf("repair notices = %d, want 1", repairNotices)
544 }
545 // The adopted answer streamed through and the canonical history keeps both
546 // turns' reasoning untouched by the provider-visible repair.
547 var answer strings.Builder
548 for _, e := range sink.kinds(event.Text) {
549 answer.WriteString(e.Text)
550 }
551 if answer.String() != "repaired answer" {
552 t.Fatalf("streamed answer = %q, want the repaired response", answer.String())
553 }
554
555 // The next run keeps the repaired prefix stripped while replaying reasoning
556 // from the newly committed assistant turn normally.
557 if err := a.Run(withNoClosedLoop(context.Background()), "again"); err != nil {
558 t.Fatalf("second Run: %v", err)
559 }
560 mu.Lock()
561 third := append([]byte(nil), bodies[len(bodies)-1]...)
562 total := len(bodies)
563 mu.Unlock()
564 if total != 3 {
565 t.Fatalf("HTTP requests = %d, want one more for the follow-up run", total)
566 }
567 if bytes.Contains(third, []byte("stale thinking")) {
568 t.Fatalf("strong projection retained stale reasoning in the next run: %s", third)
569 }
570 if !bytes.Contains(third, []byte("fresh thinking")) {
571 t.Fatalf("strong projection dropped new-turn reasoning in the next run: %s", third)
572 }
573 for _, m := range session.Snapshot() {
574 if m.Role == provider.RoleAssistant && m.Content == "old answer" && m.ReasoningContent != "stale thinking" {
575 t.Fatalf("canonical history lost its reasoning: %+v", m)
576 }
577 }
578 }
579
579 lines GO