diff --git a/src/archivematica/MCPServer/server/watch_dirs.py b/src/archivematica/MCPServer/server/watch_dirs.py index 76901c49b4..0732e576be 100644 --- a/src/archivematica/MCPServer/server/watch_dirs.py +++ b/src/archivematica/MCPServer/server/watch_dirs.py @@ -21,9 +21,74 @@ IS_LINUX = sys.platform.startswith("linux") WATCHED_BASE_DIR = os.path.abspath(settings.WATCH_DIRECTORY) +MOUNTINFO_PATH = "/proc/self/mountinfo" + +# Filesystems where inotify cannot observe changes made by other clients. +NETWORK_FILESYSTEM_TYPES = frozenset({"nfs", "nfs4"}) + logger = logging.getLogger("archivematica.mcp.server.watchdirs") +def _unescape_mountinfo_field(field): + """Decode the octal escapes mountinfo uses in path fields.""" + for escape, character in ( + ("\\040", " "), + ("\\011", "\t"), + ("\\012", "\n"), + ("\\134", "\\"), + ): + field = field.replace(escape, character) + return field + + +def _filesystem_type(path, mountinfo_path=MOUNTINFO_PATH): + """Return the type of the filesystem that ``path`` lives on. + + Resolves to the most specific mount containing ``path``, so a path under a + bind- or sub-mount reports that mount rather than its parent. Returns None + when the type cannot be determined, e.g. on a platform without + ``/proc/self/mountinfo``. + """ + try: + with open(mountinfo_path) as mountinfo: + lines = mountinfo.readlines() + except OSError: + return None + + path = os.path.abspath(path) + best_mount_point = None + best_type = None + + for line in lines: + # Optional fields sit between the mount point and the " - " separator, + # so the two halves have to be parsed separately. + before, separator, after = line.partition(" - ") + if not separator: + continue + before_fields = before.split() + after_fields = after.split() + if len(before_fields) < 5 or not after_fields: + continue + + mount_point = _unescape_mountinfo_field(before_fields[4]) + if path != mount_point and not path.startswith(mount_point.rstrip("/") + "/"): + continue + if best_mount_point is None or len(mount_point) > len(best_mount_point): + best_mount_point = mount_point + best_type = after_fields[0] + + return best_type + + +def _list_watched_dir_entries(path): + """ + Return a ``(path, name, is_dir)`` tuple for each entry in ``path``. + Raises OSError if ``path`` cannot be read. + """ + with os.scandir(path) as entries: + return [(entry.path, entry.name, entry.is_dir()) for entry in entries] + + def watch_directories_poll( watched_dirs, shutdown_event, callback, interval=settings.WATCH_DIRECTORY_INTERVAL ): @@ -34,28 +99,43 @@ def watch_directories_poll( Accepts an iterable of workflow WatchedDir objects, a shutdown event, and a callback to be called when content appears in the watched dir. """ - # paths that have already appeared in watch directories - known_paths = set() + # Paths that have already appeared in watch directories, tracked per watched directory. + known_paths = {} while not shutdown_event.is_set(): - current_paths = set() - for watched_dir in watched_dirs: path = os.path.join(WATCHED_BASE_DIR, watched_dir.path.lstrip("/")) - for item in os.scandir(path): - if watched_dir.only_dirs and not item.is_dir(): + + try: + entries = _list_watched_dir_entries(path) + except OSError: + # On a network filesystem this is expected occasionally, e.g. a + # stale handle caused by another node renaming an entry out of + # this directory mid-scan. Keep what we knew and retry. + logger.warning("Unable to scan watched dir %s", path, exc_info=True) + continue + + seen_paths = known_paths.get(path, frozenset()) + current_paths = set() + + for item_path, _item_name, is_dir in entries: + if watched_dir.only_dirs and not is_dir: continue - elif item.path in known_paths: - # Re-add to current entries, so we keep tracking it - current_paths.add(item.path) + + # Recorded before the callback runs, so that a callback which + # raises is not retried on every subsequent pass. + current_paths.add(item_path) + if item_path in seen_paths: continue - current_paths.add(item.path) - callback(item.path, watched_dir) + try: + callback(item_path, watched_dir) + except Exception: + logger.exception("Error starting chain for %s", item_path) - # Update what we know about from the last pass, so that it doesn't grow - # endlessly - known_paths = current_paths + # Update what we know about from the last pass, so that it doesn't + # grow endlessly + known_paths[path] = current_paths time.sleep(interval) @@ -65,10 +145,12 @@ def watch_directories_inotify( ): """ Watch the directories given via inotify. This is a very efficient way to handle - watches, however it requires linux, and may not work with NFS mounts. + watches, however it requires linux and a local filesystem. Accepts an iterable of workflow WatchedDir objects, a shutdown event, and a callback to be called when content appears in the watched dir. + + Raises RuntimeError if any watched directory is on a network filesystem. """ if not IS_LINUX: warnings.warn( @@ -77,6 +159,19 @@ def watch_directories_inotify( stacklevel=2, ) + for watched_dir in watched_dirs: + path = os.path.join(WATCHED_BASE_DIR, watched_dir.path.lstrip("/")) + filesystem_type = _filesystem_type(path) + if filesystem_type in NETWORK_FILESYSTEM_TYPES: + raise RuntimeError( + f'The watched directory "{path}" is on a {filesystem_type} ' + "filesystem. inotify is only told about changes made through " + "the local kernel, so it never sees a transfer moved into a " + "watched directory by another node sharing this filesystem, " + "and those units would stall without any error. Set the " + '"watch_directory_method" setting to "poll" instead.' + ) + inotify = INotify() watch_flags = flags.CREATE | flags.MOVED_TO watches = {} # descriptor: (path, WatchedDir) diff --git a/src/archivematica/archivematicaCommon/fileOperations.py b/src/archivematica/archivematicaCommon/fileOperations.py index f4ca83ef6e..1bb9d3c2fa 100644 --- a/src/archivematica/archivematicaCommon/fileOperations.py +++ b/src/archivematica/archivematicaCommon/fileOperations.py @@ -19,6 +19,7 @@ import shutil import subprocess import sys +import time import uuid from pathlib import Path @@ -147,6 +148,15 @@ def addFileToSIP( ) +# On a shared filesystem, a directory that another host has moved is still cached +# locally under its old parent. Renaming it makes the kernel relocate that cached +# entry, which needs a lock the rename itself already holds, so the lookup fails +# with ESTALE to avoid deadlocking. The failed attempt repairs the cache as it +# unwinds, so we only need two attempts to get past this error. +RENAME_ATTEMPTS = 2 +RENAME_RETRY_DELAY = 0.1 + + def rename(source, destination, printfn=print, should_exit=False): """Used to move/rename directories. This function was before used to wrap the operation with sudo.""" if source == destination: @@ -155,13 +165,25 @@ def rename(source, destination, printfn=print, should_exit=False): return 0 command = ["mv", source, destination] - exitCode, stdOut, stdError = executeOrRun("command", command, "", printing=False) - if exitCode: - printfn("exitCode:", exitCode, file=sys.stderr) - printfn(stdOut, file=sys.stderr) - printfn(stdError, file=sys.stderr) - if should_exit: - exit(exitCode) + for attempt in range(1, RENAME_ATTEMPTS + 1): + exitCode, stdOut, stdError = executeOrRun( + "command", command, "", printing=False + ) + if not exitCode: + if attempt > 1: + printfn( + f"Renamed {source} to {destination} on attempt {attempt}.", + file=sys.stderr, + ) + return exitCode + if attempt < RENAME_ATTEMPTS: + time.sleep(RENAME_RETRY_DELAY) + + printfn("exitCode:", exitCode, file=sys.stderr) + printfn(stdOut, file=sys.stderr) + printfn(stdError, file=sys.stderr) + if should_exit: + exit(exitCode) return exitCode diff --git a/tests/MCPServer/test_watch_dirs.py b/tests/MCPServer/test_watch_dirs.py new file mode 100644 index 0000000000..e879fcc34e --- /dev/null +++ b/tests/MCPServer/test_watch_dirs.py @@ -0,0 +1,274 @@ +from unittest import mock + +import pytest + +from archivematica.MCPServer.server import watch_dirs +from archivematica.MCPServer.server.watch_dirs import _filesystem_type +from archivematica.MCPServer.server.watch_dirs import _list_watched_dir_entries +from archivematica.MCPServer.server.watch_dirs import watch_directories +from archivematica.MCPServer.server.watch_dirs import watch_directories_inotify +from archivematica.MCPServer.server.watch_dirs import watch_directories_poll + + +class StubWatchedDir: + """Stands in for a workflow WatchedDir.""" + + def __init__(self, path, only_dirs=False): + self.path = path + self.only_dirs = only_dirs + + +class StubShutdownEvent: + """Reports "not set" for a fixed number of polls, then "set". + + Lets a test run the poll loop for an exact number of cycles. + """ + + def __init__(self, cycles): + self.remaining = cycles + + def is_set(self): + if self.remaining <= 0: + return True + self.remaining -= 1 + return False + + +def run_poll(tmp_path, watched_dirs, callback, cycles=1): + """Run the poll loop over ``tmp_path`` for ``cycles`` iterations.""" + with ( + mock.patch.object(watch_dirs, "WATCHED_BASE_DIR", str(tmp_path)), + mock.patch.object(watch_dirs.time, "sleep"), + ): + watch_directories_poll( + watched_dirs, StubShutdownEvent(cycles), callback, interval=0 + ) + + +# _list_watched_dir_entries + + +def test_list_watched_dir_entries_returns_paths_names_and_dir_flags(tmp_path): + (tmp_path / "a_file.txt").write_text("hello") + (tmp_path / "a_dir").mkdir() + + entries = sorted(_list_watched_dir_entries(str(tmp_path))) + + assert entries == sorted( + [ + (str(tmp_path / "a_file.txt"), "a_file.txt", False), + (str(tmp_path / "a_dir"), "a_dir", True), + ] + ) + + +def test_list_watched_dir_entries_returns_empty_list_for_empty_dir(tmp_path): + assert _list_watched_dir_entries(str(tmp_path)) == [] + + +def test_list_watched_dir_entries_raises_oserror_for_unreadable_dir(tmp_path): + with pytest.raises(OSError): + _list_watched_dir_entries(str(tmp_path / "does_not_exist")) + + +# watch_directories_poll + + +def test_poll_calls_callback_for_new_entry(tmp_path): + (tmp_path / "transfer").mkdir() + watched_dir = StubWatchedDir("/", only_dirs=True) + callback = mock.Mock() + + run_poll(tmp_path, [watched_dir], callback, cycles=1) + + callback.assert_called_once_with(str(tmp_path / "transfer"), watched_dir) + + +def test_poll_does_not_call_callback_twice_for_same_entry(tmp_path): + (tmp_path / "transfer").mkdir() + callback = mock.Mock() + + run_poll(tmp_path, [StubWatchedDir("/", only_dirs=True)], callback, cycles=3) + + callback.assert_called_once() + + +def test_poll_skips_files_when_only_dirs_is_set(tmp_path): + (tmp_path / "a_file.txt").write_text("hello") + callback = mock.Mock() + + run_poll(tmp_path, [StubWatchedDir("/", only_dirs=True)], callback, cycles=1) + + callback.assert_not_called() + + +def test_poll_keeps_running_when_a_directory_cannot_be_scanned(tmp_path): + """A transient scandir failure must not end the watch loop. + + On EFS an ESTALE from a concurrent rename would otherwise propagate out of + the watcher thread, silently stopping every watched directory. + """ + (tmp_path / "transfer").mkdir() + watched_dir = StubWatchedDir("/", only_dirs=True) + callback = mock.Mock() + listings = [ + OSError(116, "Stale file handle"), + [(str(tmp_path / "transfer"), "transfer", True)], + ] + + with mock.patch.object( + watch_dirs, "_list_watched_dir_entries", side_effect=listings + ): + run_poll(tmp_path, [watched_dir], callback, cycles=2) + + callback.assert_called_once_with(str(tmp_path / "transfer"), watched_dir) + + +def test_poll_does_not_refire_known_entries_after_a_scan_failure(tmp_path): + """A failed scan must not make already-seen entries look new again. + + Dropping them from the known set would start a second chain for a package + that is already being processed. + """ + entry = (str(tmp_path / "transfer"), "transfer", True) + callback = mock.Mock() + listings = [[entry], OSError(116, "Stale file handle"), [entry]] + + with mock.patch.object( + watch_dirs, "_list_watched_dir_entries", side_effect=listings + ): + run_poll(tmp_path, [StubWatchedDir("/", only_dirs=True)], callback, cycles=3) + + callback.assert_called_once() + + +def test_poll_keeps_running_when_the_callback_raises(tmp_path): + (tmp_path / "one").mkdir() + (tmp_path / "two").mkdir() + callback = mock.Mock(side_effect=[ValueError("boom"), None]) + + run_poll(tmp_path, [StubWatchedDir("/", only_dirs=True)], callback, cycles=1) + + assert callback.call_count == 2 + + +def test_poll_scans_remaining_directories_when_one_fails(tmp_path): + good = tmp_path / "good" + good.mkdir() + (good / "transfer").mkdir() + callback = mock.Mock() + watched_dirs = [StubWatchedDir("/missing"), StubWatchedDir("/good", only_dirs=True)] + + run_poll(tmp_path, watched_dirs, callback, cycles=1) + + callback.assert_called_once_with(str(good / "transfer"), watched_dirs[1]) + + +# _filesystem_type + + +def write_mountinfo(tmp_path, lines): + mountinfo = tmp_path / "mountinfo" + mountinfo.write_text("".join(f"{line}\n" for line in lines)) + return str(mountinfo) + + +def test_filesystem_type_returns_type_of_longest_matching_mount(tmp_path): + mountinfo = write_mountinfo( + tmp_path, + [ + "23 1 0:1 / / rw,relatime - ext4 /dev/root rw", + "44 23 0:5 / /var/archivematica/sharedDirectory rw - nfs4 fs.efs:/ rw", + ], + ) + + assert ( + _filesystem_type( + "/var/archivematica/sharedDirectory/watchedDirectories", + mountinfo_path=mountinfo, + ) + == "nfs4" + ) + + +def test_filesystem_type_returns_type_of_enclosing_mount_for_unmounted_path(tmp_path): + mountinfo = write_mountinfo( + tmp_path, ["23 1 0:1 / / rw,relatime - ext4 /dev/root rw"] + ) + + assert _filesystem_type("/var/archivematica", mountinfo_path=mountinfo) == "ext4" + + +def test_filesystem_type_ignores_mounts_that_only_share_a_name_prefix(tmp_path): + mountinfo = write_mountinfo( + tmp_path, + [ + "23 1 0:1 / / rw,relatime - ext4 /dev/root rw", + "44 23 0:5 / /var/archivematica-other rw - nfs4 fs.efs:/ rw", + ], + ) + + assert _filesystem_type("/var/archivematica", mountinfo_path=mountinfo) == "ext4" + + +def test_filesystem_type_handles_optional_fields_before_the_separator(tmp_path): + mountinfo = write_mountinfo( + tmp_path, + ["44 23 0:5 / /shared rw shared:2 master:1 - nfs4 fs.efs:/ rw"], + ) + + assert _filesystem_type("/shared", mountinfo_path=mountinfo) == "nfs4" + + +def test_filesystem_type_returns_none_when_mountinfo_is_unavailable(tmp_path): + assert _filesystem_type("/shared", mountinfo_path=str(tmp_path / "nope")) is None + + +# inotify rejection on network filesystems + + +def test_inotify_raises_when_watched_dir_is_on_a_network_filesystem(tmp_path): + watched_dirs = [StubWatchedDir("/activeTransfers")] + + with ( + mock.patch.object(watch_dirs, "WATCHED_BASE_DIR", str(tmp_path)), + mock.patch.object(watch_dirs, "_filesystem_type", return_value="nfs4"), + pytest.raises(RuntimeError, match="nfs4"), + ): + watch_directories_inotify(watched_dirs, StubShutdownEvent(0), mock.Mock()) + + +def test_inotify_error_names_the_poll_alternative(tmp_path): + with ( + mock.patch.object(watch_dirs, "WATCHED_BASE_DIR", str(tmp_path)), + mock.patch.object(watch_dirs, "_filesystem_type", return_value="nfs"), + pytest.raises(RuntimeError, match="poll"), + ): + watch_directories_inotify( + [StubWatchedDir("/activeTransfers")], StubShutdownEvent(0), mock.Mock() + ) + + +def test_watch_directories_rejects_inotify_on_a_network_filesystem(tmp_path): + with ( + mock.patch.object(watch_dirs, "WATCHED_BASE_DIR", str(tmp_path)), + mock.patch.object(watch_dirs, "_filesystem_type", return_value="nfs4"), + mock.patch.object(watch_dirs.settings, "WATCH_DIRECTORY_METHOD", "inotify"), + pytest.raises(RuntimeError), + ): + watch_directories( + [StubWatchedDir("/activeTransfers")], StubShutdownEvent(0), mock.Mock() + ) + + +def test_inotify_is_allowed_on_a_local_filesystem(tmp_path): + """inotify stays usable for local development on a real local filesystem.""" + (tmp_path / "activeTransfers").mkdir() + + with ( + mock.patch.object(watch_dirs, "WATCHED_BASE_DIR", str(tmp_path)), + mock.patch.object(watch_dirs, "_filesystem_type", return_value="ext4"), + ): + watch_directories_inotify( + [StubWatchedDir("/activeTransfers")], StubShutdownEvent(0), mock.Mock() + ) diff --git a/tests/archivematicaCommon/test_file_operations.py b/tests/archivematicaCommon/test_file_operations.py index 820cd22385..f7bb4172b5 100644 --- a/tests/archivematicaCommon/test_file_operations.py +++ b/tests/archivematicaCommon/test_file_operations.py @@ -1,15 +1,18 @@ import pathlib +import sys from unittest import mock import pytest from django.db.models import Q +from archivematica.archivematicaCommon.fileOperations import RENAME_ATTEMPTS from archivematica.archivematicaCommon.fileOperations import ( FindFileInNormalizatonCSVError, ) from archivematica.archivematicaCommon.fileOperations import addAccessionEvent from archivematica.archivematicaCommon.fileOperations import findFileInNormalizationCSV from archivematica.archivematicaCommon.fileOperations import get_extract_dir_name +from archivematica.archivematicaCommon.fileOperations import rename from archivematica.dashboard.main.models import SIP from archivematica.dashboard.main.models import Event from archivematica.dashboard.main.models import File @@ -269,3 +272,60 @@ def test_findFileInNormalizationCSV_fails_if_multiple_target_files_exist( f"More than one result found for {purpose} file ({target_file}) in DB.", file=mock.ANY, ) + + +def test_rename_is_a_noop_when_source_and_destination_match(): + printfn = mock.Mock() + + assert rename("/same", "/same", printfn=printfn) == 0 + + printfn.assert_called_once_with( + "Source and destination are the same, nothing to do." + ) + + +@mock.patch("archivematica.archivematicaCommon.fileOperations.executeOrRun") +def test_rename_does_not_retry_a_successful_move(execute_or_run): + execute_or_run.return_value = (0, "", "") + + assert rename("/src", "/dst") == 0 + + assert execute_or_run.call_count == 1 + + +@mock.patch("archivematica.archivematicaCommon.fileOperations.time.sleep") +@mock.patch("archivematica.archivematicaCommon.fileOperations.executeOrRun") +def test_rename_retries_after_a_failed_attempt(execute_or_run, sleep): + """A transient ESTALE clears once the first attempt has unwound.""" + execute_or_run.side_effect = [ + (1, "", "mv: cannot move '/src' to '/dst': Stale file handle"), + (0, "", ""), + ] + printfn = mock.Mock() + + assert rename("/src", "/dst", printfn=printfn) == 0 + + assert execute_or_run.call_count == 2 + printfn.assert_called_once_with( + "Renamed /src to /dst on attempt 2.", file=sys.stderr + ) + + +@mock.patch("archivematica.archivematicaCommon.fileOperations.time.sleep") +@mock.patch("archivematica.archivematicaCommon.fileOperations.executeOrRun") +def test_rename_gives_up_and_reports_the_exit_code(execute_or_run, sleep): + """A move that keeps failing still surfaces mv's exit code and output.""" + execute_or_run.return_value = ( + 1, + "", + "mv: cannot stat '/src': No such file or directory", + ) + printfn = mock.Mock() + + assert rename("/src", "/dst", printfn=printfn) == 1 + + assert execute_or_run.call_count == RENAME_ATTEMPTS + printfn.assert_any_call("exitCode:", 1, file=sys.stderr) + printfn.assert_any_call( + "mv: cannot stat '/src': No such file or directory", file=sys.stderr + )