Skip to content
Open
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
8 changes: 8 additions & 0 deletions backend/app/api/docs/assessment/get_dataset_rows.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
Return every row of an assessment dataset as column-keyed JSON.

Fetches and parses the underlying CSV/XLSX file in full, returning the column
`headers`, all `rows` (each a `{column: value}` dict), and `total_rows`. Intended
for the frontend to assemble a directly-usable API-batch input file client-side.

Unlike the `limit_rows` preview on `GET /datasets/{dataset_id}`, this returns the
entire dataset and always fetches the file, so responses scale with dataset size.
76 changes: 76 additions & 0 deletions backend/app/api/routes/assessment/datasets.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
from app.api.deps import AuthContextDep, SessionDep
from app.api.permissions import Permission, require_permission
from app.core.cloud import get_cloud_storage
from app.crud.assessment.batch import list_assessment_dataset_rows
from app.crud.assessment.dataset import (
delete_assessment_dataset,
get_assessment_dataset_by_id,
Expand All @@ -16,6 +17,7 @@
from app.models.assessment import (
AssessmentDatasetPreview,
AssessmentDatasetResponse,
AssessmentDatasetRows,
)
from app.models.evaluation import EvaluationDataset
from app.services.assessment.dataset import (
Expand All @@ -29,6 +31,11 @@

router = APIRouter()

# Hard cap for the rows-export endpoint. Above this the request is rejected
# up front (from the stored row count, no file download/parse) so a huge sheet
# never triggers heavy computation.
MAX_DATASET_EXPORT_ROWS = 100


def _dataset_to_response(
dataset: EvaluationDataset,
Expand Down Expand Up @@ -168,6 +175,75 @@ def get_dataset(
)


@router.get(
"/datasets/{dataset_id}/rows",
description=load_description("assessment/get_dataset_rows.md"),
response_model=APIResponse[AssessmentDatasetRows],
dependencies=[Depends(require_permission(Permission.REQUIRE_PROJECT))],
)
def get_dataset_rows(
dataset_id: int,
session: SessionDep,
auth_context: AuthContextDep,
) -> APIResponse[AssessmentDatasetRows]:
"""Return all rows of an assessment dataset as column-keyed dicts."""
dataset = get_assessment_dataset_by_id(
session=session,
dataset_id=dataset_id,
organization_id=auth_context.organization_.id,
project_id=auth_context.project_.id,
)

# Reject oversized datasets up front from the stored row count (recorded at
# upload), so no file is downloaded or parsed for a sheet we won't return.
total_items = (dataset.dataset_metadata or {}).get("total_items_count")
if isinstance(total_items, int) and total_items > MAX_DATASET_EXPORT_ROWS:
raise HTTPException(
status_code=422,
detail=(
f"Max size exceeded: dataset has {total_items} rows; the input "
f"rows export is limited to {MAX_DATASET_EXPORT_ROWS}."
),
)

try:
rows = list_assessment_dataset_rows(session, dataset)
except ValueError as e:
logger.warning(
"[get_dataset_rows] Failed to load dataset rows | dataset_id=%s | %s",
dataset_id,
e,
)
raise HTTPException(status_code=422, detail=str(e)) from e

# Drop blank-named columns and columns empty across every row (trailing
# empty spreadsheet columns) so the response carries only real fields.
# Single pass over the cells decides which columns to keep (O(rows*cols)),
# so it stays cheap even for large sheets.
ordered_columns: list[str] = []
seen_columns: set[str] = set()
non_empty_columns: set[str] = set()
for row in rows:
for column, value in row.items():
if column not in seen_columns:
seen_columns.add(column)
ordered_columns.append(column)
if column.strip() and (value or "").strip():
non_empty_columns.add(column)
kept_columns = [c for c in ordered_columns if c in non_empty_columns]
cleaned_rows = [
{column: row.get(column, "") for column in kept_columns} for row in rows
]

return APIResponse.success_response(
data=AssessmentDatasetRows(
headers=kept_columns,
rows=cleaned_rows,
total_rows=len(cleaned_rows),
)
)


@router.delete(
"/datasets/{dataset_id}",
description=load_description("assessment/delete_dataset.md"),
Expand Down
32 changes: 23 additions & 9 deletions backend/app/api/routes/assessment/runs.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@

from fastapi import APIRouter, Body, Depends, HTTPException, Query
from fastapi.responses import StreamingResponse
from pydantic import ValidationError

from app.api.deps import AuthContextDep, SessionDep
from app.api.permissions import Permission, require_permission
Expand All @@ -29,6 +30,7 @@
AssessmentRunCreate,
AssessmentRunPublic,
AssessmentRunResponse,
RunExecution,
)
from app.models.evaluation import EvaluationDataset
from app.services.assessment.service import (
Expand Down Expand Up @@ -67,6 +69,17 @@ def _build_run_public(
run.id,
)
dataset = session.get(EvaluationDataset, parent.dataset_id) if parent else None
# Runtime state was folded into the `execution` JSONB bag (migration 078); the frozen
# request input lives on the parent, not the run. Legacy list/get are scoped to
# method=RUN, but guard defensively so a non-RUN-shaped bag can never 500 the list.
try:
bag = RunExecution.model_validate(run.execution or {})
except ValidationError:
logger.warning(
"[_build_run_public] Non-RUN execution bag for run %s; returning empty bag",
run.id,
)
bag = RunExecution()
return AssessmentRunPublic(
id=run.id,
assessment_id=run.assessment_id,
Expand All @@ -78,14 +91,15 @@ def _build_run_public(
status=run.status,
total_items=run.total_items,
error_message=run.error_message,
input=run.input,
prefilter_total_rows=run.prefilter_total_rows,
prefilter_total_passed=run.prefilter_total_passed,
prefilter_total_rejected=run.prefilter_total_rejected,
stage=run.stage,
stage_status=run.stage_status,
pipeline=run.pipeline,
post_processing_config=(run.input or {}).get("post_processing_config"),
input=parent.input if parent else None,
prefilter_total_rows=bag.prefilter_total_rows,
prefilter_total_passed=bag.prefilter_total_passed,
prefilter_total_rejected=bag.prefilter_total_rejected,
stage=bag.stage,
stage_status=bag.stage_status,
pipeline=bag.pipeline,
cost=bag.cost,
post_processing_config=run.post_processing_config,
inserted_at=run.inserted_at,
updated_at=run.updated_at,
)
Expand Down Expand Up @@ -261,7 +275,7 @@ def export_assessment_run_results(
)
)

post_processing_config = (run.input or {}).get("post_processing_config") or None
post_processing_config = run.post_processing_config or None
base_label = assessment.experiment_name if assessment else f"run_{run.id}"

if export_format != "json":
Expand Down
2 changes: 2 additions & 0 deletions backend/app/crud/assessment/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
list_assessment_runs,
list_assessments,
recompute_assessment_status,
resolve_assessment_config_blob,
update_assessment_run_prefilter_stats,
update_assessment_run_status,
update_run_post_processing_config,
Expand Down Expand Up @@ -49,6 +50,7 @@
"list_assessment_datasets",
"list_assessments",
"recompute_assessment_status",
"resolve_assessment_config_blob",
"update_assessment_run_prefilter_stats",
"update_assessment_run_status",
"update_run_post_processing_config",
Expand Down
34 changes: 27 additions & 7 deletions backend/app/crud/assessment/batch.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,9 +28,9 @@
AssessmentRun,
)
from app.models.batch_job import BatchJob, BatchJobType
from app.models.config.assessment_blob import AssessmentConfigBlob
from app.models.evaluation import EvaluationDataset
from app.models.llm.constants import DEFAULT_ASSESSMENT_BATCH_MAX_TOKENS
from app.models.llm.request import ConfigBlob
from app.services.assessment.mappers import (
map_kaapi_to_anthropic_params,
map_kaapi_to_google_params,
Expand Down Expand Up @@ -80,6 +80,18 @@ def _load_dataset_rows(
return _parse_csv_rows(file_content)


def list_assessment_dataset_rows(
session: Session,
dataset: EvaluationDataset,
) -> list[dict[str, str]]:
"""Public entry point for loading all dataset rows as column-keyed dicts.

Thin wrapper over the batch loader so callers outside the batch path (e.g. the
legacy rows endpoint) don't reach into the private helper.
"""
return _load_dataset_rows(session, dataset)


def _parse_csv_rows(content: bytes) -> list[dict[str, str]]:
"""Parse CSV content into list of row dicts."""
for encoding in ("utf-8-sig", "utf-8", "latin-1"):
Expand Down Expand Up @@ -367,7 +379,7 @@ def submit_assessment_batch(
run: AssessmentRun,
assessment: Assessment,
dataset: EvaluationDataset,
config_blob: ConfigBlob,
config_blob: AssessmentConfigBlob,
assessment_input: dict[str, Any],
organization_id: int,
project_id: int,
Expand Down Expand Up @@ -407,14 +419,22 @@ def submit_assessment_batch(
"[submit_assessment_batch] Building JSONL | run_id=%s | rows=%s | provider=%s",
run.id,
len(rows),
config_blob.completion.provider,
config_blob.assessment.provider,
)

# Determine provider and build params
completion = config_blob.completion
provider_name = completion.provider or "openai"

params = dict(completion.params)
assessment_cfg = config_blob.assessment
provider_name = assessment_cfg.provider or LLMProvider.OPENAI

# Normalize assessment params for the mappers, mirroring the BATCH path's
# _stage_params: input_schema is request-validation only (not a provider param),
# and json_output_schema is the mappers' output_schema. instructions stays in
# params — the mappers consume it as the system prompt.
params = dict(assessment_cfg.params)
params.pop("input_schema", None)
json_schema = params.pop("json_output_schema", None)
if json_schema is not None:
params["output_schema"] = json_schema

# Determine the base provider (openai or google)
base_provider = provider_name.replace("-native", "")
Expand Down
Loading
Loading