From f6b2cc65b900e67300efb573f3b124d82836a80a Mon Sep 17 00:00:00 2001 From: Jelle Scholtalbers Date: Tue, 27 Feb 2018 11:05:29 +0100 Subject: [PATCH 01/24] add basic cgroup metric collection --- config/job_metrics_conf.xml.sample | 29 +++--- .../jobs/metrics/instrumenters/cgroup.py | 89 +++++++++++++++++++ 2 files changed, 106 insertions(+), 12 deletions(-) create mode 100644 lib/galaxy/jobs/metrics/instrumenters/cgroup.py diff --git a/config/job_metrics_conf.xml.sample b/config/job_metrics_conf.xml.sample index 1fdfb3849a9..f7b309ac686 100644 --- a/config/job_metrics_conf.xml.sample +++ b/config/job_metrics_conf.xml.sample @@ -15,7 +15,7 @@ - + - + + + - + # # # - dataset_collectors = map(dataset_collector, output_collection_def.dataset_collector_descriptions) - output_name = output_collection_def.name - filenames = self.find_files(output_name, collection, dataset_collectors) + if name is None: + name = "unnamed output" element_datasets = [] for filename, discovered_file in filenames.items(): @@ -241,6 +349,8 @@ class JobContext(object): # Create new primary dataset name = fields_match.name or designation + link_data = discovered_file.match.link_data + dataset = self.create_dataset( ext=ext, designation=designation, @@ -248,14 +358,15 @@ class JobContext(object): dbkey=dbkey, name=name, filename=filename, - metadata_source_name=output_collection_def.metadata_source, + metadata_source_name=metadata_source_name, + link_data=link_data, ) log.debug( "(%s) Created dynamic collection dataset for path [%s] with element identifier [%s] for output [%s] %s", self.job.id, filename, designation, - output_collection_def.name, + name, create_dataset_timer, ) element_datasets.append((element_identifiers, dataset)) @@ -270,7 +381,7 @@ class JobContext(object): log.debug( "(%s) Add dynamic collection datsets to history for output [%s] %s", self.job.id, - output_collection_def.name, + name, add_datasets_timer, ) @@ -300,12 +411,18 @@ class JobContext(object): dbkey, name, filename, - metadata_source_name, + metadata_source_name=None, + info=None, + library_folder=None, + link_data=False, ): app = self.app sa_session = self.sa_session - primary_data = _new_hda(app, sa_session, ext, designation, visible, dbkey, self.permissions) + if not library_folder: + primary_data = _new_hda(app, sa_session, ext, designation, visible, dbkey, self.permissions) + else: + primary_data = _new_ldda(self.work_context, name, ext, visible, dbkey, library_folder) # Copy metadata from one of the inputs if requested. metadata_source = None @@ -314,7 +431,11 @@ class JobContext(object): sa_session.flush() # Move data from temp location to dataset location - app.object_store.update_from_file(primary_data.dataset, file_name=filename, create=True) + if not link_data: + app.object_store.update_from_file(primary_data.dataset, file_name=filename, create=True) + else: + primary_data.link_to(filename) + primary_data.set_size() # If match specified a name use otherwise generate one from # designation. @@ -325,6 +446,9 @@ class JobContext(object): else: primary_data.init_meta() + if info is not None: + primary_data.info = info + primary_data.set_meta() primary_data.set_peek() @@ -491,6 +615,20 @@ 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) @@ -605,11 +743,12 @@ def _compose(f, g): class JsonCollectedDatasetMatch(object): - def __init__(self, as_dict, collector, filename, path=None): + 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): @@ -627,7 +766,7 @@ class JsonCollectedDatasetMatch(object): @property def element_identifiers(self): - return self.raw_element_identifiers or [self.designation] + return self._parent_identifiers + (self.raw_element_identifiers or [self.designation]) @property def raw_element_identifiers(self): @@ -664,6 +803,10 @@ class JsonCollectedDatasetMatch(object): except KeyError: return self.collector.default_visible + @property + def link_data(self): + return bool(self.as_dict.get("link_data_only", False)) + class RegexCollectedDatasetMatch(JsonCollectedDatasetMatch): @@ -676,6 +819,42 @@ class RegexCollectedDatasetMatch(JsonCollectedDatasetMatch): UNSET = object() +def _new_ldda( + trans, + name, + ext, + visible, + dbkey, + library_folder, +): + ld = trans.app.model.LibraryDataset(folder=library_folder, name=name) + trans.sa_session.add(ld) + trans.sa_session.flush() + trans.app.security_agent.copy_library_permissions(trans, library_folder, ld) + + ldda = trans.app.model.LibraryDatasetDatasetAssociation(name=name, + extension=ext, + dbkey=dbkey, + library_dataset=ld, + user=trans.user, + create_dataset=True, + sa_session=trans.sa_session) + trans.sa_session.add(ldda) + ldda.state = ldda.states.OK + # Permissions must be the same on the LibraryDatasetDatasetAssociation and the associated LibraryDataset + trans.app.security_agent.copy_library_permissions(trans, ld, ldda) + # Copy the current user's DefaultUserPermissions to the new LibraryDatasetDatasetAssociation.dataset + trans.app.security_agent.set_all_dataset_permissions(ldda.dataset, trans.app.security_agent.user_get_default_permissions(trans.user)) + library_folder.add_library_dataset(ld, genome_build=dbkey) + trans.sa_session.add(library_folder) + trans.sa_session.flush() + + ld.library_dataset_dataset_association_id = ldda.id + trans.sa_session.add(ld) + trans.sa_session.flush() + return ldda + + def _new_hda( app, sa_session, @@ -702,3 +881,4 @@ def _new_hda( DEFAULT_DATASET_COLLECTOR = DatasetCollector(DEFAULT_DATASET_COLLECTOR_DESCRIPTION) +DEFAULT_TOOL_PROVIDED_DATASET_COLLECTOR = ToolMetadataDatasetCollector(ToolProvidedMetadataDatasetCollection()) diff --git a/lib/galaxy/tools/special_tools.py b/lib/galaxy/tools/special_tools.py index 953e69dee64..129b7064a94 100644 --- a/lib/galaxy/tools/special_tools.py +++ b/lib/galaxy/tools/special_tools.py @@ -4,6 +4,7 @@ log = logging.getLogger(__name__) SPECIAL_TOOLS = { "history export": "galaxy/tools/imp_exp/exp_history_to_archive.xml", "history import": "galaxy/tools/imp_exp/imp_history_from_archive.xml", + "data fetch": "galaxy/tools/data_fetch.xml", } diff --git a/lib/galaxy/webapps/galaxy/api/_fetch_util.py b/lib/galaxy/webapps/galaxy/api/_fetch_util.py new file mode 100644 index 00000000000..76bee4efa22 --- /dev/null +++ b/lib/galaxy/webapps/galaxy/api/_fetch_util.py @@ -0,0 +1,211 @@ +import logging +import os + +from galaxy.actions.library import ( + validate_path_upload, + validate_server_directory_upload, +) +from galaxy.exceptions import ( + RequestParameterInvalidException +) +from galaxy.tools.actions.upload_common import validate_url +from galaxy.util import ( + relpath, +) + +log = logging.getLogger(__name__) + +VALID_DESTINATION_TYPES = ["library", "library_folder", "hdca"] +ELEMENTS_FROM_TYPE = ["archive", "bagit", "bagit_archive", "directory"] +# These elements_from cannot be sym linked to because they only exist during upload. +ELEMENTS_FROM_TRANSIENT_TYPES = ["archive", "bagit_archive"] + + +def validate_and_normalize_targets(trans, payload): + """Validate and normalize all src references in fetch targets. + + - Normalize ftp_import and server_dir src entries into simple path entires + with the relevant paths resolved and permissions / configuration checked. + - Check for file:// URLs in items src of "url" and convert them into path + src items - after verifying path pastes are allowed and user is admin. + - Check for valid URLs to be fetched for http and https entries. + - Based on Galaxy configuration and upload types set purge_source and in_place + as needed for each upload. + """ + 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'") + if destination_type not in VALID_DESTINATION_TYPES: + template = "Invalid target destination type [%s] encountered, must be one of %s" + 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") + description = destination.get("description", "") + synopsis = destination.get("synopsis", "") + library = trans.app.library_manager.create( + trans, library_name, description=description, synopsis=synopsis + ) + destination["type"] = "library_folder" + for key in ["name", "description", "synopsis"]: + if key in destination: + del destination[key] + destination["library_folder_id"] = trans.app.security.encode_id(library.root_folder.id) + + # Unlike upload.py we don't transmit or use run_as_real_user in the job - we just make sure + # in_place and purge_source are set on the individual upload fetch sources as needed based + # on this. + run_as_real_user = trans.app.config.external_chown_script is not None # See comment in upload.py + purge_ftp_source = getattr(trans.app.config, 'ftp_upload_purge', True) and not run_as_real_user + + payload["check_content"] = trans.app.config.check_upload_content + + def check_src(item): + # Normalize file:// URLs into paths. + if item["src"] == "url" and item["url"].startswith("file://"): + item["src"] = "path" + item["path"] = item["url"][len("file://"):] + del item["path"] + + if "in_place" in item: + raise RequestParameterInvalidException("in_place cannot be set in the upload request") + + src = item["src"] + + # Check link_data_only can only be set for certain src types and certain elements_from types. + _handle_invalid_link_data_only_elements_type(item) + if src not in ["path", "server_dir"]: + _handle_invalid_link_data_only_type(item) + elements_from = item.get("elements_from", None) + if elements_from and elements_from not in ELEMENTS_FROM_TYPE: + raise RequestParameterInvalidException("Invalid elements_from/items_from found in request") + + if src == "path" or (src == "url" and item["url"].startswith("file:")): + # Validate is admin, leave alone. + validate_path_upload(trans) + elif src == "server_dir": + # Validate and replace with path definition. + server_dir = item["server_dir"] + full_path, _ = validate_server_directory_upload(trans, server_dir) + item["src"] = "path" + item["path"] = full_path + elif src == "ftp_import": + ftp_path = item["ftp_path"] + full_path = None + + # It'd be nice if this can be de-duplicated with what is in parameters/grouping.py. + user_ftp_dir = trans.user_ftp_dir + is_directory = False + + assert not os.path.islink(user_ftp_dir), "User FTP directory cannot be a symbolic link" + for (dirpath, dirnames, filenames) in os.walk(user_ftp_dir): + for filename in filenames: + if ftp_path == filename: + path = relpath(os.path.join(dirpath, filename), user_ftp_dir) + if not os.path.islink(os.path.join(dirpath, filename)): + full_path = os.path.abspath(os.path.join(user_ftp_dir, path)) + break + + for dirname in dirnames: + if ftp_path == dirname: + path = relpath(os.path.join(dirpath, dirname), user_ftp_dir) + if not os.path.islink(os.path.join(dirpath, dirname)): + full_path = os.path.abspath(os.path.join(user_ftp_dir, path)) + is_directory = True + break + + if is_directory: + # If the target is a directory - make sure no files under it are symbolic links + for (dirpath, dirnames, filenames) in os.walk(full_path): + for filename in filenames: + if ftp_path == filename: + path = relpath(os.path.join(dirpath, filename), full_path) + if not os.path.islink(os.path.join(dirpath, filename)): + full_path = False + break + + for dirname in dirnames: + if ftp_path == dirname: + path = relpath(os.path.join(dirpath, filename), full_path) + if not os.path.islink(os.path.join(dirpath, filename)): + full_path = False + break + + if not full_path: + raise RequestParameterInvalidException("Failed to find referenced ftp_path or symbolic link was enountered") + + item["src"] = "path" + item["path"] = full_path + item["purge_source"] = purge_ftp_source + elif src == "url": + url = item["url"] + looks_like_url = False + for url_prefix in ["http://", "https://", "ftp://", "ftps://"]: + if url.startswith(url_prefix): + looks_like_url = True + break + + if not looks_like_url: + raise RequestParameterInvalidException("Invalid URL [%s] found in src definition." % url) + + validate_url(url, trans.app.config.fetch_url_whitelist_ips) + item["in_place"] = run_as_real_user + elif src == "files": + item["in_place"] = run_as_real_user + + # Small disagreement with traditional uploads - we purge less by default since whether purging + # happens varies based on upload options in non-obvious ways. + # https://github.com/galaxyproject/galaxy/issues/5361 + if "purge_source" not in item: + item["purge_source"] = False + + _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: + raise RequestParameterInvalidException("link_data_only is invalid for src type [%s]" % item.get("src")) + + +def _handle_invalid_link_data_only_elements_type(item): + link_data_only = item.get("link_data_only", False) + if link_data_only and item.get("elements_from", False) in ELEMENTS_FROM_TRANSIENT_TYPES: + 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: + _for_each_src(f, item) + if isinstance(obj, dict): + if "src" in obj: + f(obj) + for key, value in obj.items(): + _for_each_src(f, value) diff --git a/lib/galaxy/webapps/galaxy/api/tools.py b/lib/galaxy/webapps/galaxy/api/tools.py index fa0a0804b7f..08576ee812d 100644 --- a/lib/galaxy/webapps/galaxy/api/tools.py +++ b/lib/galaxy/webapps/galaxy/api/tools.py @@ -15,9 +15,14 @@ from galaxy.web import _future_expose_api_anonymous_and_sessionless as expose_ap from galaxy.web import _future_expose_api_raw_anonymous_and_sessionless as expose_api_raw_anonymous_and_sessionless from galaxy.web.base.controller import BaseAPIController from galaxy.web.base.controller import UsesVisualizationMixin +from ._fetch_util import validate_and_normalize_targets log = logging.getLogger(__name__) +# Do not allow these tools to be called directly - they (it) enforces extra security and +# provides access via a different API endpoint. +PROTECTED_TOOLS = ["__DATA_FETCH__"] + class ToolsController(BaseAPIController, UsesVisualizationMixin): """ @@ -361,12 +366,52 @@ class ToolsController(BaseAPIController, UsesVisualizationMixin): trans.response.headers["Content-Disposition"] = 'attachment; filename="%s.tgz"' % (id) return download_file + @expose_api_anonymous + def fetch(self, trans, payload, **kwd): + """Adapt clean API to tool-constrained API. + """ + log.info("Keywords are %s" % payload) + request_version = '1' + history_id = payload.pop("history_id") + clean_payload = {} + files_payload = {} + for key, value in payload.items(): + if key == "key": + continue + if key.startswith('files_') or key.startswith('__files_'): + files_payload[key] = value + continue + clean_payload[key] = value + log.info("payload %s" % clean_payload) + validate_and_normalize_targets(trans, clean_payload) + clean_payload["check_content"] = trans.app.config.check_upload_content + request = dumps(clean_payload) + log.info(request) + create_payload = { + 'tool_id': "__DATA_FETCH__", + 'history_id': history_id, + 'inputs': { + 'request_version': request_version, + 'request_json': request, + }, + } + create_payload.update(files_payload) + return self._create(trans, create_payload, **kwd) + @expose_api_anonymous def create(self, trans, payload, **kwd): """ POST /api/tools Executes tool using specified inputs and returns tool's outputs. """ + tool_id = payload.get("tool_id") + if tool_id in PROTECTED_TOOLS: + raise exceptions.RequestParameterInvalidException("Cannot execute tool [%s] directly, must use alternative endpoint." % tool_id) + if tool_id is None: + raise exceptions.RequestParameterInvalidException("Must specify a valid tool_id to use this endpoint.") + return self._create(trans, payload, **kwd) + + def _create(self, trans, payload, **kwd): # HACK: for now, if action is rerun, rerun tool. action = payload.get('action', None) if action == 'rerun': diff --git a/lib/galaxy/webapps/galaxy/buildapp.py b/lib/galaxy/webapps/galaxy/buildapp.py index f5ec41c4adf..eda00a26ee4 100644 --- a/lib/galaxy/webapps/galaxy/buildapp.py +++ b/lib/galaxy/webapps/galaxy/buildapp.py @@ -271,6 +271,7 @@ def populate_api_routes(webapp, app): # ====== TOOLS API ====== # ======================= + webapp.mapper.connect('/api/tools/fetch', action='fetch', controller='tools', conditions=dict(method=["POST"])) webapp.mapper.connect('/api/tools/all_requirements', action='all_requirements', controller="tools") webapp.mapper.connect('/api/tools/{id:.+?}/build', action='build', controller="tools") webapp.mapper.connect('/api/tools/{id:.+?}/reload', action='reload', controller="tools") diff --git a/scripts/api/fetch_to_library.py b/scripts/api/fetch_to_library.py new file mode 100644 index 00000000000..6c497bcb402 --- /dev/null +++ b/scripts/api/fetch_to_library.py @@ -0,0 +1,33 @@ +import argparse +import json + +import requests +import yaml + + +def main(): + parser = argparse.ArgumentParser(description='Upload a directory into a data library') + parser.add_argument("-u", "--url", dest="url", required=True, help="Galaxy URL") + parser.add_argument("-a", "--api", dest="api_key", required=True, help="API Key") + parser.add_argument('target', metavar='FILE', type=str, + help='file describing data library to fetch') + args = parser.parse_args() + with open(args.target, "r") as f: + target = yaml.load(f) + + histories_url = args.url + "/api/histories" + new_history_response = requests.post(histories_url, data={'key': args.api_key}) + + fetch_url = args.url + '/api/tools/fetch' + payload = { + 'key': args.api_key, + 'targets': json.dumps([target]), + 'history_id': new_history_response.json()["id"] + } + + response = requests.post(fetch_url, data=payload) + print(response.content) + + +if __name__ == '__main__': + main() diff --git a/scripts/api/fetch_to_library_example.yml b/scripts/api/fetch_to_library_example.yml new file mode 100644 index 00000000000..44bc35ef43b --- /dev/null +++ b/scripts/api/fetch_to_library_example.yml @@ -0,0 +1,42 @@ +destination: + type: library + name: Training Material + description: Data for selected tutorials from https://training.galaxyproject.org. +items: + - name: Quality Control + description: | + Data for sequence quality control tutorial at http://galaxyproject.github.io/training-material/topics/sequence-analysis/tutorials/quality-control/tutorial.html. + + 10.5281/zenodo.61771 + items: + - src: url + url: https://zenodo.org/record/61771/files/GSM461178_untreat_paired_subset_1.fastq + name: GSM461178_untreat_paired_subset_1 + ext: fastqsanger + info: Untreated subseq of GSM461178 from 10.1186/s12864-017-3692-8 + - src: url + url: https://zenodo.org/record/61771/files/GSM461182_untreat_single_subset.fastq + name: GSM461182_untreat_single_subset + ext: fastqsanger + info: Untreated subseq of GSM461182 from 10.1186/s12864-017-3692-8 + - name: Small RNA-Seq + description: | + Data for small RNA-seq tutorial available at http://galaxyproject.github.io/training-material/topics/transcriptomics/tutorials/srna/tutorial.html + + 10.5281/zenodo.826906 + items: + - src: url + url: https://zenodo.org/record/826906/files/Symp_RNAi_sRNA-seq_rep1_downsampled.fastqsanger.gz + name: Symp RNAi sRNA Rep1 + ext: fastqsanger.gz + info: Downsample rep1 from 10.1186/s12864-017-3692-8 + - src: url + url: https://zenodo.org/record/826906/files/Symp_RNAi_sRNA-seq_rep2_downsampled.fastqsanger.gz + name: Symp RNAi sRNA Rep2 + ext: fastqsanger.gz + info: Downsample rep2 from 10.1186/s12864-017-3692-8 + - src: url + url: https://zenodo.org/record/826906/files/Symp_RNAi_sRNA-seq_rep3_downsampled.fastqsanger.gz + name: Symp RNAi sRNA Rep3 + ext: fastqsanger.gz + info: Downsample rep3 from 10.1186/s12864-017-3692-8 diff --git a/test-data/example-bag.zip b/test-data/example-bag.zip new file mode 100644 index 0000000000000000000000000000000000000000..ee4de52c265be0c5ddfe54ade610957580115dac GIT binary patch literal 2966 zcmWIWW@h1H00G^m9&a!MN{BGXFqEVgm*^%Xrt7AqmLzBBW|Wi^=!b@IGBBsQ8$>Ph zG>9s#;AUWC`Nqh=z#;mYgs!XZ3B?^0f-fFm`Zt`8A z>1ZbBp7rKCq#(fdpw|B&E8D7L;)ZH=W_MR~s?QMXl4xXGwsP;@9}&Ce%>4g9t)Fe( z!ljIh+vZ-Jr?NV+JHQ+@O#^gzJP~5vR?8p#Z|NxI-ed zgitUzC8m3p=!T^h6=&w>St%IkS(v7HxNQv;}(335b>P85!j2=;G^&;>}FV z*p@wglDENthvh)nHl|ez8)f)Dsyn{^+WgMZs?#I zn&~Pev>=Ak0n7%S;31$=noE6!&7uF}kM> z`I-%Q7!J%0|IipDbGz`GmqFP>_ClHPwOS4n7M@Rk`Br?Tl(JE<=9yU=SQvU)v%PeZ zUsRvgy+8e@q3!d9f1}(FO*v@T^Te|`R83Kk;o}LpxAs>|_UHZH`15yk`}v(~@3Nx$ zVE(J~2^~P!uK{9pLOv+YNHj7vBj|}U)_jM6p181g`3qK!JehhUrHd0=r7a4396FR< z9Qc@+Q2F@q&7E0oN;^G*Ro5Ii+I#4r>cmM7iwx3&^Cvbf_?x#kbRbe=f3^ zH!bsb4l~oVyhxKwcAvXJC2rFumE^G;x)7<|xqy>>Wy0xeM`yKrE&OR8w!ZA9ZG3jj z``TT5uOI(_=25AY$@%kuo_vkRqa}&yq{No*S;K8ih8(RA|E*0jTcAE+LgP~g#<@|&pGu@?0@;X;NXNX1&_rZW-snEx@pSYJ(gCrTjqRC?DO=0o)h>IJB$KKEZ)tD z4V<-c@vgep$@4N##7)jVnXu&41_|cU*lOQNy2-!(`gHTgwy(Xq*NMUH+6SeH4=zk9 zu72_PM8)~j$rm(dEGiUV<#_i1%VFkIPRv_x{8PSH?)pmo7g|bHUa&4=cI05WaA)}i);Eug4GLru*t{E`I7t8a5vN&r@9VBz!Dp^6;7HVT zRN3)%sxXWDE53$Xc_us7KL57y<`R*{30H4jQ+@pV^S2lC*FXOL<3`_G=HhPuck^!S zSXHFUvrIKFE#sWw>^Q|MD`qr@Z%|$J(|fmU{;3P9;!ga13gJGIA_1%SOx)#|QlusB z)9?6Z;*{Xs)@6#9mwpetZhrOLs;h1*85Fl(ni6&~=DsX2xd=!u_cA9I`DDK+CYXA7~4#)+1~MX61z(M9*+o0jfu@Sb67|*c+&0%2R0zxKWSkt@^7^NRoTEWf0$nuSm zfq_K?s5k&_ggBZJe8@&1%;5u?LC~B=amj9Opy?nijA2e%W=^Ux*ij(E1TY;{T$%*5 z5TuL?;XV$a`xFeU%-Dctf-s8vKxz%aMt%VL^CJ+W`7=NdcP1$Rx*%D*_~dZUg}VhPRF&8YwzhA<=>6L}cSIBLmquXJq4` zQ3Es+6q*FgM2#4TnZTsLu%xjY!%S#25jGezW{?ejjBGG8o`7b8;t7vqG2;i>OhcfX g;o%7~6Bt`SGeNP1VJ0gmiWyjea0yUXJ;=8V0D^6pPXGV_ literal 0 HcmV?d00001 diff --git a/test/api/test_dataset_collections.py b/test/api/test_dataset_collections.py index ca8950d8bb0..87752d2b909 100644 --- a/test/api/test_dataset_collections.py +++ b/test/api/test_dataset_collections.py @@ -188,6 +188,78 @@ class DatasetCollectionApiTestCase(api.ApiTestCase): create_response = self._post("dataset_collections", payload) self._assert_status_code_is(create_response, 400) + def test_upload_collection(self): + elements = [{"src": "files", "dbkey": "hg19", "info": "my cool bed"}] + targets = [{ + "destination": {"type": "hdca"}, + "elements": elements, + "collection_type": "list", + "name": "Test upload", + }] + payload = { + "history_id": self.history_id, + "targets": json.dumps(targets), + "__files": {"files_0|file_data": open(self.test_data_resolver.get_filename("4.bed"))}, + } + self.dataset_populator.fetch(payload) + hdca = self._assert_one_collection_created_in_history() + self.assertEquals(hdca["name"], "Test upload") + assert len(hdca["elements"]) == 1, hdca + element0 = hdca["elements"][0] + assert element0["element_identifier"] == "4.bed" + assert element0["object"]["file_size"] == 61 + + def test_upload_nested(self): + elements = [{"name": "samp1", "elements": [{"src": "files", "dbkey": "hg19", "info": "my cool bed"}]}] + targets = [{ + "destination": {"type": "hdca"}, + "elements": elements, + "collection_type": "list:list", + "name": "Test upload", + }] + payload = { + "history_id": self.history_id, + "targets": json.dumps(targets), + "__files": {"files_0|file_data": open(self.test_data_resolver.get_filename("4.bed"))}, + } + self.dataset_populator.fetch(payload) + hdca = self._assert_one_collection_created_in_history() + self.assertEquals(hdca["name"], "Test upload") + assert len(hdca["elements"]) == 1, hdca + element0 = hdca["elements"][0] + assert element0["element_identifier"] == "samp1" + + def test_upload_collection_from_url(self): + elements = [{"src": "url", "url": "https://raw.githubusercontent.com/galaxyproject/galaxy/dev/test-data/4.bed", "info": "my cool bed"}] + targets = [{ + "destination": {"type": "hdca"}, + "elements": elements, + "collection_type": "list", + }] + payload = { + "history_id": self.history_id, + "targets": json.dumps(targets), + "__files": {"files_0|file_data": open(self.test_data_resolver.get_filename("4.bed"))}, + } + self.dataset_populator.fetch(payload) + hdca = self._assert_one_collection_created_in_history() + assert len(hdca["elements"]) == 1, hdca + element0 = hdca["elements"][0] + assert element0["element_identifier"] == "4.bed" + assert element0["object"]["file_size"] == 61 + + def _assert_one_collection_created_in_history(self): + contents_response = self._get("histories/%s/contents/dataset_collections" % self.history_id) + self._assert_status_code_is(contents_response, 200) + contents = contents_response.json() + assert len(contents) == 1 + hdca = contents[0] + assert hdca["history_content_type"] == "dataset_collection" + hdca_id = hdca["id"] + collection_response = self._get("histories/%s/contents/dataset_collections/%s" % (self.history_id, hdca_id)) + self._assert_status_code_is(collection_response, 200) + return collection_response.json() + def _check_create_response(self, create_response): self._assert_status_code_is(create_response, 200) dataset_collection = create_response.json() diff --git a/test/api/test_libraries.py b/test/api/test_libraries.py index 2a715f50fcb..ae0303a8074 100644 --- a/test/api/test_libraries.py +++ b/test/api/test_libraries.py @@ -1,3 +1,5 @@ +import json + from base import api from base.populators import ( DatasetCollectionPopulator, @@ -95,6 +97,94 @@ class LibrariesApiTestCase(api.ApiTestCase, TestsDatasets): assert library_dataset["peek"].find("create_test") >= 0 assert library_dataset["file_ext"] == "txt", library_dataset["file_ext"] + def test_fetch_upload_to_folder(self): + history_id, library, destination = self._setup_fetch_to_folder("flat_zip") + items = [{"src": "files", "dbkey": "hg19", "info": "my cool bed"}] + targets = [{ + "destination": destination, + "items": items + }] + payload = { + "history_id": history_id, # TODO: Shouldn't be needed :( + "targets": json.dumps(targets), + "__files": {"files_0|file_data": open(self.test_data_resolver.get_filename("4.bed"))}, + } + self.dataset_populator.fetch(payload) + dataset = self.library_populator.get_library_contents_with_path(library["id"], "/4.bed") + assert dataset["file_size"] == 61, dataset + assert dataset["genome_build"] == "hg19", dataset + assert dataset["misc_info"] == "my cool bed", dataset + assert dataset["file_ext"] == "bed", dataset + + def test_fetch_zip_to_folder(self): + history_id, library, destination = self._setup_fetch_to_folder("flat_zip") + bed_test_data_path = self.test_data_resolver.get_filename("4.bed.zip") + targets = [{ + "destination": destination, + "items_from": "archive", "src": "files", + }] + payload = { + "history_id": history_id, # TODO: Shouldn't be needed :( + "targets": json.dumps(targets), + "__files": {"files_0|file_data": open(bed_test_data_path)} + } + self.dataset_populator.fetch(payload) + dataset = self.library_populator.get_library_contents_with_path(library["id"], "/4.bed") + assert dataset["file_size"] == 61, dataset + + def test_fetch_single_url_to_folder(self): + history_id, library, destination = self._setup_fetch_to_folder("single_url") + items = [{"src": "url", "url": "https://raw.githubusercontent.com/galaxyproject/galaxy/dev/test-data/4.bed"}] + targets = [{ + "destination": destination, + "items": items + }] + payload = { + "history_id": history_id, # TODO: Shouldn't be needed :( + "targets": json.dumps(targets), + } + self.dataset_populator.fetch(payload) + dataset = self.library_populator.get_library_contents_with_path(library["id"], "/4.bed") + assert dataset["file_size"] == 61, dataset + + def test_fetch_url_archive_to_folder(self): + history_id, library, destination = self._setup_fetch_to_folder("single_url") + targets = [{ + "destination": destination, + "items_from": "archive", + "src": "url", + "url": "https://raw.githubusercontent.com/galaxyproject/galaxy/dev/test-data/4.bed.zip", + }] + payload = { + "history_id": history_id, # TODO: Shouldn't be needed :( + "targets": json.dumps(targets), + } + self.dataset_populator.fetch(payload) + dataset = self.library_populator.get_library_contents_with_path(library["id"], "/4.bed") + assert dataset["file_size"] == 61, dataset + + def test_fetch_bagit_archive_to_folder(self): + history_id, library, destination = self._setup_fetch_to_folder("bagit_archive") + example_bag_path = self.test_data_resolver.get_filename("example-bag.zip") + targets = [{ + "destination": destination, + "items_from": "bagit_archive", "src": "files", + }] + payload = { + "history_id": history_id, # TODO: Shouldn't be needed :( + "targets": json.dumps(targets), + "__files": {"files_0|file_data": open(example_bag_path)}, + } + self.dataset_populator.fetch(payload) + dataset = self.library_populator.get_library_contents_with_path(library["id"], "/README.txt") + assert dataset["file_size"] == 66, dataset + + dataset = self.library_populator.get_library_contents_with_path(library["id"], "/bdbag-profile.json") + assert dataset["file_size"] == 723, dataset + + def _setup_fetch_to_folder(self, test_name): + return self.library_populator.setup_fetch_to_folder(test_name) + def test_create_dataset_in_folder(self): library = self.library_populator.new_private_library("ForCreateDatasets") folder_response = self._create_folder(library) diff --git a/test/base/integration_util.py b/test/base/integration_util.py index 363b19a32bf..540d6bf1c10 100644 --- a/test/base/integration_util.py +++ b/test/base/integration_util.py @@ -8,6 +8,7 @@ import os from unittest import skip, TestCase from galaxy.tools.deps.commands import which +from galaxy.tools.verify.test_data import TestDataResolver from .api import UsesApiTestCaseMixin from .driver_util import GalaxyTestDriver @@ -56,6 +57,7 @@ class IntegrationTestCase(TestCase, UsesApiTestCaseMixin): cls._app_available = False def setUp(self): + self.test_data_resolver = TestDataResolver() # Setup attributes needed for API testing... server_wrapper = self._test_driver.server_wrappers[0] host = server_wrapper.host diff --git a/test/base/populators.py b/test/base/populators.py index 9bff5a08e9f..625052071b4 100644 --- a/test/base/populators.py +++ b/test/base/populators.py @@ -148,14 +148,30 @@ class BaseDatasetPopulator(object): self.wait_for_tool_run(history_id, run_response, assert_ok=kwds.get('assert_ok', True)) return run_response + def fetch(self, payload, assert_ok=True, timeout=DEFAULT_TIMEOUT): + tool_response = self._post("tools/fetch", data=payload) + if assert_ok: + job = self.check_run(tool_response) + self.wait_for_job(job["id"], timeout=timeout) + + job = tool_response.json()["jobs"][0] + details = self.get_job_details(job["id"]).json() + assert details["state"] == "ok", details + + return tool_response + def wait_for_tool_run(self, history_id, run_response, timeout=DEFAULT_TIMEOUT, assert_ok=True): - run = run_response.json() - assert run_response.status_code == 200, run - job = run["jobs"][0] + job = self.check_run(run_response) self.wait_for_job(job["id"], timeout=timeout) self.wait_for_history(history_id, assert_ok=assert_ok, timeout=timeout) return run_response + def check_run(self, run_response): + run = run_response.json() + assert run_response.status_code == 200, run + job = run["jobs"][0] + return job + def wait_for_history(self, history_id, assert_ok=False, timeout=DEFAULT_TIMEOUT): try: return wait_on_state(lambda: self._get("histories/%s" % history_id), assert_ok=assert_ok, timeout=timeout) @@ -266,8 +282,8 @@ class BaseDatasetPopulator(object): else: return tool_response - def tools_post(self, payload): - tool_response = self._post("tools", data=payload) + def tools_post(self, payload, url="tools"): + tool_response = self._post(url, data=payload) return tool_response def get_history_dataset_content(self, history_id, wait=True, filename=None, **kwds): @@ -463,6 +479,11 @@ class LibraryPopulator(object): def __init__(self, galaxy_interactor): self.galaxy_interactor = galaxy_interactor + self.dataset_populator = DatasetPopulator(galaxy_interactor) + + def get_libraries(self): + get_response = self.galaxy_interactor.get("libraries") + return get_response.json() def new_private_library(self, name): library = self.new_library(name) @@ -563,6 +584,24 @@ class LibraryPopulator(object): return library, library_dataset + def get_library_contents_with_path(self, library_id, path): + all_contents_response = self.galaxy_interactor.get("libraries/%s/contents" % library_id) + api_asserts.assert_status_code_is(all_contents_response, 200) + all_contents = all_contents_response.json() + matching = [c for c in all_contents if c["name"] == path] + if len(matching) == 0: + raise Exception("Failed to find library contents with path [%s], contents are %s" % (path, all_contents)) + get_response = self.galaxy_interactor.get(matching[0]["url"]) + api_asserts.assert_status_code_is(get_response, 200) + return get_response.json() + + def setup_fetch_to_folder(self, test_name): + history_id = self.dataset_populator.new_history() + library = self.new_private_library(test_name) + folder_id = library["root_folder_id"][1:] + destination = {"type": "library_folder", "library_folder_id": folder_id} + return history_id, library, destination + class BaseDatasetCollectionPopulator(object): diff --git a/test/integration/test_upload_configuration_options.py b/test/integration/test_upload_configuration_options.py index 441fd32f3d1..8a6d00f0afb 100644 --- a/test/integration/test_upload_configuration_options.py +++ b/test/integration/test_upload_configuration_options.py @@ -19,6 +19,7 @@ These options include: framework but tested here for FTP uploads. """ +import json import os import re import shutil @@ -54,6 +55,42 @@ class BaseUploadContentConfigurationTestCase(integration_util.IntegrationTestCas self.library_populator = LibraryPopulator(self.galaxy_interactor) self.history_id = self.dataset_populator.new_history() + def fetch_target(self, target, assert_ok=False, attach_test_file=False): + payload = { + "history_id": self.history_id, + "targets": json.dumps([target]), + } + if attach_test_file: + payload["__files"] = {"files_0|file_data": open(self.test_data_resolver.get_filename("4.bed"))} + + response = self.dataset_populator.fetch(payload, assert_ok=assert_ok) + return response + + +class InvalidFetchRequestsTestCase(BaseUploadContentConfigurationTestCase): + + def test_in_place_not_allowed(self): + elements = [{"src": "files", "in_place": False}] + target = { + "destination": {"type": "hdca"}, + "elements": elements, + "collection_type": "list", + } + response = self.fetch_target(target, attach_test_file=True) + self._assert_status_code_is(response, 400) + assert 'in_place' in response.json()["err_msg"] + + def test_files_not_attached(self): + elements = [{"src": "files"}] + target = { + "destination": {"type": "hdca"}, + "elements": elements, + "collection_type": "list", + } + response = self.fetch_target(target) + self._assert_status_code_is(response, 400) + assert 'Failed to find uploaded file matching target' in response.json()["err_msg"] + class NonAdminsCannotPasteFilePathTestCase(BaseUploadContentConfigurationTestCase): @@ -94,6 +131,26 @@ class NonAdminsCannotPasteFilePathTestCase(BaseUploadContentConfigurationTestCas response = self.library_populator.raw_library_contents_create(library["id"], payload, files=files) assert response.status_code == 403, response.json() + def test_disallowed_for_fetch(self): + elements = [{"src": "path", "path": "%s/1.txt" % TEST_DATA_DIRECTORY}] + target = { + "destination": {"type": "hdca"}, + "elements": elements, + "collection_type": "list", + } + response = self.fetch_target(target) + self._assert_status_code_is(response, 403) + + def test_disallowed_for_fetch_urls(self): + elements = [{"src": "url", "url": "file://%s/1.txt" % TEST_DATA_DIRECTORY}] + target = { + "destination": {"type": "hdca"}, + "elements": elements, + "collection_type": "list", + } + response = self.fetch_target(target) + self._assert_status_code_is(response, 403) + class AdminsCanPasteFilePathsTestCase(BaseUploadContentConfigurationTestCase): @@ -118,6 +175,26 @@ class AdminsCanPasteFilePathsTestCase(BaseUploadContentConfigurationTestCase): # Was 403 for non-admin above. assert response.status_code == 200 + def test_admin_fetch(self): + elements = [{"src": "path", "path": "%s/1.txt" % TEST_DATA_DIRECTORY}] + target = { + "destination": {"type": "hdca"}, + "elements": elements, + "collection_type": "list", + } + response = self.fetch_target(target) + self._assert_status_code_is(response, 200) + + def test_admin_fetch_file_url(self): + elements = [{"src": "url", "url": "file://%s/1.txt" % TEST_DATA_DIRECTORY}] + target = { + "destination": {"type": "hdca"}, + "elements": elements, + "collection_type": "list", + } + response = self.fetch_target(target) + self._assert_status_code_is(response, 200) + class DefaultBinaryContentFiltersTestCase(BaseUploadContentConfigurationTestCase): @@ -212,6 +289,16 @@ class LocalAddressWhitelisting(BaseUploadContentConfigurationTestCase): # the newer API decorator that handles those details. assert create_response.status_code >= 400 + def test_blocked_url_for_fetch(self): + elements = [{"src": "url", "url": "http://localhost"}] + target = { + "destination": {"type": "hdca"}, + "elements": elements, + "collection_type": "list", + } + response = self.fetch_target(target) + self._assert_status_code_is(response, 403) + class BaseFtpUploadConfigurationTestCase(BaseUploadContentConfigurationTestCase): @@ -251,6 +338,9 @@ class BaseFtpUploadConfigurationTestCase(BaseUploadContentConfigurationTestCase) if not os.path.exists(path): os.makedirs(path) + def _get_user_ftp_path(self): + return os.path.join(self.ftp_dir(), TEST_USER) + class SimpleFtpUploadConfigurationTestCase(BaseFtpUploadConfigurationTestCase): @@ -269,6 +359,24 @@ class SimpleFtpUploadConfigurationTestCase(BaseFtpUploadConfigurationTestCase): self._check_content(dataset, content) assert not os.path.exists(ftp_path) + def test_ftp_fetch(self): + content = "hello world\n" + ftp_path = self._write_ftp_file(content) + ftp_files = self.dataset_populator.get_remote_files() + assert len(ftp_files) == 1, ftp_files + assert ftp_files[0]["path"] == "test" + assert os.path.exists(ftp_path) + elements = [{"src": "ftp_import", "ftp_path": ftp_files[0]["path"]}] + target = { + "destination": {"type": "hdca"}, + "elements": elements, + "collection_type": "list", + } + response = self.fetch_target(target) + self._assert_status_code_is(response, 200) + dataset = self.dataset_populator.get_history_dataset_details(self.history_id, hid=2) + self._check_content(dataset, content) + class ExplicitEmailAsIdentifierFtpUploadConfigurationTestCase(SimpleFtpUploadConfigurationTestCase): @@ -320,6 +428,50 @@ class DisableFtpPurgeUploadConfigurationTestCase(BaseFtpUploadConfigurationTestC assert os.path.exists(ftp_path) +class AdvancedFtpUploadFetchTestCase(BaseFtpUploadConfigurationTestCase): + + def test_fetch_ftp_directory(self): + dir_path = self._get_user_ftp_path() + self._write_ftp_file(os.path.join(dir_path, "subdir"), "content 1", filename="1") + self._write_ftp_file(os.path.join(dir_path, "subdir"), "content 22", filename="2") + self._write_ftp_file(os.path.join(dir_path, "subdir"), "content 333", filename="3") + target = { + "destination": {"type": "hdca"}, + "elements_from": "directory", + "src": "ftp_import", + "ftp_path": "subdir", + "collection_type": "list", + } + response = self.fetch_target(target) + self._assert_status_code_is(response, 200) + hdca = self.dataset_populator.get_history_collection_details(self.history_id, hid=1) + assert len(hdca["elements"]) == 3, hdca + element0 = hdca["elements"][0] + assert element0["element_identifier"] == "1" + assert element0["object"]["file_size"] == 9 + + def test_fetch_nested_elements_from(self): + dir_path = self._get_user_ftp_path() + self._write_ftp_file(os.path.join(dir_path, "subdir1"), "content 1", filename="1") + self._write_ftp_file(os.path.join(dir_path, "subdir1"), "content 22", filename="2") + self._write_ftp_file(os.path.join(dir_path, "subdir2"), "content 333", filename="3") + elements = [ + {"name": "subdirel1", "src": "ftp_import", "ftp_path": "subdir1", "elements_from": "directory", "collection_type": "list"}, + {"name": "subdirel2", "src": "ftp_import", "ftp_path": "subdir2", "elements_from": "directory", "collection_type": "list"}, + ] + target = { + "destination": {"type": "hdca"}, + "elements": elements, + "collection_type": "list:list", + } + response = self.fetch_target(target) + self._assert_status_code_is(response, 200) + hdca = self.dataset_populator.get_history_collection_details(self.history_id, hid=1) + assert len(hdca["elements"]) == 2, hdca + element0 = hdca["elements"][0] + assert element0["element_identifier"] == "subdirel1" + + class UploadOptionsFtpUploadConfigurationTestCase(BaseFtpUploadConfigurationTestCase): def test_upload_api_option_space_to_tab(self): @@ -448,6 +600,8 @@ class ServerDirectoryOffByDefaultTestCase(BaseUploadContentConfigurationTestCase class ServerDirectoryValidUsageTestCase(BaseUploadContentConfigurationTestCase): + # This tests the library contents API - I think equivalent functionality is available via library datasets API + # and should also be tested. require_admin_user = True @@ -483,3 +637,84 @@ class ServerDirectoryRestrictedToAdminsUsageTestCase(BaseUploadContentConfigurat payload, files = self.library_populator.create_dataset_request(library, upload_option="upload_directory", server_dir="library") response = self.library_populator.raw_library_contents_create(library["id"], payload, files=files) assert response.status_code == 403, response.json() + + +class FetchByPathTestCase(BaseUploadContentConfigurationTestCase): + + require_admin_user = True + + @classmethod + def handle_galaxy_config_kwds(cls, config): + config["allow_path_paste"] = True + + def test_fetch_path_to_folder(self): + history_id, library, destination = self.library_populator.setup_fetch_to_folder("simple_fetch") + bed_test_data_path = self.test_data_resolver.get_filename("4.bed") + items = [{"src": "path", "path": bed_test_data_path, "info": "my cool bed"}] + targets = [{ + "destination": destination, + "items": items + }] + payload = { + "history_id": history_id, # TODO: Shouldn't be needed :( + "targets": json.dumps(targets), + } + self.dataset_populator.fetch(payload) + dataset = self.library_populator.get_library_contents_with_path(library["id"], "/4.bed") + assert dataset["file_size"] == 61, dataset + + def test_fetch_link_data_only(self): + history_id, library, destination = self.library_populator.setup_fetch_to_folder("fetch_and_link") + bed_test_data_path = self.test_data_resolver.get_filename("4.bed") + items = [{"src": "path", "path": bed_test_data_path, "info": "my cool bed", "link_data_only": True}] + targets = [{ + "destination": destination, + "items": items + }] + payload = { + "history_id": history_id, # TODO: Shouldn't be needed :( + "targets": json.dumps(targets), + } + self.dataset_populator.fetch(payload) + dataset = self.library_populator.get_library_contents_with_path(library["id"], "/4.bed") + assert dataset["file_size"] == 61, dataset + assert dataset["file_name"] == bed_test_data_path, dataset + + def test_fetch_recursive_archive(self): + history_id, library, destination = self.library_populator.setup_fetch_to_folder("recursive_archive") + bed_test_data_path = self.test_data_resolver.get_filename("testdir1.zip") + targets = [{ + "destination": destination, + "items_from": "archive", "src": "path", "path": bed_test_data_path, + }] + payload = { + "history_id": history_id, # TODO: Shouldn't be needed :( + "targets": json.dumps(targets), + } + self.dataset_populator.fetch(payload) + dataset = self.library_populator.get_library_contents_with_path(library["id"], "/file1") + assert dataset["file_size"] == 6, dataset + + dataset = self.library_populator.get_library_contents_with_path(library["id"], "/file2") + assert dataset["file_size"] == 6, dataset + + dataset = self.library_populator.get_library_contents_with_path(library["id"], "/dir1/file3") + assert dataset["file_size"] == 11, dataset + + def test_fetch_recursive_archive_to_library(self): + bed_test_data_path = self.test_data_resolver.get_filename("testdir1.zip") + targets = [{ + "destination": {"type": "library", "name": "My Cool Library"}, + "items_from": "archive", "src": "path", "path": bed_test_data_path, + }] + payload = { + "history_id": self.history_id, # TODO: Shouldn't be needed :( + "targets": json.dumps(targets), + } + self.dataset_populator.fetch(payload) + libraries = self.library_populator.get_libraries() + matching = [l for l in libraries if l["name"] == "My Cool Library"] + assert len(matching) == 1 + library = matching[0] + dataset = self.library_populator.get_library_contents_with_path(library["id"], "/file1") + assert dataset["file_size"] == 6, dataset diff --git a/tools/data_source/upload.py b/tools/data_source/upload.py index fa7f2ade5ba..bdb9ad4139a 100644 --- a/tools/data_source/upload.py +++ b/tools/data_source/upload.py @@ -18,8 +18,12 @@ from six.moves.urllib.request import urlopen from galaxy import util from galaxy.datatypes import sniff -from galaxy.datatypes.binary import Binary from galaxy.datatypes.registry import Registry +from galaxy.datatypes.upload_util import ( + handle_sniffable_binary_check, + handle_unsniffable_binary_check, + UploadProblemException, +) from galaxy.util.checkers import ( check_binary, check_bz2, @@ -36,12 +40,6 @@ else: assert sys.version_info[:2] >= (2, 7) -class UploadProblemException(Exception): - - def __init__(self, message): - self.message = message - - def file_err(msg, dataset, json_file): json_file.write(dumps(dict(type='dataset', ext='data', @@ -123,26 +121,21 @@ def add_file(dataset, registry, json_file, output_path): if dataset.type == 'url': try: - page = urlopen(dataset.path) # page will be .close()ed by sniff methods - temp_name = sniff.stream_to_file(page, prefix='url_paste', source_encoding=util.get_charset_from_http_headers(page.headers)) + dataset.path = sniff.stream_url_to_file(dataset.path) except Exception as e: raise UploadProblemException('Unable to fetch %s\n%s' % (dataset.path, str(e))) - dataset.path = temp_name + # See if we have an empty file if not os.path.exists(dataset.path): raise UploadProblemException('Uploaded temporary file (%s) does not exist.' % dataset.path) + if not os.path.getsize(dataset.path) > 0: raise UploadProblemException('The uploaded file is empty') + # Is dataset content supported sniffable binary? is_binary = check_binary(dataset.path) if is_binary: - # Sniff the data type - guessed_ext = sniff.guess_ext(dataset.path, registry.sniff_order) - # Set data_type only if guessed_ext is a binary datatype - datatype = registry.get_datatype_by_extension(guessed_ext) - if isinstance(datatype, Binary): - data_type = guessed_ext - ext = guessed_ext + data_type, ext = handle_sniffable_binary_check(data_type, ext, dataset.path, registry) if not data_type: root_datatype = registry.get_datatype_by_extension(dataset.file_type) if getattr(root_datatype, 'compressed', False): @@ -265,18 +258,9 @@ def add_file(dataset, registry, json_file, output_path): dataset.name = uncompressed_name data_type = 'zip' if not data_type: - if is_binary or registry.is_extension_unsniffable_binary(dataset.file_type): - # We have a binary dataset, but it is not Bam, Sff or Pdf - data_type = 'binary' - parts = dataset.name.split(".") - if len(parts) > 1: - ext = parts[-1].strip().lower() - is_ext_unsniffable_binary = registry.is_extension_unsniffable_binary(ext) - if check_content and not is_ext_unsniffable_binary: - raise UploadProblemException('The uploaded binary file contains inappropriate content') - elif is_ext_unsniffable_binary and dataset.file_type != ext: - err_msg = "You must manually set the 'File Format' to '%s' when uploading %s files." % (ext, ext) - raise UploadProblemException(err_msg) + data_type, ext = handle_unsniffable_binary_check( + data_type, ext, dataset.path, dataset.name, is_binary, dataset.file_type, check_content, registry + ) if not data_type: # We must have a text file if check_content and check_html(dataset.path): From ce2170a56ede27c9caa2a3c881fa5c582c25b196 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Mon, 26 Feb 2018 09:27:37 -0500 Subject: [PATCH 11/24] Test case for link_to_files during upload. --- test/base/populators.py | 4 + .../test_upload_configuration_options.py | 81 +++++++++++++------ 2 files changed, 62 insertions(+), 23 deletions(-) diff --git a/test/base/populators.py b/test/base/populators.py index 625052071b4..22d1d15d87d 100644 --- a/test/base/populators.py +++ b/test/base/populators.py @@ -535,6 +535,8 @@ class LibraryPopulator(object): "file_type": kwds.get("file_type", "auto"), "db_key": kwds.get("db_key", "?"), } + if kwds.get("link_data"): + create_data["link_data_only"] = "link_to_files" if upload_option == "upload_file": files = { @@ -553,7 +555,9 @@ class LibraryPopulator(object): library = self.new_private_library(name) payload, files = self.create_dataset_request(library, **create_dataset_kwds) dataset = self.raw_library_contents_create(library["id"], payload, files=files).json()[0] + return self.wait_on_library_dataset(library, dataset) + def wait_on_library_dataset(self, library, dataset): def show(): return self.galaxy_interactor.get("libraries/%s/contents/%s" % (library["id"], dataset["id"])) diff --git a/test/integration/test_upload_configuration_options.py b/test/integration/test_upload_configuration_options.py index 8a6d00f0afb..a4d6b8feb23 100644 --- a/test/integration/test_upload_configuration_options.py +++ b/test/integration/test_upload_configuration_options.py @@ -66,6 +66,22 @@ class BaseUploadContentConfigurationTestCase(integration_util.IntegrationTestCas response = self.dataset_populator.fetch(payload, assert_ok=assert_ok) return response + @classmethod + def temp_config_dir(cls, name): + return os.path.join(cls._test_driver.galaxy_test_tmp_dir, name) + + def _write_file(self, dir_path, content, filename="test"): + """Helper for writing ftp/server dir files.""" + self._ensure_directory(dir_path) + path = os.path.join(dir_path, filename) + with open(path, "w") as f: + f.write(content) + return path + + def _ensure_directory(self, path): + if not os.path.exists(path): + os.makedirs(path) + class InvalidFetchRequestsTestCase(BaseUploadContentConfigurationTestCase): @@ -315,7 +331,7 @@ class BaseFtpUploadConfigurationTestCase(BaseUploadContentConfigurationTestCase) @classmethod def ftp_dir(cls): - return os.path.join(cls._test_driver.galaxy_test_tmp_dir, "ftp") + return cls.temp_config_dir("ftp") def _check_content(self, dataset, content, ext="txt"): dataset = self.dataset_populator.get_history_dataset_details(self.history_id, dataset=dataset) @@ -338,9 +354,6 @@ class BaseFtpUploadConfigurationTestCase(BaseUploadContentConfigurationTestCase) if not os.path.exists(path): os.makedirs(path) - def _get_user_ftp_path(self): - return os.path.join(self.ftp_dir(), TEST_USER) - class SimpleFtpUploadConfigurationTestCase(BaseFtpUploadConfigurationTestCase): @@ -432,9 +445,9 @@ class AdvancedFtpUploadFetchTestCase(BaseFtpUploadConfigurationTestCase): def test_fetch_ftp_directory(self): dir_path = self._get_user_ftp_path() - self._write_ftp_file(os.path.join(dir_path, "subdir"), "content 1", filename="1") - self._write_ftp_file(os.path.join(dir_path, "subdir"), "content 22", filename="2") - self._write_ftp_file(os.path.join(dir_path, "subdir"), "content 333", filename="3") + self._write_file(os.path.join(dir_path, "subdir"), "content 1", filename="1") + self._write_file(os.path.join(dir_path, "subdir"), "content 22", filename="2") + self._write_file(os.path.join(dir_path, "subdir"), "content 333", filename="3") target = { "destination": {"type": "hdca"}, "elements_from": "directory", @@ -452,9 +465,9 @@ class AdvancedFtpUploadFetchTestCase(BaseFtpUploadConfigurationTestCase): def test_fetch_nested_elements_from(self): dir_path = self._get_user_ftp_path() - self._write_ftp_file(os.path.join(dir_path, "subdir1"), "content 1", filename="1") - self._write_ftp_file(os.path.join(dir_path, "subdir1"), "content 22", filename="2") - self._write_ftp_file(os.path.join(dir_path, "subdir2"), "content 333", filename="3") + self._write_file(os.path.join(dir_path, "subdir1"), "content 1", filename="1") + self._write_file(os.path.join(dir_path, "subdir1"), "content 22", filename="2") + self._write_file(os.path.join(dir_path, "subdir2"), "content 333", filename="3") elements = [ {"name": "subdirel1", "src": "ftp_import", "ftp_path": "subdir1", "elements_from": "directory", "collection_type": "list"}, {"name": "subdirel2", "src": "ftp_import", "ftp_path": "subdir2", "elements_from": "directory", "collection_type": "list"}, @@ -580,7 +593,7 @@ class UploadOptionsFtpUploadConfigurationTestCase(BaseFtpUploadConfigurationTest shutil.copyfile(input_path, os.path.join(target_dir, test_data_path)) def _write_user_ftp_file(self, path, content): - return self._write_ftp_file(content, filename=path) + return self._write_file(os.path.join(self.ftp_dir(), TEST_USER), content, filename=path) class ServerDirectoryOffByDefaultTestCase(BaseUploadContentConfigurationTestCase): @@ -607,22 +620,44 @@ class ServerDirectoryValidUsageTestCase(BaseUploadContentConfigurationTestCase): @classmethod def handle_galaxy_config_kwds(cls, config): - library_import_dir = os.path.join(cls._test_driver.galaxy_test_tmp_dir, "library_import_dir") - config["library_import_dir"] = library_import_dir - cls.dir_to_import = 'library' - full_dir_path = os.path.join(library_import_dir, cls.dir_to_import) - os.makedirs(full_dir_path) - cls.file_content = "create_test" - with tempfile.NamedTemporaryFile(dir=full_dir_path, delete=False) as fh: - fh.write(cls.file_content) - cls.file_to_import = fh.name + server_dir = cls.server_dir() + os.makedirs(server_dir) + config["library_import_dir"] = server_dir def test_valid_server_dir_uploads_okay(self): - self.library_populator.new_library_dataset("serverdirupload", upload_option="upload_directory", server_dir=self.dir_to_import) + dir_to_import = 'library' + full_dir_path = os.path.join(self.server_dir(), dir_to_import) + os.makedirs(full_dir_path) + file_content = "hello world\n" + with tempfile.NamedTemporaryFile(dir=full_dir_path, delete=False) as fh: + fh.write(file_content) + file_to_import = fh.name + + library_dataset = self.library_populator.new_library_dataset("serverdirupload", upload_option="upload_directory", server_dir=dir_to_import) # Check the file is still there and was not modified - with open(self.file_to_import) as fh: + with open(file_to_import) as fh: read_content = fh.read() - assert read_content == self.file_content + assert read_content == file_content + + assert library_dataset["file_size"] == 12, library_dataset + + def link_data_only(self): + content = "hello world\n" + dir_path = os.path.join(self.server_dir(), "lib1") + file_path = self._write_file(dir_path, content) + library = self.library_populator.new_private_library("serverdirupload") + # upload $GALAXY_ROOT/test-data/library + payload, files = self.library_populator.create_dataset_request(library, upload_option="upload_directory", server_dir="lib1", link_data=True) + response = self.library_populator.raw_library_contents_create(library["id"], payload, files=files) + assert response.status_code == 200, response.json() + dataset = response.json()[0] + ok_dataset = self.library_populator.wait_on_library_dataset(library, dataset) + assert ok_dataset["file_size"] == 12, ok_dataset + assert ok_dataset["file_name"] == file_path, ok_dataset + + @classmethod + def server_dir(cls): + return cls.temp_config_dir("server") class ServerDirectoryRestrictedToAdminsUsageTestCase(BaseUploadContentConfigurationTestCase): From b6f4bffbbb89d3f96f7a292eca38e631cd950a44 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Sat, 6 Jan 2018 14:08:29 -0500 Subject: [PATCH 12/24] Do not allow workflows to run tools that are not workflow-compatible. In the case of data-fetch there is extra validation that is done so this is somewhat important. --- lib/galaxy/workflow/modules.py | 3 +++ 1 file changed, 3 insertions(+) diff --git a/lib/galaxy/workflow/modules.py b/lib/galaxy/workflow/modules.py index a7b4303badc..d090ef51280 100644 --- a/lib/galaxy/workflow/modules.py +++ b/lib/galaxy/workflow/modules.py @@ -825,6 +825,9 @@ class ToolModule(WorkflowModule): invocation = invocation_step.workflow_invocation step = invocation_step.workflow_step tool = trans.app.toolbox.get_tool(step.tool_id, tool_version=step.tool_version) + if not tool.is_workflow_compatible: + message = "Specified tool [%s] in workflow is not workflow-compatible." % tool.id + raise Exception(message) tool_state = step.state # Not strictly needed - but keep Tool state clean by stripping runtime # metadata parameters from it. From 347f1e123bb4956361384935c28714e0734ac916 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Mon, 26 Feb 2018 09:29:26 -0500 Subject: [PATCH 13/24] Don't purge library path pastes and such in upload.py during testing. Concerning that they sometimes will get deleted in production with default settings - see 5361. --- lib/galaxy/actions/library.py | 1 + test/base/driver_util.py | 1 + .../test_upload_configuration_options.py | 46 +++++++++++++------ 3 files changed, 35 insertions(+), 13 deletions(-) diff --git a/lib/galaxy/actions/library.py b/lib/galaxy/actions/library.py index fa653b64eea..14e2eeac907 100644 --- a/lib/galaxy/actions/library.py +++ b/lib/galaxy/actions/library.py @@ -249,6 +249,7 @@ class LibraryActions(object): uploaded_dataset.to_posix_lines = params.get('to_posix_lines', None) uploaded_dataset.space_to_tab = params.get('space_to_tab', None) uploaded_dataset.tag_using_filenames = params.get('tag_using_filenames', True) + uploaded_dataset.purge_source = getattr(trans.app.config, 'ftp_upload_purge', True) if in_folder: uploaded_dataset.in_folder = in_folder uploaded_dataset.data = upload_common.new_upload(trans, 'api', uploaded_dataset, library_bunch) diff --git a/test/base/driver_util.py b/test/base/driver_util.py index 638eb7f9625..7c31e7df7ee 100644 --- a/test/base/driver_util.py +++ b/test/base/driver_util.py @@ -198,6 +198,7 @@ def setup_galaxy_config( enable_beta_tool_formats=True, expose_dataset_path=True, file_path=file_path, + ftp_upload_purge=False, galaxy_data_manager_data_path=galaxy_data_manager_data_path, id_secret='changethisinproductiontoo', job_config_file=job_config_file, diff --git a/test/integration/test_upload_configuration_options.py b/test/integration/test_upload_configuration_options.py index a4d6b8feb23..603e28d3a5e 100644 --- a/test/integration/test_upload_configuration_options.py +++ b/test/integration/test_upload_configuration_options.py @@ -186,10 +186,14 @@ class AdminsCanPasteFilePathsTestCase(BaseUploadContentConfigurationTestCase): def test_admin_path_paste_libraries(self): library = self.library_populator.new_private_library("pathpasteallowedlibraries") - payload, files = self.library_populator.create_dataset_request(library, upload_option="upload_paths", paths="%s/1.txt" % TEST_DATA_DIRECTORY) + path = "%s/1.txt" % TEST_DATA_DIRECTORY + assert os.path.exists(path) + payload, files = self.library_populator.create_dataset_request(library, upload_option="upload_paths", paths=path) response = self.library_populator.raw_library_contents_create(library["id"], payload, files=files) # Was 403 for non-admin above. assert response.status_code == 200 + # Test regression where this was getting deleted in this mode. + assert os.path.exists(path) def test_admin_fetch(self): elements = [{"src": "path", "path": "%s/1.txt" % TEST_DATA_DIRECTORY}] @@ -354,6 +358,22 @@ class BaseFtpUploadConfigurationTestCase(BaseUploadContentConfigurationTestCase) if not os.path.exists(path): os.makedirs(path) + def _run_purgable_upload(self): + # Purge setting is actually used with a fairly specific set of parameters - see: + # https://github.com/galaxyproject/galaxy/issues/5361 + content = "hello world\n" + ftp_path = self._write_ftp_file(content) + ftp_files = self.dataset_populator.get_remote_files() + assert len(ftp_files) == 1 + assert ftp_files[0]["path"] == "test" + assert os.path.exists(ftp_path) + # gotta set to_posix_lines to None currently to force purging of non-binary data. + dataset = self.dataset_populator.new_dataset( + self.history_id, ftp_files="test", file_type="txt", to_posix_lines=None, wait=True + ) + self._check_content(dataset, content) + return ftp_path + class SimpleFtpUploadConfigurationTestCase(BaseFtpUploadConfigurationTestCase): @@ -370,7 +390,6 @@ class SimpleFtpUploadConfigurationTestCase(BaseFtpUploadConfigurationTestCase): self.history_id, ftp_files="test", file_type="txt", to_posix_lines=None, wait=True ) self._check_content(dataset, content) - assert not os.path.exists(ftp_path) def test_ftp_fetch(self): content = "hello world\n" @@ -426,21 +445,22 @@ class DisableFtpPurgeUploadConfigurationTestCase(BaseFtpUploadConfigurationTestC config["ftp_upload_purge"] = "False" def test_ftp_uploads_not_purged(self): - content = "hello world\n" - ftp_path = self._write_ftp_file(content) - ftp_files = self.dataset_populator.get_remote_files() - assert len(ftp_files) == 1 - assert ftp_files[0]["path"] == "test" - assert os.path.exists(ftp_path) - # gotta set to_posix_lines to None currently to force purging of non-binary data. - dataset = self.dataset_populator.new_dataset( - self.history_id, ftp_files="test", file_type="txt", to_posix_lines=None, wait=True - ) - self._check_content(dataset, content) + ftp_path = self._run_purgable_upload() # Purge is disabled, this better still be here. assert os.path.exists(ftp_path) +class EnableFtpPurgeUploadConfigurationTestCase(BaseFtpUploadConfigurationTestCase): + + @classmethod + def handle_extra_ftp_config(cls, config): + config["ftp_upload_purge"] = "True" + + def test_ftp_uploads_not_purged(self): + ftp_path = self._run_purgable_upload() + assert not os.path.exists(ftp_path) + + class AdvancedFtpUploadFetchTestCase(BaseFtpUploadConfigurationTestCase): def test_fetch_ftp_directory(self): From 3bf5990155faea13ceaa7e43b76d61cfb7ec3cd2 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Wed, 10 Jan 2018 21:44:10 -0500 Subject: [PATCH 14/24] Allow uploading individual HDAs via fetch API. --- lib/galaxy/tools/parameters/output_collect.py | 36 +++++++++++++++++++ lib/galaxy/webapps/galaxy/api/_fetch_util.py | 2 +- .../test_upload_configuration_options.py | 21 +++++++++-- 3 files changed, 56 insertions(+), 3 deletions(-) diff --git a/lib/galaxy/tools/parameters/output_collect.py b/lib/galaxy/tools/parameters/output_collect.py index a9e4f102799..938771d515d 100644 --- a/lib/galaxy/tools/parameters/output_collect.py +++ b/lib/galaxy/tools/parameters/output_collect.py @@ -247,6 +247,42 @@ def collect_dynamic_outputs( filenames, ) collection_builder.populate() + elif destination_type == "hdas": + history = job.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, 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 + + dataset = job_context.create_dataset( + ext=ext, + designation=designation, + visible=True, + dbkey=dbkey, + name=name, + filename=discovered_file.path, + info=info, + link_data=link_data + ) + dataset.raw_set_dataset_state('ok') + datasets.append(dataset) + + collect_elements_for_history(elements) + job.history.add_datasets(job_context.sa_session, datasets) for name, has_collection in output_collections.items(): if name not in tool.output_collections: diff --git a/lib/galaxy/webapps/galaxy/api/_fetch_util.py b/lib/galaxy/webapps/galaxy/api/_fetch_util.py index 76bee4efa22..d2b2eb11a68 100644 --- a/lib/galaxy/webapps/galaxy/api/_fetch_util.py +++ b/lib/galaxy/webapps/galaxy/api/_fetch_util.py @@ -15,7 +15,7 @@ from galaxy.util import ( log = logging.getLogger(__name__) -VALID_DESTINATION_TYPES = ["library", "library_folder", "hdca"] +VALID_DESTINATION_TYPES = ["library", "library_folder", "hdca", "hdas"] ELEMENTS_FROM_TYPE = ["archive", "bagit", "bagit_archive", "directory"] # These elements_from cannot be sym linked to because they only exist during upload. ELEMENTS_FROM_TRANSIENT_TYPES = ["archive", "bagit_archive"] diff --git a/test/integration/test_upload_configuration_options.py b/test/integration/test_upload_configuration_options.py index 603e28d3a5e..fe47ec04e71 100644 --- a/test/integration/test_upload_configuration_options.py +++ b/test/integration/test_upload_configuration_options.py @@ -737,10 +737,10 @@ class FetchByPathTestCase(BaseUploadContentConfigurationTestCase): def test_fetch_recursive_archive(self): history_id, library, destination = self.library_populator.setup_fetch_to_folder("recursive_archive") - bed_test_data_path = self.test_data_resolver.get_filename("testdir1.zip") + archive_test_data_path = self.test_data_resolver.get_filename("testdir1.zip") targets = [{ "destination": destination, - "items_from": "archive", "src": "path", "path": bed_test_data_path, + "items_from": "archive", "src": "path", "path": archive_test_data_path, }] payload = { "history_id": history_id, # TODO: Shouldn't be needed :( @@ -756,6 +756,23 @@ class FetchByPathTestCase(BaseUploadContentConfigurationTestCase): dataset = self.library_populator.get_library_contents_with_path(library["id"], "/dir1/file3") assert dataset["file_size"] == 11, dataset + def test_fetch_recursive_archive_history(self): + destination = {"type": "hdas"} + archive = self.test_data_resolver.get_filename("testdir1.zip") + targets = [{ + "destination": destination, + "items_from": "archive", "src": "path", "path": archive, + }] + payload = { + "history_id": self.history_id, # TODO: Shouldn't be needed :( + "targets": json.dumps(targets), + } + self.dataset_populator.fetch(payload) + contents_response = self.dataset_populator._get_contents_request(self.history_id) + assert contents_response.status_code == 200 + contents = contents_response.json() + assert len(contents) == 3 + def test_fetch_recursive_archive_to_library(self): bed_test_data_path = self.test_data_resolver.get_filename("testdir1.zip") targets = [{ From 1720354d86c29bcc5f50664139cbf935ebc51e5b Mon Sep 17 00:00:00 2001 From: John Chilton Date: Sun, 14 Jan 2018 13:29:25 -0500 Subject: [PATCH 15/24] More upload testing, some input is getting deleted that shouldn't. --- .../test_upload_configuration_options.py | 32 +++++++++++++++---- 1 file changed, 26 insertions(+), 6 deletions(-) diff --git a/test/integration/test_upload_configuration_options.py b/test/integration/test_upload_configuration_options.py index fe47ec04e71..82c1e43496e 100644 --- a/test/integration/test_upload_configuration_options.py +++ b/test/integration/test_upload_configuration_options.py @@ -125,6 +125,8 @@ class NonAdminsCannotPasteFilePathTestCase(BaseUploadContentConfigurationTestCas @skip_without_datatype("velvet") def test_disallowed_for_composite_file(self): + path = os.path.join(TEST_DATA_DIRECTORY, "1.txt") + assert os.path.exists(path) payload = self.dataset_populator.upload_payload( self.history_id, "sequences content", @@ -132,7 +134,7 @@ class NonAdminsCannotPasteFilePathTestCase(BaseUploadContentConfigurationTestCas extra_inputs={ "files_1|url_paste": "roadmaps content", "files_1|type": "upload_dataset", - "files_2|url_paste": "file://%s/1.txt" % TEST_DATA_DIRECTORY, + "files_2|url_paste": "file://%s" % path, "files_2|type": "upload_dataset", }, ) @@ -140,15 +142,21 @@ class NonAdminsCannotPasteFilePathTestCase(BaseUploadContentConfigurationTestCas # Ideally this would be 403 but the tool API endpoint isn't using # the newer API decorator that handles those details. assert create_response.status_code >= 400 + assert os.path.exists(path) def test_disallowed_for_libraries(self): + path = os.path.join(TEST_DATA_DIRECTORY, "1.txt") + assert os.path.exists(path) library = self.library_populator.new_private_library("pathpastedisallowedlibraries") - payload, files = self.library_populator.create_dataset_request(library, upload_option="upload_paths", paths="%s/1.txt" % TEST_DATA_DIRECTORY) + payload, files = self.library_populator.create_dataset_request(library, upload_option="upload_paths", paths=path) response = self.library_populator.raw_library_contents_create(library["id"], payload, files=files) assert response.status_code == 403, response.json() + assert os.path.exists(path) def test_disallowed_for_fetch(self): - elements = [{"src": "path", "path": "%s/1.txt" % TEST_DATA_DIRECTORY}] + path = os.path.join(TEST_DATA_DIRECTORY, "1.txt") + assert os.path.exists(path) + elements = [{"src": "path", "path": path}] target = { "destination": {"type": "hdca"}, "elements": elements, @@ -156,9 +164,12 @@ class NonAdminsCannotPasteFilePathTestCase(BaseUploadContentConfigurationTestCas } response = self.fetch_target(target) self._assert_status_code_is(response, 403) + assert os.path.exists(path) def test_disallowed_for_fetch_urls(self): - elements = [{"src": "url", "url": "file://%s/1.txt" % TEST_DATA_DIRECTORY}] + path = os.path.join(TEST_DATA_DIRECTORY, "1.txt") + assert os.path.exists(path) + elements = [{"src": "url", "url": "file://%s" % path}] target = { "destination": {"type": "hdca"}, "elements": elements, @@ -166,6 +177,7 @@ class NonAdminsCannotPasteFilePathTestCase(BaseUploadContentConfigurationTestCas } response = self.fetch_target(target) self._assert_status_code_is(response, 403) + assert os.path.exists(path) class AdminsCanPasteFilePathsTestCase(BaseUploadContentConfigurationTestCase): @@ -196,7 +208,8 @@ class AdminsCanPasteFilePathsTestCase(BaseUploadContentConfigurationTestCase): assert os.path.exists(path) def test_admin_fetch(self): - elements = [{"src": "path", "path": "%s/1.txt" % TEST_DATA_DIRECTORY}] + path = os.path.join(TEST_DATA_DIRECTORY, "1.txt") + elements = [{"src": "path", "path": path}] target = { "destination": {"type": "hdca"}, "elements": elements, @@ -204,9 +217,11 @@ class AdminsCanPasteFilePathsTestCase(BaseUploadContentConfigurationTestCase): } response = self.fetch_target(target) self._assert_status_code_is(response, 200) + assert os.path.exists(path) def test_admin_fetch_file_url(self): - elements = [{"src": "url", "url": "file://%s/1.txt" % TEST_DATA_DIRECTORY}] + path = os.path.join(TEST_DATA_DIRECTORY, "1.txt") + elements = [{"src": "url", "url": "file://%s" % path}] target = { "destination": {"type": "hdca"}, "elements": elements, @@ -214,6 +229,7 @@ class AdminsCanPasteFilePathsTestCase(BaseUploadContentConfigurationTestCase): } response = self.fetch_target(target) self._assert_status_code_is(response, 200) + assert os.path.exists(path) class DefaultBinaryContentFiltersTestCase(BaseUploadContentConfigurationTestCase): @@ -705,6 +721,7 @@ class FetchByPathTestCase(BaseUploadContentConfigurationTestCase): def test_fetch_path_to_folder(self): history_id, library, destination = self.library_populator.setup_fetch_to_folder("simple_fetch") bed_test_data_path = self.test_data_resolver.get_filename("4.bed") + assert os.path.exists(bed_test_data_path) items = [{"src": "path", "path": bed_test_data_path, "info": "my cool bed"}] targets = [{ "destination": destination, @@ -717,10 +734,12 @@ class FetchByPathTestCase(BaseUploadContentConfigurationTestCase): self.dataset_populator.fetch(payload) dataset = self.library_populator.get_library_contents_with_path(library["id"], "/4.bed") assert dataset["file_size"] == 61, dataset + assert os.path.exists(bed_test_data_path) def test_fetch_link_data_only(self): history_id, library, destination = self.library_populator.setup_fetch_to_folder("fetch_and_link") bed_test_data_path = self.test_data_resolver.get_filename("4.bed") + assert os.path.exists(bed_test_data_path) items = [{"src": "path", "path": bed_test_data_path, "info": "my cool bed", "link_data_only": True}] targets = [{ "destination": destination, @@ -734,6 +753,7 @@ class FetchByPathTestCase(BaseUploadContentConfigurationTestCase): dataset = self.library_populator.get_library_contents_with_path(library["id"], "/4.bed") assert dataset["file_size"] == 61, dataset assert dataset["file_name"] == bed_test_data_path, dataset + assert os.path.exists(bed_test_data_path) def test_fetch_recursive_archive(self): history_id, library, destination = self.library_populator.setup_fetch_to_folder("recursive_archive") From 5c5dbd2b6587dee35fef329b9536cc627f2bf6a4 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Mon, 22 Jan 2018 15:38:40 -0500 Subject: [PATCH 16/24] Handle compressed datatypes appropriately in data fetch API. --- lib/galaxy/tools/data_fetch.py | 6 +++++- .../tools/sample_datatypes_conf.xml | 6 ++++++ .../test_upload_configuration_options.py | 19 +++++++++++++++++++ 3 files changed, 30 insertions(+), 1 deletion(-) diff --git a/lib/galaxy/tools/data_fetch.py b/lib/galaxy/tools/data_fetch.py index 0cb87f52485..d32640e4db1 100644 --- a/lib/galaxy/tools/data_fetch.py +++ b/lib/galaxy/tools/data_fetch.py @@ -116,7 +116,11 @@ def _fetch_target(upload_config, target): if is_binary: data_type, ext = handle_sniffable_binary_check(data_type, ext, path, registry) if data_type is None: - if is_binary: + root_datatype = registry.get_datatype_by_extension(ext) + if getattr(root_datatype, 'compressed', False): + data_type = 'compressed archive' + ext = ext + elif is_binary: data_type, ext = handle_unsniffable_binary_check( data_type, ext, path, name, is_binary, requested_ext, check_content, registry ) diff --git a/test/functional/tools/sample_datatypes_conf.xml b/test/functional/tools/sample_datatypes_conf.xml index a37897791ee..fc08309fd8d 100644 --- a/test/functional/tools/sample_datatypes_conf.xml +++ b/test/functional/tools/sample_datatypes_conf.xml @@ -14,6 +14,12 @@ + + + + + + diff --git a/test/integration/test_upload_configuration_options.py b/test/integration/test_upload_configuration_options.py index 82c1e43496e..affd58eedb5 100644 --- a/test/integration/test_upload_configuration_options.py +++ b/test/integration/test_upload_configuration_options.py @@ -776,6 +776,25 @@ class FetchByPathTestCase(BaseUploadContentConfigurationTestCase): dataset = self.library_populator.get_library_contents_with_path(library["id"], "/dir1/file3") assert dataset["file_size"] == 11, dataset + def test_fetch_history_compressed_type(self): + destination = {"type": "hdas"} + archive = self.test_data_resolver.get_filename("1.fastqsanger.gz") + targets = [{ + "destination": destination, + "items": [{"src": "path", "path": archive, "ext": "fastqsanger.gz"}], + }] + payload = { + "history_id": self.history_id, # TODO: Shouldn't be needed :( + "targets": json.dumps(targets), + } + self.dataset_populator.fetch(payload) + contents_response = self.dataset_populator._get_contents_request(self.history_id) + assert contents_response.status_code == 200 + contents = contents_response.json() + assert len(contents) == 1 + print(contents) + contents[0]["extension"] == "fastqsanger.gz" + def test_fetch_recursive_archive_history(self): destination = {"type": "hdas"} archive = self.test_data_resolver.get_filename("testdir1.zip") From 54ee573952215214a15cf21f20c79b735a394342 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Mon, 26 Feb 2018 13:26:40 -0500 Subject: [PATCH 17/24] Simplify and reduce duplication of upload actions. --- lib/galaxy/tools/actions/upload.py | 62 ++++++++++++++---------------- 1 file changed, 28 insertions(+), 34 deletions(-) diff --git a/lib/galaxy/tools/actions/upload.py b/lib/galaxy/tools/actions/upload.py index 2b6013752e5..635f3098afb 100644 --- a/lib/galaxy/tools/actions/upload.py +++ b/lib/galaxy/tools/actions/upload.py @@ -9,9 +9,9 @@ from . import ToolAction log = logging.getLogger(__name__) -class UploadToolAction(ToolAction): +class BaseUploadToolAction(ToolAction): - def execute(self, tool, trans, incoming={}, set_output_hid=True, history=None, **kwargs): + def execute(self, tool, trans, incoming={}, history=None, **kwargs): dataset_upload_inputs = [] for input_name, input in tool.inputs.items(): if input.type == "upload_dataset": @@ -21,36 +21,40 @@ class UploadToolAction(ToolAction): persisting_uploads_timer = ExecutionTimer() incoming = upload_common.persist_uploads(incoming, trans) log.debug("Persisted uploads %s" % persisting_uploads_timer) + rval = self._setup_job(tool, trans, incoming, dataset_upload_inputs, history) + return rval + + def _setup_job(self, tool, trans, incoming, dataset_upload_inputs, history): + """Take persisted uploads and create a job for given tool.""" + + def _create_job(self, *args, **kwds): + """Wrapper around upload_common.create_job with a timer.""" + create_job_timer = ExecutionTimer() + rval = upload_common.create_job(*args, **kwds) + log.debug("Created upload job %s" % create_job_timer) + return rval + + +class UploadToolAction(BaseUploadToolAction): + + def _setup_job(self, tool, trans, incoming, dataset_upload_inputs, history): check_timer = ExecutionTimer() - # We can pass an empty string as the cntrller here since it is used to check whether we - # are in an admin view, and this tool is currently not used there. uploaded_datasets = upload_common.get_uploaded_datasets(trans, '', incoming, dataset_upload_inputs, history=history) if not uploaded_datasets: return None, 'No data was entered in the upload form, please go back and choose data to upload.' - log.debug("Checked uploads %s" % check_timer) - create_job_timer = ExecutionTimer() json_file_path = upload_common.create_paramfile(trans, uploaded_datasets) data_list = [ud.data for ud in uploaded_datasets] - rval = upload_common.create_job(trans, incoming, tool, json_file_path, data_list, history=history) - log.debug("Created upload job %s" % create_job_timer) - return rval + log.debug("Checked uploads %s" % check_timer) + return self._create_job( + trans, incoming, tool, json_file_path, data_list, history=history + ) -class FetchUploadToolAction(ToolAction): - - def execute(self, tool, trans, incoming={}, set_output_hid=True, history=None, **kwargs): - dataset_upload_inputs = [] - for input_name, input in tool.inputs.items(): - if input.type == "upload_dataset": - dataset_upload_inputs.append(input) - assert dataset_upload_inputs, Exception("No dataset upload groups were found.") - - persisting_uploads_timer = ExecutionTimer() - incoming = upload_common.persist_uploads(incoming, trans) - log.debug("Persisted uploads %s" % persisting_uploads_timer) +class FetchUploadToolAction(BaseUploadToolAction): + def _setup_job(self, tool, trans, incoming, dataset_upload_inputs, history): # Now replace references in requests with these. files = incoming.get("files", []) files_iter = iter(files) @@ -76,16 +80,6 @@ class FetchUploadToolAction(ToolAction): replace_file_srcs(request) incoming["request_json"] = json.dumps(request) - check_timer = ExecutionTimer() - # We can pass an empty string as the cntrller here since it is used to check whether we - # are in an admin view, and this tool is currently not used there. - # uploaded_datasets = upload_common.get_uploaded_datasets(trans, '', incoming, dataset_upload_inputs, history=history) - - # if not uploaded_datasets: - # return None, 'No data was entered in the upload form, please go back and choose data to upload.' - - log.debug("Checked uploads %s" % check_timer) - create_job_timer = ExecutionTimer() - rval = upload_common.create_job(trans, incoming, tool, None, [], history=history) - log.debug("Created upload job %s" % create_job_timer) - return rval + return self._create_job( + trans, incoming, tool, None, [], history=history + ) From 60f632bb242c0c73831dbe0af204218f59b38d77 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Mon, 26 Feb 2018 15:20:36 -0500 Subject: [PATCH 18/24] Precreate certain outputs for upload 2.0 API. Trying to improve the user experience of this rule based uploader by placing HDAs and HDCAs in the history at the outset that the history panel can poll and that we can turn red if the upload fails. From Marius' PR review: > I can see that a job launched in my logs, but it failed and there were no visual indications of this in the UI Not every HDA for instance can be created, for example if reading them from a zip file for instance that happens on the backend still. Likewise if HDCAs don't define a collection type up front they cannot be pre-created (if for instance that is inferred from a folder structure). Library things aren't precreated at all in this commit. There is room to pre-create more but I think this is an atomic commit as it is now and it will hopefully improve the user experience for the rule based uploader considerably. --- lib/galaxy/tools/actions/upload.py | 65 ++++++++++++++++++- lib/galaxy/tools/actions/upload_common.py | 42 +++++++----- lib/galaxy/tools/data_fetch.py | 3 + lib/galaxy/tools/parameters/output_collect.py | 42 ++++++++---- lib/galaxy/webapps/galaxy/api/_fetch_util.py | 6 ++ .../test_upload_configuration_options.py | 18 +++-- 6 files changed, 142 insertions(+), 34 deletions(-) diff --git a/lib/galaxy/tools/actions/upload.py b/lib/galaxy/tools/actions/upload.py index 635f3098afb..cdcc28b1ade 100644 --- a/lib/galaxy/tools/actions/upload.py +++ b/lib/galaxy/tools/actions/upload.py @@ -1,9 +1,12 @@ import json import logging +import os +from galaxy.dataset_collections.structure import UnitializedTree from galaxy.exceptions import RequestParameterMissingException from galaxy.tools.actions import upload_common from galaxy.util import ExecutionTimer +from galaxy.util.bunch import Bunch from . import ToolAction log = logging.getLogger(__name__) @@ -79,7 +82,67 @@ class FetchUploadToolAction(BaseUploadToolAction): replace_file_srcs(request) + outputs = [] + for target in request.get("targets", []): + destination = target.get("destination") + destination_type = destination.get("type") + # Start by just pre-creating HDAs. + if destination_type == "hdas": + if target.get("elements_from"): + # Dynamic collection required I think. + continue + _precreate_fetched_hdas(trans, history, target, outputs) + + if destination_type == "hdca": + _precreate_fetched_collection_instance(trans, history, target, outputs) + incoming["request_json"] = json.dumps(request) return self._create_job( - trans, incoming, tool, None, [], history=history + trans, incoming, tool, None, outputs, history=history ) + + +def _precreate_fetched_hdas(trans, history, target, outputs): + for item in target.get("elements", []): + name = item.get("name", None) + if name is None: + src = item.get("src", None) + if src == "url": + url = item.get("url") + if name is None: + name = url.split("/")[-1] + elif src == "path": + path = item["path"] + if name is None: + name = os.path.basename(path) + + file_type = item.get("ext", "auto") + dbkey = item.get("dbkey", "?") + uploaded_dataset = Bunch( + type='file', name=name, file_type=file_type, dbkey=dbkey + ) + data = upload_common.new_upload(trans, '', uploaded_dataset, library_bunch=None, history=history) + outputs.append(data) + item["object_id"] = data.id + + +def _precreate_fetched_collection_instance(trans, history, target, outputs): + collection_type = target.get("collection_type") + if not collection_type: + # Can't precreate collections of unknown type at this time. + return + + name = target.get("name") + if not name: + return + + collections_service = trans.app.dataset_collections_service + collection_type_description = collections_service.collection_type_descriptions.for_collection_type(collection_type) + structure = UnitializedTree(collection_type_description) + hdca = collections_service.precreate_dataset_collection_instance( + trans, history, name, structure=structure + ) + outputs.append(hdca) + # Following flushed needed for an ID. + trans.sa_session.flush() + target["destination"]["object_id"] = hdca.id diff --git a/lib/galaxy/tools/actions/upload_common.py b/lib/galaxy/tools/actions/upload_common.py index ad2ebbe3200..89956b82df4 100644 --- a/lib/galaxy/tools/actions/upload_common.py +++ b/lib/galaxy/tools/actions/upload_common.py @@ -384,7 +384,7 @@ def create_paramfile(trans, uploaded_datasets): return json_file_path -def create_job(trans, params, tool, json_file_path, data_list, folder=None, history=None, job_params=None): +def create_job(trans, params, tool, json_file_path, outputs, folder=None, history=None, job_params=None): """ Create the upload job. """ @@ -412,21 +412,28 @@ def create_job(trans, params, tool, json_file_path, data_list, folder=None, hist job.add_parameter(name, value) job.add_parameter('paramfile', dumps(json_file_path)) object_store_id = None - for i, dataset in enumerate(data_list): - if folder: - job.add_output_library_dataset('output%i' % i, dataset) + for i, output_object in enumerate(outputs): + output_name = "output%i" % i + if hasattr(output_object, "collection"): + job.add_output_dataset_collection(output_name, output_object) + output_object.job = job else: - job.add_output_dataset('output%i' % i, dataset) - # Create an empty file immediately - if not dataset.dataset.external_filename: - dataset.dataset.object_store_id = object_store_id - try: - trans.app.object_store.create(dataset.dataset) - except ObjectInvalid: - raise Exception('Unable to create output dataset: object store is full') - object_store_id = dataset.dataset.object_store_id - trans.sa_session.add(dataset) - # open( dataset.file_name, "w" ).close() + dataset = output_object + if folder: + job.add_output_library_dataset(output_name, dataset) + else: + job.add_output_dataset(output_name, dataset) + # Create an empty file immediately + if not dataset.dataset.external_filename: + dataset.dataset.object_store_id = object_store_id + try: + trans.app.object_store.create(dataset.dataset) + except ObjectInvalid: + raise Exception('Unable to create output dataset: object store is full') + object_store_id = dataset.dataset.object_store_id + + trans.sa_session.add(output_object) + job.object_store_id = object_store_id job.set_state(job.states.NEW) job.set_handler(tool.get_job_handler(None)) @@ -440,8 +447,9 @@ def create_job(trans, params, tool, json_file_path, data_list, folder=None, hist trans.app.job_manager.job_queue.put(job.id, job.tool_id) trans.log_event("Added job to the job queue, id: %s" % str(job.id), tool_id=job.tool_id) output = odict() - for i, v in enumerate(data_list): - output['output%i' % i] = v + for i, v in enumerate(outputs): + if not hasattr(output_object, "collection_type"): + output['output%i' % i] = v return job, output diff --git a/lib/galaxy/tools/data_fetch.py b/lib/galaxy/tools/data_fetch.py index d32640e4db1..c2e191f1015 100644 --- a/lib/galaxy/tools/data_fetch.py +++ b/lib/galaxy/tools/data_fetch.py @@ -98,6 +98,7 @@ def _fetch_target(upload_config, target): dbkey = item.get("dbkey", "?") requested_ext = item.get("ext", "auto") info = item.get("info", None) + object_id = item.get("object_id", None) link_data_only = upload_config.link_data_only if "link_data_only" in item: # Allow overriding this on a per file basis. @@ -170,6 +171,8 @@ def _fetch_target(upload_config, target): rval = {"name": name, "filename": path, "dbkey": dbkey, "ext": ext, "link_data_only": link_data_only} if info is not None: rval["info"] = info + if object_id is not None: + rval["object_id"] = object_id return rval elements = elements_tree_map(_resolve_src, items) diff --git a/lib/galaxy/tools/parameters/output_collect.py b/lib/galaxy/tools/parameters/output_collect.py index 938771d515d..1bb14f06097 100644 --- a/lib/galaxy/tools/parameters/output_collect.py +++ b/lib/galaxy/tools/parameters/output_collect.py @@ -218,13 +218,18 @@ def collect_dynamic_outputs( elif destination_type == "hdca": history = job.history assert "collection_type" in unnamed_output_dict - name = unnamed_output_dict.get("name", "unnamed collection") - collection_type = unnamed_output_dict["collection_type"] - collection_type_description = collections_service.collection_type_descriptions.for_collection_type(collection_type) - structure = UnitializedTree(collection_type_description) - hdca = collections_service.precreate_dataset_collection_instance( - trans, history, name, structure=structure - ) + object_id = destination.get("object_id") + if object_id: + sa_session = tool.app.model.context + hdca = sa_session.query(app.model.HistoryDatasetCollectionAssociation).get(int(object_id)) + else: + name = unnamed_output_dict.get("name", "unnamed collection") + collection_type = unnamed_output_dict["collection_type"] + collection_type_description = collections_service.collection_type_descriptions.for_collection_type(collection_type) + structure = UnitializedTree(collection_type_description) + hdca = collections_service.precreate_dataset_collection_instance( + trans, history, name, structure=structure + ) filenames = odict.odict() def add_to_discovered_files(elements, parent_identifiers=[]): @@ -268,6 +273,12 @@ def collect_dynamic_outputs( # Create new primary dataset name = fields_match.name or designation + hda_id = discovered_file.match.object_id + primary_dataset = None + if hda_id: + sa_session = tool.app.model.context + primary_dataset = sa_session.query(app.model.HistoryDatasetAssociation).get(int(hda_id)) + dataset = job_context.create_dataset( ext=ext, designation=designation, @@ -276,7 +287,8 @@ def collect_dynamic_outputs( name=name, filename=discovered_file.path, info=info, - link_data=link_data + link_data=link_data, + primary_data=primary_dataset, ) dataset.raw_set_dataset_state('ok') datasets.append(dataset) @@ -451,14 +463,16 @@ class JobContext(object): info=None, library_folder=None, link_data=False, + primary_data=None, ): app = self.app sa_session = self.sa_session - if not library_folder: - primary_data = _new_hda(app, sa_session, ext, designation, visible, dbkey, self.permissions) - else: - primary_data = _new_ldda(self.work_context, name, ext, visible, dbkey, library_folder) + if primary_data is None: + if not library_folder: + primary_data = _new_hda(app, sa_session, ext, designation, visible, dbkey, self.permissions) + else: + primary_data = _new_ldda(self.work_context, name, ext, visible, dbkey, library_folder) # Copy metadata from one of the inputs if requested. metadata_source = None @@ -843,6 +857,10 @@ class JsonCollectedDatasetMatch(object): def link_data(self): return bool(self.as_dict.get("link_data_only", False)) + @property + def object_id(self): + return self.as_dict.get("object_id", None) + class RegexCollectedDatasetMatch(JsonCollectedDatasetMatch): diff --git a/lib/galaxy/webapps/galaxy/api/_fetch_util.py b/lib/galaxy/webapps/galaxy/api/_fetch_util.py index d2b2eb11a68..7c5e2ea8e54 100644 --- a/lib/galaxy/webapps/galaxy/api/_fetch_util.py +++ b/lib/galaxy/webapps/galaxy/api/_fetch_util.py @@ -37,6 +37,9 @@ def validate_and_normalize_targets(trans, payload): 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'") + if "object_id" in destination: + raise RequestParameterInvalidException("object_id not allowed to appear in the request.") + if destination_type not in VALID_DESTINATION_TYPES: template = "Invalid target destination type [%s] encountered, must be one of %s" msg = template % (destination_type, VALID_DESTINATION_TYPES) @@ -63,6 +66,9 @@ def validate_and_normalize_targets(trans, payload): payload["check_content"] = trans.app.config.check_upload_content def check_src(item): + if "object_id" in item: + raise RequestParameterInvalidException("object_id not allowed to appear in the request.") + # Normalize file:// URLs into paths. if item["src"] == "url" and item["url"].startswith("file://"): item["src"] = "path" diff --git a/test/integration/test_upload_configuration_options.py b/test/integration/test_upload_configuration_options.py index affd58eedb5..d59aa375985 100644 --- a/test/integration/test_upload_configuration_options.py +++ b/test/integration/test_upload_configuration_options.py @@ -419,9 +419,14 @@ class SimpleFtpUploadConfigurationTestCase(BaseFtpUploadConfigurationTestCase): "destination": {"type": "hdca"}, "elements": elements, "collection_type": "list", + "name": "cool collection", } response = self.fetch_target(target) self._assert_status_code_is(response, 200) + response_object = response.json() + assert "output_collections" in response_object + output_collections = response_object["output_collections"] + assert len(output_collections) == 1, response_object dataset = self.dataset_populator.get_history_dataset_details(self.history_id, hid=2) self._check_content(dataset, content) @@ -787,13 +792,18 @@ class FetchByPathTestCase(BaseUploadContentConfigurationTestCase): "history_id": self.history_id, # TODO: Shouldn't be needed :( "targets": json.dumps(targets), } - self.dataset_populator.fetch(payload) + fetch_response = self.dataset_populator.fetch(payload) + self._assert_status_code_is(fetch_response, 200) + outputs = fetch_response.json()["outputs"] + assert len(outputs) == 1 + output = outputs[0] + assert output["name"] == "1.fastqsanger.gz" contents_response = self.dataset_populator._get_contents_request(self.history_id) assert contents_response.status_code == 200 contents = contents_response.json() - assert len(contents) == 1 - print(contents) - contents[0]["extension"] == "fastqsanger.gz" + assert len(contents) == 1, contents + assert contents[0]["extension"] == "fastqsanger.gz", contents[0] + assert contents[0]["name"] == "1.fastqsanger.gz", contents[0] def test_fetch_recursive_archive_history(self): destination = {"type": "hdas"} From d783fc33091e2952c2775f13a3453517bcffb9bb Mon Sep 17 00:00:00 2001 From: John Chilton Date: Mon, 5 Mar 2018 08:42:45 -0500 Subject: [PATCH 19/24] Avoid symlinks in upload FTP tests. --- test/integration/test_upload_configuration_options.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/test/integration/test_upload_configuration_options.py b/test/integration/test_upload_configuration_options.py index d59aa375985..6f58cae9190 100644 --- a/test/integration/test_upload_configuration_options.py +++ b/test/integration/test_upload_configuration_options.py @@ -68,7 +68,8 @@ class BaseUploadContentConfigurationTestCase(integration_util.IntegrationTestCas @classmethod def temp_config_dir(cls, name): - return os.path.join(cls._test_driver.galaxy_test_tmp_dir, name) + # realpath here to get around problems with symlinks being blocked. + return os.path.realpath(os.path.join(cls._test_driver.galaxy_test_tmp_dir, name)) def _write_file(self, dir_path, content, filename="test"): """Helper for writing ftp/server dir files.""" From 33141835c9700c2fd93fc10ce6e68c6a41026377 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Mon, 5 Mar 2018 08:43:31 -0500 Subject: [PATCH 20/24] Cleanup hierarchical upload commit based on PR comments from @bgruening. - Remove seemingly unneeded hack in upload_common. - Remove stray debug statement. - Add more comments in the output collection code related to different destination types. - Restructure if/else in data_fetch to avoid assertion with constant. --- lib/galaxy/tools/actions/upload_common.py | 6 +--- lib/galaxy/tools/data_fetch.py | 28 ++++++++++--------- lib/galaxy/tools/parameters/output_collect.py | 10 ++++++- 3 files changed, 25 insertions(+), 19 deletions(-) diff --git a/lib/galaxy/tools/actions/upload_common.py b/lib/galaxy/tools/actions/upload_common.py index 89956b82df4..4a686a107d4 100644 --- a/lib/galaxy/tools/actions/upload_common.py +++ b/lib/galaxy/tools/actions/upload_common.py @@ -284,11 +284,7 @@ def new_upload(trans, cntrller, uploaded_dataset, library_bunch=None, history=No def get_uploaded_datasets(trans, cntrller, params, dataset_upload_inputs, library_bunch=None, history=None): uploaded_datasets = [] for dataset_upload_input in dataset_upload_inputs: - try: - uploaded_datasets.extend(dataset_upload_input.get_uploaded_datasets(trans, params)) - except AttributeError: - # TODO: refine... - pass + uploaded_datasets.extend(dataset_upload_input.get_uploaded_datasets(trans, params)) for uploaded_dataset in uploaded_datasets: data = new_upload(trans, cntrller, uploaded_dataset, library_bunch=library_bunch, history=history) uploaded_dataset.data = data diff --git a/lib/galaxy/tools/data_fetch.py b/lib/galaxy/tools/data_fetch.py index c2e191f1015..ddde6207925 100644 --- a/lib/galaxy/tools/data_fetch.py +++ b/lib/galaxy/tools/data_fetch.py @@ -62,19 +62,21 @@ def _fetch_target(upload_config, target): def expand_elements_from(target_or_item): elements_from = target_or_item.get("elements_from", None) items = None - assert not elements_from or elements_from in ["archive", "bagit", "bagit_archive", "directory"], elements_from - if elements_from == "archive": - decompressed_directory = _decompress_target(target_or_item) - items = _directory_to_items(decompressed_directory) - elif elements_from == "bagit": - _, elements_from_path = _has_src_to_path(target_or_item) - items = _bagit_to_items(elements_from_path) - elif elements_from == "bagit_archive": - decompressed_directory = _decompress_target(target_or_item) - items = _bagit_to_items(decompressed_directory) - elif elements_from == "directory": - _, elements_from_path = _has_src_to_path(target_or_item) - items = _directory_to_items(elements_from_path) + if elements_from: + if elements_from == "archive": + decompressed_directory = _decompress_target(target_or_item) + items = _directory_to_items(decompressed_directory) + elif elements_from == "bagit": + _, elements_from_path = _has_src_to_path(target_or_item) + items = _bagit_to_items(elements_from_path) + elif elements_from == "bagit_archive": + decompressed_directory = _decompress_target(target_or_item) + items = _bagit_to_items(decompressed_directory) + elif elements_from == "directory": + _, elements_from_path = _has_src_to_path(target_or_item) + items = _directory_to_items(elements_from_path) + else: + raise Exception("Unknown elements from type encountered [%s]" % elements_from) if items: del target_or_item["elements_from"] diff --git a/lib/galaxy/tools/parameters/output_collect.py b/lib/galaxy/tools/parameters/output_collect.py index 1bb14f06097..1ebdee5b897 100644 --- a/lib/galaxy/tools/parameters/output_collect.py +++ b/lib/galaxy/tools/parameters/output_collect.py @@ -165,7 +165,8 @@ def collect_dynamic_outputs( inp_data, input_dbkey, ) - log.info(tool_provided_metadata) + # unmapped outputs do not correspond to explicit outputs of the tool, they were inferred entirely + # from the tool provided metadata (e.g. galaxy.json). for unnamed_output_dict in tool_provided_metadata.get_unnamed_outputs(): assert "destination" in unnamed_output_dict assert "elements" in unnamed_output_dict @@ -174,9 +175,14 @@ def collect_dynamic_outputs( assert "type" in destination destination_type = destination["type"] + assert destination_type in ["library_folder", "hdca", "hdas"] trans = job_context.work_context + # three destination types we need to handle here - "library_folder" (place discovered files in a library folder), + # "hdca" (place discovered files in a history dataset collection), and "hdas" (place discovered files in a history + # as stand-alone datasets). if destination_type == "library_folder": + # populate a library folder (needs to be already have been created) library_folder_manager = app.library_folder_manager library_folder = library_folder_manager.get(trans, app.security.decode_id(destination.get("library_folder_id"))) @@ -216,6 +222,7 @@ def collect_dynamic_outputs( add_elements_to_folder(elements, library_folder) elif destination_type == "hdca": + # create or populate a dataset collection in the history history = job.history assert "collection_type" in unnamed_output_dict object_id = destination.get("object_id") @@ -253,6 +260,7 @@ def collect_dynamic_outputs( ) collection_builder.populate() elif destination_type == "hdas": + # discover files as individual datasets for the target history history = job.history datasets = [] From ca28000f8b74738e50bb33aab67eb933a2df6f7a Mon Sep 17 00:00:00 2001 From: John Chilton Date: Tue, 6 Mar 2018 11:10:29 -0500 Subject: [PATCH 21/24] Fixes for pre-creating HDAs using data fetch API. --- lib/galaxy/tools/parameters/output_collect.py | 9 +++++++-- test/integration/test_upload_configuration_options.py | 1 + 2 files changed, 8 insertions(+), 2 deletions(-) diff --git a/lib/galaxy/tools/parameters/output_collect.py b/lib/galaxy/tools/parameters/output_collect.py index 1ebdee5b897..93bf92eee08 100644 --- a/lib/galaxy/tools/parameters/output_collect.py +++ b/lib/galaxy/tools/parameters/output_collect.py @@ -285,7 +285,7 @@ def collect_dynamic_outputs( primary_dataset = None if hda_id: sa_session = tool.app.model.context - primary_dataset = sa_session.query(app.model.HistoryDatasetAssociation).get(int(hda_id)) + primary_dataset = sa_session.query(app.model.HistoryDatasetAssociation).get(hda_id) dataset = job_context.create_dataset( ext=ext, @@ -299,7 +299,8 @@ def collect_dynamic_outputs( primary_data=primary_dataset, ) dataset.raw_set_dataset_state('ok') - datasets.append(dataset) + if not hda_id: + datasets.append(dataset) collect_elements_for_history(elements) job.history.add_datasets(job_context.sa_session, datasets) @@ -481,6 +482,10 @@ class JobContext(object): primary_data = _new_hda(app, sa_session, ext, designation, visible, dbkey, self.permissions) else: primary_data = _new_ldda(self.work_context, name, ext, visible, dbkey, library_folder) + else: + primary_data.extension = ext + primary_data.visible = visible + primary_data.dbkey = dbkey # Copy metadata from one of the inputs if requested. metadata_source = None diff --git a/test/integration/test_upload_configuration_options.py b/test/integration/test_upload_configuration_options.py index 6f58cae9190..fa9548c230a 100644 --- a/test/integration/test_upload_configuration_options.py +++ b/test/integration/test_upload_configuration_options.py @@ -805,6 +805,7 @@ class FetchByPathTestCase(BaseUploadContentConfigurationTestCase): assert len(contents) == 1, contents assert contents[0]["extension"] == "fastqsanger.gz", contents[0] assert contents[0]["name"] == "1.fastqsanger.gz", contents[0] + assert contents[0]["hid"] == 1, contents[0] def test_fetch_recursive_archive_history(self): destination = {"type": "hdas"} From 495d1258fd50780ea1dffd05eeef558161dc041d Mon Sep 17 00:00:00 2001 From: John Chilton Date: Tue, 6 Mar 2018 13:23:49 -0500 Subject: [PATCH 22/24] Consistent sniffing regardless of in_place. Previously sniffing would happen on the original file (before carriage returns and tabular spaces were converted) if in_place was false and on the converted file if it was true. --- tools/data_source/upload.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tools/data_source/upload.py b/tools/data_source/upload.py index bdb9ad4139a..4e0cdcdb042 100644 --- a/tools/data_source/upload.py +++ b/tools/data_source/upload.py @@ -277,7 +277,7 @@ def add_file(dataset, registry, json_file, output_path): else: line_count, converted_path = sniff.convert_newlines(dataset.path, in_place=in_place, tmp_dir=tmpdir, tmp_prefix=tmp_prefix) if dataset.file_type == 'auto': - ext = sniff.guess_ext(dataset.path, registry.sniff_order) + ext = sniff.guess_ext(converted_path or dataset.path, registry.sniff_order) else: ext = dataset.file_type data_type = ext From d651348da663966c1c04585f557820af53c5f651 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Tue, 6 Mar 2018 11:52:34 -0500 Subject: [PATCH 23/24] More upload tests and fixes. --- lib/galaxy/datatypes/sniff.py | 15 ++-- lib/galaxy/tools/data_fetch.py | 3 + test-data/1.csv | 1 + test/api/test_tools_upload.py | 82 +++++++++++++++++-- .../tools/sample_datatypes_conf.xml | 1 + 5 files changed, 91 insertions(+), 11 deletions(-) create mode 100644 test-data/1.csv diff --git a/lib/galaxy/datatypes/sniff.py b/lib/galaxy/datatypes/sniff.py index 9cbe81270b9..c69373b0669 100644 --- a/lib/galaxy/datatypes/sniff.py +++ b/lib/galaxy/datatypes/sniff.py @@ -138,7 +138,7 @@ def convert_newlines(fname, in_place=True, tmp_dir=None, tmp_prefix="gxupload"): return (i, temp_name) -def sep2tabs(fname, in_place=True, patt="\\s+"): +def sep2tabs(fname, in_place=True, patt="\\s+", tmp_dir=None, tmp_prefix="gxupload"): """ Transforms in place a 'sep' separated file to a tab separated one @@ -150,13 +150,18 @@ def sep2tabs(fname, in_place=True, patt="\\s+"): '1\\t2\\n3\\t4\\n' """ regexp = re.compile(patt) - fd, temp_name = tempfile.mkstemp() + fd, temp_name = tempfile.mkstemp(prefix=tmp_prefix, dir=tmp_dir) with os.fdopen(fd, "wt") as fp: i = None for i, line in enumerate(open(fname)): - line = line.rstrip('\r\n') - elems = regexp.split(line) - fp.write("%s\n" % '\t'.join(elems)) + if line.endswith("\r"): + line = line.rstrip('\r') + elems = regexp.split(line) + fp.write("%s\r" % '\t'.join(elems)) + else: + line = line.rstrip('\n') + elems = regexp.split(line) + fp.write("%s\n" % '\t'.join(elems)) if i is None: i = 0 else: diff --git a/lib/galaxy/tools/data_fetch.py b/lib/galaxy/tools/data_fetch.py index ddde6207925..c6a2b7f6c61 100644 --- a/lib/galaxy/tools/data_fetch.py +++ b/lib/galaxy/tools/data_fetch.py @@ -137,6 +137,9 @@ def _fetch_target(upload_config, target): line_count, converted_path = sniff.convert_newlines_sep2tabs(path, in_place=in_place, tmp_dir=".") else: line_count, converted_path = sniff.convert_newlines(path, in_place=in_place, tmp_dir=".") + else: + if space_to_tab: + line_count, converted_path = sniff.sep2tabs(path, in_place=in_place, tmp_dir=".") if requested_ext == 'auto': ext = sniff.guess_ext(path, registry.sniff_order) diff --git a/test-data/1.csv b/test-data/1.csv new file mode 100644 index 00000000000..80d519fbe2d --- /dev/null +++ b/test-data/1.csv @@ -0,0 +1 @@ +Transaction_date,Product,Price,Payment_Type,Name,City,State,Country,Account_Created,Last_Login,Latitude,Longitude 1/2/09 6:17,Product1,1200,Mastercard,carolina,Basildon,England,United Kingdom,1/2/09 6:00,1/2/09 6:08,51.5,-1.1166667 1/2/09 4:53,Product1,1200,Visa,Betina,Parkville ,MO,United States,1/2/09 4:42,1/2/09 7:49,39.195,-94.68194 1/2/09 13:08,Product1,1200,Mastercard,Federica e Andrea,Astoria ,OR,United States,1/1/09 16:21,1/3/09 12:32,46.18806,-123.83 1/3/09 14:44,Product1,1200,Visa,Gouya,Echuca,Victoria,Australia,9/25/05 21:13,1/3/09 14:22,-36.1333333,144.75 1/4/09 12:56,Product2,3600,Visa,Gerd W ,Cahaba Heights ,AL,United States,11/15/08 15:47,1/4/09 12:45,33.52056,-86.8025 1/4/09 13:19,Product1,1200,Visa,LAURENCE,Mickleton ,NJ,United States,9/24/08 15:19,1/4/09 13:04,39.79,-75.23806 1/4/09 20:11,Product1,1200,Mastercard,Fleur,Peoria ,IL,United States,1/3/09 9:38,1/4/09 19:45,40.69361,-89.58889 1/2/09 20:09,Product1,1200,Mastercard,adam,Martin ,TN,United States,1/2/09 17:43,1/4/09 20:01,36.34333,-88.85028 1/4/09 13:17,Product1,1200,Mastercard,Renee Elisabeth,Tel Aviv,Tel Aviv,Israel,1/4/09 13:03,1/4/09 22:10,32.0666667,34.7666667 1/4/09 14:11,Product1,1200,Visa,Aidan,Chatou,Ile-de-France,France,6/3/08 4:22,1/5/09 1:17,48.8833333,2.15 1/5/09 2:42,Product1,1200,Diners,Stacy,New York ,NY,United States,1/5/09 2:23,1/5/09 4:59,40.71417,-74.00639 1/5/09 5:39,Product1,1200,Amex,Heidi,Eindhoven,Noord-Brabant,Netherlands,1/5/09 4:55,1/5/09 8:15,51.45,5.4666667 1/2/09 9:16,Product1,1200,Mastercard,Sean ,Shavano Park ,TX,United States,1/2/09 8:32,1/5/09 9:05,29.42389,-98.49333 1/5/09 10:08,Product1,1200,Visa,Georgia,Eagle ,ID,United States,11/11/08 15:53,1/5/09 10:05,43.69556,-116.35306 1/2/09 14:18,Product1,1200,Visa,Richard,Riverside ,NJ,United States,12/9/08 12:07,1/5/09 11:01,40.03222,-74.95778 1/25/09 17:58,Product2,3600,Visa,carol,Ann Arbor ,MI,United States,7/5/08 9:20,2/7/09 18:51,42.27083,-83.72639 1/9/09 14:37,Product1,1200,Visa,Nona,South Jordan ,UT,United States,1/8/09 15:14,2/7/09 19:11,40.56222,-111.92889 1/25/09 2:46,Product2,3600,Visa,Family,Dubai,Dubayy,United Arab Emirates,1/8/09 1:19,2/8/09 2:06,25.2522222,55.28 1/17/09 20:46,Product2,3600,Visa,Michelle,Dubai,Dubayy,United Arab Emirates,4/13/08 2:36,2/8/09 2:12,25.2522222,55.28 1/24/09 7:18,Product2,3600,Visa,Kathryn,Kirriemuir,Scotland,United Kingdom,1/23/09 10:31,2/8/09 2:52,56.6666667,-3 1/11/09 7:09,Product1,1200,Visa,Oswald,Tramore,Waterford,Ireland,10/13/08 16:43,2/8/09 3:02,52.1588889,-7.1463889 1/8/09 4:15,Product1,1200,Visa,Elyssa,Gdansk,Pomorskie,Poland,1/7/09 15:00,2/8/09 3:50,54.35,18.6666667 1/22/09 10:47,Product1,1200,Visa,michelle,Arklow,Wicklow,Ireland,11/18/08 1:32,2/8/09 5:07,52.7930556,-6.1413889 1/26/09 20:47,Product1,1200,Mastercard,Alicia,Lincoln ,NE,United States,6/24/08 8:05,2/8/09 7:29,40.8,-96.66667 1/12/09 12:22,Product1,1200,Mastercard,JP,Tierp,Uppsala,Sweden,1/6/09 11:34,2/8/09 11:15,60.3333333,17.5 1/26/09 1:44,Product2,3600,Visa,Geraldine,Brussels,Brussels (Bruxelles),Belgium,1/31/08 13:28,2/8/09 14:39,50.8333333,4.3333333 1/18/09 12:57,Product1,1200,Mastercard,sandra,Burr Oak ,IA,United States,1/24/08 16:11,2/8/09 15:30,43.45889,-91.86528 1/24/09 21:26,Product1,1200,Visa,Olivia,Wheaton ,IL,United States,5/8/08 16:02,2/8/09 16:00,41.86611,-88.10694 1/26/09 12:26,Product2,3600,Mastercard,Tom,Killeen ,TX,United States,1/26/09 5:23,2/8/09 17:33,31.11694,-97.7275 1/5/09 7:37,Product1,1200,Visa,Annette ,Manhattan ,NY,United States,9/26/08 4:29,2/8/09 18:42,40.71417,-74.00639 1/14/09 12:33,Product1,1200,Visa,SUSAN,Oxford,England,United Kingdom,9/11/08 23:23,2/8/09 23:00,51.75,-1.25 1/14/09 0:15,Product2,3600,Visa,Michael,Paris,Ile-de-France,France,11/28/08 0:07,2/9/09 1:30,48.8666667,2.3333333 1/1/09 12:42,Product1,1200,Visa,ashton,Exeter,England,United Kingdom,12/15/08 1:16,2/9/09 2:52,50.7,-3.5333333 1/6/09 6:07,Product1,1200,Visa,Scott,Rungsted,Frederiksborg,Denmark,12/27/08 14:29,2/9/09 4:20,55.8841667,12.5419444 1/15/09 5:11,Product2,3600,Visa,Pam,London,England,United Kingdom,7/11/06 12:43,2/9/09 4:42,51.52721,0.14559 1/17/09 4:03,Product1,1200,Visa,Lisa ,Borja,Bohol,Philippines,1/17/09 2:45,2/9/09 6:09,9.9136111,124.0927778 1/19/09 10:13,Product2,3600,Mastercard,Pavel,London,England,United Kingdom,2/28/06 5:35,2/9/09 6:57,51.51334,-0.08895 1/18/09 9:42,Product1,1200,Visa,Richard,Jamestown ,RI,United States,1/18/09 9:22,2/9/09 8:30,41.49694,-71.36778 1/9/09 11:14,Product1,1200,Visa,Jasinta Jeanne,Owings Mills ,MD,United States,1/9/09 10:43,2/9/09 9:17,39.41944,-76.78056 1/10/09 13:42,Product1,1200,Visa,Rachel,Hamilton,Ontario,Canada,1/10/09 12:22,2/9/09 9:54,43.25,-79.8333333 1/7/09 7:28,Product1,1200,Amex,Cherish ,Anchorage ,AK,United States,7/28/08 7:31,2/9/09 10:50,61.21806,-149.90028 1/18/09 6:46,Product1,1200,Visa,Shona ,Mornington,Meath,Ireland,1/15/09 9:13,2/9/09 11:55,53.7233333,-6.2825 1/30/09 12:18,Product1,1200,Mastercard,Abikay,Fullerton ,CA,United States,1/26/09 13:34,2/9/09 12:53,33.87028,-117.92444 1/6/09 5:42,Product1,1200,Amex,Abikay,Atlanta ,GA,United States,10/27/08 14:16,2/9/09 13:50,33.74889,-84.38806 1/2/09 10:58,Product2,3600,Visa,Kendra,Toronto,Ontario,Canada,1/2/09 10:38,2/9/09 13:56,43.6666667,-79.4166667 1/8/09 3:29,Product1,1200,Visa,amanda,Liverpool,England,United Kingdom,12/22/08 1:41,2/9/09 14:06,53.4166667,-3 1/12/09 13:23,Product2,3600,Amex,Leila,Ponte San Nicolo,Veneto,Italy,9/13/05 8:42,2/9/09 14:09,45.3666667,11.6166667 1/19/09 9:34,Product1,1200,Amex,amanda,Las Vegas ,NV,United States,5/10/08 8:56,2/9/09 16:44,36.175,-115.13639 1/9/09 7:49,Product1,1200,Visa,Stacy,Rochester Hills ,MI,United States,7/28/08 7:18,2/9/09 17:41,42.68056,-83.13389 1/15/09 5:27,Product2,3600,Visa,Derrick,North Bay,Ontario,Canada,1/6/09 17:42,2/9/09 18:22,46.3,-79.45 1/8/09 23:40,Product1,1200,Visa,Jacob,Lindfield,New South Wales,Australia,1/8/09 17:52,2/9/09 18:31,-33.7833333,151.1666667 1/27/09 11:02,Product1,1200,Mastercard,DOREEN,Madrid,Madrid,Spain,1/24/09 8:21,2/9/09 18:42,40.4,-3.6833333 1/14/09 13:23,Product1,1200,Diners,eugenia,Wisconsin Rapids ,WI,United States,11/15/08 13:57,2/9/09 18:44,44.38361,-89.81722 1/7/09 20:01,Product1,1200,Visa,Karen,Austin ,TX,United States,1/6/09 19:16,2/9/09 19:56,30.26694,-97.74278 1/20/09 12:32,Product1,1200,Visa,Bea,Chicago ,IL,United States,1/16/09 19:08,2/9/09 20:42,41.85,-87.65 1/6/09 14:35,Product1,1200,Diners,Hilde Karin,Las Vegas ,NV,United States,12/17/08 11:59,2/9/09 22:59,36.175,-115.13639 1/4/09 6:51,Product1,1200,Visa,Rima,Mullingar,Westmeath,Ireland,1/3/09 12:34,2/10/09 0:59,53.5333333,-7.35 1/24/09 18:30,Product1,1200,Visa,Ruangrote,Melbourne,Victoria,Australia,7/17/08 5:19,2/10/09 2:12,-37.8166667,144.9666667 1/25/09 5:57,Product1,1200,Amex,pamela,Ayacucho,Buenos Aires,Argentina,1/24/09 9:29,2/10/09 6:38,-37.15,-58.4833333 1/5/09 10:02,Product2,3600,Visa,Emillie,Eagan ,MN,United States,1/5/09 9:03,2/10/09 7:29,44.80417,-93.16667 1/13/09 9:14,Product1,1200,Visa,sangeeta,Vossevangen,Hordaland,Norway,1/9/09 9:31,2/10/09 9:04,60.6333333,6.4333333 1/22/09 7:35,Product1,1200,Visa,Anja,Ferney-Voltaire,Rhone-Alpes,France,1/22/09 6:51,2/10/09 9:18,46.25,6.1166667 1/2/09 11:06,Product1,1200,Mastercard,Andrew,Sevilla,Andalucia,Spain,3/12/06 15:02,2/10/09 10:04,37.3772222,-5.9869444 1/11/09 9:50,Product1,1200,Visa,Bato,Munchengosserstadt,Thuringia,Germany,1/7/09 11:45,2/10/09 10:28,51.05,11.65 1/21/09 20:44,Product1,1200,Mastercard,Ailsa ,Lindenhurst ,NY,United States,1/21/09 7:47,2/10/09 10:51,40.68667,-73.37389 1/5/09 9:09,Product1,1200,Visa,Sophie,Bloomfield ,MI,United States,10/23/06 6:52,2/10/09 10:58,42.53778,-83.23306 1/5/09 12:41,Product1,1200,Visa,Katrin,Calgary,Alberta,Canada,12/3/08 14:49,2/10/09 11:45,51.0833333,-114.0833333 1/28/09 12:54,Product2,3600,Mastercard,Kelly ,Vancouver,British Columbia,Canada,1/27/09 21:04,2/10/09 12:09,49.25,-123.1333333 1/21/09 4:46,Product1,1200,Visa,Tomasz,Klampenborg,Kobenhavn,Denmark,6/10/08 11:25,2/10/09 12:22,55.7666667,12.6 1/7/09 13:28,Product1,1200,Visa,Elizabeth,Calne,England,United Kingdom,1/4/09 13:07,2/10/09 12:39,51.4333333,-2 1/27/09 11:18,Product2,3600,Amex,Michael,Los Angeles ,CA,United States,1/23/09 11:47,2/10/09 13:09,34.05222,-118.24278 1/7/09 12:39,Product2,3600,Visa,Natasha,Milano,Lombardy,Italy,6/2/06 13:01,2/10/09 13:19,45.4666667,9.2 1/24/09 13:54,Product2,3600,Mastercard,Meredith,Kloten,Zurich,Switzerland,1/24/09 12:30,2/10/09 13:47,47.45,8.5833333 1/30/09 6:48,Product1,1200,Mastercard,Nicole,Fayetteville ,NC,United States,1/30/09 4:51,2/10/09 14:41,35.0525,-78.87861 1/22/09 18:07,Product1,1200,Visa,Ryan,Simpsonville ,SC,United States,1/6/09 16:59,2/10/09 15:30,34.73694,-82.25444 1/29/09 15:03,Product1,1200,Visa,Mary ,Auckland,Auckland,New Zealand,2/9/06 11:14,2/10/09 16:31,-36.8666667,174.7666667 1/2/09 14:14,Product1,1200,Diners,Aaron,Reading,England,United Kingdom,11/16/08 15:49,2/10/09 16:38,51.4333333,-1 1/19/09 11:05,Product1,1200,Visa,Bertrand,North Caldwell ,NJ,United States,10/3/08 5:55,2/10/09 18:16,40.83972,-74.27694 \ No newline at end of file diff --git a/test/api/test_tools_upload.py b/test/api/test_tools_upload.py index a05fe9dce49..cb67b84507d 100644 --- a/test/api/test_tools_upload.py +++ b/test/api/test_tools_upload.py @@ -1,3 +1,5 @@ +import json + from base import api from base.constants import ( ONE_TO_SIX_ON_WINDOWS, @@ -33,19 +35,28 @@ class ToolsUploadTestCase(api.ApiTestCase): self._assert_has_keys(create, 'err_msg') assert file_type in create['err_msg'] - def test_upload_posix_newline_fixes(self): + # upload1 rewrites content with posix lines by default but this can be disabled by setting + # to_posix_lines=None in the request. Newer fetch API does not do this by default prefering + # to keep content unaltered if possible but it can be enabled with a simple JSON boolean switch + # of the same name (to_posix_lines). + def test_upload_posix_newline_fixes_by_default(self): windows_content = ONE_TO_SIX_ON_WINDOWS result_content = self._upload_and_get_content(windows_content) self.assertEquals(result_content, ONE_TO_SIX_WITH_TABS) + def test_fetch_posix_unaltered(self): + windows_content = ONE_TO_SIX_ON_WINDOWS + result_content = self._upload_and_get_content(windows_content, api="fetch") + self.assertEquals(result_content, ONE_TO_SIX_ON_WINDOWS) + def test_upload_disable_posix_fix(self): windows_content = ONE_TO_SIX_ON_WINDOWS result_content = self._upload_and_get_content(windows_content, to_posix_lines=None) self.assertEquals(result_content, windows_content) - def test_upload_tab_to_space(self): - table = ONE_TO_SIX_WITH_SPACES - result_content = self._upload_and_get_content(table, space_to_tab="Yes") + def test_fetch_post_lines_option(self): + windows_content = ONE_TO_SIX_ON_WINDOWS + result_content = self._upload_and_get_content(windows_content, api="fetch", to_posix_lines=True) self.assertEquals(result_content, ONE_TO_SIX_WITH_TABS) def test_upload_tab_to_space_off_by_default(self): @@ -53,6 +64,21 @@ class ToolsUploadTestCase(api.ApiTestCase): result_content = self._upload_and_get_content(table) self.assertEquals(result_content, table) + def test_fetch_tab_to_space_off_by_default(self): + table = ONE_TO_SIX_WITH_SPACES + result_content = self._upload_and_get_content(table, api='fetch') + self.assertEquals(result_content, table) + + def test_upload_tab_to_space(self): + table = ONE_TO_SIX_WITH_SPACES + result_content = self._upload_and_get_content(table, space_to_tab="Yes") + self.assertEquals(result_content, ONE_TO_SIX_WITH_TABS) + + def test_fetch_tab_to_space(self): + table = ONE_TO_SIX_WITH_SPACES + result_content = self._upload_and_get_content(table, api="fetch", space_to_tab=True) + self.assertEquals(result_content, ONE_TO_SIX_WITH_TABS) + @skip_without_datatype("rdata") def test_rdata_not_decompressed(self): # Prevent regression of https://github.com/galaxyproject/galaxy/issues/753 @@ -60,6 +86,30 @@ class ToolsUploadTestCase(api.ApiTestCase): rdata_metadata = self._upload_and_get_details(open(rdata_path, "rb"), file_type="auto") self.assertEquals(rdata_metadata["file_ext"], "rdata") + @skip_without_datatype("csv") + def test_csv_upload(self): + csv_path = TestDataResolver().get_filename("1.csv") + csv_metadata = self._upload_and_get_details(open(csv_path, "rb"), file_type="csv") + self.assertEquals(csv_metadata["file_ext"], "csv") + + @skip_without_datatype("csv") + def test_csv_upload_auto(self): + csv_path = TestDataResolver().get_filename("1.csv") + csv_metadata = self._upload_and_get_details(open(csv_path, "rb"), file_type="auto") + self.assertEquals(csv_metadata["file_ext"], "csv") + + @skip_without_datatype("csv") + def test_csv_fetch(self): + csv_path = TestDataResolver().get_filename("1.csv") + csv_metadata = self._upload_and_get_details(open(csv_path, "rb"), api="fetch", ext="csv", to_posix_lines=True) + self.assertEquals(csv_metadata["file_ext"], "csv") + + @skip_without_datatype("csv") + def test_csv_sniff_fetch(self): + csv_path = TestDataResolver().get_filename("1.csv") + csv_metadata = self._upload_and_get_details(open(csv_path, "rb"), api="fetch", ext="auto", to_posix_lines=True) + self.assertEquals(csv_metadata["file_ext"], "csv") + @skip_without_datatype("velvet") def test_composite_datatype(self): with self.dataset_populator.test_history() as history_id: @@ -113,6 +163,11 @@ class ToolsUploadTestCase(api.ApiTestCase): datasets = run_response.json()["outputs"] assert datasets[0].get("genome_build") == "hg19", datasets[0] + def test_fetch_dbkey(self): + table = ONE_TO_SIX_WITH_SPACES + details = self._upload_and_get_details(table, api='fetch', dbkey="hg19") + assert details.get("genome_build") == "hg19" + def test_upload_multiple_files_1(self): with self.dataset_populator.test_history() as history_id: payload = self.dataset_populator.upload_payload(history_id, "Test123", @@ -329,8 +384,23 @@ class ToolsUploadTestCase(api.ApiTestCase): history_id, new_dataset = self._upload(content, **upload_kwds) return self.dataset_populator.get_history_dataset_details(history_id, dataset=new_dataset) - def _upload(self, content, **upload_kwds): + def _upload(self, content, api="upload1", **upload_kwds): history_id = self.dataset_populator.new_history() - new_dataset = self.dataset_populator.new_dataset(history_id, content=content, **upload_kwds) + if api == "upload1": + new_dataset = self.dataset_populator.new_dataset(history_id, content=content, **upload_kwds) + else: + assert api == "fetch" + element = dict(src="files", **upload_kwds) + target = { + "destination": {"type": "hdas"}, + "elements": [element], + } + targets = json.dumps([target]) + payload = { + "history_id": history_id, + "targets": targets, + "__files": {"files_0|file_data": content} + } + new_dataset = self.dataset_populator.fetch(payload).json()["outputs"][0] self.dataset_populator.wait_for_history(history_id, assert_ok=upload_kwds.get("assert_ok", True)) return history_id, new_dataset diff --git a/test/functional/tools/sample_datatypes_conf.xml b/test/functional/tools/sample_datatypes_conf.xml index fc08309fd8d..29917074a99 100644 --- a/test/functional/tools/sample_datatypes_conf.xml +++ b/test/functional/tools/sample_datatypes_conf.xml @@ -4,6 +4,7 @@ + From bca2c3cf173604910cc06d53afa9e0acbb02df49 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Tue, 6 Mar 2018 19:18:59 -0500 Subject: [PATCH 24/24] Fix for data_fetch sniffing. --- lib/galaxy/tools/data_fetch.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/galaxy/tools/data_fetch.py b/lib/galaxy/tools/data_fetch.py index c6a2b7f6c61..858d0c8d234 100644 --- a/lib/galaxy/tools/data_fetch.py +++ b/lib/galaxy/tools/data_fetch.py @@ -142,7 +142,7 @@ def _fetch_target(upload_config, target): line_count, converted_path = sniff.sep2tabs(path, in_place=in_place, tmp_dir=".") if requested_ext == 'auto': - ext = sniff.guess_ext(path, registry.sniff_order) + ext = sniff.guess_ext(converted_path or path, registry.sniff_order) else: ext = requested_ext