AutoTeam-F / src /autoteam /chatgpt_api.py
ZRainbow's picture
fix: complete autoteam-1 parity migration
ceeba19
Raw
History Blame Contribute Delete
82.1 kB
"""ChatGPT Team API 客户端 - 通过 Playwright 绕过 Cloudflare 调用内部 API"""
import base64
import json
import logging
import random
import re
import time
import uuid
from pathlib import Path
from playwright.sync_api import sync_playwright
import autoteam.display # noqa: F401
from autoteam.admin_state import (
get_admin_email,
get_admin_session_token,
get_chatgpt_account_id,
get_chatgpt_workspace_name,
update_admin_state,
)
from autoteam.chatgpt_transport import build_chatgpt_transport
from autoteam.config import get_playwright_context_options, get_playwright_launch_options
from autoteam.playwright_lifecycle import close_playwright_objects
from autoteam.textio import read_text
logger = logging.getLogger(__name__)
PROJECT_ROOT = Path(__file__).parent.parent.parent
BASE_DIR = PROJECT_ROOT
SCREENSHOT_DIR = PROJECT_ROOT / "screenshots"
_WORKSPACE_IGNORE_LABELS = {
"choose a workspace",
"select a workspace",
"选择一个工作空间",
"选择工作空间",
"workspace",
"terms of use",
"privacy policy",
"使用条款",
"隐私政策",
}
_WORKSPACE_FALLBACK_LABELS = (
"personal account",
"personal",
"个人账户",
"个人账号",
"free",
"免费",
"new organization",
"新组织",
"create organization",
"创建组织",
)
def _normalize_workspace_label(text: str | None) -> str:
return " ".join((text or "").split()).strip()
def _workspace_candidate_kind(text: str | None) -> str | None:
"""Classify a visible workspace label as preferred, fallback, or noise."""
label = _normalize_workspace_label(text)
if not label:
return None
lowered = label.lower()
if lowered in _WORKSPACE_IGNORE_LABELS or len(label) <= 3:
return None
if any(key in lowered for key in _WORKSPACE_FALLBACK_LABELS):
return "fallback"
return "preferred"
class ChatGPTTeamAPI:
"""通过浏览器内 fetch 调用 ChatGPT Team 内部 API。"""
EMAIL_INPUT_SELECTORS = [
'input[name="email"]',
'input[id="email-input"]',
'input[id="email"]',
'input[type="email"]',
'input[placeholder*="email" i]',
'input[placeholder*="邮箱"]',
'input[autocomplete="email"]',
'input[autocomplete="username"]',
]
PASSWORD_INPUT_SELECTORS = [
'input[name="password"]',
'input[type="password"]',
]
CODE_INPUT_SELECTORS = [
'input[name="code"]',
'input[placeholder*="验证码"]',
'input[placeholder*="code" i]',
'input[inputmode="numeric"]',
'input[autocomplete="one-time-code"]',
]
_UUID_RE = re.compile(r"^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$", re.I)
_WORKSPACE_NAME_EXCLUDES = {
"常规",
"成员",
"设置",
"帮助",
"general",
"members",
"settings",
"help",
"chat history",
"new chat",
"search chats",
"images",
"apps",
"deep research",
"see plans and pricing",
"log in",
"chatgpt",
}
_WORKSPACE_PAGE_HINTS = (
"launch a workspace",
"has access to",
"workspace",
"personal workspace",
"选择工作空间",
"选择一个工作空间",
)
def __init__(self):
self.playwright = None
self.browser = None
self.context = None
self.page = None
self.access_token = None
self.session_token = None
self.account_id = get_chatgpt_account_id()
self.workspace_name = get_chatgpt_workspace_name()
self.oai_device_id = str(uuid.uuid4())
self.login_email = None
self.login_password = None
self.workspace_options_cache = []
self.http_transport = None
self.transport_name = None
self.proxy_url = ""
self.local_proxy_url = ""
def _visible_locator_in_frames(self, selectors, timeout_ms=5000):
selector = ", ".join(selectors)
deadline = time.time() + timeout_ms / 1000
while time.time() < deadline:
frames = [self.page.main_frame]
frames.extend(frame for frame in self.page.frames if frame != self.page.main_frame)
for frame in frames:
try:
locator = frame.locator(selector).first
if locator.is_visible(timeout=250):
return locator
except Exception:
pass
time.sleep(0.2)
return None
def _launch_browser(self):
SCREENSHOT_DIR.mkdir(exist_ok=True)
if self.playwright or self.browser or self.context or self.page:
self.stop()
try:
self.playwright = sync_playwright().start()
try:
launch_options = get_playwright_launch_options(proxy_url=self.local_proxy_url)
except TypeError as exc:
if "proxy_url" not in str(exc):
raise
launch_options = get_playwright_launch_options()
self.browser = self.playwright.chromium.launch(**launch_options)
self.context = self.browser.new_context(**get_playwright_context_options())
self.page = self.context.new_page()
except Exception:
self.stop()
raise
def _ensure_admin_ipv6_proxy(self) -> None:
admin_email = (get_admin_email() or "").strip().lower()
self.proxy_url = ""
self.local_proxy_url = ""
if not admin_email:
required = False
try:
from autoteam import config as runtime_config
required = bool(getattr(runtime_config, "AUTOTEAM_IPV6_POOL_REQUIRED", False))
except Exception:
pass
if required:
raise RuntimeError("IPv6 pool is required but admin email is unavailable")
return
required = False
try:
from autoteam import config as runtime_config
from autoteam.ipv6_pool import ipv6_pool
required = bool(getattr(runtime_config, "AUTOTEAM_IPV6_POOL_REQUIRED", False))
self.proxy_url = ipv6_pool.ensure(admin_email) or ""
self.local_proxy_url = ipv6_pool.get_local_proxy_url(admin_email) or self.proxy_url
if self.local_proxy_url:
logger.info("[ChatGPT] admin %s using IPv6 proxy %s", admin_email, self.local_proxy_url)
return
if required:
raise RuntimeError("IPv6 pool is required but no admin proxy was assigned")
except Exception as exc:
if required:
logger.error("[ChatGPT] admin IPv6 proxy required but unavailable: %s", exc)
raise
self.proxy_url = ""
self.local_proxy_url = ""
logger.warning("[ChatGPT] admin IPv6 proxy unavailable, falling back to direct: %s", exc)
def _log_login_state(self, label):
try:
body_excerpt = self.page.locator("body").inner_text(timeout=1500)[:300].replace("\n", " ")
except Exception:
body_excerpt = ""
logger.info(
"[ChatGPT] %s | URL=%s | body=%s",
label,
self.page.url,
body_excerpt,
)
def _wait_for_cloudflare(self):
for i in range(12):
html = self.page.content()[:1000].lower()
if "verify you are human" not in html and "challenge" not in self.page.url:
return
logger.info("[ChatGPT] 等待 Cloudflare... (%ds)", i * 5)
time.sleep(5)
def _build_session_cookies(self, session_token, domain):
if len(session_token) > 3800:
return [
{
"name": "__Secure-next-auth.session-token.0",
"value": session_token[:3800],
"domain": domain,
"path": "/",
"httpOnly": True,
"secure": True,
"sameSite": "Lax",
},
{
"name": "__Secure-next-auth.session-token.1",
"value": session_token[3800:],
"domain": domain,
"path": "/",
"httpOnly": True,
"secure": True,
"sameSite": "Lax",
},
]
return [
{
"name": "__Secure-next-auth.session-token",
"value": session_token,
"domain": domain,
"path": "/",
"httpOnly": True,
"secure": True,
"sameSite": "Lax",
}
]
def _click_auth_button(self, field, labels):
label_re = re.compile(rf"^(?:{'|'.join(re.escape(label) for label in labels)})$", re.I)
try:
form = field.locator("xpath=ancestor::form[1]").first
btn = form.get_by_role("button", name=label_re).first
if btn.is_visible(timeout=2000):
btn.click()
return True
except Exception:
pass
try:
form = field.locator("xpath=ancestor::form[1]").first
btn = form.locator('button[type="submit"], input[type="submit"]').first
if btn.is_visible(timeout=2000):
btn.click()
return True
except Exception:
pass
try:
btn = self.page.get_by_role("button", name=label_re).last
if btn.is_visible(timeout=2000):
btn.click()
return True
except Exception:
pass
try:
field.press("Enter")
return True
except Exception:
return False
def _body_excerpt(self, limit=300):
try:
return self.page.locator("body").inner_text(timeout=1500)[:limit].replace("\n", " ")
except Exception:
return ""
def _wait_for_login_step(self, allowed_steps, timeout=15):
deadline = time.time() + timeout
while time.time() < deadline:
step, detail = self._detect_login_step()
if step in allowed_steps:
return step, detail
if "challenge" in (self.page.url or "").lower():
self._wait_for_cloudflare()
time.sleep(0.5)
return self._detect_login_step()
def _extract_session_token(self):
cookies = self.context.cookies()
session_parts = {}
session_token = None
for cookie in cookies:
name = cookie["name"]
if name == "__Secure-next-auth.session-token":
session_token = cookie["value"]
elif name.startswith("__Secure-next-auth.session-token."):
suffix = name.rsplit(".", 1)[-1]
session_parts[suffix] = cookie["value"]
if not session_token and session_parts:
session_token = "".join(session_parts[k] for k in sorted(session_parts))
self.session_token = session_token
return session_token
def _extract_account_id_from_access_token(self):
if not self.access_token:
return ""
try:
parts = self.access_token.split(".")
if len(parts) < 2:
return ""
payload_part = parts[1] + "=" * (-len(parts[1]) % 4)
payload = json.loads(base64.urlsafe_b64decode(payload_part))
except Exception:
payload = {}
auth_claims = payload.get("https://api.openai.com/auth", {}) if isinstance(payload, dict) else {}
account_id = auth_claims.get("chatgpt_account_id", "") if isinstance(auth_claims, dict) else ""
if account_id and self._UUID_RE.match(str(account_id)):
return str(account_id)
return ""
def _detect_workspace_name_from_dom(self):
if not self.page:
return ""
try:
name = self.page.evaluate(
"""(excludes) => {
const excludeSet = new Set((excludes || []).map(x => String(x).toLowerCase().trim()));
const selectors = [
'main h1', 'main h2', 'main h3',
'[role="main"] h1', '[role="main"] h2', '[role="main"] h3',
'main [class*="title"]', 'main [class*="name"]',
'[role="main"] [class*="title"]', '[role="main"] [class*="name"]',
'h1', 'h2', 'h3',
'[class*="title"]', '[class*="name"]',
];
const seen = new Set();
for (const selector of selectors) {
const nodes = document.querySelectorAll(selector);
for (const node of nodes) {
const text = (node.textContent || '').trim().replace(/\\s+/g, ' ');
const lower = text.toLowerCase();
if (!text || text.length < 2 || text.length > 60) continue;
if (excludeSet.has(lower)) continue;
if (
lower.includes('chat history') ||
lower.includes('new chat') ||
lower.includes('search chats') ||
lower.includes('see plans and pricing') ||
lower.includes('log in to get answers based on saved chats')
) continue;
if (seen.has(lower)) continue;
seen.add(lower);
return text;
}
}
return '';
}""",
sorted(self._WORKSPACE_NAME_EXCLUDES),
)
except Exception:
name = ""
return (name or "").strip()
def _is_workspace_selection_page(self):
url = (self.page.url or "").lower()
if "workspace" in url or "organization" in url:
return True
try:
body = self.page.locator("body").inner_text(timeout=1500).lower()
except Exception:
body = ""
hint_hits = sum(1 for hint in self._WORKSPACE_PAGE_HINTS if hint in body)
return hint_hits >= 2 or ("launch a workspace" in body)
def _auto_open_preferred_workspace(self):
"""在 workspace 选择页自动点击优先的 Team workspace。"""
if not self.page:
return False
try:
result = self.page.evaluate(
"""() => {
const norm = (s) => (s || '').replace(/\\s+/g, ' ').trim();
const lower = (s) => norm(s).toLowerCase();
const badKeywords = ['personal workspace', 'personal account', 'personal', 'free', '免费', '个人'];
const buttons = Array.from(document.querySelectorAll('button, a, [role="button"]'));
const candidates = [];
for (const btn of buttons) {
const btnText = lower(btn.textContent);
if (!btnText || !['open', '打开', 'launch', 'continue', '进入'].some(k => btnText.includes(k))) continue;
let card = btn;
for (let i = 0; i < 6 && card && card.parentElement; i++) {
card = card.parentElement;
const text = norm(card.textContent);
if (!text || text.length < 4) continue;
if (text.length > 200) continue;
const textLower = text.toLowerCase();
const bad = badKeywords.some(k => textLower.includes(k));
candidates.push({
text,
bad,
score: (bad ? 0 : 100) + Math.min(text.length, 80),
buttonIndex: buttons.indexOf(btn),
});
break;
}
}
if (!candidates.length) return { clicked: false, reason: 'no-candidates' };
candidates.sort((a, b) => b.score - a.score);
const chosen = candidates[0];
const btn = buttons[chosen.buttonIndex];
if (!btn) return { clicked: false, reason: 'missing-button' };
btn.click();
return { clicked: true, label: chosen.text, bad: chosen.bad };
}"""
)
except Exception as exc:
logger.warning("[ChatGPT] 自动点击 workspace 失败: %s", exc)
return False
if result and result.get("clicked"):
logger.info("[ChatGPT] 自动进入 workspace: %s", result.get("label", ""))
time.sleep(5)
self._wait_for_cloudflare()
self._log_login_state("自动进入 workspace 后")
return True
logger.warning("[ChatGPT] 未找到可自动进入的 workspace: %s", (result or {}).get("reason"))
return False
def _inject_session(self, session_token):
cookies = self._build_session_cookies(session_token, "chatgpt.com")
if self.account_id:
cookies.append(
{
"name": "_account",
"value": self.account_id,
"domain": "chatgpt.com",
"path": "/",
"secure": True,
"sameSite": "Lax",
}
)
cookies.append(
{
"name": "oai-did",
"value": self.oai_device_id,
"domain": "chatgpt.com",
"path": "/",
"secure": True,
"sameSite": "Lax",
}
)
self.context.add_cookies(cookies)
self.session_token = session_token
logger.info("[ChatGPT] 已注入 session cookies")
def _open_login_page(self):
self.page.goto("https://chatgpt.com/auth/login", wait_until="domcontentloaded", timeout=60000)
time.sleep(5)
self._wait_for_cloudflare()
self._log_login_state("打开登录页后")
try:
login_btn = self.page.locator('button:has-text("登录"), button:has-text("Log in")').first
if login_btn.is_visible(timeout=3000):
login_btn.click()
time.sleep(2)
self._log_login_state("点击登录按钮后")
except Exception:
pass
def _list_workspace_options(self):
if not self._is_workspace_selection_page():
return []
logger.info("[ChatGPT] 检测到 workspace 选择页,开始收集组织候选 | URL=%s", self.page.url)
try:
self.page.screenshot(path=str(SCREENSHOT_DIR / "admin_login_workspace_before_select.png"), full_page=True)
except Exception:
pass
candidates = []
seen_texts = set()
# 先用 JS 从 DOM 提取可见的 workspace 选项(只取叶子级别文本)
try:
js_candidates = self.page.evaluate("""() => {
const results = [];
const seen = new Set();
// 遍历所有元素,找"直接文本内容"短且有意义的
for (const el of document.querySelectorAll('*')) {
// 直接文本 = 不含子元素的纯文本
const directText = Array.from(el.childNodes)
.filter(n => n.nodeType === 3)
.map(n => n.textContent.trim())
.filter(t => t.length > 0)
.join(' ');
if (!directText || directText.length > 50 || directText.length < 2) continue;
if (seen.has(directText)) continue;
// 必须可见
const rect = el.getBoundingClientRect();
if (rect.width === 0 || rect.height === 0) continue;
// 排除标题类
const tag = el.tagName.toLowerCase();
if (['h1', 'h2', 'h3', 'title', 'head', 'script', 'style'].includes(tag)) continue;
seen.add(directText);
results.push(directText);
}
return results;
}""")
# 过滤掉标题和无关文本
title_keywords = (
"选择一个工作空间",
"select a workspace",
"选择工作空间",
"工作空间",
"workspace",
"chatgpt",
)
for text in js_candidates or []:
text = _normalize_workspace_label(text)
if text in seen_texts:
continue
# 跳过标题类文本
if text.lower() in (k.lower() for k in title_keywords):
continue
kind = _workspace_candidate_kind(text)
if not kind:
continue
seen_texts.add(text)
candidates.append({"id": str(len(candidates)), "label": text, "kind": kind})
except Exception as e:
logger.warning("[ChatGPT] JS 提取 workspace 候选失败: %s", e)
# fallback: Playwright 选择器
if not candidates:
for selector in ("button", '[role="button"]', "a", '[role="option"]', "div[class] > span", "li"):
try:
for loc in self.page.locator(selector).all():
try:
if not loc.is_visible(timeout=200):
continue
text = _normalize_workspace_label(loc.inner_text(timeout=500))
except Exception:
continue
if not text or text in seen_texts or len(text) > 80:
continue
kind = _workspace_candidate_kind(text)
if not kind:
continue
seen_texts.add(text)
candidates.append({"id": str(len(candidates)), "label": text, "kind": kind})
except Exception:
pass
logger.info("[ChatGPT] workspace 候选数: %d | candidates=%s", len(candidates), [c["label"] for c in candidates])
self.workspace_options_cache = candidates
return candidates
def list_workspace_options(self):
if self.workspace_options_cache:
return self.workspace_options_cache
return self._list_workspace_options()
def _click_workspace_option_by_label(self, label):
if not self.page:
return False
try:
result = self.page.evaluate(
"""(label) => {
const norm = (s) => (s || '').replace(/\\s+/g, ' ').trim();
const lower = (s) => norm(s).toLowerCase();
const targetLabel = lower(label);
const actionWords = ['open', '打开', 'launch', 'continue', '进入'];
const interactive = Array.from(document.querySelectorAll('button, a, [role="button"]'));
const candidates = [];
const collectCandidate = (el, cardText, depth) => {
const idx = interactive.indexOf(el);
if (idx < 0) return;
const btnText = norm(el.textContent);
const btnLower = btnText.toLowerCase();
let score = 100 - depth;
if (btnLower === targetLabel) score += 30;
if (actionWords.some(word => btnLower.includes(word))) score += 100;
if (lower(cardText).startsWith(targetLabel)) score += 15;
candidates.push({
buttonIndex: idx,
buttonText: btnText,
cardText: norm(cardText).slice(0, 160),
score,
});
};
for (const el of interactive) {
let node = el;
for (let depth = 0; depth < 7 && node && node.parentElement; depth++) {
node = node.parentElement;
const cardText = norm(node.textContent);
if (!cardText) continue;
if (!lower(cardText).includes(targetLabel)) continue;
collectCandidate(el, cardText, depth);
break;
}
}
if (!candidates.length) {
const allNodes = Array.from(document.querySelectorAll('*'));
for (const node of allNodes) {
const text = norm(node.textContent);
if (lower(text) !== targetLabel) continue;
let parent = node;
for (let depth = 0; depth < 7 && parent && parent.parentElement; depth++) {
parent = parent.parentElement;
const buttons = Array.from(parent.querySelectorAll('button, a, [role="button"]'));
if (!buttons.length) continue;
for (const btn of buttons) {
collectCandidate(btn, parent.textContent || text, depth + 10);
}
if (buttons.length) break;
}
}
}
if (!candidates.length) {
return { clicked: false, reason: 'no-match' };
}
candidates.sort((a, b) => b.score - a.score);
const chosen = candidates[0];
const btn = interactive[chosen.buttonIndex];
if (!btn) {
return { clicked: false, reason: 'missing-button', chosen };
}
btn.click();
return {
clicked: true,
buttonText: chosen.buttonText,
cardText: chosen.cardText,
candidateCount: candidates.length,
};
}""",
label,
)
except Exception as exc:
logger.warning("[ChatGPT] JS 点击 workspace(%s) 失败: %s", label, exc)
result = None
if result and result.get("clicked"):
logger.info(
"[ChatGPT] 点击 workspace 动作成功: label=%s button=%s candidates=%s card=%s",
label,
result.get("buttonText", ""),
result.get("candidateCount", 0),
result.get("cardText", ""),
)
return True
try:
action_re = re.compile(r"(open|打开|launch|continue|进入)", re.I)
container = self.page.locator(f"text={label}").first.locator(
"xpath=ancestor::*[self::div or self::li or self::section][1]"
)
action_btn = container.get_by_role("button", name=action_re).first
if action_btn.is_visible(timeout=2000):
action_btn.click(force=True)
logger.info("[ChatGPT] Playwright fallback 点击 workspace 动作: %s", label)
return True
except Exception:
pass
try:
loc = self.page.locator(f"text={label}").first
if loc.is_visible(timeout=2000):
loc.click(force=True)
logger.info("[ChatGPT] Playwright fallback 点击 workspace 文本: %s", label)
return True
except Exception:
pass
logger.warning("[ChatGPT] 未找到可点击的 workspace 动作: %s | reason=%s", label, (result or {}).get("reason"))
return False
def _wait_for_workspace_selection_exit(self, timeout=15):
deadline = time.time() + timeout
last_url = self.page.url if self.page else ""
while time.time() < deadline:
if not self.page:
return False
try:
self.page.wait_for_load_state("domcontentloaded", timeout=1000)
except Exception:
pass
if "challenge" in (self.page.url or "").lower():
self._wait_for_cloudflare()
if not self._is_workspace_selection_page():
return True
last_url = self.page.url
time.sleep(0.5)
logger.warning("[ChatGPT] workspace 点击后仍停留在选择页 | URL=%s", last_url)
return False
def _is_chatgpt_ready_url(self):
if not self.page:
return False
url = (self.page.url or "").lower()
return "chatgpt.com" in url and "auth" not in url
def _wait_for_post_workspace_ready(self, timeout=12):
deadline = time.time() + timeout
last_url = self.page.url if self.page else ""
consecutive_chatgpt_ready = 0
while time.time() < deadline:
if not self.page:
return False
try:
self.page.wait_for_load_state("domcontentloaded", timeout=1000)
except Exception:
pass
url = (self.page.url or "").lower()
if "challenge" in url:
self._wait_for_cloudflare()
session_token = self._extract_session_token()
if session_token:
return True
if self._is_chatgpt_ready_url():
consecutive_chatgpt_ready += 1
if self._body_excerpt(limit=120):
return True
if consecutive_chatgpt_ready >= 3:
logger.info("[ChatGPT] workspace 跳转后已到达 ChatGPT 页面,继续后续登录检测")
return True
else:
consecutive_chatgpt_ready = 0
last_url = self.page.url
time.sleep(0.5)
logger.warning("[ChatGPT] workspace 跳转后页面仍未稳定 | URL=%s", last_url)
return False
def select_workspace_option(self, option_id):
options = self._list_workspace_options()
for option in options:
if option["id"] != str(option_id):
continue
label = option["label"]
logger.info("[ChatGPT] 用户选择 workspace: %s", label)
clicked = self._click_workspace_option_by_label(label)
if not clicked:
raise RuntimeError(f"未找到可点击的 workspace 选项: {label}")
exited = self._wait_for_workspace_selection_exit(timeout=15)
if not exited:
logger.warning("[ChatGPT] 手动选择 workspace 后仍未离开选择页,尝试自动点击一次: %s", label)
if self._auto_open_preferred_workspace():
self._wait_for_workspace_selection_exit(timeout=10)
try:
self.page.wait_for_load_state("domcontentloaded", timeout=5000)
except Exception:
pass
self._wait_for_post_workspace_ready(timeout=12)
self.workspace_options_cache = []
self._log_login_state("选择 workspace 后")
if self._is_chatgpt_ready_url():
logger.info("[ChatGPT] 选择 workspace 后结果: completed(chatgpt) | detail=None")
return {"step": "completed", "detail": None}
step, detail = self._detect_login_step()
logger.info("[ChatGPT] 选择 workspace 后结果: %s | detail=%s", step, detail)
return {"step": step, "detail": detail}
raise RuntimeError(f"无效的 workspace 选项: {option_id}")
def _detect_login_step(self):
if "accounts.google.com" in self.page.url:
logger.warning("[ChatGPT] 登录步骤检测: 误跳转 Google | URL=%s", self.page.url)
return "error", "误跳转到了 Google 登录"
if self._is_workspace_selection_page():
logger.info("[ChatGPT] 登录步骤检测: workspace 页面 | URL=%s", self.page.url)
return "workspace_required", None
if "email-verification" in self.page.url:
logger.info("[ChatGPT] 登录步骤检测: code_required | URL=%s", self.page.url)
return "code_required", None
if self._visible_locator_in_frames(self.CODE_INPUT_SELECTORS, timeout_ms=1200):
logger.info("[ChatGPT] 登录步骤检测: code_required | URL=%s", self.page.url)
return "code_required", None
if self._visible_locator_in_frames(self.PASSWORD_INPUT_SELECTORS, timeout_ms=1200):
logger.info("[ChatGPT] 登录步骤检测: password_required | URL=%s", self.page.url)
return "password_required", None
if self._visible_locator_in_frames(self.EMAIL_INPUT_SELECTORS, timeout_ms=1200):
logger.info("[ChatGPT] 登录步骤检测: email_required | URL=%s", self.page.url)
return "email_required", None
url = (self.page.url or "").lower()
if "log-in-or-create-account" in url or url.endswith("/auth/login"):
logger.info("[ChatGPT] 登录步骤检测: email_required(url) | URL=%s", self.page.url)
return "email_required", None
session_token = self._extract_session_token()
if session_token:
logger.info("[ChatGPT] 登录步骤检测: completed(session) | URL=%s", self.page.url)
return "completed", None
if "chatgpt.com" in self.page.url and "auth" not in self.page.url:
logger.info("[ChatGPT] 登录步骤检测: completed(chatgpt) | URL=%s", self.page.url)
return "completed", None
logger.info("[ChatGPT] 登录步骤检测: unknown | URL=%s", self.page.url)
return "unknown", self.page.url
def begin_login(self, email, actor_label="账号"):
self.login_email = email
if not self.browser:
self._launch_browser()
logger.info("[ChatGPT] 开始%s登录: %s", actor_label, email)
self.page.goto("https://chatgpt.com/", wait_until="domcontentloaded", timeout=60000)
time.sleep(5)
self._wait_for_cloudflare()
self._log_login_state("进入 chatgpt.com 后")
self._open_login_page()
step, detail = self._wait_for_login_step(
{"email_required", "password_required", "code_required", "workspace_required", "completed", "error"},
timeout=12,
)
if step == "workspace_required":
self._list_workspace_options()
if step in ("password_required", "code_required", "workspace_required", "completed", "error"):
logger.info("[ChatGPT] %s登录初始步骤: %s | detail=%s", actor_label, step, detail)
return {"step": step, "detail": detail}
email_input = self._visible_locator_in_frames(self.EMAIL_INPUT_SELECTORS, timeout_ms=15000)
if not email_input:
try:
self.page.screenshot(path=str(SCREENSHOT_DIR / "admin_login_missing_email.png"), full_page=True)
except Exception:
pass
body_excerpt = ""
try:
body_excerpt = self.page.locator("body").inner_text(timeout=2000)[:300]
except Exception:
pass
raise RuntimeError(f"未找到{actor_label}邮箱输入框,当前 URL: {self.page.url},页面片段: {body_excerpt}")
final_step, final_detail = "unknown", self.page.url
for attempt in range(1, 4):
email_input = self._visible_locator_in_frames(self.EMAIL_INPUT_SELECTORS, timeout_ms=3000) or email_input
try:
email_input.fill(email)
except Exception:
try:
email_input.click(timeout=1000)
email_input.fill(email)
except Exception:
pass
time.sleep(0.5)
clicked = self._click_auth_button(email_input, ["Continue", "继续", "Log in"])
logger.info("[ChatGPT] %s邮箱已提交(第 %d 次)| clicked=%s", actor_label, attempt, clicked)
final_step, final_detail = self._wait_for_login_step(
{"email_required", "password_required", "code_required", "workspace_required", "completed", "error"},
timeout=12,
)
self._log_login_state(f"{actor_label}邮箱提交后(第 {attempt} 次)")
if final_step == "workspace_required":
self._list_workspace_options()
if final_step != "email_required":
break
logger.warning(
"[ChatGPT] %s邮箱提交后仍停留在邮箱步骤(第 %d 次)| URL=%s | body=%s",
actor_label,
attempt,
self.page.url,
self._body_excerpt(),
)
if final_step == "email_required":
raise RuntimeError(
f"{actor_label}邮箱提交后仍停留在邮箱步骤,请检查登录页是否拦截/未响应。"
f" 当前 URL: {self.page.url},页面片段: {self._body_excerpt()}"
)
logger.info("[ChatGPT] %s邮箱提交结果: %s | detail=%s", actor_label, final_step, final_detail)
return {"step": final_step, "detail": final_detail}
def begin_admin_login(self, email):
return self.begin_login(email, actor_label="管理员")
def submit_login_password(self, password, actor_label="账号"):
self.login_password = password
password_input = self._visible_locator_in_frames(self.PASSWORD_INPUT_SELECTORS, timeout_ms=5000)
if not password_input:
raise RuntimeError("当前不是密码输入步骤")
logger.info("[ChatGPT] 提交%s密码前 | URL=%s", actor_label, self.page.url)
password_input.fill(password)
time.sleep(0.5)
self._click_auth_button(password_input, ["Continue", "继续", "Log in"])
time.sleep(8)
self._log_login_state(f"{actor_label}密码提交后")
step, detail = self._detect_login_step()
if step == "workspace_required":
self._list_workspace_options()
logger.info("[ChatGPT] %s密码提交结果: %s | detail=%s", actor_label, step, detail)
return {"step": step, "detail": detail}
def submit_admin_password(self, password):
return self.submit_login_password(password, actor_label="管理员")
def submit_login_code(self, code, actor_label="账号"):
code_input = self._visible_locator_in_frames(self.CODE_INPUT_SELECTORS, timeout_ms=5000)
if not code_input:
time.sleep(3)
code_input = self._visible_locator_in_frames(self.CODE_INPUT_SELECTORS, timeout_ms=5000)
# email-verification 页面可能用单字符输入框或其他结构
if not code_input:
try:
# 尝试单字符输入框(多个 input[maxlength="1"])
single_inputs = self.page.locator('input[maxlength="1"]').all()
if len(single_inputs) >= 4:
logger.info("[ChatGPT] 检测到 %d 个单字符验证码输入框", len(single_inputs))
for i, char in enumerate(code):
if i < len(single_inputs):
single_inputs[i].fill(char)
time.sleep(0.1)
time.sleep(0.5)
# 可能自动提交,也可能需要点按钮
try:
btn = self.page.locator(
'button:has-text("Continue"), button:has-text("继续"), button:has-text("Verify"), button[type="submit"]'
).first
if btn.is_visible(timeout=2000):
btn.click()
except Exception:
pass
time.sleep(8)
self._log_login_state(f"{actor_label}验证码提交后(单字符)")
step, detail = self._detect_login_step()
if step == "workspace_required":
self._list_workspace_options()
return {"step": step, "detail": detail}
except Exception as e:
logger.warning("[ChatGPT] 单字符输入框尝试失败: %s", e)
# 尝试用 JS 直接找任何可见的 input
if not code_input:
try:
code_input = self.page.locator("input:visible").first
if not code_input.is_visible(timeout=2000):
code_input = None
except Exception:
code_input = None
if not code_input:
try:
self.page.screenshot(path=str(SCREENSHOT_DIR / "admin_login_code_not_found.png"), full_page=True)
except Exception:
pass
logger.error("[ChatGPT] 找不到%s验证码输入框 | URL=%s", actor_label, self.page.url)
raise RuntimeError("找不到验证码输入框,页面可能已跳转或验证码已过期")
logger.info("[ChatGPT] 提交%s验证码前 | URL=%s | code_len=%d", actor_label, self.page.url, len(code))
try:
self.page.screenshot(path=str(SCREENSHOT_DIR / "admin_login_code_before_submit.png"), full_page=True)
except Exception:
pass
code_input.fill(code)
time.sleep(0.5)
self._click_auth_button(code_input, ["Continue", "继续", "Verify"])
time.sleep(8)
try:
self.page.screenshot(path=str(SCREENSHOT_DIR / "admin_login_code_after_submit.png"), full_page=True)
except Exception:
pass
self._log_login_state(f"{actor_label}验证码提交后")
step, detail = self._detect_login_step()
if step == "workspace_required":
self._list_workspace_options()
logger.info("[ChatGPT] %s验证码提交结果: %s | detail=%s", actor_label, step, detail)
return {"step": step, "detail": detail}
def submit_admin_code(self, code):
return self.submit_login_code(code, actor_label="管理员")
def _list_real_workspaces(self):
"""
拉取 /backend-api/accounts,按 structure 切分为 team_accounts / personal_accounts。
这是 session 真正所属 workspace 的唯一可信来源,其他字段(_account cookie / JWT 的
chatgpt_account_id claim)都可能陈旧或被污染。
返回 (team_accounts, personal_accounts),每个 item 至少包含 id / name / structure /
current_user_role 字段。
"""
result = self._api_fetch("GET", "/backend-api/accounts")
if result.get("status") != 200:
raise RuntimeError(
f"无法获取 workspace 列表: status={result.get('status')}, body={(result.get('body') or '')[:200]}"
)
try:
data = json.loads(result["body"])
except Exception as exc:
raise RuntimeError(f"/backend-api/accounts 响应解析失败: {exc}")
items = data.get("items") or data.get("data") or data.get("accounts") or []
if not isinstance(items, list):
items = []
team = []
personal = []
for item in items:
if not isinstance(item, dict):
continue
structure = str(item.get("structure") or "").lower()
if structure == "workspace":
team.append(item)
else:
personal.append(item)
return team, personal
def _guess_account_info(self, allow_dom_fallback=True):
try:
data = self.page.evaluate(
"""async (accessToken) => {
const out = {};
const headers = accessToken ? { authorization: `Bearer ${accessToken}` } : {};
for (const path of ['/backend-api/accounts', '/backend-api/me', '/api/auth/session']) {
try {
const resp = await fetch(path, { headers });
out[path] = { status: resp.status, data: await resp.json() };
} catch (e) {
out[path] = { error: String(e) };
}
}
return out;
}""",
self.access_token,
)
except Exception:
data = {}
candidates = []
def walk(node):
if isinstance(node, dict):
# account_id 必须是 UUID 格式(排除 user-xxx 等非 account ID)
account_id = node.get("account_id")
if not account_id or not self._UUID_RE.match(str(account_id)):
account_id = node.get("id")
if not account_id or not self._UUID_RE.match(str(account_id)):
account_id = None
# workspace_name 只取 workspace_name 字段,不取 name/display_name(那可能是用户名)
name = node.get("workspace_name") or ""
if account_id:
candidates.append(
{
"account_id": account_id,
"workspace_name": name,
"type": str(node.get("type", "")),
}
)
for value in node.values():
walk(value)
elif isinstance(node, list):
for item in node:
walk(item)
walk(data)
chosen = None
for cand in candidates:
if cand["workspace_name"] and cand["workspace_name"].lower() not in ("personal",):
chosen = cand
break
if not chosen and candidates:
chosen = candidates[0]
dom_name = None
if allow_dom_fallback:
try:
self.page.goto("https://chatgpt.com/admin", wait_until="domcontentloaded", timeout=30000)
time.sleep(5)
dom_name = self._detect_workspace_name_from_dom() or None
except Exception:
dom_name = None
account_id = (chosen or {}).get("account_id") or self.account_id
workspace_name = (chosen or {}).get("workspace_name") or dom_name or self.workspace_name
return account_id, workspace_name
def complete_login(self, persist_admin_state=False):
session_token = self._extract_session_token()
if not session_token:
raise RuntimeError("登录成功后未提取到 session token")
self._fetch_access_token()
account_id, workspace_name = self._guess_account_info()
if account_id:
self.account_id = account_id
if workspace_name:
self.workspace_name = workspace_name
payload = dict(
email=self.login_email or "",
session_token=session_token,
account_id=self.account_id,
workspace_name=self.workspace_name,
)
if self.login_password:
payload["password"] = self.login_password
if persist_admin_state:
update_admin_state(**payload)
logger.info("[ChatGPT] 管理员登录状态已保存")
return {
"email": self.login_email or "",
"password": self.login_password or "",
"session_token": session_token,
"account_id": self.account_id,
"workspace_name": self.workspace_name,
"session_len": len(session_token),
}
def complete_admin_login(self):
return self.complete_login(persist_admin_state=True)
def import_admin_session(self, email, session_token):
"""手动导入管理员 session_token,并自动识别 workspace 信息。"""
email = (email or "").strip()
session_token = (session_token or "").strip()
if not email:
raise RuntimeError("管理员邮箱不能为空")
if not session_token:
raise RuntimeError("session_token 不能为空")
self.login_email = email
self.session_token = session_token
self._launch_browser()
logger.info("[ChatGPT] 开始导入管理员 session_token: %s", email)
self.page.goto("https://chatgpt.com/", wait_until="domcontentloaded", timeout=60000)
time.sleep(5)
self._wait_for_cloudflare()
self._inject_session(session_token)
self.page.goto("https://chatgpt.com/", wait_until="domcontentloaded", timeout=60000)
time.sleep(5)
self._wait_for_cloudflare()
self._log_login_state("导入 session_token 后")
if self._is_workspace_selection_page():
logger.info("[ChatGPT] 检测到 workspace 选择页,尝试自动进入 Team workspace")
self._auto_open_preferred_workspace()
token_source = self._fetch_access_token(allow_bearer_file=False)
if not token_source:
raise RuntimeError("session_token 无效或已过期,未能从当前登录态获取 access token")
logger.info("[ChatGPT] session_token 导入后 access token 来源: %s", token_source)
# 权威来源:/backend-api/accounts 列出当前 session 真实所属的所有 workspace。
# 任何其他来源(_account cookie / JWT claim)都可能是上次 logout 前的残留或 OAI 缓存的
# 陈旧值,直接写入会导致所有 admin 接口 401 "Must be part of this workspace"。
team_accounts, personal_accounts = self._list_real_workspaces()
# 优先选 Team workspace (structure=="workspace",且 current_user_role 有 admin 权限)。
# 用户可能自己就是 account-owner / admin,哪一种 role 都接受,只要不是单纯被邀请的 user。
admin_roles = ("account-owner", "admin", "org-admin", "workspace-owner")
chosen = None
chosen_reason = ""
for acc in team_accounts:
role = str(acc.get("current_user_role") or "").lower()
if role in admin_roles:
chosen = acc
chosen_reason = f"role={role}"
break
if not chosen and team_accounts:
chosen = team_accounts[0]
chosen_reason = f"role={chosen.get('current_user_role')} (非标准 admin,但接受)"
if not chosen:
raise RuntimeError(
f"当前 session ({email}) 没有可用的 Team workspace:"
f" /backend-api/accounts 只返回 {[a.get('structure') for a in personal_accounts]} "
f"结构的账号。请确认该 session_token 对应的账号已被邀请进 Team 并接受邀请。"
)
account_id = str(chosen.get("id") or "")
workspace_name = str(chosen.get("name") or "")
if not account_id or not self._UUID_RE.match(account_id):
raise RuntimeError(f"Team workspace 返回的 account id 格式异常: {account_id!r}")
# 二次确认:用这个 account_id 调 /settings,若仍 401 说明 session 与 workspace 不匹配,
# 宁可立刻失败,也不把错误 state 写进磁盘。
verify = self._api_fetch("GET", f"/backend-api/accounts/{account_id}/settings")
if verify.get("status") != 200:
body = (verify.get("body") or "")[:200]
raise RuntimeError(
f"account_id={account_id} 鉴权验证失败 status={verify.get('status')},"
f" body={body}。session_token 可能与 workspace 不匹配。"
)
logger.info(
"[ChatGPT] 已确认 Team workspace: name=%s account_id=%s (%s)",
workspace_name or "?",
account_id,
chosen_reason,
)
self.account_id = account_id
self.workspace_name = workspace_name
update_admin_state(
email=email,
session_token=session_token,
account_id=self.account_id,
workspace_name=self.workspace_name,
)
try:
from autoteam.workspace_pool import default_pool
default_pool.upsert(
f"ws-{self.account_id}",
email,
self.account_id,
workspace_name=self.workspace_name,
session_token=session_token,
enabled=True,
parallel=True,
)
except Exception as exc:
logger.warning("[ChatGPT] workspace pool 记录母号失败: %s", exc)
logger.info("[ChatGPT] 管理员 session_token 已保存")
return {
"email": email,
"password": "",
"session_token": session_token,
"account_id": self.account_id,
"workspace_name": self.workspace_name,
"session_len": len(session_token),
}
def start(self):
"""用已保存的管理员 session 启动 Team API 客户端。"""
session_token = get_admin_session_token()
self.account_id = get_chatgpt_account_id()
self.workspace_name = get_chatgpt_workspace_name()
self.start_with_session(session_token, self.account_id, self.workspace_name)
def is_started(self):
return bool(self.browser or self.http_transport)
def _build_api_headers(self, include_auth=True):
raw_headers = {
"Content-Type": "application/json",
"chatgpt-account-id": self.account_id,
"oai-device-id": self.oai_device_id,
"oai-language": "en-US",
}
if include_auth and self.access_token:
raw_headers["authorization"] = f"Bearer {self.access_token}"
# Browser fetch headers must be ISO-8859-1 (Latin1). Sanitize here so
# Playwright and optional HTTP transports share the Round 11 header fix.
headers = {}
for key, value in raw_headers.items():
if value is None:
continue
if isinstance(value, bytes):
value = value.decode("ascii", errors="replace")
else:
value = str(value)
sanitized = value.encode("latin-1", errors="replace").decode("latin-1")
if sanitized != value:
logger.warning(
"[ChatGPT] _api_fetch header %s 含非 ISO-8859-1 字符已被替换: %r -> %r",
key,
value[:60],
sanitized[:60],
)
headers[key] = sanitized
return headers
def _start_transport_session(self, session_token):
self.http_transport = build_chatgpt_transport(
session_token=session_token,
account_id=self.account_id,
oai_device_id=self.oai_device_id,
proxy_url=self.local_proxy_url,
)
if self.http_transport:
self.transport_name = getattr(self.http_transport, "name", "curl_cffi")
logger.info("[ChatGPT] Team API transport: %s", self.transport_name)
return True
return False
def _close_http_transport(self, *, reason: str = "") -> None:
transport = getattr(self, "http_transport", None)
try:
if transport is not None:
transport.close()
except Exception as exc:
if reason:
logger.debug("[ChatGPT] close curl_cffi transport failed (%s): %s", reason, exc)
else:
logger.debug("[ChatGPT] close curl_cffi transport failed: %s", exc)
finally:
self.http_transport = None
self.transport_name = None
@staticmethod
def _transport_body_looks_like_html(body):
text = str(body or "").lower()
return "<html" in text or "<!doctype" in text or "<body" in text
def _transport_response_requires_browser_fallback(self, response):
if not isinstance(response, dict):
return True
status = int(response.get("status") or 0)
body = response.get("body", "")
if status <= 0 or status >= 500:
return True
if not self._transport_body_looks_like_html(body):
return False
lower = str(body or "").lower()
return any(
marker in lower
for marker in (
"verify you are human",
"cloudflare",
"challenge",
"auth/login",
"log in",
"unable to load site",
)
)
def _start_browser_session(self, session_token):
self._launch_browser()
logger.info("[ChatGPT] 访问 chatgpt.com 过 Cloudflare...")
self.page.goto("https://chatgpt.com/", wait_until="domcontentloaded", timeout=60000)
time.sleep(5)
self._wait_for_cloudflare()
self._inject_session(session_token)
self._fetch_access_token()
self._auto_detect_workspace()
def _ensure_browser_session(self):
if self.browser:
return
if not self.session_token:
self.stop()
raise RuntimeError("缺少 session_token,无法回退到 Playwright transport")
try:
self._start_browser_session(self.session_token)
except Exception:
self.stop()
raise
def _fetch_access_token_via_transport(self, allow_bearer_file=True):
if not getattr(self, "http_transport", None):
return None
try:
result = self.http_transport.request(
"GET",
"/api/auth/session",
headers={
"Accept": "application/json, text/plain, */*",
"oai-device-id": self.oai_device_id,
},
)
except Exception as exc:
logger.warning("[ChatGPT] curl_cffi 获取 access token 失败: %s", exc)
result = None
if result and int(result.get("status") or 0) == 200:
try:
data = json.loads(result.get("body") or "{}")
except Exception:
data = {}
if "accessToken" in data:
self.access_token = data["accessToken"]
logger.info("[ChatGPT] 已通过 curl_cffi 获取 access token")
return "curl_cffi"
if allow_bearer_file:
bearer_file = BASE_DIR / "bearer_token"
if bearer_file.exists():
self.access_token = read_text(bearer_file).strip()
logger.info("[ChatGPT] 从 bearer_token 文件加载 access token")
return "file"
return None
def _auto_detect_workspace_via_transport(self):
if self.workspace_name or not getattr(self, "http_transport", None) or not self.account_id:
return self.workspace_name
try:
result = self.http_transport.request(
"GET",
f"/backend-api/accounts/{self.account_id}/settings",
headers=self._build_api_headers(),
)
except Exception as exc:
logger.warning("[ChatGPT] curl_cffi 自动检测 workspace 失败: %s", exc)
return ""
if int(result.get("status") or 0) != 200:
return ""
try:
data = json.loads(result.get("body") or "{}")
except Exception:
return ""
workspace_name = (data or {}).get("workspace_name", "") if isinstance(data, dict) else ""
if workspace_name:
self.workspace_name = workspace_name
update_admin_state(workspace_name=self.workspace_name, account_id=self.account_id)
logger.info("[ChatGPT] 自动检测到 workspace 名称: %s", self.workspace_name)
return self.workspace_name
def start_with_session(self, session_token, account_id, workspace_name="", require_browser=False):
"""用指定的 session/account 启动浏览器上下文。"""
if not session_token:
raise FileNotFoundError("缺少会话信息")
self.stop()
self.account_id = account_id or ""
self.workspace_name = workspace_name or ""
self.session_token = session_token
self.access_token = None
if not self.account_id:
raise RuntimeError("缺少 workspace/account ID")
try:
self._ensure_admin_ipv6_proxy()
if not require_browser and self._start_transport_session(session_token):
token_source = self._fetch_access_token_via_transport()
if token_source:
self._auto_detect_workspace_via_transport()
return
logger.warning("[ChatGPT] curl_cffi 未能直接获取 access token,回退 Playwright transport")
self._close_http_transport(reason="before browser fallback")
except Exception:
self.stop()
raise
try:
self._start_browser_session(session_token)
except Exception:
self.stop()
raise
def _auto_detect_workspace(self):
if self.workspace_name:
return self.workspace_name
if not self.account_id:
logger.warning("[ChatGPT] 未能自动获取 workspace 名称:account_id 缺失")
return ""
result = self.page.evaluate(
"""async ([accountId, accessToken]) => {
try {
const headers = { "chatgpt-account-id": accountId };
if (accessToken) headers["authorization"] = `Bearer ${accessToken}`;
const resp = await fetch("/backend-api/accounts/" + accountId + "/settings", {
headers
});
return await resp.json();
} catch(e) { return null; }
}""",
[self.account_id, self.access_token],
)
if result and result.get("workspace_name"):
self.workspace_name = result["workspace_name"]
update_admin_state(workspace_name=self.workspace_name, account_id=self.account_id)
logger.info("[ChatGPT] 自动检测到 workspace 名称: %s", self.workspace_name)
return self.workspace_name
try:
self.page.goto("https://chatgpt.com/admin", wait_until="domcontentloaded", timeout=30000)
time.sleep(5)
name = self._detect_workspace_name_from_dom()
if name:
self.workspace_name = name
update_admin_state(workspace_name=self.workspace_name, account_id=self.account_id)
logger.info("[ChatGPT] 自动检测到 workspace 名称: %s", name)
return name
except Exception:
pass
logger.warning("[ChatGPT] 未能自动获取 workspace 名称")
return ""
def _fetch_access_token(self, allow_bearer_file=True):
result = self.page.evaluate("""async () => {
try {
const resp = await fetch("/api/auth/session");
const data = await resp.json();
return { ok: true, data: data };
} catch(e) {
return { ok: false, error: e.message };
}
}""")
if result.get("ok") and "accessToken" in result.get("data", {}):
self.access_token = result["data"]["accessToken"]
logger.info("[ChatGPT] 已获取 access token")
return "session"
if allow_bearer_file:
bearer_file = BASE_DIR / "bearer_token"
if bearer_file.exists():
self.access_token = read_text(bearer_file).strip()
logger.info("[ChatGPT] 从 bearer_token 文件加载 access token")
return "file"
logger.info("[ChatGPT] 尝试通过页面获取 access token...")
self.page.goto("https://chatgpt.com/", wait_until="networkidle", timeout=60000)
time.sleep(10)
token = self.page.evaluate("""() => {
try {
const keys = Object.keys(localStorage);
for (const key of keys) {
const val = localStorage.getItem(key);
if (val && val.includes("eyJ") && val.length > 500) {
return val;
}
}
} catch(e) {}
return null;
}""")
if token:
self.access_token = token
logger.info("[ChatGPT] 从页面获取到 access token")
return "localstorage"
else:
logger.warning("[ChatGPT] 未能获取 access token,将尝试无 token 调用")
return None
def _browser_api_fetch(self, method, path, body=None):
if not self.page:
raise RuntimeError("Playwright 页面未初始化")
js_code = """async ([method, url, headers, body]) => {
try {
const opts = { method, headers };
if (body) opts.body = body;
const resp = await fetch(url, opts);
const text = await resp.text();
return { status: resp.status, body: text };
} catch(e) {
return { status: 0, body: e.message };
}
}"""
return self.page.evaluate(
js_code,
[method, f"https://chatgpt.com{path}", self._build_api_headers(), json.dumps(body) if body else None],
)
def _direct_api_fetch(self, method, path, body=None):
if not getattr(self, "http_transport", None):
raise RuntimeError("curl_cffi transport 未初始化")
try:
response = self.http_transport.request(
method,
path,
headers=self._build_api_headers(),
body=body,
)
except Exception as exc:
logger.warning("[ChatGPT] curl_cffi request failed,回退 Playwright transport: %s", exc)
self._close_http_transport(reason="request failed")
self._ensure_browser_session()
return self._browser_api_fetch(method, path, body)
status = int(response.get("status") or 0)
body_text = str(response.get("body") or "")
lower_body = body_text.lower()
if status == 401 and path != "/api/auth/session":
refreshed = self._fetch_access_token_via_transport(allow_bearer_file=False)
if refreshed:
try:
response = self.http_transport.request(
method,
path,
headers=self._build_api_headers(),
body=body,
)
except Exception as exc:
logger.warning("[ChatGPT] curl_cffi retry failed,回退 Playwright transport: %s", exc)
self._close_http_transport(reason="retry failed")
self._ensure_browser_session()
return self._browser_api_fetch(method, path, body)
status = int(response.get("status") or 0)
body_text = str(response.get("body") or "")
lower_body = body_text.lower()
if self._transport_response_requires_browser_fallback(response):
logger.warning("[ChatGPT] curl_cffi 返回异常响应,回退 Playwright transport")
self._close_http_transport(reason="fallback response")
self._ensure_browser_session()
return self._browser_api_fetch(method, path, body)
if status == 401 and (
"access token is missing" in lower_body or self._transport_body_looks_like_html(body_text)
):
logger.warning("[ChatGPT] curl_cffi 鉴权异常,回退 Playwright transport")
self._close_http_transport(reason="auth fallback")
self._ensure_browser_session()
return self._browser_api_fetch(method, path, body)
return response
def _api_fetch(self, method, path, body=None):
if getattr(self, "http_transport", None):
return self._direct_api_fetch(method, path, body)
if not self.page and self.session_token:
self._ensure_browser_session()
return self._browser_api_fetch(method, path, body)
# POST /invites 失败时的重试间隔(秒)。rate_limited / network 类错误按此序列退避。
_INVITE_POST_RETRY_DELAYS = (5, 15)
# PATCH /invites/{id} 失败时的重试间隔(秒)。次数 = len(_INVITE_PATCH_RETRY_DELAYS)
_INVITE_PATCH_RETRY_DELAYS = (5,)
# 明确表示"目标域名被 workspace 拒绝"的短语。注意这里**不能**单独放 "domain" 这个 token —
# 服务端有大量包含 "domain" 字面量的无关错误(rate-limit 提示里的 "your domain ...",
# 甚至 errored_emails 里 email 自身的 "@gmail.com" 都会让 "domain" 命中),会把可重试错误
# 误判成 domain_blocked 直接返回给上层,导致整批账号被错误地放弃。
_DOMAIN_BLOCKED_KEYWORDS = (
"not allowed",
"domain blocked",
"domain is not allowed",
"forbidden domain",
"domain not permitted",
)
@staticmethod
def _classify_invite_error(status, data, resp_body):
"""将 POST /invites 的响应归类为 domain_blocked / rate_limited / server_error / network / other。
- status == 0 视为 network(fetch 抛异常,可重试)
- 429 或 body 含 rate_limit/too many 视为 rate_limited(可重试)
- 5xx(500/502/503/504) 视为 server_error(可重试)
- 4xx 且 **明确字段**(detail/error/message + errored_emails[].error/code)命中
domain_blocked 关键词时归为 domain_blocked
- 其他非 200 为 other(不重试,交上层换号)
"""
if status == 0:
return "network"
if status == 429:
return "rate_limited"
if status in (500, 502, 503, 504):
return "server_error"
# 只在明确字段拼接 body_text;不再 fallthrough 到 resp_body —
# 否则 email 中的 "gmail.com" 之类会被旧逻辑的 "domain" 关键词命中。
body_text = ""
if isinstance(data, dict):
for key in ("detail", "error", "message"):
val = data.get(key)
if isinstance(val, str):
body_text += " " + val
elif isinstance(val, dict):
inner = val.get("message") or val.get("code") or ""
if isinstance(inner, str):
body_text += " " + inner
# MEDIUM-1: errored_emails 内层 error/code 字段也算明确字段
for item in data.get("errored_emails", []) or []:
if not isinstance(item, dict):
continue
for inner_key in ("error", "code", "message"):
val = item.get(inner_key)
if isinstance(val, str):
body_text += " " + val
lowered = body_text.lower()
if any(kw in lowered for kw in ("rate_limit", "rate limit", "too many", "throttle")):
return "rate_limited"
if any(kw in lowered for kw in ChatGPTTeamAPI._DOMAIN_BLOCKED_KEYWORDS):
return "domain_blocked"
return "other"
def invite_member(self, email, seat_type="usage_based", *, allow_patch_upgrade=True):
"""邀请成员加入 Team(自带 default → usage_based 兜底 + errored_emails 处理)。
返回 `(status, data)`,其中 `data` 一定是 dict(解析失败封装为 `{"_raw": <text>}`),
必含字段:
- `_seat_type` ∈ {"chatgpt","usage_based","unknown"}
* "chatgpt" POST 200 + (PATCH default 成功 或 直接 default 邀请成功)
* "usage_based" POST 200 但 PATCH 全部失败 / allow_patch_upgrade=False,仅 codex 席位
* "unknown" POST 本身失败/被业务拒绝(domain/errored)
- `_error_kind` 最后一次失败的分类(成功时不写)
- `_errored_emails` POST 200 但 errored_emails 非空时,原样保留 errored_emails 数组,
供上游记录失败原因。
参数:
- allow_patch_upgrade: True(默认) 走 usage_based POST + PATCH default 升级;
False 表示锁定 codex-only 席位,跳过 PATCH(对应 PREFERRED_SEAT_TYPE="codex")。
兜底逻辑(所有 fallback 都在本函数内完成,调用方拿到结果就是终态):
- seat_type="default" 时,若 POST 200 但 errored_emails 命中或 PATCH 失败,
自动重试一次 seat_type="usage_based"(整套 retry 计数重新算)。
- 顶层 HTTP 失败 + err_kind ∈ {network, rate_limited, server_error}:
按 _INVITE_POST_RETRY_DELAYS 退避(带 jitter)重试。
"""
# default 邀请失败时下降到 usage_based(只翻一次,避免无限递归)
return self._invite_member_with_fallback(
email, seat_type, allow_fallback=True, allow_patch_upgrade=allow_patch_upgrade
)
def _invite_member_with_fallback(self, email, seat_type, *, allow_fallback, allow_patch_upgrade=True):
status, data = self._invite_member_once(email, seat_type, allow_patch_upgrade=allow_patch_upgrade)
# default 路径:有任何业务级失败迹象都尝试 usage_based 兜底
# 1) HTTP 非 200 且不是 domain_blocked(domain_blocked 直接返回让上层换号)
# 2) HTTP 200 但 errored_emails 非空(账号被业务规则拒绝)
# 3) HTTP 200 但 _seat_type=="unknown"(理论上 _invite_member_once 不会返回这种,但兜底)
if seat_type == "default" and allow_fallback:
errored = data.get("errored_emails") if isinstance(data, dict) else None
err_kind = data.get("_error_kind") if isinstance(data, dict) else None
should_fallback = False
if status != 200 and err_kind not in ("domain_blocked",):
should_fallback = True
elif status == 200 and errored:
should_fallback = True
if should_fallback:
logger.info(
"[ChatGPT] %s default 邀请失败(status=%d kind=%s errored=%s),尝试 usage_based 兜底",
email,
status,
err_kind,
bool(errored),
)
return self._invite_member_with_fallback(
email, "usage_based", allow_fallback=False, allow_patch_upgrade=allow_patch_upgrade
)
return status, data
def _invite_member_once(self, email, seat_type, *, allow_patch_upgrade=True):
"""单一 seat_type 的完整邀请尝试(POST retry → PATCH 升级)。"""
path = f"/backend-api/accounts/{self.account_id}/invites"
body = {
"email_addresses": [email],
"role": "standard-user",
"seat_type": seat_type,
"resend_emails": True,
}
status = 0
data = {}
resp_body = ""
attempts = 1 + len(self._INVITE_POST_RETRY_DELAYS)
for attempt in range(attempts):
logger.info(
"[ChatGPT] 发送邀请到 %s (seat_type=%s, attempt=%d/%d)...",
email,
seat_type,
attempt + 1,
attempts,
)
result = self._api_fetch("POST", path, body)
status = result["status"]
resp_body = result["body"]
logger.info("[ChatGPT] 响应状态: %d", status)
try:
parsed = json.loads(resp_body)
logger.debug("[ChatGPT] 响应内容: %s", json.dumps(parsed, indent=2)[:500])
except Exception:
parsed = {"_raw": resp_body}
logger.debug("[ChatGPT] 响应内容(非 JSON): %s", (resp_body or "")[:500])
data = parsed if isinstance(parsed, dict) else {"_raw": parsed}
if status == 200:
break
err_kind = self._classify_invite_error(status, data, resp_body)
logger.warning(
"[ChatGPT] 邀请 %s 失败: status=%d kind=%s body=%s",
email,
status,
err_kind,
(resp_body or "")[:200],
)
# domain_blocked / other: 不 retry,直接返回让上层换号
if err_kind in ("domain_blocked", "other"):
data.setdefault("_seat_type", "unknown")
data.setdefault("_error_kind", err_kind)
return status, data
# rate_limited / network / server_error: 按退避表(带 jitter)retry
if attempt < len(self._INVITE_POST_RETRY_DELAYS):
base_delay = self._INVITE_POST_RETRY_DELAYS[attempt]
# MEDIUM-2: 30% jitter 避免多客户端被同一 rate-limit 窗口拒绝后同时唤醒
delay = base_delay + random.uniform(0, base_delay * 0.3)
logger.info("[ChatGPT] %s 类错误,%.1fs 后重试邀请 %s", err_kind, delay, email)
time.sleep(delay)
continue
data.setdefault("_seat_type", "unknown")
data.setdefault("_error_kind", err_kind)
return status, data
# status == 200 → 检查 errored_emails / 处理 PATCH 升级
errored = data.get("errored_emails", []) if isinstance(data, dict) else []
if errored:
err_msg = errored[0].get("error", "unknown") if isinstance(errored[0], dict) else "unknown"
logger.warning("[ChatGPT] 邀请 %s 被 errored_emails 拒绝: %s", email, err_msg)
data["_seat_type"] = "unknown"
data["_error_kind"] = "errored_emails"
data["_errored_emails"] = errored
return status, data
# 默认标 usage_based,PATCH 成功再升级为 chatgpt
data["_seat_type"] = "usage_based"
if seat_type == "usage_based" and allow_patch_upgrade:
invites = data.get("account_invites", []) if isinstance(data, dict) else []
any_patched = False
any_invite = False
for inv in invites:
invite_id = inv.get("id") if isinstance(inv, dict) else None
if not invite_id:
continue
any_invite = True
if self._update_invite_seat_type(invite_id, "default"):
any_patched = True
if any_invite and any_patched:
data["_seat_type"] = "chatgpt"
elif any_invite and not any_patched:
logger.error(
"[ChatGPT] %s PATCH seat_type 全部失败,保留 codex 席位(_seat_type=usage_based)",
email,
)
elif seat_type == "usage_based" and not allow_patch_upgrade:
# PREFERRED_SEAT_TYPE="codex" 路径:跳过 PATCH 升级,锁定 codex-only 席位
logger.info("[ChatGPT] %s 已锁 codex-only 席位(allow_patch_upgrade=False)", email)
elif seat_type == "default":
# 直接 default 邀请,只要 POST 200 + 无 errored 即完整 ChatGPT 席位
data["_seat_type"] = "chatgpt"
return status, data
def _update_invite_seat_type(self, invite_id, seat_type):
"""PATCH 修改 invite 的 seat_type。返回 True 表示成功,False 表示重试后仍失败。"""
path = f"/backend-api/accounts/{self.account_id}/invites/{invite_id}"
body = {"seat_type": seat_type}
attempts = 1 + len(self._INVITE_PATCH_RETRY_DELAYS)
for attempt in range(attempts):
logger.info(
"[ChatGPT] 修改邀请 seat_type -> %s (attempt=%d/%d)...",
seat_type,
attempt + 1,
attempts,
)
result = self._api_fetch("PATCH", path, body)
status = result["status"]
if status == 200:
logger.info("[ChatGPT] seat_type 已改为 %s", seat_type)
return True
logger.error(
"[ChatGPT] 修改 seat_type 失败: %d %s",
status,
(result.get("body") or "")[:200],
)
if attempt < len(self._INVITE_PATCH_RETRY_DELAYS):
delay = self._INVITE_PATCH_RETRY_DELAYS[attempt]
logger.info("[ChatGPT] PATCH 失败,%ds 后重试", delay)
time.sleep(delay)
return False
def list_invites(self):
path = f"/backend-api/accounts/{self.account_id}/invites"
result = self._api_fetch("GET", path)
try:
return json.loads(result["body"])
except Exception:
return result["body"]
def stop(self):
self._close_http_transport(reason="stop")
close_playwright_objects(
page=self.page,
context=self.context,
browser=self.browser,
playwright=self.playwright,
logger=logger,
label="ChatGPTTeamAPI.stop",
)
self.browser = None
self.context = None
self.page = None
self.playwright = None
self.http_transport = None
self.transport_name = None