sync: add batch_ingest.dspy from test env

This commit is contained in:
yumoqing 2026-08-03 10:40:25 +08:00
parent 8ab8987776
commit 598f5c5c23

View File

@ -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)