From 934ce239d8f08b0f333190f249605e8b48e8cc79 Mon Sep 17 00:00:00 2001 From: Vinay Sharma Date: Sat, 24 Aug 2019 20:44:34 +0530 Subject: [PATCH] bpo-37754: make shared_memory's unix implementation consistent with Windows --- Lib/multiprocessing/resource_tracker.py | 21 +++++ Lib/multiprocessing/shared_memory.py | 15 ++++ Lib/test/_test_multiprocessing.py | 76 +++++++++++++++++++ .../2019-08-24-20-43-40.bpo-37754.KAvVkE.rst | 1 + 4 files changed, 113 insertions(+) create mode 100644 Misc/NEWS.d/next/Library/2019-08-24-20-43-40.bpo-37754.KAvVkE.rst diff --git a/Lib/multiprocessing/resource_tracker.py b/Lib/multiprocessing/resource_tracker.py index 61a6dd66e72e67..8d704a3ede0173 100644 --- a/Lib/multiprocessing/resource_tracker.py +++ b/Lib/multiprocessing/resource_tracker.py @@ -19,6 +19,7 @@ import signal import sys import threading +import errno import warnings from . import spawn @@ -33,15 +34,21 @@ 'noop': lambda: None, } +_FILE_PREFIXES = {} + if os.name == 'posix': import _multiprocessing import _posixshmem + import fcntl _CLEANUP_FUNCS.update({ 'semaphore': _multiprocessing.sem_unlink, 'shared_memory': _posixshmem.shm_unlink, }) + _FILE_PREFIXES.update({ + 'shared_memory': '/dev/shm' # Directory containing memory mapped files created + }) # by SharedMemory in Unix class ResourceTracker(object): @@ -163,6 +170,7 @@ def main(fd): if _HAVE_SIGMASK: signal.pthread_sigmask(signal.SIG_UNBLOCK, _IGNORED_SIGNALS) + for f in (sys.stdin, sys.stdout): try: f.close() @@ -209,6 +217,19 @@ def main(fd): # For some reason the process which created and registered this # resource has failed to unregister it. Presumably it has # died. We therefore unlink it. + + if rtype in _FILE_PREFIXES: + try: + sh_fd = open(_FILE_PREFIXES[rtype] + name) + fcntl.flock(sh_fd, fcntl.LOCK_EX | fcntl.LOCK_NB) # Try to acquire exclusive lock on shared_memory + sh_fd.close() + except FileNotFoundError: + pass + except IOError as e: + sh_fd.close() + if e.errno == errno.EAGAIN: # Don't Cleanup if a shared flock is present + continue # implying, that a process is using it. + try: try: _CLEANUP_FUNCS[rtype](name) diff --git a/Lib/multiprocessing/shared_memory.py b/Lib/multiprocessing/shared_memory.py index 184e36704baaeb..05a3bef74d0fcc 100644 --- a/Lib/multiprocessing/shared_memory.py +++ b/Lib/multiprocessing/shared_memory.py @@ -20,6 +20,7 @@ _USE_POSIX = False else: import _posixshmem + import fcntl _USE_POSIX = True @@ -69,6 +70,7 @@ class SharedMemory: _flags = os.O_RDWR _mode = 0o600 _prepend_leading_slash = True if _USE_POSIX else False + _has_shared_lock = False def __init__(self, name=None, create=False, size=0): if not size >= 0: @@ -113,6 +115,7 @@ def __init__(self, name=None, create=False, size=0): self.unlink() raise + self._aquire_shared_lock() from .resource_tracker import register register(self._name, "shared_memory") @@ -195,6 +198,17 @@ def __reduce__(self): def __repr__(self): return f'{self.__class__.__name__}({self.name!r}, size={self.size})' + def _aquire_shared_lock(self): + if _USE_POSIX and (not self._has_shared_lock): + fcntl.flock(self._fd, fcntl.LOCK_SH | fcntl.LOCK_NB) + self._has_shared_lock = True + + + def _release_shared_lock(self): + if _USE_POSIX and self._has_shared_lock: + fcntl.flock(self._fd, fcntl.LOCK_UN | fcntl.LOCK_NB) + self._has_shared_lock = False + @property def buf(self): "A memoryview of contents of the shared memory block." @@ -217,6 +231,7 @@ def size(self): def close(self): """Closes access to the shared memory from this instance but does not destroy the shared memory block.""" + self._release_shared_lock() if self._buf is not None: self._buf.release() self._buf = None diff --git a/Lib/test/_test_multiprocessing.py b/Lib/test/_test_multiprocessing.py index 2fe0def2bcd277..9376b62626217c 100644 --- a/Lib/test/_test_multiprocessing.py +++ b/Lib/test/_test_multiprocessing.py @@ -4026,6 +4026,82 @@ def test_shared_memory_cleaned_after_process_termination(self): "resource_tracker: There appear to be 1 leaked " "shared_memory objects to clean up at shutdown", err) + + def test_shared_memory_persistence_after_one_of_multiple_processes_terminate(self): + # Test If shared memory can be attached after a process using it exits, + # but another process is still holding it. + cmd_process_1 = '''if 1: + import time, sys + from multiprocessing import shared_memory + + # Create a shared_memory segment, and send the segment name + sm = shared_memory.SharedMemory(create=True, size=10) + sys.stdout.write(sm.name + '\\n') + sys.stdout.flush() + time.sleep(100) + ''' + cmd_process_2 = '''if 1: + import time, sys + from multiprocessing import shared_memory + + # Create a shared_memory segment, and send the segment name + sm = shared_memory.SharedMemory(name={}, create=False) + sys.stdout.write(sm.name + '\\n') + sys.stdout.flush() + time.sleep(100) + ''' + + p1 = subprocess.Popen([sys.executable, '-E', '-c', cmd_process_1], + stdout=subprocess.PIPE, + stderr=subprocess.PIPE) + + name = p1.stdout.readline().strip().decode() + + p2 = subprocess.Popen([sys.executable, '-E', '-c', cmd_process_2.format(name)]) + + p2.terminate() + p2.wait() + + deadline = time.monotonic() + 60 + t = 0.1 + while time.monotonic() < deadline: + time.sleep(t) + t = min(t*2, 5) + try: + smm = shared_memory.SharedMemory(name=name, create=False) + except FileNotFoundError: + raise AssertionError("Shared Memory segment was unlinked, despite" + "the fact a process is still using it.") + + smm.close() + p1.terminate() + p1.wait() + + deadline = time.monotonic() + 60 + t = 0.1 + while time.monotonic() < deadline: + time.sleep(t) + t = min(t*2, 5) + try: + smm = shared_memory.SharedMemory(name=name, create=False) + except FileNotFoundError: + break + else: + raise AssertionError("A SharedMemory segment was leaked after" + " a process was abruptly terminated.") + + if os.name == 'posix': + # A warning was emitted by the subprocess' own + # resource_tracker (on Windows, shared memory segments + # are released automatically by the OS). + err = p1.stderr.read().decode() + self.assertIn( + "resource_tracker: There appear to be 1 leaked " + "shared_memory objects to clean up at shutdown", err) + + p1.stderr.close() + p1.stdout.close() + # # # diff --git a/Misc/NEWS.d/next/Library/2019-08-24-20-43-40.bpo-37754.KAvVkE.rst b/Misc/NEWS.d/next/Library/2019-08-24-20-43-40.bpo-37754.KAvVkE.rst new file mode 100644 index 00000000000000..c68f955396f9ca --- /dev/null +++ b/Misc/NEWS.d/next/Library/2019-08-24-20-43-40.bpo-37754.KAvVkE.rst @@ -0,0 +1 @@ +make shared_memory's unix implementation consistent with Windows