| import os |
|
|
| from uuid import UUID |
|
|
| from fastapi import APIRouter, Depends, File, HTTPException, UploadFile |
| from sqlalchemy.orm import Session, selectinload |
|
|
| from app.core.dependencies import get_current_user |
| from app.core.file_types import ALLOWED_EXTENSIONS |
| from app.database.session import get_db |
| from app.models.document import Document |
| from app.models.document_version import DocumentVersion |
| from app.models.user import User |
| from app.models.workflow_run import WorkflowRun, WorkflowStatus |
| from app.models.workspace import Workspace |
| from app.schemas.document import DocumentUploadResponse |
| from app.services.document_service import DocumentService |
| from app.services.workspace_service import WorkspaceService |
| router = APIRouter() |
|
|
|
|
| @router.get("") |
| def list_documents( |
| workspace_id: str, |
| current_user: User = Depends(get_current_user), |
| db: Session = Depends(get_db), |
| ): |
| workspace = WorkspaceService(db).get_workspace(workspace_id) |
| if workspace is None: |
| raise HTTPException(status_code=404, detail="Workspace not found.") |
| if workspace.created_by != current_user.id: |
| raise HTTPException(status_code=403, detail="You do not have access to this workspace.") |
|
|
| documents = ( |
| db.query(Document) |
| .options( |
| selectinload(Document.versions).selectinload(DocumentVersion.workflow_runs) |
| ) |
| .filter(Document.workspace_id == workspace_id) |
| .order_by(Document.created_at.desc()) |
| .all() |
| ) |
|
|
| results = [] |
| for doc in documents: |
| latest_version = doc.versions[-1] if doc.versions else None |
| workflow_runs = latest_version.workflow_runs if latest_version else [] |
| |
| latest_workflow = max(workflow_runs, key=lambda w: w.started_at, default=None) if workflow_runs else None |
|
|
| results.append({ |
| "id": str(doc.id), |
| "title": doc.title, |
| "document_type": doc.document_type, |
| "created_at": doc.created_at, |
| "version_id": str(latest_version.id) if latest_version else None, |
| "filename": latest_version.filename if latest_version else doc.title, |
| "processing_status": latest_version.status.value if latest_version else None, |
| "uploaded_at": latest_version.uploaded_at if latest_version else doc.created_at, |
| "workflow_id": str(latest_workflow.id) if latest_workflow else None, |
| "workflow_status": latest_workflow.status.value if latest_workflow else None, |
| }) |
|
|
| return results |
|
|
|
|
| @router.post("/upload") |
| async def upload_document( |
| workspace_id: str, |
| files: list[UploadFile] = File(...), |
| current_user: User = Depends(get_current_user), |
| db: Session = Depends(get_db), |
| ): |
|
|
| service = DocumentService(db) |
| workspace = WorkspaceService(db).get_workspace(workspace_id) |
|
|
| if workspace is None: |
| raise HTTPException( |
| status_code=404, |
| detail="Workspace not found.", |
| ) |
|
|
| if workspace.created_by != current_user.id: |
| raise HTTPException( |
| status_code=403, |
| detail="You do not have access to this workspace.", |
| ) |
| uploaded = [] |
| failed = [] |
|
|
| for file in files: |
|
|
| ext = os.path.splitext(file.filename)[1].lower() |
|
|
| if ext not in ALLOWED_EXTENSIONS: |
| failed.append( |
| { |
| "filename": file.filename, |
| "detail": "Unsupported file type.", |
| } |
| ) |
| continue |
|
|
| |
| DOC_TYPE_MAP = { |
| ".pdf": "PDF", ".docx": "DOCX", ".doc": "DOC", |
| ".txt": "TXT", ".md": "TXT", ".csv": "CSV", |
| ".xlsx": "XLSX", ".xls": "XLS", |
| ".pptx": "PPTX", ".ppt": "PPT", |
| ".png": "IMAGE", ".jpg": "IMAGE", ".jpeg": "IMAGE", |
| } |
| doc_type = DOC_TYPE_MAP.get(ext, "GENERAL") |
|
|
| try: |
|
|
| document, version, workflow = await service.upload_document( |
| workspace_id=workspace_id, |
| uploaded_by=current_user.id, |
| file=file, |
| document_type=doc_type, |
| ) |
|
|
| uploaded.append( |
| DocumentUploadResponse( |
| document_id=document.id, |
| version_id=version.id, |
| workflow_id=workflow.id, |
| filename=version.filename, |
| processing_stage=version.status.value, |
| uploaded_at=version.uploaded_at, |
| ) |
| ) |
|
|
| except Exception as e: |
| failed.append( |
| { |
| "filename": file.filename, |
| "detail": str(e), |
| } |
| ) |
|
|
| return { |
| "uploaded": uploaded, |
| "failed": failed, |
| } |
|
|
|
|
| @router.delete("/{document_id}", status_code=204) |
| def delete_document( |
| document_id: UUID, |
| current_user: User = Depends(get_current_user), |
| db: Session = Depends(get_db), |
| ): |
| """ |
| Delete a document and all its versions, knowledge items, proposals, |
| and cancel any active workflows. |
| """ |
| document = db.query(Document).filter(Document.id == document_id).first() |
|
|
| if document is None: |
| raise HTTPException(status_code=404, detail="Document not found.") |
|
|
| workspace = WorkspaceService(db).get_workspace(document.workspace_id) |
| if workspace is None or workspace.created_by != current_user.id: |
| raise HTTPException(status_code=403, detail="You do not have access to this document.") |
|
|
| |
| for version in document.versions: |
| for workflow in version.workflow_runs: |
| if workflow.status in (WorkflowStatus.PENDING, WorkflowStatus.RUNNING, WorkflowStatus.WAITING_FOR_REVIEW): |
| workflow.status = WorkflowStatus.CANCELLED |
| db.add(workflow) |
|
|
| |
| |
| db.delete(document) |
| db.commit() |
|
|
|
|
| @router.post("/{document_id}/retry") |
| def retry_document( |
| document_id: UUID, |
| current_user: User = Depends(get_current_user), |
| db: Session = Depends(get_db), |
| ): |
| """ |
| Retry processing a failed document. Creates a new workflow run |
| for the latest version and processes it from scratch. |
| """ |
| import threading |
| from app.services.document_service import _run_workflow_background, _active_workflows, _active_workflows_lock |
|
|
| document = db.query(Document).filter(Document.id == document_id).first() |
| if document is None: |
| raise HTTPException(status_code=404, detail="Document not found.") |
|
|
| workspace = WorkspaceService(db).get_workspace(document.workspace_id) |
| if workspace is None or workspace.created_by != current_user.id: |
| raise HTTPException(status_code=403, detail="You do not have access to this document.") |
|
|
| latest_version = document.versions[-1] if document.versions else None |
| if latest_version is None: |
| raise HTTPException(status_code=400, detail="No version to retry.") |
|
|
| |
| last_workflow = latest_version.workflow_runs[-1] if latest_version.workflow_runs else None |
| if last_workflow and last_workflow.status in (WorkflowStatus.PENDING, WorkflowStatus.RUNNING): |
| last_workflow.status = WorkflowStatus.CANCELLED |
| db.add(last_workflow) |
| db.commit() |
|
|
| |
| from app.services.workflow_service import WorkflowService |
| workflow_service = WorkflowService(db) |
| workflow = workflow_service.start_workflow( |
| workspace_id=document.workspace_id, |
| document_version_id=latest_version.id, |
| ) |
|
|
| |
| thread = threading.Thread( |
| target=_run_workflow_background, |
| args=( |
| workflow.id, |
| document.workspace_id, |
| latest_version.id, |
| latest_version.storage_path, |
| ), |
| daemon=False, |
| ) |
| with _active_workflows_lock: |
| _active_workflows[str(workflow.id)] = thread |
| thread.start() |
|
|
| return { |
| "message": "Retry started", |
| "workflow_id": str(workflow.id), |
| "document_id": str(document.id), |
| } |
|
|