From d694b31cbf43d7bece81ace62d0df9d527bc582f Mon Sep 17 00:00:00 2001 From: Hemkumar Chheda Date: Fri, 3 Jul 2026 13:04:32 +0530 Subject: [PATCH 1/3] Prevent malformed OpenSearch log entries from crashing task log fetch A single stored log entry with a non-string `event` field (for example, a task that logs a list or dict as the sole message argument) currently crashes the entire task-log-fetch request with an unhandled pydantic.ValidationError, instead of degrading gracefully. _read() built StructuredLogMessage objects from stored OpenSearch hits without catching validation failures. This catches ValidationError per hit and falls back to a stringified event, matching the existing fallback pattern in file_task_handler.py's _log_stream_to_parsed_log_stream. Related: #69304 --- .../opensearch/log/os_task_handler.py | 28 +++++++++++- .../opensearch/log/test_os_task_handler.py | 43 +++++++++++++++++++ 2 files changed, 70 insertions(+), 1 deletion(-) 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 348795bde9cf1..7ab118c9a181a 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 @@ -36,6 +36,7 @@ import pendulum from opensearchpy import OpenSearch, helpers from opensearchpy.exceptions import NotFoundError +from pydantic import ValidationError from sqlalchemy import select import airflow.logging_config as alc @@ -71,6 +72,8 @@ LOG_LINE_DEFAULTS = {"exc_text": "", "stack_info": ""} TASK_LOG_FIELDS = ["timestamp", "event", "level", "chan", "logger", "error_detail", "message", "levelname"] +logger = logging.getLogger(__name__) + def _format_error_detail(error_detail: Any) -> str | None: """Render the structured ``error_detail`` written by the Airflow 3 supervisor as a traceback string.""" @@ -120,6 +123,29 @@ def _build_log_fields(hit_dict: dict[str, Any]) -> dict[str, Any]: return fields +def _safe_build_structured_log_message(hit_dict: dict[str, Any]) -> StructuredLogMessage: + """ + Build a StructuredLogMessage from a stored OpenSearch hit, tolerating malformed fields. + + A single malformed stored log entry (for example a non-string ``event`` produced by + logging a list or dict as the sole message argument) must not fail the entire + log-fetch request. Fall back to a stringified event, mirroring the fallback used for + unparsable raw log lines in ``_log_stream_to_parsed_log_stream``. + """ + fields = _build_log_fields(hit_dict) + try: + return StructuredLogMessage(**fields) + except ValidationError: + logger.warning( + "Failed to parse stored log entry into StructuredLogMessage; falling back to " + "stringified event. Offending fields: %s", + fields, + ) + return StructuredLogMessage( + event=str(fields.get("event", hit_dict)), timestamp=fields.get("timestamp") + ) + + def getattr_nested(obj, item, default): """ Get item from obj but return default if not found. @@ -625,7 +651,7 @@ def _read( # Flatten all hits, filter to only desired fields, and construct StructuredLogMessage objects message = header + [ - StructuredLogMessage(**_build_log_fields(hit.to_dict())) + _safe_build_structured_log_message(hit.to_dict()) for hits in logs_by_host.values() for hit in hits ] 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 46a0e9805cd2f..4295f5260efe8 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 @@ -38,6 +38,7 @@ _build_log_fields, _format_error_detail, _render_log_id, + _safe_build_structured_log_message, _strip_userinfo, get_os_kwargs_from_config, getattr_nested, @@ -396,6 +397,31 @@ def test_read_with_custom_offset_and_host_fields(self, ti): assert metadata["offset"] == "1" assert not metadata["end_of_log"] + @pytest.mark.skipif(not AIRFLOW_V_3_0_PLUS, reason="StructuredLogMessage fallback is Airflow 3+ only") + @pytest.mark.db_test + def test_read_with_malformed_event_falls_back_to_stringified_event(self, ti): + malformed_event = ["not", "a", "string"] + malformed_source = { + "message": self.test_message, + "event": malformed_event, + "log_id": self.LOG_ID, + "offset": 2, + } + response = _make_os_response(self.os_task_handler.io, self.base_log_source, malformed_source) + + with patch.object(self.os_task_handler.io, "_os_read", return_value=response): + with patch("airflow.providers.opensearch.log.os_task_handler.logger") as mock_logger: + logs, metadatas = self.os_task_handler.read(ti, 1) + + metadata = _assert_log_events( + logs, + metadatas, + expected_events=[self.test_message, str(malformed_event)], + expected_sources=["http://localhost"], + ) + assert not metadata["end_of_log"] + mock_logger.warning.assert_called_once() + @pytest.mark.db_test def test_set_context(self, ti): self.os_task_handler.set_context(ti) @@ -838,3 +864,20 @@ def test_error_detail_dropped_when_empty(self): hit = {"event": "msg", "error_detail": []} result = _build_log_fields(hit) assert "error_detail" not in result + + +class TestSafeBuildStructuredLogMessage: + def test_string_event_returns_unchanged_and_does_not_warn(self): + hit = {"event": "hello", "level": "info"} + with patch("airflow.providers.opensearch.log.os_task_handler.logger") as mock_logger: + result = _safe_build_structured_log_message(hit) + assert result.event == "hello" + mock_logger.warning.assert_not_called() + + def test_non_string_event_falls_back_to_stringified_event(self): + hit = {"event": ["a", "b"], "timestamp": "2024-01-01T00:00:00Z"} + with patch("airflow.providers.opensearch.log.os_task_handler.logger") as mock_logger: + result = _safe_build_structured_log_message(hit) + assert result.event == str(["a", "b"]) + assert result.timestamp is not None + mock_logger.warning.assert_called_once() From 81c1d1efb4f0e0ac5bbd848ee26c020b85b40729 Mon Sep 17 00:00:00 2001 From: Hemkumar Chheda Date: Mon, 13 Jul 2026 18:16:40 +0530 Subject: [PATCH 2/3] Lower malformed-log fallback logging to debug level _safe_build_structured_log_message runs once per stored hit inside the _read comprehension, so logging the fallback at warning level fires once per malformed line within a single log-fetch request. Drop it to debug so a request with many malformed entries does not flood the logs. --- .../airflow/providers/opensearch/log/os_task_handler.py | 2 +- .../tests/unit/opensearch/log/test_os_task_handler.py | 8 ++++---- 2 files changed, 5 insertions(+), 5 deletions(-) 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 7ab118c9a181a..da2c22a229a8d 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 @@ -136,7 +136,7 @@ def _safe_build_structured_log_message(hit_dict: dict[str, Any]) -> StructuredLo try: return StructuredLogMessage(**fields) except ValidationError: - logger.warning( + logger.debug( "Failed to parse stored log entry into StructuredLogMessage; falling back to " "stringified event. Offending fields: %s", fields, 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 4295f5260efe8..5730583d9db1c 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 @@ -420,7 +420,7 @@ def test_read_with_malformed_event_falls_back_to_stringified_event(self, ti): expected_sources=["http://localhost"], ) assert not metadata["end_of_log"] - mock_logger.warning.assert_called_once() + mock_logger.debug.assert_called_once() @pytest.mark.db_test def test_set_context(self, ti): @@ -867,12 +867,12 @@ def test_error_detail_dropped_when_empty(self): class TestSafeBuildStructuredLogMessage: - def test_string_event_returns_unchanged_and_does_not_warn(self): + def test_string_event_returns_unchanged_and_does_not_log(self): hit = {"event": "hello", "level": "info"} with patch("airflow.providers.opensearch.log.os_task_handler.logger") as mock_logger: result = _safe_build_structured_log_message(hit) assert result.event == "hello" - mock_logger.warning.assert_not_called() + mock_logger.debug.assert_not_called() def test_non_string_event_falls_back_to_stringified_event(self): hit = {"event": ["a", "b"], "timestamp": "2024-01-01T00:00:00Z"} @@ -880,4 +880,4 @@ def test_non_string_event_falls_back_to_stringified_event(self): result = _safe_build_structured_log_message(hit) assert result.event == str(["a", "b"]) assert result.timestamp is not None - mock_logger.warning.assert_called_once() + mock_logger.debug.assert_called_once() From 5812112988ed868ce6db25dc59909fbd496ca090 Mon Sep 17 00:00:00 2001 From: Hemkumar Chheda Date: Tue, 14 Jul 2026 14:48:11 +0530 Subject: [PATCH 3/3] Set task state to SUCCESS in malformed-event read test The shared ti fixture defaults to RUNNING, and _read() now delegates running/deferred tasks to the base handler's remote-log path, which is not wired up in unit tests. Match the sibling read tests by setting the task to SUCCESS so the test exercises the stored-log fallback path. --- .../opensearch/tests/unit/opensearch/log/test_os_task_handler.py | 1 + 1 file changed, 1 insertion(+) 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 5730583d9db1c..df24545783410 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 @@ -400,6 +400,7 @@ def test_read_with_custom_offset_and_host_fields(self, ti): @pytest.mark.skipif(not AIRFLOW_V_3_0_PLUS, reason="StructuredLogMessage fallback is Airflow 3+ only") @pytest.mark.db_test def test_read_with_malformed_event_falls_back_to_stringified_event(self, ti): + ti.state = TaskInstanceState.SUCCESS malformed_event = ["not", "a", "string"] malformed_source = { "message": self.test_message,