"""端到端编排:注册 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, 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, init_db, 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 长链 — 走 payurl.ark2.cn(与 Chrome 扩展一致) 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}"), ) sub_ctx.email = email sub_ctx.password = password 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 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 add_event(email, "paypal", "ok") upsert_account(email, password, fields={"final_status": "paid"}) # 3) 重新拉 session,看 planType full_ctx.check_stop() full_ctx.set_stage(f"账号 {idx}/{total}:付款完成,重新拉 session") new_session = _refresh_session_from_page(page, access_token, log=sub_log) 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) full_ctx.set_stage(f"账号 {idx}/{total}:完成 file={cpa_result.get('fileName')}") return record 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 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} 个账号") if not cfg.cpa_url or not cfg.cpa_management_key: full_ctx.log("[full] 警告:未配置 CPA 地址/密钥,仍会注册并付款,但跳过 CPA 上传") from playwright.sync_api import sync_playwright total = max(1, int(cfg.account_count)) for idx in range(1, total + 1): full_ctx.check_stop() full_ctx.set_stage(f"启动账号 {idx}/{total} 的浏览器") with sync_playwright() as p: browser = p.chromium.launch( headless=cfg.headless, args=["--disable-blink-features=AutomationControlled"], ) ctx_browser = browser.new_context( locale="en-US", timezone_id="America/New_York", viewport={"width": 1280, "height": 900}, ) 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, idx, total) 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] 账号 {idx}/{total} 异常: {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 full_ctx.on_account_finished: try: full_ctx.on_account_finished(record) except Exception: pass full_ctx.set_stage("全部账号已处理完毕") full_ctx.state = "done" return full_ctx