kingcap/openapi
0
1#!/usr/bin/env python32import base643import hashlib4import hmac5import http.cookies6import http.server7import json8import os9import re10import secrets11import sqlite312import sys13import threading14import time15import urllib.error16import urllib.request17 18import psycopg19from psycopg.types.json import Jsonb20 21DB_PATH = os.getenv("USAGE_DB_PATH", "/data/usage.sqlite")22PG_DSN = os.getenv("USAGE_SYNC_PG_DSN") or os.getenv("DATABASE_URL") or os.getenv("PGSTORE_DSN", "")23SCHEMA = os.getenv("USAGE_SYNC_SCHEMA", "cpa_manager_usage")24INTERVAL = int(os.getenv("USAGE_SYNC_INTERVAL_SEC", "3"))25BATCH = int(os.getenv("USAGE_SYNC_BATCH", "500"))26MANAGER_BASE = os.getenv("USAGE_SYNC_MANAGER_BASE", "http://127.0.0.1:18317").rstrip("/")27WAIT_SEC = int(os.getenv("USAGE_SYNC_WAIT_SEC", "300"))28FORCE_UPSTREAM_URL = os.getenv("USAGE_SYNC_FORCE_UPSTREAM_URL", "").strip()29FORCE_MANAGEMENT_KEY = os.getenv("USAGE_SYNC_FORCE_MANAGEMENT_KEY", "").strip()30FORCE_QUEUE = os.getenv("USAGE_SYNC_FORCE_QUEUE", "").strip()31FORCE_POP_SIDE = os.getenv("USAGE_SYNC_FORCE_POP_SIDE", "").strip()32RETENTION_DAYS = int(os.getenv("USAGE_SYNC_RETENTION_DAYS", "30"))33MAX_USAGE_EVENTS = int(os.getenv("USAGE_SYNC_MAX_USAGE_EVENTS", "50000"))34MAX_DEAD_LETTERS = int(os.getenv("USAGE_SYNC_MAX_DEAD_LETTERS", "5000"))35PRUNE_INTERVAL = int(os.getenv("USAGE_SYNC_PRUNE_INTERVAL_SEC", "3600"))36PRUNE_SQLITE = os.getenv("USAGE_SYNC_PRUNE_SQLITE", "true").strip().lower() not in (37 "0",38 "false",39 "no",40 "off",41)42VACUUM_SQLITE_ON_PRUNE = os.getenv("USAGE_SYNC_VACUUM_SQLITE_ON_PRUNE", "false").strip().lower() in (43 "1",44 "true",45 "yes",46 "on",47)48PURGE_ON_START = os.getenv("USAGE_SYNC_PURGE_ON_START", "").strip().lower() in (49 "1",50 "true",51 "yes",52 "on",53)54CONFIG_PATH = os.getenv("USAGE_SYNC_CONFIG_PATH", "/data/usage_sync_config.json")55CONFIG_HTTP_HOST = os.getenv("USAGE_SYNC_CONFIG_HTTP_HOST", "127.0.0.1")56CONFIG_HTTP_PORT = int(os.getenv("USAGE_SYNC_CONFIG_HTTP_PORT", "18318"))57SESSION_TTL_SEC = int(os.getenv("USAGE_SYNC_CONFIG_SESSION_TTL_SEC", "43200"))58SESSION_COOKIE = "usage_sync_session"59CPAMP_ADMIN_CREDENTIAL_KEY = "admin_credential_v1"60CPAMP_ADMIN_KEY_FILE = os.getenv("CPA_MANAGER_ADMIN_KEY_FILE", "/data/cpamp_admin.key")61ENV_CONFIG_DEFAULTS = {62 "sync_interval_sec": INTERVAL,63 "sync_batch": BATCH,64 "retention_days": RETENTION_DAYS,65 "max_usage_events": MAX_USAGE_EVENTS,66 "max_dead_letter_events": MAX_DEAD_LETTERS,67 "prune_interval_sec": PRUNE_INTERVAL,68 "prune_sqlite": PRUNE_SQLITE,69 "vacuum_sqlite_on_prune": VACUUM_SQLITE_ON_PRUNE,70}71CONFIG_INT_LIMITS = {72 "sync_interval_sec": (1, 3600),73 "sync_batch": (1, 10000),74 "retention_days": (0, 3650),75 "max_usage_events": (0, 10_000_000),76 "max_dead_letter_events": (0, 1_000_000),77 "prune_interval_sec": (0, 86400),78}79CONFIG_BOOL_KEYS = {"prune_sqlite", "vacuum_sqlite_on_prune"}80 81if not re.fullmatch(r"[A-Za-z_][A-Za-z0-9_]*", SCHEMA):82 raise SystemExit(f"invalid USAGE_SYNC_SCHEMA: {SCHEMA}")83 84QSCHEMA = SCHEMA85CONFIG_LOCK = threading.RLock()86SESSION_LOCK = threading.RLock()87DATA_LOCK = threading.RLock()88RUNTIME_CONFIG = dict(ENV_CONFIG_DEFAULTS)89SESSIONS: dict[str, float] = {}90 91 92def log(msg: str):93 print(f"[usage-sync] {msg}", flush=True)94 95 96def now_ms() -> int:97 return int(time.time() * 1000)98 99 100def json_dumps(obj) -> str:101 return json.dumps(obj, ensure_ascii=False, sort_keys=True, separators=(",", ":"))102 103 104def row_hash(obj) -> str:105 return hashlib.sha256(json_dumps(obj).encode("utf-8")).hexdigest()106 107 108def parse_bool_value(value) -> bool:109 if isinstance(value, bool):110 return value111 if isinstance(value, (int, float)):112 return bool(value)113 return str(value).strip().lower() in ("1", "true", "yes", "on")114 115 116def validate_config(payload: dict, base: dict | None = None) -> dict:117 cfg = dict(base or ENV_CONFIG_DEFAULTS)118 errors = []119 120 for key, (min_value, max_value) in CONFIG_INT_LIMITS.items():121 if key not in payload:122 continue123 try:124 value = int(payload[key])125 except Exception:126 errors.append(f"{key} 必须是整数")127 continue128 if value < min_value or value > max_value:129 errors.append(f"{key} 必须在 {min_value} 到 {max_value} 之间")130 continue131 cfg[key] = value132 133 for key in CONFIG_BOOL_KEYS:134 if key in payload:135 cfg[key] = parse_bool_value(payload[key])136 137 if errors:138 raise ValueError(";".join(errors))139 return cfg140 141 142def load_runtime_config() -> dict:143 cfg = dict(ENV_CONFIG_DEFAULTS)144 if os.path.exists(CONFIG_PATH):145 try:146 with open(CONFIG_PATH, "r", encoding="utf-8") as f:147 raw = json.load(f)148 if isinstance(raw, dict):149 cfg = validate_config(raw, base=cfg)150 else:151 log(f"{CONFIG_PATH} 内容不是 JSON 对象,使用环境变量默认值")152 except Exception as e:153 log(f"读取同步配置失败,使用环境变量默认值:{e}")154 155 with CONFIG_LOCK:156 RUNTIME_CONFIG.clear()157 RUNTIME_CONFIG.update(cfg)158 return dict(cfg)159 160 161def get_config() -> dict:162 with CONFIG_LOCK:163 return dict(RUNTIME_CONFIG)164 165 166def save_runtime_config(payload: dict) -> dict:167 with CONFIG_LOCK:168 cfg = validate_config(payload, base=RUNTIME_CONFIG)169 parent = os.path.dirname(CONFIG_PATH)170 if parent:171 os.makedirs(parent, exist_ok=True)172 tmp_path = f"{CONFIG_PATH}.tmp"173 with open(tmp_path, "w", encoding="utf-8") as f:174 json.dump(cfg, f, ensure_ascii=False, sort_keys=True, indent=2)175 f.write("\n")176 os.replace(tmp_path, CONFIG_PATH)177 RUNTIME_CONFIG.clear()178 RUNTIME_CONFIG.update(cfg)179 return dict(cfg)180 181 182def pg_connect():183 if not PG_DSN:184 return None185 return psycopg.connect(PG_DSN, autocommit=True, prepare_threshold=None)186 187 188def sqlite_connect():189 conn = sqlite3.connect(DB_PATH, timeout=30, check_same_thread=False)190 conn.row_factory = sqlite3.Row191 conn.execute("PRAGMA busy_timeout = 5000")192 conn.execute("PRAGMA journal_mode = WAL")193 return conn194 195 196def table_exists_sqlite(conn, table: str) -> bool:197 row = conn.execute(198 "SELECT 1 FROM sqlite_master WHERE type='table' AND name=?",199 (table,),200 ).fetchone()201 return row is not None202 203 204def get_sqlite_columns(conn, table: str):205 rows = conn.execute(f"PRAGMA table_info({table})").fetchall()206 return [r["name"] for r in rows]207 208 209def fetch_sqlite_rows(conn, sql: str, params=()):210 return [dict(r) for r in conn.execute(sql, params).fetchall()]211 212 213def sqlite_count(conn, table: str) -> int:214 row = conn.execute(f"SELECT COUNT(*) AS c FROM {table}").fetchone()215 return int(row["c"] if row else 0)216 217 218def pg_get_state(pg, key: str, default="0"):219 row = pg.execute(220 f"SELECT value FROM {QSCHEMA}.sync_state WHERE key=%s",221 (key,),222 ).fetchone()223 return row[0] if row else default224 225 226def pg_set_state(pg, key: str, value: str):227 pg.execute(228 f"""229 INSERT INTO {QSCHEMA}.sync_state (key, value, updated_at_ms)230 VALUES (%s, %s, %s)231 ON CONFLICT (key) DO UPDATE232 SET value = EXCLUDED.value,233 updated_at_ms = EXCLUDED.updated_at_ms234 """,235 (key, str(value), now_ms()),236 )237 238 239def safe_rowcount(cur) -> int:240 try:241 return max(int(cur.rowcount), 0)242 except Exception:243 return 0244 245 246def ensure_pg_schema(pg):247 pg.execute(f"CREATE SCHEMA IF NOT EXISTS {QSCHEMA}")248 249 pg.execute(250 f"""251 CREATE TABLE IF NOT EXISTS {QSCHEMA}.settings (252 key text PRIMARY KEY,253 value_text text,254 row_json jsonb NOT NULL,255 synced_at_ms bigint NOT NULL256 )257 """258 )259 260 pg.execute(261 f"""262 CREATE TABLE IF NOT EXISTS {QSCHEMA}.model_prices (263 model_key text PRIMARY KEY,264 row_json jsonb NOT NULL,265 synced_at_ms bigint NOT NULL266 )267 """268 )269 270 pg.execute(271 f"""272 CREATE TABLE IF NOT EXISTS {QSCHEMA}.api_key_aliases (273 api_key_hash text PRIMARY KEY,274 alias text NOT NULL,275 row_json jsonb NOT NULL,276 synced_at_ms bigint NOT NULL277 )278 """279 )280 281 pg.execute(282 f"""283 CREATE TABLE IF NOT EXISTS {QSCHEMA}.usage_events (284 event_hash text PRIMARY KEY,285 sqlite_id bigint,286 timestamp_ms bigint,287 model text,288 auth_index text,289 failed boolean,290 row_json jsonb NOT NULL,291 synced_at_ms bigint NOT NULL292 )293 """294 )295 pg.execute(296 f"CREATE INDEX IF NOT EXISTS idx_{QSCHEMA}_usage_events_ts ON {QSCHEMA}.usage_events(timestamp_ms DESC)"297 )298 299 pg.execute(300 f"""301 CREATE TABLE IF NOT EXISTS {QSCHEMA}.dead_letter_events (302 record_hash text PRIMARY KEY,303 sqlite_id bigint,304 created_at_ms bigint,305 row_json jsonb NOT NULL,306 synced_at_ms bigint NOT NULL307 )308 """309 )310 pg.execute(311 f"CREATE INDEX IF NOT EXISTS idx_{QSCHEMA}_dead_letter_events_created ON {QSCHEMA}.dead_letter_events(created_at_ms DESC)"312 )313 314 pg.execute(315 f"""316 CREATE TABLE IF NOT EXISTS {QSCHEMA}.sync_state (317 key text PRIMARY KEY,318 value text NOT NULL,319 updated_at_ms bigint NOT NULL320 )321 """322 )323 324 325def prune_pg_by_retention(pg, cfg: dict) -> tuple[int, int]:326 retention_days = int(cfg["retention_days"])327 if retention_days <= 0:328 return 0, 0329 330 cutoff_ms = now_ms() - retention_days * 24 * 60 * 60 * 1000331 usage_deleted = safe_rowcount(332 pg.execute(333 f"""334 DELETE FROM {QSCHEMA}.usage_events335 WHERE timestamp_ms IS NOT NULL336 AND timestamp_ms > 0337 AND timestamp_ms < %s338 """,339 (cutoff_ms,),340 )341 )342 dead_deleted = safe_rowcount(343 pg.execute(344 f"""345 DELETE FROM {QSCHEMA}.dead_letter_events346 WHERE created_at_ms IS NOT NULL347 AND created_at_ms > 0348 AND created_at_ms < %s349 """,350 (cutoff_ms,),351 )352 )353 return usage_deleted, dead_deleted354 355 356def prune_pg_by_count(pg, cfg: dict) -> tuple[int, int]:357 usage_deleted = 0358 dead_deleted = 0359 max_usage_events = int(cfg["max_usage_events"])360 max_dead_letter_events = int(cfg["max_dead_letter_events"])361 362 if max_usage_events > 0:363 usage_deleted = safe_rowcount(364 pg.execute(365 f"""366 DELETE FROM {QSCHEMA}.usage_events367 WHERE event_hash IN (368 SELECT event_hash369 FROM {QSCHEMA}.usage_events370 ORDER BY COALESCE(timestamp_ms, 0) DESC,371 COALESCE(sqlite_id, 0) DESC,372 event_hash DESC373 OFFSET %s374 )375 """,376 (max_usage_events,),377 )378 )379 380 if max_dead_letter_events > 0:381 dead_deleted = safe_rowcount(382 pg.execute(383 f"""384 DELETE FROM {QSCHEMA}.dead_letter_events385 WHERE record_hash IN (386 SELECT record_hash387 FROM {QSCHEMA}.dead_letter_events388 ORDER BY COALESCE(created_at_ms, 0) DESC,389 COALESCE(sqlite_id, 0) DESC,390 record_hash DESC391 OFFSET %s392 )393 """,394 (max_dead_letter_events,),395 )396 )397 398 return usage_deleted, dead_deleted399 400 401def sqlite_delete_old(conn, table: str, ts_col: str, cutoff_ms: int) -> int:402 cur = conn.execute(403 f"""404 DELETE FROM {table}405 WHERE {ts_col} IS NOT NULL406 AND {ts_col} > 0407 AND {ts_col} < ?408 """,409 (cutoff_ms,),410 )411 return safe_rowcount(cur)412 413 414def sqlite_delete_over_limit(conn, table: str, ts_col: str, limit: int) -> int:415 if limit <= 0:416 return 0417 418 cur = conn.execute(419 f"""420 DELETE FROM {table}421 WHERE id IN (422 SELECT id423 FROM {table}424 ORDER BY COALESCE({ts_col}, 0) DESC,425 id DESC426 LIMIT -1 OFFSET ?427 )428 """,429 (limit,),430 )431 return safe_rowcount(cur)432 433 434def prune_sqlite(conn, cfg: dict) -> tuple[int, int]:435 usage_deleted = 0436 dead_deleted = 0437 retention_days = int(cfg["retention_days"])438 max_usage_events = int(cfg["max_usage_events"])439 max_dead_letter_events = int(cfg["max_dead_letter_events"])440 441 if retention_days > 0:442 cutoff_ms = now_ms() - retention_days * 24 * 60 * 60 * 1000443 usage_deleted += sqlite_delete_old(conn, "usage_events", "timestamp_ms", cutoff_ms)444 dead_deleted += sqlite_delete_old(conn, "dead_letter_events", "created_at_ms", cutoff_ms)445 446 usage_deleted += sqlite_delete_over_limit(conn, "usage_events", "timestamp_ms", max_usage_events)447 dead_deleted += sqlite_delete_over_limit(conn, "dead_letter_events", "created_at_ms", max_dead_letter_events)448 449 conn.commit()450 451 if bool(cfg["vacuum_sqlite_on_prune"]) and (usage_deleted or dead_deleted):452 try:453 conn.execute("PRAGMA wal_checkpoint(TRUNCATE)")454 conn.execute("VACUUM")455 except Exception as e:456 log(f"SQLite VACUUM 失败,已忽略:{e}")457 458 return usage_deleted, dead_deleted459 460 461def prune_usage_data(pg, reason: str = "定时裁剪", include_sqlite: bool | None = None):462 cfg = get_config()463 if include_sqlite is None:464 include_sqlite = bool(cfg["prune_sqlite"])465 466 usage_deleted = 0467 dead_deleted = 0468 469 u, d = prune_pg_by_retention(pg, cfg)470 usage_deleted += u471 dead_deleted += d472 473 u, d = prune_pg_by_count(pg, cfg)474 usage_deleted += u475 dead_deleted += d476 477 sqlite_usage_deleted = 0478 sqlite_dead_deleted = 0479 if include_sqlite:480 conn = sqlite_connect()481 try:482 sqlite_usage_deleted, sqlite_dead_deleted = prune_sqlite(conn, cfg)483 finally:484 conn.close()485 486 if usage_deleted or dead_deleted or sqlite_usage_deleted or sqlite_dead_deleted:487 log(488 f"{reason}完成:PG usage_events={usage_deleted}, PG dead_letter_events={dead_deleted}, "489 f"SQLite usage_events={sqlite_usage_deleted}, SQLite dead_letter_events={sqlite_dead_deleted}"490 )491 492 493def purge_usage_data(pg, reason: str = "启动清理"):494 if not wait_sqlite_ready(WAIT_SEC):495 raise RuntimeError("已请求清空 usage 数据,但 SQLite 未就绪;为避免旧数据重新同步,停止启动")496 497 pg_usage_deleted = safe_rowcount(pg.execute(f"DELETE FROM {QSCHEMA}.usage_events"))498 pg_dead_deleted = safe_rowcount(pg.execute(f"DELETE FROM {QSCHEMA}.dead_letter_events"))499 500 conn = sqlite_connect()501 try:502 sqlite_usage_deleted = safe_rowcount(conn.execute("DELETE FROM usage_events"))503 sqlite_dead_deleted = safe_rowcount(conn.execute("DELETE FROM dead_letter_events"))504 conn.commit()505 try:506 conn.execute("PRAGMA wal_checkpoint(TRUNCATE)")507 conn.execute("VACUUM")508 except Exception as e:509 log(f"SQLite PURGE 后 VACUUM 失败,已忽略:{e}")510 511 max_usage = conn.execute("SELECT COALESCE(MAX(id), 0) FROM usage_events").fetchone()[0]512 max_dead = conn.execute("SELECT COALESCE(MAX(id), 0) FROM dead_letter_events").fetchone()[0]513 pg_set_state(pg, "sqlite_usage_last_id", str(max_usage))514 pg_set_state(pg, "sqlite_dead_last_id", str(max_dead))515 finally:516 conn.close()517 518 log(519 f"{reason}完成:"520 f"PG usage_events={pg_usage_deleted}, PG dead_letter_events={pg_dead_deleted}, "521 f"SQLite usage_events={sqlite_usage_deleted}, SQLite dead_letter_events={sqlite_dead_deleted}"522 )523 524 525def wait_http_ok(url: str, timeout_sec: int) -> bool:526 deadline = time.time() + timeout_sec527 while time.time() < deadline:528 try:529 with urllib.request.urlopen(url, timeout=5) as resp:530 if 200 <= resp.status < 300:531 return True532 except Exception:533 pass534 time.sleep(2)535 return False536 537 538def wait_sqlite_ready(timeout_sec: int):539 deadline = time.time() + timeout_sec540 while time.time() < deadline:541 try:542 if not os.path.exists(DB_PATH):543 time.sleep(2)544 continue545 conn = sqlite_connect()546 ok = all(547 table_exists_sqlite(conn, t)548 for t in ("settings", "model_prices", "usage_events", "dead_letter_events")549 )550 conn.close()551 if ok:552 return True553 except Exception:554 pass555 time.sleep(2)556 return False557 558 559def normalize_setup_payload(raw: str):560 try:561 data = json.loads(raw)562 except Exception:563 log("PG 中 settings.setup 不是合法 JSON,跳过 setup 恢复")564 return None565 566 base = (567 FORCE_UPSTREAM_URL568 or data.get("cpaBaseUrl")569 or data.get("CPAUpstreamURL")570 or data.get("upstream")571 or ""572 ).strip().rstrip("/")573 574 key = (575 FORCE_MANAGEMENT_KEY576 or data.get("managementKey")577 or data.get("ManagementKey")578 or ""579 ).strip()580 if key.startswith("enc:v1:"):581 log("PG 中 settings.setup 的 managementKey 已由 CPA-Manager-Plus 加密,跳过 HTTP setup 恢复")582 return None583 584 queue = (585 FORCE_QUEUE586 or data.get("queue")587 or data.get("Queue")588 or "usage"589 ).strip()590 591 pop_side = (592 FORCE_POP_SIDE593 or data.get("popSide")594 or data.get("PopSide")595 or "right"596 ).strip()597 598 if not base or not key:599 log("PG 中 settings.setup 缺少 cpaBaseUrl/managementKey,跳过 setup 恢复")600 return None601 602 return {603 "cpaBaseUrl": base,604 "managementKey": key,605 "cpaManagementKey": key,606 "requestMonitoringEnabled": True,607 "ensureUsageStatisticsEnabled": True,608 "queue": queue,609 "popSide": pop_side,610 "collectorMode": os.getenv("USAGE_COLLECTOR_MODE", "auto"),611 "batchSize": int(os.getenv("USAGE_BATCH_SIZE", "100")),612 "pollIntervalMs": int(os.getenv("USAGE_POLL_INTERVAL_MS", "500")),613 "queryLimit": int(os.getenv("USAGE_QUERY_LIMIT", "50000")),614 }615 616 617def get_setup_admin_key() -> str:618 key = os.getenv("CPA_MANAGER_ADMIN_KEY", "").strip()619 if key:620 return key621 try:622 with open(CPAMP_ADMIN_KEY_FILE, "r", encoding="utf-8") as f:623 return f.read().strip()624 except Exception:625 return ""626 627 628def restore_setup_via_http(pg):629 row = pg.execute(630 f"SELECT value_text FROM {QSCHEMA}.settings WHERE key='setup'"631 ).fetchone()632 if not row or not row[0]:633 log("PG 中还没有 setup,等待你首次在面板保存后再自动镜像")634 return False635 636 payload = normalize_setup_payload(row[0])637 if not payload:638 return False639 admin_key = get_setup_admin_key()640 if not admin_key:641 log("未找到 CPA-Manager-Plus 管理员密钥文件,跳过 HTTP setup 恢复;将依赖 SQLite settings 迁移")642 return False643 644 deadline = time.time() + WAIT_SEC645 body = json.dumps(payload).encode("utf-8")646 req = urllib.request.Request(647 f"{MANAGER_BASE}/setup",648 data=body,649 headers={650 "Content-Type": "application/json",651 "Authorization": f"Bearer {admin_key}",652 },653 method="POST",654 )655 656 while time.time() < deadline:657 try:658 with urllib.request.urlopen(req, timeout=10) as resp:659 if 200 <= resp.status < 300:660 log(f"已从 PG 恢复 CPA-Manager-Plus setup,upstream={payload['cpaBaseUrl']}")661 return True662 except urllib.error.HTTPError as e:663 try:664 detail = e.read().decode("utf-8", "ignore")665 except Exception:666 detail = str(e)667 log(f"/setup 返回 {e.code},稍后重试:{detail}")668 except Exception as e:669 log(f"/setup 暂时失败,稍后重试:{e}")670 time.sleep(2)671 672 log("setup 自动恢复超时;如果面板仍未接管统计,请手动打开 management.html 保存一次")673 return False674 675 676def iso_utc_from_ms(ms: int) -> str:677 try:678 return time.strftime("%Y-%m-%dT%H:%M:%S.000Z", time.gmtime(ms / 1000))679 except Exception:680 return time.strftime("%Y-%m-%dT%H:%M:%S.000Z", time.gmtime())681 682 683def normalize_row_for_sqlite(table: str, row: dict) -> dict:684 row = dict(row)685 current_ms = now_ms()686 687 if table == "settings":688 row.setdefault("updated_at_ms", current_ms)689 if row.get("value") is None and row.get("value_text") is not None:690 row["value"] = row.get("value_text")691 return row692 693 if table == "model_prices":694 row.setdefault("model", row.get("model_key") or row_hash(row))695 for key in (696 "prompt_per_1m",697 "completion_per_1m",698 "cache_per_1m",699 "cache_read_per_1m",700 "cache_creation_per_1m",701 ):702 row.setdefault(key, 0)703 for key in (704 "prompt_configured",705 "completion_configured",706 "cache_read_configured",707 "cache_creation_configured",708 ):709 row.setdefault(key, 0)710 row.setdefault("updated_at_ms", current_ms)711 return row712 713 if table == "usage_events":714 if not row.get("event_hash"):715 row["event_hash"] = row_hash(row)716 timestamp_ms = int(row.get("timestamp_ms") or row.get("created_at_ms") or current_ms)717 row["timestamp_ms"] = timestamp_ms718 row.setdefault("timestamp", iso_utc_from_ms(timestamp_ms))719 if not row.get("model"):720 row["model"] = row.get("resolved_model") or row.get("requested_model") or "unknown"721 for key in (722 "input_tokens",723 "output_tokens",724 "reasoning_tokens",725 "cached_tokens",726 "cache_tokens",727 "cache_read_tokens",728 "cache_creation_tokens",729 "total_tokens",730 ):731 row.setdefault(key, 0)732 row["failed"] = 1 if row.get("failed") else 0733 row.setdefault("created_at_ms", timestamp_ms)734 return row735 736 if table == "dead_letter_events":737 row.setdefault("payload", row.get("raw_json") or "{}")738 row.setdefault("error", "")739 row.setdefault("created_at_ms", current_ms)740 return row741 742 if table == "api_key_aliases":743 row.setdefault("api_key_hash", row_hash(row))744 row.setdefault("alias", row.get("api_key_hash", ""))745 row.setdefault("updated_at_ms", current_ms)746 return row747 748 return row749 750 751def sqlite_insert(conn, table: str, row: dict, mode: str):752 row = normalize_row_for_sqlite(table, row)753 cols = get_sqlite_columns(conn, table)754 use_cols = [c for c in cols if c in row]755 if not use_cols:756 return757 placeholders = ",".join("?" for _ in use_cols)758 col_sql = ",".join(use_cols)759 values = [row[c] for c in use_cols]760 conn.execute(761 f"INSERT OR {mode} INTO {table} ({col_sql}) VALUES ({placeholders})",762 values,763 )764 765 766def restore_settings_to_sqlite(pg, conn):767 rows = pg.execute(768 f"SELECT row_json FROM {QSCHEMA}.settings"769 ).fetchall()770 for (row_json,) in rows:771 sqlite_insert(conn, "settings", dict(row_json), "REPLACE")772 conn.commit()773 774 775def restore_table_if_empty(pg, conn, table: str, order_sql: str):776 if sqlite_count(conn, table) > 0:777 return778 779 batch = int(get_config()["sync_batch"])780 total = 0781 while True:782 rows = pg.execute(783 f"SELECT row_json FROM {QSCHEMA}.{table} {order_sql} LIMIT %s OFFSET %s",784 (batch, total),785 ).fetchall()786 if not rows:787 break788 789 for (row_json,) in rows:790 mode = "REPLACE" if table in ("settings", "model_prices", "api_key_aliases") else "IGNORE"791 sqlite_insert(conn, table, dict(row_json), mode)792 conn.commit()793 total += len(rows)794 log(f"已从 PG 恢复 {table}: {total}")795 796 if total:797 log(f"{table} 恢复完成,共 {total} 条")798 799 800def restore_from_pg_if_needed(pg):801 if not wait_sqlite_ready(WAIT_SEC):802 log("SQLite 表长期未就绪,跳过 PG -> SQLite 恢复")803 return804 805 conn = sqlite_connect()806 try:807 restore_settings_to_sqlite(pg, conn)808 restore_table_if_empty(pg, conn, "model_prices", "ORDER BY model_key ASC")809 if table_exists_sqlite(conn, "api_key_aliases"):810 restore_table_if_empty(pg, conn, "api_key_aliases", "ORDER BY api_key_hash ASC")811 restore_table_if_empty(pg, conn, "usage_events", "ORDER BY COALESCE(sqlite_id, 0) ASC, event_hash ASC")812 restore_table_if_empty(pg, conn, "dead_letter_events", "ORDER BY COALESCE(sqlite_id, 0) ASC, record_hash ASC")813 814 max_usage = conn.execute("SELECT COALESCE(MAX(id), 0) FROM usage_events").fetchone()[0]815 max_dead = conn.execute("SELECT COALESCE(MAX(id), 0) FROM dead_letter_events").fetchone()[0]816 pg_set_state(pg, "sqlite_usage_last_id", str(max_usage))817 pg_set_state(pg, "sqlite_dead_last_id", str(max_dead))818 finally:819 conn.close()820 821 822def sync_settings_up(pg, conn):823 rows = fetch_sqlite_rows(conn, "SELECT * FROM settings")824 for row in rows:825 key = str(row.get("key", ""))826 if not key:827 continue828 value_text = row.get("value")829 pg.execute(830 f"""831 INSERT INTO {QSCHEMA}.settings (key, value_text, row_json, synced_at_ms)832 VALUES (%s, %s, %s, %s)833 ON CONFLICT (key) DO UPDATE834 SET value_text = EXCLUDED.value_text,835 row_json = EXCLUDED.row_json,836 synced_at_ms = EXCLUDED.synced_at_ms837 """,838 (key, value_text, Jsonb(row), now_ms()),839 )840 841 842def sync_model_prices_up(pg, conn):843 rows = fetch_sqlite_rows(conn, "SELECT * FROM model_prices")844 for row in rows:845 model_key = str(row.get("model") or row_hash(row))846 pg.execute(847 f"""848 INSERT INTO {QSCHEMA}.model_prices (model_key, row_json, synced_at_ms)849 VALUES (%s, %s, %s)850 ON CONFLICT (model_key) DO UPDATE851 SET row_json = EXCLUDED.row_json,852 synced_at_ms = EXCLUDED.synced_at_ms853 """,854 (model_key, Jsonb(row), now_ms()),855 )856 857 858def sync_api_key_aliases_up(pg, conn):859 if not table_exists_sqlite(conn, "api_key_aliases"):860 return861 rows = fetch_sqlite_rows(conn, "SELECT * FROM api_key_aliases")862 for row in rows:863 api_key_hash = str(row.get("api_key_hash") or row_hash(row))864 alias = row.get("alias")865 pg.execute(866 f"""867 INSERT INTO {QSCHEMA}.api_key_aliases (api_key_hash, alias, row_json, synced_at_ms)868 VALUES (%s, %s, %s, %s)869 ON CONFLICT (api_key_hash) DO UPDATE870 SET alias = EXCLUDED.alias,871 row_json = EXCLUDED.row_json,872 synced_at_ms = EXCLUDED.synced_at_ms873 """,874 (api_key_hash, alias, Jsonb(row), now_ms()),875 )876 877 878def sync_usage_events_up(pg, conn, batch: int):879 last_id = int(pg_get_state(pg, "sqlite_usage_last_id", "0"))880 rows = fetch_sqlite_rows(881 conn,882 "SELECT * FROM usage_events WHERE id > ? ORDER BY id ASC LIMIT ?",883 (last_id, batch),884 )885 if not rows:886 return 0887 888 max_id = last_id889 for row in rows:890 event_hash = str(row.get("event_hash") or row_hash(row))891 sqlite_id = int(row.get("id") or 0)892 timestamp_ms = int(row.get("timestamp_ms") or 0)893 model = row.get("model")894 auth_index = row.get("auth_index")895 failed = bool(row.get("failed") or False)896 897 pg.execute(898 f"""899 INSERT INTO {QSCHEMA}.usage_events900 (event_hash, sqlite_id, timestamp_ms, model, auth_index, failed, row_json, synced_at_ms)901 VALUES (%s, %s, %s, %s, %s, %s, %s, %s)902 ON CONFLICT (event_hash) DO UPDATE903 SET sqlite_id = EXCLUDED.sqlite_id,904 timestamp_ms = EXCLUDED.timestamp_ms,905 model = EXCLUDED.model,906 auth_index = EXCLUDED.auth_index,907 failed = EXCLUDED.failed,908 row_json = EXCLUDED.row_json,909 synced_at_ms = EXCLUDED.synced_at_ms910 """,911 (912 event_hash,913 sqlite_id,914 timestamp_ms,915 model,916 auth_index,917 failed,918 Jsonb(row),919 now_ms(),920 ),921 )922 if sqlite_id > max_id:923 max_id = sqlite_id924 925 pg_set_state(pg, "sqlite_usage_last_id", str(max_id))926 return len(rows)927 928 929def sync_dead_letters_up(pg, conn, batch: int):930 last_id = int(pg_get_state(pg, "sqlite_dead_last_id", "0"))931 rows = fetch_sqlite_rows(932 conn,933 "SELECT * FROM dead_letter_events WHERE id > ? ORDER BY id ASC LIMIT ?",934 (last_id, batch),935 )936 if not rows:937 return 0938 939 max_id = last_id940 for row in rows:941 sqlite_id = int(row.get("id") or 0)942 record_hash = str(943 row.get("record_hash")944 or row_hash(945 {946 "payload": row.get("payload"),947 "error": row.get("error"),948 "created_at_ms": row.get("created_at_ms"),949 "id": sqlite_id,950 }951 )952 )953 created_at_ms = int(row.get("created_at_ms") or 0)954 955 pg.execute(956 f"""957 INSERT INTO {QSCHEMA}.dead_letter_events958 (record_hash, sqlite_id, created_at_ms, row_json, synced_at_ms)959 VALUES (%s, %s, %s, %s, %s)960 ON CONFLICT (record_hash) DO UPDATE961 SET sqlite_id = EXCLUDED.sqlite_id,962 created_at_ms = EXCLUDED.created_at_ms,963 row_json = EXCLUDED.row_json,964 synced_at_ms = EXCLUDED.synced_at_ms965 """,966 (967 record_hash,968 sqlite_id,969 created_at_ms,970 Jsonb(row),971 now_ms(),972 ),973 )974 if sqlite_id > max_id:975 max_id = sqlite_id976 977 pg_set_state(pg, "sqlite_dead_last_id", str(max_id))978 return len(rows)979 980 981CONFIG_PAGE_HTML = r"""<!doctype html>982<html lang="zh-CN">983<head>984 <meta charset="utf-8" />985 <meta name="viewport" content="width=device-width, initial-scale=1" />986 <title>CPA Usage 同步配置</title>987 <style>988 :root { color-scheme: light dark; font-family: Inter, ui-sans-serif, system-ui, -apple-system, BlinkMacSystemFont, "Segoe UI", sans-serif; }989 body { margin: 0; background: #0f172a; color: #e2e8f0; }990 main { max-width: 880px; margin: 0 auto; padding: 32px 18px 56px; }991 .card { background: rgba(15, 23, 42, .92); border: 1px solid rgba(148, 163, 184, .28); border-radius: 18px; padding: 22px; box-shadow: 0 18px 60px rgba(0,0,0,.32); }992 h1 { margin: 0 0 8px; font-size: 28px; }993 h2 { margin: 24px 0 12px; font-size: 18px; }994 p { color: #94a3b8; line-height: 1.65; }995 label { display: block; margin: 14px 0 6px; color: #cbd5e1; font-weight: 600; }996 input[type="number"], input[type="password"] { width: 100%; box-sizing: border-box; border: 1px solid rgba(148,163,184,.35); border-radius: 12px; padding: 11px 12px; background: #020617; color: #f8fafc; font-size: 15px; }997 input[type="checkbox"] { transform: translateY(1px); margin-right: 8px; }998 .grid { display: grid; grid-template-columns: repeat(2, minmax(0, 1fr)); gap: 14px 18px; }999 .row { display: flex; gap: 10px; flex-wrap: wrap; align-items: center; }1000 button { border: 0; border-radius: 12px; padding: 10px 15px; cursor: pointer; font-weight: 700; background: #38bdf8; color: #082f49; }1001 button.secondary { background: #334155; color: #e2e8f0; }1002 button.danger { background: #fb7185; color: #450a0a; }1003 button:disabled { opacity: .55; cursor: wait; }1004 .muted { color: #94a3b8; font-size: 13px; }1005 .msg { margin-top: 14px; padding: 10px 12px; border-radius: 12px; display: none; white-space: pre-wrap; }1006 .ok { display: block; background: rgba(34,197,94,.14); color: #bbf7d0; border: 1px solid rgba(34,197,94,.3); }1007 .err { display: block; background: rgba(244,63,94,.14); color: #fecdd3; border: 1px solid rgba(244,63,94,.35); }1008 code { background: rgba(148,163,184,.15); padding: 2px 6px; border-radius: 8px; }1009 table { width: 100%; border-collapse: collapse; margin-top: 10px; }1010 td { border-bottom: 1px solid rgba(148,163,184,.18); padding: 8px 4px; color: #cbd5e1; }1011 td:last-child { text-align: right; color: #f8fafc; font-variant-numeric: tabular-nums; }1012 #app, #login { display: none; }1013 @media (max-width: 720px) { .grid { grid-template-columns: 1fr; } }1014 </style>1015</head>1016<body>1017<main>1018 <section class="card">1019 <h1>CPA Usage 同步配置</h1>1020 <p>保存后会写入 <code>/data/usage_sync_config.json</code> 并立即更新同步进程内存配置;无需重启 Space。</p>1021 1022 <div id="login">1023 <h2>登录</h2>1024 <p>请输入 <code>/management.html</code> 管理页面的登录密码。本页只校验 CPA-Manager-Plus 的管理员凭据,不接受其他环境变量密码或默认密码。</p>1025 <form id="loginForm">1026 <label for="password">管理密码</label>1027 <input id="password" type="password" autocomplete="current-password" required />1028 <p><button type="submit">进入配置</button></p>1029 </form>1030 <div id="loginMsg" class="msg"></div>1031 </div>1032 1033 <div id="app">1034 <div class="row">1035 <button id="refreshBtn" class="secondary" type="button">刷新状态</button>1036 <button id="logoutBtn" class="secondary" type="button">退出登录</button>1037 </div>1038 1039 <h2>保留策略</h2>1040 <form id="configForm">1041 <div class="grid">1042 <div>1043 <label for="retention_days">保留天数</label>1044 <input id="retention_days" name="retention_days" type="number" min="0" max="3650" required />1045 <div class="muted">0 表示不按时间删除。</div>1046 </div>1047 <div>1048 <label for="max_usage_events">最多保留 usage_events</label>1049 <input id="max_usage_events" name="max_usage_events" type="number" min="0" max="10000000" required />1050 <div class="muted">0 表示不按条数限制。</div>1051 </div>1052 <div>1053 <label for="max_dead_letter_events">最多保留 dead_letter_events</label>1054 <input id="max_dead_letter_events" name="max_dead_letter_events" type="number" min="0" max="1000000" required />1055 </div>1056 <div>1057 <label for="prune_interval_sec">裁剪间隔秒数</label>1058 <input id="prune_interval_sec" name="prune_interval_sec" type="number" min="0" max="86400" required />1059 <div class="muted">0 表示关闭定时裁剪,可手动点“立即裁剪”。</div>1060 </div>1061 <div>1062 <label for="sync_interval_sec">同步间隔秒数</label>1063 <input id="sync_interval_sec" name="sync_interval_sec" type="number" min="1" max="3600" required />1064 </div>1065 <div>1066 <label for="sync_batch">同步批量</label>1067 <input id="sync_batch" name="sync_batch" type="number" min="1" max="10000" required />1068 </div>1069 </div>1070 <p>1071 <label><input id="prune_sqlite" name="prune_sqlite" type="checkbox" /> 同时裁剪本地 SQLite</label>1072 <label><input id="vacuum_sqlite_on_prune" name="vacuum_sqlite_on_prune" type="checkbox" /> 裁剪后对 SQLite 执行 VACUUM(更慢,但更可能回收本地文件体积)</label>1073 </p>1074 <div class="row">1075 <button type="submit">保存并生效</button>1076 <button id="pruneBtn" class="secondary" type="button">立即裁剪</button>1077 <button id="purgeBtn" class="danger" type="button">清空 usage/dead</button>1078 </div>1079 </form>1080 <div id="msg" class="msg"></div>1081 1082 <h2>当前状态</h2>1083 <table id="statusTable"></table>1084 </div>1085 </section>1086</main>1087 1088<script>1089const $ = (id) => document.getElementById(id);1090const fields = ["sync_interval_sec", "sync_batch", "retention_days", "max_usage_events", "max_dead_letter_events", "prune_interval_sec"];1091function show(which) {1092 $("login").style.display = which === "login" ? "block" : "none";1093 $("app").style.display = which === "app" ? "block" : "none";1094}1095function message(el, text, ok=true) {1096 el.className = "msg " + (ok ? "ok" : "err");1097 el.textContent = text;1098}1099async function api(path, opts={}) {1100 const res = await fetch(path, {1101 credentials: "same-origin",1102 headers: {"Content-Type": "application/json"},1103 ...opts1104 });1105 let data = {};1106 try { data = await res.json(); } catch (_) {}1107 if (!res.ok) {1108 const err = new Error(data.error || ("HTTP " + res.status));1109 err.status = res.status;1110 throw err;1111 }1112 return data;1113}1114function fillConfig(cfg) {1115 for (const key of fields) $(key).value = cfg[key];1116 $("prune_sqlite").checked = !!cfg.prune_sqlite;1117 $("vacuum_sqlite_on_prune").checked = !!cfg.vacuum_sqlite_on_prune;1118}1119function readConfig() {1120 const cfg = {};1121 for (const key of fields) cfg[key] = Number($(key).value);1122 cfg.prune_sqlite = $("prune_sqlite").checked;1123 cfg.vacuum_sqlite_on_prune = $("vacuum_sqlite_on_prune").checked;1124 return cfg;1125}1126function renderStatus(data) {1127 const rows = [1128 ["schema", data.schema],1129 ["配置文件", data.config_path],1130 ["PG 状态", data.pg?.error || "ok"],1131 ["PG usage_events", data.pg?.usage_events],1132 ["PG dead_letter_events", data.pg?.dead_letter_events],1133 ["SQLite 状态", data.sqlite?.error || "ok"],1134 ["SQLite usage_events", data.sqlite?.usage_events],1135 ["SQLite dead_letter_events", data.sqlite?.dead_letter_events],1136 ["SQLite 路径", data.sqlite_path]1137 ];1138 $("statusTable").innerHTML = rows.map(([k,v]) => `<tr><td>${k}</td><td>${v ?? "-"}</td></tr>`).join("");1139}1140async function load() {1141 try {1142 const data = await api("/usage-sync-config/api/config");1143 fillConfig(data.config);1144 show("app");1145 await refreshStatus();1146 } catch (e) {1147 if (e.status === 401) show("login");1148 else { show("login"); message($("loginMsg"), e.message, false); }1149 }1150}1151async function refreshStatus() {1152 try {1153 renderStatus(await api("/usage-sync-config/api/status"));1154 } catch (e) {1155 message($("msg"), "刷新状态失败:" + e.message, false);1156 }1157}1158$("loginForm").addEventListener("submit", async (ev) => {1159 ev.preventDefault();1160 try {1161 await api("/usage-sync-config/api/login", {method: "POST", body: JSON.stringify({password: $("password").value})});1162 $("password").value = "";1163 message($("loginMsg"), "", true);1164 await load();1165 } catch (e) {1166 message($("loginMsg"), "登录失败:" + e.message, false);1167 }1168});1169$("configForm").addEventListener("submit", async (ev) => {1170 ev.preventDefault();1171 try {1172 const data = await api("/usage-sync-config/api/config", {method: "POST", body: JSON.stringify(readConfig())});1173 fillConfig(data.config);1174 message($("msg"), "已保存,配置已立即生效。", true);1175 } catch (e) {1176 message($("msg"), "保存失败:" + e.message, false);1177 }1178});1179$("refreshBtn").onclick = refreshStatus;1180$("logoutBtn").onclick = async () => { await api("/usage-sync-config/api/logout", {method: "POST", body: "{}"}).catch(()=>{}); show("login"); };1181$("pruneBtn").onclick = async () => {1182 try {1183 const data = await api("/usage-sync-config/api/prune", {method: "POST", body: "{}"});1184 message($("msg"), data.message || "裁剪完成。", true);1185 await refreshStatus();1186 } catch (e) {1187 message($("msg"), "裁剪失败:" + e.message, false);1188 }1189};1190$("purgeBtn").onclick = async () => {1191 const confirmText = prompt("此操作会清空 usage_events 和 dead_letter_events。请输入 DELETE 确认:");1192 if (confirmText !== "DELETE") return;1193 try {1194 const data = await api("/usage-sync-config/api/purge", {method: "POST", body: JSON.stringify({confirm: "DELETE"})});1195 message($("msg"), data.message || "清空完成。", true);1196 await refreshStatus();1197 } catch (e) {1198 message($("msg"), "清空失败:" + e.message, false);1199 }1200};