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

pierrejeambrun 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 6038d1156dd Scope audit log rows to a Dag only from the request path 
or body (#74232)
6038d1156dd is described below

commit 6038d1156dd99a426f56caa3af8b86786f22c84b
Author: Pierre Jeambrun <[email protected]>
AuthorDate: Tue Oct 6 11:08:05 2026 +0200

    Scope audit log rows to a Dag only from the request path or body (#74232)
    
    action_logging read dag_id from the merged query and path parameters, so a
    query filter could file an unrelated write under a Dag: POST
    /api/v2/connections?dag_id=x landed a connection write among Dag x's audit
    rows, which the per-Dag audit endpoints expose to anyone who can read Dag x.
    The Dag an audit row belongs to must come from the route the request targets
    -- its path, or the body of an endpoint that names one, such as a backfill 
--
    never a query parameter.
---
 .../src/airflow/api_fastapi/logging/decorators.py  | 22 ++++++---
 .../unit/api_fastapi/logging/test_decorators.py    | 52 ++++++++++++++++++++++
 2 files changed, 68 insertions(+), 6 deletions(-)

diff --git a/airflow-core/src/airflow/api_fastapi/logging/decorators.py 
b/airflow-core/src/airflow/api_fastapi/logging/decorators.py
index 926a4be8ed8..61cea89e3ec 100644
--- a/airflow-core/src/airflow/api_fastapi/logging/decorators.py
+++ b/airflow-core/src/airflow/api_fastapi/logging/decorators.py
@@ -152,7 +152,7 @@ def _mask_variable_entity(extra_fields):
     return result
 
 
-def _resolve_team_name(params: dict, *, session: Session) -> str | None:
+def _resolve_team_name(params: dict, *, dag_id: str | None, session: Session) 
-> str | None:
     """
     Return the team the audited action belongs to, for the resources that own 
no Dag.
 
@@ -169,7 +169,7 @@ def _resolve_team_name(params: dict, *, session: Session) 
-> str | None:
         # is committed before that runs, so recording it would fail the insert 
on a backend that
         # enforces the column width. The value stays visible in ``extra`` 
either way.
         return None if find_invalid_team_names([team_name]) else team_name
-    if params.get("dag_id"):
+    if dag_id:
         # Left to the insert-time hook on ``Log``, which covers every writer 
of an audit row rather
         # than only this one, and resolves a Dag's team through its bundle 
instead of a column.
         return None
@@ -261,6 +261,16 @@ def action_logging(event: str | None = None):
 
         extra_fields["method"] = request.method
 
+        # The dag_id/task_id/run_id columns scope an audit row to a Dag -- 
rows with a dag_id are the
+        # Dag-scoped ones the per-Dag audit endpoints expose to anyone who can 
read that Dag. They must
+        # name the resource the route acts on: the route path, or the request 
body of an endpoint that
+        # takes one (a backfill names its dag_id in the body). A query 
parameter must never scope the
+        # row, or ``POST /connections?dag_id=x`` would file that connection 
write among Dag x's rows.
+        scope = {**request.path_params}
+        if has_json_body:
+            scope.update(masked_body_json)
+        dag_id = scope.get("dag_id")
+
         # Create log entry
         log = Log(
             event=event_name,
@@ -268,10 +278,10 @@ def action_logging(event: str | None = None):
             owner=user_name,
             owner_display_name=user_display,
             extra=json.dumps(extra_fields),
-            task_id=params.get("task_id"),
-            dag_id=params.get("dag_id"),
-            run_id=params.get("run_id") or params.get("dag_run_id"),
-            team_name=_resolve_team_name(params, session=session),
+            task_id=scope.get("task_id"),
+            dag_id=dag_id,
+            run_id=scope.get("run_id") or scope.get("dag_run_id"),
+            team_name=_resolve_team_name(params, dag_id=dag_id, 
session=session),
         )
 
         if "logical_date" in request.query_params:
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 7c5775d8a26..009291c0057 100644
--- a/airflow-core/tests/unit/api_fastapi/logging/test_decorators.py
+++ b/airflow-core/tests/unit/api_fastapi/logging/test_decorators.py
@@ -444,6 +444,58 @@ class TestActionLoggingResourceTeamName:
 
         assert log.team_name == "infra"
 
+    @conf_vars({("core", "multi_team"): "True"})
+    def test_a_query_param_dag_id_does_not_hijack_the_resource_team(self, 
session):
+        """A ``?dag_id=`` query parameter must not scope the row to that Dag, 
so the team still comes
+        from the resource being acted on rather than being deferred to the 
unrelated Dag."""
+        self._create_resources_owned_by_a_team(session)
+        request = Request(
+            {
+                "type": "http",
+                "method": "DELETE",
+                "headers": [],
+                "query_string": b"dag_id=some_unrelated_dag",
+                "path_params": {"connection_id": "team_conn"},
+            }
+        )
+
+        asyncio.run(action_logging(event="delete_connection")(request=request, 
session=session, user=None))
+
+        log = session.scalar(select(Log).order_by(Log.id.desc()))
+        assert log.dag_id is None
+        assert log.team_name == "payments"
+
+
+class TestActionLoggingDagScope:
+    """``dag_id``/``task_id``/``run_id`` scope an audit row to a Dag, so a 
query parameter must not
+    set them: a query filter must not file an unrelated write (a connection, a 
variable) among a
+    Dag's audit rows, which the per-Dag audit endpoints then expose to anyone 
who can read that Dag.
+
+    Scoping from the route path and the request body is unchanged, and stays 
covered by the endpoint
+    tests -- ``test_favorite_dag`` for a path-named dag_id, 
``test_create_backfill`` for a body-named
+    one.
+    """
+
+    @staticmethod
+    def _logged_row(*, event, query_string):
+        request = Request(
+            {"type": "http", "method": "POST", "headers": [], "query_string": 
query_string, "path_params": {}}
+        )
+        session = MagicMock(spec=Session)
+        asyncio.run(action_logging(event=event)(request=request, 
session=session, user=None))
+        (logged,) = session.add.call_args.args
+        return logged
+
+    def test_query_param_dag_id_does_not_scope_the_row(self):
+        # POST /connections?dag_id=victim_dag must not land the connection 
write among victim_dag's rows.
+        logged = self._logged_row(event="post_connection", 
query_string=b"dag_id=victim_dag")
+        assert logged.dag_id is None
+
+    def test_query_param_task_id_and_run_id_do_not_scope_the_row(self):
+        logged = self._logged_row(event="post_connection", 
query_string=b"task_id=t1&run_id=r1")
+        assert logged.task_id is None
+        assert logged.run_id is None
+
 
 class TestActionLoggingUnparsableBody:
     """

Reply via email to