From 949d1224892861bd27207dd0799e2b356576c2cb Mon Sep 17 00:00:00 2001 From: John Chilton Date: Tue, 13 Dec 2016 16:40:42 -0500 Subject: [PATCH] Fix preservation of resubmits for resubmitted jobs. Create an abstraction for the explicit handling of resubmits that is done during job recovery at startup. xref https://github.com/galaxyproject/galaxy/commit/5539a083e21b601bfd2f87fa3417f23bf295a690 --- lib/galaxy/jobs/handler.py | 32 ++++++++++++++++++-------------- 1 file changed, 18 insertions(+), 14 deletions(-) diff --git a/lib/galaxy/jobs/handler.py b/lib/galaxy/jobs/handler.py index 4a4de2c1b52..0aa77ff5988 100644 --- a/lib/galaxy/jobs/handler.py +++ b/lib/galaxy/jobs/handler.py @@ -149,22 +149,27 @@ class JobHandlerQueue( object ): self.queue.put( ( job.id, job.tool_id ) ) else: # Already dispatched and running - job_wrapper = self.job_wrapper( job ) - # Use the persisted destination as its params may differ from - # what's in the job_conf xml - job_destination = JobDestination(id=job.destination_id, runner=job.job_runner_name, params=job.destination_params) - # resubmits are not persisted (it's a good thing) so they - # should be added back to the in-memory destination on startup - try: - config_job_destination = self.app.job_config.get_destination( job.destination_id ) - job_destination.resubmit = config_job_destination.resubmit - except KeyError: - log.warning( '(%s) Recovered destination id (%s) does not exist in job config (but this may be normal in the case of a dynamically generated destination)', job.id, job.destination_id ) - job_wrapper.job_runner_mapper.cached_job_destination = job_destination + job_wrapper = self.__recover_job_wrapper( job ) self.dispatcher.recover( job, job_wrapper ) if self.sa_session.dirty: self.sa_session.flush() + def __recover_job_wrapper(self, job): + # Already dispatched and running + job_wrapper = self.job_wrapper( job ) + # Use the persisted destination as its params may differ from + # what's in the job_conf xml + job_destination = JobDestination(id=job.destination_id, runner=job.job_runner_name, params=job.destination_params) + # resubmits are not persisted (it's a good thing) so they + # should be added back to the in-memory destination on startup + try: + config_job_destination = self.app.job_config.get_destination( job.destination_id ) + job_destination.resubmit = config_job_destination.resubmit + except KeyError: + log.debug( '(%s) Recovered destination id (%s) does not exist in job config (but this may be normal in the case of a dynamically generated destination)', job.id, job.destination_id ) + job_wrapper.job_runner_mapper.cached_job_destination = job_destination + return job_wrapper + def __monitor( self ): """ Continually iterate the waiting jobs, checking is each is ready to @@ -260,8 +265,7 @@ class JobHandlerQueue( object ): for job in resubmit_jobs: log.debug( '(%s) Job was resubmitted and is being dispatched immediately', job.id ) # Reassemble resubmit job destination from persisted value - jw = self.job_wrapper( job ) - jw.job_runner_mapper.cached_job_destination = JobDestination( id=job.destination_id, runner=job.job_runner_name, params=job.destination_params ) + jw = self.__recover_job_wrapper( job ) self.increase_running_job_count(job.user_id, jw.job_destination.id) self.dispatcher.put( jw ) # Iterate over new and waiting jobs and look for any that are