mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Merge branch 'release_22.01' into dev
This commit is contained in:
@@ -664,6 +664,48 @@
|
||||
:Type: str
|
||||
|
||||
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
``short_term_storage_dir``
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
|
||||
:Description:
|
||||
Location of files available for a short time as downloads (short
|
||||
term storage). This directory is exclusively used for serving
|
||||
dynamically generated downloadable content. Galaxy may uses the
|
||||
new_file_path parameter as a general temporary directory and that
|
||||
directory should be monitored by a tool such as tmpwatch in
|
||||
production environments. short_term_storage_dir on the other hand
|
||||
is monitored by Galaxy's task framework and should not require
|
||||
such external tooling.
|
||||
The value of this option will be resolved with respect to
|
||||
<cache_dir>.
|
||||
:Default: ``short_term_web_storage``
|
||||
:Type: str
|
||||
|
||||
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
``short_term_storage_default_duration``
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
|
||||
:Description:
|
||||
Default duration before short term web storage files will be
|
||||
cleaned up by Galaxy tasks (in seconds). The default duration is 1
|
||||
day.
|
||||
:Default: ``86400``
|
||||
:Type: int
|
||||
|
||||
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
``short_term_storage_cleanup_interval``
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
|
||||
:Description:
|
||||
How many seconds between instances of short term storage being
|
||||
cleaned up in default Celery task configuration.
|
||||
:Default: ``3600``
|
||||
:Type: int
|
||||
|
||||
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
``file_sources_config_file``
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
@@ -1286,6 +1328,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``
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
|
||||
+2
-1
@@ -261,6 +261,7 @@ class ConfiguresGalaxyMixin:
|
||||
"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"
|
||||
@@ -274,7 +275,7 @@ class ConfiguresGalaxyMixin:
|
||||
|
||||
def reindex_tool_search(self):
|
||||
# Call this when tools are added or removed.
|
||||
self.toolbox_search.build_index(tool_cache=self.tool_cache)
|
||||
self.toolbox_search.build_index(tool_cache=self.tool_cache, toolbox=self.toolbox)
|
||||
self.tool_cache.reset_status()
|
||||
|
||||
def _set_enabled_container_types(self):
|
||||
|
||||
@@ -456,6 +456,26 @@ galaxy:
|
||||
# option.
|
||||
#watch_tours: 'false'
|
||||
|
||||
# Location of files available for a short time as downloads (short
|
||||
# term storage). This directory is exclusively used for serving
|
||||
# dynamically generated downloadable content. Galaxy may uses the
|
||||
# new_file_path parameter as a general temporary directory and that
|
||||
# directory should be monitored by a tool such as tmpwatch in
|
||||
# production environments. short_term_storage_dir on the other hand is
|
||||
# monitored by Galaxy's task framework and should not require such
|
||||
# external tooling.
|
||||
# The value of this option will be resolved with respect to
|
||||
# <cache_dir>.
|
||||
#short_term_storage_dir: short_term_web_storage
|
||||
|
||||
# Default duration before short term web storage files will be cleaned
|
||||
# up by Galaxy tasks (in seconds). The default duration is 1 day.
|
||||
#short_term_storage_default_duration: 86400
|
||||
|
||||
# How many seconds between instances of short term storage being
|
||||
# cleaned up in default Celery task configuration.
|
||||
#short_term_storage_cleanup_interval: 3600
|
||||
|
||||
# Configured FileSource plugins.
|
||||
# The value of this option will be resolved with respect to
|
||||
# <config_dir>.
|
||||
@@ -749,6 +769,10 @@ galaxy:
|
||||
# <cache_dir>.
|
||||
#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
|
||||
|
||||
@@ -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.
|
||||
-->
|
||||
<!--
|
||||
<object_store type="distributed">
|
||||
|
||||
@@ -961,6 +961,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
|
||||
|
||||
@@ -3410,12 +3410,20 @@ class Dataset(StorableObject, Serializable, _HasTable):
|
||||
return self.state in self.ready_states
|
||||
|
||||
def get_file_name(self):
|
||||
if self.purged:
|
||||
log.warning(f"Attempt to get file name of purged dataset {self.id}")
|
||||
return ""
|
||||
if not self.external_filename:
|
||||
assert self.object_store is not None, f"Object Store has not been initialized for dataset {self.id}"
|
||||
if self.object_store.exists(self):
|
||||
return self.object_store.get_filename(self)
|
||||
file_name = self.object_store.get_filename(self)
|
||||
else:
|
||||
return ""
|
||||
file_name = ""
|
||||
if not file_name and self.state not in (self.states.NEW, self.states.QUEUED):
|
||||
# Queued datasets can be assigned an object store and have a filename, but they aren't guaranteed to.
|
||||
# Anything after queued should have a file name.
|
||||
log.warning(f"Failed to determine file name for dataset {self.id}")
|
||||
return file_name
|
||||
else:
|
||||
filename = self.external_filename
|
||||
# Make filename absolute
|
||||
|
||||
@@ -26,6 +26,7 @@ from galaxy.exceptions import (
|
||||
ObjectNotFound,
|
||||
)
|
||||
from galaxy.util import (
|
||||
asbool,
|
||||
directory_hash_id,
|
||||
force_symlink,
|
||||
parse_xml,
|
||||
@@ -850,6 +851,7 @@ class DistributedObjectStore(NestedObjectStore):
|
||||
self.original_weighted_backend_ids = []
|
||||
self.max_percent_full = {}
|
||||
self.global_max_percent_full = config_dict.get("global_max_percent_full", 0)
|
||||
self.search_for_missing = config_dict.get("search_for_missing", True)
|
||||
random.seed()
|
||||
|
||||
for backend_def in config_dict["backends"]:
|
||||
@@ -887,6 +889,7 @@ class DistributedObjectStore(NestedObjectStore):
|
||||
|
||||
backends: List[Dict[str, Any]] = []
|
||||
config_dict = {
|
||||
"search_for_missing": asbool(backends_root.get("search_for_missing", True)),
|
||||
"global_max_percent_full": float(backends_root.get("maxpctfull", 0)),
|
||||
"backends": backends,
|
||||
}
|
||||
@@ -933,6 +936,7 @@ class DistributedObjectStore(NestedObjectStore):
|
||||
def to_dict(self) -> Dict[str, Any]:
|
||||
as_dict = super().to_dict()
|
||||
as_dict["global_max_percent_full"] = self.global_max_percent_full
|
||||
as_dict["search_for_missing"] = self.search_for_missing
|
||||
backends: List[Dict[str, Any]] = []
|
||||
for backend_id, backend in self.backends.items():
|
||||
backend_as_dict = backend.to_dict()
|
||||
@@ -1003,17 +1007,17 @@ class DistributedObjectStore(NestedObjectStore):
|
||||
"The backend object store ID (%s) for %s object with ID %s is invalid"
|
||||
% (obj.object_store_id, obj.__class__.__name__, obj.id)
|
||||
)
|
||||
# if this instance has been switched from a non-distributed to a
|
||||
# distributed object store, or if the object's store id is invalid,
|
||||
# try to locate the object
|
||||
for id, store in self.backends.items():
|
||||
if store.exists(obj, **kwargs):
|
||||
log.warning(
|
||||
"%s object with ID %s found in backend object store with ID %s"
|
||||
% (obj.__class__.__name__, obj.id, id)
|
||||
)
|
||||
obj.object_store_id = id
|
||||
return id
|
||||
elif self.search_for_missing:
|
||||
# if this instance has been switched from a non-distributed to a
|
||||
# distributed object store, or if the object's store id is invalid,
|
||||
# try to locate the object
|
||||
for id, store in self.backends.items():
|
||||
if store.exists(obj, **kwargs):
|
||||
log.warning(
|
||||
f"{obj.__class__.__name__} object with ID {obj.id} found in backend object store with ID {id}"
|
||||
)
|
||||
obj.object_store_id = id
|
||||
return id
|
||||
return None
|
||||
|
||||
|
||||
|
||||
@@ -261,7 +261,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))
|
||||
@@ -368,7 +368,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
|
||||
@@ -422,7 +422,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
|
||||
|
||||
@@ -434,7 +434,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))
|
||||
@@ -567,7 +567,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
|
||||
@@ -603,7 +603,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")
|
||||
|
||||
@@ -225,7 +225,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)
|
||||
@@ -254,7 +254,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
|
||||
|
||||
@@ -288,7 +288,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(rel_path, "", content_type="application/directory")
|
||||
@@ -389,7 +389,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
|
||||
|
||||
@@ -442,7 +442,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))
|
||||
@@ -573,7 +573,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
|
||||
@@ -609,7 +609,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
|
||||
|
||||
@@ -56,7 +56,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)
|
||||
|
||||
@@ -16,6 +16,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."""
|
||||
|
||||
@@ -20,7 +20,10 @@ from galaxy.util import (
|
||||
)
|
||||
from galaxy.util.commands import shell
|
||||
from ..container_classes import CONTAINER_CLASSES
|
||||
from ..container_resolvers import ContainerResolver
|
||||
from ..container_resolvers import (
|
||||
ContainerResolver,
|
||||
ResolutionCache,
|
||||
)
|
||||
from ..docker_util import build_docker_images_command
|
||||
from ..mulled.mulled_build import (
|
||||
DEFAULT_CHANNELS,
|
||||
@@ -320,7 +323,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:
|
||||
@@ -330,15 +335,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:
|
||||
|
||||
@@ -19,6 +19,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 default_mulled_conda_channels_from_env():
|
||||
@@ -68,33 +71,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.
|
||||
@@ -106,8 +112,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)
|
||||
@@ -115,10 +126,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)
|
||||
|
||||
@@ -68,9 +68,7 @@ class ToolBoxSearch:
|
||||
for panel_view in toolbox.panel_views():
|
||||
panel_view_id = panel_view.id
|
||||
panel_index_dir = os.path.join(index_dir, panel_view_id)
|
||||
panel_searches[panel_view_id] = ToolPanelViewSearch(
|
||||
toolbox, panel_view_id, panel_index_dir, index_help=index_help
|
||||
)
|
||||
panel_searches[panel_view_id] = ToolPanelViewSearch(panel_view_id, panel_index_dir, index_help=index_help)
|
||||
self.panel_searches = panel_searches
|
||||
# We keep track of how many times the tool index has been rebuilt.
|
||||
# We start at -1, so that after the first index the count is at 0,
|
||||
@@ -78,10 +76,10 @@ class ToolBoxSearch:
|
||||
# reindexing if the index count is equal to the toolbox reload count.
|
||||
self.index_count = -1
|
||||
|
||||
def build_index(self, tool_cache, index_help: bool = True) -> None:
|
||||
def build_index(self, tool_cache, toolbox, index_help: bool = True) -> None:
|
||||
self.index_count += 1
|
||||
for panel_search in self.panel_searches.values():
|
||||
panel_search.build_index(tool_cache, index_help=index_help)
|
||||
panel_search.build_index(tool_cache, toolbox, index_help=index_help)
|
||||
|
||||
def search(self, *args, **kwd) -> List[str]:
|
||||
panel_view = kwd.pop("panel_view")
|
||||
@@ -97,7 +95,7 @@ class ToolPanelViewSearch:
|
||||
the Whoosh search library.
|
||||
"""
|
||||
|
||||
def __init__(self, toolbox, panel_view_id: str, index_dir: str, index_help: bool = True):
|
||||
def __init__(self, panel_view_id: str, index_dir: str, index_help: bool = True):
|
||||
self.schema = Schema(
|
||||
id=ID(stored=True, unique=True),
|
||||
old_id=ID,
|
||||
@@ -110,14 +108,13 @@ class ToolPanelViewSearch:
|
||||
)
|
||||
self.rex = analysis.RegexTokenizer()
|
||||
self.index_dir = index_dir
|
||||
self.toolbox = toolbox
|
||||
self.panel_view_id = panel_view_id
|
||||
self.index = self._index_setup()
|
||||
|
||||
def _index_setup(self) -> index.Index:
|
||||
return get_or_create_index(index_dir=self.index_dir, schema=self.schema)
|
||||
|
||||
def build_index(self, tool_cache, index_help: bool = True) -> None:
|
||||
def build_index(self, tool_cache, toolbox, index_help: bool = True) -> None:
|
||||
"""
|
||||
Prepare search index for tools loaded in toolbox.
|
||||
Use `tool_cache` to determine which tools need indexing and which tools should be expired.
|
||||
@@ -143,8 +140,8 @@ class ToolPanelViewSearch:
|
||||
for tool_id in tool_ids_to_remove:
|
||||
writer.delete_by_term("id", tool_id)
|
||||
for tool_id in tool_cache._new_tool_ids - indexed_tool_ids:
|
||||
tool = self.toolbox.get_tool(tool_id)
|
||||
if tool and tool.is_latest_version and self.toolbox.panel_has_tool(tool, self.panel_view_id):
|
||||
tool = toolbox.get_tool(tool_id)
|
||||
if tool and tool.is_latest_version and toolbox.panel_has_tool(tool, self.panel_view_id):
|
||||
if tool.hidden:
|
||||
# we check if there is an older tool we can return
|
||||
if tool.lineage:
|
||||
|
||||
@@ -1,2 +1,3 @@
|
||||
beaker
|
||||
pytest
|
||||
pytest-mock
|
||||
|
||||
@@ -0,0 +1,47 @@
|
||||
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.mulled.util import (
|
||||
_namespace_has_repo_name,
|
||||
mulled_tags_for,
|
||||
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"]
|
||||
Reference in New Issue
Block a user