From 2fc864cc00afebbeb3eb0f18b676580428eb6f8d Mon Sep 17 00:00:00 2001 From: Your Name Date: Tue, 1 Sep 2026 15:31:05 +0800 Subject: [PATCH] =?UTF-8?q?=E6=9B=B4=E6=96=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .workbuddy/memory/2026-08-28.md | 21 ++ backend/kefu.db-shm | Bin 32768 -> 32768 bytes backend/kefu.db-wal | Bin 4556752 -> 4556752 bytes backend/main.py | 37 ++- backend/rpa_engine/credential.py | 40 +++ backend/rpa_engine/douyin_im/auth.py | 1 + backend/rpa_engine/douyin_im/conv_util.py | 23 ++ backend/rpa_engine/douyin_im/frontier.py | 7 +- backend/rpa_engine/douyin_im/http_client.py | 55 +++- backend/rpa_engine/douyin_im/service.py | 162 +++++++++- backend/rpa_engine/douyin_im/ws_client.py | 91 +++++- backend/rpa_engine/playwright_worker.py | 127 +++++++- backend/rpa_engine/source_bound_proxy.py | 233 ++++++++++++++ backend/tests/test_account_pagination.py | 34 ++ backend/tests/test_batch_start_api.py | 90 ++++++ .../tests/test_credential_responsiveness.py | 18 +- backend/tests/test_cross_account_isolation.py | 294 ++++++++++++++++++ backend/tests/test_egress_channels.py | 45 +++ backend/tests/test_reply_queue_integration.py | 70 +++++ .../test_send_entry_and_worker_lifecycle.py | 78 ++++- 20 files changed, 1381 insertions(+), 45 deletions(-) create mode 100644 .workbuddy/memory/2026-08-28.md create mode 100644 backend/rpa_engine/source_bound_proxy.py create mode 100644 backend/tests/test_cross_account_isolation.py 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 9b62e0d149760a0feb8ddbd7a9c103015880419e..e22d7bf9ca3d24e1bb7768dac965e271226695b1 100644 GIT binary patch delta 1158 zcmb7CZ%oxy7=GX9{^7^5wY9mSQc3#07cl$S*`9Nr=Q+>&p5Na6 zl1U_!cqi*>0ZlF@%IpBq^|M4J)6V@eT;IuoVEN6ApCgCce{U=5uk645fAf>QweP)$ z)`ah(4t&)+lx9ZKmiy#=as&a#oJfnz3JZ@2M$Qn83G(+6U6M0k&U6Vq%qd@r=@Z+r8!> zXF*!rknjEPTB(W4cC@;}g;M3~OyAlI&U0}111?IH`x^N)U*)^}n954?N6jhsZJw!O zjv1`2LbcK)r^~+3{#aKBR}3K$`^?e$zNs!eiX22SL~ZDRDQ-CF_F|TnbekQG@3}pg zpK9fNNH5|Qw5aH^N8)*K6{A)}wedCkT`ys_@USgyKFZFYoQ4nb8U1>41to=&G^K_ zV%M8%TZ5r{oQXMj3JJV6aes`$+{G`CR~BtMEAyF#ClSG`*g6pzp!vL;*J&hbDz*pl ztkObrX!}m*FLp2+3)P7864{s9BRkyejiVW@>Xw?JH$oxCbj-z4Y{IsQrk_$ScXPc) zVy3C{d;b*7Kt7&E1*+7{w+(OBFv^v;_>a>9-op)O!g@7|ZPl)Cfkx#$=0JC=^EdmD zjYVp7*!UiJ+;`-^6O|Phr1$AC|HQ-m8;=?;3!Be-VpCQ>n9mAH?X_>cT?0I?718^H C=OhvU delta 719 zcmb7BO-NKx6u#$9PwQoAF;rTMKu66o!zO#pLQ`ufed1Ga{EI2+k2K448b`7kW+4%R z;F&;bk&Z2Dmz|wd6cZA(8(K{gHB?9oK_hSx)VZ&jRjVG{d%k<__kHJ_cf(;B4$Ip+ zQsqK@9$p|s-cR1X#)YWS;@b-SFL0ByekJ!q<%G33Z5*!3U-)0XZO+>=`}Xqso&M5! z2gkWcGMfy2^$=P5B2d}Ur+}-SyR)B&au|gee(oDVOE1M(T_87O_cundR-d7T`cF#W z#^mt?l*J7Ctnh-krmECu0{?cTWLIb1N;K+6P>w7(M=oZD8p-?6s@RpI+lZ497D_Aei}> z)sx4JAdUA($)F3D&N{be+Q`K7;-M&n!SsJSK+n%r%YR50#Iz`aAs*L7RToXkWpcp9 zMceSpC8z33n-MMu*aUZPT%yC=qQyipT<6pg^!7V>$VwdVALK^PKm|(=(&UwO>M4$4 z@@j$YF)KT&$qBT)*&d{s*gLod}dAx>xESkTWnKpn$UV~brI g6XK4TfOqf_zB9W7ACH<$)S-qej!@{PrFu910Lm%qOaK4? diff --git a/backend/kefu.db-wal b/backend/kefu.db-wal index 55a5ef11799d2c37bd76b0adac39518cecaa4cf5..570a865eb12019f9b0a2e8fa16d988171340a42d 100644 GIT binary patch delta 15318 zcmeHO2V7HE-@Z2qfskd;t#*9p+%RqsTea=`eZTM5+W$2p$#&wOl(r&XFt+BRPhs|L$l^Zg0)3V1E8{!S+;|)=UjFRk>(oxxlcw>}d zWOjaWZb5$9n4;{$B18Q1gQEcDEy~Q!#|0U1sYS(w8TrGq3)6D*vrO6J4DqHYL$Rw?o6TkzZ#J53(MCtK*{(Mk z<1LnWb8L*oVl|s$CyVFU=bce~iiTX(6nZFmza9$;Dmi@6+v_tbk=vCL1qgo-ek`mJ z&VoxPMIkip-zr}(8liO#`8ymq?Y>4XW)9hodM}l%lHB+N}aZ`IXn(c}+Rn*<0d64dNn-cZj!= zH<_2gYs=GcA96qEA}BRPt%ARv)5_7e>GDe)W^C$)i}XWL>Kd+1nK~hgyQB>ei7&pe zWr^Hw6?uCvrTW6bH??nZgnDQcX<8xWX!%wF+*~gAgRT8FW0Z|&(^aC5o5+F7GIdI{ zQKa1opW`#=p!Z>{yEQQk!P3K56?1}VE~ zuhLS$)6ks4?9A-kvDsN%+)SSy+FszuCtn8oA6jX_gG2>_LO4sJYrz49C^;3x1HK2B zrUL`#X?63pE4~(|-+I%-*Of7h0{g)xa0R>!27`7W48E(?HA5MDWH%)OFRq-%^^6UK zC!=*2K)0;Zjp zMa*lf<8#Ccval|S;I}tq047gW@kF%?$BqX)5GGDl`9Ws9L%Zv?bHJ>XhZU0hMMGUR zQt-qKtulZDKY(jsI;a98KpN-^IsrKlfH%=52#&bgoQu6M5PdgT=SeGlI5+p^3_Z4j ztZ)yygz|{x?ciAGGjs`QN;WMV6zE0Ka(Qg0L@r!*QXkN&c8%Ng75b66S*`6hy8}ZJ z8#S8tSJDM^CSBK30ycH*tZ&m{%P^+hqB~GmfJf>Mx`4nA5vy+U`K z<4@7JvGs78|L(MLd|mgRP%rdViQlKF17A_p{$D6+>$N%QcKUsku)?qwC>ZG7{o=KO z3p*TG@8PkSmU&6g&k*5t7+iga>JucR6HrlNwNH5TvXcIII)dj<)AFFjLhcPzJ z9;Y{%TIc2$$D8a9yTgo8+-8llz?pAqBDj{=Sh(>$P4JMooGjZ=>>&pU9UccD^uQJ-s@9cs_ z9zU7inE+!r0me04DFcjUiZXYA489?kCZ*Ai?F1kPbLTbzkleZT5a~xCa&H4f*0o+y z_J^n`5W#^D)e`@>QNwSqK3Q|koAWQKm@|RmzQbL~g|OnG`Vvb2S^X7FAA}u$RqvCC zc?z17laZT8u&oTqe^ZYTiP!G#7qxsy)j|TiDGc!PgnA)`E{u|YFW^w=Fv?$(6VJxn zh9(lsZQ5g-#^hI28nJBM{?%Ea?0gGKLv z;D;a(?(3$}funnW=`t!+lS0r#333=D=Yv9!2XavNjj~AKP0(DaBu`lNeqa<_w_8@+ zKJ)Gm>o<$cFS%ZmH^O!B1vm{p1lzzwFb<@EZXgaAffhsOk6aA+rEq(l##BLau`OT)Q<) zxI1*j6xF>Y%!{J4^;-j9tPnjAoh9h9q|x?JHpMYa{$J+`vfcd}*A|*dcZh#A*1hfG`;Bi~ax?8eM5{-ZxKxY>g78EM$ zKQZINY&Q^ZHxG0nc*^R*^TVZGe*3LcVJ`+eJyIx-6wTd13IWDWx#{f;xN)&2 zT*yL4hU=mSi#2Qi0Zle&MziLfcoBP`3kNPw-Cc9=HP=8tM|nI)2D%p>=wM~-K&L|e za&52&v$vz<<=V~vj4o?*=fPeGh}vQ{e{1KoN3pWb-!SdC%O3Yv%P&tW`W2?SoVkXf zIUaXJ{}R=p{{Z=z2C7q6fU=(a)TJouQyN^j5gmeM+f)+S+)FRd=wkCtt_9?9yarmK zMcY(7sdnZ5PhV+Ua~_a*)03X)7H`rMAzlhSNyMq4gC=FqiED^7CyxG`BTYGe(3c%G>-+aAY@oU~G5`edJ3&4dVNeKR4Z3sS>gy3^(LvYr>du?+gOeDl_9F?1R z1}`@b{_Q1)GkfKdgCvo|lVJ|dlM#IL;*t2g{CDGXyf$3%Ijap%PdR%>$mRbq{{I7v zUp!@0M%U5p<1ygU4t|1*@$(7e$6u;|_6RqSq50*C@}_Bj==mbah=HSfd`oEm5N=VE zw4Ym4PyBrd@n2j^{8?*AQC`1}>yiAAV*#w7Tv~Ulz1sTRR6^wYF+`qEh`jd&h^f58 zrJ``C1j6TAWi`Rv-P+K^?<|&}C~Nw6ve7in{D~T)zMzur&f6@zL)4fuH#G(-4x-MAC)cMwR-w4rU22-dcyDJBK)GA1j;6`DL-q9RrLhFZk0V#5YX}6 z?|Za(e#|HRm_|r0|3MHSx%?P~X~QkqN=YagJVN zA*QLkhotwnc)~#{cM+-+LsT((yNgh4Lh<lMvL=YYIVE$?+($91A#!#Zx9@&rt{kuZ2^=gLpK5A*N1gor~dA>s-V zA`Y>I!H@RbYBNE%*Ov%`J}zOfl}i}p@e}y*L>g4pN`s4tG&tWE6KwJ3@+JhrR_qX9 z(Sn7sXHUGsFqXywVa4wYgm`uTd4aG6*Jw4fB|s&p<7Ah1Kzt0}R zMG|f4JF9*kOV~pixA$Y%1GjfQe*lC(R5b91x?U+NT>OFMBhN-DYURa11zc_P9`^7K zkr&6*%ZuD;E`f3A&+5xGjflYbj@nINtZ@5QCpX?(U438C_ccgZN`Ho>c)D1MH@rAX zZm**$dO}mwl}*qL)&wdln*O~m7E_cr@%z)z7>`m?0Q*#lQ81tXjHhqqzjkYFz#NAC zxm&G*dArn_+-+fKUl(v3W&y{S1UL@66wEM2!Hj?DGXLn7_G5^G=?T8TPp`oZ3MRM& zzChO$(ayW@}QC0`NAIAR*9`n_KYPg+rBR(8j0$ z0Uw5EDBTkCY@s**T9F$NV z;7j9P2rjbU*_`ToVr*sZ^Emx3N*EuV!AHApD&G(&1#Wu$xb-tvR4VQ+$Ju^BFH8?O z$4BWwky-&M?+12y0rRg;c`r2bXNvMx!+qY0Skz{4^s9U*{Xhs%f=^^9dBq!B!(mq6 z=9Ptc$NBhXf_E7Du6y)h{+50rCpoBcU}!d<=B#oSfmW z8lfJG)2CF=X=Xtec6oVok?ukO%|ZJIDffCSy{E8WkZjsyUHjY0bdm-y9`lNVH)<3M zk)^ZJ4h-ZTMMJaq-2H}EJa18K!+JeXL=G9Yd%!m%+sS>nSb@dP$mxwO1zRer*r}Mn zft+Lz7Rl$jBUUnau-s{<>%t>*YxELz!(b~x&ae>mJ29@Y+NaU}GMabh( zR3-G?1Z2ZG`=^gOsTOsS(^h5rdtcB1E?3Z%%S z>DlnYR)vn9yJe^9dD>$ZerpOE) z!G~t2>xn+^>IXP*=@@Aisv7Jk^F)=|e!V>4q7iCYp4Ey!z&DO0AJQYfc z9g74*rL*`*TjY0w14nICDg{MF1!==`EidF$_nt2c>zz8FQ=g%hPN@S6Ge-97(y8mH z;-nW$9bd>D7h-N#T9i~WY(($j=3y{stJ1(PD$Xd*&8)5|G7245qtOv#!WU+Y)%#}k zLihJ6gK1RKUHUT@UMTnLQGKktPFtK!n!s+FLot&Zxcf#9PA`Fbue8@!Z|bUm86SdRLz&BLt|E$5PB|<&O*uSC40Tl^sS*OG^#?9 zj|dVTC}Gzkifs8%Ar?At`Z)CKOh0^CV6U`;5$;$hb-G$Ny+=Hputa(VEv*b|#S}1( zUBJ|t{IPoM8QSn+!hES2%@#x);X-3!vuz*_;iB4nf^ zM{!}*^S)A%qjJ-=Ssw)U!#OI@QeLFP8>Sx#4@XOPge_ufuV9y_-e5c4mp9!H=N^yt z?hki(*Foe0Ci8f9=4<-qo~QS4FW}4*;2i(xt|+->XrWj2(xE|c`*^>BP=7Wo7FI2i zE-c2OZ8F5Wt0Ys3zmUI_~pORe>jmk8G_m7a!D?~o4g z>&4P?c+xv00VzR}XZ3Mn>h!Sgq7x#NLgDY?NhoaRS*-)YX z0ldXF44Sp~tH*TBcjY6!M8=UC4(!6U)&u9rsjDZIMV7@gz=orI&cC=_NH}ze<2!j< zAxBGvP>NdYoD>I-ed9k#;#Rk5p0}|Jdcz^x=mEWY`YnWe=PPm%{`4gcZ}*81K;!NZ zJ!+QfSIO1mdNCllyIxYqb2nms%ovW-|0dUKA+j2=d}NO?W8KFDPi+j5!-Ro`_}}qB zx4E*u>?zpSIV>Ny@LA>xn9cZVO}x<@gVil-tR1CuWNU z_D@bj*>m$}33^Oe^*2s~7|Vrbe96cWYqwe))!!}d`J^)-6icSR8h4-E(R_v(=BA#wkZ9`qG7G;58NJ^in+o6s2$|U?T-#No=WX#*!elHW;Wy5zh*9*YDLko z@@W1|RzGNhA+8 z$RwbpgnG$l&FD=WNe4PJNGf$lQqUON{zZz|Ryl*XbBL#dSJXm(oCvm=rg1 z)Hvl2uW6j%Oig+aiY*LVIJO9Ek=UAJi^3L-Ee2Z)Y({J*Y-Vg0Y*uWs*lgJB*c{m6 hu*G9*iLDj3*4Wx$Yl|%bTOzi0*pg7vgZ49K{0}M)T?GID delta 10164 zcmeHN30PCd_RpT10J#weN?42nCN~LibF&acWyb|kK&{}48}3`(P!zESu(&}{=Jlxv zDr(URT6JvQ5Y$#}7395d$>MBru}CbWMQUkbakaQv+$}PThoz;ZmBrH{ zw|H5+Ej||C?XBf=+tp5Cv1-3_agV;U`^>Y2Z9mEWc+;_pZ(LX`7LO-j5jlb!UJg$H z4~!%qayfI{!G&jv3E-P&v;c%3B1PQ!bKF4xpB2kM#1=IPxEs9&f|4}VUEq7u-5EV} zwLWywoAo-QS*JIeO$IY4T;@KQ-!r;fLf7SKZG45a(J&(2YydxX_VRX$>ppFo-Y_g= zQi2XdE_e3^5na6If%Z4uNv9B0mEx*vd3o1Xi)206(vfcrPd0XrHD=80 z^F}Z8&`@xoyVipolRP6SY1G)5n31DR>3!2d-A3P1uF+@$k?ojdkR2}@M>=dSZhS}z zh^X)lc6(_sJ$__QIQg7+Fxz%S5SgTv3!WdbFhS-kwod?M8+_E@D;JrMq?v{cps>1) z8<0z6=9eF&Nn{gVBlzMLwO2FUwSi<fg2NkyQ9ed5}~}tgat>v1@su5!Ua(KW`>)+ID=)*(vbxp`h7(=zmp{D+Y~cf+WNz zt2AwQ4!hV06i%*Q?DgQQN5RqT^-x-%t|H~5zt}f2ZGRUDwj`Dmf7DJPpO6p9yYT#K zbvZB4gpO=B=nRZOr>F5ZeTczeq|-x&gqZ12gV~@r7|rQs+b_`Lp<)B*@sxIwnrI`l zWGv5S3M)rOf~_S!CP1%MyMum@6a#*@RU{-#jULf?;<)MYOhV@g8AGvE$fPmT$M;J} z8kI79YBB-oMhpemE}R@@fU@`1aT2|W2}N3BLd~JbZyar+`2(v>yhf_kC>*6CqX+^) z+q4pAouR$K5X8_Xn$d%_XNn}1K18oK8T2|ss38Ptm}X3Yv^kIoQP8375x3qPWHuOR za|lR^^O}TjLLg#{K_MYVBW(nAabCcs*>-JCZrzmC$zPn*>g0}lvN?IY&-?&R4mOiV z&gerdL?8G|lZ}?%ZBtQH?sW;)SE{8_|BO%av-7_1BNbyFG>xnwtH{aZ6>>TGHF=sm zMt(>xg55`ww@DXkT=Q3Ez>TIW2K< z)TbD<@sYW?G*f8vA*&LtL*UCbYG=04_PyaL?aUT5W3{QyI@H{d6+SX|yo6mtO-xN% z0co>U?)YLI^icDYV-DPMr!YHy^~c#P*3MZh7H2Xt$3QS{3w5={eilpgDWZ#r(78nq z;8#zTki*^%mi?^!Tq4dztjUO!h`hnmZPe`+t5~cWB%$U_L_5%wmGJ2YR1XfQ;i?p1 zUr$m5zvrr2v%&g2N(C6cDv&>j#UcP`{rmO1W@mw#;N&`OvI2@OMX zz}|LRA|DSp^kgj8rN=JM;cD)%OpLOnY ze#n^MFR+8L$;T}6@$1vcS7fnP)rzQT`N28A!}&en0t$Odujwbm&I6b_SQ@JP`hQ~@ z0p8tJTK~q>tYT9O(?JI%{aSvAPzbp)(dY^36H)j>5z`)qNSE1y0Re?KVO^^3zjLZ+ zt717hhOqjW;eOF4^FP2rh(l;l5F#}H^U7Mcm&@8EsT?@BlUHXB3?8X^%)>>ZE1;zs z;@|kttX}Liuhq+|cFw{sv>iNBEj|)+1zQJ)yp3Jx_UkILQzUq#QMDn022DOYBjr#! zr2mPgla+oE=EEr$PTQu}{$&-}pue^v1Cgj~1d{Hm_Wu7o|IxsXs%>$$6WwBww~$Ky zeK}=X9;9jIX?K#6xapad^HnZ35eIe*)_nFiRaLmoYO1P0#B5cR2aE0K7lXm0bKQ>J zoNB{@8*Er`1FRD(t4NQqgA< zV0I^s9PDK@A^`$Dp~eu9?WY_E6+X(nd{7mw&xSoL>I^Qpvo+>(u0B*(bYQ8hNV>Tf z^l_*C*?`YidV+7Lmfh^nD#rm0UpZ2Ns;59)e>rYCGc@|sV+;L#(6eMP?;W`pNXwHa zf(Tzy4ig)xNj#8t+er(n@48(Od!xyyn6sB!UI&+vo)4ZOq=Je%_UF=YAt>19M!w-@BUdfmI6G~g`HZc z_gV#_($mg3_`>1|rN#hMEsFNT1DBP7Y*-+oKH&oT%|I_$=^a}r2KBq$tSg{;)jlz- zm&dTU`_T<(;-|1av*jrP7$l+?K)d@(K;1B9WAQq-TcDtt&=gg7jRyO4Y73xmQ5j&) zElLP-)0A#t%Wwl%f;)ri z;nW$BT_z31Cq|jTzMpLYHw@YLGcTj@i?BXYo#F!bbNq!I(0Yz?7+AH!tv#Odz;Q~K zuh(r1>i0fsVu?gZl* z$`cDkHf$Fvi8@lX;uL!VUJSm&SFT`#Wo7=ZMb%r(FC3uwc;IZ%O{i>JG-YyJQB@(2 zt2Z05*A$-!oCR;?sVjsaF^(DvN(Z=k!k;Iq(|C9pgktQ=tPD=g?U?-IH8h?J<{4aR zxO7fHS1v?KpU(!X9w-w)+W6RbFt?LN$k#D?#snk3q(*SyRIc)}03?PgonYQRieO{O z2qW1g7xy7%C?~hFUM^g^pXx3Gt^6n#SXLH5SarOK3wq?nCW{TaP;-!;p;7m+YJI?V za3?o*S4-5YV50Q~T?`!?8H*sTF*GzbN^b&{Pn2z>iJuKV6?*;B8!u2Ch)6PBcddN| zQegkPSV(l?MfXGzAgs zSS(?Qqy8fAH@x-3YS;64>_7H^m|yS!6Jw%7jiwO25q+9LIyp&Q23=EqD&Q5m&VZ+iY5o5rwq*S7VFqf~;zroFOvtIR4Nd#cBs=hTc_0>JkOSp)z z^(bwHm9>|NYD8;q6Zbf`BU*)MB~qvx-+1AKYmkPo2_%oBKGwOS-4kDmzrRyV>Y9s?jD`2`#-(ut)0Za=PDcOz^&Z5yp&p@x5r z6>eTZ&HrQe=KmfN)N52`GIQt%a`W&hQ8xa6!KD~9-Y{SxBTHfveG zw&QMp2VjGLb@3`bsPCfN3aIX?*8YcPbLV~g@&4!A{RjiW%j9{ojy#Mu{n0CZb&iM& z!~;}s5NkNXy*IK~WI#e^S|vpN@SsZaB7XHMf_H|ivV@*?tKBns!+piQwLIHd+WdZ? zxK*QX4XQI$HK3nVBXM`vnElsW{T78-YbzWVn#oo(l{7D;U|J=$jU%o_38fZwx_?(O z^(UHY{mM=&SA2Ch#b&CTY^J*DIT)SdKl9y51A3QYD`gXXSm@6IjQ&J8fK7}9`YZ6E zkuzi<{-#tBCi^k>Oyzm>?5|DR&GG8BhJrB_!b2#yiZWo?DGKnx(ol^U)DDnGz(lr^ z5IGdv=rwP1z-l-M32DkB1qq)GT3%x2NAOS0Rw>M;I*rtp9_Ip7r!8#G-&MHAn{Qg8dUw-F3^!1ujhp$QR*Nn&H{E#Mc)u z8Pyqi1H!XO_f)sVMWtS8o)b1QT&WNfru{YexAnEKWAW5h01gd`sD1I%>9D&zFA>4i-!-DUL zwT%nh(YO*8PQ?d$0VE&33g`MSONy)r63oC-;KgLLWOD@g4)2w8RdptM1lW{G+v7lO z{$Cgelx*P2SM>(WzXEBFu*5M7t&SKtq{3Ws$>!RdZLYofxogkf*E#0>{ZI0&t}PnQ zLVpgfy+R~nLW9g9v{_F#wUu)CDip8skurIX8Z~w-q2Q4#)_0+75;yOhRK*Qs36M6B zJeZZF=}or98Qp#dX)78#28D@y?J7d0HO}Y~Tr1-G3>RF)pqD}=!C2oN3UC!oh^)X9 zSFeu16IHGDcZLF7uAtJSfvalE=AatenXl@{Sl>D#llW~}9hqhZ)@igrmgRVanv%HW zgF{dE-#LLACjL!8d@`tdqI8wl&UF-XLz_J2U%Qv0WKo+{^mF=2i25*Jz7(m%BDwXA zW>Y}*Xm@eh=M6b?(4@IgCkxCMzb-b_V}HZ1L+yH>y(p^m=Hqgt+uCf%UQdw1qWU%- zaPatmbGT`S-=ErOZGpaeHOILr4t--iXtQ34cIbE19kA+0#GM1Pp+qFX5$Ff2*gQwEOZ)_i zueX}yh%t9Ln@RwQjmjYq^`}ytVC%Qt9$H1CK8O9J_cCT}Oj)@G|0h@;I5|ZkhuOo_ z4Nkuhi2!PNJ;K_0+fjy^b_yovJ>jD91z=uBXd9A|)afA_|cjkp__#Q6QoqM8Sw?L^?zaB0VAlA|oObq7X!8 iM4^bJC<;+DjA)E0So~i?>~awR 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()