diff --git a/lib/galaxy/jobs/runners/__init__.py b/lib/galaxy/jobs/runners/__init__.py index 1be761b9b9c..2676e6231a7 100644 --- a/lib/galaxy/jobs/runners/__init__.py +++ b/lib/galaxy/jobs/runners/__init__.py @@ -10,6 +10,7 @@ import threading import time import traceback import typing +import uuid from queue import ( Empty, Queue, @@ -33,10 +34,7 @@ from galaxy.jobs.runners.util.job_script import ( job_script, write_script, ) -from galaxy.model.base import ( - check_database_connection, - transaction, -) +from galaxy.model.base import transaction from galaxy.tool_util.deps.dependencies import ( JobInfo, ToolInfo, @@ -830,12 +828,18 @@ class AsynchronousJobRunner(BaseJobRunner, Monitors): self.watched.append(async_job_state) except Empty: pass + # Ideally we'd construct a sqlalchemy session now and pass it into `check_watched_items` + # and have that be the only session being used. The next best thing is to scope + # the session and discard it after each check_watched_item loop + scoped_id = str(uuid.uuid4()) + self.app.model.set_request_id(scoped_id) # Iterate over the list of watched jobs and check state try: - check_database_connection(self.sa_session) self.check_watched_items() except Exception: log.exception("Unhandled exception checking active jobs") + finally: + self.app.model.unset_request_id(scoped_id) # Sleep a bit before the next state check time.sleep(self.app.config.job_runner_monitor_sleep)