diff --git a/.jules/bolt.md b/.jules/bolt.md new file mode 100644 index 0000000..9ebe50c --- /dev/null +++ b/.jules/bolt.md @@ -0,0 +1,3 @@ +## 2024-07-18 - Prevent O(N) memory loading for sequence calculation +**Learning:** Found an anti-pattern where calculating the next sequence number for a job event loaded up to 1000 events into memory via `list_for_job(...)[-1]` just to inspect the final sequence number, causing unnecessary memory consumption and query overhead. +**Action:** Replaced array indexing on list results with direct O(1) SQL aggregation (`func.max(sequence_no)`) via a new `get_max_sequence_no` repository method. Always use database aggregation instead of loading entities into application memory just to compute aggregates. diff --git a/apps/api/src/cortex_api/services/jobs.py b/apps/api/src/cortex_api/services/jobs.py index f99216e..73b6ddc 100644 --- a/apps/api/src/cortex_api/services/jobs.py +++ b/apps/api/src/cortex_api/services/jobs.py @@ -110,8 +110,8 @@ 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 + max_seq = await uow.job_events.get_max_sequence_no(job_id) + next_sequence = max_seq + 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..45d4e31 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 ( @@ -1637,6 +1637,15 @@ async def add(self, record: JobEventRecord) -> JobEventRecord: await self._session.refresh(model) return _job_event_from_model(model) + async def get_max_sequence_no(self, job_id: str) -> int: + """Efficiently gets the maximum sequence number for a job's events.""" + result = await self._session.execute( + select(func.coalesce(func.max(JobEventModel.sequence_no), 0)).where( + JobEventModel.job_id == job_id + ) + ) + return result.scalar_one() + async def list_for_job( self, job_id: str, diff --git a/packages/evaluation/src/cortex_evaluation/jobs.py b/packages/evaluation/src/cortex_evaluation/jobs.py index 6b861f5..0d2a692 100644 --- a/packages/evaluation/src/cortex_evaluation/jobs.py +++ b/packages/evaluation/src/cortex_evaluation/jobs.py @@ -561,8 +561,8 @@ 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 + max_seq = await uow.job_events.get_max_sequence_no(job.job_id) + next_sequence = max_seq + 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..1b66cb3 100644 --- a/packages/knowledge/src/cortex_knowledge/jobs.py +++ b/packages/knowledge/src/cortex_knowledge/jobs.py @@ -480,8 +480,8 @@ 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 + max_seq = await uow.job_events.get_max_sequence_no(job.job_id) + next_sequence = max_seq + 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..5097f1f 100644 --- a/packages/parse/src/cortex_parse/jobs.py +++ b/packages/parse/src/cortex_parse/jobs.py @@ -412,8 +412,8 @@ 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 + max_seq = await uow.job_events.get_max_sequence_no(job.job_id) + next_sequence = max_seq + 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..57d410c 100644 --- a/packages/synthesis/src/cortex_synthesis/jobs.py +++ b/packages/synthesis/src/cortex_synthesis/jobs.py @@ -401,8 +401,8 @@ 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 + max_seq = await uow.job_events.get_max_sequence_no(job.job_id) + next_sequence = max_seq + 1 await uow.job_events.add( JobEventRecord( job_id=job.job_id, diff --git a/workers/evaluation-worker/src/cortex_worker_evaluation/artifacts.py b/workers/evaluation-worker/src/cortex_worker_evaluation/artifacts.py index 6aed23c..d62094e 100644 --- a/workers/evaluation-worker/src/cortex_worker_evaluation/artifacts.py +++ b/workers/evaluation-worker/src/cortex_worker_evaluation/artifacts.py @@ -2,6 +2,8 @@ from cortex_evaluation import ( EvaluationStorageCaller as WorkerStorageCaller, +) +from cortex_evaluation import ( persist_evaluation_report, )