From f9954224f5978e456147bd933f77062b1b65b038 Mon Sep 17 00:00:00 2001 From: roshanprabu Date: Sun, 9 Aug 2026 01:45:30 +0530 Subject: [PATCH] Guard StackdriverRemoteLogIO's transport.send() against errors The log processor's transport.send() call ran unguarded. This processor is installed for every supervised component (scheduler, dag-processor, triggerer, workers) whenever REMOTE_TASK_LOG routes through Stackdriver, so any IAM permission error, gRPC connectivity failure, or quota exception from Cloud Logging propagated straight out of the structlog processor and crashed the entire process -- observed as a dag-processor stuck in CrashLoopBackOff from a missing logging.logEntries.create IAM binding. The read path already has this exact guard one function down (see the try/except around read_logs() a few lines below, with the same rationale in its comment); the write path was simply missing the equivalent. Log delivery is best-effort, so this logs a warning via the handler's own logger instead of raising. --- .../cloud/log/stackdriver_task_handler.py | 9 +++++- .../log/test_stackdriver_task_handler.py | 29 +++++++++++++++++++ 2 files changed, 37 insertions(+), 1 deletion(-) diff --git a/providers/google/src/airflow/providers/google/cloud/log/stackdriver_task_handler.py b/providers/google/src/airflow/providers/google/cloud/log/stackdriver_task_handler.py index 6c1772a88490c..1f1d6c6bbbf08 100644 --- a/providers/google/src/airflow/providers/google/cloud/log/stackdriver_task_handler.py +++ b/providers/google/src/airflow/providers/google/cloud/log/stackdriver_task_handler.py @@ -219,7 +219,14 @@ def proc( if map_index := event.get("map_index"): labels["map_index"] = str(map_index) - _transport.send(record, str(msg.get("event", "")), resource=self.resource, labels=labels) + try: + _transport.send(record, str(msg.get("event", "")), resource=self.resource, labels=labels) + except Exception: + # Cloud Logging unavailable / IAM glitch / gRPC error. This processor runs for + # every supervised component (scheduler, dag-processor, triggerer, workers), so + # letting this propagate would crash the whole process on a logging + # misconfiguration. Log delivery is best-effort; degrade gracefully instead. + _logger.warning("Failed to send log entry to Cloud Logging", exc_info=True) return event return (proc,) diff --git a/providers/google/tests/unit/google/cloud/log/test_stackdriver_task_handler.py b/providers/google/tests/unit/google/cloud/log/test_stackdriver_task_handler.py index c66ff24421b9d..922b3c8a6402f 100644 --- a/providers/google/tests/unit/google/cloud/log/test_stackdriver_task_handler.py +++ b/providers/google/tests/unit/google/cloud/log/test_stackdriver_task_handler.py @@ -367,6 +367,35 @@ def test_processors_sends_to_transport(self, mock_client, mock_get_creds_and_pro record = mock_transport.send.call_args[0][0] assert record.levelno == logging.INFO + @pytest.mark.skipif(not AIRFLOW_V_3_0_PLUS, reason="airflow.sdk.log only exists in Airflow 3+") + @mock.patch("airflow.providers.google.cloud.log.stackdriver_task_handler.get_credentials_and_project_id") + @mock.patch("airflow.providers.google.cloud.log.stackdriver_task_handler.gcp_logging.Client") + def test_processors_survives_transport_send_failure(self, mock_client, mock_get_creds_and_project_id): + """A Cloud Logging IAM/gRPC error must not propagate out of the log processor. + + This processor runs for every supervised component (scheduler, dag-processor, + triggerer, workers); letting a transport error propagate would crash the whole + process on a logging misconfiguration. + """ + mock_get_creds_and_project_id.return_value = ("creds", "project_id") + + mock_transport_type = mock.MagicMock() + mock_transport_type.return_value.send.side_effect = Exception("IAM permission denied") + with mock.patch("airflow.sdk.log.relative_path_from_logger", return_value="dag/task/1.log"): + io = StackdriverRemoteLogIO( + base_log_folder=self.local_log_location, + gcp_log_name="airflow", + transport_type=mock_transport_type, + ) + proc = io.processors[0] + + event = {"event": "hello world", "logger_name": "airflow.task"} + # Must not raise, despite transport.send() failing. + result = proc(mock.MagicMock(), "info", event) + + assert result is event + mock_transport_type.return_value.send.assert_called_once() + @pytest.mark.skipif(not AIRFLOW_V_3_0_PLUS, reason="airflow.sdk.log only exists in Airflow 3+") @mock.patch("airflow.providers.google.cloud.log.stackdriver_task_handler.get_credentials_and_project_id") @mock.patch("airflow.providers.google.cloud.log.stackdriver_task_handler.gcp_logging.Client")