From f494eaa6865525a2b79a41250ddf1dfb843bed8e Mon Sep 17 00:00:00 2001 From: yumoqing Date: Mon, 27 Jul 2026 13:49:41 +0800 Subject: [PATCH] feat: uapi-based RAG ingest pipeline + improved delete cleanup MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - _rag_ingest_async(): chunk → embed(VDB) → VDB upsert → NER → graph save → DB records - _split_text(): paragraph-based text chunker - doc_upload_handler: replace pipeline.ingest with _rag_ingest_async - doc_delete_handler: use uapi vdb-delete + graph-delete (was broken graph/save) --- rag/__pycache__/init.cpython-310.pyc | Bin 1616 -> 23423 bytes rag/init.py | 121 +++++++++++++++++++++++---- 2 files changed, 104 insertions(+), 17 deletions(-) diff --git a/rag/__pycache__/init.cpython-310.pyc b/rag/__pycache__/init.cpython-310.pyc index d457eafa6e00262ca9de6af0fbda2e6f6da40821..4094703658c8e689a271eb00cc87a143be0257ef 100644 GIT binary patch literal 23423 zcmb_^32+?OncnnF&rHvO!QdnaQe=zwP#{23Tv^hJvLzCdX;F|&P@>IR7j}R*00x|k z-902R>`^FD*5NGIlI8V5(QDA>Y7={%tk6e+96Ro+)i(Q@o~jQ69Zcv(BHl@rGkd>kt!%czZk%FP~0Qdh@oed<7>oPF^rla)T|ThQL`3%w}}nd8^+#7p<{0y_HGxOu(w`p z7F*E%HnCM~!*+w%E_PtMQRKuyG?vR+>Pya5uXg@?yY|qvzI)QpPMWgrwWFCaPskw>mSZnM5$ny`Y!#*0}tz0 zmtMc}!}D*we(}nuE?+g4z5c1UzP~hf<%Ol|U;J9~I{FIb z6el@bHOhHws&HloM+w>Wc;yrhvPTMYMbj!657jF(R%al)FNMVr>NDpRaU1G;#Rjp_t=CDw6ct>#87)$L-4<@M)o(K-w%J|nb#q~{{ZeEu zB6g4fBH<%;@txK4iZ=o=iu;b?xBF6TP7!-9sm*R%TT0s9mlY*+akyTv`1BDO|%Yi6y~LfX!Xd&RF`2pUv38+^K|`JxZuCY-PlKb`NHDztw#{^lIo!7`+m*QcG!Z zKc1rJyn=VE91F$Jhrl&Y^-Q|k)5_^k?t>Ff1T^b(pQuz%mkMI8Fg2Yw3#OBR#Hg3S zD5euVF+Eij2J}*>h~I}Ak+q1g}zzr zF6^8AKrZIQW=a%&*!g`Si@UI-vtAg^IYylot)6ve7;0#K;5%@_tv9$Z&8fJ znm6{*RqAOYUwh~^#h^YK35=hf>|Xswctpo|4Wnx0vQB);s8^;crqe~Lw+^LHHReFT zaq2=j;)L;Vky4@JB(18IFHJEuWDDi#0+Z8}RDMpxv_MTXK~4;l4ozavLm5QQPGa0s z@dmwIgS1Po*XW^rlHEoh75yljNTG5HZ_pbC`UUN(Mqvh=MGRm%is5Krw<)e8#>0E9 zS1jYBbh$>uo1(GA+<~cdl@l`xX1!!NN`du|EC%U{Zg*ZhrLz{5=Hn>B=VD4&$te9w zFO-j_4B?l=uh;)2!$}+s^O^N%7kfAStDUjfjdWZ^ZNHLM_Fc)YA5=Cd+PUz%hCrun zYKs_#a4xH@zr~KPwkF;~ZS8KxM7S9zQYp6S#eU5Z6Y=K zRFXw1wk%RJ6UGppym4-F{P6gJqxw$$p(Bqys>_#096MW_GxcL19zQa!OLVY*JMq}5 z!gl?EiG%tMP%GrEerQskc;fJ3{jnputl7WitlVjC$)Wb(_!0fVk9++Re~#!iqgXYH z)}nrJeDZ)3lg(#Ha(n-XPziKS;%5?l&V}sICt^*&PFp#zyc|9mo|MQ27*`d=8SAyM z6CT}d{07=~Qd;sp;=h?{xpx0J^+%oehZ>t*_~;Sf=zXQ?Oul5^j~%Z)!1NOb4vdc< z96t#8qZ-xIUU*L`U^4eSY!bJQFMt6 z!4v`|03PwhBNC2qzug4h24sZ>P4d<5eld_CJllc-6!UDK~wGmF;xVK%!IEOb5Wy=of79ln_f`-J@*E@~reQczv^v4YUx^zanKOG@LIuObw(Y(UOo5f%|5 zZr*3oCKa)qu{ zhH4LqR5NKOUJVh+!b(Z;`81bZTxKkwezZVXQBVCgxT|9SPBEM?0Q_w4$pP* zj%JVDvot7rZP;*j+|JtFqR&pjs>y#6nE?Gfi`i+cLQwTEyIR8 z>UV=Lg3c)YZ=U;%{?-q_bmd!L;N61{zWwd*T=~=A;e)ro`Hd@|ei=@!E1u=RxG z^E$e460^mUlxmHn^!bbele5UYkzrfT7nzN8rK(eNoiwen3720dw9ruL5+_Qh%%hO0 zTHy>d_f++S6FxoN7y=%KVUi~o*I0GNDp(M*a1fOl`Q+h!3>})796xeYKQwXlG2J&M z^&LeqqVo_&__i!#-Cc=bc&n(-LsSiZ66TW^d8;J!RIN#uK^8#@XArbdBY z&&(5#AAI2Gxb7=7X|?N<<40+>rTwtq-@9Yah~Bbw_P2KL*m71rr;A660wB|!CxT^k zQjU|%v_`aU&F;F}xEK4RXW8H~*)c$emmJcXk15qH=vf}jC= zAXro@1pHH8KjxPvk@&(*77n6!UMsp1tlU}Fs+f-bj;uRDR8C8^xG_{cI z=%IiO0m%9^HSeZ^?Mv8)dEuz_IvL_r?sFpZg)>gL3fBOWgl2*T*{7i{z{3f2Tr}8N zjpwJ~hXB!NQ&UXmQ;t^7pMdAoz_nEnENp5R2?fT(PR|q^0{}j1K@ zhD`nQG%$LY00dzmO^81?xH<%3_)Ojgo*9?AS=M|Gm9UNBlS!KSQ(#&iJkz{6g|Ql+ zrUuGvU^W?scgfcPW0}lfb{{w`WaHA{N}#5<^#VL#LjcJ90tx_|unJsG(1w6bCZdM3 zN^gWL-VLaiXV%};B1xq`q5}A6VGS+Scay(TH5@7FcX#}qnn|*KokmH{-PpVRDzzO_ zlJCa&sx0=x=XULClMQ`SjfBy{`%`bIaqSn`gb$#;jls16C?XNEvVV-9`4dck z06^_$3jm=5fRHHg|DY9yBN(n;9|Wm@6j2va!~iM&ik-z#5f`up(29Tp>H)Ii<_F!Y z6Y^@fDr3H@5^bcb5?}@336LuW>o_Hpwq6o;!gyN7mFxuD?nygoJl)7xF_8hFfGsV$ zFNM#$ccmUHYuSc42c-%v0!5twxJk7s>R2=NR0@ikx1>7x-f$yRkAY_Z95cR$9+Ko< zmgiQu(;o%yH1Iy|l!BYOZLR=-VoNE+#5(Tp)(@>xzjno5hB3~x1qb=>n|&`Z)F#HmU_ADXZydx?;R4dL!rht8oLaec>-)D z&&rMD*r`2kcn&%|4#Sp3YosX5M&Btht;;@qMFwzJv5Ho)V6wWqAwzO(1CNntcA~uB z7~~E7^@V(;Uh^++j2(od3t)uZTphFIs%;uMV+&H@ynDMXS^Ks?NMrA^wu=z` zC5E1!7p?P}t-))&0eOvE!}!K6&;kvZov`~KJ~9rZHFe7mk z&_lj-4|=SV(4!r7NX0*aKAr9Y-d>J=^{M%uhRKTRQxU#U#4P-ipwY*jlxdS+0Ld00y;*aTty(4s9(A|x>v^; z%;tm-*hA+;@Uj>eI!T|=OzGCiI5~ z0UOfF&@Z*JY%dph8N&7m#7+U&OYPyFIC>@R}r* zp>Ve4JNXfA)~5PNz|ijie&$gWlqs+jwnOkqC~HGYk;|cG*oM#K6hNdvwgC23><@}4 z@5iveo{}X{uVH_aNbr90MA+Ct?L`V12I^KilC!r#yo}cOu2-rc?+8eBH<_&Og?F%dNGe`b#$>@VgQ$m z<6xmSM}5=^b2z9~sGKX#qMR<*aN=}fx`qSu)SMG7mf=&&#c)`z-HGya{uElwA_9hk z>Kt8K-j_$2&(Ao~d|ec)xhPfE_EF`jBJEdcl;sn2CHgKq(NjfHsG_bW5O}Lv^BkUP zQSF!T#7?wSJ&pQ^LGzC|80=Ex_psYRLB`in_fz~>8T2z&B)k~>sS3Bco;xD z$UjDvL;&V*!j^uT9lj+*|1IUbjQK}J#D#HTz_{&h%wM(P#1~Q5Q-U)Ls-*Nyai){XS!M@g8B=y{O6RVx>)`6}EU`BAcW^t`EC2*gtvcPD{0g41=N zRK)!z;Vq}M9ry6491aZLD*OuF=sDMA{9C%gpHo3&Fn&zMPpIHhl6uF_dlOk0g&QFE z`)>~398yMa8h?rAZL*q_q=p~!Fa?c#9)`SbM#xaR96A{K#KzwUH&qyG5%@$Swz?)H z@Usy9JP%_G?WjJ=^)b{V@Ci>ZjI|47;zdrTY|05sH13})VonSCM~an&{9gFt7I9z| zPzQGyur5tLEMvTML(();q}M=ta3&hzvu7LOX0y?C_AI^SQBH;0(TL=lO&RwTtt}aifj@*&=5mn+gKT+9?z=ISKg7TqaSA&DJ57L7 z7-A89ycU({$0u7S<=zCj95qSO1865FX^{5bK(P+<^A3?x_*^;+l2f?$|6yi`pjMli zbyECM*lXF`zzYZ?z_4gQS5QFt{V_I;ENA600vzK&oSV8QYve>VC|PkA*Gw2!sfTM+ zbUfypsN00cG_KQ0I%<17issS;9Ki1#8U|%YLL0)hL0S!(N3|mh$^tUsDC%P?mlcWz zQ79C8CUzl0{Dt!|I|}3zzYzX(H^QZnB}z8XE(1LUb&CQU4zZa#m+1(tPNBa^VBvYA zTB+7eod}6gu>LpaJ`0F=Y& zPt*83brbKT!kFw;O34WNHvWR{+5V24g#SQDJqx`}D?rbIVQRf1`5M zH)s5JoXu(VN3jXc;=@dI?=g=|LDCX=|1){d?tAvGG?T{vzz__|PU(EeRKTj)M+EtQ zsL|h{XcO+_e<=)FqWYL}>4NXVQRJU4Vd+tTi9paSQU* zFc89#`H3Kr12k#u$4M+c2QLm~!VJQ9mXYC5PAcT_E{oi2O_>41TMlJ)*VCApD=)ly z<;^#)z4~p0tfcu^M=j^mLe;|^fTX9Q!ERsan_0&OFPj+NiB3PwCiKww-b%40*DE6< zOk}DDB2CiJro7)XUU-4b8thA$GXm~4XN)3e{RJaR%cDr!Ind);<|F!UM1U|>_T?<{ zDG9zP8iZ^*JSPQ-I2x2gOyn^PUrxnQnhHv|LlbGhyBPy9dh{aM?AXfK*1fu$FCcY^ zuWcpFBKNiNMJg^*!B+T_v_}x+JyQhMQKO$z@e3*#KK(uI%}~)H1A1Fy@*Gh2gZf>t z6bn;20R3su;xfk!HR0s1GU~n?nK)&=aMsu_(RCX?k#P6<>A&E>=cn;Deu|m{@9ELr zu(C(n{2k_QP3%oWplB zaILu1^DdgSQ4h%)pL&QMed;;$52l`v`qa}N$sni@g$!=5)sR+?^Z>y49!ecqirGnu zX)I{?C2HeHUID#9KErQ*@*W%6i=W7}w?m7R-{MiG?RM$i;qUjQH_>Q)MH)>m{;unT zAGA`C-hS*va*l-*l@1z3G^Ne=-d!i5aOH*1UwQrwKX>=q#b+hGUR`>*CCpiF8!w=_ z^rJVgcG=jK@|;BY{($stOa^Z52lxZG z-@g48clF%E9ruGLu}E3UB)Q@Iq32+62`d)2cT~DhP<|VKjfRaR2}+ro42$(J&Ck71a@*o^VtSH2@-q<`5K8@R&ccF`R z##6zBF@!5s)!jxk7-ky0l^JE<-)6dfhQvCwTQKA(&iV{lYGX(r9xa98xuhH~H+_+k z7Y$?-nOof~E8FzEnani>*-Rgnk)rcq>d#R=0dB*3gm&hsM1gvSYB3oLREZ2RD_aA%f=$H^#*?rqwfH%g@rJFm=;q*8*6Iobzhov5 zSh9@)H40~f0eQ#V=oFUrnqzZ?!B{jN$|9$1)kVZ>m@}> z7ZqprG|mNq6D^`@%g3S1h-D6YF_}M>R!Xn|D6k4r zY7N(8vnH|GNPx|iOH@sXJitkla`$&deIQD-}Iwb>rTdOmY~tc|%=z)D~qm6<0B z=83_b3%9V#JgB+2lbpp%$QC9>p^SQF5^4+ELnH(F_u+*I7%Blq3R6Zra)f2pebl!# zaFxstU_^x@^N9|oK}dol{m+1DysZ0f&KAzLWedxQac9oN^thvi!9>Wk;D zUV3pE8$CkwL5fTA2gEtp0ijryrv=n63-6K_g2TJpu*2cqZW-Q1aF@coVR-G71^CuQ z=#nafy*b6h{66s@&zKida*mqI4gN{#O{45NV-7dzfZrSs{t{6U5$(26V+#=(X&H~i zR*D8QU%?HAkO;4c)E=keMID= zA@cJhY6oN#W#^HGWgae$nS-n`cP&Fm5Z$K~;~ux(P2PfYiNt-<756EK`w?IzFDE`F z(mqy#^C*CoGA>ra6)skSkQP~=uv5Os=N7(!f+N%FAZ(xAV)I&f0T#Ywznk>6LRw{h*wFJAfm&%E`guZ+F*#v5-x`})|GuYTv+8-H^3#aB9o zzLRu`m5{q8_W?4P5J(?1Mi@gO+gPqLXBUyz4EY}I%z&kdv}li0f?x+=@g*Y2RjfX5 zY-?Swm+PjbPb0nBZ8=5_$GDxFdPD-Q@u?s(b1(12XRThn5t5Z9M|EqbRU1$vyG5e2 zvX>Zb$sN~C6~t(;ToJxOLnnepK{j>L{xBgTB)`h#`kq3`wD{GDz)jxCv|qr3ukN^N z@$h>Avu=#h$Wr06v=c3pgK@8V)_eT(nJqBSnrHdhwGGwI6m)y`U6 z*~aPQi13{ew}5A>aJt9ca}n_NHs;}=gP3?Y9l4=LDsVtVqV}% z$xzggxB|g0;sFEeT!6=@Xt~ZZ)YSV(g-i9X1!2@0YTv?Aw^mij)s6fWD11fI?hg*g zBkn)PAj}<@d{1qWhc?5$H41DeVdAGFZL1L};@C#i@7LxC5Va)aIe16SjZDTIjFY*d zZI(wIs%X1wc>qVj@<_Or2hPCq;IGrz-pXGh>2?oU+*d$oAFw%aw9Mv!FBS?&61B8B zG#TFz-OGqt1yQ%JP>j1a2XC3EvvsmLuG0GUuU!4|OHvKNQ`Tm4bPzbxc0j{$xWgmF zPO^T+^=Mg00fDIJ91q&Ei7mZm0b(y(GK&OMfLs$wA}L#=E;&;VEHL)FmbAczP-7=OMe9jD z(q7fxG19oWk!`=qz4E@0if;GMdfkYrmtRFOYGh6xr5f>q* zT!K~f$A=IC&SFqL+ zS#N03I`bHvdB)C~U*@M+gKs5aZK6y@$3lg)FBC8eDHQM>8G3>&kE9pt7h6YtVtwGK zAI~*_b&CI$p9!C{vU{OHy(b#92Qi8@%{BG_!mu0dHOO{bgS@si_^QjAg)aQCN)%-l zegORv)|SHc|7Ug|=1}4OZbyH8f&PZ{B+Y7J0LvGC0pS|*ui*~Ml_1v}Ze&uK1zj*bxUUV_LslWZh=dZr`wV!?e zx2`;Q;q5>8)?0t_qHHE}Z5R}h2L*X~Y%FzzfYnJ6OfZI zCVkU~R=AnQIUKb;6O;jfj8whvwB9FD2^A!O!vP7wibQZv@fwt3$1e~nnBQl(Y`KUQ zy~tNPNPD#6q>r8pOtAVbjcquKtkGVtTVuPwH$dI^ef7WZzWV$gR_klP>#H%dQk-wK zMZ|J(Y0y4yk2ZQX0^ku6aCx;bR-5LGKcWYriOa<$0OAlGW0wnb>SZd}79f9Z<4j;o zGLRAoc4BO>`lni77~=4#?+2CxLLvyS!SPpsw2g;Jo|`Q4I_Ujacx)g^b34BU#8nui z`D1iaf!iDJXN2sM`2Rda?sL0A@3xJ8@g`TIUy?^6v=-@oOuz28XgmPw6Y;8u@eZFt zKIpy;F9iL5S0?F)4?TM5sD38{RowV<7$_~E;$-M$kKm(U@~^T;&3KQ12SIi{N_=TT zM0^bFW&rtA1XlN7nOrc;{Ii1KVb;7856F;s%j)xK>?E3XL3@isSL<@0H)qS9&J!?g z9k$qXwGMyH>kxc1Ow$)7+Q(`hCI(Y3*0ylM*noi=8)>vM>St`F3cuD}Hcb-JRe$)p zQI*zcxZWnIZ9H$&P|Mq7JVV`gQ$cBp90o*$kd{kJPw^9|#piqqmYn_PsSz{vZ_^$# zH3vF462mVmM??E5oxOwho}@j_ufo%Ya+_B--T6sbiKpDBpig;8MNZEy10(d%PF5;J ztkRpB29xLKq`!fJCbZB5-wBJek|+pA!M@DuF&&L%3!P3#9>%VmN|pMLQZYisKcPYI zqB|+4jCZKl9vTud7_bx^>#r$uz|m;M7wsAUf|`GpihoAMs!sBDrT@Rszxhvb0eKr} zX%ziYq?59F29Z!eCJ?QHp*7P<Ut80Cv$<@C)MQVg&Z3P(df&Ggf%Ypk&5wv>pV zMMiKvT`7;Yv>k23k;YfZ-pQ|_MTW7kin69@MYjBl-581w{aJSZT|dg+?;Bsgj#1MC z`v2&)5(6B-&!l8rD+CiH#LI68HHC#@VP!7X`0__;EaYIi`wF*~jd`iY`leF;38FxQrb@}^gE938X;o~%38?4^-_q(kge?M#W z(ms7gGS@TVq6}awiXl32_uEx6(LG7zCBIhS@V<2Tyr*d*X99Wb=PDdjG zYt?#vP_6X9IvnTmd}t1P)LM2vKG5f*%3PmZT0(Ed(G2(7yZqr){tTUA;)L-W72l@f zM^sRzp+N~_2DxGl3Z@$rDKscTW03RKAZM3B#rdKhY+%{mHD>pG+p>NiB*0+JM%r St;3lfEvpS_ecC{J;{O5^w%nrt literal 1616 zcmZ`(&2Jnv6u0MdcXns92_a3|999qsO+gzWapX&+t>Dn8YLF0fSt zah15lV*#(lHPE@lXLTONA#cQua$aFGycsu3`zo8|bMag`4%kb4KAwM0Ut=%xg?NEq ziLda*c#$u~OJGwc;R6(3B@HqI-)qF)KvDB4oE}BCb_wkk1MN~Fa-p3x>*o>H)eS0k zsJNN!tVfm(*H}t3wLzs!^XzQ;ams04sHfO$gmrK?8Qe-4y*&`%Qj=75h9Vz`gzM%2 zCL$S}DI}OLOm&QL+3D_edWoVeAHYiUE*+#&QE{utdJ29BI{(H8ck-Oc$k7hTds0`7 z`)8ypb**1WN+eB0Z&L?~-cxBNsZc$0uCrEwK^5dr5P1mj1sVV~^NlZ{!}kBab*EC&ee$)WaY7pCAjo z9q0%4M2pb|dULvid~X+HWcsHdBOij{u4^0Ux&OMPjP_LPdaEyT-h%73b~mY@t$a75 zA|Mk*EERD6&o_MFS|%vxY;SlY-F>>1GM1Iy5Ta z6LJAI(RI1aaE_d zD^_5!2Xfxb^K?MgKqLWB=m3q8f_QhZvH&e*1KP>~bd(F|Di6?8K3JNJs1U3mq7WT3 zwp8_eY^&gW?5LU<6Bi=r6*u5lLG|>jA5MSz=JfEp*7*ba`}e~qzkPf9%Teog%619~ zY-SK-$m`aVqX&;4?mzzi;px{u{&n!jqx~;~s+fU|MAN{mfjNNnsBTITnyi^T7N$mt znt{3jLxTtbbVW>G&wuuaEOE`OTcTgaVFIs~E-X)od7s9hq}y{me7V=L=}jzN1B+-` zhskhwr|7b@*M>JM4Clf&V3^l)$tGkFl#ll0136M2U;9a#Zz?qup0V|ft+OSeq6Gb% z5y)^TKEhm22XLBc^|-jIu7eipMW;|HljX(r)r*syUA#Xn8xualGKJfuJ2drUlK5u# xHb^o|UqttWO>@;tOqtvj6T04>yqR}RfsxAq^QaEl*#t2pO~|U_MO?A!;lJYLr)B^E diff --git a/rag/init.py b/rag/init.py index a3025ee..3676ba2 100644 --- a/rag/init.py +++ b/rag/init.py @@ -138,26 +138,27 @@ async def doc_upload_handler(request, params_kw, *args, **kwargs): "UPDATE knowledge_bases SET doc_count=doc_count+1, total_size=total_size+${size}$ WHERE id=${kb_id}$", {"size": file_size, "kb_id": kb_id}) - # Trigger async ingest for text-based files + # Trigger async ingest for text-based files via uapi ingest_result = None if file_type == "text": try: text = file_data.decode("utf-8", errors="replace") - from pipeline import ingest as pipeline_ingest - ingest_result = pipeline_ingest( - text, pipeline_name="kg-rag-standard", - collection=kb_id, graph_name=kb_id, llm_func=None) + ingest_result = await _rag_ingest_async(env, text, kb_id, doc_id) # Update status to done async with get_sor_context(env, 'rag') as sor: + chunks_n = ingest_result.get("chunks", 0) if ingest_result else 0 await sor.sqlExe( "UPDATE documents SET status='done', chunk_count=${chunks}$ WHERE id=${id}$", - {"chunks": ingest_result.get("chunks", 0) if ingest_result else 0, "id": doc_id}) - if ingest_result: + {"chunks": chunks_n, "id": doc_id}) + if chunks_n: await sor.sqlExe( "UPDATE knowledge_bases SET chunk_count=chunk_count+${n}$ WHERE id=${kb_id}$", - {"n": ingest_result.get("chunks", 0), "kb_id": kb_id}) + {"n": chunks_n, "kb_id": kb_id}) except Exception as e: - exception(f"async ingest failed: {e}") + exception(f"uapi ingest failed: {e}") + async with get_sor_context(env, 'rag') as sor: + await sor.sqlExe( + "UPDATE documents SET status='error' WHERE id=${id}$", {"id": doc_id}) return json.dumps({ "status": "SUCCEEDED", @@ -196,17 +197,16 @@ async def doc_delete_handler(request, params_kw, *args, **kwargs): vector_ids = [c.vector_id for c in chunks if c.vector_id] if vector_ids: try: - await _call_vdb_async("/v1/delete", {"colname": doc.kb_id, "ids": vector_ids}) + await _call_uapi("rag-vdb", "delete", + {"colname": doc.kb_id, "ids": vector_ids}) except Exception as e: exception(f"vdb delete failed: {e}") - # Delete from graph - entities = await sor.R("entities", {"kb_id": doc.kb_id}) - if entities: - try: - await _call_graph_async("/api/graph/save", {"graph": doc.kb_id}) - except Exception as e: - exception(f"graph cleanup failed: {e}") + # Delete entities from graph + try: + await _call_uapi("rag-graph", "delete", {"graph": doc.kb_id}) + except Exception as e: + exception(f"graph delete failed: {e}") # Delete DB records await sor.sqlExe("DELETE FROM document_chunks WHERE doc_id=${id}$", {"id": doc_id}) @@ -270,6 +270,93 @@ async def _call_uapi(upappid, apiname, data, timeout=10): return await resp.json() +async def _rag_ingest_async(env, text, kb_id, doc_id): + """RAG ingestion pipeline via uapi: chunk → embed → VDB → NER → graph""" + chunks = _split_text(text, chunk_size=512, overlap=64) + if not chunks: + return {"chunks": 0} + chunk_count = len(chunks) + + # 1. Embedding + try: + emb_resp = await _call_uapi("rag-embedding", "embed", + {"texts": chunks, "model": "CLIP-ViT-H-14"}) + embeddings = emb_resp.get("embeddings", []) if isinstance(emb_resp, dict) else [] + except Exception as e: + exception(f"embedding failed: {e}") + embeddings = [] + + # 2. VDB upsert + vector_ids = [] + if embeddings: + try: + vdb_data = { + "collection": kb_id, + "data": [{"id": f"{doc_id}_{i}", "vector": emb, "text": chunks[i]} + for i, emb in enumerate(embeddings)] + } + vdb_resp = await _call_uapi("rag-vdb", "upsert", vdb_data) + vector_ids = [f"{doc_id}_{i}" for i in range(len(embeddings))] + except Exception as e: + exception(f"vdb upsert failed: {e}") + + # 3. NER entity extraction + entities_found = [] + try: + full_text = " ".join(chunks[:20]) # first 20 chunks for NER + ner_resp = await _call_uapi("rag-ner", "entities", {"text": full_text}) + entities_found = ner_resp.get("entities", []) if isinstance(ner_resp, dict) else [] + except Exception as e: + exception(f"ner failed: {e}") + + # 4. Save to graph + if entities_found: + try: + graph_data = { + "graph": kb_id, + "data": {"entities": entities_found, "source_doc": doc_id} + } + await _call_uapi("rag-graph", "save", graph_data) + except Exception as e: + exception(f"graph save failed: {e}") + + # 5. Record chunks in DB + async with get_sor_context(env, 'rag') as sor: + for i, (chunk_text, vid) in enumerate(zip(chunks, vector_ids)): + await sor.sqlExe( + "INSERT INTO document_chunks (id, doc_id, kb_id, chunk_index, content, vector_id, created_at) " + "VALUES (${id}$, ${doc_id}$, ${kb_id}$, ${idx}$, ${content}$, ${vid}$, NOW())", + {"id": f"{doc_id}_c{i}", "doc_id": doc_id, "kb_id": kb_id, + "idx": i, "content": chunk_text[:2000], "vid": vid}) + + return {"chunks": chunk_count, "vectors": len(vector_ids), + "entities": len(entities_found)} + + +def _split_text(text, chunk_size=512, overlap=64): + """Simple text chunker: paragraph-based with size limits""" + paragraphs = text.split('\n') + chunks = [] + current = "" + for p in paragraphs: + p = p.strip() + if not p: + continue + if len(current) + len(p) < chunk_size: + current = (current + " " + p).strip() + else: + if current: + chunks.append(current) + current = p + if current: + chunks.append(current) + # If still no chunks (single giant paragraph), force-split by size + if not chunks and text.strip(): + for i in range(0, len(text), chunk_size - overlap): + chunks.append(text[i:i + chunk_size]) + return chunks + + async def _render_tmpl(tmpl, data): """Simple Jinja2-style template rendering for uapi data templates""" import re