29 changed files with 1748 additions and 199 deletions
+53 -7
View File
@@ -1,8 +1,54 @@
PROVIDER=openrouter # --- NiceGUI Server ---
OPENROUTER_API_KEY=sk-or-... # HOST=`0.0.0.0` (default)
# PROVIDER_MODEL= # optional: OpenRouter adapter supplies default # PORT=8000 (default)
# LOG_LEVEL: [`critical`, `error`, `warning`, `info` (default), `debug`, `trace`]
# RELOAD=false (default)
# --- AI provider ---
# PROVIDER=[`openrouter`(default), `google_genai`]
PROVIDER=openrouter
# OPENROUTER_API_KEY - Required when `PROVIDER=openrouter`
OPENROUTER_API_KEY=your-api-key-goes-here
# GEMINI_API_KEY - Required when `PROVIDER=google_genai`
# PROVIDER_MODEL= specify model. If left blank OpenRouter will supply default.
PROVIDER_MODEL=google/gemini-2.5-flash
# OPENROUTER_HTTP_REFERER=https://example.com # OPENROUTER_HTTP_REFERER=https://example.com
# OPENROUTER_APP_TITLE=Historical Transcription MVP # OPENROUTER_APP_TITLE="Google: Gemini 2.5 Flash (openrouter)"
# DATABASE_URL=sqlite:///./transcription.db
# UPLOAD_DIR=./uploads # --- runtime environment ---
# PROMPT_DIR=./prompts # ENVIRONMENT: [`development`(default), `test`, `production`]
# --- persistence ---
# Use nested settings with double underscore because env_nested_delimiter="__".
# SQLite example:
# DATABASE__DRIVER=sqlite
# DATABASE__PATH=app.db
#
# SQLite with custom relative path:
# DATABASE__DRIVER=sqlite
DATABASE__PATH=./data/transcription.db
#
# Postgres example:
# DATABASE__DRIVER=postgres
# DATABASE__HOST=localhost
# DATABASE__PORT=5432
# DATABASE__DATABASE=transcription
# DATABASE__USER=postgres
# DATABASE__PASSWORD=change-me
#
# Optional persistence flags:
# BOOTSTRAP_SCHEMA_ON_STARTUP=false
# SQLITE_CHECK_SAME_THREAD=false
# --- filesystem paths ---
UPLOAD_DIR="./data"
PROMPT_DIR="./prompts"
# --- worker reliability ---
WORKER_MAX_RETRIES=0
WORKER_RETRY_BACKOFF_SECONDS=0
# WORKER_PROVIDER_TIMEOUT_SECONDS=[0-20]
WORKER_PROVIDER_TIMEOUT_SECONDS=20
WORKER_MIN_TRANSCRIPTION_CHARS=0
WORKER_MIN_TRANSCRIPTION_LINES=0
WORKER_FAIL_ON_FINISH_REASON_LENGTH=false
+1
View File
@@ -17,3 +17,4 @@ wheels/
# Document images # Document images
uploads/* uploads/*
data/*
+2 -1
View File
@@ -50,7 +50,7 @@ erDiagram
JOB { JOB {
UUID id PK UUID id PK
UUID document_id FK UUID document_id FK
VARCHAR status "queued | processing | completed | partial_success | failed" VARCHAR status "queued | processing | transcribed | completed | partial_success | failed"
INTEGER retry_count INTEGER retry_count
TEXT provider TEXT provider
TEXT model TEXT model
@@ -97,6 +97,7 @@ erDiagram
### Page-Level Execution & AI Outputs ### Page-Level Execution & AI Outputs
* Execution Granularity: Every single image execution by an AI model produces a dedicated record in job_source. * Execution Granularity: Every single image execution by an AI model produces a dedicated record in job_source.
* Source vs Execution Status: `source` does not carry a `status` column. Per-source execution state is tracked in `job_source.status` (`pending`, `transcribed`, `failed`).
* Point-in-Time Auditability: job_source.raw_api_response stores the unparsed REST response envelope for that specific image page call. job_source.ai_metadata stores spatial bounding boxes, token usage, and layout details for that specific image page call. * Point-in-Time Auditability: job_source.raw_api_response stores the unparsed REST response envelope for that specific image page call. job_source.ai_metadata stores spatial bounding boxes, token usage, and layout details for that specific image page call.
* Active Output Caching: Upon successful completion of an image call, source.raw_transcription is updated with the latest output string from job_source.raw_transcription for fast UI rendering. * Active Output Caching: Upon successful completion of an image call, source.raw_transcription is updated with the latest output string from job_source.raw_transcription for fast UI rendering.
+2 -2
View File
@@ -86,11 +86,11 @@ def create_app(settings: Settings | None = None) -> FastAPI:
@app.get("/", include_in_schema=False) @app.get("/", include_in_schema=False)
async def root_redirect() -> RedirectResponse: async def root_redirect() -> RedirectResponse:
return RedirectResponse(url="/ui", status_code=status.HTTP_307_TEMPORARY_REDIRECT) return RedirectResponse(url="/ui/homepage", status_code=status.HTTP_307_TEMPORARY_REDIRECT)
@app.get("/ui", include_in_schema=False) @app.get("/ui", include_in_schema=False)
async def ui_redirect() -> RedirectResponse: async def ui_redirect() -> RedirectResponse:
return RedirectResponse(url="/ui/documents", status_code=status.HTTP_307_TEMPORARY_REDIRECT) return RedirectResponse(url="/ui/homepage", status_code=status.HTTP_307_TEMPORARY_REDIRECT)
@app.get("/healthz") @app.get("/healthz")
def health() -> dict[str, str]: def health() -> dict[str, str]:
+18 -6
View File
@@ -4,6 +4,7 @@ from dataclasses import dataclass
from datetime import UTC from datetime import UTC
from datetime import datetime from datetime import datetime
from pathlib import Path from pathlib import Path
import shutil
from uuid import UUID from uuid import UUID
from sqlalchemy.exc import IntegrityError from sqlalchemy.exc import IntegrityError
@@ -119,6 +120,7 @@ class DocumentService(ServiceBase):
async def delete_document(self, document: Document, *, session: AsyncSession | None = None) -> None: async def delete_document(self, document: Document, *, session: AsyncSession | None = None) -> None:
"""Delete a document from the database.""" """Delete a document from the database."""
document_id = document.id
async with self._session_scope(session) as _session: async with self._session_scope(session) as _session:
existing = await _session.get( existing = await _session.get(
Document, Document,
@@ -152,6 +154,20 @@ class DocumentService(ServiceBase):
await _session.delete(existing) await _session.delete(existing)
await self._finalize(session=_session, caller_session=session) await self._finalize(session=_session, caller_session=session)
self._delete_document_storage_folder(document_id=document_id)
def _delete_document_storage_folder(self, *, document_id: UUID) -> None:
"""Best-effort cleanup for document-scoped source storage."""
document_dir = self.settings.upload_dir / "documents" / str(document_id)
if not document_dir.exists():
return
try:
shutil.rmtree(document_dir)
logger.info("Deleted document storage folder: %s", document_dir)
except OSError:
logger.warning("Failed to delete document storage folder: %s", document_dir)
async def create_person(self, person: Person, *, session: AsyncSession | None = None) -> Person: async def create_person(self, person: Person, *, session: AsyncSession | None = None) -> Person:
"""Create a new person in the database.""" """Create a new person in the database."""
async with self._session_scope(session) as _session: async with self._session_scope(session) as _session:
@@ -217,12 +233,8 @@ class DocumentService(ServiceBase):
suggestion="Verify the person id and retry.", suggestion="Verify the person id and retry.",
) )
if existing.document_people: for link in list(existing.document_people):
raise PersonDeleteBlockedError( await _session.delete(link)
"Person delete blocked by linked documents",
category=ErrorCategory.VALIDATION,
suggestion="Remove linked DocumentPerson records first, then retry deletion.",
)
await _session.delete(existing) await _session.delete(existing)
await self._finalize(session=_session, caller_session=session) await self._finalize(session=_session, caller_session=session)
+93
View File
@@ -11,6 +11,7 @@ from ..errors import AppError
from ..errors import ErrorCategory from ..errors import ErrorCategory
from ..db.models import Job from ..db.models import Job
from ..db.models import JobSource from ..db.models import JobSource
from ..db.models import JobSourceStatus
from ..db.models import JobStatus from ..db.models import JobStatus
from ..db.models import Source from ..db.models import Source
from .base import ServiceBase from .base import ServiceBase
@@ -20,6 +21,14 @@ class JobDeleteBlockedError(AppError):
"""Raised when a job delete operation is blocked by lifecycle policy.""" """Raised when a job delete operation is blocked by lifecycle policy."""
class JobCancelBlockedError(AppError):
"""Raised when a job cancel operation is blocked by lifecycle policy."""
class JobResubmitBlockedError(AppError):
"""Raised when a job resubmit operation is blocked by lifecycle policy."""
class JobService(ServiceBase): class JobService(ServiceBase):
"""Thin service class for managing jobs in the database.""" """Thin service class for managing jobs in the database."""
@@ -221,3 +230,87 @@ class JobService(ServiceBase):
await _session.delete(job) await _session.delete(job)
await self._finalize(session=_session, caller_session=session) await self._finalize(session=_session, caller_session=session)
async def cancel_job(self, *, job_id: UUID, session: AsyncSession | None = None) -> Job:
"""Cancel a queued/processing job and stop remaining source work."""
async with self._session_scope(session) as _session:
query = (
select(Job)
.options(
selectinload(Job.job_sources).selectinload(JobSource.source), # pyright: ignore[reportArgumentType]
)
.where(Job.id == job_id)
.execution_options(populate_existing=True)
)
job = (await _session.exec(query)).first()
if job is None:
raise ValueError(f"Job with id {job_id} not found")
if job.status in {JobStatus.TRANSCRIBED, JobStatus.COMPLETED}:
raise JobCancelBlockedError(
"Job cancel is not allowed for transcribed/completed jobs",
category=ErrorCategory.VALIDATION,
suggestion="Use resubmit for reprocessing needs, or leave the terminal job unchanged.",
)
now = datetime.now(UTC)
job.status = JobStatus.FAILED
job.date_updated = now
for job_source in job.job_sources:
if job_source.status == JobSourceStatus.TRANSCRIBED:
continue
job_source.status = JobSourceStatus.FAILED
job_source.raw_transcription = None
job_source.error_detail = "Cancelled by user"
job_source.executed_at = now
if job_source.source is not None:
job_source.source.raw_transcription = None
await self._finalize(session=_session, caller_session=session, refresh=(job,))
return job
async def resubmit_non_transcribed_sources(self, *, job_id: UUID, session: AsyncSession | None = None) -> int:
"""Reset non-transcribed source executions and queue the job for reprocessing."""
async with self._session_scope(session) as _session:
query = (
select(Job)
.options(
selectinload(Job.job_sources).selectinload(JobSource.source), # pyright: ignore[reportArgumentType]
)
.where(Job.id == job_id)
.execution_options(populate_existing=True)
)
job = (await _session.exec(query)).first()
if job is None:
raise ValueError(f"Job with id {job_id} not found")
if job.status == JobStatus.PROCESSING:
raise JobResubmitBlockedError(
"Job resubmit is blocked while processing is active",
category=ErrorCategory.VALIDATION,
suggestion="Cancel processing first, then resubmit remaining sources.",
)
candidates = [job_source for job_source in job.job_sources if job_source.status != JobSourceStatus.TRANSCRIBED]
if not candidates:
raise JobResubmitBlockedError(
"Job has no non-transcribed sources to resubmit",
category=ErrorCategory.VALIDATION,
suggestion="Only failed or pending sources can be resubmitted.",
)
now = datetime.now(UTC)
for job_source in candidates:
job_source.status = JobSourceStatus.PENDING
job_source.raw_transcription = None
job_source.error_detail = None
job_source.executed_at = now
if job_source.source is not None:
job_source.source.raw_transcription = None
job.status = JobStatus.QUEUED
job.date_updated = now
await self._finalize(session=_session, caller_session=session, refresh=(job,))
return len(candidates)
+64 -19
View File
@@ -41,6 +41,15 @@ class JobCreateResult:
source_ids: tuple[UUID, ...] source_ids: tuple[UUID, ...]
@dataclass(frozen=True)
class PendingStoredUpload:
"""Pre-staged upload artifact tied to a source id."""
source_id: UUID
original_filename: str
stored_path: Path
async def create_upload_job( async def create_upload_job(
*, *,
filename: str, filename: str,
@@ -50,14 +59,20 @@ async def create_upload_job(
) -> UploadJobResult: ) -> UploadJobResult:
"""Create upload-backed document and queued job records.""" """Create upload-backed document and queued job records."""
runtime_settings = settings or get_settings() runtime_settings = settings or get_settings()
document_id = uuid4()
source_id = uuid4()
stored_path = store_file( stored_path = store_file(
filename=filename, filename=filename,
file_bytes=file_bytes, file_bytes=file_bytes,
settings=runtime_settings, settings=runtime_settings,
relative_directory=Path("documents") / str(document_id),
filename_stem=str(source_id),
) )
try: try:
document, job = await _create_upload_records( document, job = await _create_upload_records(
session=session, session=session,
document_id=document_id,
source_id=source_id,
original_filename=filename, original_filename=filename,
stored_path=stored_path, stored_path=stored_path,
) )
@@ -99,15 +114,19 @@ async def create_job_for_document(
runtime_settings = settings or get_settings() runtime_settings = settings or get_settings()
sorted_uploads = sorted(uploads, key=lambda item: Path(item[0]).name.casefold()) sorted_uploads = sorted(uploads, key=lambda item: Path(item[0]).name.casefold())
stored_uploads: list[tuple[str, Path]] = [] stored_uploads: list[PendingStoredUpload] = []
for filename, file_bytes in sorted_uploads: for filename, file_bytes in sorted_uploads:
source_id = uuid4()
stored_uploads.append( stored_uploads.append(
( PendingStoredUpload(
filename, source_id=source_id,
store_file( original_filename=filename,
stored_path=store_file(
filename=filename, filename=filename,
file_bytes=file_bytes, file_bytes=file_bytes,
settings=runtime_settings, settings=runtime_settings,
relative_directory=Path("documents") / str(document_id),
filename_stem=str(source_id),
), ),
) )
) )
@@ -122,8 +141,8 @@ async def create_job_for_document(
prompt_name=prompt_name, prompt_name=prompt_name,
) )
except Exception as exc: except Exception as exc:
for _, stored_path in stored_uploads: for upload in stored_uploads:
_best_effort_delete(stored_path) _best_effort_delete(upload.stored_path)
raise UploadError( raise UploadError(
"Failed to create job records from uploads", "Failed to create job records from uploads",
category=ErrorCategory.INFRA_TRANSIENT, category=ErrorCategory.INFRA_TRANSIENT,
@@ -142,10 +161,13 @@ async def create_job_for_document(
async def _create_upload_records( async def _create_upload_records(
*, *,
session: AsyncSession, session: AsyncSession,
document_id: UUID,
source_id: UUID,
original_filename: str, original_filename: str,
stored_path: Path, stored_path: Path,
) -> tuple[Document, Job]: ) -> tuple[Document, Job]:
document = Document( document = Document(
id=document_id,
name=Path(original_filename).name, name=Path(original_filename).name,
) )
session.add(document) session.add(document)
@@ -156,6 +178,7 @@ async def _create_upload_records(
await session.flush() await session.flush()
source = Source( source = Source(
id=source_id,
document_id=document.id, document_id=document.id,
page_number=1, page_number=1,
upload_name=Path(original_filename).name, upload_name=Path(original_filename).name,
@@ -183,7 +206,7 @@ async def _create_job_for_document_records(
*, *,
session: AsyncSession, session: AsyncSession,
document_id: UUID, document_id: UUID,
stored_uploads: Sequence[tuple[str, Path]], stored_uploads: Sequence[PendingStoredUpload],
provider: str | None, provider: str | None,
model: str | None, model: str | None,
prompt_name: str | None, prompt_name: str | None,
@@ -211,13 +234,14 @@ async def _create_job_for_document_records(
await session.flush() await session.flush()
source_ids: list[UUID] = [] source_ids: list[UUID] = []
for page_offset, (original_filename, stored_path) in enumerate(stored_uploads): for page_offset, upload in enumerate(stored_uploads):
source = Source( source = Source(
id=upload.source_id,
document_id=document_id, document_id=document_id,
page_number=next_page_number + page_offset, page_number=next_page_number + page_offset,
upload_name=Path(original_filename).name, upload_name=Path(upload.original_filename).name,
filename=stored_path.name, filename=upload.stored_path.name,
file_path=str(stored_path), file_path=str(upload.stored_path),
) )
session.add(source) session.add(source)
await session.flush() await session.flush()
@@ -244,22 +268,41 @@ def _best_effort_delete(path: Path) -> None:
logger.warning("Failed to clean up upload file after DB error: %s", path) logger.warning("Failed to clean up upload file after DB error: %s", path)
def store_file(*, filename: str, file_bytes: bytes, settings: Settings | None = None) -> Path: def store_file(
*,
filename: str,
file_bytes: bytes,
settings: Settings | None = None,
relative_directory: Path | None = None,
filename_stem: str | None = None,
) -> Path:
"""Persist an uploaded file to the configured upload directory.""" """Persist an uploaded file to the configured upload directory."""
runtime_settings = settings or get_settings() runtime_settings = settings or get_settings()
_validate_upload(filename=filename, file_bytes=file_bytes, supported_extensions=SUPPORTED_UPLOAD_EXTENSIONS) _validate_upload(filename=filename, file_bytes=file_bytes, supported_extensions=SUPPORTED_UPLOAD_EXTENSIONS)
return _store_file_bytes(filename=filename, file_bytes=file_bytes, settings=runtime_settings) return _store_file_bytes(
filename=filename,
file_bytes=file_bytes,
settings=runtime_settings,
relative_directory=relative_directory,
filename_stem=filename_stem,
)
def store_person_portrait(*, filename: str, file_bytes: bytes, settings: Settings | None = None) -> Path: def store_person_portrait(
"""Persist a portrait upload under uploads/portraits/person.""" *,
person_id: UUID,
filename: str,
file_bytes: bytes,
settings: Settings | None = None,
) -> Path:
"""Persist a portrait upload under persons/<person_id>."""
runtime_settings = settings or get_settings() runtime_settings = settings or get_settings()
_validate_upload(filename=filename, file_bytes=file_bytes, supported_extensions=SUPPORTED_PORTRAIT_EXTENSIONS) _validate_upload(filename=filename, file_bytes=file_bytes, supported_extensions=SUPPORTED_PORTRAIT_EXTENSIONS)
return _store_file_bytes( return _store_file_bytes(
filename=filename, filename=filename,
file_bytes=file_bytes, file_bytes=file_bytes,
settings=runtime_settings, settings=runtime_settings,
relative_directory=Path("portraits") / "person", relative_directory=Path("persons") / str(person_id),
) )
@@ -269,12 +312,13 @@ def _store_file_bytes(
file_bytes: bytes, file_bytes: bytes,
settings: Settings, settings: Settings,
relative_directory: Path | None = None, relative_directory: Path | None = None,
filename_stem: str | None = None,
) -> Path: ) -> Path:
upload_dir = settings.upload_dir upload_dir = settings.upload_dir
target_dir = upload_dir if relative_directory is None else upload_dir / relative_directory target_dir = upload_dir if relative_directory is None else upload_dir / relative_directory
target_dir.mkdir(parents=True, exist_ok=True) target_dir.mkdir(parents=True, exist_ok=True)
stored_name = _build_stored_filename(filename) stored_name = _build_stored_filename(filename=filename, filename_stem=filename_stem)
stored_path = target_dir / stored_name stored_path = target_dir / stored_name
try: try:
@@ -315,7 +359,8 @@ def _validate_upload(*, filename: str, file_bytes: bytes, supported_extensions:
) )
def _build_stored_filename(filename: str) -> str: def _build_stored_filename(*, filename: str, filename_stem: str | None = None) -> str:
safe_name = Path(filename).name safe_name = Path(filename).name
suffix = Path(safe_name).suffix.lower() suffix = Path(safe_name).suffix.lower()
return f"{uuid4()}{suffix}" stem = filename_stem or str(uuid4())
return f"{stem}{suffix}"
+118 -6
View File
@@ -113,10 +113,43 @@ class TranscriptionService(ServiceBase):
async def delete_source(self, source: Source, *, session: AsyncSession | None = None) -> None: async def delete_source(self, source: Source, *, session: AsyncSession | None = None) -> None:
"""Delete a source page record.""" """Delete a source page record."""
source_file_path = source.file_path
async with self._session_scope(session) as _session: async with self._session_scope(session) as _session:
await _session.delete(source) await _session.delete(source)
await self._finalize(session=_session, caller_session=session) await self._finalize(session=_session, caller_session=session)
self._delete_source_file(source_file_path=source_file_path)
async def delete_unlinked_source(self, *, source_id: UUID, session: AsyncSession | None = None) -> None:
"""Delete a source only when no JobSource links exist."""
async with self._session_scope(session) as _session:
source = await _session.get(
Source,
source_id,
options=(
selectinload(Source.job_sources), # pyright: ignore[reportArgumentType]
),
)
if source is None:
raise TranscriptionNotFoundError(
f"Source with id {source_id} not found",
category=ErrorCategory.NOT_FOUND,
suggestion="Verify the source id and retry.",
)
if source.job_sources:
raise SourceDeleteBlockedError(
"Source delete blocked because it is linked to one or more jobs",
category=ErrorCategory.VALIDATION,
suggestion="Remove JobSource links first, then retry deletion.",
)
source_file_path = source.file_path
await _session.delete(source)
await self._finalize(session=_session, caller_session=session)
self._delete_source_file(source_file_path=source_file_path)
async def list_sources( async def list_sources(
self, self,
*, *,
@@ -236,9 +269,26 @@ class TranscriptionService(ServiceBase):
for job_source in matching_links: for job_source in matching_links:
await _session.delete(job_source) await _session.delete(job_source)
source_file_path = source.file_path
await _session.delete(source) await _session.delete(source)
await self._finalize(session=_session, caller_session=session) await self._finalize(session=_session, caller_session=session)
self._delete_source_file(source_file_path=source_file_path)
def _delete_source_file(self, *, source_file_path: str) -> None:
"""Best-effort cleanup for source media files."""
candidate_path = Path(source_file_path)
resolved_path = candidate_path if candidate_path.is_absolute() else self.settings.upload_dir / candidate_path
if not resolved_path.exists():
return
try:
resolved_path.unlink()
logger.info("Deleted source file: %s", resolved_path)
except OSError:
logger.warning("Failed to delete source file: %s", resolved_path)
async def list_job_sources( async def list_job_sources(
self, self,
*, *,
@@ -292,7 +342,11 @@ class TranscriptionService(ServiceBase):
prompt_name: str = DEFAULT_PROMPT_FILE, prompt_name: str = DEFAULT_PROMPT_FILE,
session: AsyncSession | None = None, session: AsyncSession | None = None,
) -> Job: ) -> Job:
"""Persist original transcription output fields on a job.""" """Persist transcription output for the first ordered source in a job's document.
This compatibility helper keeps legacy single-source workflows working.
New multi-source flows should use ``update_job_source_transcription``.
"""
async with self._session_scope(session) as _session: async with self._session_scope(session) as _session:
job = await _session.get(Job, job_id) job = await _session.get(Job, job_id)
if job is None: if job is None:
@@ -314,14 +368,72 @@ class TranscriptionService(ServiceBase):
) )
source_row = source.first() source_row = source.first()
if source_row is not None: if source_row is not None:
await self.update_job_source_transcription(
job_id=job.id,
source_id=source_row.id,
text=text,
error_detail=error_detail,
provider=provider,
model=model,
prompt_name=prompt_name,
session=_session,
)
await self._finalize(session=_session, caller_session=session, refresh=(job,))
return job
async def update_job_source_transcription(
self,
*,
job_id: UUID,
source_id: UUID,
text: str | None,
error_detail: str | None = None,
provider: str | None = None,
model: str | None = None,
prompt_name: str = DEFAULT_PROMPT_FILE,
session: AsyncSession | None = None,
) -> JobSource:
"""Persist transcription fields for one source within a specific job."""
async with self._session_scope(session) as _session:
job = await _session.get(Job, job_id)
if job is None:
raise TranscriptionNotFoundError(
f"Job with id {job_id} not found",
category=ErrorCategory.NOT_FOUND,
suggestion="Verify the job id and retry.",
)
source = await _session.get(Source, source_id)
if source is None:
raise TranscriptionNotFoundError(
f"Source with id {source_id} not found",
category=ErrorCategory.NOT_FOUND,
suggestion="Verify the source id and retry.",
)
if source.document_id != job.document_id:
raise TranscriptionError(
f"Source {source_id} does not belong to job {job_id}",
category=ErrorCategory.VALIDATION,
suggestion="Link the source to the same document as the job and retry.",
)
job.provider = provider or job.provider or self.settings.provider.value
job.model = model or job.model or _resolve_transcript_model(provider=self.provider, settings=self.settings)
job.prompt_name = prompt_name or job.prompt_name or DEFAULT_PROMPT_FILE
job.date_updated = datetime.now(UTC)
source.raw_transcription = text
existing_job_source = await _session.exec( existing_job_source = await _session.exec(
select(JobSource).where(JobSource.job_id == job.id).where(JobSource.source_id == source_row.id) select(JobSource).where(JobSource.job_id == job_id).where(JobSource.source_id == source_id)
) )
job_source = existing_job_source.first() job_source = existing_job_source.first()
if job_source is None: if job_source is None:
job_source = JobSource( job_source = JobSource(
job_id=job.id, job_id=job_id,
source_id=source_row.id, source_id=source_id,
status=JobSourceStatus.TRANSCRIBED if text is not None else JobSourceStatus.FAILED, status=JobSourceStatus.TRANSCRIBED if text is not None else JobSourceStatus.FAILED,
raw_transcription=text, raw_transcription=text,
error_detail=error_detail, error_detail=error_detail,
@@ -333,8 +445,8 @@ class TranscriptionService(ServiceBase):
job_source.status = JobSourceStatus.TRANSCRIBED if text is not None else JobSourceStatus.FAILED job_source.status = JobSourceStatus.TRANSCRIBED if text is not None else JobSourceStatus.FAILED
job_source.executed_at = datetime.now(UTC) job_source.executed_at = datetime.now(UTC)
await self._finalize(session=_session, caller_session=session, refresh=(job,)) await self._finalize(session=_session, caller_session=session, refresh=(job, source, job_source))
return job return job_source
async def upsert_revision_for_source( async def upsert_revision_for_source(
self, self,
+141 -25
View File
@@ -6,6 +6,7 @@ from sqlmodel.ext.asyncio.session import AsyncSession
from ..config import Settings from ..config import Settings
from ..config import get_settings from ..config import get_settings
from ..db.models import Job from ..db.models import Job
from ..db.models import JobSourceStatus
from ..db.models import JobStatus from ..db.models import JobStatus
from ..db.models import Source from ..db.models import Source
from ..errors import AppError from ..errors import AppError
@@ -69,20 +70,24 @@ async def process_queued_job(
await session.commit() await session.commit()
source_job = await services.jobs.read_job(job_id=job.id, session=session) source_job = await services.jobs.read_job(job_id=job.id, session=session)
source = _resolve_primary_source(source_job) sources = _resolve_job_sources(source_job)
if source is None: if not sources and not source_job.job_sources:
candidate_sources = await services.transcriptions.list_sources(document_id=job.document_id, session=session) candidate_sources = await services.transcriptions.list_sources(document_id=job.document_id, session=session)
source = next(iter(sorted(candidate_sources, key=lambda item: item.page_number)), None) sources = list(sorted(candidate_sources, key=lambda item: (item.page_number, item.upload_name.casefold())))
if not sources:
return await services.jobs.mark_job_status(job.id, JobStatus.TRANSCRIBED, session=session)
successful_pages: list[tuple[Source, TranscriptionResult]] = []
failed_pages: list[tuple[Source, AppError]] = []
externally_stopped = False
for source in sources:
if await _job_no_longer_processing(job_id=job.id, services=services, session=session):
externally_stopped = True
break
if source is None:
error = AppError(
f"Job {job.id} has no associated source record.",
category=ErrorCategory.VALIDATION,
suggestion="Attach at least one source to the job and retry.",
)
return await _finalize_failed(job=job, services=services, error=error, session=session)
started_at = asyncio.get_running_loop().time() started_at = asyncio.get_running_loop().time()
try: try:
result = await asyncio.wait_for( result = await asyncio.wait_for(
transcribe_document_image(source.file_path), transcribe_document_image(source.file_path),
@@ -109,15 +114,7 @@ async def process_queued_job(
) )
_validate_transcription_quality(result=result, settings=runtime_settings) _validate_transcription_quality(result=result, settings=runtime_settings)
successful_pages.append((source, result))
job = await _finalize_transcribed(job=job, services=services, result=result, session=session)
logger.info(
"Job transcribed operation=worker.process_job job_id=%s document_id=%s source_id=%s provider=%s",
job.id,
job.document_id,
source.id,
result.provider,
)
except TimeoutError: except TimeoutError:
error = AppError( error = AppError(
f"Provider call timed out after {runtime_settings.worker_provider_timeout_seconds:.1f}s", f"Provider call timed out after {runtime_settings.worker_provider_timeout_seconds:.1f}s",
@@ -125,9 +122,9 @@ async def process_queued_job(
suggestion="Retry the job. If this repeats, verify provider latency and request payload size.", suggestion="Retry the job. If this repeats, verify provider latency and request payload size.",
retriable=True, retriable=True,
) )
job = await _finalize_failed(job=job, services=services, error=error, session=session) failed_pages.append((source, error))
logger.error( logger.error(
"Job failed operation=worker.process_job job_id=%s document_id=%s source_id=%s error_id=%s category=%s", "Source failed operation=worker.process_job job_id=%s document_id=%s source_id=%s error_id=%s category=%s",
job.id, job.id,
job.document_id, job.document_id,
source.id, source.id,
@@ -141,16 +138,46 @@ async def process_queued_job(
case _: case _:
error = classify_unexpected_error(exc, operation="worker.process_job") error = classify_unexpected_error(exc, operation="worker.process_job")
job = await _finalize_failed(job=job, services=services, error=error, session=session) failed_pages.append((source, error))
logger.error( logger.error(
"Job failed operation=worker.process_job job_id=%s document_id=%s source_id=%s error_id=%s category=%s", "Source failed operation=worker.process_job job_id=%s document_id=%s source_id=%s error_id=%s category=%s",
job.id, job.id,
job.document_id, job.document_id,
source.id, source.id,
error.error_id, error.error_id,
error.category.value, error.category.value,
) )
return job
if await _job_no_longer_processing(job_id=job.id, services=services, session=session):
externally_stopped = True
break
terminal_status = JobStatus.TRANSCRIBED
if externally_stopped:
terminal_status = JobStatus.FAILED
elif failed_pages and successful_pages:
terminal_status = JobStatus.PARTIAL_SUCCESS
elif failed_pages and not successful_pages:
terminal_status = JobStatus.FAILED
updated_job = await _finalize_batch_outcome(
job=job,
services=services,
successful_pages=successful_pages,
failed_pages=failed_pages,
status=terminal_status,
session=session,
)
logger.info(
"Job finished operation=worker.process_job job_id=%s document_id=%s status=%s success_pages=%s failed_pages=%s",
updated_job.id,
updated_job.document_id,
updated_job.status.value,
len(successful_pages),
len(failed_pages),
)
return updated_job
async def process_next_queued_job( async def process_next_queued_job(
@@ -307,6 +334,95 @@ def _resolve_primary_source(job: Job) -> Source | None:
return next((job_source.source for job_source in job.job_sources if job_source.source is not None), None) return next((job_source.source for job_source in job.job_sources if job_source.source is not None), None)
def _resolve_job_sources(job: Job) -> list[Source]:
"""Resolve non-transcribed linked sources for a job in deterministic page order."""
if not job.job_sources:
return []
sources = [
job_source.source
for job_source in job.job_sources
if job_source.source is not None and job_source.status != JobSourceStatus.TRANSCRIBED
]
return list(sorted(sources, key=lambda item: (item.page_number, item.upload_name.casefold())))
async def _job_no_longer_processing(
*,
job_id,
services: ServiceBundle,
session: AsyncSession | None = None,
) -> bool:
"""Return True when job status changed externally from PROCESSING."""
latest_job = await services.jobs.read_job(job_id=job_id, session=session)
return latest_job.status != JobStatus.PROCESSING
async def _finalize_batch_outcome(
*,
job: Job,
services: ServiceBundle,
successful_pages: list[tuple[Source, TranscriptionResult]],
failed_pages: list[tuple[Source, AppError]],
status: JobStatus,
session: AsyncSession | None = None,
) -> Job:
"""Transaction B: write per-source outcomes and terminal job status atomically."""
if session is None:
async with services.jobs._session_scope() as local_session:
for source, result in successful_pages:
await services.transcriptions.update_job_source_transcription(
job_id=job.id,
source_id=source.id,
text=result.text,
error_detail=None,
provider=result.provider,
model=result.model,
prompt_name=result.prompt_name,
session=local_session,
)
for source, error in failed_pages:
await services.transcriptions.update_job_source_transcription(
job_id=job.id,
source_id=source.id,
text=None,
error_detail=format_error_detail(error),
prompt_name=DEFAULT_PROMPT_FILE,
session=local_session,
)
updated_job = await services.jobs.mark_job_status(job.id, status, session=local_session)
await local_session.commit()
return updated_job
for source, result in successful_pages:
await services.transcriptions.update_job_source_transcription(
job_id=job.id,
source_id=source.id,
text=result.text,
error_detail=None,
provider=result.provider,
model=result.model,
prompt_name=result.prompt_name,
session=session,
)
for source, error in failed_pages:
await services.transcriptions.update_job_source_transcription(
job_id=job.id,
source_id=source.id,
text=None,
error_detail=format_error_detail(error),
prompt_name=DEFAULT_PROMPT_FILE,
session=session,
)
updated_job = await services.jobs.mark_job_status(job.id, status, session=session)
await session.commit()
return updated_job
def _validate_transcription_quality(*, result: TranscriptionResult, settings: Settings) -> None: def _validate_transcription_quality(*, result: TranscriptionResult, settings: Settings) -> None:
text_chars = len(result.text) text_chars = len(result.text)
text_lines = _line_count(result.text) text_lines = _line_count(result.text)
+2
View File
@@ -3,6 +3,7 @@
from fastapi import FastAPI from fastapi import FastAPI
from nicegui import ui from nicegui import ui
from transcription.ui.pages.home_page import register_page as register_home_page
from transcription.ui.pages.documents_page import register_page as register_documents_page from transcription.ui.pages.documents_page import register_page as register_documents_page
from transcription.ui.pages.jobs_page import register_page as register_jobs_page from transcription.ui.pages.jobs_page import register_page as register_jobs_page
from transcription.ui.pages.people_page import register_page as register_people_page from transcription.ui.pages.people_page import register_page as register_people_page
@@ -25,6 +26,7 @@ def _register_global_styles(app: FastAPI) -> None:
def register_pages(app: FastAPI) -> None: def register_pages(app: FastAPI) -> None:
"""Register all NiceGUI pages and mount them onto the FastAPI app.""" """Register all NiceGUI pages and mount them onto the FastAPI app."""
_register_global_styles(app) _register_global_styles(app)
register_home_page()
register_upload_page() register_upload_page()
register_documents_page() register_documents_page()
register_people_page() register_people_page()
+4 -2
View File
@@ -43,7 +43,7 @@ def _render_nav_button(*, label: str, path: str, icon: str, current_path: str) -
def _normalize_path(current_path: str | None) -> str: def _normalize_path(current_path: str | None) -> str:
normalized = (current_path or "").strip() normalized = (current_path or "").strip()
if not normalized: if not normalized:
return "/jobs" return "/homepage"
return normalized.rstrip("/") or "/" return normalized.rstrip("/") or "/"
@@ -53,7 +53,9 @@ def render_app_shell(*, current_path: str | None = None) -> None:
normalized_path = _normalize_path(current_path) normalized_path = _normalize_path(current_path)
with ui.header().classes("app-shell"), ui.element("div").classes("app-shell__inner"): with ui.header().classes("app-shell"), ui.element("div").classes("app-shell__inner"):
with ui.row().classes("app-shell__brand no-wrap"): with ui.element("a").props('href="/ui/homepage"').style(
"display:flex; align-items:center; gap:0.75rem; text-decoration:none; color:inherit;"
).classes("app-shell__brand no-wrap"):
ui.html(VIBESCRIBE_LOGO_SVG).classes("app-shell__brand-mark") ui.html(VIBESCRIBE_LOGO_SVG).classes("app-shell__brand-mark")
ui.label("VibeScribe").classes("app-shell__brand-name") ui.label("VibeScribe").classes("app-shell__brand-name")
@@ -23,6 +23,8 @@ class SourceTableRow:
upload_name: str upload_name: str
filename: str filename: str
document_id: UUID document_id: UUID
job_source_status: str | None = None
job_source_error_detail: str | None = None
def _serialize_rows(rows: Sequence[SourceTableRow]) -> list[dict[str, Any]]: def _serialize_rows(rows: Sequence[SourceTableRow]) -> list[dict[str, Any]]:
@@ -33,6 +35,8 @@ def _serialize_rows(rows: Sequence[SourceTableRow]) -> list[dict[str, Any]]:
"upload_name": row.upload_name, "upload_name": row.upload_name,
"filename": row.filename, "filename": row.filename,
"document_id": str(row.document_id), "document_id": str(row.document_id),
"job_source_status": row.job_source_status or "-",
"job_source_error_detail": row.job_source_error_detail or "-",
} }
for row in rows for row in rows
] ]
@@ -51,6 +55,20 @@ def render_sources_table(rows: Sequence[SourceTableRow]) -> None:
{"name": "page_number", "label": "Page", "field": "page_number", "sortable": True}, {"name": "page_number", "label": "Page", "field": "page_number", "sortable": True},
{"name": "upload_name", "label": "Upload Title", "field": "upload_name", "sortable": True, "classes": "font-serif"}, {"name": "upload_name", "label": "Upload Title", "field": "upload_name", "sortable": True, "classes": "font-serif"},
{"name": "filename", "label": "Stored Filename", "field": "filename", "sortable": True, "classes": "font-mono"}, {"name": "filename", "label": "Stored Filename", "field": "filename", "sortable": True, "classes": "font-mono"},
{
"name": "job_source_status",
"label": "Job Source Status",
"field": "job_source_status",
"sortable": True,
"classes": "font-mono",
},
{
"name": "job_source_error_detail",
"label": "Job Source Error Detail",
"field": "job_source_error_detail",
"sortable": False,
"classes": "font-mono text-xs",
},
{"name": "document_id", "label": "Document ID", "field": "document_id", "sortable": True, "classes": "font-mono"}, {"name": "document_id", "label": "Document ID", "field": "document_id", "sortable": True, "classes": "font-mono"},
], ],
default_sort_by="page_number", default_sort_by="page_number",
+62
View File
@@ -0,0 +1,62 @@
"""File-backed storage helpers for the homepage content."""
from __future__ import annotations
from pathlib import Path
HOME_PAGE_DIR = Path(__file__).resolve().parents[3] / "data" / "homepage"
HOME_PAGE_MARKDOWN_PATH = HOME_PAGE_DIR / "homepage.md"
SUPPORTED_IMAGE_SUFFIXES = {".jpg", ".jpeg", ".png", ".gif", ".webp", ".bmp", ".tif", ".tiff"}
def ensure_homepage_storage() -> None:
"""Create the homepage storage directory when needed."""
HOME_PAGE_DIR.mkdir(parents=True, exist_ok=True)
def read_homepage_markdown() -> str:
"""Read the saved homepage markdown text."""
ensure_homepage_storage()
if not HOME_PAGE_MARKDOWN_PATH.exists():
return ""
return HOME_PAGE_MARKDOWN_PATH.read_text(encoding="utf-8")
def save_homepage_markdown(markdown_text: str) -> None:
"""Persist the homepage markdown text."""
ensure_homepage_storage()
HOME_PAGE_MARKDOWN_PATH.write_text(markdown_text, encoding="utf-8")
def store_homepage_image(*, filename: str, file_bytes: bytes) -> Path:
"""Persist an uploaded homepage image in the shared homepage folder."""
ensure_homepage_storage()
safe_name = Path(filename).name
if not safe_name:
msg = "Homepage image filename is required"
raise ValueError(msg)
stored_path = HOME_PAGE_DIR / safe_name
stored_path.write_bytes(file_bytes)
return stored_path
def list_homepage_images() -> list[Path]:
"""List stored homepage images in the order they were last updated."""
ensure_homepage_storage()
image_paths = [
path
for path in HOME_PAGE_DIR.iterdir()
if path.is_file() and path.suffix.lower() in SUPPORTED_IMAGE_SUFFIXES
]
return sorted(image_paths, key=lambda path: (path.stat().st_mtime, path.name))
def latest_homepage_image() -> Path | None:
"""Return the most recently updated homepage image, if one exists."""
image_paths = list_homepage_images()
if not image_paths:
return None
return image_paths[-1]
+110
View File
@@ -0,0 +1,110 @@
"""Homepage registration and handlers."""
from __future__ import annotations
from nicegui import ui
from transcription.ui.components.app_shell import render_navigation_header
from transcription.ui.components.cards import archival_card
from transcription.ui.components.primitives import render_empty_state
from transcription.ui.components.primitives import section_header_row
from transcription.ui.components.viewers import dark_room_viewer
from transcription.ui.homepage_store import latest_homepage_image
from transcription.ui.homepage_store import read_homepage_markdown
from transcription.ui.homepage_store import save_homepage_markdown
from transcription.ui.homepage_store import store_homepage_image
from transcription.ui.theme import apply_archival_theme
from transcription.ui.theme import page_header
def _render_homepage_view(*, markdown_text: str, image_path) -> None:
with ui.grid().classes("w-full grid-cols-12 gap-4"):
with ui.column().classes("col-span-12 lg:col-span-4"):
dark_room_viewer(str(image_path) if image_path else None, count_label="Homepage Image")
with ui.column().classes("col-span-12 lg:col-span-5 gap-4"), archival_card(title="Home Text"):
if markdown_text:
ui.markdown(markdown_text)
else:
render_empty_state("No homepage text saved yet.")
with ui.column().classes("col-span-12 lg:col-span-3"):
ui.element("div")
def _render_homepage_editor(*, render_image_panel, markdown_input, on_upload) -> None:
with ui.grid().classes("w-full grid-cols-12 gap-4"):
with ui.column().classes("col-span-12 lg:col-span-4 gap-4"):
with archival_card(title="Homepage Image"):
ui.upload(on_upload=on_upload, auto_upload=True, label="Upload image").props(
'accept=".jpg,.jpeg,.png,.gif,.webp,.bmp,.tif,.tiff"'
).classes("w-full")
render_image_panel()
with ui.column().classes("col-span-12 lg:col-span-5 gap-4"), archival_card(title="Home Text"):
markdown_input[0] = ui.textarea(
label="Homepage markdown",
value=read_homepage_markdown(),
).props("outlined autogrow").classes("w-full")
with ui.column().classes("col-span-12 lg:col-span-3"):
ui.element("div")
def register_page() -> None:
"""Register the homepage routes."""
@ui.page("/homepage", title="VibeScribe Home")
def homepage_page() -> None:
apply_archival_theme()
render_navigation_header(current_path="/homepage")
with ui.column().classes("w-full max-w-[1800px] mx-auto p-4 gap-4"):
with section_header_row():
page_header("Home")
ui.button(
"Edit Home Page",
on_click=lambda: ui.navigate.to("/homepage/edit"),
icon="edit",
).classes("ui-btn-primary text-xs")
_render_homepage_view(
markdown_text=read_homepage_markdown().strip(),
image_path=latest_homepage_image(),
)
@ui.page("/homepage/edit", title="Edit Homepage")
def homepage_edit_page() -> None:
apply_archival_theme()
render_navigation_header(current_path="/homepage")
preview_image = [latest_homepage_image()]
markdown_input = [None]
@ui.refreshable
def render_image_panel() -> None:
dark_room_viewer(str(preview_image[0]) if preview_image[0] else None, count_label="Homepage Image")
async def on_upload(event) -> None:
payload = await event.file.read()
preview_image[0] = store_homepage_image(filename=event.file.name, file_bytes=payload)
ui.notify(f"Uploaded {event.file.name}", type="positive")
render_image_panel.refresh()
async def save_homepage() -> None:
save_homepage_markdown((markdown_input[0].value if markdown_input[0] is not None else "") or "")
ui.notify("Homepage saved", type="positive")
ui.navigate.to("/homepage")
with ui.column().classes("w-full max-w-[1800px] mx-auto p-4 gap-4"):
with section_header_row():
page_header("Edit Home Page")
with ui.row().classes("items-center gap-2"):
ui.button("Save", on_click=save_homepage, icon="save").classes("ui-btn-primary text-xs")
ui.button("Cancel", on_click=lambda: ui.navigate.to("/homepage"), icon="close").props("flat")
_render_homepage_editor(
render_image_panel=render_image_panel,
markdown_input=markdown_input,
on_upload=on_upload,
)
+126 -1
View File
@@ -8,10 +8,13 @@ from uuid import UUID
from fastapi import Request from fastapi import Request
from nicegui import ui from nicegui import ui
from transcription.db.models import JobSourceStatus
from transcription.db.models import JobStatus from transcription.db.models import JobStatus
from transcription.db.session import session_scope from transcription.db.session import session_scope
from transcription.services.documents import DocumentService from transcription.services.documents import DocumentService
from transcription.services.jobs import JobDeleteBlockedError from transcription.services.jobs import JobDeleteBlockedError
from transcription.services.jobs import JobCancelBlockedError
from transcription.services.jobs import JobResubmitBlockedError
from transcription.services.jobs import JobService from transcription.services.jobs import JobService
from transcription.services.store import create_job_for_document from transcription.services.store import create_job_for_document
from transcription.ui.components.app_shell import render_navigation_header from transcription.ui.components.app_shell import render_navigation_header
@@ -203,7 +206,7 @@ def register_page() -> None: # noqa: PLR0915
ui.button("Back to Jobs", on_click=lambda: ui.navigate.to("/jobs"), icon="arrow_back").props("flat") ui.button("Back to Jobs", on_click=lambda: ui.navigate.to("/jobs"), icon="arrow_back").props("flat")
@ui.page("/jobs/{job_id}") @ui.page("/jobs/{job_id}")
async def job_detail_page(job_id: str, session_factory: SessionFactoryDep) -> None: async def job_detail_page(job_id: str, request: Request, session_factory: SessionFactoryDep) -> None:
apply_archival_theme() apply_archival_theme()
jobs_service = JobService(session_factory=session_factory) jobs_service = JobService(session_factory=session_factory)
render_navigation_header(current_path="/jobs") render_navigation_header(current_path="/jobs")
@@ -225,6 +228,20 @@ def register_page() -> None: # noqa: PLR0915
page_header(f"Job Record: {job.id}") page_header(f"Job Record: {job.id}")
with ui.row().classes("items-center gap-2"): with ui.row().classes("items-center gap-2"):
archival_badge(job.status.value.upper()) archival_badge(job.status.value.upper())
if job.status in {JobStatus.QUEUED, JobStatus.PROCESSING}:
destructive_button(
"Cancel",
on_click=lambda: ui.navigate.to(f"/jobs/{job.id}/cancel"),
icon="stop_circle",
extra_classes="text-xs",
)
if job.status != JobStatus.TRANSCRIBED:
ui.button("Resubmit", on_click=lambda: ui.navigate.to(f"/jobs/{job.id}/resubmit"), icon="replay").props(
"outlined"
).classes("text-xs")
destructive_button( destructive_button(
"Delete Job", "Delete Job",
on_click=lambda: ui.navigate.to(f"/jobs/{job.id}/delete"), on_click=lambda: ui.navigate.to(f"/jobs/{job.id}/delete"),
@@ -254,6 +271,114 @@ def register_page() -> None: # noqa: PLR0915
icon="description", icon="description",
).props("flat text-xs").classes("ui-link-primary w-full") ).props("flat text-xs").classes("ui-link-primary w-full")
@ui.page("/jobs/{job_id}/cancel")
async def job_cancel_page(job_id: str, request: Request, session_factory: SessionFactoryDep) -> None:
apply_archival_theme()
jobs_service = JobService(session_factory=session_factory)
render_navigation_header(current_path="/jobs")
try:
parsed_job_id = UUID(job_id)
except ValueError:
ui.label("Invalid job id").classes("text-h6 text-red-800 p-4")
return
try:
job = await jobs_service.read_job(job_id=parsed_job_id)
except ValueError:
ui.label("Job not found").classes("text-h6 text-red-800 p-4")
return
with ui.column().classes("w-full max-w-xl mx-auto p-4 gap-4"):
page_header("Cancel Processing Job")
with archival_card(extra_classes="gap-2"):
ui.label(f"Job ID: {job.id}").classes("text-sm font-semibold font-mono ui-text-primary")
metadata_row("Current Status:", job.status.value)
ui.label("Cancel stops processing and marks remaining non-transcribed sources as failed.").classes(
"text-xs ui-text-muted"
)
async def submit_cancel() -> None:
try:
await jobs_service.cancel_job(job_id=job.id)
except JobCancelBlockedError as exc:
ui.notify(exc.message, type="warning")
return
except ValueError:
ui.notify("Job not found.", type="warning")
ui.navigate.to("/jobs")
return
except Exception as exc: # noqa: BLE001
show_error(exc, title="Cancel job failed", operation="jobs.cancel")
return
resolve_worker_notifier(request.app.state).notify()
ui.notify("Job cancelled", type="positive")
ui.navigate.to(f"/jobs/{job.id}")
with ui.row().classes("w-full items-center gap-2 mt-2"):
destructive_button(
"Cancel job",
on_click=submit_cancel,
icon="stop_circle",
variant="solid",
)
ui.button("Back to Job", on_click=lambda: ui.navigate.to(f"/jobs/{job.id}"), icon="arrow_back").props("flat")
@ui.page("/jobs/{job_id}/resubmit")
async def job_resubmit_page(job_id: str, request: Request, session_factory: SessionFactoryDep) -> None:
apply_archival_theme()
jobs_service = JobService(session_factory=session_factory)
render_navigation_header(current_path="/jobs")
try:
parsed_job_id = UUID(job_id)
except ValueError:
ui.label("Invalid job id").classes("text-h6 text-red-800 p-4")
return
try:
job = await jobs_service.read_job(job_id=parsed_job_id)
except ValueError:
ui.label("Job not found").classes("text-h6 text-red-800 p-4")
return
non_transcribed_count = sum(1 for job_source in job.job_sources if job_source.status != JobSourceStatus.TRANSCRIBED)
with ui.column().classes("w-full max-w-xl mx-auto p-4 gap-4"):
page_header("Resubmit Job")
with archival_card(extra_classes="gap-2"):
ui.label(f"Job ID: {job.id}").classes("text-sm font-semibold font-mono ui-text-primary")
metadata_row("Current Status:", job.status.value)
metadata_row("Non-Transcribed Sources:", str(non_transcribed_count))
ui.label("Resubmit queues all non-transcribed linked sources. New results overwrite prior page-level results.").classes(
"text-xs ui-text-muted"
)
async def submit_resubmit() -> None:
try:
resubmitted_count = await jobs_service.resubmit_non_transcribed_sources(job_id=job.id)
except JobResubmitBlockedError as exc:
ui.notify(exc.message, type="warning")
return
except ValueError:
ui.notify("Job not found.", type="warning")
ui.navigate.to("/jobs")
return
except Exception as exc: # noqa: BLE001
show_error(exc, title="Resubmit failed", operation="jobs.resubmit")
return
resolve_worker_notifier(request.app.state).notify()
ui.notify(f"Resubmitted {resubmitted_count} source(s)", type="positive")
ui.navigate.to(f"/jobs/{job.id}")
with ui.row().classes("w-full items-center gap-2 mt-2"):
ui.button("Resubmit now", on_click=submit_resubmit, icon="replay").classes("ui-btn-primary")
ui.button("Back to Job", on_click=lambda: ui.navigate.to(f"/jobs/{job.id}"), icon="arrow_back").props("flat")
@ui.page("/jobs/{job_id}/delete") @ui.page("/jobs/{job_id}/delete")
async def job_delete_page(job_id: str, session_factory: SessionFactoryDep) -> None: async def job_delete_page(job_id: str, session_factory: SessionFactoryDep) -> None:
apply_archival_theme() apply_archival_theme()
+20 -22
View File
@@ -5,6 +5,7 @@ from __future__ import annotations
from datetime import date from datetime import date
from urllib.parse import quote from urllib.parse import quote
from uuid import UUID from uuid import UUID
from uuid import uuid4
from fastapi import Request from fastapi import Request
from nicegui import ui from nicegui import ui
@@ -15,7 +16,6 @@ from transcription.errors import ErrorCategory
from transcription.services.documents import ( from transcription.services.documents import (
DocumentError, DocumentError,
DocumentService, DocumentService,
PersonDeleteBlockedError,
) )
from transcription.services.store import UploadError, store_person_portrait from transcription.services.store import UploadError, store_person_portrait
from transcription.ui.components.app_shell import render_navigation_header from transcription.ui.components.app_shell import render_navigation_header
@@ -44,11 +44,12 @@ def _parse_optional_date(value: str | None, *, label: str) -> date | None:
raise ValueError(f"{label} must use YYYY-MM-DD.") from exc raise ValueError(f"{label} must use YYYY-MM-DD.") from exc
def _bind_portrait_file_picker(portrait_path_input: ui.input, *, settings: Settings) -> None: def _bind_portrait_file_picker(portrait_path_input: ui.input, *, settings: Settings, person_id: UUID) -> None:
async def on_portrait_selected(event) -> None: async def on_portrait_selected(event) -> None:
payload = await event.file.read() payload = await event.file.read()
try: try:
stored_path = store_person_portrait( stored_path = store_person_portrait(
person_id=person_id,
filename=event.file.name, filename=event.file.name,
file_bytes=payload, file_bytes=payload,
settings=settings, settings=settings,
@@ -73,7 +74,8 @@ def _bind_portrait_file_picker(portrait_path_input: ui.input, *, settings: Setti
auto_upload=True, auto_upload=True,
label="Choose portrait file", label="Choose portrait file",
).props('accept=".jpg,.jpeg,.png,.gif,.webp,.bmp,.tif,.tiff"').classes("w-full") ).props('accept=".jpg,.jpeg,.png,.gif,.webp,.bmp,.tif,.tiff"').classes("w-full")
ui.label("Portraits are stored under uploads/portraits/person.").classes("text-xs ui-text-muted") portrait_dir = settings.upload_dir / "persons" / str(person_id)
ui.label(f"Portraits are stored under {portrait_dir}.").classes("text-xs ui-text-muted")
def _resolve_portrait_src(path: str | None) -> str | None: def _resolve_portrait_src(path: str | None) -> str | None:
@@ -145,6 +147,7 @@ def register_page() -> None: # noqa: PLR0915
apply_archival_theme() apply_archival_theme()
people_service = DocumentService(session_factory=session_factory) people_service = DocumentService(session_factory=session_factory)
render_navigation_header(current_path="/people") render_navigation_header(current_path="/people")
draft_person_id = uuid4()
with ui.column().classes("w-full max-w-4xl mx-auto p-4 gap-4"): with ui.column().classes("w-full max-w-4xl mx-auto p-4 gap-4"):
page_header("Create Person Record", subtitle="Full name is required.") page_header("Create Person Record", subtitle="Full name is required.")
@@ -167,7 +170,11 @@ def register_page() -> None: # noqa: PLR0915
biography_input = ui.textarea(label="Biography").props("outlined bg-white autogrow").classes("w-full") biography_input = ui.textarea(label="Biography").props("outlined bg-white autogrow").classes("w-full")
portrait_path_input = ui.input(label="Portrait path").props("outlined bg-white").classes("w-full") portrait_path_input = ui.input(label="Portrait path").props("outlined bg-white").classes("w-full")
_bind_portrait_file_picker(portrait_path_input, settings=_resolve_runtime_settings(request)) _bind_portrait_file_picker(
portrait_path_input,
settings=_resolve_runtime_settings(request),
person_id=draft_person_id,
)
async def submit_create() -> None: async def submit_create() -> None:
full_name = (full_name_input.value or "").strip() full_name = (full_name_input.value or "").strip()
@@ -183,6 +190,7 @@ def register_page() -> None: # noqa: PLR0915
return return
candidate = Person( candidate = Person(
id=draft_person_id,
full_name=full_name, full_name=full_name,
display_name=(display_name_input.value or "").strip() or None, display_name=(display_name_input.value or "").strip() or None,
maiden_name=(maiden_name_input.value or "").strip() or None, maiden_name=(maiden_name_input.value or "").strip() or None,
@@ -357,7 +365,11 @@ def register_page() -> None: # noqa: PLR0915
portrait_path_input = ( portrait_path_input = (
ui.input(label="Portrait path", value=person.portrait_path or "").props("outlined bg-white").classes("w-full") ui.input(label="Portrait path", value=person.portrait_path or "").props("outlined bg-white").classes("w-full")
) )
_bind_portrait_file_picker(portrait_path_input, settings=_resolve_runtime_settings(request)) _bind_portrait_file_picker(
portrait_path_input,
settings=_resolve_runtime_settings(request),
person_id=person.id,
)
async def submit_edit() -> None: async def submit_edit() -> None:
full_name = (full_name_input.value or "").strip() full_name = (full_name_input.value or "").strip()
@@ -431,29 +443,15 @@ def register_page() -> None: # noqa: PLR0915
ui.label(f"Person: {person.full_name}").classes("text-sm font-semibold ui-text-primary") ui.label(f"Person: {person.full_name}").classes("text-sm font-semibold ui-text-primary")
if person.document_people: if person.document_people:
ui.label("Delete is blocked because linked documents exist.").classes("text-xs text-red-800 font-bold mt-2") ui.label(
ui.label(f"Linked documents: {len(person.document_people)}").classes("text-xs ui-text-muted") f"This will also remove {len(person.document_people)} linked document relationship(s)."
ui.label("Remove document links first, then retry deletion.").classes("text-xs ui-text-muted italic") ).classes("text-xs text-red-800 font-bold mt-2")
with ui.row().classes("w-full items-center gap-2 mt-4"):
ui.button(
"Back to Person",
on_click=lambda: ui.navigate.to(f"/people/{person.id}"),
icon="arrow_back",
).classes("ui-btn-primary text-xs")
ui.button("Go to Documents", on_click=lambda: ui.navigate.to("/documents"), icon="description").props(
"flat text-xs"
)
return
ui.label("This action permanently deletes the person record.").classes("text-xs text-red-800 font-medium") ui.label("This action permanently deletes the person record.").classes("text-xs text-red-800 font-medium")
async def submit_delete() -> None: async def submit_delete() -> None:
try: try:
await people_service.delete_person(person) await people_service.delete_person(person)
except PersonDeleteBlockedError as exc:
ui.notify(exc.message, type="warning")
ui.navigate.to(f"/people/{person.id}/delete")
return
except DocumentError as exc: except DocumentError as exc:
if exc.category == ErrorCategory.NOT_FOUND: if exc.category == ErrorCategory.NOT_FOUND:
ui.notify("Person not found.", type="warning") ui.notify("Person not found.", type="warning")
+103
View File
@@ -13,6 +13,7 @@ from transcription.db.models import JobSource, Source
from transcription.services.documents import DocumentError, DocumentService from transcription.services.documents import DocumentError, DocumentService
from transcription.services.jobs import JobService from transcription.services.jobs import JobService
from transcription.services.transcription import ( from transcription.services.transcription import (
SourceDeleteBlockedError,
TranscriptionNotFoundError, TranscriptionNotFoundError,
TranscriptionService, TranscriptionService,
) )
@@ -21,6 +22,7 @@ from transcription.ui.components.cards import archival_card
from transcription.ui.components.data_display import metadata_row from transcription.ui.components.data_display import metadata_row
from transcription.ui.components.document_panzoom import render_document_panzoom from transcription.ui.components.document_panzoom import render_document_panzoom
from transcription.ui.components.error_presenter import show_error from transcription.ui.components.error_presenter import show_error
from transcription.ui.components.primitives import destructive_button
from transcription.ui.components.primitives import section_header_row from transcription.ui.components.primitives import section_header_row
from transcription.ui.components.table.sources import SourceTableRow, render_sources_table from transcription.ui.components.table.sources import SourceTableRow, render_sources_table
from transcription.ui.theme import apply_archival_theme from transcription.ui.theme import apply_archival_theme
@@ -50,6 +52,7 @@ def register_page() -> None:
job_label = None job_label = None
back_path = None back_path = None
sources: list[Source] = [] sources: list[Source] = []
job_source_by_source_id: dict[UUID, JobSource] = {}
try: try:
if document_id is not None: if document_id is not None:
@@ -62,6 +65,10 @@ def register_page() -> None:
job_label = str(job.id) job_label = str(job.id)
back_path = f"/jobs/{job.id}" back_path = f"/jobs/{job.id}"
job_sources = await sources_service.list_job_sources(job_id=job.id) job_sources = await sources_service.list_job_sources(job_id=job.id)
job_source_by_source_id = {
job_source.source_id: job_source
for job_source in job_sources
}
sources = [job_source.source for job_source in job_sources if job_source.source is not None] sources = [job_source.source for job_source in job_sources if job_source.source is not None]
sources.sort(key=lambda item: (item.page_number, item.upload_name.casefold())) sources.sort(key=lambda item: (item.page_number, item.upload_name.casefold()))
else: else:
@@ -101,6 +108,16 @@ def register_page() -> None:
upload_name=source.upload_name, upload_name=source.upload_name,
filename=source.filename, filename=source.filename,
document_id=source.document_id, document_id=source.document_id,
job_source_status=(
job_source_by_source_id[source.id].status.value
if source.id in job_source_by_source_id
else None
),
job_source_error_detail=(
job_source_by_source_id[source.id].error_detail
if source.id in job_source_by_source_id
else None
),
) )
for source in sources for source in sources
] ]
@@ -133,6 +150,7 @@ def register_page() -> None:
with section_header_row(): with section_header_row():
page_header(f"Source Page {source.page_number}: {source.upload_name}", subtitle=f"Source ID: {source.id}") page_header(f"Source Page {source.page_number}: {source.upload_name}", subtitle=f"Source ID: {source.id}")
with ui.row().classes("items-center gap-2"):
if back_path is not None: if back_path is not None:
back_label = ( back_label = (
"Back to Document" "Back to Document"
@@ -149,6 +167,13 @@ def register_page() -> None:
"flat text-xs" "flat text-xs"
) )
destructive_button(
"Delete Source",
on_click=lambda: ui.navigate.to(f"/sources/{source.id}/delete{_back_query(request.query_params)}"),
icon="delete",
extra_classes="text-xs",
)
with ui.grid().classes("w-full grid-cols-12 gap-4"): with ui.grid().classes("w-full grid-cols-12 gap-4"):
with ui.column().classes("col-span-12 lg:col-span-7 gap-4"): with ui.column().classes("col-span-12 lg:col-span-7 gap-4"):
with archival_card(title="Source Inspection Viewer", extra_classes="p-2"): with archival_card(title="Source Inspection Viewer", extra_classes="p-2"):
@@ -166,6 +191,17 @@ def register_page() -> None:
source.date_revised.isoformat() if source.date_revised else "Not revised", source.date_revised.isoformat() if source.date_revised else "Not revised",
) )
with archival_card(title="Job Source Outcomes"):
if not source.job_sources:
ui.label("No job-source execution records found for this source.").classes("text-xs ui-text-muted")
else:
for job_source in sorted(source.job_sources, key=lambda item: item.executed_at, reverse=True):
with ui.column().classes("w-full gap-1 p-2 ui-row-surface rounded"):
metadata_row("Job ID:", str(job_source.job_id))
metadata_row("Status:", job_source.status.value)
metadata_row("Executed At:", job_source.executed_at.isoformat())
metadata_row("Error Detail:", job_source.error_detail or "None")
with archival_card(title="Automated Raw Transcription"): with archival_card(title="Automated Raw Transcription"):
ui.textarea(value=_source_transcription_text(source) or "").props("outlined autogrow readonly bg-white").classes( ui.textarea(value=_source_transcription_text(source) or "").props("outlined autogrow readonly bg-white").classes(
"w-full text-xs font-mono" "w-full text-xs font-mono"
@@ -196,6 +232,73 @@ def register_page() -> None:
with ui.row().classes("w-full items-center gap-2 mt-2"): with ui.row().classes("w-full items-center gap-2 mt-2"):
ui.button("Save Revision", on_click=save_revision, icon="save").classes("ui-btn-primary text-xs") ui.button("Save Revision", on_click=save_revision, icon="save").classes("ui-btn-primary text-xs")
@ui.page("/sources/{source_id}/delete")
async def source_delete_page(source_id: str, request: Request, session_factory: SessionFactoryDep) -> None:
apply_archival_theme()
sources_service = TranscriptionService(session_factory=session_factory)
render_navigation_header(current_path="/sources")
try:
parsed_source_id = UUID(source_id)
except ValueError:
ui.label("Invalid source id").classes("text-h6 text-red-800 p-4")
return
try:
source = await sources_service.read_source_detail(source_id=parsed_source_id)
except TranscriptionNotFoundError:
ui.label("Source not found").classes("text-h6 text-red-800 p-4")
return
except Exception as exc: # noqa: BLE001
show_error(exc, title="Load failed", operation="sources.delete.load")
return
back_path = _back_path_from_query(request.query_params) or "/sources"
next_sources_path = f"/sources{_back_query(request.query_params)}"
linked_count = len(source.job_sources)
with ui.column().classes("w-full max-w-xl mx-auto p-4 gap-4"):
page_header("Delete Source Record")
with archival_card(extra_classes="gap-2"):
metadata_row("Source ID:", str(source.id))
metadata_row("Upload Name:", source.upload_name)
metadata_row("Linked Jobs:", str(linked_count))
if linked_count > 0:
ui.label("Delete is only available for unlinked sources.").classes("text-xs text-red-800 font-bold mt-2")
ui.label("This source is linked to one or more jobs and cannot be deleted from this view.").classes(
"text-xs ui-text-muted italic"
)
else:
ui.label("This action permanently deletes the source record.").classes("text-xs text-red-800 font-medium")
async def submit_delete() -> None:
try:
await sources_service.delete_unlinked_source(source_id=source.id)
except SourceDeleteBlockedError as exc:
ui.notify(exc.message, type="warning")
return
except TranscriptionNotFoundError:
ui.notify("Source not found.", type="warning")
ui.navigate.to(next_sources_path)
return
except Exception as exc: # noqa: BLE001
show_error(exc, title="Delete source failed", operation="sources.delete")
return
ui.notify("Source deleted", type="positive")
ui.navigate.to(next_sources_path)
with ui.row().classes("w-full items-center gap-2 mt-2"):
destructive_button(
"Delete source permanently",
on_click=submit_delete,
icon="delete_forever",
variant="solid",
)
ui.button("Cancel", on_click=lambda route=back_path: ui.navigate.to(route), icon="arrow_back").props("flat")
@ui.page("/documents/{document_id}/sources") @ui.page("/documents/{document_id}/sources")
async def document_sources_page(document_id: str) -> RedirectResponse: async def document_sources_page(document_id: str) -> RedirectResponse:
return RedirectResponse(url=f"/ui/sources?document_id={document_id}") return RedirectResponse(url=f"/ui/sources?document_id={document_id}")
+8 -1
View File
@@ -152,8 +152,15 @@ async def run_worker_loop(
wake_event.clear() wake_event.clear()
processed_any = False processed_any = False
while await process_next_queued_job(session_factory=session_factory): while True:
with handle_worker_exceptions(operation="worker.process_next_queued_job"):
processed = await process_next_queued_job(session_factory=session_factory)
if not processed:
break
processed_any = True processed_any = True
continue
break
if wake_event is None and not processed_any: if wake_event is None and not processed_any:
await asyncio.sleep(poll_interval_seconds) await asyncio.sleep(poll_interval_seconds)
+231 -8
View File
@@ -1,25 +1,49 @@
"""Integration tests for end-to-end upload and worker pipeline behavior.""" """Integration tests for end-to-end upload and worker pipeline behavior."""
from pathlib import Path from pathlib import Path
from uuid import uuid4
import pytest import pytest
from transcription.config import Settings from transcription.config import Settings
from transcription.db.models import Document
from transcription.db.models import Job from transcription.db.models import Job
from transcription.db.models import JobSourceStatus
from transcription.db.models import JobStatus from transcription.db.models import JobStatus
from transcription.providers.base import TranscriptionResult from transcription.providers.base import TranscriptionResult
from transcription.services import ServiceBundle from transcription.services import ServiceBundle
from transcription.services.store import create_job_for_document
from transcription.services.store import create_upload_job from transcription.services.store import create_upload_job
from transcription.services.workflows import advance_job from transcription.services.workflows import advance_job
def _build_services(default_session_factory) -> ServiceBundle:
services = ServiceBundle()
object.__setattr__(
services,
"documents",
services.documents.__class__(session_factory=default_session_factory),
)
object.__setattr__(
services,
"jobs",
services.jobs.__class__(session_factory=default_session_factory),
)
object.__setattr__(
services,
"transcriptions",
services.transcriptions.__class__(session_factory=default_session_factory),
)
return services
@pytest.mark.integration @pytest.mark.integration
class TestPipelineSuccessFlow: class TestPipelineSuccessFlow:
"""Verify end-to-end success lifecycle behavior.""" """Verify end-to-end success lifecycle behavior."""
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_upload_then_worker_persists_transcribed_terminal_state( async def test_upload_then_worker_persists_transcribed_terminal_state(
self, async_session, tmp_path: Path, monkeypatch self, async_session, default_session_factory, tmp_path: Path, monkeypatch
): ):
"""Upload followed by worker processing persists job transcription and transcribed status.""" """Upload followed by worker processing persists job transcription and transcribed status."""
settings = Settings(openrouter_api_key="test-key", upload_dir=tmp_path) settings = Settings(openrouter_api_key="test-key", upload_dir=tmp_path)
@@ -59,12 +83,12 @@ class TestPipelineSuccessFlow:
_fake_transcribe_document_image, _fake_transcribe_document_image,
) )
services = ServiceBundle() services = _build_services(default_session_factory)
queued_job = await services.jobs.read_next_queued_job(session=async_session) queued_job = await services.jobs.read_job(job_id=upload_result.job_id, session=async_session)
processed = queued_job is not None processed = queued_job is not None
if queued_job is not None: if queued_job is not None:
await advance_job(job=queued_job, services=services, session=async_session) await advance_job(job=queued_job, services=services, session=async_session)
job = await async_session.get(Job, upload_result.job_id) job = await services.jobs.read_job(job_id=upload_result.job_id, session=async_session)
assert processed is True assert processed is True
assert job is not None assert job is not None
@@ -72,13 +96,212 @@ class TestPipelineSuccessFlow:
assert any(job_source.raw_transcription == "Pipeline transcript" for job_source in job.job_sources) assert any(job_source.raw_transcription == "Pipeline transcript" for job_source in job.job_sources)
assert all(job_source.error_detail is None for job_source in job.job_sources) assert all(job_source.error_detail is None for job_source in job.job_sources)
@pytest.mark.asyncio
async def test_worker_transcribes_all_sources_for_multi_page_job(
self,
async_session,
default_session_factory,
tmp_path: Path,
monkeypatch,
):
"""Worker stores transcription output for every source linked to the queued job."""
settings = Settings(openrouter_api_key="test-key", upload_dir=tmp_path)
document = Document(id=uuid4(), name="multi-page-document")
async_session.add(document)
await async_session.commit()
create_result = await create_job_for_document(
document_id=document.id,
uploads=[
("page-01.jpg", b"one"),
("page-02.jpg", b"two"),
("page-03.jpg", b"three"),
],
session=async_session,
settings=settings,
)
async def _fake_transcribe_document_image(
image_path,
*,
prompt_name="transcribe_document.md",
settings=None,
provider=None,
) -> TranscriptionResult:
page_name = Path(image_path).name
_ = (prompt_name, settings, provider)
return TranscriptionResult(
text=f"Transcript for {page_name}",
provider="openrouter",
model="test-model",
prompt_name="transcribe_document.md",
)
monkeypatch.setattr(
"transcription.services.workflows.transcribe_document_image",
_fake_transcribe_document_image,
)
services = _build_services(default_session_factory)
queued_job = await services.jobs.read_job(job_id=create_result.job_id, session=async_session)
assert queued_job is not None
await advance_job(job=queued_job, services=services, session=async_session)
job = await services.jobs.read_job(job_id=create_result.job_id, session=async_session)
assert job.status == JobStatus.TRANSCRIBED
assert len(job.job_sources) == 3
assert all(job_source.status == JobSourceStatus.TRANSCRIBED for job_source in job.job_sources)
assert all(job_source.raw_transcription for job_source in job.job_sources)
assert all(job_source.source is not None and job_source.source.raw_transcription for job_source in job.job_sources)
@pytest.mark.asyncio
async def test_worker_marks_partial_success_when_some_sources_fail(
self,
async_session,
default_session_factory,
tmp_path: Path,
monkeypatch,
):
"""Mixed page outcomes produce PARTIAL_SUCCESS and preserve per-source status."""
settings = Settings(openrouter_api_key="test-key", upload_dir=tmp_path)
document = Document(id=uuid4(), name="partial-page-document")
async_session.add(document)
await async_session.commit()
create_result = await create_job_for_document(
document_id=document.id,
uploads=[
("page-01.jpg", b"one"),
("page-02.jpg", b"two"),
],
session=async_session,
settings=settings,
)
call_count = 0
async def _fake_transcribe_document_image(
image_path,
*,
prompt_name="transcribe_document.md",
settings=None,
provider=None,
) -> TranscriptionResult:
nonlocal call_count
call_count += 1
_ = (prompt_name, settings, provider)
if call_count == 2:
raise RuntimeError("simulated page failure")
return TranscriptionResult(
text="Transcript for first page",
provider="openrouter",
model="test-model",
prompt_name="transcribe_document.md",
)
monkeypatch.setattr(
"transcription.services.workflows.transcribe_document_image",
_fake_transcribe_document_image,
)
services = _build_services(default_session_factory)
queued_job = await services.jobs.read_job(job_id=create_result.job_id, session=async_session)
assert queued_job is not None
await advance_job(job=queued_job, services=services, session=async_session)
job = await services.jobs.read_job(job_id=create_result.job_id, session=async_session)
assert job.status == JobStatus.PARTIAL_SUCCESS
assert len(job.job_sources) == 2
statuses = {job_source.status for job_source in job.job_sources}
assert statuses == {JobSourceStatus.TRANSCRIBED, JobSourceStatus.FAILED}
assert any(job_source.error_detail is not None for job_source in job.job_sources)
@pytest.mark.asyncio
async def test_worker_skips_already_transcribed_sources_on_resubmit(
self,
async_session,
default_session_factory,
tmp_path: Path,
monkeypatch,
):
"""Queued jobs only process non-transcribed JobSource records."""
settings = Settings(openrouter_api_key="test-key", upload_dir=tmp_path)
document = Document(id=uuid4(), name="resubmit-filter-document")
async_session.add(document)
await async_session.commit()
create_result = await create_job_for_document(
document_id=document.id,
uploads=[
("page-01.jpg", b"one"),
("page-02.jpg", b"two"),
],
session=async_session,
settings=settings,
)
services = _build_services(default_session_factory)
job = await services.jobs.read_job(job_id=create_result.job_id, session=async_session)
page_one = next(js for js in job.job_sources if js.source is not None and js.source.page_number == 1)
page_two = next(js for js in job.job_sources if js.source is not None and js.source.page_number == 2)
page_one.status = JobSourceStatus.TRANSCRIBED
page_one.raw_transcription = "existing transcript"
page_two.status = JobSourceStatus.PENDING
page_two.raw_transcription = None
await services.transcriptions.update_job_source(job_source=page_one, session=async_session)
await services.transcriptions.update_job_source(job_source=page_two, session=async_session)
await services.jobs.update_job_state(job_id=job.id, status=JobStatus.QUEUED, session=async_session)
await async_session.commit()
call_count = 0
async def _fake_transcribe_document_image(
image_path,
*,
prompt_name="transcribe_document.md",
settings=None,
provider=None,
) -> TranscriptionResult:
nonlocal call_count
_ = (image_path, prompt_name, settings, provider)
call_count += 1
return TranscriptionResult(
text="new transcript",
provider="openrouter",
model="test-model",
prompt_name="transcribe_document.md",
)
monkeypatch.setattr(
"transcription.services.workflows.transcribe_document_image",
_fake_transcribe_document_image,
)
queued_job = await services.jobs.read_job(job_id=create_result.job_id, session=async_session)
assert queued_job is not None
await advance_job(job=queued_job, services=services, session=async_session)
refreshed = await services.jobs.read_job(job_id=create_result.job_id, session=async_session)
assert call_count == 1
statuses = {js.status for js in refreshed.job_sources}
assert statuses == {JobSourceStatus.TRANSCRIBED}
@pytest.mark.integration @pytest.mark.integration
class TestPipelineFailureFlow: class TestPipelineFailureFlow:
"""Verify end-to-end failure lifecycle behavior.""" """Verify end-to-end failure lifecycle behavior."""
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_upload_then_worker_persists_failed_terminal_state(self, async_session, tmp_path: Path, monkeypatch): async def test_upload_then_worker_persists_failed_terminal_state(
self,
async_session,
default_session_factory,
tmp_path: Path,
monkeypatch,
):
"""Upload followed by worker processing persists error detail and failed status on the job.""" """Upload followed by worker processing persists error detail and failed status on the job."""
settings = Settings(openrouter_api_key="test-key", upload_dir=tmp_path) settings = Settings(openrouter_api_key="test-key", upload_dir=tmp_path)
upload_result = await create_upload_job( upload_result = await create_upload_job(
@@ -103,12 +326,12 @@ class TestPipelineFailureFlow:
_fake_transcribe_document_image, _fake_transcribe_document_image,
) )
services = ServiceBundle() services = _build_services(default_session_factory)
queued_job = await services.jobs.read_next_queued_job(session=async_session) queued_job = await services.jobs.read_job(job_id=upload_result.job_id, session=async_session)
processed = queued_job is not None processed = queued_job is not None
if queued_job is not None: if queued_job is not None:
await advance_job(job=queued_job, services=services, session=async_session) await advance_job(job=queued_job, services=services, session=async_session)
job = await async_session.get(Job, upload_result.job_id) job = await services.jobs.read_job(job_id=upload_result.job_id, session=async_session)
assert processed is True assert processed is True
assert job is not None assert job is not None
+43 -4
View File
@@ -2,6 +2,7 @@ from __future__ import annotations
from datetime import UTC from datetime import UTC
from datetime import datetime from datetime import datetime
from pathlib import Path
from uuid import uuid4 from uuid import uuid4
import pytest import pytest
@@ -14,7 +15,6 @@ from transcription.db.models import Person
from transcription.db.models import Source from transcription.db.models import Source
from transcription.services.documents import DocumentDeleteBlockedError from transcription.services.documents import DocumentDeleteBlockedError
from transcription.services.documents import DocumentError from transcription.services.documents import DocumentError
from transcription.services.documents import PersonDeleteBlockedError
from transcription.services.documents import DocumentService from transcription.services.documents import DocumentService
@@ -88,8 +88,9 @@ async def test_delete_document_blocks_when_dependencies_exist(default_session_fa
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_delete_document_succeeds_when_unlinked(default_session_factory): async def test_delete_document_succeeds_when_unlinked(default_session_factory, tmp_path):
service = DocumentService(session_factory=default_session_factory) service = DocumentService(session_factory=default_session_factory)
service.settings.upload_dir = tmp_path
document = await service.create_document( document = await service.create_document(
Document( Document(
@@ -99,12 +100,45 @@ async def test_delete_document_succeeds_when_unlinked(default_session_factory):
) )
) )
document_dir = service.settings.upload_dir / "documents" / str(document.id)
document_dir.mkdir(parents=True, exist_ok=True)
(document_dir / "leftover.txt").write_text("orphan", encoding="utf-8")
await service.delete_document(document) await service.delete_document(document)
assert not document_dir.exists()
with pytest.raises(DocumentError): with pytest.raises(DocumentError):
await service.read_document_detail(document.id) await service.read_document_detail(document.id)
@pytest.mark.asyncio
async def test_delete_document_removes_populated_storage_tree(default_session_factory, tmp_path):
service = DocumentService(session_factory=default_session_factory)
service.settings.upload_dir = tmp_path
document = await service.create_document(
Document(
id=uuid4(),
name="tree-delete",
document_type="memo",
)
)
document_dir = service.settings.upload_dir / "documents" / str(document.id)
(document_dir / "page-1.jpg").parent.mkdir(parents=True, exist_ok=True)
(document_dir / "page-1.jpg").write_bytes(b"one")
(document_dir / "page-2.jpg").write_bytes(b"two")
(document_dir / "nested" / "manifest.json").parent.mkdir(parents=True, exist_ok=True)
(document_dir / "nested" / "manifest.json").write_text('{"ok": true}', encoding="utf-8")
assert document_dir.exists()
await service.delete_document(document)
assert not document_dir.exists()
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_read_person_detail_loads_document_links(default_session_factory): async def test_read_person_detail_loads_document_links(default_session_factory):
service = DocumentService(session_factory=default_session_factory) service = DocumentService(session_factory=default_session_factory)
@@ -153,7 +187,7 @@ async def test_update_person_refreshes_updated_timestamp(default_session_factory
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_delete_person_blocks_when_linked_documents_exist(default_session_factory): async def test_delete_person_removes_links_when_linked_documents_exist(default_session_factory):
service = DocumentService(session_factory=default_session_factory) service = DocumentService(session_factory=default_session_factory)
document = await service.create_document( document = await service.create_document(
@@ -172,9 +206,14 @@ async def test_delete_person_blocks_when_linked_documents_exist(default_session_
) )
) )
with pytest.raises(PersonDeleteBlockedError):
await service.delete_person(person) await service.delete_person(person)
links = await service.list_document_people(person_id=person.id)
assert links == []
with pytest.raises(DocumentError):
await service.read_person_detail(person.id)
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_delete_person_succeeds_when_unlinked(default_session_factory): async def test_delete_person_succeeds_when_unlinked(default_session_factory):
+156
View File
@@ -13,6 +13,8 @@ from transcription.db.models import JobStatus
from transcription.db.models import Source from transcription.db.models import Source
from transcription.services.documents import DocumentService from transcription.services.documents import DocumentService
from transcription.services.jobs import JobDeleteBlockedError from transcription.services.jobs import JobDeleteBlockedError
from transcription.services.jobs import JobCancelBlockedError
from transcription.services.jobs import JobResubmitBlockedError
from transcription.services.jobs import JobService from transcription.services.jobs import JobService
@@ -224,3 +226,157 @@ class TestJobService:
with pytest.raises(ValueError): with pytest.raises(ValueError):
await job_service.read_job(job_id=job.id) await job_service.read_job(job_id=job.id)
@pytest.mark.asyncio
async def test_cancel_job_marks_non_transcribed_sources_failed(
self,
job_service: JobService,
document_service: DocumentService,
):
document = Document(id=uuid4(), name="cancel-job-doc")
await document_service.create_document(document=document)
job = Job(document_id=document.id, status=JobStatus.QUEUED)
await job_service.create_job(job=job)
async with job_service._session_scope() as session:
source_one = Source(
document_id=document.id,
page_number=1,
upload_name="cancel-1.jpg",
filename="stored-cancel-1.jpg",
file_path="/uploads/stored-cancel-1.jpg",
)
source_two = Source(
document_id=document.id,
page_number=2,
upload_name="cancel-2.jpg",
filename="stored-cancel-2.jpg",
file_path="/uploads/stored-cancel-2.jpg",
)
session.add(source_one)
session.add(source_two)
await session.flush()
session.add(
JobSource(
job_id=job.id,
source_id=source_one.id,
status=JobSourceStatus.TRANSCRIBED,
raw_transcription="done",
)
)
session.add(
JobSource(
job_id=job.id,
source_id=source_two.id,
status=JobSourceStatus.PENDING,
)
)
await session.commit()
cancelled = await job_service.cancel_job(job_id=job.id)
assert cancelled.status == JobStatus.FAILED
refreshed = await job_service.read_job(job_id=job.id)
statuses = {item.status for item in refreshed.job_sources}
assert JobSourceStatus.TRANSCRIBED in statuses
assert JobSourceStatus.FAILED in statuses
pending_entry = next(item for item in refreshed.job_sources if item.status == JobSourceStatus.FAILED)
assert pending_entry.error_detail == "Cancelled by user"
@pytest.mark.asyncio
async def test_resubmit_non_transcribed_sources_resets_only_non_transcribed(
self,
job_service: JobService,
document_service: DocumentService,
):
document = Document(id=uuid4(), name="resubmit-job-doc")
await document_service.create_document(document=document)
job = Job(document_id=document.id, status=JobStatus.FAILED)
await job_service.create_job(job=job)
async with job_service._session_scope() as session:
source_one = Source(
document_id=document.id,
page_number=1,
upload_name="resubmit-1.jpg",
filename="stored-resubmit-1.jpg",
file_path="/uploads/stored-resubmit-1.jpg",
raw_transcription="existing text",
)
source_two = Source(
document_id=document.id,
page_number=2,
upload_name="resubmit-2.jpg",
filename="stored-resubmit-2.jpg",
file_path="/uploads/stored-resubmit-2.jpg",
raw_transcription="done text",
)
session.add(source_one)
session.add(source_two)
await session.flush()
session.add(
JobSource(
job_id=job.id,
source_id=source_one.id,
status=JobSourceStatus.FAILED,
raw_transcription=None,
error_detail="prior error",
)
)
session.add(
JobSource(
job_id=job.id,
source_id=source_two.id,
status=JobSourceStatus.TRANSCRIBED,
raw_transcription="done text",
)
)
await session.commit()
count = await job_service.resubmit_non_transcribed_sources(job_id=job.id)
assert count == 1
refreshed = await job_service.read_job(job_id=job.id)
assert refreshed.status == JobStatus.QUEUED
failed_entry = next(item for item in refreshed.job_sources if item.source is not None and item.source.page_number == 1)
transcribed_entry = next(item for item in refreshed.job_sources if item.source is not None and item.source.page_number == 2)
assert failed_entry.status == JobSourceStatus.PENDING
assert failed_entry.error_detail is None
assert failed_entry.source is not None
assert failed_entry.source.raw_transcription is None
assert transcribed_entry.status == JobSourceStatus.TRANSCRIBED
@pytest.mark.asyncio
async def test_resubmit_non_transcribed_sources_blocks_when_processing(
self,
job_service: JobService,
document_service: DocumentService,
):
document = Document(id=uuid4(), name="resubmit-blocked-doc")
await document_service.create_document(document=document)
job = Job(document_id=document.id, status=JobStatus.PROCESSING)
await job_service.create_job(job=job)
with pytest.raises(JobResubmitBlockedError):
await job_service.resubmit_non_transcribed_sources(job_id=job.id)
@pytest.mark.asyncio
async def test_cancel_job_blocks_transcribed_terminal_jobs(
self,
job_service: JobService,
document_service: DocumentService,
):
document = Document(id=uuid4(), name="cancel-blocked-doc")
await document_service.create_document(document=document)
job = Job(document_id=document.id, status=JobStatus.TRANSCRIBED)
await job_service.create_job(job=job)
with pytest.raises(JobCancelBlockedError):
await job_service.cancel_job(job_id=job.id)
+48
View File
@@ -1,3 +1,4 @@
from pathlib import Path
from uuid import uuid4 from uuid import uuid4
import pytest import pytest
@@ -9,7 +10,9 @@ from transcription.db.models import Job
from transcription.db.models import JobSource from transcription.db.models import JobSource
from transcription.db.models import Source from transcription.db.models import Source
from transcription.services.store import UploadError from transcription.services.store import UploadError
from transcription.services.store import create_upload_job
from transcription.services.store import create_job_for_document from transcription.services.store import create_job_for_document
from transcription.services.store import store_person_portrait
@pytest.mark.asyncio @pytest.mark.asyncio
@@ -66,7 +69,52 @@ async def test_create_job_for_document_sorts_uploads_and_creates_links(async_ses
assert [source.upload_name for source in sources] == ["A_page.pdf", "b_page.pdf"] assert [source.upload_name for source in sources] == ["A_page.pdf", "b_page.pdf"]
assert all(source.filename.endswith(".pdf") for source in sources) assert all(source.filename.endswith(".pdf") for source in sources)
assert all("A_page" not in source.filename and "b_page" not in source.filename for source in sources) assert all("A_page" not in source.filename and "b_page" not in source.filename for source in sources)
assert all(Path(source.filename).stem == str(source.id) for source in sources)
assert all(Path(source.file_path).parent == (tmp_path / "documents" / str(document.id)) for source in sources)
job_sources = (await async_session.exec(select(JobSource).where(JobSource.job_id == result.job_id))).all() job_sources = (await async_session.exec(select(JobSource).where(JobSource.job_id == result.job_id))).all()
assert len(job_sources) == 2 assert len(job_sources) == 2
assert set(result.source_ids) == {job_source.source_id for job_source in job_sources} assert set(result.source_ids) == {job_source.source_id for job_source in job_sources}
@pytest.mark.asyncio
async def test_create_upload_job_stores_source_under_document_id_directory(async_session, tmp_path):
settings = Settings(openrouter_api_key="test-key", upload_dir=tmp_path)
result = await create_upload_job(
filename="single-page.jpg",
file_bytes=b"image-bytes",
session=async_session,
settings=settings,
)
expected_parent = tmp_path / "documents" / str(result.document_id)
assert result.stored_path.parent == expected_parent
assert result.stored_path.exists()
source = (
await async_session.exec(
select(Source)
.where(Source.document_id == result.document_id)
.order_by(Source.page_number) # pyright: ignore[reportArgumentType]
)
).first()
assert source is not None
assert Path(source.filename).stem == str(source.id)
assert result.stored_path.name == source.filename
assert Path(source.file_path).parent == expected_parent
def test_store_person_portrait_stores_file_under_person_id_directory(tmp_path):
settings = Settings(openrouter_api_key="test-key", upload_dir=tmp_path)
person_id = uuid4()
stored_path = store_person_portrait(
person_id=person_id,
filename="portrait.png",
file_bytes=b"portrait-bytes",
settings=settings,
)
assert stored_path.parent == (tmp_path / "persons" / str(person_id))
assert stored_path.exists()
+65 -2
View File
@@ -93,10 +93,11 @@ class TestTranscriptionServiceRevisionUpsert:
assert revisions[0].revised_text == "Revision v2" assert revisions[0].revised_text == "Revision v2"
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_delete_source_from_job_context_removes_source_and_single_link(self, default_session_factory): async def test_delete_source_from_job_context_removes_source_and_single_link(self, default_session_factory, tmp_path):
documents = DocumentService(session_factory=default_session_factory) documents = DocumentService(session_factory=default_session_factory)
jobs = JobService(session_factory=default_session_factory) jobs = JobService(session_factory=default_session_factory)
transcriptions = TranscriptionService(session_factory=default_session_factory) transcriptions = TranscriptionService(session_factory=default_session_factory)
transcriptions.settings.upload_dir = tmp_path
document = Document(id=uuid4(), name="delete-source-success") document = Document(id=uuid4(), name="delete-source-success")
await documents.create_document(document=document) await documents.create_document(document=document)
@@ -104,12 +105,16 @@ class TestTranscriptionServiceRevisionUpsert:
job = Job(document_id=document.id, status=JobStatus.QUEUED) job = Job(document_id=document.id, status=JobStatus.QUEUED)
await jobs.create_job(job=job) await jobs.create_job(job=job)
stored_path = tmp_path / "documents" / str(document.id) / "delete.jpg"
stored_path.parent.mkdir(parents=True, exist_ok=True)
stored_path.write_bytes(b"data")
source = Source( source = Source(
document_id=document.id, document_id=document.id,
page_number=1, page_number=1,
upload_name="delete.jpg", upload_name="delete.jpg",
filename="delete.jpg", filename="delete.jpg",
file_path="uploads/delete.jpg", file_path=str(stored_path),
) )
async with transcriptions._session_scope() as session: async with transcriptions._session_scope() as session:
session.add(source) session.add(source)
@@ -122,6 +127,7 @@ class TestTranscriptionServiceRevisionUpsert:
with pytest.raises(TranscriptionNotFoundError): with pytest.raises(TranscriptionNotFoundError):
await transcriptions.read_source(source.id) await transcriptions.read_source(source.id)
assert not stored_path.exists()
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_delete_source_from_job_context_blocks_when_other_job_links_exist(self, default_session_factory): async def test_delete_source_from_job_context_blocks_when_other_job_links_exist(self, default_session_factory):
@@ -154,3 +160,60 @@ class TestTranscriptionServiceRevisionUpsert:
with pytest.raises(SourceDeleteBlockedError): with pytest.raises(SourceDeleteBlockedError):
await transcriptions.delete_source_from_job_context(job_id=job_one.id, source_id=source.id) await transcriptions.delete_source_from_job_context(job_id=job_one.id, source_id=source.id)
@pytest.mark.asyncio
async def test_delete_unlinked_source_succeeds(self, default_session_factory, tmp_path):
documents = DocumentService(session_factory=default_session_factory)
transcriptions = TranscriptionService(session_factory=default_session_factory)
transcriptions.settings.upload_dir = tmp_path
document = Document(id=uuid4(), name="delete-unlinked-source")
await documents.create_document(document=document)
stored_path = tmp_path / "documents" / str(document.id) / "orphan.jpg"
stored_path.parent.mkdir(parents=True, exist_ok=True)
stored_path.write_bytes(b"data")
source = Source(
document_id=document.id,
page_number=1,
upload_name="orphan.jpg",
filename="orphan.jpg",
file_path=str(stored_path),
)
await transcriptions.create_source(source=source)
await transcriptions.delete_unlinked_source(source_id=source.id)
with pytest.raises(TranscriptionNotFoundError):
await transcriptions.read_source(source.id)
assert not stored_path.exists()
@pytest.mark.asyncio
async def test_delete_unlinked_source_blocks_when_linked(self, default_session_factory):
documents = DocumentService(session_factory=default_session_factory)
jobs = JobService(session_factory=default_session_factory)
transcriptions = TranscriptionService(session_factory=default_session_factory)
document = Document(id=uuid4(), name="delete-unlinked-blocked")
await documents.create_document(document=document)
job = Job(document_id=document.id, status=JobStatus.QUEUED)
await jobs.create_job(job=job)
source = Source(
document_id=document.id,
page_number=1,
upload_name="linked.jpg",
filename="linked.jpg",
file_path="uploads/linked.jpg",
)
async with transcriptions._session_scope() as session:
session.add(source)
await session.flush()
session.add(JobSource(job_id=job.id, source_id=source.id, status=JobSourceStatus.PENDING))
await session.commit()
await session.refresh(source)
with pytest.raises(SourceDeleteBlockedError):
await transcriptions.delete_unlinked_source(source_id=source.id)
+29
View File
@@ -0,0 +1,29 @@
import asyncio
import logging
import pytest
from transcription.worker import run_worker_loop
@pytest.mark.asyncio
async def test_run_worker_loop_survives_process_next_exception(monkeypatch, caplog):
calls = 0
stop_event = asyncio.Event()
async def _fake_process_next_queued_job(*, session=None, session_factory=None):
nonlocal calls
_ = (session, session_factory)
calls += 1
if calls == 1:
raise RuntimeError("boom")
stop_event.set()
return False
monkeypatch.setattr("transcription.worker.process_next_queued_job", _fake_process_next_queued_job)
with caplog.at_level(logging.ERROR):
await run_worker_loop(stop_event=stop_event, poll_interval_seconds=0)
assert calls == 2
assert "Worker loop exception" in caplog.text
+28 -6
View File
@@ -86,7 +86,7 @@ class TestPageRendering:
assert "document links" in response.text.lower() assert "document links" in response.text.lower()
assert "Sources" in response.text assert "Sources" in response.text
assert "Jobs" in response.text assert "Jobs" in response.text
assert "Delete job" not in response.text assert "Delete Job" in response.text
def test_job_detail_page_rejects_invalid_id(self, app_client): def test_job_detail_page_rejects_invalid_id(self, app_client):
"""GET /ui/jobs/{job_id} shows validation feedback for malformed IDs.""" """GET /ui/jobs/{job_id} shows validation feedback for malformed IDs."""
@@ -105,19 +105,41 @@ class TestPageRendering:
assert response.status_code == 200 assert response.status_code == 200
assert "Job not found" in response.text assert "Job not found" in response.text
def test_job_detail_page_hides_delete_action(self, app_client, seed_job): def test_job_detail_page_shows_cancel_and_resubmit_when_queued(self, app_client, seed_job):
"""GET /ui/jobs/{job_id} does not expose job deletion controls in this revision.""" """GET /ui/jobs/{job_id} exposes cancel/resubmit controls for queued jobs."""
_, client = app_client _, client = app_client
job_id = seed_job( job_id = seed_job(
filename="no-revision.pdf", filename="no-revision.pdf",
status=JobStatus.TRANSCRIBED, status=JobStatus.QUEUED,
transcription_text="original text", transcription_text=None,
) )
response = client.get(f"/ui/jobs/{job_id}") response = client.get(f"/ui/jobs/{job_id}")
assert response.status_code == 200 assert response.status_code == 200
assert "Delete job" not in response.text assert "Cancel" in response.text
assert "Resubmit" in response.text
assert "Delete Job" in response.text
def test_job_cancel_page_renders_confirmation(self, app_client, seed_job):
_, client = app_client
job_id = seed_job(filename="cancel-ready.pdf", status=JobStatus.PROCESSING)
response = client.get(f"/ui/jobs/{job_id}/cancel")
assert response.status_code == 200
assert "Cancel Processing Job" in response.text
assert "Cancel job" in response.text
def test_job_resubmit_page_renders_confirmation(self, app_client, seed_job):
_, client = app_client
job_id = seed_job(filename="resubmit-ready.pdf", status=JobStatus.FAILED, transcription_text=None)
response = client.get(f"/ui/jobs/{job_id}/resubmit")
assert response.status_code == 200
assert "Resubmit Job" in response.text
assert "Resubmit now" in response.text
def test_job_delete_page_shows_confirmation_when_not_processing(self, app_client, seed_job): def test_job_delete_page_shows_confirmation_when_not_processing(self, app_client, seed_job):
_, client = app_client _, client = app_client
+2
View File
@@ -11,12 +11,14 @@ class TestPageRegistration:
"""Mounted UI routes respond successfully when the full app is created.""" """Mounted UI routes respond successfully when the full app is created."""
_, client = app_client _, client = app_client
homepage_response = client.get("/ui/homepage")
upload_response = client.get("/ui/upload", follow_redirects=False) upload_response = client.get("/ui/upload", follow_redirects=False)
documents_response = client.get("/ui/documents") documents_response = client.get("/ui/documents")
people_response = client.get("/ui/people") people_response = client.get("/ui/people")
sources_response = client.get("/ui/sources") sources_response = client.get("/ui/sources")
jobs_response = client.get("/ui/jobs") jobs_response = client.get("/ui/jobs")
assert homepage_response.status_code == 200
assert upload_response.status_code == 307 assert upload_response.status_code == 307
assert documents_response.status_code == 200 assert documents_response.status_code == 200
assert people_response.status_code == 200 assert people_response.status_code == 200
+3 -4
View File
@@ -206,7 +206,7 @@ class TestPeoplePageRendering:
assert "This action permanently deletes the person record." in response.text assert "This action permanently deletes the person record." in response.text
assert "Delete person permanently" in response.text assert "Delete person permanently" in response.text
def test_person_delete_page_shows_blocked_state_when_linked_documents_exist(self, app_client): def test_person_delete_page_warns_links_will_be_removed_when_linked_documents_exist(self, app_client):
_, client = app_client _, client = app_client
async def _seed_links() -> str: async def _seed_links() -> str:
@@ -233,6 +233,5 @@ class TestPeoplePageRendering:
response = client.get(f"/ui/people/{person_id}/delete") response = client.get(f"/ui/people/{person_id}/delete")
assert response.status_code == 200 assert response.status_code == 200
assert "Delete is blocked because linked documents exist." in response.text assert "This will also remove 1 linked document relationship(s)." in response.text
assert "Linked documents: 1" in response.text assert "Delete person permanently" in response.text
assert "Go to Documents" in response.text
+96
View File
@@ -9,6 +9,7 @@ from sqlmodel import select
from transcription.db import session_scope from transcription.db import session_scope
from transcription.db.models import Document from transcription.db.models import Document
from transcription.db.models import Job from transcription.db.models import Job
from transcription.db.models import JobStatus
from transcription.db.models import Source from transcription.db.models import Source
@@ -104,6 +105,23 @@ class TestSourcesPageRendering:
assert "Sources for Job" in response.text assert "Sources for Job" in response.text
assert "Back to Job" in response.text assert "Back to Job" in response.text
assert "job-page.png" in response.text assert "job-page.png" in response.text
assert "Job Source Status" in response.text
def test_sources_page_job_context_shows_job_source_status_and_error_detail(self, app_client, seed_job):
_, client = app_client
job_id = seed_job(
filename="job-failed-page.png",
status=JobStatus.FAILED,
transcription_text=None,
error_detail="Provider timed out",
)
response = client.get(f"/ui/sources?job_id={job_id}")
assert response.status_code == 200
assert "job-failed-page.png" in response.text
assert "failed" in response.text.lower()
assert "Provider timed out" in response.text
def test_source_detail_page_renders_preview_and_revision_box(self, app_client, seed_job): def test_source_detail_page_renders_preview_and_revision_box(self, app_client, seed_job):
_, client = app_client _, client = app_client
@@ -138,3 +156,81 @@ class TestSourcesPageRendering:
assert "human revision text" in response.text assert "human revision text" in response.text
assert "Page Number:" in response.text assert "Page Number:" in response.text
assert "Stored Filename:" in response.text assert "Stored Filename:" in response.text
assert "Delete Source" in response.text
def test_source_detail_page_displays_job_source_status_and_error_detail(self, app_client, seed_job):
_, client = app_client
job_id = seed_job(
filename="failed-source.png",
status=JobStatus.FAILED,
transcription_text=None,
error_detail="Provider timed out",
)
async def _get_source_id() -> str:
async with session_scope() as session:
job = await session.get(Job, job_id)
assert job is not None
source = (
await session.exec(select(Source).where(Source.document_id == job.document_id))
).first()
assert source is not None
return str(source.id)
source_id = asyncio.run(_get_source_id())
response = client.get(f"/ui/sources/{source_id}")
assert response.status_code == 200
assert "JOB SOURCE OUTCOMES" in response.text
assert "Status:" in response.text
assert "failed" in response.text.lower()
assert "Error Detail:" in response.text
assert "Provider timed out" in response.text
def test_source_delete_page_blocks_when_source_is_job_linked(self, app_client, seed_job):
_, client = app_client
job_id = seed_job(filename="linked-source.png", transcription_text="linked text")
async def _get_source_id() -> str:
async with session_scope() as session:
job = await session.get(Job, job_id)
assert job is not None
source = (
await session.exec(select(Source).where(Source.document_id == job.document_id))
).first()
assert source is not None
return str(source.id)
source_id = asyncio.run(_get_source_id())
response = client.get(f"/ui/sources/{source_id}/delete")
assert response.status_code == 200
assert "Delete Source Record" in response.text
assert "Delete is only available for unlinked sources." in response.text
def test_source_delete_page_allows_unlinked_source(self, app_client):
_, client = app_client
async def _seed_unlinked_source() -> str:
async with session_scope() as session:
document = Document(name="Unlinked Source Doc", document_type="memo")
session.add(document)
await session.flush()
source = Source(
document_id=document.id,
page_number=1,
upload_name="orphan-source.png",
filename="orphan-source.png",
file_path="/tmp/orphan-source.png",
)
session.add(source)
await session.commit()
return str(source.id)
source_id = asyncio.run(_seed_unlinked_source())
response = client.get(f"/ui/sources/{source_id}/delete")
assert response.status_code == 200
assert "Delete Source Record" in response.text
assert "Delete source permanently" in response.text
assert "Delete is only available for unlinked sources." not in response.text
+23 -4
View File
@@ -13,15 +13,34 @@ class TestPageRendering:
response = client.get("/", follow_redirects=False) response = client.get("/", follow_redirects=False)
assert response.status_code == 307 assert response.status_code == 307
assert response.headers["location"] == "/ui" assert response.headers["location"] == "/ui/homepage"
def test_ui_redirects_to_documents(self, app_client): def test_ui_redirects_to_homepage(self, app_client):
"""GET /ui redirects to the documents page.""" """GET /ui redirects to the homepage."""
_, client = app_client _, client = app_client
response = client.get("/ui", follow_redirects=False) response = client.get("/ui", follow_redirects=False)
assert response.status_code == 307 assert response.status_code == 307
assert response.headers["location"] == "/ui/documents" assert response.headers["location"] == "/ui/homepage"
def test_homepage_page_renders(self, app_client):
"""GET /ui/homepage renders the homepage page."""
_, client = app_client
response = client.get("/ui/homepage")
assert response.status_code == 200
assert "Home" in response.text
assert "Edit Home Page" in response.text
assert '/ui/homepage' in response.text
def test_homepage_edit_page_renders(self, app_client):
"""GET /ui/homepage/edit renders the edit page."""
_, client = app_client
response = client.get("/ui/homepage/edit")
assert response.status_code == 200
assert "Edit Home Page" in response.text
assert "Homepage markdown" in response.text
def test_upload_page_renders_expected_controls(self, app_client): def test_upload_page_renders_expected_controls(self, app_client):
"""GET /ui/upload redirects to the job-create flow.""" """GET /ui/upload redirects to the job-create flow."""