mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Support external and celery managed object store caches.
This commit is contained in:
@@ -215,6 +215,9 @@ def setup_periodic_tasks(config, celery_app):
|
||||
schedule_task("prune_history_audit_table", config.history_audit_table_prune_interval)
|
||||
schedule_task("cleanup_short_term_storage", config.short_term_storage_cleanup_interval)
|
||||
|
||||
if config.object_store_cache_monitor_driver in ["auto", "celery"]:
|
||||
schedule_task("clean_object_store_caches", config.object_store_cache_monitor_interval)
|
||||
|
||||
if beat_schedule:
|
||||
celery_app.conf.beat_schedule = beat_schedule
|
||||
|
||||
|
||||
@@ -35,6 +35,7 @@ from galaxy.managers.tool_data import ToolDataImportManager
|
||||
from galaxy.metadata.set_metadata import set_metadata_portable
|
||||
from galaxy.model.scoped_session import galaxy_scoped_session
|
||||
from galaxy.objectstore import BaseObjectStore
|
||||
from galaxy.objectstore.caching import check_caches
|
||||
from galaxy.schema.tasks import (
|
||||
ComputeDatasetHashTaskRequest,
|
||||
GenerateHistoryContentDownload,
|
||||
@@ -406,3 +407,8 @@ def prune_history_audit_table(sa_session: galaxy_scoped_session):
|
||||
def cleanup_short_term_storage(storage_monitor: ShortTermStorageMonitor):
|
||||
"""Cleanup short term storage."""
|
||||
storage_monitor.cleanup()
|
||||
|
||||
|
||||
@galaxy_task(action="prune object store cache directories")
|
||||
def clean_object_store_caches(object_store: BaseObjectStore):
|
||||
check_caches(object_store.cache_targets())
|
||||
|
||||
@@ -974,6 +974,35 @@ mapping:
|
||||
Configuration file for the object store
|
||||
If this is set and exists, it overrides any other objectstore settings.
|
||||
|
||||
object_store_cache_monitor_driver:
|
||||
type: str
|
||||
default: 'auto'
|
||||
required: false
|
||||
enum: ['auto', 'inprocess', 'celery', 'external']
|
||||
desc: |
|
||||
Specify where cache monitoring is driven for caching object stores
|
||||
such as S3, Azure, and iRODS. This option has no affect on disk object stores.
|
||||
For production instances, the cache should be monitored by external tools such
|
||||
as tmpwatch and this value should be set to 'external'. This will disable all
|
||||
cache monitoring in Galaxy. Alternatively, 'celery' can monitor caches using
|
||||
a periodic task or an 'inprocess' thread can be used - but this last option
|
||||
seriously limits Galaxy's ability to scale. The default of 'auto' will use
|
||||
'celery' if 'enable_celery_tasks' is set to true or 'inprocess' otherwise.
|
||||
This option serves as the default for all object stores and can be overridden
|
||||
on a per object store basis (but don't - just setup tmpwatch for all relevant
|
||||
cache paths).
|
||||
|
||||
object_store_cache_monitor_interval:
|
||||
type: int
|
||||
default: 600
|
||||
required: false
|
||||
desc: |
|
||||
For object store cache monitoring done by Galaxy, this is the interval between
|
||||
cache checking steps. This is used by both inprocess cache monitors (which we
|
||||
recommend you do not use) and by the celery task if it is configured (by setting
|
||||
enable_celery_tasks to true and not setting object_store_cache_monitor_driver to
|
||||
external).
|
||||
|
||||
object_store_cache_path:
|
||||
type: str
|
||||
default: object_store_cache
|
||||
|
||||
@@ -48,6 +48,7 @@ from .badges import (
|
||||
serialize_badges,
|
||||
StoredBadgeDict,
|
||||
)
|
||||
from .caching import CacheTarget
|
||||
|
||||
NO_SESSION_ERROR_MESSAGE = (
|
||||
"Attempted to 'create' object store entity in configuration with no database session present."
|
||||
@@ -301,9 +302,14 @@ class ObjectStore(metaclass=abc.ABCMeta):
|
||||
"""
|
||||
raise NotImplementedError()
|
||||
|
||||
@abc.abstractmethod
|
||||
def cache_targets(self) -> List[CacheTarget]:
|
||||
"""Return a list of CacheTargets used by this object store."""
|
||||
raise NotImplementedError()
|
||||
|
||||
@abc.abstractmethod
|
||||
def to_dict(self) -> Dict[str, Any]:
|
||||
raise NotImplementedError
|
||||
raise NotImplementedError()
|
||||
|
||||
@abc.abstractmethod
|
||||
def get_quota_source_map(self):
|
||||
@@ -445,6 +451,9 @@ class BaseObjectStore(ObjectStore):
|
||||
def is_private(self, obj):
|
||||
return self._invoke("is_private", obj)
|
||||
|
||||
def cache_targets(self) -> List[CacheTarget]:
|
||||
return []
|
||||
|
||||
@classmethod
|
||||
def parse_private_from_config_xml(clazz, config_xml):
|
||||
private = DEFAULT_PRIVATE
|
||||
@@ -546,6 +555,14 @@ class ConcreteObjectStore(BaseObjectStore):
|
||||
def _is_private(self, obj):
|
||||
return self.private
|
||||
|
||||
@property
|
||||
def cache_target(self) -> Optional[CacheTarget]:
|
||||
return None
|
||||
|
||||
def cache_targets(self) -> List[CacheTarget]:
|
||||
cache_target = self.cache_target
|
||||
return [cache_target] if cache_target is not None else []
|
||||
|
||||
def get_quota_source_map(self):
|
||||
quota_source_map = QuotaSourceMap(
|
||||
self.quota_source,
|
||||
@@ -899,6 +916,13 @@ class NestedObjectStore(BaseObjectStore):
|
||||
objectstore = random.choice(list(self.backends.values()))
|
||||
return objectstore.create(obj, **kwargs)
|
||||
|
||||
def cache_targets(self) -> List[CacheTarget]:
|
||||
cache_targets = []
|
||||
for backend in self.backends.values():
|
||||
cache_targets.extend(backend.cache_targets())
|
||||
# TODO: merge more intelligently - de-duplicate paths and handle conflicting sizes/percents
|
||||
return cache_targets
|
||||
|
||||
def _empty(self, obj, **kwargs):
|
||||
"""For the first backend that has this `obj`, determine if it is empty."""
|
||||
return self._call_method("_empty", obj, True, False, **kwargs)
|
||||
|
||||
@@ -99,7 +99,7 @@ class AzureBlobObjectStore(ConcreteObjectStore):
|
||||
auth_dict = config_dict["auth"]
|
||||
container_dict = config_dict["container"]
|
||||
cache_dict = config_dict.get("cache") or {}
|
||||
self.enable_cache_monitor = enable_cache_monitor(config, config_dict)
|
||||
self.enable_cache_monitor, self.cache_monitor_interval = enable_cache_monitor(config, config_dict)
|
||||
|
||||
self.account_name = auth_dict.get("account_name")
|
||||
self.account_key = auth_dict.get("account_key")
|
||||
@@ -122,7 +122,7 @@ class AzureBlobObjectStore(ConcreteObjectStore):
|
||||
if self.cache_size != -1 and self.enable_cache_monitor:
|
||||
# Convert GBs to bytes for comparison
|
||||
self.cache_size = self.cache_size * 1073741824
|
||||
self.cache_monitor = InProcessCacheMonitor(self.cache_target, 30)
|
||||
self.cache_monitor = InProcessCacheMonitor(self.cache_target, self.cache_monitor_interval)
|
||||
|
||||
def to_dict(self):
|
||||
as_dict = super().to_dict()
|
||||
|
||||
@@ -30,6 +30,11 @@ class CacheTarget(NamedTuple):
|
||||
limit: float # cache limit as a percent
|
||||
|
||||
|
||||
def check_caches(targets: List[CacheTarget]):
|
||||
for target in targets:
|
||||
check_cache(target)
|
||||
|
||||
|
||||
def check_cache(cache_target: CacheTarget):
|
||||
"""Run a step of the cache monitor."""
|
||||
total_size, file_list = _get_cache_size_files(cache_target.path)
|
||||
@@ -103,23 +108,36 @@ def parse_caching_config_dict_from_xml(config_xml):
|
||||
if len(cache_els) > 0:
|
||||
c_xml = config_xml.findall("cache")[0]
|
||||
cache_size = float(c_xml.get("size", -1))
|
||||
|
||||
staging_path = c_xml.get("path", None)
|
||||
monitor = c_xml.get("monitor", "auto")
|
||||
|
||||
cache_dict = {
|
||||
"size": cache_size,
|
||||
"path": staging_path,
|
||||
"monitor": monitor,
|
||||
}
|
||||
else:
|
||||
cache_dict = {}
|
||||
return cache_dict
|
||||
|
||||
|
||||
def enable_cache_monitor(config, config_dict):
|
||||
if getattr(config, "disable_process_management", False):
|
||||
return True
|
||||
return config_dict.get("enable_cache_monitor", True)
|
||||
def enable_cache_monitor(config, config_dict) -> Tuple[bool, int]:
|
||||
cache_config_dict = config_dict.get("cache") or {}
|
||||
default_interval = getattr(config, "object_store_cache_monitor_interval", 600)
|
||||
interval = cache_config_dict.get("monitor_interval") or default_interval
|
||||
|
||||
if getattr(config, "disable_process_management", False):
|
||||
return True, interval
|
||||
|
||||
if config_dict.get("enable_cache_monitor", False) is False:
|
||||
return False, interval
|
||||
|
||||
default_cache_driver = getattr(config, "object_store_cache_monitor_driver", "auto")
|
||||
monitor = cache_config_dict.get("monitor", default_cache_driver)
|
||||
if monitor == "auto":
|
||||
monitor = "celery" if getattr(config, "enable_celery_tasks", False) else "inprocess"
|
||||
|
||||
return monitor == "inprocess", interval
|
||||
|
||||
|
||||
class InProcessCacheMonitor:
|
||||
|
||||
@@ -87,7 +87,7 @@ class Cloud(ConcreteObjectStore, CloudConfigMixin):
|
||||
bucket_dict = config_dict["bucket"]
|
||||
connection_dict = config_dict.get("connection", {})
|
||||
cache_dict = config_dict.get("cache") or {}
|
||||
self.enable_cache_monitor = enable_cache_monitor(config, config_dict)
|
||||
self.enable_cache_monitor, self.cache_monitor_interval = enable_cache_monitor(config, config_dict)
|
||||
|
||||
self.provider = config_dict["provider"]
|
||||
self.credentials = config_dict["auth"]
|
||||
@@ -125,7 +125,7 @@ class Cloud(ConcreteObjectStore, CloudConfigMixin):
|
||||
if self.cache_size != -1 and self.enable_cache_monitor:
|
||||
# Convert GBs to bytes for comparison
|
||||
self.cache_size = self.cache_size * 1073741824
|
||||
self.cache_monitor = InProcessCacheMonitor(self.cache_target, 30)
|
||||
self.cache_monitor = InProcessCacheMonitor(self.cache_target, self.cache_monitor_interval)
|
||||
|
||||
@staticmethod
|
||||
def _get_connection(provider, credentials):
|
||||
|
||||
@@ -154,7 +154,7 @@ class S3ObjectStore(ConcreteObjectStore, CloudConfigMixin):
|
||||
bucket_dict = config_dict["bucket"]
|
||||
connection_dict = config_dict.get("connection", {})
|
||||
cache_dict = config_dict.get("cache") or {}
|
||||
self.enable_cache_monitor = enable_cache_monitor(config, config_dict)
|
||||
self.enable_cache_monitor, self.cache_monitor_interval = enable_cache_monitor(config, config_dict)
|
||||
|
||||
self.access_key = auth_dict.get("access_key")
|
||||
self.secret_key = auth_dict.get("secret_key")
|
||||
@@ -207,7 +207,7 @@ class S3ObjectStore(ConcreteObjectStore, CloudConfigMixin):
|
||||
if self.cache_size != -1 and self.enable_cache_monitor:
|
||||
# Convert GBs to bytes for comparison
|
||||
self.cache_size = self.cache_size * 1073741824
|
||||
self.cache_monitor = InProcessCacheMonitor(self.cache_target, 30)
|
||||
self.cache_monitor = InProcessCacheMonitor(self.cache_target, self.cache_monitor_interval)
|
||||
|
||||
def _configure_connection(self):
|
||||
log.debug("Configuring S3 Connection")
|
||||
|
||||
@@ -0,0 +1,53 @@
|
||||
from contextlib import contextmanager
|
||||
from typing import (
|
||||
Iterator,
|
||||
List,
|
||||
)
|
||||
|
||||
from galaxy.celery import set_thread_app
|
||||
from galaxy.celery.tasks import clean_object_store_caches
|
||||
from galaxy.di import Container
|
||||
from galaxy.objectstore import BaseObjectStore
|
||||
from galaxy.objectstore.caching import CacheTarget
|
||||
|
||||
|
||||
class MockObjectStore:
|
||||
def __init__(self, cache_targets: List[CacheTarget]):
|
||||
self._cache_targets = cache_targets
|
||||
|
||||
def cache_targets(self) -> List[CacheTarget]:
|
||||
return self._cache_targets
|
||||
|
||||
|
||||
def test_clean_object_store_caches(tmp_path):
|
||||
with celery_injected_app_container() as container:
|
||||
cache_targets: List[CacheTarget] = []
|
||||
container[BaseObjectStore] = MockObjectStore(cache_targets) # type: ignore[assignment]
|
||||
|
||||
# similar code used in object store unit tests
|
||||
cache_dir = tmp_path
|
||||
path = cache_dir / "a_file_0"
|
||||
path.write_text("this is an example file")
|
||||
|
||||
# works fine on an empty list of cache targets...
|
||||
clean_object_store_caches()
|
||||
|
||||
assert path.exists()
|
||||
|
||||
# place the file in mock object store's cache targets and
|
||||
# run the task again and the above file should be gone.
|
||||
cache_targets.append(CacheTarget(cache_dir, 1, 0.000000001))
|
||||
# works fine on an empty list of cache targets...
|
||||
clean_object_store_caches()
|
||||
|
||||
assert not path.exists()
|
||||
|
||||
|
||||
@contextmanager
|
||||
def celery_injected_app_container() -> Iterator[Container]:
|
||||
container = Container()
|
||||
set_thread_app(container)
|
||||
try:
|
||||
yield container
|
||||
finally:
|
||||
set_thread_app(None)
|
||||
@@ -4,6 +4,7 @@ from tempfile import (
|
||||
mkdtemp,
|
||||
mkstemp,
|
||||
)
|
||||
from unittest.mock import patch
|
||||
from uuid import uuid4
|
||||
|
||||
import pytest
|
||||
@@ -29,6 +30,29 @@ from galaxy.util import (
|
||||
)
|
||||
|
||||
|
||||
# Unit testing the cloud and advanced infrastructure object stores is difficult, but
|
||||
# we can at least stub out initializing and test the configuration of these things from
|
||||
# XML and dicts.
|
||||
class UninitializedPithosObjectStore(PithosObjectStore):
|
||||
def _initialize(self):
|
||||
pass
|
||||
|
||||
|
||||
class UninitializedS3ObjectStore(S3ObjectStore):
|
||||
def _initialize(self):
|
||||
pass
|
||||
|
||||
|
||||
class UninitializedAzureBlobObjectStore(AzureBlobObjectStore):
|
||||
def _initialize(self):
|
||||
pass
|
||||
|
||||
|
||||
class UninitializedCloudObjectStore(Cloud):
|
||||
def _initialize(self):
|
||||
pass
|
||||
|
||||
|
||||
def test_unlink_path():
|
||||
with pytest.raises(FileNotFoundError):
|
||||
unlink(uuid4().hex)
|
||||
@@ -389,6 +413,14 @@ def test_mixed_private():
|
||||
assert as_dict["private"] is True
|
||||
|
||||
|
||||
def test_empty_cache_targets_for_disk_nested_stores():
|
||||
with TestConfig(MIXED_STORE_BY_DISTRIBUTED_TEST_CONFIG) as (directory, object_store):
|
||||
assert len(object_store.cache_targets()) == 0
|
||||
|
||||
with TestConfig(MIXED_STORE_BY_HIERARCHICAL_TEST_CONFIG) as (directory, object_store):
|
||||
assert len(object_store.cache_targets()) == 0
|
||||
|
||||
|
||||
BADGES_TEST_1_CONFIG_XML = """<?xml version="1.0"?>
|
||||
<object_store type="disk">
|
||||
<files_dir path="${temp_directory}/files1"/>
|
||||
@@ -567,6 +599,57 @@ def test_distributed_store():
|
||||
assert len(extra_dirs) == 2
|
||||
|
||||
|
||||
def test_distributed_store_empty_cache_targets():
|
||||
for config_str in [DISTRIBUTED_TEST_CONFIG, DISTRIBUTED_TEST_CONFIG_YAML]:
|
||||
with TestConfig(config_str) as (directory, object_store):
|
||||
assert len(object_store.cache_targets()) == 0
|
||||
|
||||
|
||||
DISTRIBUTED_TEST_S3_CONFIG_YAML = """
|
||||
type: distributed
|
||||
backends:
|
||||
- id: files1
|
||||
weight: 1
|
||||
type: s3
|
||||
auth:
|
||||
access_key: access_moo
|
||||
secret_key: secret_cow
|
||||
|
||||
bucket:
|
||||
name: unique_bucket_name_all_lowercase
|
||||
use_reduced_redundancy: false
|
||||
|
||||
extra_dirs:
|
||||
- type: job_work
|
||||
path: ${temp_directory}/job_working_directory_s3
|
||||
- type: temp
|
||||
path: ${temp_directory}/tmp_s3
|
||||
- id: files2
|
||||
weight: 1
|
||||
type: s3
|
||||
auth:
|
||||
access_key: access_moo
|
||||
secret_key: secret_cow
|
||||
|
||||
bucket:
|
||||
name: unique_bucket_name_all_lowercase_2
|
||||
use_reduced_redundancy: false
|
||||
|
||||
extra_dirs:
|
||||
- type: job_work
|
||||
path: ${temp_directory}/job_working_directory_s3_2
|
||||
- type: temp
|
||||
path: ${temp_directory}/tmp_s3_2
|
||||
"""
|
||||
|
||||
|
||||
@patch("galaxy.objectstore.s3.S3ObjectStore", UninitializedS3ObjectStore)
|
||||
def test_distributed_store_with_cache_targets():
|
||||
for config_str in [DISTRIBUTED_TEST_S3_CONFIG_YAML]:
|
||||
with TestConfig(config_str) as (directory, object_store):
|
||||
assert len(object_store.cache_targets()) == 2
|
||||
|
||||
|
||||
HIERARCHICAL_MUST_HAVE_UNIFIED_QUOTA_SOURCE = """<?xml version="1.0"?>
|
||||
<object_store type="hierarchical" private="true">
|
||||
<backends>
|
||||
@@ -597,29 +680,6 @@ def test_hiercachical_backend_must_share_quota_source():
|
||||
assert the_exception is not None
|
||||
|
||||
|
||||
# Unit testing the cloud and advanced infrastructure object stores is difficult, but
|
||||
# we can at least stub out initializing and test the configuration of these things from
|
||||
# XML and dicts.
|
||||
class UninitializedPithosObjectStore(PithosObjectStore):
|
||||
def _initialize(self):
|
||||
pass
|
||||
|
||||
|
||||
class UninitializedS3ObjectStore(S3ObjectStore):
|
||||
def _initialize(self):
|
||||
pass
|
||||
|
||||
|
||||
class UninitializedAzureBlobObjectStore(AzureBlobObjectStore):
|
||||
def _initialize(self):
|
||||
pass
|
||||
|
||||
|
||||
class UninitializedCloudObjectStore(Cloud):
|
||||
def _initialize(self):
|
||||
pass
|
||||
|
||||
|
||||
PITHOS_TEST_CONFIG = """<?xml version="1.0"?>
|
||||
<object_store type="pithos">
|
||||
<auth url="http://example.org/" token="extoken123" />
|
||||
@@ -966,8 +1026,6 @@ def test_config_parse_cloud():
|
||||
_assert_key_has_value(cache_dict, "size", 1000.0)
|
||||
_assert_key_has_value(cache_dict, "path", "database/object_store_cache")
|
||||
|
||||
_assert_key_has_value(as_dict, "enable_cache_monitor", False)
|
||||
|
||||
extra_dirs = as_dict["extra_dirs"]
|
||||
assert len(extra_dirs) == 2
|
||||
|
||||
|
||||
Reference in New Issue
Block a user