Merge pull request #7294 from hexylena/s3-upload-issue

[20.01] Fix S3 issue discovered during admin training
This commit is contained in:
Nicola Soranzo
2020-03-03 22:31:53 +01:00
committed by GitHub
3 changed files with 153 additions and 30 deletions
+3 -1
View File
@@ -82,7 +82,7 @@ class ObjectStore(object):
000/obj.id)
"""
def __init__(self, config, config_dict={}, **kwargs):
def __init__(self, config, config_dict=None, **kwargs):
"""
:type config: object
:param config: An object, most likely populated from
@@ -95,6 +95,8 @@ class ObjectStore(object):
parent directory those directories will be created.
* new_file_path -- Used to set the 'temp' extra_dir.
"""
if config_dict is None:
config_dict = {}
self.running = True
self.config = config
self.check_old_style = config.object_store_check_old_style
+37 -29
View File
@@ -125,7 +125,8 @@ class CloudConfigMixin(object):
'cache': {
'size': self.cache_size,
'path': self.staging_path,
}
},
'enable_cache_monitor': False,
}
@@ -138,7 +139,7 @@ class S3ObjectStore(ObjectStore, CloudConfigMixin):
store_type = 's3'
def __init__(self, config, config_dict):
super(S3ObjectStore, self).__init__(config)
super(S3ObjectStore, self).__init__(config, config_dict)
self.transfer_progress = 0
@@ -146,6 +147,7 @@ class S3ObjectStore(ObjectStore, CloudConfigMixin):
bucket_dict = config_dict['bucket']
connection_dict = config_dict.get('connection', {})
cache_dict = config_dict['cache']
self.enable_cache_monitor = config_dict.get('enable_cache_monitor', True)
self.access_key = auth_dict.get('access_key')
self.secret_key = auth_dict.get('secret_key')
@@ -167,9 +169,6 @@ class S3ObjectStore(ObjectStore, CloudConfigMixin):
(e['type'], e['path']) for e in config_dict.get('extra_dirs', []))
self.extra_dirs.update(extra_dirs)
log.debug("Object cache dir: %s", self.staging_path)
log.debug(" job work dir: %s", self.extra_dirs['job_work'])
self._initialize()
def _initialize(self):
@@ -187,9 +186,17 @@ class S3ObjectStore(ObjectStore, CloudConfigMixin):
'conn_path': self.conn_path}
self._configure_connection()
self.bucket = self._get_bucket(self.bucket)
self._bucket = self._get_bucket(self.bucket)
self.start_cache_monitor()
# Test if 'axel' is available for parallel download and pull the key into cache
if which('axel'):
self.use_axel = True
else:
self.use_axel = False
def start_cache_monitor(self):
# Clean cache only if value is set in galaxy.ini
if self.cache_size != -1:
if self.cache_size != -1 and self.enable_cache_monitor:
# Convert GBs to bytes for comparison
self.cache_size = self.cache_size * 1073741824
# Helper for interruptable sleep
@@ -197,11 +204,6 @@ class S3ObjectStore(ObjectStore, CloudConfigMixin):
self.cache_monitor_thread = threading.Thread(target=self.__cache_monitor)
self.cache_monitor_thread.start()
log.info("Cache cleaner manager started")
# Test if 'axel' is available for parallel download and pull the key into cache
if which('axel'):
self.use_axel = True
else:
self.use_axel = False
def _configure_connection(self):
log.debug("Configuring S3 Connection")
@@ -325,7 +327,7 @@ class S3ObjectStore(ObjectStore, 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)
@@ -334,7 +336,7 @@ class S3ObjectStore(ObjectStore, 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)
@@ -343,7 +345,7 @@ class S3ObjectStore(ObjectStore, 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):
@@ -354,7 +356,7 @@ class S3ObjectStore(ObjectStore, CloudConfigMixin):
def _get_size_in_s3(self, rel_path):
try:
key = self.bucket.get_key(rel_path)
key = self._bucket.get_key(rel_path)
if key:
return key.size
except S3ResponseError:
@@ -367,19 +369,17 @@ class S3ObjectStore(ObjectStore, CloudConfigMixin):
# A hackish way of testing if the rel_path is a folder vs a file
is_dir = rel_path[-1] == '/'
if is_dir:
keyresult = self.bucket.get_all_keys(prefix=rel_path)
keyresult = self._bucket.get_all_keys(prefix=rel_path)
if len(keyresult) > 0:
exists = True
else:
exists = False
else:
key = Key(self.bucket, rel_path)
key = Key(self._bucket, rel_path)
exists = key.exists()
except S3ResponseError:
log.exception("Trouble checking existence of S3 key '%s'", rel_path)
return False
if rel_path[0] == '/':
raise
return exists
def _in_cache(self, rel_path):
@@ -427,7 +427,7 @@ class S3ObjectStore(ObjectStore, CloudConfigMixin):
def _download(self, rel_path):
try:
log.debug("Pulling key '%s' into cache to %s", rel_path, self._get_cache_path(rel_path))
key = self.bucket.get_key(rel_path)
key = self._bucket.get_key(rel_path)
# Test if cache is large enough to hold the new file
if self.cache_size > 0 and key.size > self.cache_size:
log.critical("File %s is larger (%s) than the cache size (%s). Cannot download.",
@@ -446,7 +446,7 @@ class S3ObjectStore(ObjectStore, CloudConfigMixin):
key.get_contents_to_filename(self._get_cache_path(rel_path), cb=self._transfer_cb, num_cb=10)
return True
except S3ResponseError:
log.exception("Problem downloading key '%s' from S3 bucket '%s'", rel_path, self.bucket.name)
log.exception("Problem downloading key '%s' from S3 bucket '%s'", rel_path, self._bucket.name)
return False
def _push_to_os(self, rel_path, source_file=None, from_string=None):
@@ -460,7 +460,7 @@ class S3ObjectStore(ObjectStore, CloudConfigMixin):
try:
source_file = source_file if source_file else self._get_cache_path(rel_path)
if os.path.exists(source_file):
key = Key(self.bucket, rel_path)
key = Key(self._bucket, rel_path)
if os.path.getsize(source_file) == 0 and key.exists():
log.debug("Wanted to push file '%s' to S3 key '%s' but its size is 0; skipping.", source_file, rel_path)
return True
@@ -478,7 +478,7 @@ class S3ObjectStore(ObjectStore, CloudConfigMixin):
cb=self._transfer_cb,
num_cb=10)
else:
multipart_upload(self.s3server, self.bucket, key.name, source_file, mb_size)
multipart_upload(self.s3server, self._bucket, key.name, source_file, mb_size)
end_time = datetime.now()
log.debug("Pushed cache file '%s' to key '%s' (%s bytes transfered in %s sec)",
source_file, rel_path, os.path.getsize(source_file), end_time - start_time)
@@ -547,7 +547,7 @@ class S3ObjectStore(ObjectStore, 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 S3ObjectStore(ObjectStore, CloudConfigMixin):
# 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='')
@@ -609,7 +609,7 @@ class S3ObjectStore(ObjectStore, CloudConfigMixin):
# but requires iterating through each individual key in S3 and deleing it.
if entire_dir and extra_dir:
shutil.rmtree(self._get_cache_path(rel_path))
results = self.bucket.get_all_keys(prefix=rel_path)
results = self._bucket.get_all_keys(prefix=rel_path)
for key in results:
log.debug("Deleting key %s", key.name)
key.delete()
@@ -619,7 +619,7 @@ class S3ObjectStore(ObjectStore, CloudConfigMixin):
os.unlink(self._get_cache_path(rel_path))
# Delete from S3 as well
if self._key_exists(rel_path):
key = Key(self.bucket, rel_path)
key = Key(self._bucket, rel_path)
log.debug("Deleting key %s", key.name)
key.delete()
return True
@@ -707,7 +707,7 @@ class S3ObjectStore(ObjectStore, CloudConfigMixin):
if self.exists(obj, **kwargs):
rel_path = self._construct_path(obj, **kwargs)
try:
key = Key(self.bucket, rel_path)
key = Key(self._bucket, rel_path)
return key.generate_url(expires_in=86400) # 24hrs
except S3ResponseError:
log.exception("Trouble generating URL for dataset '%s'", rel_path)
@@ -716,6 +716,14 @@ class S3ObjectStore(ObjectStore, CloudConfigMixin):
def get_store_usage_percent(self):
return 0.0
def shutdown(self):
self.running = False
thread = getattr(self, 'cache_monitor_thread', None)
if thread:
log.debug("Shutting down thread")
self.sleeper.wake()
thread.join(5)
class SwiftObjectStore(S3ObjectStore):
"""
@@ -0,0 +1,113 @@
import os
import string
import subprocess
from galaxy_test.driver import integration_util
OBJECT_STORE_HOST = os.environ.get('GALAXY_INTEGRATION_OBJECT_STORE_HOST', '127.0.0.1')
OBJECT_STORE_PORT = int(os.environ.get('GALAXY_INTEGRATION_OBJECT_STORE_PORT', 9000))
OBJECT_STORE_ACCESS_KEY = os.environ.get('GALAXY_INTEGRATION_OBJECT_STORE_ACCESS_KEY', 'minioadmin')
OBJECT_STORE_SECRET_KEY = os.environ.get('GALAXY_INTEGRATION_OBJECT_STORE_SECRET_KEY', 'minioadmin')
OBJECT_STORE_CONFIG = string.Template("""
<object_store type="hierarchical" id="primary">
<backends>
<object_store id="swifty" type="swift" weight="1" order="0">
<auth access_key="${access_key}" secret_key="${secret_key}" />
<bucket name="galaxy" use_reduced_redundancy="False" max_chunk_size="250"/>
<connection host="${host}" port="${port}" is_secure="False" conn_path="" multipart="True"/>
<cache path="${temp_directory}/object_store_cache" size="1000" />
<extra_dir type="job_work" path="${temp_directory}/job_working_directory_swift"/>
<extra_dir type="temp" path="${temp_directory}/tmp_swift"/>
</object_store>
</backends>
</object_store>
""")
TEST_TOOL_IDS = [
"multi_output",
"multi_output_configured",
"multi_output_assign_primary",
"multi_output_recurse",
"tool_provided_metadata_1",
"tool_provided_metadata_2",
"tool_provided_metadata_3",
"tool_provided_metadata_4",
"tool_provided_metadata_5",
"tool_provided_metadata_6",
"tool_provided_metadata_7",
"tool_provided_metadata_8",
"tool_provided_metadata_9",
"tool_provided_metadata_10",
"tool_provided_metadata_11",
"tool_provided_metadata_12",
"composite_output",
"composite_output_tests",
"metadata",
"metadata_bam",
"output_format",
"output_auto_format",
]
def start_minio(container_name):
minio_start_args = [
'docker',
'run',
'-p',
'{port}:9000'.format(port=OBJECT_STORE_PORT),
'-d',
'--name',
container_name,
# '--rm',
'minio/minio:latest',
'server',
'/tmp/data']
subprocess.check_call(minio_start_args)
def stop_minio(container_name):
subprocess.check_call(['docker', 'rm', '-f', container_name])
@integration_util.skip_unless_docker()
class SwiftObjectStoreIntegrationTestCase(integration_util.IntegrationTestCase):
@classmethod
def setUpClass(cls):
cls.container_name = "%s_container" % cls.__name__
start_minio(cls.container_name)
super(SwiftObjectStoreIntegrationTestCase, cls).setUpClass()
@classmethod
def tearDownClass(cls):
stop_minio(cls.container_name)
super(SwiftObjectStoreIntegrationTestCase, 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")
config["object_store_store_by"] = "uuid"
# 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,
"access_key": OBJECT_STORE_ACCESS_KEY,
"secret_key": OBJECT_STORE_SECRET_KEY,
}
)
)
config["object_store_config_file"] = config_path
instance = integration_util.integration_module_instance(SwiftObjectStoreIntegrationTestCase)
test_tools = integration_util.integration_tool_runner(TEST_TOOL_IDS)