| """Token 用量存储与聚合(SQLite,单文件 logs/usage.db)。 |
| |
| 每次 /v1 请求拿到 usage 后记一行(时间戳 + 路径 + 模型 + token 数)。 |
| 面板按 24h / 1d / 3d / 7d / 30d 时间窗聚合:分桶时间序列 + 分模型汇总 + 总计。 |
| |
| - 纯标准库 sqlite3,无外部依赖;WAL 模式并发友好,进程内 threading 锁串行写。 |
| - 查询失败/库损坏不阻断主流程(记 usage 是旁路,绝不影响 API 响应)。 |
| """ |
| from __future__ import annotations |
|
|
| import sqlite3 |
| import threading |
| import time |
| from pathlib import Path |
| from typing import Any |
|
|
| _LOCK = threading.Lock() |
| _CONN: sqlite3.Connection | None = None |
| _DB_PATH: Path | None = None |
|
|
| |
| _WINDOWS: dict[str, tuple[int, int]] = { |
| "24h": (24 * 3600, 3600), |
| "1d": (24 * 3600, 3600), |
| "3d": (3 * 24 * 3600, 3 * 3600), |
| "7d": (7 * 24 * 3600, 6 * 3600), |
| "30d": (30 * 24 * 3600, 24 * 3600), |
| } |
|
|
|
|
| def init(db_path: str | Path) -> None: |
| """初始化数据库连接与表结构(应用启动时调一次)。""" |
| global _CONN, _DB_PATH |
| with _LOCK: |
| _DB_PATH = Path(db_path) |
| _DB_PATH.parent.mkdir(parents=True, exist_ok=True) |
| _CONN = sqlite3.connect(str(_DB_PATH), check_same_thread=False) |
| _CONN.execute("PRAGMA journal_mode=WAL") |
| _CONN.execute("PRAGMA synchronous=NORMAL") |
| _CONN.execute( |
| """CREATE TABLE IF NOT EXISTS usage ( |
| id INTEGER PRIMARY KEY AUTOINCREMENT, |
| ts REAL NOT NULL, |
| path TEXT NOT NULL DEFAULT '', |
| model TEXT NOT NULL DEFAULT '', |
| input_tokens INTEGER NOT NULL DEFAULT 0, |
| output_tokens INTEGER NOT NULL DEFAULT 0, |
| total_tokens INTEGER NOT NULL DEFAULT 0 |
| )""" |
| ) |
| _CONN.execute("CREATE INDEX IF NOT EXISTS idx_usage_ts ON usage(ts)") |
| _CONN.commit() |
|
|
|
|
| def _int(v: Any) -> int: |
| try: |
| return int(v or 0) |
| except (TypeError, ValueError): |
| return 0 |
|
|
|
|
| def record(usage: dict[str, Any] | None, *, path: str = "", model: str = "") -> None: |
| """记录一次用量(旁路,异常静默)。usage 支持 OpenAI/Anthropic 两种字段名。""" |
| if _CONN is None or not usage: |
| return |
| inp = _int(usage.get("input_tokens") or usage.get("prompt_tokens")) |
| out = _int(usage.get("output_tokens") or usage.get("completion_tokens")) |
| tot = _int(usage.get("total_tokens")) or (inp + out) |
| if not (inp or out or tot): |
| return |
| try: |
| with _LOCK: |
| _CONN.execute( |
| "INSERT INTO usage(ts, path, model, input_tokens, output_tokens, total_tokens) " |
| "VALUES(?,?,?,?,?,?)", |
| (time.time(), path or "", model or "", inp, out, tot), |
| ) |
| _CONN.commit() |
| except Exception: |
| pass |
|
|
|
|
| def aggregate(window: str = "24h") -> dict[str, Any]: |
| """按时间窗聚合:返回分桶时间序列 + 分模型汇总 + 总计。 |
| |
| 返回 ``{window, since, buckets:[{t, input, output, total, requests}], |
| by_model:[{model, input, output, total, requests}], total:{...}}``。 |
| """ |
| span, gran = _WINDOWS.get(window, _WINDOWS["24h"]) |
| now = time.time() |
| since = now - span |
| empty = { |
| "window": window, "since": since, "granularity": gran, |
| "buckets": [], "by_model": [], "requests": 0, |
| "total": {"input": 0, "output": 0, "total": 0, "requests": 0}, |
| } |
| if _CONN is None: |
| return empty |
| try: |
| with _LOCK: |
| rows = _CONN.execute( |
| "SELECT ts, model, input_tokens, output_tokens, total_tokens " |
| "FROM usage WHERE ts >= ? ORDER BY ts", |
| (since,), |
| ).fetchall() |
| except Exception: |
| return empty |
|
|
| |
| n_buckets = max(1, int(span // gran)) |
| start = now - n_buckets * gran |
| buckets = [ |
| {"t": start + i * gran, "input": 0, "output": 0, "total": 0, "requests": 0} |
| for i in range(n_buckets) |
| ] |
| by_model: dict[str, dict[str, int]] = {} |
| tot_in = tot_out = tot_tot = tot_req = 0 |
|
|
| for ts, model, inp, out, tot in rows: |
| idx = int((ts - start) // gran) |
| if 0 <= idx < n_buckets: |
| b = buckets[idx] |
| b["input"] += inp; b["output"] += out; b["total"] += tot; b["requests"] += 1 |
| m = by_model.setdefault(model or "(unknown)", |
| {"input": 0, "output": 0, "total": 0, "requests": 0}) |
| m["input"] += inp; m["output"] += out; m["total"] += tot; m["requests"] += 1 |
| tot_in += inp; tot_out += out; tot_tot += tot; tot_req += 1 |
|
|
| by_model_list = sorted( |
| [{"model": k, **v} for k, v in by_model.items()], |
| key=lambda x: x["total"], reverse=True, |
| ) |
| return { |
| "window": window, "since": since, "granularity": gran, |
| "buckets": buckets, "by_model": by_model_list, "requests": tot_req, |
| "total": {"input": tot_in, "output": tot_out, "total": tot_tot, "requests": tot_req}, |
| } |
|
|
|
|
| def purge_older_than(seconds: float) -> int: |
| """删除超过保留期的旧记录,返回删除行数(可选维护,默认不调用)。""" |
| if _CONN is None: |
| return 0 |
| try: |
| with _LOCK: |
| cur = _CONN.execute("DELETE FROM usage WHERE ts < ?", (time.time() - seconds,)) |
| _CONN.commit() |
| return cur.rowcount |
| except Exception: |
| return 0 |
|
|
|
|
| def close() -> None: |
| """关闭连接(测试/关停用)。""" |
| global _CONN |
| with _LOCK: |
| if _CONN is not None: |
| try: |
| _CONN.close() |
| except Exception: |
| pass |
| _CONN = None |
|
|