From 0516279b7c5a5b6c11118017517ec06e4e3745d6 Mon Sep 17 00:00:00 2001 From: vahid Date: Thu, 20 Jul 2017 22:06:12 -0700 Subject: [PATCH 01/26] Added CloudBridge to the requirements list. --- lib/galaxy/dependencies/conditional-requirements.txt | 1 + 1 file changed, 1 insertion(+) diff --git a/lib/galaxy/dependencies/conditional-requirements.txt b/lib/galaxy/dependencies/conditional-requirements.txt index 173a16fecb6..565ce7dd23f 100644 --- a/lib/galaxy/dependencies/conditional-requirements.txt +++ b/lib/galaxy/dependencies/conditional-requirements.txt @@ -13,6 +13,7 @@ graphitesend azure-storage==0.32.0 # PyRods not in PyPI python-ldap==2.4.27 +cloudbridge==0.3.1 # Synnefo / Pithos+ object store client kamaki From aaad00162f3ea84798031de798fe99fb0fafe05b Mon Sep 17 00:00:00 2001 From: vahid Date: Thu, 20 Jul 2017 23:00:34 -0700 Subject: [PATCH 02/26] Added CloudBridge to the dependencies check. --- lib/galaxy/dependencies/__init__.py | 3 +++ 1 file changed, 3 insertions(+) diff --git a/lib/galaxy/dependencies/__init__.py b/lib/galaxy/dependencies/__init__.py index 5d950c3fa16..94cd8099e0e 100644 --- a/lib/galaxy/dependencies/__init__.py +++ b/lib/galaxy/dependencies/__init__.py @@ -109,6 +109,9 @@ class ConditionalDependencies( object ): def check_azure_storage( self ): return 'azure_blob' in self.object_stores + def check_cloudbridge( self ): + return 's3' in self.object_stores + def check_kamaki(self): return 'pithos' in self.object_stores From 55c976cb85237bb0906379d692e212f05a327c2b Mon Sep 17 00:00:00 2001 From: vahid Date: Fri, 21 Jul 2017 22:23:07 -0700 Subject: [PATCH 03/26] (1) Added Cloud to ObjectStore; (2) updated CloudBridge dependencies. --- lib/galaxy/dependencies/__init__.py | 2 +- lib/galaxy/objectstore/cloud.py | 615 ++++++++++++++++++++++++++++ 2 files changed, 616 insertions(+), 1 deletion(-) create mode 100644 lib/galaxy/objectstore/cloud.py diff --git a/lib/galaxy/dependencies/__init__.py b/lib/galaxy/dependencies/__init__.py index 94cd8099e0e..10ebb8cdd20 100644 --- a/lib/galaxy/dependencies/__init__.py +++ b/lib/galaxy/dependencies/__init__.py @@ -110,7 +110,7 @@ class ConditionalDependencies( object ): return 'azure_blob' in self.object_stores def check_cloudbridge( self ): - return 's3' in self.object_stores + return 'cloud' in self.object_stores def check_kamaki(self): return 'pithos' in self.object_stores diff --git a/lib/galaxy/objectstore/cloud.py b/lib/galaxy/objectstore/cloud.py new file mode 100644 index 00000000000..809e92ad95a --- /dev/null +++ b/lib/galaxy/objectstore/cloud.py @@ -0,0 +1,615 @@ +""" +Object Store plugin for Cloud storage. +""" + +import logging +import multiprocessing +import os +import shutil +import subprocess +import threading +import time + +from datetime import datetime + +from galaxy.exceptions import ObjectInvalid, ObjectNotFound +from galaxy.util import ( + directory_hash_id, + safe_relpath, + string_as_bool, + umask_fix_perms, +) +from galaxy.util.sleeper import Sleeper + +from ..objectstore import convert_bytes, ObjectStore +from cloudbridge.cloud.factory import CloudProviderFactory, ProviderList + +try: + # Imports are done this way to allow objectstore code to be used outside of Galaxy. + import boto + + from boto.exception import S3ResponseError + from boto.s3.key import Key + from boto.s3.connection import S3Connection +except ImportError: + boto = None + +NO_BOTO_ERROR_MESSAGE = ("Cloud object store is configured, but no boto dependency available." + "Please install and properly configure boto or modify object store configuration.") + +log = logging.getLogger( __name__ ) +logging.getLogger('boto').setLevel(logging.INFO) # Otherwise boto is quite noisy + + +class Cloud( ObjectStore ): + """ + Object store that stores objects as items in an cloud storage. A local + cache exists that is used as an intermediate location for files between + Galaxy and the cloud storage. + """ + def __init__( self, config, config_xml ): + super( Cloud, self ).__init__( config ) + self.staging_path = self.config.file_path + self.transfer_progress = 0 + self._parse_config_xml( config_xml ) + self._configure_connection() + self.bucket = self._get_bucket(self.bucket) + # Clean cache only if value is set in galaxy.ini + if self.cache_size != -1: + # Convert GBs to bytes for comparison + self.cache_size = self.cache_size * 1073741824 + # Helper for interruptable sleep + self.sleeper = Sleeper() + 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 + try: + subprocess.call('axel') + self.use_axel = True + except OSError: + self.use_axel = False + + def _configure_connection( self ): + log.debug("Configuring AWS-S3 Connection") + aws_config = {'aws_access_key': self.access_key, + 'aws_secret_key': self.secret_key} + self.conn = CloudProviderFactory().create_provider(ProviderList.AWS, aws_config) + + def _parse_config_xml(self, config_xml): + try: + a_xml = config_xml.findall('auth')[0] + self.access_key = a_xml.get('access_key') + self.secret_key = a_xml.get('secret_key') + b_xml = config_xml.findall('bucket')[0] + self.bucket = b_xml.get('name') + self.use_rr = string_as_bool(b_xml.get('use_reduced_redundancy', "False")) + self.max_chunk_size = int(b_xml.get('max_chunk_size', 250)) + cn_xml = config_xml.findall('connection') + if not cn_xml: + cn_xml = {} + else: + cn_xml = cn_xml[0] + self.host = cn_xml.get('host', None) + self.port = int(cn_xml.get('port', 6000)) + self.multipart = string_as_bool(cn_xml.get('multipart', 'True')) + self.is_secure = string_as_bool(cn_xml.get('is_secure', 'True')) + self.conn_path = cn_xml.get('conn_path', '/') + c_xml = config_xml.findall('cache')[0] + self.cache_size = float(c_xml.get('size', -1)) + self.staging_path = c_xml.get('path', self.config.object_store_cache_path) + + for d_xml in config_xml.findall('extra_dir'): + self.extra_dirs[d_xml.get('type')] = d_xml.get('path') + + log.debug("Object cache dir: %s", self.staging_path) + log.debug(" job work dir: %s", self.extra_dirs['job_work']) + + except Exception: + # Toss it back up after logging, we can't continue loading at this point. + log.exception("Malformed ObjectStore Configuration XML -- unable to continue") + raise + + def __cache_monitor(self): + time.sleep(2) # Wait for things to load before starting the monitor + while self.running: + total_size = 0 + # Is this going to be too expensive of an operation to be done frequently? + file_list = [] + for dirpath, _, filenames in os.walk(self.staging_path): + for filename in filenames: + filepath = os.path.join(dirpath, filename) + file_size = os.path.getsize(filepath) + total_size += file_size + # Get the time given file was last accessed + last_access_time = time.localtime(os.stat(filepath)[7]) + # Compose a tuple of the access time and the file path + file_tuple = last_access_time, filepath, file_size + file_list.append(file_tuple) + # Sort the file list (based on access time) + file_list.sort() + # Initiate cleaning once within 10% of the defined cache size? + cache_limit = self.cache_size * 0.9 + if total_size > cache_limit: + log.info("Initiating cache cleaning: current cache size: %s; clean until smaller than: %s", + convert_bytes(total_size), convert_bytes(cache_limit)) + # How much to delete? If simply deleting up to the cache-10% limit, + # is likely to be deleting frequently and may run the risk of hitting + # the limit - maybe delete additional #%? + # For now, delete enough to leave at least 10% of the total cache free + delete_this_much = total_size - cache_limit + self.__clean_cache(file_list, delete_this_much) + self.sleeper.sleep(30) # Test cache size every 30 seconds? + + def __clean_cache(self, file_list, delete_this_much): + """ Keep deleting files from the file_list until the size of the deleted + files is greater than the value in delete_this_much parameter. + + :type file_list: list + :param file_list: List of candidate files that can be deleted. This method + will start deleting files from the beginning of the list so the list + should be sorted accordingly. The list must contains 3-element tuples, + positioned as follows: position 0 holds file last accessed timestamp + (as time.struct_time), position 1 holds file path, and position 2 has + file size (e.g., (, /mnt/data/dataset_1.dat), 472394) + + :type delete_this_much: int + :param delete_this_much: Total size of files, in bytes, that should be deleted. + """ + # Keep deleting datasets from file_list until deleted_amount does not + # exceed delete_this_much; start deleting from the front of the file list, + # which assumes the oldest files come first on the list. + deleted_amount = 0 + for entry in enumerate(file_list): + if deleted_amount < delete_this_much: + deleted_amount += entry[2] + os.remove(entry[1]) + # Debugging code for printing deleted files' stats + # folder, file_name = os.path.split(f[1]) + # file_date = time.strftime("%m/%d/%y %H:%M:%S", f[0]) + # log.debug("%s. %-25s %s, size %s (deleted %s/%s)" \ + # % (i, file_name, convert_bytes(f[2]), file_date, \ + # convert_bytes(deleted_amount), convert_bytes(delete_this_much))) + else: + log.debug("Cache cleaning done. Total space freed: %s", convert_bytes(deleted_amount)) + return + + def _get_bucket(self, bucket_name): + """ Sometimes a handle to a bucket is not established right away so try + it a few times. Raise error if connection is not established. """ + for i in range(5): + try: + bucket = self.conn.object_store.get(bucket_name) + if bucket is None: + log.debug("Bucket not found, creating a bucket with handle '%s'", bucket_name) + bucket = self.conn.object_store.create(bucket_name) + log.debug("Using cloud object store with bucket '%s'", bucket.name) + return bucket + except S3ResponseError: + log.exception("Could not get bucket '%s', attempt %s/5", bucket_name, i + 1) + time.sleep(2) + # All the attempts have been exhausted and connection was not established, + # raise error + raise S3ResponseError + + def _fix_permissions(self, rel_path): + """ Set permissions on rel_path""" + for basedir, _, files in os.walk(rel_path): + umask_fix_perms(basedir, self.config.umask, 0o777, self.config.gid) + for filename in files: + path = os.path.join(basedir, filename) + # Ignore symlinks + if os.path.islink(path): + continue + umask_fix_perms(path, self.config.umask, 0o666, self.config.gid) + + def _construct_path(self, obj, base_dir=None, dir_only=None, extra_dir=None, extra_dir_at_root=False, alt_name=None, + obj_dir=False, **kwargs): + # extra_dir should never be constructed from provided data but just + # make sure there are no shenannigans afoot + if extra_dir and extra_dir != os.path.normpath(extra_dir): + log.warning('extra_dir is not normalized: %s', extra_dir) + raise ObjectInvalid("The requested object is invalid") + # ensure that any parent directory references in alt_name would not + # result in a path not contained in the directory path constructed here + if alt_name: + if not safe_relpath(alt_name): + log.warning('alt_name would locate path outside dir: %s', alt_name) + raise ObjectInvalid("The requested object is invalid") + # 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)) + if extra_dir is not None: + if extra_dir_at_root: + rel_path = os.path.join(extra_dir, rel_path) + else: + rel_path = os.path.join(rel_path, extra_dir) + + # for JOB_WORK directory + if obj_dir: + rel_path = os.path.join(rel_path, str(obj.id)) + if base_dir: + base = self.extra_dirs.get(base_dir) + return os.path.join(base, rel_path) + + # S3 folders are marked by having trailing '/' so add it now + 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) + return rel_path + + def _get_cache_path(self, rel_path): + return os.path.abspath(os.path.join(self.staging_path, rel_path)) + + def _get_transfer_progress(self): + return self.transfer_progress + + def _get_size_in_cloud(self, rel_path): + try: + obj = self.bucket.get(rel_path) + if obj: + return obj.size + except S3ResponseError: + log.exception("Could not get size of key '%s' from S3", rel_path) + return -1 + + def _key_exists(self, rel_path): + exists = False + try: + # 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.list(prefix=rel_path) + if len(keyresult) > 0: + exists = True + else: + exists = False + else: + exists = self.bucket.exists(rel_path) + 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): + """ Check if the given dataset is in the local cache and return True if so. """ + # log.debug("------ Checking cache for rel_path %s" % rel_path) + cache_path = self._get_cache_path(rel_path) + return os.path.exists(cache_path) + # TODO: Part of checking if a file is in cache should be to ensure the + # size of the cached file matches that on S3. Once the upload tool explicitly + # creates, this check sould be implemented- in the mean time, it's not + # looking likely to be implementable reliably. + # if os.path.exists(cache_path): + # # print "***1 %s exists" % cache_path + # if self._key_exists(rel_path): + # # print "***2 %s exists in S3" % rel_path + # # Make sure the size in cache is available in its entirety + # # print "File '%s' cache size: %s, S3 size: %s" % (cache_path, os.path.getsize(cache_path), self._get_size_in_cloud(rel_path)) + # if os.path.getsize(cache_path) == self._get_size_in_cloud(rel_path): + # # print "***2.1 %s exists in S3 and the size is the same as in cache (in_cache=True)" % rel_path + # exists = True + # else: + # # print "***2.2 %s exists but differs in size from cache (in_cache=False)" % cache_path + # exists = False + # else: + # # Although not perfect decision making, this most likely means + # # that the file is currently being uploaded + # # print "***3 %s found in cache but not in S3 (in_cache=True)" % cache_path + # exists = True + # else: + # return False + + def _pull_into_cache(self, rel_path): + # Ensure the cache directory structure exists (e.g., dataset_#_files/) + rel_path_dir = os.path.dirname(rel_path) + if not os.path.exists(self._get_cache_path(rel_path_dir)): + os.makedirs(self._get_cache_path(rel_path_dir)) + # Now pull in the file + file_ok = self._download(rel_path) + self._fix_permissions(self._get_cache_path(rel_path_dir)) + return file_ok + + def _transfer_cb(self, complete, total): + self.transfer_progress += 10 + + 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(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.", + rel_path, key.size, self.cache_size) + return False + if self.use_axel: + log.debug("Parallel pulled key '%s' into cache to %s", rel_path, self._get_cache_path(rel_path)) + ncores = multiprocessing.cpu_count() + url = key.generate_url(7200) + ret_code = subprocess.call("axel -a -n %s '%s'" % (ncores, url)) + if ret_code == 0: + return True + else: + log.debug("Pulled key '%s' into cache to %s", rel_path, self._get_cache_path(rel_path)) + self.transfer_progress = 0 # Reset transfer progress counter + with open(self._get_cache_path(rel_path), "w+") as downloaded_file_handle: + key.save_content(downloaded_file_handle) + return True + except S3ResponseError: + 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): + """ + Push the file pointed to by ``rel_path`` to the object store naming the key + ``rel_path``. If ``source_file`` is provided, push that file instead while + still using ``rel_path`` as the key name. + If ``from_string`` is provided, set contents of the file to the value of + the string. + """ + try: + source_file = source_file if source_file else self._get_cache_path(rel_path) + if os.path.exists(source_file): + if os.path.getsize(source_file) == 0 and self.bucket.exists(rel_path): + log.debug("Wanted to push file '%s' to S3 key '%s' but its size is 0; skipping.", source_file, + rel_path) + return True + # FIXME: don't need to differenciate between uploading from a string or file, + # because CloudBridge handles this internally. + if from_string: + if not self.bucket.get(rel_path): + created_obj = self.bucket.create_object(rel_path) + created_obj.upload(source_file) + else: + self.bucket.get(rel_path).upload(source_file) + log.debug("Pushed data from string '%s' to key '%s'", from_string, rel_path) + else: + start_time = datetime.now() + log.debug("Pushing cache file '%s' of size %s bytes to key '%s'", source_file, + os.path.getsize(source_file), rel_path) + self.transfer_progress = 0 # Reset transfer progress counter + if not self.bucket.get(rel_path): + created_obj = self.bucket.create_object(rel_path) + created_obj.upload(source_file) + else: + self.bucket.get(rel_path).upload(source_file) + 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) + return True + else: + log.error("Tried updating key '%s' from source file '%s', but source file does not exist.", + rel_path, source_file) + except S3ResponseError: + log.exception("Trouble pushing S3 key '%s' from file '%s'", rel_path, source_file) + return False + + def file_ready(self, obj, **kwargs): + """ + A helper method that checks if a file corresponding to a dataset is + ready and available to be used. Return ``True`` if so, ``False`` otherwise. + """ + rel_path = self._construct_path(obj, **kwargs) + # Make sure the size in cache is available in its entirety + if self._in_cache(rel_path): + if os.path.getsize(self._get_cache_path(rel_path)) == self._get_size_in_cloud(rel_path): + return True + log.debug("Waiting for dataset %s to transfer from OS: %s/%s", rel_path, + os.path.getsize(self._get_cache_path(rel_path)), self._get_size_in_cloud(rel_path)) + return False + + def exists(self, obj, **kwargs): + in_cache = False + rel_path = self._construct_path(obj, **kwargs) + + # Check cache + if self._in_cache(rel_path): + in_cache = True + # Check cloud + in_cloud = self._key_exists(rel_path) + # log.debug("~~~~~~ File '%s' exists in cache: %s; in s3: %s" % (rel_path, in_cache, in_s3)) + # dir_only does not get synced so shortcut the decision + dir_only = kwargs.get('dir_only', False) + base_dir = kwargs.get('base_dir', None) + if dir_only: + if in_cache or in_cloud: + return True + # for JOB_WORK directory + elif base_dir: + if not os.path.exists(rel_path): + os.makedirs(rel_path) + return True + else: + return False + + # TODO: Sync should probably not be done here. Add this to an async upload stack? + if in_cache and not in_cloud: + self._push_to_os(rel_path, source_file=self._get_cache_path(rel_path)) + return True + elif in_cloud: + return True + else: + return False + + def create(self, obj, **kwargs): + if not self.exists(obj, **kwargs): + + # Pull out locally used fields + extra_dir = kwargs.get('extra_dir', None) + extra_dir_at_root = kwargs.get('extra_dir_at_root', False) + dir_only = kwargs.get('dir_only', False) + alt_name = kwargs.get('alt_name', None) + + # Construct hashed path + rel_path = os.path.join(*directory_hash_id(obj.id)) + + # Optionally append extra_dir + if extra_dir is not None: + if extra_dir_at_root: + rel_path = os.path.join(extra_dir, rel_path) + else: + rel_path = os.path.join(rel_path, extra_dir) + + # Create given directory in cache + cache_dir = os.path.join(self.staging_path, rel_path) + if not os.path.exists(cache_dir): + os.makedirs(cache_dir) + + # Although not really necessary to create S3 folders (because S3 has + # flat namespace), do so for consistency with the regular file system + # S3 folders are marked by having trailing '/' so add it now + # s3_dir = '%s/' % rel_path + # 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) + open(os.path.join(self.staging_path, rel_path), 'w').close() + self._push_to_os(rel_path, from_string='') + + def empty(self, obj, **kwargs): + if self.exists(obj, **kwargs): + return bool(self.size(obj, **kwargs) > 0) + else: + raise ObjectNotFound( 'objectstore.empty, object does not exist: %s, kwargs: %s' + % ( str( obj ), str( kwargs ) ) ) + + def size(self, obj, **kwargs): + rel_path = self._construct_path(obj, **kwargs) + if self._in_cache(rel_path): + try: + return os.path.getsize(self._get_cache_path(rel_path)) + except OSError as ex: + log.info("Could not get size of file '%s' in local cache, will try cloud. Error: %s", rel_path, ex) + elif self.exists(obj, **kwargs): + return self._get_size_in_cloud(rel_path) + log.warning("Did not find dataset '%s', returning 0 for size", rel_path) + return 0 + + def delete(self, obj, entire_dir=False, **kwargs): + rel_path = self._construct_path(obj, **kwargs) + extra_dir = kwargs.get('extra_dir', None) + base_dir = kwargs.get('base_dir', None) + dir_only = kwargs.get('dir_only', False) + obj_dir = kwargs.get('obj_dir', False) + try: + # Remove temparory data in JOB_WORK directory + if base_dir and dir_only and obj_dir: + shutil.rmtree(os.path.abspath(rel_path)) + return True + + # For the case of extra_files, because we don't have a reference to + # individual files/keys we need to remove the entire directory structure + # with all the files in it. This is easy for the local file system, + # 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.list(prefix=rel_path) + for key in results: + log.debug("Deleting key %s", key.name) + key.delete() + return True + else: + # Delete from cache first + os.unlink(self._get_cache_path(rel_path)) + # Delete from S3 as well + if self._key_exists(rel_path): + key = self.bucket.get(rel_path) + log.debug("Deleting key %s", key.name) + key.delete() + return True + except S3ResponseError: + log.exception("Could not delete key '%s' from cloud", rel_path) + except OSError: + log.exception('%s delete error', self.get_filename(obj, **kwargs)) + return False + + def get_data(self, obj, start=0, count=-1, **kwargs): + rel_path = self._construct_path(obj, **kwargs) + # Check cache first and get file if not there + if not self._in_cache(rel_path): + self._pull_into_cache(rel_path) + # Read the file content from cache + data_file = open(self._get_cache_path(rel_path), 'r') + data_file.seek(start) + content = data_file.read(count) + data_file.close() + return content + + def get_filename(self, obj, **kwargs): + base_dir = kwargs.get('base_dir', None) + dir_only = kwargs.get('dir_only', False) + obj_dir = kwargs.get('obj_dir', False) + rel_path = self._construct_path(obj, **kwargs) + + # for JOB_WORK directory + if base_dir and dir_only and obj_dir: + return os.path.abspath(rel_path) + + cache_path = self._get_cache_path(rel_path) + # S3 does not recognize directories as files so cannot check if those exist. + # So, if checking dir only, ensure given dir exists in cache and return + # the expected cache path. + # dir_only = kwargs.get('dir_only', False) + # if dir_only: + # if not os.path.exists(cache_path): + # os.makedirs(cache_path) + # return cache_path + # Check if the file exists in the cache first + if self._in_cache(rel_path): + return cache_path + # Check if the file exists in persistent storage and, if it does, pull it into cache + elif self.exists(obj, **kwargs): + if dir_only: # Directories do not get pulled into cache + return cache_path + else: + if self._pull_into_cache(rel_path): + return cache_path + # For the case of retrieving a directory only, return the expected path + # even if it does not exist. + # if dir_only: + # return cache_path + raise ObjectNotFound( 'objectstore.get_filename, no cache_path: %s, kwargs: %s' + % ( str( obj ), str( kwargs ) ) ) + # return cache_path # Until the upload tool does not explicitly create the dataset, return expected path + + def update_from_file(self, obj, file_name=None, create=False, **kwargs): + if create: + self.create(obj, **kwargs) + if self.exists(obj, **kwargs): + rel_path = self._construct_path(obj, **kwargs) + # Chose whether to use the dataset file itself or an alternate file + if file_name: + source_file = os.path.abspath(file_name) + # Copy into cache + cache_file = self._get_cache_path(rel_path) + try: + if source_file != cache_file: + # FIXME? Should this be a `move`? + shutil.copy2(source_file, cache_file) + self._fix_permissions(cache_file) + except OSError: + log.exception("Trouble copying source file '%s' to cache '%s'", source_file, cache_file) + else: + source_file = self._get_cache_path(rel_path) + # Update the file on cloud + self._push_to_os(rel_path, source_file) + else: + raise ObjectNotFound( 'objectstore.update_from_file, object does not exist: %s, kwargs: %s' + % ( str( obj ), str( kwargs ) ) ) + + def get_object_url(self, obj, **kwargs): + if self.exists(obj, **kwargs): + rel_path = self._construct_path(obj, **kwargs) + try: + key = self.bucket.get(rel_path) + return key.generate_url(expires_in=86400) # 24hrs + except S3ResponseError: + log.exception("Trouble generating URL for dataset '%s'", rel_path) + return None + + def get_store_usage_percent(self): + return 0.0 From e4c00dd695bf23a506ffeec15f98a6f6d1a05312 Mon Sep 17 00:00:00 2001 From: vahid Date: Tue, 25 Jul 2017 18:42:57 -0700 Subject: [PATCH 04/26] - Updated `Cloud` to account for changes in `CloudBridge` interface. - Added `Cloud` to objectstore import. --- lib/galaxy/objectstore/__init__.py | 3 +++ lib/galaxy/objectstore/cloud.py | 4 ++-- 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/lib/galaxy/objectstore/__init__.py b/lib/galaxy/objectstore/__init__.py index 30fc086c3cc..50b4a7937bb 100644 --- a/lib/galaxy/objectstore/__init__.py +++ b/lib/galaxy/objectstore/__init__.py @@ -722,6 +722,9 @@ def build_object_store_from_config(config, fsmon=False, config_xml=None): elif store == 's3': from .s3 import S3ObjectStore return S3ObjectStore(config=config, config_xml=config_xml) + elif store == 'cloud': + from .cloud import Cloud + return Cloud(config=config, config_xml=config_xml) elif store == 'swift': from .s3 import SwiftObjectStore return SwiftObjectStore(config=config, config_xml=config_xml) diff --git a/lib/galaxy/objectstore/cloud.py b/lib/galaxy/objectstore/cloud.py index 809e92ad95a..06136f89385 100644 --- a/lib/galaxy/objectstore/cloud.py +++ b/lib/galaxy/objectstore/cloud.py @@ -267,7 +267,7 @@ class Cloud( ObjectStore ): else: exists = False else: - exists = self.bucket.exists(rel_path) + exists = True if self.bucket.get(rel_path) is not None else False except S3ResponseError: log.exception("Trouble checking existence of S3 key '%s'", rel_path) return False @@ -354,7 +354,7 @@ class Cloud( ObjectStore ): try: source_file = source_file if source_file else self._get_cache_path(rel_path) if os.path.exists(source_file): - if os.path.getsize(source_file) == 0 and self.bucket.exists(rel_path): + if os.path.getsize(source_file) == 0 and (self.bucket.get(rel_path) is not None): log.debug("Wanted to push file '%s' to S3 key '%s' but its size is 0; skipping.", source_file, rel_path) return True From ab34fc9bd92041065dd8f140194c56e0f7b2954f Mon Sep 17 00:00:00 2001 From: vahid Date: Wed, 26 Jul 2017 15:24:08 -0700 Subject: [PATCH 05/26] Updated `Babel` and `requests` packages version to meet the minimum requirements of `cloudbridge` package. --- lib/galaxy/dependencies/pinned-requirements.txt | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/lib/galaxy/dependencies/pinned-requirements.txt b/lib/galaxy/dependencies/pinned-requirements.txt index 5cb2c8a1d9c..c5111a5db92 100644 --- a/lib/galaxy/dependencies/pinned-requirements.txt +++ b/lib/galaxy/dependencies/pinned-requirements.txt @@ -22,7 +22,7 @@ WebOb==1.4.1 WebHelpers==1.3 Mako==1.0.2 pytz==2015.4 -Babel==2.1.1 +Babel==2.4.0 Beaker==1.7.0 dictobj==0.3.1 nose==1.3.7 @@ -39,7 +39,7 @@ Markdown==2.6.3 # BioBlend and dependencies bioblend==0.7.0 boto==2.38.0 -requests==2.8.1 +requests==2.10.0 requests-toolbelt==0.4.0 # kombu and dependencies From e8541a71a4fdd7f975b3e8d0c6b054bdc33480ab Mon Sep 17 00:00:00 2001 From: vahid Date: Wed, 26 Jul 2017 16:07:42 -0700 Subject: [PATCH 06/26] Reverted the changes on package versions, as these changes are requests via a different PR. --- lib/galaxy/dependencies/pinned-requirements.txt | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/lib/galaxy/dependencies/pinned-requirements.txt b/lib/galaxy/dependencies/pinned-requirements.txt index c5111a5db92..5cb2c8a1d9c 100644 --- a/lib/galaxy/dependencies/pinned-requirements.txt +++ b/lib/galaxy/dependencies/pinned-requirements.txt @@ -22,7 +22,7 @@ WebOb==1.4.1 WebHelpers==1.3 Mako==1.0.2 pytz==2015.4 -Babel==2.4.0 +Babel==2.1.1 Beaker==1.7.0 dictobj==0.3.1 nose==1.3.7 @@ -39,7 +39,7 @@ Markdown==2.6.3 # BioBlend and dependencies bioblend==0.7.0 boto==2.38.0 -requests==2.10.0 +requests==2.8.1 requests-toolbelt==0.4.0 # kombu and dependencies From 7ad7d881046ddd20276998d6e1fca495d1ac9f50 Mon Sep 17 00:00:00 2001 From: vahid Date: Fri, 28 Jul 2017 09:26:35 -0700 Subject: [PATCH 07/26] Changed boto import style. --- lib/galaxy/objectstore/cloud.py | 4 ---- 1 file changed, 4 deletions(-) diff --git a/lib/galaxy/objectstore/cloud.py b/lib/galaxy/objectstore/cloud.py index 06136f89385..d72ae0d589b 100644 --- a/lib/galaxy/objectstore/cloud.py +++ b/lib/galaxy/objectstore/cloud.py @@ -25,10 +25,6 @@ from ..objectstore import convert_bytes, ObjectStore from cloudbridge.cloud.factory import CloudProviderFactory, ProviderList try: - # Imports are done this way to allow objectstore code to be used outside of Galaxy. - import boto - - from boto.exception import S3ResponseError from boto.s3.key import Key from boto.s3.connection import S3Connection except ImportError: From 9c57adb90481b871cc6f98f58227caca4e414467 Mon Sep 17 00:00:00 2001 From: vahid Date: Wed, 2 Aug 2017 10:05:14 -0700 Subject: [PATCH 08/26] Updated CloudBridge to its current latest version. --- lib/galaxy/dependencies/conditional-requirements.txt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/galaxy/dependencies/conditional-requirements.txt b/lib/galaxy/dependencies/conditional-requirements.txt index 565ce7dd23f..466b2e9ab2e 100644 --- a/lib/galaxy/dependencies/conditional-requirements.txt +++ b/lib/galaxy/dependencies/conditional-requirements.txt @@ -13,7 +13,7 @@ graphitesend azure-storage==0.32.0 # PyRods not in PyPI python-ldap==2.4.27 -cloudbridge==0.3.1 +cloudbridge==0.3.2 # Synnefo / Pithos+ object store client kamaki From 2589c745ae7f88ada7a08265f1238ae00274425e Mon Sep 17 00:00:00 2001 From: vahid Date: Wed, 2 Aug 2017 11:29:25 -0700 Subject: [PATCH 09/26] Commented on boto import, and removed a `use_reduced_redundancy` config. --- lib/galaxy/objectstore/cloud.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/lib/galaxy/objectstore/cloud.py b/lib/galaxy/objectstore/cloud.py index d72ae0d589b..7254f076e5c 100644 --- a/lib/galaxy/objectstore/cloud.py +++ b/lib/galaxy/objectstore/cloud.py @@ -24,7 +24,9 @@ from galaxy.util.sleeper import Sleeper from ..objectstore import convert_bytes, ObjectStore from cloudbridge.cloud.factory import CloudProviderFactory, ProviderList +# boto is only used to handle exceptions; it will be removed once CloudBridge wraps and throws proper exceptions. try: + # Imports are done this way to allow objectstore code to be used outside of Galaxy. from boto.s3.key import Key from boto.s3.connection import S3Connection except ImportError: @@ -79,7 +81,6 @@ class Cloud( ObjectStore ): self.secret_key = a_xml.get('secret_key') b_xml = config_xml.findall('bucket')[0] self.bucket = b_xml.get('name') - self.use_rr = string_as_bool(b_xml.get('use_reduced_redundancy', "False")) self.max_chunk_size = int(b_xml.get('max_chunk_size', 250)) cn_xml = config_xml.findall('connection') if not cn_xml: From 89be432770d3cf9442f9c003ead8c6ea6a16f395 Mon Sep 17 00:00:00 2001 From: vahid Date: Wed, 2 Aug 2017 11:54:13 -0700 Subject: [PATCH 10/26] Removed boto import, and replaced all S3-specific exception catches with a generic exception--a temporary solution till CloudBridge wraps the exceptions properly. --- lib/galaxy/objectstore/cloud.py | 24 ++++++++---------------- 1 file changed, 8 insertions(+), 16 deletions(-) diff --git a/lib/galaxy/objectstore/cloud.py b/lib/galaxy/objectstore/cloud.py index 7254f076e5c..def78bbe62d 100644 --- a/lib/galaxy/objectstore/cloud.py +++ b/lib/galaxy/objectstore/cloud.py @@ -24,14 +24,6 @@ from galaxy.util.sleeper import Sleeper from ..objectstore import convert_bytes, ObjectStore from cloudbridge.cloud.factory import CloudProviderFactory, ProviderList -# boto is only used to handle exceptions; it will be removed once CloudBridge wraps and throws proper exceptions. -try: - # Imports are done this way to allow objectstore code to be used outside of Galaxy. - from boto.s3.key import Key - from boto.s3.connection import S3Connection -except ImportError: - boto = None - NO_BOTO_ERROR_MESSAGE = ("Cloud object store is configured, but no boto dependency available." "Please install and properly configure boto or modify object store configuration.") @@ -182,12 +174,12 @@ class Cloud( ObjectStore ): bucket = self.conn.object_store.create(bucket_name) log.debug("Using cloud object store with bucket '%s'", bucket.name) return bucket - except S3ResponseError: + except Exception: log.exception("Could not get bucket '%s', attempt %s/5", bucket_name, i + 1) time.sleep(2) # All the attempts have been exhausted and connection was not established, # raise error - raise S3ResponseError + raise Exception def _fix_permissions(self, rel_path): """ Set permissions on rel_path""" @@ -248,7 +240,7 @@ class Cloud( ObjectStore ): obj = self.bucket.get(rel_path) if obj: return obj.size - except S3ResponseError: + except Exception: log.exception("Could not get size of key '%s' from S3", rel_path) return -1 @@ -265,7 +257,7 @@ class Cloud( ObjectStore ): exists = False else: exists = True if self.bucket.get(rel_path) is not None else False - except S3ResponseError: + except Exception: log.exception("Trouble checking existence of S3 key '%s'", rel_path) return False if rel_path[0] == '/': @@ -336,7 +328,7 @@ class Cloud( ObjectStore ): with open(self._get_cache_path(rel_path), "w+") as downloaded_file_handle: key.save_content(downloaded_file_handle) return True - except S3ResponseError: + except Exception: log.exception("Problem downloading key '%s' from S3 bucket '%s'", rel_path, self.bucket.name) return False @@ -381,7 +373,7 @@ class Cloud( ObjectStore ): else: log.error("Tried updating key '%s' from source file '%s', but source file does not exist.", rel_path, source_file) - except S3ResponseError: + except Exception: log.exception("Trouble pushing S3 key '%s' from file '%s'", rel_path, source_file) return False @@ -518,7 +510,7 @@ class Cloud( ObjectStore ): log.debug("Deleting key %s", key.name) key.delete() return True - except S3ResponseError: + except Exception: log.exception("Could not delete key '%s' from cloud", rel_path) except OSError: log.exception('%s delete error', self.get_filename(obj, **kwargs)) @@ -604,7 +596,7 @@ class Cloud( ObjectStore ): try: key = self.bucket.get(rel_path) return key.generate_url(expires_in=86400) # 24hrs - except S3ResponseError: + except Exception: log.exception("Trouble generating URL for dataset '%s'", rel_path) return None From fafb019f4be3b33ad9cab4bfd2ded6788b419592 Mon Sep 17 00:00:00 2001 From: vahid Date: Wed, 2 Aug 2017 12:05:45 -0700 Subject: [PATCH 11/26] Consolidated two upload methods, i.e., upload from file and string, because CloudBridge internally maps to appropriate functions depending on the input type. --- lib/galaxy/objectstore/cloud.py | 32 +++++++++++--------------------- 1 file changed, 11 insertions(+), 21 deletions(-) diff --git a/lib/galaxy/objectstore/cloud.py b/lib/galaxy/objectstore/cloud.py index def78bbe62d..dac5dd6baa6 100644 --- a/lib/galaxy/objectstore/cloud.py +++ b/lib/galaxy/objectstore/cloud.py @@ -347,28 +347,18 @@ class Cloud( ObjectStore ): log.debug("Wanted to push file '%s' to S3 key '%s' but its size is 0; skipping.", source_file, rel_path) return True - # FIXME: don't need to differenciate between uploading from a string or file, - # because CloudBridge handles this internally. - if from_string: - if not self.bucket.get(rel_path): - created_obj = self.bucket.create_object(rel_path) - created_obj.upload(source_file) - else: - self.bucket.get(rel_path).upload(source_file) - log.debug("Pushed data from string '%s' to key '%s'", from_string, rel_path) + start_time = datetime.now() + log.debug("Pushing cache file '%s' of size %s bytes to key '%s'", source_file, + os.path.getsize(source_file), rel_path) + self.transfer_progress = 0 # Reset transfer progress counter + if not self.bucket.get(rel_path): + created_obj = self.bucket.create_object(rel_path) + created_obj.upload(source_file) else: - start_time = datetime.now() - log.debug("Pushing cache file '%s' of size %s bytes to key '%s'", source_file, - os.path.getsize(source_file), rel_path) - self.transfer_progress = 0 # Reset transfer progress counter - if not self.bucket.get(rel_path): - created_obj = self.bucket.create_object(rel_path) - created_obj.upload(source_file) - else: - self.bucket.get(rel_path).upload(source_file) - 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) + self.bucket.get(rel_path).upload(source_file) + 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) return True else: log.error("Tried updating key '%s' from source file '%s', but source file does not exist.", From fb25495d191c2172ce22496900db2310d6fc9b69 Mon Sep 17 00:00:00 2001 From: vahid Date: Wed, 2 Aug 2017 12:11:26 -0700 Subject: [PATCH 12/26] Removed a S3-related comment, and swapped an if/else condition. --- lib/galaxy/objectstore/cloud.py | 12 +++--------- 1 file changed, 3 insertions(+), 9 deletions(-) diff --git a/lib/galaxy/objectstore/cloud.py b/lib/galaxy/objectstore/cloud.py index dac5dd6baa6..432f1b01dd9 100644 --- a/lib/galaxy/objectstore/cloud.py +++ b/lib/galaxy/objectstore/cloud.py @@ -351,11 +351,11 @@ class Cloud( ObjectStore ): log.debug("Pushing cache file '%s' of size %s bytes to key '%s'", source_file, os.path.getsize(source_file), rel_path) self.transfer_progress = 0 # Reset transfer progress counter - if not self.bucket.get(rel_path): + if self.bucket.get(rel_path): + self.bucket.get(rel_path).upload(source_file) + else: created_obj = self.bucket.create_object(rel_path) created_obj.upload(source_file) - else: - self.bucket.get(rel_path).upload(source_file) 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) @@ -438,12 +438,6 @@ class Cloud( ObjectStore ): if not os.path.exists(cache_dir): os.makedirs(cache_dir) - # Although not really necessary to create S3 folders (because S3 has - # flat namespace), do so for consistency with the regular file system - # S3 folders are marked by having trailing '/' so add it now - # s3_dir = '%s/' % rel_path - # 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) open(os.path.join(self.staging_path, rel_path), 'w').close() From 7727c0bf627af5905c6e029417602389238f120e Mon Sep 17 00:00:00 2001 From: vahid Date: Wed, 2 Aug 2017 12:33:53 -0700 Subject: [PATCH 13/26] Removed a comment section. --- lib/galaxy/objectstore/cloud.py | 23 ----------------------- 1 file changed, 23 deletions(-) diff --git a/lib/galaxy/objectstore/cloud.py b/lib/galaxy/objectstore/cloud.py index 432f1b01dd9..f6d0b3c7217 100644 --- a/lib/galaxy/objectstore/cloud.py +++ b/lib/galaxy/objectstore/cloud.py @@ -269,29 +269,6 @@ class Cloud( ObjectStore ): # log.debug("------ Checking cache for rel_path %s" % rel_path) cache_path = self._get_cache_path(rel_path) return os.path.exists(cache_path) - # TODO: Part of checking if a file is in cache should be to ensure the - # size of the cached file matches that on S3. Once the upload tool explicitly - # creates, this check sould be implemented- in the mean time, it's not - # looking likely to be implementable reliably. - # if os.path.exists(cache_path): - # # print "***1 %s exists" % cache_path - # if self._key_exists(rel_path): - # # print "***2 %s exists in S3" % rel_path - # # Make sure the size in cache is available in its entirety - # # print "File '%s' cache size: %s, S3 size: %s" % (cache_path, os.path.getsize(cache_path), self._get_size_in_cloud(rel_path)) - # if os.path.getsize(cache_path) == self._get_size_in_cloud(rel_path): - # # print "***2.1 %s exists in S3 and the size is the same as in cache (in_cache=True)" % rel_path - # exists = True - # else: - # # print "***2.2 %s exists but differs in size from cache (in_cache=False)" % cache_path - # exists = False - # else: - # # Although not perfect decision making, this most likely means - # # that the file is currently being uploaded - # # print "***3 %s found in cache but not in S3 (in_cache=True)" % cache_path - # exists = True - # else: - # return False def _pull_into_cache(self, rel_path): # Ensure the cache directory structure exists (e.g., dataset_#_files/) From 37efbd03ce4215086d29e8f4e9fd977b8e1299e0 Mon Sep 17 00:00:00 2001 From: vahid Date: Sun, 6 Aug 2017 22:48:42 -0700 Subject: [PATCH 14/26] Changed indentation of two lines. --- lib/galaxy/objectstore/cloud.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/lib/galaxy/objectstore/cloud.py b/lib/galaxy/objectstore/cloud.py index f6d0b3c7217..18f5da4d507 100644 --- a/lib/galaxy/objectstore/cloud.py +++ b/lib/galaxy/objectstore/cloud.py @@ -326,7 +326,7 @@ class Cloud( ObjectStore ): return True start_time = datetime.now() log.debug("Pushing cache file '%s' of size %s bytes to key '%s'", source_file, - os.path.getsize(source_file), rel_path) + os.path.getsize(source_file), rel_path) self.transfer_progress = 0 # Reset transfer progress counter if self.bucket.get(rel_path): self.bucket.get(rel_path).upload(source_file) @@ -335,7 +335,7 @@ class Cloud( ObjectStore ): created_obj.upload(source_file) 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) + source_file, rel_path, os.path.getsize(source_file), end_time - start_time) return True else: log.error("Tried updating key '%s' from source file '%s', but source file does not exist.", From 9a5016fdcf4f36d5c20991f44d4b02fe7904ef66 Mon Sep 17 00:00:00 2001 From: Vahid Date: Thu, 17 Aug 2017 01:55:52 -0700 Subject: [PATCH 15/26] Update CloudBridge to its current latest Should use the wheel created at [this PR](https://github.com/galaxyproject/starforge/pull/139). --- lib/galaxy/dependencies/conditional-requirements.txt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/galaxy/dependencies/conditional-requirements.txt b/lib/galaxy/dependencies/conditional-requirements.txt index 466b2e9ab2e..8aaeb1746ea 100644 --- a/lib/galaxy/dependencies/conditional-requirements.txt +++ b/lib/galaxy/dependencies/conditional-requirements.txt @@ -13,7 +13,7 @@ graphitesend azure-storage==0.32.0 # PyRods not in PyPI python-ldap==2.4.27 -cloudbridge==0.3.2 +cloudbridge==0.3.3 # Synnefo / Pithos+ object store client kamaki From 4e17c8e3b790d2d21bcf12c3d4a855b1226d810d Mon Sep 17 00:00:00 2001 From: vahid Date: Thu, 17 Aug 2017 19:05:25 -0700 Subject: [PATCH 16/26] Add Cloud configuration to object_store_conf.xml.sample --- config/object_store_conf.xml.sample | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/config/object_store_conf.xml.sample b/config/object_store_conf.xml.sample index 655e7271fbe..59170e046ca 100644 --- a/config/object_store_conf.xml.sample +++ b/config/object_store_conf.xml.sample @@ -52,6 +52,16 @@ --> + + From a8784369e4b243856b9582f9417039bf0c6bc83a Mon Sep 17 00:00:00 2001 From: vahid Date: Wed, 23 Aug 2017 12:41:15 -0700 Subject: [PATCH 17/26] Removed white spaces. --- lib/galaxy/objectstore/cloud.py | 24 ++++++++++++------------ 1 file changed, 12 insertions(+), 12 deletions(-) diff --git a/lib/galaxy/objectstore/cloud.py b/lib/galaxy/objectstore/cloud.py index 18f5da4d507..7ffdb32274e 100644 --- a/lib/galaxy/objectstore/cloud.py +++ b/lib/galaxy/objectstore/cloud.py @@ -27,21 +27,21 @@ from cloudbridge.cloud.factory import CloudProviderFactory, ProviderList NO_BOTO_ERROR_MESSAGE = ("Cloud object store is configured, but no boto dependency available." "Please install and properly configure boto or modify object store configuration.") -log = logging.getLogger( __name__ ) +log = logging.getLogger(__name__) logging.getLogger('boto').setLevel(logging.INFO) # Otherwise boto is quite noisy -class Cloud( ObjectStore ): +class Cloud(ObjectStore): """ Object store that stores objects as items in an cloud storage. A local cache exists that is used as an intermediate location for files between Galaxy and the cloud storage. """ - def __init__( self, config, config_xml ): - super( Cloud, self ).__init__( config ) + def __init__(self, config, config_xml): + super(Cloud, self).__init__(config) self.staging_path = self.config.file_path self.transfer_progress = 0 - self._parse_config_xml( config_xml ) + self._parse_config_xml(config_xml) self._configure_connection() self.bucket = self._get_bucket(self.bucket) # Clean cache only if value is set in galaxy.ini @@ -60,7 +60,7 @@ class Cloud( ObjectStore ): except OSError: self.use_axel = False - def _configure_connection( self ): + def _configure_connection(self): log.debug("Configuring AWS-S3 Connection") aws_config = {'aws_access_key': self.access_key, 'aws_secret_key': self.secret_key} @@ -424,8 +424,8 @@ class Cloud( ObjectStore ): if self.exists(obj, **kwargs): return bool(self.size(obj, **kwargs) > 0) else: - raise ObjectNotFound( 'objectstore.empty, object does not exist: %s, kwargs: %s' - % ( str( obj ), str( kwargs ) ) ) + raise ObjectNotFound('objectstore.empty, object does not exist: %s, kwargs: %s' + % (str(obj), str(kwargs))) def size(self, obj, **kwargs): rel_path = self._construct_path(obj, **kwargs) @@ -522,8 +522,8 @@ class Cloud( ObjectStore ): # even if it does not exist. # if dir_only: # return cache_path - raise ObjectNotFound( 'objectstore.get_filename, no cache_path: %s, kwargs: %s' - % ( str( obj ), str( kwargs ) ) ) + raise ObjectNotFound('objectstore.get_filename, no cache_path: %s, kwargs: %s' + % (str(obj), str(kwargs))) # return cache_path # Until the upload tool does not explicitly create the dataset, return expected path def update_from_file(self, obj, file_name=None, create=False, **kwargs): @@ -548,8 +548,8 @@ class Cloud( ObjectStore ): # Update the file on cloud self._push_to_os(rel_path, source_file) else: - raise ObjectNotFound( 'objectstore.update_from_file, object does not exist: %s, kwargs: %s' - % ( str( obj ), str( kwargs ) ) ) + raise ObjectNotFound('objectstore.update_from_file, object does not exist: %s, kwargs: %s' + % (str(obj), str(kwargs))) def get_object_url(self, obj, **kwargs): if self.exists(obj, **kwargs): From eaa5931d195d4710d939e6f8d53bfe2fed38ba2f Mon Sep 17 00:00:00 2001 From: vahid Date: Wed, 23 Aug 2017 13:29:43 -0700 Subject: [PATCH 18/26] Updated requirements: sqlalchemy-migrate and pbr. --- lib/galaxy/dependencies/pinned-requirements.txt | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/lib/galaxy/dependencies/pinned-requirements.txt b/lib/galaxy/dependencies/pinned-requirements.txt index 9c8eab78bd0..2801b8905ec 100644 --- a/lib/galaxy/dependencies/pinned-requirements.txt +++ b/lib/galaxy/dependencies/pinned-requirements.txt @@ -58,11 +58,11 @@ psutil==4.1.0 pulsar-galaxy-lib==0.7.0.dev5 # sqlalchemy-migrate and dependencies -sqlalchemy-migrate==0.10.0 +sqlalchemy-migrate==0.11.0 decorator==4.0.2 Tempita==0.5.3dev sqlparse==0.1.16 -pbr==1.8.0 +pbr==2.0.0 # svgwrite and dependencies svgwrite==1.1.6 From 90c444d9a136400af608cb7fd7f8e8d59de3ba90 Mon Sep 17 00:00:00 2001 From: vahid Date: Wed, 23 Aug 2017 18:05:00 -0700 Subject: [PATCH 19/26] Removed a white space between parenthesis and arguments --- lib/galaxy/dependencies/__init__.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/galaxy/dependencies/__init__.py b/lib/galaxy/dependencies/__init__.py index eee90ce5ddb..4bceea79620 100644 --- a/lib/galaxy/dependencies/__init__.py +++ b/lib/galaxy/dependencies/__init__.py @@ -109,7 +109,7 @@ class ConditionalDependencies(object): def check_azure_storage(self): return 'azure_blob' in self.object_stores - def check_cloudbridge( self ): + def check_cloudbridge(self): return 'cloud' in self.object_stores def check_kamaki(self): From f6586a8e6d861c8c6803f96340641abf6470c2fc Mon Sep 17 00:00:00 2001 From: vahid Date: Thu, 24 Aug 2017 11:43:21 -0700 Subject: [PATCH 20/26] Surrounded CloudBridge import in a try-catch block. --- lib/galaxy/objectstore/cloud.py | 13 ++++++++----- 1 file changed, 8 insertions(+), 5 deletions(-) diff --git a/lib/galaxy/objectstore/cloud.py b/lib/galaxy/objectstore/cloud.py index 7ffdb32274e..94510bd000a 100644 --- a/lib/galaxy/objectstore/cloud.py +++ b/lib/galaxy/objectstore/cloud.py @@ -20,16 +20,19 @@ from galaxy.util import ( umask_fix_perms, ) from galaxy.util.sleeper import Sleeper - from ..objectstore import convert_bytes, ObjectStore -from cloudbridge.cloud.factory import CloudProviderFactory, ProviderList - -NO_BOTO_ERROR_MESSAGE = ("Cloud object store is configured, but no boto dependency available." - "Please install and properly configure boto or modify object store configuration.") log = logging.getLogger(__name__) logging.getLogger('boto').setLevel(logging.INFO) # Otherwise boto is quite noisy +try: + from cloudbridge.cloud.factory import CloudProviderFactory, ProviderList +except ImportError: + log.error("Could not import CloudBridge.") + +NO_BOTO_ERROR_MESSAGE = ("Cloud object store is configured, but no boto dependency available." + "Please install and properly configure boto or modify object store configuration.") + class Cloud(ObjectStore): """ From 26cbf04ed0f30f72c994af0ffab54eab7dd8ac0b Mon Sep 17 00:00:00 2001 From: John Chilton Date: Tue, 5 Sep 2017 09:36:39 -0400 Subject: [PATCH 21/26] Rework conditional dependency handling in galaxy.objectstore.cloud. This file may be loaded (e.g. for testing) but not used, so we shouldn't log a generic error about cloudbridge being unavailable, if someone attempts to actually use the object store and it isn't available then raise an informative exception. --- lib/galaxy/objectstore/cloud.py | 18 ++++++++++++------ 1 file changed, 12 insertions(+), 6 deletions(-) diff --git a/lib/galaxy/objectstore/cloud.py b/lib/galaxy/objectstore/cloud.py index 94510bd000a..316b9493d4a 100644 --- a/lib/galaxy/objectstore/cloud.py +++ b/lib/galaxy/objectstore/cloud.py @@ -22,16 +22,20 @@ from galaxy.util import ( from galaxy.util.sleeper import Sleeper from ..objectstore import convert_bytes, ObjectStore -log = logging.getLogger(__name__) -logging.getLogger('boto').setLevel(logging.INFO) # Otherwise boto is quite noisy - try: from cloudbridge.cloud.factory import CloudProviderFactory, ProviderList except ImportError: - log.error("Could not import CloudBridge.") + CloudProviderFactory = None + ProviderList = None -NO_BOTO_ERROR_MESSAGE = ("Cloud object store is configured, but no boto dependency available." - "Please install and properly configure boto or modify object store configuration.") +log = logging.getLogger(__name__) + +logging.getLogger('boto').setLevel(logging.INFO) # Otherwise boto is quite noisy + +NO_CLOUDBRIDGE_ERROR_MESSAGE = ( + "ObjectStore configured, but no cloudbridge dependency available." + "Please install cloudbridge or modify Object Store configuration." +) class Cloud(ObjectStore): @@ -42,6 +46,8 @@ class Cloud(ObjectStore): """ def __init__(self, config, config_xml): super(Cloud, self).__init__(config) + if CloudProviderFactory is None: + raise Exception(NO_CLOUDBRIDGE_ERROR_MESSAGE) self.staging_path = self.config.file_path self.transfer_progress = 0 self._parse_config_xml(config_xml) From 6d65e902b12b4fca9af613085b5bc828c8a6d2fc Mon Sep 17 00:00:00 2001 From: vahid Date: Tue, 5 Sep 2017 10:57:36 -0700 Subject: [PATCH 22/26] Updated an error message. --- lib/galaxy/objectstore/cloud.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/lib/galaxy/objectstore/cloud.py b/lib/galaxy/objectstore/cloud.py index 316b9493d4a..191bcf8d4bd 100644 --- a/lib/galaxy/objectstore/cloud.py +++ b/lib/galaxy/objectstore/cloud.py @@ -33,8 +33,8 @@ log = logging.getLogger(__name__) logging.getLogger('boto').setLevel(logging.INFO) # Otherwise boto is quite noisy NO_CLOUDBRIDGE_ERROR_MESSAGE = ( - "ObjectStore configured, but no cloudbridge dependency available." - "Please install cloudbridge or modify Object Store configuration." + "Cloud ObjectStore is configured, but no CloudBridge dependency available." + "Please install CloudBridge or modify ObjectStore configuration." ) From cfce6c221c3f453e75542d53262f53c7e3411f86 Mon Sep 17 00:00:00 2001 From: vahid Date: Tue, 5 Sep 2017 10:58:23 -0700 Subject: [PATCH 23/26] Removed `boto` logging level set. --- lib/galaxy/objectstore/cloud.py | 2 -- 1 file changed, 2 deletions(-) diff --git a/lib/galaxy/objectstore/cloud.py b/lib/galaxy/objectstore/cloud.py index 191bcf8d4bd..7edef0a0144 100644 --- a/lib/galaxy/objectstore/cloud.py +++ b/lib/galaxy/objectstore/cloud.py @@ -30,8 +30,6 @@ except ImportError: log = logging.getLogger(__name__) -logging.getLogger('boto').setLevel(logging.INFO) # Otherwise boto is quite noisy - NO_CLOUDBRIDGE_ERROR_MESSAGE = ( "Cloud ObjectStore is configured, but no CloudBridge dependency available." "Please install CloudBridge or modify ObjectStore configuration." From b55583701cbf33f52db88f147dda6b04499243d7 Mon Sep 17 00:00:00 2001 From: vahid Date: Tue, 5 Sep 2017 11:04:26 -0700 Subject: [PATCH 24/26] Removed retries for getting a bucket. --- lib/galaxy/objectstore/cloud.py | 26 +++++++++++--------------- 1 file changed, 11 insertions(+), 15 deletions(-) diff --git a/lib/galaxy/objectstore/cloud.py b/lib/galaxy/objectstore/cloud.py index 7edef0a0144..7cba1715b85 100644 --- a/lib/galaxy/objectstore/cloud.py +++ b/lib/galaxy/objectstore/cloud.py @@ -171,21 +171,17 @@ class Cloud(ObjectStore): return def _get_bucket(self, bucket_name): - """ Sometimes a handle to a bucket is not established right away so try - it a few times. Raise error if connection is not established. """ - for i in range(5): - try: - bucket = self.conn.object_store.get(bucket_name) - if bucket is None: - log.debug("Bucket not found, creating a bucket with handle '%s'", bucket_name) - bucket = self.conn.object_store.create(bucket_name) - log.debug("Using cloud object store with bucket '%s'", bucket.name) - return bucket - except Exception: - log.exception("Could not get bucket '%s', attempt %s/5", bucket_name, i + 1) - time.sleep(2) - # All the attempts have been exhausted and connection was not established, - # raise error + try: + bucket = self.conn.object_store.get(bucket_name) + if bucket is None: + log.debug("Bucket not found, creating a bucket with handle '%s'", bucket_name) + bucket = self.conn.object_store.create(bucket_name) + log.debug("Using cloud ObjectStore with bucket '%s'", bucket.name) + return bucket + except Exception: + # These two generic exceptions will be replaced by specific exceptions + # once proper exceptions are exposed by CloudBridge. + log.exception("Could not get bucket '%s'.", bucket_name) raise Exception def _fix_permissions(self, rel_path): From 5977056715a5d2e12b638012b6fc99d885ea034a Mon Sep 17 00:00:00 2001 From: vahid Date: Mon, 11 Sep 2017 22:20:22 -0700 Subject: [PATCH 25/26] Resolved a bug with uploading a dataset using Cloud (not differentiating between uploading from a file vs. string). --- lib/galaxy/objectstore/cloud.py | 34 ++++++++++++++++++++++----------- 1 file changed, 23 insertions(+), 11 deletions(-) diff --git a/lib/galaxy/objectstore/cloud.py b/lib/galaxy/objectstore/cloud.py index 7cba1715b85..01ee1e9a006 100644 --- a/lib/galaxy/objectstore/cloud.py +++ b/lib/galaxy/objectstore/cloud.py @@ -327,18 +327,30 @@ class Cloud(ObjectStore): log.debug("Wanted to push file '%s' to S3 key '%s' but its size is 0; skipping.", source_file, rel_path) return True - start_time = datetime.now() - log.debug("Pushing cache file '%s' of size %s bytes to key '%s'", source_file, - os.path.getsize(source_file), rel_path) - self.transfer_progress = 0 # Reset transfer progress counter - if self.bucket.get(rel_path): - self.bucket.get(rel_path).upload(source_file) + if from_string: + if not self.bucket.get(rel_path): + created_obj = self.bucket.create_object(rel_path) + created_obj.upload(source_file) + else: + self.bucket.get(rel_path).upload(source_file) + log.debug("Pushed data from string '%s' to key '%s'", from_string, rel_path) else: - created_obj = self.bucket.create_object(rel_path) - created_obj.upload(source_file) - 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) + start_time = datetime.now() + log.debug("Pushing cache file '%s' of size %s bytes to key '%s'", source_file, + os.path.getsize(source_file), rel_path) + mb_size = os.path.getsize(source_file) / 1e6 + self.transfer_progress = 0 # Reset transfer progress counter + if not self.bucket.get(rel_path): + created_obj = self.bucket.create_object(rel_path) + created_obj.upload_from_file(source_file) + else: + self.bucket.get(rel_path).upload_from_file(source_file) + # else: + # multipart_upload(self.s3server, bucket, bucket.get(rel_path).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) return True else: log.error("Tried updating key '%s' from source file '%s', but source file does not exist.", From 45ac767a4eb6ab9e15a1e46a1e9e4ae70591e188 Mon Sep 17 00:00:00 2001 From: vahid Date: Mon, 11 Sep 2017 23:02:27 -0700 Subject: [PATCH 26/26] Removed an unused variable (i.e., mb_size) from Cloud --- lib/galaxy/objectstore/cloud.py | 3 --- 1 file changed, 3 deletions(-) diff --git a/lib/galaxy/objectstore/cloud.py b/lib/galaxy/objectstore/cloud.py index 01ee1e9a006..2c5569dd07a 100644 --- a/lib/galaxy/objectstore/cloud.py +++ b/lib/galaxy/objectstore/cloud.py @@ -338,15 +338,12 @@ class Cloud(ObjectStore): start_time = datetime.now() log.debug("Pushing cache file '%s' of size %s bytes to key '%s'", source_file, os.path.getsize(source_file), rel_path) - mb_size = os.path.getsize(source_file) / 1e6 self.transfer_progress = 0 # Reset transfer progress counter if not self.bucket.get(rel_path): created_obj = self.bucket.create_object(rel_path) created_obj.upload_from_file(source_file) else: self.bucket.get(rel_path).upload_from_file(source_file) - # else: - # multipart_upload(self.s3server, bucket, bucket.get(rel_path).name, source_file, mb_size) end_time = datetime.now() log.debug("Pushed cache file '%s' to key '%s' (%s bytes transfered in %s sec)",