diff --git a/.dockerignore b/.dockerignore index 43f5832..f00a0d4 100644 --- a/.dockerignore +++ b/.dockerignore @@ -58,6 +58,21 @@ Makefiles dev_context playground +# Per-developer Pipelex overrides. These files are untracked, so CI builds the published image +# without them — but a local `make docker-build` would otherwise bake one machine's settings into +# the image: its storage backend, its log level, its telemetry credentials. The image would then +# behave differently from the published one, and nothing in the diff would say so. An operator +# supplies overrides to a container by mounting them at /root/.pipelex, the documented path. +# +# The runtime layers more than the `_override` tier over a base file: `_local`, `_{environment}` +# selected by PIPELEX_ENV, and `_temporary_override`, for every configuration family in the +# directory rather than for `pipelex` alone. Environment names are open-ended, so the exclusion is +# by shape — any suffixed variant of a base file — and the one tracked variant is named back in. +# The nested inference overrides need their own line because a pattern's `*` does not cross a `/`. +.pipelex/*_*.toml +!.pipelex/pipelex_service.toml +.pipelex/inference/*_override.toml + # OS / system files .DS_Store Thumbs.db diff --git a/.pipelex/pipelex.toml b/.pipelex/pipelex.toml index f336f8e..dbb1fc1 100644 --- a/.pipelex/pipelex.toml +++ b/.pipelex/pipelex.toml @@ -155,9 +155,24 @@ signed_urls_lifespan_seconds = 3600 # Set to "disabled [runtime.log] # Default logging level: "DEBUG", "INFO", "WARNING", "ERROR" default_log_level = "INFO" -# Log output target: "stdout" or "stderr" -console_log_target = "stdout" +# The registered log sink this server installs. "json" writes one JSON object per line, with every +# field, the bound run identifiers and the message flat beside each other — the shape a log agent +# in front of a container ingests without a parser. The alternative, "console", is the Rich +# renderer meant for a terminal; this server does not ask for the `cli` extra that declares Rich, +# which is what makes the renderer unused here rather than unavailable — see pretty_print_mode. +sink = "json" +# The stream the json sink writes to: "stdout" or "stderr". Logs are diagnostics and belong off +# the data channel. +console_log_target = "stderr" console_print_target = "stdout" +# The panels `pretty_print(...)` draws, the "Output of pipe" one after every operator pipe among +# them. A server has no terminal to draw into and must not spend time rendering on the thread +# serving a request, so nothing is printed and no renderable is built. This is what keeps Rich off +# the request path, and it is not the same as keeping it out of the image: `typer` and `instructor` +# are core pipelex dependencies that require Rich unconditionally, so an `import pipelex` loads it +# whatever this file says. Selecting "rich" here, or the "console" sink above, would therefore work +# rather than refuse — these two keys are the whole of what keeps the renderer unused. +pretty_print_mode = "silent" [runtime.log.package_log_levels] # Log levels for specific packages (use "-" instead of "." in package names) diff --git a/CHANGELOG.md b/CHANGELOG.md index dce8668..fb221dc 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,21 @@ # Changelog +## [v0.28.0] - 2026-09-25 + +### Highlights + +**The runner's logs are structured.** Every line it writes to stderr is one JSON object with each value under a key of its own, the request id rides every record a request emits, the Pipelex runtime's included, and secrets are scrubbed before a line is written. + +### Changed + +- **The server's logs are structured, one JSON object per line on stderr (Breaking)**: `[runtime.log] sink = "json"` replaces the Rich console renderer, `console_log_target` moves to `stderr`, and `pretty_print_mode = "silent"` stops an operator pipe drawing its "Output of pipe" panel on the thread serving a request. Every value an error line carries — `route`, `status`, `error_type`, `error_domain`, `retryable`, `user_id` / `pipe_code` / `pipeline_run_id` when the request bound them, and `detail` on the failures this API authors itself — is a key of its own now rather than part of a `key=value` run inside the message, and `request_id` rides the runtime's request-scoped log context, so it lands on every record emitted during a request, including the ones Pipelex emits from inside a run. The message is a short sentence built only from the status and the error type, so no caller-supplied string reaches it and the API's own escaping is gone: the sink is what serializes a value now. A log query matching `event=api_error` as text has to move to the `event` field. The new `docs/logging.md` documents the line and every field on it. Uvicorn's own banner and access log are unchanged and still plain text. +- **Pinned `pipelex` 0.66.0 (Breaking)**: up from `==0.65.0`, exactly, the release that carries the structured-log seam the entry above rides on: named fields and the run-scoped log context, the `json` sink selected by `[runtime.log] sink`, the redaction of secrets before any sink sees a record, and Rich behind the `cli` extra, which this server's extras leave out. On this server's lines that means a credential echoed into `detail` reads `[REDACTED]`, a control character in a field's value reads as its printable escape (`\n` where a caller sent a newline), and a line logged inside a traced run carries `trace_id`, `span_id` and `trace_flags`; `docs/logging.md` says so. The `.pipelex/` config shipped here already sits at the schema that release migrates to, so no migration is required. One change reaches an operator beyond the logs: the S3 storage provider now signs links and reads and writes objects on the bucket's own regional host, `.s3..amazonaws.com`, so a deployment whose egress rules allow S3 by hostname must allow `*.s3..amazonaws.com`. Nothing on the wire moves. +- **`POST /v1/codegen` stamps `engine_version` `0.66.0`**: the stamp is the pinned `pipelex` version, so a `codegen.lock` committed against `0.65.0` no longer matches until it is regenerated. `POST /v1/build/runner` carries the same stamp. + +### Fixed + +- **A local image build no longer bakes the builder's own Pipelex overrides in**: `.pipelex/pipelex_override.toml` and `.pipelex/telemetry_override.toml` are untracked per-developer files, so CI never had them, but `make docker-build` copied whatever the developer had into the image — their storage backend, their log level, their telemetry credentials. A locally built image then behaved differently from the published one, with nothing in the diff to say so. `.dockerignore` excludes them; an operator still supplies overrides to a container by mounting them at `/root/.pipelex`. + ## [v0.27.5] - 2026-09-25 ### Changed diff --git a/CLAUDE.md b/CLAUDE.md index 2b994de..0b1c5fa 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -169,7 +169,7 @@ Every error is rendered as RFC 7807 `application/problem+json` by the global han - **Domain errors** (pipelex `PipelexError` subclasses) — raise from your code and let them propagate. The `PipelexError` global handler obtains an `ErrorReport` via `to_error_report()` and renders it into a problem document. Do not wrap, classify, or re-shape. - **API-authored 4xx/5xx** — use the helpers in `api/errors.py`: `raise_validation_error`, `raise_bad_request`, `raise_forbidden`, `raise_unauthenticated`, `raise_payload_too_large`, `raise_internal_server_error`. Each raises an `ApiError` carrying a pre-built problem document; the global handler emits it. **Do not raise `HTTPException` directly** — FastAPI's default handler wraps the body as `{"detail": }` and cannot emit a flat RFC 7807 document. - **Auth errors** — the helpers set `WWW-Authenticate: Bearer` automatically on 401. -- **Logging** — the global handlers emit one structured log line per error (`event=api_error`) with `request_id`, `route`, `error_type`, `error_domain`, `retryable`, `status`, and `user_id` when authenticated. Log disposition follows the final HTTP status: 4xx logs at `warning` (caller mistakes, the provider-429 passthrough, and API-level 4xx overrides like the 409 conflict); 5xx logs at `error` with traceback. Routes should not log error tracebacks themselves. +- **Logging** — the global handlers emit one structured record per error, carrying `event: "api_error"`, `route`, `error_type`, `error_domain`, `retryable`, `status`, and `user_id` / `pipe_code` / `pipeline_run_id` when the request bound them. `detail` rides the API-authored path only: a Pipelex `ErrorReport`'s body text has been through disclosure redaction, so it is not the cause and is deliberately not logged. They travel as `fields=` on the runtime's log call, never interpolated into the message, and `request_id` is not among them: `RequestIdMiddleware` binds it on the runtime's log context, so every record emitted under the request already carries it. The server selects the `json` sink, so a record reaches stderr as one JSON object per line — see `docs/logging.md`. Log disposition follows the final HTTP status: 4xx logs at `warning` (caller mistakes, the provider-429 passthrough, and API-level 4xx overrides like the 409 conflict); 5xx logs at `error` with traceback. Routes should not log error tracebacks themselves. - **Documenting a failure in OpenAPI** — the shared, typed `responses=` declarations live in `api/openapi_responses.py` (`ProblemDocument` + one constant per status). Every auth-wrapped `/v1` route already documents `401`/`413`/`422`/`500` via the composite router's `responses=` (`api/routes/__init__.py`); a route declares on its own decorator only the statuses **it alone** can produce. Never hand-write an error `content` block: `api/openapi_schema.py` re-keys the generated schema onto `application/problem+json`, because FastAPI renders a response `model` under the route's response-class media type and offers no per-response override. Adding a new status means adding a constant there and referencing it — then `make openapi-export`. Typical route: @@ -177,12 +177,13 @@ Typical route: ```python @router.post("/start", response_model=PipelexStartAck, status_code=202) async def start( - request: Annotated[RunRequest, Depends(request_deserialization)], + request: Request, + run_request: Annotated[RunRequest, Depends(request_deserialization)], user: Annotated[RequestUser | None, Depends(get_optional_user)], - request_id: Annotated[str, Depends(get_request_id)], ) -> PipelexStartAck: # Let PipelexError / EnvVarNotFoundError / etc. propagate to the global handler. - return await api_runner.start(request, user=user, request_id=request_id) + # `request_id_of` (api/middleware.py) reads back what RequestIdMiddleware put on the request. + return await api_runner.start(run_request, user=user, request_id=request_id_of(request)) ``` For an API-authored failure that has no `PipelexError`: diff --git a/Dockerfile b/Dockerfile index 8659a3e..b88201e 100644 --- a/Dockerfile +++ b/Dockerfile @@ -13,7 +13,11 @@ RUN pip install --no-cache-dir uv WORKDIR /app -# Copy lockfile + project metadata first so dependency-install layer is cacheable +# Copy lockfile + project metadata first so dependency-install layer is cacheable. +# The extras this installs are the ones declared in pyproject.toml, and `cli` is deliberately not +# among them: the server selects the `json` log sink and a Rich-free pretty-print mode, so nothing +# it does on a request renders a terminal. (Rich itself is still in the image — typer and +# instructor both require it unconditionally — it is simply never reached.) COPY pyproject.toml uv.lock ./ RUN uv sync --frozen --no-dev --no-install-project diff --git a/api/errors.py b/api/errors.py index 5f427ae..d94fb09 100644 --- a/api/errors.py +++ b/api/errors.py @@ -20,7 +20,6 @@ from pipelex.base_exceptions import ErrorDomain from api.error_types import ErrorType -from api.logging_context import get_request_id, get_route_path from api.problem_document import build_problem_document_from_api_error @@ -31,7 +30,9 @@ class ApiError(Exception): `api.exception_handlers` as `application/problem+json`. Distinct from a pipelex `PipelexError`: there is no `ErrorReport` behind it — the failure is the API's own request validation, auth, or configuration check. The problem - document is built at raise time so the handler only has to serialize it. + document is built at raise time, less the request context: `instance` and + `request_id` are stamped by the handler, which is the frame that holds the + `Request`. """ def __init__(self, *, status_code: int, document: dict[str, Any], headers: dict[str, str] | None = None) -> None: @@ -51,16 +52,19 @@ def _raise_api_error( ) -> NoReturn: """Build the RFC 7807 document and raise `ApiError`. - `instance` and `request_id` come from the request-scoped logging - contextvars (`api.logging_context`), bound by `RequestIdMiddleware`, so the - helpers stay parameter-clean and call sites need no `Request`. + The document is built without the request context: `instance` and + `request_id` are stamped by `handle_api_error`, from the `Request` it is + handed. That keeps these helpers parameter-clean — a call site deep inside + a route still needs no `Request` — while leaving the API with no ambient + request state of its own, and it is how the three error paths end up + reading the route and the id from exactly one place. """ document = build_problem_document_from_api_error( error_type, message, status, - instance=get_route_path(), - request_id=get_request_id(), + instance=None, + request_id=None, error_domain=error_domain, ) raise ApiError(status_code=status, document=document, headers=headers) diff --git a/api/exception_handlers.py b/api/exception_handlers.py index 44f7543..9773faa 100644 --- a/api/exception_handlers.py +++ b/api/exception_handlers.py @@ -27,8 +27,7 @@ """ import math -import re -from collections.abc import Awaitable, Callable +from collections.abc import Awaitable, Callable, Mapping from typing import TYPE_CHECKING, Any, cast from fastapi import FastAPI, Request, Response @@ -40,7 +39,13 @@ from api.error_types import ErrorType from api.errors import ApiError -from api.problem_document import PROBLEM_JSON_MEDIA_TYPE, build_problem_document, build_problem_document_from_api_error +from api.middleware import request_id_of +from api.problem_document import ( + PROBLEM_JSON_MEDIA_TYPE, + build_problem_document, + build_problem_document_from_api_error, + with_request_context, +) if TYPE_CHECKING: from api.security import RequestUser @@ -48,14 +53,9 @@ # A Starlette/FastAPI async exception handler: `(request, exc) -> response`. _ExceptionHandler = Callable[[Request, Exception], Awaitable[Response]] - -def _request_id_of(request: Request) -> str | None: - """Return the correlation id `RequestIdMiddleware` stored on the request. - - Read defensively with `getattr`: a request that never went through the - middleware (a unit test, a non-HTTP scope) simply has no id. - """ - return getattr(request.state, "request_id", None) +# The value of the `event` field every error record carries: the one key a log sink filters this +# server's error stream on, whichever of the three handlers below produced the record. +API_ERROR_EVENT = "api_error" def _user_id_of(request: Request) -> str | None: @@ -68,9 +68,9 @@ def _user_id_of(request: Request) -> str | None: multiple log lines by `request_id` — and without each route having to log `user=` itself before raising (which Phase 3 removed). Returns `None` pre-auth, on the static-API-key surface (no per-caller identity), or for - `AUTH_MODE=none` without the forwarded-identity opt-in: `emit_error_log` - drops `None`-valued fields, so the field is absent from the rendered line - rather than `user_id=None`. + `AUTH_MODE=none` without the forwarded-identity opt-in: `_emit_api_error` + drops `None`-valued fields, so the attribute is absent from the record + rather than carried as `user_id: null`. """ user: RequestUser | None = getattr(request.state, "user", None) return user.user_id if user is not None else None @@ -83,11 +83,11 @@ def _pipe_code_of(request: Request) -> str | None: `request.state` right after the body decodes (before `_validate_extras` / `from_body`), normalized through `_coerce_correlation_field` — empty / non-string / oversized inputs become - `None` so a caller cannot inject a bare `pipe_code=` token or inflate the - log line. Returns `None` for routes that don't use `_parse_request`, for - requests whose body never decoded (`_decode_body` raised 422), and for - bodies that legitimately omitted `pipe_code` (the `mthds_contents`-only - invocation). `emit_error_log` drops `None`-valued fields. + `None`, so a caller cannot inflate every record the request emits. Returns + `None` for routes that don't use `_parse_request`, for requests whose body + never decoded (`_decode_body` raised 422), and for bodies that legitimately + omitted `pipe_code` (the `mthds_contents`-only invocation). + `_emit_api_error` drops `None`-valued fields. """ return getattr(request.state, "pipe_code", None) @@ -106,17 +106,27 @@ def _pipeline_run_id_of(request: Request) -> str | None: return getattr(request.state, "pipeline_run_id", None) -def _request_correlation_fields(request: Request) -> dict[str, str | None]: - """Return the request-scoped correlation fields every error log carries. +def _request_fields(request: Request) -> dict[str, Any]: + """Return the request-scoped attributes every error record carries. + + Single source of truth for the `route` / `user_id` / `pipe_code` / + `pipeline_run_id` set, so the three log paths (`_log_error_report`, + `_log_api_authored_error`, `handle_unexpected_error`) cannot drift. - Single source of truth for the `user_id` / `pipe_code` / `pipeline_run_id` - field set so the three log paths (`_log_error_report`, - `_log_api_authored_error`, `handle_unexpected_error`) cannot drift. Each - value is `None` when the corresponding state is not bound on this request; - `emit_error_log` drops `None`-valued fields, so unbound identifiers are - absent from the rendered line rather than appearing as `pipe_code=None`. + `request_id` is deliberately NOT here. `RequestIdMiddleware` binds it on the + runtime's log context for the whole request, so it is already an attribute of + every record emitted underneath — including the ones pipelex emits from inside + a run, which this module never sees. Repeating it would give one value two + sources. `route` is here rather than on that context because the runtime + reserves the context for its three run identifiers, and a route path is not + one of them; a field is the seam it offers for everything else. + + Each value is `None` when the corresponding state is not bound on this request; + `_emit_api_error` drops those, so an unbound identifier is absent from the + record rather than carried as `pipe_code: null`. """ return { + "route": request.url.path, "user_id": _user_id_of(request), "pipe_code": _pipe_code_of(request), "pipeline_run_id": _pipeline_run_id_of(request), @@ -145,52 +155,42 @@ def _retry_after_header(report: ErrorReport) -> dict[str, str]: return {"Retry-After": str(max(0, math.ceil(seconds)))} -_LOGFMT_NEEDS_QUOTING = re.compile(r'[\s"=]') +def _error_summary(attributes: Mapping[str, Any]) -> str: + """Return the human-readable message an `api_error` record carries. + + Built from server-authored values only — the HTTP status and the error type, + both of which the API or pipelex chose. Nothing a caller supplied ever reaches + the message: `detail` is caller-controlled on several routes, and the route + path is percent-decoded from the request line, so both ride fields instead, + where a structured sink serializes them as values and a crafted newline stays + inside one. The message says which failure it is; the fields say everything + about it. + """ + status = attributes.get("status") + error_type = attributes.get("error_type") + return f"API error {status}: {error_type}" if error_type else f"API error {status}" -def _logfmt_value(value: Any) -> str: - """Render a value for the ``key=value`` log format, escaping caller input. +def _emit_api_error(*, fields: dict[str, Any], as_error: bool) -> None: + """Emit one `event=api_error` record: a summary message, everything else a record attribute. - Several `emit_error_log` callers ship caller-controlled strings into the - field map — `_log_api_authored_error` forwards `document["detail"]` - (a validation message, a callback-URL rejection reason, etc.) and the - catch-all handler ships `type(exc).__name__` which can be anything pipelex - raised. Without escaping, a crafted body with a newline or whitespace - inside `detail` would break the field separator (`key=value key2=value2`) - or forge extra log fields (`detail=ok status=200 event=fake`) — both are - log-injection vectors. + The fields are handed to the runtime's log call as `fields=`, so each one + becomes an attribute of the record and the selected sink decides how it goes + on the wire — a key of its own on the JSON sink's line, an OTLP attribute on + the collector's. Nothing is flattened into the message here any more, which + is what retires the API's own `key=value` rendering and its escaping with it. - Non-string values (`int`, `bool`, `StrEnum` instances) have no injection - surface and render bare. Strings get logfmt-style treatment: control chars - are backslash-escaped first, then values containing whitespace, `=`, or - `"` are wrapped in double quotes with embedded `"` doubled. The shape - survives both grep (the line stays single-line and the keys remain at the - same offsets) and a future JSON log sink (each field is recoverable). - """ - if not isinstance(value, str): - return str(value) - escaped = value.encode("unicode_escape").decode("ascii") - if _LOGFMT_NEEDS_QUOTING.search(escaped): - return '"' + escaped.replace('"', '""') + '"' - return escaped - - -def emit_error_log(*, fields: dict[str, Any], as_error: bool) -> None: - """Emit one structured error-log line from a flat field map. - - The pipelex `log` object renders a single message string rather than - indexed key/value fields, so the fields are flattened to a `key=value` - run: greppable today, and a clean migration target once a JSON log sink - lands. `None`-valued fields are dropped. `as_error` picks the level — - `error` (with traceback) for operator-actionable failures, `warning` for - `INPUT`-domain caller mistakes. Caller-controlled values go through - `_logfmt_value` so a crafted `detail` can't forge log fields. + `None`-valued fields are dropped, so an identifier this request never bound is + absent rather than carried as a null. `as_error` picks the level — `error` + (with the traceback) for operator-actionable failures, `warning` for + `INPUT`-domain caller mistakes. """ - rendered = " ".join(f"{key}={_logfmt_value(value)}" for key, value in fields.items() if value is not None) + attributes = {key: value for key, value in fields.items() if value is not None} + summary = _error_summary(attributes) if as_error: - log.error(rendered, include_exception=True) + log.error(summary, fields=attributes, include_exception=True) else: - log.warning(rendered) + log.warning(summary, fields=attributes) def _emit_at_error_level(status: int) -> bool: @@ -207,31 +207,29 @@ def _emit_at_error_level(status: int) -> bool: return status >= 500 -def _log_error_report(report: ErrorReport, *, request: Request, request_id: str | None, status: int | None = None) -> None: - """Emit the structured log entry for a handled `ErrorReport`. +def _log_error_report(report: ErrorReport, *, request: Request, status: int | None = None) -> None: + """Emit the structured error record for a handled `ErrorReport`. Disposition follows the final HTTP status (see ``_emit_at_error_level``): a 4xx is client-facing and logs at `warning` without a traceback (caller mistakes, the provider-429 passthrough, and API-level 4xx overrides like the 409 conflict), a 5xx logs at `error` with the traceback. The fields mirror the response so the two never drift. - `user_id` rides every line when auth bound a caller — without it, the + `user_id` rides every record when auth bound a caller — without it, the storage / pipeline-backend leg of a failure carries only `request_id` and `route`, and tying the failure to the caller requires correlating the - request id across unrelated log lines (Phase 3 deleted the per-route + request id across unrelated records (Phase 3 deleted the per-route `log.error(... user=...)` lines those failures used to emit). ``status`` defaults to ``report.http_status`` and exists so the caller can pass the post-override value (see ``_ERROR_TYPE_STATUS_OVERRIDES``) — the - log line then agrees with the HTTP status actually sent rather than the + record then agrees with the HTTP status actually sent rather than the domain default. """ effective_status = status if status is not None else report.http_status fields: dict[str, Any] = { - "event": "api_error", - "request_id": request_id, - "route": request.url.path, - **_request_correlation_fields(request), + "event": API_ERROR_EVENT, + **_request_fields(request), "error_type": report.error_type, "error_category": report.error_category, "error_domain": report.error_domain, @@ -244,16 +242,16 @@ def _log_error_report(report: ErrorReport, *, request: Request, request_id: str if metadata is not None: fields["provider_status_code"] = metadata.status_code fields["provider_request_id"] = metadata.request_id - emit_error_log(fields=fields, as_error=_emit_at_error_level(effective_status)) + _emit_api_error(fields=fields, as_error=_emit_at_error_level(effective_status)) -def _log_api_authored_error(*, document: dict[str, Any], status: int, request: Request, request_id: str | None) -> None: - """Emit the structured log entry for an API-authored error response. +def _log_api_authored_error(*, document: dict[str, Any], status: int, request: Request) -> None: + """Emit the structured error record for an API-authored error response. Shares the disposition rule and the common-key set of `_log_error_report` so every error response — a pipelex `ErrorReport` translated to RFC 7807 *or* an API-authored 4xx/5xx raised by an `api.errors` helper — produces - one `event=api_error` line a downstream sink can grep uniformly on + one `event=api_error` record a downstream sink can filter uniformly on `event`, `request_id`, `route`, `error_type`, `error_domain`, `retryable`, and `status`. Without this, an API-owned 500 (a `raise_internal_server_error` site — `/version`'s missing-package case is the canonical example) @@ -265,7 +263,7 @@ def _log_api_authored_error(*, document: dict[str, Any], status: int, request: R - API-authored docs add `detail`, always safe because `build_problem_document_from_api_error` does not apply strict-disclosure redaction (only `build_problem_document` does for pipelex domain errors). - Carrying the message preserves the operator-facing cause in the log line. + Carrying the message preserves the operator-facing cause in the record. - API-authored docs omit `error_category`, `provider`, `model`, and `provider_metadata.*` — those are inference-domain classifiers pipelex sets only on classifiable failures and the API never authors itself. @@ -275,19 +273,16 @@ def _log_api_authored_error(*, document: dict[str, Any], status: int, request: R uses (see ``_emit_at_error_level``), so a sink dedup'ing by level sees one shape. """ - error_domain = document.get("error_domain") fields: dict[str, Any] = { - "event": "api_error", - "request_id": request_id, - "route": request.url.path, - **_request_correlation_fields(request), + "event": API_ERROR_EVENT, + **_request_fields(request), "error_type": document.get("error_type"), - "error_domain": error_domain, + "error_domain": document.get("error_domain"), "retryable": document.get("retryable"), "status": status, "detail": document.get("detail"), } - emit_error_log(fields=fields, as_error=_emit_at_error_level(status)) + _emit_api_error(fields=fields, as_error=_emit_at_error_level(status)) def _json_safe_report(report: ErrorReport) -> ErrorReport: @@ -376,7 +371,7 @@ def _problem_response(report: ErrorReport, *, request: Request, disclosure_mode: mode into the handlers it registers). """ report = _json_safe_report(report) - request_id = _request_id_of(request) + request_id = request_id_of(request) status = _http_status_for(report) document = build_problem_document( report, @@ -388,7 +383,7 @@ def _problem_response(report: ErrorReport, *, request: Request, disclosure_mode: # when an API-layer override has bumped it away from ``report.http_status``; # otherwise the two surfaces (header vs body) would silently disagree. document["status"] = status - _log_error_report(report, request=request, request_id=request_id, status=status) + _log_error_report(report, request=request, status=status) return JSONResponse( status_code=status, content=document, @@ -481,13 +476,11 @@ async def handle_unexpected_error(request: Request, exc: Exception) -> Response: `to_error_report()`): Starlette's `ServerErrorMiddleware` wraps the others, so such a failure still lands here rather than a bodyless default. """ - request_id = _request_id_of(request) - emit_error_log( + request_id = request_id_of(request) + _emit_api_error( fields={ - "event": "api_error", - "request_id": request_id, - "route": request.url.path, - **_request_correlation_fields(request), + "event": API_ERROR_EVENT, + **_request_fields(request), "error_type": type(exc).__name__, "error_category": "unknown", "error_domain": ErrorDomain.RUNTIME, @@ -522,25 +515,27 @@ async def handle_api_error(request: Request, exc: Exception) -> Response: `ApiError` is raised by the `raise_*` helpers in `api.errors` for the API's own 4xx/5xx — request validation, auth, payload limits, misconfiguration. - The problem document is built at raise time (under the request-scoped - logging contextvars); this handler serializes it, re-attaches any - `WWW-Authenticate` challenge header, and emits the structured `event=api_error` - log line so an API-owned 500 from any route is observable (a - `raise_internal_server_error` site like `version.py`'s missing-package case - is the canonical example — it has no preceding `log.error` at the call - site). `exc` is typed `Exception` to match Starlette's handler contract; - FastAPI only routes an `ApiError` here, so the cast is sound. + The helper builds the problem document at raise time, where it holds no + `Request`, so this handler stamps the two request-scoped members (`instance` + and `request_id`) onto a copy of it — the same two values, read from the same + place, as the pipelex-error and request-validation paths use. It then + serializes that document, re-attaches any `WWW-Authenticate` challenge header, + and emits the structured `event=api_error` record so an API-owned 500 from any + route is observable (a `raise_internal_server_error` site like `version.py`'s + missing-package case is the canonical example — it has no preceding + `log.error` at the call site). `exc` is typed `Exception` to match Starlette's + handler contract; FastAPI only routes an `ApiError` here, so the cast is sound. """ api_error = cast("ApiError", exc) - _log_api_authored_error( - document=api_error.document, - status=api_error.status_code, - request=request, - request_id=_request_id_of(request), + document = with_request_context( + api_error.document, + instance=request.url.path, + request_id=request_id_of(request), ) + _log_api_authored_error(document=document, status=api_error.status_code, request=request) return JSONResponse( status_code=api_error.status_code, - content=api_error.document, + content=document, media_type=PROBLEM_JSON_MEDIA_TYPE, headers=api_error.headers, ) @@ -605,18 +600,17 @@ async def handle_request_validation_error(request: Request, exc: Exception) -> R is sound. """ validation_error = cast("RequestValidationError", exc) - request_id = _request_id_of(request) document = build_problem_document_from_api_error( _request_validation_error_type(validation_error), _summarize_request_validation_error(validation_error), 422, instance=request.url.path, - request_id=request_id, + request_id=request_id_of(request), error_domain=ErrorDomain.INPUT, ) - # Same `event=api_error` line as an explicit `raise_validation_error`, + # Same `event=api_error` record as an explicit `raise_validation_error`, # so FastAPI's automatic-validation 422s aren't silent in operator logs. - _log_api_authored_error(document=document, status=422, request=request, request_id=request_id) + _log_api_authored_error(document=document, status=422, request=request) return JSONResponse(status_code=422, content=document, media_type=PROBLEM_JSON_MEDIA_TYPE) diff --git a/api/logging_context.py b/api/logging_context.py deleted file mode 100644 index afa18b6..0000000 --- a/api/logging_context.py +++ /dev/null @@ -1,45 +0,0 @@ -"""Per-request logging context — request id and route path as contextvars. - -`api.middleware.RequestIdMiddleware` binds both values at request entry, for -the duration of the request. Downstream code reads them through the getters -below without threading a `Request` object through its signatures — the global -exception handlers, the 4xx helpers in `api.errors`, and structured-log call -sites all rely on this. - -Both contextvars default to `None`, so the getters are safe to call outside a -request scope (a CLI import of an `api` module, a unit test that issues no -request). -""" - -import contextvars -from collections.abc import Generator -from contextlib import contextmanager - -_request_id_ctxvar: contextvars.ContextVar[str | None] = contextvars.ContextVar("pipelex_api_request_id", default=None) -_route_path_ctxvar: contextvars.ContextVar[str | None] = contextvars.ContextVar("pipelex_api_route_path", default=None) - - -def get_request_id() -> str | None: - """Return the current request's correlation id, or `None` outside a request scope.""" - return _request_id_ctxvar.get() - - -def get_route_path() -> str | None: - """Return the current request's URL path, or `None` outside a request scope.""" - return _route_path_ctxvar.get() - - -@contextmanager -def bound_request_context(*, request_id: str, route_path: str) -> Generator[None]: - """Bind the request-scoped logging contextvars for the duration of the `with` block. - - Resets both contextvars to their prior state on exit — including when the - wrapped request raises — so a value never leaks into an unrelated context. - """ - request_id_token = _request_id_ctxvar.set(request_id) - route_path_token = _route_path_ctxvar.set(route_path) - try: - yield - finally: - _route_path_ctxvar.reset(route_path_token) - _request_id_ctxvar.reset(request_id_token) diff --git a/api/middleware.py b/api/middleware.py index bcb20ed..d862f24 100644 --- a/api/middleware.py +++ b/api/middleware.py @@ -7,31 +7,43 @@ from fastapi import Request, Response from fastapi.responses import JSONResponse +from pipelex import log from starlette.datastructures import MutableHeaders from starlette.types import ASGIApp, Message, Receive, Scope, Send from api.error_types import ErrorType from api.limits import MAX_REQUEST_BODY_BYTES, MAX_REQUEST_BODY_MIB -from api.logging_context import bound_request_context, get_request_id, get_route_path from api.problem_document import PROBLEM_JSON_MEDIA_TYPE, build_problem_document_from_api_error -def _too_large_response() -> JSONResponse: +def request_id_of(request: Request) -> str | None: + """Return the correlation id `RequestIdMiddleware` stored on this request, or `None` when it did not run. + + `getattr` rather than attribute access: a request that never went through the middleware — a unit + test issuing a bare ASGI call, a non-HTTP scope — simply has no id, and a reader of the id is not + the place to discover that. This is the one way the API reads the id back; the runtime's bound log + context carries the same value onto every record but is not a lookup table for anyone else. + """ + return getattr(request.state, "request_id", None) + + +def _too_large_response(*, request: Request) -> JSONResponse: """Build the 413 RFC 7807 problem response for an over-limit request body. The body-size check runs in middleware, before routing, so it cannot go through the `api.errors` helpers — a middleware must `return` a response, - not raise. It builds the same problem document directly. - `RequestIdMiddleware` runs outermost, so the request-scoped contextvars are - already bound and feed `instance` / `request_id`; that middleware's `send` - wrapper also stamps the `X-Request-ID` header onto this response. + not raise. It builds the same problem document directly, reading the request + context off the `Request` it was handed. `RequestIdMiddleware` runs outermost, + so the id is already on `request.state` by the time this is reached; that + middleware's `send` wrapper also stamps the `X-Request-ID` header onto this + response. """ document = build_problem_document_from_api_error( ErrorType.PAYLOAD_TOO_LARGE, f"Request body exceeds {MAX_REQUEST_BODY_MIB} MiB limit", 413, - instance=get_route_path(), - request_id=get_request_id(), + instance=request.url.path, + request_id=request_id_of(request), ) return JSONResponse(status_code=413, content=document, media_type=PROBLEM_JSON_MEDIA_TYPE) @@ -57,7 +69,7 @@ async def request_body_size_middleware(request: Request, call_next: Callable[[Re except ValueError: declared = -1 if declared > MAX_REQUEST_BODY_BYTES: - return _too_large_response() + return _too_large_response(request=request) original_receive = request.receive bytes_seen = 0 @@ -86,7 +98,7 @@ async def counting_receive() -> Message: response = await call_next(request) if too_large: - return _too_large_response() + return _too_large_response(request=request) return response @@ -149,17 +161,23 @@ class RequestIdMiddleware: """Pure-ASGI middleware that assigns a correlation id to every HTTP request. For each request it reuses a valid inbound `X-Request-ID` or mints a fresh - ULID, stores it on `request.state.request_id`, binds the `request_id` / - `route_path` logging contextvars for the duration of the request, and - echoes `X-Request-ID` on the response (success and error alike). + ULID, stores it on `request.state.request_id`, binds the runtime's log + context to that id for the duration of the request, and echoes + `X-Request-ID` on the response (success and error alike). + + Binding the runtime's context rather than an API-owned contextvar is what + puts `request_id` on *every* record emitted underneath — the API's own + error lines, and equally the ones pipelex emits from inside a run — as an + attribute a structured sink indexes, with no call site having to pass it + and no message having to interpolate it. Applied in `api.main` by wrapping the whole FastAPI app (`app = RequestIdMiddleware(app)`), NOT via `app.add_middleware()`. `add_middleware` always nests a middleware *inside* Starlette's - `ServerErrorMiddleware`, which would leave it unable to bind the contextvars + `ServerErrorMiddleware`, which would leave it unable to bind the context for — or set a header on — the catch-all 500 that `ServerErrorMiddleware` emits. Wrapping the app puts this middleware genuinely outermost, outside - `ServerErrorMiddleware`, so the contextvars are bound and `X-Request-ID` is + `ServerErrorMiddleware`, so the context is bound and `X-Request-ID` is echoed on every response, the catch-all 500 included. Raw ASGI (rather than `BaseHTTPMiddleware`) keeps a single contextvar context across the whole stack and lets the `send` wrapper inject the header on any response. @@ -181,5 +199,9 @@ async def send_with_request_id(message: Message) -> None: MutableHeaders(scope=message)[REQUEST_ID_HEADER] = request_id await send(message) - with bound_request_context(request_id=request_id, route_path=scope.get("path", "")): + # The runtime's own context, not an API-owned one: `request_id` is one of the three run + # identifiers it reserves, so every record emitted underneath carries it as an attribute. + # The route path is deliberately not bound here — it is not a run identifier, and the API + # ships it as a `route` field on the lines that want it (see `api.exception_handlers`). + with log.context(request_id=request_id): await self.app(scope, receive, send_with_request_id) diff --git a/api/problem_document.py b/api/problem_document.py index 8a2744f..49d5789 100644 --- a/api/problem_document.py +++ b/api/problem_document.py @@ -115,3 +115,20 @@ def build_problem_document_from_api_error( document["error_domain"] = error_domain document["retryable"] = retryable return document + + +def with_request_context(document: dict[str, Any], *, instance: str | None, request_id: str | None) -> dict[str, Any]: + """Return a copy of a problem document carrying the two request-scoped members. + + The `api.errors` helpers build their document where no `Request` is in hand — a validation + check deep inside a route, an auth refusal — so `instance` and `request_id` are stamped later, + by the handler that renders the response and does hold the request. A copy rather than an + in-place edit, because the `ApiError` a caller catches and inspects must read exactly as it was + raised. `None` is dropped rather than emitted as `null`, the rule every builder here follows. + """ + stamped = dict(document) + if instance is not None: + stamped["instance"] = instance + if request_id is not None: + stamped["request_id"] = request_id + return stamped diff --git a/api/routes/pipelex/pipeline.py b/api/routes/pipelex/pipeline.py index bdfb349..8717c78 100644 --- a/api/routes/pipelex/pipeline.py +++ b/api/routes/pipelex/pipeline.py @@ -31,8 +31,8 @@ from api.bundle import ParsedBundle, materialize_parsed, parse_bundle from api.error_types import ErrorType from api.errors import raise_bad_request, raise_forbidden, raise_validation_error -from api.logging_context import get_request_id from api.method_source import fetched_method_source +from api.middleware import request_id_of from api.openapi_responses import ( PROBLEM_400_START_REQUIRES_ASYNC, PROBLEM_403_RUN_POLICY, @@ -543,7 +543,7 @@ def _validate_extras(request_data: dict[str, Any]) -> PipelineApiExtras: # Per-field bound applied at the request.state binding site so an oversized -# caller-supplied `pipe_code` cannot blow up downstream log-line size. +# caller-supplied `pipe_code` cannot blow up the size of every record the request emits. # `RunRequest.pipe_code` carries no Pydantic `max_length`; this is the # narrow cap that protects the structured error log without changing the # upstream type contract. 256 covers any realistic pipe code (kebab-case @@ -556,12 +556,11 @@ def _coerce_correlation_field(value: Any) -> str | None: Returns `None` when the value is missing, empty, or non-string — so the handler's `_pipe_code_of` / `_pipeline_run_id_of` getters see a uniform - `None` and `emit_error_log` drops the field rather than rendering a bare - `pipe_code=` token (the empty-string case would otherwise pass the - `is not None` filter in `emit_error_log` and look like a logfmt parse - error to downstream sinks). Truncates oversized strings to - `_MAX_CORRELATION_FIELD_LEN` so a caller cannot inflate every error log - line for the request by sending a megabyte-long pipe_code. + `None` and the error record carries no attribute at all, rather than an + empty string a downstream query would read as a real value. Truncates + oversized strings to `_MAX_CORRELATION_FIELD_LEN` so a caller cannot + inflate every error record the request emits by sending a megabyte-long + pipe_code. """ if not isinstance(value, str) or not value: return None @@ -838,7 +837,7 @@ async def start( dynamic_output_concept_ref=run_request.dynamic_output_concept_ref, pipeline_run_id=extras.pipeline_run_id, callback_urls=extras.callback_urls, - request_id=get_request_id(), + request_id=request_id_of(request), requested_orchestration_mode=extras.orchestration_mode, ) # The ack plus this server's provenance extension: `(address, tag, commit_sha)` for a diff --git a/docs/configuration.md b/docs/configuration.md index 7f2eaa1..dc3d120 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -200,6 +200,19 @@ services: Use this when you want full control — for example, to ship your own inference backends, model deck, or routing profiles. You're now responsible for keeping the contents in sync with the version of Pipelex inside the image (the image's bundled `inference/`, `pipelex.toml`, etc. are no longer present at runtime). +!!! warning "Your `pipelex.toml` must re-supply the logging keys" + + The shipped `pipelex.toml` is what puts this server on structured logs, and replacing the directory replaces it. Without these three keys Pipelex falls back to its own defaults, which are written for someone at a terminal: the Rich console sink instead of JSON lines, and the "Output of pipe" panel rendered on the thread serving each request. Copy them into your own file: + + ```toml + [runtime.log] + sink = "json" + console_log_target = "stderr" + pretty_print_mode = "silent" + ``` + + See [Logging](logging.md#configuration) for what each one does. Option 1 above does not have this problem — a layered override keeps the keys you did not mention. + ### `docker run` equivalent ```bash diff --git a/docs/error-responses.md b/docs/error-responses.md index 9ebc02c..f945e5e 100644 --- a/docs/error-responses.md +++ b/docs/error-responses.md @@ -170,7 +170,7 @@ API-authored errors (`ValidationError`, `BadRequest`, `Unauthenticated`, etc.) f ## Request correlation -Every response carries `X-Request-ID`. The middleware respects an inbound `X-Request-ID` header if present, otherwise generates a UUID. The same id rides through onto `JobMetadata.request_id`, so it correlates the inbound HTTP call with the API-side log line — and, on a distributed-execution flavor, with every orchestrator worker-side log record produced during the run. +Every response carries `X-Request-ID`. The middleware respects an inbound `X-Request-ID` header if present, otherwise mints one. It also binds that id onto the Pipelex runtime's log context for the duration of the request, so every record emitted underneath carries it as a `request_id` field — the server's own error lines and the runtime's lines from inside a run alike. The same id rides onto `JobMetadata.request_id`, so on a distributed-execution flavor it correlates with every orchestrator worker-side record too. What those lines look like and which fields they carry is in [Logging](logging.md). When opening an issue, include the `request_id` from the response (or response headers) and the timestamp. diff --git a/docs/logging.md b/docs/logging.md new file mode 100644 index 0000000..5444cf9 --- /dev/null +++ b/docs/logging.md @@ -0,0 +1,81 @@ +# Logging + +The server writes structured logs: one JSON object per line, on **stderr**, with every value the line carries sitting as a key of its own. That is the shape a log agent in front of a container ingests without a parser — the CloudWatch agent, the Google Cloud Logging agent, an OTLP collector — so a query filters on a field rather than matching a substring of a message. + +## What a line looks like + +An error response produces exactly one line. A caller mistake: + +```json +{"time": "2026-09-18T13:53:17.919Z", "severity": "WARNING", "logger": "api.exception_handlers", "message": "API error 422: InvalidModelCategory", "request_id": "01M2TCP02W04RZG6DTM8AR508C", "event": "api_error", "route": "/v1/models", "error_type": "InvalidModelCategory", "error_domain": "input", "retryable": false, "status": 422, "detail": "Invalid model category. Valid values: extract, img_gen, llm, search"} +``` + +A server fault looks the same at `ERROR`, and carries the traceback under an `exception` key. + +The `message` is a short, stable sentence built from the HTTP status and the error type — both of them values the server chose. Nothing a caller supplied ever reaches it: the caller-facing explanation rides the `detail` field instead, so a body crafted with newlines or quotes cannot break the line or forge a key. Two layers keep the `detail` value safe to write. Before any sink sees a record, the Pipelex runtime's redaction processor replaces a control character in a field's value with its printable escape, so a newline a caller sent reads as `\n` in `detail`, and replaces a credential it recognises, an `Authorization` header's token or an API key, with `[REDACTED]`. The sink then escapes what it writes, so a quote or an `=` stays inside its value. The processor is configured under `[runtime.log.redaction]` and is on by default. + +## The fields a line carries + +These keys come from the sink itself and are on every line written by the process, this server's and the Pipelex runtime's alike: + +| Key | What it is | +| --- | --- | +| `time` | When the record was created — ISO 8601, UTC, milliseconds, `Z` | +| `severity` | The level name: `WARNING`, `ERROR`, … | +| `logger` | The module that emitted it | +| `message` | The human-readable summary | +| `exception` | The traceback, when the record carries one | +| `trace_id`, `span_id`, `trace_flags` | The trace context, in lowercase hex, when the record was logged inside a span, which is the case for a line the runtime emits from inside a traced run | + +`request_id` comes from the request-scoped context the request-id middleware binds, so **every** record emitted while a request is in flight carries it. The example above is one of the server's own lines, but a line the Pipelex runtime emits from inside a pipeline run carries the same id, which is what ties the two together without any call site passing it along. The value is the one echoed in the response's `X-Request-ID` header and in the problem document's `request_id` member, so a caller reporting a failure hands you the key to its log lines. + +The remaining keys are what the error handlers attach to an `event: "api_error"` record: + +| Key | What it is | +| --- | --- | +| `event` | Always `api_error` on an error line — the one key a query filters this server's error stream on | +| `route` | The request's URL path | +| `status` | The HTTP status actually sent, including any API-level override | +| `error_type` | The class or `ErrorType` name the response reports | +| `error_domain` | `input`, `config`, `runtime`, … — who fixes it | +| `error_category` | A Pipelex classification, on the failures it classifies; `unknown` on the catch-all 500 | +| `retryable` | Whether a blind retry can help, when the failure says | +| `detail` | The operator-facing explanation, the same text the response body carries — on an API-authored failure only, see below | +| `user_id` | The authenticated caller, when auth bound one | +| `pipe_code` | The pipe the request named, on a run route whose body parsed | +| `pipeline_run_id` | The run the request named, on a run route whose body parsed | +| `provider`, `model`, `provider_status_code`, `provider_request_id` | The inference provider's own identifiers, on a failure that reached one | + +A key whose value is not set for this request is **absent** from the line rather than written as `null`, so a query filtering on presence gets an honest answer. + +`detail` is the one field whose absence follows the failure's origin rather than the request's shape, and it is worth knowing which way round. A failure this API authored itself — a validation error, an unknown model category — carries `detail` on the record. A failure that arrives as a Pipelex `ErrorReport`, which is most `5xx` and every domain error, does not: the response body still carries a `detail`, but the record does not, because the body's text has been through disclosure redaction and the cause has not. So a `4xx` from that path logs as `API error 422: SomeError` with no explanation and no traceback, and the response is where the explanation is. Build an operator query on `error_type` and `route`, which every line carries, rather than on `detail`. + +Every field in that table is a **record attribute**, which is a different thing from the message. A structured sink — `json` here, and the OTLP sink — writes them beside the message as keys. The Rich console sink renders the message alone and ignores them entirely, so a deployment that switches the sink to `"console"` does not get a more readable version of these lines: it gets `API error 500: PipelexConfigError` and nothing else. That is a reason to keep `sink = "json"` wherever the lines are being read by anything other than a person at a terminal. + +Disposition follows the HTTP status, not the error domain: a `4xx` is a caller mistake and logs at `WARNING` without a traceback; a `5xx` is a server fault and logs at `ERROR` with one. See [Error Responses](error-responses.md) for the response side of the same failure. + +## Configuration + +The keys live in `[runtime.log]` of the `.pipelex/pipelex.toml` this repository ships, which the image copies to `/root/.pipelex/`: + +```toml +[runtime.log] +default_log_level = "INFO" +sink = "json" +console_log_target = "stderr" +pretty_print_mode = "silent" + +[runtime.log.package_log_levels] +pipelex = "INFO" +``` + +- `sink = "json"` selects the one-object-per-line renderer. The alternative, `"console"`, is the Rich renderer meant for a terminal. This server does not ask for Pipelex's `cli` extra, but that does not put the renderer out of reach: `typer` and `instructor` are core Pipelex dependencies that require Rich unconditionally, so Rich is installed in the image and an `import pipelex` loads it. Selecting `"console"` here would therefore give you a working console sink, not a refusal. This key is what keeps the renderer unused, and it is the only thing that does. +- `console_log_target = "stderr"` keeps logs off the data channel. +- `pretty_print_mode = "silent"` suppresses the "Output of pipe" panel an operator pipe would otherwise draw: there is no terminal to draw into, and a server should not spend time rendering on the thread serving a request. +- `[runtime.log.package_log_levels]` raises or lowers one package's level independently of `default_log_level`. + +To change any of this for a deployment, mount a `pipelex_override.toml` into `/root/.pipelex/` with just the keys you want different — see [Configuration](configuration.md#providing-your-own-configuration-to-docker). Note that replacing the whole config directory, which that page documents as Option 2, drops the shipped `pipelex.toml` along with the three keys above and silently returns the server to Pipelex's terminal-facing defaults; a layered override does not. + +## Uvicorn's own lines + +Uvicorn logs through its own handlers, not through the sink above, so its startup banner and its access log are plain text rather than JSON: the banner on stderr, the access log on stdout. Pass `--no-access-log` to turn the access log off, which is what the hosted deployment does; configure uvicorn's `--log-config` if you want its lines structured too. diff --git a/docs/openapi/pipelex-api.openapi.yaml b/docs/openapi/pipelex-api.openapi.yaml index 0173fa9..06ab0fd 100644 --- a/docs/openapi/pipelex-api.openapi.yaml +++ b/docs/openapi/pipelex-api.openapi.yaml @@ -11,7 +11,7 @@ info: license: name: Elastic License 2.0 identifier: Elastic-2.0 - version: 0.27.5 + version: 0.28.0 paths: /health: get: diff --git a/mkdocs.yml b/mkdocs.yml index b600a9b..353d81f 100644 --- a/mkdocs.yml +++ b/mkdocs.yml @@ -109,6 +109,7 @@ nav: - MTHDS Tools: mthds-tools.md - Storage Transport: storage-transport.md - Error Responses: error-responses.md + - Logging: logging.md - Contribute: - Contributing: contributing.md - Code of Conduct: CODE_OF_CONDUCT.md diff --git a/pyproject.toml b/pyproject.toml index f99f3d1..efc00e0 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "pipelex-api" -version = "0.27.5" +version = "0.28.0" description = "Pipelex API" authors = [{ name = "Evotis S.A.S.", email = "oss@pipelex.com" }] maintainers = [{ name = "Pipelex staff", email = "oss@pipelex.com" }] @@ -16,7 +16,9 @@ classifiers = [ ] dependencies = [ - "pipelex[mistralai,anthropic,google,google-genai,bedrock,fal]==0.65.0", + # The extras deliberately leave out `cli`: the runner selects the `json` log sink and a Rich-free + # pretty-print mode, so nothing it does on a request renders a terminal. + "pipelex[mistralai,anthropic,google,google-genai,bedrock,fal]==0.66.0", "fastapi>=0.118.0", "pyjwt>=2.10.1", "uvicorn>=0.37.0", diff --git a/tests/unit/test_error_log_records.py b/tests/unit/test_error_log_records.py new file mode 100644 index 0000000..89a48ed --- /dev/null +++ b/tests/unit/test_error_log_records.py @@ -0,0 +1,178 @@ +"""What an error record actually carries, and what it looks like once the `json` sink has written it. + +The tests in `test_exception_handlers.py` read the `fields=` mapping the handlers hand the runtime, +with the runtime's `log` replaced by a spy. These ones let the real thing run end to end: a request +goes through `RequestIdMiddleware`, the handler emits, `caplog` catches the record the runtime built, +and the runner's configured sink formatter renders it. That is the whole path this server's operator +output takes in production, minus the process stream it is written to. +""" + +import json +import logging +from pathlib import Path +from typing import Any, cast + +import pytest +from fastapi import APIRouter, FastAPI +from fastapi.testclient import TestClient +from pipelex.base_exceptions import PipelexConfigError +from pipelex.system.console_target import ConsoleTarget +from pipelex.tools.log.json_log_sink import JsonLogFormatter +from pipelex.tools.log.log_sink import LogSinkMethod +from pipelex.tools.misc.pretty import PrettyPrintMode +from pipelex.tools.misc.toml_utils import load_toml_from_path + +from api.error_types import ErrorType +from api.errors import raise_validation_error +from api.exception_handlers import API_ERROR_EVENT, register_exception_handlers +from api.middleware import REQUEST_ID_HEADER, RequestIdMiddleware + +_SHIPPED_PIPELEX_CONFIG = Path(__file__).parents[2] / ".pipelex" / "pipelex.toml" + +_router = APIRouter() + + +@_router.get("/pipelex-failure") +async def pipelex_failure_route() -> None: + # An operator-actionable 500 — the `error`-level branch, with a traceback on the record. + msg = "the gateway config is missing" + raise PipelexConfigError(msg) + + +@_router.get("/caller-mistake") +async def caller_mistake_route(detail: str) -> None: + # A caller-input 422 whose `detail` is exactly what the caller sent — the `warning`-level + # branch, and the one field on the record that a caller controls. + raise_validation_error(detail, ErrorType.VALIDATION_ERROR) + + +def _build_client() -> TestClient: + """Wire a throwaway app with the production handlers and the request-id middleware.""" + app = FastAPI() + register_exception_handlers(app) + app.include_router(_router) + return TestClient(RequestIdMiddleware(app)) + + +def _api_error_record(caplog: pytest.LogCaptureFixture) -> logging.LogRecord: + """The one `api_error` record the request emitted.""" + records = [record for record in caplog.records if getattr(record, "event", None) == API_ERROR_EVENT] + assert len(records) == 1, f"expected exactly one api_error record, got {len(records)}" + return records[0] + + +def _rendered_json(record: logging.LogRecord) -> dict[str, Any]: + """The record as the `json` sink writes it: one line, parsed back.""" + line = JsonLogFormatter().format(record) + assert "\n" not in line, f"the sink wrote more than one line: {line!r}" + assert "\r" not in line, f"the sink wrote a carriage return: {line!r}" + payload = json.loads(line) + assert isinstance(payload, dict) + return cast("dict[str, Any]", payload) + + +class TestErrorLogRecords: + def test_the_runner_ships_a_configuration_that_selects_the_json_sink(self): + # Reads the shipped file rather than `get_config()`, deliberately. The runtime layers + # `_local`, `_{environment}` and `_override` files over a base, so the merged config is + # partly this machine's: a developer who sets `sink = "console"` locally for a readable + # `make run` would red this test while CI, which has no such file, stayed green. What the + # image boots on is this file, because nothing mounts an override into it. + # + # The json sink writes to stderr so stdout stays the data channel, and the pretty-print + # mode is silent because a server has no terminal and must not render on the request + # thread. Neither is a claim that Rich is absent: `typer` and `instructor` require it + # unconditionally, so it is installed and reachable — these keys are what keeps it unused. + log_section = load_toml_from_path(_SHIPPED_PIPELEX_CONFIG)["runtime"]["log"] + assert log_section["sink"] == LogSinkMethod.JSON + assert log_section["console_log_target"] == ConsoleTarget.STDERR + assert log_section["pretty_print_mode"] == PrettyPrintMode.SILENT + + def test_request_id_and_route_ride_the_error_record(self, caplog: pytest.LogCaptureFixture): + # The two identifiers an operator starts from. `request_id` arrives from the log context + # the middleware bound for the request; `route` from the field set the handler ships. + # Neither is interpolated into the message, which is the whole point of the change. + with caplog.at_level(logging.WARNING): + response = _build_client().get("/pipelex-failure") + assert response.status_code == 500 + record = _api_error_record(caplog) + assert record.levelno == logging.ERROR + assert getattr(record, "request_id", None) == response.headers[REQUEST_ID_HEADER] + assert getattr(record, "route", None) == "/pipelex-failure" + assert getattr(record, "error_type", None) == "PipelexConfigError" + assert record.getMessage() == "API error 500: PipelexConfigError" + + def test_a_caller_mistake_records_a_warning_carrying_the_same_identifiers(self, caplog: pytest.LogCaptureFixture): + # A 4xx is a caller mistake, so it lands at `warning` — but it carries the same two + # identifiers, so one query over `event` returns the whole error stream of a request. + with caplog.at_level(logging.WARNING): + response = _build_client().get("/caller-mistake", params={"detail": "the input is malformed"}) + assert response.status_code == 422 + record = _api_error_record(caplog) + assert record.levelno == logging.WARNING + assert getattr(record, "request_id", None) == response.headers[REQUEST_ID_HEADER] + assert getattr(record, "route", None) == "/caller-mistake" + assert getattr(record, "status", None) == 422 + + def test_the_json_sink_writes_one_object_per_line_with_the_fields_flat(self, caplog: pytest.LogCaptureFixture): + # The shape the runner's configured sink puts on stderr: the sink's own keys, then every + # field and the bound context identifiers flat beside them, one JSON object on one line. + with caplog.at_level(logging.WARNING): + response = _build_client().get("/pipelex-failure") + payload = _rendered_json(_api_error_record(caplog)) + assert payload["severity"] == "ERROR" + assert payload["message"] == "API error 500: PipelexConfigError" + assert payload["event"] == API_ERROR_EVENT + assert payload["request_id"] == response.headers[REQUEST_ID_HEADER] + assert payload["route"] == "/pipelex-failure" + assert payload["status"] == 500 + assert payload["error_type"] == "PipelexConfigError" + assert payload["error_domain"] == "config" + # `retryable` is absent rather than false: pipelex populates it only on a classifiable + # failure, and a field with no value is dropped rather than written as a null a query + # filtering on presence would read as an answer. + assert "retryable" not in payload + # `time` and `logger` come from the sink itself; an operator-actionable failure carries + # its traceback under `exception` rather than spilling it across the line. + assert payload["time"].endswith("Z") + assert "PipelexConfigError" in payload["exception"] + + @pytest.mark.parametrize( + ("crafted_detail", "logged_detail"), + [ + # Back when the API flattened its own fields into a `key=value` run, each of these + # forged either a sibling field or a whole second line. The API renders nothing now: + # the runtime escapes a control character in a field's value and the sink escapes what + # it writes, so each must survive as one field's value. + ("legit\nstatus=200 event=auth_success", "legit\\nstatus=200 event=auth_success"), + ("hijack status=200 event=fake", "hijack status=200 event=fake"), + ("legit\rstatus=200", "legit\\rstatus=200"), + ('has " a quote = inside', 'has " a quote = inside'), + ], + ) + def test_a_crafted_detail_survives_as_one_value_and_forges_nothing( + self, caplog: pytest.LogCaptureFixture, crafted_detail: str, logged_detail: str + ): + # The escaping the API used to do itself is the runtime's job now, in two layers. The + # redaction processor every sink sits behind turns a control character in a field's value + # into its printable escape, so a newline the caller sent reads `\n` in `detail`; then + # `json.dumps` writes each value inside one string, so a quote or an `=` stays part of it. + # Either way the line stays one object and the value comes back as one field. + with caplog.at_level(logging.WARNING): + response = _build_client().get("/caller-mistake", params={"detail": crafted_detail}) + assert response.status_code == 422 + payload = _rendered_json(_api_error_record(caplog)) + assert payload["detail"] == logged_detail + assert payload["event"] == API_ERROR_EVENT, "a crafted detail forged or overwrote a field" + assert payload["status"] == 422, "a crafted detail forged or overwrote a field" + assert crafted_detail not in payload["message"], "caller input reached the message" + + def test_a_credential_echoed_into_a_detail_is_redacted_on_the_line(self, caplog: pytest.LogCaptureFixture): + # A validation message can echo what the caller sent, header text included. The runtime's + # redaction processor scrubs the credential out of the field before the sink writes it, and + # keeps the scheme, so the line still says which kind of credential went. + with caplog.at_level(logging.WARNING): + response = _build_client().get("/caller-mistake", params={"detail": "rejected Authorization: Bearer placeholder-not-a-token"}) + assert response.status_code == 422 + payload = _rendered_json(_api_error_record(caplog)) + assert payload["detail"] == "rejected Authorization: Bearer [REDACTED]" diff --git a/tests/unit/test_error_responses.py b/tests/unit/test_error_responses.py index 64c4822..b87748d 100644 --- a/tests/unit/test_error_responses.py +++ b/tests/unit/test_error_responses.py @@ -136,6 +136,10 @@ def test_payload_too_large_is_rfc7807(self): assert body["error_type"] == "PayloadTooLarge" assert body["error_domain"] == "input" assert body["status"] == 413 + # The body-size middleware has to return a response rather than raise, so it builds its own + # problem document — and it reads the route and the id off the `Request` it was handed, + # like every other error path. + assert body["instance"] == "/v1/version" assert body["request_id"] == response.headers[REQUEST_ID_HEADER] assert body["retryable"] is False diff --git a/tests/unit/test_exception_handlers.py b/tests/unit/test_exception_handlers.py index 788bdc3..5978ab8 100644 --- a/tests/unit/test_exception_handlers.py +++ b/tests/unit/test_exception_handlers.py @@ -9,7 +9,7 @@ """ import re -from typing import Any +from typing import Any, cast import pytest from fastapi import APIRouter, FastAPI, Request @@ -36,7 +36,7 @@ from api.error_types import ErrorType from api.errors import raise_internal_server_error, raise_validation_error -from api.exception_handlers import emit_error_log, register_exception_handlers +from api.exception_handlers import API_ERROR_EVENT, register_exception_handlers from api.middleware import REQUEST_ID_HEADER, RequestIdMiddleware from api.problem_document import PROBLEM_JSON_MEDIA_TYPE from api.security import RequestUser @@ -45,48 +45,17 @@ _ULID_RE = re.compile(r"\A[0-9A-HJKMNP-TV-Z]{26}\Z") -def _parse_logfmt(line: str) -> dict[str, str]: - """Minimal logfmt parser used by the injection-guard test. +def _emitted_fields(log_spy: Any, *, as_error: bool) -> dict[str, Any]: + """The record attributes the handler shipped on its one `api_error` call. - Splits ``key=value`` pairs separated by spaces, honoring ``"..."`` quoting - with embedded ``""`` doubled (the same convention `_logfmt_value` writes). - Returns the top-level field map a downstream parser would extract — used - to assert that a crafted value did NOT introduce extra top-level keys. + The handlers hand everything but the summary message to the runtime as `fields=`, so the + assertions below read that mapping rather than parsing a rendered line: there is no rendering + left to parse, and a sink is what decides how an attribute reaches the wire. """ - fields: dict[str, str] = {} - index = 0 - length = len(line) - while index < length: - while index < length and line[index] == " ": - index += 1 - if index >= length: - break - equals_position = line.find("=", index) - if equals_position == -1: - break - key = line[index:equals_position] - index = equals_position + 1 - if index < length and line[index] == '"': - index += 1 - value_chars: list[str] = [] - while index < length: - if line[index] == '"': - if index + 1 < length and line[index + 1] == '"': - value_chars.append('"') - index += 2 - continue - index += 1 - break - value_chars.append(line[index]) - index += 1 - fields[key] = "".join(value_chars) - else: - value_end = line.find(" ", index) - if value_end == -1: - value_end = length - fields[key] = line[index:value_end] - index = value_end - return fields + call = log_spy.error.call_args if as_error else log_spy.warning.call_args + fields = call.kwargs["fields"] + assert isinstance(fields, dict) + return cast("dict[str, Any]", fields) class _SimulatedLLMError(PipelexError): @@ -312,6 +281,13 @@ async def api_input_error_route() -> None: raise_validation_error("a caller-side mistake") +@_router.get("/crafted-detail") +async def crafted_detail_route(detail: str) -> None: + # An API-authored 422 whose `detail` is exactly what the caller sent — the real shape of the + # `raise_validation_error(str(exc))` sites, where the message is built from caller input. + raise_validation_error(detail) + + @_router.get("/api-config-error") async def api_config_error_route() -> None: # An API-authored 500 — the CONFIG-domain branch of `handle_api_error` @@ -638,31 +614,31 @@ def test_inbound_request_id_is_echoed(self): assert response.headers[REQUEST_ID_HEADER] == "inbound-correlation-007" assert response.json()["request_id"] == "inbound-correlation-007" - def test_api_authored_500_emits_structured_error_log(self, mocker: MockerFixture): + def test_api_authored_500_emits_structured_error_record(self, mocker: MockerFixture): # Without `handle_api_error` logging, an `ApiError`-shaped 500 produces # zero operator output (`/version`'s `PackageNotFoundError → # raise_internal_server_error` is the canonical silent case). Assert - # the handler emits one `event=api_error` line at `error` level with - # the response fields — same shape `_log_error_report` emits for a - # pipelex-derived 500, so a downstream log sink sees them uniformly. + # the handler emits one `api_error` record at `error` level whose fields + # mirror the response — the same set `_log_error_report` ships for a + # pipelex-derived 500, so a downstream sink queries them uniformly. log_spy = mocker.patch("api.exception_handlers.log") response = _build_client().get("/api-config-error") assert response.status_code == 500 log_spy.error.assert_called_once() - rendered = log_spy.error.call_args.args[0] - assert "event=api_error" in rendered - assert "status=500" in rendered - assert "error_type=ServerMisconfigured" in rendered - assert "error_domain=config" in rendered - assert "retryable=False" in rendered - assert "route=/api-config-error" in rendered - # Values containing whitespace are logfmt-quoted by `_logfmt_value` - # so a crafted `detail` cannot forge sibling fields. - assert 'detail="the configuration is broken"' in rendered - assert log_spy.error.call_args.kwargs == {"include_exception": True} + fields = _emitted_fields(log_spy, as_error=True) + assert fields["event"] == API_ERROR_EVENT + assert fields["status"] == 500 + assert fields["error_type"] == "ServerMisconfigured" + assert fields["error_domain"] == "config" + assert fields["retryable"] is False + assert fields["route"] == "/api-config-error" + # The operator-facing cause rides a field of its own and never reaches the message, which + # is what stops a crafted value from shaping the rendered line at all. + assert fields["detail"] == "the configuration is broken" + assert log_spy.error.call_args.kwargs["include_exception"] is True log_spy.warning.assert_not_called() - def test_api_authored_4xx_emits_structured_warning_log(self, mocker: MockerFixture): + def test_api_authored_4xx_emits_structured_warning_record(self, mocker: MockerFixture): # Mirror at the warning level: an INPUT-domain `ApiError` is a caller # mistake, not an operator fault, so it logs at `warning` without a # traceback — same disposition rule `_log_error_report` uses. @@ -670,30 +646,37 @@ def test_api_authored_4xx_emits_structured_warning_log(self, mocker: MockerFixtu response = _build_client().get("/api-input-error") assert response.status_code == 422 log_spy.warning.assert_called_once() - rendered = log_spy.warning.call_args.args[0] - assert "event=api_error" in rendered - assert "status=422" in rendered - assert "error_type=ValidationError" in rendered - assert "error_domain=input" in rendered - assert "retryable=False" in rendered - assert 'detail="a caller-side mistake"' in rendered + fields = _emitted_fields(log_spy, as_error=False) + assert fields["event"] == API_ERROR_EVENT + assert fields["status"] == 422 + assert fields["error_type"] == "ValidationError" + assert fields["error_domain"] == "input" + assert fields["retryable"] is False + assert fields["detail"] == "a caller-side mistake" log_spy.error.assert_not_called() - def test_user_id_rides_authenticated_pipelex_error_log(self, mocker: MockerFixture): + def test_the_summary_message_names_the_status_and_the_error_type(self, mocker: MockerFixture): + # The message is built from server-authored values alone, so it stays a stable sentence + # whatever a caller sent; everything variable is a field beside it. + log_spy = mocker.patch("api.exception_handlers.log") + assert _build_client().get("/api-config-error").status_code == 500 + assert log_spy.error.call_args.args[0] == "API error 500: ServerMisconfigured" + + def test_user_id_rides_authenticated_pipelex_error_record(self, mocker: MockerFixture): # Phase 3 deleted the per-route `log.error(... user=...)` lines on # storage / pipeline-backend failures; the global handler now reads # `request.state.user` so the operator can still tie a `PipelexError` - # to the caller via the structured log alone — no cross-line grep on - # `request_id` required. + # to the caller through the record alone — no correlating across + # records on `request_id` required. log_spy = mocker.patch("api.exception_handlers.log") response = _build_client().get("/authenticated-config-error") assert response.status_code == 500 log_spy.error.assert_called_once() - rendered = log_spy.error.call_args.args[0] - assert f"user_id={_TEST_USER_ID}" in rendered - assert "error_type=PipelexConfigError" in rendered + fields = _emitted_fields(log_spy, as_error=True) + assert fields["user_id"] == _TEST_USER_ID + assert fields["error_type"] == "PipelexConfigError" - def test_user_id_rides_authenticated_api_authored_log(self, mocker: MockerFixture): + def test_user_id_rides_authenticated_api_authored_record(self, mocker: MockerFixture): # `handle_api_error` covers the API-authored 4xx surface (storage's # `raise_bad_request`, `raise_forbidden`, `raise_payload_too_large`; # uploader's same set). It must ride the same `user_id` correlation @@ -703,11 +686,11 @@ def test_user_id_rides_authenticated_api_authored_log(self, mocker: MockerFixtur response = _build_client().get("/authenticated-api-input-error") assert response.status_code == 422 log_spy.warning.assert_called_once() - rendered = log_spy.warning.call_args.args[0] - assert f"user_id={_TEST_USER_ID}" in rendered - assert "error_type=ValidationError" in rendered + fields = _emitted_fields(log_spy, as_error=False) + assert fields["user_id"] == _TEST_USER_ID + assert fields["error_type"] == "ValidationError" - def test_user_id_rides_authenticated_unexpected_error_log(self, mocker: MockerFixture): + def test_user_id_rides_authenticated_unexpected_error_record(self, mocker: MockerFixture): # The catch-all 500 (`handle_unexpected_error`) is the one place a # missing `user_id` is most expensive — by definition the failure # was not classifiable upstream — so the same enrichment fires here. @@ -715,23 +698,22 @@ def test_user_id_rides_authenticated_unexpected_error_log(self, mocker: MockerFi response = _build_client(raise_server_exceptions=False).get("/authenticated-unexpected-error") assert response.status_code == 500 log_spy.error.assert_called_once() - rendered = log_spy.error.call_args.args[0] - assert f"user_id={_TEST_USER_ID}" in rendered - assert "error_type=RuntimeError" in rendered + fields = _emitted_fields(log_spy, as_error=True) + assert fields["user_id"] == _TEST_USER_ID + assert fields["error_type"] == "RuntimeError" - def test_user_id_absent_from_unauthenticated_error_log(self, mocker: MockerFixture): + def test_user_id_absent_from_unauthenticated_error_record(self, mocker: MockerFixture): # Pre-auth paths, and a deployment with no user model, have no `request.state.user`; - # `_user_id_of` returns `None` and `emit_error_log` drops `None`- - # valued fields, so the rendered line carries no `user_id=` token — - # never `user_id=None`, which would be misleading noise. + # `_user_id_of` returns `None` and `_emit_api_error` drops `None`-valued fields, so the + # record carries no `user_id` attribute at all — never a null, which a query filtering on + # presence would read as an answer. log_spy = mocker.patch("api.exception_handlers.log") response = _build_client().get("/config-error") assert response.status_code == 500 log_spy.error.assert_called_once() - rendered = log_spy.error.call_args.args[0] - assert "user_id=" not in rendered + assert "user_id" not in _emitted_fields(log_spy, as_error=True) - def test_run_state_rides_pipelex_error_log(self, mocker: MockerFixture): + def test_run_state_rides_pipelex_error_record(self, mocker: MockerFixture): # `_parse_request` binds `pipe_code` / `pipeline_run_id` on # `request.state` so a downstream pipelex failure is tied to the # specific pipe and run id without each backend frame having to forward @@ -742,62 +724,59 @@ def test_run_state_rides_pipelex_error_log(self, mocker: MockerFixture): response = _build_client().get("/pipeline-state-pipelex-error") assert response.status_code == 500 log_spy.error.assert_called_once() - rendered = log_spy.error.call_args.args[0] - assert f"pipe_code={_TEST_PIPE_CODE}" in rendered - assert f"pipeline_run_id={_TEST_RUN_ID}" in rendered - assert "error_type=PipelexConfigError" in rendered + fields = _emitted_fields(log_spy, as_error=True) + assert fields["pipe_code"] == _TEST_PIPE_CODE + assert fields["pipeline_run_id"] == _TEST_RUN_ID + assert fields["error_type"] == "PipelexConfigError" - def test_run_state_rides_api_authored_log(self, mocker: MockerFixture): + def test_run_state_rides_api_authored_record(self, mocker: MockerFixture): # API-authored 4xx surface: `raise_validation_error` raised from a # route that has already bound pipe state must ship the same fields, - # so the operator log is uniform across the validation and the + # so the operator record is uniform across the validation and the # pipelex-domain failure surfaces. log_spy = mocker.patch("api.exception_handlers.log") response = _build_client().get("/pipeline-state-api-input-error") assert response.status_code == 422 log_spy.warning.assert_called_once() - rendered = log_spy.warning.call_args.args[0] - assert f"pipe_code={_TEST_PIPE_CODE}" in rendered - assert f"pipeline_run_id={_TEST_RUN_ID}" in rendered - assert "error_type=ValidationError" in rendered + fields = _emitted_fields(log_spy, as_error=False) + assert fields["pipe_code"] == _TEST_PIPE_CODE + assert fields["pipeline_run_id"] == _TEST_RUN_ID + assert fields["error_type"] == "ValidationError" - def test_run_state_rides_unexpected_error_log(self, mocker: MockerFixture): + def test_run_state_rides_unexpected_error_record(self, mocker: MockerFixture): # The catch-all 500 is the most operationally expensive case for a # missing pipe-state field — by definition the failure was not - # classifiable upstream, so the identifiers in the log are the + # classifiable upstream, so the identifiers on the record are the # operator's only starting point for which pipe and which run were # in flight. log_spy = mocker.patch("api.exception_handlers.log") response = _build_client(raise_server_exceptions=False).get("/pipeline-state-unexpected-error") assert response.status_code == 500 log_spy.error.assert_called_once() - rendered = log_spy.error.call_args.args[0] - assert f"pipe_code={_TEST_PIPE_CODE}" in rendered - assert f"pipeline_run_id={_TEST_RUN_ID}" in rendered - assert "error_type=RuntimeError" in rendered + fields = _emitted_fields(log_spy, as_error=True) + assert fields["pipe_code"] == _TEST_PIPE_CODE + assert fields["pipeline_run_id"] == _TEST_RUN_ID + assert fields["error_type"] == "RuntimeError" def test_run_state_absent_when_parse_request_did_not_bind(self, mocker: MockerFixture): # A route that didn't go through `_parse_request` (here: `/config-error`, # a GET that never touches the body) has no `pipe_code` / # `pipeline_run_id` on `request.state`. The getters return `None` and - # `emit_error_log` drops `None`-valued fields, so the rendered line - # carries no `pipe_code=` or `pipeline_run_id=` token — never - # `pipe_code=None`, which would be misleading noise. Same defensive - # posture as `test_user_id_absent_from_unauthenticated_error_log`. + # `_emit_api_error` drops those, so the record carries neither attribute. + # Same defensive posture as `test_user_id_absent_from_unauthenticated_error_record`. log_spy = mocker.patch("api.exception_handlers.log") response = _build_client().get("/config-error") assert response.status_code == 500 log_spy.error.assert_called_once() - rendered = log_spy.error.call_args.args[0] - assert "pipe_code=" not in rendered - assert "pipeline_run_id=" not in rendered + fields = _emitted_fields(log_spy, as_error=True) + assert "pipe_code" not in fields + assert "pipeline_run_id" not in fields - def test_request_validation_error_emits_structured_warning_log(self, mocker: MockerFixture): + def test_request_validation_error_emits_structured_warning_record(self, mocker: MockerFixture): # FastAPI's automatic-validation 422 goes through # `handle_request_validation_error`, not `handle_api_error`, but emits - # the same `event=api_error` warning line so the surface is uniform — - # whichever code path rejects a caller-input failure, the operator - # log shape is identical. + # the same `api_error` record so the surface is uniform — whichever code + # path rejects a caller-input failure, the operator record is identical. log_spy = mocker.patch("api.exception_handlers.log") # `{"wrong": 1}` triggers FastAPI's automatic validation failure on # `_RequestValidationBody`: `field` is missing AND `wrong` is an @@ -807,59 +786,48 @@ def test_request_validation_error_emits_structured_warning_log(self, mocker: Moc response = _build_client().post("/needs-body", json={"wrong": 1}) assert response.status_code == 422 log_spy.warning.assert_called_once() - rendered = log_spy.warning.call_args.args[0] - assert "event=api_error" in rendered - assert "status=422" in rendered - assert "error_type=ValidationError" in rendered - assert "error_domain=input" in rendered + fields = _emitted_fields(log_spy, as_error=False) + assert fields["event"] == API_ERROR_EVENT + assert fields["status"] == 422 + assert fields["error_type"] == "ValidationError" + assert fields["error_domain"] == "input" # The summary covers both per-field failures the request triggered. - assert "field" in rendered - assert "wrong" in rendered + assert "field" in fields["detail"] + assert "wrong" in fields["detail"] log_spy.error.assert_not_called() @pytest.mark.parametrize( "crafted_detail", [ - # A newline in the detail would forge a second log line entirely — - # the worst case for `\n`-delimited log shippers (Loki, journald, - # CloudWatch). `_logfmt_value` encodes it as `\\n` so the line - # stays single-line. + # A newline in the detail used to forge a second log line entirely — the worst case + # for newline-delimited shippers (Loki, journald, CloudWatch) back when the API + # flattened its own fields into one `key=value` run. "legit\nstatus=200 event=auth_success", - # Whitespace + key=value runs are the canonical logfmt forge: a - # crafted detail must not produce a top-level `event` field that - # logfmt-aware parsers would treat as a real field. + # Whitespace plus `key=value` runs were the canonical logfmt forge. "hijack status=200 event=fake", - # A carriage return alone is enough to corrupt a `\r\n` shipper. + # A carriage return alone was enough to corrupt a CRLF shipper. "legit\rstatus=200", - # A bare `"` inside the value would close the quoted region early - # in a permissive parser; the escape doubles it. + # A bare quote used to close the quoted region early in a permissive parser. 'has " a quote = inside', ], ) - def test_log_injection_in_detail_is_neutralized(self, mocker: MockerFixture, crafted_detail: str) -> None: - # `_log_api_authored_error` ships `document["detail"]` into the log - # line, and `detail` originates from caller input on several routes - # (`raise_validation_error(str(exc))`, callback-URL rejections, - # `_summarize_request_validation_error`'s per-field messages). Without - # escaping, a crafted body could forge sibling log fields or whole - # new log lines (Greptile P1 finding). Pin the neutralization - # contract directly at the formatter so a future caller that ships a - # new caller-controlled field is also covered without re-auditing. + def test_crafted_detail_stays_one_field_and_shapes_no_other(self, mocker: MockerFixture, crafted_detail: str) -> None: + # `_log_api_authored_error` ships `document["detail"]`, and `detail` originates from caller + # input on several routes (`raise_validation_error(str(exc))`, callback-URL rejections, + # `_summarize_request_validation_error`'s per-field messages). The API renders none of it + # any more: the value is handed to the runtime as one field, so it cannot forge a sibling + # field and cannot reach the message, and the selected sink is what escapes it on the way + # out. Pin that at the handler, so a route shipping some new caller-controlled field later + # is covered without re-auditing this surface. log_spy = mocker.patch("api.exception_handlers.log") - emit_error_log( - fields={"event": "api_error", "detail": crafted_detail, "status": 422}, - as_error=False, - ) - rendered = log_spy.warning.call_args.args[0] - # Line stays single-line — no shipper-delimiter forgery. - assert "\n" not in rendered, f"newline survived escaping: rendered={rendered!r}" - assert "\r" not in rendered, f"carriage return survived escaping: rendered={rendered!r}" - # The top-level field set, as a logfmt-aware parser would read it, - # is exactly the three keys we shipped — no forged siblings. - parsed = _parse_logfmt(rendered) - assert set(parsed.keys()) == {"event", "detail", "status"}, f"forged field surfaced: rendered={rendered!r}, parsed={parsed!r}" - assert parsed["event"] == "api_error", f"legitimate `event` field was overwritten: rendered={rendered!r}" - assert parsed["status"] == "422", f"legitimate `status` field was overwritten: rendered={rendered!r}" + response = _build_client().get("/crafted-detail", params={"detail": crafted_detail}) + assert response.status_code == 422 + fields = _emitted_fields(log_spy, as_error=False) + assert fields["detail"] == crafted_detail, f"the value was altered on the way to the record: {fields['detail']!r}" + assert fields["event"] == API_ERROR_EVENT, "a crafted detail overwrote a legitimate field" + assert fields["status"] == 422, "a crafted detail overwrote a legitimate field" + message = log_spy.warning.call_args.args[0] + assert crafted_detail not in message, f"caller input reached the message: {message!r}" def test_async_execution_not_enabled_maps_to_501(self): # The pipelex-side ``AsyncExecutionNotEnabledError`` reports as @@ -885,19 +853,19 @@ def test_async_execution_not_enabled_maps_to_501(self): def test_async_execution_not_enabled_logs_post_override_status(self, mocker: MockerFixture): # ``_log_error_report`` receives the overridden status so the operator - # log line agrees with what the client actually saw. Without the - # ``status=`` plumbing the log would read ``status=500`` (the report's - # domain default) while the response shipped 501 — a silent disagree + # record agrees with what the client actually saw. Without the + # ``status=`` plumbing the record would carry 500 (the report's domain + # default) while the response shipped 501 — a silent disagreement # between the two surfaces. Pin the alignment. log_spy = mocker.patch("api.exception_handlers.log") response = _build_client().get("/async-execution-not-enabled") assert response.status_code == 501 log_spy.error.assert_called_once() - rendered = log_spy.error.call_args.args[0] - assert "event=api_error" in rendered - assert "status=501" in rendered - assert "error_type=AsyncExecutionNotEnabledError" in rendered - assert "error_domain=config" in rendered + fields = _emitted_fields(log_spy, as_error=True) + assert fields["event"] == API_ERROR_EVENT + assert fields["status"] == 501 + assert fields["error_type"] == "AsyncExecutionNotEnabledError" + assert fields["error_domain"] == "config" def test_pipeline_run_id_conflict_maps_to_409(self): # ``PipelineManagerAlreadyExistsError`` carries no ``error_domain`` @@ -923,14 +891,14 @@ def test_pipeline_run_id_conflict_logs_at_warning_post_override_status(self, moc # A duplicate pipeline_run_id is a client-visible 409 conflict, not a # server fault: disposition is keyed off the final HTTP status, so it # logs at `warning` without a traceback — never `error` (which would - # page or pollute error dashboards for a normal conflict). The line + # page or pollute error dashboards for a normal conflict). The record # still agrees with the status the client saw. log_spy = mocker.patch("api.exception_handlers.log") response = _build_client().get("/pipeline-run-id-conflict") assert response.status_code == 409 log_spy.warning.assert_called_once() log_spy.error.assert_not_called() - rendered = log_spy.warning.call_args.args[0] - assert "event=api_error" in rendered - assert "status=409" in rendered - assert "error_type=PipelineManagerAlreadyExistsError" in rendered + fields = _emitted_fields(log_spy, as_error=False) + assert fields["event"] == API_ERROR_EVENT + assert fields["status"] == 409 + assert fields["error_type"] == "PipelineManagerAlreadyExistsError" diff --git a/tests/unit/test_logging_context.py b/tests/unit/test_logging_context.py deleted file mode 100644 index 52176df..0000000 --- a/tests/unit/test_logging_context.py +++ /dev/null @@ -1,33 +0,0 @@ -"""Unit tests for the per-request logging contextvars.""" - -import pytest - -from api.logging_context import bound_request_context, get_request_id, get_route_path - - -class TestLoggingContext: - def test_getters_return_none_outside_request(self): - assert get_request_id() is None - assert get_route_path() is None - - def test_getters_return_bound_values(self): - with bound_request_context(request_id="REQ123", route_path="/v1/start"): - assert get_request_id() == "REQ123" - assert get_route_path() == "/v1/start" - - def test_context_resets_on_clean_exit(self): - with bound_request_context(request_id="REQ123", route_path="/v1/start"): - pass - assert get_request_id() is None - assert get_route_path() is None - - def test_context_resets_when_body_raises(self): - def _raise_inside_context() -> None: - with bound_request_context(request_id="REQ123", route_path="/x"): - msg = "boom" - raise RuntimeError(msg) - - with pytest.raises(RuntimeError): - _raise_inside_context() - assert get_request_id() is None - assert get_route_path() is None diff --git a/tests/unit/test_pipeline_routes.py b/tests/unit/test_pipeline_routes.py index bd7ab7c..1b6f1b6 100644 --- a/tests/unit/test_pipeline_routes.py +++ b/tests/unit/test_pipeline_routes.py @@ -236,11 +236,9 @@ def test_parse_request_binds_pipe_code_and_pipeline_run_id_to_state(self, mocker # is logged with both fields. The unit-level tests pin the # handler->getter->log path; this one pins that `_parse_request` itself # actually writes to `request.state` against the production route. - # Both values are kept free of logfmt-active characters (whitespace, - # `=`, `"`) so the substring assertions below match the unquoted - # rendering. Future test additions that exercise quoted values should - # parse the logfmt line via the `_parse_logfmt` helper in - # `test_exception_handlers.py` instead of substring matching. + # The two values reach the record as attributes of their own, so nothing about their + # spelling matters here any more — the assertions read the `fields=` mapping the handler + # handed the runtime, not a rendered line. client, _, start_mock = _build_client(mocker) body_pipe_code = "echo" body_pipeline_run_id = "run-end-to-end-0001" @@ -257,16 +255,15 @@ def test_parse_request_binds_pipe_code_and_pipeline_run_id_to_state(self, mocker ) assert response.status_code == 500 log_spy.error.assert_called_once() - rendered = log_spy.error.call_args.args[0] - assert f"pipe_code={body_pipe_code}" in rendered - assert f"pipeline_run_id={body_pipeline_run_id}" in rendered + fields = log_spy.error.call_args.kwargs["fields"] + assert fields["pipe_code"] == body_pipe_code + assert fields["pipeline_run_id"] == body_pipeline_run_id def test_parse_request_drops_empty_correlation_fields(self, mocker: MockerFixture): - # An empty-string `pipe_code` / `pipeline_run_id` in the body must NOT - # render as a bare `pipe_code=` token in the operator log — the bare - # token reads as a logfmt parse error to downstream sinks and defeats - # grep-by-value. `_coerce_correlation_field` normalizes empty strings - # to `None`, and `emit_error_log` drops `None`-valued fields. + # An empty-string `pipe_code` / `pipeline_run_id` in the body must NOT reach the record + # as an empty attribute, which a downstream query filtering on presence would read as a + # real value. `_coerce_correlation_field` normalizes empty strings to `None`, and + # `_emit_api_error` drops those. client, _, start_mock = _build_client(mocker) start_mock.side_effect = PipelexConfigError("simulated config fault") log_spy = mocker.patch("api.exception_handlers.log") @@ -281,19 +278,17 @@ def test_parse_request_drops_empty_correlation_fields(self, mocker: MockerFixtur ) assert response.status_code == 500 log_spy.error.assert_called_once() - rendered = log_spy.error.call_args.args[0] - # No bare token of either kind — neither `pipe_code= ` nor at end-of-line. - assert "pipe_code=" not in rendered - assert "pipeline_run_id=" not in rendered + fields = log_spy.error.call_args.kwargs["fields"] + assert "pipe_code" not in fields + assert "pipeline_run_id" not in fields def test_parse_request_caps_oversized_pipe_code(self, mocker: MockerFixture): # `RunRequest.pipe_code` carries no Pydantic max_length, so a # caller can in principle send a megabyte-long string. The binding - # site caps the value rendered into operator logs so a single failed - # request cannot blow per-line log-sink budgets. 256 is the limit; - # anything longer is silently truncated for the log surface (the - # actual `run_request.pipe_code` passed to the runner is - # unchanged — only the `request.state` mirror is capped). + # site caps the value that reaches the operator record, so a single failed request cannot + # blow a log sink's per-record budget. 256 is the limit; anything longer is silently + # truncated for the log surface (the actual `run_request.pipe_code` passed to the runner + # is unchanged — only the `request.state` mirror is capped). client, _, start_mock = _build_client(mocker) start_mock.side_effect = PipelexConfigError("simulated config fault") log_spy = mocker.patch("api.exception_handlers.log") @@ -308,16 +303,15 @@ def test_parse_request_caps_oversized_pipe_code(self, mocker: MockerFixture): ) assert response.status_code == 500 log_spy.error.assert_called_once() - rendered = log_spy.error.call_args.args[0] - # The capped value (256 x's) appears in the log; the original 5000-x - # string does NOT — proves the cap fires and bounds the per-line cost. - assert f"pipe_code={'x' * 256}" in rendered - assert "x" * 5000 not in rendered + fields = log_spy.error.call_args.kwargs["fields"] + # The capped value carries; the original oversized string does not — proof the cap fires + # and bounds what one request costs a sink. + assert fields["pipe_code"] == "x" * 256 def test_parse_request_binds_pipe_code_before_extras_validation(self, mocker: MockerFixture): # The binding must run BEFORE `_validate_extras` so an SSRF-rejected # callback URL (or any other extras-validation 422) still rides the - # caller's `pipe_code` into the operator log. The unit-level tests + # caller's `pipe_code` onto the operator record. The unit-level tests # cannot exercise this ordering — only an end-to-end POST does. client, _, _ = _build_client(mocker) log_spy = mocker.patch("api.exception_handlers.log") @@ -334,12 +328,11 @@ def test_parse_request_binds_pipe_code_before_extras_validation(self, mocker: Mo assert response.status_code == 422 # An INPUT-domain 422 logs at `warning`, not `error`. log_spy.warning.assert_called_once() - rendered = log_spy.warning.call_args.args[0] - assert f"pipe_code={body_pipe_code}" in rendered + assert log_spy.warning.call_args.kwargs["fields"]["pipe_code"] == body_pipe_code def test_start_propagates_request_id_to_runner(self, mocker: MockerFixture): - # The middleware binds the inbound `X-Request-ID` onto the request-scoped - # contextvar; the route reads it via `get_request_id()` and passes it as + # The middleware stores the inbound `X-Request-ID` on `request.state`; the route reads it + # back via `request_id_of(request)` and passes it as # `request_id=` to `ApiRunner.start`, which forwards it to # `pipeline_run_setup(...)` so it lands on `JobMetadata.request_id`. # Without this hop the worker's `WorkflowLog` would carry `None`. diff --git a/tests/unit/test_request_id_middleware.py b/tests/unit/test_request_id_middleware.py index 335ac26..d41b359 100644 --- a/tests/unit/test_request_id_middleware.py +++ b/tests/unit/test_request_id_middleware.py @@ -1,30 +1,44 @@ """Unit tests for RequestIdMiddleware and request-id propagation.""" +import logging import re import time +import pytest from fastapi import APIRouter, FastAPI, HTTPException, Request from fastapi.testclient import TestClient +from pipelex import log +from pipelex.tools.log.log_context import get_log_context from starlette.middleware.base import BaseHTTPMiddleware +from starlette.types import Message, Receive, Scope, Send -from api.logging_context import get_request_id, get_route_path -from api.middleware import REQUEST_ID_HEADER, RequestIdMiddleware, generate_request_id, request_body_size_middleware +from api.middleware import REQUEST_ID_HEADER, RequestIdMiddleware, generate_request_id, request_body_size_middleware, request_id_of # Crockford Base32, 26 chars — the ULID alphabet (no I, L, O, U). _ULID_RE = re.compile(r"\A[0-9A-HJKMNP-TV-Z]{26}\Z") +# The message the `/emits-a-log-line` route logs, matched back out of the captured records. +_ROUTE_LOG_MESSAGE = "a line emitted from inside the request" + _router = APIRouter() @_router.get("/probe") async def probe(request: Request) -> dict[str, str | None]: + bound = get_log_context() return { - "ctx_request_id": get_request_id(), - "ctx_route_path": get_route_path(), + "ctx_request_id": bound.request_id if bound is not None else None, "state_request_id": request.state.request_id, + "helper_request_id": request_id_of(request), } +@_router.get("/emits-a-log-line") +async def emits_a_log_line() -> dict[str, str]: + log.warning(_ROUTE_LOG_MESSAGE) + return {"status": "logged"} + + @_router.get("/boom") async def boom() -> None: raise HTTPException(status_code=400, detail="deliberate") @@ -59,7 +73,51 @@ def test_generates_ulid_when_absent(self): body = response.json() assert body["ctx_request_id"] == request_id assert body["state_request_id"] == request_id - assert body["ctx_route_path"] == "/probe" + assert body["helper_request_id"] == request_id + + def test_bound_context_puts_the_request_id_on_every_record(self, caplog: pytest.LogCaptureFixture): + # The middleware binds the runtime's own log context for the request, so a record emitted + # anywhere underneath — a route's own line here, but equally one from deep inside pipelex — + # carries `request_id` as a record attribute. The runner no longer interpolates the id into + # a message, which is what makes the id a field a structured sink indexes. + with caplog.at_level(logging.WARNING): + response = _build_client().get("/emits-a-log-line") + assert response.status_code == 200 + request_id = response.headers[REQUEST_ID_HEADER] + emitted = [record for record in caplog.records if record.getMessage() == _ROUTE_LOG_MESSAGE] + assert len(emitted) == 1, f"expected exactly one captured record, got {[record.getMessage() for record in caplog.records]}" + assert getattr(emitted[0], "request_id", None) == request_id + # The id rides the record, never the message — a sink indexes the field, and an operator + # grepping the text of a line is not the contract any more. + assert request_id not in emitted[0].getMessage() + + @pytest.mark.asyncio + async def test_bound_context_is_released_when_the_request_ends(self): + # The binding is per-request: once the request is over, nothing emitted afterwards may inherit + # its id, or a shared worker would attribute later work to it. The middleware is driven on this + # test's own event loop rather than through `TestClient`, which runs the app in a portal thread + # whose context never reaches the test's — through it, a binding that leaked would look released. + bound_while_running: list[str | None] = [] + + async def inner(_scope: Scope, _receive: Receive, send: Send) -> None: + bound = get_log_context() + bound_while_running.append(bound.request_id if bound is not None else None) + await send({"type": "http.response.start", "status": 200, "headers": []}) + await send({"type": "http.response.body", "body": b""}) + + async def receive() -> Message: + return {"type": "http.request", "body": b"", "more_body": False} + + async def send(message: Message) -> None: + pass + + scope: Scope = {"type": "http", "method": "GET", "path": "/", "headers": []} + assert get_log_context() is None + await RequestIdMiddleware(inner)(scope, receive, send) + + assert len(bound_while_running) == 1 + assert bound_while_running[0] is not None, "the id must be bound while the request runs" + assert get_log_context() is None, "the binding outlived the request" def test_echoes_valid_inbound_id(self): response = _build_client().get("/probe", headers={REQUEST_ID_HEADER: "client-supplied-123"}) diff --git a/uv.lock b/uv.lock index f6a6947..3ce5b59 100644 --- a/uv.lock +++ b/uv.lock @@ -2311,7 +2311,7 @@ wheels = [ [[package]] name = "pipelex" -version = "0.65.0" +version = "0.66.0" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "aiofiles" }, @@ -2344,7 +2344,6 @@ dependencies = [ { name = "pyyaml" }, { name = "reportlab" }, { name = "requests" }, - { name = "rich" }, { name = "semantic-version" }, { name = "shortuuid" }, { name = "tomli" }, @@ -2353,9 +2352,9 @@ dependencies = [ { name = "typing-extensions" }, { name = "urllib3" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/a0/7d/d23f9899579135e8f681d25f926226a78ca3aa7224fcbc4fcc2b3ca37c64/pipelex-0.65.0.tar.gz", hash = "sha256:f394d9715c2d47f2b40bae65ccfcb1f419ea392a99a081ad0d1fa011f4b48960", size = 1606912, upload-time = "2026-09-25T02:01:48.511Z" } +sdist = { url = "https://files.pythonhosted.org/packages/6c/ee/2e85605ccf04ab984eb7c69cd268f7b0e6c28bc1206e6817a0457c894179/pipelex-0.66.0.tar.gz", hash = "sha256:5043cc6068358c38e7cef02b3853e0e2d30c9342c07025c9f5dcd62c1ee3686e", size = 1667402, upload-time = "2026-09-25T16:35:44.747Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/35/0f/106beef7bb5f582e18f6778582fd16c956f0381e5291e8acac82d265fc8e/pipelex-0.65.0-py3-none-any.whl", hash = "sha256:303a75cd815452e115f00e332579b30eb5b79b890439c398ab2f8067a55867f1", size = 2226337, upload-time = "2026-09-25T02:01:45.792Z" }, + { url = "https://files.pythonhosted.org/packages/e9/ee/790c5d8861a512afe5ceab5e8d945f5802e6599c7e90a09392c2169b8c01/pipelex-0.66.0-py3-none-any.whl", hash = "sha256:cc0294ad349929690bf20d8fe828dc6716a77a9d674d6c7759885ea2cf098883", size = 2302727, upload-time = "2026-09-25T16:35:42.515Z" }, ] [package.optional-dependencies] @@ -2384,7 +2383,7 @@ mistralai = [ [[package]] name = "pipelex-api" -version = "0.27.5" +version = "0.28.0" source = { editable = "." } dependencies = [ { name = "fastapi" }, @@ -2436,7 +2435,7 @@ requires-dist = [ { name = "mkdocs-meta-manager", marker = "extra == 'docs'", specifier = "==1.1.0" }, { name = "mypy", marker = "extra == 'dev'", specifier = ">=1.11.2" }, { name = "pandas-stubs", marker = "extra == 'dev'", specifier = ">=2.2.3.241126" }, - { name = "pipelex", extras = ["mistralai", "anthropic", "google", "google-genai", "bedrock", "fal"], specifier = "==0.65.0" }, + { name = "pipelex", extras = ["mistralai", "anthropic", "google", "google-genai", "bedrock", "fal"], specifier = "==0.66.0" }, { name = "pyjwt", specifier = ">=2.10.1" }, { name = "pylint", marker = "extra == 'dev'", specifier = ">=3.3.8" }, { name = "pyright", marker = "extra == 'dev'", specifier = ">=1.1.405" }, diff --git a/wip/structured-logs/review-deferrals.md b/wip/structured-logs/review-deferrals.md new file mode 100644 index 0000000..b34a076 --- /dev/null +++ b/wip/structured-logs/review-deferrals.md @@ -0,0 +1,51 @@ +# Findings deferred on the runner's structured-logs member + +The per-round traces for L-260916-b7759f, the pipelex-api member of sprint L-260916-a8cdd2; the campaign documents live at the workspace root under `wip/structured-logs/`. A finding lands here when a round read it as real and did not act on it at that pass's bar. An entry marked *unverified* rests on the reviewer's word alone: nobody reproduced it. Findings owned by another repo become ledger items instead, and those are named here rather than described twice. + +## L-260916-b7759f, round 1 (bar `open`) + +Three reviewers ran against `origin/dev`: cubic, `codex:review` and `code-review` at level medium. Codex returned no findings. Eight distinct findings merged; seven were verified in one batch and all seven confirmed, and they were fixed in this round. What follows is the one deferral and the two findings routed to `pipelex`. + +- **The `[tool.uv.sources]` pin names a bare commit, so a rebase or a deletion upstream breaks every job at once** (*unverified*, code-review). CI runs `uv lock --check` and `uv sync --frozen` on every pull request and the Docker build resolves the same rev, so an unreachable commit would fail all of them together with nothing in this repository to fix. Deferred because the exposure is the sprint's own and already has an owner: the pin is temporary by construction, the sprint's pin table records the collapse — replace the git source with the shipped version and re-lock — and the upstream branch is a member of the same sprint rather than a stranger's. Pinning a tag instead would mean cutting one per runtime commit the runner needs, which costs more than the window it closes. Left unverified because it describes a future state of the remote rather than a property of the code. + +Two findings were confirmed here but belong to the runtime, and both carry the verification in their notes: + +- The console sink renders only the message and ignores every carried attribute, so the fields this member moved out of the message are invisible on it — `L-260918-f7062d`. Its own documentation was corrected in this round to say the fields are structured-sink-only. +- The json sink writes Rich markup verbatim into the message, so a structured line carries console styling — `L-260918-83431f`, found on the sibling `pipelex-server` member and visible in this repository's own test output. + +### What the round fixed, for the record + +The image no longer copies any of the runtime's per-developer configuration tiers, verified by building a throwaway context and listing what survived; the configuration test reads the shipped file instead of the merged one, so a developer's local `sink = "console"` no longer reds a suite that CI keeps green; three places claiming a Rich-free boot were corrected, Rich being installed unconditionally through `typer` and `instructor`; `detail` is documented as riding the API-authored path only; the documented whole-directory config mount now warns that it drops the logging keys; and the boot refuses a `pipelex` without `log.context` by name, since the branch and the published release declare the same version and no specifier can tell them apart. + +## L-260916-b7759f, round 2 (bar `defects`) + +Two reviewers ran against `origin/dev`: cubic and `code-review` at level medium; Codex is disabled in this machine's review policy. Three findings; two were verified in one batch, both confirmed and fixed in this round. One is deferred. + +- **The release and bump-pipelex skills do not know about the git source** (*unverified*, cubic). `[tool.uv.sources]` overrides the `pipelex` specifier, but `.claude/skills/release/SKILL.md` still asserts an exact `==A.B.C` pin and `.claude/skills/bump-pipelex/SKILL.md` edits only that specifier, so neither removes the source, and no CI check refuses a release whose runtime resolves from git. Deferred as an improvement rather than a defect at this bar: the collapse is a step of the sprint's own train — the merge gesture replaces the git source with the shipped version and re-locks before this branch reaches `dev` — and a sprint pin found on a base branch is refused by the sprint frontier until it comes off. A CI guard or a skill step would make the collapse independent of the sprint, which is worth doing once this repository's pins are no longer only a sprint's. + +### What the round fixed, for the record + +The middleware test meant to prove the per-request log binding is released could not fail: `TestClient` runs the app in a portal thread whose context never reaches the test's, and a mutation that never released the binding left the whole unit suite green. It now drives the middleware on the test's own event loop and goes red under that mutation. The runtime-contract check was unreachable on the install it exists for: `api.main` validates the config at import, and a published `pipelex` refuses this server's `sink` key there, before `lifespan` ever ran the check. The check now runs at import, above the first config read, a structural test pins its place, and the docstrings, the comment and the changelog entry say what actually happens. + +## L-260916-b7759f, round 3 (bar `necessity`) + +cubic and `code-review` at level medium ran against `origin/dev`; code-review returned no findings. One finding, deferred at this bar because the previous pass did not introduce it and shipping without it loses no data. + +- **The body-size middleware's 413 writes no `api_error` record** (*unverified*, cubic). `api/middleware.py::_too_large_response` builds the problem document and returns it without going through `api.exception_handlers._emit_api_error`, so an oversized-body rejection never reaches the structured error stream, while `docs/logging.md` says an error response produces exactly one line. The cure is either to emit the record from the middleware — it holds the `Request` for `route`, and `request_id` comes from the bound context — or to qualify the doc so an operator does not query for a line that is never written. + +## L-260916-b7759f, round 4 (bar `freeze`) + +Four reviewers ran against `origin/dev` at profile 4, after the branch absorbed `dev` through v0.27.5 and moved its `pipelex` pin: cubic, `codex:review`, `codex:adversarial` and `code-review` at level medium. code-review returned no findings. Nine findings merged; the two that could have been critical, or that might have belonged to the runtime, were verified in one batch. Neither was critical, so this pass changed no code and every finding it did not reject is deferred here. + +Rejected: **`.dockerignore` drops a flavor's baked `.pipelex/api_.toml`** (`codex:review`), refuted. The pattern would drop such a file, but no tracked file is lost against `origin/dev`, and the one real flavor, `pipelex-server`'s `api-hosted` image, builds from its own tree with its own ignore file and copies its env files from there, so this repository's `.dockerignore` never applies to it. + +- **A distributed `/execute` loses `request_id` at the worker** (*verified*, `codex:adversarial`). The route calls `ApiRunner.execute` with no request id, and the runtime's `PipelexMTHDSProtocol.execute` takes none, so the job's `RunMetadata.request_id` stays `None`. In direct mode the middleware's binding shows through because the run happens in-process, but a Temporal worker binds its log context from the serialized job alone, so every worker-side record of a Temporal `/execute` lacks the id; `/start` threads it and is unaffected. The gap predates this branch — `origin/dev` makes the same call — and it loses a log attribute rather than data. The fix fits in this repository: give `ApiRunner.execute` a `request_id`, pass `request_id_of(request)` from the route, and stamp it on the built job's run metadata in `_OrchestratorPipeRun` before dispatch; a `request_id` parameter on the runtime's `execute` would be the cleaner long-term seam. `docs/logging.md` meanwhile says every record emitted during a request carries the id, which holds in-process and not on a Temporal worker, and wants a one-line qualification. +- **The body-size 413, and Starlette's own 404 and 405, write no `api_error` record** (*unverified*, cubic and `codex:adversarial`). This extends the round-3 deferral: on top of the 413 bypassing `_emit_api_error`, the streamed oversized path lets the app run on a truncated body before the response is replaced, so a downstream 422 record can precede a 413 the caller receives. `docs/logging.md` says an error response produces exactly one line. +- **Uvicorn's catch-all 500 traceback is plain text on the same stderr as the JSON lines** (*unverified*, cubic). Starlette's `ServerErrorMiddleware` re-raises after the handler logs, and `uvicorn.error` then prints a multi-line traceback to stderr, where `console_log_target = "stderr"` now sends the json sink too; `docs/logging.md` and the changelog name only the banner and the access log as uvicorn's plain-text output. +- **With no version specifier on `pipelex`, a non-`uv` install is unbounded, and the bump-pipelex collapse has no `==old` to swap** (*unverified*, cubic). The git source overrides a specifier under `uv`, so keeping `==` beside it would cost nothing and would bound a `pip install .`; the runtime-contract check refuses only a build without `log.context`. The collapse to a published version at release restores the specifier either way. +- **The runtime-contract test that the message names the installed version cannot fail** (*unverified*, cubic). `tests/unit/test_runtime_contract.py` asserts `"pipelex" in str(...)`, which the fixed wording satisfies whether or not the installed version is interpolated; patching `_installed_runtime_version` to a sentinel would make it bite. The same file's comment still speaks of "a published 0.59.0" (code-review, cosmetic). +- **The round-2 deferral's premise no longer holds** (*unverified*, cubic). `dev` now carries a `version-check.yml` that fails a release pull request whose manifest keeps `[tool.uv.sources]` or whose lock resolves `pipelex` from a non-registry source, and `.claude/skills/bump-pipelex/SKILL.md` deletes the block, so the gap that entry describes is closed on the release path; what remains is the missing specifier above. + +## The pin collapse + +The collapse onto the released `pipelex` 0.66.0 settles the three entries above that turn on the git source. `pyproject.toml` carries `==0.66.0` again with no `[tool.uv.sources]` section, so a non-`uv` install is bounded by the same exact specifier as every other path, and the runtime-contract guard and its test were deleted with the source they existed to cover, as the guard's own docstring said they would be.