This is an automated email from the ASF dual-hosted git repository.

potiuk 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 a472c83d715 Add WasbRemoteLogIO.from_config and register wasb remote 
logging scheme (#70301)
a472c83d715 is described below

commit a472c83d715bc2c81bcfa2c19c5933e5fe7dba63
Author: Andrew Chang <[email protected]>
AuthorDate: Sat Aug 1 18:04:28 2026 +0800

    Add WasbRemoteLogIO.from_config and register wasb remote logging scheme 
(#70301)
    
    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 ++
 .../providers/microsoft/azure/get_provider_info.py |  6 ++
 .../microsoft/azure/log/wasb_task_handler.py       | 29 ++++++++
 .../microsoft/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 d4a5e84a2ff..d1a471c3995 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 bf7ab5e96a8..34509081d0e 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 5c631f4d612..51b7f7813c2 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 349ca250cbd..88f3be1c68d 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 @@ pytestmark = pytest.mark.db_test
 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):

Reply via email to