"""本地 Web 控制台:网页配置 + 一键全自动注册→付款→上传 CPA。"""
from __future__ import annotations
import json
import queue
import threading
import time
from dataclasses import asdict
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from chatgpt_flow import FullRunContext, run_full
from config import AppConfig
from cpa_uploader import build_cpa_auth_payload
from recheck import recheck_account
from storage import (
create_task,
get_account,
get_task,
init_db,
list_accounts,
list_events,
list_tasks,
)
from task_runner import get_runner, make_task_id
HOST = "127.0.0.1" # main() 会读 cfg.api_host 覆盖
PORT = 7791
class JobManager:
def __init__(self):
self.lock = threading.Lock()
self.full_ctx: FullRunContext | None = None
self.thread: threading.Thread | None = None
self.log_queue: queue.Queue[str] = queue.Queue()
self.history: list[str] = []
self.stage: str = ""
def _log(self, msg: str):
line = f"[{time.strftime('%H:%M:%S')}] {msg}"
self.history.append(line)
if len(self.history) > 4000:
self.history = self.history[-3000:]
self.log_queue.put(line)
def _on_stage(self, name: str):
self.stage = name
# stage 也写到日志,便于复盘
self._log(f"[STAGE] {name}")
def start(self, cfg: AppConfig) -> str:
with self.lock:
if self.thread and self.thread.is_alive():
return "已有任务在运行"
self.history.clear()
while not self.log_queue.empty():
self.log_queue.get_nowait()
self.stage = ""
def runner():
try:
self.full_ctx = run_full(cfg, log=self._log, on_stage=self._on_stage)
except Exception as exc:
import traceback
self._log(f"[server] 任务异常: {exc!r}")
self._log(traceback.format_exc())
self.thread = threading.Thread(target=runner, daemon=True)
self.thread.start()
return ""
def stop(self):
if self.full_ctx:
self.full_ctx.state = "stopped"
self._log("[user] 已请求停止")
def status(self) -> dict:
running = bool(self.thread and self.thread.is_alive())
ctx = self.full_ctx
accounts = []
state = "idle"
if ctx:
state = ctx.state
for a in ctx.accounts:
accounts.append({
"email": a.get("email"),
"stage": a.get("stage"),
"planType": a.get("planType"),
"error": a.get("error"),
"cpaFile": (a.get("cpa") or {}).get("fileName") if a.get("cpa") else None,
})
return {
"running": running,
"state": state,
"stage": self.stage,
"accounts": accounts,
}
JOB = JobManager()
INDEX_HTML = r"""
ChatGPT Plus 全自动注册
ChatGPT Plus 全自动注册 + CPA 上传
流程:a4sky 邮箱注册 → 拿 Plus 长链 → PayPal 创建账号付款 → 校验 plan=plus → 上传 CPA。手机号统一 +15822201173。
配置
外网 API 接入
在线 API 文档:/docs ·
OpenAPI 规范:/openapi.json
空闲
已注册账号库
数据库:data/accounts.db
| 邮箱 |
plan |
状态 |
CPA 文件 |
注册时间 |
更新时间 |
错误 |
动作 |
点击查看选中账号详情
"""
def _read_json(handler) -> dict:
length = int(handler.headers.get("content-length") or "0")
if length <= 0:
return {}
raw = handler.rfile.read(length).decode("utf-8", errors="replace")
return json.loads(raw or "{}")
class Handler(BaseHTTPRequestHandler):
def do_GET(self):
path = self.path.split("?", 1)[0]
query = self.path.split("?", 1)[1] if "?" in self.path else ""
if path in ("/", "/index.html"):
self._send(200, INDEX_HTML.encode("utf-8"), "text/html; charset=utf-8")
return
if path in ("/docs", "/docs/"):
self._send(200, SWAGGER_HTML.encode("utf-8"), "text/html; charset=utf-8")
return
if path == "/openapi.json":
self._send_json(200, _build_openapi_spec())
return
if not self._check_auth():
return
if path == "/api/status":
self._send_json(200, JOB.status())
return
if path == "/api/config":
self._send_json(200, asdict(AppConfig.load()))
return
if path == "/api/log":
self._stream_log()
return
# ===== 任务化 API(GET)=====
if path == "/api/tasks":
from urllib.parse import parse_qs
q = parse_qs(query)
status = (q.get("status") or [""])[0] or None
limit = int((q.get("limit") or ["100"])[0])
try:
tasks = list_tasks(limit=limit, status=status)
self._send_json(200, {"tasks": tasks})
except Exception as exc:
self._send_json(500, {"error": str(exc)})
return
if path.startswith("/api/tasks/"):
from urllib.parse import unquote
task_id = unquote(path[len("/api/tasks/"):])
t = get_task(task_id)
if not t:
self._send_json(404, {"error": "task not found"})
return
self._send_json(200, {"task": t})
return
if path == "/api/accounts":
try:
from urllib.parse import parse_qs
q = parse_qs(query)
status = (q.get("status") or [""])[0] or None
limit = int((q.get("limit") or ["200"])[0])
accounts = list_accounts(limit=limit, status=status)
# 别把 session 全文 dump 给列表,太大;列表只回主要字段
slim = []
for a in accounts:
slim.append({k: a.get(k) for k in (
"email", "plan_type", "final_status", "cpa_file_name",
"long_link", "last_error", "created_at", "updated_at",
"cpa_uploaded_at"
)})
self._send_json(200, {"accounts": slim})
except Exception as exc:
self._send_json(500, {"error": str(exc)})
return
if path.startswith("/api/account/") and path.endswith("/cpa.json"):
from urllib.parse import unquote
email = unquote(path[len("/api/account/"):-len("/cpa.json")])
acc = get_account(email)
if not acc:
self._send_json(404, {"error": "account not found"})
return
session = acc.get("plus_session") or acc.get("initial_session")
if not session:
self._send_json(404, {"error": "该账号没有可下载的 session"})
return
try:
payload = build_cpa_auth_payload(session, email_hint=email)
except Exception as exc:
self._send_json(500, {"error": f"构造 CPA auth JSON 失败: {exc}"})
return
file_name = acc.get("cpa_file_name") or payload["fileName"]
content = json.dumps(payload["authJson"], ensure_ascii=False, indent=2).encode("utf-8")
self.send_response(200)
self.send_header("Content-Type", "application/json; charset=utf-8")
self.send_header("Content-Disposition", f'attachment; filename="{file_name}"')
self.send_header("Cache-Control", "no-store")
self.send_header("Content-Length", str(len(content)))
self.end_headers()
self.wfile.write(content)
return
if path.startswith("/api/account/"):
from urllib.parse import unquote
email = unquote(path[len("/api/account/"):])
acc = get_account(email)
if not acc:
self._send_json(404, {"error": "account not found"})
return
events = list_events(email, limit=200)
self._send_json(200, {"account": acc, "events": events})
return
self._send_json(404, {"error": "not found"})
def do_POST(self):
path = self.path.split("?", 1)[0]
if not self._check_auth():
return
if path == "/api/config":
try:
body = _read_json(self)
cfg = AppConfig.load().update(body or {})
self._send_json(200, asdict(cfg))
except Exception as exc:
self._send_json(500, {"error": str(exc)})
return
# ===== 任务化 API =====
if path == "/api/tasks":
try:
body = _read_json(self) or {}
mode = (body.get("mode") or "full").strip().lower()
if mode not in ("full", "pay_only"):
self._send_json(400, {"error": "mode 必须是 full 或 pay_only"})
return
params = body.get("params") or {}
if mode == "pay_only":
sess = params.get("session")
if not isinstance(sess, dict) or not sess.get("accessToken"):
self._send_json(400, {"error": "pay_only 需要 params.session 是 JSON 且包含 accessToken"})
return
max_attempts = int(body.get("max_attempts") or 3)
max_attempts = max(1, min(10, max_attempts))
task_id = make_task_id()
t = create_task(task_id, mode, params, max_attempts=max_attempts)
# 启动 runner(幂等)
get_runner(log=lambda m: JOB._log(m))
self._send_json(200, {"task_id": task_id, "task": t})
except Exception as exc:
self._send_json(500, {"error": str(exc)})
return
if path.startswith("/api/tasks/") and path.endswith("/cancel"):
from urllib.parse import unquote
task_id = unquote(path[len("/api/tasks/"):-len("/cancel")])
t = get_task(task_id)
if not t:
self._send_json(404, {"error": "task not found"})
return
runner = get_runner(log=lambda m: JOB._log(m))
runner.cancel(task_id)
self._send_json(200, {"ok": True, "task_id": task_id})
return
if path == "/api/start":
try:
cfg = AppConfig.load()
err = JOB.start(cfg)
if err:
self._send_json(409, {"error": err})
else:
self._send_json(200, {"ok": True})
except Exception as exc:
self._send_json(500, {"error": str(exc)})
return
if path.startswith("/api/account/") and path.endswith("/recheck"):
from urllib.parse import unquote
email = unquote(path[len("/api/account/"):-len("/recheck")])
cfg = AppConfig.load()
try:
result = recheck_account(
email,
cpa_url=cfg.cpa_url,
cpa_management_key=cfg.cpa_management_key,
log=lambda msg: JOB._log(f"[acc:{email[:24]}] {msg}"),
)
self._send_json(200, result)
except Exception as exc:
self._send_json(500, {"error": str(exc)})
return
if path == "/api/stop":
JOB.stop()
self._send_json(200, {"ok": True})
return
self._send_json(404, {"error": "not found"})
def _stream_log(self):
self.send_response(200)
self.send_header("Content-Type", "text/event-stream; charset=utf-8")
self.send_header("Cache-Control", "no-cache")
self.send_header("Connection", "keep-alive")
self.end_headers()
try:
for line in JOB.history[-300:]:
self._sse_send(line)
while True:
try:
line = JOB.log_queue.get(timeout=15)
self._sse_send(line)
except queue.Empty:
self.wfile.write(b": ping\n\n")
self.wfile.flush()
except (BrokenPipeError, ConnectionResetError):
return
def _sse_send(self, line: str):
for piece in line.splitlines() or [""]:
self.wfile.write(b"data: " + piece.encode("utf-8") + b"\n")
self.wfile.write(b"\n")
self.wfile.flush()
def _send_json(self, status: int, payload: dict):
self._send(status, json.dumps(payload, ensure_ascii=False).encode("utf-8"), "application/json; charset=utf-8")
def _send(self, status: int, content: bytes, content_type: str):
self.send_response(status)
self.send_header("Content-Type", content_type)
self.send_header("Cache-Control", "no-store")
self.send_header("Content-Length", str(len(content)))
# CORS(仅 /api/* 需要时由调用方决定,但统一发也无害)
try:
cfg = AppConfig.load()
origin = (cfg.api_cors_origin or "*").strip()
self.send_header("Access-Control-Allow-Origin", origin)
self.send_header("Access-Control-Allow-Credentials", "true")
except Exception:
self.send_header("Access-Control-Allow-Origin", "*")
self.end_headers()
self.wfile.write(content)
def do_OPTIONS(self):
# CORS preflight
self.send_response(204)
try:
cfg = AppConfig.load()
origin = (cfg.api_cors_origin or "*").strip()
except Exception:
origin = "*"
self.send_header("Access-Control-Allow-Origin", origin)
self.send_header("Access-Control-Allow-Methods", "GET, POST, OPTIONS")
self.send_header("Access-Control-Allow-Headers", "Content-Type, Authorization")
self.send_header("Access-Control-Max-Age", "86400")
self.send_header("Access-Control-Allow-Credentials", "true")
self.end_headers()
def _check_auth(self) -> bool:
"""非空 api_token 时校验 Authorization: Bearer。返回 True 表示放行。"""
try:
cfg = AppConfig.load()
token = (cfg.api_token or "").strip()
except Exception:
token = ""
if not token:
return True
# 公开接口豁免:根页面、OpenAPI 文档、Swagger UI、static 静态
path = self.path.split("?", 1)[0]
public = ("/", "/index.html", "/docs", "/docs/", "/openapi.json", "/openapi.yaml")
if path in public:
return True
auth = self.headers.get("Authorization", "")
if auth == f"Bearer {token}":
return True
# 也支持 ?token=xxx
if "token=" in (self.path.split("?", 1)[1] if "?" in self.path else ""):
from urllib.parse import parse_qs
q = parse_qs(self.path.split("?", 1)[1])
if (q.get("token") or [""])[0] == token:
return True
self._send_json(401, {"error": "missing or invalid Bearer token"})
return False
def log_message(self, fmt, *args):
return
def handle_one_request(self):
try:
return super().handle_one_request()
except (ConnectionResetError, BrokenPipeError):
# 浏览器主动断开 SSE / fetch 时打印栈很碍眼,直接静音
self.close_connection = True
SWAGGER_HTML = r"""
API 文档 · ChatGPT Plus 自动化
"""
def _build_openapi_spec() -> dict:
"""生成 OpenAPI 3.1 规范。"""
cfg = AppConfig.load()
return {
"openapi": "3.1.0",
"info": {
"title": "ChatGPT Plus 自动化 API",
"version": "1.0.0",
"description": (
"ChatGPT Plus 全自动注册 + PayPal 付款 + CPA 上传。\n\n"
"**两种模式**:\n"
"- `full` — 全自动注册新 ChatGPT 账号,注册→付款→上传 CPA\n"
"- `pay_only` — 传入已有 session JSON,跳过注册直接付款→上传 CPA\n\n"
"**调用流程**:\n"
"1. POST `/api/tasks` 创建任务,立即拿到 `task_id`\n"
"2. 轮询 GET `/api/tasks/{task_id}` 看 `status` 和 `stage`\n"
"3. `status` 变 `success` 时可调 GET `/api/account/{email}/cpa.json` 下载 CPA 文件\n\n"
"**重试**:每个任务整体失败会重试 `max_attempts` 次(默认 3)。"
),
},
"servers": [
{"url": f"http://{cfg.api_host or '127.0.0.1'}:{cfg.api_port or 7791}", "description": "当前实例"},
],
"components": {
"securitySchemes": {
"BearerAuth": {
"type": "http",
"scheme": "bearer",
"description": "如果配置了 `api_token`,所有 /api/* 请求需带 `Authorization: Bearer `。也支持 `?token=xxx` 查询参数。",
}
},
"schemas": {
"Task": {
"type": "object",
"properties": {
"task_id": {"type": "string", "example": "t-1779470219-65e44d55"},
"mode": {"type": "string", "enum": ["full", "pay_only"]},
"status": {"type": "string", "enum": ["queued", "running", "success", "failed", "cancelled"]},
"stage": {"type": "string", "description": "当前阶段描述"},
"attempts": {"type": "integer"},
"max_attempts": {"type": "integer"},
"params": {"type": "object", "description": "创建任务时传入的参数(脱敏后)"},
"result": {"type": "object", "nullable": True, "description": "成功时的结果(含 CPA 文件名等)"},
"last_error": {"type": "string", "nullable": True},
"email": {"type": "string", "nullable": True, "description": "注册成功的 ChatGPT 邮箱"},
"plan_type": {"type": "string", "nullable": True, "example": "plus"},
"cpa_file_name": {"type": "string", "nullable": True, "example": "codex-foo@example.com-plus.json"},
"created_at": {"type": "integer", "description": "毫秒时间戳"},
"updated_at": {"type": "integer"},
"started_at": {"type": "integer", "nullable": True},
"finished_at": {"type": "integer", "nullable": True},
},
},
"CreateTaskRequest": {
"type": "object",
"required": ["mode"],
"properties": {
"mode": {"type": "string", "enum": ["full", "pay_only"]},
"max_attempts": {"type": "integer", "default": 3, "minimum": 1, "maximum": 10},
"params": {
"type": "object",
"description": "可覆盖全局配置;pay_only 模式必须包含 session 字段",
"properties": {
"session": {
"type": "object",
"description": "ChatGPT /api/auth/session 完整 JSON(仅 pay_only 模式必填)",
"properties": {
"accessToken": {"type": "string"},
"user": {"type": "object"},
"account": {"type": "object"},
},
},
"headless": {"type": "boolean"},
"use_promo": {"type": "boolean"},
"phone_e164": {"type": "string", "example": "+15822201173"},
"sms_api_url": {"type": "string"},
"cpa_url": {"type": "string"},
"cpa_management_key": {"type": "string"},
"proxy_url": {"type": "string"},
"paypal_only_proxy": {"type": "string"},
"mail_helper_url": {"type": "string"},
"mail_domain": {"type": "string"},
},
},
},
},
"Account": {
"type": "object",
"properties": {
"email": {"type": "string"},
"plan_type": {"type": "string", "nullable": True},
"final_status": {"type": "string"},
"cpa_file_name": {"type": "string", "nullable": True},
"long_link": {"type": "string", "nullable": True},
"last_error": {"type": "string", "nullable": True},
"created_at": {"type": "integer"},
"updated_at": {"type": "integer"},
"cpa_uploaded_at": {"type": "integer", "nullable": True},
},
},
"Error": {
"type": "object",
"properties": {"error": {"type": "string"}},
},
},
},
"security": [{"BearerAuth": []}] if cfg.api_token else [],
"paths": {
"/api/tasks": {
"post": {
"tags": ["Tasks"],
"summary": "创建任务",
"description": "创建一个 full 或 pay_only 任务,立即返回 task_id,任务异步执行。",
"requestBody": {
"required": True,
"content": {
"application/json": {
"schema": {"$ref": "#/components/schemas/CreateTaskRequest"},
"examples": {
"full": {
"summary": "全自动注册",
"value": {
"mode": "full",
"max_attempts": 3,
"params": {},
},
},
"pay_only": {
"summary": "传入 session 直接付款",
"value": {
"mode": "pay_only",
"max_attempts": 3,
"params": {
"session": {
"accessToken": "eyJxxx...",
"user": {"email": "user@example.com"},
"account": {"planType": "free"},
}
},
},
},
},
}
},
},
"responses": {
"200": {
"description": "任务已创建",
"content": {
"application/json": {
"schema": {
"type": "object",
"properties": {
"task_id": {"type": "string"},
"task": {"$ref": "#/components/schemas/Task"},
},
}
}
},
},
"400": {"description": "参数错误", "content": {"application/json": {"schema": {"$ref": "#/components/schemas/Error"}}}},
"401": {"description": "Bearer token 缺失或无效"},
},
},
"get": {
"tags": ["Tasks"],
"summary": "列出任务",
"parameters": [
{"name": "status", "in": "query", "schema": {"type": "string", "enum": ["queued", "running", "success", "failed", "cancelled"]}},
{"name": "limit", "in": "query", "schema": {"type": "integer", "default": 100}},
],
"responses": {
"200": {
"content": {
"application/json": {
"schema": {
"type": "object",
"properties": {"tasks": {"type": "array", "items": {"$ref": "#/components/schemas/Task"}}},
}
}
}
}
},
},
},
"/api/tasks/{task_id}": {
"get": {
"tags": ["Tasks"],
"summary": "查询任务进度",
"description": "轮询此接口查任务实时 status 和 stage。建议 5-10 秒间隔。",
"parameters": [{"name": "task_id", "in": "path", "required": True, "schema": {"type": "string"}}],
"responses": {
"200": {"content": {"application/json": {"schema": {"type": "object", "properties": {"task": {"$ref": "#/components/schemas/Task"}}}}}},
"404": {"description": "任务不存在"},
},
}
},
"/api/tasks/{task_id}/cancel": {
"post": {
"tags": ["Tasks"],
"summary": "取消任务",
"description": "请求取消任务。如果任务已经在跑,会在下一个 stop 检查点退出。",
"parameters": [{"name": "task_id", "in": "path", "required": True, "schema": {"type": "string"}}],
"responses": {"200": {"description": "已请求取消"}, "404": {"description": "任务不存在"}},
}
},
"/api/accounts": {
"get": {
"tags": ["Accounts"],
"summary": "列出已注册账号",
"parameters": [
{"name": "status", "in": "query", "schema": {"type": "string"}, "description": "如 cpa_uploaded / plus_check_failed"},
{"name": "limit", "in": "query", "schema": {"type": "integer", "default": 200}},
],
"responses": {
"200": {
"content": {
"application/json": {
"schema": {
"type": "object",
"properties": {"accounts": {"type": "array", "items": {"$ref": "#/components/schemas/Account"}}},
}
}
}
}
},
}
},
"/api/account/{email}": {
"get": {
"tags": ["Accounts"],
"summary": "查询账号详情(含完整 session 和事件流)",
"parameters": [{"name": "email", "in": "path", "required": True, "schema": {"type": "string"}}],
"responses": {"200": {"description": "OK"}, "404": {"description": "账号不存在"}},
}
},
"/api/account/{email}/cpa.json": {
"get": {
"tags": ["Accounts"],
"summary": "下载 CPA codex auth JSON",
"description": "返回该账号当时上传给 CPA 的完整 codex auth JSON 文件。带 Content-Disposition 头,浏览器会自动下载。",
"parameters": [{"name": "email", "in": "path", "required": True, "schema": {"type": "string"}}],
"responses": {
"200": {"description": "OK", "content": {"application/json": {}}},
"404": {"description": "账号或 session 不存在"},
},
}
},
"/api/account/{email}/recheck": {
"post": {
"tags": ["Accounts"],
"summary": "对失败账号补救",
"description": "用 DB 里存的 access_token 调 backend-api/me,若已 plus 则自动重传 CPA。",
"parameters": [{"name": "email", "in": "path", "required": True, "schema": {"type": "string"}}],
"responses": {"200": {"description": "OK"}, "404": {"description": "账号不存在"}},
}
},
"/api/config": {
"get": {"tags": ["Config"], "summary": "读取当前配置", "responses": {"200": {"description": "OK"}}},
"post": {
"tags": ["Config"],
"summary": "更新配置",
"requestBody": {"content": {"application/json": {"schema": {"type": "object"}}}},
"responses": {"200": {"description": "OK"}},
},
},
"/api/status": {
"get": {"tags": ["Misc"], "summary": "(旧)读取 UI 任务状态", "responses": {"200": {"description": "OK"}}}
},
"/api/log": {
"get": {"tags": ["Misc"], "summary": "实时日志(Server-Sent Events)", "responses": {"200": {"description": "text/event-stream"}}}
},
},
"tags": [
{"name": "Tasks", "description": "任务化 API(推荐用法)"},
{"name": "Accounts", "description": "账号库"},
{"name": "Config", "description": "服务配置"},
{"name": "Misc", "description": "其他"},
],
}
def _silence_threading_excepthook():
"""ThreadingHTTPServer 在 worker 线程里仍可能抛 ConnectionResetError;接住它。"""
import threading
prev = threading.excepthook
def hook(args):
if isinstance(args.exc_value, (ConnectionResetError, BrokenPipeError)):
return
prev(args)
threading.excepthook = hook
def main():
init_db()
_silence_threading_excepthook()
# 启动后台任务 worker(幂等)
get_runner(log=lambda m: JOB._log(m))
cfg = AppConfig.load()
host = (cfg.api_host or HOST).strip() or HOST
port = int(cfg.api_port or PORT)
server = ThreadingHTTPServer((host, port), Handler)
print(f"ChatGPT Plus Auto Console:")
print(f" Web UI: http://{host}:{port}/")
print(f" Docs: http://{host}:{port}/docs")
print(f" OpenAPI: http://{host}:{port}/openapi.json")
if cfg.api_token:
print(f" Auth: Bearer (已启用)")
if host == "0.0.0.0":
print(f" ⚠️ 当前监听所有网卡,外网可访问。建议设置 api_token。")
try:
server.serve_forever()
except KeyboardInterrupt:
print("\nStopped.")
if __name__ == "__main__":
main()