返回 DeepSeek-Reasonix
run_output.go
根目录 / internal / cli / run_output.go
1 package cli
2
3 import (
4 "encoding/json"
5 "errors"
6 "fmt"
7 "io"
8 "strings"
9 "sync"
10 "time"
11
12 "reasonix/internal/billing"
13 "reasonix/internal/control"
14 "reasonix/internal/event"
15 "reasonix/internal/eventwire"
16 )
17
18 type runOutputFormat string
19
20 const (
21 runOutputText runOutputFormat = "text"
22 runOutputJSON runOutputFormat = "json"
23 runOutputStreamJSON runOutputFormat = "stream-json"
24 runOutputEventsJSONL runOutputFormat = "events-jsonl"
25 )
26
27 func parseRunOutputFormat(value string) (runOutputFormat, error) {
28 switch runOutputFormat(strings.ToLower(strings.TrimSpace(value))) {
29 case runOutputText:
30 return runOutputText, nil
31 case runOutputJSON:
32 return runOutputJSON, nil
33 case runOutputStreamJSON:
34 return runOutputStreamJSON, nil
35 default:
36 return "", fmt.Errorf("unknown output format %q (want text, json, or stream-json)", value)
37 }
38 }
39
40 // runOutputSessionID preserves the established json/stream-json contract while
41 // keeping the redacted events-jsonl surface independent from transcript names.
42 func runOutputSessionID(format runOutputFormat, rawSessionID string, identityKey []byte) string {
43 if format == runOutputEventsJSONL {
44 return machineSessionIDWithKey(rawSessionID, identityKey)
45 }
46 return rawSessionID
47 }
48
49 type runResultUsage struct {
50 InputTokens int `json:"input_tokens"`
51 OutputTokens int `json:"output_tokens"`
52 CacheReadInputTokens int `json:"cache_read_input_tokens"`
53 CacheCreationInputTokens int `json:"cache_creation_input_tokens"`
54 Estimated bool `json:"estimated,omitempty"`
55 }
56
57 type runResult struct {
58 Type string `json:"type"`
59 Subtype string `json:"subtype"`
60 IsError bool `json:"is_error"`
61 DurationMS int64 `json:"duration_ms"`
62 NumTurns int `json:"num_turns"`
63 Result string `json:"result"`
64 // ResultFromReasoning marks a Result taken from the turn's reasoning
65 // because the model emitted no visible text. Omitted when false.
66 ResultFromReasoning bool `json:"result_from_reasoning,omitempty"`
67 SessionID string `json:"session_id,omitempty"`
68 TotalCost float64 `json:"total_cost,omitempty"`
69 Currency string `json:"currency,omitempty"`
70 // TotalCostUSD is the released compatibility alias. It mirrors TotalCost;
71 // new consumers must pair TotalCost with Currency instead of assuming USD.
72 TotalCostUSD float64 `json:"total_cost_usd,omitempty"`
73 // CostComplete is false when mixed originals lack a shared display valuation.
74 CostComplete bool `json:"cost_complete"`
75 DisplayComplete bool `json:"display_complete"`
76 DisplayStatus string `json:"display_status,omitempty"`
77 AggregateMode string `json:"aggregate_mode,omitempty"`
78 // OriginalCosts lists per-ISO original totals (never cross-added).
79 OriginalCosts map[string]float64 `json:"original_costs,omitempty"`
80 OriginalTotals []billing.Money `json:"original_totals,omitempty"`
81 CostQuote *billing.CostQuote `json:"cost_quote,omitempty"`
82 Usage runResultUsage `json:"usage"`
83 ErrorCode string `json:"error_code,omitempty"`
84 Authentication string `json:"authentication_status,omitempty"`
85 Recovery []string `json:"recovery_actions,omitempty"`
86 }
87
88 type machineEventUsage struct {
89 InputTokens int `json:"input_tokens"`
90 OutputTokens int `json:"output_tokens"`
91 CacheHitTokens int `json:"cache_hit_tokens"`
92 CacheMissTokens int `json:"cache_miss_tokens"`
93 Estimated bool `json:"estimated,omitempty"`
94 }
95
96 // machineEventRecord is deliberately content-free. The existing stream-json
97 // format is a rich UI transport and includes prompts, tool arguments, results,
98 // and reasoning; this contract is for automation that must not receive them.
99 type machineEventRecord struct {
100 SchemaVersion int `json:"schema_version"`
101 Sequence uint64 `json:"sequence"`
102 Kind string `json:"kind"`
103 Code string `json:"code,omitempty"`
104 Level string `json:"level,omitempty"`
105 ToolID string `json:"tool_id,omitempty"`
106 ToolName string `json:"tool_name,omitempty"`
107 ToolReadOnly bool `json:"tool_read_only,omitempty"`
108 ToolError bool `json:"tool_error,omitempty"`
109 ToolTruncated bool `json:"tool_truncated,omitempty"`
110 ToolDurationMS int64 `json:"tool_duration_ms,omitempty"`
111 Usage *machineEventUsage `json:"usage,omitempty"`
112 ApprovalID string `json:"approval_id,omitempty"`
113 ApprovalKind string `json:"approval_kind,omitempty"`
114 AskID string `json:"ask_id,omitempty"`
115 Outcome string `json:"outcome,omitempty"`
116 Cancelled bool `json:"cancelled,omitempty"`
117 Error bool `json:"error,omitempty"`
118 Recovery *event.RecoveryStatus `json:"recovery,omitempty"`
119 RetryAttempt int `json:"retry_attempt,omitempty"`
120 RetryMax int `json:"retry_max,omitempty"`
121 CompactionType string `json:"compaction_type,omitempty"`
122 CompactionMsgs int `json:"compaction_messages,omitempty"`
123 GuardianResult string `json:"guardian_result,omitempty"`
124 GuardianRisk string `json:"guardian_risk,omitempty"`
125 }
126
127 type machineRunDone struct {
128 SchemaVersion int `json:"schema_version"`
129 Sequence uint64 `json:"sequence"`
130 Kind string `json:"kind"`
131 SessionID string `json:"session_id,omitempty"`
132 OK bool `json:"ok"`
133 DurationMS int64 `json:"duration_ms"`
134 NumTurns int `json:"num_turns"`
135 Usage machineEventUsage `json:"usage"`
136 ErrorCode string `json:"error_code,omitempty"`
137 Authentication string `json:"authentication_status,omitempty"`
138 Recovery []string `json:"recovery_actions,omitempty"`
139 }
140
141 func runAuthenticationMetadata(err error) (code, status string, actions []string) {
142 var authErr *control.AuthenticationError
143 if !errors.As(err, &authErr) || authErr == nil {
144 return "", "", nil
145 }
146 state := authErr.State
147 code = state.Code
148 if code == "" {
149 code = string(state.Status)
150 }
151 status = string(state.Status)
152 switch state.Status {
153 case control.AuthenticationMissingCredential:
154 actions = []string{"configure_credentials", "select_model", "diagnose_credentials"}
155 case control.AuthenticationRejected:
156 actions = []string{"update_credentials", "select_model", "test_connection", "retry_authentication"}
157 case control.AuthenticationCredentialStoreUnavailable:
158 actions = []string{"diagnose_credentials", "select_model"}
159 }
160 return code, status, actions
161 }
162
163 // finalAnswer is the turn's answer together with where it came from. A
164 // thinking model may answer entirely in the reasoning channel, finishing with
165 // non-empty reasoning and an empty visible message - an accepted shape. The
166 // transcript renderer prints a "thinking" marker for it, and the reasoning
167 // itself under --show-thinking; these sinks print nothing at all and still
168 // report success, so the answer is lost and nothing says so.
169 //
170 // The two fields are one value because they are only ever correct together: a
171 // fromReasoning left over from an earlier turn would misreport the provenance
172 // of a later visible answer.
173 type finalAnswer struct {
174 text string
175 fromReasoning bool
176 }
177
178 type runOutputSink struct {
179 mu sync.Mutex
180 format runOutputFormat
181 out io.Writer
182 encoder *json.Encoder
183 final finalAnswer
184 usage runResultUsage
185 cost float64
186 currency string
187 costComplete bool
188 displayComplete bool
189 displayStatus string
190 aggregateMode string
191 originalTotals []billing.Money
192 sawQuote bool
193 originalCosts map[string]float64
194 quoteLedger *billing.Ledger
195 turns int
196 sequence uint64
197 machineToolIDs map[string]string
198 machineToolNames map[string]string
199 nextMachineToolID uint64
200 nextMachineToolName uint64
201 err error
202 }
203
204 func newRunOutputSink(out io.Writer, format runOutputFormat) *runOutputSink {
205 return &runOutputSink{
206 format: format,
207 out: out,
208 encoder: json.NewEncoder(out),
209 machineToolIDs: make(map[string]string),
210 machineToolNames: make(map[string]string),
211 }
212 }
213
214 func (s *runOutputSink) Emit(e event.Event) {
215 s.mu.Lock()
216 defer s.mu.Unlock()
217 if e.Kind == event.Message {
218 s.final = finalAnswer{text: e.Text}
219 if e.Text == "" && e.Reasoning != "" {
220 s.final = finalAnswer{text: e.Reasoning, fromReasoning: true}
221 }
222 }
223 if e.Kind == event.Usage && e.Usage != nil {
224 s.usage.InputTokens += e.Usage.PromptTokens
225 s.usage.OutputTokens += e.Usage.CompletionTokens
226 s.usage.CacheReadInputTokens += e.Usage.CacheHitTokens
227 s.usage.CacheCreationInputTokens += e.Usage.CacheMissTokens
228 s.usage.Estimated = s.usage.Estimated || e.Usage.Estimated
229 q := e.CostQuote
230 if q == nil && e.Pricing != nil {
231 q = event.EnsureCostQuote(e, nil)
232 }
233 if q != nil {
234 s.sawQuote = true
235 if !q.CostComplete {
236 s.costComplete = false
237 }
238 // First complete quote establishes complete=true.
239 if q.Complete && s.quoteLedger == nil {
240 s.costComplete = true
241 }
242 if s.originalCosts == nil {
243 s.originalCosts = map[string]float64{}
244 }
245 if cur := billing.NormalizeCurrency(q.Original.Currency); cur != "" {
246 s.originalCosts[cur] += q.Original.Float64()
247 }
248 if q.Selected != nil && (s.currency == "" || s.currency == q.LegacyCurrencyCode()) {
249 s.cost += q.Selected.Float64()
250 s.currency = q.LegacyCurrencyCode()
251 } else if q.Selected != nil {
252 s.currency = ""
253 s.cost = 0
254 }
255 if s.quoteLedger == nil {
256 s.quoteLedger = billing.NewLedger()
257 }
258 s.quoteLedger.Add(*q, billing.UsageTokens{
259 PromptTokens: e.Usage.PromptTokens,
260 CompletionTokens: e.Usage.CompletionTokens,
261 CacheHitTokens: e.Usage.CacheHitTokens,
262 CacheMissTokens: e.Usage.CacheMissTokens,
263 CacheWriteTokens: e.Usage.CacheWriteTokens,
264 CacheWriteBilledTokens: e.Usage.CacheWriteBilledTokens,
265 Estimated: e.Usage.Estimated,
266 }, time.Now().UTC())
267 }
268 }
269 if e.Kind == event.TurnDone {
270 s.turns++
271 }
272 if s.format == runOutputStreamJSON && s.err == nil {
273 s.err = s.encoder.Encode(eventwire.ToWire(e))
274 } else if s.format == runOutputEventsJSONL && s.err == nil {
275 s.sequence++
276 s.err = s.encoder.Encode(s.machineEventRecordFor(e, s.sequence))
277 }
278 }
279
280 func (s *runOutputSink) Finalize(sessionID string, started time.Time, runErr error) error {
281 s.mu.Lock()
282 defer s.mu.Unlock()
283 if s.err != nil {
284 return s.err
285 }
286 // Mixed original currencies no longer error: totals use shared display
287 // valuations when complete, otherwise cost_complete=false with original_costs.
288 if s.format == runOutputText {
289 if s.final.text != "" {
290 _, s.err = fmt.Fprintln(s.out, s.final.text)
291 }
292 return s.err
293 }
294 completion := classifyRunCompletion(runErr)
295 errorCode, authentication, recovery := runAuthenticationMetadata(runErr)
296 if s.format == runOutputEventsJSONL {
297 s.sequence++
298 turns := s.turns
299 if turns == 0 && !completion.isError {
300 turns = 1
301 }
302 return s.encoder.Encode(machineRunDone{
303 SchemaVersion: machineSchemaVersion,
304 Sequence: s.sequence,
305 Kind: "run_done",
306 SessionID: sessionID,
307 OK: !completion.isError,
308 DurationMS: time.Since(started).Milliseconds(),
309 NumTurns: turns,
310 Usage: machineEventUsage{InputTokens: s.usage.InputTokens, OutputTokens: s.usage.OutputTokens, CacheHitTokens: s.usage.CacheReadInputTokens, CacheMissTokens: s.usage.CacheCreationInputTokens},
311 ErrorCode: errorCode,
312 Authentication: authentication,
313 Recovery: recovery,
314 })
315 }
316 answer := s.final
317 if runErr != nil && answer.text == "" {
318 answer = finalAnswer{text: runErr.Error()}
319 }
320 turns := s.turns
321 if turns == 0 && !completion.isError {
322 turns = 1
323 }
324 var aggQuote *billing.CostQuote
325 if s.quoteLedger != nil && len(s.quoteLedger.Entries) > 0 {
326 agg := s.quoteLedger.Total("")
327 aggQuote = &agg
328 if agg.Selected != nil {
329 s.cost = agg.Selected.Float64()
330 s.currency = agg.LegacyCurrencyCode()
331 }
332 if agg.Selected == nil {
333 s.cost = 0
334 s.currency = ""
335 }
336 s.costComplete = agg.CostComplete
337 s.displayComplete = agg.DisplayComplete
338 s.displayStatus = agg.DisplayStatus
339 s.aggregateMode = agg.AggregateMode
340 if agg.OriginalTotals != nil {
341 s.originalTotals = append([]billing.Money(nil), agg.OriginalTotals...)
342 }
343 }
344 return s.encoder.Encode(runResult{
345 Type: "result",
346 Subtype: completion.subtype,
347 IsError: completion.isError,
348 DurationMS: time.Since(started).Milliseconds(),
349 NumTurns: turns,
350 Result: answer.text,
351 ResultFromReasoning: answer.fromReasoning,
352 SessionID: sessionID,
353 TotalCost: s.cost,
354 Currency: s.currency,
355 TotalCostUSD: s.cost,
356 CostComplete: s.costComplete || (!s.sawQuote && s.currency != ""),
357 DisplayComplete: s.displayComplete,
358 DisplayStatus: s.displayStatus,
359 AggregateMode: s.aggregateMode,
360 OriginalCosts: s.originalCosts,
361 OriginalTotals: s.originalTotals,
362 CostQuote: aggQuote,
363 Usage: s.usage,
364 ErrorCode: errorCode,
365 Authentication: authentication,
366 Recovery: recovery,
367 })
368 }
369
370 func (s *runOutputSink) machineEventRecordFor(e event.Event, sequence uint64) machineEventRecord {
371 record := machineEventRecord{SchemaVersion: machineSchemaVersion, Sequence: sequence, Kind: machineEventKind(e.Kind)}
372 switch e.Kind {
373 case event.Notice:
374 record.Code = e.Code
375 if e.Level == event.LevelWarn {
376 record.Level = "warn"
377 } else {
378 record.Level = "info"
379 }
380 case event.ToolDispatch, event.ToolResult, event.ToolProgress:
381 // Tool-call IDs and names originate in the provider stream. Treat both as
382 // untrusted content: an OpenAI-compatible endpoint may put arbitrary prompt
383 // or argument text in either field. Per-run opaque aliases preserve event
384 // correlation without exposing the provider-controlled values.
385 record.ToolID = machineOpaqueValue(s.machineToolIDs, &s.nextMachineToolID, "tool", e.Tool.ID)
386 record.ToolName = machineOpaqueValue(s.machineToolNames, &s.nextMachineToolName, "tool_name", e.Tool.Name)
387 record.ToolReadOnly = e.Tool.ReadOnly
388 record.ToolError = e.Tool.Err != ""
389 record.ToolTruncated = e.Tool.Truncated
390 record.ToolDurationMS = e.Tool.DurationMs
391 case event.Usage:
392 if e.Usage != nil {
393 record.Usage = &machineEventUsage{InputTokens: e.Usage.PromptTokens, OutputTokens: e.Usage.CompletionTokens, CacheHitTokens: e.Usage.CacheHitTokens, CacheMissTokens: e.Usage.CacheMissTokens, Estimated: e.Usage.Estimated}
394 }
395 case event.ApprovalRequest:
396 record.ApprovalID = e.Approval.ID
397 record.ApprovalKind = e.Approval.Kind
398 case event.AskRequest:
399 record.AskID = e.Ask.ID
400 case event.TurnDone:
401 record.Outcome = e.Outcome
402 record.Cancelled = e.Cancelled
403 record.Error = e.Err != nil
404 case event.CompactionStarted, event.CompactionDone:
405 record.CompactionType = e.Compaction.Trigger
406 record.CompactionMsgs = e.Compaction.Messages
407 case event.GuardianAssessment:
408 record.GuardianResult = e.Guardian.Outcome
409 record.GuardianRisk = e.Guardian.RiskLevel
410 case event.Retrying:
411 record.Recovery = e.Recovery
412 record.RetryAttempt = e.RetryAttempt
413 record.RetryMax = e.RetryMax
414 }
415 return record
416 }
417
418 func machineOpaqueValue(values map[string]string, next *uint64, prefix, value string) string {
419 if value == "" {
420 return ""
421 }
422 if opaque := values[value]; opaque != "" {
423 return opaque
424 }
425 (*next)++
426 opaque := fmt.Sprintf("%s_%d", prefix, *next)
427 values[value] = opaque
428 return opaque
429 }
430
431 func machineEventKind(kind event.Kind) string {
432 names := eventwire.KindNames()
433 if int(kind) >= 0 && int(kind) < len(names) {
434 return names[kind]
435 }
436 return "unknown"
437 }
438
438 lines GO