diff --git a/.workbuddy/memory/2026-08-28.md b/.workbuddy/memory/2026-08-28.md new file mode 100644 index 0000000..9c723f8 --- /dev/null +++ b/.workbuddy/memory/2026-08-28.md @@ -0,0 +1,21 @@ +# 2026-08-28 工作日志 + +## 抖音 IM「decision=KICK」排查(服务器 116.62.23.103) + +**结论**:不是抖音改了 IM 规则/签名算法,是账号被安全网关风控踢下线。签名链路正常(同服务器另一账号可正常发送、读接口正常返回、凭证完整)。 + +**关键证据**(服务器 MySQL `kefu` 库 + `/www/wwwlogs/python/douyin/error.log`): +- 活跃账号:id=11「随安尔乐」my_uid=7670159096859706425、id=12「抖音账号_12」my_uid=7670157997767050299,均为 Chrome/148 UA,keys/web_protect 凭证完整(len 533/459)。 +- 两个账号反复 KICK,且都在回复同一测试号「Huhao」(peer_uid=66578464308),回复内容只是 "ss"/"jjj" 测试文本(非导流话术)。 +- 时序:KICK → create_conversation INVALID_REQUEST → 「用户未登录」→ 自动重登录 → 恢复 → 再发 → 再 KICK,形成死循环(约每 10~30 分钟一次)。 +- 读接口(get_by_user_init)正常,只有 signed 的发送(message/send)被 KICK → 签名没问题,是账号级风控。 + +**根因判断**:抖音 2026 风控收紧(内容+设备+IP+行为+账号五维)。触发点最可能是:多账号同服务器 IP + 对同一陌生 peer 的自动回复行为;且 KICK→重登→再发 的紧循环本身会加重风控。 + +**处理方向**: +1. 停掉受影响账号的托管/自动回复,冷却几小时~一天。 +2. 浏览器模式手动重登,先用互关好友或「对方先发」的真实用户测发送,别再用小号对冷门 peer 反复自动回。 +3. 确认账号没在手机端同时登录(并发登录会吊销 web ticket)。 +4. 后续可选代码改进:KICK 后对同一会话加冷却(重登后暂不重发),打断紧循环。 + +**环境备注**:后端用 MySQL(`KEFU_DB_TYPE=mysql`),`kefu.db`/`kefu1.db` 是遗留 SQLite(kefu.db 已损坏,非活跃库,可忽略)。 diff --git a/backend/kefu.db-shm b/backend/kefu.db-shm index 9b62e0d..e22d7bf 100644 Binary files a/backend/kefu.db-shm and b/backend/kefu.db-shm differ diff --git a/backend/kefu.db-wal b/backend/kefu.db-wal index 55a5ef1..570a865 100644 Binary files a/backend/kefu.db-wal and b/backend/kefu.db-wal differ diff --git a/backend/main.py b/backend/main.py index c6974ba..12ffee4 100644 --- a/backend/main.py +++ b/backend/main.py @@ -84,6 +84,7 @@ from rpa_engine.credential import ( assess_account_credential, build_im_session_from_storage, build_cookie_credential_detail, + credential_egress_mismatch, ) from utils.cookie_store import ( write_cookie_file, @@ -1864,6 +1865,7 @@ async def update_account( if body.user_agent is not None: ua = (body.user_agent or "").strip() account.user_agent = ua or None + egress_changed = False if "egress_public_ip" in body.model_fields_set: selected_public_ip = str(body.egress_public_ip or "").strip() if selected_public_ip: @@ -1873,12 +1875,27 @@ async def update_account( raise HTTPException(status_code=400, detail="公网通道必须是有效的 IPv4 地址") if parsed_ip.version != 4: raise HTTPException(status_code=400, detail="公网通道目前仅支持 IPv4") + previous_public_ip = str(account.egress_public_ip or "").strip() + egress_changed = previous_public_ip != selected_public_ip account.egress_public_ip = selected_public_ip or None if body.egress_auto_attempts is not None: account.egress_auto_attempts = clamp_attempts(body.egress_auto_attempts) account.updated_at = datetime.utcnow() await db.commit() await db.refresh(account) + if egress_changed and manager.is_running(account_id): + await manager.stop_worker(account_id) + await db.execute( + update(Account).where(Account.id == account_id).values( + status="offline", + error_message=( + "公网通道已变更,请重新启动托管以使用新通道;" + "已保留登录凭证,校验通过后无需重新扫码" + ), + ) + ) + await db.commit() + await db.refresh(account) if follow_config_changed: worker = manager.workers.get(account_id) invalidate = getattr(worker, "invalidate_follow_welcome_config", None) @@ -1888,8 +1905,9 @@ async def update_account( runtime_service = getattr(worker, "_im_service", None) if worker else None if runtime_service: runtime_session = runtime_service.session - runtime_session.egress_public_ip = str(account.egress_public_ip or "").strip() - runtime_session.egress_source_ip = "" + # Reconnect after a public-IP change so HTTP and the existing WS do + # not use different routes. Reconnecting does not invalidate cookies. + # Only the retry-count can be hot-updated without reconnecting. runtime_session.egress_auto_attempts = clamp_attempts(account.egress_auto_attempts) return _build_account_response(account) @@ -2350,13 +2368,26 @@ async def _start_account_rpa_impl( # Credential assessment issues real network requests to Douyin, and a bulk # start runs it for every queued account. Release the connection first. await _release_db_connection(db) + reset_performed = False + + selected_public_ip = str(getattr(account, "egress_public_ip", "") or "").strip() + if credential_egress_mismatch(cookie_data, selected_public_ip): + # The stored IP is local metadata, not a platform authentication + # verdict. Imported/legacy cookies may not have it at all. Keep the + # credentials and use normal validation on the selected route. + logger.info( + "Account %s egress marker differs; preserving credentials and " + "validating on selected channel %s", + account_id, + selected_public_ip or "default", + ) assessment = await assess_account_credential( cookie_data, account.im_session_data, startup_priority=True, + egress_public_ip=selected_public_ip, ) login_mode = requested_login_mode or assessment["login_mode"] - reset_performed = False if assessment.get("should_reset") and login_mode != "im_direct": account = await _reset_account_credentials(account_id, db) diff --git a/backend/rpa_engine/credential.py b/backend/rpa_engine/credential.py index 92fc883..b77d3f8 100644 --- a/backend/rpa_engine/credential.py +++ b/backend/rpa_engine/credential.py @@ -10,6 +10,33 @@ from utils.cookie_store import analyze_cookie logger = logging.getLogger("credential") +CREDENTIAL_EGRESS_PUBLIC_IP_KEY = "credential_egress_public_ip" + + +def credential_egress_mismatch( + cookie_data: Optional[str], + selected_public_ip: str = "", +) -> bool: + """Compare historical browser egress metadata for diagnostics only. + + This is not an authentication check: a different or missing local marker + cannot prove that cookies are invalid. Callers must keep the credentials + and use normal validation instead of forcing a reset or browser login. + """ + if not cookie_data: + return False + try: + storage = json.loads(cookie_data) + except (TypeError, ValueError): + return False + if not isinstance(storage, dict): + return False + selected = str(selected_public_ip or "").strip() + if CREDENTIAL_EGRESS_PUBLIC_IP_KEY not in storage: + return bool(selected) + stored = str(storage.get(CREDENTIAL_EGRESS_PUBLIC_IP_KEY) or "").strip() + return stored != selected + def _should_reset_credentials(assessment: dict) -> bool: """凭证全面失效时需清空 Cookie/IM 数据并重新登录。""" @@ -202,6 +229,7 @@ async def assess_account_credential( im_session_data: Optional[str] = None, *, startup_priority: bool = False, + egress_public_ip: str = "", ) -> dict: cookie_info = analyze_cookie(cookie_data) result = { @@ -227,6 +255,18 @@ async def assess_account_credential( return result session = build_im_session_from_storage(storage, im_session_data) + selected_public_ip = str(egress_public_ip or "").strip() + if selected_public_ip: + try: + from rpa_engine.egress_channels import resolve_fixed_channel + + route = await resolve_fixed_channel(selected_public_ip) + session.egress_public_ip = selected_public_ip + session.egress_source_ip = str(route.source_ip or "") + except Exception as exc: + result["message"] = f"指定公网通道 {selected_public_ip} 当前不可用:{exc}" + result["login_mode"] = "browser" + return result result["has_sessionid"] = has_im_session_token(session) if not cookie_info.get("cookie_valid"): diff --git a/backend/rpa_engine/douyin_im/auth.py b/backend/rpa_engine/douyin_im/auth.py index 734df1b..92d54ca 100644 --- a/backend/rpa_engine/douyin_im/auth.py +++ b/backend/rpa_engine/douyin_im/auth.py @@ -129,6 +129,7 @@ class DouyinAuth: auth.device_id = resolve_proto_device_id( session.device_id, session.web_id, session.my_uid ) + auth.source_ip = str(getattr(session, "egress_source_ip", "") or "") # web_protect 缺 client_cert 时,才用 frontier 抓包证书兜底(不覆盖 ts_sign) if not auth.client_cert and getattr(session, "sdk_cert", ""): auth.client_cert = normalize_client_cert(session.sdk_cert) diff --git a/backend/rpa_engine/douyin_im/conv_util.py b/backend/rpa_engine/douyin_im/conv_util.py index 16e8aba..f736038 100644 --- a/backend/rpa_engine/douyin_im/conv_util.py +++ b/backend/rpa_engine/douyin_im/conv_util.py @@ -45,3 +45,26 @@ def normalize_conversation_id(conversation_id: str, my_uid: int) -> str: if peer_uid and my_uid: return build_conversation_id(my_uid, peer_uid) return (conversation_id or "").strip() + + +def conversation_belongs_to(conversation_id: str, my_uid: int) -> bool: + """判断单聊会话是否属于 my_uid 本人。 + + 托管多个账号时,一条属于别的账号的会话(例如 frontier 长连接按设备号寻址 + 造成的跨账号推送)一旦流进本账号的处理链路,resolve_peer_uid 会把末段当成 + 「对方」、normalize_conversation_id 再拼成 0:1:{本账号}:{别人的好友},于是 + 本账号就把消息发给了另一个账号的好友。这里给出唯一的归属判据。 + + 无法判定时一律返回 True(保守放行):缺 my_uid、群聊、裸 UID 等形态本来就 + 不带参与方信息。只有两个参与方都已知、且都不是本账号时才判定为不属于本账号。 + """ + try: + uid = int(my_uid or 0) + except (TypeError, ValueError): + return True + if not uid: + return True + parts = parse_conversation_parts(conversation_id) + if not parts: + return True + return uid in parts diff --git a/backend/rpa_engine/douyin_im/frontier.py b/backend/rpa_engine/douyin_im/frontier.py index 1e8a971..6b97cda 100644 --- a/backend/rpa_engine/douyin_im/frontier.py +++ b/backend/rpa_engine/douyin_im/frontier.py @@ -102,11 +102,16 @@ def resolve_frontier_device_id(session: DouyinImSession) -> str: return "" -def _ws_device_id(url: str) -> str: +def ws_device_id(url: str) -> str: + """frontier 推送的寻址键:设备号(不是账号 UID)。""" m = re.search(r"[?&]device_id=([^&\s]+)", url or "") return unquote(m.group(1)) if m else "" +# 兼容内部旧引用 +_ws_device_id = ws_device_id + + def _ws_device_matches_session(session: DouyinImSession, url: str) -> bool: ws_dev = _ws_device_id(url) if not ws_dev or not ws_dev.isdigit(): diff --git a/backend/rpa_engine/douyin_im/http_client.py b/backend/rpa_engine/douyin_im/http_client.py index 8ed6b90..bd47de6 100644 --- a/backend/rpa_engine/douyin_im/http_client.py +++ b/backend/rpa_engine/douyin_im/http_client.py @@ -14,7 +14,12 @@ from rpa_engine.egress_channels import ( resolve_fixed_channel, resolve_send_channels, ) -from .conv_util import build_conversation_id, normalize_conversation_id, resolve_peer_uid +from .conv_util import ( + build_conversation_id, + conversation_belongs_to, + normalize_conversation_id, + resolve_peer_uid, +) from .message_content import format_im_message, serialize_message_content from .peer_profile import enrich_conversation_item, fetch_peer_profile, is_generic_peer_name from .protocol import normalize_im_payload_from_bytes, _pick_avatar_url @@ -1144,8 +1149,17 @@ class DouyinImHttpClient: self.session.my_uid, uid, ) + previous = int(self.session.my_uid or 0) self.session.my_uid = uid self.session.uid_verified = True + # 托管注册表按 UID 记录「本系统正在托管谁」。纠正后必须迁移,否则回环 + # 防护会认错人:旧 UID 永远留在表里,真实 UID 从未登记。只迁移确实已登记 + # 的托管身份,避免 API 侧的临时客户端把自己也登记进去。 + from . import hosted_registry + + if hosted_registry.is_hosted(previous): + hosted_registry.unregister(previous) + hosted_registry.register(uid) async def get_conversations( self, @@ -1383,6 +1397,26 @@ class DouyinImHttpClient: self._set_error("无法获取当前账号 UID") self._log_send_failure(conversation_id, "无法获取当前账号 UID(Cookie 可能已失效)") return False + + # 跨账号写入闸门:normalize_conversation_id 会把任何会话 ID 改写成 + # 0:1:{本账号}:{末段 UID},所以一条属于别的账号的会话流到这里会被 + # 静默改写并发给对方的好友。发送前先确认本账号确实是该会话的参与方。 + if not conversation_belongs_to(conversation_id, my_uid): + detail = ( + f"会话 {conversation_id} 的参与方都不是本账号(uid={my_uid})," + "拒绝发送:这条会话属于另一个账号,继续发送会把消息发给别人的好友。" + ) + self._set_error(detail) + self.last_send_channel_retryable = False + self._log_send_failure(conversation_id, detail) + logger.error( + "Account %s refused cross-account send to %s (my_uid=%s)", + self.account_id, + conversation_id, + my_uid, + ) + return False + if not auth.is_sign_ready(): self._set_error("缺少 IM 签名密钥,请用浏览器登录补全 localStorage") self._log_send_failure( @@ -1510,10 +1544,14 @@ class DouyinImHttpClient: decision = str(result.get("decision") or "").strip().upper() if decision == "KICK": - self.last_send_channel_retryable = True + # KICK is a terminal, account-session decision. Retrying the + # same authenticated write from another source address cannot + # repair the session and only adds another high-risk request. + self.last_send_needs_refresh = False + self.last_send_channel_retryable = False detail = ( "抖音安全网关返回 decision=KICK,当前登录/安全会话已被服务端踢下线;" - "系统正在自动重登录,请留意账号卡片上的二维码并扫码" + "已停止本次发送及公网通道重试,系统正在自动重登录,请留意账号卡片上的二维码并扫码" ) elif decision: detail = f"抖音安全网关拒绝发送 decision={decision}" @@ -1527,7 +1565,11 @@ class DouyinImHttpClient: hint = _BUSINESS_REJECT_FALLBACK # 7911 属于“签名凭证失效/安全校验未过”,标记为可刷新后重试 self.last_send_needs_refresh = status_code in _CREDENTIAL_EXPIRED_CODES - self.last_send_channel_retryable = self.last_send_needs_refresh + # 7911 is a credential/signature problem. It may be retried + # once only after refreshing the credentials on the same + # session; switching egress mid-session makes the fingerprint + # less consistent and must not be used as the recovery path. + self.last_send_channel_retryable = False detail = f"抖音拒绝投递 status_code={status_code}" if status_reason: detail += f";抖音提示:{status_reason}" @@ -1550,7 +1592,10 @@ class DouyinImHttpClient: detail = ";".join(reason_bits) or "接口返回但未确认投递(无 server_message_id)" if "INVALID_REQUEST" in detail.upper(): - self.last_send_channel_retryable = True + # INVALID_REQUEST is a protocol/session rejection, not a + # transport failure. A second public IP sends the same invalid + # request and can invalidate an otherwise recoverable login. + self.last_send_channel_retryable = False full_detail = f"{detail};{target};resp[{result.get('summary')}]" if self.last_request_debug: full_detail += f"\n--- 请求详情 ---\n{self.last_request_debug}" diff --git a/backend/rpa_engine/douyin_im/service.py b/backend/rpa_engine/douyin_im/service.py index 870b572..738388c 100644 --- a/backend/rpa_engine/douyin_im/service.py +++ b/backend/rpa_engine/douyin_im/service.py @@ -17,7 +17,8 @@ from .reply_queue import AccountReplyQueue from .traffic_control import get_traffic_controller from .reply_payload import format_reply_display, serialize_reply_log -from .conv_util import resolve_peer_uid +from . import hosted_registry +from .conv_util import conversation_belongs_to, resolve_peer_uid from .peer_profile import ( enrich_conversation_item, fetch_peer_profile, @@ -338,6 +339,10 @@ class DouyinImService: self._ready_notified = False self._session_invalid_strikes = 0 self._session_invalid_fired = False + # A keepalive browser may refresh cookies/security material while an + # outbound reply is being prepared. Serialize the short credential + # hand-off with sends so one request never mixes old and new state. + self._session_lock = asyncio.Lock() self.reply_delay_seconds = max(0, int(reply_delay_seconds or 0)) # 实时解析账号排队间隔:账号专属优先,否则使用系统默认值。 self._reply_delay_resolver = reply_delay_resolver @@ -355,15 +360,19 @@ class DouyinImService: self._cooldown_resolver = cooldown_resolver # 由 worker 注入:触发后台重新采集 web_protect/keys(刷新 ts_sign),返回是否刷新成功 self.refresh_credentials = refresh_credentials - # 由 worker 注入的第二套发送方案:当 HTTP 签名发送被安全网关拒绝 - # (decision=KICK / 7911 / INVALID_REQUEST)时,用浏览器页面上下文 - # 重新发送(真实 JS 生成 a_bogus/bd-ticket-guard,可自愈被踢的会话)。 + # 由 worker 注入的第二套发送方案:仅当 HTTP 返回非终态的 7911 + # 签名错误时,可在同一账号/同一出口的浏览器页面上下文重试一次。 + # KICK 与 INVALID_REQUEST 不得重放,避免在已失效会话上继续写请求。 # 签名: async (conversation_id, content) -> (ok, detail) self.send_fallback = send_fallback self._running = False self._replied_keys: set[str] = set() self._logged_keys: set[str] = set() self._received_logged_keys: set[str] = set() + # 已告警过的「不属于本账号」的会话,避免同一条串号会话刷屏 + self._foreign_conv_logged: set[str] = set() + # 已告警过的「对方也是本系统托管账号」的 peer,避免同一对账号刷屏 + self._hosted_peer_logged: set[str] = set() # 每个对话/用户最近一次自动回复的时间戳(monotonic 秒),用于冷却窗口去重 self._last_reply_at: dict[str, float] = {} self._conv_previews: dict[str, str] = {} @@ -510,6 +519,38 @@ class DouyinImService: return f"用户{sender_uid[-6:]}" if len(sender_uid) > 6 else f"用户{sender_uid}" return "未知用户" + def _conversation_is_mine(self, conv_id: str) -> bool: + """本账号是否为该单聊会话的参与方;不是就丢弃,绝不改写后发送。""" + my_uid = int(self.session.my_uid or 0) + if conversation_belongs_to(conv_id, my_uid): + return True + conv_key = str(conv_id or "") + logger.warning( + "Account %s dropped a message from foreign conversation %s " + "(my_uid=%s); two accounts most likely share one set of credentials", + self.account_id, + conv_key, + my_uid, + ) + if conv_key not in self._foreign_conv_logged: + if len(self._foreign_conv_logged) > 200: + self._foreign_conv_logged.clear() + self._foreign_conv_logged.add(conv_key) + system_logger.record( + "已丢弃不属于本账号的私信", + detail=( + f"会话 {conv_key} 的参与方都不是本账号(uid={my_uid})," + "该消息属于另一个账号,已丢弃且不会自动回复。" + "常见原因:多个账号的凭证来自同一台机器/同一个浏览器," + "frontier 长连接按设备号寻址导致两个账号互相收到对方的私信。" + "请为每个账号单独采集凭证(独立浏览器配置/设备)。" + ), + level="warning", + category="recv", + account_id=self.account_id, + ) + return False + def _is_self_message(self, msg: dict) -> bool: sender_uid = str(msg.get("sender_uid") or "").strip() if not sender_uid or not self.session.my_uid: @@ -637,10 +678,18 @@ class DouyinImService: self, msg: dict, ) -> Optional[Callable[[], Awaitable[None]]]: + conv_id = msg.get("conversation_id") or "" + # 跨账号隔离:只处理本账号自己的会话。frontier 按设备号寻址推送, + # 同一台机器/同一浏览器采集出来的多个账号 device_id 可能相同,两条长连接 + # 会订阅到同一个地址并互相收到对方的私信。若不在这里拦住, + # normalize_conversation_id 会把别人的会话改写成 + # 0:1:{本账号}:{别人的好友},本账号就把自动回复发给了另一个账号的好友。 + if not self._conversation_is_mine(conv_id): + return + if self._is_self_message(msg): return - conv_id = msg.get("conversation_id") or "" sender_uid = str(msg.get("sender_uid") or "") sender = self._resolve_sender_name(msg) sender_avatar = str(msg.get("sender_avatar") or "").strip() @@ -769,6 +818,44 @@ class DouyinImService: # 防止延迟排队期间被重复加入发送队列。 self._replied_keys.add(key) + # 对方也是本系统托管的账号:双方都会自动回复,一来一回就是无限回环。 + # 这种高频互发是触发抖音风控(7911)/业务拒绝(8004)的常见根因,因此消息 + # 照常记录,但不再自动回复。需要回复请用消息页手动发送。 + if peer_uid and hosted_registry.is_hosted(peer_uid): + await self.log_fn( + **log_kwargs, + reply=None, + status="ignored", + error=( + f"对方(UID {peer_uid})也是本系统托管中的账号," + "自动回复会在两个账号之间形成无限回环并触发抖音风控,已跳过;" + "如需回复请在消息页手动发送" + ), + ) + if content: + self._conv_previews[sender] = content + if peer_uid not in self._hosted_peer_logged: + if len(self._hosted_peer_logged) > 200: + self._hosted_peer_logged.clear() + self._hosted_peer_logged.add(peer_uid) + logger.info( + "Account %s skipped auto-reply to hosted account %s", + self.account_id, + peer_uid, + ) + system_logger.record( + "自动回复已跳过(对方也是托管账号)", + detail=( + f"{sender}(UID {peer_uid})是本系统托管中的另一个账号。" + "两个托管账号互相自动回复会形成无限回环," + "属于抖音风控(7911/8004)的高发场景,因此只记录消息、不自动回复。" + ), + level="warning", + category="send", + account_id=self.account_id, + ) + return + # 同账号、同会话只保留一个尚未发送的回复任务。后续来信只追加到 # 原任务详情,不改变它的发送时间、位置或已经匹配好的回复。 queue_merge_keys = self._reply_queue_merge_keys(conv_id, peer_uid) @@ -1428,11 +1515,60 @@ class DouyinImService: """把指定自动回复任务移入账号紧急队列;实际发送仍由单消费者串行执行。""" return await self._reply_queue.send_now(job_id) + async def replace_session(self, fresh: DouyinImSession) -> None: + """Atomically install a freshly harvested login/security session. + + The running WebSocket can keep its current connection, but future + reconnects and every HTTP send must see the same refreshed object. + Account egress selection lives outside persisted IM credentials, so it + is deliberately carried over from the current runtime session. + """ + async with self._session_lock: + current = self.session + current_uid = int(getattr(current, "my_uid", 0) or 0) + fresh_uid = int(getattr(fresh, "my_uid", 0) or 0) + if current_uid and fresh_uid and current_uid != fresh_uid: + raise ValueError( + f"refusing cross-account session refresh: {current_uid} != {fresh_uid}" + ) + + fresh.conv_meta = { + **dict(getattr(current, "conv_meta", {}) or {}), + **dict(getattr(fresh, "conv_meta", {}) or {}), + } + if not fresh.ws_urls: + fresh.ws_urls = list(getattr(current, "ws_urls", []) or []) + fresh.egress_public_ip = str( + getattr(current, "egress_public_ip", "") or "" + ) + fresh.egress_source_ip = str( + getattr(current, "egress_source_ip", "") or "" + ) + fresh.egress_auto_attempts = int( + getattr(current, "egress_auto_attempts", 1) or 1 + ) + self.session = fresh + if self._ws_client is not None: + self._ws_client.session = fresh + async def _send_text( self, conversation_id: str, content: str, conversation_short_id: str = "", + ) -> tuple[bool, Optional[dict]]: + async with self._session_lock: + return await self._send_text_unlocked( + conversation_id, + content, + conversation_short_id=conversation_short_id, + ) + + async def _send_text_unlocked( + self, + conversation_id: str, + content: str, + conversation_short_id: str = "", ) -> tuple[bool, Optional[dict]]: """发送一条私信;若因签名凭证失效(7911)失败,刷新 web_protect 后自动重试一次。 @@ -1468,16 +1604,14 @@ class DouyinImService: continue break - # 第二套发送方案(浏览器页面内发送): - # HTTP 签名发送被安全网关拒绝(KICK/7911/INVALID_REQUEST)时,交给 worker - # 用浏览器页面上下文重发——由抖音页面自带的 security-sdk 在真实环境生成 - # a_bogus/bd-ticket-guard,绕开我们 Node execjs 的签名模拟,可自愈被踢会话。 + # 第二套发送方案(浏览器页面内发送):仅处理非终态 7911。 + # KICK/INVALID_REQUEST 会停止发送并进入下线处理,不在失效会话上重放。 upper_err = (self.last_error or "").upper() - if self.send_fallback and ( - "DECISION=KICK" in upper_err - or "STATUS_CODE=7911" in upper_err - or "INVALID_REQUEST" in upper_err - ): + # KICK already invalidated the login and INVALID_REQUEST is a + # protocol/session rejection. Replaying either through a browser + # fetch cannot heal it and creates another risky write. 7911 is the + # only non-terminal signing failure eligible for the browser fallback. + if self.send_fallback and "STATUS_CODE=7911" in upper_err: try: fb_ok, fb_detail = await self.send_fallback(conversation_id, content) except Exception as exc: diff --git a/backend/rpa_engine/douyin_im/ws_client.py b/backend/rpa_engine/douyin_im/ws_client.py index fca9711..4b4519b 100644 --- a/backend/rpa_engine/douyin_im/ws_client.py +++ b/backend/rpa_engine/douyin_im/ws_client.py @@ -103,6 +103,12 @@ _LOOP_STATES: "weakref.WeakKeyDictionary[asyncio.AbstractEventLoop, _LoopWsState ) +# frontier 按 device_id 寻址推送:两个托管账号共用同一个设备号时,两条长连接会 +# 订阅到同一个地址并互相收到对方的私信。真正的拦截在 service 的会话归属校验里, +# 这里只负责把「为什么会串号」明确告诉用户。持弱引用,账号停管后自动失效。 +_FRONTIER_DEVICE_OWNERS: "dict[str, weakref.ref[DouyinImWsClient]]" = {} + + def _get_loop_state() -> _LoopWsState: loop = asyncio.get_running_loop() state = _LOOP_STATES.get(loop) @@ -155,6 +161,8 @@ class DouyinImWsClient: self._dispatcher_task: Optional[asyncio.Task] = None self._received_frame_count = 0 self._heartbeat_ack_logged = False + self._frontier_device_id = "" + self._blocked_device_owner_id: Optional[int] = None async def start(self): if self._task and not self._task.done(): @@ -210,6 +218,7 @@ class DouyinImWsClient: if self._task is task: self._task = None self._connection = None + self._release_frontier_device() await self._stop_dispatcher() def _record_connection_system_event( @@ -247,6 +256,7 @@ class DouyinImWsClient: account_key = int(self.account_id or 0) state.system_log_last_at.pop((account_key, "connected"), None) state.system_log_last_at.pop((account_key, "retry"), None) + state.system_log_last_at.pop((account_key, "device_taken"), None) def _ensure_dispatcher(self) -> None: if self._dispatcher_task and not self._dispatcher_task.done(): @@ -311,8 +321,15 @@ class DouyinImWsClient: first_attempt = False if not connect_url: raise RuntimeError("frontier WebSocket URL is unavailable") - logger.info("Connecting IM WebSocket: %s...", connect_url[:100]) - await self._run_connection(connect_url) + if self._claim_frontier_device(connect_url): + logger.info("Connecting IM WebSocket: %s...", connect_url[:100]) + await self._run_connection(connect_url) + else: + # 设备号已被另一个在跑的账号占用:绝不并连同一个推送地址, + # 本账号本轮退回 HTTP 轮询兜底(connected 保持 False, + # service 会自动切到更快的会话对账节奏),并在退避后重试, + # 等占用方停管时自动接管。 + self._report_frontier_device_taken(connect_url) except asyncio.CancelledError: break except Exception as exc: @@ -348,6 +365,76 @@ class DouyinImWsClient: except asyncio.CancelledError: break + self._release_frontier_device() + + def _frontier_device_owner(self, device_id: str) -> "Optional[DouyinImWsClient]": + """当前仍活着的设备号占用方(run 循环任务还在跑才算数)。""" + reference = _FRONTIER_DEVICE_OWNERS.get(device_id) + owner = reference() if reference is not None else None + if owner is None or owner is self: + return None + task = owner._task + if not owner._running or task is None or task.done(): + return None + return owner + + def _claim_frontier_device(self, url: str) -> bool: + """独占本账号的 frontier 设备地址;已被别的账号占用时返回 False。 + + frontier 按 device_id 寻址推送。两个账号共用同一个设备号时,同时建连 + 会让两条连接互相收到对方的私信(串号的根因),且抖音也可能只保留最后 + 一条连接、把先连上的那个账号踢成「连着但收不到」。所以同一个设备地址 + 永远只允许一个账号建连,另一个账号走 HTTP 轮询兜底。 + """ + from .frontier import ws_device_id + + device_id = ws_device_id(url) + if not device_id: + # 判不出设备号(自建地址/异常格式)时不阻断连接,交给会话归属校验兜底。 + return True + owner = self._frontier_device_owner(device_id) + if owner is not None and int(owner.account_id or 0) != int(self.account_id or 0): + self._blocked_device_owner_id = owner.account_id + return False + _FRONTIER_DEVICE_OWNERS[device_id] = weakref.ref(self) + self._frontier_device_id = device_id + self._blocked_device_owner_id = None + return True + + def _report_frontier_device_taken(self, url: str) -> None: + from .frontier import ws_device_id + + device_id = ws_device_id(url) + owner_id = self._blocked_device_owner_id + logger.error( + "Account %s cannot open frontier device_id %s: already held by " + "account %s; falling back to HTTP polling this round", + self.account_id, + device_id, + owner_id, + ) + self._record_connection_system_event( + "device_taken", + "实时接收已让出:与另一个账号共用长连接设备号", + detail=( + f"本账号与账号 {owner_id} 的 frontier 设备号相同(device_id={device_id})。" + "同一个设备地址只允许一个账号建立长连接,否则两个账号会互相收到对方的" + "私信。本账号本轮不建连,改由 HTTP 会话轮询接收(有几十秒级延迟)," + "并在对方停止托管后自动接管。" + "根治办法:为每个账号在独立的浏览器配置/设备上重新采集凭证。" + ), + level="error", + ) + + def _release_frontier_device(self) -> None: + device_id = self._frontier_device_id + self._frontier_device_id = "" + if not device_id: + return + reference = _FRONTIER_DEVICE_OWNERS.get(device_id) + if reference is not None and reference() is self: + _FRONTIER_DEVICE_OWNERS.pop(device_id, None) + def _connection_headers(self) -> list[tuple[str, str]]: headers = [ ("Pragma", "no-cache"), diff --git a/backend/rpa_engine/playwright_worker.py b/backend/rpa_engine/playwright_worker.py index ae6b387..6570144 100644 --- a/backend/rpa_engine/playwright_worker.py +++ b/backend/rpa_engine/playwright_worker.py @@ -28,7 +28,11 @@ from rpa_engine.douyin_im.session import DouyinImSession from rpa_engine.douyin_im.frontier import ensure_frontier_ws from rpa_engine.douyin_im.http_client import DouyinImHttpClient from rpa_engine.douyin_im.traffic_control import get_traffic_controller -from rpa_engine.credential import validate_im_session, build_im_session_from_storage +from rpa_engine.credential import ( + CREDENTIAL_EGRESS_PUBLIC_IP_KEY, + validate_im_session, + build_im_session_from_storage, +) from rpa_engine.device_profiles import resolve_user_agent from rpa_engine.runtime_config import ( resolve_headless, @@ -36,6 +40,7 @@ from rpa_engine.runtime_config import ( playwright_proxy, ) from rpa_engine.egress_channels import clamp_attempts, resolve_fixed_channel +from rpa_engine.source_bound_proxy import playwright_proxy_for_source logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s") logger = logging.getLogger("rpa_engine") @@ -55,7 +60,12 @@ async def _start_playwright_for_browser(headless: Optional[bool] = None): return await async_playwright().start(), headless -async def _launch_chromium(pw, args: list[str], headless: Optional[bool] = None): +async def _launch_chromium( + pw, + args: list[str], + headless: Optional[bool] = None, + source_ip: str = "", +): """统一的 Chromium 启动入口:自动处理无头/有头、虚拟显示与住宅代理。 - headless 默认由 KEFU_BROWSER_HEADLESS 决定(缺省有头,避免抖音安全 SDK 判定)。 @@ -66,10 +76,21 @@ async def _launch_chromium(pw, args: list[str], headless: Optional[bool] = None) headless = resolve_headless(default=False) await ensure_browser_display(headless) launch_kwargs: dict = {"headless": headless, "args": args} - proxy = playwright_proxy() + proxy = ( + await playwright_proxy_for_source(source_ip) + if source_ip + else playwright_proxy() + ) if proxy: launch_kwargs["proxy"] = proxy - logger.info("浏览器将通过代理启动:%s", proxy.get("server")) + if source_ip: + logger.info( + "浏览器将固定使用本机源地址 %s(本地代理 %s)", + source_ip, + proxy.get("server"), + ) + else: + logger.info("浏览器将通过代理启动:%s", proxy.get("server")) return await pw.chromium.launch(**launch_kwargs) @@ -223,6 +244,44 @@ class DouyinWorker: async def get_db(self): return AsyncSessionLocal() + async def _resolve_browser_source_ip(self) -> str: + """Resolve the account's fixed browser egress source address. + + Browser login, keepalive and credential refresh must use the same + public channel as IM HTTP/WS. Silently falling back to the default + route for a selected-but-unavailable channel would create a mixed-IP + security session, so that case intentionally fails closed. + """ + selected = await self._load_selected_egress_public_ip() + if not selected: + return "" + route = await resolve_fixed_channel(selected) + source_ip = str(route.source_ip or "").strip() + if not source_ip: + raise RuntimeError( + f"账号指定公网通道 {selected} 无法绑定到本机网卡,已停止浏览器登录/刷新," + "避免登录出口与发送出口不一致" + ) + return source_ip + + async def _load_selected_egress_public_ip(self) -> str: + db = await self.get_db() + try: + result = await db.execute( + select(Account.egress_public_ip).where(Account.id == self.account_id) + ) + selected = str(result.scalar_one_or_none() or "").strip() + except Exception as exc: + logger.debug( + "Account %s browser egress config unavailable: %s", + self.account_id, + exc, + ) + return "" + finally: + await db.close() + return selected + def _mark_startup_ready(self) -> None: self._startup_error = "" self._startup_ready.set() @@ -640,6 +699,7 @@ class DouyinWorker: "--no-sandbox", "--disable-setuid-sandbox", ], + source_ip=await self._resolve_browser_source_ip(), ) if storage_state: self.context = await self.browser.new_context( @@ -679,10 +739,13 @@ class DouyinWorker: async def _persist_cookies(self): """登录成功或运行中将 Cookie 同步到文件和数据库(合并 HttpOnly sessionid)""" if not self.context: - return + return None storage = await self.context.storage_state() live_cookies = await self.context.cookies() storage = merge_playwright_cookies(storage, live_cookies) + storage[CREDENTIAL_EGRESS_PUBLIC_IP_KEY] = ( + await self._load_selected_egress_public_ip() + ) cookie_json = json.dumps(storage, ensure_ascii=False, indent=2) with open(self.cookie_path, "w", encoding="utf-8") as f: f.write(cookie_json) @@ -704,6 +767,7 @@ class DouyinWorker: await db.rollback() finally: await db.close() + return storage async def _build_im_session_from_storage( self, @@ -1172,8 +1236,8 @@ class DouyinWorker: # 不在发送链路上自动开浏览器刷新:实测重载页面并不会重生 web_protect, # 反而每次失败阻塞 ~22s("反应特别慢"),且无法解决 7911 风控。 refresh_credentials=None, - # 第二套发送方案:HTTP 签名发送被 KICK/7911/INVALID_REQUEST 拒绝时, - # 用浏览器页面上下文重发(真实 JS 签名,可自愈被踢会话)。 + # 第二套发送方案:仅在非终态 7911 时用同出口浏览器页面重试; + # KICK/INVALID_REQUEST 必须停发并下线,不能继续重放。 send_fallback=self.send_im_via_browser_page, ) self._im_service = im_service @@ -1332,7 +1396,10 @@ class DouyinWorker: if sys.platform == "win32": args.append("--start-minimized") browser = await _launch_chromium( - pw, args, headless=browser_headless + pw, + args, + headless=browser_headless, + source_ip=await self._resolve_browser_source_ip(), ) ua = self._user_agent or resolve_user_agent(None) context = await browser.new_context( @@ -1397,13 +1464,26 @@ class DouyinWorker: f"after[{self._fmt_expires_map(after_exp)}] " f"renewed={','.join(renewed) or 'none'}" ) - # 活跃访问后 cookie(msToken 等)可能更新,重新落库 + # 活跃访问后 cookie(msToken 等)可能更新。Cookie、ticket、 + # ts_sign、private key 是一套安全会话,不能只更新数据库里的 + # cookie 而让正在发送的内存会话继续使用旧值;否则 WS 仍能收, + # 下一次写请求却会因新旧凭证混用被安全网关 KICK。 try: await self._persist_cookies() + service = self._im_service + if service is not None: + fresh_session = await self._build_im_session() + await service.replace_session(fresh_session) + await self._persist_im_session(service.session) + logger.info( + "Account %s: keepalive credentials synchronized " + "to active IM session", + self.account_id, + ) except Exception as exc: logger.warning( - f"Account {self.account_id}: keepalive persist cookies " - f"failed: {exc}" + f"Account {self.account_id}: keepalive credential sync " + f"failed; active session left unchanged: {exc}" ) return True, f"已访问 {target_url} 并刷新登录态" except asyncio.CancelledError: @@ -1445,16 +1525,19 @@ class DouyinWorker: ) -> tuple[bool, str]: """第二套发送方案:浏览器页面上下文内重发私信。 - HTTP 签名发送被抖音安全网关拒绝(decision=KICK / 7911 / INVALID_REQUEST) - 时的兜底:用已保存的登录态打开抖音页面,由页面自带 security-sdk 在真实 - 浏览器环境里生成 a_bogus / bd-ticket-guard 并完成发送——绕开 Node execjs - 的签名模拟;浏览器重新加载页面也会重建安全会话,可自愈被服务端踢掉的 - 登录态。仅文本/表情/卡片内容可用,图片需先走 HTTP 上传链路。 + 仅供非终态 7911 签名错误使用:用已保存的登录态打开抖音页面,在与账号 + 相同的固定出口中完成一次页面内发送。KICK/INVALID_REQUEST 不会调用此 + 方法,避免对已经失效的登录态继续重放。仅文本/表情/卡片内容可用,图片 + 需先走 HTTP 上传链路。 返回 (是否成功, 详情)。失败不会抛异常,只记录日志。 """ from rpa_engine.douyin_im.auth import DouyinAuth - from rpa_engine.douyin_im.conv_util import normalize_conversation_id, resolve_peer_uid + from rpa_engine.douyin_im.conv_util import ( + conversation_belongs_to, + normalize_conversation_id, + resolve_peer_uid, + ) from rpa_engine.douyin_im.pb_decode import analyze_send_response from rpa_engine.douyin_im.proto_builder import ProtoBuilder from rpa_engine.douyin_im.reply_payload import build_msg_payload, parse_reply_content @@ -1476,6 +1559,13 @@ class DouyinWorker: if not my_uid: return False, "无法获取 my_uid" + # 与 HTTP 发送同一道跨账号闸门:不是本账号的会话绝不改写后重发。 + if not conversation_belongs_to(conversation_id, my_uid): + return False, ( + f"会话 {conversation_id} 的参与方都不是本账号(uid={my_uid})," + "拒绝发送:这条会话属于另一个账号" + ) + conv_id = normalize_conversation_id(conversation_id, my_uid) peer_uid = resolve_peer_uid(conv_id, my_uid) if not peer_uid: @@ -1545,6 +1635,7 @@ class DouyinWorker: pw, token_args, headless=browser_headless, + source_ip=await self._resolve_browser_source_ip(), ) storage_state = await self._load_storage_state() context = await browser.new_context( @@ -1780,6 +1871,7 @@ class DouyinWorker: pw, token_args, headless=browser_headless, + source_ip=await self._resolve_browser_source_ip(), ) context = await browser.new_context( storage_state=storage_state, @@ -2382,6 +2474,7 @@ class DouyinWorker: self.playwright, args, headless=browser_headless, + source_ip=await self._resolve_browser_source_ip(), ) logger.info(f"Account {self.account_id}: opening browser for IM setup (minimized)") except RuntimeError: diff --git a/backend/rpa_engine/source_bound_proxy.py b/backend/rpa_engine/source_bound_proxy.py new file mode 100644 index 0000000..161e6a1 --- /dev/null +++ b/backend/rpa_engine/source_bound_proxy.py @@ -0,0 +1,233 @@ +"""Loopback HTTP proxy whose outbound sockets bind to one local IPv4. + +Playwright does not expose a ``local_address`` option. Accounts that select a +specific server egress channel therefore use this tiny process-local proxy so +their browser login/refresh traffic leaves through the same interface as IM +HTTP and WebSocket traffic. The listener is loopback-only and does not rotate +or retry public addresses. +""" + +from __future__ import annotations + +import asyncio +import ipaddress +import logging +import socket +import weakref +from urllib.parse import urlsplit + +logger = logging.getLogger("rpa_engine.source_proxy") + +_MAX_HEADER_BYTES = 64 * 1024 +_HEADER_TIMEOUT_SECONDS = 20.0 + + +class SourceBoundProxy: + """Minimal HTTP/HTTPS CONNECT proxy bound to a fixed source address.""" + + def __init__(self, source_ip: str): + address = ipaddress.ip_address(str(source_ip or "").strip()) + if address.version != 4 or address.is_unspecified or address.is_multicast: + raise ValueError(f"invalid IPv4 source address: {source_ip!r}") + self.source_ip = str(address) + self._server: asyncio.AbstractServer | None = None + + @property + def server_url(self) -> str: + if self._server is None or not self._server.sockets: + raise RuntimeError("source-bound proxy has not started") + port = int(self._server.sockets[0].getsockname()[1]) + return f"http://127.0.0.1:{port}" + + async def start(self) -> "SourceBoundProxy": + if self._server is None: + self._server = await asyncio.start_server( + self._handle_client, + host="127.0.0.1", + port=0, + family=socket.AF_INET, + ) + logger.info( + "source-bound browser proxy ready: %s -> source %s", + self.server_url, + self.source_ip, + ) + return self + + async def close(self) -> None: + server = self._server + self._server = None + if server is not None: + server.close() + await server.wait_closed() + + async def _open_upstream( + self, + host: str, + port: int, + ) -> tuple[asyncio.StreamReader, asyncio.StreamWriter]: + return await asyncio.open_connection( + host=host, + port=port, + family=socket.AF_INET, + local_addr=(self.source_ip, 0), + ) + + @staticmethod + async def _relay( + source: asyncio.StreamReader, + destination: asyncio.StreamWriter, + ) -> None: + try: + while True: + chunk = await source.read(64 * 1024) + if not chunk: + break + destination.write(chunk) + await destination.drain() + except (ConnectionError, asyncio.CancelledError): + pass + finally: + try: + destination.write_eof() + except (AttributeError, OSError, RuntimeError): + pass + + @classmethod + async def _bridge( + cls, + client_reader: asyncio.StreamReader, + client_writer: asyncio.StreamWriter, + upstream_reader: asyncio.StreamReader, + upstream_writer: asyncio.StreamWriter, + ) -> None: + tasks = ( + asyncio.create_task(cls._relay(client_reader, upstream_writer)), + asyncio.create_task(cls._relay(upstream_reader, client_writer)), + ) + try: + await asyncio.gather(*tasks) + finally: + for task in tasks: + if not task.done(): + task.cancel() + await asyncio.gather(*tasks, return_exceptions=True) + + @staticmethod + def _parse_authority(authority: str, default_port: int) -> tuple[str, int]: + parsed = urlsplit(f"//{authority}") + host = str(parsed.hostname or "").strip() + if not host: + raise ValueError("proxy request is missing a host") + return host, int(parsed.port or default_port) + + async def _handle_client( + self, + client_reader: asyncio.StreamReader, + client_writer: asyncio.StreamWriter, + ) -> None: + upstream_writer: asyncio.StreamWriter | None = None + try: + header = await asyncio.wait_for( + client_reader.readuntil(b"\r\n\r\n"), + timeout=_HEADER_TIMEOUT_SECONDS, + ) + if len(header) > _MAX_HEADER_BYTES: + raise ValueError("proxy request headers are too large") + lines = header.decode("latin-1").split("\r\n") + request_line = lines[0].split(" ", 2) + if len(request_line) != 3: + raise ValueError("malformed proxy request line") + method, target, version = request_line + + if method.upper() == "CONNECT": + host, port = self._parse_authority(target, 443) + upstream_reader, upstream_writer = await self._open_upstream(host, port) + client_writer.write(b"HTTP/1.1 200 Connection Established\r\n\r\n") + await client_writer.drain() + else: + parsed = urlsplit(target) + host_header = next( + ( + line.partition(":")[2].strip() + for line in lines[1:] + if line.lower().startswith("host:") + ), + "", + ) + authority = parsed.netloc or host_header + host, port = self._parse_authority( + authority, + 443 if parsed.scheme.lower() == "https" else 80, + ) + upstream_reader, upstream_writer = await self._open_upstream(host, port) + origin_target = parsed.path or "/" + if parsed.query: + origin_target += f"?{parsed.query}" + forwarded = [f"{method} {origin_target} {version}"] + forwarded.extend( + line for line in lines[1:] + if line and not line.lower().startswith("proxy-connection:") + ) + upstream_writer.write(("\r\n".join(forwarded) + "\r\n\r\n").encode("latin-1")) + await upstream_writer.drain() + + await self._bridge( + client_reader, + client_writer, + upstream_reader, + upstream_writer, + ) + except asyncio.IncompleteReadError: + pass + except asyncio.CancelledError: + # Event-loop shutdown may cancel an in-flight browser tunnel. + # Closing both writers below is sufficient; do not leak a noisy + # cancelled handler callback into the server log. + pass + except Exception as exc: + logger.warning("source-bound browser proxy request failed: %s", exc) + try: + client_writer.write( + b"HTTP/1.1 502 Bad Gateway\r\nConnection: close\r\n\r\n" + ) + await client_writer.drain() + except (ConnectionError, RuntimeError): + pass + finally: + for writer in (upstream_writer, client_writer): + if writer is None: + continue + try: + writer.close() + await writer.wait_closed() + except (ConnectionError, RuntimeError): + pass + + +class _LoopProxyState: + def __init__(self) -> None: + self.lock = asyncio.Lock() + self.proxies: dict[str, SourceBoundProxy] = {} + + +_loop_states: weakref.WeakKeyDictionary[ + asyncio.AbstractEventLoop, _LoopProxyState +] = weakref.WeakKeyDictionary() + + +async def playwright_proxy_for_source(source_ip: str) -> dict[str, str]: + """Return a Playwright proxy config fixed to ``source_ip``.""" + + loop = asyncio.get_running_loop() + state = _loop_states.get(loop) + if state is None: + state = _LoopProxyState() + _loop_states[loop] = state + normalized = str(ipaddress.ip_address(str(source_ip or "").strip())) + async with state.lock: + proxy = state.proxies.get(normalized) + if proxy is None: + proxy = await SourceBoundProxy(normalized).start() + state.proxies[normalized] = proxy + return {"server": proxy.server_url} diff --git a/backend/tests/test_account_pagination.py b/backend/tests/test_account_pagination.py index a65837d..cb86e49 100644 --- a/backend/tests/test_account_pagination.py +++ b/backend/tests/test_account_pagination.py @@ -112,6 +112,40 @@ class AccountPaginationTests(unittest.IsolatedAsyncioTestCase): finally: main.manager.workers = original_workers + async def test_account_channel_change_stops_running_worker(self): + account = SimpleNamespace( + id=2202, + egress_public_ip="116.62.23.103", + status="online", + ) + db = SimpleNamespace( + commit=AsyncMock(), + refresh=AsyncMock(), + execute=AsyncMock(), + ) + + with ( + patch.object(main, "get_owned_account", AsyncMock(return_value=account)), + patch.object(main.manager, "is_running", return_value=True), + patch.object(main.manager, "stop_worker", AsyncMock(return_value=True)) as stop, + patch.object(main, "_build_account_response", return_value={"id": 2202}), + ): + response = await main.update_account( + account_id=2202, + body=main.AccountUpdate(egress_public_ip="47.96.154.74"), + db=db, + user=SimpleNamespace(id=7, role="operator"), + ) + + self.assertEqual(response, {"id": 2202}) + self.assertEqual(account.egress_public_ip, "47.96.154.74") + stop.assert_awaited_once_with(2202) + self.assertEqual(db.commit.await_count, 2) + values = db.execute.await_args.args[0].compile().params + self.assertIn("已保留登录凭证", values["error_message"]) + self.assertNotIn("cookie_data", values) + self.assertNotIn("im_session_data", values) + async def test_log_stats_uses_one_aggregate_and_respects_ownership(self): engine = create_async_engine("sqlite+aiosqlite:///:memory:") async with engine.begin() as connection: diff --git a/backend/tests/test_batch_start_api.py b/backend/tests/test_batch_start_api.py index 40ad239..584b49a 100644 --- a/backend/tests/test_batch_start_api.py +++ b/backend/tests/test_batch_start_api.py @@ -1,6 +1,7 @@ from __future__ import annotations import asyncio +import json import os import sys import unittest @@ -262,6 +263,95 @@ class BatchStartApiTests(unittest.IsolatedAsyncioTestCase): # return the connection: one before validation, one before the wait. self.assertEqual(events, ["release", "assess", "release", "start-worker"]) + async def test_changed_egress_preserves_valid_credentials(self): + ready_assessment = { + "login_mode": "im_direct", + "should_reset": False, + "can_skip_browser": True, + "message": "ready", + "cookie_valid": True, + "im_ready": True, + } + scenarios = ( + ({"cookies": []}, "47.96.154.74"), + ({"cookies": [], "credential_egress_public_ip": "116.62.23.103"}, "47.96.154.74"), + ({"cookies": [], "credential_egress_public_ip": "47.96.154.74"}, ""), + ) + modes = (("im_direct", False), (None, False), (None, True)) + for storage, selected_ip in scenarios: + for requested_mode, wait_for_ready in modes: + with self.subTest(storage=storage, mode=requested_mode, batch=wait_for_ready): + cookie_data = json.dumps(storage) + account = SimpleNamespace( + id=506, + status="offline", + qr_code_base64=None, + error_message="old channel warning", + cookie_data=cookie_data, + im_session_data="saved-session", + egress_public_ip=selected_ip, + ) + db = SimpleNamespace(commit=AsyncMock()) + with ( + patch.object(main.manager, "is_running", return_value=False), + patch.object(main.manager, "start_worker", AsyncMock(return_value=True)) as start, + patch.object(main, "_get_account_cookie_data", return_value=cookie_data), + patch.object(main, "_reset_account_credentials", AsyncMock()) as reset, + patch.object(main, "assess_account_credential", AsyncMock(return_value=ready_assessment)) as assess, + ): + result = await main._start_account_rpa_impl( + account, db, requested_mode, wait_for_ready=wait_for_ready + ) + + reset.assert_not_awaited() + assess.assert_awaited_once_with( + cookie_data, "saved-session", + startup_priority=True, egress_public_ip=selected_ip, + ) + start.assert_awaited_once_with( + 506, login_mode="im_direct", + wait_until_ready=wait_for_ready, credential_prevalidated=True, + ) + self.assertEqual(account.cookie_data, cookie_data) + self.assertEqual(account.im_session_data, "saved-session") + self.assertIsNone(account.error_message) + self.assertTrue(result["skip_qr"]) + self.assertTrue(result["skip_browser"]) + + async def test_changed_egress_still_rejects_invalid_im_credentials(self): + account = SimpleNamespace( + id=506, + status="offline", + qr_code_base64=None, + error_message=None, + im_session_data="saved-session", + egress_public_ip="47.96.154.74", + ) + db = SimpleNamespace(commit=AsyncMock()) + invalid_assessment = { + "login_mode": "browser", + "should_reset": False, + "can_skip_browser": False, + "message": "缺少 IM 签名密钥(web_protect/keys),请用浏览器登录补全", + "cookie_valid": True, + "im_ready": False, + } + + with ( + patch.object(main.manager, "is_running", return_value=False), + patch.object(main.manager, "start_worker", AsyncMock()) as start, + patch.object(main, "_get_account_cookie_data", return_value='{"cookies": []}'), + patch.object(main, "_reset_account_credentials", AsyncMock()) as reset, + patch.object(main, "assess_account_credential", AsyncMock(return_value=invalid_assessment)), + ): + with self.assertRaises(main.HTTPException) as error: + await main._start_account_rpa_impl(account, db, "im_direct") + + self.assertEqual(error.exception.status_code, 400) + self.assertEqual(error.exception.detail, invalid_assessment["message"]) + reset.assert_not_awaited() + start.assert_not_awaited() + async def test_batch_start_does_not_launch_interactive_browser_login(self): account = SimpleNamespace( id=504, diff --git a/backend/tests/test_credential_responsiveness.py b/backend/tests/test_credential_responsiveness.py index 8c2810c..685c6de 100644 --- a/backend/tests/test_credential_responsiveness.py +++ b/backend/tests/test_credential_responsiveness.py @@ -3,10 +3,26 @@ import unittest from types import SimpleNamespace from unittest.mock import patch -from rpa_engine.credential import validate_im_session +from rpa_engine.credential import credential_egress_mismatch, validate_im_session class CredentialResponsivenessTests(unittest.IsolatedAsyncioTestCase): + def test_legacy_egress_marker_comparison_is_diagnostic(self): + legacy = '{"cookies": []}' + + self.assertFalse(credential_egress_mismatch(legacy, "")) + self.assertTrue(credential_egress_mismatch(legacy, "47.96.154.74")) + + def test_stamped_egress_marker_comparison(self): + stamped = ( + '{"cookies": [], ' + '"credential_egress_public_ip": "47.96.154.74"}' + ) + + self.assertFalse(credential_egress_mismatch(stamped, "47.96.154.74")) + self.assertTrue(credential_egress_mismatch(stamped, "116.62.23.103")) + self.assertTrue(credential_egress_mismatch(stamped, "")) + async def test_uid_lookup_does_not_block_event_loop(self): event_loop_thread_id = threading.get_ident() lookup_thread_ids = [] diff --git a/backend/tests/test_cross_account_isolation.py b/backend/tests/test_cross_account_isolation.py new file mode 100644 index 0000000..99f6acc --- /dev/null +++ b/backend/tests/test_cross_account_isolation.py @@ -0,0 +1,294 @@ +"""托管多个账号时的会话归属隔离回归测试。 + +复现的缺陷:账号 A 的处理链路收到属于账号 B 的会话(0:1:B:B的好友)后, +resolve_peer_uid 把末段当成「对方」、normalize_conversation_id 再拼成 +0:1:A:B的好友,于是账号 A 用自己的凭证把自动回复发给了账号 B 的好友。 +""" +from __future__ import annotations + +import asyncio +import os +import sys +import unittest +from pathlib import Path +from types import SimpleNamespace +from unittest.mock import AsyncMock, Mock, patch + + +BACKEND_DIR = Path(__file__).resolve().parents[1] +os.environ.setdefault("KEFU_DB_TYPE", "sqlite") +os.environ.setdefault("KEFU_DATABASE_URL", "") +os.environ.setdefault("KEFU_DB_PATH", str(BACKEND_DIR / "kefu.db")) +if str(BACKEND_DIR) not in sys.path: + sys.path.insert(0, str(BACKEND_DIR)) + +from rpa_engine.douyin_im import hosted_registry +from rpa_engine.douyin_im import ws_client as ws_module +from rpa_engine.douyin_im.conv_util import conversation_belongs_to +from rpa_engine.douyin_im.http_client import DouyinImHttpClient +from rpa_engine.douyin_im.service import DouyinImService +from rpa_engine.douyin_im.session import DouyinImSession +from rpa_engine.douyin_im.ws_client import DouyinImWsClient + +ACCOUNT_A_UID = 7670159096859706425 +ACCOUNT_B_UID = 7670157997767050299 +PEER_OF_B = 66578464308 + + +class ConversationOwnershipTests(unittest.TestCase): + def test_foreign_single_chat_is_rejected(self): + self.assertFalse( + conversation_belongs_to( + f"0:1:{ACCOUNT_B_UID}:{PEER_OF_B}", ACCOUNT_A_UID + ) + ) + + def test_own_conversation_in_either_position(self): + self.assertTrue( + conversation_belongs_to(f"0:1:{ACCOUNT_A_UID}:{PEER_OF_B}", ACCOUNT_A_UID) + ) + self.assertTrue( + conversation_belongs_to(f"0:1:{PEER_OF_B}:{ACCOUNT_A_UID}", ACCOUNT_A_UID) + ) + + def test_undecidable_shapes_pass_through(self): + # 缺 my_uid / 群聊 / 裸 UID:本来就判不了归属,保守放行 + self.assertTrue(conversation_belongs_to(f"0:1:{ACCOUNT_B_UID}:{PEER_OF_B}", 0)) + self.assertTrue(conversation_belongs_to("0:2:123:456", ACCOUNT_A_UID)) + self.assertTrue(conversation_belongs_to(str(PEER_OF_B), ACCOUNT_A_UID)) + self.assertTrue(conversation_belongs_to("", ACCOUNT_A_UID)) + + +class ForeignMessageDropTests(unittest.IsolatedAsyncioTestCase): + def _service(self) -> DouyinImService: + service = DouyinImService( + session=DouyinImSession(cookies={"sessionid": "a"}, my_uid=ACCOUNT_A_UID), + match_reply=AsyncMock(return_value=["自动回复"]), + log_fn=AsyncMock(), + account_id=1, + ) + service._running = True + return service + + async def test_message_from_another_account_never_schedules_a_reply(self): + service = self._service() + service._resolve_peer_profile = AsyncMock( + return_value=("B 的好友", "", str(PEER_OF_B)) + ) + + with patch( + "rpa_engine.douyin_im.service.system_logger.record", Mock() + ) as record: + result = await service._prepare_incoming( + { + "conversation_id": f"0:1:{ACCOUNT_B_UID}:{PEER_OF_B}", + "sender_uid": str(PEER_OF_B), + "content": "在吗", + "server_message_id": "7665317099296081465", + } + ) + + self.assertIsNone(result) + service.match_reply.assert_not_awaited() + service.log_fn.assert_not_awaited() + self.assertEqual(service._conv_meta, {}) + self.assertTrue(record.called) + + async def test_own_message_is_still_processed(self): + service = self._service() + conv_id = f"0:1:{ACCOUNT_A_UID}:{PEER_OF_B}" + service._resolve_peer_profile = AsyncMock( + return_value=("我的好友", "", str(PEER_OF_B)) + ) + service._resolve_cooldown_seconds = AsyncMock(return_value=0) + service._resolve_reply_delay_seconds = AsyncMock(return_value=0) + service._send_auto_reply = AsyncMock() + + with patch("rpa_engine.douyin_im.service.system_logger.record", Mock()): + send_reply = await service._prepare_incoming( + { + "conversation_id": conv_id, + "sender_uid": str(PEER_OF_B), + "content": "在吗", + "server_message_id": "7665317099296081466", + } + ) + + self.assertIsNotNone(send_reply) + service.match_reply.assert_awaited() + self.assertIn(conv_id, service._conv_meta) + + +class ForeignSendRefusalTests(unittest.IsolatedAsyncioTestCase): + async def test_send_refuses_a_conversation_owned_by_another_account(self): + client = DouyinImHttpClient( + DouyinImSession(cookies={"sessionid": "a"}, my_uid=ACCOUNT_A_UID), + account_id=1, + ) + resolve_meta = AsyncMock() + + with ( + patch.object( + DouyinImHttpClient, + "_resolve_authoritative_uid", + return_value=ACCOUNT_A_UID, + ), + patch.object( + DouyinImHttpClient, "resolve_conversation_meta", resolve_meta + ), + patch("rpa_engine.douyin_im.http_client.system_logger.record", Mock()), + ): + sent = await client.send_text_message( + f"0:1:{ACCOUNT_B_UID}:{PEER_OF_B}", + "你好", + _bypass_global_queue=True, + ) + + self.assertFalse(sent) + # 关键断言:拒发必须发生在解析 ticket / 真正写出去之前 + resolve_meta.assert_not_awaited() + self.assertIn("不是本账号", client.last_error) + self.assertFalse(client.last_send_channel_retryable) + + +class HostedPeerLoopTests(unittest.IsolatedAsyncioTestCase): + """两个本系统托管的账号之间不得互相自动回复(无限回环 → 抖音风控)。""" + + def _service(self) -> DouyinImService: + service = DouyinImService( + session=DouyinImSession(cookies={"sessionid": "a"}, my_uid=ACCOUNT_A_UID), + match_reply=AsyncMock(return_value=["自动回复"]), + log_fn=AsyncMock(), + account_id=1, + ) + service._running = True + service._resolve_cooldown_seconds = AsyncMock(return_value=0) + service._resolve_reply_delay_seconds = AsyncMock(return_value=0) + return service + + def tearDown(self): + hosted_registry.unregister(ACCOUNT_B_UID) + + async def _incoming_from(self, service, peer_uid: int, message_id: str): + service._resolve_peer_profile = AsyncMock( + return_value=("对方", "", str(peer_uid)) + ) + with patch("rpa_engine.douyin_im.service.system_logger.record", Mock()): + return await service._prepare_incoming( + { + "conversation_id": f"0:1:{ACCOUNT_A_UID}:{peer_uid}", + "sender_uid": str(peer_uid), + "content": "在吗", + "server_message_id": message_id, + } + ) + + async def test_no_auto_reply_to_another_hosted_account(self): + hosted_registry.register(ACCOUNT_B_UID) + service = self._service() + + result = await self._incoming_from(service, ACCOUNT_B_UID, "1") + + self.assertIsNone(result) + service.match_reply.assert_not_awaited() + # 消息本身照常入库,只是标记为未回复 + statuses = [ + call.kwargs.get("status") for call in service.log_fn.await_args_list + ] + self.assertIn("received", statuses) + self.assertIn("ignored", statuses) + + async def test_ordinary_follower_still_gets_a_reply(self): + hosted_registry.register(ACCOUNT_B_UID) + service = self._service() + service._send_auto_reply = AsyncMock() + + result = await self._incoming_from(service, PEER_OF_B, "2") + + self.assertIsNotNone(result) + service.match_reply.assert_awaited() + + +class FrontierDeviceExclusivityTests(unittest.IsolatedAsyncioTestCase): + """同一个 frontier 设备号同时只允许一个账号建连。""" + + WS_URL = ( + "wss://frontier-im.douyin.com/ws/v2?fpid=9&device_id=987654321&" + "token=shared-token" + ) + + def setUp(self): + ws_module._FRONTIER_DEVICE_OWNERS.clear() + + def tearDown(self): + ws_module._FRONTIER_DEVICE_OWNERS.clear() + + def _client(self, account_id: int) -> DouyinImWsClient: + client = DouyinImWsClient( + DouyinImSession(cookies={"sessionid": "s"}, ws_urls=[self.WS_URL]), + AsyncMock(), + account_id=account_id, + ) + client._running = True + client._task = SimpleNamespace(done=lambda: False) + return client + + def test_second_account_is_denied_while_the_first_holds_the_device(self): + first = self._client(11) + second = self._client(12) + + self.assertTrue(first._claim_frontier_device(self.WS_URL)) + self.assertFalse(second._claim_frontier_device(self.WS_URL)) + self.assertEqual(second._blocked_device_owner_id, 11) + # 让出方不会被误标为已占用,重连时仍是 HTTP 轮询兜底 + self.assertFalse(second.connected) + + def test_device_is_taken_over_after_the_owner_stops(self): + first = self._client(11) + second = self._client(12) + self.assertTrue(first._claim_frontier_device(self.WS_URL)) + + first._running = False + first._release_frontier_device() + + self.assertTrue(second._claim_frontier_device(self.WS_URL)) + + def test_same_account_reconnect_keeps_its_own_device(self): + client = self._client(11) + self.assertTrue(client._claim_frontier_device(self.WS_URL)) + self.assertTrue(client._claim_frontier_device(self.WS_URL)) + + async def test_run_loop_does_not_open_a_second_connection(self): + owner = self._client(11) + self.assertTrue(owner._claim_frontier_device(self.WS_URL)) + + blocked = self._client(12) + blocked._prepare_url = AsyncMock(return_value=self.WS_URL) + run_connection = AsyncMock() + blocked._run_connection = run_connection + + async def stop_after_first_backoff(_seconds): + blocked._running = False + + with ( + patch.object(ws_module, "_reconnect_delay", return_value=0.0), + patch.object(ws_module.system_logger, "record") as record, + patch.object(ws_module.asyncio, "sleep", stop_after_first_backoff), + ): + await asyncio.wait_for(blocked._run_loop(self.WS_URL), timeout=1.0) + + run_connection.assert_not_awaited() + self.assertFalse(blocked.connected) + self.assertTrue(record.called) + + def test_url_without_device_id_is_not_blocked(self): + first = self._client(11) + second = self._client(12) + url = "wss://frontier-im.douyin.com/ws/v2?fpid=9&token=t" + + self.assertTrue(first._claim_frontier_device(url)) + self.assertTrue(second._claim_frontier_device(url)) + + +if __name__ == "__main__": + unittest.main() diff --git a/backend/tests/test_egress_channels.py b/backend/tests/test_egress_channels.py index 708d71c..de9d6de 100644 --- a/backend/tests/test_egress_channels.py +++ b/backend/tests/test_egress_channels.py @@ -1,5 +1,6 @@ from __future__ import annotations +import asyncio import os import sys import time @@ -25,6 +26,7 @@ from rpa_engine.egress_channels import ( reset_egress_cache_for_tests, resolve_send_channels, ) +from rpa_engine.source_bound_proxy import SourceBoundProxy from models.db_migrate import migrate_accounts_table from models.models import Account @@ -99,6 +101,49 @@ class EgressChannelTests(unittest.IsolatedAsyncioTestCase): with self.assertRaises(EgressChannelUnavailable): await resolve_send_channels("198.51.100.99", 2) + async def test_browser_proxy_binds_selected_source_address(self): + observed_peer = asyncio.get_running_loop().create_future() + + async def target_handler(reader, writer): + if not observed_peer.done(): + observed_peer.set_result(writer.get_extra_info("peername")[0]) + payload = await reader.readexactly(4) + writer.write(payload) + await writer.drain() + writer.close() + await writer.wait_closed() + + target = await asyncio.start_server(target_handler, "127.0.0.1", 0) + target_port = target.sockets[0].getsockname()[1] + proxy = await SourceBoundProxy("127.0.0.2").start() + writer = None + try: + reader, writer = await asyncio.open_connection( + "127.0.0.1", + int(proxy.server_url.rpartition(":")[2]), + ) + writer.write( + ( + f"CONNECT 127.0.0.1:{target_port} HTTP/1.1\r\n" + f"Host: 127.0.0.1:{target_port}\r\n\r\n" + ).encode("ascii") + ) + await writer.drain() + response = await reader.readuntil(b"\r\n\r\n") + self.assertIn(b"200 Connection Established", response) + + writer.write(b"ping") + await writer.drain() + self.assertEqual(await reader.readexactly(4), b"ping") + self.assertEqual(await asyncio.wait_for(observed_peer, 1), "127.0.0.2") + finally: + if writer is not None: + writer.close() + await writer.wait_closed() + await proxy.close() + target.close() + await target.wait_closed() + class EgressMigrationTests(unittest.TestCase): def test_mysql_accounts_uses_longtext_for_browser_payloads(self): diff --git a/backend/tests/test_reply_queue_integration.py b/backend/tests/test_reply_queue_integration.py index 84d631e..24d1286 100644 --- a/backend/tests/test_reply_queue_integration.py +++ b/backend/tests/test_reply_queue_integration.py @@ -18,6 +18,8 @@ if str(BACKEND_DIR) not in sys.path: from auth.system_settings import SystemSettingsData, set_cached_settings from rpa_engine.douyin_im.service import DouyinImService +from rpa_engine.douyin_im.session import DouyinImSession +from rpa_engine.douyin_im import service as service_module from rpa_engine.playwright_worker import DouyinWorker @@ -114,6 +116,74 @@ def _build_service(delay_seconds: int = 60): class ReplyQueueIntegrationTests(unittest.IsolatedAsyncioTestCase): + async def test_kick_does_not_replay_through_browser_fallback(self): + callback = AsyncMock() + fallback = AsyncMock(return_value=(True, "must not run")) + session = DouyinImSession(cookies={"sessionid": "test"}, my_uid=999) + service = DouyinImService( + session=session, + match_reply=AsyncMock(), + log_fn=AsyncMock(), + account_id=1, + send_fallback=fallback, + on_session_invalid=callback, + ) + service._running = True + kicked_http = SimpleNamespace( + send_text_message=AsyncMock(return_value=False), + last_error="decision=KICK", + last_send_needs_refresh=False, + ) + context = AsyncMock() + context.__aenter__.return_value = kicked_http + context.__aexit__.return_value = None + + with ( + unittest.mock.patch.object( + service_module, "DouyinImHttpClient", return_value=context + ), + unittest.mock.patch( + "rpa_engine.douyin_im.service.system_logger.record", Mock() + ), + ): + sent, _ = await service._send_text("0:1:999:123", "hello") + + self.assertFalse(sent) + fallback.assert_not_awaited() + callback.assert_awaited_once() + + async def test_fresh_session_replaces_send_and_ws_state_atomically(self): + current = DouyinImSession( + cookies={"sessionid": "old"}, + my_uid=999, + conv_meta={"old": {"ticket": "one"}}, + ) + current.egress_public_ip = "203.0.113.10" + current.egress_source_ip = "10.0.0.10" + current.egress_auto_attempts = 2 + service = DouyinImService( + session=current, + match_reply=AsyncMock(), + log_fn=AsyncMock(), + account_id=1, + ) + service._ws_client = SimpleNamespace(session=current) + fresh = DouyinImSession( + cookies={"sessionid": "fresh"}, + my_uid=999, + conv_meta={"new": {"ticket": "two"}}, + ) + + await service.replace_session(fresh) + + self.assertIs(service.session, fresh) + self.assertIs(service._ws_client.session, fresh) + self.assertEqual(service.session.cookies["sessionid"], "fresh") + self.assertEqual(set(service.session.conv_meta), {"old", "new"}) + self.assertEqual(service.session.egress_public_ip, "203.0.113.10") + self.assertEqual(service.session.egress_source_ip, "10.0.0.10") + self.assertEqual(service.session.egress_auto_attempts, 2) + async def test_kick_response_takes_account_offline_immediately(self): callback = AsyncMock() service, _, _ = _build_service() diff --git a/backend/tests/test_send_entry_and_worker_lifecycle.py b/backend/tests/test_send_entry_and_worker_lifecycle.py index dff39a3..2828046 100644 --- a/backend/tests/test_send_entry_and_worker_lifecycle.py +++ b/backend/tests/test_send_entry_and_worker_lifecycle.py @@ -205,7 +205,7 @@ class SendTextMessageEntryTests(unittest.IsolatedAsyncioTestCase): self.assertTrue(client.last_send_needs_refresh) self.assertEqual(client.last_request_debug, "response status=401") - async def test_retryable_rejection_switches_channels_serially(self): + async def test_retryable_network_failure_switches_channels_serially(self): client = self._make_client(account_id=89) client.session.egress_auto_attempts = 2 routes = [ @@ -216,7 +216,7 @@ class SendTextMessageEntryTests(unittest.IsolatedAsyncioTestCase): first = SimpleNamespace( send_text_message=AsyncMock(return_value=False), last_send_meta={}, - last_error="decision=KICK", + last_error="connect timeout", last_send_needs_refresh=False, last_send_channel_retryable=True, last_request_debug="first route", @@ -264,8 +264,82 @@ class SendTextMessageEntryTests(unittest.IsolatedAsyncioTestCase): second.send_text_message.assert_awaited_once() self.assertEqual(client.last_request_debug, "second route") + async def test_kick_never_switches_public_channels(self): + client = self._make_client(account_id=90) + client.session.egress_auto_attempts = 2 + routes = [ + EgressChannel("198.51.100.10", "10.0.0.10", "eth0", True), + EgressChannel("198.51.100.11", "10.0.0.11", "eth1", False), + ] + kicked = SimpleNamespace( + send_text_message=AsyncMock(return_value=False), + last_send_meta={}, + last_error="decision=KICK", + last_send_needs_refresh=False, + last_send_channel_retryable=False, + last_request_debug="terminal kick", + ) + context = MagicMock() + context.__aenter__ = AsyncMock(return_value=kicked) + context.__aexit__ = AsyncMock(return_value=None) + queued_factory = MagicMock(return_value=context) + + async def execute_submission(account_id, operation, description=""): + return await operation() + + with ( + patch( + "rpa_engine.douyin_im.traffic_control.submit_outbound", + AsyncMock(side_effect=execute_submission), + ), + patch( + "rpa_engine.douyin_im.http_client.resolve_send_channels", + AsyncMock(return_value=routes), + ), + patch.object(http_client_module, "DouyinImHttpClient", queued_factory), + ): + sent = await client.send_text_message("0:1:10001:20002", "hello") + + self.assertFalse(sent) + queued_factory.assert_called_once() + kicked.send_text_message.assert_awaited_once() + class WorkerLifecycleTests(unittest.IsolatedAsyncioTestCase): + async def test_browser_launch_uses_selected_source_proxy(self): + launch = AsyncMock(return_value="browser") + pw = SimpleNamespace(chromium=SimpleNamespace(launch=launch)) + source_proxy = AsyncMock(return_value={"server": "http://127.0.0.1:43210"}) + + with ( + patch.object( + playwright_worker_module, + "ensure_browser_display", + AsyncMock(), + ), + patch.object( + playwright_worker_module, + "playwright_proxy_for_source", + source_proxy, + ), + patch.object(playwright_worker_module, "playwright_proxy") as global_proxy, + ): + browser = await playwright_worker_module._launch_chromium( + pw, + ["--no-sandbox"], + headless=True, + source_ip="10.0.0.6", + ) + + self.assertEqual(browser, "browser") + source_proxy.assert_awaited_once_with("10.0.0.6") + global_proxy.assert_not_called() + launch.assert_awaited_once_with( + headless=True, + args=["--no-sandbox"], + proxy={"server": "http://127.0.0.1:43210"}, + ) + async def test_virtual_display_starts_before_playwright_driver(self): ensure_display = AsyncMock()