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