From 72e361bcf14191dae58af8e3963b89a4e8acd38a Mon Sep 17 00:00:00 2001 From: John Chilton Date: Wed, 16 Aug 2017 12:03:07 -0400 Subject: [PATCH] Allow discovering datasets directly from galaxy.json... without needing to specify a discovered dataset pattern. If the ``discovered_datasets`` tag includes ``from_tool_provided_metadata="true"`` the datsets listed in galaxy.json will just be used directly without needing to be "discovered" using a dataset pattern. Metadata and such can continue to be specified either in that file or on the discovered_dataset element except things like pattern (which is not used) and sort_by (since the json should describe the order I suppose). --- lib/galaxy/jobs/__init__.py | 1 + lib/galaxy/tools/__init__.py | 4 +- lib/galaxy/tools/parameters/output_collect.py | 185 +++++++++++++++--- .../tools/parser/output_collection_def.py | 53 +++-- lib/galaxy/tools/xsd/galaxy.xsd | 5 + ...ction_creates_dynamic_nested_from_json.xml | 78 ++++++++ test/functional/tools/samples_tool_conf.xml | 2 + .../tools/tool_provided_metadata_7.xml | 43 ++++ 8 files changed, 332 insertions(+), 39 deletions(-) create mode 100644 test/functional/tools/collection_creates_dynamic_nested_from_json.xml create mode 100644 test/functional/tools/tool_provided_metadata_7.xml diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index bc3687eef2e..6522f2c3035 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -1391,6 +1391,7 @@ class JobWrapper(object, HasResourceParameters): } self.tool.collect_dynamic_collections( out_collections, + self.get_tool_provided_job_metadata(), job_working_directory=tool_working_directory, inp_data=inp_data, job=job, diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py index 6f4ca1c275d..0d62277bc33 100755 --- a/lib/galaxy/tools/__init__.py +++ b/lib/galaxy/tools/__init__.py @@ -1638,10 +1638,10 @@ class Tool(object, Dictifiable): """ return output_collect.collect_primary_datasets(self, output, tool_provided_metadata, job_working_directory, input_ext, input_dbkey=input_dbkey) - def collect_dynamic_collections(self, output, **kwds): + def collect_dynamic_collections(self, output, tool_provided_metadata, **kwds): """ Find files corresponding to dynamically structured collections. """ - return output_collect.collect_dynamic_collections(self, output, **kwds) + return output_collect.collect_dynamic_collections(self, output, tool_provided_metadata, **kwds) def to_archive(self): tool = self diff --git a/lib/galaxy/tools/parameters/output_collect.py b/lib/galaxy/tools/parameters/output_collect.py index fc457a0c394..9b30b937abb 100644 --- a/lib/galaxy/tools/parameters/output_collect.py +++ b/lib/galaxy/tools/parameters/output_collect.py @@ -7,6 +7,8 @@ import operator import os import re +from collections import namedtuple + from galaxy import util from galaxy.tools.parser.output_collection_def import ( DEFAULT_DATASET_COLLECTOR_DESCRIPTION, @@ -23,7 +25,9 @@ log = logging.getLogger(__name__) class NullToolProvidedMetadata(object): - pass + + def get_new_datasets(self, output_name): + return [] def get_new_dataset_meta_by_basename(self, output_name, basename): return {} @@ -66,6 +70,10 @@ class LegacyToolProvidedMetadata(object): if meta['type'] == 'new_primary_dataset' and meta['filename'] == basename: return meta + def get_new_datasets(self, output_name): + log.warning("Called get_new_datasets with legacy tool metadata provider - that is unimplemented.") + return [] + class ToolProvidedMetadata(object): @@ -83,10 +91,15 @@ class ToolProvidedMetadata(object): if meta['filename'] == basename: return meta + def get_new_datasets(self, output_name): + datasets = self.tool_provided_job_metadata.get(output_name, {}).get("datasets", []) + return datasets + def collect_dynamic_collections( tool, output_collections, + tool_provided_metadata, job_working_directory, inp_data={}, job=None, @@ -95,6 +108,7 @@ def collect_dynamic_collections( collections_service = tool.app.dataset_collections_service job_context = JobContext( tool, + tool_provided_metadata, job, job_working_directory, inp_data, @@ -132,13 +146,14 @@ def collect_dynamic_collections( class JobContext(object): - def __init__(self, tool, job, job_working_directory, inp_data, input_dbkey): + def __init__(self, tool, tool_provided_metadata, job, job_working_directory, inp_data, input_dbkey): self.inp_data = inp_data self.input_dbkey = input_dbkey self.app = tool.app self.sa_session = tool.sa_session self.job = job self.job_working_directory = job_working_directory + self.tool_provided_metadata = tool_provided_metadata @property def permissions(self): @@ -151,10 +166,10 @@ class JobContext(object): permissions = self.app.security_agent.history_get_default_permissions(self.job.history) return permissions - def find_files(self, collection, dataset_collectors): + def find_files(self, output_name, collection, dataset_collectors): filenames = odict.odict() - for path, extra_file_collector in walk_over_extra_files(dataset_collectors, self.job_working_directory, collection): - filenames[path] = extra_file_collector + for discovered_file in discover_files(output_name, self.tool_provided_metadata, dataset_collectors, self.job_working_directory, collection): + filenames[discovered_file.path] = discovered_file return filenames def populate_collection_elements(self, collection, root_collection_builder, output_collection_def): @@ -164,12 +179,13 @@ class JobContext(object): # # dataset_collectors = map(dataset_collector, output_collection_def.dataset_collector_descriptions) - filenames = self.find_files(collection, dataset_collectors) + output_name = output_collection_def.name + filenames = self.find_files(output_name, collection, dataset_collectors) element_datasets = [] - for filename, extra_file_collector in filenames.items(): + for filename, discovered_file in filenames.items(): create_dataset_timer = ExecutionTimer() - fields_match = extra_file_collector.match(collection, os.path.basename(filename)) + fields_match = discovered_file.match if not fields_match: raise Exception("Problem parsing metadata fields for file %s" % filename) element_identifiers = fields_match.element_identifiers @@ -292,8 +308,7 @@ def collect_primary_datasets(tool, output, tool_provided_metadata, job_working_d # This should not be considered an error or warning condition, this file is optional pass # Loop through output file names, looking for generated primary - # datasets in form of: - # 'primary_associatedWithDatasetID_designation_visibility_extension(_DBKEY)' + # datasets in form specified by discover dataset patterns or in tool provided metadata. primary_output_assigned = False new_outdata_name = None primary_datasets = {} @@ -306,12 +321,17 @@ def collect_primary_datasets(tool, output, tool_provided_metadata, job_working_d # only use old-style matching (glob instead of regex and only # using default collector - if enabled). for filename in glob.glob(os.path.join(app.config.new_file_path, "primary_%i_*" % outdata.id)): - filenames[filename] = DEFAULT_DATASET_COLLECTOR + filenames[filename] = DiscoveredFile( + filename, + DEFAULT_DATASET_COLLECTOR, + DEFAULT_DATASET_COLLECTOR.match(outdata, os.path.basename(filename)) + ) if 'job_working_directory' in app.config.collect_outputs_from: - for path, extra_file_collector in walk_over_extra_files(dataset_collectors, job_working_directory, outdata): - filenames[path] = extra_file_collector - for filename_index, (filename, extra_file_collector) in enumerate(filenames.items()): - fields_match = extra_file_collector.match(outdata, os.path.basename(filename)) + for discovered_file in discover_files(name, tool_provided_metadata, dataset_collectors, job_working_directory, outdata): + filenames[discovered_file.path] = discovered_file + for filename_index, (filename, discovered_file) in enumerate(filenames.items()): + extra_file_collector = discovered_file.collector + fields_match = discovered_file.match if not fields_match: # Before I guess pop() would just have thrown an IndexError raise Exception("Problem parsing metadata fields for file %s" % filename) @@ -411,14 +431,39 @@ def collect_primary_datasets(tool, output, tool_provided_metadata, job_working_d return primary_datasets +DiscoveredFile = namedtuple('DiscoveredFile', ['path', 'collector', 'match']) + + +def discover_files(output_name, tool_provided_metadata, extra_file_collectors, job_working_directory, matchable): + if extra_file_collectors and extra_file_collectors[0].discover_via == "tool_provided_metadata": + # just load entries from tool provided metadata... + assert len(extra_file_collectors) == 1 + extra_file_collector = extra_file_collectors[0] + target_directory = discover_target_directory(extra_file_collector, job_working_directory) + for dataset in tool_provided_metadata.get_new_datasets(output_name): + filename = dataset["filename"] + path = os.path.join(target_directory, filename) + yield DiscoveredFile(path, extra_file_collector, JsonCollectedDatasetMatch(dataset, extra_file_collector, filename, path=path)) + else: + for (match, collector) in walk_over_extra_files(extra_file_collectors, job_working_directory, matchable): + yield DiscoveredFile(match.path, collector, match) + + +def discover_target_directory(extra_file_collector, job_working_directory): + directory = job_working_directory + if extra_file_collector.directory: + directory = os.path.join(directory, extra_file_collector.directory) + if not util.in_directory(directory, job_working_directory): + raise Exception("Problem with tool configuration, attempting to pull in datasets from outside working directory.") + return directory + + def walk_over_extra_files(extra_file_collectors, job_working_directory, matchable): + for extra_file_collector in extra_file_collectors: + assert extra_file_collector.discover_via == "pattern" matches = [] - directory = job_working_directory - if extra_file_collector.directory: - directory = os.path.join(directory, extra_file_collector.directory) - if not util.in_directory(directory, job_working_directory): - raise Exception("Problem with tool configuration, attempting to pull in datasets from outside working directory.") + directory = discover_target_directory(extra_file_collector, job_working_directory) if not os.path.isdir(directory): continue for filename in os.listdir(directory): @@ -430,7 +475,7 @@ def walk_over_extra_files(extra_file_collectors, job_working_directory, matchabl matches.append(match) for match in extra_file_collector.sort(matches): - yield match.path, extra_file_collector + yield match, extra_file_collector def dataset_collector(dataset_collection_description): @@ -439,12 +484,27 @@ def dataset_collector(dataset_collection_description): # treated like a singleton. return DEFAULT_DATASET_COLLECTOR else: - return DatasetCollector(dataset_collection_description) + if dataset_collection_description.discover_via == "pattern": + return DatasetCollector(dataset_collection_description) + else: + return ToolMetadataDatasetCollector(dataset_collection_description) + + +class ToolMetadataDatasetCollector(object): + + def __init__(self, dataset_collection_description): + self.discover_via = dataset_collection_description.discover_via + self.default_dbkey = dataset_collection_description.default_dbkey + self.default_ext = dataset_collection_description.default_ext + self.default_visible = dataset_collection_description.default_visible + self.directory = dataset_collection_description.directory + self.assign_primary_output = dataset_collection_description.assign_primary_output class DatasetCollector(object): def __init__(self, dataset_collection_description): + self.discover_via = dataset_collection_description.discover_via # dataset_collection_description is an abstract description # built from the tool parsing module - see galaxy.tools.parser.output_colleciton_def self.sort_key = dataset_collection_description.sort_key @@ -457,18 +517,18 @@ class DatasetCollector(object): self.directory = dataset_collection_description.directory self.assign_primary_output = dataset_collection_description.assign_primary_output - def pattern_for_dataset(self, dataset_instance=None): + def _pattern_for_dataset(self, dataset_instance=None): token_replacement = r'\d+' if dataset_instance: token_replacement = str(dataset_instance.id) return self.pattern.replace(DATASET_ID_TOKEN, token_replacement) def match(self, dataset_instance, filename, path=None): - pattern = self.pattern_for_dataset(dataset_instance) + pattern = self._pattern_for_dataset(dataset_instance) re_match = re.match(pattern, filename) match_object = None if re_match: - match_object = CollectedDatasetMatch(re_match, self, filename, path=path) + match_object = RegexCollectedDatasetMatch(re_match, self, filename, path=path) return match_object def sort(self, matches): @@ -488,7 +548,80 @@ def _compose(f, g): return lambda x: f(g(x)) -class CollectedDatasetMatch(object): +class JsonCollectedDatasetMatch(object): + + def __init__(self, as_dict, collector, filename, path=None): + self.as_dict = as_dict + self.collector = collector + self.filename = filename + self.path = path + + @property + def designation(self): + # If collecting nested collection, grap identifier_0, + # identifier_1, etc... and join on : to build designation. + element_identifiers = self.raw_element_identifiers + if element_identifiers: + return ":".join(element_identifiers) + elif "designation" in self.as_dict: + return self.as_dict.get("designation") + elif "name" in self.as_dict: + return self.as_dict.get("name") + else: + return None + + @property + def element_identifiers(self): + return self.raw_element_identifiers or [self.designation] + + @property + def raw_element_identifiers(self): + identifiers = [] + i = 0 + while True: + key = "identifier_%d" % i + if key in self.as_dict: + identifiers.append(self.as_dict.get(key)) + else: + break + i += 1 + + return identifiers + + @property + def name(self): + """ Return name or None if not defined by the discovery pattern. + """ + name = None + if "name" in self.as_dict: + name = self.as_dict.get("name") + return name + + @property + def dbkey(self): + try: + return self.as_dict["dbkey"] + except KeyError: + return self.collector.default_dbkey + + @property + def ext(self): + try: + return self.as_dict["ext"] + except KeyError: + return self.collector.default_ext + + @property + def visible(self): + try: + return self.as_dict["visible"].lower() == "visible" + except KeyError: + return self.collector.default_visible + + +class RegexCollectedDatasetMatch(object): + # TODO: This could probably subclass JsonCollectedDatasetMatch if group dict is + # treated the same as the JSON dict. def __init__(self, re_match, collector, filename, path=None): self.re_match = re_match diff --git a/lib/galaxy/tools/parser/output_collection_def.py b/lib/galaxy/tools/parser/output_collection_def.py index 8c79d49e5f7..42c07faf533 100644 --- a/lib/galaxy/tools/parser/output_collection_def.py +++ b/lib/galaxy/tools/parser/output_collection_def.py @@ -25,31 +25,62 @@ LEGACY_DEFAULT_DBKEY = None # don't use __input__ for legacy default collection def dataset_collector_descriptions_from_elem(elem, legacy=True): primary_dataset_elems = elem.findall("discover_datasets") - if len(primary_dataset_elems) == 0 and legacy: - return [DEFAULT_DATASET_COLLECTOR_DESCRIPTION] + num_discover_dataset_blocks = len(primary_dataset_elems) + if num_discover_dataset_blocks == 0 and legacy: + collectors = [DEFAULT_DATASET_COLLECTOR_DESCRIPTION] else: - return map(lambda elem: DatasetCollectionDescription(**elem.attrib), primary_dataset_elems) + collectors = map(lambda elem: dataset_collection_description(**elem.attrib), primary_dataset_elems) + + if num_discover_dataset_blocks > 1: + for collector in collectors: + if collector.discover_via == "tool_provided_metadata": + raise Exception("Cannot specify more than one discover dataset condition if any of them specify tool_provided_metadata.") + + return collectors def dataset_collector_descriptions_from_list(discover_datasets_dicts): - return map(lambda kwds: DatasetCollectionDescription(**kwds), discover_datasets_dicts) + return map(lambda kwds: dataset_collection_description(**kwds), discover_datasets_dicts) + + +def dataset_collection_description(**kwargs): + if asbool(kwargs.get("from_provided_metadata", False)): + for key in ["pattern", "sort_by"]: + if kwargs.get(key): + raise Exception("Cannot specify attribute [%s] if from_provided_metadata is True" % key) + return ToolProvidedMetadataDatasetCollection(**kwargs) + else: + return FilePatternDatasetCollectionDescription(**kwargs) class DatasetCollectionDescription(object): def __init__(self, **kwargs): - pattern = kwargs.get("pattern", "__default__") - if pattern in NAMED_PATTERNS: - pattern = NAMED_PATTERNS.get(pattern) - self.pattern = pattern self.default_dbkey = kwargs.get("dbkey", INPUT_DBKEY_TOKEN) self.default_ext = kwargs.get("ext", None) if self.default_ext is None and "format" in kwargs: self.default_ext = kwargs.get("format") self.default_visible = asbool(kwargs.get("visible", None)) - self.directory = kwargs.get("directory", None) self.assign_primary_output = asbool(kwargs.get('assign_primary_output', False)) - sort_by = kwargs.get("sort_by", DEFAULT_SORT_BY) + self.directory = kwargs.get("directory", None) + + +class ToolProvidedMetadataDatasetCollection(DatasetCollectionDescription): + + discover_via = "tool_provided_metadata" + + +class FilePatternDatasetCollectionDescription(DatasetCollectionDescription): + + discover_via = "pattern" + + def __init__( self, **kwargs ): + super(FilePatternDatasetCollectionDescription, self).__init__( **kwargs ) + pattern = kwargs.get( "pattern", "__default__" ) + if pattern in NAMED_PATTERNS: + pattern = NAMED_PATTERNS.get( pattern ) + self.pattern = pattern + sort_by = kwargs.get( "sort_by", DEFAULT_SORT_BY ) if sort_by.startswith("reverse_"): self.sort_reverse = True sort_by = sort_by[len("reverse_"):] @@ -70,6 +101,6 @@ class DatasetCollectionDescription(object): self.sort_comp = sort_comp -DEFAULT_DATASET_COLLECTOR_DESCRIPTION = DatasetCollectionDescription( +DEFAULT_DATASET_COLLECTOR_DESCRIPTION = FilePatternDatasetCollectionDescription( default_dbkey=LEGACY_DEFAULT_DBKEY, ) diff --git a/lib/galaxy/tools/xsd/galaxy.xsd b/lib/galaxy/tools/xsd/galaxy.xsd index 0abeed028e0..a6e350a2925 100644 --- a/lib/galaxy/tools/xsd/galaxy.xsd +++ b/lib/galaxy/tools/xsd/galaxy.xsd @@ -3873,6 +3873,11 @@ More information can be found on Planemo's documentation for [multiple output files](https://planemo.readthedocs.io/en/latest/writing_advanced.html#multiple-output-files). ]]> + + + Indicate that dataset filenames should simply be read from the provided metadata file (e.g. galaxy.json). If this is set - pattern and sort must not be set. + + Regular expression used to find filenames and parse dynamic properties. diff --git a/test/functional/tools/collection_creates_dynamic_nested_from_json.xml b/test/functional/tools/collection_creates_dynamic_nested_from_json.xml new file mode 100644 index 00000000000..5cd6391f9b7 --- /dev/null +++ b/test/functional/tools/collection_creates_dynamic_nested_from_json.xml @@ -0,0 +1,78 @@ + + + echo "A" > oe1_ie1.fq ; + echo "B" > oe1_ie2.fq ; + echo "C" > oe2_ie1.fq ; + echo "D" > oe2_ie2.fq ; + echo "E" > oe3_ie1.fq ; + echo "F" > oe3_ie2.fq ; + cp $c1 galaxy.json + + + {"list_output": { + "datasets": [ + {"identifier_0": "oe1", "identifier_1": "ie1", "filename": "oe1_ie1.fq"}, + {"identifier_0": "oe1", "identifier_1": "ie2", "filename": "oe1_ie2.fq"}, + {"identifier_0": "oe2", "identifier_1": "ie1", "filename": "oe2_ie1.fq"}, + {"identifier_0": "oe2", "identifier_1": "ie2", "filename": "oe2_ie2.fq"}, + {"identifier_0": "oe3", "identifier_1": "ie1", "filename": "oe3_ie1.fq"}, + {"identifier_0": "oe3", "identifier_1": "ie2", "filename": "oe3_ie2.fq"} +]}} + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + diff --git a/test/functional/tools/samples_tool_conf.xml b/test/functional/tools/samples_tool_conf.xml index dc612199ce8..4fa4feee5f0 100644 --- a/test/functional/tools/samples_tool_conf.xml +++ b/test/functional/tools/samples_tool_conf.xml @@ -21,6 +21,7 @@ + @@ -103,6 +104,7 @@ + diff --git a/test/functional/tools/tool_provided_metadata_7.xml b/test/functional/tools/tool_provided_metadata_7.xml new file mode 100644 index 00000000000..a50adbfe2c3 --- /dev/null +++ b/test/functional/tools/tool_provided_metadata_7.xml @@ -0,0 +1,43 @@ + + + + echo "1" > sample1.report.tsv; + echo "2" > sample2.report.tsv; + cp $c1 galaxy.json; + + + {"sample": { +"datasets": [ +{"filename": "sample1.report.tsv", "designation": "sample1", "name": "cool name 1", "ext": "txt", "info": "cool 1 info", "dbkey": "hg19"}, +{"filename": "sample2.report.tsv", "designation": "sample2", "name": "cool name 2", "ext": "txt", "info": "cool 2 info", "dbkey": "hg19"} +] +}} + + + + + + + + + + + + + + + + + + + + + + + + + + + + +