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

potiuk 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 caef4f804dc Do not forward an empty ca_certs to the OpenSearch log 
client (#72552)
caef4f804dc is described below

commit caef4f804dc38a064bf5c3c8635a3f2e5462328b
Author: PoAn Yang <[email protected]>
AuthorDate: Wed Sep 9 10:05:54 2026 +0900

    Do not forward an empty ca_certs to the OpenSearch log client (#72552)
    
    Signed-off-by: PoAn Yang <[email protected]>
---
 .../providers/opensearch/log/os_task_handler.py    |  5 +++++
 .../unit/opensearch/log/test_os_task_handler.py    | 23 ++++++++++++++++++++++
 2 files changed, 28 insertions(+)

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 0a7f797ab09..adda50dd0fd 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
@@ -186,6 +186,11 @@ def _ensure_ti(ti: TaskInstanceKey | TaskInstance, 
session) -> TaskInstance:
 def get_os_kwargs_from_config() -> dict[str, Any]:
     open_search_config = conf.getsection("opensearch_configs")
     kwargs_dict = {key: value for key, value in open_search_config.items()} if 
open_search_config else {}
+    # ``ca_certs`` defaults to an empty string, which means "not configured". 
opensearch-py only uses its
+    # default CA bundle when the argument is left out entirely. Passing an 
empty string instead makes it
+    # raise ImproperlyConfigured when both ``use_ssl`` and ``verify_certs`` 
are enabled.
+    if not kwargs_dict.get("ca_certs"):
+        kwargs_dict.pop("ca_certs", None)
     return kwargs_dict
 
 
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 c877ab5cbf1..39ab54f913e 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
@@ -625,6 +625,29 @@ class TestTaskHandlerHelpers:
             assert "http_compress" in args_from_config
             assert "self" not in args_from_config
 
+    @conf_vars(
+        {
+            ("opensearch_configs", "use_ssl"): "True",
+            ("opensearch_configs", "verify_certs"): "True",
+            ("opensearch_configs", "ca_certs"): "",
+        }
+    )
+    def test_empty_ca_certs_is_not_forwarded_to_client(self):
+        """The provider default for ``ca_certs`` is an empty string, which 
opensearch-py treats as
+        "no root certificates" rather than "use certifi" once TLS verification 
is enabled."""
+        from airflow.providers.opensearch.log.os_task_handler import 
_create_opensearch_client
+
+        os_kwargs = get_os_kwargs_from_config()
+
+        assert "ca_certs" not in os_kwargs
+        # Raises ImproperlyConfigured("Root certificates are missing ...") if 
ca_certs="" is forwarded.
+        client = _create_opensearch_client("localhost", 9200, "admin", 
"admin", os_kwargs)
+        assert isinstance(client, opensearchpy.OpenSearch)
+
+    @conf_vars({("opensearch_configs", "ca_certs"): "/etc/ssl/certs/ca.pem"})
+    def test_configured_ca_certs_is_forwarded_to_client(self):
+        assert get_os_kwargs_from_config()["ca_certs"] == 
"/etc/ssl/certs/ca.pem"
+
 
 class TestOpensearchRemoteLogIO:
     @pytest.fixture(autouse=True)

Reply via email to