From 92b48412c5cbb50c31b26e97912c073b7c0f29e2 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Wed, 20 Sep 2017 19:17:59 +0200 Subject: [PATCH 01/17] Extend job search for HDCA input --- lib/galaxy/webapps/galaxy/api/jobs.py | 41 ++++++++++++--------- test/api/test_jobs.py | 52 +++++++++++++++++++++++++-- 2 files changed, 74 insertions(+), 19 deletions(-) diff --git a/lib/galaxy/webapps/galaxy/api/jobs.py b/lib/galaxy/webapps/galaxy/api/jobs.py index 63367938600..49d954cd3f3 100644 --- a/lib/galaxy/webapps/galaxy/api/jobs.py +++ b/lib/galaxy/webapps/galaxy/api/jobs.py @@ -28,6 +28,7 @@ class JobController(BaseAPIController, UsesLibraryMixinItems): def __init__(self, app): super(JobController, self).__init__(app) self.hda_manager = managers.hdas.HDAManager(app) + self.dataset_collection_manager = managers.collections.DatasetCollectionManager(app) self.dataset_manager = managers.datasets.DatasetManager(app) @expose_api @@ -267,9 +268,7 @@ class JobController(BaseAPIController, UsesLibraryMixinItems): recycle the old results. """ - tool_id = None - if 'tool_id' in payload: - tool_id = payload.get('tool_id') + tool_id = payload.get('tool_id') if tool_id is None: raise exceptions.ObjectAttributeMissingException("No tool id") @@ -288,12 +287,19 @@ class JobController(BaseAPIController, UsesLibraryMixinItems): if 'id' in v: if 'src' not in v or v['src'] == 'hda': hda_id = self.decode_id(v['id']) - dataset = self.hda_manager.get_accessible(hda_id, trans.user) + datasets = self.hda_manager.get_accessible(hda_id, trans.user) else: - dataset = self.get_library_dataset_dataset_association(trans, v['id']) - if dataset is None: + if v['src'] == 'hdca': + hdca = self.dataset_collection_manager.get_dataset_collection_instance(trans=trans, + instance_type='history', + id=v['id'], + ) + datasets = hdca.dataset_instances + else: + dataset = self.get_library_dataset_dataset_association(trans, v['id']) + if datasets is None: raise exceptions.ObjectNotFound("Dataset %s not found" % (v['id'])) - input_data[k] = dataset.dataset_id + input_data[k] = [datasets.dataset_id] if isinstance(datasets, model.HistoryDatasetAssociation) else [d.dataset_id for d in datasets] else: input_param[k] = json.dumps(str(v)) @@ -331,20 +337,23 @@ class JobController(BaseAPIController, UsesLibraryMixinItems): a.value == v )) - for k, v in input_data.items(): + for datasets in input_data.values(): # Here we are attempting to link the inputs to the underlying # dataset (not the dataset association). # This way, if the calculation was done using a copied HDA # (copied from the library or another history), the search will # still find the job - a = aliased(trans.app.model.JobToInputDatasetAssociation) - b = aliased(trans.app.model.HistoryDatasetAssociation) - query = query.filter(and_( - trans.app.model.Job.id == a.job_id, - a.dataset_id == b.id, - b.deleted == false(), - b.dataset_id == v - )) + for dataset in datasets: + a = aliased(trans.app.model.JobToInputDatasetAssociation) + b = aliased(trans.app.model.HistoryDatasetAssociation) + # TODO: I don't think this checks for the actual parameter name that the input matches to + # we may get a wrong result here + query = query.filter(and_( + trans.app.model.Job.id == a.job_id, + a.dataset_id == b.id, + b.deleted == false(), + b.dataset_id == dataset + )) out = [] for job in query.all(): diff --git a/test/api/test_jobs.py b/test/api/test_jobs.py index 6cb1565abc3..d8fd9df6779 100644 --- a/test/api/test_jobs.py +++ b/test/api/test_jobs.py @@ -8,6 +8,7 @@ from operator import itemgetter from base import api from base.api_asserts import assert_status_code_is_ok from base.populators import ( + DatasetCollectionPopulator, DatasetPopulator, wait_on, wait_on_state, @@ -21,6 +22,7 @@ class JobsApiTestCase(api.ApiTestCase): def setUp(self): super(JobsApiTestCase, self).setUp() self.dataset_populator = DatasetPopulator(self.galaxy_interactor) + self.dataset_collection_populator = DatasetCollectionPopulator(self.galaxy_interactor) def test_index(self): # Create HDA to ensure at least one job exists... @@ -295,16 +297,53 @@ class JobsApiTestCase(api.ApiTestCase): self.__run_cat_tool(history_id, dataset_id) self.dataset_populator.wait_for_history(history_id, assert_ok=True) + search_count = self._search(search_payload) + + self.assertEquals(search_count, 1) + + def test_search_with_hdca_input(self): + history_id, collection_id = self.__history_with_ok_collection() + + inputs = json.dumps( + dict( + f1=dict( + src='hdca', + id=collection_id + ), + f2=dict( + src='hdca', + id=collection_id + ), + ) + ) + + search_payload = dict( + tool_id = "multi_data_param", + inputs=inputs, + history_id=history_id, + state='ok' + ) + + self._post("tools", data=search_payload) + self.dataset_populator.wait_for_history(history_id, assert_ok=True) + search_count = self._search(search_payload) + + self.assertEquals(search_count, 1) + + # inputs = {"tool_id":"multi_data_param","tool_version":"0.1.0","inputs":{"f1":{"values":[{"src":"hdca","name":"collection","tags":[],"keep":False,"hid":32,"id":"f2db41e1fa331b3e"}]},"f2":{"values":[{"src":"hdca","name":"collection","tags":[],"keep":false,"hid":32,"id":"f2db41e1fa331b3e"}]}}} + + pass + + def _search(self, payload): search_count = -1 # in case job and history aren't updated at exactly the same # time give time to wait for i in range(5): - search_count = self._search_count(search_payload) + search_count = self._search_count(payload) if search_count == 1: break time.sleep(.1) - - self.assertEquals(search_count, 1) + return search_count def _search_count(self, search_payload): search_response = self._post("jobs/search", data=search_payload) @@ -357,6 +396,13 @@ class JobsApiTestCase(api.ApiTestCase): dataset_id = self.dataset_populator.new_dataset(history_id, wait=True)["id"] return history_id, dataset_id + def __history_with_ok_collection(self): + contents = ["a\tb\nc\td", "e\tf\ng\th"] + history_id = self.dataset_populator.new_history() + create_reposonse = self.dataset_collection_populator.create_list_in_history(history_id, contents=contents).json() + self.dataset_collection_populator.wait_for_dataset_collection(create_reposonse) + return history_id, create_reposonse['id'] + def __jobs_index(self, **kwds): jobs_response = self._get("jobs", **kwds) self._assert_status_code_is(jobs_response, 200) From 1d37c7d458f156e391d969bf7630044a476b9856 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Thu, 21 Sep 2017 10:34:00 +0200 Subject: [PATCH 02/17] Check JobParameter match Matching the JobParameters will prevent returning jobs where a job uses the same input datasets, but where the inputs are used in a different order (and hence the result may differ) --- lib/galaxy/webapps/galaxy/api/jobs.py | 17 ++++--- test/api/test_jobs.py | 64 +++++++++++++++------------ 2 files changed, 45 insertions(+), 36 deletions(-) diff --git a/lib/galaxy/webapps/galaxy/api/jobs.py b/lib/galaxy/webapps/galaxy/api/jobs.py index 49d954cd3f3..ea9d18c6bfc 100644 --- a/lib/galaxy/webapps/galaxy/api/jobs.py +++ b/lib/galaxy/webapps/galaxy/api/jobs.py @@ -281,13 +281,14 @@ class JobController(BaseAPIController, UsesLibraryMixinItems): inputs = payload['inputs'] input_data = {} + decoded_input_data = {} input_param = {} for k, v in inputs.items(): if isinstance(v, dict): if 'id' in v: + decoded_id = self.decode_id(v['id']) if 'src' not in v or v['src'] == 'hda': - hda_id = self.decode_id(v['id']) - datasets = self.hda_manager.get_accessible(hda_id, trans.user) + datasets = self.hda_manager.get_accessible(decoded_id, trans.user) else: if v['src'] == 'hdca': hdca = self.dataset_collection_manager.get_dataset_collection_instance(trans=trans, @@ -296,10 +297,14 @@ class JobController(BaseAPIController, UsesLibraryMixinItems): ) datasets = hdca.dataset_instances else: - dataset = self.get_library_dataset_dataset_association(trans, v['id']) + datasets = self.get_library_dataset_dataset_association(trans, v['id']) if datasets is None: raise exceptions.ObjectNotFound("Dataset %s not found" % (v['id'])) input_data[k] = [datasets.dataset_id] if isinstance(datasets, model.HistoryDatasetAssociation) else [d.dataset_id for d in datasets] + decoded_input_data[k] = v.copy() + decoded_input_data[k]['id'] = decoded_id + decoded_input_data[k] = json.dumps({'values': [decoded_input_data[k]]}) + input_param.update(decoded_input_data) else: input_param[k] = json.dumps(str(v)) @@ -337,7 +342,7 @@ class JobController(BaseAPIController, UsesLibraryMixinItems): a.value == v )) - for datasets in input_data.values(): + for k, datasets in input_data.items(): # Here we are attempting to link the inputs to the underlying # dataset (not the dataset association). # This way, if the calculation was done using a copied HDA @@ -346,13 +351,11 @@ class JobController(BaseAPIController, UsesLibraryMixinItems): for dataset in datasets: a = aliased(trans.app.model.JobToInputDatasetAssociation) b = aliased(trans.app.model.HistoryDatasetAssociation) - # TODO: I don't think this checks for the actual parameter name that the input matches to - # we may get a wrong result here query = query.filter(and_( trans.app.model.Job.id == a.job_id, a.dataset_id == b.id, b.deleted == false(), - b.dataset_id == dataset + b.dataset_id == dataset, )) out = [] diff --git a/test/api/test_jobs.py b/test/api/test_jobs.py index d8fd9df6779..703d7f06d41 100644 --- a/test/api/test_jobs.py +++ b/test/api/test_jobs.py @@ -302,37 +302,38 @@ class JobsApiTestCase(api.ApiTestCase): self.assertEquals(search_count, 1) def test_search_with_hdca_input(self): - history_id, collection_id = self.__history_with_ok_collection() - - inputs = json.dumps( - dict( - f1=dict( - src='hdca', - id=collection_id - ), - f2=dict( - src='hdca', - id=collection_id - ), - ) - ) - - search_payload = dict( - tool_id = "multi_data_param", - inputs=inputs, - history_id=history_id, - state='ok' - ) + history_id, list_id_a = self.__history_with_ok_collection(collection_type='list') + history_id, list_id_b = self.__history_with_ok_collection(collection_type='list', history_id=history_id) + inputs = json.dumps({ + 'f1': {'src': 'hdca', 'id': list_id_a}, + 'f2': {'src': 'hdca', 'id': list_id_b}, + }) + search_payload = self._search_payload(history_id=history_id, tool_id='multi_data_param', inputs=inputs) + empty_search_response = self._post("jobs/search", data=search_payload) + self._assert_status_code_is(empty_search_response, 200) + self.assertEquals(len(empty_search_response.json()), 0) self._post("tools", data=search_payload) self.dataset_populator.wait_for_history(history_id, assert_ok=True) search_count = self._search(search_payload) - self.assertEquals(search_count, 1) + # We switch the inputs, this should not return a match + inputs = json.dumps({ + 'f2': {'src': 'hdca', 'id': list_id_a}, + 'f1': {'src': 'hdca', 'id': list_id_b}, + }) + search_payload = self._search_payload(history_id=history_id, tool_id='multi_data_param', inputs=inputs) + search_count = self._search(search_payload) + self.assertEquals(search_count, 0) - # inputs = {"tool_id":"multi_data_param","tool_version":"0.1.0","inputs":{"f1":{"values":[{"src":"hdca","name":"collection","tags":[],"keep":False,"hid":32,"id":"f2db41e1fa331b3e"}]},"f2":{"values":[{"src":"hdca","name":"collection","tags":[],"keep":false,"hid":32,"id":"f2db41e1fa331b3e"}]}}} - - pass + def _search_payload(self, history_id, tool_id, inputs, state='ok'): + search_payload = dict( + tool_id=tool_id, + inputs=inputs, + history_id=history_id, + state=state + ) + return search_payload def _search(self, payload): search_count = -1 @@ -396,10 +397,15 @@ class JobsApiTestCase(api.ApiTestCase): dataset_id = self.dataset_populator.new_dataset(history_id, wait=True)["id"] return history_id, dataset_id - def __history_with_ok_collection(self): - contents = ["a\tb\nc\td", "e\tf\ng\th"] - history_id = self.dataset_populator.new_history() - create_reposonse = self.dataset_collection_populator.create_list_in_history(history_id, contents=contents).json() + def __history_with_ok_collection(self, collection_type='list', history_id=None): + if not history_id: + history_id = self.dataset_populator.new_history() + if collection_type == 'list': + create_reposonse = self.dataset_collection_populator.create_list_in_history(history_id).json() + elif collection_type == 'pair': + create_reposonse = self.dataset_collection_populator.create_pair_in_history(history_id).json() + elif collection_type == 'list:pair': + create_reposonse = self.dataset_collection_populator.create_list_of_pairs_in_history(history_id).json() self.dataset_collection_populator.wait_for_dataset_collection(create_reposonse) return history_id, create_reposonse['id'] From 2b7a1631280e19468021931ca827b94567bca457 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Thu, 21 Sep 2017 11:03:18 +0200 Subject: [PATCH 03/17] Add job_search tests for other collection types --- test/api/test_jobs.py | 38 ++++++++++++++++++++++++++++---------- 1 file changed, 28 insertions(+), 10 deletions(-) diff --git a/test/api/test_jobs.py b/test/api/test_jobs.py index 703d7f06d41..f8eac212a45 100644 --- a/test/api/test_jobs.py +++ b/test/api/test_jobs.py @@ -301,22 +301,14 @@ class JobsApiTestCase(api.ApiTestCase): self.assertEquals(search_count, 1) - def test_search_with_hdca_input(self): + def test_search_with_hdca_list_input(self): history_id, list_id_a = self.__history_with_ok_collection(collection_type='list') history_id, list_id_b = self.__history_with_ok_collection(collection_type='list', history_id=history_id) - inputs = json.dumps({ 'f1': {'src': 'hdca', 'id': list_id_a}, 'f2': {'src': 'hdca', 'id': list_id_b}, }) - search_payload = self._search_payload(history_id=history_id, tool_id='multi_data_param', inputs=inputs) - empty_search_response = self._post("jobs/search", data=search_payload) - self._assert_status_code_is(empty_search_response, 200) - self.assertEquals(len(empty_search_response.json()), 0) - self._post("tools", data=search_payload) - self.dataset_populator.wait_for_history(history_id, assert_ok=True) - search_count = self._search(search_payload) - self.assertEquals(search_count, 1) + self._collection_job_search(tool_id='multi_data_param', history_id=history_id, inputs=inputs) # We switch the inputs, this should not return a match inputs = json.dumps({ 'f2': {'src': 'hdca', 'id': list_id_a}, @@ -326,6 +318,32 @@ class JobsApiTestCase(api.ApiTestCase): search_count = self._search(search_payload) self.assertEquals(search_count, 0) + def test_search_with_hdca_pair_input(self): + history_id, list_id_a = self.__history_with_ok_collection(collection_type='pair') + inputs = json.dumps({ + 'f1': {'src': 'hdca', 'id': list_id_a}, + 'f2': {'src': 'hdca', 'id': list_id_a}, + }) + self._collection_job_search(tool_id='multi_data_param', history_id=history_id, inputs=inputs) + + def test_search_with_hdca_list_pair_input(self): + history_id, list_id_a = self.__history_with_ok_collection(collection_type='list:pair') + inputs = json.dumps({ + 'f1': {'src': 'hdca', 'id': list_id_a}, + 'f2': {'src': 'hdca', 'id': list_id_a}, + }) + self._collection_job_search(tool_id='multi_data_param', history_id=history_id, inputs=inputs) + + def _collection_job_search(self, tool_id, history_id, inputs): + search_payload = self._search_payload(history_id=history_id, tool_id=tool_id, inputs=inputs) + empty_search_response = self._post("jobs/search", data=search_payload) + self._assert_status_code_is(empty_search_response, 200) + self.assertEquals(len(empty_search_response.json()), 0) + self._post("tools", data=search_payload) + self.dataset_populator.wait_for_history(history_id, assert_ok=True) + search_count = self._search(search_payload) + self.assertEquals(search_count, 1) + def _search_payload(self, history_id, tool_id, inputs, state='ok'): search_payload = dict( tool_id=tool_id, From 7089de8142199e1b5dadaeef215ae015042f018b Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Thu, 21 Sep 2017 11:19:11 +0200 Subject: [PATCH 04/17] Reduce job_search test copypaste --- test/api/test_jobs.py | 51 +++++++------------------------------------ 1 file changed, 8 insertions(+), 43 deletions(-) diff --git a/test/api/test_jobs.py b/test/api/test_jobs.py index f8eac212a45..1c1383de4a2 100644 --- a/test/api/test_jobs.py +++ b/test/api/test_jobs.py @@ -275,31 +275,10 @@ class JobsApiTestCase(api.ApiTestCase): def test_search(self): history_id, dataset_id = self.__history_with_ok_dataset() - - inputs = json.dumps( - dict( - input1=dict( - src='hda', - id=dataset_id, - ) - ) - ) - search_payload = dict( - tool_id="cat1", - inputs=inputs, - state="ok", - ) - - empty_search_response = self._post("jobs/search", data=search_payload) - self._assert_status_code_is(empty_search_response, 200) - self.assertEquals(len(empty_search_response.json()), 0) - - self.__run_cat_tool(history_id, dataset_id) - self.dataset_populator.wait_for_history(history_id, assert_ok=True) - - search_count = self._search(search_payload) - - self.assertEquals(search_count, 1) + inputs = json.dumps({ + 'input1': {'src': 'hda', 'id': dataset_id} + }) + self._job_search(tool_id='cat1', history_id=history_id, inputs=inputs) def test_search_with_hdca_list_input(self): history_id, list_id_a = self.__history_with_ok_collection(collection_type='list') @@ -308,7 +287,7 @@ class JobsApiTestCase(api.ApiTestCase): 'f1': {'src': 'hdca', 'id': list_id_a}, 'f2': {'src': 'hdca', 'id': list_id_b}, }) - self._collection_job_search(tool_id='multi_data_param', history_id=history_id, inputs=inputs) + self._job_search(tool_id='multi_data_param', history_id=history_id, inputs=inputs) # We switch the inputs, this should not return a match inputs = json.dumps({ 'f2': {'src': 'hdca', 'id': list_id_a}, @@ -324,7 +303,7 @@ class JobsApiTestCase(api.ApiTestCase): 'f1': {'src': 'hdca', 'id': list_id_a}, 'f2': {'src': 'hdca', 'id': list_id_a}, }) - self._collection_job_search(tool_id='multi_data_param', history_id=history_id, inputs=inputs) + self._job_search(tool_id='multi_data_param', history_id=history_id, inputs=inputs) def test_search_with_hdca_list_pair_input(self): history_id, list_id_a = self.__history_with_ok_collection(collection_type='list:pair') @@ -332,9 +311,9 @@ class JobsApiTestCase(api.ApiTestCase): 'f1': {'src': 'hdca', 'id': list_id_a}, 'f2': {'src': 'hdca', 'id': list_id_a}, }) - self._collection_job_search(tool_id='multi_data_param', history_id=history_id, inputs=inputs) + self._job_search(tool_id='multi_data_param', history_id=history_id, inputs=inputs) - def _collection_job_search(self, tool_id, history_id, inputs): + def _job_search(self, tool_id, history_id, inputs): search_payload = self._search_payload(history_id=history_id, tool_id=tool_id, inputs=inputs) empty_search_response = self._post("jobs/search", data=search_payload) self._assert_status_code_is(empty_search_response, 200) @@ -370,20 +349,6 @@ class JobsApiTestCase(api.ApiTestCase): search_json = search_response.json() return len(search_json) - def __run_cat_tool(self, history_id, dataset_id): - # Code duplication with test_jobs.py, eliminate - payload = self.dataset_populator.run_tool_payload( - tool_id='cat1', - inputs=dict( - input1=dict( - src='hda', - id=dataset_id - ), - ), - history_id=history_id, - ) - self._post("tools", data=payload) - def __run_randomlines_tool(self, lines, history_id, dataset_id): payload = self.dataset_populator.run_tool_payload( tool_id="random_lines1", From 04754c4c4d8c4c287a54e13bc0263f1fdaeaf00c Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Thu, 21 Sep 2017 16:11:06 +0200 Subject: [PATCH 05/17] Allow job_search to find job even if current HDCA was copied from another HDCA and make sure that jobs with same inputs matched to different job parameter inputs are not considered as a match --- lib/galaxy/webapps/galaxy/api/jobs.py | 59 +++++++++++++++++++-------- test/api/test_jobs.py | 25 ++++++++++++ 2 files changed, 68 insertions(+), 16 deletions(-) diff --git a/lib/galaxy/webapps/galaxy/api/jobs.py b/lib/galaxy/webapps/galaxy/api/jobs.py index ea9d18c6bfc..bf1fe412254 100644 --- a/lib/galaxy/webapps/galaxy/api/jobs.py +++ b/lib/galaxy/webapps/galaxy/api/jobs.py @@ -283,28 +283,43 @@ class JobController(BaseAPIController, UsesLibraryMixinItems): input_data = {} decoded_input_data = {} input_param = {} + decoded_ids = [] for k, v in inputs.items(): if isinstance(v, dict): if 'id' in v: decoded_id = self.decode_id(v['id']) if 'src' not in v or v['src'] == 'hda': - datasets = self.hda_manager.get_accessible(decoded_id, trans.user) + datasets = [self.hda_manager.get_accessible(decoded_id, trans.user)] else: if v['src'] == 'hdca': hdca = self.dataset_collection_manager.get_dataset_collection_instance(trans=trans, instance_type='history', id=v['id'], ) + copied_from = getattr(hdca, 'copied_from_history_dataset_collection_association_id', None) + if copied_from: + # TODO: fetch these recursively or come up with a DB query to fetch all possible hdca ids + decoded_ids = [decoded_id, copied_from] + else: + decoded_ids = [decoded_id] datasets = hdca.dataset_instances else: - datasets = self.get_library_dataset_dataset_association(trans, v['id']) + datasets = [self.get_library_dataset_dataset_association(trans, v['id'])] if datasets is None: raise exceptions.ObjectNotFound("Dataset %s not found" % (v['id'])) - input_data[k] = [datasets.dataset_id] if isinstance(datasets, model.HistoryDatasetAssociation) else [d.dataset_id for d in datasets] - decoded_input_data[k] = v.copy() - decoded_input_data[k]['id'] = decoded_id - decoded_input_data[k] = json.dumps({'values': [decoded_input_data[k]]}) - input_param.update(decoded_input_data) + if v['src'] != 'hdca': + all_history_associations = [] + for dataset in datasets: + instances = trans.sa_session.query(trans.app.model.HistoryDatasetAssociation).filter(trans.app.model.HistoryDatasetAssociation.dataset_id == dataset.dataset_id).all() + all_history_associations.extend(instances) + decoded_ids = [h.id for h in all_history_associations] + parameter_values = [] + for decoded_id in decoded_ids: + parameter_value = v.copy() + parameter_value['id'] = decoded_id + parameter_values.append(json.dumps({'values': [parameter_value]})) + decoded_input_data[k] = parameter_values + input_data[k] = [d.dataset_id for d in datasets] else: input_param[k] = json.dumps(str(v)) @@ -342,21 +357,33 @@ class JobController(BaseAPIController, UsesLibraryMixinItems): a.value == v )) + for k, v in decoded_input_data.items(): + # Here we make sure that the incoming hdca/hda id + # (v is a list of possible text representations of a input data parameter) + # matches the right input data parameter for a job to be equivalent. + # If we don't do this we return false positive jobs in case + # a different combination of identical inputs has been used in a job. + a = aliased(trans.app.model.JobParameter) + query = query.filter(and_( + trans.app.model.Job.id == a.job_id, + a.name == k, + a.value.in_(v) + )) + for k, datasets in input_data.items(): # Here we are attempting to link the inputs to the underlying # dataset (not the dataset association). # This way, if the calculation was done using a copied HDA # (copied from the library or another history), the search will # still find the job - for dataset in datasets: - a = aliased(trans.app.model.JobToInputDatasetAssociation) - b = aliased(trans.app.model.HistoryDatasetAssociation) - query = query.filter(and_( - trans.app.model.Job.id == a.job_id, - a.dataset_id == b.id, - b.deleted == false(), - b.dataset_id == dataset, - )) + a = aliased(trans.app.model.JobToInputDatasetAssociation) + b = aliased(trans.app.model.HistoryDatasetAssociation) + query = query.filter(and_( + trans.app.model.Job.id == a.job_id, + a.dataset_id == b.id, + b.deleted == false(), + b.dataset_id.in_(datasets), + )) out = [] for job in query.all(): diff --git a/test/api/test_jobs.py b/test/api/test_jobs.py index 1c1383de4a2..068c1d40ef4 100644 --- a/test/api/test_jobs.py +++ b/test/api/test_jobs.py @@ -279,6 +279,19 @@ class JobsApiTestCase(api.ApiTestCase): 'input1': {'src': 'hda', 'id': dataset_id} }) self._job_search(tool_id='cat1', history_id=history_id, inputs=inputs) + # We test that a job can be found even if the dataset has been copied to another history + new_history_id = self.dataset_populator.new_history() + copy_payload = {"content": dataset_id, "source": "hda", "type": "dataset"} + copy_response = self._post("histories/%s/contents" % new_history_id, data=copy_payload) + self._assert_status_code_is(copy_response, 200) + new_dataset_id = copy_response.json()['id'] + copied_inputs = json.dumps({ + 'input1': {'src': 'hda', 'id': new_dataset_id} + }) + search_payload = self._search_payload(history_id=history_id, tool_id='cat1', inputs=copied_inputs) + search_count = self._search(search_payload) + self.assertEquals(search_count, 1) + def test_search_with_hdca_list_input(self): history_id, list_id_a = self.__history_with_ok_collection(collection_type='list') @@ -304,6 +317,18 @@ class JobsApiTestCase(api.ApiTestCase): 'f2': {'src': 'hdca', 'id': list_id_a}, }) self._job_search(tool_id='multi_data_param', history_id=history_id, inputs=inputs) + new_history_id = self.dataset_populator.new_history() + copy_payload = {"content": list_id_a, "source": "hdca", "type": "dataset_collection"} + copy_response = self._post("histories/%s/contents" % new_history_id, data=copy_payload) + self._assert_status_code_is(copy_response, 200) + new_list_a = copy_response.json()['id'] + copied_inputs = json.dumps({ + 'f1': {'src': 'hdca', 'id': new_list_a}, + 'f2': {'src': 'hdca', 'id': new_list_a}, + }) + search_payload = self._search_payload(history_id=new_history_id, tool_id='multi_data_param', inputs=copied_inputs) + search_count = self._search(search_payload) + self.assertEquals(search_count, 1) def test_search_with_hdca_list_pair_input(self): history_id, list_id_a = self.__history_with_ok_collection(collection_type='list:pair') From 64cf60b502daea924fbb7b1106c0040976cdbdc1 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Thu, 21 Sep 2017 19:15:29 +0200 Subject: [PATCH 06/17] Test deleting input datasets and recovering jobs from search This revealed a wrong query filter (comparing JobToInputDatasetAssociation.dataset_id == HistoryDatasetAssociation.id, instead of JobToInputDatasetAssociation.dataset_id == HistoryDatasetAssociation.dataset_id). --- lib/galaxy/webapps/galaxy/api/jobs.py | 13 ++++++------- test/api/test_jobs.py | 9 ++++++++- 2 files changed, 14 insertions(+), 8 deletions(-) diff --git a/lib/galaxy/webapps/galaxy/api/jobs.py b/lib/galaxy/webapps/galaxy/api/jobs.py index bf1fe412254..43ea71f6753 100644 --- a/lib/galaxy/webapps/galaxy/api/jobs.py +++ b/lib/galaxy/webapps/galaxy/api/jobs.py @@ -290,6 +290,11 @@ class JobController(BaseAPIController, UsesLibraryMixinItems): decoded_id = self.decode_id(v['id']) if 'src' not in v or v['src'] == 'hda': datasets = [self.hda_manager.get_accessible(decoded_id, trans.user)] + all_history_associations = [] + for dataset in datasets: + instances = trans.sa_session.query(trans.app.model.HistoryDatasetAssociation).filter(trans.app.model.HistoryDatasetAssociation.dataset_id == dataset.dataset_id).all() + all_history_associations.extend(instances) + decoded_ids = set(h.id for h in all_history_associations) else: if v['src'] == 'hdca': hdca = self.dataset_collection_manager.get_dataset_collection_instance(trans=trans, @@ -307,12 +312,6 @@ class JobController(BaseAPIController, UsesLibraryMixinItems): datasets = [self.get_library_dataset_dataset_association(trans, v['id'])] if datasets is None: raise exceptions.ObjectNotFound("Dataset %s not found" % (v['id'])) - if v['src'] != 'hdca': - all_history_associations = [] - for dataset in datasets: - instances = trans.sa_session.query(trans.app.model.HistoryDatasetAssociation).filter(trans.app.model.HistoryDatasetAssociation.dataset_id == dataset.dataset_id).all() - all_history_associations.extend(instances) - decoded_ids = [h.id for h in all_history_associations] parameter_values = [] for decoded_id in decoded_ids: parameter_value = v.copy() @@ -380,7 +379,7 @@ class JobController(BaseAPIController, UsesLibraryMixinItems): b = aliased(trans.app.model.HistoryDatasetAssociation) query = query.filter(and_( trans.app.model.Job.id == a.job_id, - a.dataset_id == b.id, + a.dataset_id == b.dataset_id, b.deleted == false(), b.dataset_id.in_(datasets), )) diff --git a/test/api/test_jobs.py b/test/api/test_jobs.py index 068c1d40ef4..3211c351e93 100644 --- a/test/api/test_jobs.py +++ b/test/api/test_jobs.py @@ -291,7 +291,14 @@ class JobsApiTestCase(api.ApiTestCase): search_payload = self._search_payload(history_id=history_id, tool_id='cat1', inputs=copied_inputs) search_count = self._search(search_payload) self.assertEquals(search_count, 1) - + delete_respone = self._delete("histories/%s/contents/%s" % (history_id, dataset_id)) + self._assert_status_code_is(delete_respone, 200) + search_count = self._search(search_payload) + self.assertEquals(search_count, 1) + delete_respone = self._delete("histories/%s/contents/%s" % (history_id, new_dataset_id)) + self._assert_status_code_is(delete_respone, 200) + search_count = self._search(search_payload) + self.assertEquals(search_count, 0) def test_search_with_hdca_list_input(self): history_id, list_id_a = self.__history_with_ok_collection(collection_type='list') From df91074a3f587e6bd3270e67c4d00ffc321518af Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Fri, 22 Sep 2017 08:32:17 +0200 Subject: [PATCH 07/17] Test deleting HDCAs and finding equivalent jobs --- test/api/test_jobs.py | 13 +++++++++++++ 1 file changed, 13 insertions(+) diff --git a/test/api/test_jobs.py b/test/api/test_jobs.py index 3211c351e93..24c20ca599f 100644 --- a/test/api/test_jobs.py +++ b/test/api/test_jobs.py @@ -291,10 +291,12 @@ class JobsApiTestCase(api.ApiTestCase): search_payload = self._search_payload(history_id=history_id, tool_id='cat1', inputs=copied_inputs) search_count = self._search(search_payload) self.assertEquals(search_count, 1) + # Now we delete the original HDA that was used -- we should still be able to find the job delete_respone = self._delete("histories/%s/contents/%s" % (history_id, dataset_id)) self._assert_status_code_is(delete_respone, 200) search_count = self._search(search_payload) self.assertEquals(search_count, 1) + # Now we also delete the copy -- we shouldn't find a job delete_respone = self._delete("histories/%s/contents/%s" % (history_id, new_dataset_id)) self._assert_status_code_is(delete_respone, 200) search_count = self._search(search_payload) @@ -324,6 +326,7 @@ class JobsApiTestCase(api.ApiTestCase): 'f2': {'src': 'hdca', 'id': list_id_a}, }) self._job_search(tool_id='multi_data_param', history_id=history_id, inputs=inputs) + # We test that a job can be found even if the collection has been copied to another history new_history_id = self.dataset_populator.new_history() copy_payload = {"content": list_id_a, "source": "hdca", "type": "dataset_collection"} copy_response = self._post("histories/%s/contents" % new_history_id, data=copy_payload) @@ -336,6 +339,16 @@ class JobsApiTestCase(api.ApiTestCase): search_payload = self._search_payload(history_id=new_history_id, tool_id='multi_data_param', inputs=copied_inputs) search_count = self._search(search_payload) self.assertEquals(search_count, 1) + # Now we delete the original HDCA that was used -- we should still be able to find the job + delete_respone = self._delete("histories/%s/contents/%s" % (history_id, list_id_a)) + self._assert_status_code_is(delete_respone, 200) + search_count = self._search(search_payload) + self.assertEquals(search_count, 1) + # Now we also delete the copy -- we shouldn't find a job + delete_respone = self._delete("histories/%s/contents/%s" % (history_id, new_list_a)) + self._assert_status_code_is(delete_respone, 200) + search_count = self._search(search_payload) + self.assertEquals(search_count, 0) def test_search_with_hdca_list_pair_input(self): history_id, list_id_a = self.__history_with_ok_collection(collection_type='list:pair') From 9e3ceaa7f3410127c7ee8c2c1c3854f4208d5afc Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Fri, 22 Sep 2017 16:52:19 +0200 Subject: [PATCH 08/17] Don't return jobs with deleted HDAs/HDCAs from job_search --- lib/galaxy/webapps/galaxy/api/jobs.py | 90 ++++++++++++---------- test/api/test_jobs.py | 106 +++++++++++++++----------- 2 files changed, 110 insertions(+), 86 deletions(-) diff --git a/lib/galaxy/webapps/galaxy/api/jobs.py b/lib/galaxy/webapps/galaxy/api/jobs.py index 43ea71f6753..dc1772faf5e 100644 --- a/lib/galaxy/webapps/galaxy/api/jobs.py +++ b/lib/galaxy/webapps/galaxy/api/jobs.py @@ -8,7 +8,7 @@ import json import logging from six import string_types -from sqlalchemy import and_, false, or_ +from sqlalchemy import and_, or_ from sqlalchemy.orm import aliased from galaxy import exceptions @@ -283,18 +283,14 @@ class JobController(BaseAPIController, UsesLibraryMixinItems): input_data = {} decoded_input_data = {} input_param = {} - decoded_ids = [] for k, v in inputs.items(): if isinstance(v, dict): if 'id' in v: decoded_id = self.decode_id(v['id']) if 'src' not in v or v['src'] == 'hda': - datasets = [self.hda_manager.get_accessible(decoded_id, trans.user)] - all_history_associations = [] - for dataset in datasets: - instances = trans.sa_session.query(trans.app.model.HistoryDatasetAssociation).filter(trans.app.model.HistoryDatasetAssociation.dataset_id == dataset.dataset_id).all() - all_history_associations.extend(instances) - decoded_ids = set(h.id for h in all_history_associations) + dataset = self.hda_manager.get_accessible(decoded_id, trans.user) + datasets = trans.sa_session.query(model.HistoryDatasetAssociation).filter(model.HistoryDatasetAssociation.dataset_id == dataset.dataset_id).all() + hda_ids = decoded_ids = set(h.id for h in datasets) else: if v['src'] == 'hdca': hdca = self.dataset_collection_manager.get_dataset_collection_instance(trans=trans, @@ -302,14 +298,25 @@ class JobController(BaseAPIController, UsesLibraryMixinItems): id=v['id'], ) copied_from = getattr(hdca, 'copied_from_history_dataset_collection_association_id', None) + if hdca.deleted and not copied_from: + # We don't return the job if an input HDCA is marked as deleted -- it isn't necessary to be that strict, + # we could also just rely on the individual datasets to not be deleted. + return [] if copied_from: # TODO: fetch these recursively or come up with a DB query to fetch all possible hdca ids + if hdca.deleted: + hdca = trans.sa_session.query(model.HistoryDatasetCollectionAssociation).filter( + model.HistoryDatasetCollectionAssociation.id == copied_from).first() + if hdca.deleted: + return [] decoded_ids = [decoded_id, copied_from] else: decoded_ids = [decoded_id] datasets = hdca.dataset_instances + hda_ids = set(h.id for h in datasets) else: - datasets = [self.get_library_dataset_dataset_association(trans, v['id'])] + dataset = self.get_library_dataset_dataset_association(trans, v['id']) + hda_ids = decoded_ids = dataset.id if datasets is None: raise exceptions.ObjectNotFound("Dataset %s not found" % (v['id'])) parameter_values = [] @@ -317,41 +324,45 @@ class JobController(BaseAPIController, UsesLibraryMixinItems): parameter_value = v.copy() parameter_value['id'] = decoded_id parameter_values.append(json.dumps({'values': [parameter_value]})) - decoded_input_data[k] = parameter_values - input_data[k] = [d.dataset_id for d in datasets] + decoded_input_data[k] = parameter_values + # We store the possible input hda ids keyed on input parameter if any of the + # corresponding HDAs exists + input_data[k] = hda_ids if any([d for d in datasets if not d.deleted]) else [] + if not input_data[k]: + return [] else: input_param[k] = json.dumps(str(v)) - query = trans.sa_session.query(trans.app.model.Job).filter( - trans.app.model.Job.tool_id == tool_id, - trans.app.model.Job.user == trans.user + query = trans.sa_session.query(model.Job).filter( + model.Job.tool_id == tool_id, + model.Job.user == trans.user ) if 'state' not in payload: query = query.filter( or_( - trans.app.model.Job.state == 'running', - trans.app.model.Job.state == 'queued', - trans.app.model.Job.state == 'waiting', - trans.app.model.Job.state == 'running', - trans.app.model.Job.state == 'ok', + model.Job.state == 'running', + model.Job.state == 'queued', + model.Job.state == 'waiting', + model.Job.state == 'running', + model.Job.state == 'ok', ) ) else: if isinstance(payload['state'], string_types): - query = query.filter(trans.app.model.Job.state == payload['state']) + query = query.filter(model.Job.state == payload['state']) elif isinstance(payload['state'], list): o = [] for s in payload['state']: - o.append(trans.app.model.Job.state == s) + o.append(model.Job.state == s) query = query.filter( or_(*o) ) for k, v in input_param.items(): - a = aliased(trans.app.model.JobParameter) + a = aliased(model.JobParameter) query = query.filter(and_( - trans.app.model.Job.id == a.job_id, + model.Job.id == a.job_id, a.name == k, a.value == v )) @@ -362,32 +373,27 @@ class JobController(BaseAPIController, UsesLibraryMixinItems): # matches the right input data parameter for a job to be equivalent. # If we don't do this we return false positive jobs in case # a different combination of identical inputs has been used in a job. - a = aliased(trans.app.model.JobParameter) + a = aliased(model.JobParameter) query = query.filter(and_( - trans.app.model.Job.id == a.job_id, + model.Job.id == a.job_id, a.name == k, a.value.in_(v) )) - for k, datasets in input_data.items(): - # Here we are attempting to link the inputs to the underlying - # dataset (not the dataset association). - # This way, if the calculation was done using a copied HDA - # (copied from the library or another history), the search will - # still find the job - a = aliased(trans.app.model.JobToInputDatasetAssociation) - b = aliased(trans.app.model.HistoryDatasetAssociation) - query = query.filter(and_( - trans.app.model.Job.id == a.job_id, - a.dataset_id == b.dataset_id, - b.deleted == false(), - b.dataset_id.in_(datasets), - )) - out = [] for job in query.all(): - # check to make sure none of the output files have been deleted - if all(list(a.dataset.deleted is False for a in job.output_datasets)): + # check to make sure none of the output datasets or collections have been deleted + outputs_deleted = False + for hda in job.output_datasets: + if hda.dataset.deleted: + outputs_deleted = True + break + if not outputs_deleted: + for collection_instance in job.output_dataset_collection_instances: + if collection_instance.dataset_collection_instance.deleted: + outputs_deleted = True + break + if not outputs_deleted: out.append(self.encode_all_ids(trans, job.to_dict('element'), True)) return out diff --git a/test/api/test_jobs.py b/test/api/test_jobs.py index 24c20ca599f..15c97c2dd88 100644 --- a/test/api/test_jobs.py +++ b/test/api/test_jobs.py @@ -289,18 +289,27 @@ class JobsApiTestCase(api.ApiTestCase): 'input1': {'src': 'hda', 'id': new_dataset_id} }) search_payload = self._search_payload(history_id=history_id, tool_id='cat1', inputs=copied_inputs) - search_count = self._search(search_payload) - self.assertEquals(search_count, 1) - # Now we delete the original HDA that was used -- we should still be able to find the job + self._search(search_payload, expected_search_count=1) + # Now we delete the original input HDA that was used -- we should still be able to find the job delete_respone = self._delete("histories/%s/contents/%s" % (history_id, dataset_id)) self._assert_status_code_is(delete_respone, 200) - search_count = self._search(search_payload) - self.assertEquals(search_count, 1) + self._search(search_payload, expected_search_count=1) # Now we also delete the copy -- we shouldn't find a job - delete_respone = self._delete("histories/%s/contents/%s" % (history_id, new_dataset_id)) + delete_respone = self._delete("histories/%s/contents/%s" % (new_history_id, new_dataset_id)) self._assert_status_code_is(delete_respone, 200) - search_count = self._search(search_payload) - self.assertEquals(search_count, 0) + 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({ + 'input1': {'src': 'hda', 'id': dataset_id} + }) + tool_response = self._job_search(tool_id='cat1', history_id=history_id, inputs=inputs) + output_id = tool_response.json()['outputs'][0]['id'] + delete_respone = self._delete("histories/%s/contents/%s" % (history_id, output_id)) + self._assert_status_code_is(delete_respone, 200) + search_payload = self._search_payload(history_id=history_id, tool_id='cat1', inputs=inputs) + self._search(search_payload, expected_search_count=0) def test_search_with_hdca_list_input(self): history_id, list_id_a = self.__history_with_ok_collection(collection_type='list') @@ -309,15 +318,41 @@ class JobsApiTestCase(api.ApiTestCase): 'f1': {'src': 'hdca', 'id': list_id_a}, 'f2': {'src': 'hdca', 'id': list_id_b}, }) - self._job_search(tool_id='multi_data_param', history_id=history_id, inputs=inputs) + tool_response = self._job_search(tool_id='multi_data_param', history_id=history_id, inputs=inputs) # We switch the inputs, this should not return a match - inputs = json.dumps({ + inputs_switched = json.dumps({ 'f2': {'src': 'hdca', 'id': list_id_a}, 'f1': {'src': 'hdca', 'id': list_id_b}, }) + search_payload = self._search_payload(history_id=history_id, tool_id='multi_data_param', inputs=inputs_switched) + self._search(search_payload, expected_search_count=0) + # We delete the ouput (this is a HDA, as multi_data_param reduces collections) + # and use the correct input job definition, the job should not be found + output_id = tool_response.json()['outputs'][0]['id'] + delete_respone = self._delete("histories/%s/contents/%s" % (history_id, output_id)) + self._assert_status_code_is(delete_respone, 200) search_payload = self._search_payload(history_id=history_id, tool_id='multi_data_param', inputs=inputs) - search_count = self._search(search_payload) - self.assertEquals(search_count, 0) + self._search(search_payload, expected_search_count=0) + + def test_search_delete_hdca_output(self): + history_id, list_id_a = self.__history_with_ok_collection(collection_type='list') + inputs = json.dumps({ + 'input1': {'src': 'hdca', 'id': list_id_a}, + }) + tool_response = self._job_search(tool_id='collection_creates_list', history_id=history_id, inputs=inputs) + output_id = tool_response.json()['outputs'][0]['id'] + # We delete a single tool output, no job should be returned + delete_respone = self._delete("histories/%s/contents/%s" % (history_id, output_id)) + self._assert_status_code_is(delete_respone, 200) + search_payload = self._search_payload(history_id=history_id, tool_id='collection_creates_list', inputs=inputs) + self._search(search_payload, expected_search_count=0) + tool_response = self._job_search(tool_id='collection_creates_list', history_id=history_id, inputs=inputs) + output_collection_id = tool_response.json()['output_collections'][0]['id'] + # We delete a collection output, no job should be returned + delete_respone = self._delete("histories/%s/contents/dataset_collections/%s" % (history_id, output_collection_id)) + self._assert_status_code_is(delete_respone, 200) + search_payload = self._search_payload(history_id=history_id, tool_id='collection_creates_list', inputs=inputs) + self._search(search_payload, expected_search_count=0) def test_search_with_hdca_pair_input(self): history_id, list_id_a = self.__history_with_ok_collection(collection_type='pair') @@ -337,18 +372,15 @@ class JobsApiTestCase(api.ApiTestCase): 'f2': {'src': 'hdca', 'id': new_list_a}, }) search_payload = self._search_payload(history_id=new_history_id, tool_id='multi_data_param', inputs=copied_inputs) - search_count = self._search(search_payload) - self.assertEquals(search_count, 1) - # Now we delete the original HDCA that was used -- we should still be able to find the job - delete_respone = self._delete("histories/%s/contents/%s" % (history_id, list_id_a)) + self._search(search_payload, expected_search_count=1) + # Now we delete the original input HDCA that was used -- we should still be able to find the job + delete_respone = self._delete("histories/%s/contents/dataset_collections/%s" % (history_id, list_id_a)) self._assert_status_code_is(delete_respone, 200) - search_count = self._search(search_payload) - self.assertEquals(search_count, 1) + self._search(search_payload, expected_search_count=1) # Now we also delete the copy -- we shouldn't find a job - delete_respone = self._delete("histories/%s/contents/%s" % (history_id, new_list_a)) + delete_respone = self._delete("histories/%s/contents/dataset_collections/%s" % (history_id, new_list_a)) self._assert_status_code_is(delete_respone, 200) - search_count = self._search(search_payload) - self.assertEquals(search_count, 0) + self._search(search_payload, expected_search_count=0) def test_search_with_hdca_list_pair_input(self): history_id, list_id_a = self.__history_with_ok_collection(collection_type='list:pair') @@ -363,10 +395,10 @@ class JobsApiTestCase(api.ApiTestCase): empty_search_response = self._post("jobs/search", data=search_payload) self._assert_status_code_is(empty_search_response, 200) self.assertEquals(len(empty_search_response.json()), 0) - self._post("tools", data=search_payload) - self.dataset_populator.wait_for_history(history_id, assert_ok=True) - search_count = self._search(search_payload) - self.assertEquals(search_count, 1) + tool_response = self._post("tools", data=search_payload) + self.dataset_populator.wait_for_tool_run(history_id, run_response=tool_response) + self._search(search_payload, expected_search_count=1) + return tool_response def _search_payload(self, history_id, tool_id, inputs, state='ok'): search_payload = dict( @@ -377,15 +409,15 @@ class JobsApiTestCase(api.ApiTestCase): ) return search_payload - def _search(self, payload): - search_count = -1 + def _search(self, payload, expected_search_count=1): # in case job and history aren't updated at exactly the same # time give time to wait - for i in range(5): + for i in range(15): search_count = self._search_count(payload) - if search_count == 1: + if search_count == expected_search_count: break - time.sleep(.1) + time.sleep(1) + assert search_count == expected_search_count, "expected to find %d jobs, got %d jobs" % (expected_search_count, search_count) return search_count def _search_count(self, search_payload): @@ -394,20 +426,6 @@ class JobsApiTestCase(api.ApiTestCase): search_json = search_response.json() return len(search_json) - def __run_randomlines_tool(self, lines, history_id, dataset_id): - payload = self.dataset_populator.run_tool_payload( - tool_id="random_lines1", - inputs=dict( - num_lines=lines, - input=dict( - src='hda', - id=dataset_id, - ), - ), - history_id=history_id, - ) - self._post("tools", data=payload) - def __uploads_with_state(self, *states): jobs_response = self._get("jobs", data=dict(state=states)) self._assert_status_code_is(jobs_response, 200) From 07daa4d248ba5f42facb86a0e38370c54fb9e7b6 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sat, 23 Sep 2017 18:29:22 +0200 Subject: [PATCH 09/17] Move job_search out of API for re-use --- lib/galaxy/jobs/search.py | 143 ++++++++++++++++++++++++++ lib/galaxy/webapps/galaxy/api/jobs.py | 132 +----------------------- 2 files changed, 148 insertions(+), 127 deletions(-) create mode 100644 lib/galaxy/jobs/search.py diff --git a/lib/galaxy/jobs/search.py b/lib/galaxy/jobs/search.py new file mode 100644 index 00000000000..60f5717dabd --- /dev/null +++ b/lib/galaxy/jobs/search.py @@ -0,0 +1,143 @@ +import json +from six import string_types +from sqlalchemy import and_, or_ +from sqlalchemy.orm import aliased + +from galaxy import model +from galaxy.managers.hdas import HDAManager +from galaxy.managers.collections import DatasetCollectionManager +from galaxy.managers.lddas import LDDAManager + + +class JobSearch(object): + """Search for jobs using tool inputs or other jobs""" + def __init__(self, app): + self.app = app + self.sa_session = app.model.context + self.hda_manager = HDAManager(app) + self.dataset_collection_manager = DatasetCollectionManager(app) + 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'): + """Search for jobs producing same results using the 'inputs' part of a tool POST.""" + user = trans.user + input_param = {} + 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 [] + parameter_values = [] + for item in all_items: + parameter_value = v.copy() + parameter_value['id'] = item.id + parameter_values.append(json.dumps({'values': [parameter_value]})) + input_param[k] = parameter_values + else: + input_param[k] = json.dumps(str(v)) + return self.__search(tool_id=tool_id, user=user, input_param=input_param, job_state=job_state) + + def one_not_deleted(self, items): + return any((True for item in items if not item.deleted)) + + 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`, find other instances that refer to the same dataset.""" + ldda = self.ldda_manager.get(trans=trans, id=ldda_id) + hdas = self.sa_session.query(model.HistoryDatasetAssociation).filter( + model.HistoryDatasetAssociation.dataset_id == ldda.dataset_id).all() + hdas.append(ldda) + return hdas + + 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, job_state=None): + + query = self.sa_session.query(model.Job).filter( + model.Job.tool_id == tool_id, + model.Job.user == user + ) + + if job_state is None: + query = query.filter( + or_( + model.Job.state == 'running', + model.Job.state == 'queued', + model.Job.state == 'waiting', + model.Job.state == 'running', + model.Job.state == 'ok', + ) + ) + else: + if isinstance(job_state, string_types): + query = query.filter(model.Job.state == job_state) + elif isinstance(job_state, list): + o = [] + for s in job_state: + o.append(model.Job.state == s) + query = query.filter( + 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 + )) + elif isinstance(v, list): + query = query.filter(and_( + model.Job.id == a.job_id, + a.name == k, + a.value.in_(v) + )) + + jobs = [] + for job in query.all(): + # check to make sure none of the output datasets or collections have been deleted + outputs_deleted = False + for hda in job.output_datasets: + if hda.dataset.deleted: + outputs_deleted = True + break + if not outputs_deleted: + for collection_instance in job.output_dataset_collection_instances: + if collection_instance.dataset_collection_instance.deleted: + outputs_deleted = True + break + if not outputs_deleted: + jobs.append(job) + return jobs diff --git a/lib/galaxy/webapps/galaxy/api/jobs.py b/lib/galaxy/webapps/galaxy/api/jobs.py index dc1772faf5e..48a604b7a00 100644 --- a/lib/galaxy/webapps/galaxy/api/jobs.py +++ b/lib/galaxy/webapps/galaxy/api/jobs.py @@ -4,17 +4,15 @@ API operations on a jobs. .. seealso:: :class:`galaxy.model.Jobs` """ -import json import logging from six import string_types -from sqlalchemy import and_, or_ -from sqlalchemy.orm import aliased +from sqlalchemy import or_ from galaxy import exceptions -from galaxy import managers from galaxy import model from galaxy import util +from galaxy.jobs.search import JobSearch 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 @@ -27,9 +25,7 @@ class JobController(BaseAPIController, UsesLibraryMixinItems): def __init__(self, app): super(JobController, self).__init__(app) - self.hda_manager = managers.hdas.HDAManager(app) - self.dataset_collection_manager = managers.collections.DatasetCollectionManager(app) - self.dataset_manager = managers.datasets.DatasetManager(app) + self.job_search = JobSearch(app) @expose_api def index(self, trans, **kwd): @@ -267,135 +263,17 @@ class JobController(BaseAPIController, UsesLibraryMixinItems): the exact some input parameters and datasets. This can be used to minimize the amount of repeated work, and simply recycle the old results. """ - tool_id = payload.get('tool_id') if tool_id is None: raise exceptions.ObjectAttributeMissingException("No tool id") - tool = trans.app.toolbox.get_tool(tool_id) if tool is None: raise exceptions.ObjectNotFound("Requested tool not found") if 'inputs' not in payload: raise exceptions.ObjectAttributeMissingException("No inputs defined") - inputs = payload['inputs'] - - input_data = {} - decoded_input_data = {} - input_param = {} - for k, v in inputs.items(): - if isinstance(v, dict): - if 'id' in v: - decoded_id = self.decode_id(v['id']) - if 'src' not in v or v['src'] == 'hda': - dataset = self.hda_manager.get_accessible(decoded_id, trans.user) - datasets = trans.sa_session.query(model.HistoryDatasetAssociation).filter(model.HistoryDatasetAssociation.dataset_id == dataset.dataset_id).all() - hda_ids = decoded_ids = set(h.id for h in datasets) - else: - if v['src'] == 'hdca': - hdca = self.dataset_collection_manager.get_dataset_collection_instance(trans=trans, - instance_type='history', - id=v['id'], - ) - copied_from = getattr(hdca, 'copied_from_history_dataset_collection_association_id', None) - if hdca.deleted and not copied_from: - # We don't return the job if an input HDCA is marked as deleted -- it isn't necessary to be that strict, - # we could also just rely on the individual datasets to not be deleted. - return [] - if copied_from: - # TODO: fetch these recursively or come up with a DB query to fetch all possible hdca ids - if hdca.deleted: - hdca = trans.sa_session.query(model.HistoryDatasetCollectionAssociation).filter( - model.HistoryDatasetCollectionAssociation.id == copied_from).first() - if hdca.deleted: - return [] - decoded_ids = [decoded_id, copied_from] - else: - decoded_ids = [decoded_id] - datasets = hdca.dataset_instances - hda_ids = set(h.id for h in datasets) - else: - dataset = self.get_library_dataset_dataset_association(trans, v['id']) - hda_ids = decoded_ids = dataset.id - if datasets is None: - raise exceptions.ObjectNotFound("Dataset %s not found" % (v['id'])) - parameter_values = [] - for decoded_id in decoded_ids: - parameter_value = v.copy() - parameter_value['id'] = decoded_id - parameter_values.append(json.dumps({'values': [parameter_value]})) - decoded_input_data[k] = parameter_values - # We store the possible input hda ids keyed on input parameter if any of the - # corresponding HDAs exists - input_data[k] = hda_ids if any([d for d in datasets if not d.deleted]) else [] - if not input_data[k]: - return [] - else: - input_param[k] = json.dumps(str(v)) - - query = trans.sa_session.query(model.Job).filter( - model.Job.tool_id == tool_id, - model.Job.user == trans.user - ) - - if 'state' not in payload: - query = query.filter( - or_( - model.Job.state == 'running', - model.Job.state == 'queued', - model.Job.state == 'waiting', - model.Job.state == 'running', - model.Job.state == 'ok', - ) - ) - else: - if isinstance(payload['state'], string_types): - query = query.filter(model.Job.state == payload['state']) - elif isinstance(payload['state'], list): - o = [] - for s in payload['state']: - o.append(model.Job.state == s) - query = query.filter( - or_(*o) - ) - - for k, v in input_param.items(): - a = aliased(model.JobParameter) - query = query.filter(and_( - model.Job.id == a.job_id, - a.name == k, - a.value == v - )) - - for k, v in decoded_input_data.items(): - # Here we make sure that the incoming hdca/hda id - # (v is a list of possible text representations of a input data parameter) - # matches the right input data parameter for a job to be equivalent. - # If we don't do this we return false positive jobs in case - # a different combination of identical inputs has been used in a job. - a = aliased(model.JobParameter) - query = query.filter(and_( - model.Job.id == a.job_id, - a.name == k, - a.value.in_(v) - )) - - out = [] - for job in query.all(): - # check to make sure none of the output datasets or collections have been deleted - outputs_deleted = False - for hda in job.output_datasets: - if hda.dataset.deleted: - outputs_deleted = True - break - if not outputs_deleted: - for collection_instance in job.output_dataset_collection_instances: - if collection_instance.dataset_collection_instance.deleted: - outputs_deleted = True - break - if not outputs_deleted: - out.append(self.encode_all_ids(trans, job.to_dict('element'), True)) - return out + 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] @expose_api def error(self, trans, id, **kwd): From 054cdd146a84460db07f4725a5251036bb42278f Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sun, 24 Sep 2017 10:22:43 +0200 Subject: [PATCH 10/17] When querying jobs for data input parameters drop all other parameters --- lib/galaxy/jobs/search.py | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/lib/galaxy/jobs/search.py b/lib/galaxy/jobs/search.py index 60f5717dabd..c8516c3c470 100644 --- a/lib/galaxy/jobs/search.py +++ b/lib/galaxy/jobs/search.py @@ -45,9 +45,7 @@ class JobSearch(object): return [] parameter_values = [] for item in all_items: - parameter_value = v.copy() - parameter_value['id'] = item.id - parameter_values.append(json.dumps({'values': [parameter_value]})) + parameter_values.append(json.dumps({'values': [{'src': src, 'id': item.id}]})) input_param[k] = parameter_values else: input_param[k] = json.dumps(str(v)) From 19818cb380889e584741caf4a78462e0502f30c1 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sun, 24 Sep 2017 19:41:19 +0200 Subject: [PATCH 11/17] Fix loading of model.JobToInputDatasetCollectionAssociation If we don't load this relationship we do not properly record a job's input_dataset_collections in `DefaultToolAction._record_inputs`. --- lib/galaxy/model/mapping.py | 1 + 1 file changed, 1 insertion(+) diff --git a/lib/galaxy/model/mapping.py b/lib/galaxy/model/mapping.py index d4f0e0c245e..95b5ebdf329 100644 --- a/lib/galaxy/model/mapping.py +++ b/lib/galaxy/model/mapping.py @@ -2137,6 +2137,7 @@ mapper(model.Job, model.Job.table, properties=dict( library_folder=relation(model.LibraryFolder, lazy=True), parameters=relation(model.JobParameter, lazy=True), input_datasets=relation(model.JobToInputDatasetAssociation), + input_dataset_collections=relation(model.JobToInputDatasetCollectionAssociation, lazy=True), output_datasets=relation(model.JobToOutputDatasetAssociation, lazy=True), output_dataset_collection_instances=relation(model.JobToOutputDatasetCollectionAssociation, lazy=True), output_dataset_collections=relation(model.JobToImplicitOutputDatasetCollectionAssociation, lazy=True), From 87439012602cd500346a1fe442292b2e7cbaaeab Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sun, 24 Sep 2017 19:56:01 +0200 Subject: [PATCH 12/17] Filter jobs on input data using JobToInput*Association. This has the advantage that for complex parameters we don't need to decompose parameter names like `queries_0|input2` in order to compare JobParameters. --- lib/galaxy/jobs/search.py | 50 ++++++++++++++++++++++++++------------- 1 file changed, 33 insertions(+), 17 deletions(-) diff --git a/lib/galaxy/jobs/search.py b/lib/galaxy/jobs/search.py index c8516c3c470..3a47da04dfa 100644 --- a/lib/galaxy/jobs/search.py +++ b/lib/galaxy/jobs/search.py @@ -23,6 +23,7 @@ class JobSearch(object): """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: @@ -43,13 +44,10 @@ class JobSearch(object): return [] else: return [] - parameter_values = [] - for item in all_items: - parameter_values.append(json.dumps({'values': [{'src': src, 'id': item.id}]})) - input_param[k] = parameter_values + input_data[k] = {src: {item.id for item in all_items}} else: input_param[k] = json.dumps(str(v)) - return self.__search(tool_id=tool_id, user=user, input_param=input_param, job_state=job_state) + return self.__search(tool_id=tool_id, user=user, input_param=input_param, input_data=input_data, job_state=job_state) def one_not_deleted(self, items): return any((True for item in items if not item.deleted)) @@ -61,12 +59,10 @@ class JobSearch(object): model.HistoryDatasetAssociation.dataset_id == hda.dataset_id).all() def _get_all_lddas(self, trans, ldda_id): - """Given a decoded `ldda_id`, find other instances that refer to the same dataset.""" + """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) - hdas = self.sa_session.query(model.HistoryDatasetAssociation).filter( - model.HistoryDatasetAssociation.dataset_id == ldda.dataset_id).all() - hdas.append(ldda) - return hdas + 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.""" @@ -80,7 +76,7 @@ class JobSearch(object): hdcas.append(copied_from) return hdcas - def __search(self, tool_id, user, input_param, job_state=None): + def __search(self, tool_id, user, input_param, input_data, job_state=None): query = self.sa_session.query(model.Job).filter( model.Job.tool_id == tool_id, @@ -116,16 +112,36 @@ class JobSearch(object): a.name == k, a.value == v )) - elif isinstance(v, list): - query = query.filter(and_( - model.Job.id == a.job_id, - a.name == k, - a.value.in_(v) - )) + for k, type_values in input_data.items(): + for t, v in type_values.items(): + if t == 'hda': + a = aliased(model.JobToInputDatasetAssociation) + query = query.filter(and_( + model.Job.id == a.job_id, + a.name == k, + a.dataset_id.in_(v) + )) + elif t == 'ldda': + a = aliased(model.JobToInputDatasetCollectionAssociation) + query = query.filter(and_( + model.Job.id == a.job_id, + a.name == k, + a.ldda_id.in_(v) + )) + elif t == 'hdca': + a = aliased(model.JobToInputLibraryDatasetAssociation) + query = query.filter(and_( + model.Job.id == a.job_id, + a.name == k, + a.dataset_collection_id.in_(v) + )) + else: + return [] jobs = [] for job in query.all(): # check to make sure none of the output datasets or collections have been deleted + # TODO: find copies if output is deleted outputs_deleted = False for hda in job.output_datasets: if hda.dataset.deleted: From 5fc360a0829ae6819dc96ad0228940cda59d4fa7 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sun, 24 Sep 2017 20:50:15 +0200 Subject: [PATCH 13/17] Move JobSearch to galaxy.managers.jobs --- lib/galaxy/{jobs/search.py => managers/jobs.py} | 4 ++-- lib/galaxy/webapps/galaxy/api/jobs.py | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) rename lib/galaxy/{jobs/search.py => managers/jobs.py} (97%) diff --git a/lib/galaxy/jobs/search.py b/lib/galaxy/managers/jobs.py similarity index 97% rename from lib/galaxy/jobs/search.py rename to lib/galaxy/managers/jobs.py index 3a47da04dfa..d6d00c2756e 100644 --- a/lib/galaxy/jobs/search.py +++ b/lib/galaxy/managers/jobs.py @@ -122,14 +122,14 @@ class JobSearch(object): a.dataset_id.in_(v) )) elif t == 'ldda': - a = aliased(model.JobToInputDatasetCollectionAssociation) + a = aliased(model.JobToInputLibraryDatasetAssociation) query = query.filter(and_( model.Job.id == a.job_id, a.name == k, a.ldda_id.in_(v) )) elif t == 'hdca': - a = aliased(model.JobToInputLibraryDatasetAssociation) + a = aliased(model.JobToInputDatasetCollectionAssociation) query = query.filter(and_( model.Job.id == a.job_id, a.name == k, diff --git a/lib/galaxy/webapps/galaxy/api/jobs.py b/lib/galaxy/webapps/galaxy/api/jobs.py index 48a604b7a00..7cc56b125ff 100644 --- a/lib/galaxy/webapps/galaxy/api/jobs.py +++ b/lib/galaxy/webapps/galaxy/api/jobs.py @@ -12,7 +12,7 @@ from sqlalchemy import or_ from galaxy import exceptions from galaxy import model from galaxy import util -from galaxy.jobs.search import JobSearch +from galaxy.managers.jobs import JobSearch 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 0e287d38ca1b779c306eb8e0a4c9d3c8bc140882 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Sun, 24 Sep 2017 20:55:34 +0200 Subject: [PATCH 14/17] Return empty list if type is unknown --- lib/galaxy/managers/jobs.py | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/lib/galaxy/managers/jobs.py b/lib/galaxy/managers/jobs.py index d6d00c2756e..311ca32d1a6 100644 --- a/lib/galaxy/managers/jobs.py +++ b/lib/galaxy/managers/jobs.py @@ -44,6 +44,10 @@ class JobSearch(object): 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(str(v)) From 480ee7980251848d09e837faef1cd9873b2e2e33 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Mon, 25 Sep 2017 09:06:35 +0200 Subject: [PATCH 15/17] Skip recording JobToInputDatasetCollection relations if the item to record is not a HistoryDatasetCollectionAssociation. This fixes ``` Error executing tool: Attempting to flush an item of type as a member of collection "JobToInputDatasetCollectionAssociation.dataset_collection". Expected an object of type or a polymorphic subclass of this type.' ``` It isn't ideal, but before 19818cb we didn't record Collection-like elements at all. --- lib/galaxy/tools/actions/__init__.py | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/lib/galaxy/tools/actions/__init__.py b/lib/galaxy/tools/actions/__init__.py index 17f64e88f5f..a71581552f4 100644 --- a/lib/galaxy/tools/actions/__init__.py +++ b/lib/galaxy/tools/actions/__init__.py @@ -559,7 +559,11 @@ class DefaultToolAction(object): reductions[name].append(dataset_collection) # TODO: verify can have multiple with same name, don't want to lose traceability - job.add_input_dataset_collection(name, dataset_collection) + if isinstance(dataset_collection, model.HistoryDatasetCollectionAssociation): + # FIXME: when recording inputs for special tools (e.g. ModelOperationToolAction), + # dataset_collection is actually a DatasetCollectionElement, which can't be added + # to a jobs' input_dataset_collection relation, which expects HDCA instances + job.add_input_dataset_collection(name, dataset_collection) # If this an input collection is a reduction, we expanded it for dataset security, type # checking, and such, but the persisted input must be the original collection From 8f5125194630def8530c0e3e30f1b9a4439bd0fb Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Mon, 25 Sep 2017 12:00:22 +0200 Subject: [PATCH 16/17] Linting fixes and drop unnecessary internal parenthesis (thanks @nsoranzo!). --- .ci/flake8_lint_include_list.txt | 1 + lib/galaxy/managers/jobs.py | 9 +++++---- 2 files changed, 6 insertions(+), 4 deletions(-) diff --git a/.ci/flake8_lint_include_list.txt b/.ci/flake8_lint_include_list.txt index 4cb3e453782..63578af8248 100644 --- a/.ci/flake8_lint_include_list.txt +++ b/.ci/flake8_lint_include_list.txt @@ -52,6 +52,7 @@ lib/galaxy/managers/collections_util.py lib/galaxy/managers/context.py lib/galaxy/managers/deletable.py lib/galaxy/managers/__init__.py +lib/galaxy/managers/jobs.py lib/galaxy/managers/lddas.py lib/galaxy/managers/libraries.py lib/galaxy/managers/secured.py diff --git a/lib/galaxy/managers/jobs.py b/lib/galaxy/managers/jobs.py index 311ca32d1a6..4e5c3b2187a 100644 --- a/lib/galaxy/managers/jobs.py +++ b/lib/galaxy/managers/jobs.py @@ -1,11 +1,12 @@ import json + from six import string_types from sqlalchemy import and_, or_ from sqlalchemy.orm import aliased from galaxy import model -from galaxy.managers.hdas import HDAManager from galaxy.managers.collections import DatasetCollectionManager +from galaxy.managers.hdas import HDAManager from galaxy.managers.lddas import LDDAManager @@ -40,7 +41,7 @@ class JobSearch(object): 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)): + if any(True for hda in all_items[0].dataset_instances if hda.deleted): return [] else: return [] @@ -50,11 +51,11 @@ class JobSearch(object): return [] input_data[k] = {src: {item.id for item in all_items}} else: - input_param[k] = json.dumps(str(v)) + 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) def one_not_deleted(self, items): - return any((True for item in items if not item.deleted)) + return any(True for item in items if not item.deleted) def _get_all_hdas(self, trans, hda_id): """Given a decoded `hda_id`, find other instances that refer to the same dataset.""" From 7d260d0b50a5858af533fe136105ce357b4ef490 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Wed, 27 Sep 2017 18:29:04 +0200 Subject: [PATCH 17/17] 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):