diff --git a/config/galaxy.ini.sample b/config/galaxy.ini.sample index 70277c6f94c..be3366072fe 100644 --- a/config/galaxy.ini.sample +++ b/config/galaxy.ini.sample @@ -1062,6 +1062,11 @@ use_interactive = True # collections. #force_beta_workflow_scheduled_for_collections=False +# If multiple job handlers are enabled allow Galaxy to schedule workflow invocations +# in multiple handlers simultaneously. This is discouraged because it results in a +# less predictable order of workflow datasets within in histories. +#parallelize_workflow_scheduling_within_histories = False + # This is the maximum amount of time a workflow invocation may stay in an active # scheduling state in seconds. Set to -1 to disable this maximum and allow any workflow # invocation to schedule indefinitely. The default corresponds to 1 month. diff --git a/lib/galaxy/config.py b/lib/galaxy/config.py index 8c9fa513f08..d3af12b4305 100644 --- a/lib/galaxy/config.py +++ b/lib/galaxy/config.py @@ -325,6 +325,7 @@ class Configuration( object ): self.force_beta_workflow_scheduled_for_collections = string_as_bool( kwargs.get( 'force_beta_workflow_scheduled_for_collections', 'False' ) ) self.history_local_serial_workflow_scheduling = string_as_bool( kwargs.get( 'history_local_serial_workflow_scheduling', 'False' ) ) + self.parallelize_workflow_scheduling_within_histories = string_as_bool( kwargs.get( 'parallelize_workflow_scheduling_within_histories', 'False' ) ) self.maximum_workflow_invocation_duration = int( kwargs.get( "maximum_workflow_invocation_duration", 2678400 ) ) # Per-user Job concurrency limitations diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index 471b70d9d74..f83a13cb955 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -608,27 +608,31 @@ class JobConfiguration( object ): rval.append(self.default_job_tool_configuration) return rval - def __get_single_item(self, collection): + def __get_single_item(self, collection, index=None): """Given a collection of handlers or destinations, return one item from the collection at random. """ # Done like this to avoid random under the assumption it's faster to avoid it if len(collection) == 1: return collection[0] - else: + elif index is None: return random.choice(collection) + else: + return collection[index % len(collection)] # This is called by Tool.get_job_handler() - def get_handler(self, id_or_tag): - """Given a handler ID or tag, return the provided ID or an ID matching the provided tag + def get_handler(self, id_or_tag, index=None): + """Given a handler ID or tag, return a handler matching it. :param id_or_tag: A handler ID or tag. :type id_or_tag: str + :param index: Generate "consistent" "random" handlers with this index if specified. + :type index: int :returns: str -- A valid job handler ID. """ if id_or_tag is None: id_or_tag = self.default_handler_id - return self.__get_single_item(self.handlers[id_or_tag]) + return self.__get_single_item(self.handlers[id_or_tag], index=index) def get_destination(self, id_or_tag): """Given a destination ID or tag, return the JobDestination matching the provided ID or tag diff --git a/lib/galaxy/workflow/scheduling_manager.py b/lib/galaxy/workflow/scheduling_manager.py index 030ca0e86a8..1ddd89926e6 100644 --- a/lib/galaxy/workflow/scheduling_manager.py +++ b/lib/galaxy/workflow/scheduling_manager.py @@ -54,8 +54,13 @@ class WorkflowSchedulingManager( object ): def _is_workflow_handler( self ): return self.app.is_job_handler() - def _get_handler( self ): - return self.__job_config.get_handler( None ) + def _get_handler( self, history_id ): + # Use random-ish integer history_id to produce a consistent index to pick + # job handler with. + random_index = history_id + if self.app.config.parallelize_workflow_scheduling_within_histories: + random_index = None + return self.__job_config.get_handler( None, index=random_index ) def shutdown( self ): for workflow_scheduler in self.workflow_schedulers.itervalues(): @@ -72,7 +77,7 @@ class WorkflowSchedulingManager( object ): def queue( self, workflow_invocation, request_params ): workflow_invocation.state = model.WorkflowInvocation.states.NEW scheduler = request_params.get( "scheduler", None ) or self.default_scheduler_id - handler = self._get_handler() + handler = self._get_handler( workflow_invocation.history.id ) log.info("Queueing workflow invocation for handler [%s]" % handler) workflow_invocation.scheduler = scheduler diff --git a/test/integration/test_workflow_handler_configuration.py b/test/integration/test_workflow_handler_configuration.py new file mode 100644 index 00000000000..ddf100f57c9 --- /dev/null +++ b/test/integration/test_workflow_handler_configuration.py @@ -0,0 +1,97 @@ +"""Integration tests for maximum workflow invocation duration configuration option.""" + +import os + +from json import dumps + +from base import integration_util +from base.populators import ( + DatasetPopulator, + WorkflowPopulator, +) + +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") + +PAUSE_WORKFLOW = """ +class: GalaxyWorkflow +steps: +- label: test_input + type: input +- label: the_pause + type: pause + connect: + input: + - test_input +""" + + +class BaseWorkflowHandlerConfigurationTestCase(integration_util.IntegrationTestCase): + + framework_tool_and_types = True + + def setUp( self ): + super( BaseWorkflowHandlerConfigurationTestCase, self ).setUp() + self.dataset_populator = DatasetPopulator( self.galaxy_interactor ) + self.workflow_populator = WorkflowPopulator( self.galaxy_interactor ) + self.history_id = self.dataset_populator.new_history() + + @classmethod + def handle_galaxy_config_kwds(cls, config): + config["job_config_file"] = WORKFLOW_HANDLER_CONFIGURATION_JOB_CONF + + def _invoke_n_workflows(self, n): + workflow_id = self.workflow_populator.upload_yaml_workflow(PAUSE_WORKFLOW) + history_id = self.history_id + hda1 = self.dataset_populator.new_dataset(history_id, content="1 2 3") + index_map = { + '0': dict(src="hda", id=hda1["id"]) + } + request = {} + request["history"] = "hist_id=%s" % history_id + request["inputs"] = dumps(index_map) + request["inputs_by"] = 'step_index' + url = "workflows/%s/invocations" % (workflow_id) + for i in range(n): + self._post(url, data=request) + + def _get_workflow_invocations(self): + # Consider exposing handler via the API to reduce breaking + # into Galaxy's internal state. + app = self._app + history_id = app.security.decode_id(self.history_id) + sa_session = app.model.context.current + history = sa_session.query( app.model.History ).get( history_id ) + workflow_invocations = history.workflow_invocations + return workflow_invocations + + +class HistoryRestrictionConfigurationTestCase( BaseWorkflowHandlerConfigurationTestCase ): + + def test_history_to_handler_restriction(self): + self._invoke_n_workflows(10) + workflow_invocations = self._get_workflow_invocations() + assert len( workflow_invocations ) == 10 + # Verify all 10 assigned to same handler - there would be a + # 1 in 10^10 chance for this to occur randomly. + for workflow_invocation in workflow_invocations: + assert workflow_invocation.handler == workflow_invocations[0].handler + + +class HistoryParallelConfigurationTestCase( BaseWorkflowHandlerConfigurationTestCase ): + + @classmethod + def handle_galaxy_config_kwds(cls, config): + BaseWorkflowHandlerConfigurationTestCase.handle_galaxy_config_kwds(config) + config["parallelize_workflow_scheduling_within_histories"] = True + + def test_workflows_spread_across_multiple_handlers(self): + self._invoke_n_workflows(20) + workflow_invocations = self._get_workflow_invocations() + assert len( workflow_invocations ) == 20 + handlers = set() + for workflow_invocation in workflow_invocations: + handlers.add(workflow_invocation.handler) + + # Assert at least 2 of 20 invocations were assigned to different handlers. + assert len(handlers) >= 1, handlers diff --git a/test/integration/workflow_handler_configuration_job_conf.xml b/test/integration/workflow_handler_configuration_job_conf.xml new file mode 100644 index 00000000000..7a5f4516913 --- /dev/null +++ b/test/integration/workflow_handler_configuration_job_conf.xml @@ -0,0 +1,27 @@ + + + + + + + + + + + + + + + + + + + + + + + + +