recheck.py 8.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234
  1. """对历史失败账号做补救:用 access_token 调 backend-api 检查 plan,若已 plus 则上传 CPA。
  2. 注意:/api/auth/session 走 NextAuth 只认 cookie,不能用 Bearer 直连。
  3. 真正的鉴权接口是 /backend-api/me(也叫 accounts/check),它接受 Authorization: Bearer。
  4. """
  5. from __future__ import annotations
  6. import copy
  7. import json
  8. import time
  9. import urllib.error
  10. import urllib.request
  11. from typing import Callable
  12. from cpa_uploader import (
  13. get_session_plan_type,
  14. is_plus_session,
  15. upload_session_to_cpa,
  16. )
  17. from storage import add_event, get_account, upsert_account
  18. ME_URL = "https://chatgpt.com/backend-api/me"
  19. ACCOUNTS_CHECK_URL = "https://chatgpt.com/backend-api/accounts/check/v4-2023-04-27"
  20. def _http_get_json(url: str, *, headers: dict, timeout: int = 20, log: Callable[[str], None] = print) -> tuple[int, str, dict]:
  21. started = time.time()
  22. text = ""
  23. status = 0
  24. try:
  25. from curl_cffi import requests as curl_requests # type: ignore
  26. r = curl_requests.get(url, headers=headers, impersonate="chrome136", timeout=timeout)
  27. text = r.text
  28. status = r.status_code
  29. except Exception as exc:
  30. log(f"[recheck] curl_cffi 不可用: {exc!r},回退 urllib")
  31. req = urllib.request.Request(url, method="GET")
  32. for k, v in headers.items():
  33. req.add_header(k, v)
  34. try:
  35. with urllib.request.urlopen(req, timeout=timeout) as resp:
  36. text = resp.read().decode("utf-8", errors="replace")
  37. status = resp.status
  38. except urllib.error.HTTPError as exc2:
  39. text = exc2.read().decode("utf-8", errors="replace") if hasattr(exc2, "read") else ""
  40. status = exc2.code
  41. elapsed_ms = int((time.time() - started) * 1000)
  42. log(f"[recheck] GET {url} HTTP {status} 耗时 {elapsed_ms}ms 长度={len(text)}")
  43. parsed: dict = {}
  44. try:
  45. parsed = json.loads(text or "{}")
  46. except Exception:
  47. parsed = {}
  48. return status, text, parsed
  49. def _extract_plan_from_me(payload: dict) -> tuple[str, str, str]:
  50. """从 /backend-api/me 或 /accounts/check 响应里抽 (plan_type, account_id, email)。"""
  51. if not isinstance(payload, dict):
  52. return "", "", ""
  53. # /me 里直接给 email 和 chat_id;plan 通常在 accounts.<id>.account.plan_type
  54. email = ""
  55. if isinstance(payload.get("email"), str):
  56. email = payload["email"]
  57. accounts = payload.get("accounts")
  58. if isinstance(accounts, dict):
  59. # 优先找 plan_type='plus'
  60. best_id = ""
  61. best_plan = ""
  62. for acc_id, item in accounts.items():
  63. if not isinstance(item, dict):
  64. continue
  65. account_obj = item.get("account") if isinstance(item.get("account"), dict) else item
  66. plan = (account_obj or {}).get("plan_type") or (account_obj or {}).get("planType") or ""
  67. plan = str(plan or "").strip()
  68. if not plan:
  69. continue
  70. if plan.lower() == "plus" and not best_plan:
  71. best_plan = plan
  72. best_id = acc_id
  73. break
  74. if not best_plan:
  75. best_plan = plan
  76. best_id = acc_id
  77. if best_plan:
  78. return best_plan, best_id, email
  79. # 兜底:顶层 plan_type
  80. plan = str(payload.get("plan_type") or payload.get("planType") or "").strip()
  81. return plan, "", email
  82. def _merge_plus_into_session(base_session: dict, plan: str, account_id: str, email: str) -> dict:
  83. """把 backend-api 拉到的 plan/account_id 合并进 DB 里的 session,构造给 CPA 用的 plus session。"""
  84. sess = copy.deepcopy(base_session) if isinstance(base_session, dict) else {}
  85. if "account" not in sess or not isinstance(sess.get("account"), dict):
  86. sess["account"] = {}
  87. if plan:
  88. sess["account"]["planType"] = plan
  89. sess["planType"] = plan
  90. if account_id and not sess["account"].get("id"):
  91. sess["account"]["id"] = account_id
  92. if email:
  93. if "user" not in sess or not isinstance(sess.get("user"), dict):
  94. sess["user"] = {}
  95. if not sess["user"].get("email"):
  96. sess["user"]["email"] = email
  97. return sess
  98. def recheck_account(
  99. email: str,
  100. *,
  101. cpa_url: str,
  102. cpa_management_key: str,
  103. log: Callable[[str], None] = print,
  104. ) -> dict:
  105. """对单个账号做补救。返回 { ok, planType, action, cpa, error }。"""
  106. acc = get_account(email)
  107. if not acc:
  108. return {"ok": False, "error": "账号不存在"}
  109. plus_sess = acc.get("plus_session") or {}
  110. init_sess = acc.get("initial_session") or {}
  111. access_token = (plus_sess.get("accessToken") if isinstance(plus_sess, dict) else None) \
  112. or (init_sess.get("accessToken") if isinstance(init_sess, dict) else None)
  113. if not access_token:
  114. return {"ok": False, "error": "数据库里找不到 accessToken"}
  115. log(f"[recheck] {email} 开始补救(访问 backend-api/me)")
  116. add_event(email, "recheck", "info", "begin")
  117. headers = {
  118. "Authorization": f"Bearer {access_token}",
  119. "Accept": "application/json",
  120. "Origin": "https://chatgpt.com",
  121. "Referer": "https://chatgpt.com/",
  122. }
  123. plan = ""
  124. account_id = ""
  125. me_email = ""
  126. last_payload: dict = {}
  127. # 先试 /backend-api/me(轻量),失败回退 /accounts/check
  128. for url in (ME_URL, ACCOUNTS_CHECK_URL):
  129. try:
  130. status, text, payload = _http_get_json(url, headers=headers, log=log)
  131. except Exception as exc:
  132. log(f"[recheck] {url} 异常: {exc!r}")
  133. continue
  134. if status >= 400:
  135. log(f"[recheck] {url} 返回 {status},预览={text[:200]}")
  136. continue
  137. last_payload = payload
  138. plan, account_id, me_email = _extract_plan_from_me(payload)
  139. log(f"[recheck] {url} 提取 plan={plan!r} account_id={account_id!r} email={me_email!r}")
  140. if plan:
  141. break
  142. if not plan:
  143. msg = "重拉 session 失败:所有 backend-api 接口都未返回 plan_type"
  144. log(f"[recheck] {email} {msg}")
  145. upsert_account(email, acc.get("password") or "", fields={"last_error": msg})
  146. add_event(email, "recheck", "error", msg)
  147. return {"ok": False, "error": msg}
  148. # 用 base session(plus_session 优先,否则 init_session)合并新 plan
  149. base_for_merge = plus_sess if isinstance(plus_sess, dict) and plus_sess else init_sess
  150. merged = _merge_plus_into_session(base_for_merge or {}, plan, account_id, me_email or email)
  151. # accessToken 用我们手头那个(merged 不一定包含)
  152. if not merged.get("accessToken"):
  153. merged["accessToken"] = access_token
  154. upsert_account(email, acc.get("password") or "", fields={
  155. "plan_type": plan,
  156. "plus_session": merged,
  157. })
  158. if not is_plus_session(merged):
  159. upsert_account(email, acc.get("password") or "", fields={
  160. "final_status": "plus_check_failed",
  161. "last_error": f"补救后 planType 仍为 {plan!r}",
  162. })
  163. add_event(email, "recheck", "warn", f"plan={plan}")
  164. return {"ok": False, "planType": plan, "action": "still_not_plus"}
  165. upsert_account(email, acc.get("password") or "", fields={
  166. "final_status": "plus", "last_error": ""
  167. })
  168. add_event(email, "plus_check", "ok", f"plan={plan} (recheck)")
  169. if not (cpa_url and cpa_management_key):
  170. upsert_account(email, acc.get("password") or "", fields={
  171. "final_status": "cpa_skipped"
  172. })
  173. add_event(email, "cpa", "warn", "未配置 CPA")
  174. return {"ok": True, "planType": plan, "action": "plus_no_cpa_config"}
  175. try:
  176. cpa_result = upload_session_to_cpa(
  177. merged,
  178. cpa_url=cpa_url,
  179. management_key=cpa_management_key,
  180. email_hint=email,
  181. log=log,
  182. )
  183. except Exception as exc:
  184. msg = f"CPA 上传失败: {exc}"
  185. log(f"[recheck] {email} {msg}")
  186. upsert_account(email, acc.get("password") or "", fields={
  187. "final_status": "cpa_failed",
  188. "last_error": msg,
  189. })
  190. add_event(email, "cpa", "error", msg)
  191. return {"ok": False, "planType": plan, "action": "cpa_failed", "error": msg}
  192. upsert_account(email, acc.get("password") or "", fields={
  193. "final_status": "cpa_uploaded",
  194. "cpa_file_name": cpa_result.get("fileName"),
  195. "cpa_uploaded_at": int(time.time() * 1000),
  196. "last_error": "",
  197. })
  198. add_event(email, "cpa", "ok", cpa_result.get("fileName"), payload=cpa_result)
  199. return {
  200. "ok": True,
  201. "planType": plan,
  202. "action": "cpa_uploaded",
  203. "cpa": cpa_result,
  204. }