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]

Reply via email to