DealDocumentScreening/src/contract_check/api/services.py

161 lines
5.9 KiB
Python

"""Shared api upload/enqueue service (used by Telegram and B2B adapters)."""
from __future__ import annotations
import uuid
from typing import TYPE_CHECKING
from fastapi import HTTPException
from sqlalchemy import text
from src.contract_check.api.schemas import DocumentUploadResponse
from src.contract_check.core.billing.errors import NoCredits
from src.contract_check.core.billing.quota import (
release_document_slot,
reserve_document_slot,
)
from src.contract_check.core.config import get_settings
from src.contract_check.core.credits import adjust_credits, reserve_credit
from src.contract_check.core.db.repositories import DocumentRepository, JobRepository
from src.contract_check.core.db.repositories.credits import CreditsRepository
from src.contract_check.core.extraction.formats import SUPPORTED_SUFFIXES
from src.contract_check.core.logging import get_logger, new_correlation_id
from src.contract_check.core.metrics import credits_reserved, documents_uploaded
from src.contract_check.core.mq.messages import DocumentUploaded
from src.contract_check.core.mq.topology import RK_EXTRACT
from src.contract_check.core.s3 import original_key
if TYPE_CHECKING:
from fastapi import UploadFile
from sqlalchemy.ext.asyncio import AsyncSession
from src.contract_check.core.mq.publisher import Publisher
from src.contract_check.core.s3.port import Storage
log = get_logger(__name__)
def _content_type_from_suffix(suffix: str) -> str:
return {
".pdf": "application/pdf",
".docx": "application/vnd.openxmlformats-officedocument.wordprocessingml.document",
}.get(suffix, "application/octet-stream")
async def upload_and_enqueue(
session: AsyncSession,
storage: Storage,
publisher: Publisher,
user_id: uuid.UUID,
file: UploadFile,
) -> DocumentUploadResponse:
"""Validate, store in MinIO, create document/job rows, reserve credit, publish.
Returns 202 payload: {document_id, correlation_id, credits_left}.
Raises HTTPException on validation, credit, storage, DB, or MQ failure.
"""
if file.filename is None:
raise HTTPException(status_code=400, detail="filename is required")
suffix = file.filename.lower().split(".")[-1] if "." in file.filename else ""
suffix = f".{suffix}" if suffix else ""
if suffix not in SUPPORTED_SUFFIXES:
raise HTTPException(
status_code=400,
detail=f"unsupported format {suffix!r}; supported: {sorted(SUPPORTED_SUFFIXES)}",
)
credits_repo = CreditsRepository(session)
doc_repo = DocumentRepository(session)
job_repo = JobRepository(session)
document_id = uuid.uuid4()
correlation_id = new_correlation_id()
s3_key = original_key(str(user_id), str(document_id), suffix)
content_type = file.content_type or _content_type_from_suffix(suffix)
try:
data = await file.read()
if len(data) == 0:
raise HTTPException(status_code=400, detail="empty file")
await storage.put(s3_key, data, content_type=content_type)
except HTTPException:
raise
except Exception as exc:
log.error("s3_upload_failed", document_id=str(document_id), error=str(exc))
raise HTTPException(status_code=500, detail="failed to store document") from exc
try:
await doc_repo.create(
document_id=document_id,
user_id=user_id,
s3_key=s3_key,
filename=file.filename or "document",
mime=content_type,
bytes_=len(data),
status="queued",
)
await job_repo.create(
document_id=document_id,
correlation_id=uuid.UUID(correlation_id),
queue="extract",
status="pending",
)
except Exception as exc:
log.error("db_enqueue_failed", document_id=str(document_id), error=str(exc))
raise HTTPException(status_code=500, detail="failed to enqueue document") from exc
# Reserve credit *after* the document row exists so the ledger can link it.
result = await session.execute(
text("SELECT billing_hold FROM users WHERE id = :u"),
{"u": user_id},
)
row = result.first()
if row is not None and row[0]:
raise HTTPException(status_code=402, detail="billing hold")
source = "credits"
if get_settings().plans_enabled:
try:
source = await reserve_document_slot(session, user_id, document_id)
except NoCredits:
raise HTTPException(status_code=402, detail="no credits available") from None
else:
if not await reserve_credit(session, user_id, document_id=document_id):
raise HTTPException(status_code=402, detail="no credits available")
await session.commit()
credits_reserved.inc()
try:
msg = DocumentUploaded(
correlation_id=uuid.UUID(correlation_id),
document_id=document_id,
user_id=user_id,
s3_key=s3_key,
filename=file.filename,
mime=content_type,
)
await publisher.publish(msg, routing_key=RK_EXTRACT)
documents_uploaded.inc()
except Exception as exc:
log.error("mq_publish_failed", document_id=str(document_id), error=str(exc))
if source == "quota":
await release_document_slot(session, document_id)
else:
await adjust_credits(
session,
user_id,
1,
kind="refund_auto",
document_id=document_id,
)
await doc_repo.mark_failed(document_id, stage="publish_failed")
await session.commit()
raise HTTPException(status_code=500, detail="failed to publish job") from exc
credits_left = await credits_repo.get_balance(user_id)
return DocumentUploadResponse(
document_id=str(document_id),
correlation_id=str(correlation_id),
credits_left=credits_left,
)