193 lines
8.1 KiB
Python
193 lines
8.1 KiB
Python
import hashlib
|
|
import hmac
|
|
import json
|
|
import os
|
|
from contextlib import asynccontextmanager
|
|
from datetime import datetime, timezone
|
|
from pathlib import Path
|
|
from typing import Annotated, Any, Callable
|
|
|
|
from fastapi import Depends, FastAPI, File, Form, Header, HTTPException, Response, UploadFile
|
|
|
|
from .engine import OcrEngine
|
|
from .jobs import JobActiveError, JobQueue, QueueFullError, StoragePressureError
|
|
from .models import load_runtime_engine
|
|
from .version import release_version
|
|
|
|
|
|
MAX_UPLOAD_BYTES = 50 * 1024 * 1024
|
|
MAX_PAGES = 100
|
|
ALLOWED_CONFIG = {
|
|
"languages": ["es", "en"],
|
|
"dpi": 200,
|
|
"engine": "paddleocr",
|
|
"engineVersion": "3.4.0",
|
|
"runtimeVersion": "3.2.2",
|
|
"configVersion": "ocr-v2",
|
|
"returnLayout": True,
|
|
}
|
|
REQUEST_FIELDS = {"documentSha256", "pages", *ALLOWED_CONFIG}
|
|
|
|
|
|
def fail(status: int, code: str, message: str, retryable: bool = False, headers: dict[str, str] | None = None) -> None:
|
|
raise HTTPException(status, {"code": code, "message": message, "retryable": retryable}, headers)
|
|
|
|
|
|
def page_hash(pages: list[int]) -> str:
|
|
value = json.dumps(pages, separators=(",", ":")).encode()
|
|
return hashlib.sha256(value).hexdigest()
|
|
|
|
|
|
def request_identity(idempotency_key: str, document_sha256: str, pages: list[int], config_version: str) -> str:
|
|
canonical = json.dumps({
|
|
"configVersion": config_version,
|
|
"documentSha256": document_sha256,
|
|
"idempotencyKey": idempotency_key,
|
|
"requestedPages": pages,
|
|
}, sort_keys=True, separators=(",", ":")).encode()
|
|
return hashlib.sha256(canonical).hexdigest()
|
|
|
|
|
|
def utc_now() -> datetime:
|
|
return datetime.now(timezone.utc)
|
|
|
|
|
|
def create_app(
|
|
token: str,
|
|
db_path: str | Path = ":memory:",
|
|
max_upload_bytes: int = MAX_UPLOAD_BYTES,
|
|
engine_ready: bool = False,
|
|
engine: OcrEngine | None = None,
|
|
now: Callable[[], datetime] = utc_now,
|
|
) -> FastAPI:
|
|
version = release_version()
|
|
queue = JobQueue(db_path, now)
|
|
|
|
@asynccontextmanager
|
|
async def lifespan(_application: FastAPI):
|
|
queue.start_worker(engine)
|
|
try:
|
|
yield
|
|
finally:
|
|
queue.stop()
|
|
|
|
application = FastAPI(title="Private OCR Service", docs_url=None, redoc_url=None, lifespan=lifespan)
|
|
application.state.queue = queue
|
|
|
|
def authorize(authorization: Annotated[str | None, Header()] = None) -> None:
|
|
scheme, _, supplied = (authorization or "").partition(" ")
|
|
if not token or scheme != "Bearer" or not hmac.compare_digest(supplied, token):
|
|
fail(401, "UNAUTHORIZED", "Valid bearer authorization is required", headers={"WWW-Authenticate": "Bearer"})
|
|
|
|
@application.get("/health/live")
|
|
def live() -> dict[str, str]:
|
|
return {"status": "ok", "service": "ocr", "version": version,
|
|
"revision": os.getenv("BUILD_REVISION", "unknown")}
|
|
|
|
@application.get("/health/ready")
|
|
def ready(response: Response) -> dict[str, Any]:
|
|
state = queue.operational_state()
|
|
storage_available = queue.storage_available()
|
|
ready_state = engine_ready and state["workerOperational"] and state["sweeperOperational"] and storage_available
|
|
if not ready_state:
|
|
response.status_code = 503
|
|
return {"ready": ready_state, "engineLoaded": engine_ready, "storageAvailable": storage_available,
|
|
**state, "service": "ocr", "version": version,
|
|
"revision": os.getenv("BUILD_REVISION", "unknown")}
|
|
|
|
@application.post("/v1/jobs", status_code=202, dependencies=[Depends(authorize)])
|
|
async def create_job(
|
|
file: Annotated[UploadFile, File()], request: Annotated[str, Form()],
|
|
idempotency_key: Annotated[str | None, Header(alias="Idempotency-Key")] = None,
|
|
) -> dict[str, Any]:
|
|
content = await file.read(max_upload_bytes + 1)
|
|
if len(content) > max_upload_bytes:
|
|
fail(413, "UPLOAD_LIMIT_EXCEEDED", "PDF exceeds the upload limit")
|
|
try:
|
|
payload = json.loads(request)
|
|
except (json.JSONDecodeError, TypeError):
|
|
fail(400, "INVALID_REQUEST", "Request must be valid JSON")
|
|
if not isinstance(payload, dict) or set(payload) != REQUEST_FIELDS:
|
|
fail(400, "INVALID_REQUEST", "Request fields do not match the contract")
|
|
pages = payload["pages"]
|
|
if not isinstance(pages, list) or not pages or any(type(page) is not int or page < 1 for page in pages):
|
|
fail(422, "INVALID_PAGES", "Pages must be positive one-based integers")
|
|
if len(pages) > MAX_PAGES:
|
|
fail(413, "PAGE_LIMIT_EXCEEDED", "OCR jobs accept at most 100 pages")
|
|
if pages != sorted(set(pages)):
|
|
fail(422, "INVALID_PAGES", "Pages must be unique and ordered")
|
|
if any(payload[name] != value for name, value in ALLOWED_CONFIG.items()):
|
|
fail(400, "CONFIG_NOT_ALLOWED", "OCR configuration is not allowlisted")
|
|
digest = hashlib.sha256(content).hexdigest()
|
|
if payload["documentSha256"] != digest:
|
|
fail(422, "INTEGRITY_MISMATCH", "PDF bytes do not match documentSha256")
|
|
if file.content_type != "application/pdf" or not content.startswith(b"%PDF-") or b"/Encrypt" in content:
|
|
fail(422, "UNSUPPORTED_PDF", "PDF is corrupt, encrypted, or unsupported")
|
|
expected_key = f'{digest}:ocr-v2:{page_hash(pages)}'
|
|
if idempotency_key != expected_key and not queue.contains(idempotency_key or ""):
|
|
fail(400, "INVALID_IDEMPOTENCY_KEY", "Idempotency-Key does not match request identity")
|
|
try:
|
|
ack, _created = queue.submit(idempotency_key or "", payload, content)
|
|
except ValueError:
|
|
fail(409, "IDEMPOTENCY_CONFLICT", "The idempotency key is already bound to another request")
|
|
except QueueFullError:
|
|
fail(429, "QUEUE_FULL", "The OCR queue is full", True, {"Retry-After": "2"})
|
|
except StoragePressureError:
|
|
fail(503, "OCR_STORAGE_PRESSURE", "OCR storage capacity is unavailable", True, {"Retry-After": "15"})
|
|
return ack
|
|
|
|
@application.get("/v1/jobs/{job_id}", dependencies=[Depends(authorize)])
|
|
def get_job(job_id: str) -> dict[str, Any]:
|
|
status = queue.status(job_id)
|
|
if status is None:
|
|
fail(404, "JOB_NOT_FOUND", "OCR job does not exist")
|
|
return status
|
|
|
|
@application.get("/v1/jobs/{job_id}/result", dependencies=[Depends(authorize)])
|
|
def get_result(job_id: str) -> dict[str, Any]:
|
|
stored = queue.result(job_id)
|
|
if stored is None:
|
|
fail(404, "JOB_NOT_FOUND", "OCR job does not exist")
|
|
if stored[0] != "succeeded" or stored[1] is None:
|
|
fail(409, "RESULT_NOT_READY", "OCR job result is not ready")
|
|
return stored[1]
|
|
|
|
@application.get("/v1/jobs/{job_id}/pages/{page}/image", dependencies=[Depends(authorize)])
|
|
def get_review_image(job_id: str, page: int) -> Response:
|
|
try:
|
|
stored = queue.review_image(job_id, page)
|
|
except RuntimeError:
|
|
fail(409, "RESULT_NOT_READY", "OCR review image is not ready")
|
|
except ValueError:
|
|
fail(422, "INVALID_PAGE", "OCR review image does not exist")
|
|
if stored is None:
|
|
fail(404, "JOB_NOT_FOUND", "OCR job does not exist")
|
|
identity, png = stored
|
|
return Response(png, media_type="image/png", headers={
|
|
"X-Ocr-Job-Id": job_id,
|
|
"X-Ocr-Identity-Sha256": identity["requestIdentitySha256"],
|
|
"X-Ocr-Config-Version": identity["configVersion"],
|
|
"X-Document-Sha256": identity["documentSha256"],
|
|
"X-Content-Sha256": hashlib.sha256(png).hexdigest(),
|
|
"X-Page-Number": str(page),
|
|
})
|
|
|
|
@application.delete("/v1/jobs/{job_id}", status_code=204, dependencies=[Depends(authorize)])
|
|
def delete_job(job_id: str) -> Response:
|
|
try:
|
|
queue.delete(job_id)
|
|
except JobActiveError:
|
|
fail(409, "JOB_ACTIVE", "An active OCR job cannot be deleted")
|
|
return Response(status_code=204)
|
|
|
|
return application
|
|
|
|
|
|
runtime_engine = load_runtime_engine()
|
|
|
|
app = create_app(
|
|
os.getenv("OCR_INTERNAL_TOKEN", ""),
|
|
os.getenv("OCR_JOBS_DB", ":memory:"),
|
|
engine_ready=runtime_engine is not None,
|
|
engine=runtime_engine,
|
|
)
|