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:
"""