diff --git a/lib/galaxy/celery/__init__.py b/lib/galaxy/celery/__init__.py index 5507a0873ec..3bf7a94f39a 100644 --- a/lib/galaxy/celery/__init__.py +++ b/lib/galaxy/celery/__init__.py @@ -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 diff --git a/lib/galaxy/celery/tasks.py b/lib/galaxy/celery/tasks.py index 2d8be88c603..186c38be766 100644 --- a/lib/galaxy/celery/tasks.py +++ b/lib/galaxy/celery/tasks.py @@ -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()) diff --git a/lib/galaxy/config/schemas/config_schema.yml b/lib/galaxy/config/schemas/config_schema.yml index 6d32efcfa7b..fc2e56ca282 100644 --- a/lib/galaxy/config/schemas/config_schema.yml +++ b/lib/galaxy/config/schemas/config_schema.yml @@ -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 diff --git a/lib/galaxy/objectstore/__init__.py b/lib/galaxy/objectstore/__init__.py index 0e9ca240d9d..c0e634ccfd0 100644 --- a/lib/galaxy/objectstore/__init__.py +++ b/lib/galaxy/objectstore/__init__.py @@ -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) diff --git a/lib/galaxy/objectstore/azure_blob.py b/lib/galaxy/objectstore/azure_blob.py index 7f4a212fceb..e7be7f835c2 100644 --- a/lib/galaxy/objectstore/azure_blob.py +++ b/lib/galaxy/objectstore/azure_blob.py @@ -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() diff --git a/lib/galaxy/objectstore/caching.py b/lib/galaxy/objectstore/caching.py index 01a12488b0a..8061af95baf 100644 --- a/lib/galaxy/objectstore/caching.py +++ b/lib/galaxy/objectstore/caching.py @@ -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: diff --git a/lib/galaxy/objectstore/cloud.py b/lib/galaxy/objectstore/cloud.py index f7060d269b9..ece9e440117 100644 --- a/lib/galaxy/objectstore/cloud.py +++ b/lib/galaxy/objectstore/cloud.py @@ -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): diff --git a/lib/galaxy/objectstore/s3.py b/lib/galaxy/objectstore/s3.py index 94fdfab27a8..8c10c6e6856 100644 --- a/lib/galaxy/objectstore/s3.py +++ b/lib/galaxy/objectstore/s3.py @@ -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") diff --git a/test/unit/app/test_tasks.py b/test/unit/app/test_tasks.py new file mode 100644 index 00000000000..72ee650c22b --- /dev/null +++ b/test/unit/app/test_tasks.py @@ -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) diff --git a/test/unit/objectstore/test_objectstore.py b/test/unit/objectstore/test_objectstore.py index 8deba551543..409ee89aeb5 100644 --- a/test/unit/objectstore/test_objectstore.py +++ b/test/unit/objectstore/test_objectstore.py @@ -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 = """ @@ -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 = """ @@ -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 = """ @@ -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