package usage import ( "bufio" "context" "encoding/json" "errors" "fmt" "os" "path/filepath" "sort" "strings" "time" ) // RetentionConfig controls the tiered usage history. DetailDays keeps raw // per-request JSONL journals. Older detail is converted to daily rollups. // DailyDays is the maximum age of daily rollups; older days are folded into // idempotent monthly rollups. MonthlyMonths=0 keeps monthly rollups forever. type RetentionConfig struct { DetailDays int DailyDays int MonthlyMonths int CompactionInterval time.Duration } type Aggregate struct { Requests uint64 `json:"requests"` Errors uint64 `json:"errors,omitempty"` PromptTokens int64 `json:"prompt_tokens"` CompletionTokens int64 `json:"completion_tokens"` CachedPromptTokens int64 `json:"cached_prompt_tokens,omitempty"` Credits float64 `json:"credits"` QueueMS int64 `json:"queue_ms"` ServiceMS int64 `json:"service_ms"` PromptEvalNS int64 `json:"prompt_eval_ns,omitempty"` EvalNS int64 `json:"eval_ns,omitempty"` BytesIn int64 `json:"bytes_in,omitempty"` BytesOut int64 `json:"bytes_out,omitempty"` LastRequest time.Time `json:"last_request,omitempty"` } func (a Aggregate) PromptTPS() float64 { if a.PromptEvalNS <= 0 { return 0 } return float64(a.PromptTokens) / (float64(a.PromptEvalNS) / 1e9) } func (a Aggregate) OutputTPS() float64 { if a.EvalNS <= 0 { return 0 } return float64(a.CompletionTokens) / (float64(a.EvalNS) / 1e9) } func (a *Aggregate) addEvent(e Event) { a.Requests++ if e.Status >= 400 { a.Errors++ } a.PromptTokens += e.Usage.PromptTokens a.CompletionTokens += e.Usage.CompletionTokens a.CachedPromptTokens += e.Usage.CachedPromptTokens a.Credits += e.ActualCredits a.QueueMS += e.QueueMS a.ServiceMS += e.ServiceMS a.PromptEvalNS += e.Usage.PromptEvalNS a.EvalNS += e.Usage.EvalNS a.BytesIn += e.BytesIn a.BytesOut += e.BytesOut if e.Time.After(a.LastRequest) { a.LastRequest = e.Time } } func (a *Aggregate) merge(b Aggregate) { a.Requests += b.Requests a.Errors += b.Errors a.PromptTokens += b.PromptTokens a.CompletionTokens += b.CompletionTokens a.CachedPromptTokens += b.CachedPromptTokens a.Credits += b.Credits a.QueueMS += b.QueueMS a.ServiceMS += b.ServiceMS a.PromptEvalNS += b.PromptEvalNS a.EvalNS += b.EvalNS a.BytesIn += b.BytesIn a.BytesOut += b.BytesOut if b.LastRequest.After(a.LastRequest) { a.LastRequest = b.LastRequest } } type RollupData struct { Global Aggregate `json:"global"` Tenants map[string]Aggregate `json:"tenants,omitempty"` Actors map[string]Aggregate `json:"actors,omitempty"` Applications map[string]Aggregate `json:"applications,omitempty"` Models map[string]Aggregate `json:"models,omitempty"` Workers map[string]Aggregate `json:"workers,omitempty"` } func newRollupData() RollupData { return RollupData{Tenants: map[string]Aggregate{}, Actors: map[string]Aggregate{}, Applications: map[string]Aggregate{}, Models: map[string]Aggregate{}, Workers: map[string]Aggregate{}} } func addMap(m map[string]Aggregate, key string, e Event) { if key == "" { return } x := m[key] x.addEvent(e) m[key] = x } func mergeMap(dst map[string]Aggregate, src map[string]Aggregate) { for k, v := range src { x := dst[k] x.merge(v) dst[k] = x } } func (d *RollupData) addEvent(e Event) { d.Global.addEvent(e) addMap(d.Tenants, e.Tenant, e) actor := eventActor(e) if e.Tenant != "" && actor != "" { addMap(d.Actors, e.Tenant+"\x00"+actor, e) } addMap(d.Applications, e.Application, e) addMap(d.Models, e.Model, e) addMap(d.Workers, e.Worker, e) } func (d *RollupData) merge(o RollupData) { d.Global.merge(o.Global) mergeMap(d.Tenants, o.Tenants) mergeMap(d.Actors, o.Actors) mergeMap(d.Applications, o.Applications) mergeMap(d.Models, o.Models) mergeMap(d.Workers, o.Workers) } type DailyRollup struct { Version int `json:"version"` Granularity string `json:"granularity"` Day string `json:"day"` GeneratedAt time.Time `json:"generated_at"` Data RollupData `json:"data"` } type MonthlyRollup struct { Version int `json:"version"` Granularity string `json:"granularity"` Month string `json:"month"` GeneratedAt time.Time `json:"generated_at"` // Days makes the monthly update idempotent if the process crashes after // writing the month but before removing the source daily file. Days map[string]RollupData `json:"days"` } func (m MonthlyRollup) Total() RollupData { out := newRollupData() keys := make([]string, 0, len(m.Days)) for k := range m.Days { keys = append(keys, k) } sort.Strings(keys) for _, k := range keys { out.merge(m.Days[k]) } return out } type RollupPoint struct { Period string `json:"period"` Requests uint64 `json:"requests"` Errors uint64 `json:"errors"` PromptTokens int64 `json:"prompt_tokens"` CompletionTokens int64 `json:"completion_tokens"` Credits float64 `json:"credits"` QueueMS int64 `json:"queue_ms"` ServiceMS int64 `json:"service_ms"` PromptTPS float64 `json:"prompt_tps"` OutputTPS float64 `json:"output_tps"` } func point(period string, a Aggregate) RollupPoint { return RollupPoint{Period: period, Requests: a.Requests, Errors: a.Errors, PromptTokens: a.PromptTokens, CompletionTokens: a.CompletionTokens, Credits: a.Credits, QueueMS: a.QueueMS, ServiceMS: a.ServiceMS, PromptTPS: a.PromptTPS(), OutputTPS: a.OutputTPS()} } type RetentionStatus struct { Enabled bool `json:"enabled"` DetailDays int `json:"detail_days"` DailyDays int `json:"daily_days"` MonthlyMonths int `json:"monthly_months"` CompactionInterval string `json:"compaction_interval"` RawFiles int `json:"raw_files"` RawBytes int64 `json:"raw_bytes"` DailyFiles int `json:"daily_files"` DailyBytes int64 `json:"daily_bytes"` MonthlyFiles int `json:"monthly_files"` MonthlyBytes int64 `json:"monthly_bytes"` LastCompaction time.Time `json:"last_compaction,omitempty"` LastError string `json:"last_error,omitempty"` LastReclaimedBytes int64 `json:"last_reclaimed_bytes,omitempty"` LastRawCompacted int `json:"last_raw_compacted,omitempty"` LastDailyCompacted int `json:"last_daily_compacted,omitempty"` } func normalizeRetention(c RetentionConfig) RetentionConfig { if c.DetailDays <= 0 { c.DetailDays = 30 } if c.DailyDays <= 0 { c.DailyDays = 400 } if c.DailyDays < c.DetailDays { c.DailyDays = c.DetailDays } if c.CompactionInterval <= 0 { c.CompactionInterval = 6 * time.Hour } return c } func (r *Recorder) rollupDirs() (string, string) { return filepath.Join(r.dir, "rollups", "daily"), filepath.Join(r.dir, "rollups", "monthly") } func atomicJSON(path string, v any) error { if err := os.MkdirAll(filepath.Dir(path), 0750); err != nil { return err } b, err := json.MarshalIndent(v, "", " ") if err != nil { return err } b = append(b, '\n') f, err := os.CreateTemp(filepath.Dir(path), ".tmp-*") if err != nil { return err } tmp := f.Name() defer os.Remove(tmp) if err = f.Chmod(0640); err == nil { _, err = f.Write(b) } if err == nil { err = f.Sync() } cerr := f.Close() if err == nil { err = cerr } if err != nil { return err } if err = os.Rename(tmp, path); err != nil { return err } if d, e := os.Open(filepath.Dir(path)); e == nil { _ = d.Sync() _ = d.Close() } return nil } func readJSON(path string, v any) error { b, err := os.ReadFile(path) if err != nil { return err } return json.Unmarshal(b, v) } func parseDayFromRaw(path string) (time.Time, bool) { base := filepath.Base(path) if !strings.HasPrefix(base, "usage-") || !strings.HasSuffix(base, ".jsonl") { return time.Time{}, false } s := strings.TrimSuffix(strings.TrimPrefix(base, "usage-"), ".jsonl") t, err := time.Parse("2006-01-02", s) return t, err == nil } func parseDayFromDaily(path string) (time.Time, bool) { base := filepath.Base(path) if !strings.HasPrefix(base, "rollup-daily-") || !strings.HasSuffix(base, ".json") { return time.Time{}, false } s := strings.TrimSuffix(strings.TrimPrefix(base, "rollup-daily-"), ".json") t, err := time.Parse("2006-01-02", s) return t, err == nil } func cutoffDays(now time.Time, days int) time.Time { today := time.Date(now.UTC().Year(), now.UTC().Month(), now.UTC().Day(), 0, 0, 0, 0, time.UTC) return today.AddDate(0, 0, -days+1) } func cutoffMonths(now time.Time, months int) time.Time { n := now.UTC() first := time.Date(n.Year(), n.Month(), 1, 0, 0, 0, 0, time.UTC) return first.AddDate(0, -months+1, 0) } func aggregateRaw(path string) (DailyRollup, error) { day, ok := parseDayFromRaw(path) if !ok { return DailyRollup{}, fmt.Errorf("invalid journal filename %s", path) } d := newRollupData() f, err := os.Open(path) if err != nil { return DailyRollup{}, err } defer f.Close() sc := bufio.NewScanner(f) sc.Buffer(make([]byte, 64<<10), 2<<20) for sc.Scan() { var e Event if json.Unmarshal(sc.Bytes(), &e) == nil { d.addEvent(e) } } if err := sc.Err(); err != nil { return DailyRollup{}, err } return DailyRollup{Version: 1, Granularity: "daily", Day: day.Format("2006-01-02"), GeneratedAt: time.Now().UTC(), Data: d}, nil } func (r *Recorder) compactDisk(ctx context.Context, now time.Time) (RetentionStatus, error) { if r.dir == "" { return RetentionStatus{}, nil } r.compactMu.Lock() defer r.compactMu.Unlock() c := r.retention dailyDir, monthlyDir := r.rollupDirs() _ = os.MkdirAll(dailyDir, 0750) _ = os.MkdirAll(monthlyDir, 0750) before := dirSize(r.dir) stat := RetentionStatus{Enabled: true, DetailDays: c.DetailDays, DailyDays: c.DailyDays, MonthlyMonths: c.MonthlyMonths, CompactionInterval: c.CompactionInterval.String()} detailCut := cutoffDays(now, c.DetailDays) raws, _ := filepath.Glob(filepath.Join(r.dir, "usage-*.jsonl")) sort.Strings(raws) for _, p := range raws { select { case <-ctx.Done(): return stat, ctx.Err() default: } day, ok := parseDayFromRaw(p) if !ok || !day.Before(detailCut) { continue } roll, err := aggregateRaw(p) if err != nil { return stat, err } out := filepath.Join(dailyDir, "rollup-daily-"+roll.Day+".json") if err := atomicJSON(out, roll); err != nil { return stat, err } if err := os.Remove(p); err != nil { return stat, err } stat.LastRawCompacted++ } dailyCut := cutoffDays(now, c.DailyDays) ds, _ := filepath.Glob(filepath.Join(dailyDir, "rollup-daily-*.json")) sort.Strings(ds) for _, p := range ds { select { case <-ctx.Done(): return stat, ctx.Err() default: } day, ok := parseDayFromDaily(p) if !ok || !day.Before(dailyCut) { continue } var d DailyRollup if err := readJSON(p, &d); err != nil { return stat, err } month := day.Format("2006-01") mp := filepath.Join(monthlyDir, "rollup-monthly-"+month+".json") m := MonthlyRollup{Version: 1, Granularity: "monthly", Month: month, GeneratedAt: time.Now().UTC(), Days: map[string]RollupData{}} if err := readJSON(mp, &m); err != nil && !errors.Is(err, os.ErrNotExist) { return stat, err } if m.Days == nil { m.Days = map[string]RollupData{} } m.Version = 1 m.Granularity = "monthly" m.Month = month m.GeneratedAt = time.Now().UTC() m.Days[d.Day] = d.Data if err := atomicJSON(mp, m); err != nil { return stat, err } if err := os.Remove(p); err != nil { return stat, err } stat.LastDailyCompacted++ } if c.MonthlyMonths > 0 { cut := cutoffMonths(now, c.MonthlyMonths) ms, _ := filepath.Glob(filepath.Join(monthlyDir, "rollup-monthly-*.json")) for _, p := range ms { base := strings.TrimSuffix(strings.TrimPrefix(filepath.Base(p), "rollup-monthly-"), ".json") mt, err := time.Parse("2006-01", base) if err != nil || !mt.Before(cut) { continue } var expired MonthlyRollup _ = readJSON(p, &expired) if err := os.Remove(p); err != nil { return stat, err } if r.loaded && expired.Days != nil { r.removeAggregate(expired.Total()) } } } after := dirSize(r.dir) if before > after { stat.LastReclaimedBytes = before - after } stat.LastCompaction = time.Now().UTC() r.setRetentionStatus(stat) return stat, nil } func dirSize(dir string) int64 { var n int64 _ = filepath.Walk(dir, func(_ string, info os.FileInfo, err error) error { if err == nil && info.Mode().IsRegular() { n += info.Size() } return nil }) return n } func fileCounts(pattern string) (int, int64) { ps, _ := filepath.Glob(pattern) var n int64 for _, p := range ps { if s, e := os.Stat(p); e == nil && s.Mode().IsRegular() { n += s.Size() } } return len(ps), n } func (r *Recorder) setRetentionStatus(s RetentionStatus) { r.retentionMu.Lock() r.retentionStatus = s r.retentionMu.Unlock() } func (r *Recorder) RetentionStatus() RetentionStatus { r.retentionMu.RLock() s := r.retentionStatus r.retentionMu.RUnlock() if r.dir == "" { return s } d, m := r.rollupDirs() s.RawFiles, s.RawBytes = fileCounts(filepath.Join(r.dir, "usage-*.jsonl")) s.DailyFiles, s.DailyBytes = fileCounts(filepath.Join(d, "rollup-daily-*.json")) s.MonthlyFiles, s.MonthlyBytes = fileCounts(filepath.Join(m, "rollup-monthly-*.json")) return s } func (r *Recorder) Compact(ctx context.Context) (RetentionStatus, error) { if r == nil || r.dir == "" { return RetentionStatus{}, nil } if err := r.Flush(ctx); err != nil { return r.RetentionStatus(), err } stat, err := r.compactDisk(ctx, time.Now().UTC()) if err != nil { stat.LastError = err.Error() r.setRetentionStatus(stat) return stat, err } // Refresh only the rollup index used for historical charts. All-time // counters are already correct in memory and must not be applied twice. _ = r.reloadRollupIndex() return r.RetentionStatus(), nil } func (r *Recorder) retentionLoop() { defer r.wg.Done() t := time.NewTicker(r.retention.CompactionInterval) defer t.Stop() for { select { case <-r.closeCh: return case <-t.C: ctx, c := context.WithTimeout(context.Background(), 30*time.Minute) _, err := r.compactDisk(ctx, time.Now().UTC()) if err == nil { _ = r.reloadRollupIndex() } c() if err != nil { st := r.RetentionStatus() st.LastError = err.Error() r.setRetentionStatus(st) } } } } func (r *Recorder) loadRollupsForReplay() error { dailyDir, monthlyDir := r.rollupDirs() ms, _ := filepath.Glob(filepath.Join(monthlyDir, "rollup-monthly-*.json")) sort.Strings(ms) for _, p := range ms { var m MonthlyRollup if err := readJSON(p, &m); err != nil { return err } d := m.Total() r.applyAggregate(d) r.monthlyAgg[m.Month] = d } ds, _ := filepath.Glob(filepath.Join(dailyDir, "rollup-daily-*.json")) sort.Strings(ds) for _, p := range ds { var d DailyRollup if err := readJSON(p, &d); err != nil { return err } r.applyAggregate(d.Data) r.dailyAgg[d.Day] = d.Data } return nil } func (r *Recorder) reloadRollupIndex() error { daily := map[string]RollupData{} monthly := map[string]RollupData{} dd, md := r.rollupDirs() ms, _ := filepath.Glob(filepath.Join(md, "rollup-monthly-*.json")) for _, p := range ms { var m MonthlyRollup if readJSON(p, &m) == nil { monthly[m.Month] = m.Total() } } ds, _ := filepath.Glob(filepath.Join(dd, "rollup-daily-*.json")) for _, p := range ds { var d DailyRollup if readJSON(p, &d) == nil { daily[d.Day] = d.Data } } // Preserve in-memory raw/detail days; replace only days that are now on disk rollups. r.rollupMu.Lock() for day, v := range r.dailyAgg { if _, ok := daily[day]; !ok { if t, e := time.Parse("2006-01-02", day); e == nil && !t.Before(cutoffDays(time.Now().UTC(), r.retention.DetailDays)) { daily[day] = v } } } r.dailyAgg = daily r.monthlyAgg = monthly r.rollupMu.Unlock() return nil } func (r *Recorder) applyAggregate(d RollupData) { r.global = mergeSummary(r.global, d.Global) for k, a := range d.Tenants { r.byTenant[k] = mergeSummary(r.byTenant[k], a) } for k, a := range d.Actors { r.byActor[k] = mergeSummary(r.byActor[k], a) } } func (r *Recorder) removeAggregate(d RollupData) { r.mu.Lock() defer r.mu.Unlock() r.global = subtractSummary(r.global, d.Global) for k, a := range d.Tenants { x := subtractSummary(r.byTenant[k], a) if x.Requests == 0 { delete(r.byTenant, k) } else { r.byTenant[k] = x } } for k, a := range d.Actors { x := subtractSummary(r.byActor[k], a) if x.Requests == 0 { delete(r.byActor, k) } else { r.byActor[k] = x } } } func subtractSummary(s Summary, a Aggregate) Summary { if a.Requests >= s.Requests { s.Requests = 0 } else { s.Requests -= a.Requests } s.PromptTokens -= a.PromptTokens if s.PromptTokens < 0 { s.PromptTokens = 0 } s.CompletionTokens -= a.CompletionTokens if s.CompletionTokens < 0 { s.CompletionTokens = 0 } s.Credits -= a.Credits if s.Credits < 0 { s.Credits = 0 } s.QueueMS -= a.QueueMS if s.QueueMS < 0 { s.QueueMS = 0 } s.ServiceMS -= a.ServiceMS if s.ServiceMS < 0 { s.ServiceMS = 0 } if !a.LastRequest.IsZero() && a.LastRequest.Equal(s.LastRequest) { s.LastRequest = time.Time{} } return s } func mergeSummary(s Summary, a Aggregate) Summary { s.Requests += a.Requests s.PromptTokens += a.PromptTokens s.CompletionTokens += a.CompletionTokens s.Credits += a.Credits s.QueueMS += a.QueueMS s.ServiceMS += a.ServiceMS if a.LastRequest.After(s.LastRequest) { s.LastRequest = a.LastRequest } return s } func chooseAggregate(d RollupData, dimension, name string) Aggregate { switch dimension { case "tenant": return d.Tenants[name] case "actor": return d.Actors[name] case "application": return d.Applications[name] case "model": return d.Models[name] case "worker": return d.Workers[name] default: return d.Global } } func (r *Recorder) Series(granularity, dimension, name string, limit int) []RollupPoint { if limit <= 0 { limit = 90 } if limit > 2000 { limit = 2000 } r.rollupMu.RLock() defer r.rollupMu.RUnlock() src := r.dailyAgg if granularity == "monthly" { src = r.monthlyAgg } keys := make([]string, 0, len(src)) for k := range src { keys = append(keys, k) } sort.Strings(keys) if len(keys) > limit { keys = keys[len(keys)-limit:] } out := make([]RollupPoint, 0, len(keys)) for _, k := range keys { out = append(out, point(k, chooseAggregate(src[k], dimension, name))) } return out } // CompactAt is intentionally exported for deterministic retention tests. func (r *Recorder) CompactAt(ctx context.Context, now time.Time) (RetentionStatus, error) { if err := r.Flush(ctx); err != nil { return r.RetentionStatus(), err } s, e := r.compactDisk(ctx, now) if e == nil { _ = r.reloadRollupIndex() } return s, e }