From d93bdab754d2cccef91f52fdea389d1420d95da5 Mon Sep 17 00:00:00 2001 From: AutomationDev85 Date: Fri, 22 May 2026 14:30:29 +0200 Subject: [PATCH 1/9] Bring back edge worker DualStatusManager compability --- .../edge3/src/airflow/providers/edge3/models/edge_worker.py | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/providers/edge3/src/airflow/providers/edge3/models/edge_worker.py b/providers/edge3/src/airflow/providers/edge3/models/edge_worker.py index fb800b8b2457b..58b7361aa2e78 100644 --- a/providers/edge3/src/airflow/providers/edge3/models/edge_worker.py +++ b/providers/edge3/src/airflow/providers/edge3/models/edge_worker.py @@ -34,6 +34,11 @@ from airflow.utils.session import NEW_SESSION, provide_session from airflow.utils.sqlalchemy import UtcDateTime +try: + from airflow.sdk.observability.stats import DualStatsManager +except ImportError: + DualStatsManager = None # type: ignore[assignment,misc] # Airflow < 3.2 compat + if TYPE_CHECKING: from collections.abc import Sequence From b0b6a1484616c14610379dc623ab6a70c8c5ac07 Mon Sep 17 00:00:00 2001 From: AutomationDev85 Date: Mon, 8 Jun 2026 13:33:19 +0200 Subject: [PATCH 2/9] Reworked statsd taging for pre airflow 3.3 versions --- .../providers/edge3/models/edge_worker.py | 38 ++++++++++++++++--- 1 file changed, 33 insertions(+), 5 deletions(-) diff --git a/providers/edge3/src/airflow/providers/edge3/models/edge_worker.py b/providers/edge3/src/airflow/providers/edge3/models/edge_worker.py index 58b7361aa2e78..d662a45a6d62b 100644 --- a/providers/edge3/src/airflow/providers/edge3/models/edge_worker.py +++ b/providers/edge3/src/airflow/providers/edge3/models/edge_worker.py @@ -28,17 +28,16 @@ from airflow.providers.common.compat.sdk import AirflowException, Stats, timezone from airflow.providers.common.compat.sqlalchemy.orm import mapped_column from airflow.providers.edge3.models.edge_base import Base +<<<<<<< HEAD from airflow.utils.helpers import prune_dict +======= +from airflow.providers.edge3.version_compat import AIRFLOW_V_3_3_PLUS +>>>>>>> 366552b881 (Reworked statsd taging for pre airflow 3.3 versions) from airflow.utils.log.logging_mixin import LoggingMixin from airflow.utils.providers_configuration_loader import providers_configuration_loaded from airflow.utils.session import NEW_SESSION, provide_session from airflow.utils.sqlalchemy import UtcDateTime -try: - from airflow.sdk.observability.stats import DualStatsManager -except ImportError: - DualStatsManager = None # type: ignore[assignment,misc] # Airflow < 3.2 compat - if TYPE_CHECKING: from collections.abc import Sequence @@ -188,6 +187,7 @@ def set_metrics( metric_tags = prune_dict({"worker_name": worker_name, "team_name": team_name}) Stats.gauge( +<<<<<<< HEAD "edge_worker.status", sysinfo.get("status", logging.NOTSET), # type: ignore tags=metric_tags, @@ -201,6 +201,19 @@ def set_metrics( "edge_worker.num_queues", len(queues), tags={**metric_tags, "queues": ",".join(queues)}, +======= + "edge_worker.status", sysinfo.get("status", logging.NOTSET), tags={"worker_name": worker_name} + ) # type: ignore + Stats.gauge("edge_worker.connected", int(connected), tags={"worker_name": worker_name}) + Stats.gauge("edge_worker.maintenance", int(maintenance), tags={"worker_name": worker_name}) + Stats.gauge("edge_worker.jobs_active", jobs_active, tags={"worker_name": worker_name}) + Stats.gauge("edge_worker.concurrency", concurrency, tags={"worker_name": worker_name}) + Stats.gauge("edge_worker.free_concurrency", free_concurrency, tags={"worker_name": worker_name}) + Stats.gauge( + "edge_worker.num_queues", + len(queues), + tags={"worker_name": worker_name, "queues": ",".join(queues)}, +>>>>>>> 366552b881 (Reworked statsd taging for pre airflow 3.3 versions) ) for key in additional_keys: @@ -208,6 +221,21 @@ def set_metrics( if isinstance(value, (int, float)): Stats.gauge(f"edge_worker.{key}", value, tags=metric_tags) + if not AIRFLOW_V_3_3_PLUS: + # Airflow < 3.3: export legacy per-worker metrics (no auto-tag expansion). + Stats.gauge(f"edge_worker.status.{worker_name}", sysinfo.get("status", logging.NOTSET)) # type: ignore + Stats.gauge(f"edge_worker.connected.{worker_name}", int(connected)) + Stats.gauge(f"edge_worker.maintenance.{worker_name}", int(maintenance)) + Stats.gauge(f"edge_worker.jobs_active.{worker_name}", jobs_active) + Stats.gauge(f"edge_worker.concurrency.{worker_name}", concurrency) + Stats.gauge(f"edge_worker.free_concurrency.{worker_name}", free_concurrency) + Stats.gauge(f"edge_worker.num_queues.{worker_name}", len(queues)) + + for key in additional_keys: + value = sysinfo.get(key) + if isinstance(value, (int, float)): + Stats.gauge(f"edge_worker.{key}.{worker_name}", value) + def reset_metrics(worker_name: str, team_name: str | None = None) -> None: """Reset metrics of worker.""" From 2bbe36b16bc63bee60bdd1d6675d5dc44ae4c4f9 Mon Sep 17 00:00:00 2001 From: AutomationDev85 Date: Mon, 8 Jun 2026 16:15:17 +0200 Subject: [PATCH 3/9] Fix mypy --- .../providers/edge3/models/edge_worker.py | 25 +++++-------------- 1 file changed, 6 insertions(+), 19 deletions(-) diff --git a/providers/edge3/src/airflow/providers/edge3/models/edge_worker.py b/providers/edge3/src/airflow/providers/edge3/models/edge_worker.py index d662a45a6d62b..db459391239e5 100644 --- a/providers/edge3/src/airflow/providers/edge3/models/edge_worker.py +++ b/providers/edge3/src/airflow/providers/edge3/models/edge_worker.py @@ -185,13 +185,13 @@ def set_metrics( "free_concurrency", } metric_tags = prune_dict({"worker_name": worker_name, "team_name": team_name}) + status = sysinfo.get("status", logging.NOTSET) + if not isinstance(status, (int, float)): + status = logging.NOTSET Stats.gauge( -<<<<<<< HEAD - "edge_worker.status", - sysinfo.get("status", logging.NOTSET), # type: ignore - tags=metric_tags, - ) + "edge_worker.status", status, tags=metric_tags + ) # type: ignore Stats.gauge("edge_worker.connected", int(connected), tags=metric_tags) Stats.gauge("edge_worker.maintenance", int(maintenance), tags=metric_tags) Stats.gauge("edge_worker.jobs_active", jobs_active, tags=metric_tags) @@ -201,19 +201,6 @@ def set_metrics( "edge_worker.num_queues", len(queues), tags={**metric_tags, "queues": ",".join(queues)}, -======= - "edge_worker.status", sysinfo.get("status", logging.NOTSET), tags={"worker_name": worker_name} - ) # type: ignore - Stats.gauge("edge_worker.connected", int(connected), tags={"worker_name": worker_name}) - Stats.gauge("edge_worker.maintenance", int(maintenance), tags={"worker_name": worker_name}) - Stats.gauge("edge_worker.jobs_active", jobs_active, tags={"worker_name": worker_name}) - Stats.gauge("edge_worker.concurrency", concurrency, tags={"worker_name": worker_name}) - Stats.gauge("edge_worker.free_concurrency", free_concurrency, tags={"worker_name": worker_name}) - Stats.gauge( - "edge_worker.num_queues", - len(queues), - tags={"worker_name": worker_name, "queues": ",".join(queues)}, ->>>>>>> 366552b881 (Reworked statsd taging for pre airflow 3.3 versions) ) for key in additional_keys: @@ -223,7 +210,7 @@ def set_metrics( if not AIRFLOW_V_3_3_PLUS: # Airflow < 3.3: export legacy per-worker metrics (no auto-tag expansion). - Stats.gauge(f"edge_worker.status.{worker_name}", sysinfo.get("status", logging.NOTSET)) # type: ignore + Stats.gauge(f"edge_worker.status.{worker_name}", int(status)) Stats.gauge(f"edge_worker.connected.{worker_name}", int(connected)) Stats.gauge(f"edge_worker.maintenance.{worker_name}", int(maintenance)) Stats.gauge(f"edge_worker.jobs_active.{worker_name}", jobs_active) From 1c28cefe325d216f89d9a30b62286b5b4f0e1985 Mon Sep 17 00:00:00 2001 From: AutomationDev85 Date: Tue, 9 Jun 2026 09:53:12 +0200 Subject: [PATCH 4/9] document metrics compability issues --- providers/edge3/docs/edge_executor.rst | 27 ++++++++++++++++++++++++++ 1 file changed, 27 insertions(+) diff --git a/providers/edge3/docs/edge_executor.rst b/providers/edge3/docs/edge_executor.rst index 4452b86958966..673251d025ac4 100644 --- a/providers/edge3/docs/edge_executor.rst +++ b/providers/edge3/docs/edge_executor.rst @@ -217,3 +217,30 @@ Current Limitations Edge Executor - Multi-team isolation is logical only — all teams share a single authentication secret. A worker administrator could change the team name and access another team's jobs. See :ref:`edge_executor:multi_team` for details and planned improvements. + + +Metrics Export Compatibility +----------------------------- + +The Edge Worker integrates with Airflow's statsd metrics system to export runtime metrics. Compatibility +between Edge provider versions and Airflow versions varies due to changes in the metrics initialization +pipeline. The table below documents known compatibility issues and workarounds: + +.. list-table:: + :header-rows: 1 + + * - Provider version + - Airflow 3.3 + - Airflow 3.2 + - Airflow 3.1 + * - >= 3.6.0 + - Working + - Broken: webserver statsd initialization is missing (DualStatsManager removed) + - Broken: statsd metric tags + * - <= 3.5.0 + - Working + - Requires manual patch of ``metrics_template.yaml`` to match previous statsd export schema + - Working + +**For Airflow 3.2 users:** If upgrading to Edge provider >= 3.6.0 breaks metrics export, either +(1) upgrade Airflow to 3.3+, or (2) downgrade to Edge provider <= 3.5.0 with the workaround above. From 739c3ab5a0c89a6cdf1f6327ff72ea48bb9d52fd Mon Sep 17 00:00:00 2001 From: AutomationDev85 Date: Tue, 9 Jun 2026 13:13:09 +0200 Subject: [PATCH 5/9] fix doc --- providers/edge3/docs/edge_executor.rst | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/providers/edge3/docs/edge_executor.rst b/providers/edge3/docs/edge_executor.rst index 673251d025ac4..9c05997f2de79 100644 --- a/providers/edge3/docs/edge_executor.rst +++ b/providers/edge3/docs/edge_executor.rst @@ -222,7 +222,7 @@ Current Limitations Edge Executor Metrics Export Compatibility ----------------------------- -The Edge Worker integrates with Airflow's statsd metrics system to export runtime metrics. Compatibility +The Edge Worker integrates with Airflow's metrics system to export runtime metrics. Compatibility between Edge provider versions and Airflow versions varies due to changes in the metrics initialization pipeline. The table below documents known compatibility issues and workarounds: From 5450f8a38b477c5c287b7f04029cf32e64d9765b Mon Sep 17 00:00:00 2001 From: AutomationDev85 Date: Tue, 23 Jun 2026 11:26:35 +0200 Subject: [PATCH 6/9] Fix docu and added unit test --- providers/edge3/docs/edge_executor.rst | 6 +- .../unit/edge3/models/test_edge_worker.py | 64 +++++++++++++++++++ 2 files changed, 67 insertions(+), 3 deletions(-) create mode 100644 providers/edge3/tests/unit/edge3/models/test_edge_worker.py diff --git a/providers/edge3/docs/edge_executor.rst b/providers/edge3/docs/edge_executor.rst index 9c05997f2de79..b29dafbe3afbd 100644 --- a/providers/edge3/docs/edge_executor.rst +++ b/providers/edge3/docs/edge_executor.rst @@ -235,11 +235,11 @@ pipeline. The table below documents known compatibility issues and workarounds: - Airflow 3.1 * - >= 3.6.0 - Working - - Broken: webserver statsd initialization is missing (DualStatsManager removed) - - Broken: statsd metric tags + - Broken: webserver missing a `stats.initialize(...)` call -- DualStatsManager removed, missing legacy metric names + - Broken: metric tags * - <= 3.5.0 - Working - - Requires manual patch of ``metrics_template.yaml`` to match previous statsd export schema + - Requires manual patch of ``metrics_template.yaml`` to match previous export schema - Working **For Airflow 3.2 users:** If upgrading to Edge provider >= 3.6.0 breaks metrics export, either diff --git a/providers/edge3/tests/unit/edge3/models/test_edge_worker.py b/providers/edge3/tests/unit/edge3/models/test_edge_worker.py new file mode 100644 index 0000000000000..05c6e16a3e33b --- /dev/null +++ b/providers/edge3/tests/unit/edge3/models/test_edge_worker.py @@ -0,0 +1,64 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. +from __future__ import annotations + +from unittest import mock + +from airflow.providers.common.compat.sdk import Stats +from airflow.providers.edge3.models.edge_worker import EdgeWorkerState, set_metrics + +from tests_common.test_utils.version_compat import AIRFLOW_V_3_3_PLUS + +stats_reference = f"{Stats.__module__}.Stats" + + +def test_set_metrics(): + worker_name = "test_worker1" + if AIRFLOW_V_3_3_PLUS: + with mock.patch("airflow.sdk._shared.observability.metrics.stats._get_backend") as mock_get_backend: + mock_backend = mock.MagicMock() + mock_get_backend.return_value = mock_backend + + set_metrics( + worker_name=worker_name, + state=EdgeWorkerState.IDLE, + jobs_active=0, + concurrency=1, + free_concurrency=1, + queues=None, + sysinfo={"status": 1}, + ) + + metric_names = [call.args[0] for call in mock_backend.gauge.call_args_list] + else: + with mock.patch(f"{stats_reference}.gauge") as mock_gauge: + set_metrics( + worker_name=worker_name, + state=EdgeWorkerState.IDLE, + jobs_active=0, + concurrency=1, + free_concurrency=1, + queues=None, + sysinfo={"status": 1}, + ) + + metric_names = [call.args[0] for call in mock_gauge.call_args_list] + + assert "edge_worker.status" in metric_names + + legacy_metric_name = f"edge_worker.status.{worker_name}" + assert legacy_metric_name in metric_names From aef0652ff46562370a29a7b6bf67c225545395c7 Mon Sep 17 00:00:00 2001 From: AutomationDev85 Date: Mon, 13 Jul 2026 12:46:12 +0200 Subject: [PATCH 7/9] Fix imports --- .../src/airflow/providers/edge3/models/edge_worker.py | 9 ++------- 1 file changed, 2 insertions(+), 7 deletions(-) diff --git a/providers/edge3/src/airflow/providers/edge3/models/edge_worker.py b/providers/edge3/src/airflow/providers/edge3/models/edge_worker.py index db459391239e5..ffc6e20383cde 100644 --- a/providers/edge3/src/airflow/providers/edge3/models/edge_worker.py +++ b/providers/edge3/src/airflow/providers/edge3/models/edge_worker.py @@ -28,11 +28,8 @@ from airflow.providers.common.compat.sdk import AirflowException, Stats, timezone from airflow.providers.common.compat.sqlalchemy.orm import mapped_column from airflow.providers.edge3.models.edge_base import Base -<<<<<<< HEAD -from airflow.utils.helpers import prune_dict -======= from airflow.providers.edge3.version_compat import AIRFLOW_V_3_3_PLUS ->>>>>>> 366552b881 (Reworked statsd taging for pre airflow 3.3 versions) +from airflow.utils.helpers import prune_dict from airflow.utils.log.logging_mixin import LoggingMixin from airflow.utils.providers_configuration_loader import providers_configuration_loaded from airflow.utils.session import NEW_SESSION, provide_session @@ -189,9 +186,7 @@ def set_metrics( if not isinstance(status, (int, float)): status = logging.NOTSET - Stats.gauge( - "edge_worker.status", status, tags=metric_tags - ) # type: ignore + Stats.gauge("edge_worker.status", status, tags=metric_tags) Stats.gauge("edge_worker.connected", int(connected), tags=metric_tags) Stats.gauge("edge_worker.maintenance", int(maintenance), tags=metric_tags) Stats.gauge("edge_worker.jobs_active", jobs_active, tags=metric_tags) From 88930521af7abe3bbf56d0ada5128ba3eecacf94 Mon Sep 17 00:00:00 2001 From: AutomationDev85 Date: Mon, 13 Jul 2026 13:29:06 +0200 Subject: [PATCH 8/9] Fix unit test --- .../unit/edge3/worker_api/routes/test_worker.py | 17 ++++++++++++++++- 1 file changed, 16 insertions(+), 1 deletion(-) diff --git a/providers/edge3/tests/unit/edge3/worker_api/routes/test_worker.py b/providers/edge3/tests/unit/edge3/worker_api/routes/test_worker.py index f42d93950f3b7..939cf5d5f54ff 100644 --- a/providers/edge3/tests/unit/edge3/worker_api/routes/test_worker.py +++ b/providers/edge3/tests/unit/edge3/worker_api/routes/test_worker.py @@ -45,6 +45,7 @@ ) from tests_common.test_utils.config import conf_vars +from tests_common.test_utils.version_compat import AIRFLOW_V_3_3_PLUS if TYPE_CHECKING: from sqlalchemy.orm import Session @@ -423,7 +424,21 @@ def test_set_state_metrics_team_name_tags( tags={**expected_worker_tags, "queues": ",".join(queues)}, ) mock_stats_gauge.assert_any_call("edge_worker.disk_usage", 42.5, tags=expected_worker_tags) - assert mock_stats_gauge.call_count == 8 + if AIRFLOW_V_3_3_PLUS: + assert mock_stats_gauge.call_count == 8 + else: + mock_stats_gauge.assert_any_call( + "edge_worker.status.test2_worker", + self.MOCK_SYSINFO["status"], + ) + mock_stats_gauge.assert_any_call("edge_worker.connected.test2_worker", 1) + mock_stats_gauge.assert_any_call("edge_worker.maintenance.test2_worker", 0) + mock_stats_gauge.assert_any_call("edge_worker.jobs_active.test2_worker", 1) + mock_stats_gauge.assert_any_call("edge_worker.concurrency.test2_worker", 8) + mock_stats_gauge.assert_any_call("edge_worker.free_concurrency.test2_worker", 8) + mock_stats_gauge.assert_any_call("edge_worker.num_queues.test2_worker", len(queues)) + mock_stats_gauge.assert_any_call("edge_worker.disk_usage.test2_worker", 42.5) + assert mock_stats_gauge.call_count == 16 def test_set_state_returns_concurrency(self, session: Session, cli_worker: EdgeWorker): """set_state includes the DB-stored concurrency override in its response.""" From aa7509104b68f312f165f811e3f574fda502f3b4 Mon Sep 17 00:00:00 2001 From: AutomationDev85 Date: Tue, 14 Jul 2026 15:51:37 +0200 Subject: [PATCH 9/9] Fix flaky unit tests by reset after test execution --- providers/edge3/tests/unit/edge3/cli/test_worker.py | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/providers/edge3/tests/unit/edge3/cli/test_worker.py b/providers/edge3/tests/unit/edge3/cli/test_worker.py index bbbf2690faf4a..716afe024022e 100644 --- a/providers/edge3/tests/unit/edge3/cli/test_worker.py +++ b/providers/edge3/tests/unit/edge3/cli/test_worker.py @@ -164,6 +164,13 @@ def setup_parser(self): importlib.reload(cli_parser) self.parser = cli_parser.get_parser() + @pytest.fixture(autouse=True) + def reset_edge_worker_class_attrs(self): + yield + EdgeWorker.jobs = [] + EdgeWorker.drain = False + EdgeWorker.maintenance_mode = False + @pytest.fixture def cli_worker_with_team(self, tmp_path: Path) -> EdgeWorker: test_worker = EdgeWorker(str(tmp_path / "mock.pid"), "mock", None, 8, team_name="team_a")