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"}
+]
+}}
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+