From 5423af707d869ed6133214071ee0b114b8dda486 Mon Sep 17 00:00:00 2001 From: HycJack <772403255@qq.com> Date: Thu, 30 Apr 2026 04:46:54 +0800 Subject: [PATCH 1/2] Logging error in Loguru Handler #2 --- .../{message_ok_zh.mp3 => message_ok.mp3} | Bin .../audio/{tts_error_zh.mp3 => tts_error.mp3} | Bin .../assets/roles/zhuli/zh/tts_error.mp3 | Bin 5121 -> 12609 bytes talkingq-url/banban/dao/device_setting.py | 21 +++++++++++++++--- talkingq-url/banban/security.py | 7 +++++- talkingq-url/banban/service/avatar_storage.py | 7 +++++- talkingq-url/banban/service/device_setting.py | 8 +++---- 7 files changed, 34 insertions(+), 9 deletions(-) rename talkingq-url/assets/audio/{message_ok_zh.mp3 => message_ok.mp3} (100%) rename talkingq-url/assets/audio/{tts_error_zh.mp3 => tts_error.mp3} (100%) diff --git a/talkingq-url/assets/audio/message_ok_zh.mp3 b/talkingq-url/assets/audio/message_ok.mp3 similarity index 100% rename from talkingq-url/assets/audio/message_ok_zh.mp3 rename to talkingq-url/assets/audio/message_ok.mp3 diff --git a/talkingq-url/assets/audio/tts_error_zh.mp3 b/talkingq-url/assets/audio/tts_error.mp3 similarity index 100% rename from talkingq-url/assets/audio/tts_error_zh.mp3 rename to talkingq-url/assets/audio/tts_error.mp3 diff --git a/talkingq-url/assets/roles/zhuli/zh/tts_error.mp3 b/talkingq-url/assets/roles/zhuli/zh/tts_error.mp3 index 080cd1b69deed0ee92b8ae5c67e5e33422c5ecfa..67e49073de10a38e10f745ab6087dade4376febf 100644 GIT binary patch literal 12609 zcmc)RRa6^M)F|KtcY?bIm*Vd3?xi?IO7SAa-QC?SP$=3$i@Q4%FSIzNf)$-hdVg8} z+kLtB?uX3GN>^V`9=R*R)gF#Yblu#DD41Qpx>L$h zpjUXjeSex^{AnHUyE+HuxDW>1Lqmh{LL0OW;m22AHOb81`1`Y>n)-(5(&8xC#|Il) zRyoiISzw4!CLtEOeSG4Sabax9yW0uj6xAR6~o%U>uoR^ z#d{Adq)9CLD+-cQoGRD)JTDj$+wI6Y22ue7qEN)2mQQ7$-uI$uIwY$ss0OF!v41AE zD|7wK4EHcnd2Wujstki;!+ES&z8`3*I$mr5C;54Lp4Z~np6A5{P+8CIuyfUWsC#hW z9NKj(U^X#-IfXgfh5qegL6%Rxrz21AZy|qGq>scDYc;{m!4B46H)i@Yg(OIS1Q`;U zJZT8yFhtd+MJd!;L8ODepM^fZNhpXS_?vsjV|;wLyPF|qU&q;BZyLCV7do$qhHN4Y zj^B{Q4EOv?$miyOjf{+k^#>WEgI~S*C>RVWv zA1XbW?`n!s=m_J5)Dc=lOk9MteZ6o{aFO61I>uHE+p&TiaISl2=_qy#jU?yb^92=~ zEb5x|XV^4RjnJ`_za#Ip6-qDAdi{SdlE3#R!Da#Rjs120uJw>Ezaw00xj81w4iqx4 z&I-LnOAc#6-*DroFa?0eNZxK&^*U}f;Fbm`s-8H@<#uQj`q%X~UuhOa6ib?=jR4mj z>bf9QR&YE$D%^uj{YNncyk-m{4~7!BP&|t%WfM${>tnC}pd8RE4+j6*xVz5JPuc(G zd-9qapqCc=HV`%GK=DipH#$C=-Lb`wW9dNyAF+lc${<@3L z)8Ag`*RT-)-0l7JpN4WE)&syj#0NMAU^@*+{tN=Zjb$J5|T7GV%q?l3LCLhw02Ds|jtruC-bWI4DLYw7F%{@C_1b3_f4?5b_-r+Hy zOoJb3A4u`Pg*MyGN0C_}=b!}+|Ebo2$_ll%zS@L)Ak?=QwHL%#+85OY52m(^o17;& zZ$x|mr|)qbMO`iV2AZq>9mhpZ-q-YEkK5mF-@OFfHwlwO^>y}*Rd^NTCUq^wDwHSi zSn|N~*kj!9#08mH&bYPZC8mvl-QUB7WaS%3f8#Hgy~Wi zrDZf!K4g01;oD}8+K6`Kdc8ih^GJ!#eGc@D>%UE%0)XqYkxR$S!ce6=&yAl;z`o7P zS9_5oZS_Sgry8&9w!|59&&0Qjzg{bjX%3HQK0Q1(*I_hO&_BVy~z%y zj)Hqw8NG-Dz{w6EkC1D zhS1?6s2~0n!HSQHq(9QQyLZrO6k8BVU`h1cYJs#j_Vh3=C&E3%l;zVk6MYcy!=)^N zZF?qtIsQEHXocGj23~)RXVI{YM6fTHi1a3VSYMETTwK0^zi!^p1~xX6|Z{ zh;fKD0!|-VdY&IeW4C#7_4*5;OsWLSD#{%~k4l`F8+N4x1>b19TTXUqoXF)Bw|yh$ zy@qdEr12Oi-s3b#f2ZRO;dl$_l1;FkVf6zOaGO1KM_Oa2MA(qP6e6H`8!|8Ey}zVQ zPfzznhS=jVZd@Bnm_G`UAhct3_GCkcOFFajJ4 zGhh#^O1Tu2z&(;w1h;wGL^!!$oWl?_>)rRu88xUWh}bYv3>zQI>wrfO(dq|0%s2%4 zDk-3nss~U593E~Z$o zSAebw)N|Y2;^wzNbOSL0!QfWq)zpJLn+yV}h>(Iexq5L&U_Ty;OI^oKA>5-*X+1wJ zX{@l{+%3}xI1P`sxF=QM-De-TMdJ^XUzth#3hk3v$e(ThrU`sT${ab7)6|9_zw;pC zJDadcxm;fk=5R~v{KicGsgS+SpQ9%ynEMVZ0ZI7cQi$uTdc8*{wH0YNA|f9daEDh2 z+>*`!w?!kts_G0hG5s|+CPzOcQhc2>9#(W^D+b}77Zh>BH6RqifQsSf4NK!6kMEN& zgXu6@zK^oSlqdv0Ac(8l`1Xz!pC}N}c~BzChy0<7(__{jZC0kBwO~L4nTPOVd1fWQ zeen{xe}`P8;wboUFJ15hcO%f6`t#OVZ)$49_dU9_Jb{KvgU{pjBso!um`b%@XGF=$ zP{bDQ8hQfGO-d}pptIG$67F%Q7r8xr4{%$S~>Oj6Dy&AWB+NNaNoacW@L28hZXy zIwIT?K{=N+Ky5OD6|yWJ9~ZtAyxKRM237g9ZeMfJk4rC~dbsk6n^w|!&vNv~v1~>q z{-(p}Ol}j>1o5@2d80+$ImfqDtyQ&3AfT!0V6+oNKnG7 zhe&F~S)x-_5$I%V=$hAqpC40#p@M>!FQVPS z{Ar|nu(%D1GSSNH{+T1ICnFj#h?U@lYrU}xMtjkWYkcyfkV<_}mOlX~<%8FzaAK*>+04wYfz$gLnZ_1hWQsz2DV}ztL*A0t})c%&2#W-`fk25cs4X-K1Dwk1FXwM_Iw` z7ydk|1SVifTH&zcETq|3!*$@mIi^+zgf+F@v>V>X zGub>kzWBj#&O5q#mX6DnU`$&L&(AN$4B`q9OG1DtW1CfpL6nTDXyXd(w>gCl+t*Z7 zdtR`7W@Ns}&Fd0N*)%o2WpU#gP!Ht0D$d1tR&wTsC8Tw~ zWDTJY6Q#uyHc+84>04&^XnwW$VR=3djErvqj7}J$7;rVnmsL4IbA3?dN`oK%pvucD zC4?(zalkzrly}437c5G#0vx{P?~P*ORvA&so}b|-;4J!Y|IKtEDU0cPI0o-H_MSO2 zJ}l`3*rAak;Swv4U>$*E<`guoB|*v%YaTHYlexAu6M0@4^e9I`yEf;z1f&v*-x7Jr z`xzu*ZpVIg$300S+Dv)%srw(iPrpbdmRA+7oQqjacR0@kvqPPoot(Vko)bE9l;(iN zXsqsB*<2ac>Gn3Jg3tmU)$KtIpb@C~YDBbK9kl(fm;I-hRXn4Idls|1!`&RlxUkd|D#Nqew1LGVO$Vgs4&<~-lz=cR^_ zW<=s@7t3%4+E7x=2iXp*P_bexBG(toc+c6VR=!sp^0tYdcYTH;!;(wI%dC}GyCUxv zk2il+97tD4f19>s`VH>+%D}!B#!XyfOHCpvJzB<$x1gEJLgDMo$e(6&qS`f-FuBzh zj6(w?Vv5H;^!7&pv%<2RJx#!a)Ps5O&Ct-awa|GSzI7jcQ1CIUq4TzlL4{@oWAt$B z6}M^jF_yMNcO?~9E^#=WFoErDmjf@8^$uYSmd6^aRP9T}U3nd^!mZZHnH(JDU%akx z&oaXWPL^p!Jhg_)bl74lH^U$2276aMOJ(F(m!6{XvRO?0^+%^?4>-bgxesQ$is`|} zj`e21gBteCzO1t{3zBdg>k07SlZiuD2+UL{-1w78>i=N zJ_KH7qV^+JRp@oLk)kjGyMtH*Z{DseDC|UuKUSpGe6fxJyEPH5CF}U!_T0od0hAWJu^p=h151}R#msTRM37}PSD0dBcVehN~_LPng6Q*S$E{efjV(e z&zJ^qCOflmD2v46BetdgOkET>eeQajeH-*fLVQ74x2GjJsHtV|~f0YNj7I&k4 zHe+;GeEpntHv|HT)DHO0nd7({LW;b$5|}bX5tk2iOaZf{ANX=5P`LJ-BD0ai(|j<`71dNNq8}Xne@NDf|%Rb@E@Fb#m#vznaJo>3i+i#6ksN_NDR}g zP{C&)?!=?mq<0$thm`1&em#atgzEd=kMkRD)JUQHnX9KY{AeH+;#iL4^O50w?5OGP zvpAEBrH9P)OB#w;>n7=s%RzBMByq^}RxHmn0Q0UIW^p7)UFep)Pl~W)O!>v$D^sY0@OAGufVWRf8(=Kua6KyXsQ?1Rw zXRVa&`yaqZb!C&+x3LHraE}dbyk?bn`9#^f2PjX^T;%xnrHD!JhPr$uKpVn_@h2ql zzN7zz?XUqz#*+C0R34o1t&aZXj#O>ewjvdW7b>fY()JK^r%!B0F~4^cmR`L+5B>!Ojkm9GdPXUc;&UI)6Ds;cB_htUiQwm-7ni(5 zsmMUovdZs^r^#}WkC7`Bw8W1Lpl31pp?K--i&;+U7{2aFX>(2)l?}Xg3UzW7X$aHI zdVIENud3pp$H0*$W=scAKZF954fmA)f_q%4jg>3BG0(jUfUjdxJqXLSdO$2vKY*K;l zrnGg2ygEITsmgS4P_>Jcy{`h8J|IOugMSD-k7v$W(1_H^$zzI^^2HJ@JHY zYx>kd$h7J&^B?l@SsWwYvNjsACYL#F@-zLo?aB+UrxTcuIKte$=)+GePbV!6l9T*jUU&~zbNa->}HHcVHSVQUX zF;&F+1Zz)o)j@F^rRhkFs+ge}h-+g!6-|ab8X|Ay2;iR2grm(+!f|tMLuzlcfucOb zF_9UNi%dU6y3!?lZKH^djYgFn0r94DzB~p>IMmPhTuzimBgFD64k9mxE|2t~8e8`K zrU(ClB8x!Rm8cR8M)bHyV#0Ubbms|Stc2w8Uz@8xIkJc%*OapEqGdcWT#Rz<1ndu= zowM4qvOW$p zC8a>js;OkzIFN96JfYQCy57AU#4EKeHRd_~C{=AYi(MYw1PiZ*Tk=? z#Zb0=w5;j36b+Hf@j~nk)xXbY#7pt6vq}$|_5*t`>@dVhqXLM1nM`Kkp5N4?-4KKY zYLPUgz-4PTeW!1=8;OCscc1#DU^me)5ngo}41wT0qiE?6GJN44_fPJ3nYy|x=rN%2 zQmMK6-|vvQk7WF#HVfHq6Vqnim^S`s*xDp&H=A;*J9|yg5PMS9aP`9MN+6wKe&W@; z32X1QEJd~Yj(VwAo-Jod=bUM={+*La0;+_8;~sF&Uk0{eB+xpFpLx*dt7}iO|LnWm z22X*st0@A-qDQ2RtSG3LYp;6sad+0;TK&+DPVnjm!IdyMb#Ek}oGQiyH(Ry1aNO*U z&PQthyx6QXzpYw|zktH#XZ7ca^y)pkb_6cHj*O}FZvCWt)?A=vNR;H>Ua)$R5BFYO z;A0)uo1>sHQ$e`rl;SRml7BAF-Tzz2Wwk#BdVO-CwKCeF7h2mWOPZJo<|rX1&hr!- zx{{}hM^=4|jQCb`bZc-FvZ+BtS~w0>Wu}w;vx@I#NtHaWTiAvKM+c}-XDj>}t%bbx z`}x@v(H5t62Sfux>zw7iI=#ALT=7>$PH>2c>Fne?ZzQ`$Kg0?pW}gPV2@bNjfP4N> zJx76>FASr*?9sTL3ZrDHvV_ER5-dTY(CS&dEdNHCaai|hd|P1&r&Rc!sH4GVmr5Ju zRF$qFVI+5=0;{;C7H3}-i8L{4J0j2Mn#8YXO+Ss!j5N8;7|-yd$CxX5DdW7ctoasLrYwNnNJJtdSS5N4qRGj+Z9)N+j0fev`)yCqwJ+2=) z>EWTmXZg0sR>$`>OPuI8ZdsNO9B6T^ng~v#yIe?O*!C1YTtibLx-q(aUrdQ4xHCI@6lRnEPw5N%8s^tFT#Ug_dVjRf@frWg?0*kDr2^Z`Iz8bB)72mvpwoXjp5CZpb#*2Ys(in+-2h$9bJXVo=?r z5?|g_OW$pK`7A&5$jHp`zD`x4&b_j1OQ7+C;XXC>p2g6j-MV|n2P5{+O>agdWV%y5 zlVF~%?aMd0fBqdW+FNrZ&{KX$qN8KWHuQK|4%Uong!&p0-p;2Nvj}RDv)fwnqPV^` zgAH+jFd)_-YivZM_NY@|`Ut6@jnxY`#Fm*)2Z=W!lv6^zw}pt@4qAfcUnsZ{$3pWF zsAzghuv-Xngt!sGpjoYV*vQenb5Q z39(J3!4wtmy58K|MY53Aa|Jktd-m@PR6}T%0HA+%keW9_-^|Jl;j1_c=yh*XF9tgB z;nLq?5JHc1gGGtp!#vqmY(x}^^49L`+{-54sH-DmHR zF%Cp|v^UX67+}tj5Bs!&B$=8$po<=Ntg#R|RZ?Ta}B1cqmA zZOw{jk%;okiO!V`?)+5zTWQmkj}==FbEZyt@8Mn73^bnlO7UpQc$D>8dc+I*^gk z@Y(S8*?!c3j>0ZqHNKhmD(XB780)RI)W2Q~w_QD?=1zUJ@x}+~R5_=HjhqCk9&2qy z$Rofa%?)yRQ(1ifCP~Esmj|D?E6F>RaT+pZV|d}-pLFcQ^GNyb02A(EX41Qg9ad$ZYG2y$veqr>Kig#37!S7t1&N&7|ia(!aWe$^Qa>80i=&%Fs*ICDQ3ldLN1OY zlh!dnR)z7^U?`c-^b~Gb`c_e$?Z-y5+xN!!2}KRl1dD+$ZTAs)J|DiALlJmemeDmm z3gk*1Iy6O+0q1-HkW@vUgamKCAXmdToVkSw7tezMh3gj5O1%(dEJZUuQ|1zUHK+iV zHuf1Rqs5f}vt^Pa#eeq~xU{-1Y36-H@;dmkOz8s|fjz?D1{+g=aT{X9A9oWx`@)VmYFCi1V+%xj_iB+<3(`Zp1Q+F-M-ih27TkQ(17C=eo^;&fJ-0 zpuC{pD-;E}%U|6HCOnDrBpf08UgwiaBV97wR|Z(zIJicJdls6P9SpP88( z1KRjT6#rf6*ki|SK0n=Nl~BNTtStGy+a=FcCQXOmX^eBMnF`n&P~svdS#+@Qg3cZ> z#WcysL+Ec(`Cr-t5B#&w^8|46t6uOjX=rGuRaK00>12@rz~03*7k#0SsY{q+^QNL6 zQI&!wwfv3zM;9FdZXUZX(DSqh9duSK&yJHNxCfiEPqW)%!WV?_l>>bwv0#6~xkOz2 zCjjuWLk9d3vwH6O?K=w^iHTY7O?f()iLhY;vUT?(&T^Pm&djJ*tSES9QB(ClzV(_J z#el5e?I6<83Wspu@goD@lI5Nrh*aZH006mTpCxHuNXsx~7j{KrYl)@sEi5h{r>#p&LJ&7e@1zTCWL7l@ zH)Z9{w}PchYKBv8(xhj3Ic%X4%ss{=_z{Ejw6i7KkE8lf6WY-TtLtDJQ~6lC9Kknb z^nJS0JR*#fIx;Aaw!ia=E+(oXY5mws4TaR3LOVy~9f>k;K7sf5J|BFxXIH{|qxi3$ zz)ndXR@ogtm{Z{D9WXeBa(Q!(OmG=#2UhCR4AC4BJxhvu9t^Om?4v za&0vRnRWD%yK6_Lt$bI{REz;r;yE4T--e6WZkM%!G_u4{czz@(J-hp8dD1>==Zhh;%;v&_1zbnQ9OZNIVx&~KbnyLF&YvWwU!V{w@* zjh=ut!AaUMZ_Uj2 zS9f>v+rNu0bHNc;%r6I1B`eGBCjvLxXMpDZu&)nO7oyDmf^V~lElt_BFHuVz?O^-! zFvRDGf@^E22cl?SQ^Q~T9h9uD{7@eJ_7jJ%nta=*^S5v-<668y<@bRb`HeN*_NEL_ z^XiwG*03*3G5@aT#p%{url}<%m#UP>%}yp}g5yhK+NopZMy_g%85!rqnnbkKoX|gn zwYr&DMoNmpW==V3(V^>+NNoI{WZGPX7B@#WPE!r@WWN>A3X8apudQwMv{{SM629)B zas}=(J|!42DYo9*q@sg9*9sN%l!hD)j|2@^r7m9}a1&&TvY=(~O68#Gy%p9yhaWE? zN;T1OLBGN2*2qo3o6*>v=n^XHIjtFgZ;eOfm^CoG5s+`q zJSZ+`U(E>ekRQ%iG|-$|5O{#5mmG2+(-#}>6L6%3mQVkyhbxn(@5jg}Vxd^U~gtB$uaik{?TJsq_~Rh&{d#@y5buwB_Tb5+4UxW3pc4 z%F>M+lIq#29d0&jR(mNrb{yXb&yOYT!fnd9!$Q-uyZjoi%5QG&;-~|dWvwqsbPXmk zhA`cu@X!pq)%|`u;u*lu!~vjHN-yEBeT71++PAN%U@5R#jCg@XI-+2p`B?JcH>oe( z#V5kyfy(*Su3m%_L17Mar(6xM(9y7-3EslMgk@NI0(9ie&)qk&K#-W5Z4_M1bqmQs zyT8|@68vzF1*KD5ZAP9(1%)?luC_=bVDsX0Bd`MSK5!_2(P%zAjB@lSAO3nBdNXBc z4B5@DRVL|TKd2AL!Xy#4B9q*nhxp18fI?*2*V9$YJoraf^8A>XMPyeb6Z+nN(asQz zQO3ezQ*_rSR+xlh(Mj7iv|xU0J~KLgp(b5AFK_hN-&Rr#02i=OtPnT2$BleUx8F<> zzhbSQs`%X-@+dmc$4=MhJAT>{1j$A;ZS>k_1nDQmMdV+U4g`s;8CbMXrh1i#NCo#7 zkrtKDeMV^=sSv|_YJhO1Nj`&WjWU_l(slvzLH|XXBM2qxgNhUc0d$}DNE3q@og77c zek@EgxQ8R!iGci~p_it#i55+uAXXvM?zPa?zy2IeS*qG@8aAKNg88LN#Kk#^>U|^+ z9xn-jhVuKQ#FQd)%elrYR9=h9T25$277hN?x6#+rev^sG?>kItDx>VRtH(%ur6UBY zJVrZ{&2V0a6-(Xd$7alw3^X%2$#yiAoJsV zNL!L%0e6L>p1Fy_XBtGkVF82;$LNf4jJBk^hW$-T8x(2~&XmYSDjN7M&r!rv)R;A6 zRL=zNfl-kKjZhK~`^KEJnJzB5*G)}h83Ek3ECfqk)F)V2TEJ2wP)aBV+>YHEdiDT4 z`3o=bDUG#9oY+q4SB?Ong$ECt{w&%{(Wb~f(LU1}f~Ek))sOANb>p=nOrr~%wMmBAOfZIY*FLR$KQ33AxNWE~|^ciO`$XiT| zjq&c(yvRz1pU@whX*hBy{UHAq6Axw4PF6<^R25okg(>&gCJxHOOvj0dDKhm4Vd?s? z^kALal>k|q;ScD16_3&D&T>H<)vJupEbx$CFX?aiYAW!h{c$x3S7JV2N&=D>JI|g~ z9XafA6$T6$JKIg=eA&}8fgRUxKQ?Eika1UVc_~8c~d>Y@G45|1OYWSy)65vxb37JF(n8sTuC*&~Djl z@$h9YnYb$!4BjY-VWMp*k;=t-EaK>pKwBKXv%6HDMP1()B_z8nb-l_w&606$9l40U zc1dy7pKE(N>;FuI^RuI7un^dFyOOZS*_Cru!cVBu!h5;4F8&+v0`A!-@258t4^mCA zxqrs_@Tj1WSE@bpY_b$Jk}GDdz_wQxUmj=iZP4zgE*G<2i>NLo{QL}sn*I>0kPJDD zsB!f5Y>YUV+auYt1C@5Eo}q(ucl)A6SR2#CtHL{C!C{4gBc6HE_6M?j-effc7qe0cpftkdsD`QOdHjn779^84gCaN(pqHKODk9-&hYDG z0drR85z;?#8jZ~XrE2l*Fry^fYBKo&rpi0ARH{V%+)^4>0~6VF9ZsZ@?;7Qmes6uAe91Jk-lWpub-?As}ur^AKxf zQ+u%?!X}(b&oV~nHqR#ygC#*zD?D#6U+fFIcey$>EIkGTt^nY#={Ix#xcFydGi$Y& z6e?E;SQ@a{bb_`Yw$YPNn;tKdMzmEIr%+R>EHDs zKzlOvh1W02Kr?ns%;0S^BNT%_RYr@Csa(U>K(QR-zq%95^K_I<;pTNYPYwF;0>Z1t z)27Ep74kOYt({mw!zIfzNvqtRyD#PS`rcEGZMQ2M{Mg8~kKMJ0u3O^J54Y2z?+(3u zQx}B~rcJ7vW)8I3a7+NX9=#yo5I8x;jqmiwsfFj~5$c6q1L{jazWi2XbYGV$%%$h8 zEW~3&ifFLgTU-C9;gBc$g|}VO<@)+h-a6ff0@!BPU(-%uV;PS5cftFS5%{?L#@TfZ z)+cv)uZTD#a3c6n&+*O)a4qhI-(rL#hgeWJ2_7i t!T;~azy8B;^8GJp4U_6Wn)`oD_#Zn7_XIrsM?o1rAOC;6^Z#e#e*qGeR$~AF literal 5121 zcmbu@c{o&W-vID4LyTD%j9m=I8X8Pe#9*u;YuWdbHA|#^3R$uiiA1tRNM*~CYLH#o z6^g-=l2WpjW#+tdJlFI7_gvR|UGM$JoH=Li^SQs@GuM5e^F1aMHC6y>QnofW2FxuN z0I-;OM4VG6Dyyg~s}PC%zyA9JT(@ug-(8cSa~GI9%>Doc0C>j%G?tfFKtx1LQc^}< zUP(oTq@kgyqoZ$VYHDF+WoK{i?C$R6 zs;g`2>zi6zUw3qM^}TyHJU%}8@gr?v;mh*!%J1KsfB)`*(BRDAyqLkMsX#mbX#|$G ze}Cm2zhm(J)qf5%x8S;5Cv(GM3>~w|xSoQA#r+C>vl{$aXP~_&SWh;(8oPM#jcD=E zq7FlFK@A*h6bHGg7(UJ`7e{0nJ@I0S>^fTw6Of~cG(OOEOIHw}`NbLwcYL993c$~Q z<@DJ-`LOvjg$71`)(Sd?;RW~ZkV5Alwt#epQb7SMk>te%g6oDvWmSag0mv~Eu#qIf z?wtoZ`9Kid4q?sEeTy}28{SD*I6+`x;ksQ$oGYnc<)_-sd_S}K0ES1P2c$dWIyy`& zmc_WxD%5LULXabmt5nv4<;ekPLgj9y6`G*285CyGENrsy$wLWgHIK8k%#p5adR@h4 z3^ye}9Xa5*nZiAs2+=y`S5>C$7P&1g5hBW7zO zR`O{q&MgH}MUaExI)6?Bcr(h{9Mnq)wpi+V_(an@Ee;l(w7#JB?(yT4;(+&c2)*gw z9|CuHtWur+Y*dQ7?dCnsa>3?qMQq?*&5YUi<3SfWY5{N1hFO4eaddU@X?QGknwA00v4!~ZxvkH#op1;MU1OELPNbVBv>L&zywfL zDOSZT;FSlEBZUH7$v}r0YwP#th0)Bv{yR^KC6M%8g{<-0wmks5F|X-{Nt2XNvM)BZ zZ48((7@3Z)YILiA-oPBDPz}tG0_-4171d^F%->-Kq~I$C9+@~w$TrEh25Vq|8ysQqLQM)T=QIOpP4iAhfedQK<4YoHxjyC$QDg|RN2pqml z16OBHA{4br{ZeV62DBs1VXrcHz#uVXwsqM^>k_B-{yNbvDfc;75l0If}`j@K(Zk5Nube{N^9MV@8tI1l^_6 zJO4#9k29-MGd`(IT^-U-I3WVsLXJ49Vr>#O5DCl|73=3uJl~A04!*@dR6XHgTEz1Z z%K{3&xj5y0sKed;=(si<<~05p#`~2Ow|EMHJ9wAGb|Tdh{~U4-p^K&&TE3LfqJ)WP zaA&+rAB{&xFBGjutai=p_J}i(HFxS(D$J!-_;PG1`F~#?Hpn_*R5s{I)^|JE3W9sG zbXUkxL$7X4W%!RLU9Mw26ksn0VC4oZV=9mR+{Wtqmg-OM-8p6&)Wm5{Q^dGiA0a({ zp&P(ADY*o$6>HUy_DJIK$k7*&a~z|QGw~XQCEx{+I=DTTGl+t)h;zxH@eESUq$8H-Vs~&3Z=Ca4#)`_eZXC-qSrjKwq5_}*>5S@{@n4$T$2>12uNb83u ztgb)6-AI7<{FA-0WS z+jLa9UT!(EpESQF*<%p*&|qeIp#?i$@#jj;cO7u*`2IXSj2;gCqRpKK5UAi{wm2TH z>%39-%B!;`zWL_6>aSKphRVUUwWY|kna#ZqZku-Kx7|ISDk^kdR)Ys@+v4CJC1HP_ zs$(Q%DsUc&0dV>&j>jKTzIbJdtn+s%SDo+UX|^3j-Sr&sduLGWuXZe`-(@4uxf+!u z>QKpCV7%?n#|S0UL;$%B<)wwLRgXblVGWMM(60rLWpUh(PmOv+3CtZV2f<878-6to zugvS56T{EIMT2xKBIdzpRS^=6G|m+ExXAG->Q~jQHsnyyQ`#-Bx#O#hExzRo;^jqU zQYGztNlkZeW@l$#mni2USC|Ze;FNhK5Mznf<1ss9NG7)nBRV503z1hWKW5L)Cqd3h zH2JnY%!CM_osR3Ay-7pheH`M7}EB9aK+U zc&+~HR2kwiZOeK~KRX|C?D%T(Y%zI}09FCcp-zCaO!qCxJbf!9)!}}_jQbxUz{0mA zN7o}()>S1NKj%4jI4%%zKAh~=&W8%t#1QQzQsk~f&M`r2M#0&B(7!%!z2aX64`r6Hc@v`_RtZMgd zQ^GP}DG^vC2svz=*x3a!cqfS!4h!*nL$1h6AGiED@>qAqTE_8>h$6g`Bjs>222`y89RiN@MVxK8MvXkN)k9@V_K-ViiXn z7Ehc|4W~wo*MH)RxkOzA7kPC49CTY&_GBcJZe05#zZx|>N3d1o#u zmKax{N0!ul#FU4&6Ga^v-o%F`SGmva$Aiq=8VAm^V|kLFFJH{b57*hAz4*fD(L?@6 z54{|mEJZnV`J;Dy>tMd7Rd-|#d|&{~tKw*5)`zcoDiL2^gJ6r=f*e6~&9zxqk#52c z;WyFt-GSJq;HOIhAl*m5{8;&u_E9Rxa13ry%;|%}hT#&%CtZoT6C8r~ZXCbndx5Km zV<_Im^)BS#F@mR;U5&fp8X%=PmI|(3rGnsJJlD{pl}_I5i|L!RPm=EN(XhtZbL6%k z1`ZmMS;n_3K{uPuI+N1SA2+5Gp8cskm8t7hj%)cNa+8c; zkzINorh<%bP+Y_A94T{}tN`pfB)~T}2j7j%(|$^=rC1|D#;HZfQASVp&%DgkaGfqpB{8I>9)_2^l{3V$E@z+F4%fQ4FKp##P7 z_Zw99=cyXMKe7!FVLOZHbs5llOWpQ}iqOYmm@|nM6ye6OyF0BpoajHX^|-Cgu|n3; zRK`dspyw69XPWXP%HZ@}n}!4M{pSQ6dU)aeU;R#F)|SJzLfRFbvddG(r#^&-06caD z_x73?#RGiowDq%orTR**SAK-fu}T-VR2K3iNW8ZhW+avLa$Me@?;^aCPj_5xOqYl1 z-jCnEuP2Scf@uPr>_Xr;@`Q8Yv{x3GR6yP2MzYH?Vi%d0rJuIfzQ_wuw~v$EZN)BZ zmIC6?b>`sGRUdRk4SKH?#+P$GcPm^E8NQ;WdkeHNEL|D^H59NYF6#PC8S*{-0trCDJT)92cw=B4| zUk^&5SD!Wm+JcHb^TYCo#7P@bCzuI?;Izjaj4OaWMyV?|&9~}FT)h6!@doycuawhn ziZLr)Ywv7ewG>@^Zh+p004N>>bo1=geNPHMbKz!`$3F1}JJ1B6Ezf_ed)mnJ zb9K&NJ0y2=IWy4<6OC^XHcgdqV~)hYN&3@!9Keli$WcX4KKL8?Q-#bo!$-inwaNw4!uZn|C2Y1@YXRaxrIG#T+H)VdH71f@rTydf7Opk#$M4U> zn@$Nq+5O&sEjY(2_Uls$aVJ@d2f9I<+h`ke1b84I!_e$%Pp z*g~Vii+|d3{&jI&X=t(BQeHlVJWx~G@hq~1!hG<6;;}*I-nInB-V$j}7yRlEio2dT z9v_`a;!B>c7F9gtc%g~6 zDZoHCF8NT4pqv7$c}j0zzC^DCuZ#AlJFIx?@;4pwk1ztE-WAr^8&|KrntApSI{Qg* z>`WY@_d2{R=9L)a1fo`Rm~|x^t}iDhm{KOH!RDTnCNm#y$ubit>|)Ls+N=jBPJIOF z(Z$lcce;k*K6-jJa_Q-sBz>3XHvU%ZgDFVJi9{2b^=XGSa3>}8fkUaOSy<)lv-*vQ zrAx&hb--xwLYc6KJMyV$=8aSoXCE%MkzeU?p_|V=K^%)p1gqoOuQ`)DkaHas?6_Z_ z23^=OlW+JO%%uU*5z-Hu+=GfS?Cc0bc%5DS5 zsEP$xa~Y6x3mpvAxz<2oMk43H4&AumV|CfqK%?rVM=5)3{4WS_fBM)sMkqZ=L z^oRv`ba$q8nqZtMQq7LF<%=(z+poKGQAy0a-UtHLrTxg3@&|S*x?VR%*B+gg)Uv-1 zUJJW++5lnccY~+0VUc93>7nm%euwld@XRy+$J$Y5a?4pTAa?(` int: result = await self.execute( """ INSERT INTO device_settings ( device_id, sleep_mode, disable_time_start, disable_time_end, - timezone, volume, brightness, disable_weekdays + timezone, volume, brightness, disable_weekdays, `signal`, `version`, power ) VALUES ( :device_id, :sleep_mode, :disable_time_start, :disable_time_end, - :timezone, :volume, :brightness, :disable_weekdays + :timezone, :volume, :brightness, :disable_weekdays, :signal, :version, :power ) """, { @@ -39,6 +42,9 @@ class DeviceSettingDAO(BaseDAO): "volume": volume, "brightness": brightness, "disable_weekdays": disable_weekdays, + "signal": signal, + "version": version, + "power": power, }, ) await self.commit() @@ -70,6 +76,9 @@ class DeviceSettingDAO(BaseDAO): volume: Optional[int] = None, brightness: Optional[int] = None, disable_weekdays: Optional[str] = None, + signal: Optional[int] = None, + version: Optional[str] = None, + power: Optional[int] = None, ) -> None: await self.execute( """ @@ -80,7 +89,10 @@ class DeviceSettingDAO(BaseDAO): timezone = COALESCE(:timezone, timezone), volume = COALESCE(:volume, volume), brightness = COALESCE(:brightness, brightness), - disable_weekdays = COALESCE(:disable_weekdays, disable_weekdays) + disable_weekdays = COALESCE(:disable_weekdays, disable_weekdays), + `signal` = COALESCE(:signal, `signal`), + `version` = COALESCE(:version, `version`), + power = COALESCE(:power, power) WHERE device_id = :device_id """, { @@ -92,6 +104,9 @@ class DeviceSettingDAO(BaseDAO): "volume": volume, "brightness": brightness, "disable_weekdays": disable_weekdays, + "signal": signal, + "version": version, + "power": power, }, ) await self.commit() diff --git a/talkingq-url/banban/security.py b/talkingq-url/banban/security.py index b4e0f8c..47bc736 100644 --- a/talkingq-url/banban/security.py +++ b/talkingq-url/banban/security.py @@ -1,4 +1,9 @@ -from datetime import UTC, datetime, timedelta +from datetime import datetime, timedelta +try: + from datetime import UTC # Python 3.11+ +except ImportError: + from datetime import timezone + UTC = timezone.utc # Python 3.10 及更早版本 import jwt from fastapi import HTTPException, Request, status diff --git a/talkingq-url/banban/service/avatar_storage.py b/talkingq-url/banban/service/avatar_storage.py index 1163fc8..a1ef034 100644 --- a/talkingq-url/banban/service/avatar_storage.py +++ b/talkingq-url/banban/service/avatar_storage.py @@ -1,6 +1,11 @@ import asyncio from dataclasses import dataclass -from datetime import UTC, datetime +from datetime import datetime +try: + from datetime import UTC # Python 3.11+ +except ImportError: + from datetime import timezone + UTC = timezone.utc # Python 3.10 及更早版本 from pathlib import Path from uuid import uuid4 diff --git a/talkingq-url/banban/service/device_setting.py b/talkingq-url/banban/service/device_setting.py index 3052be5..6e0a289 100644 --- a/talkingq-url/banban/service/device_setting.py +++ b/talkingq-url/banban/service/device_setting.py @@ -38,8 +38,8 @@ class DeviceSettingService(DatabaseServiceBase): volume=volume, brightness=brightness, disable_weekdays=disable_weekdays, - signal_strength=signal_strength, - version_str=version_str, + signal=signal_strength, + version=version_str, ) finally: await db_session.close() @@ -79,8 +79,8 @@ class DeviceSettingService(DatabaseServiceBase): brightness=brightness, disable_weekdays=disable_weekdays, power=power, - signal_strength=signal_strength, - version_str=version_str, + signal=signal_strength, + version=version_str, ) finally: await db_session.close() From b3687fbc71b5aa70a0073da5df1e0d915de01a35 Mon Sep 17 00:00:00 2001 From: HycJack <772403255@qq.com> Date: Thu, 30 Apr 2026 05:11:58 +0800 Subject: [PATCH 2/2] add code merge from test-clean --- talkingq-url/banban/dao/__init__.py | 3 +- talkingq-url/banban/dao/binding.py | 237 +++++++++++---- talkingq-url/banban/dao/child.py | 4 +- talkingq-url/banban/dao/im.py | 30 +- talkingq-url/banban/routers/bindings.py | 86 ++++-- talkingq-url/banban/routers/im.py | 6 +- talkingq-url/banban/service/binding.py | 98 ++++++- talkingq-url/banban/service/child.py | 21 +- talkingq-url/banban/service/im.py | 274 ++++++++++-------- .../banban/service/message_audio_storage.py | 129 +++++++++ talkingq-url/config.py | 1 + talkingq-url/handlers/audio_file_handler.py | 43 +-- talkingq-url/handlers/mqtt_handler.py | 10 +- talkingq-url/services/card_service.py | 262 ++++++++++------- 14 files changed, 847 insertions(+), 357 deletions(-) create mode 100644 talkingq-url/banban/service/message_audio_storage.py diff --git a/talkingq-url/banban/dao/__init__.py b/talkingq-url/banban/dao/__init__.py index fd6faa8..793454d 100644 --- a/talkingq-url/banban/dao/__init__.py +++ b/talkingq-url/banban/dao/__init__.py @@ -7,7 +7,8 @@ class BaseDAO: self.db = db async def execute(self, query, params: dict = None): - return await self.db.execute(text(query), params or {}) + statement = query if hasattr(query, "_execute_on_connection") else text(query) + return await self.db.execute(statement, params or {}) async def commit(self): await self.db.commit() diff --git a/talkingq-url/banban/dao/binding.py b/talkingq-url/banban/dao/binding.py index 9ef916a..73d5383 100644 --- a/talkingq-url/banban/dao/binding.py +++ b/talkingq-url/banban/dao/binding.py @@ -4,13 +4,21 @@ from collections.abc import Mapping from datetime import datetime, timedelta from typing import Optional -from sqlalchemy import text -from sqlalchemy.exc import IntegrityError - from banban.dao import BaseDAO logger = logging.getLogger("banban.dao.binding") +SESSION_STATUS_PENDING = 1 +SESSION_STATUS_COMPLETED = 2 +SESSION_STATUS_EXPIRED = 3 +SESSION_STATUS_FAILED = 4 +SESSION_STATUS_CANCELLED = 5 + +BIND_SOURCE_SESSION_CONFIRM = 1 +BIND_SOURCE_DIRECT = 2 +BIND_SOURCE_SET_CHILD = 3 +BIND_SOURCE_NFC = 4 + class BindingDAO(BaseDAO): async def get_device_auth(self, device_id: str) -> Optional[Mapping]: @@ -21,7 +29,7 @@ class BindingDAO(BaseDAO): FROM device_auth WHERE device_id = :device_id LIMIT 1 - """, + """, {"device_id": device_id}, ) ).mappings().first() @@ -39,7 +47,7 @@ class BindingDAO(BaseDAO): unbound_at = NULL, updated_at = CURRENT_TIMESTAMP WHERE id = :id - """, + """, {"id": row_id}, ) return @@ -50,51 +58,28 @@ class BindingDAO(BaseDAO): SET child_id = NULL, updated_at = CURRENT_TIMESTAMP WHERE id = :id - """, + """, {"id": row_id}, ) async def _upsert_parent_child_relation(self, user_id: int, child_id: int) -> None: - updated = await self.execute( + await self.execute( """ - UPDATE parent_child_relations - SET status = 1, + INSERT INTO parent_child_relations (user_id, child_id, relation_type, is_primary, status) + VALUES (:user_id, :child_id, 9, 0, 1) + ON DUPLICATE KEY UPDATE + status = VALUES(status), updated_at = CURRENT_TIMESTAMP - WHERE user_id = :user_id - AND child_id = :child_id - """, + """, {"user_id": user_id, "child_id": child_id}, ) - if updated.rowcount and updated.rowcount > 0: - return - - try: - await self.execute( - """ - INSERT INTO parent_child_relations (user_id, child_id, relation_type, is_primary, status) - VALUES (:user_id, :child_id, 9, 0, 1) - """, - {"user_id": user_id, "child_id": child_id}, - ) - except IntegrityError: - await self.db.rollback() - await self.execute( - """ - UPDATE parent_child_relations - SET status = 1, - updated_at = CURRENT_TIMESTAMP - WHERE user_id = :user_id - AND child_id = :child_id - """, - {"user_id": user_id, "child_id": child_id}, - ) async def _insert_bind_history(self, device_id: str, child_id: Optional[int], user_id: int, bind_source: int) -> None: await self.execute( """ INSERT INTO device_bind_history (device_id, child_id, bound_by_user_id, bind_source, bound_at) VALUES (:device_id, :child_id, :user_id, :bind_source, CURRENT_TIMESTAMP) - """, + """, { "device_id": device_id, "child_id": child_id, @@ -122,7 +107,7 @@ class BindingDAO(BaseDAO): unbound_at = NULL, updated_at = CURRENT_TIMESTAMP WHERE device_id = :device_id - """, + """, {"owner_user_id": user_id, "device_id": device_id}, ) else: @@ -130,7 +115,7 @@ class BindingDAO(BaseDAO): """ INSERT INTO device_bindings (device_id, owner_user_id, child_id, status, bound_at) VALUES (:device_id, :owner_user_id, NULL, 1, CURRENT_TIMESTAMP) - """, + """, {"device_id": device_id, "owner_user_id": user_id}, ) return @@ -156,8 +141,8 @@ class BindingDAO(BaseDAO): unbound_at = NULL, updated_at = CURRENT_TIMESTAMP WHERE device_id = :device_id - """, - {"owner_user_id": user_id, "child_id": child_id, "device_id": device_id}, + """, + {"owner_user_id": user_id, "child_id": child_id, "device_id": device_id}, ) return @@ -168,57 +153,160 @@ class BindingDAO(BaseDAO): """ INSERT INTO device_bindings (device_id, owner_user_id, child_id, status, bound_at) VALUES (:device_id, :owner_user_id, :child_id, 1, CURRENT_TIMESTAMP) - """, + """, {"device_id": device_id, "owner_user_id": user_id, "child_id": child_id}, ) - async def start_bind(self, user_id: int, device_id: str, child_id: Optional[int]) -> str: + async def start_bind(self, user_id: int, device_id: str, child_id: Optional[int]) -> tuple[str, datetime]: bind_token = str(uuid.uuid4()) expires_at = datetime.utcnow() + timedelta(minutes=10) + await self.execute( + """ + UPDATE device_bind_sessions + SET status = :cancelled_status, + consumed_at = COALESCE(consumed_at, CURRENT_TIMESTAMP), + updated_at = CURRENT_TIMESTAMP + WHERE device_id = :device_id + AND status = :pending_status + """, + { + "device_id": device_id, + "pending_status": SESSION_STATUS_PENDING, + "cancelled_status": SESSION_STATUS_CANCELLED, + }, + ) + await self.execute( """ INSERT INTO device_bind_sessions (bind_token, device_id, initiator_user_id, target_child_id, expires_at, status) - VALUES (:bind_token, :device_id, :initiator_user_id, :target_child_id, :expires_at, 1) - """, + VALUES (:bind_token, :device_id, :initiator_user_id, :target_child_id, :expires_at, :status) + """, { "bind_token": bind_token, "device_id": device_id, "initiator_user_id": user_id, "target_child_id": child_id, "expires_at": expires_at, + "status": SESSION_STATUS_PENDING, }, ) - await self.commit() - return bind_token + return bind_token, expires_at async def get_session(self, bind_token: str, user_id: int) -> Optional[Mapping]: return ( await self.execute( - "SELECT * FROM device_bind_sessions WHERE bind_token = :bind_token AND initiator_user_id = :user_id", - {"bind_token": bind_token, "user_id": user_id}, + """ + SELECT + s.*, + CASE + WHEN s.status = :completed_status + AND s.confirmed_at IS NOT NULL + AND c.updated_at >= s.confirmed_at + THEN c.card_uuid + ELSE NULL + END AS card_uuid + FROM device_bind_sessions AS s + LEFT JOIN cards AS c + ON c.device_id = s.device_id + WHERE s.bind_token = :bind_token + AND s.initiator_user_id = :user_id + LIMIT 1 + """, + { + "bind_token": bind_token, + "user_id": user_id, + "completed_status": SESSION_STATUS_COMPLETED, + }, ) ).mappings().first() + async def get_latest_pending_session_by_device(self, device_id: str) -> Optional[Mapping]: + return ( + await self.execute( + """ + SELECT * + FROM device_bind_sessions + WHERE device_id = :device_id + AND status = :status + ORDER BY id DESC + LIMIT 1 + """, + {"device_id": device_id, "status": SESSION_STATUS_PENDING}, + ) + ).mappings().first() + + async def mark_session_status(self, session_id: int, status: int) -> None: + await self.execute( + """ + UPDATE device_bind_sessions + SET status = :status, + consumed_at = COALESCE(consumed_at, CURRENT_TIMESTAMP), + updated_at = CURRENT_TIMESTAMP + WHERE id = :id + """, + {"id": session_id, "status": status}, + ) + async def confirm_bind(self, session_id: int, device_id: str, child_id: Optional[int], user_id: int) -> None: if child_id is not None: await self._upsert_parent_child_relation(user_id=user_id, child_id=child_id) await self.execute( - "UPDATE device_bind_sessions SET status = 2, confirmed_at = CURRENT_TIMESTAMP WHERE id = :id", - {"id": session_id}, + """ + UPDATE device_bind_sessions + SET status = :status, + confirmed_at = CURRENT_TIMESTAMP, + consumed_at = CURRENT_TIMESTAMP, + updated_at = CURRENT_TIMESTAMP + WHERE id = :id + """, + {"id": session_id, "status": SESSION_STATUS_COMPLETED}, ) await self._bind_device(device_id=device_id, user_id=user_id, child_id=child_id) - await self._insert_bind_history(device_id=device_id, child_id=child_id, user_id=user_id, bind_source=1) - await self.commit() + await self._insert_bind_history( + device_id=device_id, + child_id=child_id, + user_id=user_id, + bind_source=BIND_SOURCE_SESSION_CONFIRM, + ) + + async def complete_nfc_bind(self, session_id: int, device_id: str, child_id: Optional[int], user_id: int) -> None: + if child_id is not None: + await self._upsert_parent_child_relation(user_id=user_id, child_id=child_id) + + await self.execute( + """ + UPDATE device_bind_sessions + SET status = :status, + confirmed_at = CURRENT_TIMESTAMP, + consumed_at = CURRENT_TIMESTAMP, + updated_at = CURRENT_TIMESTAMP + WHERE id = :id + """, + {"id": session_id, "status": SESSION_STATUS_COMPLETED}, + ) + + await self._bind_device(device_id=device_id, user_id=user_id, child_id=child_id) + await self._insert_bind_history( + device_id=device_id, + child_id=child_id, + user_id=user_id, + bind_source=BIND_SOURCE_NFC, + ) async def direct_bind(self, device_id: str, child_id: Optional[int], user_id: int) -> None: if child_id is not None: await self._upsert_parent_child_relation(user_id=user_id, child_id=child_id) await self._bind_device(device_id=device_id, user_id=user_id, child_id=child_id) - await self._insert_bind_history(device_id=device_id, child_id=child_id, user_id=user_id, bind_source=2) + await self._insert_bind_history( + device_id=device_id, + child_id=child_id, + user_id=user_id, + bind_source=BIND_SOURCE_DIRECT, + ) await self.commit() async def get_current_by_user(self, user_id: int) -> Optional[Mapping]: @@ -231,7 +319,7 @@ class BindingDAO(BaseDAO): AND status = 1 ORDER BY bound_at DESC LIMIT 1 - """, + """, {"user_id": user_id}, ) ).mappings().first() @@ -260,7 +348,7 @@ class BindingDAO(BaseDAO): WHERE {where} ORDER BY db.id DESC LIMIT :limit - """, + """, params, ) ).mappings().all() @@ -275,7 +363,7 @@ class BindingDAO(BaseDAO): WHERE device_id = :device_id AND owner_user_id = :user_id AND status = 1 - """, + """, {"device_id": device_id, "user_id": user_id}, ) ).mappings().first() @@ -287,7 +375,12 @@ class BindingDAO(BaseDAO): await self._upsert_parent_child_relation(user_id=user_id, child_id=child_id) await self._bind_device(device_id=device_id, user_id=user_id, child_id=child_id) - await self._insert_bind_history(device_id=device_id, child_id=child_id, user_id=user_id, bind_source=3) + await self._insert_bind_history( + device_id=device_id, + child_id=child_id, + user_id=user_id, + bind_source=BIND_SOURCE_SET_CHILD, + ) await self.commit() return True @@ -302,7 +395,30 @@ class BindingDAO(BaseDAO): ) await self.execute( - "INSERT INTO device_bind_history (device_id, child_id, bound_by_user_id, unbound_by_user_id, bind_source, bound_at, unbound_at, unbind_reason) SELECT device_id, child_id, bound_by_user_id, :user_id, bind_source, bound_at, CURRENT_TIMESTAMP, 'user_unbind' FROM device_bind_history WHERE device_id = :device_id AND unbound_at IS NULL", + """ + INSERT INTO device_bind_history ( + device_id, + child_id, + bound_by_user_id, + unbound_by_user_id, + bind_source, + bound_at, + unbound_at, + unbind_reason + ) + SELECT + device_id, + child_id, + bound_by_user_id, + :user_id, + bind_source, + bound_at, + CURRENT_TIMESTAMP, + 'user_unbind' + FROM device_bind_history + WHERE device_id = :device_id + AND unbound_at IS NULL + """, {"device_id": device_id, "user_id": user_id}, ) await self.commit() @@ -318,11 +434,12 @@ class BindingDAO(BaseDAO): rows = ( await self.execute( f""" - SELECT * FROM device_bind_history + SELECT * + FROM device_bind_history WHERE {where} ORDER BY bound_at DESC LIMIT :limit - """, + """, params, ) ).mappings().all() diff --git a/talkingq-url/banban/dao/child.py b/talkingq-url/banban/dao/child.py index 65e072e..adf3c7c 100644 --- a/talkingq-url/banban/dao/child.py +++ b/talkingq-url/banban/dao/child.py @@ -24,6 +24,7 @@ class ChildDAO(BaseDAO): child_name: str, child_gender: int = 2, child_birthday: Optional[date] = None, + auto_commit: bool = True, ) -> int: result = await self.execute( """ @@ -34,7 +35,8 @@ class ChildDAO(BaseDAO): ) child_id = inserted_primary_key(result) await self._create_relation(user_id, child_id) - await self.commit() + if auto_commit: + await self.commit() return child_id async def get_by_id(self, child_id: int) -> Optional[Mapping]: diff --git a/talkingq-url/banban/dao/im.py b/talkingq-url/banban/dao/im.py index b4ab869..dcc0447 100644 --- a/talkingq-url/banban/dao/im.py +++ b/talkingq-url/banban/dao/im.py @@ -25,6 +25,14 @@ class ConversationMessageCreateResult: conversation_type: int message: dict + @property + def conversation_type_name(self) -> str: + if self.conversation_type == 1: + return "child_peer" + if self.conversation_type == 2: + return "parent_child" + return f"unknown_{self.conversation_type}" + class ImDAO(BaseDAO): async def assert_parent_child_access(self, *, user_id: int, child_id: int) -> Mapping[str, Any]: @@ -179,6 +187,22 @@ class ImDAO(BaseDAO): return "[image]" return "[json]" + async def ensure_parent_child_conversation(self, *, parent_user_id: int, child_id: int) -> int: + child_row = await self.assert_parent_child_access(user_id=parent_user_id, child_id=child_id) + parent_row = await self._get_parent_row(parent_user_id) + if not parent_row: + from fastapi import HTTPException + raise HTTPException(status_code=404, detail="parent not found") + + return await self._get_or_create_conversation( + conversation_type=2, + participant_a_type=2, + participant_a_id=str(child_id), + participant_b_type=1, + participant_b_id=str(parent_user_id), + pair_key=f"{child_id}:{parent_user_id}", + ) + async def create_message( self, *, @@ -428,8 +452,8 @@ class ImDAO(BaseDAO): FROM im_conversations WHERE conversation_type = :conversation_type AND pair_key = :pair_key - {lock_clause} LIMIT 1 + {lock_clause} """ ), {"conversation_type": conversation_type, "pair_key": pair_key}, @@ -449,8 +473,8 @@ class ImDAO(BaseDAO): SELECT id, conversation_type, last_seq, status FROM im_conversations WHERE id = :conversation_id - {lock_clause} LIMIT 1 + {lock_clause} """ ), {"conversation_id": conversation_id}, @@ -484,4 +508,4 @@ class ImDAO(BaseDAO): async def _next_primary_key(self, table_name: str) -> int | None: result = await self.execute(text(f"SELECT 1")) - return None \ No newline at end of file + return None diff --git a/talkingq-url/banban/routers/bindings.py b/talkingq-url/banban/routers/bindings.py index 8db2d62..3104c04 100644 --- a/talkingq-url/banban/routers/bindings.py +++ b/talkingq-url/banban/routers/bindings.py @@ -4,12 +4,15 @@ from datetime import datetime from fastapi import APIRouter, Depends, HTTPException, Query, Request, status from pydantic import BaseModel -try: - from banban.security import get_current_user_id - from banban.service.binding import BindingError, BindingService -except ModuleNotFoundError: - from banban.security import get_current_user_id - from banban.service.binding import BindingError, BindingService +from banban.dao.binding import ( + SESSION_STATUS_CANCELLED, + SESSION_STATUS_COMPLETED, + SESSION_STATUS_EXPIRED, + SESSION_STATUS_FAILED, + SESSION_STATUS_PENDING, +) +from banban.security import get_current_user_id +from banban.service.binding import BindingError, BindingService router = APIRouter(prefix="/bindings", tags=["bindings"]) @@ -25,6 +28,16 @@ class BindStartRequest(BaseModel): class BindStartResponse(BaseModel): bind_token: str expires_at: str + status: int + + +class BindSessionResponse(BaseModel): + bind_token: str + device_id: str + child_id: int | None = None + status: int + expires_at: str + card_uuid: str | None = None class BindConfirmRequest(BaseModel): @@ -77,6 +90,7 @@ async def start_bind( request: Request, current_user_id: int = Depends(get_current_user_id), ) -> BindStartResponse: + del request service = BindingService() try: bind_token, expires_at = await service.start_bind( @@ -85,9 +99,35 @@ async def start_bind( payload.serial_number, payload.child_id, ) - except BindingError as e: - raise HTTPException(status_code=e.status_code, detail=str(e)) - return BindStartResponse(bind_token=bind_token, expires_at=expires_at.isoformat()) + except BindingError as exc: + raise HTTPException(status_code=exc.status_code, detail=str(exc)) + return BindStartResponse( + bind_token=bind_token, + expires_at=expires_at.isoformat(), + status=SESSION_STATUS_PENDING, + ) + + +@router.get("/sessions/{bind_token}", response_model=BindSessionResponse) +async def get_bind_session( + bind_token: str, + request: Request, + current_user_id: int = Depends(get_current_user_id), +) -> BindSessionResponse: + del request + service = BindingService() + session = await service.get_bind_session(bind_token, current_user_id) + if not session: + raise HTTPException(status_code=404, detail="bind session not found") + + return BindSessionResponse( + bind_token=session["bind_token"], + device_id=session["device_id"], + child_id=session["target_child_id"], + status=int(session["status"]), + expires_at=session["expires_at"].isoformat(), + card_uuid=session.get("card_uuid"), + ) @router.post("/confirm", response_model=BindConfirmResponse) @@ -96,13 +136,14 @@ async def confirm_bind( request: Request, current_user_id: int = Depends(get_current_user_id), ) -> BindConfirmResponse: + del request service = BindingService() try: result = await service.confirm_bind(payload.bind_token, current_user_id) - except BindingError as e: - raise HTTPException(status_code=e.status_code, detail=str(e)) - except ValueError as e: - raise HTTPException(status_code=400, detail=str(e)) + except BindingError as exc: + raise HTTPException(status_code=exc.status_code, detail=str(exc)) + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) return BindConfirmResponse(**result) @@ -127,6 +168,7 @@ async def direct_bind( request: Request, current_user_id: int = Depends(get_current_user_id), ) -> DirectBindResponse: + del request service = BindingService() try: result = await service.direct_bind( @@ -135,8 +177,8 @@ async def direct_bind( payload.child_id, current_user_id, ) - except BindingError as e: - raise HTTPException(status_code=e.status_code, detail=str(e)) + except BindingError as exc: + raise HTTPException(status_code=exc.status_code, detail=str(exc)) return DirectBindResponse(**result) @@ -147,13 +189,14 @@ async def set_binding_child( request: Request, current_user_id: int = Depends(get_current_user_id), ) -> DirectBindResponse: + del request service = BindingService() try: result = await service.set_binding_child(device_id=device_id, child_id=payload.child_id, user_id=current_user_id) - except BindingError as e: - raise HTTPException(status_code=e.status_code, detail=str(e)) - except ValueError as e: - raise HTTPException(status_code=404, detail=str(e)) + except BindingError as exc: + raise HTTPException(status_code=exc.status_code, detail=str(exc)) + except ValueError as exc: + raise HTTPException(status_code=404, detail=str(exc)) return DirectBindResponse(**result) @@ -162,6 +205,7 @@ async def get_current_binding( request: Request, current_user_id: int = Depends(get_current_user_id), ): + del request service = BindingService() binding = await service.get_current_binding(current_user_id) if not binding: @@ -176,6 +220,7 @@ async def list_bindings( limit: int = Query(default=20, ge=1, le=100), current_user_id: int = Depends(get_current_user_id), ) -> BindingListResponse: + del request service = BindingService() rows, has_more = await service.list_bindings(current_user_id, limit, cursor) next_cursor = int(rows[-1]["id"]) if has_more and rows else None @@ -198,6 +243,7 @@ async def get_binding( request: Request, current_user_id: int = Depends(get_current_user_id), ) -> BindingGetResponse: + del request service = BindingService() binding = await service.get_binding(device_id, current_user_id) if not binding: @@ -211,6 +257,7 @@ async def unbind_device( request: Request, current_user_id: int = Depends(get_current_user_id), ) -> None: + del request service = BindingService() if not await service.unbind(device_id, current_user_id): raise HTTPException(status_code=404, detail="binding not found") @@ -224,6 +271,7 @@ async def get_bind_history( limit: int = 20, current_user_id: int = Depends(get_current_user_id), ) -> BindHistoryResponse: + del request, current_user_id cursor_dt = datetime.fromisoformat(cursor) if cursor else None service = BindingService() rows, has_more = await service.list_history(device_id, limit, cursor_dt) diff --git a/talkingq-url/banban/routers/im.py b/talkingq-url/banban/routers/im.py index 757d1a2..958973d 100644 --- a/talkingq-url/banban/routers/im.py +++ b/talkingq-url/banban/routers/im.py @@ -16,7 +16,7 @@ try: ConversationMessageCreateResponse, ParentChildMessageCreateRequest, ) - from banban.service.im import ImService, im_service + from banban.service.im import ImService, im_service, present_message_item except ModuleNotFoundError: from banban.security import get_current_user_id from banban.schemas.im import ( @@ -27,7 +27,7 @@ except ModuleNotFoundError: ConversationMessageCreateResponse, ParentChildMessageCreateRequest, ) - from banban.service.im import ImService, im_service + from banban.service.im import ImService, im_service, present_message_item router = APIRouter(prefix="/children", tags=["im"]) @@ -506,7 +506,7 @@ async def list_child_conversation_messages( rows = rows[:limit] rows = list(rows) rows.reverse() - items = [_row_to_message_item(row) for row in rows] + items = [await present_message_item(row, audio_storage=im_service.audio_storage) for row in rows] next_cursor_seq = items[0].seq if has_more and items else None logger.info( diff --git a/talkingq-url/banban/service/binding.py b/talkingq-url/banban/service/binding.py index 34bc87a..cc7b364 100644 --- a/talkingq-url/banban/service/binding.py +++ b/talkingq-url/banban/service/binding.py @@ -2,7 +2,15 @@ from collections.abc import Mapping from datetime import datetime from typing import Optional -from banban.dao.binding import BindingDAO +from banban.dao.binding import ( + SESSION_STATUS_CANCELLED, + SESSION_STATUS_COMPLETED, + SESSION_STATUS_EXPIRED, + SESSION_STATUS_FAILED, + SESSION_STATUS_PENDING, + BindingDAO, +) +from services.card_service import card_service from services.database_service_base import DatabaseServiceBase @@ -26,6 +34,12 @@ class BindingService(DatabaseServiceBase): if int(row["is_active"]) != 1: raise BindingError("device is inactive", status_code=400) + def _normalize_session_status(self, session: Mapping) -> int: + status = int(session["status"]) + if status == SESSION_STATUS_PENDING and datetime.utcnow() > session["expires_at"]: + return SESSION_STATUS_EXPIRED + return status + async def start_bind( self, user_id: int, @@ -37,9 +51,15 @@ class BindingService(DatabaseServiceBase): try: await self._ensure_bindable_device(db_session, device_id, serial_number) dao = BindingDAO(db_session) - bind_token = await dao.start_bind(user_id, device_id, child_id) + bind_token, expires_at = await dao.start_bind(user_id, device_id, child_id) await db_session.commit() - return bind_token, datetime.utcnow() + from handlers.mqtt_handler import TalkingQMQTTService + + service = await TalkingQMQTTService.get_instance() + if service is None: + raise BindingError("MQTT service is unavailable", status_code=503) + await service.send_bind_nfc_command(device_id) + return bind_token, expires_at finally: await db_session.close() @@ -52,7 +72,7 @@ class BindingService(DatabaseServiceBase): raise ValueError("Bind session not found") if datetime.utcnow() > session["expires_at"]: raise ValueError("Bind session expired") - if session["status"] != 1: + if int(session["status"]) != SESSION_STATUS_PENDING: raise ValueError("Bind session already processed") await dao.confirm_bind(session["id"], session["device_id"], session["target_child_id"], user_id) @@ -61,6 +81,76 @@ class BindingService(DatabaseServiceBase): finally: await db_session.close() + async def get_bind_session(self, bind_token: str, user_id: int) -> Optional[Mapping]: + db_session = await self.get_session() + try: + dao = BindingDAO(db_session) + session = await dao.get_session(bind_token, user_id) + if session is None: + return None + + normalized_status = self._normalize_session_status(session) + if normalized_status == SESSION_STATUS_EXPIRED and int(session["status"]) != SESSION_STATUS_EXPIRED: + await dao.mark_session_status(int(session["id"]), SESSION_STATUS_EXPIRED) + await db_session.commit() + session = await dao.get_session(bind_token, user_id) + if session is None: + return None + normalized_status = SESSION_STATUS_EXPIRED + + payload = dict(session) + payload["status"] = normalized_status + payload["card_uuid"] = payload.get("card_uuid") + return payload + finally: + await db_session.close() + + async def finalize_nfc_bind(self, device_id: str, card_uuid: str) -> Optional[Mapping]: + db_session = await self.get_session() + try: + dao = BindingDAO(db_session) + session = await dao.get_latest_pending_session_by_device(device_id) + if not session: + await db_session.rollback() + return None + + if datetime.utcnow() > session["expires_at"]: + await dao.mark_session_status(int(session["id"]), SESSION_STATUS_EXPIRED) + await db_session.commit() + return { + "device_id": device_id, + "bind_token": session["bind_token"], + "status": SESSION_STATUS_EXPIRED, + } + + try: + await card_service.activate_card( + card_uuid=card_uuid, + device_id=device_id, + db_session=db_session, + ) + await dao.complete_nfc_bind( + session_id=int(session["id"]), + device_id=device_id, + child_id=session["target_child_id"], + user_id=int(session["initiator_user_id"]), + ) + await db_session.commit() + except Exception: + await dao.mark_session_status(int(session["id"]), SESSION_STATUS_FAILED) + await db_session.commit() + raise + + return { + "device_id": device_id, + "bind_token": session["bind_token"], + "status": SESSION_STATUS_COMPLETED, + "child_id": session["target_child_id"], + "card_uuid": card_uuid, + } + finally: + await db_session.close() + async def get_binding(self, device_id: str, user_id: int) -> Optional[Mapping]: db_session = await self.get_session() try: diff --git a/talkingq-url/banban/service/child.py b/talkingq-url/banban/service/child.py index 269d6d2..d304cdf 100644 --- a/talkingq-url/banban/service/child.py +++ b/talkingq-url/banban/service/child.py @@ -3,6 +3,7 @@ from datetime import date from typing import Optional from banban.dao.child import ChildDAO +from banban.dao.im import ImDAO from services.database_service_base import DatabaseServiceBase @@ -19,10 +20,24 @@ class ChildService(DatabaseServiceBase): ) -> Mapping: db_session = await self.get_session() try: - dao = ChildDAO(db_session) - child_id = await dao.create(user_id, child_name, child_gender, child_birthday) + child_dao = ChildDAO(db_session) + im_dao = ImDAO(db_session) + child_id = await child_dao.create( + user_id, + child_name, + child_gender, + child_birthday, + auto_commit=False, + ) + await im_dao.ensure_parent_child_conversation( + parent_user_id=user_id, + child_id=child_id, + ) await db_session.commit() - return await self.get(child_id) + return await child_dao.get_by_id(child_id) + except Exception: + await db_session.rollback() + raise finally: await db_session.close() diff --git a/talkingq-url/banban/service/im.py b/talkingq-url/banban/service/im.py index 9c5574d..74d3b86 100644 --- a/talkingq-url/banban/service/im.py +++ b/talkingq-url/banban/service/im.py @@ -1,10 +1,12 @@ from dataclasses import dataclass +import hashlib import json from collections.abc import Mapping from typing import Any -from fastapi import HTTPException, status +from fastapi import HTTPException from services.database_service_base import DatabaseServiceBase +from banban.service.message_audio_storage import MessageAudioStorageService, MessageAudioStorageError try: from banban.dao.im import ImDAO, DeviceIdentity, ConversationMessageCreateResult @@ -43,6 +45,12 @@ def participant_type_name(participant_type: int) -> str: return PARTICIPANT_TYPE_NAMES.get(participant_type, f"unknown_{participant_type}") +def build_device_audio_client_msg_id(*, device_id: str, target_device_id: str, audio_url: str) -> str: + raw = f"{device_id}|{target_device_id}|{audio_url}" + digest = hashlib.sha256(raw.encode("utf-8")).hexdigest() + return f"device-audio-{digest[:32]}" + + def normalize_content_json(value: Any) -> dict[str, Any] | None: if value is None: return None @@ -85,9 +93,24 @@ def row_to_message_item(row: Mapping[str, Any]) -> ChildConversationMessageItem: ) +async def present_message_item( + row: Mapping[str, Any], + *, + audio_storage: MessageAudioStorageService, +) -> ChildConversationMessageItem: + item = row_to_message_item(row) + if item.content_type == 2 and item.media_file_key: + try: + item.media_file_key = await audio_storage.get_audio_url(item.media_file_key) + except MessageAudioStorageError: + pass + return item + + class ImService(DatabaseServiceBase): def __init__(self): super().__init__(service_name="im_service") + self.audio_storage = MessageAudioStorageService() async def assert_parent_child_access(self, *, user_id: int, child_id: int) -> Mapping[str, Any]: db_session = await self.get_session() @@ -105,6 +128,34 @@ class ImService(DatabaseServiceBase): finally: await db_session.close() + async def _build_message_create_result( + self, + *, + dao: ImDAO, + idempotent: bool, + conversation_id: int, + conversation_type: int, + client_msg_id: str, + ) -> ConversationMessageCreateResult: + message_row = await dao._get_message_by_conversation_client_id( + conversation_id=conversation_id, + client_msg_id=client_msg_id, + ) + if not message_row: + raise RuntimeError("message was not found after insert") + + presented_message = await present_message_item( + message_row, + audio_storage=self.audio_storage, + ) + + return ConversationMessageCreateResult( + idempotent=idempotent, + conversation_id=conversation_id, + conversation_type=conversation_type, + message=presented_message, + ) + async def create_parent_child_message( self, *, @@ -137,18 +188,12 @@ class ImService(DatabaseServiceBase): receiver_avatar_snapshot=None, payload=payload, ) - message_row = await dao._get_message_by_conversation_client_id( - conversation_id=conversation_id, - client_msg_id=payload.client_msg_id, - ) - if not message_row: - raise RuntimeError("message was not found after insert") - - return ConversationMessageCreateResult( + return await self._build_message_create_result( + dao=dao, idempotent=idempotent, conversation_id=conversation_id, conversation_type=PARENT_CHILD_CONVERSATION_TYPE, - message=row_to_message_item(message_row), + client_msg_id=payload.client_msg_id, ) except Exception: await db_session.rollback() @@ -156,13 +201,87 @@ class ImService(DatabaseServiceBase): finally: await db_session.close() + async def _create_device_message_with_payload( + self, + *, + dao: ImDAO, + device_identity: DeviceIdentity, + payload: DeviceMessageCreateRequest, + ) -> ConversationMessageCreateResult: + if payload.conversation_type == CHILD_PEER_CONVERSATION_TYPE: + if payload.peer_child_id == device_identity.child_id: + raise HTTPException(status_code=400, detail="peer_child_id must be different from current child") + sender_child_row = await dao.assert_child_exists(child_id=device_identity.child_id) + receiver_child_row = await dao.assert_child_exists(child_id=payload.peer_child_id) + participant_a_id, participant_b_id, pair_key = dao._build_child_peer_pair( + device_identity.child_id, + payload.peer_child_id, + ) + + conversation_id, idempotent = await dao.create_message( + conversation_type=CHILD_PEER_CONVERSATION_TYPE, + participant_a_type=CHILD_PARTICIPANT_TYPE, + participant_a_id=participant_a_id, + participant_b_type=CHILD_PARTICIPANT_TYPE, + participant_b_id=participant_b_id, + pair_key=pair_key, + sender_type=CHILD_PARTICIPANT_TYPE, + sender_id=str(device_identity.child_id), + receiver_type=CHILD_PARTICIPANT_TYPE, + receiver_id=str(payload.peer_child_id), + sender_name_snapshot=sender_child_row["child_name"], + sender_avatar_snapshot=None, + receiver_name_snapshot=receiver_child_row["child_name"], + receiver_avatar_snapshot=None, + payload=payload, + ) + else: + sender_child_row = await dao.assert_child_exists(child_id=device_identity.child_id) + parent_row = await dao._get_parent_row(payload.parent_user_id) + if not parent_row: + raise HTTPException(status_code=404, detail="parent not found") + await dao.assert_parent_child_access(user_id=payload.parent_user_id, child_id=device_identity.child_id) + + conversation_id, idempotent = await dao.create_message( + conversation_type=PARENT_CHILD_CONVERSATION_TYPE, + participant_a_type=CHILD_PARTICIPANT_TYPE, + participant_a_id=str(device_identity.child_id), + participant_b_type=PARENT_PARTICIPANT_TYPE, + participant_b_id=str(payload.parent_user_id), + pair_key=f"{device_identity.child_id}:{payload.parent_user_id}", + sender_type=CHILD_PARTICIPANT_TYPE, + sender_id=str(device_identity.child_id), + receiver_type=PARENT_PARTICIPANT_TYPE, + receiver_id=str(payload.parent_user_id), + sender_name_snapshot=sender_child_row["child_name"], + sender_avatar_snapshot=None, + receiver_name_snapshot=parent_row["nickname"], + receiver_avatar_snapshot=parent_row["avatar_url"], + payload=payload, + ) + + return await self._build_message_create_result( + dao=dao, + idempotent=idempotent, + conversation_id=conversation_id, + conversation_type=payload.conversation_type, + client_msg_id=payload.client_msg_id, + ) + async def create_device_message( self, *, device_id: str, serial_number: str, - payload: DeviceMessageCreateRequest, + payload: DeviceMessageCreateRequest | None = None, + target_device_id: str | None = None, + audio_url: str | None = None, ) -> tuple[DeviceIdentity, ConversationMessageCreateResult]: + if payload is None and (not target_device_id or not audio_url): + raise ValueError("payload or target_device_id/audio_url is required") + if payload is not None and (target_device_id is not None or audio_url is not None): + raise ValueError("payload and target_device_id/audio_url cannot be used together") + db_session = await self.get_session() try: dao = ImDAO(db_session) @@ -171,70 +290,31 @@ class ImService(DatabaseServiceBase): serial_number=serial_number, ) - if payload.conversation_type == CHILD_PEER_CONVERSATION_TYPE: - if payload.peer_child_id == device_identity.child_id: - raise HTTPException(status_code=400, detail="peer_child_id must be different from current child") - sender_child_row = await dao.assert_child_exists(child_id=device_identity.child_id) - receiver_child_row = await dao.assert_child_exists(child_id=payload.peer_child_id) - participant_a_id, participant_b_id, pair_key = dao._build_child_peer_pair( - device_identity.child_id, - payload.peer_child_id, - ) - - conversation_id, idempotent = await dao.create_message( + resolved_payload = payload + if resolved_payload is None: + target_device_identity = await dao.get_device_by_id(device_id=target_device_id) + resolved_payload = DeviceMessageCreateRequest( conversation_type=CHILD_PEER_CONVERSATION_TYPE, - participant_a_type=CHILD_PARTICIPANT_TYPE, - participant_a_id=participant_a_id, - participant_b_type=CHILD_PARTICIPANT_TYPE, - participant_b_id=participant_b_id, - pair_key=pair_key, - sender_type=CHILD_PARTICIPANT_TYPE, - sender_id=str(device_identity.child_id), - receiver_type=CHILD_PARTICIPANT_TYPE, - receiver_id=str(payload.peer_child_id), - sender_name_snapshot=sender_child_row["child_name"], - sender_avatar_snapshot=None, - receiver_name_snapshot=receiver_child_row["child_name"], - receiver_avatar_snapshot=None, - payload=payload, - ) - else: - sender_child_row = await dao.assert_child_exists(child_id=device_identity.child_id) - parent_row = await dao._get_parent_row(payload.parent_user_id) - if not parent_row: - raise HTTPException(status_code=404, detail="parent not found") - await dao.assert_parent_child_access(user_id=payload.parent_user_id, child_id=device_identity.child_id) - - conversation_id, idempotent = await dao.create_message( - conversation_type=PARENT_CHILD_CONVERSATION_TYPE, - participant_a_type=CHILD_PARTICIPANT_TYPE, - participant_a_id=str(device_identity.child_id), - participant_b_type=PARENT_PARTICIPANT_TYPE, - participant_b_id=str(payload.parent_user_id), - pair_key=f"{device_identity.child_id}:{payload.parent_user_id}", - sender_type=CHILD_PARTICIPANT_TYPE, - sender_id=str(device_identity.child_id), - receiver_type=PARENT_PARTICIPANT_TYPE, - receiver_id=str(payload.parent_user_id), - sender_name_snapshot=sender_child_row["child_name"], - sender_avatar_snapshot=None, - receiver_name_snapshot=parent_row["nickname"], - receiver_avatar_snapshot=parent_row["avatar_url"], - payload=payload, + peer_child_id=target_device_identity.child_id, + content_type=2, + media_file_key=audio_url, + media_mime_type="audio/mpeg", + client_msg_id=build_device_audio_client_msg_id( + device_id=device_id, + target_device_id=target_device_id, + audio_url=audio_url, + ), + ext_json={ + "source": "device_audio_message", + "source_device_id": device_id, + "target_device_id": target_device_id, + }, ) - message_row = await dao._get_message_by_conversation_client_id( - conversation_id=conversation_id, - client_msg_id=payload.client_msg_id, - ) - if not message_row: - raise RuntimeError("message was not found after insert") - - result = ConversationMessageCreateResult( - idempotent=idempotent, - conversation_id=conversation_id, - conversation_type=payload.conversation_type, - message=row_to_message_item(message_row), + result = await self._create_device_message_with_payload( + dao=dao, + device_identity=device_identity, + payload=resolved_payload, ) return device_identity, result except Exception: @@ -251,55 +331,5 @@ class ImService(DatabaseServiceBase): finally: await db_session.close() - ''' - Todo 创建设备消息, 还不完善 - ''' - async def create_device_message( - self, - *, - device_id: str, - serial_number: str, - target_device_id: str, - audio_url: str, - ): - db_session = await self.get_session() - try: - dao = ImDAO(db_session) - device_identity = await dao.authenticate_device_identity( - device_id=device_id, - serial_number=serial_number, - ) - target_device_identity = await dao.get_device_by_id(device_id=target_device_id) - sender_child_row = await dao.assert_child_exists(child_id=device_identity.child_id) - receiver_child_row = await dao.assert_child_exists(child_id=target_device_identity.child_id) - - conversation_id, idempotent = await dao.create_message( - conversation_type=PARENT_CHILD_CONVERSATION_TYPE, - participant_a_type=CHILD_PARTICIPANT_TYPE, - participant_a_id=str(device_identity.child_id), - participant_b_type=PARENT_PARTICIPANT_TYPE, - participant_b_id=str(receiver_child_row.child_id), - pair_key=f"{device_identity.child_id}:{target_device_identity.child_id}", - sender_type=CHILD_PARTICIPANT_TYPE, - sender_id=str(device_identity.child_id), - receiver_type=PARENT_PARTICIPANT_TYPE, - receiver_id=str(receiver_child_row.child_id), - sender_name_snapshot=sender_child_row["child_name"], - sender_avatar_snapshot=None, - receiver_name_snapshot=receiver_child_row["child_name"], - receiver_avatar_snapshot=None, - payload=None, - ) - - - return device_identity - except Exception: - await db_session.rollback() - raise - finally: - await db_session.close() - - -# 创建全局 ImService 实例 im_service = ImService() diff --git a/talkingq-url/banban/service/message_audio_storage.py b/talkingq-url/banban/service/message_audio_storage.py new file mode 100644 index 0000000..eb749d3 --- /dev/null +++ b/talkingq-url/banban/service/message_audio_storage.py @@ -0,0 +1,129 @@ +import asyncio +from dataclasses import dataclass +from datetime import UTC, datetime +from urllib.parse import urlparse +from uuid import uuid4 + +try: + from qcloud_cos import CosConfig, CosS3Client +except ModuleNotFoundError: # pragma: no cover - exercised in runtime env + CosConfig = None + CosS3Client = None + +from config import settings + + +class MessageAudioStorageError(Exception): + pass + + +@dataclass(frozen=True) +class StoredMessageAudio: + file_key: str + public_url: str + + +class MessageAudioStorageService: + def __init__(self) -> None: + self._client = None + + def _assert_ready(self) -> None: + if CosConfig is None or CosS3Client is None: + raise MessageAudioStorageError("COS SDK is not installed") + + required_pairs = { + "COS_SECRET_ID": settings.cos_secret_id, + "COS_SECRET_KEY": settings.cos_secret_key, + "COS_REGION": settings.cos_region, + "COS_BUCKET_MESSAGE": settings.cos_bucket_message, + } + missing = [key for key, value in required_pairs.items() if not value] + if missing: + raise MessageAudioStorageError(f"missing COS message config: {', '.join(missing)}") + + def _get_client(self): + if self._client is None: + config = CosConfig( + Region=settings.cos_region, + SecretId=settings.cos_secret_id, + SecretKey=settings.cos_secret_key, + Scheme="https", + ) + self._client = CosS3Client(config) + return self._client + + def _build_key(self, *, device_id: str, extension: str) -> str: + prefix = settings.cos_message_prefix.strip("/") or "messages/audio" + now = datetime.now(UTC) + return ( + f"{prefix}/{device_id}/{now.strftime('%Y/%m/%d')}/" + f"{uuid4().hex}.{extension}" + ) + + def _build_public_url(self, *, file_key: str) -> str: + base_url = settings.cos_public_base_url.strip().rstrip("/") + if not base_url: + base_url = f"https://{settings.cos_bucket_message}.cos.{settings.cos_region}.myqcloud.com" + return f"{base_url}/{file_key.lstrip('/')}" + + def _normalize_file_key(self, file_key_or_url: str) -> str: + value = (file_key_or_url or "").strip() + if not value: + raise MessageAudioStorageError("audio file key is required") + if value.startswith("http://") or value.startswith("https://"): + parsed = urlparse(value) + path = parsed.path.lstrip("/") + if not path: + raise MessageAudioStorageError("audio file key is invalid") + return path + return value.lstrip("/") + + async def upload_audio( + self, + *, + device_id: str, + content: bytes, + content_type: str = "audio/mpeg", + extension: str = "mp3", + ) -> StoredMessageAudio: + self._assert_ready() + if not content: + raise MessageAudioStorageError("audio content is empty") + + file_key = self._build_key(device_id=device_id, extension=extension) + await asyncio.to_thread( + self._upload_audio_sync, + file_key=file_key, + content=content, + content_type=content_type, + ) + return StoredMessageAudio( + file_key=file_key, + public_url=self._build_public_url(file_key=file_key), + ) + + async def get_audio_url(self, file_key_or_url: str) -> str: + self._assert_ready() + normalized_key = self._normalize_file_key(file_key_or_url) + return await asyncio.to_thread( + self._get_client().get_presigned_url, + Bucket=settings.cos_bucket_message, + Key=normalized_key, + Method="GET", + Expired=settings.cos_avatar_url_expire_seconds, + ) + + def _upload_audio_sync( + self, + *, + file_key: str, + content: bytes, + content_type: str, + ) -> None: + self._get_client().put_object( + Bucket=settings.cos_bucket_message, + Body=content, + Key=file_key, + ContentType=content_type, + EnableMD5=False, + ) diff --git a/talkingq-url/config.py b/talkingq-url/config.py index 28ed176..7b175e9 100644 --- a/talkingq-url/config.py +++ b/talkingq-url/config.py @@ -89,6 +89,7 @@ class Settings(BaseSettings): cos_bucket_message: str = Field(default="", validation_alias="COS_BUCKET_MESSAGE") cos_bucket_ava: str = Field(default="", validation_alias="COS_BUCKET_AVA") cos_public_base_url: str = Field(default="", validation_alias="COS_PUBLIC_BASE_URL") + cos_message_prefix: str = Field(default="messages/audio/", validation_alias="COS_MESSAGE_PREFIX") cos_avatar_prefix: str = Field(default="avatars/", validation_alias="COS_AVATAR_PREFIX") cos_avatar_url_expire_seconds: int = Field( default=86400, diff --git a/talkingq-url/handlers/audio_file_handler.py b/talkingq-url/handlers/audio_file_handler.py index efd34ab..e695baf 100644 --- a/talkingq-url/handlers/audio_file_handler.py +++ b/talkingq-url/handlers/audio_file_handler.py @@ -1,36 +1,19 @@ -import os -import uuid -from config import settings +from banban.service.message_audio_storage import MessageAudioStorageService from utils.logger import session_logger -# from utils.audio_denoiser import reduce_background_noise + + +message_audio_storage_service = MessageAudioStorageService() async def save_audio_file(audio_data: bytes, device_id: str) -> str: - """ - 保存音频数据到 assets/audio 目录 - - Args: - audio_data: 音频二进制数据 - device_id: 设备ID - - Returns: - 音频文件的相对路径 - """ + """Upload device audio to COS and return its object key.""" try: - audio_dir = os.path.join(settings.assets_dir, "audio") - os.makedirs(audio_dir, exist_ok=True) - - filename = f"{device_id}_{uuid.uuid4().hex[:8]}.mp3" - filepath = os.path.join(audio_dir, filename) - - with open(filepath, 'wb') as f: - f.write(audio_data) - - # relative_path = f"assets/audio/{filename}" - session_logger.info(device_id, "audio", f"音频文件已保存: {filepath}") - - # reduce_background_noise(filepath, relative_path,noise_path='assets/audio/noise_sample.wav',normalize_volume=True) - return filepath + stored = await message_audio_storage_service.upload_audio( + device_id=device_id, + content=audio_data, + ) + session_logger.info(device_id, "audio", f"audio uploaded to COS: {stored.file_key}") + return stored.file_key except Exception as e: - session_logger.error(device_id, "audio", f"保存音频文件时出错: {e}", exc_info=True) - raise \ No newline at end of file + session_logger.error(device_id, "audio", f"failed to store audio: {e}", exc_info=True) + raise diff --git a/talkingq-url/handlers/mqtt_handler.py b/talkingq-url/handlers/mqtt_handler.py index d70b2fb..d4ca8c4 100644 --- a/talkingq-url/handlers/mqtt_handler.py +++ b/talkingq-url/handlers/mqtt_handler.py @@ -3,6 +3,7 @@ import time import asyncio from typing import Optional, Dict, Callable, Awaitable from banban.service.device_setting import device_setting_service +from banban.service.binding import BindingService from config import settings from services.card_service import card_service from services.offline_audio_cache import offline_audio_cache @@ -166,8 +167,13 @@ class TalkingQMQTTService: nfc_uuid = params.get("uuid") # 卡片与设备绑定 - await card_service.activate_card(device_id=device_id, card_uuid=nfc_uuid) - logger.info(device_id, "", f"[NFC绑定] 设备 {device_id} 请求绑定卡片, UUID={nfc_uuid}") + # await card_service.activate_card(device_id=device_id, card_uuid=nfc_uuid) + service = BindingService() + result = await service.finalize_nfc_bind(device_id=device_id, card_uuid=nfc_uuid) + if result is None: + logger.warning(device_id, "", "[NFC bind] no pending bind session found") + return + logger.info(device_id, "", f"[NFC绑定] 设备 {device_id} 请求绑定卡片, UUID={nfc_uuid}, result={result}") async def _handle_open_response(self, device_id: str, payload: dict): params = payload.get("params", {}) diff --git a/talkingq-url/services/card_service.py b/talkingq-url/services/card_service.py index 737f32b..41957cf 100644 --- a/talkingq-url/services/card_service.py +++ b/talkingq-url/services/card_service.py @@ -1,12 +1,14 @@ import asyncio -from typing import Optional, Dict -import time -from sqlalchemy import select, update, insert -from sqlalchemy.ext.asyncio import AsyncSession -from utils.logger import session_logger +from typing import Dict, Optional + from pydantic import BaseModel +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession + from database.models import Card as DBCard from services.database_service_base import DatabaseServiceBase +from utils.logger import session_logger + class Card(BaseModel): card_id: Optional[int] = None @@ -17,10 +19,9 @@ class Card(BaseModel): total_swaps: int = 0 created_at: Optional[float] = None updated_at: Optional[float] = None - + @classmethod def from_db_model(cls, db_model: DBCard): - """从数据库模型创建卡片对象""" return cls( card_id=db_model.card_id, card_uuid=db_model.card_uuid, @@ -29,162 +30,205 @@ class Card(BaseModel): status=db_model.status, total_swaps=db_model.total_swaps, created_at=db_model.created_at.timestamp() if db_model.created_at else None, - updated_at=db_model.updated_at.timestamp() if db_model.updated_at else None + updated_at=db_model.updated_at.timestamp() if db_model.updated_at else None, ) + class CardService(DatabaseServiceBase): def __init__(self): super().__init__(service_name="card") self.cards: Dict[str, Card] = {} self.lock = asyncio.Lock() - + + async def _cache_card(self, card: Card) -> None: + async with self.lock: + self.cards[card.card_uuid] = card + async def _load_card_from_db(self, card_uuid: str, async_session: AsyncSession) -> Optional[Card]: - """从数据库加载卡片信息""" try: query = select(DBCard).where(DBCard.card_uuid == card_uuid) result = await async_session.execute(query) db_card = result.scalar_one_or_none() - - if db_card: - card = Card.from_db_model(db_card) - async with self.lock: - self.cards[card_uuid] = card - return card + + if not db_card: + return None + + card = Card.from_db_model(db_card) + await self._cache_card(card) + return card + except Exception as exc: + session_logger.error("card", "service", f"load card failed: {exc}") return None - except Exception as e: - session_logger.error("card", "service", f"从数据库加载卡片失败: {str(e)}") - return None - - async def _save_card_to_db(self, card: Card, async_session: AsyncSession): - """保存卡片信息到数据库""" + + async def _clear_existing_device_card( + self, + device_id: str, + card_uuid: str, + async_session: AsyncSession, + ) -> None: + query = select(DBCard).where( + DBCard.device_id == device_id, + DBCard.card_uuid != card_uuid, + ) + result = await async_session.execute(query) + existing_cards = result.scalars().all() + + for db_card in existing_cards: + db_card.device_id = None + db_card.status = 0 + + async with self.lock: + cached = self.cards.get(db_card.card_uuid) + if cached is not None: + cached.device_id = None + cached.status = 0 + + if existing_cards: + await async_session.flush() + + async def _save_card_to_db(self, card: Card, async_session: AsyncSession, commit: bool = True) -> None: try: query = select(DBCard).where(DBCard.card_uuid == card.card_uuid) result = await async_session.execute(query) existing_card = result.scalar_one_or_none() - + if existing_card: - stmt = update(DBCard).where( - DBCard.card_uuid == card.card_uuid - ).values( - device_id=card.device_id, - card_name=card.card_name, - status=card.status, - total_swaps=card.total_swaps - ) + existing_card.device_id = card.device_id + existing_card.card_name = card.card_name + existing_card.status = card.status + existing_card.total_swaps = card.total_swaps + db_card = existing_card else: - stmt = insert(DBCard).values( + db_card = DBCard( card_uuid=card.card_uuid, device_id=card.device_id, card_name=card.card_name, status=card.status, - total_swaps=card.total_swaps + total_swaps=card.total_swaps, ) - - await async_session.execute(stmt) - await async_session.commit() - - session_logger.info("card", "service", f"卡片已保存到数据库: {card.card_uuid}") - except Exception as e: - await async_session.rollback() - session_logger.error("card", "service", f"保存卡片到数据库失败: {str(e)}") + async_session.add(db_card) + + await async_session.flush() + if commit: + await async_session.commit() + + card.card_id = db_card.card_id + await self._cache_card(card) + session_logger.info("card", "service", f"card saved: {card.card_uuid}") + except Exception as exc: + if commit: + await async_session.rollback() + session_logger.error("card", "service", f"save card failed: {exc}") raise - - async def get_card_by_uuid(self, card_uuid: str, force_refresh: bool = False) -> Optional[Card]: - """根据UUID获取卡片信息""" + + async def get_card_by_uuid( + self, + card_uuid: str, + force_refresh: bool = False, + db_session: Optional[AsyncSession] = None, + ) -> Optional[Card]: await self._init_database() - + card = None - if not force_refresh: async with self.lock: card = self.cards.get(card_uuid) - - if not card: - db_session = await self.db_manager.get_session() - try: - card = await self._load_card_from_db(card_uuid, db_session) - finally: - await db_session.close() - - return card - + + if card: + return card + + if db_session is not None: + return await self._load_card_from_db(card_uuid, db_session) + + session = await self.db_manager.get_session() + try: + return await self._load_card_from_db(card_uuid, session) + finally: + await session.close() + async def get_card_by_device_id(self, device_id: str) -> list[Card]: - """根据设备ID获取卡片列表""" await self._init_database() - + db_session = await self.db_manager.get_session() try: query = select(DBCard).where(DBCard.device_id == device_id) result = await db_session.execute(query) db_cards = result.scalars().all() - + cards = [] for db_card in db_cards: card = Card.from_db_model(db_card) - async with self.lock: - self.cards[card.card_uuid] = card + await self._cache_card(card) cards.append(card) - return cards - except Exception as e: - session_logger.error(device_id, "card", f"根据设备ID获取卡片失败: {str(e)}") + except Exception as exc: + session_logger.error(device_id, "card", f"load device cards failed: {exc}") return [] finally: await db_session.close() - - async def activate_card(self, card_uuid: str, device_id: str, card_name: Optional[str] = None) -> Card: - """激活卡片并绑定到设备""" + + async def activate_card( + self, + card_uuid: str, + device_id: str, + card_name: Optional[str] = None, + db_session: Optional[AsyncSession] = None, + ) -> Card: await self._init_database() - - db_session = await self.db_manager.get_session() + + owns_session = db_session is None + if db_session is None: + db_session = await self.db_manager.get_session() + try: - # 检查卡片是否已存在 - existing_card = await self.get_card_by_uuid(card_uuid) - + existing_card = await self.get_card_by_uuid(card_uuid, db_session=db_session) + await self._clear_existing_device_card(device_id=device_id, card_uuid=card_uuid, async_session=db_session) + if existing_card: - # 更新现有卡片 existing_card.device_id = device_id existing_card.card_name = card_name - existing_card.status = 1 # 激活状态 - await self._save_card_to_db(existing_card, db_session) - session_logger.info(device_id, "card", f"卡片已激活并绑定到设备: {card_uuid}") + existing_card.status = 1 + await self._save_card_to_db(existing_card, db_session, commit=owns_session) + session_logger.info(device_id, "card", f"card activated: {card_uuid}") return existing_card - else: - # 创建新卡片 - new_card = Card( - card_uuid=card_uuid, - device_id=device_id, - card_name=card_name, - status=1, # 激活状态 - total_swaps=0 - ) - await self._save_card_to_db(new_card, db_session) - session_logger.info(device_id, "card", f"新卡片已创建并激活: {card_uuid}") - return new_card + + new_card = Card( + card_uuid=card_uuid, + device_id=device_id, + card_name=card_name, + status=1, + total_swaps=0, + ) + await self._save_card_to_db(new_card, db_session, commit=owns_session) + session_logger.info(device_id, "card", f"new card activated: {card_uuid}") + return new_card + except Exception: + if owns_session: + await db_session.rollback() + raise + finally: + if owns_session: + await db_session.close() + + async def increment_swap_count(self, card_uuid: str) -> Optional[Card]: + await self._init_database() + + card = await self.get_card_by_uuid(card_uuid) + if not card: + return None + + card.total_swaps += 1 + db_session = await self.db_manager.get_session() + try: + await self._save_card_to_db(card, db_session) + session_logger.info("card", "service", f"swap count incremented: {card_uuid}, total={card.total_swaps}") + return card finally: await db_session.close() - - async def increment_swap_count(self, card_uuid: str) -> Optional[Card]: - """增加卡片交换次数""" - await self._init_database() - - card = await self.get_card_by_uuid(card_uuid) - if card: - card.total_swaps += 1 - db_session = await self.db_manager.get_session() - try: - await self._save_card_to_db(card, db_session) - session_logger.info("card", "service", f"卡片交换次数已增加: {card_uuid}, 总次数: {card.total_swaps}") - return card - finally: - await db_session.close() - return None - + async def check_card_ownership(self, card_uuid: str, device_id: str) -> bool: - """检查卡片是否属于指定设备""" card = await self.get_card_by_uuid(card_uuid) - if card and card.device_id == device_id: - return True - return False + return bool(card and card.device_id == device_id) + card_service = CardService()