vincbeck commented on code in PR #73580:
URL: https://github.com/apache/airflow/pull/73580#discussion_r4207624040


##########
airflow-core/src/airflow/utils/log/callback_log_reader.py:
##########
@@ -0,0 +1,110 @@
+# 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.
+"""Reader for callback execution logs stored in remote or local storage."""
+
+from __future__ import annotations
+
+import os
+import re
+from collections.abc import Generator
+from contextlib import suppress
+from pathlib import Path
+from typing import TYPE_CHECKING
+
+from airflow.configuration import conf
+from airflow.utils.log.file_task_handler import (
+    FileTaskHandler,
+    StructuredLogMessage,
+    _get_compatible_log_stream,
+    _interleave_logs,
+)
+
+if TYPE_CHECKING:
+    from airflow._shared.logging.remote import LogSourceInfo, RawLogStream, 
StreamingLogResponse
+
+_SAFE_PATH_COMPONENT = re.compile(r"[A-Za-z0-9._:+\-~@]+")
+
+
+def validate_log_path_component(component: str) -> str:
+    """Validate a single log path component, raising ValueError if it could 
escape the log folder."""
+    if component in (".", "..") or not 
_SAFE_PATH_COMPONENT.fullmatch(component):
+        raise ValueError(f"Invalid log path component: {component!r}")
+    return component
+
+
+def read_callback_log(

Review Comment:
   Just to be sure I understand because I think this is important. You are 
using the task log if it is available to log the callback logs, and as a 
fallback use the local system (a file). Am I correct?
   
   You are basically trying to address a bug issue in Airflow, having logs 
available to something unrelated to a task. I am pretty sure we will have this 
issue later for something else (e.g. trigger). In that case, and that might 
require 2 different PRs to make it simpler. Can we have a common and generic 
mechanism to do so. Example, the name of this function is `read_callback_log`, 
which makes it specific to callbacks, but could we use it for something like as 
well (e.g. triggers)?



##########
airflow-core/src/airflow/utils/log/callback_log_reader.py:
##########
@@ -0,0 +1,110 @@
+# 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.
+"""Reader for callback execution logs stored in remote or local storage."""
+
+from __future__ import annotations
+
+import os
+import re
+from collections.abc import Generator
+from contextlib import suppress
+from pathlib import Path
+from typing import TYPE_CHECKING
+
+from airflow.configuration import conf
+from airflow.utils.log.file_task_handler import (
+    FileTaskHandler,
+    StructuredLogMessage,
+    _get_compatible_log_stream,
+    _interleave_logs,
+)
+
+if TYPE_CHECKING:
+    from airflow._shared.logging.remote import LogSourceInfo, RawLogStream, 
StreamingLogResponse
+
+_SAFE_PATH_COMPONENT = re.compile(r"[A-Za-z0-9._:+\-~@]+")
+
+
+def validate_log_path_component(component: str) -> str:
+    """Validate a single log path component, raising ValueError if it could 
escape the log folder."""
+    if component in (".", "..") or not 
_SAFE_PATH_COMPONENT.fullmatch(component):
+        raise ValueError(f"Invalid log path component: {component!r}")
+    return component
+
+
+def read_callback_log(
+    dag_id: str, run_id: str, callback_id: str
+) -> Generator[StructuredLogMessage, None, None]:
+    """
+    Stream callback logs, trying remote storage first and then the local 
filesystem.
+
+    Executor callbacks log to ``executor_callbacks/...`` (see 
``ExecuteCallback.make()``) and
+    triggerer callbacks to ``triggerer_callbacks/...`` (see 
``TriggerLoggingFactory``).
+    """
+    for component in (dag_id, run_id, callback_id):
+        validate_log_path_component(component)
+
+    sources: LogSourceInfo = []
+    log_streams: list[RawLogStream] = []
+
+    for prefix in ("executor_callbacks", "triggerer_callbacks"):
+        relative_path = f"{prefix}/{dag_id}/{run_id}/{callback_id}"
+        with suppress(Exception):
+            remote_sources, remote_log_streams = 
_read_callback_remote_logs(relative_path)
+            sources.extend(remote_sources)
+            log_streams.extend(remote_log_streams)
+
+        if not log_streams:
+            local_sources, local_log_streams = 
_read_callback_local_logs(relative_path)
+            sources.extend(local_sources)
+            log_streams.extend(local_log_streams)
+
+        if log_streams:
+            break
+
+    if not log_streams:
+        yield StructuredLogMessage(event="No callback logs found.")
+        return
+
+    yield StructuredLogMessage(event="::group::Log message source details", 
sources=sources)  # type: ignore[call-arg]
+    yield StructuredLogMessage(event="::endgroup::")
+    yield from _interleave_logs(*log_streams)
+
+
+def _read_callback_remote_logs(relative_path: str) -> StreamingLogResponse:
+    from airflow.logging_config import get_remote_task_log
+
+    remote_io = get_remote_task_log()
+    if remote_io is None:
+        return [], []
+
+    # Callbacks have no TaskInstance; remote handlers only use ``ti`` for 
optional metadata.
+    if stream_method := getattr(remote_io, "stream", None):

Review Comment:
   Can you explain more/better why checking if `stream` is available make the 
remote logger a good fit?



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