diff --git a/lib/galaxy/jobs/__init__.py b/lib/galaxy/jobs/__init__.py index 71d17ec9928..86ac45c92b5 100644 --- a/lib/galaxy/jobs/__init__.py +++ b/lib/galaxy/jobs/__init__.py @@ -89,9 +89,16 @@ class JobWrapper( object ): self.__user_system_pwent = None self.__galaxy_system_pwent = None + def can_split( self ): + # Should the job handler split this job up? + return self.app.config.use_tasked_jobs and self.tool.parallelism + def get_job_runner_url( self ): return self.job_runner_mapper.get_job_runner_url( self.params ) + def get_parallelism(self): + return self.tool.parallelism + # legacy naming get_job_runner = get_job_runner_url @@ -941,6 +948,11 @@ class TaskWrapper(JobWrapper): self.prepare_input_files_cmds = None self.status = task.states.NEW + def can_split( self ): + # Should the job handler split this job up? TaskWrapper should + # always return False as the job has already been split. + return False + def get_job( self ): if self.job_id: return self.sa_session.query( model.Job ).get( self.job_id ) @@ -1160,3 +1172,20 @@ class NoopQueue( object ): return def shutdown( self ): return + +class ParallelismInfo(object): + """ + Stores the information (if any) for running multiple instances of the tool in parallel + on the same set of inputs. + """ + def __init__(self, tag): + self.method = tag.get('method') + if isinstance(tag, dict): + items = tag.iteritems() + else: + items = tag.attrib.items() + self.attributes = dict([item for item in items if item[0] != 'method' ]) + if len(self.attributes) == 0: + # legacy basic mode - provide compatible defaults + self.attributes['split_size'] = 20 + self.attributes['split_mode'] = 'number_of_parts' diff --git a/lib/galaxy/jobs/handler.py b/lib/galaxy/jobs/handler.py index c208c7e9ccf..c2464240c6e 100644 --- a/lib/galaxy/jobs/handler.py +++ b/lib/galaxy/jobs/handler.py @@ -460,7 +460,7 @@ class DefaultJobDispatcher( object ): log.debug( 'Loaded job runner: %s' % display_name ) def __get_runner_name( self, job_wrapper ): - if self.app.config.use_tasked_jobs and job_wrapper.tool.parallelism is not None and not isinstance(job_wrapper, TaskWrapper): + if job_wrapper.can_split(): runner_name = "tasks" else: runner_name = ( job_wrapper.get_job_runner_url().split(":", 1) )[0] @@ -478,7 +478,7 @@ class DefaultJobDispatcher( object ): job_wrapper.fail( failure_message ) return try: - if self.app.config.use_tasked_jobs and job_wrapper.tool.parallelism is not None and isinstance(job_wrapper, TaskWrapper): + if isinstance(job_wrapper, TaskWrapper): #DBTODO Refactor log.debug( "dispatching task %s, of job %d, to %s runner" %( job_wrapper.task_id, job_wrapper.job_id, runner_name ) ) else: diff --git a/lib/galaxy/jobs/runners/tasks.py b/lib/galaxy/jobs/runners/tasks.py index 3f05c3cfe45..8caa985d9f4 100644 --- a/lib/galaxy/jobs/runners/tasks.py +++ b/lib/galaxy/jobs/runners/tasks.py @@ -72,12 +72,13 @@ class TaskedJobRunner( object ): try: job_wrapper.change_state( model.Job.states.RUNNING ) self.sa_session.flush() - # Split with the tool-defined method. + # Split with the defined method. + parallelism = job_wrapper.get_parallelism() try: - splitter = getattr(__import__('galaxy.jobs.splitters', globals(), locals(), [job_wrapper.tool.parallelism.method]), job_wrapper.tool.parallelism.method) + splitter = getattr(__import__('galaxy.jobs.splitters', globals(), locals(), [parallelism.method]), parallelism.method) except: job_wrapper.change_state( model.Job.states.ERROR ) - job_wrapper.fail("Job Splitting Failed, no match for '%s'" % job_wrapper.tool.parallelism) + job_wrapper.fail("Job Splitting Failed, no match for '%s'" % parallelism) return tasks = splitter.do_split(job_wrapper) # Not an option for now. Task objects don't *do* anything diff --git a/lib/galaxy/jobs/splitters/basic.py b/lib/galaxy/jobs/splitters/basic.py index 5ec1bf24dfb..50d7151691e 100644 --- a/lib/galaxy/jobs/splitters/basic.py +++ b/lib/galaxy/jobs/splitters/basic.py @@ -5,8 +5,9 @@ log = logging.getLogger( __name__ ) def set_basic_defaults(job_wrapper): parent_job = job_wrapper.get_job() - job_wrapper.tool.parallelism.attributes['split_inputs'] = parent_job.input_datasets[0].name - job_wrapper.tool.parallelism.attributes['merge_outputs'] = job_wrapper.get_output_hdas_and_fnames().keys()[0] + parallelism = job_wrapper.get_parallelism() + parallelism.attributes['split_inputs'] = parent_job.input_datasets[0].name + parallelism.attributes['merge_outputs'] = job_wrapper.get_output_hdas_and_fnames().keys()[0] def do_split (job_wrapper): if len(job_wrapper.get_input_fnames()) > 1 or len(job_wrapper.get_output_fnames()) > 1: diff --git a/lib/galaxy/jobs/splitters/multi.py b/lib/galaxy/jobs/splitters/multi.py index 3c08baa6c9f..8cb8dc9a79a 100644 --- a/lib/galaxy/jobs/splitters/multi.py +++ b/lib/galaxy/jobs/splitters/multi.py @@ -9,7 +9,7 @@ def do_split (job_wrapper): parent_job = job_wrapper.get_job() working_directory = os.path.abspath(job_wrapper.working_directory) - parallel_settings = job_wrapper.tool.parallelism.attributes + parallel_settings = job_wrapper.get_parallelism().attributes # Syntax: split_inputs="input1,input2" shared_inputs="genome" # Designates inputs to be split or shared split_inputs=parallel_settings.get("split_inputs") @@ -92,7 +92,7 @@ def do_split (job_wrapper): def do_merge( job_wrapper, task_wrappers): - parallel_settings = job_wrapper.tool.parallelism.attributes + parallel_settings = job_wrapper.get_parallelism().attributes # Syntax: merge_outputs="export" pickone_outputs="genomesize" # Designates outputs to be merged, or selected from as a representative merge_outputs = parallel_settings.get("merge_outputs") diff --git a/lib/galaxy/tools/__init__.py b/lib/galaxy/tools/__init__.py index 7c56d36c30b..32405803b9c 100755 --- a/lib/galaxy/tools/__init__.py +++ b/lib/galaxy/tools/__init__.py @@ -15,6 +15,7 @@ from galaxy.util.odict import odict from galaxy.util.bunch import Bunch from galaxy.util.template import fill_template from galaxy import util, jobs, model +from galaxy.jobs import ParallelismInfo from elementtree import ElementTree from parameters import * from parameters.grouping import * @@ -802,19 +803,6 @@ class ToolRequirement( object ): self.type = type self.version = version -class ToolParallelismInfo(object): - """ - Stores the information (if any) for running multiple instances of the tool in parallel - on the same set of inputs. - """ - def __init__(self, tag): - self.method = tag.get('method') - self.attributes = dict([item for item in tag.attrib.items() if item[0] != 'method' ]) - if len(self.attributes) == 0: - # legacy basic mode - provide compatible defaults - self.attributes['split_size'] = 20 - self.attributes['split_mode'] = 'number_of_parts' - class Tool( object ): """ Represents a computational tool that can be executed through Galaxy. @@ -994,7 +982,7 @@ class Tool( object ): # Parallelism for tasks, read from tool config. parallelism = root.find("parallelism") if parallelism is not None and parallelism.get("method"): - self.parallelism = ToolParallelismInfo(parallelism) + self.parallelism = ParallelismInfo(parallelism) else: self.parallelism = None # Set job handler(s). Each handler is a dict with 'url' and, optionally, 'params'.