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 b156582855c Fix OpenSearch provider log pagination (#73947)
b156582855c is described below

commit b156582855ce557f8fdc2c0f2c8fd8f66e7e124c
Author: BHUMIKA KADU✨ <[email protected]>
AuthorDate: Sun Oct 4 19:08:45 2026 +0530

    Fix OpenSearch provider log pagination (#73947)
    
    * Fix OpenSearch provider log pagination
    
    * breeze: trigger OpenSearch provider integration tests
    
    * test: fix selective checks formatting
    
    * test: update OpenSearch selective check expectation
---
 dev/breeze/src/airflow_breeze/global_constants.py  |  4 ++-
 dev/breeze/tests/test_selective_checks.py          | 15 +++++++-
 .../providers/opensearch/log/os_task_handler.py    | 24 +++++++++++--
 .../opensearch/log/test_os_remote_log_io.py        | 18 ++++++++++
 .../unit/opensearch/log/test_os_task_handler.py    | 42 +++++++++++++++++++++-
 5 files changed, 97 insertions(+), 6 deletions(-)

diff --git a/dev/breeze/src/airflow_breeze/global_constants.py 
b/dev/breeze/src/airflow_breeze/global_constants.py
index 08ce0cada3d..cac0d88b617 100644
--- a/dev/breeze/src/airflow_breeze/global_constants.py
+++ b/dev/breeze/src/airflow_breeze/global_constants.py
@@ -89,6 +89,7 @@ TESTABLE_PROVIDERS_INTEGRATIONS = [
     "cassandra",
     "drill",
     "elasticsearch",
+    "opensearch",
     "tinkerpop",
     "kafka",
     "localstack",
@@ -120,6 +121,7 @@ TESTABLE_PROVIDERS_INTEGRATION_OWNERS = {
     "cassandra": "apache.cassandra",
     "drill": "apache.drill",
     "elasticsearch": "elasticsearch",
+    "opensearch": "opensearch",
     "tinkerpop": "apache.tinkerpop",
     "kafka": "apache.kafka",
     "localstack": "amazon",
@@ -137,7 +139,7 @@ OTEL_INTEGRATION = "otel"
 OPENLINEAGE_INTEGRATION = "openlineage"
 OPENSEARCH_INTEGRATION = "opensearch"
 OTHER_CORE_INTEGRATIONS = [STATSD_INTEGRATION, KEYCLOAK_INTEGRATION]
-OTHER_PROVIDERS_INTEGRATIONS = [OPENLINEAGE_INTEGRATION, 
OPENSEARCH_INTEGRATION]
+OTHER_PROVIDERS_INTEGRATIONS = [OPENLINEAGE_INTEGRATION]
 ALLOWED_DEBIAN_VERSIONS = ["bookworm"]
 ALL_CORE_INTEGRATIONS = sorted(
     [
diff --git a/dev/breeze/tests/test_selective_checks.py 
b/dev/breeze/tests/test_selective_checks.py
index 8c7d74079a7..7e8d6e8d712 100644
--- a/dev/breeze/tests/test_selective_checks.py
+++ b/dev/breeze/tests/test_selective_checks.py
@@ -1353,7 +1353,7 @@ def assert_outputs_are_printed(expected_outputs: 
dict[str, str], stderr: str):
                     "core-test-types-list-as-strings-in-json": 
ALL_CI_SELECTIVE_TEST_TYPES_AS_JSON,
                     "providers-test-types-list-as-strings-in-json": 
ALL_PROVIDERS_SELECTIVE_TEST_TYPES_AS_JSON,
                     "testable-core-integrations": "['kerberos', 'otel', 
'redis']",
-                    "testable-providers-integrations": "['celery', 
'cassandra', 'drill', 'elasticsearch', 'tinkerpop', 'kafka', "
+                    "testable-providers-integrations": "['celery', 
'cassandra', 'drill', 'elasticsearch', 'opensearch', 'tinkerpop', 'kafka', "
                     "'mongo', 'pinot', 'qdrant', 'redis', 'trino', 'ydb']",
                     "run-mypy-providers": "true",
                 },
@@ -3734,6 +3734,19 @@ def 
test_testable_providers_integrations_gated_by_affected_provider():
     assert "ydb" not in result
 
 
+def test_opensearch_provider_integration_triggered_by_affected_provider():
+    """Verify that changes to the OpenSearch provider trigger its integration 
test."""
+    selective_checks = SelectiveChecks(
+        
files=("providers/opensearch/src/airflow/providers/opensearch/log/os_task_handler.py",),
+        commit_ref=NEUTRAL_COMMIT,
+        github_event=GithubEvents.PULL_REQUEST,
+        platform=CI_AMD_PLATFORM,
+    )
+    result = selective_checks.testable_providers_integrations
+    assert "opensearch" in result
+    assert "cassandra" not in result
+
+
 def test_individual_providers_excludes_platform_excluded_on_arm():
     """ibm.mq and ibm.db2 declare `excluded-platforms: [linux/arm64]`, so they 
must be
     absent from the ARM individual-providers matrix (used by the Low-dep ARM 
canary job)
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 bdc34643b40..8561a11eb02 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
@@ -999,9 +999,27 @@ class OpensearchRemoteLogIO(LoggingMixin):  # noqa: D101
     def read(self, _relative_path: str, ti: RuntimeTI) -> tuple[LogSourceInfo, 
LogMessages]:
         log_id = _render_log_id(self.log_id_template, ti, ti.try_number)  # 
type: ignore[arg-type]
         self.log.info("Reading log %s from Opensearch", log_id)
-        response = self._os_read(log_id, 0, ti)
-        if response is not None and response.hits:
-            logs_by_host = self._group_logs_by_host(response)
+        responses = []
+        offset = 0
+
+        while True:
+            response = self._os_read(log_id, offset, ti)
+            if response is None or not response.hits:
+                break
+
+            responses.append(response)
+
+            next_offset = attrgetter(self.offset_field)(response[-1])
+            if next_offset == offset:
+                break
+            offset = next_offset
+
+        if responses:
+            grouped_logs = defaultdict(list)
+            for response in responses:
+                for host, hits in self._group_logs_by_host(response).items():
+                    grouped_logs[host].extend(hits)
+            logs_by_host = grouped_logs
         else:
             logs_by_host = None
 
diff --git 
a/providers/opensearch/tests/integration/opensearch/log/test_os_remote_log_io.py
 
b/providers/opensearch/tests/integration/opensearch/log/test_os_remote_log_io.py
index 9aef4777618..a81c3eda1c5 100644
--- 
a/providers/opensearch/tests/integration/opensearch/log/test_os_remote_log_io.py
+++ 
b/providers/opensearch/tests/integration/opensearch/log/test_os_remote_log_io.py
@@ -104,6 +104,24 @@ class TestOpensearchRemoteLogIOIntegration:
             assert "event" in log_entry
             assert log_entry["event"] == expected
 
+    @patch(
+        "airflow.providers.opensearch.log.os_task_handler.TASK_LOG_FIELDS",
+        ["message"],
+    )
+    def test_read_returns_all_logs_when_exceeding_page_size(self, ti, 
tmp_path):
+        log_file = tmp_path / "large.log"
+        sample_logs = [{"message": f"log line {i}"} for i in range(1500)]
+        log_file.write_text("\n".join(json.dumps(log) for log in sample_logs) 
+ "\n")
+
+        self.opensearch_io.upload(log_file, ti)
+        self.opensearch_io.client.indices.refresh(index=self.target_index)
+
+        _, log_messages = self.opensearch_io.read("", ti)
+
+        assert len(log_messages) == 1500
+        assert json.loads(log_messages[0])["event"] == "log line 0"
+        assert json.loads(log_messages[-1])["event"] == "log line 1499"
+
     def test_read_missing_log(self, ti):
         self.opensearch_io.client.indices.create(index=self.target_index)
 
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 39ab54f913e..3af00e588ca 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
@@ -24,7 +24,7 @@ import os
 import re
 from io import StringIO
 from pathlib import Path
-from unittest.mock import Mock, patch
+from unittest.mock import Mock, call, patch
 
 import pendulum
 import pytest
@@ -792,6 +792,46 @@ class TestOpensearchRemoteLogIO:
         assert log_source_info == []
         assert f"*** Log {log_id} not found in Opensearch" in log_messages[0]
 
+    def test_read_returns_all_logs_when_exceeding_page_size(self, ti):
+        log_id = _render_log_id(self.opensearch_io.log_id_template, ti, 
ti.try_number)
+
+        first_page = [
+            {
+                "event": f"log line {i}",
+                "log_id": log_id,
+                "offset": i + 1,
+            }
+            for i in range(1000)
+        ]
+        second_page = [
+            {
+                "event": f"log line {1000 + i}",
+                "log_id": log_id,
+                "offset": 1001 + i,
+            }
+            for i in range(500)
+        ]
+
+        responses = [
+            _make_os_response(self.opensearch_io, *first_page),
+            _make_os_response(self.opensearch_io, *second_page),
+            None,
+        ]
+
+        with patch.object(self.opensearch_io, "_os_read", 
side_effect=responses) as mock_os_read:
+            log_source_info, log_messages = self.opensearch_io.read("", ti)
+
+        assert log_source_info == ["http://localhost";]
+        assert len(log_messages) == 1500
+        assert json.loads(log_messages[0])["event"] == "log line 0"
+        assert json.loads(log_messages[-1])["event"] == "log line 1499"
+
+        assert mock_os_read.call_args_list == [
+            call(log_id, 0, ti),
+            call(log_id, 1000, ti),
+            call(log_id, 1500, ti),
+        ]
+
     def test_get_index_patterns_with_callable(self):
         with 
patch("airflow.providers.opensearch.log.os_task_handler.import_string") as 
mock_import_string:
             mock_callable = Mock(return_value="callable_index_pattern")

Reply via email to