Files
ModelRouter/internal/gateway/stats.go
JianFeeeee a2e1adc2d8 fix(gemini): endpoint 自相矛盾导致预置模板必失败
gemini.lua 里 adapter.endpoint = "/v1/models",而它自己的注释写的是
  POST /v1/models/{model}:generateContent
两者矛盾,而 Go 侧是静态拼接(provider.URL = base_url + endpoint),拼不出
模型名。预置模板 "Google Gemini"(base_url=.../v1beta)于是会 POST 到
  https://generativelanguage.googleapis.com/v1beta/v1/models
既多一段 /v1,又缺 :generateContent——那是 Gemini 的模型**列表**端点,
对 POST 返 405。所以任何用户从模板建这个源,拿到的都是必定失败的源。

实测确认影响范围:线上 21 个源里没有 gemini(openai×15 / trae / sensenova /
opencodezen / deepseek / anthropic / agentrouter),所以是潜伏缺陷。

修法:endpoint 改成模板 `/v1beta/models/{model}:generateContent`,新增
provider.ChatURL(model, stream):
  - 用 **PathEscape** 替换 {model}——模型 id 进的是 URL 路径,不转义的话
    一个 "/" 就会静默指向另一个资源(判据里用 RequestURI 而非 URL.Path
    断言,因为后者是解码后的,看不出 %2F);
  - 流式把 ":generateContent" 换成 ":streamGenerateContent"(同一个路径、
    不同动词,也在路径里)。替换刻意只认这个精确后缀,免得别的适配器
    仅仅提到这个词就被改写;
  - source 自己设的 endpoint: 仍然优先,模板被整体跳过。
Chat / ChatStream / probeChat 三处调用点改为传本次请求真实的 model——AUTO
按槽位把 req.Model 钉死,所以 URL 必须跟随**请求**的模型,用源默认模型会让
多模型源每次都打同一个(还记到别的模型的账上)。

判定静态 endpoint 的其他 10 个适配器零影响(TestNonGeminiEndpointsAreUntouched)。

顺带:Stats 的 mutex 不是可重入的,导出方法自己加锁、*Locked 后缀要求调用
方持锁。持锁调导出方法会死锁——我的探针真卡死过一次(直到 10 分钟超时)。
补上 LOCKING 注释,并加判据把这条规则钉住(含一个 20 秒上限的行为判据,
让未来的重构撞死锁时快速失败而不是拖满整个套件)。
2026-10-01 23:37:17 +08:00

1197 lines
35 KiB
Go

package gateway
import (
"bufio"
"bytes"
"encoding/json"
"fmt"
"io"
"log"
"os"
"path/filepath"
"sort"
"strconv"
"strings"
"sync"
"time"
)
// Req is one recorded gateway request (audit trail + per-key/per-model stats).
type Req struct {
Time int64 `json:"time"` // unix milliseconds
Key string `json:"key"` // gateway key id
Type string `json:"type"` // chat | stream | image
Model string `json:"model"` // effective model used upstream
// Source provider/source name
Source string `json:"source"`
// Prompt prompt tokens
Prompt int64 `json:"prompt_tokens"`
// Compl completion tokens
Compl int64 `json:"completion_tokens"`
// LatMs total handling time ms
LatMs int64 `json:"latency_ms"`
// FirstByteMs time-to-first-byte for streaming (ms from request start
// to the first SSE chunk sent to the client); for non-streaming it
// equals LatMs. 0 when unmeasured (legacy records).
FirstByteMs int64 `json:"first_byte_ms,omitempty"`
// CacheHit / CacheMiss carry the upstream prompt-cache accounting
// (DeepSeek-style hit/miss tokens) when the upstream reports it.
// Both 0 = upstream gave no cache data.
CacheHit int64 `json:"cache_hit_tokens,omitempty"`
CacheMiss int64 `json:"cache_miss_tokens,omitempty"`
// CacheReported marks that the upstream usage reported cache
// accounting at all (even when the hit count is 0). The WebUI shows
// "0%" instead of "—" for such rows.
CacheReported bool `json:"cache_reported,omitempty"`
OK bool `json:"ok"`
// Status http status code
Status int `json:"status"`
// Err short error message
Err string `json:"error,omitempty"`
// ImageCount is how many images an image-generation request produced.
// Kept separate from Compl/Prompt: image generation has no token concept,
// so counting images as "completion tokens" would corrupt the token totals.
ImageCount int `json:"image_count,omitempty"`
}
// Stat aggregates counters for one dimension row.
type Stat struct {
Reqs int64 `json:"reqs"`
OK int64 `json:"ok"`
Err int64 `json:"err"`
Tokens int64 `json:"tokens"`
Prompt int64 `json:"prompt_tokens"`
Compl int64 `json:"completion_tokens"`
LatSum int64 `json:"latency_sum_ms"`
LatMax int64 `json:"latency_max_ms"`
FirstByteSum int64 `json:"first_byte_sum_ms,omitempty"`
}
type agrRow struct {
Name string `json:"name"`
Stat
}
// Stats collects per-key / per-model / per-source aggregates plus a bounded
// ring of raw request records, all guarded by one mutex.
//
// LOCKING: mu is a plain sync.Mutex and is NOT reentrant. The *Locked methods
// (aggregateLocked, addKeyTokenLocked, addKeyHourLocked, addKeyReqLocked,
// wantPinnedBuckets, rotateAuditLocked, …) assume the caller already holds it,
// while every other exported method takes it itself.
//
// Calling an exported method while already holding mu DEADLOCKS. This is not
// hypothetical: a test that did KeyWindowReqs under s.mu.Lock() hung until the
// 10-minute panic timeout. Always reach for the *Locked variant when the lock
// is already held, and prefer the exported method when it is not.
type Stats struct {
mu sync.Mutex
active int64
byKey map[string]*Stat
byModel map[string]*Stat
bySrc map[string]*Stat
byKeyModel map[string]map[string]*Stat
byKeySrc map[string]map[string]*Stat
byStatus map[int]*Stat // per http status code aggregates (incl. 402/400)
recs []Req
maxRecs int
auditPath string
replayPartial bool // aggregates built from a bounded audit tail
modelHour map[string]map[int64]int64 // model -> unix-hour bucket -> tokens
// keyModelHour buckets the same tokens as modelHour but keyed by
// (gateway key id, model) so per-key model quotas are isolated from each
// other. modelHour stays key-blind on purpose: it backs the AUTO slot
// quota, which is a gateway-wide limit on a slot, not a per-key one.
keyModelHour map[string]map[string]map[int64]int64 // key -> model -> unix-hour -> tokens
// keyHour buckets a key's total tokens per unix hour, backing the
// whole-key quota (all models of one key share one budget).
keyHour map[string]map[int64]int64
// keyReqHour buckets a key's request count per unix hour, backing the
// key-wide request quota. Counted for every request type (chat, stream,
// image) including failed ones, so a failing client cannot loop for free.
keyReqHour map[string]map[int64]int64
// pinnedWanted records which (key, model) pairs have a quota that pins a
// source, so the pinned bucket is maintained only where it is read.
pinnedWanted map[string]bool
}
const hourSec = 3600
// auditRotateBytes rotates the audit file once it grows past this size (the
// file is renamed to <path>.<unix>.old and a fresh one is started); pruning
// keeps at most auditKeepOld rotated files. Both are vars so tests can shrink
// the threshold.
//
// The single-file threshold is deliberately small relative to the total budget
// (16 MB x 16 = 256 MB): records are paged by reading a file backwards, and a
// smaller file keeps the first page from seeking through a huge one.
var (
auditRotateBytes int64 = 16 << 20
auditKeepOld = 16
)
// defaultRingSize is how many recent request records stay resident. It has to
// cover two consumers: the status page's 5-minute source windows
// (SourceRecent/SourceAverages) and the first screen of the records table.
// Everything beyond that is paged from the audit file.
const defaultRingSize = 500
// firstScreenRecords is how many records a dashboard load ships by default — one
// screen's worth. Scrolling pulls the rest through /api/stats/records.
const firstScreenRecords = 100
func NewStats(maxRecords int) *Stats {
if maxRecords <= 0 {
maxRecords = defaultRingSize
}
return &Stats{
byKey: map[string]*Stat{},
byModel: map[string]*Stat{},
bySrc: map[string]*Stat{},
byKeyModel: map[string]map[string]*Stat{},
byKeySrc: map[string]map[string]*Stat{},
byStatus: map[int]*Stat{},
modelHour: map[string]map[int64]int64{},
keyModelHour: map[string]map[string]map[int64]int64{},
keyHour: map[string]map[int64]int64{},
keyReqHour: map[string]map[int64]int64{},
pinnedWanted: map[string]bool{},
maxRecs: maxRecords,
}
}
// Begin accounts an in-flight request; the returned func must be called once
// the request finished (defer ok).
func (s *Stats) Begin() func() {
s.mu.Lock()
s.active++
s.mu.Unlock()
return func() {
s.mu.Lock()
s.active--
s.mu.Unlock()
}
}
func inc(m map[string]*Stat, name string, r Req) {
a := m[name]
if a == nil {
a = &Stat{}
m[name] = a
}
incStatus(a, name, r)
}
func incStatus(a *Stat, name string, r Req) {
a.Reqs++
if r.OK {
a.OK++
} else {
a.Err++
}
a.Tokens += r.Prompt + r.Compl
a.Prompt += r.Prompt
a.Compl += r.Compl
a.LatSum += r.LatMs
if r.LatMs > a.LatMax {
a.LatMax = r.LatMs
}
if r.FirstByteMs > 0 {
a.FirstByteSum += r.FirstByteMs
}
}
// LoadAudit primes the in-memory state from the audit file.
//
// The aggregates (totals, per-key/model/source rows, status counts and the
// hourly quota buckets) are built from the FULL history: every audit file is
// streamed oldest-first so the dashboard shows real all-time numbers rather than
// whatever happened to fit in a replay window. This is affordable because the
// scan keeps nothing per record — aggregate maps are keyed by key/model/source,
// so their size is bounded by cardinality, not by request count. Measured on the
// production host: 29 MB / 221k lines / 37k requests in ~260 ms.
//
// The raw-record ring is what stays bounded: only the newest maxRecs records are
// retained, and everything older is paged from disk on demand by AuditPage /
// StreamAuditRecords. The old behaviour — appending EVERY record into a slice
// and truncating at the end — is what cost ~25 MB of resident memory.
func (s *Stats) LoadAudit(path string) {
s.mu.Lock()
s.auditPath = path
s.mu.Unlock()
// Oldest-first: the hourly-bucket retention prunes relative to the newest
// hour seen so far, so replaying in chronological order keeps exactly the
// intended trailing window.
chain := s.auditChain()
files := make([]string, 0, len(chain))
for i := len(chain) - 1; i >= 0; i-- {
files = append(files, chain[i])
}
ring := newReqRing(s.ringSize())
scanned := 0
s.mu.Lock()
for _, p := range files {
n, err := scanAuditFile(p, func(r Req) {
s.aggregateLocked(r)
ring.push(r)
})
scanned += n
if err != nil {
// A truncated or unreadable tail is not fatal: keep whatever was
// aggregated and mark the numbers as incomplete.
s.replayPartial = true
}
}
s.recs = ring.slice()
s.mu.Unlock()
if scanned > 0 {
log.Printf("[stats] replayed %d audit records from %d file(s) for aggregates; keeping the newest %d in memory",
scanned, len(files), len(s.recs))
}
}
func (s *Stats) ringSize() int {
s.mu.Lock()
defer s.mu.Unlock()
if s.maxRecs <= 0 {
return defaultRingSize
}
return s.maxRecs
}
// scanAuditFile streams one audit file, invoking fn for every request row, and
// returns how many request rows it saw. Access/event rows and malformed lines
// are skipped.
func scanAuditFile(path string, fn func(Req)) (int, error) {
f, err := os.Open(path)
if err != nil {
return 0, err
}
defer f.Close()
sc := bufio.NewScanner(f)
// tolerate long error summaries / oversized junk lines
sc.Buffer(make([]byte, 64*1024), 16*1024*1024)
n := 0
for sc.Scan() {
var r Req
if json.Unmarshal(sc.Bytes(), &r) != nil || r.Type == "" {
continue
}
fn(r)
n++
}
return n, sc.Err()
}
// reqRing keeps the newest n records seen, in chronological order, without
// growing with the number of records pushed through it.
type reqRing struct {
buf []Req
next int
full bool
limit int
}
func newReqRing(n int) *reqRing {
if n <= 0 {
n = defaultRingSize
}
return &reqRing{buf: make([]Req, n), limit: n}
}
func (r *reqRing) push(rec Req) {
r.buf[r.next] = rec
r.next++
if r.next == r.limit {
r.next = 0
r.full = true
}
}
// slice returns the retained records oldest-first.
func (r *reqRing) slice() []Req {
if !r.full {
out := make([]Req, r.next)
copy(out, r.buf[:r.next])
return out
}
out := make([]Req, 0, r.limit)
out = append(out, r.buf[r.next:]...)
out = append(out, r.buf[:r.next]...)
return out
}
// ReplayPartial reports whether the resident aggregates are known to be
// incomplete (an audit file could not be read in full). Under normal operation
// the aggregates cover the entire audit history, so this is false.
func (s *Stats) ReplayPartial() bool {
s.mu.Lock()
defer s.mu.Unlock()
return s.replayPartial
}
// Record appends a finished request to the aggregates and ring buffer.
func (s *Stats) Record(r Req) {
s.mu.Lock()
defer s.mu.Unlock()
s.aggregateLocked(r)
s.recs = append(s.recs, r)
if len(s.recs) > s.maxRecs {
s.recs = s.recs[len(s.recs)-s.maxRecs:]
}
if s.auditPath != "" {
s.rotateAuditLocked()
appendAuditLine(s.auditPath, r)
}
}
// aggregateLocked folds r into every aggregate row and the quota window
// bucket. Caller must hold s.mu.
func (s *Stats) aggregateLocked(r Req) {
if r.Key != "" {
s.addKeyReqLocked(r.Key, (r.Time/1000)/hourSec, 1)
}
inc(s.byKey, r.Key, r)
if r.Model != "" {
inc(s.byModel, r.Model, r)
}
inc(s.bySrc, r.Source, r)
km := s.byKeyModel[r.Key]
if km == nil {
km = map[string]*Stat{}
s.byKeyModel[r.Key] = km
}
if r.Model != "" {
inc(km, r.Model, r)
}
ks := s.byKeySrc[r.Key]
if ks == nil {
ks = map[string]*Stat{}
s.byKeySrc[r.Key] = ks
}
inc(ks, r.Source, r)
if r.Status > 0 {
a := s.byStatus[r.Status]
if a == nil {
a = &Stat{}
s.byStatus[r.Status] = a
}
incStatus(a, strconv.Itoa(r.Status), r)
}
// window bucket for quota enforcement (per source-model pair, per unix
// hour). r.Time is unix MILLISECONDS (audit format); hourSec is seconds,
// so convert before bucketing — otherwise the bucket width would be
// 3.6s and every WindowTokens cutoff comparison would be off by ~1000x.
tok := r.Prompt + r.Compl
if tok > 0 && r.Model != "" {
key := r.Model
if r.Source != "" {
key = r.Source + "::" + r.Model
}
h := (r.Time / 1000) / hourSec
hm := s.modelHour[key]
if hm == nil {
hm = map[int64]int64{}
s.modelHour[key] = hm
}
hm[h] += tok
// per-key buckets: same token split, but isolated per key so one
// key's quota cannot be drained by another key's traffic
if r.Key != "" {
s.addKeyTokenLocked(r.Key, r.Model, r.Source, h, tok)
s.addKeyHourLocked(r.Key, h, tok)
}
// retention: 24*40 = 960 hourly buckets ≈ 40 days of history (covers
// the longest "month" quota window)
if len(hm) > 24*40 {
for k := range hm {
if k < h-24*40 {
delete(hm, k)
}
}
}
}
}
// quotaRetentionHours is how much hourly history the per-key quota buckets
// keep. It matches the modelHour retention (40 days) so the longest "month"
// window is fully covered after a restart, and it is deliberately applied per
// key so an idle key's buckets are reclaimed instead of pinning memory.
const quotaRetentionHours = 24 * 40
// addKeyTokenLocked adds tok to one key's (model[, source]) hourly bucket.
//
// The bare-model bucket is always written: it is what a scope entry without a
// source pin reads, and writing it is what makes that quota work at all (every
// recorded request carries a resolved source, so a source::model-only bucket
// would leave an unpinned quota reading an empty bucket forever).
//
// The pinned "source::model" bucket is written only for (key, model) pairs
// some scope entry actually pins a source for. Maintaining it unconditionally
// doubles the bucket count — measured at 26.7 MB for 20 keys x 8 models x 3
// sources at full retention, against a documented ~32 MB total budget — to
// serve a lookup that almost nobody performs.
func (s *Stats) addKeyTokenLocked(key, model, source string, h, tok int64) {
if model == "" {
return
}
byModel := s.keyModelHour[key]
if byModel == nil {
byModel = map[string]map[int64]int64{}
s.keyModelHour[key] = byModel
}
addBucket(byModel, model, h, tok)
if source == "" {
return
}
pinned := source + "::" + model
if pinned == model || !s.pinnedWanted[key+"\x00"+model] {
return
}
addBucket(byModel, pinned, h, tok)
}
// addBucket adds tok to one named hourly bucket and enforces retention.
func addBucket(byModel map[string]map[int64]int64, name string, h, tok int64) {
hm := byModel[name]
if hm == nil {
hm = map[int64]int64{}
byModel[name] = hm
}
hm[h] += tok
if len(hm) > quotaRetentionHours {
for k := range hm {
if k < h-quotaRetentionHours {
delete(hm, k)
}
}
}
}
// wantPinnedBuckets starts maintaining the pinned (source::model) bucket for a
// key's model, so a quota that pins a source counts traffic from the moment it
// is configured. Usage recorded before that cannot be attributed to a source
// after the fact, so a pinned quota may under-count its first window.
func (s *Stats) wantPinnedBuckets(key, model string) {
if key == "" || model == "" {
return
}
s.pinnedWanted[key+"\x00"+model] = true
}
// addKeyHourLocked adds tok to one key's all-model hourly total.
func (s *Stats) addKeyHourLocked(key string, h, tok int64) {
hm := s.keyHour[key]
if hm == nil {
hm = map[int64]int64{}
s.keyHour[key] = hm
}
hm[h] += tok
if len(hm) > quotaRetentionHours {
for k := range hm {
if k < h-quotaRetentionHours {
delete(hm, k)
}
}
}
}
// KeyWindowReqs returns how many requests one key issued within sec seconds;
// sec <= 0 means all retained history.
func (s *Stats) KeyWindowReqs(key string, sec int64) int64 {
s.mu.Lock()
defer s.mu.Unlock()
return sumBuckets(s.keyReqHour[key], time.Now().Unix(), sec)
}
// addKeyReqLocked adds n requests to a key's hourly count bucket.
func (s *Stats) addKeyReqLocked(key string, h, n int64) {
hm := s.keyReqHour[key]
if hm == nil {
hm = map[int64]int64{}
s.keyReqHour[key] = hm
}
hm[h] += n
if len(hm) > quotaRetentionHours {
for k := range hm {
if k < h-quotaRetentionHours {
delete(hm, k)
}
}
}
}
// sumBuckets totals the hourly buckets inside the window ending at `now`
// (unix seconds). sec <= 0 means "all retained history" (no reset).
//
// It only visits the hour buckets a window can actually cover. The naive full
// scan measured 5.9us/op at the full 40-day retention, and a per-key quota
// check runs on every request (2-3 of these), so the scan cost lands directly
// on the request path; visiting at most sec/3600 + 1 buckets makes it
// proportional to the window instead of to the retained history. The map keys
// are whole unix hours, so the covered range is exact rather than a sample.
func sumBuckets(hm map[int64]int64, now, sec int64) int64 {
if len(hm) == 0 {
return 0
}
if sec <= 0 {
var total int64
for _, v := range hm {
total += v
}
return total
}
nowHour := now / hourSec
// A bucket at hour H covers [H*3600, (H+1)*3600). The pre-existing
// WindowTokens counts a bucket when h*3600 >= now-sec, so the first hour
// that qualifies is ceil((now-sec)/3600) — NOT floor, which would pull in
// one bucket too many (caught by TestSumBucketsMatchesFullScan: a 30-day
// window reported 8878 vs 8649).
firstHour := (now - sec + hourSec - 1) / hourSec
if firstHour < 0 {
firstHour = 0
}
var total int64
for h := firstHour; h <= nowHour; h++ {
total += hm[h]
}
return total
}
// KeyWindowModelTokens returns the tokens one key consumed for one model
// (optionally pinned to a source) within sec seconds; sec <= 0 means all
// retained history. Unlike WindowTokens this is isolated per key, which is
// what a per-key model quota needs.
//
// A source pin matches the pinned bucket exactly. Without a pin the call
// totals the key's bare-model bucket, which counts traffic on every source —
// a quota on "this model" should not stop applying just because the request
// happened to be served by a different upstream.
func (s *Stats) KeyWindowModelTokens(key, model, source string, sec int64) int64 {
s.mu.Lock()
defer s.mu.Unlock()
byModel := s.keyModelHour[key]
if len(byModel) == 0 {
return 0
}
mk := model
if source != "" {
mk = source + "::" + model
// from now on this pair keeps a pinned bucket, so a quota configured
// against this source counts the traffic it is meant to cap
s.wantPinnedBuckets(key, model)
}
return sumBuckets(byModel[mk], time.Now().Unix(), sec)
}
// KeyWindowTokens returns one key's total tokens across all models within sec
// seconds; sec <= 0 means all retained history. It backs the whole-key quota.
func (s *Stats) KeyWindowTokens(key string, sec int64) int64 {
s.mu.Lock()
defer s.mu.Unlock()
return sumBuckets(s.keyHour[key], time.Now().Unix(), sec)
}
// rotateAuditLocked renames the audit file to <path>.<unix>.old once it
// exceeds auditRotateBytes and prunes old files beyond auditKeepOld, keeping
// the newest ones. Caller must hold s.mu.
func (s *Stats) rotateAuditLocked() {
if s.auditPath == "" || auditRotateBytes <= 0 {
return
}
if fi, err := os.Stat(s.auditPath); err == nil && fi.Size() < auditRotateBytes {
return
}
ts := time.Now().Unix()
if os.Rename(s.auditPath, fmt.Sprintf("%s.%d.old", s.auditPath, ts)) == nil {
old, _ := filepath.Glob(s.auditPath + ".*.old")
sort.Sort(sort.Reverse(sort.StringSlice(old)))
for i := auditKeepOld; i < len(old); i++ {
_ = os.Remove(old[i])
}
}
}
func appendAuditLine(path string, row interface{}) {
f, err := os.OpenFile(path, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0644)
if err != nil {
return
}
defer f.Close()
if b, err := json.Marshal(row); err == nil {
_, _ = f.Write(append(b, '\n'))
}
}
// AppendAudit writes a generic event line (access log entry, login event,
// config change, …) to the same audit file without touching the aggregates.
func (s *Stats) AppendAudit(obj string, data map[string]interface{}) {
row := map[string]interface{}{"obj": obj, "time": time.Now().UnixMilli()}
for k, v := range data {
row[k] = v
}
s.mu.Lock()
defer s.mu.Unlock()
if s.auditPath == "" {
return
}
s.rotateAuditLocked()
appendAuditLine(s.auditPath, row)
}
// ModelTokens returns the tokens consumed per model for one gateway key id
// (used for per-model token quota enforcement).
func (s *Stats) ModelTokens(key string) map[string]int64 {
s.mu.Lock()
defer s.mu.Unlock()
out := map[string]int64{}
for k, v := range s.byKeyModel[key] {
out[k] = v.Tokens
}
return out
}
// KeyTokens returns the total tokens consumed by one gateway key id.
func (s *Stats) KeyTokens(key string) int64 {
s.mu.Lock()
defer s.mu.Unlock()
a := s.byKey[key]
if a == nil {
return 0
}
return a.Tokens
}
// AutoPeriodSeconds maps a quota reset period to its window length in
// seconds. "" → 0 (never resets); "hour" → 1h; "week" → 7d; "month" → 30d;
// "nhour" → Hours (>=1) hours.
func AutoPeriodSeconds(period string, hours int64) int64 {
switch period {
case "hour":
return hourSec
case "week":
return 7 * 24 * hourSec
case "month":
return 30 * 24 * hourSec
case "nhour":
if hours < 1 {
hours = 1
}
return hours * hourSec
}
return 0
}
// AutoSecondsToReset returns how many seconds remain until the quota window
// rolls over, for the Retry-After header on a 429. Buckets are whole unix
// hours, so the value is rounded up to the next hour boundary and never
// exceeds one window. 0 means "unknown" (never resets, or no period set), in
// which case the caller should not advertise a retry time.
func AutoSecondsToReset(period string, hours int64) int64 {
sec := AutoPeriodSeconds(period, hours)
if sec <= 0 {
return 0
}
now := time.Now().Unix()
// window = the last sec seconds, bucketed by whole hours; the oldest
// bucket still inside the window expires at its own hour boundary
elapsed := now % hourSec
left := sec - elapsed
if left <= 0 || left > sec {
left = sec
}
return left
}
// WindowTokens returns the tokens consumed for one model (optionally pinned
// to a single source) within the window; sec <= 0 means all time. Buckets are
// whole unix hours, so a sliding window overcounts by up to one hour — an
// accepted truncation for quota enforcement.
func (s *Stats) WindowTokens(model, source string, sec int64) int64 {
s.mu.Lock()
defer s.mu.Unlock()
key := model
if source != "" {
key = source + "::" + model
}
now := time.Now().Unix()
hm := s.modelHour[key]
if len(hm) == 0 {
return 0
}
var total int64
if sec <= 0 {
for _, v := range hm {
total += v
}
return total
}
cut := now - sec
for h, v := range hm {
if h*hourSec >= cut {
total += v
}
}
return total
}
func rows(m map[string]*Stat) []StatsRow {
out := make([]StatsRow, 0, len(m))
for k, v := range m {
out = append(out, StatsRow{Name: k, Stat: *v})
}
for i := 1; i < len(out); i++ {
for j := i; j > 0 && out[j].Reqs > out[j-1].Reqs; j-- {
out[j], out[j-1] = out[j-1], out[j]
}
}
return out
}
// StatsRow is one aggregated row for the dashboard.
type StatsRow struct {
Name string `json:"name"`
Stat
}
// AuditRecords returns every request row from the audit file plus its rotated
// .old files within the [from,to] unix-millisecond window, optionally for one
// masked key id, sorted oldest-first.
//
// It materializes the whole window, so it is only for tests and small windows.
// Production paths must use StreamAuditRecords (O(1) memory, for CSV export) or
// AuditPage (bounded pages, for the records view).
func (s *Stats) AuditRecords(from, to int64, key string) []Req {
var out []Req
_ = s.StreamAuditRecords(from, to, key, func(r Req) error {
out = append(out, r)
return nil
})
sort.Slice(out, func(i, j int) bool { return out[i].Time < out[j].Time })
return out
}
// ---- on-demand audit paging ----
//
// The records view is paged straight off disk instead of being held in memory:
// the newest page is read from the tail of the newest audit file backwards, and
// each response hands back a cursor for the next page. Nothing is cached between
// requests, so "release when the user leaves the page" is guaranteed by never
// retaining anything in the first place.
// auditCursor points at a byte offset inside one audit file. Paging walks
// backwards, so the cursor means "continue reading BEFORE this offset".
type auditCursor struct {
File string
Offset int64
}
// String encodes the cursor for the wire as "<file>:<offset>".
func (c auditCursor) String() string {
if c.File == "" {
return ""
}
return c.File + ":" + strconv.FormatInt(c.Offset, 10)
}
// parseAuditCursor decodes a wire cursor. An empty or malformed value means
// "start at the newest record".
func parseAuditCursor(s string) (auditCursor, bool) {
if s == "" {
return auditCursor{}, false
}
i := strings.LastIndex(s, ":")
if i <= 0 || i == len(s)-1 {
return auditCursor{}, false
}
off, err := strconv.ParseInt(s[i+1:], 10, 64)
if err != nil || off < 0 {
return auditCursor{}, false
}
return auditCursor{File: s[:i], Offset: off}, true
}
// auditChain lists the audit files newest-first: the live file, then the rotated
// ones in descending timestamp order.
func (s *Stats) auditChain() []string {
s.mu.Lock()
path := s.auditPath
s.mu.Unlock()
if path == "" {
return nil
}
files := []string{path}
if olds, err := filepath.Glob(path + ".*.old"); err == nil {
sort.Sort(sort.Reverse(sort.StringSlice(olds)))
files = append(files, olds...)
}
return files
}
// AuditPageResult is one page of records plus the cursor to continue from.
type AuditPageResult struct {
Records []Req `json:"records"`
Next string `json:"next_cursor,omitempty"`
HasMore bool `json:"has_more"`
// Rotated marks that the requested cursor's file no longer exists (the
// audit log rotated under the reader), so the client should restart from
// the newest page instead of silently skipping records.
Rotated bool `json:"rotated,omitempty"`
}
// revChunk is how much is read per backwards seek. Lines are far shorter than
// this, so a page is typically satisfied by one or two chunks.
const revChunk = 64 * 1024
// AuditPage returns up to limit request records ending at cursor, newest first.
// It reads the audit files backwards and stops as soon as the page is full, so
// cost and memory are proportional to limit rather than to the file size. The
// file handle is closed before returning: no state is kept between calls.
func (s *Stats) AuditPage(cursor string, limit int, key string) AuditPageResult {
if limit <= 0 || limit > 1000 {
limit = 100
}
chain := s.auditChain()
if len(chain) == 0 {
return AuditPageResult{Records: []Req{}}
}
start := 0
var offset int64 = -1 // -1 = start at EOF
rotated := false
if c, ok := parseAuditCursor(cursor); ok {
idx := -1
for i, f := range chain {
if f == c.File {
idx = i
break
}
}
if idx < 0 {
// the cursor's file rotated away: restart from the newest page and
// tell the client so it can reset its view
rotated = true
} else {
start, offset = idx, c.Offset
}
}
out := make([]Req, 0, limit)
for i := start; i < len(chain); i++ {
next, err := s.pageFile(chain[i], offset, limit, key, &out)
offset = -1 // subsequent files always start at their EOF
if err != nil {
continue
}
if len(out) >= limit {
res := AuditPageResult{Records: out, Rotated: rotated}
if next > 0 {
res.Next = auditCursor{File: chain[i], Offset: next}.String()
res.HasMore = true
} else if i+1 < len(chain) {
res.Next = auditCursor{File: chain[i+1], Offset: fileSize(chain[i+1])}.String()
res.HasMore = true
}
return res
}
}
return AuditPageResult{Records: out, HasMore: false, Rotated: rotated}
}
func fileSize(path string) int64 {
fi, err := os.Stat(path)
if err != nil {
return 0
}
return fi.Size()
}
// pageFile reads one audit file backwards from end (or EOF when end < 0),
// appending matching records to out (newest first) until it holds limit entries.
// It returns the offset to continue from within this file, or 0 when the file is
// exhausted.
func (s *Stats) pageFile(path string, end int64, limit int, key string, out *[]Req) (int64, error) {
f, err := os.Open(path)
if err != nil {
return 0, err
}
defer f.Close() // released before the response is written: nothing is retained
fi, err := f.Stat()
if err != nil {
return 0, err
}
if end < 0 || end > fi.Size() {
end = fi.Size()
}
buf := make([]byte, 0, revChunk)
var tail []byte // bytes of a line whose start was not in this chunk
pos := end
for pos > 0 && len(*out) < limit {
size := int64(revChunk)
if size > pos {
size = pos
}
pos -= size
buf = buf[:size]
if _, err := f.ReadAt(buf, pos); err != nil && err != io.EOF {
return 0, err
}
chunk := buf
if len(tail) > 0 {
chunk = append(append(make([]byte, 0, len(buf)+len(tail)), buf...), tail...)
tail = tail[:0]
}
// walk the chunk's complete lines from the end towards the start
lineEnd := len(chunk)
for i := len(chunk) - 1; i >= 0; i-- {
if chunk[i] != '\n' {
continue
}
if rec, ok := decodeAuditLine(chunk[i+1:lineEnd], key); ok {
*out = append(*out, rec)
if len(*out) >= limit {
// resume at the newline we just consumed
return pos + int64(i) + 1, nil
}
}
lineEnd = i
}
// chunk[:lineEnd] is an incomplete line: carry it into the next chunk
if lineEnd > 0 {
tail = append(tail[:0], chunk[:lineEnd]...)
}
if pos == 0 {
// start of file: whatever is carried is a complete first line
if rec, ok := decodeAuditLine(tail, key); ok {
*out = append(*out, rec)
}
}
}
if pos <= 0 {
return 0, nil // file exhausted
}
return pos, nil
}
// decodeAuditLine parses one audit line as a request record, filtering by key.
// Access/event rows and malformed lines are rejected.
func decodeAuditLine(line []byte, key string) (Req, bool) {
line = bytes.TrimSpace(line)
if len(line) == 0 {
return Req{}, false
}
var r Req
if json.Unmarshal(line, &r) != nil || r.Type == "" {
return Req{}, false
}
if key != "" && r.Key != key {
return Req{}, false
}
return r, true
}
// StreamAuditRecords walks every request record in [from, to] (unix millis,
// 0 = unbounded) newest file first and hands each one to fn. Nothing is
// accumulated: an export of an arbitrarily long period costs O(1) memory, and
// the file handles are closed as the walk proceeds. fn returning an error stops
// the walk (used to abort on a broken client connection).
func (s *Stats) StreamAuditRecords(from, to int64, key string, fn func(Req) error) error {
for _, path := range s.auditChain() {
if err := streamAuditFile(path, from, to, key, fn); err != nil {
return err
}
}
return nil
}
func streamAuditFile(path string, from, to int64, key string, fn func(Req) error) error {
f, err := os.Open(path)
if err != nil {
return nil // a rotated-away file is not an export failure
}
defer f.Close()
sc := bufio.NewScanner(f)
sc.Buffer(make([]byte, 64*1024), 16*1024*1024)
for sc.Scan() {
var r Req
if json.Unmarshal(sc.Bytes(), &r) != nil || r.Type == "" {
continue
}
if key != "" && r.Key != key {
continue
}
if from > 0 && r.Time < from {
continue
}
if to > 0 && r.Time > to {
continue
}
if err := fn(r); err != nil {
return err
}
}
return nil
}
// SourceRecent counts real gateway requests per source within the last
// window (unix seconds). It reads only the in-memory ring, so it is cheap and
// reflects live traffic — used by /api/status so the source status column is
// driven by what actually happens, not just a probe.
func (s *Stats) SourceRecent(windowSec int64) map[string][2]int64 {
s.mu.Lock()
defer s.mu.Unlock()
cut := time.Now().Unix() - windowSec
out := map[string][2]int64{}
for _, r := range s.recs {
if r.Time/1000 < cut {
continue
}
if r.Source == "" {
continue
}
v := out[r.Source]
if r.OK {
v[0]++
} else {
v[1]++
}
out[r.Source] = v
}
return out
}
// SourceAvg carries per-source performance averages for the status page.
type SourceAvg struct {
// AvgFirstByteMs is the mean time-to-first-byte over successful
// requests in the window (0 when no measured samples).
AvgFirstByteMs int64 `json:"avg_first_byte_ms"`
// AvgTokPerS is the aggregate completion throughput:
// sum(completion_tokens) / sum(latency_seconds) (0 when no samples).
AvgTokPerS int64 `json:"avg_tok_per_s"`
// Samples is the number of successful requests the averages cover.
Samples int64 `json:"samples"`
}
// SourceAverages computes TTFB and tokens/s averages per source from the
// in-memory ring within the window (unix seconds). Only successful chat/
// stream rows count; image and failed rows are skipped.
func (s *Stats) SourceAverages(windowSec int64) map[string]SourceAvg {
s.mu.Lock()
defer s.mu.Unlock()
cut := time.Now().Unix() - windowSec
type acc struct {
fbSum, latSum, complSum, n int64
}
accs := map[string]*acc{}
for _, r := range s.recs {
if !r.OK || r.Source == "" || r.Time/1000 < cut {
continue
}
if r.Type != "chat" && r.Type != "stream" {
continue
}
if r.LatMs <= 0 {
continue
}
a := accs[r.Source]
if a == nil {
a = &acc{}
accs[r.Source] = a
}
a.latSum += r.LatMs
a.complSum += r.Compl
a.fbSum += r.FirstByteMs
if r.FirstByteMs > 0 {
a.n++
}
}
out := make(map[string]SourceAvg, len(accs))
for src, a := range accs {
avg := SourceAvg{Samples: a.n}
if a.n > 0 {
avg.AvgFirstByteMs = a.fbSum / a.n
}
if a.latSum > 0 && a.complSum > 0 {
avg.AvgTokPerS = a.complSum * 1000 / a.latSum
}
out[src] = avg
}
return out
}
// Snapshot returns the whole dashboard payload; when key != "" the records
// and aggregate views are restricted to that gateway key.
// Snapshot returns the aggregate rows plus the FIRST SCREEN of request records
// (newest last, matching the ring order). limit bounds the record slice only;
// older records are not included here at all — the dashboard pages them from
// /api/stats/records as the user scrolls, so a dashboard load never serializes
// the whole ring. limit <= 0 means the default first screen.
func (s *Stats) Snapshot(limit int, key string) map[string]interface{} {
s.mu.Lock()
defer s.mu.Unlock()
if limit <= 0 {
limit = firstScreenRecords
}
if limit > s.maxRecs {
limit = s.maxRecs
}
start := 0
if len(s.recs) > limit {
start = len(s.recs) - limit
}
recs := s.recs[start:]
if key != "" {
filt := recs[:0]
for _, r := range recs {
if r.Key == key {
filt = append(filt, r)
}
}
recs = filt
}
var total Stat
var byModel, byKey, bySrc map[string]*Stat
if key == "" {
byKey, byModel, bySrc = s.byKey, s.byModel, s.bySrc
} else {
byKey = map[string]*Stat{key: s.byKey[key]}
byModel = s.byKeyModel[key]
bySrc = s.byKeySrc[key]
}
for _, a := range byModel {
total.Reqs += a.Reqs
total.OK += a.OK
total.Err += a.Err
total.Tokens += a.Tokens
total.Prompt += a.Prompt
total.Compl += a.Compl
total.LatSum += a.LatSum
if a.LatMax > total.LatMax {
total.LatMax = a.LatMax
}
}
bs := make([]agrRow, 0, len(s.byStatus))
for code := range s.byStatus {
bs = append(bs, agrRow{Name: strconv.Itoa(code), Stat: *s.byStatus[code]})
}
sort.Slice(bs, func(i, j int) bool {
ci, _ := strconv.Atoi(bs[i].Name)
cj, _ := strconv.Atoi(bs[j].Name)
return ci < cj
})
return map[string]interface{}{
"active": s.active,
"total": total,
"by_key": rows(byKey),
"by_model": rows(byModel),
"by_source": rows(bySrc),
"by_status": bs,
"records": append([]Req(nil), recs...),
"replay_partial": s.replayPartial,
}
}