kaxil commented on code in PR #71220:
URL: https://github.com/apache/airflow/pull/71220#discussion_r4112788638
##########
airflow-core/src/airflow/api_fastapi/core_api/datamodels/pools.py:
##########
@@ -35,6 +36,28 @@ def _call_function(function: Callable[[], int]) -> int:
return function()
+def _apply_include_deferred_override(value: bool) -> bool:
+ override = Pool.get_include_deferred_override()
+ return value if override is None else override
+
+
+def _reject_user_provided_include_deferred(provided_fields: set[str]) -> bool
| None:
+ """
+ Get the cluster-wide ``include_deferred`` value, rejecting any value
provided by the user.
+
+ :param provided_fields: names of the fields explicitly present in the
request body.
+ :return: the cluster-wide value, or None when pools choose the flag
themselves.
+ :raises ValueError: if the flag is fixed cluster-wide but present in the
request body.
+ """
+ override = Pool.get_include_deferred_override()
+ if override is not None and "include_deferred" in provided_fields:
Review Comment:
With the config set, this 422s every pool write from airflowctl, not just
the ones where the user asked for a value. airflowctl's generated `PoolBody`
defaults `include_deferred` to `False`, its generated boolean flags default to
`False` as well, and `pools.create` dumps with `exclude_none=True`, so the key
is always sent. `pools.update` and `pools.bulk` dump without `exclude_unset`,
so an update sends `"include_deferred": null` (the parametrized `None` case in
`test_patch_pool_rejects_include_deferred_when_fixed_by_config` pins that as a
422), and `pools import` passes `include_deferred=pool_config.get(..., False)`
even when the file has no such key. The released airflowctl 1.0.0b1 does the
same. Since `airflow pools set/import/export` are `deprecated_for_airflowctl`,
the recommended CLI can't create, update or import pools on a cluster with this
set, which also contradicts the "whether via the API or the CLI" line in
`pools.rst`.
Could this treat `null` as not provided (and possibly accept a value equal
to the override), with airflowctl changed to only send fields the user passed?
That would also let a GET-then-PATCH round trip through, since GET now returns
the override value. Whichever contract you pick, an airflowctl test with the
option set would pin it.
##########
airflow-core/docs/administration-and-deployment/pools.rst:
##########
@@ -46,6 +46,14 @@ descendants.
Note that if tasks are not given a pool, they are assigned to a default pool
``default_pool``, which is
initialized with 128 slots and can be modified through the UI or CLI (but
cannot be removed).
+Whether deferred tasks occupy pool slots is normally decided per pool via its
``include_deferred`` flag.
+A Deployment Manager can instead fix this behavior for the whole cluster with
+:ref:`config:core__pool_include_deferred`. When that option is set to ``True``
or ``False``, the configured
Review Comment:
Should this say which components need the option? The scheduler applies it
in `slots_stats` and `open_slots`, while the UI checkbox and `PoolResponse`
reflect the API server's config. If it's set only on the API server, the UI
shows every pool as fixed to the configured value while the scheduler keeps
counting from the stored column.
##########
airflow-core/src/airflow/models/pool.py:
##########
@@ -92,6 +92,26 @@ class Pool(Base):
def __repr__(self):
return str(self.pool)
+ @staticmethod
+ def get_include_deferred_override() -> bool | None:
+ """
+ Get the cluster-wide ``include_deferred`` value fixed via config, if
any.
+
+ When ``[core] pool_include_deferred`` is set, its value applies to
every pool and takes
+ precedence over the per-pool ``include_deferred`` column. Returns None
when unset.
+ """
+ from airflow.configuration import conf
Review Comment:
There's no cycle to avoid here: `pool.py` already imports
`airflow.models.base`, which imports `conf` at module level, and
`configuration.py` doesn't import `airflow.models` at the top. Could this move
to the top of the file? The existing one in `create_or_update_pool` can go with
it.
##########
airflow-core/tests/unit/models/test_pool.py:
##########
@@ -343,6 +344,87 @@ def test_get_name_to_team_name_mapping(self, testing_team:
Team, session: Sessio
}
+class TestPoolIncludeDeferredOverride:
+ @staticmethod
+ def clean_db():
+ clear_db_dags()
+ clear_db_runs()
+ clear_db_pools()
+
+ def setup_method(self):
+ self.clean_db()
+
+ def teardown_method(self):
+ self.clean_db()
+
+ @pytest.mark.parametrize(
+ ("conf_value", "expected"),
+ [("", None), ("True", True), ("False", False)],
+ )
+ def test_get_include_deferred_override(self, conf_value, expected):
+ with conf_vars({("core", "pool_include_deferred"): conf_value}):
+ assert Pool.get_include_deferred_override() is expected
+
+ @pytest.mark.parametrize(
+ ("pool_value", "conf_value", "expected"),
+ [(False, "", False), (True, "", True), (False, "True", True), (True,
"False", False)],
+ )
+ def test_effective_include_deferred(self, pool_value, conf_value,
expected):
+ pool = Pool(pool="test_pool", slots=5, include_deferred=pool_value)
+ with conf_vars({("core", "pool_include_deferred"): conf_value}):
+ assert pool.effective_include_deferred is expected
+
+ @pytest.mark.parametrize(
+ ("pool_value", "conf_value", "deferred_is_occupied"),
+ [(False, "True", True), (True, "False", False)],
+ )
+ def test_get_occupied_states_uses_override(self, pool_value, conf_value,
deferred_is_occupied):
+ pool = Pool(pool="test_pool", slots=5, include_deferred=pool_value)
+ with conf_vars({("core", "pool_include_deferred"): conf_value}):
+ assert (TaskInstanceState.DEFERRED in pool.get_occupied_states())
is deferred_is_occupied
+
+ @conf_vars({("core", "pool_include_deferred"): "True"})
+ def test_slots_stats_use_override(self, dag_maker):
Review Comment:
This covers stored `False` with the option `True` only (same for
`test_create_or_update_pool_stores_override` with input `False`), so a
regression to `stored or override` would still pass while stored `True` under
`False` counts deferred tasks. Could both be parametrized over `(False,
"True")` and `(True, "False")` the way `test_get_occupied_states_uses_override`
is? `test_create_or_update_pool_stores_override` also takes `session` without
using it; asserting the stored row through it would check what the name says.
##########
airflow-core/src/airflow/cli/commands/pool_command.py:
##########
@@ -153,14 +172,17 @@ def pool_import_helper(filepath):
def pool_export_helper(filepath):
"""Help export all the pools to the json file."""
api_client = get_current_api_client()
+ # Omit the field when it is fixed cluster-wide, otherwise the exported
file cannot be imported back
+ export_include_deferred = Pool.get_include_deferred_override() is None
Review Comment:
Leaving the key out keeps a same-cluster round trip working, but the export
still prints "Exported N pools" with no hint that `include_deferred` was
dropped. Importing that file on a cluster without the option then creates every
pool with `include_deferred=False` through the `v.get("include_deferred",
False)` default, where before this PR the value travelled with the file. Worth
printing a line when it's omitted, and a sentence in `pools.rst` saying exports
carry no `include_deferred` while the option is set?
##########
airflow-core/src/airflow/api_fastapi/core_api/datamodels/pools.py:
##########
@@ -108,3 +139,10 @@ def validate_team_name(self) -> PoolBody:
"team_name cannot be set when multi_team mode is disabled.
Please contact your administrator."
)
return self
+
+ @model_validator(mode="after")
+ def enforce_include_deferred_override(self) -> PoolBody:
+ override =
_reject_user_provided_include_deferred(self.model_fields_set)
+ if override is not None:
+ self.include_deferred = override
Review Comment:
Assigning here adds `include_deferred` to `model_fields_set` (pydantic v2
`__setattr__` does that without `validate_assignment`), so the server-injected
value then looks user-provided to everything downstream. With the option set, a
bulk update of `default_pool` fails on its own: `PATCH /pools` with
`{"actions":[{"action":"update","update_mask":["slots"],"entities":[{"name":"default_pool","slots":64}]}]}`
hits `update_orm_from_pydantic`, which dumps
`include={"slots","include_deferred"}, exclude_unset=True`, gets
`include_deferred` back, and re-validates it through `PoolPatchBody`, which
rejects it. That's a `RequestValidationError`, not an `HTTPException`, so
`handle_bulk_update` doesn't catch it and the whole bulk request returns 422,
other actions included. Single-pool `PATCH /pools/default_pool` is fine, since
`PoolPatchBody` never assigns.
Could the value be set without marking it as provided? A
`Field(default_factory=lambda: Pool.get_include_deferred_override() or False)`
with only the reject check left in the validator does that, as does applying it
where the ORM row is built. The same leak is why bulk overwrite and mask-less
bulk update write the override into existing pools' column while single PATCH
doesn't; whichever way option 3 lands, it'd be good for those paths to agree.
There's no bulk-update test under the option today, only bulk create, so one
for `default_pool` with `update_mask=["slots"]` would catch this.
##########
airflow-core/src/airflow/api_fastapi/core_api/datamodels/pools.py:
##########
@@ -35,6 +36,28 @@ def _call_function(function: Callable[[], int]) -> int:
return function()
+def _apply_include_deferred_override(value: bool) -> bool:
+ override = Pool.get_include_deferred_override()
+ return value if override is None else override
Review Comment:
This is the same `value if override is None else override` as
`Pool.effective_include_deferred`, which this PR adds as the single place for
it (and it's repeated again in `slots_stats` and `create_or_update_pool`).
Since `PoolResponse` is built from the ORM object, `include_deferred: bool =
Field(validation_alias=AliasChoices("effective_include_deferred",
"include_deferred"))` would reuse the property and drop this helper, without
changing the OpenAPI schema.
##########
airflow-core/src/airflow/config_templates/config.yml:
##########
@@ -474,6 +474,18 @@ core:
type: integer
example: ~
default: "128"
+ pool_include_deferred:
+ description: |
+ Cluster-wide setting for the ``include_deferred`` flag of pools. When
left empty (the default),
+ each pool keeps its own ``include_deferred`` value, configurable per
pool via the UI, API or CLI.
+ When set to ``True`` or ``False``, the configured value is used for
**every** pool (including
+ pre-existing pools, whatever their stored value) when calculating
occupied slots, and users can
+ no longer choose the flag per pool: any attempt to set
``include_deferred`` when creating or
+ updating a pool is rejected, whether or not the requested value
matches the configured one.
+ version_added: 3.4.0
+ type: string
Review Comment:
`rerun_with_latest_version` in this same section is also unset/True/False,
and it's declared `type: boolean`, `default: ~` and read with
`conf.has_option(...)` (`routes/ui/config.py:70`). Could this follow the same
pattern rather than `type: string` with `""` as unset? It'd show as a boolean
in the config reference, and it would line up with how the other tri-state is
read. As it stands a value like `yes` makes `getboolean` raise on every read,
including `/ui/config`, which the UI loads at startup.
##########
airflow-core/tests/unit/cli/commands/test_pool_command.py:
##########
@@ -79,10 +81,81 @@ def test_pool_update_deferred(self):
pool_command.pool_set(self.parser.parse_args(["pools", "set", "foo",
"1", "test"]))
assert self.session.scalar(select(Pool).where(Pool.pool ==
"foo")).include_deferred is False
+ @pytest.mark.parametrize("conf_value", ["True", "False"])
+ def test_pool_set_include_deferred_rejected_when_fixed_by_config(self,
conf_value):
+ with conf_vars({("core", "pool_include_deferred"): conf_value}):
+ with pytest.raises(
+ SystemExit, match=f"include_deferred cannot be set because it
is fixed to {conf_value}"
+ ):
+ pool_command.pool_set(
+ self.parser.parse_args(["pools", "set", "locked_pool",
"1", "test", "--include-deferred"])
+ )
+ assert self.session.scalar(select(Pool).where(Pool.pool ==
"locked_pool")) is None
+
+ def test_pool_set_without_include_deferred_stores_config_value(self):
+ try:
Review Comment:
`TestCliPools.tearDown` is the unittest name, which pytest never calls on a
plain class, and that's why these tests wrap themselves in `try/finally:
self._cleanup()`. Renaming it to `teardown_method` would let the three new
wrappers go, and it would also clean up after the two new rejection tests,
which have no cleanup if a regression lets `locked_pool` get created.
##########
airflow-core/src/airflow/cli/commands/pool_command.py:
##########
@@ -27,11 +27,23 @@
from airflow.cli.simple_table import AirflowConsole
from airflow.cli.utils import deprecated_for_airflowctl
from airflow.exceptions import PoolNotFound
+from airflow.models.pool import Pool
from airflow.utils import cli as cli_utils
from airflow.utils.cli import suppress_logs_and_warning
from airflow.utils.providers_configuration_loader import
providers_configuration_loaded
+def get_include_deferred_rejection() -> str | None:
Review Comment:
This message is also built in `_reject_user_provided_include_deferred` in
`datamodels/pools.py`, and the two copies have already drifted (the API one
adds "Please contact your administrator."). A single helper next to
`Pool.get_include_deferred_override` returning the message, with the API
raising `ValueError` and the CLI `SystemExit`, keeps them in step. Nothing
outside this module uses it, so `_get_include_deferred_rejection` would match
`_show_pools` and friends.
##########
airflow-core/src/airflow/ui/src/pages/Pools/PoolForm.tsx:
##########
@@ -56,9 +57,13 @@ const PoolForm = ({ error, initialPool, isPending,
manageMutate, setError }: Poo
mode: "onChange",
});
const multiTeamEnabled = Boolean(useConfig("multi_team"));
+ const includeDeferredConfig = useConfig("pool_include_deferred");
+ // A boolean means include_deferred is fixed cluster-wide and cannot be
chosen per pool
+ const includeDeferredOverride =
+ typeof includeDeferredConfig === "boolean" ? includeDeferredConfig :
undefined;
const onSubmit = (data: PoolBody) => {
- manageMutate(data);
+ manageMutate(includeDeferredOverride === undefined ? data : { ...data,
include_deferred: undefined });
Review Comment:
Nothing tests this branch or the matching `update_mask` guard in
`useEditPool.ts`. If a later change sends `include_deferred: null` or passes
`data` through unchanged, every UI create and edit starts returning 422 with
the option set. A `PoolForm.test.tsx` mocking `useConfig` (as
`useRerunWithLatestVersion.test.tsx` does) that checks the checkbox is disabled
and `manageMutate` gets `include_deferred: undefined` would pin it.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]