Merge pull request #7684 from jmchilton/build_objects_script

Script to populate datasets/libraries directly into an objectstore and model store.
This commit is contained in:
Martin Cech
2019-04-16 22:16:31 -04:00
committed by GitHub
14 changed files with 1153 additions and 472 deletions
+34 -10
View File
@@ -152,12 +152,17 @@ class HasTags(object):
class SerializationOptions(object):
def __init__(self, for_edit, serialize_dataset_objects=None, serialize_files_handler=None):
def __init__(self, for_edit, serialize_dataset_objects=None, serialize_files_handler=None, strip_metadata_files=None):
self.for_edit = for_edit
if serialize_dataset_objects is None:
serialize_dataset_objects = for_edit
self.serialize_dataset_objects = serialize_dataset_objects
self.serialize_files_handler = serialize_files_handler
if strip_metadata_files is None:
# If we're editing datasets - keep MetadataFile(s) in tact. For pure export
# expect metadata tool to be rerun.
strip_metadata_files = not for_edit
self.strip_metadata_files = strip_metadata_files
def attach_identifier(self, id_encoder, obj, ret_val):
if self.for_edit and obj.id:
@@ -2293,6 +2298,10 @@ class Dataset(StorableObject, RepresentById):
def serialize(self, id_encoder, serialization_options):
# serialize Dataset objects only for jobs that can actually modify these models.
assert serialization_options.serialize_dataset_objects
def to_int(n):
return int(n) if n is not None else 0
rval = dict_for(
self,
state=self.state,
@@ -2300,9 +2309,9 @@ class Dataset(StorableObject, RepresentById):
purged=self.purged,
external_filename=self.external_filename,
_extra_files_path=self._extra_files_path,
file_size=self.file_size,
file_size=to_int(self.file_size),
object_store_id=self.object_store_id,
total_size=self.total_size,
total_size=to_int(self.total_size),
uuid=str(self.uuid or '') or None,
hashes=list(map(lambda h: h.serialize(id_encoder, serialization_options), self.hashes))
)
@@ -2404,8 +2413,10 @@ class DatasetInstance(object):
def set_dataset_state(self, state):
if self.raw_set_dataset_state(state):
object_session(self).add(self.dataset)
object_session(self).flush() # flush here, because hda.flush() won't flush the Dataset object
sa_session = object_session(self)
if sa_session:
object_session(self).add(self.dataset)
object_session(self).flush() # flush here, because hda.flush() won't flush the Dataset object
state = property(get_dataset_state, set_dataset_state)
def get_file_name(self):
@@ -2796,6 +2807,7 @@ class DatasetInstance(object):
serialization_options.attach_identifier(id_encoder, self, rval)
return rval
metadata = _prepare_metadata_for_serialization(id_encoder, serialization_options, self.metadata)
rval = dict_for(
self,
create_time=self.create_time.__str__(),
@@ -2805,7 +2817,7 @@ class DatasetInstance(object):
blurb=self.blurb,
peek=self.peek,
extension=self.extension,
metadata=_prepare_metadata_for_serialization(dict(self.metadata.items())),
metadata=metadata,
designation=self.designation,
deleted=self.deleted,
visible=self.visible,
@@ -5188,6 +5200,12 @@ class MetadataFile(StorableObject, RepresentById):
# Return filename inside hashed directory
return os.path.abspath(os.path.join(path, "metadata_%d.dat" % self.id))
def serialize(self, id_encoder, serialization_options):
as_dict = dict_for(self)
serialization_options.attach_identifier(id_encoder, self, as_dict)
as_dict["uuid"] = str(self.uuid or '') or None
return as_dict
class FormDefinition(Dictifiable, RepresentById):
# The following form_builder classes are supported by the FormDefinition class.
@@ -5862,11 +5880,17 @@ def copy_list(lst, *args, **kwds):
return [el.copy(*args, **kwds) for el in lst]
def _prepare_metadata_for_serialization(metadata):
def _prepare_metadata_for_serialization(id_encoder, serialization_options, metadata):
""" Prepare metatdata for exporting. """
for name, value in list(metadata.items()):
processed_metadata = {}
for name, value in metadata.items():
# Metadata files are not needed for export because they can be
# regenerated.
if isinstance(value, MetadataFile):
del metadata[name]
return metadata
if serialization_options.strip_metadata_files:
continue
else:
value = value.serialize(id_encoder, serialization_options)
processed_metadata[name] = value
return processed_metadata
+8 -3
View File
@@ -568,10 +568,15 @@ class FileParameter(MetadataParameter):
return value
def new_file(self, dataset=None, **kwds):
if object_session(dataset):
# If there is a place to store the file (i.e. an object_store has been bound to
# Dataset) then use a MetadataFile and assume it is accessible. Otherwise use
# a MetadataTempFile.
if getattr(dataset.dataset, "object_store", False):
mf = galaxy.model.MetadataFile(name=self.spec.name, dataset=dataset, **kwds)
object_session(dataset).add(mf)
object_session(dataset).flush() # flush to assign id
sa_session = object_session(dataset)
if sa_session:
sa_session.add(mf)
sa_session.flush() # flush to assign id
return mf
else:
# we need to make a tmp file that is accessable to the head node,
+126 -49
View File
@@ -6,6 +6,7 @@ import shutil
import tarfile
import tempfile
from json import dump, dumps, load
from uuid import uuid4
import six
from bdbag import bdbag_api as bdb
@@ -18,7 +19,6 @@ from galaxy.security.idencoding import IdEncodingHelper
from galaxy.util import FILENAME_VALID_CHARS
from galaxy.util import in_directory
from galaxy.util.bunch import Bunch
from galaxy.version import VERSION_MAJOR
from ..item_attrs import add_item_annotation, get_item_annotation_str
from ... import model
@@ -29,6 +29,7 @@ ATTRS_FILENAME_IMPLICIT_COLLECTION_JOBS = 'implicit_collection_jobs_attrs.txt'
ATTRS_FILENAME_COLLECTIONS = 'collections_attrs.txt'
ATTRS_FILENAME_EXPORT = 'export_attrs.txt'
ATTRS_FILENAME_LIBRARIES = 'libraries_attrs.txt'
GALAXY_EXPORT_VERSION = "2"
class ImportOptions(object):
@@ -181,6 +182,7 @@ class ModelImportStore(object):
object_key = self.object_key
for dataset_attrs in datasets_attrs:
def handle_dataset_object_edit(dataset_instance):
if "dataset" in dataset_attrs:
assert self.import_options.allow_dataset_object_edit
@@ -224,7 +226,16 @@ class ModelImportStore(object):
]
for attribute in attributes:
if attribute in dataset_attrs:
setattr(hda, attribute, dataset_attrs[attribute])
value = dataset_attrs[attribute]
if attribute == "metadata":
def remap_objects(p, k, obj):
if isinstance(obj, dict) and "model_class" in obj and obj["model_class"] == "MetadataFile":
return (k, model.MetadataFile(dataset=hda, uuid=obj["uuid"]))
return (k, obj)
value = remap(value, remap_objects)
setattr(hda, attribute, value)
handle_dataset_object_edit(hda)
self._flush()
@@ -248,8 +259,6 @@ class ModelImportStore(object):
create_dataset=True,
flush=False,
sa_session=self.sa_session)
if 'id' in dataset_attrs and self.import_options.allow_edit:
dataset_instance.id = dataset_attrs['id']
elif model_class == "LibraryDatasetDatasetAssociation":
# Create dataset and HDA.
dataset_instance = model.LibraryDatasetDatasetAssociation(name=dataset_attrs['name'],
@@ -266,6 +275,7 @@ class ModelImportStore(object):
sa_session=self.sa_session)
else:
raise Exception("Unknown dataset instance type encountered")
self._attach_raw_id_if_editing(dataset_instance, dataset_attrs)
# Older style...
if 'uuid' in dataset_attrs:
@@ -333,7 +343,10 @@ class ModelImportStore(object):
if model_class == "HistoryDatasetAssociation" and self.user:
add_item_annotation(self.sa_session, self.user, dataset_instance, dataset_attrs['annotation'])
# TODO: Set tags.
tag_list = dataset_attrs.get('tags')
if tag_list:
tag_handler = model.tags.GalaxyTagHandler(sa_session=self.sa_session)
tag_handler.set_tags_from_list(user=self.user, item=dataset_instance, new_tags_list=tag_list)
if self.app:
self.app.datatypes_registry.set_external_metadata_tool.regenerate_imported_metadata_if_needed(
@@ -417,8 +430,19 @@ class ModelImportStore(object):
element_identifier=element_attrs['element_identifier'])
if 'hda' in element_attrs:
hda_attrs = element_attrs['hda']
hda_key = hda_attrs[object_key]
hda = object_import_tracker.hdas_by_key[hda_key]
if object_key in hda_attrs:
hda_key = hda_attrs[object_key]
hdas_by_key = object_import_tracker.hdas_by_key
if hda_key in hdas_by_key:
hda = hdas_by_key[hda_key]
else:
raise KeyError("Failed to find exported hda with key [%s] of type [%s] in [%s]" % (hda_key, object_key, hdas_by_key))
else:
hda_id = hda_attrs["id"]
hdas_by_id = object_import_tracker.hdas_by_id
if hda_id not in hdas_by_id:
raise Exception("Failed to find HDA with id [%s] in [%s]" % (hda_id, hdas_by_id))
hda = hdas_by_id[hda_id]
dce.hda = hda
elif 'child_collection' in element_attrs:
dce.child_collection = import_collection(element_attrs['child_collection'])
@@ -442,6 +466,7 @@ class ModelImportStore(object):
# create collection
dc = model.DatasetCollection(collection_type=collection_attrs['type'])
dc.populated_state = collection_attrs["populated_state"]
self._attach_raw_id_if_editing(dc, collection_attrs)
# TODO: element_count...
materialize_elements(dc)
@@ -458,8 +483,7 @@ class ModelImportStore(object):
visible=True,
name=collection_attrs['display_name'],
implicit_output_name=collection_attrs.get("implicit_output_name"))
if 'id' in collection_attrs and self.import_options.allow_edit:
hdca.id = collection_attrs['id']
self._attach_raw_id_if_editing(hdca, collection_attrs)
hdca.history = history
if new_history and self.trust_hid(collection_attrs):
@@ -474,6 +498,10 @@ class ModelImportStore(object):
assert 'id' in collection_attrs
object_import_tracker.hdcas_by_id[collection_attrs['id']] = hdca
def _attach_raw_id_if_editing(self, obj, attrs):
if self.sessionless and 'id' in attrs and self.import_options.allow_edit:
obj.id = attrs['id']
def _import_collection_implicit_input_associations(self, object_import_tracker, collections_attrs):
object_key = self.object_key
@@ -572,13 +600,38 @@ class ModelImportStore(object):
def _import_jobs(self, object_import_tracker, history):
object_key = self.object_key
def _find_hda(input_key):
hda = None
if input_key in object_import_tracker.hdas_by_key:
hda = object_import_tracker.hdas_by_key[input_key]
if input_key in object_import_tracker.hda_copied_from_sinks:
hda = object_import_tracker.hdas_by_key[object_import_tracker.hda_copied_from_sinks[input_key]]
return hda
def _find_hdca(input_key):
hdca = None
if input_key in object_import_tracker.hdcas_by_key:
hdca = object_import_tracker.hdcas_by_key[input_key]
if input_key in object_import_tracker.hdca_copied_from_sinks:
hdca = object_import_tracker.hdca_copied_from_sinks[input_key]
return hdca
#
# Create jobs.
#
jobs_attrs = self.jobs_properties()
# Create each job.
for job_attrs in jobs_attrs:
if 'id' in job_attrs:
# only thing we allow editing currently is associations for incoming jobs.
assert self.import_options.allow_edit
assert not self.sessionless
job = self.sa_session.query(model.Job).get(job_attrs["id"])
self._connect_job_io(job, job_attrs, _find_hda, _find_hdca)
self._flush()
continue
imported_job = model.Job()
imported_job.user = self.user
imported_job.history = history
@@ -612,22 +665,6 @@ class ModelImportStore(object):
self._session_add(imported_job)
self._flush()
def _find_hda(input_key):
hda = None
if input_key in object_import_tracker.hdas_by_key:
hda = object_import_tracker.hdas_by_key[input_key]
if input_key in object_import_tracker.hda_copied_from_sinks:
hda = object_import_tracker.hdas_by_key[object_import_tracker.hda_copied_from_sinks[input_key]]
return hda
def _find_hdca(input_key):
hdca = None
if input_key in object_import_tracker.hdcas_by_key:
hdca = object_import_tracker.hdcas_by_key[input_key]
if input_key in object_import_tracker.hdca_copied_from_sinks:
hdca = object_import_tracker.hdca_copied_from_sinks[input_key]
return hdca
# Connect jobs to input and output datasets.
params = self._normalize_job_parameters(imported_job, job_attrs, _find_hda, _find_hdca)
for name, value in params.items():
@@ -688,7 +725,7 @@ def _copied_from_object_key(copied_from_chain, objects_by_key):
class ObjectImportTracker(object):
"""Keep track or new and existing imported objects.
"""Keep track of new and existing imported objects.
Needed to re-establish connections and such in multiple passes.
"""
@@ -920,26 +957,40 @@ class ModelExportStore(object):
"""Export store should be used as context manager."""
@abc.abstractmethod
def __exit__(self):
def __exit__(self, exc_type, exc_val, exc_tb):
"""Export store should be used as context manager."""
class DirectoryModelExportStore(ModelExportStore):
def __init__(self, export_directory, app=None, for_edit=False, serialize_dataset_objects=None, export_files=None):
def __init__(self, export_directory, app=None, for_edit=False, serialize_dataset_objects=None, export_files=None, strip_metadata_files=True):
"""
:param export_directory: path to export directory. Will be created if it does not exist.
:param app: Galaxy App or app-like object. Must be provided if `for_edit` and/or `serialize_dataset_objects` are True
:param for_edit: Allow modifying existing HDA and dataset metadata during import.
:param serialize_dataset_objects: If True will encode IDs using the host secret. Defaults `for_edit`.
:param export_files: How files should be exported, can be 'symlink', 'copy' or None, in which case files
will not be serialized.
"""
if not os.path.exists(export_directory):
os.makedirs(export_directory)
if app is not None:
self.app = app
self.security = app.security
security = app.security
sessionless = False
else:
self.security = IdEncodingHelper(id_secret="randomdoesntmatter")
sessionless = True
security = IdEncodingHelper(id_secret="randomdoesntmatter")
self.sessionless = sessionless
self.security = security
self.export_directory = export_directory
self.serialization_options = model.SerializationOptions(
for_edit=for_edit,
serialize_dataset_objects=serialize_dataset_objects,
strip_metadata_files=strip_metadata_files,
serialize_files_handler=self,
)
self.export_files = export_files
@@ -951,6 +1002,8 @@ class DirectoryModelExportStore(ModelExportStore):
self.collections_attrs = []
self.dataset_id_to_path = {}
self.job_output_dataset_associations = {}
def serialize_files(self, dataset, as_dict):
if self.export_files is None:
return None
@@ -964,23 +1017,23 @@ class DirectoryModelExportStore(ModelExportStore):
shutil.copyfile(src, dest)
export_directory = self.export_directory
if dataset.id in self.included_datasets:
_, include_files = self.included_datasets[dataset.id]
if not include_files:
return
file_name, extra_files_path = None, None
try:
_file_name = dataset.file_name
if os.path.exists(_file_name):
file_name = _file_name
except ObjectNotFound:
pass
_, include_files = self.included_datasets[dataset.id]
if not include_files:
return
if dataset.extra_files_path_exists():
extra_files_path = dataset.extra_files_path
else:
pass
file_name, extra_files_path = None, None
try:
_file_name = dataset.file_name
if os.path.exists(_file_name):
file_name = _file_name
except ObjectNotFound:
pass
if dataset.extra_files_path_exists():
extra_files_path = dataset.extra_files_path
else:
pass
dir_name = 'datasets'
dir_path = os.path.join(export_directory, dir_name)
@@ -1023,7 +1076,7 @@ class DirectoryModelExportStore(ModelExportStore):
self.dataset_id_to_path[dataset.dataset.id] = (as_dict.get("file_name"), as_dict.get("extra_files_path"))
def exported_key(self, obj):
return self.security.encode_id(obj.id, kind='model_export')
return self.serialization_options.get_identifier(self.security, obj)
def __enter__(self):
return self
@@ -1096,12 +1149,27 @@ class DirectoryModelExportStore(ModelExportStore):
collect_datasets(library.root_folder)
def add_job_output_dataset_associations(self, job_id, name, dataset_instance):
job_output_dataset_associations = self.job_output_dataset_associations
if job_id not in job_output_dataset_associations:
job_output_dataset_associations[job_id] = {}
job_output_dataset_associations[job_id][name] = dataset_instance
def add_dataset_collection(self, collection):
self.collections_attrs.append(collection)
self.included_collections.append(collection)
def add_dataset(self, dataset, include_files=True):
self.included_datasets[dataset.id] = (dataset, include_files)
dataset_id = dataset.id
if dataset_id is None:
# Better be a sessionless export, just assign a random ID
# won't be able to de-duplicate datasets. This could be fixed
# by using object identity or attaching something to the object
# like temp_id used in serialization.
assert self.sessionless
dataset_id = uuid4().hex
self.included_datasets[dataset_id] = (dataset, include_files)
def _finalize(self):
export_directory = self.export_directory
@@ -1238,6 +1306,14 @@ class DirectoryModelExportStore(ModelExportStore):
jobs_attrs.append(job_attrs)
for job_id, job_output_dataset_associations in self.job_output_dataset_associations.items():
output_dataset_mapping = {}
for name, dataset in job_output_dataset_associations.items():
if name not in output_dataset_mapping:
output_dataset_mapping[name] = []
output_dataset_mapping[name].append(self.exported_key(dataset))
jobs_attrs.append({"id": job_id, 'output_dataset_mapping': output_dataset_mapping})
icjs_attrs = []
for icj_id, icj in implicit_collection_jobs_dict.items():
icj_attrs = icj.serialize(self.security, self.serialization_options)
@@ -1253,12 +1329,13 @@ class DirectoryModelExportStore(ModelExportStore):
export_attrs_filename = os.path.join(export_directory, ATTRS_FILENAME_EXPORT)
with open(export_attrs_filename, 'w') as export_attrs_out:
dump({"galaxy_version": VERSION_MAJOR}, export_attrs_out)
dump({"galaxy_export_version": GALAXY_EXPORT_VERSION}, export_attrs_out)
def __exit__(self, exc_type, exc_val, exc_tb):
if exc_type is None:
self._finalize()
# http://effbot.org/zone/python-with-statement.htm
# Ignores TypeError exceptions
return isinstance(exc_val, TypeError)
+130
View File
@@ -0,0 +1,130 @@
import argparse
import logging
import os
import sys
import yaml
import galaxy.model
from galaxy.datatypes.registry import example_datatype_registry_for_sample
from galaxy.model import store
from galaxy.model.store.discover import persist_target_to_export_store
from galaxy.objectstore import build_object_store_from_config
from galaxy.util.bunch import Bunch
DESCRIPTION = """Build import ready model objects from YAML description of files.
The positional argument to this script should be a a YAML file containing a data
fetch API-like YAML file describing files. This script will then import this data
into a defined object store and populate metadata corresponding to these files as
datasets into a "model store".
The YAML file should contain a dictionary of destination+elements objects or a list
of such dictionaries. Each such destination+elements dictionary should contain at least
two keys - 'destination' and 'items'.
Examples of destinations for creating libraries and just populating history datasets
are as follows:
```
destination:
type: library
name: Training Material
description: Data for selected tutorials from https://training.galaxyproject.org.
```
```
destination:
type: hdas
````
The 'items' definition should be a list of files or library folders. If library folders
need to be setup they should each be defined with a name and recursive set of items (
again files or folders). The following code fragment describes both a library folder
definition and a file entry:
```
items:
- name: "Example Folder 1"
description: "Description of what is in Example Folder 1"
items:
- url: https://raw.githubusercontent.com/eteriSokhoyan/test-data/master/cliques-high-representatives.fa
filename: cliques-high-representatives.fa
ext: fasta
info: "A cool longer description."
dbkey: "hg19"
md5: e5d21b1ea57fc9a31f8ea0110531bf3d
```
Currently for each file a filename must be supplied, url information will be stored if
provided but not fetched on-demand like the upload 2.0/data fetch endpoint.
Differences with respect to the Upload 2.0 YAML/JSON format: currently various src tags
such as 'url' are not supported, data cleaning options such as to_posix_lines, and
space_to_tab are not supported, unpacking zip files and walking directoriesm etc.. are
not supported either. The format consumed by this script should continue to evolve to
converge with the upload 2.0 format.
"""
logging.basicConfig()
log = logging.getLogger(__name__)
def main(argv=None):
if argv is None:
argv = sys.argv[1:]
args = _arg_parser().parse_args(argv)
object_store_config = Bunch(
object_store_store_by="uuid",
object_store_config_file=args.object_store_config,
object_store_check_old_style=False,
jobs_directory=None,
new_file_path=None,
umask=os.umask(0o77),
gid=os.getgid(),
)
object_store = build_object_store_from_config(object_store_config)
galaxy.model.Dataset.object_store = object_store
galaxy.model.set_datatypes_registry(example_datatype_registry_for_sample())
from galaxy.model import mapping
mapping.init("/tmp", "sqlite:///:memory:", create_tables=True, object_store=object_store)
with open(args.objects, "r") as f:
targets = yaml.load(f)
if not isinstance(targets, list):
targets = [targets]
export_path = args.export
export_type = args.export_type
if export_type is None:
export_type = "directory" if not export_path.endswith(".tgz") else "bag_archive"
export_types = {
"directory": store.DirectoryModelExportStore,
"tar": store.TarModelExportStore,
"bag_directory": store.BagDirectoryModelExportStore,
"bag_archive": store.BagArchiveModelExportStore,
}
store_class = export_types[export_type]
export_kwds = {
"serialize_dataset_objects": True,
}
with store_class(export_path, **export_kwds) as export_store:
for target in targets:
persist_target_to_export_store(target, export_store, object_store, ".")
def _arg_parser():
parser = argparse.ArgumentParser(description=DESCRIPTION)
parser.add_argument('objects', metavar='OBJECT_CONFIG', help='config file describing files to build objects for')
parser.add_argument('--object-store-config', help="object store configuration file")
parser.add_argument('-e', '--export', default="export", help='export path')
parser.add_argument('--export-type', default=None, help='export type (if needed)')
return parser
if __name__ == "__main__":
main()
+578
View File
@@ -0,0 +1,578 @@
"""Utilities for discovering files to add to a model store.
Working with input "JSON" format used for Fetch API, galaxy.json
imports, etc... High-level utilities in this file can be used during
job output discovery or for persisting Galaxy model objects
corresponding to files in other contexts.
"""
import abc
import os
from collections import namedtuple
import six
import galaxy.model
from galaxy import util
from galaxy.exceptions import (
RequestParameterInvalidException
)
from galaxy.util.hash_util import HASH_NAME_MAP
UNSET = object()
@six.add_metaclass(abc.ABCMeta)
class ModelPersistenceContext(object):
"""Class for creating datasets while finding files.
This class implement the create_dataset method that takes care of populating metadata
required for datasets and other potential model objects.
"""
def create_dataset(
self,
ext,
designation,
visible,
dbkey,
name,
filename,
metadata_source_name=None,
info=None,
library_folder=None,
link_data=False,
primary_data=None,
init_from=None,
dataset_attributes=None,
tag_list=[],
sources=[],
hashes=[],
):
sa_session = self.sa_session
# You can initialize a dataset or initialize from a dataset but not both.
if init_from:
assert primary_data is None
if primary_data:
assert init_from is None
if metadata_source_name:
assert init_from is None
if init_from:
assert metadata_source_name is None
if primary_data is not None:
primary_data.extension = ext
primary_data.visible = visible
primary_data.dbkey = dbkey
else:
if not library_folder:
primary_data = galaxy.model.HistoryDatasetAssociation(extension=ext,
designation=designation,
visible=visible,
dbkey=dbkey,
create_dataset=True,
flush=False,
sa_session=sa_session)
self.persist_object(primary_data)
if init_from:
self.permission_provider.copy_dataset_permissions(init_from, primary_data)
primary_data.state = init_from.state
else:
self.permission_provider.set_default_hda_permissions(primary_data)
else:
ld = galaxy.model.LibraryDataset(folder=library_folder, name=name)
ldda = galaxy.model.LibraryDatasetDatasetAssociation(name=name,
extension=ext,
dbkey=dbkey,
# library_dataset=ld,
user=self.user,
create_dataset=True,
flush=False,
sa_session=sa_session)
ld.library_dataset_dataset_association = ldda
ldda.raw_set_dataset_state(ldda.states.OK)
self.add_library_dataset_to_folder(library_folder, ld)
primary_data = ldda
for source_dict in sources:
source = galaxy.model.DatasetSource()
source.source_uri = source_dict["source_uri"]
primary_data.dataset.sources.append(source)
for hash_dict in hashes:
hash_object = galaxy.model.DatasetHash()
hash_object.hash_function = hash_dict["hash_function"]
hash_object.hash_value = hash_dict["hash_value"]
primary_data.dataset.hashes.append(hash_object)
self.flush()
if tag_list:
self.tag_handler.add_tags_from_list(self.job.user, primary_data, tag_list)
# Move data from temp location to dataset location
if not link_data:
self.object_store.update_from_file(primary_data.dataset, file_name=filename, create=True)
else:
primary_data.link_to(filename)
# We are sure there are no extra files, so optimize things that follow by settting total size also.
primary_data.set_size(no_extra_files=True)
# If match specified a name use otherwise generate one from
# designation.
primary_data.name = name
# Copy metadata from one of the inputs if requested.
if metadata_source_name:
metadata_source = self.metadata_source_provider.get_metadata_source(metadata_source_name)
primary_data.init_meta(copy_from=metadata_source)
elif init_from:
metadata_source = init_from
primary_data.init_meta(copy_from=init_from)
# when coming from primary dataset - respect pattern of output - this makes sense
primary_data.dbkey = dbkey
else:
primary_data.init_meta()
if info is not None:
primary_data.info = info
# add tool/metadata provided information
dataset_attributes = dataset_attributes or {}
if dataset_attributes:
# TODO: discover_files should produce a match that encorporates this -
# would simplify ToolProvidedMetadata interface and eliminate this
# crap path.
dataset_att_by_name = dict(ext='extension')
for att_set in ['name', 'info', 'ext', 'dbkey']:
dataset_att_name = dataset_att_by_name.get(att_set, att_set)
setattr(primary_data, dataset_att_name, dataset_attributes.get(att_set, getattr(primary_data, dataset_att_name)))
metadata_dict = dataset_attributes.get('metadata', None)
if metadata_dict:
if "dbkey" in dataset_attributes:
metadata_dict["dbkey"] = dataset_attributes["dbkey"]
# branch tested with tool_provided_metadata_3 / tool_provided_metadata_10
primary_data.metadata.from_JSON_dict(json_dict=metadata_dict)
else:
primary_data.set_meta()
primary_data.set_peek()
return primary_data
@abc.abstractproperty
def tag_handler(self):
"""Return a galaxy.model.tags.TagHandler-like object for persisting tags."""
@abc.abstractproperty
def user(self):
"""If bound to a database, return the user the datasets should be created for.
Return None otherwise.
"""
@abc.abstractmethod
def add_library_dataset_to_folder(self, library_folder, ld):
"""Add library dataset to persisted library folder."""
@abc.abstractmethod
def create_library_folder(self, parent_folder, name, description):
"""Create a library folder ready from supplied attributes for supplied parent."""
def add_datasets_to_history(self, datasets, for_output_dataset=None):
"""Add datasets to the history this context points at."""
def persist_object(self, obj):
"""Add the target to the persistence layer."""
def flush(self):
"""If database bound, flush the persisted objects to ensure IDs."""
@six.add_metaclass(abc.ABCMeta)
class PermissionProvider(object):
"""Interface for working with permissions while importing datasets with ModelPersistenceContext."""
@property
def permissions(self):
return UNSET
def set_default_hda_permissions(self, primary_data):
return
@abc.abstractmethod
def copy_dataset_permissions(self, init_from, primary_data):
"""Copy dataset permissions from supplied input dataset."""
class UnusedPermissionProvider(PermissionProvider):
def copy_dataset_permissions(self, init_from, primary_data):
"""Throws NotImplementedError.
This should only be called as part of job output collection where
there should be a session available to initialize this from.
"""
raise NotImplementedError()
@six.add_metaclass(abc.ABCMeta)
class MetadataSourceProvider(object):
"""Interface for working with fetching input dataset metadata with ModelPersistenceContext."""
@abc.abstractmethod
def get_metadata_source(self, input_name):
"""Get metadata for supplied input_name."""
class UnusedMetadataSourceProvider(MetadataSourceProvider):
def get_metadata_source(self, input_name):
"""Throws NotImplementedError.
This should only be called as part of job output collection where
one can actually collect metadata from inputs, this is unused in the
context of SessionlessModelPersistenceContext.
"""
raise NotImplementedError()
class SessionlessModelPersistenceContext(ModelPersistenceContext):
"""A variant of ModelPersistenceContext that persists to an export store instead of database directly."""
def __init__(self, object_store, export_store, working_directory):
self.permission_provider = UnusedPermissionProvider()
self.metadata_source_provider = UnusedMetadataSourceProvider()
self.sa_session = None
self.object_store = object_store
self.export_store = export_store
self.job_working_directory = working_directory # TODO: rename...
@property
def tag_handler(self):
raise NotImplementedError()
@property
def user(self):
return None
def add_library_dataset_to_folder(self, library_folder, ld):
library_folder.datasets.append(ld)
ld.order_id = library_folder.item_count
library_folder.item_count += 1
def create_library_folder(self, parent_folder, name, description):
nested_folder = galaxy.model.LibraryFolder(name=name, description=description, order_id=parent_folder.item_count)
parent_folder.item_count += 1
parent_folder.folders.append(nested_folder)
return nested_folder
def add_datasets_to_history(self, datasets, for_output_dataset=None):
# Consider copying these datasets to for_output_dataset copied histories
# somehow. Not sure it is worth the effort/complexity?
for dataset in datasets:
self.export_store.add_dataset(dataset)
def persist_object(self, obj):
"""No-op right now for the sessionless variant of this.
This works currently because either things are added to a target history with add_datasets_to_history
or the parent LibraryFolder was added to the export store in persist_target_to_export_store.
"""
def flush(self):
"""No-op for the sessionless variant of this, no database to flush."""
def persist_target_to_export_store(target_dict, export_store, object_store, work_directory):
replace_request_syntax_sugar(target_dict)
model_persistence_context = SessionlessModelPersistenceContext(object_store, export_store, work_directory)
assert "destination" in target_dict
assert "elements" in target_dict
destination = target_dict["destination"]
elements = target_dict["elements"]
assert "type" in destination
destination_type = destination["type"]
assert destination_type in ["library", "hdas"]
if destination_type == "library":
name = get_required_item(destination, "name", "Must specify a library name")
description = destination.get("description", "")
synopsis = destination.get("synopsis", "")
root_folder = galaxy.model.LibraryFolder(name=name, description='')
library = galaxy.model.Library(
name=name,
description=description,
synopsis=synopsis,
root_folder=root_folder,
)
persist_elements_to_folder(model_persistence_context, elements, root_folder)
export_store.export_library(library)
elif destination_type == "hdas":
persist_hdas(elements, model_persistence_context)
def persist_elements_to_folder(model_persistence_context, elements, library_folder):
for element in elements:
if "elements" in element:
assert "name" in element
name = element["name"]
description = element.get("description")
nested_folder = model_persistence_context.create_library_folder(library_folder, name, description)
persist_elements_to_folder(model_persistence_context, element["elements"], nested_folder)
else:
discovered_file = discovered_file_for_element(element, model_persistence_context.job_working_directory)
fields_match = discovered_file.match
designation = fields_match.designation
visible = fields_match.visible
ext = fields_match.ext
dbkey = fields_match.dbkey
info = element.get("info", None)
link_data = discovered_file.match.link_data
# Create new primary dataset
name = fields_match.name or designation
sources = fields_match.sources
hashes = fields_match.hashes
model_persistence_context.create_dataset(
ext=ext,
designation=designation,
visible=visible,
dbkey=dbkey,
name=name,
filename=discovered_file.path,
info=info,
library_folder=library_folder,
link_data=link_data,
sources=sources,
hashes=hashes,
)
def persist_hdas(elements, model_persistence_context):
# discover files as individual datasets for the target history
datasets = []
def collect_elements_for_history(elements):
for element in elements:
if "elements" in element:
collect_elements_for_history(element["elements"])
else:
discovered_file = discovered_file_for_element(element, model_persistence_context.job_working_directory)
fields_match = discovered_file.match
designation = fields_match.designation
ext = fields_match.ext
dbkey = fields_match.dbkey
info = element.get("info", None)
link_data = discovered_file.match.link_data
# Create new primary dataset
name = fields_match.name or designation
hda_id = discovered_file.match.object_id
primary_dataset = None
if hda_id:
primary_dataset = model_persistence_context.sa_session.query(galaxy.model.HistoryDatasetAssociation).get(hda_id)
sources = fields_match.sources
hashes = fields_match.hashes
dataset = model_persistence_context.create_dataset(
ext=ext,
designation=designation,
visible=True,
dbkey=dbkey,
name=name,
filename=discovered_file.path,
info=info,
link_data=link_data,
primary_data=primary_dataset,
sources=sources,
hashes=hashes,
)
dataset.raw_set_dataset_state('ok')
if not hda_id:
datasets.append(dataset)
collect_elements_for_history(elements)
model_persistence_context.add_datasets_to_history(datasets)
def add_datasets_to_history(self, datasets, for_output_dataset=None):
if for_output_dataset is not None:
raise NotImplementedError()
for dataset in datasets:
self.export_store.add_dataset(dataset)
def persist_object(self, obj):
pass
def flush(self):
pass
def get_required_item(from_dict, key, message):
if key not in from_dict:
raise RequestParameterInvalidException(message)
return from_dict[key]
def validate_and_normalize_target(obj):
replace_request_syntax_sugar(obj)
def replace_request_syntax_sugar(obj):
# For data libraries and hdas to make sense - allow items and items_from in place of elements
# and elements_from. This is destructive and modifies the supplied request.
if isinstance(obj, list):
for el in obj:
replace_request_syntax_sugar(el)
elif isinstance(obj, dict):
if "items" in obj:
obj["elements"] = obj["items"]
del obj["items"]
if "items_from" in obj:
obj["elements_from"] = obj["items_from"]
del obj["items_from"]
for value in obj.values():
replace_request_syntax_sugar(value)
if "src" in obj or "filename" in obj:
# item...
new_hashes = []
for key in HASH_NAME_MAP.keys():
if key in obj:
new_hashes.append({"hash_function": key, "hash_value": obj[key]})
del obj[key]
if key.lower() in obj:
new_hashes.append({"hash_function": key, "hash_value": obj[key.lower()]})
del obj[key.lower()]
if "hashes" not in obj:
obj["hashes"] = []
obj["hashes"].extend(new_hashes)
DiscoveredFile = namedtuple('DiscoveredFile', ['path', 'collector', 'match'])
def discovered_file_for_element(dataset, job_working_directory, parent_identifiers=[], collector=None):
target_directory = discover_target_directory(getattr(collector, "directory", None), job_working_directory)
filename = dataset["filename"]
# handle link_data_only here, verify filename is in directory if not linking...
if not dataset.get("link_data_only"):
path = os.path.join(target_directory, filename)
if not util.in_directory(path, target_directory):
raise Exception("Problem with tool configuration, attempting to pull in datasets from outside working directory.")
else:
path = filename
return DiscoveredFile(path, collector, JsonCollectedDatasetMatch(dataset, collector, filename, path=path, parent_identifiers=parent_identifiers))
def discover_target_directory(dir_name, job_working_directory):
if dir_name:
directory = os.path.join(job_working_directory, dir_name)
if not util.in_directory(directory, job_working_directory):
raise Exception("Problem with tool configuration, attempting to pull in datasets from outside working directory.")
return directory
else:
return job_working_directory
class JsonCollectedDatasetMatch(object):
def __init__(self, as_dict, collector, filename, path=None, parent_identifiers=[]):
self.as_dict = as_dict
self.collector = collector
self.filename = filename
self.path = path
self._parent_identifiers = parent_identifiers
@property
def designation(self):
# If collecting nested collection, grab identifier_0,
# identifier_1, etc... and join on : to build designation.
element_identifiers = self.raw_element_identifiers
if element_identifiers:
return ":".join(element_identifiers)
elif "designation" in self.as_dict:
return self.as_dict.get("designation")
elif "name" in self.as_dict:
return self.as_dict.get("name")
else:
return None
@property
def element_identifiers(self):
return self._parent_identifiers + (self.raw_element_identifiers or [self.designation])
@property
def raw_element_identifiers(self):
identifiers = []
i = 0
while True:
key = "identifier_%d" % i
if key in self.as_dict:
identifiers.append(self.as_dict.get(key))
else:
break
i += 1
return identifiers
@property
def name(self):
""" Return name or None if not defined by the discovery pattern.
"""
return self.as_dict.get("name")
@property
def dbkey(self):
return self.as_dict.get("dbkey", getattr(self.collector, "default_dbkey", "?"))
@property
def ext(self):
return self.as_dict.get("ext", getattr(self.collector, "default_ext", "data"))
@property
def visible(self):
try:
return self.as_dict["visible"].lower() == "visible"
except KeyError:
return getattr(self.collector, "default_visible", True)
@property
def link_data(self):
return bool(self.as_dict.get("link_data_only", False))
@property
def tag_list(self):
return self.as_dict.get("tags", [])
@property
def object_id(self):
return self.as_dict.get("object_id", None)
@property
def sources(self):
return self.as_dict.get("sources", [])
@property
def hashes(self):
return self.as_dict.get("hashes", [])
class RegexCollectedDatasetMatch(JsonCollectedDatasetMatch):
def __init__(self, re_match, collector, filename, path=None):
super(RegexCollectedDatasetMatch, self).__init__(
re_match.groupdict(), collector, filename, path=path
)
+5 -6
View File
@@ -99,12 +99,11 @@ def _fetch_target(upload_config, target):
url = item.get("url")
if url:
sources.append({"source_uri": url})
hashes = []
for hash_function in HASH_NAMES:
hash_value = item.get(hash_function)
if hash_value:
hashes.append({"hash_function": hash_function, "hash_value": hash_value})
_handle_hash_validation(upload_config, hash_function, hash_value, path)
hashes = item.get("hashes", [])
for hash_dict in hashes:
hash_function = hash_dict.get("hash_function")
hash_value = hash_dict.get("hash_value")
_handle_hash_validation(upload_config, hash_function, hash_value, path)
dbkey = item.get("dbkey", "?")
requested_ext = item.get("ext", "auto")
+12 -368
View File
@@ -4,11 +4,20 @@ import logging
import operator
import os
import re
from collections import namedtuple
import galaxy.model
from galaxy import util
from galaxy.model.dataset_collections.structure import UninitializedTree
from galaxy.model.store.discover import (
discover_target_directory,
discovered_file_for_element,
DiscoveredFile,
JsonCollectedDatasetMatch,
ModelPersistenceContext,
persist_elements_to_folder,
persist_hdas,
RegexCollectedDatasetMatch,
UNSET,
)
from galaxy.tools.parser.output_collection_def import (
DEFAULT_DATASET_COLLECTOR_DESCRIPTION,
INPUT_DBKEY_TOKEN,
@@ -58,19 +67,6 @@ class PermissionProvider(object):
self._security_agent.copy_dataset_permissions(init_from.dataset, primary_data.dataset)
class UnusedPermissionProvider(object):
@property
def permissions(self):
return UNSET
def set_default_hda_permissions(self, primary_data):
return
def copy_dataset_permissions(self, init_from, primary_data):
raise NotImplementedError()
class MetadataSourceProvider(object):
def __init__(self, inp_data):
@@ -80,12 +76,6 @@ class MetadataSourceProvider(object):
return self._inp_data[input_name]
class UnusedMetadataSourceProvider(object):
def get_metadata_source(self, input_name):
raise NotImplementedError()
def collect_dynamic_outputs(
job_context,
output_collections,
@@ -140,7 +130,7 @@ def collect_dynamic_outputs(
if "elements" in element:
add_to_discovered_files(element["elements"], parent_identifiers + [element["name"]])
else:
discovered_file = discovered_file_for_unnamed_output(element, job_working_directory, parent_identifiers)
discovered_file = discovered_file_for_element(element, job_working_directory, parent_identifiers, collector=DEFAULT_DATASET_COLLECTOR)
filenames[discovered_file.path] = discovered_file
add_to_discovered_files(elements)
@@ -196,232 +186,6 @@ def collect_dynamic_outputs(
collection.handle_population_failed("Problem building datasets for collection.")
def persist_elements_to_folder(job_context, elements, library_folder):
for element in elements:
if "elements" in element:
assert "name" in element
name = element["name"]
description = element.get("description")
nested_folder = job_context.create_library_folder(library_folder, name, description)
persist_elements_to_folder(job_context, element["elements"], nested_folder)
else:
discovered_file = discovered_file_for_unnamed_output(element, job_context.job_working_directory)
fields_match = discovered_file.match
designation = fields_match.designation
visible = fields_match.visible
ext = fields_match.ext
dbkey = fields_match.dbkey
info = element.get("info", None)
link_data = discovered_file.match.link_data
# Create new primary dataset
name = fields_match.name or designation
sources = fields_match.sources
hashes = fields_match.hashes
job_context.create_dataset(
ext=ext,
designation=designation,
visible=visible,
dbkey=dbkey,
name=name,
filename=discovered_file.path,
info=info,
library_folder=library_folder,
link_data=link_data,
sources=sources,
hashes=hashes,
)
def persist_hdas(elements, model_create_context):
# discover files as individual datasets for the target history
datasets = []
def collect_elements_for_history(elements):
for element in elements:
if "elements" in element:
collect_elements_for_history(element["elements"])
else:
discovered_file = discovered_file_for_unnamed_output(element, model_create_context.job_working_directory)
fields_match = discovered_file.match
designation = fields_match.designation
ext = fields_match.ext
dbkey = fields_match.dbkey
info = element.get("info", None)
link_data = discovered_file.match.link_data
# Create new primary dataset
name = fields_match.name or designation
hda_id = discovered_file.match.object_id
primary_dataset = None
if hda_id:
primary_dataset = model_create_context.sa_session.query(galaxy.model.HistoryDatasetAssociation).get(hda_id)
sources = fields_match.sources
hashes = fields_match.hashes
dataset = model_create_context.create_dataset(
ext=ext,
designation=designation,
visible=True,
dbkey=dbkey,
name=name,
filename=discovered_file.path,
info=info,
link_data=link_data,
primary_data=primary_dataset,
sources=sources,
hashes=hashes,
)
dataset.raw_set_dataset_state('ok')
if not hda_id:
datasets.append(dataset)
collect_elements_for_history(elements)
model_create_context.add_datasets_to_history(datasets)
class ModelPersistenceContext(object):
def create_dataset(
self,
ext,
designation,
visible,
dbkey,
name,
filename,
metadata_source_name=None,
info=None,
library_folder=None,
link_data=False,
primary_data=None,
init_from=None,
dataset_attributes=None,
tag_list=[],
sources=[],
hashes=[],
):
sa_session = self.sa_session
# You can initialize a dataset or initialize from a dataset but not both.
if init_from:
assert primary_data is None
if primary_data:
assert init_from is None
if metadata_source_name:
assert init_from is None
if init_from:
assert metadata_source_name is None
if primary_data is not None:
primary_data.extension = ext
primary_data.visible = visible
primary_data.dbkey = dbkey
else:
if not library_folder:
primary_data = galaxy.model.HistoryDatasetAssociation(extension=ext,
designation=designation,
visible=visible,
dbkey=dbkey,
create_dataset=True,
flush=False,
sa_session=sa_session)
self.persist_object(primary_data)
if init_from:
self.permission_provider.copy_dataset_permissions(init_from, primary_data)
primary_data.state = init_from.state
else:
self.permission_provider.set_default_hda_permissions(primary_data)
else:
ld = galaxy.model.LibraryDataset(folder=library_folder, name=name)
ldda = galaxy.model.LibraryDatasetDatasetAssociation(name=name,
extension=ext,
dbkey=dbkey,
# library_dataset=ld,
user=self.user,
create_dataset=True,
flush=False,
sa_session=sa_session)
ld.library_dataset_dataset_association = ldda
ldda.raw_set_dataset_state(ldda.states.OK)
self.add_library_dataset_to_folder(library_folder, ld)
primary_data = ldda
for source_dict in sources:
source = galaxy.model.DatasetSource()
source.source_uri = source_dict["source_uri"]
primary_data.dataset.sources.append(source)
for hash_dict in hashes:
hash_object = galaxy.model.DatasetHash()
hash_object.hash_function = hash_dict["hash_function"]
hash_object.hash_value = hash_dict["hash_value"]
primary_data.dataset.hashes.append(hash_object)
self.flush()
if tag_list:
self.tag_handler.add_tags_from_list(self.job.user, primary_data, tag_list)
# Move data from temp location to dataset location
if not link_data:
self.object_store.update_from_file(primary_data.dataset, file_name=filename, create=True)
else:
primary_data.link_to(filename)
# We are sure there are no extra files, so optimize things that follow by settting total size also.
primary_data.set_size(no_extra_files=True)
# If match specified a name use otherwise generate one from
# designation.
primary_data.name = name
# Copy metadata from one of the inputs if requested.
if metadata_source_name:
metadata_source = self.metadata_source_provider.get_metadata_source(metadata_source_name)
primary_data.init_meta(copy_from=metadata_source)
elif init_from:
metadata_source = init_from
primary_data.init_meta(copy_from=init_from)
# when coming from primary dataset - respect pattern of output - this makes sense
primary_data.dbkey = dbkey
else:
primary_data.init_meta()
if info is not None:
primary_data.info = info
# add tool/metadata provided information
dataset_attributes = dataset_attributes or {}
if dataset_attributes:
# TODO: discover_files should produce a match that encorporates this -
# would simplify ToolProvidedMetadata interface and eliminate this
# crap path.
dataset_att_by_name = dict(ext='extension')
for att_set in ['name', 'info', 'ext', 'dbkey']:
dataset_att_name = dataset_att_by_name.get(att_set, att_set)
setattr(primary_data, dataset_att_name, dataset_attributes.get(att_set, getattr(primary_data, dataset_att_name)))
metadata_dict = dataset_attributes.get('metadata', None)
if metadata_dict:
if "dbkey" in dataset_attributes:
metadata_dict["dbkey"] = dataset_attributes["dbkey"]
# branch tested with tool_provided_metadata_3 / tool_provided_metadata_10
primary_data.metadata.from_JSON_dict(json_dict=metadata_dict)
else:
primary_data.set_meta()
primary_data.set_peek()
return primary_data
class JobContext(ModelPersistenceContext):
def __init__(self, tool, tool_provided_metadata, job, job_working_directory, permission_provider, metadata_source_provider, input_dbkey, object_store):
@@ -685,9 +449,6 @@ def collect_primary_datasets(job_context, output, input_ext):
return primary_datasets
DiscoveredFile = namedtuple('DiscoveredFile', ['path', 'collector', 'match'])
def discover_files(output_name, tool_provided_metadata, extra_file_collectors, job_working_directory, matchable):
extra_file_collectors = extra_file_collectors
if extra_file_collectors and extra_file_collectors[0].discover_via == "tool_provided_metadata":
@@ -704,30 +465,6 @@ def discover_files(output_name, tool_provided_metadata, extra_file_collectors, j
yield DiscoveredFile(match.path, collector, match)
def discovered_file_for_unnamed_output(dataset, job_working_directory, parent_identifiers=[]):
extra_file_collector = DEFAULT_TOOL_PROVIDED_DATASET_COLLECTOR
target_directory = discover_target_directory(extra_file_collector.directory, job_working_directory)
filename = dataset["filename"]
# handle link_data_only here, verify filename is in directory if not linking...
if not dataset.get("link_data_only"):
path = os.path.join(target_directory, filename)
if not util.in_directory(path, target_directory):
raise Exception("Problem with tool configuration, attempting to pull in datasets from outside working directory.")
else:
path = filename
return DiscoveredFile(path, extra_file_collector, JsonCollectedDatasetMatch(dataset, extra_file_collector, filename, path=path, parent_identifiers=parent_identifiers))
def discover_target_directory(dir_name, job_working_directory):
if dir_name:
directory = os.path.join(job_working_directory, dir_name)
if not util.in_directory(directory, job_working_directory):
raise Exception("Problem with tool configuration, attempting to pull in datasets from outside working directory.")
return directory
else:
return job_working_directory
def walk_over_file_collectors(extra_file_collectors, job_working_directory, matchable):
for extra_file_collector in extra_file_collectors:
assert extra_file_collector.discover_via == "pattern"
@@ -830,98 +567,5 @@ def _compose(f, g):
return lambda x: f(g(x))
class JsonCollectedDatasetMatch(object):
def __init__(self, as_dict, collector, filename, path=None, parent_identifiers=[]):
self.as_dict = as_dict
self.collector = collector
self.filename = filename
self.path = path
self._parent_identifiers = parent_identifiers
@property
def designation(self):
# If collecting nested collection, grab identifier_0,
# identifier_1, etc... and join on : to build designation.
element_identifiers = self.raw_element_identifiers
if element_identifiers:
return ":".join(element_identifiers)
elif "designation" in self.as_dict:
return self.as_dict.get("designation")
elif "name" in self.as_dict:
return self.as_dict.get("name")
else:
return None
@property
def element_identifiers(self):
return self._parent_identifiers + (self.raw_element_identifiers or [self.designation])
@property
def raw_element_identifiers(self):
identifiers = []
i = 0
while True:
key = "identifier_%d" % i
if key in self.as_dict:
identifiers.append(self.as_dict.get(key))
else:
break
i += 1
return identifiers
@property
def name(self):
""" Return name or None if not defined by the discovery pattern.
"""
return self.as_dict.get("name")
@property
def dbkey(self):
return self.as_dict.get("dbkey", self.collector.default_dbkey)
@property
def ext(self):
return self.as_dict.get("ext", self.collector.default_ext)
@property
def visible(self):
try:
return self.as_dict["visible"].lower() == "visible"
except KeyError:
return self.collector.default_visible
@property
def link_data(self):
return bool(self.as_dict.get("link_data_only", False))
@property
def tag_list(self):
return self.as_dict.get("tags", [])
@property
def object_id(self):
return self.as_dict.get("object_id", None)
@property
def sources(self):
return self.as_dict.get("sources", [])
@property
def hashes(self):
return self.as_dict.get("hashes", [])
class RegexCollectedDatasetMatch(JsonCollectedDatasetMatch):
def __init__(self, re_match, collector, filename, path=None):
super(RegexCollectedDatasetMatch, self).__init__(
re_match.groupdict(), collector, filename, path=path
)
UNSET = object()
DEFAULT_DATASET_COLLECTOR = DatasetCollector(DEFAULT_DATASET_COLLECTOR_DESCRIPTION)
DEFAULT_TOOL_PROVIDED_DATASET_COLLECTOR = ToolMetadataDatasetCollector(ToolProvidedMetadataDatasetCollection())
+8 -27
View File
@@ -8,6 +8,10 @@ from galaxy.actions.library import (
from galaxy.exceptions import (
RequestParameterInvalidException
)
from galaxy.model.store.discover import (
get_required_item,
replace_request_syntax_sugar,
)
from galaxy.tools.actions.upload_common import validate_url
from galaxy.util import (
relpath,
@@ -35,8 +39,8 @@ def validate_and_normalize_targets(trans, payload):
targets = payload.get("targets", [])
for target in targets:
destination = _get_required_item(target, "destination", "Each target must specify a 'destination'")
destination_type = _get_required_item(destination, "type", "Each target destination must specify a 'type'")
destination = get_required_item(target, "destination", "Each target must specify a 'destination'")
destination_type = get_required_item(destination, "type", "Each target destination must specify a 'type'")
if "object_id" in destination:
raise RequestParameterInvalidException("object_id not allowed to appear in the request.")
@@ -45,7 +49,7 @@ def validate_and_normalize_targets(trans, payload):
msg = template % (destination_type, VALID_DESTINATION_TYPES)
raise RequestParameterInvalidException(msg)
if destination_type == "library":
library_name = _get_required_item(destination, "name", "Must specify a library name")
library_name = get_required_item(destination, "name", "Must specify a library name")
description = destination.get("description", "")
synopsis = destination.get("synopsis", "")
library = trans.app.library_manager.create(
@@ -167,27 +171,10 @@ def validate_and_normalize_targets(trans, payload):
if "purge_source" not in item:
item["purge_source"] = False
_replace_request_syntax_sugar(targets)
replace_request_syntax_sugar(targets)
_for_each_src(check_src, targets)
def _replace_request_syntax_sugar(obj):
# For data libraries and hdas to make sense - allow items and items_from in place of elements
# and elements_from. This is destructive and modifies the supplied request.
if isinstance(obj, list):
for el in obj:
_replace_request_syntax_sugar(el)
elif isinstance(obj, dict):
if "items" in obj:
obj["elements"] = obj["items"]
del obj["items"]
if "items_from" in obj:
obj["elements_from"] = obj["items_from"]
del obj["items_from"]
for value in obj.values():
_replace_request_syntax_sugar(value)
def _handle_invalid_link_data_only_type(item):
link_data_only = item.get("link_data_only", False)
if link_data_only:
@@ -200,12 +187,6 @@ def _handle_invalid_link_data_only_elements_type(item):
raise RequestParameterInvalidException("link_data_only is invalid for derived elements from [%s]" % item.get("elements_from"))
def _get_required_item(from_dict, key, message):
if key not in from_dict:
raise RequestParameterInvalidException(message)
return from_dict[key]
def _for_each_src(f, obj):
if isinstance(obj, list):
for item in obj:
+1
View File
@@ -1,5 +1,6 @@
galaxy-objectstore
galaxy-util
bdbag
bx
h5py
isa-rwval
+2
View File
@@ -39,11 +39,13 @@ PACKAGES = [
'galaxy.model.dataset_collections',
'galaxy.model.migrate',
'galaxy.model.orm',
'galaxy.model.store',
'galaxy.model.tool_shed_install',
'galaxy.security',
]
ENTRY_POINTS = '''
[console_scripts]
gx-build-objects=galaxy.model.store.build_objects:main
'''
PACKAGE_DATA = {
# Be sure to update MANIFEST.in for source dist.
+154
View File
@@ -0,0 +1,154 @@
import os
from tempfile import mkdtemp
from galaxy import model
from galaxy.model import store
from galaxy.model.store.discover import persist_target_to_export_store
from .tools.test_history_imp_exp import _mock_app
def test_model_create_context_persist_hdas():
work_directory = mkdtemp()
with open(os.path.join(work_directory, "file1.txt"), "w") as f:
f.write("hello world\nhello world line 2")
target = {
"destination": {
"type": "hdas",
},
"elements": [{
"filename": "file1.txt",
"ext": "txt",
"dbkey": "hg19",
"name": "my file",
"md5": "e5d21b1ea57fc9a31f8ea0110531bf3d",
}],
}
app = _mock_app(store_by="uuid")
temp_directory = mkdtemp()
with store.DirectoryModelExportStore(temp_directory, serialize_dataset_objects=True) as export_store:
persist_target_to_export_store(target, export_store, app.object_store, work_directory)
u = model.User(email="collection@example.com", password="password")
import_history = model.History(name="Test History for Import", user=u)
sa_session = app.model.context
sa_session.add(u)
sa_session.add(import_history)
sa_session.flush()
assert len(import_history.datasets) == 0
import_options = store.ImportOptions(allow_dataset_object_edit=True)
import_model_store = store.get_import_model_store_for_directory(temp_directory, app=app, user=u, import_options=import_options)
with import_model_store.target_history(default_history=import_history):
import_model_store.perform_import(import_history)
assert len(import_history.datasets) == 1
imported_hda = import_history.datasets[0]
assert imported_hda.ext == "txt"
assert imported_hda.name == "my file"
assert imported_hda.metadata.data_lines == 2
assert len(imported_hda.dataset.hashes) == 1
assert imported_hda.dataset.hashes[0].hash_value == "e5d21b1ea57fc9a31f8ea0110531bf3d"
with open(imported_hda.file_name, "r") as f:
assert f.read().startswith("hello world\n")
def test_persist_target_library_dataset():
work_directory = mkdtemp()
with open(os.path.join(work_directory, "file1.txt"), "w") as f:
f.write("hello world\nhello world line 2")
target = {
"destination": {
"type": "library",
"name": "Example Library",
"description": "Example Library Description",
"synopsis": "Example Library Synopsis",
},
"elements": [{
"filename": "file1.txt",
"ext": "txt",
"dbkey": "hg19",
"name": "my file",
}],
}
sa_session = _import_library_target(target, work_directory)
new_library = _assert_one_library_created(sa_session)
assert new_library.name == "Example Library"
assert new_library.description == "Example Library Description"
assert new_library.synopsis == "Example Library Synopsis"
new_root = new_library.root_folder
assert new_root
assert new_root.name == "Example Library"
assert len(new_root.datasets) == 1
ldda = new_root.datasets[0].library_dataset_dataset_association
assert ldda.metadata.data_lines == 2
with open(ldda.file_name, "r") as f:
assert f.read().startswith("hello world\n")
def test_persist_target_library_folder():
work_directory = mkdtemp()
with open(os.path.join(work_directory, "file1.txt"), "w") as f:
f.write("hello world\nhello world line 2")
target = {
"destination": {
"type": "library",
"name": "Example Library",
"description": "Example Library Description",
"synopsis": "Example Library Synopsis",
},
"items": [{
"name": "Folder 1",
"description": "Folder 1 Description",
"items": [{
"filename": "file1.txt",
"ext": "txt",
"dbkey": "hg19",
"info": "dataset info",
"name": "my file",
}]
}],
}
sa_session = _import_library_target(target, work_directory)
new_library = _assert_one_library_created(sa_session)
new_root = new_library.root_folder
assert len(new_root.datasets) == 0
assert len(new_root.folders) == 1
child_folder = new_root.folders[0]
assert child_folder.name == "Folder 1"
assert child_folder.description == "Folder 1 Description"
assert len(child_folder.folders) == 0
assert len(child_folder.datasets) == 1
ldda = child_folder.datasets[0].library_dataset_dataset_association
assert ldda.metadata.data_lines == 2
with open(ldda.file_name, "r") as f:
assert f.read().startswith("hello world\n")
def _assert_one_library_created(sa_session):
all_libraries = sa_session.query(model.Library).all()
assert len(all_libraries) == 1, len(all_libraries)
new_library = all_libraries[0]
return new_library
def _import_library_target(target, work_directory):
app = _mock_app(store_by="uuid")
temp_directory = mkdtemp()
with store.DirectoryModelExportStore(temp_directory, app=app, serialize_dataset_objects=True) as export_store:
persist_target_to_export_store(target, export_store, app.object_store, work_directory)
u = model.User(email="library@example.com", password="password")
import_options = store.ImportOptions(allow_dataset_object_edit=True, allow_library_creation=True)
import_model_store = store.get_import_model_store_for_directory(temp_directory, app=app, user=u, import_options=import_options)
import_model_store.perform_import()
sa_session = app.model.context
return sa_session
+30 -1
View File
@@ -1,9 +1,11 @@
"""Unit tests for importing and exporting data from model stores."""
import json
import os
from tempfile import mkdtemp
from tempfile import mkdtemp, NamedTemporaryFile
from galaxy import model
from galaxy.model import store
from galaxy.model.metadata import MetadataTempFile
from galaxy.tools.imp_exp import unpack_tar_gz_archive
from .tools.test_history_imp_exp import _create_datasets, _mock_app, Dummy
@@ -273,6 +275,33 @@ def test_import_export_edit_collection():
assert len(c1.elements) == 2
def test_edit_metadata_files():
app = _mock_app(store_by="uuid")
sa_session = app.model.context
u = model.User(email="collection@example.com", password="password")
h = model.History(name="Test History", user=u)
d1 = _create_datasets(sa_session, h, 1, extension="bam")[0]
sa_session.add_all((h, d1))
sa_session.flush()
index = NamedTemporaryFile("w")
index.write("cool bam index")
metadata_dict = {"bam_index": MetadataTempFile.from_JSON({"kwds": {}, "filename": index.name})}
d1.metadata.from_JSON_dict(json_dict=metadata_dict)
assert d1.metadata.bam_index
assert isinstance(d1.metadata.bam_index, model.MetadataFile)
temp_directory = mkdtemp()
with store.DirectoryModelExportStore(temp_directory, app=app, for_edit=True, strip_metadata_files=False) as export_store:
export_store.add_dataset(d1)
import_history = model.History(name="Test History for Import", user=u)
sa_session.add(import_history)
sa_session.flush()
_perform_import_from_directory(temp_directory, app, u, import_history, store.ImportOptions(allow_edit=True))
def test_sessionless_import_edit_datasets():
app, h, temp_directory, import_history = _setup_simple_export({"for_edit": True})
# Create a model store without a session and import it.
+4 -3
View File
@@ -740,14 +740,14 @@ def test_config_parse_azure():
class TestConfig(object):
def __init__(self, config_str=DISK_TEST_CONFIG, clazz=None):
def __init__(self, config_str=DISK_TEST_CONFIG, clazz=None, store_by="id"):
self.temp_directory = mkdtemp()
if config_str.startswith("<"):
config_file = "store.xml"
else:
config_file = "store.yaml"
self.write(config_str, config_file)
config = MockConfig(self.temp_directory, config_file)
config = MockConfig(self.temp_directory, config_file, store_by=store_by)
if clazz is None:
self.object_store = objectstore.build_object_store_from_config(config)
elif config_file == "store.xml":
@@ -774,11 +774,12 @@ class TestConfig(object):
class MockConfig(object):
def __init__(self, temp_directory, config_file):
def __init__(self, temp_directory, config_file, store_by="id"):
self.file_path = temp_directory
self.object_store_config_file = os.path.join(temp_directory, config_file)
self.object_store_check_old_style = False
self.object_store_cache_path = os.path.join(temp_directory, "staging")
self.object_store_store_by = store_by
self.jobs_directory = temp_directory
self.new_file_path = temp_directory
self.umask = 0000
+61 -5
View File
@@ -38,9 +38,9 @@ def _run_jihaw_cleanup(archive_dir, app=None):
return app, jihaw.cleanup_after_job()
def _mock_app():
def _mock_app(store_by="id"):
app = MockApp()
test_object_store_config = TestConfig()
test_object_store_config = TestConfig(store_by=store_by)
app.object_store = test_object_store_config.object_store
app.model.Dataset.object_store = app.object_store
app.datatypes_registry.set_external_metadata_tool = MockSetExternalTool()
@@ -461,7 +461,7 @@ def test_export_collection_with_copied_datasets_and_overlapping_hids():
def test_export_copied_collection():
app, sa_session, h = _setup_history_for_export("Collection History with dataset from other history")
app, sa_session, h = _setup_history_for_export("Collection History with copied collection")
d1, d2 = _create_datasets(sa_session, h, 2)
@@ -496,6 +496,62 @@ def test_export_copied_collection():
assert imported_by_hid[6].copied_from_history_dataset_collection_association == imported_by_hid[3]
def test_export_copied_objects_copied_outside_history():
app, sa_session, h = _setup_history_for_export("Collection History with copied objects")
d1, d2 = _create_datasets(sa_session, h, 2)
c1 = model.DatasetCollection(collection_type="paired")
hc1 = model.HistoryDatasetCollectionAssociation(history=h, hid=3, collection=c1, name="HistoryCollectionTest1")
h.hid_counter = 4
dce1 = model.DatasetCollectionElement(collection=c1, element=d1, element_identifier="forward", element_index=0)
dce2 = model.DatasetCollectionElement(collection=c1, element=d2, element_identifier="reverse", element_index=1)
sa_session.add_all((dce1, dce2, d1, d2, hc1))
sa_session.flush()
hc2 = hc1.copy(element_destination=h)
h.add_dataset_collection(hc2)
sa_session.add(hc2)
other_h = model.History(name=h.name + "-other", user=h.user)
sa_session.add(other_h)
hc3 = hc2.copy(element_destination=other_h)
other_h.add_dataset_collection(hc3)
sa_session.add(hc3)
sa_session.flush()
hc4 = hc3.copy(element_destination=h)
h.add_dataset_collection(hc4)
sa_session.add(hc4)
sa_session.flush()
assert h.hid_counter == 10
original_by_hid = _hid_dict(h)
assert original_by_hid[7].copied_from_history_dataset_association != original_by_hid[4]
assert original_by_hid[8].copied_from_history_dataset_association != original_by_hid[5]
assert original_by_hid[9].copied_from_history_dataset_collection_association != original_by_hid[6]
imported_history = _import_export(app, h)
assert imported_history.hid_counter == 10
assert len(imported_history.dataset_collections) == 3
assert len(imported_history.datasets) == 6
_assert_distinct_hids(imported_history)
imported_by_hid = _hid_dict(imported_history)
assert imported_by_hid[4].copied_from_history_dataset_association == imported_by_hid[1]
assert imported_by_hid[5].copied_from_history_dataset_association == imported_by_hid[2]
assert imported_by_hid[6].copied_from_history_dataset_collection_association == imported_by_hid[3]
assert imported_by_hid[7].copied_from_history_dataset_association == imported_by_hid[4]
assert imported_by_hid[8].copied_from_history_dataset_association == imported_by_hid[5]
assert imported_by_hid[9].copied_from_history_dataset_collection_association == imported_by_hid[6]
def test_export_collection_hids():
app, sa_session, h = _setup_history_for_export("Collection History with dataset from this history")
@@ -548,8 +604,8 @@ def _assert_distinct(l):
assert len(l) == len(set(l))
def _create_datasets(sa_session, history, n):
return [model.HistoryDatasetAssociation(extension="txt", history=history, create_dataset=True, sa_session=sa_session, hid=i + 1) for i in range(n)]
def _create_datasets(sa_session, history, n, extension="txt"):
return [model.HistoryDatasetAssociation(extension=extension, history=history, create_dataset=True, sa_session=sa_session, hid=i + 1) for i in range(n)]
def _setup_history_for_export(history_name):