From 7d260d0b50a5858af533fe136105ce357b4ef490 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Wed, 27 Sep 2017 18:29:04 +0200 Subject: [PATCH] Tighten job search result We now expand incoming tool parameters, so that we have a dictionary of all tool parameters (i.e specified/posted parameters and defaults). We start the search by quering the database for jobs using the same tool id for the same owner with the same input datasets, taking into account copied dataset(-collection)s. We then replace all dataset ids (HDA/LDA or HDCA) in the parameter dict with the input data used for the jobs queried in the first step. This allows us to check every parameter for a match with the job. --- .../dependencies/pinned-requirements.txt | 1 + lib/galaxy/dependencies/requirements.txt | 1 + lib/galaxy/managers/jobs.py | 235 ++++++++++++------ lib/galaxy/tools/__init__.py | 27 +- lib/galaxy/util/__init__.py | 3 + lib/galaxy/webapps/galaxy/api/jobs.py | 23 +- 6 files changed, 200 insertions(+), 90 deletions(-) diff --git a/lib/galaxy/dependencies/pinned-requirements.txt b/lib/galaxy/dependencies/pinned-requirements.txt index 4c4f7601c3b..50723015782 100644 --- a/lib/galaxy/dependencies/pinned-requirements.txt +++ b/lib/galaxy/dependencies/pinned-requirements.txt @@ -14,6 +14,7 @@ uWSGI==2.0.15 # pure Python packages bz2file==0.98; python_version < '3.3' +boltons==17.1.0 Paste==2.0.2 PasteDeploy==1.5.2 docutils==0.12 diff --git a/lib/galaxy/dependencies/requirements.txt b/lib/galaxy/dependencies/requirements.txt index f28cbc732c0..b75706accb5 100644 --- a/lib/galaxy/dependencies/requirements.txt +++ b/lib/galaxy/dependencies/requirements.txt @@ -14,6 +14,7 @@ pycrypto # pure Python packages bz2file; python_version < '3.3' +boltons Paste PasteDeploy docutils diff --git a/lib/galaxy/managers/jobs.py b/lib/galaxy/managers/jobs.py index 4e5c3b2187a..a291d377b23 100644 --- a/lib/galaxy/managers/jobs.py +++ b/lib/galaxy/managers/jobs.py @@ -1,13 +1,40 @@ import json +import logging +from boltons.iterutils import remap from six import string_types -from sqlalchemy import and_, or_ +from sqlalchemy import and_, false, or_ from sqlalchemy.orm import aliased from galaxy import model from galaxy.managers.collections import DatasetCollectionManager from galaxy.managers.hdas import HDAManager from galaxy.managers.lddas import LDDAManager +from galaxy.util import ( + defaultdict, + ExecutionTimer +) + +log = logging.getLogger(__name__) + + +def get_path_key(path_tuple): + path_key = "" + tuple_elements = len(path_tuple) + for i, p in enumerate(path_tuple): + if isinstance(p, int): + sep = '_' + else: + sep = '|' + if i == (tuple_elements - 2) and p == 'values': + # dataset inputs are always wrapped in lists. To avoid 'rep_factorName_0|rep_factorLevel_2|countsFile|values_0', + # we remove the last 2 items of the path tuple (values and list index) + return path_key + if path_key: + path_key = "%s%s%s" % (path_key, sep, p) + else: + path_key = p + return path_key class JobSearch(object): @@ -20,69 +47,36 @@ class JobSearch(object): self.ldda_manager = LDDAManager(app) self.decode_id = self.app.security.decode_id - def by_tool_input(self, trans, tool_id, inputs, job_state='ok'): + def by_tool_input(self, trans, tool_id, param_dump=None, job_state='ok', is_workflow_step=False): """Search for jobs producing same results using the 'inputs' part of a tool POST.""" user = trans.user - input_param = {} - input_data = {} - for k, v in inputs.items(): - if isinstance(v, dict): - if 'id' in v: - decoded_id = self.decode_id(v['id']) - src = v.get('src', 'hda') - if src == 'hda': - all_items = self._get_all_hdas(trans=trans, hda_id=decoded_id) - if not self.one_not_deleted(all_items): - return [] - elif src == 'ldda': - all_items = self._get_all_lddas(trans=trans, ldda_id=v['id']) - if not self.one_not_deleted(all_items): - return [] - elif src == 'hdca': - all_items = self._get_all_hdcas(trans=trans, hdca_id=v['id']) - if self.one_not_deleted(all_items): - if any(True for hda in all_items[0].dataset_instances if hda.deleted): - return [] - else: - return [] - else: - # We don't know how to deal with inputs that are not HDCA/HDA/LDDA - # and possibly LDCA in the future. - return [] - input_data[k] = {src: {item.id for item in all_items}} - else: - input_param[k] = json.dumps(v) if not isinstance(v, int) else json.dumps(str(v)) - return self.__search(tool_id=tool_id, user=user, input_param=input_param, input_data=input_data, job_state=job_state) + input_data = defaultdict(list) + input_ids = defaultdict(dict) - def one_not_deleted(self, items): - return any(True for item in items if not item.deleted) + def populate_input_data_input_id(path, key, value): + """Traverses expanded incoming using remap and collects input_ids and input_data.""" + if key == 'id': + path_key = get_path_key(path[:-2]) + current_case = param_dump + for p in path: + current_case = current_case[p] + src = current_case['src'] + input_data[path_key].append({'src': src, 'id': value}) + input_ids[src][value] = True + return key, value + return key, value - def _get_all_hdas(self, trans, hda_id): - """Given a decoded `hda_id`, find other instances that refer to the same dataset.""" - hda = self.hda_manager.get_accessible(hda_id, trans.user) - return self.sa_session.query(model.HistoryDatasetAssociation).filter( - model.HistoryDatasetAssociation.dataset_id == hda.dataset_id).all() - - def _get_all_lddas(self, trans, ldda_id): - """Given a decoded `ldda_id` return the corresponding ldda""" - # TODO: implement looking up HDAs that point to the same dataset - ldda = self.ldda_manager.get(trans=trans, id=ldda_id) - return [ldda] - - def _get_all_hdcas(self, trans, hdca_id): - """Given an hdca, returns a list of other hdcas that were copied from this hdca.""" - # TODO: would be great if we can find all identical hdcas - hdca = self.dataset_collection_manager.get_dataset_collection_instance(trans=trans, - instance_type='history', - id=hdca_id) - copied_from = getattr(hdca, 'copied_from_history_dataset_collection_association', None) - hdcas = [hdca] - if copied_from: - hdcas.append(copied_from) - return hdcas - - def __search(self, tool_id, user, input_param, input_data, job_state=None): + remap(param_dump, visit=populate_input_data_input_id) + return self.__search(tool_id=tool_id, + user=user, + input_data=input_data, + job_state=job_state, + param_dump=param_dump, + input_ids=input_ids, + is_workflow_step=is_workflow_step) + def __search(self, tool_id, user, input_data, input_ids=None, job_state=None, param_dump=None, is_workflow_step=False): + search_timer = ExecutionTimer() query = self.sa_session.query(model.Job).filter( model.Job.tool_id == tool_id, model.Job.user == user @@ -109,44 +103,132 @@ class JobSearch(object): or_(*o) ) - for k, v in input_param.items(): - a = aliased(model.JobParameter) - if isinstance(v, string_types): - query = query.filter(and_( - model.Job.id == a.job_id, - a.name == k, - a.value == v - )) - for k, type_values in input_data.items(): - for t, v in type_values.items(): + for k, input_list in input_data.items(): + for type_values in input_list: + t = type_values['src'] + v = type_values['id'] if t == 'hda': a = aliased(model.JobToInputDatasetAssociation) + b = aliased(model.HistoryDatasetAssociation) + c = aliased(model.HistoryDatasetAssociation) query = query.filter(and_( model.Job.id == a.job_id, a.name == k, - a.dataset_id.in_(v) + a.dataset_id == b.id, + c.dataset_id == b.dataset_id, + c.id == v, + or_(b.deleted == false(), c.deleted == false()) )) elif t == 'ldda': a = aliased(model.JobToInputLibraryDatasetAssociation) query = query.filter(and_( model.Job.id == a.job_id, a.name == k, - a.ldda_id.in_(v) + a.ldda_id == v )) elif t == 'hdca': a = aliased(model.JobToInputDatasetCollectionAssociation) + b = aliased(model.HistoryDatasetCollectionAssociation) + c = aliased(model.HistoryDatasetCollectionAssociation) query = query.filter(and_( model.Job.id == a.job_id, a.name == k, - a.dataset_collection_id.in_(v) + b.id == a.dataset_collection_id, + c.id == v, + or_(and_(b.deleted == false(), b.id == v), + and_(or_(c.copied_from_history_dataset_collection_association_id == b.id, + b.copied_from_history_dataset_collection_association_id == c.id), + c.deleted == false() + ) + ) )) else: return [] - jobs = [] for job in query.all(): + # We found a job that is equal in terms of tool_id, user, state and input datasets, + # but to be able to verify that the parameters match we need to modify all instances of + # dataset_ids (HDA, LDDA, HDCA) in the incoming param_dump to point to those used by the + # possibly equivalent job, which may have been run on copies of the original input data. + replacement_timer = ExecutionTimer() + job_input_ids = {} + for src, items in input_ids.items(): + for dataset_id in items: + if src in job_input_ids and dataset_id in job_input_ids[src]: + continue + if src == 'hda': + a = aliased(model.JobToInputDatasetAssociation) + b = aliased(model.HistoryDatasetAssociation) + c = aliased(model.HistoryDatasetAssociation) + + (job_dataset_id,) = self.sa_session.query(b.id).filter( + and_( + a.job_id == job.id, + b.id == a.dataset_id, + c.dataset_id == b.dataset_id, + c.id == dataset_id + ) + ).first() + elif src == 'hdca': + a = aliased(model.JobToInputDatasetCollectionAssociation) + b = aliased(model.HistoryDatasetCollectionAssociation) + c = aliased(model.HistoryDatasetCollectionAssociation) + + (job_dataset_id,) = self.sa_session.query(b.id).filter( + and_( + a.job_id == job.id, + b.id == a.dataset_collection_id, + c.id == dataset_id, + or_(b.id == c.id, or_(c.copied_from_history_dataset_collection_association_id == b.id, + b.copied_from_history_dataset_collection_association_id == c.id) + ) + ) + ).first() + elif src == 'ldda': + job_dataset_id = dataset_id + else: + return [] + if src not in job_input_ids: + job_input_ids[src] = {dataset_id: job_dataset_id} + else: + job_input_ids[src][dataset_id] = job_dataset_id + + def replace_dataset_ids(path, key, value): + """Exchanges dataset_ids (HDA, LDA, HDCA, not Dataset) in param_dump with dataset ids used in job.""" + if key == 'id': + current_case = param_dump + for p in path: + current_case = current_case[p] + src = current_case['src'] + value = job_input_ids[src][value] + return key, value + return key, value + + new_param_dump = remap(param_dump, visit=replace_dataset_ids) + log.info("Parameter replacement finished %s", replacement_timer) + # new_param_dump has its dataset ids remapped to those used by the job. + # We now ask if the remapped job parameters match the current job. + query = self.sa_session.query(model.Job).filter(model.Job.id == job.id) + for k, v in new_param_dump.items(): + a = aliased(model.JobParameter) + query = query.filter(and_( + a.job_id == job.id, + a.name == k, + a.value == json.dumps(v) + )) + if query.first() is None: + continue + if is_workflow_step: + add_n_parameters = 3 + else: + add_n_parameters = 2 + if not len(job.parameters) == (len(new_param_dump) + add_n_parameters): + # Verify that equivalent jobs had the same number of job parameters + # We add 2 or 3 to new_param_dump because chrominfo and dbkey (and __workflow_invocation_uuid__) are not passed + # as input parameters + continue # check to make sure none of the output datasets or collections have been deleted - # TODO: find copies if output is deleted + # TODO: refactors this into the initial job query outputs_deleted = False for hda in job.output_datasets: if hda.dataset.deleted: @@ -158,5 +240,6 @@ class JobSearch(object): outputs_deleted = True break if not outputs_deleted: - jobs.append(job) - return jobs + log.info("Searching jobs finished %s", search_timer) + return job + return None diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py index 95fb7492b90..08e2b7eaab7 100755 --- a/lib/galaxy/tools/__init__.py +++ b/lib/galaxy/tools/__init__.py @@ -1235,14 +1235,7 @@ class Tool(object, Dictifiable): if self.check_values: visit_input_values(self.inputs, values, callback) - def handle_input(self, trans, incoming, history=None): - """ - Process incoming parameters for this tool from the dict `incoming`, - update the tool state (or create if none existed), and either return - to the form or execute the tool (only if 'execute' was clicked and - there were no errors). - """ - request_context = WorkRequestContext(app=trans.app, user=trans.user, history=history or trans.history) + def expand_incoming(self, trans, incoming, request_context): rerun_remap_job_id = None if 'rerun_remap_job_id' in incoming: try: @@ -1260,7 +1253,8 @@ class Tool(object, Dictifiable): # Remapping a single job to many jobs doesn't make sense, so disable # remap if multi-runs of tools are being used. if rerun_remap_job_id and len(expanded_incomings) > 1: - raise exceptions.MessageException('Failure executing tool (cannot create multiple jobs when remapping existing job).') + raise exceptions.MessageException( + 'Failure executing tool (cannot create multiple jobs when remapping existing job).') # Process incoming data validation_timer = ExecutionTimer() @@ -1288,6 +1282,17 @@ class Tool(object, Dictifiable): all_errors.append(errors) all_params.append(params) log.debug('Validated and populated state for tool request %s' % validation_timer) + return all_params, all_errors, rerun_remap_job_id, collection_info + + def handle_input(self, trans, incoming, history=None): + """ + Process incoming parameters for this tool from the dict `incoming`, + update the tool state (or create if none existed), and either return + to the form or execute the tool (only if 'execute' was clicked and + there were no errors). + """ + request_context = WorkRequestContext(app=trans.app, user=trans.user, history=history or trans.history) + all_params, all_errors, rerun_remap_job_id, collection_info = self.expand_incoming(trans=trans, incoming=incoming, request_context=request_context) # If there were errors, we stay on the same page and display them if any(all_errors): err_data = {key: value for d in all_errors for (key, value) in d.items()} @@ -1392,8 +1397,8 @@ class Tool(object, Dictifiable): """ return self.tool_action.execute(self, trans, incoming=incoming, set_output_hid=set_output_hid, history=history, **kwargs) - def params_to_strings(self, params, app): - return params_to_strings(self.inputs, params, app) + def params_to_strings(self, params, app, nested=False): + return params_to_strings(self.inputs, params, app, nested) def params_from_strings(self, params, app, ignore_errors=False): return params_from_strings(self.inputs, params, app, ignore_errors) diff --git a/lib/galaxy/util/__init__.py b/lib/galaxy/util/__init__.py index fe4da02c155..6bea695ea64 100644 --- a/lib/galaxy/util/__init__.py +++ b/lib/galaxy/util/__init__.py @@ -68,6 +68,9 @@ BINARY_CHARS = [NULL_CHAR] FILENAME_VALID_CHARS = '.,^_-()[]0123456789abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ' +defaultdict = collections.defaultdict + + def remove_protocol_from_url(url): """ Supplied URL may be null, if not ensure http:// or https:// etc... is stripped off. diff --git a/lib/galaxy/webapps/galaxy/api/jobs.py b/lib/galaxy/webapps/galaxy/api/jobs.py index 7cc56b125ff..df2b0360efd 100644 --- a/lib/galaxy/webapps/galaxy/api/jobs.py +++ b/lib/galaxy/webapps/galaxy/api/jobs.py @@ -17,6 +17,7 @@ from galaxy.web import _future_expose_api as expose_api from galaxy.web import _future_expose_api_anonymous as expose_api_anonymous from galaxy.web.base.controller import BaseAPIController from galaxy.web.base.controller import UsesLibraryMixinItems +from galaxy.work.context import WorkRequestContext log = logging.getLogger(__name__) @@ -271,9 +272,25 @@ class JobController(BaseAPIController, UsesLibraryMixinItems): raise exceptions.ObjectNotFound("Requested tool not found") if 'inputs' not in payload: raise exceptions.ObjectAttributeMissingException("No inputs defined") - inputs = payload['inputs'] - jobs = self.job_search.by_tool_input(trans=trans, tool_id=tool_id, inputs=inputs, job_state=payload.get('state')) - return [self.encode_all_ids(trans, job.to_dict('element'), True) for job in jobs] + inputs = payload.get('inputs', {}) + # Find files coming in as multipart file data and add to inputs. + for k, v in payload.iteritems(): + if k.startswith('files_') or k.startswith('__files_'): + inputs[k] = v + request_context = WorkRequestContext(app=trans.app, user=trans.user, history=trans.history) + all_params, all_errors, _, _ = tool.expand_incoming(trans=trans, incoming=inputs, request_context=request_context) + if any(all_errors): + return [] + params_dump = [tool.params_to_strings(param, self.app, nested=True) for param in all_params] + jobs = [] + for param_dump in params_dump: + job = self.job_search.by_tool_input(trans=trans, + tool_id=tool_id, + param_dump=param_dump, + job_state=payload.get('state')) + if job: + jobs.append(job) + return [self.encode_all_ids(trans, single_job.to_dict('element'), True) for single_job in jobs] @expose_api def error(self, trans, id, **kwd):