mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Dynamic destinations - pass job_conf.xml params to rule functions.
In other words, send extra job destination parameters to dynamic rule functions as arguments (in addition to those dynamically populated by Galaxy itself). This enables greater parameterization of rule functions and should lead to cleaner separation of logic and data (i.e. sites can program rules that restrict access to users, but which users can be populated at a higher level in `job_conf.xml`).
For example the following dynamic job rule:
def cluster1(app, memory="4096", cores="1", hours="48"):
native_spec = "--time=%s:00:00 --nodes=1 --ntasks=%s --mem=%s" % ( hours, cores, memory )
return JobDestination( "cluster1", params=dict( native_specification=native_spec ) )
Could then be called with various parameters in job_conf.xml as follows:
<destination id="short_job" type="dyanmic">
<param id="function">cluster1</param>
<param id="hours">1</param>
</destination>
<destination id="big_job" type="dynamic">
<param id="function">cluster1</param>
<param id="cores">8</param>
<param id="memory">32768</param>
</destination>
This commit is contained in:
@@ -69,7 +69,7 @@ class JobRunnerMapper( object ):
|
||||
names.append( rule_module_name )
|
||||
return names
|
||||
|
||||
def __invoke_expand_function( self, expand_function ):
|
||||
def __invoke_expand_function( self, expand_function, destination_params ):
|
||||
function_arg_names = inspect.getargspec( expand_function ).args
|
||||
app = self.job_wrapper.app
|
||||
possible_args = {
|
||||
@@ -83,6 +83,11 @@ class JobRunnerMapper( object ):
|
||||
|
||||
actual_args = {}
|
||||
|
||||
# Send through any job_conf.xml defined args to function
|
||||
for destination_param in destination_params.keys():
|
||||
if destination_param in function_arg_names:
|
||||
actual_args[ destination_param ] = destination_params[ destination_param ]
|
||||
|
||||
# Populate needed args
|
||||
for possible_arg_name in possible_args:
|
||||
if possible_arg_name in function_arg_names:
|
||||
@@ -179,12 +184,12 @@ class JobRunnerMapper( object ):
|
||||
raise Exception( message )
|
||||
|
||||
expand_function = self.__get_expand_function( expand_function_name )
|
||||
return self.__handle_rule( expand_function )
|
||||
return self.__handle_rule( expand_function, destination )
|
||||
else:
|
||||
raise Exception( "Unhandled dynamic job runner type specified - %s" % expand_type )
|
||||
|
||||
def __handle_rule( self, rule_function ):
|
||||
job_destination = self.__invoke_expand_function( rule_function )
|
||||
def __handle_rule( self, rule_function, destination ):
|
||||
job_destination = self.__invoke_expand_function( rule_function, destination.params )
|
||||
if not isinstance(job_destination, galaxy.jobs.JobDestination):
|
||||
job_destination_rep = str(job_destination) # Should be either id or url
|
||||
if '://' in job_destination_rep:
|
||||
|
||||
@@ -46,6 +46,12 @@ def test_dynamic_mapping_defaults_to_tool_id_as_rule():
|
||||
assert mapper.job_config.rule_response == "tool1_dest_id"
|
||||
|
||||
|
||||
def test_dynamic_mapping_job_conf_params():
|
||||
mapper = __mapper( __dynamic_destination( dict( function="check_job_conf_params", param1="7" ) ) )
|
||||
assert mapper.get_job_destination( {} ) is DYNAMICALLY_GENERATED_DESTINATION
|
||||
assert mapper.job_config.rule_response == "sent_7_dest_id"
|
||||
|
||||
|
||||
def test_dynamic_mapping_function_parameters():
|
||||
mapper = __mapper( __dynamic_destination( dict( function="check_rule_params" ) ) )
|
||||
assert mapper.get_job_destination( {} ) is DYNAMICALLY_GENERATED_DESTINATION
|
||||
|
||||
@@ -40,6 +40,11 @@ def check_rule_params(
|
||||
return "all_passed"
|
||||
|
||||
|
||||
def check_job_conf_params( param1 ):
|
||||
assert param1 == "7"
|
||||
return "sent_7_dest_id"
|
||||
|
||||
|
||||
def check_resource_params( resource_params ):
|
||||
assert resource_params["memory"] == "8gb"
|
||||
return "have_resource_params"
|
||||
|
||||
Reference in New Issue
Block a user