From 5d769ac065bbeccf8879ba5ec5ae4e16ca86680f Mon Sep 17 00:00:00 2001 From: Liam Lloyd-Tucker Date: Mon, 17 Aug 2026 17:22:09 -0700 Subject: [PATCH 1/2] Make Archivematica EFS compatible Currently, the greater latency of EFS operations can cause race-condition-driven errors in Archivematica's watch_dirs.py logic, which in the worst case can cause the watch_dirs process to fail silently, leading processing to hang. This commit updates watch_dirs to avoid that. It also makes watch_dirs fail explicitly if it attempts to use inotify mode against an NFS-based filesystem, because this will not work. --- .../MCPServer/server/watch_dirs.py | 125 +++++++- tests/MCPServer/test_watch_dirs.py | 274 ++++++++++++++++++ 2 files changed, 384 insertions(+), 15 deletions(-) create mode 100644 tests/MCPServer/test_watch_dirs.py 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/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() + ) From a0ce17188344cc03119e1c8f6325523786357b13 Mon Sep 17 00:00:00 2001 From: Liam Lloyd-Tucker Date: Wed, 26 Aug 2026 15:26:20 -0700 Subject: [PATCH 2/2] Add retry to rename to handle ESTALE errors on NFS When running against EFS, Archivematica can run into an issue in the NFS client where a move operation fails because as it looks up the directory to be moved, it realizes that another client has moved it from the place where this client last saw it to another place. This triggers an effort to update where that directory lives in the client's internal state, which fails because it needs a lock that the original move operation is holding. This commit updates Archivematica to retry failed move operations once, which will deterministicly fix this particular bug, because mv's error handling reads the directory it failed to move, which successful triggers the client to update where that directory lives in its internal representation. --- .../archivematicaCommon/fileOperations.py | 36 ++++++++--- .../test_file_operations.py | 60 +++++++++++++++++++ 2 files changed, 89 insertions(+), 7 deletions(-) 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/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 + )