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

eladkal 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 c8b136358ce Allow OpensearchTaskHandler and OpensearchRemoteLogIO to 
take empty username and password (#71692)
c8b136358ce is described below

commit c8b136358ced95df3c073e64b04f9985dcaec504
Author: Christopher Anderson 
<[email protected]>
AuthorDate: Sun Oct 4 15:40:21 2026 +0200

    Allow OpensearchTaskHandler and OpensearchRemoteLogIO to take empty 
username and password (#71692)
    
    * Allow OpensearchTaskHandler and OpensearchRemoteLogIO to take empty 
username and password
    
    * add tests for username and/or password being set
    
    * change product code to pass http_auth if either or both username and/or 
password are non-empty
    
    * add OpensearchIO tests for username and/or password being set
    
    * ruff fix
    
    * consolidate auth and no-auth behavior into single test
    
    * use the correct variable name
    
    * don't know how I missed this one too
---
 .../providers/opensearch/log/os_task_handler.py    |  8 ++-
 .../unit/opensearch/log/test_os_task_handler.py    | 60 ++++++++++++++++++++++
 2 files changed, 66 insertions(+), 2 deletions(-)

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 8561a11eb02..1f164cf56c8 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
@@ -220,9 +220,13 @@ def _create_opensearch_client(
 ) -> OpenSearch:
     parsed_url = urlparse(_format_url(host))
     resolved_port = port if port is not None else (parsed_url.port or 9200)
+    connection_kwargs: dict[str, Any] = {
+        "hosts": [{"host": parsed_url.hostname, "port": resolved_port, 
"scheme": parsed_url.scheme}]
+    }
+    if username or password:
+        connection_kwargs["http_auth"] = (username, password)
     return OpenSearch(
-        hosts=[{"host": parsed_url.hostname, "port": resolved_port, "scheme": 
parsed_url.scheme}],
-        http_auth=(username, password),
+        **connection_kwargs,
         **os_kwargs,
     )
 
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 3af00e588ca..b9e6582c663 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
@@ -288,6 +288,37 @@ class TestOpensearchTaskHandler:
         )
         assert handler.index_patterns == patterns
 
+    @pytest.mark.parametrize(
+        ("username", "password", "expect_http_auth"),
+        [
+            ("admin", "secret", True),
+            ("admin", "", True),
+            ("", "secret", True),
+        ],
+    )
+    def test_client_with_auth(self, username, password, expect_http_auth):
+        """If either username or password are provided, the handler should 
pass http_auth to the client."""
+        handler = OpensearchTaskHandler(
+            base_log_folder=self.local_log_location,
+            end_of_log_mark=self.end_of_log_mark,
+            write_stdout=self.write_stdout,
+            host="localhost",
+            port=9200,
+            username=username,
+            password=password,
+            json_format=self.json_format,
+            json_fields=self.json_fields,
+            host_field=self.host_field,
+            offset_field=self.offset_field,
+        )
+
+        transport_args = handler.client.transport.kwargs
+        if expect_http_auth:
+            assert "http_auth" in transport_args
+            assert transport_args["http_auth"] == (username, password)
+        else:
+            assert "http_auth" not in handler.client.transport.kwargs
+
     @pytest.mark.db_test
     @pytest.mark.parametrize("metadata_mode", ["provided", "none", "empty"])
     def test_read(self, ti, metadata_mode):
@@ -849,6 +880,35 @@ class TestOpensearchRemoteLogIO:
         log_file.write_text('{"message": "test"}\n')
         self.opensearch_io.upload(log_file, ti=None)
 
+    @pytest.mark.parametrize(
+        ("username", "password", "expect_http_auth"),
+        [
+            ("admin", "secret", True),
+            ("admin", "", True),
+            ("", "secret", True),
+        ],
+    )
+    def test_client_with_auth(self, username, password, expect_http_auth):
+        """If either username or password are provided, the IO should pass 
http_auth to the client."""
+        opensearch_io = OpensearchRemoteLogIO(
+            write_to_opensearch=True,
+            write_stdout=True,
+            delete_local_copy=True,
+            host="localhost",
+            port=9200,
+            username=username,
+            password=password,
+            base_log_folder=self.opensearch_io.base_log_folder,
+            
log_id_template="{dag_id}-{task_id}-{run_id}-{map_index}-{try_number}",
+        )
+
+        transport_args = opensearch_io.client.transport.kwargs
+        if expect_http_auth:
+            assert "http_auth" in transport_args
+            assert transport_args["http_auth"] == (username, password)
+        else:
+            assert "http_auth" not in opensearch_io.client.transport.kwargs
+
 
 class TestOpensearchRemoteLogIOFromConfig:
     @conf_vars(

Reply via email to