Merge pull request #7446 from mvdbeek/fix_concurrent_add_remove

Use a lock on cleanup, expire_tool and cache_tool
This commit is contained in:
Nicola Soranzo
2019-03-04 18:21:53 +00:00
committed by GitHub
3 changed files with 72 additions and 55 deletions
+54 -40
View File
@@ -1,10 +1,16 @@
import logging
import os
from threading import local
from threading import (
local,
Lock,
)
from sqlalchemy.orm.exc import DetachedInstanceError
from galaxy.util.hash_util import md5_hash_file
log = logging.getLogger(__name__)
class ToolCache(object):
"""
@@ -13,6 +19,7 @@ class ToolCache(object):
"""
def __init__(self):
self._lock = Lock()
self._hash_by_tool_paths = {}
self._tools_by_path = {}
self._tool_paths_by_id = {}
@@ -30,26 +37,30 @@ class ToolCache(object):
"""
removed_tool_ids = []
try:
paths_to_cleanup = {path: tool.all_ids for path, tool in self._tools_by_path.items() if self._should_cleanup(path)}
for config_filename, tool_ids in paths_to_cleanup.items():
del self._hash_by_tool_paths[config_filename]
if os.path.exists(config_filename):
# This tool has probably been broken while editing on disk
# We record it here, so that we can recover it
self._removed_tools_by_path[config_filename] = self._tools_by_path[config_filename]
del self._tools_by_path[config_filename]
for tool_id in tool_ids:
if tool_id in self._tool_paths_by_id:
del self._tool_paths_by_id[tool_id]
removed_tool_ids.extend(tool_ids)
for tool_id in removed_tool_ids:
self._removed_tool_ids.add(tool_id)
if tool_id in self._new_tool_ids:
self._new_tool_ids.remove(tool_id)
except Exception:
with self._lock:
paths_to_cleanup = {path: tool.all_ids for path, tool in self._tools_by_path.items() if self._should_cleanup(path)}
for config_filename, tool_ids in paths_to_cleanup.items():
del self._hash_by_tool_paths[config_filename]
if os.path.exists(config_filename):
# This tool has probably been broken while editing on disk
# We record it here, so that we can recover it
self._removed_tools_by_path[config_filename] = self._tools_by_path[config_filename]
del self._tools_by_path[config_filename]
for tool_id in tool_ids:
if tool_id in self._tool_paths_by_id:
del self._tool_paths_by_id[tool_id]
removed_tool_ids.extend(tool_ids)
for tool_id in removed_tool_ids:
self._removed_tool_ids.add(tool_id)
if tool_id in self._new_tool_ids:
self._new_tool_ids.remove(tool_id)
except Exception as e:
log.debug("Exception while checking tools to remove from cache: %s" % e)
# If by chance the file is being removed while calculating the hash or modtime
# we don't want the thread to die.
pass
if removed_tool_ids:
log.debug("Removed the following tools from cache: %s" % removed_tool_ids)
return removed_tool_ids
def _should_cleanup(self, config_filename):
@@ -79,39 +90,42 @@ class ToolCache(object):
return self.get_tool(self._tool_paths_by_id.get(tool_id))
def expire_tool(self, tool_id):
if tool_id in self._tool_paths_by_id:
config_filename = self._tool_paths_by_id[tool_id]
del self._hash_by_tool_paths[config_filename]
del self._tool_paths_by_id[tool_id]
del self._tools_by_path[config_filename]
del self._mod_time_by_path[config_filename]
if tool_id in self._new_tool_ids:
self._new_tool_ids.remove(tool_id)
with self._lock:
if tool_id in self._tool_paths_by_id:
config_filename = self._tool_paths_by_id[tool_id]
del self._hash_by_tool_paths[config_filename]
del self._tool_paths_by_id[tool_id]
del self._tools_by_path[config_filename]
del self._mod_time_by_path[config_filename]
if tool_id in self._new_tool_ids:
self._new_tool_ids.remove(tool_id)
def cache_tool(self, config_filename, tool):
tool_hash = md5_hash_file(config_filename)
if tool_hash is None:
return
tool_id = str(tool.id)
self._hash_by_tool_paths[config_filename] = tool_hash
self._mod_time_by_path[config_filename] = os.path.getmtime(config_filename)
self._tool_paths_by_id[tool_id] = config_filename
self._tools_by_path[config_filename] = tool
self._new_tool_ids.add(tool_id)
for macro_path in tool._macro_paths:
self._mod_time_by_path[macro_path] = os.path.getmtime(macro_path)
if tool_id not in self._macro_paths_by_id:
self._macro_paths_by_id[tool_id] = {macro_path}
else:
self._macro_paths_by_id[tool_id].add(macro_path)
with self._lock:
self._hash_by_tool_paths[config_filename] = tool_hash
self._mod_time_by_path[config_filename] = os.path.getmtime(config_filename)
self._tool_paths_by_id[tool_id] = config_filename
self._tools_by_path[config_filename] = tool
self._new_tool_ids.add(tool_id)
for macro_path in tool._macro_paths:
self._mod_time_by_path[macro_path] = os.path.getmtime(macro_path)
if tool_id not in self._macro_paths_by_id:
self._macro_paths_by_id[tool_id] = {macro_path}
else:
self._macro_paths_by_id[tool_id].add(macro_path)
def reset_status(self):
"""
Reset tracking of new and newly disabled tools.
"""
self._new_tool_ids = set()
self._removed_tool_ids = set()
self._removed_tools_by_path = {}
with self._lock:
self._new_tool_ids = set()
self._removed_tool_ids = set()
self._removed_tools_by_path = {}
class ToolShedRepositoryCache(object):
+16 -14
View File
@@ -51,11 +51,9 @@ class BaseToolBoxTestCase(unittest.TestCase, UsesApp, UsesTools):
@property
def toolbox(self):
if self.__toolbox is None:
self.__toolbox = SimplifiedToolBox(self)
# wire app with this new toolbox
self.app.toolbox = self.__toolbox
return self.__toolbox
if self._toolbox is None:
self.app.toolbox = self._toolbox = SimplifiedToolBox(self)
return self._toolbox
def setUp(self):
self.reindexed = False
@@ -67,9 +65,12 @@ class BaseToolBoxTestCase(unittest.TestCase, UsesApp, UsesTools):
itp_config = os.path.join(self.test_directory, "integrated_tool_panel.xml")
self.app.config.integrated_tool_panel_config = itp_config
self.app.watchers = ConfigWatchers(self.app)
self.__toolbox = None
self._toolbox = None
self.config_files = []
def tearDown(self):
self.app.watchers.shutdown()
def _repo_install(self, changeset, config_filename=None):
metadata = {
'tools': [{
@@ -193,13 +194,14 @@ class ToolBoxTestCase(BaseToolBoxTestCase):
assert tool is not None
assert len(tool._macro_paths) == 1
macro_path = tool._macro_paths[0]
time.sleep(1.5)
with open(macro_path, 'w') as macro_out:
macro_out.write(SIMPLE_MACRO.substitute(tool_version="3.0"))
time.sleep(1.5)
tool = self.app.toolbox.get_tool("tool_with_macro")
assert tool.version == "3.0"
def check_tool_macro():
tool = self.toolbox.get_tool("tool_with_macro")
assert tool.version == "3.0"
self._try_until_no_errors(check_tool_macro)
def test_tool_reload_for_broken_tool(self):
self._init_tool(filename="simple_tool.xml", version="1.0")
@@ -215,7 +217,7 @@ class ToolBoxTestCase(BaseToolBoxTestCase):
out.write('certainly not a valid tool')
def check_tool_errors():
tool = self.app.toolbox.get_tool("test_tool")
tool = self.toolbox.get_tool("test_tool")
assert tool is not None
assert tool.version == "1.0"
assert tool.tool_errors == 'Current on-disk tool is not valid'
@@ -226,7 +228,7 @@ class ToolBoxTestCase(BaseToolBoxTestCase):
self._init_tool(filename="simple_tool.xml", version="2.0")
def check_no_tool_errors():
tool = self.app.toolbox.get_tool("test_tool")
tool = self.toolbox.get_tool("test_tool")
assert tool is not None
assert tool.version == "2.0"
assert tool.tool_errors is None
@@ -236,7 +238,7 @@ class ToolBoxTestCase(BaseToolBoxTestCase):
def _try_until_no_errors(self, f):
e = None
for i in range(300):
for i in range(10):
try:
f()
return
@@ -561,4 +563,4 @@ class SimplifiedToolBox(ToolBox):
def reload_callback(test_case):
test_case.app.tool_cache.cleanup()
test_case.__toolbox = test_case.app.toolbox = SimplifiedToolBox(test_case)
test_case._toolbox = test_case.app.toolbox = SimplifiedToolBox(test_case)
+2 -1
View File
@@ -102,7 +102,8 @@ class UsesTools(object):
def __write_tool(self, contents, path=None):
path = path or self.tool_file
open(path, "w").write(contents)
with open(path, "w") as out:
out.write(contents)
class MockContext(object):