feat: async upload — return immediately, ingest in background
This commit is contained in:
parent
cdc02e31e4
commit
8dad4e9b5e
@ -22,27 +22,38 @@ file_size = len(file_data)
|
||||
ext = '.' + file_name.rsplit('.', 1)[1] if '.' in file_name else '.bin'
|
||||
ext_l = ext.lower()
|
||||
|
||||
text_exts = {'.txt', '.md', '.csv', '.json', '.xml', '.html', '.htm', '.py', '.js', '.css', '.yaml', '.yml', '.log', '.rst'}
|
||||
image_exts = {'.jpg', '.jpeg', '.png', '.bmp', '.gif', '.webp'}
|
||||
audio_exts = {'.mp3', '.wav', '.flac', '.ogg', '.m4a', '.aac'}
|
||||
video_exts = {'.mp4', '.avi', '.mov', '.mkv', '.webm'}
|
||||
# ============================================================
|
||||
# BACKGROUND INGESTION — runs in a separate asyncio task after
|
||||
# the response is returned (background_reco = create_task wrapper).
|
||||
# SELF-CONTAINED: never touches env/request (invalid after response),
|
||||
# creates its own DBPools connection. All args are primitives.
|
||||
# ============================================================
|
||||
async def ingest_doc(doc_id, kb_id, file_name, ext_l, real_path):
|
||||
db = DBPools()
|
||||
text = ''
|
||||
chunks_n = 0
|
||||
face_count = 0
|
||||
voice_speakers = 0
|
||||
meta_parts = {}
|
||||
try:
|
||||
with open(real_path, 'rb') as f:
|
||||
file_data = f.read()
|
||||
|
||||
text = ''
|
||||
chunks_n = 0
|
||||
face_count = 0
|
||||
voice_speakers = 0
|
||||
meta_parts = {}
|
||||
text_exts = {'.txt', '.md', '.csv', '.json', '.xml', '.html', '.htm', '.py', '.js', '.css', '.yaml', '.yml', '.log', '.rst'}
|
||||
image_exts = {'.jpg', '.jpeg', '.png', '.bmp', '.gif', '.webp'}
|
||||
audio_exts = {'.mp3', '.wav', '.flac', '.ogg', '.m4a', '.aac'}
|
||||
video_exts = {'.mp4', '.avi', '.mov', '.mkv', '.webm'}
|
||||
|
||||
# --- OFFICE DOCS: text extraction (PDF/DOCX/PPTX/XLSX) ---
|
||||
if ext_l == '.pdf' and not text:
|
||||
# --- OFFICE DOCS: text extraction (PDF/DOCX/PPTX/XLSX) ---
|
||||
if ext_l == '.pdf' and not text:
|
||||
import io; from PyPDF2 import PdfReader
|
||||
reader = PdfReader(io.BytesIO(file_data))
|
||||
text = '\n'.join(p.extract_text() or '' for p in reader.pages)
|
||||
elif ext_l == '.docx' and not text:
|
||||
elif ext_l == '.docx' and not text:
|
||||
import io; from docx import Document
|
||||
doc = Document(io.BytesIO(file_data))
|
||||
text = '\n'.join(p.text for p in doc.paragraphs)
|
||||
elif ext_l == '.pptx' and not text:
|
||||
elif ext_l == '.pptx' and not text:
|
||||
import io; from pptx import Presentation
|
||||
prs = Presentation(io.BytesIO(file_data))
|
||||
parts = []
|
||||
@ -51,7 +62,7 @@ elif ext_l == '.pptx' and not text:
|
||||
if hasattr(shape, 'text') and shape.text:
|
||||
parts.append(shape.text)
|
||||
text = '\n'.join(parts)
|
||||
elif ext_l == '.xlsx' and not text:
|
||||
elif ext_l == '.xlsx' and not text:
|
||||
import io; from openpyxl import load_workbook
|
||||
wb = load_workbook(io.BytesIO(file_data), data_only=True)
|
||||
parts = []
|
||||
@ -60,12 +71,12 @@ elif ext_l == '.xlsx' and not text:
|
||||
parts.append('\t'.join(str(c or '') for c in row))
|
||||
text = '\n'.join(parts)
|
||||
|
||||
# --- TEXT EXTRACTION ---
|
||||
if ext_l in text_exts:
|
||||
# --- TEXT EXTRACTION ---
|
||||
if ext_l in text_exts:
|
||||
text = file_data.decode('utf-8', errors='replace')
|
||||
|
||||
# --- IMAGE: face detection ---
|
||||
if ext_l in image_exts:
|
||||
# --- IMAGE: face detection ---
|
||||
if ext_l in image_exts:
|
||||
img_b64 = base64.b64encode(file_data).decode()
|
||||
try:
|
||||
client = StreamHttpClient()
|
||||
@ -77,8 +88,8 @@ if ext_l in image_exts:
|
||||
meta_parts['face'] = face_count
|
||||
except: pass
|
||||
|
||||
# --- AUDIO: voiceprint ---
|
||||
if ext_l in audio_exts:
|
||||
# --- AUDIO: voiceprint ---
|
||||
if ext_l in audio_exts:
|
||||
try:
|
||||
client = StreamHttpClient()
|
||||
resp = await client.request('POST', 'https://media.opencomputing.net/voiceprint/extract/submit',
|
||||
@ -88,8 +99,8 @@ if ext_l in audio_exts:
|
||||
meta_parts['voiceprint'] = voice_speakers
|
||||
except: pass
|
||||
|
||||
# --- VIDEO: frame extraction ---
|
||||
if ext_l in video_exts:
|
||||
# --- VIDEO: frame extraction ---
|
||||
if ext_l in video_exts:
|
||||
meta_parts['video'] = 'pending'
|
||||
try:
|
||||
tmp_img = '/tmp/' + doc_id + '_frame.jpg'
|
||||
@ -123,7 +134,7 @@ if ext_l in video_exts:
|
||||
{"id": doc_id + "_frame0", "vector": img_embeddings[0], "text": file_name}
|
||||
]}
|
||||
await client3.request('POST', 'https://vectordb.opencomputing.net/v1/upsert', json=vdb_data)
|
||||
async with get_sor_context(env, 'rag') as sor:
|
||||
async with db.sqlorContext('rag') as sor:
|
||||
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}$, 0, ${content}$, ${vid}$, NOW())",
|
||||
@ -134,8 +145,8 @@ if ext_l in video_exts:
|
||||
os.remove(tmp_img)
|
||||
except: pass
|
||||
|
||||
# --- RAG INGEST for text ---
|
||||
if text and len(text.strip()) > 10:
|
||||
# --- RAG INGEST for text ---
|
||||
if text and len(text.strip()) > 10:
|
||||
paragraphs = text.split('\n')
|
||||
chunks = []
|
||||
cur = ''
|
||||
@ -173,7 +184,7 @@ if text and len(text.strip()) > 10:
|
||||
except:
|
||||
pass
|
||||
|
||||
async with get_sor_context(env, 'rag') as sor:
|
||||
async with db.sqlorContext('rag') as sor:
|
||||
for i, chunk_text in enumerate(chunks):
|
||||
vid = vector_ids[i] if i < len(vector_ids) else ''
|
||||
await sor.sqlExe(
|
||||
@ -182,27 +193,45 @@ if text and len(text.strip()) > 10:
|
||||
{"id": doc_id + "_c" + str(i), "doc_id": doc_id, "kb_id": kb_id,
|
||||
"idx": i, "content": chunk_text[:2000], "vid": vid})
|
||||
chunks_n = len(chunks)
|
||||
except Exception as e:
|
||||
try:
|
||||
async with db.sqlorContext('rag') as sor:
|
||||
await sor.sqlExe(
|
||||
"UPDATE documents SET status='failed', metadata=${meta}$, updated_at=NOW() WHERE id=${id}$",
|
||||
{"id": doc_id, "meta": json.dumps({"error": str(e)[:300]}, ensure_ascii=False)})
|
||||
except: pass
|
||||
return
|
||||
|
||||
# --- SAVE TO DB ---
|
||||
meta_json = json.dumps(meta_parts, ensure_ascii=False)
|
||||
status = 'done'
|
||||
async with get_sor_context(env, 'rag') as sor:
|
||||
# --- finalize: mark document done + update KB chunk counts ---
|
||||
meta_json = json.dumps(meta_parts, ensure_ascii=False)
|
||||
try:
|
||||
async with db.sqlorContext('rag') as sor:
|
||||
await sor.sqlExe(
|
||||
"INSERT INTO documents (id, kb_id, folder_id, file_name, file_type, file_size, file_path, mime_type, status, chunk_count, metadata, org_id, created_at, updated_at) "
|
||||
"VALUES (${id}$, ${kb_id}$, ${folder_id}$, ${file_name}$, 'other', ${file_size}$, ${file_path}$, 'application/octet-stream', ${status}$, ${chunks}$, ${meta}$, ${org_id}$, NOW(), NOW())",
|
||||
{"id": doc_id, "kb_id": kb_id, "folder_id": folder_id, "file_name": file_name,
|
||||
"file_size": file_size, "file_path": web_path,
|
||||
"status": status, "chunks": chunks_n, "meta": meta_json, "org_id": userorgid})
|
||||
await sor.sqlExe(
|
||||
"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})
|
||||
"UPDATE documents SET status='done', chunk_count=${chunks}$, metadata=${meta}$, updated_at=NOW() WHERE id=${id}$",
|
||||
{"id": doc_id, "chunks": chunks_n, "meta": meta_json})
|
||||
if chunks_n:
|
||||
await sor.sqlExe(
|
||||
"UPDATE knowledge_bases SET chunk_count=chunk_count+${n}$ WHERE id=${kb_id}$",
|
||||
{"n": chunks_n, "kb_id": kb_id})
|
||||
except: pass
|
||||
|
||||
# ============================================================
|
||||
# SYNC PART — record document as 'pending', update KB counts,
|
||||
# fire background ingestion, return immediately.
|
||||
# ============================================================
|
||||
async with get_sor_context(env, 'rag') as sor:
|
||||
await sor.sqlExe(
|
||||
"INSERT INTO documents (id, kb_id, folder_id, file_name, file_type, file_size, file_path, mime_type, status, chunk_count, metadata, org_id, created_at, updated_at) "
|
||||
"VALUES (${id}$, ${kb_id}$, ${folder_id}$, ${file_name}$, 'other', ${file_size}$, ${file_path}$, 'application/octet-stream', 'pending', 0, '{}', ${org_id}$, NOW(), NOW())",
|
||||
{"id": doc_id, "kb_id": kb_id, "folder_id": folder_id, "file_name": file_name,
|
||||
"file_size": file_size, "file_path": web_path, "org_id": userorgid})
|
||||
await sor.sqlExe(
|
||||
"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})
|
||||
|
||||
# Fire background ingestion — pass primitives only (no env/request/proxy objects)
|
||||
background_reco(ingest_doc, doc_id, kb_id, file_name, ext_l, real_path)
|
||||
|
||||
result = {"status": "SUCCEEDED", "doc_id": doc_id, "file_name": file_name,
|
||||
"file_size": file_size, "folder_id": folder_id,
|
||||
"text_len": len(text), "chunks": chunks_n,
|
||||
"faces": face_count, "speakers": voice_speakers}
|
||||
"file_size": file_size, "folder_id": folder_id, "ingest": "pending"}
|
||||
return json.dumps(result, ensure_ascii=False, default=str)
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user