|
| 1 | +"""Admin routes for the evaluation page. |
| 2 | +
|
| 3 | +Datasets are uploaded once and replayed by runs. Every route is admin-only: |
| 4 | +a run indexes a corpus, spends grader tokens, and occupies the single runner |
| 5 | +slot, so it is not something a partition editor should be able to trigger. |
| 6 | +""" |
| 7 | + |
| 8 | +from __future__ import annotations |
| 9 | + |
| 10 | +from dataclasses import asdict |
| 11 | +from typing import Any |
| 12 | + |
| 13 | +from api.dependencies.auth import current_user, require_admin |
| 14 | +from api.schemas.admin.evaluation_schemas import ( |
| 15 | + EvalDatasetResponse, |
| 16 | + EvalRunResponse, |
| 17 | + EvalRunSummaryResponse, |
| 18 | + StartRunRequest, |
| 19 | +) |
| 20 | +from di.providers import get_evaluation_service |
| 21 | +from fastapi import APIRouter, Depends, File, Form, UploadFile, status |
| 22 | + |
| 23 | +router = APIRouter(dependencies=[Depends(require_admin)]) |
| 24 | + |
| 25 | + |
| 26 | +def _run_summary(run: Any) -> EvalRunSummaryResponse: |
| 27 | + return EvalRunSummaryResponse( |
| 28 | + id=run.id, |
| 29 | + dataset_id=run.dataset_id, |
| 30 | + status=run.status.value, |
| 31 | + started_at=run.started_at, |
| 32 | + finished_at=run.finished_at, |
| 33 | + hit_rate=run.retrieval.hit_rate if run.retrieval else None, |
| 34 | + mrr=run.retrieval.mrr if run.retrieval else None, |
| 35 | + answer_pass_rate=run.answer.pass_rate if run.answer else None, |
| 36 | + files_per_minute=run.indexing.files_per_minute if run.indexing else None, |
| 37 | + error=run.error, |
| 38 | + ) |
| 39 | + |
| 40 | + |
| 41 | +def _run_detail(run: Any) -> EvalRunResponse: |
| 42 | + return EvalRunResponse( |
| 43 | + id=run.id, |
| 44 | + dataset_id=run.dataset_id, |
| 45 | + status=run.status.value, |
| 46 | + started_at=run.started_at, |
| 47 | + finished_at=run.finished_at, |
| 48 | + indexing=asdict(run.indexing) if run.indexing else None, |
| 49 | + retrieval=asdict(run.retrieval) if run.retrieval else None, |
| 50 | + answer=asdict(run.answer) if run.answer else None, |
| 51 | + cases=[asdict(case) for case in run.cases], |
| 52 | + error=run.error, |
| 53 | + created_by=run.created_by, |
| 54 | + ) |
| 55 | + |
| 56 | + |
| 57 | +@router.get("/datasets", response_model=list[EvalDatasetResponse]) |
| 58 | +async def list_datasets(service=Depends(get_evaluation_service)): |
| 59 | + """List stored evaluation datasets, newest first.""" |
| 60 | + return [asdict(dataset) for dataset in await service.list_datasets()] |
| 61 | + |
| 62 | + |
| 63 | +@router.post( |
| 64 | + "/datasets", |
| 65 | + response_model=EvalDatasetResponse, |
| 66 | + status_code=status.HTTP_201_CREATED, |
| 67 | +) |
| 68 | +async def create_dataset( |
| 69 | + name: str = Form(..., description="Human-readable dataset name"), |
| 70 | + testset: UploadFile = File(..., description="CSV: question,expected_answer,expected_file_ids"), |
| 71 | + corpus: list[UploadFile] = File(..., description="Documents to index for the run"), |
| 72 | + user=Depends(current_user), |
| 73 | + service=Depends(get_evaluation_service), |
| 74 | +): |
| 75 | + """Upload a corpus and its test set. |
| 76 | +
|
| 77 | + The CSV is validated here, so a bad test set fails now rather than after a |
| 78 | + run has already indexed the corpus. |
| 79 | + """ |
| 80 | + # Pass the open streams, not the bytes: large uploads are already spooled |
| 81 | + # to disk, and reading them here would pull them into memory unbounded. |
| 82 | + dataset = await service.create_dataset( |
| 83 | + name=name, |
| 84 | + corpus=[(upload.filename or "unnamed", upload.file) for upload in corpus], |
| 85 | + testset=testset.file, |
| 86 | + user_id=user.get("id") if isinstance(user, dict) else None, |
| 87 | + ) |
| 88 | + return asdict(dataset) |
| 89 | + |
| 90 | + |
| 91 | +@router.delete("/datasets/{dataset_id}", status_code=status.HTTP_204_NO_CONTENT) |
| 92 | +async def delete_dataset(dataset_id: str, service=Depends(get_evaluation_service)): |
| 93 | + """Delete a dataset and its stored files.""" |
| 94 | + await service.delete_dataset(dataset_id) |
| 95 | + |
| 96 | + |
| 97 | +@router.get("/runs", response_model=list[EvalRunSummaryResponse]) |
| 98 | +async def list_runs(limit: int = 50, service=Depends(get_evaluation_service)): |
| 99 | + """Run history, newest first.""" |
| 100 | + return [_run_summary(run) for run in await service.list_runs(limit)] |
| 101 | + |
| 102 | + |
| 103 | +@router.post("/runs", response_model=EvalRunResponse, status_code=status.HTTP_202_ACCEPTED) |
| 104 | +async def start_run( |
| 105 | + body: StartRunRequest, |
| 106 | + user=Depends(current_user), |
| 107 | + service=Depends(get_evaluation_service), |
| 108 | +): |
| 109 | + """Queue a run against a dataset. |
| 110 | +
|
| 111 | + Returns ``409`` when a run is already in flight — runs execute one at a |
| 112 | + time so that indexing timings stay comparable between them. |
| 113 | + """ |
| 114 | + run = await service.start_run( |
| 115 | + body.dataset_id, |
| 116 | + user.get("id") if isinstance(user, dict) else None, |
| 117 | + ) |
| 118 | + return _run_detail(run) |
| 119 | + |
| 120 | + |
| 121 | +@router.get("/runs/{run_id}", response_model=EvalRunResponse) |
| 122 | +async def get_run(run_id: str, service=Depends(get_evaluation_service)): |
| 123 | + """One run with its metrics and per-question detail.""" |
| 124 | + return _run_detail(await service.get_run(run_id)) |
| 125 | + |
| 126 | + |
| 127 | +@router.post("/runs/{run_id}/cancel", response_model=EvalRunResponse) |
| 128 | +async def cancel_run(run_id: str, service=Depends(get_evaluation_service)): |
| 129 | + """Ask the runner to abandon an in-flight run.""" |
| 130 | + return _run_detail(await service.cancel_run(run_id)) |
0 commit comments