diff --git a/validator/src/eval_backend/api/routes.py b/validator/src/eval_backend/api/routes.py index ba57cff..e36448b 100644 --- a/validator/src/eval_backend/api/routes.py +++ b/validator/src/eval_backend/api/routes.py @@ -8,11 +8,19 @@ from uuid import uuid4 from fastapi import APIRouter, Depends, File, Form, Header, HTTPException, Request, UploadFile -from sqlalchemy import select +from sqlalchemy import or_, select from sqlalchemy.orm import Session, selectinload from ..core.config import Settings -from ..models import AdminUser, CompetitionRuntimeConfig, EvaluationRun, JobQueue, Submission, TrainRun +from ..models import ( + AdminUser, + Artifact, + CompetitionRuntimeConfig, + EvaluationRun, + JobQueue, + Submission, + TrainRun, +) from ..schemas import ( AdminLoginRequest, AdminLoginResponse, @@ -1141,6 +1149,87 @@ def admin_create_provider_evaluations( session.close() +@admin_router.delete("/evaluations/{evaluation_id}", status_code=204) +def admin_delete_evaluation( + request: Request, + evaluation_id: int, + user: AdminUser = Depends(_require_admin_user), +) -> None: + session = get_session(request) + try: + run = session.get(EvaluationRun, evaluation_id) + if run is None: + raise HTTPException(status_code=404, detail="evaluation not found") + if run.submission_id is not None: + raise HTTPException( + status_code=400, + detail="only standalone provider evaluations can be deleted here", + ) + session.execute( + Artifact.__table__.delete().where(Artifact.evaluation_id == evaluation_id) + ) + session.delete(run) + session.commit() + except HTTPException: + session.rollback() + raise + except Exception as exc: + session.rollback() + raise HTTPException(status_code=400, detail=str(exc)) from exc + finally: + session.close() + + +@admin_router.delete("/submissions/{submission_id}", status_code=204) +def admin_delete_submission( + request: Request, + submission_id: str, + user: AdminUser = Depends(_require_admin_user), +) -> None: + session = get_session(request) + try: + submission = session.get(Submission, submission_id) + if submission is None: + raise HTTPException(status_code=404, detail="submission not found") + + train_ids = [ + row[0] + for row in session.execute( + select(TrainRun.id).where(TrainRun.submission_id == submission_id) + ) + ] + evaluation_ids = [ + row[0] + for row in session.execute( + select(EvaluationRun.id).where(EvaluationRun.submission_id == submission_id) + ) + ] + + session.execute(JobQueue.__table__.delete().where(JobQueue.submission_id == submission_id)) + + artifact_filters = [Artifact.submission_id == submission_id] + if train_ids: + artifact_filters.append(Artifact.train_id.in_(train_ids)) + if evaluation_ids: + artifact_filters.append(Artifact.evaluation_id.in_(evaluation_ids)) + if submission.submission_artifact_id is not None: + artifact_filters.append(Artifact.id == submission.submission_artifact_id) + session.execute(Artifact.__table__.delete().where(or_(*artifact_filters))) + + submission.submission_artifact_id = None + session.flush() + session.delete(submission) + session.commit() + except HTTPException: + session.rollback() + raise + except Exception as exc: + session.rollback() + raise HTTPException(status_code=400, detail=str(exc)) from exc + finally: + session.close() + + @admin_router.post("/trains", response_model=TrainCreateResponse) def admin_create_train_job( request: Request, diff --git a/validator/src/eval_backend/db.py b/validator/src/eval_backend/db.py index 457d1d9..c6f38ee 100644 --- a/validator/src/eval_backend/db.py +++ b/validator/src/eval_backend/db.py @@ -3,7 +3,7 @@ from collections.abc import Iterator from contextlib import contextmanager -from sqlalchemy import create_engine +from sqlalchemy import create_engine, text from sqlalchemy.engine import make_url from sqlalchemy.orm import Session, declarative_base, sessionmaker @@ -29,48 +29,76 @@ def build_engine(settings: Settings): ) +def _column_exists(conn, table: str, column: str) -> bool: + row = conn.execute( + text( + "SELECT 1 FROM information_schema.columns " + "WHERE table_schema = current_schema() AND table_name = :table AND column_name = :column" + ), + {"table": table, "column": column}, + ).first() + return row is not None + + +def _column_nullable(conn, table: str, column: str) -> bool: + row = conn.execute( + text( + "SELECT is_nullable FROM information_schema.columns " + "WHERE table_schema = current_schema() AND table_name = :table AND column_name = :column" + ), + {"table": table, "column": column}, + ).first() + return bool(row) and row[0] == "YES" + + +def _add_column_if_missing(conn, table: str, column: str, ddl_type: str) -> None: + # Only take the exclusive ALTER TABLE lock when the column is actually + # missing. In steady state (the common case, after the first migration) + # this is a plain catalog SELECT that only needs a non-blocking + # AccessShareLock, so starting a new process no longer has to fight a + # busy job transaction for an exclusive lock it doesn't really need. + if not _column_exists(conn, table, column): + conn.exec_driver_sql(f"ALTER TABLE {table} ADD COLUMN IF NOT EXISTS {column} {ddl_type}") + + +def _drop_not_null_if_needed(conn, table: str, column: str) -> None: + if not _column_nullable(conn, table, column): + conn.exec_driver_sql(f"ALTER TABLE {table} ALTER COLUMN {column} DROP NOT NULL") + + def ensure_schema(engine) -> None: # Minimal, Postgres-only schema upgrades for the running validator. + # Every process (API and worker) runs this at startup. ALTER TABLE takes + # an exclusive lock and Postgres queues later requests (even plain + # SELECTs) behind a pending exclusive lock request, so a long-running job + # transaction here can otherwise block not just this migration but every + # already-running process's unrelated queries. The lock_timeout bounds + # the wait for the rare case a migration is genuinely needed while + # something's busy; the _if_missing/_if_needed checks avoid requesting + # the exclusive lock at all once the schema is already up to date. with engine.begin() as conn: - conn.exec_driver_sql( - "ALTER TABLE submissions ADD COLUMN IF NOT EXISTS miner_id VARCHAR(255)" - ) - conn.exec_driver_sql( - "ALTER TABLE submissions ADD COLUMN IF NOT EXISTS benchmark_names_json JSON" - ) - conn.exec_driver_sql( - "ALTER TABLE submissions ADD COLUMN IF NOT EXISTS submission_artifact_id VARCHAR(36)" - ) - conn.exec_driver_sql( - "ALTER TABLE submissions ADD COLUMN IF NOT EXISTS latest_train_id INTEGER" - ) - conn.exec_driver_sql( - "ALTER TABLE submissions ADD COLUMN IF NOT EXISTS latest_eval_id INTEGER" - ) - conn.exec_driver_sql( - "ALTER TABLE submissions ADD COLUMN IF NOT EXISTS best_eval_id INTEGER" - ) - conn.exec_driver_sql( - "ALTER TABLE submissions ADD COLUMN IF NOT EXISTS finished_at TIMESTAMPTZ" - ) - conn.exec_driver_sql( - "ALTER TABLE submissions ADD COLUMN IF NOT EXISTS duration_seconds DOUBLE PRECISION" - ) - conn.exec_driver_sql("ALTER TABLE submissions ADD COLUMN IF NOT EXISTS cost_usd DOUBLE PRECISION") - conn.exec_driver_sql( - "ALTER TABLE competition_runtime_config ADD COLUMN IF NOT EXISTS default_eval_execution_mode VARCHAR(32)" - ) - conn.exec_driver_sql( - "ALTER TABLE competition_runtime_config ADD COLUMN IF NOT EXISTS default_eval_batch_size INTEGER" - ) - conn.exec_driver_sql( - "ALTER TABLE competition_runtime_config ADD COLUMN IF NOT EXISTS king_score DOUBLE PRECISION" - ) - conn.exec_driver_sql("ALTER TABLE submissions ALTER COLUMN artifact_name DROP NOT NULL") - conn.exec_driver_sql("ALTER TABLE submissions ALTER COLUMN artifact_path DROP NOT NULL") - conn.exec_driver_sql("ALTER TABLE submissions ALTER COLUMN artifact_sha256 DROP NOT NULL") - conn.exec_driver_sql("ALTER TABLE submissions ALTER COLUMN checkpoint_path DROP NOT NULL") - conn.exec_driver_sql("ALTER TABLE submissions ALTER COLUMN benchmark DROP NOT NULL") + conn.exec_driver_sql("SET LOCAL lock_timeout = '5s'") + _add_column_if_missing(conn, "submissions", "miner_id", "VARCHAR(255)") + _add_column_if_missing(conn, "submissions", "benchmark_names_json", "JSON") + _add_column_if_missing(conn, "submissions", "submission_artifact_id", "VARCHAR(36)") + _add_column_if_missing(conn, "submissions", "latest_train_id", "INTEGER") + _add_column_if_missing(conn, "submissions", "latest_eval_id", "INTEGER") + _add_column_if_missing(conn, "submissions", "best_eval_id", "INTEGER") + _add_column_if_missing(conn, "submissions", "finished_at", "TIMESTAMPTZ") + _add_column_if_missing(conn, "submissions", "duration_seconds", "DOUBLE PRECISION") + _add_column_if_missing(conn, "submissions", "cost_usd", "DOUBLE PRECISION") + _add_column_if_missing( + conn, "competition_runtime_config", "default_eval_execution_mode", "VARCHAR(32)" + ) + _add_column_if_missing( + conn, "competition_runtime_config", "default_eval_batch_size", "INTEGER" + ) + _add_column_if_missing(conn, "competition_runtime_config", "king_score", "DOUBLE PRECISION") + _drop_not_null_if_needed(conn, "submissions", "artifact_name") + _drop_not_null_if_needed(conn, "submissions", "artifact_path") + _drop_not_null_if_needed(conn, "submissions", "artifact_sha256") + _drop_not_null_if_needed(conn, "submissions", "checkpoint_path") + _drop_not_null_if_needed(conn, "submissions", "benchmark") conn.exec_driver_sql( "UPDATE submissions SET miner_id = COALESCE(miner_id, team_name) " "WHERE miner_id IS NULL AND team_name IS NOT NULL"