diff --git a/providers/alibaba/newsfragments/70479.bugfix.rst b/providers/alibaba/newsfragments/70479.bugfix.rst new file mode 100644 index 0000000000000..b11913df8fb18 --- /dev/null +++ b/providers/alibaba/newsfragments/70479.bugfix.rst @@ -0,0 +1 @@ +Fix remote log handler crashing the triggerer's logging thread when delete_local_copy is set: delete the uploaded log file and prune empty parent dirs instead of shutil.rmtree of the shared directory, which raced concurrent handler closes and raised FileNotFoundError diff --git a/providers/alibaba/src/airflow/providers/alibaba/cloud/log/oss_task_handler.py b/providers/alibaba/src/airflow/providers/alibaba/cloud/log/oss_task_handler.py index a000c6ff01f36..ce03ff406791a 100644 --- a/providers/alibaba/src/airflow/providers/alibaba/cloud/log/oss_task_handler.py +++ b/providers/alibaba/src/airflow/providers/alibaba/cloud/log/oss_task_handler.py @@ -19,7 +19,6 @@ import contextlib import os -import shutil from functools import cached_property from pathlib import Path from typing import TYPE_CHECKING @@ -59,7 +58,16 @@ def upload(self, path: os.PathLike | str, ti: RuntimeTI | None = None) -> None: log = local_loc.read_text() has_uploaded = self.oss_write(log, relative_path) if has_uploaded and self.delete_local_copy: - shutil.rmtree(os.path.dirname(local_loc)) + # Delete the file and prune empty parents (idempotent); a bare + # rmtree of the dir races concurrent handler closes in the triggerer. + local_loc.unlink(missing_ok=True) + parent = local_loc.parent + while parent != self.base_log_folder and parent.is_dir(): + if any(parent.iterdir()): + break + with contextlib.suppress(OSError): + parent.rmdir() + parent = parent.parent @cached_property def base_folder(self): diff --git a/providers/amazon/newsfragments/70479.bugfix.rst b/providers/amazon/newsfragments/70479.bugfix.rst new file mode 100644 index 0000000000000..b11913df8fb18 --- /dev/null +++ b/providers/amazon/newsfragments/70479.bugfix.rst @@ -0,0 +1 @@ +Fix remote log handler crashing the triggerer's logging thread when delete_local_copy is set: delete the uploaded log file and prune empty parent dirs instead of shutil.rmtree of the shared directory, which raced concurrent handler closes and raised FileNotFoundError diff --git a/providers/amazon/src/airflow/providers/amazon/aws/log/s3_task_handler.py b/providers/amazon/src/airflow/providers/amazon/aws/log/s3_task_handler.py index 6ef82f301da4b..b72672938c076 100644 --- a/providers/amazon/src/airflow/providers/amazon/aws/log/s3_task_handler.py +++ b/providers/amazon/src/airflow/providers/amazon/aws/log/s3_task_handler.py @@ -17,11 +17,11 @@ # under the License. from __future__ import annotations +import contextlib import inspect import logging import os import pathlib -import shutil from functools import cached_property from typing import TYPE_CHECKING @@ -87,7 +87,16 @@ def upload(self, path: os.PathLike | str, ti: RuntimeTI | None = None) -> None: log = local_loc.read_text() has_uploaded = self.write(log, remote_loc) if has_uploaded and self.delete_local_copy: - shutil.rmtree(os.path.dirname(local_loc)) + # Delete the file and prune empty parents (idempotent); a bare + # rmtree of the dir races concurrent handler closes in the triggerer. + local_loc.unlink(missing_ok=True) + parent = local_loc.parent + while parent != self.base_log_folder and parent.is_dir(): + if any(parent.iterdir()): + break + with contextlib.suppress(OSError): + parent.rmdir() + parent = parent.parent elif has_uploaded: local_loc.write_text("") diff --git a/providers/amazon/tests/unit/amazon/aws/log/test_s3_task_handler.py b/providers/amazon/tests/unit/amazon/aws/log/test_s3_task_handler.py index 0576756c1965f..05421d0a08ca9 100644 --- a/providers/amazon/tests/unit/amazon/aws/log/test_s3_task_handler.py +++ b/providers/amazon/tests/unit/amazon/aws/log/test_s3_task_handler.py @@ -486,6 +486,32 @@ def test_close_with_delete_local_logs_conf(self, delete_local_copy, expected_exi handler.close() assert os.path.exists(handler.handler.baseFilename) == expected_existence_of_local_copy + def test_upload_delete_local_copy_is_idempotent_and_keeps_shared_parent(self): + # The triggerer runs many trigger log handlers concurrently and close() is not + # atomic, so upload() can run twice for the same path. Deleting the whole parent + # dir (shutil.rmtree) made the second call raise FileNotFoundError, killing the + # logging listener thread. Deleting the file and pruning empty parents must be + # idempotent and must not remove a parent that still holds a sibling's log. + log_io = S3RemoteLogIO( + remote_base=self.remote_log_base, + base_log_folder=self.local_log_location, + delete_local_copy=True, + ) + ti_dir = pathlib.Path(self.local_log_location) / "dag_id=d" / "run_id=r" / "task_id=t" + ti_dir.mkdir(parents=True, exist_ok=True) + local_loc = ti_dir / "attempt=1.log" + local_loc.write_text("log content") + sibling = ti_dir / "attempt=1.log.trigger.99.log" + sibling.write_text("sibling still open") + + log_io.upload(local_loc, ti=None) + assert not local_loc.exists(), "uploaded log file should be deleted" + assert ti_dir.is_dir(), "parent must survive while a sibling log remains" + assert sibling.exists(), "sibling log must be untouched" + + # Second upload of the already-removed file must not raise (idempotent). + log_io.upload(local_loc, ti=None) + def test_filename_template_for_backward_compatibility(self): # filename_template arg support for running the latest provider on airflow 2 S3TaskHandler(self.local_log_location, self.remote_log_base, filename_template=None) diff --git a/providers/apache/hdfs/newsfragments/70479.bugfix.rst b/providers/apache/hdfs/newsfragments/70479.bugfix.rst new file mode 100644 index 0000000000000..b11913df8fb18 --- /dev/null +++ b/providers/apache/hdfs/newsfragments/70479.bugfix.rst @@ -0,0 +1 @@ +Fix remote log handler crashing the triggerer's logging thread when delete_local_copy is set: delete the uploaded log file and prune empty parent dirs instead of shutil.rmtree of the shared directory, which raced concurrent handler closes and raised FileNotFoundError diff --git a/providers/apache/hdfs/src/airflow/providers/apache/hdfs/log/hdfs_task_handler.py b/providers/apache/hdfs/src/airflow/providers/apache/hdfs/log/hdfs_task_handler.py index 22304637db237..23ac062caa56a 100644 --- a/providers/apache/hdfs/src/airflow/providers/apache/hdfs/log/hdfs_task_handler.py +++ b/providers/apache/hdfs/src/airflow/providers/apache/hdfs/log/hdfs_task_handler.py @@ -17,9 +17,9 @@ # under the License. from __future__ import annotations +import contextlib import logging import os -import shutil from functools import cached_property from pathlib import Path from typing import TYPE_CHECKING @@ -59,7 +59,16 @@ def upload(self, path: os.PathLike | str, ti: RuntimeTI | None = None) -> None: if local_loc.is_file(): self.hook.load_file(local_loc, remote_loc) if self.delete_local_copy: - shutil.rmtree(os.path.dirname(local_loc)) + # Delete the file and prune empty parents (idempotent); a bare + # rmtree of the dir races concurrent handler closes in the triggerer. + local_loc.unlink(missing_ok=True) + parent = local_loc.parent + while parent != self.base_log_folder and parent.is_dir(): + if any(parent.iterdir()): + break + with contextlib.suppress(OSError): + parent.rmdir() + parent = parent.parent @cached_property def hook(self): diff --git a/providers/elasticsearch/newsfragments/70479.bugfix.rst b/providers/elasticsearch/newsfragments/70479.bugfix.rst new file mode 100644 index 0000000000000..b11913df8fb18 --- /dev/null +++ b/providers/elasticsearch/newsfragments/70479.bugfix.rst @@ -0,0 +1 @@ +Fix remote log handler crashing the triggerer's logging thread when delete_local_copy is set: delete the uploaded log file and prune empty parent dirs instead of shutil.rmtree of the shared directory, which raced concurrent handler closes and raised FileNotFoundError diff --git a/providers/elasticsearch/src/airflow/providers/elasticsearch/log/es_task_handler.py b/providers/elasticsearch/src/airflow/providers/elasticsearch/log/es_task_handler.py index 2fb610186e9d9..d3a315061ca39 100644 --- a/providers/elasticsearch/src/airflow/providers/elasticsearch/log/es_task_handler.py +++ b/providers/elasticsearch/src/airflow/providers/elasticsearch/log/es_task_handler.py @@ -22,7 +22,6 @@ import json import logging import os -import shutil import sys import time import warnings @@ -753,7 +752,16 @@ def upload(self, path: os.PathLike | str, ti: RuntimeTI | None = None) -> None: log_lines = self._parse_raw_log(local_loc.read_text(), log_id) success = self._write_to_es(log_lines) if success and self.delete_local_copy: - shutil.rmtree(os.path.dirname(local_loc)) + # Delete the file and prune empty parents (idempotent); a bare + # rmtree of the dir races concurrent handler closes in the triggerer. + local_loc.unlink(missing_ok=True) + parent = local_loc.parent + while parent != self.base_log_folder and parent.is_dir(): + if any(parent.iterdir()): + break + with contextlib.suppress(OSError): + parent.rmdir() + parent = parent.parent def _parse_raw_log(self, log: str, log_id: str) -> list[dict[str, Any]]: logs = log.split("\n") diff --git a/providers/google/newsfragments/70479.bugfix.rst b/providers/google/newsfragments/70479.bugfix.rst new file mode 100644 index 0000000000000..b11913df8fb18 --- /dev/null +++ b/providers/google/newsfragments/70479.bugfix.rst @@ -0,0 +1 @@ +Fix remote log handler crashing the triggerer's logging thread when delete_local_copy is set: delete the uploaded log file and prune empty parent dirs instead of shutil.rmtree of the shared directory, which raced concurrent handler closes and raised FileNotFoundError diff --git a/providers/google/src/airflow/providers/google/cloud/log/gcs_task_handler.py b/providers/google/src/airflow/providers/google/cloud/log/gcs_task_handler.py index 42c4ba9a61b41..265319b669aad 100644 --- a/providers/google/src/airflow/providers/google/cloud/log/gcs_task_handler.py +++ b/providers/google/src/airflow/providers/google/cloud/log/gcs_task_handler.py @@ -17,9 +17,9 @@ # under the License. from __future__ import annotations +import contextlib import logging import os -import shutil from collections.abc import Collection from functools import cached_property from pathlib import Path @@ -85,7 +85,16 @@ def upload(self, path: os.PathLike | str, ti: RuntimeTI | None = None) -> None: log = local_loc.read_text() has_uploaded = self.write(log, remote_loc) if has_uploaded and self.delete_local_copy: - shutil.rmtree(os.path.dirname(local_loc)) + # Delete the file and prune empty parents (idempotent); a bare + # rmtree of the dir races concurrent handler closes in the triggerer. + local_loc.unlink(missing_ok=True) + parent = local_loc.parent + while parent != self.base_log_folder and parent.is_dir(): + if any(parent.iterdir()): + break + with contextlib.suppress(OSError): + parent.rmdir() + parent = parent.parent @cached_property def hook(self) -> GCSHook | None: diff --git a/providers/microsoft/azure/newsfragments/70479.bugfix.rst b/providers/microsoft/azure/newsfragments/70479.bugfix.rst new file mode 100644 index 0000000000000..b11913df8fb18 --- /dev/null +++ b/providers/microsoft/azure/newsfragments/70479.bugfix.rst @@ -0,0 +1 @@ +Fix remote log handler crashing the triggerer's logging thread when delete_local_copy is set: delete the uploaded log file and prune empty parent dirs instead of shutil.rmtree of the shared directory, which raced concurrent handler closes and raised FileNotFoundError diff --git a/providers/microsoft/azure/src/airflow/providers/microsoft/azure/log/wasb_task_handler.py b/providers/microsoft/azure/src/airflow/providers/microsoft/azure/log/wasb_task_handler.py index 32cd3ee374ec1..b6bf0806abace 100644 --- a/providers/microsoft/azure/src/airflow/providers/microsoft/azure/log/wasb_task_handler.py +++ b/providers/microsoft/azure/src/airflow/providers/microsoft/azure/log/wasb_task_handler.py @@ -17,8 +17,8 @@ # under the License. from __future__ import annotations +import contextlib import os -import shutil from functools import cached_property from pathlib import Path from typing import TYPE_CHECKING @@ -65,7 +65,16 @@ def upload(self, path: str | os.PathLike, ti: RuntimeTI | None = None) -> None: log = local_loc.read_text() has_uploaded = self.write(log, remote_loc) if has_uploaded and self.delete_local_copy: - shutil.rmtree(os.path.dirname(local_loc)) + # Delete the file and prune empty parents (idempotent); a bare + # rmtree of the dir races concurrent handler closes in the triggerer. + local_loc.unlink(missing_ok=True) + parent = local_loc.parent + while parent != self.base_log_folder and parent.is_dir(): + if any(parent.iterdir()): + break + with contextlib.suppress(OSError): + parent.rmdir() + parent = parent.parent @cached_property def hook(self):