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(