generated from john/python-template
claude-sonnet-5 review: Phase 1 implemented by gpt-5.3-codex
Quality Gate / gate (push) Failing after 12s
Quality Gate / gate (push) Failing after 12s
This commit is contained in:
@@ -360,6 +360,7 @@ class JobSource(SQLModel, table=True):
|
||||
"""A single AI execution record for one source page."""
|
||||
|
||||
__tablename__ = "job_source"
|
||||
__table_args__ = (UniqueConstraint("job_id", "source_id", name="uq_job_source_job_source"),)
|
||||
|
||||
id: UUID = Field(default_factory=uuid4, primary_key=True)
|
||||
job_id: UUID = Field(foreign_key="job.id", index=True)
|
||||
|
||||
@@ -5,6 +5,7 @@ from datetime import datetime
|
||||
from uuid import UUID
|
||||
|
||||
from sqlalchemy import func
|
||||
from sqlalchemy import update
|
||||
from sqlmodel import col
|
||||
from sqlmodel import select
|
||||
from sqlmodel.ext.asyncio.session import AsyncSession
|
||||
@@ -182,21 +183,46 @@ class JobService(ServiceBase):
|
||||
concurrent workers never contend for the same job.
|
||||
"""
|
||||
async with self._session_scope(session) as _session:
|
||||
query = (
|
||||
select(Job)
|
||||
.where(Job.status == JobStatus.QUEUED)
|
||||
# Break ties by id so "next" is stable when two rows share close timestamps.
|
||||
dialect = _session.get_bind().dialect.name
|
||||
if dialect == "postgresql":
|
||||
query = (
|
||||
select(Job)
|
||||
.where(Job.status == JobStatus.QUEUED)
|
||||
# Break ties by id so "next" is stable when two rows share close timestamps.
|
||||
.order_by(col(Job.date_created), col(Job.id))
|
||||
.limit(1)
|
||||
.with_for_update(skip_locked=True)
|
||||
)
|
||||
job = (await _session.exec(query)).first()
|
||||
if job is None:
|
||||
return None
|
||||
job.status = JobStatus.PROCESSING
|
||||
await self._finalize(session=_session, caller_session=session, refresh=(job,))
|
||||
return job
|
||||
|
||||
now = datetime.now(UTC)
|
||||
queued_job_id = (
|
||||
select(col(Job.id))
|
||||
.where(col(Job.status) == JobStatus.QUEUED)
|
||||
.order_by(col(Job.date_created), col(Job.id))
|
||||
.limit(1)
|
||||
.scalar_subquery()
|
||||
)
|
||||
if _session.get_bind().dialect.name == "postgresql":
|
||||
query = query.with_for_update(skip_locked=True)
|
||||
claim_statement = (
|
||||
update(Job)
|
||||
.where(col(Job.id) == queued_job_id)
|
||||
.where(col(Job.status) == JobStatus.QUEUED)
|
||||
.values(status=JobStatus.PROCESSING, date_updated=now)
|
||||
.returning(col(Job.id))
|
||||
)
|
||||
claimed_row = (await _session.exec(claim_statement)).first()
|
||||
if claimed_row is None:
|
||||
return None
|
||||
claimed_job_id = claimed_row if isinstance(claimed_row, UUID) else claimed_row[0]
|
||||
|
||||
job = (await _session.exec(query)).first()
|
||||
job = (await _session.exec(select(Job).where(Job.id == claimed_job_id))).first()
|
||||
if job is None:
|
||||
return None
|
||||
|
||||
job.status = JobStatus.PROCESSING
|
||||
await self._finalize(session=_session, caller_session=session, refresh=(job,))
|
||||
return job
|
||||
|
||||
|
||||
@@ -23,6 +23,7 @@ from pydantic import ValidationError
|
||||
from sqlalchemy import func
|
||||
from sqlalchemy import literal
|
||||
from sqlalchemy import tuple_
|
||||
from sqlalchemy.exc import IntegrityError
|
||||
from sqlalchemy.ext.asyncio import async_sessionmaker
|
||||
from sqlmodel import col
|
||||
from sqlmodel import select
|
||||
@@ -317,7 +318,10 @@ class SourceService(ServiceBase):
|
||||
"""Create a new job_source execution record in the database."""
|
||||
async with self._session_scope(session) as _session:
|
||||
_session.add(job_source)
|
||||
await self._finalize(session=_session, caller_session=session, refresh=(job_source,))
|
||||
try:
|
||||
await self._finalize(session=_session, caller_session=session, refresh=(job_source,))
|
||||
except IntegrityError as exc:
|
||||
raise self._job_source_conflict(job_id=job_source.job_id, source_id=job_source.source_id) from exc
|
||||
return job_source
|
||||
|
||||
async def read_job_source(self, job_source_id: UUID, *, session: AsyncSession | None = None) -> JobSource:
|
||||
@@ -524,6 +528,10 @@ class SourceService(ServiceBase):
|
||||
if job_source is None:
|
||||
job_source = JobSource(job_id=job_id, source_id=source_id, status=outcome)
|
||||
_session.add(job_source)
|
||||
try:
|
||||
await _session.flush()
|
||||
except IntegrityError as exc:
|
||||
raise self._job_source_conflict(job_id=job_id, source_id=source_id) from exc
|
||||
else:
|
||||
job_source.status = outcome
|
||||
|
||||
@@ -589,6 +597,14 @@ class SourceService(ServiceBase):
|
||||
await self._finalize(session=_session, caller_session=session, refresh=(job, source, job_source, attempt))
|
||||
return job_source
|
||||
|
||||
@staticmethod
|
||||
def _job_source_conflict(*, job_id: UUID, source_id: UUID) -> TranscriptionError:
|
||||
return TranscriptionError(
|
||||
f"Source {source_id} is already linked to Job {job_id}",
|
||||
category=ErrorCategory.CONFLICT,
|
||||
suggestion="Use the existing job-source link instead of creating a duplicate.",
|
||||
)
|
||||
|
||||
async def upsert_revision_for_source(
|
||||
self,
|
||||
*,
|
||||
|
||||
Reference in New Issue
Block a user