| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791179217931794179517961797179817991800180118021803180418051806180718081809181018111812181318141815181618171818181918201821182218231824182518261827182818291830183118321833183418351836183718381839184018411842184318441845184618471848184918501851185218531854185518561857185818591860186118621863186418651866186718681869187018711872187318741875187618771878187918801881188218831884188518861887188818891890189118921893189418951896189718981899190019011902190319041905190619071908190919101911191219131914191519161917191819191920192119221923192419251926192719281929193019311932193319341935193619371938193919401941194219431944194519461947194819491950195119521953195419551956195719581959196019611962196319641965196619671968196919701971197219731974197519761977197819791980198119821983198419851986198719881989199019911992199319941995199619971998199920002001200220032004200520062007200820092010201120122013201420152016201720182019202020212022202320242025202620272028202920302031203220332034203520362037203820392040204120422043204420452046204720482049205020512052205320542055205620572058205920602061206220632064206520662067206820692070207120722073207420752076207720782079 |
- import json
- from dataclasses import replace
- from pathlib import Path
- from core.base_mailbox import MailboxAccount
- from core.settings import AppSettings
- from core.cfmail_domain_rotation import DomainHealthTracker
- from core.cfmail_provisioner import ProvisionResult
- from main import RegistrationBurstScheduler, RegistrationLoop
- def _base_settings(**overrides) -> AppSettings:
- settings = AppSettings(
- cleanup_enabled=False,
- register_sleep_min=0,
- register_sleep_max=0,
- register_mail_provider="mailtm,mailgw",
- )
- return replace(settings, **overrides)
- def test_registration_loop_switches_provider_after_configured_failures(monkeypatch) -> None:
- settings = _base_settings(register_max_consecutive_failures=2)
- loop = RegistrationLoop(settings)
- loop._providers = ["mailtm", "mailgw"]
- calls: list[str] = []
- def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
- calls.append(kwargs["mail_provider"])
- if len(calls) >= 3:
- loop._stop_event.set()
- return {"success": False, "error_message": "failed"}
- monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
- loop._worker(thread_id=1, initial_provider="mailtm")
- assert calls == ["mailtm", "mailtm", "mailgw"]
- def test_registration_loop_stops_when_target_count_reached(monkeypatch) -> None:
- settings = _base_settings(register_target_count=2, register_mail_provider="mailtm")
- loop = RegistrationLoop(settings)
- loop._providers = ["mailtm"]
- calls: list[str] = []
- def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
- calls.append(kwargs["mail_provider"])
- return {"success": True, "email": f"user{len(calls)}@example.com"}
- monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
- monkeypatch.setattr(loop, "_sync_cpa_from_success", lambda result, thread_id: (True, "", str(result.get("email") or "")))
- loop._worker(thread_id=1, initial_provider="mailtm")
- assert calls == ["mailtm", "mailtm"]
- assert loop._target_reached.is_set() is True
- def test_registration_loop_uses_register_proxy_when_proxy_pool_disabled(monkeypatch) -> None:
- settings = _base_settings(register_proxy="http://127.0.0.1:7890", register_mail_provider="mailtm")
- loop = RegistrationLoop(settings)
- loop._providers = ["mailtm"]
- proxies: list[str | None] = []
- def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
- proxies.append(kwargs["proxy"])
- loop._stop_event.set()
- return {"success": False, "error_message": "failed"}
- monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
- loop._worker(thread_id=1, initial_provider="mailtm")
- assert proxies == ["http://127.0.0.1:7890"]
- def test_registration_loop_passes_selected_cfmail_profile_into_register_flow(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- loop._providers = ["cfmail"]
- recorded: dict[str, object] = {}
- monkeypatch.setattr(
- loop,
- "_current_cfmail_active_accounts",
- lambda: [{"name": "cfmail-tw", "domain": "tw.example.test"}],
- )
- def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
- recorded["cfmail_profile_name"] = kwargs["cfmail_profile_name"]
- loop._stop_event.set()
- return {"success": True, "email": "demo@example.com", "metadata": {"email_domain": "tw.example.test"}}
- monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
- monkeypatch.setattr(loop, "_sync_cpa_from_success", lambda result, thread_id: (True, "", str(result.get("email") or "")))
- loop._worker(thread_id=1, initial_provider="cfmail")
- assert recorded["cfmail_profile_name"] == "cfmail-tw"
- def test_registration_loop_snapshot_exposes_cfmail_domain_pool(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- monkeypatch.setattr(
- loop,
- "_current_cfmail_active_accounts",
- lambda: [
- {"name": "cfmail-tw", "domain": "tw.example.test"},
- {"name": "cfmail-sg", "domain": "sg.example.test"},
- ],
- )
- loop._recent_attempts.extend(
- [
- {
- "timestamp": "2026-03-30T22:50:00",
- "success": True,
- "stage": "completed",
- "signal": "",
- "error_message": "",
- "email_domain": "tw.example.test",
- "proxy_key": "tw-node",
- "email": "ok@tw.example.test",
- },
- {
- "timestamp": "2026-03-30T22:50:10",
- "success": False,
- "stage": "create_account",
- "signal": "mailbox_reused",
- "error_message": "create account failed",
- "email_domain": "sg.example.test",
- "create_account_error_code": "user_already_exists",
- "proxy_key": "sg-node",
- "email": "dup@sg.example.test",
- },
- ]
- )
- loop._cfmail_flow_state["inflight_by_thread"] = {
- 1: {"domain": "tw.example.test", "profile_name": "cfmail-tw"},
- }
- loop._cfmail_flow_state["last_started_by_domain"] = {
- "tw.example.test": 100.0,
- "sg.example.test": 90.0,
- }
- snapshot = loop.snapshot()
- domain_pool = snapshot["cfmail_domain_pool"]
- assert domain_pool["active_count"] == 2
- assert [item["domain"] for item in domain_pool["active_domains"]] == [
- "tw.example.test",
- "sg.example.test",
- ]
- assert domain_pool["active_domains"][0]["inflight"] == 1
- assert domain_pool["active_domains"][0]["recent_success"] == 1
- assert domain_pool["active_domains"][1]["recent_failure"] == 1
- assert domain_pool["active_domains"][1]["failure_signals"]["mailbox_reused"] == 1
- def test_registration_loop_uses_multi_domain_flow_defaults(monkeypatch) -> None:
- monkeypatch.delenv("ZHUCE6_CFMAIL_START_INTERVAL_SECONDS", raising=False)
- monkeypatch.delenv("ZHUCE6_CFMAIL_MAX_INFLIGHT", raising=False)
- monkeypatch.delenv("ZHUCE6_CFMAIL_ACTIVE_DOMAIN_COUNT", raising=False)
- monkeypatch.delenv("ZHUCE6_CFMAIL_FRESH_DOMAIN_ATTEMPT_BUDGET", raising=False)
- loop = RegistrationLoop(_base_settings(register_mail_provider="cfmail"))
- assert loop._cfmail_start_interval_seconds == 8
- assert loop._cfmail_max_inflight == 4
- assert loop._cfmail_active_domain_count == 3
- assert loop._cfmail_fresh_domain_attempt_budget == 2
- assert loop._cfmail_add_phone_window == 10
- assert loop._cfmail_add_phone_threshold == 3
- assert loop._cfmail_wait_otp_window == 6
- assert loop._cfmail_wait_otp_threshold == 2
- def test_registration_loop_caps_pending_token_retry_delay_to_one_minute(monkeypatch) -> None:
- monkeypatch.setenv("ZHUCE6_PENDING_TOKEN_RETRY_DELAY_SECONDS", "300")
- loop = RegistrationLoop(_base_settings())
- assert loop._pending_token_retry_delay_seconds == 60
- def test_registration_loop_uses_proxy_pool_and_releases_lease(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="mailtm")
- loop = RegistrationLoop(settings)
- loop._providers = ["mailtm"]
- recorded: dict[str, object] = {}
- class FakeLease:
- name = "sg-node"
- local_port = 17891
- proxy_url = "socks5://127.0.0.1:17891"
- class FakePool:
- def acquire(self, timeout=5.0, preferred_name=None, preferred_regions=()): # type: ignore[no-untyped-def]
- recorded["timeout"] = timeout
- recorded["preferred_name"] = preferred_name
- recorded["preferred_regions"] = tuple(preferred_regions)
- return FakeLease()
- def release(self, lease, *, success, stage=None): # type: ignore[no-untyped-def]
- recorded["released"] = (lease.proxy_url, success, stage)
- def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
- recorded["proxy"] = kwargs["proxy"]
- loop._stop_event.set()
- return {"success": True, "email": "demo@example.com"}
- loop._proxy_pool = FakePool()
- monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
- monkeypatch.setattr(loop, "_sync_cpa_from_success", lambda result, thread_id: (True, "", str(result.get("email") or "")))
- loop._worker(thread_id=1, initial_provider="mailtm")
- assert recorded["proxy"] == "socks5://127.0.0.1:17891"
- assert recorded["released"] == ("socks5://127.0.0.1:17891", True, "completed")
- def test_registration_loop_prefers_fresh_proxy_regions(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="mailtm", register_fresh_proxy_regions=("tw", "sg"))
- loop = RegistrationLoop(settings)
- loop._providers = ["mailtm"]
- recorded: dict[str, object] = {}
- class FakeLease:
- name = "tw-node"
- local_port = 17891
- proxy_url = "socks5://127.0.0.1:17891"
- class FakePool:
- def acquire(self, timeout=5.0, preferred_name=None, preferred_regions=()): # type: ignore[no-untyped-def]
- recorded["timeout"] = timeout
- recorded["preferred_name"] = preferred_name
- recorded["preferred_regions"] = tuple(preferred_regions)
- return FakeLease()
- def release(self, lease, *, success, stage=None): # type: ignore[no-untyped-def]
- recorded["released"] = (lease.proxy_url, success, stage)
- def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
- recorded["proxy"] = kwargs["proxy"]
- loop._stop_event.set()
- return {"success": True, "email": "demo@example.com"}
- loop._proxy_pool = FakePool()
- monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
- monkeypatch.setattr(loop, "_sync_cpa_from_success", lambda result, thread_id: (True, "", str(result.get("email") or "")))
- loop._worker(thread_id=1, initial_provider="mailtm")
- assert recorded["preferred_regions"] == ("tw", "sg")
- def test_registration_loop_starts_proxy_pool_for_direct_urls(monkeypatch) -> None:
- settings = _base_settings(
- register_mail_provider="mailtm",
- register_threads=1,
- proxy_pool_direct_urls="http://5.6.7.8:8080",
- )
- loop = RegistrationLoop(settings)
- recorded: dict[str, object] = {}
- class FakePool:
- def start(self) -> None:
- recorded["pool_started"] = True
- class FakeThread:
- def __init__(self, target=None, args=(), daemon=None, name=None): # type: ignore[no-untyped-def]
- recorded["thread_args"] = {"target": target, "args": args, "daemon": daemon, "name": name}
- def start(self) -> None:
- recorded["thread_started"] = True
- monkeypatch.setattr("main.load_all", lambda: None)
- monkeypatch.setattr("core.proxy_pool.ProxyPool.from_settings", lambda settings: FakePool())
- monkeypatch.setattr("main.threading.Thread", FakeThread)
- monkeypatch.setattr("core.registration._maybe_reconcile_cpa_runtime", lambda **kwargs: None)
- monkeypatch.setattr(loop, "_write_runtime_state", lambda: None)
- loop.start()
- assert recorded["pool_started"] is True
- assert recorded["thread_started"] is True
- assert loop._proxy_pool is not None
- def test_registration_loop_start_reconciles_pool_backups_to_cpa(monkeypatch, tmp_path: Path) -> None:
- settings = _base_settings(
- register_mail_provider="mailtm",
- register_threads=1,
- backend="cpa",
- cpa_runtime_reconcile_enabled=True,
- pool_dir=tmp_path / "pool",
- )
- loop = RegistrationLoop(settings)
- recorded: dict[str, object] = {}
- class FakeThread:
- def __init__(self, target=None, args=(), daemon=None, name=None): # type: ignore[no-untyped-def]
- recorded["thread_args"] = {"target": target, "args": args, "daemon": daemon, "name": name}
- def start(self) -> None:
- recorded["thread_started"] = True
- monkeypatch.setattr("main.load_all", lambda: None)
- monkeypatch.setattr("main.threading.Thread", FakeThread)
- monkeypatch.setattr("core.registration.create_backend_client", lambda settings: "backend-client")
- monkeypatch.setattr(
- "core.registration._maybe_reconcile_cpa_runtime",
- lambda **kwargs: recorded.setdefault("reconcile", kwargs),
- )
- monkeypatch.setattr(loop, "_write_runtime_state", lambda: None)
- loop.start()
- assert recorded["thread_started"] is True
- assert recorded["reconcile"]["pool_dir"] == settings.pool_dir
- assert recorded["reconcile"]["client"] == "backend-client"
- def test_registration_loop_start_normalizes_cfmail_to_domain_pool(monkeypatch) -> None:
- settings = _base_settings(
- register_mail_provider="cfmail",
- register_threads=1,
- backend="cpa",
- cpa_runtime_reconcile_enabled=False,
- )
- loop = RegistrationLoop(settings)
- recorded: dict[str, object] = {}
- class FakeThread:
- def __init__(self, target=None, args=(), daemon=None, name=None): # type: ignore[no-untyped-def]
- recorded["thread_args"] = {"target": target, "args": args, "daemon": daemon, "name": name}
- def start(self) -> None:
- recorded["thread_started"] = True
- class FakeProvisioner:
- def __init__(self, *, proxy_url=None, **_kwargs): # type: ignore[no-untyped-def]
- recorded["proxy_url"] = proxy_url
- def normalize_to_domain_pool(self, target_count): # type: ignore[no-untyped-def]
- recorded["normalized"] = True
- recorded["target_count"] = target_count
- return {
- "active_domains": ["auto-live.example.test", "auto-next.example.test"],
- "provisioned_domains": ["auto-next.example.test"],
- "retired_domains": [],
- }
- class FakeManager:
- def reload_if_needed(self, force=False): # type: ignore[no-untyped-def]
- recorded["reload_force"] = force
- return True
- monkeypatch.setattr("main.load_all", lambda: None)
- monkeypatch.setattr("main.threading.Thread", FakeThread)
- monkeypatch.setattr("core.cfmail_provisioner.CfmailProvisioner", FakeProvisioner)
- monkeypatch.setattr("core.cfmail.DEFAULT_CFMAIL_MANAGER", FakeManager())
- monkeypatch.setattr(loop, "_ensure_cfmail_active_domain_ready", lambda: True)
- monkeypatch.setattr(loop, "_write_runtime_state", lambda: None)
- monkeypatch.setattr(loop, "_log", lambda message: recorded.setdefault("logs", []).append(message))
- loop.start()
- assert recorded["thread_started"] is True
- assert recorded["proxy_url"] == settings.register_proxy
- assert recorded["normalized"] is True
- assert recorded["target_count"] == 3
- assert recorded["reload_force"] is True
- assert any("normalized active domain pool" in line for line in recorded["logs"])
- def test_registration_loop_retries_device_id_once_with_new_proxy(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="mailtm")
- loop = RegistrationLoop(settings)
- loop._providers = ["mailtm"]
- recorded: dict[str, object] = {"releases": []}
- class FakeLease:
- def __init__(self, name: str, local_port: int) -> None:
- self.name = name
- self.local_port = local_port
- self.proxy_url = f"socks5://127.0.0.1:{local_port}"
- class FakePool:
- def __init__(self) -> None:
- self._leases = [FakeLease("tw-bad", 17891), FakeLease("sg-good", 17892)]
- def acquire(self, timeout=5.0, preferred_name=None, preferred_regions=()): # type: ignore[no-untyped-def]
- recorded.setdefault("timeouts", []).append(timeout)
- return self._leases.pop(0)
- def release(self, lease, *, success, stage=None): # type: ignore[no-untyped-def]
- recorded["releases"].append((lease.name, success, stage))
- attempts: list[str] = []
- def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
- attempts.append(kwargs["proxy"])
- if len(attempts) == 1:
- return {
- "success": False,
- "stage": "device_id",
- "error_message": "device id acquisition failed",
- }
- loop._stop_event.set()
- return {"success": True, "stage": "completed", "email": "demo@example.com"}
- loop._proxy_pool = FakePool()
- monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
- monkeypatch.setattr(loop, "_sync_cpa_from_success", lambda result, thread_id: (True, "", str(result.get("email") or "")))
- loop._worker(thread_id=1, initial_provider="mailtm")
- assert attempts == ["socks5://127.0.0.1:17891", "socks5://127.0.0.1:17892"]
- assert recorded["releases"] == [
- ("tw-bad", False, "device_id"),
- ("sg-good", True, "completed"),
- ]
- snapshot = loop.snapshot()
- assert snapshot["total_attempts"] == 1
- assert snapshot["total_success"] == 1
- assert snapshot["total_failure"] == 0
- def test_registration_loop_syncs_cpa_immediately_after_success(monkeypatch, tmp_path: Path) -> None:
- settings = _base_settings(
- register_mail_provider="mailtm",
- )
- loop = RegistrationLoop(settings)
- loop._providers = ["mailtm"]
- pool_file = tmp_path / "fresh@example.com.json"
- pool_file.write_text(json.dumps({"email": "fresh@example.com", "access_token": "tok", "account_id": "acct"}), encoding="utf-8")
- recorded: dict[str, object] = {}
- def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
- del kwargs
- loop._stop_event.set()
- return {
- "success": True,
- "email": "fresh@example.com",
- "pool_file": str(pool_file),
- "written_to_pool": True,
- }
- monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
- monkeypatch.setattr("core.registration.get_management_key", lambda: "secret") # type: ignore[no-untyped-def]
- 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]
- def fake_upload_to_cpa(token_data, api_url=None, api_key=None, proxy=None): # type: ignore[no-untyped-def]
- recorded["token_data"] = token_data
- recorded["api_url"] = api_url
- recorded["api_key"] = api_key
- recorded["proxy"] = proxy
- return True, "upload success"
- monkeypatch.setattr("platforms.chatgpt.cpa_upload.upload_to_cpa", fake_upload_to_cpa)
- loop._worker(thread_id=1, initial_provider="mailtm")
- assert recorded["token_data"]["email"] == "fresh@example.com"
- assert recorded["api_url"] == "http://127.0.0.1:8317"
- assert recorded["api_key"] == "secret"
- assert recorded["proxy"] is None
- payload = json.loads(pool_file.read_text(encoding="utf-8"))
- assert payload["cpa_sync_status"] == "synced"
- snapshot = loop.snapshot()
- assert snapshot["total_success"] == 1
- assert snapshot["total_success_registered"] == 1
- assert snapshot["total_cpa_sync_success"] == 1
- assert snapshot["total_cpa_sync_failure"] == 0
- assert snapshot["registered_success_rate"] == 100.0
- assert snapshot["cpa_sync_success_rate"] == 100.0
- def test_registration_loop_syncs_cpa_before_breaking_on_target_reached(monkeypatch, tmp_path: Path) -> None:
- settings = _base_settings(
- register_mail_provider="mailtm",
- register_target_count=1,
- )
- loop = RegistrationLoop(settings)
- loop._providers = ["mailtm"]
- pool_file = tmp_path / "target@example.com.json"
- pool_file.write_text(json.dumps({"email": "target@example.com", "access_token": "tok", "account_id": "acct"}), encoding="utf-8")
- recorded: dict[str, object] = {}
- def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
- del kwargs
- return {
- "success": True,
- "email": "target@example.com",
- "pool_file": str(pool_file),
- "written_to_pool": True,
- }
- monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
- monkeypatch.setattr("core.registration.get_management_key", lambda: "secret") # type: ignore[no-untyped-def]
- 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]
- def fake_upload_to_cpa(token_data, api_url=None, api_key=None, proxy=None): # type: ignore[no-untyped-def]
- recorded["token_data"] = token_data
- recorded["api_url"] = api_url
- recorded["api_key"] = api_key
- recorded["proxy"] = proxy
- return True, "upload success"
- monkeypatch.setattr("platforms.chatgpt.cpa_upload.upload_to_cpa", fake_upload_to_cpa)
- loop._worker(thread_id=1, initial_provider="mailtm")
- assert loop._target_reached.is_set() is True
- assert recorded["token_data"]["email"] == "target@example.com"
- payload = json.loads(pool_file.read_text(encoding="utf-8"))
- assert payload["cpa_sync_status"] == "synced"
- snapshot = loop.snapshot()
- assert snapshot["total_attempts"] == 1
- assert snapshot["total_success"] == 1
- assert snapshot["total_cpa_sync_success"] == 1
- assert snapshot["total_cpa_sync_failure"] == 0
- def test_registration_loop_marks_pool_backup_when_direct_cpa_sync_fails(monkeypatch, tmp_path: Path) -> None:
- settings = _base_settings(
- register_mail_provider="mailtm",
- )
- loop = RegistrationLoop(settings)
- loop._providers = ["mailtm"]
- pool_file = tmp_path / "fresh@example.com.json"
- pool_file.write_text(json.dumps({"email": "fresh@example.com", "access_token": "tok", "account_id": "acct"}), encoding="utf-8")
- def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
- del kwargs
- loop._stop_event.set()
- return {
- "success": True,
- "email": "fresh@example.com",
- "pool_file": str(pool_file),
- "written_to_pool": True,
- }
- monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
- monkeypatch.setattr("core.registration.get_management_key", lambda: "secret") # type: ignore[no-untyped-def]
- 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]
- monkeypatch.setattr("platforms.chatgpt.cpa_upload.upload_to_cpa", lambda *args, **kwargs: (False, "unexpected EOF")) # type: ignore[no-untyped-def]
- loop._worker(thread_id=1, initial_provider="mailtm")
- payload = json.loads(pool_file.read_text(encoding="utf-8"))
- assert payload["backup_written"] is True
- assert payload["cpa_sync_status"] == "failed"
- assert payload["last_cpa_sync_error"] == "unexpected EOF"
- snapshot = loop.snapshot()
- assert snapshot["total_success"] == 0
- assert snapshot["total_failure"] == 1
- assert snapshot["total_cpa_sync_success"] == 0
- assert snapshot["total_cpa_sync_failure"] == 1
- assert snapshot["failure_by_stage"]["cpa_sync"] == 1
- assert snapshot["failure_signals"]["cpa_sync_failed"] == 1
- def test_registration_loop_holds_add_phone_gated_success_in_warmup(monkeypatch, tmp_path: Path) -> None:
- settings = _base_settings(
- register_mail_provider="cfmail",
- )
- loop = RegistrationLoop(settings)
- loop._providers = ["cfmail"]
- pool_file = tmp_path / "fresh@example.com.json"
- pool_file.write_text(json.dumps({"email": "fresh@example.com", "access_token": "tok", "account_id": "acct"}), encoding="utf-8")
- uploaded: list[str] = []
- def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
- del kwargs
- loop._stop_event.set()
- return {
- "success": True,
- "email": "fresh@example.com",
- "pool_file": str(pool_file),
- "written_to_pool": True,
- "metadata": {
- "mail_provider": "cfmail",
- "email_domain": "demo.example.test",
- "post_create_gate": "add_phone",
- },
- }
- monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
- monkeypatch.setattr("core.registration.get_management_key", lambda: "secret") # type: ignore[no-untyped-def]
- 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]
- monkeypatch.setattr(
- "platforms.chatgpt.cpa_upload.upload_to_cpa",
- lambda token_data, api_url=None, api_key=None, proxy=None: uploaded.append(token_data["email"]) or (True, "ok"),
- )
- loop._worker(thread_id=1, initial_provider="cfmail")
- assert uploaded == []
- payload = json.loads(pool_file.read_text(encoding="utf-8"))
- assert payload["warmup_required"] is True
- assert payload["warmup_state"] == "pending"
- assert payload["cpa_sync_status"] == "warmup_pending"
- snapshot = loop.snapshot()
- assert snapshot["total_success"] == 0
- assert snapshot["total_success_registered"] == 0
- assert snapshot["total_cpa_sync_success"] == 0
- assert snapshot["total_cpa_sync_failure"] == 0
- assert snapshot["total_warmup_pending"] == 1
- def test_registration_loop_marks_add_phone_backup_as_warmup_pending(monkeypatch, tmp_path: Path) -> None:
- settings = _base_settings(
- register_mail_provider="cfmail",
- )
- loop = RegistrationLoop(settings)
- loop._providers = ["cfmail"]
- pool_file = tmp_path / "fresh@example.com.json"
- pool_file.write_text(
- json.dumps({"email": "fresh@example.com", "access_token": "tok", "account_id": "acct"}),
- encoding="utf-8",
- )
- def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
- del kwargs
- loop._stop_event.set()
- return {
- "success": True,
- "email": "fresh@example.com",
- "pool_file": str(pool_file),
- "written_to_pool": True,
- "metadata": {
- "mail_provider": "cfmail",
- "email_domain": "demo.example.test",
- "post_create_gate": "add_phone",
- },
- }
- monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
- monkeypatch.setattr("core.registration.get_management_key", lambda: "secret") # type: ignore[no-untyped-def]
- 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]
- 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]
- loop._worker(thread_id=1, initial_provider="cfmail")
- payload = json.loads(pool_file.read_text(encoding="utf-8"))
- assert payload["backup_written"] is True
- assert payload["cpa_sync_status"] == "warmup_pending"
- assert payload["last_cpa_sync_error"] == "warmup pending"
- def test_registration_loop_enqueue_pending_token_preserves_proxy_provenance() -> None:
- loop = RegistrationLoop(_base_settings())
- loop._enqueue_pending_token(
- {
- "metadata": {
- "deferred_credentials": {
- "email": "fresh@example.com",
- "password": "pw-secret",
- "registration_proxy_key": "台湾♣备用-1",
- "registration_proxy_region": "tw",
- "registration_proxy_url": "socks5://127.0.0.1:17891",
- "registration_fingerprint_profile": "chrome120_win",
- "cfmail_profile_name": "cfmail-tw",
- "add_phone_trace_path": "/tmp/add-phone-trace.json",
- }
- }
- },
- thread_id=1,
- )
- assert len(loop._pending_token_queue) == 1
- entry = loop._pending_token_queue[0]
- assert entry["registration_proxy_key"] == "台湾♣备用-1"
- assert entry["registration_proxy_region"] == "tw"
- assert entry["registration_proxy_url"] == "socks5://127.0.0.1:17891"
- assert entry["registration_fingerprint_profile"] == "chrome120_win"
- assert entry["cfmail_profile_name"] == "cfmail-tw"
- assert entry["add_phone_trace_path"] == "/tmp/add-phone-trace.json"
- def test_registration_loop_retry_pending_token_reuses_registration_proxy(monkeypatch, tmp_path: Path) -> None:
- settings = _base_settings(pool_dir=tmp_path)
- loop = RegistrationLoop(settings)
- recorded: dict[str, object] = {}
- class FakeLease:
- name = "台湾♣备用-1"
- proxy_url = "socks5://127.0.0.1:17891"
- class FakePool:
- def acquire(self, timeout=5.0, preferred_name=None, preferred_regions=()): # type: ignore[no-untyped-def]
- recorded["preferred_name"] = preferred_name
- recorded["preferred_regions"] = tuple(preferred_regions)
- return FakeLease()
- def release(self, lease, *, success, stage=None): # type: ignore[no-untyped-def]
- recorded["released"] = (lease.proxy_url, success, stage)
- class FakeMailbox:
- def __init__(self, manager=None): # type: ignore[no-untyped-def]
- self.manager = manager
- class FakeAdapter:
- def __init__(self, mailbox): # type: ignore[no-untyped-def]
- self.mailbox = mailbox
- self._account = None
- class FakeEngine:
- def __init__(self, email_service, proxy_url=None): # type: ignore[no-untyped-def]
- recorded["engine_proxy_url"] = proxy_url
- self.email_service = email_service
- self.email = ""
- self.password = ""
- def _login_for_token(self): # type: ignore[no-untyped-def]
- return {
- "access_token": "access-token",
- "refresh_token": "refresh-token",
- "id_token": "",
- "account_id": "acct-123",
- "last_refresh": "2026-03-31T00:00:00Z",
- "expired": "2026-04-01T00:00:00Z",
- }
- def fake_write_token_record(token_data, pool_dir): # type: ignore[no-untyped-def]
- recorded["token_data"] = dict(token_data)
- path = Path(pool_dir) / "fresh@example.com.json"
- path.write_text(json.dumps(token_data), encoding="utf-8")
- return path
- loop._proxy_pool = FakePool()
- monkeypatch.setattr("core.registration.CfMailMailbox", FakeMailbox, raising=False)
- monkeypatch.setattr("core.registration.MailboxEmailServiceAdapter", FakeAdapter, raising=False)
- monkeypatch.setattr("platforms.chatgpt.register.RegistrationEngine", FakeEngine)
- monkeypatch.setattr("platforms.chatgpt.pool.write_token_record", fake_write_token_record)
- monkeypatch.setattr(loop, "_sync_cpa_from_success", lambda result, thread_id: (True, "", "fresh@example.com"))
- entry = {
- "email": "fresh@example.com",
- "password": "pw-secret",
- "mailbox_jwt": "",
- "mailbox_extra": {},
- "registration_proxy_key": "台湾♣备用-1",
- "registration_proxy_region": "tw",
- "registration_proxy_url": "socks5://127.0.0.1:17891",
- "registration_fingerprint_profile": "chrome120_win",
- "cfmail_profile_name": "cfmail-tw",
- "retry_count": 0,
- }
- loop._retry_pending_token(entry)
- assert recorded["preferred_name"] == "台湾♣备用-1"
- assert recorded["preferred_regions"] == ("tw",)
- assert recorded["engine_proxy_url"] == "socks5://127.0.0.1:17891"
- assert recorded["released"] == ("socks5://127.0.0.1:17891", True, "deferred_retry")
- token_data = recorded["token_data"]
- assert token_data["registration_proxy_key"] == "台湾♣备用-1"
- assert token_data["registration_proxy_region"] == "tw"
- assert token_data["registration_proxy_url"] == "socks5://127.0.0.1:17891"
- assert token_data["registration_fingerprint_profile"] == "chrome120_win"
- assert token_data["registration_cfmail_profile_name"] == "cfmail-tw"
- def test_registration_loop_dumps_engine_logs_on_failure(monkeypatch) -> None:
- settings = _base_settings(register_max_consecutive_failures=1, register_mail_provider="mailtm")
- loop = RegistrationLoop(settings)
- loop._providers = ["mailtm"]
- logged: list[str] = []
- original_log = loop._log
- def capture_log(msg: str) -> None:
- logged.append(msg)
- loop._log = capture_log
- def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
- loop._stop_event.set()
- return {
- "success": False,
- "stage": "signup",
- "error_message": "HTTP 403: access denied",
- "logs": [
- "[09:35:20] check_ip_location: JP",
- "[09:35:21] created mailbox: test@example.com",
- "[09:35:22] signup form status: 403",
- ],
- }
- monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
- loop._worker(thread_id=1, initial_provider="mailtm")
- # Verify stage appears in the failure line
- fail_lines = [line for line in logged if "failed" in line and "stage=" in line]
- assert len(fail_lines) == 1
- assert "[stage=signup]" in fail_lines[0]
- assert "HTTP 403: access denied" in fail_lines[0]
- # Verify engine logs are dumped with ↳ prefix
- engine_lines = [line for line in logged if "\u21b3" in line]
- assert len(engine_lines) == 3
- assert "signup form status: 403" in engine_lines[2]
- snapshot = loop.snapshot()
- assert snapshot["failure_by_stage"]["signup"] == 1
- assert snapshot["recent_failure_hotspots"] == [{"key": "signup", "stage": "signup", "count": 1}]
- def test_registration_loop_rotates_cfmail_domain_after_blacklist_threshold(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail", register_max_consecutive_failures=5)
- loop = RegistrationLoop(settings)
- loop._providers = ["cfmail"]
- loop._cfmail_tracker = DomainHealthTracker(window_size=2, blacklist_threshold=2, rotation_cooldown_seconds=1)
- monkeypatch.setattr(loop, "_ensure_cfmail_domain_pool_target", lambda **_: None)
- rotation_calls: list[str] = []
- class FakeProvisioner:
- def rotate_active_domain(self): # type: ignore[no-untyped-def]
- rotation_calls.append("rotate")
- return ProvisionResult(
- success=True,
- step="completed",
- old_domain="nova.example.test",
- new_domain="auto0322.example.test",
- )
- def provision_additional_domain(self): # type: ignore[no-untyped-def]
- return ProvisionResult(success=False, step="provision_additional_domain", error="not needed")
- loop._cfmail_provisioner = FakeProvisioner()
- attempts = {"count": 0}
- def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
- del kwargs
- attempts["count"] += 1
- if attempts["count"] >= 2:
- loop._stop_event.set()
- return {
- "success": False,
- "stage": "create_account",
- "error_message": "create account failed",
- "metadata": {
- "mail_provider": "cfmail",
- "email_domain": "nova.example.test",
- "create_account_error_code": "registration_disallowed",
- "create_account_error_message": "blocked",
- },
- }
- monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
- loop._worker(thread_id=1, initial_provider="cfmail")
- assert rotation_calls == ["rotate"]
- assert loop.snapshot()["cfmail_rotation"]["last_new_domain"] == "auto0322.example.test"
- def test_registration_loop_does_not_rotate_cfmail_on_mailbox_failure(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail", register_max_consecutive_failures=2)
- loop = RegistrationLoop(settings)
- loop._providers = ["cfmail"]
- loop._cfmail_tracker = DomainHealthTracker(window_size=2, blacklist_threshold=2, rotation_cooldown_seconds=1)
- monkeypatch.setattr(loop, "_ensure_cfmail_domain_pool_target", lambda **_: None)
- class FakeProvisioner:
- def rotate_active_domain(self): # type: ignore[no-untyped-def]
- raise AssertionError("rotation should not be called")
- def provision_additional_domain(self): # type: ignore[no-untyped-def]
- return ProvisionResult(success=False, step="provision_additional_domain", error="not needed")
- loop._cfmail_provisioner = FakeProvisioner()
- def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
- del kwargs
- loop._stop_event.set()
- return {
- "success": False,
- "stage": "mailbox",
- "error_message": "create email failed",
- "metadata": {"mail_provider": "cfmail"},
- }
- monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
- loop._worker(thread_id=1, initial_provider="cfmail")
- assert loop.snapshot()["cfmail_rotation"]["last_new_domain"] == ""
- def test_registration_loop_rotates_cfmail_on_invalid_domain_mailbox_failure() -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- loop._providers = ["cfmail"]
- loop._cfmail_tracker = DomainHealthTracker(window_size=2, blacklist_threshold=2, rotation_cooldown_seconds=1)
- class FakeProvisioner:
- def rotate_active_domain(self): # type: ignore[no-untyped-def]
- return ProvisionResult(
- success=True,
- step="rotate",
- old_domain="bad.example.test",
- new_domain="auto-new.example.test",
- )
- reload_calls: list[bool] = []
- class FakeManager:
- def reload_if_needed(self, force=False): # type: ignore[no-untyped-def]
- reload_calls.append(bool(force))
- return True
- loop._cfmail_manager = FakeManager()
- loop._cfmail_provisioner = FakeProvisioner()
- loop._cfmail_wait_otp_cooldown_seconds = 60
- loop._cfmail_add_phone_cooldown_seconds = 60
- result = {
- "success": False,
- "stage": "mailbox",
- "error_message": "create email failed",
- "logs": ["[10:00:00] create_email failed: 创建邮箱地址失败: 无效的域名"],
- "metadata": {
- "mail_provider": "cfmail",
- "email_domain": "bad.example.test",
- },
- }
- assert loop._force_rotate_cfmail_for_invalid_mailbox(thread_id=1, result=result) is True
- snapshot = loop.snapshot()
- assert snapshot["cfmail_rotation"]["last_new_domain"] == "auto-new.example.test"
- assert snapshot["cfmail_wait_otp_stoploss"]["active_domain"] == "auto-new.example.test"
- assert reload_calls == [True]
- def test_registration_loop_startup_reports_active_domain_pool_ready(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- logs: list[str] = []
- class FakeManager:
- accounts = [
- type("Account", (), {"name": "cfmail-a", "email_domain": "a.example.test"})(),
- type("Account", (), {"name": "cfmail-b", "email_domain": "b.example.test"})(),
- ]
- def skip_remaining_seconds(self, account_name): # type: ignore[no-untyped-def]
- del account_name
- return 0
- loop._cfmail_manager = FakeManager()
- loop._cfmail_provisioner = object()
- monkeypatch.setattr(loop, "_log", logs.append)
- assert loop._ensure_cfmail_active_domain_ready() is True
- assert any("active-domain pool ready" in line for line in logs)
- def test_registration_loop_forces_cfmail_rotation_when_all_accounts_in_cooldown(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail", register_max_consecutive_failures=2)
- loop = RegistrationLoop(settings)
- loop._providers = ["cfmail"]
- loop._cfmail_tracker = DomainHealthTracker(window_size=2, blacklist_threshold=2, rotation_cooldown_seconds=1)
- run_calls: list[str] = []
- rotation_calls: list[str] = []
- class FakeManager:
- def reload_if_needed(self) -> bool:
- return False
- def select_account(self, profile_name=None): # type: ignore[no-untyped-def]
- del profile_name
- return None
- class FakeProvisioner:
- def rotate_active_domain(self): # type: ignore[no-untyped-def]
- rotation_calls.append("rotate")
- loop._stop_event.set()
- return ProvisionResult(
- success=True,
- step="completed",
- old_domain="auto-old.example.test",
- new_domain="auto-new.example.test",
- )
- def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
- del kwargs
- loop._stop_event.set()
- run_calls.append("called")
- raise AssertionError("registration should not run while cfmail is fully unavailable")
- loop._cfmail_manager = FakeManager()
- loop._cfmail_provisioner = FakeProvisioner()
- monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
- loop._worker(thread_id=1, initial_provider="cfmail")
- assert rotation_calls == ["rotate"]
- assert run_calls == []
- assert loop.snapshot()["cfmail_rotation"]["last_new_domain"] == "auto-new.example.test"
- def test_registration_loop_does_not_penalize_proxy_for_blacklist_failure(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- loop._providers = ["cfmail"]
- recorded: dict[str, object] = {}
- class FakeLease:
- name = "sg-node"
- local_port = 17891
- proxy_url = "socks5://127.0.0.1:17891"
- class FakePool:
- def acquire(self, timeout=5.0, preferred_name=None, preferred_regions=()): # type: ignore[no-untyped-def]
- return FakeLease()
- def release(self, lease, *, success, stage=None): # type: ignore[no-untyped-def]
- recorded["released"] = (success, stage)
- def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
- del kwargs
- loop._stop_event.set()
- return {
- "success": False,
- "stage": "create_account",
- "error_message": "create account failed",
- "metadata": {
- "mail_provider": "cfmail",
- "email_domain": "nova.example.test",
- "create_account_error_code": "registration_disallowed",
- },
- }
- loop._proxy_pool = FakePool()
- monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
- loop._worker(thread_id=1, initial_provider="cfmail")
- assert recorded["released"] == (None, "create_account")
- def test_registration_loop_records_add_phone_gate_signal(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- loop._providers = ["cfmail"]
- def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
- del kwargs
- loop._stop_event.set()
- return {
- "success": False,
- "stage": "add_phone_gate",
- "error_message": "post-create flow requires phone gate",
- "metadata": {
- "mail_provider": "cfmail",
- "email_domain": "demo.example.test",
- "post_create_gate": "add_phone",
- },
- }
- monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
- loop._worker(thread_id=1, initial_provider="cfmail")
- snapshot = loop.snapshot()
- assert snapshot["failure_by_stage"]["add_phone_gate"] == 1
- assert snapshot["failure_signals"]["add_phone_gate"] == 1
- assert snapshot["recent_failure_hotspots"] == [{"key": "add_phone_gate", "stage": "add_phone_gate", "count": 1}]
- assert snapshot["recent_attempts"][0]["post_create_gate"] == "add_phone"
- def test_registration_loop_records_signup_invalid_auth_step_signal() -> None:
- loop = RegistrationLoop(_base_settings(register_mail_provider="cfmail"))
- loop._record_attempt(
- success=False,
- stage="signup",
- error_message='HTTP 400: {"error":{"code":"invalid_auth_step"}}',
- metadata={
- "mail_provider": "cfmail",
- "email_domain": "demo.example.test",
- "signup_error_code": "invalid_auth_step",
- "signup_http_status": 400,
- },
- proxy_key="台湾原生-01",
- email="demo@example.test",
- )
- snapshot = loop.snapshot()
- assert snapshot["failure_signals"]["invalid_auth_step"] == 1
- attempt = snapshot["recent_attempts"][0]
- assert attempt["signup_error_code"] == "invalid_auth_step"
- assert attempt["signup_http_status"] == 400
- def test_registration_loop_classifies_mailbox_transport_failures() -> None:
- loop = RegistrationLoop(_base_settings(register_mail_provider="cfmail"))
- loop._record_attempt(
- success=False,
- stage="mailbox",
- error_message="create email failed",
- metadata={
- "mail_provider": "cfmail",
- "email_domain": "demo.example.test",
- "mailbox_error_kind": "transport_error",
- "mailbox_error_stage": "create_email",
- },
- proxy_key="新加坡原生-02",
- )
- snapshot = loop.snapshot()
- assert snapshot["failure_signals"]["mailbox_create_transport_error"] == 1
- attempt = snapshot["recent_attempts"][0]
- assert attempt["mailbox_error_kind"] == "transport_error"
- assert attempt["mailbox_error_stage"] == "create_email"
- def test_registration_loop_classifies_user_already_exists_as_mailbox_reused(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- loop._providers = ["cfmail"]
- def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
- del kwargs
- loop._stop_event.set()
- return {
- "success": False,
- "stage": "create_account",
- "error_message": "create account failed",
- "email": "dup@example.com",
- "metadata": {
- "mail_provider": "cfmail",
- "email_domain": "demo.example.test",
- "create_account_error_code": "user_already_exists",
- },
- }
- monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
- loop._worker(thread_id=1, initial_provider="cfmail")
- snapshot = loop.snapshot()
- assert snapshot["failure_by_stage"]["create_account"] == 1
- assert snapshot["failure_signals"]["mailbox_reused"] == 1
- assert snapshot["recent_failure_hotspots"] == [{"key": "mailbox_reused", "stage": "create_account", "count": 1}]
- def test_registration_loop_activates_add_phone_stoploss(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- loop._providers = ["cfmail"]
- monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "demo.example.test")
- loop._cfmail_add_phone_window = 2
- loop._cfmail_add_phone_threshold = 2
- loop._cfmail_add_phone_max_successes = 0
- loop._cfmail_add_phone_cooldown_seconds = 60
- attempts = {"count": 0}
- def fake_run_chatgpt_register_once(**kwargs): # type: ignore[no-untyped-def]
- del kwargs
- attempts["count"] += 1
- if attempts["count"] >= 2:
- loop._stop_event.set()
- return {
- "success": False,
- "stage": "add_phone_gate",
- "error_message": "post-create flow requires phone gate",
- "metadata": {
- "mail_provider": "cfmail",
- "email_domain": "demo.example.test",
- "post_create_gate": "add_phone",
- },
- }
- monkeypatch.setattr("main.run_chatgpt_register_once", fake_run_chatgpt_register_once)
- loop._worker(thread_id=1, initial_provider="cfmail")
- snapshot = loop.snapshot()
- stoploss = snapshot["cfmail_add_phone_stoploss"]
- assert stoploss["active_domain"] == "demo.example.test"
- assert stoploss["in_cooldown"] is True
- assert stoploss["last_add_phone_failures"] == 2
- assert stoploss["last_successes"] == 0
- def test_registration_loop_snapshot_exposes_active_domain_only_attempts() -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- loop._recent_attempts.extend(
- [
- {
- "timestamp": "2026-03-24T18:00:00",
- "success": False,
- "stage": "create_account",
- "signal": "registration_disallowed",
- "error_message": "failed",
- "email_domain": "old.example.test",
- "proxy_key": "old-node",
- "email": "old@old.example.test",
- },
- {
- "timestamp": "2026-03-24T18:00:10",
- "success": True,
- "stage": "completed",
- "signal": "",
- "error_message": "",
- "email_domain": "new.example.test",
- "proxy_key": "new-node-1",
- "email": "ok@new.example.test",
- },
- {
- "timestamp": "2026-03-24T18:00:20",
- "success": False,
- "stage": "create_account",
- "signal": "registration_disallowed",
- "error_message": "failed",
- "email_domain": "new.example.test",
- "proxy_key": "new-node-2",
- "email": "bad@new.example.test",
- },
- ]
- )
- class FakeTracker:
- def snapshot(self) -> dict[str, object]:
- return {
- "active_domain": "new.example.test",
- "last_new_domain": "new.example.test",
- "last_blacklisted_domain": "old.example.test",
- }
- loop._cfmail_tracker = FakeTracker()
- snapshot = loop.snapshot()
- assert [item["email_domain"] for item in snapshot["active_domain_recent_attempts"]] == [
- "new.example.test",
- "new.example.test",
- ]
- assert snapshot["active_domain_failure_by_stage"] == {"create_account": 1}
- assert snapshot["active_domain_failure_signals"] == {"registration_disallowed": 1}
- assert snapshot["active_domain_recent_failure_hotspots"] == [
- {"key": "registration_disallowed", "stage": "create_account", "count": 1}
- ]
- def test_registration_loop_snapshot_infers_active_domain_from_recent_attempts() -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- loop._recent_attempts.extend(
- [
- {
- "timestamp": "2026-03-24T18:10:00",
- "success": False,
- "stage": "create_account",
- "signal": "registration_disallowed",
- "error_message": "failed",
- "email_domain": "old.example.test",
- "proxy_key": "old-node",
- "email": "old@old.example.test",
- },
- {
- "timestamp": "2026-03-24T18:10:10",
- "success": True,
- "stage": "completed",
- "signal": "",
- "error_message": "",
- "email_domain": "new.example.test",
- "proxy_key": "new-node",
- "email": "ok@new.example.test",
- },
- ]
- )
- snapshot = loop.snapshot()
- assert [item["email_domain"] for item in snapshot["active_domain_recent_attempts"]] == [
- "new.example.test"
- ]
- def test_registration_loop_disables_add_phone_cooldown_when_configured_zero(monkeypatch) -> None:
- monkeypatch.setenv("ZHUCE6_CFMAIL_ADD_PHONE_COOLDOWN_SECONDS", "0")
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- loop._cfmail_add_phone_window = 2
- loop._cfmail_add_phone_threshold = 2
- loop._cfmail_add_phone_max_successes = 0
- result = {
- "success": False,
- "stage": "add_phone_gate",
- "error_message": "post-create flow requires phone gate",
- "metadata": {
- "mail_provider": "cfmail",
- "email_domain": "demo.example.test",
- "post_create_gate": "add_phone",
- },
- }
- loop._update_cfmail_add_phone_stoploss(result)
- loop._update_cfmail_add_phone_stoploss(result)
- stoploss = loop.snapshot()["cfmail_add_phone_stoploss"]
- assert stoploss["in_cooldown"] is False
- assert stoploss["cooldown_remaining_seconds"] == 0
- def test_registration_loop_add_phone_stoploss_retires_triggered_domain_without_blocking_other_domains(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- loop._cfmail_add_phone_window = 2
- loop._cfmail_add_phone_threshold = 2
- loop._cfmail_add_phone_max_successes = 0
- loop._cfmail_add_phone_cooldown_seconds = 60
- monkeypatch.setattr(
- loop,
- "_current_cfmail_active_accounts",
- lambda: [
- {"name": "cfmail-tw", "domain": "tw.example.test"},
- {"name": "cfmail-sg", "domain": "sg.example.test"},
- ],
- )
- retired: list[str] = []
- scheduled: list[str] = []
- class FakeProvisioner:
- def retire_domain(self, domain: str) -> ProvisionResult:
- retired.append(domain)
- return ProvisionResult(success=True, step="retire_domain", old_domain=domain)
- loop._cfmail_provisioner = FakeProvisioner() # type: ignore[assignment]
- loop._cfmail_tracker = DomainHealthTracker(window_size=2, blacklist_threshold=2, rotation_cooldown_seconds=1)
- monkeypatch.setattr(loop, "_reload_cfmail_manager_after_rotation", lambda: None)
- monkeypatch.setattr(
- loop,
- "_schedule_cfmail_domain_pool_replenish",
- lambda *, trigger_thread_id, reason: scheduled.append(f"{trigger_thread_id}:{reason}"),
- raising=False,
- )
- result = {
- "success": False,
- "stage": "add_phone_gate",
- "error_message": "post-create flow requires phone gate",
- "metadata": {
- "mail_provider": "cfmail",
- "email_domain": "sg.example.test",
- "post_create_gate": "add_phone",
- },
- }
- loop._update_cfmail_add_phone_stoploss(result)
- loop._update_cfmail_add_phone_stoploss(result)
- assert loop._wait_if_cfmail_add_phone_stopped(thread_id=3, provider="cfmail") is False
- assert retired == ["sg.example.test"]
- assert scheduled == ["3:add_phone stoploss"]
- def test_registration_loop_activates_wait_otp_stoploss_for_no_message_timeouts(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "demo.example.test")
- loop._cfmail_wait_otp_window = 2
- loop._cfmail_wait_otp_threshold = 2
- loop._cfmail_wait_otp_cooldown_seconds = 60
- result = {
- "success": False,
- "stage": "wait_otp",
- "error_message": "otp retrieval failed",
- "metadata": {
- "mail_provider": "cfmail",
- "email_domain": "demo.example.test",
- "otp_wait_failure_reason": "mailbox_timeout_no_message",
- "otp_mailbox_message_scan_count": 0,
- },
- }
- loop._update_cfmail_wait_otp_stoploss(result)
- loop._update_cfmail_wait_otp_stoploss(result)
- stoploss = loop.snapshot()["cfmail_wait_otp_stoploss"]
- assert stoploss["active_domain"] == "demo.example.test"
- assert stoploss["in_cooldown"] is True
- assert stoploss["last_no_message_timeouts"] == 2
- def test_registration_loop_wait_otp_stoploss_retires_triggered_domain_without_blocking_other_domains(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- loop._cfmail_wait_otp_window = 2
- loop._cfmail_wait_otp_threshold = 2
- loop._cfmail_wait_otp_cooldown_seconds = 60
- monkeypatch.setattr(
- loop,
- "_current_cfmail_active_accounts",
- lambda: [
- {"name": "cfmail-tw", "domain": "tw.example.test"},
- {"name": "cfmail-sg", "domain": "sg.example.test"},
- ],
- )
- retired: list[str] = []
- scheduled: list[str] = []
- class FakeProvisioner:
- def retire_domain(self, domain: str) -> ProvisionResult:
- retired.append(domain)
- return ProvisionResult(success=True, step="retire_domain", old_domain=domain)
- loop._cfmail_provisioner = FakeProvisioner() # type: ignore[assignment]
- loop._cfmail_tracker = DomainHealthTracker(window_size=2, blacklist_threshold=2, rotation_cooldown_seconds=1)
- monkeypatch.setattr(loop, "_reload_cfmail_manager_after_rotation", lambda: None)
- monkeypatch.setattr(
- loop,
- "_schedule_cfmail_domain_pool_replenish",
- lambda *, trigger_thread_id, reason: scheduled.append(f"{trigger_thread_id}:{reason}"),
- raising=False,
- )
- result = {
- "success": False,
- "stage": "wait_otp",
- "error_message": "otp retrieval failed",
- "metadata": {
- "mail_provider": "cfmail",
- "email_domain": "sg.example.test",
- "otp_wait_failure_reason": "mailbox_timeout_no_message",
- "otp_mailbox_message_scan_count": 0,
- },
- }
- loop._update_cfmail_wait_otp_stoploss(result)
- loop._update_cfmail_wait_otp_stoploss(result)
- assert loop._wait_if_cfmail_wait_otp_stopped(thread_id=4, provider="cfmail") is False
- assert retired == ["sg.example.test"]
- assert scheduled == ["4:wait_otp stoploss"]
- def test_registration_loop_replenish_worker_retries_transient_provision_failure(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- loop._cfmail_active_domain_count = 3
- monkeypatch.setattr("core.registration.time.sleep", lambda _: None)
- active_accounts = [
- {"name": "cfmail-tw", "domain": "tw.example.test"},
- {"name": "cfmail-sg", "domain": "sg.example.test"},
- ]
- monkeypatch.setattr(loop, "_current_cfmail_active_accounts", lambda: list(active_accounts))
- monkeypatch.setattr(loop, "_reload_cfmail_manager_after_rotation", lambda: None)
- attempts = {"count": 0}
- class FakeProvisioner:
- def provision_additional_domain(self) -> ProvisionResult:
- attempts["count"] += 1
- if attempts["count"] == 1:
- return ProvisionResult(success=False, step="provision_additional_domain", error="tls connect error")
- active_accounts.append({"name": "cfmail-jp", "domain": "jp.example.test"})
- return ProvisionResult(success=True, step="provision_additional_domain", new_domain="jp.example.test")
- loop._cfmail_provisioner = FakeProvisioner() # type: ignore[assignment]
- loop._cfmail_tracker = DomainHealthTracker(window_size=2, blacklist_threshold=2, rotation_cooldown_seconds=1)
- loop._cfmail_domain_pool_replenish_worker(trigger_thread_id=2, reason="fresh domain budget reached")
- assert attempts["count"] == 2
- assert len(active_accounts) == 3
- def test_registration_loop_schedules_replenish_when_usable_domain_pool_below_target(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- loop._cfmail_active_domain_count = 3
- monkeypatch.setattr(
- loop,
- "_current_cfmail_active_accounts",
- lambda: [
- {"name": "cfmail-tw", "domain": "tw.example.test"},
- {"name": "cfmail-sg", "domain": "sg.example.test"},
- ],
- )
- loop._cfmail_provisioner = object() # type: ignore[assignment]
- scheduled: list[str] = []
- monkeypatch.setattr(
- loop,
- "_schedule_cfmail_domain_pool_replenish",
- lambda *, trigger_thread_id, reason: scheduled.append(f"{trigger_thread_id}:{reason}"),
- raising=False,
- )
- loop._ensure_cfmail_domain_pool_target(trigger_thread_id=5, reason="startup")
- assert scheduled == ["5:startup"]
- def test_registration_loop_does_not_activate_wait_otp_stoploss_when_window_contains_success(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "demo.example.test")
- loop._cfmail_wait_otp_window = 2
- loop._cfmail_wait_otp_threshold = 2
- loop._cfmail_wait_otp_cooldown_seconds = 60
- timeout_result = {
- "success": False,
- "stage": "wait_otp",
- "error_message": "otp retrieval failed",
- "metadata": {
- "mail_provider": "cfmail",
- "email_domain": "demo.example.test",
- "otp_wait_failure_reason": "mailbox_timeout_no_message",
- "otp_mailbox_message_scan_count": 0,
- },
- }
- success_result = {
- "success": True,
- "stage": "completed",
- "email": "ok@demo.example.test",
- "metadata": {
- "mail_provider": "cfmail",
- "email_domain": "demo.example.test",
- },
- }
- loop._update_cfmail_wait_otp_stoploss(timeout_result)
- loop._update_cfmail_wait_otp_stoploss(success_result)
- stoploss = loop.snapshot()["cfmail_wait_otp_stoploss"]
- assert stoploss["active_domain"] == "demo.example.test"
- assert stoploss["in_cooldown"] is False
- def test_registration_loop_does_not_activate_wait_otp_stoploss_when_window_contains_message_seen(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "demo.example.test")
- loop._cfmail_wait_otp_window = 2
- loop._cfmail_wait_otp_threshold = 2
- loop._cfmail_wait_otp_cooldown_seconds = 60
- timeout_result = {
- "success": False,
- "stage": "wait_otp",
- "error_message": "otp retrieval failed",
- "metadata": {
- "mail_provider": "cfmail",
- "email_domain": "demo.example.test",
- "otp_wait_failure_reason": "mailbox_timeout_no_message",
- "otp_mailbox_message_scan_count": 0,
- },
- }
- message_seen_result = {
- "success": False,
- "stage": "wait_otp",
- "error_message": "otp retrieval failed",
- "metadata": {
- "mail_provider": "cfmail",
- "email_domain": "demo.example.test",
- "otp_wait_failure_reason": "mailbox_timeout_no_match",
- "otp_mailbox_message_scan_count": 1,
- },
- }
- loop._update_cfmail_wait_otp_stoploss(timeout_result)
- loop._update_cfmail_wait_otp_stoploss(message_seen_result)
- stoploss = loop.snapshot()["cfmail_wait_otp_stoploss"]
- assert stoploss["active_domain"] == "demo.example.test"
- assert stoploss["in_cooldown"] is False
- def test_registration_loop_disables_cfmail_canary_gate() -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- assert loop._wait_if_cfmail_canary_pending(thread_id=2, provider="cfmail") is False
- assert loop._wait_if_cfmail_canary_pending(thread_id=1, provider="cfmail") is False
- def test_registration_loop_cfmail_canary_snapshot_is_disabled(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- loop._arm_cfmail_canary("demo.example.test")
- loop._mark_cfmail_canary_ready("demo.example.test", reason="mailbox_message_seen")
- loop._update_cfmail_canary_after_result(thread_id=3, result={})
- snapshot = loop.snapshot()["cfmail_canary"]
- assert snapshot["pending"] is False
- assert snapshot["active_domain"] == ""
- assert snapshot["last_ready_reason"] == "disabled"
- def test_registration_loop_tracks_fresh_domain_budget_after_mail_seen(monkeypatch) -> None:
- monkeypatch.setenv("ZHUCE6_CFMAIL_FRESH_DOMAIN_ATTEMPT_BUDGET", "2")
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "demo.example.test")
- first_result = {
- "success": False,
- "stage": "add_phone_gate",
- "error_message": "post-create flow requires phone gate",
- "metadata": {
- "mail_provider": "cfmail",
- "email_domain": "demo.example.test",
- "otp_mailbox_message_scan_count": 1,
- },
- }
- second_result = {
- "success": False,
- "stage": "wait_otp",
- "error_message": "otp retrieval failed",
- "metadata": {
- "mail_provider": "cfmail",
- "email_domain": "demo.example.test",
- "otp_mailbox_message_scan_count": 0,
- },
- }
- loop._update_cfmail_fresh_domain_budget(first_result)
- snapshot = loop.snapshot()["cfmail_fresh_domain_budget"]
- assert snapshot["completed_attempts"] == 1
- assert snapshot["mail_seen_attempts"] == 1
- assert snapshot["last_triggered_at"] == ""
- loop._update_cfmail_fresh_domain_budget(second_result)
- snapshot = loop.snapshot()["cfmail_fresh_domain_budget"]
- assert snapshot["completed_attempts"] == 2
- assert snapshot["mail_seen_attempts"] == 1
- assert snapshot["last_reason"] == "fresh_domain_attempt_budget_reached"
- def test_registration_loop_rotates_domain_when_fresh_domain_budget_reached(monkeypatch) -> None:
- monkeypatch.setenv("ZHUCE6_CFMAIL_FRESH_DOMAIN_ATTEMPT_BUDGET", "2")
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "demo.example.test")
- loop._cfmail_tracker = DomainHealthTracker(window_size=2, blacklist_threshold=2, rotation_cooldown_seconds=1)
- loop._cfmail_fresh_domain_state.update(
- {
- "active_domain": "demo.example.test",
- "completed_attempts": 2,
- "mail_seen_attempts": 1,
- "successes": 0,
- "last_triggered_at": "2026-03-29T10:00:00",
- "last_rotation_attempted_at": "",
- "last_reason": "fresh_domain_attempt_budget_reached",
- }
- )
- class FakeProvisioner:
- def rotate_active_domain(self): # type: ignore[no-untyped-def]
- return ProvisionResult(
- success=True,
- step="rotate",
- old_domain="demo.example.test",
- new_domain="auto-new.example.test",
- )
- reload_calls: list[bool] = []
- class FakeManager:
- def reload_if_needed(self, force=False): # type: ignore[no-untyped-def]
- reload_calls.append(bool(force))
- return True
- loop._cfmail_manager = FakeManager()
- loop._cfmail_provisioner = FakeProvisioner()
- assert loop._rotate_cfmail_for_fresh_domain_budget(thread_id=5) is True
- snapshot = loop.snapshot()
- assert snapshot["cfmail_rotation"]["last_new_domain"] == "auto-new.example.test"
- assert snapshot["cfmail_fresh_domain_budget"]["active_domain"] == "auto-new.example.test"
- assert snapshot["cfmail_fresh_domain_budget"]["completed_attempts"] == 0
- assert snapshot["cfmail_canary"]["pending"] is False
- assert reload_calls == [True]
- def test_registration_loop_throttles_cfmail_inflight_attempts(monkeypatch) -> None:
- monkeypatch.setenv("ZHUCE6_CFMAIL_MAX_INFLIGHT", "1")
- monkeypatch.setenv("ZHUCE6_CFMAIL_START_INTERVAL_SECONDS", "0")
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- monkeypatch.setattr(
- loop,
- "_current_cfmail_active_accounts",
- lambda: [{"name": "cfmail-demo", "domain": "demo.example.test"}],
- )
- assert loop._wait_if_cfmail_flow_throttled(thread_id=1, provider="cfmail") is False
- assert loop._wait_if_cfmail_flow_throttled(thread_id=2, provider="cfmail") is True
- loop._release_cfmail_flow_slot(thread_id=1)
- assert loop._wait_if_cfmail_flow_throttled(thread_id=2, provider="cfmail") is False
- def test_registration_loop_throttles_cfmail_start_interval(monkeypatch) -> None:
- monkeypatch.setenv("ZHUCE6_CFMAIL_MAX_INFLIGHT", "2")
- monkeypatch.setenv("ZHUCE6_CFMAIL_START_INTERVAL_SECONDS", "15")
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- monkeypatch.setattr(
- loop,
- "_current_cfmail_active_accounts",
- lambda: [{"name": "cfmail-demo", "domain": "demo.example.test"}],
- )
- current_time = {"value": 100.0}
- monkeypatch.setattr("core.registration.time.time", lambda: current_time["value"])
- assert loop._wait_if_cfmail_flow_throttled(thread_id=1, provider="cfmail") is False
- loop._release_cfmail_flow_slot(thread_id=1)
- current_time["value"] = 105.0
- assert loop._wait_if_cfmail_flow_throttled(thread_id=2, provider="cfmail") is True
- current_time["value"] = 116.0
- assert loop._wait_if_cfmail_flow_throttled(thread_id=2, provider="cfmail") is False
- def test_registration_loop_activates_live_wait_otp_stoploss_on_stalled_waits(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "demo.example.test")
- loop._cfmail_wait_otp_cooldown_seconds = 60
- loop._cfmail_wait_otp_live_threshold = 2
- loop._cfmail_wait_otp_live_age_seconds = 30
- account_a = MailboxAccount(
- email="a@demo.example.test",
- account_id="jwt-a",
- extra={"email_domain": "demo.example.test"},
- )
- account_b = MailboxAccount(
- email="b@demo.example.test",
- account_id="jwt-b",
- extra={"email_domain": "demo.example.test"},
- )
- loop._on_cfmail_wait_progress(
- account_a,
- {"message_scan_count": 0, "elapsed_seconds": 35},
- )
- loop._on_cfmail_wait_progress(
- account_b,
- {"message_scan_count": 0, "elapsed_seconds": 36},
- )
- stoploss = loop.snapshot()["cfmail_wait_otp_stoploss"]
- assert stoploss["active_domain"] == "demo.example.test"
- assert stoploss["in_cooldown"] is True
- assert stoploss["last_reason"] == "live wait_otp no-message threshold reached"
- assert stoploss["last_no_message_timeouts"] == 2
- def test_registration_loop_live_wait_otp_stoploss_can_be_disabled(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "demo.example.test")
- loop._cfmail_wait_otp_cooldown_seconds = 60
- loop._cfmail_wait_otp_live_threshold = 0
- loop._cfmail_wait_otp_live_age_seconds = 30
- account_a = MailboxAccount(
- email="a@demo.example.test",
- account_id="jwt-a",
- extra={"email_domain": "demo.example.test"},
- )
- account_b = MailboxAccount(
- email="b@demo.example.test",
- account_id="jwt-b",
- extra={"email_domain": "demo.example.test"},
- )
- loop._on_cfmail_wait_progress(
- account_a,
- {"message_scan_count": 0, "elapsed_seconds": 35},
- )
- loop._on_cfmail_wait_progress(
- account_b,
- {"message_scan_count": 0, "elapsed_seconds": 36},
- )
- stoploss = loop.snapshot()["cfmail_wait_otp_stoploss"]
- assert stoploss["in_cooldown"] is False
- def test_registration_loop_does_not_activate_wait_otp_stoploss_for_non_mailbox_timeout() -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- loop._cfmail_wait_otp_window = 2
- loop._cfmail_wait_otp_threshold = 2
- loop._cfmail_wait_otp_cooldown_seconds = 60
- result = {
- "success": False,
- "stage": "wait_otp",
- "error_message": "otp retrieval failed",
- "metadata": {
- "mail_provider": "cfmail",
- "email_domain": "demo.example.test",
- "otp_wait_failure_reason": "mailbox_timeout_no_match",
- "otp_mailbox_message_scan_count": 1,
- },
- }
- loop._update_cfmail_wait_otp_stoploss(result)
- loop._update_cfmail_wait_otp_stoploss(result)
- stoploss = loop.snapshot()["cfmail_wait_otp_stoploss"]
- assert stoploss["in_cooldown"] is False
- def test_registration_loop_rotates_cfmail_domain_when_wait_otp_stoploss_active(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- loop._providers = ["cfmail"]
- monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "demo.example.test")
- loop._cfmail_tracker = DomainHealthTracker(window_size=2, blacklist_threshold=2, rotation_cooldown_seconds=1)
- loop._cfmail_wait_otp_cooldown_seconds = 60
- loop._cfmail_wait_otp_state.update(
- {
- "active_domain": "demo.example.test",
- "in_cooldown": True,
- "cooldown_until": 9999999999.0,
- "last_triggered_at": "2026-03-29T10:00:00",
- "last_rotation_attempted_at": "",
- }
- )
- class FakeProvisioner:
- def rotate_active_domain(self): # type: ignore[no-untyped-def]
- return ProvisionResult(
- success=True,
- step="rotate",
- old_domain="demo.example.test",
- new_domain="auto-new.example.test",
- )
- reload_calls: list[bool] = []
- class FakeManager:
- def reload_if_needed(self, force=False): # type: ignore[no-untyped-def]
- reload_calls.append(bool(force))
- return True
- loop._cfmail_manager = FakeManager()
- loop._cfmail_provisioner = FakeProvisioner()
- assert loop._wait_if_cfmail_wait_otp_stopped(thread_id=1, provider="cfmail") is True
- snapshot = loop.snapshot()
- assert snapshot["cfmail_rotation"]["last_new_domain"] == "auto-new.example.test"
- assert snapshot["cfmail_wait_otp_stoploss"]["in_cooldown"] is False
- assert snapshot["cfmail_wait_otp_stoploss"]["active_domain"] == "auto-new.example.test"
- assert reload_calls == [True]
- def test_registration_loop_rotates_cfmail_domain_when_add_phone_stoploss_active(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- loop._providers = ["cfmail"]
- monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "demo.example.test")
- loop._cfmail_tracker = DomainHealthTracker(window_size=2, blacklist_threshold=2, rotation_cooldown_seconds=1)
- loop._cfmail_add_phone_cooldown_seconds = 60
- loop._cfmail_add_phone_state.update(
- {
- "active_domain": "demo.example.test",
- "in_cooldown": True,
- "cooldown_until": 9999999999.0,
- "last_triggered_at": "2026-03-29T10:00:00",
- "last_rotation_attempted_at": "",
- }
- )
- class FakeProvisioner:
- def rotate_active_domain(self): # type: ignore[no-untyped-def]
- return ProvisionResult(
- success=True,
- step="rotate",
- old_domain="demo.example.test",
- new_domain="auto-new.example.test",
- )
- loop._cfmail_provisioner = FakeProvisioner()
- assert loop._wait_if_cfmail_add_phone_stopped(thread_id=1, provider="cfmail") is True
- snapshot = loop.snapshot()
- assert snapshot["cfmail_rotation"]["last_new_domain"] == "auto-new.example.test"
- assert snapshot["cfmail_add_phone_stoploss"]["in_cooldown"] is False
- assert snapshot["cfmail_add_phone_stoploss"]["active_domain"] == "auto-new.example.test"
- def test_registration_loop_ignores_stale_wait_otp_result_from_old_domain(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "auto-new.example.test")
- result = {
- "success": False,
- "stage": "wait_otp",
- "error_message": "otp retrieval failed",
- "metadata": {
- "mail_provider": "cfmail",
- "email_domain": "auto-old.example.test",
- "otp_wait_failure_reason": "mailbox_timeout_no_message",
- "otp_mailbox_message_scan_count": 0,
- },
- }
- loop._update_cfmail_wait_otp_stoploss(result)
- stoploss = loop.snapshot()["cfmail_wait_otp_stoploss"]
- assert stoploss["active_domain"] == ""
- assert stoploss["in_cooldown"] is False
- def test_registration_loop_does_not_abort_old_domain_wait_after_rotation(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "auto-new.example.test")
- loop._cfmail_wait_otp_state.update(
- {
- "active_domain": "auto-old.example.test",
- "in_cooldown": True,
- "cooldown_until": 9999999999.0,
- "last_triggered_at": "2026-03-29T10:00:00",
- }
- )
- account = MailboxAccount(
- email="a@auto-old.example.test",
- account_id="jwt-old",
- extra={
- "email_domain": "auto-old.example.test",
- "otp_wait_started_at": 1743242390.0,
- },
- )
- assert loop._should_abort_cfmail_wait(account) is False
- def test_registration_loop_aborts_wait_for_domain_in_wait_otp_cooldown(monkeypatch) -> None:
- settings = _base_settings(register_mail_provider="cfmail")
- loop = RegistrationLoop(settings)
- monkeypatch.setattr(loop, "_current_cfmail_active_domain", lambda: "demo.example.test")
- loop._cfmail_wait_otp_state.update(
- {
- "active_domain": "demo.example.test",
- "in_cooldown": True,
- "cooldown_until": 9999999999.0,
- }
- )
- account = MailboxAccount(
- email="a@demo.example.test",
- account_id="jwt-demo",
- extra={
- "email_domain": "demo.example.test",
- "otp_wait_started_at": 1743242410.0,
- },
- )
- assert loop._should_abort_cfmail_wait(account) is True
- def test_registration_burst_scheduler_runs_batch_with_batch_config(monkeypatch, tmp_path: Path) -> None:
- settings = _base_settings(
- register_log_file="",
- runtime_state_file=tmp_path / "runtime_state.json",
- register_batch_threads=1,
- register_batch_target_count=20,
- register_batch_interval_seconds=60,
- )
- observed: dict[str, int] = {}
- class FakeLoop:
- def __init__(self, batch_settings): # type: ignore[no-untyped-def]
- observed["threads"] = batch_settings.register_threads
- observed["target"] = batch_settings.register_target_count
- self._threads = []
- def start(self) -> None:
- return None
- def stop(self) -> None:
- return None
- def snapshot(self) -> dict[str, object]:
- return {
- "name": "register",
- "status": "stopped",
- "threads_alive": 0,
- "threads_total": 1,
- "total_attempts": 22,
- "total_success": 20,
- "total_failure": 2,
- "success_rate": 90.9,
- "target_count": 20,
- "target_reached": True,
- "last_error": "",
- "proxy": None,
- "proxy_pool_enabled": False,
- "mail_provider": "mailtm",
- "interval_seconds": 5,
- "run_count": 22,
- "success_count": 20,
- "failure_count": 2,
- "is_running": False,
- "last_started_at": None,
- "last_finished_at": None,
- "last_duration_seconds": None,
- "next_run_at": None,
- "failure_by_stage": {"add_phone_gate": 2},
- "failure_signals": {"add_phone_gate": 2},
- "recent_failure_hotspots": [{"key": "add_phone_gate", "stage": "add_phone_gate", "count": 2}],
- "recent_attempts": [
- {
- "timestamp": "2026-03-27T10:00:00",
- "success": False,
- "stage": "add_phone_gate",
- "signal": "add_phone_gate",
- "error_message": "phone gate",
- "email_domain": "demo.example.test",
- "post_create_gate": "add_phone",
- "create_account_error_code": "",
- "proxy_key": "",
- "email": "",
- }
- ],
- "active_domain_recent_attempts": [],
- "active_domain_failure_by_stage": {"add_phone_gate": 2},
- "active_domain_failure_signals": {"add_phone_gate": 2},
- "active_domain_recent_failure_hotspots": [{"key": "add_phone_gate", "stage": "add_phone_gate", "count": 2}],
- "cfmail_rotation": None,
- "cfmail_add_phone_stoploss": {"active_domain": "demo.example.test"}
- }
- monkeypatch.setattr("main.RegistrationLoop", FakeLoop)
- scheduler = RegistrationBurstScheduler(settings)
- original_absorb = scheduler._absorb_batch_snapshot
- def absorb_and_stop(snapshot: dict[str, object], *, duration_seconds: float) -> None:
- original_absorb(snapshot, duration_seconds=duration_seconds)
- scheduler.stop()
- monkeypatch.setattr(scheduler, "_absorb_batch_snapshot", absorb_and_stop)
- scheduler.run()
- payload = json.loads((tmp_path / "runtime_state.json").read_text(encoding="utf-8"))
- register_snapshot = payload["register_snapshot"]
- assert observed == {"threads": 1, "target": 20}
- assert register_snapshot["status"] == "stopped"
- assert register_snapshot["scheduler_mode"] == "burst"
- assert register_snapshot["run_count"] == 1
- assert register_snapshot["total_success"] == 20
- assert register_snapshot["batch_target_count"] == 20
|