返回 DeepSeek-Reasonix
stream_eof_test.go
根目录 / internal / provider / anthropic / stream_eof_test.go
1 package anthropic
2
3 import (
4 "context"
5 "errors"
6 "io"
7 "net/http"
8 "strings"
9 "testing"
10 "time"
11
12 "reasonix/internal/provider"
13 )
14
15 // Compatible Anthropic gateways often omit message_stop after message_delta
16 // the same way OpenAI-compatible gateways omit [DONE] after finish_reason.
17 // A stop_reason is the provider's commit signal; EOF after that must finalize.
18
19 func TestReadStreamAcceptsStopReasonWithoutMessageStop(t *testing.T) {
20 sse := `event: message_start
21 data: {"type":"message_start","message":{"id":"msg_1","usage":{"input_tokens":5}}}
22
23 event: content_block_delta
24 data: {"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"hello"}}
25
26 event: message_delta
27 data: {"type":"message_delta","delta":{"stop_reason":"end_turn"},"usage":{"output_tokens":1}}
28 `
29 c := &client{name: "cc"}
30 resp := &http.Response{Body: io.NopCloser(strings.NewReader(sse))}
31 ch := make(chan provider.Chunk)
32 go c.readStream(context.Background(), resp, ch)
33
34 var text strings.Builder
35 var sawDone, sawErr bool
36 var usage *provider.Usage
37 for ck := range ch {
38 switch ck.Type {
39 case provider.ChunkText:
40 text.WriteString(ck.Text)
41 case provider.ChunkUsage:
42 usage = ck.Usage
43 case provider.ChunkDone:
44 sawDone = true
45 case provider.ChunkError:
46 sawErr = true
47 t.Fatalf("stop_reason without message_stop should complete cleanly: %v", ck.Err)
48 }
49 }
50 if text.String() != "hello" || !sawDone || sawErr {
51 t.Fatalf("text=%q done=%v err=%v", text.String(), sawDone, sawErr)
52 }
53 if usage == nil || usage.FinishReason != "stop" {
54 t.Fatalf("usage = %+v, want finish_reason=stop", usage)
55 }
56 }
57
58 func TestReadStreamAcceptsToolUseStopReasonWithoutMessageStop(t *testing.T) {
59 sse := `event: message_start
60 data: {"type":"message_start","message":{"id":"msg_1","usage":{"input_tokens":10}}}
61
62 event: content_block_start
63 data: {"type":"content_block_start","index":0,"content_block":{"type":"tool_use","id":"toolu_1","name":"bash"}}
64
65 event: content_block_delta
66 data: {"type":"content_block_delta","index":0,"delta":{"type":"input_json_delta","partial_json":"{\"cmd\":\"ls\"}"}}
67
68 event: content_block_stop
69 data: {"type":"content_block_stop","index":0}
70
71 event: message_delta
72 data: {"type":"message_delta","delta":{"stop_reason":"tool_use"},"usage":{"output_tokens":8}}
73 `
74 c := &client{name: "cc"}
75 resp := &http.Response{Body: io.NopCloser(strings.NewReader(sse))}
76 ch := make(chan provider.Chunk)
77 go c.readStream(context.Background(), resp, ch)
78
79 var sawDone, sawErr bool
80 var call *provider.ToolCall
81 var usage *provider.Usage
82 for ck := range ch {
83 switch ck.Type {
84 case provider.ChunkToolCall:
85 call = ck.ToolCall
86 case provider.ChunkUsage:
87 usage = ck.Usage
88 case provider.ChunkDone:
89 sawDone = true
90 case provider.ChunkError:
91 sawErr = true
92 t.Fatalf("tool_use stop_reason without message_stop should complete: %v", ck.Err)
93 }
94 }
95 if sawErr || !sawDone {
96 t.Fatalf("done=%v err=%v", sawDone, sawErr)
97 }
98 if call == nil || call.ID != "toolu_1" || call.Name != "bash" || call.Arguments != `{"cmd":"ls"}` {
99 t.Fatalf("tool call = %+v", call)
100 }
101 if usage == nil || usage.FinishReason != "tool_calls" {
102 t.Fatalf("usage = %+v, want finish_reason=tool_calls", usage)
103 }
104 }
105
106 func TestStreamScanEndErrorClassifiesTerminal(t *testing.T) {
107 if err := streamScanEndError("cc", time.Second, false, nil, "end_turn"); err != nil {
108 t.Fatalf("stop_reason must commit: %v", err)
109 }
110 err := streamScanEndError("cc", time.Second, false, nil, "")
111 if !provider.IsStreamInterrupted(err) {
112 t.Fatalf("empty stop_reason = %v, want interrupt", err)
113 }
114 if !strings.Contains(err.Error(), "stream ended before message_stop") {
115 t.Fatalf("interrupt = %v", err)
116 }
117 if err := streamScanEndError("cc", time.Second, true, nil, ""); !provider.IsStreamInterrupted(err) ||
118 provider.StreamInterruptReason(err) != provider.StreamInterruptIdleTimeout {
119 t.Fatalf("stall = %v", err)
120 }
121 }
122
123 func TestReadStreamStillInterruptsWithoutStopReasonOrMessageStop(t *testing.T) {
124 sse := `event: message_start
125 data: {"type":"message_start","message":{"id":"msg_1","usage":{"input_tokens":5}}}
126
127 event: content_block_delta
128 data: {"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"partial"}}
129 `
130 c := &client{name: "cc"}
131 resp := &http.Response{Body: io.NopCloser(strings.NewReader(sse))}
132 ch := make(chan provider.Chunk)
133 go c.readStream(context.Background(), resp, ch)
134
135 var gotInterrupted, sawDone bool
136 for ck := range ch {
137 switch ck.Type {
138 case provider.ChunkDone:
139 sawDone = true
140 case provider.ChunkError:
141 var interrupted *provider.StreamInterruptedError
142 gotInterrupted = errors.As(ck.Err, &interrupted)
143 }
144 }
145 if sawDone {
146 t.Fatal("must not emit ChunkDone without message_stop or stop_reason")
147 }
148 if !gotInterrupted {
149 t.Fatal("EOF without stop_reason must stay StreamInterruptedError")
150 }
151 }
152
152 lines GO