diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index ce20074c2a1..fe26813a4e2 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -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 diff --git a/lib/galaxy/model/metadata.py b/lib/galaxy/model/metadata.py index bc12efb852e..1cd05399320 100644 --- a/lib/galaxy/model/metadata.py +++ b/lib/galaxy/model/metadata.py @@ -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, diff --git a/lib/galaxy/model/store/__init__.py b/lib/galaxy/model/store/__init__.py index 14451ff74de..5eddf60123e 100644 --- a/lib/galaxy/model/store/__init__.py +++ b/lib/galaxy/model/store/__init__.py @@ -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) diff --git a/lib/galaxy/model/store/build_objects.py b/lib/galaxy/model/store/build_objects.py new file mode 100644 index 00000000000..1e3b2701d04 --- /dev/null +++ b/lib/galaxy/model/store/build_objects.py @@ -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() diff --git a/lib/galaxy/model/store/discover.py b/lib/galaxy/model/store/discover.py new file mode 100644 index 00000000000..c048a37cbd3 --- /dev/null +++ b/lib/galaxy/model/store/discover.py @@ -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 + ) diff --git a/lib/galaxy/tools/data_fetch.py b/lib/galaxy/tools/data_fetch.py index 60a010e3b38..75eaa204d48 100644 --- a/lib/galaxy/tools/data_fetch.py +++ b/lib/galaxy/tools/data_fetch.py @@ -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") diff --git a/lib/galaxy/tools/parameters/output_collect.py b/lib/galaxy/tools/parameters/output_collect.py index d40e5ba5218..8f9a0078bbb 100644 --- a/lib/galaxy/tools/parameters/output_collect.py +++ b/lib/galaxy/tools/parameters/output_collect.py @@ -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()) diff --git a/lib/galaxy/webapps/galaxy/api/_fetch_util.py b/lib/galaxy/webapps/galaxy/api/_fetch_util.py index 085fef13dc7..d979ee38861 100644 --- a/lib/galaxy/webapps/galaxy/api/_fetch_util.py +++ b/lib/galaxy/webapps/galaxy/api/_fetch_util.py @@ -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: diff --git a/packages/data/requirements.txt b/packages/data/requirements.txt index 67a3b0d0a59..0c13610d96a 100644 --- a/packages/data/requirements.txt +++ b/packages/data/requirements.txt @@ -1,5 +1,6 @@ galaxy-objectstore galaxy-util +bdbag bx h5py isa-rwval diff --git a/packages/data/setup.py b/packages/data/setup.py index e9fe1fea3bb..9d63ce30e1c 100644 --- a/packages/data/setup.py +++ b/packages/data/setup.py @@ -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. diff --git a/test/unit/test_model_discovery.py b/test/unit/test_model_discovery.py new file mode 100644 index 00000000000..5482ffe80ea --- /dev/null +++ b/test/unit/test_model_discovery.py @@ -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 diff --git a/test/unit/test_model_store.py b/test/unit/test_model_store.py index 79f133c7cb0..b7be6a19ae7 100644 --- a/test/unit/test_model_store.py +++ b/test/unit/test_model_store.py @@ -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. diff --git a/test/unit/test_objectstore.py b/test/unit/test_objectstore.py index afd8998e7bc..dd2da0009f7 100644 --- a/test/unit/test_objectstore.py +++ b/test/unit/test_objectstore.py @@ -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 diff --git a/test/unit/tools/test_history_imp_exp.py b/test/unit/tools/test_history_imp_exp.py index a62062acab5..7f9f63610ee 100644 --- a/test/unit/tools/test_history_imp_exp.py +++ b/test/unit/tools/test_history_imp_exp.py @@ -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):