potiuk commented on code in PR #66633:
URL: https://github.com/apache/airflow/pull/66633#discussion_r4064350912
##########
task-sdk/src/airflow/sdk/log.py:
##########
@@ -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_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
+ and (remote := load_remote_log_handler())
+ and (remote_processors := getattr(remote, "processors", None))
Review Comment:
This is the half that got fixed. The identical line in
`logging_processors()` didn't:
```python
77: if (remote := load_remote_log_handler()) and (remote_processors :=
getattr(remote, "processors")):
```
Your description says the PR adds *"a `None` default to `getattr` to handle
third-party `RemoteLogIO` objects missing the `processors` attribute"* — but
`logging_processors()` is reached from `supervisor.py:2593`,
`callback_supervisor.py:391`, `dag_processing/manager.py:1437` and
`triggerer_job_runner.py:1035`, so a third-party handler without `.processors`
still raises `AttributeError` in the supervisor, the Dag processor and the
triggerer. Those are arguably the paths where it matters most.
One token (`, None`) at line 77 and the claim in the description becomes
true.
##########
task-sdk/tests/task_sdk/test_log.py:
##########
@@ -0,0 +1,122 @@
+#
+# 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:
Review Comment:
`upload_to_remote` is at `log.py:232` and this PR doesn't touch it; it's
also already covered indirectly by `test_supervisor.py` and
`test_callback_supervisor.py`. Per `AGENTS.md` § Testing Standards: *"Do not
add tests for pre-existing logic that was already present before the PR."*
I realise the pull here is that you're creating `test_log.py` from scratch
and it feels bare with one test in it — but scope is scope, and the effort is
better spent on the two gaps in my comment on line 119.
Style nits while you're in the file, both from the same section:
`@mock.patch` decorators are preferred over `with mock.patch(...)` context
managers, and mocks should carry `spec`/`autospec` — `mock.MagicMock()` for the
handler, the TI and the logger will accept any attribute you ask of it,
including ones the real object doesn't have.
##########
task-sdk/tests/task_sdk/test_log.py:
##########
@@ -0,0 +1,122 @@
+#
+# 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_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"])
+
+ 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()
+ 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, (
Review Comment:
This assertion is the weakest property the test could have checked. It
confirms `load_remote_log_handler` was *called* after `dictConfig`, and nothing
else.
Two consequences:
- The insertion at `log.py:155` — `current_processors[:-1] +
list(remote_processors) + [current_processors[-1]]` — is the fiddliest line in
the diff and is completely unpinned. This test passes just as happily if the
processors are appended *after* the renderer, where they'd silently never run.
Asserting the mock processor's presence and position in
`structlog.get_config()["processors"]` would catch that, and it's barely more
code than you already have.
- The `not sending_to_supervisor` gate at `log.py:147` is new behaviour, not
a move, and nothing exercises it. Inverting or dropping it would be caught by
nothing. Per `AGENTS.md` § Testing Standards, *"Every changed or added
behaviour must have a test"* — a parametrized case over
`sending_to_supervisor=True/False` would cover it.
##########
task-sdk/src/airflow/sdk/log.py:
##########
@@ -119,8 +119,13 @@ 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.
Review Comment:
Two things about this comment block.
First, it explains code that is no longer here, and then the same
explanation is given again at line 142. A regression guard is legitimate
context — *"don't move this back above `configure_logging()`"* is worth
recording — but one of the two blocks should carry it, not both. `AGENTS.md` §
Coding Standards: *"Do not narrate code, repeat code."*
Second, a question about the last line: *"Remote processors are injected
AFTER dictConfig"* — true here, but `logging_processors()` at line 62 is the
same bug's other door. It's `@cache`d and its first call runs
`load_remote_log_handler()` → `_load_logging_config()`, which builds the
handler and registers it in `logging._handlerList`. Anything reaching it before
this function's `dictConfig()` gets the original bug back. I checked the four
callers and they all run after startup configuration, so I don't think it's
live — but nothing records that invariant, and your comment calling that path
"lazy" reads as if it's inherently safe. Lazy only means "on first call".
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]