From a55067011f87209a281c26d197dd18c7efd3b8e8 Mon Sep 17 00:00:00 2001 From: "google-labs-jules[bot]" <161369871+google-labs-jules[bot]@users.noreply.github.com> Date: Mon, 20 Jul 2026 20:06:05 +0000 Subject: [PATCH] Optimize JobEvent sequence number generation by using SQL aggregation in get_max_sequence_no Co-authored-by: crabcanon <3458947+crabcanon@users.noreply.github.com> --- apps/api/src/cortex_api/services/jobs.py | 3 +-- packages/db/src/cortex_db/repositories.py | 10 +++++++++- packages/evaluation/src/cortex_evaluation/jobs.py | 3 +-- packages/knowledge/src/cortex_knowledge/jobs.py | 3 +-- packages/parse/src/cortex_parse/jobs.py | 7 ++----- packages/synthesis/src/cortex_synthesis/jobs.py | 3 +-- 6 files changed, 15 insertions(+), 14 deletions(-) diff --git a/apps/api/src/cortex_api/services/jobs.py b/apps/api/src/cortex_api/services/jobs.py index f99216e..1d64866 100644 --- a/apps/api/src/cortex_api/services/jobs.py +++ b/apps/api/src/cortex_api/services/jobs.py @@ -110,8 +110,7 @@ async def cancel_job(uow: CortexUnitOfWork, job_id: str) -> JobStatusDetail: if updated is None: # pragma: no cover - defensive raise NotFoundError(f"Job `{job_id}` was not found.") - existing_events = await uow.job_events.list_for_job(job_id, limit=1000) - next_sequence = existing_events[-1].sequence_no + 1 if existing_events else 1 + next_sequence = await uow.job_events.get_max_sequence_no(job_id) + 1 await uow.job_events.add( JobEventRecord( job_id=job_id, diff --git a/packages/db/src/cortex_db/repositories.py b/packages/db/src/cortex_db/repositories.py index 1b1d253..bdb042b 100644 --- a/packages/db/src/cortex_db/repositories.py +++ b/packages/db/src/cortex_db/repositories.py @@ -49,7 +49,7 @@ SynthesisType, TenantRecord, ) -from sqlalchemy import desc, or_, select, update +from sqlalchemy import desc, func, or_, select, update from sqlalchemy.ext.asyncio import AsyncSession from .models import ( @@ -1651,6 +1651,14 @@ async def list_for_job( ) return [_job_event_from_model(model) for model in result.scalars().all()] + async def get_max_sequence_no(self, job_id: str) -> int: + result = await self._session.execute( + select(func.coalesce(func.max(JobEventModel.sequence_no), 0)).where( + JobEventModel.job_id == job_id + ) + ) + return result.scalar_one() + class ParseRunRepository: def __init__(self, session: AsyncSession) -> None: diff --git a/packages/evaluation/src/cortex_evaluation/jobs.py b/packages/evaluation/src/cortex_evaluation/jobs.py index 6b861f5..ff66a13 100644 --- a/packages/evaluation/src/cortex_evaluation/jobs.py +++ b/packages/evaluation/src/cortex_evaluation/jobs.py @@ -561,8 +561,7 @@ async def _append_event( message: str, details: dict[str, Any] | None = None, ) -> None: - existing_events = await uow.job_events.list_for_job(job.job_id, limit=1000) - next_sequence = existing_events[-1].sequence_no + 1 if existing_events else 1 + next_sequence = await uow.job_events.get_max_sequence_no(job.job_id) + 1 await uow.job_events.add( JobEventRecord( job_id=job.job_id, diff --git a/packages/knowledge/src/cortex_knowledge/jobs.py b/packages/knowledge/src/cortex_knowledge/jobs.py index 63ef80a..1be8170 100644 --- a/packages/knowledge/src/cortex_knowledge/jobs.py +++ b/packages/knowledge/src/cortex_knowledge/jobs.py @@ -480,8 +480,7 @@ async def _append_event( message: str, details: dict[str, Any] | None = None, ) -> None: - existing_events = await uow.job_events.list_for_job(job.job_id, limit=1000) - next_sequence = existing_events[-1].sequence_no + 1 if existing_events else 1 + next_sequence = await uow.job_events.get_max_sequence_no(job.job_id) + 1 await uow.job_events.add( JobEventRecord( job_id=job.job_id, diff --git a/packages/parse/src/cortex_parse/jobs.py b/packages/parse/src/cortex_parse/jobs.py index 63820aa..c0d7fe1 100644 --- a/packages/parse/src/cortex_parse/jobs.py +++ b/packages/parse/src/cortex_parse/jobs.py @@ -33,9 +33,7 @@ async def submit( idempotency_key: str | None = None, request_id: str | None = None, ) -> JobAccepted: - normalized_key = ( - normalize_idempotency_key(idempotency_key) if idempotency_key else None - ) + normalized_key = normalize_idempotency_key(idempotency_key) if idempotency_key else None if normalized_key is not None: existing = await uow.jobs.get_by_idempotency( caller.tenant_id, @@ -412,8 +410,7 @@ async def _append_event( message: str, details: dict[str, Any] | None = None, ) -> None: - existing_events = await uow.job_events.list_for_job(job.job_id, limit=1000) - next_sequence = existing_events[-1].sequence_no + 1 if existing_events else 1 + next_sequence = await uow.job_events.get_max_sequence_no(job.job_id) + 1 await uow.job_events.add( JobEventRecord( job_id=job.job_id, diff --git a/packages/synthesis/src/cortex_synthesis/jobs.py b/packages/synthesis/src/cortex_synthesis/jobs.py index 547cd85..2937cee 100644 --- a/packages/synthesis/src/cortex_synthesis/jobs.py +++ b/packages/synthesis/src/cortex_synthesis/jobs.py @@ -401,8 +401,7 @@ async def _append_event( message: str, details: dict[str, Any] | None = None, ) -> None: - existing_events = await uow.job_events.list_for_job(job.job_id, limit=1000) - next_sequence = existing_events[-1].sequence_no + 1 if existing_events else 1 + next_sequence = await uow.job_events.get_max_sequence_no(job.job_id) + 1 await uow.job_events.add( JobEventRecord( job_id=job.job_id,