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

henry3260 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 0def984b047 Fix bulk endpoints returning 500 for a malformed action 
entry (#73453)
0def984b047 is described below

commit 0def984b047662074680e2ab9b30082d99204c3d
Author: Y-C <[email protected]>
AuthorDate: Tue Sep 22 04:22:53 2026 +0800

    Fix bulk endpoints returning 500 for a malformed action entry (#73453)
---
 .../api_fastapi/core_api/datamodels/common.py      | 19 ++++++--
 .../api_fastapi/core_api/datamodels/test_common.py | 50 +++++++++++++++++++++-
 2 files changed, 65 insertions(+), 4 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 3caad7b34a6..32570ee1fbb 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
@@ -24,6 +24,7 @@ from __future__ import annotations
 
 import enum
 import logging
+from collections.abc import Mapping
 from typing import Annotated, Any, Generic, Literal, TypeVar, Union
 
 from pydantic import BeforeValidator, Discriminator, Field, Tag, TypeAdapter, 
ValidationError
@@ -224,8 +225,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]):
@@ -238,7 +247,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
index d1c1515d6a9..8213b40461b 100644
--- 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
@@ -19,7 +19,13 @@ from __future__ import annotations
 import pytest
 from pydantic import TypeAdapter, ValidationError
 
-from airflow.api_fastapi.core_api.datamodels.common import AssetExpression, 
MaybeAssetExpression
+from airflow.api_fastapi.core_api.datamodels.common import (
+    AssetExpression,
+    BulkBody,
+    BulkCreateAction,
+    MaybeAssetExpression,
+)
+from airflow.api_fastapi.core_api.datamodels.variables import VariableBody
 
 # A single adapter is enough to validate/serialize the discriminated union.
 _adapter: TypeAdapter[AssetExpression] = TypeAdapter(AssetExpression)
@@ -125,3 +131,45 @@ def test_field_preserves_current_shapes(expression):
         assert validated is None
     else:
         assert _field_adapter.dump_python(validated, by_alias=True, 
exclude_unset=True) == expression
+
+
+# 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'"

Reply via email to