This is an automated email from the ASF dual-hosted git repository.
bbovenzi 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 77e0906dc31 Scope plugin views on any entity field, not a fixed
criteria set (#74278)
77e0906dc31 is described below
commit 77e0906dc31667c44ce66b24bb6eed14731535a4
Author: Brent Bovenzi <[email protected]>
AuthorDate: Wed Oct 7 15:22:35 2026 -0400
Scope plugin views on any entity field, not a fixed criteria set (#74278)
The five criteria `applies_to` accepted could not express the scoping plugin
authors actually want. The most-asked-for case -- offer a tab or panel only
on a
failed Dag Run or Task Instance -- was impossible, since state was not among
them, and no finite list of criteria will stay ahead of what Airflow
exposes on
an entity.
`applies_to` has not shipped in an airflow-core release yet, so the closed
set
can be replaced outright rather than grown, and no plugin author has to
migrate
off it later.
This is a trade rather than a pure win, and it is worth being explicit about
what is given up. A fixed vocabulary could be checked exhaustively; paths
are
resolved against the API response, so they can only be checked as far as the
response models describe themselves -- below a field typed as a bare dict,
such
as `class_ref`, a wrong path cannot be told from a right one. A path that
names
nothing is skipped at match time, which widens the scope rather than
narrowing
it, so the failure is a view appearing where it should not. Validating paths
when plugins load, with a suggested correction, is what keeps that
debuggable.
Matching stays equality-only. "Not successful" and "longer than five
minutes"
are the obvious next asks and neither is expressible here; they need
evaluation
the browser cannot do from the records on the page.
Co-authored-by: Claude <[email protected]>
---
.../docs/administration-and-deployment/plugins.rst | 172 +++++++++---
airflow-core/newsfragments/69148.feature.rst | 1 -
airflow-core/newsfragments/74278.feature.rst | 1 +
.../api_fastapi/core_api/datamodels/plugins.py | 31 ++-
.../core_api/openapi/v2-rest-api-generated.yaml | 61 ++---
airflow-core/src/airflow/plugins_manager.py | 216 +++++++++++++--
.../airflow/ui/openapi-gen/requests/schemas.gen.ts | 84 +-----
.../airflow/ui/openapi-gen/requests/types.gen.ts | 15 +-
.../ui/src/hooks/usePluginAppliesToContext.test.ts | 12 +
.../ui/src/hooks/usePluginAppliesToContext.ts | 12 +-
.../airflow/ui/src/hooks/usePluginTabs.test.tsx | 10 +-
.../ui/src/pages/Dag/Overview/Overview.test.tsx | 2 +-
.../ui/src/pages/Task/Overview/Overview.test.tsx | 2 +-
.../airflow/ui/src/utils/pluginAppliesTo.test.ts | 303 ++++++++++++++-------
.../src/airflow/ui/src/utils/pluginAppliesTo.ts | 202 +++++++++-----
.../core_api/routes/public/test_plugins.py | 30 +-
.../unit/cli/commands/test_plugins_command.py | 4 +-
airflow-core/tests/unit/plugins/test_plugin.py | 4 +-
.../tests/unit/plugins/test_plugins_manager.py | 220 +++++++++++++--
.../src/airflowctl/api/datamodels/generated.py | 20 +-
.../plugins_manager/plugins_manager.py | 12 +-
21 files changed, 977 insertions(+), 437 deletions(-)
diff --git a/airflow-core/docs/administration-and-deployment/plugins.rst
b/airflow-core/docs/administration-and-deployment/plugins.rst
index c0cd215b1cd..57a0722c78e 100644
--- a/airflow-core/docs/administration-and-deployment/plugins.rst
+++ b/airflow-core/docs/administration-and-deployment/plugins.rst
@@ -338,11 +338,12 @@ definitions in Airflow.
# are still grouped into the submenu; a single remaining non-promoted
item is also shown on the toolbar.
# Defaults to False.
"nav_top_level": True,
- # Optional scoping, limiting where this view is shown. Omit it
entirely to show the view
- # everywhere (the default). See "Scoping a view to specific Dags and
tasks" below.
+ # Optional scoping, limiting where this view is shown. Keys are dotted
field paths into
+ # the records the page has. Omit it entirely to show the view
everywhere (the default).
+ # See "Scoping a view to specific Dags and tasks" below.
"applies_to": {
- "dag_tags": ["production", "ml"],
- "dag_ids": ["my_dag", "my_other_dag"],
+ "dag.tags.name": ["production", "ml"],
+ "dag.dag_id": ["my_dag", "my_other_dag"],
},
}
@@ -374,11 +375,12 @@ definitions in Airflow.
# are still grouped into the submenu; a single remaining non-promoted
item is also shown on the toolbar.
# Defaults to False.
"nav_top_level": True,
- # Optional scoping, limiting where this app is shown. Omit it entirely
to show the app
- # everywhere (the default). See "Scoping a view to specific Dags and
tasks" below.
+ # Optional scoping, limiting where this app is shown. Keys are dotted
field paths into
+ # the records the page has. Omit it entirely to show the app
everywhere (the default).
+ # See "Scoping a view to specific Dags and tasks" below.
"applies_to": {
- "dag_tags": ["production", "ml"],
- "operators": ["KubernetesPodOperator"],
+ "dag.tags.name": ["production", "ml"],
+ "operator_name": ["KubernetesPodOperator"],
},
}
@@ -404,54 +406,134 @@ relevant instead of appearing on every Dag:
.. code-block:: python
"applies_to": {
- "dag_tags": ["ml"], # Dag carries any of these tags
- "dag_ids": ["train_pipeline"], # exact dag_id
- "task_ids": ["train_model"], # exact task_id
- "operators": ["KubernetesPodOperator"], # operator class name
+ "state": ["failed", "upstream_failed"], # the entity's own field
+ "dag.tags.name": ["ml"], # a related record, array-aware
+ "dag.dag_id": ["train_pipeline"],
}
-All keys are optional. ``operators`` and ``operator_names`` are matched
separately, the same
-way the task instance filters treat them: ``operators`` is the operator class
name, while
-``operator_names`` is the display name shown in the UI (an operator's
-``custom_operator_name``). For a plain operator the two are identical, so
either key works.
-They differ for decorator-based tasks: a ``@task.bash`` task has the display
name
-``@task.bash`` but the private class name ``_BashDecoratedOperator``, so use
-``operator_names`` to target it.
+Keys are **dotted field paths** into the records the page has, not a fixed set
of criteria, so
+anything the REST API returns for an entity is addressable — ``state``,
``operator``, ``pool``,
+``queue``, ``try_number``, and so on. Values are matched for equality, and are
compared as
+strings, so numeric and boolean fields work without quoting rules of their own
+(``"try_number": ["2"]``, ``"is_paused": ["false"]``).
-Criteria combine like Kubernetes label selectors — **OR within a key, AND
across keys**. A
-Dag matching any listed tag satisfies ``dag_tags``, and a view configured with
both
-``dag_tags`` and ``operators`` requires both to match.
+An **unqualified path is rooted at the entity the destination is about** — on
``dag_run``,
+``state`` is the run's state; on ``task_instance``, the task instance's. A
path may instead
+name a related record as its first segment: ``dag``, ``dag_run``, ``task`` or
+``task_instance``. Traversing a list fans out across it, so ``dag.tags.name``
collects every
+tag name and matches if any of them is listed.
-Crucially, the AND applies **only across criteria the current page can
evaluate**. A
-``task_ids`` criterion cannot be judged on a Dag-level page, so it is skipped
there rather
+Paths combine like Kubernetes label selectors — **OR within a path, AND across
paths**. A Dag
+matching any listed tag satisfies ``dag.tags.name``, and a view configured
with both
+``dag.tags.name`` and ``state`` requires both to match.
+
+Crucially, the AND applies **only across paths the current page can
evaluate**. A
+``task_instance.*`` path cannot be judged on a Dag-level page, so it is
skipped there rather
than failing the match. This lets one ``applies_to`` block be shared by a
plugin's Dag- and
-task-level destinations. Which criteria each destination can evaluate:
+task-level destinations. Which records each destination resolves:
.. list-table::
:header-rows: 1
* - Destination
- - ``dag_tags`` / ``dag_ids``
- - ``task_ids`` / ``operators`` / ``operator_names``
- * - ``dag``, ``dag_run``, ``dag_overview``
- - evaluated
- - skipped
- * - ``task``, ``task_overview``, ``task_instance``
- - evaluated
- - evaluated
+ - Unqualified path is rooted at
+ - Records a qualified path can reach
+ * - ``dag``, ``dag_overview``
+ - ``dag``
+ - ``dag``
+ * - ``dag_run``
+ - ``dag_run``
+ - ``dag``, ``dag_run``
+ * - ``task``, ``task_overview``
+ - ``task``
+ - ``dag``, ``task``
+ * - ``task_instance``
+ - ``task_instance``
+ - ``dag``, ``dag_run``, ``task``, ``task_instance``
* - ``nav``, ``base``, ``dashboard``, ``asset``
- - skipped
- - skipped
-
-If none of the configured criteria can be evaluated on a given page, the view
is shown. On
-task group pages the task-level criteria are skipped, since a group is not a
task.
-
-A malformed ``applies_to`` — one that is not a dictionary, names an unknown
criterion, or
-gives a criterion something other than a list of strings — is reported as a
warning when
-plugins are loaded, and ignored, so the view still loads unscoped. Configuring
a criterion
-the ``destination`` cannot evaluate (for example ``task_ids`` on a ``dag``
view) is also
-warned about, since it has no effect there. Check the API server log for these
warnings if a
-view is not scoped the way you expect.
+ - —
+ - none, so every path is skipped
+
+If none of the configured paths can be evaluated on a given page, the view is
shown. On task
+group pages the task-level records are absent, since a group is not a task.
+
+A path is also skipped when the record exists but has no such field. That
means a **bad path
+widens the scope rather than narrowing it**, so paths are checked against the
API response
+models when plugins load and a path naming no field is logged with a suggested
correction.
+The check stops at a field whose contents the models do not describe — the
dict behind
+``class_ref``, or a Dag Run's ``conf`` — so a path below one of those is still
accepted and
+still widens silently. If a view appears in more places than you expect, run
+``airflow plugins list --verbose`` first, then check the path against the REST
API response for
+that entity.
+
+Segments are the field names **the REST API returns**, which are not always
the Python
+attribute names: a task instance's run is ``dag_run_id``, not ``run_id``, and
computed fields
+such as a Dag's ``is_backfillable`` are addressable like any other. If a field
appears in the
+API response for an entity, it can be named here.
+
+An empty list is different from a missing field: a Dag with no tags has
definitively answered
+``dag.tags.name``, so the view is not shown. A path that stops short of a leaf
is likewise a
+decided answer rather than a skip: ``dag.tags`` resolves to a list of objects,
which have no
+comparable value, so the view is hidden. Address the field you mean to compare
+(``dag.tags.name``).
+
+Targeting operators
+^^^^^^^^^^^^^^^^^^^
+
+``operator_name`` is spelled the same way on a task and on a task instance, so
**one
+unqualified path targets an operator on either page**:
+
+.. code-block:: python
+
+ "applies_to": {"operator_name": ["KubernetesPodOperator"]}
+
+``operator_name`` is the display name shown in the UI (an operator's
+``custom_operator_name``). For a plain operator it is the class name, but the
two differ for
+decorator-based tasks: a ``@task.bash`` task has the display name
``@task.bash`` and the
+private class name ``_BashDecoratedOperator``, so ``operator_name`` is the one
to match on.
+
+If you specifically need the operator *class* name, the two records spell it
differently — a
+task instance has ``operator``, while a task carries it through
``class_ref.class_name``.
+Leave both paths **unqualified** so that each page reads its own record:
+
+.. code-block:: python
+
+ "applies_to": {
+ # On a task page the first is skipped and the second decides; on a
task instance
+ # page, the reverse. Prefer `operator_name` unless you need the
private class name.
+ "operator": ["KubernetesPodOperator"],
+ "class_ref.class_name": ["KubernetesPodOperator"],
+ }
+
+.. warning::
+ Do not write these as ``task_instance.operator`` and
``task.class_ref.class_name``. A task
+ instance page resolves **both** records, so both paths would be evaluated
and AND-ed — and
+ the task record always carries the Dag's *current* definition. An instance
that ran under
+ ``KubernetesPodOperator`` would stop matching the moment the task was
changed to a
+ different operator, hiding the view on historical runs.
+
+This is the general rule where records overlap: **prefer an unqualified
path**, which always
+reads the entity the page is about. Qualify a path only when you genuinely
mean the related
+record — ``dag.tags.name`` from a task instance, say — and not as a way of
naming the same
+concept twice.
+
+A malformed ``applies_to`` — one that is not a dictionary, has a non-string or
empty path, or
+gives a path something other than a list of strings — is reported as a warning
when plugins
+are loaded, and ignored, so the view still loads unscoped. Configuring a path
whose root the
+``destination`` cannot resolve (for example ``task_instance.state`` on a
``dag`` view) is also
+warned about, since it has no effect there.
+
+All of these are logged when the UI plugins are first collected, which happens
in the API
+server the first time anything requests ``/api/v2/plugins`` — and only once
per process, so
+reloading the page will not log them again. To see them on demand instead of
searching the API
+server log, run:
+
+.. code-block:: bash
+
+ airflow plugins list --verbose
+
+Each invocation is a fresh process, so every warning for every plugin is
reported. The
+``--verbose`` flag is required: without it the command suppresses log output.
.. note::
``applies_to`` is a display convenience, not an authorization boundary. It
controls
diff --git a/airflow-core/newsfragments/69148.feature.rst
b/airflow-core/newsfragments/69148.feature.rst
deleted file mode 100644
index fa359adfd51..00000000000
--- a/airflow-core/newsfragments/69148.feature.rst
+++ /dev/null
@@ -1 +0,0 @@
-Plugin external views and React apps can be scoped to specific Dags and tasks
with an optional ``applies_to`` block accepting ``dag_tags``, ``dag_ids``,
``task_ids``, ``operators``, and ``operator_names``, so a tab or panel is only
shown where it is relevant instead of on every Dag. Criteria are combined with
OR within a key and AND across keys, evaluated only against the keys the
current page can judge; omitting ``applies_to`` shows the item everywhere,
which is the previous behaviour. [...]
diff --git a/airflow-core/newsfragments/74278.feature.rst
b/airflow-core/newsfragments/74278.feature.rst
new file mode 100644
index 00000000000..69f483b4d0f
--- /dev/null
+++ b/airflow-core/newsfragments/74278.feature.rst
@@ -0,0 +1 @@
+Plugin external views and React apps can be scoped with an optional
``applies_to`` block: a map of dotted field path to the values that path may
take, such as ``{"state": ["failed"], "dag.tags.name": ["ml"]}``. Any field the
REST API returns for the entity is addressable, so a tab or panel appears only
where it is relevant instead of on every Dag. An unqualified path is rooted at
the entity the ``destination`` is about; a path may instead name a related
record (``dag``, ``dag_run``, ``ta [...]
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/plugins.py
b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/plugins.py
index ecc5c487cfd..f34a7f7cb90 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/plugins.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/plugins.py
@@ -19,7 +19,7 @@ from __future__ import annotations
from typing import Annotated, Any, Literal
-from pydantic import BeforeValidator, ConfigDict, Field, field_validator,
model_validator
+from pydantic import BeforeValidator, ConfigDict, Field, RootModel,
field_validator, model_validator
from airflow.api_fastapi.core_api.base import BaseModel
from airflow.plugins_manager import AirflowPluginSource, BaseDestinationLiteral
@@ -69,16 +69,19 @@ class AppBuilderMenuItemResponse(BaseModel):
category: str | None = None
-class PluginAppliesToResponse(BaseModel):
- """Serializer for the optional Dag/task scoping criteria of a UI plugin."""
+class PluginAppliesToResponse(RootModel[dict[str, list[str]]]):
+ """
+ Serializer for the optional scoping criteria of a UI plugin.
- model_config = ConfigDict(extra="forbid")
+ An open map of dotted field path to the values that path may take -- not a
closed set of
+ criteria. ``{"state": ["failed"], "dag.tags.name": ["ml"]}`` scopes to
failed entities of
+ ml-tagged Dags. An unqualified path is rooted at the entity the
``destination`` is about;
+ a path may instead name a related record (``dag``, ``dag_run``, ``task``,
``task_instance``)
+ as its first segment. Matching is equality against the listed values, OR
within a path and
+ AND across paths, and is evaluated client-side.
+ """
- dag_tags: list[str] | None = None
- dag_ids: list[str] | None = None
- task_ids: list[str] | None = None
- operators: list[str] | None = None
- operator_names: list[str] | None = None
+ root: dict[str, list[str]] = Field(default_factory=dict)
class BaseUIResponse(BaseModel):
@@ -92,11 +95,11 @@ class BaseUIResponse(BaseModel):
url_route: str | None = None
category: str | None = None
nav_top_level: bool | None = False
- # Optional visibility scoping, evaluated client-side. Criteria are OR-ed
within a
- # key and AND-ed across keys, but only across keys the current destination
can
- # actually evaluate (a `task_ids` criterion cannot be judged on a
Dag-level page,
- # so it is skipped there rather than failing the match). Omitting
`applies_to`
- # shows the item everywhere. Display gating only, not an authorization
boundary.
+ # Optional visibility scoping, evaluated client-side. Values are OR-ed
within a path and
+ # AND-ed across paths, but only across paths whose root record the current
destination
+ # actually has (a `task_instance.*` path cannot be judged on a Dag-level
page, so it is
+ # skipped there rather than failing the match). Omitting `applies_to`
shows the item
+ # everywhere. Display gating only, not an authorization boundary.
applies_to: PluginAppliesToResponse | None = None
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
index d74464351d5..671ac401abc 100644
---
a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
+++
b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml
@@ -16065,46 +16065,31 @@ components:
title: PatchTaskInstanceBody
description: Request body for patching task instance state.
PluginAppliesToResponse:
- properties:
- dag_tags:
- anyOf:
- - items:
- type: string
- type: array
- - type: 'null'
- title: Dag Tags
- dag_ids:
- anyOf:
- - items:
- type: string
- type: array
- - type: 'null'
- title: Dag Ids
- task_ids:
- anyOf:
- - items:
- type: string
- type: array
- - type: 'null'
- title: Task Ids
- operators:
- anyOf:
- - items:
- type: string
- type: array
- - type: 'null'
- title: Operators
- operator_names:
- anyOf:
- - items:
- type: string
- type: array
- - type: 'null'
- title: Operator Names
- additionalProperties: false
+ additionalProperties:
+ items:
+ type: string
+ type: array
type: object
title: PluginAppliesToResponse
- description: Serializer for the optional Dag/task scoping criteria of a
UI plugin.
+ description: 'Serializer for the optional scoping criteria of a UI
plugin.
+
+
+ An open map of dotted field path to the values that path may take --
not a
+ closed set of
+
+ criteria. ``{"state": ["failed"], "dag.tags.name": ["ml"]}`` scopes to
failed
+ entities of
+
+ ml-tagged Dags. An unqualified path is rooted at the entity the
``destination``
+ is about;
+
+ a path may instead name a related record (``dag``, ``dag_run``,
``task``,
+ ``task_instance``)
+
+ as its first segment. Matching is equality against the listed values,
OR within
+ a path and
+
+ AND across paths, and is evaluated client-side.'
PluginCollectionResponse:
properties:
plugins:
diff --git a/airflow-core/src/airflow/plugins_manager.py
b/airflow-core/src/airflow/plugins_manager.py
index b652f9a58f3..599b22fc201 100644
--- a/airflow-core/src/airflow/plugins_manager.py
+++ b/airflow-core/src/airflow/plugins_manager.py
@@ -19,13 +19,15 @@
from __future__ import annotations
+import difflib
import inspect
import json
import logging
+import types
from collections.abc import Iterable
from functools import cache
from pathlib import Path
-from typing import TYPE_CHECKING, Any
+from typing import TYPE_CHECKING, Annotated, Any, Union, get_args, get_origin
from airflow import settings
from airflow._shared.module_loading import import_string, qualname
@@ -159,40 +161,176 @@ def _get_plugins() -> tuple[list[AirflowPlugin],
dict[str, str]]:
return plugins, import_errors
-_DAG_APPLIES_TO_CRITERIA = frozenset({"dag_tags", "dag_ids"})
-_TASK_APPLIES_TO_CRITERIA = frozenset({"task_ids", "operators",
"operator_names"})
-_APPLIES_TO_CRITERIA = _DAG_APPLIES_TO_CRITERIA | _TASK_APPLIES_TO_CRITERIA
-
-# Which `applies_to` criteria each destination can resolve a record for. A
destination
-# missing from this mapping cannot evaluate any criterion. Kept in sync with
the table in
+# The records each destination can resolve, and therefore the path roots it
can evaluate. A
+# destination missing from this mapping can evaluate nothing. Kept in sync
with the table in
# docs/administration-and-deployment/plugins.rst.
-_EVALUABLE_CRITERIA_BY_DESTINATION: dict[str, frozenset[str]] = {
- "dag": _DAG_APPLIES_TO_CRITERIA,
- "dag_run": _DAG_APPLIES_TO_CRITERIA,
- "dag_overview": _DAG_APPLIES_TO_CRITERIA,
- "task": _APPLIES_TO_CRITERIA,
- "task_overview": _APPLIES_TO_CRITERIA,
- "task_instance": _APPLIES_TO_CRITERIA,
+_APPLIES_TO_ROOTS: dict[str, frozenset[str]] = {
+ "dag": frozenset({"dag"}),
+ "dag_overview": frozenset({"dag"}),
+ "dag_run": frozenset({"dag", "dag_run"}),
+ "task": frozenset({"dag", "task"}),
+ "task_overview": frozenset({"dag", "task"}),
+ "task_instance": frozenset({"dag", "dag_run", "task", "task_instance"}),
"nav": frozenset(),
"base": frozenset(),
"dashboard": frozenset(),
"asset": frozenset(),
}
+# Which record an unqualified path is rooted at -- the entity the destination
is about.
+_APPLIES_TO_ENTITY_ROOT: dict[str, str] = {
+ "dag": "dag",
+ "dag_overview": "dag",
+ "dag_run": "dag_run",
+ "task": "task",
+ "task_overview": "task",
+ "task_instance": "task_instance",
+}
+
+_APPLIES_TO_ROOT_NAMES = frozenset({"dag", "dag_run", "task", "task_instance"})
+
+
+def _applies_to_path_root(path: str, destination: str) -> str | None:
+ """
+ Return the record a path is rooted at, or ``None`` if the destination has
no entity.
+
+ A path may name a related record as its first segment; otherwise it is
rooted at the
+ entity the destination is about.
+ """
+ head, _, rest = path.partition(".")
+ if head in _APPLIES_TO_ROOT_NAMES and rest:
+ return head
+ return _APPLIES_TO_ENTITY_ROOT.get(destination)
+
+
+# Sentinel for an annotation that does not describe what it contains, so a
path cannot be
+# checked past it.
+_OPAQUE = object()
+
+
+@cache
+def _applies_to_root_models() -> dict[str, Any]:
+ """
+ Return the response model backing each path root, for validating paths at
plugin load.
+
+ Imported lazily because ``datamodels.plugins`` imports this module; a
module-level import
+ would be circular. Only ``_get_ui_plugins`` reaches this, so components
that load plugins
+ without serving the UI never pay for it.
+
+ These are the models the UI actually fetches for the ``applies_to``
context -- keep them in
+ step with ``AppliesToContext`` in ``src/utils/pluginAppliesTo.ts``.
+ """
+ from airflow.api_fastapi.core_api.datamodels.dag_run import DAGRunResponse
+ from airflow.api_fastapi.core_api.datamodels.dags import DAGResponse
+ from airflow.api_fastapi.core_api.datamodels.task_instances import
TaskInstanceResponse
+ from airflow.api_fastapi.core_api.datamodels.tasks import TaskResponse
+
+ return {
+ "dag": DAGResponse,
+ "dag_run": DAGRunResponse,
+ "task": TaskResponse,
+ "task_instance": TaskInstanceResponse,
+ }
+
+
+def _unwrap_applies_to_annotation(annotation: Any) -> Any:
+ """
+ Reduce a field annotation to the type a further path segment reads through.
+
+ ``Annotated`` and ``X | None`` wrappers are stripped, and a list is
stepped into, because
+ traversing one fans out across its elements. Returns ``_OPAQUE`` for a
union of several
+ real types, whose fields depend on which member a record actually holds.
+ """
+ while True:
+ origin = get_origin(annotation)
+ if origin is Annotated:
+ annotation = get_args(annotation)[0]
+ elif origin in (Union, types.UnionType):
+ members = [arg for arg in get_args(annotation) if arg is not
type(None)]
+ if len(members) != 1:
+ return _OPAQUE
+ annotation = members[0]
+ elif origin in (list, set, frozenset, tuple):
+ args = get_args(annotation)
+ if not args:
+ return _OPAQUE
+ annotation = args[0]
+ else:
+ return annotation
+
+
+def _applies_to_serialized_fields(model: Any) -> dict[str, Any] | None:
+ """
+ Map each field name *as the API serializes it* to the annotation behind it.
+
+ ``applies_to`` paths are resolved in the browser against the JSON a
response model produces,
+ so validation has to use the serialized names rather than the Python
attribute names:
+ ``TaskInstanceResponse.run_id`` reaches the browser as ``dag_run_id``, and
computed fields
+ such as ``DAGResponse.is_backfillable`` have no entry in ``model_fields``
at all.
+
+ Returns ``None`` for anything that is not a model.
+ """
+ fields = getattr(model, "model_fields", None)
+ if fields is None:
+ return None
+
+ serialized = {
+ field.serialization_alias or field.alias or name: field.annotation for
name, field in fields.items()
+ }
+ for name, computed in getattr(model, "model_computed_fields", {}).items():
+ serialized[computed.alias or name] = computed.return_type
+ return serialized
+
+
+def _describe_applies_to_path_error(path: str, root: str) -> str | None:
+ """
+ Return why ``path`` names no field on its root record, or ``None`` if it
is not knowably wrong.
+
+ The walk stops -- accepting whatever follows -- at a field the models do
not describe the
+ contents of, such as the bare ``dict`` behind ``class_ref`` or a Dag Run's
``conf``. So this
+ catches a misspelling of a modelled field, not every bad path.
+ """
+ segments = path.split(".")
+ # A qualified path's first segment names the record, which `root` has
already resolved.
+ if segments[0] == root and len(segments) > 1:
+ segments = segments[1:]
+
+ current: Any = _applies_to_root_models()[root]
+ for segment in segments:
+ if current is _OPAQUE or current is Any:
+ return None
+ base = get_origin(current) or current
+ if isinstance(base, type) and issubclass(base, dict):
+ return None
+
+ fields = _applies_to_serialized_fields(current)
+ if fields is None:
+ name = getattr(current, "__name__", repr(current))
+ return f"'{path}' reads '{segment}' from {name}, which has no
fields"
+ if segment not in fields:
+ model = getattr(current, "__name__", repr(current))
+ hint = ""
+ if close := difflib.get_close_matches(segment, list(fields), n=1):
+ hint = f" (did you mean '{close[0]}'?)"
+ return f"'{path}' names no field '{segment}' on {model}{hint}"
+
+ current = _unwrap_applies_to_annotation(fields[segment])
+ return None
+
def _describe_applies_to_error(applies_to: Any) -> str | None:
"""Return a description of why ``applies_to`` is malformed, or ``None`` if
it is valid."""
if not isinstance(applies_to, dict):
return f"expected a dictionary, got {type(applies_to).__name__}"
if non_string_keys := [key for key in applies_to if not isinstance(key,
str)]:
- return f"criterion names must be strings, got {sorted(non_string_keys,
key=repr)!r}"
- if unknown_keys := set(applies_to) - _APPLIES_TO_CRITERIA:
- return f"unknown criteria {sorted(unknown_keys)}, expected any of
{sorted(_APPLIES_TO_CRITERIA)}"
- for criterion, values in applies_to.items():
+ return f"field paths must be strings, got {sorted(non_string_keys,
key=repr)!r}"
+ if empty_keys := [key for key in applies_to if not key.strip()]:
+ return f"field paths must not be empty, got {empty_keys!r}"
+ for path, values in applies_to.items():
if values is None:
continue
if not isinstance(values, (list, tuple)) or not all(isinstance(value,
str) for value in values):
- return f"'{criterion}' must be a list of strings, got {values!r}"
+ return f"'{path}' must be a list of strings, got {values!r}"
return None
@@ -223,24 +361,54 @@ def _validate_applies_to(plugin_name: str | None, view:
ExternalViewDict | React
del view["applies_to"]
return
+ # A null value means "no values configured", which the matcher already
ignores the same way
+ # it ignores an empty list. Leaving it in place would fail
`PluginAppliesToResponse`
+ # serialization and drop the whole plugin -- including its other, valid
views -- from the
+ # plugins API. Dropping the path keeps the rest of the block working.
+ for path in [path for path, values in applies_to.items() if values is
None]:
+ del applies_to[path]
+
destination = view.get("destination", "nav")
- if destination not in _EVALUABLE_CRITERIA_BY_DESTINATION:
+ if destination not in _APPLIES_TO_ROOTS:
# An unrecognised destination already fails serialization; warning
here too would
# only add noise pointing at the wrong problem.
return
- configured = {criterion for criterion, values in applies_to.items() if
values}
- if unevaluable := configured -
_EVALUABLE_CRITERIA_BY_DESTINATION[destination]:
+ available = _APPLIES_TO_ROOTS[destination]
+ unevaluable = sorted(
+ path
+ for path, values in applies_to.items()
+ if values and _applies_to_path_root(path, destination) not in available
+ )
+ if unevaluable:
log.warning(
"Plugin '%s' has %s '%s' with destination '%s', which cannot
evaluate %s. "
- "Those criteria will be ignored.",
+ "Those paths will be ignored.",
plugin_name,
kind,
view.get("name"),
destination,
- sorted(unevaluable),
+ unevaluable,
)
+ # A path naming a field no record has is indistinguishable at match time
from one the page
+ # simply cannot judge, so the UI skips it -- which *widens* the scope
instead of narrowing
+ # it. Catching the misspelling here is the only place it can be told apart.
+ for path, values in applies_to.items():
+ root = _applies_to_path_root(path, destination)
+ if not values or root is None or root not in available:
+ continue
+ if error := _describe_applies_to_path_error(path, root):
+ log.warning(
+ "Plugin '%s' has %s '%s' with an 'applies_to' path that
matches no field: %s. "
+ "That path will be ignored, so the %s will appear in more
places than intended.",
+ plugin_name,
+ kind,
+ view.get("name"),
+ error,
+ kind.split()[-1],
+ )
+
@cache
def _get_ui_plugins() -> tuple[list[ExternalViewDict], list[ReactAppDict]]:
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
index 5f25205592d..08e49ded767 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts
@@ -6777,82 +6777,22 @@ export const $PatchTaskInstanceBody = {
} as const;
export const $PluginAppliesToResponse = {
- properties: {
- dag_tags: {
- anyOf: [
- {
- items: {
- type: 'string'
- },
- type: 'array'
- },
- {
- type: 'null'
- }
- ],
- title: 'Dag Tags'
+ additionalProperties: {
+ items: {
+ type: 'string'
},
- dag_ids: {
- anyOf: [
- {
- items: {
- type: 'string'
- },
- type: 'array'
- },
- {
- type: 'null'
- }
- ],
- title: 'Dag Ids'
- },
- task_ids: {
- anyOf: [
- {
- items: {
- type: 'string'
- },
- type: 'array'
- },
- {
- type: 'null'
- }
- ],
- title: 'Task Ids'
- },
- operators: {
- anyOf: [
- {
- items: {
- type: 'string'
- },
- type: 'array'
- },
- {
- type: 'null'
- }
- ],
- title: 'Operators'
- },
- operator_names: {
- anyOf: [
- {
- items: {
- type: 'string'
- },
- type: 'array'
- },
- {
- type: 'null'
- }
- ],
- title: 'Operator Names'
- }
+ type: 'array'
},
- additionalProperties: false,
type: 'object',
title: 'PluginAppliesToResponse',
- description: 'Serializer for the optional Dag/task scoping criteria of a
UI plugin.'
+ description: `Serializer for the optional scoping criteria of a UI plugin.
+
+An open map of dotted field path to the values that path may take -- not a
closed set of
+criteria. \`\`{"state": ["failed"], "dag.tags.name": ["ml"]}\`\` scopes to
failed entities of
+ml-tagged Dags. An unqualified path is rooted at the entity the
\`\`destination\`\` is about;
+a path may instead name a related record (\`\`dag\`\`, \`\`dag_run\`\`,
\`\`task\`\`, \`\`task_instance\`\`)
+as its first segment. Matching is equality against the listed values, OR
within a path and
+AND across paths, and is evaluated client-side.`
} as const;
export const $PluginCollectionResponse = {
diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
index ba3a3bd7a5f..e6bea43bda0 100644
--- a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
+++ b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts
@@ -1850,14 +1850,17 @@ export type PatchTaskInstanceBody = {
};
/**
- * Serializer for the optional Dag/task scoping criteria of a UI plugin.
+ * Serializer for the optional scoping criteria of a UI plugin.
+ *
+ * An open map of dotted field path to the values that path may take -- not a
closed set of
+ * criteria. ``{"state": ["failed"], "dag.tags.name": ["ml"]}`` scopes to
failed entities of
+ * ml-tagged Dags. An unqualified path is rooted at the entity the
``destination`` is about;
+ * a path may instead name a related record (``dag``, ``dag_run``, ``task``,
``task_instance``)
+ * as its first segment. Matching is equality against the listed values, OR
within a path and
+ * AND across paths, and is evaluated client-side.
*/
export type PluginAppliesToResponse = {
- dag_tags?: Array<(string)> | null;
- dag_ids?: Array<(string)> | null;
- task_ids?: Array<(string)> | null;
- operators?: Array<(string)> | null;
- operator_names?: Array<(string)> | null;
+ [key: string]: Array<(string)>;
};
/**
diff --git
a/airflow-core/src/airflow/ui/src/hooks/usePluginAppliesToContext.test.ts
b/airflow-core/src/airflow/ui/src/hooks/usePluginAppliesToContext.test.ts
index 0bc74811005..aaf515ea077 100644
--- a/airflow-core/src/airflow/ui/src/hooks/usePluginAppliesToContext.test.ts
+++ b/airflow-core/src/airflow/ui/src/hooks/usePluginAppliesToContext.test.ts
@@ -22,6 +22,7 @@ import { beforeEach, describe, expect, it, vi } from "vitest";
import type * as OpenapiQueries from "openapi/queries";
import {
+ UseDagRunServiceGetDagRunKeyFn,
UseDagServiceGetDagKeyFn,
UseTaskInstanceServiceGetMappedTaskInstanceKeyFn,
UseTaskServiceGetTaskKeyFn,
@@ -55,6 +56,7 @@ const { calls, record } = vi.hoisted(() => {
vi.mock("openapi/queries", async (importOriginal) => ({
...(await importOriginal<typeof OpenapiQueries>()),
+ useDagRunServiceGetDagRun: record("dagRun"),
useDagServiceGetDag: record("dag"),
useTaskInstanceServiceGetMappedTaskInstance: record("taskInstance"),
useTaskServiceGetTask: record("task"),
@@ -69,6 +71,7 @@ describe("usePluginAppliesToContext", () => {
renderHook(() => usePluginAppliesToContext(false));
expect(calls.dag?.options?.enabled).toBe(false);
+ expect(calls.dagRun?.options?.enabled).toBe(false);
expect(calls.task?.options?.enabled).toBe(false);
expect(calls.taskInstance?.options?.enabled).toBe(false);
});
@@ -77,6 +80,7 @@ describe("usePluginAppliesToContext", () => {
renderHook(() => usePluginAppliesToContext(true));
expect(calls.dag?.options?.enabled).toBe(true);
+ expect(calls.dagRun?.options?.enabled).toBe(true);
expect(calls.task?.options?.enabled).toBe(true);
expect(calls.taskInstance?.options?.enabled).toBe(true);
});
@@ -89,12 +93,16 @@ describe("usePluginAppliesToContext", () => {
renderHook(() => usePluginAppliesToContext(true));
expect(calls.dag?.key).toBeUndefined();
+ expect(calls.dagRun?.key).toBeUndefined();
expect(calls.task?.key).toBeUndefined();
expect(calls.taskInstance?.key).toBeUndefined();
expect(UseDagServiceGetDagKeyFn(calls.dag?.params as { dagId: string
})).toStrictEqual(
UseDagServiceGetDagKeyFn({ dagId }),
);
+ expect(
+ UseDagRunServiceGetDagRunKeyFn(calls.dagRun?.params as { dagId: string;
dagRunId: string }),
+ ).toStrictEqual(UseDagRunServiceGetDagRunKeyFn({ dagId, dagRunId: runId
}));
expect(
UseTaskServiceGetTaskKeyFn(calls.task?.params as { dagId: string;
taskId: unknown }),
).toStrictEqual(UseTaskServiceGetTaskKeyFn({ dagId, taskId }));
@@ -123,6 +131,8 @@ describe("usePluginAppliesToContext", () => {
renderHook(() => usePluginAppliesToContext(true));
expect(calls.dag?.options?.enabled).toBe(true);
+ // A group route still has a run, so the run record stays resolvable.
+ expect(calls.dagRun?.options?.enabled).toBe(true);
expect(calls.task?.options?.enabled).toBe(false);
expect(calls.taskInstance?.options?.enabled).toBe(false);
});
@@ -141,6 +151,8 @@ describe("usePluginAppliesToContext", () => {
renderHook(() => usePluginAppliesToContext(true));
expect(calls.dag?.options?.enabled).toBe(true);
+ // No run in the route, so a `dag_run.*` path is unevaluable here.
+ expect(calls.dagRun?.options?.enabled).toBe(false);
expect(calls.task?.options?.enabled).toBe(false);
expect(calls.taskInstance?.options?.enabled).toBe(false);
});
diff --git a/airflow-core/src/airflow/ui/src/hooks/usePluginAppliesToContext.ts
b/airflow-core/src/airflow/ui/src/hooks/usePluginAppliesToContext.ts
index 4d650ab97e0..829247235f6 100644
--- a/airflow-core/src/airflow/ui/src/hooks/usePluginAppliesToContext.ts
+++ b/airflow-core/src/airflow/ui/src/hooks/usePluginAppliesToContext.ts
@@ -19,6 +19,7 @@
import { useParams } from "react-router-dom";
import {
+ useDagRunServiceGetDagRun,
useDagServiceGetDag,
useTaskInstanceServiceGetMappedTaskInstance,
useTaskServiceGetTask,
@@ -46,6 +47,12 @@ export const usePluginAppliesToContext = (enabled: boolean):
AppliesToContext =>
enabled: enabled && Boolean(dagId),
});
+ const { data: dagRun, isLoading: isDagRunLoading } =
useDagRunServiceGetDagRun(
+ { dagId, dagRunId: runId },
+ undefined,
+ { enabled: enabled && Boolean(dagId) && Boolean(runId) },
+ );
+
const { data: task, isLoading: isTaskLoading } = useTaskServiceGetTask({
dagId, taskId }, undefined, {
enabled: enabled && Boolean(dagId) && Boolean(taskId) && groupId ===
undefined,
});
@@ -67,10 +74,11 @@ export const usePluginAppliesToContext = (enabled:
boolean): AppliesToContext =>
return {
dag,
+ dagRun,
// `isLoading` (not `isPending`) is deliberate: a disabled query reports
// `isPending` forever, which would withhold scoped views indefinitely on
- // destinations that legitimately have no task or task instance.
- isLoading: isDagLoading || isTaskLoading || isTaskInstanceLoading,
+ // destinations that legitimately have no run, task or task instance.
+ isLoading: isDagLoading || isDagRunLoading || isTaskLoading ||
isTaskInstanceLoading,
task,
taskInstance,
};
diff --git a/airflow-core/src/airflow/ui/src/hooks/usePluginTabs.test.tsx
b/airflow-core/src/airflow/ui/src/hooks/usePluginTabs.test.tsx
index e94fec1e7d4..0796bd4df22 100644
--- a/airflow-core/src/airflow/ui/src/hooks/usePluginTabs.test.tsx
+++ b/airflow-core/src/airflow/ui/src/hooks/usePluginTabs.test.tsx
@@ -67,8 +67,8 @@ describe("usePluginTabs", () => {
});
it.each([
- ["a matching applies_to", { dag_tags: ["ml"] }, 1],
- ["a non-matching applies_to", { dag_ids: ["etl_orders"] }, 0],
+ ["a matching applies_to", { "dag.tags.name": ["ml"] }, 1],
+ ["a non-matching applies_to", { "dag.dag_id": ["etl_orders"] }, 0],
])("includes the right tabs for %s", (_label, appliesTo:
PluginAppliesToResponse, expected) => {
setPlugins([makeView({ applies_to: appliesTo })]);
@@ -96,7 +96,7 @@ describe("usePluginTabs", () => {
it("withholds a scoped tab until its context resolves, avoiding a flicker",
() => {
mockUseContext.mockReturnValue({ dag: undefined, isLoading: true });
- setPlugins([makeView({ applies_to: { dag_tags: ["ml"] } }), makeView({
url_route: "always" })]);
+ setPlugins([makeView({ applies_to: { "dag.tags.name": ["ml"] } }),
makeView({ url_route: "always" })]);
const { result } = renderHook(() => usePluginTabs("dag_run"));
@@ -108,10 +108,10 @@ describe("usePluginTabs", () => {
["no view is scoped", [{}], false],
[
"a view for another destination is scoped",
- [{ applies_to: { dag_ids: ["x"] }, destination: "task" as const }],
+ [{ applies_to: { "dag.dag_id": ["x"] }, destination: "task" as const }],
false,
],
- ["a view for this destination is scoped", [{ applies_to: { dag_ids: ["x"]
} }], true],
+ ["a view for this destination is scoped", [{ applies_to: { "dag.dag_id":
["x"] } }], true],
])("resolves the context only when %s", (_label, views:
Array<Partial<ExternalViewResponse>>, enabled) => {
setPlugins(views.map(makeView));
diff --git
a/airflow-core/src/airflow/ui/src/pages/Dag/Overview/Overview.test.tsx
b/airflow-core/src/airflow/ui/src/pages/Dag/Overview/Overview.test.tsx
index 62260498a40..e035fddc669 100644
--- a/airflow-core/src/airflow/ui/src/pages/Dag/Overview/Overview.test.tsx
+++ b/airflow-core/src/airflow/ui/src/pages/Dag/Overview/Overview.test.tsx
@@ -37,7 +37,7 @@ vi.mock("openapi/queries", () => ({
{ bundle_url: "/dag.js", destination: "dag_overview", name: "Dag
overview plugin" },
{ bundle_url: "/task.js", destination: "task_overview", name:
"Task overview plugin" },
{
- applies_to: { dag_tags: ["finance"] },
+ applies_to: { "dag.tags.name": ["finance"] },
bundle_url: "/scoped.js",
destination: "dag_overview",
name: "Scoped overview plugin",
diff --git
a/airflow-core/src/airflow/ui/src/pages/Task/Overview/Overview.test.tsx
b/airflow-core/src/airflow/ui/src/pages/Task/Overview/Overview.test.tsx
index 60c2f62be89..d419c5ca172 100644
--- a/airflow-core/src/airflow/ui/src/pages/Task/Overview/Overview.test.tsx
+++ b/airflow-core/src/airflow/ui/src/pages/Task/Overview/Overview.test.tsx
@@ -55,7 +55,7 @@ vi.mock("openapi/queries", () => ({
{ bundle_url: "/dag.js", destination: "dag_overview", name: "Dag
overview plugin" },
{ bundle_url: "/task.js", destination: "task_overview", name:
"Task overview plugin" },
{
- applies_to: { operators: ["PythonOperator"] },
+ applies_to: { "class_ref.class_name": ["PythonOperator"] },
bundle_url: "/scoped.js",
destination: "task_overview",
name: "Scoped overview plugin",
diff --git a/airflow-core/src/airflow/ui/src/utils/pluginAppliesTo.test.ts
b/airflow-core/src/airflow/ui/src/utils/pluginAppliesTo.test.ts
index dd3695f0220..3db8e67ef77 100644
--- a/airflow-core/src/airflow/ui/src/utils/pluginAppliesTo.test.ts
+++ b/airflow-core/src/airflow/ui/src/utils/pluginAppliesTo.test.ts
@@ -20,6 +20,7 @@ import { describe, expect, it } from "vitest";
import type {
DAGResponse,
+ DAGRunResponse,
ExternalViewResponse,
PluginAppliesToResponse,
TaskInstanceResponse,
@@ -33,31 +34,43 @@ import {
matchesAppliesTo,
} from "./pluginAppliesTo";
-// These fixtures carry only the fields the matcher reads, so they are cast
through
-// `unknown` rather than spelling out every field of the full response types.
+// These fixtures carry only the fields the tests address, so they are cast
through `unknown`
+// rather than spelling out every field of the full response types.
const makeDag = (dagId: string, tagNames: Array<string>): DAGResponse =>
({
dag_id: dagId,
+ is_paused: false,
tags: tagNames.map((name) => ({ dag_display_name: dagId, dag_id: dagId,
name })),
}) as unknown as DAGResponse;
-const makeTask = (taskId: string, className: string, operatorName?: string):
TaskResponse =>
+const makeDagRun = (state: string | null): DAGRunResponse =>
+ ({ dag_id: "etl_sales", dag_run_id: "manual__1", state }) as unknown as
DAGRunResponse;
+
+const makeTask = (taskId: string, className: string): TaskResponse =>
({
class_ref: { class_name: className, module_path: "some.module" },
- operator_name: operatorName ?? className,
+ operator_name: className,
task_id: taskId,
}) as unknown as TaskResponse;
-const makeTaskInstance = (taskId: string, operator: string, operatorName?:
string): TaskInstanceResponse =>
+const makeTaskInstance = (overrides: Record<string, unknown> = {}):
TaskInstanceResponse =>
({
- operator,
- operator_name: operatorName ?? operator,
- task_id: taskId,
+ dag_id: "etl_sales",
+ map_index: -1,
+ operator: "KubernetesPodOperator",
+ operator_name: "KubernetesPodOperator",
+ state: "failed",
+ task_id: "train_model",
+ try_number: 2,
+ ...overrides,
}) as unknown as TaskInstanceResponse;
-const makeView = (appliesTo?: PluginAppliesToResponse): ExternalViewResponse
=> ({
+const makeView = (
+ appliesTo?: PluginAppliesToResponse,
+ destination: ExternalViewResponse["destination"] = "dag_run",
+): ExternalViewResponse => ({
applies_to: appliesTo,
- destination: "dag_run",
+ destination,
href: "/plugin/example",
name: "Example",
url_route: "example",
@@ -65,147 +78,229 @@ const makeView = (appliesTo?: PluginAppliesToResponse):
ExternalViewResponse =>
const dag = makeDag("etl_sales", ["ml", "prod"]);
-// A Dag-level page: no task or task instance in scope.
const dagContext: AppliesToContext = { dag, isLoading: false };
-
+const dagRunContext: AppliesToContext = { dag, dagRun: makeDagRun("failed"),
isLoading: false };
const taskContext: AppliesToContext = {
dag,
isLoading: false,
task: makeTask("train_model", "KubernetesPodOperator"),
};
-
const taskInstanceContext: AppliesToContext = {
dag,
+ dagRun: makeDagRun("failed"),
isLoading: false,
- taskInstance: makeTaskInstance("train_model", "KubernetesPodOperator"),
+ taskInstance: makeTaskInstance(),
};
-describe("matchesAppliesTo", () => {
- it("shows a view with no applies_to everywhere", () => {
- expect(matchesAppliesTo(makeView(), dagContext)).toBe(true);
- expect(matchesAppliesTo(makeView(), { isLoading: false })).toBe(true);
+describe("matchesAppliesTo — unqualified paths", () => {
+ it("shows a contribution with no applies_to everywhere", () => {
+ expect(matchesAppliesTo(makeView(), dagRunContext)).toBe(true);
+ expect(matchesAppliesTo(makeView({}), dagRunContext)).toBe(true);
});
- it("shows a view whose applies_to has no criteria", () => {
- expect(matchesAppliesTo(makeView({}), dagContext)).toBe(true);
+ // The case that motivated replacing the closed criteria set.
+ it("matches a Dag Run's own state", () => {
+ expect(matchesAppliesTo(makeView({ state: ["failed"] }),
dagRunContext)).toBe(true);
+ expect(matchesAppliesTo(makeView({ state: ["success"] }),
dagRunContext)).toBe(false);
});
- it("treats an empty criteria list as unset", () => {
- expect(matchesAppliesTo(makeView({ dag_ids: [] }), dagContext)).toBe(true);
+ it("roots an unqualified path at the entity the destination is about", () =>
{
+ const view = makeView({ state: ["failed"] }, "task_instance");
+
+ // Same path, different root: the task instance's state, not the run's.
+ expect(matchesAppliesTo(view, taskInstanceContext)).toBe(true);
+ expect(
+ matchesAppliesTo(view, {
+ ...taskInstanceContext,
+ taskInstance: makeTaskInstance({ state: "success" }),
+ }),
+ ).toBe(false);
});
- it.each([
- ["a matching tag", { dag_tags: ["ml"] }, true],
- ["a non-matching tag", { dag_tags: ["finance"] }, false],
- ["a matching dag_id", { dag_ids: ["etl_sales", "etl_orders"] }, true],
- ["a non-matching dag_id", { dag_ids: ["etl_orders"] }, false],
- ])("matches on %s", (_label, appliesTo, expected) => {
- expect(matchesAppliesTo(makeView(appliesTo), dagContext)).toBe(expected);
- });
-
- it("ANDs across keys, requiring every evaluable criterion to match", () => {
- expect(matchesAppliesTo(makeView({ dag_ids: ["etl_sales"], dag_tags:
["ml"] }), dagContext)).toBe(true);
- expect(matchesAppliesTo(makeView({ dag_ids: ["etl_sales"], dag_tags:
["finance"] }), dagContext)).toBe(
- false,
+ it("matches any of the listed values", () => {
+ const view = makeView({ state: ["failed", "upstream_failed"] });
+
+ expect(matchesAppliesTo(view, dagRunContext)).toBe(true);
+ expect(matchesAppliesTo(view, { ...dagRunContext, dagRun:
makeDagRun("queued") })).toBe(false);
+ });
+
+ it("compares non-string fields as strings", () => {
+ expect(matchesAppliesTo(makeView({ try_number: ["2"] }, "task_instance"),
taskInstanceContext)).toBe(
+ true,
+ );
+ expect(matchesAppliesTo(makeView({ map_index: ["-1"] }, "task_instance"),
taskInstanceContext)).toBe(
+ true,
);
+ expect(matchesAppliesTo(makeView({ is_paused: ["false"] }, "dag"),
dagContext)).toBe(true);
});
+});
- it("skips criteria the current context cannot evaluate", () => {
- // task_ids is unjudgeable on a Dag-level page, so the dag_tags match
decides.
- const view = makeView({ dag_tags: ["ml"], task_ids: ["train_model"] });
+describe("matchesAppliesTo — qualified paths", () => {
+ it("reaches a related record by naming it", () => {
+ expect(matchesAppliesTo(makeView({ "dag.dag_id": ["etl_sales"] }),
dagRunContext)).toBe(true);
+ expect(matchesAppliesTo(makeView({ "dag.dag_id": ["other"] }),
dagRunContext)).toBe(false);
+ });
- expect(matchesAppliesTo(view, dagContext)).toBe(true);
+ it("fans out across an array, matching if any element does", () => {
+ expect(matchesAppliesTo(makeView({ "dag.tags.name": ["ml"] }),
dagRunContext)).toBe(true);
+ expect(matchesAppliesTo(makeView({ "dag.tags.name": ["finance"] }),
dagRunContext)).toBe(false);
});
- it("shows a view when no configured criterion is evaluable at all", () => {
- expect(matchesAppliesTo(makeView({ task_ids: ["train_model"] }),
dagContext)).toBe(true);
- expect(matchesAppliesTo(makeView({ dag_ids: ["etl_sales"] }), { isLoading:
false })).toBe(true);
+ it("treats an empty array as no match, not as unevaluable", () => {
+ // A Dag with no tags has definitively answered the question; skipping
here would show the
+ // contribution on every untagged Dag.
+ expect(
+ matchesAppliesTo(makeView({ "dag.tags.name": ["ml"] }), {
+ dag: makeDag("etl_sales", []),
+ isLoading: false,
+ }),
+ ).toBe(false);
});
- it.each([
- ["task", taskContext],
- ["task instance", taskInstanceContext],
- ])("matches task_ids against the %s in scope", (_label, context) => {
- expect(matchesAppliesTo(makeView({ task_ids: ["train_model"] }),
context)).toBe(true);
- expect(matchesAppliesTo(makeView({ task_ids: ["evaluate"] }),
context)).toBe(false);
+ it("matches an operator class name through the task's class_ref", () => {
+ expect(
+ matchesAppliesTo(
+ makeView({ "task.class_ref.class_name": ["KubernetesPodOperator"] },
"task"),
+ taskContext,
+ ),
+ ).toBe(true);
});
+});
- it.each([
- ["task", taskContext],
- ["task instance", taskInstanceContext],
- ])("matches operators by class name against the %s in scope", (_label,
context) => {
- expect(matchesAppliesTo(makeView({ operators: ["KubernetesPodOperator"]
}), context)).toBe(true);
- expect(matchesAppliesTo(makeView({ operators: ["PythonOperator"] }),
context)).toBe(false);
+describe("matchesAppliesTo — the three verdicts", () => {
+ it("skips a path whose root record the surface does not have", () => {
+ // A task-instance path on a Dag page: unevaluable, so it must not fail
the match.
+ expect(matchesAppliesTo(makeView({ "task_instance.state": ["failed"] },
"dag"), dagContext)).toBe(true);
});
- it.each([
- ["task", { dag, isLoading: false, task: makeTask("submit",
"SparkSubmitOperator", "Spark Submit") }],
- [
- "task instance",
+ it("does not match when the field resolves to null", () => {
+ expect(
+ matchesAppliesTo(makeView({ state: ["failed"] }), {
+ ...dagRunContext,
+ dagRun: makeDagRun(null),
+ }),
+ ).toBe(false);
+ });
+
+ it("skips a path the record has no such field for", () => {
+ // Indistinguishable at runtime from an author typo, so it has to be
skipped: `operator`
+ // exists on a task instance but not on a task. Typos are caught at plugin
load instead.
+ expect(matchesAppliesTo(makeView({ operator: ["KubernetesPodOperator"] },
"task"), taskContext)).toBe(
+ true,
+ );
+ });
+
+ it("skips a segment naming an inherited member rather than an own field", ()
=> {
+ // Only own fields count. `in` would have found `toString` on the record's
prototype and
+ // resolved the path to a function -- which has no comparable form, so the
view would have
+ // been hidden on the strength of a field the record does not actually
have.
+ expect(matchesAppliesTo(makeView({ toString: ["[object Object]"] }),
dagRunContext)).toBe(true);
+ });
+
+ it("does not match when a path stops on an object instead of a leaf", () => {
+ // `dag.tags` is a list of objects, so there is nothing to compare. Unlike
a bad segment this
+ // narrows rather than widens: the path is answerable, and the answer is
no.
+ expect(matchesAppliesTo(makeView({ "dag.tags": ["ml"] }),
dagRunContext)).toBe(false);
+ });
+
+ it("skips everything on a destination with no entity record", () => {
+ expect(matchesAppliesTo(makeView({ state: ["failed"] }, "nav"),
dagRunContext)).toBe(true);
+ });
+});
+
+describe("matchesAppliesTo — combining paths", () => {
+ it("ANDs across paths the surface can evaluate", () => {
+ expect(matchesAppliesTo(makeView({ "dag.tags.name": ["ml"], state:
["failed"] }), dagRunContext)).toBe(
+ true,
+ );
+ expect(
+ matchesAppliesTo(makeView({ "dag.tags.name": ["finance"], state:
["failed"] }), dagRunContext),
+ ).toBe(false);
+ });
+
+ // The reason the skip rule exists: an author names both operator sources
and the block works
+ // on a task page and a task-instance page alike.
+ it("lets one block serve a task and a task instance", () => {
+ const view = makeView(
{
- dag,
- isLoading: false,
- taskInstance: makeTaskInstance("submit", "SparkSubmitOperator", "Spark
Submit"),
+ operator: ["KubernetesPodOperator"],
+ "task.class_ref.class_name": ["KubernetesPodOperator"],
},
- ],
- ])(
- "keeps class name and display name on separate keys for the %s in scope",
- (_label, context: AppliesToContext) => {
- expect(matchesAppliesTo(makeView({ operators: ["SparkSubmitOperator"]
}), context)).toBe(true);
- expect(matchesAppliesTo(makeView({ operator_names: ["Spark Submit"] }),
context)).toBe(true);
- // Each key sees only its own field, so a display name given to
`operators` does not match.
- expect(matchesAppliesTo(makeView({ operators: ["Spark Submit"] }),
context)).toBe(false);
- expect(matchesAppliesTo(makeView({ operator_names:
["SparkSubmitOperator"] }), context)).toBe(false);
- },
- );
-
- it("targets a decorated task by its display name, which is the only name it
exposes", () => {
- const context: AppliesToContext = {
- dag,
- isLoading: false,
- taskInstance: makeTaskInstance("run_script", "_BashDecoratedOperator",
"@task.bash"),
- };
+ "task",
+ );
+
+ expect(matchesAppliesTo(view, taskContext)).toBe(true);
+ expect(matchesAppliesTo({ ...view, destination: "task_instance" },
taskInstanceContext)).toBe(true);
+ });
+
+ // `operator_name` is spelled the same on a task and a task instance, so it
needs no second
+ // path -- the portable way to target an operator, and what the docs lead
with.
+ it("targets an operator on either destination with one portable path", () =>
{
+ const view = makeView({ operator_name: ["KubernetesPodOperator"] },
"task");
- expect(matchesAppliesTo(makeView({ operator_names: ["@task.bash"] }),
context)).toBe(true);
- expect(matchesAppliesTo(makeView({ operators: ["BashOperator"] }),
context)).toBe(false);
+ expect(matchesAppliesTo(view, taskContext)).toBe(true);
+ expect(matchesAppliesTo({ ...view, destination: "task_instance" },
taskInstanceContext)).toBe(true);
});
- it("uses the task instance operator instead of the latest Dag task
operator", () => {
- const context: AppliesToContext = {
+ // A task instance page resolves the task record too, and that record always
carries the Dag's
+ // current definition -- which may differ from what this instance actually
ran.
+ it("reads the page's own operator when the task has since been changed", ()
=> {
+ const staleContext: AppliesToContext = {
dag,
+ dagRun: makeDagRun("failed"),
isLoading: false,
- task: makeTask("run_it", "PythonOperator", "Current Python"),
- taskInstance: makeTaskInstance("run_it", "BashOperator", "Historical
Bash"),
+ task: makeTask("train_model", "PythonOperator"),
+ taskInstance: makeTaskInstance(),
};
- expect(matchesAppliesTo(makeView({ operators: ["BashOperator"] }),
context)).toBe(true);
- expect(matchesAppliesTo(makeView({ operators: ["PythonOperator"] }),
context)).toBe(false);
- expect(matchesAppliesTo(makeView({ operator_names: ["Historical Bash"] }),
context)).toBe(true);
- expect(matchesAppliesTo(makeView({ operator_names: ["Current Python"] }),
context)).toBe(false);
- });
-
- it("combines Dag- and task-level criteria on a task page", () => {
- const view = makeView({ dag_tags: ["ml"], operators:
["KubernetesPodOperator"] });
+ // Unqualified: the instance's own `operator` decides, and
`class_ref.class_name` is absent
+ // on a task instance and so is skipped.
+ expect(
+ matchesAppliesTo(
+ makeView(
+ {
+ "class_ref.class_name": ["KubernetesPodOperator"],
+ operator: ["KubernetesPodOperator"],
+ },
+ "task_instance",
+ ),
+ staleContext,
+ ),
+ ).toBe(true);
- expect(matchesAppliesTo(view, taskContext)).toBe(true);
+ // Qualified, which the docs warn against: both records resolve, so the
task's current
+ // (changed) class is AND-ed in and hides the view on this historical
instance.
expect(
- matchesAppliesTo(makeView({ dag_tags: ["finance"], task_ids:
["train_model"] }), taskContext),
+ matchesAppliesTo(
+ makeView(
+ {
+ "task.class_ref.class_name": ["KubernetesPodOperator"],
+ "task_instance.operator": ["KubernetesPodOperator"],
+ },
+ "task_instance",
+ ),
+ staleContext,
+ ),
).toBe(false);
});
+
+ it("ignores a path configured with no values", () => {
+ expect(matchesAppliesTo(makeView({ state: [] }),
dagRunContext)).toBe(true);
+ });
});
describe("isAppliesToPending", () => {
- it("never withholds a view without criteria", () => {
+ it("never withholds a contribution without scoping", () => {
expect(isAppliesToPending(makeView(), { isLoading: true })).toBe(false);
expect(isAppliesToPending(makeView({}), { isLoading: true })).toBe(false);
});
- it("withholds a scoped view while its context is loading", () => {
- expect(isAppliesToPending(makeView({ dag_tags: ["ml"] }), { isLoading:
true })).toBe(true);
+ it("withholds a scoped contribution while its context is loading", () => {
+ expect(isAppliesToPending(makeView({ state: ["failed"] }), { isLoading:
true })).toBe(true);
});
- it("releases a scoped view once its context has resolved", () => {
- expect(isAppliesToPending(makeView({ dag_tags: ["ml"] }),
dagContext)).toBe(false);
+ it("releases a scoped contribution once its context has resolved", () => {
+ expect(isAppliesToPending(makeView({ state: ["failed"] }),
dagRunContext)).toBe(false);
});
});
@@ -213,12 +308,8 @@ describe("hasAppliesToCriteria", () => {
it.each([
["no applies_to", undefined, false],
["an empty applies_to", {}, false],
- ["only empty criteria lists", { dag_ids: [], operator_names: [] }, false],
- ["dag_tags", { dag_tags: ["ml"] }, true],
- ["dag_ids", { dag_ids: ["etl_sales"] }, true],
- ["task_ids", { task_ids: ["train_model"] }, true],
- ["operators", { operators: ["KubernetesPodOperator"] }, true],
- ["operator_names", { operator_names: ["@task.bash"] }, true],
+ ["only empty value lists", { "dag.dag_id": [], state: [] }, false],
+ ["a populated path", { state: ["failed"] }, true],
])("reports %s as %s", (_label, appliesTo, expected) => {
expect(hasAppliesToCriteria(makeView(appliesTo))).toBe(expected);
});
diff --git a/airflow-core/src/airflow/ui/src/utils/pluginAppliesTo.ts
b/airflow-core/src/airflow/ui/src/utils/pluginAppliesTo.ts
index 86fce62e4b3..179a14f3502 100644
--- a/airflow-core/src/airflow/ui/src/utils/pluginAppliesTo.ts
+++ b/airflow-core/src/airflow/ui/src/utils/pluginAppliesTo.ts
@@ -18,6 +18,7 @@
*/
import type {
DAGResponse,
+ DAGRunResponse,
ExternalViewResponse,
ReactAppResponse,
TaskInstanceResponse,
@@ -27,78 +28,157 @@ import type {
export type PluginView = ExternalViewResponse | ReactAppResponse;
/**
- * The records the current route resolved to, used to evaluate `applies_to`.
+ * The records the current route resolved to, which `applies_to` paths are
evaluated against.
*
- * A field is `undefined` when the current destination has no such record
(e.g. no
- * `task` on a Dag-level page). `isLoading` is true while a record the route
*does*
- * have is still being fetched.
+ * A field is `undefined` when the route has no such record (e.g. no `task` on
a Dag-level
+ * page), which makes every path rooted at it unevaluable. `isLoading` is true
while a record
+ * the route *does* have is still being fetched.
*/
export type AppliesToContext = {
dag?: DAGResponse;
+ dagRun?: DAGRunResponse;
isLoading: boolean;
task?: TaskResponse;
taskInstance?: TaskInstanceResponse;
};
-const isNonEmpty = (value: Array<string> | null | undefined): value is
Array<string> =>
- value !== undefined && value !== null && value.length > 0;
+type RootName = "dag" | "dagRun" | "task" | "taskInstance";
+
+// A path may name a related record as its first segment. Keyed by the wire
spelling, since
+// that is what a plugin author writes.
+const ROOT_BY_PREFIX: Record<string, RootName> = {
+ dag: "dag",
+ dag_run: "dagRun",
+ task: "task",
+ task_instance: "taskInstance",
+};
-const knownNames = (...values: Array<string | null | undefined>):
Array<string> =>
- values.filter((value): value is string => value !== undefined && value !==
null && value !== "");
+// Which record an unqualified path is rooted at: the entity the destination
is about. A
+// destination missing here can evaluate nothing, so all of its paths are
skipped.
+// Kept in sync with `_APPLIES_TO_ENTITY_ROOT` in `airflow/plugins_manager.py`.
+const ENTITY_ROOT_BY_DESTINATION: Record<string, RootName> = {
+ dag: "dag",
+ dag_overview: "dag",
+ dag_run: "dagRun",
+ task: "task",
+ task_instance: "taskInstance",
+ task_overview: "task",
+};
+
+type Resolution = { found: false } | { found: true; values: Array<unknown> };
-const matchesAnyName = (criteria: Array<string>, names: Array<string>):
boolean | undefined =>
- names.length === 0 ? undefined : names.some((name) =>
criteria.includes(name));
+/**
+ * Walk a dotted path through the context's records.
+ *
+ * Traversing a list fans out across its elements, so `dag.tags.name` collects
every tag name --
+ * the general form of the hardcoded `dag.tags.some(...)` this replaces.
+ */
+const resolvePath = (
+ path: string,
+ destination: string | undefined,
+ context: AppliesToContext,
+): Resolution => {
+ const segments = path.split(".");
+ const [head, ...tail] = segments;
+ const prefixRoot = head === undefined ? undefined : ROOT_BY_PREFIX[head];
+ const isQualified = prefixRoot !== undefined && tail.length > 0;
+
+ const rootName = isQualified
+ ? prefixRoot
+ : destination === undefined
+ ? undefined
+ : ENTITY_ROOT_BY_DESTINATION[destination];
+ const root = rootName === undefined ? undefined : context[rootName];
+
+ if (root === undefined) {
+ return { found: false };
+ }
-// A criterion resolves to `undefined` when the current context cannot judge
it, which
-// is distinct from `false` (context available, nothing matched).
-const matchesDagTags = (criteria: Array<string>, { dag }: AppliesToContext):
boolean | undefined =>
- dag === undefined ? undefined : dag.tags.some((tag) =>
criteria.includes(tag.name));
+ let nodes: Array<unknown> = [root];
+
+ for (const segment of isQualified ? tail : segments) {
+ if (nodes.length === 0) {
+ // The previous segment was present but fanned out to nothing -- a Dag
with no tags, say.
+ // That is a definite "no values to match", not a path the page cannot
judge.
+ return { found: true, values: [] };
+ }
+
+ // `Object.hasOwn`, not `in`: `in` walks the prototype chain, so a path
segment naming an
+ // inherited member (`constructor`, `toString`) would read as a field the
record has.
+ const readable = nodes.filter(
+ (node): node is Record<string, unknown> =>
+ node !== null && typeof node === "object" && Object.hasOwn(node,
segment),
+ );
+
+ if (readable.length === 0) {
+ // The record has no such field. This is indistinguishable at runtime
from an author
+ // typo, so it has to be treated as unevaluable: `operator` exists on a
task instance but
+ // not on a task, and a block naming both sources must still work on
both pages. The
+ // cost is that a typo widens the scope rather than narrowing it.
+ return { found: false };
+ }
+
+ nodes = readable.flatMap((node) => {
+ const value = node[segment];
+
+ return Array.isArray(value) ? (value as Array<unknown>) : [value];
+ });
+ }
-const matchesDagIds = (criteria: Array<string>, { dag }: AppliesToContext):
boolean | undefined =>
- dag === undefined ? undefined : criteria.includes(dag.dag_id);
+ return { found: true, values: nodes };
+};
-const matchesTaskIds = (
- criteria: Array<string>,
- { task, taskInstance }: AppliesToContext,
-): boolean | undefined => {
- const taskId = taskInstance?.task_id ?? task?.task_id;
+// Stringified so a numeric field (`map_index`, `try_number`) or a boolean
(`is_paused`) is
+// addressable without the author having to think about JSON types. Anything
that is not a
+// primitive has no sensible string form: a path stopping on an object or
`null` yields no
+// comparable value, so it reads as "no match" rather than as unevaluable.
Stopping a path short
+// of a leaf therefore hides the view instead of widening it -- the opposite
of a bad segment.
+const toComparable = (value: unknown): string | undefined => {
+ if (typeof value === "string") {
+ return value;
+ }
- return taskId === undefined || taskId === null ? undefined :
criteria.includes(taskId);
+ return typeof value === "boolean" || typeof value === "number" ?
String(value) : undefined;
};
-// Matched separately, the way the task instance filters are: the class name
comes from
-// `task_type` and the display name from `custom_operator_name`. A decorated
task
-// (`@task.bash`) is reachable only by display name, since its class name is
private.
-const matchesOperators = (
- criteria: Array<string>,
- { task, taskInstance }: AppliesToContext,
-): boolean | undefined => {
- if (taskInstance !== undefined) {
- return matchesAnyName(criteria, knownNames(taskInstance.operator));
+// `undefined` means the page cannot judge this path, which is distinct from
`false` (records
+// available, nothing matched).
+const matchesPath = ({
+ allowed,
+ context,
+ destination,
+ path,
+}: {
+ allowed: Array<string>;
+ context: AppliesToContext;
+ destination: string | undefined;
+ path: string;
+}): boolean | undefined => {
+ const resolution = resolvePath(path, destination, context);
+
+ if (!resolution.found) {
+ return undefined;
}
- const classRef = task?.class_ref as { class_name?: string } | null |
undefined;
+ return resolution.values.some((value) => {
+ const comparable = toComparable(value);
- return matchesAnyName(criteria, knownNames(classRef?.class_name));
+ return comparable !== undefined && allowed.includes(comparable);
+ });
};
-const matchesOperatorNames = (
- criteria: Array<string>,
- { task, taskInstance }: AppliesToContext,
-): boolean | undefined =>
- taskInstance === undefined
- ? matchesAnyName(criteria, knownNames(task?.operator_name))
- : matchesAnyName(criteria, knownNames(taskInstance.operator_name));
+const isNonEmpty = (value: Array<string> | null | undefined): value is
Array<string> =>
+ value !== undefined && value !== null && value.length > 0;
/**
* Decide whether a plugin view should be shown for the current route.
*
- * Criteria are OR-ed within a key and AND-ed across keys, but only across
keys the
- * current destination can actually evaluate — a `task_ids` criterion cannot
be judged
- * on a Dag-level page, so it is skipped there rather than failing the match.
This lets
- * one `applies_to` block be shared by a plugin's Dag- and task-level
destinations.
+ * Values are OR-ed within a path and AND-ed across paths, but only across
paths whose root
+ * record the current destination actually has -- a `task_instance.*` path
cannot be judged on a
+ * Dag-level page, so it is skipped there rather than failing the match. That
is what lets one
+ * `applies_to` block be shared by a plugin's Dag- and task-level destinations.
*
- * Omitting `applies_to` (or giving it no criteria) shows the view everywhere.
+ * Omitting `applies_to` (or giving it no paths) shows the view everywhere.
*/
export const matchesAppliesTo = (view: PluginView, context: AppliesToContext):
boolean => {
const { applies_to: appliesTo } = view;
@@ -107,28 +187,18 @@ export const matchesAppliesTo = (view: PluginView,
context: AppliesToContext): b
return true;
}
- const {
- dag_ids: dagIds,
- dag_tags: dagTags,
- operator_names: operatorNames,
- operators,
- task_ids: taskIds,
- } = appliesTo;
-
- const verdicts = [
- isNonEmpty(dagTags) ? matchesDagTags(dagTags, context) : undefined,
- isNonEmpty(dagIds) ? matchesDagIds(dagIds, context) : undefined,
- isNonEmpty(taskIds) ? matchesTaskIds(taskIds, context) : undefined,
- isNonEmpty(operators) ? matchesOperators(operators, context) : undefined,
- isNonEmpty(operatorNames) ? matchesOperatorNames(operatorNames, context) :
undefined,
- ].filter((verdict) => verdict !== undefined);
-
- // No criterion was evaluable (either none configured, or none judgeable
here).
+ const verdicts = Object.entries(appliesTo)
+ .filter(([, allowed]) => isNonEmpty(allowed))
+ .map(([path, allowed]) => matchesPath({ allowed, context, destination:
view.destination, path }))
+ .filter((verdict) => verdict !== undefined);
+
+ // An empty list means no path was evaluable (either none configured, or
none judgeable
+ // here), and `every` is vacuously true for it -- so the view shows, which
is the default.
return verdicts.every(Boolean);
};
/**
- * Whether a view configures any scoping criterion at all.
+ * Whether a view configures any scoping at all.
*
* Callers use this to skip fetching the context records entirely when no view
needs them.
*/
@@ -139,10 +209,10 @@ export const hasAppliesToCriteria = (view: PluginView):
boolean => {
};
/**
- * Whether a view should be withheld while the records its criteria need are
in flight.
+ * Whether a view should be withheld while the records its paths need are in
flight.
*
- * Without this, a scoped view would render on first paint and disappear once
the
- * queries resolve. Unscoped views never wait.
+ * Without this, a scoped view would render on first paint and disappear once
the queries
+ * resolve. Unscoped views never wait.
*/
export const isAppliesToPending = (view: PluginView, context:
AppliesToContext): boolean =>
hasAppliesToCriteria(view) && context.isLoading;
diff --git
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_plugins.py
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_plugins.py
index 2290438faf9..ac58f76e334 100644
--- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_plugins.py
+++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_plugins.py
@@ -105,11 +105,8 @@ class TestGetPlugins:
"category": "browse",
"nav_top_level": False,
"applies_to": {
- "dag_tags": ["ml", "production"],
- "dag_ids": ["example_dag"],
- "task_ids": None,
- "operators": None,
- "operator_names": None,
+ "dag.tags.name": ["ml", "production"],
+ "dag.dag_id": ["example_dag"],
},
},
]
@@ -165,27 +162,28 @@ class TestGetPlugins:
assert view.applies_to is None
- def test_applies_to_parses_nested_criteria(self):
+ def test_applies_to_accepts_any_field_path(self):
from airflow.api_fastapi.core_api.datamodels.plugins import
ExternalViewResponse
view = ExternalViewResponse(
name="Scoped",
href="https://example.com/",
applies_to={
- "dag_tags": ["ml"],
- "operators": ["KubernetesPodOperator"],
- "operator_names": ["@task.bash"],
+ "dag.tags.name": ["ml"],
+ "operator_name": ["@task.bash"],
+ "state": ["failed", "upstream_failed"],
},
)
assert view.applies_to is not None
- assert view.applies_to.dag_tags == ["ml"]
- assert view.applies_to.operators == ["KubernetesPodOperator"]
- assert view.applies_to.operator_names == ["@task.bash"]
- assert view.applies_to.dag_ids is None
- assert view.applies_to.task_ids is None
+ # An open map: the key set is whatever the entity records expose, not
a fixed list.
+ assert view.applies_to.root == {
+ "dag.tags.name": ["ml"],
+ "operator_name": ["@task.bash"],
+ "state": ["failed", "upstream_failed"],
+ }
- def test_applies_to_rejects_unknown_criteria(self):
+ def test_applies_to_rejects_non_list_values(self):
from pydantic import ValidationError
from airflow.api_fastapi.core_api.datamodels.plugins import
ExternalViewResponse
@@ -194,7 +192,7 @@ class TestGetPlugins:
ExternalViewResponse(
name="Scoped",
href="https://example.com/",
- applies_to={"dag_tag": ["ml"]},
+ applies_to={"state": "failed"},
)
def
test_invalid_external_view_destination_should_log_warning_and_continue(self,
test_client, caplog):
diff --git a/airflow-core/tests/unit/cli/commands/test_plugins_command.py
b/airflow-core/tests/unit/cli/commands/test_plugins_command.py
index 00027fc26af..fbead422d62 100644
--- a/airflow-core/tests/unit/cli/commands/test_plugins_command.py
+++ b/airflow-core/tests/unit/cli/commands/test_plugins_command.py
@@ -118,8 +118,8 @@ class TestPluginsCommand:
"url_route": "test_iframe_plugin",
"category": "browse",
"applies_to": {
- "dag_tags": ["ml", "production"],
- "dag_ids": ["example_dag"],
+ "dag.tags.name": ["ml", "production"],
+ "dag.dag_id": ["example_dag"],
},
},
],
diff --git a/airflow-core/tests/unit/plugins/test_plugin.py
b/airflow-core/tests/unit/plugins/test_plugin.py
index 269a95eceed..1f22da97f75 100644
--- a/airflow-core/tests/unit/plugins/test_plugin.py
+++ b/airflow-core/tests/unit/plugins/test_plugin.py
@@ -136,8 +136,8 @@ external_view_with_metadata: ExternalViewDict = {
"destination": "dag",
"category": "browse",
"applies_to": {
- "dag_tags": ["ml", "production"],
- "dag_ids": ["example_dag"],
+ "dag.tags.name": ["ml", "production"],
+ "dag.dag_id": ["example_dag"],
},
}
diff --git a/airflow-core/tests/unit/plugins/test_plugins_manager.py
b/airflow-core/tests/unit/plugins/test_plugins_manager.py
index f4f361481cc..56dc6095e97 100644
--- a/airflow-core/tests/unit/plugins/test_plugins_manager.py
+++ b/airflow-core/tests/unit/plugins/test_plugins_manager.py
@@ -337,15 +337,14 @@ class TestPluginsManager:
id="not-a-dict",
),
pytest.param(
- {"dag_tag": ["ml"]},
- "unknown criteria ['dag_tag'], expected any of "
- "['dag_ids', 'dag_tags', 'operator_names', 'operators',
'task_ids']",
- id="unknown-key",
+ {1: ["ml"]},
+ "field paths must be strings, got [1]",
+ id="non-string-key",
),
pytest.param(
- {1: ["ml"], "dag_tag": ["ml"]},
- "criterion names must be strings, got [1]",
- id="non-string-and-unknown-keys",
+ {"": ["ml"]},
+ "field paths must not be empty, got ['']",
+ id="empty-key",
),
pytest.param(
{"dag_ids": "my_dag"},
@@ -404,7 +403,11 @@ class TestPluginsManager:
"bundle_url": "/scoped.js",
"url_route": "/scoped",
"destination": "dag_run",
- "applies_to": {"dag_tags": ["ml"], "task_ids": ["train"],
"operators": ["Op"]},
+ "applies_to": {
+ "dag.tags.name": ["ml"],
+ "task.class_ref.class_name": ["Op"],
+ "task_instance.operator": ["Op"],
+ },
}
]
@@ -416,15 +419,19 @@ class TestPluginsManager:
_, react_apps = plugins_manager._get_ui_plugins()
- # The block is only warned about, never stripped: task criteria
are ignored on a
- # Dag-level page by design so one block can be shared across
destinations.
+ # The block is only warned about, never stripped: a path whose
root the page lacks
+ # is skipped by design, so one block can be shared across
destinations.
assert react_apps == [
{
"name": "Scoped",
"bundle_url": "/scoped.js",
"url_route": "/scoped",
"destination": "dag_run",
- "applies_to": {"dag_tags": ["ml"], "task_ids": ["train"],
"operators": ["Op"]},
+ "applies_to": {
+ "dag.tags.name": ["ml"],
+ "task.class_ref.class_name": ["Op"],
+ "task_instance.operator": ["Op"],
+ },
}
]
@@ -433,10 +440,187 @@ class TestPluginsManager:
"airflow.plugins_manager",
logging.WARNING,
"Plugin 'test_plugin' has a React App 'Scoped' with
destination 'dag_run', which cannot "
- "evaluate ['operators', 'task_ids']. Those criteria will be
ignored.",
+ "evaluate ['task.class_ref.class_name',
'task_instance.operator']. Those paths will be "
+ "ignored.",
),
]
+ @pytest.mark.parametrize(
+ ("path", "error"),
+ [
+ pytest.param(
+ "dag.tags.nme",
+ "'dag.tags.nme' names no field 'nme' on DagTagResponse (did
you mean 'name'?)",
+ id="misspelled-nested-leaf",
+ ),
+ pytest.param(
+ "stat",
+ "'stat' names no field 'stat' on DAGRunResponse (did you mean
'state'?)",
+ id="misspelled-own-field",
+ ),
+ pytest.param(
+ "state.length",
+ "'state.length' reads 'length' from DagRunState, which has no
fields",
+ id="path-past-a-scalar",
+ ),
+ ],
+ )
+ def test_warns_about_a_path_matching_no_field(self, path, error, caplog):
+ class TestPlugin(AirflowPlugin):
+ name = "test_plugin"
+
+ external_views = [
+ {
+ "name": "Scoped",
+ "href": "/scoped",
+ "url_route": "/scoped",
+ "destination": "dag_run",
+ "applies_to": {path: ["x"]},
+ }
+ ]
+
+ with (
+ mock_plugin_manager(plugins=[TestPlugin()]),
+ caplog.at_level(logging.WARNING, logger="airflow.plugins_manager"),
+ ):
+ from airflow import plugins_manager
+
+ external_views, _ = plugins_manager._get_ui_plugins()
+
+ # Warned about, not stripped: the path is inert either way, and
dropping it would
+ # change the block the UI receives.
+ assert external_views[0]["applies_to"] == {path: ["x"]}
+
+ assert caplog.record_tuples == [
+ (
+ "airflow.plugins_manager",
+ logging.WARNING,
+ f"Plugin 'test_plugin' has an external view 'Scoped' with an
'applies_to' path that "
+ f"matches no field: {error}. That path will be ignored, so the
view will appear in "
+ f"more places than intended.",
+ ),
+ ]
+
+ @pytest.mark.parametrize(
+ ("destination", "path"),
+ [
+ # `TaskInstanceResponse.run_id` is serialized as `dag_run_id`; the
browser only ever
+ # sees the alias.
+ pytest.param("task_instance", "dag_run_id",
id="serialization-alias"),
+ pytest.param("task_instance", "queued_when",
id="serialization-alias-datetime"),
+ # A computed field has no `model_fields` entry but is in the
response.
+ pytest.param("dag", "is_backfillable", id="computed-field"),
+ ],
+ )
+ def test_accepts_fields_as_the_api_serializes_them(self, destination,
path, caplog):
+ class TestPlugin(AirflowPlugin):
+ name = "test_plugin"
+
+ external_views = [
+ {
+ "name": "Scoped",
+ "href": "/scoped",
+ "url_route": "/scoped",
+ "destination": destination,
+ "applies_to": {path: ["x"]},
+ }
+ ]
+
+ with (
+ mock_plugin_manager(plugins=[TestPlugin()]),
+ caplog.at_level(logging.WARNING, logger="airflow.plugins_manager"),
+ ):
+ from airflow import plugins_manager
+
+ plugins_manager._get_ui_plugins()
+
+ assert caplog.record_tuples == []
+
+ def test_rejects_a_python_attribute_name_that_is_not_in_the_response(self,
caplog):
+ """`run_id` is the Python attribute; the response carries
`dag_run_id`."""
+
+ class TestPlugin(AirflowPlugin):
+ name = "test_plugin"
+
+ external_views = [
+ {
+ "name": "Scoped",
+ "href": "/scoped",
+ "url_route": "/scoped",
+ "destination": "task_instance",
+ "applies_to": {"run_id": ["manual__1"]},
+ }
+ ]
+
+ with (
+ mock_plugin_manager(plugins=[TestPlugin()]),
+ caplog.at_level(logging.WARNING, logger="airflow.plugins_manager"),
+ ):
+ from airflow import plugins_manager
+
+ plugins_manager._get_ui_plugins()
+
+ assert len(caplog.record_tuples) == 1
+ assert "names no field 'run_id' on TaskInstanceResponse" in
caplog.record_tuples[0][2]
+
+ def test_drops_a_null_valued_path_keeping_the_rest_of_the_block(self):
+ """A null value is not a ``list[str]``; leaving it in drops the whole
plugin."""
+
+ class TestPlugin(AirflowPlugin):
+ name = "test_plugin"
+
+ external_views = [
+ {
+ "name": "Scoped",
+ "href": "/scoped",
+ "url_route": "/scoped",
+ "destination": "dag_run",
+ "applies_to": {"state": None, "dag.tags.name": ["ml"]},
+ }
+ ]
+
+ with mock_plugin_manager(plugins=[TestPlugin()]):
+ from airflow import plugins_manager
+
+ external_views, _ = plugins_manager._get_ui_plugins()
+
+ assert external_views[0]["applies_to"] == {"dag.tags.name": ["ml"]}
+
+ # The point of dropping it: the block still serializes, so the
view survives.
+ from airflow.api_fastapi.core_api.datamodels.plugins import
ExternalViewResponse
+
+ assert ExternalViewResponse(**external_views[0]).applies_to.root
== {"dag.tags.name": ["ml"]}
+
+ def test_accepts_a_path_through_a_field_the_models_do_not_describe(self,
caplog):
+ """A path cannot be checked past a bare ``dict``, so everything below
it is accepted."""
+
+ class TestPlugin(AirflowPlugin):
+ name = "test_plugin"
+
+ external_views = [
+ {
+ "name": "Scoped",
+ "href": "/scoped",
+ "url_route": "/scoped",
+ "destination": "task_instance",
+ # `TaskResponse.class_ref` is a bare dict, and `conf` is
dict[str, Any].
+ "applies_to": {
+ "task.class_ref.class_name": ["KubernetesPodOperator"],
+ "dag_run.conf.environment": ["prod"],
+ },
+ }
+ ]
+
+ with (
+ mock_plugin_manager(plugins=[TestPlugin()]),
+ caplog.at_level(logging.WARNING, logger="airflow.plugins_manager"),
+ ):
+ from airflow import plugins_manager
+
+ plugins_manager._get_ui_plugins()
+
+ assert caplog.record_tuples == []
+
def test_does_not_warn_about_valid_applies_to(self, caplog):
class TestPlugin(AirflowPlugin):
name = "test_plugin"
@@ -448,9 +632,9 @@ class TestPluginsManager:
"url_route": "/scoped",
"destination": "task",
"applies_to": {
- "dag_tags": ["ml"],
- "task_ids": ["train"],
- "operator_names": ["@task.bash"],
+ "dag.tags.name": ["ml"],
+ "operator_name": ["@task.bash"],
+ "task_id": ["train"],
},
}
]
@@ -464,9 +648,9 @@ class TestPluginsManager:
external_views, _ = plugins_manager._get_ui_plugins()
assert external_views[0]["applies_to"] == {
- "dag_tags": ["ml"],
- "task_ids": ["train"],
- "operator_names": ["@task.bash"],
+ "dag.tags.name": ["ml"],
+ "operator_name": ["@task.bash"],
+ "task_id": ["train"],
}
assert caplog.record_tuples == []
diff --git a/airflow-ctl/src/airflowctl/api/datamodels/generated.py
b/airflow-ctl/src/airflowctl/api/datamodels/generated.py
index 1908f0ce980..df3c9b5dbf9 100644
--- a/airflow-ctl/src/airflowctl/api/datamodels/generated.py
+++ b/airflow-ctl/src/airflowctl/api/datamodels/generated.py
@@ -1093,19 +1093,19 @@ class NewTaskResponse(BaseModel):
task_display_name: Annotated[str, Field(title="Task Display Name")]
-class PluginAppliesToResponse(BaseModel):
+class PluginAppliesToResponse(RootModel[dict[str, list[str]]]):
"""
- Serializer for the optional Dag/task scoping criteria of a UI plugin.
+ Serializer for the optional scoping criteria of a UI plugin.
+
+ An open map of dotted field path to the values that path may take -- not a
closed set of
+ criteria. ``{"state": ["failed"], "dag.tags.name": ["ml"]}`` scopes to
failed entities of
+ ml-tagged Dags. An unqualified path is rooted at the entity the
``destination`` is about;
+ a path may instead name a related record (``dag``, ``dag_run``, ``task``,
``task_instance``)
+ as its first segment. Matching is equality against the listed values, OR
within a path and
+ AND across paths, and is evaluated client-side.
"""
- model_config = ConfigDict(
- extra="forbid",
- )
- dag_tags: Annotated[list[str] | None, Field(title="Dag Tags")] = None
- dag_ids: Annotated[list[str] | None, Field(title="Dag Ids")] = None
- task_ids: Annotated[list[str] | None, Field(title="Task Ids")] = None
- operators: Annotated[list[str] | None, Field(title="Operators")] = None
- operator_names: Annotated[list[str] | None, Field(title="Operator Names")]
= None
+ root: dict[str, list[str]]
class PluginImportErrorResponse(BaseModel):
diff --git
a/shared/plugins_manager/src/airflow_shared/plugins_manager/plugins_manager.py
b/shared/plugins_manager/src/airflow_shared/plugins_manager/plugins_manager.py
index 58f9c043670..6f8cf5c495f 100644
---
a/shared/plugins_manager/src/airflow_shared/plugins_manager/plugins_manager.py
+++
b/shared/plugins_manager/src/airflow_shared/plugins_manager/plugins_manager.py
@@ -89,14 +89,10 @@ class AirflowPluginException(Exception):
BaseDestinationLiteral = Literal["nav", "dag", "dag_run", "task",
"task_instance", "asset", "base"]
-class AppliesToDict(TypedDict):
- """Dictionary structure for the optional ``applies_to`` scoping block on
UI plugins."""
-
- dag_tags: NotRequired[list[str] | None]
- dag_ids: NotRequired[list[str] | None]
- task_ids: NotRequired[list[str] | None]
- operators: NotRequired[list[str] | None]
- operator_names: NotRequired[list[str] | None]
+# The optional ``applies_to`` scoping block: an open map of dotted field path
to the values
+# that path may take, e.g. ``{"state": ["failed"], "dag.tags.name": ["ml"]}``.
Not a closed set
+# of criteria -- any field of the records the destination resolves is
addressable.
+AppliesToDict = dict[str, list[str]]
class _BaseUIDict(TypedDict):