registration.py 127 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188218921902191219221932194219521962197219821992200220122022203220422052206220722082209221022112212221322142215221622172218221922202221222222232224222522262227222822292230223122322233223422352236223722382239224022412242224322442245224622472248224922502251225222532254225522562257225822592260226122622263226422652266226722682269227022712272227322742275227622772278227922802281228222832284228522862287228822892290229122922293229422952296229722982299230023012302230323042305230623072308230923102311231223132314231523162317231823192320232123222323232423252326232723282329233023312332233323342335233623372338233923402341234223432344234523462347234823492350235123522353235423552356235723582359236023612362236323642365236623672368236923702371237223732374237523762377237823792380238123822383238423852386238723882389239023912392239323942395239623972398239924002401240224032404240524062407240824092410241124122413241424152416241724182419242024212422242324242425242624272428242924302431243224332434243524362437243824392440244124422443244424452446244724482449245024512452245324542455245624572458245924602461246224632464246524662467246824692470247124722473247424752476247724782479248024812482248324842485248624872488248924902491249224932494249524962497249824992500250125022503250425052506250725082509251025112512251325142515251625172518251925202521252225232524252525262527252825292530253125322533253425352536253725382539254025412542254325442545254625472548254925502551255225532554255525562557255825592560256125622563256425652566256725682569257025712572257325742575257625772578257925802581258225832584258525862587258825892590259125922593259425952596259725982599260026012602260326042605260626072608260926102611261226132614261526162617261826192620262126222623262426252626262726282629263026312632263326342635263626372638263926402641264226432644264526462647264826492650265126522653265426552656265726582659266026612662266326642665266626672668266926702671267226732674267526762677267826792680268126822683268426852686268726882689269026912692269326942695269626972698269927002701270227032704270527062707
  1. """Registration runtime loops for zhuce6."""
  2. from __future__ import annotations
  3. from collections import deque
  4. from dataclasses import replace
  5. from datetime import datetime
  6. import json
  7. import logging
  8. import os
  9. from pathlib import Path
  10. import random
  11. import sys
  12. import threading
  13. import time
  14. from typing import Any
  15. from urllib.parse import urlsplit, urlunsplit
  16. from core.registry import load_all
  17. from core.settings import AppSettings
  18. from dashboard.api import _count_cpa_files, _fetch_management_auth_files, _is_regular_free_account
  19. from ops.common import create_backend_client, get_management_key
  20. from ops.rotate_runtime import _maybe_reconcile_cpa_runtime
  21. from platforms.chatgpt.fingerprint import build_registration_provenance, infer_proxy_region
  22. from platforms.chatgpt.pool import is_warmup_pending_record, now_iso, update_token_record
  23. from core.chatgpt_flow_runner import run_chatgpt_register_once
  24. DEFAULT_ADD_PHONE_STOPLOSS_WINDOW = 10
  25. DEFAULT_ADD_PHONE_STOPLOSS_THRESHOLD = 3
  26. DEFAULT_ADD_PHONE_STOPLOSS_COOLDOWN_SECONDS = 300
  27. DEFAULT_ADD_PHONE_STOPLOSS_MAX_SUCCESSES = 2
  28. DEFAULT_WAIT_OTP_STOPLOSS_WINDOW = 6
  29. DEFAULT_WAIT_OTP_STOPLOSS_THRESHOLD = 2
  30. DEFAULT_WAIT_OTP_STOPLOSS_COOLDOWN_SECONDS = 300
  31. DEFAULT_WAIT_OTP_LIVE_ABORT_THRESHOLD = 0
  32. DEFAULT_WAIT_OTP_LIVE_ABORT_AGE_SECONDS = 90
  33. DEFAULT_CFMAIL_FRESH_DOMAIN_ATTEMPT_BUDGET = 2
  34. def _classify_token_file(*args, **kwargs): # type: ignore[no-untyped-def]
  35. from ops.scan import classify_token_file
  36. return classify_token_file(*args, **kwargs)
  37. def _compat_main_attr(name: str, default: object) -> object:
  38. main_module = sys.modules.get("main")
  39. if main_module is None:
  40. return default
  41. return getattr(main_module, name, default)
  42. class RegistrationLoop:
  43. """Multi-threaded continuous registration with fallback, target count, and logging."""
  44. def __init__(self, settings: AppSettings) -> None:
  45. self.settings = settings
  46. self._threads: list[threading.Thread] = []
  47. self._stop_event = threading.Event()
  48. self._target_reached = threading.Event()
  49. self._lock = threading.RLock()
  50. self._total_attempts = 0
  51. self._total_success = 0
  52. self._total_warmup_pending = 0
  53. self._total_cpa_sync_success = 0
  54. self._total_cpa_sync_failure = 0
  55. self._total_failure = 0
  56. self._last_error: str | None = None
  57. self._started_at: float | None = None
  58. self._failure_by_stage: dict[str, int] = {}
  59. self._failure_signals: dict[str, int] = {}
  60. self._recent_attempts: deque[dict[str, object]] = deque(maxlen=80)
  61. self._providers: list[str] = []
  62. self._proxy_pool = None
  63. self._logger = self._setup_logger()
  64. self._cfmail_tracker = None
  65. self._cfmail_provisioner = None
  66. self._cfmail_manager: Any = None
  67. self._cfmail_rotation_lock = threading.Lock()
  68. self._cfmail_rotation_pause = threading.Event()
  69. self._cfmail_rotation_pause.set()
  70. self._cpa_management_key_cache: str | None | bool = False
  71. self._cfmail_add_phone_window = max(
  72. 1,
  73. int(str(os.getenv("ZHUCE6_CFMAIL_ADD_PHONE_WINDOW", DEFAULT_ADD_PHONE_STOPLOSS_WINDOW)).strip() or str(DEFAULT_ADD_PHONE_STOPLOSS_WINDOW)),
  74. )
  75. self._cfmail_add_phone_threshold = max(
  76. 1,
  77. int(
  78. str(
  79. os.getenv(
  80. "ZHUCE6_CFMAIL_ADD_PHONE_THRESHOLD",
  81. DEFAULT_ADD_PHONE_STOPLOSS_THRESHOLD,
  82. )
  83. ).strip()
  84. or str(DEFAULT_ADD_PHONE_STOPLOSS_THRESHOLD)
  85. ),
  86. )
  87. self._cfmail_add_phone_cooldown_seconds = max(
  88. 1,
  89. int(
  90. str(
  91. os.getenv(
  92. "ZHUCE6_CFMAIL_ADD_PHONE_COOLDOWN_SECONDS",
  93. DEFAULT_ADD_PHONE_STOPLOSS_COOLDOWN_SECONDS,
  94. )
  95. ).strip()
  96. or str(DEFAULT_ADD_PHONE_STOPLOSS_COOLDOWN_SECONDS)
  97. ),
  98. )
  99. self._cfmail_add_phone_max_successes = max(
  100. 0,
  101. int(
  102. str(
  103. os.getenv(
  104. "ZHUCE6_CFMAIL_ADD_PHONE_MAX_SUCCESSES",
  105. DEFAULT_ADD_PHONE_STOPLOSS_MAX_SUCCESSES,
  106. )
  107. ).strip()
  108. or str(DEFAULT_ADD_PHONE_STOPLOSS_MAX_SUCCESSES)
  109. ),
  110. )
  111. self._cfmail_add_phone_events: dict[str, deque[dict[str, object]]] = {}
  112. self._cfmail_add_phone_state: dict[str, object] = {
  113. "active_domain": "",
  114. "in_cooldown": False,
  115. "cooldown_until": 0.0,
  116. "last_triggered_at": "",
  117. "last_rotation_attempted_at": "",
  118. "last_reason": "",
  119. "last_add_phone_failures": 0,
  120. "last_successes": 0,
  121. "last_window_size": 0,
  122. "last_logged_at": 0.0,
  123. }
  124. self._cfmail_wait_otp_window = max(
  125. 1,
  126. int(str(os.getenv("ZHUCE6_CFMAIL_WAIT_OTP_WINDOW", DEFAULT_WAIT_OTP_STOPLOSS_WINDOW)).strip() or str(DEFAULT_WAIT_OTP_STOPLOSS_WINDOW)),
  127. )
  128. self._cfmail_wait_otp_threshold = max(
  129. 1,
  130. int(
  131. str(
  132. os.getenv(
  133. "ZHUCE6_CFMAIL_WAIT_OTP_THRESHOLD",
  134. DEFAULT_WAIT_OTP_STOPLOSS_THRESHOLD,
  135. )
  136. ).strip()
  137. or str(DEFAULT_WAIT_OTP_STOPLOSS_THRESHOLD)
  138. ),
  139. )
  140. self._cfmail_wait_otp_cooldown_seconds = max(
  141. 0,
  142. int(
  143. str(
  144. os.getenv(
  145. "ZHUCE6_CFMAIL_WAIT_OTP_COOLDOWN_SECONDS",
  146. DEFAULT_WAIT_OTP_STOPLOSS_COOLDOWN_SECONDS,
  147. )
  148. ).strip()
  149. or str(DEFAULT_WAIT_OTP_STOPLOSS_COOLDOWN_SECONDS)
  150. ),
  151. )
  152. self._cfmail_wait_otp_events: dict[str, deque[dict[str, object]]] = {}
  153. self._cfmail_wait_otp_state: dict[str, object] = {
  154. "active_domain": "",
  155. "in_cooldown": False,
  156. "cooldown_until": 0.0,
  157. "last_triggered_at": "",
  158. "last_rotation_attempted_at": "",
  159. "last_reason": "",
  160. "last_no_message_timeouts": 0,
  161. "last_window_size": 0,
  162. "last_logged_at": 0.0,
  163. }
  164. try:
  165. self._cfmail_wait_otp_live_threshold = max(
  166. 0,
  167. int(
  168. str(
  169. os.getenv(
  170. "ZHUCE6_CFMAIL_WAIT_OTP_LIVE_ABORT_THRESHOLD",
  171. DEFAULT_WAIT_OTP_LIVE_ABORT_THRESHOLD,
  172. )
  173. ).strip()
  174. or str(DEFAULT_WAIT_OTP_LIVE_ABORT_THRESHOLD)
  175. ),
  176. )
  177. except Exception:
  178. self._cfmail_wait_otp_live_threshold = DEFAULT_WAIT_OTP_LIVE_ABORT_THRESHOLD
  179. try:
  180. self._cfmail_wait_otp_live_age_seconds = max(
  181. 30,
  182. int(
  183. str(
  184. os.getenv(
  185. "ZHUCE6_CFMAIL_WAIT_OTP_LIVE_ABORT_AGE_SECONDS",
  186. DEFAULT_WAIT_OTP_LIVE_ABORT_AGE_SECONDS,
  187. )
  188. ).strip()
  189. or str(DEFAULT_WAIT_OTP_LIVE_ABORT_AGE_SECONDS)
  190. ),
  191. )
  192. except Exception:
  193. self._cfmail_wait_otp_live_age_seconds = DEFAULT_WAIT_OTP_LIVE_ABORT_AGE_SECONDS
  194. self._cfmail_wait_otp_live_lock = threading.RLock()
  195. self._cfmail_wait_otp_live_progress: dict[str, dict[str, dict[str, object]]] = {}
  196. self._cfmail_canary_state: dict[str, object] = {
  197. "active_domain": "",
  198. "pending": False,
  199. "owner_thread_id": 0,
  200. "attempt_started_at": 0.0,
  201. "last_ready_at": "",
  202. "last_ready_reason": "",
  203. "last_logged_at": 0.0,
  204. }
  205. try:
  206. self._cfmail_start_interval_seconds = max(
  207. 0,
  208. int(str(os.getenv("ZHUCE6_CFMAIL_START_INTERVAL_SECONDS", "8")).strip() or "8"),
  209. )
  210. except Exception:
  211. self._cfmail_start_interval_seconds = 8
  212. try:
  213. self._cfmail_max_inflight = max(
  214. 1,
  215. int(str(os.getenv("ZHUCE6_CFMAIL_MAX_INFLIGHT", "4")).strip() or "4"),
  216. )
  217. except Exception:
  218. self._cfmail_max_inflight = 4
  219. self._cfmail_flow_state: dict[str, object] = {
  220. "inflight_by_thread": {},
  221. "last_started_by_domain": {},
  222. "selected_profile_by_thread": {},
  223. "last_logged_at": 0.0,
  224. }
  225. try:
  226. self._cfmail_active_domain_count = max(
  227. 1,
  228. int(str(os.getenv("ZHUCE6_CFMAIL_ACTIVE_DOMAIN_COUNT", "3")).strip() or "3"),
  229. )
  230. except Exception:
  231. self._cfmail_active_domain_count = 3
  232. try:
  233. self._cfmail_fresh_domain_attempt_budget = max(
  234. 0,
  235. int(
  236. str(
  237. os.getenv(
  238. "ZHUCE6_CFMAIL_FRESH_DOMAIN_ATTEMPT_BUDGET",
  239. DEFAULT_CFMAIL_FRESH_DOMAIN_ATTEMPT_BUDGET,
  240. )
  241. ).strip()
  242. or str(DEFAULT_CFMAIL_FRESH_DOMAIN_ATTEMPT_BUDGET)
  243. ),
  244. )
  245. except Exception:
  246. self._cfmail_fresh_domain_attempt_budget = DEFAULT_CFMAIL_FRESH_DOMAIN_ATTEMPT_BUDGET
  247. self._cfmail_fresh_domain_state: dict[str, object] = {
  248. "active_domain": "",
  249. "completed_attempts": 0,
  250. "mail_seen_attempts": 0,
  251. "successes": 0,
  252. "last_triggered_at": "",
  253. "last_rotation_attempted_at": "",
  254. "last_reason": "",
  255. }
  256. # Solution B: deferred retry queue for add_phone_gate accounts
  257. self._pending_token_queue: list[dict[str, Any]] = []
  258. self._pending_token_lock = threading.Lock()
  259. try:
  260. raw_pending_retry_delay = int(
  261. str(os.getenv("ZHUCE6_PENDING_TOKEN_RETRY_DELAY_SECONDS", "600")).strip() or "600"
  262. )
  263. except Exception:
  264. raw_pending_retry_delay = 600
  265. # add_phone accounts are only useful if the deferred retry happens inside the
  266. # same short observation window as the registration loop. Cap the first retry
  267. # base delay so a large env value cannot postpone every retry past 5 minutes.
  268. self._pending_token_retry_delay_seconds = max(60, min(raw_pending_retry_delay, 60))
  269. try:
  270. self._pending_token_max_retries = max(
  271. 1,
  272. int(str(os.getenv("ZHUCE6_PENDING_TOKEN_MAX_RETRIES", "3")).strip() or "3"),
  273. )
  274. except Exception:
  275. self._pending_token_max_retries = 3
  276. self._pending_token_total_enqueued = 0
  277. self._pending_token_total_success = 0
  278. self._pending_token_total_failed = 0
  279. self._cfmail_replenish_thread: threading.Thread | None = None
  280. self._cfmail_replenish_reason = ""
  281. def _write_runtime_state(self) -> None:
  282. state_file = Path(self.settings.runtime_state_file)
  283. try:
  284. state_file.parent.mkdir(parents=True, exist_ok=True)
  285. payload = {
  286. "updated_at": datetime.now().isoformat(timespec="seconds"),
  287. "register_snapshot": self.snapshot(),
  288. "proxy_pool": self._proxy_pool_snapshot(),
  289. }
  290. tmp_file = state_file.with_name(
  291. f"{state_file.name}.{os.getpid()}.{threading.get_ident()}.tmp"
  292. )
  293. tmp_file.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8")
  294. tmp_file.replace(state_file)
  295. except Exception as exc:
  296. self._log(f"[zhuce6:register] runtime state write failed: {exc}")
  297. def _cfmail_canary_snapshot(self) -> dict[str, object]:
  298. return {
  299. "active_domain": "",
  300. "pending": False,
  301. "owner_thread_id": 0,
  302. "attempt_started_at": 0.0,
  303. "last_ready_at": "",
  304. "last_ready_reason": "disabled",
  305. }
  306. def _arm_cfmail_canary(self, domain: str, *, pending: bool = True) -> None:
  307. del domain, pending
  308. return
  309. def _mark_cfmail_canary_ready(self, domain: str, *, reason: str) -> None:
  310. del domain, reason
  311. return
  312. def _update_cfmail_canary_after_result(self, *, thread_id: int, result: dict[str, object]) -> None:
  313. del thread_id, result
  314. return
  315. def _wait_if_cfmail_canary_pending(self, thread_id: int, provider: str) -> bool:
  316. del thread_id, provider
  317. return False
  318. def _release_cfmail_flow_slot(self, thread_id: int) -> None:
  319. with self._lock:
  320. inflight_by_thread = self._cfmail_flow_state.setdefault("inflight_by_thread", {})
  321. if isinstance(inflight_by_thread, dict):
  322. inflight_by_thread.pop(thread_id, None)
  323. selected_profile_by_thread = self._cfmail_flow_state.setdefault("selected_profile_by_thread", {})
  324. if isinstance(selected_profile_by_thread, dict):
  325. selected_profile_by_thread.pop(thread_id, None)
  326. def _cfmail_active_domain_set(self) -> set[str]:
  327. domains = {
  328. str(item.get("domain") or "").strip().lower()
  329. for item in self._current_cfmail_active_accounts()
  330. if str(item.get("domain") or "").strip()
  331. }
  332. if domains:
  333. return domains
  334. current_domain = str(self._current_cfmail_active_domain() or "").strip().lower()
  335. return {current_domain} if current_domain else set()
  336. def _is_cfmail_active_domain(self, domain: str) -> bool:
  337. domain_key = str(domain or "").strip().lower()
  338. if not domain_key:
  339. return False
  340. return domain_key in self._cfmail_active_domain_set()
  341. def _schedule_cfmail_domain_pool_replenish(self, *, trigger_thread_id: int, reason: str) -> None:
  342. if self._cfmail_provisioner is None:
  343. return
  344. with self._lock:
  345. thread = self._cfmail_replenish_thread
  346. if thread is not None and thread.is_alive():
  347. return
  348. self._cfmail_replenish_reason = str(reason or "").strip()
  349. thread = threading.Thread(
  350. target=self._cfmail_domain_pool_replenish_worker,
  351. kwargs={
  352. "trigger_thread_id": trigger_thread_id,
  353. "reason": self._cfmail_replenish_reason,
  354. },
  355. daemon=True,
  356. name="zhuce6-cfmail-replenish",
  357. )
  358. self._cfmail_replenish_thread = thread
  359. thread.start()
  360. def _ensure_cfmail_domain_pool_target(self, *, trigger_thread_id: int, reason: str) -> None:
  361. if self._cfmail_provisioner is None:
  362. return
  363. if len(self._current_cfmail_active_accounts()) >= self._cfmail_active_domain_count:
  364. return
  365. self._schedule_cfmail_domain_pool_replenish(
  366. trigger_thread_id=trigger_thread_id,
  367. reason=reason,
  368. )
  369. def _cfmail_domain_pool_replenish_worker(self, *, trigger_thread_id: int, reason: str) -> None:
  370. provisioner = self._cfmail_provisioner
  371. if provisioner is None:
  372. return
  373. while not self._stop_event.is_set():
  374. active_accounts = self._current_cfmail_active_accounts()
  375. if len(active_accounts) >= self._cfmail_active_domain_count:
  376. return
  377. result = provisioner.provision_additional_domain()
  378. if not result.success:
  379. self._log(
  380. f"[zhuce6:register] [thread-{trigger_thread_id}] [cfmail] replenish failed after {reason}: "
  381. f"{result.error}"
  382. )
  383. if self._stop_event.wait(5.0):
  384. return
  385. continue
  386. self._reload_cfmail_manager_after_rotation()
  387. if self._cfmail_tracker is not None:
  388. try:
  389. self._cfmail_tracker.mark_rotation_completed("", result.new_domain)
  390. except Exception:
  391. pass
  392. self._log(
  393. f"[zhuce6:register] [thread-{trigger_thread_id}] [cfmail] replenished domain pool "
  394. f"after {reason}: +{result.new_domain}"
  395. )
  396. def _clear_cfmail_domain_state(self, domain: str) -> None:
  397. domain_key = str(domain or "").strip().lower()
  398. if not domain_key:
  399. return
  400. with self._lock:
  401. self._cfmail_add_phone_events.pop(domain_key, None)
  402. if str(self._cfmail_add_phone_state.get("active_domain") or "").strip().lower() == domain_key:
  403. self._cfmail_add_phone_state = {
  404. "active_domain": "",
  405. "in_cooldown": False,
  406. "cooldown_until": 0.0,
  407. "last_triggered_at": "",
  408. "last_rotation_attempted_at": "",
  409. "last_reason": "",
  410. "last_add_phone_failures": 0,
  411. "last_successes": 0,
  412. "last_window_size": 0,
  413. "last_logged_at": 0.0,
  414. }
  415. self._cfmail_wait_otp_events.pop(domain_key, None)
  416. if str(self._cfmail_wait_otp_state.get("active_domain") or "").strip().lower() == domain_key:
  417. self._cfmail_wait_otp_state = {
  418. "active_domain": "",
  419. "in_cooldown": False,
  420. "cooldown_until": 0.0,
  421. "last_triggered_at": "",
  422. "last_rotation_attempted_at": "",
  423. "last_reason": "",
  424. "last_no_message_timeouts": 0,
  425. "last_successes": 0,
  426. "last_message_seen": 0,
  427. "last_window_size": 0,
  428. "last_logged_at": 0.0,
  429. }
  430. self._cfmail_wait_otp_live_progress.pop(domain_key, None)
  431. if str(self._cfmail_fresh_domain_state.get("active_domain") or "").strip().lower() == domain_key:
  432. self._cfmail_fresh_domain_state = {
  433. "active_domain": "",
  434. "completed_attempts": 0,
  435. "mail_seen_attempts": 0,
  436. "successes": 0,
  437. "last_triggered_at": "",
  438. "last_rotation_attempted_at": "",
  439. "last_reason": "",
  440. }
  441. def _replace_cfmail_domain(
  442. self,
  443. *,
  444. thread_id: int,
  445. domain: str,
  446. reason_label: str,
  447. ) -> bool:
  448. domain_key = str(domain or "").strip().lower()
  449. if not domain_key or self._cfmail_provisioner is None:
  450. return False
  451. active_accounts = self._current_cfmail_active_accounts()
  452. if not active_accounts:
  453. active_accounts = [{"name": "", "domain": domain_key}]
  454. if not any(item["domain"] == domain_key for item in active_accounts):
  455. self._clear_cfmail_domain_state(domain_key)
  456. return False
  457. if not self._cfmail_rotation_lock.acquire(blocking=False):
  458. return False
  459. self._cfmail_rotation_pause.clear()
  460. try:
  461. if self._cfmail_tracker is not None:
  462. self._cfmail_tracker.mark_rotation_started(domain_key, reason_label)
  463. if len(active_accounts) <= 1:
  464. provision_result = self._cfmail_provisioner.rotate_active_domain()
  465. if not provision_result.success:
  466. if self._cfmail_tracker is not None:
  467. self._cfmail_tracker.mark_rotation_failed(domain_key, provision_result.error)
  468. self._log(
  469. f"[zhuce6:register] [thread-{thread_id}] [cfmail] {reason_label} rotation failed: "
  470. f"{provision_result.error}"
  471. )
  472. return False
  473. self._reload_cfmail_manager_after_rotation()
  474. self._clear_cfmail_domain_state(provision_result.old_domain or domain_key)
  475. self._reset_cfmail_add_phone_stoploss(provision_result.new_domain)
  476. self._reset_cfmail_wait_otp_stoploss(provision_result.new_domain)
  477. self._reset_cfmail_fresh_domain_budget(provision_result.new_domain)
  478. if self._cfmail_tracker is not None:
  479. self._cfmail_tracker.mark_rotation_completed(
  480. provision_result.old_domain,
  481. provision_result.new_domain,
  482. )
  483. self._log(
  484. f"[zhuce6:register] [thread-{thread_id}] [cfmail] {reason_label} rotation completed: "
  485. f"{provision_result.old_domain} -> {provision_result.new_domain}"
  486. )
  487. return True
  488. retire_result = self._cfmail_provisioner.retire_domain(domain_key)
  489. if not retire_result.success:
  490. if self._cfmail_tracker is not None:
  491. self._cfmail_tracker.mark_rotation_failed(domain_key, retire_result.error)
  492. self._log(
  493. f"[zhuce6:register] [thread-{thread_id}] [cfmail] {reason_label} retire failed: "
  494. f"{retire_result.error}"
  495. )
  496. return False
  497. self._reload_cfmail_manager_after_rotation()
  498. self._clear_cfmail_domain_state(domain_key)
  499. self._log(
  500. f"[zhuce6:register] [thread-{thread_id}] [cfmail] retired domain {domain_key} "
  501. f"because {reason_label}; scheduling replenish"
  502. )
  503. self._schedule_cfmail_domain_pool_replenish(
  504. trigger_thread_id=thread_id,
  505. reason=reason_label,
  506. )
  507. return True
  508. finally:
  509. self._cfmail_rotation_pause.set()
  510. self._cfmail_rotation_lock.release()
  511. def _wait_if_cfmail_flow_throttled(self, thread_id: int, provider: str) -> bool:
  512. if provider != "cfmail":
  513. return False
  514. self._ensure_cfmail_domain_pool_target(
  515. trigger_thread_id=thread_id,
  516. reason="usable domain pool below target",
  517. )
  518. active_accounts = self._current_cfmail_active_accounts()
  519. if not active_accounts:
  520. return False
  521. with self._lock:
  522. inflight_by_thread = self._cfmail_flow_state.setdefault("inflight_by_thread", {})
  523. if not isinstance(inflight_by_thread, dict):
  524. inflight_by_thread = {}
  525. self._cfmail_flow_state["inflight_by_thread"] = inflight_by_thread
  526. last_started_by_domain = self._cfmail_flow_state.setdefault("last_started_by_domain", {})
  527. if not isinstance(last_started_by_domain, dict):
  528. last_started_by_domain = {}
  529. self._cfmail_flow_state["last_started_by_domain"] = last_started_by_domain
  530. selected_profile_by_thread = self._cfmail_flow_state.setdefault("selected_profile_by_thread", {})
  531. if not isinstance(selected_profile_by_thread, dict):
  532. selected_profile_by_thread = {}
  533. self._cfmail_flow_state["selected_profile_by_thread"] = selected_profile_by_thread
  534. tracked = inflight_by_thread.get(thread_id)
  535. if isinstance(tracked, dict):
  536. tracked_domain = str(tracked.get("domain") or "").strip().lower()
  537. tracked_profile = str(tracked.get("profile_name") or "").strip()
  538. if tracked_domain and tracked_profile:
  539. selected_profile_by_thread[thread_id] = tracked_profile
  540. return False
  541. now = time.time()
  542. candidates: list[tuple[int, float, str, str]] = []
  543. min_wait_seconds = 2.0
  544. wait_domain = ""
  545. wait_reason = "inflight_limit"
  546. for account in active_accounts:
  547. domain = account["domain"]
  548. profile_name = account["name"]
  549. active_inflight = sum(
  550. 1
  551. for value in inflight_by_thread.values()
  552. if isinstance(value, dict)
  553. and str(value.get("domain") or "").strip().lower() == domain
  554. )
  555. if active_inflight >= self._cfmail_max_inflight:
  556. wait_domain = wait_domain or domain
  557. continue
  558. last_started_at = float(last_started_by_domain.get(domain) or 0.0)
  559. remaining = 0.0
  560. if self._cfmail_start_interval_seconds > 0 and last_started_at > 0.0:
  561. remaining = self._cfmail_start_interval_seconds - (now - last_started_at)
  562. if remaining > 0.0:
  563. wait_domain = wait_domain or domain
  564. wait_reason = "start_interval"
  565. min_wait_seconds = min(min_wait_seconds, min(2.0, max(0.5, remaining)))
  566. continue
  567. candidates.append((active_inflight, last_started_at, profile_name, domain))
  568. if candidates:
  569. candidates.sort(key=lambda item: (item[0], item[1], item[3], item[2]))
  570. _active_inflight, _last_started_at, profile_name, domain = candidates[0]
  571. inflight_by_thread[thread_id] = {
  572. "domain": domain,
  573. "profile_name": profile_name,
  574. "started_at": now,
  575. }
  576. last_started_by_domain[domain] = now
  577. selected_profile_by_thread[thread_id] = profile_name
  578. return False
  579. last_logged_at = float(self._cfmail_flow_state.get("last_logged_at") or 0.0)
  580. should_log = now - last_logged_at >= 15.0
  581. if should_log:
  582. self._cfmail_flow_state["last_logged_at"] = now
  583. if should_log:
  584. if wait_reason == "inflight_limit":
  585. self._log(
  586. f"[zhuce6:register] [thread-{thread_id}] [cfmail] flow throttle for {wait_domain or '-'}: "
  587. f"inflight_limit reached on all active domains"
  588. )
  589. else:
  590. self._log(
  591. f"[zhuce6:register] [thread-{thread_id}] [cfmail] flow throttle for {wait_domain or '-'}: "
  592. f"start_interval_remaining={min_wait_seconds:.1f}s"
  593. )
  594. self._stop_event.wait(min_wait_seconds)
  595. return True
  596. def _proxy_pool_snapshot(self) -> dict[str, object]:
  597. pool = self._proxy_pool
  598. nodes: list[dict[str, object]] = []
  599. snapshot_error: str | None = None
  600. if pool is not None:
  601. try:
  602. snapshot = pool.snapshot()
  603. except Exception as exc:
  604. snapshot_error = str(exc)
  605. else:
  606. if isinstance(snapshot, list):
  607. nodes = [item for item in snapshot if isinstance(item, dict)]
  608. return {
  609. "configured": bool(self.settings.proxy_pool_configured or pool is not None),
  610. "enabled": pool is not None,
  611. "snapshot_error": snapshot_error,
  612. "node_count": len(nodes),
  613. "in_use_count": sum(1 for item in nodes if item.get("in_use")),
  614. "disabled_count": sum(1 for item in nodes if item.get("disabled")),
  615. "nodes": nodes,
  616. }
  617. def _setup_logger(self) -> Any:
  618. import logging
  619. from logging.handlers import RotatingFileHandler
  620. logger = logging.getLogger("zhuce6.register")
  621. logger.setLevel(logging.INFO)
  622. logger.propagate = False
  623. if not logger.handlers:
  624. console = logging.StreamHandler()
  625. console.setFormatter(logging.Formatter("%(message)s"))
  626. logger.addHandler(console)
  627. if self.settings.register_log_file:
  628. fh = RotatingFileHandler(
  629. self.settings.register_log_file,
  630. maxBytes=2 * 1024 * 1024, # 2MB
  631. backupCount=5,
  632. encoding="utf-8",
  633. )
  634. fh.setFormatter(logging.Formatter("%(asctime)s %(message)s", datefmt="%Y-%m-%d %H:%M:%S"))
  635. logger.addHandler(fh)
  636. return logger
  637. def _log(self, msg: str) -> None:
  638. self._logger.info(msg)
  639. def _record_attempt(
  640. self,
  641. *,
  642. success: bool,
  643. stage: str,
  644. error_message: str,
  645. metadata: dict[str, object] | None = None,
  646. proxy_key: str = "",
  647. email: str = "",
  648. ) -> None:
  649. meta = metadata if isinstance(metadata, dict) else {}
  650. stage_key = str(stage or "?").strip() or "?"
  651. signal = self._classify_failure_signal(stage=stage_key, metadata=meta)
  652. timestamp = datetime.now().isoformat(timespec="seconds")
  653. event = {
  654. "timestamp": timestamp,
  655. "success": success,
  656. "stage": stage_key if not success else "completed",
  657. "signal": signal,
  658. "error_message": str(error_message or "").strip(),
  659. "email_domain": str(meta.get("email_domain") or "").strip(),
  660. "post_create_gate": str(meta.get("post_create_gate") or "").strip(),
  661. "create_account_error_code": str(meta.get("create_account_error_code") or "").strip(),
  662. "signup_error_code": str(meta.get("signup_error_code") or "").strip(),
  663. "signup_http_status": meta.get("signup_http_status"),
  664. "mailbox_error_kind": str(meta.get("mailbox_error_kind") or "").strip(),
  665. "mailbox_error_stage": str(meta.get("mailbox_error_stage") or "").strip(),
  666. "proxy_key": proxy_key,
  667. "email": email,
  668. }
  669. self._recent_attempts.append(event)
  670. if success or stage_key == "warmup_pending":
  671. return
  672. self._failure_by_stage[stage_key] = self._failure_by_stage.get(stage_key, 0) + 1
  673. if signal:
  674. self._failure_signals[signal] = self._failure_signals.get(signal, 0) + 1
  675. def _classify_failure_signal(self, *, stage: str, metadata: dict[str, object]) -> str:
  676. code = str(metadata.get("create_account_error_code") or "").strip().lower()
  677. signup_code = str(metadata.get("signup_error_code") or "").strip().lower()
  678. post_gate = str(metadata.get("post_create_gate") or "").strip().lower()
  679. if stage == "cpa_sync":
  680. return "cpa_sync_failed"
  681. if stage == "signup" and signup_code:
  682. return signup_code
  683. if stage == "add_phone_gate" or post_gate == "add_phone":
  684. return "add_phone_gate"
  685. if stage == "create_account" and code == "user_already_exists":
  686. return "mailbox_reused"
  687. if stage == "create_account" and code in {"registration_disallowed", "unsupported_email"}:
  688. return code
  689. if stage == "mailbox":
  690. provider = str(metadata.get("mail_provider") or "").strip().lower()
  691. if provider == "cfmail":
  692. mailbox_error_kind = str(metadata.get("mailbox_error_kind") or "").strip().lower()
  693. mailbox_error_stage = str(metadata.get("mailbox_error_stage") or "").strip().lower()
  694. if mailbox_error_stage == "create_email":
  695. mailbox_error_stage = "create"
  696. elif mailbox_error_stage == "fetch_email":
  697. mailbox_error_stage = "fetch"
  698. if mailbox_error_kind == "transport_error":
  699. suffix = mailbox_error_stage or "backend"
  700. return f"mailbox_{suffix}_transport_error"
  701. if mailbox_error_kind == "provider_error":
  702. suffix = mailbox_error_stage or "backend"
  703. return f"mailbox_{suffix}_provider_error"
  704. return "mailbox_backend_failure"
  705. return "mailbox_failure"
  706. return ""
  707. def _recent_failure_hotspots(self, limit: int = 5) -> list[dict[str, object]]:
  708. return self._recent_failure_hotspots_from_attempts(self._recent_attempts, limit=limit)
  709. def _recent_failure_hotspots_from_attempts(
  710. self,
  711. attempts: list[dict[str, object]] | deque[dict[str, object]],
  712. *,
  713. limit: int = 5,
  714. ) -> list[dict[str, object]]:
  715. counts: dict[tuple[str, str], int] = {}
  716. for item in attempts:
  717. if item.get("success"):
  718. continue
  719. stage = str(item.get("stage") or "?").strip() or "?"
  720. signal = str(item.get("signal") or "").strip()
  721. key = (signal or stage, stage)
  722. counts[key] = counts.get(key, 0) + 1
  723. ordered = sorted(counts.items(), key=lambda kv: (-kv[1], kv[0][0], kv[0][1]))
  724. return [
  725. {"key": key, "stage": stage, "count": count}
  726. for (key, stage), count in ordered[:limit]
  727. ]
  728. def _failure_counts_from_attempts(
  729. self,
  730. attempts: list[dict[str, object]] | deque[dict[str, object]],
  731. ) -> tuple[dict[str, int], dict[str, int]]:
  732. failure_by_stage: dict[str, int] = {}
  733. failure_signals: dict[str, int] = {}
  734. for item in attempts:
  735. if item.get("success"):
  736. continue
  737. stage = str(item.get("stage") or "?").strip() or "?"
  738. failure_by_stage[stage] = failure_by_stage.get(stage, 0) + 1
  739. signal = str(item.get("signal") or "").strip()
  740. if signal:
  741. failure_signals[signal] = failure_signals.get(signal, 0) + 1
  742. return (
  743. dict(sorted(failure_by_stage.items(), key=lambda item: (-item[1], item[0]))),
  744. dict(sorted(failure_signals.items(), key=lambda item: (-item[1], item[0]))),
  745. )
  746. def _active_domain_attempts(
  747. self,
  748. recent_attempts: list[dict[str, object]],
  749. active_domain: str,
  750. ) -> list[dict[str, object]]:
  751. domain = str(active_domain or "").strip().lower()
  752. if not domain:
  753. return list(recent_attempts)
  754. filtered = [
  755. item
  756. for item in recent_attempts
  757. if str(item.get("email_domain") or "").strip().lower() == domain
  758. ]
  759. return filtered
  760. def _infer_active_domain(
  761. self,
  762. recent_attempts: list[dict[str, object]],
  763. cfmail_rotation: dict[str, object] | None,
  764. stoploss: dict[str, object],
  765. ) -> str:
  766. if isinstance(cfmail_rotation, dict):
  767. domain = str(cfmail_rotation.get("active_domain") or "").strip().lower()
  768. if domain:
  769. return domain
  770. domain = str(stoploss.get("active_domain") or "").strip().lower()
  771. if domain:
  772. return domain
  773. for item in reversed(recent_attempts):
  774. domain = str(item.get("email_domain") or "").strip().lower()
  775. if domain:
  776. return domain
  777. return ""
  778. def _extract_email_domain(self, result: dict[str, object]) -> str:
  779. metadata = result.get("metadata") if isinstance(result.get("metadata"), dict) else {}
  780. domain = str(metadata.get("email_domain") or "").strip().lower()
  781. if domain:
  782. return domain
  783. email = str(result.get("email") or "").strip().lower()
  784. if "@" not in email:
  785. return ""
  786. return email.rsplit("@", 1)[-1].strip().lower()
  787. def _update_cfmail_add_phone_stoploss(self, result: dict[str, object]) -> None:
  788. metadata = result.get("metadata") if isinstance(result.get("metadata"), dict) else {}
  789. provider = str(metadata.get("mail_provider") or result.get("mail_provider") or "").strip().lower()
  790. if provider not in {"", "cfmail"}:
  791. return
  792. domain = self._extract_email_domain(result)
  793. if not domain:
  794. return
  795. if not self._is_cfmail_active_domain(domain):
  796. return
  797. success = bool(result.get("success"))
  798. stage = str(result.get("stage") or "").strip().lower()
  799. post_gate = str(metadata.get("post_create_gate") or "").strip().lower()
  800. is_add_phone = stage == "add_phone_gate" or post_gate == "add_phone"
  801. with self._lock:
  802. events = self._cfmail_add_phone_events.setdefault(
  803. domain,
  804. deque(maxlen=self._cfmail_add_phone_window),
  805. )
  806. events.append(
  807. {
  808. "success": success,
  809. "is_add_phone": is_add_phone,
  810. }
  811. )
  812. state = self._cfmail_add_phone_state
  813. state["active_domain"] = domain
  814. if self._cfmail_add_phone_cooldown_seconds <= 0:
  815. state["in_cooldown"] = False
  816. state["cooldown_until"] = 0.0
  817. return
  818. cooldown_until = float(state.get("cooldown_until") or 0.0)
  819. if time.time() < cooldown_until:
  820. state["in_cooldown"] = True
  821. return
  822. state["in_cooldown"] = False
  823. if len(events) < self._cfmail_add_phone_window:
  824. return
  825. add_phone_failures = sum(1 for item in events if item.get("is_add_phone"))
  826. successes = sum(1 for item in events if item.get("success"))
  827. if (
  828. add_phone_failures >= self._cfmail_add_phone_threshold
  829. and successes <= self._cfmail_add_phone_max_successes
  830. ):
  831. state["in_cooldown"] = True
  832. state["cooldown_until"] = time.time() + self._cfmail_add_phone_cooldown_seconds
  833. state["last_triggered_at"] = datetime.now().isoformat(timespec="seconds")
  834. state["last_rotation_attempted_at"] = ""
  835. state["last_reason"] = "add_phone threshold reached"
  836. state["last_add_phone_failures"] = add_phone_failures
  837. state["last_successes"] = successes
  838. state["last_window_size"] = len(events)
  839. self._log(
  840. f"[zhuce6:register] [cfmail] add_phone stoploss activated for {domain} "
  841. f"(add_phone_failures={add_phone_failures}, successes={successes}, window={len(events)})"
  842. )
  843. def _is_cfmail_wait_otp_no_message_timeout(self, result: dict[str, object]) -> bool:
  844. metadata = result.get("metadata") if isinstance(result.get("metadata"), dict) else {}
  845. provider = str(metadata.get("mail_provider") or result.get("mail_provider") or "").strip().lower()
  846. if provider not in {"", "cfmail"}:
  847. return False
  848. stage = str(result.get("stage") or "").strip().lower()
  849. if stage != "wait_otp":
  850. return False
  851. failure_reason = str(metadata.get("otp_wait_failure_reason") or "").strip().lower()
  852. if failure_reason:
  853. return failure_reason == "mailbox_timeout_no_message"
  854. try:
  855. return int(metadata.get("otp_mailbox_message_scan_count") or 0) <= 0
  856. except Exception:
  857. return False
  858. def _is_cfmail_invalid_domain_mailbox_failure(self, result: dict[str, object]) -> bool:
  859. metadata = result.get("metadata") if isinstance(result.get("metadata"), dict) else {}
  860. if str(result.get("stage") or "").strip().lower() != "mailbox":
  861. return False
  862. provider = str(metadata.get("mail_provider") or result.get("mail_provider") or "").strip().lower()
  863. if provider not in {"", "cfmail"}:
  864. return False
  865. haystacks = [str(result.get("error_message") or "")]
  866. haystacks.extend(str(item or "") for item in (result.get("logs") or []))
  867. text = "\n".join(haystacks).lower()
  868. return "invalid domain" in text or "无效的域名" in text
  869. def _on_cfmail_wait_progress(self, account: object, diagnostics: dict[str, object]) -> None:
  870. if self._cfmail_wait_otp_live_threshold <= 0:
  871. return
  872. email = str(getattr(account, "email", "") or "").strip().lower()
  873. extra = getattr(account, "extra", {}) or {}
  874. domain = str(extra.get("email_domain") or "").strip().lower()
  875. if not domain and "@" in email:
  876. domain = email.rsplit("@", 1)[-1].strip().lower()
  877. if not domain:
  878. return
  879. with self._cfmail_wait_otp_live_lock:
  880. now = time.time()
  881. for tracked_domain in list(self._cfmail_wait_otp_live_progress.keys()):
  882. entries = self._cfmail_wait_otp_live_progress.get(tracked_domain) or {}
  883. fresh_entries = {
  884. key: value
  885. for key, value in entries.items()
  886. if now - float(value.get("updated_at") or 0.0) <= 15.0
  887. }
  888. if fresh_entries:
  889. self._cfmail_wait_otp_live_progress[tracked_domain] = fresh_entries
  890. else:
  891. self._cfmail_wait_otp_live_progress.pop(tracked_domain, None)
  892. if not self._is_cfmail_active_domain(domain):
  893. self._cfmail_wait_otp_live_progress.pop(domain, None)
  894. return
  895. key = str(getattr(account, "account_id", "") or email or id(account))
  896. domain_entries = self._cfmail_wait_otp_live_progress.setdefault(domain, {})
  897. domain_entries[key] = {
  898. "scan_count": int(diagnostics.get("message_scan_count") or 0),
  899. "elapsed_seconds": float(diagnostics.get("elapsed_seconds") or 0.0),
  900. "updated_at": now,
  901. }
  902. if int(diagnostics.get("message_scan_count") or 0) > 0:
  903. self._mark_cfmail_canary_ready(domain, reason="live_mailbox_message_seen")
  904. state = self._cfmail_wait_otp_state
  905. if bool(state.get("in_cooldown")) and str(state.get("active_domain") or "").strip().lower() == domain:
  906. return
  907. stalled = [
  908. value
  909. for value in domain_entries.values()
  910. if int(value.get("scan_count") or 0) <= 0
  911. and float(value.get("elapsed_seconds") or 0.0) >= self._cfmail_wait_otp_live_age_seconds
  912. ]
  913. if len(stalled) < self._cfmail_wait_otp_live_threshold:
  914. return
  915. with self._lock:
  916. state = self._cfmail_wait_otp_state
  917. if bool(state.get("in_cooldown")) and str(state.get("active_domain") or "").strip().lower() == domain:
  918. return
  919. self._cfmail_wait_otp_state = {
  920. "active_domain": domain,
  921. "in_cooldown": True,
  922. "cooldown_until": time.time() + self._cfmail_wait_otp_cooldown_seconds,
  923. "last_triggered_at": now_iso(),
  924. "last_rotation_attempted_at": "",
  925. "last_reason": "live wait_otp no-message threshold reached",
  926. "last_no_message_timeouts": len(stalled),
  927. "last_window_size": len(stalled),
  928. "last_logged_at": 0.0,
  929. }
  930. self._log(
  931. f"[zhuce6:register] [cfmail] wait_otp live stoploss activated for {domain} "
  932. f"(stalled_waits={len(stalled)}, age>={self._cfmail_wait_otp_live_age_seconds}s)"
  933. )
  934. def _update_cfmail_wait_otp_stoploss(self, result: dict[str, object]) -> None:
  935. metadata = result.get("metadata") if isinstance(result.get("metadata"), dict) else {}
  936. provider = str(metadata.get("mail_provider") or result.get("mail_provider") or "").strip().lower()
  937. if provider not in {"", "cfmail"}:
  938. return
  939. domain = self._extract_email_domain(result)
  940. if not domain:
  941. return
  942. if not self._is_cfmail_active_domain(domain):
  943. return
  944. is_no_message_timeout = self._is_cfmail_wait_otp_no_message_timeout(result)
  945. try:
  946. message_scan_count = int(metadata.get("otp_mailbox_message_scan_count") or 0)
  947. except Exception:
  948. message_scan_count = 0
  949. has_message_seen = message_scan_count > 0
  950. with self._lock:
  951. events = self._cfmail_wait_otp_events.setdefault(
  952. domain,
  953. deque(maxlen=self._cfmail_wait_otp_window),
  954. )
  955. events.append(
  956. {
  957. "success": bool(result.get("success")),
  958. "is_no_message_timeout": is_no_message_timeout,
  959. "has_message_seen": has_message_seen,
  960. }
  961. )
  962. state = self._cfmail_wait_otp_state
  963. state["active_domain"] = domain
  964. if self._cfmail_wait_otp_cooldown_seconds <= 0:
  965. state["in_cooldown"] = False
  966. state["cooldown_until"] = 0.0
  967. return
  968. cooldown_until = float(state.get("cooldown_until") or 0.0)
  969. if time.time() < cooldown_until:
  970. state["in_cooldown"] = True
  971. return
  972. state["in_cooldown"] = False
  973. if len(events) < self._cfmail_wait_otp_window:
  974. return
  975. no_message_timeouts = sum(1 for item in events if item.get("is_no_message_timeout"))
  976. successes = sum(1 for item in events if item.get("success"))
  977. message_seen = sum(1 for item in events if item.get("has_message_seen"))
  978. if (
  979. no_message_timeouts >= self._cfmail_wait_otp_threshold
  980. and successes <= 0
  981. and message_seen <= 0
  982. ):
  983. state["in_cooldown"] = True
  984. state["cooldown_until"] = time.time() + self._cfmail_wait_otp_cooldown_seconds
  985. state["last_triggered_at"] = datetime.now().isoformat(timespec="seconds")
  986. state["last_rotation_attempted_at"] = ""
  987. state["last_reason"] = "wait_otp no-message threshold reached"
  988. state["last_no_message_timeouts"] = no_message_timeouts
  989. state["last_successes"] = successes
  990. state["last_message_seen"] = message_seen
  991. state["last_window_size"] = len(events)
  992. self._log(
  993. f"[zhuce6:register] [cfmail] wait_otp stoploss activated for {domain} "
  994. f"(no_message_timeouts={no_message_timeouts}, successes={successes}, "
  995. f"message_seen={message_seen}, window={len(events)})"
  996. )
  997. def _update_cfmail_fresh_domain_budget(self, result: dict[str, object]) -> None:
  998. if self._cfmail_fresh_domain_attempt_budget <= 0:
  999. return
  1000. metadata = result.get("metadata") if isinstance(result.get("metadata"), dict) else {}
  1001. provider = str(metadata.get("mail_provider") or result.get("mail_provider") or "").strip().lower()
  1002. if provider not in {"", "cfmail"}:
  1003. return
  1004. domain = self._extract_email_domain(result)
  1005. if not domain:
  1006. return
  1007. if not self._is_cfmail_active_domain(domain):
  1008. return
  1009. try:
  1010. message_scan_count = int(metadata.get("otp_mailbox_message_scan_count") or 0)
  1011. except Exception:
  1012. message_scan_count = 0
  1013. success = bool(result.get("success"))
  1014. with self._lock:
  1015. state = self._cfmail_fresh_domain_state
  1016. tracked_domain = str(state.get("active_domain") or "").strip().lower()
  1017. if tracked_domain != domain:
  1018. self._cfmail_fresh_domain_state = {
  1019. "active_domain": domain,
  1020. "completed_attempts": 0,
  1021. "mail_seen_attempts": 0,
  1022. "successes": 0,
  1023. "last_triggered_at": "",
  1024. "last_rotation_attempted_at": "",
  1025. "last_reason": "",
  1026. }
  1027. state = self._cfmail_fresh_domain_state
  1028. state["completed_attempts"] = int(state.get("completed_attempts") or 0) + 1
  1029. if message_scan_count > 0:
  1030. state["mail_seen_attempts"] = int(state.get("mail_seen_attempts") or 0) + 1
  1031. if success:
  1032. state["successes"] = int(state.get("successes") or 0) + 1
  1033. if (
  1034. int(state.get("mail_seen_attempts") or 0) > 0
  1035. and int(state.get("completed_attempts") or 0) >= self._cfmail_fresh_domain_attempt_budget
  1036. and not str(state.get("last_rotation_attempted_at") or "").strip()
  1037. ):
  1038. state["last_triggered_at"] = datetime.now().isoformat(timespec="seconds")
  1039. state["last_reason"] = "fresh_domain_attempt_budget_reached"
  1040. self._log(
  1041. f"[zhuce6:register] [cfmail] fresh-domain budget reached for {domain} "
  1042. f"(completed_attempts={int(state.get('completed_attempts') or 0)}, "
  1043. f"mail_seen_attempts={int(state.get('mail_seen_attempts') or 0)}, "
  1044. f"budget={self._cfmail_fresh_domain_attempt_budget})"
  1045. )
  1046. def _cfmail_add_phone_stoploss_snapshot(self) -> dict[str, object]:
  1047. with self._lock:
  1048. state = dict(self._cfmail_add_phone_state)
  1049. if self._cfmail_add_phone_cooldown_seconds <= 0:
  1050. return {
  1051. "active_domain": str(state.get("active_domain") or ""),
  1052. "in_cooldown": False,
  1053. "cooldown_remaining_seconds": 0,
  1054. "last_triggered_at": str(state.get("last_triggered_at") or ""),
  1055. "last_reason": str(state.get("last_reason") or ""),
  1056. "last_add_phone_failures": int(state.get("last_add_phone_failures") or 0),
  1057. "last_successes": int(state.get("last_successes") or 0),
  1058. "last_window_size": int(state.get("last_window_size") or 0),
  1059. "window_size": self._cfmail_add_phone_window,
  1060. "threshold": self._cfmail_add_phone_threshold,
  1061. "max_successes_in_window": self._cfmail_add_phone_max_successes,
  1062. }
  1063. cooldown_until = float(state.get("cooldown_until") or 0.0)
  1064. remaining = max(0, int(cooldown_until - time.time()))
  1065. if remaining <= 0:
  1066. state["in_cooldown"] = False
  1067. return {
  1068. "active_domain": str(state.get("active_domain") or ""),
  1069. "in_cooldown": bool(state.get("in_cooldown")),
  1070. "cooldown_remaining_seconds": remaining,
  1071. "last_triggered_at": str(state.get("last_triggered_at") or ""),
  1072. "last_reason": str(state.get("last_reason") or ""),
  1073. "last_add_phone_failures": int(state.get("last_add_phone_failures") or 0),
  1074. "last_successes": int(state.get("last_successes") or 0),
  1075. "last_window_size": int(state.get("last_window_size") or 0),
  1076. "window_size": self._cfmail_add_phone_window,
  1077. "threshold": self._cfmail_add_phone_threshold,
  1078. "max_successes_in_window": self._cfmail_add_phone_max_successes,
  1079. }
  1080. def _reset_cfmail_add_phone_stoploss(self, new_domain: str = "") -> None:
  1081. with self._lock:
  1082. self._cfmail_add_phone_state = {
  1083. "active_domain": new_domain,
  1084. "in_cooldown": False,
  1085. "cooldown_until": 0.0,
  1086. "last_triggered_at": "",
  1087. "last_rotation_attempted_at": "",
  1088. "last_reason": "",
  1089. "last_add_phone_failures": 0,
  1090. "last_successes": 0,
  1091. "last_window_size": 0,
  1092. "last_logged_at": 0.0,
  1093. }
  1094. def _cfmail_wait_otp_stoploss_snapshot(self) -> dict[str, object]:
  1095. with self._lock:
  1096. state = dict(self._cfmail_wait_otp_state)
  1097. if self._cfmail_wait_otp_cooldown_seconds <= 0:
  1098. return {
  1099. "active_domain": str(state.get("active_domain") or ""),
  1100. "in_cooldown": False,
  1101. "cooldown_remaining_seconds": 0,
  1102. "last_triggered_at": str(state.get("last_triggered_at") or ""),
  1103. "last_reason": str(state.get("last_reason") or ""),
  1104. "last_no_message_timeouts": int(state.get("last_no_message_timeouts") or 0),
  1105. "last_successes": int(state.get("last_successes") or 0),
  1106. "last_message_seen": int(state.get("last_message_seen") or 0),
  1107. "last_window_size": int(state.get("last_window_size") or 0),
  1108. "window_size": self._cfmail_wait_otp_window,
  1109. "threshold": self._cfmail_wait_otp_threshold,
  1110. }
  1111. cooldown_until = float(state.get("cooldown_until") or 0.0)
  1112. remaining = max(0, int(cooldown_until - time.time()))
  1113. if remaining <= 0:
  1114. state["in_cooldown"] = False
  1115. return {
  1116. "active_domain": str(state.get("active_domain") or ""),
  1117. "in_cooldown": bool(state.get("in_cooldown")),
  1118. "cooldown_remaining_seconds": remaining,
  1119. "last_triggered_at": str(state.get("last_triggered_at") or ""),
  1120. "last_reason": str(state.get("last_reason") or ""),
  1121. "last_no_message_timeouts": int(state.get("last_no_message_timeouts") or 0),
  1122. "last_successes": int(state.get("last_successes") or 0),
  1123. "last_message_seen": int(state.get("last_message_seen") or 0),
  1124. "last_window_size": int(state.get("last_window_size") or 0),
  1125. "window_size": self._cfmail_wait_otp_window,
  1126. "threshold": self._cfmail_wait_otp_threshold,
  1127. }
  1128. def _cfmail_fresh_domain_budget_snapshot(self) -> dict[str, object]:
  1129. with self._lock:
  1130. state = dict(self._cfmail_fresh_domain_state)
  1131. return {
  1132. "active_domain": str(state.get("active_domain") or ""),
  1133. "completed_attempts": int(state.get("completed_attempts") or 0),
  1134. "mail_seen_attempts": int(state.get("mail_seen_attempts") or 0),
  1135. "successes": int(state.get("successes") or 0),
  1136. "last_triggered_at": str(state.get("last_triggered_at") or ""),
  1137. "last_rotation_attempted_at": str(state.get("last_rotation_attempted_at") or ""),
  1138. "last_reason": str(state.get("last_reason") or ""),
  1139. "attempt_budget": self._cfmail_fresh_domain_attempt_budget,
  1140. }
  1141. def _reset_cfmail_wait_otp_stoploss(self, new_domain: str = "") -> None:
  1142. with self._lock:
  1143. self._cfmail_wait_otp_state = {
  1144. "active_domain": new_domain,
  1145. "in_cooldown": False,
  1146. "cooldown_until": 0.0,
  1147. "last_triggered_at": "",
  1148. "last_rotation_attempted_at": "",
  1149. "last_reason": "",
  1150. "last_no_message_timeouts": 0,
  1151. "last_successes": 0,
  1152. "last_message_seen": 0,
  1153. "last_window_size": 0,
  1154. "last_logged_at": 0.0,
  1155. }
  1156. with self._cfmail_wait_otp_live_lock:
  1157. if new_domain:
  1158. self._cfmail_wait_otp_live_progress = {
  1159. str(new_domain).strip().lower(): {}
  1160. }
  1161. else:
  1162. self._cfmail_wait_otp_live_progress = {}
  1163. def _reset_cfmail_fresh_domain_budget(self, new_domain: str = "") -> None:
  1164. with self._lock:
  1165. self._cfmail_fresh_domain_state = {
  1166. "active_domain": str(new_domain or "").strip().lower(),
  1167. "completed_attempts": 0,
  1168. "mail_seen_attempts": 0,
  1169. "successes": 0,
  1170. "last_triggered_at": "",
  1171. "last_rotation_attempted_at": "",
  1172. "last_reason": "",
  1173. }
  1174. def _reload_cfmail_manager_after_rotation(self) -> None:
  1175. manager = self._cfmail_manager
  1176. if manager is None:
  1177. return
  1178. try:
  1179. reload_if_needed = getattr(manager, "reload_if_needed", None)
  1180. if callable(reload_if_needed):
  1181. try:
  1182. reload_if_needed(force=True)
  1183. except TypeError:
  1184. reload_if_needed()
  1185. except Exception as exc:
  1186. self._log(f"[zhuce6:register] [cfmail] manager reload after rotation failed: {exc}")
  1187. def _ensure_cfmail_active_domain_ready(self) -> bool:
  1188. if self._cfmail_provisioner is None or self._cfmail_manager is None:
  1189. return False
  1190. active_accounts = self._current_cfmail_active_accounts()
  1191. active_domains = [item["domain"] for item in active_accounts]
  1192. if not active_domains:
  1193. return False
  1194. self._log(
  1195. f"[zhuce6:register] [cfmail] active-domain pool ready: "
  1196. f"{', '.join(active_domains)}"
  1197. )
  1198. if self._cfmail_tracker is not None:
  1199. for domain in active_domains:
  1200. try:
  1201. self._cfmail_tracker.active_domain = domain
  1202. except Exception:
  1203. pass
  1204. return True
  1205. def _force_rotate_cfmail_for_invalid_mailbox(self, thread_id: int, result: dict[str, object]) -> bool:
  1206. if not self._is_cfmail_invalid_domain_mailbox_failure(result):
  1207. return False
  1208. if self._cfmail_tracker is None or self._cfmail_provisioner is None:
  1209. return False
  1210. current_domain = self._extract_email_domain(result) or self._current_cfmail_active_domain()
  1211. if not current_domain:
  1212. return False
  1213. self._log(
  1214. f"[zhuce6:register] [thread-{thread_id}] [cfmail] rotating invalid domain "
  1215. f"{current_domain} after mailbox bootstrap failure"
  1216. )
  1217. return self._replace_cfmail_domain(
  1218. thread_id=thread_id,
  1219. domain=current_domain,
  1220. reason_label="mailbox invalid domain",
  1221. )
  1222. def _rotate_cfmail_for_failed_canary(self, thread_id: int, result: dict[str, object]) -> bool:
  1223. del thread_id, result
  1224. return False
  1225. def _rotate_cfmail_for_fresh_domain_budget(self, thread_id: int) -> bool:
  1226. if self._cfmail_fresh_domain_attempt_budget <= 0:
  1227. return False
  1228. if self._cfmail_tracker is None or self._cfmail_provisioner is None:
  1229. return False
  1230. with self._lock:
  1231. state = dict(self._cfmail_fresh_domain_state)
  1232. domain = str(state.get("active_domain") or "").strip().lower()
  1233. completed_attempts = int(state.get("completed_attempts") or 0)
  1234. mail_seen_attempts = int(state.get("mail_seen_attempts") or 0)
  1235. if (
  1236. not domain
  1237. or mail_seen_attempts <= 0
  1238. or completed_attempts < self._cfmail_fresh_domain_attempt_budget
  1239. or str(state.get("last_rotation_attempted_at") or "").strip()
  1240. ):
  1241. return False
  1242. self._cfmail_fresh_domain_state["last_rotation_attempted_at"] = datetime.now().isoformat(timespec="seconds")
  1243. if not self._is_cfmail_active_domain(domain):
  1244. self._clear_cfmail_domain_state(domain)
  1245. return False
  1246. self._log(
  1247. f"[zhuce6:register] [thread-{thread_id}] [cfmail] rotating domain {domain} "
  1248. f"because fresh domain budget reached"
  1249. )
  1250. if self._replace_cfmail_domain(
  1251. thread_id=thread_id,
  1252. domain=domain,
  1253. reason_label="fresh domain budget reached",
  1254. ):
  1255. return True
  1256. with self._lock:
  1257. self._cfmail_fresh_domain_state["last_rotation_attempted_at"] = ""
  1258. return False
  1259. def _rotate_cfmail_for_stoploss(
  1260. self,
  1261. *,
  1262. thread_id: int,
  1263. state_attr: str,
  1264. reason_label: str,
  1265. ) -> bool:
  1266. if self._cfmail_tracker is None or self._cfmail_provisioner is None:
  1267. return False
  1268. with self._lock:
  1269. state_obj = getattr(self, state_attr, None)
  1270. if not isinstance(state_obj, dict):
  1271. return False
  1272. domain = str(state_obj.get("active_domain") or "").strip().lower()
  1273. if not state_obj.get("in_cooldown") or not domain:
  1274. return False
  1275. if str(state_obj.get("last_rotation_attempted_at") or "").strip():
  1276. return False
  1277. state_obj["last_rotation_attempted_at"] = datetime.now().isoformat(timespec="seconds")
  1278. return self._replace_cfmail_domain(
  1279. thread_id=thread_id,
  1280. domain=domain,
  1281. reason_label=reason_label,
  1282. )
  1283. def _wait_if_cfmail_add_phone_stopped(self, thread_id: int, provider: str) -> bool:
  1284. if provider != "cfmail":
  1285. return False
  1286. if self._cfmail_add_phone_cooldown_seconds <= 0:
  1287. return False
  1288. state = self._cfmail_add_phone_stoploss_snapshot()
  1289. if not state.get("in_cooldown"):
  1290. return False
  1291. had_spare_domains = len(self._cfmail_active_domain_set()) > 1
  1292. domain = str(state.get("active_domain") or "").strip().lower()
  1293. if not self._is_cfmail_active_domain(domain):
  1294. self._clear_cfmail_domain_state(domain)
  1295. return False
  1296. if self._rotate_cfmail_for_stoploss(
  1297. thread_id=thread_id,
  1298. state_attr="_cfmail_add_phone_state",
  1299. reason_label="add_phone stoploss",
  1300. ):
  1301. return not had_spare_domains
  1302. remaining = int(state.get("cooldown_remaining_seconds") or 0)
  1303. should_log = False
  1304. with self._lock:
  1305. last_logged_at = float(self._cfmail_add_phone_state.get("last_logged_at") or 0.0)
  1306. now = time.time()
  1307. if now - last_logged_at >= 15:
  1308. self._cfmail_add_phone_state["last_logged_at"] = now
  1309. should_log = True
  1310. if should_log:
  1311. self._log(
  1312. f"[zhuce6:register] [thread-{thread_id}] [cfmail] add_phone stoploss active for {domain}, "
  1313. f"remaining={remaining}s"
  1314. )
  1315. wait_seconds = min(max(remaining, 1), 5)
  1316. self._stop_event.wait(wait_seconds)
  1317. return True
  1318. def _wait_if_cfmail_wait_otp_stopped(self, thread_id: int, provider: str) -> bool:
  1319. if provider != "cfmail":
  1320. return False
  1321. if self._cfmail_wait_otp_cooldown_seconds <= 0:
  1322. return False
  1323. state = self._cfmail_wait_otp_stoploss_snapshot()
  1324. if not state.get("in_cooldown"):
  1325. return False
  1326. had_spare_domains = len(self._cfmail_active_domain_set()) > 1
  1327. domain = str(state.get("active_domain") or "").strip().lower()
  1328. if not self._is_cfmail_active_domain(domain):
  1329. self._clear_cfmail_domain_state(domain)
  1330. return False
  1331. if self._rotate_cfmail_for_stoploss(
  1332. thread_id=thread_id,
  1333. state_attr="_cfmail_wait_otp_state",
  1334. reason_label="wait_otp stoploss",
  1335. ):
  1336. return not had_spare_domains
  1337. remaining = int(state.get("cooldown_remaining_seconds") or 0)
  1338. should_log = False
  1339. with self._lock:
  1340. last_logged_at = float(self._cfmail_wait_otp_state.get("last_logged_at") or 0.0)
  1341. now = time.time()
  1342. if now - last_logged_at >= 15:
  1343. self._cfmail_wait_otp_state["last_logged_at"] = now
  1344. should_log = True
  1345. if should_log:
  1346. self._log(
  1347. f"[zhuce6:register] [thread-{thread_id}] [cfmail] wait_otp stoploss active for {domain}, "
  1348. f"remaining={remaining}s"
  1349. )
  1350. wait_seconds = min(max(remaining, 1), 5)
  1351. self._stop_event.wait(wait_seconds)
  1352. return True
  1353. def start(self) -> None:
  1354. if self._threads:
  1355. return
  1356. _compat_main_attr("load_all", load_all)()
  1357. self._stop_event.clear()
  1358. self._target_reached.clear()
  1359. self._started_at = time.time()
  1360. num = self.settings.register_threads
  1361. self._providers = [p.strip() for p in self.settings.register_mail_provider.split(",") if p.strip()]
  1362. if not self._providers:
  1363. self._providers = ["cfmail"]
  1364. if "cfmail" in self._providers:
  1365. from core.cfmail_domain_rotation import DomainHealthTracker
  1366. import core.cfmail as cfmail_module
  1367. from core.cfmail import DEFAULT_CFMAIL_MANAGER
  1368. from core.cfmail_provisioner import CfmailProvisioner
  1369. self._cfmail_tracker = DomainHealthTracker()
  1370. self._cfmail_provisioner = CfmailProvisioner(proxy_url=self.settings.register_proxy)
  1371. self._cfmail_manager = DEFAULT_CFMAIL_MANAGER
  1372. cfmail_module.CFMAIL_WAIT_ABORT_PREDICATE = self._should_abort_cfmail_wait
  1373. cfmail_module.CFMAIL_WAIT_PROGRESS_CALLBACK = self._on_cfmail_wait_progress
  1374. try:
  1375. normalize_result = self._cfmail_provisioner.normalize_to_domain_pool(self._cfmail_active_domain_count)
  1376. self._cfmail_manager.reload_if_needed(force=True)
  1377. provisioned_domains = list(normalize_result.get("provisioned_domains") or [])
  1378. retired_domains = list(normalize_result.get("retired_domains") or [])
  1379. active_domains = list(normalize_result.get("active_domains") or [])
  1380. if provisioned_domains or retired_domains:
  1381. self._log(
  1382. "[zhuce6:register] [cfmail] normalized active domain pool: "
  1383. f"active={','.join(active_domains) or '-'} "
  1384. f"provisioned={','.join(provisioned_domains) or '-'} "
  1385. f"retired={','.join(retired_domains) or '-'}"
  1386. )
  1387. except Exception as exc:
  1388. self._log(f"[zhuce6:register] [cfmail] normalize active domain pool failed: {exc}")
  1389. try:
  1390. self._ensure_cfmail_active_domain_ready()
  1391. self._ensure_cfmail_domain_pool_target(
  1392. trigger_thread_id=0,
  1393. reason="startup",
  1394. )
  1395. except Exception as exc:
  1396. self._log(f"[zhuce6:register] [cfmail] startup active-domain preflight failed: {exc}")
  1397. if self.settings.backend == "cpa" and self.settings.cpa_runtime_reconcile_enabled:
  1398. try:
  1399. _maybe_reconcile_cpa_runtime(
  1400. pool_dir=self.settings.pool_dir,
  1401. management_base_url=self.settings.cpa_management_base_url,
  1402. enabled=True,
  1403. cooldown_seconds=self.settings.cpa_runtime_reconcile_cooldown_seconds,
  1404. restart_enabled=self.settings.cpa_runtime_reconcile_restart_enabled,
  1405. state_file=self.settings.pool_dir / "cpa_runtime_reconcile_state.json",
  1406. client=create_backend_client(self.settings),
  1407. management_key=self.settings.cpa_management_key,
  1408. )
  1409. except Exception as exc:
  1410. self._log(f"[zhuce6:register] startup reconcile failed: {exc}")
  1411. if self.settings.proxy_pool_configured:
  1412. from core.proxy_pool import ProxyPool
  1413. self._proxy_pool = ProxyPool.from_settings(self.settings)
  1414. if self._proxy_pool is not None:
  1415. self._proxy_pool.start()
  1416. target_msg = f", target={self.settings.register_target_count}" if self.settings.register_target_count > 0 else ""
  1417. self._log(
  1418. f"[zhuce6:register] starting {num} threads, "
  1419. f"providers={','.join(self._providers)}, "
  1420. f"proxy={self.settings.register_proxy or 'none'}, "
  1421. f"sleep={self.settings.register_sleep_min}-{self.settings.register_sleep_max}s"
  1422. f"{target_msg}"
  1423. )
  1424. threading_module = _compat_main_attr("threading", threading)
  1425. for i in range(num):
  1426. provider = self._providers[i % len(self._providers)]
  1427. t = threading_module.Thread(
  1428. target=self._worker,
  1429. args=(i + 1, provider),
  1430. daemon=True,
  1431. name=f"zhuce6-register-{i + 1}",
  1432. )
  1433. t.start()
  1434. self._threads.append(t)
  1435. # Solution B: start deferred retry worker thread
  1436. pending_t = threading_module.Thread(
  1437. target=self._pending_token_retry_worker,
  1438. daemon=True,
  1439. name="zhuce6-pending-token-retry",
  1440. )
  1441. pending_t.start()
  1442. self._threads.append(pending_t)
  1443. self._write_runtime_state()
  1444. def stop(self) -> None:
  1445. self._stop_event.set()
  1446. for t in self._threads:
  1447. t.join(timeout=2)
  1448. self._threads.clear()
  1449. try:
  1450. import core.cfmail as cfmail_module
  1451. cfmail_module.CFMAIL_WAIT_ABORT_PREDICATE = None
  1452. cfmail_module.CFMAIL_WAIT_PROGRESS_CALLBACK = None
  1453. except Exception:
  1454. pass
  1455. if self._proxy_pool is not None:
  1456. self._proxy_pool.close()
  1457. self._proxy_pool = None
  1458. self._write_runtime_state()
  1459. def _should_stop(self) -> bool:
  1460. return self._stop_event.is_set() or self._target_reached.is_set()
  1461. def _check_target(self) -> bool:
  1462. """Return True if target reached and threads should stop."""
  1463. if self.settings.register_target_count <= 0:
  1464. return False
  1465. with self._lock:
  1466. if self._total_success >= self.settings.register_target_count:
  1467. self._target_reached.set()
  1468. return True
  1469. return False
  1470. def _cpa_api_root(self) -> str:
  1471. parsed = urlsplit(self.settings.cpa_management_base_url)
  1472. path = parsed.path or ""
  1473. suffix = "/v0/management"
  1474. if path.endswith(suffix):
  1475. path = path[: -len(suffix)]
  1476. return urlunsplit((parsed.scheme, parsed.netloc, path, "", "")).rstrip("/")
  1477. def _get_cpa_management_key(self) -> str | None:
  1478. cached = self._cpa_management_key_cache
  1479. if cached is not False:
  1480. return str(cached or "") or None
  1481. key = str(get_management_key() or "").strip() or None
  1482. self._cpa_management_key_cache = key or None
  1483. return key
  1484. def _sync_cpa_from_success(self, result: dict[str, object], thread_id: int) -> tuple[bool, str, str]:
  1485. """Persist a registered account to CPA immediately while keeping pool as backup."""
  1486. pool_file_raw = str(result.get("pool_file") or "").strip()
  1487. if not pool_file_raw:
  1488. return False, "missing pool file", ""
  1489. pool_file = Path(pool_file_raw)
  1490. if not pool_file.is_file():
  1491. return False, f"pool file missing: {pool_file.name}", ""
  1492. sync_started_at = now_iso()
  1493. key = self._get_cpa_management_key()
  1494. if not key:
  1495. update_token_record(
  1496. pool_file,
  1497. backup_written=True,
  1498. cpa_sync_status="failed",
  1499. last_cpa_sync_at=sync_started_at,
  1500. last_cpa_sync_error="CPA management key unavailable",
  1501. )
  1502. return False, "CPA management key unavailable", ""
  1503. try:
  1504. from platforms.chatgpt.pool import load_token_record
  1505. token_data = load_token_record(pool_file)
  1506. except Exception as exc:
  1507. return False, f"invalid pool file {pool_file.name}: {exc}", ""
  1508. metadata = result.get("metadata") if isinstance(result.get("metadata"), dict) else {}
  1509. post_create_gate = str((metadata or {}).get("post_create_gate") or token_data.get("registration_post_create_gate") or "").strip().lower()
  1510. if post_create_gate == "add_phone" and not bool(token_data.get("warmup_required")):
  1511. token_data = update_token_record(
  1512. pool_file,
  1513. warmup_required=True,
  1514. warmup_state="pending",
  1515. warmup_passed=False,
  1516. registration_post_create_gate="add_phone",
  1517. )
  1518. if not isinstance(token_data, dict) or not str(token_data.get("email") or "").strip():
  1519. update_token_record(
  1520. pool_file,
  1521. backup_written=True,
  1522. cpa_sync_status="failed",
  1523. last_cpa_sync_at=sync_started_at,
  1524. last_cpa_sync_error="missing email in pool record",
  1525. )
  1526. return False, f"missing email in {pool_file.name}", ""
  1527. if is_warmup_pending_record(token_data):
  1528. update_token_record(
  1529. pool_file,
  1530. backup_written=True,
  1531. cpa_sync_status="warmup_pending",
  1532. last_cpa_sync_at=sync_started_at,
  1533. last_cpa_sync_error="warmup pending",
  1534. )
  1535. return (
  1536. False,
  1537. "warmup pending",
  1538. str(token_data.get("email") or pool_file.name).strip() or pool_file.name,
  1539. )
  1540. from platforms.chatgpt.cpa_upload import upload_to_cpa
  1541. ok, message = upload_to_cpa(
  1542. token_data,
  1543. api_url=self._cpa_api_root(),
  1544. api_key=key,
  1545. proxy=None,
  1546. )
  1547. email = str(token_data.get("email") or pool_file.name).strip() or pool_file.name
  1548. if ok:
  1549. update_token_record(
  1550. pool_file,
  1551. health_status="good",
  1552. backup_written=True,
  1553. cpa_sync_status="synced",
  1554. last_cpa_sync_at=sync_started_at,
  1555. last_cpa_sync_error="",
  1556. )
  1557. return True, "", email
  1558. else:
  1559. update_token_record(
  1560. pool_file,
  1561. backup_written=True,
  1562. cpa_sync_status="failed",
  1563. last_cpa_sync_at=sync_started_at,
  1564. last_cpa_sync_error=message,
  1565. )
  1566. return False, message, email
  1567. # ── Solution B: Deferred retry queue ──────────────────────────
  1568. def _enqueue_pending_token(
  1569. self,
  1570. result: dict[str, object],
  1571. thread_id: int,
  1572. *,
  1573. proxy_key: str = "",
  1574. proxy_url: str = "",
  1575. ) -> None:
  1576. """Save an add_phone_gate account for deferred token acquisition retry."""
  1577. metadata = result.get("metadata") if isinstance(result.get("metadata"), dict) else {}
  1578. deferred = metadata.get("deferred_credentials")
  1579. if not isinstance(deferred, dict):
  1580. return
  1581. email = str(deferred.get("email") or "").strip()
  1582. password = str(deferred.get("password") or "").strip()
  1583. if not email or not password:
  1584. return
  1585. entry = {
  1586. "email": email,
  1587. "password": password,
  1588. "mailbox_jwt": str(deferred.get("mailbox_jwt") or "").strip(),
  1589. "mailbox_extra": dict(deferred.get("mailbox_extra") or {}),
  1590. "registration_proxy_key": str(deferred.get("registration_proxy_key") or proxy_key or "").strip(),
  1591. "registration_proxy_region": str(
  1592. deferred.get("registration_proxy_region")
  1593. or infer_proxy_region(
  1594. str(deferred.get("registration_proxy_key") or proxy_key or "").strip()
  1595. or str(deferred.get("registration_proxy_url") or proxy_url or "").strip()
  1596. )
  1597. or ""
  1598. ).strip(),
  1599. "registration_proxy_url": str(deferred.get("registration_proxy_url") or proxy_url or "").strip(),
  1600. "registration_fingerprint_profile": str(
  1601. deferred.get("registration_fingerprint_profile") or "chrome120_win"
  1602. ).strip(),
  1603. "cfmail_profile_name": str(
  1604. deferred.get("cfmail_profile_name") or metadata.get("cfmail_profile_name") or ""
  1605. ).strip(),
  1606. "add_phone_trace_path": str(
  1607. deferred.get("add_phone_trace_path") or metadata.get("add_phone_trace_path") or ""
  1608. ).strip(),
  1609. "created_at": time.time(),
  1610. "retry_count": 0,
  1611. "last_retry_at": 0.0,
  1612. }
  1613. with self._pending_token_lock:
  1614. self._pending_token_queue.append(entry)
  1615. self._pending_token_total_enqueued += 1
  1616. origin_proxy = str(entry.get("registration_proxy_key") or entry.get("registration_proxy_region") or entry.get("registration_proxy_url") or "").strip()
  1617. self._log(
  1618. f"[zhuce6:register] [thread-{thread_id}] 📥 deferred token retry enqueued: {email} "
  1619. f"(origin_proxy={origin_proxy or '-'}, queue_size={len(self._pending_token_queue)})"
  1620. )
  1621. def _pending_token_retry_worker(self) -> None:
  1622. """Background thread that retries token acquisition for queued add_phone_gate accounts."""
  1623. while not self._should_stop():
  1624. if self._stop_event.wait(30):
  1625. break
  1626. batch: list[dict[str, Any]] = []
  1627. now = time.time()
  1628. with self._pending_token_lock:
  1629. remaining: list[dict[str, Any]] = []
  1630. for entry in self._pending_token_queue:
  1631. created_at = float(entry.get("created_at") or 0)
  1632. retry_count = int(entry.get("retry_count") or 0)
  1633. last_retry = float(entry.get("last_retry_at") or 0)
  1634. age = now - created_at
  1635. since_last = now - last_retry if last_retry > 0 else age
  1636. delay = self._pending_token_retry_delay_seconds * (retry_count + 1)
  1637. if retry_count >= self._pending_token_max_retries:
  1638. self._pending_token_total_failed += 1
  1639. self._log(
  1640. f"[zhuce6:register] [pending] ❌ exhausted retries for "
  1641. f"{entry.get('email')}, discarding"
  1642. )
  1643. continue
  1644. if since_last >= delay:
  1645. batch.append(entry)
  1646. else:
  1647. remaining.append(entry)
  1648. self._pending_token_queue = remaining
  1649. if not batch:
  1650. continue
  1651. for entry in batch:
  1652. if self._should_stop():
  1653. break
  1654. self._retry_pending_token(entry)
  1655. def _retry_pending_token(self, entry: dict[str, Any]) -> None:
  1656. """Attempt token acquisition for a single deferred account."""
  1657. email = str(entry.get("email") or "").strip()
  1658. password = str(entry.get("password") or "").strip()
  1659. mailbox_jwt = str(entry.get("mailbox_jwt") or "").strip()
  1660. mailbox_extra = dict(entry.get("mailbox_extra") or {})
  1661. registration_proxy_key = str(entry.get("registration_proxy_key") or "").strip()
  1662. registration_proxy_region = str(entry.get("registration_proxy_region") or "").strip()
  1663. registration_proxy_url = str(entry.get("registration_proxy_url") or "").strip()
  1664. registration_fingerprint_profile = str(
  1665. entry.get("registration_fingerprint_profile") or ""
  1666. ).strip()
  1667. cfmail_profile_name = str(entry.get("cfmail_profile_name") or "").strip()
  1668. add_phone_trace_path = str(entry.get("add_phone_trace_path") or "").strip()
  1669. retry_count = int(entry.get("retry_count") or 0) + 1
  1670. self._log(
  1671. f"[zhuce6:register] [pending] 🔄 retrying token acquisition "
  1672. f"for {email} (attempt {retry_count}/{self._pending_token_max_retries})"
  1673. )
  1674. proxy_lease = None
  1675. proxy_url = registration_proxy_url or self.settings.register_proxy
  1676. proxy_release_success = False
  1677. try:
  1678. from core.base_mailbox import MailboxAccount
  1679. from core.cfmail import CfMailMailbox, DEFAULT_CFMAIL_MANAGER
  1680. from platforms.chatgpt.plugin import MailboxEmailServiceAdapter
  1681. from platforms.chatgpt.register import RegistrationEngine
  1682. from platforms.chatgpt.pool import write_token_record
  1683. if self._proxy_pool is not None:
  1684. try:
  1685. proxy_lease = self._proxy_pool.acquire(
  1686. timeout=5.0,
  1687. preferred_name=registration_proxy_key or None,
  1688. preferred_regions=(registration_proxy_region,) if registration_proxy_region else (),
  1689. )
  1690. proxy_url = str(proxy_lease.proxy_url or "").strip() or proxy_url
  1691. self._log(
  1692. "[zhuce6:register] [pending] "
  1693. f"using retry proxy {proxy_lease.name} for {email} "
  1694. f"(origin={registration_proxy_key or registration_proxy_region or registration_proxy_url or '-'})"
  1695. )
  1696. except Exception as exc:
  1697. self._log(
  1698. "[zhuce6:register] [pending] "
  1699. f"proxy acquire failed for {email}: {exc}"
  1700. )
  1701. mailbox = CfMailMailbox(manager=DEFAULT_CFMAIL_MANAGER)
  1702. adapter = MailboxEmailServiceAdapter(mailbox)
  1703. # Reconstruct the mailbox account so _wait_for_mailbox_code can poll
  1704. if mailbox_jwt and mailbox_extra:
  1705. adapter._account = MailboxAccount(
  1706. email=email,
  1707. account_id=mailbox_jwt,
  1708. extra=dict(mailbox_extra),
  1709. )
  1710. engine = RegistrationEngine(
  1711. email_service=adapter,
  1712. proxy_url=proxy_url,
  1713. )
  1714. engine.email = email
  1715. engine.password = password
  1716. token_info = engine._login_for_token()
  1717. if token_info:
  1718. proxy_release_success = True
  1719. self._log(f"[zhuce6:register] [pending] ✅ deferred token acquired for {email}")
  1720. # Write to pool
  1721. token_data = {
  1722. "type": "codex",
  1723. "email": email,
  1724. "password": password,
  1725. "mail_provider": "cfmail",
  1726. "expired": str(token_info.get("expired") or ""),
  1727. "id_token": str(token_info.get("id_token") or ""),
  1728. "account_id": str(token_info.get("account_id") or ""),
  1729. "access_token": str(token_info.get("access_token") or ""),
  1730. "last_refresh": str(token_info.get("last_refresh") or ""),
  1731. "refresh_token": str(token_info.get("refresh_token") or ""),
  1732. "source": "deferred_retry",
  1733. "add_phone_trace_path": add_phone_trace_path,
  1734. }
  1735. token_data.update(
  1736. build_registration_provenance(
  1737. {
  1738. "location": "",
  1739. "mail_provider": "cfmail",
  1740. "post_create_gate": "add_phone",
  1741. },
  1742. proxy_url=registration_proxy_url or proxy_url,
  1743. proxy_key=registration_proxy_key or getattr(proxy_lease, "name", ""),
  1744. proxy_region=registration_proxy_region,
  1745. cfmail_profile_name=cfmail_profile_name,
  1746. )
  1747. )
  1748. if registration_fingerprint_profile:
  1749. token_data["registration_fingerprint_profile"] = registration_fingerprint_profile
  1750. pool_file = write_token_record(token_data, self.settings.pool_dir)
  1751. sync_ok, sync_error, synced_email = self._sync_cpa_from_success(
  1752. {"pool_file": str(pool_file), "success": True, "stage": "deferred_retry"},
  1753. thread_id=0,
  1754. )
  1755. with self._lock:
  1756. if sync_ok:
  1757. self._total_success += 1
  1758. self._total_cpa_sync_success += 1
  1759. self._pending_token_total_success += 1
  1760. elif sync_error == "warmup pending":
  1761. self._total_warmup_pending += 1
  1762. else:
  1763. self._total_failure += 1
  1764. self._total_cpa_sync_failure += 1
  1765. self._pending_token_total_failed += 1
  1766. if sync_ok:
  1767. self._log(f"[zhuce6:register] [pending] ✅ CPA sync success: {synced_email or email}")
  1768. elif sync_error == "warmup pending":
  1769. self._log(f"[zhuce6:register] [pending] ⏳ warmup pending: {synced_email or email}")
  1770. else:
  1771. self._log(f"[zhuce6:register] [pending] ❌ failed [stage=cpa_sync]: {sync_error or 'unknown'}")
  1772. return
  1773. else:
  1774. self._log(f"[zhuce6:register] [pending] ⏳ deferred retry failed for {email}")
  1775. except Exception as exc:
  1776. self._log(f"[zhuce6:register] [pending] ⚠️ deferred retry error for {email}: {exc}")
  1777. finally:
  1778. if proxy_lease is not None and self._proxy_pool is not None:
  1779. try:
  1780. self._proxy_pool.release(
  1781. proxy_lease,
  1782. success=proxy_release_success,
  1783. stage="deferred_retry",
  1784. )
  1785. except Exception as exc:
  1786. self._log(f"[zhuce6:register] [pending] proxy release failed for {email}: {exc}")
  1787. # Re-enqueue with incremented retry count
  1788. entry["retry_count"] = retry_count
  1789. entry["last_retry_at"] = time.time()
  1790. with self._pending_token_lock:
  1791. self._pending_token_queue.append(entry)
  1792. def _worker(self, thread_id: int, initial_provider: str) -> None:
  1793. provider = initial_provider
  1794. consecutive_failures = 0
  1795. max_failures = self.settings.register_max_consecutive_failures
  1796. while not self._should_stop():
  1797. if provider == "cfmail" and self._cfmail_tracker is not None:
  1798. self._cfmail_rotation_pause.wait()
  1799. if self._wait_if_cfmail_add_phone_stopped(thread_id, provider):
  1800. continue
  1801. if self._wait_if_cfmail_wait_otp_stopped(thread_id, provider):
  1802. continue
  1803. if self._wait_if_cfmail_canary_pending(thread_id, provider):
  1804. continue
  1805. if provider == "cfmail" and self._cfmail_manager is not None:
  1806. if self._check_cfmail_all_cooldown_rotation(thread_id):
  1807. continue
  1808. if self._cfmail_all_accounts_in_cooldown():
  1809. self._stop_event.wait(10)
  1810. continue
  1811. if self._wait_if_cfmail_flow_throttled(thread_id, provider):
  1812. continue
  1813. try:
  1814. result: dict[str, object] = {}
  1815. proxy_key = ""
  1816. proxy_outcome: bool | None = False
  1817. result_metadata: dict[str, object] = {}
  1818. result_stage = "?"
  1819. result_error = ""
  1820. result_email = ""
  1821. max_proxy_attempts = 2 if self._proxy_pool is not None else 1
  1822. for proxy_attempt in range(1, max_proxy_attempts + 1):
  1823. proxy_lease = None
  1824. proxy_url = self.settings.register_proxy
  1825. proxy_key = proxy_url or ""
  1826. proxy_outcome = False
  1827. release_stage = "exception"
  1828. try:
  1829. if self._proxy_pool is not None:
  1830. proxy_lease = self._proxy_pool.acquire(
  1831. timeout=5,
  1832. preferred_regions=tuple(self.settings.register_fresh_proxy_regions or ()),
  1833. )
  1834. proxy_url = proxy_lease.proxy_url
  1835. proxy_key = proxy_lease.name or proxy_url or ""
  1836. self._log(
  1837. f"[zhuce6:register] [thread-{thread_id}] acquired proxy {proxy_lease.local_port} ({proxy_lease.name})"
  1838. )
  1839. self._log(f"[zhuce6:register] [thread-{thread_id}] attempting (provider={provider})")
  1840. cfmail_profile_name = self._selected_cfmail_profile(thread_id) if provider == "cfmail" else "auto"
  1841. result = _compat_main_attr("run_chatgpt_register_once", run_chatgpt_register_once)(
  1842. email=None,
  1843. password=None,
  1844. mail_provider=provider,
  1845. cfmail_profile_name=cfmail_profile_name,
  1846. proxy=proxy_url,
  1847. write_pool=True,
  1848. pool_dir=self.settings.pool_dir,
  1849. )
  1850. proxy_outcome = self._classify_proxy_outcome(result)
  1851. result_metadata = result.get("metadata") if isinstance(result.get("metadata"), dict) else {}
  1852. result_stage = str(result.get("stage") or "?").strip() or "?"
  1853. if bool(result.get("success")) and result_stage == "?":
  1854. result_stage = "completed"
  1855. result_error = str(result.get("error_message") or "").strip()
  1856. result_email = str(result.get("email") or "").strip()
  1857. release_stage = result_stage
  1858. should_retry_device_id = (
  1859. self._proxy_pool is not None
  1860. and proxy_attempt < max_proxy_attempts
  1861. and not bool(result.get("success"))
  1862. and result_stage == "device_id"
  1863. )
  1864. if should_retry_device_id:
  1865. self._log(
  1866. f"[zhuce6:register] [thread-{thread_id}] device_id failed on proxy {proxy_key}; "
  1867. "rotating proxy and retrying once"
  1868. )
  1869. continue
  1870. break
  1871. finally:
  1872. if proxy_lease is not None and self._proxy_pool is not None:
  1873. try:
  1874. self._proxy_pool.release(
  1875. proxy_lease,
  1876. success=proxy_outcome,
  1877. stage=release_stage,
  1878. )
  1879. except Exception as exc:
  1880. self._log(f"[zhuce6:register] [thread-{thread_id}] proxy release failed: {exc}")
  1881. success = bool(result.get("success"))
  1882. should_break_after_iteration = False
  1883. with self._lock:
  1884. self._total_attempts += 1
  1885. if success:
  1886. pool_file_raw = str(result.get("pool_file") or "").strip()
  1887. if pool_file_raw:
  1888. try:
  1889. update_token_record(
  1890. Path(pool_file_raw),
  1891. **build_registration_provenance(
  1892. result_metadata,
  1893. proxy_url=proxy_url,
  1894. proxy_key=proxy_key,
  1895. proxy_region=infer_proxy_region(proxy_key or proxy_url),
  1896. cfmail_profile_name=str(result_metadata.get("cfmail_profile_name") or ""),
  1897. ),
  1898. )
  1899. except Exception as exc:
  1900. self._log(
  1901. f"[zhuce6:register] [thread-{thread_id}] provenance update failed: {exc}"
  1902. )
  1903. sync_ok, sync_error, synced_email = self._sync_cpa_from_success(result, thread_id)
  1904. with self._lock:
  1905. if sync_ok:
  1906. self._total_success += 1
  1907. self._total_cpa_sync_success += 1
  1908. self._last_error = None
  1909. consecutive_failures = 0
  1910. self._record_attempt(
  1911. success=True,
  1912. stage=result_stage,
  1913. error_message="",
  1914. metadata=result_metadata,
  1915. proxy_key=proxy_key,
  1916. email=result_email or synced_email,
  1917. )
  1918. email = result_email or synced_email or "?"
  1919. self._log(f"[zhuce6:register] [thread-{thread_id}] \u2705 success: {email}")
  1920. self._log(f"[zhuce6:register] [thread-{thread_id}] ✅ CPA sync success: {email}")
  1921. if self._check_target():
  1922. self._log(
  1923. f"[zhuce6:register] target reached ({self.settings.register_target_count}), stopping"
  1924. )
  1925. should_break_after_iteration = True
  1926. elif sync_error == "warmup pending":
  1927. self._total_warmup_pending += 1
  1928. self._last_error = None
  1929. consecutive_failures = 0
  1930. self._record_attempt(
  1931. success=False,
  1932. stage="warmup_pending",
  1933. error_message="warmup pending",
  1934. metadata=result_metadata,
  1935. proxy_key=proxy_key,
  1936. email=result_email or synced_email,
  1937. )
  1938. email = result_email or synced_email or "?"
  1939. self._log(f"[zhuce6:register] [thread-{thread_id}] ⏳ warmup pending: {email}")
  1940. else:
  1941. self._total_failure += 1
  1942. self._total_cpa_sync_failure += 1
  1943. consecutive_failures += 1
  1944. err = sync_error or "CPA sync failed"
  1945. self._last_error = err
  1946. self._record_attempt(
  1947. success=False,
  1948. stage="cpa_sync",
  1949. error_message=err,
  1950. metadata=result_metadata,
  1951. proxy_key=proxy_key,
  1952. email=result_email or synced_email,
  1953. )
  1954. self._log(
  1955. f"[zhuce6:register] [thread-{thread_id}] \u274c failed ({consecutive_failures}/{max_failures}) "
  1956. f"[stage=cpa_sync]: {err}"
  1957. )
  1958. else:
  1959. with self._lock:
  1960. self._total_failure += 1
  1961. consecutive_failures += 1
  1962. err = result_error or "unknown"
  1963. stage = result_stage
  1964. self._last_error = err
  1965. self._record_attempt(
  1966. success=False,
  1967. stage=stage,
  1968. error_message=err,
  1969. metadata=result_metadata,
  1970. proxy_key=proxy_key,
  1971. email=result_email,
  1972. )
  1973. self._log(f"[zhuce6:register] [thread-{thread_id}] \u274c failed ({consecutive_failures}/{max_failures}) [stage={stage}]: {err}")
  1974. for log_line in result.get("logs", []):
  1975. self._log(f"[zhuce6:register] [thread-{thread_id}] \u21b3 {log_line}")
  1976. # Solution B: enqueue add_phone_gate accounts for deferred retry
  1977. if result_stage == "add_phone_gate":
  1978. self._enqueue_pending_token(
  1979. result,
  1980. thread_id,
  1981. proxy_key=proxy_key,
  1982. proxy_url=proxy_url,
  1983. )
  1984. invalid_mailbox_rotation = self._force_rotate_cfmail_for_invalid_mailbox(
  1985. thread_id,
  1986. result,
  1987. )
  1988. canary_rotation = self._rotate_cfmail_for_failed_canary(thread_id, result)
  1989. rotation_success = self._handle_cfmail_rotation(
  1990. thread_id=thread_id,
  1991. result=result,
  1992. proxy_key=proxy_key,
  1993. )
  1994. self._update_cfmail_add_phone_stoploss(result)
  1995. self._update_cfmail_wait_otp_stoploss(result)
  1996. self._update_cfmail_canary_after_result(thread_id=thread_id, result=result)
  1997. self._update_cfmail_fresh_domain_budget(result)
  1998. fresh_domain_rotation = self._rotate_cfmail_for_fresh_domain_budget(thread_id)
  1999. if invalid_mailbox_rotation or canary_rotation or rotation_success or fresh_domain_rotation:
  2000. consecutive_failures = 0
  2001. self._write_runtime_state()
  2002. if should_break_after_iteration:
  2003. break
  2004. except Exception as exc:
  2005. with self._lock:
  2006. self._total_attempts += 1
  2007. self._total_failure += 1
  2008. self._last_error = str(exc)
  2009. self._record_attempt(
  2010. success=False,
  2011. stage="exception",
  2012. error_message=str(exc),
  2013. metadata={"mail_provider": provider},
  2014. proxy_key=proxy_key,
  2015. email="",
  2016. )
  2017. consecutive_failures += 1
  2018. self._log(f"[zhuce6:register] [thread-{thread_id}] \u274c exception ({consecutive_failures}/{max_failures}): {exc}")
  2019. self._write_runtime_state()
  2020. finally:
  2021. self._release_cfmail_flow_slot(thread_id)
  2022. # Fallback: switch provider after N consecutive failures
  2023. if consecutive_failures >= max_failures and len(self._providers) > 1:
  2024. old_provider = provider
  2025. current_idx = self._providers.index(provider) if provider in self._providers else 0
  2026. provider = self._providers[(current_idx + 1) % len(self._providers)]
  2027. consecutive_failures = 0
  2028. self._log(
  2029. f"[zhuce6:register] [thread-{thread_id}] [fallback] switching {old_provider} -> {provider}"
  2030. )
  2031. # Random sleep between attempts
  2032. sleep_sec = random.randint(
  2033. self.settings.register_sleep_min,
  2034. max(self.settings.register_sleep_min, self.settings.register_sleep_max),
  2035. )
  2036. if self._stop_event.wait(sleep_sec) or self._target_reached.is_set():
  2037. break
  2038. def _cfmail_all_accounts_in_cooldown(self) -> bool:
  2039. manager = self._cfmail_manager
  2040. if manager is None:
  2041. return False
  2042. try:
  2043. manager.reload_if_needed()
  2044. except Exception:
  2045. pass
  2046. return manager.select_account() is None
  2047. def _current_cfmail_active_domain(self) -> str:
  2048. accounts = self._current_cfmail_active_accounts()
  2049. if accounts:
  2050. return str(accounts[0]["domain"]).strip().lower()
  2051. return ""
  2052. def _should_abort_cfmail_wait(self, account: Any) -> bool:
  2053. try:
  2054. extra = account.extra if hasattr(account, "extra") and isinstance(account.extra, dict) else {}
  2055. domain = str(extra.get("email_domain") or "").strip().lower()
  2056. if not domain:
  2057. email = str(getattr(account, "email", "") or "").strip().lower()
  2058. if "@" in email:
  2059. domain = email.rsplit("@", 1)[-1].strip().lower()
  2060. if not domain:
  2061. return False
  2062. stoploss = self._cfmail_wait_otp_stoploss_snapshot()
  2063. if not bool(stoploss.get("in_cooldown")):
  2064. return False
  2065. if str(stoploss.get("active_domain") or "").strip().lower() != domain:
  2066. return False
  2067. wait_started_at = float(extra.get("otp_wait_started_at") or 0.0)
  2068. triggered_at_raw = str(stoploss.get("last_triggered_at") or "").strip()
  2069. if wait_started_at > 0.0 and triggered_at_raw:
  2070. try:
  2071. triggered_at = datetime.fromisoformat(triggered_at_raw).timestamp()
  2072. except Exception:
  2073. triggered_at = 0.0
  2074. if triggered_at > 0.0 and wait_started_at < triggered_at:
  2075. return False
  2076. return True
  2077. except Exception:
  2078. return False
  2079. def _check_cfmail_all_cooldown_rotation(self, thread_id: int) -> bool:
  2080. """When all cfmail accounts are in cooldown, proactively trigger domain rotation."""
  2081. if self._cfmail_tracker is None or self._cfmail_provisioner is None:
  2082. return False
  2083. if not self._cfmail_all_accounts_in_cooldown():
  2084. return False
  2085. if not self._cfmail_rotation_lock.acquire(blocking=False):
  2086. return False
  2087. self._cfmail_rotation_pause.clear()
  2088. try:
  2089. self._log(
  2090. f"[zhuce6:register] [thread-{thread_id}] [cfmail] all accounts in cooldown, "
  2091. "forcing domain rotation to break deadlock"
  2092. )
  2093. provision_result = self._cfmail_provisioner.rotate_active_domain()
  2094. if not provision_result.success:
  2095. if provision_result.old_domain:
  2096. self._cfmail_tracker.mark_rotation_failed(
  2097. provision_result.old_domain,
  2098. provision_result.error,
  2099. )
  2100. self._log(
  2101. f"[zhuce6:register] [thread-{thread_id}] [cfmail] deadlock rotation failed: "
  2102. f"{provision_result.error}"
  2103. )
  2104. return False
  2105. self._cfmail_tracker.mark_rotation_completed(
  2106. provision_result.old_domain,
  2107. provision_result.new_domain,
  2108. )
  2109. self._reload_cfmail_manager_after_rotation()
  2110. self._reset_cfmail_add_phone_stoploss(provision_result.new_domain)
  2111. self._reset_cfmail_wait_otp_stoploss(provision_result.new_domain)
  2112. self._reset_cfmail_fresh_domain_budget(provision_result.new_domain)
  2113. self._arm_cfmail_canary(provision_result.new_domain)
  2114. self._log(
  2115. f"[zhuce6:register] [thread-{thread_id}] [cfmail] deadlock rotation completed: "
  2116. f"{provision_result.old_domain} -> {provision_result.new_domain}"
  2117. )
  2118. return True
  2119. finally:
  2120. self._cfmail_rotation_pause.set()
  2121. self._cfmail_rotation_lock.release()
  2122. def _current_cfmail_active_accounts(self) -> list[dict[str, str]]:
  2123. manager = self._cfmail_manager
  2124. accounts = list(getattr(manager, "accounts", []) or [])
  2125. active_accounts: list[dict[str, str]] = []
  2126. for item in accounts:
  2127. if not bool(getattr(item, "email_domain", "") or ""):
  2128. continue
  2129. if manager is not None and callable(getattr(manager, "skip_remaining_seconds", None)):
  2130. try:
  2131. if int(manager.skip_remaining_seconds(getattr(item, "name", "")) or 0) > 0:
  2132. continue
  2133. except Exception:
  2134. pass
  2135. active_accounts.append(
  2136. {
  2137. "name": str(getattr(item, "name", "") or "").strip(),
  2138. "domain": str(getattr(item, "email_domain", "") or "").strip().lower(),
  2139. }
  2140. )
  2141. return [item for item in active_accounts if item["name"] and item["domain"]]
  2142. def _cfmail_domain_pool_snapshot(
  2143. self,
  2144. recent_attempts: list[dict[str, object]],
  2145. ) -> dict[str, object]:
  2146. active_accounts = self._current_cfmail_active_accounts()
  2147. with self._lock:
  2148. inflight_by_thread = dict(self._cfmail_flow_state.get("inflight_by_thread") or {})
  2149. last_started_by_domain = dict(self._cfmail_flow_state.get("last_started_by_domain") or {})
  2150. replenish_thread = self._cfmail_replenish_thread
  2151. replenish_reason = self._cfmail_replenish_reason
  2152. domains: list[dict[str, object]] = []
  2153. now_ts = time.time()
  2154. manager = self._cfmail_manager
  2155. for account in active_accounts:
  2156. domain = account["domain"]
  2157. profile_name = account["name"]
  2158. domain_attempts = self._active_domain_attempts(recent_attempts, domain)
  2159. failure_by_stage, failure_signals = self._failure_counts_from_attempts(domain_attempts)
  2160. recent_success = sum(1 for item in domain_attempts if bool(item.get("success")))
  2161. recent_failure = sum(1 for item in domain_attempts if not bool(item.get("success")))
  2162. inflight = sum(
  2163. 1
  2164. for value in inflight_by_thread.values()
  2165. if isinstance(value, dict)
  2166. and str(value.get("domain") or "").strip().lower() == domain
  2167. )
  2168. last_started_at = float(last_started_by_domain.get(domain) or 0.0)
  2169. start_interval_remaining = 0.0
  2170. if self._cfmail_start_interval_seconds > 0 and last_started_at > 0:
  2171. start_interval_remaining = max(
  2172. 0.0,
  2173. (last_started_at + self._cfmail_start_interval_seconds) - now_ts,
  2174. )
  2175. skip_remaining_seconds = 0
  2176. if manager is not None and callable(getattr(manager, "skip_remaining_seconds", None)):
  2177. try:
  2178. skip_remaining_seconds = int(manager.skip_remaining_seconds(profile_name) or 0)
  2179. except Exception:
  2180. skip_remaining_seconds = 0
  2181. add_phone_state = self._cfmail_add_phone_stoploss_snapshot()
  2182. wait_otp_state = self._cfmail_wait_otp_stoploss_snapshot()
  2183. domains.append(
  2184. {
  2185. "name": profile_name,
  2186. "domain": domain,
  2187. "inflight": inflight,
  2188. "recent_attempts": len(domain_attempts),
  2189. "recent_success": recent_success,
  2190. "recent_failure": recent_failure,
  2191. "failure_by_stage": failure_by_stage,
  2192. "failure_signals": failure_signals,
  2193. "recent_failure_hotspots": self._recent_failure_hotspots_from_attempts(domain_attempts),
  2194. "last_started_at": datetime.fromtimestamp(last_started_at).isoformat(timespec="seconds")
  2195. if last_started_at > 0
  2196. else "",
  2197. "start_interval_remaining_seconds": round(start_interval_remaining, 1),
  2198. "skip_remaining_seconds": skip_remaining_seconds,
  2199. "add_phone_cooldown": bool(
  2200. add_phone_state.get("in_cooldown")
  2201. and str(add_phone_state.get("active_domain") or "").strip().lower() == domain
  2202. ),
  2203. "wait_otp_cooldown": bool(
  2204. wait_otp_state.get("in_cooldown")
  2205. and str(wait_otp_state.get("active_domain") or "").strip().lower() == domain
  2206. ),
  2207. }
  2208. )
  2209. return {
  2210. "target_count": self._cfmail_active_domain_count,
  2211. "active_count": len(domains),
  2212. "active_domains": domains,
  2213. "replenishing": bool(replenish_thread is not None and replenish_thread.is_alive()),
  2214. "replenish_reason": str(replenish_reason or ""),
  2215. }
  2216. def _selected_cfmail_profile(self, thread_id: int) -> str:
  2217. with self._lock:
  2218. selected_profile_by_thread = self._cfmail_flow_state.setdefault("selected_profile_by_thread", {})
  2219. if not isinstance(selected_profile_by_thread, dict):
  2220. return "auto"
  2221. value = str(selected_profile_by_thread.get(thread_id) or "").strip()
  2222. return value or "auto"
  2223. def _handle_cfmail_rotation(
  2224. self,
  2225. *,
  2226. thread_id: int,
  2227. result: dict[str, object],
  2228. proxy_key: str,
  2229. ) -> bool:
  2230. if self._cfmail_tracker is None or self._cfmail_provisioner is None:
  2231. return False
  2232. metadata = result.get("metadata")
  2233. if not isinstance(metadata, dict):
  2234. metadata = {}
  2235. if str(metadata.get("mail_provider") or result.get("mail_provider") or "").strip().lower() not in {"", "cfmail"}:
  2236. return False
  2237. from core.cfmail_domain_rotation import classify_domain_attempt
  2238. attempt = classify_domain_attempt(result, proxy_key=proxy_key)
  2239. if attempt is None:
  2240. return False
  2241. decision = self._cfmail_tracker.record(attempt)
  2242. if attempt.backend_failure and not decision.should_rotate:
  2243. self._log(
  2244. f"[zhuce6:register] [thread-{thread_id}] [cfmail] backend unhealthy for domain={attempt.domain}; rotation skipped"
  2245. )
  2246. return False
  2247. if not decision.should_rotate:
  2248. return False
  2249. self._log(
  2250. f"[zhuce6:register] [thread-{thread_id}] [cfmail] rotating domain {decision.domain} "
  2251. f"(reason={decision.reason}, failures={decision.blacklist_failures}, window={decision.window_size})"
  2252. )
  2253. return self._replace_cfmail_domain(
  2254. thread_id=thread_id,
  2255. domain=decision.domain,
  2256. reason_label=decision.reason,
  2257. )
  2258. def _classify_proxy_outcome(self, result: dict[str, object]) -> bool | None:
  2259. if bool(result.get("success")):
  2260. return True
  2261. stage = str(result.get("stage") or "").strip().lower()
  2262. metadata = result.get("metadata") if isinstance(result.get("metadata"), dict) else {}
  2263. code = str(metadata.get("create_account_error_code") or "").strip().lower()
  2264. post_create_gate = str(metadata.get("post_create_gate") or "").strip().lower()
  2265. if stage == "create_account" and code in {"registration_disallowed", "unsupported_email"}:
  2266. return None
  2267. if stage == "add_phone_gate" or post_create_gate == "add_phone":
  2268. return None
  2269. if stage in {"mailbox", "device_id", "password", "token_acquisition"}:
  2270. return False
  2271. return False
  2272. def snapshot(self) -> dict[str, object]:
  2273. with self._lock:
  2274. alive = sum(1 for t in self._threads if t.is_alive())
  2275. target = self.settings.register_target_count
  2276. cfmail_rotation = self._cfmail_tracker.snapshot() if self._cfmail_tracker is not None else None
  2277. recent_attempts = list(self._recent_attempts)
  2278. add_phone_stoploss = self._cfmail_add_phone_stoploss_snapshot()
  2279. wait_otp_stoploss = self._cfmail_wait_otp_stoploss_snapshot()
  2280. cfmail_canary = self._cfmail_canary_snapshot()
  2281. cfmail_fresh_domain_budget = self._cfmail_fresh_domain_budget_snapshot()
  2282. cfmail_domain_pool = self._cfmail_domain_pool_snapshot(recent_attempts)
  2283. active_domain = self._infer_active_domain(recent_attempts, cfmail_rotation, add_phone_stoploss)
  2284. active_domain_attempts = self._active_domain_attempts(recent_attempts, active_domain)
  2285. active_failure_by_stage, active_failure_signals = self._failure_counts_from_attempts(active_domain_attempts)
  2286. return {
  2287. "name": "register",
  2288. "status": "running" if alive > 0 else ("pending" if self._total_attempts == 0 else "stopped"),
  2289. "threads_alive": alive,
  2290. "threads_total": len(self._threads),
  2291. "total_attempts": self._total_attempts,
  2292. "total_success": self._total_success,
  2293. "total_success_registered": self._total_success,
  2294. "total_success_direct": self._total_success,
  2295. "total_warmup_pending": self._total_warmup_pending,
  2296. "total_cpa_sync_success": self._total_cpa_sync_success,
  2297. "total_cpa_sync_success_direct": self._total_cpa_sync_success,
  2298. "total_cpa_sync_failure": self._total_cpa_sync_failure,
  2299. "total_failure": self._total_failure,
  2300. "success_rate": round(self._total_success / max(self._total_attempts, 1) * 100, 1),
  2301. "registered_success_rate": round(self._total_success / max(self._total_attempts, 1) * 100, 1),
  2302. "cpa_sync_success_rate": round(self._total_cpa_sync_success / max(self._total_attempts, 1) * 100, 1),
  2303. "target_count": target if target > 0 else None,
  2304. "target_reached": self._target_reached.is_set(),
  2305. "last_error": self._last_error,
  2306. "proxy": self.settings.register_proxy,
  2307. "proxy_pool_enabled": self._proxy_pool is not None,
  2308. "mail_provider": self.settings.register_mail_provider,
  2309. "interval_seconds": self.settings.register_interval,
  2310. "run_count": self._total_attempts,
  2311. "success_count": self._total_success,
  2312. "failure_count": self._total_failure,
  2313. "is_running": alive > 0,
  2314. "last_started_at": datetime.fromtimestamp(self._started_at).isoformat(timespec="seconds") if self._started_at else None,
  2315. "last_finished_at": None,
  2316. "last_duration_seconds": None,
  2317. "next_run_at": None,
  2318. "failure_by_stage": dict(sorted(self._failure_by_stage.items(), key=lambda item: (-item[1], item[0]))),
  2319. "failure_signals": dict(sorted(self._failure_signals.items(), key=lambda item: (-item[1], item[0]))),
  2320. "recent_failure_hotspots": self._recent_failure_hotspots(),
  2321. "recent_attempts": recent_attempts,
  2322. "active_domain_recent_attempts": active_domain_attempts,
  2323. "active_domain_failure_by_stage": active_failure_by_stage,
  2324. "active_domain_failure_signals": active_failure_signals,
  2325. "active_domain_recent_failure_hotspots": self._recent_failure_hotspots_from_attempts(active_domain_attempts),
  2326. "cfmail_rotation": cfmail_rotation,
  2327. "cfmail_add_phone_stoploss": add_phone_stoploss,
  2328. "cfmail_wait_otp_stoploss": wait_otp_stoploss,
  2329. "cfmail_canary": cfmail_canary,
  2330. "cfmail_fresh_domain_budget": cfmail_fresh_domain_budget,
  2331. "cfmail_domain_pool": cfmail_domain_pool,
  2332. "pending_token_queue": {
  2333. "queue_size": len(self._pending_token_queue),
  2334. "total_enqueued": self._pending_token_total_enqueued,
  2335. "total_success": self._pending_token_total_success,
  2336. "total_failed": self._pending_token_total_failed,
  2337. "retry_delay_seconds": self._pending_token_retry_delay_seconds,
  2338. "max_retries": self._pending_token_max_retries,
  2339. },
  2340. }
  2341. class RegistrationBurstScheduler:
  2342. """Run registration in timed batches and expose scheduler state via runtime_state.json."""
  2343. def __init__(self, settings: AppSettings) -> None:
  2344. self.settings = settings
  2345. self._stop_event = threading.Event()
  2346. self._lock = threading.RLock()
  2347. self._active_loop: RegistrationLoop | None = None
  2348. self._active_started_at: float | None = None
  2349. self._next_run_at_ts: float | None = None
  2350. self._run_count = 0
  2351. self._total_attempts = 0
  2352. self._total_success = 0
  2353. self._total_cpa_sync_success = 0
  2354. self._total_cpa_sync_failure = 0
  2355. self._total_failure = 0
  2356. self._last_error: str | None = None
  2357. self._recent_attempts: deque[dict[str, object]] = deque(maxlen=80)
  2358. self._failure_by_stage: dict[str, int] = {}
  2359. self._failure_signals: dict[str, int] = {}
  2360. self._last_batch_started_at: str | None = None
  2361. self._last_batch_finished_at: str | None = None
  2362. self._last_batch_duration_seconds: float | None = None
  2363. self._last_cfmail_add_phone_stoploss: dict[str, object] | None = None
  2364. self._last_cfmail_wait_otp_stoploss: dict[str, object] | None = None
  2365. self._last_proxy_pool_snapshot: dict[str, object] = {
  2366. "configured": bool(self.settings.proxy_pool_configured),
  2367. "enabled": False,
  2368. "snapshot_error": None,
  2369. "node_count": 0,
  2370. "in_use_count": 0,
  2371. "disabled_count": 0,
  2372. "nodes": [],
  2373. }
  2374. self._logger = self._setup_logger()
  2375. def _setup_logger(self) -> Any:
  2376. from logging.handlers import RotatingFileHandler
  2377. logger = logging.getLogger("zhuce6.register")
  2378. logger.setLevel(logging.INFO)
  2379. logger.propagate = False
  2380. if not logger.handlers:
  2381. console = logging.StreamHandler()
  2382. console.setFormatter(logging.Formatter("%(message)s"))
  2383. logger.addHandler(console)
  2384. if self.settings.register_log_file:
  2385. fh = RotatingFileHandler(
  2386. self.settings.register_log_file,
  2387. maxBytes=2 * 1024 * 1024,
  2388. backupCount=5,
  2389. encoding="utf-8",
  2390. )
  2391. fh.setFormatter(logging.Formatter("%(asctime)s %(message)s", datefmt="%Y-%m-%d %H:%M:%S"))
  2392. logger.addHandler(fh)
  2393. return logger
  2394. def _log(self, msg: str) -> None:
  2395. self._logger.info(msg)
  2396. def _merge_counts(self, target: dict[str, int], incoming: dict[str, object] | None) -> None:
  2397. if not isinstance(incoming, dict):
  2398. return
  2399. for key, value in incoming.items():
  2400. try:
  2401. inc = int(value or 0)
  2402. except Exception:
  2403. continue
  2404. target[str(key)] = target.get(str(key), 0) + inc
  2405. def _write_runtime_state(self) -> None:
  2406. state_file = Path(self.settings.runtime_state_file)
  2407. try:
  2408. state_file.parent.mkdir(parents=True, exist_ok=True)
  2409. payload = {
  2410. "updated_at": datetime.now().isoformat(timespec="seconds"),
  2411. "register_snapshot": self.snapshot(),
  2412. "proxy_pool": self._current_proxy_pool_snapshot(),
  2413. }
  2414. tmp_file = state_file.with_name(
  2415. f"{state_file.name}.{os.getpid()}.{threading.get_ident()}.tmp"
  2416. )
  2417. tmp_file.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8")
  2418. tmp_file.replace(state_file)
  2419. except Exception as exc:
  2420. self._log(f"[zhuce6:register] burst runtime state write failed: {exc}")
  2421. def _current_proxy_pool_snapshot(self) -> dict[str, object]:
  2422. with self._lock:
  2423. active_loop = self._active_loop
  2424. last_snapshot = dict(self._last_proxy_pool_snapshot)
  2425. if active_loop is not None:
  2426. try:
  2427. return active_loop._proxy_pool_snapshot()
  2428. except Exception as exc:
  2429. last_snapshot["snapshot_error"] = str(exc)
  2430. return last_snapshot
  2431. return last_snapshot
  2432. def _absorb_batch_snapshot(self, snapshot: dict[str, object], *, duration_seconds: float) -> None:
  2433. with self._lock:
  2434. self._run_count += 1
  2435. self._total_attempts += int(snapshot.get("total_attempts") or 0)
  2436. self._total_success += int(snapshot.get("total_success") or 0)
  2437. self._total_cpa_sync_success += int(snapshot.get("total_cpa_sync_success") or 0)
  2438. self._total_cpa_sync_failure += int(snapshot.get("total_cpa_sync_failure") or 0)
  2439. self._total_failure += int(snapshot.get("total_failure") or 0)
  2440. self._last_error = str(snapshot.get("last_error") or "").strip() or None
  2441. self._last_batch_started_at = snapshot.get("last_started_at") if isinstance(snapshot.get("last_started_at"), str) else None
  2442. self._last_batch_finished_at = datetime.now().isoformat(timespec="seconds")
  2443. self._last_batch_duration_seconds = round(duration_seconds, 3)
  2444. self._merge_counts(self._failure_by_stage, snapshot.get("failure_by_stage") if isinstance(snapshot.get("failure_by_stage"), dict) else None)
  2445. self._merge_counts(self._failure_signals, snapshot.get("failure_signals") if isinstance(snapshot.get("failure_signals"), dict) else None)
  2446. attempts = snapshot.get("recent_attempts")
  2447. if isinstance(attempts, list):
  2448. for item in attempts:
  2449. if isinstance(item, dict):
  2450. self._recent_attempts.append(item)
  2451. if isinstance(snapshot.get("cfmail_add_phone_stoploss"), dict):
  2452. self._last_cfmail_add_phone_stoploss = dict(snapshot.get("cfmail_add_phone_stoploss") or {})
  2453. if isinstance(snapshot.get("cfmail_wait_otp_stoploss"), dict):
  2454. self._last_cfmail_wait_otp_stoploss = dict(snapshot.get("cfmail_wait_otp_stoploss") or {})
  2455. def snapshot(self) -> dict[str, object]:
  2456. with self._lock:
  2457. active_loop = self._active_loop
  2458. next_run_at_ts = self._next_run_at_ts
  2459. total_attempts = self._total_attempts
  2460. total_success = self._total_success
  2461. total_cpa_sync_success = self._total_cpa_sync_success
  2462. total_cpa_sync_failure = self._total_cpa_sync_failure
  2463. total_failure = self._total_failure
  2464. run_count = self._run_count
  2465. failure_by_stage = dict(sorted(self._failure_by_stage.items(), key=lambda item: (-item[1], item[0])))
  2466. failure_signals = dict(sorted(self._failure_signals.items(), key=lambda item: (-item[1], item[0])))
  2467. recent_attempts = list(self._recent_attempts)
  2468. last_error = self._last_error
  2469. last_batch_started_at = self._last_batch_started_at
  2470. last_batch_finished_at = self._last_batch_finished_at
  2471. last_batch_duration_seconds = self._last_batch_duration_seconds
  2472. add_phone_stoploss = dict(self._last_cfmail_add_phone_stoploss or {})
  2473. wait_otp_stoploss = dict(self._last_cfmail_wait_otp_stoploss or {})
  2474. if active_loop is not None:
  2475. current = dict(active_loop.snapshot())
  2476. current.update(
  2477. {
  2478. "scheduler_mode": "burst",
  2479. "batch_threads": self.settings.register_batch_threads,
  2480. "batch_target_count": self.settings.register_batch_target_count,
  2481. "batch_interval_seconds": self.settings.register_batch_interval_seconds,
  2482. "run_count": run_count,
  2483. "next_run_at": None,
  2484. }
  2485. )
  2486. return current
  2487. status = "scheduled" if next_run_at_ts and not self._stop_event.is_set() else ("stopped" if run_count > 0 or self._stop_event.is_set() else "pending")
  2488. counts: dict[tuple[str, str], int] = {}
  2489. for item in recent_attempts:
  2490. if item.get("success"):
  2491. continue
  2492. stage = str(item.get("stage") or "?").strip() or "?"
  2493. signal = str(item.get("signal") or "").strip()
  2494. key = (signal or stage, stage)
  2495. counts[key] = counts.get(key, 0) + 1
  2496. recent_failure_hotspots = [
  2497. {"key": key, "stage": stage, "count": count}
  2498. for (key, stage), count in sorted(
  2499. counts.items(),
  2500. key=lambda kv: (-kv[1], kv[0][0], kv[0][1]),
  2501. )[:5]
  2502. ]
  2503. return {
  2504. "name": "register",
  2505. "status": status,
  2506. "scheduler_mode": "burst",
  2507. "threads_alive": 0,
  2508. "threads_total": self.settings.register_batch_threads,
  2509. "total_attempts": total_attempts,
  2510. "total_success": total_success,
  2511. "total_success_registered": total_success,
  2512. "total_success_direct": total_success,
  2513. "total_cpa_sync_success": total_cpa_sync_success,
  2514. "total_cpa_sync_success_direct": total_cpa_sync_success,
  2515. "total_cpa_sync_failure": total_cpa_sync_failure,
  2516. "total_failure": total_failure,
  2517. "success_rate": round(total_success / max(total_attempts, 1) * 100, 1),
  2518. "registered_success_rate": round(total_success / max(total_attempts, 1) * 100, 1),
  2519. "cpa_sync_success_rate": round(total_cpa_sync_success / max(total_attempts, 1) * 100, 1),
  2520. "target_count": self.settings.register_batch_target_count,
  2521. "target_reached": False,
  2522. "last_error": last_error,
  2523. "proxy": self.settings.register_proxy,
  2524. "proxy_pool_enabled": bool(self.settings.proxy_pool_configured),
  2525. "mail_provider": self.settings.register_mail_provider,
  2526. "interval_seconds": self.settings.register_interval,
  2527. "run_count": run_count,
  2528. "success_count": total_success,
  2529. "failure_count": total_failure,
  2530. "is_running": False,
  2531. "last_started_at": last_batch_started_at,
  2532. "last_finished_at": last_batch_finished_at,
  2533. "last_duration_seconds": last_batch_duration_seconds,
  2534. "next_run_at": datetime.fromtimestamp(next_run_at_ts).isoformat(timespec="seconds") if next_run_at_ts else None,
  2535. "failure_by_stage": failure_by_stage,
  2536. "failure_signals": failure_signals,
  2537. "recent_failure_hotspots": recent_failure_hotspots,
  2538. "recent_attempts": recent_attempts,
  2539. "active_domain_recent_attempts": [],
  2540. "active_domain_failure_by_stage": {},
  2541. "active_domain_failure_signals": {},
  2542. "active_domain_recent_failure_hotspots": [],
  2543. "cfmail_rotation": None,
  2544. "cfmail_add_phone_stoploss": add_phone_stoploss,
  2545. "cfmail_wait_otp_stoploss": wait_otp_stoploss,
  2546. "batch_threads": self.settings.register_batch_threads,
  2547. "batch_target_count": self.settings.register_batch_target_count,
  2548. "batch_interval_seconds": self.settings.register_batch_interval_seconds,
  2549. }
  2550. def stop(self) -> None:
  2551. self._stop_event.set()
  2552. with self._lock:
  2553. active_loop = self._active_loop
  2554. if active_loop is not None:
  2555. active_loop.stop()
  2556. self._write_runtime_state()
  2557. def run(self) -> None:
  2558. self._next_run_at_ts = time.time()
  2559. self._write_runtime_state()
  2560. while not self._stop_event.is_set():
  2561. now = time.time()
  2562. next_run_at_ts = self._next_run_at_ts or now
  2563. if now < next_run_at_ts:
  2564. self._write_runtime_state()
  2565. self._stop_event.wait(min(max(next_run_at_ts - now, 1), 5))
  2566. continue
  2567. batch_settings = replace(
  2568. self.settings,
  2569. register_enabled=True,
  2570. register_threads=self.settings.register_batch_threads,
  2571. register_target_count=self.settings.register_batch_target_count,
  2572. )
  2573. loop = _compat_main_attr("RegistrationLoop", RegistrationLoop)(batch_settings)
  2574. started_at = time.time()
  2575. with self._lock:
  2576. self._active_loop = loop
  2577. self._active_started_at = started_at
  2578. self._write_runtime_state()
  2579. loop.start()
  2580. try:
  2581. while not self._stop_event.is_set():
  2582. snapshot = loop.snapshot()
  2583. if int(snapshot.get("threads_alive") or 0) <= 0:
  2584. break
  2585. self._write_runtime_state()
  2586. self._stop_event.wait(1)
  2587. finally:
  2588. loop.stop()
  2589. batch_snapshot = loop.snapshot()
  2590. with self._lock:
  2591. self._active_loop = None
  2592. self._active_started_at = None
  2593. proxy_snapshot_fn = getattr(loop, "_proxy_pool_snapshot", None)
  2594. if callable(proxy_snapshot_fn):
  2595. self._last_proxy_pool_snapshot = proxy_snapshot_fn()
  2596. self._absorb_batch_snapshot(batch_snapshot, duration_seconds=time.time() - started_at)
  2597. self._next_run_at_ts = started_at + self.settings.register_batch_interval_seconds
  2598. self._write_runtime_state()
  2599. if self._stop_event.is_set():
  2600. break