Pass runners instead of task objects/uuids.
Runners contain tasks and provide the identifying characteristics for the tasks, pass these along to the underlying components instead of tasks. Change-Id: I795a52d3a218e508c38ea209d757750ae6a936ea
This commit is contained in:
@@ -160,8 +160,7 @@ class Flow(base.Flow):
|
||||
self.task_notifier.notify(states.STARTED, details={
|
||||
'context': context,
|
||||
'flow': self,
|
||||
'task': runner.task,
|
||||
'task_uuid': runner.uuid,
|
||||
'runner': runner,
|
||||
})
|
||||
if not simulate_run:
|
||||
result = runner(context, *args, **kwargs)
|
||||
@@ -201,20 +200,17 @@ class Flow(base.Flow):
|
||||
self.task_notifier.notify(states.SUCCESS, details={
|
||||
'context': context,
|
||||
'flow': self,
|
||||
'result': result,
|
||||
'task': runner.task,
|
||||
'task_uuid': runner.uuid,
|
||||
'runner': runner,
|
||||
})
|
||||
except Exception as e:
|
||||
runner.result = e
|
||||
cause = utils.FlowFailure(runner, self, e)
|
||||
with excutils.save_and_reraise_exception():
|
||||
# Notify any listeners that the task has errored.
|
||||
self.task_notifier.notify(states.FAILURE, details={
|
||||
'context': context,
|
||||
'flow': self,
|
||||
'result': e,
|
||||
'task': runner.task,
|
||||
'task_uuid': runner.uuid,
|
||||
'runner': runner,
|
||||
})
|
||||
self.rollback(context, cause)
|
||||
|
||||
|
||||
@@ -36,24 +36,23 @@ class Resumption(object):
|
||||
def _task_listener(state, details):
|
||||
"""Store the result of the task under the given flow in the log
|
||||
book so that it can be retrieved later."""
|
||||
task_id = details['task_uuid']
|
||||
task = details['task']
|
||||
runner = details['runner']
|
||||
flow = details['flow']
|
||||
LOG.debug("Recording %s:%s of %s has finished state %s",
|
||||
utils.get_task_name(task), task_id, flow, state)
|
||||
LOG.debug("Recording %s of %s has finished state %s",
|
||||
runner, flow, state)
|
||||
# TODO(harlowja): switch to using uuids
|
||||
flow_id = flow.name
|
||||
metadata = {}
|
||||
flow_details = self._logbook[flow_id]
|
||||
if state in (states.SUCCESS, states.FAILURE):
|
||||
metadata['result'] = details['result']
|
||||
if task_id not in flow_details:
|
||||
metadata['result'] = runner.result
|
||||
if runner.uuid not in flow_details:
|
||||
metadata['states'] = [state]
|
||||
metadata['version'] = utils.get_task_version(task)
|
||||
flow_details.add_task(task_id, metadata)
|
||||
metadata['version'] = runner.version
|
||||
flow_details.add_task(runner.uuid, metadata)
|
||||
else:
|
||||
details = flow_details[task_id]
|
||||
immediate_version = utils.get_task_version(task)
|
||||
details = flow_details[runner.uuid]
|
||||
immediate_version = runner.version
|
||||
recorded_version = details.metadata.get('version')
|
||||
if recorded_version is not None:
|
||||
if not utils.is_version_compatible(recorded_version,
|
||||
|
||||
+9
-1
@@ -200,11 +200,19 @@ class Runner(object):
|
||||
self.runs_before = []
|
||||
self.result = None
|
||||
|
||||
@property
|
||||
def version(self):
|
||||
return get_task_version(self.task)
|
||||
|
||||
@property
|
||||
def name(self):
|
||||
return get_task_name(self.task)
|
||||
|
||||
def reset(self):
|
||||
self.result = None
|
||||
|
||||
def __str__(self):
|
||||
return "%s:%s" % (self.task, self.uuid)
|
||||
return "Runner %s: %s; %s" % (self.name, self.uuid, self.version)
|
||||
|
||||
def __call__(self, *args, **kwargs):
|
||||
# Find all of our inputs first.
|
||||
|
||||
Reference in New Issue
Block a user