From 598f5c5c23a841b35beb2e2dc8cf0f05dea70640 Mon Sep 17 00:00:00 2001 From: yumoqing Date: Mon, 3 Aug 2026 10:40:25 +0800 Subject: [PATCH] sync: add batch_ingest.dspy from test env --- .../knowledge_bases_list/batch_ingest.dspy | 204 ++++++++++++++++++ 1 file changed, 204 insertions(+) create mode 100644 wwwroot/knowledge_bases_list/batch_ingest.dspy diff --git a/wwwroot/knowledge_bases_list/batch_ingest.dspy b/wwwroot/knowledge_bases_list/batch_ingest.dspy new file mode 100644 index 0000000..00e1631 --- /dev/null +++ b/wwwroot/knowledge_bases_list/batch_ingest.dspy @@ -0,0 +1,204 @@ +ns = params_kw.copy() +kb_id = ns.get('kb_id', '41c9f1cd45e0') +env = request._run_ns + +# Get unprocessed documents +rows = [] +async with get_sor_context(env, 'rag') as sor: + recs = await sor.sqlExe( + "SELECT id, file_name, file_size, file_path FROM documents WHERE kb_id=${kb_id}$ AND (chunk_count IS NULL OR chunk_count=0 OR status='pending') LIMIT 20", + {"kb_id": kb_id}) + rows = [dict(r) for r in recs] + +if not rows: + return json.dumps({"widgettype": "Text", "options": {"text": "No unprocessed documents found"}}, ensure_ascii=False) + +import aiohttp, io, os + +FS = '/d/rag/ragserver/files' +LOCAL = '/d/rag/ragserver/pkgs/rag/rag/files' +results = [] + +async with aiohttp.ClientSession(timeout=aiohttp.ClientTimeout(total=60)) as session: + for r in rows: + fid = r['id'] + fn = r['file_name'] + fp = r['file_path'] or '' + + # Resolve file + fpath = None + if fp.startswith('/idfile/'): + fname = fp.replace('/idfile/', '') + if os.path.exists(os.path.join(LOCAL, fname)): + fpath = os.path.join(LOCAL, fname) + else: + full = os.path.join(FS, fp.lstrip('/')) + if os.path.exists(full): + fpath = full + + if not fpath: + results.append(f"✗ {fn}: not found") + continue + + ext = os.path.splitext(fpath)[1].lower() + text_exts = {'.txt','.md','.csv','.json','.xml','.html','.htm','.py','.js','.css','.yaml','.yml','.log','.rst'} + with open(fpath, 'rb') as f: + data = f.read() + + text = '' + if ext in text_exts: + text = data.decode('utf-8', errors='replace') + elif ext == '.pdf': + try: + from PyPDF2 import PdfReader; reader = PdfReader(io.BytesIO(data)) + text = '\n'.join(p.extract_text() or '' for p in reader.pages) + except: text = '' + elif ext == '.docx': + try: + from docx import Document; doc = Document(io.BytesIO(data)) + text = '\n'.join(p.text for p in doc.paragraphs if p.text.strip()) + except: text = '' + elif ext == '.pptx': + try: + from pptx import Presentation; prs = Presentation(io.BytesIO(data)) + parts = [sh.text for s in prs.slides for sh in s.shapes if hasattr(sh,'text') and sh.text.strip()] + text = '\n'.join(parts) + except: text = '' + elif ext in {'.xlsx','.xls'}: + try: + from openpyxl import load_workbook + wb = load_workbook(io.BytesIO(data), read_only=True, data_only=True) + rows2 = [' | '.join(str(c) if c else '' for c in row) for ws in wb.worksheets for row in ws.iter_rows(values_only=True) if any(c for c in row)] + text = '\n'.join(rows2) + except: text = '' + + if not text or len(text.strip()) < 10: + ext_l = os.path.splitext(fpath)[1].lower() + image_exts = {'.jpg','.jpeg','.png','.bmp','.gif','.webp'} + audio_exts = {'.mp3','.wav','.flac','.ogg','.m4a','.aac'} + video_exts = {'.mp4','.avi','.mov','.mkv','.webm'} + face_count = 0 + speakers = 0 + import base64, subprocess + + # IMAGE: face detection + if ext_l in image_exts: + try: + img_b64 = base64.b64encode(data).decode() + r_img = await session.post('https://media.opencomputing.net:10443/face/api/detect', + json={"images": [img_b64]}) + if r_img.status == 200: + fd = await r_img.json() + res = fd.get("results", []) + if res and isinstance(res[0], dict): + face_count = len(res[0].get("faces", res[0].get("detections", []))) + except: pass + + # VIDEO: frame → face + audio → voiceprint + if ext_l in video_exts: + try: + tmp_img = '/tmp/' + fid + '_frame.jpg' + subprocess.run(['ffmpeg','-y','-i',fpath,'-vframes','1','-q:v','2',tmp_img], + capture_output=True, timeout=30) + if os.path.exists(tmp_img): + with open(tmp_img, 'rb') as fi: + fd2 = fi.read() + img_b64 = base64.b64encode(fd2).decode() + r_vid = await session.post('https://media.opencomputing.net:10443/face/api/detect', + json={"images": [img_b64]}) + if r_vid.status == 200: + fdd = await r_vid.json() + res = fdd.get("results", []) + if res and isinstance(res[0], dict): + face_count = len(res[0].get("faces", res[0].get("detections", []))) + os.remove(tmp_img) + except: pass + try: + tmp_audio = '/tmp/' + fid + '_audio.wav' + subprocess.run(['ffmpeg','-y','-i',fpath,'-vn','-acodec','pcm_s16le','-ar','16000','-ac','1',tmp_audio], + capture_output=True, timeout=60) + if os.path.exists(tmp_audio) and os.path.getsize(tmp_audio) > 44: + with open(tmp_audio, 'rb') as fa: + aud = fa.read() + form = aiohttp.FormData() + form.add_field('file', aud, filename='audio.wav') + r_au = await session.post('https://media.opencomputing.net:10443/voiceprint/extract/submit', + data=form) + if r_au.status == 200: + vd = await r_au.json() + if vd.get('status') == 'SUCCEEDED' or vd.get('embedding'): + speakers = vd.get('speakers', 1) + os.remove(tmp_audio) + except: pass + + # AUDIO: voiceprint + if ext_l in audio_exts: + try: + form = aiohttp.FormData() + form.add_field('file', data, filename=fn) + r_au = await session.post('https://media.opencomputing.net:10443/voiceprint/extract/submit', + data=form) + if r_au.status == 200: + vd = await r_au.json() + if vd.get('status') == 'SUCCEEDED' or vd.get('embedding'): + speakers = vd.get('speakers', 1) + except: pass + + import json as j2 + meta = j2.dumps({"face": face_count, "speakers": speakers}, ensure_ascii=False) + async with get_sor_context(env, 'rag') as sor: + await sor.sqlExe( + "UPDATE documents SET status='done',chunk_count=0,metadata=${meta}$ WHERE id=${id}$", + {"meta": meta, "id": fid}) + results.append(f"🎬 {fn}: faces={face_count} speakers={speakers}") + continue + + # Chunk + paragraphs = text.split('\n') + chunks, cur = [], '' + for p in paragraphs: + p = p.strip() + if not p: + if cur: chunks.append(cur); cur = '' + continue + if len(cur) + len(p) < 500: cur = (cur+'\n'+p).strip() + else: + if cur: chunks.append(cur) + cur = p + if cur: chunks.append(cur) + + if not chunks: + results.append(f"~ {fn}: empty chunks") + continue + + # Embed + try: + r2 = await session.post('https://embedding.opencomputing.net:10443/api/embed', + json={"texts": chunks, "model": "CLIP-ViT-H-14"}) + emb = await r2.json() if r2.status == 200 else {} + embeddings = emb.get("text_embeddings", emb.get("embeddings", [])) + except: + embeddings = [] + + # VDB + if embeddings: + try: + await session.post('https://vectordb.opencomputing.net:10443/v1/upsert', + json={"colname": kb_id, "data": [{"id": f"{fid}_{i}", "vector": e, "text": chunks[i]} for i, e in enumerate(embeddings)]}) + except: pass + + # DB + async with get_sor_context(env, 'rag') as sor: + for i, ct in enumerate(chunks): + 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"{fid}_c{i}", "doc_id": fid, "kb_id": kb_id, "idx": i, "content": ct[:2000], "vid": ""}) + await sor.sqlExe("UPDATE documents SET status='done',chunk_count=${n}$ WHERE id=${id}$", {"n": len(chunks), "id": fid}) + await sor.sqlExe("UPDATE knowledge_bases SET chunk_count=chunk_count+${n}$ WHERE id=${kb_id}$", {"n": len(chunks), "kb_id": kb_id}) + + results.append(f"✓ {fn}: {len(text)} chars → {len(chunks)} chunks") + +return json.dumps({ + "widgettype": "VBox", "options": {"padding": "20px"}, + "subwidgets": [{"widgettype": "Text", "options": {"text": r, "cfontsize": 13, "color": "#333"}} for r in results] +}, ensure_ascii=False)