diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index efd7850c7e0..3d6b0dd66eb 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -1029,26 +1029,33 @@ class JobWrapper( object ): self.sa_session.add( job ) self.sa_session.flush() - def change_state( self, state, info=False ): - job = self.get_job() - self.sa_session.refresh( job ) + def change_state( self, state, info=False, flush=True, job=None ): + job_supplied = job is not None + if not job_supplied: + job = self.get_job() + self.sa_session.refresh( job ) + # Else: + # If this is a new job (e.g. initially queued) - we are in the same + # thread and no other threads are working on the job yet - so don't refresh. + if job.state in model.Job.terminal_states: log.warning( "(%s) Ignoring state change from '%s' to '%s' for job " "that is already terminal", job.id, job.state, state ) return for dataset_assoc in job.output_datasets + job.output_library_datasets: dataset = dataset_assoc.dataset - self.sa_session.refresh( dataset ) - dataset.state = state + if not job_supplied: + self.sa_session.refresh( dataset ) + dataset.raw_set_dataset_state( state ) if info: dataset.info = info self.sa_session.add( dataset ) - self.sa_session.flush() if info: job.info = info job.set_state( state ) self.sa_session.add( job ) - self.sa_session.flush() + if flush: + self.sa_session.flush() def get_state( self ): job = self.get_job() @@ -1059,22 +1066,23 @@ class JobWrapper( object ): log.warning('set_runner() is deprecated, use set_job_destination()') self.set_job_destination(self.job_destination, external_id) - def set_job_destination( self, job_destination, external_id=None ): + def set_job_destination( self, job_destination, external_id=None, flush=True, job=None ): """ Persist job destination params in the database for recovery. self.job_destination is not used because a runner may choose to rewrite parts of the destination (e.g. the params). """ - job = self.get_job() - self.sa_session.refresh(job) + if job is None: + job = self.get_job() log.debug('(%s) Persisting job destination (destination id: %s)' % (job.id, job_destination.id)) job.destination_id = job_destination.id job.destination_params = job_destination.params job.job_runner_name = job_destination.runner job.job_runner_external_id = external_id self.sa_session.add(job) - self.sa_session.flush() + if flush: + self.sa_session.flush() def finish( self, stdout, stderr, tool_exit_code=None, remote_working_directory=None ): """ @@ -1759,7 +1767,7 @@ class TaskWrapper(JobWrapper): self.status = 'error' # How do we want to handle task failure? Fail the job and let it clean up? - def change_state( self, state, info=False ): + def change_state( self, state, info=False, flush=True, job=None ): task = self.get_task() self.sa_session.refresh( task ) if info: diff --git a/lib/galaxy/jobs/runners/__init__.py b/lib/galaxy/jobs/runners/__init__.py index 3468ea693f6..406e38ddccf 100644 --- a/lib/galaxy/jobs/runners/__init__.py +++ b/lib/galaxy/jobs/runners/__init__.py @@ -106,10 +106,12 @@ class BaseJobRunner( object ): """Add a job to the queue (by job identifier), indicate that the job is ready to run. """ put_timer = ExecutionTimer() + job = job_wrapper.get_job() # Change to queued state before handing to worker thread so the runner won't pick it up again - job_wrapper.change_state( model.Job.states.QUEUED ) + job_wrapper.change_state( model.Job.states.QUEUED, flush=False, job=job ) # Persist the destination so that the job will be included in counts if using concurrency limits - job_wrapper.set_job_destination( job_wrapper.job_destination, None ) + job_wrapper.set_job_destination( job_wrapper.job_destination, None, flush=False, job=job ) + self.sa_session.flush() self.mark_as_queued(job_wrapper) log.debug("Job [%s] queued %s" % (job_wrapper.job_id, put_timer)) diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index a32007e95db..e45702a41dd 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -1789,10 +1789,17 @@ class DatasetInstance( object ): return self._state return self.dataset.state + def raw_set_dataset_state( self, state ): + if state != self.dataset.state: + self.dataset.state = state + return True + else: + return False + def set_dataset_state( self, state ): - self.dataset.state = state - object_session( self ).add( self.dataset ) - object_session( self ).flush() # flush here, because hda.flush() won't flush the Dataset object + if self.raw_set_dataset_state( state ): + object_session( self ).add( self.dataset ) + object_session( self ).flush() # flush here, because hda.flush() won't flush the Dataset object state = property( get_dataset_state, set_dataset_state ) def get_file_name( self ):