From 517f248ac49caca3f0933c4b1f48cd4d74cf89fe Mon Sep 17 00:00:00 2001 From: John Chilton Date: Wed, 16 Aug 2017 08:48:32 -0400 Subject: [PATCH] Overhaul of tool provided job metadata. - Add another test to demonstrate referencing entries in galaxy.json for datasets by dataset path basename instead of ID. - Add a bit more documentation to test tools. - Redo the structure of galaxy.json for profile >= 17.09 tools with two new test tools to demonstrate the new structure - for both output dataset metadata and discovered dataset metadata. - Improved abstractions for interaction between tools and jobs to enable different tool provided metadata collection depending on tool profile version. --- lib/galaxy/jobs/__init__.py | 25 ++++++- lib/galaxy/tools/__init__.py | 22 +++++- lib/galaxy/tools/parameters/output_collect.py | 71 +++++++++++++++++-- test/functional/tools/samples_tool_conf.xml | 3 + .../tools/tool_provided_metadata_1.xml | 1 + .../tools/tool_provided_metadata_2.xml | 1 + .../tools/tool_provided_metadata_3.xml | 1 + .../tools/tool_provided_metadata_4.xml | 31 ++++++++ .../tools/tool_provided_metadata_5.xml | 36 ++++++++++ .../tools/tool_provided_metadata_6.xml | 43 +++++++++++ .../tools/test_collect_primary_datasets.py | 9 ++- 11 files changed, 234 insertions(+), 9 deletions(-) create mode 100644 test/functional/tools/tool_provided_metadata_4.xml create mode 100644 test/functional/tools/tool_provided_metadata_5.xml create mode 100644 test/functional/tools/tool_provided_metadata_6.xml 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