diff --git a/lib/galaxy/jobs/runners/__init__.py b/lib/galaxy/jobs/runners/__init__.py index 095f6e3527b..de6aa6cb184 100644 --- a/lib/galaxy/jobs/runners/__init__.py +++ b/lib/galaxy/jobs/runners/__init__.py @@ -131,6 +131,11 @@ class BaseJobRunner(object): method(arg) except Exception: log.exception("(%s) Unhandled exception calling %s" % (job_id, name)) + if not isinstance(arg, JobState): + job_state = JobState(job_wrapper=arg, job_destination={}) + else: + job_state = arg + self.work_queue.put((self.fail_job, job_state)) # Causes a runner's `queue_job` method to be called from a worker thread def put(self, job_wrapper): diff --git a/test/integration/test_work_queue_put_failure.py b/test/integration/test_work_queue_put_failure.py new file mode 100644 index 00000000000..b155d01abb3 --- /dev/null +++ b/test/integration/test_work_queue_put_failure.py @@ -0,0 +1,41 @@ +import tempfile + +from base import integration_util +from base.populators import DatasetPopulator + +job_conf_yaml = """ +runners: + local: + load: galaxy.jobs.runners.pulsar:PulsarRESTJobRunner + workers: 1 +execution: + default: local + environments: + local: + runner: local +""" + + +class WorkQueuePutFailureTestCase(integration_util.IntegrationTestCase): + + def setUp(self): + super(WorkQueuePutFailureTestCase, self).setUp() + self.dataset_populator = DatasetPopulator(self.galaxy_interactor) + self.history_id = self.dataset_populator.new_history() + + @classmethod + def handle_galaxy_config_kwds(cls, config, ): + # config["jobs_directory"] = cls.jobs_directory + fd, path = tempfile.mkstemp(suffix='job_conf.yml') + with open(path, 'w') as job_conf: + job_conf.write(job_conf_yaml) + config["job_config_file"] = path + # Disable tool dependency resolution. + config["tool_dependency_dir"] = "none" + + def test_job_fails(self): + self.dataset_populator.new_dataset(self.history_id, content="1 2 3") + self.dataset_populator.wait_for_history(self.history_id, assert_ok=False) + state_details = self.galaxy_interactor.get('histories/%s' % self.history_id).json()['state_details'] + assert state_details['running'] == 0 + assert state_details['error'] == 1