返回 DeepSeek-Reasonix
client.go
根目录 / internal / telemetry / client.go
1 package telemetry
2
3 import (
4 "bytes"
5 "context"
6 "encoding/hex"
7 "encoding/json"
8 "fmt"
9 "io"
10 "net/http"
11 "os"
12 "path/filepath"
13 "runtime"
14 "sort"
15 "strings"
16 "time"
17
18 "reasonix/internal/fileutil"
19 "reasonix/internal/netclient"
20 )
21
22 var endpoint = "https://crash.reasonix.io/v1"
23
24 var uploadSignals = map[string]bool{
25 "finish_reason": true, "empty_final": true, "provider_error": true,
26 "cache_hit": true, "tool_error": true, "compaction": true, "turns": true,
27 "recovery_failure": true, "recovery_rule_continue": true,
28 "recovery_review_continue": true, "recovery_human_prompt": true,
29 "recovery_human_continue": true, "recovery_human_revise": true,
30 "recovery_review_error": true, "recovery_repeat_prompt": true,
31 "recovery_review_latency": true, "client_surface": true,
32 "client_version": true, "settings_language": true, "cli_mode": true,
33 "cli_permission_mode": true, "cli_session_mode": true,
34 "cli_turn_latency": true, "cli_exit": true,
35 "completion_validation_outcome": true, "completion_validation_latency": true,
36 "completion_validation_error": true, "completion_validation_attempt": true,
37 "completion_evaluator_finish_reason": true, "completion_evaluator_cache_hit": true,
38 }
39
40 type Client struct {
41 home string
42 version string
43 installID string
44 http *http.Client
45 }
46
47 func newClient(home, version string, proxy netclient.ProxySpec) (*Client, error) {
48 if strings.TrimSpace(home) == "" {
49 return nil, fmt.Errorf("telemetry: empty home")
50 }
51 client, err := netclient.NewHTTPClient(proxy, netclient.TransportOptions{
52 DialTimeout: 500 * time.Millisecond,
53 TLSHandshakeTimeout: 500 * time.Millisecond,
54 ResponseHeaderTimeout: 750 * time.Millisecond,
55 })
56 if err != nil {
57 return nil, err
58 }
59 client.Timeout = time.Second
60 id, err := installID(home)
61 if err != nil {
62 return nil, err
63 }
64 return &Client{home: home, version: version, installID: id, http: client}, nil
65 }
66
67 func installID(home string) (string, error) {
68 if err := os.MkdirAll(home, 0o700); err != nil {
69 return "", err
70 }
71 path := filepath.Join(home, "cli-telemetry-install-id")
72 if b, err := os.ReadFile(path); err == nil {
73 id := strings.TrimSpace(string(b))
74 if validInstallID(id) {
75 return id, nil
76 }
77 replacement, err := randomHex(16)
78 if err != nil {
79 return "", err
80 }
81 if err := fileutil.AtomicWriteFile(path, []byte(replacement+"\n"), 0o600); err != nil {
82 return "", err
83 }
84 return replacement, nil
85 } else if !os.IsNotExist(err) {
86 return "", err
87 }
88 id, err := randomHex(16)
89 if err != nil {
90 return "", err
91 }
92 f, err := os.OpenFile(path, os.O_WRONLY|os.O_CREATE|os.O_EXCL, 0o600)
93 if err != nil {
94 if b, readErr := os.ReadFile(path); readErr == nil && validInstallID(strings.TrimSpace(string(b))) {
95 return strings.TrimSpace(string(b)), nil
96 }
97 return "", err
98 }
99 if _, err = io.WriteString(f, id+"\n"); err != nil {
100 _ = f.Close()
101 _ = os.Remove(path)
102 return "", err
103 }
104 if err := f.Close(); err != nil {
105 return "", err
106 }
107 return id, nil
108 }
109
110 func validInstallID(id string) bool {
111 if len(id) != 32 {
112 return false
113 }
114 _, err := hex.DecodeString(id)
115 return err == nil
116 }
117
118 func (c *Client) backgroundFlush() {
119 ctx, cancel := context.WithTimeout(context.Background(), time.Second)
120 defer cancel()
121 _ = c.sendDailyPing(ctx)
122 _ = c.flushPending(ctx)
123 }
124
125 func (c *Client) sendDailyPing(ctx context.Context) error {
126 day := time.Now().UTC().Format("2006-01-02")
127 claim := filepath.Join(c.home, "cli-telemetry-ping-"+day)
128 f, err := os.OpenFile(claim, os.O_WRONLY|os.O_CREATE|os.O_EXCL, 0o600)
129 if err != nil {
130 if os.IsExist(err) {
131 state, readErr := os.ReadFile(claim)
132 if readErr == nil && strings.TrimSpace(string(state)) == "sent" {
133 return nil
134 }
135 if info, statErr := os.Stat(claim); statErr == nil && time.Since(info.ModTime()) < 2*time.Minute {
136 return nil
137 }
138 if os.Remove(claim) == nil {
139 return c.sendDailyPing(ctx)
140 }
141 }
142 return err
143 }
144 if _, err := io.WriteString(f, "sending\n"); err != nil {
145 _ = f.Close()
146 _ = os.Remove(claim)
147 return err
148 }
149 if err := f.Close(); err != nil {
150 _ = os.Remove(claim)
151 return err
152 }
153 err = c.post(ctx, "/ping", pingPayload{
154 InstallID: c.installID, Version: c.version, OS: runtime.GOOS, Arch: runtime.GOARCH, Surface: "cli",
155 })
156 if err != nil {
157 _ = os.Remove(claim)
158 } else if writeErr := os.WriteFile(claim, []byte("sent\n"), 0o600); writeErr != nil {
159 _ = os.Remove(claim)
160 err = writeErr
161 }
162 c.prunePingClaims(day)
163 return err
164 }
165
166 func (c *Client) prunePingClaims(current string) {
167 entries, _ := os.ReadDir(c.home)
168 for _, entry := range entries {
169 if strings.HasPrefix(entry.Name(), "cli-telemetry-ping-") && entry.Name() != "cli-telemetry-ping-"+current {
170 _ = os.Remove(filepath.Join(c.home, entry.Name()))
171 }
172 }
173 }
174
175 func (c *Client) flushPending(ctx context.Context) error {
176 dir := filepath.Join(c.home, pendingDirName)
177 paths, err := claimPendingFiles(dir, time.Now())
178 if err != nil {
179 if os.IsNotExist(err) {
180 return nil
181 }
182 return err
183 }
184 claimed := make(map[string]bool, len(paths))
185 for _, path := range paths {
186 claimed[path] = true
187 }
188 defer func() {
189 for path, active := range claimed {
190 if active {
191 _ = os.Rename(path, strings.TrimSuffix(path, ".uploading"))
192 }
193 }
194 }()
195 type group struct {
196 payload pendingPayload
197 paths []string
198 counts map[string]int
199 }
200 groups := map[string]*group{}
201 for _, path := range paths {
202 b, err := os.ReadFile(path)
203 if err != nil {
204 continue
205 }
206 var p pendingPayload
207 if json.Unmarshal(b, &p) != nil || !validPendingPayload(p) {
208 removeClaim(path, claimed)
209 continue
210 }
211 key := p.Version + "\x00" + p.OS
212 g := groups[key]
213 if g == nil {
214 g = &group{payload: pendingPayload{Version: p.Version, OS: p.OS}, counts: map[string]int{}}
215 groups[key] = g
216 }
217 g.paths = append(g.paths, path)
218 for _, counter := range p.Counters {
219 if validCounter(counter) {
220 g.counts[counter.Signal+"\x00"+counter.Bucket] += counter.Count
221 }
222 }
223 }
224 for _, g := range groups {
225 counters := make([]Counter, 0, len(g.counts))
226 for key, count := range g.counts {
227 signal, bucket, _ := strings.Cut(key, "\x00")
228 if count > 1_000_000 {
229 count = 1_000_000
230 }
231 counters = append(counters, Counter{Signal: signal, Bucket: bucket, Count: count})
232 }
233 sort.Slice(counters, func(i, j int) bool {
234 if counters[i].Signal == counters[j].Signal {
235 return counters[i].Bucket < counters[j].Bucket
236 }
237 return counters[i].Signal < counters[j].Signal
238 })
239 if len(counters) == 0 {
240 for _, path := range g.paths {
241 removeClaim(path, claimed)
242 }
243 continue
244 }
245 if err := c.post(ctx, "/metrics", metricsPayload{
246 Version: g.payload.Version, OS: g.payload.OS, Surface: "cli", Counters: counters,
247 }); err != nil {
248 return err
249 }
250 for _, path := range g.paths {
251 removeClaim(path, claimed)
252 }
253 }
254 return nil
255 }
256
257 func removeClaim(path string, claimed map[string]bool) {
258 if err := os.Remove(path); err == nil || os.IsNotExist(err) {
259 claimed[path] = false
260 }
261 }
262
263 func claimPendingFiles(dir string, now time.Time) ([]string, error) {
264 entries, err := os.ReadDir(dir)
265 if err != nil {
266 return nil, err
267 }
268 for _, entry := range entries {
269 if entry.IsDir() || !strings.HasSuffix(entry.Name(), ".json.uploading") {
270 continue
271 }
272 path := filepath.Join(dir, entry.Name())
273 info, err := entry.Info()
274 if err != nil || now.Sub(info.ModTime()) < 2*time.Minute {
275 continue
276 }
277 _ = os.Rename(path, strings.TrimSuffix(path, ".uploading"))
278 }
279 entries, err = os.ReadDir(dir)
280 if err != nil {
281 return nil, err
282 }
283 var claimed []string
284 for _, entry := range entries {
285 if entry.IsDir() || !strings.HasSuffix(entry.Name(), ".json") {
286 continue
287 }
288 path := filepath.Join(dir, entry.Name())
289 claim := path + ".uploading"
290 if os.Rename(path, claim) == nil {
291 claimed = append(claimed, claim)
292 }
293 }
294 return claimed, nil
295 }
296
297 func validPendingPayload(p pendingPayload) bool {
298 if !releaseVersionPattern.MatchString(strings.TrimSpace(p.Version)) || len(p.Version) > 64 {
299 return false
300 }
301 switch p.OS {
302 case "android", "darwin", "freebsd", "linux", "windows":
303 default:
304 return false
305 }
306 return len(p.Counters) > 0 && len(p.Counters) <= 128
307 }
308
309 func validCounter(c Counter) bool {
310 return uploadSignals[c.Signal] && c.Count > 0 && c.Count <= 1_000_000 &&
311 len(c.Bucket) > 0 && len(c.Bucket) <= 96 && !unsafeBucketChars.MatchString(c.Bucket)
312 }
313
314 func (c *Client) post(ctx context.Context, path string, payload any) error {
315 b, err := json.Marshal(payload)
316 if err != nil {
317 return err
318 }
319 req, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint+path, bytes.NewReader(b))
320 if err != nil {
321 return err
322 }
323 req.Header.Set("Content-Type", "application/json")
324 resp, err := c.http.Do(req)
325 if err != nil {
326 return err
327 }
328 defer resp.Body.Close()
329 _, _ = io.Copy(io.Discard, io.LimitReader(resp.Body, 1024))
330 if resp.StatusCode < 200 || resp.StatusCode >= 300 {
331 return fmt.Errorf("telemetry: HTTP %d", resp.StatusCode)
332 }
333 return nil
334 }
335
335 lines GO