Merge pull request #9609 from kxk302/dev

Enabled various object store types to store data objects by uuid.
This commit is contained in:
Marius van den Beek
2020-04-22 17:43:23 +02:00
committed by GitHub
9 changed files with 235 additions and 165 deletions
+21 -36
View File
@@ -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
+5 -5
View File
@@ -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='')
+5 -5
View File
@@ -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='')
+5 -5
View File
@@ -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='')
+5 -5
View File
@@ -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, '')
+4 -5
View File
@@ -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]
@@ -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("""
<object_store type="distributed">
<backends>
<backend id="files1" type="disk" weight="1">
<files_dir path="${temp_directory}/database/files1"/>
<extra_dir type="temp" path="${temp_directory}/database/tmp1"/>
<extra_dir type="job_work" path="${temp_directory}/database/job_working_directory1"/>
</backend>
<backend id="files2" type="disk" weight="1">
<files_dir path="${temp_directory}/database/files2"/>
<extra_dir type="temp" path="${temp_directory}/database/tmp2"/>
<extra_dir type="job_work" path="${temp_directory}/database/job_working_directory2"/>
</backend>
</backends>
</object_store>
""")
DISTRIBUTED_IRODS_OBJECT_STORE_CONFIG = string.Template("""
<object_store type="distributed">
<backends>
<backend id="files1" type="disk" weight="1">
<files_dir path="${temp_directory}/database/files1"/>
<extra_dir type="temp" path="${temp_directory}/database/tmp1"/>
<extra_dir type="job_work" path="${temp_directory}/database/job_working_directory1"/>
</backend>
<backend id="files2" type="irods" weight="1">
<auth username="${username}" password="${password}"/>
<resource name="${resource}"/>
<zone name="${zone}"/>
<connection host="${host}" port="${port}" timeout="${timeout}"/>
<cache path="${temp_directory}/object_store_cache" size="1000"/>
<extra_dir type="job_work" path="${temp_directory}/job_working_directory_irods"/>
<extra_dir type="temp" path="${temp_directory}/tmp_irods"/>
</backend>
</backends>
</object_store>
""")
IRODS_OBJECT_STORE_CONFIG = string.Template("""<object_store type="irods">
<auth username="${username}" password="${password}"/>
<resource name="${resource}"/>
<zone name="${zone}"/>
<connection host="${host}" port="${port}" timeout="${timeout}"/>
<cache path="${temp_directory}/object_store_cache" size="1000"/>
<extra_dir type="job_work" path="${temp_directory}/job_working_directory_irods"/>
<extra_dir type="temp" path="${temp_directory}/tmp_irods"/>
</object_store>
""")
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)
+2
View File
@@ -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)
@@ -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("""
<object_store type="irods">
<auth username="${username}" password="${password}"/>
<resource name="${resource}"/>
<zone name="${zone}"/>
<connection host="${host}" port="${port}" timeout="${timeout}"/>
<cache path="${temp_directory}/object_store_cache" size="1000"/>
<extra_dir type="job_work" path="${temp_directory}/job_working_directory_irods"/>
<extra_dir type="temp" path="${temp_directory}/tmp_irods"/>
</object_store>
""")
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)