返回 DeepSeek-Reasonix
sink.go
根目录 / internal / telemetry / sink.go
1 package telemetry
2
3 import (
4 "context"
5 "errors"
6 "net"
7 "regexp"
8 "runtime"
9 "strings"
10 "time"
11
12 "reasonix/internal/event"
13 "reasonix/internal/evidence"
14 "reasonix/internal/netclient"
15 "reasonix/internal/provider"
16 "reasonix/internal/recovery"
17 )
18
19 type Options struct {
20 Mode string
21 Version string
22 HomeDir string
23 Interactive bool
24 Proxy netclient.ProxySpec
25 CLIMode string
26 Profile string
27 PermissionMode string
28 SessionMode string
29 Language string
30 }
31
32 type Reporter struct {
33 client *Client
34 version string
35 home string
36 static []Counter
37 }
38
39 func Start(opts Options) *Reporter {
40 if !Enabled(opts.Mode, opts.Version, opts.Interactive) {
41 if strings.EqualFold(strings.TrimSpace(opts.Mode), "off") || envOptOut() {
42 _ = Cleanup(opts.HomeDir)
43 }
44 return nil
45 }
46 client, err := newClient(opts.HomeDir, opts.Version, opts.Proxy)
47 if err != nil {
48 return nil
49 }
50 r := &Reporter{
51 client: client,
52 version: opts.Version,
53 home: opts.HomeDir,
54 static: []Counter{
55 {Signal: "client_surface", Bucket: "cli", Count: 1},
56 {Signal: "client_version", Bucket: safeBucket(opts.Version, "other"), Count: 1},
57 {Signal: "cli_mode", Bucket: enumBucket(opts.CLIMode, "run", "tui"), Count: 1},
58 {Signal: "cli_profile", Bucket: enumBucket(opts.Profile, "economy", "balanced", "delivery"), Count: 1},
59 {Signal: "cli_permission_mode", Bucket: permissionBucket(opts.PermissionMode), Count: 1},
60 {Signal: "cli_session_mode", Bucket: enumBucket(opts.SessionMode, "fresh", "resume", "continue", "copy"), Count: 1},
61 {Signal: "settings_language", Bucket: languageBucket(opts.Language), Count: 1},
62 },
63 }
64 go client.backgroundFlush()
65 return r
66 }
67
68 func (r *Reporter) Wrap(inner event.Sink) event.Sink {
69 if r == nil {
70 return inner
71 }
72 return &sink{inner: inner, reporter: r, counts: countersFrom(r.static)}
73 }
74
75 func (r *Reporter) RecordRecovery(m recovery.Metrics) {
76 if r == nil {
77 return
78 }
79 counts := map[string]int{}
80 addMetric(counts, "recovery_failure", "count", m.FailureEvents)
81 addMetric(counts, "recovery_rule_continue", "count", m.RuleContinues)
82 addMetric(counts, "recovery_review_continue", "count", m.ReviewContinues)
83 addMetric(counts, "recovery_human_prompt", "count", m.HumanPrompts)
84 addMetric(counts, "recovery_human_continue", "count", m.HumanContinues)
85 addMetric(counts, "recovery_human_revise", "count", m.HumanRevises)
86 addMetric(counts, "recovery_review_error", "count", m.ReviewErrors)
87 addMetric(counts, "recovery_repeat_prompt", "count", m.RepeatPrompts)
88 if m.ReviewLatencyCount > 0 {
89 add(counts, "recovery_review_latency", latencyBucket(time.Duration(m.ReviewLatencyMsSum/m.ReviewLatencyCount)*time.Millisecond), int(m.ReviewLatencyCount))
90 }
91 r.append(counts)
92 }
93
94 func addMetric(counts map[string]int, signal, bucket string, count int64) {
95 if count > 0 {
96 add(counts, signal, bucket, int(count))
97 }
98 }
99
100 func (r *Reporter) append(counts map[string]int) {
101 if r == nil || len(counts) == 0 {
102 return
103 }
104 counters := make([]Counter, 0, len(counts))
105 for key, count := range counts {
106 signal, bucket, _ := strings.Cut(key, "\x00")
107 if count > 1_000_000 {
108 count = 1_000_000
109 }
110 counters = append(counters, Counter{Signal: signal, Bucket: bucket, Count: count})
111 }
112 _ = appendPending(r.home, pendingPayload{Version: r.version, OS: runtime.GOOS, Counters: counters})
113 }
114
115 type sink struct {
116 inner event.Sink
117 reporter *Reporter
118 counts map[string]int
119 started time.Time
120 hasText bool
121 emptyFinalSeen bool
122 }
123
124 func (s *sink) Emit(e event.Event) {
125 s.observe(e)
126 s.inner.Emit(e)
127 }
128
129 func (s *sink) RecordReadinessAudit(a evidence.ReadinessAudit) {
130 event.RecordReadinessAudit(s.inner, a)
131 }
132
133 func (s *sink) RecordProtocolRecovery(a event.ProtocolRecoveryAudit) {
134 add(s.counts, "tool_call_reasoning_recovery", string(a.Kind), 1)
135 event.RecordProtocolRecovery(s.inner, a)
136 }
137
138 func (s *sink) observe(e event.Event) {
139 switch e.Kind {
140 case event.TurnStarted:
141 s.started = time.Now()
142 s.hasText = false
143 s.emptyFinalSeen = false
144 add(s.counts, "turns", "count", 1)
145 case event.Text:
146 if e.Text != "" {
147 s.hasText = true
148 }
149 case event.Message:
150 if e.Text != "" {
151 s.hasText = true
152 }
153 case event.Usage:
154 if e.Usage != nil {
155 add(s.counts, "finish_reason", finishReasonBucket(e.Usage.FinishReason), 1)
156 add(s.counts, "cache_hit", cacheBucket(e.Usage.CacheHitTokens, e.Usage.CacheMissTokens), 1)
157 }
158 case event.ToolResult:
159 if e.Tool.Err != "" {
160 add(s.counts, "tool_error", toolErrorBucket(e.Tool.Err), 1)
161 }
162 case event.Notice:
163 if e.Code == event.NoticeCodeEmptyFinal {
164 add(s.counts, "empty_final", "yes", 1)
165 s.emptyFinalSeen = true
166 }
167 case event.CompactionStarted:
168 add(s.counts, "compaction", enumBucket(e.Compaction.Trigger, "auto", "manual"), 1)
169 case event.TurnDone:
170 if !s.hasText && e.Err == nil && !s.emptyFinalSeen {
171 add(s.counts, "empty_final", "yes", 1)
172 }
173 if bucket := providerErrorBucket(e.Err); bucket != "" {
174 add(s.counts, "provider_error", bucket, 1)
175 }
176 add(s.counts, "cli_exit", exitBucket(e), 1)
177 if !s.started.IsZero() {
178 add(s.counts, "cli_turn_latency", latencyBucket(time.Since(s.started)), 1)
179 }
180 s.reporter.append(s.counts)
181 s.counts = map[string]int{}
182 s.started = time.Time{}
183 s.hasText = false
184 s.emptyFinalSeen = false
185 }
186 }
187
188 func countersFrom(in []Counter) map[string]int {
189 out := map[string]int{}
190 for _, c := range in {
191 add(out, c.Signal, c.Bucket, c.Count)
192 }
193 return out
194 }
195
196 func add(counts map[string]int, signal, bucket string, count int) {
197 if count <= 0 || signal == "" || bucket == "" {
198 return
199 }
200 counts[signal+"\x00"+bucket] += count
201 }
202
203 var unsafeBucketChars = regexp.MustCompile(`[^a-z0-9_]+`)
204
205 func safeBucket(value, fallback string) string {
206 value = strings.ToLower(strings.TrimSpace(value))
207 value = unsafeBucketChars.ReplaceAllString(value, "_")
208 value = strings.Trim(value, "_")
209 if value == "" {
210 return fallback
211 }
212 if len(value) > 96 {
213 value = value[:96]
214 }
215 return value
216 }
217
218 func enumBucket(value string, allowed ...string) string {
219 value = strings.ToLower(strings.TrimSpace(value))
220 for _, item := range allowed {
221 if value == item {
222 return value
223 }
224 }
225 return "other"
226 }
227
228 func permissionBucket(value string) string {
229 switch strings.ToLower(strings.TrimSpace(value)) {
230 case "manual", "ask":
231 return "ask"
232 case "auto", "acceptedits":
233 return "auto"
234 case "dontask":
235 return "dont_ask"
236 case "plan":
237 return "plan"
238 case "bypasspermissions", "yolo":
239 return "yolo"
240 default:
241 return "other"
242 }
243 }
244
245 func languageBucket(value string) string {
246 value = strings.ToLower(strings.TrimSpace(value))
247 if strings.HasPrefix(value, "zh") {
248 return "zh"
249 }
250 if strings.HasPrefix(value, "en") {
251 return "en"
252 }
253 if value == "" || value == "auto" {
254 return "auto"
255 }
256 return "other"
257 }
258
259 func finishReasonBucket(value string) string {
260 switch strings.ToLower(strings.TrimSpace(value)) {
261 case "stop", "tool_calls", "length", "content_filter", "repetition_truncation":
262 return safeBucket(value, "unknown")
263 case "":
264 return "unknown"
265 default:
266 return "other"
267 }
268 }
269
270 func cacheBucket(hit, miss int) string {
271 total := hit + miss
272 if total <= 0 {
273 return "unknown"
274 }
275 pct := hit * 100 / total
276 switch {
277 case pct == 0:
278 return "0"
279 case pct < 25:
280 return "1_24"
281 case pct < 50:
282 return "25_49"
283 case pct < 75:
284 return "50_74"
285 case pct < 90:
286 return "75_89"
287 default:
288 return "90_100"
289 }
290 }
291
292 func toolErrorBucket(value string) string {
293 v := strings.ToLower(value)
294 switch {
295 case strings.Contains(v, "permission"), strings.Contains(v, "blocked"), strings.Contains(v, "denied"):
296 return "permission"
297 case strings.Contains(v, "timeout"), strings.Contains(v, "deadline"):
298 return "timeout"
299 case strings.Contains(v, "cancel"):
300 return "cancelled"
301 case strings.Contains(v, "not found"), strings.Contains(v, "no such"):
302 return "not_found"
303 default:
304 return "other"
305 }
306 }
307
308 func providerErrorBucket(err error) string {
309 if err == nil {
310 return ""
311 }
312 var auth *provider.AuthError
313 if errors.As(err, &auth) {
314 return "auth"
315 }
316 var api *provider.APIError
317 if errors.As(err, &api) {
318 switch {
319 case api.Status == 429:
320 return "rate_limit"
321 case api.Status >= 500:
322 return "server"
323 case api.Status >= 400:
324 return "request"
325 default:
326 return "http"
327 }
328 }
329 if errors.Is(err, context.DeadlineExceeded) {
330 return "timeout"
331 }
332 if errors.Is(err, context.Canceled) {
333 return "cancelled"
334 }
335 var netErr net.Error
336 if errors.As(err, &netErr) {
337 return "network"
338 }
339 if provider.IsStreamInterrupted(err) {
340 return "interrupted"
341 }
342 return ""
343 }
344
345 func latencyBucket(d time.Duration) string {
346 switch {
347 case d < time.Second:
348 return "lt_1s"
349 case d < 5*time.Second:
350 return "s_1_5"
351 case d < 15*time.Second:
352 return "s_5_15"
353 case d < time.Minute:
354 return "s_15_60"
355 case d < 5*time.Minute:
356 return "m_1_5"
357 case d < 15*time.Minute:
358 return "m_5_15"
359 default:
360 return "m_15_plus"
361 }
362 }
363
364 func exitBucket(e event.Event) string {
365 if e.Cancelled || errors.Is(e.Err, context.Canceled) {
366 return "cancelled"
367 }
368 if e.Outcome == event.TurnOutcomeRecoveryPaused {
369 return "recovery_paused"
370 }
371 if e.Err != nil {
372 return "error"
373 }
374 return "success"
375 }
376
376 lines GO