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

Reply via email to