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

jason810496 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 8a15796a93c Add OpensearchRemoteLogIO.from_config and register 
opensearch scheme (#70295)
8a15796a93c is described below

commit 8a15796a93c69becf158ce09a57ae9a5c3ca6132
Author: PoAn Yang <[email protected]>
AuthorDate: Tue Aug 4 23:34:53 2026 +0900

    Add OpensearchRemoteLogIO.from_config and register opensearch scheme 
(#70295)
    
    Signed-off-by: PoAn Yang <[email protected]>
---
 providers/opensearch/docs/logging/index.rst        |  16 +++
 providers/opensearch/provider.yaml                 |   4 +
 .../providers/opensearch/get_provider_info.py      |   6 +
 .../providers/opensearch/log/os_task_handler.py    |  38 +++++++
 .../unit/opensearch/log/test_os_task_handler.py    | 121 +++++++++++++++++++++
 5 files changed, 185 insertions(+)

diff --git a/providers/opensearch/docs/logging/index.rst 
b/providers/opensearch/docs/logging/index.rst
index 43f84921711..f64aee61b77 100644
--- a/providers/opensearch/docs/logging/index.rst
+++ b/providers/opensearch/docs/logging/index.rst
@@ -44,6 +44,22 @@ First, to use the handler, ``airflow.cfg`` must be 
configured as follows:
     username = <username>
     password = <password>
 
+On Airflow 3.3.0 or above 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
diff --git a/providers/opensearch/provider.yaml 
b/providers/opensearch/provider.yaml
index e0d5fed38ba..65996ac1a03 100644
--- a/providers/opensearch/provider.yaml
+++ b/providers/opensearch/provider.yaml
@@ -96,6 +96,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: ~
diff --git 
a/providers/opensearch/src/airflow/providers/opensearch/get_provider_info.py 
b/providers/opensearch/src/airflow/providers/opensearch/get_provider_info.py
index 70a3fac8ec2..c6a0368e2e9 100644
--- a/providers/opensearch/src/airflow/providers/opensearch/get_provider_info.py
+++ b/providers/opensearch/src/airflow/providers/opensearch/get_provider_info.py
@@ -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,
diff --git 
a/providers/opensearch/src/airflow/providers/opensearch/log/os_task_handler.py 
b/providers/opensearch/src/airflow/providers/opensearch/log/os_task_handler.py
index 6d9722478a5..2ad4db39daa 100644
--- 
a/providers/opensearch/src/airflow/providers/opensearch/log/os_task_handler.py
+++ 
b/providers/opensearch/src/airflow/providers/opensearch/log/os_task_handler.py
@@ -873,6 +873,44 @@ class OpensearchRemoteLogIO(LoggingMixin):  # noqa: D101
 
     processors = ()
 
+    @classmethod
+    def from_config(cls) -> OpensearchRemoteLogIO:
+        """
+        Build the remote log IO from Airflow logging and ``[opensearch]`` 
configuration.
+
+        The ``opensearch://`` value in ``[logging] remote_base_log_folder`` is 
only a routing
+        marker, so every connection and behaviour parameter is read from the 
``[opensearch]``
+        section here.
+
+        This does not merge ``[logging] remote_task_handler_kwargs`` 
IO-kwargs, matching the
+        legacy behavior for OpenSearch.
+        """
+        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)}"
+            )
+        # ``[opensearch] port`` declares an empty-string default, so the key 
is always present and
+        # ``conf.getint`` raises on ``int("")`` instead of falling back to 
9200.
+        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 9200,
+            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}",
+        )
+
     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)
diff --git 
a/providers/opensearch/tests/unit/opensearch/log/test_os_task_handler.py 
b/providers/opensearch/tests/unit/opensearch/log/test_os_task_handler.py
index ca49103b8a5..c877ab5cbf1 100644
--- a/providers/opensearch/tests/unit/opensearch/log/test_os_task_handler.py
+++ b/providers/opensearch/tests/unit/opensearch/log/test_os_task_handler.py
@@ -20,6 +20,7 @@ from __future__ import annotations
 import dataclasses
 import json
 import logging
+import os
 import re
 from io import StringIO
 from pathlib import Path
@@ -786,6 +787,126 @@ class TestOpensearchRemoteLogIO:
         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_ignores_remote_task_handler_kwargs(self):
+        """Unlike the object-storage backends, OpenSearch does not merge IO 
kwargs (legacy parity)."""
+        subject = OpensearchRemoteLogIO.from_config()
+
+        # ``delete_local_copy`` stays at the ``[logging] delete_local_logs`` 
value.
+        assert subject.delete_local_copy is False
+        # ``max_bytes`` belongs to FileTaskHandler and must not reach the IO 
class.
+        assert not hasattr(subject, "max_bytes")
+
+    @conf_vars(
+        {
+            ("logging", "base_log_folder"): "/tmp/airflow/logs",
+            ("opensearch", "host"): "https://opensearch.example.com:9201";,
+            ("opensearch", "port"): "",
+            ("opensearch", "username"): "admin",
+            ("opensearch", "password"): "secret",
+        }
+    )
+    def test_from_config_defaults_port_when_unset(self):
+        """An unset port falls back to 9200 (legacy intent) rather than the 
host URL's port."""
+        subject = OpensearchRemoteLogIO.from_config()
+
+        assert subject.port == 9200
+
+    @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

Reply via email to