Skip to content
Merged
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
93 changes: 91 additions & 2 deletions validator/src/eval_backend/api/routes.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down
108 changes: 68 additions & 40 deletions validator/src/eval_backend/db.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -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"
Expand Down
Loading