PascalSmile8358/openapi
0
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 