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)
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..f6c33aeadb6 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,42 @@ 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_HANDLER_JOB_CONFIG_TEMPLATE = string.Template("""
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+""")
+
+WORKFLOW_SCHEDULERS_CONFIG_TEMPLATE = string.Template("""
+
+
+
+
+
+
+
+""")
JOB_HANDLER_PATTERN = re.compile(r"handler\d")
WORKFLOW_SCHEDULER_HANDLER_PATTERN = re.compile(r"work\d")
@@ -30,9 +67,20 @@ steps:
"""
+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(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()
@@ -42,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)
@@ -114,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_HANDLER_CONFIGURATION_WORKFLOW_SCHEDULER_CONF
+ 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)
@@ -122,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):
@@ -150,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
@@ -157,29 +227,49 @@ 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
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"] = config_file(WORKFLOW_SCHEDULERS_CONFIG_TEMPLATE, assign_with=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"] = 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_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 @@
-
-
-
-
-
-
-
-