From 4fb622562aa760c8486d0c2708d5a1ccf31a6ba4 Mon Sep 17 00:00:00 2001 From: Idris Akorede Ibrahim Date: Sat, 1 Aug 2026 10:43:27 +0100 Subject: [PATCH 1/6] Fix remote processor injection happening before dictConfig runs in configure_logging MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit In configure_logging(), getattr(remote, 'processors') was called before dictConfig() ran. dictConfig() calls _clearExistingHandlers() which closes every handler in logging._handlerList — including the remote handler built moments earlier. This caused CloudWatch/Watchtower logs to be silently dropped when using the Task SDK with remote logging on ECS/Kubernetes workers. Fix: move the remote processor injection to after dictConfig() has run via a second structlog.configure() call. Also gate on not sending_to_supervisor to avoid unmasked events from the task subprocess, and add a None default to getattr() to handle third-party RemoteLogIO objects. Provider-side fix: #68779 Closes #66475 --- task-sdk/src/airflow/sdk/log.py | 84 +++++++++----- task-sdk/tests/task_sdk/test_log.py | 166 ++++++++++++++++++++++++++++ 2 files changed, 221 insertions(+), 29 deletions(-) create mode 100644 task-sdk/tests/task_sdk/test_log.py diff --git a/task-sdk/src/airflow/sdk/log.py b/task-sdk/src/airflow/sdk/log.py index e105fd70d762e..5151b00ed4a49 100644 --- a/task-sdk/src/airflow/sdk/log.py +++ b/task-sdk/src/airflow/sdk/log.py @@ -17,7 +17,6 @@ # under the License. from __future__ import annotations -from contextlib import suppress from functools import cache from pathlib import Path from typing import TYPE_CHECKING, Any, BinaryIO, TextIO @@ -36,7 +35,7 @@ from airflow.sdk.types import Logger, RuntimeTaskInstanceProtocol as RuntimeTI -from airflow.sdk._shared.secrets_masker import _secrets_masker, redact +from airflow.sdk._shared.secrets_masker import redact class _ActiveLoggingConfig: @@ -119,8 +118,14 @@ def configure_logging( if mask_secrets: extra_processors += (mask_logs,) - if (remote := load_remote_log_handler()) and (remote_processors := getattr(remote, "processors")): - extra_processors += remote_processors + # NOTE: Do NOT call getattr(remote, "processors") here. + # Accessing remote.processors triggers creation of, for example, the watchtower CloudWatchLogHandler + # via a cached_property. The configure_logging() call below runs dictConfig() internally, + # which calls _clearExistingHandlers() -> logging.shutdown() on ALL existing handlers — + # including the watchtower handler we would have just built. The handler ends up dead + # (shutting_down=True) before a single task log is emitted. + # See: https://github.com/apache/airflow/issues/66475 + # Remote processors are injected AFTER dictConfig via a second structlog.configure() call below. configure_logging( json_output=json_output, @@ -134,6 +139,22 @@ def configure_logging( callsite_parameters=callsite_params, ) + # dictConfig has now run, so it is safe to build the remote handler. Re-inject the remote + # processors into the global structlog chain (before the final renderer) for parity with the old + # extra_process layout. Task-log streaming itself does not rely on this: it uses the + # file-backed logger built from logging_processors(), which loads remote.processors lazily. + if ( + not sending_to_supervisor + and (remote := load_remote_log_handler()) + and (remote_processors := getattr(remote, "processors", None)) + ): + current_processors = list(structlog.get_config()["processors"]) + # Insert before the final renderer + # NOTE: unlike the old extra_processors path, stdlib-routed records (ProcessorFormatter's + # foreign_pre_chain) intentionally do NOT pass through the remote processors. + updated_processors = current_processors[:-1] + list(remote_processors) + [current_processors[-1]] + structlog.configure(processors=updated_processors) + def logger_at_level(name: str, level: int) -> Logger: """Create a new logger at the given level.""" @@ -173,12 +194,7 @@ def init_log_file(local_relative_path: str) -> Path: def _load_logging_config() -> None: - """ - Load and cache the remote logging configuration from SDK config. - - SDK mirror of :func:`airflow.logging_config.load_logging_config` — see that - function for the ``logging_config_class`` / ``REMOTE_TASK_LOG`` contract. - """ + """Load and cache the remote logging configuration from SDK config.""" from airflow.sdk._shared.logging.factory import resolve_remote_task_log from airflow.sdk._shared.module_loading import import_string from airflow.sdk.configuration import conf @@ -231,20 +247,44 @@ def relative_path_from_logger(logger) -> Path | None: def upload_to_remote(logger: FilteringBoundLogger, ti: RuntimeTI | None = None): raw_logger = getattr(logger, "_logger") + # Dedicated logger for remote-upload visibility — operators relying on + # remote log handlers need a way to see when those handlers fail to load + # or fail to upload. + upload_log = structlog.get_logger("airflow.logging.remote") + + ti_id = str(ti.id) if ti else None handler = load_remote_log_handler() if not handler: + upload_log.warning( + "remote_log_handler_unavailable", + ti_id=ti_id, + note="Remote log handler could not be loaded; logs will be available locally only.", + ) return try: relative_path = relative_path_from_logger(raw_logger) - except Exception: + except Exception as exc: + upload_log.warning( + "remote_log_path_resolution_failed", + ti_id=ti_id, + exc_info=exc, + ) return if not relative_path: return log_relative_path = relative_path.as_posix() - handler.upload(log_relative_path, ti) + try: + handler.upload(log_relative_path, ti) + except Exception as exc: + upload_log.warning( + "remote_log_upload_failed", + ti_id=ti_id, + log_relative_path=log_relative_path, + exc_info=exc, + ) def mask_secret(secret: JsonValue, name: str | None = None) -> None: @@ -255,24 +295,10 @@ def mask_secret(secret: JsonValue, name: str | None = None) -> None: they're masked in both the task subprocess AND supervisor's log output. Works safely in both sync and async contexts. """ - _secrets_masker().add_mask(secret, name) + from contextlib import suppress - with suppress(Exception): - # Try to tell supervisor (only if in task execution context) - from airflow.sdk.execution_time import task_runner - from airflow.sdk.execution_time.comms import MaskSecret - - if comms := getattr(task_runner, "SUPERVISOR_COMMS", None): - comms.send(MaskSecret(value=secret, name=name)) + from airflow.sdk._shared.secrets_masker import _secrets_masker - -async def amask_secret(secret: JsonValue, name: str | None = None) -> None: - """ - Async version of mask_secret for use in async contexts. - - Uses asend() instead of send() to avoid deadlock when called from within - an async task that already has an asend() in flight. - """ _secrets_masker().add_mask(secret, name) with suppress(Exception): @@ -281,7 +307,7 @@ async def amask_secret(secret: JsonValue, name: str | None = None) -> None: from airflow.sdk.execution_time.comms import MaskSecret if comms := getattr(task_runner, "SUPERVISOR_COMMS", None): - await comms.asend(MaskSecret(value=secret, name=name)) + comms.send(MaskSecret(value=secret, name=name)) def reset_logging(): diff --git a/task-sdk/tests/task_sdk/test_log.py b/task-sdk/tests/task_sdk/test_log.py new file mode 100644 index 0000000000000..59b0b46cfb307 --- /dev/null +++ b/task-sdk/tests/task_sdk/test_log.py @@ -0,0 +1,166 @@ +# +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. +from __future__ import annotations + +from unittest import mock + +import structlog +import structlog.testing +from uuid6 import uuid7 + +from airflow.sdk import log as sdk_log + + +def _make_ti(): + ti = mock.MagicMock() + ti.id = uuid7() + return ti + + +def _make_logger(): + """Build a FilteringBoundLogger-like object exposing ``_logger``.""" + logger = mock.MagicMock() + logger._logger = mock.MagicMock() + return logger + + +class TestUploadToRemote: + def test_warns_when_handler_unavailable(self): + ti = _make_ti() + with ( + mock.patch.object(sdk_log, "load_remote_log_handler", return_value=None), + structlog.testing.capture_logs() as captured, + ): + sdk_log.upload_to_remote(_make_logger(), ti) + + events = [e for e in captured if e["event"] == "remote_log_handler_unavailable"] + assert len(events) == 1 + assert events[0]["log_level"] == "warning" + assert events[0]["ti_id"] == str(ti.id) + + def test_warns_when_path_resolution_fails(self): + ti = _make_ti() + handler = mock.MagicMock() + boom = RuntimeError("cannot resolve path") + with ( + mock.patch.object(sdk_log, "load_remote_log_handler", return_value=handler), + mock.patch.object(sdk_log, "relative_path_from_logger", side_effect=boom), + structlog.testing.capture_logs() as captured, + ): + sdk_log.upload_to_remote(_make_logger(), ti) + + events = [e for e in captured if e["event"] == "remote_log_path_resolution_failed"] + assert len(events) == 1 + assert events[0]["log_level"] == "warning" + assert events[0]["ti_id"] == str(ti.id) + assert events[0]["exc_info"] is boom + handler.upload.assert_not_called() + + def test_warns_when_upload_fails(self, tmp_path): + ti = _make_ti() + handler = mock.MagicMock() + boom = RuntimeError("s3 unreachable") + handler.upload.side_effect = boom + relative = tmp_path / "dag_id" / "run_id" / "task.log" + with ( + mock.patch.object(sdk_log, "load_remote_log_handler", return_value=handler), + mock.patch.object(sdk_log, "relative_path_from_logger", return_value=relative), + structlog.testing.capture_logs() as captured, + ): + sdk_log.upload_to_remote(_make_logger(), ti) + + events = [e for e in captured if e["event"] == "remote_log_upload_failed"] + assert len(events) == 1 + assert events[0]["log_level"] == "warning" + assert events[0]["ti_id"] == str(ti.id) + assert events[0]["log_relative_path"] == relative.as_posix() + assert events[0]["exc_info"] is boom + handler.upload.assert_called_once_with(relative.as_posix(), ti) + + def test_silent_when_relative_path_is_none(self): + ti = _make_ti() + handler = mock.MagicMock() + with ( + mock.patch.object(sdk_log, "load_remote_log_handler", return_value=handler), + mock.patch.object(sdk_log, "relative_path_from_logger", return_value=None), + structlog.testing.capture_logs() as captured, + ): + sdk_log.upload_to_remote(_make_logger(), ti) + + assert captured == [] + handler.upload.assert_not_called() + + def test_silent_on_success(self, tmp_path): + ti = _make_ti() + handler = mock.MagicMock() + relative = tmp_path / "dag_id" / "run_id" / "task.log" + with ( + mock.patch.object(sdk_log, "load_remote_log_handler", return_value=handler), + mock.patch.object(sdk_log, "relative_path_from_logger", return_value=relative), + structlog.testing.capture_logs() as captured, + ): + sdk_log.upload_to_remote(_make_logger(), ti) + + assert captured == [] + handler.upload.assert_called_once_with(relative.as_posix(), ti) + + +class TestConfigureLogging: + def test_remote_processors_injected_after_dictconfig(self): + """ + Regression test: remote processor injection must happen AFTER dictConfig() runs. + + dictConfig()'s non-incremental reset closes every handler in + logging._handlerList. If the remote handler is built before dictConfig + runs, it is closed before any task log is emitted and silently drops + all records. + """ + import airflow.sdk._shared.logging as shared_logging + + call_order = [] + + mock_handler = mock.MagicMock() + mock_handler.processors = (mock.MagicMock(),) + + def track_load_remote(): + call_order.append("load_remote_log_handler") + return mock_handler + + original_inner = shared_logging.configure_logging + + def track_inner_configure(*args, **kwargs): + call_order.append("dictConfig") + return original_inner(*args, **kwargs) + + # configure_logging is @cache decorated — clear it so the test actually runs the function + sdk_log.configure_logging.cache_clear() + + with ( + mock.patch.object(sdk_log, "load_remote_log_handler", side_effect=track_load_remote), + mock.patch.object(shared_logging, "configure_logging", side_effect=track_inner_configure), + ): + sdk_log.configure_logging() + + assert "dictConfig" in call_order, "inner configure_logging was never called" + assert "load_remote_log_handler" in call_order, "load_remote_log_handler() was never called" + dictconfig_pos = call_order.index("dictConfig") + load_remote_pos = call_order.index("load_remote_log_handler") + assert dictconfig_pos < load_remote_pos, ( + "load_remote_log_handler() must be called AFTER dictConfig() runs, " + "otherwise dictConfig closes the just-built handler before any task log is emitted" + ) From 565d30b6c957c9739df9c3a68945f2e7c8a28cfe Mon Sep 17 00:00:00 2001 From: Idris Akorede Ibrahim Date: Sun, 2 Aug 2026 06:53:57 +0100 Subject: [PATCH 2/6] Fix remote processor injection happening before dictConfig runs in configure_logging MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit In configure_logging(), getattr(remote, 'processors') was called before dictConfig() ran. dictConfig() calls _clearExistingHandlers() which closes every handler in logging._handlerList — including the remote handler built moments earlier. This caused CloudWatch/Watchtower logs to be silently dropped when using the Task SDK with remote logging on ECS/Kubernetes workers. Fix: move the remote processor injection to after dictConfig() has run via a second structlog.configure() call. Also gate on not sending_to_supervisor to avoid unmasked events from the task subprocess, and add a None default to getattr() to handle third-party RemoteLogIO objects. Closes #66475 --- task-sdk/src/airflow/sdk/log.py | 64 ++++++++++++++++----------------- 1 file changed, 30 insertions(+), 34 deletions(-) diff --git a/task-sdk/src/airflow/sdk/log.py b/task-sdk/src/airflow/sdk/log.py index 5151b00ed4a49..94f5b85dbb411 100644 --- a/task-sdk/src/airflow/sdk/log.py +++ b/task-sdk/src/airflow/sdk/log.py @@ -17,6 +17,7 @@ # under the License. from __future__ import annotations +from contextlib import suppress from functools import cache from pathlib import Path from typing import TYPE_CHECKING, Any, BinaryIO, TextIO @@ -35,7 +36,7 @@ from airflow.sdk.types import Logger, RuntimeTaskInstanceProtocol as RuntimeTI -from airflow.sdk._shared.secrets_masker import redact +from airflow.sdk._shared.secrets_masker import _secrets_masker, redact class _ActiveLoggingConfig: @@ -119,10 +120,10 @@ def configure_logging( extra_processors += (mask_logs,) # NOTE: Do NOT call getattr(remote, "processors") here. - # Accessing remote.processors triggers creation of, for example, the watchtower CloudWatchLogHandler + # Accessing remote.processors triggers creation of the remote handler # via a cached_property. The configure_logging() call below runs dictConfig() internally, # which calls _clearExistingHandlers() -> logging.shutdown() on ALL existing handlers — - # including the watchtower handler we would have just built. The handler ends up dead + # including the remote handler we would have just built. The handler ends up dead # (shutting_down=True) before a single task log is emitted. # See: https://github.com/apache/airflow/issues/66475 # Remote processors are injected AFTER dictConfig via a second structlog.configure() call below. @@ -141,7 +142,7 @@ def configure_logging( # dictConfig has now run, so it is safe to build the remote handler. Re-inject the remote # processors into the global structlog chain (before the final renderer) for parity with the old - # extra_process layout. Task-log streaming itself does not rely on this: it uses the + # extra_processors layout. Task-log streaming itself does not rely on this: it uses the # file-backed logger built from logging_processors(), which loads remote.processors lazily. if ( not sending_to_supervisor @@ -194,7 +195,12 @@ def init_log_file(local_relative_path: str) -> Path: def _load_logging_config() -> None: - """Load and cache the remote logging configuration from SDK config.""" + """ + Load and cache the remote logging configuration from SDK config. + + SDK mirror of :func:`airflow.logging_config.load_logging_config` — see that + function for the ``logging_config_class`` / ``REMOTE_TASK_LOG`` contract. + """ from airflow.sdk._shared.logging.factory import resolve_remote_task_log from airflow.sdk._shared.module_loading import import_string from airflow.sdk.configuration import conf @@ -247,44 +253,20 @@ def relative_path_from_logger(logger) -> Path | None: def upload_to_remote(logger: FilteringBoundLogger, ti: RuntimeTI | None = None): raw_logger = getattr(logger, "_logger") - # Dedicated logger for remote-upload visibility — operators relying on - # remote log handlers need a way to see when those handlers fail to load - # or fail to upload. - upload_log = structlog.get_logger("airflow.logging.remote") - - ti_id = str(ti.id) if ti else None handler = load_remote_log_handler() if not handler: - upload_log.warning( - "remote_log_handler_unavailable", - ti_id=ti_id, - note="Remote log handler could not be loaded; logs will be available locally only.", - ) return try: relative_path = relative_path_from_logger(raw_logger) - except Exception as exc: - upload_log.warning( - "remote_log_path_resolution_failed", - ti_id=ti_id, - exc_info=exc, - ) + except Exception: return if not relative_path: return log_relative_path = relative_path.as_posix() - try: - handler.upload(log_relative_path, ti) - except Exception as exc: - upload_log.warning( - "remote_log_upload_failed", - ti_id=ti_id, - log_relative_path=log_relative_path, - exc_info=exc, - ) + handler.upload(log_relative_path, ti) def mask_secret(secret: JsonValue, name: str | None = None) -> None: @@ -295,10 +277,24 @@ def mask_secret(secret: JsonValue, name: str | None = None) -> None: they're masked in both the task subprocess AND supervisor's log output. Works safely in both sync and async contexts. """ - from contextlib import suppress + _secrets_masker().add_mask(secret, name) + + with suppress(Exception): + # Try to tell supervisor (only if in task execution context) + from airflow.sdk.execution_time import task_runner + from airflow.sdk.execution_time.comms import MaskSecret + + if comms := getattr(task_runner, "SUPERVISOR_COMMS", None): + comms.send(MaskSecret(value=secret, name=name)) - from airflow.sdk._shared.secrets_masker import _secrets_masker +async def amask_secret(secret: JsonValue, name: str | None = None) -> None: + """ + Async version of mask_secret for use in async contexts. + + Uses asend() instead of send() to avoid deadlock when called from within + an async task that already has an asend() in flight. + """ _secrets_masker().add_mask(secret, name) with suppress(Exception): @@ -307,7 +303,7 @@ def mask_secret(secret: JsonValue, name: str | None = None) -> None: from airflow.sdk.execution_time.comms import MaskSecret if comms := getattr(task_runner, "SUPERVISOR_COMMS", None): - comms.send(MaskSecret(value=secret, name=name)) + await comms.asend(MaskSecret(value=secret, name=name)) def reset_logging(): From 24adaf745fe865ad4cc243a0cdaad979d3523702 Mon Sep 17 00:00:00 2001 From: Idris Akorede Ibrahim Date: Mon, 3 Aug 2026 23:13:41 +0100 Subject: [PATCH 3/6] Fix remote processor injection happening before dictConfig runs in configure_logging MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit In configure_logging(), getattr(remote, 'processors') was called before dictConfig() ran. dictConfig() calls _clearExistingHandlers() which closes every handler in logging._handlerList — including the remote handler built moments earlier. This caused CloudWatch/Watchtower logs to be silently dropped when using the Task SDK with remote logging on ECS/Kubernetes workers. Fix: move the remote processor injection to after dictConfig() has run via a second structlog.configure() call. Also gate on not sending_to_supervisor to avoid unmasked events from the task subprocess, and add a None default to getattr() to handle third-party RemoteLogIO objects. Closes #66475 --- task-sdk/src/airflow/sdk/log.py | 3 +- task-sdk/tests/task_sdk/test_log.py | 72 ++++++----------------------- 2 files changed, 15 insertions(+), 60 deletions(-) diff --git a/task-sdk/src/airflow/sdk/log.py b/task-sdk/src/airflow/sdk/log.py index 94f5b85dbb411..696ece269256f 100644 --- a/task-sdk/src/airflow/sdk/log.py +++ b/task-sdk/src/airflow/sdk/log.py @@ -120,8 +120,7 @@ def configure_logging( extra_processors += (mask_logs,) # NOTE: Do NOT call getattr(remote, "processors") here. - # Accessing remote.processors triggers creation of the remote handler - # via a cached_property. The configure_logging() call below runs dictConfig() internally, + # The configure_logging() call below runs dictConfig() internally, # which calls _clearExistingHandlers() -> logging.shutdown() on ALL existing handlers — # including the remote handler we would have just built. The handler ends up dead # (shutting_down=True) before a single task log is emitted. diff --git a/task-sdk/tests/task_sdk/test_log.py b/task-sdk/tests/task_sdk/test_log.py index 59b0b46cfb307..4d9ba44812ef0 100644 --- a/task-sdk/tests/task_sdk/test_log.py +++ b/task-sdk/tests/task_sdk/test_log.py @@ -40,58 +40,6 @@ def _make_logger(): class TestUploadToRemote: - def test_warns_when_handler_unavailable(self): - ti = _make_ti() - with ( - mock.patch.object(sdk_log, "load_remote_log_handler", return_value=None), - structlog.testing.capture_logs() as captured, - ): - sdk_log.upload_to_remote(_make_logger(), ti) - - events = [e for e in captured if e["event"] == "remote_log_handler_unavailable"] - assert len(events) == 1 - assert events[0]["log_level"] == "warning" - assert events[0]["ti_id"] == str(ti.id) - - def test_warns_when_path_resolution_fails(self): - ti = _make_ti() - handler = mock.MagicMock() - boom = RuntimeError("cannot resolve path") - with ( - mock.patch.object(sdk_log, "load_remote_log_handler", return_value=handler), - mock.patch.object(sdk_log, "relative_path_from_logger", side_effect=boom), - structlog.testing.capture_logs() as captured, - ): - sdk_log.upload_to_remote(_make_logger(), ti) - - events = [e for e in captured if e["event"] == "remote_log_path_resolution_failed"] - assert len(events) == 1 - assert events[0]["log_level"] == "warning" - assert events[0]["ti_id"] == str(ti.id) - assert events[0]["exc_info"] is boom - handler.upload.assert_not_called() - - def test_warns_when_upload_fails(self, tmp_path): - ti = _make_ti() - handler = mock.MagicMock() - boom = RuntimeError("s3 unreachable") - handler.upload.side_effect = boom - relative = tmp_path / "dag_id" / "run_id" / "task.log" - with ( - mock.patch.object(sdk_log, "load_remote_log_handler", return_value=handler), - mock.patch.object(sdk_log, "relative_path_from_logger", return_value=relative), - structlog.testing.capture_logs() as captured, - ): - sdk_log.upload_to_remote(_make_logger(), ti) - - events = [e for e in captured if e["event"] == "remote_log_upload_failed"] - assert len(events) == 1 - assert events[0]["log_level"] == "warning" - assert events[0]["ti_id"] == str(ti.id) - assert events[0]["log_relative_path"] == relative.as_posix() - assert events[0]["exc_info"] is boom - handler.upload.assert_called_once_with(relative.as_posix(), ti) - def test_silent_when_relative_path_is_none(self): ti = _make_ti() handler = mock.MagicMock() @@ -147,14 +95,22 @@ def track_inner_configure(*args, **kwargs): call_order.append("dictConfig") return original_inner(*args, **kwargs) - # configure_logging is @cache decorated — clear it so the test actually runs the function + # Save global structlog state so we can restore it after the test + original_processors = list(structlog.get_config()["processors"]) + sdk_log.configure_logging.cache_clear() - with ( - mock.patch.object(sdk_log, "load_remote_log_handler", side_effect=track_load_remote), - mock.patch.object(shared_logging, "configure_logging", side_effect=track_inner_configure), - ): - sdk_log.configure_logging() + try: + with ( + mock.patch.object(sdk_log, "load_remote_log_handler", side_effect=track_load_remote), + mock.patch.object(shared_logging, "configure_logging", side_effect=track_inner_configure), + ): + sdk_log.configure_logging() + finally: + # Restore global structlog processor chain and clear the cache so + # subsequent tests start from a clean state + structlog.configure(processors=original_processors) + sdk_log.configure_logging.cache_clear() assert "dictConfig" in call_order, "inner configure_logging was never called" assert "load_remote_log_handler" in call_order, "load_remote_log_handler() was never called" From 30dc125378c65e749cf81be5398d7475e7ff4d7a Mon Sep 17 00:00:00 2001 From: Idris Akorede Ibrahim Date: Tue, 22 Sep 2026 10:38:28 +0100 Subject: [PATCH 4/6] Keep remote task log handlers alive during setup dictConfig can close a remote handler before it receives its first task-log record, silently dropping remote logs. The logging setup must avoid that lifecycle conflict while remaining compatible with remote handlers that do not provide custom processors. --- airflow-core/newsfragments/66633.bugfix.rst | 1 + task-sdk/src/airflow/sdk/log.py | 18 +-- task-sdk/tests/task_sdk/test_log.py | 136 +++++++------------- 3 files changed, 51 insertions(+), 104 deletions(-) create mode 100644 airflow-core/newsfragments/66633.bugfix.rst diff --git a/airflow-core/newsfragments/66633.bugfix.rst b/airflow-core/newsfragments/66633.bugfix.rst new file mode 100644 index 0000000000000..5ea85be3152d7 --- /dev/null +++ b/airflow-core/newsfragments/66633.bugfix.rst @@ -0,0 +1 @@ +Fix remote task logs being silently dropped when ``dictConfig`` closes a remote logging handler during configuration. diff --git a/task-sdk/src/airflow/sdk/log.py b/task-sdk/src/airflow/sdk/log.py index 696ece269256f..849a7eca20932 100644 --- a/task-sdk/src/airflow/sdk/log.py +++ b/task-sdk/src/airflow/sdk/log.py @@ -74,7 +74,7 @@ def logging_processors( if mask_secrets: extra_processors += (mask_logs,) - if (remote := load_remote_log_handler()) and (remote_processors := getattr(remote, "processors")): + if (remote := load_remote_log_handler()) and (remote_processors := getattr(remote, "processors", None)): extra_processors += remote_processors procs, _, final_writer = structlog_processors( @@ -119,14 +119,6 @@ def configure_logging( if mask_secrets: extra_processors += (mask_logs,) - # NOTE: Do NOT call getattr(remote, "processors") here. - # The configure_logging() call below runs dictConfig() internally, - # which calls _clearExistingHandlers() -> logging.shutdown() on ALL existing handlers — - # including the remote handler we would have just built. The handler ends up dead - # (shutting_down=True) before a single task log is emitted. - # See: https://github.com/apache/airflow/issues/66475 - # Remote processors are injected AFTER dictConfig via a second structlog.configure() call below. - configure_logging( json_output=json_output, log_level=log_level, @@ -139,19 +131,13 @@ def configure_logging( callsite_parameters=callsite_params, ) - # dictConfig has now run, so it is safe to build the remote handler. Re-inject the remote - # processors into the global structlog chain (before the final renderer) for parity with the old - # extra_processors layout. Task-log streaming itself does not rely on this: it uses the - # file-backed logger built from logging_processors(), which loads remote.processors lazily. + # Build the remote handler after dictConfig(), which closes every previously registered handler. if ( not sending_to_supervisor and (remote := load_remote_log_handler()) and (remote_processors := getattr(remote, "processors", None)) ): current_processors = list(structlog.get_config()["processors"]) - # Insert before the final renderer - # NOTE: unlike the old extra_processors path, stdlib-routed records (ProcessorFormatter's - # foreign_pre_chain) intentionally do NOT pass through the remote processors. updated_processors = current_processors[:-1] + list(remote_processors) + [current_processors[-1]] structlog.configure(processors=updated_processors) diff --git a/task-sdk/tests/task_sdk/test_log.py b/task-sdk/tests/task_sdk/test_log.py index 4d9ba44812ef0..4203c6d7b4fb7 100644 --- a/task-sdk/tests/task_sdk/test_log.py +++ b/task-sdk/tests/task_sdk/test_log.py @@ -19,104 +19,64 @@ from unittest import mock +import pytest import structlog -import structlog.testing -from uuid6 import uuid7 from airflow.sdk import log as sdk_log -def _make_ti(): - ti = mock.MagicMock() - ti.id = uuid7() - return ti - - -def _make_logger(): - """Build a FilteringBoundLogger-like object exposing ``_logger``.""" - logger = mock.MagicMock() - logger._logger = mock.MagicMock() - return logger - - -class TestUploadToRemote: - def test_silent_when_relative_path_is_none(self): - ti = _make_ti() - handler = mock.MagicMock() - with ( - mock.patch.object(sdk_log, "load_remote_log_handler", return_value=handler), - mock.patch.object(sdk_log, "relative_path_from_logger", return_value=None), - structlog.testing.capture_logs() as captured, - ): - sdk_log.upload_to_remote(_make_logger(), ti) - - assert captured == [] - handler.upload.assert_not_called() - - def test_silent_on_success(self, tmp_path): - ti = _make_ti() - handler = mock.MagicMock() - relative = tmp_path / "dag_id" / "run_id" / "task.log" - with ( - mock.patch.object(sdk_log, "load_remote_log_handler", return_value=handler), - mock.patch.object(sdk_log, "relative_path_from_logger", return_value=relative), - structlog.testing.capture_logs() as captured, - ): - sdk_log.upload_to_remote(_make_logger(), ti) - - assert captured == [] - handler.upload.assert_called_once_with(relative.as_posix(), ti) - - class TestConfigureLogging: - def test_remote_processors_injected_after_dictconfig(self): - """ - Regression test: remote processor injection must happen AFTER dictConfig() runs. - - dictConfig()'s non-incremental reset closes every handler in - logging._handlerList. If the remote handler is built before dictConfig - runs, it is closed before any task log is emitted and silently drops - all records. - """ - import airflow.sdk._shared.logging as shared_logging - - call_order = [] - - mock_handler = mock.MagicMock() - mock_handler.processors = (mock.MagicMock(),) - - def track_load_remote(): - call_order.append("load_remote_log_handler") - return mock_handler - - original_inner = shared_logging.configure_logging - - def track_inner_configure(*args, **kwargs): - call_order.append("dictConfig") - return original_inner(*args, **kwargs) - - # Save global structlog state so we can restore it after the test - original_processors = list(structlog.get_config()["processors"]) - + @pytest.mark.parametrize("sending_to_supervisor", [True, False]) + @mock.patch("airflow.sdk.log.load_remote_log_handler") + @mock.patch("airflow.sdk._shared.logging.configure_logging") + def test_injects_remote_processors_after_dictconfig( + self, mock_configure_logging, mock_load_remote_log_handler, sending_to_supervisor + ): + initial_processor = mock.Mock() + final_renderer = mock.Mock() + remote_processor = mock.Mock() + remote_handler = mock.Mock(spec=["processors"]) + remote_handler.processors = (remote_processor,) + calls = [] + + def configure_structlog(**_): + calls.append("dictConfig") + structlog.configure(processors=[initial_processor, final_renderer]) + + def load_remote_handler(): + calls.append("load_remote_log_handler") + return remote_handler + + mock_configure_logging.side_effect = configure_structlog + mock_load_remote_log_handler.side_effect = load_remote_handler + original_processors = structlog.get_config()["processors"] sdk_log.configure_logging.cache_clear() try: - with ( - mock.patch.object(sdk_log, "load_remote_log_handler", side_effect=track_load_remote), - mock.patch.object(shared_logging, "configure_logging", side_effect=track_inner_configure), - ): - sdk_log.configure_logging() + sdk_log.configure_logging(sending_to_supervisor=sending_to_supervisor) + + assert mock_configure_logging.called + if sending_to_supervisor: + assert calls == ["dictConfig"] + assert structlog.get_config()["processors"] == [initial_processor, final_renderer] + else: + assert calls == ["dictConfig", "load_remote_log_handler"] + assert structlog.get_config()["processors"] == [ + initial_processor, + remote_processor, + final_renderer, + ] finally: - # Restore global structlog processor chain and clear the cache so - # subsequent tests start from a clean state structlog.configure(processors=original_processors) sdk_log.configure_logging.cache_clear() - assert "dictConfig" in call_order, "inner configure_logging was never called" - assert "load_remote_log_handler" in call_order, "load_remote_log_handler() was never called" - dictconfig_pos = call_order.index("dictConfig") - load_remote_pos = call_order.index("load_remote_log_handler") - assert dictconfig_pos < load_remote_pos, ( - "load_remote_log_handler() must be called AFTER dictConfig() runs, " - "otherwise dictConfig closes the just-built handler before any task log is emitted" - ) + @mock.patch("airflow.sdk.log.load_remote_log_handler", return_value=object()) + def test_allows_remote_handler_without_processors(self, mock_load_remote_log_handler): + sdk_log.logging_processors.cache_clear() + + try: + sdk_log.logging_processors(json_output=False) + finally: + sdk_log.logging_processors.cache_clear() + + mock_load_remote_log_handler.assert_called_once_with() From 5b4fe19933f3c59bec5b02788616a2bf9d7f9591 Mon Sep 17 00:00:00 2001 From: Idris Akorede Ibrahim Date: Wed, 23 Sep 2026 09:46:36 +0100 Subject: [PATCH 5/6] Cover remote handlers without custom processors Handlers supplied by third-party remote logging implementations may not define custom processors. Keep the regression coverage explicit so that compatibility does not silently change. --- task-sdk/tests/task_sdk/test_log.py | 11 +++++++++-- 1 file changed, 9 insertions(+), 2 deletions(-) diff --git a/task-sdk/tests/task_sdk/test_log.py b/task-sdk/tests/task_sdk/test_log.py index 4203c6d7b4fb7..cab5f6548c95a 100644 --- a/task-sdk/tests/task_sdk/test_log.py +++ b/task-sdk/tests/task_sdk/test_log.py @@ -70,13 +70,20 @@ def load_remote_handler(): structlog.configure(processors=original_processors) sdk_log.configure_logging.cache_clear() + @mock.patch("airflow.sdk._shared.logging.structlog.structlog_processors") @mock.patch("airflow.sdk.log.load_remote_log_handler", return_value=object()) - def test_allows_remote_handler_without_processors(self, mock_load_remote_log_handler): + def test_allows_remote_handler_without_processors( + self, mock_load_remote_log_handler, mock_structlog_processors + ): + initial_processor = mock.Mock() + final_renderer = mock.Mock() + mock_structlog_processors.return_value = ([initial_processor], None, final_renderer) sdk_log.logging_processors.cache_clear() try: - sdk_log.logging_processors(json_output=False) + processors = sdk_log.logging_processors(json_output=False) finally: sdk_log.logging_processors.cache_clear() mock_load_remote_log_handler.assert_called_once_with() + assert processors == (initial_processor, final_renderer) From 5e939720824d321d84a2a2fd1be76617cafa8b6a Mon Sep 17 00:00:00 2001 From: Idris Akorede Ibrahim Date: Wed, 23 Sep 2026 11:25:52 +0100 Subject: [PATCH 6/6] Cover remote handlers without custom processors Handlers supplied by third-party remote logging implementations may not define custom processors. Keep the regression coverage explicit so that compatibility does not silently change. --- task-sdk/tests/task_sdk/test_log.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/task-sdk/tests/task_sdk/test_log.py b/task-sdk/tests/task_sdk/test_log.py index cab5f6548c95a..aa955a87a60a4 100644 --- a/task-sdk/tests/task_sdk/test_log.py +++ b/task-sdk/tests/task_sdk/test_log.py @@ -86,4 +86,4 @@ def test_allows_remote_handler_without_processors( sdk_log.logging_processors.cache_clear() mock_load_remote_log_handler.assert_called_once_with() - assert processors == (initial_processor, final_renderer) + assert processors == (initial_processor, sdk_log.mask_logs, final_renderer)