diff --git a/lib/galaxy/managers/jobs.py b/lib/galaxy/managers/jobs.py index f7a11dd82d8..ed1e3fceae0 100644 --- a/lib/galaxy/managers/jobs.py +++ b/lib/galaxy/managers/jobs.py @@ -49,7 +49,7 @@ class JobSearch(object): self.ldda_manager = LDDAManager(app) self.decode_id = self.app.security.decode_id - def by_tool_input(self, trans, tool_id, tool_version, param_dump=None, job_state='ok'): + def by_tool_input(self, trans, tool_id, tool_version, param=None, param_dump=None, job_state='ok'): """Search for jobs producing same results using the 'inputs' part of a tool POST.""" user = trans.user input_data = defaultdict(list) @@ -63,7 +63,19 @@ class JobSearch(object): for p in path: current_case = current_case[p] src = current_case['src'] - input_data[path_key].append({'src': src, 'id': value}) + current_case = param + for i, p in enumerate(path): + if p == 'values' and i == len(path) - 2: + continue + if isinstance(current_case, (list, dict)): + current_case = current_case[p] + identifier = getattr(current_case, "element_identifier", None) + name = current_case.name + input_data[path_key].append({'src': src, + 'id': value, + 'identifier': identifier, + 'name': name, + }) input_ids[src][value] = "__id_wildcard__" return key, value return key, value @@ -111,18 +123,32 @@ class JobSearch(object): for type_values in input_list: t = type_values['src'] v = type_values['id'] + identifier = type_values['identifier'] if t == 'hda': a = aliased(model.JobToInputDatasetAssociation) b = aliased(model.HistoryDatasetAssociation) c = aliased(model.HistoryDatasetAssociation) + d = aliased(model.JobParameter) query = query.filter(and_( model.Job.id == a.job_id, a.name == k, a.dataset_id == b.id, c.dataset_id == b.dataset_id, - c.id == v, + c.id == v, # c is the requested job input HDA + # We can only compare input dataset names and metadata for a job + # if we know that the input dataset hasn't changed since the job was run. + # This is relatively strict, we may be able to lift this requirement if we record the jobs' + # relevant parameters as JobParameters in the database + b.update_time < model.Job.create_time, + b.name == c.name, + b.extension == c.extension, + b._metadata == c._metadata, or_(b.deleted == false(), c.deleted == false()) )) + if identifier: + query = query.filter(model.Job.id == d.job_id, + d.name == "%s|__identifier__" % k, + d.value == json.dumps(identifier)) elif t == 'ldda': a = aliased(model.JobToInputLibraryDatasetAssociation) query = query.filter(and_( @@ -134,15 +160,17 @@ class JobSearch(object): a = aliased(model.JobToInputDatasetCollectionAssociation) b = aliased(model.HistoryDatasetCollectionAssociation) c = aliased(model.HistoryDatasetCollectionAssociation) + name = type_values['name'] query = query.filter(and_( model.Job.id == a.job_id, a.name == k, b.id == a.dataset_collection_id, c.id == v, - or_(and_(b.deleted == false(), b.id == v), + or_(and_(b.deleted == false(), b.id == v, b.name == name), 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() + c.deleted == false(), + c.name == name, ) ) )) diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py index 09fe7d22562..e62410ccc94 100755 --- a/lib/galaxy/tools/__init__.py +++ b/lib/galaxy/tools/__init__.py @@ -1302,6 +1302,7 @@ class Tool(object, Dictifiable): trans=trans, tool_id=self.id, tool_version=self.version, + param=param, param_dump=self.params_to_strings(param, self.app, nested=True), ) execution_tracker = execute_job(trans, self, mapping_params, history=request_context.history, rerun_remap_job_id=rerun_remap_job_id, collection_info=collection_info, completed_jobs=completed_jobs) diff --git a/lib/galaxy/webapps/galaxy/api/jobs.py b/lib/galaxy/webapps/galaxy/api/jobs.py index 39cfaa9cc5a..9330359ef55 100644 --- a/lib/galaxy/webapps/galaxy/api/jobs.py +++ b/lib/galaxy/webapps/galaxy/api/jobs.py @@ -302,10 +302,11 @@ class JobController(BaseAPIController, UsesLibraryMixinItems): return [] params_dump = [tool.params_to_strings(param, self.app, nested=True) for param in all_params] jobs = [] - for param_dump in params_dump: + for param_dump, param in zip(params_dump, all_params): job = self.job_search.by_tool_input(trans=trans, tool_id=tool_id, tool_version=tool.version, + param=param, param_dump=param_dump, job_state=payload.get('state')) if job: diff --git a/lib/galaxy/workflow/modules.py b/lib/galaxy/workflow/modules.py index a408da498d5..2e64319ed4b 100644 --- a/lib/galaxy/workflow/modules.py +++ b/lib/galaxy/workflow/modules.py @@ -879,6 +879,7 @@ class ToolModule(WorkflowModule): trans=trans, tool_id=tool.id, tool_version=tool.version, + param=param, param_dump=tool.params_to_strings(param, trans.app, nested=True), ) try: diff --git a/test/api/test_jobs.py b/test/api/test_jobs.py index 1220a5c198b..805a91f34d1 100644 --- a/test/api/test_jobs.py +++ b/test/api/test_jobs.py @@ -298,6 +298,21 @@ class JobsApiTestCase(api.ApiTestCase): self._assert_status_code_is(delete_respone, 200) self._search(search_payload, expected_search_count=0) + def test_search_handle_identifiers(self): + # Test that input name and element identifier of a jobs' output must match for a job to be returned. + history_id, dataset_id = self.__history_with_ok_dataset() + inputs = json.dumps({ + 'input1': {'src': 'hda', 'id': dataset_id} + }) + self._job_search(tool_id='identifier_single', history_id=history_id, inputs=inputs) + dataset_details = self._get("histories/%s/contents/%s" % (history_id, dataset_id)).json() + dataset_details['name'] = 'Renamed Test Dataset' + dataset_update_response = self._put("histories/%s/contents/%s" % (history_id, dataset_id), data=dict(name='Renamed Test Dataset')) + self._assert_status_code_is(dataset_update_response, 200) + assert dataset_update_response.json()['name'] == 'Renamed Test Dataset' + search_payload = self._search_payload(history_id=history_id, tool_id='identifier_single', inputs=inputs) + self._search(search_payload, expected_search_count=0) + def test_search_delete_outputs(self): history_id, dataset_id = self.__history_with_ok_dataset() inputs = json.dumps({ diff --git a/test/base/api.py b/test/base/api.py index abff78b4bb6..4cc36ac09d0 100644 --- a/test/base/api.py +++ b/test/base/api.py @@ -76,6 +76,9 @@ class UsesApiTestCaseMixin: def _delete(self, *args, **kwds): return self.galaxy_interactor.delete(*args, **kwds) + def _put(self, *args, **kwds): + return self.galaxy_interactor.put(*args, **kwds) + def _patch(self, *args, **kwds): return self.galaxy_interactor.patch(*args, **kwds) @@ -129,3 +132,6 @@ class ApiTestInteractor(BaseInteractor): def patch(self, *args, **kwds): return self._patch(*args, **kwds) + + def put(self, *args, **kwds): + return self._put(*args, **kwds) diff --git a/test/base/interactor.py b/test/base/interactor.py index 736527477f5..9ea5ed0318c 100644 --- a/test/base/interactor.py +++ b/test/base/interactor.py @@ -6,7 +6,13 @@ import time from json import dumps from logging import getLogger -from requests import delete, get, patch, post +from requests import ( + delete, + get, + patch, + post, + put, +) from six import StringIO, text_type from galaxy import util @@ -479,6 +485,17 @@ class GalaxyInteractorApi(object): params = {} return patch("%s/%s" % (self.api_url, path), params=params, data=data) + def _put(self, path, data={}, key=None, admin=False, anon=False): + if not anon: + if not key: + key = self.api_key if not admin else self.master_api_key + params = dict(key=key) + data = data.copy() + data['key'] = key + else: + params = {} + return put("%s/%s" % (self.api_url, path), params=params, data=data) + def _get(self, path, data={}, key=None, admin=False, anon=False): if not anon: if not key: