mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Merge pull request #3820 from jmchilton/fixed_history_handler
[17.01] Restrict workflow scheduling within a history to a fixed, random handler.
This commit is contained in:
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
@@ -0,0 +1,27 @@
|
||||
<?xml version="1.0"?>
|
||||
<!--
|
||||
job_conf used by test_workflow_handler_configuration.py
|
||||
-->
|
||||
<job_conf>
|
||||
<plugins>
|
||||
<plugin id="local" type="runner" load="galaxy.jobs.runners.local:LocalJobRunner" workers="2"/>
|
||||
</plugins>
|
||||
|
||||
<handlers default="handlers">
|
||||
<handler id="handler0" tags="handlers"/>
|
||||
<handler id="handler1" tags="handlers"/>
|
||||
<handler id="handler2" tags="handlers" />
|
||||
<handler id="handler3" tags="handlers" />
|
||||
<handler id="handler4" tags="handlers" />
|
||||
<handler id="handler5" tags="handlers" />
|
||||
<handler id="handler6" tags="handlers" />
|
||||
<handler id="handler7" tags="handlers" />
|
||||
<handler id="handler8" tags="handlers" />
|
||||
<handler id="handler9" tags="handlers" />
|
||||
</handlers>
|
||||
|
||||
<destinations default="local">
|
||||
<destination id="local" runner="local">
|
||||
</destination>
|
||||
</destinations>
|
||||
</job_conf>
|
||||
Reference in New Issue
Block a user