Files
grok-keygen-new/cpa/push_queue.py
T
chaos df9213463d Add bblbb mail provider and CPA push queue; sync runtime artifacts.
Switch default email provider to mail.bblbb.com, queue failed CPA auth pushes for retry, and refresh local auth/account outputs.
2026-07-14 09:14:06 +08:00

280 lines
8.3 KiB
Python

"""CPA 远端推送失败队列。
本地 auth 已写成功但 POST /v0/management/auth-files 失败时,把文件名记入
`cpa_push_pending.txt`;下次开启推送时先重试队列里的文件,成功则剔除。
行格式(兼容只写文件名):
xai-email@x.com.json----reason----unix_ts
"""
from __future__ import annotations
import json
import threading
import time
from pathlib import Path
from typing import Any, Callable
from .client import CpaPushError, push_auth_file
PENDING_NAME = "cpa_push_pending.txt"
_lock = threading.RLock()
def pending_path(out_dir: str | Path) -> Path:
return Path(out_dir).expanduser().resolve() / PENDING_NAME
def _normalize_filename(name: str) -> str:
name = (name or "").strip().replace("\\", "/").split("/")[-1]
if name and not name.endswith(".json"):
name += ".json"
return name
def _parse_line(line: str) -> tuple[str, str, int] | None:
line = (line or "").strip()
if not line or line.startswith("#"):
return None
parts = line.split("----")
filename = _normalize_filename(parts[0])
if not filename:
return None
reason = parts[1].strip() if len(parts) > 1 else ""
ts = 0
if len(parts) > 2:
try:
ts = int(parts[2].strip())
except ValueError:
ts = 0
return filename, reason, ts
def load_pending(out_dir: str | Path) -> list[tuple[str, str, int]]:
"""按文件名去重,保留最后一次记录。"""
path = pending_path(out_dir)
if not path.is_file():
return []
ordered: dict[str, tuple[str, str, int]] = {}
try:
text = path.read_text(encoding="utf-8", errors="replace")
except OSError:
return []
for line in text.splitlines():
row = _parse_line(line)
if row:
ordered[row[0]] = row
return list(ordered.values())
def _rewrite(out_dir: str | Path, rows: list[tuple[str, str, int]]) -> None:
out = Path(out_dir).expanduser().resolve()
out.mkdir(parents=True, exist_ok=True)
path = pending_path(out)
lines = [
f"{fn}----{reason}----{ts or int(time.time())}\n"
for fn, reason, ts in rows
if fn
]
tmp = path.with_suffix(".tmp")
tmp.write_text("".join(lines), encoding="utf-8")
tmp.replace(path)
def record_push_failure(
out_dir: str | Path,
filename: str,
reason: str = "",
) -> None:
"""推送失败:加入待重试队列(同名覆盖为最新原因)。"""
fn = _normalize_filename(filename)
if not fn:
return
reason = (reason or "").replace("\n", " ").replace("\r", " ").strip()[:300]
now = int(time.time())
with _lock:
rows = {r[0]: r for r in load_pending(out_dir)}
rows[fn] = (fn, reason, now)
_rewrite(out_dir, list(rows.values()))
def remove_push_pending(out_dir: str | Path, filename: str) -> None:
"""推送成功:从队列移除。"""
fn = _normalize_filename(filename)
if not fn:
return
with _lock:
rows = [r for r in load_pending(out_dir) if r[0] != fn]
path = pending_path(out_dir)
if not rows:
if path.is_file():
try:
path.unlink()
except OSError:
_rewrite(out_dir, [])
return
_rewrite(out_dir, rows)
def push_local_auth_file(
out_dir: str | Path,
filename: str,
*,
remote_base: str,
secret: str,
proxy: str | None = None,
verify_tls: bool = True,
) -> tuple[bool, int, str]:
"""读取本地 xai-*.json 并推送。文件缺失返回 (False, 0, 'missing')."""
fn = _normalize_filename(filename)
path = Path(out_dir).expanduser().resolve() / fn
if not path.is_file():
return False, 0, "missing local file"
try:
payload = json.loads(path.read_text(encoding="utf-8"))
except Exception as exc: # noqa: BLE001
return False, 0, f"read/json: {exc}"
return push_auth_file(
remote_base=remote_base,
secret=secret,
filename=fn,
payload=payload,
proxy=proxy,
verify_tls=verify_tls,
)
def flush_push_pending(
out_dir: str | Path,
*,
remote_base: str,
secret: str,
proxy: str | None = None,
verify_tls: bool = True,
log: Callable[[str], None] | None = None,
) -> dict[str, Any]:
"""重试队列中所有仍存在的本地 auth 文件。成功剔除,失败保留。"""
log = log or (lambda m: None)
remote_base = (remote_base or "").strip()
secret = (secret or "").strip()
if not remote_base or not secret:
return {"ok": 0, "fail": 0, "missing": 0, "skipped": True, "reason": "no remote config"}
with _lock:
pending = load_pending(out_dir)
if not pending:
return {"ok": 0, "fail": 0, "missing": 0, "total": 0}
log(f"[cpa] 重试待推送队列 count={len(pending)}")
ok_n = fail_n = missing_n = 0
for filename, old_reason, _ts in pending:
try:
ok, status, text = push_local_auth_file(
out_dir,
filename,
remote_base=remote_base,
secret=secret,
proxy=proxy,
verify_tls=verify_tls,
)
except CpaPushError as exc:
ok, status, text = False, -1, str(exc)
except Exception as exc: # noqa: BLE001
ok, status, text = False, -1, str(exc)
if status == 0 and text.startswith("missing"):
missing_n += 1
remove_push_pending(out_dir, filename)
log(f"[cpa] 待推送文件已不存在,移出队列: {filename}")
continue
if ok:
ok_n += 1
remove_push_pending(out_dir, filename)
log(f"[cpa] 队列重推成功: {filename} (HTTP {status})")
else:
fail_n += 1
reason = f"HTTP {status}: {text[:200]}" if status > 0 else text[:200]
record_push_failure(out_dir, filename, reason or old_reason)
log(f"[!] [cpa] 队列重推失败: {filename} {reason}")
return {
"ok": ok_n,
"fail": fail_n,
"missing": missing_n,
"total": len(pending),
"remaining": len(load_pending(out_dir)),
}
def push_with_queue(
out_dir: str | Path,
filename: str,
payload: dict | bytes | str | None,
*,
remote_base: str,
secret: str,
proxy: str | None = None,
verify_tls: bool = True,
flush_first: bool = True,
log: Callable[[str], None] | None = None,
) -> dict[str, Any]:
"""推送当前文件;可选先 flush 队列。失败记入 pending,成功移出。
返回 {pushed, push_status?, push_error?, flush?}。
"""
log = log or (lambda m: None)
result: dict[str, Any] = {"pushed": False}
remote_base = (remote_base or "").strip()
secret = (secret or "").strip()
if not remote_base or not secret:
result["push_error"] = "cpa_remote_base/secret 未配置"
return result
if flush_first:
result["flush"] = flush_push_pending(
out_dir,
remote_base=remote_base,
secret=secret,
proxy=proxy,
verify_tls=verify_tls,
log=log,
)
fn = _normalize_filename(filename)
try:
if payload is None:
ok, status, text = push_local_auth_file(
out_dir,
fn,
remote_base=remote_base,
secret=secret,
proxy=proxy,
verify_tls=verify_tls,
)
else:
ok, status, text = push_auth_file(
remote_base=remote_base,
secret=secret,
filename=fn,
payload=payload,
proxy=proxy,
verify_tls=verify_tls,
)
result["push_status"] = status
if ok:
result["pushed"] = True
remove_push_pending(out_dir, fn)
log(f"[Debug] 已推送远端 (HTTP {status}): {fn}")
else:
result["push_error"] = text[:300]
record_push_failure(out_dir, fn, f"HTTP {status}: {text[:200]}")
log(f"[!] [cpa] 推送远端失败 HTTP {status}: {text[:200]}")
except Exception as exc: # noqa: BLE001
result["push_error"] = str(exc)
record_push_failure(out_dir, fn, str(exc))
log(f"[!] [cpa] 推送远端异常: {exc}")
return result