From 62dbf67fcc49b5298b5e580cb2e8fc391815789b Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Wed, 2 Sep 2020 12:01:14 +0200 Subject: [PATCH 1/3] Add workflow invocation grabbing with db-skipped-lock and db-transaction-isolation. Closes https://github.com/galaxyproject/galaxy/issues/8209. --- lib/galaxy/jobs/handler.py | 137 ++++++++++++---------- lib/galaxy/workflow/scheduling_manager.py | 25 +++- 2 files changed, 98 insertions(+), 64 deletions(-) diff --git a/lib/galaxy/jobs/handler.py b/lib/galaxy/jobs/handler.py index 1d309bbb458..047727e993d 100644 --- a/lib/galaxy/jobs/handler.py +++ b/lib/galaxy/jobs/handler.py @@ -62,6 +62,71 @@ class JobHandler: self.job_stop_queue.shutdown() +class ItemGrabber: + + def __init__(self, app, grab_type='Job', handler_assignment_method=None, max_grab=None, self_handler_tags=None, handler_tags=None): + self.app = app + self.sa_session = app.model.context + self.grab_this = getattr(model, grab_type) + self.grab_type = grab_type + self._grab_conn_opts = {'autocommit': False} + subq = select([self.grab_this.id]) \ + .where(and_( + self.grab_this.table.c.handler.in_(self_handler_tags), + self.grab_this.table.c.state == self.grab_this.states.NEW)) \ + .order_by(self.grab_this.table.c.id) + if max_grab: + subq = subq.limit(max_grab) + if handler_assignment_method == HANDLER_ASSIGNMENT_METHODS.DB_SKIP_LOCKED: + subq = subq.with_for_update(skip_locked=True) + self._grab_query = self.grab_this.table.update() \ + .returning(self.grab_this.table.c.id) \ + .where(self.grab_this.table.c.id.in_(subq)) \ + .values(handler=self.app.config.server_name) + if handler_assignment_method == HANDLER_ASSIGNMENT_METHODS.DB_TRANSACTION_ISOLATION: + self._grab_conn_opts['isolation_level'] = 'SERIALIZABLE' + log.info( + "Handler job grabber initialized with '%s' assignment method for handler '%s', tag(s): %s", handler_assignment_method, + self.app.config.server_name, ', '.join(str(x) for x in handler_tags) + ) + + @staticmethod + def get_grabbable_handler_assignment_method(handler_assignment_methods): + grabbable_methods = { + HANDLER_ASSIGNMENT_METHODS.DB_TRANSACTION_ISOLATION, + HANDLER_ASSIGNMENT_METHODS.DB_SKIP_LOCKED, + } + try: + return [m for m in handler_assignment_methods if m in grabbable_methods][0] + except IndexError: + return + + def grab_unhandled_items(self): + """ + Attempts to assign unassigned jobs or invocaions to itself using DB serialization methods, if enabled. This + simply sets `Job.handler` or `WorkflowInvocation.handler` to the current server name, which causes the job to be picked up by + the appropriate handler. + """ + # an excellent discussion on PostgreSQL concurrency safety: + # https://blog.2ndquadrant.com/what-is-select-skip-locked-for-in-postgresql-9-5/ + self.sa_session.expunge_all() + conn = self.sa_session.connection(execution_options=self._grab_conn_opts) + with conn.begin() as trans: + try: + rows = conn.execute(self._grab_query).fetchall() + if rows: + log.debug('Grabbed %s(s): %s', self.grab_type, ', '.join(str(row[0]) for row in rows)) + trans.commit() + else: + trans.rollback() + 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('Grabbing %s failed (serialization failures are ok): %s', self.grab_type, unicodify(e)) + trans.rollback() + + class JobHandlerQueue(Monitors): """ Job Handler's Internal Queue, this is what actually implements waiting for @@ -91,38 +156,17 @@ class JobHandlerQueue(Monitors): self.job_wrappers = {} name = "JobHandlerQueue.monitor_thread" self._init_monitor_thread(name, target=self.__monitor, config=app.config) - self.__grab_query = None - self.__grab_conn_opts = {'autocommit': False} - self.__initialize_job_grabbing() - - def __initialize_job_grabbing(self): - grabbable_methods = { - HANDLER_ASSIGNMENT_METHODS.DB_TRANSACTION_ISOLATION, - HANDLER_ASSIGNMENT_METHODS.DB_SKIP_LOCKED, - } - try: - method = [m for m in self.app.job_config.handler_assignment_methods if m in grabbable_methods][0] - except IndexError: - return - subq = select([model.Job.id]) \ - .where(and_( - model.Job.table.c.handler.in_(self.app.job_config.self_handler_tags), - model.Job.table.c.state == model.Job.states.NEW)) \ - .order_by(model.Job.table.c.id) - if self.app.job_config.handler_max_grab: - subq = subq.limit(self.app.job_config.handler_max_grab) - if method == HANDLER_ASSIGNMENT_METHODS.DB_SKIP_LOCKED: - subq = subq.with_for_update(skip_locked=True) - self.__grab_query = model.Job.table.update() \ - .returning(model.Job.table.c.id) \ - .where(model.Job.table.c.id.in_(subq)) \ - .values(handler=self.app.config.server_name) - if method == HANDLER_ASSIGNMENT_METHODS.DB_TRANSACTION_ISOLATION: - self.__grab_conn_opts['isolation_level'] = 'SERIALIZABLE' - log.info( - "Handler job grabber initialized with '%s' assignment method for handler '%s', tag(s): %s", method, - self.app.config.server_name, ', '.join(str(x) for x in self.app.job_config.handler_tags) - ) + self.job_grabber = None + handler_assignment_method = ItemGrabber.get_grabbable_handler_assignment_method(self.app.job_config.handler_assignment_methods) + if handler_assignment_method: + self.job_grabber = ItemGrabber( + app=app, + grab_type='Job', + handler_assignment_method=handler_assignment_method, + max_grab=self.app.job_config.handler_max_grab, + self_handler_tags=self.app.job_config.self_handler_tags, + handler_tags=self.app.job_config.handler_tags, + ) def start(self): """ @@ -248,36 +292,11 @@ class JobHandlerQueue(Monitors): 'internal.galaxy.jobs.handlers.monitor_step', 'Job handler monitor step complete.' ) - if self.__grab_query is not None: - self.__grab_unhandled_jobs() + if self.job_grabber is not None: + self.job_grabber.grab_unhandled_items() self.__handle_waiting_jobs() log.trace(monitor_step_timer.to_str()) - def __grab_unhandled_jobs(self): - """ - Attempts to "grab" jobs (assign unassigned jobs to itself) using DB serialization methods, if enabled. This - simply sets `Job.handler` to the current server name, which causes the job to be picked up by - `__handle_waiting_jobs()`. - """ - # an excellent discussion on PostgreSQL concurrency safety: - # https://blog.2ndquadrant.com/what-is-select-skip-locked-for-in-postgresql-9-5/ - self.sa_session.expunge_all() - conn = self.sa_session.connection(execution_options=self.__grab_conn_opts) - with conn.begin() as trans: - try: - rows = conn.execute(self.__grab_query).fetchall() - if rows: - log.debug('Grabbed job(s): %s', ', '.join(str(row[0]) for row in rows)) - trans.commit() - else: - trans.rollback() - 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('Grabbing job failed (serialization failures are ok): %s', unicodify(e)) - trans.rollback() - def __handle_waiting_jobs(self): """ Gets any new jobs (either from the database or from its own queue), then iterates over all new and waiting jobs diff --git a/lib/galaxy/workflow/scheduling_manager.py b/lib/galaxy/workflow/scheduling_manager.py index ecfd591ca96..031950cfba4 100644 --- a/lib/galaxy/workflow/scheduling_manager.py +++ b/lib/galaxy/workflow/scheduling_manager.py @@ -4,6 +4,7 @@ from functools import partial import galaxy.workflow.schedulers from galaxy import model from galaxy.exceptions import HandlerAssignmentError +from galaxy.jobs.handler import ItemGrabber from galaxy.util import ( parse_xml, plugin_config, @@ -31,11 +32,9 @@ class WorkflowSchedulingManager(ConfiguresHandlers): processes. """ DEFAULT_BASE_HANDLER_POOLS = ('workflow-schedulers', 'job-handlers') - UNSUPPORTED_HANDLER_ASSIGNMENT_METHODS = ( - HANDLER_ASSIGNMENT_METHODS.DB_TRANSACTION_ISOLATION, - HANDLER_ASSIGNMENT_METHODS.DB_SKIP_LOCKED, + UNSUPPORTED_HANDLER_ASSIGNMENT_METHODS = { HANDLER_ASSIGNMENT_METHODS.UWSGI_MULE_MESSAGE, - ) + } def __init__(self, app): self.app = app @@ -267,10 +266,27 @@ class WorkflowRequestMonitor(Monitors): self.app = app self.workflow_scheduling_manager = workflow_scheduling_manager self._init_monitor_thread(name="WorkflowRequestMonitor.monitor_thread", target=self.__monitor, config=app.config) + self.invocation_grabber = None + if self.workflow_scheduling_manager.handler_assignment_methods_configured: + self_handler_tags = set(self.app.job_config.self_handler_tags) + self_handler_tags.add(self.workflow_scheduling_manager.default_handler_id) + handler_assignment_method = ItemGrabber.get_grabbable_handler_assignment_method(self.workflow_scheduling_manager.handler_assignment_methods) + if handler_assignment_method: + self.invocation_grabber = ItemGrabber( + app=app, + grab_type='WorkflowInvocation', + handler_assignment_method=handler_assignment_method, + max_grab=self.workflow_scheduling_manager.handler_max_grab, + self_handler_tags=self_handler_tags, + handler_tags=self_handler_tags, + ) def __monitor(self): to_monitor = self.workflow_scheduling_manager.active_workflow_schedulers while self.monitor_running: + if self.invocation_grabber: + self.invocation_grabber.grab_unhandled_items() + monitor_step_timer = self.app.execution_timer_factory.get_timer( 'internal.galaxy.workflows.scheduling_manager.monitor_step', 'Workflow scheduling manager monitor step complete.' @@ -280,7 +296,6 @@ class WorkflowRequestMonitor(Monitors): return self.__schedule(workflow_scheduler_id, workflow_scheduler) - log.trace(monitor_step_timer.to_str()) self._monitor_sleep(1) From 51a3c10c68e8eef9028f9ac59dc95175d9748faa Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Wed, 2 Sep 2020 20:15:19 +0200 Subject: [PATCH 2/3] Add test for db-skip-locked/db-transaction-isolation --- lib/galaxy_test/driver/integration_util.py | 7 +++ .../test_workflow_handler_configuration.py | 53 ++++++++++++++++--- ..._configuration_workflow_scheduler_conf.xml | 8 --- 3 files changed, 54 insertions(+), 14 deletions(-) delete mode 100644 test/integration/workflow_handler_configuration_workflow_scheduler_conf.xml diff --git a/lib/galaxy_test/driver/integration_util.py b/lib/galaxy_test/driver/integration_util.py index c7f30fbdb22..444c6278a15 100644 --- a/lib/galaxy_test/driver/integration_util.py +++ b/lib/galaxy_test/driver/integration_util.py @@ -17,6 +17,7 @@ from .driver_util import GalaxyTestDriver NO_APP_MESSAGE = "test_case._app called though no Galaxy has been configured." # Following should be for Homebrew Rabbitmq and Docker on Mac "amqp://guest:guest@localhost:5672//" AMQP_URL = os.environ.get("GALAXY_TEST_AMQP_URL", None) +POSTGRES_CONFIGURED = 'postgres' in os.environ.get("GALAXY_TEST_DBURI", '') def _identity(func): @@ -36,6 +37,12 @@ def skip_unless_amqp(): return pytest.mark.skip("AMQP_URL is not set, required for this test.") +def skip_unless_postgres(): + if POSTGRES_CONFIGURED: + return _identity + return pytest.mark.skip("GALAXY_TEST_DBURI does not point to postgres database, required for this test.") + + def skip_unless_executable(executable): if which(executable): return _identity diff --git a/test/integration/test_workflow_handler_configuration.py b/test/integration/test_workflow_handler_configuration.py index 423c67f7cd4..7794323c9e4 100644 --- a/test/integration/test_workflow_handler_configuration.py +++ b/test/integration/test_workflow_handler_configuration.py @@ -2,6 +2,9 @@ import os import re +import string +import tempfile +import time from json import dumps from galaxy_test.base.populators import ( @@ -12,8 +15,16 @@ from galaxy_test.driver import integration_util SCRIPT_DIRECTORY = os.path.abspath(os.path.dirname(__file__)) WORKFLOW_HANDLER_CONFIGURATION_JOB_CONF = os.path.join(SCRIPT_DIRECTORY, "workflow_handler_configuration_job_conf.xml") -WORKFLOW_HANDLER_CONFIGURATION_WORKFLOW_SCHEDULER_CONF = os.path.join(SCRIPT_DIRECTORY, "workflow_handler_configuration_workflow_scheduler_conf.xml") +WORKFLOW_SCHEDULERS_CONFIG_TEMPLATE = string.Template(""" + + + + + + + +""") JOB_HANDLER_PATTERN = re.compile(r"handler\d") WORKFLOW_SCHEDULER_HANDLER_PATTERN = re.compile(r"work\d") @@ -30,6 +41,16 @@ steps: """ +def workflow_schedulers_config_file(assign_with=''): + fd, path = tempfile.mkstemp(suffix=".xml", prefix="workflow_handler_config_") + os.close(fd) + with open(path, 'w') as config: + if assign_with: + assign_with = 'assign_with="{}"'.format(assign_with) + config.write(WORKFLOW_SCHEDULERS_CONFIG_TEMPLATE.substitute(assign_with=assign_with)) + return path + + class BaseWorkflowHandlerConfigurationTestCase(integration_util.IntegrationTestCase): framework_tool_and_types = True @@ -114,7 +135,7 @@ class WorkflowSchedulerHandlerAssignment(BaseWorkflowHandlerConfigurationTestCas @classmethod def handle_galaxy_config_kwds(cls, config): BaseWorkflowHandlerConfigurationTestCase.handle_galaxy_config_kwds(config) - config["workflow_schedulers_config_file"] = WORKFLOW_HANDLER_CONFIGURATION_WORKFLOW_SCHEDULER_CONF + config["workflow_schedulers_config_file"] = workflow_schedulers_config_file() def test_handler_assignment(self): self._invoke_n_workflows(1) @@ -163,23 +184,43 @@ class DefaultWorkflowHandlerIfJobHandlerOffTestCase(BaseWorkflowHandlerConfigura class ExplicitWorkflowHandlersOnTestCase(BaseWorkflowHandlerConfigurationTestCase): + assign_with = "" + @classmethod def handle_galaxy_config_kwds(cls, config): BaseWorkflowHandlerConfigurationTestCase.handle_galaxy_config_kwds(config) - config["workflow_schedulers_config_file"] = WORKFLOW_HANDLER_CONFIGURATION_WORKFLOW_SCHEDULER_CONF + config["workflow_schedulers_config_file"] = workflow_schedulers_config_file(cls.assign_with) config["server_name"] = "work1" - def test_workflows_spread_across_multiple_handlers(self): + def test_app_is_workflow_scheduler(self): assert self.is_app_workflow_scheduler +@integration_util.skip_unless_postgres() +class WorkflowSchedulerHandlerAssignmentDbSkipLocked(ExplicitWorkflowHandlersOnTestCase): + + assign_with = 'db-skip-locked' + + def test_handler_assignment(self): + self._invoke_n_workflows(1) + time.sleep(2) + workflow_invocations = self._get_workflow_invocations() + assert WORKFLOW_SCHEDULER_HANDLER_PATTERN.match(workflow_invocations[0].handler) + + +@integration_util.skip_unless_postgres() +class WorkflowSchedulerHandlerAssignmentDbTransactionIsolation(WorkflowSchedulerHandlerAssignmentDbSkipLocked): + + assign_with = 'db-transaction-isolation' + + class ExplicitWorkflowHandlersOffTestCase(BaseWorkflowHandlerConfigurationTestCase): @classmethod def handle_galaxy_config_kwds(cls, config): BaseWorkflowHandlerConfigurationTestCase.handle_galaxy_config_kwds(config) - config["workflow_schedulers_config_file"] = WORKFLOW_HANDLER_CONFIGURATION_WORKFLOW_SCHEDULER_CONF + config["workflow_schedulers_config_file"] = workflow_schedulers_config_file() config["server_name"] = "handler0" # Configured as a job handler but not a workflow handler. - def test_workflows_spread_across_multiple_handlers(self): + def test_app_is_not_workflow_scheduler(self): assert not self.is_app_workflow_scheduler diff --git a/test/integration/workflow_handler_configuration_workflow_scheduler_conf.xml b/test/integration/workflow_handler_configuration_workflow_scheduler_conf.xml deleted file mode 100644 index 7fff1a5e703..00000000000 --- a/test/integration/workflow_handler_configuration_workflow_scheduler_conf.xml +++ /dev/null @@ -1,8 +0,0 @@ - - - - - - - - From 78daef3a8bc9c1114a7b7cbe0bc41a7277755724 Mon Sep 17 00:00:00 2001 From: mvdbeek Date: Fri, 18 Sep 2020 16:05:44 +0200 Subject: [PATCH 3/3] Add a test for job handler only with db-skip-locked --- .../test_workflow_handler_configuration.py | 65 ++++++++++++++++--- 1 file changed, 57 insertions(+), 8 deletions(-) diff --git a/test/integration/test_workflow_handler_configuration.py b/test/integration/test_workflow_handler_configuration.py index 7794323c9e4..f6c33aeadb6 100644 --- a/test/integration/test_workflow_handler_configuration.py +++ b/test/integration/test_workflow_handler_configuration.py @@ -16,6 +16,32 @@ from galaxy_test.driver import integration_util SCRIPT_DIRECTORY = os.path.abspath(os.path.dirname(__file__)) WORKFLOW_HANDLER_CONFIGURATION_JOB_CONF = os.path.join(SCRIPT_DIRECTORY, "workflow_handler_configuration_job_conf.xml") +WORKFLOW_HANDLER_JOB_CONFIG_TEMPLATE = string.Template(""" + + + + + + + + + + + + + + + + + + + + + + + +""") + WORKFLOW_SCHEDULERS_CONFIG_TEMPLATE = string.Template(""" @@ -41,19 +67,20 @@ steps: """ -def workflow_schedulers_config_file(assign_with=''): +def config_file(template, assign_with=''): fd, path = tempfile.mkstemp(suffix=".xml", prefix="workflow_handler_config_") os.close(fd) with open(path, 'w') as config: if assign_with: assign_with = 'assign_with="{}"'.format(assign_with) - config.write(WORKFLOW_SCHEDULERS_CONFIG_TEMPLATE.substitute(assign_with=assign_with)) + config.write(template.substitute(assign_with=assign_with)) return path class BaseWorkflowHandlerConfigurationTestCase(integration_util.IntegrationTestCase): framework_tool_and_types = True + assign_with = "" def setUp(self): super(BaseWorkflowHandlerConfigurationTestCase, self).setUp() @@ -63,7 +90,7 @@ class BaseWorkflowHandlerConfigurationTestCase(integration_util.IntegrationTestC @classmethod def handle_galaxy_config_kwds(cls, config): - config["job_config_file"] = WORKFLOW_HANDLER_CONFIGURATION_JOB_CONF + config["job_config_file"] = config_file(WORKFLOW_HANDLER_JOB_CONFIG_TEMPLATE, assign_with=cls.assign_with) def _invoke_n_workflows(self, n): workflow_id = self.workflow_populator.upload_yaml_workflow(PAUSE_WORKFLOW) @@ -135,7 +162,7 @@ class WorkflowSchedulerHandlerAssignment(BaseWorkflowHandlerConfigurationTestCas @classmethod def handle_galaxy_config_kwds(cls, config): BaseWorkflowHandlerConfigurationTestCase.handle_galaxy_config_kwds(config) - config["workflow_schedulers_config_file"] = workflow_schedulers_config_file() + config["workflow_schedulers_config_file"] = config_file(WORKFLOW_SCHEDULERS_CONFIG_TEMPLATE, assign_with=cls.assign_with) def test_handler_assignment(self): self._invoke_n_workflows(1) @@ -143,11 +170,14 @@ class WorkflowSchedulerHandlerAssignment(BaseWorkflowHandlerConfigurationTestCas assert WORKFLOW_SCHEDULER_HANDLER_PATTERN.match(workflow_invocations[0].handler) -# Follow five classes test 5 different ways Galaxy processes can be workflow schedulers or not. +# Following 8 classes test 8 different ways Galaxy processes can be workflow schedulers or not. # - In single process mode, the process is a workflow scheduler. # - If no workflow schedulers conf is configured and the process is a job handler, it is a workflow scheduler as well. +# - If no workflow schedulers conf is configured and the process is a job handler, it is a workflow scheduler as well (with db-dkip-locked). # - If no workflow schedulers conf is configured and the process is not a job handler, it is not a workflow scheduler as well. # - If a workflow scheduler conf is defined and the process is listed as a handler, it is workflow scheduler. +# - If a workflow scheduler conf is defined and assign_with is set to db-skip-locked, invocation handler is correctly set +# - If a workflow scheduler conf is defined and assign_with is set to db-transaction-isolation, invocation handler is correctly set # - If a workflow scheduler conf is defined and the process is not listed as a handler, it is not workflow scheduler. class DefaultWorkflowHandlerOnTestCase(BaseWorkflowHandlerConfigurationTestCase): @@ -171,6 +201,25 @@ class DefaultWorkflowHandlerIfJobHandlerOnTestCase(BaseWorkflowHandlerConfigurat assert self.is_app_workflow_scheduler +class JobHandlerAsWorkflowHandlerWithDbSkipLocked(BaseWorkflowHandlerConfigurationTestCase): + + assign_with = 'db-skip-locked' + + @classmethod + def handle_galaxy_config_kwds(cls, config): + BaseWorkflowHandlerConfigurationTestCase.handle_galaxy_config_kwds(config) + config["server_name"] = "handler0" + + def test_handler_assignment(self): + self._invoke_n_workflows(1) + time.sleep(2) + workflow_invocations = self._get_workflow_invocations() + assert JOB_HANDLER_PATTERN.match(workflow_invocations[0].handler) + + def test_default_job_handler_is_workflow_handler(self): + assert self.is_app_workflow_scheduler + + class DefaultWorkflowHandlerIfJobHandlerOffTestCase(BaseWorkflowHandlerConfigurationTestCase): @classmethod @@ -178,7 +227,7 @@ class DefaultWorkflowHandlerIfJobHandlerOffTestCase(BaseWorkflowHandlerConfigura BaseWorkflowHandlerConfigurationTestCase.handle_galaxy_config_kwds(config) config["server_name"] = "web0" - def test_default_job_handler_is_workflow_handler(self): + def test_default_job_handler_is_not_workflow_handler(self): assert not self.is_app_workflow_scheduler @@ -189,7 +238,7 @@ class ExplicitWorkflowHandlersOnTestCase(BaseWorkflowHandlerConfigurationTestCas @classmethod def handle_galaxy_config_kwds(cls, config): BaseWorkflowHandlerConfigurationTestCase.handle_galaxy_config_kwds(config) - config["workflow_schedulers_config_file"] = workflow_schedulers_config_file(cls.assign_with) + config["workflow_schedulers_config_file"] = config_file(WORKFLOW_SCHEDULERS_CONFIG_TEMPLATE, assign_with=cls.assign_with) config["server_name"] = "work1" def test_app_is_workflow_scheduler(self): @@ -219,7 +268,7 @@ class ExplicitWorkflowHandlersOffTestCase(BaseWorkflowHandlerConfigurationTestCa @classmethod def handle_galaxy_config_kwds(cls, config): BaseWorkflowHandlerConfigurationTestCase.handle_galaxy_config_kwds(config) - config["workflow_schedulers_config_file"] = workflow_schedulers_config_file() + config["workflow_schedulers_config_file"] = config_file(WORKFLOW_SCHEDULERS_CONFIG_TEMPLATE, assign_with=cls.assign_with) config["server_name"] = "handler0" # Configured as a job handler but not a workflow handler. def test_app_is_not_workflow_scheduler(self):