Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions providers/alibaba/newsfragments/70479.bugfix.rst
Original file line number Diff line number Diff line change
@@ -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
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,6 @@

import contextlib
import os
import shutil
from functools import cached_property
from pathlib import Path
from typing import TYPE_CHECKING
Expand Down Expand Up @@ -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):
Expand Down
1 change: 1 addition & 0 deletions providers/amazon/newsfragments/70479.bugfix.rst
Original file line number Diff line number Diff line change
@@ -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
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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("")

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
1 change: 1 addition & 0 deletions providers/apache/hdfs/newsfragments/70479.bugfix.rst
Original file line number Diff line number Diff line change
@@ -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
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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):
Expand Down
1 change: 1 addition & 0 deletions providers/elasticsearch/newsfragments/70479.bugfix.rst
Original file line number Diff line number Diff line change
@@ -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
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,6 @@
import json
import logging
import os
import shutil
import sys
import time
import warnings
Expand Down Expand Up @@ -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")
Expand Down
1 change: 1 addition & 0 deletions providers/google/newsfragments/70479.bugfix.rst
Original file line number Diff line number Diff line change
@@ -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
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand Down
1 change: 1 addition & 0 deletions providers/microsoft/azure/newsfragments/70479.bugfix.rst
Original file line number Diff line number Diff line change
@@ -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
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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):
Expand Down
Loading