diff --git a/lib/galaxy/objectstore/__init__.py b/lib/galaxy/objectstore/__init__.py index b986a0e81eb..2d5608b3c06 100644 --- a/lib/galaxy/objectstore/__init__.py +++ b/lib/galaxy/objectstore/__init__.py @@ -721,26 +721,22 @@ class DistributedObjectStore(NestedObjectStore): self.global_max_percent_full = config_dict.get("global_max_percent_full", 0) random.seed() - backends_def = config_dict["backends"] - for backend_def in backends_def: + for backend_def in config_dict["backends"]: backened_id = backend_def["id"] - file_path = backend_def["files_dir"] - extra_dirs = backend_def.get("extra_dirs", []) maxpctfull = backend_def.get("max_percent_full", 0) weight = backend_def["weight"] - store_by = backend_def.get("store_by") - disk_config_dict = dict(files_dir=file_path, extra_dirs=extra_dirs) - if store_by is not None: - disk_config_dict['store_by'] = store_by - self.backends[backened_id] = DiskObjectStore(config, disk_config_dict) + + backend = build_object_store_from_config(config, config_dict=backend_def, fsmon=fsmon) + + self.backends[backened_id] = backend self.max_percent_full[backened_id] = maxpctfull - log.debug("Loaded disk backend '%s' with weight %s and file_path: %s" % (backened_id, weight, file_path)) for i in range(0, weight): # The simplest way to do weighting: add backend ids to a # sequence the number of times equalling weight, then randomly # choose a backend from that sequence at creation self.weighted_backend_ids.append(backened_id) + self.original_weighted_backend_ids = self.weighted_backend_ids self.sleeper = None @@ -764,33 +760,22 @@ class DistributedObjectStore(NestedObjectStore): 'backends': backends, } - for elem in [e for e in backends_root if e.tag == 'backend']: - id = elem.get('id') - weight = int(elem.get('weight', 1)) - maxpctfull = float(elem.get('maxpctfull', 0)) - elem_type = elem.get('type', 'disk') - store_by = elem.get('store_by', None) - if elem_type: - path = None - extra_dirs = [] - for sub in elem: - if sub.tag == 'files_dir': - path = sub.get('path') - elif sub.tag == 'extra_dir': - type = sub.get('type') - extra_dirs.append({"type": type, "path": sub.get('path')}) + for b in [e for e in backends_root if e.tag == 'backend']: + store_id = b.get("id") + store_weight = int(b.get("weight", 1)) + store_maxpctfull = float(b.get('maxpctfull', 0)) + store_type = b.get("type", "disk") + store_by = b.get('store_by', None) - backend_dict = { - 'id': id, - 'weight': weight, - 'max_percent_full': maxpctfull, - 'files_dir': path, - 'extra_dirs': extra_dirs, - 'type': elem_type, - } - if store_by is not None: - backend_dict['store_by'] = store_by - backends.append(backend_dict) + objectstore_class, _ = type_to_object_store_class(store_type) + backend_config_dict = objectstore_class.parse_xml(b) + backend_config_dict["id"] = store_id + backend_config_dict["weight"] = store_weight + backend_config_dict["max_percent_full"] = store_maxpctfull + backend_config_dict["type"] = store_type + if store_by is not None: + backend_config_dict["store_by"] = store_by + backends.append(backend_config_dict) return config_dict diff --git a/lib/galaxy/objectstore/azure_blob.py b/lib/galaxy/objectstore/azure_blob.py index b47c152edb2..f0431e67640 100644 --- a/lib/galaxy/objectstore/azure_blob.py +++ b/lib/galaxy/objectstore/azure_blob.py @@ -173,7 +173,7 @@ class AzureBlobObjectStore(ConcreteObjectStore): # follow them, so if they are valid we normalize them out alt_name = os.path.normpath(alt_name) - rel_path = os.path.join(*directory_hash_id(obj.id)) + rel_path = os.path.join(*directory_hash_id(self._get_object_id(obj))) if extra_dir is not None: if extra_dir_at_root: @@ -183,7 +183,7 @@ class AzureBlobObjectStore(ConcreteObjectStore): # for JOB_WORK directory if obj_dir: - rel_path = os.path.join(rel_path, str(obj.id)) + rel_path = os.path.join(rel_path, str(self._get_object_id(obj))) if base_dir: base = self.extra_dirs.get(base_dir) return os.path.join(base, rel_path) @@ -192,7 +192,7 @@ class AzureBlobObjectStore(ConcreteObjectStore): # rel_path = '%s/' % rel_path # assume for now we don't need this in Azure blob storage. if not dir_only: - rel_path = os.path.join(rel_path, alt_name if alt_name else "dataset_%s.dat" % obj.id) + rel_path = os.path.join(rel_path, alt_name if alt_name else "dataset_%s.dat" % self._get_object_id(obj)) return rel_path @@ -368,7 +368,7 @@ class AzureBlobObjectStore(ConcreteObjectStore): alt_name = kwargs.get('alt_name', None) # Construct hashed path - rel_path = os.path.join(*directory_hash_id(obj.id)) + rel_path = os.path.join(*directory_hash_id(self._get_object_id(obj))) # Optionally append extra_dir if extra_dir is not None: @@ -389,7 +389,7 @@ class AzureBlobObjectStore(ConcreteObjectStore): # self._push_to_os(s3_dir, from_string='') # If instructed, create the dataset in cache & in S3 if not dir_only: - rel_path = os.path.join(rel_path, alt_name if alt_name else "dataset_%s.dat" % obj.id) + rel_path = os.path.join(rel_path, alt_name if alt_name else "dataset_%s.dat" % self._get_object_id(obj)) open(os.path.join(self.staging_path, rel_path), 'w').close() self._push_to_os(rel_path, from_string='') diff --git a/lib/galaxy/objectstore/cloud.py b/lib/galaxy/objectstore/cloud.py index b94c9eba9ea..c3ef638b438 100644 --- a/lib/galaxy/objectstore/cloud.py +++ b/lib/galaxy/objectstore/cloud.py @@ -351,7 +351,7 @@ class Cloud(ConcreteObjectStore, CloudConfigMixin): # alt_name can contain parent directory references, but S3 will not # follow them, so if they are valid we normalize them out alt_name = os.path.normpath(alt_name) - rel_path = os.path.join(*directory_hash_id(obj.id)) + rel_path = os.path.join(*directory_hash_id(self._get_object_id(obj))) if extra_dir is not None: if extra_dir_at_root: rel_path = os.path.join(extra_dir, rel_path) @@ -360,7 +360,7 @@ class Cloud(ConcreteObjectStore, CloudConfigMixin): # for JOB_WORK directory if obj_dir: - rel_path = os.path.join(rel_path, str(obj.id)) + rel_path = os.path.join(rel_path, str(self._get_object_id(obj))) if base_dir: base = self.extra_dirs.get(base_dir) return os.path.join(base, rel_path) @@ -369,7 +369,7 @@ class Cloud(ConcreteObjectStore, CloudConfigMixin): rel_path = '%s/' % rel_path if not dir_only: - rel_path = os.path.join(rel_path, alt_name if alt_name else "dataset_%s.dat" % obj.id) + rel_path = os.path.join(rel_path, alt_name if alt_name else "dataset_%s.dat" % self._get_object_id(obj)) return rel_path def _get_cache_path(self, rel_path): @@ -553,7 +553,7 @@ class Cloud(ConcreteObjectStore, CloudConfigMixin): alt_name = kwargs.get('alt_name', None) # Construct hashed path - rel_path = os.path.join(*directory_hash_id(obj.id)) + rel_path = os.path.join(*directory_hash_id(self._get_object_id(obj))) # Optionally append extra_dir if extra_dir is not None: @@ -568,7 +568,7 @@ class Cloud(ConcreteObjectStore, CloudConfigMixin): os.makedirs(cache_dir) if not dir_only: - rel_path = os.path.join(rel_path, alt_name if alt_name else "dataset_%s.dat" % obj.id) + rel_path = os.path.join(rel_path, alt_name if alt_name else "dataset_%s.dat" % self._get_object_id(obj)) open(os.path.join(self.staging_path, rel_path), 'w').close() self._push_to_os(rel_path, from_string='') diff --git a/lib/galaxy/objectstore/irods.py b/lib/galaxy/objectstore/irods.py index fed2a2fa184..8abd8f2d352 100644 --- a/lib/galaxy/objectstore/irods.py +++ b/lib/galaxy/objectstore/irods.py @@ -258,7 +258,7 @@ class IRODSObjectStore(DiskObjectStore, CloudConfigMixin): # alt_name can contain parent directory references, but S3 will not # follow them, so if they are valid we normalize them out alt_name = os.path.normpath(alt_name) - rel_path = os.path.join(*directory_hash_id(obj.id)) + rel_path = os.path.join(*directory_hash_id(self._get_object_id(obj))) if extra_dir is not None: if extra_dir_at_root: rel_path = os.path.join(extra_dir, rel_path) @@ -267,13 +267,13 @@ class IRODSObjectStore(DiskObjectStore, CloudConfigMixin): # for JOB_WORK directory if obj_dir: - rel_path = os.path.join(rel_path, str(obj.id)) + rel_path = os.path.join(rel_path, str(self._get_object_id(obj))) if base_dir: base = self.extra_dirs.get(base_dir) return os.path.join(base, rel_path) if not dir_only: - rel_path = os.path.join(rel_path, alt_name if alt_name else "dataset_%s.dat" % obj.id) + rel_path = os.path.join(rel_path, alt_name if alt_name else "dataset_%s.dat" % self._get_object_id(obj)) return rel_path def _get_cache_path(self, rel_path): @@ -460,7 +460,7 @@ class IRODSObjectStore(DiskObjectStore, CloudConfigMixin): alt_name = kwargs.get('alt_name', None) # Construct hashed path - rel_path = os.path.join(*directory_hash_id(obj.id)) + rel_path = os.path.join(*directory_hash_id(self._get_object_id(obj))) # Optionally append extra_dir if extra_dir is not None: @@ -475,7 +475,7 @@ class IRODSObjectStore(DiskObjectStore, CloudConfigMixin): os.makedirs(cache_dir) if not dir_only: - rel_path = os.path.join(rel_path, alt_name if alt_name else "dataset_%s.dat" % obj.id) + rel_path = os.path.join(rel_path, alt_name if alt_name else "dataset_%s.dat" % self._get_object_id(obj)) open(os.path.join(self.staging_path, rel_path), 'w').close() self._push_to_irods(rel_path, from_string='') diff --git a/lib/galaxy/objectstore/pithos.py b/lib/galaxy/objectstore/pithos.py index 63bc8e23e94..50aaa79e589 100644 --- a/lib/galaxy/objectstore/pithos.py +++ b/lib/galaxy/objectstore/pithos.py @@ -163,7 +163,7 @@ class PithosObjectStore(ConcreteObjectStore): # alt_name can contain parent directory references, but S3 will not # follow them, so if they are valid we normalize them out alt_name = os.path.normpath(alt_name) - rel_path = os.path.join(*directory_hash_id(obj.id)) + rel_path = os.path.join(*directory_hash_id(self._get_object_id(obj))) if extra_dir is not None: if extra_dir_at_root: rel_path = os.path.join(extra_dir, rel_path) @@ -172,7 +172,7 @@ class PithosObjectStore(ConcreteObjectStore): # for JOB_WORK directory if obj_dir: - rel_path = os.path.join(rel_path, str(obj.id)) + rel_path = os.path.join(rel_path, str(self._get_object_id(obj))) if base_dir: base = self.extra_dirs.get(base_dir) return os.path.join(base, rel_path) @@ -181,7 +181,7 @@ class PithosObjectStore(ConcreteObjectStore): rel_path = '{0}/'.format(rel_path) if not dir_only: - an = alt_name if alt_name else 'dataset_{0}.dat'.format(obj.id) + an = alt_name if alt_name else 'dataset_{0}.dat'.format(self._get_object_id(obj)) rel_path = os.path.join(rel_path, an) return rel_path @@ -263,7 +263,7 @@ class PithosObjectStore(ConcreteObjectStore): alt_name = kwargs.get('alt_name', None) # Construct hashed path - rel_path = os.path.join(*directory_hash_id(obj.id)) + rel_path = os.path.join(*directory_hash_id(self._get_object_id(obj))) # Optionally append extra_dir if extra_dir is not None: @@ -283,7 +283,7 @@ class PithosObjectStore(ConcreteObjectStore): else: rel_path = os.path.join( rel_path, - alt_name if alt_name else 'dataset_{0}.dat'.format(obj.id)) + alt_name if alt_name else 'dataset_{0}.dat'.format(self._get_object_id(obj))) new_file = os.path.join(self.staging_path, rel_path) open(new_file, 'w').close() self.pithos.upload_from_string(rel_path, '') diff --git a/lib/tool_shed/webapp/model/__init__.py b/lib/tool_shed/webapp/model/__init__.py index d1981ed661a..f087bcfcef3 100644 --- a/lib/tool_shed/webapp/model/__init__.py +++ b/lib/tool_shed/webapp/model/__init__.py @@ -6,11 +6,6 @@ from datetime import ( timedelta ) -from mercurial import ( - hg, - ui -) - import tool_shed.repository_types.util as rt_util from galaxy import util from galaxy.model.orm.now import now @@ -202,6 +197,10 @@ class Repository(Dictifiable): @property def hg_repo(self): + from mercurial import ( + hg, + ui + ) if not WEAK_HG_REPO_CACHE.get(self): WEAK_HG_REPO_CACHE[self] = hg.cachedlocalrepo(hg.repository(ui.ui(), self.repo_path().encode('utf-8'))) return WEAK_HG_REPO_CACHE[self].fetch()[0] diff --git a/test/integration/objectstore/test_objectstore_datatype_upload.py b/test/integration/objectstore/test_objectstore_datatype_upload.py new file mode 100644 index 00000000000..8651813f8c4 --- /dev/null +++ b/test/integration/objectstore/test_objectstore_datatype_upload.py @@ -0,0 +1,188 @@ +"""Tests objectstores by exercising the datatype upload integration tests.""" +import os +import string +import subprocess +import time + +import pytest + +from galaxy_test.driver import integration_util +from ..test_datatype_upload import ( + TEST_CASES, + upload_datatype_helper, + UploadTestDatatypeDataTestCase +) + +SCRIPT_DIRECTORY = os.path.abspath(os.path.dirname(__file__)) +IRODS_OBJECT_STORE_HOST = os.environ.get('GALAXY_INTEGRATION_IRODS_OBJECT_STORE_HOST', 'localhost') +IRODS_OBJECT_STORE_PORT = int(os.environ.get('GALAXY_INTEGRATION_IRODS_OBJECT_STORE_PORT', 1247)) +IRODS_OBJECT_STORE_TIMEOUT = int(os.environ.get('GALAXY_INTEGRATION_IRODS_OBJECT_STORE_TIMEOUT', 30)) +IRODS_OBJECT_STORE_USERNAME = os.environ.get('GALAXY_INTEGRATION_IRODS_OBJECT_STORE_USERNAME', 'rods') +IRODS_OBJECT_STORE_PASSWORD = os.environ.get('GALAXY_INTEGRATION_IRODS_OBJECT_STORE_PASSWORD', 'rods') +IRODS_OBJECT_STORE_RESOURCE = os.environ.get('GALAXY_INTEGRATION_IRODS_OBJECT_STORE_RESOURCE', 'demoResc') +IRODS_OBJECT_STORE_ZONE = os.environ.get('GALAXY_INTEGRATION_IRODS_OBJECT_STORE_ZONE', 'tempZone') +# Run test for only the first 10 test files +TEST_CASES = dict(list(TEST_CASES.items())[0:10]) +DISTRIBUTED_OBJECT_STORE_CONFIG = string.Template(""" + + + + + + + + + + + + + + +""") +DISTRIBUTED_IRODS_OBJECT_STORE_CONFIG = string.Template(""" + + + + + + + + + + + + + + + + + + +""") +IRODS_OBJECT_STORE_CONFIG = string.Template(""" + + + + + + + + +""") + + +def check_container_active(container_name): + return subprocess.call([ + 'docker', + 'inspect', + '-f', + '{{.State.Running}}', + container_name + ]) == 0 + + +def start_irods(container_name): + if not check_container_active(container_name): + irods_start_args = [ + 'docker', + 'run', + '-p', + '1247:1247', + '-d', + '--name', + container_name, + 'kxk302/irods-server:0.1'] + subprocess.check_call(irods_start_args) + + # Sleep so integration test's iRODS instance is up and running + time.sleep(20) + + +def stop_irods(container_name): + if check_container_active(container_name): + subprocess.check_call(['docker', 'rm', '-f', container_name]) + + +class BaseObjectstoreUploadTest(UploadTestDatatypeDataTestCase): + + object_store_template = None + + @classmethod + def handle_galaxy_config_kwds(cls, config): + temp_directory = cls._test_driver.mkdtemp() + cls.object_stores_parent = temp_directory + cls.object_store_config_path = os.path.join(temp_directory, "object_store_conf.xml") + # This doesn't quite work yet, fails with extra_files_path + # config["metadata_strategy"] = "extended" + config["outpus_to_working_dir"] = True + config["retry_metadata_internally"] = False + config["object_store_store_by"] = "uuid" + with open(cls.object_store_config_path, "w") as f: + f.write(cls.object_store_template.safe_substitute(**cls.get_object_store_kwargs())) + + @classmethod + def get_object_store_kwargs(cls): + return {} + + +class IrodsUploadTestDatatypeDataTestCase(BaseObjectstoreUploadTest): + + object_store_template = IRODS_OBJECT_STORE_CONFIG + + @classmethod + def setUpClass(cls): + cls.container_name = "irods_integration_container" + start_irods(cls.container_name) + super(UploadTestDatatypeDataTestCase, cls).setUpClass() + + @classmethod + def tearDownClass(cls): + stop_irods(cls.container_name) + super(UploadTestDatatypeDataTestCase, cls).tearDownClass() + + @classmethod + def get_object_store_kwargs(cls): + return { + "temp_directory": cls.object_stores_parent, + "host": IRODS_OBJECT_STORE_HOST, + "port": IRODS_OBJECT_STORE_PORT, + "timeout": IRODS_OBJECT_STORE_TIMEOUT, + "username": IRODS_OBJECT_STORE_USERNAME, + "password": IRODS_OBJECT_STORE_PASSWORD, + "resource": IRODS_OBJECT_STORE_RESOURCE, + "zone": IRODS_OBJECT_STORE_ZONE + } + + +class UploadTestDosDiskAndDiskTestCase(BaseObjectstoreUploadTest): + + object_store_template = DISTRIBUTED_OBJECT_STORE_CONFIG + + @classmethod + def get_object_store_kwargs(cls): + return {'temp_directory': cls.object_stores_parent} + + +class UploadTestDosIrodsAndDiskTestCase(IrodsUploadTestDatatypeDataTestCase): + + object_store_template = DISTRIBUTED_IRODS_OBJECT_STORE_CONFIG + + +distributed_instance = integration_util.integration_module_instance(UploadTestDosDiskAndDiskTestCase) +irods_instance = integration_util.integration_module_instance(IrodsUploadTestDatatypeDataTestCase) +distributed_and_irods_instance = integration_util.integration_module_instance(UploadTestDosIrodsAndDiskTestCase) + + +@pytest.mark.parametrize('test_data', TEST_CASES.values(), ids=list(TEST_CASES.keys())) +def test_upload_datatype_dos_disk_and_disk(distributed_instance, test_data, temp_file): + upload_datatype_helper(distributed_instance, test_data, temp_file) + + +@pytest.mark.parametrize('test_data', TEST_CASES.values(), ids=list(TEST_CASES.keys())) +def test_upload_datatype_irods(irods_instance, test_data, temp_file): + upload_datatype_helper(irods_instance, test_data, temp_file) + + +@pytest.mark.parametrize('test_data', TEST_CASES.values(), ids=list(TEST_CASES.keys())) +def test_upload_datatype_dos_irods_and_disk(distributed_and_irods_instance, test_data, temp_file): + upload_datatype_helper(distributed_and_irods_instance, test_data, temp_file) diff --git a/test/integration/test_datatype_upload.py b/test/integration/test_datatype_upload.py index 0250d3d3a6d..c90ce258302 100644 --- a/test/integration/test_datatype_upload.py +++ b/test/integration/test_datatype_upload.py @@ -42,6 +42,8 @@ def collect_test_data(registry): class UploadTestDatatypeDataTestCase(BaseUploadContentConfigurationInstance): framework_tool_and_types = False datatypes_conf_override = DATATYPES_CONFIG + object_store_config = None + object_store_config_path = None instance = integration_util.integration_module_instance(UploadTestDatatypeDataTestCase) diff --git a/test/integration/test_datatype_upload_irods.py b/test/integration/test_datatype_upload_irods.py deleted file mode 100644 index 3001da5d3a1..00000000000 --- a/test/integration/test_datatype_upload_irods.py +++ /dev/null @@ -1,104 +0,0 @@ - -import os -import string -import subprocess -import time - -import pytest - -from galaxy_test.driver import integration_util -from .test_datatype_upload import ( - TEST_CASES, - upload_datatype_helper, - UploadTestDatatypeDataTestCase -) - -SCRIPT_DIRECTORY = os.path.abspath(os.path.dirname(__file__)) -# Run test for only the first 10 test files -IRODS_TEST_CASES = dict(list(TEST_CASES.items())[0:10]) -OBJECT_STORE_HOST = os.environ.get('GALAXY_INTEGRATION_IRODS_OBJECT_STORE_HOST', 'localhost') -OBJECT_STORE_PORT = int(os.environ.get('GALAXY_INTEGRATION_IRODS_OBJECT_STORE_PORT', 1247)) -OBJECT_STORE_TIMEOUT = int(os.environ.get('GALAXY_INTEGRATION_IRODS_OBJECT_STORE_TIMEOUT', 30)) -OBJECT_STORE_USERNAME = os.environ.get('GALAXY_INTEGRATION_IRODS_OBJECT_STORE_USERNAME', 'rods') -OBJECT_STORE_PASSWORD = os.environ.get('GALAXY_INTEGRATION_IRODS_OBJECT_STORE_PASSWORD', 'rods') -OBJECT_STORE_RESOURCE = os.environ.get('GALAXY_INTEGRATION_IRODS_OBJECT_STORE_RESOURCE', 'demoResc') -OBJECT_STORE_ZONE = os.environ.get('GALAXY_INTEGRATION_IRODS_OBJECT_STORE_ZONE', 'tempZone') -OBJECT_STORE_CONFIG = string.Template(""" - - - - - - - - - -""") - - -def start_irods(container_name): - irods_start_args = [ - 'docker', - 'run', - '-p', - '1247:1247', - '-d', - '--name', - container_name, - 'kxk302/irods-server:0.1'] - subprocess.check_call(irods_start_args) - - # Sleep so integration test's iRODS instance is up and running - time.sleep(20) - - -def stop_irods(container_name): - subprocess.check_call(['docker', 'rm', '-f', container_name]) - - -class UploadTestDatatypeDataTestCase(UploadTestDatatypeDataTestCase): - - @classmethod - def setUpClass(cls): - cls.container_name = "%s_container" % cls.__name__ - start_irods(cls.container_name) - super(UploadTestDatatypeDataTestCase, cls).setUpClass() - - @classmethod - def tearDownClass(cls): - stop_irods(cls.container_name) - super(UploadTestDatatypeDataTestCase, cls).tearDownClass() - - @classmethod - def handle_galaxy_config_kwds(cls, config): - temp_directory = cls._test_driver.mkdtemp() - cls.object_stores_parent = temp_directory - config_path = os.path.join(temp_directory, "object_store_conf.xml") - # This doesn't quite work yet, fails with extra_files_path - # config["metadata_strategy"] = "extended" - config["outpus_to_working_dir"] = True - config["retry_metadata_internally"] = False - with open(config_path, "w") as f: - f.write( - OBJECT_STORE_CONFIG.safe_substitute( - { - "temp_directory": temp_directory, - "host": OBJECT_STORE_HOST, - "port": OBJECT_STORE_PORT, - "timeout": OBJECT_STORE_TIMEOUT, - "username": OBJECT_STORE_USERNAME, - "password": OBJECT_STORE_PASSWORD, - "resource": OBJECT_STORE_RESOURCE, - "zone": OBJECT_STORE_ZONE, - } - ) - ) - config["object_store_config_file"] = config_path - - -instance = integration_util.integration_module_instance(UploadTestDatatypeDataTestCase) - - -@pytest.mark.parametrize('test_data', IRODS_TEST_CASES.values(), ids=list(IRODS_TEST_CASES.keys())) -def test_upload_datatype_irods(instance, test_data, temp_file): - upload_datatype_helper(instance, test_data, temp_file)