CoolFace
Apppublic

PascalSmile8358/openapi

sourceHugging Faceupdated 4mo agoView on Hugging Face
0likes
sync_usage.py557 linesDownload Raw Back to docker
1#!/usr/bin/env python32import hashlib3import json4import os5import re6import sqlite37import sys8import time9import urllib.error10import urllib.request11 12import psycopg13from psycopg.types.json import Jsonb14 15DB_PATH = os.getenv("USAGE_DB_PATH", "/data/usage.sqlite")16PG_DSN = os.getenv("USAGE_SYNC_PG_DSN") or os.getenv("PGSTORE_DSN", "")17SCHEMA = os.getenv("USAGE_SYNC_SCHEMA", "cpa_manager_usage")18INTERVAL = int(os.getenv("USAGE_SYNC_INTERVAL_SEC", "3"))19BATCH = int(os.getenv("USAGE_SYNC_BATCH", "500"))20MANAGER_BASE = os.getenv("USAGE_SYNC_MANAGER_BASE", "http://127.0.0.1:18317").rstrip("/")21WAIT_SEC = int(os.getenv("USAGE_SYNC_WAIT_SEC", "300"))22FORCE_UPSTREAM_URL = os.getenv("USAGE_SYNC_FORCE_UPSTREAM_URL", "").strip()23FORCE_MANAGEMENT_KEY = os.getenv("USAGE_SYNC_FORCE_MANAGEMENT_KEY", "").strip()24FORCE_QUEUE = os.getenv("USAGE_SYNC_FORCE_QUEUE", "").strip()25FORCE_POP_SIDE = os.getenv("USAGE_SYNC_FORCE_POP_SIDE", "").strip()26 27if not re.fullmatch(r"[A-Za-z_][A-Za-z0-9_]*", SCHEMA):28    raise SystemExit(f"invalid USAGE_SYNC_SCHEMA: {SCHEMA}")29 30QSCHEMA = SCHEMA31 32 33def log(msg: str):34    print(f"[usage-sync] {msg}", flush=True)35 36 37def now_ms() -> int:38    return int(time.time() * 1000)39 40 41def json_dumps(obj) -> str:42    return json.dumps(obj, ensure_ascii=False, sort_keys=True, separators=(",", ":"))43 44 45def row_hash(obj) -> str:46    return hashlib.sha256(json_dumps(obj).encode("utf-8")).hexdigest()47 48 49def pg_connect():50    if not PG_DSN:51        return None52    return psycopg.connect(PG_DSN, autocommit=True, prepare_threshold=None)53 54 55def sqlite_connect():56    conn = sqlite3.connect(DB_PATH, timeout=30, check_same_thread=False)57    conn.row_factory = sqlite3.Row58    conn.execute("PRAGMA busy_timeout = 5000")59    conn.execute("PRAGMA journal_mode = WAL")60    return conn61 62 63def table_exists_sqlite(conn, table: str) -> bool:64    row = conn.execute(65        "SELECT 1 FROM sqlite_master WHERE type='table' AND name=?",66        (table,),67    ).fetchone()68    return row is not None69 70 71def get_sqlite_columns(conn, table: str):72    rows = conn.execute(f"PRAGMA table_info({table})").fetchall()73    return [r["name"] for r in rows]74 75 76def fetch_sqlite_rows(conn, sql: str, params=()):77    return [dict(r) for r in conn.execute(sql, params).fetchall()]78 79 80def sqlite_count(conn, table: str) -> int:81    row = conn.execute(f"SELECT COUNT(*) AS c FROM {table}").fetchone()82    return int(row["c"] if row else 0)83 84 85def pg_get_state(pg, key: str, default="0"):86    row = pg.execute(87        f"SELECT value FROM {QSCHEMA}.sync_state WHERE key=%s",88        (key,),89    ).fetchone()90    return row[0] if row else default91 92 93def pg_set_state(pg, key: str, value: str):94    pg.execute(95        f"""96        INSERT INTO {QSCHEMA}.sync_state (key, value, updated_at_ms)97        VALUES (%s, %s, %s)98        ON CONFLICT (key) DO UPDATE99        SET value = EXCLUDED.value,100            updated_at_ms = EXCLUDED.updated_at_ms101        """,102        (key, str(value), now_ms()),103    )104 105 106def ensure_pg_schema(pg):107    pg.execute(f"CREATE SCHEMA IF NOT EXISTS {QSCHEMA}")108 109    pg.execute(110        f"""111        CREATE TABLE IF NOT EXISTS {QSCHEMA}.settings (112            key text PRIMARY KEY,113            value_text text,114            row_json jsonb NOT NULL,115            synced_at_ms bigint NOT NULL116        )117        """118    )119 120    pg.execute(121        f"""122        CREATE TABLE IF NOT EXISTS {QSCHEMA}.model_prices (123            model_key text PRIMARY KEY,124            row_json jsonb NOT NULL,125            synced_at_ms bigint NOT NULL126        )127        """128    )129 130    pg.execute(131        f"""132        CREATE TABLE IF NOT EXISTS {QSCHEMA}.usage_events (133            event_hash text PRIMARY KEY,134            sqlite_id bigint,135            timestamp_ms bigint,136            model text,137            auth_index text,138            failed boolean,139            row_json jsonb NOT NULL,140            synced_at_ms bigint NOT NULL141        )142        """143    )144    pg.execute(145        f"CREATE INDEX IF NOT EXISTS idx_{QSCHEMA}_usage_events_ts ON {QSCHEMA}.usage_events(timestamp_ms DESC)"146    )147 148    pg.execute(149        f"""150        CREATE TABLE IF NOT EXISTS {QSCHEMA}.dead_letter_events (151            record_hash text PRIMARY KEY,152            sqlite_id bigint,153            created_at_ms bigint,154            row_json jsonb NOT NULL,155            synced_at_ms bigint NOT NULL156        )157        """158    )159 160    pg.execute(161        f"""162        CREATE TABLE IF NOT EXISTS {QSCHEMA}.sync_state (163            key text PRIMARY KEY,164            value text NOT NULL,165            updated_at_ms bigint NOT NULL166        )167        """168    )169 170 171def wait_http_ok(url: str, timeout_sec: int) -> bool:172    deadline = time.time() + timeout_sec173    while time.time() < deadline:174        try:175            with urllib.request.urlopen(url, timeout=5) as resp:176                if 200 <= resp.status < 300:177                    return True178        except Exception:179            pass180        time.sleep(2)181    return False182 183 184def wait_sqlite_ready(timeout_sec: int):185    deadline = time.time() + timeout_sec186    while time.time() < deadline:187        try:188            if not os.path.exists(DB_PATH):189                time.sleep(2)190                continue191            conn = sqlite_connect()192            ok = all(193                table_exists_sqlite(conn, t)194                for t in ("settings", "model_prices", "usage_events", "dead_letter_events")195            )196            conn.close()197            if ok:198                return True199        except Exception:200            pass201        time.sleep(2)202    return False203 204 205def normalize_setup_payload(raw: str):206    try:207        data = json.loads(raw)208    except Exception:209        log("PG 中 settings.setup 不是合法 JSON,跳过 setup 恢复")210        return None211 212    base = (213        FORCE_UPSTREAM_URL214        or data.get("cpaBaseUrl")215        or data.get("CPAUpstreamURL")216        or data.get("upstream")217        or ""218    ).strip().rstrip("/")219 220    key = (221        FORCE_MANAGEMENT_KEY222        or data.get("managementKey")223        or data.get("ManagementKey")224        or ""225    ).strip()226 227    queue = (228        FORCE_QUEUE229        or data.get("queue")230        or data.get("Queue")231        or "usage"232    ).strip()233 234    pop_side = (235        FORCE_POP_SIDE236        or data.get("popSide")237        or data.get("PopSide")238        or "right"239    ).strip()240 241    if not base or not key:242        log("PG 中 settings.setup 缺少 cpaBaseUrl/managementKey,跳过 setup 恢复")243        return None244 245    return {246        "cpaBaseUrl": base,247        "managementKey": key,248        "queue": queue,249        "popSide": pop_side,250    }251 252 253def restore_setup_via_http(pg):254    row = pg.execute(255        f"SELECT value_text FROM {QSCHEMA}.settings WHERE key='setup'"256    ).fetchone()257    if not row or not row[0]:258        log("PG 中还没有 setup,等待你首次在面板保存后再自动镜像")259        return False260 261    payload = normalize_setup_payload(row[0])262    if not payload:263        return False264 265    deadline = time.time() + WAIT_SEC266    body = json.dumps(payload).encode("utf-8")267    req = urllib.request.Request(268        f"{MANAGER_BASE}/setup",269        data=body,270        headers={"Content-Type": "application/json"},271        method="POST",272    )273 274    while time.time() < deadline:275        try:276            with urllib.request.urlopen(req, timeout=10) as resp:277                if 200 <= resp.status < 300:278                    log(f"已从 PG 恢复 cpa-manager setup,upstream={payload['cpaBaseUrl']}")279                    return True280        except urllib.error.HTTPError as e:281            try:282                detail = e.read().decode("utf-8", "ignore")283            except Exception:284                detail = str(e)285            log(f"/setup 返回 {e.code},稍后重试:{detail}")286        except Exception as e:287            log(f"/setup 暂时失败,稍后重试:{e}")288        time.sleep(2)289 290    log("setup 自动恢复超时;如果面板仍未接管统计,请手动打开 management.html 保存一次")291    return False292 293 294def sqlite_insert(conn, table: str, row: dict, mode: str):295    cols = get_sqlite_columns(conn, table)296    use_cols = [c for c in cols if c in row]297    if not use_cols:298        return299    placeholders = ",".join("?" for _ in use_cols)300    col_sql = ",".join(use_cols)301    values = [row[c] for c in use_cols]302    conn.execute(303        f"INSERT OR {mode} INTO {table} ({col_sql}) VALUES ({placeholders})",304        values,305    )306 307 308def restore_settings_to_sqlite(pg, conn):309    rows = pg.execute(310        f"SELECT row_json FROM {QSCHEMA}.settings WHERE key <> 'setup'"311    ).fetchall()312    for (row_json,) in rows:313        sqlite_insert(conn, "settings", dict(row_json), "REPLACE")314    conn.commit()315 316 317def restore_table_if_empty(pg, conn, table: str, order_sql: str):318    if sqlite_count(conn, table) > 0:319        return320 321    total = 0322    while True:323        rows = pg.execute(324            f"SELECT row_json FROM {QSCHEMA}.{table} {order_sql} LIMIT %s OFFSET %s",325            (BATCH, total),326        ).fetchall()327        if not rows:328            break329 330        for (row_json,) in rows:331            mode = "REPLACE" if table in ("settings", "model_prices") else "IGNORE"332            sqlite_insert(conn, table, dict(row_json), mode)333        conn.commit()334        total += len(rows)335        log(f"已从 PG 恢复 {table}: {total}")336 337    if total:338        log(f"{table} 恢复完成,共 {total} 条")339 340 341def restore_from_pg_if_needed(pg):342    if not wait_sqlite_ready(WAIT_SEC):343        log("SQLite 表长期未就绪,跳过 PG -> SQLite 恢复")344        return345 346    conn = sqlite_connect()347    try:348        restore_settings_to_sqlite(pg, conn)349        restore_table_if_empty(pg, conn, "model_prices", "ORDER BY model_key ASC")350        restore_table_if_empty(pg, conn, "usage_events", "ORDER BY COALESCE(sqlite_id, 0) ASC, event_hash ASC")351        restore_table_if_empty(pg, conn, "dead_letter_events", "ORDER BY COALESCE(sqlite_id, 0) ASC, record_hash ASC")352 353        max_usage = conn.execute("SELECT COALESCE(MAX(id), 0) FROM usage_events").fetchone()[0]354        max_dead = conn.execute("SELECT COALESCE(MAX(id), 0) FROM dead_letter_events").fetchone()[0]355        pg_set_state(pg, "sqlite_usage_last_id", str(max_usage))356        pg_set_state(pg, "sqlite_dead_last_id", str(max_dead))357    finally:358        conn.close()359 360 361def sync_settings_up(pg, conn):362    rows = fetch_sqlite_rows(conn, "SELECT * FROM settings")363    for row in rows:364        key = str(row.get("key", ""))365        if not key:366            continue367        value_text = row.get("value")368        pg.execute(369            f"""370            INSERT INTO {QSCHEMA}.settings (key, value_text, row_json, synced_at_ms)371            VALUES (%s, %s, %s, %s)372            ON CONFLICT (key) DO UPDATE373            SET value_text = EXCLUDED.value_text,374                row_json = EXCLUDED.row_json,375                synced_at_ms = EXCLUDED.synced_at_ms376            """,377            (key, value_text, Jsonb(row), now_ms()),378        )379 380 381def sync_model_prices_up(pg, conn):382    rows = fetch_sqlite_rows(conn, "SELECT * FROM model_prices")383    for row in rows:384        model_key = str(row.get("model") or row_hash(row))385        pg.execute(386            f"""387            INSERT INTO {QSCHEMA}.model_prices (model_key, row_json, synced_at_ms)388            VALUES (%s, %s, %s)389            ON CONFLICT (model_key) DO UPDATE390            SET row_json = EXCLUDED.row_json,391                synced_at_ms = EXCLUDED.synced_at_ms392            """,393            (model_key, Jsonb(row), now_ms()),394        )395 396 397def sync_usage_events_up(pg, conn):398    last_id = int(pg_get_state(pg, "sqlite_usage_last_id", "0"))399    rows = fetch_sqlite_rows(400        conn,401        "SELECT * FROM usage_events WHERE id > ? ORDER BY id ASC LIMIT ?",402        (last_id, BATCH),403    )404    if not rows:405        return 0406 407    max_id = last_id408    for row in rows:409        event_hash = str(row.get("event_hash") or row_hash(row))410        sqlite_id = int(row.get("id") or 0)411        timestamp_ms = int(row.get("timestamp_ms") or 0)412        model = row.get("model")413        auth_index = row.get("auth_index")414        failed = bool(row.get("failed") or False)415 416        pg.execute(417            f"""418            INSERT INTO {QSCHEMA}.usage_events419            (event_hash, sqlite_id, timestamp_ms, model, auth_index, failed, row_json, synced_at_ms)420            VALUES (%s, %s, %s, %s, %s, %s, %s, %s)421            ON CONFLICT (event_hash) DO UPDATE422            SET sqlite_id = EXCLUDED.sqlite_id,423                timestamp_ms = EXCLUDED.timestamp_ms,424                model = EXCLUDED.model,425                auth_index = EXCLUDED.auth_index,426                failed = EXCLUDED.failed,427                row_json = EXCLUDED.row_json,428                synced_at_ms = EXCLUDED.synced_at_ms429            """,430            (431                event_hash,432                sqlite_id,433                timestamp_ms,434                model,435                auth_index,436                failed,437                Jsonb(row),438                now_ms(),439            ),440        )441        if sqlite_id > max_id:442            max_id = sqlite_id443 444    pg_set_state(pg, "sqlite_usage_last_id", str(max_id))445    return len(rows)446 447 448def sync_dead_letters_up(pg, conn):449    last_id = int(pg_get_state(pg, "sqlite_dead_last_id", "0"))450    rows = fetch_sqlite_rows(451        conn,452        "SELECT * FROM dead_letter_events WHERE id > ? ORDER BY id ASC LIMIT ?",453        (last_id, BATCH),454    )455    if not rows:456        return 0457 458    max_id = last_id459    for row in rows:460        sqlite_id = int(row.get("id") or 0)461        record_hash = str(462            row.get("record_hash")463            or row_hash(464                {465                    "payload": row.get("payload"),466                    "error": row.get("error"),467                    "created_at_ms": row.get("created_at_ms"),468                    "id": sqlite_id,469                }470            )471        )472        created_at_ms = int(row.get("created_at_ms") or 0)473 474        pg.execute(475            f"""476            INSERT INTO {QSCHEMA}.dead_letter_events477            (record_hash, sqlite_id, created_at_ms, row_json, synced_at_ms)478            VALUES (%s, %s, %s, %s, %s)479            ON CONFLICT (record_hash) DO UPDATE480            SET sqlite_id = EXCLUDED.sqlite_id,481                created_at_ms = EXCLUDED.created_at_ms,482                row_json = EXCLUDED.row_json,483                synced_at_ms = EXCLUDED.synced_at_ms484            """,485            (486                record_hash,487                sqlite_id,488                created_at_ms,489                Jsonb(row),490                now_ms(),491            ),492        )493        if sqlite_id > max_id:494            max_id = sqlite_id495 496    pg_set_state(pg, "sqlite_dead_last_id", str(max_id))497    return len(rows)498 499 500def sync_loop(pg):501    if not wait_sqlite_ready(WAIT_SEC):502        log("SQLite 长期未就绪,退出同步进程")503        sys.exit(1)504 505    while True:506        try:507            conn = sqlite_connect()508            try:509                sync_settings_up(pg, conn)510                sync_model_prices_up(pg, conn)511 512                moved_usage = 0513                while True:514                    n = sync_usage_events_up(pg, conn)515                    moved_usage += n516                    if n < BATCH:517                        break518 519                moved_dead = 0520                while True:521                    n = sync_dead_letters_up(pg, conn)522                    moved_dead += n523                    if n < BATCH:524                        break525 526                if moved_usage or moved_dead:527                    log(f"本轮已同步 usage_events={moved_usage}, dead_letter_events={moved_dead}")528            finally:529                conn.close()530        except Exception as e:531            log(f"同步失败:{e}")532 533        time.sleep(INTERVAL)534 535 536def main():537    if not PG_DSN:538        log("未设置 PGSTORE_DSN / USAGE_SYNC_PG_DSN,仅使用本地 SQLite,不做 PG 镜像")539        while True:540            time.sleep(3600)541 542    pg = pg_connect()543    ensure_pg_schema(pg)544    log(f"PG 镜像表已就绪,schema={SCHEMA}")545 546    if not wait_http_ok(f"{MANAGER_BASE}/health", WAIT_SEC):547        log("cpa-manager /health 长期未就绪,退出")548        sys.exit(1)549 550    restore_setup_via_http(pg)551    restore_from_pg_if_needed(pg)552    sync_loop(pg)553 554 555if __name__ == "__main__":556    main()557