From 4beb98e973000e0bd51b06ad35301ddb8b905ec1 Mon Sep 17 00:00:00 2001 From: John Chilton Date: Thu, 23 Feb 2017 14:40:05 -0500 Subject: [PATCH] [17.01] By default, do not allow workflow invocations to schedule indefinitely. Give up after a month but allow admins to reduce this amount as well. --- config/galaxy.ini.sample | 4 ++ lib/galaxy/config.py | 1 + lib/galaxy/model/__init__.py | 5 ++ lib/galaxy/workflow/run.py | 11 +++++ test/api/test_workflows.py | 20 ++------ test/api/test_workflows_from_yaml.py | 1 + test/base/data/test_workflow_pause.ga | 2 +- test/base/populators.py | 22 ++++++++- ...st_maximum_worklfow_invocation_duration.py | 48 +++++++++++++++++++ 9 files changed, 96 insertions(+), 18 deletions(-) create mode 100644 test/integration/test_maximum_worklfow_invocation_duration.py diff --git a/config/galaxy.ini.sample b/config/galaxy.ini.sample index 336be53393b..70277c6f94c 100644 --- a/config/galaxy.ini.sample +++ b/config/galaxy.ini.sample @@ -1062,6 +1062,10 @@ use_interactive = True # collections. #force_beta_workflow_scheduled_for_collections=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. +#maximum_workflow_invocation_duration = 2678400 # Force serial scheduling of workflows within the context of a particular history #history_local_serial_workflow_scheduling=False diff --git a/lib/galaxy/config.py b/lib/galaxy/config.py index de66ec47a00..8c9fa513f08 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.maximum_workflow_invocation_duration = int( kwargs.get( "maximum_workflow_invocation_duration", 2678400 ) ) # Per-user Job concurrency limitations self.cache_user_job_count = string_as_bool( kwargs.get( 'cache_user_job_count', False ) ) diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index b3913fc12cb..caea87c346c 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -4063,6 +4063,11 @@ class WorkflowInvocation( object, Dictifiable ): return True return False + @property + def seconds_since_created( self ): + create_time = self.create_time or galaxy.model.orm.now.now() # In case not flushed yet + return (galaxy.model.orm.now.now() - create_time).total_seconds() + class WorkflowInvocationToSubworkflowInvocationAssociation( object, Dictifiable ): dict_collection_visible_keys = ( 'id', 'workflow_step_id', 'workflow_invocation_id', 'subworkflow_invocation_id' ) diff --git a/lib/galaxy/workflow/run.py b/lib/galaxy/workflow/run.py index 3afb3ed6414..9a3c51a8b6a 100644 --- a/lib/galaxy/workflow/run.py +++ b/lib/galaxy/workflow/run.py @@ -145,6 +145,17 @@ class WorkflowInvoker( object ): def invoke( self ): workflow_invocation = self.workflow_invocation + maximum_duration = getattr( self.trans.app.config, "maximum_workflow_invocation_duration", -1 ) + if maximum_duration > 0 and workflow_invocation.seconds_since_created > maximum_duration: + log.debug("Workflow invocation [%s] exceeded maximum number of seconds allowed for scheduling [%s], failing." % (workflow_invocation.id, maximum_duration)) + workflow_invocation.state = model.WorkflowInvocation.states.FAILED + # All jobs ran successfully, so we can save now + self.trans.sa_session.add( workflow_invocation ) + + # Not flushing in here, because web controller may create multiple + # invocations. + return self.progress.outputs + remaining_steps = self.progress.remaining_steps() delayed_steps = False for step in remaining_steps: diff --git a/test/api/test_workflows.py b/test/api/test_workflows.py index 1338fc46597..25c08ba2763 100644 --- a/test/api/test_workflows.py +++ b/test/api/test_workflows.py @@ -18,10 +18,6 @@ from base.populators import ( from galaxy.exceptions import error_codes from galaxy.tools.verify.test_data import TestDataResolver -from base.workflows_format_2 import ( - convert_and_import_workflow, - ImporterGalaxyInterface, -) SIMPLE_NESTED_WORKFLOW_YAML = """ class: GalaxyWorkflow @@ -69,7 +65,7 @@ test_data: """ -class BaseWorkflowsApiTestCase( api.ApiTestCase, ImporterGalaxyInterface ): +class BaseWorkflowsApiTestCase( api.ApiTestCase ): # TODO: Find a new file for this class. def setUp( self ): @@ -88,20 +84,12 @@ class BaseWorkflowsApiTestCase( api.ApiTestCase, ImporterGalaxyInterface ): names = [w[ "name" ] for w in index_response.json()] return names - # Import importer interface... def import_workflow(self, workflow, **kwds): - workflow_str = dumps(workflow, indent=4) - data = { - 'workflow': workflow_str, - } - data.update(**kwds) - upload_response = self._post( "workflows", data=data ) - self._assert_status_code_is( upload_response, 200 ) - return upload_response.json() + upload_response = self.workflow_populator.import_workflow(workflow, **kwds) + return upload_response def _upload_yaml_workflow(self, has_yaml, **kwds): - workflow = convert_and_import_workflow(has_yaml, galaxy_interface=self, **kwds) - return workflow[ "id" ] + return self.workflow_populator.upload_yaml_workflow(has_yaml, **kwds) def _setup_workflow_run( self, workflow, inputs_by='step_id', history_id=None ): uploaded_workflow_id = self.workflow_populator.create_workflow( workflow ) diff --git a/test/api/test_workflows_from_yaml.py b/test/api/test_workflows_from_yaml.py index cc269146a54..d1a41027156 100644 --- a/test/api/test_workflows_from_yaml.py +++ b/test/api/test_workflows_from_yaml.py @@ -279,6 +279,7 @@ steps: def _steps_by_label(self, workflow_as_dict): by_label = {} + assert "steps" in workflow_as_dict, workflow_as_dict for step in workflow_as_dict["steps"].values(): by_label[step['label']] = step return by_label diff --git a/test/base/data/test_workflow_pause.ga b/test/base/data/test_workflow_pause.ga index 6f5b4cd4580..02b9797ddd8 100644 --- a/test/base/data/test_workflow_pause.ga +++ b/test/base/data/test_workflow_pause.ga @@ -107,7 +107,7 @@ }, "post_job_actions": {}, "tool_errors": null, - "tool_id": "cat1", + "tool_id": "cat", "tool_state": "{\"__page__\": 0, \"__rerun_remap_job_id__\": null, \"input1\": \"null\", \"queries\": \"[]\"}", "tool_version": "1.0.0", "type": "tool", diff --git a/test/base/populators.py b/test/base/populators.py index 183d58c9903..fd3730f9ab7 100644 --- a/test/base/populators.py +++ b/test/base/populators.py @@ -9,6 +9,10 @@ from pkg_resources import resource_string from six import StringIO from base import api_asserts +from base.workflows_format_2 import ( + convert_and_import_workflow, + ImporterGalaxyInterface, +) # Simple workflow that takes an input and call cat wrapper on it. workflow_str = resource_string( __name__, "data/test_workflow_1.ga" ) @@ -253,6 +257,10 @@ class BaseWorkflowPopulator( object ): upload_response = self._post( "workflows/upload", data=data ) return upload_response + def upload_yaml_workflow(self, has_yaml, **kwds): + workflow = convert_and_import_workflow(has_yaml, galaxy_interface=self, **kwds) + return workflow[ "id" ] + def wait_for_invocation( self, workflow_id, invocation_id, timeout=DEFAULT_TIMEOUT ): url = "workflows/%s/usage/%s" % ( workflow_id, invocation_id ) return wait_on_state( lambda: self._get( url ), timeout=timeout ) @@ -264,7 +272,7 @@ class BaseWorkflowPopulator( object ): self.dataset_populator.wait_for_history( history_id, assert_ok=assert_ok, timeout=timeout ) -class WorkflowPopulator( BaseWorkflowPopulator ): +class WorkflowPopulator( BaseWorkflowPopulator, ImporterGalaxyInterface ): def __init__( self, galaxy_interactor ): self.galaxy_interactor = galaxy_interactor @@ -276,6 +284,18 @@ class WorkflowPopulator( BaseWorkflowPopulator ): def _get( self, route ): return self.galaxy_interactor.get( route ) + # Required for ImporterGalaxyInterface interface - so we can recurisvely import + # nested workflows. + def import_workflow(self, workflow, **kwds): + workflow_str = json.dumps(workflow, indent=4) + data = { + 'workflow': workflow_str, + } + data.update(**kwds) + upload_response = self._post( "workflows", data=data ) + assert upload_response.status_code == 200, upload_response + return upload_response.json() + class LibraryPopulator( object ): diff --git a/test/integration/test_maximum_worklfow_invocation_duration.py b/test/integration/test_maximum_worklfow_invocation_duration.py new file mode 100644 index 00000000000..b3544504352 --- /dev/null +++ b/test/integration/test_maximum_worklfow_invocation_duration.py @@ -0,0 +1,48 @@ +"""Integration tests for maximum workflow invocation duration configuration option.""" + +import time + +from json import dumps + +from base import integration_util +from base.populators import ( + DatasetPopulator, + WorkflowPopulator, +) + + +class MaximumWorkflowInvocationDurationTestCase(integration_util.IntegrationTestCase): + """Start a Pulsar job.""" + + framework_tool_and_types = True + + def setUp( self ): + super( MaximumWorkflowInvocationDurationTestCase, self ).setUp() + self.dataset_populator = DatasetPopulator( self.galaxy_interactor ) + self.workflow_populator = WorkflowPopulator( self.galaxy_interactor ) + + @classmethod + def handle_galaxy_config_kwds(cls, config): + config["maximum_workflow_invocation_duration"] = 20 + + def do_test(self): + workflow = self.workflow_populator.load_workflow_from_resource("test_workflow_pause") + workflow_id = self.workflow_populator.create_workflow(workflow) + history_id = self.dataset_populator.new_history() + 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) + invocation_response = self._post(url, data=request) + invocation_url = url + "/" + invocation_response.json()["id"] + time.sleep(5) + state = self._get(invocation_url).json()["state"] + assert state != "failed", state + time.sleep(35) + state = self._get(invocation_url).json()["state"] + assert state == "failed", state