From 28d0b1223baddf53c2907b23f7c6aa2934889092 Mon Sep 17 00:00:00 2001 From: Eomdahyeon <213566566+Eomdahyeon@users.noreply.github.com> Date: Sun, 4 Oct 2026 23:49:52 +0900 Subject: [PATCH 1/5] fix(builds): let the owner read a run a restart interrupted Closes #996 Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 1 + contract/builder-api.yaml | 17 +- contract/fixtures/responses.json | 2 +- src/kpubdata_builder/events/store.py | 65 ++++++- src/kpubdata_builder/service/app.py | 4 +- .../service/build_runs_api.py | 54 +++++- .../service/routes/_guards.py | 14 +- tests/unit/test_interrupted_run_status.py | 176 ++++++++++++++++++ 8 files changed, 323 insertions(+), 10 deletions(-) create mode 100644 tests/unit/test_interrupted_run_status.py diff --git a/CHANGELOG.md b/CHANGELOG.md index fe61faec..e11e9f00 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,7 @@ ### Fixed +- The owner of a run a restart interrupted can read why it ended (#996, API contract 1.78.0, additive). Such a run is failed as `credentials_required` in the event store only (#683); with the job registry empty after the restart and no manifest, `GET /builds/{run_id}` and `GET /builds/{run_id}/events` both answered 404, so a client polling its build was told the run did not exist instead of to submit it again. The submitter of an async run is now recorded next to its events, the status route reads the terminal event back as `failed` with `error` and the new optional `BuildJob.code` `credentials_required`, and the events route admits the owner. Another user gets what a missing run gets, as before. A run submitted before this version has no recorded submitter and still answers 404. - A cross-origin page can read `Content-Disposition`, `X-Request-ID` and `Retry-After`, and the overload 503 (#995). Responses to an allowed origin carried no `Access-Control-Expose-Headers`, so a browser hid every header outside the CORS-safelisted set: Studio could not read a download's file name, the request id, or how long to wait. They are now exposed. The overload 503 is written to the socket before the request is read, so it had no CORS header at all and a browser reported it as a network error; it now carries the deployment's one allowed origin, or `*` when several are allowed (readable by a request made without credentials — a bearer header is not one), with `Retry-After` exposed, and still nothing when no origin is allowed. - A run whose table commit failed has one job status (#997). Such a build answers 409 with `warehouse_failures`, so `GET /builds/{run_id}` said `failed` while the job registry held it; once the registry evicted the job or the service restarted, the status was derived from the manifest, whose `status` is `ok`, and the same run read `succeeded`. The manifest path now reads a recorded `warehouse_failures` as `failed` too. The manifest is unchanged, and the contract's `BuildJob.status` says which status this case has. - An async build reads its submitter's uploads (#998). `POST /builds` accepted a spec with a `kind: file` source and the job then failed with `file source requires an authenticated, stable principal owner`: the worker passed the submitting owner to the manifest and to credential resolution but not to the file resolver, so only the synchronous `POST /build` could build an uploaded file. The worker now passes the same owner there too. Uploads stay isolated per owner — a job naming another owner's `upload_id` fails as it does synchronously — and a submission with no stable owner still fails with that message. The contract's `submitBuild` description says so; the wire is unchanged. diff --git a/contract/builder-api.yaml b/contract/builder-api.yaml index 9eafa5c4..e86181b0 100644 --- a/contract/builder-api.yaml +++ b/contract/builder-api.yaml @@ -1,7 +1,7 @@ openapi: 3.1.0 info: title: KPubData Builder Service API - version: "1.77.0" + version: "1.78.0" description: >- KPubData Builder(패키지 `kpubdata-builder`) 서비스 계약. BuildSpec 검증, preview, build 실행, artifact 조회, 계약 버전 조회를 제공한다. CLI와 HTTP service mode(#36)가 동일한 도메인 계약을 @@ -20,6 +20,7 @@ info: 추가한다(#504, additive). v1.8.0은 stable principal별 encrypted Provider credential CRUD와 Provider status/test API를 추가한다(#492, additive). v1.9.0은 `/query` 응답에 child startup과 Polars engine 실행 시간을 추가한다(#523, additive). + v1.78.0 lets the owner read a run a restart interrupted (#996, additive): `GET /builds/{run_id}` answers such a run as `failed` with the new optional `BuildJob.code` `credentials_required` and the reason in `error`, and `GET /builds/{run_id}/events` shows its `run_failed` event, where both answered 404 before; another user still gets what a missing run gets. The submitter of an async run is now recorded with its events, so this holds for runs submitted from this version on. v1.77.0 stops a snapshot profile's range from disclosing a single record's value (#903, additive): `ColumnRange.min`/`max` are read after removing the `range_trim` lowest and the `range_trim` highest values, the range's `status` is `trimmed` instead of `exact`, and it carries `trimmed_count`. `SnapshotProfile` gains `range_trim`, and `min_range_values` is at least `2 * range_trim + 1`. A client that showed only `exact` ranges shows none until it reads `trimmed`. v1.76.0 records which SQL a saved analysis is written in (#875, additive): `SavedAnalysis` gains `sql_dialect` (`duckdb`, or `legacy-polars` for one saved before the DuckDB cutover), `engine`, `engine_version`, `query_contract_version` and `migration_required`. A legacy analysis is never re-run on DuckDB: `POST /analyses/{analysis_id}/run` answers 409 `analysis_migration_required` until its SQL is reviewed and saved as a new analysis. Existing stores gain the columns on start; their rows are the legacy ones. v1.74.0 runs SQL, rows, aggregates, profiles and exports on DuckDB (#874, ADR 0021 step 13): each in a child process over a DuckDB connection that may read only the pinned snapshot file and spill into its own directory, with no other file, network or extension access and its configuration locked. The validator now parses the DuckDB dialect and refuses introspective functions (settings, variables, versions, sequences …) and nondeterministic SQL (random values, UUIDs, the current time, `USING SAMPLE`/`TABLESAMPLE`) with `unsafe_query`. Results keep the types DuckDB gives them, in Builder's spelling, where they differ from the Polars engine: `COUNT(*)` is `int64` (was `uint32`) and an unnamed aggregate is named as DuckDB names it (`count_star()`, `sum(v)` — they were `len` or failed); a `SUM` of an integer column is `int128`, sent as a number or as exact decimal text by its values; `ORDER BY … DESC` puts nulls last (Polars put them first; rows and aggregates always did); `date + INTERVAL` now works and is a datetime; a column that is Null in Builder is `int32` (all null) in SQL results, while a page of rows still reports it Null; a zoned datetime is sent in UTC. A table with a column SQL cannot name (an empty name, or two names one letter case apart) answers `query_execution_failed` to SQL, while rows, aggregates and profiles still read it. @@ -7501,8 +7502,10 @@ components: type: object required: [run_id, status, created_at, updated_at] description: >- - 비동기 build job 스냅샷 (#482, #480). 프로세스 메모리 상주 registry - 기반이라 서버 재시작 후 active/terminal 잡은 더 이상 조회되지 않는다. + 비동기 build job 스냅샷 (#482, #480). 프로세스 메모리 상주 registry 가 + 먼저 답하고, 거기서 밀려나거나 서버가 재시작한 뒤에는 매니페스트로, + 매니페스트도 없이 재시작으로 중단된 run 은 이벤트 저장소의 기록으로 + 답한다(`failed`, `code: credentials_required`, #996). `response`는 잡이 만든 응답 본문이다 — 성공한 잡은 최종 build 응답, 실패한 잡은 오류 본문일 수 있다(빌드 본문이 생기기 전에 실패하면 `{"error": ...}` 뿐이다, #921). 형태는 `object | null` 만 보장하며, 실패 @@ -7543,6 +7546,14 @@ components: error: type: [string, "null"] description: 실패 사유 (실패 시에만) + code: + type: string + enum: [credentials_required] + description: >- + A stable reason for a failure the service itself caused (#996), when there + is one. `credentials_required`: the server restarted while the job was + queued or running; its provider keys were held only in memory (#683), so + it cannot go on — submit it again. BuildsResponse: type: object diff --git a/contract/fixtures/responses.json b/contract/fixtures/responses.json index 82428350..a15cec92 100644 --- a/contract/fixtures/responses.json +++ b/contract/fixtures/responses.json @@ -1,6 +1,6 @@ { "fixture_format": 1, - "contract_version": "1.77.0", + "contract_version": "1.78.0", "probe_field": "future_optional_field", "rules": "A response may gain optional fields in any minor version; a client ignores fields it does not know. A required field keeps its name and type until the next major version; a client rejects a body whose required field is missing or mistyped.", "fixtures": [ diff --git a/src/kpubdata_builder/events/store.py b/src/kpubdata_builder/events/store.py index ad67a729..dbf9ef53 100644 --- a/src/kpubdata_builder/events/store.py +++ b/src/kpubdata_builder/events/store.py @@ -28,7 +28,7 @@ import threading from collections.abc import Iterator from contextlib import contextmanager -from dataclasses import replace +from dataclasses import dataclass, replace from datetime import datetime, timezone from pathlib import Path from typing import cast @@ -56,6 +56,29 @@ _SELECT_COLUMNS = "seq, run_id, timestamp, event, status, source_key, stage, message, metrics" +# Who submitted an async run, and when (#996). The event timeline says what happened to +# a run but not whose it is; ownership lived in the job registry (memory) until the run +# wrote a manifest. A run a restart interrupts has neither, so its owner could not read +# why it ended. Kept here, next to the events it belongs with; rows are only added. +_CREATE_SUBMISSIONS_SQL = """ +CREATE TABLE IF NOT EXISTS run_submissions ( + run_id TEXT PRIMARY KEY, + owner_id TEXT, + created_by TEXT, + submitted_at TEXT NOT NULL +) +""" + + +@dataclass(frozen=True, slots=True) +class RunSubmission: + """Who submitted an async run and when (#996).""" + + run_id: str + owner_id: str | None + created_by: str | None + submitted_at: str + class BuildEventStore: """Append-only event store for all runs under single output_root. @@ -104,6 +127,7 @@ def _init_db(self) -> None: f"INSERT INTO schema_version (version) VALUES ({SCHEMA_VERSION})" ) self._conn.execute(_CREATE_TABLE_SQL) + self._conn.execute(_CREATE_SUBMISSIONS_SQL) self._conn.execute( "CREATE INDEX IF NOT EXISTS idx_build_events_run_seq ON build_events(run_id, seq)" ) @@ -154,6 +178,43 @@ def append(self, event: BuildEvent) -> BuildEvent: raise RuntimeError("build event insert did not return a row id") return replace(event, seq=seq) + def record_submission( + self, run_id: str, *, owner_id: str | None, created_by: str | None, submitted_at: datetime + ) -> None: + """Remember who submitted ``run_id`` (#996). A second submission of the id is ignored.""" + if submitted_at.tzinfo is None or submitted_at.utcoffset() is None: + raise ValueError("submitted_at must be timezone-aware") + with self._transaction(): + self._conn.execute( + "INSERT OR IGNORE INTO run_submissions (run_id, owner_id, created_by, submitted_at)" + " VALUES (?, ?, ?, ?)", + (run_id, owner_id, created_by, submitted_at.astimezone(timezone.utc).isoformat()), + ) + + def submission(self, run_id: str) -> RunSubmission | None: + """The recorded submitter of ``run_id``, or None when it was never recorded.""" + row = self._conn.execute( + "SELECT owner_id, created_by, submitted_at FROM run_submissions WHERE run_id = ?", + (run_id,), + ).fetchone() + if row is None: + return None + return RunSubmission( + run_id=run_id, + owner_id=None if row[0] is None else str(row[0]), + created_by=None if row[1] is None else str(row[1]), + submitted_at=str(row[2]), + ) + + def terminal_event(self, run_id: str) -> BuildEvent | None: + """The run's last terminal event (finished, failed or cancelled), if it has one.""" + events = [ + event + for event in self.list_for_run(run_id, limit=50, tail=True) + if event.event in ("run_finished", "run_failed", "run_cancelled") + ] + return events[-1] if events else None + def unfinished_runs(self) -> tuple[str, ...]: """Runs with a submission or start event and no terminal event (#683).""" rows = self._conn.execute( @@ -224,4 +285,4 @@ def _row_to_event(row: tuple[object, ...]) -> BuildEvent: ) -__all__ = ["BuildEventStore", "SCHEMA_VERSION"] +__all__ = ["SCHEMA_VERSION", "BuildEventStore", "RunSubmission"] diff --git a/src/kpubdata_builder/service/app.py b/src/kpubdata_builder/service/app.py index 0250c26e..62a821ac 100644 --- a/src/kpubdata_builder/service/app.py +++ b/src/kpubdata_builder/service/app.py @@ -397,7 +397,9 @@ def _enforce_ownership() -> bool: # 1.76.0 -> 1.77.0: a snapshot profile's range is trimmed (ColumnRange.status # `trimmed`, trimmed_count; SnapshotProfile.range_trim) so that one record's extreme # is not disclosed (#903, additive). -API_CONTRACT_VERSION = "1.77.0" +# 1.77.0 -> 1.78.0: a run a restart interrupted is readable by its owner — BuildJob +# gains the optional `code` (`credentials_required`) (#996, additive). +API_CONTRACT_VERSION = "1.78.0" #: manifest status vocabulary (ok/failed/cancelled) → publish status vocabulary diff --git a/src/kpubdata_builder/service/build_runs_api.py b/src/kpubdata_builder/service/build_runs_api.py index 78d57836..b89a8205 100644 --- a/src/kpubdata_builder/service/build_runs_api.py +++ b/src/kpubdata_builder/service/build_runs_api.py @@ -84,6 +84,11 @@ def _declared_dataset_id(spec_yaml: str) -> str | None: return dataset_id if isinstance(dataset_id, str) and dataset_id else None +#: The stable code of a run a restart interrupted (#683, #996). It starts the failure +#: event's message and is the ``code`` of the job status read back from that event. +INTERRUPTED_CODE = "credentials_required" + + class BuildRunsApiService: """Runs a build, queues one, reports on it and cancels it.""" @@ -383,16 +388,26 @@ def _record_run_submitted() -> None: # fails to record, job also never created: "event lost but job running" # contradiction never happens (recorder absorption differs; here no real # side effect yet to compromise other canonical). + submitted_at = datetime.now(tz=timezone.utc) self._event_store().append( BuildEvent( seq=0, - timestamp=datetime.now(tz=timezone.utc), + timestamp=submitted_at, run_id=resolved_run_id, event="run_submitted", status="ok", message="build accepted for async execution", ) ) + # Whose run this is, kept where it survives a restart (#996): the registry + # that holds the owner is memory, and a run that never writes a manifest + # has nothing else to say who may read why it ended. + self._event_store().record_submission( + resolved_run_id, + owner_id=owner_id, + created_by=created_by, + submitted_at=submitted_at, + ) if job_credentials is not None: job_credentials.bind(resolved_run_id, owner_id, request_credentials.current_keys()) @@ -489,8 +504,43 @@ def build_status(self, run_id: str) -> ServiceResponse: evicted = self._build_status_from_manifest(run_id) if evicted is not None: return ServiceResponse(200, evicted) + interrupted = self._build_status_from_events(run_id) + if interrupted is not None: + return ServiceResponse(200, interrupted) return ServiceResponse(404, {"error": f"build job not found: {run_id}"}) + def _build_status_from_events(self, run_id: str) -> dict[str, JsonValue] | None: + """Status of a run that ended without a manifest and is no longer in the registry. + + A restart leaves such a run: ``mark_interrupted_runs`` records its failure in + the event store only (#683), so neither the registry nor a manifest knows it + and the owner polling it got 404 instead of "submit it again" (#996). The + submission record gives the run's start, the terminal event its end and reason. + Only a failed or cancelled ending is reported from here — a run that finished + has a manifest, and that path is the authority for it. + """ + store = self._event_store() + submission = store.submission(run_id) + if submission is None: + return None + terminal = store.terminal_event(run_id) + if terminal is None or terminal.event == "run_finished": + return None + body: dict[str, JsonValue] = { + "run_id": run_id, + "status": "cancelled" if terminal.event == "run_cancelled" else "failed", + "created_at": submission.submitted_at, + "updated_at": terminal.timestamp.astimezone(timezone.utc).isoformat(), + } + if submission.created_by is not None: + body["created_by"] = submission.created_by + if terminal.event == "run_failed": + message = terminal.message or "the run did not finish" + body["error"] = message + if message.startswith(f"{INTERRUPTED_CODE}:"): + body["code"] = INTERRUPTED_CODE + return body + def _build_status_from_manifest(self, run_id: str) -> dict[str, JsonValue] | None: """Restore evicted terminal job status from persisted manifest. @@ -583,7 +633,7 @@ def mark_interrupted_runs(self) -> tuple[str, ...]: run_id=run_id, event="run_failed", status="fail", - message="credentials_required: the server restarted and the job's " + message=f"{INTERRUPTED_CODE}: the server restarted and the job's " "provider keys, held only in memory, are gone; submit it again", ) ) diff --git a/src/kpubdata_builder/service/routes/_guards.py b/src/kpubdata_builder/service/routes/_guards.py index bc97d92b..d1277a4b 100644 --- a/src/kpubdata_builder/service/routes/_guards.py +++ b/src/kpubdata_builder/service/routes/_guards.py @@ -83,7 +83,9 @@ def check_active_run_access( enqueue failure), decide ownership by that snapshot's stable ``owner_id`` (#505 canonical identity — ``created_by``/``Principal.label`` are legacy fallbacks only; we pass them as-is to ``ownership_allows`` for priority). - 3. If neither exist, 404. + 3. If neither exist, the recorded submission (event store) decides, by the + same ownership rule: this is a run a restart interrupted (#996). + 4. If none exist, 404. Snapshot ``owner_id`` is the value ``BuilderService.submit_build`` preserved in registry — not exposed in wire response (``BuildJobSnapshot.to_body()`` never exports @@ -106,6 +108,16 @@ def check_active_run_access( ): return None return not_owner(run_id) + # A run a restart interrupted has no manifest and is gone from the registry; who + # submitted it is in the event store (#996). Someone else still gets what a run + # that does not exist gets. + submission = service._event_store.submission(run_id) + if submission is not None: + if ownership_module.ownership_allows( + created_by=submission.created_by, owner_id=submission.owner_id, principal=principal + ): + return None + return not_owner(run_id) return ServiceResponse(404, {"error": f"run not found: {run_id}"}) diff --git a/tests/unit/test_interrupted_run_status.py b/tests/unit/test_interrupted_run_status.py new file mode 100644 index 00000000..be7491e2 --- /dev/null +++ b/tests/unit/test_interrupted_run_status.py @@ -0,0 +1,176 @@ +"""A run a restart interrupted is readable by its owner, and only by its owner (#996). + +``mark_interrupted_runs`` fails such a run in the event store alone. After the restart +the job registry is empty and no manifest exists, so the status and events routes +answered 404: a client polling its build was told the run did not exist instead of to +submit it again. +""" + +from __future__ import annotations + +import threading +from collections.abc import Callable, Iterator +from pathlib import Path +from typing import cast + +import pytest + +from kpubdata_builder.service import BuilderService, ServiceResponse, dispatch +from kpubdata_builder.service import app as app_module +from kpubdata_builder.service.auth import Principal + +_ALICE = Principal("oidc", "alice", "oidc:alice") +_BOB = Principal("oidc", "bob", "oidc:bob") +_SPEC = """\ +dataset_id: interrupted.run +title: Interrupted +description: d +sources: + - provider: datago + dataset: air_quality +exports: + - kind: jsonl + output_path: data.jsonl +""" + + +class _Result: + items = [{"id": "1"}] + + +def _factory(gate: threading.Event | None = None) -> Callable[..., object]: + """A client factory with the keywords the service asks a factory for (#683, #786). + + The fetch waits on ``gate``, so the first process's job stays in flight. + """ + + class _Dataset: + def list(self, **_params: object) -> _Result: + if gate is not None: + gate.wait(timeout=10) + return _Result() + + class _Client: + def dataset(self, _key: str) -> _Dataset: + return _Dataset() + + def create( + *, + provider_keys: dict[str, str] | None = None, + timeout: float | None = None, + cache: bool | None = None, + environment_keys: bool = True, + ) -> object: + return _Client() + + return create + + +def _as(monkeypatch: pytest.MonkeyPatch, principal: Principal) -> None: + monkeypatch.setattr(app_module, "authenticate", lambda **_: principal) + + +def _get(service: BuilderService, path: str) -> ServiceResponse: + response = dispatch(service, "GET", path, None) + assert isinstance(response, ServiceResponse) + return response + + +@pytest.fixture() +def restarted(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> Iterator[BuilderService]: + """A service started after another one was stopped with Alice's job in flight. + + The first service's worker stays blocked until the test is over: a real restart + kills it, and letting it run on would add events and a manifest behind the test. + """ + monkeypatch.setenv("ENFORCE_OWNERSHIP", "true") + gate = threading.Event() + first = BuilderService(output_root=tmp_path, client_factory=_factory(gate)) + _as(monkeypatch, _ALICE) + submitted = dispatch( + first, + "POST", + "/builds", + {"spec": _SPEC, "run_id": "in-flight"}, + provider_key_headers=["datago=alice-key"], + ) + assert isinstance(submitted, ServiceResponse) + assert submitted.status_code == 202 + + second = BuilderService(output_root=tmp_path, client_factory=_factory()) + assert "in-flight" in second.mark_interrupted_runs() + try: + yield second + finally: + gate.set() + first._async_builds.shutdown() + second._async_builds.shutdown() + + +def test_the_owner_reads_failed_with_a_stable_code( + restarted: BuilderService, monkeypatch: pytest.MonkeyPatch +) -> None: + _as(monkeypatch, _ALICE) + + response = _get(restarted, "/builds/in-flight") + + assert response.status_code == 200 + assert response.body["status"] == "failed" + assert response.body["code"] == "credentials_required" + assert str(response.body["error"]).startswith("credentials_required: the server restarted") + assert response.body["run_id"] == "in-flight" + assert response.body["created_at"] <= response.body["updated_at"] + + +def test_the_owner_sees_the_failure_event( + restarted: BuilderService, monkeypatch: pytest.MonkeyPatch +) -> None: + _as(monkeypatch, _ALICE) + + response = _get(restarted, "/builds/in-flight/events") + + assert response.status_code == 200 + events = cast(list[dict[str, object]], response.body["events"]) + names = [event["event"] for event in events] + assert names[0] == "run_submitted" + assert "run_failed" in names + + +@pytest.mark.parametrize("path", ["/builds/in-flight", "/builds/in-flight/events"]) +def test_another_user_gets_what_a_missing_run_gets( + restarted: BuilderService, monkeypatch: pytest.MonkeyPatch, path: str +) -> None: + _as(monkeypatch, _BOB) + + theirs = _get(restarted, path) + missing = _get(restarted, path.replace("in-flight", "never-existed")) + + assert theirs.status_code == missing.status_code == 404 + assert "credentials_required" not in str(theirs.body) + + +def test_a_run_nobody_submitted_is_still_404( + restarted: BuilderService, monkeypatch: pytest.MonkeyPatch +) -> None: + _as(monkeypatch, _ALICE) + + assert _get(restarted, "/builds/never-existed").status_code == 404 + + +def test_the_submission_record_is_written_once_and_read_back(tmp_path: Path) -> None: + import datetime as dt + + from kpubdata_builder.events import BuildEventStore + + store = BuildEventStore(tmp_path) + at = dt.datetime(2026, 10, 4, 12, 0, tzinfo=dt.timezone.utc) + + store.record_submission("r1", owner_id="oidc:alice", created_by="alice", submitted_at=at) + store.record_submission("r1", owner_id="oidc:mallory", created_by="mallory", submitted_at=at) + + found = store.submission("r1") + assert found is not None + assert (found.owner_id, found.created_by) == ("oidc:alice", "alice") + assert found.submitted_at == "2026-10-04T12:00:00+00:00" + assert store.submission("r2") is None + assert store.terminal_event("r1") is None From 090adc5e480fa2536e29b69c7f87fb4f6d9627d2 Mon Sep 17 00:00:00 2001 From: Eomdahyeon <213566566+Eomdahyeon@users.noreply.github.com> Date: Sun, 4 Oct 2026 23:58:22 +0900 Subject: [PATCH 2/5] feat(api): give the full build queue, auth failures and the overload 503 a stable code Closes #1000 Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 1 + contract/builder-api.yaml | 26 +++++++++++--- contract/fixtures/responses.json | 35 +++++++++++++++--- src/kpubdata_builder/service/app.py | 9 +++-- src/kpubdata_builder/service/auth.py | 15 ++++++++ .../service/build_runs_api.py | 6 +++- src/kpubdata_builder/service/http.py | 2 +- tests/unit/test_oidc_auth.py | 35 ++++++++++++++++++ tests/unit/test_service.py | 4 +-- tests/unit/test_stable_error_codes.py | 36 +++++++++++++++++++ 10 files changed, 155 insertions(+), 14 deletions(-) create mode 100644 tests/unit/test_stable_error_codes.py diff --git a/CHANGELOG.md b/CHANGELOG.md index e11e9f00..837fb378 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -41,6 +41,7 @@ ### Added +- Four failures a client had to tell apart by their sentence have a stable `code` (#1000, API contract 1.79.0, additive). The full async build queue answers 429 `build_queue_full` — it shared its status with `auth_throttled` and had no code. An authentication failure answers `unauthorized`, or `token_expired` when the bearer token's `exp` has passed, or 503 `auth_unavailable` when the JWKS could not be fetched. The overload 503 written to the socket before a request is read answers `server_overloaded`. Every body keeps its `error` text. A run a restart interrupted got its code, `credentials_required`, with #996. - `card.json` is declared in the contract as `DatasetCard` and states the licence mismatch and "no processing" as fields (#955, API contract 1.75.0, additive). The card sections every Gold output has carried since 1.71.0 (#694) had no schema, and two facts lived only in sentences a client had to match: `…; the provider declares: …` and `No transformation declared: values are as the source gave them.` Each `provenance` entry now also has `license_declared` (the BuildSpec's licence as written, with its link; null when none), `license_provider` (the kpubdata catalog's terms; null when none) and `license_mismatch` (both known and the provider's terms are not in the declared licence — the condition the sentence already used), and the card has `processing_declared` (false when only the no-transformation sentence is listed; a composed output's join counts as processing). The sentences are unchanged for the README. `GET /artifacts/{run_id}/{file_path}` says `gold//card.json` is a `DatasetCard` and gains an `application/json` example of one, so `contract/fixtures/responses.json` carries card fixtures. Tests validate real `card_sections` output, from the functions and from single and composed builds, against the schema, and fail when the writer sends a field the schema does not name. Cards written before 1.75.0 lack the four fields. - `GET /admin/runs` gives `total`, how many runs there are before `limit` (#948, API contract 1.73.0, additive). `count` was always the runs in the response — `len(runs)` after the `limit` cut — and is now documented as that, not a total; `total` is the build index's runs plus the queued and running jobs it does not hold yet, each counted once, with one `COUNT(*)` and a key lookup rather than reading every row (`BuildIndex.count_builds`, SQLite and CUBRID). `count < total` means older runs were left out, so an admin screen can say "N of M". A client talking to an older Builder finds `total` absent. - Saved analyses record the SQL they are written in (#875, step 14 of ADR 0021, API contract 1.76.0): `sql_dialect`, `engine`, `engine_version`, `query_contract_version` and `migration_required`. An existing analysis store gains the columns when Builder starts, and its analyses — saved for the Polars engine — are `legacy-polars`. Those are never re-run on DuckDB as if nothing had changed: `POST /analyses/{analysis_id}/run` answers 409 `analysis_migration_required` (Studio: kpubdata-lab/kpubdata-studio#565) until the SQL is reviewed and saved as a new analysis; reading and deleting them (which releases the snapshot hold) work as before. diff --git a/contract/builder-api.yaml b/contract/builder-api.yaml index e86181b0..79d6bda4 100644 --- a/contract/builder-api.yaml +++ b/contract/builder-api.yaml @@ -1,7 +1,7 @@ openapi: 3.1.0 info: title: KPubData Builder Service API - version: "1.78.0" + version: "1.79.0" description: >- KPubData Builder(패키지 `kpubdata-builder`) 서비스 계약. BuildSpec 검증, preview, build 실행, artifact 조회, 계약 버전 조회를 제공한다. CLI와 HTTP service mode(#36)가 동일한 도메인 계약을 @@ -20,6 +20,7 @@ info: 추가한다(#504, additive). v1.8.0은 stable principal별 encrypted Provider credential CRUD와 Provider status/test API를 추가한다(#492, additive). v1.9.0은 `/query` 응답에 child startup과 Polars engine 실행 시간을 추가한다(#523, additive). + v1.79.0 gives four failures a stable `code` (#1000, additive): the full async build queue's 429 is `build_queue_full` (it was told from 429 `auth_throttled` only by its sentence), an authentication failure is `unauthorized`, `token_expired` (the bearer token's `exp` has passed — get a new token) or, with 503, `auth_unavailable` (the JWKS could not be fetched), and the overload 503 written before a request is read is `server_overloaded`. Each body keeps its `error` text. v1.78.0 lets the owner read a run a restart interrupted (#996, additive): `GET /builds/{run_id}` answers such a run as `failed` with the new optional `BuildJob.code` `credentials_required` and the reason in `error`, and `GET /builds/{run_id}/events` shows its `run_failed` event, where both answered 404 before; another user still gets what a missing run gets. The submitter of an async run is now recorded with its events, so this holds for runs submitted from this version on. v1.77.0 stops a snapshot profile's range from disclosing a single record's value (#903, additive): `ColumnRange.min`/`max` are read after removing the `range_trim` lowest and the `range_trim` highest values, the range's `status` is `trimmed` instead of `exact`, and it carries `trimmed_count`. `SnapshotProfile` gains `range_trim`, and `min_range_values` is at least `2 * range_trim + 1`. A client that showed only `exact` ranges shows none until it reads `trimmed`. v1.76.0 records which SQL a saved analysis is written in (#875, additive): `SavedAnalysis` gains `sql_dialect` (`duckdb`, or `legacy-polars` for one saved before the DuckDB cutover), `engine`, `engine_version`, `query_contract_version` and `migration_required`. A legacy analysis is never re-run on DuckDB: `POST /analyses/{analysis_id}/run` answers 409 `analysis_migration_required` until its SQL is reviewed and saved as a new analysis. Existing stores gain the columns on start; their rows are the legacy ones. @@ -1454,11 +1455,17 @@ paths: schema: $ref: "#/components/schemas/Error" "429": - description: async build queue가 가득 참 + description: async build queue가 가득 참 (`build_queue_full`, #1000) content: application/json: schema: $ref: "#/components/schemas/Error" + examples: + BuildQueueFull: + summary: The async build queue is full + value: + error: async build queue is full + code: build_queue_full "500": description: job 접수 실패(event 기록/큐잉 불가 — job은 생성되지 않음) content: @@ -4184,7 +4191,10 @@ components: code: signup_rejected Unauthorized: - description: X-API-Key가 누락되었거나 서버에 설정된 키와 일치하지 않음 + description: >- + X-API-Key가 누락되었거나 서버에 설정된 키와 일치하지 않음, 또는 Bearer 토큰이 + 검증되지 않음. `code` 는 `unauthorized`, 토큰의 `exp` 가 지났으면 + `token_expired` 다(#1000). content: application/json: schema: @@ -4193,7 +4203,8 @@ components: MissingOrInvalidApiKey: summary: 인증 정보가 없거나 일치하지 않음 value: - error: unauthorized + error: invalid api key + code: unauthorized PiiDeclarationUnavailable: description: >- @@ -4408,6 +4419,13 @@ components: (`SignupNotApprovedError`, `components.responses.SignupNotApproved`). `revision_conflict`(#820, #947, 409) — a revision save or revert based on a stale revision; the body carries `current_revision` (`RevisionConflictError`). + `build_queue_full`(#1000, 429) — `POST /builds` when the async build queue + is full; try again later. `unauthorized` / `token_expired`(#1000, 401) — any + operation: the credentials were refused, or the bearer token's `exp` has + passed and a new token is needed. `auth_unavailable`(#1000, 503) — the + JWKS could not be fetched, so the credentials were not judged. + `server_overloaded`(#1000, 503) — too many requests in flight; the response + carries `Retry-After`. redistribution: description: >- `code: redistribution_forbidden`일 때만 온다(#688): 판정과 source별 이유. diff --git a/contract/fixtures/responses.json b/contract/fixtures/responses.json index a15cec92..2ff29d91 100644 --- a/contract/fixtures/responses.json +++ b/contract/fixtures/responses.json @@ -1,6 +1,6 @@ { "fixture_format": 1, - "contract_version": "1.78.0", + "contract_version": "1.79.0", "probe_field": "future_optional_field", "rules": "A response may gain optional fields in any minor version; a client ignores fields it does not know. A required field keeps its name and type until the next major version; a client rejects a body whose required field is missing or mistyped.", "fixtures": [ @@ -5097,6 +5097,30 @@ }, "broken_path": "$.error" }, + { + "operation_id": "submitBuild", + "method": "POST", + "path": "/builds", + "status": 429, + "example": "BuildQueueFull", + "current": { + "error": "async build queue is full", + "code": "build_queue_full" + }, + "with_additive_fields": { + "error": "async build queue is full", + "code": "build_queue_full", + "future_optional_field": "added by a later minor contract version" + }, + "additive_paths": [ + "$" + ], + "required_type_broken": { + "error": 12345, + "code": "build_queue_full" + }, + "broken_path": "$.error" + }, { "operation_id": "listBuilds", "method": "GET", @@ -5694,17 +5718,20 @@ "status": 401, "example": "MissingOrInvalidApiKey", "current": { - "error": "unauthorized" + "error": "invalid api key", + "code": "unauthorized" }, "with_additive_fields": { - "error": "unauthorized", + "error": "invalid api key", + "code": "unauthorized", "future_optional_field": "added by a later minor contract version" }, "additive_paths": [ "$" ], "required_type_broken": { - "error": 12345 + "error": 12345, + "code": "unauthorized" }, "broken_path": "$.error" }, diff --git a/src/kpubdata_builder/service/app.py b/src/kpubdata_builder/service/app.py index 62a821ac..c0b21356 100644 --- a/src/kpubdata_builder/service/app.py +++ b/src/kpubdata_builder/service/app.py @@ -399,7 +399,10 @@ def _enforce_ownership() -> bool: # is not disclosed (#903, additive). # 1.77.0 -> 1.78.0: a run a restart interrupted is readable by its owner — BuildJob # gains the optional `code` (`credentials_required`) (#996, additive). -API_CONTRACT_VERSION = "1.78.0" +# 1.78.0 -> 1.79.0: stable codes for the full build queue (429 build_queue_full), +# authentication failures (unauthorized, token_expired, auth_unavailable) and the +# overload 503 (server_overloaded) (#1000, additive). +API_CONTRACT_VERSION = "1.79.0" #: manifest status vocabulary (ok/failed/cancelled) → publish status vocabulary @@ -1455,7 +1458,9 @@ def _dispatch_impl( # fault). if principal.status_code == 401: service._auth_throttle.record_failure(client_id) - return ServiceResponse(principal.status_code, {"error": principal.reason}) + return ServiceResponse( + principal.status_code, {"error": principal.reason, "code": principal.code} + ) # Successful authentication clears failure record — normal client that received # a few 401s due to token expiry doesn't get throttled during subsequent normal diff --git a/src/kpubdata_builder/service/auth.py b/src/kpubdata_builder/service/auth.py index 03bd8f00..14892346 100644 --- a/src/kpubdata_builder/service/auth.py +++ b/src/kpubdata_builder/service/auth.py @@ -168,6 +168,21 @@ class AuthError: reason: str status_code: int = 401 + @property + def code(self) -> str: + """A stable code for the failure, so a client does not branch on ``reason`` (#1000). + + ``token_expired`` — the bearer token's ``exp`` has passed: get a new token and + send the request again. ``auth_unavailable`` — the JWKS could not be fetched + (503): the credentials were not judged, try again. ``unauthorized`` — every + other refusal: a missing or wrong API key, a token that does not verify. + """ + if self.status_code == 503: + return "auth_unavailable" + if self.reason.endswith("ExpiredSignatureError"): + return "token_expired" + return "unauthorized" + def _is_dev_mode() -> bool: """Check if in local development mode (#321, ADR 0006). diff --git a/src/kpubdata_builder/service/build_runs_api.py b/src/kpubdata_builder/service/build_runs_api.py index b89a8205..47189b0a 100644 --- a/src/kpubdata_builder/service/build_runs_api.py +++ b/src/kpubdata_builder/service/build_runs_api.py @@ -478,7 +478,11 @@ def _record_enqueue_failure() -> None: raise RuntimeError("existing async build is missing snapshot") return ServiceResponse(200, result.snapshot.to_body()) case "queue_full": - return ServiceResponse(429, {"error": "async build queue is full"}) + # A code of its own (#1000): the other 429 on this route's way in is + # `auth_throttled`, and the sentence was the only way to tell them apart. + return ServiceResponse( + 429, {"error": "async build queue is full", "code": "build_queue_full"} + ) case unreachable: assert_never(unreachable) diff --git a/src/kpubdata_builder/service/http.py b/src/kpubdata_builder/service/http.py index ec08273a..fe28168a 100644 --- a/src/kpubdata_builder/service/http.py +++ b/src/kpubdata_builder/service/http.py @@ -79,7 +79,7 @@ def _overloaded_response(allowed_origins: frozenset[str] = frozenset()) -> bytes (a bearer header is not one). The body is a constant and says nothing a page from another origin could not learn by being refused. """ - body = b'{"error": "server overloaded"}' + body = b'{"error": "server overloaded", "code": "server_overloaded"}' cors = b"" # This is written to the socket as bytes, so a configured origin with a line break # or a non-ASCII character would corrupt the response: send no CORS header then. diff --git a/tests/unit/test_oidc_auth.py b/tests/unit/test_oidc_auth.py index 3726ea96..c9250fe7 100644 --- a/tests/unit/test_oidc_auth.py +++ b/tests/unit/test_oidc_auth.py @@ -8,6 +8,7 @@ import logging import time +from pathlib import Path import jwt import pytest @@ -271,6 +272,8 @@ def test_expired_token(self, oidc_env: bytes) -> None: token = _make_token(oidc_env, exp=now - 120) result = authenticate(bearer_token=f"Bearer {token}") assert isinstance(result, AuthError) + # A code of its own (#1000): the client gets a new token instead of reading the reason. + assert result.code == "token_expired" def test_wrong_audience(self, oidc_env: bytes) -> None: token = _make_token(oidc_env, aud="wrong-client") @@ -768,3 +771,35 @@ def test_ambiguous_record_with_neither_field_fails_closed(self) -> None: def test_non_owner_denied_via_legacy_path(self) -> None: principal = Principal(kind="oidc", identifier="b", owner_id="oidc:deadbeef") assert not principal_owns(created_by="oidc:a", owner_id=None, principal=principal) + + +class TestStableAuthCodes: + """An authentication failure says what kind it is without its sentence (#1000).""" + + def test_a_refused_credential_is_unauthorized(self, oidc_env: bytes) -> None: + result = authenticate(bearer_token=f"Bearer {_make_token(oidc_env, aud='wrong-client')}") + + assert isinstance(result, AuthError) + assert (result.status_code, result.code) == (401, "unauthorized") + + def test_an_unreachable_jwks_is_auth_unavailable(self) -> None: + assert AuthError(reason="auth service unavailable (jwks)", status_code=503).code == ( + "auth_unavailable" + ) + + def test_the_response_body_carries_the_code( + self, oidc_env: bytes, tmp_path: Path, monkeypatch: pytest.MonkeyPatch + ) -> None: + from kpubdata_builder.service import BuilderService, ServiceResponse, dispatch + + service = BuilderService(output_root=tmp_path, client_factory=lambda **_kw: None) + expired = _make_token(oidc_env, exp=int(time.time()) - 120) + + response = dispatch(service, "GET", "/datasets", None, bearer_token=f"Bearer {expired}") + + assert isinstance(response, ServiceResponse) + assert response.status_code == 401 + assert response.body == { + "error": "invalid token: ExpiredSignatureError", + "code": "token_expired", + } diff --git a/tests/unit/test_service.py b/tests/unit/test_service.py index 7af29120..ce8cd2f1 100644 --- a/tests/unit/test_service.py +++ b/tests/unit/test_service.py @@ -1716,7 +1716,7 @@ def test_overloaded_response_is_a_well_formed_http_message(self) -> None: head, _, body = _OVERLOADED_RESPONSE.partition(b"\r\n\r\n") headers = dict(line.split(b": ", 1) for line in head.split(b"\r\n")[1:]) assert int(headers[b"Content-Length"]) == len(body) - assert json.loads(body) == {"error": "server overloaded"} + assert json.loads(body) == {"error": "server overloaded", "code": "server_overloaded"} @staticmethod def _overloaded_headers(allowed: frozenset[str]) -> dict[bytes, bytes]: @@ -1725,7 +1725,7 @@ def _overloaded_headers(allowed: frozenset[str]) -> dict[bytes, bytes]: head, _, body = _overloaded_response(allowed).partition(b"\r\n\r\n") headers = dict(line.split(b": ", 1) for line in head.split(b"\r\n")[1:]) assert int(headers[b"Content-Length"]) == len(body) - assert json.loads(body) == {"error": "server overloaded"} + assert json.loads(body) == {"error": "server overloaded", "code": "server_overloaded"} return headers def test_overloaded_response_has_no_cors_header_without_an_allowed_origin(self) -> None: diff --git a/tests/unit/test_stable_error_codes.py b/tests/unit/test_stable_error_codes.py new file mode 100644 index 00000000..5da20a11 --- /dev/null +++ b/tests/unit/test_stable_error_codes.py @@ -0,0 +1,36 @@ +"""Failures a client had to tell apart by their sentence carry a stable code (#1000).""" + +from __future__ import annotations + +import json +import threading +from pathlib import Path + +from kpubdata_builder.service.http import _overloaded_response +from tests.unit.test_service_jobs import VALID_SPEC_YAML, _BlockingBuildService + + +def test_a_full_build_queue_says_so_with_a_code(tmp_path: Path) -> None: + """429 was shared with ``auth_throttled`` and had no code of its own.""" + entered = threading.Event() + release = threading.Event() + service = _BlockingBuildService( + output_root=tmp_path, entered=entered, release=release, async_max_queue_size=1 + ) + try: + service.submit_build(VALID_SPEC_YAML, run_id="run1", created_by="tester") + assert entered.wait(timeout=5) + service.submit_build(VALID_SPEC_YAML, run_id="run2", created_by="tester") + + saturated = service.submit_build(VALID_SPEC_YAML, run_id="run3", created_by="tester") + finally: + release.set() + + assert saturated.status_code == 429 + assert saturated.body == {"error": "async build queue is full", "code": "build_queue_full"} + + +def test_the_overload_response_carries_its_code() -> None: + _head, _, body = _overloaded_response().partition(b"\r\n\r\n") + + assert json.loads(body)["code"] == "server_overloaded" From e55fa050400f64613dec88e394c23324f058b838 Mon Sep 17 00:00:00 2001 From: Eomdahyeon <213566566+Eomdahyeon@users.noreply.github.com> Date: Mon, 5 Oct 2026 00:06:11 +0900 Subject: [PATCH 3/5] docs(contract): declare the provider key header, the shared 429 and 503, and the request id Refs #994 Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 1 + contract/builder-api.yaml | 75 ++++++++++- contract/fixtures/responses.json | 49 ++++++- src/kpubdata_builder/service/app.py | 4 +- .../unit/test_contract_shared_declarations.py | 125 ++++++++++++++++++ 5 files changed, 251 insertions(+), 3 deletions(-) create mode 100644 tests/unit/test_contract_shared_declarations.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 837fb378..a0124633 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -26,6 +26,7 @@ ### Changed +- The contract declares four things the service already sent (#994, API contract 1.80.0, description only). `X-Provider-Key` is a header parameter of the five operations that call a provider (`previewBuild`, `createBuild`, `submitBuild`, `getProviderStatus`, `testProviderConnection`) — it was mentioned only in `info.description`, so a client's drift test had nothing to check it against. The 429 `auth_throttled` and the overload 503 are shared responses any operation can answer (`components.responses.AuthThrottled`, `ServerOverloaded`), and `X-Request-ID` is `components.headers.RequestId`. `tests/unit/test_contract_shared_declarations.py` compares each with what the code sends. Generating the route table from the route adapter, the remaining part of #994, is not done. - Four documents describe the service as it is (#1001). `BUILD_STATE.md` §8 no longer calls the async job state machine a future plan outside the contract; the contract's `info.description` no longer says async and publish endpoints are excluded; `API_CONTRACT.md` lists the preflight methods and headers the server sends (`GET, POST, PUT, DELETE, OPTIONS`; `X-Provider-Key` and `X-Publish-Credential` among the headers); and `API_CONTRACT.md` states by deployment mode whether a provider key is stored — not stored in a multi-user deployment, stored encrypted in a single-user one. No behaviour or schema changed, so the contract version stays. - The package metadata names its repository (kpubdata#791): `[project.urls]` gives the homepage, documentation, repository, issue tracker and changelog under `kpubdata-lab/kpubdata-builder`. There was no `[project.urls]`. - The repository moved from `yeongseon/kpubdata-builder` to `kpubdata-lab/kpubdata-builder`. Links, the documentation site (`https://kpubdata-lab.github.io/kpubdata-builder/`) and the shared GitHub Actions references now use the new owner. Container images are now published to `ghcr.io/kpubdata-lab/kpubdata-builder`; images already pushed to `ghcr.io/yeongseon/kpubdata-builder` stay where they are. diff --git a/contract/builder-api.yaml b/contract/builder-api.yaml index 79d6bda4..62757f7c 100644 --- a/contract/builder-api.yaml +++ b/contract/builder-api.yaml @@ -1,7 +1,7 @@ openapi: 3.1.0 info: title: KPubData Builder Service API - version: "1.79.0" + version: "1.80.0" description: >- KPubData Builder(패키지 `kpubdata-builder`) 서비스 계약. BuildSpec 검증, preview, build 실행, artifact 조회, 계약 버전 조회를 제공한다. CLI와 HTTP service mode(#36)가 동일한 도메인 계약을 @@ -20,6 +20,7 @@ info: 추가한다(#504, additive). v1.8.0은 stable principal별 encrypted Provider credential CRUD와 Provider status/test API를 추가한다(#492, additive). v1.9.0은 `/query` 응답에 child startup과 Polars engine 실행 시간을 추가한다(#523, additive). + v1.80.0 declares what the service already sends (#994, description only — the wire is unchanged): the `X-Provider-Key` request header as a parameter of the five operations that call a provider (`previewBuild`, `createBuild`, `submitBuild`, `getProviderStatus`, `testProviderConnection`), the 429 `auth_throttled` and the overload 503 as shared responses any operation can answer (`components.responses.AuthThrottled`, `ServerOverloaded`), and the `X-Request-ID` response header (`components.headers.RequestId`). v1.79.0 gives four failures a stable `code` (#1000, additive): the full async build queue's 429 is `build_queue_full` (it was told from 429 `auth_throttled` only by its sentence), an authentication failure is `unauthorized`, `token_expired` (the bearer token's `exp` has passed — get a new token) or, with 503, `auth_unavailable` (the JWKS could not be fetched), and the overload 503 written before a request is read is `server_overloaded`. Each body keeps its `error` text. v1.78.0 lets the owner read a run a restart interrupted (#996, additive): `GET /builds/{run_id}` answers such a run as `failed` with the new optional `BuildJob.code` `credentials_required` and the reason in `error`, and `GET /builds/{run_id}/events` shows its `run_failed` event, where both answered 404 before; another user still gets what a missing run gets. The submitter of an async run is now recorded with its events, so this holds for runs submitted from this version on. v1.77.0 stops a snapshot profile's range from disclosing a single record's value (#903, additive): `ColumnRange.min`/`max` are read after removing the `range_trim` lowest and the `range_trim` highest values, the range's `status` is `trimmed` instead of `exact`, and it carries `trimmed_count`. `SnapshotProfile` gains `range_trim`, and `min_range_values` is at least `2 * range_trim + 1`. A client that showed only `exact` ranges shows none until it reads `trimmed`. @@ -288,6 +289,7 @@ paths: operationId: getProviderStatus summary: 현재 principal credential로 Provider 연결 상태 확인 parameters: + - $ref: "#/components/parameters/ProviderKey" - $ref: "#/components/parameters/ProviderName" responses: "200": @@ -336,6 +338,7 @@ paths: operationId: testProviderConnection summary: 현재 principal credential로 lightweight connection test 실행 parameters: + - $ref: "#/components/parameters/ProviderKey" - $ref: "#/components/parameters/ProviderName" responses: "200": @@ -682,6 +685,8 @@ paths: post: operationId: previewBuild summary: 각 소스의 스키마와 샘플 행 산출 (파일 미기록, 동기식) + parameters: + - $ref: "#/components/parameters/ProviderKey" requestBody: required: true content: @@ -805,6 +810,8 @@ paths: post: operationId: createBuild summary: 파이프라인 실행 및 결과 반환 (동기식) + parameters: + - $ref: "#/components/parameters/ProviderKey" requestBody: required: true content: @@ -1376,6 +1383,8 @@ paths: 제출한 principal 의 업로드를 읽고, 다른 소유자의 `upload_id` 나 안정적인 소유자가 없는 요청은 job 이 `failed` 로 끝나며 그 source 의 `error` 에 사유가 적힌다. 제출 시점(202)에는 업로드를 확인하지 않는다. + parameters: + - $ref: "#/components/parameters/ProviderKey" requestBody: required: true content: @@ -4162,7 +4171,57 @@ components: each operation: an operation has one 403 response, and many already use it for their own refusal, whose `Error.code` tells the two apart. + headers: + RequestId: + description: >- + The id this request is logged under (#994). Sent on every response a handler + writes — success or error — so it can be quoted when reporting a problem; + absent only from the overload 503, which is written before a request is read. + Readable from a cross-origin page (`Access-Control-Expose-Headers`, #995). + schema: + type: string + responses: + AuthThrottled: + x-status: 429 + description: >- + Too many failed authentication attempts from this client address (#994) — any + operation can answer this before it runs. `retry_after_seconds` says how long + the window has left. `x-status` names the status code, since no single + operation references this response. + content: + application/json: + schema: + $ref: "#/components/schemas/Error" + examples: + AuthThrottled: + summary: The failure limit for this client address was reached + value: + error: too many failed authentication attempts + code: auth_throttled + retry_after_seconds: 42 + ServerOverloaded: + x-status: 503 + description: >- + Too many requests are in flight (#994, #1000). The connection is answered + before the request is read and closed; `Retry-After` says when to try again. + Any operation can answer this. `x-status` names the status code, since no + single operation references this response. + headers: + Retry-After: + description: Seconds to wait before sending the request again. + schema: + type: string + content: + application/json: + schema: + $ref: "#/components/schemas/Error" + examples: + ServerOverloaded: + summary: The request was refused before it was read + value: + error: server overloaded + code: server_overloaded SignupNotApproved: x-status: 403 description: >- @@ -4247,6 +4306,20 @@ components: reason: "the dataset declares redistribution: forbidden" parameters: + ProviderKey: + name: X-Provider-Key + in: header + required: false + description: >- + Multi-user deployment only (#683): the requester's own provider key for this + request or the job it submits, `=`, repeated or comma-separated. + Held in memory until the request ends, or the async job ends, is cancelled or + its TTL passes; never stored, logged or returned. A malformed header answers + 400 `invalid_provider_key` on any route. A single-user deployment ignores it + and keeps its stored and server keys. Declared on the operations that call a + provider (#994). + schema: + type: string PublishCredential: name: X-Publish-Credential in: header diff --git a/contract/fixtures/responses.json b/contract/fixtures/responses.json index 2ff29d91..adc490c1 100644 --- a/contract/fixtures/responses.json +++ b/contract/fixtures/responses.json @@ -1,6 +1,6 @@ { "fixture_format": 1, - "contract_version": "1.79.0", + "contract_version": "1.80.0", "probe_field": "future_optional_field", "rules": "A response may gain optional fields in any minor version; a client ignores fields it does not know. A required field keeps its name and type until the next major version; a client rejects a body whose required field is missing or mistyped.", "fixtures": [ @@ -5669,6 +5669,53 @@ }, "broken_path": "$.error" }, + { + "response": "AuthThrottled", + "status": 429, + "example": "AuthThrottled", + "current": { + "error": "too many failed authentication attempts", + "code": "auth_throttled", + "retry_after_seconds": 42 + }, + "with_additive_fields": { + "error": "too many failed authentication attempts", + "code": "auth_throttled", + "retry_after_seconds": 42, + "future_optional_field": "added by a later minor contract version" + }, + "additive_paths": [ + "$" + ], + "required_type_broken": { + "error": 12345, + "code": "auth_throttled", + "retry_after_seconds": 42 + }, + "broken_path": "$.error" + }, + { + "response": "ServerOverloaded", + "status": 503, + "example": "ServerOverloaded", + "current": { + "error": "server overloaded", + "code": "server_overloaded" + }, + "with_additive_fields": { + "error": "server overloaded", + "code": "server_overloaded", + "future_optional_field": "added by a later minor contract version" + }, + "additive_paths": [ + "$" + ], + "required_type_broken": { + "error": 12345, + "code": "server_overloaded" + }, + "broken_path": "$.error" + }, { "response": "SignupNotApproved", "status": 403, diff --git a/src/kpubdata_builder/service/app.py b/src/kpubdata_builder/service/app.py index c0b21356..2808af38 100644 --- a/src/kpubdata_builder/service/app.py +++ b/src/kpubdata_builder/service/app.py @@ -402,7 +402,9 @@ def _enforce_ownership() -> bool: # 1.78.0 -> 1.79.0: stable codes for the full build queue (429 build_queue_full), # authentication failures (unauthorized, token_expired, auth_unavailable) and the # overload 503 (server_overloaded) (#1000, additive). -API_CONTRACT_VERSION = "1.79.0" +# 1.79.0 -> 1.80.0: the X-Provider-Key parameter, the shared 429 auth_throttled and +# overload 503 responses and the X-Request-ID header are declared (#994, description). +API_CONTRACT_VERSION = "1.80.0" #: manifest status vocabulary (ok/failed/cancelled) → publish status vocabulary diff --git a/tests/unit/test_contract_shared_declarations.py b/tests/unit/test_contract_shared_declarations.py new file mode 100644 index 00000000..edebc564 --- /dev/null +++ b/tests/unit/test_contract_shared_declarations.py @@ -0,0 +1,125 @@ +"""What the service sends on every route is declared in the contract (#994). + +The ``X-Provider-Key`` request header, the 429 ``auth_throttled``, the overload 503 and +the ``X-Request-ID`` response header existed in the service and only in prose in the +contract. These compare the declarations with what the code actually sends. +""" + +from __future__ import annotations + +import json +from pathlib import Path +from typing import Any + +import pytest +import yaml + +from kpubdata_builder.service.http import _overloaded_response +from kpubdata_builder.service.request_credentials import PROVIDER_KEY_HEADER + +_CONTRACT = Path(__file__).resolve().parents[2] / "contract" / "builder-api.yaml" + +#: The operations that call a provider with the requester's key. +_PROVIDER_OPERATIONS = { + "previewBuild", + "createBuild", + "submitBuild", + "getProviderStatus", + "testProviderConnection", +} + + +@pytest.fixture(scope="module") +def contract() -> dict[str, Any]: + loaded: dict[str, Any] = yaml.safe_load(_CONTRACT.read_text(encoding="utf-8")) + return loaded + + +def _operations(contract: dict[str, Any]) -> dict[str, dict[str, Any]]: + return { + operation["operationId"]: operation + for item in contract["paths"].values() + for method, operation in item.items() + if method in ("get", "post", "put", "delete") and "operationId" in operation + } + + +def test_the_provider_key_header_is_a_declared_parameter(contract: dict[str, Any]) -> None: + parameter = contract["components"]["parameters"]["ProviderKey"] + + assert (parameter["name"], parameter["in"], parameter["required"]) == ( + PROVIDER_KEY_HEADER, + "header", + False, + ) + + +def test_it_is_declared_on_exactly_the_operations_that_call_a_provider( + contract: dict[str, Any], +) -> None: + reference = {"$ref": "#/components/parameters/ProviderKey"} + declaring = { + name + for name, operation in _operations(contract).items() + if reference in operation.get("parameters", []) + } + + assert declaring == _PROVIDER_OPERATIONS + + +def test_the_overload_response_is_the_declared_one(contract: dict[str, Any]) -> None: + declared = contract["components"]["responses"]["ServerOverloaded"] + head, _, body = _overloaded_response().partition(b"\r\n\r\n") + headers = dict(line.split(b": ", 1) for line in head.split(b"\r\n")[1:]) + + assert head.startswith(b"HTTP/1.1 %d " % declared["x-status"]) + example = declared["content"]["application/json"]["examples"]["ServerOverloaded"]["value"] + assert json.loads(body) == example + assert set(declared["headers"]) <= {name.decode() for name in headers} + + +def test_the_auth_throttle_response_is_the_declared_one( + contract: dict[str, Any], tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + from kpubdata_builder.service import BuilderService, ServiceResponse, dispatch + + declared = contract["components"]["responses"]["AuthThrottled"] + # The suite runs in dev mode, which skips authentication; this needs it on. + monkeypatch.delenv("KPUBDATA_BUILDER_DEV_MODE", raising=False) + monkeypatch.setenv("KPUBDATA_BUILDER_API_KEY", "the-right-key") + monkeypatch.setenv("KPUBDATA_BUILDER_AUTH_FAILURE_LIMIT", "1") + service = BuilderService(output_root=tmp_path, client_factory=lambda **_kw: None) + + first = dispatch(service, "GET", "/datasets", None, api_key="wrong", client_id="10.0.0.9") + second = dispatch(service, "GET", "/datasets", None, api_key="wrong", client_id="10.0.0.9") + + assert isinstance(first, ServiceResponse) and first.status_code == 401 + assert isinstance(second, ServiceResponse) + assert second.status_code == declared["x-status"] + example = declared["content"]["application/json"]["examples"]["AuthThrottled"]["value"] + assert set(second.body) == set(example) + assert (second.body["error"], second.body["code"]) == (example["error"], example["code"]) + + +def test_the_request_id_header_is_declared_and_sent( + contract: dict[str, Any], tmp_path: Path +) -> None: + import threading + import urllib.request + from http.server import HTTPServer + + from kpubdata_builder.service import BuilderService + from kpubdata_builder.service.http import make_handler + + assert "RequestId" in contract["components"]["headers"] + service = BuilderService(output_root=tmp_path, client_factory=lambda **_kw: None) + server = HTTPServer(("127.0.0.1", 0), make_handler(service)) + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + try: + url = f"http://127.0.0.1:{server.server_address[1]}/healthz" + with urllib.request.urlopen(url, timeout=5.0) as response: + assert response.headers["X-Request-ID"] + finally: + server.shutdown() + server.server_close() From 7df0ed4fa601fc78a53192d7336a53dd4f29760c Mon Sep 17 00:00:00 2001 From: Eomdahyeon <213566566+Eomdahyeon@users.noreply.github.com> Date: Mon, 5 Oct 2026 00:17:18 +0900 Subject: [PATCH 4/5] test: wait for the released worker before the next test runs The fixture released the first service's worker and returned. The worker then ran its build behind the following tests, and test_json_array_reader measures peak memory with tracemalloc, which counts every thread. Co-Authored-By: Claude Opus 5.5 --- tests/unit/test_interrupted_run_status.py | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/tests/unit/test_interrupted_run_status.py b/tests/unit/test_interrupted_run_status.py index be7491e2..eae8d0cf 100644 --- a/tests/unit/test_interrupted_run_status.py +++ b/tests/unit/test_interrupted_run_status.py @@ -103,8 +103,10 @@ def restarted(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> Iterator[Build yield second finally: gate.set() - first._async_builds.shutdown() - second._async_builds.shutdown() + # Wait for the released worker: left running, it allocates behind whichever + # test comes next, and one of those measures peak memory. + first._async_builds._executor.shutdown(wait=True) + second._async_builds._executor.shutdown(wait=True) def test_the_owner_reads_failed_with_a_stable_code( From 8fcb40d7782c8c074c4ef72d6d2514e5a6979c49 Mon Sep 17 00:00:00 2001 From: Eomdahyeon <213566566+Eomdahyeon@users.noreply.github.com> Date: Mon, 5 Oct 2026 00:18:17 +0900 Subject: [PATCH 5/5] fix(auth): set token_expired where the failure is made The code was read back from the end of the reason sentence, so rewording the reason would have turned token_expired into unauthorized without a failure anywhere. AuthError now carries the fact. Also pin that a missing API key and a wrong one both answer with the contract example's sentence. Co-Authored-By: Claude Opus 5.5 --- src/kpubdata_builder/service/auth.py | 10 ++++++++-- tests/unit/test_oidc_auth.py | 19 +++++++++++++++++++ 2 files changed, 27 insertions(+), 2 deletions(-) diff --git a/src/kpubdata_builder/service/auth.py b/src/kpubdata_builder/service/auth.py index 14892346..7e38f598 100644 --- a/src/kpubdata_builder/service/auth.py +++ b/src/kpubdata_builder/service/auth.py @@ -167,6 +167,9 @@ class AuthError: reason: str status_code: int = 401 + # Set where the failure is made, so ``code`` does not depend on how ``reason`` is + # worded. + expired: bool = False @property def code(self) -> str: @@ -179,7 +182,7 @@ def code(self) -> str: """ if self.status_code == 503: return "auth_unavailable" - if self.reason.endswith("ExpiredSignatureError"): + if self.expired: return "token_expired" return "unauthorized" @@ -472,7 +475,10 @@ def _verify_bearer_token(token: str) -> Principal | AuthError: options={"require": ["exp", "iat", "iss", "sub"]}, ) except jwt.PyJWTError as exc: - return AuthError(reason=f"invalid token: {type(exc).__name__}") + return AuthError( + reason=f"invalid token: {type(exc).__name__}", + expired=isinstance(exc, jwt.ExpiredSignatureError), + ) if not payload.get("email_verified", False): return AuthError(reason="email not verified") diff --git a/tests/unit/test_oidc_auth.py b/tests/unit/test_oidc_auth.py index c9250fe7..adaecb97 100644 --- a/tests/unit/test_oidc_auth.py +++ b/tests/unit/test_oidc_auth.py @@ -782,6 +782,25 @@ def test_a_refused_credential_is_unauthorized(self, oidc_env: bytes) -> None: assert isinstance(result, AuthError) assert (result.status_code, result.code) == (401, "unauthorized") + def test_the_code_does_not_read_the_sentence(self) -> None: + # ``token_expired`` is set where the failure is made; rewording ``reason`` — or a + # reason that happens to end the same way — does not change the code. + assert AuthError(reason="the token is past its exp", expired=True).code == "token_expired" + assert AuthError(reason="invalid token: ExpiredSignatureError").code == "unauthorized" + + @pytest.mark.parametrize("sent", [None, "wrong-key"]) + def test_a_missing_key_and_a_wrong_key_get_the_contract_example( + self, sent: str | None, monkeypatch: pytest.MonkeyPatch + ) -> None: + # The contract's one ``Unauthorized`` example stands for both, so both must + # answer with its sentence. + monkeypatch.setenv("KPUBDATA_BUILDER_API_KEY", "secret-key") + + result = authenticate(api_key=sent) + + assert isinstance(result, AuthError) + assert (result.reason, result.code) == ("invalid api key", "unauthorized") + def test_an_unreachable_jwks_is_auth_unavailable(self) -> None: assert AuthError(reason="auth service unavailable (jwks)", status_code=503).code == ( "auth_unavailable"