"""端到端编排:注册 ChatGPT → 生成 Plus 长链 → PayPal 付款 → 校验 Plus → 上传 CPA。""" from __future__ import annotations import json import time import traceback from dataclasses import dataclass, field from typing import Callable, Optional from automation import ( RunContext, StripeNonFreeDetected, generate_long_link_local, generate_long_link_payurl, run_paypal_flow, _dump_page, ) from chatgpt_signup import fetch_current_session, signup_chatgpt from config import AppConfig from cpa_uploader import ( get_session_plan_type, is_plus_session, upload_session_to_cpa, ) from mail_provider import build_a4sky_email # noqa: F401 (re-exported for tests) from storage import add_event, get_account, init_db, set_account_trial_eligibility, upsert_account @dataclass class FullRunContext: cfg: AppConfig log: Callable[[str], None] = print on_stage: Optional[Callable[[str], None]] = None on_account_finished: Optional[Callable[[dict], None]] = None state: str = "running" stage: str = "" accounts: list[dict] = field(default_factory=list) run_id: str = field(default_factory=lambda: time.strftime("%Y%m%d-%H%M%S")) def set_stage(self, name: str): self.stage = name self.log(f"[stage:full] {name}") if self.on_stage: try: self.on_stage(name) except Exception: pass def check_stop(self): if self.state == "stopped": raise RuntimeError("STOPPED_BY_USER") def _refresh_session(page, log: Callable[[str], None]) -> dict: """支付完成后重新拉一次 /api/auth/session 看 planType。""" return fetch_current_session(page, log) def _run_one_account(full_ctx: FullRunContext, page, idx: int, total: int) -> dict: cfg = full_ctx.cfg full_ctx.check_stop() full_ctx.set_stage(f"账号 {idx}/{total}:开始注册") sub_log = lambda msg: full_ctx.log(f"[acc{idx}] {msg}") signup_result = signup_chatgpt( page, helper_url=cfg.mail_helper_url, mail_domain=cfg.mail_domain, mail_poll_interval_sec=cfg.mail_poll_interval_sec, mail_poll_max_attempts=cfg.mail_poll_max_attempts, log=sub_log, on_stage=lambda name: full_ctx.set_stage(f"账号 {idx}/{total}:注册-{name}"), ) email = signup_result["email"] password = signup_result["password"] session = signup_result["session"] access_token = session.get("accessToken") or "" if not access_token: upsert_account(email, password, fields={ "final_status": "failed", "last_error": "注册成功但未拿到 accessToken", "initial_session": session, }) add_event(email, "register", "error", "未拿到 accessToken") raise RuntimeError("注册成功但未拿到 accessToken") initial_plan = get_session_plan_type(session) upsert_account(email, password, fields={ "final_status": "registered", "plan_type": initial_plan, "initial_session": session, }) add_event(email, "register", "ok", f"plan={initial_plan}") record = { "email": email, "password": password, "stage": "registered", "planType": initial_plan, "cpa": None, "error": None, } full_ctx.set_stage(f"账号 {idx}/{total}:注册成功 email={email}") # 1) 生成 Plus 长链 full_ctx.check_stop() full_ctx.set_stage(f"账号 {idx}/{total}:生成 Plus 长链") sub_ctx = RunContext( token=access_token, plan="plus", country="US", currency="USD", use_promo=cfg.use_promo, headless=cfg.headless, log=sub_log, on_stage=lambda name: full_ctx.set_stage(f"账号 {idx}/{total}:长链-{name}"), phone_e164=cfg.phone_e164, sms_api_url=cfg.sms_api_url, paypal_proxy=cfg.effective_paypal_proxy, long_link_mode=cfg.long_link_mode, long_link_proxy=cfg.long_link_proxy, ) sub_ctx.email = email sub_ctx.password = password sub_ctx._stop_hook = full_ctx.check_stop sub_ctx._on_trial_eligibility_detected = lambda eligible: set_account_trial_eligibility(email, password, eligible) # type: ignore if cfg.long_link_mode == "local": long_link = generate_long_link_local(sub_ctx, proxy=cfg.long_link_proxy) else: long_link = generate_long_link_payurl(sub_ctx) sub_ctx.long_link = long_link record["longLink"] = long_link upsert_account(email, password, fields={"long_link": long_link}) add_event(email, "long_link", "ok", long_link[:200]) # 2) PayPal 付款 — 复用同一个浏览器 page,避免嵌套 sync_playwright full_ctx.check_stop() full_ctx.set_stage(f"账号 {idx}/{total}:进入 PayPal 付款流") sub_ctx._reuse_account_for_paypal = True # type: ignore try: run_paypal_flow(sub_ctx, page=page) except StripeNonFreeDetected as exc: full_ctx.log(f"[acc{idx}] 非免费金额,跳过自动付款,标记为 trial 待手动付款: {exc}") upsert_account(email, password, fields={ "final_status": "trial", "last_error": f"非免费金额需手动付款: {exc}", "initial_session": session, }) add_event(email, "paypal", "warn", f"非免费金额,跳过自动付款: {exc}") record["stage"] = "trial" record["error"] = None return record except Exception as exc: upsert_account(email, password, fields={ "final_status": "failed", "last_error": f"PayPal 流异常: {exc!r}", }) add_event(email, "paypal", "error", repr(exc)) raise is_free_trial = getattr(sub_ctx, "_is_free_trial", False) if is_free_trial: add_event(email, "paypal", "ok", "免费试用付款完成") upsert_account(email, password, fields={"final_status": "paid"}) else: add_event(email, "paypal", "ok") upsert_account(email, password, fields={"final_status": "paid"}) # 3) 重新拉 session,看 planType(带轮询:付款后 plus 状态可能延迟到位) full_ctx.check_stop() full_ctx.set_stage(f"账号 {idx}/{total}:付款完成,等待 plus 状态到位(最长 3 分钟)") # PayPal 阶段如果切了代理 context,结束后会把切回的原 page 写到 sub_ctx._post_paypal_page session_page = getattr(sub_ctx, "_post_paypal_page", None) or page new_session = _wait_for_plus_session( session_page, access_token, log=sub_log, max_wait_sec=180, interval_sec=6, ) new_plan = get_session_plan_type(new_session) record["planType"] = new_plan record["sessionRefreshed"] = True upsert_account(email, password, fields={ "plan_type": new_plan, "plus_session": new_session, }) if not is_plus_session(new_session): record["stage"] = "plus_check_failed" record["error"] = f"planType 不是 plus,实际为 {record['planType']!r}" full_ctx.log(f"[acc{idx}] 失败:{record['error']}") upsert_account(email, password, fields={ "final_status": "plus_check_failed", "last_error": record["error"], }) add_event(email, "plus_check", "error", record["error"]) return record full_ctx.set_stage(f"账号 {idx}/{total}:Plus 校验通过") upsert_account(email, password, fields={"final_status": "plus"}) add_event(email, "plus_check", "ok", f"plan={new_plan}") # 4) 上传 CPA full_ctx.check_stop() full_ctx.set_stage(f"账号 {idx}/{total}:上传 CPA") if not (cfg.cpa_url and cfg.cpa_management_key): full_ctx.log(f"[acc{idx}] 跳过 CPA 上传(未配置 cpa_url / cpa_management_key)") record["stage"] = "cpa_skipped" upsert_account(email, password, fields={"final_status": "cpa_skipped"}) add_event(email, "cpa", "warn", "未配置 CPA") return record try: cpa_result = upload_session_to_cpa( new_session, cpa_url=cfg.cpa_url, management_key=cfg.cpa_management_key, email_hint=email, log=sub_log, ) except Exception as exc: upsert_account(email, password, fields={ "final_status": "cpa_failed", "last_error": f"CPA 上传异常: {exc!r}", }) add_event(email, "cpa", "error", repr(exc)) raise record["cpa"] = cpa_result record["stage"] = "cpa_uploaded" upsert_account(email, password, fields={ "final_status": "cpa_uploaded", "cpa_file_name": cpa_result.get("fileName"), "cpa_uploaded_at": int(time.time() * 1000), }) add_event(email, "cpa", "ok", cpa_result.get("fileName"), payload=cpa_result) if is_free_trial: full_ctx.log(f"[acc{idx}] 免费试用账号,CPA 已上传,标记为 trial(后续到期可重新付款)") full_ctx.set_stage(f"账号 {idx}/{total}:完成 file={cpa_result.get('fileName')}") return record def _wait_for_plus_session( page, access_token: str, *, log: Callable[[str], None], max_wait_sec: int = 180, interval_sec: int = 6, ) -> dict: """付款完成后轮询 /api/auth/session,直到 planType=='plus' 或超时。 返回最终拿到的 session(若一直没 plus,也返回最后一次的 session 给上层判定)。 """ log(f"[session] 开始轮询等待 planType=plus,最长 {max_wait_sec}s,间隔 {interval_sec}s") deadline = time.time() + max_wait_sec last_session: dict = {} last_plan = "" attempt = 0 while time.time() < deadline: attempt += 1 try: sess = _refresh_session_from_page(page, access_token, log=log) except Exception as exc: log(f"[session] 第 {attempt} 次拉取异常: {exc!r},{interval_sec}s 后重试") time.sleep(interval_sec) continue last_session = sess or {} plan = (last_session.get("account") or {}).get("planType") or "" if plan != last_plan: log(f"[session] 第 {attempt} 次:planType={plan!r}") last_plan = plan if plan.lower() == "plus": log(f"[session] planType=plus 已到位,用时 ~{attempt * interval_sec}s") return last_session remaining = max(0, int(deadline - time.time())) log(f"[session] planType 仍为 {plan!r}(非 plus),{interval_sec}s 后重试,剩余 {remaining}s") time.sleep(interval_sec) log(f"[session] 等待 plus 超时,最后 planType={last_plan!r}") return last_session def _refresh_session_from_page(page, access_token: str, *, log: Callable[[str], None]) -> dict: """付款完成后用同一个浏览器 page 拉 session,避免嵌套 sync_playwright。""" log("[session] 浏览器内拉 /api/auth/session ...") try: page.goto("https://chatgpt.com/", wait_until="domcontentloaded", timeout=45000) page.wait_for_timeout(2000) except Exception as exc: log(f"[session] 跳回 chatgpt.com 异常: {exc!r}") deadline = time.time() + 60 last = "" while time.time() < deadline: try: data = page.evaluate( """async (token) => { try { const r = await fetch('/api/auth/session', { credentials: 'include', headers: { 'Authorization': 'Bearer ' + token, 'Accept': 'application/json' }, }); const t = await r.text(); try { return { ok: true, data: JSON.parse(t), status: r.status }; } catch (_) { return { ok: false, raw: t, status: r.status }; } } catch (e) { return { ok: false, error: String(e) }; } }""", access_token, ) if isinstance(data, dict) and data.get("ok") and isinstance(data.get("data"), dict): sess = data["data"] if sess.get("accessToken"): plan = (sess.get("account") or {}).get("planType") log(f"[session] 拉到 session planType={plan} status={data.get('status')}") return sess preview = json.dumps(data, ensure_ascii=False)[:200] if isinstance(data, dict) else str(data)[:200] if preview != last: log(f"[session] 暂无可用 session,预览={preview}") last = preview except Exception as exc: log(f"[session] page.evaluate 异常: {exc!r}") time.sleep(2) raise TimeoutError("拉取 session 超时(60s 内未取到 accessToken)") def _open_and_fetch_session_with_token(access_token: str, *, log: Callable[[str], None]) -> dict: """[已废弃] 旧实现会嵌套 sync_playwright 导致流程静默退出。 保留空壳避免外部 import 报错;新流程请用 _refresh_session_from_page。 """ raise RuntimeError("_open_and_fetch_session_with_token 已废弃,请使用 _refresh_session_from_page(page, access_token)") def _build_proxy_cfg(cfg: AppConfig, log: Callable[[str], None]): proxy_str = cfg.effective_global_proxy if not proxy_str: log("[full] 全局代理未配置,ChatGPT 注册 / 长链 直连") return None from automation import _parse_proxy_url proxy_cfg = _parse_proxy_url(proxy_str) if not proxy_cfg: log(f"[full] 警告:proxy_url={proxy_str!r} 解析失败,将直连") return None masked = dict(proxy_cfg) if masked.get("password"): masked["password"] = "***" log(f"[full] 全局代理(ChatGPT/长链/默认): {masked}") return proxy_cfg def _open_browser(cfg: AppConfig, proxy_cfg, log: Callable[[str], None]): """打开 chromium browser + context + page,返回 (p_ctx_mgr, browser, context, page)。 调用方负责 close browser 和 退出 sync_playwright 上下文。""" from patchright.sync_api import sync_playwright from geo_fingerprint import detect_openai_geo_fingerprint geo = detect_openai_geo_fingerprint(cfg.effective_global_proxy, log=log) p_ctx = sync_playwright().__enter__() browser = p_ctx.chromium.launch( headless=cfg.headless, args=["--disable-blink-features=AutomationControlled"], ) ctx_browser = browser.new_context( locale=geo.locale, timezone_id=geo.timezone_id, viewport={"width": 1280, "height": 900}, proxy=proxy_cfg, ) page = ctx_browser.new_page() page.on("console", lambda m: log(f"[browser-console:{m.type}] {m.text[:300]}")) page.on("pageerror", lambda e: log(f"[browser-pageerror] {e}")) return p_ctx, browser, ctx_browser, page def _close_browser(p_ctx, browser, log: Callable[[str], None]): try: browser.close() except Exception as exc: log(f"[full] browser.close 异常: {exc!r}") try: from patchright.sync_api import sync_playwright sync_playwright().__exit__(None, None, None) except Exception: pass # 直接调 __exit__ 在新实例上不对,只能依赖 GC try: p_ctx.__exit__(None, None, None) except Exception: pass def run_pay_only( cfg: AppConfig, *, session: dict, log: Callable[[str], None] = print, on_stage: Optional[Callable[[str], None]] = None, stop_check: Optional[Callable[[], None]] = None, ) -> dict: """传入已有 ChatGPT session JSON,直接走付款 → Plus 校验 → CPA 上传。 返回 record(同 _run_one_account 的格式)。 """ full_ctx = FullRunContext(cfg=cfg, log=log, on_stage=on_stage) if stop_check: # 把外部 stop 钩子接入 check_stop orig_check = full_ctx.check_stop def _check_combined(): stop_check() orig_check() full_ctx.check_stop = _check_combined # type: ignore if not isinstance(session, dict) or not session.get("accessToken"): raise RuntimeError("pay_only 需要 session JSON 且包含 accessToken") access_token = session["accessToken"] email = ( ((session.get("user") or {}).get("email")) or session.get("email") or "" ) password = "" # pay_only 不知道原密码,PayPal 注册新邮箱不需要 account_key = email or f"unknown-{int(time.time())}" full_ctx.set_stage(f"pay_only:开始(email={email or '(unknown)'})") upsert_account(account_key, password, fields={ "final_status": "registered", "plan_type": get_session_plan_type(session), "initial_session": session, }) add_event(email or "", "pay_only", "info", "begin") existing_account = get_account(email) if email else None proxy_cfg = _build_proxy_cfg(cfg, full_ctx.log) p_ctx, browser, ctx_browser, page = _open_browser(cfg, proxy_cfg, full_ctx.log) try: # 用 sub_ctx 走付款 full_ctx.set_stage("pay_only:生成 Plus 长链") sub_ctx = RunContext( token=access_token, plan="plus", country="US", currency="USD", use_promo=cfg.use_promo, headless=cfg.headless, log=lambda m: full_ctx.log(f"[pay] {m}"), on_stage=lambda name: full_ctx.set_stage(f"pay_only:长链-{name}"), phone_e164=cfg.phone_e164, sms_api_url=cfg.sms_api_url, paypal_proxy=cfg.effective_paypal_proxy, long_link_mode=cfg.long_link_mode, long_link_proxy=cfg.long_link_proxy, ) sub_ctx.email = email sub_ctx.password = password sub_ctx._stop_hook = full_ctx.check_stop # type: ignore sub_ctx._on_trial_eligibility_detected = lambda eligible: set_account_trial_eligibility(account_key, password, eligible) # type: ignore if existing_account and existing_account.get("trial_eligible"): sub_ctx._trial_eligible = True if cfg.long_link_mode == "local": long_link = generate_long_link_local(sub_ctx, proxy=cfg.long_link_proxy) else: long_link = generate_long_link_payurl(sub_ctx) sub_ctx.long_link = long_link full_ctx.set_stage("pay_only:进入 PayPal") sub_ctx._reuse_account_for_paypal = True # type: ignore run_paypal_flow(sub_ctx, page=page) full_ctx.set_stage("pay_only:付款完成,等 plus 状态") session_page = getattr(sub_ctx, "_post_paypal_page", None) or page new_session = _wait_for_plus_session( session_page, access_token, log=full_ctx.log, max_wait_sec=180, interval_sec=6, ) plan = get_session_plan_type(new_session) upsert_account(email or "", password, fields={ "plan_type": plan, "plus_session": new_session, }) if not is_plus_session(new_session): upsert_account(email or "", password, fields={ "final_status": "plus_check_failed", "last_error": f"planType={plan!r}", }) return {"stage": "plus_check_failed", "planType": plan, "error": f"planType={plan!r}"} upsert_account(email or "", password, fields={"final_status": "plus", "last_error": ""}) full_ctx.set_stage("pay_only:Plus 校验通过") if not (cfg.cpa_url and cfg.cpa_management_key): upsert_account(email or "", password, fields={"final_status": "cpa_skipped"}) return {"stage": "cpa_skipped", "planType": plan} cpa_result = upload_session_to_cpa( new_session, cpa_url=cfg.cpa_url, management_key=cfg.cpa_management_key, email_hint=email, log=full_ctx.log, ) upsert_account(email or "", password, fields={ "final_status": "cpa_uploaded", "cpa_file_name": cpa_result.get("fileName"), "cpa_uploaded_at": int(time.time() * 1000), "last_error": "", }) add_event(email or "", "cpa", "ok", cpa_result.get("fileName"), payload=cpa_result) return { "stage": "cpa_uploaded", "email": email, "planType": plan, "cpa": cpa_result, } finally: try: browser.close() except Exception: pass def run_full(cfg: AppConfig, *, log: Callable[[str], None] = print, on_stage: Optional[Callable[[str], None]] = None) -> FullRunContext: full_ctx = FullRunContext(cfg=cfg, log=log, on_stage=on_stage) full_ctx.set_stage(f"开始全自动流程 目标 {cfg.account_count} 个 CPA 完成") if not cfg.cpa_url or not cfg.cpa_management_key: full_ctx.log("[full] 警告:未配置 CPA 地址/密钥,仍会注册并付款,但跳过 CPA 上传") from patchright.sync_api import sync_playwright from geo_fingerprint import detect_geo_fingerprint proxy_str = cfg.effective_global_proxy proxy_cfg = None if proxy_str: from automation import _parse_proxy_url proxy_cfg = _parse_proxy_url(proxy_str) if proxy_cfg: masked = dict(proxy_cfg) if masked.get("password"): masked["password"] = "***" full_ctx.log(f"[full] 全局代理(ChatGPT/长链/默认): {masked}") else: full_ctx.log(f"[full] 警告:proxy_url={proxy_str!r} 解析失败,将直连") else: full_ctx.log("[full] 全局代理未配置,ChatGPT 注册 / 长链 直连") from geo_fingerprint import detect_openai_geo_fingerprint geo = detect_openai_geo_fingerprint(proxy_str, log=full_ctx.log) paypal_only_str = cfg.paypal_only_proxy.strip() if cfg.paypal_only_proxy else "" if paypal_only_str: full_ctx.log(f"[full] PayPal 独立代理已配置(PayPal 阶段会切到该代理)") target = max(1, int(cfg.account_count)) success_count = 0 attempt_idx = 0 while success_count < target: attempt_idx += 1 full_ctx.check_stop() full_ctx.set_stage(f"启动账号 第{attempt_idx}次尝试(已成功 {success_count}/{target})的浏览器") with sync_playwright() as p: browser = p.chromium.launch( headless=cfg.headless, args=["--disable-blink-features=AutomationControlled"], ) ctx_browser = browser.new_context( locale=geo.locale, timezone_id=geo.timezone_id, viewport={"width": 1280, "height": 900}, proxy=proxy_cfg, ) page = ctx_browser.new_page() page.on("console", lambda m: full_ctx.log(f"[browser-console:{m.type}] {m.text[:300]}")) page.on("pageerror", lambda e: full_ctx.log(f"[browser-pageerror] {e}")) try: record = _run_one_account(full_ctx, page, attempt_idx, target) except Exception as exc: if str(exc) == "STOPPED_BY_USER": full_ctx.state = "stopped" full_ctx.set_stage("用户停止") full_ctx.accounts.append({"stage": "stopped", "error": "user stopped"}) return full_ctx full_ctx.log(f"[full] 第{attempt_idx}次尝试异常: {exc!r}") full_ctx.log(traceback.format_exc()) record = {"stage": "error", "error": repr(exc)} finally: try: browser.close() except Exception: pass full_ctx.accounts.append(record) if record.get("stage") == "cpa_uploaded": success_count += 1 full_ctx.log(f"[full] CPA 上传成功 {success_count}/{target}") if full_ctx.on_account_finished: try: full_ctx.on_account_finished(record) except Exception: pass full_ctx.set_stage(f"已完成 {success_count}/{target} 个 CPA 上传") full_ctx.state = "done" return full_ctx