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