diff --git a/.ci/flake8_lint_include_list.txt b/.ci/flake8_lint_include_list.txt index 6d90bc4ab6f..3e14bb58636 100644 --- a/.ci/flake8_lint_include_list.txt +++ b/.ci/flake8_lint_include_list.txt @@ -49,6 +49,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/dependencies/pinned-requirements.txt b/lib/galaxy/dependencies/pinned-requirements.txt index 80b2bcffa29..7171eb56fac 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 new file mode 100644 index 00000000000..a291d377b23 --- /dev/null +++ b/lib/galaxy/managers/jobs.py @@ -0,0 +1,245 @@ +import json +import logging + +from boltons.iterutils import remap +from six import string_types +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): + """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, 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_data = defaultdict(list) + input_ids = defaultdict(dict) + + 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 + + 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 + ) + + 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, 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 == 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 == 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, + 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 [] + + 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: refactors this into the initial job query + 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: + log.info("Searching jobs finished %s", search_timer) + return job + return None 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), 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/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 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 63367938600..df2b0360efd 100644 --- a/lib/galaxy/webapps/galaxy/api/jobs.py +++ b/lib/galaxy/webapps/galaxy/api/jobs.py @@ -4,21 +4,20 @@ API operations on a jobs. .. seealso:: :class:`galaxy.model.Jobs` """ -import json import logging from six import string_types -from sqlalchemy import and_, false, 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.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 galaxy.web.base.controller import UsesLibraryMixinItems +from galaxy.work.context import WorkRequestContext log = logging.getLogger(__name__) @@ -27,8 +26,7 @@ class JobController(BaseAPIController, UsesLibraryMixinItems): def __init__(self, app): super(JobController, self).__init__(app) - self.hda_manager = managers.hdas.HDAManager(app) - self.dataset_manager = managers.datasets.DatasetManager(app) + self.job_search = JobSearch(app) @expose_api def index(self, trans, **kwd): @@ -266,92 +264,33 @@ 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 = 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") - 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 = {} - input_param = {} - for k, v in inputs.items(): - if isinstance(v, dict): - 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) - else: - dataset = self.get_library_dataset_dataset_association(trans, v['id']) - if dataset is None: - raise exceptions.ObjectNotFound("Dataset %s not found" % (v['id'])) - input_data[k] = dataset.dataset_id - 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 - ) - - 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', - ) - ) - else: - if isinstance(payload['state'], string_types): - query = query.filter(trans.app.model.Job.state == payload['state']) - elif isinstance(payload['state'], list): - o = [] - for s in payload['state']: - o.append(trans.app.model.Job.state == s) - query = query.filter( - or_(*o) - ) - - for k, v in input_param.items(): - a = aliased(trans.app.model.JobParameter) - query = query.filter(and_( - trans.app.model.Job.id == a.job_id, - a.name == k, - a.value == v - )) - - for k, v 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.id, - b.deleted == false(), - b.dataset_id == v - )) - - 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)): - out.append(self.encode_all_ids(trans, job.to_dict('element'), True)) - return out + 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): diff --git a/test/api/test_jobs.py b/test/api/test_jobs.py index 6cb1565abc3..15c97c2dd88 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... @@ -273,38 +275,150 @@ class JobsApiTestCase(api.ApiTestCase): def test_search(self): history_id, dataset_id = self.__history_with_ok_dataset() + inputs = json.dumps({ + '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) + 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) + 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" % (new_history_id, new_dataset_id)) + self._assert_status_code_is(delete_respone, 200) + self._search(search_payload, expected_search_count=0) - inputs = json.dumps( - dict( - input1=dict( - src='hda', - id=dataset_id, - ) - ) - ) - search_payload = dict( - tool_id="cat1", - inputs=inputs, - state="ok", - ) + 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') + 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}, + }) + 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_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) + 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') + inputs = json.dumps({ + 'f1': {'src': 'hdca', 'id': list_id_a}, + '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) + 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) + 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) + 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/dataset_collections/%s" % (history_id, new_list_a)) + self._assert_status_code_is(delete_respone, 200) + 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') + inputs = json.dumps({ + 'f1': {'src': 'hdca', 'id': list_id_a}, + 'f2': {'src': 'hdca', 'id': list_id_a}, + }) + self._job_search(tool_id='multi_data_param', history_id=history_id, inputs=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) self.assertEquals(len(empty_search_response.json()), 0) + 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 - self.__run_cat_tool(history_id, dataset_id) - self.dataset_populator.wait_for_history(history_id, assert_ok=True) + 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 - 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): - search_count = self._search_count(search_payload) - if search_count == 1: + for i in range(15): + search_count = self._search_count(payload) + if search_count == expected_search_count: break - time.sleep(.1) - - self.assertEquals(search_count, 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): search_response = self._post("jobs/search", data=search_payload) @@ -312,34 +426,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", - 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) @@ -357,6 +443,18 @@ 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, 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'] + def __jobs_index(self, **kwds): jobs_response = self._get("jobs", **kwds) self._assert_status_code_is(jobs_response, 200)