From 24663d60f19051d1d6e252b2ebe71986c41d441c Mon Sep 17 00:00:00 2001 From: LIU ZHE YOU Date: Mon, 13 Jul 2026 06:38:42 +0000 Subject: [PATCH] Add S3RemoteLogIO.from_config and register s3 remote logging scheme Airflow resolves remote task log handlers through provider dispatch on the remote_base_log_folder URL scheme since #67056; providers must expose from_config so core and the Task SDK no longer depend on the hardcoded branches in airflow_local_settings.py. This migrates the s3 scheme as the first adopter. --- providers/amazon/provider.yaml | 2 + .../amazon/aws/log/s3_task_handler.py | 26 ++++++ .../providers/amazon/get_provider_info.py | 3 +- .../amazon/aws/log/test_s3_task_handler.py | 83 ++++++++++++++++++- 4 files changed, 112 insertions(+), 2 deletions(-) diff --git a/providers/amazon/provider.yaml b/providers/amazon/provider.yaml index 4c00876e56518..2230bf8f345e1 100644 --- a/providers/amazon/provider.yaml +++ b/providers/amazon/provider.yaml @@ -1131,6 +1131,8 @@ logging: remote-logging: - classpath: airflow.providers.amazon.aws.log.cloudwatch_task_handler.CloudWatchRemoteLogIO scheme: cloudwatch + - classpath: airflow.providers.amazon.aws.log.s3_task_handler.S3RemoteLogIO + scheme: s3 config: aws: diff --git a/providers/amazon/src/airflow/providers/amazon/aws/log/s3_task_handler.py b/providers/amazon/src/airflow/providers/amazon/aws/log/s3_task_handler.py index b1ce26098fb34..6ef82f301da4b 100644 --- a/providers/amazon/src/airflow/providers/amazon/aws/log/s3_task_handler.py +++ b/providers/amazon/src/airflow/providers/amazon/aws/log/s3_task_handler.py @@ -17,6 +17,7 @@ # under the License. from __future__ import annotations +import inspect import logging import os import pathlib @@ -46,6 +47,31 @@ class S3RemoteLogIO(LoggingMixin): # noqa: D101 processors = () + @classmethod + def from_config(cls) -> S3RemoteLogIO: + """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)}" + ) + # remote_task_handler_kwargs mixes FileTaskHandler kwargs with IO kwargs; only the + # latter belong to this class (same split as airflow_local_settings.py). + 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"), + "delete_local_copy": conf.getboolean("logging", "delete_local_logs"), + } + | io_kwargs, + ) + def upload(self, path: os.PathLike | str, ti: RuntimeTI | None = None) -> None: """Upload the given log path to the remote storage.""" path = pathlib.Path(path) diff --git a/providers/amazon/src/airflow/providers/amazon/get_provider_info.py b/providers/amazon/src/airflow/providers/amazon/get_provider_info.py index 9bc764b09af08..5314709a5df68 100644 --- a/providers/amazon/src/airflow/providers/amazon/get_provider_info.py +++ b/providers/amazon/src/airflow/providers/amazon/get_provider_info.py @@ -1258,7 +1258,8 @@ def get_provider_info(): { "classpath": "airflow.providers.amazon.aws.log.cloudwatch_task_handler.CloudWatchRemoteLogIO", "scheme": "cloudwatch", - } + }, + {"classpath": "airflow.providers.amazon.aws.log.s3_task_handler.S3RemoteLogIO", "scheme": "s3"}, ], "config": { "aws": { diff --git a/providers/amazon/tests/unit/amazon/aws/log/test_s3_task_handler.py b/providers/amazon/tests/unit/amazon/aws/log/test_s3_task_handler.py index 01614ff732b52..0576756c1965f 100644 --- a/providers/amazon/tests/unit/amazon/aws/log/test_s3_task_handler.py +++ b/providers/amazon/tests/unit/amazon/aws/log/test_s3_task_handler.py @@ -21,6 +21,7 @@ import copy import logging import os +import pathlib from unittest import mock import boto3 @@ -30,7 +31,7 @@ from airflow.models import DAG, DagRun, TaskInstance from airflow.providers.amazon.aws.hooks.s3 import S3Hook -from airflow.providers.amazon.aws.log.s3_task_handler import S3TaskHandler +from airflow.providers.amazon.aws.log.s3_task_handler import S3RemoteLogIO, S3TaskHandler from airflow.utils.state import State, TaskInstanceState from tests_common.test_utils.compat import EmptyOperator @@ -52,6 +53,86 @@ def s3mock(): yield +class TestS3RemoteLogIOFromConfig: + @conf_vars( + { + ("logging", "base_log_folder"): "~/airflow/logs", + ("logging", "remote_base_log_folder"): "s3://bucket/remote/log/location", + ("logging", "delete_local_logs"): "True", + } + ) + def test_from_config(self): + subject = S3RemoteLogIO.from_config() + + assert subject.remote_base == "s3://bucket/remote/log/location" + assert subject.base_log_folder == pathlib.Path(os.path.expanduser("~/airflow/logs")) + assert subject.delete_local_copy is True + + @conf_vars( + { + ("logging", "base_log_folder"): "/tmp/airflow/logs", + ("logging", "remote_base_log_folder"): "s3://bucket/remote/log/location", + ("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 = S3RemoteLogIO.from_config() + + assert subject.delete_local_copy is True + assert not hasattr(subject, "max_bytes") + + @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"): + S3RemoteLogIO.from_config() + + def test_provider_registers_s3_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("s3") + + assert info is not None + assert info.classpath == "airflow.providers.amazon.aws.log.s3_task_handler.S3RemoteLogIO" + + @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"): "s3://bucket/remote/log/location", + ("logging", "remote_log_conn_id"): "aws_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, S3RemoteLogIO) + assert remote_task_log.remote_base == "s3://bucket/remote/log/location" + assert conn_id == "aws_default" + legacy_discover.assert_not_called() + + @pytest.mark.db_test class TestS3RemoteLogIO: def clear_db(self):