From da61e77068f115ebe0afcd10a6d34b9dd43aaae7 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Wed, 13 Apr 2022 17:00:05 +0200 Subject: [PATCH 1/6] Fix persistent mulled cache unused for repo_has_name --- doc/source/admin/galaxy_options.rst | 14 +++- lib/galaxy/app.py | 7 +- lib/galaxy/config/sample/galaxy.yml.sample | 4 ++ lib/galaxy/config/schemas/config_schema.yml | 7 ++ .../deps/container_resolvers/__init__.py | 2 + .../deps/container_resolvers/mulled.py | 11 +++- lib/galaxy/tool_util/deps/mulled/util.py | 66 ++++++++++++------- packages/tool_util/test-requirements.txt | 1 + test/unit/tool_util/test_resolution_cache.py | 61 +++++++++++++++++ 9 files changed, 139 insertions(+), 34 deletions(-) create mode 100644 test/unit/tool_util/test_resolution_cache.py diff --git a/doc/source/admin/galaxy_options.rst b/doc/source/admin/galaxy_options.rst index 7235cb45ab9..676b3e08fa3 100644 --- a/doc/source/admin/galaxy_options.rst +++ b/doc/source/admin/galaxy_options.rst @@ -1292,6 +1292,17 @@ :Type: str +~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ +``mulled_resolution_cache_expire`` +~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ + +:Description: + Seconds until the beaker cache is considered old and a new value + is created. +:Default: ``3600`` +:Type: int + + ~~~~~~~~~~~~~~~~~~~~~~~~~~~~ ``object_store_config_file`` ~~~~~~~~~~~~~~~~~~~~~~~~~~~~ @@ -4835,6 +4846,3 @@ . :Default: ``vault_conf.yml`` :Type: str - - - diff --git a/lib/galaxy/app.py b/lib/galaxy/app.py index f80c7b131cf..31731dbfef9 100644 --- a/lib/galaxy/app.py +++ b/lib/galaxy/app.py @@ -218,9 +218,10 @@ class ConfiguresGalaxyMixin: mulled_resolution_cache = None if self.config.mulled_resolution_cache_type: cache_opts = { - 'cache.type': self.config.mulled_resolution_cache_type, - 'cache.data_dir': self.config.mulled_resolution_cache_data_dir, - 'cache.lock_dir': self.config.mulled_resolution_cache_lock_dir, + "cache.type": self.config.mulled_resolution_cache_type, + "cache.data_dir": self.config.mulled_resolution_cache_data_dir, + "cache.lock_dir": self.config.mulled_resolution_cache_lock_dir, + "cache.expire": self.config.mulled_resolution_cache_expire, } mulled_resolution_cache = CacheManager(**parse_cache_config_options(cache_opts)).get_cache('mulled_resolution') self.container_finder = containers.ContainerFinder(app_info, mulled_resolution_cache=mulled_resolution_cache) diff --git a/lib/galaxy/config/sample/galaxy.yml.sample b/lib/galaxy/config/sample/galaxy.yml.sample index b64f2527938..7bff2d9b9b3 100644 --- a/lib/galaxy/config/sample/galaxy.yml.sample +++ b/lib/galaxy/config/sample/galaxy.yml.sample @@ -847,6 +847,10 @@ galaxy: # . #mulled_resolution_cache_lock_dir: mulled/locks + # Seconds until the beaker cache is considered old and a new value is + # created. + #mulled_resolution_cache_expire: 3600 + # Configuration file for the object store If this is set and exists, # it overrides any other objectstore settings. # The value of this option will be resolved with respect to diff --git a/lib/galaxy/config/schemas/config_schema.yml b/lib/galaxy/config/schemas/config_schema.yml index 28a61527bcd..cc4b3abad73 100644 --- a/lib/galaxy/config/schemas/config_schema.yml +++ b/lib/galaxy/config/schemas/config_schema.yml @@ -930,6 +930,13 @@ mapping: desc: | Lock directory used by beaker for caching mulled resolution requests. + mulled_resolution_cache_expire: + type: int + default: 3600 + required: false + desc: | + Seconds until the beaker cache is considered old and a new value is created. + object_store_config_file: type: str default: object_store_conf.xml diff --git a/lib/galaxy/tool_util/deps/container_resolvers/__init__.py b/lib/galaxy/tool_util/deps/container_resolvers/__init__.py index 654f6b9d994..5965b79f3ea 100644 --- a/lib/galaxy/tool_util/deps/container_resolvers/__init__.py +++ b/lib/galaxy/tool_util/deps/container_resolvers/__init__.py @@ -17,6 +17,8 @@ class ResolutionCache(Bunch): one resolution at a time in a single thread. """ + mulled_resolution_cache = None + class ContainerResolver(Dictifiable, metaclass=ABCMeta): """Description of a technique for resolving container images for tool execution.""" diff --git a/lib/galaxy/tool_util/deps/container_resolvers/mulled.py b/lib/galaxy/tool_util/deps/container_resolvers/mulled.py index bbae6793514..f92ad6aa053 100644 --- a/lib/galaxy/tool_util/deps/container_resolvers/mulled.py +++ b/lib/galaxy/tool_util/deps/container_resolvers/mulled.py @@ -19,6 +19,7 @@ from galaxy.util.commands import shell from ..container_classes import CONTAINER_CLASSES from ..container_resolvers import ( ContainerResolver, + ResolutionCache, ) from ..docker_util import build_docker_images_command from ..mulled.mulled_build import ( @@ -314,7 +315,9 @@ def singularity_cached_container_description(targets, cache_directory, hash_func return container -def targets_to_mulled_name(targets, hash_func, namespace, resolution_cache=None, session=None): +def targets_to_mulled_name( + targets, hash_func, namespace, resolution_cache: Optional[ResolutionCache] = None, session=None +): unresolved_cache_key = "galaxy.tool_util.deps.container_resolvers.mulled:unresolved" if resolution_cache is not None: if unresolved_cache_key not in resolution_cache: @@ -324,15 +327,17 @@ def targets_to_mulled_name(targets, hash_func, namespace, resolution_cache=None, unresolved_cache = set() mulled_resolution_cache = None - if resolution_cache and hasattr(resolution_cache, 'mulled_resolution_cache'): + if resolution_cache and resolution_cache.mulled_resolution_cache: mulled_resolution_cache = resolution_cache.mulled_resolution_cache name = None def cached_name(cache_key): if mulled_resolution_cache: - if cache_key in mulled_resolution_cache: + try: return resolution_cache.get(cache_key) + except KeyError: + return None return None if len(targets) == 1: diff --git a/lib/galaxy/tool_util/deps/mulled/util.py b/lib/galaxy/tool_util/deps/mulled/util.py index 307da3c49e7..7379cc5fc2f 100644 --- a/lib/galaxy/tool_util/deps/mulled/util.py +++ b/lib/galaxy/tool_util/deps/mulled/util.py @@ -18,6 +18,9 @@ QUAY_REPOSITORY_API_ENDPOINT = 'https://quay.io/api/v1/repository' BUILD_NUMBER_REGEX = re.compile(r'\d+$') PARSED_TAG = collections.namedtuple('PARSED_TAG', 'tag version build_string build_number') MULLED_SOCKET_TIMEOUT = 12 +QUAY_VERSIONS_CACHE_EXPIRY = 300 +NAMESPACE_HAS_REPO_NAME_KEY = "galaxy.tool_util.deps.container_resolvers.mulled.util:namespace_repo_names" +TAG_CACHE_KEY = "galaxy.tool_util.deps.container_resolvers.mulled.util:tag_cache" def create_repository(namespace, repo_name, oauth_token): @@ -60,29 +63,36 @@ def _namespace_has_repo_name(namespace, repo_name, resolution_cache): """ Get all quay containers in the biocontainers repo """ - cache_key = "galaxy.tool_util.deps.container_resolvers.mulled.util:namespace_repo_names" - if resolution_cache is not None and cache_key in resolution_cache: - repo_names = resolution_cache.get(cache_key) - else: - next_page = None - repo_names = [] - repos_headers = {"Accept-encoding": "gzip", "Accept": "application/json"} - while True: - repos_parameters = {"public": "true", "namespace": namespace, "next_page": next_page} - repos_response = requests.get( - QUAY_REPOSITORY_API_ENDPOINT, headers=repos_headers, params=repos_parameters, timeout=MULLED_SOCKET_TIMEOUT) - repos_response_json = repos_response.json() - repos = repos_response_json["repositories"] - repo_names += [r["name"] for r in repos] - next_page = repos_response_json.get("next_page") - if not next_page: - break - if resolution_cache is not None: - resolution_cache[cache_key] = repo_names + # resolution_cache.mulled_resolution_cache is the persistent variant of the resolution cache + resolution_cache = resolution_cache.mulled_resolution_cache or resolution_cache + cache_key = NAMESPACE_HAS_REPO_NAME_KEY + if resolution_cache is not None: + try: + return repo_name in resolution_cache.get(cache_key) + except KeyError: + pass + next_page = None + repo_names = [] + repos_headers = {"Accept-encoding": "gzip", "Accept": "application/json"} + while True: + repos_parameters = {"public": "true", "namespace": namespace, "next_page": next_page} + repos_response = requests.get( + QUAY_REPOSITORY_API_ENDPOINT, headers=repos_headers, params=repos_parameters, timeout=MULLED_SOCKET_TIMEOUT + ) + repos_response_json = repos_response.json() + repos = repos_response_json["repositories"] + repo_names += [r["name"] for r in repos] + next_page = repos_response_json.get("next_page") + if not next_page: + break + if resolution_cache is not None: + resolution_cache[cache_key] = repo_names return repo_name in repo_names -def mulled_tags_for(namespace, image, tag_prefix=None, resolution_cache=None, session=None): +def mulled_tags_for( + namespace, image, tag_prefix=None, resolution_cache=None, session=None, expire=QUAY_VERSIONS_CACHE_EXPIRY +): """Fetch remote tags available for supplied image name. The result will be sorted so newest tags are first. @@ -94,8 +104,13 @@ def mulled_tags_for(namespace, image, tag_prefix=None, resolution_cache=None, se log.info(f"skipping mulled_tags_for [{image}] no repository") return [] - cache_key = "galaxy.tool_util.deps.container_resolvers.mulled.util:tag_cache" + cache_key = TAG_CACHE_KEY if resolution_cache is not None: + if resolution_cache.mulled_resolution_cache is not None: + # Use persistent cache if possible. Since tags query is lightweight use a relatively short expiry time. + resolution_cache = resolution_cache.mulled_resolution_cache._get_cache( + "mulled_tag_cache", {"expire": expire} + ) if cache_key not in resolution_cache: resolution_cache[cache_key] = collections.defaultdict(dict) tag_cache = resolution_cache.get(cache_key) @@ -103,10 +118,11 @@ def mulled_tags_for(namespace, image, tag_prefix=None, resolution_cache=None, se tag_cache = collections.defaultdict(dict) tags_cached = False - if namespace in tag_cache: - if image in tag_cache[namespace]: - tags = tag_cache[namespace][image] - tags_cached = True + try: + tags = tag_cache[namespace][image] + tags_cached = True + except KeyError: + pass if not tags_cached: tags = quay_versions(namespace, image, session) diff --git a/packages/tool_util/test-requirements.txt b/packages/tool_util/test-requirements.txt index 12136494321..c32c10526b8 100644 --- a/packages/tool_util/test-requirements.txt +++ b/packages/tool_util/test-requirements.txt @@ -1,2 +1,3 @@ +beaker pytest pytest-mock diff --git a/test/unit/tool_util/test_resolution_cache.py b/test/unit/tool_util/test_resolution_cache.py new file mode 100644 index 00000000000..3ca5ae163e0 --- /dev/null +++ b/test/unit/tool_util/test_resolution_cache.py @@ -0,0 +1,61 @@ +import time + +import pytest +from beaker.cache import CacheManager +from beaker.util import parse_cache_config_options + +from galaxy.tool_util.deps.container_resolvers import ResolutionCache +from galaxy.tool_util.deps.container_resolvers.mulled import mulled_tags_for +from galaxy.tool_util.deps.mulled.util import ( + _namespace_has_repo_name, + NAMESPACE_HAS_REPO_NAME_KEY, + TAG_CACHE_KEY, +) + + +@pytest.fixture() +def resolution_cache(tmpdir): + resolution_cache = ResolutionCache() + cache_opts = { + "cache.type": "file", + "cache.data_dir": str(tmpdir / "data"), + "cache.lock_dir": str(tmpdir / "lock"), + "cache.expire": "1", + } + cm = CacheManager(**parse_cache_config_options(cache_opts)).get_cache( + "mulled_resolution" + ) + resolution_cache.mulled_resolution_cache = cm + return resolution_cache + + +def test_resolution_cache_namepace_has_repo_name(resolution_cache): + resolution_cache.mulled_resolution_cache[NAMESPACE_HAS_REPO_NAME_KEY] = [ + "mytool3000" + ] + assert _namespace_has_repo_name( + "bioconda", "mytool3000", resolution_cache=resolution_cache + ) + + +def test_resolution_cache_expires(resolution_cache): + resolution_cache.mulled_resolution_cache[NAMESPACE_HAS_REPO_NAME_KEY] = [ + "mytool3000" + ] + assert NAMESPACE_HAS_REPO_NAME_KEY in resolution_cache.mulled_resolution_cache + time.sleep(1.2) + assert NAMESPACE_HAS_REPO_NAME_KEY not in resolution_cache.mulled_resolution_cache + + +def test_targets_to_mulled_name(resolution_cache): + resolution_cache.mulled_resolution_cache[NAMESPACE_HAS_REPO_NAME_KEY] = [ + "mytool3000" + ] + cache = resolution_cache.mulled_resolution_cache._get_cache( + "mulled_tag_cache", {"expire": 1} + ) + cache[TAG_CACHE_KEY] = {"bioconda": {"mytool3000": ["1.0", "1.1"]}} + tags = mulled_tags_for( + namespace="bioconda", image="mytool3000", resolution_cache=resolution_cache + ) + assert tags == ["1.1", "1.0"] From d99f26decd353cd0db1cbd6a4830546f1eaa0fc1 Mon Sep 17 00:00:00 2001 From: Kaivan Kamali Date: Wed, 20 Apr 2022 09:57:11 -0400 Subject: [PATCH 2/6] Added exist_ok flag to other object stores --- lib/galaxy/objectstore/azure_blob.py | 6 +++--- lib/galaxy/objectstore/cloud.py | 6 +++--- lib/galaxy/objectstore/pithos.py | 8 ++++---- lib/galaxy/objectstore/s3.py | 6 +++--- lib/galaxy/objectstore/unittest_utils/__init__.py | 2 +- 5 files changed, 14 insertions(+), 14 deletions(-) diff --git a/lib/galaxy/objectstore/azure_blob.py b/lib/galaxy/objectstore/azure_blob.py index 51405dc4fa7..613938f338f 100644 --- a/lib/galaxy/objectstore/azure_blob.py +++ b/lib/galaxy/objectstore/azure_blob.py @@ -246,7 +246,7 @@ class AzureBlobObjectStore(ConcreteObjectStore): # Ensure the cache directory structure exists (e.g., dataset_#_files/) rel_path_dir = os.path.dirname(rel_path) if not os.path.exists(self._get_cache_path(rel_path_dir)): - os.makedirs(self._get_cache_path(rel_path_dir)) + os.makedirs(self._get_cache_path(rel_path_dir), exist_ok=True) # Now pull in the file file_ok = self._download(rel_path) self._fix_permissions(self._get_cache_path(rel_path_dir)) @@ -327,7 +327,7 @@ class AzureBlobObjectStore(ConcreteObjectStore): # for JOB_WORK directory elif base_dir: if not os.path.exists(rel_path): - os.makedirs(rel_path) + os.makedirs(rel_path, exist_ok=True) return True else: return False @@ -381,7 +381,7 @@ class AzureBlobObjectStore(ConcreteObjectStore): # Create given directory in cache cache_dir = os.path.join(self.staging_path, rel_path) if not os.path.exists(cache_dir): - os.makedirs(cache_dir) + os.makedirs(cache_dir, exist_ok=True) # Although not really necessary to create S3 folders (because S3 has # flat namespace), do so for consistency with the regular file system diff --git a/lib/galaxy/objectstore/cloud.py b/lib/galaxy/objectstore/cloud.py index 920f42d315e..50e4ad31e57 100644 --- a/lib/galaxy/objectstore/cloud.py +++ b/lib/galaxy/objectstore/cloud.py @@ -417,7 +417,7 @@ class Cloud(ConcreteObjectStore, CloudConfigMixin): # Ensure the cache directory structure exists (e.g., dataset_#_files/) rel_path_dir = os.path.dirname(rel_path) if not os.path.exists(self._get_cache_path(rel_path_dir)): - os.makedirs(self._get_cache_path(rel_path_dir)) + os.makedirs(self._get_cache_path(rel_path_dir), exist_ok=True) # Now pull in the file file_ok = self._download(rel_path) self._fix_permissions(self._get_cache_path(rel_path_dir)) @@ -529,7 +529,7 @@ class Cloud(ConcreteObjectStore, CloudConfigMixin): # for JOB_WORK directory elif base_dir: if not os.path.exists(rel_path): - os.makedirs(rel_path) + os.makedirs(rel_path, exist_ok=True) return True else: return False @@ -565,7 +565,7 @@ class Cloud(ConcreteObjectStore, CloudConfigMixin): # Create given directory in cache cache_dir = os.path.join(self.staging_path, rel_path) if not os.path.exists(cache_dir): - os.makedirs(cache_dir) + os.makedirs(cache_dir, exist_ok=True) if not dir_only: rel_path = os.path.join(rel_path, alt_name if alt_name else f"dataset_{self._get_object_id(obj)}.dat") diff --git a/lib/galaxy/objectstore/pithos.py b/lib/galaxy/objectstore/pithos.py index 501e21afd17..b4b29a03119 100644 --- a/lib/galaxy/objectstore/pithos.py +++ b/lib/galaxy/objectstore/pithos.py @@ -210,7 +210,7 @@ class PithosObjectStore(ConcreteObjectStore): rel_path_dir = os.path.dirname(rel_path) rel_cache_path_dir = self._get_cache_path(rel_path_dir) if not os.path.exists(rel_cache_path_dir): - os.makedirs(self._get_cache_path(rel_path_dir)) + os.makedirs(self._get_cache_path(rel_path_dir), exist_ok=True) # Now pull in the file cache_path = self._get_cache_path(rel_path_dir) self.pithos.download_object(rel_path, cache_path) @@ -239,7 +239,7 @@ class PithosObjectStore(ConcreteObjectStore): return True elif base_dir: # for JOB_WORK directory if not os.path.exists(path): - os.makedirs(path) + os.makedirs(path, exist_ok=True) return True return False @@ -273,7 +273,7 @@ class PithosObjectStore(ConcreteObjectStore): # Create given directory in cache cache_dir = os.path.join(self.staging_path, rel_path) if not os.path.exists(cache_dir): - os.makedirs(cache_dir) + os.makedirs(cache_dir, exist_ok=True) if dir_only: self.pithos.upload_from_string( @@ -378,7 +378,7 @@ class PithosObjectStore(ConcreteObjectStore): cache_path = self._get_cache_path(path) if dir_only: if not os.path.exists(cache_path): - os.makedirs(cache_path) + os.makedirs(cache_path, exist_ok=True) return cache_path if self._in_cache(path): return cache_path diff --git a/lib/galaxy/objectstore/s3.py b/lib/galaxy/objectstore/s3.py index 40b91a567ce..254076f0739 100644 --- a/lib/galaxy/objectstore/s3.py +++ b/lib/galaxy/objectstore/s3.py @@ -420,7 +420,7 @@ class S3ObjectStore(ConcreteObjectStore, CloudConfigMixin): # Ensure the cache directory structure exists (e.g., dataset_#_files/) rel_path_dir = os.path.dirname(rel_path) if not os.path.exists(self._get_cache_path(rel_path_dir)): - os.makedirs(self._get_cache_path(rel_path_dir)) + os.makedirs(self._get_cache_path(rel_path_dir), exist_ok=True) # Now pull in the file file_ok = self._download(rel_path) self._fix_permissions(self._get_cache_path(rel_path_dir)) @@ -529,7 +529,7 @@ class S3ObjectStore(ConcreteObjectStore, CloudConfigMixin): # for JOB_WORK directory elif base_dir: if not os.path.exists(rel_path): - os.makedirs(rel_path) + os.makedirs(rel_path, exist_ok=True) return True else: return False @@ -565,7 +565,7 @@ class S3ObjectStore(ConcreteObjectStore, CloudConfigMixin): # Create given directory in cache cache_dir = os.path.join(self.staging_path, rel_path) if not os.path.exists(cache_dir): - os.makedirs(cache_dir) + os.makedirs(cache_dir, exist_ok=True) # Although not really necessary to create S3 folders (because S3 has # flat namespace), do so for consistency with the regular file system diff --git a/lib/galaxy/objectstore/unittest_utils/__init__.py b/lib/galaxy/objectstore/unittest_utils/__init__.py index 938ff2c6917..aa9c4d3016c 100644 --- a/lib/galaxy/objectstore/unittest_utils/__init__.py +++ b/lib/galaxy/objectstore/unittest_utils/__init__.py @@ -57,7 +57,7 @@ class Config: path = os.path.join(self.temp_directory, name) directory = os.path.dirname(path) if not os.path.exists(directory): - os.makedirs(directory) + os.makedirs(directory, exist_ok=True) contents_template = Template(contents) expanded_contents = contents_template.safe_substitute(temp_directory=self.temp_directory) open(path, "w").write(expanded_contents) From 8a810387012088dd051f1be6fd4e7c38282e9df8 Mon Sep 17 00:00:00 2001 From: Nate Coraor Date: Tue, 19 Apr 2022 14:28:01 -0400 Subject: [PATCH 3/6] Add an option to disable searching for distributed object store datasets that don't have an object_store_id --- .../sample/object_store_conf.xml.sample | 8 ++++++- lib/galaxy/model/__init__.py | 12 ++++++++-- lib/galaxy/objectstore/__init__.py | 23 +++++++++++-------- 3 files changed, 31 insertions(+), 12 deletions(-) diff --git a/lib/galaxy/config/sample/object_store_conf.xml.sample b/lib/galaxy/config/sample/object_store_conf.xml.sample index 401e27a4655..db5163a46a7 100644 --- a/lib/galaxy/config/sample/object_store_conf.xml.sample +++ b/lib/galaxy/config/sample/object_store_conf.xml.sample @@ -59,10 +59,16 @@ In distributed and hierarchical world, you can choose that some backends are automatically unused whenever they become too full. Setting the maxpctfull - attribute (on top level object_store it behaves as a global default) enables + attribute (on top level backends tag it behaves as a global default) enables this, or it can be applied to individual backends to override a global setting. This only applies to disk based backends and not remote object stores. + + By default, if a dataset should exist but its object_store_id is null, all + backends will be searched until it is found. This is to aid in Galaxy + servers moving from non-distributed to distributed object stores, but this + behavior can be disabled by setting search_for_missing="false" on the top + level backends tag. -->