From 28756c05d1c5b937fe0bd5ffc221ed9f0e237172 Mon Sep 17 00:00:00 2001 From: Elizabeth Esswein Date: Sun, 7 Jun 2026 10:03:00 -0400 Subject: [PATCH 1/4] revert timer handling --- Makefile | 4 - .../bpmn/serializer/default/workflow.py | 1 - SpiffWorkflow/bpmn/util/subworkflow.py | 7 +- SpiffWorkflow/bpmn/workflow.py | 170 ++---------------- SpiffWorkflow/task.py | 3 - SpiffWorkflow/workflow.py | 2 - .../bpmn/WaitingTaskStressBenchmark.py | 141 --------------- .../bpmn/WaitingTaskStressTest.py | 110 ------------ .../bpmn/events/TimerCycleTest.py | 2 +- .../bpmn/events/TimerDurationTest.py | 2 +- .../bpmn/events/TimerIntermediateTest.py | 5 +- .../SpiffWorkflow/bpmn/waiting_task_stress.py | 152 ---------------- 12 files changed, 23 insertions(+), 576 deletions(-) delete mode 100644 tests/SpiffWorkflow/bpmn/WaitingTaskStressBenchmark.py delete mode 100644 tests/SpiffWorkflow/bpmn/WaitingTaskStressTest.py delete mode 100644 tests/SpiffWorkflow/bpmn/waiting_task_stress.py 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/serializer/default/workflow.py b/SpiffWorkflow/bpmn/serializer/default/workflow.py index 9284c1d1..d2c34c16 100644 --- a/SpiffWorkflow/bpmn/serializer/default/workflow.py +++ b/SpiffWorkflow/bpmn/serializer/default/workflow.py @@ -201,7 +201,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/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..4568f1fb 100644 --- a/SpiffWorkflow/bpmn/workflow.py +++ b/SpiffWorkflow/bpmn/workflow.py @@ -17,9 +17,6 @@ # Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, MA # 02110-1301 USA -import heapq -from datetime import datetime, timezone - from SpiffWorkflow.task import Task from SpiffWorkflow.util.task import TaskState from SpiffWorkflow.exceptions import WorkflowException @@ -27,8 +24,6 @@ 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 @@ -38,98 +33,6 @@ 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 +51,6 @@ 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 super().__init__(spec, **kwargs) for obj in self.spec.data_objects: @@ -228,7 +129,7 @@ def catch(self, event): for task in tasks: task.task_spec.catch(task, event) if len(tasks) > 0: - self._refresh_caught_tasks(tasks) + self.refresh_waiting_tasks() def send_event(self, event): """Allows this workflow to catch an externally generated event.""" @@ -241,7 +142,7 @@ def send_event(self, event): 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) + self.refresh_waiting_tasks() def get_events(self): """Returns the list of events that cannot be handled from within this workflow.""" @@ -263,10 +164,8 @@ 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) def _do_engine_steps(self, will_complete_task=None, did_complete_task=None): @@ -298,62 +197,25 @@ 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: + def update_task(task): + if will_refresh_task is not None: + will_refresh_task(task) task.task_spec._update(task) - self._waiting_task_index.reschedule_timer(task) + if did_refresh_task is not None: + did_refresh_task(task) + + for subprocess in sorted(self.get_active_subprocesses(), key=lambda v: v.depth, reverse=True): + for task in subprocess.get_tasks_iterator(skip_subprocesses=True, state=TaskState.WAITING): + update_task(task) + + for task in self.get_tasks_iterator(skip_subprocesses=True, state=TaskState.WAITING): + update_task(task) def get_task_from_id(self, task_id): if task_id not in self.tasks: 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/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/events/TimerCycleTest.py b/tests/SpiffWorkflow/bpmn/events/TimerCycleTest.py index 1d287753..75a8236b 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_waiting_tasks() 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/TimerDurationTest.py b/tests/SpiffWorkflow/bpmn/events/TimerDurationTest.py index 77886aba..b29a787f 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_waiting_tasks() 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..91326e0b 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_waiting_tasks() 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/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()) From a44597c2359ef31d2919fa4e1c116aca663b335b Mon Sep 17 00:00:00 2001 From: Elizabeth Esswein Date: Mon, 8 Jun 2026 18:01:53 -0400 Subject: [PATCH 2/4] allow bpmn start task to trigger new subprocesses --- SpiffWorkflow/bpmn/parser/BpmnParser.py | 19 ++------- SpiffWorkflow/bpmn/serializer/config.py | 3 +- .../bpmn/serializer/default/process_spec.py | 3 ++ .../bpmn/serializer/default/task_spec.py | 12 ++++++ SpiffWorkflow/bpmn/specs/bpmn_process_spec.py | 8 ++++ SpiffWorkflow/bpmn/specs/control.py | 28 ++++++++++++- SpiffWorkflow/bpmn/util/event.py | 16 +++++++- SpiffWorkflow/bpmn/workflow.py | 40 +++++-------------- tests/SpiffWorkflow/bpmn/CollaborationTest.py | 10 ++--- .../bpmn/events/MultipleThrowEventTest.py | 5 +-- tests/SpiffWorkflow/spiff/CorrelationTest.py | 5 +-- 11 files changed, 89 insertions(+), 60 deletions(-) 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/serializer/config.py b/SpiffWorkflow/bpmn/serializer/config.py index dc31382a..c2f8c8f0 100644 --- a/SpiffWorkflow/bpmn/serializer/config.py +++ b/SpiffWorkflow/bpmn/serializer/config.py @@ -96,6 +96,7 @@ EventConverter, BoundaryEventConverter, IOSpecificationConverter, + BpmnStartTaskConverter, ) from .default.event_definition import ( TimerConditionalEventDefinitionConverter, @@ -115,7 +116,7 @@ BpmnIoSpecification: IOSpecificationConverter, BpmnProcessSpec: BpmnProcessSpecConverter, SimpleBpmnTask: BpmnTaskSpecConverter, - BpmnStartTask: BpmnTaskSpecConverter, + BpmnStartTask: BpmnStartTaskConverter, _EndJoin: BpmnTaskSpecConverter, NoneTask: BpmnTaskSpecConverter, ManualTask: BpmnTaskSpecConverter, 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/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/util/event.py b/SpiffWorkflow/bpmn/util/event.py index a885d037..510e484d 100644 --- a/SpiffWorkflow/bpmn/util/event.py +++ b/SpiffWorkflow/bpmn/util/event.py @@ -11,4 +11,18 @@ 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): + 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, None) + + diff --git a/SpiffWorkflow/bpmn/workflow.py b/SpiffWorkflow/bpmn/workflow.py index 4568f1fb..c66dbce5 100644 --- a/SpiffWorkflow/bpmn/workflow.py +++ b/SpiffWorkflow/bpmn/workflow.py @@ -21,13 +21,10 @@ from SpiffWorkflow.util.task import TaskState from SpiffWorkflow.exceptions import WorkflowException +from SpiffWorkflow.bpmn.specs.control import BoundaryEventSplit 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.item_aware_event import CodeEventDefinition -from SpiffWorkflow.bpmn.specs.control import BoundaryEventSplit - from SpiffWorkflow.bpmn.util.subworkflow import BpmnBaseWorkflow, BpmnSubWorkflow from .script_engine.python_engine import PythonScriptEngine @@ -167,6 +164,7 @@ def do_engine_steps(self, will_complete_task=None, did_complete_task=None): count = self._do_engine_steps(will_complete_task, did_complete_task) while count > 0: count = self._do_engine_steps(will_complete_task, did_complete_task) + self.refresh_waiting_tasks() def _do_engine_steps(self, will_complete_task=None, did_complete_task=None): @@ -271,33 +269,17 @@ def cancel(self, workflow=None): 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: + if sp.get_next_task(state=TaskState.WAITING, spec_name=task_spec) 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] + child_id = self.spec.start.trigger_wf(self.task_tree, wf_spec) + return self.subprocesses[child_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) + for name in self.spec.start.trigger_specs: + sp_spec = self.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/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/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/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() From 8ad6a7469f938fd49c976b2fc3b55308beeb25b2 Mon Sep 17 00:00:00 2001 From: Elizabeth Esswein Date: Mon, 8 Jun 2026 23:40:33 -0400 Subject: [PATCH 3/4] split event management from workflow --- .../bpmn/serializer/default/workflow.py | 6 ++ .../bpmn/specs/mixins/events/event_types.py | 10 +++ .../specs/mixins/events/intermediate_event.py | 3 +- SpiffWorkflow/bpmn/util/event.py | 47 +++++++++- SpiffWorkflow/bpmn/workflow.py | 89 +++++-------------- .../bpmn/BpmnWorkflowTestCase.py | 2 +- .../bpmn/ParallelMultiInstanceTest.py | 2 - .../bpmn/ResetTokenOnBoundaryEventTest.py | 2 - .../bpmn/SequentialMultiInstanceTest.py | 3 - tests/SpiffWorkflow/bpmn/StandardLoopTest.py | 4 - .../bpmn/events/ActionManagementTest.py | 8 +- .../bpmn/events/EventBasedGatewayTest.py | 15 ++-- .../bpmn/events/MultipleCatchEventTest.py | 3 - .../events/NITimerDurationBoundaryTest.py | 4 +- .../bpmn/events/TimerCycleStartTest.py | 2 +- .../bpmn/events/TimerCycleTest.py | 2 +- .../bpmn/events/TimerDateTest.py | 2 +- .../events/TimerDurationBoundaryOnTaskTest.py | 2 +- .../bpmn/events/TimerDurationBoundaryTest.py | 2 +- .../bpmn/events/TimerDurationTest.py | 2 +- .../bpmn/events/TimerIntermediateTest.py | 2 +- .../bpmn/serializer/VersionMigrationTest.py | 13 +-- .../camunda/MessageBoundaryEventTest.py | 4 +- .../camunda/StartMessageEventTest.py | 1 - 24 files changed, 110 insertions(+), 120 deletions(-) diff --git a/SpiffWorkflow/bpmn/serializer/default/workflow.py b/SpiffWorkflow/bpmn/serializer/default/workflow.py index d2c34c16..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): 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/util/event.py b/SpiffWorkflow/bpmn/util/event.py index 510e484d..1d35c297 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 @@ -16,13 +18,54 @@ def __init__(self, name, event_type, value=None, correlations=None): class EventManager: - def __init__(self): + 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, None) + 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/workflow.py b/SpiffWorkflow/bpmn/workflow.py index c66dbce5..a676709c 100644 --- a/SpiffWorkflow/bpmn/workflow.py +++ b/SpiffWorkflow/bpmn/workflow.py @@ -17,15 +17,17 @@ # Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, MA # 02110-1301 USA +import warnings + from SpiffWorkflow.task import Task from SpiffWorkflow.util.task import TaskState from SpiffWorkflow.exceptions import WorkflowException from SpiffWorkflow.bpmn.specs.control import BoundaryEventSplit -from SpiffWorkflow.bpmn.specs.mixins.events.event_types import CatchingEvent -from SpiffWorkflow.bpmn.specs.event_definitions.item_aware_event import CodeEventDefinition +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 @@ -48,6 +50,7 @@ def __init__(self, spec, subprocess_specs=None, script_engine=None, **kwargs): self.subprocesses = {} self.bpmn_events = [] self.correlations = {} + self.event_manager = EventManager(self) super().__init__(spec, **kwargs) for obj in self.spec.data_objects: @@ -102,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_waiting_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_waiting_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.""" @@ -148,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): """ @@ -164,7 +135,7 @@ def do_engine_steps(self, will_complete_task=None, did_complete_task=None): count = self._do_engine_steps(will_complete_task, did_complete_task) while count > 0: count = self._do_engine_steps(will_complete_task, did_complete_task) - self.refresh_waiting_tasks() + self.refresh_timers() def _do_engine_steps(self, will_complete_task=None, did_complete_task=None): @@ -201,19 +172,17 @@ def refresh_waiting_tasks(self, will_refresh_task=None, did_refresh_task=None): :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 """ - def update_task(task): - if will_refresh_task is not None: - will_refresh_task(task) - task.task_spec._update(task) - if did_refresh_task is not None: - did_refresh_task(task) - - for subprocess in sorted(self.get_active_subprocesses(), key=lambda v: v.depth, reverse=True): - for task in subprocess.get_tasks_iterator(skip_subprocesses=True, state=TaskState.WAITING): - update_task(task) - - for task in self.get_tasks_iterator(skip_subprocesses=True, state=TaskState.WAITING): - update_task(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: @@ -266,20 +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) is not None: - return sp - child_id = self.spec.start.trigger_wf(self.task_tree, wf_spec) - return self.subprocesses[child_id] - - # Start a subprocess for known specs with start events that catch this - for name in self.spec.start.trigger_specs: - sp_spec = self.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/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/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/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/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/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 75a8236b..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.refresh_waiting_tasks() + 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 b29a787f..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.refresh_waiting_tasks() + 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 91326e0b..637890f2 100644 --- a/tests/SpiffWorkflow/bpmn/events/TimerIntermediateTest.py +++ b/tests/SpiffWorkflow/bpmn/events/TimerIntermediateTest.py @@ -31,7 +31,7 @@ def testRunThroughHappy(self): time.sleep(0.02) self.assertEqual(1, len(self.workflow.get_tasks(state=TaskState.WAITING))) - self.workflow.refresh_waiting_tasks() + 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))) 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/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) From ac48a3c656bb5591a0766520ee9d06e81e53ce32 Mon Sep 17 00:00:00 2001 From: Elizabeth Esswein Date: Tue, 9 Jun 2026 14:51:36 -0400 Subject: [PATCH 4/4] add event subprocess --- SpiffWorkflow/bpmn/parser/ProcessParser.py | 5 ++ SpiffWorkflow/bpmn/parser/TaskParser.py | 5 +- SpiffWorkflow/bpmn/parser/task_parsers.py | 6 +- SpiffWorkflow/bpmn/serializer/config.py | 2 + SpiffWorkflow/bpmn/specs/defaults.py | 4 + SpiffWorkflow/bpmn/specs/mixins/__init__.py | 3 +- .../bpmn/specs/mixins/subworkflow_task.py | 4 + SpiffWorkflow/bpmn/util/event.py | 1 - SpiffWorkflow/spiff/parser/task_spec.py | 11 ++- SpiffWorkflow/spiff/serializer/config.py | 4 + SpiffWorkflow/spiff/specs/defaults.py | 4 + doc/bpmn/supported.rst | 1 + .../bpmn/data/event-subprocess.bpmn | 89 +++++++++++++++++++ .../bpmn/events/EventSubprocessTest.py | 30 +++++++ 14 files changed, 163 insertions(+), 6 deletions(-) create mode 100644 tests/SpiffWorkflow/bpmn/data/event-subprocess.bpmn create mode 100644 tests/SpiffWorkflow/bpmn/events/EventSubprocessTest.py 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 c2f8c8f0..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, @@ -128,6 +129,7 @@ SubWorkflowTask: SubWorkflowConverter, CallActivity: SubWorkflowConverter, TransactionSubprocess: SubWorkflowConverter, + EventSubprocess: SubWorkflowConverter, BoundaryEventSplit: BpmnTaskSpecConverter, BoundaryEventJoin: EventJoinConverter, ExclusiveGateway: ExclusiveGatewayConverter, 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/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 1d35c297..cbd17c2f 100644 --- a/SpiffWorkflow/bpmn/util/event.py +++ b/SpiffWorkflow/bpmn/util/event.py @@ -68,4 +68,3 @@ def get_or_create_subprocess(task_spec, wf_spec): 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/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/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/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/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) +