""" 多进程采集 Worker — 每个子进程处理单一平台的任务 每个进程独立加载平台 Client / WASM / cookie, 互不干扰。 主进程通过 Queue 分发任务、收集结果。 """ import json import os import random import re import sys import time from multiprocessing import Event from typing import Optional BASE_DIR = os.path.dirname(os.path.abspath(__file__)) if BASE_DIR not in sys.path: sys.path.insert(0, BASE_DIR) DEPS_DIR = os.path.join(BASE_DIR, "deps") if DEPS_DIR not in sys.path: sys.path.insert(0, DEPS_DIR) import requests as req_lib from config import COOKIE_PLATFORM_IDS, QUOTA_KEYWORDS from cookie_pool import CookiePool, parse_cookie_string NO_PROXY = {"http": None, "https": None} for _k in ("HTTP_PROXY", "HTTPS_PROXY", "http_proxy", "https_proxy"): os.environ.pop(_k, None) def _parse_and_check(plat_key, raw_body, raw_req, prompt, deep): """解析 + 基本校验 (答案非空、无限流话术), 不保存原始文件""" from collector import parse_only if isinstance(raw_body, bytes): raw_str = raw_body.decode("utf-8", errors="replace") else: raw_str = raw_body or "" if not raw_str: return None parsed = parse_only(plat_key, raw_str, raw_req) if not parsed: return None answer = (parsed.get("result") or {}).get("answer") or "" if not answer.strip(): return None if any(re.search(kw, answer) for kw in QUOTA_KEYWORDS): return {"_quota": True} parsed["deep_thinking"] = 1 if deep else 0 parsed["question"] = prompt parsed["_raw_req"] = raw_req or "" return parsed def _run_with_concurrency(worker_fn, worker_id, task_queue, result_queue, stop_event, cfg, platform_name): """通用并发包装器: 从配置读取并发数, 用线程池并发处理任务""" from concurrent.futures import ThreadPoolExecutor, as_completed import threading conc = max(1, int((cfg.get("concurrency") or {}).get(platform_name, 1))) print(f"[{platform_name}-{worker_id}] 并发={conc} pid={os.getpid()}") if conc == 1: # 单线程: 直接调用原始 worker worker_fn(worker_id, task_queue, result_queue, stop_event, cfg) return # 多线程: 线程池 thread_counter = [0] counter_lock = threading.Lock() def thread_worker(): with counter_lock: tid = thread_counter[0] thread_counter[0] += 1 thread_name = f"{platform_name}-{worker_id}-t{tid}" # 每个线程跑独立的 worker (各自独立 client/cookie) worker_fn(f"{worker_id}.{tid}", task_queue, result_queue, stop_event, cfg) with ThreadPoolExecutor(max_workers=conc) as executor: futures = [executor.submit(thread_worker) for _ in range(conc)] for f in as_completed(futures): try: f.result() except Exception as e: print(f"[{platform_name}-{worker_id}] 线程异常: {e}") # ============================================================ # DeepSeek Worker # ============================================================ def _ds_worker_single(worker_id, task_queue, result_queue, stop_event, cfg): """DS 单线程 worker 逻辑 (PoW WASM 不支持多线程, 但可多进程)""" from platforms.deepseek import DeepSeekClient, get_ds_session_id pool = CookiePool.from_config(cfg) client = None session_id = None def ensure_client(): nonlocal client, session_id if client: return client sess = pool.get_for_internal("ds") if not sess: print(f" [DS-{worker_id}] cookie 池无可用 session") return None client = DeepSeekClient(sess["cookie"], str(sess.get("id", ""))) session_id = sess.get("id") print(f" [DS-{worker_id}] 就绪 (pool id={session_id})") return client def drop_client(reason="invalid", error_msg=""): nonlocal client, session_id if session_id is not None: reload_time = "" if reason == "muted" and error_msg: m = re.search(r"""['"]mute_until['"]\s*:\s*([0-9.]+)""", error_msg) if m: try: reload_time = str(int(float(m.group(1)))) except Exception: pass pool.invalidate(session_id, reason) if reload_time: pool.update_session(session_id, reload_time, "3") client = None session_id = None print(f"[DS-{worker_id}] 启动 pid={os.getpid()}") while not stop_event.is_set(): try: task = task_queue.get(timeout=2) except Exception: continue if task is None: break did = task.get("dispatchId") or task.get("id") or "?" prompt = (task.get("prompt") or "").strip() deep = bool(int(task.get("deepThinking") or 0)) print(f" [DS-{worker_id}] 采集 DT={int(deep)} did={did} | {prompt[:40]}...") for attempt in range(3): c = ensure_client() if not c: result_queue.put({"task": task, "ok": False, "error": "no_client"}) break try: raw_body, raw_req = c.chat(prompt, deep) except Exception as e: err = str(e) print(f" [DS-{worker_id}] chat 异常: {err[:100]}") if "Authorization Failed" in err or "invalid token" in err: drop_client("invalid") elif "user is muted" in err: drop_client("muted", error_msg=err) else: drop_client("invalid") if attempt < 2: time.sleep(random.uniform(3, 8)) continue result_queue.put({"task": task, "ok": False, "error": err[:200]}) break # 限流检测 (在解析前, 因为限流时答案为空) raw_str = raw_body if isinstance(raw_body, str) else raw_body.decode("utf-8", errors="replace") if "rate_limit" in raw_str or "过于频繁" in raw_str: if attempt == 0: print(f" [DS-{worker_id}] ⏳ 限流, 等待 10s 后同号重试") time.sleep(10) continue else: print(f" [DS-{worker_id}] ⏳ 限流, 换号重试") drop_client("invalid") time.sleep(random.uniform(3, 5)) if attempt < 2: continue result_queue.put({"task": task, "ok": False, "error": "rate_limit"}) break parsed = _parse_and_check("ds", raw_body, raw_req, prompt, deep) if not parsed: if attempt < 2: print(f" [DS-{worker_id}] 解析为空, 重试 ({attempt + 1}/3)") time.sleep(random.uniform(3, 8)) continue result_queue.put({"task": task, "ok": False, "error": "parse_failed"}) break if parsed.get("_quota"): drop_client("quota") if attempt < 2: time.sleep(random.uniform(5, 15)) continue result_queue.put({"task": task, "ok": False, "error": "quota"}) break raw_str = raw_body if isinstance(raw_body, str) else raw_body.decode("utf-8", errors="replace") src_count = len(parsed.get("sources") or []) if src_count == 0: search_down = "搜索暂不可用" in raw_str or "搜索不可用" in raw_str if search_down: print(f" [DS-{worker_id}] ⚠ 联网搜索暂不可用 (平台故障)") result_queue.put({"task": task, "ok": False, "error": "search_unavailable"}) break print(f" [DS-{worker_id}] 信源为 0, 换号重试 ({attempt + 1}/3)") drop_client("invalid") if attempt < 2: time.sleep(random.uniform(3, 8)) continue result_queue.put({"task": task, "ok": False, "error": "no_sources"}) break # 分享链接 sid = get_ds_session_id(parsed) if sid and c: try: sr = req_lib.post( "https://chat.deepseek.com/api/v0/share/create", headers=c.headers, json={"chat_session_id": sid, "message_ids": [1, 2]}, timeout=30, verify=False, proxies=NO_PROXY, ) biz = ((sr.json() or {}).get("data") or {}).get("biz_data") or {} sh_id = biz.get("share_id", "") if isinstance(biz, dict) else "" if sh_id: parsed["result"]["share_link"] = f"https://chat.deepseek.com/share/{sh_id}" except Exception: pass think_len = len((parsed["result"].get("thinking_process") or "")) print(f" [DS-{worker_id}] ✓ 答案 {len(parsed['result']['answer'])} 字 | 思考 {think_len} 字 | 信源 {src_count}") result_queue.put({"task": task, "ok": True, "parsed": parsed}) break else: result_queue.put({"task": task, "ok": False, "error": "max_retry"}) print(f"[DS-{worker_id}] 退出") def ds_worker(worker_id, task_queue, result_queue, stop_event, cfg): _run_with_concurrency(_ds_worker_single, worker_id, task_queue, result_queue, stop_event, cfg, "deepseek") # ============================================================ # Kimi Worker # ============================================================ def _kimi_worker_single(worker_id, task_queue, result_queue, stop_event, cfg): from platforms.kimi import KimiClient pool = CookiePool.from_config(cfg) client = None session_id = None def ensure_client(): nonlocal client, session_id if client: return client sess = pool.get_for_internal("kimi") token = sess.get("cookie") if sess else None if not token: print(f" [Kimi-{worker_id}] cookie 池无可用 session") return None try: client = KimiClient(token) session_id = sess.get("id") client.cookie_id = session_id print(f" [Kimi-{worker_id}] 就绪 (pool id={session_id})") return client except Exception as e: print(f" [Kimi-{worker_id}] 初始化失败: {e}") return None def drop_client(reason="invalid"): nonlocal client, session_id client = None session_id = None def diag_response(raw_body): if isinstance(raw_body, bytes): raw_str = raw_body.decode("utf-8", errors="replace") else: raw_str = raw_body or "" if "unauthenticated" in raw_str or "请登录" in raw_str: return "unauthenticated" if "OVERLOADED" in raw_str or "Kimi有点累了" in raw_str or "聊的人太多了" in raw_str or "resource_exhausted" in raw_str: return "overloaded" if len(raw_str) < 10000 and any(kw in raw_str for kw in ("算力不足", "前往升级", "对话次数已达上限", "免费次数用完", "次数已用完")): return "quota" if "account_abnormal" in raw_str or "账号状态异常" in raw_str: return "account_abnormal" return "ok" print(f"[Kimi-{worker_id}] 启动 pid={os.getpid()}") max_retry = 2 max_overload = 5 while not stop_event.is_set(): try: task = task_queue.get(timeout=2) except Exception: continue if task is None: break did = task.get("dispatchId") or task.get("id") or "?" prompt = (task.get("prompt") or "").strip() deep = bool(int(task.get("deepThinking") or 0)) print(f" [Kimi-{worker_id}] 采集 DT={int(deep)} did={did} | {prompt[:40]}...") overload_retries = 0 for attempt in range(max_retry + 1 + max_overload): c = ensure_client() if not c: result_queue.put({"task": task, "ok": False, "error": "no_client"}) break try: raw_body, raw_req = c.chat(prompt, bool(deep)) except Exception as e: err = str(e) print(f" [Kimi-{worker_id}] chat 异常: {err[:100]}") if "unauthenticated" in err or "请登录" in err: drop_client("invalid") if attempt < max_retry: time.sleep(random.uniform(3, 8)) continue result_queue.put({"task": task, "ok": False, "error": err[:200]}) break diag = diag_response(raw_body) if diag in ("unauthenticated", "account_abnormal"): raw_len = len(raw_body) if raw_body else 0 print(f" [Kimi-{worker_id}] ✗ {diag} (响应 {raw_len}B)") # 先尝试 refresh_token 刷新 if diag == "unauthenticated" and c and hasattr(c, "refresh_access_token"): if c.refresh_access_token(): print(f" [Kimi-{worker_id}] ↻ token 已刷新, 重试") time.sleep(1) continue drop_client("invalid") if attempt < max_retry: print(f" [Kimi-{worker_id}] ↻ 换号重试 ({attempt + 1}/{max_retry})") time.sleep(random.uniform(2, 5)) continue result_queue.put({"task": task, "ok": False, "error": diag}) break if diag == "overloaded": overload_retries += 1 if overload_retries <= max_overload: wait = random.uniform(3, 5) print(f" [Kimi-{worker_id}] ⏳ 算力过载, 等待 {wait:.0f}s ({overload_retries}/{max_overload})") time.sleep(wait) continue result_queue.put({"task": task, "ok": False, "error": "overloaded"}) break if diag == "quota": # 保存原始响应用于调试 try: from datetime import datetime as _dt debug_dir = os.path.join(BASE_DIR, "debug_kimi_quota") os.makedirs(debug_dir, exist_ok=True) ts = _dt.now().strftime("%Y%m%d_%H%M%S") raw_bytes = raw_body if isinstance(raw_body, bytes) else raw_body.encode("utf-8", errors="replace") dump_path = os.path.join(debug_dir, f"kimi_quota_{ts}_W{worker_id}.bin") with open(dump_path, "wb") as df: df.write(raw_bytes) print(f" [Kimi-{worker_id}] 配额调试: 已保存 {dump_path} ({len(raw_bytes)}B)") except Exception as de: print(f" [Kimi-{worker_id}] 配额调试保存失败: {de}") drop_client("quota") if attempt < max_retry: print(f" [Kimi-{worker_id}] ↻ 配额用完换号 ({attempt + 1}/{max_retry})") time.sleep(random.uniform(5, 15)) continue result_queue.put({"task": task, "ok": False, "error": "quota"}) break parsed = _parse_and_check("kimi", raw_body, raw_req, prompt, deep) if not parsed: if attempt < max_retry: drop_client("invalid") time.sleep(random.uniform(3, 6)) continue result_queue.put({"task": task, "ok": False, "error": "parse_failed"}) break if parsed.get("_quota"): drop_client("quota") if attempt < max_retry: time.sleep(random.uniform(5, 15)) continue result_queue.put({"task": task, "ok": False, "error": "quota"}) break # 信源为 0 换号 src_count = len(parsed.get("sources") or []) if src_count == 0 and attempt < max_retry: print(f" [Kimi-{worker_id}] ⚠ 信源为 0, 换号重试 ({attempt + 1}/{max_retry})") drop_client("invalid") time.sleep(random.uniform(3, 8)) continue # 分享链接 chat_id = parsed.get("chat_id", "") msg_ids = parsed.get("message_ids", []) if chat_id and msg_ids and c: try: headers = c._build_headers() headers["content-type"] = "application/json" headers.pop("connect-protocol-version", None) resp = c.requests.post( "https://www.kimi.com/apiv2/kimi.gateway.chat.v1.ChatService/CreateChatShare", headers=headers, json={"chat_id": chat_id, "message_ids": msg_ids}, timeout=30, impersonate="chrome146", ) if resp.status_code == 200 and resp.content: raw = resp.content js = raw.find(b"{") if js >= 0: resp_data = json.loads(raw[js: raw.rfind(b"}") + 1]) share_id = (resp_data.get("share") or {}).get("id", "") or resp_data.get("id", "") if share_id: parsed["result"]["share_link"] = f"https://www.kimi.com/share/{share_id}" except Exception: pass think_len = len((parsed["result"].get("thinking_process") or "")) print(f" [Kimi-{worker_id}] ✓ 答案 {len(parsed['result']['answer'])} 字 | 思考 {think_len} 字 | 信源 {src_count}") result_queue.put({"task": task, "ok": True, "parsed": parsed}) break else: result_queue.put({"task": task, "ok": False, "error": "max_retry"}) print(f"[Kimi-{worker_id}] 退出") def kimi_worker(worker_id, task_queue, result_queue, stop_event, cfg): _run_with_concurrency(_kimi_worker_single, worker_id, task_queue, result_queue, stop_event, cfg, "kimi") # ============================================================ # 千问 Worker # ============================================================ def _qianwen_worker_single(worker_id, task_queue, result_queue, stop_event, cfg): import asyncio from platforms.qianwen import QianwenClient, get_qw_session_id pool = CookiePool.from_config(cfg) client = None session_id = None def ensure_client(): nonlocal client, session_id if client and client._cookie_str: return client c = QianwenClient() sess = pool.get_for_internal("qianwen") if sess and sess.get("cookie"): cookie_str = sess["cookie"] c._cookie_str = cookie_str if "=" in cookie_str else cookie_str jar = parse_cookie_string(cookie_str) if jar: c._cookie_str = "; ".join(f"{k}={v}" for k, v in jar.items()) c._xsrf_token = jar.get("XSRF-TOKEN", "") session_id = sess.get("id") print(f" [QW-{worker_id}] 就绪 (pool id={session_id})") else: print(f" [QW-{worker_id}] cookie 池无数据, 尝试 HMAC 匿名") client = c return client def drop_client(reason="invalid"): nonlocal client, session_id if session_id is not None: pool.invalidate(session_id, reason) client = None session_id = None print(f"[QW-{worker_id}] 启动 pid={os.getpid()}") loop = asyncio.new_event_loop() while not stop_event.is_set(): try: task = task_queue.get(timeout=2) except Exception: continue if task is None: break did = task.get("dispatchId") or task.get("id") or "?" prompt = (task.get("prompt") or "").strip() deep = bool(int(task.get("deepThinking") or 0)) print(f" [QW-{worker_id}] 采集 DT={int(deep)} did={did} | {prompt[:40]}...") c = ensure_client() try: raw_body, raw_req = loop.run_until_complete(c.chat(prompt, thinking=deep)) except Exception as e: err = str(e).lower() if "401" in err or "login" in err or "auth" in err: drop_client("invalid") result_queue.put({"task": task, "ok": False, "error": str(e)[:200]}) continue parsed = _parse_and_check("qianwen", raw_body, raw_req, prompt, deep) if not parsed: result_queue.put({"task": task, "ok": False, "error": "parse_failed"}) continue if parsed.get("_quota"): drop_client("quota") result_queue.put({"task": task, "ok": False, "error": "quota"}) continue # 分享链接 if c._cookie_str: qw_sid = get_qw_session_id(parsed) if qw_sid: try: share_id = loop.run_until_complete(c.create_share(qw_sid)) if share_id: parsed["result"]["share_link"] = f"https://www.qianwen.com/share/chat/{share_id}" except Exception: pass src_count = len(parsed.get("sources") or []) think_len = len((parsed["result"].get("thinking_process") or "")) print(f" [QW-{worker_id}] ✓ 答案 {len(parsed['result']['answer'])} 字 | 思考 {think_len} 字 | 信源 {src_count}") result_queue.put({"task": task, "ok": True, "parsed": parsed}) loop.close() print(f"[QW-{worker_id}] 退出") def qianwen_worker(worker_id, task_queue, result_queue, stop_event, cfg): _run_with_concurrency(_qianwen_worker_single, worker_id, task_queue, result_queue, stop_event, cfg, "tongyi") # ============================================================ # 文心 Worker # ============================================================ def _wenxin_worker_single(worker_id, task_queue, result_queue, stop_event, cfg): import asyncio from platforms.wenxin import WenxinClient pool = CookiePool.from_config(cfg) client = None session_id = None def ensure_client(): nonlocal client, session_id if client and client._cookie_str: return client c = WenxinClient() sess = pool.get_for_internal("wenxin") if sess and sess.get("cookie"): cookie_str = sess["cookie"] jar = parse_cookie_string(cookie_str) c._cookie_str = "; ".join(f"{k}={v}" for k, v in jar.items()) if jar else cookie_str session_id = sess.get("id") print(f" [WX-{worker_id}] 就绪 (pool id={session_id})") else: print(f" [WX-{worker_id}] cookie 池无数据") client = c return client def drop_client(reason="invalid"): nonlocal client, session_id if session_id is not None: pool.invalidate(session_id, reason) client = None session_id = None print(f"[WX-{worker_id}] 启动 pid={os.getpid()}") loop = asyncio.new_event_loop() while not stop_event.is_set(): try: task = task_queue.get(timeout=2) except Exception: continue if task is None: break did = task.get("dispatchId") or task.get("id") or "?" prompt = (task.get("prompt") or "").strip() deep = bool(int(task.get("deepThinking") or 0)) print(f" [WX-{worker_id}] 采集 DT={int(deep)} did={did} | {prompt[:40]}...") c = ensure_client() try: raw_body, raw_req = loop.run_until_complete(c.chat(prompt, thinking=deep)) except Exception as e: result_queue.put({"task": task, "ok": False, "error": str(e)[:200]}) continue parsed = _parse_and_check("wenxin", raw_body, raw_req, prompt, deep) if not parsed: result_queue.put({"task": task, "ok": False, "error": "parse_failed"}) continue if parsed.get("_quota"): drop_client("quota") result_queue.put({"task": task, "ok": False, "error": "quota"}) continue # 分享链接 (多进程模式不用浏览器, 用页面链接) lid = parsed.get("conversation_lid", "") if lid: parsed["result"]["share_link"] = f"https://chat.baidu.com/search/{lid}" src_count = len(parsed.get("sources") or []) think_len = len((parsed["result"].get("thinking_process") or "")) print(f" [WX-{worker_id}] ✓ 答案 {len(parsed['result']['answer'])} 字 | 思考 {think_len} 字 | 信源 {src_count}") result_queue.put({"task": task, "ok": True, "parsed": parsed}) loop.close() print(f"[WX-{worker_id}] 退出") def wenxin_worker(worker_id, task_queue, result_queue, stop_event, cfg): _run_with_concurrency(_wenxin_worker_single, worker_id, task_queue, result_queue, stop_event, cfg, "baiduai") # ============================================================ # 豆包 Worker (每进程独立浏览器, 多 Tab 并发, 验证码打码+重启) # doubao 的并发由 doubao.accounts × doubao.tabs_per_account 控制, 不走线程池 # ============================================================ def doubao_worker(worker_id, task_queue, result_queue, stop_event, cfg): import asyncio import socket import subprocess dc = cfg.get("doubao") or {} chrome_path = dc.get("chrome_path") or "" profile_base = dc.get("profile_dir") or "doubao_profiles" if not os.path.isabs(profile_base): profile_base = os.path.join(BASE_DIR, profile_base) profile_dir = os.path.join(profile_base, f"w{worker_id}") proxy = dc.get("proxy") or "" wait_timeout = int(dc.get("wait_timeout") or 180) share_enabled = bool(dc.get("share", True)) max_retry = int(dc.get("max_retry") or 2) captcha_timeout = int(dc.get("captcha_timeout") or 120) pool = CookiePool.from_config(cfg) DOUBAO_DOMAINS = [".doubao.com", "www.doubao.com", "doubao.com"] # 打码平台 captcha_client = None captcha_server = dc.get("captcha_server") or "" if captcha_server: try: from client_sdk import CaptchaClient captcha_client = CaptchaClient(captcha_server, dc.get("captcha_key") or "") except ImportError: pass print(f"[DB-{worker_id}] 启动 pid={os.getpid()}") if not chrome_path or not os.path.isfile(chrome_path): print(f"[DB-{worker_id}] ✗ ungoogled-chromium 未找到: {chrome_path}") return async def _run(): os.environ["NO_PROXY"] = "*" os.makedirs(profile_dir, exist_ok=True) with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s: s.bind(("", 0)) port = s.getsockname()[1] chrome_args = [ chrome_path, f"--remote-debugging-port={port}", f"--user-data-dir={profile_dir}", "--disable-blink-features=AutomationControlled", "--no-first-run", "--no-default-browser-check", "--disk-cache-size=104857600", f"--fingerprint={random.getrandbits(32)}", "--fingerprint-platform=windows", "--fingerprint-brand=Chrome", "--fingerprint-canvas=noise", "--fingerprint-webgl=noise", "--ignore-certificate-errors", "--remote-allow-origins=*", ] if proxy: chrome_args.append(f"--proxy-server={proxy}") proc = subprocess.Popen(chrome_args, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL) print(f"[DB-{worker_id}] Chrome pid={proc.pid} port={port}") for i in range(15): await asyncio.sleep(1) try: _, w = await asyncio.open_connection("127.0.0.1", port) w.close(); await w.wait_closed() print(f"[DB-{worker_id}] CDP ready ({i+1}s)"); break except Exception: pass else: print(f"[DB-{worker_id}] CDP 超时"); proc.kill(); return from playwright.async_api import async_playwright from platforms.doubao import DoubaoAdapter from cookie_pool import cookie_str_to_playwright pw = await async_playwright().start() browser = await pw.chromium.connect_over_cdp(f"http://127.0.0.1:{port}") ctx = browser.contexts[0] async def check_login(pg): try: return await pg.evaluate("""() => { const btns = document.querySelectorAll('button, a, div[role="button"]'); for (const b of btns) { const t = (b.textContent || '').trim(); const rect = b.getBoundingClientRect(); if (t === '登录' && rect.right > window.innerWidth * 0.7 && rect.top < 80) return false; } return true; }""") except Exception: return False # 注入 cookie + 登录检测 (最多尝试 5 个 cookie) current_sid = None page = await ctx.new_page() for _try in range(5): sess = pool.get_for_internal("doubao") if not sess or not sess.get("cookie"): print(f"[DB-{worker_id}] cookie 池无数据"); break await ctx.clear_cookies() cookies = cookie_str_to_playwright(sess["cookie"], DOUBAO_DOMAINS) await ctx.add_cookies(cookies) current_sid = sess.get("id") await page.goto("https://www.doubao.com/chat", wait_until="domcontentloaded", timeout=60000) await asyncio.sleep(3) if await check_login(page): print(f"[DB-{worker_id}] 登录态 OK (pool id={current_sid})"); break print(f"[DB-{worker_id}] cookie 无效 (id={current_sid}), 换号 ({_try+1}/5)") current_sid = None if not current_sid: print(f"[DB-{worker_id}] 无可用 cookie, 退出"); await pw.stop(); proc.kill(); return cdp = await ctx.new_cdp_session(page) await cdp.send("Network.enable") adapter = DoubaoAdapter(page, cdp) print(f"[DB-{worker_id}] 就绪") # 换号重启 async def restart_instance(): nonlocal current_sid, page, adapter, cdp # 提交打码 if captcha_client and adapter and adapter.captcha_decision: fp = "" try: for c in await ctx.cookies(): if c.get("name") == "s_v_web_id": fp = c.get("value", ""); break except Exception: pass if fp: print(f" [DB-{worker_id}] 提交打码 fp={fp[:30]}...") try: loop = asyncio.get_running_loop() tid = await loop.run_in_executor(None, captcha_client.submit_task, fp, adapter.captcha_decision, "v2-db") result = await loop.run_in_executor(None, captcha_client.wait_for_completion, tid, captcha_timeout) print(f" [DB-{worker_id}] 打码{'成功' if result.get('success') else '失败'}") except Exception as e: print(f" [DB-{worker_id}] 打码异常: {e}") else: await asyncio.sleep(10) else: await asyncio.sleep(10) # 关闭旧页面 for p in ctx.pages: try: await p.close() except Exception: pass # 换号 sess = pool.get_for_internal("doubao") if sess and sess.get("cookie"): await ctx.clear_cookies() cks = cookie_str_to_playwright(sess["cookie"], DOUBAO_DOMAINS) await ctx.add_cookies(cks) current_sid = sess.get("id") print(f" [DB-{worker_id}] 换号 → id={current_sid}") # 新建页面 page = await ctx.new_page() await page.goto("https://www.doubao.com/chat", wait_until="domcontentloaded", timeout=60000) await asyncio.sleep(3) try: if cdp: await cdp.detach() except Exception: pass cdp = await ctx.new_cdp_session(page) await cdp.send("Network.enable") adapter = DoubaoAdapter(page, cdp) # 主循环 while not stop_event.is_set(): try: task = task_queue.get(timeout=2) except Exception: continue if task is None: break did = task.get("dispatchId") or task.get("id") or "?" prompt = (task.get("prompt") or "").strip() deep = bool(int(task.get("deepThinking") or 0)) print(f" [DB-{worker_id}] 采集 DT={int(deep)} did={did} | {prompt[:40]}...") parsed = None for attempt in range(max_retry + 1): if stop_event.is_set(): break try: print(f" [DB-{worker_id}] attempt={attempt+1} sid={current_sid}") await adapter.new_chat() if not await check_login(page): print(f" [DB-{worker_id}] 未登录, 换号") await restart_instance() if attempt < max_retry: continue break await adapter.setup_intercept() await adapter.set_deep_thinking(deep) await adapter.send_question(prompt) # 等待 + 验证码检测 captcha_detected = False wait_start = time.time() while not adapter._finished_event.is_set(): if stop_event.is_set(): break if time.time() - wait_start > wait_timeout: print(f" [DB-{worker_id}] 响应超时"); break try: has_captcha = await page.evaluate("""() => { var c = document.getElementById('captcha_container'); return c && c.offsetHeight > 0 && getComputedStyle(c).display !== 'none'; }""") if has_captcha and not captcha_detected: captcha_detected = True print(f" [DB-{worker_id}] ⚠ 验证码!") break except Exception: pass await asyncio.sleep(3) if captcha_detected: await restart_instance() continue await asyncio.sleep(2) raw_body, raw_req = adapter.get_last_response() if not raw_body: page_txt = "" try: page_txt = await page.inner_text("body", timeout=3000) except Exception: pass if any(k in page_txt for k in ("免费额度用完", "得休息一阵子了")): print(f" [DB-{worker_id}] 配额用完, 换号") await restart_instance() if attempt < max_retry: continue break if attempt < max_retry: continue break raw_str = raw_body.decode("utf-8", errors="replace") if isinstance(raw_body, bytes) else raw_body if any(k in raw_str for k in ("免费额度用完了", "得休息一阵子了")): print(f" [DB-{worker_id}] 配额用完 (SSE), 换号") await restart_instance() if attempt < max_retry: continue break if any(k in raw_str[:500] for k in ('"code":401', '"code":403', '"未登录"')): print(f" [DB-{worker_id}] 鉴权错误, 换号") await restart_instance() if attempt < max_retry: continue break parsed = _parse_and_check("doubao", raw_body, raw_req, prompt, deep) if not parsed: if attempt < max_retry: continue break # 分享链接 if share_enabled: conv_id = parsed.get("conversation_id", "") if conv_id: try: share_url = await page.evaluate("""async (convId) => { const aid = "497858"; try { const r1 = await fetch("/samantha/thread/share/info?aid=" + aid, { method: "POST", headers: {"Content-Type": "application/json"}, body: JSON.stringify({conversation_id: convId}), credentials: "include" }); const d1 = await r1.json(); const shareToken = d1?.data?.share_token || ""; const preShareId = d1?.data?.pre_share_id || ""; const r2 = await fetch("/samantha/thread/share/save?aid=" + aid, { method: "POST", headers: {"Content-Type": "application/json"}, body: JSON.stringify({conversation_id: convId, share_token: shareToken, share_id: preShareId, message_index_end: 2, share_type: 1, is_allow_seo: false}), credentials: "include" }); const d2 = await r2.json(); const sid = d2?.data?.share_id || ""; return sid ? "https://www.doubao.com/thread/" + sid : ""; } catch { return ""; } }""", conv_id) if share_url: parsed["result"]["share_link"] = share_url print(f" [DB-{worker_id}] 分享链接 OK") except Exception: pass src_count = len(parsed.get("sources") or []) think_len = len((parsed["result"].get("thinking_process") or "")) print(f" [DB-{worker_id}] ✓ 答案 {len(parsed['result']['answer'])} 字 | 思考 {think_len} 字 | 信源 {src_count}") result_queue.put({"task": task, "ok": True, "parsed": parsed}) break except Exception as e: print(f" [DB-{worker_id}] 异常: {e}") if attempt < max_retry: await asyncio.sleep(2) continue break if parsed is None and not stop_event.is_set(): result_queue.put({"task": task, "ok": False, "error": "collect_failed"}) try: await pw.stop() except Exception: pass proc.kill() print(f"[DB-{worker_id}] 退出") loop = asyncio.new_event_loop() try: loop.run_until_complete(_run()) finally: loop.close() # ============================================================ # Worker 注册表 # ============================================================ WORKER_FUNCS = { "deepseek": ds_worker, "kimi": kimi_worker, "tongyi": qianwen_worker, "baiduai": wenxin_worker, "doubao": doubao_worker, }