4 Commits
Author SHA1 Message Date
zoltan57 8580aa7a1a Phase 6 negative test: deliberate lint error to prove the CI gate blocks.
Quality Gate / gate (push) Failing after 44s
This commit and its branch are scratch and will be deleted.
2026-08-18 16:18:58 -05:00
zoltan57andCopilot App fca959fa5d Phase 5: contain worker faults instead of discarding the classification
Review log [8]. classify_unexpected_error already returned retriable=False and
the verdict was logged and then thrown away. Measured across src/: retriable
was assigned in 9 places and read in none.

The plan asks for a test that a programming error "does not silently retry".
Probing with an injected AttributeError showed that is not what happens, and
the two real failure modes need different fixes.

Mode A, raised after the claim commits (inside advance_job): raised exactly
once, job left at PROCESSING, retry_count 0, never re-claimed, because
claim_next_queued_job filters status == QUEUED. A permanently stranded job
with one swallowed log line, not a retry. advance_job's PROCESSING branch,
commented "Recover mid-flight jobs", is unreachable from the worker for the
same reason.

Mode B, raised before or during the claim: 20 raises in 1.2s, an unbounded hot
spin at the poll interval. It never reaches the per-job retry machinery, so
WORKER_MAX_RETRIES does not cap it and the plan's 60s worst case understates
this path.

services/workflows.py
  _advance_job_with_containment wraps advance_job. Any escaping exception is
  classified and the job driven to terminal FAILED, which is visible in the UI
  and resubmittable. The caller session is rolled back first and the terminal
  write runs in its own transaction, so it stays atomic even when the failure
  left that session dirty (plan task 3). The loop continues, so one poison job
  cannot halt transcription for every other job.

worker.py
  handle_worker_exceptions re-raises non-retriable faults rather than
  suppressing them; retriable ones are still suppressed so transient
  conditions do not stop work. run_worker_loop catches that, logs CRITICAL and
  returns cleanly. Returning rather than propagating matters: the exception
  would otherwise surface only at app shutdown, through the wait_for in
  worker_consumer_lifespan.

tests
  test_run_worker_loop_survives_process_next_exception asserted the loop
  SURVIVES a RuntimeError and continues, which is the Mode B defect written
  down as an expectation. Replaced by
  test_run_worker_loop_stops_on_non_retriable_exception, with a new
  test_run_worker_loop_survives_retriable_exception so suppression of genuinely
  transient faults stays covered, and
  test_error_after_claim_fails_the_job_instead_of_stranding_it for Mode A.

  All three were verified to fail on pre-fix code. The Mode B guard fails by
  timing out, which is the infinite spin made visible.

Verified: 295 passed, 4 skipped, 0 ruff, 0 ty.

Co-authored-by: Copilot App <[email protected]>
2026-08-18 16:12:52 -05:00
zoltan57andCopilot App 110f40a28b Phase 4: measure only the provider call in duration_ms
Review log [55]. Three historical local_timeout rows recorded 0.4-2.0s more
than the configured budget because the measurement window opened before the
provider call.

The plan named two causes, and both were already gone. Diffed against
f86c0ff~1: at V4.6 the window held resolve_provider_input (async;
normalization + artifact write + DB work) and a session.commit(). Phase 1
deleted both. What remains between the clock and the wait_for is
build_provider_input, now pure field copying because normalization moved to
ingest and file_hash is already stored: 6.2 us per call, zero awaits, so it
cannot yield to the event loop.

A third cause was still there and is not in the plan. The regression test
below measured 890ms where ~200ms was expected. services.sources.provider is
a lazy property that appears as an argument expression to _call_transcriber,
so it is evaluated after the clock starts but before wait_for begins timing.
Constructing OpenRouterTranscriptionProvider costs 475ms on first access and
0.001ms after, so the first attempt of every worker process booked half a
second of HTTP client construction as provider latency. That plausibly
accounts for the low end of the historical overshoot.

workflows.py
  - Re-capture monotonic_started_at immediately before the wait_for, reusing
    the same variable. The pre-loop assignment stays as the fallback: binding
    a new name inside the try would leave the general-exception handler
    referencing an unbound variable when build_provider_input raises. All
    three duration write sites (success, TimeoutError, general failure) then
    measure the correct window with no further change.
  - Hoist the provider property above the per-source loop. It is
    loop-invariant, so this also removes the repeated lookup from the two
    evidence-capture sites.

tests/services/test_workflows_reliability.py
  test_timeout_duration_excludes_pre_call_setup simulates 400ms of blocking
  setup against a 200ms budget and asserts the recorded duration sits near
  the budget and well clear of budget+setup. Confirmed to fail on the pre-fix
  code (assert 625 < 540) and pass after, so it guards behaviour rather than
  restating it. This is the plan's verification criterion as a test.

ui/pages/sources_page.py
  _format_duration renders >=1s as "27.6 s" and below that as "612 ms",
  replacing the raw "27612 ms". No test asserted the old format.

Plan task 3 (record preprocessing as its own value) declined and logged as a
deviation: after Phase 1 there is no preprocessing left to record, and a
preprocessing_ms column to measure 6 us of attribute copying is complexity
without a reader.

Verified: 293 passed, 4 skipped, 0 ruff, 0 ty.

Co-authored-by: Copilot App <[email protected]>
2026-08-18 16:02:58 -05:00
zoltan57andCopilot App 7dd0d2c9bf Phase 3: extract EvidenceService and rewrite the service ownership rule
Decompose SourceService along the aggregate boundary and then correct the
instruction file that caused it to grow, in that order. The refactor is the
empirical test of the rule.

services/evidence.py (new)
  EvidenceService owns ExecutionAttempt: read_latest_execution_attempt,
  list_execution_attempts, promote_machine_attempt, build_evidence_export,
  plus the LatestExecutionAttempt projection. Moved verbatim from sources.py.

services/errors.py (new)
  The five-class error hierarchy (PromptLoadError, TranscriptionError,
  TranscriptionNotFoundError, SourceDeleteBlockedError,
  CandidatePromotionError) moved out of sources.py. evidence.py needs
  TranscriptionNotFoundError, and test_service_boundaries.py correctly
  rejected the sibling import. errors.py defines no *Service class, so it is
  a legal shared home. This was the boundary test doing its job, not an
  obstacle to route around.

sources.py 1,389 -> 885 lines (1,063 after Phase 2).

services/__init__.py
  ServiceBundle and from_session_factory register evidence. Note that
  field-by-field ServiceBundle construction silently binds services to the
  process-global session factory via default_factory; from_session_factory is
  the only safe constructor. Two test bundles were fixed for this.

.github/instructions/services.instructions.md
  Rewritten to describe the boundaries the decomposition actually produced,
  per plan Phase 3 task 7 and review log [59].

  - "1 service class per data model" -> one service class per aggregate.
    The table-shaped rule is the measured cause of sources.py reaching
    1,389 lines; DocumentType has no lifecycle without Document.
  - New Model Ownership section. Junctions are owned by their lifecycle
    owner, the service that creates and deletes the rows: document_person
    to PeopleService (sole writer, measured), job_source to SourceService.
    Two carve-outs are stated rather than left as silent violations:
    cascade deletion when a service deletes its own aggregate root, and
    status transitions that create and delete nothing (cancel_job,
    resubmit_failed_sources), which are Job lifecycle events on the work
    queue. EvidenceService.promote_machine_attempt's two-field write to
    Source is named and scoped.
  - Mandatory CRUD softened to intent. It was already false: five modules
    define no service class, EvidenceService has no create/delete because
    ExecutionAttempt is append-only, RegistryService uses <op>_entry.
  - Separated reading across models via eager loads from the owning root,
    which is allowed, from importing another service, which is not. The old
    line 13 and lines 75-77 read as contradictory.
  - Typo: picutre.

  No code was moved to satisfy the rule.

tests/test_service_boundaries.py
  Docstring no longer cites the instruction file by line number; that anchor
  would desynchronise silently. errors.py added to the neutral-module list.

Verified: 292 passed, 4 skipped, 0 ruff, 0 ty. All 25 /ui/* routes walked
against the live app; 24x 200. /ui/documents/{id}/sources 404s via a 307 that
drops the /ui prefix, confirmed pre-existing (last touched in 6a3ee26) and
left alone as out of scope.

Co-authored-by: Copilot App <[email protected]>
2026-08-18 15:54:27 -05:00
20 changed files with 692 additions and 251 deletions
+81 -14
View File
@@ -7,29 +7,89 @@ applyTo: 'src/transcription/services/*.py'
## Structure ## Structure
- Project core data models defined in [models](../../src/transcription/db/models.py) - Project core data models are defined in [models](../../src/transcription/db/models.py)
- 1 service class per data model - One service class per **aggregate**, not per table. An aggregate is a root model plus
- Only services directly interact with the database, and only through async methods the models that have no independent lifecycle of their own. `DocumentType` has no
- 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. 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 ## Error Handling
- Service-specific errors defined at the top of the respective module and inherit from `AppError` - Errors used by a single service are defined at the top of that 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 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 ## Checklist
- [ ] Uses `ServiceBase` for common logic - [ ] Uses `ServiceBase` for common logic
- [ ] CRUD methods created at the top - [ ] Session kwarg for `AsyncSession` to pass a session object into each method
- [ ] 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 through
- [ ] 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
- 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 ## CRUD Methods
- Create, read, update, and delete, created in that order - Name format `<operation>_<model>`, for example `create_document` or `update_job`.
- Name format `<operation>_<model >`, for example `create_document` or `update_job` - Where a service exposes create/read/update/delete for its root model, define them at the
- All services must define these 4 methods first, and in that order 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 `<operation>_entry`
naming instead.
## Transaction Finalization ## Transaction Finalization
@@ -74,4 +134,11 @@ Separation of concerns:
# Service Composition # 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.
+45
View File
@@ -0,0 +1,45 @@
name: Quality Gate
# V4.7 Phase 6 / review log [40]. Before this, ruff, ty and pytest were enforced
# only by .pre-commit-config.yaml, and only for developers who had actually run
# `pre-commit install`.
on:
push:
pull_request:
jobs:
gate:
runs-on: ubuntu-latest
steps:
- name: Check out the commit under test
uses: actions/checkout@v4
- name: Install uv
run: |
curl -LsSf https://astral.sh/uv/install.sh | sh
echo "$HOME/.local/bin" >> "$GITHUB_PATH"
- name: Install dependencies from the lockfile
# --locked fails if uv.lock has drifted from pyproject.toml, so a stale
# lockfile is caught here rather than producing an untested dependency set.
run: uv sync --locked
- name: Write placeholder configuration
# Settings requires openrouter_api_key and 115 tests cannot construct
# Settings without it. This is written to a .env file rather than exported
# as an environment variable on purpose: the external tests guard on
# os.getenv("OPENROUTER_API_KEY"), which reads the process environment and
# not the file, so writing the file reproduces the local result exactly -
# the 4 external tests skip instead of running against a fake key and
# failing. Exporting it instead produces 3 failures.
run: echo "OPENROUTER_API_KEY=ci-placeholder-not-a-real-key" > .env
- name: Lint and type check
# Runs the hooks defined in .pre-commit-config.yaml instead of repeating
# "ruff check" and "ty check" here. The commands then have one definition,
# so the local and CI gates cannot drift apart.
run: uv run pre-commit run --all-files --show-diff-on-failure
- name: Tests
run: uv run pytest
+7
View File
@@ -0,0 +1,7 @@
"""Scratch file for the Phase 6 CI negative test. Deleted with this branch."""
import os
def broken() -> int:
return undefined_name_that_does_not_exist
+12 -1
View File
@@ -9,12 +9,21 @@ from sqlmodel.ext.asyncio.session import AsyncSession
from ..config import Settings from ..config import Settings
from .documents import DocumentService from .documents import DocumentService
from .evidence import EvidenceService
from .jobs import JobService from .jobs import JobService
from .people import PeopleService from .people import PeopleService
from .prompts import PromptStore from .prompts import PromptStore
from .sources import SourceService from .sources import SourceService
__all__ = ["DocumentService", "JobService", "PeopleService", "PromptStore", "ServiceBundle", "SourceService"] __all__ = [
"DocumentService",
"EvidenceService",
"JobService",
"PeopleService",
"PromptStore",
"ServiceBundle",
"SourceService",
]
@dataclass(frozen=True, slots=True) @dataclass(frozen=True, slots=True)
@@ -25,6 +34,7 @@ class ServiceBundle:
sources: SourceService = field(default_factory=SourceService) sources: SourceService = field(default_factory=SourceService)
jobs: JobService = field(default_factory=JobService) jobs: JobService = field(default_factory=JobService)
people: PeopleService = field(default_factory=PeopleService) people: PeopleService = field(default_factory=PeopleService)
evidence: EvidenceService = field(default_factory=EvidenceService)
@classmethod @classmethod
def from_session_factory( def from_session_factory(
@@ -41,6 +51,7 @@ class ServiceBundle:
sources=SourceService(session_factory=session_factory, settings=settings), sources=SourceService(session_factory=session_factory, settings=settings),
jobs=JobService(session_factory=session_factory, settings=settings), jobs=JobService(session_factory=session_factory, settings=settings),
people=PeopleService(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: async def aclose(self) -> None:
+32
View File
@@ -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."""
+210
View File
@@ -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,
)
+4 -182
View File
@@ -3,7 +3,6 @@
from __future__ import annotations from __future__ import annotations
import asyncio import asyncio
import base64
import hashlib import hashlib
import logging import logging
from collections.abc import Sequence from collections.abc import Sequence
@@ -22,7 +21,6 @@ from pydantic import JsonValue
from pydantic import TypeAdapter from pydantic import TypeAdapter
from pydantic import ValidationError from pydantic import ValidationError
from sqlalchemy import func from sqlalchemy import func
from sqlalchemy import inspect as sqlalchemy_inspect
from sqlalchemy import literal from sqlalchemy import literal
from sqlalchemy import tuple_ from sqlalchemy import tuple_
from sqlalchemy.ext.asyncio import async_sessionmaker 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 JobSource
from transcription.db.models import JobSourceStatus from transcription.db.models import JobSourceStatus
from transcription.db.models import Source from transcription.db.models import Source
from transcription.errors import AppError
from transcription.errors import ErrorCategory from transcription.errors import ErrorCategory
from transcription.providers import ProviderAuthError from transcription.providers import ProviderAuthError
from transcription.providers import ProviderError from transcription.providers import ProviderError
@@ -50,10 +47,13 @@ from transcription.providers import TranscriptionResult
from transcription.providers import TransportEvidence from transcription.providers import TransportEvidence
from transcription.providers import get_transcription_provider from transcription.providers import get_transcription_provider
from ..db.loading import defer
from ..db.loading import orm_attribute from ..db.loading import orm_attribute
from ..db.loading import selectinload from ..db.loading import selectinload
from .base import ServiceBase 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 lookup_source_mime_type
from .source_media import supported_source_formats 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) 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) @dataclass(frozen=True, slots=True)
class SourceNavigation: class SourceNavigation:
"""Adjacent Source identifiers within one ordered Document.""" """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): class SourceService(ServiceBase):
"""Manage source records, media payloads, revisions, and page execution output.""" """Manage source records, media payloads, revisions, and page execution output."""
@@ -214,35 +186,6 @@ class SourceService(ServiceBase):
) )
return source 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( async def read_source_navigation(
self, self,
source_id: UUID, source_id: UUID,
@@ -640,127 +583,6 @@ class SourceService(ServiceBase):
await self._finalize(session=_session, caller_session=session, refresh=(job, source, job_source, attempt)) await self._finalize(session=_session, caller_session=session, refresh=(job, source, job_source, attempt))
return job_source 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( async def upsert_revision_for_source(
self, self,
*, *,
+1 -1
View File
@@ -23,10 +23,10 @@ from ..db.models import JobSourceStatus
from ..db.models import Source from ..db.models import Source
from ..db.session import SessionFactory from ..db.session import SessionFactory
from ..db.session import session_scope from ..db.session import session_scope
from .errors import TranscriptionError
from .media_storage import build_stored_filename from .media_storage import build_stored_filename
from .media_storage import write_media_bytes from .media_storage import write_media_bytes
from .normalization import normalize_orientation_async from .normalization import normalize_orientation_async
from .sources import TranscriptionError
from .sources import build_prompt_execution from .sources import build_prompt_execution
from .sources import source_mime_type from .sources import source_mime_type
from .sources import validate_source_content from .sources import validate_source_content
+50 -4
View File
@@ -205,6 +205,10 @@ async def process_queued_job( # noqa: PLR0915
externally_stopped = False externally_stopped = False
prompt_execution = _resolve_job_prompt_execution(source_job=source_job, settings=runtime_settings) prompt_execution = _resolve_job_prompt_execution(source_job=source_job, settings=runtime_settings)
# Resolve the lazy provider property once, outside the timed region. First access
# constructs the HTTP client (~0.5s), which would otherwise be booked as provider
# latency on the first attempt of every worker process (review log [55]).
provider = services.sources.provider
for source in sources: for source in sources:
if await _job_no_longer_processing(job_id=job.id, services=services, session=session): if await _job_no_longer_processing(job_id=job.id, services=services, session=session):
@@ -212,6 +216,8 @@ async def process_queued_job( # noqa: PLR0915
break break
started_at = datetime.now(UTC) started_at = datetime.now(UTC)
# Fallback start for failures raised before the provider call; reset to the
# true call boundary immediately before the wait_for below.
monotonic_started_at = asyncio.get_running_loop().time() monotonic_started_at = asyncio.get_running_loop().time()
result: TranscriptionResult | None = None result: TranscriptionResult | None = None
provider_input = None provider_input = None
@@ -225,12 +231,16 @@ async def process_queued_job( # noqa: PLR0915
media_type=provider_input.media_type, media_type=provider_input.media_type,
page_number=source.page_number, page_number=source.page_number,
) )
# Restart the clock so duration_ms covers only what the wait_for below
# governs. The pre-loop assignment stays as the fallback for failures
# raised before this point, which would otherwise leave it unbound.
monotonic_started_at = asyncio.get_running_loop().time()
result = await asyncio.wait_for( result = await asyncio.wait_for(
_call_transcriber( _call_transcriber(
input_path=provider_input.path, input_path=provider_input.path,
prompt_execution=prompt_execution, prompt_execution=prompt_execution,
settings=runtime_settings, settings=runtime_settings,
provider=services.sources.provider, provider=provider,
source_reference=source_reference, source_reference=source_reference,
requested_model=source_job.model, requested_model=source_job.model,
), ),
@@ -283,8 +293,8 @@ async def process_queued_job( # noqa: PLR0915
0, 0,
int((asyncio.get_running_loop().time() - monotonic_started_at) * 1000), int((asyncio.get_running_loop().time() - monotonic_started_at) * 1000),
), ),
request_manifest=services.sources.provider.current_request_manifest, request_manifest=provider.current_request_manifest,
transport_evidence=services.sources.provider.current_transport_evidence, transport_evidence=provider.current_transport_evidence,
failure_phase="local_timeout", failure_phase="local_timeout",
) )
failed_pages.append(page_outcome) failed_pages.append(page_outcome)
@@ -406,10 +416,46 @@ async def process_next_queued_job(
if session is not None: if session is not None:
await session.commit() await session.commit()
await advance_job(job=job, services=services, settings=settings, session=session) await _advance_job_with_containment(job=job, services=services, settings=settings, session=session)
return True return True
async def _advance_job_with_containment(
*,
job: Job,
services: ServiceBundle,
settings: Settings | None,
session: AsyncSession | None,
) -> None:
"""Advance a claimed job, guaranteeing it never stays stuck in PROCESSING.
The job has already been committed as PROCESSING at this point, and
``claim_next_queued_job`` only ever selects QUEUED rows. So an exception
escaping ``advance_job`` used to strand the job in PROCESSING permanently,
with a single swallowed log line and no recovery path (review log [8]).
Any escaping exception is therefore classified and the job driven to the
terminal FAILED state, which is visible in the UI and resubmittable. The
terminal write runs in its own transaction so it is atomic even when the
caller's session was left dirty by the failure.
"""
try:
await advance_job(job=job, services=services, settings=settings, session=session)
except Exception as exc:
error = exc if isinstance(exc, AppError) else classify_unexpected_error(exc, operation="worker.advance_job")
logger.exception(
"Job processing failed outside page handling operation=worker.advance_job "
"job_id=%s error_id=%s category=%s retriable=%s",
job.id,
error.error_id,
error.category.value,
error.retriable,
)
if session is not None:
await session.rollback()
await services.jobs.update_job_state(job_id=job.id, status=JobStatus.FAILED)
def _resolve_job_sources(job: Job) -> list[Source]: def _resolve_job_sources(job: Job) -> list[Source]:
"""Resolve pending linked sources for a job in deterministic page order. """Resolve pending linked sources for a job in deterministic page order.
+21 -12
View File
@@ -14,10 +14,11 @@ from transcription.db.models import ExecutionAttempt
from transcription.db.models import JobSource from transcription.db.models import JobSource
from transcription.db.models import JobSourceStatus from transcription.db.models import JobSourceStatus
from transcription.db.models import Source from transcription.db.models import Source
from transcription.services.sources import LatestExecutionAttempt from transcription.services.errors import SourceDeleteBlockedError
from transcription.services.sources 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 SourceService
from transcription.services.sources import TranscriptionNotFoundError
from transcription.ui.components.app_shell import render_navigation_header from transcription.ui.components.app_shell import render_navigation_header
from transcription.ui.components.cards import archival_card from transcription.ui.components.cards import archival_card
from transcription.ui.components.confirm_delete import render_delete_actions 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}") @ui.page("/sources/{source_id}")
async def source_detail_page(source_id: str, request: Request, session_factory: SessionFactoryDep) -> None: async def source_detail_page(source_id: str, request: Request, session_factory: SessionFactoryDep) -> None:
sources_service = SourceService(session_factory=session_factory) sources_service = SourceService(session_factory=session_factory)
evidence_service = EvidenceService(session_factory=session_factory)
render_navigation_header(current_path="/sources") render_navigation_header(current_path="/sources")
parsed_source_id = parsed_record_id(source_id, noun="Source") 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) navigation = await sources_service.read_source_navigation(parsed_source_id)
latest_job_source = source.latest_job_source latest_job_source = source.latest_job_source
latest_attempt = ( 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 if latest_job_source is not None
else 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: except TranscriptionNotFoundError:
render_record_not_found("Source") render_record_not_found("Source")
return return
@@ -158,7 +160,7 @@ def register_page() -> None: # noqa: PLR0915
"Export Evidence", "Export Evidence",
on_click=lambda: _download_evidence( on_click=lambda: _download_evidence(
source_id=source.id, source_id=source.id,
sources_service=sources_service, evidence_service=evidence_service,
), ),
icon="download", icon="download",
).props("flat") ).props("flat")
@@ -185,7 +187,7 @@ def register_page() -> None: # noqa: PLR0915
_render_machine_candidates( _render_machine_candidates(
source=source, source=source,
attempts=attempts, attempts=attempts,
sources_service=sources_service, evidence_service=evidence_service,
) )
_render_source_metadata_column( _render_source_metadata_column(
source=source, source=source,
@@ -364,6 +366,13 @@ def _render_source_job_metadata_zone(
_render_provider_evidence(latest_attempt=latest_attempt) _render_provider_evidence(latest_attempt=latest_attempt)
def _format_duration(duration_ms: int) -> str:
"""Render an attempt duration with a unit that suits its magnitude."""
if duration_ms >= 1000:
return f"{duration_ms / 1000:.1f} s"
return f"{duration_ms} ms"
def _render_provider_evidence(*, latest_attempt: LatestExecutionAttempt | None) -> None: def _render_provider_evidence(*, latest_attempt: LatestExecutionAttempt | None) -> None:
ui.label("Provider Evidence").classes("text-xs font-semibold ui-text-primary mt-3") ui.label("Provider Evidence").classes("text-xs font-semibold ui-text-primary mt-3")
if latest_attempt is None: if latest_attempt is None:
@@ -372,7 +381,7 @@ def _render_provider_evidence(*, latest_attempt: LatestExecutionAttempt | None)
attempt = latest_attempt.attempt attempt = latest_attempt.attempt
metadata_row("Attempt:", str(attempt.attempt_number)) metadata_row("Attempt:", str(attempt.attempt_number))
metadata_row("Duration:", f"{attempt.duration_ms} ms") metadata_row("Duration:", _format_duration(attempt.duration_ms))
_render_json_evidence("Request Manifest", attempt.request_manifest) _render_json_evidence("Request Manifest", attempt.request_manifest)
_render_json_evidence("Transport Response", _transport_display(latest_attempt)) _render_json_evidence("Transport Response", _transport_display(latest_attempt))
_render_json_evidence("OpenRouter SDK Response Snapshot", attempt.sdk_response_snapshot) _render_json_evidence("OpenRouter SDK Response Snapshot", attempt.sdk_response_snapshot)
@@ -409,9 +418,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: 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 except Exception as exc: # noqa: BLE001
show_error(exc, title="Export failed", operation="sources.evidence_export") show_error(exc, title="Export failed", operation="sources.evidence_export")
return return
@@ -521,7 +530,7 @@ def _render_machine_candidates(
*, *,
source: Source, source: Source,
attempts: list[ExecutionAttempt], attempts: list[ExecutionAttempt],
sources_service: SourceService, evidence_service: EvidenceService,
) -> None: ) -> None:
successful = [ successful = [
attempt attempt
@@ -581,7 +590,7 @@ def _render_machine_candidates(
async def promote(candidate_id: UUID = attempt.id) -> None: async def promote(candidate_id: UUID = attempt.id) -> None:
try: try:
await sources_service.promote_machine_attempt( await evidence_service.promote_machine_attempt(
source_id=source.id, source_id=source.id,
execution_attempt_id=candidate_id, execution_attempt_id=candidate_id,
) )
+29 -2
View File
@@ -94,16 +94,30 @@ async def worker_consumer_lifespan(
@contextmanager @contextmanager
def handle_worker_exceptions(operation: str = "worker.loop"): def handle_worker_exceptions(operation: str = "worker.loop"):
"""Context manager to log and suppress exceptions in the worker loop.""" """Log worker-loop exceptions, suppressing only the retriable ones.
``classify_unexpected_error`` marks unknown exceptions non-retriable, but that
verdict used to be logged and then discarded, so a programming error raised
before a job could be claimed spun the loop at the poll interval indefinitely.
Nothing capped it, because it never reached the per-job retry machinery
(review log [8]).
A non-retriable fault is a defect rather than a transient condition, so it is
re-raised for the caller to stop on. Faults raised after a job is claimed are
contained at the job level instead, and never reach here.
"""
try: try:
yield yield
except Exception as exc: except Exception as exc:
error = exc if isinstance(exc, AppError) else classify_unexpected_error(exc, operation=operation) error = exc if isinstance(exc, AppError) else classify_unexpected_error(exc, operation=operation)
logger.exception( logger.exception(
"Worker loop exception error_id=%s category=%s", "Worker loop exception error_id=%s category=%s retriable=%s",
error.error_id, error.error_id,
error.category.value, error.category.value,
error.retriable,
) )
if not error.retriable:
raise error from exc
async def run_worker_loop( async def run_worker_loop(
@@ -121,6 +135,10 @@ async def run_worker_loop(
The service bundle and with it the provider's pooled HTTP client — is built The service bundle and with it the provider's pooled HTTP client — is built
once for the lifetime of the loop, so consecutive jobs reuse one connection once for the lifetime of the loop, so consecutive jobs reuse one connection
instead of paying a fresh TLS handshake each time. instead of paying a fresh TLS handshake each time.
The loop returns early on a non-retriable error raised before a job could be
claimed. Faults raised after a claim are contained by marking that job FAILED,
so a single poison job cannot stop transcription for every other job.
""" """
services = ServiceBundle.from_session_factory(session_factory) services = ServiceBundle.from_session_factory(session_factory)
try: try:
@@ -150,6 +168,15 @@ async def run_worker_loop(
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)
except AppError as error:
# Stop rather than spin. A non-retriable fault here means no job could be
# claimed, so there is no row to mark FAILED and no reason to expect the
# next poll to behave differently.
logger.critical(
"Worker loop stopped after a non-retriable error error_id=%s category=%s",
error.error_id,
error.category.value,
)
finally: finally:
await services.aclose() await services.aclose()
+3 -1
View File
@@ -13,6 +13,7 @@ from transcription.db.models import JobSourceStatus
from transcription.db.models import JobStatus 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.evidence import EvidenceService
from transcription.services.jobs import JobCancelBlockedError from transcription.services.jobs import JobCancelBlockedError
from transcription.services.jobs import JobDeleteBlockedError from transcription.services.jobs import JobDeleteBlockedError
from transcription.services.jobs import JobNotFoundError from transcription.services.jobs import JobNotFoundError
@@ -274,6 +275,7 @@ class TestJobService:
document_service: DocumentService, document_service: DocumentService,
): ):
source_service = SourceService(session_factory=job_service.session_factory) 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")) 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)) job = await job_service.create_job(Job(document_id=document.id, status=JobStatus.FAILED))
source = await source_service.create_source( source = await source_service.create_source(
@@ -299,7 +301,7 @@ class TestJobService:
with pytest.raises(JobNotFoundError): with pytest.raises(JobNotFoundError):
await job_service.read_job(job_id=job.id) 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 assert (await source_service.read_source(source.id)).id == source.id
@pytest.mark.asyncio @pytest.mark.asyncio
+2 -2
View File
@@ -13,10 +13,10 @@ from transcription.db.models import JobSourceStatus
from transcription.db.models import JobStatus 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.errors import SourceDeleteBlockedError
from transcription.services.errors import TranscriptionNotFoundError
from transcription.services.jobs import JobService from transcription.services.jobs import JobService
from transcription.services.sources import SourceDeleteBlockedError
from transcription.services.sources import SourceService from transcription.services.sources import SourceService
from transcription.services.sources import TranscriptionNotFoundError
@pytest.mark.integration @pytest.mark.integration
+4 -2
View File
@@ -13,10 +13,11 @@ from transcription.db.models import Source
from transcription.errors import ErrorCategory from transcription.errors import ErrorCategory
from transcription.services.documents import DocumentDeleteBlockedError from transcription.services.documents import DocumentDeleteBlockedError
from transcription.services.documents import DocumentService 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.jobs import JobService
from transcription.services.people import PeopleError from transcription.services.people import PeopleError
from transcription.services.people import PeopleService from transcription.services.people import PeopleService
from transcription.services.sources import SourceDeleteBlockedError
from transcription.services.sources import SourceService 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 len(stored_rows) == 1
assert stored_rows[0].status == JobSourceStatus.TRANSCRIBED 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 is not None
assert attempt.attempt.raw_transcription == "provider transcript" assert attempt.attempt.raw_transcription == "provider transcript"
assert attempt.attempt.normalized_metadata == metadata assert attempt.attempt.normalized_metadata == metadata
+5 -12
View File
@@ -11,19 +11,12 @@ from transcription.db.models import JobPurpose
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 import ServiceBundle from transcription.services import ServiceBundle
from transcription.services.documents import DocumentService from transcription.services.errors import CandidatePromotionError
from transcription.services.jobs import JobService
from transcription.services.sources import CandidatePromotionError
from transcription.services.sources import SourceService
from transcription.services.workflows import create_source_retranscription_job from transcription.services.workflows import create_source_retranscription_job
def _services(default_session_factory, settings: Settings) -> ServiceBundle: def _services(default_session_factory, settings: Settings) -> ServiceBundle:
return ServiceBundle( return ServiceBundle.from_session_factory(default_session_factory, settings=settings)
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),
)
async def _seed_source(services: ServiceBundle) -> Source: 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", model="model-b",
) )
unchanged = await services.sources.read_source(source.id) 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.raw_transcription == "first result"
assert unchanged.preferred_execution_attempt_id == first_attempt_id assert unchanged.preferred_execution_attempt_id == first_attempt_id
assert {attempt.raw_transcription for attempt in attempts} == {"first result", "candidate result"} 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") 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, source_id=source.id,
execution_attempt_id=candidate.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) source = await _seed_source(services)
with pytest.raises(CandidatePromotionError): with pytest.raises(CandidatePromotionError):
await services.sources.promote_machine_attempt( await services.evidence.promote_machine_attempt(
source_id=source.id, source_id=source.id,
execution_attempt_id=uuid4(), execution_attempt_id=uuid4(),
) )
+134 -7
View File
@@ -1,6 +1,7 @@
"""Reliability tests for worker workflow timeout behavior.""" """Reliability tests for worker workflow timeout behavior."""
import asyncio import asyncio
import time
from pathlib import Path from pathlib import Path
from uuid import uuid4 from uuid import uuid4
@@ -18,6 +19,7 @@ from transcription.db.models import JobStatus
from transcription.db.models import Source from transcription.db.models import Source
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 import workflows as workflows_module
from transcription.services.workflows import process_queued_job from transcription.services.workflows import process_queued_job
@@ -103,18 +105,143 @@ class TestWorkflowReliability:
assert "timed out" in error_detail.lower() assert "timed out" in error_detail.lower()
assert "20.0s" in error_detail assert "20.0s" in error_detail
@pytest.mark.asyncio
async def test_timeout_duration_excludes_pre_call_setup(self, default_session_factory, monkeypatch):
"""duration_ms covers only the provider call, not the setup preceding it.
Regression guard for review log [55]: three historical ``local_timeout`` rows
recorded 0.4-2.0 s more than the configured budget because the measurement
window opened before payload resolution. Blocking setup is simulated here so
the assertion fails if that window ever reopens.
"""
services = ServiceBundle.from_session_factory(default_session_factory)
async with services.jobs._session_scope() as session:
document = Document(id=uuid4(), name="window-doc")
session.add(document)
await session.flush()
job = Job(document_id=document.id, status=JobStatus.QUEUED)
session.add(job)
await session.flush()
source = Source(
document_id=document.id,
page_number=1,
upload_name="window.jpg",
filename="window.jpg",
file_path=str(Path("tests/fixtures/images/real/Book Two - page 02.jpg")),
file_hash="d" * 64,
file_size_bytes=1,
)
session.add(source)
await session.flush()
session.add(JobSource(job_id=job.id, source_id=source.id, status=JobSourceStatus.PENDING))
await session.commit()
loaded = await services.jobs.read_job(job_id=job.id, session=session)
setup_seconds = 0.40
budget_seconds = 0.20
real_build = workflows_module.build_provider_input
def _slow_build(source_arg):
time.sleep(setup_seconds)
return real_build(source_arg)
async def _never_returns(*args, **kwargs):
_ = (args, kwargs)
await asyncio.sleep(budget_seconds * 20)
monkeypatch.setattr("transcription.services.workflows.build_provider_input", _slow_build)
monkeypatch.setattr("transcription.services.workflows.transcribe_document_image", _never_returns)
result = await process_queued_job(
job=loaded,
services=services,
settings=Settings(openrouter_api_key="test-key", worker_provider_timeout_seconds=budget_seconds),
)
assert result is not None
assert result.status == JobStatus.FAILED
async with services.jobs._session_scope() as session:
attempts = (
(
await session.exec(
select(ExecutionAttempt).where(
col(ExecutionAttempt.job_source_id).in_([js.id for js in result.job_sources])
)
)
)
.all()
)
assert len(attempts) == 1
duration_ms = attempts[0].duration_ms
# At or just above the budget, and well clear of budget + setup.
assert duration_ms >= int(budget_seconds * 1000 * 0.9)
assert duration_ms < int((budget_seconds + setup_seconds) * 1000 * 0.9)
@pytest.mark.asyncio
async def test_error_after_claim_fails_the_job_instead_of_stranding_it(
self,
default_session_factory,
monkeypatch,
):
"""A non-retriable fault after the claim drives the job terminal, not stuck.
Regression guard for review log [8]. The claim commits PROCESSING before any
provider work, and claim_next_queued_job only ever selects QUEUED, so an
exception escaping advance_job used to strand the job in PROCESSING forever
with one swallowed log line. Measured before the fix: raised once, job left
processing, retry_count 0, never re-claimed.
"""
services = ServiceBundle.from_session_factory(default_session_factory)
async with services.jobs._session_scope() as session:
document = Document(id=uuid4(), name="strand-doc")
session.add(document)
await session.flush()
job = Job(document_id=document.id, status=JobStatus.QUEUED)
session.add(job)
await session.flush()
source = Source(
document_id=document.id,
page_number=1,
upload_name="strand.jpg",
filename="strand.jpg",
file_path=str(Path("tests/fixtures/images/real/Book Two - page 02.jpg")),
file_hash="e" * 64,
file_size_bytes=1,
)
session.add(source)
await session.flush()
session.add(JobSource(job_id=job.id, source_id=source.id, status=JobSourceStatus.PENDING))
await session.commit()
job_id = job.id
async def _succeeds(*args, **kwargs):
_ = (args, kwargs)
return TranscriptionResult(text="page text", provider="test", model="test-model")
async def _boom(**kwargs):
_ = kwargs
raise AttributeError("deliberate programming error")
monkeypatch.setattr("transcription.services.workflows.transcribe_document_image", _succeeds)
monkeypatch.setattr("transcription.services.workflows._finalize_batch_outcome", _boom)
processed = await workflows_module.process_next_queued_job(services=services)
assert processed is True
async with services.jobs._session_scope() as session:
final = await session.get(Job, job_id)
assert final is not None
# Terminal and resubmittable, rather than stranded in PROCESSING.
assert final.status == JobStatus.FAILED
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_completed_page_is_committed_before_next_provider_call_finishes( async def test_completed_page_is_committed_before_next_provider_call_finishes(
self, self,
default_session_factory, default_session_factory,
monkeypatch, monkeypatch,
): ):
services = ServiceBundle( services = ServiceBundle.from_session_factory(default_session_factory)
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),
)
async with services.jobs._session_scope() as session: async with services.jobs._session_scope() as session:
document = Document(id=uuid4(), name="durability-doc") document = Document(id=uuid4(), name="durability-doc")
session.add(document) session.add(document)
@@ -155,7 +282,7 @@ class TestWorkflowReliability:
task = asyncio.create_task(process_queued_job(job=loaded, services=services)) task = asyncio.create_task(process_queued_job(job=loaded, services=services))
await asyncio.wait_for(second_started.wait(), timeout=2) 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 len(attempts) == 1
assert attempts[0].raw_transcription == "page 1" assert attempts[0].raw_transcription == "page 1"
+1 -1
View File
@@ -7,8 +7,8 @@ import pytest
from pydantic import ValidationError from pydantic import ValidationError
from transcription.config import Settings from transcription.config import Settings
from transcription.services.errors import PromptLoadError
from transcription.services.sources import PromptExecution 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 build_prompt_execution
from transcription.services.sources import load_prompt_text from transcription.services.sources import load_prompt_text
+6 -5
View File
@@ -1,9 +1,10 @@
"""Structural rules for the services package. """Structural rules for the services package.
`.github/instructions/services.instructions.md:13` requires that service classes The "Structure" section of `.github/instructions/services.instructions.md` requires
stay independent of one another. Shared behavior belongs in a neutral module that service modules stay independent of one another. Shared behavior belongs in a
(`base.py`, `registry.py`, `source_media.py`, `media_storage.py`), and any neutral module that defines no service class (`base.py`, `errors.py`, `registry.py`,
operation spanning two services belongs in an orchestration module. `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 from __future__ import annotations
@@ -13,7 +14,7 @@ from pathlib import Path
SERVICES_DIR = Path(__file__).resolve().parents[1] / "src" / "transcription" / "services" 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__"}) ORCHESTRATION_MODULES = frozenset({"store", "workflows", "__init__"})
+5 -3
View File
@@ -25,6 +25,7 @@ from transcription.providers.base import TranscriptionResult
from transcription.providers.evidence import SourceEvidenceReference from transcription.providers.evidence import SourceEvidenceReference
from transcription.providers.openrouter import OpenRouterTranscriptionProvider from transcription.providers.openrouter import OpenRouterTranscriptionProvider
from transcription.services.documents import DocumentService from transcription.services.documents import DocumentService
from transcription.services.evidence import EvidenceService
from transcription.services.jobs import JobDeleteBlockedError from transcription.services.jobs import JobDeleteBlockedError
from transcription.services.jobs import JobService from transcription.services.jobs import JobService
from transcription.services.sources import SourceService 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) documents = DocumentService(session_factory=default_session_factory)
jobs = JobService(session_factory=default_session_factory) jobs = JobService(session_factory=default_session_factory)
sources = SourceService(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")) document = await documents.create_document(Document(name="Evidence"))
job = await jobs.create_job(Job(document_id=document.id)) job = await jobs.create_job(Job(document_id=document.id))
source = await sources.create_source( 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, 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 [attempt.attempt_number for attempt in attempts] == [1, 2]
assert attempts[0].status == JobSourceStatus.FAILED assert attempts[0].status == JobSourceStatus.FAILED
assert attempts[0].error_detail == "first failed" assert attempts[0].error_detail == "first failed"
assert attempts[1].status == JobSourceStatus.TRANSCRIBED assert attempts[1].status == JobSourceStatus.TRANSCRIBED
assert attempts[1].raw_transcription == "second succeeded" 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(export["source"])["digest_sha256"] == "a" * 64
assert [_json_object(item)["attempt_number"] for item in _json_array(export["attempts"])] == [1, 2] assert [_json_object(item)["attempt_number"] for item in _json_array(export["attempts"])] == [1, 2]
assert "file_path" not in json.dumps(export) 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 latest_job_source = detail.latest_job_source
assert latest_job_source is not None assert latest_job_source is not None
assert latest_job_source.execution_attempts == [] 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 is not None
assert latest_attempt.attempt.attempt_number == 2 assert latest_attempt.attempt.attempt_number == 2
+40 -2
View File
@@ -4,6 +4,8 @@ from typing import cast
import pytest import pytest
from transcription.errors import AppError
from transcription.errors import ErrorCategory
from transcription.services import ServiceBundle from transcription.services import ServiceBundle
from transcription.services.sources import SourceService from transcription.services.sources import SourceService
from transcription.worker import process_next_queued_job from transcription.worker import process_next_queued_job
@@ -11,7 +13,38 @@ from transcription.worker import run_worker_loop
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_run_worker_loop_survives_process_next_exception(monkeypatch, caplog): async def test_run_worker_loop_stops_on_non_retriable_exception(monkeypatch, caplog):
"""A programming error before a job is claimed stops the loop instead of spinning.
Regression guard for review log [8]. This previously spun at the poll interval
forever: the fault was classified non-retriable, logged, and then discarded, and
it never reached the per-job retry machinery so nothing capped it. Measured at 20
iterations in 1.2s before the fix.
"""
calls = 0
stop_event = asyncio.Event()
async def _fake_process_next_queued_job(*, session=None, session_factory=None, services=None):
nonlocal calls
_ = (session, session_factory, services)
calls += 1
raise RuntimeError("boom")
monkeypatch.setattr("transcription.worker.process_next_queued_job", _fake_process_next_queued_job)
with caplog.at_level(logging.CRITICAL):
await asyncio.wait_for(
run_worker_loop(stop_event=stop_event, poll_interval_seconds=0),
timeout=5,
)
assert calls == 1
assert "Worker loop stopped after a non-retriable error" in caplog.text
@pytest.mark.asyncio
async def test_run_worker_loop_survives_retriable_exception(monkeypatch, caplog):
"""A retriable fault is still suppressed so transient conditions do not stop work."""
calls = 0 calls = 0
stop_event = asyncio.Event() stop_event = asyncio.Event()
@@ -20,7 +53,12 @@ async def test_run_worker_loop_survives_process_next_exception(monkeypatch, capl
_ = (session, session_factory, services) _ = (session, session_factory, services)
calls += 1 calls += 1
if calls == 1: if calls == 1:
raise RuntimeError("boom") raise AppError(
"transient",
category=ErrorCategory.EXTERNAL_PROVIDER,
suggestion="retry",
retriable=True,
)
stop_event.set() stop_event.set()
return False return False