mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Merge branch 'release_19.05' into dev
This commit is contained in:
@@ -42,7 +42,7 @@ const applyRegex = function(regex, target, data, replacement, groupCount) {
|
||||
return null;
|
||||
}
|
||||
if (!replacement) {
|
||||
groupCount = groupCount && parseInt(groupCount);
|
||||
groupCount = groupCount && parseInt(groupCount, 10);
|
||||
if (groupCount) {
|
||||
if (match.length != groupCount + 1) {
|
||||
failedCount++;
|
||||
@@ -108,7 +108,7 @@ const RULES = {
|
||||
}
|
||||
},
|
||||
save: (component, rule) => {
|
||||
rule.start = parseInt(component.addColumnRownumStart);
|
||||
rule.start = parseInt(component.addColumnRownumStart, 10);
|
||||
},
|
||||
apply: (rule, data, sources, columns) => {
|
||||
let rownum = rule.start;
|
||||
@@ -169,7 +169,7 @@ const RULES = {
|
||||
const ruleValue = rule.value;
|
||||
let newRow;
|
||||
if (ruleValue.indexOf("identifier") == 0) {
|
||||
const identifierIndex = parseInt(ruleValue.substring("identifier".length));
|
||||
const identifierIndex = parseInt(ruleValue.substring("identifier".length), 10);
|
||||
newRow = (row, index) => {
|
||||
const newRow = row.slice();
|
||||
newRow.push(sources[index]["identifiers"][identifierIndex]);
|
||||
@@ -253,7 +253,7 @@ const RULES = {
|
||||
component.addColumnRegexTarget = rule.target_column;
|
||||
component.addColumnRegexExpression = rule.expression;
|
||||
component.addColumnRegexReplacement = rule.replacement;
|
||||
component.addColumnRegexGroupCount = rule.group_count;
|
||||
component.addColumnRegexGroupCount = parseInt(rule.group_count);
|
||||
}
|
||||
let addColumnRegexType = "global";
|
||||
if (component.addColumnRegexGroupCount) {
|
||||
@@ -347,7 +347,7 @@ const RULES = {
|
||||
},
|
||||
save: (component, rule) => {
|
||||
rule.target_column = component.addColumnSubstrTarget;
|
||||
rule.length = parseInt(component.addColumnSubstrLength);
|
||||
rule.length = parseInt(component.addColumnSubstrLength, 10);
|
||||
rule.substr_type = component.addColumnSubstrType;
|
||||
},
|
||||
apply: (rule, data, sources, columns) => {
|
||||
@@ -403,7 +403,7 @@ const RULES = {
|
||||
function newRow(row) {
|
||||
const newRow = [];
|
||||
for (const index in row) {
|
||||
if (targets.indexOf(parseInt(index)) == -1) {
|
||||
if (targets.indexOf(parseInt(index, 10)) == -1) {
|
||||
newRow.push(row[index]);
|
||||
}
|
||||
}
|
||||
@@ -448,7 +448,7 @@ const RULES = {
|
||||
const target = rule.target_column;
|
||||
const invert = rule.invert;
|
||||
const filterFunction = function(el, index) {
|
||||
const row = data[parseInt(index)];
|
||||
const row = data[parseInt(index, 10)];
|
||||
return regExp.exec(row[target]) ? !invert : invert;
|
||||
};
|
||||
sources = sources.filter(filterFunction);
|
||||
@@ -477,13 +477,13 @@ const RULES = {
|
||||
component.addFilterCountWhich = "first";
|
||||
component.addFilterCountInvert = false;
|
||||
} else {
|
||||
component.addFilterCountN = parseInt(rule.count);
|
||||
component.addFilterCountN = parseInt(rule.count, 10);
|
||||
component.addFilterCountWhich = rule.which;
|
||||
component.addFilterCountInvert = rule.inverse;
|
||||
}
|
||||
},
|
||||
save: (component, rule) => {
|
||||
rule.count = parseInt(component.addFilterCountN);
|
||||
rule.count = parseInt(component.addFilterCountN, 10);
|
||||
rule.which = component.addFilterCountWhich;
|
||||
rule.invert = component.addFilterCountInvert;
|
||||
},
|
||||
@@ -528,7 +528,7 @@ const RULES = {
|
||||
const target = rule.target_column;
|
||||
const invert = rule.invert;
|
||||
const filterFunction = function(el, index) {
|
||||
const row = data[parseInt(index)];
|
||||
const row = data[parseInt(index, 10)];
|
||||
return row[target].length ? !invert : invert;
|
||||
};
|
||||
sources = sources.filter(filterFunction);
|
||||
@@ -562,7 +562,7 @@ const RULES = {
|
||||
const invert = rule.invert;
|
||||
const value = rule.value;
|
||||
const filterFunction = function(el, index) {
|
||||
const row = data[parseInt(index)];
|
||||
const row = data[parseInt(index, 10)];
|
||||
return row[target] == value ? !invert : invert;
|
||||
};
|
||||
sources = sources.filter(filterFunction);
|
||||
@@ -598,7 +598,7 @@ const RULES = {
|
||||
const compare_type = rule.compare_type;
|
||||
const value = rule.value;
|
||||
const filterFunction = function(el, index) {
|
||||
const row = data[parseInt(index)];
|
||||
const row = data[parseInt(index, 10)];
|
||||
const targetValue = parseFloat(row[target]);
|
||||
let matches;
|
||||
if (compare_type == "less_than") {
|
||||
@@ -730,7 +730,7 @@ const RULES = {
|
||||
const newRow0 = [],
|
||||
newRow1 = [];
|
||||
for (let index in row) {
|
||||
index = parseInt(index);
|
||||
index = parseInt(index, 10);
|
||||
if (targets0.indexOf(index) > -1) {
|
||||
newRow0.push(row[index]);
|
||||
} else if (targets1.indexOf(index) > -1) {
|
||||
|
||||
@@ -886,7 +886,10 @@ class JobHandlerStopQueue(Monitors):
|
||||
.filter((model.Job.state == model.Job.states.DELETED_NEW) &
|
||||
(model.Job.handler == self.app.config.server_name)).all()
|
||||
for job in newly_deleted_jobs:
|
||||
jobs_to_check.append((job, job.stderr))
|
||||
# job.stderr is always a string (job.job_stderr + job.tool_stderr, possibly `''`),
|
||||
# while any `not None` message returned in self.queue.get_nowait() is interpreted
|
||||
# as an error, so here we use None if job.stderr is false-y
|
||||
jobs_to_check.append((job, job.stderr or None))
|
||||
# Also pull from the queue (in the case of Administrative stopped jobs)
|
||||
try:
|
||||
while 1:
|
||||
|
||||
@@ -4,6 +4,7 @@ $headers
|
||||
|
||||
_galaxy_setup_environment() {
|
||||
local _use_framework_galaxy="$1"
|
||||
_GALAXY_JOB_DIR="$working_directory"
|
||||
_GALAXY_JOB_HOME_DIR="$working_directory/home"
|
||||
_GALAXY_JOB_TMP_DIR=$tmp_dir_creation_statement
|
||||
$env_setup_commands
|
||||
|
||||
@@ -62,6 +62,7 @@ class JobManager(object):
|
||||
for data_assoc in job.output_datasets:
|
||||
if not self.dataset_manager.is_accessible(data_assoc.dataset.dataset, trans.user):
|
||||
raise ItemAccessibilityException("You are not allowed to rerun this job.")
|
||||
trans.sa_session.refresh(job)
|
||||
return job
|
||||
|
||||
|
||||
|
||||
@@ -322,27 +322,28 @@ class JobLike(object):
|
||||
|
||||
return "%s[%s,tool_id=%s]" % (self.__class__.__name__, extra, self.tool_id)
|
||||
|
||||
def get_stdout(self):
|
||||
stdout = self.tool_stdout
|
||||
@property
|
||||
def stdout(self):
|
||||
stdout = self.tool_stdout or ''
|
||||
if self.job_stdout:
|
||||
stdout += "\n" + self.job_stdout
|
||||
return stdout
|
||||
|
||||
def set_stdout(self, stdout):
|
||||
@stdout.setter
|
||||
def stdout(self, stdout):
|
||||
raise NotImplementedError("Attempt to set stdout, must set tool_stdout or job_stdout")
|
||||
|
||||
def get_stderr(self):
|
||||
stderr = self.tool_stderr
|
||||
@property
|
||||
def stderr(self):
|
||||
stderr = self.tool_stderr or ''
|
||||
if self.job_stderr:
|
||||
stderr += "\n" + self.job_stderr
|
||||
return stderr
|
||||
|
||||
def set_stderr(self, stderr):
|
||||
@stderr.setter
|
||||
def stderr(self, stderr):
|
||||
raise NotImplementedError("Attempt to set stdout, must set tool_stderr or job_stderr")
|
||||
|
||||
stdout = property(get_stdout, set_stdout)
|
||||
stderr = property(get_stderr, set_stderr)
|
||||
|
||||
|
||||
class User(Dictifiable, RepresentById):
|
||||
use_pbkdf2 = True
|
||||
|
||||
@@ -45,7 +45,10 @@ def upgrade(migrate_engine):
|
||||
|
||||
OldWorkflowStepConnection_table = Table("workflow_step_connection", metadata, autoload=True)
|
||||
for index in OldWorkflowStepConnection_table.indexes:
|
||||
index.drop()
|
||||
try:
|
||||
index.drop()
|
||||
except Exception:
|
||||
log.exception("Dropping index '%s' from table '%s' failed", index, OldWorkflowStepConnection_table)
|
||||
OldWorkflowStepConnection_table.rename("workflow_step_connection_preupgrade145")
|
||||
# Try to deregister that table to work around some caching problems it seems.
|
||||
OldWorkflowStepConnection_table.deregister()
|
||||
|
||||
@@ -521,7 +521,8 @@ class ToolEvaluator(object):
|
||||
os.close(fd)
|
||||
self.__write_workdir_file(config_filename, environment_variable_template, param_dict)
|
||||
config_file_basename = os.path.basename(config_filename)
|
||||
environment_variable["value"] = "`cat %s`" % config_file_basename
|
||||
# environment setup in job file template happens before `cd $working_directory`
|
||||
environment_variable["value"] = '`cat "$_GALAXY_JOB_DIR/%s"`' % config_file_basename
|
||||
environment_variable["raw"] = True
|
||||
environment_variables.append(environment_variable)
|
||||
|
||||
|
||||
@@ -244,7 +244,7 @@ class GalaxyInteractorApi(object):
|
||||
def wait_for_job(self, job_id, history_id, maxseconds):
|
||||
self.wait_for(lambda: not self.__job_ready(job_id, history_id), maxseconds=maxseconds)
|
||||
|
||||
def wait_for(self, func, **kwd):
|
||||
def wait_for(self, func, what='Tool test run', **kwd):
|
||||
sleep_amount = 0.2
|
||||
slept = 0
|
||||
walltime_exceeded = int(kwd.get("maxseconds", DEFAULT_TOOL_TEST_WAIT))
|
||||
@@ -258,7 +258,7 @@ class GalaxyInteractorApi(object):
|
||||
else:
|
||||
return
|
||||
|
||||
message = 'Tool test run exceeded walltime [total %s, max %s], terminating.' % (slept, walltime_exceeded)
|
||||
message = '%s exceeded walltime [total %s, max %s], terminating.' % (what, slept, walltime_exceeded)
|
||||
log.info(message)
|
||||
raise AssertionError(message)
|
||||
|
||||
|
||||
@@ -352,10 +352,16 @@ class WorkflowProgress(object):
|
||||
try:
|
||||
replacement = step_outputs[output_name]
|
||||
except KeyError:
|
||||
# Must resolve.
|
||||
template = "Workflow evaluation problem - failed to find output_name %s in step_outputs %s"
|
||||
message = template % (output_name, step_outputs)
|
||||
raise Exception(message)
|
||||
replacement = self.inputs_by_step_id.get(output_step_id)
|
||||
if connection.output_step.type == 'parameter_input' and output_step_id is not None:
|
||||
# FIXME: parameter_input step outputs should be properly recorded as step outputs, but for now we can
|
||||
# short-circuit and just pick the input value
|
||||
pass
|
||||
else:
|
||||
# Must resolve.
|
||||
template = "Workflow evaluation problem - failed to find output_name %s in step_outputs %s"
|
||||
message = template % (output_name, step_outputs)
|
||||
raise Exception(message)
|
||||
if isinstance(replacement, model.HistoryDatasetCollectionAssociation):
|
||||
if not replacement.collection.populated:
|
||||
if not replacement.collection.waiting_for_elements:
|
||||
|
||||
@@ -18,29 +18,53 @@ class LocalJobCancellationTestCase(integration_util.IntegrationTestCase):
|
||||
super(LocalJobCancellationTestCase, self).setUp()
|
||||
self.dataset_populator = DatasetPopulator(self.galaxy_interactor)
|
||||
|
||||
def setup_cat_data_and_sleep(self, history_id):
|
||||
hda1 = self.dataset_populator.new_dataset(history_id, content="1 2 3")
|
||||
running_inputs = {
|
||||
"input1": {"src": "hda", "id": hda1["id"]},
|
||||
"sleep_time": 240,
|
||||
}
|
||||
running_response = self.dataset_populator.run_tool(
|
||||
"cat_data_and_sleep",
|
||||
running_inputs,
|
||||
history_id,
|
||||
assert_ok=False,
|
||||
).json()
|
||||
job_dict = running_response["jobs"][0]
|
||||
return job_dict
|
||||
|
||||
def test_cancel_job_with_admin_message(self):
|
||||
with self.dataset_populator.test_history() as history_id:
|
||||
job_dict = self.setup_cat_data_and_sleep(history_id)
|
||||
self.galaxy_interactor.wait_for(lambda: self._get("jobs/%s" % job_dict['id']).json()['state'] != 'running',
|
||||
what="Wait for job to start running",
|
||||
maxseconds=60)
|
||||
app = self._app
|
||||
sa_session = app.model.context.current
|
||||
Job = app.model.Job
|
||||
job = sa_session.query(Job).filter_by(tool_id="cat_data_and_sleep").order_by(Job.create_time.desc()).first()
|
||||
# This is how the admin controller code cancels a job
|
||||
job.job_stderr = 'admin cancelled job'
|
||||
job.set_state(app.model.Job.states.DELETED_NEW)
|
||||
sa_session.add(job)
|
||||
sa_session.flush()
|
||||
self.galaxy_interactor.wait_for(lambda: self._get("jobs/%s" % job_dict['id']).json()['state'] != 'error',
|
||||
what="Wait for job to end in error",
|
||||
maxseconds=60)
|
||||
|
||||
def test_kill_process(self):
|
||||
"""
|
||||
"""
|
||||
with self.dataset_populator.test_history() as history_id:
|
||||
hda1 = self.dataset_populator.new_dataset(history_id, content="1 2 3")
|
||||
running_inputs = {
|
||||
"input1": {"src": "hda", "id": hda1["id"]},
|
||||
"sleep_time": 240,
|
||||
}
|
||||
running_response = self.dataset_populator.run_tool(
|
||||
"cat_data_and_sleep",
|
||||
running_inputs,
|
||||
history_id,
|
||||
assert_ok=False,
|
||||
).json()
|
||||
job_dict = running_response["jobs"][0]
|
||||
job_dict = self.setup_cat_data_and_sleep(history_id)
|
||||
|
||||
app = self._app
|
||||
sa_session = app.model.context.current
|
||||
external_id = None
|
||||
state = False
|
||||
Job = app.model.Job
|
||||
|
||||
job = sa_session.query(app.model.Job).filter_by(tool_id="cat_data_and_sleep").one()
|
||||
job = sa_session.query(Job).filter_by(tool_id="cat_data_and_sleep").order_by(Job.create_time.desc()).first()
|
||||
# Not checking the state here allows the change from queued to running to overwrite
|
||||
# the change from queued to deleted_new in the API thread - this is a problem because
|
||||
# the job will still run. See issue https://github.com/galaxyproject/galaxy/issues/4960.
|
||||
|
||||
Reference in New Issue
Block a user