This is an automated email from the ASF dual-hosted git repository.
ashb 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 07197c29a2f Fix Elasticsearch and OpenSearch response wrapper bugs
(#73725)
07197c29a2f is described below
commit 07197c29a2f4b1598615014f7b36bd8237a8de2c
Author: Subhramit Basu <[email protected]>
AuthorDate: Fri Sep 25 20:26:09 2026 +0530
Fix Elasticsearch and OpenSearch response wrapper bugs (#73725)
---
.../airflow/providers/elasticsearch/log/es_response.py | 6 ++----
.../tests/unit/elasticsearch/log/test_es_response.py | 13 +++++++++++++
.../src/airflow/providers/opensearch/log/os_response.py | 6 ++----
.../tests/unit/opensearch/log/test_os_response.py | 15 +++++++++++++++
4 files changed, 32 insertions(+), 8 deletions(-)
diff --git
a/providers/elasticsearch/src/airflow/providers/elasticsearch/log/es_response.py
b/providers/elasticsearch/src/airflow/providers/elasticsearch/log/es_response.py
index fc14e971e68..ea9d4b6c21b 100644
---
a/providers/elasticsearch/src/airflow/providers/elasticsearch/log/es_response.py
+++
b/providers/elasticsearch/src/airflow/providers/elasticsearch/log/es_response.py
@@ -61,7 +61,7 @@ class AttributeList:
def __getitem__(self, k):
"""Retrieve an item or a slice from the list. If the item is a
dictionary, it is wrapped in an AttributeDict."""
val = self._l_[k]
- if isinstance(val, slice):
+ if isinstance(k, slice):
return AttributeList(val)
return _wrap(val)
@@ -104,9 +104,7 @@ class Hit(AttributeDict):
"""
def __init__(self, document):
- data = {}
- if "_source" in document:
- data = document["_source"]
+ data = dict(document.get("_source", {}))
if "fields" in document:
data.update(document["fields"])
diff --git
a/providers/elasticsearch/tests/unit/elasticsearch/log/test_es_response.py
b/providers/elasticsearch/tests/unit/elasticsearch/log/test_es_response.py
index 30c41f8d920..449314712e1 100644
--- a/providers/elasticsearch/tests/unit/elasticsearch/log/test_es_response.py
+++ b/providers/elasticsearch/tests/unit/elasticsearch/log/test_es_response.py
@@ -65,6 +65,12 @@ class TestAttributeList:
assert attr_list[1].key1 == "value1"
assert attr_list[2] == 3
+ def test_slice_access_returns_attribute_list(self):
+ result = AttributeList([1, 2, 3, 4])[1:3]
+
+ assert isinstance(result, AttributeList)
+ assert list(result) == [2, 3]
+
def test_iteration(self):
test_list = [1, {"key": "value"}, 3]
attr_list = AttributeList(test_list)
@@ -174,6 +180,13 @@ class TestHitAndHitMetaAndElasticSearchResponse:
assert isinstance(hit.meta, HitMeta)
assert hit.to_dict() == self.HIT_DOCUMENT["_source"]
+ def test_hit_does_not_mutate_source_when_merging_fields(self):
+ source = {"a": 1}
+ hit = Hit({"_source": source, "fields": {"b": 2}})
+
+ assert hit.to_dict() == {"a": 1, "b": 2}
+ assert source == {"a": 1}
+
def test_hitmeta_initialization_and_to_dict(self):
hitmeta = HitMeta(self.HIT_DOCUMENT)
diff --git
a/providers/opensearch/src/airflow/providers/opensearch/log/os_response.py
b/providers/opensearch/src/airflow/providers/opensearch/log/os_response.py
index 2827cb0f045..937a11f8e70 100644
--- a/providers/opensearch/src/airflow/providers/opensearch/log/os_response.py
+++ b/providers/opensearch/src/airflow/providers/opensearch/log/os_response.py
@@ -36,7 +36,7 @@ class AttributeList:
def __getitem__(self, k):
"""Retrieve an item or a slice from the list. If the item is a
dictionary, it is wrapped in an AttributeDict."""
val = self._l_[k]
- if isinstance(val, slice):
+ if isinstance(k, slice):
return AttributeList(val)
return _wrap(val)
@@ -79,9 +79,7 @@ class Hit(AttributeDict):
"""
def __init__(self, document):
- data = {}
- if "_source" in document:
- data = document["_source"]
+ data = dict(document.get("_source", {}))
if "fields" in document:
data.update(document["fields"])
diff --git a/providers/opensearch/tests/unit/opensearch/log/test_os_response.py
b/providers/opensearch/tests/unit/opensearch/log/test_os_response.py
index f7f36b6732f..920d47b1ad5 100644
--- a/providers/opensearch/tests/unit/opensearch/log/test_os_response.py
+++ b/providers/opensearch/tests/unit/opensearch/log/test_os_response.py
@@ -33,6 +33,14 @@ from airflow.providers.opensearch.log.os_task_handler import
OpensearchTaskHandl
opensearchpy = pytest.importorskip("opensearchpy")
+class TestAttributeList:
+ def test_slice_access_returns_attribute_list(self):
+ result = AttributeList([1, 2, 3, 4])[1:3]
+
+ assert isinstance(result, AttributeList)
+ assert list(result) == [2, 3]
+
+
class TestHitAndHitMetaAndOpenSearchResponse:
OS_DOCUMENT: dict[str, Any] = {
"_shards": {"failed": 0, "skipped": 0, "successful": 7, "total": 7},
@@ -93,6 +101,13 @@ class TestHitAndHitMetaAndOpenSearchResponse:
assert isinstance(hit.meta, HitMeta)
assert hit.to_dict() == self.HIT_DOCUMENT["_source"]
+ def test_hit_does_not_mutate_source_when_merging_fields(self):
+ source = {"a": 1}
+ hit = Hit({"_source": source, "fields": {"b": 2}})
+
+ assert hit.to_dict() == {"a": 1, "b": 2}
+ assert source == {"a": 1}
+
def test_hitmeta_initialization_and_to_dict(self):
hitmeta = HitMeta(self.HIT_DOCUMENT)