This is an automated email from the ASF dual-hosted git repository.
vincbeck 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 e4c58e48381 Authorize POST /assets/events on the asset named in the
body (#73007) (#73139)
e4c58e48381 is described below
commit e4c58e48381d07769c618611e0566facc96a4836
Author: Vincent <[email protected]>
AuthorDate: Mon Sep 14 14:52:51 2026 -0400
Authorize POST /assets/events on the asset named in the body (#73007)
(#73139)
The route dependency reads the asset id from the path, but this endpoint
carries it in the request body, so the auth manager was only asked whether
the caller may post to any asset at all. An auth manager that scopes assets
by id, name, or uri could not deny an event for an asset the caller may not
touch, and the response still returned that asset's name and uri.
Co-authored-by: Henry Chen <[email protected]>
---
.../api_fastapi/core_api/routes/public/assets.py | 5 ++-
.../src/airflow/api_fastapi/core_api/security.py | 32 ++++++++++++---
.../core_api/routes/public/test_assets.py | 47 ++++++++++++++++++++--
3 files changed, 74 insertions(+), 10 deletions(-)
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/assets.py
b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/assets.py
index fd359fd86bc..de780f24dba 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/assets.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/assets.py
@@ -366,7 +366,10 @@ def get_asset_events(
@assets_router.post(
"/assets/events",
responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]),
- dependencies=[Depends(requires_access_asset(method="POST")),
Depends(action_logging())],
+ dependencies=[
+ Depends(requires_access_asset(method="POST", asset_id_from_body=True)),
+ Depends(action_logging()),
+ ],
)
def create_asset_event(
body: CreateAssetEventsBody,
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 a78187be582..7d180d59067 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/security.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/security.py
@@ -56,6 +56,7 @@ from
airflow.api_fastapi.auth.managers.models.resource_details import (
)
from airflow.api_fastapi.common.db.common import SessionDep
from airflow.api_fastapi.core_api.base import OrmClause
+from airflow.api_fastapi.core_api.datamodels.assets import
CreateAssetEventsBody
from airflow.api_fastapi.core_api.datamodels.common import (
BulkAction,
BulkActionOnExistence,
@@ -967,19 +968,38 @@ def requires_access_dag_run_clear_bulk() ->
Callable[[BulkDAGRunClearBody, BaseU
return inner
-def requires_access_asset(method: ResourceMethod) -> Callable[[Request,
BaseUser], None]:
- def inner(
- request: Request,
- user: GetUserDep,
- ) -> None:
- asset_id = request.path_params.get("asset_id")
+def requires_access_asset(method: ResourceMethod, *, asset_id_from_body: bool
= False) -> Callable[..., None]:
+ """
+ Authorize the caller on the asset targeted by the request.
+
+ :param method: the method to perform
+ :param asset_id_from_body: read ``asset_id`` from a
``CreateAssetEventsBody`` request body instead of
+ the path. The dependency parameter must be named ``body`` to share the
route's body.
+ """
+ def _authorize(asset_id: str | None, user: BaseUser) -> None:
_requires_access(
is_authorized_callback=lambda:
get_auth_manager().is_authorized_asset(
method=method, details=AssetDetails(id=asset_id), user=user
),
)
+ if asset_id_from_body:
+
+ def inner_from_body(
+ body: CreateAssetEventsBody,
+ user: GetUserDep,
+ ) -> None:
+ _authorize(str(body.asset_id), user)
+
+ return inner_from_body
+
+ def inner(
+ request: Request,
+ user: GetUserDep,
+ ) -> None:
+ _authorize(request.path_params.get("asset_id"), user)
+
return inner
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_assets.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_assets.py
index e1b66fb7a65..e5a19448302 100644
--- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_assets.py
+++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_assets.py
@@ -26,7 +26,11 @@ from sqlalchemy import delete, func, select, update
from airflow._shared.timezones import timezone
from airflow.api_fastapi.auth.managers.base_auth_manager import BaseAuthManager
-from airflow.api_fastapi.auth.managers.models.resource_details import
DagAccessEntity, DagDetails
+from airflow.api_fastapi.auth.managers.models.resource_details import (
+ AssetDetails,
+ DagAccessEntity,
+ DagDetails,
+)
from airflow.api_fastapi.core_api.security import PermittedAssetEventFilter
from airflow.models import DagModel
from airflow.models.asset import (
@@ -1532,12 +1536,49 @@ class TestPostAssetEvents(TestAssets):
}
check_last_log(session, dag_id=None, event="create_asset_event",
logical_date=None)
+ @mock.patch(
+
"airflow.api_fastapi.auth.managers.simple.simple_auth_manager.SimpleAuthManager.is_authorized_asset",
+ autospec=True,
+ )
+ def test_should_authorize_on_the_asset_named_in_the_body(
+ self, mock_is_authorized_asset, test_client, session
+ ):
+ """The asset id lives in the body, so the route must read it from
there to authorize."""
+ (asset,) = self.create_assets(num=1, session=session)
+ mock_is_authorized_asset.return_value = True
+
+ response = test_client.post("/assets/events", json={"asset_id":
asset.id})
+
+ assert response.status_code == 200
+ mock_is_authorized_asset.assert_called_once_with(
+ mock.ANY,
+ method="POST",
+ details=AssetDetails(id=str(asset.id)),
+ user=mock.ANY,
+ )
+
+ @mock.patch(
+
"airflow.api_fastapi.auth.managers.simple.simple_auth_manager.SimpleAuthManager.is_authorized_asset",
+ autospec=True,
+ )
+ def test_should_respond_403_when_not_authorized_on_the_asset(
+ self, mock_is_authorized_asset, test_client, session
+ ):
+ (asset,) = self.create_assets(num=1, session=session)
+ mock_is_authorized_asset.return_value = False
+
+ response = test_client.post("/assets/events", json={"asset_id":
asset.id})
+
+ assert response.status_code == 403
+ assert session.scalar(select(func.count()).select_from(AssetEvent)) == 0
+
def test_should_respond_401(self, unauthenticated_test_client):
response = unauthenticated_test_client.post("/assets/events",
json={"asset_uri": "s3://bucket/key/1"})
assert response.status_code == 401
- def test_should_respond_403(self, unauthorized_test_client):
- response = unauthorized_test_client.post("/assets/events",
json={"asset_uri": "s3://bucket/key/1"})
+ def test_should_respond_403(self, unauthorized_test_client, session):
+ (asset,) = self.create_assets(num=1, session=session)
+ response = unauthorized_test_client.post("/assets/events",
json={"asset_id": asset.id})
assert response.status_code == 403
def test_invalid_attr_not_allowed(self, test_client, session):