This is an automated email from the ASF dual-hosted git repository.
henry3260 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 ca622fa4f89 Fix bulk endpoints returning 500 for a malformed action
entry (#73453) (#73505)
ca622fa4f89 is described below
commit ca622fa4f899ed69bacfe16f18456ba569e13079
Author: Henry Chen <[email protected]>
AuthorDate: Tue Sep 22 17:11:46 2026 +0800
Fix bulk endpoints returning 500 for a malformed action entry (#73453)
(#73505)
(cherry picked from commit 0def984b047662074680e2ab9b30082d99204c3d)
Co-authored-by: Y-C <[email protected]>
---
.../api_fastapi/core_api/datamodels/common.py | 19 ++++++-
.../api_fastapi/core_api/datamodels/test_common.py | 64 ++++++++++++++++++++++
2 files changed, 80 insertions(+), 3 deletions(-)
diff --git a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/common.py
b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/common.py
index 48b55f2f1e6..770e6440d9f 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/common.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/common.py
@@ -23,6 +23,7 @@ Common Data Models for Airflow REST API.
from __future__ import annotations
import enum
+from collections.abc import Mapping
from typing import Annotated, Any, Generic, Literal, TypeVar, Union
from pydantic import Discriminator, Field, Tag
@@ -98,8 +99,16 @@ class BulkDeleteAction(BulkBaseAction[T]):
action_on_non_existence: BulkActionNotOnExistence =
BulkActionNotOnExistence.FAIL
-def _action_discriminator(action: Any) -> str:
- return BulkAction(action["action"]).value
+def _action_discriminator(action: Any) -> str | None:
+ """Select a bulk action variant, returning ``None`` for anything
unrecognised."""
+ value = action.get("action") if isinstance(action, Mapping) else
getattr(action, "action", None)
+ try:
+ return BulkAction(value).value
+ except ValueError:
+ return None
+
+
+_BULK_ACTION_TAGS = ", ".join(repr(action.value) for action in BulkAction)
class BulkBody(StrictBaseModel, Generic[T]):
@@ -112,7 +121,11 @@ class BulkBody(StrictBaseModel, Generic[T]):
Annotated[BulkUpdateAction[T], Tag(BulkAction.UPDATE.value)],
Annotated[BulkDeleteAction[T], Tag(BulkAction.DELETE.value)],
],
- Discriminator(_action_discriminator),
+ Discriminator(
+ _action_discriminator,
+ custom_error_type="bulk_action_invalid",
+ custom_error_message=f"Each entry needs an 'action' of
{_BULK_ACTION_TAGS}",
+ ),
]
]
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/datamodels/test_common.py
b/airflow-core/tests/unit/api_fastapi/core_api/datamodels/test_common.py
new file mode 100644
index 00000000000..ef315e39f9e
--- /dev/null
+++ b/airflow-core/tests/unit/api_fastapi/core_api/datamodels/test_common.py
@@ -0,0 +1,64 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+from __future__ import annotations
+
+import pytest
+from pydantic import TypeAdapter, ValidationError
+
+from airflow.api_fastapi.core_api.datamodels.common import BulkBody,
BulkCreateAction
+from airflow.api_fastapi.core_api.datamodels.variables import VariableBody
+
+# The bulk body is generic; ``VariableBody`` is the smallest entity to
instantiate it with. The
+# discriminator under test is shared by every bulk endpoint (variables, pools,
connections,
+# Dag runs, task instances), so one concrete instantiation covers all of them.
+_bulk_adapter: TypeAdapter[BulkBody[VariableBody]] =
TypeAdapter(BulkBody[VariableBody])
+
+
+def test_bulk_action_discriminator_reads_the_tag_off_a_model_instance():
+ """
+ Validating an already-built action -- what ``model_validate`` on a model
instance does -- hands
+ the discriminator the instance rather than a mapping, so the tag has to be
read as an attribute.
+ """
+ action = BulkCreateAction[VariableBody](action="create",
entities=[VariableBody(key="k", value="v")])
+ validated = _bulk_adapter.validate_python({"actions": [action]})
+ assert isinstance(validated.actions[0], BulkCreateAction)
+
+
[email protected](
+ "action",
+ [
+ pytest.param("x", id="not_a_mapping"),
+ pytest.param(None, id="null"),
+ pytest.param(5, id="number"),
+ pytest.param({"entities": []}, id="missing_action_key"),
+ pytest.param({"action": "bogus", "entities": []},
id="unknown_action_value"),
+ pytest.param({"action": None, "entities": []}, id="null_action_value"),
+ pytest.param({"action": ["create"], "entities": []},
id="unhashable_action_value"),
+ ],
+)
+def
test_bulk_action_discriminator_reports_invalid_actions_as_validation_errors(action):
+ """
+ A callable discriminator is handed the *raw, unvalidated* input and
pydantic does not wrap what
+ it raises, so any exception escaping it surfaces as a 500 instead of a
422. Every malformed
+ ``action`` entry must instead come back as a ``ValidationError`` naming
the accepted tags.
+ """
+ with pytest.raises(ValidationError) as exc_info:
+ _bulk_adapter.validate_python({"actions": [action]})
+
+ (error,) = exc_info.value.errors()
+ assert error["loc"] == ("actions", 0)
+ assert error["msg"] == "Each entry needs an 'action' of 'create',
'delete', 'update'"