diff --git a/Makefile b/Makefile
index 19306fe7..928be160 100644
--- a/Makefile
+++ b/Makefile
@@ -50,10 +50,6 @@ tests-ind:
tests-timing:
@make tests-ind 2>&1 | ./scripts/test_times.py
-.PHONY : waiting-task-stress
-waiting-task-stress:
- $(RUN) python -m unittest -v tests.SpiffWorkflow.bpmn.WaitingTaskStressBenchmark
-
wheel: clean
$(RUN) python -m build --sdist --wheel --outdir dist/
diff --git a/SpiffWorkflow/bpmn/parser/BpmnParser.py b/SpiffWorkflow/bpmn/parser/BpmnParser.py
index 3661c2bb..7b5cc082 100644
--- a/SpiffWorkflow/bpmn/parser/BpmnParser.py
+++ b/SpiffWorkflow/bpmn/parser/BpmnParser.py
@@ -45,11 +45,8 @@
BoundaryEvent,
EventBasedGateway
)
-from SpiffWorkflow.bpmn.specs.event_definitions.simple import NoneEventDefinition
-from SpiffWorkflow.bpmn.specs.event_definitions.timer import TimerEventDefinition
from SpiffWorkflow.bpmn.specs.event_definitions.message import CorrelationProperty
from SpiffWorkflow.bpmn.specs.mixins.subworkflow_task import SubWorkflowTask as SubWorkflowTaskMixin
-from SpiffWorkflow.bpmn.specs.mixins.events.start_event import StartEvent as StartEventMixin
from .ValidationException import ValidationException
from .ProcessParser import ProcessParser
@@ -448,23 +445,13 @@ def get_collaboration(self, name):
spec = BpmnProcessSpec(name)
subprocesses = {}
participant_type = self._get_parser_class(full_tag('callActivity'))[1]
- start_type = self._get_parser_class(full_tag('startEvent'))[1]
- end_type = self._get_parser_class(full_tag('endEvent'))[1]
- start = start_type(spec, 'Start Collaboration', NoneEventDefinition())
- spec.start.connect(start)
- end = end_type(spec, 'End Collaboration', NoneEventDefinition())
- end.connect(spec.end)
for process in self.collaborations[name]:
process_parser = self.get_process_parser(process)
if process_parser and process_parser.process_executable:
sp_spec = self.get_spec(process)
subprocesses[process] = sp_spec
subprocesses.update(self.get_subprocess_specs(process))
- if len([s for s in sp_spec.task_specs.values() if
- isinstance(s, StartEventMixin) and
- isinstance(s.event_definition, (NoneEventDefinition, TimerEventDefinition))
- ]):
- participant = participant_type(spec, process, process)
- start.connect(participant)
- participant.connect(end)
+ participant = participant_type(spec, process, process)
+ spec.start.connect_or_add_trigger(participant, sp_spec)
+
return spec, subprocesses
diff --git a/SpiffWorkflow/bpmn/parser/ProcessParser.py b/SpiffWorkflow/bpmn/parser/ProcessParser.py
index b04feb84..f5e72e30 100644
--- a/SpiffWorkflow/bpmn/parser/ProcessParser.py
+++ b/SpiffWorkflow/bpmn/parser/ProcessParser.py
@@ -191,6 +191,11 @@ def _parse(self):
self.spec.start.outputs = [split_task]
split_task.inputs = [self.spec.start]
+ for node in self.xpath("./bpmn:subProcess[@triggeredByEvent='true']"):
+ task_spec = self.parse_node(node)
+ sp_spec = self.parser.process_parsers[task_spec.spec].get_spec()
+ self.spec.start.connect_or_add_trigger(task_spec, sp_spec)
+
def parse_data_object(self, obj):
return self.create_data_spec(obj, DataObject)
diff --git a/SpiffWorkflow/bpmn/parser/TaskParser.py b/SpiffWorkflow/bpmn/parser/TaskParser.py
index d04f77ac..34e93b2b 100644
--- a/SpiffWorkflow/bpmn/parser/TaskParser.py
+++ b/SpiffWorkflow/bpmn/parser/TaskParser.py
@@ -23,7 +23,8 @@
from SpiffWorkflow.bpmn.specs.defaults import (
StandardLoopTask,
SequentialMultiInstanceTask,
- ParallelMultiInstanceTask
+ ParallelMultiInstanceTask,
+ EventSubprocess,
)
from SpiffWorkflow.bpmn.specs.control import BoundaryEventSplit, BoundaryEventJoin
from SpiffWorkflow.bpmn.specs.event_definitions.simple import CancelEventDefinition
@@ -47,6 +48,8 @@ class TaskParser(NodeParser):
STANDARD_LOOP_CLASS = StandardLoopTask
PARALLEL_MI_CLASS = ParallelMultiInstanceTask
SEQUENTIAL_MI_CLASS = SequentialMultiInstanceTask
+ # I have to add another attribute here. This parser is so stupid.
+ EVENT_SUBPROCESS_CLASS = EventSubprocess
def __init__(self, process_parser, spec_class, node, nsmap=None, lane=None):
"""
diff --git a/SpiffWorkflow/bpmn/parser/task_parsers.py b/SpiffWorkflow/bpmn/parser/task_parsers.py
index c5339cfc..2f252e0a 100644
--- a/SpiffWorkflow/bpmn/parser/task_parsers.py
+++ b/SpiffWorkflow/bpmn/parser/task_parsers.py
@@ -94,7 +94,11 @@ class SubWorkflowParser(TaskParser):
def create_task(self):
subworkflow_spec = SubprocessParser.get_subprocess_spec(self)
- return self.spec_class(self.spec, self.bpmn_id, subworkflow_spec=subworkflow_spec, **self.bpmn_attributes)
+ if self.attribute('triggeredByEvent'):
+ spec_class = self.EVENT_SUBPROCESS_CLASS
+ else:
+ spec_class = self.spec_class
+ return spec_class(self.spec, self.bpmn_id, subworkflow_spec=subworkflow_spec, **self.bpmn_attributes)
class CallActivityParser(TaskParser):
diff --git a/SpiffWorkflow/bpmn/serializer/config.py b/SpiffWorkflow/bpmn/serializer/config.py
index dc31382a..19807d3c 100644
--- a/SpiffWorkflow/bpmn/serializer/config.py
+++ b/SpiffWorkflow/bpmn/serializer/config.py
@@ -37,6 +37,7 @@
SubWorkflowTask,
CallActivity,
TransactionSubprocess,
+ EventSubprocess,
StartEvent,
EndEvent,
IntermediateCatchEvent,
@@ -96,6 +97,7 @@
EventConverter,
BoundaryEventConverter,
IOSpecificationConverter,
+ BpmnStartTaskConverter,
)
from .default.event_definition import (
TimerConditionalEventDefinitionConverter,
@@ -115,7 +117,7 @@
BpmnIoSpecification: IOSpecificationConverter,
BpmnProcessSpec: BpmnProcessSpecConverter,
SimpleBpmnTask: BpmnTaskSpecConverter,
- BpmnStartTask: BpmnTaskSpecConverter,
+ BpmnStartTask: BpmnStartTaskConverter,
_EndJoin: BpmnTaskSpecConverter,
NoneTask: BpmnTaskSpecConverter,
ManualTask: BpmnTaskSpecConverter,
@@ -127,6 +129,7 @@
SubWorkflowTask: SubWorkflowConverter,
CallActivity: SubWorkflowConverter,
TransactionSubprocess: SubWorkflowConverter,
+ EventSubprocess: SubWorkflowConverter,
BoundaryEventSplit: BpmnTaskSpecConverter,
BoundaryEventJoin: EventJoinConverter,
ExclusiveGateway: ExclusiveGatewayConverter,
diff --git a/SpiffWorkflow/bpmn/serializer/default/process_spec.py b/SpiffWorkflow/bpmn/serializer/default/process_spec.py
index 9d574d87..2c705972 100644
--- a/SpiffWorkflow/bpmn/serializer/default/process_spec.py
+++ b/SpiffWorkflow/bpmn/serializer/default/process_spec.py
@@ -36,6 +36,7 @@ def to_dict(self, spec):
'io_specification': self.registry.convert(spec.io_specification),
'data_objects': {name: self.registry.convert(obj) for name, obj in spec.data_objects.items()},
'correlation_keys': spec.correlation_keys,
+ 'bpmn_start_events': [ts.name for ts in spec.bpmn_start_events],
}
for name, task_spec in spec.task_specs.items():
task_dict = self.registry.convert(task_spec)
@@ -82,4 +83,6 @@ def from_dict(self, dct):
child_spec = spec.task_specs.get(task_spec.task_spec)
child_spec.completed_event.connect(task_spec.merge_child)
+ spec.bpmn_start_events = [spec.task_specs.get(name) for name in dct.get('bpmn_start_events', [])]
+
return spec
diff --git a/SpiffWorkflow/bpmn/serializer/default/task_spec.py b/SpiffWorkflow/bpmn/serializer/default/task_spec.py
index 392ff988..b45a8767 100644
--- a/SpiffWorkflow/bpmn/serializer/default/task_spec.py
+++ b/SpiffWorkflow/bpmn/serializer/default/task_spec.py
@@ -201,6 +201,18 @@ def to_dict(self, spec):
def from_dict(self, dct):
return self.task_spec_from_dict(dct)
+class BpmnStartTaskConverter(BpmnTaskSpecConverter):
+
+ def to_dict(self, spec):
+ dct = super().to_dict(spec)
+ dct['trigger_specs'] = self.registry.convert(spec.trigger_specs)
+ return dct
+
+ def from_dict(self, dct):
+ trigger_specs = dct.pop('trigger_specs', [])
+ spec = super().from_dict(dct)
+ spec.trigger_specs = trigger_specs
+ return spec
class EventConverter(BpmnTaskSpecConverter):
"""The default converter for BPMN events"""
diff --git a/SpiffWorkflow/bpmn/serializer/default/workflow.py b/SpiffWorkflow/bpmn/serializer/default/workflow.py
index 9284c1d1..214026d4 100644
--- a/SpiffWorkflow/bpmn/serializer/default/workflow.py
+++ b/SpiffWorkflow/bpmn/serializer/default/workflow.py
@@ -19,7 +19,9 @@
from uuid import UUID
+from SpiffWorkflow.task import TaskState
from SpiffWorkflow.bpmn.specs.mixins.subworkflow_task import SubWorkflowTask
+from SpiffWorkflow.bpmn.specs.mixins.events.event_types import CatchingEvent
from SpiffWorkflow.util.deep_merge import DeepMerge
from ..helpers.bpmn_converter import BpmnConverter
@@ -62,6 +64,8 @@ def from_dict(self, dct, workflow):
task.last_state_change = dct['last_state_change']
task.triggered = dct['triggered']
task.internal_data = self.registry.restore(dct['internal_data'])
+ if isinstance(task_spec, CatchingEvent) and task.has_state(TaskState.WAITING):
+ task.workflow.top_workflow.event_manager.add_task(task)
delta = dct.get('delta')
if delta and task.parent is not None:
@@ -101,6 +105,8 @@ def from_dict(self, dct, workflow):
task.triggered = dct['triggered']
task.internal_data = self.registry.restore(dct['internal_data'])
task.data = self.registry.restore(dct['data'])
+ if isinstance(task_spec, CatchingEvent) and task.has_state(TaskState.DEFINITE_MASK):
+ task.workflow.top_workflow.event_manager.add_task(task)
return task
class BpmnEventConverter(BpmnConverter):
@@ -201,7 +207,6 @@ def from_dict(self, dct):
# Handle the remaining top workflow attributes
self.subprocesses_from_dict(dct['subprocesses'], workflow)
workflow.bpmn_events = self.registry.restore(dct.pop('bpmn_events', []))
- workflow._rebuild_waiting_task_index()
return workflow
diff --git a/SpiffWorkflow/bpmn/specs/bpmn_process_spec.py b/SpiffWorkflow/bpmn/specs/bpmn_process_spec.py
index 2aec984a..bc7afc38 100644
--- a/SpiffWorkflow/bpmn/specs/bpmn_process_spec.py
+++ b/SpiffWorkflow/bpmn/specs/bpmn_process_spec.py
@@ -19,6 +19,7 @@
from SpiffWorkflow.specs.WorkflowSpec import WorkflowSpec
from SpiffWorkflow.bpmn.specs.control import _EndJoin, BpmnStartTask, SimpleBpmnTask
+from SpiffWorkflow.bpmn.specs.mixins.events.start_event import StartEvent
class BpmnProcessSpec(WorkflowSpec):
@@ -45,3 +46,10 @@ def __init__(self, name=None, description=None, filename=None, svg=None):
self.data_objects = {}
self.data_stores = {}
self.correlation_keys = {}
+ self.bpmn_start_events = []
+
+ def _add_notify(self, task_spec):
+ super()._add_notify(task_spec)
+ if isinstance(task_spec, (StartEvent, )):
+ self.bpmn_start_events.append(task_spec)
+
diff --git a/SpiffWorkflow/bpmn/specs/control.py b/SpiffWorkflow/bpmn/specs/control.py
index 2896a168..3a01db60 100644
--- a/SpiffWorkflow/bpmn/specs/control.py
+++ b/SpiffWorkflow/bpmn/specs/control.py
@@ -27,9 +27,35 @@
from SpiffWorkflow.bpmn.specs.mixins.events.intermediate_event import BoundaryEvent
from SpiffWorkflow.bpmn.specs.mixins.events.start_event import StartEvent
+from SpiffWorkflow.bpmn.specs.event_definitions.simple import NoneEventDefinition
+from SpiffWorkflow.bpmn.specs.event_definitions.timer import TimerEventDefinition
+
class BpmnStartTask(BpmnTaskSpec, StartTask):
- pass
+
+ def __init__(self, wf_spec, name, **kwargs):
+ super().__init__(wf_spec, name, **kwargs)
+ self.trigger_specs = []
+
+ def connect_or_add_trigger(self, task_spec, sp_spec):
+ for ts in sp_spec.bpmn_start_events:
+ if isinstance(ts.event_definition, (NoneEventDefinition, TimerEventDefinition)):
+ self.connect(task_spec)
+ task_spec.connect(self._wf_spec.end)
+ else:
+ self._wf_spec.task_specs[sp_spec.name] = task_spec
+ self.trigger_specs.append(sp_spec.name)
+
+ def trigger_wf(self, my_task, spec_name):
+ spec = self._wf_spec.task_specs.get(spec_name)
+ if self not in spec.inputs:
+ self.connect(spec)
+ spec.connect(self._wf_spec.end)
+ child = my_task._add_child(spec, TaskState.FUTURE)
+ child.triggered = True
+ child.task_spec._update(child)
+ return child.id
+
class SimpleBpmnTask(BpmnTaskSpec):
pass
diff --git a/SpiffWorkflow/bpmn/specs/defaults.py b/SpiffWorkflow/bpmn/specs/defaults.py
index 83c767c4..af39fd09 100644
--- a/SpiffWorkflow/bpmn/specs/defaults.py
+++ b/SpiffWorkflow/bpmn/specs/defaults.py
@@ -33,6 +33,7 @@
SubWorkflowTaskMixin,
CallActivityMixin,
TransactionSubprocessMixin,
+ EventSubprocessMixin,
StartEventMixin,
EndEventMixin,
IntermediateCatchEventMixin,
@@ -88,6 +89,9 @@ class CallActivity(CallActivityMixin, BpmnSpecMixin):
class TransactionSubprocess(TransactionSubprocessMixin, BpmnSpecMixin):
pass
+class EventSubprocess(EventSubprocessMixin, BpmnSpecMixin):
+ pass
+
class StartEvent(StartEventMixin, BpmnSpecMixin):
pass
diff --git a/SpiffWorkflow/bpmn/specs/mixins/__init__.py b/SpiffWorkflow/bpmn/specs/mixins/__init__.py
index 08452d08..26bbdf58 100644
--- a/SpiffWorkflow/bpmn/specs/mixins/__init__.py
+++ b/SpiffWorkflow/bpmn/specs/mixins/__init__.py
@@ -35,6 +35,7 @@
SubWorkflowTask as SubWorkflowTaskMixin,
CallActivity as CallActivityMixin,
TransactionSubprocess as TransactionSubprocessMixin,
+ EventSubprocess as EventSubprocessMixin,
)
from .events.start_event import StartEvent as StartEventMixin
@@ -46,4 +47,4 @@
EventBasedGateway as EventBasedGatewayMixin,
SendTask as SendTaskMixin,
ReceiveTask as ReceiveTaskMixin,
-)
\ No newline at end of file
+)
diff --git a/SpiffWorkflow/bpmn/specs/mixins/events/event_types.py b/SpiffWorkflow/bpmn/specs/mixins/events/event_types.py
index 3fbf3f75..fdcf9d68 100644
--- a/SpiffWorkflow/bpmn/specs/mixins/events/event_types.py
+++ b/SpiffWorkflow/bpmn/specs/mixins/events/event_types.py
@@ -58,8 +58,12 @@ def _update_hook(self, my_task):
elif my_task.state != TaskState.WAITING:
my_task._set_state(TaskState.WAITING)
+ my_task.workflow.top_workflow.event_manager.add_task(my_task)
self.event_definition.update_task(my_task)
+ def _on_ready_hook(self, my_task):
+ my_task.workflow.top_workflow.event_manager.remove_task(my_task)
+
def _run_hook(self, my_task):
self.event_definition.update_task_data(my_task)
@@ -72,6 +76,12 @@ def _predict_hook(self, my_task):
if not isinstance(self.event_definition, CycleTimerEventDefinition):
super()._predict_hook(my_task)
+ def _on_cancel(self, my_task):
+ for child in my_task:
+ child.workflow.top_workflow.event_manager.remove_task(child)
+ my_task.workflow.top_workflow.event_manager.remove_task(my_task)
+ super()._on_cancel(my_task)
+
class ThrowingEvent(TaskSpec):
"""Base Task Spec for Throwing Event nodes."""
diff --git a/SpiffWorkflow/bpmn/specs/mixins/events/intermediate_event.py b/SpiffWorkflow/bpmn/specs/mixins/events/intermediate_event.py
index c3e97a9f..62861154 100644
--- a/SpiffWorkflow/bpmn/specs/mixins/events/intermediate_event.py
+++ b/SpiffWorkflow/bpmn/specs/mixins/events/intermediate_event.py
@@ -54,9 +54,10 @@ def catches(self, my_task, event):
class EventBasedGateway(CatchingEvent):
def _predict_hook(self, my_task):
- my_task._sync_children(self.outputs, state=TaskState.MAYBE)
+ my_task._sync_children(self.outputs, state=TaskState.WAITING)
def _on_ready_hook(self, my_task):
for child in my_task.children:
if not child.internal_data.get('event_fired'):
child.cancel()
+
diff --git a/SpiffWorkflow/bpmn/specs/mixins/subworkflow_task.py b/SpiffWorkflow/bpmn/specs/mixins/subworkflow_task.py
index 9c3aea19..19b44349 100644
--- a/SpiffWorkflow/bpmn/specs/mixins/subworkflow_task.py
+++ b/SpiffWorkflow/bpmn/specs/mixins/subworkflow_task.py
@@ -131,3 +131,7 @@ class TransactionSubprocess(SubWorkflowTask):
def __init__(self, wf_spec, bpmn_id, subworkflow_spec, **kwargs):
super().__init__(wf_spec, bpmn_id, subworkflow_spec, True, **kwargs)
+
+class EventSubprocess(SubWorkflowTask):
+ pass
+
diff --git a/SpiffWorkflow/bpmn/util/event.py b/SpiffWorkflow/bpmn/util/event.py
index a885d037..cbd17c2f 100644
--- a/SpiffWorkflow/bpmn/util/event.py
+++ b/SpiffWorkflow/bpmn/util/event.py
@@ -1,3 +1,5 @@
+from SpiffWorkflow.task import TaskState
+
class BpmnEvent:
def __init__(self, event_definition, payload=None, correlations=None, target=None):
self.event_definition = event_definition
@@ -11,4 +13,58 @@ def __init__(self, name, event_type, value=None, correlations=None):
self.name = name
self.event_type = event_type
self.value = value
- self.correlations = correlations or {}
\ No newline at end of file
+ self.correlations = correlations or {}
+
+
+class EventManager:
+
+ def __init__(self, workflow):
+ self.workflow = workflow
+ self.tasks = {}
+
+ def add_task(self, my_task):
+ self.tasks[my_task.id] = my_task
+
+ def remove_task(self, my_task):
+ self.tasks.pop(my_task.id, None)
+
+ def get_waiting_tasks(self):
+ return [t.task_spec.event_definition.details(t) for t in self.tasks.values()]
+
+ def catch(self, event, internal=True):
+ if event.target is not None:
+ # This limits results to tasks in the specified workflow
+ tasks = [t for t in self.tasks.values() if t.workflow == event.target
+ and t.task_spec.event_definition.catches(t, event)]
+ if len(tasks) == 0:
+ event.target = event.target.parent_workflow
+ self.catch(event)
+ else:
+ self.update_collaboration(event)
+ tasks = [t for t in self.tasks.values() if t.task_spec.catches(t, event)]
+ # Figure out if we need to create an external event
+ if len(tasks) == 0 and internal:
+ self.workflow.bpmn_events.append(event)
+
+ for task in tasks:
+ task.task_spec.catch(task, event)
+ task.task_spec._update(task)
+
+ return len(tasks)
+
+ def update_collaboration(self, event):
+
+ def get_or_create_subprocess(task_spec, wf_spec):
+ for sp in self.workflow.subprocesses.values():
+ if sp.get_next_task(state=TaskState.WAITING, spec_name=task_spec) is not None:
+ return sp
+ child_id = self.workflow.spec.start.trigger_wf(self.workflow.task_tree, wf_spec)
+ return self.workflow.subprocesses[child_id]
+
+ # Start a subprocess for known specs with start events that catch this
+ for name in self.workflow.spec.start.trigger_specs:
+ sp_spec = self.workflow.subprocess_specs.get(name)
+ for ts in sp_spec.bpmn_start_events:
+ if ts.event_definition == event.event_definition:
+ subprocess = get_or_create_subprocess(ts.name, sp_spec.name)
+ subprocess.correlations.update(event.correlations)
diff --git a/SpiffWorkflow/bpmn/util/subworkflow.py b/SpiffWorkflow/bpmn/util/subworkflow.py
index addc17d0..06823ca4 100644
--- a/SpiffWorkflow/bpmn/util/subworkflow.py
+++ b/SpiffWorkflow/bpmn/util/subworkflow.py
@@ -34,12 +34,6 @@ def data_objects(self):
def get_tasks_iterator(self, first_task=None, **kwargs):
return BpmnTaskIterator(first_task or self.task_tree, **kwargs)
- def update_waiting_tasks(self):
- self.top_workflow._refresh_internal_waiting_tasks()
-
- def _task_state_changed_notify(self, task, old_state, new_state):
- self.top_workflow._waiting_task_state_changed(task, old_state, new_state)
-
class BpmnSubWorkflow(BpmnBaseWorkflow):
@@ -73,3 +67,4 @@ def collect_log_extras(self, dct=None):
dct = super().collect_log_extras()
dct.update({'parent_task_id': self.parent_task_id})
return dct
+
diff --git a/SpiffWorkflow/bpmn/workflow.py b/SpiffWorkflow/bpmn/workflow.py
index ca45c406..a676709c 100644
--- a/SpiffWorkflow/bpmn/workflow.py
+++ b/SpiffWorkflow/bpmn/workflow.py
@@ -17,119 +17,21 @@
# Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, MA
# 02110-1301 USA
-import heapq
-from datetime import datetime, timezone
+import warnings
from SpiffWorkflow.task import Task
from SpiffWorkflow.util.task import TaskState
from SpiffWorkflow.exceptions import WorkflowException
-from SpiffWorkflow.bpmn.specs.mixins.events.event_types import CatchingEvent
-from SpiffWorkflow.bpmn.specs.mixins.events.start_event import StartEvent
-from SpiffWorkflow.bpmn.specs.mixins.subworkflow_task import CallActivity
-from SpiffWorkflow.bpmn.specs.event_definitions.multiple import MultipleEventDefinition
-from SpiffWorkflow.bpmn.specs.event_definitions.timer import TimerEventDefinition
-from SpiffWorkflow.bpmn.specs.event_definitions.item_aware_event import CodeEventDefinition
-
from SpiffWorkflow.bpmn.specs.control import BoundaryEventSplit
+from SpiffWorkflow.bpmn.specs.event_definitions.timer import TimerEventDefinition
from SpiffWorkflow.bpmn.util.subworkflow import BpmnBaseWorkflow, BpmnSubWorkflow
+from SpiffWorkflow.bpmn.util.event import EventManager
from .script_engine.python_engine import PythonScriptEngine
-class _WaitingTaskIndex:
-
- def __init__(self):
- self.waiting_tasks = {}
- self.timer_tasks = {}
- self.timer_due_at = {}
- self.timer_heap = []
- self._sequence = 0
-
- def task_state_changed(self, task, old_state, new_state):
- if old_state == TaskState.WAITING:
- self._remove(task)
- if new_state == TaskState.WAITING:
- self._add(task)
-
- def refresh_internal_tasks(self, refresh_task):
- for task in list(self.waiting_tasks.values()):
- if task.id not in self.timer_tasks and task.state == TaskState.WAITING:
- refresh_task(task)
- self.refresh_due_timers(refresh_task)
-
- def refresh_due_timers(self, refresh_task):
- self._schedule_missing_timer_due_times(refresh_task)
- now = datetime.now(timezone.utc).timestamp()
- while self.timer_heap and self.timer_heap[0][0] <= now:
- due_at, _sequence, task_id = heapq.heappop(self.timer_heap)
- if self.timer_due_at.get(task_id) != due_at:
- continue
- task = self.timer_tasks.get(task_id)
- if task is None or task.state != TaskState.WAITING:
- continue
- refresh_task(task)
-
- def refresh_tasks(self, tasks, refresh_task):
- for task in tasks:
- refresh_task(task)
-
- def reschedule_timer(self, task):
- if task.id in self.timer_tasks:
- self._schedule_timer(task)
-
- def _add(self, task):
- self.waiting_tasks[task.id] = task
- if self._is_timer_task(task):
- self.timer_tasks[task.id] = task
- self._schedule_timer(task)
-
- def _remove(self, task):
- self.waiting_tasks.pop(task.id, None)
- self.timer_tasks.pop(task.id, None)
- self.timer_due_at.pop(task.id, None)
-
- def _schedule_missing_timer_due_times(self, refresh_task):
- for task in list(self.timer_tasks.values()):
- if task.state != TaskState.WAITING:
- continue
- if task.id not in self.timer_due_at:
- refresh_task(task)
-
- def _schedule_timer(self, task):
- due_at = self._get_timer_due_at(task)
- if due_at is None:
- self.timer_due_at.pop(task.id, None)
- return
- self.timer_due_at[task.id] = due_at
- self._sequence += 1
- heapq.heappush(self.timer_heap, (due_at, self._sequence, task.id))
-
- def _is_timer_task(self, task):
- return isinstance(task.task_spec, CatchingEvent) and self._has_timer_definition(task.task_spec.event_definition)
-
- def _has_timer_definition(self, event_definition):
- if isinstance(event_definition, TimerEventDefinition):
- return True
- if isinstance(event_definition, MultipleEventDefinition):
- return any(self._has_timer_definition(definition) for definition in event_definition.event_definitions)
- return False
-
- def _get_timer_due_at(self, task):
- event_value = task._get_internal_data('event_value')
- if event_value is None:
- return None
- if isinstance(event_value, dict):
- if event_value.get('cycles') == 0:
- return 0
- next_event = event_value.get('next')
- if next_event is None:
- return None
- return TimerEventDefinition.get_datetime(next_event).timestamp()
- return TimerEventDefinition.get_datetime(event_value).timestamp()
-
-
class BpmnWorkflow(BpmnBaseWorkflow):
"""
The engine that executes a BPMN workflow. This specialises the standard
@@ -148,8 +50,7 @@ def __init__(self, spec, subprocess_specs=None, script_engine=None, **kwargs):
self.subprocesses = {}
self.bpmn_events = []
self.correlations = {}
- self._waiting_task_index = _WaitingTaskIndex()
- self._refreshing_waiting_tasks = False
+ self.event_manager = EventManager(self)
super().__init__(spec, **kwargs)
for obj in self.spec.data_objects:
@@ -204,44 +105,13 @@ def get_active_subprocesses(self):
return [sp for sp in self.subprocesses.values() if not sp.completed]
def catch(self, event):
- """
- Tasks can always catch events, regardless of their state. The event information is stored in the task's
- internal data and processed when the task is reached in the workflow. If a task should only receive messages
- while it is running (eg a boundary event), the task should call the event_definition's reset method before
- executing to clear out a stale message.
-
- :param event: the thrown event
- """
- if event.target is not None:
- # This limits results to tasks in the specified workflow
- tasks = event.target.get_tasks(skip_subprocesses=True, state=TaskState.NOT_FINISHED_MASK, catches_event=event)
- if isinstance(event.event_definition, CodeEventDefinition) and len(tasks) == 0:
- event.target = event.target.parent_workflow
- self.catch(event)
- else:
- self.update_collaboration(event)
- tasks = self.get_tasks(state=TaskState.NOT_FINISHED_MASK, catches_event=event)
- # Figure out if we need to create an external event
- if len(tasks) == 0:
- self.bpmn_events.append(event)
-
- for task in tasks:
- task.task_spec.catch(task, event)
- if len(tasks) > 0:
- self._refresh_caught_tasks(tasks)
+ self.event_manager.catch(event)
def send_event(self, event):
"""Allows this workflow to catch an externally generated event."""
- if event.target is not None:
- self.catch(event)
- else:
- tasks = self.get_tasks(state=TaskState.NOT_FINISHED_MASK, catches_event=event)
- if len(tasks) == 0:
- raise WorkflowException(f"This process is not waiting for {event.event_definition.name}")
- for task in tasks:
- task.task_spec.catch(task, event)
- self._refresh_caught_tasks(tasks)
+ if self.event_manager.catch(event, internal=False) == 0:
+ raise WorkflowException(f"This process is not waiting for {event.event_definition.name}")
def get_events(self):
"""Returns the list of events that cannot be handled from within this workflow."""
@@ -250,8 +120,7 @@ def get_events(self):
return events
def waiting_events(self):
- iter = self.get_tasks_iterator(state=TaskState.WAITING, spec_class=CatchingEvent)
- return [t.task_spec.event_definition.details(t) for t in iter]
+ return self.event_manager.get_waiting_tasks()
def do_engine_steps(self, will_complete_task=None, did_complete_task=None):
"""
@@ -263,11 +132,10 @@ def do_engine_steps(self, will_complete_task=None, did_complete_task=None):
:param will_complete_task: Callback that will be called prior to completing a task
:param did_complete_task: Callback that will be called after completing a task
"""
- self._refresh_due_waiting_tasks()
count = self._do_engine_steps(will_complete_task, did_complete_task)
while count > 0:
- self._refresh_due_waiting_tasks()
count = self._do_engine_steps(will_complete_task, did_complete_task)
+ self.refresh_timers()
def _do_engine_steps(self, will_complete_task=None, did_complete_task=None):
@@ -298,62 +166,23 @@ def update_workflow(wf):
def refresh_waiting_tasks(self, will_refresh_task=None, did_refresh_task=None):
"""
- Compatibility no-op.
-
- BPMN workflows now refresh WAITING task internals through engine steps,
- targeted event catches, and task completion notifications.
+ Refresh the state of all WAITING tasks. This will, for example, update
+ Catching Timer Events whose waiting time has passed.
:param will_refresh_task: Callback that will be called prior to refreshing a task
:param did_refresh_task: Callback that will be called after refreshing a task
"""
- pass
-
- def refresh_due_waiting_tasks(self):
- """Refresh WAITING timer tasks that are currently due."""
- self._refresh_due_waiting_tasks()
-
- def _waiting_task_state_changed(self, task, old_state, new_state):
- self._waiting_task_index.task_state_changed(task, old_state, new_state)
-
- def _rebuild_waiting_task_index(self):
- self._waiting_task_index = _WaitingTaskIndex()
- workflows = [self] + list(self.subprocesses.values())
- for workflow in workflows:
- for task in workflow.tasks.values():
- if task.state == TaskState.WAITING:
- self._waiting_task_index.task_state_changed(task, None, TaskState.WAITING)
-
- def _refresh_internal_waiting_tasks(self):
- if self._refreshing_waiting_tasks:
- return
- self._refreshing_waiting_tasks = True
- try:
- self._waiting_task_index.refresh_internal_tasks(self._refresh_waiting_task)
- finally:
- self._refreshing_waiting_tasks = False
-
- def _refresh_due_waiting_tasks(self):
- if self._refreshing_waiting_tasks:
- return
- self._refreshing_waiting_tasks = True
- try:
- self._waiting_task_index.refresh_due_timers(self._refresh_waiting_task)
- finally:
- self._refreshing_waiting_tasks = False
-
- def _refresh_caught_tasks(self, tasks):
- if self._refreshing_waiting_tasks:
- return
- self._refreshing_waiting_tasks = True
- try:
- self._waiting_task_index.refresh_tasks(tasks, self._refresh_waiting_task)
- finally:
- self._refreshing_waiting_tasks = False
-
- def _refresh_waiting_task(self, task):
- if task.state == TaskState.WAITING:
- task.task_spec._update(task)
- self._waiting_task_index.reschedule_timer(task)
+ warnings.warn(
+ DeprecationWarning(f'BpmnWorkflow.refresh_waiting_tasks will be removed in future versions; use refresh_timers')
+ )
+ self.refresh_timers()
+
+ def refresh_timers(self):
+ # Ideally this would go in event manager but I can't import the necessary classes there
+ # Eventually I'll move it
+ for task in list(self.event_manager.tasks.values()):
+ if isinstance(task.task_spec.event_definition, (TimerEventDefinition, )):
+ task.task_spec._update(task)
def get_task_from_id(self, task_id):
if task_id not in self.tasks:
@@ -406,36 +235,4 @@ def cancel(self, workflow=None):
return cancelled
- def update_collaboration(self, event):
-
- def get_or_create_subprocess(task_spec, wf_spec):
-
- for sp in self.subprocesses.values():
- if sp.get_next_task(state=TaskState.WAITING, spec_name=task_spec.name) is not None:
- return sp
-
- # This creates a new task associated with a process when an event that kicks of a process is received
- # I need to know what class is being used to create new processes in this case, and this seems slightly
- # less bad than adding yet another argument. Still sucks though.
- # TODO: Make collaborations a class rather than trying to shoehorn them into a process.
- for spec in self.spec.task_specs.values():
- if isinstance(spec, CallActivity):
- spec_class = spec.__class__
- break
- else:
- # Default to the mixin class, which will probably fail in many cases.
- spec_class = CallActivity
-
- new = spec_class(self.spec, f'{wf_spec.name}_{len(self.subprocesses)}', wf_spec.name)
- self.spec.start.connect(new)
- task = Task(self, new, parent=self.task_tree)
- # This (indirectly) calls create_subprocess
- task.task_spec._update(task)
- return self.subprocesses[task.id]
-
- # Start a subprocess for known specs with start events that catch this
- for spec in self.subprocess_specs.values():
- for task_spec in spec.task_specs.values():
- if isinstance(task_spec, StartEvent) and task_spec.event_definition == event.event_definition:
- subprocess = get_or_create_subprocess(task_spec, spec)
- subprocess.correlations.update(event.correlations)
+
diff --git a/SpiffWorkflow/spiff/parser/task_spec.py b/SpiffWorkflow/spiff/parser/task_spec.py
index e7467b79..c67713b5 100644
--- a/SpiffWorkflow/spiff/parser/task_spec.py
+++ b/SpiffWorkflow/spiff/parser/task_spec.py
@@ -29,6 +29,7 @@
SequentialMultiInstanceTask,
BusinessRuleTask,
UserTask,
+ EventSubprocess,
)
SPIFFWORKFLOW_NSMAP = {'spiffworkflow': 'http://spiffworkflow.org/bpmn/schema/1.0/core'}
@@ -39,6 +40,7 @@ class SpiffTaskParser(TaskParser):
STANDARD_LOOP_CLASS = StandardLoopTask
PARALLEL_MI_CLASS = ParallelMultiInstanceTask
SEQUENTIAL_MI_CLASS = SequentialMultiInstanceTask
+ EVENT_SUBPROCESS_CLASS = EventSubprocess
def parse_extensions(self, node=None):
if node is None:
@@ -154,13 +156,18 @@ def create_task(self):
prescript = extensions.get('preScript')
postscript = extensions.get('postScript')
subworkflow_spec = SubprocessParser.get_subprocess_spec(self)
- return self.spec_class(
+ if self.attribute('triggeredByEvent'):
+ spec_class = self.EVENT_SUBPROCESS_CLASS
+ else:
+ spec_class = self.spec_class
+ return spec_class(
self.spec,
self.bpmn_id,
subworkflow_spec=subworkflow_spec,
prescript=prescript,
postscript=postscript,
- **self.bpmn_attributes)
+ **self.bpmn_attributes
+ )
class ScriptTaskParser(SpiffTaskParser):
diff --git a/SpiffWorkflow/spiff/serializer/config.py b/SpiffWorkflow/spiff/serializer/config.py
index 07159baf..1d2012b3 100644
--- a/SpiffWorkflow/spiff/serializer/config.py
+++ b/SpiffWorkflow/spiff/serializer/config.py
@@ -29,6 +29,7 @@
ScriptTask as DefaultScriptTask,
SubWorkflowTask as DefaultSubWorkflowTask,
TransactionSubprocess as DefaultTransactionSubprocess,
+ EventSubprocess as DefaultEventSubprocess,
CallActivity as DefaultCallActivity,
StandardLoopTask as DefaultStandardLoopTask,
ParallelMultiInstanceTask as DefaultParallelMultiInstanceTask,
@@ -46,6 +47,7 @@
ServiceTask,
SubWorkflowTask,
TransactionSubprocess,
+ EventSubprocess,
CallActivity,
StandardLoopTask,
ParallelMultiInstanceTask,
@@ -88,6 +90,7 @@
SPIFF_CONFIG.pop(DefaultReceiveTask)
SPIFF_CONFIG.pop(DefaultSubWorkflowTask)
SPIFF_CONFIG.pop(DefaultTransactionSubprocess)
+SPIFF_CONFIG.pop(DefaultEventSubprocess)
SPIFF_CONFIG.pop(DefaultCallActivity)
SPIFF_CONFIG.pop(DefaultStandardLoopTask)
SPIFF_CONFIG.pop(DefaultParallelMultiInstanceTask)
@@ -104,6 +107,7 @@
SPIFF_CONFIG[SubWorkflowTask] = SubWorkflowTaskConverter
SPIFF_CONFIG[CallActivity] = SubWorkflowTaskConverter
SPIFF_CONFIG[TransactionSubprocess] = SubWorkflowTaskConverter
+SPIFF_CONFIG[EventSubprocess] = SubWorkflowTaskConverter
SPIFF_CONFIG[ParallelMultiInstanceTask] = SpiffMultiInstanceConverter
SPIFF_CONFIG[SequentialMultiInstanceTask] = SpiffMultiInstanceConverter
SPIFF_CONFIG[StandardLoopTask] = StandardLoopTaskConverter
diff --git a/SpiffWorkflow/spiff/specs/defaults.py b/SpiffWorkflow/spiff/specs/defaults.py
index 877a2baf..9809f381 100644
--- a/SpiffWorkflow/spiff/specs/defaults.py
+++ b/SpiffWorkflow/spiff/specs/defaults.py
@@ -24,6 +24,7 @@
SubWorkflowTaskMixin,
CallActivityMixin,
TransactionSubprocessMixin,
+ EventSubprocessMixin,
StandardLoopTaskMixin,
ParallelMultiInstanceTaskMixin,
SequentialMultiInstanceTaskMixin,
@@ -76,5 +77,8 @@ class CallActivity(CallActivityMixin, SpiffBpmnTask):
class TransactionSubprocess(TransactionSubprocessMixin, SpiffBpmnTask):
pass
+class EventSubprocess(EventSubprocessMixin, SpiffBpmnTask):
+ pass
+
class ServiceTask(ServiceTaskMixin, SpiffBpmnTask):
pass
diff --git a/SpiffWorkflow/task.py b/SpiffWorkflow/task.py
index a7838b19..3b602db1 100644
--- a/SpiffWorkflow/task.py
+++ b/SpiffWorkflow/task.py
@@ -305,12 +305,9 @@ def _set_state(self, value: int) -> None:
"""Force set the state on a task"""
if value != self.state:
- old_state = self._state
elapsed = time.time() - self.last_state_change
self.last_state_change = time.time()
self._state = value
- if hasattr(self.workflow, '_task_state_changed_notify'):
- self.workflow._task_state_changed_notify(self, old_state, value)
logger.info(
f'State changed to {TaskState.get_name(value)}',
extra=self.collect_log_extras({'elapsed': elapsed})
diff --git a/SpiffWorkflow/workflow.py b/SpiffWorkflow/workflow.py
index 17ba3f46..3a0efd2a 100644
--- a/SpiffWorkflow/workflow.py
+++ b/SpiffWorkflow/workflow.py
@@ -293,8 +293,6 @@ def _remove_task(self, task_id: UUID) -> None:
task = self.tasks[task_id]
for child in task.children:
self._remove_task(child.id)
- if hasattr(self, '_task_state_changed_notify'):
- self._task_state_changed_notify(task, task.state, None)
task.parent._children.remove(task.id)
self.tasks.pop(task_id)
diff --git a/doc/bpmn/supported.rst b/doc/bpmn/supported.rst
index 523fdf4a..6f05b7ca 100644
--- a/doc/bpmn/supported.rst
+++ b/doc/bpmn/supported.rst
@@ -29,6 +29,7 @@ Subrocesses and Call Activities
* Subprocess
* Call Activity
* Transaction Subprocess
+* Event Subprocess
Events
------
diff --git a/tests/SpiffWorkflow/bpmn/BpmnWorkflowTestCase.py b/tests/SpiffWorkflow/bpmn/BpmnWorkflowTestCase.py
index 5ac1d27d..3798ae9c 100644
--- a/tests/SpiffWorkflow/bpmn/BpmnWorkflowTestCase.py
+++ b/tests/SpiffWorkflow/bpmn/BpmnWorkflowTestCase.py
@@ -141,5 +141,5 @@ def restore(self, state):
def _get_workflow_state(self, do_steps=True):
if do_steps:
self.workflow.do_engine_steps()
- self.workflow.refresh_waiting_tasks()
+ self.workflow.refresh_timers()
return self.serializer.to_dict(self.workflow)
diff --git a/tests/SpiffWorkflow/bpmn/CollaborationTest.py b/tests/SpiffWorkflow/bpmn/CollaborationTest.py
index c42b0705..61cbecc5 100644
--- a/tests/SpiffWorkflow/bpmn/CollaborationTest.py
+++ b/tests/SpiffWorkflow/bpmn/CollaborationTest.py
@@ -76,9 +76,8 @@ def testBpmnMessage(self):
def testCorrelation(self):
- specs = self.get_all_specs('correlation.bpmn')
- proc_1 = specs['proc_1']
- self.workflow = BpmnWorkflow(proc_1, specs)
+ spec, subprocesses = self.load_collaboration('correlation.bpmn', 'correlation_test')
+ self.workflow = BpmnWorkflow(spec, subprocesses)
self.workflow.do_engine_steps()
for idx, task in enumerate(self.get_ready_user_tasks()):
task.data['task_num'] = idx
@@ -102,9 +101,8 @@ def testCorrelation(self):
def testTwoCorrelationKeys(self):
- specs = self.get_all_specs('correlation_two_conversations.bpmn')
- proc_1 = specs['proc_1']
- self.workflow = BpmnWorkflow(proc_1, specs)
+ spec, subprocesses = self.load_collaboration('correlation_two_conversations.bpmn', 'correlation_test')
+ self.workflow = BpmnWorkflow(spec, subprocesses)
self.workflow.do_engine_steps()
for idx, task in enumerate(self.get_ready_user_tasks()):
task.data['task_num'] = idx
diff --git a/tests/SpiffWorkflow/bpmn/ParallelMultiInstanceTest.py b/tests/SpiffWorkflow/bpmn/ParallelMultiInstanceTest.py
index 14415343..ef306a14 100644
--- a/tests/SpiffWorkflow/bpmn/ParallelMultiInstanceTest.py
+++ b/tests/SpiffWorkflow/bpmn/ParallelMultiInstanceTest.py
@@ -39,7 +39,6 @@ def set_io_and_run_workflow(self, data, data_input=None, data_output=None, save_
if save_restore:
self.save_restore()
ready_tasks = self.get_ready_user_tasks()
- self.workflow.refresh_waiting_tasks()
self.workflow.do_engine_steps()
any_task = self.workflow.get_next_task(spec_name='any_task')
@@ -64,7 +63,6 @@ def run_workflow_with_condition(self, data):
task.data['output_item'] = task.data['input_item'] * 2
task.run()
self.workflow.do_engine_steps()
- self.workflow.refresh_waiting_tasks()
self.assertTrue(self.workflow.completed)
self.assertEqual(len([ t for t in ready_tasks if t.state == TaskState.CANCELLED]), 2)
diff --git a/tests/SpiffWorkflow/bpmn/ResetTokenOnBoundaryEventTest.py b/tests/SpiffWorkflow/bpmn/ResetTokenOnBoundaryEventTest.py
index 1ae3a8eb..99aad8ed 100644
--- a/tests/SpiffWorkflow/bpmn/ResetTokenOnBoundaryEventTest.py
+++ b/tests/SpiffWorkflow/bpmn/ResetTokenOnBoundaryEventTest.py
@@ -83,7 +83,6 @@ def advance_to_task(self, name):
while ready_tasks[0].task_spec.name != name:
ready_tasks[0].run()
self.workflow.do_engine_steps()
- self.workflow.refresh_waiting_tasks()
ready_tasks = self.workflow.get_tasks(state=TaskState.READY)
def complete_workflow(self):
@@ -92,5 +91,4 @@ def complete_workflow(self):
while len(ready_tasks) > 0:
ready_tasks[0].run()
self.workflow.do_engine_steps()
- self.workflow.refresh_waiting_tasks()
ready_tasks = self.workflow.get_tasks(state=TaskState.READY)
diff --git a/tests/SpiffWorkflow/bpmn/SequentialMultiInstanceTest.py b/tests/SpiffWorkflow/bpmn/SequentialMultiInstanceTest.py
index c5dcd5ad..9657cf79 100644
--- a/tests/SpiffWorkflow/bpmn/SequentialMultiInstanceTest.py
+++ b/tests/SpiffWorkflow/bpmn/SequentialMultiInstanceTest.py
@@ -17,7 +17,6 @@ def set_io_and_run_workflow(self, data, data_input=None, data_output=None, save_
any_task.task_spec.data_output = TaskDataReference(data_output) if data_output is not None else None
self.workflow.do_engine_steps()
- self.workflow.refresh_waiting_tasks()
ready_tasks = self.get_ready_user_tasks()
task_info = any_task.task_spec.task_info(any_task)
@@ -64,7 +63,6 @@ def run_workflow_with_condition(self, data, condition):
task.task_spec.condition = condition
self.workflow.do_engine_steps()
- self.workflow.refresh_waiting_tasks()
ready_tasks = self.get_ready_user_tasks()
while len(ready_tasks) > 0:
@@ -74,7 +72,6 @@ def run_workflow_with_condition(self, data, condition):
ready.data['output_item'] = ready.data['input_item'] * 2
ready.run()
self.workflow.do_engine_steps()
- self.workflow.refresh_waiting_tasks()
ready_tasks = self.get_ready_user_tasks()
self.workflow.do_engine_steps()
diff --git a/tests/SpiffWorkflow/bpmn/StandardLoopTest.py b/tests/SpiffWorkflow/bpmn/StandardLoopTest.py
index 2a0c30de..2c35d1ea 100644
--- a/tests/SpiffWorkflow/bpmn/StandardLoopTest.py
+++ b/tests/SpiffWorkflow/bpmn/StandardLoopTest.py
@@ -24,7 +24,6 @@ def testLoopMaximum(self):
for idx in range(3):
self.workflow.do_engine_steps()
- self.workflow.refresh_waiting_tasks()
ready_tasks = self.get_ready_user_tasks()
self.assertEqual(len(ready_tasks), 1)
ready_tasks[0].data[str(idx)] = True
@@ -47,7 +46,6 @@ def testLoopCondition(self):
start[0].data['done'] = False
self.workflow.do_engine_steps()
- self.workflow.refresh_waiting_tasks()
ready_tasks = self.get_ready_user_tasks()
self.assertEqual(len(ready_tasks), 1)
ready_tasks[0].data['done'] = True
@@ -62,8 +60,6 @@ def testSkipLoop(self):
start = self.workflow.get_tasks(end_at_spec='StartEvent_1')
start[0].data['done'] = True
self.workflow.do_engine_steps()
- self.workflow.refresh_waiting_tasks()
- self.workflow.do_engine_steps()
self.assertTrue(self.workflow.completed)
diff --git a/tests/SpiffWorkflow/bpmn/WaitingTaskStressBenchmark.py b/tests/SpiffWorkflow/bpmn/WaitingTaskStressBenchmark.py
deleted file mode 100644
index fc18e811..00000000
--- a/tests/SpiffWorkflow/bpmn/WaitingTaskStressBenchmark.py
+++ /dev/null
@@ -1,141 +0,0 @@
-"""
-Run with:
- make RUN='uv run' waiting-task-stress
-
-Useful scale knobs:
- SPIFF_WAITING_STRESS_TIMERS=500
- SPIFF_WAITING_STRESS_READY_STEPS=500
- SPIFF_WAITING_STRESS_DUE_TIMERS=50
- SPIFF_WAITING_STRESS_FUTURE_TIMERS=450
-
-Optional guard for optimized branches:
- SPIFF_WAITING_STRESS_MAX_TIMER_CHECKS=500
-"""
-
-import os
-import time
-from unittest.mock import patch
-
-from SpiffWorkflow import TaskState
-from SpiffWorkflow.bpmn.specs.event_definitions.timer import DurationTimerEventDefinition
-
-from .BpmnWorkflowTestCase import BpmnWorkflowTestCase
-from .waiting_task_stress import StressBpmnKind, WaitingTaskStressConfig, load_stress_workflow
-
-
-class WaitingTaskStressBenchmark(BpmnWorkflowTestCase):
-
- def test_ready_hot_path_with_many_dormant_timers(self):
- config = WaitingTaskStressConfig(
- waiting_timers=_env_int("SPIFF_WAITING_STRESS_TIMERS", 100),
- ready_steps=_env_int("SPIFF_WAITING_STRESS_READY_STEPS", 100),
- )
- workflow = load_stress_workflow(self, StressBpmnKind.READY_HOT_PATH, config)
-
- with _count_duration_timer_checks() as counter:
- started_at = time.perf_counter()
- workflow.do_engine_steps()
- elapsed = time.perf_counter() - started_at
-
- waiting_timers = _tasks_with_bpmn_id_prefix(workflow, TaskState.WAITING, "timer_wait_")
- completed_hot_steps = _tasks_with_bpmn_id_prefix(workflow, TaskState.COMPLETED, "hot_step_")
-
- self.assertEqual(config.waiting_timers, len(waiting_timers))
- self.assertEqual(config.ready_steps, len(completed_hot_steps))
- _print_metrics(
- "READY HOT PATH WITH DORMANT TIMERS",
- {
- "waiting_timers": config.waiting_timers,
- "ready_steps": config.ready_steps,
- "timer_has_fired_calls": counter.calls,
- "elapsed_seconds": f"{elapsed:.6f}",
- },
- )
- _assert_optional_max("SPIFF_WAITING_STRESS_MAX_TIMER_CHECKS", counter.calls, self)
-
- def test_staggered_timers_refresh_cost(self):
- due_timers = _env_int("SPIFF_WAITING_STRESS_DUE_TIMERS", 10)
- waiting_timers = _env_int("SPIFF_WAITING_STRESS_FUTURE_TIMERS", 90)
- config = WaitingTaskStressConfig(
- waiting_timers=waiting_timers,
- due_timers=due_timers,
- due_duration="PT0.01S",
- )
- workflow = load_stress_workflow(self, StressBpmnKind.STAGGERED_TIMERS, config)
- workflow.do_engine_steps()
-
- time.sleep(0.02)
- with _count_duration_timer_checks() as counter:
- started_at = time.perf_counter()
- workflow.refresh_waiting_tasks()
- workflow.do_engine_steps()
- elapsed = time.perf_counter() - started_at
-
- waiting_timer_tasks = _tasks_with_bpmn_id_prefix(workflow, TaskState.WAITING, "timer_wait_")
- completed_timer_tasks = _tasks_with_bpmn_id_prefix(workflow, TaskState.COMPLETED, "timer_wait_")
-
- self.assertEqual(waiting_timers, len(waiting_timer_tasks))
- self.assertEqual(due_timers, len(completed_timer_tasks))
- _print_metrics(
- "STAGGERED TIMER REFRESH",
- {
- "due_timers": due_timers,
- "future_timers": waiting_timers,
- "timer_has_fired_calls": counter.calls,
- "elapsed_seconds": f"{elapsed:.6f}",
- },
- )
- _assert_optional_max("SPIFF_WAITING_STRESS_MAX_TIMER_CHECKS", counter.calls, self)
-
-
-class _TimerCheckCounter:
- def __init__(self):
- self.calls = 0
-
-
-def _count_duration_timer_checks():
- counter = _TimerCheckCounter()
- original = DurationTimerEventDefinition.has_fired
-
- def counted_has_fired(event_definition, task):
- counter.calls += 1
- return original(event_definition, task)
-
- patcher = patch.object(DurationTimerEventDefinition, "has_fired", counted_has_fired)
-
- class TimerCheckContext:
- def __enter__(self):
- patcher.start()
- return counter
-
- def __exit__(self, exc_type, exc_value, traceback):
- patcher.stop()
-
- return TimerCheckContext()
-
-
-def _tasks_with_bpmn_id_prefix(workflow, state, prefix):
- return [
- task for task in workflow.get_tasks(state=state)
- if task.task_spec.bpmn_id is not None and task.task_spec.bpmn_id.startswith(prefix)
- ]
-
-
-def _env_int(name, default):
- value = os.environ.get(name)
- return default if value is None else int(value)
-
-
-def _assert_optional_max(env_name, actual, test_case):
- expected = os.environ.get(env_name)
- if expected is not None:
- test_case.assertLessEqual(actual, int(expected))
-
-
-def _print_metrics(title, metrics):
- print("\n" + "=" * 80)
- print(f"WAITING TASK STRESS: {title}")
- print("=" * 80)
- for key, value in metrics.items():
- print(f" {key}: {value}")
- print("=" * 80)
diff --git a/tests/SpiffWorkflow/bpmn/WaitingTaskStressTest.py b/tests/SpiffWorkflow/bpmn/WaitingTaskStressTest.py
deleted file mode 100644
index 6fe25ba4..00000000
--- a/tests/SpiffWorkflow/bpmn/WaitingTaskStressTest.py
+++ /dev/null
@@ -1,110 +0,0 @@
-import time
-from unittest.mock import patch
-
-from SpiffWorkflow import TaskState
-from SpiffWorkflow.bpmn.specs.event_definitions.timer import DurationTimerEventDefinition
-
-from .BpmnWorkflowTestCase import BpmnWorkflowTestCase
-from .waiting_task_stress import StressBpmnKind, WaitingTaskStressConfig, load_stress_workflow
-
-
-class WaitingTaskStressTest(BpmnWorkflowTestCase):
-
- def test_ready_hot_path_stress_fixture_keeps_many_dormant_waiting_timers(self):
- config = WaitingTaskStressConfig(waiting_timers=6, ready_steps=5)
- workflow = load_stress_workflow(self, StressBpmnKind.READY_HOT_PATH, config)
-
- workflow.do_engine_steps()
-
- waiting_timer_tasks = [
- task for task in workflow.get_tasks(state=TaskState.WAITING)
- if task.task_spec.bpmn_id is not None and task.task_spec.bpmn_id.startswith("timer_wait_")
- ]
- completed_hot_steps = [
- task for task in workflow.get_tasks(state=TaskState.COMPLETED)
- if task.task_spec.bpmn_id is not None and task.task_spec.bpmn_id.startswith("hot_step_")
- ]
-
- self.assertEqual(config.waiting_timers, len(waiting_timer_tasks))
- self.assertEqual(config.ready_steps, len(completed_hot_steps))
- self.assertFalse(workflow.completed)
-
- def test_ready_hot_path_does_not_poll_dormant_timers_per_step(self):
- config = WaitingTaskStressConfig(waiting_timers=8, ready_steps=6)
- workflow = load_stress_workflow(self, StressBpmnKind.READY_HOT_PATH, config)
-
- with _count_duration_timer_checks() as counter:
- workflow.do_engine_steps()
-
- self.assertLessEqual(counter.calls, config.waiting_timers)
-
- def test_refresh_waiting_tasks_is_noop_and_engine_steps_refresh_due_timers(self):
- config = WaitingTaskStressConfig(waiting_timers=0, due_timers=1, due_duration="PT0.01S")
- workflow = load_stress_workflow(self, StressBpmnKind.STAGGERED_TIMERS, config)
-
- workflow.do_engine_steps()
- timer_task = workflow.get_tasks(state=TaskState.WAITING, spec_name="timer_wait_0")[0]
- callbacks = []
- time.sleep(0.02)
-
- workflow.refresh_waiting_tasks(callbacks.append, callbacks.append)
-
- self.assertEqual(TaskState.WAITING, timer_task.state)
- self.assertEqual([], callbacks)
-
- workflow.do_engine_steps()
-
- self.assertEqual(TaskState.COMPLETED, timer_task.state)
-
- def test_get_tasks_does_not_refresh_due_timers_by_inspection(self):
- config = WaitingTaskStressConfig(waiting_timers=0, due_timers=1, due_duration="PT0.01S")
- workflow = load_stress_workflow(self, StressBpmnKind.STAGGERED_TIMERS, config)
-
- workflow.do_engine_steps()
- time.sleep(0.02)
-
- waiting_tasks = workflow.get_tasks(state=TaskState.WAITING, spec_name="timer_wait_0")
- ready_tasks = workflow.get_tasks(state=TaskState.READY, spec_name="timer_wait_0")
-
- self.assertEqual(1, len(waiting_tasks))
- self.assertEqual(0, len(ready_tasks))
-
- def test_due_timer_survives_save_restore_without_public_refresh(self):
- config = WaitingTaskStressConfig(waiting_timers=0, due_timers=1, due_duration="PT0.01S")
- self.workflow = load_stress_workflow(self, StressBpmnKind.STAGGERED_TIMERS, config)
-
- self.workflow.do_engine_steps()
- timer_task = self.workflow.get_tasks(state=TaskState.WAITING, spec_name="timer_wait_0")[0]
- self.save_restore()
- time.sleep(0.02)
-
- self.workflow.do_engine_steps()
-
- timer_task = self.workflow.get_task_from_id(timer_task.id)
- self.assertEqual(TaskState.COMPLETED, timer_task.state)
-
-
-class _TimerCheckCounter:
- def __init__(self):
- self.calls = 0
-
-
-def _count_duration_timer_checks():
- counter = _TimerCheckCounter()
- original = DurationTimerEventDefinition.has_fired
-
- def counted_has_fired(event_definition, task):
- counter.calls += 1
- return original(event_definition, task)
-
- patcher = patch.object(DurationTimerEventDefinition, "has_fired", counted_has_fired)
-
- class TimerCheckContext:
- def __enter__(self):
- patcher.start()
- return counter
-
- def __exit__(self, exc_type, exc_value, traceback):
- patcher.stop()
-
- return TimerCheckContext()
diff --git a/tests/SpiffWorkflow/bpmn/data/event-subprocess.bpmn b/tests/SpiffWorkflow/bpmn/data/event-subprocess.bpmn
new file mode 100644
index 00000000..3218f3fd
--- /dev/null
+++ b/tests/SpiffWorkflow/bpmn/data/event-subprocess.bpmn
@@ -0,0 +1,89 @@
+
+
+
+
+ Flow_17oshlz
+
+
+ Flow_17oshlz
+ Flow_107abhx
+
+
+
+
+ Flow_107abhx
+ Flow_1m0lq2g
+
+
+
+ Flow_1m0lq2g
+
+
+
+
+ Flow_19eg09d
+ Flow_1gvlnfw
+
+
+ Flow_1gvlnfw
+
+
+
+
+ Flow_19eg09d
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
diff --git a/tests/SpiffWorkflow/bpmn/events/ActionManagementTest.py b/tests/SpiffWorkflow/bpmn/events/ActionManagementTest.py
index b98a013c..b2224e5a 100644
--- a/tests/SpiffWorkflow/bpmn/events/ActionManagementTest.py
+++ b/tests/SpiffWorkflow/bpmn/events/ActionManagementTest.py
@@ -38,7 +38,7 @@ def testRunThroughHappy(self):
self.assertEqual('Cancel Action (if necessary)', self.workflow.get_tasks(state=TaskState.READY)[0].task_spec.bpmn_name)
time.sleep(self.START_TIME_DELTA)
- self.workflow.refresh_waiting_tasks()
+ self.workflow.refresh_timers()
self.workflow.do_engine_steps()
self.assertEqual(1, len(self.workflow.get_tasks(state=TaskState.WAITING)))
self.assertEqual(1, len(self.workflow.get_tasks(state=TaskState.STARTED)))
@@ -61,7 +61,7 @@ def testRunThroughOverdue(self):
self.assertEqual('Cancel Action (if necessary)', self.workflow.get_tasks(state=TaskState.READY)[0].task_spec.bpmn_name)
time.sleep(self.START_TIME_DELTA)
- self.workflow.refresh_waiting_tasks()
+ self.workflow.refresh_timers()
self.workflow.do_engine_steps()
self.assertEqual(1, len(self.workflow.get_tasks(state=TaskState.WAITING)))
self.assertEqual(1, len(self.workflow.get_tasks(state=TaskState.STARTED)))
@@ -74,7 +74,7 @@ def testRunThroughOverdue(self):
self.assertEqual(1, len(self.workflow.get_tasks(state=TaskState.STARTED)))
self.assertEqual('Finish Time', self.workflow.get_next_task(state=TaskState.WAITING).task_spec.bpmn_name)
time.sleep(self.FINISH_TIME_DELTA)
- self.workflow.refresh_waiting_tasks()
+ self.workflow.refresh_timers()
self.workflow.do_engine_steps()
self.assertEqual(2, len(self.workflow.get_tasks(state=TaskState.WAITING)))
self.assertEqual(1, len(self.workflow.get_tasks(state=TaskState.STARTED)))
@@ -115,7 +115,7 @@ def testRunThroughCancelAfterWorkStarted(self):
self.assertEqual(1, len(self.workflow.get_tasks(state=TaskState.READY)))
time.sleep(self.START_TIME_DELTA)
- self.workflow.refresh_waiting_tasks()
+ self.workflow.refresh_timers()
self.workflow.do_engine_steps()
self.assertEqual(1, len(self.workflow.get_tasks(state=TaskState.WAITING)))
self.assertEqual(1, len(self.workflow.get_tasks(state=TaskState.STARTED)))
diff --git a/tests/SpiffWorkflow/bpmn/events/EventBasedGatewayTest.py b/tests/SpiffWorkflow/bpmn/events/EventBasedGatewayTest.py
index ae46a163..8238f078 100644
--- a/tests/SpiffWorkflow/bpmn/events/EventBasedGatewayTest.py
+++ b/tests/SpiffWorkflow/bpmn/events/EventBasedGatewayTest.py
@@ -27,11 +27,14 @@ def actual_test(self, save_restore=False):
if save_restore:
self.save_restore()
self.workflow.script_engine = self.script_engine
- self.assertEqual(len(waiting_tasks), 2)
+ self.assertEqual(len(waiting_tasks), 4)
self.workflow.catch(BpmnEvent(MessageEventDefinition('message_1'), {}))
self.workflow.do_engine_steps()
- self.workflow.refresh_waiting_tasks()
- self.assertEqual(self.workflow.completed, True)
+ # This needs to be fixed -- it shouldn't be necessary to call this method
+ # Unfortunately that requires completely rewriting event based gateways
+ # I really don't understand why the bpmn spec dictates that both the gateway and the children
+ # have duplicate event definitions, but it sure makes things difficult
+ self.assertEqual(self.workflow.is_completed(), True)
self.assertEqual(self.workflow.get_next_task(spec_name='message_1_event').state, TaskState.COMPLETED)
self.assertEqual(self.workflow.get_next_task(spec_name='message_2_event').state, TaskState.CANCELLED)
self.assertEqual(self.workflow.get_next_task(spec_name='timer_event').state, TaskState.CANCELLED)
@@ -40,12 +43,11 @@ def testTimeout(self):
self.workflow.do_engine_steps()
waiting_tasks = self.workflow.get_tasks(state=TaskState.WAITING)
- self.assertEqual(len(waiting_tasks), 2)
+ self.assertEqual(len(waiting_tasks), 4)
timer_event_definition = waiting_tasks[0].task_spec.event_definition.event_definitions[-1]
self.workflow.catch(BpmnEvent(timer_event_definition))
- self.workflow.refresh_waiting_tasks()
self.workflow.do_engine_steps()
- self.assertEqual(self.workflow.completed, True)
+ self.assertEqual(self.workflow.is_completed(), True)
self.assertEqual(self.workflow.get_next_task(spec_name='message_1_event').state, TaskState.CANCELLED)
self.assertEqual(self.workflow.get_next_task(spec_name='message_2_event').state, TaskState.CANCELLED)
self.assertEqual(self.workflow.get_next_task(spec_name='timer_event').state, TaskState.COMPLETED)
@@ -56,5 +58,4 @@ def testMultipleStart(self):
workflow.do_engine_steps()
workflow.catch(BpmnEvent(MessageEventDefinition('message_1'), {}))
workflow.catch(BpmnEvent(MessageEventDefinition('message_2'), {}))
- workflow.refresh_waiting_tasks()
workflow.do_engine_steps()
diff --git a/tests/SpiffWorkflow/bpmn/events/EventSubprocessTest.py b/tests/SpiffWorkflow/bpmn/events/EventSubprocessTest.py
new file mode 100644
index 00000000..5a88ff5e
--- /dev/null
+++ b/tests/SpiffWorkflow/bpmn/events/EventSubprocessTest.py
@@ -0,0 +1,30 @@
+from SpiffWorkflow import TaskState
+from SpiffWorkflow.bpmn import BpmnWorkflow
+
+from ..BpmnWorkflowTestCase import BpmnWorkflowTestCase
+
+class EventBasedGatewayTest(BpmnWorkflowTestCase):
+
+ def setUp(self):
+ self.spec, self.subprocesses = self.load_workflow_spec('event-subprocess.bpmn', 'main')
+ self.workflow = BpmnWorkflow(self.spec, self.subprocesses)
+
+ def testEventSubprocess(self):
+ self.actual_test()
+
+ def testEventSubprocessSaveRestore(self):
+ self.actual_test(True)
+
+ def actual_test(self, save_restore=False):
+ self.workflow.do_engine_steps()
+ set_data = self.workflow.get_next_task(spec_name='set_data')
+ set_data.data.update(v1=True, v2=False)
+ set_data.run()
+ self.workflow.do_engine_steps()
+ task = self.workflow.get_next_task(spec_name='task')
+ self.assertEqual(task.state, TaskState.READY)
+ self.assertDictEqual(task.data, {'v1': True, 'v2': False})
+ task.run()
+ self.workflow.do_engine_steps()
+ self.assertTrue(self.workflow.completed)
+
diff --git a/tests/SpiffWorkflow/bpmn/events/MultipleCatchEventTest.py b/tests/SpiffWorkflow/bpmn/events/MultipleCatchEventTest.py
index 2e5140be..b8fcdc90 100644
--- a/tests/SpiffWorkflow/bpmn/events/MultipleCatchEventTest.py
+++ b/tests/SpiffWorkflow/bpmn/events/MultipleCatchEventTest.py
@@ -30,7 +30,6 @@ def actual_test(self, save_restore=False):
self.assertEqual(waiting_tasks[0].task_spec.name, 'StartEvent_1')
self.workflow.catch(BpmnEvent(MessageEventDefinition('message_1'), {}))
- self.workflow.refresh_waiting_tasks()
self.workflow.do_engine_steps()
# Now the first task should be ready
@@ -64,7 +63,6 @@ def actual_test(self, save_restore=False):
self.assertEqual(waiting_tasks[0].task_spec.name, 'StartEvent_1')
self.workflow.catch(BpmnEvent(MessageEventDefinition('message_1'), {}))
- self.workflow.refresh_waiting_tasks()
self.workflow.do_engine_steps()
# It should still be waiting because it has to receive both messages
@@ -73,7 +71,6 @@ def actual_test(self, save_restore=False):
self.assertEqual(waiting_tasks[0].task_spec.name, 'StartEvent_1')
self.workflow.catch(BpmnEvent(MessageEventDefinition('message_2'), {}))
- self.workflow.refresh_waiting_tasks()
self.workflow.do_engine_steps()
# Now the first task should be ready
diff --git a/tests/SpiffWorkflow/bpmn/events/MultipleThrowEventTest.py b/tests/SpiffWorkflow/bpmn/events/MultipleThrowEventTest.py
index 7dda3436..898a20c0 100644
--- a/tests/SpiffWorkflow/bpmn/events/MultipleThrowEventTest.py
+++ b/tests/SpiffWorkflow/bpmn/events/MultipleThrowEventTest.py
@@ -27,9 +27,8 @@ def actual_test(self, save_restore=False):
class MultipleThrowEventStartsEventTest(BpmnWorkflowTestCase):
def setUp(self):
- specs = self.get_all_specs('multiple-throw-start.bpmn')
- self.spec = specs.pop('initiate')
- self.workflow = BpmnWorkflow(self.spec, specs)
+ self.spec, subprocesses = self.load_collaboration('multiple-throw-start.bpmn', 'top')
+ self.workflow = BpmnWorkflow(self.spec, subprocesses)
def testMultipleThrowEventStartEvent(self):
self.actual_test()
diff --git a/tests/SpiffWorkflow/bpmn/events/NITimerDurationBoundaryTest.py b/tests/SpiffWorkflow/bpmn/events/NITimerDurationBoundaryTest.py
index bdc72051..5352f407 100644
--- a/tests/SpiffWorkflow/bpmn/events/NITimerDurationBoundaryTest.py
+++ b/tests/SpiffWorkflow/bpmn/events/NITimerDurationBoundaryTest.py
@@ -45,7 +45,7 @@ def actual_test(self,save_restore = False):
ready_tasks = self.workflow.get_tasks(state=TaskState.READY)
# There should be one ready task until the boundary event fires
self.assertEqual(len(self.get_ready_user_tasks()), 1)
- self.workflow.refresh_waiting_tasks()
+ self.workflow.refresh_timers()
self.workflow.do_engine_steps()
loopcount += 1
@@ -63,7 +63,7 @@ def actual_test(self,save_restore = False):
elif task.task_spec.name == 'Activity_Work':
task.data['work_done'] = 'Yes'
task.run()
- self.workflow.refresh_waiting_tasks()
+ self.workflow.refresh_timers()
self.workflow.do_engine_steps()
self.workflow.do_engine_steps()
self.assertEqual(self.workflow.completed, True)
diff --git a/tests/SpiffWorkflow/bpmn/events/TimerCycleStartTest.py b/tests/SpiffWorkflow/bpmn/events/TimerCycleStartTest.py
index 34eba577..3fc61c9a 100644
--- a/tests/SpiffWorkflow/bpmn/events/TimerCycleStartTest.py
+++ b/tests/SpiffWorkflow/bpmn/events/TimerCycleStartTest.py
@@ -54,7 +54,7 @@ def actual_test(self,save_restore = False):
if save_restore:
self.save_restore()
time.sleep(0.1)
- self.workflow.refresh_waiting_tasks()
+ self.workflow.refresh_timers()
self.assertEqual(counter, 2)
self.assertTrue(self.workflow.completed)
diff --git a/tests/SpiffWorkflow/bpmn/events/TimerCycleTest.py b/tests/SpiffWorkflow/bpmn/events/TimerCycleTest.py
index 1d287753..0459dba2 100644
--- a/tests/SpiffWorkflow/bpmn/events/TimerCycleTest.py
+++ b/tests/SpiffWorkflow/bpmn/events/TimerCycleTest.py
@@ -48,7 +48,7 @@ def actual_test(self,save_restore = False):
self.workflow.do_engine_steps()
if save_restore:
self.save_restore()
- self.workflow.do_engine_steps()
+ self.workflow.refresh_timers()
events = self.workflow.waiting_events()
refill = self.workflow.get_tasks(spec_name='Refill_Coffee')
# Wait time is 0.1s, with a limit of 2 children, so by the 3rd iteration, the event should be complete
diff --git a/tests/SpiffWorkflow/bpmn/events/TimerDateTest.py b/tests/SpiffWorkflow/bpmn/events/TimerDateTest.py
index e05f99fb..aa6a50fe 100644
--- a/tests/SpiffWorkflow/bpmn/events/TimerDateTest.py
+++ b/tests/SpiffWorkflow/bpmn/events/TimerDateTest.py
@@ -37,7 +37,7 @@ def actual_test(self,save_restore = False):
self.save_restore()
self.workflow.script_engine = self.script_engine
time.sleep(0.01)
- self.workflow.refresh_waiting_tasks()
+ self.workflow.refresh_timers()
loopcount += 1
endtime = datetime.datetime.now()
self.workflow.do_engine_steps()
diff --git a/tests/SpiffWorkflow/bpmn/events/TimerDurationBoundaryOnTaskTest.py b/tests/SpiffWorkflow/bpmn/events/TimerDurationBoundaryOnTaskTest.py
index 581e572a..e24df6d0 100644
--- a/tests/SpiffWorkflow/bpmn/events/TimerDurationBoundaryOnTaskTest.py
+++ b/tests/SpiffWorkflow/bpmn/events/TimerDurationBoundaryOnTaskTest.py
@@ -29,7 +29,7 @@ def actual_test(self,save_restore = False):
self.save_restore()
self.workflow.script_engine = self.script_engine
time.sleep(1)
- self.workflow.refresh_waiting_tasks()
+ self.workflow.refresh_timers()
self.workflow.do_engine_steps()
# Make sure the timer got called
diff --git a/tests/SpiffWorkflow/bpmn/events/TimerDurationBoundaryTest.py b/tests/SpiffWorkflow/bpmn/events/TimerDurationBoundaryTest.py
index f073fa07..b0b8b197 100644
--- a/tests/SpiffWorkflow/bpmn/events/TimerDurationBoundaryTest.py
+++ b/tests/SpiffWorkflow/bpmn/events/TimerDurationBoundaryTest.py
@@ -33,7 +33,7 @@ def actual_test(self,save_restore = False):
self.save_restore()
time.sleep(0.01)
self.assertEqual(len(self.workflow.get_tasks(state=TaskState.READY)), 1)
- self.workflow.refresh_waiting_tasks()
+ self.workflow.refresh_timers()
self.workflow.do_engine_steps()
loopcount += 1
diff --git a/tests/SpiffWorkflow/bpmn/events/TimerDurationTest.py b/tests/SpiffWorkflow/bpmn/events/TimerDurationTest.py
index 77886aba..1f195d8e 100644
--- a/tests/SpiffWorkflow/bpmn/events/TimerDurationTest.py
+++ b/tests/SpiffWorkflow/bpmn/events/TimerDurationTest.py
@@ -35,7 +35,7 @@ def actual_test(self,save_restore = False):
self.save_restore()
self.workflow.script_engine = self.script_engine
time.sleep(0.1)
- self.workflow.do_engine_steps()
+ self.workflow.refresh_timers()
loopcount += 1
endtime = datetime.now()
duration = endtime - starttime
diff --git a/tests/SpiffWorkflow/bpmn/events/TimerIntermediateTest.py b/tests/SpiffWorkflow/bpmn/events/TimerIntermediateTest.py
index 3e78088d..637890f2 100644
--- a/tests/SpiffWorkflow/bpmn/events/TimerIntermediateTest.py
+++ b/tests/SpiffWorkflow/bpmn/events/TimerIntermediateTest.py
@@ -31,6 +31,9 @@ def testRunThroughHappy(self):
time.sleep(0.02)
self.assertEqual(1, len(self.workflow.get_tasks(state=TaskState.WAITING)))
- self.workflow.do_engine_steps()
+ self.workflow.refresh_timers()
self.assertEqual(0, len(self.workflow.get_tasks(state=TaskState.WAITING)))
+ self.assertEqual(1, len(self.workflow.get_tasks(state=TaskState.READY)))
+
+ self.workflow.do_engine_steps()
self.assertEqual(0, len(self.workflow.get_tasks(state=TaskState.READY|TaskState.WAITING)))
diff --git a/tests/SpiffWorkflow/bpmn/serializer/VersionMigrationTest.py b/tests/SpiffWorkflow/bpmn/serializer/VersionMigrationTest.py
index a85d752a..d4624526 100644
--- a/tests/SpiffWorkflow/bpmn/serializer/VersionMigrationTest.py
+++ b/tests/SpiffWorkflow/bpmn/serializer/VersionMigrationTest.py
@@ -18,10 +18,6 @@ def test_convert_subprocess(self):
self.assertEqual('Action3', ready_tasks[0].task_spec.bpmn_name)
ready_tasks[0].run()
wf.do_engine_steps()
- wf.refresh_waiting_tasks()
- wf.do_engine_steps()
- wf.refresh_waiting_tasks()
- wf.do_engine_steps()
self.assertEqual(True, wf.completed)
@@ -30,17 +26,13 @@ class Version_1_1_Test(BaseTestCase):
def test_timers(self):
wf = self.deserialize_workflow('v1.1-timers.json')
wf.script_engine = PythonScriptEngine(environment=TaskDataEnvironment({"time": time}))
- wf.refresh_waiting_tasks()
- wf.do_engine_steps()
- wf.refresh_waiting_tasks()
+ wf.refresh_timers()
wf.do_engine_steps()
self.assertTrue(wf.completed)
def test_convert_data_specs(self):
wf = self.deserialize_workflow('v1.1-data.json')
wf.do_engine_steps()
- wf.refresh_waiting_tasks()
- wf.do_engine_steps()
self.assertTrue(wf.completed)
def test_convert_exclusive_gateway(self):
@@ -66,7 +58,7 @@ def test_remove_loop_reset(self):
end = time.time() + 3
while not wf.completed and time.time() < end:
wf.do_engine_steps()
- wf.refresh_waiting_tasks()
+ wf.refresh_timers()
self.assertTrue(wf.completed)
self.assertEqual(wf.last_task.data['counter'], 20)
@@ -195,7 +187,6 @@ def test_update_mi_states(self):
task.data['output_item'] = task.data['input_item'] * 2
task.run()
ready_tasks = wf.get_tasks(state=TaskState.READY, manual=True)
- wf.refresh_waiting_tasks()
wf.do_engine_steps()
any_task = wf.get_next_task(spec_name='any_task')
diff --git a/tests/SpiffWorkflow/bpmn/waiting_task_stress.py b/tests/SpiffWorkflow/bpmn/waiting_task_stress.py
deleted file mode 100644
index a5379d85..00000000
--- a/tests/SpiffWorkflow/bpmn/waiting_task_stress.py
+++ /dev/null
@@ -1,152 +0,0 @@
-import os
-from dataclasses import dataclass
-from enum import Enum
-from uuid import uuid4
-
-from SpiffWorkflow.bpmn.workflow import BpmnWorkflow
-
-
-class StressBpmnKind(Enum):
- READY_HOT_PATH = "ready_hot_path"
- STAGGERED_TIMERS = "staggered_timers"
-
-
-@dataclass(frozen=True)
-class WaitingTaskStressConfig:
- waiting_timers: int = 100
- ready_steps: int = 100
- due_timers: int = 0
- future_duration: str = "PT24H"
- due_duration: str = "PT0S"
-
-
-def build_stress_bpmn(kind, config):
- if kind == StressBpmnKind.READY_HOT_PATH:
- return _build_ready_hot_path_bpmn(config)
- if kind == StressBpmnKind.STAGGERED_TIMERS:
- return _build_staggered_timers_bpmn(config)
- raise ValueError(f"Unsupported stress BPMN kind: {kind}")
-
-
-def load_stress_workflow(test_case, kind, config):
- filename = write_stress_bpmn(test_case, kind, config)
- try:
- spec, subprocesses = test_case.load_workflow_spec(filename, "waiting_task_stress", validate=False)
- return BpmnWorkflow(spec, subprocesses)
- finally:
- path = _data_path(filename)
- if os.path.exists(path):
- os.unlink(path)
-
-
-def write_stress_bpmn(test_case, kind, config):
- filename = f"_generated_waiting_task_stress_{kind.value}_{uuid4().hex}.bpmn"
- path = _data_path(filename)
- with open(path, "w") as bpmn_file:
- bpmn_file.write(build_stress_bpmn(kind, config))
- return filename
-
-
-def _data_path(filename):
- return os.path.join(os.path.dirname(__file__), "data", filename)
-
-
-def _build_ready_hot_path_bpmn(config):
- timer_branches = [
- _timer_branch(idx, config.future_duration)
- for idx in range(config.waiting_timers)
- ]
- timer_flows = [
- f'flow_split_timer_{idx}'
- for idx in range(config.waiting_timers)
- ]
- return _definitions(
- "\n".join([
- _start_and_split(timer_flows + ["flow_split_hot_0"]),
- "\n".join(timer_branches),
- _hot_path(config.ready_steps),
- ])
- )
-
-
-def _build_staggered_timers_bpmn(config):
- timer_count = config.waiting_timers + config.due_timers
- timer_flows = [
- f'flow_split_timer_{idx}'
- for idx in range(timer_count)
- ]
- branches = []
- for idx in range(timer_count):
- duration = config.due_duration if idx < config.due_timers else config.future_duration
- branches.append(_timer_branch(idx, duration))
- return _definitions("\n".join([
- _start_and_split(timer_flows),
- "\n".join(branches),
- ]))
-
-
-def _definitions(process_body):
- return f"""
-
-
-{_indent(process_body, 4)}
-
-
-"""
-
-
-def _start_and_split(split_outgoing):
- outgoing = "\n".join(split_outgoing)
- return f"""
- flow_start_split
-
-
- flow_start_split
-{_indent(outgoing, 2)}
-
-"""
-
-
-def _timer_branch(idx, duration):
- return f"""
- flow_split_timer_{idx}
- flow_timer_{idx}_end
-
- "{duration}"
-
-
-
- flow_timer_{idx}_end
-
-
-"""
-
-
-def _hot_path(ready_steps):
- if ready_steps < 1:
- raise ValueError("ready_steps must be at least 1")
-
- tasks = []
- flows = ['']
- for idx in range(ready_steps):
- incoming = "flow_split_hot_0" if idx == 0 else f"flow_hot_{idx - 1}_{idx}"
- outgoing = "flow_hot_last_end" if idx == ready_steps - 1 else f"flow_hot_{idx}_{idx + 1}"
- tasks.append(f"""
- {incoming}
- {outgoing}
- hot_path_steps = hot_path_steps + 1 if 'hot_path_steps' in locals() else 1
-""")
- if idx < ready_steps - 1:
- flows.append(
- f''
- )
- flows.append('flow_hot_last_end')
- flows.append(
- f''
- )
- return "\n".join(tasks + flows)
-
-
-def _indent(text, spaces):
- prefix = " " * spaces
- return "\n".join(f"{prefix}{line}" if line else line for line in text.splitlines())
diff --git a/tests/SpiffWorkflow/camunda/MessageBoundaryEventTest.py b/tests/SpiffWorkflow/camunda/MessageBoundaryEventTest.py
index 8d0627a4..b2031583 100644
--- a/tests/SpiffWorkflow/camunda/MessageBoundaryEventTest.py
+++ b/tests/SpiffWorkflow/camunda/MessageBoundaryEventTest.py
@@ -40,12 +40,12 @@ def actual_test(self,save_restore = False):
self.workflow.run_task_from_id(task.id)
self.workflow.do_engine_steps()
time.sleep(.01)
- self.workflow.refresh_waiting_tasks()
+ self.workflow.refresh_timers()
if save_restore:
self.save_restore()
ready_tasks = self.workflow.get_tasks(state=TaskState.READY)
time.sleep(.01)
- self.workflow.refresh_waiting_tasks()
+ self.workflow.refresh_timers()
self.workflow.do_engine_steps()
self.assertEqual(self.workflow.completed, True, 'Expected the workflow to be complete at this point')
diff --git a/tests/SpiffWorkflow/camunda/StartMessageEventTest.py b/tests/SpiffWorkflow/camunda/StartMessageEventTest.py
index a6849465..831b59f1 100644
--- a/tests/SpiffWorkflow/camunda/StartMessageEventTest.py
+++ b/tests/SpiffWorkflow/camunda/StartMessageEventTest.py
@@ -46,7 +46,6 @@ def actual_test(self,save_restore = False):
current_task.set_data(**step[1])
current_task.run()
self.workflow.do_engine_steps()
- self.workflow.refresh_waiting_tasks()
if save_restore:
self.save_restore()
ready_tasks = self.workflow.get_tasks(state=TaskState.READY)
diff --git a/tests/SpiffWorkflow/spiff/CorrelationTest.py b/tests/SpiffWorkflow/spiff/CorrelationTest.py
index 5fdf1a9f..d3ba1ef8 100644
--- a/tests/SpiffWorkflow/spiff/CorrelationTest.py
+++ b/tests/SpiffWorkflow/spiff/CorrelationTest.py
@@ -13,9 +13,8 @@ def testMessagePayloadSaveRestore(self):
def actual_test(self,save_restore):
- specs = self.get_all_specs('correlation.bpmn')
- proc_1 = specs['proc_1']
- self.workflow = BpmnWorkflow(proc_1, specs)
+ spec, subprocesses = self.load_collaboration('correlation.bpmn', 'correlation_test')
+ self.workflow = BpmnWorkflow(spec, subprocesses)
if save_restore:
self.save_restore()
self.workflow.do_engine_steps()