Merge pull request #4690 from mvdbeek/skip_job_execution

Skip job execution if equivalent job exists
This commit is contained in:
John Chilton
2018-01-02 10:29:52 -05:00
committed by GitHub
34 changed files with 803 additions and 191 deletions
@@ -200,6 +200,7 @@ var View = Backbone.View.extend({
this._renderMessage();
this._renderParameters();
this._renderHistory();
this._renderUseCachedJob();
_.each(this.steps, step => {
self._renderStep(step);
});
@@ -311,6 +312,36 @@ var View = Backbone.View.extend({
this._append(this.$steps, this.history_form.$el);
},
/** Render job caching option */
_renderUseCachedJob: function() {
var extra_user_preferences = {};
if (Galaxy.user.attributes.preferences && 'extra_user_preferences' in Galaxy.user.attributes.preferences){
extra_user_preferences = JSON.parse(Galaxy.user.attributes.preferences.extra_user_preferences);
}
var display_use_cached_job_checkbox = 'use_cached_job|use_cached_job_checkbox' in extra_user_preferences ? extra_user_preferences['use_cached_job|use_cached_job_checkbox'] : false ;
this.display_use_cached_job_checkbox = display_use_cached_job_checkbox === 'true';
if (this.display_use_cached_job_checkbox){
this.job_options_form = new Form({
cls: "ui-portlet-narrow",
title: "<b>Job re-use Options</b>",
inputs: [
{
type: "conditional",
name: "use_cached_job",
test_param: {
name: "check",
label: "BETA: Attempt to reuse jobs with identical parameters?",
type: "boolean",
value: "false",
help: "This may skip executing jobs that you have already run."
},
}
]
});
this._append(this.$steps, this.job_options_form.$el);
}
},
/** Render step */
_renderStep: function(step) {
var self = this;
@@ -510,6 +541,9 @@ var View = Backbone.View.extend({
// so that inputs can be batched.
batch: true
};
if (this.display_use_cached_job_checkbox) {
job_def['use_cached_job'] = this.job_options_form.data.create()["use_cached_job|check"] === 'true';
}
var validated = true;
for (var i in this.forms) {
var form = this.forms[i];
@@ -148,6 +148,25 @@ var View = Backbone.View.extend({
"The previous run of this tool failed and other tools were waiting for it to finish successfully. Use this option to resume those tools using the new output(s) of this tool run."
});
}
// Job Re-use Options
var extra_user_preferences = {};
if (Galaxy.user.attributes.preferences && 'extra_user_preferences' in Galaxy.user.attributes.preferences) {
extra_user_preferences = JSON.parse(Galaxy.user.attributes.preferences.extra_user_preferences);
}
var use_cached_job = 'use_cached_job|use_cached_job_checkbox' in extra_user_preferences ? extra_user_preferences['use_cached_job|use_cached_job_checkbox'] : false ;
if (use_cached_job === 'true'){
options.inputs.push({
label: "BETA: Attempt to re-use jobs with identical parameters ?",
help: "This may skip executing jobs that you have already run",
name: "use_cached_job",
type: "select",
display: "radio",
ignore: "__ignore__",
value: "__ignore__",
options: [["No", false], ["Yes", true]],
});
}
},
/** Submit a regular job.
@@ -39,3 +39,13 @@ preferences:
label: Maximum number of search results
type: text
required: False
use_cached_job:
description: Do you want to be able to re-use previously run jobs ?
inputs:
- name: use_cached_job_checkbox
label: Do you want to be able to re-use equivalent jobs ?
type: boolean
checked: false
value: false
help: If you select yes, you will be able to select for each tool and workflow run if you would like to use this feature.
+15
View File
@@ -305,6 +305,21 @@ class JobHandlerQueue(Monitors, object):
try:
# Check the job's dependencies, requeue if they're not done.
# Some of these states will only happen when using the in-memory job queue
if job.copied_from_job_id:
copied_from_job = self.sa_session.query(model.Job).get(job.copied_from_job_id)
job.numeric_metrics = copied_from_job.numeric_metrics
job.text_metrics = copied_from_job.text_metrics
job.dependencies = copied_from_job.dependencies
job.state = copied_from_job.state
job.stderr = copied_from_job.stderr
job.stdout = copied_from_job.stdout
job.command_line = copied_from_job.command_line
job.traceback = copied_from_job.traceback
job.tool_version = copied_from_job.tool_version
job.exit_code = copied_from_job.exit_code
job.job_runner_name = copied_from_job.job_runner_name
job.job_runner_external_id = copied_from_job.job_runner_external_id
continue
job_state = self.__check_job_state(job)
if job_state == JOB_WAIT:
new_waiting_jobs.append(job.id)
+132 -115
View File
@@ -49,11 +49,10 @@ class JobSearch(object):
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):
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)
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."""
@@ -63,187 +62,205 @@ 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})
input_ids[src][value] = True
return key, 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)
input_data[path_key].append({'src': src,
'id': value,
'identifier': identifier,
})
return key, "__id_wildcard__"
return key, value
remap(param_dump, visit=populate_input_data_input_id)
wildcard_param_dump = remap(param_dump, visit=populate_input_data_input_id)
return self.__search(tool_id=tool_id,
tool_version=tool_version,
user=user,
input_data=input_data,
job_state=job_state,
param_dump=param_dump,
input_ids=input_ids,
is_workflow_step=is_workflow_step)
wildcard_param_dump=wildcard_param_dump)
def __search(self, tool_id, user, input_data, input_ids=None, job_state=None, param_dump=None, is_workflow_step=False):
def __search(self, tool_id, tool_version, user, input_data, job_state=None, param_dump=None, wildcard_param_dump=None):
search_timer = ExecutionTimer()
query = self.sa_session.query(model.Job).filter(
model.Job.tool_id == tool_id,
model.Job.user == user
)
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
conditions = [and_(model.Job.tool_id == tool_id,
model.Job.user == user)]
if tool_version:
conditions.append(model.Job.tool_version == str(tool_version))
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',
)
conditions.append(
model.Job.state.in_([model.Job.states.NEW,
model.Job.states.QUEUED,
model.Job.states.WAITING,
model.Job.states.RUNNING,
model.Job.states.OK])
)
else:
if isinstance(job_state, string_types):
query = query.filter(model.Job.state == job_state)
conditions.append(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(
conditions.append(
or_(*o)
)
# We now build the query filters that relate to the input datasets
# that this job uses. We keep track of the requested dataset id in `requested_ids`,
# the type (hda, hdca or lda) in `data_types`
# and the ids that have been used in the job that has already been run in `used_ids`.
requested_ids = []
data_types = []
used_ids = []
for k, input_list in input_data.items():
for type_values in input_list:
t = type_values['src']
v = type_values['id']
requested_ids.append(v)
data_types.append(t)
identifier = type_values['identifier']
if t == 'hda':
a = aliased(model.JobToInputDatasetAssociation)
b = aliased(model.HistoryDatasetAssociation)
c = aliased(model.HistoryDatasetAssociation)
query = query.filter(and_(
d = aliased(model.JobParameter)
e = aliased(model.HistoryDatasetAssociationHistory)
conditions.append(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 compare input dataset names and metadata for a job
# if we know that the input dataset hasn't changed since the job was run,
# or if the job recorded a dataset_version and name and metadata of the current
# job request matches those that were recorded for the job in question (introduced in release 18.01)
or_(and_(
b.update_time < model.Job.create_time,
b.name == c.name,
b.extension == c.extension,
b.metadata == c.metadata,
), and_(
b.id == e.history_dataset_association_id,
a.dataset_version == e.version,
e.name == c.name,
e.extension == c.extension,
e._metadata == c._metadata,
)),
or_(b.deleted == false(), c.deleted == false())
))
if identifier:
conditions.append(and_(model.Job.id == d.job_id,
d.name == "%s|__identifier__" % k,
d.value == json.dumps(identifier)))
used_ids.append(a.dataset_id)
elif t == 'ldda':
a = aliased(model.JobToInputLibraryDatasetAssociation)
query = query.filter(and_(
conditions.append(and_(
model.Job.id == a.job_id,
a.name == k,
a.ldda_id == v
))
used_ids.append(a.ldda_id)
elif t == 'hdca':
a = aliased(model.JobToInputDatasetCollectionAssociation)
b = aliased(model.HistoryDatasetCollectionAssociation)
c = aliased(model.HistoryDatasetCollectionAssociation)
query = query.filter(and_(
conditions.append(and_(
model.Job.id == a.job_id,
a.name == k,
b.id == a.dataset_collection_id,
c.id == v,
b.name == c.name,
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()
c.deleted == false(),
)
)
))
used_ids.append(a.dataset_collection_id)
else:
return []
for k, v in wildcard_param_dump.items():
wildcard_value = json.dumps(v).replace('"id": "__id_wildcard__"', '"id": %')
a = aliased(model.JobParameter)
conditions.append(and_(
model.Job.id == a.job_id,
a.name == k,
a.value.like(wildcard_value)
))
conditions.append(and_(
model.Job.any_output_dataset_collection_instances_deleted == false(),
model.Job.any_output_dataset_deleted == false()
))
query = self.sa_session.query(model.Job.id, *used_ids).filter(and_(*conditions))
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 len(job) > 1:
# We do have datasets to check
job_id, current_jobs_data_ids = job[0], job[1:]
job_parameter_conditions = [model.Job.id == job_id]
for src, requested_id, used_id in zip(data_types, requested_ids, current_jobs_data_ids):
if src not in job_input_ids:
job_input_ids[src] = {dataset_id: job_dataset_id}
job_input_ids[src] = {requested_id: used_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
job_input_ids[src][requested_id] = used_id
new_param_dump = remap(param_dump, visit=replace_dataset_ids)
# 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.
for k, v in new_param_dump.items():
a = aliased(model.JobParameter)
job_parameter_conditions.append(and_(
a.name == k,
a.value == json.dumps(v)
))
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
job_parameter_conditions = [model.Job.id == job]
query = self.sa_session.query(model.Job).filter(*job_parameter_conditions)
job = query.first()
if job is None:
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
n_parameters = 0
# Verify that equivalent jobs had the same number of job parameters
# We skip chrominfo, dbkey, __workflow_invocation_uuid__ and identifer
# parameter as these are not passed along when expanding tool parameters
# and they can differ without affecting the resulting dataset.
for parameter in job.parameters:
if parameter.name in {'__workflow_invocation_uuid__', 'chromInfo', 'dbkey'} or parameter.name.endswith('|__identifier__'):
continue
n_parameters += 1
if not n_parameters == len(param_dump):
continue
log.info("Found equivalent job %s", search_timer)
return job
log.info("No equivalent jobs found %s", search_timer)
return None
+84 -5
View File
@@ -5,6 +5,7 @@ Naming: try to use class names that have a distinct plural form so that
the relationship cardinalities are obvious (e.g. prefer Dataset to Data)
"""
import errno
import json
import logging
import numbers
import operator
@@ -16,11 +17,21 @@ from string import Template
from uuid import UUID, uuid4
from six import string_types
from sqlalchemy import (and_, func, join, not_, or_, select, true, type_coerce,
types)
from sqlalchemy import (
and_,
func,
inspect,
join,
not_,
or_,
select,
true,
type_coerce,
types)
from sqlalchemy.ext import hybrid
from sqlalchemy.orm import aliased, joinedload, object_session
import galaxy.model.metadata
import galaxy.model.orm.now
import galaxy.security.passwords
@@ -201,6 +212,17 @@ class User(object, Dictifiable):
self.credentials = []
# ? self.roles = []
@property
def extra_preferences(self):
data = {}
extra_user_preferences = self.preferences.get('extra_user_preferences')
if extra_user_preferences:
try:
data = json.loads(extra_user_preferences)
except Exception:
pass
return data
def set_password_cleartext(self, cleartext):
"""
Set user password to the digest of `cleartext`.
@@ -474,6 +496,7 @@ class Job(object, JobLike, UsesCreateAndUpdateTime, Dictifiable):
self.user_id = None
self.tool_id = None
self.tool_version = None
self.copied_from_job_id = None
self.command_line = None
self.dependencies = []
self.param_filename = None
@@ -542,6 +565,9 @@ class Job(object, JobLike, UsesCreateAndUpdateTime, Dictifiable):
def get_parameters(self):
return self.parameters
def get_copied_from_job_id(self):
return self.copied_from_job_id
def get_input_datasets(self):
return self.input_datasets
@@ -625,6 +651,9 @@ class Job(object, JobLike, UsesCreateAndUpdateTime, Dictifiable):
def set_parameters(self, parameters):
self.parameters = parameters
def set_copied_from_job_id(self, job_id):
self.copied_from_job_id = job_id
def set_input_datasets(self, input_datasets):
self.input_datasets = input_datasets
@@ -993,6 +1022,7 @@ class JobToInputDatasetAssociation(object):
def __init__(self, name, dataset):
self.name = name
self.dataset = dataset
self.dataset_version = dataset.version if dataset else None
class JobToOutputDatasetAssociation(object):
@@ -1957,7 +1987,7 @@ class DatasetInstance(object):
self.metadata = metadata or dict()
self.extended_metadata = extended_metadata
if dbkey: # dbkey is stored in metadata, only set if non-zero, or else we could clobber one supplied by input 'metadata'
self.dbkey = dbkey
self._metadata['dbkey'] = dbkey
self.deleted = deleted
self.visible = visible
# Relationships
@@ -2043,8 +2073,6 @@ class DatasetInstance(object):
if "dbkey" in self.datatype.metadata_spec:
if not isinstance(value, list):
self.metadata.dbkey = [value]
else:
self.metadata.dbkey = value
dbkey = property(get_dbkey, set_dbkey)
def change_datatype(self, new_ext):
@@ -2394,6 +2422,36 @@ class HistoryDatasetAssociation(DatasetInstance, HasTags, Dictifiable, UsesAnnot
self.copied_from_history_dataset_association = copied_from_history_dataset_association
self.copied_from_library_dataset_dataset_association = copied_from_library_dataset_dataset_association
def __create_version__(self, session):
state = inspect(self)
changes = {}
for attr in state.attrs:
hist = state.get_history(attr.key, True)
if not hist.has_changes():
continue
# hist.deleted holds old value(s)
changes[attr.key] = hist.deleted
if self.update_time and self.state == self.states.OK:
# We only record changes to HDAs that exist in the database and have a update_time
new_values = {}
new_values['name'] = changes.get('name', self.name)
new_values['dbkey'] = changes.get('dbkey', self.dbkey)
new_values['extension'] = changes.get('extension', self.extension)
new_values['extended_metadata_id'] = changes.get('extended_metadata_id', self.extended_metadata_id)
for k, v in new_values.items():
if isinstance(v, list):
new_values[k] = v[0]
new_values['update_time'] = self.update_time
new_values['version'] = self.version or 1
new_values['metadata'] = self._metadata
past_hda = HistoryDatasetAssociationHistory(history_dataset_association_id=self.id,
**new_values)
self.version = self.version + 1 if self.version else 1
session.add(past_hda)
def copy(self, parent_id=None):
"""
Create a copy of this HDA.
@@ -2609,6 +2667,27 @@ class HistoryDatasetAssociation(DatasetInstance, HasTags, Dictifiable, UsesAnnot
self.tags.append(new_tag_assoc)
class HistoryDatasetAssociationHistory(object):
def __init__(self,
history_dataset_association_id,
name,
dbkey,
update_time,
version,
extension,
extended_metadata_id,
metadata,
):
self.history_dataset_association_id = history_dataset_association_id
self.name = name
self.dbkey = dbkey
self.update_time = update_time
self.version = version
self.extension = extension
self.extended_metadata_id = extended_metadata_id
self._metadata = metadata
class HistoryDatasetAssociationDisplayAtAuthorization(object):
def __init__(self, hda=None, user=None, site=None):
self.history_dataset_association = hda
+19 -1
View File
@@ -7,6 +7,7 @@ from inspect import (
isclass
)
from sqlalchemy import event
from sqlalchemy.orm import (
scoped_session,
sessionmaker
@@ -20,7 +21,9 @@ class ModelMapping(Bunch):
def __init__(self, model_modules, engine):
self.engine = engine
context = scoped_session(sessionmaker(autoflush=False, autocommit=True))
Session = sessionmaker(autoflush=False, autocommit=True)
versioned_session(Session)
context = scoped_session(Session)
# For backward compatibility with "context.current"
# deprecated?
context.current = context
@@ -44,3 +47,18 @@ class ModelMapping(Bunch):
For backward compat., deprecated.
"""
return self.context
def versioned_objects(iter):
for obj in iter:
if hasattr(obj, '__create_version__'):
yield obj
def versioned_session(session):
@event.listens_for(session, 'before_flush')
def before_flush(session, flush_context, instances):
for obj in versioned_objects(session.dirty):
obj.__create_version__(session)
for obj in versioned_objects(session.deleted):
obj.__create_version__(session, deleted=True)
+35 -1
View File
@@ -29,8 +29,9 @@ from sqlalchemy import (
)
from sqlalchemy.ext.associationproxy import association_proxy
from sqlalchemy.ext.orderinglist import ordering_list
from sqlalchemy.orm import backref, class_mapper, deferred, mapper, object_session, relation
from sqlalchemy.orm import backref, class_mapper, column_property, deferred, mapper, object_session, relation
from sqlalchemy.orm.collections import attribute_mapped_collection
from sqlalchemy.sql import exists
from sqlalchemy.types import BigInteger
from galaxy import model
@@ -142,11 +143,26 @@ model.HistoryDatasetAssociation.table = Table(
Column("deleted", Boolean, index=True, default=False),
Column("visible", Boolean),
Column("extended_metadata_id", Integer, ForeignKey("extended_metadata.id"), index=True),
Column("version", Integer, default=1, nullable=True, index=True),
Column("hid", Integer),
Column("purged", Boolean, index=True, default=False),
Column("hidden_beneath_collection_instance_id",
ForeignKey("history_dataset_collection_association.id"), nullable=True))
model.HistoryDatasetAssociationHistory.table = Table(
"history_dataset_association_history", metadata,
Column("id", Integer, primary_key=True),
Column("history_dataset_association_id", Integer, ForeignKey("history_dataset_association.id"), index=True),
Column("update_time", DateTime, default=now),
Column("version", Integer),
Column("name", TrimmedString(255)),
Column("extension", TrimmedString(64)),
Column("metadata", MetadataType(), key="_metadata"),
Column("extended_metadata_id", Integer, ForeignKey("extended_metadata.id"), index=True),
)
model.Dataset.table = Table(
"dataset", metadata,
Column("id", Integer, primary_key=True),
@@ -464,6 +480,7 @@ model.Job.table = Table(
Column("tool_version", TEXT, default="1.0.0"),
Column("state", String(64), index=True),
Column("info", TrimmedString(255)),
Column("copied_from_job_id", Integer, nullable=True),
Column("command_line", TEXT),
Column("dependencies", JSONType, nullable=True),
Column("param_filename", String(1024)),
@@ -504,6 +521,7 @@ model.JobToInputDatasetAssociation.table = Table(
Column("id", Integer, primary_key=True),
Column("job_id", Integer, ForeignKey("job.id"), index=True),
Column("dataset_id", Integer, ForeignKey("history_dataset_association.id"), index=True),
Column("dataset_version", Integer),
Column("name", String(255)))
model.JobToOutputDatasetAssociation.table = Table(
@@ -1474,6 +1492,8 @@ simple_mapping(model.Dataset,
backref='datasets')
)
mapper(model.HistoryDatasetAssociationHistory, model.HistoryDatasetAssociationHistory.table)
mapper(model.HistoryDatasetAssociationDisplayAtAuthorization, model.HistoryDatasetAssociationDisplayAtAuthorization.table, properties=dict(
history_dataset_association=relation(model.HistoryDatasetAssociation),
user=relation(model.User)
@@ -1965,6 +1985,20 @@ mapper(model.Job, model.Job.table, properties=dict(
input_datasets=relation(model.JobToInputDatasetAssociation),
input_dataset_collections=relation(model.JobToInputDatasetCollectionAssociation, lazy=True),
output_datasets=relation(model.JobToOutputDatasetAssociation, lazy=True),
any_output_dataset_deleted=column_property(
exists([model.HistoryDatasetAssociation],
and_(model.Job.table.c.id == model.JobToOutputDatasetAssociation.table.c.job_id,
model.HistoryDatasetAssociation.table.c.id == model.JobToOutputDatasetAssociation.table.c.dataset_id,
model.HistoryDatasetAssociation.table.c.deleted == true())
)
),
any_output_dataset_collection_instances_deleted=column_property(
exists([model.HistoryDatasetCollectionAssociation.table.c.id],
and_(model.Job.table.c.id == model.JobToOutputDatasetCollectionAssociation.table.c.job_id,
model.HistoryDatasetCollectionAssociation.table.c.id == model.JobToOutputDatasetCollectionAssociation.table.c.dataset_collection_id,
model.HistoryDatasetCollectionAssociation.table.c.deleted == true())
)
),
output_dataset_collection_instances=relation(model.JobToOutputDatasetCollectionAssociation, lazy=True),
output_dataset_collections=relation(model.JobToImplicitOutputDatasetCollectionAssociation, lazy=True),
post_job_actions=relation(model.PostJobActionAssociation, lazy=False),
@@ -0,0 +1,40 @@
"""
Add copied_from_job_id column to jobs table
"""
from __future__ import print_function
import logging
from sqlalchemy import Column, Integer, MetaData, Table
log = logging.getLogger(__name__)
copied_from_job_id_column = Column("copied_from_job_id", Integer, nullable=True)
def upgrade(migrate_engine):
print(__doc__)
metadata = MetaData()
metadata.bind = migrate_engine
metadata.reflect()
# Add the copied_from_job_id column to the job table
try:
jobs_table = Table("job", metadata, autoload=True)
copied_from_job_id_column.create(jobs_table)
assert copied_from_job_id_column is jobs_table.c.copied_from_job_id
except Exception:
log.exception("Adding column 'copied_from_job_id_column' to job table failed.")
def downgrade(migrate_engine):
metadata = MetaData()
metadata.bind = migrate_engine
metadata.reflect()
# Drop the job table's copied_from_job_id column.
try:
jobs_table = Table("job", metadata, autoload=True)
copied_from_job_id = jobs_table.c.copied_from_job_id_column
copied_from_job_id.drop()
except Exception:
log.exception("Dropping 'copied_from_job_id_column' column from job table failed.")
@@ -0,0 +1,40 @@
"""
Add version column to history_dataset_association table
"""
from __future__ import print_function
import logging
from sqlalchemy import Column, Integer, MetaData, Table
log = logging.getLogger(__name__)
version_column = Column("version", Integer, default=1)
def upgrade(migrate_engine):
print(__doc__)
metadata = MetaData()
metadata.bind = migrate_engine
metadata.reflect()
# Add the version column to the history_dataset_association table
try:
hda_table = Table("history_dataset_association", metadata, autoload=True)
version_column.create(hda_table)
assert version_column is hda_table.c.version
except Exception:
log.exception("Adding column 'copied_from_job_id_column' to job table failed.")
def downgrade(migrate_engine):
metadata = MetaData()
metadata.bind = migrate_engine
metadata.reflect()
# Drop the history_dataset_association table's version column.
try:
hda_table = Table("history_dataset_association", metadata, autoload=True)
version_column = hda_table.c.version
version_column.drop()
except Exception:
log.exception("Dropping 'copied_from_job_id_column' column from job table failed.")
@@ -0,0 +1,49 @@
"""
Migration script to add the history_dataset_association_history table.
"""
from __future__ import print_function
import datetime
import logging
from sqlalchemy import Column, DateTime, ForeignKey, Integer, MetaData, Table
from galaxy.model.custom_types import MetadataType, TrimmedString
now = datetime.datetime.utcnow
log = logging.getLogger(__name__)
log.setLevel(logging.DEBUG)
metadata = MetaData()
HistoryDatasetAssociationHistory_table = Table(
"history_dataset_association_history", metadata,
Column("id", Integer, primary_key=True),
Column("history_dataset_association_id", Integer, ForeignKey("history_dataset_association.id"), index=True),
Column("update_time", DateTime, default=now),
Column("version", Integer, index=True),
Column("name", TrimmedString(255)),
Column("extension", TrimmedString(64)),
Column("metadata", MetadataType(), key='_metadata'),
Column("extended_metadata_id", Integer, ForeignKey("extended_metadata.id"), index=True),
)
def upgrade(migrate_engine):
print(__doc__)
metadata.bind = migrate_engine
metadata.reflect()
try:
HistoryDatasetAssociationHistory_table.create()
log.debug("Created history_dataset_association_history table")
except Exception:
log.exception("Creating history_dataset_association_history table failed.")
def downgrade(migrate_engine):
metadata.bind = migrate_engine
metadata.reflect()
try:
HistoryDatasetAssociationHistory_table.drop()
log.debug("Dropped history_dataset_association_history table")
except Exception:
log.exception("Dropping history_dataset_association_history table failed.")
@@ -0,0 +1,40 @@
"""
Add dataset_version column to job_to_input_dataset table
"""
from __future__ import print_function
import logging
from sqlalchemy import Column, Integer, MetaData, Table
log = logging.getLogger(__name__)
dataset_version_column = Column("dataset_version", Integer)
def upgrade(migrate_engine):
print(__doc__)
metadata = MetaData()
metadata.bind = migrate_engine
metadata.reflect()
# Add the version column to the job_to_input_dataset table
try:
job_to_input_dataset_table = Table("job_to_input_dataset", metadata, autoload=True)
dataset_version_column.create(job_to_input_dataset_table)
assert dataset_version_column is job_to_input_dataset_table.c.dataset_version
except Exception:
log.exception("Adding column 'dataset_history_id' to job_to_input_dataset table failed.")
def downgrade(migrate_engine):
metadata = MetaData()
metadata.bind = migrate_engine
metadata.reflect()
# Drop the job_to_input_dataset table's version column.
try:
job_to_input_dataset_table = Table("job_to_input_dataset", metadata, autoload=True)
dataset_version_column = job_to_input_dataset_table.c.dataset_version
dataset_version_column.drop()
except Exception:
log.exception("Dropping 'dataset_version' column from job_to_input_dataset table failed.")
+19 -3
View File
@@ -26,6 +26,7 @@ from galaxy import (
)
from galaxy.datatypes.metadata import JobExternalOutputMetadataWrapper
from galaxy.managers import histories
from galaxy.managers.jobs import JobSearch
from galaxy.queue_worker import send_control_task
from galaxy.tools.actions import DefaultToolAction
from galaxy.tools.actions.data_manager import DataManagerToolAction
@@ -438,6 +439,7 @@ class Tool(object, Dictifiable):
raise e
self.history_manager = histories.HistoryManager(app)
self._view = views.DependencyResolversView(app)
self.job_search = JobSearch(app=self.app)
@property
def version_object(self):
@@ -1279,7 +1281,7 @@ class Tool(object, Dictifiable):
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):
def handle_input(self, trans, incoming, history=None, use_cached_job=False):
"""
Process incoming parameters for this tool from the dict `incoming`,
update the tool state (or create if none existed), and either return
@@ -1294,7 +1296,20 @@ class Tool(object, Dictifiable):
raise exceptions.MessageException(', '.join(msg for msg in err_data.values()), err_data=err_data)
else:
mapping_params = MappingParameters(incoming, all_params)
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 = {}
for i, param in enumerate(all_params):
if use_cached_job:
completed_jobs[i] = self.job_search.by_tool_input(
trans=trans,
tool_id=self.id,
tool_version=self.version,
param=param,
param_dump=self.params_to_strings(param, self.app, nested=True),
job_state=None,
)
else:
completed_jobs[i] = None
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)
# Raise an exception if there were jobs to execute and none of them were submitted,
# if at least one is submitted or there are no jobs to execute - return aggregate
# information including per-job errors. Arguably we should just always return the
@@ -1311,7 +1326,7 @@ class Tool(object, Dictifiable):
output_collections=execution_tracker.output_collections,
implicit_collections=execution_tracker.implicit_collections)
def handle_single_execution(self, trans, rerun_remap_job_id, execution_slice, history, execution_cache=None):
def handle_single_execution(self, trans, rerun_remap_job_id, execution_slice, history, execution_cache=None, completed_job=None):
"""
Return a pair with whether execution is successful as well as either
resulting output data or an error message indicating the problem.
@@ -1324,6 +1339,7 @@ class Tool(object, Dictifiable):
rerun_remap_job_id=rerun_remap_job_id,
execution_cache=execution_cache,
dataset_collection_elements=execution_slice.dataset_collection_elements,
completed_job=completed_job,
)
except httpexceptions.HTTPFound as e:
# if it's a paste redirect exception, pass it up the stack
+35 -13
View File
@@ -195,7 +195,7 @@ class DefaultToolAction(object):
return history, inp_data, inp_dataset_collections
def execute(self, tool, trans, incoming={}, return_job=False, set_output_hid=True, history=None, job_params=None, rerun_remap_job_id=None, execution_cache=None, dataset_collection_elements=None):
def execute(self, tool, trans, incoming={}, return_job=False, set_output_hid=True, history=None, job_params=None, rerun_remap_job_id=None, execution_cache=None, dataset_collection_elements=None, completed_job=None):
"""
Executes a tool, creating job and tool outputs, associating them, and
submitting the job to the job queue. If history is not specified, use
@@ -223,7 +223,7 @@ class DefaultToolAction(object):
continue
# Convert LDDA to an HDA.
if isinstance(data, LibraryDatasetDatasetAssociation):
if isinstance(data, LibraryDatasetDatasetAssociation) and not completed_job:
data = data.to_history_dataset_association(None)
inp_data[name] = data
@@ -246,13 +246,14 @@ class DefaultToolAction(object):
inp_data.update({"chromInfo": db_dataset})
incoming["chromInfo"] = chrom_info
# Determine output dataset permission/roles list
existing_datasets = [inp for inp in inp_data.values() if inp]
if existing_datasets:
output_permissions = app.security_agent.guess_derived_permissions_for_datasets(existing_datasets)
else:
# No valid inputs, we will use history defaults
output_permissions = app.security_agent.history_get_default_permissions(history)
if not completed_job:
# Determine output dataset permission/roles list
existing_datasets = [inp for inp in inp_data.values() if inp]
if existing_datasets:
output_permissions = app.security_agent.guess_derived_permissions_for_datasets(existing_datasets)
else:
# No valid inputs, we will use history defaults
output_permissions = app.security_agent.history_get_default_permissions(history)
# Add the dbkey to the incoming parameters
incoming["dbkey"] = input_dbkey
@@ -301,7 +302,18 @@ class DefaultToolAction(object):
inp_dataset_collections,
input_ext
)
data = app.model.HistoryDatasetAssociation(extension=ext, create_dataset=True, flush=False)
create_datasets = True
dataset = None
if completed_job:
for output_dataset in completed_job.output_datasets:
if output_dataset.name == name:
create_datasets = False
completed_data = output_dataset.dataset
dataset = output_dataset.dataset.dataset
break
data = app.model.HistoryDatasetAssociation(extension=ext, dataset=dataset, create_dataset=create_datasets, flush=False)
if hidden is None:
hidden = output.hidden
if not hidden and dataset_collection_elements is not None: # Mapping over a collection - hide datasets
@@ -311,14 +323,17 @@ class DefaultToolAction(object):
if dataset_collection_elements is not None and name in dataset_collection_elements:
dataset_collection_elements[name].hda = data
trans.sa_session.add(data)
trans.app.security_agent.set_all_dataset_permissions(data.dataset, output_permissions, new=True)
if not completed_job:
trans.app.security_agent.set_all_dataset_permissions(data.dataset, output_permissions, new=True)
for _, tag in preserved_tags.items():
data.tags.append(tag.copy())
# Must flush before setting object store id currently.
# TODO: optimize this.
trans.sa_session.flush()
object_store_populator.set_object_store_id(data)
if not completed_job:
object_store_populator.set_object_store_id(data)
# This may not be neccesary with the new parent/child associations
data.designation = name
@@ -338,7 +353,12 @@ class DefaultToolAction(object):
# Take dbkey from LAST input
data.dbkey = str(input_dbkey)
# Set state
data.blurb = "queued"
if completed_job:
data.blurb = completed_data.blurb
data.peek = completed_data.peek
data._metadata = completed_data._metadata
else:
data.blurb = "queued"
# Set output label
data.name = self.get_output_name(output, data, tool, on_text, trans, incoming, history, wrapped_params.params, job_params)
# Store output
@@ -447,6 +467,8 @@ class DefaultToolAction(object):
if job_params:
job.params = dumps(job_params)
job.set_handler(tool.get_job_handler(job_params))
if completed_job:
job.set_copied_from_job_id(completed_job.id)
trans.sa_session.add(job)
# Now that we have a job id, we can remap any outputs if this is a rerun and the user chose to continue dependent jobs
# This functionality requires tracking jobs in the database.
+7 -8
View File
@@ -31,7 +31,7 @@ class PartialJobExecution(Exception):
MappingParameters = collections.namedtuple("MappingParameters", ["param_template", "param_combinations"])
def execute(trans, tool, mapping_params, history, rerun_remap_job_id=None, collection_info=None, workflow_invocation_uuid=None, invocation_step=None, max_num_jobs=None, job_callback=None):
def execute(trans, tool, mapping_params, history, rerun_remap_job_id=None, collection_info=None, workflow_invocation_uuid=None, invocation_step=None, max_num_jobs=None, job_callback=None, completed_jobs=None):
"""
Execute a tool and return object containing summary (output data, number of
failures, etc...).
@@ -49,7 +49,7 @@ def execute(trans, tool, mapping_params, history, rerun_remap_job_id=None, colle
app = trans.app
execution_cache = ToolExecutionCache(trans)
def execute_single_job(execution_slice):
def execute_single_job(execution_slice, completed_job):
job_timer = ExecutionTimer()
params = execution_slice.param_combination
if workflow_invocation_uuid:
@@ -59,7 +59,7 @@ def execute(trans, tool, mapping_params, history, rerun_remap_job_id=None, colle
# values or rerun parameters.
del params['__workflow_invocation_uuid__']
job, result = tool.handle_single_execution(trans, rerun_remap_job_id, execution_slice, history, execution_cache)
job, result = tool.handle_single_execution(trans, rerun_remap_job_id, execution_slice, history, execution_cache, completed_job)
if job:
message = EXECUTION_SUCCESS_MESSAGE % (tool.id, job.id, job_timer)
log.debug(message)
@@ -89,13 +89,12 @@ def execute(trans, tool, mapping_params, history, rerun_remap_job_id=None, colle
has_remaining_jobs = False
if (job_count < burst_at or burst_threads < 2):
for execution_slice in execution_tracker.new_execution_slices():
for i, execution_slice in enumerate(execution_tracker.new_execution_slices()):
if max_num_jobs and jobs_executed >= max_num_jobs:
has_remaining_jobs = True
break
else:
execute_single_job(execution_slice)
jobs_executed += 1
execute_single_job(execution_slice, completed_jobs[i])
else:
# TODO: re-record success...
q = Queue()
@@ -111,12 +110,12 @@ def execute(trans, tool, mapping_params, history, rerun_remap_job_id=None, colle
t.daemon = True
t.start()
for execution_slice in execution_tracker.new_execution_slices():
for i, execution_slice in enumerate(execution_tracker.new_execution_slices()):
if max_num_jobs and jobs_executed >= max_num_jobs:
has_remaining_jobs = True
break
else:
q.put(execution_slice)
q.put(execution_slice, completed_jobs[i])
jobs_executed += 1
q.join()
+3
View File
@@ -3,6 +3,7 @@ Basic tool parameters.
"""
from __future__ import print_function
import cgi
import logging
import os
import os.path
@@ -543,6 +544,8 @@ class FileToolParameter(ToolParameter):
return value['local_filename']
except KeyError:
return None
elif isinstance(value, cgi.FieldStorage):
return value.filename
raise Exception("FileToolParameter cannot be persisted")
def to_python(self, value, app):
+3 -1
View File
@@ -302,9 +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:
+6 -1
View File
@@ -346,7 +346,12 @@ class ToolsController(BaseAPIController, UsesVisualizationMixin):
# TODO: handle dbkeys
params = util.Params(inputs, sanitize=False)
incoming = params.__dict__
vars = tool.handle_input(trans, incoming, history=target_history)
# use_cached_job can be passed in via the top-level payload or among the tool inputs.
# I think it should be a top-level parameter, but because the selector is implemented
# as a regular tool parameter we accept both.
use_cached_job = payload.get('use_cached_job', False) or util.string_as_bool(inputs.get('use_cached_job', 'false'))
vars = tool.handle_input(trans, incoming, history=target_history, use_cached_job=use_cached_job)
# TODO: check for errors and ensure that output dataset(s) are available.
output_datasets = vars.get('out_data', [])
+12 -12
View File
@@ -296,21 +296,21 @@ class UserAPIController(BaseAPIController, UsesTagsMixin, CreatesUsersMixin, Cre
"""
if not preferences:
return []
data = []
# Get data if present
data_key = "extra_user_preferences"
if data_key in user.preferences:
data = json.loads(user.preferences[data_key])
extra_pref_inputs = list()
# Build sections for different categories of inputs
for item, value in preferences.items():
if value is not None:
for input in value["inputs"]:
input['help'] = 'Required' if input['required'] else ''
help = input.get('help', '')
required = 'Required' if util.string_as_bool(input.get('required')) else ''
if help:
input['help'] = "%s %s" % (help, required)
else:
input['help'] = required
field = item + '|' + input['name']
for data_item in data:
for data_item in user.extra_preferences:
if field in data_item:
input['value'] = data[data_item]
input['value'] = user.extra_preferences[data_item]
extra_pref_inputs.append({'type': 'section', 'title': value['description'], 'name': item, 'expanded': True, 'inputs': value['inputs']})
return extra_pref_inputs
@@ -455,9 +455,9 @@ class UserAPIController(BaseAPIController, UsesTagsMixin, CreatesUsersMixin, Cre
# Update values for extra user preference items
extra_user_pref_data = dict()
get_extra_pref_keys = self._get_extra_user_preferences(trans)
if get_extra_pref_keys is not None:
for key in get_extra_pref_keys:
extra_pref_keys = self._get_extra_user_preferences(trans)
if extra_pref_keys is not None:
for key in extra_pref_keys:
key_prefix = key + '|'
for item in payload:
if item.startswith(key_prefix):
@@ -465,7 +465,7 @@ class UserAPIController(BaseAPIController, UsesTagsMixin, CreatesUsersMixin, Cre
if payload[item] == "":
# Raise an exception when a required field is empty while saving the form
keys = item.split("|")
section = get_extra_pref_keys[keys[0]]
section = extra_pref_keys[keys[0]]
for input in section['inputs']:
if input['name'] == keys[1] and input['required']:
raise MessageException("Please fill the required field")
+9 -4
View File
@@ -255,6 +255,9 @@ class WorkflowsAPIController(BaseAPIController, UsesStoredWorkflowMixin, UsesAnn
:param allow_tool_state_corrections: If set to True, any Tool parameter changes will not prevent running workflow, defaults to False
:type allow_tool_state_corrections: bool
:param use_cached_job: If set to True galaxy will attempt to find previously executed steps for all workflow steps with the exact same parameter combinations
and will copy the outputs of the previously executed step.
"""
ways_to_create = set([
'workflow_id',
@@ -336,10 +339,12 @@ class WorkflowsAPIController(BaseAPIController, UsesStoredWorkflowMixin, UsesAnn
rval = {}
rval['history'] = trans.security.encode_id(history.id)
rval['outputs'] = []
for step in workflow.steps:
if step.type == 'tool' or step.type is None:
for v in outputs[step.id].values():
rval['outputs'].append(trans.security.encode_id(v.id))
if outputs:
# Newer outputs don't necessarily fill outputs (?)
for step in workflow.steps:
if step.type == 'tool' or step.type is None:
for v in outputs[step.id].values():
rval['outputs'].append(trans.security.encode_id(v.id))
# Newer version of this API just returns the invocation as a dict, to
# facilitate migration - produce the newer style response and blend in
+21 -7
View File
@@ -214,7 +214,7 @@ class WorkflowModule(object):
state.decode(runtime_state, Bunch(inputs=self.get_runtime_inputs()), self.trans.app)
return state
def execute(self, trans, progress, invocation_step):
def execute(self, trans, progress, invocation_step, use_cached_job=False):
""" Execute the given workflow invocation step.
Use the supplied workflow progress object to track outputs, find
@@ -333,13 +333,13 @@ class SubWorkflowModule(WorkflowModule):
def get_content_id(self):
return self.trans.security.encode_id(self.subworkflow.id)
def execute(self, trans, progress, invocation_step):
def execute(self, trans, progress, invocation_step, use_cached_job=False):
""" Execute the given workflow step in the given workflow invocation.
Use the supplied workflow progress object to track outputs, find
inputs, etc...
"""
step = invocation_step.workflow_step
subworkflow_invoker = progress.subworkflow_invoker(trans, step)
subworkflow_invoker = progress.subworkflow_invoker(trans, step, use_cached_job=use_cached_job)
subworkflow_invoker.invoke()
subworkflow = subworkflow_invoker.workflow
subworkflow_progress = subworkflow_invoker.progress
@@ -367,7 +367,7 @@ class InputModule(WorkflowModule):
def get_data_inputs(self):
return []
def execute(self, trans, progress, invocation_step):
def execute(self, trans, progress, invocation_step, use_cached_job=False):
invocation = invocation_step.workflow_invocation
step = invocation_step.workflow_step
step_outputs = dict(output=step.state.inputs['input'])
@@ -504,7 +504,7 @@ class InputParameterModule(WorkflowModule):
def get_data_inputs(self):
return []
def execute(self, trans, progress, invocation_step):
def execute(self, trans, progress, invocation_step, use_cached_job=False):
step = invocation_step.workflow_step
step_outputs = dict(output=step.state.inputs['input'])
progress.set_outputs_for_input(invocation_step, step_outputs)
@@ -535,7 +535,7 @@ class PauseModule(WorkflowModule):
state.inputs = dict()
return state
def execute(self, trans, progress, invocation_step):
def execute(self, trans, progress, invocation_step, use_cached_job=False):
step = invocation_step.workflow_step
progress.mark_step_outputs_delayed(step, why="executing pause step")
@@ -805,7 +805,7 @@ class ToolModule(WorkflowModule):
else:
raise ToolMissingException("Tool %s missing. Cannot recover runtime state." % self.tool_id)
def execute(self, trans, progress, invocation_step):
def execute(self, trans, progress, invocation_step, use_cached_job=False):
invocation = invocation_step.workflow_invocation
step = invocation_step.workflow_step
tool = trans.app.toolbox.get_tool(step.tool_id, tool_version=step.tool_version)
@@ -873,6 +873,19 @@ class ToolModule(WorkflowModule):
param_combinations.append(execution_state.inputs)
complete = False
completed_jobs = {}
for i, param in enumerate(param_combinations):
if use_cached_job:
completed_jobs[i] = tool.job_search.by_tool_input(
trans=trans,
tool_id=tool.id,
tool_version=tool.version,
param=param,
param_dump=tool.params_to_strings(param, trans.app, nested=True),
job_state=None,
)
else:
completed_jobs[i] = None
try:
mapping_params = MappingParameters(tool_state.inputs, param_combinations)
max_num_jobs = progress.maximum_jobs_to_schedule_or_none
@@ -886,6 +899,7 @@ class ToolModule(WorkflowModule):
invocation_step=invocation_step,
max_num_jobs=max_num_jobs,
job_callback=lambda job: self._handle_post_job_actions(step, job, invocation.replacement_dict),
completed_jobs=completed_jobs
)
complete = True
except PartialJobExecution as pje:
+7 -2
View File
@@ -136,6 +136,7 @@ class WorkflowInvoker(object):
self.workflow_invocation = workflow_invocation
self.workflow_invocation.copy_inputs_to_history = workflow_run_config.copy_inputs_to_history
self.workflow_invocation.use_cached_job = workflow_run_config.use_cached_job
self.workflow_invocation.replacement_dict = workflow_run_config.replacement_dict
module_injector = modules.WorkflowModuleInjector(trans)
@@ -255,7 +256,10 @@ class WorkflowInvoker(object):
pass
def _invoke_step(self, invocation_step):
incomplete_or_none = invocation_step.workflow_step.module.execute(self.trans, self.progress, invocation_step)
incomplete_or_none = invocation_step.workflow_step.module.execute(self.trans,
self.progress,
invocation_step,
use_cached_job=self.workflow_invocation.use_cached_job)
return incomplete_or_none
@@ -437,7 +441,7 @@ class WorkflowProgress(object):
raise Exception("Failed to find persisted workflow invocation for step [%s]" % step.id)
return subworkflow_invocation
def subworkflow_invoker(self, trans, step):
def subworkflow_invoker(self, trans, step, use_cached_job=False):
subworkflow_progress = self.subworkflow_progress(step)
subworkflow_invocation = subworkflow_progress.workflow_invocation
workflow_run_config = WorkflowRunConfig(
@@ -446,6 +450,7 @@ class WorkflowProgress(object):
inputs={},
param_map={},
copy_inputs_to_history=False,
use_cached_job=use_cached_job
)
return WorkflowInvoker(
trans,
+17 -2
View File
@@ -41,13 +41,20 @@ class WorkflowRunConfig(object):
:type param_map: dict
"""
def __init__(self, target_history, replacement_dict, copy_inputs_to_history=False, inputs={}, param_map={}, allow_tool_state_corrections=False):
def __init__(self, target_history,
replacement_dict,
copy_inputs_to_history=False,
inputs={},
param_map={},
allow_tool_state_corrections=False,
use_cached_job=False):
self.target_history = target_history
self.replacement_dict = replacement_dict
self.copy_inputs_to_history = copy_inputs_to_history
self.inputs = inputs
self.param_map = param_map
self.allow_tool_state_corrections = allow_tool_state_corrections
self.use_cached_job = use_cached_job
def _normalize_inputs(steps, inputs, inputs_by):
@@ -197,6 +204,7 @@ def _get_target_history(trans, workflow, payload, param_keys=[], index=0):
def build_workflow_run_configs(trans, workflow, payload):
app = trans.app
allow_tool_state_corrections = payload.get('allow_tool_state_corrections', False)
use_cached_job = payload.get('use_cached_job', False)
# Sanity checks.
if len(workflow.steps) == 0:
@@ -301,7 +309,8 @@ def build_workflow_run_configs(trans, workflow, payload):
replacement_dict=payload.get('replacement_params', {}),
inputs=normalized_inputs,
param_map=param_map,
allow_tool_state_corrections=allow_tool_state_corrections
allow_tool_state_corrections=allow_tool_state_corrections,
use_cached_job=use_cached_job,
))
return run_configs
@@ -338,6 +347,7 @@ def workflow_run_config_to_request(trans, run_config, workflow):
target_history=run_config.target_history,
replacement_dict=run_config.replacement_dict,
copy_inputs_to_history=False,
use_cached_job=run_config.use_cached_job,
inputs={},
param_map={},
allow_tool_state_corrections=run_config.allow_tool_state_corrections
@@ -363,6 +373,7 @@ def workflow_run_config_to_request(trans, run_config, workflow):
workflow_invocation.add_input(content, step_id)
add_parameter("copy_inputs_to_history", "true" if run_config.copy_inputs_to_history else "false", param_types.META_PARAMETERS)
add_parameter("use_cached_job", "true" if run_config.use_cached_job else "false", param_types.META_PARAMETERS)
return workflow_invocation
@@ -373,6 +384,7 @@ def workflow_request_to_run_config(work_request_context, workflow_invocation):
inputs = {}
param_map = {}
copy_inputs_to_history = None
use_cached_job = False
for parameter in workflow_invocation.input_parameters:
parameter_type = parameter.type
@@ -381,6 +393,8 @@ def workflow_request_to_run_config(work_request_context, workflow_invocation):
elif parameter_type == param_types.META_PARAMETERS:
if parameter.name == "copy_inputs_to_history":
copy_inputs_to_history = (parameter.value == "true")
if parameter.name == 'use_cached_job':
use_cached_job = (parameter.value == 'true')
for input_association in workflow_invocation.input_datasets:
inputs[input_association.workflow_step_id] = input_association.dataset
for input_association in workflow_invocation.input_dataset_collections:
@@ -395,6 +409,7 @@ def workflow_request_to_run_config(work_request_context, workflow_invocation):
inputs=inputs,
param_map=param_map,
copy_inputs_to_history=copy_inputs_to_history,
use_cached_job=use_cached_job,
)
return workflow_run_config
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
+21 -5
View File
@@ -274,15 +274,16 @@ class JobsApiTestCase(api.ApiTestCase):
def test_search(self):
history_id, dataset_id = self.__history_with_ok_dataset()
# We first copy the datasets, so that the update time is lower than the job creation time
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)
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}
@@ -298,6 +299,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({
@@ -411,7 +427,7 @@ class JobsApiTestCase(api.ApiTestCase):
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(15):
for i in range(5):
search_count = self._search_count(payload)
if search_count == expected_search_count:
break
+26 -3
View File
@@ -303,6 +303,25 @@ class ToolsTestCase(api.ApiTestCase):
output1_content = self.dataset_populator.get_history_dataset_content(history_id, dataset=output1)
self.assertEqual(output1_content.strip(), "Cat1Test")
@skip_without_tool("cat1")
def test_run_cat1_use_cached_job(self):
with self.dataset_populator.test_history() as history_id:
# Run simple non-upload tool with an input data parameter.
new_dataset = self.dataset_populator.new_dataset(history_id, content='Cat1Test')
inputs = dict(
input1=dataset_to_param(new_dataset),
)
outputs_one = self._run_cat1(history_id, inputs=inputs, assert_ok=True, wait_for_job=True)
outputs_two = self._run_cat1(history_id, inputs=inputs, use_cached_job=False, assert_ok=True, wait_for_job=True)
outputs_three = self._run_cat1(history_id, inputs=inputs, use_cached_job=True, assert_ok=True, wait_for_job=True)
dataset_details = []
for output in [outputs_one, outputs_two, outputs_three]:
output_id = output['outputs'][0]['id']
dataset_details.append(self._get("datasets/%s" % output_id).json())
filenames = [dd['file_name'] for dd in dataset_details]
assert len(filenames) == 3, filenames
assert len(set(filenames)) <= 2, filenames
@skip_without_tool("cat1")
def test_run_cat1_listified_param(self):
# Run simple non-upload tool with an input data parameter.
@@ -1384,10 +1403,10 @@ class ToolsTestCase(api.ApiTestCase):
self._assert_status_code_is(create_response, 200)
return create_response.json()['outputs']
def _run_cat1(self, history_id, inputs, assert_ok=False):
return self._run('cat1', history_id, inputs, assert_ok=assert_ok)
def _run_cat1(self, history_id, inputs, assert_ok=False, **kwargs):
return self._run('cat1', history_id, inputs, assert_ok=assert_ok, **kwargs)
def _run(self, tool_id, history_id, inputs, assert_ok=False, tool_version=None):
def _run(self, tool_id, history_id, inputs, assert_ok=False, tool_version=None, use_cached_job=False, wait_for_job=False):
payload = self.dataset_populator.run_tool_payload(
tool_id=tool_id,
inputs=inputs,
@@ -1395,7 +1414,11 @@ class ToolsTestCase(api.ApiTestCase):
)
if tool_version is not None:
payload["tool_version"] = tool_version
if use_cached_job:
payload['use_cached_job'] = True
create_response = self._post("tools", data=payload)
if wait_for_job:
self.dataset_populator.wait_for_job(job_id=create_response.json()['jobs'][0]['id'])
if assert_ok:
self._assert_status_code_is(create_response, 200)
create = create_response.json()
+69 -1
View File
@@ -1,5 +1,6 @@
from __future__ import print_function
import json
import time
from collections import namedtuple
from json import dumps
@@ -254,6 +255,8 @@ class BaseWorkflowsApiTestCase(api.ApiTestCase):
invocation_id=invocation_id,
inputs=inputs,
jobs=jobs,
invocation=invocation,
workflow_request=workflow_request
)
def _history_jobs(self, history_id):
@@ -1817,6 +1820,71 @@ test_data:
self.dataset_populator.wait_for_history_jobs(history_id, assert_ok=assert_ok)
time.sleep(.5)
@skip_without_tool('cat1')
def test_workflow_rerun_with_use_cached_job(self):
workflow = self.workflow_populator.load_workflow(name="test_for_run")
# We launch a workflow
with self.dataset_populator.test_history() as history_id:
workflow_request, _ = self._setup_workflow_run(workflow, history_id=history_id)
run_workflow_response = self._post("workflows", data=workflow_request).json()
# We copy the workflow inputs to a new history
new_workflow_request = workflow_request.copy()
new_ds_map = json.loads(new_workflow_request['ds_map'])
with self.dataset_populator.test_history() as new_history_id:
for key, input_values in run_workflow_response['inputs'].items():
copy_payload = {"content": input_values['id'], "source": "hda", "type": "dataset"}
copy_response = self._post("histories/%s/contents" % new_history_id, data=copy_payload).json()
new_ds_map[key]['id'] = copy_response['id']
new_workflow_request['ds_map'] = json.dumps(new_ds_map)
new_workflow_request['history'] = "hist_id=%s" % new_history_id
new_workflow_request['use_cached_job'] = True
# We run the workflow again, it should not produce any new outputs
new_workflow_response = self._post("workflows", data=new_workflow_request).json()
first_wf_output = self._get("datasets/%s" % run_workflow_response['outputs'][0]).json()
second_wf_output = self._get("datasets/%s" % new_workflow_response['outputs'][0]).json()
assert first_wf_output['file_name'] == second_wf_output['file_name'], \
"first output :\n%s\nsecond output: %s" % (first_wf_output, second_wf_output)
@skip_without_tool('cat1')
def test_nested_workflow_rerun_with_use_cached_job(self):
with self.dataset_populator.test_history() as history_id_one, self.dataset_populator.test_history() as history_id_two:
workflow_run_description = """%s
test_data:
outer_input:
value: 1.bed
type: File
""" % SIMPLE_NESTED_WORKFLOW_YAML
run_jobs_summary = self._run_jobs(workflow_run_description, history_id=history_id_one)
self.dataset_populator.wait_for_history(history_id_one, assert_ok=True)
workflow_request = run_jobs_summary.workflow_request
# We copy the inputs to a new history and re-reun the workflow
inputs = json.loads(workflow_request['inputs'])
dataset_type = inputs['outer_input']['src']
dataset_id = inputs['outer_input']['id']
copy_payload = {"content": dataset_id, "source": dataset_type, "type": "dataset"}
copy_response = self._post("histories/%s/contents" % history_id_two, data=copy_payload)
self._assert_status_code_is(copy_response, 200)
new_dataset_id = copy_response.json()['id']
inputs['outer_input']['id'] = new_dataset_id
workflow_request['use_cached_job'] = True
workflow_request['history'] = "hist_id={history_id_two}".format(history_id_two=history_id_two)
workflow_request['inputs'] = json.dumps(inputs)
run_workflow_response = self._post("workflows", data=run_jobs_summary.workflow_request).json()
self.workflow_populator.wait_for_workflow(workflow_request['workflow_id'],
run_workflow_response['id'],
history_id_two,
assert_ok=True)
# Now make sure that the HDAs in each history point to the same dataset instances
history_one_contents = self.__history_contents(history_id_one)
history_two_contents = self.__history_contents(history_id_two)
assert len(history_one_contents) == len(history_two_contents)
for i, (item_one, item_two) in enumerate(zip(history_one_contents, history_two_contents)):
assert item_one['dataset_id'] == item_two['dataset_id'], \
'Dataset ids should match, but "%s" and "%s" are not the same for History item %i.' % (item_one['dataset_id'],
item_two['dataset_id'],
i + 1)
def test_cannot_run_inaccessible_workflow(self):
workflow = self.workflow_populator.load_workflow(name="test_for_run_cannot_access")
workflow_request, history_id = self._setup_workflow_run(workflow)
@@ -2702,4 +2770,4 @@ steps:
)
RunJobsSummary = namedtuple('RunJobsSummary', ['history_id', 'workflow_id', 'invocation_id', 'inputs', 'jobs'])
RunJobsSummary = namedtuple('RunJobsSummary', ['history_id', 'workflow_id', 'invocation_id', 'inputs', 'jobs', 'invocation', 'workflow_request'])
+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)
+2 -2
View File
@@ -78,10 +78,10 @@
"post_job_actions": {},
"tool_errors": null,
"tool_id": "cat1",
"tool_state": "{\"__page__\": 0, \"__rerun_remap_job_id__\": null, \"input1\": \"null\", \"chromInfo\": \"\\\"/home/john/workspace/galaxy-central/tool-data/shared/ucsc/chrom/?.len\\\"\", \"queries\": \"[{\\\"input2\\\": null, \\\"__index__\\\": 0}]\"}",
"tool_state": "{\"__page__\": 0, \"__rerun_remap_job_id__\": null, \"input1\": \"null\", \"queries\": \"[{\\\"input2\\\": null, \\\"__index__\\\": 0}]\"}",
"tool_version": "1.0.0",
"type": "tool",
"user_outputs": []
}
}
}
}
+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: