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
+        )

Reply via email to