diff --git a/Dockerfile b/Dockerfile index 119d28a..dd7d141 100644 --- a/Dockerfile +++ b/Dockerfile @@ -6,6 +6,9 @@ RUN apt-get update && apt-get install -y --no-install-recommends \ build-essential \ libsndfile1 \ curl \ + tesseract-ocr \ + tesseract-ocr-rus \ + tesseract-ocr-eng \ && rm -rf /var/lib/apt/lists/* # Рабочая директория diff --git a/Dockerfile.rag b/Dockerfile.rag index 9a5f4af..94e4dcf 100644 --- a/Dockerfile.rag +++ b/Dockerfile.rag @@ -8,4 +8,8 @@ RUN pip install --no-cache-dir --timeout 300 \ python-dotenv>=1.0.0 \ sentence-transformers>=3.0.0 \ bcrypt>=4.0.0 \ - "python-jose[cryptography]" + "python-jose[cryptography]" \ + pymupdf>=1.24.0 \ + openpyxl>=3.1.0 \ + Pillow>=10.0.0 \ + pytesseract>=0.3.10 diff --git a/backend/auth/models.py b/backend/auth/models.py index 154325b..333ecb1 100644 --- a/backend/auth/models.py +++ b/backend/auth/models.py @@ -1,7 +1,7 @@ """Auth data models.""" from dataclasses import dataclass, field -from typing import List, Optional +from typing import Any, Dict, List, Optional ROLES = ("admin", "director", "user") @@ -51,3 +51,11 @@ class UserContext: def can_global_search(self) -> bool: return self.has_all_projects_access + + def can_see_task(self, task: Dict[str, Any]) -> bool: + """Whether realtime/API task updates for this task may be shown to the user.""" + if task.get("org_slug") != self.org_slug: + return False + if self.has_all_projects_access: + return True + return task.get("user_id") == self.user_id diff --git a/backend/auth/routes.py b/backend/auth/routes.py index aeca3f7..3719ec8 100644 --- a/backend/auth/routes.py +++ b/backend/auth/routes.py @@ -15,6 +15,7 @@ from backend.auth.service import ( delete_personal_project, get_user_context, list_accessible_projects, + normalize_project_slug, update_user_projects, user_to_dict, ) @@ -134,11 +135,13 @@ async def admin_list_projects(admin: UserContext = Depends(require_admin)): @admin_router.post("/projects") async def admin_create_project(payload: CreateProjectRequest, admin: UserContext = Depends(require_admin)): - slug = payload.slug.strip().lower() - if not slug: - raise HTTPException(status_code=400, detail="slug обязателен") try: - project = db.create_project(admin.org_id, slug, payload.name, owner_user_id=None, config=load_config()) + slug = normalize_project_slug(payload.slug) + except ValueError as e: + raise HTTPException(status_code=400, detail=str(e)) from e + display_name = payload.name.strip() or slug + try: + project = db.create_project(admin.org_id, slug, display_name, owner_user_id=None, config=load_config()) return {"project": project} except Exception as e: if "UNIQUE" in str(e): diff --git a/backend/ingest_worker.py b/backend/ingest_worker.py new file mode 100644 index 0000000..893bea8 --- /dev/null +++ b/backend/ingest_worker.py @@ -0,0 +1,133 @@ +"""Document ingestion worker pipeline.""" + +import asyncio +import json +import shutil +from datetime import datetime +from pathlib import Path +from typing import Any, Dict + +from backend.paths import org_documents_dir, org_rag_index_dir, write_folder_project_meta +from src.config import load_config, resolve_opencode_credentials +from src.ingest.classify import classify_document +from src.ingest.formatter import format_global_index_document, format_index_document +from src.ingest.router import extract_document +from src.rag.indexer import index_meeting + + +async def process_document_ingest(job: Dict[str, Any], tasks: dict, send_progress): + task_id = job["task_id"] + file_path = Path(job["file_path"]) + org_slug = job["org_slug"] + project_slug = job["project_slug"] + doc_type = job.get("doc_type", "other") + display_name = job.get("display_name", file_path.name) + + tasks[task_id].update({"status": "processing", "message": "Извлечение текста...", "progress": 10}) + await send_progress(task_id, 10, "Извлечение текста...", "processing") + + try: + config = load_config() + ingest_cfg = config.get("ingest", {}) + pdf_ocr = ingest_cfg.get("pdf_ocr", True) + + doc = await asyncio.to_thread( + extract_document, + file_path, + project_slug, + doc_type, + None, + pdf_ocr, + ) + + if not doc.full_text.strip(): + raise ValueError("Не удалось извлечь текст из документа") + + documents_dir = org_documents_dir(org_slug) + output_dir = documents_dir / doc.document_id + await asyncio.to_thread(output_dir.mkdir, parents=True, exist_ok=True) + + original_dest = output_dir / file_path.name + await asyncio.to_thread(shutil.copy2, file_path, original_dest) + await asyncio.to_thread( + (output_dir / "extracted.md").write_text, + doc.full_text, + encoding="utf-8", + ) + await asyncio.to_thread(write_folder_project_meta, output_dir, project_slug) + + tasks[task_id].update({"status": "postprocessing", "message": "Анализ документа...", "progress": 40}) + await send_progress(task_id, 40, "Анализ документа...", "postprocessing") + + metadata = doc.to_metadata_dict() + rag_cfg = config.get("rag", {}) + api_key, base_url = resolve_opencode_credentials(config) + + if api_key and ingest_cfg.get("auto_classify", True): + metadata = await classify_document( + text=doc.full_text, + project=project_slug, + doc_type_hint=doc_type, + api_key=api_key, + base_url=base_url, + model=rag_cfg.get("index_model", "mimo-v2.5-free"), + chunk_size=int(rag_cfg.get("classify_chunk_size", 7000)), + ) + metadata["filename"] = doc.filename + metadata["document_id"] = doc.document_id + + await asyncio.to_thread( + (output_dir / "metadata.json").write_text, + json.dumps(metadata, ensure_ascii=False, indent=2), + encoding="utf-8", + ) + + doc_text = format_index_document(doc, metadata) + index_path = output_dir / "index.txt" + await asyncio.to_thread(index_path.write_text, doc_text, encoding="utf-8") + + result_data = { + "document_id": doc.document_id, + "dir": str(output_dir), + "rel_dir": str(output_dir.relative_to(documents_dir)), + "extracted": str(output_dir / "extracted.md"), + "index": str(index_path), + "project": project_slug, + "doc_type": metadata.get("doc_type", doc_type), + "kind": "document", + } + + if rag_cfg.get("enabled", False) and rag_cfg.get("auto_index", True): + tasks[task_id].update({"message": "Индексация в RAG...", "progress": 75}) + await send_progress(task_id, 75, "Индексация в RAG...", "postprocessing") + global_doc_text = format_global_index_document(doc_text, metadata) + await index_meeting( + doc_text=doc_text, + global_doc_text=global_doc_text, + project_name=project_slug, + working_dir_base=org_rag_index_dir(org_slug), + model=rag_cfg.get("index_model", "mimo-v2.5-free"), + api_key=api_key, + base_url=base_url, + ) + + from backend.queue import _cleanup_upload + await asyncio.to_thread(_cleanup_upload, file_path) + + tasks[task_id].update({ + "status": "completed", + "progress": 100, + "message": "Документ проиндексирован", + "result": result_data, + "finished": datetime.now().isoformat(), + }) + await send_progress(task_id, 100, "Документ проиндексирован", "completed", result=result_data) + + except Exception as e: + error_msg = str(e) + tasks[task_id].update({ + "status": "error", + "message": f"Ошибка: {error_msg}", + "error": error_msg, + }) + await send_progress(task_id, 0, f"Ошибка: {error_msg}", "error", error=error_msg) diff --git a/backend/main.py b/backend/main.py index 477d67e..bae11a4 100644 --- a/backend/main.py +++ b/backend/main.py @@ -16,7 +16,7 @@ from backend.auth.models import UserContext from backend.auth.routes import admin_router, router as auth_router from backend.auth import database as auth_db from backend.auth.service import ensure_project_access, list_accessible_projects -from backend.paths import org_meetings_dir, org_rag_index_dir, resolve_meeting_path +from backend.paths import org_documents_dir, org_meetings_dir, org_rag_index_dir, resolve_document_path, resolve_meeting_path from backend.queue import ( delete_folder, get_all_tasks, @@ -29,6 +29,7 @@ from backend.queue import ( set_progress_callback, start_workers, stop_workers, + tasks, ) sys.path.insert(0, str(Path(__file__).parent.parent)) @@ -55,7 +56,14 @@ class ConnectionManager: self.active_connections = [(ws, user if ws is websocket else u) for ws, u in self.active_connections] async def broadcast(self, message: dict): - for conn, _user in self.active_connections: + """Send only to connections allowed to see the task (org + ACL).""" + task_id = message.get("task_id") + task = tasks.get(task_id) if task_id else None + for conn, user in self.active_connections: + if user is None: + continue + if task is not None and not user.can_see_task(task): + continue try: await conn.send_json(message) except Exception: @@ -77,8 +85,13 @@ async def lifespan(app: FastAPI): queue_cfg = config.get("queue", {}) transcribe_workers = int(queue_cfg.get("transcribe_workers", 2)) postprocess_workers = int(queue_cfg.get("postprocess_workers", 1)) + ingest_workers = int(queue_cfg.get("ingest_workers", 1)) print("🚀 Запуск рабочих процессов...") - start_workers(transcribe_workers=transcribe_workers, postprocess_workers=postprocess_workers) + start_workers( + transcribe_workers=transcribe_workers, + postprocess_workers=postprocess_workers, + ingest_workers=ingest_workers, + ) yield print("🛑 Остановка рабочих процессов...") stop_workers() @@ -107,12 +120,24 @@ async def _list_rag_project_slugs(user: UserContext) -> List[str]: return user.filter_projects(projects) -async def _rag_chat_for_user(user: UserContext, question: str, history: list, project_name: Optional[str], mode: str): +async def _rag_chat_for_user( + user: UserContext, + question: str, + history: list, + project_name: Optional[str], + chat_mode: str = "hybrid", + retrieval_mode: str = "hybrid", +): if project_name: ensure_project_access(user, project_name) elif not user.can_global_search(): raise HTTPException(status_code=403, detail="Глобальный поиск доступен только администратору") + if chat_mode not in ("hybrid", "compare", "timeline"): + chat_mode = "hybrid" + if retrieval_mode not in ("naive", "local", "global", "hybrid"): + retrieval_mode = "hybrid" + config = load_config() rag_cfg = config.get("rag", {}) api_key, base_url = resolve_opencode_credentials(config) @@ -124,7 +149,8 @@ async def _rag_chat_for_user(user: UserContext, question: str, history: list, pr project_name=project_name, base_url=base_url, chat_model=rag_cfg.get("chat_model", "deepseek-v4-flash-free"), - mode=mode, + mode=retrieval_mode, + chat_mode=chat_mode, index_model=rag_cfg.get("index_model", "mimo-v2.5-free"), ) @@ -149,13 +175,16 @@ async def login_page(): async def upload_file( file: UploadFile = File(...), project: str = Form(...), + doc_type: str = Form("other"), user: UserContext = Depends(get_current_user), ): content = await file.read() try: - task_id, _ = await save_upload(content, file.filename or "upload.bin", user, project) + task_id, _ = await save_upload(content, file.filename or "upload.bin", user, project, doc_type) except PermissionError as e: raise HTTPException(status_code=403, detail=str(e)) from e + except ValueError as e: + raise HTTPException(status_code=400, detail=str(e)) from e return { "task_id": task_id, @@ -167,10 +196,38 @@ async def upload_file( } +@app.post("/upload-document") +async def upload_document( + file: UploadFile = File(...), + project: str = Form(...), + doc_type: str = Form("other"), + user: UserContext = Depends(get_current_user), +): + """Явная загрузка документа (MD, PDF, DOCX, XLSX, TXT).""" + content = await file.read() + try: + task_id, _ = await save_upload(content, file.filename or "document.bin", user, project, doc_type) + except PermissionError as e: + raise HTTPException(status_code=403, detail=str(e)) from e + except ValueError as e: + raise HTTPException(status_code=400, detail=str(e)) from e + + return { + "task_id": task_id, + "file": file.filename, + "project": project, + "doc_type": doc_type, + "status": "queued", + "message": "Документ добавлен в очередь ingest", + "queue": get_queue_info(user), + } + + @app.post("/upload-batch") async def upload_batch( files: List[UploadFile] = File(...), project: str = Form(...), + doc_type: str = Form("other"), user: UserContext = Depends(get_current_user), ): if not files: @@ -180,9 +237,11 @@ async def upload_batch( for file in files: content = await file.read() try: - task_id, _ = await save_upload(content, file.filename or "upload.bin", user, project) + task_id, _ = await save_upload(content, file.filename or "upload.bin", user, project, doc_type) except PermissionError as e: raise HTTPException(status_code=403, detail=str(e)) from e + except ValueError as e: + raise HTTPException(status_code=400, detail=f"{file.filename}: {e}") from e results.append({"task_id": task_id, "file": file.filename, "project": project, "status": "queued"}) return { @@ -298,7 +357,8 @@ async def api_rag_query(payload: dict, user: UserContext = Depends(get_current_u payload.get("question", ""), payload.get("history", []), payload.get("project"), - payload.get("mode", "hybrid"), + chat_mode=payload.get("chat_mode", payload.get("mode", "hybrid")), + retrieval_mode=payload.get("retrieval_mode", "hybrid"), ) return {"answer": result["answer"], "context": result["context"], "project": result["project"]} except HTTPException: @@ -317,7 +377,8 @@ async def api_rag_query_global(payload: dict, user: UserContext = Depends(get_cu payload.get("question", ""), payload.get("history", []), None, - payload.get("mode", "hybrid"), + chat_mode=payload.get("chat_mode", payload.get("mode", "hybrid")), + retrieval_mode=payload.get("retrieval_mode", "hybrid"), ) return {"answer": result["answer"], "context": result["context"], "project": None} except HTTPException: @@ -329,19 +390,29 @@ async def api_rag_query_global(payload: dict, user: UserContext = Depends(get_cu @app.post("/api/rag/index/{folder_name:path}") async def api_rag_index_folder(folder_name: str, user: UserContext = Depends(get_current_user)): try: - folder_path = resolve_meeting_path(user.org_slug, folder_name) + from backend.queue import _folder_project_slug + + if folder_name.startswith("documents/"): + folder_path = resolve_document_path(user.org_slug, folder_name[len("documents/"):]) + base_dir = org_documents_dir(user.org_slug) + index_files = list(folder_path.glob("index.txt")) + txt_files = index_files + else: + rel = folder_name[len("meetings/"):] if folder_name.startswith("meetings/") else folder_name + folder_path = resolve_meeting_path(user.org_slug, rel) + base_dir = org_meetings_dir(user.org_slug) + txt_files = list(folder_path.glob("*.txt")) + if not folder_path.exists(): return {"error": "Folder not found"} - from backend.queue import _folder_project_slug - project = _folder_project_slug(folder_path.name, org_meetings_dir(user.org_slug)) + project = _folder_project_slug(folder_path.name, base_dir) if not project: return {"error": "Project metadata not found"} ensure_project_access(user, project) - txt_files = list(folder_path.glob("*.txt")) if not txt_files: - return {"error": "No .txt protocol found in folder"} + return {"error": "No index.txt or .txt protocol found in folder"} doc_text = txt_files[0].read_text(encoding="utf-8") config = load_config() @@ -376,7 +447,8 @@ async def _handle_rag_query_ws(websocket: WebSocket, msg: dict, user: UserContex msg.get("question", ""), msg.get("history", []), project, - msg.get("mode", "hybrid"), + chat_mode=msg.get("chat_mode", msg.get("mode", "hybrid")), + retrieval_mode=msg.get("retrieval_mode", "hybrid"), ) await websocket.send_json({ "type": "rag_response", diff --git a/backend/paths.py b/backend/paths.py index 3b3b04a..7fb4c40 100644 --- a/backend/paths.py +++ b/backend/paths.py @@ -9,6 +9,7 @@ UPLOAD_ROOT = Path("uploads") PROCESSED_ROOT = Path("processed") RAG_CACHE_DIRNAME = "lightrag_caches" MEETINGS_DIRNAME = "meetings" +DOCUMENTS_DIRNAME = "documents" def org_upload_dir(org_slug: str, user_id: int) -> Path: @@ -29,6 +30,20 @@ def org_rag_index_dir(org_slug: str) -> Path: return path +def org_documents_dir(org_slug: str) -> Path: + path = PROCESSED_ROOT / org_slug / DOCUMENTS_DIRNAME + path.mkdir(parents=True, exist_ok=True) + return path + + +def resolve_document_path(org_slug: str, rel_path: str) -> Path: + base = org_documents_dir(org_slug).resolve() + full = (base / rel_path).resolve() + if not str(full).startswith(str(base)): + raise ValueError("Invalid path") + return full + + def resolve_meeting_path(org_slug: str, rel_path: str) -> Path: """Resolve relative path under org meetings dir; reject traversal.""" base = org_meetings_dir(org_slug).resolve() diff --git a/backend/queue.py b/backend/queue.py index 9c970ca..255dc9d 100644 --- a/backend/queue.py +++ b/backend/queue.py @@ -13,7 +13,8 @@ sys.path.insert(0, str(Path(__file__).parent.parent)) from backend.auth.models import UserContext from backend.auth.service import ensure_project_access -from backend.paths import org_meetings_dir, org_rag_index_dir, org_upload_dir, resolve_meeting_path, write_folder_project_meta +from backend.paths import org_documents_dir, org_meetings_dir, org_rag_index_dir, org_upload_dir, resolve_document_path, resolve_meeting_path, write_folder_project_meta +from src.ingest.router import is_audio_file, is_document_file from src.audio_utils import prepare_audio_input from src.config import load_config, resolve_opencode_credentials from src.document import build_document @@ -31,6 +32,7 @@ tasks: Dict[str, Dict[str, Any]] = {} _progress_callback: Optional[Callable] = None _transcribe_queue: asyncio.Queue = asyncio.Queue() _postprocess_queue: asyncio.Queue = asyncio.Queue() +_ingest_queue: asyncio.Queue = asyncio.Queue() _workers: List[asyncio.Task] = [] @@ -40,11 +42,7 @@ def set_progress_callback(callback: Callable): def _task_visible(task: Dict[str, Any], user: UserContext) -> bool: - if task.get("org_slug") != user.org_slug: - return False - if user.has_all_projects_access: - return True - return task.get("user_id") == user.user_id + return user.can_see_task(task) def _filter_tasks_for_user(user: UserContext) -> List[Dict[str, Any]]: @@ -62,7 +60,8 @@ def _filter_queue_info(user: UserContext) -> Dict[str, Any]: "postprocessing": postprocessing, "pending_transcribe": _transcribe_queue.qsize(), "pending_postprocess": _postprocess_queue.qsize(), - "pending_in_queue": _transcribe_queue.qsize(), + "pending_ingest": _ingest_queue.qsize(), + "pending_in_queue": _transcribe_queue.qsize() + _ingest_queue.qsize(), } @@ -77,6 +76,8 @@ async def _send_progress(task_id: str, progress: int, message: str, status: str, "status": status, "file": task_info.get("file", ""), "project": task_info.get("project_slug", ""), + "org_slug": task_info.get("org_slug"), + "user_id": task_info.get("user_id"), "queue_position": task_info.get("queue_position"), "result": result, "error": error, @@ -345,14 +346,35 @@ async def _postprocess_worker_loop(worker_id: int): print(f"[Postprocess Worker {worker_id} Error] {e}") -def start_workers(transcribe_workers: int = 2, postprocess_workers: int = 1): +async def _ingest_worker_loop(worker_id: int): + from backend.ingest_worker import process_document_ingest + + print(f"[Ingest Worker {worker_id}] запущен") + while True: + try: + job = await _ingest_queue.get() + print(f"[Ingest Worker {worker_id}] задача {job.get('task_id')}") + await process_document_ingest(job, tasks, _send_progress) + _ingest_queue.task_done() + except asyncio.CancelledError: + break + except Exception as e: + print(f"[Ingest Worker {worker_id} Error] {e}") + + +def start_workers(transcribe_workers: int = 2, postprocess_workers: int = 1, ingest_workers: int = 1): global _workers _workers.clear() for i in range(transcribe_workers): _workers.append(asyncio.create_task(_transcribe_worker_loop(i + 1))) + for i in range(ingest_workers): + _workers.append(asyncio.create_task(_ingest_worker_loop(i + 1))) for i in range(postprocess_workers): _workers.append(asyncio.create_task(_postprocess_worker_loop(i + 1))) - print(f"[Queue] transcribe_workers={transcribe_workers}, postprocess_workers={postprocess_workers}") + print( + f"[Queue] transcribe={transcribe_workers}, ingest={ingest_workers}, " + f"postprocess={postprocess_workers}" + ) def stop_workers(): @@ -365,10 +387,17 @@ async def save_upload( filename: str, user: UserContext, project_slug: str, + doc_type: str = "other", ) -> tuple[str, Path]: ensure_project_access(user, project_slug) - slug = project_slug.strip().lower() safe_name = Path(filename).name or "upload.bin" + + if is_document_file(safe_name): + return await save_document_upload(content, safe_name, user, project_slug, doc_type) + if not is_audio_file(safe_name): + raise ValueError(f"Неподдерживаемый формат файла: {safe_name}") + + slug = project_slug.strip().lower() task_id = f"task_{datetime.now().strftime('%Y%m%d_%H%M%S')}_{uuid.uuid4().hex[:8]}" task_dir = org_upload_dir(user.org_slug, user.user_id) / task_id task_dir.mkdir(parents=True, exist_ok=True) @@ -383,6 +412,7 @@ async def save_upload( "progress": 0, "message": f"В очереди транскрибации (№{queue_position})", "file": safe_name, + "task_type": "transcribe", "project_slug": slug, "org_slug": user.org_slug, "user_id": user.user_id, @@ -398,6 +428,53 @@ async def save_upload( return task_id, file_path +async def save_document_upload( + content: bytes, + filename: str, + user: UserContext, + project_slug: str, + doc_type: str = "other", +) -> tuple[str, Path]: + slug = project_slug.strip().lower() + task_id = f"doc_{datetime.now().strftime('%Y%m%d_%H%M%S')}_{uuid.uuid4().hex[:8]}" + task_dir = org_upload_dir(user.org_slug, user.user_id) / task_id + task_dir.mkdir(parents=True, exist_ok=True) + file_path = task_dir / filename + + await asyncio.to_thread(file_path.write_bytes, content) + + queue_position = _ingest_queue.qsize() + 1 + tasks[task_id] = { + "task_id": task_id, + "status": "queued", + "progress": 0, + "message": f"В очереди ingest (№{queue_position})", + "file": filename, + "task_type": "ingest", + "doc_type": doc_type, + "project_slug": slug, + "org_slug": user.org_slug, + "user_id": user.user_id, + "username": user.username, + "queue_position": queue_position, + "result": None, + "error": None, + "started": datetime.now().isoformat(), + } + + job = { + "task_id": task_id, + "file_path": str(file_path), + "display_name": filename, + "org_slug": user.org_slug, + "project_slug": slug, + "doc_type": doc_type, + } + await _ingest_queue.put(job) + await _send_progress(task_id, 0, f"В очереди ingest (№{queue_position})", "queued") + return task_id, file_path + + def get_queue_info(user: Optional[UserContext] = None) -> Dict[str, Any]: if user: return _filter_queue_info(user) @@ -410,10 +487,126 @@ def get_queue_info(user: Optional[UserContext] = None) -> Dict[str, Any]: "postprocessing": postprocessing, "pending_transcribe": _transcribe_queue.qsize(), "pending_postprocess": _postprocess_queue.qsize(), - "pending_in_queue": _transcribe_queue.qsize(), + "pending_ingest": _ingest_queue.qsize(), + "pending_in_queue": _transcribe_queue.qsize() + _ingest_queue.qsize(), } +def _build_tree_section(base_dir: Path, rel_prefix: str, user: UserContext, kind_default: str) -> List[Dict[str, Any]]: + tree = [] + if not base_dir.exists(): + return tree + + skip_suffixes = ("_segments.json",) + skip_names = {".project.json"} + + for item in sorted(base_dir.iterdir()): + if not item.is_dir(): + continue + project_slug = _folder_project_slug(item.name, base_dir) + if project_slug and not user.can_access_project(project_slug): + continue + + files = [] + for f in sorted(item.iterdir(), key=lambda p: _file_sort_key({"name": p.name})): + if f.is_file() and not f.name.endswith(skip_suffixes) and f.name not in skip_names: + kind = kind_default + if "_summary" in f.name: + kind = "summary" + elif f.name in ("extracted.md", "index.txt", "metadata.json"): + kind = "document" + files.append({ + "name": f.name, + "path": f"{rel_prefix}/{f.relative_to(base_dir).as_posix()}", + "size": f.stat().st_size, + "ext": f.suffix.lower(), + "kind": kind, + }) + meta_path = item / "metadata.json" + doc_type = None + if meta_path.exists(): + try: + doc_type = json.loads(meta_path.read_text(encoding="utf-8")).get("doc_type") + except Exception: + pass + + tree.append({ + "name": item.name, + "path": f"{rel_prefix}/{item.relative_to(base_dir).as_posix()}", + "project": project_slug, + "doc_type": doc_type, + "section": rel_prefix, + "files": files, + "created": datetime.fromtimestamp(item.stat().st_ctime).isoformat(), + }) + return tree + + +def get_processed_tree(user: UserContext) -> List[Dict[str, Any]]: + meetings_dir = org_meetings_dir(user.org_slug) + documents_dir = org_documents_dir(user.org_slug) + tree = _build_tree_section(meetings_dir, "meetings", user, "protocol") + tree.extend(_build_tree_section(documents_dir, "documents", user, "document")) + tree.sort(key=lambda x: x.get("created", ""), reverse=True) + return tree + + +def _resolve_content_path(user: UserContext, rel_path: str) -> Path: + if rel_path.startswith("documents/"): + return resolve_document_path(user.org_slug, rel_path[len("documents/"):]) + if rel_path.startswith("meetings/"): + return resolve_meeting_path(user.org_slug, rel_path[len("meetings/"):]) + full = resolve_meeting_path(user.org_slug, rel_path) + if full.exists(): + return full + return resolve_document_path(user.org_slug, rel_path) + + +def read_file_content(user: UserContext, rel_path: str) -> str: + full_path = _resolve_content_path(user, rel_path) + if not full_path.exists() or not full_path.is_file(): + raise FileNotFoundError(f"Файл не найден: {rel_path}") + + base_dir = full_path.parent.parent + folder = full_path.parent.name + project_slug = _folder_project_slug(folder, base_dir) + if project_slug and not user.can_access_project(project_slug): + raise PermissionError("Нет доступа к этому файлу") + + with open(full_path, "r", encoding="utf-8") as f: + return f.read() + + +def get_download_path(user: UserContext, rel_path: str) -> Path: + full_path = _resolve_content_path(user, rel_path) + if not full_path.exists() or not full_path.is_file(): + raise FileNotFoundError(f"Файл не найден: {rel_path}") + base_dir = full_path.parent.parent + project_slug = _folder_project_slug(full_path.parent.name, base_dir) + if project_slug and not user.can_access_project(project_slug): + raise PermissionError("Нет доступа к этому файлу") + return full_path + + +def delete_folder(user: UserContext, folder_rel: str) -> None: + if folder_rel.startswith("documents/"): + folder_path = resolve_document_path(user.org_slug, folder_rel[len("documents/"):]) + base_dir = org_documents_dir(user.org_slug) + elif folder_rel.startswith("meetings/"): + folder_path = resolve_meeting_path(user.org_slug, folder_rel[len("meetings/"):]) + base_dir = org_meetings_dir(user.org_slug) + else: + folder_path = resolve_meeting_path(user.org_slug, folder_rel) + base_dir = org_meetings_dir(user.org_slug) + + if not folder_path.is_dir(): + raise FileNotFoundError("Folder not found") + project_slug = _folder_project_slug(folder_path.name, base_dir) + if project_slug and not user.can_access_project(project_slug): + raise PermissionError("Нет доступа к этой папке") + shutil.rmtree(folder_path) + + def get_task_status(task_id: str, user: Optional[UserContext] = None) -> Optional[Dict[str, Any]]: task = tasks.get(task_id) if not task: @@ -442,8 +635,8 @@ def _file_sort_key(file_info: Dict[str, Any]) -> tuple: return (4, name) -def _folder_project_slug(folder_name: str, meetings_dir: Path) -> Optional[str]: - meta_path = meetings_dir / folder_name / ".project.json" +def _folder_project_slug(folder_name: str, base_dir: Path) -> Optional[str]: + meta_path = base_dir / folder_name / ".project.json" if meta_path.exists(): try: data = json.loads(meta_path.read_text(encoding="utf-8")) @@ -451,75 +644,3 @@ def _folder_project_slug(folder_name: str, meetings_dir: Path) -> Optional[str]: except Exception: pass return None - - -def get_processed_tree(user: UserContext) -> List[Dict[str, Any]]: - tree = [] - meetings_dir = org_meetings_dir(user.org_slug) - if not meetings_dir.exists(): - return tree - - skip_suffixes = ("_segments.json",) - - for item in sorted(meetings_dir.iterdir()): - if not item.is_dir(): - continue - project_slug = _folder_project_slug(item.name, meetings_dir) - if project_slug and not user.can_access_project(project_slug): - continue - - files = [] - for f in sorted(item.iterdir(), key=lambda p: _file_sort_key({"name": p.name})): - if f.is_file() and not f.name.endswith(skip_suffixes) and f.name != ".project.json": - files.append({ - "name": f.name, - "path": str(f.relative_to(meetings_dir)), - "size": f.stat().st_size, - "ext": f.suffix.lower(), - "kind": "summary" if "_summary" in f.name else "protocol", - }) - tree.append({ - "name": item.name, - "path": str(item.relative_to(meetings_dir)), - "project": project_slug, - "files": files, - "created": datetime.fromtimestamp(item.stat().st_ctime).isoformat(), - }) - return tree - - -def read_file_content(user: UserContext, rel_path: str) -> str: - full_path = resolve_meeting_path(user.org_slug, rel_path) - if not full_path.exists() or not full_path.is_file(): - raise FileNotFoundError(f"Файл не найден: {rel_path}") - - folder = full_path.parent.name - meetings_dir = org_meetings_dir(user.org_slug) - project_slug = _folder_project_slug(folder, meetings_dir) - if project_slug and not user.can_access_project(project_slug): - raise PermissionError("Нет доступа к этому файлу") - - with open(full_path, "r", encoding="utf-8") as f: - return f.read() - - -def get_download_path(user: UserContext, rel_path: str) -> Path: - full_path = resolve_meeting_path(user.org_slug, rel_path) - if not full_path.exists() or not full_path.is_file(): - raise FileNotFoundError(f"Файл не найден: {rel_path}") - folder = full_path.parent.name - meetings_dir = org_meetings_dir(user.org_slug) - project_slug = _folder_project_slug(folder, meetings_dir) - if project_slug and not user.can_access_project(project_slug): - raise PermissionError("Нет доступа к этому файлу") - return full_path - - -def delete_folder(user: UserContext, folder_rel: str) -> None: - folder_path = resolve_meeting_path(user.org_slug, folder_rel) - if not folder_path.is_dir(): - raise FileNotFoundError("Folder not found") - project_slug = _folder_project_slug(folder_path.name, org_meetings_dir(user.org_slug)) - if project_slug and not user.can_access_project(project_slug): - raise PermissionError("Нет доступа к этой папке") - shutil.rmtree(folder_path) diff --git a/backend/static/app.js b/backend/static/app.js index 47553ec..b116b40 100644 --- a/backend/static/app.js +++ b/backend/static/app.js @@ -269,6 +269,7 @@ class TranscriptionApp { if (!files.length) return; const project = document.getElementById('uploadProjectSelect')?.value; + const docType = document.getElementById('uploadDocTypeSelect')?.value || 'other'; if (!project) { this.showToast('Выберите проект перед загрузкой', 'error'); return; @@ -276,6 +277,7 @@ class TranscriptionApp { const formData = new FormData(); formData.append('project', project); + formData.append('doc_type', docType); for (const file of files) { formData.append('files', file); } @@ -328,6 +330,7 @@ class TranscriptionApp { const parts = []; if (queue.processing) parts.push(`транскрибация: ${queue.processing}`); if (queue.pending_transcribe) parts.push(`в очереди ASR: ${queue.pending_transcribe}`); + if (queue.pending_ingest) parts.push(`ingest: ${queue.pending_ingest}`); if (queue.postprocessing) parts.push(`summary/RAG: ${queue.postprocessing}`); if (queue.pending_postprocess) parts.push(`в очереди post: ${queue.pending_postprocess}`); el.textContent = parts.length ? parts.join(' · ') : 'очередь пуста'; @@ -381,7 +384,7 @@ class TranscriptionApp {