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")