Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions .jules/bolt.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@

## 2024-07-27 - Optimize Job Event Sequence Calculation
**Learning:** Loading lists of entities into memory (e.g., using `list_for_job(limit=1000)[-1]`) just to determine the maximum sequence number is an anti-pattern that behaves like an N+1 query memory leak.
**Action:** Always use direct SQL aggregations (e.g., `select(func.coalesce(func.max(Model.field), 0))`) in repository methods to fetch scalar maximums in O(1) time rather than loading unneeded objects into memory.
4 changes: 2 additions & 2 deletions apps/api/src/cortex_api/services/jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
10 changes: 9 additions & 1 deletion packages/db/src/cortex_db/repositories.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand Down Expand Up @@ -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:
Expand Down
4 changes: 2 additions & 2 deletions packages/evaluation/src/cortex_evaluation/jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
4 changes: 2 additions & 2 deletions packages/knowledge/src/cortex_knowledge/jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
8 changes: 3 additions & 5 deletions packages/parse/src/cortex_parse/jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -412,8 +410,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,
Expand Down
4 changes: 2 additions & 2 deletions packages/synthesis/src/cortex_synthesis/jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Loading