Merge remote-tracking branch 'origin/codex/direct-mihomo-fallback' into HEAD

# Conflicts:
#	DouYinSparkFlow/core/friends.py
#	DouYinSparkFlow/login_desktop_server.py
This commit is contained in:
Rixuan Shao
2026-08-25 18:59:46 +08:00
8 changed files with 395 additions and 79 deletions
+76 -10
View File
@@ -8,7 +8,7 @@ from pathlib import Path
from playwright.async_api import async_playwright
from rich.console import Console
from utils.config import DEBUG, Environment, get_environment
from utils.config import DEBUG, Environment, get_app_settings, get_environment
console = Console()
@@ -52,6 +52,75 @@ def _browser_args():
]
def _douyin_network_mode():
settings = get_app_settings(force_reload=True)
return str(
os.getenv("SPARKFLOW_DOUYIN_NETWORK_MODE")
or settings.get("douyin_network_mode", "direct")
).strip().lower()
def douyin_network_modes():
# Direct is the default; Mihomo is the fallback unless explicitly selected.
mode = _douyin_network_mode()
if mode == "mihomo":
return ("mihomo",)
return ("direct", "mihomo")
def _douyin_browser_proxy(network_mode=None):
# Return an explicit proxy URL for Douyin traffic, or None for direct.
settings = get_app_settings(force_reload=True)
mode = str(network_mode or _douyin_network_mode()).strip().lower()
if mode != "mihomo":
return None
return str(
os.getenv("SPARKFLOW_DOUYIN_PROXY_URL")
or settings.get("douyin_proxy_url", "http://proxy:7890")
).strip() or None
def _browser_launch_options(GUI=False, network_mode=None):
args = _browser_args()
proxy = _douyin_browser_proxy(network_mode=network_mode)
if proxy:
return {
"headless": _headless_for(GUI),
"args": args,
"proxy": {"server": proxy},
}
args.append("--no-proxy-server")
return {
"headless": _headless_for(GUI),
"args": args,
}
async def select_douyin_network_mode(target_url):
# Select the first route that can load the target before a task starts.
failures = []
for network_mode in douyin_network_modes():
playwright = browser = page = None
try:
playwright, browser = await get_browser(network_mode=network_mode)
page = await browser.new_page()
response = await page.goto(target_url, wait_until="commit", timeout=30000)
status = response.status if response is not None else None
if status is not None and status < 500:
return network_mode
failures.append(f"{network_mode}: HTTP {status}")
except Exception as exc:
failures.append(f"{network_mode}: {exc}")
finally:
if page:
await page.close()
if browser:
await browser.close()
if playwright:
await playwright.stop()
raise RuntimeError(f"Douyin network preflight failed: {'; '.join(failures)}")
def sanitize_profile_name(value):
raw = str(value or "").strip()
if not raw:
@@ -78,15 +147,12 @@ async def install_browser():
console.print(f"[bold red]Browser install failed: {exc}[/bold red]")
async def get_browser(GUI=False):
async def get_browser(GUI=False, network_mode=None):
configure_playwright_environment()
try:
playwright = await async_playwright().start()
browser = await playwright.chromium.launch(
headless=_headless_for(GUI),
args=_browser_args(),
)
browser = await playwright.chromium.launch(**_browser_launch_options(GUI, network_mode=network_mode))
return playwright, browser
except Exception as exc:
if "Executable doesn't exist" in str(exc) and get_environment() != Environment.GITHUBACTION:
@@ -97,7 +163,7 @@ async def get_browser(GUI=False):
raise
async def get_persistent_browser_context(profile_name, GUI=False, root=None):
async def get_persistent_browser_context(profile_name, GUI=False, root=None, network_mode=None):
configure_playwright_environment()
profile_dir = browser_profile_root(root) / sanitize_profile_name(profile_name)
@@ -105,11 +171,11 @@ async def get_persistent_browser_context(profile_name, GUI=False, root=None):
try:
playwright = await async_playwright().start()
launch_options = _browser_launch_options(GUI, network_mode=network_mode)
launch_options["viewport"] = {"width": 1600, "height": 1000}
context = await playwright.chromium.launch_persistent_context(
str(profile_dir),
headless=_headless_for(GUI),
viewport={"width": 1600, "height": 1000},
args=_browser_args(),
**launch_options,
)
return playwright, context, profile_dir
except Exception as exc:
+43 -7
View File
@@ -1,5 +1,10 @@
import asyncio
from core.browser import get_browser
import logging
from core.browser import douyin_network_modes, get_browser
logger = logging.getLogger(__name__)
CHAT_PAGE_URL = "https://creator.douyin.com/creator-micro/data/following/chat"
@@ -218,19 +223,15 @@ async def collect_friend_names(page):
return found_names
async def fetch_account_friends(account):
async def _fetch_account_friends_once(account, network_mode):
cookies = list(account.get("cookies") or [])
if not cookies:
raise RuntimeError("账号没有可用 cookies,请重新扫码登录")
playwright = browser = context = page = None
try:
playwright, browser = await get_browser(GUI=False)
playwright, browser = await get_browser(GUI=False, network_mode=network_mode)
context = await browser.new_context()
context.set_default_navigation_timeout(120000)
context.set_default_timeout(120000)
page = await context.new_page()
await context.add_cookies(cookies)
await page.goto(CHAT_PAGE_URL, wait_until="commit", timeout=30000)
await asyncio.sleep(1)
@@ -250,3 +251,38 @@ async def fetch_account_friends(account):
await browser.close()
if playwright:
await playwright.stop()
async def fetch_account_friends(account):
cookies = list(account.get("cookies") or [])
if not cookies:
raise RuntimeError("account has no cookies; scan login QR code first")
modes = douyin_network_modes()
last_error = None
for index, network_mode in enumerate(modes):
try:
friends = await _fetch_account_friends_once(account, network_mode)
logger.info(
"Friend refresh route=%s count=%s attempt=%s/%s",
network_mode,
len(friends),
index + 1,
len(modes),
)
if friends or index == len(modes) - 1:
return friends
logger.warning(
"Friend refresh route=%s returned zero friends; trying next route",
network_mode,
)
except RuntimeError as exc:
text = str(exc).lower()
if any(marker in text for marker in ("login", "cookie", "scan", "登录", "扫码")):
raise
last_error = exc
logger.warning("Friend refresh route=%s failed; trying next route: %s", network_mode, exc)
except Exception as exc:
last_error = exc
logger.warning("Friend refresh route=%s failed; trying next route: %s", network_mode, exc)
raise RuntimeError(f"friend refresh failed after routes {modes}: {last_error}")
+15 -6
View File
@@ -12,7 +12,12 @@ from datetime import datetime, timedelta, timezone
from pathlib import Path
from zoneinfo import ZoneInfo
from core.browser import get_browser, get_persistent_browser_context, sanitize_profile_name
from core.browser import (
get_browser,
get_persistent_browser_context,
sanitize_profile_name,
select_douyin_network_mode,
)
from core.msg_builder import build_message, build_message_candidates
from core.protocol_dispatch import run_protocol_tasks
from core.send_state import parse_sent_at, target_is_strong_confirmed_today
@@ -2354,6 +2359,8 @@ async def run_browser_tasks(active_config, browser_user_data):
send_strategy = _normalize_send_strategy(active_config)
friend_scan_config = _normalize_friend_list_scan_config(active_config)
profile_config = _normalize_persistent_profile_config(active_config)
network_mode = await select_douyin_network_mode(CREATOR_HOME_URL)
logger.info("Selected Douyin task network route=%s", network_mode)
semaphore = asyncio.Semaphore(active_config["taskCount"] if active_config["multiTask"] else 1)
tasks = []
@@ -2371,11 +2378,11 @@ async def run_browser_tasks(active_config, browser_user_data):
user.get("username", "unknown"),
len(user["targets"]),
)
tasks.append(do_user_task(None, user, semaphore, send_strategy, profile_config, friend_scan_config))
tasks.append(do_user_task(None, user, semaphore, send_strategy, profile_config, friend_scan_config, network_mode))
await asyncio.gather(*tasks)
return
playwright, browser = await get_browser()
playwright, browser = await get_browser(network_mode=network_mode)
try:
for user in browser_user_data:
logger.info(
@@ -2383,7 +2390,7 @@ async def run_browser_tasks(active_config, browser_user_data):
user.get("username", "unknown"),
len(user["targets"]),
)
tasks.append(do_user_task(browser, user, semaphore, send_strategy, profile_config, friend_scan_config))
tasks.append(do_user_task(browser, user, semaphore, send_strategy, profile_config, friend_scan_config, network_mode))
await asyncio.gather(*tasks)
finally:
@@ -2391,7 +2398,7 @@ async def run_browser_tasks(active_config, browser_user_data):
await browser.close()
async def do_user_task(browser, user, semaphore, send_strategy, profile_config, friend_scan_config):
async def do_user_task(browser, user, semaphore, send_strategy, profile_config, friend_scan_config, network_mode):
async with semaphore:
account_name = user.get("username", "unknown")
account_lock_handle = None
@@ -2416,6 +2423,7 @@ async def do_user_task(browser, user, semaphore, send_strategy, profile_config,
profile_config,
friend_scan_config,
account_name,
network_mode,
),
timeout=timeout_seconds,
)
@@ -2437,7 +2445,7 @@ async def do_user_task(browser, user, semaphore, send_strategy, profile_config,
_release_browser_account_lock(account_lock_handle, account_lock_path, account_name)
async def _do_user_task_locked(browser, user, send_strategy, profile_config, friend_scan_config, account_name):
async def _do_user_task_locked(browser, user, send_strategy, profile_config, friend_scan_config, account_name, network_mode):
cookies = user["cookies"]
targets = user["targets"]
start_delay = _random_delay_seconds(
@@ -2453,6 +2461,7 @@ async def _do_user_task_locked(browser, user, send_strategy, profile_config, fri
owned_playwright, context, profile_dir = await get_persistent_browser_context(
_account_profile_name(user),
root=profile_config["root"],
network_mode=network_mode,
)
logger.info("Opened persistent browser profile for %s at %s", account_name, profile_dir)
if profile_config["syncStoredCookiesBeforeRun"]:
+198 -10
View File
@@ -2,8 +2,10 @@ import asyncio
import os
import shutil
import time
import urllib.request
from contextlib import asynccontextmanager
from pathlib import Path
from urllib.parse import urlsplit
from fastapi import FastAPI, HTTPException, Request, Response
import uvicorn
@@ -28,6 +30,16 @@ if LOGIN_DESKTOP_MODE not in {"native", "novnc"}:
IDLE_TIMEOUT_SECONDS = max(300, int(os.getenv("LOGIN_DESKTOP_IDLE_TIMEOUT_SECONDS", "1800")))
STOP_AFTER_EXPORT_SECONDS = max(0, int(os.getenv("LOGIN_DESKTOP_STOP_AFTER_EXPORT_SECONDS", "60")))
STATUS_CACHE_SECONDS = max(1, int(os.getenv("LOGIN_DESKTOP_STATUS_CACHE_SECONDS", "15")))
LOGIN_NETWORK_MODE = str(os.getenv("LOGIN_DESKTOP_PROXY_MODE", "auto")).strip().lower()
if LOGIN_NETWORK_MODE not in {"auto", "direct", "proxy"}:
LOGIN_NETWORK_MODE = "auto"
LOGIN_PROXY_SERVER = str(os.getenv("LOGIN_DESKTOP_PROXY", "http://proxy:7890")).strip()
LOGIN_PREFLIGHT_TIMEOUT_SECONDS = max(
3, int(os.getenv("LOGIN_DESKTOP_PREFLIGHT_TIMEOUT_SECONDS", "15"))
)
LOGIN_NETWORK_CACHE_SECONDS = max(
0, int(os.getenv("LOGIN_DESKTOP_NETWORK_CACHE_SECONDS", "30"))
)
GENERIC_WWW_NAMES = {
"",
"我的",
@@ -51,6 +63,64 @@ GENERIC_WWW_NAMES = {
}
class LoginNetworkError(RuntimeError):
"""Raised when the login browser cannot reach Douyin."""
def __init__(self, message, *, checks=None):
super().__init__(message)
self.checks = checks or {}
def _safe_proxy_label(proxy_server):
if not proxy_server:
return ""
try:
parsed = urlsplit(proxy_server)
host = parsed.hostname or ""
port = f":{parsed.port}" if parsed.port else ""
return f"{host}{port}" if host else "configured proxy"
except ValueError:
return "configured proxy"
def _probe_login_target(proxy_server=None, timeout_seconds=15):
"""Probe Douyin without inheriting the process proxy environment."""
if proxy_server:
handlers = [
urllib.request.ProxyHandler(
{"http": proxy_server, "https": proxy_server}
)
]
else:
handlers = [urllib.request.ProxyHandler({})]
opener = urllib.request.build_opener(*handlers)
request = urllib.request.Request(
REMOTE_LOGIN_URL,
headers={"User-Agent": "DouYinSparkFlow-login-preflight/1"},
)
started = time.monotonic()
try:
with opener.open(request, timeout=timeout_seconds) as response:
response.read(256)
status = int(getattr(response, "status", 200))
if status >= 400:
raise RuntimeError(f"HTTP {status}")
return {
"ok": True,
"status": status,
"latency_ms": round((time.monotonic() - started) * 1000),
}
except Exception as exc:
error = str(exc) or exc.__class__.__name__
if proxy_server:
error = error.replace(proxy_server, _safe_proxy_label(proxy_server))
return {
"ok": False,
"error": error[:240],
"latency_ms": round((time.monotonic() - started) * 1000),
}
class LoginDesktopManager:
def __init__(self):
self._lock = asyncio.Lock()
@@ -63,6 +133,83 @@ class LoginDesktopManager:
self._status_checked_at = 0.0
self._idle_monitor_task = None
self._scheduled_stop_task = None
self._network_route = None
self._network_checks = {}
self._network_checked_at = 0.0
def _network_payload(self):
route = dict(self._network_route or {})
return {
"mode": LOGIN_NETWORK_MODE,
"selected": route.get("mode", ""),
"proxy": _safe_proxy_label(LOGIN_PROXY_SERVER) if route.get("mode") == "proxy" else "",
"checked_at": route.get("checked_at", ""),
"checks": dict(self._network_checks),
}
def _invalidate_network_route(self):
self._network_route = None
self._network_checked_at = 0.0
async def _select_network_route(self, *, force=False):
now = time.monotonic()
if (
not force
and self._network_route
and now - self._network_checked_at < LOGIN_NETWORK_CACHE_SECONDS
):
return dict(self._network_route)
checks = {}
candidates = []
if LOGIN_NETWORK_MODE in {"auto", "direct"}:
candidates.append(("direct", None))
if LOGIN_NETWORK_MODE in {"auto", "proxy"} and LOGIN_PROXY_SERVER:
candidates.append(("proxy", LOGIN_PROXY_SERVER))
for mode, proxy_server in candidates:
result = await asyncio.to_thread(
_probe_login_target,
proxy_server,
LOGIN_PREFLIGHT_TIMEOUT_SECONDS,
)
checks[mode] = result
if result.get("ok"):
route = {
"mode": mode,
"proxy": _safe_proxy_label(proxy_server),
"checked_at": time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()),
}
self._network_checks = checks
self._network_route = route
self._network_checked_at = now
return dict(route)
self._network_checks = checks
self._network_route = None
self._network_checked_at = now
if LOGIN_NETWORK_MODE == "direct":
message = "无法直连抖音创作者中心,请检查服务器网络出口"
elif LOGIN_NETWORK_MODE == "proxy":
message = f"代理 {_safe_proxy_label(LOGIN_PROXY_SERVER)} 无法访问抖音创作者中心"
else:
message = "直连和代理都无法访问抖音创作者中心"
raise LoginNetworkError(message, checks=checks)
async def network_preflight(self, *, force=True):
try:
route = await self._select_network_route(force=force)
return {
"ok": True,
"route": route,
"network": self._network_payload(),
}
except LoginNetworkError as exc:
return {
"ok": False,
"error": str(exc),
"network": self._network_payload(),
}
def mark_activity(self):
self._last_activity = time.monotonic()
@@ -136,6 +283,7 @@ class LoginDesktopManager:
except Exception:
pass
self.playwright = None
route = await self._select_network_route()
self.playwright = await async_playwright().start()
launch_args = [
"--start-maximized",
@@ -158,11 +306,18 @@ class LoginDesktopManager:
"--renderer-process-limit=2",
]
)
launch_options = {
"headless": False,
"viewport": {"width": 1600, "height": 1000},
"args": launch_args,
}
if route["mode"] == "proxy":
launch_options["proxy"] = {"server": LOGIN_PROXY_SERVER}
else:
launch_args.append("--no-proxy-server")
self.context = await self.playwright.chromium.launch_persistent_context(
str(PROFILE_DIR),
headless=False,
viewport={"width": 1600, "height": 1000},
args=launch_args,
**launch_options,
)
self.page = self.context.pages[0] if self.context.pages else await self.context.new_page()
@@ -227,7 +382,12 @@ class LoginDesktopManager:
await page.bring_to_front()
except Exception:
pass
return {"ok": True, "url": page.url, "mode": LOGIN_DESKTOP_MODE}
return {
"ok": True,
"url": page.url,
"mode": LOGIN_DESKTOP_MODE,
"network": self._network_payload(),
}
async def status(self):
now = time.monotonic()
@@ -247,6 +407,7 @@ class LoginDesktopManager:
"unique_id": "",
"current_url": "",
"profile_dir": str(PROFILE_DIR),
"network": self._network_payload(),
}
self._status_cache = payload
self._status_checked_at = now
@@ -272,6 +433,7 @@ class LoginDesktopManager:
"unique_id": "",
"current_url": "",
"profile_dir": str(PROFILE_DIR),
"network": self._network_payload(),
}
self._status_cache = payload
self._status_checked_at = now
@@ -295,6 +457,7 @@ class LoginDesktopManager:
"unique_id": unique_id,
"current_url": current_url,
"profile_dir": str(PROFILE_DIR),
"network": self._network_payload(),
}
self._status_cache = payload
self._status_checked_at = now
@@ -309,7 +472,7 @@ class LoginDesktopManager:
try:
page = await self._get_active_page()
if page.url.startswith(REMOTE_LOGIN_URL):
return {"ok": True, "url": page.url}
return {"ok": True, "url": page.url, "network": self._network_payload()}
refresh_url = f"{REMOTE_LOGIN_URL}?qr_refresh={int(time.time() * 1000)}"
try:
await page.goto(refresh_url, wait_until="commit", timeout=30000)
@@ -318,7 +481,7 @@ class LoginDesktopManager:
await self.start()
page = await self._get_active_page()
await page.goto(refresh_url, wait_until="commit", timeout=30000)
return {"ok": True, "url": page.url}
return {"ok": True, "url": page.url, "network": self._network_payload()}
finally:
self._page_operation_lock.release()
@@ -516,6 +679,14 @@ async def health():
return {"ok": True}
@app.get("/preflight")
async def preflight():
result = await manager.network_preflight(force=True)
if not result.get("ok"):
raise HTTPException(status_code=502, detail=result)
return result
@app.get("/status")
async def status():
return await manager.status()
@@ -523,8 +694,13 @@ async def status():
@app.post("/open-login")
async def open_login():
await manager.open_login()
return {"ok": True}
try:
return await manager.open_login()
except LoginNetworkError as exc:
raise HTTPException(
status_code=502,
detail={"code": "LOGIN_NETWORK_UNAVAILABLE", "message": str(exc), "checks": exc.checks},
) from exc
@app.post("/reset")
@@ -546,7 +722,13 @@ async def focus():
@app.post("/refresh-qr")
async def refresh_qr():
return await manager.refresh_login_qr()
try:
return await manager.refresh_login_qr()
except LoginNetworkError as exc:
raise HTTPException(
status_code=502,
detail={"code": "LOGIN_NETWORK_UNAVAILABLE", "message": str(exc), "checks": exc.checks},
) from exc
@app.post("/export")
@@ -566,7 +748,13 @@ async def export():
async def login_qr():
if manager._page_operation_lock.locked():
raise HTTPException(status_code=503, detail="login page is busy; retry shortly")
page = await manager._get_active_page()
try:
page = await manager._get_active_page()
except LoginNetworkError as exc:
raise HTTPException(
status_code=502,
detail={"code": "LOGIN_NETWORK_UNAVAILABLE", "message": str(exc), "checks": exc.checks},
) from exc
expired = await page.locator('[class*="qrcode_expired"]').count()
if expired and await page.locator('[class*="qrcode_expired"]').first.is_visible():
raise HTTPException(status_code=409, detail="login QR code has expired")
@@ -0,0 +1,32 @@
import unittest
from pathlib import Path
REPO_ROOT = Path(__file__).resolve().parents[2]
SOURCE_ROOT = REPO_ROOT / "DouYinSparkFlow"
class NetworkFallbackContractTests(unittest.TestCase):
def test_browser_exposes_direct_first_routes_and_preflight(self):
browser = (SOURCE_ROOT / "core" / "browser.py").read_text(encoding="utf-8")
self.assertIn("def douyin_network_modes", browser)
self.assertIn("return (\"direct\", \"mihomo\")", browser)
self.assertIn("async def select_douyin_network_mode", browser)
self.assertIn("get_browser(network_mode=network_mode)", browser)
def test_friend_refresh_can_try_the_next_route(self):
friends = (SOURCE_ROOT / "core" / "friends.py").read_text(encoding="utf-8")
self.assertIn("for index, network_mode in enumerate(modes)", friends)
self.assertIn("returned zero friends; trying next route", friends)
self.assertIn("get_browser(GUI=False, network_mode=network_mode)", friends)
def test_tasks_select_route_before_browser_creation(self):
tasks = (SOURCE_ROOT / "core" / "tasks.py").read_text(encoding="utf-8")
self.assertIn("select_douyin_network_mode(CREATOR_HOME_URL)", tasks)
self.assertIn("get_browser(network_mode=network_mode)", tasks)
self.assertIn("network_mode=network_mode", tasks)
def test_login_defaults_to_auto_in_compose(self):
compose = (REPO_ROOT / "docker-compose.yml").read_text(encoding="utf-8")
self.assertIn("LOGIN_DESKTOP_PROXY_MODE: ${LOGIN_DESKTOP_PROXY_MODE:-auto}", compose)
if __name__ == "__main__":
unittest.main()