diff --git a/backend/app/models/assessment/__init__.py b/backend/app/models/assessment/__init__.py index 019ab08c8..0e52cdc37 100644 --- a/backend/app/models/assessment/__init__.py +++ b/backend/app/models/assessment/__init__.py @@ -38,6 +38,7 @@ AssessmentOutput, AssessmentResult, AssessmentResultData, + AssessmentResultFiles, AssessmentResultRow, AssessmentSubmitResponse, AssessmentSummary, @@ -99,6 +100,7 @@ "AssessmentCreate", "AssessmentSubmitResponse", "AssessmentCallback", + "AssessmentResultFiles", "PreFilterVerdict", "PreFilter", "AssessmentOutput", diff --git a/backend/app/models/assessment/assessment_api.py b/backend/app/models/assessment/assessment_api.py index 8ffcaf0f3..38229d4f9 100644 --- a/backend/app/models/assessment/assessment_api.py +++ b/backend/app/models/assessment/assessment_api.py @@ -8,7 +8,7 @@ from typing import Annotated, Any, NotRequired, TypedDict from uuid import UUID -from pydantic import BaseModel, ConfigDict, Field, HttpUrl, model_validator +from pydantic import BaseModel, ConfigDict, Field, HttpUrl, JsonValue, model_validator from sqlmodel import SQLModel from app.models.assessment.assessment import ( @@ -140,7 +140,7 @@ class AssessmentCreate(BaseModel): "GET /assessments/{assessment_id} instead" ), ) - request_metadata: dict[str, Any] | None = Field( + request_metadata: dict[str, JsonValue] | None = Field( default=None, description="Passed through unchanged in the callback for correlation", ) @@ -260,4 +260,12 @@ class AssessmentCallback(BaseModel): assessment_id: UUID status: AssessmentStatus data: AssessmentResultData | None = None - request_metadata: dict[str, Any] | None = None + request_metadata: dict[str, JsonValue] | None = None + + +class AssessmentResultFiles(BaseModel): + """URLs of the provider dumps for each stage, if any.""" + + topic_relevance: str | None = None + assessment: str | None = None + errors: str | None = None diff --git a/backend/app/services/assessment/api/callbacks.py b/backend/app/services/assessment/api/callbacks.py index dbea7341c..137a34837 100644 --- a/backend/app/services/assessment/api/callbacks.py +++ b/backend/app/services/assessment/api/callbacks.py @@ -6,8 +6,8 @@ """ import logging -from typing import Any +from pydantic import JsonValue from sqlmodel import Session from app.models.assessment import ( @@ -16,7 +16,7 @@ AssessmentCallback, AssessmentStatus, ) -from app.services.assessment.api.result_files import build_callback_metadata +from app.services.assessment.api.result_files import presign_result_files from app.utils import get_webhook_secret, send_callback logger = logging.getLogger(__name__) @@ -28,7 +28,7 @@ def deliver( assessment: Assessment, result: AssessmentBatchResult, callback_url: str, - request_metadata: dict[str, Any] | None, + request_metadata: dict[str, JsonValue] | None, failure_message: str | None, ) -> bool: """POST the assessment result to ``callback_url`` (HMAC-signed). Returns whether it was sent. @@ -47,7 +47,7 @@ def deliver( ) try: - metadata = build_callback_metadata(session=session, assessment=assessment) + files = presign_result_files(session=session, assessment=assessment) except Exception: # A metadata bug must never cost the client its result. logger.error( @@ -55,7 +55,7 @@ def deliver( assessment.id, exc_info=True, ) - metadata = None + files = None sent = send_callback( callback_url, @@ -63,15 +63,15 @@ def deliver( "success": assessment.status != AssessmentStatus.FAILED, "data": callback.model_dump(mode="json"), "error": failure_message, - "metadata": metadata, + "metadata": {"files": files.model_dump()} if files else None, }, webhook_secret=webhook_secret, ) logger.info( - "[deliver] Callback %s | assessment_id=%s | status=%s | result_files=%s", + "[deliver] Callback %s | assessment_id=%s | status=%s | files=%s", "sent" if sent else "failed", assessment.id, assessment.status, - sorted((metadata or {}).get("result_files", {})), + files, ) return sent diff --git a/backend/app/services/assessment/api/result_files.py b/backend/app/services/assessment/api/result_files.py index 582ded0e1..7a08ae56d 100644 --- a/backend/app/services/assessment/api/result_files.py +++ b/backend/app/services/assessment/api/result_files.py @@ -1,12 +1,11 @@ """Durable result dumps for the BATCH API-client path. -Records every provider dump on ``assessment.result_files`` as ``{kind: {object_store_url}}``, -builds ``errors.jsonl`` at terminal time, and presigns both into the callback envelope. -Nothing here raises into the terminal path: a missing dump degrades to a missing key. +Records every provider dump on ``assessment.result_files`` as ``{stage: {object_store_url}}``, +builds ``errors.jsonl`` only when there is something to report, and presigns both into the +callback envelope. Nothing here raises: a missing dump degrades to a missing key. """ import logging -from datetime import timedelta from enum import StrEnum from typing import Any @@ -15,17 +14,16 @@ from app.core.cloud import get_cloud_storage from app.core.config import settings from app.core.storage_utils import upload_jsonl_to_object_store -from app.core.util import now from app.crud.assessment import api from app.crud.job import get_batch_job from app.models.assessment import ( Assessment, + AssessmentResultFiles, AssessmentRun, BatchRunState, ) from app.models.batch_job import BatchJob from app.services.assessment.api.batch import ( - ApiStage, _build_batch_provider, error_file_entries, ) @@ -36,7 +34,6 @@ logger = logging.getLogger(__name__) -RESULTS_FILE_KIND = "results" ERRORS_FILE_KIND = "errors" # 86400 is the storage layer's own ceiling, so the presigned urls live exactly one day. @@ -52,13 +49,6 @@ class ErrorRecordEnum(StrEnum): PROVIDER_ERROR_FILE_UNAVAILABLE = "provider_error_file_unavailable" -def stage_file_kind(stage: str) -> str: - """Result-file kind for a stage's dump; the assessment dump is the run's ``results``.""" - if stage == ApiStage.ASSESSMENT.value: - return RESULTS_FILE_KIND - return f"{stage}_results" - - def record_stage_dump( *, session: Session, @@ -81,7 +71,7 @@ def record_stage_dump( api.set_result_files( session=session, assessment=assessment, - files={stage_file_kind(stage): {"object_store_url": url}}, + files={stage: {"object_store_url": url}}, ) @@ -162,11 +152,8 @@ def build_and_upload_errors( bag: BatchRunState, failure_message: str | None, ) -> str | None: - """Assemble and upload the run's ``errors.jsonl``. Returns its object-store url. - - Uploaded even when there are no rows, so the "both a results and an errors url" - promise holds on the clean-success path too. - """ + """Assemble and upload the run's ``errors.jsonl``; ``None`` when there is nothing + to report, so a clean run leaves no empty object behind.""" rows: list[dict[str, Any]] = [] if failure_message: rows.append( @@ -193,6 +180,9 @@ def build_and_upload_errors( ) ) + if not rows: + return None + try: storage = get_cloud_storage(session=session, project_id=assessment.project_id) except Exception: @@ -239,7 +229,7 @@ def finalize_result_files( for stage, url in (bag.get("stage_output_urls") or {}).items(): if not url: continue - files[stage_file_kind(stage)] = {"object_store_url": url} + files[stage] = {"object_store_url": url} errors_url = build_and_upload_errors( session=session, @@ -263,27 +253,25 @@ def finalize_result_files( ) -def build_callback_metadata( +def presign_result_files( *, session: Session, assessment: Assessment -) -> dict[str, Any]: +) -> AssessmentResultFiles: """Presign every recorded result file for the callback envelope's ``metadata``. - Always returns ``{"result_files": ..., "expires_at": ...}``; a per-key presign - failure drops that entry rather than the whole envelope key. + An unset field means no such dump; a presign failure leaves just that field null. """ - expires_at = (now() + timedelta(seconds=SIGNED_URL_EXPIRY_SECONDS)).isoformat() - signed: dict[str, dict[str, Any]] = {} + files: dict[str, str] = {} try: storage = get_cloud_storage(session=session, project_id=assessment.project_id) except Exception: logger.error( - "[build_callback_metadata] Storage unavailable, sending empty result_files | " + "[presign_result_files] Storage unavailable, sending empty files | " "assessment_id=%s", assessment.id, exc_info=True, ) - return {"result_files": signed, "expires_at": expires_at} + return AssessmentResultFiles() for kind, record in assessment.result_files.items(): entry: dict[str, Any] = record or {} @@ -296,13 +284,13 @@ def build_callback_metadata( ) except Exception: logger.error( - "[build_callback_metadata] Presign failed, dropping kind | " + "[presign_result_files] Presign failed, dropping kind | " "assessment_id=%s | kind=%s", assessment.id, kind, exc_info=True, ) continue - signed[kind] = {"signed_url": signed_url} + files[kind] = signed_url - return {"result_files": signed, "expires_at": expires_at} + return AssessmentResultFiles.model_validate(files) diff --git a/backend/app/tests/assessment/test_callbacks.py b/backend/app/tests/assessment/test_callbacks.py new file mode 100644 index 000000000..95c0c5223 --- /dev/null +++ b/backend/app/tests/assessment/test_callbacks.py @@ -0,0 +1,75 @@ +"""Tests for webhook delivery (app/services/assessment/api/callbacks.py).""" + +from unittest.mock import patch + +from app.models.assessment import AssessmentBatchResult +from app.services.assessment.api.callbacks import deliver +from app.tests.assessment.test_result_files import _seed +from app.tests.utils.auth import get_user_test_auth_context + + +def _result() -> AssessmentBatchResult: + return AssessmentBatchResult(total_items=1) + + +class TestDeliver: + def test_sends_presigned_files_in_metadata(self, db) -> None: + auth = get_user_test_auth_context(db) + assessment, _ = _seed(db, auth) + + with ( + patch( + "app.services.assessment.api.callbacks.get_webhook_secret", + return_value=None, + ), + patch( + "app.services.assessment.api.callbacks.presign_result_files", + return_value=None, + ) as presign, + patch( + "app.services.assessment.api.callbacks.send_callback", + return_value=True, + ) as send_callback, + ): + sent = deliver( + session=db, + assessment=assessment, + result=_result(), + callback_url="https://client.example/webhook", + request_metadata=None, + failure_message=None, + ) + + assert sent is True + presign.assert_called_once_with(session=db, assessment=assessment) + assert send_callback.call_args.args[1]["metadata"] is None + + def test_a_metadata_bug_still_delivers_the_result(self, db) -> None: + auth = get_user_test_auth_context(db) + assessment, _ = _seed(db, auth) + + with ( + patch( + "app.services.assessment.api.callbacks.get_webhook_secret", + return_value=None, + ), + patch( + "app.services.assessment.api.callbacks.presign_result_files", + side_effect=RuntimeError("s3 unreachable"), + ), + patch( + "app.services.assessment.api.callbacks.send_callback", + return_value=True, + ) as send_callback, + ): + sent = deliver( + session=db, + assessment=assessment, + result=_result(), + callback_url="https://client.example/webhook", + request_metadata=None, + failure_message=None, + ) + + assert sent is True + assert send_callback.call_args.args[1]["metadata"] is None diff --git a/backend/app/tests/assessment/test_result_files.py b/backend/app/tests/assessment/test_result_files.py index 3a9629a1e..f30423f4e 100644 --- a/backend/app/tests/assessment/test_result_files.py +++ b/backend/app/tests/assessment/test_result_files.py @@ -5,10 +5,8 @@ """ import json -from datetime import datetime, timedelta from unittest.mock import MagicMock, patch -from app.core.util import now from app.crud.assessment import api from app.models.assessment import AssessmentMethod, BatchRunState from app.models.batch_job import BatchJob, BatchJobType @@ -17,10 +15,9 @@ from app.services.assessment.api.batch import ApiStage from app.services.assessment.api.result_files import ( build_and_upload_errors, - build_callback_metadata, finalize_result_files, + presign_result_files, record_stage_dump, - stage_file_kind, ) from app.tests.utils.auth import get_user_test_auth_context from app.tests.utils.test_data import create_test_config @@ -136,16 +133,6 @@ def _upload_patch(uploads: _Uploads): ) -class TestStageFileKind: - def test_assessment_stage_is_the_runs_results(self) -> None: - assert stage_file_kind(ApiStage.ASSESSMENT.value) == "results" - - def test_prefilter_stage_is_suffixed(self) -> None: - assert stage_file_kind(ApiStage.TOPIC_RELEVANCE.value) == ( - "topic_relevance_results" - ) - - class TestRecordStageDump: def test_dump_is_on_the_parent_row_before_any_terminal_state(self, db) -> None: auth = get_user_test_auth_context(db) @@ -160,7 +147,7 @@ def test_dump_is_on_the_parent_row_before_any_terminal_state(self, db) -> None: db.refresh(assessment) assert assessment.result_files == { - "topic_relevance_results": { + "topic_relevance": { "object_store_url": "s3://bucket/batch-1170/output.jsonl", } } @@ -181,7 +168,7 @@ def test_missing_url_records_nothing(self, db) -> None: class TestBuildAndUploadErrors: - def test_clean_success_still_uploads_an_empty_file(self, db) -> None: + def test_clean_success_uploads_nothing(self, db) -> None: auth = get_user_test_auth_context(db) assessment, execution = _seed(db, auth) uploads = _Uploads() @@ -195,10 +182,8 @@ def test_clean_success_still_uploads_an_empty_file(self, db) -> None: failure_message=None, ) - assert url == "s3://bucket/errors.jsonl" - assert uploads.rows == [] - assert uploads.calls[0]["filename"] == "errors.jsonl" - assert uploads.calls[0]["subdirectory"] == f"assessment/{assessment.id}" + assert url is None + assert uploads.calls == [] def test_row_errors_are_flattened_per_stage(self, db) -> None: auth = get_user_test_auth_context(db) @@ -357,7 +342,7 @@ def test_pre_provider_failure_records_only_a_synthetic_execution_error( } ] - def test_completed_run_carries_both_a_results_and_an_errors_record( + def test_completed_run_with_no_errors_carries_only_the_results_record( self, db ) -> None: auth = get_user_test_auth_context(db) @@ -384,8 +369,7 @@ def test_completed_run_carries_both_a_results_and_an_errors_record( db.refresh(assessment) assert assessment.result_files == { - "results": {"object_store_url": "s3://bucket/batch-1173/output.jsonl"}, - "errors": {"object_store_url": "s3://bucket/errors.jsonl"}, + "assessment": {"object_store_url": "s3://bucket/batch-1173/output.jsonl"}, } def test_prefilter_and_assessment_dumps_coexist(self, db) -> None: @@ -414,11 +398,9 @@ def test_prefilter_and_assessment_dumps_coexist(self, db) -> None: db.refresh(assessment) assert set(assessment.result_files) == { - "topic_relevance_results", - "results", - "errors", + "topic_relevance", + "assessment", } - assert "topic_relevance_results" in assessment.result_files def test_a_second_tick_does_not_duplicate_or_lose_records(self, db) -> None: auth = get_user_test_auth_context(db) @@ -439,7 +421,7 @@ def test_a_second_tick_does_not_duplicate_or_lose_records(self, db) -> None: ) db.refresh(assessment) - assert set(assessment.result_files) == {"results", "errors"} + assert set(assessment.result_files) == {"assessment"} def test_upload_failure_does_not_raise_into_the_terminal_path(self, db) -> None: auth = get_user_test_auth_context(db) @@ -462,13 +444,14 @@ def test_upload_failure_does_not_raise_into_the_terminal_path(self, db) -> None: ApiStage.ASSESSMENT.value: "s3://bucket/out.jsonl" } ), + failure_message="kaboom", ) db.refresh(assessment) assert assessment.result_files == {} -class TestBuildCallbackMetadata: +class TestPresignResultFiles: def _signing_storage(self, failing_url: str | None = None) -> MagicMock: storage = MagicMock() @@ -487,27 +470,18 @@ def test_every_kind_is_signed(self, db) -> None: session=db, assessment=assessment, files={ - "results": {"object_store_url": "s3://bucket/out.jsonl"}, + "assessment": {"object_store_url": "s3://bucket/out.jsonl"}, "errors": {"object_store_url": "s3://bucket/errors.jsonl"}, }, ) with _storage_patch(self._signing_storage()): - metadata = build_callback_metadata(session=db, assessment=assessment) - - assert metadata["result_files"]["results"] == { - "signed_url": f"https://signed.example/s3://bucket/out.jsonl?exp={ONE_DAY_SECONDS}", - } - - def test_expires_at_is_one_day_out(self, db) -> None: - auth = get_user_test_auth_context(db) - assessment, _ = _seed(db, auth) - - with _storage_patch(self._signing_storage()): - metadata = build_callback_metadata(session=db, assessment=assessment) + files = presign_result_files(session=db, assessment=assessment) - expires_at = datetime.fromisoformat(metadata["expires_at"]) - assert timedelta(hours=23, minutes=59) < expires_at - now() <= timedelta(days=1) + assert ( + files.assessment + == f"https://signed.example/s3://bucket/out.jsonl?exp={ONE_DAY_SECONDS}" + ) def test_a_failing_presign_drops_only_its_own_kind(self, db) -> None: auth = get_user_test_auth_context(db) @@ -516,31 +490,30 @@ def test_a_failing_presign_drops_only_its_own_kind(self, db) -> None: session=db, assessment=assessment, files={ - "results": {"object_store_url": "s3://bucket/out.jsonl"}, + "assessment": {"object_store_url": "s3://bucket/out.jsonl"}, "errors": {"object_store_url": "s3://bucket/errors.jsonl"}, }, ) with _storage_patch(self._signing_storage(failing_url="s3://bucket/out.jsonl")): - metadata = build_callback_metadata(session=db, assessment=assessment) + files = presign_result_files(session=db, assessment=assessment) - assert set(metadata["result_files"]) == {"errors"} - assert metadata["expires_at"] + assert files.assessment is None + assert files.errors is not None - def test_storage_outage_still_returns_the_envelope_keys(self, db) -> None: + def test_storage_outage_returns_an_empty_envelope(self, db) -> None: auth = get_user_test_auth_context(db) assessment, _ = _seed(db, auth) api.set_result_files( session=db, assessment=assessment, - files={"results": {"object_store_url": "s3://bucket/out.jsonl"}}, + files={"assessment": {"object_store_url": "s3://bucket/out.jsonl"}}, ) with patch( "app.services.assessment.api.result_files.get_cloud_storage", side_effect=RuntimeError("s3 unreachable"), ): - metadata = build_callback_metadata(session=db, assessment=assessment) + files = presign_result_files(session=db, assessment=assessment) - assert metadata["result_files"] == {} - assert metadata["expires_at"] + assert files.model_dump(exclude_none=True) == {} diff --git a/docs/wiki/modules/assessment.md b/docs/wiki/modules/assessment.md index bb09c0e47..dd3c1511d 100644 --- a/docs/wiki/modules/assessment.md +++ b/docs/wiki/modules/assessment.md @@ -21,9 +21,9 @@ Config version (tag=ASSESSMENT, `models/config/assessment_blob.py`) owns system `assessment.id` is a **UUID** (like config/job/llm_call). Per-item result = `AssessmentResult {output: {assessment, pre_filter}, error}` (no `metadata` — the provider/model/usage block was removed from the API-client output) where `output.assessment` = the LLM output parsed to an object when the config has a `json_output_schema`, else string (null for gated/failed rows), and `output.pre_filter` holds the `{topic_relevance}` verdict (`{verdict, reasoning}` or null) and is itself null when no pre-filter ran. The API-client BATCH rows do **not** live in `assessment.input` (that column is now RESPONSE/RUN only, NULL for BATCH). They are uploaded to `submission.jsonl` at submit and `assessment.submission_input` holds its `s3://` url: 3-6MB of JSONB was dragged along by every full-row `SELECT` of the assessment. `services/assessment/api/submission_store.py` owns the round trip: `upload_submission_rows` at submit, and the `open_submission_rows` context manager on read, which streams the JSONL line by line and closes the body on exit. Rows are read only inside `_submit_stage`, never on a poll, and only the rows the stage's subset needs are held (the count comes from `execution.total_items`, so an empty subset never opens the file). A storage read failure raises `SubmissionUnavailableError`, which requeues the task instead of failing the execution; a corrupt line is a `ValueError` and terminal. An upload failure at submit is a 503 (an assessment without its rows is unrunnable). `build_result` reads `execution.total_items` rather than re-deriving the count from the rows, so the terminal path never fetches them. -`assessment.result_files` is the durable record of every provider batch dump held: JSONB, **NOT NULL default `{}`**, CHECK-constrained to a JSON object (`ck_assessment_result_files_is_object`), keyed by file kind (`results` / `errors` / `_results`, derived by `stage_file_kind`) with `{object_store_url}` per kind. It stores the raw `s3://` url (presigning is per-delivery, so the column never holds an expiring link); writes go through `crud/assessment/api.py::set_result_files`, which merges **server-side** (`result_files || :files::jsonb`) rather than read-modify-write, because two drivers can touch the row in the same second. `batch_job.provider_error_file_id` persists OpenAI's `error_file_id` (OpenAI only; Anthropic and Gemini report per-item errors inline) so the error dump is still fetchable at terminal time, in a later task than the poll that surfaced it. +`assessment.result_files` is the durable record of every provider batch dump held: JSONB, **NOT NULL default `{}`**, CHECK-constrained to a JSON object (`ck_assessment_result_files_is_object`), keyed by stage name (`assessment` / `topic_relevance`, the `ApiStage` values) plus `errors`, with `{object_store_url}` per key. It stores the raw `s3://` url (presigning is per-delivery, so the column never holds an expiring link); writes go through `crud/assessment/api.py::set_result_files`, which merges **server-side** (`result_files || :files::jsonb`) rather than read-modify-write, because two drivers can touch the row in the same second. `batch_job.provider_error_file_id` persists OpenAI's `error_file_id` (OpenAI only; Anthropic and Gemini report per-item errors inline) so the error dump is still fetchable at terminal time, in a later task than the poll that surfaced it. -Delivery is poll-or-push: the `POST /assessments` ack is the flat `AssessmentSubmitResponse {assessment_id, status, message, inserted_at, updated_at}`, the result is always readable from `GET /assessments/{assessment_id}`, and when the request carried a `callback_url` the `AssessmentCallback {assessment_id, status, data, request_metadata}` is also POSTed there on completion — where `data` is a single `AssessmentResult` (RESPONSE) or an `AssessmentBatchResult {total_items, counts, items}` (BATCH); `status` lives on the envelope only. The outer `send_callback` envelope carries `{success, data, error, metadata}`: `error` is `_fail`'s failure message (null on a success path) and `metadata` is `{result_files: {kind: {signed_url}}, expires_at}` — the column's `object_store_url` presigned for 1 day at delivery time, never stored, so a client whose receiver rejects the (large, item-inlining) body can still fetch the dumps. Both were hardcoded `None` before. Delivery is one inline attempt with no retry, so a rejected callback is still lost. Pre-filter `stop_on_fail` flag (`config/assessment_blob.py`) drives which filters hard-stop the chain on a failing verdict vs pass-through (record only). +Delivery is poll-or-push: the `POST /assessments` ack is the flat `AssessmentSubmitResponse {assessment_id, status, message, inserted_at, updated_at}`, the result is always readable from `GET /assessments/{assessment_id}`, and when the request carried a `callback_url` the `AssessmentCallback {assessment_id, status, data, request_metadata}` is also POSTed there on completion — where `data` is a single `AssessmentResult` (RESPONSE) or an `AssessmentBatchResult {total_items, counts, items}` (BATCH); `status` lives on the envelope only. The outer `send_callback` envelope carries `{success, data, error, metadata}`: `error` is `_fail`'s failure message (null on a success path) and `metadata` is `{files: AssessmentResultFiles}` — a flat `{: url | null, errors: url | null}` map, so a client iterates it without branching on shape; every key is always present (null when the run produced no such file) — the column's `object_store_url` presigned for 1 day at delivery time, never stored, so a client whose receiver rejects the (large, item-inlining) body can still fetch the dumps. Both were hardcoded `None` before. Delivery is one inline attempt with no retry, so a rejected callback is still lost. Pre-filter `stop_on_fail` flag (`config/assessment_blob.py`) drives which filters hard-stop the chain on a failing verdict vs pass-through (record only). ## Services / CRUD - `services/assessment/utils/attachments.py` — cell→provider attachment conversion (Drive URL normalization; OpenAI/Anthropic/Gemini part builders). `rewrite_gcs_attachment_urls` bulk-resolves `gs://` cells to provider-reachable URLs via `services/buckets/` (Path A native passthrough for google-gcp / Path B signed HTTPS otherwise) **before** JSONL build. Called in: `services/assessment/api/batch.py::_submit_provider_batch` (API pipeline, all stages), `crud/assessment/batch.py::submit_assessment_batch` (legacy L2), `services/assessment/tasks.py` (legacy prefilter). Submit-time validation (`services/assessment/api/submission.py`) allows `gs://` alongside `http(s)://`. @@ -33,7 +33,7 @@ Delivery is poll-or-push: the `POST /assessments` ack is the flat `AssessmentSub - **Provider param mapping lives in `services/llm/mappers.py`, not in the assessment tree.** The assessment fork (`services/assessment/mappers.py`) is deleted; `crud/assessment/batch.py`, `services/assessment/api/batch.py` and `prefilter/request_builder.py` all import the shared mappers. The structured-output key is `output_schema` on every provider (the wire name stays `json_output_schema` on `AssessmentTextParams`; `api/batch.py::_stage_params` renames it). `normalize_llm_text` is exported from the same module but no mapper calls it — assessment normalises `instructions` at its own call sites, so prod LLM prompts are untouched. - **`evaluation_dataset` is no longer used anywhere in the assessment domain.** Assessment rows used to be `type='assessment'` entries in that shared table, which multiplexes four surfaces behind a discriminator, carries eval-only columns (`language_id`, `langfuse_dataset_id`), and whose name uniqueness ignores `type`, so an eval dataset name blocked an assessment one. `assessment.dataset_id` became `assessment.submission_id` (UUID FK) in migration 084. - `services/assessment/api/` — API-client pipeline: `submission.py` (submit), `submission_store.py` (submission rows in object storage), `batch.py` (staged provider batches — gate pre-filters → pass-through → assessment, over `core/batch`; `PREFILTER_VERDICT_SCHEMA`), `results.py` (builds `AssessmentBatchResult`), `result_files.py` (dump bookkeeping + `errors.jsonl` + presigned callback `metadata`), `callbacks.py` (webhook) -- `result_files.py` — `stage_file_kind` / `record_stage_dump` (called per completed stage from `batch.py`), `build_and_upload_errors` (one self-describing JSONL of `execution_error` / `row_error` / `provider_error_file` records, uploaded **even when empty** so both urls always exist; per-row errors come from the exec bag's `stage_errors`; OpenAI's error-file rows are merged into the stage's parsed results at poll time via `batch.error_file_entries`, so they land in `stage_errors` and `build_result` reports them per row instead of as silent "no output, no error" gaps; the same helper feeds the errors.jsonl dump), `finalize_result_files` (called from `_finalize`/`_fail` *before* the callback_url check, so durability never depends on a callback), `build_callback_metadata` (presigns each kind at `expires_in=86400`; never raises). Storage keys land under one per-assessment prefix, `/assessment//` (`assessment_subdirectory`), holding `errors.jsonl` plus `batch-/results.jsonl` per stage; the API path passes that prefix into `process_completed_batch`'s optional `subdirectory` (legacy RUN and evaluations keep the default `/batch-`). Uploaded submission files sit at `assessment/submissions//.` beside their parsed `submission.jsonl`. Only an **inline** BATCH gets a per-assessment copy at `/submission.jsonl` (`assessment.submission_input`); a run submitted by `submission_doc_id` streams the uploaded submission's own rows, because a submission is immutable (create/read/delete only, no update crud) and the copy would never diverge. +- `result_files.py` — `record_stage_dump` (called per completed stage from `batch.py`), `build_and_upload_errors` (one self-describing JSONL of `execution_error` / `row_error` / `provider_error_file` records, **not uploaded at all when there are no rows**, so a clean run leaves no empty object and the callback reports `errors: null`; per-row errors come from the exec bag's `stage_errors`; OpenAI's error-file rows are merged into the stage's parsed results at poll time via `batch.error_file_entries`, so they land in `stage_errors` and `build_result` reports them per row instead of as silent "no output, no error" gaps; the same helper feeds the errors.jsonl dump), `finalize_result_files` (called from `_finalize`/`_fail` *before* the callback_url check, so durability never depends on a callback), `presign_result_files` (presigns each key at `expires_in=86400` into `AssessmentResultFiles`; never raises). Storage keys land under one per-assessment prefix, `/assessment//` (`assessment_subdirectory`), holding `errors.jsonl` plus `batch-/results.jsonl` per stage; the API path passes that prefix into `process_completed_batch`'s optional `subdirectory` (legacy RUN and evaluations keep the default `/batch-`). Uploaded submission files sit at `assessment/submissions//.` beside their parsed `submission.jsonl`. Only an **inline** BATCH gets a per-assessment copy at `/submission.jsonl` (`assessment.submission_input`); a run submitted by `submission_doc_id` streams the uploaded submission's own rows, because a submission is immutable (create/read/delete only, no update crud) and the copy would never diverge. - `crud/assessment/submission.py` — `create_submission`, `get_submission_by_name`, `get_submission_by_id`, `list_submissions`, `delete_submission` (refuses while an assessment still references it) - `crud/assessment/api.py` — new API-client crud (method-based Assessment/AssessmentRun writes): `create_assessment` (takes a pre-generated `assessment_id` + `submission_input`, so submit uploads `submission.jsonl` first and then does one insert — no FAILED orphan when the upload fails), `set_assessment_job`, `create_execution`, `set_execution_batch_job`, `save_execution_state`, `set_result_files`, `update_status`, `list_executions`, `list_assessments_with_execution` (BATCH only; the detail route reads through the legacy `core.get_assessment_by_id`). Namespaced under `api` (`from app.crud.assessment import api`) to avoid colliding with the legacy `create_assessment`. - `crud/assessment/{core,cron,processing,batch}.py` — legacy RUN pipeline crud. RUN runtime for the dropped columns (`stage`/`stage_status`/`pipeline`/`stage_batches`/`prefilter_total_*`/`object_store_url`) now lives in the `assessment_run.execution` bag via `core._read_exec`/`_write_exec`; `status` is the `AssessmentStatus` enum; the RUN input binding lives on the parent `assessment.input`. `batch.py` builds per-row prompts (`{column}` substitution from the parent `InputBinding.prompt`).