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
2 changes: 2 additions & 0 deletions providers/amazon/provider.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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:
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 logging
import os
import pathlib
Expand Down Expand Up @@ -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) - {
Comment thread
ferruzzi marked this conversation as resolved.
"self",
"base_log_folder",
}
Comment thread
jason810496 marked this conversation as resolved.
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,
)
Comment thread
jason810496 marked this conversation as resolved.

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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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": {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
import copy
import logging
import os
import pathlib
from unittest import mock

import boto3
Expand All @@ -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
Expand All @@ -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):
Comment thread
jason810496 marked this conversation as resolved.
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):
Expand Down
Loading