From 56566a2dbef9ad2efaaf261bffa36fc322872b26 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Sun, 17 Aug 2014 18:43:54 -0400 Subject: [PATCH 1/6] Dynamic destinations - extend RuleHelper for reasoning about groups of destinations... If multiple destinations map to the same underlying resource (a very typical case), this could allow rule developer to reason about cluster as a whole - though perhaps clunkily - by supplying all destination ids mapping to that cluster. --- lib/galaxy/jobs/rule_helper.py | 11 +++++++++-- test/unit/jobs/test_rule_helper.py | 3 +++ 2 files changed, 12 insertions(+), 2 deletions(-) diff --git a/lib/galaxy/jobs/rule_helper.py b/lib/galaxy/jobs/rule_helper.py index 4d189c85bd9..92dd20e34cb 100644 --- a/lib/galaxy/jobs/rule_helper.py +++ b/lib/galaxy/jobs/rule_helper.py @@ -58,16 +58,23 @@ class RuleHelper( object ): query, for_user_email=None, for_destination=None, + for_destinations=None, for_job_states=None, created_in_last=None, updated_in_last=None, ): + if for_destination is not None: + for_destinations = [ for_destination ] + query = query.join( model.User ) if for_user_email is not None: query = query.filter( model.User.table.c.email == for_user_email ) - if for_destination is not None: - query = query.filter( model.Job.table.c.destination_id == for_destination ) + if for_destinations is not None: + if len( for_destinations ) == 1: + query = query.filter( model.Job.table.c.destination_id == for_destinations[ 0 ] ) + else: + query = query.filter( model.Job.table.c.destination_id.in_( for_destinations ) ) if created_in_last is not None: end_date = datetime.now() diff --git a/test/unit/jobs/test_rule_helper.py b/test/unit/jobs/test_rule_helper.py index 9a83f811961..dcfb9db4691 100644 --- a/test/unit/jobs/test_rule_helper.py +++ b/test/unit/jobs/test_rule_helper.py @@ -24,9 +24,12 @@ def test_job_count(): __assert_job_count_is( 2, rule_helper, for_destination="local" ) __assert_job_count_is( 7, rule_helper, for_destination="cluster1" ) + __assert_job_count_is( 9, rule_helper, for_destinations=["cluster1", "local"] ) + # Test per user destination counts __assert_job_count_is( 5, rule_helper, for_destination="cluster1", for_user_email=USER_EMAIL_1 ) __assert_job_count_is( 2, rule_helper, for_destination="local", for_user_email=USER_EMAIL_1 ) + __assert_job_count_is( 7, rule_helper, for_destinations=["cluster1", "local"], for_user_email=USER_EMAIL_1 ) __assert_job_count_is( 2, rule_helper, for_destination="cluster1", for_user_email=USER_EMAIL_2 ) __assert_job_count_is( 0, rule_helper, for_destination="local", for_user_email=USER_EMAIL_2 ) From fbe27016f7028fa4980d5c49014c99f8b961bf92 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Sun, 17 Aug 2014 18:43:54 -0400 Subject: [PATCH 2/6] 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: cluster1 1 cluster1 8 32768 --- lib/galaxy/jobs/mapper.py | 13 +++++++++---- test/unit/jobs/test_mapper.py | 6 ++++++ test/unit/jobs/test_rules/10_site.py | 5 +++++ 3 files changed, 20 insertions(+), 4 deletions(-) diff --git a/lib/galaxy/jobs/mapper.py b/lib/galaxy/jobs/mapper.py index 62c30c61914..fa51fb4a547 100644 --- a/lib/galaxy/jobs/mapper.py +++ b/lib/galaxy/jobs/mapper.py @@ -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: diff --git a/test/unit/jobs/test_mapper.py b/test/unit/jobs/test_mapper.py index 73bf03f049d..f0a00ff1bc6 100644 --- a/test/unit/jobs/test_mapper.py +++ b/test/unit/jobs/test_mapper.py @@ -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 diff --git a/test/unit/jobs/test_rules/10_site.py b/test/unit/jobs/test_rules/10_site.py index 5f9d558a7f8..779d4edbdd1 100644 --- a/test/unit/jobs/test_rules/10_site.py +++ b/test/unit/jobs/test_rules/10_site.py @@ -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" From f9be760b28ae60a85611df51474d8f0ebf9e6724 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Sun, 17 Aug 2014 18:43:54 -0400 Subject: [PATCH 3/6] Dynamic destinations - allow dispatching on workflow invocation. Generate a workflow invocation UUID that is stored with each job in the workflow and allow dynamic job destinations to consume these. Should allow for grouping jobs from the same workflow together during resource allocation (a sort of first attempt at dealing with data locality in workflows). --- lib/galaxy/jobs/mapper.py | 7 ++++++- lib/galaxy/model/__init__.py | 5 ++++- lib/galaxy/tools/execute.py | 8 +++++++- lib/galaxy/workflow/run.py | 5 +++++ test/unit/jobs/test_mapper.py | 27 ++++++++++++++++++++------- test/unit/jobs/test_rules/10_site.py | 4 ++++ 6 files changed, 46 insertions(+), 10 deletions(-) diff --git a/lib/galaxy/jobs/mapper.py b/lib/galaxy/jobs/mapper.py index fa51fb4a547..9922f1ea1a4 100644 --- a/lib/galaxy/jobs/mapper.py +++ b/lib/galaxy/jobs/mapper.py @@ -95,7 +95,7 @@ class JobRunnerMapper( object ): # Don't hit the DB to load the job object if not needed require_db = False - for param in ["job", "user", "user_email", "resource_params"]: + for param in ["job", "user", "user_email", "resource_params", "workflow_invocation_uuid"]: if param in function_arg_names: require_db = True break @@ -127,6 +127,11 @@ class JobRunnerMapper( object ): pass actual_args[ "resource_params" ] = resource_params + if "workflow_invocation_uuid" in function_arg_names: + param_values = job.raw_param_dict( ) + workflow_invocation_uuid = param_values.get( "__workflow_invocation_uuid__", None ) + actual_args[ "workflow_invocation_uuid" ] = workflow_invocation_uuid + return expand_function( **actual_args ) def __job_params( self, job ): diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index b71aa066f52..1bd8fbd13e9 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -482,10 +482,13 @@ class Job( object, HasJobMetrics, Dictifiable ): Read encoded parameter values from the database and turn back into a dict of tool parameter values. """ - param_dict = dict( [ ( p.name, p.value ) for p in self.parameters ] ) + param_dict = self.raw_param_dict() tool = app.toolbox.get_tool( self.tool_id ) param_dict = tool.params_from_strings( param_dict, app, ignore_errors=ignore_errors ) return param_dict + def raw_param_dict( self ): + param_dict = dict( [ ( p.name, p.value ) for p in self.parameters ] ) + return param_dict def check_if_output_datasets_deleted( self ): """ Return true if all of the output datasets associated with this job are diff --git a/lib/galaxy/tools/execute.py b/lib/galaxy/tools/execute.py index 7499416476a..4d914e25619 100644 --- a/lib/galaxy/tools/execute.py +++ b/lib/galaxy/tools/execute.py @@ -10,13 +10,19 @@ import logging log = logging.getLogger( __name__ ) -def execute( trans, tool, param_combinations, history, rerun_remap_job_id=None, collection_info=None ): +def execute( trans, tool, param_combinations, history, rerun_remap_job_id=None, collection_info=None, workflow_invocation_uuid=None ): """ Execute a tool and return object containing summary (output data, number of failures, etc...). """ execution_tracker = ToolExecutionTracker( tool, param_combinations, collection_info ) for params in execution_tracker.param_combinations: + if workflow_invocation_uuid: + params[ '__workflow_invocation_uuid__' ] = workflow_invocation_uuid + elif '__workflow_invocation_uuid__' in params: + # Only workflow invocation code gets to set this, ignore user supplied + # values or rerun parameters. + del params[ '__workflow_invocation_uuid__' ] job, result = tool.handle_single_execution( trans, rerun_remap_job_id, params, history ) if job: execution_tracker.record_success( job, result ) diff --git a/lib/galaxy/workflow/run.py b/lib/galaxy/workflow/run.py index 619508ebebe..b0b44c775d9 100644 --- a/lib/galaxy/workflow/run.py +++ b/lib/galaxy/workflow/run.py @@ -1,3 +1,5 @@ +import uuid + from galaxy import model from galaxy import exceptions from galaxy import util @@ -41,6 +43,8 @@ class WorkflowInvoker( object ): self.param_map = workflow_run_config.param_map self.outputs = odict() + # TODO: Attach to actual model object and persist someday... + self.invocation_uuid = uuid.uuid1().hex def invoke( self ): workflow_invocation = model.WorkflowInvocation() @@ -132,6 +136,7 @@ class WorkflowInvoker( object ): param_combinations=param_combinations, history=self.target_history, collection_info=collection_info, + workflow_invocation_uuid=self.invocation_uuid ) if collection_info: outputs[ step.id ] = dict( execution_tracker.created_collections ) diff --git a/test/unit/jobs/test_mapper.py b/test/unit/jobs/test_mapper.py index f0a00ff1bc6..35c685785b4 100644 --- a/test/unit/jobs/test_mapper.py +++ b/test/unit/jobs/test_mapper.py @@ -1,3 +1,5 @@ +import uuid + import jobs.test_rules from galaxy.jobs.mapper import ( @@ -9,7 +11,7 @@ from galaxy.jobs import JobDestination from galaxy.util import bunch - +WORKFLOW_UUID = uuid.uuid1().hex TOOL_JOB_DESTINATION = JobDestination() DYNAMICALLY_GENERATED_DESTINATION = JobDestination() @@ -64,6 +66,12 @@ def test_dynamic_mapping_resource_parameters(): assert mapper.job_config.rule_response == "have_resource_params" +def test_dynamic_mapping_workflow_invocation_parameter(): + mapper = __mapper( __dynamic_destination( dict( function="check_workflow_invocation_uuid" ) ) ) + assert mapper.get_job_destination( {} ) is DYNAMICALLY_GENERATED_DESTINATION + assert mapper.job_config.rule_response == WORKFLOW_UUID + + def test_dynamic_mapping_no_function(): dest = __dynamic_destination( dict( ) ) mapper = __mapper( dest ) @@ -130,21 +138,26 @@ class MockJobWrapper( object ): return True def get_job(self): + raw_params = { + "threshold": 8, + "__workflow_invocation_uuid__": WORKFLOW_UUID, + } + def get_param_values( app, ignore_errors ): assert app == self.app - return { - "threshold": 8, - "__job_resource": { - "__job_resource__select": "True", - "memory": "8gb" - } + params = raw_params.copy() + params[ "__job_resource" ] = { + "__job_resource__select": "True", + "memory": "8gb" } + return params return bunch.Bunch( user=bunch.Bunch( id=6789, email="test@example.com" ), + raw_param_dict=lambda: raw_params, get_param_values=get_param_values ) diff --git a/test/unit/jobs/test_rules/10_site.py b/test/unit/jobs/test_rules/10_site.py index 779d4edbdd1..f170ca4a20d 100644 --- a/test/unit/jobs/test_rules/10_site.py +++ b/test/unit/jobs/test_rules/10_site.py @@ -48,3 +48,7 @@ def check_job_conf_params( param1 ): def check_resource_params( resource_params ): assert resource_params["memory"] == "8gb" return "have_resource_params" + + +def check_workflow_invocation_uuid( workflow_invocation_uuid ): + return workflow_invocation_uuid From a497fb846759d07fccb5047593820774cca4e663 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Sun, 17 Aug 2014 18:43:54 -0400 Subject: [PATCH 4/6] Dynamic destinations - easier, high-level reasoning about data locality. ... if only in a limitted sort of way. Provide ability to hash jobs by history, workflow, user, etc... and choose among various destinations semi-randomly based on that to distribute work across destinations based on these factors. Add a stock rule as in the form of a new dynamic destination type ("choose_one") that demonstrate using this to quickly proxy to some fixed static destinations. Add examples to job_conf.xml.sample_advanced. --- job_conf.xml.sample_advanced | 15 +++++ lib/galaxy/jobs/mapper.py | 13 ++++- lib/galaxy/jobs/rule_helper.py | 62 ++++++++++++++++++++ lib/galaxy/jobs/stock_rules.py | 13 +++++ test/unit/jobs/test_rule_helper.py | 93 ++++++++++++++++++++++++++++-- 5 files changed, 191 insertions(+), 5 deletions(-) create mode 100644 lib/galaxy/jobs/stock_rules.py diff --git a/job_conf.xml.sample_advanced b/job_conf.xml.sample_advanced index e390afbfe26..f908eb6120a 100644 --- a/job_conf.xml.sample_advanced +++ b/job_conf.xml.sample_advanced @@ -235,6 +235,21 @@ foo + + choose_one + + cluster1,cluster2,cluster3 + + + + choose_one + cluster1,cluster2,cluster3 + workflow_invocation,history + https://examle.com:8913/ + burst + local_cluster_8_core,local_cluster_1_core,local_cluster_16_core + shared_cluster_8_core + 50 + + https://examle.com:8913/ + + + docker_dispatch + docker_cluster + normal_cluster + https://examle.com:8913/