Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions providers/microsoft/azure/provider.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
# under the License.
from __future__ import annotations

import inspect
import os
import shutil
from functools import cached_property
Expand Down Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down