From b98a2b95074ad02fe70df4f96d33eebbb1a0b3c8 Mon Sep 17 00:00:00 2001 From: stefanfeng Date: Fri, 4 Sep 2026 11:53:56 +0800 Subject: [PATCH] fix(avatar): index knowledge documents asynchronously --- digital-avatar-app/backend/database.py | 1 + digital-avatar-app/backend/main.py | 2 + digital-avatar-app/backend/models.py | 2 + .../backend/routers/knowledge.py | 87 +++++------- .../backend/services/knowledge_vectorizer.py | 125 ++++++++++++++++++ .../backend/tests/test_knowledge_storage.py | 106 +++++++++++++-- digital-avatar-app/src/api/index.ts | 7 +- .../src/views/KnowledgeManage.vue | 51 ++++++- 8 files changed, 312 insertions(+), 69 deletions(-) create mode 100644 digital-avatar-app/backend/services/knowledge_vectorizer.py diff --git a/digital-avatar-app/backend/database.py b/digital-avatar-app/backend/database.py index 9937efa..db5ebf6 100644 --- a/digital-avatar-app/backend/database.py +++ b/digital-avatar-app/backend/database.py @@ -53,6 +53,7 @@ def init_db(): ("knowledge_docs", "embedding_model", "VARCHAR DEFAULT ''"), ("knowledge_docs", "chunk_count", "INTEGER DEFAULT 0"), ("knowledge_docs", "vectorized_at", "TIMESTAMP"), + ("knowledge_docs", "error_message", "VARCHAR DEFAULT ''"), ("avatars", "owner_id", "VARCHAR DEFAULT ''"), ("authorizations", "takeover_enabled", "BOOLEAN DEFAULT 0"), ("authorizations", "takeover_mode", "VARCHAR DEFAULT 'immediate'"), diff --git a/digital-avatar-app/backend/main.py b/digital-avatar-app/backend/main.py index 02e5f14..7bd8418 100644 --- a/digital-avatar-app/backend/main.py +++ b/digital-avatar-app/backend/main.py @@ -20,6 +20,7 @@ import routers.chat import routers.takeover from responses import ok from services.chat_attachment_service import purge_expired_chat_attachments +from services.knowledge_vectorizer import knowledge_vectorizer from services.token_billing import DEFAULT_TOKEN_GRANT, release_stale_reservations logger = logging.getLogger(__name__) @@ -131,6 +132,7 @@ def on_startup(): init_db() seed() + knowledge_vectorizer.start() # Release stale resources when startup is invoked again by a reload/test. stop_takeover_scheduler() diff --git a/digital-avatar-app/backend/models.py b/digital-avatar-app/backend/models.py index 0233eee..71a5969 100644 --- a/digital-avatar-app/backend/models.py +++ b/digital-avatar-app/backend/models.py @@ -190,6 +190,7 @@ class KnowledgeDoc(Base): file_size = Column(Integer, default=0) file_url = Column(String, default="") status = Column(String, default="uploaded") # uploaded | parsing | ready | failed + error_message = Column(String, default="") # 建立索引失败原因 vectorized = Column(Boolean, default=False) # 是否已向量化 embedding_model = Column(String, default="") # 向量模型标识 chunk_count = Column(Integer, default=0) # 切片数量 @@ -205,6 +206,7 @@ class KnowledgeDoc(Base): "fileSize": self.file_size, "fileUrl": self.file_url, "status": self.status, + "errorMessage": self.error_message or "", "vectorized": bool(self.vectorized), "embeddingModel": self.embedding_model, "chunkCount": self.chunk_count, diff --git a/digital-avatar-app/backend/routers/knowledge.py b/digital-avatar-app/backend/routers/knowledge.py index d336a2a..2b2a4fa 100644 --- a/digital-avatar-app/backend/routers/knowledge.py +++ b/digital-avatar-app/backend/routers/knowledge.py @@ -1,8 +1,5 @@ import os -import json -import logging import uuid -from datetime import datetime, timezone from fastapi import APIRouter, UploadFile, File, Depends, Header, HTTPException from pydantic import BaseModel @@ -12,10 +9,9 @@ from database import get_db from models import KnowledgeDoc, QAPair, KnowledgeChunk, Avatar, User from responses import ok, fail import embeddings +from services.knowledge_vectorizer import knowledge_vectorizer router = APIRouter() -logger = logging.getLogger(__name__) - BASE_DIR = os.path.dirname(os.path.abspath(__file__)) UPLOAD_DIR = os.path.abspath(os.getenv("UPLOAD_DIR", os.path.join(BASE_DIR, "uploads"))) os.makedirs(UPLOAD_DIR, exist_ok=True) @@ -71,15 +67,6 @@ def list_docs(avatar_id: str, authorization: str = Header(None), db: Session = D .order_by(KnowledgeDoc.created_at.desc()) .all() ) - # Older synchronous uploads could be interrupted after persisting "parsing". - # New uploads are committed only after indexing finishes, so these rows are stale. - stale_docs = [doc for doc in docs if doc.status == "parsing"] - if stale_docs: - for doc in stale_docs: - doc.status = "failed" - doc.vectorized = False - doc.chunk_count = 0 - db.commit() return ok([_doc_payload(d) for d in docs]) @@ -108,50 +95,42 @@ async def upload_doc(avatar_id: str, file: UploadFile = File(...), authorization status="parsing", ) - # Complete extraction and embedding before the first database commit so a - # process restart cannot leave a permanent "parsing" row behind. - try: - text = embeddings.extract_text(path, ext) - chunks = embeddings.chunk_text(text) - if not chunks: - raise ValueError("文档没有可建立索引的文字内容") - vectors = embeddings.embed(chunks) - if len(vectors) != len(chunks): - raise ValueError("向量服务返回数量与文档分段不一致") - doc.vectorized = True - doc.embedding_model = embeddings.MODEL - doc.chunk_count = len(chunks) - doc.vectorized_at = datetime.now(timezone.utc) - doc.status = "ready" - db.add(doc) - for i, (chunk, vector) in enumerate(zip(chunks, vectors)): - db.add( - KnowledgeChunk( - doc_id=doc.id, - avatar_id=avatar_id, - content=chunk, - vector=json.dumps(vector), - chunk_index=i, - embedding_model=embeddings.MODEL, - ) - ) - db.commit() - db.refresh(doc) - except Exception as exc: - db.rollback() - doc.status = "failed" - doc.vectorized = False - doc.embedding_model = "" - doc.chunk_count = 0 - doc.vectorized_at = None - db.add(doc) - db.commit() - db.refresh(doc) - logger.exception("knowledge vectorization failed for %s: %s", doc.id, exc) + # Persist and acknowledge the upload first. Extraction and embeddings may take + # minutes for a PDF and must never consume the browser request timeout. + db.add(doc) + db.commit() + db.refresh(doc) + knowledge_vectorizer.enqueue(doc.id) return ok(_doc_payload(doc)) +@router.post("/avatar/{avatar_id}/knowledge/docs/{doc_id}/retry") +def retry_doc(avatar_id: str, doc_id: str, authorization: str = Header(None), db: Session = Depends(get_db)): + _require_owned_avatar(db, avatar_id, authorization) + doc = db.query(KnowledgeDoc).filter( + KnowledgeDoc.id == doc_id, KnowledgeDoc.avatar_id == avatar_id + ).first() + if not doc: + return fail("文档不存在", code=404) + if doc.vectorized and doc.status == "ready": + return ok(_doc_payload(doc)) + stored_name = os.path.basename(doc.file_url or "") + if not stored_name or not os.path.isfile(os.path.join(UPLOAD_DIR, avatar_id, stored_name)): + return fail("原文件不可用,请重新上传", code=400) + db.query(KnowledgeChunk).filter(KnowledgeChunk.doc_id == doc.id).delete() + doc.status = "parsing" + doc.vectorized = False + doc.embedding_model = "" + doc.chunk_count = 0 + doc.vectorized_at = None + doc.error_message = "" + db.commit() + db.refresh(doc) + knowledge_vectorizer.enqueue(doc.id) + return ok(_doc_payload(doc)) + + @router.delete("/avatar/{avatar_id}/knowledge/docs/{doc_id}") def delete_doc(avatar_id: str, doc_id: str, authorization: str = Header(None), db: Session = Depends(get_db)): _require_owned_avatar(db, avatar_id, authorization) diff --git a/digital-avatar-app/backend/services/knowledge_vectorizer.py b/digital-avatar-app/backend/services/knowledge_vectorizer.py new file mode 100644 index 0000000..5e1624c --- /dev/null +++ b/digital-avatar-app/backend/services/knowledge_vectorizer.py @@ -0,0 +1,125 @@ +"""Durable, serial knowledge-document indexing for the avatar knowledge base.""" + +import json +import logging +import os +import queue +import threading +from datetime import datetime, timezone + +from database import SessionLocal +from models import KnowledgeChunk, KnowledgeDoc +import embeddings + +logger = logging.getLogger(__name__) + +BACKEND_DIR = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) +UPLOAD_DIR = os.path.abspath( + os.getenv("UPLOAD_DIR", os.path.join(BACKEND_DIR, "routers", "uploads")) +) + + +class KnowledgeVectorizer: + """Indexes one document at a time so slow providers cannot block uploads.""" + + def __init__(self): + self._queue: queue.Queue[str] = queue.Queue() + self._queued: set[str] = set() + self._lock = threading.Lock() + self._thread: threading.Thread | None = None + + def start(self): + if self._thread and self._thread.is_alive(): + return + self._thread = threading.Thread( + target=self._run, name="knowledge-vectorizer", daemon=True + ) + self._thread.start() + db = SessionLocal() + try: + # A process restart must not abandon documents already accepted by upload. + for (doc_id,) in db.query(KnowledgeDoc.id).filter(KnowledgeDoc.status == "parsing"): + self.enqueue(doc_id) + finally: + db.close() + + def enqueue(self, doc_id: str): + with self._lock: + if doc_id in self._queued: + return + self._queued.add(doc_id) + self._queue.put(doc_id) + + def _run(self): + while True: + doc_id = self._queue.get() + try: + self.vectorize_document(doc_id) + except Exception: + logger.exception("Unexpected knowledge vectorizer failure for %s", doc_id) + finally: + with self._lock: + self._queued.discard(doc_id) + self._queue.task_done() + + def vectorize_document(self, doc_id: str): + db = SessionLocal() + try: + doc = db.get(KnowledgeDoc, doc_id) + if not doc or doc.status != "parsing": + return + + stored_name = os.path.basename(doc.file_url or "") + path = os.path.join(UPLOAD_DIR, doc.avatar_id, stored_name) + if not stored_name or not os.path.isfile(path): + raise FileNotFoundError("原文件不可用,请重新上传") + + text = embeddings.extract_text(path, f".{doc.file_type}") + chunks = embeddings.chunk_text(text) + if not chunks: + raise ValueError("文档没有可建立索引的文字内容") + vectors = embeddings.embed(chunks) + if len(vectors) != len(chunks): + raise ValueError("向量服务返回数量与文档分段不一致") + + # Commit the document and every chunk together. Chat only sees complete indexes. + db.query(KnowledgeChunk).filter(KnowledgeChunk.doc_id == doc.id).delete() + db.add_all( + [ + KnowledgeChunk( + doc_id=doc.id, + avatar_id=doc.avatar_id, + content=chunk, + vector=json.dumps(vector), + chunk_index=index, + embedding_model=embeddings.MODEL, + ) + for index, (chunk, vector) in enumerate(zip(chunks, vectors)) + ] + ) + doc.vectorized = True + doc.embedding_model = embeddings.MODEL + doc.chunk_count = len(chunks) + doc.vectorized_at = datetime.now(timezone.utc) + doc.status = "ready" + doc.error_message = "" + db.commit() + logger.info("Knowledge document %s indexed with %s chunks", doc.id, len(chunks)) + except Exception as exc: + db.rollback() + failed_doc = db.get(KnowledgeDoc, doc_id) + if failed_doc: + db.query(KnowledgeChunk).filter(KnowledgeChunk.doc_id == failed_doc.id).delete() + failed_doc.status = "failed" + failed_doc.vectorized = False + failed_doc.embedding_model = "" + failed_doc.chunk_count = 0 + failed_doc.vectorized_at = None + failed_doc.error_message = str(exc)[:300] or "建立知识索引失败" + db.commit() + logger.exception("Knowledge vectorization failed for %s: %s", doc_id, exc) + finally: + db.close() + + +knowledge_vectorizer = KnowledgeVectorizer() diff --git a/digital-avatar-app/backend/tests/test_knowledge_storage.py b/digital-avatar-app/backend/tests/test_knowledge_storage.py index eede9e3..50c86a3 100644 --- a/digital-avatar-app/backend/tests/test_knowledge_storage.py +++ b/digital-avatar-app/backend/tests/test_knowledge_storage.py @@ -8,6 +8,7 @@ from database import SessionLocal from main import app from models import Avatar, KnowledgeChunk, KnowledgeDoc, QAPair from routers.knowledge import _doc_payload +from services.knowledge_vectorizer import knowledge_vectorizer client = TestClient(app) @@ -31,14 +32,14 @@ def test_doc_payload_reports_whether_the_persisted_file_exists(tmp_path: Path): assert _doc_payload(doc)["filePresent"] is True -def test_upload_marks_vectorization_failure_instead_of_staying_processing( +def test_upload_returns_before_background_vectorization( tmp_path: Path, authorization_context, ): context = authorization_context with ( patch("routers.knowledge.UPLOAD_DIR", str(tmp_path)), - patch("routers.knowledge.embeddings.embed", side_effect=RuntimeError("provider unavailable")), + patch("routers.knowledge.knowledge_vectorizer.enqueue") as enqueue, ): response = client.post( f"/api/avatar/{context['avatar'].id}/knowledge/docs", @@ -47,14 +48,15 @@ def test_upload_marks_vectorization_failure_instead_of_staying_processing( ) payload = response.json()["data"] - assert payload["status"] == "failed" + assert payload["status"] == "parsing" assert payload["vectorized"] is False assert payload["chunkCount"] == 0 + enqueue.assert_called_once_with(payload["id"]) db = SessionLocal() try: stored = db.query(KnowledgeDoc).filter(KnowledgeDoc.id == payload["id"]).one() - assert stored.status == "failed" + assert stored.status == "parsing" assert db.query(KnowledgeChunk).filter(KnowledgeChunk.doc_id == stored.id).count() == 0 db.delete(stored) db.commit() @@ -62,14 +64,14 @@ def test_upload_marks_vectorization_failure_instead_of_staying_processing( db.close() -def test_markdown_upload_commits_ready_document_and_chunks_together( +def test_background_vectorizer_commits_ready_document_and_chunks_together( tmp_path: Path, authorization_context, ): context = authorization_context with ( patch("routers.knowledge.UPLOAD_DIR", str(tmp_path)), - patch("routers.knowledge.embeddings.embed", return_value=[[1.0, 0.0]]), + patch("routers.knowledge.knowledge_vectorizer.enqueue"), ): response = client.post( f"/api/avatar/{context['avatar'].id}/knowledge/docs", @@ -78,14 +80,19 @@ def test_markdown_upload_commits_ready_document_and_chunks_together( ) payload = response.json()["data"] - assert payload["status"] == "ready" - assert payload["vectorized"] is True - assert payload["chunkCount"] == 1 + assert payload["status"] == "parsing" + with ( + patch("services.knowledge_vectorizer.UPLOAD_DIR", str(tmp_path)), + patch("services.knowledge_vectorizer.embeddings.embed", return_value=[[1.0, 0.0]]), + ): + knowledge_vectorizer.vectorize_document(payload["id"]) db = SessionLocal() try: stored = db.query(KnowledgeDoc).filter(KnowledgeDoc.id == payload["id"]).one() assert stored.status == "ready" + assert stored.vectorized is True + assert stored.chunk_count == 1 assert db.query(KnowledgeChunk).filter(KnowledgeChunk.doc_id == stored.id).count() == 1 db.query(KnowledgeChunk).filter(KnowledgeChunk.doc_id == stored.id).delete() db.delete(stored) @@ -94,6 +101,87 @@ def test_markdown_upload_commits_ready_document_and_chunks_together( db.close() +def test_background_vectorizer_keeps_failure_reason_for_retry( + tmp_path: Path, + authorization_context, +): + context = authorization_context + with ( + patch("routers.knowledge.UPLOAD_DIR", str(tmp_path)), + patch("routers.knowledge.knowledge_vectorizer.enqueue"), + ): + response = client.post( + f"/api/avatar/{context['avatar'].id}/knowledge/docs", + headers=context["owner_headers"], + files={"file": ("knowledge.md", b"# Knowledge\n\nTest content", "text/markdown")}, + ) + + payload = response.json()["data"] + with ( + patch("services.knowledge_vectorizer.UPLOAD_DIR", str(tmp_path)), + patch("services.knowledge_vectorizer.embeddings.embed", side_effect=RuntimeError("provider unavailable")), + ): + knowledge_vectorizer.vectorize_document(payload["id"]) + + db = SessionLocal() + try: + stored = db.query(KnowledgeDoc).filter(KnowledgeDoc.id == payload["id"]).one() + assert stored.status == "failed" + assert stored.error_message == "provider unavailable" + assert db.query(KnowledgeChunk).filter(KnowledgeChunk.doc_id == stored.id).count() == 0 + db.delete(stored) + db.commit() + finally: + db.close() + + +def test_retry_queues_a_failed_document_again( + tmp_path: Path, + authorization_context, +): + context = authorization_context + document_id = f"retry-doc-{context['suffix']}" + avatar_dir = tmp_path / context["avatar"].id + avatar_dir.mkdir() + (avatar_dir / "retry.md").write_text("retry content", encoding="utf-8") + db = SessionLocal() + try: + db.add( + KnowledgeDoc( + id=document_id, + avatar_id=context["avatar"].id, + filename="retry.md", + file_type="md", + file_url=f"/api/files/{context['avatar'].id}/retry.md", + status="failed", + error_message="provider unavailable", + ) + ) + db.commit() + finally: + db.close() + + with ( + patch("routers.knowledge.UPLOAD_DIR", str(tmp_path)), + patch("routers.knowledge.knowledge_vectorizer.enqueue") as enqueue, + ): + response = client.post( + f"/api/avatar/{context['avatar'].id}/knowledge/docs/{document_id}/retry", + headers=context["owner_headers"], + ) + + payload = response.json()["data"] + assert payload["status"] == "parsing" + assert payload["errorMessage"] == "" + enqueue.assert_called_once_with(document_id) + db = SessionLocal() + try: + db.query(KnowledgeDoc).filter(KnowledgeDoc.id == document_id).delete() + db.commit() + finally: + db.close() + + def test_each_avatar_has_an_independent_document_and_qa_scope(authorization_context): context = authorization_context first_avatar_id = context["avatar"].id diff --git a/digital-avatar-app/src/api/index.ts b/digital-avatar-app/src/api/index.ts index be7a4a8..9a0fba6 100644 --- a/digital-avatar-app/src/api/index.ts +++ b/digital-avatar-app/src/api/index.ts @@ -305,6 +305,7 @@ export interface KnowledgeDoc { vectorized?: boolean embeddingModel?: string chunkCount?: number + errorMessage?: string createdAt: string } @@ -335,7 +336,8 @@ export const uploadKnowledgeDoc = (avatarId: string, file: File) => { const form = new FormData() form.append('file', file) return request.post(`/avatar/${avatarId}/knowledge/docs`, form, { - headers: { 'Content-Type': 'multipart/form-data' } + headers: { 'Content-Type': 'multipart/form-data' }, + timeout: 120000 }) } @@ -343,6 +345,9 @@ export const uploadKnowledgeDoc = (avatarId: string, file: File) => { export const deleteKnowledgeDoc = (avatarId: string, docId: string) => request.delete(`/avatar/${avatarId}/knowledge/docs/${docId}`) +export const retryKnowledgeDoc = (avatarId: string, docId: string) => + request.post(`/avatar/${avatarId}/knowledge/docs/${docId}/retry`) + // 标准问答对列表 export const getQAPairs = (avatarId: string) => request.get(`/avatar/${avatarId}/knowledge/qa`) diff --git a/digital-avatar-app/src/views/KnowledgeManage.vue b/digital-avatar-app/src/views/KnowledgeManage.vue index e20e3e2..82dc2d0 100644 --- a/digital-avatar-app/src/views/KnowledgeManage.vue +++ b/digital-avatar-app/src/views/KnowledgeManage.vue @@ -27,7 +27,7 @@

支持 MD / TXT / PDF / DOC / DOCX / XLSX,上传后自动向量化

-

上传并向量化中…

+

文件上传中…

{{ uploadError }}

@@ -42,7 +42,10 @@

{{ doc.fileType.toUpperCase() }} · {{ formatSize(doc.fileSize) }} · {{ formatDate(doc.createdAt) }}

{{ documentState(doc).detail }}

- +
+ + +
📂 暂无文档,先上传一个知识文件
@@ -78,7 +81,7 @@