diff --git a/taskflow/patterns/linear_flow.py b/taskflow/patterns/linear_flow.py index 6ab083871..47c9b1852 100644 --- a/taskflow/patterns/linear_flow.py +++ b/taskflow/patterns/linear_flow.py @@ -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) diff --git a/taskflow/patterns/resumption/logbook.py b/taskflow/patterns/resumption/logbook.py index a77cd6659..604522d2d 100644 --- a/taskflow/patterns/resumption/logbook.py +++ b/taskflow/patterns/resumption/logbook.py @@ -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, diff --git a/taskflow/utils.py b/taskflow/utils.py index e13324469..114c01ecf 100644 --- a/taskflow/utils.py +++ b/taskflow/utils.py @@ -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.