mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Merge pull request #7154 from jmchilton/objectstore_uuid
Allow storing objects in the object store by UUID.
This commit is contained in:
@@ -528,6 +528,11 @@ galaxy:
|
||||
# it overrides any other objectstore settings.
|
||||
#object_store_config_file: config/object_store_conf.xml
|
||||
|
||||
# What Dataset attribute is used to reference files in an ObjectStore
|
||||
# implementation, default is 'id' but can also be set to 'uuid' for
|
||||
# more de-centralized usage.
|
||||
#object_store_store_by: id
|
||||
|
||||
# Galaxy sends mail for various things: subscribing users to the
|
||||
# mailing list if they request it, password resets, reporting dataset
|
||||
# errors, and sending activation emails. To do this, it needs to send
|
||||
|
||||
@@ -945,6 +945,18 @@
|
||||
:Type: str
|
||||
|
||||
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
``object_store_store_by``
|
||||
~~~~~~~~~~~~~~~~~~~~~~~~~
|
||||
|
||||
:Description:
|
||||
What Dataset attribute is used to reference files in an
|
||||
ObjectStore implementation, default is 'id' but can also be set to
|
||||
'uuid' for more de-centralized usage.
|
||||
:Default: ``id``
|
||||
:Type: str
|
||||
|
||||
|
||||
~~~~~~~~~~~~~~~
|
||||
``smtp_server``
|
||||
~~~~~~~~~~~~~~~
|
||||
|
||||
@@ -543,6 +543,8 @@ class Configuration(object):
|
||||
self.object_store = kwargs.get('object_store', 'disk')
|
||||
self.object_store_check_old_style = string_as_bool(kwargs.get('object_store_check_old_style', False))
|
||||
self.object_store_cache_path = resolve_path(kwargs.get("object_store_cache_path", "database/object_store_cache"), self.root)
|
||||
self.object_store_store_by = kwargs.get("object_store_store_by", "id")
|
||||
|
||||
# Handle AWS-specific config options for backward compatibility
|
||||
if kwargs.get('aws_access_key', None) is not None:
|
||||
self.os_access_key = kwargs.get('aws_access_key', None)
|
||||
|
||||
@@ -2003,7 +2003,7 @@ class Dataset(StorableObject, RepresentById):
|
||||
# actual database column so if SA instantiates this object - the
|
||||
# attribute won't exist yet.
|
||||
if not getattr(self, "external_extra_files_path", None):
|
||||
return self.object_store.get_filename(self, dir_only=True, extra_dir=self._extra_files_path or "dataset_%d_files" % self.id)
|
||||
return self.object_store.get_filename(self, dir_only=True, extra_dir=self._extra_files_rel_path)
|
||||
else:
|
||||
return os.path.abspath(self.external_extra_files_path)
|
||||
|
||||
@@ -2015,7 +2015,12 @@ class Dataset(StorableObject, RepresentById):
|
||||
extra_files_path = property(get_extra_files_path, set_extra_files_path)
|
||||
|
||||
def extra_files_path_exists(self):
|
||||
return self.object_store.exists(self, extra_dir=self._extra_files_path or "dataset_%d_files" % self.id, dir_only=True)
|
||||
return self.object_store.exists(self, extra_dir=self._extra_files_rel_path, dir_only=True)
|
||||
|
||||
@property
|
||||
def _extra_files_rel_path(self):
|
||||
store_by = getattr(self.object_store, "store_by", "id")
|
||||
return self._extra_files_path or "dataset_%s_files" % getattr(self, store_by)
|
||||
|
||||
def _calculate_size(self):
|
||||
if self.external_filename:
|
||||
@@ -2064,7 +2069,7 @@ class Dataset(StorableObject, RepresentById):
|
||||
if self.file_size is None:
|
||||
self.set_size()
|
||||
self.total_size = self.file_size or 0
|
||||
if self.object_store.exists(self, extra_dir=self._extra_files_path or "dataset_%d_files" % self.id, dir_only=True):
|
||||
if self.object_store.exists(self, extra_dir=self._extra_files_rel_path, dir_only=True):
|
||||
for root, dirs, files in os.walk(self.extra_files_path):
|
||||
self.total_size += sum([os.path.getsize(os.path.join(root, file)) for file in files if os.path.exists(os.path.join(root, file))])
|
||||
|
||||
@@ -2090,8 +2095,8 @@ class Dataset(StorableObject, RepresentById):
|
||||
"""Remove the file and extra files, marks deleted and purged"""
|
||||
# os.unlink( self.file_name )
|
||||
self.object_store.delete(self)
|
||||
if self.object_store.exists(self, extra_dir=self._extra_files_path or "dataset_%d_files" % self.id, dir_only=True):
|
||||
self.object_store.delete(self, entire_dir=True, extra_dir=self._extra_files_path or "dataset_%d_files" % self.id, dir_only=True)
|
||||
if self.object_store.exists(self, extra_dir=self._extra_files_rel_path, dir_only=True):
|
||||
self.object_store.delete(self, entire_dir=True, extra_dir=self._extra_files_rel_path, dir_only=True)
|
||||
# if os.path.exists( self.extra_files_path ):
|
||||
# shutil.rmtree( self.extra_files_path )
|
||||
# TODO: purge metadata files
|
||||
|
||||
@@ -81,7 +81,7 @@ class ObjectStore(object):
|
||||
000/obj.id)
|
||||
"""
|
||||
|
||||
def __init__(self, config, **kwargs):
|
||||
def __init__(self, config, config_dict={}, **kwargs):
|
||||
"""
|
||||
:type config: object
|
||||
:param config: An object, most likely populated from
|
||||
@@ -95,11 +95,16 @@ class ObjectStore(object):
|
||||
* new_file_path -- Used to set the 'temp' extra_dir.
|
||||
"""
|
||||
self.running = True
|
||||
self.extra_dirs = {}
|
||||
self.config = config
|
||||
self.check_old_style = config.object_store_check_old_style
|
||||
self.extra_dirs['job_work'] = config.jobs_directory
|
||||
self.extra_dirs['temp'] = config.new_file_path
|
||||
self.store_by = config_dict.get("store_by", None) or getattr(config, "object_store_store_by", "id")
|
||||
assert self.store_by in ["id", "uuid"]
|
||||
extra_dirs = {}
|
||||
extra_dirs['job_work'] = config.jobs_directory
|
||||
extra_dirs['temp'] = config.new_file_path
|
||||
extra_dirs.update(dict(
|
||||
(e['type'], e['path']) for e in config_dict.get('extra_dirs', [])))
|
||||
self.extra_dirs = extra_dirs
|
||||
|
||||
def shutdown(self):
|
||||
"""Close any connections for this ObjectStore."""
|
||||
@@ -227,9 +232,13 @@ class ObjectStore(object):
|
||||
return {
|
||||
'config': config_to_dict(self.config),
|
||||
'extra_dirs': extra_dirs,
|
||||
'store_by': self.store_by,
|
||||
'type': self.store_type,
|
||||
}
|
||||
|
||||
def _get_object_id(self, obj):
|
||||
return getattr(obj, self.store_by)
|
||||
|
||||
|
||||
class DiskObjectStore(ObjectStore):
|
||||
"""
|
||||
@@ -265,28 +274,21 @@ class DiskObjectStore(ObjectStore):
|
||||
:type extra_dirs: dict
|
||||
:param extra_dirs: Keys are string, values are directory paths.
|
||||
"""
|
||||
super(DiskObjectStore, self).__init__(config)
|
||||
super(DiskObjectStore, self).__init__(config, config_dict)
|
||||
self.file_path = config_dict.get("files_dir") or config.file_path
|
||||
extra_dirs = dict(
|
||||
(e['type'], e['path']) for e in config_dict.get('extra_dirs', []))
|
||||
self.extra_dirs.update(extra_dirs)
|
||||
|
||||
@classmethod
|
||||
def parse_xml(clazz, config_xml):
|
||||
extra_dirs = []
|
||||
file_path = None
|
||||
|
||||
config_dict = {}
|
||||
if config_xml is not None:
|
||||
for e in config_xml:
|
||||
if e.tag == 'files_dir':
|
||||
file_path = e.get('path')
|
||||
config_dict["files_dir"] = e.get('path')
|
||||
else:
|
||||
extra_dirs.append({"type": e.get('type'), "path": e.get('path')})
|
||||
|
||||
config_dict = {
|
||||
"files_dir": file_path,
|
||||
"extra_dirs": extra_dirs,
|
||||
}
|
||||
config_dict["extra_dirs"] = extra_dirs
|
||||
return config_dict
|
||||
|
||||
def to_dict(self):
|
||||
@@ -359,10 +361,11 @@ class DiskObjectStore(ObjectStore):
|
||||
path = base
|
||||
else:
|
||||
# Construct hashed path
|
||||
rel_path = os.path.join(*directory_hash_id(obj.id))
|
||||
obj_id = self._get_object_id(obj)
|
||||
rel_path = os.path.join(*directory_hash_id(obj_id))
|
||||
# Create a subdirectory for the object ID
|
||||
if obj_dir:
|
||||
rel_path = os.path.join(rel_path, str(obj.id))
|
||||
rel_path = os.path.join(rel_path, str(obj_id))
|
||||
# Optionally append extra_dir
|
||||
if extra_dir is not None:
|
||||
if extra_dir_at_root:
|
||||
@@ -371,7 +374,7 @@ class DiskObjectStore(ObjectStore):
|
||||
rel_path = os.path.join(rel_path, extra_dir)
|
||||
path = os.path.join(base, rel_path)
|
||||
if not dir_only:
|
||||
path = os.path.join(path, alt_name if alt_name else "dataset_%s.dat" % obj.id)
|
||||
path = os.path.join(path, alt_name if alt_name else "dataset_%s.dat" % obj_id)
|
||||
return os.path.abspath(path)
|
||||
|
||||
def exists(self, obj, **kwargs):
|
||||
@@ -559,7 +562,8 @@ class NestedObjectStore(ObjectStore):
|
||||
def _repr_object_for_exception(self, obj):
|
||||
try:
|
||||
# there are a few objects in python that don't have __class__
|
||||
return '{}(id={})'.format(obj.__class__.__name__, obj.id)
|
||||
obj_id = self._get_object_id(obj)
|
||||
return '{}({}={})'.format(obj.__class__.__name__, self.store_by, obj_id)
|
||||
except AttributeError:
|
||||
return str(obj)
|
||||
|
||||
@@ -602,7 +606,7 @@ class DistributedObjectStore(NestedObjectStore):
|
||||
:param fsmon: If True, monitor the file system for free space,
|
||||
removing backends when they get too full.
|
||||
"""
|
||||
super(DistributedObjectStore, self).__init__(config)
|
||||
super(DistributedObjectStore, self).__init__(config, config_dict)
|
||||
|
||||
self.backends = {}
|
||||
self.weighted_backend_ids = []
|
||||
@@ -788,7 +792,7 @@ class HierarchicalObjectStore(NestedObjectStore):
|
||||
|
||||
def __init__(self, config, config_dict, fsmon=False):
|
||||
"""The default contructor. Extends `NestedObjectStore`."""
|
||||
super(HierarchicalObjectStore, self).__init__(config)
|
||||
super(HierarchicalObjectStore, self).__init__(config, config_dict)
|
||||
|
||||
backends = odict()
|
||||
for order, backend_def in enumerate(config_dict["backends"]):
|
||||
|
||||
@@ -90,7 +90,7 @@ class AzureBlobObjectStore(ObjectStore):
|
||||
store_type = 'azure_blob'
|
||||
|
||||
def __init__(self, config, config_dict):
|
||||
super(AzureBlobObjectStore, self).__init__(config)
|
||||
super(AzureBlobObjectStore, self).__init__(config, config_dict)
|
||||
|
||||
self.transfer_progress = 0
|
||||
|
||||
@@ -107,10 +107,6 @@ class AzureBlobObjectStore(ObjectStore):
|
||||
self.cache_size = cache_dict.get('size', -1)
|
||||
self.staging_path = cache_dict.get('path') or self.config.object_store_cache_path
|
||||
|
||||
extra_dirs = dict(
|
||||
(e['type'], e['path']) for e in config_dict.get('extra_dirs', []))
|
||||
self.extra_dirs.update(extra_dirs)
|
||||
|
||||
self._initialize()
|
||||
|
||||
def _initialize(self):
|
||||
|
||||
@@ -44,7 +44,7 @@ class Cloud(ObjectStore, CloudConfigMixin):
|
||||
store_type = 'cloud'
|
||||
|
||||
def __init__(self, config, config_dict):
|
||||
super(Cloud, self).__init__(config)
|
||||
super(Cloud, self).__init__(config, config_dict)
|
||||
self.transfer_progress = 0
|
||||
|
||||
auth_dict = config_dict['auth']
|
||||
@@ -68,10 +68,6 @@ class Cloud(ObjectStore, CloudConfigMixin):
|
||||
self.cache_size = cache_dict.get('size', -1)
|
||||
self.staging_path = cache_dict.get('path') or self.config.object_store_cache_path
|
||||
|
||||
extra_dirs = dict(
|
||||
(e['type'], e['path']) for e in config_dict.get('extra_dirs', []))
|
||||
self.extra_dirs.update(extra_dirs)
|
||||
|
||||
self._initialize()
|
||||
|
||||
def _initialize(self):
|
||||
|
||||
@@ -87,16 +87,12 @@ class PithosObjectStore(ObjectStore):
|
||||
store_type = 'pithos'
|
||||
|
||||
def __init__(self, config, config_dict):
|
||||
super(PithosObjectStore, self).__init__(config)
|
||||
super(PithosObjectStore, self).__init__(config, config_dict)
|
||||
self.staging_path = self.config.file_path
|
||||
log.info('Parse config_xml for pithos object store')
|
||||
self.config_dict = config_dict
|
||||
log.debug(self.config_dict)
|
||||
|
||||
log.info('Define extra_dirs')
|
||||
extra_dirs = dict(
|
||||
(e['type'], e['path']) for e in config_dict.get('extra_dirs', []))
|
||||
self.extra_dirs.update(extra_dirs)
|
||||
self._initialize()
|
||||
|
||||
def _initialize(self):
|
||||
|
||||
@@ -720,6 +720,14 @@ mapping:
|
||||
Configuration file for the object store
|
||||
If this is set and exists, it overrides any other objectstore settings.
|
||||
|
||||
object_store_store_by:
|
||||
type: str
|
||||
default: id
|
||||
required: false
|
||||
desc: |
|
||||
What Dataset attribute is used to reference files in an ObjectStore implementation,
|
||||
default is 'id' but can also be set to 'uuid' for more de-centralized usage.
|
||||
|
||||
smtp_server:
|
||||
type: str
|
||||
default: ''
|
||||
|
||||
@@ -3,6 +3,7 @@ from contextlib import contextmanager
|
||||
from shutil import rmtree
|
||||
from string import Template
|
||||
from tempfile import mkdtemp
|
||||
from uuid import uuid4
|
||||
from xml.etree import ElementTree
|
||||
|
||||
import yaml
|
||||
@@ -14,6 +15,7 @@ from galaxy.objectstore.azure_blob import AzureBlobObjectStore
|
||||
from galaxy.objectstore.cloud import Cloud
|
||||
from galaxy.objectstore.pithos import PithosObjectStore
|
||||
from galaxy.objectstore.s3 import S3ObjectStore
|
||||
from galaxy.util import directory_hash_id
|
||||
|
||||
|
||||
DISK_TEST_CONFIG = """<?xml version="1.0"?>
|
||||
@@ -93,6 +95,75 @@ def test_disk_store():
|
||||
assert not os.path.exists(to_delete_real_path)
|
||||
|
||||
|
||||
DISK_TEST_CONFIG_BY_UUID_YAML = """
|
||||
type: disk
|
||||
files_dir: "${temp_directory}/files1"
|
||||
store_by: uuid
|
||||
extra_dirs:
|
||||
- type: temp
|
||||
path: "${temp_directory}/tmp1"
|
||||
- type: job_work
|
||||
path: "${temp_directory}/job_working_directory1"
|
||||
"""
|
||||
|
||||
|
||||
def test_disk_store_by_uuid():
|
||||
for config_str in [DISK_TEST_CONFIG_BY_UUID_YAML]:
|
||||
with TestConfig(config_str) as (directory, object_store):
|
||||
# Test no dataset with id 1 exists.
|
||||
absent_dataset = MockDataset(1)
|
||||
assert not object_store.exists(absent_dataset)
|
||||
|
||||
# Write empty dataset 2 in second backend, ensure it is empty and
|
||||
# exists.
|
||||
empty_dataset = MockDataset(2)
|
||||
directory.write("", "files1/%s/dataset_%s.dat" % (empty_dataset.rel_path_for_uuid_test(), empty_dataset.uuid))
|
||||
assert object_store.exists(empty_dataset)
|
||||
assert object_store.empty(empty_dataset)
|
||||
|
||||
# Write non-empty dataset in backend 1, test it is not emtpy & exists.
|
||||
hello_world_dataset = MockDataset(3)
|
||||
directory.write("Hello World!", "files1/%s/dataset_%s.dat" % (hello_world_dataset.rel_path_for_uuid_test(), hello_world_dataset.uuid))
|
||||
assert object_store.exists(hello_world_dataset)
|
||||
assert not object_store.empty(hello_world_dataset)
|
||||
|
||||
# Test get_data
|
||||
data = object_store.get_data(hello_world_dataset)
|
||||
assert data == "Hello World!"
|
||||
|
||||
data = object_store.get_data(hello_world_dataset, start=1, count=6)
|
||||
assert data == "ello W"
|
||||
|
||||
# Test Size
|
||||
|
||||
# Test absent and empty datasets yield size of 0.
|
||||
assert object_store.size(absent_dataset) == 0
|
||||
assert object_store.size(empty_dataset) == 0
|
||||
# Elsewise
|
||||
assert object_store.size(hello_world_dataset) > 0 # Should this always be the number of bytes?
|
||||
|
||||
# Test percent used (to some degree)
|
||||
percent_store_used = object_store.get_store_usage_percent()
|
||||
assert percent_store_used > 0.0
|
||||
assert percent_store_used < 100.0
|
||||
|
||||
# Test update_from_file test
|
||||
output_dataset = MockDataset(4)
|
||||
output_real_path = os.path.join(directory.temp_directory, "files1", output_dataset.rel_path_for_uuid_test(), "dataset_%s.dat" % output_dataset.uuid)
|
||||
assert not os.path.exists(output_real_path)
|
||||
output_working_path = directory.write("NEW CONTENTS", "job_working_directory1/example_output")
|
||||
object_store.update_from_file(output_dataset, file_name=output_working_path, create=True)
|
||||
assert os.path.exists(output_real_path)
|
||||
|
||||
# Test delete
|
||||
to_delete_dataset = MockDataset(5)
|
||||
to_delete_real_path = directory.write("content to be deleted!", "files1/%s/dataset_%s.dat" % (to_delete_dataset.rel_path_for_uuid_test(), to_delete_dataset.uuid))
|
||||
assert object_store.exists(to_delete_dataset)
|
||||
assert object_store.delete(to_delete_dataset)
|
||||
assert not object_store.exists(to_delete_dataset)
|
||||
assert not os.path.exists(to_delete_real_path)
|
||||
|
||||
|
||||
def test_disk_store_alt_name_relpath():
|
||||
""" Test that alt_name cannot be used to access arbitrary paths using a
|
||||
relative path
|
||||
@@ -636,8 +707,13 @@ class MockDataset(object):
|
||||
def __init__(self, id):
|
||||
self.id = id
|
||||
self.object_store_id = None
|
||||
self.uuid = uuid4()
|
||||
self.tags = []
|
||||
|
||||
def rel_path_for_uuid_test(self):
|
||||
rel_path = os.path.join(*directory_hash_id(self.uuid))
|
||||
return rel_path
|
||||
|
||||
|
||||
# Poor man's mocking. Need to get a real mocking library as real Galaxy development
|
||||
# dependnecy.
|
||||
|
||||
Reference in New Issue
Block a user