This is an automated email from the ASF dual-hosted git repository.
kaxil pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new 896bde50ba5 Fix 500 from the audit logger on a request body it cannot
parse (#73226)
896bde50ba5 is described below
commit 896bde50ba54d7c5472c888fce8713b22a024b65
Author: Kaxil Naik <[email protected]>
AuthorDate: Wed Sep 16 08:22:45 2026 +0100
Fix 500 from the audit logger on a request body it cannot parse (#73226)
---
.../src/airflow/api_fastapi/core_api/security.py | 6 +-
.../src/airflow/api_fastapi/logging/decorators.py | 26 ++--
.../unit/api_fastapi/logging/test_decorators.py | 137 ++++++++++++++++++++-
3 files changed, 158 insertions(+), 11 deletions(-)
diff --git a/airflow-core/src/airflow/api_fastapi/core_api/security.py
b/airflow-core/src/airflow/api_fastapi/core_api/security.py
index dc326924c59..3217fccc33f 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/security.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/security.py
@@ -585,8 +585,10 @@ def requires_access_backfill(
# Left: the routes naming their Dag in the body (create, dry run) or
in the query string
# (list, read by ``requires_access_dag``), and ids the handler's own
parser will reject.
dag_id = None
- # Not a json body, ignore
- with suppress(JSONDecodeError):
+ # An unreadable body names no Dag, the same state as no body, so it
falls through to the
+ # authorization below. Broad because ``json.loads`` also raises
UnicodeDecodeError, bare
+ # ValueError and RecursionError, none of them JSONDecodeError.
+ with suppress(Exception):
body = await request.json()
if isinstance(body, dict):
dag_id = body.get("dag_id")
diff --git a/airflow-core/src/airflow/api_fastapi/logging/decorators.py
b/airflow-core/src/airflow/api_fastapi/logging/decorators.py
index b459d784b13..926a4be8ed8 100644
--- a/airflow-core/src/airflow/api_fastapi/logging/decorators.py
+++ b/airflow-core/src/airflow/api_fastapi/logging/decorators.py
@@ -203,8 +203,19 @@ def action_logging(event: str | None = None):
masked_body_json = {}
if has_json_body:
- # Non-dict bodies fall through to the endpoint's own 422.
- parsed_body = await request.json()
+ # This runs before the route so an access is logged whether or not
the route succeeds,
+ # so it must never fail the request. On a route declaring no body
FastAPI parses none,
+ # leaving this the only parse, where an error became a 500.
+ #
+ # Broad like FastAPI's own catch here: ``json.loads`` also raises
UnicodeDecodeError,
+ # bare ValueError (``int_max_str_digits``) and RecursionError, and
enumerating types is
+ # how this line collected three fixes that each missed the next
shape.
+ try:
+ parsed_body = await request.json()
+ except Exception:
+ logger.debug("Audit log could not parse the request body;
logging without it")
+ parsed_body = None
+ # A non-dict body is left to the endpoint, which rejects it where
it declares one.
if isinstance(parsed_body, dict):
request_body = parsed_body
masked_body_json = {k: secrets_masker.redact(v, k) for k, v in
request_body.items()}
@@ -228,14 +239,13 @@ def action_logging(event: str | None = None):
for k, v in itertools.chain(request.query_params.items(),
request.path_params.items())
if k not in fields_skip_logging
}
+ # ``request_body`` is empty with no body, or one that could not be
read. The path/query
+ # fallback keeps the target in the row: a delete names it only there,
and ``Log`` has no
+ # column for it.
if "variable" in event_name:
- extra_fields = _mask_variable_fields(
- {k: v for k, v in request_body.items()} if has_json_body else
extra_fields
- )
+ extra_fields = _mask_variable_fields(request_body or extra_fields)
elif "connection" in event_name:
- extra_fields = _mask_connection_fields(
- {k: v for k, v in request_body.items()} if has_json_body else
extra_fields
- )
+ extra_fields = _mask_connection_fields(request_body or
extra_fields)
elif has_json_body:
extra_fields = {**extra_fields, **masked_body_json}
diff --git a/airflow-core/tests/unit/api_fastapi/logging/test_decorators.py
b/airflow-core/tests/unit/api_fastapi/logging/test_decorators.py
index e1b97ebcd2d..7c5775d8a26 100644
--- a/airflow-core/tests/unit/api_fastapi/logging/test_decorators.py
+++ b/airflow-core/tests/unit/api_fastapi/logging/test_decorators.py
@@ -18,11 +18,14 @@ from __future__ import annotations
import asyncio
import json
+import re
from unittest.mock import MagicMock
import pytest
from fastapi import Request
-from sqlalchemy import select
+from fastapi.dependencies.utils import get_flat_dependant
+from fastapi.routing import APIRoute
+from sqlalchemy import delete, func, select
from sqlalchemy.orm import Session
from airflow.api_fastapi.auth.managers.models.base_user import BaseUser
@@ -440,3 +443,135 @@ class TestActionLoggingResourceTeamName:
log = self._log_action(session, {"dag_id": "dag_owned_by_infra",
"pool_name": "team_pool"})
assert log.team_name == "infra"
+
+
+class TestActionLoggingUnparsableBody:
+ """
+ An unparsable body is logged as an access with no body, never raised.
+
+ The dependency runs before the route, so it must never fail the request.
On a route
+ declaring no body it is the only thing that parses one, where an error
became a 500.
+ """
+
+ @staticmethod
+ def _request(body: bytes) -> Request:
+ async def receive():
+ return {"type": "http.request", "body": body, "more_body": False}
+
+ return Request(
+ {
+ "type": "http",
+ "method": "DELETE",
+ "headers": [(b"content-type", b"application/json")],
+ "query_string": b"",
+ "path_params": {},
+ },
+ receive=receive,
+ )
+
+ @pytest.mark.parametrize(
+ "body",
+ [
+ pytest.param(b"{bad", id="malformed_json"),
+ pytest.param(b'{"a": ', id="truncated_json"),
+ pytest.param(b'"\xff"', id="invalid_utf8_in_string"),
+ pytest.param(b"\xff\xfe\xfd", id="invalid_utf8_bare"),
+ ],
+ )
+ def test_the_access_is_logged_with_no_body(self, body):
+ session = MagicMock(spec=Session)
+
+ asyncio.run(
+ action_logging(event="test_event")(
+ request=self._request(body), session=session,
user=MagicMock(spec=BaseUser)
+ )
+ )
+
+ (logged,) = session.add.call_args.args
+ assert logged.event == "test_event"
+
+
+def _bodyless_action_logging_routes(app) -> list[tuple[str, str]]:
+ """Every ``action_logging`` route that declares no request body of its
own."""
+ found = set()
+
+ def _collect(route):
+ if not isinstance(route, APIRoute):
+ return
+ uses_logging = any(
+ getattr(dep.call, "__qualname__", "").startswith("action_logging")
+ for dep in route.dependant.dependencies
+ )
+ if not uses_logging:
+ return
+ if get_flat_dependant(route.dependant, skip_repeats=True).body_params:
+ return
+ for method in route.methods or []:
+ found.add((method, route.path))
+
+ for route in app.routes:
+ _collect(route)
+ for sub in getattr(getattr(route, "app", None), "routes", []) or []:
+ _collect(sub)
+ return sorted(found)
+
+
[email protected]_test
+class TestNoActionLoggingRouteRejectsAnUnparsableBody:
+ """
+ No route may turn an unparsable body into a 500, asserted across all of
them.
+
+ Three earlier fixes to this line each patched the shape just reported --
an empty body
+ (#49035), a list (#62354), then non-dict bodies again -- each covered by a
test naming that
+ report's endpoint. Enumerating routes from the app, over every shape that
makes
+ ``json.loads`` raise, is what makes the next one fail here instead of in
production.
+ """
+
+ # One per exception type: JSONDecodeError, UnicodeDecodeError, bare
ValueError
+ # (int_max_str_digits) and RecursionError.
+ PAYLOADS = {
+ "malformed": b"{bad",
+ "invalid_utf8": b"\xff\xfe\xfd",
+ "oversized_int": b"9" * 4301,
+ "deeply_nested": b"[" * 10000 + b"]" * 10000,
+ }
+
+ # Authorized before ``action_logging`` and rejected on the placeholder
path value, so it never
+ # reaches the parse. Listed rather than skipped, since it must still not
500.
+ REJECTED_BEFORE_LOGGING = {("PUT", "/api/v2/parseDagFile/{file_token}")}
+
+ @pytest.mark.parametrize("payload_name", list(PAYLOADS))
+ def test_every_bodyless_route_tolerates_it(self, test_client, session,
payload_name):
+ routes = _bodyless_action_logging_routes(test_client.app)
+ assert routes, "found no bodyless action_logging routes, so this would
pass vacuously"
+
+ failures = []
+ for method, path in routes:
+ url = re.sub(r"\{[^}]+\}", "does-not-exist", path)
+ session.execute(delete(Log))
+ session.commit()
+ try:
+ response = test_client.request(
+ method,
+ url,
+ content=self.PAYLOADS[payload_name],
+ headers={"Content-Type": "application/json"},
+ )
+ except BaseException as exc:
+ # The client re-raises server exceptions (anyio may group
them), so catch per
+ # route to report every offender instead of stopping at the
first.
+ failures.append(f"{method} {path} -> raised
{type(exc).__name__}")
+ continue
+ if response.status_code == 500:
+ failures.append(f"{method} {path} -> 500")
+ continue
+ # A route rejected by an earlier dependency never reached the
parse, so requiring its
+ # audit row keeps it from passing vacuously -- one of them
silently did.
+ if (method, path) in self.REJECTED_BEFORE_LOGGING:
+ continue
+ if not session.scalar(select(func.count()).select_from(Log)):
+ failures.append(f"{method} {path} -> {response.status_code}
but wrote no audit row")
+
+ assert not failures, f"an unparsable body ({payload_name}) must not
fail these routes:\n" + "\n".join(
+ failures
+ )