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):