This is an automated email from the ASF dual-hosted git repository.

potiuk pushed a commit to branch v3-3-test
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/v3-3-test by this push:
     new 4140090d787 [v3-3-test] Mask nested entities when writing bulk 
audit-log entries (#70890) (#71043)
4140090d787 is described below

commit 4140090d7870e678b38f1e573651fea56f000614
Author: Jarek Potiuk <[email protected]>
AuthorDate: Tue Aug 4 03:10:29 2026 +0200

    [v3-3-test] Mask nested entities when writing bulk audit-log entries 
(#70890) (#71043)
    
    The audit-log maskers for Variables and Connections dispatch on top-level 
key
    names. A bulk request body has exactly one top-level key, `actions`, and the
    entities sit two levels down in `actions[].entities[]` -- so neither
    `val`/`value` nor `extra` was ever seen, and the payload was recorded as
    supplied.
    
    The masker also runs on the raw request body before validation, so `extra` 
can
    arrive as any JSON type. `json.loads` raises `TypeError` rather than
    `JSONDecodeError` for a non-string, which escaped the audit-log path 
entirely.
    Bulk bodies newly reach this branch, so the shape is now reachable where it
    previously was not.
    
    (cherry picked from commit fc8d6d8b6194674f4043fb4f632e7700876c0fd7)
---
 .../src/airflow/api_fastapi/logging/decorators.py  |  59 +++++++-
 .../unit/api_fastapi/logging/test_decorators.py    | 150 +++++++++++++++++++++
 2 files changed, 207 insertions(+), 2 deletions(-)

diff --git a/airflow-core/src/airflow/api_fastapi/logging/decorators.py 
b/airflow-core/src/airflow/api_fastapi/logging/decorators.py
index f32a2bffc70..71942cafe41 100644
--- a/airflow-core/src/airflow/api_fastapi/logging/decorators.py
+++ b/airflow-core/src/airflow/api_fastapi/logging/decorators.py
@@ -45,8 +45,51 @@ def _sanitize_for_stdlib_log(value: str) -> str:
     return value.replace("\r", " ").replace("\n", " ")
 
 
+def _mask_bulk_entities(extra_fields, mask_entity):
+    """
+    Apply per-entity masking to a bulk request body.
+
+    A ``BulkBody`` has exactly one top-level field, ``actions``; the entities 
carrying the
+    secrets sit two levels down, in ``actions[].entities[]``. The per-entity 
maskers below
+    inspect top-level key names, so handing them a bulk body means they see 
only the key
+    ``actions`` and pass its whole payload through untouched. Reach the 
entities first.
+
+    Returns ``None`` when the body is not bulk-shaped, so callers fall back to 
flat masking.
+    """
+    actions = extra_fields.get("actions")
+    if not isinstance(actions, list):
+        return None
+
+    masked_actions = []
+    for action in actions:
+        if not isinstance(action, dict):
+            masked_actions.append(action)
+            continue
+        entities = action.get("entities")
+        if not isinstance(entities, list):
+            masked_actions.append(action)
+            continue
+        # ``delete`` actions may list bare id/key strings rather than entity 
objects;
+        # those carry no secret and are left as they are.
+        masked_actions.append(
+            {
+                **action,
+                "entities": [mask_entity(e) if isinstance(e, dict) else e for 
e in entities],
+            }
+        )
+    return {**extra_fields, "actions": masked_actions}
+
+
 def _mask_connection_fields(extra_fields):
-    """Mask connection fields."""
+    """Mask connection fields, for either a single-entity or a bulk request 
body."""
+    bulk = _mask_bulk_entities(extra_fields, _mask_connection_entity)
+    if bulk is not None:
+        return bulk
+    return _mask_connection_entity(extra_fields)
+
+
+def _mask_connection_entity(extra_fields):
+    """Mask the fields of one connection."""
     result = {}
     for k, v in extra_fields.items():
         if k == "extra" and v:
@@ -59,7 +102,11 @@ def _mask_connection_fields(extra_fields):
                     result[k] = {ek: "***" for ek in parsed_extra}
                 else:
                     result[k] = "Expected JSON object in `extra` field, got 
non-dict JSON"
-            except json.JSONDecodeError:
+            except (json.JSONDecodeError, TypeError):
+                # ``extra`` is declared as a string, but this runs on the raw 
body before
+                # validation, so it can arrive as any JSON type -- a number or 
an already-decoded
+                # object makes ``json.loads`` raise TypeError rather than 
JSONDecodeError. Both
+                # are recorded without the value, instead of raising out of 
the audit-log path.
                 result[k] = "Encountered non-JSON in `extra` field"
         else:
             result[k] = secrets_masker.redact(v, k)
@@ -67,6 +114,14 @@ def _mask_connection_fields(extra_fields):
 
 
 def _mask_variable_fields(extra_fields):
+    """Mask variable values, for either a single-entity or a bulk request 
body."""
+    bulk = _mask_bulk_entities(extra_fields, _mask_variable_entity)
+    if bulk is not None:
+        return bulk
+    return _mask_variable_entity(extra_fields)
+
+
+def _mask_variable_entity(extra_fields):
     """
     Mask the variable value.
 
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 41b9883d6d3..ab9218e1dd2 100644
--- a/airflow-core/tests/unit/api_fastapi/logging/test_decorators.py
+++ b/airflow-core/tests/unit/api_fastapi/logging/test_decorators.py
@@ -122,3 +122,153 @@ class TestMaskVariableFields:
     def test_value_without_key_is_still_masked(self):
         result = _mask_variable_fields({"value": "secretval"})
         assert result == {"value": "***"}
+
+
+class TestMaskBulkFields:
+    """The bulk endpoints nest their entities, and the masking has to reach 
them."""
+
+    def test_bulk_variable_values_are_masked(self):
+        result = _mask_variable_fields(
+            {
+                "actions": [
+                    {
+                        "action": "create",
+                        "entities": [
+                            {"key": "campaign_signing_material", "value": 
"VARVAL_LEAK_token"},
+                            {"key": "other", "val": "VARVAL_LEAK_alias"},
+                        ],
+                        "action_on_existence": "overwrite",
+                    }
+                ]
+            }
+        )
+        assert result == {
+            "actions": [
+                {
+                    "action": "create",
+                    "entities": [
+                        {"key": "campaign_signing_material", "value": "***"},
+                        {"key": "other", "val": "***"},
+                    ],
+                    "action_on_existence": "overwrite",
+                }
+            ]
+        }
+
+    def test_bulk_connection_extra_is_masked(self):
+        """``extra`` is the load-bearing case for connections.
+
+        Key-name redaction already covered a nested ``password`` -- but only 
while
+        ``hide_sensitive_var_conn_fields`` is enabled, which is a deployment 
setting and is off
+        in this test environment. ``extra`` was never covered by it at all: 
the name is not a
+        recognised sensitive field, and the value is a JSON *string*, which 
``redact`` returns
+        unchanged. Masking ``extra`` is structural here, so it does not depend 
on that setting.
+        """
+        result = _mask_connection_fields(
+            {
+                "actions": [
+                    {
+                        "action": "create",
+                        "entities": [
+                            {
+                                "connection_id": "c1",
+                                "conn_type": "http",
+                                "password": "CONN_LEAK_pw",
+                                "extra": json.dumps({"token": 
"CONN_LEAK_token", "region": "eu"}),
+                            }
+                        ],
+                    }
+                ]
+            }
+        )
+        entity = result["actions"][0]["entities"][0]
+        assert entity["extra"] == {"token": "***", "region": "***"}
+        assert entity["connection_id"] == "c1"
+        # every value is gone, only the key names of ``extra`` remain
+        assert "CONN_LEAK_token" not in json.dumps(result)
+
+    @pytest.mark.parametrize(
+        ("body", "masker"),
+        [
+            (
+                {"actions": [{"action": "create", "entities": [{"key": "k", 
"value": "LEAK_v"}]}]},
+                _mask_variable_fields,
+            ),
+            (
+                {
+                    "actions": [
+                        {
+                            "action": "update",
+                            "entities": [{"connection_id": "c", "extra": 
json.dumps({"t": "LEAK_v"})}],
+                        }
+                    ]
+                },
+                _mask_connection_fields,
+            ),
+        ],
+        ids=["variables", "connections"],
+    )
+    def test_no_secret_survives_in_the_serialized_entry(self, body, masker):
+        """The audit entry is serialized whole, so assert on the 
serialization, not one field."""
+        assert "LEAK_v" not in json.dumps(masker(body))
+
+    def test_multiple_actions_and_entities_are_all_masked(self):
+        result = _mask_variable_fields(
+            {
+                "actions": [
+                    {
+                        "action": "create",
+                        "entities": [{"key": "a", "value": "s1"}, {"key": "b", 
"value": "s2"}],
+                    },
+                    {"action": "update", "entities": [{"key": "c", "value": 
"s3"}]},
+                ]
+            }
+        )
+        values = [e["value"] for a in result["actions"] for e in a["entities"]]
+        assert values == ["***", "***", "***"]
+
+    def test_delete_by_key_entities_are_left_alone(self):
+        """``delete`` may list bare keys rather than entity objects; those 
carry no secret."""
+        body = {"actions": [{"action": "delete", "entities": ["key_one", 
"key_two"]}]}
+        assert _mask_variable_fields(body) == body
+
+    @pytest.mark.parametrize(
+        "body",
+        [
+            {"key": "k", "value": "secret"},
+            {"actions": "not-a-list"},
+            {"actions": [{"action": "create"}]},
+            {"actions": [{"action": "create", "entities": "not-a-list"}]},
+            {"actions": ["not-a-dict"]},
+        ],
+        ids=["flat-body", "actions-not-list", "no-entities", 
"entities-not-list", "action-not-dict"],
+    )
+    def test_non_bulk_and_malformed_shapes_do_not_raise(self, body):
+        """The masker runs on request bodies before validation, so it must not 
add a failure mode."""
+        _mask_variable_fields(body)
+        _mask_connection_fields(body)
+
+    @pytest.mark.parametrize(
+        "extra",
+        [
+            pytest.param({"token": "CONN_LEAK_token"}, 
id="already-decoded-object"),
+            pytest.param(["CONN_LEAK_token"], id="already-decoded-array"),
+            pytest.param(123, id="number"),
+            pytest.param(True, id="bool"),
+        ],
+    )
+    def test_non_string_extra_does_not_raise_and_does_not_leak(self, extra):
+        """``extra`` reaches this before validation, so it need not be a 
string.
+
+        ``json.loads`` raises ``TypeError`` rather than ``JSONDecodeError`` 
for a non-string, which
+        would otherwise escape the audit-log decorator.
+        """
+        body = {"actions": [{"action": "create", "entities": 
[{"connection_id": "c1", "extra": extra}]}]}
+
+        result = _mask_connection_fields(body)
+
+        assert "CONN_LEAK_token" not in json.dumps(result)
+        assert result["actions"][0]["entities"][0]["connection_id"] == "c1"
+
+    def test_flat_body_still_takes_the_single_entity_path(self):
+        assert _mask_variable_fields({"key": "k", "value": "secret"}) == 
{"key": "k", "value": "***"}

Reply via email to