gcli2api / internal /reqlog /reqlog.go
a3216's picture
chore: 同步到上游 1.12.0-panel + 凭证同步/独立启动器/保活
6d60378 verified
Raw History Blame Contribute Delete
16.6 kB
// Package reqlog 记录脱敏的请求级指标与可选 JSONL 归档。
//
// 内存指标有界保存最近 100 条并维护进程级计数;磁盘归档只写请求元数据,
// 不写提示词、响应正文、Authorization 或其它凭证。可选的调用来源(客户端 IP /
// User-Agent,见 Event.ClientIP/UserAgent)由 server 按配置开关决定是否填充。
// 归档队列满时丢弃并计数,不允许日志写盘阻塞模型请求。
package reqlog
import (
"bufio"
"crypto/rand"
"encoding/json"
"fmt"
"os"
"path/filepath"
"sort"
"strconv"
"strings"
"sync"
"sync/atomic"
"time"
)
const (
recentLimit = 100
defaultFileMax = int64(16 << 20)
defaultQueue = 1024
defaultReadMax = 1000
archiveFileGlob = "requests-*.jsonl"
)
// NewRequestID 生成不含用户信息的本地请求 ID。
func NewRequestID() string {
var b [8]byte
if _, err := rand.Read(b[:]); err == nil {
return fmt.Sprintf("req-%x", b[:])
}
return fmt.Sprintf("req-%d", time.Now().UnixNano())
}
const (
OutcomeSuccess = "success"
OutcomeHTTPError = "http_error"
OutcomeStreamError = "stream_error"
OutcomeInterrupted = "interrupted"
)
// Config 归档参数。Enabled=false 时仍保留内存指标。
type Config struct {
Dir string
Enabled bool
RetentionDays int
MaxBytes int64
FileMaxBytes int64
QueueSize int
}
// Event 是一条脱敏请求记录。Account 只保存“昵称(uid8)”标签,不保存完整 UID。
//
// ClientIP / UserAgent 是**调用来源**:面板「运行日志」用它回答"这条请求是谁打进来的"。
// 二者由 server 侧按 logging.request_client_info 开关决定是否填充(关掉即保持空串,
// 归档里不会出现来源字段)——来源信息比 token 计数敏感,运营可自行决定是否落盘。
// 仍然不写提示词、响应正文、Authorization 或其它凭证。
type Event struct {
Time time.Time `json:"time"`
RequestID string `json:"request_id"`
Path string `json:"path"`
Account string `json:"account,omitempty"`
Model string `json:"model,omitempty"`
Status int `json:"status"`
OK bool `json:"ok"`
Outcome string `json:"outcome"`
DurationMs int64 `json:"duration_ms"`
TTFBMs int64 `json:"ttfb_ms,omitempty"`
Attempts int `json:"attempts,omitempty"`
PromptTokens int64 `json:"prompt_tokens,omitempty"`
CompletionTokens int64 `json:"completion_tokens,omitempty"`
TotalTokens int64 `json:"total_tokens,omitempty"`
Credit float64 `json:"credit,omitempty"`
HasCredit bool `json:"credit_known"`
// CacheHitTokens / CacheMissTokens 上游前缀缓存命中/未命中 token(issue #92)。
// 上游未回该维度时两者皆零值省略;hit=0 + miss>0 即整段未命中。
CacheHitTokens int64 `json:"cache_hit_tokens,omitempty"`
CacheMissTokens int64 `json:"cache_miss_tokens,omitempty"`
ClientIP string `json:"client_ip,omitempty"`
UserAgent string `json:"user_agent,omitempty"`
}
// Filter 用于从归档中筛选最近记录。字符串字段一律「包含」匹配(大小写不敏感),
// 便于面板用一段 IP 前缀或 UA 片段捞请求;From/To 是闭区间(零值 = 该侧不设界),
// 供「今天 / 近 7 天 / 自定义区间」这类时间查询使用。
type Filter struct {
Outcome string
Account string
Model string
ClientIP string
UserAgent string
From time.Time
To time.Time
}
// ArchiveStats 归档存储状态。
type ArchiveStats struct {
Enabled bool `json:"enabled"`
Dir string `json:"dir,omitempty"`
Files int `json:"files"`
Bytes int64 `json:"bytes"`
DroppedWrites uint64 `json:"dropped_writes"`
LastError string `json:"last_error,omitempty"`
}
// Snapshot 一次面板读取的完整指标快照。
type Snapshot struct {
StartedAt time.Time `json:"started_at"`
Completed int64 `json:"completed"`
InFlight int64 `json:"in_flight"`
Succeeded int64 `json:"succeeded"`
Failed int64 `json:"failed"`
SuccessRate float64 `json:"success_rate"`
HTTPSuccessRate float64 `json:"http_success_rate"`
AvgDurationMs float64 `json:"avg_duration_ms"`
Recent []Event `json:"recent"`
Archive ArchiveStats `json:"archive"`
}
// Recorder 并发安全的有界请求指标与归档记录器。
type Recorder struct {
mu sync.Mutex
started time.Time
inFlight int64
completed int64
succeeded int64
httpSuccess int64
durationSum int64
recent []Event
archive *archiveWriter
}
// New 创建记录器;Dir 为空或 Enabled=false 时只启用内存指标。
func New(cfg Config) *Recorder {
if cfg.FileMaxBytes <= 0 {
cfg.FileMaxBytes = defaultFileMax
}
if cfg.QueueSize <= 0 {
cfg.QueueSize = defaultQueue
}
r := &Recorder{started: time.Now(), recent: make([]Event, 0, recentLimit)}
r.archive = newArchiveWriter(cfg)
return r
}
// Begin 标记一个请求进入处理。
func (r *Recorder) Begin() {
if r == nil {
return
}
r.mu.Lock()
r.inFlight++
r.mu.Unlock()
}
// Record 记录一个请求完成事件并写入归档队列。
func (r *Recorder) Record(e Event) {
if r == nil {
return
}
if e.Time.IsZero() {
e.Time = time.Now()
}
if e.Outcome == "" {
if e.Status == 200 && e.OK {
e.Outcome = OutcomeSuccess
} else {
e.Outcome = OutcomeHTTPError
}
}
r.mu.Lock()
if r.inFlight > 0 {
r.inFlight--
}
r.completed++
if e.OK {
r.succeeded++
}
if e.Status >= 200 && e.Status < 300 {
r.httpSuccess++
}
r.durationSum += e.DurationMs
r.recent = append([]Event{e}, r.recent...)
if len(r.recent) > recentLimit {
r.recent = r.recent[:recentLimit]
}
r.mu.Unlock()
if r.archive != nil {
r.archive.enqueue(e)
}
}
// Snapshot 返回进程内指标和归档状态。
func (r *Recorder) Snapshot() Snapshot {
if r == nil {
return Snapshot{}
}
r.mu.Lock()
s := Snapshot{
StartedAt: r.started,
Completed: r.completed,
InFlight: r.inFlight,
Succeeded: r.succeeded,
Recent: append([]Event(nil), r.recent...),
}
if r.completed > 0 {
s.Failed = r.completed - r.succeeded
s.SuccessRate = round1(float64(r.succeeded) / float64(r.completed) * 100)
s.HTTPSuccessRate = round1(float64(r.httpSuccess) / float64(r.completed) * 100)
s.AvgDurationMs = float64(r.durationSum) / float64(r.completed)
}
r.mu.Unlock()
if r.archive != nil {
s.Archive = r.archive.stats()
}
return s
}
// ReadArchive 返回最近的归档事件(按时间倒序)。limit<=0 时回落 200,最大 1000。
func (r *Recorder) ReadArchive(limit int, filter Filter) ([]Event, error) {
if r == nil || r.archive == nil {
return nil, nil
}
return r.archive.read(limit, filter)
}
// Close 刷盘并停止后台归档。
func (r *Recorder) Close() {
if r == nil || r.archive == nil {
return
}
r.archive.close()
}
func round1(v float64) float64 {
return float64(int(v*10+0.5)) / 10
}
type archiveWriter struct {
cfg Config
ch chan Event
stop chan struct{}
done chan struct{}
closeOnce sync.Once
dropped atomic.Uint64
lastErrMu sync.Mutex
lastErr string
file *os.File
buf *bufio.Writer
path string
day string
size int64
}
func newArchiveWriter(cfg Config) *archiveWriter {
if !cfg.Enabled || strings.TrimSpace(cfg.Dir) == "" {
return &archiveWriter{cfg: cfg}
}
if cfg.RetentionDays <= 0 {
cfg.RetentionDays = 7
}
if cfg.MaxBytes <= 0 {
cfg.MaxBytes = 100 << 20
}
if cfg.FileMaxBytes <= 0 {
cfg.FileMaxBytes = defaultFileMax
}
if cfg.QueueSize <= 0 {
cfg.QueueSize = defaultQueue
}
w := &archiveWriter{
cfg: cfg,
ch: make(chan Event, cfg.QueueSize),
stop: make(chan struct{}),
done: make(chan struct{}),
}
if err := os.MkdirAll(cfg.Dir, 0o700); err != nil {
w.setErr(err)
w.cfg.Enabled = false
close(w.done)
return w
}
go w.run()
return w
}
func (w *archiveWriter) enabled() bool {
return w != nil && w.cfg.Enabled && w.cfg.Dir != "" && w.done != nil
}
func (w *archiveWriter) enqueue(e Event) {
if !w.enabled() {
return
}
select {
case w.ch <- e:
default:
w.dropped.Add(1)
}
}
func (w *archiveWriter) run() {
defer close(w.done)
ticker := time.NewTicker(time.Second)
defer ticker.Stop()
for {
select {
case e := <-w.ch:
w.writeEvent(e)
case <-ticker.C:
w.flush()
w.prune()
case <-w.stop:
for {
select {
case e := <-w.ch:
w.writeEvent(e)
default:
w.flush()
w.closeFile()
return
}
}
}
}
}
func (w *archiveWriter) writeEvent(e Event) {
raw, err := json.Marshal(e)
if err != nil {
w.setErr(err)
return
}
now := e.Time
if now.IsZero() {
now = time.Now()
}
day := now.Format("2006-01-02")
if w.file == nil || w.day != day {
if err := w.openFile(now, false); err != nil {
w.setErr(err)
return
}
}
if w.size > 0 && w.size+int64(len(raw))+1 > w.cfg.FileMaxBytes {
if err := w.openFile(now, true); err != nil {
w.setErr(err)
return
}
}
if _, err := w.buf.Write(raw); err != nil {
w.setErr(err)
return
}
if err := w.buf.WriteByte('\n'); err != nil {
w.setErr(err)
return
}
w.size += int64(len(raw)) + 1
}
func (w *archiveWriter) openFile(now time.Time, rotate bool) error {
w.closeFile()
day := now.Format("2006-01-02")
if err := os.MkdirAll(w.cfg.Dir, 0o700); err != nil {
return err
}
path := filepath.Join(w.cfg.Dir, "requests-"+day+".jsonl")
if rotate {
path = nextArchivePath(w.cfg.Dir, day)
} else if latest := latestArchiveForDay(w.cfg.Dir, day); latest != "" {
path = filepath.Join(w.cfg.Dir, latest)
}
f, err := os.OpenFile(path, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0o600)
if err != nil {
return err
}
info, err := f.Stat()
if err != nil {
_ = f.Close()
return err
}
w.file = f
w.buf = bufio.NewWriterSize(f, 64<<10)
w.path = path
w.day = day
w.size = info.Size()
return nil
}
func latestArchiveForDay(dir, day string) string {
entries, err := os.ReadDir(dir)
if err != nil {
return ""
}
prefix := "requests-" + day
best := ""
bestIndex := -1
for _, entry := range entries {
if entry.IsDir() {
continue
}
name := entry.Name()
if !strings.HasPrefix(name, prefix) || !strings.HasSuffix(name, ".jsonl") {
continue
}
idx := 0
if name[len(prefix):len(name)-len(".jsonl")] != "" {
part := strings.TrimSuffix(strings.TrimPrefix(name[len(prefix):], "."), ".jsonl")
n, err := strconv.Atoi(part)
if err != nil {
continue
}
idx = n
}
if idx > bestIndex {
bestIndex = idx
best = name
}
}
if best == "" {
return ""
}
return best
}
func nextArchivePath(dir, day string) string {
base := filepath.Join(dir, "requests-"+day)
path := base + ".jsonl"
for i := 1; fileExists(path); i++ {
path = fmt.Sprintf("%s.%d.jsonl", base, i)
}
return path
}
func fileExists(path string) bool {
st, err := os.Stat(path)
return err == nil && !st.IsDir()
}
func (w *archiveWriter) flush() {
if w.buf == nil {
return
}
if err := w.buf.Flush(); err != nil {
w.setErr(err)
}
}
func (w *archiveWriter) closeFile() {
if w.buf != nil {
_ = w.buf.Flush()
}
if w.file != nil {
_ = w.file.Close()
}
w.file = nil
w.buf = nil
w.path = ""
w.day = ""
w.size = 0
}
func (w *archiveWriter) prune() {
if !w.enabled() {
return
}
entries, err := os.ReadDir(w.cfg.Dir)
if err != nil {
w.setErr(err)
return
}
type item struct {
path string
size int64
mtime time.Time
}
items := make([]item, 0, len(entries))
var total int64
cutoff := time.Now().AddDate(0, 0, -w.cfg.RetentionDays)
for _, entry := range entries {
if entry.IsDir() || !strings.HasPrefix(entry.Name(), "requests-") || !strings.HasSuffix(entry.Name(), ".jsonl") {
continue
}
path := filepath.Join(w.cfg.Dir, entry.Name())
info, err := entry.Info()
if err != nil {
continue
}
items = append(items, item{path: path, size: info.Size(), mtime: info.ModTime()})
total += info.Size()
}
sort.Slice(items, func(i, j int) bool {
if items[i].mtime.Equal(items[j].mtime) {
return items[i].path < items[j].path
}
return items[i].mtime.Before(items[j].mtime)
})
for _, it := range items {
if it.path != w.path && (it.mtime.Before(cutoff) || total > w.cfg.MaxBytes) {
if err := os.Remove(it.path); err == nil {
total -= it.size
}
}
}
}
func (w *archiveWriter) read(limit int, filter Filter) ([]Event, error) {
if !w.enabled() {
return nil, nil
}
if limit <= 0 {
limit = 200
}
if limit > defaultReadMax {
limit = defaultReadMax
}
entries, err := os.ReadDir(w.cfg.Dir)
if err != nil {
return nil, err
}
// 文件读取顺序无关紧要:结果一律按事件时间排序。不能依赖归档文件的 mtime 还原时间
// 顺序——同一秒内连续轮转写出的多个文件 mtime 经常完全相同(Linux 文件时间戳粒度粗),
// os.ReadDir 的字典序又会把装着最早事件的基准文件 requests-<day>.jsonl 排在
// requests-<day>.N.jsonl 之后;目录被整体拷贝 / 恢复备份后 mtime 更不可信。
var out []Event
for _, entry := range entries {
name := entry.Name()
if entry.IsDir() || !strings.HasPrefix(name, "requests-") || !strings.HasSuffix(name, ".jsonl") {
continue
}
rows, err := readFile(filepath.Join(w.cfg.Dir, name), filter)
if err != nil {
return out, err
}
out = append(out, rows...)
}
// 契约:按事件时间倒序返回最近 limit 条(同一时刻用 Stable 保留落盘先后)。
sort.SliceStable(out, func(i, j int) bool { return out[i].Time.After(out[j].Time) })
if len(out) > limit {
out = out[:limit]
}
return out, nil
}
func readFile(path string, filter Filter) ([]Event, error) {
f, err := os.Open(path)
if err != nil {
return nil, err
}
defer f.Close()
scanner := bufio.NewScanner(f)
scanner.Buffer(make([]byte, 64<<10), 1<<20)
var out []Event
for scanner.Scan() {
var e Event
if json.Unmarshal(scanner.Bytes(), &e) != nil || !filter.match(e) {
continue
}
out = append(out, e)
}
return out, scanner.Err()
}
func (f Filter) match(e Event) bool {
if f.Outcome != "" && e.Outcome != f.Outcome {
return false
}
if !containsFold(e.Account, f.Account) {
return false
}
if !containsFold(e.Model, f.Model) {
return false
}
if !containsFold(e.ClientIP, f.ClientIP) {
return false
}
if !containsFold(e.UserAgent, f.UserAgent) {
return false
}
if !f.From.IsZero() && e.Time.Before(f.From) {
return false
}
if !f.To.IsZero() && e.Time.After(f.To) {
return false
}
return true
}
// containsFold 大小写不敏感的子串匹配;needle 为空视为命中(不筛该字段)。
func containsFold(haystack, needle string) bool {
if needle == "" {
return true
}
return strings.Contains(strings.ToLower(haystack), strings.ToLower(needle))
}
func (w *archiveWriter) stats() ArchiveStats {
s := ArchiveStats{Enabled: w.enabled(), Dir: w.cfg.Dir, DroppedWrites: w.dropped.Load(), LastError: w.errString()}
if !s.Enabled {
return s
}
entries, err := os.ReadDir(w.cfg.Dir)
if err != nil {
s.LastError = err.Error()
return s
}
for _, entry := range entries {
if entry.IsDir() || !strings.HasPrefix(entry.Name(), "requests-") || !strings.HasSuffix(entry.Name(), ".jsonl") {
continue
}
if info, err := entry.Info(); err == nil {
s.Files++
s.Bytes += info.Size()
}
}
return s
}
func (w *archiveWriter) close() {
if !w.enabled() {
return
}
w.closeOnce.Do(func() {
close(w.stop)
<-w.done
})
}
func (w *archiveWriter) setErr(err error) {
if err == nil {
return
}
w.lastErrMu.Lock()
w.lastErr = err.Error()
w.lastErrMu.Unlock()
}
func (w *archiveWriter) errString() string {
w.lastErrMu.Lock()
defer w.lastErrMu.Unlock()
return w.lastErr
}