diff --git a/.github/instructions/services.instructions.md b/.github/instructions/services.instructions.md index 981b01e..a668b44 100644 --- a/.github/instructions/services.instructions.md +++ b/.github/instructions/services.instructions.md @@ -7,29 +7,89 @@ applyTo: 'src/transcription/services/*.py' ## Structure -- Project core data models defined in [models](../../src/transcription/db/models.py) -- 1 service class per data model -- Only services directly interact with the database, and only through async methods -- Services are completely independent of one another. Any operation that needs to use more than a single service, which is most of them, needs to have a separate orchestration function. +- Project core data models are defined in [models](../../src/transcription/db/models.py) +- One service class per **aggregate**, not per table. An aggregate is a root model plus + the models that have no independent lifecycle of their own. `DocumentType` has no + meaning without `Document`, so it belongs to `DocumentService`; it does not get its + own service. Splitting per table produces services that must reach across each other + for every real operation, which is what line 13 forbids. +- Only services interact with the database, and only through async methods. +- **A service module must not import another service module.** This is enforced by + [test_service_boundaries](../../tests/test_service_boundaries.py). Shared types go in a + neutral module that defines no service class (see [errors](../../src/transcription/services/errors.py)). +- Not every module in this package is a service. Helper modules that define no `*Service` + class (`base`, `errors`, `normalization`, `prompts`, `quality`, `media_storage`, + `source_media`) are free-function modules and are exempt from the service rules below. + +## Model Ownership + +Every model has exactly one owning service. The owner defines that model's invariants and +is the only service that may **create or delete** its rows. + +| Model | Owner | +| --- | --- | +| `Document`, `DocumentType` | `DocumentService` | +| `Source`, `JobSource` | `SourceService` | +| `Job` | `JobService` | +| `Person`, `PersonRole`, `DocumentPerson` | `PeopleService` | +| `ExecutionAttempt` | `EvidenceService` | + +### Junction tables + +A junction table is owned by the service that **creates and deletes its rows** — its +lifecycle owner. The service on the other side may read through the junction (via +`selectinload`) but must not create rows in it. + +- `document_person` -> `PeopleService`. Every write is there; `DocumentService` only + eager-loads through it. +- `job_source` -> `SourceService`, which creates the row, records each page's outcome, + and deletes it. + +Two consequences follow, and both are deliberate: + +- **Cascade deletion is not a violation.** A service deleting the aggregate root it owns + may delete junction rows referencing that root, because they cannot outlive it + (`JobService.delete_job_with_guardrails`). +- **Ownership governs creation and deletion, not every state transition.** `job_source` is + both a link and the transcription work queue. `JobService.cancel_job` and + `resubmit_failed_sources` transition `job_source.status` across a whole job, because that + transition is a Job lifecycle event, not a per-page outcome. They create and delete + nothing. + +`EvidenceService.promote_machine_attempt` writes two fields on `Source` +(`preferred_execution_attempt_id`, `raw_transcription`). This is allowed on the same +principle: selecting which attempt a Source presents is an evidence decision that happens +to land on `Source`. It is scoped to those two projection fields. + +If a new operation cannot be expressed within one owner, it belongs in an orchestration +module, not in a cross-service import. ## Error Handling -- Service-specific errors defined at the top of the respective module and inherit from `AppError` -- Use a context manager for large `try/except` blocks like `handle_transcription_errors` in [sources](../../src/transcription/services/sources.py) +- Errors used by a single service are defined at the top of that module and inherit from `AppError`. +- Errors shared by more than one service go in [errors](../../src/transcription/services/errors.py), + which defines no service class and is therefore importable by any of them. +- Use a context manager for large `try/except` blocks, like `handle_transcription_errors` in + [sources](../../src/transcription/services/sources.py). ## Checklist - [ ] Uses `ServiceBase` for common logic -- [ ] CRUD methods created at the top -- [ ] Session kwarg for `AsyncSession` to pass in a session object to each method -- [ ] Services use `self._session_scope` in their methods to pass the session thru. - - Multiple operations on the same object(s) require sharing a session between all the methods used. +- [ ] Session kwarg for `AsyncSession` to pass a session object into each method +- [ ] Services use `self._session_scope` in their methods to pass the session through + - Multiple operations on the same object(s) require sharing a session between all the methods used +- [ ] Every model the module touches is either owned by it or reached read-only ## CRUD Methods -- Create, read, update, and delete, created in that order -- Name format `_`, for example `create_document` or `update_job` -- All services must define these 4 methods first, and in that order +- Name format `_`, for example `create_document` or `update_job`. +- Where a service exposes create/read/update/delete for its root model, define them at the + top of the class in that order, before derived reads and workflow helpers. +- Not every aggregate needs all four. `ExecutionAttempt` is append-only evidence written by + `workflows.py`, so `EvidenceService` deliberately exposes reads and no create or delete. + Do not add unused CRUD methods to satisfy symmetry. +- `RegistryService` is generic across small lookup models and uses `_entry` + naming instead. ## Transaction Finalization @@ -74,4 +134,11 @@ Separation of concerns: # Service Composition -Some operations, like uploading a picutre, require modifications to multiple tables, which can be done by composing methods from the service object into a separate function. +A service method may read across models it does not own, using eager loads from its own +aggregate root. What it may not do is import another service. + +Operations that must **write** models owned by more than one service — uploading a picture, +for example — are composed in an orchestration module +([store](../../src/transcription/services/store.py), +[workflows](../../src/transcription/services/workflows.py)). Orchestration modules define no +service class, may import any service, and own the commit boundary. diff --git a/src/transcription/services/__init__.py b/src/transcription/services/__init__.py index 256f49a..b9b4b1b 100644 --- a/src/transcription/services/__init__.py +++ b/src/transcription/services/__init__.py @@ -9,12 +9,21 @@ from sqlmodel.ext.asyncio.session import AsyncSession from ..config import Settings from .documents import DocumentService +from .evidence import EvidenceService from .jobs import JobService from .people import PeopleService from .prompts import PromptStore from .sources import SourceService -__all__ = ["DocumentService", "JobService", "PeopleService", "PromptStore", "ServiceBundle", "SourceService"] +__all__ = [ + "DocumentService", + "EvidenceService", + "JobService", + "PeopleService", + "PromptStore", + "ServiceBundle", + "SourceService", +] @dataclass(frozen=True, slots=True) @@ -25,6 +34,7 @@ class ServiceBundle: sources: SourceService = field(default_factory=SourceService) jobs: JobService = field(default_factory=JobService) people: PeopleService = field(default_factory=PeopleService) + evidence: EvidenceService = field(default_factory=EvidenceService) @classmethod def from_session_factory( @@ -41,6 +51,7 @@ class ServiceBundle: sources=SourceService(session_factory=session_factory, settings=settings), jobs=JobService(session_factory=session_factory, settings=settings), people=PeopleService(session_factory=session_factory, settings=settings), + evidence=EvidenceService(session_factory=session_factory, settings=settings), ) async def aclose(self) -> None: diff --git a/src/transcription/services/errors.py b/src/transcription/services/errors.py new file mode 100644 index 0000000..22fb166 --- /dev/null +++ b/src/transcription/services/errors.py @@ -0,0 +1,32 @@ +"""Error vocabulary shared across the source, evidence, and prompt services. + +These live in a neutral module rather than in the service that raises them +because more than one service raises them, and ``services.instructions.md`` +forbids a service module from importing a sibling. Orchestration modules and +the UI import from here, so the exception a caller catches does not change when +an operation moves between services. +""" + +from __future__ import annotations + +from transcription.errors import AppError + + +class PromptLoadError(AppError): + """Raised when prompt artifacts cannot be loaded safely.""" + + +class TranscriptionError(AppError): + """Raised when transcription execution fails.""" + + +class TranscriptionNotFoundError(TranscriptionError): + """Raised when a transcription-related resource is not found.""" + + +class SourceDeleteBlockedError(TranscriptionError): + """Raised when source deletion is blocked by dependency policy.""" + + +class CandidatePromotionError(TranscriptionError): + """Raised when a machine attempt cannot be selected for its Source.""" diff --git a/src/transcription/services/evidence.py b/src/transcription/services/evidence.py new file mode 100644 index 0000000..687cf59 --- /dev/null +++ b/src/transcription/services/evidence.py @@ -0,0 +1,210 @@ +"""Read and export the immutable execution evidence trail. + +``ExecutionAttempt`` is append-only: one row per provider call, written once by +the transcription workflow and never updated. Everything here is therefore a +read, a projection, or an export, with one exception - ``promote_machine_attempt`` +selects which attempt a ``Source`` presents, which is an evidence decision even +though the write lands on ``Source``. +""" + +from __future__ import annotations + +import base64 +import hashlib +from collections.abc import Sequence +from dataclasses import dataclass +from uuid import UUID + +from pydantic import JsonValue +from sqlalchemy import inspect as sqlalchemy_inspect +from sqlmodel import col +from sqlmodel import select +from sqlmodel.ext.asyncio.session import AsyncSession + +from transcription.db.models import ExecutionAttempt +from transcription.db.models import JobSourceStatus +from transcription.db.models import Source +from transcription.errors import ErrorCategory + +from ..db.loading import defer +from .base import ServiceBase +from .errors import CandidatePromotionError +from .errors import TranscriptionNotFoundError + + +@dataclass(frozen=True, slots=True) +class LatestExecutionAttempt: + """One execution attempt plus the loader facts a caller needs to render it.""" + + attempt: ExecutionAttempt + transport_body_deferred: bool + + +class EvidenceService(ServiceBase): + """Read, project, and export execution attempt evidence.""" + + async def read_latest_execution_attempt( + self, + *, + job_source_id: UUID, + session: AsyncSession | None = None, + ) -> LatestExecutionAttempt | None: + """Read only the latest immutable attempt for one compatibility projection. + + The transport body is deferred because it can be arbitrarily large; the + returned read model reports that as a plain flag so callers never have to + inspect ORM loader state. + """ + async with self._session_scope(session) as _session: + query = ( + select(ExecutionAttempt) + .options(defer(ExecutionAttempt.transport_body)) + .where(ExecutionAttempt.job_source_id == job_source_id) + .order_by( + col(ExecutionAttempt.attempt_number).desc(), + col(ExecutionAttempt.id).desc(), + ) + .limit(1) + ) + attempt = (await _session.exec(query)).first() + if attempt is None: + return None + deferred = "transport_body" in sqlalchemy_inspect(attempt).unloaded + return LatestExecutionAttempt(attempt=attempt, transport_body_deferred=deferred) + + async def list_execution_attempts( + self, + *, + source_id: UUID | None = None, + job_id: UUID | None = None, + session: AsyncSession | None = None, + ) -> Sequence[ExecutionAttempt]: + """List immutable execution evidence in stable attempt order.""" + async with self._session_scope(session) as _session: + query = select(ExecutionAttempt) + if source_id is not None: + query = query.where(ExecutionAttempt.source_id == source_id) + if job_id is not None: + query = query.where(ExecutionAttempt.job_id == job_id) + query = query.order_by( + col(ExecutionAttempt.job_id), + col(ExecutionAttempt.source_id), + col(ExecutionAttempt.attempt_number), + col(ExecutionAttempt.id), + ) + return (await _session.exec(query)).all() + + async def promote_machine_attempt( + self, + *, + source_id: UUID, + execution_attempt_id: UUID, + session: AsyncSession | None = None, + ) -> Source: + """Atomically select one successful machine attempt as the Source projection.""" + async with self._session_scope(session) as _session: + source = await self._read_source( + session=_session, + source_id=source_id, + suggestion="Refresh Source Detail and retry.", + ) + attempt = await _session.get(ExecutionAttempt, execution_attempt_id) + if ( + attempt is None + or attempt.source_id != source_id + or attempt.status != JobSourceStatus.TRANSCRIBED + or not attempt.raw_transcription + ): + raise CandidatePromotionError( + "Only a successful transcription attempt belonging to this Source can be selected", + category=ErrorCategory.VALIDATION, + suggestion="Select an available successful candidate from Source Detail.", + ) + source.preferred_execution_attempt_id = attempt.id + source.raw_transcription = attempt.raw_transcription + await self._finalize(session=_session, caller_session=session, refresh=(source,)) + return source + + async def build_evidence_export( + self, + *, + source_id: UUID, + session: AsyncSession | None = None, + ) -> dict[str, JsonValue]: + """Build a versioned, source-reference-only evidence export.""" + async with self._session_scope(session) as _session: + source = await self._read_source(session=_session, source_id=source_id) + attempts = list(await self.list_execution_attempts(source_id=source_id, session=_session)) + + attempt_payloads = [ + { + "id": str(attempt.id), + "job_id": str(attempt.job_id), + "source_id": str(attempt.source_id), + "attempt_number": attempt.attempt_number, + "status": attempt.status.value, + "provider": attempt.provider, + "model": attempt.model, + "request_manifest": attempt.request_manifest, + "request_manifest_sha256": attempt.request_manifest_sha256, + "request_manifest_schema_version": attempt.request_manifest_schema_version, + "transport": { + "response_received": attempt.response_received, + "status_code": attempt.transport_status_code, + "body_base64": ( + base64.b64encode(attempt.transport_body).decode("ascii") + if attempt.transport_body is not None + else None + ), + "body_sha256": ( + hashlib.sha256(attempt.transport_body).hexdigest() + if attempt.transport_body is not None + else None + ), + "content_type": attempt.transport_content_type, + "content_encoding": attempt.transport_content_encoding, + "safe_headers": attempt.transport_safe_headers, + "request_id": attempt.router_request_id, + "generation_id": attempt.router_generation_id, + }, + "sdk_response_snapshot": attempt.sdk_response_snapshot, + "normalized_metadata": attempt.normalized_metadata, + "software_context": attempt.software_context, + "raw_transcription": attempt.raw_transcription, + "error_category": attempt.error_category, + "error_detail": attempt.error_detail, + "failure_phase": attempt.failure_phase, + "started_at": attempt.started_at.isoformat(), + "finished_at": attempt.finished_at.isoformat(), + "duration_ms": attempt.duration_ms, + } + for attempt in attempts + ] + return { + "schema_name": "transcription.evidence-export", + "schema_version": "1", + "source": { + "id": str(source.id), + "digest_sha256": source.file_hash, + "byte_size": source.file_size_bytes, + "page_number": source.page_number, + "upload_name": source.upload_name, + }, + "attempts": attempt_payloads, + } + + async def _read_source( + self, + *, + session: AsyncSession, + source_id: UUID, + suggestion: str = "Verify the source id and retry.", + ) -> Source: + return await self._get_or_raise( + Source, + source_id, + session=session, + error=TranscriptionNotFoundError, + noun="Source", + suggestion=suggestion, + ) diff --git a/src/transcription/services/sources.py b/src/transcription/services/sources.py index f1c50c0..779da08 100644 --- a/src/transcription/services/sources.py +++ b/src/transcription/services/sources.py @@ -3,7 +3,6 @@ from __future__ import annotations import asyncio -import base64 import hashlib import logging from collections.abc import Sequence @@ -22,7 +21,6 @@ from pydantic import JsonValue from pydantic import TypeAdapter from pydantic import ValidationError from sqlalchemy import func -from sqlalchemy import inspect as sqlalchemy_inspect from sqlalchemy import literal from sqlalchemy import tuple_ from sqlalchemy.ext.asyncio import async_sessionmaker @@ -37,7 +35,6 @@ from transcription.db.models import Job from transcription.db.models import JobSource from transcription.db.models import JobSourceStatus from transcription.db.models import Source -from transcription.errors import AppError from transcription.errors import ErrorCategory from transcription.providers import ProviderAuthError from transcription.providers import ProviderError @@ -50,10 +47,13 @@ from transcription.providers import TranscriptionResult from transcription.providers import TransportEvidence from transcription.providers import get_transcription_provider -from ..db.loading import defer from ..db.loading import orm_attribute from ..db.loading import selectinload from .base import ServiceBase +from .errors import PromptLoadError +from .errors import SourceDeleteBlockedError +from .errors import TranscriptionError +from .errors import TranscriptionNotFoundError from .source_media import lookup_source_mime_type from .source_media import supported_source_formats @@ -76,26 +76,6 @@ class PromptExecution(BaseModel): top_p: float | None = Field(ge=0.0, le=1.0) -class PromptLoadError(AppError): - """Raised when prompt artifacts cannot be loaded safely.""" - - -class TranscriptionError(AppError): - """Raised when transcription execution fails.""" - - -class TranscriptionNotFoundError(TranscriptionError): - """Raised when a transcription-related resource is not found.""" - - -class SourceDeleteBlockedError(TranscriptionError): - """Raised when source deletion is blocked by dependency policy.""" - - -class CandidatePromotionError(TranscriptionError): - """Raised when a machine attempt cannot be selected for its Source.""" - - @dataclass(frozen=True, slots=True) class SourceNavigation: """Adjacent Source identifiers within one ordered Document.""" @@ -128,14 +108,6 @@ def build_provider_input(source: Source) -> ProviderInput: ) -@dataclass(frozen=True, slots=True) -class LatestExecutionAttempt: - """One execution attempt plus the loader facts a caller needs to render it.""" - - attempt: ExecutionAttempt - transport_body_deferred: bool - - class SourceService(ServiceBase): """Manage source records, media payloads, revisions, and page execution output.""" @@ -214,35 +186,6 @@ class SourceService(ServiceBase): ) return source - async def read_latest_execution_attempt( - self, - *, - job_source_id: UUID, - session: AsyncSession | None = None, - ) -> LatestExecutionAttempt | None: - """Read only the latest immutable attempt for one compatibility projection. - - The transport body is deferred because it can be arbitrarily large; the - returned read model reports that as a plain flag so callers never have to - inspect ORM loader state. - """ - async with self._session_scope(session) as _session: - query = ( - select(ExecutionAttempt) - .options(defer(ExecutionAttempt.transport_body)) - .where(ExecutionAttempt.job_source_id == job_source_id) - .order_by( - col(ExecutionAttempt.attempt_number).desc(), - col(ExecutionAttempt.id).desc(), - ) - .limit(1) - ) - attempt = (await _session.exec(query)).first() - if attempt is None: - return None - deferred = "transport_body" in sqlalchemy_inspect(attempt).unloaded - return LatestExecutionAttempt(attempt=attempt, transport_body_deferred=deferred) - async def read_source_navigation( self, source_id: UUID, @@ -640,127 +583,6 @@ class SourceService(ServiceBase): await self._finalize(session=_session, caller_session=session, refresh=(job, source, job_source, attempt)) return job_source - async def promote_machine_attempt( - self, - *, - source_id: UUID, - execution_attempt_id: UUID, - session: AsyncSession | None = None, - ) -> Source: - """Atomically select one successful machine attempt as the Source projection.""" - async with self._session_scope(session) as _session: - source = await self._read_source( - session=_session, - source_id=source_id, - suggestion="Refresh Source Detail and retry.", - ) - attempt = await _session.get(ExecutionAttempt, execution_attempt_id) - if ( - attempt is None - or attempt.source_id != source_id - or attempt.status != JobSourceStatus.TRANSCRIBED - or not attempt.raw_transcription - ): - raise CandidatePromotionError( - "Only a successful transcription attempt belonging to this Source can be selected", - category=ErrorCategory.VALIDATION, - suggestion="Select an available successful candidate from Source Detail.", - ) - source.preferred_execution_attempt_id = attempt.id - source.raw_transcription = attempt.raw_transcription - await self._finalize(session=_session, caller_session=session, refresh=(source,)) - return source - - async def list_execution_attempts( - self, - *, - source_id: UUID | None = None, - job_id: UUID | None = None, - session: AsyncSession | None = None, - ) -> Sequence[ExecutionAttempt]: - """List immutable execution evidence in stable attempt order.""" - async with self._session_scope(session) as _session: - query = select(ExecutionAttempt) - if source_id is not None: - query = query.where(ExecutionAttempt.source_id == source_id) - if job_id is not None: - query = query.where(ExecutionAttempt.job_id == job_id) - query = query.order_by( - col(ExecutionAttempt.job_id), - col(ExecutionAttempt.source_id), - col(ExecutionAttempt.attempt_number), - col(ExecutionAttempt.id), - ) - return (await _session.exec(query)).all() - - async def build_evidence_export( - self, - *, - source_id: UUID, - session: AsyncSession | None = None, - ) -> dict[str, JsonValue]: - """Build a versioned, source-reference-only evidence export.""" - async with self._session_scope(session) as _session: - source = await self._read_source(session=_session, source_id=source_id) - attempts = list(await self.list_execution_attempts(source_id=source_id, session=_session)) - - attempt_payloads = [ - { - "id": str(attempt.id), - "job_id": str(attempt.job_id), - "source_id": str(attempt.source_id), - "attempt_number": attempt.attempt_number, - "status": attempt.status.value, - "provider": attempt.provider, - "model": attempt.model, - "request_manifest": attempt.request_manifest, - "request_manifest_sha256": attempt.request_manifest_sha256, - "request_manifest_schema_version": attempt.request_manifest_schema_version, - "transport": { - "response_received": attempt.response_received, - "status_code": attempt.transport_status_code, - "body_base64": ( - base64.b64encode(attempt.transport_body).decode("ascii") - if attempt.transport_body is not None - else None - ), - "body_sha256": ( - hashlib.sha256(attempt.transport_body).hexdigest() - if attempt.transport_body is not None - else None - ), - "content_type": attempt.transport_content_type, - "content_encoding": attempt.transport_content_encoding, - "safe_headers": attempt.transport_safe_headers, - "request_id": attempt.router_request_id, - "generation_id": attempt.router_generation_id, - }, - "sdk_response_snapshot": attempt.sdk_response_snapshot, - "normalized_metadata": attempt.normalized_metadata, - "software_context": attempt.software_context, - "raw_transcription": attempt.raw_transcription, - "error_category": attempt.error_category, - "error_detail": attempt.error_detail, - "failure_phase": attempt.failure_phase, - "started_at": attempt.started_at.isoformat(), - "finished_at": attempt.finished_at.isoformat(), - "duration_ms": attempt.duration_ms, - } - for attempt in attempts - ] - return { - "schema_name": "transcription.evidence-export", - "schema_version": "1", - "source": { - "id": str(source.id), - "digest_sha256": source.file_hash, - "byte_size": source.file_size_bytes, - "page_number": source.page_number, - "upload_name": source.upload_name, - }, - "attempts": attempt_payloads, - } - async def upsert_revision_for_source( self, *, diff --git a/src/transcription/services/store.py b/src/transcription/services/store.py index 29a298f..26a803d 100644 --- a/src/transcription/services/store.py +++ b/src/transcription/services/store.py @@ -23,10 +23,10 @@ from ..db.models import JobSourceStatus from ..db.models import Source from ..db.session import SessionFactory from ..db.session import session_scope +from .errors import TranscriptionError from .media_storage import build_stored_filename from .media_storage import write_media_bytes from .normalization import normalize_orientation_async -from .sources import TranscriptionError from .sources import build_prompt_execution from .sources import source_mime_type from .sources import validate_source_content diff --git a/src/transcription/ui/pages/sources_page.py b/src/transcription/ui/pages/sources_page.py index fb7b983..271a2fc 100644 --- a/src/transcription/ui/pages/sources_page.py +++ b/src/transcription/ui/pages/sources_page.py @@ -14,10 +14,11 @@ from transcription.db.models import ExecutionAttempt from transcription.db.models import JobSource from transcription.db.models import JobSourceStatus from transcription.db.models import Source -from transcription.services.sources import LatestExecutionAttempt -from transcription.services.sources import SourceDeleteBlockedError +from transcription.services.errors import SourceDeleteBlockedError +from transcription.services.errors import TranscriptionNotFoundError +from transcription.services.evidence import EvidenceService +from transcription.services.evidence import LatestExecutionAttempt from transcription.services.sources import SourceService -from transcription.services.sources import TranscriptionNotFoundError from transcription.ui.components.app_shell import render_navigation_header from transcription.ui.components.cards import archival_card from transcription.ui.components.confirm_delete import render_delete_actions @@ -112,6 +113,7 @@ def register_page() -> None: # noqa: PLR0915 @ui.page("/sources/{source_id}") async def source_detail_page(source_id: str, request: Request, session_factory: SessionFactoryDep) -> None: sources_service = SourceService(session_factory=session_factory) + evidence_service = EvidenceService(session_factory=session_factory) render_navigation_header(current_path="/sources") parsed_source_id = parsed_record_id(source_id, noun="Source") @@ -123,11 +125,11 @@ def register_page() -> None: # noqa: PLR0915 navigation = await sources_service.read_source_navigation(parsed_source_id) latest_job_source = source.latest_job_source latest_attempt = ( - await sources_service.read_latest_execution_attempt(job_source_id=latest_job_source.id) + await evidence_service.read_latest_execution_attempt(job_source_id=latest_job_source.id) if latest_job_source is not None else None ) - attempts = list(await sources_service.list_execution_attempts(source_id=parsed_source_id)) + attempts = list(await evidence_service.list_execution_attempts(source_id=parsed_source_id)) except TranscriptionNotFoundError: render_record_not_found("Source") return @@ -158,7 +160,7 @@ def register_page() -> None: # noqa: PLR0915 "Export Evidence", on_click=lambda: _download_evidence( source_id=source.id, - sources_service=sources_service, + evidence_service=evidence_service, ), icon="download", ).props("flat") @@ -185,7 +187,7 @@ def register_page() -> None: # noqa: PLR0915 _render_machine_candidates( source=source, attempts=attempts, - sources_service=sources_service, + evidence_service=evidence_service, ) _render_source_metadata_column( source=source, @@ -409,9 +411,9 @@ def _transport_display(latest_attempt: LatestExecutionAttempt) -> dict[str, obje } -async def _download_evidence(*, source_id: UUID, sources_service: SourceService) -> None: +async def _download_evidence(*, source_id: UUID, evidence_service: EvidenceService) -> None: try: - payload = await sources_service.build_evidence_export(source_id=source_id) + payload = await evidence_service.build_evidence_export(source_id=source_id) except Exception as exc: # noqa: BLE001 show_error(exc, title="Export failed", operation="sources.evidence_export") return @@ -521,7 +523,7 @@ def _render_machine_candidates( *, source: Source, attempts: list[ExecutionAttempt], - sources_service: SourceService, + evidence_service: EvidenceService, ) -> None: successful = [ attempt @@ -581,7 +583,7 @@ def _render_machine_candidates( async def promote(candidate_id: UUID = attempt.id) -> None: try: - await sources_service.promote_machine_attempt( + await evidence_service.promote_machine_attempt( source_id=source.id, execution_attempt_id=candidate_id, ) diff --git a/tests/services/test_job_service.py b/tests/services/test_job_service.py index a0dddf0..ad7f9bb 100644 --- a/tests/services/test_job_service.py +++ b/tests/services/test_job_service.py @@ -13,6 +13,7 @@ from transcription.db.models import JobSourceStatus from transcription.db.models import JobStatus from transcription.db.models import Source from transcription.services.documents import DocumentService +from transcription.services.evidence import EvidenceService from transcription.services.jobs import JobCancelBlockedError from transcription.services.jobs import JobDeleteBlockedError from transcription.services.jobs import JobNotFoundError @@ -274,6 +275,7 @@ class TestJobService: document_service: DocumentService, ): source_service = SourceService(session_factory=job_service.session_factory) + evidence_service = EvidenceService(session_factory=job_service.session_factory) document = await document_service.create_document(Document(name="evidence-delete-doc")) job = await job_service.create_job(Job(document_id=document.id, status=JobStatus.FAILED)) source = await source_service.create_source( @@ -299,7 +301,7 @@ class TestJobService: with pytest.raises(JobNotFoundError): await job_service.read_job(job_id=job.id) - assert await source_service.list_execution_attempts(job_id=job.id) == [] + assert await evidence_service.list_execution_attempts(job_id=job.id) == [] assert (await source_service.read_source(source.id)).id == source.id @pytest.mark.asyncio diff --git a/tests/services/test_transcription_service.py b/tests/services/test_transcription_service.py index 18e278d..8e72629 100644 --- a/tests/services/test_transcription_service.py +++ b/tests/services/test_transcription_service.py @@ -13,10 +13,10 @@ from transcription.db.models import JobSourceStatus from transcription.db.models import JobStatus from transcription.db.models import Source from transcription.services.documents import DocumentService +from transcription.services.errors import SourceDeleteBlockedError +from transcription.services.errors import TranscriptionNotFoundError from transcription.services.jobs import JobService -from transcription.services.sources import SourceDeleteBlockedError from transcription.services.sources import SourceService -from transcription.services.sources import TranscriptionNotFoundError @pytest.mark.integration diff --git a/tests/services/test_v2_crud.py b/tests/services/test_v2_crud.py index 90101aa..10960ba 100644 --- a/tests/services/test_v2_crud.py +++ b/tests/services/test_v2_crud.py @@ -13,10 +13,11 @@ from transcription.db.models import Source from transcription.errors import ErrorCategory from transcription.services.documents import DocumentDeleteBlockedError from transcription.services.documents import DocumentService +from transcription.services.errors import SourceDeleteBlockedError +from transcription.services.evidence import EvidenceService from transcription.services.jobs import JobService from transcription.services.people import PeopleError from transcription.services.people import PeopleService -from transcription.services.sources import SourceDeleteBlockedError from transcription.services.sources import SourceService @@ -324,7 +325,8 @@ async def test_update_job_source_transcription_persists_provider_json_payloads(d assert len(stored_rows) == 1 assert stored_rows[0].status == JobSourceStatus.TRANSCRIBED - attempt = await transcriptions.read_latest_execution_attempt(job_source_id=stored_rows[0].id) + evidence = EvidenceService(session_factory=transcriptions.session_factory) + attempt = await evidence.read_latest_execution_attempt(job_source_id=stored_rows[0].id) assert attempt is not None assert attempt.attempt.raw_transcription == "provider transcript" assert attempt.attempt.normalized_metadata == metadata diff --git a/tests/services/test_v45_candidates.py b/tests/services/test_v45_candidates.py index 26e271c..a54b61f 100644 --- a/tests/services/test_v45_candidates.py +++ b/tests/services/test_v45_candidates.py @@ -11,19 +11,12 @@ from transcription.db.models import JobPurpose from transcription.db.models import JobSource from transcription.db.models import Source from transcription.services import ServiceBundle -from transcription.services.documents import DocumentService -from transcription.services.jobs import JobService -from transcription.services.sources import CandidatePromotionError -from transcription.services.sources import SourceService +from transcription.services.errors import CandidatePromotionError from transcription.services.workflows import create_source_retranscription_job def _services(default_session_factory, settings: Settings) -> ServiceBundle: - return ServiceBundle( - documents=DocumentService(session_factory=default_session_factory, settings=settings), - jobs=JobService(session_factory=default_session_factory, settings=settings), - sources=SourceService(session_factory=default_session_factory, settings=settings), - ) + return ServiceBundle.from_session_factory(default_session_factory, settings=settings) async def _seed_source(services: ServiceBundle) -> Source: @@ -71,14 +64,14 @@ async def test_first_success_is_preferred_and_later_success_remains_candidate(de model="model-b", ) unchanged = await services.sources.read_source(source.id) - attempts = await services.sources.list_execution_attempts(source_id=source.id) + attempts = await services.evidence.list_execution_attempts(source_id=source.id) assert unchanged.raw_transcription == "first result" assert unchanged.preferred_execution_attempt_id == first_attempt_id assert {attempt.raw_transcription for attempt in attempts} == {"first result", "candidate result"} candidate = next(attempt for attempt in attempts if attempt.raw_transcription == "candidate result") - promoted = await services.sources.promote_machine_attempt( + promoted = await services.evidence.promote_machine_attempt( source_id=source.id, execution_attempt_id=candidate.id, ) @@ -95,7 +88,7 @@ async def test_promotion_rejects_unrelated_attempt(default_session_factory): source = await _seed_source(services) with pytest.raises(CandidatePromotionError): - await services.sources.promote_machine_attempt( + await services.evidence.promote_machine_attempt( source_id=source.id, execution_attempt_id=uuid4(), ) diff --git a/tests/services/test_workflows_reliability.py b/tests/services/test_workflows_reliability.py index 9efdc44..e88cb84 100644 --- a/tests/services/test_workflows_reliability.py +++ b/tests/services/test_workflows_reliability.py @@ -109,12 +109,7 @@ class TestWorkflowReliability: default_session_factory, monkeypatch, ): - services = ServiceBundle( - documents=ServiceBundle().documents.__class__(session_factory=default_session_factory), - jobs=ServiceBundle().jobs.__class__(session_factory=default_session_factory), - sources=ServiceBundle().sources.__class__(session_factory=default_session_factory), - people=ServiceBundle().people.__class__(session_factory=default_session_factory), - ) + services = ServiceBundle.from_session_factory(default_session_factory) async with services.jobs._session_scope() as session: document = Document(id=uuid4(), name="durability-doc") session.add(document) @@ -155,7 +150,7 @@ class TestWorkflowReliability: task = asyncio.create_task(process_queued_job(job=loaded, services=services)) await asyncio.wait_for(second_started.wait(), timeout=2) - attempts = await services.sources.list_execution_attempts(job_id=job.id) + attempts = await services.evidence.list_execution_attempts(job_id=job.id) assert len(attempts) == 1 assert attempts[0].raw_transcription == "page 1" diff --git a/tests/test_prompts.py b/tests/test_prompts.py index 15af84b..a1be55e 100644 --- a/tests/test_prompts.py +++ b/tests/test_prompts.py @@ -7,8 +7,8 @@ import pytest from pydantic import ValidationError from transcription.config import Settings +from transcription.services.errors import PromptLoadError from transcription.services.sources import PromptExecution -from transcription.services.sources import PromptLoadError from transcription.services.sources import build_prompt_execution from transcription.services.sources import load_prompt_text diff --git a/tests/test_service_boundaries.py b/tests/test_service_boundaries.py index 815d4ce..e38615f 100644 --- a/tests/test_service_boundaries.py +++ b/tests/test_service_boundaries.py @@ -1,9 +1,10 @@ """Structural rules for the services package. -`.github/instructions/services.instructions.md:13` requires that service classes -stay independent of one another. Shared behavior belongs in a neutral module -(`base.py`, `registry.py`, `source_media.py`, `media_storage.py`), and any -operation spanning two services belongs in an orchestration module. +The "Structure" section of `.github/instructions/services.instructions.md` requires +that service modules stay independent of one another. Shared behavior belongs in a +neutral module that defines no service class (`base.py`, `errors.py`, `registry.py`, +`source_media.py`, `media_storage.py`), and any operation that writes models owned by +two services belongs in an orchestration module. """ from __future__ import annotations @@ -13,7 +14,7 @@ from pathlib import Path SERVICES_DIR = Path(__file__).resolve().parents[1] / "src" / "transcription" / "services" -# Modules that intentionally compose several services rather than owning one table. +# Modules that intentionally compose several services rather than owning one aggregate. ORCHESTRATION_MODULES = frozenset({"store", "workflows", "__init__"}) diff --git a/tests/test_v42_evidence.py b/tests/test_v42_evidence.py index a551304..bd9dbf0 100644 --- a/tests/test_v42_evidence.py +++ b/tests/test_v42_evidence.py @@ -25,6 +25,7 @@ from transcription.providers.base import TranscriptionResult from transcription.providers.evidence import SourceEvidenceReference from transcription.providers.openrouter import OpenRouterTranscriptionProvider from transcription.services.documents import DocumentService +from transcription.services.evidence import EvidenceService from transcription.services.jobs import JobDeleteBlockedError from transcription.services.jobs import JobService from transcription.services.sources import SourceService @@ -203,6 +204,7 @@ async def test_attempts_are_append_only_and_exported_with_integrity(default_sess documents = DocumentService(session_factory=default_session_factory) jobs = JobService(session_factory=default_session_factory) sources = SourceService(session_factory=default_session_factory) + evidence = EvidenceService(session_factory=default_session_factory) document = await documents.create_document(Document(name="Evidence")) job = await jobs.create_job(Job(document_id=document.id)) source = await sources.create_source( @@ -239,14 +241,14 @@ async def test_attempts_are_append_only_and_exported_with_integrity(default_sess finished_at=now, ) - attempts = await sources.list_execution_attempts(source_id=source.id) + attempts = await evidence.list_execution_attempts(source_id=source.id) assert [attempt.attempt_number for attempt in attempts] == [1, 2] assert attempts[0].status == JobSourceStatus.FAILED assert attempts[0].error_detail == "first failed" assert attempts[1].status == JobSourceStatus.TRANSCRIBED assert attempts[1].raw_transcription == "second succeeded" - export = await sources.build_evidence_export(source_id=source.id) + export = await evidence.build_evidence_export(source_id=source.id) assert _json_object(export["source"])["digest_sha256"] == "a" * 64 assert [_json_object(item)["attempt_number"] for item in _json_array(export["attempts"])] == [1, 2] assert "file_path" not in json.dumps(export) @@ -257,7 +259,7 @@ async def test_attempts_are_append_only_and_exported_with_integrity(default_sess latest_job_source = detail.latest_job_source assert latest_job_source is not None assert latest_job_source.execution_attempts == [] - latest_attempt = await sources.read_latest_execution_attempt(job_source_id=latest_job_source.id) + latest_attempt = await evidence.read_latest_execution_attempt(job_source_id=latest_job_source.id) assert latest_attempt is not None assert latest_attempt.attempt.attempt_number == 2