rotate.py 8.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203
  1. """Single-pool rotation: probe main pool and hard-delete unhealthy accounts."""
  2. from __future__ import annotations
  3. import argparse
  4. import time
  5. from dataclasses import dataclass
  6. from datetime import datetime
  7. from pathlib import Path
  8. from core.settings import AppSettings
  9. from platforms.chatgpt.pool import load_token_record
  10. from .common import CpaClient, DEFAULT_MANAGEMENT_BASE_URL, DEFAULT_POOL_DIR, now
  11. from .rotate_probe import _collect_quota_probe_results, classify_status_message
  12. from .rotate_promote import handle_unhealthy_entries
  13. from .rotate_runtime import _fetch_main_pool_entries, _maybe_reconcile_cpa_runtime
  14. @dataclass
  15. class RotateResult:
  16. main_pool_before: int = 0
  17. main_pool_after: int = 0
  18. deleted_401: int = 0
  19. deleted_429: int = 0
  20. quota_probed: int = 0
  21. quota_probe_401: int = 0
  22. quota_probe_429: int = 0
  23. quota_probe_skipped: int = 0
  24. def rotate_once(
  25. pool_dir: Path,
  26. *,
  27. client: object | None = None,
  28. management_base_url: str = DEFAULT_MANAGEMENT_BASE_URL,
  29. cpa_management_key: str | None = None,
  30. rotate_probe_workers: int = 8,
  31. fresh_grace_seconds: int = 0,
  32. cpa_runtime_reconcile_enabled: bool = True,
  33. cpa_runtime_reconcile_cooldown_seconds: int = 300,
  34. cpa_runtime_reconcile_restart_enabled: bool = False,
  35. ) -> RotateResult:
  36. result = RotateResult()
  37. backend_client = client or CpaClient(management_base_url, management_key=cpa_management_key)
  38. if not getattr(backend_client, "health_check")():
  39. return result
  40. pool_dir = Path(pool_dir).expanduser().resolve()
  41. pool_dir.mkdir(parents=True, exist_ok=True)
  42. entries = _fetch_main_pool_entries(
  43. management_base_url,
  44. client=backend_client,
  45. management_key=cpa_management_key,
  46. )
  47. if entries is None:
  48. return result
  49. reg_entries = [entry for entry in entries if "@" in str(entry.get("name", ""))]
  50. result.main_pool_before = len(reg_entries)
  51. management_key = None
  52. if isinstance(backend_client, CpaClient):
  53. management_key = backend_client._resolve_key() # noqa: SLF001
  54. elif hasattr(backend_client, "_resolve_key"):
  55. try:
  56. management_key = getattr(backend_client, "_resolve_key")()
  57. except Exception:
  58. management_key = None
  59. probe_results: dict[str, tuple[int, str, bool]] = {}
  60. if management_key:
  61. initial_classified = []
  62. for entry in reg_entries:
  63. status_message = str(entry.get("status_message", ""))
  64. classified_code = classify_status_message(status_message)
  65. if classified_code in {401, 429}:
  66. continue
  67. initial_classified.append(entry)
  68. probe_candidates: list[dict] = []
  69. grace_skipped = 0
  70. for entry in initial_classified:
  71. pool_file = pool_dir / str(entry.get("name") or "").strip()
  72. if _is_fresh_pool_entry(pool_file, fresh_grace_seconds):
  73. grace_skipped += 1
  74. continue
  75. probe_candidates.append(entry)
  76. if probe_candidates:
  77. probe_results, probe_counters = _collect_quota_probe_results(
  78. probe_candidates,
  79. management_key=management_key,
  80. management_base_url=management_base_url,
  81. max_count=0,
  82. workers=rotate_probe_workers,
  83. )
  84. result.quota_probed = probe_counters["probed"]
  85. result.quota_probe_401 = probe_counters["probe_401"]
  86. result.quota_probe_429 = probe_counters["probe_429"]
  87. result.quota_probe_skipped = probe_counters["probe_skipped"] + grace_skipped
  88. else:
  89. result.quota_probe_skipped = grace_skipped
  90. handle_unhealthy_entries(
  91. result=result,
  92. reg_entries=reg_entries,
  93. probe_results=probe_results,
  94. pool_dir=pool_dir,
  95. backend_client=backend_client,
  96. now_func=now,
  97. classify_status_message_func=classify_status_message,
  98. is_deactivated_status_message_func=lambda message: "deactivated" in str(message or "").lower(),
  99. )
  100. result.main_pool_after = result.main_pool_before - result.deleted_401 - result.deleted_429
  101. _maybe_reconcile_cpa_runtime(
  102. pool_dir=pool_dir,
  103. management_base_url=management_base_url,
  104. enabled=cpa_runtime_reconcile_enabled,
  105. cooldown_seconds=cpa_runtime_reconcile_cooldown_seconds,
  106. restart_enabled=cpa_runtime_reconcile_restart_enabled,
  107. state_file=pool_dir / "cpa_runtime_reconcile_state.json",
  108. client=backend_client,
  109. management_key=cpa_management_key,
  110. )
  111. return result
  112. def _is_fresh_pool_entry(pool_file: Path, fresh_grace_seconds: int) -> bool:
  113. if int(fresh_grace_seconds or 0) <= 0 or not pool_file.is_file():
  114. return False
  115. try:
  116. payload = load_token_record(pool_file)
  117. except Exception:
  118. return False
  119. created_at = str(payload.get("created_at") or "").strip()
  120. if not created_at:
  121. return False
  122. try:
  123. created_dt = datetime.fromisoformat(created_at)
  124. except Exception:
  125. return False
  126. age_seconds = (datetime.now().astimezone() - created_dt).total_seconds()
  127. return age_seconds < max(0, int(fresh_grace_seconds))
  128. def print_rotate_summary(r: RotateResult) -> None:
  129. print(
  130. f"[{now()}] [rotate] summary"
  131. f" | 主池: {r.main_pool_before} → {r.main_pool_after}"
  132. f" | 401删除: {r.deleted_401}"
  133. f" | quota探测: {r.quota_probed}"
  134. f" | probe401: {r.quota_probe_401}"
  135. f" | probe429: {r.quota_probe_429}"
  136. f" | probe跳过: {r.quota_probe_skipped}"
  137. )
  138. def main() -> None:
  139. env_settings = AppSettings.from_env()
  140. parser = argparse.ArgumentParser(description="Rotate zhuce6 backend main pool in single-pool mode")
  141. parser.add_argument("--pool-dir", default=str(env_settings.pool_dir or DEFAULT_POOL_DIR), help="本地 pool 目录")
  142. parser.add_argument("--interval", type=int, default=env_settings.rotate_interval, help="轮换间隔秒数")
  143. parser.add_argument("--once", action="store_true", help="只执行一轮")
  144. parser.add_argument("--management-base-url", default=env_settings.cpa_management_base_url or DEFAULT_MANAGEMENT_BASE_URL, help="CPA management base url")
  145. parser.add_argument("--management-key", default=env_settings.cpa_management_key, help="可选 CPA management key")
  146. parser.add_argument("--rotate-probe-workers", type=int, default=env_settings.rotate_probe_workers, help="quota probe 并发数")
  147. parser.add_argument("--fresh-grace-seconds", type=int, default=env_settings.rotate_fresh_grace_seconds, help="fresh 账号在 grace 窗口内跳过 rotate 探测")
  148. parser.add_argument("--cpa-runtime-reconcile-enabled", action="store_true" if not env_settings.cpa_runtime_reconcile_enabled else "store_false", default=env_settings.cpa_runtime_reconcile_enabled, help="是否启用 CPA runtime drift 检测")
  149. parser.add_argument("--cpa-runtime-reconcile-cooldown-seconds", type=int, default=env_settings.cpa_runtime_reconcile_cooldown_seconds, help="CPA runtime drift 观测 cooldown")
  150. parser.add_argument("--cpa-runtime-reconcile-restart-enabled", action="store_true" if not env_settings.cpa_runtime_reconcile_restart_enabled else "store_false", default=env_settings.cpa_runtime_reconcile_restart_enabled, help="保留兼容字段, API-only 模式下不会自动重启")
  151. args = parser.parse_args()
  152. interval = max(1, args.interval)
  153. pool_dir = Path(args.pool_dir).expanduser().resolve()
  154. print(
  155. "[rotate] 启动"
  156. f" | pool: {pool_dir}"
  157. f" | management_base_url: {args.management_base_url}"
  158. f" | quota probe workers: {max(1, int(args.rotate_probe_workers))}"
  159. f" | interval: {interval}s"
  160. )
  161. while True:
  162. started = time.time()
  163. result = rotate_once(
  164. pool_dir=pool_dir,
  165. management_base_url=str(args.management_base_url or "").strip() or DEFAULT_MANAGEMENT_BASE_URL,
  166. cpa_management_key=str(args.management_key or "").strip() or None,
  167. rotate_probe_workers=max(1, int(args.rotate_probe_workers)),
  168. fresh_grace_seconds=max(0, int(args.fresh_grace_seconds)),
  169. cpa_runtime_reconcile_enabled=bool(args.cpa_runtime_reconcile_enabled),
  170. cpa_runtime_reconcile_cooldown_seconds=max(0, int(args.cpa_runtime_reconcile_cooldown_seconds)),
  171. cpa_runtime_reconcile_restart_enabled=bool(args.cpa_runtime_reconcile_restart_enabled),
  172. )
  173. print_rotate_summary(result)
  174. elapsed = time.time() - started
  175. if args.once:
  176. break
  177. time.sleep(max(0, interval - elapsed))
  178. if __name__ == "__main__":
  179. main()