test_registration_loop.py 81 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791179217931794179517961797179817991800180118021803180418051806180718081809181018111812181318141815181618171818181918201821182218231824182518261827182818291830183118321833183418351836183718381839184018411842184318441845184618471848184918501851185218531854185518561857185818591860186118621863186418651866186718681869187018711872187318741875187618771878187918801881188218831884188518861887188818891890189118921893189418951896189718981899190019011902190319041905190619071908190919101911191219131914191519161917191819191920192119221923192419251926192719281929193019311932193319341935193619371938193919401941194219431944194519461947194819491950195119521953195419551956195719581959196019611962196319641965196619671968196919701971197219731974197519761977197819791980198119821983198419851986198719881989199019911992199319941995199619971998199920002001200220032004200520062007200820092010201120122013201420152016201720182019202020212022202320242025202620272028202920302031203220332034203520362037203820392040204120422043204420452046204720482049205020512052205320542055205620572058205920602061206220632064206520662067206820692070207120722073207420752076207720782079
  1. import json
  2. from dataclasses import replace
  3. from pathlib import Path
  4. from core.base_mailbox import MailboxAccount
  5. from core.settings import AppSettings
  6. from core.cfmail_domain_rotation import DomainHealthTracker
  7. from core.cfmail_provisioner import ProvisionResult
  8. from main import RegistrationBurstScheduler, RegistrationLoop
  9. def _base_settings(**overrides) -> AppSettings:
  10. settings = AppSettings(
  11. cleanup_enabled=False,
  12. register_sleep_min=0,
  13. register_sleep_max=0,
  14. register_mail_provider="mailtm,mailgw",
  15. )
  16. return replace(settings, **overrides)
  17. def test_registration_loop_switches_provider_after_configured_failures(monkeypatch) -> None:
  18. settings = _base_settings(register_max_consecutive_failures=2)
  19. loop = RegistrationLoop(settings)
  20. loop._providers = ["mailtm", "mailgw"]
  21. calls: list[str] = []
  22. def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
  23. calls.append(kwargs["mail_provider"])
  24. if len(calls) >= 3:
  25. loop._stop_event.set()
  26. return {"success": False, "error_message": "failed"}
  27. monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
  28. loop._worker(thread_id=1, initial_provider="mailtm")
  29. assert calls == ["mailtm", "mailtm", "mailgw"]
  30. def test_registration_loop_stops_when_target_count_reached(monkeypatch) -> None:
  31. settings = _base_settings(register_target_count=2, register_mail_provider="mailtm")
  32. loop = RegistrationLoop(settings)
  33. loop._providers = ["mailtm"]
  34. calls: list[str] = []
  35. def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
  36. calls.append(kwargs["mail_provider"])
  37. return {"success": True, "email": f"user{len(calls)}@example.com"}
  38. monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
  39. monkeypatch.setattr(loop, "_sync_cpa_from_success", lambda result, thread_id: (True, "", str(result.get("email") or "")))
  40. loop._worker(thread_id=1, initial_provider="mailtm")
  41. assert calls == ["mailtm", "mailtm"]
  42. assert loop._target_reached.is_set() is True
  43. def test_registration_loop_uses_register_proxy_when_proxy_pool_disabled(monkeypatch) -> None:
  44. settings = _base_settings(register_proxy="http://127.0.0.1:7890", register_mail_provider="mailtm")
  45. loop = RegistrationLoop(settings)
  46. loop._providers = ["mailtm"]
  47. proxies: list[str | None] = []
  48. def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
  49. proxies.append(kwargs["proxy"])
  50. loop._stop_event.set()
  51. return {"success": False, "error_message": "failed"}
  52. monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
  53. loop._worker(thread_id=1, initial_provider="mailtm")
  54. assert proxies == ["http://127.0.0.1:7890"]
  55. def test_registration_loop_passes_selected_cfmail_profile_into_register_flow(monkeypatch) -> None:
  56. settings = _base_settings(register_mail_provider="cfmail")
  57. loop = RegistrationLoop(settings)
  58. loop._providers = ["cfmail"]
  59. recorded: dict[str, object] = {}
  60. monkeypatch.setattr(
  61. loop,
  62. "_current_cfmail_active_accounts",
  63. lambda: [{"name": "cfmail-tw", "domain": "tw.example.test"}],
  64. )
  65. def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
  66. recorded["cfmail_profile_name"] = kwargs["cfmail_profile_name"]
  67. loop._stop_event.set()
  68. return {"success": True, "email": "demo@example.com", "metadata": {"email_domain": "tw.example.test"}}
  69. monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
  70. monkeypatch.setattr(loop, "_sync_cpa_from_success", lambda result, thread_id: (True, "", str(result.get("email") or "")))
  71. loop._worker(thread_id=1, initial_provider="cfmail")
  72. assert recorded["cfmail_profile_name"] == "cfmail-tw"
  73. def test_registration_loop_snapshot_exposes_cfmail_domain_pool(monkeypatch) -> None:
  74. settings = _base_settings(register_mail_provider="cfmail")
  75. loop = RegistrationLoop(settings)
  76. monkeypatch.setattr(
  77. loop,
  78. "_current_cfmail_active_accounts",
  79. lambda: [
  80. {"name": "cfmail-tw", "domain": "tw.example.test"},
  81. {"name": "cfmail-sg", "domain": "sg.example.test"},
  82. ],
  83. )
  84. loop._recent_attempts.extend(
  85. [
  86. {
  87. "timestamp": "2026-03-30T22:50:00",
  88. "success": True,
  89. "stage": "completed",
  90. "signal": "",
  91. "error_message": "",
  92. "email_domain": "tw.example.test",
  93. "proxy_key": "tw-node",
  94. "email": "ok@tw.example.test",
  95. },
  96. {
  97. "timestamp": "2026-03-30T22:50:10",
  98. "success": False,
  99. "stage": "create_account",
  100. "signal": "mailbox_reused",
  101. "error_message": "create account failed",
  102. "email_domain": "sg.example.test",
  103. "create_account_error_code": "user_already_exists",
  104. "proxy_key": "sg-node",
  105. "email": "dup@sg.example.test",
  106. },
  107. ]
  108. )
  109. loop._cfmail_flow_state["inflight_by_thread"] = {
  110. 1: {"domain": "tw.example.test", "profile_name": "cfmail-tw"},
  111. }
  112. loop._cfmail_flow_state["last_started_by_domain"] = {
  113. "tw.example.test": 100.0,
  114. "sg.example.test": 90.0,
  115. }
  116. snapshot = loop.snapshot()
  117. domain_pool = snapshot["cfmail_domain_pool"]
  118. assert domain_pool["active_count"] == 2
  119. assert [item["domain"] for item in domain_pool["active_domains"]] == [
  120. "tw.example.test",
  121. "sg.example.test",
  122. ]
  123. assert domain_pool["active_domains"][0]["inflight"] == 1
  124. assert domain_pool["active_domains"][0]["recent_success"] == 1
  125. assert domain_pool["active_domains"][1]["recent_failure"] == 1
  126. assert domain_pool["active_domains"][1]["failure_signals"]["mailbox_reused"] == 1
  127. def test_registration_loop_uses_multi_domain_flow_defaults(monkeypatch) -> None:
  128. monkeypatch.delenv("ZHUCE6_CFMAIL_START_INTERVAL_SECONDS", raising=False)
  129. monkeypatch.delenv("ZHUCE6_CFMAIL_MAX_INFLIGHT", raising=False)
  130. monkeypatch.delenv("ZHUCE6_CFMAIL_ACTIVE_DOMAIN_COUNT", raising=False)
  131. monkeypatch.delenv("ZHUCE6_CFMAIL_FRESH_DOMAIN_ATTEMPT_BUDGET", raising=False)
  132. loop = RegistrationLoop(_base_settings(register_mail_provider="cfmail"))
  133. assert loop._cfmail_start_interval_seconds == 8
  134. assert loop._cfmail_max_inflight == 4
  135. assert loop._cfmail_active_domain_count == 3
  136. assert loop._cfmail_fresh_domain_attempt_budget == 2
  137. assert loop._cfmail_add_phone_window == 10
  138. assert loop._cfmail_add_phone_threshold == 3
  139. assert loop._cfmail_wait_otp_window == 6
  140. assert loop._cfmail_wait_otp_threshold == 2
  141. def test_registration_loop_caps_pending_token_retry_delay_to_one_minute(monkeypatch) -> None:
  142. monkeypatch.setenv("ZHUCE6_PENDING_TOKEN_RETRY_DELAY_SECONDS", "300")
  143. loop = RegistrationLoop(_base_settings())
  144. assert loop._pending_token_retry_delay_seconds == 60
  145. def test_registration_loop_uses_proxy_pool_and_releases_lease(monkeypatch) -> None:
  146. settings = _base_settings(register_mail_provider="mailtm")
  147. loop = RegistrationLoop(settings)
  148. loop._providers = ["mailtm"]
  149. recorded: dict[str, object] = {}
  150. class FakeLease:
  151. name = "sg-node"
  152. local_port = 17891
  153. proxy_url = "socks5://127.0.0.1:17891"
  154. class FakePool:
  155. def acquire(self, timeout=5.0, preferred_name=None, preferred_regions=()): # type: ignore[no-untyped-def]
  156. recorded["timeout"] = timeout
  157. recorded["preferred_name"] = preferred_name
  158. recorded["preferred_regions"] = tuple(preferred_regions)
  159. return FakeLease()
  160. def release(self, lease, *, success, stage=None): # type: ignore[no-untyped-def]
  161. recorded["released"] = (lease.proxy_url, success, stage)
  162. def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
  163. recorded["proxy"] = kwargs["proxy"]
  164. loop._stop_event.set()
  165. return {"success": True, "email": "demo@example.com"}
  166. loop._proxy_pool = FakePool()
  167. monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
  168. monkeypatch.setattr(loop, "_sync_cpa_from_success", lambda result, thread_id: (True, "", str(result.get("email") or "")))
  169. loop._worker(thread_id=1, initial_provider="mailtm")
  170. assert recorded["proxy"] == "socks5://127.0.0.1:17891"
  171. assert recorded["released"] == ("socks5://127.0.0.1:17891", True, "completed")
  172. def test_registration_loop_prefers_fresh_proxy_regions(monkeypatch) -> None:
  173. settings = _base_settings(register_mail_provider="mailtm", register_fresh_proxy_regions=("tw", "sg"))
  174. loop = RegistrationLoop(settings)
  175. loop._providers = ["mailtm"]
  176. recorded: dict[str, object] = {}
  177. class FakeLease:
  178. name = "tw-node"
  179. local_port = 17891
  180. proxy_url = "socks5://127.0.0.1:17891"
  181. class FakePool:
  182. def acquire(self, timeout=5.0, preferred_name=None, preferred_regions=()): # type: ignore[no-untyped-def]
  183. recorded["timeout"] = timeout
  184. recorded["preferred_name"] = preferred_name
  185. recorded["preferred_regions"] = tuple(preferred_regions)
  186. return FakeLease()
  187. def release(self, lease, *, success, stage=None): # type: ignore[no-untyped-def]
  188. recorded["released"] = (lease.proxy_url, success, stage)
  189. def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
  190. recorded["proxy"] = kwargs["proxy"]
  191. loop._stop_event.set()
  192. return {"success": True, "email": "demo@example.com"}
  193. loop._proxy_pool = FakePool()
  194. monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
  195. monkeypatch.setattr(loop, "_sync_cpa_from_success", lambda result, thread_id: (True, "", str(result.get("email") or "")))
  196. loop._worker(thread_id=1, initial_provider="mailtm")
  197. assert recorded["preferred_regions"] == ("tw", "sg")
  198. def test_registration_loop_starts_proxy_pool_for_direct_urls(monkeypatch) -> None:
  199. settings = _base_settings(
  200. register_mail_provider="mailtm",
  201. register_threads=1,
  202. proxy_pool_direct_urls="http://5.6.7.8:8080",
  203. )
  204. loop = RegistrationLoop(settings)
  205. recorded: dict[str, object] = {}
  206. class FakePool:
  207. def start(self) -> None:
  208. recorded["pool_started"] = True
  209. class FakeThread:
  210. def __init__(self, target=None, args=(), daemon=None, name=None): # type: ignore[no-untyped-def]
  211. recorded["thread_args"] = {"target": target, "args": args, "daemon": daemon, "name": name}
  212. def start(self) -> None:
  213. recorded["thread_started"] = True
  214. monkeypatch.setattr("main.load_all", lambda: None)
  215. monkeypatch.setattr("core.proxy_pool.ProxyPool.from_settings", lambda settings: FakePool())
  216. monkeypatch.setattr("main.threading.Thread", FakeThread)
  217. monkeypatch.setattr("core.registration._maybe_reconcile_cpa_runtime", lambda **kwargs: None)
  218. monkeypatch.setattr(loop, "_write_runtime_state", lambda: None)
  219. loop.start()
  220. assert recorded["pool_started"] is True
  221. assert recorded["thread_started"] is True
  222. assert loop._proxy_pool is not None
  223. def test_registration_loop_start_reconciles_pool_backups_to_cpa(monkeypatch, tmp_path: Path) -> None:
  224. settings = _base_settings(
  225. register_mail_provider="mailtm",
  226. register_threads=1,
  227. backend="cpa",
  228. cpa_runtime_reconcile_enabled=True,
  229. pool_dir=tmp_path / "pool",
  230. )
  231. loop = RegistrationLoop(settings)
  232. recorded: dict[str, object] = {}
  233. class FakeThread:
  234. def __init__(self, target=None, args=(), daemon=None, name=None): # type: ignore[no-untyped-def]
  235. recorded["thread_args"] = {"target": target, "args": args, "daemon": daemon, "name": name}
  236. def start(self) -> None:
  237. recorded["thread_started"] = True
  238. monkeypatch.setattr("main.load_all", lambda: None)
  239. monkeypatch.setattr("main.threading.Thread", FakeThread)
  240. monkeypatch.setattr("core.registration.create_backend_client", lambda settings: "backend-client")
  241. monkeypatch.setattr(
  242. "core.registration._maybe_reconcile_cpa_runtime",
  243. lambda **kwargs: recorded.setdefault("reconcile", kwargs),
  244. )
  245. monkeypatch.setattr(loop, "_write_runtime_state", lambda: None)
  246. loop.start()
  247. assert recorded["thread_started"] is True
  248. assert recorded["reconcile"]["pool_dir"] == settings.pool_dir
  249. assert recorded["reconcile"]["client"] == "backend-client"
  250. def test_registration_loop_start_normalizes_cfmail_to_domain_pool(monkeypatch) -> None:
  251. settings = _base_settings(
  252. register_mail_provider="cfmail",
  253. register_threads=1,
  254. backend="cpa",
  255. cpa_runtime_reconcile_enabled=False,
  256. )
  257. loop = RegistrationLoop(settings)
  258. recorded: dict[str, object] = {}
  259. class FakeThread:
  260. def __init__(self, target=None, args=(), daemon=None, name=None): # type: ignore[no-untyped-def]
  261. recorded["thread_args"] = {"target": target, "args": args, "daemon": daemon, "name": name}
  262. def start(self) -> None:
  263. recorded["thread_started"] = True
  264. class FakeProvisioner:
  265. def __init__(self, *, proxy_url=None, **_kwargs): # type: ignore[no-untyped-def]
  266. recorded["proxy_url"] = proxy_url
  267. def normalize_to_domain_pool(self, target_count): # type: ignore[no-untyped-def]
  268. recorded["normalized"] = True
  269. recorded["target_count"] = target_count
  270. return {
  271. "active_domains": ["auto-live.example.test", "auto-next.example.test"],
  272. "provisioned_domains": ["auto-next.example.test"],
  273. "retired_domains": [],
  274. }
  275. class FakeManager:
  276. def reload_if_needed(self, force=False): # type: ignore[no-untyped-def]
  277. recorded["reload_force"] = force
  278. return True
  279. monkeypatch.setattr("main.load_all", lambda: None)
  280. monkeypatch.setattr("main.threading.Thread", FakeThread)
  281. monkeypatch.setattr("core.cfmail_provisioner.CfmailProvisioner", FakeProvisioner)
  282. monkeypatch.setattr("core.cfmail.DEFAULT_CFMAIL_MANAGER", FakeManager())
  283. monkeypatch.setattr(loop, "_ensure_cfmail_active_domain_ready", lambda: True)
  284. monkeypatch.setattr(loop, "_write_runtime_state", lambda: None)
  285. monkeypatch.setattr(loop, "_log", lambda message: recorded.setdefault("logs", []).append(message))
  286. loop.start()
  287. assert recorded["thread_started"] is True
  288. assert recorded["proxy_url"] == settings.register_proxy
  289. assert recorded["normalized"] is True
  290. assert recorded["target_count"] == 3
  291. assert recorded["reload_force"] is True
  292. assert any("normalized active domain pool" in line for line in recorded["logs"])
  293. def test_registration_loop_retries_device_id_once_with_new_proxy(monkeypatch) -> None:
  294. settings = _base_settings(register_mail_provider="mailtm")
  295. loop = RegistrationLoop(settings)
  296. loop._providers = ["mailtm"]
  297. recorded: dict[str, object] = {"releases": []}
  298. class FakeLease:
  299. def __init__(self, name: str, local_port: int) -> None:
  300. self.name = name
  301. self.local_port = local_port
  302. self.proxy_url = f"socks5://127.0.0.1:{local_port}"
  303. class FakePool:
  304. def __init__(self) -> None:
  305. self._leases = [FakeLease("tw-bad", 17891), FakeLease("sg-good", 17892)]
  306. def acquire(self, timeout=5.0, preferred_name=None, preferred_regions=()): # type: ignore[no-untyped-def]
  307. recorded.setdefault("timeouts", []).append(timeout)
  308. return self._leases.pop(0)
  309. def release(self, lease, *, success, stage=None): # type: ignore[no-untyped-def]
  310. recorded["releases"].append((lease.name, success, stage))
  311. attempts: list[str] = []
  312. def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
  313. attempts.append(kwargs["proxy"])
  314. if len(attempts) == 1:
  315. return {
  316. "success": False,
  317. "stage": "device_id",
  318. "error_message": "device id acquisition failed",
  319. }
  320. loop._stop_event.set()
  321. return {"success": True, "stage": "completed", "email": "demo@example.com"}
  322. loop._proxy_pool = FakePool()
  323. monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
  324. monkeypatch.setattr(loop, "_sync_cpa_from_success", lambda result, thread_id: (True, "", str(result.get("email") or "")))
  325. loop._worker(thread_id=1, initial_provider="mailtm")
  326. assert attempts == ["socks5://127.0.0.1:17891", "socks5://127.0.0.1:17892"]
  327. assert recorded["releases"] == [
  328. ("tw-bad", False, "device_id"),
  329. ("sg-good", True, "completed"),
  330. ]
  331. snapshot = loop.snapshot()
  332. assert snapshot["total_attempts"] == 1
  333. assert snapshot["total_success"] == 1
  334. assert snapshot["total_failure"] == 0
  335. def test_registration_loop_syncs_cpa_immediately_after_success(monkeypatch, tmp_path: Path) -> None:
  336. settings = _base_settings(
  337. register_mail_provider="mailtm",
  338. )
  339. loop = RegistrationLoop(settings)
  340. loop._providers = ["mailtm"]
  341. pool_file = tmp_path / "fresh@example.com.json"
  342. pool_file.write_text(json.dumps({"email": "fresh@example.com", "access_token": "tok", "account_id": "acct"}), encoding="utf-8")
  343. recorded: dict[str, object] = {}
  344. def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
  345. del kwargs
  346. loop._stop_event.set()
  347. return {
  348. "success": True,
  349. "email": "fresh@example.com",
  350. "pool_file": str(pool_file),
  351. "written_to_pool": True,
  352. }
  353. monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
  354. monkeypatch.setattr("core.registration.get_management_key", lambda: "secret") # type: ignore[no-untyped-def]
  355. monkeypatch.setattr("main.classify_token_file", lambda *args, **kwargs: (_ for _ in ()).throw(AssertionError("readiness probe should not gate direct CPA sync"))) # type: ignore[no-untyped-def]
  356. def fake_upload_to_cpa(token_data, api_url=None, api_key=None, proxy=None): # type: ignore[no-untyped-def]
  357. recorded["token_data"] = token_data
  358. recorded["api_url"] = api_url
  359. recorded["api_key"] = api_key
  360. recorded["proxy"] = proxy
  361. return True, "upload success"
  362. monkeypatch.setattr("platforms.chatgpt.cpa_upload.upload_to_cpa", fake_upload_to_cpa)
  363. loop._worker(thread_id=1, initial_provider="mailtm")
  364. assert recorded["token_data"]["email"] == "fresh@example.com"
  365. assert recorded["api_url"] == "http://127.0.0.1:8317"
  366. assert recorded["api_key"] == "secret"
  367. assert recorded["proxy"] is None
  368. payload = json.loads(pool_file.read_text(encoding="utf-8"))
  369. assert payload["cpa_sync_status"] == "synced"
  370. snapshot = loop.snapshot()
  371. assert snapshot["total_success"] == 1
  372. assert snapshot["total_success_registered"] == 1
  373. assert snapshot["total_cpa_sync_success"] == 1
  374. assert snapshot["total_cpa_sync_failure"] == 0
  375. assert snapshot["registered_success_rate"] == 100.0
  376. assert snapshot["cpa_sync_success_rate"] == 100.0
  377. def test_registration_loop_syncs_cpa_before_breaking_on_target_reached(monkeypatch, tmp_path: Path) -> None:
  378. settings = _base_settings(
  379. register_mail_provider="mailtm",
  380. register_target_count=1,
  381. )
  382. loop = RegistrationLoop(settings)
  383. loop._providers = ["mailtm"]
  384. pool_file = tmp_path / "target@example.com.json"
  385. pool_file.write_text(json.dumps({"email": "target@example.com", "access_token": "tok", "account_id": "acct"}), encoding="utf-8")
  386. recorded: dict[str, object] = {}
  387. def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
  388. del kwargs
  389. return {
  390. "success": True,
  391. "email": "target@example.com",
  392. "pool_file": str(pool_file),
  393. "written_to_pool": True,
  394. }
  395. monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
  396. monkeypatch.setattr("core.registration.get_management_key", lambda: "secret") # type: ignore[no-untyped-def]
  397. monkeypatch.setattr("main.classify_token_file", lambda *args, **kwargs: (_ for _ in ()).throw(AssertionError("readiness probe should not gate direct CPA sync"))) # type: ignore[no-untyped-def]
  398. def fake_upload_to_cpa(token_data, api_url=None, api_key=None, proxy=None): # type: ignore[no-untyped-def]
  399. recorded["token_data"] = token_data
  400. recorded["api_url"] = api_url
  401. recorded["api_key"] = api_key
  402. recorded["proxy"] = proxy
  403. return True, "upload success"
  404. monkeypatch.setattr("platforms.chatgpt.cpa_upload.upload_to_cpa", fake_upload_to_cpa)
  405. loop._worker(thread_id=1, initial_provider="mailtm")
  406. assert loop._target_reached.is_set() is True
  407. assert recorded["token_data"]["email"] == "target@example.com"
  408. payload = json.loads(pool_file.read_text(encoding="utf-8"))
  409. assert payload["cpa_sync_status"] == "synced"
  410. snapshot = loop.snapshot()
  411. assert snapshot["total_attempts"] == 1
  412. assert snapshot["total_success"] == 1
  413. assert snapshot["total_cpa_sync_success"] == 1
  414. assert snapshot["total_cpa_sync_failure"] == 0
  415. def test_registration_loop_marks_pool_backup_when_direct_cpa_sync_fails(monkeypatch, tmp_path: Path) -> None:
  416. settings = _base_settings(
  417. register_mail_provider="mailtm",
  418. )
  419. loop = RegistrationLoop(settings)
  420. loop._providers = ["mailtm"]
  421. pool_file = tmp_path / "fresh@example.com.json"
  422. pool_file.write_text(json.dumps({"email": "fresh@example.com", "access_token": "tok", "account_id": "acct"}), encoding="utf-8")
  423. def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
  424. del kwargs
  425. loop._stop_event.set()
  426. return {
  427. "success": True,
  428. "email": "fresh@example.com",
  429. "pool_file": str(pool_file),
  430. "written_to_pool": True,
  431. }
  432. monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
  433. monkeypatch.setattr("core.registration.get_management_key", lambda: "secret") # type: ignore[no-untyped-def]
  434. monkeypatch.setattr("main.classify_token_file", lambda *args, **kwargs: (_ for _ in ()).throw(AssertionError("readiness probe should not gate direct CPA sync"))) # type: ignore[no-untyped-def]
  435. monkeypatch.setattr("platforms.chatgpt.cpa_upload.upload_to_cpa", lambda *args, **kwargs: (False, "unexpected EOF")) # type: ignore[no-untyped-def]
  436. loop._worker(thread_id=1, initial_provider="mailtm")
  437. payload = json.loads(pool_file.read_text(encoding="utf-8"))
  438. assert payload["backup_written"] is True
  439. assert payload["cpa_sync_status"] == "failed"
  440. assert payload["last_cpa_sync_error"] == "unexpected EOF"
  441. snapshot = loop.snapshot()
  442. assert snapshot["total_success"] == 0
  443. assert snapshot["total_failure"] == 1
  444. assert snapshot["total_cpa_sync_success"] == 0
  445. assert snapshot["total_cpa_sync_failure"] == 1
  446. assert snapshot["failure_by_stage"]["cpa_sync"] == 1
  447. assert snapshot["failure_signals"]["cpa_sync_failed"] == 1
  448. def test_registration_loop_holds_add_phone_gated_success_in_warmup(monkeypatch, tmp_path: Path) -> None:
  449. settings = _base_settings(
  450. register_mail_provider="cfmail",
  451. )
  452. loop = RegistrationLoop(settings)
  453. loop._providers = ["cfmail"]
  454. pool_file = tmp_path / "fresh@example.com.json"
  455. pool_file.write_text(json.dumps({"email": "fresh@example.com", "access_token": "tok", "account_id": "acct"}), encoding="utf-8")
  456. uploaded: list[str] = []
  457. def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
  458. del kwargs
  459. loop._stop_event.set()
  460. return {
  461. "success": True,
  462. "email": "fresh@example.com",
  463. "pool_file": str(pool_file),
  464. "written_to_pool": True,
  465. "metadata": {
  466. "mail_provider": "cfmail",
  467. "email_domain": "demo.example.test",
  468. "post_create_gate": "add_phone",
  469. },
  470. }
  471. monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
  472. monkeypatch.setattr("core.registration.get_management_key", lambda: "secret") # type: ignore[no-untyped-def]
  473. monkeypatch.setattr("main.classify_token_file", lambda *args, **kwargs: (_ for _ in ()).throw(AssertionError("warmup gating should not use legacy readiness probe"))) # type: ignore[no-untyped-def]
  474. monkeypatch.setattr(
  475. "platforms.chatgpt.cpa_upload.upload_to_cpa",
  476. lambda token_data, api_url=None, api_key=None, proxy=None: uploaded.append(token_data["email"]) or (True, "ok"),
  477. )
  478. loop._worker(thread_id=1, initial_provider="cfmail")
  479. assert uploaded == []
  480. payload = json.loads(pool_file.read_text(encoding="utf-8"))
  481. assert payload["warmup_required"] is True
  482. assert payload["warmup_state"] == "pending"
  483. assert payload["cpa_sync_status"] == "warmup_pending"
  484. snapshot = loop.snapshot()
  485. assert snapshot["total_success"] == 0
  486. assert snapshot["total_success_registered"] == 0
  487. assert snapshot["total_cpa_sync_success"] == 0
  488. assert snapshot["total_cpa_sync_failure"] == 0
  489. assert snapshot["total_warmup_pending"] == 1
  490. def test_registration_loop_marks_add_phone_backup_as_warmup_pending(monkeypatch, tmp_path: Path) -> None:
  491. settings = _base_settings(
  492. register_mail_provider="cfmail",
  493. )
  494. loop = RegistrationLoop(settings)
  495. loop._providers = ["cfmail"]
  496. pool_file = tmp_path / "fresh@example.com.json"
  497. pool_file.write_text(
  498. json.dumps({"email": "fresh@example.com", "access_token": "tok", "account_id": "acct"}),
  499. encoding="utf-8",
  500. )
  501. def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
  502. del kwargs
  503. loop._stop_event.set()
  504. return {
  505. "success": True,
  506. "email": "fresh@example.com",
  507. "pool_file": str(pool_file),
  508. "written_to_pool": True,
  509. "metadata": {
  510. "mail_provider": "cfmail",
  511. "email_domain": "demo.example.test",
  512. "post_create_gate": "add_phone",
  513. },
  514. }
  515. monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
  516. monkeypatch.setattr("core.registration.get_management_key", lambda: "secret") # type: ignore[no-untyped-def]
  517. monkeypatch.setattr("main.classify_token_file", lambda *args, **kwargs: (_ for _ in ()).throw(AssertionError("readiness probe should not gate direct CPA sync"))) # type: ignore[no-untyped-def]
  518. monkeypatch.setattr("platforms.chatgpt.cpa_upload.upload_to_cpa", lambda *args, **kwargs: (_ for _ in ()).throw(AssertionError("warmup pending account should not upload immediately"))) # type: ignore[no-untyped-def]
  519. loop._worker(thread_id=1, initial_provider="cfmail")
  520. payload = json.loads(pool_file.read_text(encoding="utf-8"))
  521. assert payload["backup_written"] is True
  522. assert payload["cpa_sync_status"] == "warmup_pending"
  523. assert payload["last_cpa_sync_error"] == "warmup pending"
  524. def test_registration_loop_enqueue_pending_token_preserves_proxy_provenance() -> None:
  525. loop = RegistrationLoop(_base_settings())
  526. loop._enqueue_pending_token(
  527. {
  528. "metadata": {
  529. "deferred_credentials": {
  530. "email": "fresh@example.com",
  531. "password": "pw-secret",
  532. "registration_proxy_key": "台湾♣备用-1",
  533. "registration_proxy_region": "tw",
  534. "registration_proxy_url": "socks5://127.0.0.1:17891",
  535. "registration_fingerprint_profile": "chrome120_win",
  536. "cfmail_profile_name": "cfmail-tw",
  537. "add_phone_trace_path": "/tmp/add-phone-trace.json",
  538. }
  539. }
  540. },
  541. thread_id=1,
  542. )
  543. assert len(loop._pending_token_queue) == 1
  544. entry = loop._pending_token_queue[0]
  545. assert entry["registration_proxy_key"] == "台湾♣备用-1"
  546. assert entry["registration_proxy_region"] == "tw"
  547. assert entry["registration_proxy_url"] == "socks5://127.0.0.1:17891"
  548. assert entry["registration_fingerprint_profile"] == "chrome120_win"
  549. assert entry["cfmail_profile_name"] == "cfmail-tw"
  550. assert entry["add_phone_trace_path"] == "/tmp/add-phone-trace.json"
  551. def test_registration_loop_retry_pending_token_reuses_registration_proxy(monkeypatch, tmp_path: Path) -> None:
  552. settings = _base_settings(pool_dir=tmp_path)
  553. loop = RegistrationLoop(settings)
  554. recorded: dict[str, object] = {}
  555. class FakeLease:
  556. name = "台湾♣备用-1"
  557. proxy_url = "socks5://127.0.0.1:17891"
  558. class FakePool:
  559. def acquire(self, timeout=5.0, preferred_name=None, preferred_regions=()): # type: ignore[no-untyped-def]
  560. recorded["preferred_name"] = preferred_name
  561. recorded["preferred_regions"] = tuple(preferred_regions)
  562. return FakeLease()
  563. def release(self, lease, *, success, stage=None): # type: ignore[no-untyped-def]
  564. recorded["released"] = (lease.proxy_url, success, stage)
  565. class FakeMailbox:
  566. def __init__(self, manager=None): # type: ignore[no-untyped-def]
  567. self.manager = manager
  568. class FakeAdapter:
  569. def __init__(self, mailbox): # type: ignore[no-untyped-def]
  570. self.mailbox = mailbox
  571. self._account = None
  572. class FakeEngine:
  573. def __init__(self, email_service, proxy_url=None): # type: ignore[no-untyped-def]
  574. recorded["engine_proxy_url"] = proxy_url
  575. self.email_service = email_service
  576. self.email = ""
  577. self.password = ""
  578. def _login_for_token(self): # type: ignore[no-untyped-def]
  579. return {
  580. "access_token": "access-token",
  581. "refresh_token": "refresh-token",
  582. "id_token": "",
  583. "account_id": "acct-123",
  584. "last_refresh": "2026-03-31T00:00:00Z",
  585. "expired": "2026-04-01T00:00:00Z",
  586. }
  587. def fake_write_token_record(token_data, pool_dir): # type: ignore[no-untyped-def]
  588. recorded["token_data"] = dict(token_data)
  589. path = Path(pool_dir) / "fresh@example.com.json"
  590. path.write_text(json.dumps(token_data), encoding="utf-8")
  591. return path
  592. loop._proxy_pool = FakePool()
  593. monkeypatch.setattr("core.registration.CfMailMailbox", FakeMailbox, raising=False)
  594. monkeypatch.setattr("core.registration.MailboxEmailServiceAdapter", FakeAdapter, raising=False)
  595. monkeypatch.setattr("platforms.chatgpt.register.RegistrationEngine", FakeEngine)
  596. monkeypatch.setattr("platforms.chatgpt.pool.write_token_record", fake_write_token_record)
  597. monkeypatch.setattr(loop, "_sync_cpa_from_success", lambda result, thread_id: (True, "", "fresh@example.com"))
  598. entry = {
  599. "email": "fresh@example.com",
  600. "password": "pw-secret",
  601. "mailbox_jwt": "",
  602. "mailbox_extra": {},
  603. "registration_proxy_key": "台湾♣备用-1",
  604. "registration_proxy_region": "tw",
  605. "registration_proxy_url": "socks5://127.0.0.1:17891",
  606. "registration_fingerprint_profile": "chrome120_win",
  607. "cfmail_profile_name": "cfmail-tw",
  608. "retry_count": 0,
  609. }
  610. loop._retry_pending_token(entry)
  611. assert recorded["preferred_name"] == "台湾♣备用-1"
  612. assert recorded["preferred_regions"] == ("tw",)
  613. assert recorded["engine_proxy_url"] == "socks5://127.0.0.1:17891"
  614. assert recorded["released"] == ("socks5://127.0.0.1:17891", True, "deferred_retry")
  615. token_data = recorded["token_data"]
  616. assert token_data["registration_proxy_key"] == "台湾♣备用-1"
  617. assert token_data["registration_proxy_region"] == "tw"
  618. assert token_data["registration_proxy_url"] == "socks5://127.0.0.1:17891"
  619. assert token_data["registration_fingerprint_profile"] == "chrome120_win"
  620. assert token_data["registration_cfmail_profile_name"] == "cfmail-tw"
  621. def test_registration_loop_dumps_engine_logs_on_failure(monkeypatch) -> None:
  622. settings = _base_settings(register_max_consecutive_failures=1, register_mail_provider="mailtm")
  623. loop = RegistrationLoop(settings)
  624. loop._providers = ["mailtm"]
  625. logged: list[str] = []
  626. original_log = loop._log
  627. def capture_log(msg: str) -> None:
  628. logged.append(msg)
  629. loop._log = capture_log
  630. def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
  631. loop._stop_event.set()
  632. return {
  633. "success": False,
  634. "stage": "signup",
  635. "error_message": "HTTP 403: access denied",
  636. "logs": [
  637. "[09:35:20] check_ip_location: JP",
  638. "[09:35:21] created mailbox: test@example.com",
  639. "[09:35:22] signup form status: 403",
  640. ],
  641. }
  642. monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
  643. loop._worker(thread_id=1, initial_provider="mailtm")
  644. # Verify stage appears in the failure line
  645. fail_lines = [line for line in logged if "failed" in line and "stage=" in line]
  646. assert len(fail_lines) == 1
  647. assert "[stage=signup]" in fail_lines[0]
  648. assert "HTTP 403: access denied" in fail_lines[0]
  649. # Verify engine logs are dumped with ↳ prefix
  650. engine_lines = [line for line in logged if "\u21b3" in line]
  651. assert len(engine_lines) == 3
  652. assert "signup form status: 403" in engine_lines[2]
  653. snapshot = loop.snapshot()
  654. assert snapshot["failure_by_stage"]["signup"] == 1
  655. assert snapshot["recent_failure_hotspots"] == [{"key": "signup", "stage": "signup", "count": 1}]
  656. def test_registration_loop_rotates_cfmail_domain_after_blacklist_threshold(monkeypatch) -> None:
  657. settings = _base_settings(register_mail_provider="cfmail", register_max_consecutive_failures=5)
  658. loop = RegistrationLoop(settings)
  659. loop._providers = ["cfmail"]
  660. loop._cfmail_tracker = DomainHealthTracker(window_size=2, blacklist_threshold=2, rotation_cooldown_seconds=1)
  661. monkeypatch.setattr(loop, "_ensure_cfmail_domain_pool_target", lambda **_: None)
  662. rotation_calls: list[str] = []
  663. class FakeProvisioner:
  664. def rotate_active_domain(self): # type: ignore[no-untyped-def]
  665. rotation_calls.append("rotate")
  666. return ProvisionResult(
  667. success=True,
  668. step="completed",
  669. old_domain="nova.example.test",
  670. new_domain="auto0322.example.test",
  671. )
  672. def provision_additional_domain(self): # type: ignore[no-untyped-def]
  673. return ProvisionResult(success=False, step="provision_additional_domain", error="not needed")
  674. loop._cfmail_provisioner = FakeProvisioner()
  675. attempts = {"count": 0}
  676. def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
  677. del kwargs
  678. attempts["count"] += 1
  679. if attempts["count"] >= 2:
  680. loop._stop_event.set()
  681. return {
  682. "success": False,
  683. "stage": "create_account",
  684. "error_message": "create account failed",
  685. "metadata": {
  686. "mail_provider": "cfmail",
  687. "email_domain": "nova.example.test",
  688. "create_account_error_code": "registration_disallowed",
  689. "create_account_error_message": "blocked",
  690. },
  691. }
  692. monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
  693. loop._worker(thread_id=1, initial_provider="cfmail")
  694. assert rotation_calls == ["rotate"]
  695. assert loop.snapshot()["cfmail_rotation"]["last_new_domain"] == "auto0322.example.test"
  696. def test_registration_loop_does_not_rotate_cfmail_on_mailbox_failure(monkeypatch) -> None:
  697. settings = _base_settings(register_mail_provider="cfmail", register_max_consecutive_failures=2)
  698. loop = RegistrationLoop(settings)
  699. loop._providers = ["cfmail"]
  700. loop._cfmail_tracker = DomainHealthTracker(window_size=2, blacklist_threshold=2, rotation_cooldown_seconds=1)
  701. monkeypatch.setattr(loop, "_ensure_cfmail_domain_pool_target", lambda **_: None)
  702. class FakeProvisioner:
  703. def rotate_active_domain(self): # type: ignore[no-untyped-def]
  704. raise AssertionError("rotation should not be called")
  705. def provision_additional_domain(self): # type: ignore[no-untyped-def]
  706. return ProvisionResult(success=False, step="provision_additional_domain", error="not needed")
  707. loop._cfmail_provisioner = FakeProvisioner()
  708. def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
  709. del kwargs
  710. loop._stop_event.set()
  711. return {
  712. "success": False,
  713. "stage": "mailbox",
  714. "error_message": "create email failed",
  715. "metadata": {"mail_provider": "cfmail"},
  716. }
  717. monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
  718. loop._worker(thread_id=1, initial_provider="cfmail")
  719. assert loop.snapshot()["cfmail_rotation"]["last_new_domain"] == ""
  720. def test_registration_loop_rotates_cfmail_on_invalid_domain_mailbox_failure() -> None:
  721. settings = _base_settings(register_mail_provider="cfmail")
  722. loop = RegistrationLoop(settings)
  723. loop._providers = ["cfmail"]
  724. loop._cfmail_tracker = DomainHealthTracker(window_size=2, blacklist_threshold=2, rotation_cooldown_seconds=1)
  725. class FakeProvisioner:
  726. def rotate_active_domain(self): # type: ignore[no-untyped-def]
  727. return ProvisionResult(
  728. success=True,
  729. step="rotate",
  730. old_domain="bad.example.test",
  731. new_domain="auto-new.example.test",
  732. )
  733. reload_calls: list[bool] = []
  734. class FakeManager:
  735. def reload_if_needed(self, force=False): # type: ignore[no-untyped-def]
  736. reload_calls.append(bool(force))
  737. return True
  738. loop._cfmail_manager = FakeManager()
  739. loop._cfmail_provisioner = FakeProvisioner()
  740. loop._cfmail_wait_otp_cooldown_seconds = 60
  741. loop._cfmail_add_phone_cooldown_seconds = 60
  742. result = {
  743. "success": False,
  744. "stage": "mailbox",
  745. "error_message": "create email failed",
  746. "logs": ["[10:00:00] create_email failed: 创建邮箱地址失败: 无效的域名"],
  747. "metadata": {
  748. "mail_provider": "cfmail",
  749. "email_domain": "bad.example.test",
  750. },
  751. }
  752. assert loop._force_rotate_cfmail_for_invalid_mailbox(thread_id=1, result=result) is True
  753. snapshot = loop.snapshot()
  754. assert snapshot["cfmail_rotation"]["last_new_domain"] == "auto-new.example.test"
  755. assert snapshot["cfmail_wait_otp_stoploss"]["active_domain"] == "auto-new.example.test"
  756. assert reload_calls == [True]
  757. def test_registration_loop_startup_reports_active_domain_pool_ready(monkeypatch) -> None:
  758. settings = _base_settings(register_mail_provider="cfmail")
  759. loop = RegistrationLoop(settings)
  760. logs: list[str] = []
  761. class FakeManager:
  762. accounts = [
  763. type("Account", (), {"name": "cfmail-a", "email_domain": "a.example.test"})(),
  764. type("Account", (), {"name": "cfmail-b", "email_domain": "b.example.test"})(),
  765. ]
  766. def skip_remaining_seconds(self, account_name): # type: ignore[no-untyped-def]
  767. del account_name
  768. return 0
  769. loop._cfmail_manager = FakeManager()
  770. loop._cfmail_provisioner = object()
  771. monkeypatch.setattr(loop, "_log", logs.append)
  772. assert loop._ensure_cfmail_active_domain_ready() is True
  773. assert any("active-domain pool ready" in line for line in logs)
  774. def test_registration_loop_forces_cfmail_rotation_when_all_accounts_in_cooldown(monkeypatch) -> None:
  775. settings = _base_settings(register_mail_provider="cfmail", register_max_consecutive_failures=2)
  776. loop = RegistrationLoop(settings)
  777. loop._providers = ["cfmail"]
  778. loop._cfmail_tracker = DomainHealthTracker(window_size=2, blacklist_threshold=2, rotation_cooldown_seconds=1)
  779. run_calls: list[str] = []
  780. rotation_calls: list[str] = []
  781. class FakeManager:
  782. def reload_if_needed(self) -> bool:
  783. return False
  784. def select_account(self, profile_name=None): # type: ignore[no-untyped-def]
  785. del profile_name
  786. return None
  787. class FakeProvisioner:
  788. def rotate_active_domain(self): # type: ignore[no-untyped-def]
  789. rotation_calls.append("rotate")
  790. loop._stop_event.set()
  791. return ProvisionResult(
  792. success=True,
  793. step="completed",
  794. old_domain="auto-old.example.test",
  795. new_domain="auto-new.example.test",
  796. )
  797. def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
  798. del kwargs
  799. loop._stop_event.set()
  800. run_calls.append("called")
  801. raise AssertionError("registration should not run while cfmail is fully unavailable")
  802. loop._cfmail_manager = FakeManager()
  803. loop._cfmail_provisioner = FakeProvisioner()
  804. monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
  805. loop._worker(thread_id=1, initial_provider="cfmail")
  806. assert rotation_calls == ["rotate"]
  807. assert run_calls == []
  808. assert loop.snapshot()["cfmail_rotation"]["last_new_domain"] == "auto-new.example.test"
  809. def test_registration_loop_does_not_penalize_proxy_for_blacklist_failure(monkeypatch) -> None:
  810. settings = _base_settings(register_mail_provider="cfmail")
  811. loop = RegistrationLoop(settings)
  812. loop._providers = ["cfmail"]
  813. recorded: dict[str, object] = {}
  814. class FakeLease:
  815. name = "sg-node"
  816. local_port = 17891
  817. proxy_url = "socks5://127.0.0.1:17891"
  818. class FakePool:
  819. def acquire(self, timeout=5.0, preferred_name=None, preferred_regions=()): # type: ignore[no-untyped-def]
  820. return FakeLease()
  821. def release(self, lease, *, success, stage=None): # type: ignore[no-untyped-def]
  822. recorded["released"] = (success, stage)
  823. def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
  824. del kwargs
  825. loop._stop_event.set()
  826. return {
  827. "success": False,
  828. "stage": "create_account",
  829. "error_message": "create account failed",
  830. "metadata": {
  831. "mail_provider": "cfmail",
  832. "email_domain": "nova.example.test",
  833. "create_account_error_code": "registration_disallowed",
  834. },
  835. }
  836. loop._proxy_pool = FakePool()
  837. monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
  838. loop._worker(thread_id=1, initial_provider="cfmail")
  839. assert recorded["released"] == (None, "create_account")
  840. def test_registration_loop_records_add_phone_gate_signal(monkeypatch) -> None:
  841. settings = _base_settings(register_mail_provider="cfmail")
  842. loop = RegistrationLoop(settings)
  843. loop._providers = ["cfmail"]
  844. def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
  845. del kwargs
  846. loop._stop_event.set()
  847. return {
  848. "success": False,
  849. "stage": "add_phone_gate",
  850. "error_message": "post-create flow requires phone gate",
  851. "metadata": {
  852. "mail_provider": "cfmail",
  853. "email_domain": "demo.example.test",
  854. "post_create_gate": "add_phone",
  855. },
  856. }
  857. monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
  858. loop._worker(thread_id=1, initial_provider="cfmail")
  859. snapshot = loop.snapshot()
  860. assert snapshot["failure_by_stage"]["add_phone_gate"] == 1
  861. assert snapshot["failure_signals"]["add_phone_gate"] == 1
  862. assert snapshot["recent_failure_hotspots"] == [{"key": "add_phone_gate", "stage": "add_phone_gate", "count": 1}]
  863. assert snapshot["recent_attempts"][0]["post_create_gate"] == "add_phone"
  864. def test_registration_loop_records_signup_invalid_auth_step_signal() -> None:
  865. loop = RegistrationLoop(_base_settings(register_mail_provider="cfmail"))
  866. loop._record_attempt(
  867. success=False,
  868. stage="signup",
  869. error_message='HTTP 400: {"error":{"code":"invalid_auth_step"}}',
  870. metadata={
  871. "mail_provider": "cfmail",
  872. "email_domain": "demo.example.test",
  873. "signup_error_code": "invalid_auth_step",
  874. "signup_http_status": 400,
  875. },
  876. proxy_key="台湾原生-01",
  877. email="demo@example.test",
  878. )
  879. snapshot = loop.snapshot()
  880. assert snapshot["failure_signals"]["invalid_auth_step"] == 1
  881. attempt = snapshot["recent_attempts"][0]
  882. assert attempt["signup_error_code"] == "invalid_auth_step"
  883. assert attempt["signup_http_status"] == 400
  884. def test_registration_loop_classifies_mailbox_transport_failures() -> None:
  885. loop = RegistrationLoop(_base_settings(register_mail_provider="cfmail"))
  886. loop._record_attempt(
  887. success=False,
  888. stage="mailbox",
  889. error_message="create email failed",
  890. metadata={
  891. "mail_provider": "cfmail",
  892. "email_domain": "demo.example.test",
  893. "mailbox_error_kind": "transport_error",
  894. "mailbox_error_stage": "create_email",
  895. },
  896. proxy_key="新加坡原生-02",
  897. )
  898. snapshot = loop.snapshot()
  899. assert snapshot["failure_signals"]["mailbox_create_transport_error"] == 1
  900. attempt = snapshot["recent_attempts"][0]
  901. assert attempt["mailbox_error_kind"] == "transport_error"
  902. assert attempt["mailbox_error_stage"] == "create_email"
  903. def test_registration_loop_classifies_user_already_exists_as_mailbox_reused(monkeypatch) -> None:
  904. settings = _base_settings(register_mail_provider="cfmail")
  905. loop = RegistrationLoop(settings)
  906. loop._providers = ["cfmail"]
  907. def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
  908. del kwargs
  909. loop._stop_event.set()
  910. return {
  911. "success": False,
  912. "stage": "create_account",
  913. "error_message": "create account failed",
  914. "email": "dup@example.com",
  915. "metadata": {
  916. "mail_provider": "cfmail",
  917. "email_domain": "demo.example.test",
  918. "create_account_error_code": "user_already_exists",
  919. },
  920. }
  921. monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
  922. loop._worker(thread_id=1, initial_provider="cfmail")
  923. snapshot = loop.snapshot()
  924. assert snapshot["failure_by_stage"]["create_account"] == 1
  925. assert snapshot["failure_signals"]["mailbox_reused"] == 1
  926. assert snapshot["recent_failure_hotspots"] == [{"key": "mailbox_reused", "stage": "create_account", "count": 1}]
  927. def test_registration_loop_activates_add_phone_stoploss(monkeypatch) -> None:
  928. settings = _base_settings(register_mail_provider="cfmail")
  929. loop = RegistrationLoop(settings)
  930. loop._providers = ["cfmail"]
  931. monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "demo.example.test")
  932. loop._cfmail_add_phone_window = 2
  933. loop._cfmail_add_phone_threshold = 2
  934. loop._cfmail_add_phone_max_successes = 0
  935. loop._cfmail_add_phone_cooldown_seconds = 60
  936. attempts = {"count": 0}
  937. def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
  938. del kwargs
  939. attempts["count"] += 1
  940. if attempts["count"] >= 2:
  941. loop._stop_event.set()
  942. return {
  943. "success": False,
  944. "stage": "add_phone_gate",
  945. "error_message": "post-create flow requires phone gate",
  946. "metadata": {
  947. "mail_provider": "cfmail",
  948. "email_domain": "demo.example.test",
  949. "post_create_gate": "add_phone",
  950. },
  951. }
  952. monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
  953. loop._worker(thread_id=1, initial_provider="cfmail")
  954. snapshot = loop.snapshot()
  955. stoploss = snapshot["cfmail_add_phone_stoploss"]
  956. assert stoploss["active_domain"] == "demo.example.test"
  957. assert stoploss["in_cooldown"] is True
  958. assert stoploss["last_add_phone_failures"] == 2
  959. assert stoploss["last_successes"] == 0
  960. def test_registration_loop_snapshot_exposes_active_domain_only_attempts() -> None:
  961. settings = _base_settings(register_mail_provider="cfmail")
  962. loop = RegistrationLoop(settings)
  963. loop._recent_attempts.extend(
  964. [
  965. {
  966. "timestamp": "2026-03-24T18:00:00",
  967. "success": False,
  968. "stage": "create_account",
  969. "signal": "registration_disallowed",
  970. "error_message": "failed",
  971. "email_domain": "old.example.test",
  972. "proxy_key": "old-node",
  973. "email": "old@old.example.test",
  974. },
  975. {
  976. "timestamp": "2026-03-24T18:00:10",
  977. "success": True,
  978. "stage": "completed",
  979. "signal": "",
  980. "error_message": "",
  981. "email_domain": "new.example.test",
  982. "proxy_key": "new-node-1",
  983. "email": "ok@new.example.test",
  984. },
  985. {
  986. "timestamp": "2026-03-24T18:00:20",
  987. "success": False,
  988. "stage": "create_account",
  989. "signal": "registration_disallowed",
  990. "error_message": "failed",
  991. "email_domain": "new.example.test",
  992. "proxy_key": "new-node-2",
  993. "email": "bad@new.example.test",
  994. },
  995. ]
  996. )
  997. class FakeTracker:
  998. def snapshot(self) -> dict[str, object]:
  999. return {
  1000. "active_domain": "new.example.test",
  1001. "last_new_domain": "new.example.test",
  1002. "last_blacklisted_domain": "old.example.test",
  1003. }
  1004. loop._cfmail_tracker = FakeTracker()
  1005. snapshot = loop.snapshot()
  1006. assert [item["email_domain"] for item in snapshot["active_domain_recent_attempts"]] == [
  1007. "new.example.test",
  1008. "new.example.test",
  1009. ]
  1010. assert snapshot["active_domain_failure_by_stage"] == {"create_account": 1}
  1011. assert snapshot["active_domain_failure_signals"] == {"registration_disallowed": 1}
  1012. assert snapshot["active_domain_recent_failure_hotspots"] == [
  1013. {"key": "registration_disallowed", "stage": "create_account", "count": 1}
  1014. ]
  1015. def test_registration_loop_snapshot_infers_active_domain_from_recent_attempts() -> None:
  1016. settings = _base_settings(register_mail_provider="cfmail")
  1017. loop = RegistrationLoop(settings)
  1018. loop._recent_attempts.extend(
  1019. [
  1020. {
  1021. "timestamp": "2026-03-24T18:10:00",
  1022. "success": False,
  1023. "stage": "create_account",
  1024. "signal": "registration_disallowed",
  1025. "error_message": "failed",
  1026. "email_domain": "old.example.test",
  1027. "proxy_key": "old-node",
  1028. "email": "old@old.example.test",
  1029. },
  1030. {
  1031. "timestamp": "2026-03-24T18:10:10",
  1032. "success": True,
  1033. "stage": "completed",
  1034. "signal": "",
  1035. "error_message": "",
  1036. "email_domain": "new.example.test",
  1037. "proxy_key": "new-node",
  1038. "email": "ok@new.example.test",
  1039. },
  1040. ]
  1041. )
  1042. snapshot = loop.snapshot()
  1043. assert [item["email_domain"] for item in snapshot["active_domain_recent_attempts"]] == [
  1044. "new.example.test"
  1045. ]
  1046. def test_registration_loop_disables_add_phone_cooldown_when_configured_zero(monkeypatch) -> None:
  1047. monkeypatch.setenv("ZHUCE6_CFMAIL_ADD_PHONE_COOLDOWN_SECONDS", "0")
  1048. settings = _base_settings(register_mail_provider="cfmail")
  1049. loop = RegistrationLoop(settings)
  1050. loop._cfmail_add_phone_window = 2
  1051. loop._cfmail_add_phone_threshold = 2
  1052. loop._cfmail_add_phone_max_successes = 0
  1053. result = {
  1054. "success": False,
  1055. "stage": "add_phone_gate",
  1056. "error_message": "post-create flow requires phone gate",
  1057. "metadata": {
  1058. "mail_provider": "cfmail",
  1059. "email_domain": "demo.example.test",
  1060. "post_create_gate": "add_phone",
  1061. },
  1062. }
  1063. loop._update_cfmail_add_phone_stoploss(result)
  1064. loop._update_cfmail_add_phone_stoploss(result)
  1065. stoploss = loop.snapshot()["cfmail_add_phone_stoploss"]
  1066. assert stoploss["in_cooldown"] is False
  1067. assert stoploss["cooldown_remaining_seconds"] == 0
  1068. def test_registration_loop_add_phone_stoploss_retires_triggered_domain_without_blocking_other_domains(monkeypatch) -> None:
  1069. settings = _base_settings(register_mail_provider="cfmail")
  1070. loop = RegistrationLoop(settings)
  1071. loop._cfmail_add_phone_window = 2
  1072. loop._cfmail_add_phone_threshold = 2
  1073. loop._cfmail_add_phone_max_successes = 0
  1074. loop._cfmail_add_phone_cooldown_seconds = 60
  1075. monkeypatch.setattr(
  1076. loop,
  1077. "_current_cfmail_active_accounts",
  1078. lambda: [
  1079. {"name": "cfmail-tw", "domain": "tw.example.test"},
  1080. {"name": "cfmail-sg", "domain": "sg.example.test"},
  1081. ],
  1082. )
  1083. retired: list[str] = []
  1084. scheduled: list[str] = []
  1085. class FakeProvisioner:
  1086. def retire_domain(self, domain: str) -> ProvisionResult:
  1087. retired.append(domain)
  1088. return ProvisionResult(success=True, step="retire_domain", old_domain=domain)
  1089. loop._cfmail_provisioner = FakeProvisioner() # type: ignore[assignment]
  1090. loop._cfmail_tracker = DomainHealthTracker(window_size=2, blacklist_threshold=2, rotation_cooldown_seconds=1)
  1091. monkeypatch.setattr(loop, "_reload_cfmail_manager_after_rotation", lambda: None)
  1092. monkeypatch.setattr(
  1093. loop,
  1094. "_schedule_cfmail_domain_pool_replenish",
  1095. lambda *, trigger_thread_id, reason: scheduled.append(f"{trigger_thread_id}:{reason}"),
  1096. raising=False,
  1097. )
  1098. result = {
  1099. "success": False,
  1100. "stage": "add_phone_gate",
  1101. "error_message": "post-create flow requires phone gate",
  1102. "metadata": {
  1103. "mail_provider": "cfmail",
  1104. "email_domain": "sg.example.test",
  1105. "post_create_gate": "add_phone",
  1106. },
  1107. }
  1108. loop._update_cfmail_add_phone_stoploss(result)
  1109. loop._update_cfmail_add_phone_stoploss(result)
  1110. assert loop._wait_if_cfmail_add_phone_stopped(thread_id=3, provider="cfmail") is False
  1111. assert retired == ["sg.example.test"]
  1112. assert scheduled == ["3:add_phone stoploss"]
  1113. def test_registration_loop_activates_wait_otp_stoploss_for_no_message_timeouts(monkeypatch) -> None:
  1114. settings = _base_settings(register_mail_provider="cfmail")
  1115. loop = RegistrationLoop(settings)
  1116. monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "demo.example.test")
  1117. loop._cfmail_wait_otp_window = 2
  1118. loop._cfmail_wait_otp_threshold = 2
  1119. loop._cfmail_wait_otp_cooldown_seconds = 60
  1120. result = {
  1121. "success": False,
  1122. "stage": "wait_otp",
  1123. "error_message": "otp retrieval failed",
  1124. "metadata": {
  1125. "mail_provider": "cfmail",
  1126. "email_domain": "demo.example.test",
  1127. "otp_wait_failure_reason": "mailbox_timeout_no_message",
  1128. "otp_mailbox_message_scan_count": 0,
  1129. },
  1130. }
  1131. loop._update_cfmail_wait_otp_stoploss(result)
  1132. loop._update_cfmail_wait_otp_stoploss(result)
  1133. stoploss = loop.snapshot()["cfmail_wait_otp_stoploss"]
  1134. assert stoploss["active_domain"] == "demo.example.test"
  1135. assert stoploss["in_cooldown"] is True
  1136. assert stoploss["last_no_message_timeouts"] == 2
  1137. def test_registration_loop_wait_otp_stoploss_retires_triggered_domain_without_blocking_other_domains(monkeypatch) -> None:
  1138. settings = _base_settings(register_mail_provider="cfmail")
  1139. loop = RegistrationLoop(settings)
  1140. loop._cfmail_wait_otp_window = 2
  1141. loop._cfmail_wait_otp_threshold = 2
  1142. loop._cfmail_wait_otp_cooldown_seconds = 60
  1143. monkeypatch.setattr(
  1144. loop,
  1145. "_current_cfmail_active_accounts",
  1146. lambda: [
  1147. {"name": "cfmail-tw", "domain": "tw.example.test"},
  1148. {"name": "cfmail-sg", "domain": "sg.example.test"},
  1149. ],
  1150. )
  1151. retired: list[str] = []
  1152. scheduled: list[str] = []
  1153. class FakeProvisioner:
  1154. def retire_domain(self, domain: str) -> ProvisionResult:
  1155. retired.append(domain)
  1156. return ProvisionResult(success=True, step="retire_domain", old_domain=domain)
  1157. loop._cfmail_provisioner = FakeProvisioner() # type: ignore[assignment]
  1158. loop._cfmail_tracker = DomainHealthTracker(window_size=2, blacklist_threshold=2, rotation_cooldown_seconds=1)
  1159. monkeypatch.setattr(loop, "_reload_cfmail_manager_after_rotation", lambda: None)
  1160. monkeypatch.setattr(
  1161. loop,
  1162. "_schedule_cfmail_domain_pool_replenish",
  1163. lambda *, trigger_thread_id, reason: scheduled.append(f"{trigger_thread_id}:{reason}"),
  1164. raising=False,
  1165. )
  1166. result = {
  1167. "success": False,
  1168. "stage": "wait_otp",
  1169. "error_message": "otp retrieval failed",
  1170. "metadata": {
  1171. "mail_provider": "cfmail",
  1172. "email_domain": "sg.example.test",
  1173. "otp_wait_failure_reason": "mailbox_timeout_no_message",
  1174. "otp_mailbox_message_scan_count": 0,
  1175. },
  1176. }
  1177. loop._update_cfmail_wait_otp_stoploss(result)
  1178. loop._update_cfmail_wait_otp_stoploss(result)
  1179. assert loop._wait_if_cfmail_wait_otp_stopped(thread_id=4, provider="cfmail") is False
  1180. assert retired == ["sg.example.test"]
  1181. assert scheduled == ["4:wait_otp stoploss"]
  1182. def test_registration_loop_replenish_worker_retries_transient_provision_failure(monkeypatch) -> None:
  1183. settings = _base_settings(register_mail_provider="cfmail")
  1184. loop = RegistrationLoop(settings)
  1185. loop._cfmail_active_domain_count = 3
  1186. monkeypatch.setattr("core.registration.time.sleep", lambda _: None)
  1187. active_accounts = [
  1188. {"name": "cfmail-tw", "domain": "tw.example.test"},
  1189. {"name": "cfmail-sg", "domain": "sg.example.test"},
  1190. ]
  1191. monkeypatch.setattr(loop, "_current_cfmail_active_accounts", lambda: list(active_accounts))
  1192. monkeypatch.setattr(loop, "_reload_cfmail_manager_after_rotation", lambda: None)
  1193. attempts = {"count": 0}
  1194. class FakeProvisioner:
  1195. def provision_additional_domain(self) -> ProvisionResult:
  1196. attempts["count"] += 1
  1197. if attempts["count"] == 1:
  1198. return ProvisionResult(success=False, step="provision_additional_domain", error="tls connect error")
  1199. active_accounts.append({"name": "cfmail-jp", "domain": "jp.example.test"})
  1200. return ProvisionResult(success=True, step="provision_additional_domain", new_domain="jp.example.test")
  1201. loop._cfmail_provisioner = FakeProvisioner() # type: ignore[assignment]
  1202. loop._cfmail_tracker = DomainHealthTracker(window_size=2, blacklist_threshold=2, rotation_cooldown_seconds=1)
  1203. loop._cfmail_domain_pool_replenish_worker(trigger_thread_id=2, reason="fresh domain budget reached")
  1204. assert attempts["count"] == 2
  1205. assert len(active_accounts) == 3
  1206. def test_registration_loop_schedules_replenish_when_usable_domain_pool_below_target(monkeypatch) -> None:
  1207. settings = _base_settings(register_mail_provider="cfmail")
  1208. loop = RegistrationLoop(settings)
  1209. loop._cfmail_active_domain_count = 3
  1210. monkeypatch.setattr(
  1211. loop,
  1212. "_current_cfmail_active_accounts",
  1213. lambda: [
  1214. {"name": "cfmail-tw", "domain": "tw.example.test"},
  1215. {"name": "cfmail-sg", "domain": "sg.example.test"},
  1216. ],
  1217. )
  1218. loop._cfmail_provisioner = object() # type: ignore[assignment]
  1219. scheduled: list[str] = []
  1220. monkeypatch.setattr(
  1221. loop,
  1222. "_schedule_cfmail_domain_pool_replenish",
  1223. lambda *, trigger_thread_id, reason: scheduled.append(f"{trigger_thread_id}:{reason}"),
  1224. raising=False,
  1225. )
  1226. loop._ensure_cfmail_domain_pool_target(trigger_thread_id=5, reason="startup")
  1227. assert scheduled == ["5:startup"]
  1228. def test_registration_loop_does_not_activate_wait_otp_stoploss_when_window_contains_success(monkeypatch) -> None:
  1229. settings = _base_settings(register_mail_provider="cfmail")
  1230. loop = RegistrationLoop(settings)
  1231. monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "demo.example.test")
  1232. loop._cfmail_wait_otp_window = 2
  1233. loop._cfmail_wait_otp_threshold = 2
  1234. loop._cfmail_wait_otp_cooldown_seconds = 60
  1235. timeout_result = {
  1236. "success": False,
  1237. "stage": "wait_otp",
  1238. "error_message": "otp retrieval failed",
  1239. "metadata": {
  1240. "mail_provider": "cfmail",
  1241. "email_domain": "demo.example.test",
  1242. "otp_wait_failure_reason": "mailbox_timeout_no_message",
  1243. "otp_mailbox_message_scan_count": 0,
  1244. },
  1245. }
  1246. success_result = {
  1247. "success": True,
  1248. "stage": "completed",
  1249. "email": "ok@demo.example.test",
  1250. "metadata": {
  1251. "mail_provider": "cfmail",
  1252. "email_domain": "demo.example.test",
  1253. },
  1254. }
  1255. loop._update_cfmail_wait_otp_stoploss(timeout_result)
  1256. loop._update_cfmail_wait_otp_stoploss(success_result)
  1257. stoploss = loop.snapshot()["cfmail_wait_otp_stoploss"]
  1258. assert stoploss["active_domain"] == "demo.example.test"
  1259. assert stoploss["in_cooldown"] is False
  1260. def test_registration_loop_does_not_activate_wait_otp_stoploss_when_window_contains_message_seen(monkeypatch) -> None:
  1261. settings = _base_settings(register_mail_provider="cfmail")
  1262. loop = RegistrationLoop(settings)
  1263. monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "demo.example.test")
  1264. loop._cfmail_wait_otp_window = 2
  1265. loop._cfmail_wait_otp_threshold = 2
  1266. loop._cfmail_wait_otp_cooldown_seconds = 60
  1267. timeout_result = {
  1268. "success": False,
  1269. "stage": "wait_otp",
  1270. "error_message": "otp retrieval failed",
  1271. "metadata": {
  1272. "mail_provider": "cfmail",
  1273. "email_domain": "demo.example.test",
  1274. "otp_wait_failure_reason": "mailbox_timeout_no_message",
  1275. "otp_mailbox_message_scan_count": 0,
  1276. },
  1277. }
  1278. message_seen_result = {
  1279. "success": False,
  1280. "stage": "wait_otp",
  1281. "error_message": "otp retrieval failed",
  1282. "metadata": {
  1283. "mail_provider": "cfmail",
  1284. "email_domain": "demo.example.test",
  1285. "otp_wait_failure_reason": "mailbox_timeout_no_match",
  1286. "otp_mailbox_message_scan_count": 1,
  1287. },
  1288. }
  1289. loop._update_cfmail_wait_otp_stoploss(timeout_result)
  1290. loop._update_cfmail_wait_otp_stoploss(message_seen_result)
  1291. stoploss = loop.snapshot()["cfmail_wait_otp_stoploss"]
  1292. assert stoploss["active_domain"] == "demo.example.test"
  1293. assert stoploss["in_cooldown"] is False
  1294. def test_registration_loop_disables_cfmail_canary_gate() -> None:
  1295. settings = _base_settings(register_mail_provider="cfmail")
  1296. loop = RegistrationLoop(settings)
  1297. assert loop._wait_if_cfmail_canary_pending(thread_id=2, provider="cfmail") is False
  1298. assert loop._wait_if_cfmail_canary_pending(thread_id=1, provider="cfmail") is False
  1299. def test_registration_loop_cfmail_canary_snapshot_is_disabled(monkeypatch) -> None:
  1300. settings = _base_settings(register_mail_provider="cfmail")
  1301. loop = RegistrationLoop(settings)
  1302. loop._arm_cfmail_canary("demo.example.test")
  1303. loop._mark_cfmail_canary_ready("demo.example.test", reason="mailbox_message_seen")
  1304. loop._update_cfmail_canary_after_result(thread_id=3, result={})
  1305. snapshot = loop.snapshot()["cfmail_canary"]
  1306. assert snapshot["pending"] is False
  1307. assert snapshot["active_domain"] == ""
  1308. assert snapshot["last_ready_reason"] == "disabled"
  1309. def test_registration_loop_tracks_fresh_domain_budget_after_mail_seen(monkeypatch) -> None:
  1310. monkeypatch.setenv("ZHUCE6_CFMAIL_FRESH_DOMAIN_ATTEMPT_BUDGET", "2")
  1311. settings = _base_settings(register_mail_provider="cfmail")
  1312. loop = RegistrationLoop(settings)
  1313. monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "demo.example.test")
  1314. first_result = {
  1315. "success": False,
  1316. "stage": "add_phone_gate",
  1317. "error_message": "post-create flow requires phone gate",
  1318. "metadata": {
  1319. "mail_provider": "cfmail",
  1320. "email_domain": "demo.example.test",
  1321. "otp_mailbox_message_scan_count": 1,
  1322. },
  1323. }
  1324. second_result = {
  1325. "success": False,
  1326. "stage": "wait_otp",
  1327. "error_message": "otp retrieval failed",
  1328. "metadata": {
  1329. "mail_provider": "cfmail",
  1330. "email_domain": "demo.example.test",
  1331. "otp_mailbox_message_scan_count": 0,
  1332. },
  1333. }
  1334. loop._update_cfmail_fresh_domain_budget(first_result)
  1335. snapshot = loop.snapshot()["cfmail_fresh_domain_budget"]
  1336. assert snapshot["completed_attempts"] == 1
  1337. assert snapshot["mail_seen_attempts"] == 1
  1338. assert snapshot["last_triggered_at"] == ""
  1339. loop._update_cfmail_fresh_domain_budget(second_result)
  1340. snapshot = loop.snapshot()["cfmail_fresh_domain_budget"]
  1341. assert snapshot["completed_attempts"] == 2
  1342. assert snapshot["mail_seen_attempts"] == 1
  1343. assert snapshot["last_reason"] == "fresh_domain_attempt_budget_reached"
  1344. def test_registration_loop_rotates_domain_when_fresh_domain_budget_reached(monkeypatch) -> None:
  1345. monkeypatch.setenv("ZHUCE6_CFMAIL_FRESH_DOMAIN_ATTEMPT_BUDGET", "2")
  1346. settings = _base_settings(register_mail_provider="cfmail")
  1347. loop = RegistrationLoop(settings)
  1348. monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "demo.example.test")
  1349. loop._cfmail_tracker = DomainHealthTracker(window_size=2, blacklist_threshold=2, rotation_cooldown_seconds=1)
  1350. loop._cfmail_fresh_domain_state.update(
  1351. {
  1352. "active_domain": "demo.example.test",
  1353. "completed_attempts": 2,
  1354. "mail_seen_attempts": 1,
  1355. "successes": 0,
  1356. "last_triggered_at": "2026-03-29T10:00:00",
  1357. "last_rotation_attempted_at": "",
  1358. "last_reason": "fresh_domain_attempt_budget_reached",
  1359. }
  1360. )
  1361. class FakeProvisioner:
  1362. def rotate_active_domain(self): # type: ignore[no-untyped-def]
  1363. return ProvisionResult(
  1364. success=True,
  1365. step="rotate",
  1366. old_domain="demo.example.test",
  1367. new_domain="auto-new.example.test",
  1368. )
  1369. reload_calls: list[bool] = []
  1370. class FakeManager:
  1371. def reload_if_needed(self, force=False): # type: ignore[no-untyped-def]
  1372. reload_calls.append(bool(force))
  1373. return True
  1374. loop._cfmail_manager = FakeManager()
  1375. loop._cfmail_provisioner = FakeProvisioner()
  1376. assert loop._rotate_cfmail_for_fresh_domain_budget(thread_id=5) is True
  1377. snapshot = loop.snapshot()
  1378. assert snapshot["cfmail_rotation"]["last_new_domain"] == "auto-new.example.test"
  1379. assert snapshot["cfmail_fresh_domain_budget"]["active_domain"] == "auto-new.example.test"
  1380. assert snapshot["cfmail_fresh_domain_budget"]["completed_attempts"] == 0
  1381. assert snapshot["cfmail_canary"]["pending"] is False
  1382. assert reload_calls == [True]
  1383. def test_registration_loop_throttles_cfmail_inflight_attempts(monkeypatch) -> None:
  1384. monkeypatch.setenv("ZHUCE6_CFMAIL_MAX_INFLIGHT", "1")
  1385. monkeypatch.setenv("ZHUCE6_CFMAIL_START_INTERVAL_SECONDS", "0")
  1386. settings = _base_settings(register_mail_provider="cfmail")
  1387. loop = RegistrationLoop(settings)
  1388. monkeypatch.setattr(
  1389. loop,
  1390. "_current_cfmail_active_accounts",
  1391. lambda: [{"name": "cfmail-demo", "domain": "demo.example.test"}],
  1392. )
  1393. assert loop._wait_if_cfmail_flow_throttled(thread_id=1, provider="cfmail") is False
  1394. assert loop._wait_if_cfmail_flow_throttled(thread_id=2, provider="cfmail") is True
  1395. loop._release_cfmail_flow_slot(thread_id=1)
  1396. assert loop._wait_if_cfmail_flow_throttled(thread_id=2, provider="cfmail") is False
  1397. def test_registration_loop_throttles_cfmail_start_interval(monkeypatch) -> None:
  1398. monkeypatch.setenv("ZHUCE6_CFMAIL_MAX_INFLIGHT", "2")
  1399. monkeypatch.setenv("ZHUCE6_CFMAIL_START_INTERVAL_SECONDS", "15")
  1400. settings = _base_settings(register_mail_provider="cfmail")
  1401. loop = RegistrationLoop(settings)
  1402. monkeypatch.setattr(
  1403. loop,
  1404. "_current_cfmail_active_accounts",
  1405. lambda: [{"name": "cfmail-demo", "domain": "demo.example.test"}],
  1406. )
  1407. current_time = {"value": 100.0}
  1408. monkeypatch.setattr("core.registration.time.time", lambda: current_time["value"])
  1409. assert loop._wait_if_cfmail_flow_throttled(thread_id=1, provider="cfmail") is False
  1410. loop._release_cfmail_flow_slot(thread_id=1)
  1411. current_time["value"] = 105.0
  1412. assert loop._wait_if_cfmail_flow_throttled(thread_id=2, provider="cfmail") is True
  1413. current_time["value"] = 116.0
  1414. assert loop._wait_if_cfmail_flow_throttled(thread_id=2, provider="cfmail") is False
  1415. def test_registration_loop_activates_live_wait_otp_stoploss_on_stalled_waits(monkeypatch) -> None:
  1416. settings = _base_settings(register_mail_provider="cfmail")
  1417. loop = RegistrationLoop(settings)
  1418. monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "demo.example.test")
  1419. loop._cfmail_wait_otp_cooldown_seconds = 60
  1420. loop._cfmail_wait_otp_live_threshold = 2
  1421. loop._cfmail_wait_otp_live_age_seconds = 30
  1422. account_a = MailboxAccount(
  1423. email="a@demo.example.test",
  1424. account_id="jwt-a",
  1425. extra={"email_domain": "demo.example.test"},
  1426. )
  1427. account_b = MailboxAccount(
  1428. email="b@demo.example.test",
  1429. account_id="jwt-b",
  1430. extra={"email_domain": "demo.example.test"},
  1431. )
  1432. loop._on_cfmail_wait_progress(
  1433. account_a,
  1434. {"message_scan_count": 0, "elapsed_seconds": 35},
  1435. )
  1436. loop._on_cfmail_wait_progress(
  1437. account_b,
  1438. {"message_scan_count": 0, "elapsed_seconds": 36},
  1439. )
  1440. stoploss = loop.snapshot()["cfmail_wait_otp_stoploss"]
  1441. assert stoploss["active_domain"] == "demo.example.test"
  1442. assert stoploss["in_cooldown"] is True
  1443. assert stoploss["last_reason"] == "live wait_otp no-message threshold reached"
  1444. assert stoploss["last_no_message_timeouts"] == 2
  1445. def test_registration_loop_live_wait_otp_stoploss_can_be_disabled(monkeypatch) -> None:
  1446. settings = _base_settings(register_mail_provider="cfmail")
  1447. loop = RegistrationLoop(settings)
  1448. monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "demo.example.test")
  1449. loop._cfmail_wait_otp_cooldown_seconds = 60
  1450. loop._cfmail_wait_otp_live_threshold = 0
  1451. loop._cfmail_wait_otp_live_age_seconds = 30
  1452. account_a = MailboxAccount(
  1453. email="a@demo.example.test",
  1454. account_id="jwt-a",
  1455. extra={"email_domain": "demo.example.test"},
  1456. )
  1457. account_b = MailboxAccount(
  1458. email="b@demo.example.test",
  1459. account_id="jwt-b",
  1460. extra={"email_domain": "demo.example.test"},
  1461. )
  1462. loop._on_cfmail_wait_progress(
  1463. account_a,
  1464. {"message_scan_count": 0, "elapsed_seconds": 35},
  1465. )
  1466. loop._on_cfmail_wait_progress(
  1467. account_b,
  1468. {"message_scan_count": 0, "elapsed_seconds": 36},
  1469. )
  1470. stoploss = loop.snapshot()["cfmail_wait_otp_stoploss"]
  1471. assert stoploss["in_cooldown"] is False
  1472. def test_registration_loop_does_not_activate_wait_otp_stoploss_for_non_mailbox_timeout() -> None:
  1473. settings = _base_settings(register_mail_provider="cfmail")
  1474. loop = RegistrationLoop(settings)
  1475. loop._cfmail_wait_otp_window = 2
  1476. loop._cfmail_wait_otp_threshold = 2
  1477. loop._cfmail_wait_otp_cooldown_seconds = 60
  1478. result = {
  1479. "success": False,
  1480. "stage": "wait_otp",
  1481. "error_message": "otp retrieval failed",
  1482. "metadata": {
  1483. "mail_provider": "cfmail",
  1484. "email_domain": "demo.example.test",
  1485. "otp_wait_failure_reason": "mailbox_timeout_no_match",
  1486. "otp_mailbox_message_scan_count": 1,
  1487. },
  1488. }
  1489. loop._update_cfmail_wait_otp_stoploss(result)
  1490. loop._update_cfmail_wait_otp_stoploss(result)
  1491. stoploss = loop.snapshot()["cfmail_wait_otp_stoploss"]
  1492. assert stoploss["in_cooldown"] is False
  1493. def test_registration_loop_rotates_cfmail_domain_when_wait_otp_stoploss_active(monkeypatch) -> None:
  1494. settings = _base_settings(register_mail_provider="cfmail")
  1495. loop = RegistrationLoop(settings)
  1496. loop._providers = ["cfmail"]
  1497. monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "demo.example.test")
  1498. loop._cfmail_tracker = DomainHealthTracker(window_size=2, blacklist_threshold=2, rotation_cooldown_seconds=1)
  1499. loop._cfmail_wait_otp_cooldown_seconds = 60
  1500. loop._cfmail_wait_otp_state.update(
  1501. {
  1502. "active_domain": "demo.example.test",
  1503. "in_cooldown": True,
  1504. "cooldown_until": 9999999999.0,
  1505. "last_triggered_at": "2026-03-29T10:00:00",
  1506. "last_rotation_attempted_at": "",
  1507. }
  1508. )
  1509. class FakeProvisioner:
  1510. def rotate_active_domain(self): # type: ignore[no-untyped-def]
  1511. return ProvisionResult(
  1512. success=True,
  1513. step="rotate",
  1514. old_domain="demo.example.test",
  1515. new_domain="auto-new.example.test",
  1516. )
  1517. reload_calls: list[bool] = []
  1518. class FakeManager:
  1519. def reload_if_needed(self, force=False): # type: ignore[no-untyped-def]
  1520. reload_calls.append(bool(force))
  1521. return True
  1522. loop._cfmail_manager = FakeManager()
  1523. loop._cfmail_provisioner = FakeProvisioner()
  1524. assert loop._wait_if_cfmail_wait_otp_stopped(thread_id=1, provider="cfmail") is True
  1525. snapshot = loop.snapshot()
  1526. assert snapshot["cfmail_rotation"]["last_new_domain"] == "auto-new.example.test"
  1527. assert snapshot["cfmail_wait_otp_stoploss"]["in_cooldown"] is False
  1528. assert snapshot["cfmail_wait_otp_stoploss"]["active_domain"] == "auto-new.example.test"
  1529. assert reload_calls == [True]
  1530. def test_registration_loop_rotates_cfmail_domain_when_add_phone_stoploss_active(monkeypatch) -> None:
  1531. settings = _base_settings(register_mail_provider="cfmail")
  1532. loop = RegistrationLoop(settings)
  1533. loop._providers = ["cfmail"]
  1534. monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "demo.example.test")
  1535. loop._cfmail_tracker = DomainHealthTracker(window_size=2, blacklist_threshold=2, rotation_cooldown_seconds=1)
  1536. loop._cfmail_add_phone_cooldown_seconds = 60
  1537. loop._cfmail_add_phone_state.update(
  1538. {
  1539. "active_domain": "demo.example.test",
  1540. "in_cooldown": True,
  1541. "cooldown_until": 9999999999.0,
  1542. "last_triggered_at": "2026-03-29T10:00:00",
  1543. "last_rotation_attempted_at": "",
  1544. }
  1545. )
  1546. class FakeProvisioner:
  1547. def rotate_active_domain(self): # type: ignore[no-untyped-def]
  1548. return ProvisionResult(
  1549. success=True,
  1550. step="rotate",
  1551. old_domain="demo.example.test",
  1552. new_domain="auto-new.example.test",
  1553. )
  1554. loop._cfmail_provisioner = FakeProvisioner()
  1555. assert loop._wait_if_cfmail_add_phone_stopped(thread_id=1, provider="cfmail") is True
  1556. snapshot = loop.snapshot()
  1557. assert snapshot["cfmail_rotation"]["last_new_domain"] == "auto-new.example.test"
  1558. assert snapshot["cfmail_add_phone_stoploss"]["in_cooldown"] is False
  1559. assert snapshot["cfmail_add_phone_stoploss"]["active_domain"] == "auto-new.example.test"
  1560. def test_registration_loop_ignores_stale_wait_otp_result_from_old_domain(monkeypatch) -> None:
  1561. settings = _base_settings(register_mail_provider="cfmail")
  1562. loop = RegistrationLoop(settings)
  1563. monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "auto-new.example.test")
  1564. result = {
  1565. "success": False,
  1566. "stage": "wait_otp",
  1567. "error_message": "otp retrieval failed",
  1568. "metadata": {
  1569. "mail_provider": "cfmail",
  1570. "email_domain": "auto-old.example.test",
  1571. "otp_wait_failure_reason": "mailbox_timeout_no_message",
  1572. "otp_mailbox_message_scan_count": 0,
  1573. },
  1574. }
  1575. loop._update_cfmail_wait_otp_stoploss(result)
  1576. stoploss = loop.snapshot()["cfmail_wait_otp_stoploss"]
  1577. assert stoploss["active_domain"] == ""
  1578. assert stoploss["in_cooldown"] is False
  1579. def test_registration_loop_does_not_abort_old_domain_wait_after_rotation(monkeypatch) -> None:
  1580. settings = _base_settings(register_mail_provider="cfmail")
  1581. loop = RegistrationLoop(settings)
  1582. monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "auto-new.example.test")
  1583. loop._cfmail_wait_otp_state.update(
  1584. {
  1585. "active_domain": "auto-old.example.test",
  1586. "in_cooldown": True,
  1587. "cooldown_until": 9999999999.0,
  1588. "last_triggered_at": "2026-03-29T10:00:00",
  1589. }
  1590. )
  1591. account = MailboxAccount(
  1592. email="a@auto-old.example.test",
  1593. account_id="jwt-old",
  1594. extra={
  1595. "email_domain": "auto-old.example.test",
  1596. "otp_wait_started_at": 1743242390.0,
  1597. },
  1598. )
  1599. assert loop._should_abort_cfmail_wait(account) is False
  1600. def test_registration_loop_aborts_wait_for_domain_in_wait_otp_cooldown(monkeypatch) -> None:
  1601. settings = _base_settings(register_mail_provider="cfmail")
  1602. loop = RegistrationLoop(settings)
  1603. monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "demo.example.test")
  1604. loop._cfmail_wait_otp_state.update(
  1605. {
  1606. "active_domain": "demo.example.test",
  1607. "in_cooldown": True,
  1608. "cooldown_until": 9999999999.0,
  1609. }
  1610. )
  1611. account = MailboxAccount(
  1612. email="a@demo.example.test",
  1613. account_id="jwt-demo",
  1614. extra={
  1615. "email_domain": "demo.example.test",
  1616. "otp_wait_started_at": 1743242410.0,
  1617. },
  1618. )
  1619. assert loop._should_abort_cfmail_wait(account) is True
  1620. def test_registration_burst_scheduler_runs_batch_with_batch_config(monkeypatch, tmp_path: Path) -> None:
  1621. settings = _base_settings(
  1622. register_log_file="",
  1623. runtime_state_file=tmp_path / "runtime_state.json",
  1624. register_batch_threads=1,
  1625. register_batch_target_count=20,
  1626. register_batch_interval_seconds=60,
  1627. )
  1628. observed: dict[str, int] = {}
  1629. class FakeLoop:
  1630. def __init__(self, batch_settings): # type: ignore[no-untyped-def]
  1631. observed["threads"] = batch_settings.register_threads
  1632. observed["target"] = batch_settings.register_target_count
  1633. self._threads = []
  1634. def start(self) -> None:
  1635. return None
  1636. def stop(self) -> None:
  1637. return None
  1638. def snapshot(self) -> dict[str, object]:
  1639. return {
  1640. "name": "register",
  1641. "status": "stopped",
  1642. "threads_alive": 0,
  1643. "threads_total": 1,
  1644. "total_attempts": 22,
  1645. "total_success": 20,
  1646. "total_failure": 2,
  1647. "success_rate": 90.9,
  1648. "target_count": 20,
  1649. "target_reached": True,
  1650. "last_error": "",
  1651. "proxy": None,
  1652. "proxy_pool_enabled": False,
  1653. "mail_provider": "mailtm",
  1654. "interval_seconds": 5,
  1655. "run_count": 22,
  1656. "success_count": 20,
  1657. "failure_count": 2,
  1658. "is_running": False,
  1659. "last_started_at": None,
  1660. "last_finished_at": None,
  1661. "last_duration_seconds": None,
  1662. "next_run_at": None,
  1663. "failure_by_stage": {"add_phone_gate": 2},
  1664. "failure_signals": {"add_phone_gate": 2},
  1665. "recent_failure_hotspots": [{"key": "add_phone_gate", "stage": "add_phone_gate", "count": 2}],
  1666. "recent_attempts": [
  1667. {
  1668. "timestamp": "2026-03-27T10:00:00",
  1669. "success": False,
  1670. "stage": "add_phone_gate",
  1671. "signal": "add_phone_gate",
  1672. "error_message": "phone gate",
  1673. "email_domain": "demo.example.test",
  1674. "post_create_gate": "add_phone",
  1675. "create_account_error_code": "",
  1676. "proxy_key": "",
  1677. "email": "",
  1678. }
  1679. ],
  1680. "active_domain_recent_attempts": [],
  1681. "active_domain_failure_by_stage": {"add_phone_gate": 2},
  1682. "active_domain_failure_signals": {"add_phone_gate": 2},
  1683. "active_domain_recent_failure_hotspots": [{"key": "add_phone_gate", "stage": "add_phone_gate", "count": 2}],
  1684. "cfmail_rotation": None,
  1685. "cfmail_add_phone_stoploss": {"active_domain": "demo.example.test"}
  1686. }
  1687. monkeypatch.setattr("main.RegistrationLoop", FakeLoop)
  1688. scheduler = RegistrationBurstScheduler(settings)
  1689. original_absorb = scheduler._absorb_batch_snapshot
  1690. def absorb_and_stop(snapshot: dict[str, object], *, duration_seconds: float) -> None:
  1691. original_absorb(snapshot, duration_seconds=duration_seconds)
  1692. scheduler.stop()
  1693. monkeypatch.setattr(scheduler, "_absorb_batch_snapshot", absorb_and_stop)
  1694. scheduler.run()
  1695. payload = json.loads((tmp_path / "runtime_state.json").read_text(encoding="utf-8"))
  1696. register_snapshot = payload["register_snapshot"]
  1697. assert observed == {"threads": 1, "target": 20}
  1698. assert register_snapshot["status"] == "stopped"
  1699. assert register_snapshot["scheduler_mode"] == "burst"
  1700. assert register_snapshot["run_count"] == 1
  1701. assert register_snapshot["total_success"] == 20
  1702. assert register_snapshot["batch_target_count"] == 20