This is an automated email from the ASF dual-hosted git repository.
bbovenzi 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 9fbec0321e5 Require a request body on the asset materialize endpoint
(#73328)
9fbec0321e5 is described below
commit 9fbec0321e5dea498256f0fe56272f7c0dff6a3a
Author: Pierre Jeambrun <[email protected]>
AuthorDate: Tue Sep 22 16:17:29 2026 +0200
Require a request body on the asset materialize endpoint (#73328)
* Require a request body on the asset materialize endpoint
The materialize endpoint accepted a bodyless POST and applied defaults,
unlike the other mutating endpoints -- including its sibling trigger-run route
-- which require an explicit JSON body. Requiring one here makes the contract
consistent and unambiguous, and updates the in-tree UI and airflowctl clients
accordingly.
* Add release note for the required materialize request body
---
airflow-core/newsfragments/73328.significant.rst | 6 ++++
.../core_api/openapi/v2-rest-api-generated.yaml | 6 ++--
.../api_fastapi/core_api/routes/public/assets.py | 14 ++++----
.../src/airflow/ui/openapi-gen/queries/queries.ts | 4 +--
.../airflow/ui/openapi-gen/requests/types.gen.ts | 2 +-
.../core_api/routes/public/test_assets.py | 41 +++++++++++++++++-----
airflow-ctl/src/airflowctl/api/operations.py | 2 +-
.../tests/airflow_ctl/api/test_operations.py | 2 ++
8 files changed, 53 insertions(+), 24 deletions(-)
diff --git a/airflow-core/newsfragments/73328.significant.rst
b/airflow-core/newsfragments/73328.significant.rst
new file mode 100644
index 00000000000..78349399960
--- /dev/null
+++ b/airflow-core/newsfragments/73328.significant.rst
@@ -0,0 +1,6 @@
+The ``POST /assets/{asset_id}/materialize`` endpoint now requires a JSON
request body.
+
+Previously a request with no body was accepted and the asset was materialized
with default
+values. The endpoint now requires an explicit JSON body, consistent with the
other mutating
+endpoints: a request without a body is rejected with ``422``, and an empty
HTML form is not a
+valid input anymore. All body fields remain optional, so an empty object
(``{}``) is a valid body.
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
index a8799ec436a..203a8e3ab54 100644
---
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
+++
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
@@ -591,13 +591,11 @@ paths:
type: integer
title: Asset Id
requestBody:
+ required: true
content:
application/json:
schema:
- anyOf:
- - $ref: '#/components/schemas/MaterializeAssetBody'
- - type: 'null'
- title: Body
+ $ref: '#/components/schemas/MaterializeAssetBody'
responses:
'200':
description: Successful Response
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 f1fcd55f001..bf29a043153 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
@@ -450,7 +450,7 @@ def materialize_asset(
dag_bag: DagBagDep,
user: GetUserDep,
session: SessionDep,
- body: MaterializeAssetBody | None = None,
+ body: MaterializeAssetBody,
) -> DAGRunResponse:
"""Materialize an asset by triggering a Dag run that produces it."""
dag_id_it = iter(
@@ -486,18 +486,16 @@ def materialize_asset(
dag = get_latest_version_of_dag(dag_bag, dag_id, session)
- resolved_body = body or MaterializeAssetBody()
-
try:
preloaded_dag_version = None
context_dag = dag
- if resolved_body.bundle_version is not None and not
dag.disable_bundle_versioning:
+ if body.bundle_version is not None and not
dag.disable_bundle_versioning:
preloaded_dag_version = DagVersion.get_latest_version(
- dag_id, bundle_version=resolved_body.bundle_version,
load_serialized_dag=True, session=session
+ dag_id, bundle_version=body.bundle_version,
load_serialized_dag=True, session=session
)
if not preloaded_dag_version:
raise DagVersionNotFound(
- f"DAG with dag_id: '{dag_id}' does not have a version for
bundle_version '{resolved_body.bundle_version}'"
+ f"DAG with dag_id: '{dag_id}' does not have a version for
bundle_version '{body.bundle_version}'"
)
context_dag = preloaded_dag_version.serialized_dag.dag
@@ -510,7 +508,7 @@ def materialize_asset(
f"Dag with dag_id: '{dag_id}' does not allow asset
materialization runs",
)
- params = resolved_body.validate_context(context_dag)
+ params = body.validate_context(context_dag)
return dag.create_dagrun(
run_id=params["run_id"],
logical_date=params["logical_date"],
@@ -525,7 +523,7 @@ def materialize_asset(
partition_date=params["partition_date"],
note=params["note"],
session=session,
- bundle_version=resolved_body.bundle_version,
+ bundle_version=body.bundle_version,
dag_version=preloaded_dag_version,
)
except (ParamValidationError, ValueError) as e:
diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
b/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
index a20b8210436..2fd053cb65c 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts
@@ -2285,10 +2285,10 @@ export const useAssetServiceCreateAssetEvent = <TData =
Common.AssetServiceCreat
*/
export const useAssetServiceMaterializeAsset = <TData =
Common.AssetServiceMaterializeAssetMutationResult, TError = unknown, TContext =
unknown>(options?: Omit<UseMutationOptions<TData, TError, {
assetId: number;
- requestBody?: MaterializeAssetBody;
+ requestBody: MaterializeAssetBody;
}, TContext>, "mutationFn">) => useMutation<TData, TError, {
assetId: number;
- requestBody?: MaterializeAssetBody;
+ requestBody: MaterializeAssetBody;
}, TContext>({ mutationFn: ({ assetId, requestBody }) =>
AssetService.materializeAsset({ assetId, requestBody }) as unknown as
Promise<TData>, ...options });
/**
* Create Backfill
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
index 8cf3d7bb7d5..e175c7203aa 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
@@ -3135,7 +3135,7 @@ export type CreateAssetEventResponse = AssetEventResponse;
export type MaterializeAssetData = {
assetId: number;
- requestBody?: MaterializeAssetBody | null;
+ requestBody: MaterializeAssetBody;
};
export type MaterializeAssetResponse = DAGRunResponse;
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 f48df57083c..47eff810abc 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
@@ -2305,13 +2305,13 @@ class TestPostAssetMaterialize(TestAssets):
return_value="Jane Doe",
)
def test_materialize_records_triggering_user_display_name(self,
mock_display_name, test_client):
- response = test_client.post("/assets/1/materialize")
+ response = test_client.post("/assets/1/materialize", json={})
assert response.status_code == 200
assert response.json()["triggering_user_name"] == "Jane Doe"
@pytest.mark.usefixtures("configure_git_connection_for_dag_bundle")
def test_should_respond_200(self, test_client):
- response = test_client.post("/assets/1/materialize")
+ response = test_client.post("/assets/1/materialize", json={})
assert response.status_code == 200
assert response.json() == {
"bundle_version": None,
@@ -2339,6 +2339,31 @@ class TestPostAssetMaterialize(TestAssets):
"team_name": None,
}
+ @pytest.mark.usefixtures("configure_git_connection_for_dag_bundle")
+ @pytest.mark.parametrize(
+ ("headers", "content"),
+ [
+ pytest.param(None, None, id="no-body"),
+ pytest.param(
+ {"content-type": "application/x-www-form-urlencoded"}, b"",
id="empty-urlencoded-form"
+ ),
+ pytest.param(
+ {"content-type": "application/x-www-form-urlencoded"},
+ b"partition_key=x",
+ id="urlencoded-form-with-field",
+ ),
+ pytest.param({"content-type": "multipart/form-data; boundary=x"},
b"", id="empty-multipart-form"),
+ pytest.param({"content-type": "text/plain"}, b"",
id="empty-text-plain"),
+ ],
+ )
+ def test_should_reject_missing_or_non_json_body(self, test_client,
session, headers, content):
+ # The request body is required and must be JSON, like the other
mutating endpoints. A bodyless
+ # request and a browser HTML form submission
(url-encoded/multipart/text, which is all a plain
+ # cross-origin <form> can send) are both rejected without queuing a
Dag run.
+ response = test_client.post("/assets/1/materialize", content=content,
headers=headers)
+ assert response.status_code == 422
+ assert session.scalar(select(func.count()).select_from(DagRun)) == 0
+
@pytest.mark.usefixtures("configure_git_connection_for_dag_bundle")
def test_should_respond_200_with_partition_key(self, test_client):
partition_key = "2026-03-23"
@@ -2401,12 +2426,12 @@ class TestPostAssetMaterialize(TestAssets):
assert response.status_code == 403
def test_should_respond_409_on_multiple_dags(self, test_client):
- response = test_client.post("/assets/2/materialize")
+ response = test_client.post("/assets/2/materialize", json={})
assert response.status_code == 409
assert response.json()["detail"] == "More than one Dag materializes
asset with ID: 2"
def test_should_respond_404_on_multiple_dags(self, test_client):
- response = test_client.post("/assets/3/materialize")
+ response = test_client.post("/assets/3/materialize", json={})
assert response.status_code == 404
assert response.json()["detail"] == "No Dag materializes asset with
ID: 3"
@@ -2422,7 +2447,7 @@ class TestPostAssetMaterialize(TestAssets):
.values(_data=data)
)
session.commit()
- response = test_client.post("/assets/1/materialize")
+ response = test_client.post("/assets/1/materialize", json={})
assert response.status_code == 400
assert (
response.json()["detail"]
@@ -2459,7 +2484,7 @@ class TestPostAssetMaterialize(TestAssets):
assert response.json()["bundle_version"] == "v1"
# Without bundle_version the latest (v2) governs and rejects the run.
- response = test_client.post("/assets/1/materialize")
+ response = test_client.post("/assets/1/materialize", json={})
assert response.status_code == 400
assert (
response.json()["detail"]
@@ -2474,7 +2499,7 @@ class TestPostAssetMaterialize(TestAssets):
) as mock_get_auth_manager:
mock_get_auth_manager.return_value.is_authorized_dag.return_value
= False
- response = test_client.post("/assets/1/materialize")
+ response = test_client.post("/assets/1/materialize", json={})
assert response.status_code == 403
assert response.json()["detail"] == (
@@ -2603,7 +2628,7 @@ class TestPostAssetMaterialize(TestAssets):
DagModel, "get_team_name", return_value=team_name,
autospec=True
) as mock_get_team_name,
):
- test_client.post("/assets/1/materialize")
+ test_client.post("/assets/1/materialize", json={})
assert len(recorded) == 1, "expected exactly one authorization check"
details = recorded[0]["details"]
diff --git a/airflow-ctl/src/airflowctl/api/operations.py
b/airflow-ctl/src/airflowctl/api/operations.py
index 0b8a0d9670a..c5f0dc392f0 100644
--- a/airflow-ctl/src/airflowctl/api/operations.py
+++ b/airflow-ctl/src/airflowctl/api/operations.py
@@ -293,7 +293,7 @@ class AssetsOperations(BaseOperations):
def materialize(self, asset_id: str) -> DAGRunResponse |
ServerResponseError:
"""Materialize an asset."""
- self.response = self.client.post(f"assets/{asset_id}/materialize")
+ self.response = self.client.post(f"assets/{asset_id}/materialize",
json={})
return DAGRunResponse.model_validate_json(self.response.content)
def get_queued_events(self, asset_id: str) ->
QueuedEventCollectionResponse | ServerResponseError:
diff --git a/airflow-ctl/tests/airflow_ctl/api/test_operations.py
b/airflow-ctl/tests/airflow_ctl/api/test_operations.py
index cdf78859058..f6a9802af26 100644
--- a/airflow-ctl/tests/airflow_ctl/api/test_operations.py
+++ b/airflow-ctl/tests/airflow_ctl/api/test_operations.py
@@ -468,6 +468,8 @@ class TestAssetsOperations:
def test_materialize(self):
def handle_request(request: httpx.Request) -> httpx.Response:
assert request.url.path ==
f"/api/v2/assets/{self.asset_id}/materialize"
+ # The endpoint requires a request body, so the client must send
one.
+ assert json.loads(request.content) == {}
return httpx.Response(200,
json=json.loads(self.dag_run_response.model_dump_json()))
client = make_api_client(transport=httpx.MockTransport(handle_request))