Take into account HDA name, metadata, element_identifier and datatype

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