From 850e96e0ea3a8775a5274cda78d9980d715068d5 Mon Sep 17 00:00:00 2001 From: yumoqing Date: Mon, 27 Jul 2026 14:08:47 +0800 Subject: [PATCH] feat: uapi-based unified search with text+media, multi-KB, rerank MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - search_handler: query(opt) + media(opt) → embed → VDB search → rerank → enriched - _resolve_search_kbs: specific KB or all org KBs - _build_search_vector: text + image/audio/video/text file → CLIP embedding - _parse_vdb_hits/_apply_rerank/_enrich_search_results: result pipeline --- rag/__pycache__/init.cpython-310.pyc | Bin 23423 -> 25980 bytes rag/init.py | 218 +++++++++++++++++++++++++-- 2 files changed, 204 insertions(+), 14 deletions(-) diff --git a/rag/__pycache__/init.cpython-310.pyc b/rag/__pycache__/init.cpython-310.pyc index 4094703658c8e689a271eb00cc87a143be0257ef..f655d8a58d9b59c8188e27c33df64443ccfe28df 100644 GIT binary patch delta 6414 zcma(#ZE#yvcJDq(PqMx&S+XU|mi_z{M|K?Baboge$0Uw}&4O)WlMn*cD!NLd#Fm|_ z=RlMvWN=7Am(ajvKL7*H4AY(LPHCuicDMV3?xx$$HoK*}l$~kc%oesYZ5P^M%dn*d zmZaz0Cp#u&rnN`s9o>8G`MT%p>H2rc@2?W?jK|}Y;BV@KLq|{LpZ3<0T~D_!6$WUC zhEJ3GL$sPkPD}TP>FqR1>tLmt)>9Q$B6NT@&_-CP0j!BO16B)I3ta?WA4sY^2?E13a5(58VjQX1a;?;(-==8{G`gHFOKz3eQ%$ zjrPHFE$ye<;Tfa$2`RRtc#!zLBv;hR7YKQr|Bo-~8@EW(qePdEN_p!Ak~>vw@jpij z$8-1Tj-y}}nsi>y2ehzm)dRYsmn~RFBua**e1%r2Td99R-YA(X&v)hhdIkG$y<$nG zHMI7koFTL}AJqLzRWzyx^)gzwC|{IKY9i^u3xt`w^{}f@b$ORBRD=$=yGwj8kZj9C zgh=?mA3pn}%#c;msxOd}vqZ{Q5-GP?i_iwW@&)Nd87z~&pjj1YhFUbzCcP|8c)lVU zTFpOkgFk;E5FKgO%XGV5uKV-~+Onw7HChdAJ#X13O-Lol-z-UvfaJ9(nZ~r*1qn2; z3JokvR;p+;&)5xVa;%66DLwMo=HrnOoD8jrWtz*3CN+6V7{X|TIMiqjc52TJe7Q* zI!qS%V)Z>6f+Q#hiIvF2L7dkV%Q7aOs7Jt3JioY89($HhnO+BY2y=<4XK=MxA$sWsda39av4v|J>mL~uX~K2S0}0{&8U7bHT17G9MIc-*2@b8qw*d6ws* zuam#xTk3k8b_lwT6ZD{=@Z)vu>L+FR{V+?+9x-glM;|0Wii9;9AbnnWh-3o$Lu*6G>qaNx`1wX|=TwG2?z7cgqp@5ZDon0x+EMWMXP6 zj!7suvc0&aAUKXdWYaieBFvpx$OtX_P+B{}p2yV}5WI-sB?K=cSOQ=;4kogx3~z0y zBVPXWNCcdaXb6O!!9C}RR8pH^kg{2p#1l{P3k}`mHGZX`&94wA@m}}J5#pGauglh2 z{#iqH-P*U6H5DWRdsf(UlCbQPt4JNvTw_=G#>N-PIsRtjzqxK`wa9<57JZ?|qNAd<8KSpEwD;)c%xe6ouo z2Q^72lu(&k&Rd?4&PiD>wNeQE2}vX8F#D{z8cn7OwV#uygO-)@&IwJ#rdc4gtP3&( zt!>FhU56!X8;7O50%g9&T97p5Mo9a6Gfvs1TjU#{{d(SC_K^Y%yV8Jqchh&&|0 zjzd@2Q@K`MJ}aFihh*xZ-qQ{ZRu>2s1X{i*=gUATAN50AV;MkooIpz;*O9hNBf3&p zAm)L*1N6ZPp`1KlQpW@81a7Fa51P|0RdMB-XmH61wnI^PT&|WblDr#4Ln)wP-CZgu zDBc6!@Ca{Q);-{j(BZDkm&I+IdQq!TJ;{Z>`< zD)pbwJfo(j52mQV3VV0$6v*h%9`O`4>0Jrp{@vSwX-qbp;#@qN&ZP|d%!w46O3X4$ zZ$<#*JXqHRLkca{un9Eh8N6#`ykl>AvSYZTdn0hMY0z1-hHZKV>fNZ}ES)H14rg-} zE5}y`gCz1^{(})t!b*jynQa4bSkk2j{or(xmxa3_LDAj13>QPHVSZb*JS`uW!8fu& z;xUUDB@D$ZN2~)_uu1?iOMxAMr$N%t`KJ$Ff?);sa1;;(tO~I-fL2f_i%KYxLh7#P zTD~w;tnZteff}3L4oEtyC1<8*x5M;v3OR`SQv}RV=4P`grY+l7fHeudE!p5B42%FG z3@pn-b%9ta8CNohuzv*tQ+-)z&wc|7F>))d+fi%e9+7eZXhy0jeZ|e66o~q65cYF+ zyv9t58IDv&OKa&=R=6qWhP4%PFz012jwNAknn}Xa_hHGj@O!vv#px**Fq<3V zE5}(FC0!aBofz6Psg8_Jj;Z{!po1qm+q^Rw4Fs=K%`4&$>})ES!(P65o!Wfr*g!&Dz%hsp(X8yI z=&Ud?<~0hgcy@Lwt;I2NnP`pJEDUElo6bOc%p_Cn5G=5H1S{>-C`)DLrc*4Tr3@QO zz`Sd=O>D)6C6}HR{Wcb0`|y|uMO4fjKV4?7C1>U`S^=(N`e2-;va{@8kV16O9|C4L z;3~z?YGO!Yc$e=jeh5xe7+mL|-=pq^1MX1pF`VL7%z1Y|^2E?t&TwJp?{F1!gB9LG zC=%e00rO)hz2F%gDSZQ9avYg!Q(DX)@@KwWo(@x$hBh< zYt?PkQh+Z?ND`#B<1XFC*6Y?KXy-UJ>Pp^R+q`Pcn%xURfi>H!S*aCRLMPUt1%#ot zgJ@6QE~q7V_p<9C%xsoKIzFshp;p;)n=fIo&L$u`adhm!VT-Cb zmZ?~f)v5G!8p6O% zYZ-C>62*ub=p_8V);2`nf`uc<$qR@)?*LvtIHBU1xA@&{y*wQoFK7Rb=N?Az9RU1| zSfr41cAn_&G-qap6?GJf_B@hca$d5TgQ}4A1}qpZXnpu17oUc%C1ybp+yY>18|CI; zj2R`~4bgocVlc_`ZiwO_VHo+#p~mlPOE$9qgl(oFK)qs?TzMym=@eJgi6xJ}+19e1 zVL2RSUO3Lm5!@u%Z%Q^(Qpi4uR22vUq}b6e6Cb__vY#XP1%kij1M7Cjz6R(e`^}ar zb8KoVeUSCRu3=%Rk+7gtKQn8(J>|0&*S-yYu`Zd9DjX7nCH7xUjrU&L?Eh9Gyma8n3wn1 z_|+~i`6~a%t{&3D6aCR*TlZ-~=J}s)*lAy5lJg&pgn%?azQvn*T1hK^e;`^s(xVY_ zmj7zwZ{-HU-`G@B4X?IR&IzyUacsg{)tIOBVg$UxJksm5$B-XTxAzXpF2ZMf+sU8v z#ohtZDOi5on~-C|9t8<+7cQJkM#AE4g73hBUt{ezApog z?rG}xllS|Lweba@~^KIK0ih>_mn`vZRh zM~aI(Y6xlL-`}~j8I3b6@HV*0xLeufx(!)wMgTuP4BjTYWd8EtX0oIB4}&|1PwdX2 zQ%WSUp=VreNZp5^pa1nxn;m7b0m3hDXaFspym5Dpe-&FK+%Ji%P3^9fn`HjMNPz#* z?x;LyM%UHdUBuMGl*Dbr-&0H>{L1hI*|qFGd6&#Lj%*_v`Rqsy*~y>T%wua^rFe{H;m?;mfj z+JakH5iVIwG4iy=`uN-9A>}BtyOUoYk5-JL9)xXOCT3%T;@B4`w%tTAWl{{^lWzYm z%13cAI}@%O!-Yu%`12XNfM6Oy4}WA&6Tg1Xfo2pKv)x2)xM(Vygl?k28-^f$glrT3 zsGewc3)Us8_-kWT7}UyRrW!{l?kam0C42~g|K-G37$ukH@>Qx-f`81D&Od~cVsvfh zOOrK$J*W`8MWdM|Z!D84UonNeG#ULIS;-|UWKWu8znzS%5mUY>;JcCG9mw2j`rvL{ zFjsl!y;~r|KHC%?w&AMyq2yjr+_1n6o0H`kv`V<~JYvF)6NvTT3E|{>5UWD)p~QcB z?}!V}Ge9{but;+yn@L@D+?iKrj+WB(|EMEX&@7Y}{<(-4`h7 z3rP%ub_PX`Qqj;RGxZO~5p`CvRqIHlQ^#8SrM0s&75$^NAKG?m{isstIqwCG=s(@e z``z=-J@=e*&pqd!m-DZ(i>H}m&~CR#_^T-F>i^-;{f-b@a(~10*eYJZeFs=?1^4s7 z0jbx=SMeYZA;*tYm`9KbAT^IikqRPJ$*YhG@oHX!@?l=f=i?aRb^IzE=kXY?$1%zq z_yQa&`9i)3$12{)7vosXoA?qOYxq*$d_eN`*79Y1IWp$+6?`R*b-cJ+inZhw=5nyl z@&`*N89N9TSJ-JbNYVsTrG9DDc!Z7IoL}KO&c-&6#8q=YZh|&&SRQqg`x)PWf7HVy z)uXz3ST(68&UoZ-(P)|CRgK&=DK|^)5@**-4=)*YsjlgAK2I$Z7uB+9nMZkWM$R%G zR7-f}lsqHrlBlih2opg~y6Tdoc}7O%)l)K@amU6M>}5=%zaD&Uma+k(qy&zzeH~01 z^)hM1t^|3F>U~5yA>+1Cq6eKA_3rm(i}Z0X+bb!Cej~3{OVnbuRCTIleEyV)*C`=> z)nUWUQnxg#b8=Q^0CmPlC$IM?;YkUV8%eohm^N~g5`t}>O4g7c^?c0OAyCTC)Cq=p z5Mu>XI=efSQ^Y=cHHvTJAXI`;g8i_jqSIC|CHsc6x94KAcmR%8>?rppa{IE$sJ=A1 zE0NBmcuN$n^?BWgXxY64hY0`yGQNO4B0A`OGrox439wp=1}Tit~h$ z5Xa!KKkB%TxNO&~MtI&Ic9ssMhf>(%(+_=cbWM*V0X1d@p8jU+bai=-J%;WM!qXCWv7KKaEfEAx#;*)zr{!(_e!XE-n>>Vf#Ce}@G z8@C^j)ROhmBpbD;7BMm6P)yu`iC3&zcMmP|5O8g|vCXLlG- zN_bSYVm4PQ{yE7^D-22iZV%tij=@{u7ui`j7Fk+kY24G)2)~UqEH9W-StYHcQ#ldD z1O2SDF)@^G+$|DAeT})qo>ZYoqesfMA)3sj64~J)h|JrSzh|C_HQbMj#DfI1e8k+C z=9H0*;yc9aM~E530MWt2i9HD)SJt~i%)#{9@M+Z( zfTk+HvHopS9c-!cvzMW->IAz7Ox5q<(coImuM_ZT^#b-N)YjCpNAsI%3SF7D(65# z9yiNItYaUGVvUP8C79p53|21f>MxjvhjJ;Qh;N}-!NgOUl#+@WMzplhr*t}(N{D2i z$fMXuu}JJ*h)1>Mq%0BPC>0JN_-=PbBdkDRe%MWECnDe<9xU zQr@==FRWkD%dZK3L-0D(FW(URE|MpUuXNHhEHat&PH{EP77QYFhj<%~FAulWlb23p znq>?6i1*^fGW53QqWB1zNo=K;f5k}Ae%9e5{wBdcmn*?5J)o)g3H~L)*oug?nSB13 z7@t6sBb@)~ii`|f;~win#5+##Fm$!7bbd#0` zxCYL)ES1kouxVo`?_GU>k+P$0MbOc1ffufh%Kw$%?W-GEH#l1BqK(iX~hwtZl#CFhP?=*ca17f?9;2FzVflLM7&IF}vvy34^~{ zLisi8*N|sp*FDIraB*vhoq#19TA2lgHdHZefg5(ISu<<&E z+Kya*AXtpfj@XpM?zo=PEm3g;+KyQ+?S~@Mh{8=>VL0CIWASN3i5!Wg#acZ9+MltB19D&?+ttBqhO4QvP<%! z4k=B@^A#cYSB1Q~P}8<_!!zUv*^LzoEk2#ZxQRgf(9$}W)|srZ`_^js^R{hfG8-!r z3^v^w#2nmx>wY5{!EO6NzO6=&vkM;IUTJF3Ep^}4VW$pKC{7M$73`&OYHI*~a@#UD z8NthJu2nd-J<66s!}c0G-8A@hs5w?f3QL+bg|pTiVy=Lj1U4+YB!C+VkjbK5iDW8Q zFzp#k#;OAMV0p7^Td1V!Oaf-g#yvKtF{D3ex%)@8pz3_~uZ z=zBOs;#%ai_Nzs%n|S!SHXGL%Q6CVmfx)SGdlmVhUtwgt;2cQtbb=aceCNJ6)+rW% zf5)n7%{?Q@77;8Uc#PC6q*ZK>_)vo39Z~B`#G~FUp28b33_srCkpIo#7dxJ2tMhjz zUSX{VY4Q}o%LL~ME)bYVX94ElUvmwv kHZq&R#%$Q#ZP?w-2D8mvYj&8+&0EZ^7N^+*+f$MM0r5_ 10: + # Check if it's a form-encoded request + if not query and b'=' in body[:100]: + pass # form data, not a file + elif not body.startswith(b'{') and not body.startswith(b'['): + file_data = body + file_name = params_kw.get("file_name", "search_upload") + + # Process query + media → embedding vector + query_vec = None + if query or file_data: + query_vec = await _build_search_vector(query, file_data, file_name) + + if not query_vec: + return json.dumps({"status": "SUCCEEDED", "data": {"results": [], "total": 0, "message": "no query or file provided"}}, ensure_ascii=False) + + # Multi-KB VDB search + all_hits = [] + for kid in kb_ids: + try: + vdb_resp = await _call_uapi("rag-vdb", "search", { + "collection": kid, + "vector": query_vec, + "topK": recall_k + }) + hits = _parse_vdb_hits(vdb_resp, kid) + all_hits.extend(hits) + except Exception as e: + exception(f"vdb search kb={kid}: {e}") + + # Deduplicate + sort by score + seen = set() + unique_hits = [] + for h in sorted(all_hits, key=lambda x: x.get("score", 0), reverse=True): + hid = h.get("id", h.get("text", "")) + if hid not in seen: + seen.add(hid) + unique_hits.append(h) + + # Rerank if query text provided + if query and unique_hits: + try: + documents = [h.get("text", h.get("content", "")) for h in unique_hits[:recall_k]] + rerank_resp = await _call_uapi("rag-reranker", "rerank", { + "query": query, + "documents": documents + }) + reranked = _apply_rerank(unique_hits, rerank_resp) + unique_hits = reranked + except Exception as e: + exception(f"rerank failed: {e}") + + # Limit + enrich with DB metadata + final = unique_hits[:top_k] + enriched = await _enrich_search_results(env, final) + + return json.dumps({ + "status": "SUCCEEDED", + "data": {"results": enriched, "total": len(enriched), + "recall": len(all_hits), "kbs_searched": len(kb_ids)} + }, ensure_ascii=False, default=str) + except Exception as e: exception(f"search: {e}, {format_exc()}") return json.dumps({"error": str(e)}) +async def _resolve_search_kbs(env, userorgid, kb_id): + """Resolve KB IDs: specific or all org KBs""" + async with get_sor_context(env, 'rag') as sor: + if kb_id: + recs = await sor.R("knowledge_bases", {"id": kb_id}) + return [r.id for r in recs] + # All org KBs (global + org-specific) + sql = "SELECT id FROM knowledge_bases WHERE org_id IS NULL OR org_id=${org_id}$" + recs = await sor.sqlExe(sql, {"org_id": userorgid}) + return [r.id for r in recs] + + +async def _build_search_vector(query, file_data, file_name): + """Build search embedding from text query + media file""" + texts = [] + if query: + texts.append(query) + + if file_data and file_name: + ext = os.path.splitext(file_name)[1].lower() + if ext in ('.jpg', '.jpeg', '.png', '.gif', '.bmp', '.webp'): + # Image → embed as-is (CLIP multimodal) + texts.append(f"[IMAGE:{file_name}]") + elif ext in ('.txt', '.md', '.json', '.csv', '.html', '.py'): + # Text file → extract content + try: + content = file_data.decode("utf-8", errors="replace")[:4000] + texts.append(content) + except Exception: + pass + elif ext in ('.mp3', '.wav', '.flac', '.ogg'): + # Audio → placeholder (would need ASR service) + texts.append(f"[AUDIO:{file_name}]") + elif ext in ('.mp4', '.avi', '.mov', '.mkv'): + texts.append(f"[VIDEO:{file_name}]") + else: + # Try as text + try: + texts.append(file_data.decode("utf-8", errors="replace")[:2000]) + except Exception: + pass + + if not texts: + return None + + combined = " ".join(texts) + try: + resp = await _call_uapi("rag-embedding", "embed", { + "texts": [combined], + "model": "CLIP-ViT-H-14" + }) + vecs = resp.get("embeddings", []) if isinstance(resp, dict) else [] + return vecs[0] if vecs else None + except Exception as e: + exception(f"query embedding failed: {e}") + return None + + +def _parse_vdb_hits(vdb_resp, kb_id): + """Parse VDB search response into uniform hit format""" + hits = [] + data = vdb_resp + if isinstance(data, dict): + for key in ("results", "data", "hits", "rows"): + candidates = data.get(key) + if isinstance(candidates, list): + data = candidates + break + if not isinstance(data, list): + return hits + for item in data: + if isinstance(item, dict): + hits.append({ + "id": item.get("id", item.get("doc_id", "")), + "text": item.get("text", item.get("content", "")), + "score": item.get("score", item.get("distance", 0)), + "kb_id": kb_id, + "metadata": item.get("metadata", item.get("meta", {})) + }) + return hits + + +def _apply_rerank(hits, rerank_resp): + """Apply reranker scores to reorder hits""" + scores = [] + if isinstance(rerank_resp, dict): + scores = rerank_resp.get("scores", rerank_resp.get("results", [])) + if isinstance(scores, list) and len(scores) == len(hits): + for i, s in enumerate(scores): + if isinstance(s, dict): + hits[i]["rerank_score"] = s.get("score", s.get("relevance_score", 0)) + else: + hits[i]["rerank_score"] = float(s) if s else 0 + hits.sort(key=lambda x: x.get("rerank_score", 0), reverse=True) + return hits + + +async def _enrich_search_results(env, hits): + """Enrich hits with document metadata from DB""" + doc_ids = list(set(h.get("id", "") for h in hits if h.get("id"))) + if not doc_ids: + return hits + + async with get_sor_context(env, 'rag') as sor: + recs = await sor.sqlExe( + "SELECT id, file_name, file_type, file_size, status, kb_id, created_at " + "FROM documents WHERE id IN (" + ",".join(repr(d) for d in doc_ids) + ")", {}) + doc_map = {r.id: dict(r) for r in recs} + + for h in hits: + did = h.get("id", "") + if did in doc_map: + h["document"] = doc_map[did] + return hits + + async def doc_upload_handler(request, params_kw, *args, **kwargs): """文件上传 → 保存 → DB记录 → 触发RAG入库""" env = request._run_ns