CoolFace
Apppublic

kingcap/openapi

sourceHugging Faceupdated 2mo agoView on Hugging Face
0likes
sync_usage.py1606 linesDownload Raw Back to docker
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};

Showing the first 1,200 of 1606 lines. Download the file for the rest.