-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy pathdecoder.go
More file actions
137 lines (126 loc) · 3.83 KB
/
Copy pathdecoder.go
File metadata and controls
137 lines (126 loc) · 3.83 KB
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
131
132
133
134
135
136
137
package codexcw
import (
"encoding/json"
"fmt"
"time"
)
type wireEvent struct {
Type EventType `json:"type"`
ThreadID string `json:"thread_id"`
Usage Usage `json:"usage"`
Item json.RawMessage `json:"item"`
Error json.RawMessage `json:"error"`
Message string `json:"message"`
}
type wireItem struct {
ID string `json:"id"`
Type ItemType `json:"type"`
Status string `json:"status"`
Text string `json:"text"`
Message string `json:"message"`
Command string `json:"command"`
AggregatedOutput string `json:"aggregated_output"`
ExitCode *int `json:"exit_code"`
Changes []FileChange `json:"changes"`
Tool string `json:"tool"`
SenderThreadID string `json:"sender_thread_id"`
ReceiverThreadIDs []string `json:"receiver_thread_ids"`
}
func decodeEvent(line []byte, runID string, threadID string, now time.Time) (Event, error) {
raw := append(json.RawMessage(nil), line...)
var wire wireEvent
if err := json.Unmarshal(line, &wire); err != nil {
return Event{}, err
}
if wire.Type == "" {
return Event{}, fmt.Errorf("missing event type")
}
event := Event{
Type: wire.Type,
RunID: runID,
ThreadID: threadID,
ReceivedAt: now,
Raw: raw,
}
switch wire.Type {
case EventThreadStarted:
event.ThreadID = wire.ThreadID
event.ThreadStarted = &ThreadStartedEvent{ThreadID: wire.ThreadID}
case EventTurnStarted:
event.TurnStarted = &TurnStartedEvent{}
case EventTurnCompleted:
event.TurnCompleted = &TurnCompletedEvent{Usage: normalizeCodexUsage(wire.Usage)}
case EventTurnFailed:
event.TurnFailed = &TurnFailedEvent{
Error: decodeEventError(wire.Error),
Usage: normalizeCodexUsage(wire.Usage),
}
case EventItemStarted:
item, err := decodeItem(wire.Item)
if err != nil {
return Event{}, err
}
event.ItemStarted = &ItemEvent{Item: item}
case EventItemCompleted:
item, err := decodeItem(wire.Item)
if err != nil {
return Event{}, err
}
event.ItemCompleted = &ItemEvent{Item: item}
case EventError:
event.Error = &CodexErrorEvent{Message: wire.Message, Raw: append(json.RawMessage(nil), wire.Error...)}
if event.Error.Message == "" && len(wire.Error) > 0 {
event.Error.Message = string(wire.Error)
}
}
return event, nil
}
func decodeItem(raw json.RawMessage) (Item, error) {
if len(raw) == 0 {
return Item{}, fmt.Errorf("missing item payload")
}
var wire wireItem
if err := json.Unmarshal(raw, &wire); err != nil {
return Item{}, err
}
return Item{
ID: wire.ID,
Type: wire.Type,
Status: wire.Status,
Raw: append(json.RawMessage(nil), raw...),
Text: wire.Text,
Message: wire.Message,
Command: wire.Command,
AggregatedOutput: wire.AggregatedOutput,
ExitCode: wire.ExitCode,
Changes: wire.Changes,
Tool: wire.Tool,
SenderThreadID: wire.SenderThreadID,
ReceiverThreadIDs: wire.ReceiverThreadIDs,
}, nil
}
// normalizeCodexUsage derives the total when Codex omits total_tokens.
// cached_input_tokens is a subset of input_tokens on the codex wire, so the
// derived total is input plus output; an explicit total is preserved.
func normalizeCodexUsage(usage Usage) Usage {
if usage.TotalTokens == 0 {
usage.TotalTokens = usage.InputTokens + usage.OutputTokens
}
return usage
}
func decodeEventError(raw json.RawMessage) ErrorPayload {
err := ErrorPayload{Raw: append(json.RawMessage(nil), raw...)}
if len(raw) == 0 {
return err
}
var wire struct {
Message string `json:"message"`
}
if json.Unmarshal(raw, &wire) == nil {
err.Message = wire.Message
}
if err.Message == "" {
err.Message = string(raw)
}
return err
}