1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
|
package usage
import (
"testing"
"time"
"aigw/internal/domain"
)
func TestObserverReadsOpenAIJSONUsage(t *testing.T) {
observer := NewObserver(domain.ProtocolOpenAI, false)
_, _ = observer.Write([]byte(`{"choices":[],"usage":{"prompt_tokens":11,"prompt_tokens_details":{"cached_tokens":3,"cache_write_tokens":2},"completion_tokens":7,"total_tokens":18}}`))
got := observer.Usage()
if !observer.Reported() {
t.Fatal("expected usage to be marked as reported")
}
if got.InputTokens != 6 || got.CacheReadInputTokens != 3 || got.CacheCreationInputTokens != 2 || got.OutputTokens != 7 || got.TotalTokens != 18 {
t.Fatalf("unexpected usage: %+v", got)
}
}
func TestObserverCombinesAnthropicSSEUsage(t *testing.T) {
observer := NewObserver(domain.ProtocolAnthropic, true)
_, _ = observer.Write([]byte("event: message_start\ndata: {\"type\":\"message_start\",\"message\":{\"usage\":{\"input_tokens\":12,\"output_tokens\":1,\"cache_read_input_tokens\":5,\"cache_creation_input_tokens\":2}}}\n\n"))
_, _ = observer.Write([]byte("event: message_delta\ndata: {\"type\":\"message_delta\",\"usage\":{\"output_tokens\":8}}\n\n"))
got := observer.Usage()
if !observer.Reported() {
t.Fatal("expected streaming usage to be marked as reported")
}
if got.InputTokens != 12 || got.OutputTokens != 8 || got.TotalTokens != 20 || got.CacheReadInputTokens != 5 || got.CacheCreationInputTokens != 2 {
t.Fatalf("unexpected usage: %+v", got)
}
}
func TestObserverDistinguishesMissingUsageFromReportedZero(t *testing.T) {
missing := NewObserver(domain.ProtocolOpenAI, false)
_, _ = missing.Write([]byte(`{"choices":[]}`))
if missing.Reported() {
t.Fatal("response without usage must not be reported")
}
reported := NewObserver(domain.ProtocolOpenAI, false)
_, _ = reported.Write([]byte(`{"choices":[],"usage":{"prompt_tokens":0,"completion_tokens":0,"total_tokens":0}}`))
if !reported.Reported() {
t.Fatal("explicit zero usage must be distinguished from a missing usage object")
}
}
func TestObserverReadsResponsesUsage(t *testing.T) {
nonStream := NewObserver(domain.ProtocolOpenAIResponses, false)
_, _ = nonStream.Write([]byte(`{"object":"response","usage":{"input_tokens":11,"input_tokens_details":{"cached_tokens":3,"cache_write_tokens":2},"output_tokens":7,"total_tokens":18}}`))
if got := nonStream.Usage(); got.InputTokens != 6 || got.OutputTokens != 7 || got.TotalTokens != 18 || got.CacheReadInputTokens != 3 || got.CacheCreationInputTokens != 2 || !nonStream.Reported() {
t.Fatalf("unexpected non-stream Responses usage: %+v reported=%v", got, nonStream.Reported())
}
stream := NewObserver(domain.ProtocolOpenAIResponses, true)
_, _ = stream.Write([]byte("event: response.completed\ndata: {\"type\":\"response.completed\",\"response\":{\"usage\":{\"input_tokens\":13,\"output_tokens\":5,\"total_tokens\":18}}}\n\n"))
if got := stream.Usage(); got.InputTokens != 13 || got.OutputTokens != 5 || got.TotalTokens != 18 || !stream.Reported() {
t.Fatalf("unexpected streaming Responses usage: %+v reported=%v", got, stream.Reported())
}
}
func TestObserverReadsEmbeddingsUsage(t *testing.T) {
observer := NewObserver(domain.ProtocolOpenAIEmbeddings, false)
_, _ = observer.Write([]byte(`{"object":"list","data":[],"usage":{"prompt_tokens":17,"total_tokens":17}}`))
got := observer.Usage()
if !observer.Reported() || got.InputTokens != 17 || got.OutputTokens != 0 || got.TotalTokens != 17 {
t.Fatalf("unexpected Embeddings usage: %+v reported=%v", got, observer.Reported())
}
}
func TestObserverMarksFirstVisibleStreamingOutput(t *testing.T) {
base := time.Date(2026, 8, 6, 1, 2, 3, 0, time.UTC)
tests := []struct {
name string
protocol domain.Protocol
metadata string
output string
}{
{
name: "openai chat", protocol: domain.ProtocolOpenAI,
metadata: "data: {\"choices\":[{\"delta\":{\"role\":\"assistant\",\"content\":null}}]}\n\n",
output: "data: {\"choices\":[{\"delta\":{\"content\":\"Hello\"}}]}\n\n",
},
{
name: "responses", protocol: domain.ProtocolOpenAIResponses,
metadata: "data: {\"type\":\"response.created\"}\n\n",
output: "data: {\"type\":\"response.output_text.delta\",\"delta\":\"Hello\"}\n\n",
},
{
name: "anthropic", protocol: domain.ProtocolAnthropic,
metadata: "data: {\"type\":\"content_block_start\",\"content_block\":{\"type\":\"text\",\"text\":\"\"}}\n\n",
output: "data: {\"type\":\"content_block_delta\",\"delta\":{\"type\":\"text_delta\",\"text\":\"Hello\"}}\n\n",
},
}
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
observer := NewObserver(test.protocol, true)
current := base
observer.now = func() time.Time { return current }
_, _ = observer.Write([]byte(test.metadata))
if got := observer.FirstOutputAt(); !got.IsZero() {
t.Fatalf("metadata marked first output at %v", got)
}
current = base.Add(275 * time.Millisecond)
_, _ = observer.Write([]byte(test.output))
if got := observer.FirstOutputAt(); !got.Equal(current) {
t.Fatalf("first output = %v, want %v", got, current)
}
current = base.Add(time.Second)
_, _ = observer.Write([]byte(test.output))
if got := observer.FirstOutputAt(); !got.Equal(base.Add(275 * time.Millisecond)) {
t.Fatalf("first output changed to %v", got)
}
})
}
}
func TestObserverMarksFirstNonStreamingBodyWrite(t *testing.T) {
base := time.Date(2026, 8, 6, 1, 2, 3, 0, time.UTC)
observer := NewObserver(domain.ProtocolOpenAI, false)
observer.now = func() time.Time { return base }
_, _ = observer.Write(nil)
if !observer.FirstOutputAt().IsZero() {
t.Fatal("empty write must not mark first output")
}
_, _ = observer.Write([]byte(`{"choices":[]}`))
if got := observer.FirstOutputAt(); !got.Equal(base) {
t.Fatalf("first output = %v, want %v", got, base)
}
}
|