返回 DeepSeek-Reasonix
stream_fragment_recovery_test.go
根目录 / internal / agent / stream_fragment_recovery_test.go
1 package agent
2
3 import (
4 "bytes"
5 "context"
6 "io"
7 "net/http"
8 "net/http/httptest"
9 "sync"
10 "testing"
11
12 "reasonix/internal/event"
13 "reasonix/internal/provider"
14 "reasonix/internal/provider/anthropic"
15 "reasonix/internal/provider/openai"
16 "reasonix/internal/provider/responses"
17 )
18
19 func TestTruncatedJSONStreamFailsWithoutRetryOrExecutingPartialTools(t *testing.T) {
20 for _, protocol := range []string{"chat", "anthropic"} {
21 for _, cut := range []bool{true, false} {
22 name := protocol + "/malformed_line"
23 if cut {
24 name = protocol + "/unterminated_fragment"
25 }
26 t.Run(name, func(t *testing.T) {
27 var mu sync.Mutex
28 var bodies [][]byte
29 srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
30 body, _ := io.ReadAll(r.Body)
31 mu.Lock()
32 bodies = append(bodies, body)
33 n := len(bodies)
34 mu.Unlock()
35 w.Header().Set("Content-Type", "text/event-stream")
36 if n == 1 {
37 // Complete tool arguments received before EOF still must not execute.
38 first := `data: {"choices":[{"delta":{"reasoning_content":"partial","tool_calls":[{"index":0,"id":"uncommitted","type":"function","function":{"name":"echo","arguments":"{\"text\":\"unsafe\"}"}}]}}]}` + "\n\n"
39 if protocol == "anthropic" {
40 first = `data: {"type":"content_block_start","index":0,"content_block":{"type":"tool_use","id":"uncommitted","name":"echo","input":{"text":"unsafe"}}}` + "\n\n"
41 }
42 _, _ = io.WriteString(w, first+`data: {"delta":nu`)
43 if !cut {
44 _, _ = io.WriteString(w, "\n\n")
45 }
46 return
47 }
48 if n > 2 {
49 w.WriteHeader(http.StatusUnauthorized)
50 return
51 }
52 reply := `data: {"choices":[{"delta":{"content":"recovered"},"finish_reason":"stop"}],"usage":{"prompt_tokens":12,"completion_tokens":2,"total_tokens":14}}` + "\n\ndata: [DONE]\n\n"
53 if protocol == "anthropic" {
54 reply = finalAnswerSSE
55 }
56 _, _ = io.WriteString(w, reply)
57 }))
58 defer srv.Close()
59 cfg := provider.Config{Name: "fragment", BaseURL: srv.URL, APIKey: "fixture", Model: "deepseek-v4-flash", Extra: map[string]any{"reasoning_protocol": "deepseek", "thinking": "enabled"}}
60 var p provider.Provider
61 var err error
62 if protocol == "chat" {
63 p, err = openai.New(cfg)
64 } else {
65 p, err = anthropic.New(cfg)
66 }
67 if err != nil {
68 t.Fatal(err)
69 }
70 sink := &recordSink{}
71 a := New(p, echoRegistry(), NewSession(""), Options{MissingReasoningWarnStateDir: t.TempDir()}, sink)
72 err = a.Run(withNoClosedLoop(context.Background()), "go")
73 if err == nil {
74 t.Fatalf("cut=%v error=%v", cut, err)
75 }
76 if len(sink.kinds(event.ToolResult)) != 0 {
77 t.Fatal("uncommitted tool executed")
78 }
79 mu.Lock()
80 defer mu.Unlock()
81 if len(bodies) != 1 {
82 t.Fatalf("failed stream retried: %d requests", len(bodies))
83 }
84 usages := sink.kinds(event.Usage)
85 if len(usages) != 1 || usages[0].Usage == nil || !usages[0].Usage.Unknown || usages[0].Usage.RequestCount != 1 {
86 t.Fatalf("usage=%+v", usages)
87 }
88 })
89 }
90 }
91 }
92
93 func TestResponsesPassbackRejectionRepairsHistoryWithoutRepeatingTool(t *testing.T) {
94 var mu sync.Mutex
95 var bodies [][]byte
96 srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
97 body, _ := io.ReadAll(r.Body)
98 mu.Lock()
99 bodies = append(bodies, body)
100 n := len(bodies)
101 mu.Unlock()
102 if n == 2 {
103 w.WriteHeader(http.StatusBadRequest)
104 _, _ = io.WriteString(w, "{\"error\":{\"message\":\"The `reasoning_text` in the thinking mode must be passed back to the API.\",\"type\":\"invalid_request_error\"}}")
105 return
106 }
107 if n > 4 {
108 w.WriteHeader(http.StatusUnauthorized)
109 return
110 }
111 w.Header().Set("Content-Type", "text/event-stream")
112 if n == 1 {
113 _, _ = io.WriteString(w, responsesToolWithoutReasoningSSE)
114 } else {
115 _, _ = io.WriteString(w, responsesFinalAnswerSSE)
116 }
117 }))
118 defer srv.Close()
119 p := responses.New(responses.Config{Name: "responses", BaseURL: "https://api.deepseek.com", RequestURL: srv.URL, Model: "deepseek-v4-flash", APIKey: "fixture", Effort: "high", Mode: "stateless"})
120 sink := &recordSink{}
121 a := New(p, echoRegistry(), NewSession(""), Options{MissingReasoningWarnStateDir: t.TempDir()}, sink)
122 if err := a.Run(withNoClosedLoop(context.Background()), "go"); err != nil {
123 t.Fatal(err)
124 }
125 if err := a.Run(withNoClosedLoop(context.Background()), "A new question: acknowledge the previous result without calling tools."); err != nil {
126 t.Fatal(err)
127 }
128 mu.Lock()
129 defer mu.Unlock()
130 if len(bodies) != 4 || len(sink.kinds(event.ToolResult)) != 1 {
131 t.Fatalf("requests=%d executions=%d", len(bodies), len(sink.kinds(event.ToolResult)))
132 }
133 if bytes.Contains(bodies[3], []byte(`"type":"function_call"`)) || !bytes.Contains(bodies[3], []byte("echoed: hi")) {
134 t.Fatal("later turn restored invalid tool history or lost completed results")
135 }
136 if !bytes.Contains(bodies[2], []byte("echoed: hi")) {
137 t.Fatal("repair dropped the actual completed tool output")
138 }
139 if !bytes.Contains(bodies[2], []byte("completed_tools")) || bytes.Contains(bodies[2], []byte(`"type":"function_call"`)) {
140 t.Fatal("repair lost completed facts or retained invalid call")
141 }
142 }
143
143 lines GO