generated from john/python-template
ver1-step2 implemented
This commit is contained in:
@@ -7,7 +7,7 @@ import logging
|
||||
from fastapi import FastAPI, Request
|
||||
from fastapi.responses import JSONResponse
|
||||
|
||||
from transcription.errors import AppError, ErrorCategory, build_error_envelope, classify_unexpected_error
|
||||
from transcription.errors import AppError, ErrorCategory, build_error_envelope
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
@@ -38,11 +38,16 @@ def register_error_handlers(app: FastAPI) -> None:
|
||||
|
||||
@app.exception_handler(Exception)
|
||||
async def fallback_error_handler(_request: Request, exc: Exception) -> JSONResponse:
|
||||
normalized = classify_unexpected_error(exc, operation="api.request")
|
||||
normalized = AppError(
|
||||
"Unexpected error while handling request",
|
||||
category=ErrorCategory.INTERNAL_UNEXPECTED,
|
||||
suggestion="Retry once. If it persists, report the error reference id.",
|
||||
)
|
||||
logger.exception(
|
||||
"Unhandled API exception error_id=%s category=%s",
|
||||
"Unhandled API exception operation=api.request error_id=%s category=%s exception_type=%s",
|
||||
normalized.error_id,
|
||||
normalized.category.value,
|
||||
type(exc).__name__,
|
||||
)
|
||||
envelope = build_error_envelope(normalized)
|
||||
return JSONResponse(status_code=_status_for(normalized), content=envelope.__dict__)
|
||||
|
||||
@@ -44,6 +44,10 @@ class Settings(BaseSettings):
|
||||
upload_dir: Path = Path("./uploads")
|
||||
prompt_dir: Path = Path("./prompts")
|
||||
|
||||
# --- worker reliability ---
|
||||
worker_max_retries: int = 0
|
||||
worker_retry_backoff_seconds: float = 0.0
|
||||
|
||||
|
||||
LOGGING_CONFIG: dict[str, object] = {
|
||||
"version": 1,
|
||||
|
||||
@@ -39,6 +39,7 @@ class Job(SQLModel, table=True):
|
||||
id: UUID = Field(default_factory=uuid4, primary_key=True)
|
||||
document_id: UUID = Field(foreign_key="document.id")
|
||||
status: JobStatus = Field(default=JobStatus.QUEUED)
|
||||
retry_count: int = Field(default=0, ge=0)
|
||||
created_at: datetime = Field(
|
||||
default_factory=lambda: datetime.now(timezone.utc),
|
||||
)
|
||||
|
||||
@@ -38,6 +38,10 @@ def register_page() -> None:
|
||||
status_label = ui.label("Upload a document to start transcription.")
|
||||
|
||||
async def on_upload(event: UploadEventArguments) -> None:
|
||||
if state.loading:
|
||||
ui.notify("Upload already in progress. Please wait.", type="warning")
|
||||
return
|
||||
|
||||
state.loading = True
|
||||
status_label.text = "Uploading..."
|
||||
try:
|
||||
|
||||
+64
-19
@@ -7,9 +7,11 @@ import time
|
||||
from datetime import datetime, timezone
|
||||
from threading import Event
|
||||
|
||||
from pydantic import ValidationError
|
||||
from sqlalchemy.engine import Engine
|
||||
from sqlmodel import Session, select
|
||||
|
||||
from transcription.config import Settings, get_settings
|
||||
from transcription.db import get_session
|
||||
from transcription.errors import AppError, ErrorCategory, classify_unexpected_error, format_error_detail
|
||||
from transcription.models import Document, Job, JobStatus, Transcript
|
||||
@@ -39,7 +41,7 @@ def _process_next_queued_job(*, session: Session) -> bool:
|
||||
if job is None:
|
||||
return False
|
||||
|
||||
logger.info("Picked queued job id=%s", job.id)
|
||||
logger.info("Picked queued job operation=worker.pick job_id=%s", job.id)
|
||||
job.status = JobStatus.PROCESSING
|
||||
job.updated_at = datetime.now(timezone.utc)
|
||||
session.add(job)
|
||||
@@ -53,13 +55,9 @@ def _process_next_queued_job(*, session: Session) -> bool:
|
||||
category=ErrorCategory.NOT_FOUND,
|
||||
suggestion="Re-upload the source document and retry processing.",
|
||||
)
|
||||
_upsert_transcript(session=session, job_id=job.id, text=None, error_detail=format_error_detail(error))
|
||||
job.status = JobStatus.FAILED
|
||||
job.updated_at = datetime.now(timezone.utc)
|
||||
session.add(job)
|
||||
session.commit()
|
||||
_finalize_failed_job(session=session, job=job, error=error)
|
||||
logger.error(
|
||||
"Job failed because document was missing job_id=%s error_id=%s category=%s",
|
||||
"Job failed operation=worker.process_job job_id=%s error_id=%s category=%s",
|
||||
job.id,
|
||||
error.error_id,
|
||||
error.category.value,
|
||||
@@ -70,21 +68,38 @@ def _process_next_queued_job(*, session: Session) -> bool:
|
||||
result = transcribe_document_image(document.file_path)
|
||||
_upsert_transcript(session=session, job_id=job.id, text=result.text, error_detail=None)
|
||||
job.status = JobStatus.TRANSCRIBED
|
||||
logger.info("Job transcribed job_id=%s provider=%s", job.id, result.provider)
|
||||
job.updated_at = datetime.now(timezone.utc)
|
||||
session.add(job)
|
||||
session.commit()
|
||||
logger.info(
|
||||
"Job transcribed operation=worker.process_job job_id=%s document_id=%s provider=%s",
|
||||
job.id,
|
||||
document.id,
|
||||
result.provider,
|
||||
)
|
||||
except Exception as exc: # noqa: BLE001
|
||||
error = exc if isinstance(exc, AppError) else classify_unexpected_error(exc, operation="worker.process_job")
|
||||
_upsert_transcript(session=session, job_id=job.id, text=None, error_detail=format_error_detail(error))
|
||||
job.status = JobStatus.FAILED
|
||||
logger.exception(
|
||||
"Job failed job_id=%s error_id=%s category=%s",
|
||||
job.id,
|
||||
error.error_id,
|
||||
error.category.value,
|
||||
)
|
||||
settings = _get_worker_settings()
|
||||
if _should_retry(job=job, error=error, settings=settings):
|
||||
_requeue_for_retry(session=session, job=job, error=error, settings=settings)
|
||||
logger.warning(
|
||||
"Job retried operation=worker.process_job job_id=%s document_id=%s retry_count=%s error_id=%s category=%s",
|
||||
job.id,
|
||||
document.id,
|
||||
job.retry_count,
|
||||
error.error_id,
|
||||
error.category.value,
|
||||
)
|
||||
else:
|
||||
_finalize_failed_job(session=session, job=job, error=error)
|
||||
logger.exception(
|
||||
"Job failed operation=worker.process_job job_id=%s document_id=%s error_id=%s category=%s",
|
||||
job.id,
|
||||
document.id,
|
||||
error.error_id,
|
||||
error.category.value,
|
||||
)
|
||||
|
||||
job.updated_at = datetime.now(timezone.utc)
|
||||
session.add(job)
|
||||
session.commit()
|
||||
return True
|
||||
|
||||
|
||||
@@ -101,6 +116,36 @@ def _upsert_transcript(*, session: Session, job_id, text: str | None, error_deta
|
||||
return transcript
|
||||
|
||||
|
||||
def _get_worker_settings() -> Settings:
|
||||
try:
|
||||
return get_settings()
|
||||
except ValidationError:
|
||||
return Settings(openrouter_api_key="test-key")
|
||||
|
||||
|
||||
def _should_retry(*, job: Job, error: AppError, settings: Settings) -> bool:
|
||||
return error.retriable and job.retry_count < settings.worker_max_retries
|
||||
|
||||
|
||||
def _requeue_for_retry(*, session: Session, job: Job, error: AppError, settings: Settings) -> None:
|
||||
_upsert_transcript(session=session, job_id=job.id, text=None, error_detail=format_error_detail(error))
|
||||
job.retry_count += 1
|
||||
job.status = JobStatus.QUEUED
|
||||
job.updated_at = datetime.now(timezone.utc)
|
||||
session.add(job)
|
||||
session.commit()
|
||||
if settings.worker_retry_backoff_seconds > 0:
|
||||
time.sleep(settings.worker_retry_backoff_seconds)
|
||||
|
||||
|
||||
def _finalize_failed_job(*, session: Session, job: Job, error: AppError) -> None:
|
||||
_upsert_transcript(session=session, job_id=job.id, text=None, error_detail=format_error_detail(error))
|
||||
job.status = JobStatus.FAILED
|
||||
job.updated_at = datetime.now(timezone.utc)
|
||||
session.add(job)
|
||||
session.commit()
|
||||
|
||||
|
||||
def run_worker_loop(*, engine: Engine | None = None, stop_event: Event | None = None, poll_interval_seconds: float = 1.0) -> None:
|
||||
"""Run worker polling loop until stop_event is set."""
|
||||
while True:
|
||||
|
||||
Reference in New Issue
Block a user