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)
 

Reply via email to