| 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 |