返回 DeepSeek-Reasonix
sink_test.go
根目录 / internal / telemetry / sink_test.go
1 package telemetry
2
3 import (
4 "encoding/json"
5 "errors"
6 "os"
7 "path/filepath"
8 "strings"
9 "testing"
10
11 "reasonix/internal/event"
12 "reasonix/internal/evidence"
13 "reasonix/internal/provider"
14 )
15
16 type readinessSink struct {
17 events int
18 audits int
19 recovery int
20 }
21
22 func TestIncompleteReadIsPausedNotProviderError(t *testing.T) {
23 if got := exitBucket(event.Event{Outcome: event.TurnOutcomeIncompleteRead, Err: errors.New("incomplete read")}); got != "incomplete_read" {
24 t.Fatalf("exit bucket = %q", got)
25 }
26 }
27
28 func (s *readinessSink) Emit(event.Event) { s.events++ }
29 func (s *readinessSink) RecordReadinessAudit(evidence.ReadinessAudit) {
30 s.audits++
31 }
32 func (s *readinessSink) RecordProtocolRecovery(event.ProtocolRecoveryAudit) {
33 s.recovery++
34 }
35
36 func TestSinkWritesOnlyWhitelistedContentFreeCounters(t *testing.T) {
37 home := t.TempDir()
38 reporter := &Reporter{
39 home: home,
40 version: "v1.20.0",
41 static: []Counter{
42 {Signal: "client_surface", Bucket: "cli", Count: 1},
43 {Signal: "cli_mode", Bucket: "run", Count: 1},
44 },
45 }
46 inner := &readinessSink{}
47 sink := reporter.Wrap(inner)
48 secret := "PRIVATE_PROMPT_TOKEN_123"
49 sink.Emit(event.Event{Kind: event.TurnStarted})
50 sink.Emit(event.Event{Kind: event.Text, Text: secret})
51 sink.Emit(event.Event{Kind: event.Message, Text: secret, Reasoning: secret})
52 sink.Emit(event.Event{Kind: event.Usage, Usage: &provider.Usage{
53 FinishReason: "stop", CacheHitTokens: 90, CacheMissTokens: 10,
54 }})
55 sink.Emit(event.Event{Kind: event.ToolResult, Tool: event.Tool{
56 Name: secret, Args: secret, Output: secret, Err: "permission denied: " + secret,
57 }})
58 event.RecordProtocolRecovery(sink, event.ProtocolRecoveryAudit{Kind: event.ProtocolRecoveryMissingReasoningRetryReplaced})
59 sink.Emit(event.Event{Kind: event.TurnDone, Err: &provider.APIError{
60 Provider: secret, Status: 429, Body: secret, TraceID: secret,
61 }})
62 event.RecordReadinessAudit(sink, evidence.ReadinessAudit{})
63
64 entries, err := os.ReadDir(filepath.Join(home, pendingDirName))
65 if err != nil || len(entries) != 1 {
66 t.Fatalf("pending files = %d, err = %v", len(entries), err)
67 }
68 b, err := os.ReadFile(filepath.Join(home, pendingDirName, entries[0].Name()))
69 if err != nil {
70 t.Fatal(err)
71 }
72 if strings.Contains(string(b), secret) {
73 t.Fatalf("pending payload leaked private content: %s", b)
74 }
75 var payload pendingPayload
76 if err := json.Unmarshal(b, &payload); err != nil {
77 t.Fatal(err)
78 }
79 got := map[string]string{}
80 for _, counter := range payload.Counters {
81 got[counter.Signal] = counter.Bucket
82 }
83 for signal, bucket := range map[string]string{
84 "client_surface": "cli",
85 "cli_mode": "run",
86 "turns": "count",
87 "finish_reason": "stop",
88 "cache_hit": "90_100",
89 "tool_error": "permission",
90 "provider_error": "rate_limit",
91 "cli_exit": "error",
92 "tool_call_reasoning_recovery": "missing_reasoning_retry_replaced_response",
93 } {
94 if got[signal] != bucket {
95 t.Errorf("%s bucket = %q, want %q", signal, got[signal], bucket)
96 }
97 }
98 if inner.events != 6 || inner.audits != 1 || inner.recovery != 1 {
99 t.Fatalf("forwarding events=%d audits=%d recovery=%d", inner.events, inner.audits, inner.recovery)
100 }
101 }
102
103 func TestCleanupRemovesPendingQueueOnly(t *testing.T) {
104 home := t.TempDir()
105 if err := appendPending(home, pendingPayload{
106 Version: "v1.20.0", OS: "linux", Counters: []Counter{{Signal: "turns", Bucket: "count", Count: 1}},
107 }); err != nil {
108 t.Fatal(err)
109 }
110 idPath := filepath.Join(home, "cli-telemetry-install-id")
111 if err := os.WriteFile(idPath, []byte(strings.Repeat("a", 32)), 0o600); err != nil {
112 t.Fatal(err)
113 }
114 if err := Cleanup(home); err != nil {
115 t.Fatal(err)
116 }
117 if _, err := os.Stat(filepath.Join(home, pendingDirName)); !errors.Is(err, os.ErrNotExist) {
118 t.Fatalf("pending directory still exists: %v", err)
119 }
120 if _, err := os.Stat(idPath); err != nil {
121 t.Fatalf("install id should remain stable after opt-out cleanup: %v", err)
122 }
123 }
124
125 func TestEnvironmentOptOutRemovesPendingQueue(t *testing.T) {
126 clearPolicyEnv(t)
127 home := t.TempDir()
128 if err := appendPending(home, pendingPayload{
129 Version: "v1.20.0", OS: "linux", Counters: []Counter{{Signal: "turns", Bucket: "count", Count: 1}},
130 }); err != nil {
131 t.Fatal(err)
132 }
133 t.Setenv("DO_NOT_TRACK", "1")
134 if reporter := Start(Options{Mode: "on", Version: "v1.20.0", HomeDir: home, Interactive: true}); reporter != nil {
135 t.Fatal("environment opt-out unexpectedly started telemetry")
136 }
137 if _, err := os.Stat(filepath.Join(home, pendingDirName)); !errors.Is(err, os.ErrNotExist) {
138 t.Fatalf("environment opt-out did not remove pending queue: %v", err)
139 }
140 }
141
141 lines GO