diff --git a/client/galaxy/scripts/mvc/rules/rule-definitions.js b/client/galaxy/scripts/mvc/rules/rule-definitions.js index 23ccb754b8f..29909912306 100644 --- a/client/galaxy/scripts/mvc/rules/rule-definitions.js +++ b/client/galaxy/scripts/mvc/rules/rule-definitions.js @@ -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) { diff --git a/lib/galaxy/jobs/handler.py b/lib/galaxy/jobs/handler.py index 942eac7d16c..823f185535d 100644 --- a/lib/galaxy/jobs/handler.py +++ b/lib/galaxy/jobs/handler.py @@ -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: diff --git a/lib/galaxy/jobs/runners/util/job_script/DEFAULT_JOB_FILE_TEMPLATE.sh b/lib/galaxy/jobs/runners/util/job_script/DEFAULT_JOB_FILE_TEMPLATE.sh index 81107c568b5..a588fb12b5a 100644 --- a/lib/galaxy/jobs/runners/util/job_script/DEFAULT_JOB_FILE_TEMPLATE.sh +++ b/lib/galaxy/jobs/runners/util/job_script/DEFAULT_JOB_FILE_TEMPLATE.sh @@ -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 diff --git a/lib/galaxy/managers/jobs.py b/lib/galaxy/managers/jobs.py index 3940e20c4eb..ad2d1bacfeb 100644 --- a/lib/galaxy/managers/jobs.py +++ b/lib/galaxy/managers/jobs.py @@ -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 diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index 6ad851608f2..448f42f442b 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -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 diff --git a/lib/galaxy/model/migrate/versions/0145_add_workflow_step_input.py b/lib/galaxy/model/migrate/versions/0145_add_workflow_step_input.py index 4944fc3ec9c..a5e799104d5 100644 --- a/lib/galaxy/model/migrate/versions/0145_add_workflow_step_input.py +++ b/lib/galaxy/model/migrate/versions/0145_add_workflow_step_input.py @@ -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() diff --git a/lib/galaxy/tools/evaluation.py b/lib/galaxy/tools/evaluation.py index 9d95987e8cf..dc1b6b00dae 100644 --- a/lib/galaxy/tools/evaluation.py +++ b/lib/galaxy/tools/evaluation.py @@ -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) diff --git a/lib/galaxy/tools/verify/interactor.py b/lib/galaxy/tools/verify/interactor.py index 72ddd5b4ac4..168931157f3 100644 --- a/lib/galaxy/tools/verify/interactor.py +++ b/lib/galaxy/tools/verify/interactor.py @@ -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) diff --git a/lib/galaxy/workflow/run.py b/lib/galaxy/workflow/run.py index e9120f2529b..3336139eab3 100644 --- a/lib/galaxy/workflow/run.py +++ b/lib/galaxy/workflow/run.py @@ -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: diff --git a/test/integration/test_local_job_cancellation.py b/test/integration/test_local_job_cancellation.py index abcaa0a1576..d01cea8bd2e 100644 --- a/test/integration/test_local_job_cancellation.py +++ b/test/integration/test_local_job_cancellation.py @@ -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.