Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
18 commits
Select commit Hold shift + click to select a range
aaab3af
feat: add disqualified_agents table and disqualified_agent_ids view
jmnmv12 Jul 24, 2026
b33d618
feat: add DisqualifiedAgent Pydantic and ORM models
jmnmv12 Jul 24, 2026
246f556
feat: add disqualify_agent and get_disqualified_agent queries
jmnmv12 Jul 24, 2026
41675a9
feat: add admin endpoint to disqualify an agent from emission
jmnmv12 Jul 24, 2026
ec21cdd
feat: exclude disqualified agents via unified view in ranking queries
jmnmv12 Jul 24, 2026
73f302a
feat: surface disqualified agents on leaderboard and evaluation-run s…
jmnmv12 Jul 24, 2026
d773591
test: verify disqualified agents drop from emission candidates
jmnmv12 Jul 24, 2026
1e7aadd
chore: rename disqualified-agents endpoint test to test_admin.py
jmnmv12 Jul 24, 2026
2e64f86
feat: add disqualification_jobs table, ORM and model
jmnmv12 Jul 27, 2026
f481fee
feat: add disqualification job queries
jmnmv12 Jul 27, 2026
006d076
refactor: extract _apply_incentive_decision from _insert_incentive_ap…
jmnmv12 Jul 27, 2026
cd74dd6
feat: add disqualification reapproval replay
jmnmv12 Jul 27, 2026
d9e09cd
feat: enqueue and drain disqualification reapproval jobs
jmnmv12 Jul 27, 2026
e5f2efd
fix: bound disqualification drain to one pass per invocation and reta…
jmnmv12 Jul 27, 2026
fb79ee1
fix: exclude attempted job ids from claim to prevent starvation behin…
jmnmv12 Jul 27, 2026
daeb2ae
fix: preserve kept leader's approved_at during disqualification reapp…
jmnmv12 Jul 27, 2026
cecd4eb
refactor: :card_file_box: Merge migration files
jmnmv12 Jul 27, 2026
43c657f
feat: update surviving approved agent's reward snapshot against new l…
jmnmv12 Jul 27, 2026
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
89 changes: 89 additions & 0 deletions alembic/versions/2026_07_24_add_disqualified_agents.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,89 @@
"""Add disqualified_agents table, disqualified_agent_ids view, and disqualification_jobs table.

Revision ID: e5c8a1f0b942
Revises: b3f1a9c4d210
Create Date: 2026-07-24 00:00:00.000000

"""

from typing import Sequence, Union

import sqlalchemy as sa

from alembic import op

revision: str = "e5c8a1f0b942"
down_revision: Union[str, Sequence[str], None] = "b3f1a9c4d210"
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None


_CREATE_VIEW = """
CREATE VIEW disqualified_agent_ids AS
SELECT a.agent_id
FROM agents a
JOIN banned_coldkeys bc ON bc.miner_coldkey = a.miner_coldkey
UNION
SELECT agent_id
FROM disqualified_agents;
"""


def upgrade() -> None:
op.create_table(
"disqualified_agents",
sa.Column(
"agent_id",
sa.UUID(),
sa.ForeignKey("agents.agent_id", ondelete="CASCADE"),
primary_key=True,
),
sa.Column("reason", sa.Text(), nullable=False),
sa.Column(
"disqualified_at",
sa.TIMESTAMP(timezone=True),
server_default=sa.text("NOW()"),
nullable=False,
),
)
op.execute(_CREATE_VIEW)

op.create_table(
"disqualification_jobs",
sa.Column(
"id",
sa.UUID(),
server_default=sa.text("gen_random_uuid()"),
primary_key=True,
),
sa.Column(
"agent_id",
sa.UUID(),
sa.ForeignKey("agents.agent_id", ondelete="CASCADE"),
nullable=False,
),
sa.Column("set_id", sa.Integer(), nullable=False),
sa.Column(
"created_at",
sa.TIMESTAMP(timezone=True),
server_default=sa.text("NOW()"),
nullable=False,
),
sa.Column("processed_at", sa.TIMESTAMP(timezone=True), nullable=True),
sa.Column("attempts", sa.Integer(), server_default=sa.text("0"), nullable=False),
sa.Column("error", sa.Text(), nullable=True),
)
op.create_index(
"uq_disqualification_jobs_pending",
"disqualification_jobs",
["agent_id"],
unique=True,
postgresql_where=sa.text("processed_at IS NULL"),
)


def downgrade() -> None:
op.drop_index("uq_disqualification_jobs_pending", table_name="disqualification_jobs")
op.drop_table("disqualification_jobs")
op.execute("DROP VIEW IF EXISTS disqualified_agent_ids")
op.drop_table("disqualified_agents")
2 changes: 2 additions & 0 deletions api/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -207,6 +207,7 @@
PRE_SCREENING_PROJECTOR_POLL_INTERVAL_SECONDS = int(os.getenv("PRE_SCREENING_PROJECTOR_POLL_INTERVAL_SECONDS", "5"))
AUTO_APPROVAL_ENABLED = os.getenv("AUTO_APPROVAL_ENABLED", "false").lower() == "true"
AUTO_APPROVAL_RUN_LOOP = SHOULD_RUN_LOOPS and AUTO_APPROVAL_ENABLED
DISQUALIFICATION_REAPPROVAL_RUN = SHOULD_RUN_LOOPS
AUTO_APPROVAL_POLICY_VERSION = os.getenv("AUTO_APPROVAL_POLICY_VERSION", "approval-v1")
APPROVAL_PROJECTOR_POLL_INTERVAL_SECONDS = int(os.getenv("APPROVAL_PROJECTOR_POLL_INTERVAL_SECONDS", "5"))

Expand Down Expand Up @@ -309,6 +310,7 @@ def _fraction_setting(name: str, default: str) -> float:
logger.info(f"Pre-Screening Projector Poll Interval: {PRE_SCREENING_PROJECTOR_POLL_INTERVAL_SECONDS} second(s)")
logger.info(f"Auto Approval Enabled: {AUTO_APPROVAL_ENABLED}")
logger.info(f"Auto Approval Projector Loop Enabled: {AUTO_APPROVAL_RUN_LOOP}")
logger.info(f"Disqualification Reapproval Enabled: {DISQUALIFICATION_REAPPROVAL_RUN}")
logger.info(f"Approval Projector Poll Interval: {APPROVAL_PROJECTOR_POLL_INTERVAL_SECONDS} second(s)")
logger.info(f"Earliest SET ID with good data: {EARLIEST_SET_ID_WITH_GOOD_DATA}")
logger.info(f"Incentive Start Set ID: {INCENTIVE_START_SET_ID}")
Expand Down
60 changes: 60 additions & 0 deletions api/endpoints/admin.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,8 @@
import asyncio
import logging
import secrets
from typing import Annotated
from uuid import UUID

from bittensor_wallet.keypair import Keypair
from fastapi import APIRouter, Depends, HTTPException, Response, status
Expand All @@ -8,12 +11,24 @@

import api.config as config
from models.banned_coldkey import BannedColdkey
from models.disqualified_agent import DisqualifiedAgent
from queries.agent import get_agent_by_id
from queries.approval import process_pending_disqualification_jobs
from queries.banned_coldkey import ban_coldkey, unban_coldkey
from queries.disqualification_job import enqueue_disqualification_job
from queries.disqualified_agent import disqualify_agent
from utils.database import DatabaseConnection, db_operation
from utils.ttl import clear_all_ttl_caches

logger = logging.getLogger(__name__)

router = APIRouter(tags=["admin"])
admin_bearer = HTTPBearer(auto_error=False)

# Retains references to fire-and-forget drain tasks so they can't be garbage-collected
# before completion (asyncio only holds a weak reference to a bare create_task result).
_background_tasks: set[asyncio.Task[None]] = set()


class ColdkeyBanRequest(BaseModel):
reason: Annotated[str, StringConstraints(strip_whitespace=True, min_length=1, max_length=1000)]
Expand Down Expand Up @@ -62,3 +77,48 @@ async def delete_banned_coldkey(miner_coldkey: str) -> Response:
await unban_coldkey(miner_coldkey)
clear_all_ttl_caches()
return Response(status_code=status.HTTP_204_NO_CONTENT)


@router.put(
"/disqualified-agents/{agent_id}",
response_model=DisqualifiedAgent,
dependencies=[Depends(require_coldkey_ban_admin)],
)
async def put_disqualified_agent(agent_id: UUID, request: ColdkeyBanRequest) -> DisqualifiedAgent:
agent = await get_agent_by_id(agent_id)
if agent is None:
raise HTTPException(status_code=404, detail="Agent not found")

disqualified = await disqualify_agent(agent_id, request.reason)

set_id = await _enqueue_disqualification_job_operation(agent_id=agent_id)
if set_id is not None:
_fire_disqualification_drain()

clear_all_ttl_caches()
return disqualified


@db_operation
async def _enqueue_disqualification_job_operation(conn: DatabaseConnection, *, agent_id: UUID) -> int | None:
"""Enqueue a reapproval job for the agent's set. Returns the set_id, or None if the agent has none."""
async with conn.conn.transaction():
set_id = await conn.fetchval("SELECT set_id FROM agents WHERE agent_id = $1", agent_id)
if set_id is None:
return None
await enqueue_disqualification_job(conn, agent_id=agent_id, set_id=set_id)
return set_id


async def _run_disqualification_drain() -> None:
try:
await process_pending_disqualification_jobs()
except Exception as exc: # noqa: BLE001
logger.error(f"Disqualification drain task failed: {type(exc).__name__}: {exc}")


def _fire_disqualification_drain() -> None:
"""Fire the drain as a background task, retaining a reference until it completes."""
task = asyncio.create_task(_run_disqualification_drain())
_background_tasks.add(task)
task.add_done_callback(_background_tasks.discard)
9 changes: 9 additions & 0 deletions api/src/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@
from api.src.endpoints.upload import router as upload_router
from api.src.middleware.request_interceptor import RequestInterceptorMiddleware
from api.src.utils.sentry import initialize_sentry
from queries.approval import process_pending_disqualification_jobs
from queries.evaluation import set_all_unfinished_evaluation_runs_to_errored
from utils.bittensor import subtensor_client
from utils.database import deinitialize_database, initialize_database
Expand Down Expand Up @@ -96,6 +97,14 @@ async def lifespan(app: FastAPI):
error_message="Platform crashed while running this evaluation"
)

if config.DISQUALIFICATION_REAPPROVAL_RUN:
try:
drained = await process_pending_disqualification_jobs()
if drained:
logger.info(f"Drained {drained} pending disqualification job(s) on startup")
except Exception as exc: # noqa: BLE001
logger.error(f"Startup disqualification drain failed: {type(exc).__name__}: {exc}")

yield

tasks_to_cancel = tuple(background_tasks)
Expand Down
34 changes: 34 additions & 0 deletions db/models/agent.py
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,40 @@ class BannedColdkey(Base):
)


class DisqualifiedAgent(Base):
__tablename__ = "disqualified_agents"

agent_id: Mapped[UUID] = mapped_column(
PG_UUID(as_uuid=True),
sa.ForeignKey("agents.agent_id", ondelete="CASCADE"),
primary_key=True,
)
reason: Mapped[str] = mapped_column(sa.Text, nullable=False)
disqualified_at: Mapped[datetime] = mapped_column(
sa.TIMESTAMP(timezone=True),
nullable=False,
server_default=sa.text("NOW()"),
)


class DisqualificationJob(Base):
__tablename__ = "disqualification_jobs"

id: Mapped[UUID] = mapped_column(
PG_UUID(as_uuid=True), primary_key=True, server_default=sa.text("gen_random_uuid()")
)
agent_id: Mapped[UUID] = mapped_column(
PG_UUID(as_uuid=True), sa.ForeignKey("agents.agent_id", ondelete="CASCADE"), nullable=False
)
set_id: Mapped[int] = mapped_column(sa.Integer, nullable=False)
created_at: Mapped[datetime] = mapped_column(
sa.TIMESTAMP(timezone=True), nullable=False, server_default=sa.text("NOW()")
)
processed_at: Mapped[Optional[datetime]] = mapped_column(sa.TIMESTAMP(timezone=True))
attempts: Mapped[int] = mapped_column(sa.Integer, nullable=False, server_default=sa.text("0"))
error: Mapped[Optional[str]] = mapped_column(sa.Text)


class BenchmarkAgentId(Base):
__tablename__ = "benchmark_agent_ids"

Expand Down
14 changes: 14 additions & 0 deletions models/disqualification_job.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
from datetime import datetime
from uuid import UUID

from pydantic import BaseModel


class DisqualificationJob(BaseModel):
id: UUID
agent_id: UUID
set_id: int
created_at: datetime
processed_at: datetime | None = None
attempts: int = 0
error: str | None = None
10 changes: 10 additions & 0 deletions models/disqualified_agent.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
from datetime import datetime
from uuid import UUID

from pydantic import BaseModel


class DisqualifiedAgent(BaseModel):
agent_id: UUID
reason: str
disqualified_at: datetime
8 changes: 4 additions & 4 deletions queries/agent.py
Original file line number Diff line number Diff line change
Expand Up @@ -544,8 +544,8 @@ async def get_top_agents(conn: DatabaseConnection, number_of_agents: int = 10, p
and ass.agent_id not in (select agent_id from benchmark_agent_ids)
and not exists (
select 1
from banned_coldkeys bc
where bc.miner_coldkey = a.miner_coldkey
from disqualified_agent_ids dq
where dq.agent_id = a.agent_id
)
and ass.status::text <> 'cancelled'
and (
Expand Down Expand Up @@ -590,8 +590,8 @@ async def get_code_hiding_score_cutoff(
AND ass.agent_id NOT IN (SELECT agent_id FROM benchmark_agent_ids)
AND NOT EXISTS (
SELECT 1
FROM banned_coldkeys bc
WHERE bc.miner_coldkey = a.miner_coldkey
FROM disqualified_agent_ids dq
WHERE dq.agent_id = a.agent_id
)
AND review.approval_review_status IS DISTINCT FROM 'rejected'
)
Expand Down
Loading
Loading