diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py
index 8b80400599b..bc3687eef2e 100644
--- a/lib/galaxy/jobs/__init__.py
+++ b/lib/galaxy/jobs/__init__.py
@@ -1213,7 +1213,7 @@ class JobWrapper(object, HasResourceParameters):
job_context = ExpressionContext(dict(stdout=job.stdout, stderr=job.stderr))
for dataset_assoc in job.output_datasets + job.output_library_datasets:
- context = self.get_dataset_finish_context(job_context, dataset_assoc.dataset.dataset)
+ context = self.get_dataset_finish_context(job_context, dataset_assoc)
# should this also be checking library associations? - can a library item be added from a history before the job has ended? -
# lets not allow this to occur
# need to update all associated output hdas, i.e. history was shared with job running
@@ -1387,7 +1387,7 @@ class JobWrapper(object, HasResourceParameters):
tool_working_directory = self.working_directory
collected_datasets = {
'children': self.tool.collect_child_datasets(out_data, tool_working_directory),
- 'primary': self.tool.collect_primary_datasets(out_data, tool_working_directory, input_ext, input_dbkey)
+ 'primary': self.tool.collect_primary_datasets(out_data, self.get_tool_provided_job_metadata(), tool_working_directory, input_ext, input_dbkey)
}
self.tool.collect_dynamic_collections(
out_collections,
@@ -1621,6 +1621,7 @@ class JobWrapper(object, HasResourceParameters):
if self.tool_provided_job_metadata is not None:
return self.tool_provided_job_metadata
+<<<<<<< HEAD
# Look for JSONified job metadata
self.tool_provided_job_metadata = []
meta_file = os.path.join(self.tool_working_directory, TOOL_PROVIDED_JOB_METADATA_FILE)
@@ -1655,6 +1656,21 @@ class JobWrapper(object, HasResourceParameters):
for meta in self.get_tool_provided_job_metadata():
if meta['type'] == 'dataset' and meta['dataset_id'] == dataset.id:
return ExpressionContext(meta, job_context)
+=======
+ self.tool_provided_job_metadata = self.tool.tool_provided_metadata( self )
+ return self.tool_provided_job_metadata
+
+ def get_dataset_finish_context( self, job_context, output_dataset_assoc ):
+ meta = {}
+ tool_provided_metadata = self.get_tool_provided_job_metadata()
+ if hasattr(tool_provided_metadata, "get_meta_by_dataset_id"):
+ meta = tool_provided_metadata.get_meta_by_dataset_id(output_dataset_assoc.dataset.dataset.id)
+ elif hasattr(tool_provided_metadata, "get_meta_by_name"):
+ meta = tool_provided_metadata.get_meta_by_name(output_dataset_assoc.name)
+
+ if meta:
+ return ExpressionContext( meta, job_context )
+>>>>>>> a73a0275a9... Overhaul of tool provided job metadata.
return job_context
def invalidate_external_metadata(self):
@@ -1674,8 +1690,13 @@ class JobWrapper(object, HasResourceParameters):
if set_extension:
for output_dataset_assoc in job.output_datasets:
if output_dataset_assoc.dataset.ext == 'auto':
+<<<<<<< HEAD
context = self.get_dataset_finish_context(dict(), output_dataset_assoc.dataset.dataset)
output_dataset_assoc.dataset.extension = context.get('ext', 'data')
+=======
+ context = self.get_dataset_finish_context( dict(), output_dataset_assoc )
+ output_dataset_assoc.dataset.extension = context.get( 'ext', 'data' )
+>>>>>>> a73a0275a9... Overhaul of tool provided job metadata.
self.sa_session.flush()
if tmp_dir is None:
# this dir should should relative to the exec_dir
diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py
index 77d7b840d82..6f4ca1c275d 100755
--- a/lib/galaxy/tools/__init__.py
+++ b/lib/galaxy/tools/__init__.py
@@ -805,6 +805,24 @@ class Tool(object, Dictifiable):
self.__tests_populated = True
return self.__tests
+ @property
+ def _legacy_tool_provided_metadata_format(self):
+ return self.profile < 17.09
+
+ def tool_provided_metadata(self, job_wrapper):
+ meta_file = os.path.join(job_wrapper.tool_working_directory, galaxy.jobs.TOOL_PROVIDED_JOB_METADATA_FILE)
+ # LEGACY: Remove in 17.XX
+ if not os.path.exists(meta_file):
+ # Maybe this is a legacy job, use the job working directory instead
+ meta_file = os.path.join(job_wrapper.working_directory, galaxy.jobs.TOOL_PROVIDED_JOB_METADATA_FILE)
+
+ if not os.path.exists(meta_file):
+ return output_collect.NullToolProvidedMetadata()
+ if self._legacy_tool_provided_metadata_format:
+ return output_collect.LegacyToolProvidedMetadata(job_wrapper, meta_file)
+ else:
+ return output_collect.ToolProvidedMetadata(job_wrapper, meta_file)
+
def parse_command(self, tool_source):
"""
"""
@@ -1613,12 +1631,12 @@ class Tool(object, Dictifiable):
self.sa_session.flush()
return children
- def collect_primary_datasets(self, output, job_working_directory, input_ext, input_dbkey="?"):
+ def collect_primary_datasets(self, output, tool_provided_metadata, job_working_directory, input_ext, input_dbkey="?"):
"""
Find any additional datasets generated by a tool and attach (for
cases where number of outputs is not known in advance).
"""
- return output_collect.collect_primary_datasets(self, output, job_working_directory, input_ext, input_dbkey=input_dbkey)
+ 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):
""" Find files corresponding to dynamically structured collections.
diff --git a/lib/galaxy/tools/parameters/output_collect.py b/lib/galaxy/tools/parameters/output_collect.py
index 3b68868aa42..fc457a0c394 100644
--- a/lib/galaxy/tools/parameters/output_collect.py
+++ b/lib/galaxy/tools/parameters/output_collect.py
@@ -7,7 +7,6 @@ import operator
import os
import re
-from galaxy import jobs
from galaxy import util
from galaxy.tools.parser.output_collection_def import (
DEFAULT_DATASET_COLLECTOR_DESCRIPTION,
@@ -23,6 +22,68 @@ DATASET_ID_TOKEN = "DATASET_ID"
log = logging.getLogger(__name__)
+class NullToolProvidedMetadata(object):
+ pass
+
+ def get_new_dataset_meta_by_basename(self, output_name, basename):
+ return {}
+
+
+class LegacyToolProvidedMetadata(object):
+
+ def __init__(self, job_wrapper, meta_file):
+ self.job_wrapper = job_wrapper
+ self.tool_provided_job_metadata = []
+
+ with open(meta_file, 'r') as f:
+ for line in f:
+ try:
+ line = json.loads(line)
+ assert 'type' in line
+ except:
+ log.exception('(%s) Got JSON data from tool, but data is improperly formatted or no "type" key in data' % job_wrapper.job_id)
+ log.debug('Offending data was: %s' % line)
+ continue
+ # Set the dataset id if it's a dataset entry and isn't set.
+ # This isn't insecure. We loop the job's output datasets in
+ # the finish method, so if a tool writes out metadata for a
+ # dataset id that it doesn't own, it'll just be ignored.
+ if line['type'] == 'dataset' and 'dataset_id' not in line:
+ try:
+ line['dataset_id'] = job_wrapper.get_output_file_id(line['dataset'])
+ except KeyError:
+ log.warning('(%s) Tool provided job dataset-specific metadata without specifying a dataset' % job_wrapper.job_id)
+ continue
+ self.tool_provided_job_metadata.append(line)
+
+ def get_meta_by_dataset_id(self, dataset_id):
+ for meta in self.tool_provided_job_metadata:
+ if meta['type'] == 'dataset' and meta['dataset_id'] == dataset_id:
+ return meta
+
+ def get_new_dataset_meta_by_basename(self, output_name, basename):
+ for meta in self.tool_provided_job_metadata:
+ if meta['type'] == 'new_primary_dataset' and meta['filename'] == basename:
+ return meta
+
+
+class ToolProvidedMetadata(object):
+
+ def __init__(self, job_wrapper, meta_file):
+ self.job_wrapper = job_wrapper
+ with open(meta_file, 'r') as f:
+ self.tool_provided_job_metadata = json.load(f)
+
+ def get_meta_by_name(self, name):
+ return self.tool_provided_job_metadata.get(name, {})
+
+ def get_new_dataset_meta_by_basename(self, output_name, basename):
+ datasets = self.tool_provided_job_metadata.get(output_name, {}).get("datasets", [])
+ for meta in datasets:
+ if meta['filename'] == basename:
+ return meta
+
+
def collect_dynamic_collections(
tool,
output_collections,
@@ -164,7 +225,7 @@ class JobContext(object):
# Associate new dataset with job
if job:
element_identifier_str = ":".join(element_identifiers)
- # Below was changed from '__new_primary_file_%s|%s__' % ( name, designation )
+ # Below was changed from '__new_primary_file_%s|%s__' % (name, designation )
assoc = app.model.JobToOutputDatasetAssociation('__new_primary_file_%s|%s__' % (name, element_identifier_str), dataset)
assoc.job = self.job
sa_session.add(assoc)
@@ -212,7 +273,7 @@ class JobContext(object):
return primary_data
-def collect_primary_datasets(tool, output, job_working_directory, input_ext, input_dbkey="?"):
+def collect_primary_datasets(tool, output, tool_provided_metadata, job_working_directory, input_ext, input_dbkey="?"):
app = tool.app
sa_session = tool.sa_session
new_primary_datasets = {}
@@ -294,8 +355,10 @@ def collect_primary_datasets(tool, output, job_working_directory, input_ext, inp
sa_session.add(assoc)
sa_session.flush()
primary_data.state = outdata.state
+ # TODO: should be able to disambiguate files in different directories...
+ new_primary_filename = os.path.split(filename)[-1]
+ new_primary_datasets_attributes = tool_provided_metadata.get_new_dataset_meta_by_basename(name, new_primary_filename)
# add tool/metadata provided information
- new_primary_datasets_attributes = new_primary_datasets.get(os.path.split(filename)[-1], {})
if new_primary_datasets_attributes:
dataset_att_by_name = dict(ext='extension')
for att_set in ['name', 'info', 'ext', 'dbkey']:
diff --git a/test/functional/tools/samples_tool_conf.xml b/test/functional/tools/samples_tool_conf.xml
index 916ce3f5f17..dc612199ce8 100644
--- a/test/functional/tools/samples_tool_conf.xml
+++ b/test/functional/tools/samples_tool_conf.xml
@@ -18,6 +18,9 @@
+
+
+
diff --git a/test/functional/tools/tool_provided_metadata_1.xml b/test/functional/tools/tool_provided_metadata_1.xml
index 02d971a85f4..a573a7091a9 100644
--- a/test/functional/tools/tool_provided_metadata_1.xml
+++ b/test/functional/tools/tool_provided_metadata_1.xml
@@ -1,4 +1,5 @@
+
echo "This is a line of text." > $out1;
cp $c1 galaxy.json;
diff --git a/test/functional/tools/tool_provided_metadata_2.xml b/test/functional/tools/tool_provided_metadata_2.xml
index 8879d0f62af..9eca149cb81 100644
--- a/test/functional/tools/tool_provided_metadata_2.xml
+++ b/test/functional/tools/tool_provided_metadata_2.xml
@@ -1,4 +1,5 @@
+
echo "1" > sample1.report.tsv;
echo "2" > sample2.report.tsv;
diff --git a/test/functional/tools/tool_provided_metadata_3.xml b/test/functional/tools/tool_provided_metadata_3.xml
index a9a64f6ee82..f3ec2ac29e4 100644
--- a/test/functional/tools/tool_provided_metadata_3.xml
+++ b/test/functional/tools/tool_provided_metadata_3.xml
@@ -1,4 +1,5 @@
+
echo "1" > sample1.report.tsv;
echo "2" > sample2.report.tsv;
diff --git a/test/functional/tools/tool_provided_metadata_4.xml b/test/functional/tools/tool_provided_metadata_4.xml
new file mode 100644
index 00000000000..15cafddcc6a
--- /dev/null
+++ b/test/functional/tools/tool_provided_metadata_4.xml
@@ -0,0 +1,31 @@
+
+
+
+ echo "This is a line of text." > '$out1';
+ cp $c1 galaxy.json;
+
+
+ #import os
+{"type": "dataset", "dataset": "${os.path.basename(str($out1))}", "name": "my dynamic name", "ext": "txt", "info": "my dynamic info", "dbkey": "cust1"}
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
diff --git a/test/functional/tools/tool_provided_metadata_5.xml b/test/functional/tools/tool_provided_metadata_5.xml
new file mode 100644
index 00000000000..33b527c4f30
--- /dev/null
+++ b/test/functional/tools/tool_provided_metadata_5.xml
@@ -0,0 +1,36 @@
+
+
+
+ echo "This is a line of text." > $out1;
+ cp $c1 galaxy.json;
+
+
+ {"out1": {
+ "name": "my dynamic name",
+ "ext": "txt",
+ "info": "my dynamic info",
+ "dbkey": "cust1"
+}}
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
diff --git a/test/functional/tools/tool_provided_metadata_6.xml b/test/functional/tools/tool_provided_metadata_6.xml
new file mode 100644
index 00000000000..78e9f3c36ca
--- /dev/null
+++ b/test/functional/tools/tool_provided_metadata_6.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", "name": "cool name 1", "ext": "txt", "info": "cool 1 info", "dbkey": "hg19"},
+{"filename": "sample2.report.tsv", "name": "cool name 2", "ext": "txt", "info": "cool 2 info", "dbkey": "hg19"}
+]
+}}
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
diff --git a/test/unit/tools/test_collect_primary_datasets.py b/test/unit/tools/test_collect_primary_datasets.py
index 51e1dbdde5c..b5317e1b358 100644
--- a/test/unit/tools/test_collect_primary_datasets.py
+++ b/test/unit/tools/test_collect_primary_datasets.py
@@ -7,6 +7,7 @@ from galaxy import (
model,
util
)
+from galaxy.tools.parameters.output_collect import LegacyToolProvidedMetadata, NullToolProvidedMetadata
from galaxy.tools.parser import output_collection_def
@@ -250,7 +251,13 @@ class CollectPrimaryDatasetsTestCase(unittest.TestCase, tools_support.UsesApp, t
def _collect(self, job_working_directory=None):
if not job_working_directory:
job_working_directory = self.test_directory
- return self.tool.collect_primary_datasets(self.outputs, job_working_directory, "txt", input_dbkey="btau")
+ meta_file = os.path.join(self.test_directory, "galaxy.json")
+ if not os.path.exists(meta_file):
+ tool_provided_metadata = NullToolProvidedMetadata()
+ else:
+ tool_provided_metadata = LegacyToolProvidedMetadata(None, meta_file)
+
+ return self.tool.collect_primary_datasets(self.outputs, tool_provided_metadata, job_working_directory, "txt", input_dbkey="btau")
def _replace_output_collectors(self, xml_str):
# Rewrite tool as if it had been created with output containing