Merge pull request #10437 from kxk302/irods-connection-2

Let python-irodsclient session manage connections
This commit is contained in:
Marius van den Beek
2020-10-20 15:45:54 +02:00
committed by GitHub
3 changed files with 177 additions and 193 deletions
@@ -8,7 +8,7 @@ drmaa
statsd
docker
azure-storage==0.32.0
python-irodsclient==0.8.3
python-irodsclient==0.8.4
python-ldap==3.2.0
python-pam
galaxycloudrunner
@@ -41,7 +41,7 @@ pytest-mock==3.1.1
pytest-postgresql==2.3.0
pytest-pythonpath==0.7.3
pytest==5.4.3
python-irodsclient==0.8.3
python-irodsclient==0.8.4
pytz==2020.1
recommonmark==0.6.0
requests==2.24.0
+175 -191
View File
@@ -4,7 +4,6 @@ Object Store plugin for the Integrated Rule-Oriented Data System (iRODS)
import logging
import os
import shutil
from contextlib import contextmanager
from datetime import datetime
from functools import partial
try:
@@ -111,31 +110,6 @@ def parse_config_xml(config_xml):
raise
def acquire_session(host='localhost', port='1247', user='rods', password='rods', zone='tempZone', timeout='30'):
session = iRODSSession(host=host, port=port, user=user, password=password, zone=zone)
# Set connection timeout
session.connection_timeout = timeout
return session
def release_session(session):
# This call will cleanup all the connections in the connection pool
# OSError sometimes happens on GitHub Actions, after the test has successfully completed. Ignore it if it happens.
try:
session.cleanup()
except OSError:
pass
@contextmanager
def managed_session(host='localhost', port='1247', user='rods', password='rods', zone='tempZone', timeout='30'):
session = acquire_session(host=host, port=port, user=user, password=password, zone=zone, timeout=timeout)
try:
yield session
finally:
release_session(session)
class CloudConfigMixin:
def _config_to_dict(self):
@@ -235,8 +209,23 @@ class IRODSObjectStore(DiskObjectStore, CloudConfigMixin):
raise Exception(IRODS_IMPORT_MESSAGE)
self.home = "/" + self.zone + "/home/" + self.username
if irods is None:
raise Exception(IRODS_IMPORT_MESSAGE)
self.session = iRODSSession(host=self.host, port=self.port, user=self.username, password=self.password, zone=self.zone)
# Set connection timeout
self.session.connection_timeout = self.timeout
log.debug("irods __init__ %s", reload_timer)
def shutdown(self):
# This call will cleanup all the connections in the connection pool
# OSError sometimes happens on GitHub Actions, after the test has successfully completed. Ignore it if it happens.
try:
self.session.cleanup()
except OSError:
pass
@classmethod
def parse_xml(cls, config_xml):
return parse_config_xml(config_xml)
@@ -295,43 +284,41 @@ class IRODSObjectStore(DiskObjectStore, CloudConfigMixin):
# rel_path is file or folder?
def _get_size_in_irods(self, rel_path):
with managed_session(host=self.host, port=self.port, user=self.username, password=self.password, zone=self.zone, timeout=self.timeout) as session:
p = Path(rel_path)
data_object_name = p.stem + p.suffix
subcollection_name = p.parent
p = Path(rel_path)
data_object_name = p.stem + p.suffix
subcollection_name = p.parent
collection_path = self.home + "/" + str(subcollection_name)
data_object_path = collection_path + "/" + str(data_object_name)
collection_path = self.home + "/" + str(subcollection_name)
data_object_path = collection_path + "/" + str(data_object_name)
try:
data_obj = session.data_objects.get(data_object_path)
return data_obj.__sizeof__()
except (DataObjectDoesNotExist, CollectionDoesNotExist):
log.warn("Collection or data object (%s) does not exist", data_object_path)
return -1
except NetworkException as e:
log.exception(e)
return -1
try:
data_obj = self.session.data_objects.get(data_object_path)
return data_obj.__sizeof__()
except (DataObjectDoesNotExist, CollectionDoesNotExist):
log.warn("Collection or data object (%s) does not exist", data_object_path)
return -1
except NetworkException as e:
log.exception(e)
return -1
# rel_path is file or folder?
def _data_object_exists(self, rel_path):
with managed_session(host=self.host, port=self.port, user=self.username, password=self.password, zone=self.zone, timeout=self.timeout) as session:
p = Path(rel_path)
data_object_name = p.stem + p.suffix
subcollection_name = p.parent
p = Path(rel_path)
data_object_name = p.stem + p.suffix
subcollection_name = p.parent
collection_path = self.home + "/" + str(subcollection_name)
data_object_path = collection_path + "/" + str(data_object_name)
collection_path = self.home + "/" + str(subcollection_name)
data_object_path = collection_path + "/" + str(data_object_name)
try:
session.data_objects.get(data_object_path)
return True
except (DataObjectDoesNotExist, CollectionDoesNotExist):
log.debug("Collection or data object (%s) does not exist", data_object_path)
return False
except NetworkException as e:
log.exception(e)
return False
try:
self.session.data_objects.get(data_object_path)
return True
except (DataObjectDoesNotExist, CollectionDoesNotExist):
log.debug("Collection or data object (%s) does not exist", data_object_path)
return False
except NetworkException as e:
log.exception(e)
return False
def _in_cache(self, rel_path):
""" Check if the given dataset is in the local cache and return True if so. """
@@ -349,37 +336,36 @@ class IRODSObjectStore(DiskObjectStore, CloudConfigMixin):
return file_ok
def _download(self, rel_path):
with managed_session(host=self.host, port=self.port, user=self.username, password=self.password, zone=self.zone, timeout=self.timeout) as session:
log.debug("Pulling data object '%s' into cache to %s", rel_path, self._get_cache_path(rel_path))
log.debug("Pulling data object '%s' into cache to %s", rel_path, self._get_cache_path(rel_path))
p = Path(rel_path)
data_object_name = p.stem + p.suffix
subcollection_name = p.parent
p = Path(rel_path)
data_object_name = p.stem + p.suffix
subcollection_name = p.parent
collection_path = self.home + "/" + str(subcollection_name)
data_object_path = collection_path + "/" + str(data_object_name)
data_obj = None
collection_path = self.home + "/" + str(subcollection_name)
data_object_path = collection_path + "/" + str(data_object_name)
data_obj = None
try:
data_obj = session.data_objects.get(data_object_path)
except (DataObjectDoesNotExist, CollectionDoesNotExist):
log.warn("Collection or data object (%s) does not exist", data_object_path)
return False
except NetworkException as e:
log.exception(e)
return False
try:
data_obj = self.session.data_objects.get(data_object_path)
except (DataObjectDoesNotExist, CollectionDoesNotExist):
log.warn("Collection or data object (%s) does not exist", data_object_path)
return False
except NetworkException as e:
log.exception(e)
return False
if self.cache_size > 0 and data_obj.__sizeof__() > self.cache_size:
log.critical("File %s is larger (%s) than the cache size (%s). Cannot download.",
rel_path, data_obj.__sizeof__(), self.cache_size)
return False
if self.cache_size > 0 and data_obj.__sizeof__() > self.cache_size:
log.critical("File %s is larger (%s) than the cache size (%s). Cannot download.",
rel_path, data_obj.__sizeof__(), self.cache_size)
return False
log.debug("Pulled data object '%s' into cache to %s", rel_path, self._get_cache_path(rel_path))
log.debug("Pulled data object '%s' into cache to %s", rel_path, self._get_cache_path(rel_path))
with data_obj.open('r') as data_obj_fp, open(self._get_cache_path(rel_path), "wb") as cache_fp:
for chunk in iter(partial(data_obj_fp.read, CHUNK_SIZE), b''):
cache_fp.write(chunk)
return True
with data_obj.open('r') as data_obj_fp, open(self._get_cache_path(rel_path), "wb") as cache_fp:
for chunk in iter(partial(data_obj_fp.read, CHUNK_SIZE), b''):
cache_fp.write(chunk)
return True
def _push_to_irods(self, rel_path, source_file=None, from_string=None):
"""
@@ -390,54 +376,53 @@ class IRODSObjectStore(DiskObjectStore, CloudConfigMixin):
still using ``rel_path`` for collection and object store names.
If ``from_string`` is provided, set contents of the file to the value of the string.
"""
with managed_session(host=self.host, port=self.port, user=self.username, password=self.password, zone=self.zone, timeout=self.timeout) as session:
p = Path(rel_path)
data_object_name = p.stem + p.suffix
subcollection_name = p.parent
p = Path(rel_path)
data_object_name = p.stem + p.suffix
subcollection_name = p.parent
source_file = source_file if source_file else self._get_cache_path(rel_path)
options = {kw.FORCE_FLAG_KW: ''}
source_file = source_file if source_file else self._get_cache_path(rel_path)
options = {kw.FORCE_FLAG_KW: ''}
if os.path.exists(source_file):
# Check if the data object exists in iRODS
collection_path = self.home + "/" + str(subcollection_name)
data_object_path = collection_path + "/" + str(data_object_name)
exists = False
try:
exists = session.data_objects.exists(data_object_path)
if os.path.exists(source_file):
# Check if the data object exists in iRODS
collection_path = self.home + "/" + str(subcollection_name)
data_object_path = collection_path + "/" + str(data_object_name)
exists = False
try:
exists = self.session.data_objects.exists(data_object_path)
if os.path.getsize(source_file) == 0 and exists:
log.debug("Wanted to push file '%s' to iRODS collection '%s' but its size is 0; skipping.", source_file, rel_path)
return True
if from_string:
data_obj = session.data_objects.create(data_object_path, self.resource, **options)
with data_obj.open('w') as data_obj_fp:
data_obj_fp.write(from_string)
log.debug("Pushed data from string '%s' to collection '%s'", from_string, data_object_path)
else:
start_time = datetime.now()
log.debug("Pushing cache file '%s' of size %s bytes to collection '%s'", source_file, os.path.getsize(source_file), rel_path)
# Create sub-collection first
session.collections.create(collection_path, recurse=True)
data_obj = session.data_objects.create(data_object_path, self.resource, **options)
# Write to file in subcollection created above
with open(source_file, 'rb') as content_file, data_obj.open('w') as data_obj_fp:
for chunk in iter(partial(content_file.read, CHUNK_SIZE), b''):
data_obj_fp.write(chunk)
end_time = datetime.now()
log.debug("Pushed cache file '%s' to collection '%s' (%s bytes transfered in %s sec)",
source_file, rel_path, os.path.getsize(source_file), end_time - start_time)
if os.path.getsize(source_file) == 0 and exists:
log.debug("Wanted to push file '%s' to iRODS collection '%s' but its size is 0; skipping.", source_file, rel_path)
return True
except NetworkException as e:
log.exception(e)
return False
log.error("Tried updating key '%s' from source file '%s', but source file does not exist.", rel_path, source_file)
return False
if from_string:
data_obj = self.session.data_objects.create(data_object_path, self.resource, **options)
with data_obj.open('w') as data_obj_fp:
data_obj_fp.write(from_string)
log.debug("Pushed data from string '%s' to collection '%s'", from_string, data_object_path)
else:
start_time = datetime.now()
log.debug("Pushing cache file '%s' of size %s bytes to collection '%s'", source_file, os.path.getsize(source_file), rel_path)
# Create sub-collection first
self.session.collections.create(collection_path, recurse=True)
data_obj = self.session.data_objects.create(data_object_path, self.resource, **options)
# Write to file in subcollection created above
with open(source_file, 'rb') as content_file, data_obj.open('w') as data_obj_fp:
for chunk in iter(partial(content_file.read, CHUNK_SIZE), b''):
data_obj_fp.write(chunk)
end_time = datetime.now()
log.debug("Pushed cache file '%s' to collection '%s' (%s bytes transfered in %s sec)",
source_file, rel_path, os.path.getsize(source_file), end_time - start_time)
return True
except NetworkException as e:
log.exception(e)
return False
log.error("Tried updating key '%s' from source file '%s', but source file does not exist.", rel_path, source_file)
return False
def file_ready(self, obj, **kwargs):
"""
@@ -518,76 +503,75 @@ class IRODSObjectStore(DiskObjectStore, CloudConfigMixin):
return 0
def _delete(self, obj, entire_dir=False, **kwargs):
with managed_session(host=self.host, port=self.port, user=self.username, password=self.password, zone=self.zone, timeout=self.timeout) as session:
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)
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 temporary data in JOB_WORK directory
if base_dir and dir_only and obj_dir:
shutil.rmtree(os.path.abspath(rel_path))
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 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 irods and deleing it.
if entire_dir and extra_dir:
shutil.rmtree(self._get_cache_path(rel_path))
col_path = self.home + "/" + str(rel_path)
col = None
try:
col = self.session.collections.get(col_path)
except CollectionDoesNotExist:
log.warn("Collection (%s) does not exist!", col_path)
return False
except NetworkException as e:
log.exception(e)
return False
cols = col.walk()
# Traverse the tree only one level deep
for _ in range(2):
# get next result
_, _, data_objects = next(cols)
# Delete data objects
for data_object in data_objects:
data_object.unlink(force=True)
return True
else:
# Delete from cache first
os.unlink(self._get_cache_path(rel_path))
# Delete from irods as well
p = Path(rel_path)
data_object_name = p.stem + p.suffix
subcollection_name = p.parent
collection_path = self.home + "/" + str(subcollection_name)
data_object_path = collection_path + "/" + str(data_object_name)
try:
data_obj = self.session.data_objects.get(data_object_path)
# remove object
data_obj.unlink(force=True)
return True
# For the case of extra_files, because we don't have a reference to
# individual files 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 irods and deleting it.
if entire_dir and extra_dir:
shutil.rmtree(self._get_cache_path(rel_path))
col_path = self.home + "/" + str(rel_path)
col = None
try:
col = session.collections.get(col_path)
except CollectionDoesNotExist:
log.warn("Collection (%s) does not exist!", col_path)
return False
except NetworkException as e:
log.exception(e)
return False
cols = col.walk()
# Traverse the tree only one level deep
for _ in range(2):
# get next result
_, _, data_objects = next(cols)
# Delete data objects
for data_object in data_objects:
data_object.unlink(force=True)
except (DataObjectDoesNotExist, CollectionDoesNotExist):
log.info("Collection or data object (%s) does not exist", data_object_path)
return True
else:
# Delete from cache first
os.unlink(self._get_cache_path(rel_path))
# Delete from irods as well
p = Path(rel_path)
data_object_name = p.stem + p.suffix
subcollection_name = p.parent
collection_path = self.home + "/" + str(subcollection_name)
data_object_path = collection_path + "/" + str(data_object_name)
try:
data_obj = session.data_objects.get(data_object_path)
# remove object
data_obj.unlink(force=True)
return True
except (DataObjectDoesNotExist, CollectionDoesNotExist):
log.info("Collection or data object (%s) does not exist", data_object_path)
return True
except NetworkException as e:
log.exception(e)
return False
except OSError:
log.exception('%s delete error', self._get_filename(obj, **kwargs))
except NetworkException as e:
log.exception(e)
return False
except NetworkException as e:
log.exception(e)
return False
except OSError:
log.exception('%s delete error', self._get_filename(obj, **kwargs))
except NetworkException as e:
log.exception(e)
return False
def _get_data(self, obj, start=0, count=-1, **kwargs):
rel_path = self._construct_path(obj, **kwargs)