Merge pull request #4665 from mvdbeek/extend_job_search

Extend job search for hdcas
This commit is contained in:
John Chilton
2017-09-29 08:38:28 -04:00
committed by GitHub
10 changed files with 445 additions and 147 deletions
+1
View File
@@ -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
@@ -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
+1
View File
@@ -14,6 +14,7 @@ pycrypto
# pure Python packages
bz2file; python_version < '3.3'
boltons
Paste
PasteDeploy
docutils
+245
View File
@@ -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
+1
View File
@@ -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),
+16 -11
View File
@@ -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)
+5 -1
View File
@@ -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
+3
View File
@@ -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.
+24 -85
View File
@@ -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):
+148 -50
View File
@@ -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)