From 103658314386d1b9de2623bcf656560610698b61 Mon Sep 17 00:00:00 2001 From: Andrew Chang Date: Thu, 23 Jul 2026 21:30:00 +0800 Subject: [PATCH] Add WasbRemoteLogIO.from_config and register wasb remote logging scheme Core resolves remote log handlers by URL scheme through ProvidersManager dispatch (#67056); s3 and cloudwatch already migrated. This moves wasb onto the same path, so Azure Blob remote logging is built by the provider's from_config() instead of the hardcoded branch in airflow_local_settings.py. Existing wasb:// configs resolve to an equivalent handler, and a from_config failure falls back to the legacy path, so behaviour is unchanged. Part of #70265. closes #70268. --- providers/microsoft/azure/provider.yaml | 4 + .../microsoft/azure/get_provider_info.py | 6 ++ .../microsoft/azure/log/wasb_task_handler.py | 29 +++++++ .../azure/log/test_wasb_task_handler.py | 83 +++++++++++++++++++ 4 files changed, 122 insertions(+) diff --git a/providers/microsoft/azure/provider.yaml b/providers/microsoft/azure/provider.yaml index d4a5e84a2ff59..d1a471c39951c 100644 --- a/providers/microsoft/azure/provider.yaml +++ b/providers/microsoft/azure/provider.yaml @@ -988,6 +988,10 @@ secrets-backends: logging: - airflow.providers.microsoft.azure.log.wasb_task_handler.WasbTaskHandler +remote-logging: + - classpath: airflow.providers.microsoft.azure.log.wasb_task_handler.WasbRemoteLogIO + scheme: wasb + extra-links: - airflow.providers.microsoft.azure.operators.data_factory.AzureDataFactoryPipelineRunLink - airflow.providers.microsoft.azure.operators.synapse.AzureSynapsePipelineRunLink diff --git a/providers/microsoft/azure/src/airflow/providers/microsoft/azure/get_provider_info.py b/providers/microsoft/azure/src/airflow/providers/microsoft/azure/get_provider_info.py index bf7ab5e96a808..34509081d0eb4 100644 --- a/providers/microsoft/azure/src/airflow/providers/microsoft/azure/get_provider_info.py +++ b/providers/microsoft/azure/src/airflow/providers/microsoft/azure/get_provider_info.py @@ -957,6 +957,12 @@ def get_provider_info(): ], "secrets-backends": ["airflow.providers.microsoft.azure.secrets.key_vault.AzureKeyVaultBackend"], "logging": ["airflow.providers.microsoft.azure.log.wasb_task_handler.WasbTaskHandler"], + "remote-logging": [ + { + "classpath": "airflow.providers.microsoft.azure.log.wasb_task_handler.WasbRemoteLogIO", + "scheme": "wasb", + } + ], "extra-links": [ "airflow.providers.microsoft.azure.operators.data_factory.AzureDataFactoryPipelineRunLink", "airflow.providers.microsoft.azure.operators.synapse.AzureSynapsePipelineRunLink", diff --git a/providers/microsoft/azure/src/airflow/providers/microsoft/azure/log/wasb_task_handler.py b/providers/microsoft/azure/src/airflow/providers/microsoft/azure/log/wasb_task_handler.py index 5c631f4d61248..51b7f7813c2d9 100644 --- a/providers/microsoft/azure/src/airflow/providers/microsoft/azure/log/wasb_task_handler.py +++ b/providers/microsoft/azure/src/airflow/providers/microsoft/azure/log/wasb_task_handler.py @@ -17,6 +17,7 @@ # under the License. from __future__ import annotations +import inspect import os import shutil from functools import cached_property @@ -50,6 +51,34 @@ class WasbRemoteLogIO(LoggingMixin): # noqa: D101 processors = () + @classmethod + def from_config(cls) -> WasbRemoteLogIO: + """Build the remote log IO from Airflow logging configuration.""" + remote_task_handler_kwargs = conf.getjson("logging", "remote_task_handler_kwargs", fallback={}) + if not isinstance(remote_task_handler_kwargs, dict): + raise ValueError( + "logging/remote_task_handler_kwargs must be a JSON object (a python dict), we got " + f"{type(remote_task_handler_kwargs)}" + ) + fth_params = frozenset(inspect.signature(FileTaskHandler.__init__).parameters) - { + "self", + "base_log_folder", + } + io_kwargs = {k: v for k, v in remote_task_handler_kwargs.items() if k not in fth_params} + return cls( + **{ + "base_log_folder": os.path.expanduser(conf.get_mandatory_value("logging", "base_log_folder")), + "remote_base": conf.get_mandatory_value("logging", "remote_base_log_folder").removeprefix( + "wasb://" + ), + "delete_local_copy": conf.getboolean("logging", "delete_local_logs"), + "wasb_container": conf.get_mandatory_value( + "azure_remote_logging", "remote_wasb_log_container", fallback="airflow-logs" + ), + } + | io_kwargs, + ) + def upload(self, path: str | os.PathLike, ti: RuntimeTI | None = None) -> None: """Upload the given log path to the remote storage.""" path = Path(path) diff --git a/providers/microsoft/azure/tests/unit/microsoft/azure/log/test_wasb_task_handler.py b/providers/microsoft/azure/tests/unit/microsoft/azure/log/test_wasb_task_handler.py index 349ca250cbdb4..88f3be1c68d35 100644 --- a/providers/microsoft/azure/tests/unit/microsoft/azure/log/test_wasb_task_handler.py +++ b/providers/microsoft/azure/tests/unit/microsoft/azure/log/test_wasb_task_handler.py @@ -41,6 +41,89 @@ DEFAULT_DATE = datetime(2020, 8, 10) +class TestWasbRemoteLogIOFromConfig: + @conf_vars( + { + ("logging", "base_log_folder"): "~/airflow/logs", + ("logging", "remote_base_log_folder"): "wasb://path/to/logs", + ("logging", "delete_local_logs"): "True", + ("azure_remote_logging", "remote_wasb_log_container"): "my-container", + } + ) + def test_from_config(self): + subject = WasbRemoteLogIO.from_config() + + assert subject.remote_base == "path/to/logs" + assert subject.base_log_folder == Path(os.path.expanduser("~/airflow/logs")) + assert subject.delete_local_copy is True + assert subject.wasb_container == "my-container" + + @conf_vars( + { + ("logging", "base_log_folder"): "/tmp/airflow/logs", + ("logging", "remote_base_log_folder"): "wasb://path/to/logs", + ("logging", "delete_local_logs"): "False", + ("logging", "remote_task_handler_kwargs"): '{"delete_local_copy": true, "max_bytes": 1024}', + } + ) + def test_from_config_applies_io_kwargs_and_filters_file_handler_kwargs(self): + subject = WasbRemoteLogIO.from_config() + + assert subject.delete_local_copy is True + assert not hasattr(subject, "max_bytes") + assert subject.wasb_container == "airflow-logs" + + @conf_vars({("logging", "remote_task_handler_kwargs"): '["not", "a", "dict"]'}) + def test_from_config_rejects_non_dict_remote_task_handler_kwargs(self): + with pytest.raises(ValueError, match="remote_task_handler_kwargs"): + WasbRemoteLogIO.from_config() + + def test_provider_registers_wasb_scheme(self): + from airflow.providers_manager import ProvidersManager + + manager = ProvidersManager() + if not hasattr(manager, "remote_logging_handler_by_scheme"): + pytest.skip("Airflow core does not support remote logging provider dispatch") + + info = manager.remote_logging_handler_by_scheme("wasb") + + assert info is not None + assert info.classpath == "airflow.providers.microsoft.azure.log.wasb_task_handler.WasbRemoteLogIO" + + @pytest.mark.parametrize( + "manager_classpath", + [ + pytest.param("airflow.providers_manager.ProvidersManager", id="core"), + pytest.param( + "airflow.sdk.providers_manager_runtime.ProvidersManagerTaskRuntime", id="task-runtime" + ), + ], + ) + @conf_vars( + { + ("logging", "remote_logging"): "True", + ("logging", "remote_base_log_folder"): "wasb://path/to/logs", + ("logging", "remote_log_conn_id"): "wasb_default", + } + ) + def test_resolve_remote_task_log_uses_provider_dispatch_not_local_settings(self, manager_classpath): + factory = pytest.importorskip("airflow._shared.logging.factory") + from airflow._shared.module_loading import import_string + from airflow.configuration import conf + + with mock.patch.object(factory, "discover_remote_log_handler", autospec=True) as legacy_discover: + remote_task_log, conn_id = factory.resolve_remote_task_log( + conf=conf, + providers_manager=import_string(manager_classpath)(), + import_string=import_string, + ) + + assert isinstance(remote_task_log, WasbRemoteLogIO) + assert remote_task_log.remote_base == "path/to/logs" + assert conn_id == "wasb_default" + legacy_discover.assert_not_called() + + class TestWasbTaskHandler: @pytest.fixture(autouse=True) def ti(self, create_task_instance, create_log_template):