返回 DeepSeek-Reasonix
recorder.go
根目录 / internal / stats / recorder.go
1 package stats
2
3 import (
4 "context"
5 "strings"
6 "sync"
7 "time"
8
9 "reasonix/internal/billing"
10 "reasonix/internal/event"
11 "reasonix/internal/evidence"
12 "reasonix/internal/provider"
13 )
14
15 // Recorder is a passthrough event.Sink that snapshots token usage (event.Usage)
16 // and completed turns (event.TurnDone) into the daily stats files. It observes
17 // only; it never alters the event stream.
18 //
19 // Wire it around the frontend sink at the boot layer so every entry point
20 // (desktop, CLI, serve) records consistently; Source distinguishes them.
21 type Recorder struct {
22 inner event.Sink
23 writer *Writer
24 dispatcher *recordDispatcher
25 source string
26 }
27
28 var _ event.OptionalSinkCapabilities = (*Recorder)(nil)
29
30 const recorderQueueSize = 2048
31
32 type dispatchItem struct {
33 record record
34 flush chan struct{}
35 }
36
37 // recordDispatcher keeps filesystem latency off provider/UI event goroutines.
38 // Dispatchers are shared per state directory, so controller rebuilds do not
39 // create one goroutine per recorder instance; CloseRecordDispatchers retires
40 // them once their directory's process lifetime ends.
41 type recordDispatcher struct {
42 writer *Writer
43 queue chan dispatchItem
44 // stop requests the run goroutine to drain the queue and exit. The queue
45 // channel itself is never closed: Recorders can outlive a shutdown, and
46 // enqueue/flush must never send on a closed channel.
47 stop chan struct{}
48 stopped chan struct{}
49 }
50
51 var recorderDispatchers = struct {
52 sync.Mutex
53 byDir map[string]*recordDispatcher
54 }{byDir: map[string]*recordDispatcher{}}
55
56 func dispatcherFor(writer *Writer) *recordDispatcher {
57 if writer == nil || writer.dir == "" {
58 return nil
59 }
60 recorderDispatchers.Lock()
61 defer recorderDispatchers.Unlock()
62 if dispatcher := recorderDispatchers.byDir[writer.dir]; dispatcher != nil {
63 return dispatcher
64 }
65 dispatcher := &recordDispatcher{writer: writer, queue: make(chan dispatchItem, recorderQueueSize), stop: make(chan struct{}), stopped: make(chan struct{})}
66 recorderDispatchers.byDir[writer.dir] = dispatcher
67 go dispatcher.run()
68 return dispatcher
69 }
70
71 func existingDispatcher(dir string) *recordDispatcher {
72 if strings.TrimSpace(dir) == "" {
73 return nil
74 }
75 recorderDispatchers.Lock()
76 defer recorderDispatchers.Unlock()
77 return recorderDispatchers.byDir[dir]
78 }
79
80 func (d *recordDispatcher) run() {
81 defer close(d.stopped)
82 for {
83 select {
84 case item := <-d.queue:
85 d.process(item)
86 case <-d.stop:
87 // Drain what was already accepted, then exit so a process moving
88 // through many state directories does not keep one goroutine per
89 // dir. Later enqueues are dropped, matching the best-effort queue.
90 for {
91 select {
92 case item := <-d.queue:
93 d.process(item)
94 default:
95 return
96 }
97 }
98 }
99 }
100 }
101
102 func (d *recordDispatcher) process(item dispatchItem) {
103 if item.flush != nil {
104 close(item.flush)
105 return
106 }
107 _ = d.writer.Append(item.record)
108 }
109
110 func (d *recordDispatcher) enqueue(rec record) {
111 if d == nil {
112 return
113 }
114 // Statistics are observational. A full queue may lose a record, but it must
115 // never apply backpressure to model streaming or turn completion.
116 select {
117 case d.queue <- dispatchItem{record: rec}:
118 default:
119 }
120 }
121
122 func (d *recordDispatcher) flush(ctx context.Context) error {
123 if d == nil {
124 return nil
125 }
126 if ctx == nil {
127 ctx = context.Background()
128 }
129 done := make(chan struct{})
130 select {
131 case d.queue <- dispatchItem{flush: done}:
132 case <-d.stopped:
133 return nil
134 case <-ctx.Done():
135 return ctx.Err()
136 }
137 select {
138 case <-done:
139 return nil
140 case <-d.stopped:
141 return nil
142 case <-ctx.Done():
143 return ctx.Err()
144 }
145 }
146
147 // CloseRecordDispatchers stops every process-local record dispatcher, drains
148 // what was already queued, and forgets the per-directory cache so a process
149 // (or test binary) moving through many state directories does not accumulate
150 // one goroutine per dir. Mirrors CloseUsageCatalogs for shutdown and test
151 // isolation boundaries. Each dispatcher is captured exactly once under the
152 // registry lock, so stop is closed exactly once per generation.
153 func CloseRecordDispatchers(ctx context.Context) error {
154 if ctx == nil {
155 ctx = context.Background()
156 }
157 recorderDispatchers.Lock()
158 dispatchers := make([]*recordDispatcher, 0, len(recorderDispatchers.byDir))
159 for dir, dispatcher := range recorderDispatchers.byDir {
160 dispatchers = append(dispatchers, dispatcher)
161 delete(recorderDispatchers.byDir, dir)
162 }
163 recorderDispatchers.Unlock()
164 for _, dispatcher := range dispatchers {
165 close(dispatcher.stop)
166 }
167 var first error
168 for _, dispatcher := range dispatchers {
169 select {
170 case <-dispatcher.stopped:
171 case <-ctx.Done():
172 if first == nil {
173 first = ctx.Err()
174 }
175 }
176 }
177 return first
178 }
179
180 // NewRecorder wraps inner with usage recording. source labels every record
181 // (desktop/cli/serve/...); an empty source keeps records unlabelled.
182 func NewRecorder(inner event.Sink, dir, source string) *Recorder {
183 writer := NewWriter(dir)
184 writer.usage = managerForUsage(writer.dir)
185 return &Recorder{
186 inner: inner, writer: writer, dispatcher: dispatcherFor(writer), source: strings.TrimSpace(source),
187 }
188 }
189
190 // Emit forwards user-visible events unchanged, then queues any usage/turn
191 // record without waiting for filesystem I/O. Request-only usage is internal
192 // accounting for failed provider calls, so it is persisted without surfacing a
193 // zero-token receipt in the wrapped frontend.
194 func (r *Recorder) Emit(e event.Event) {
195 requestOnly := e.Kind == event.Usage && e.Usage != nil && e.Usage.TotalTokens <= 0 && e.Usage.RequestCount > 0
196 if r != nil && r.inner != nil && !requestOnly {
197 r.inner.Emit(e)
198 }
199 if r != nil && r.writer != nil && e.Kind == event.Usage {
200 r.recordUsage(e)
201 } else if r != nil && r.writer != nil && e.Kind == event.GuardianAssessment && e.Guardian.Usage != nil {
202 r.recordProviderUsage(e.ModelRef, e.Guardian.Usage, nil, "")
203 } else if r != nil && r.writer != nil && e.Kind == event.TurnDone {
204 r.recordTurnCompletion()
205 }
206 }
207
208 // RecordTurnCompletion records synchronous controller runs that deliberately do
209 // not emit TurnDone into the UI event stream.
210 func (r *Recorder) RecordTurnCompletion() {
211 r.recordTurnCompletion()
212 if r != nil {
213 event.RecordTurnCompletion(r.inner)
214 }
215 }
216
217 func (r *Recorder) recordTurnCompletion() {
218 if r == nil || r.dispatcher == nil {
219 return
220 }
221 r.dispatcher.enqueue(record{Timestamp: time.Now(), Source: r.source, Turn: true})
222 }
223
224 // Flush waits until records already accepted by this recorder's shared queue
225 // have been written. Production event paths never call Flush; it exists for
226 // shutdown/verification boundaries that can explicitly tolerate waiting.
227 func (r *Recorder) Flush(ctx context.Context) error {
228 if r == nil {
229 return nil
230 }
231 if err := r.dispatcher.flush(ctx); err != nil {
232 return err
233 }
234 if r.writer != nil && r.writer.usage != nil {
235 if catalog := r.writer.usage.catalog.Load(); catalog != nil {
236 return catalog.Flush(ctx)
237 }
238 }
239 return nil
240 }
241
242 // Flush waits for records already queued for dir. It is primarily useful when
243 // a caller must read its own just-recorded statistics deterministically.
244 func Flush(ctx context.Context, dir string) error {
245 dir = strings.TrimSpace(dir)
246 if err := existingDispatcher(dir).flush(ctx); err != nil {
247 return err
248 }
249 if manager := existingUsageManager(dir); manager != nil {
250 if catalog := manager.catalog.Load(); catalog != nil {
251 return catalog.Flush(ctx)
252 }
253 }
254 return nil
255 }
256
257 // RecordReadinessAudit forwards audit receipts to the wrapped sink.
258 func (r *Recorder) RecordReadinessAudit(a evidence.ReadinessAudit) {
259 event.RecordReadinessAudit(r.inner, a)
260 }
261
262 func (r *Recorder) RecordAnchorSafetyAudit(a event.AnchorSafetyAudit) {
263 event.RecordAnchorSafetyAudit(r.inner, a)
264 }
265
266 // RecordProtocolRecovery preserves the wrapped sink's audit capability.
267 func (r *Recorder) RecordProtocolRecovery(a event.ProtocolRecoveryAudit) {
268 event.RecordProtocolRecovery(r.inner, a)
269 }
270
271 // RecordContractShadow preserves the wrapped sink's audit capability.
272 func (r *Recorder) RecordContractShadow(a event.ContractShadowAudit) {
273 event.RecordContractShadow(r.inner, a)
274 }
275
276 // RecordCompletionReport preserves the wrapped sink's audit capability.
277 func (r *Recorder) RecordDelegationAudit(a evidence.DelegationAudit) {
278 event.RecordDelegationAudit(r.inner, a)
279 }
280
281 func (r *Recorder) RecordCompletionReport(a event.CompletionReportAudit) {
282 event.RecordCompletionReport(r.inner, a)
283 }
284
285 // RecordOutcomeProgress preserves the wrapped sink's audit capability.
286 func (r *Recorder) RecordOutcomeProgress(sample evidence.OutcomeSample) {
287 event.RecordOutcomeProgress(r.inner, sample)
288 }
289
290 // RecordMemoryRecall preserves the wrapped sink's audit capability.
291 func (r *Recorder) RecordMemoryRecall(a event.MemoryRecallAudit) {
292 event.RecordMemoryRecall(r.inner, a)
293 }
294
295 // RecordDelegationAdmission preserves the wrapped sink's audit capability.
296 func (r *Recorder) RecordDelegationAdmission(a event.DelegationAdmissionAudit) {
297 event.RecordDelegationAdmission(r.inner, a)
298 }
299
300 func (r *Recorder) RecordWorkspaceMutation(m event.WorkspaceMutation) {
301 event.RecordWorkspaceMutation(r.inner, m)
302 }
303
304 func (r *Recorder) RecordRunBudget(sample event.RunBudgetSample) {
305 event.RecordRunBudget(r.inner, sample)
306 }
307
308 func (r *Recorder) RecordSubagentLifecycle(info event.SubagentLifecycleInfo) {
309 event.RecordSubagentLifecycle(r.inner, info)
310 }
311
312 func (r *Recorder) recordUsage(e event.Event) {
313 r.recordProviderUsage(e.ModelRef, e.Usage, e.CostQuote, e.UsageSource)
314 }
315
316 func (r *Recorder) recordProviderUsage(modelRef string, usage *provider.Usage, quote *billing.CostQuote, usageSource string) {
317 if usage == nil || (usage.TotalTokens <= 0 && usage.RequestCount <= 0) {
318 return
319 }
320 // Recording is best-effort: a stats file failure (disk full, permissions)
321 // must never interrupt the event stream, matching telemetry's append idiom.
322 rec := record{
323 Timestamp: time.Now(),
324 ModelRef: modelRef,
325 Source: r.source,
326 Prompt: usage.PromptTokens,
327 Completion: usage.CompletionTokens,
328 Reasoning: usage.ReasoningTokens,
329 CacheHit: usage.CacheHitTokens,
330 CacheMiss: usage.CacheMissTokens,
331 Total: usage.TotalTokens,
332 Requests: usageRequestCount(usage),
333 UsageSource: strings.TrimSpace(usageSource),
334 }
335 if quote != nil {
336 rec.CostAmount = quote.Original.Amount
337 rec.CostCurrency = quote.Original.Currency
338 rec.PricingFingerprint = quote.PricingFingerprint
339 rec.RateDate = quote.RateDate
340 rec.RateBand = quote.RateBand
341 rec.RatedAt = quote.RatedAt
342 rec.IncompleteReason = quote.IncompleteReason
343 rec.BillingMode = quote.BillingMode
344 rec.CostEstimated = quote.Estimated
345 rec.LegacyEstimate = quote.LegacyEstimate
346 costComplete := quote.CostComplete
347 displayComplete := quote.DisplayComplete
348 rec.CostComplete = &costComplete
349 rec.DisplayComplete = &displayComplete
350 rec.DisplayStatus = quote.DisplayStatus
351 rec.AggregateMode = quote.AggregateMode
352 for _, total := range quote.OriginalTotals {
353 rec.OriginalTotals = append(rec.OriginalTotals, total.Currency+":"+total.Amount)
354 }
355 if quote.Selected != nil {
356 rec.SelectedAmount = quote.Selected.Amount
357 rec.SelectedCurrency = quote.Selected.Currency
358 rec.SelectedCost = quote.Selected.Float64()
359 }
360 if v, ok := quote.Valuations["CNY"]; ok {
361 rec.ValuationCNY = v.Money.Amount
362 }
363 if v, ok := quote.Valuations["USD"]; ok {
364 rec.ValuationUSD = v.Money.Amount
365 }
366 }
367 r.dispatcher.enqueue(rec)
368 }
369
370 func usageRequestCount(usage *provider.Usage) int {
371 if usage != nil && usage.RequestCount > 0 {
372 return usage.RequestCount
373 }
374 return 1
375 }
376
376 lines GO