diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index 2ef3f7a3f63..94395ca2cc7 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -1321,7 +1321,7 @@ class JobWrapper(HasResourceParameters): # Pause any dependent jobs (and those jobs' outputs) for dep_job_assoc in dataset.dependent_jobs: self.pause(dep_job_assoc.job, "Execution of this dataset's job is paused because its input datasets are in an error state.") - job.set_final_state(job.states.ERROR) + job.set_final_state(job.states.ERROR, supports_skip_locked=self.app.application_stack.supports_skip_locked()) job.command_line = unicodify(self.command_line) job.info = message # TODO: Put setting the stdout, stderr, and exit code in one place @@ -1408,7 +1408,7 @@ class JobWrapper(HasResourceParameters): job.info = info job.set_state(state) self.sa_session.add(job) - job.update_output_states() + job.update_output_states(self.app.application_stack.supports_skip_locked()) if flush: self.sa_session.flush() @@ -1776,7 +1776,7 @@ class JobWrapper(HasResourceParameters): # Finally set the job state. This should only happen *after* all # dataset creation, and will allow us to eliminate force_history_refresh. - job.set_final_state(final_job_state) + job.set_final_state(final_job_state, supports_skip_locked=self.app.application_stack.supports_skip_locked()) if not job.tasks: # If job was composed of tasks, don't attempt to recollect statistics self._collect_metrics(job, job_metrics_directory) diff --git a/lib/galaxy/jobs/handler.py b/lib/galaxy/jobs/handler.py index 8e341fa294c..e1d3551aad1 100644 --- a/lib/galaxy/jobs/handler.py +++ b/lib/galaxy/jobs/handler.py @@ -924,7 +924,7 @@ class JobHandlerStopQueue(Monitors): if error_msg is not None: final_state = job.states.ERROR job.info = error_msg - job.set_final_state(final_state) + job.set_final_state(final_state, supports_skip_locked=self.app.application_stack.supports_skip_locked()) self.sa_session.add(job) self.sa_session.flush() diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index a02783f3357..b47ef0bcdbe 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -37,6 +37,7 @@ from sqlalchemy import ( true, type_coerce, types) +from sqlalchemy.exc import OperationalError from sqlalchemy.ext import hybrid from sqlalchemy.orm import ( aliased, @@ -1177,22 +1178,36 @@ class Job(JobLike, UsesCreateAndUpdateTime, Dictifiable, RepresentById): return rval - def update_hdca_update_time_for_job(self, update_time, sa_session): + def update_hdca_update_time_for_job(self, update_time, sa_session, supports_skip_locked): subq = sa_session.query(HistoryDatasetCollectionAssociation.id) \ .join(ImplicitCollectionJobs) \ .join(ImplicitCollectionJobsJobAssociation) \ - .filter(ImplicitCollectionJobsJobAssociation.job_id == self.id) \ - .with_for_update(skip_locked=True).subquery() + .filter(ImplicitCollectionJobsJobAssociation.job_id == self.id) + if supports_skip_locked: + subq = subq.with_for_update(skip_locked=True).subquery() implicit_statement = HistoryDatasetCollectionAssociation.table.update() \ .where(HistoryDatasetCollectionAssociation.table.c.id.in_(subq)) \ .values(update_time=update_time) explicit_statement = HistoryDatasetCollectionAssociation.table.update() \ .where(HistoryDatasetCollectionAssociation.table.c.job_id == self.id) \ .values(update_time=update_time) - sa_session.execute(implicit_statement) sa_session.execute(explicit_statement) + if supports_skip_locked: + sa_session.execute(implicit_statement) + else: + conn = sa_session.connection(execution_options={'isolation_level': 'SERIALIZABLE'}) + with conn.begin() as trans: + try: + conn.execute(implicit_statement) + trans.commit() + except OperationalError as e: + # If this is a serialization failure on PostgreSQL, then e.orig is a psycopg2 TransactionRollbackError + # and should have attribute `code`. Other engines should just report the message and move on. + if int(getattr(e.orig, 'pgcode', -1)) != 40001: + log.debug(f"Updating implicit collection uptime_time for job {self.id} failed (this is expected for large collections and not a problem): {unicodify(e)}") + trans.rollback() - def set_final_state(self, final_state): + def set_final_state(self, final_state, supports_skip_locked): self.set_state(final_state) # TODO: migrate to where-in subqueries? statement = ''' @@ -1202,8 +1217,7 @@ class Job(JobLike, UsesCreateAndUpdateTime, Dictifiable, RepresentById): ''' sa_session = object_session(self) update_time = galaxy.model.orm.now.now() - log.debug(f'Final state update time: {update_time}') - self.update_hdca_update_time_for_job(update_time=update_time, sa_session=sa_session) + self.update_hdca_update_time_for_job(update_time=update_time, sa_session=sa_session, supports_skip_locked=supports_skip_locked) params = { 'job_id': self.id, 'update_time': update_time @@ -1230,7 +1244,7 @@ class Job(JobLike, UsesCreateAndUpdateTime, Dictifiable, RepresentById): for dataset_assoc in self.output_datasets: return dataset_assoc.dataset.tool_version - def update_output_states(self): + def update_output_states(self, supports_skip_locked): # TODO: migrate to where-in subqueries? statements = [''' UPDATE dataset @@ -1275,7 +1289,7 @@ class Job(JobLike, UsesCreateAndUpdateTime, Dictifiable, RepresentById): '''] sa_session = object_session(self) update_time = galaxy.model.orm.now.now() - self.update_hdca_update_time_for_job(update_time=update_time, sa_session=sa_session) + self.update_hdca_update_time_for_job(update_time=update_time, sa_session=sa_session, supports_skip_locked=supports_skip_locked) params = { 'job_id': self.id, 'state': self.state, diff --git a/test/unit/managers/test_HistoryContentsManager.py b/test/unit/managers/test_HistoryContentsManager.py index eeb6c34d9fc..e43ba2d06e2 100644 --- a/test/unit/managers/test_HistoryContentsManager.py +++ b/test/unit/managers/test_HistoryContentsManager.py @@ -241,15 +241,15 @@ class HistoryAsContainerTestCase(HistoryAsContainerBaseTestCase): contents.append(self.add_list_collection_to_history(history, contents[4:6])) self.log("should allow filtering by update_time") - # in the case of collections we have to change the collection.collection (ugh) to change the update_time - contents[3].collection.populated_state = 'big ball of mud' + # change the update_time by updating the name + contents[3].name = 'big ball of mud' self.app.model.context.flush() - update_time = contents[3].collection.update_time + update_time = contents[3].update_time def get_update_time(item): update_time = getattr(item, 'update_time', None) if not update_time: - update_time = item.collection.update_time + update_time = item.update_time return update_time results = self.contents_manager.contents(history, filters=[parsed_filter("orm", column('update_time') >= update_time)])