32 lines
1.4 KiB
Python
32 lines
1.4 KiB
Python
|
|
"""Selection previews run on the existing durable job runner, outside request transactions."""
|
||
|
|
|
||
|
|
import asyncio
|
||
|
|
from uuid import uuid4
|
||
|
|
|
||
|
|
from sqlalchemy import select
|
||
|
|
|
||
|
|
from ..alphas import sanitize
|
||
|
|
from ..backtests.contracts import fingerprint
|
||
|
|
from ..models import Job, SuperSelectionSnapshot, now
|
||
|
|
from .contracts import SelectionPreview
|
||
|
|
from .evidence import parse_components
|
||
|
|
|
||
|
|
|
||
|
|
async def run_selection(runner, job_id, payload):
|
||
|
|
async with runner.sessions() as db:
|
||
|
|
existing = await db.scalar(select(SuperSelectionSnapshot).where(SuperSelectionSnapshot.job_id == job_id))
|
||
|
|
if existing:
|
||
|
|
return # Restart after snapshot commit must not replace the original observation.
|
||
|
|
request = SelectionPreview.model_validate(payload)
|
||
|
|
raw = await runner.client.run_super_selection(request.platform_query())
|
||
|
|
parsed = parse_components(raw)
|
||
|
|
async with runner.sessions.begin() as db:
|
||
|
|
job = await db.get(Job, job_id)
|
||
|
|
if job.cancel_requested:
|
||
|
|
raise asyncio.CancelledError()
|
||
|
|
snapshot = SuperSelectionSnapshot(id=str(uuid4()), job_id=job_id, source="preview", request=payload,
|
||
|
|
request_hash=fingerprint(request.platform_query()), raw=sanitize(raw), **parsed)
|
||
|
|
db.add(snapshot)
|
||
|
|
job.processed, job.total, job.updated_at = 1, 1, now()
|
||
|
|
job.checkpoint = {"snapshot_id": snapshot.id, "complete": parsed["complete"]}
|