Merge pull request #8161 from mvdbeek/fail_job_if_queueing_job_fails

Fail job if queueing job fails
This commit is contained in:
Dannon
2019-06-19 13:52:00 -04:00
committed by GitHub
2 changed files with 46 additions and 0 deletions
+5
View File
@@ -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):
@@ -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