Skip to content
Open
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
16 changes: 16 additions & 0 deletions providers/opensearch/docs/logging/index.rst
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,22 @@ First, to use the handler, ``airflow.cfg`` must be configured as follows:
username = <username>
password = <password>

On Airflow 3.x you can also route remote logging to OpenSearch through the provider
dispatch mechanism by adding an ``opensearch://`` scheme to
``[logging] remote_base_log_folder``:

.. code-block:: ini

[logging]
remote_logging = True
remote_base_log_folder = opensearch://

[opensearch]
host = <host>
port = <port>
username = <username>
password = <password>

To output task logs to stdout in JSON format, the following config could be used:

.. code-block:: ini
Expand Down
4 changes: 4 additions & 0 deletions providers/opensearch/provider.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,10 @@ connection-types:
logging:
- airflow.providers.opensearch.log.os_task_handler.OpensearchTaskHandler

remote-logging:
- classpath: airflow.providers.opensearch.log.os_task_handler.OpensearchRemoteLogIO
scheme: opensearch

config:
opensearch:
description: ~
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,12 @@ def get_provider_info():
}
],
"logging": ["airflow.providers.opensearch.log.os_task_handler.OpensearchTaskHandler"],
"remote-logging": [
{
"classpath": "airflow.providers.opensearch.log.os_task_handler.OpensearchRemoteLogIO",
"scheme": "opensearch",
}
],
"config": {
"opensearch": {
"description": None,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
from __future__ import annotations

import contextlib
import inspect
import json
import logging
import os
Expand Down Expand Up @@ -873,6 +874,43 @@ class OpensearchRemoteLogIO(LoggingMixin): # noqa: D101

processors = ()

@classmethod
def from_config(cls) -> OpensearchRemoteLogIO:
"""Build the remote log IO from Airflow logging and ``[opensearch]`` 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}
port = conf.get("opensearch", "port", fallback="")
return cls(
**{
"base_log_folder": os.path.expanduser(conf.get_mandatory_value("logging", "base_log_folder")),
"delete_local_copy": conf.getboolean("logging", "delete_local_logs"),
"host": conf.get("opensearch", "host", fallback=""),
"port": int(port) if port else None,
"username": conf.get_mandatory_value("opensearch", "username"),
"password": conf.get_mandatory_value("opensearch", "password"),
"write_stdout": conf.getboolean("opensearch", "write_stdout"),
"write_to_opensearch": conf.getboolean("opensearch", "write_to_os"),
"json_format": conf.getboolean("opensearch", "json_format"),
"target_index": conf.get_mandatory_value("opensearch", "target_index"),
"host_field": conf.get_mandatory_value("opensearch", "host_field"),
"offset_field": conf.get_mandatory_value("opensearch", "offset_field"),
"log_id_template": conf.get("opensearch", "log_id_template", fallback="")
or "{dag_id}-{task_id}-{run_id}-{map_index}-{try_number}",
}
| io_kwargs,
)

def __attrs_post_init__(self):
self.host = _format_url(self.host)
self.port = self.port if self.port is not None else (urlparse(self.host).port or 9200)
Expand Down
105 changes: 105 additions & 0 deletions providers/opensearch/tests/unit/opensearch/log/test_os_task_handler.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
import dataclasses
import json
import logging
import os
import re
from io import StringIO
from pathlib import Path
Expand Down Expand Up @@ -788,6 +789,110 @@ def test_upload_returns_early_when_ti_is_none(self, tmp_path):
self.opensearch_io.upload(log_file, ti=None)


class TestOpensearchRemoteLogIOFromConfig:
@conf_vars(
{
("logging", "base_log_folder"): "~/airflow/logs",
("logging", "delete_local_logs"): "True",
("opensearch", "host"): "https://opensearch.example.com:9200",
("opensearch", "port"): "9201",
("opensearch", "username"): "admin",
("opensearch", "password"): "secret",
("opensearch", "write_stdout"): "True",
("opensearch", "write_to_os"): "True",
("opensearch", "json_format"): "True",
("opensearch", "target_index"): "my-logs",
("opensearch", "host_field"): "host.name",
("opensearch", "offset_field"): "log.offset",
("opensearch", "log_id_template"): "{dag_id}-{task_id}-{run_id}",
}
)
def test_from_config(self):
subject = OpensearchRemoteLogIO.from_config()

assert subject.base_log_folder == Path(os.path.expanduser("~/airflow/logs"))
assert subject.delete_local_copy is True
assert subject.host == "https://opensearch.example.com:9200"
assert subject.port == 9201
assert subject.username == "admin"
assert subject.password == "secret"
assert subject.write_stdout is True
assert subject.write_to_opensearch is True
assert subject.json_format is True
assert subject.target_index == "my-logs"
assert subject.host_field == "host.name"
assert subject.offset_field == "log.offset"
assert subject.log_id_template == "{dag_id}-{task_id}-{run_id}"

@conf_vars(
{
("logging", "base_log_folder"): "/tmp/airflow/logs",
("logging", "delete_local_logs"): "False",
("opensearch", "host"): "https://opensearch.example.com:9200",
("opensearch", "username"): "admin",
("opensearch", "password"): "secret",
("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 = OpensearchRemoteLogIO.from_config()

# ``delete_local_copy`` is an IO kwarg, so it overrides the config value.
assert subject.delete_local_copy is True
# ``max_bytes`` belongs to FileTaskHandler, so it must not reach the IO class.
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"):
OpensearchRemoteLogIO.from_config()

def test_provider_registers_opensearch_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("opensearch")

assert info is not None
assert info.classpath == "airflow.providers.opensearch.log.os_task_handler.OpensearchRemoteLogIO"

@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"): "opensearch://",
("opensearch", "host"): "https://opensearch.example.com:9200",
("opensearch", "username"): "admin",
("opensearch", "password"): "secret",
}
)
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 patch.object(factory, "discover_remote_log_handler", autospec=True) as legacy_discover:
remote_task_log, _ = factory.resolve_remote_task_log(
conf=conf,
providers_manager=import_string(manager_classpath)(),
import_string=import_string,
)

assert isinstance(remote_task_log, OpensearchRemoteLogIO)
legacy_discover.assert_not_called()


class TestFormatErrorDetail:
def test_returns_none_for_empty(self):
assert _format_error_detail(None) is None
Expand Down