Skip to content

Commit daddaa3

Browse files
Allow OpensearchTaskHandler and OpensearchRemoteLogIO to take empty username and password
1 parent ef32e9f commit daddaa3

2 files changed

Lines changed: 36 additions & 2 deletions

File tree

providers/opensearch/src/airflow/providers/opensearch/log/os_task_handler.py

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -219,9 +219,13 @@ def _create_opensearch_client(
219219
) -> OpenSearch:
220220
parsed_url = urlparse(_format_url(host))
221221
resolved_port = port if port is not None else (parsed_url.port or 9200)
222+
connection_kwargs: dict[str, Any] = {
223+
"hosts": [{"host": parsed_url.hostname, "port": resolved_port, "scheme": parsed_url.scheme}]
224+
}
225+
if username and password:
226+
connection_kwargs["http_auth"] = (username, password)
222227
return OpenSearch(
223-
hosts=[{"host": parsed_url.hostname, "port": resolved_port, "scheme": parsed_url.scheme}],
224-
http_auth=(username, password),
228+
**connection_kwargs,
225229
**os_kwargs,
226230
)
227231

providers/opensearch/tests/unit/opensearch/log/test_os_task_handler.py

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -288,6 +288,22 @@ def test_client_with_patterns(self):
288288
)
289289
assert handler.index_patterns == patterns
290290

291+
def test_client_no_auth(self):
292+
handler = OpensearchTaskHandler(
293+
base_log_folder=self.local_log_location,
294+
end_of_log_mark=self.end_of_log_mark,
295+
write_stdout=self.write_stdout,
296+
host="localhost",
297+
port=9200,
298+
username="",
299+
password="",
300+
json_format=self.json_format,
301+
json_fields=self.json_fields,
302+
host_field=self.host_field,
303+
offset_field=self.offset_field,
304+
)
305+
assert "http_auth" not in handler.client.transport.kwargs
306+
291307
@pytest.mark.db_test
292308
@pytest.mark.parametrize("metadata_mode", ["provided", "none", "empty"])
293309
def test_read(self, ti, metadata_mode):
@@ -786,6 +802,20 @@ def test_upload_returns_early_when_ti_is_none(self, tmp_path):
786802
log_file.write_text('{"message": "test"}\n')
787803
self.opensearch_io.upload(log_file, ti=None)
788804

805+
def test_client_no_auth(self):
806+
opensearch_io = OpensearchRemoteLogIO(
807+
write_to_opensearch=True,
808+
write_stdout=True,
809+
delete_local_copy=True,
810+
host="localhost",
811+
port=9200,
812+
username="",
813+
password="",
814+
base_log_folder=self.opensearch_io.base_log_folder,
815+
log_id_template="{dag_id}-{task_id}-{run_id}-{map_index}-{try_number}",
816+
)
817+
assert "http_auth" not in opensearch_io.client.transport.kwargs
818+
789819

790820
class TestOpensearchRemoteLogIOFromConfig:
791821
@conf_vars(

0 commit comments

Comments
 (0)