| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188218921902191219221932194219521962197219821992200220122022203220422052206220722082209221022112212221322142215221622172218221922202221222222232224222522262227222822292230223122322233223422352236223722382239224022412242224322442245224622472248224922502251225222532254225522562257225822592260226122622263226422652266226722682269227022712272227322742275227622772278227922802281228222832284228522862287228822892290229122922293229422952296229722982299230023012302230323042305230623072308230923102311231223132314231523162317231823192320232123222323232423252326232723282329233023312332233323342335233623372338233923402341234223432344234523462347234823492350235123522353235423552356235723582359236023612362236323642365236623672368236923702371237223732374237523762377237823792380238123822383238423852386238723882389239023912392239323942395239623972398239924002401240224032404240524062407240824092410241124122413241424152416241724182419242024212422242324242425242624272428242924302431243224332434243524362437243824392440244124422443244424452446244724482449245024512452245324542455245624572458245924602461246224632464246524662467246824692470247124722473247424752476247724782479248024812482248324842485248624872488248924902491249224932494249524962497249824992500250125022503250425052506250725082509251025112512251325142515251625172518251925202521252225232524252525262527252825292530253125322533253425352536253725382539254025412542254325442545254625472548254925502551255225532554255525562557255825592560256125622563256425652566256725682569257025712572257325742575257625772578257925802581258225832584258525862587258825892590259125922593259425952596259725982599260026012602260326042605260626072608260926102611261226132614261526162617261826192620262126222623262426252626262726282629263026312632263326342635263626372638263926402641264226432644264526462647264826492650265126522653265426552656265726582659266026612662266326642665266626672668266926702671267226732674267526762677267826792680268126822683268426852686268726882689269026912692269326942695269626972698269927002701270227032704270527062707 |
- """Registration runtime loops for zhuce6."""
- from __future__ import annotations
- from collections import deque
- from dataclasses import replace
- from datetime import datetime
- import json
- import logging
- import os
- from pathlib import Path
- import random
- import sys
- import threading
- import time
- from typing import Any
- from urllib.parse import urlsplit, urlunsplit
- from core.registry import load_all
- from core.settings import AppSettings
- from dashboard.api import _count_cpa_files, _fetch_management_auth_files, _is_regular_free_account
- from ops.common import create_backend_client, get_management_key
- from ops.rotate_runtime import _maybe_reconcile_cpa_runtime
- from platforms.chatgpt.fingerprint import build_registration_provenance, infer_proxy_region
- from platforms.chatgpt.pool import is_warmup_pending_record, now_iso, update_token_record
- from core.chatgpt_flow_runner import run_chatgpt_register_once
- DEFAULT_ADD_PHONE_STOPLOSS_WINDOW = 10
- DEFAULT_ADD_PHONE_STOPLOSS_THRESHOLD = 3
- DEFAULT_ADD_PHONE_STOPLOSS_COOLDOWN_SECONDS = 300
- DEFAULT_ADD_PHONE_STOPLOSS_MAX_SUCCESSES = 2
- DEFAULT_WAIT_OTP_STOPLOSS_WINDOW = 6
- DEFAULT_WAIT_OTP_STOPLOSS_THRESHOLD = 2
- DEFAULT_WAIT_OTP_STOPLOSS_COOLDOWN_SECONDS = 300
- DEFAULT_WAIT_OTP_LIVE_ABORT_THRESHOLD = 0
- DEFAULT_WAIT_OTP_LIVE_ABORT_AGE_SECONDS = 90
- DEFAULT_CFMAIL_FRESH_DOMAIN_ATTEMPT_BUDGET = 2
- def _classify_token_file(*args, **kwargs): # type: ignore[no-untyped-def]
- from ops.scan import classify_token_file
- return classify_token_file(*args, **kwargs)
- def _compat_main_attr(name: str, default: object) -> object:
- main_module = sys.modules.get("main")
- if main_module is None:
- return default
- return getattr(main_module, name, default)
- class RegistrationLoop:
- """Multi-threaded continuous registration with fallback, target count, and logging."""
- def __init__(self, settings: AppSettings) -> None:
- self.settings = settings
- self._threads: list[threading.Thread] = []
- self._stop_event = threading.Event()
- self._target_reached = threading.Event()
- self._lock = threading.RLock()
- self._total_attempts = 0
- self._total_success = 0
- self._total_warmup_pending = 0
- self._total_cpa_sync_success = 0
- self._total_cpa_sync_failure = 0
- self._total_failure = 0
- self._last_error: str | None = None
- self._started_at: float | None = None
- self._failure_by_stage: dict[str, int] = {}
- self._failure_signals: dict[str, int] = {}
- self._recent_attempts: deque[dict[str, object]] = deque(maxlen=80)
- self._providers: list[str] = []
- self._proxy_pool = None
- self._logger = self._setup_logger()
- self._cfmail_tracker = None
- self._cfmail_provisioner = None
- self._cfmail_manager: Any = None
- self._cfmail_rotation_lock = threading.Lock()
- self._cfmail_rotation_pause = threading.Event()
- self._cfmail_rotation_pause.set()
- self._cpa_management_key_cache: str | None | bool = False
- self._cfmail_add_phone_window = max(
- 1,
- int(str(os.getenv("ZHUCE6_CFMAIL_ADD_PHONE_WINDOW", DEFAULT_ADD_PHONE_STOPLOSS_WINDOW)).strip() or str(DEFAULT_ADD_PHONE_STOPLOSS_WINDOW)),
- )
- self._cfmail_add_phone_threshold = max(
- 1,
- int(
- str(
- os.getenv(
- "ZHUCE6_CFMAIL_ADD_PHONE_THRESHOLD",
- DEFAULT_ADD_PHONE_STOPLOSS_THRESHOLD,
- )
- ).strip()
- or str(DEFAULT_ADD_PHONE_STOPLOSS_THRESHOLD)
- ),
- )
- self._cfmail_add_phone_cooldown_seconds = max(
- 1,
- int(
- str(
- os.getenv(
- "ZHUCE6_CFMAIL_ADD_PHONE_COOLDOWN_SECONDS",
- DEFAULT_ADD_PHONE_STOPLOSS_COOLDOWN_SECONDS,
- )
- ).strip()
- or str(DEFAULT_ADD_PHONE_STOPLOSS_COOLDOWN_SECONDS)
- ),
- )
- self._cfmail_add_phone_max_successes = max(
- 0,
- int(
- str(
- os.getenv(
- "ZHUCE6_CFMAIL_ADD_PHONE_MAX_SUCCESSES",
- DEFAULT_ADD_PHONE_STOPLOSS_MAX_SUCCESSES,
- )
- ).strip()
- or str(DEFAULT_ADD_PHONE_STOPLOSS_MAX_SUCCESSES)
- ),
- )
- self._cfmail_add_phone_events: dict[str, deque[dict[str, object]]] = {}
- self._cfmail_add_phone_state: dict[str, object] = {
- "active_domain": "",
- "in_cooldown": False,
- "cooldown_until": 0.0,
- "last_triggered_at": "",
- "last_rotation_attempted_at": "",
- "last_reason": "",
- "last_add_phone_failures": 0,
- "last_successes": 0,
- "last_window_size": 0,
- "last_logged_at": 0.0,
- }
- self._cfmail_wait_otp_window = max(
- 1,
- int(str(os.getenv("ZHUCE6_CFMAIL_WAIT_OTP_WINDOW", DEFAULT_WAIT_OTP_STOPLOSS_WINDOW)).strip() or str(DEFAULT_WAIT_OTP_STOPLOSS_WINDOW)),
- )
- self._cfmail_wait_otp_threshold = max(
- 1,
- int(
- str(
- os.getenv(
- "ZHUCE6_CFMAIL_WAIT_OTP_THRESHOLD",
- DEFAULT_WAIT_OTP_STOPLOSS_THRESHOLD,
- )
- ).strip()
- or str(DEFAULT_WAIT_OTP_STOPLOSS_THRESHOLD)
- ),
- )
- self._cfmail_wait_otp_cooldown_seconds = max(
- 0,
- int(
- str(
- os.getenv(
- "ZHUCE6_CFMAIL_WAIT_OTP_COOLDOWN_SECONDS",
- DEFAULT_WAIT_OTP_STOPLOSS_COOLDOWN_SECONDS,
- )
- ).strip()
- or str(DEFAULT_WAIT_OTP_STOPLOSS_COOLDOWN_SECONDS)
- ),
- )
- self._cfmail_wait_otp_events: dict[str, deque[dict[str, object]]] = {}
- self._cfmail_wait_otp_state: dict[str, object] = {
- "active_domain": "",
- "in_cooldown": False,
- "cooldown_until": 0.0,
- "last_triggered_at": "",
- "last_rotation_attempted_at": "",
- "last_reason": "",
- "last_no_message_timeouts": 0,
- "last_window_size": 0,
- "last_logged_at": 0.0,
- }
- try:
- self._cfmail_wait_otp_live_threshold = max(
- 0,
- int(
- str(
- os.getenv(
- "ZHUCE6_CFMAIL_WAIT_OTP_LIVE_ABORT_THRESHOLD",
- DEFAULT_WAIT_OTP_LIVE_ABORT_THRESHOLD,
- )
- ).strip()
- or str(DEFAULT_WAIT_OTP_LIVE_ABORT_THRESHOLD)
- ),
- )
- except Exception:
- self._cfmail_wait_otp_live_threshold = DEFAULT_WAIT_OTP_LIVE_ABORT_THRESHOLD
- try:
- self._cfmail_wait_otp_live_age_seconds = max(
- 30,
- int(
- str(
- os.getenv(
- "ZHUCE6_CFMAIL_WAIT_OTP_LIVE_ABORT_AGE_SECONDS",
- DEFAULT_WAIT_OTP_LIVE_ABORT_AGE_SECONDS,
- )
- ).strip()
- or str(DEFAULT_WAIT_OTP_LIVE_ABORT_AGE_SECONDS)
- ),
- )
- except Exception:
- self._cfmail_wait_otp_live_age_seconds = DEFAULT_WAIT_OTP_LIVE_ABORT_AGE_SECONDS
- self._cfmail_wait_otp_live_lock = threading.RLock()
- self._cfmail_wait_otp_live_progress: dict[str, dict[str, dict[str, object]]] = {}
- self._cfmail_canary_state: dict[str, object] = {
- "active_domain": "",
- "pending": False,
- "owner_thread_id": 0,
- "attempt_started_at": 0.0,
- "last_ready_at": "",
- "last_ready_reason": "",
- "last_logged_at": 0.0,
- }
- try:
- self._cfmail_start_interval_seconds = max(
- 0,
- int(str(os.getenv("ZHUCE6_CFMAIL_START_INTERVAL_SECONDS", "8")).strip() or "8"),
- )
- except Exception:
- self._cfmail_start_interval_seconds = 8
- try:
- self._cfmail_max_inflight = max(
- 1,
- int(str(os.getenv("ZHUCE6_CFMAIL_MAX_INFLIGHT", "4")).strip() or "4"),
- )
- except Exception:
- self._cfmail_max_inflight = 4
- self._cfmail_flow_state: dict[str, object] = {
- "inflight_by_thread": {},
- "last_started_by_domain": {},
- "selected_profile_by_thread": {},
- "last_logged_at": 0.0,
- }
- try:
- self._cfmail_active_domain_count = max(
- 1,
- int(str(os.getenv("ZHUCE6_CFMAIL_ACTIVE_DOMAIN_COUNT", "3")).strip() or "3"),
- )
- except Exception:
- self._cfmail_active_domain_count = 3
- try:
- self._cfmail_fresh_domain_attempt_budget = max(
- 0,
- int(
- str(
- os.getenv(
- "ZHUCE6_CFMAIL_FRESH_DOMAIN_ATTEMPT_BUDGET",
- DEFAULT_CFMAIL_FRESH_DOMAIN_ATTEMPT_BUDGET,
- )
- ).strip()
- or str(DEFAULT_CFMAIL_FRESH_DOMAIN_ATTEMPT_BUDGET)
- ),
- )
- except Exception:
- self._cfmail_fresh_domain_attempt_budget = DEFAULT_CFMAIL_FRESH_DOMAIN_ATTEMPT_BUDGET
- self._cfmail_fresh_domain_state: dict[str, object] = {
- "active_domain": "",
- "completed_attempts": 0,
- "mail_seen_attempts": 0,
- "successes": 0,
- "last_triggered_at": "",
- "last_rotation_attempted_at": "",
- "last_reason": "",
- }
- # Solution B: deferred retry queue for add_phone_gate accounts
- self._pending_token_queue: list[dict[str, Any]] = []
- self._pending_token_lock = threading.Lock()
- try:
- raw_pending_retry_delay = int(
- str(os.getenv("ZHUCE6_PENDING_TOKEN_RETRY_DELAY_SECONDS", "600")).strip() or "600"
- )
- except Exception:
- raw_pending_retry_delay = 600
- # add_phone accounts are only useful if the deferred retry happens inside the
- # same short observation window as the registration loop. Cap the first retry
- # base delay so a large env value cannot postpone every retry past 5 minutes.
- self._pending_token_retry_delay_seconds = max(60, min(raw_pending_retry_delay, 60))
- try:
- self._pending_token_max_retries = max(
- 1,
- int(str(os.getenv("ZHUCE6_PENDING_TOKEN_MAX_RETRIES", "3")).strip() or "3"),
- )
- except Exception:
- self._pending_token_max_retries = 3
- self._pending_token_total_enqueued = 0
- self._pending_token_total_success = 0
- self._pending_token_total_failed = 0
- self._cfmail_replenish_thread: threading.Thread | None = None
- self._cfmail_replenish_reason = ""
- def _write_runtime_state(self) -> None:
- state_file = Path(self.settings.runtime_state_file)
- try:
- state_file.parent.mkdir(parents=True, exist_ok=True)
- payload = {
- "updated_at": datetime.now().isoformat(timespec="seconds"),
- "register_snapshot": self.snapshot(),
- "proxy_pool": self._proxy_pool_snapshot(),
- }
- tmp_file = state_file.with_name(
- f"{state_file.name}.{os.getpid()}.{threading.get_ident()}.tmp"
- )
- tmp_file.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8")
- tmp_file.replace(state_file)
- except Exception as exc:
- self._log(f"[zhuce6:register] runtime state write failed: {exc}")
- def _cfmail_canary_snapshot(self) -> dict[str, object]:
- return {
- "active_domain": "",
- "pending": False,
- "owner_thread_id": 0,
- "attempt_started_at": 0.0,
- "last_ready_at": "",
- "last_ready_reason": "disabled",
- }
- def _arm_cfmail_canary(self, domain: str, *, pending: bool = True) -> None:
- del domain, pending
- return
- def _mark_cfmail_canary_ready(self, domain: str, *, reason: str) -> None:
- del domain, reason
- return
- def _update_cfmail_canary_after_result(self, *, thread_id: int, result: dict[str, object]) -> None:
- del thread_id, result
- return
- def _wait_if_cfmail_canary_pending(self, thread_id: int, provider: str) -> bool:
- del thread_id, provider
- return False
- def _release_cfmail_flow_slot(self, thread_id: int) -> None:
- with self._lock:
- inflight_by_thread = self._cfmail_flow_state.setdefault("inflight_by_thread", {})
- if isinstance(inflight_by_thread, dict):
- inflight_by_thread.pop(thread_id, None)
- selected_profile_by_thread = self._cfmail_flow_state.setdefault("selected_profile_by_thread", {})
- if isinstance(selected_profile_by_thread, dict):
- selected_profile_by_thread.pop(thread_id, None)
- def _cfmail_active_domain_set(self) -> set[str]:
- domains = {
- str(item.get("domain") or "").strip().lower()
- for item in self._current_cfmail_active_accounts()
- if str(item.get("domain") or "").strip()
- }
- if domains:
- return domains
- current_domain = str(self._current_cfmail_active_domain() or "").strip().lower()
- return {current_domain} if current_domain else set()
- def _is_cfmail_active_domain(self, domain: str) -> bool:
- domain_key = str(domain or "").strip().lower()
- if not domain_key:
- return False
- return domain_key in self._cfmail_active_domain_set()
- def _schedule_cfmail_domain_pool_replenish(self, *, trigger_thread_id: int, reason: str) -> None:
- if self._cfmail_provisioner is None:
- return
- with self._lock:
- thread = self._cfmail_replenish_thread
- if thread is not None and thread.is_alive():
- return
- self._cfmail_replenish_reason = str(reason or "").strip()
- thread = threading.Thread(
- target=self._cfmail_domain_pool_replenish_worker,
- kwargs={
- "trigger_thread_id": trigger_thread_id,
- "reason": self._cfmail_replenish_reason,
- },
- daemon=True,
- name="zhuce6-cfmail-replenish",
- )
- self._cfmail_replenish_thread = thread
- thread.start()
- def _ensure_cfmail_domain_pool_target(self, *, trigger_thread_id: int, reason: str) -> None:
- if self._cfmail_provisioner is None:
- return
- if len(self._current_cfmail_active_accounts()) >= self._cfmail_active_domain_count:
- return
- self._schedule_cfmail_domain_pool_replenish(
- trigger_thread_id=trigger_thread_id,
- reason=reason,
- )
- def _cfmail_domain_pool_replenish_worker(self, *, trigger_thread_id: int, reason: str) -> None:
- provisioner = self._cfmail_provisioner
- if provisioner is None:
- return
- while not self._stop_event.is_set():
- active_accounts = self._current_cfmail_active_accounts()
- if len(active_accounts) >= self._cfmail_active_domain_count:
- return
- result = provisioner.provision_additional_domain()
- if not result.success:
- self._log(
- f"[zhuce6:register] [thread-{trigger_thread_id}] [cfmail] replenish failed after {reason}: "
- f"{result.error}"
- )
- if self._stop_event.wait(5.0):
- return
- continue
- self._reload_cfmail_manager_after_rotation()
- if self._cfmail_tracker is not None:
- try:
- self._cfmail_tracker.mark_rotation_completed("", result.new_domain)
- except Exception:
- pass
- self._log(
- f"[zhuce6:register] [thread-{trigger_thread_id}] [cfmail] replenished domain pool "
- f"after {reason}: +{result.new_domain}"
- )
- def _clear_cfmail_domain_state(self, domain: str) -> None:
- domain_key = str(domain or "").strip().lower()
- if not domain_key:
- return
- with self._lock:
- self._cfmail_add_phone_events.pop(domain_key, None)
- if str(self._cfmail_add_phone_state.get("active_domain") or "").strip().lower() == domain_key:
- self._cfmail_add_phone_state = {
- "active_domain": "",
- "in_cooldown": False,
- "cooldown_until": 0.0,
- "last_triggered_at": "",
- "last_rotation_attempted_at": "",
- "last_reason": "",
- "last_add_phone_failures": 0,
- "last_successes": 0,
- "last_window_size": 0,
- "last_logged_at": 0.0,
- }
- self._cfmail_wait_otp_events.pop(domain_key, None)
- if str(self._cfmail_wait_otp_state.get("active_domain") or "").strip().lower() == domain_key:
- self._cfmail_wait_otp_state = {
- "active_domain": "",
- "in_cooldown": False,
- "cooldown_until": 0.0,
- "last_triggered_at": "",
- "last_rotation_attempted_at": "",
- "last_reason": "",
- "last_no_message_timeouts": 0,
- "last_successes": 0,
- "last_message_seen": 0,
- "last_window_size": 0,
- "last_logged_at": 0.0,
- }
- self._cfmail_wait_otp_live_progress.pop(domain_key, None)
- if str(self._cfmail_fresh_domain_state.get("active_domain") or "").strip().lower() == domain_key:
- self._cfmail_fresh_domain_state = {
- "active_domain": "",
- "completed_attempts": 0,
- "mail_seen_attempts": 0,
- "successes": 0,
- "last_triggered_at": "",
- "last_rotation_attempted_at": "",
- "last_reason": "",
- }
- def _replace_cfmail_domain(
- self,
- *,
- thread_id: int,
- domain: str,
- reason_label: str,
- ) -> bool:
- domain_key = str(domain or "").strip().lower()
- if not domain_key or self._cfmail_provisioner is None:
- return False
- active_accounts = self._current_cfmail_active_accounts()
- if not active_accounts:
- active_accounts = [{"name": "", "domain": domain_key}]
- if not any(item["domain"] == domain_key for item in active_accounts):
- self._clear_cfmail_domain_state(domain_key)
- return False
- if not self._cfmail_rotation_lock.acquire(blocking=False):
- return False
- self._cfmail_rotation_pause.clear()
- try:
- if self._cfmail_tracker is not None:
- self._cfmail_tracker.mark_rotation_started(domain_key, reason_label)
- if len(active_accounts) <= 1:
- provision_result = self._cfmail_provisioner.rotate_active_domain()
- if not provision_result.success:
- if self._cfmail_tracker is not None:
- self._cfmail_tracker.mark_rotation_failed(domain_key, provision_result.error)
- self._log(
- f"[zhuce6:register] [thread-{thread_id}] [cfmail] {reason_label} rotation failed: "
- f"{provision_result.error}"
- )
- return False
- self._reload_cfmail_manager_after_rotation()
- self._clear_cfmail_domain_state(provision_result.old_domain or domain_key)
- self._reset_cfmail_add_phone_stoploss(provision_result.new_domain)
- self._reset_cfmail_wait_otp_stoploss(provision_result.new_domain)
- self._reset_cfmail_fresh_domain_budget(provision_result.new_domain)
- if self._cfmail_tracker is not None:
- self._cfmail_tracker.mark_rotation_completed(
- provision_result.old_domain,
- provision_result.new_domain,
- )
- self._log(
- f"[zhuce6:register] [thread-{thread_id}] [cfmail] {reason_label} rotation completed: "
- f"{provision_result.old_domain} -> {provision_result.new_domain}"
- )
- return True
- retire_result = self._cfmail_provisioner.retire_domain(domain_key)
- if not retire_result.success:
- if self._cfmail_tracker is not None:
- self._cfmail_tracker.mark_rotation_failed(domain_key, retire_result.error)
- self._log(
- f"[zhuce6:register] [thread-{thread_id}] [cfmail] {reason_label} retire failed: "
- f"{retire_result.error}"
- )
- return False
- self._reload_cfmail_manager_after_rotation()
- self._clear_cfmail_domain_state(domain_key)
- self._log(
- f"[zhuce6:register] [thread-{thread_id}] [cfmail] retired domain {domain_key} "
- f"because {reason_label}; scheduling replenish"
- )
- self._schedule_cfmail_domain_pool_replenish(
- trigger_thread_id=thread_id,
- reason=reason_label,
- )
- return True
- finally:
- self._cfmail_rotation_pause.set()
- self._cfmail_rotation_lock.release()
- def _wait_if_cfmail_flow_throttled(self, thread_id: int, provider: str) -> bool:
- if provider != "cfmail":
- return False
- self._ensure_cfmail_domain_pool_target(
- trigger_thread_id=thread_id,
- reason="usable domain pool below target",
- )
- active_accounts = self._current_cfmail_active_accounts()
- if not active_accounts:
- return False
- with self._lock:
- inflight_by_thread = self._cfmail_flow_state.setdefault("inflight_by_thread", {})
- if not isinstance(inflight_by_thread, dict):
- inflight_by_thread = {}
- self._cfmail_flow_state["inflight_by_thread"] = inflight_by_thread
- last_started_by_domain = self._cfmail_flow_state.setdefault("last_started_by_domain", {})
- if not isinstance(last_started_by_domain, dict):
- last_started_by_domain = {}
- self._cfmail_flow_state["last_started_by_domain"] = last_started_by_domain
- selected_profile_by_thread = self._cfmail_flow_state.setdefault("selected_profile_by_thread", {})
- if not isinstance(selected_profile_by_thread, dict):
- selected_profile_by_thread = {}
- self._cfmail_flow_state["selected_profile_by_thread"] = selected_profile_by_thread
- tracked = inflight_by_thread.get(thread_id)
- if isinstance(tracked, dict):
- tracked_domain = str(tracked.get("domain") or "").strip().lower()
- tracked_profile = str(tracked.get("profile_name") or "").strip()
- if tracked_domain and tracked_profile:
- selected_profile_by_thread[thread_id] = tracked_profile
- return False
- now = time.time()
- candidates: list[tuple[int, float, str, str]] = []
- min_wait_seconds = 2.0
- wait_domain = ""
- wait_reason = "inflight_limit"
- for account in active_accounts:
- domain = account["domain"]
- profile_name = account["name"]
- active_inflight = sum(
- 1
- for value in inflight_by_thread.values()
- if isinstance(value, dict)
- and str(value.get("domain") or "").strip().lower() == domain
- )
- if active_inflight >= self._cfmail_max_inflight:
- wait_domain = wait_domain or domain
- continue
- last_started_at = float(last_started_by_domain.get(domain) or 0.0)
- remaining = 0.0
- if self._cfmail_start_interval_seconds > 0 and last_started_at > 0.0:
- remaining = self._cfmail_start_interval_seconds - (now - last_started_at)
- if remaining > 0.0:
- wait_domain = wait_domain or domain
- wait_reason = "start_interval"
- min_wait_seconds = min(min_wait_seconds, min(2.0, max(0.5, remaining)))
- continue
- candidates.append((active_inflight, last_started_at, profile_name, domain))
- if candidates:
- candidates.sort(key=lambda item: (item[0], item[1], item[3], item[2]))
- _active_inflight, _last_started_at, profile_name, domain = candidates[0]
- inflight_by_thread[thread_id] = {
- "domain": domain,
- "profile_name": profile_name,
- "started_at": now,
- }
- last_started_by_domain[domain] = now
- selected_profile_by_thread[thread_id] = profile_name
- return False
- last_logged_at = float(self._cfmail_flow_state.get("last_logged_at") or 0.0)
- should_log = now - last_logged_at >= 15.0
- if should_log:
- self._cfmail_flow_state["last_logged_at"] = now
- if should_log:
- if wait_reason == "inflight_limit":
- self._log(
- f"[zhuce6:register] [thread-{thread_id}] [cfmail] flow throttle for {wait_domain or '-'}: "
- f"inflight_limit reached on all active domains"
- )
- else:
- self._log(
- f"[zhuce6:register] [thread-{thread_id}] [cfmail] flow throttle for {wait_domain or '-'}: "
- f"start_interval_remaining={min_wait_seconds:.1f}s"
- )
- self._stop_event.wait(min_wait_seconds)
- return True
- def _proxy_pool_snapshot(self) -> dict[str, object]:
- pool = self._proxy_pool
- nodes: list[dict[str, object]] = []
- snapshot_error: str | None = None
- if pool is not None:
- try:
- snapshot = pool.snapshot()
- except Exception as exc:
- snapshot_error = str(exc)
- else:
- if isinstance(snapshot, list):
- nodes = [item for item in snapshot if isinstance(item, dict)]
- return {
- "configured": bool(self.settings.proxy_pool_configured or pool is not None),
- "enabled": pool is not None,
- "snapshot_error": snapshot_error,
- "node_count": len(nodes),
- "in_use_count": sum(1 for item in nodes if item.get("in_use")),
- "disabled_count": sum(1 for item in nodes if item.get("disabled")),
- "nodes": nodes,
- }
- def _setup_logger(self) -> Any:
- import logging
- from logging.handlers import RotatingFileHandler
- logger = logging.getLogger("zhuce6.register")
- logger.setLevel(logging.INFO)
- logger.propagate = False
- if not logger.handlers:
- console = logging.StreamHandler()
- console.setFormatter(logging.Formatter("%(message)s"))
- logger.addHandler(console)
- if self.settings.register_log_file:
- fh = RotatingFileHandler(
- self.settings.register_log_file,
- maxBytes=2 * 1024 * 1024, # 2MB
- backupCount=5,
- encoding="utf-8",
- )
- fh.setFormatter(logging.Formatter("%(asctime)s %(message)s", datefmt="%Y-%m-%d %H:%M:%S"))
- logger.addHandler(fh)
- return logger
- def _log(self, msg: str) -> None:
- self._logger.info(msg)
- def _record_attempt(
- self,
- *,
- success: bool,
- stage: str,
- error_message: str,
- metadata: dict[str, object] | None = None,
- proxy_key: str = "",
- email: str = "",
- ) -> None:
- meta = metadata if isinstance(metadata, dict) else {}
- stage_key = str(stage or "?").strip() or "?"
- signal = self._classify_failure_signal(stage=stage_key, metadata=meta)
- timestamp = datetime.now().isoformat(timespec="seconds")
- event = {
- "timestamp": timestamp,
- "success": success,
- "stage": stage_key if not success else "completed",
- "signal": signal,
- "error_message": str(error_message or "").strip(),
- "email_domain": str(meta.get("email_domain") or "").strip(),
- "post_create_gate": str(meta.get("post_create_gate") or "").strip(),
- "create_account_error_code": str(meta.get("create_account_error_code") or "").strip(),
- "signup_error_code": str(meta.get("signup_error_code") or "").strip(),
- "signup_http_status": meta.get("signup_http_status"),
- "mailbox_error_kind": str(meta.get("mailbox_error_kind") or "").strip(),
- "mailbox_error_stage": str(meta.get("mailbox_error_stage") or "").strip(),
- "proxy_key": proxy_key,
- "email": email,
- }
- self._recent_attempts.append(event)
- if success or stage_key == "warmup_pending":
- return
- self._failure_by_stage[stage_key] = self._failure_by_stage.get(stage_key, 0) + 1
- if signal:
- self._failure_signals[signal] = self._failure_signals.get(signal, 0) + 1
- def _classify_failure_signal(self, *, stage: str, metadata: dict[str, object]) -> str:
- code = str(metadata.get("create_account_error_code") or "").strip().lower()
- signup_code = str(metadata.get("signup_error_code") or "").strip().lower()
- post_gate = str(metadata.get("post_create_gate") or "").strip().lower()
- if stage == "cpa_sync":
- return "cpa_sync_failed"
- if stage == "signup" and signup_code:
- return signup_code
- if stage == "add_phone_gate" or post_gate == "add_phone":
- return "add_phone_gate"
- if stage == "create_account" and code == "user_already_exists":
- return "mailbox_reused"
- if stage == "create_account" and code in {"registration_disallowed", "unsupported_email"}:
- return code
- if stage == "mailbox":
- provider = str(metadata.get("mail_provider") or "").strip().lower()
- if provider == "cfmail":
- mailbox_error_kind = str(metadata.get("mailbox_error_kind") or "").strip().lower()
- mailbox_error_stage = str(metadata.get("mailbox_error_stage") or "").strip().lower()
- if mailbox_error_stage == "create_email":
- mailbox_error_stage = "create"
- elif mailbox_error_stage == "fetch_email":
- mailbox_error_stage = "fetch"
- if mailbox_error_kind == "transport_error":
- suffix = mailbox_error_stage or "backend"
- return f"mailbox_{suffix}_transport_error"
- if mailbox_error_kind == "provider_error":
- suffix = mailbox_error_stage or "backend"
- return f"mailbox_{suffix}_provider_error"
- return "mailbox_backend_failure"
- return "mailbox_failure"
- return ""
- def _recent_failure_hotspots(self, limit: int = 5) -> list[dict[str, object]]:
- return self._recent_failure_hotspots_from_attempts(self._recent_attempts, limit=limit)
- def _recent_failure_hotspots_from_attempts(
- self,
- attempts: list[dict[str, object]] | deque[dict[str, object]],
- *,
- limit: int = 5,
- ) -> list[dict[str, object]]:
- counts: dict[tuple[str, str], int] = {}
- for item in attempts:
- if item.get("success"):
- continue
- stage = str(item.get("stage") or "?").strip() or "?"
- signal = str(item.get("signal") or "").strip()
- key = (signal or stage, stage)
- counts[key] = counts.get(key, 0) + 1
- ordered = sorted(counts.items(), key=lambda kv: (-kv[1], kv[0][0], kv[0][1]))
- return [
- {"key": key, "stage": stage, "count": count}
- for (key, stage), count in ordered[:limit]
- ]
- def _failure_counts_from_attempts(
- self,
- attempts: list[dict[str, object]] | deque[dict[str, object]],
- ) -> tuple[dict[str, int], dict[str, int]]:
- failure_by_stage: dict[str, int] = {}
- failure_signals: dict[str, int] = {}
- for item in attempts:
- if item.get("success"):
- continue
- stage = str(item.get("stage") or "?").strip() or "?"
- failure_by_stage[stage] = failure_by_stage.get(stage, 0) + 1
- signal = str(item.get("signal") or "").strip()
- if signal:
- failure_signals[signal] = failure_signals.get(signal, 0) + 1
- return (
- dict(sorted(failure_by_stage.items(), key=lambda item: (-item[1], item[0]))),
- dict(sorted(failure_signals.items(), key=lambda item: (-item[1], item[0]))),
- )
- def _active_domain_attempts(
- self,
- recent_attempts: list[dict[str, object]],
- active_domain: str,
- ) -> list[dict[str, object]]:
- domain = str(active_domain or "").strip().lower()
- if not domain:
- return list(recent_attempts)
- filtered = [
- item
- for item in recent_attempts
- if str(item.get("email_domain") or "").strip().lower() == domain
- ]
- return filtered
- def _infer_active_domain(
- self,
- recent_attempts: list[dict[str, object]],
- cfmail_rotation: dict[str, object] | None,
- stoploss: dict[str, object],
- ) -> str:
- if isinstance(cfmail_rotation, dict):
- domain = str(cfmail_rotation.get("active_domain") or "").strip().lower()
- if domain:
- return domain
- domain = str(stoploss.get("active_domain") or "").strip().lower()
- if domain:
- return domain
- for item in reversed(recent_attempts):
- domain = str(item.get("email_domain") or "").strip().lower()
- if domain:
- return domain
- return ""
- def _extract_email_domain(self, result: dict[str, object]) -> str:
- metadata = result.get("metadata") if isinstance(result.get("metadata"), dict) else {}
- domain = str(metadata.get("email_domain") or "").strip().lower()
- if domain:
- return domain
- email = str(result.get("email") or "").strip().lower()
- if "@" not in email:
- return ""
- return email.rsplit("@", 1)[-1].strip().lower()
- def _update_cfmail_add_phone_stoploss(self, result: dict[str, object]) -> None:
- metadata = result.get("metadata") if isinstance(result.get("metadata"), dict) else {}
- provider = str(metadata.get("mail_provider") or result.get("mail_provider") or "").strip().lower()
- if provider not in {"", "cfmail"}:
- return
- domain = self._extract_email_domain(result)
- if not domain:
- return
- if not self._is_cfmail_active_domain(domain):
- return
- success = bool(result.get("success"))
- stage = str(result.get("stage") or "").strip().lower()
- post_gate = str(metadata.get("post_create_gate") or "").strip().lower()
- is_add_phone = stage == "add_phone_gate" or post_gate == "add_phone"
- with self._lock:
- events = self._cfmail_add_phone_events.setdefault(
- domain,
- deque(maxlen=self._cfmail_add_phone_window),
- )
- events.append(
- {
- "success": success,
- "is_add_phone": is_add_phone,
- }
- )
- state = self._cfmail_add_phone_state
- state["active_domain"] = domain
- if self._cfmail_add_phone_cooldown_seconds <= 0:
- state["in_cooldown"] = False
- state["cooldown_until"] = 0.0
- return
- cooldown_until = float(state.get("cooldown_until") or 0.0)
- if time.time() < cooldown_until:
- state["in_cooldown"] = True
- return
- state["in_cooldown"] = False
- if len(events) < self._cfmail_add_phone_window:
- return
- add_phone_failures = sum(1 for item in events if item.get("is_add_phone"))
- successes = sum(1 for item in events if item.get("success"))
- if (
- add_phone_failures >= self._cfmail_add_phone_threshold
- and successes <= self._cfmail_add_phone_max_successes
- ):
- state["in_cooldown"] = True
- state["cooldown_until"] = time.time() + self._cfmail_add_phone_cooldown_seconds
- state["last_triggered_at"] = datetime.now().isoformat(timespec="seconds")
- state["last_rotation_attempted_at"] = ""
- state["last_reason"] = "add_phone threshold reached"
- state["last_add_phone_failures"] = add_phone_failures
- state["last_successes"] = successes
- state["last_window_size"] = len(events)
- self._log(
- f"[zhuce6:register] [cfmail] add_phone stoploss activated for {domain} "
- f"(add_phone_failures={add_phone_failures}, successes={successes}, window={len(events)})"
- )
- def _is_cfmail_wait_otp_no_message_timeout(self, result: dict[str, object]) -> bool:
- metadata = result.get("metadata") if isinstance(result.get("metadata"), dict) else {}
- provider = str(metadata.get("mail_provider") or result.get("mail_provider") or "").strip().lower()
- if provider not in {"", "cfmail"}:
- return False
- stage = str(result.get("stage") or "").strip().lower()
- if stage != "wait_otp":
- return False
- failure_reason = str(metadata.get("otp_wait_failure_reason") or "").strip().lower()
- if failure_reason:
- return failure_reason == "mailbox_timeout_no_message"
- try:
- return int(metadata.get("otp_mailbox_message_scan_count") or 0) <= 0
- except Exception:
- return False
- def _is_cfmail_invalid_domain_mailbox_failure(self, result: dict[str, object]) -> bool:
- metadata = result.get("metadata") if isinstance(result.get("metadata"), dict) else {}
- if str(result.get("stage") or "").strip().lower() != "mailbox":
- return False
- provider = str(metadata.get("mail_provider") or result.get("mail_provider") or "").strip().lower()
- if provider not in {"", "cfmail"}:
- return False
- haystacks = [str(result.get("error_message") or "")]
- haystacks.extend(str(item or "") for item in (result.get("logs") or []))
- text = "\n".join(haystacks).lower()
- return "invalid domain" in text or "无效的域名" in text
- def _on_cfmail_wait_progress(self, account: object, diagnostics: dict[str, object]) -> None:
- if self._cfmail_wait_otp_live_threshold <= 0:
- return
- email = str(getattr(account, "email", "") or "").strip().lower()
- extra = getattr(account, "extra", {}) or {}
- domain = str(extra.get("email_domain") or "").strip().lower()
- if not domain and "@" in email:
- domain = email.rsplit("@", 1)[-1].strip().lower()
- if not domain:
- return
- with self._cfmail_wait_otp_live_lock:
- now = time.time()
- for tracked_domain in list(self._cfmail_wait_otp_live_progress.keys()):
- entries = self._cfmail_wait_otp_live_progress.get(tracked_domain) or {}
- fresh_entries = {
- key: value
- for key, value in entries.items()
- if now - float(value.get("updated_at") or 0.0) <= 15.0
- }
- if fresh_entries:
- self._cfmail_wait_otp_live_progress[tracked_domain] = fresh_entries
- else:
- self._cfmail_wait_otp_live_progress.pop(tracked_domain, None)
- if not self._is_cfmail_active_domain(domain):
- self._cfmail_wait_otp_live_progress.pop(domain, None)
- return
- key = str(getattr(account, "account_id", "") or email or id(account))
- domain_entries = self._cfmail_wait_otp_live_progress.setdefault(domain, {})
- domain_entries[key] = {
- "scan_count": int(diagnostics.get("message_scan_count") or 0),
- "elapsed_seconds": float(diagnostics.get("elapsed_seconds") or 0.0),
- "updated_at": now,
- }
- if int(diagnostics.get("message_scan_count") or 0) > 0:
- self._mark_cfmail_canary_ready(domain, reason="live_mailbox_message_seen")
- state = self._cfmail_wait_otp_state
- if bool(state.get("in_cooldown")) and str(state.get("active_domain") or "").strip().lower() == domain:
- return
- stalled = [
- value
- for value in domain_entries.values()
- if int(value.get("scan_count") or 0) <= 0
- and float(value.get("elapsed_seconds") or 0.0) >= self._cfmail_wait_otp_live_age_seconds
- ]
- if len(stalled) < self._cfmail_wait_otp_live_threshold:
- return
- with self._lock:
- state = self._cfmail_wait_otp_state
- if bool(state.get("in_cooldown")) and str(state.get("active_domain") or "").strip().lower() == domain:
- return
- self._cfmail_wait_otp_state = {
- "active_domain": domain,
- "in_cooldown": True,
- "cooldown_until": time.time() + self._cfmail_wait_otp_cooldown_seconds,
- "last_triggered_at": now_iso(),
- "last_rotation_attempted_at": "",
- "last_reason": "live wait_otp no-message threshold reached",
- "last_no_message_timeouts": len(stalled),
- "last_window_size": len(stalled),
- "last_logged_at": 0.0,
- }
- self._log(
- f"[zhuce6:register] [cfmail] wait_otp live stoploss activated for {domain} "
- f"(stalled_waits={len(stalled)}, age>={self._cfmail_wait_otp_live_age_seconds}s)"
- )
- def _update_cfmail_wait_otp_stoploss(self, result: dict[str, object]) -> None:
- metadata = result.get("metadata") if isinstance(result.get("metadata"), dict) else {}
- provider = str(metadata.get("mail_provider") or result.get("mail_provider") or "").strip().lower()
- if provider not in {"", "cfmail"}:
- return
- domain = self._extract_email_domain(result)
- if not domain:
- return
- if not self._is_cfmail_active_domain(domain):
- return
- is_no_message_timeout = self._is_cfmail_wait_otp_no_message_timeout(result)
- try:
- message_scan_count = int(metadata.get("otp_mailbox_message_scan_count") or 0)
- except Exception:
- message_scan_count = 0
- has_message_seen = message_scan_count > 0
- with self._lock:
- events = self._cfmail_wait_otp_events.setdefault(
- domain,
- deque(maxlen=self._cfmail_wait_otp_window),
- )
- events.append(
- {
- "success": bool(result.get("success")),
- "is_no_message_timeout": is_no_message_timeout,
- "has_message_seen": has_message_seen,
- }
- )
- state = self._cfmail_wait_otp_state
- state["active_domain"] = domain
- if self._cfmail_wait_otp_cooldown_seconds <= 0:
- state["in_cooldown"] = False
- state["cooldown_until"] = 0.0
- return
- cooldown_until = float(state.get("cooldown_until") or 0.0)
- if time.time() < cooldown_until:
- state["in_cooldown"] = True
- return
- state["in_cooldown"] = False
- if len(events) < self._cfmail_wait_otp_window:
- return
- no_message_timeouts = sum(1 for item in events if item.get("is_no_message_timeout"))
- successes = sum(1 for item in events if item.get("success"))
- message_seen = sum(1 for item in events if item.get("has_message_seen"))
- if (
- no_message_timeouts >= self._cfmail_wait_otp_threshold
- and successes <= 0
- and message_seen <= 0
- ):
- state["in_cooldown"] = True
- state["cooldown_until"] = time.time() + self._cfmail_wait_otp_cooldown_seconds
- state["last_triggered_at"] = datetime.now().isoformat(timespec="seconds")
- state["last_rotation_attempted_at"] = ""
- state["last_reason"] = "wait_otp no-message threshold reached"
- state["last_no_message_timeouts"] = no_message_timeouts
- state["last_successes"] = successes
- state["last_message_seen"] = message_seen
- state["last_window_size"] = len(events)
- self._log(
- f"[zhuce6:register] [cfmail] wait_otp stoploss activated for {domain} "
- f"(no_message_timeouts={no_message_timeouts}, successes={successes}, "
- f"message_seen={message_seen}, window={len(events)})"
- )
- def _update_cfmail_fresh_domain_budget(self, result: dict[str, object]) -> None:
- if self._cfmail_fresh_domain_attempt_budget <= 0:
- return
- metadata = result.get("metadata") if isinstance(result.get("metadata"), dict) else {}
- provider = str(metadata.get("mail_provider") or result.get("mail_provider") or "").strip().lower()
- if provider not in {"", "cfmail"}:
- return
- domain = self._extract_email_domain(result)
- if not domain:
- return
- if not self._is_cfmail_active_domain(domain):
- return
- try:
- message_scan_count = int(metadata.get("otp_mailbox_message_scan_count") or 0)
- except Exception:
- message_scan_count = 0
- success = bool(result.get("success"))
- with self._lock:
- state = self._cfmail_fresh_domain_state
- tracked_domain = str(state.get("active_domain") or "").strip().lower()
- if tracked_domain != domain:
- self._cfmail_fresh_domain_state = {
- "active_domain": domain,
- "completed_attempts": 0,
- "mail_seen_attempts": 0,
- "successes": 0,
- "last_triggered_at": "",
- "last_rotation_attempted_at": "",
- "last_reason": "",
- }
- state = self._cfmail_fresh_domain_state
- state["completed_attempts"] = int(state.get("completed_attempts") or 0) + 1
- if message_scan_count > 0:
- state["mail_seen_attempts"] = int(state.get("mail_seen_attempts") or 0) + 1
- if success:
- state["successes"] = int(state.get("successes") or 0) + 1
- if (
- int(state.get("mail_seen_attempts") or 0) > 0
- and int(state.get("completed_attempts") or 0) >= self._cfmail_fresh_domain_attempt_budget
- and not str(state.get("last_rotation_attempted_at") or "").strip()
- ):
- state["last_triggered_at"] = datetime.now().isoformat(timespec="seconds")
- state["last_reason"] = "fresh_domain_attempt_budget_reached"
- self._log(
- f"[zhuce6:register] [cfmail] fresh-domain budget reached for {domain} "
- f"(completed_attempts={int(state.get('completed_attempts') or 0)}, "
- f"mail_seen_attempts={int(state.get('mail_seen_attempts') or 0)}, "
- f"budget={self._cfmail_fresh_domain_attempt_budget})"
- )
- def _cfmail_add_phone_stoploss_snapshot(self) -> dict[str, object]:
- with self._lock:
- state = dict(self._cfmail_add_phone_state)
- if self._cfmail_add_phone_cooldown_seconds <= 0:
- return {
- "active_domain": str(state.get("active_domain") or ""),
- "in_cooldown": False,
- "cooldown_remaining_seconds": 0,
- "last_triggered_at": str(state.get("last_triggered_at") or ""),
- "last_reason": str(state.get("last_reason") or ""),
- "last_add_phone_failures": int(state.get("last_add_phone_failures") or 0),
- "last_successes": int(state.get("last_successes") or 0),
- "last_window_size": int(state.get("last_window_size") or 0),
- "window_size": self._cfmail_add_phone_window,
- "threshold": self._cfmail_add_phone_threshold,
- "max_successes_in_window": self._cfmail_add_phone_max_successes,
- }
- cooldown_until = float(state.get("cooldown_until") or 0.0)
- remaining = max(0, int(cooldown_until - time.time()))
- if remaining <= 0:
- state["in_cooldown"] = False
- return {
- "active_domain": str(state.get("active_domain") or ""),
- "in_cooldown": bool(state.get("in_cooldown")),
- "cooldown_remaining_seconds": remaining,
- "last_triggered_at": str(state.get("last_triggered_at") or ""),
- "last_reason": str(state.get("last_reason") or ""),
- "last_add_phone_failures": int(state.get("last_add_phone_failures") or 0),
- "last_successes": int(state.get("last_successes") or 0),
- "last_window_size": int(state.get("last_window_size") or 0),
- "window_size": self._cfmail_add_phone_window,
- "threshold": self._cfmail_add_phone_threshold,
- "max_successes_in_window": self._cfmail_add_phone_max_successes,
- }
- def _reset_cfmail_add_phone_stoploss(self, new_domain: str = "") -> None:
- with self._lock:
- self._cfmail_add_phone_state = {
- "active_domain": new_domain,
- "in_cooldown": False,
- "cooldown_until": 0.0,
- "last_triggered_at": "",
- "last_rotation_attempted_at": "",
- "last_reason": "",
- "last_add_phone_failures": 0,
- "last_successes": 0,
- "last_window_size": 0,
- "last_logged_at": 0.0,
- }
- def _cfmail_wait_otp_stoploss_snapshot(self) -> dict[str, object]:
- with self._lock:
- state = dict(self._cfmail_wait_otp_state)
- if self._cfmail_wait_otp_cooldown_seconds <= 0:
- return {
- "active_domain": str(state.get("active_domain") or ""),
- "in_cooldown": False,
- "cooldown_remaining_seconds": 0,
- "last_triggered_at": str(state.get("last_triggered_at") or ""),
- "last_reason": str(state.get("last_reason") or ""),
- "last_no_message_timeouts": int(state.get("last_no_message_timeouts") or 0),
- "last_successes": int(state.get("last_successes") or 0),
- "last_message_seen": int(state.get("last_message_seen") or 0),
- "last_window_size": int(state.get("last_window_size") or 0),
- "window_size": self._cfmail_wait_otp_window,
- "threshold": self._cfmail_wait_otp_threshold,
- }
- cooldown_until = float(state.get("cooldown_until") or 0.0)
- remaining = max(0, int(cooldown_until - time.time()))
- if remaining <= 0:
- state["in_cooldown"] = False
- return {
- "active_domain": str(state.get("active_domain") or ""),
- "in_cooldown": bool(state.get("in_cooldown")),
- "cooldown_remaining_seconds": remaining,
- "last_triggered_at": str(state.get("last_triggered_at") or ""),
- "last_reason": str(state.get("last_reason") or ""),
- "last_no_message_timeouts": int(state.get("last_no_message_timeouts") or 0),
- "last_successes": int(state.get("last_successes") or 0),
- "last_message_seen": int(state.get("last_message_seen") or 0),
- "last_window_size": int(state.get("last_window_size") or 0),
- "window_size": self._cfmail_wait_otp_window,
- "threshold": self._cfmail_wait_otp_threshold,
- }
- def _cfmail_fresh_domain_budget_snapshot(self) -> dict[str, object]:
- with self._lock:
- state = dict(self._cfmail_fresh_domain_state)
- return {
- "active_domain": str(state.get("active_domain") or ""),
- "completed_attempts": int(state.get("completed_attempts") or 0),
- "mail_seen_attempts": int(state.get("mail_seen_attempts") or 0),
- "successes": int(state.get("successes") or 0),
- "last_triggered_at": str(state.get("last_triggered_at") or ""),
- "last_rotation_attempted_at": str(state.get("last_rotation_attempted_at") or ""),
- "last_reason": str(state.get("last_reason") or ""),
- "attempt_budget": self._cfmail_fresh_domain_attempt_budget,
- }
- def _reset_cfmail_wait_otp_stoploss(self, new_domain: str = "") -> None:
- with self._lock:
- self._cfmail_wait_otp_state = {
- "active_domain": new_domain,
- "in_cooldown": False,
- "cooldown_until": 0.0,
- "last_triggered_at": "",
- "last_rotation_attempted_at": "",
- "last_reason": "",
- "last_no_message_timeouts": 0,
- "last_successes": 0,
- "last_message_seen": 0,
- "last_window_size": 0,
- "last_logged_at": 0.0,
- }
- with self._cfmail_wait_otp_live_lock:
- if new_domain:
- self._cfmail_wait_otp_live_progress = {
- str(new_domain).strip().lower(): {}
- }
- else:
- self._cfmail_wait_otp_live_progress = {}
- def _reset_cfmail_fresh_domain_budget(self, new_domain: str = "") -> None:
- with self._lock:
- self._cfmail_fresh_domain_state = {
- "active_domain": str(new_domain or "").strip().lower(),
- "completed_attempts": 0,
- "mail_seen_attempts": 0,
- "successes": 0,
- "last_triggered_at": "",
- "last_rotation_attempted_at": "",
- "last_reason": "",
- }
- def _reload_cfmail_manager_after_rotation(self) -> None:
- manager = self._cfmail_manager
- if manager is None:
- return
- try:
- reload_if_needed = getattr(manager, "reload_if_needed", None)
- if callable(reload_if_needed):
- try:
- reload_if_needed(force=True)
- except TypeError:
- reload_if_needed()
- except Exception as exc:
- self._log(f"[zhuce6:register] [cfmail] manager reload after rotation failed: {exc}")
- def _ensure_cfmail_active_domain_ready(self) -> bool:
- if self._cfmail_provisioner is None or self._cfmail_manager is None:
- return False
- active_accounts = self._current_cfmail_active_accounts()
- active_domains = [item["domain"] for item in active_accounts]
- if not active_domains:
- return False
- self._log(
- f"[zhuce6:register] [cfmail] active-domain pool ready: "
- f"{', '.join(active_domains)}"
- )
- if self._cfmail_tracker is not None:
- for domain in active_domains:
- try:
- self._cfmail_tracker.active_domain = domain
- except Exception:
- pass
- return True
- def _force_rotate_cfmail_for_invalid_mailbox(self, thread_id: int, result: dict[str, object]) -> bool:
- if not self._is_cfmail_invalid_domain_mailbox_failure(result):
- return False
- if self._cfmail_tracker is None or self._cfmail_provisioner is None:
- return False
- current_domain = self._extract_email_domain(result) or self._current_cfmail_active_domain()
- if not current_domain:
- return False
- self._log(
- f"[zhuce6:register] [thread-{thread_id}] [cfmail] rotating invalid domain "
- f"{current_domain} after mailbox bootstrap failure"
- )
- return self._replace_cfmail_domain(
- thread_id=thread_id,
- domain=current_domain,
- reason_label="mailbox invalid domain",
- )
- def _rotate_cfmail_for_failed_canary(self, thread_id: int, result: dict[str, object]) -> bool:
- del thread_id, result
- return False
- def _rotate_cfmail_for_fresh_domain_budget(self, thread_id: int) -> bool:
- if self._cfmail_fresh_domain_attempt_budget <= 0:
- return False
- if self._cfmail_tracker is None or self._cfmail_provisioner is None:
- return False
- with self._lock:
- state = dict(self._cfmail_fresh_domain_state)
- domain = str(state.get("active_domain") or "").strip().lower()
- completed_attempts = int(state.get("completed_attempts") or 0)
- mail_seen_attempts = int(state.get("mail_seen_attempts") or 0)
- if (
- not domain
- or mail_seen_attempts <= 0
- or completed_attempts < self._cfmail_fresh_domain_attempt_budget
- or str(state.get("last_rotation_attempted_at") or "").strip()
- ):
- return False
- self._cfmail_fresh_domain_state["last_rotation_attempted_at"] = datetime.now().isoformat(timespec="seconds")
- if not self._is_cfmail_active_domain(domain):
- self._clear_cfmail_domain_state(domain)
- return False
- self._log(
- f"[zhuce6:register] [thread-{thread_id}] [cfmail] rotating domain {domain} "
- f"because fresh domain budget reached"
- )
- if self._replace_cfmail_domain(
- thread_id=thread_id,
- domain=domain,
- reason_label="fresh domain budget reached",
- ):
- return True
- with self._lock:
- self._cfmail_fresh_domain_state["last_rotation_attempted_at"] = ""
- return False
- def _rotate_cfmail_for_stoploss(
- self,
- *,
- thread_id: int,
- state_attr: str,
- reason_label: str,
- ) -> bool:
- if self._cfmail_tracker is None or self._cfmail_provisioner is None:
- return False
- with self._lock:
- state_obj = getattr(self, state_attr, None)
- if not isinstance(state_obj, dict):
- return False
- domain = str(state_obj.get("active_domain") or "").strip().lower()
- if not state_obj.get("in_cooldown") or not domain:
- return False
- if str(state_obj.get("last_rotation_attempted_at") or "").strip():
- return False
- state_obj["last_rotation_attempted_at"] = datetime.now().isoformat(timespec="seconds")
- return self._replace_cfmail_domain(
- thread_id=thread_id,
- domain=domain,
- reason_label=reason_label,
- )
- def _wait_if_cfmail_add_phone_stopped(self, thread_id: int, provider: str) -> bool:
- if provider != "cfmail":
- return False
- if self._cfmail_add_phone_cooldown_seconds <= 0:
- return False
- state = self._cfmail_add_phone_stoploss_snapshot()
- if not state.get("in_cooldown"):
- return False
- had_spare_domains = len(self._cfmail_active_domain_set()) > 1
- domain = str(state.get("active_domain") or "").strip().lower()
- if not self._is_cfmail_active_domain(domain):
- self._clear_cfmail_domain_state(domain)
- return False
- if self._rotate_cfmail_for_stoploss(
- thread_id=thread_id,
- state_attr="_cfmail_add_phone_state",
- reason_label="add_phone stoploss",
- ):
- return not had_spare_domains
- remaining = int(state.get("cooldown_remaining_seconds") or 0)
- should_log = False
- with self._lock:
- last_logged_at = float(self._cfmail_add_phone_state.get("last_logged_at") or 0.0)
- now = time.time()
- if now - last_logged_at >= 15:
- self._cfmail_add_phone_state["last_logged_at"] = now
- should_log = True
- if should_log:
- self._log(
- f"[zhuce6:register] [thread-{thread_id}] [cfmail] add_phone stoploss active for {domain}, "
- f"remaining={remaining}s"
- )
- wait_seconds = min(max(remaining, 1), 5)
- self._stop_event.wait(wait_seconds)
- return True
- def _wait_if_cfmail_wait_otp_stopped(self, thread_id: int, provider: str) -> bool:
- if provider != "cfmail":
- return False
- if self._cfmail_wait_otp_cooldown_seconds <= 0:
- return False
- state = self._cfmail_wait_otp_stoploss_snapshot()
- if not state.get("in_cooldown"):
- return False
- had_spare_domains = len(self._cfmail_active_domain_set()) > 1
- domain = str(state.get("active_domain") or "").strip().lower()
- if not self._is_cfmail_active_domain(domain):
- self._clear_cfmail_domain_state(domain)
- return False
- if self._rotate_cfmail_for_stoploss(
- thread_id=thread_id,
- state_attr="_cfmail_wait_otp_state",
- reason_label="wait_otp stoploss",
- ):
- return not had_spare_domains
- remaining = int(state.get("cooldown_remaining_seconds") or 0)
- should_log = False
- with self._lock:
- last_logged_at = float(self._cfmail_wait_otp_state.get("last_logged_at") or 0.0)
- now = time.time()
- if now - last_logged_at >= 15:
- self._cfmail_wait_otp_state["last_logged_at"] = now
- should_log = True
- if should_log:
- self._log(
- f"[zhuce6:register] [thread-{thread_id}] [cfmail] wait_otp stoploss active for {domain}, "
- f"remaining={remaining}s"
- )
- wait_seconds = min(max(remaining, 1), 5)
- self._stop_event.wait(wait_seconds)
- return True
- def start(self) -> None:
- if self._threads:
- return
- _compat_main_attr("load_all", load_all)()
- self._stop_event.clear()
- self._target_reached.clear()
- self._started_at = time.time()
- num = self.settings.register_threads
- self._providers = [p.strip() for p in self.settings.register_mail_provider.split(",") if p.strip()]
- if not self._providers:
- self._providers = ["cfmail"]
- if "cfmail" in self._providers:
- from core.cfmail_domain_rotation import DomainHealthTracker
- import core.cfmail as cfmail_module
- from core.cfmail import DEFAULT_CFMAIL_MANAGER
- from core.cfmail_provisioner import CfmailProvisioner
- self._cfmail_tracker = DomainHealthTracker()
- self._cfmail_provisioner = CfmailProvisioner(proxy_url=self.settings.register_proxy)
- self._cfmail_manager = DEFAULT_CFMAIL_MANAGER
- cfmail_module.CFMAIL_WAIT_ABORT_PREDICATE = self._should_abort_cfmail_wait
- cfmail_module.CFMAIL_WAIT_PROGRESS_CALLBACK = self._on_cfmail_wait_progress
- try:
- normalize_result = self._cfmail_provisioner.normalize_to_domain_pool(self._cfmail_active_domain_count)
- self._cfmail_manager.reload_if_needed(force=True)
- provisioned_domains = list(normalize_result.get("provisioned_domains") or [])
- retired_domains = list(normalize_result.get("retired_domains") or [])
- active_domains = list(normalize_result.get("active_domains") or [])
- if provisioned_domains or retired_domains:
- self._log(
- "[zhuce6:register] [cfmail] normalized active domain pool: "
- f"active={','.join(active_domains) or '-'} "
- f"provisioned={','.join(provisioned_domains) or '-'} "
- f"retired={','.join(retired_domains) or '-'}"
- )
- except Exception as exc:
- self._log(f"[zhuce6:register] [cfmail] normalize active domain pool failed: {exc}")
- try:
- self._ensure_cfmail_active_domain_ready()
- self._ensure_cfmail_domain_pool_target(
- trigger_thread_id=0,
- reason="startup",
- )
- except Exception as exc:
- self._log(f"[zhuce6:register] [cfmail] startup active-domain preflight failed: {exc}")
- if self.settings.backend == "cpa" and self.settings.cpa_runtime_reconcile_enabled:
- try:
- _maybe_reconcile_cpa_runtime(
- pool_dir=self.settings.pool_dir,
- management_base_url=self.settings.cpa_management_base_url,
- enabled=True,
- cooldown_seconds=self.settings.cpa_runtime_reconcile_cooldown_seconds,
- restart_enabled=self.settings.cpa_runtime_reconcile_restart_enabled,
- state_file=self.settings.pool_dir / "cpa_runtime_reconcile_state.json",
- client=create_backend_client(self.settings),
- management_key=self.settings.cpa_management_key,
- )
- except Exception as exc:
- self._log(f"[zhuce6:register] startup reconcile failed: {exc}")
- if self.settings.proxy_pool_configured:
- from core.proxy_pool import ProxyPool
- self._proxy_pool = ProxyPool.from_settings(self.settings)
- if self._proxy_pool is not None:
- self._proxy_pool.start()
- target_msg = f", target={self.settings.register_target_count}" if self.settings.register_target_count > 0 else ""
- self._log(
- f"[zhuce6:register] starting {num} threads, "
- f"providers={','.join(self._providers)}, "
- f"proxy={self.settings.register_proxy or 'none'}, "
- f"sleep={self.settings.register_sleep_min}-{self.settings.register_sleep_max}s"
- f"{target_msg}"
- )
- threading_module = _compat_main_attr("threading", threading)
- for i in range(num):
- provider = self._providers[i % len(self._providers)]
- t = threading_module.Thread(
- target=self._worker,
- args=(i + 1, provider),
- daemon=True,
- name=f"zhuce6-register-{i + 1}",
- )
- t.start()
- self._threads.append(t)
- # Solution B: start deferred retry worker thread
- pending_t = threading_module.Thread(
- target=self._pending_token_retry_worker,
- daemon=True,
- name="zhuce6-pending-token-retry",
- )
- pending_t.start()
- self._threads.append(pending_t)
- self._write_runtime_state()
- def stop(self) -> None:
- self._stop_event.set()
- for t in self._threads:
- t.join(timeout=2)
- self._threads.clear()
- try:
- import core.cfmail as cfmail_module
- cfmail_module.CFMAIL_WAIT_ABORT_PREDICATE = None
- cfmail_module.CFMAIL_WAIT_PROGRESS_CALLBACK = None
- except Exception:
- pass
- if self._proxy_pool is not None:
- self._proxy_pool.close()
- self._proxy_pool = None
- self._write_runtime_state()
- def _should_stop(self) -> bool:
- return self._stop_event.is_set() or self._target_reached.is_set()
- def _check_target(self) -> bool:
- """Return True if target reached and threads should stop."""
- if self.settings.register_target_count <= 0:
- return False
- with self._lock:
- if self._total_success >= self.settings.register_target_count:
- self._target_reached.set()
- return True
- return False
- def _cpa_api_root(self) -> str:
- parsed = urlsplit(self.settings.cpa_management_base_url)
- path = parsed.path or ""
- suffix = "/v0/management"
- if path.endswith(suffix):
- path = path[: -len(suffix)]
- return urlunsplit((parsed.scheme, parsed.netloc, path, "", "")).rstrip("/")
- def _get_cpa_management_key(self) -> str | None:
- cached = self._cpa_management_key_cache
- if cached is not False:
- return str(cached or "") or None
- key = str(get_management_key() or "").strip() or None
- self._cpa_management_key_cache = key or None
- return key
- def _sync_cpa_from_success(self, result: dict[str, object], thread_id: int) -> tuple[bool, str, str]:
- """Persist a registered account to CPA immediately while keeping pool as backup."""
- pool_file_raw = str(result.get("pool_file") or "").strip()
- if not pool_file_raw:
- return False, "missing pool file", ""
- pool_file = Path(pool_file_raw)
- if not pool_file.is_file():
- return False, f"pool file missing: {pool_file.name}", ""
- sync_started_at = now_iso()
- key = self._get_cpa_management_key()
- if not key:
- update_token_record(
- pool_file,
- backup_written=True,
- cpa_sync_status="failed",
- last_cpa_sync_at=sync_started_at,
- last_cpa_sync_error="CPA management key unavailable",
- )
- return False, "CPA management key unavailable", ""
- try:
- from platforms.chatgpt.pool import load_token_record
- token_data = load_token_record(pool_file)
- except Exception as exc:
- return False, f"invalid pool file {pool_file.name}: {exc}", ""
- metadata = result.get("metadata") if isinstance(result.get("metadata"), dict) else {}
- post_create_gate = str((metadata or {}).get("post_create_gate") or token_data.get("registration_post_create_gate") or "").strip().lower()
- if post_create_gate == "add_phone" and not bool(token_data.get("warmup_required")):
- token_data = update_token_record(
- pool_file,
- warmup_required=True,
- warmup_state="pending",
- warmup_passed=False,
- registration_post_create_gate="add_phone",
- )
- if not isinstance(token_data, dict) or not str(token_data.get("email") or "").strip():
- update_token_record(
- pool_file,
- backup_written=True,
- cpa_sync_status="failed",
- last_cpa_sync_at=sync_started_at,
- last_cpa_sync_error="missing email in pool record",
- )
- return False, f"missing email in {pool_file.name}", ""
- if is_warmup_pending_record(token_data):
- update_token_record(
- pool_file,
- backup_written=True,
- cpa_sync_status="warmup_pending",
- last_cpa_sync_at=sync_started_at,
- last_cpa_sync_error="warmup pending",
- )
- return (
- False,
- "warmup pending",
- str(token_data.get("email") or pool_file.name).strip() or pool_file.name,
- )
- from platforms.chatgpt.cpa_upload import upload_to_cpa
- ok, message = upload_to_cpa(
- token_data,
- api_url=self._cpa_api_root(),
- api_key=key,
- proxy=None,
- )
- email = str(token_data.get("email") or pool_file.name).strip() or pool_file.name
- if ok:
- update_token_record(
- pool_file,
- health_status="good",
- backup_written=True,
- cpa_sync_status="synced",
- last_cpa_sync_at=sync_started_at,
- last_cpa_sync_error="",
- )
- return True, "", email
- else:
- update_token_record(
- pool_file,
- backup_written=True,
- cpa_sync_status="failed",
- last_cpa_sync_at=sync_started_at,
- last_cpa_sync_error=message,
- )
- return False, message, email
- # ── Solution B: Deferred retry queue ──────────────────────────
- def _enqueue_pending_token(
- self,
- result: dict[str, object],
- thread_id: int,
- *,
- proxy_key: str = "",
- proxy_url: str = "",
- ) -> None:
- """Save an add_phone_gate account for deferred token acquisition retry."""
- metadata = result.get("metadata") if isinstance(result.get("metadata"), dict) else {}
- deferred = metadata.get("deferred_credentials")
- if not isinstance(deferred, dict):
- return
- email = str(deferred.get("email") or "").strip()
- password = str(deferred.get("password") or "").strip()
- if not email or not password:
- return
- entry = {
- "email": email,
- "password": password,
- "mailbox_jwt": str(deferred.get("mailbox_jwt") or "").strip(),
- "mailbox_extra": dict(deferred.get("mailbox_extra") or {}),
- "registration_proxy_key": str(deferred.get("registration_proxy_key") or proxy_key or "").strip(),
- "registration_proxy_region": str(
- deferred.get("registration_proxy_region")
- or infer_proxy_region(
- str(deferred.get("registration_proxy_key") or proxy_key or "").strip()
- or str(deferred.get("registration_proxy_url") or proxy_url or "").strip()
- )
- or ""
- ).strip(),
- "registration_proxy_url": str(deferred.get("registration_proxy_url") or proxy_url or "").strip(),
- "registration_fingerprint_profile": str(
- deferred.get("registration_fingerprint_profile") or "chrome120_win"
- ).strip(),
- "cfmail_profile_name": str(
- deferred.get("cfmail_profile_name") or metadata.get("cfmail_profile_name") or ""
- ).strip(),
- "add_phone_trace_path": str(
- deferred.get("add_phone_trace_path") or metadata.get("add_phone_trace_path") or ""
- ).strip(),
- "created_at": time.time(),
- "retry_count": 0,
- "last_retry_at": 0.0,
- }
- with self._pending_token_lock:
- self._pending_token_queue.append(entry)
- self._pending_token_total_enqueued += 1
- origin_proxy = str(entry.get("registration_proxy_key") or entry.get("registration_proxy_region") or entry.get("registration_proxy_url") or "").strip()
- self._log(
- f"[zhuce6:register] [thread-{thread_id}] 📥 deferred token retry enqueued: {email} "
- f"(origin_proxy={origin_proxy or '-'}, queue_size={len(self._pending_token_queue)})"
- )
- def _pending_token_retry_worker(self) -> None:
- """Background thread that retries token acquisition for queued add_phone_gate accounts."""
- while not self._should_stop():
- if self._stop_event.wait(30):
- break
- batch: list[dict[str, Any]] = []
- now = time.time()
- with self._pending_token_lock:
- remaining: list[dict[str, Any]] = []
- for entry in self._pending_token_queue:
- created_at = float(entry.get("created_at") or 0)
- retry_count = int(entry.get("retry_count") or 0)
- last_retry = float(entry.get("last_retry_at") or 0)
- age = now - created_at
- since_last = now - last_retry if last_retry > 0 else age
- delay = self._pending_token_retry_delay_seconds * (retry_count + 1)
- if retry_count >= self._pending_token_max_retries:
- self._pending_token_total_failed += 1
- self._log(
- f"[zhuce6:register] [pending] ❌ exhausted retries for "
- f"{entry.get('email')}, discarding"
- )
- continue
- if since_last >= delay:
- batch.append(entry)
- else:
- remaining.append(entry)
- self._pending_token_queue = remaining
- if not batch:
- continue
- for entry in batch:
- if self._should_stop():
- break
- self._retry_pending_token(entry)
- def _retry_pending_token(self, entry: dict[str, Any]) -> None:
- """Attempt token acquisition for a single deferred account."""
- email = str(entry.get("email") or "").strip()
- password = str(entry.get("password") or "").strip()
- mailbox_jwt = str(entry.get("mailbox_jwt") or "").strip()
- mailbox_extra = dict(entry.get("mailbox_extra") or {})
- registration_proxy_key = str(entry.get("registration_proxy_key") or "").strip()
- registration_proxy_region = str(entry.get("registration_proxy_region") or "").strip()
- registration_proxy_url = str(entry.get("registration_proxy_url") or "").strip()
- registration_fingerprint_profile = str(
- entry.get("registration_fingerprint_profile") or ""
- ).strip()
- cfmail_profile_name = str(entry.get("cfmail_profile_name") or "").strip()
- add_phone_trace_path = str(entry.get("add_phone_trace_path") or "").strip()
- retry_count = int(entry.get("retry_count") or 0) + 1
- self._log(
- f"[zhuce6:register] [pending] 🔄 retrying token acquisition "
- f"for {email} (attempt {retry_count}/{self._pending_token_max_retries})"
- )
- proxy_lease = None
- proxy_url = registration_proxy_url or self.settings.register_proxy
- proxy_release_success = False
- try:
- from core.base_mailbox import MailboxAccount
- from core.cfmail import CfMailMailbox, DEFAULT_CFMAIL_MANAGER
- from platforms.chatgpt.plugin import MailboxEmailServiceAdapter
- from platforms.chatgpt.register import RegistrationEngine
- from platforms.chatgpt.pool import write_token_record
- if self._proxy_pool is not None:
- try:
- proxy_lease = self._proxy_pool.acquire(
- timeout=5.0,
- preferred_name=registration_proxy_key or None,
- preferred_regions=(registration_proxy_region,) if registration_proxy_region else (),
- )
- proxy_url = str(proxy_lease.proxy_url or "").strip() or proxy_url
- self._log(
- "[zhuce6:register] [pending] "
- f"using retry proxy {proxy_lease.name} for {email} "
- f"(origin={registration_proxy_key or registration_proxy_region or registration_proxy_url or '-'})"
- )
- except Exception as exc:
- self._log(
- "[zhuce6:register] [pending] "
- f"proxy acquire failed for {email}: {exc}"
- )
- mailbox = CfMailMailbox(manager=DEFAULT_CFMAIL_MANAGER)
- adapter = MailboxEmailServiceAdapter(mailbox)
- # Reconstruct the mailbox account so _wait_for_mailbox_code can poll
- if mailbox_jwt and mailbox_extra:
- adapter._account = MailboxAccount(
- email=email,
- account_id=mailbox_jwt,
- extra=dict(mailbox_extra),
- )
- engine = RegistrationEngine(
- email_service=adapter,
- proxy_url=proxy_url,
- )
- engine.email = email
- engine.password = password
- token_info = engine._login_for_token()
- if token_info:
- proxy_release_success = True
- self._log(f"[zhuce6:register] [pending] ✅ deferred token acquired for {email}")
- # Write to pool
- token_data = {
- "type": "codex",
- "email": email,
- "password": password,
- "mail_provider": "cfmail",
- "expired": str(token_info.get("expired") or ""),
- "id_token": str(token_info.get("id_token") or ""),
- "account_id": str(token_info.get("account_id") or ""),
- "access_token": str(token_info.get("access_token") or ""),
- "last_refresh": str(token_info.get("last_refresh") or ""),
- "refresh_token": str(token_info.get("refresh_token") or ""),
- "source": "deferred_retry",
- "add_phone_trace_path": add_phone_trace_path,
- }
- token_data.update(
- build_registration_provenance(
- {
- "location": "",
- "mail_provider": "cfmail",
- "post_create_gate": "add_phone",
- },
- proxy_url=registration_proxy_url or proxy_url,
- proxy_key=registration_proxy_key or getattr(proxy_lease, "name", ""),
- proxy_region=registration_proxy_region,
- cfmail_profile_name=cfmail_profile_name,
- )
- )
- if registration_fingerprint_profile:
- token_data["registration_fingerprint_profile"] = registration_fingerprint_profile
- pool_file = write_token_record(token_data, self.settings.pool_dir)
- sync_ok, sync_error, synced_email = self._sync_cpa_from_success(
- {"pool_file": str(pool_file), "success": True, "stage": "deferred_retry"},
- thread_id=0,
- )
- with self._lock:
- if sync_ok:
- self._total_success += 1
- self._total_cpa_sync_success += 1
- self._pending_token_total_success += 1
- elif sync_error == "warmup pending":
- self._total_warmup_pending += 1
- else:
- self._total_failure += 1
- self._total_cpa_sync_failure += 1
- self._pending_token_total_failed += 1
- if sync_ok:
- self._log(f"[zhuce6:register] [pending] ✅ CPA sync success: {synced_email or email}")
- elif sync_error == "warmup pending":
- self._log(f"[zhuce6:register] [pending] ⏳ warmup pending: {synced_email or email}")
- else:
- self._log(f"[zhuce6:register] [pending] ❌ failed [stage=cpa_sync]: {sync_error or 'unknown'}")
- return
- else:
- self._log(f"[zhuce6:register] [pending] ⏳ deferred retry failed for {email}")
- except Exception as exc:
- self._log(f"[zhuce6:register] [pending] ⚠️ deferred retry error for {email}: {exc}")
- finally:
- if proxy_lease is not None and self._proxy_pool is not None:
- try:
- self._proxy_pool.release(
- proxy_lease,
- success=proxy_release_success,
- stage="deferred_retry",
- )
- except Exception as exc:
- self._log(f"[zhuce6:register] [pending] proxy release failed for {email}: {exc}")
- # Re-enqueue with incremented retry count
- entry["retry_count"] = retry_count
- entry["last_retry_at"] = time.time()
- with self._pending_token_lock:
- self._pending_token_queue.append(entry)
- def _worker(self, thread_id: int, initial_provider: str) -> None:
- provider = initial_provider
- consecutive_failures = 0
- max_failures = self.settings.register_max_consecutive_failures
- while not self._should_stop():
- if provider == "cfmail" and self._cfmail_tracker is not None:
- self._cfmail_rotation_pause.wait()
- if self._wait_if_cfmail_add_phone_stopped(thread_id, provider):
- continue
- if self._wait_if_cfmail_wait_otp_stopped(thread_id, provider):
- continue
- if self._wait_if_cfmail_canary_pending(thread_id, provider):
- continue
- if provider == "cfmail" and self._cfmail_manager is not None:
- if self._check_cfmail_all_cooldown_rotation(thread_id):
- continue
- if self._cfmail_all_accounts_in_cooldown():
- self._stop_event.wait(10)
- continue
- if self._wait_if_cfmail_flow_throttled(thread_id, provider):
- continue
- try:
- result: dict[str, object] = {}
- proxy_key = ""
- proxy_outcome: bool | None = False
- result_metadata: dict[str, object] = {}
- result_stage = "?"
- result_error = ""
- result_email = ""
- max_proxy_attempts = 2 if self._proxy_pool is not None else 1
- for proxy_attempt in range(1, max_proxy_attempts + 1):
- proxy_lease = None
- proxy_url = self.settings.register_proxy
- proxy_key = proxy_url or ""
- proxy_outcome = False
- release_stage = "exception"
- try:
- if self._proxy_pool is not None:
- proxy_lease = self._proxy_pool.acquire(
- timeout=5,
- preferred_regions=tuple(self.settings.register_fresh_proxy_regions or ()),
- )
- proxy_url = proxy_lease.proxy_url
- proxy_key = proxy_lease.name or proxy_url or ""
- self._log(
- f"[zhuce6:register] [thread-{thread_id}] acquired proxy {proxy_lease.local_port} ({proxy_lease.name})"
- )
- self._log(f"[zhuce6:register] [thread-{thread_id}] attempting (provider={provider})")
- cfmail_profile_name = self._selected_cfmail_profile(thread_id) if provider == "cfmail" else "auto"
- result = _compat_main_attr("run_chatgpt_register_once", run_chatgpt_register_once)(
- email=None,
- password=None,
- mail_provider=provider,
- cfmail_profile_name=cfmail_profile_name,
- proxy=proxy_url,
- write_pool=True,
- pool_dir=self.settings.pool_dir,
- )
- proxy_outcome = self._classify_proxy_outcome(result)
- result_metadata = result.get("metadata") if isinstance(result.get("metadata"), dict) else {}
- result_stage = str(result.get("stage") or "?").strip() or "?"
- if bool(result.get("success")) and result_stage == "?":
- result_stage = "completed"
- result_error = str(result.get("error_message") or "").strip()
- result_email = str(result.get("email") or "").strip()
- release_stage = result_stage
- should_retry_device_id = (
- self._proxy_pool is not None
- and proxy_attempt < max_proxy_attempts
- and not bool(result.get("success"))
- and result_stage == "device_id"
- )
- if should_retry_device_id:
- self._log(
- f"[zhuce6:register] [thread-{thread_id}] device_id failed on proxy {proxy_key}; "
- "rotating proxy and retrying once"
- )
- continue
- break
- finally:
- if proxy_lease is not None and self._proxy_pool is not None:
- try:
- self._proxy_pool.release(
- proxy_lease,
- success=proxy_outcome,
- stage=release_stage,
- )
- except Exception as exc:
- self._log(f"[zhuce6:register] [thread-{thread_id}] proxy release failed: {exc}")
- success = bool(result.get("success"))
- should_break_after_iteration = False
- with self._lock:
- self._total_attempts += 1
- if success:
- pool_file_raw = str(result.get("pool_file") or "").strip()
- if pool_file_raw:
- try:
- update_token_record(
- Path(pool_file_raw),
- **build_registration_provenance(
- result_metadata,
- proxy_url=proxy_url,
- proxy_key=proxy_key,
- proxy_region=infer_proxy_region(proxy_key or proxy_url),
- cfmail_profile_name=str(result_metadata.get("cfmail_profile_name") or ""),
- ),
- )
- except Exception as exc:
- self._log(
- f"[zhuce6:register] [thread-{thread_id}] provenance update failed: {exc}"
- )
- sync_ok, sync_error, synced_email = self._sync_cpa_from_success(result, thread_id)
- with self._lock:
- if sync_ok:
- self._total_success += 1
- self._total_cpa_sync_success += 1
- self._last_error = None
- consecutive_failures = 0
- self._record_attempt(
- success=True,
- stage=result_stage,
- error_message="",
- metadata=result_metadata,
- proxy_key=proxy_key,
- email=result_email or synced_email,
- )
- email = result_email or synced_email or "?"
- self._log(f"[zhuce6:register] [thread-{thread_id}] \u2705 success: {email}")
- self._log(f"[zhuce6:register] [thread-{thread_id}] ✅ CPA sync success: {email}")
- if self._check_target():
- self._log(
- f"[zhuce6:register] target reached ({self.settings.register_target_count}), stopping"
- )
- should_break_after_iteration = True
- elif sync_error == "warmup pending":
- self._total_warmup_pending += 1
- self._last_error = None
- consecutive_failures = 0
- self._record_attempt(
- success=False,
- stage="warmup_pending",
- error_message="warmup pending",
- metadata=result_metadata,
- proxy_key=proxy_key,
- email=result_email or synced_email,
- )
- email = result_email or synced_email or "?"
- self._log(f"[zhuce6:register] [thread-{thread_id}] ⏳ warmup pending: {email}")
- else:
- self._total_failure += 1
- self._total_cpa_sync_failure += 1
- consecutive_failures += 1
- err = sync_error or "CPA sync failed"
- self._last_error = err
- self._record_attempt(
- success=False,
- stage="cpa_sync",
- error_message=err,
- metadata=result_metadata,
- proxy_key=proxy_key,
- email=result_email or synced_email,
- )
- self._log(
- f"[zhuce6:register] [thread-{thread_id}] \u274c failed ({consecutive_failures}/{max_failures}) "
- f"[stage=cpa_sync]: {err}"
- )
- else:
- with self._lock:
- self._total_failure += 1
- consecutive_failures += 1
- err = result_error or "unknown"
- stage = result_stage
- self._last_error = err
- self._record_attempt(
- success=False,
- stage=stage,
- error_message=err,
- metadata=result_metadata,
- proxy_key=proxy_key,
- email=result_email,
- )
- self._log(f"[zhuce6:register] [thread-{thread_id}] \u274c failed ({consecutive_failures}/{max_failures}) [stage={stage}]: {err}")
- for log_line in result.get("logs", []):
- self._log(f"[zhuce6:register] [thread-{thread_id}] \u21b3 {log_line}")
- # Solution B: enqueue add_phone_gate accounts for deferred retry
- if result_stage == "add_phone_gate":
- self._enqueue_pending_token(
- result,
- thread_id,
- proxy_key=proxy_key,
- proxy_url=proxy_url,
- )
- invalid_mailbox_rotation = self._force_rotate_cfmail_for_invalid_mailbox(
- thread_id,
- result,
- )
- canary_rotation = self._rotate_cfmail_for_failed_canary(thread_id, result)
- rotation_success = self._handle_cfmail_rotation(
- thread_id=thread_id,
- result=result,
- proxy_key=proxy_key,
- )
- self._update_cfmail_add_phone_stoploss(result)
- self._update_cfmail_wait_otp_stoploss(result)
- self._update_cfmail_canary_after_result(thread_id=thread_id, result=result)
- self._update_cfmail_fresh_domain_budget(result)
- fresh_domain_rotation = self._rotate_cfmail_for_fresh_domain_budget(thread_id)
- if invalid_mailbox_rotation or canary_rotation or rotation_success or fresh_domain_rotation:
- consecutive_failures = 0
- self._write_runtime_state()
- if should_break_after_iteration:
- break
- except Exception as exc:
- with self._lock:
- self._total_attempts += 1
- self._total_failure += 1
- self._last_error = str(exc)
- self._record_attempt(
- success=False,
- stage="exception",
- error_message=str(exc),
- metadata={"mail_provider": provider},
- proxy_key=proxy_key,
- email="",
- )
- consecutive_failures += 1
- self._log(f"[zhuce6:register] [thread-{thread_id}] \u274c exception ({consecutive_failures}/{max_failures}): {exc}")
- self._write_runtime_state()
- finally:
- self._release_cfmail_flow_slot(thread_id)
- # Fallback: switch provider after N consecutive failures
- if consecutive_failures >= max_failures and len(self._providers) > 1:
- old_provider = provider
- current_idx = self._providers.index(provider) if provider in self._providers else 0
- provider = self._providers[(current_idx + 1) % len(self._providers)]
- consecutive_failures = 0
- self._log(
- f"[zhuce6:register] [thread-{thread_id}] [fallback] switching {old_provider} -> {provider}"
- )
- # Random sleep between attempts
- sleep_sec = random.randint(
- self.settings.register_sleep_min,
- max(self.settings.register_sleep_min, self.settings.register_sleep_max),
- )
- if self._stop_event.wait(sleep_sec) or self._target_reached.is_set():
- break
- def _cfmail_all_accounts_in_cooldown(self) -> bool:
- manager = self._cfmail_manager
- if manager is None:
- return False
- try:
- manager.reload_if_needed()
- except Exception:
- pass
- return manager.select_account() is None
- def _current_cfmail_active_domain(self) -> str:
- accounts = self._current_cfmail_active_accounts()
- if accounts:
- return str(accounts[0]["domain"]).strip().lower()
- return ""
- def _should_abort_cfmail_wait(self, account: Any) -> bool:
- try:
- extra = account.extra if hasattr(account, "extra") and isinstance(account.extra, dict) else {}
- domain = str(extra.get("email_domain") or "").strip().lower()
- if not domain:
- email = str(getattr(account, "email", "") or "").strip().lower()
- if "@" in email:
- domain = email.rsplit("@", 1)[-1].strip().lower()
- if not domain:
- return False
- stoploss = self._cfmail_wait_otp_stoploss_snapshot()
- if not bool(stoploss.get("in_cooldown")):
- return False
- if str(stoploss.get("active_domain") or "").strip().lower() != domain:
- return False
- wait_started_at = float(extra.get("otp_wait_started_at") or 0.0)
- triggered_at_raw = str(stoploss.get("last_triggered_at") or "").strip()
- if wait_started_at > 0.0 and triggered_at_raw:
- try:
- triggered_at = datetime.fromisoformat(triggered_at_raw).timestamp()
- except Exception:
- triggered_at = 0.0
- if triggered_at > 0.0 and wait_started_at < triggered_at:
- return False
- return True
- except Exception:
- return False
- def _check_cfmail_all_cooldown_rotation(self, thread_id: int) -> bool:
- """When all cfmail accounts are in cooldown, proactively trigger domain rotation."""
- if self._cfmail_tracker is None or self._cfmail_provisioner is None:
- return False
- if not self._cfmail_all_accounts_in_cooldown():
- return False
- if not self._cfmail_rotation_lock.acquire(blocking=False):
- return False
- self._cfmail_rotation_pause.clear()
- try:
- self._log(
- f"[zhuce6:register] [thread-{thread_id}] [cfmail] all accounts in cooldown, "
- "forcing domain rotation to break deadlock"
- )
- provision_result = self._cfmail_provisioner.rotate_active_domain()
- if not provision_result.success:
- if provision_result.old_domain:
- self._cfmail_tracker.mark_rotation_failed(
- provision_result.old_domain,
- provision_result.error,
- )
- self._log(
- f"[zhuce6:register] [thread-{thread_id}] [cfmail] deadlock rotation failed: "
- f"{provision_result.error}"
- )
- return False
- self._cfmail_tracker.mark_rotation_completed(
- provision_result.old_domain,
- provision_result.new_domain,
- )
- self._reload_cfmail_manager_after_rotation()
- self._reset_cfmail_add_phone_stoploss(provision_result.new_domain)
- self._reset_cfmail_wait_otp_stoploss(provision_result.new_domain)
- self._reset_cfmail_fresh_domain_budget(provision_result.new_domain)
- self._arm_cfmail_canary(provision_result.new_domain)
- self._log(
- f"[zhuce6:register] [thread-{thread_id}] [cfmail] deadlock rotation completed: "
- f"{provision_result.old_domain} -> {provision_result.new_domain}"
- )
- return True
- finally:
- self._cfmail_rotation_pause.set()
- self._cfmail_rotation_lock.release()
- def _current_cfmail_active_accounts(self) -> list[dict[str, str]]:
- manager = self._cfmail_manager
- accounts = list(getattr(manager, "accounts", []) or [])
- active_accounts: list[dict[str, str]] = []
- for item in accounts:
- if not bool(getattr(item, "email_domain", "") or ""):
- continue
- if manager is not None and callable(getattr(manager, "skip_remaining_seconds", None)):
- try:
- if int(manager.skip_remaining_seconds(getattr(item, "name", "")) or 0) > 0:
- continue
- except Exception:
- pass
- active_accounts.append(
- {
- "name": str(getattr(item, "name", "") or "").strip(),
- "domain": str(getattr(item, "email_domain", "") or "").strip().lower(),
- }
- )
- return [item for item in active_accounts if item["name"] and item["domain"]]
- def _cfmail_domain_pool_snapshot(
- self,
- recent_attempts: list[dict[str, object]],
- ) -> dict[str, object]:
- active_accounts = self._current_cfmail_active_accounts()
- with self._lock:
- inflight_by_thread = dict(self._cfmail_flow_state.get("inflight_by_thread") or {})
- last_started_by_domain = dict(self._cfmail_flow_state.get("last_started_by_domain") or {})
- replenish_thread = self._cfmail_replenish_thread
- replenish_reason = self._cfmail_replenish_reason
- domains: list[dict[str, object]] = []
- now_ts = time.time()
- manager = self._cfmail_manager
- for account in active_accounts:
- domain = account["domain"]
- profile_name = account["name"]
- domain_attempts = self._active_domain_attempts(recent_attempts, domain)
- failure_by_stage, failure_signals = self._failure_counts_from_attempts(domain_attempts)
- recent_success = sum(1 for item in domain_attempts if bool(item.get("success")))
- recent_failure = sum(1 for item in domain_attempts if not bool(item.get("success")))
- inflight = sum(
- 1
- for value in inflight_by_thread.values()
- if isinstance(value, dict)
- and str(value.get("domain") or "").strip().lower() == domain
- )
- last_started_at = float(last_started_by_domain.get(domain) or 0.0)
- start_interval_remaining = 0.0
- if self._cfmail_start_interval_seconds > 0 and last_started_at > 0:
- start_interval_remaining = max(
- 0.0,
- (last_started_at + self._cfmail_start_interval_seconds) - now_ts,
- )
- skip_remaining_seconds = 0
- if manager is not None and callable(getattr(manager, "skip_remaining_seconds", None)):
- try:
- skip_remaining_seconds = int(manager.skip_remaining_seconds(profile_name) or 0)
- except Exception:
- skip_remaining_seconds = 0
- add_phone_state = self._cfmail_add_phone_stoploss_snapshot()
- wait_otp_state = self._cfmail_wait_otp_stoploss_snapshot()
- domains.append(
- {
- "name": profile_name,
- "domain": domain,
- "inflight": inflight,
- "recent_attempts": len(domain_attempts),
- "recent_success": recent_success,
- "recent_failure": recent_failure,
- "failure_by_stage": failure_by_stage,
- "failure_signals": failure_signals,
- "recent_failure_hotspots": self._recent_failure_hotspots_from_attempts(domain_attempts),
- "last_started_at": datetime.fromtimestamp(last_started_at).isoformat(timespec="seconds")
- if last_started_at > 0
- else "",
- "start_interval_remaining_seconds": round(start_interval_remaining, 1),
- "skip_remaining_seconds": skip_remaining_seconds,
- "add_phone_cooldown": bool(
- add_phone_state.get("in_cooldown")
- and str(add_phone_state.get("active_domain") or "").strip().lower() == domain
- ),
- "wait_otp_cooldown": bool(
- wait_otp_state.get("in_cooldown")
- and str(wait_otp_state.get("active_domain") or "").strip().lower() == domain
- ),
- }
- )
- return {
- "target_count": self._cfmail_active_domain_count,
- "active_count": len(domains),
- "active_domains": domains,
- "replenishing": bool(replenish_thread is not None and replenish_thread.is_alive()),
- "replenish_reason": str(replenish_reason or ""),
- }
- def _selected_cfmail_profile(self, thread_id: int) -> str:
- with self._lock:
- selected_profile_by_thread = self._cfmail_flow_state.setdefault("selected_profile_by_thread", {})
- if not isinstance(selected_profile_by_thread, dict):
- return "auto"
- value = str(selected_profile_by_thread.get(thread_id) or "").strip()
- return value or "auto"
- def _handle_cfmail_rotation(
- self,
- *,
- thread_id: int,
- result: dict[str, object],
- proxy_key: str,
- ) -> bool:
- if self._cfmail_tracker is None or self._cfmail_provisioner is None:
- return False
- metadata = result.get("metadata")
- if not isinstance(metadata, dict):
- metadata = {}
- if str(metadata.get("mail_provider") or result.get("mail_provider") or "").strip().lower() not in {"", "cfmail"}:
- return False
- from core.cfmail_domain_rotation import classify_domain_attempt
- attempt = classify_domain_attempt(result, proxy_key=proxy_key)
- if attempt is None:
- return False
- decision = self._cfmail_tracker.record(attempt)
- if attempt.backend_failure and not decision.should_rotate:
- self._log(
- f"[zhuce6:register] [thread-{thread_id}] [cfmail] backend unhealthy for domain={attempt.domain}; rotation skipped"
- )
- return False
- if not decision.should_rotate:
- return False
- self._log(
- f"[zhuce6:register] [thread-{thread_id}] [cfmail] rotating domain {decision.domain} "
- f"(reason={decision.reason}, failures={decision.blacklist_failures}, window={decision.window_size})"
- )
- return self._replace_cfmail_domain(
- thread_id=thread_id,
- domain=decision.domain,
- reason_label=decision.reason,
- )
- def _classify_proxy_outcome(self, result: dict[str, object]) -> bool | None:
- if bool(result.get("success")):
- return True
- stage = str(result.get("stage") or "").strip().lower()
- metadata = result.get("metadata") if isinstance(result.get("metadata"), dict) else {}
- code = str(metadata.get("create_account_error_code") or "").strip().lower()
- post_create_gate = str(metadata.get("post_create_gate") or "").strip().lower()
- if stage == "create_account" and code in {"registration_disallowed", "unsupported_email"}:
- return None
- if stage == "add_phone_gate" or post_create_gate == "add_phone":
- return None
- if stage in {"mailbox", "device_id", "password", "token_acquisition"}:
- return False
- return False
- def snapshot(self) -> dict[str, object]:
- with self._lock:
- alive = sum(1 for t in self._threads if t.is_alive())
- target = self.settings.register_target_count
- cfmail_rotation = self._cfmail_tracker.snapshot() if self._cfmail_tracker is not None else None
- recent_attempts = list(self._recent_attempts)
- add_phone_stoploss = self._cfmail_add_phone_stoploss_snapshot()
- wait_otp_stoploss = self._cfmail_wait_otp_stoploss_snapshot()
- cfmail_canary = self._cfmail_canary_snapshot()
- cfmail_fresh_domain_budget = self._cfmail_fresh_domain_budget_snapshot()
- cfmail_domain_pool = self._cfmail_domain_pool_snapshot(recent_attempts)
- active_domain = self._infer_active_domain(recent_attempts, cfmail_rotation, add_phone_stoploss)
- active_domain_attempts = self._active_domain_attempts(recent_attempts, active_domain)
- active_failure_by_stage, active_failure_signals = self._failure_counts_from_attempts(active_domain_attempts)
- return {
- "name": "register",
- "status": "running" if alive > 0 else ("pending" if self._total_attempts == 0 else "stopped"),
- "threads_alive": alive,
- "threads_total": len(self._threads),
- "total_attempts": self._total_attempts,
- "total_success": self._total_success,
- "total_success_registered": self._total_success,
- "total_success_direct": self._total_success,
- "total_warmup_pending": self._total_warmup_pending,
- "total_cpa_sync_success": self._total_cpa_sync_success,
- "total_cpa_sync_success_direct": self._total_cpa_sync_success,
- "total_cpa_sync_failure": self._total_cpa_sync_failure,
- "total_failure": self._total_failure,
- "success_rate": round(self._total_success / max(self._total_attempts, 1) * 100, 1),
- "registered_success_rate": round(self._total_success / max(self._total_attempts, 1) * 100, 1),
- "cpa_sync_success_rate": round(self._total_cpa_sync_success / max(self._total_attempts, 1) * 100, 1),
- "target_count": target if target > 0 else None,
- "target_reached": self._target_reached.is_set(),
- "last_error": self._last_error,
- "proxy": self.settings.register_proxy,
- "proxy_pool_enabled": self._proxy_pool is not None,
- "mail_provider": self.settings.register_mail_provider,
- "interval_seconds": self.settings.register_interval,
- "run_count": self._total_attempts,
- "success_count": self._total_success,
- "failure_count": self._total_failure,
- "is_running": alive > 0,
- "last_started_at": datetime.fromtimestamp(self._started_at).isoformat(timespec="seconds") if self._started_at else None,
- "last_finished_at": None,
- "last_duration_seconds": None,
- "next_run_at": None,
- "failure_by_stage": dict(sorted(self._failure_by_stage.items(), key=lambda item: (-item[1], item[0]))),
- "failure_signals": dict(sorted(self._failure_signals.items(), key=lambda item: (-item[1], item[0]))),
- "recent_failure_hotspots": self._recent_failure_hotspots(),
- "recent_attempts": recent_attempts,
- "active_domain_recent_attempts": active_domain_attempts,
- "active_domain_failure_by_stage": active_failure_by_stage,
- "active_domain_failure_signals": active_failure_signals,
- "active_domain_recent_failure_hotspots": self._recent_failure_hotspots_from_attempts(active_domain_attempts),
- "cfmail_rotation": cfmail_rotation,
- "cfmail_add_phone_stoploss": add_phone_stoploss,
- "cfmail_wait_otp_stoploss": wait_otp_stoploss,
- "cfmail_canary": cfmail_canary,
- "cfmail_fresh_domain_budget": cfmail_fresh_domain_budget,
- "cfmail_domain_pool": cfmail_domain_pool,
- "pending_token_queue": {
- "queue_size": len(self._pending_token_queue),
- "total_enqueued": self._pending_token_total_enqueued,
- "total_success": self._pending_token_total_success,
- "total_failed": self._pending_token_total_failed,
- "retry_delay_seconds": self._pending_token_retry_delay_seconds,
- "max_retries": self._pending_token_max_retries,
- },
- }
- class RegistrationBurstScheduler:
- """Run registration in timed batches and expose scheduler state via runtime_state.json."""
- def __init__(self, settings: AppSettings) -> None:
- self.settings = settings
- self._stop_event = threading.Event()
- self._lock = threading.RLock()
- self._active_loop: RegistrationLoop | None = None
- self._active_started_at: float | None = None
- self._next_run_at_ts: float | None = None
- self._run_count = 0
- self._total_attempts = 0
- self._total_success = 0
- self._total_cpa_sync_success = 0
- self._total_cpa_sync_failure = 0
- self._total_failure = 0
- self._last_error: str | None = None
- self._recent_attempts: deque[dict[str, object]] = deque(maxlen=80)
- self._failure_by_stage: dict[str, int] = {}
- self._failure_signals: dict[str, int] = {}
- self._last_batch_started_at: str | None = None
- self._last_batch_finished_at: str | None = None
- self._last_batch_duration_seconds: float | None = None
- self._last_cfmail_add_phone_stoploss: dict[str, object] | None = None
- self._last_cfmail_wait_otp_stoploss: dict[str, object] | None = None
- self._last_proxy_pool_snapshot: dict[str, object] = {
- "configured": bool(self.settings.proxy_pool_configured),
- "enabled": False,
- "snapshot_error": None,
- "node_count": 0,
- "in_use_count": 0,
- "disabled_count": 0,
- "nodes": [],
- }
- self._logger = self._setup_logger()
- def _setup_logger(self) -> Any:
- from logging.handlers import RotatingFileHandler
- logger = logging.getLogger("zhuce6.register")
- logger.setLevel(logging.INFO)
- logger.propagate = False
- if not logger.handlers:
- console = logging.StreamHandler()
- console.setFormatter(logging.Formatter("%(message)s"))
- logger.addHandler(console)
- if self.settings.register_log_file:
- fh = RotatingFileHandler(
- self.settings.register_log_file,
- maxBytes=2 * 1024 * 1024,
- backupCount=5,
- encoding="utf-8",
- )
- fh.setFormatter(logging.Formatter("%(asctime)s %(message)s", datefmt="%Y-%m-%d %H:%M:%S"))
- logger.addHandler(fh)
- return logger
- def _log(self, msg: str) -> None:
- self._logger.info(msg)
- def _merge_counts(self, target: dict[str, int], incoming: dict[str, object] | None) -> None:
- if not isinstance(incoming, dict):
- return
- for key, value in incoming.items():
- try:
- inc = int(value or 0)
- except Exception:
- continue
- target[str(key)] = target.get(str(key), 0) + inc
- def _write_runtime_state(self) -> None:
- state_file = Path(self.settings.runtime_state_file)
- try:
- state_file.parent.mkdir(parents=True, exist_ok=True)
- payload = {
- "updated_at": datetime.now().isoformat(timespec="seconds"),
- "register_snapshot": self.snapshot(),
- "proxy_pool": self._current_proxy_pool_snapshot(),
- }
- tmp_file = state_file.with_name(
- f"{state_file.name}.{os.getpid()}.{threading.get_ident()}.tmp"
- )
- tmp_file.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8")
- tmp_file.replace(state_file)
- except Exception as exc:
- self._log(f"[zhuce6:register] burst runtime state write failed: {exc}")
- def _current_proxy_pool_snapshot(self) -> dict[str, object]:
- with self._lock:
- active_loop = self._active_loop
- last_snapshot = dict(self._last_proxy_pool_snapshot)
- if active_loop is not None:
- try:
- return active_loop._proxy_pool_snapshot()
- except Exception as exc:
- last_snapshot["snapshot_error"] = str(exc)
- return last_snapshot
- return last_snapshot
- def _absorb_batch_snapshot(self, snapshot: dict[str, object], *, duration_seconds: float) -> None:
- with self._lock:
- self._run_count += 1
- self._total_attempts += int(snapshot.get("total_attempts") or 0)
- self._total_success += int(snapshot.get("total_success") or 0)
- self._total_cpa_sync_success += int(snapshot.get("total_cpa_sync_success") or 0)
- self._total_cpa_sync_failure += int(snapshot.get("total_cpa_sync_failure") or 0)
- self._total_failure += int(snapshot.get("total_failure") or 0)
- self._last_error = str(snapshot.get("last_error") or "").strip() or None
- self._last_batch_started_at = snapshot.get("last_started_at") if isinstance(snapshot.get("last_started_at"), str) else None
- self._last_batch_finished_at = datetime.now().isoformat(timespec="seconds")
- self._last_batch_duration_seconds = round(duration_seconds, 3)
- self._merge_counts(self._failure_by_stage, snapshot.get("failure_by_stage") if isinstance(snapshot.get("failure_by_stage"), dict) else None)
- self._merge_counts(self._failure_signals, snapshot.get("failure_signals") if isinstance(snapshot.get("failure_signals"), dict) else None)
- attempts = snapshot.get("recent_attempts")
- if isinstance(attempts, list):
- for item in attempts:
- if isinstance(item, dict):
- self._recent_attempts.append(item)
- if isinstance(snapshot.get("cfmail_add_phone_stoploss"), dict):
- self._last_cfmail_add_phone_stoploss = dict(snapshot.get("cfmail_add_phone_stoploss") or {})
- if isinstance(snapshot.get("cfmail_wait_otp_stoploss"), dict):
- self._last_cfmail_wait_otp_stoploss = dict(snapshot.get("cfmail_wait_otp_stoploss") or {})
- def snapshot(self) -> dict[str, object]:
- with self._lock:
- active_loop = self._active_loop
- next_run_at_ts = self._next_run_at_ts
- total_attempts = self._total_attempts
- total_success = self._total_success
- total_cpa_sync_success = self._total_cpa_sync_success
- total_cpa_sync_failure = self._total_cpa_sync_failure
- total_failure = self._total_failure
- run_count = self._run_count
- failure_by_stage = dict(sorted(self._failure_by_stage.items(), key=lambda item: (-item[1], item[0])))
- failure_signals = dict(sorted(self._failure_signals.items(), key=lambda item: (-item[1], item[0])))
- recent_attempts = list(self._recent_attempts)
- last_error = self._last_error
- last_batch_started_at = self._last_batch_started_at
- last_batch_finished_at = self._last_batch_finished_at
- last_batch_duration_seconds = self._last_batch_duration_seconds
- add_phone_stoploss = dict(self._last_cfmail_add_phone_stoploss or {})
- wait_otp_stoploss = dict(self._last_cfmail_wait_otp_stoploss or {})
- if active_loop is not None:
- current = dict(active_loop.snapshot())
- current.update(
- {
- "scheduler_mode": "burst",
- "batch_threads": self.settings.register_batch_threads,
- "batch_target_count": self.settings.register_batch_target_count,
- "batch_interval_seconds": self.settings.register_batch_interval_seconds,
- "run_count": run_count,
- "next_run_at": None,
- }
- )
- return current
- 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")
- counts: dict[tuple[str, str], int] = {}
- for item in recent_attempts:
- if item.get("success"):
- continue
- stage = str(item.get("stage") or "?").strip() or "?"
- signal = str(item.get("signal") or "").strip()
- key = (signal or stage, stage)
- counts[key] = counts.get(key, 0) + 1
- recent_failure_hotspots = [
- {"key": key, "stage": stage, "count": count}
- for (key, stage), count in sorted(
- counts.items(),
- key=lambda kv: (-kv[1], kv[0][0], kv[0][1]),
- )[:5]
- ]
- return {
- "name": "register",
- "status": status,
- "scheduler_mode": "burst",
- "threads_alive": 0,
- "threads_total": self.settings.register_batch_threads,
- "total_attempts": total_attempts,
- "total_success": total_success,
- "total_success_registered": total_success,
- "total_success_direct": total_success,
- "total_cpa_sync_success": total_cpa_sync_success,
- "total_cpa_sync_success_direct": total_cpa_sync_success,
- "total_cpa_sync_failure": total_cpa_sync_failure,
- "total_failure": total_failure,
- "success_rate": round(total_success / max(total_attempts, 1) * 100, 1),
- "registered_success_rate": round(total_success / max(total_attempts, 1) * 100, 1),
- "cpa_sync_success_rate": round(total_cpa_sync_success / max(total_attempts, 1) * 100, 1),
- "target_count": self.settings.register_batch_target_count,
- "target_reached": False,
- "last_error": last_error,
- "proxy": self.settings.register_proxy,
- "proxy_pool_enabled": bool(self.settings.proxy_pool_configured),
- "mail_provider": self.settings.register_mail_provider,
- "interval_seconds": self.settings.register_interval,
- "run_count": run_count,
- "success_count": total_success,
- "failure_count": total_failure,
- "is_running": False,
- "last_started_at": last_batch_started_at,
- "last_finished_at": last_batch_finished_at,
- "last_duration_seconds": last_batch_duration_seconds,
- "next_run_at": datetime.fromtimestamp(next_run_at_ts).isoformat(timespec="seconds") if next_run_at_ts else None,
- "failure_by_stage": failure_by_stage,
- "failure_signals": failure_signals,
- "recent_failure_hotspots": recent_failure_hotspots,
- "recent_attempts": recent_attempts,
- "active_domain_recent_attempts": [],
- "active_domain_failure_by_stage": {},
- "active_domain_failure_signals": {},
- "active_domain_recent_failure_hotspots": [],
- "cfmail_rotation": None,
- "cfmail_add_phone_stoploss": add_phone_stoploss,
- "cfmail_wait_otp_stoploss": wait_otp_stoploss,
- "batch_threads": self.settings.register_batch_threads,
- "batch_target_count": self.settings.register_batch_target_count,
- "batch_interval_seconds": self.settings.register_batch_interval_seconds,
- }
- def stop(self) -> None:
- self._stop_event.set()
- with self._lock:
- active_loop = self._active_loop
- if active_loop is not None:
- active_loop.stop()
- self._write_runtime_state()
- def run(self) -> None:
- self._next_run_at_ts = time.time()
- self._write_runtime_state()
- while not self._stop_event.is_set():
- now = time.time()
- next_run_at_ts = self._next_run_at_ts or now
- if now < next_run_at_ts:
- self._write_runtime_state()
- self._stop_event.wait(min(max(next_run_at_ts - now, 1), 5))
- continue
- batch_settings = replace(
- self.settings,
- register_enabled=True,
- register_threads=self.settings.register_batch_threads,
- register_target_count=self.settings.register_batch_target_count,
- )
- loop = _compat_main_attr("RegistrationLoop", RegistrationLoop)(batch_settings)
- started_at = time.time()
- with self._lock:
- self._active_loop = loop
- self._active_started_at = started_at
- self._write_runtime_state()
- loop.start()
- try:
- while not self._stop_event.is_set():
- snapshot = loop.snapshot()
- if int(snapshot.get("threads_alive") or 0) <= 0:
- break
- self._write_runtime_state()
- self._stop_event.wait(1)
- finally:
- loop.stop()
- batch_snapshot = loop.snapshot()
- with self._lock:
- self._active_loop = None
- self._active_started_at = None
- proxy_snapshot_fn = getattr(loop, "_proxy_pool_snapshot", None)
- if callable(proxy_snapshot_fn):
- self._last_proxy_pool_snapshot = proxy_snapshot_fn()
- self._absorb_batch_snapshot(batch_snapshot, duration_seconds=time.time() - started_at)
- self._next_run_at_ts = started_at + self.settings.register_batch_interval_seconds
- self._write_runtime_state()
- if self._stop_event.is_set():
- break
|