Download src/api/v2/documents.py from DataEyond/Agentic-Service-Data-Eyond-Catalog: direct link, hf CLI and curl.
- Browser
- Download file 9.68 kB
-
https://huggingface.co/spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/resolve/main/src/api/v2/documents.py
- Command line
-
hf download hf://spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/src/api/v2/documents.py
-
curl -L -o documents.py https://huggingface.co/spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/resolve/main/src/api/v2/documents.py
9.68 kB
| """v2 document routes — unstructured files (pdf / docx / txt) only. | |
| UNSTRUCTURED_V2_PLAN U3. The FE calls these instead of Go's `/api/v1/document*` | |
| for unstructured types; tabular files (csv / xlsx) stay entirely on Go, and so | |
| does **delete** for every type (U-D7 — Go's delete is authenticated and already | |
| removes the original and every embedding by `document_id`, v2 rows included). | |
| Request/response shapes, messages and status codes MIRROR Go's handlers | |
| (`Orchestrator-Agent-Service/internal/documents/handler.go`, read 2026-09-23) so | |
| the FE changes only the URL prefix. Envelope: `{status, message, data}`. | |
| ⚠️ Unlike Go, these routes do not verify a bearer token — the whole Python | |
| surface is unauthenticated by design until DEV_PLAN #43. Recorded under U-D7. | |
| """ | |
| # No `from __future__ import annotations` here: FastAPI must resolve the | |
| # `UploadFile | None` parameter annotation at import time. | |
| import uuid | |
| from pathlib import Path | |
| from types import SimpleNamespace | |
| from typing import Any | |
| from fastapi import APIRouter, BackgroundTasks, Depends, File, Form, Request, UploadFile | |
| from fastapi.responses import JSONResponse | |
| from sqlalchemy import select | |
| from sqlalchemy.ext.asyncio import AsyncSession | |
| from src.db.postgres.connection import get_db | |
| from src.db.postgres.models import Document | |
| from src.document.image_token import resolve_image_urls | |
| from src.document.ingest_v2 import UNSTRUCTURED_TYPES, document_ingest_v2 | |
| from src.middlewares.logging import get_logger | |
| from src.middlewares.rate_limit import limiter | |
| logger = get_logger("documents_v2_api") | |
| router = APIRouter(prefix="/api/v2", tags=["Documents v2"]) | |
| MAX_FILE_SIZE_BYTES = 10 * 1024 * 1024 # Go's maxFileSizeBytes | |
| _CONTENT_TYPES = { | |
| "pdf": "application/pdf", | |
| "docx": "application/vnd.openxmlformats-officedocument.wordprocessingml.document", | |
| "txt": "text/plain; charset=utf-8", | |
| } | |
| def _envelope(code: int, status: str, message: str, data: Any = None) -> JSONResponse: | |
| body: dict[str, Any] = {"status": status, "message": message} | |
| if data is not None: | |
| body["data"] = data | |
| return JSONResponse(status_code=code, content=body) | |
| def _error(code: int, message: str) -> JSONResponse: | |
| return _envelope(code, "error", message) | |
| def _doc_json(d: Document) -> dict[str, Any]: | |
| """Go's Document model. `processed_at` / `error_message` are omitempty there.""" | |
| out: dict[str, Any] = { | |
| "id": d.id, | |
| "user_id": d.user_id, | |
| "filename": d.filename, | |
| "blob_name": d.blob_name, | |
| "file_size": d.file_size or 0, | |
| "file_type": d.file_type or "", | |
| "status": d.status or "", | |
| "chunks_count": getattr(d, "chunks_count", None) or 0, | |
| "created_at": d.created_at.isoformat() if d.created_at else None, | |
| } | |
| if d.processed_at: | |
| out["processed_at"] = d.processed_at.isoformat() | |
| if d.error_message: | |
| out["error_message"] = d.error_message | |
| return out | |
| def _extension(filename: str) -> str: | |
| return Path(filename or "").suffix.lstrip(".").lower() | |
| async def list_doctypes(): | |
| data = [ | |
| {"type": t, "max_size_mb": 10, "status": "active", "message": None} | |
| for t in UNSTRUCTURED_TYPES | |
| ] | |
| return _envelope(200, "success", "supported document types", data) | |
| async def list_documents(user_id: str, db: AsyncSession = Depends(get_db)): | |
| """Every document the user owns (all types — same table and order as Go's list).""" | |
| try: | |
| rows = ( | |
| await db.execute( | |
| select(Document) | |
| .where(Document.user_id == user_id) | |
| .order_by(Document.created_at.desc()) | |
| ) | |
| ).scalars().all() | |
| except Exception as e: | |
| logger.error("list documents failed", user_id=user_id, error=repr(e)) | |
| return _error(500, "internal server error") | |
| return _envelope(200, "success", "documents", [_doc_json(d) for d in rows]) | |
| # Go: WithRateLimit(10, …) on this route | |
| async def upload_document( | |
| request: Request, | |
| user_id: str = Form(""), | |
| file: UploadFile | None = File(None), | |
| db: AsyncSession = Depends(get_db), | |
| ): | |
| if not user_id.strip(): | |
| return _error(400, "user_id is required") | |
| if file is None: | |
| return _error(400, "file is required") | |
| ext = _extension(file.filename or "") | |
| if ext not in UNSTRUCTURED_TYPES: | |
| return _error(400, "unsupported file type") | |
| content = await file.read(MAX_FILE_SIZE_BYTES + 1) | |
| if len(content) > MAX_FILE_SIZE_BYTES: | |
| return _error(400, "file too large (max 10 MB)") | |
| from src.storage.object_storage.supabase_s3 import object_storage | |
| doc_id = str(uuid.uuid4()) | |
| doc = Document( | |
| id=doc_id, | |
| user_id=user_id, | |
| filename=file.filename, | |
| blob_name=f"{user_id}/{doc_id}.{ext}", # Go's exact blob_name format | |
| file_size=len(content), | |
| file_type=ext, | |
| status="uploading", | |
| ) | |
| try: | |
| db.add(doc) | |
| await db.commit() | |
| await db.refresh(doc) | |
| except Exception as e: | |
| logger.error("upload: db record failed", user_id=user_id, error=repr(e)) | |
| return _error(500, "upload failed") | |
| try: | |
| await object_storage.upload_bytes(content, doc.blob_name, _CONTENT_TYPES[ext]) | |
| except Exception as e: | |
| doc.status, doc.error_message = "failed", repr(e)[:500] | |
| await db.commit() | |
| return _error(500, "upload failed") | |
| doc.status = "uploaded" | |
| await db.commit() | |
| await db.refresh(doc) | |
| logger.info("upload: completed", document_id=doc_id, user_id=user_id, size=len(content)) | |
| return _envelope(201, "success", "file uploaded", _doc_json(doc)) | |
| def _json_body(properties: dict[str, Any], example: dict[str, Any]) -> dict[str, Any]: | |
| """OpenAPI-only body description (added 2026-09-24). | |
| These routes read `request.json()` themselves so their errors stay Go-shaped (400 | |
| `invalid request body`, not FastAPI's 422). The side effect: FastAPI declared no | |
| request body, so Swagger's "Try it out" showed no editor and POSTed nothing — | |
| every Swagger test answered 400. This documents the body; validation is unchanged. | |
| """ | |
| return {"requestBody": {"required": True, "content": {"application/json": { | |
| "schema": {"type": "object", "required": list(properties), "properties": properties}, | |
| "example": example, | |
| }}}} | |
| async def process_document( | |
| request: Request, | |
| background: BackgroundTasks, | |
| db: AsyncSession = Depends(get_db), | |
| ): | |
| try: | |
| body = await request.json() | |
| document_id = str(body.get("document_id") or "") | |
| user_id = str(body.get("user_id") or "") | |
| except Exception: | |
| return _error(400, "invalid request body") | |
| if not document_id or not user_id: | |
| return _error(400, "document_id and user_id are required") | |
| doc = await db.get(Document, document_id) | |
| if doc is None: | |
| return _error(404, "document not found") | |
| if doc.user_id != user_id: | |
| return _error(403, "access denied") | |
| if (doc.file_type or "").lower() not in UNSTRUCTURED_TYPES: | |
| # Tabular processing (parquet + catalog) is Go's; never half-do it here. | |
| return _error(400, "unsupported file type") | |
| doc.status, doc.error_message = "processing", None | |
| await db.commit() | |
| snapshot = SimpleNamespace( | |
| id=doc.id, user_id=doc.user_id, filename=doc.filename, | |
| file_type=(doc.file_type or "").lower(), blob_name=doc.blob_name, | |
| ) | |
| background.add_task(document_ingest_v2.process, snapshot) | |
| return _envelope( | |
| 202, "success", "document processing started", | |
| {"document_id": doc.id, "file_type": snapshot.file_type, "status": "processing"}, | |
| ) | |
| async def image_token(request: Request): | |
| """Resolve `asset://<id>` markers to 1-hour signed URLs (U4). See image_token.py.""" | |
| try: | |
| body = await request.json() | |
| user_id = str(body.get("user_id") or "") | |
| analysis_id = str(body.get("analysis_id") or "") | |
| asset_ids = body.get("asset_ids") | |
| except Exception: | |
| return _error(400, "invalid request body") | |
| if not user_id or not analysis_id or not isinstance(asset_ids, list): | |
| return _error(400, "user_id, analysis_id and asset_ids are required") | |
| if len(asset_ids) > 50: | |
| return _error(400, "at most 50 asset_ids per request") | |
| try: | |
| urls, expires_at = await resolve_image_urls(user_id, analysis_id, asset_ids) | |
| except Exception as e: | |
| logger.error("image-token failed", user_id=user_id, error=repr(e)) | |
| return _error(500, "internal server error") | |
| return _envelope(200, "success", "image urls", {"urls": urls, "expires_at": expires_at}) | |