Files
2026-09-11 06:14:38 +02:00

694 lines
19 KiB
Go

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
}