Skip to content

Commit 9e6ee7f

Browse files
authored
Merge pull request #481 from sartography/bugfix/event-based-gateway-fixes
Bugfix/event based gateway fixes
2 parents 73ee57e + 0907bbe commit 9e6ee7f

11 files changed

Lines changed: 617 additions & 83 deletions

File tree

SpiffWorkflow/bpmn/parser/event_parsers.py

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -306,6 +306,3 @@ def create_task(self):
306306
def handles_multiple_outgoing(self):
307307
return True
308308

309-
def connect_outgoing(self, outgoing_task, sequence_flow_node, is_default):
310-
self.task.event_definition.event_definitions.append(outgoing_task.event_definition)
311-
self.task.connect(outgoing_task)
Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,18 @@
1+
2+
def update_event_gateway_children(dct):
3+
4+
def update(tasks, specs):
5+
for task in tasks:
6+
task_spec = specs.get(task['task_spec'], {})
7+
if task_spec['typename'] == 'EventBasedGateway':
8+
for name in task_spec['outputs']:
9+
child = specs.get(name)
10+
child['event_definition'] = {
11+
"description": "Default",
12+
"name": None,
13+
"typename": "NoneEventDefinition"
14+
}
15+
16+
for up in dct['subprocesses'].values():
17+
update(sp['tasks'].values(), sp['spec']['task_specs'])
18+
update(dct['tasks'].values(), dct['spec']['task_specs'])

SpiffWorkflow/bpmn/serializer/migration/version_migration.py

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,19 @@
3636
update_data_objects,
3737
)
3838
from .version_1_4 import update_mi_states
39+
from .version_1_5 import update_event_gateway_children
40+
41+
def from_version_1_4(dct):
42+
"""Upgrade serialization from v1.4 to v1.5
43+
44+
Event based gateways now manage events for their children. Duplicating the events in
45+
multiple task specs was very difficult to manage and introduced a lot of problems with
46+
task states and predictions. The gateway will handle dropping the branches for the
47+
alternate events and the child branch with the matched event will proceed with a none
48+
event.
49+
"""
50+
dct['VERSION'] = "1.5"
51+
update_event_gateway_children(dct)
3952

4053
def from_version_1_3(dct):
4154
"""Upgrade serialization from v1.3 to v1.4
@@ -122,4 +135,5 @@ def from_version_1_0(dct):
122135
'1.1': from_version_1_1,
123136
'1.2': from_version_1_2,
124137
'1.3': from_version_1_3,
138+
'1.4': from_version_1_4,
125139
}

SpiffWorkflow/bpmn/serializer/workflow.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -26,7 +26,7 @@
2626
from .config import DEFAULT_CONFIG
2727

2828
# This is the default version set on the workflow, it can be overridden in init
29-
VERSION = "1.4"
29+
VERSION = "1.5"
3030

3131

3232
class BpmnWorkflowSerializer:
@@ -181,4 +181,4 @@ def from_dict(self, dct, **kwargs):
181181
Returns:
182182
a restored object
183183
"""
184-
return self.registry.restore(dct, **kwargs)
184+
return self.registry.restore(dct, **kwargs)

SpiffWorkflow/bpmn/specs/event_definitions/multiple.py

Lines changed: 20 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,4 @@
1+
from SpiffWorkflow.bpmn.util.event import BpmnEvent
12
from .timer import TimerEventDefinition, EventDefinition
23

34
class MultipleEventDefinition(EventDefinition):
@@ -11,27 +12,34 @@ def has_fired(self, my_task):
1112

1213
event_definitions = list(self.event_definitions)
1314
seen_events = my_task.internal_data.get('seen_events', [])
14-
for event_definition in self.event_definitions:
15+
for event in seen_events:
16+
if event.event_definition in event_definitions:
17+
event_definitions.remove(event.event_definition)
18+
19+
for event_definition in event_definitions:
1520
if isinstance(event_definition, TimerEventDefinition):
16-
child = [c for c in my_task.children if c.task_spec.event_definition == event_definition]
17-
child[0].task_spec._update_hook(child[0])
18-
if event_definition.has_fired(child[0]) and event_definition in event_definitions:
21+
if event_definition.has_fired(my_task):
1922
event_definitions.remove(event_definition)
20-
else:
21-
for event in seen_events:
22-
if event_definition.catches(my_task, event) and event_definition in event_definitions:
23-
event_definitions.remove(event_definition)
23+
seen_events.append(BpmnEvent(event_definition))
24+
25+
my_task.internal_data['seen_events'] = seen_events
2426

2527
if self.parallel:
2628
# Parallel multiple need to match all events
2729
return len(event_definitions) == 0
2830
else:
2931
return len(seen_events) > 0
3032

33+
def catches(self, my_task, event=None):
34+
for item in self.event_definitions:
35+
if item.catches(my_task, event):
36+
return True
37+
3138
def catch(self, my_task, event=None):
32-
event.event_definition.catch(my_task, event)
33-
seen_events = my_task.internal_data.get('seen_events', []) + [event]
34-
my_task._set_internal_data(seen_events=seen_events)
39+
for item in self.event_definitions:
40+
if item.catches(my_task, event):
41+
seen_events = my_task.internal_data.get('seen_events', []) + [event]
42+
my_task._set_internal_data(seen_events=seen_events)
3543

3644
def reset(self, my_task):
3745
my_task.internal_data.pop('seen_events', None)
@@ -47,4 +55,4 @@ def __eq__(self, other):
4755
def throw(self, my_task):
4856
# Mutiple events throw all associated events when they fire
4957
for event_definition in self.event_definitions:
50-
event_definition.throw(my_task)
58+
event_definition.throw(my_task)

SpiffWorkflow/bpmn/specs/mixins/events/intermediate_event.py

Lines changed: 19 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@
1818
# 02110-1301 USA
1919

2020
from SpiffWorkflow.util.task import TaskState
21+
from SpiffWorkflow.bpmn.specs.event_definitions import NoneEventDefinition
2122
from .event_types import ThrowingEvent, CatchingEvent
2223

2324

@@ -53,11 +54,24 @@ def catches(self, my_task, event):
5354

5455
class EventBasedGateway(CatchingEvent):
5556

57+
def connect(self, child):
58+
# Having the events duplicated in the gateway and the child is difficult to manage
59+
# Therefore, I am going to have the gateway manage the events and remove the event from the child
60+
super().connect(child)
61+
self.event_definition.event_definitions.append(child.event_definition)
62+
child.event_definition = NoneEventDefinition()
63+
5664
def _predict_hook(self, my_task):
57-
my_task._sync_children(self.outputs, state=TaskState.WAITING)
65+
my_task._sync_children(self.outputs, state=TaskState.MAYBE)
5866

59-
def _on_ready_hook(self, my_task):
60-
for child in my_task.children:
61-
if not child.internal_data.get('event_fired'):
62-
child.cancel()
67+
def _run_hook(self, my_task):
68+
seen_events = my_task.internal_data.get('seen_events', [])
69+
matches = []
70+
for event in seen_events:
71+
idx = self.event_definition.event_definitions.index(event.event_definition)
72+
matches.append(my_task.children[idx].task_spec)
6373

74+
my_task._sync_children(matches, TaskState.FUTURE)
75+
for child in my_task.children:
76+
child.task_spec._predict(child, mask=TaskState.FUTURE|TaskState.PREDICTED_MASK)
77+
return True

SpiffWorkflow/bpmn/workflow.py

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@
2525

2626
from SpiffWorkflow.bpmn.specs.control import BoundaryEventSplit
2727
from SpiffWorkflow.bpmn.specs.event_definitions.timer import TimerEventDefinition
28+
from SpiffWorkflow.bpmn.specs.event_definitions.multiple import MultipleEventDefinition
2829

2930
from SpiffWorkflow.bpmn.util.subworkflow import BpmnBaseWorkflow, BpmnSubWorkflow
3031
from SpiffWorkflow.bpmn.util.event import EventManager
@@ -181,7 +182,7 @@ def refresh_timers(self):
181182
# Ideally this would go in event manager but I can't import the necessary classes there
182183
# Eventually I'll move it
183184
for task in list(self.event_manager.tasks.values()):
184-
if isinstance(task.task_spec.event_definition, (TimerEventDefinition, )):
185+
if isinstance(task.task_spec.event_definition, (TimerEventDefinition, MultipleEventDefinition)):
185186
task.task_spec._update(task)
186187

187188
def get_task_from_id(self, task_id):

tests/SpiffWorkflow/bpmn/data/event-gateway.bpmn

Lines changed: 54 additions & 45 deletions
Original file line numberDiff line numberDiff line change
@@ -4,16 +4,16 @@
44
<bpmn:startEvent id="StartEvent_1">
55
<bpmn:outgoing>Flow_0w4b5t2</bpmn:outgoing>
66
</bpmn:startEvent>
7-
<bpmn:sequenceFlow id="Flow_0w4b5t2" sourceRef="StartEvent_1" targetRef="Gateway_1434v9l" />
7+
<bpmn:sequenceFlow id="Flow_0w4b5t2" sourceRef="StartEvent_1" targetRef="Gateway_1dxbnbw" />
88
<bpmn:eventBasedGateway id="Gateway_1434v9l">
9-
<bpmn:incoming>Flow_0w4b5t2</bpmn:incoming>
9+
<bpmn:incoming>Flow_04a2n0x</bpmn:incoming>
1010
<bpmn:outgoing>Flow_0gge7fn</bpmn:outgoing>
1111
<bpmn:outgoing>Flow_0px7ksu</bpmn:outgoing>
1212
<bpmn:outgoing>Flow_1rfbrlf</bpmn:outgoing>
1313
</bpmn:eventBasedGateway>
1414
<bpmn:intermediateCatchEvent id="message_1_event">
1515
<bpmn:incoming>Flow_0gge7fn</bpmn:incoming>
16-
<bpmn:outgoing>Flow_1g4g85l</bpmn:outgoing>
16+
<bpmn:outgoing>Flow_1xrb5hw</bpmn:outgoing>
1717
<bpmn:messageEventDefinition id="MessageEventDefinition_158nhox" messageRef="Message_0lyfmat" />
1818
</bpmn:intermediateCatchEvent>
1919
<bpmn:sequenceFlow id="Flow_0gge7fn" sourceRef="Gateway_1434v9l" targetRef="message_1_event" />
@@ -27,59 +27,32 @@
2727
<bpmn:incoming>Flow_1rfbrlf</bpmn:incoming>
2828
<bpmn:outgoing>Flow_0mppjk9</bpmn:outgoing>
2929
<bpmn:timerEventDefinition id="TimerEventDefinition_0reo0gl">
30-
<bpmn:timeDuration xsi:type="bpmn:tFormalExpression">"PT1S"</bpmn:timeDuration>
30+
<bpmn:timeDuration xsi:type="bpmn:tFormalExpression">"PT0.1S"</bpmn:timeDuration>
3131
</bpmn:timerEventDefinition>
3232
</bpmn:intermediateCatchEvent>
3333
<bpmn:sequenceFlow id="Flow_1rfbrlf" sourceRef="Gateway_1434v9l" targetRef="timer_event" />
34-
<bpmn:endEvent id="timer">
34+
<bpmn:endEvent id="timer_end">
3535
<bpmn:incoming>Flow_0mppjk9</bpmn:incoming>
3636
</bpmn:endEvent>
37-
<bpmn:sequenceFlow id="Flow_0mppjk9" sourceRef="timer_event" targetRef="timer" />
38-
<bpmn:endEvent id="message_1_end">
39-
<bpmn:incoming>Flow_1g4g85l</bpmn:incoming>
40-
</bpmn:endEvent>
41-
<bpmn:sequenceFlow id="Flow_1g4g85l" sourceRef="message_1_event" targetRef="message_1_end" />
37+
<bpmn:sequenceFlow id="Flow_0mppjk9" sourceRef="timer_event" targetRef="timer_end" />
4238
<bpmn:endEvent id="message_2_end">
4339
<bpmn:incoming>Flow_18v90rx</bpmn:incoming>
4440
</bpmn:endEvent>
4541
<bpmn:sequenceFlow id="Flow_18v90rx" sourceRef="message_2_event" targetRef="message_2_end" />
42+
<bpmn:exclusiveGateway id="Gateway_1dxbnbw" default="Flow_04a2n0x">
43+
<bpmn:incoming>Flow_0w4b5t2</bpmn:incoming>
44+
<bpmn:incoming>Flow_1xrb5hw</bpmn:incoming>
45+
<bpmn:outgoing>Flow_04a2n0x</bpmn:outgoing>
46+
</bpmn:exclusiveGateway>
47+
<bpmn:sequenceFlow id="Flow_04a2n0x" sourceRef="Gateway_1dxbnbw" targetRef="Gateway_1434v9l" />
48+
<bpmn:sequenceFlow id="Flow_1xrb5hw" sourceRef="message_1_event" targetRef="Gateway_1dxbnbw" />
4649
</bpmn:process>
4750
<bpmn:message id="Message_0lyfmat" name="message_1" />
4851
<bpmn:message id="Message_1ntpwce" name="message_2" />
4952
<bpmndi:BPMNDiagram id="BPMNDiagram_1">
5053
<bpmndi:BPMNPlane id="BPMNPlane_1" bpmnElement="Process_0pvx19v">
51-
<bpmndi:BPMNEdge id="Flow_18v90rx_di" bpmnElement="Flow_18v90rx">
52-
<di:waypoint x="408" y="230" />
53-
<di:waypoint x="472" y="230" />
54-
</bpmndi:BPMNEdge>
55-
<bpmndi:BPMNEdge id="Flow_1g4g85l_di" bpmnElement="Flow_1g4g85l">
56-
<di:waypoint x="408" y="117" />
57-
<di:waypoint x="472" y="117" />
58-
</bpmndi:BPMNEdge>
59-
<bpmndi:BPMNEdge id="Flow_0mppjk9_di" bpmnElement="Flow_0mppjk9">
60-
<di:waypoint x="408" y="340" />
61-
<di:waypoint x="472" y="340" />
62-
</bpmndi:BPMNEdge>
63-
<bpmndi:BPMNEdge id="Flow_1rfbrlf_di" bpmnElement="Flow_1rfbrlf">
64-
<di:waypoint x="290" y="142" />
65-
<di:waypoint x="290" y="340" />
66-
<di:waypoint x="372" y="340" />
67-
</bpmndi:BPMNEdge>
68-
<bpmndi:BPMNEdge id="Flow_0px7ksu_di" bpmnElement="Flow_0px7ksu">
69-
<di:waypoint x="290" y="142" />
70-
<di:waypoint x="290" y="230" />
71-
<di:waypoint x="372" y="230" />
72-
</bpmndi:BPMNEdge>
73-
<bpmndi:BPMNEdge id="Flow_0gge7fn_di" bpmnElement="Flow_0gge7fn">
74-
<di:waypoint x="315" y="117" />
75-
<di:waypoint x="372" y="117" />
76-
</bpmndi:BPMNEdge>
77-
<bpmndi:BPMNEdge id="Flow_0w4b5t2_di" bpmnElement="Flow_0w4b5t2">
78-
<di:waypoint x="215" y="117" />
79-
<di:waypoint x="265" y="117" />
80-
</bpmndi:BPMNEdge>
8154
<bpmndi:BPMNShape id="_BPMNShape_StartEvent_2" bpmnElement="StartEvent_1">
82-
<dc:Bounds x="179" y="99" width="36" height="36" />
55+
<dc:Bounds x="52" y="99" width="36" height="36" />
8356
</bpmndi:BPMNShape>
8457
<bpmndi:BPMNShape id="Gateway_0gplu2e_di" bpmnElement="Gateway_1434v9l">
8558
<dc:Bounds x="265" y="92" width="50" height="50" />
@@ -93,15 +66,51 @@
9366
<bpmndi:BPMNShape id="Event_1ea7gov_di" bpmnElement="timer_event">
9467
<dc:Bounds x="372" y="322" width="36" height="36" />
9568
</bpmndi:BPMNShape>
96-
<bpmndi:BPMNShape id="Event_0emrepu_di" bpmnElement="timer">
69+
<bpmndi:BPMNShape id="Event_0emrepu_di" bpmnElement="timer_end">
9770
<dc:Bounds x="472" y="322" width="36" height="36" />
9871
</bpmndi:BPMNShape>
99-
<bpmndi:BPMNShape id="Event_0x07cac_di" bpmnElement="message_1_end">
100-
<dc:Bounds x="472" y="99" width="36" height="36" />
101-
</bpmndi:BPMNShape>
10272
<bpmndi:BPMNShape id="Event_1p7slpj_di" bpmnElement="message_2_end">
10373
<dc:Bounds x="472" y="212" width="36" height="36" />
10474
</bpmndi:BPMNShape>
75+
<bpmndi:BPMNShape id="Gateway_1dxbnbw_di" bpmnElement="Gateway_1dxbnbw" isMarkerVisible="true">
76+
<dc:Bounds x="155" y="92" width="50" height="50" />
77+
</bpmndi:BPMNShape>
78+
<bpmndi:BPMNEdge id="Flow_0w4b5t2_di" bpmnElement="Flow_0w4b5t2">
79+
<di:waypoint x="88" y="117" />
80+
<di:waypoint x="155" y="117" />
81+
</bpmndi:BPMNEdge>
82+
<bpmndi:BPMNEdge id="Flow_0gge7fn_di" bpmnElement="Flow_0gge7fn">
83+
<di:waypoint x="315" y="117" />
84+
<di:waypoint x="372" y="117" />
85+
</bpmndi:BPMNEdge>
86+
<bpmndi:BPMNEdge id="Flow_0px7ksu_di" bpmnElement="Flow_0px7ksu">
87+
<di:waypoint x="290" y="142" />
88+
<di:waypoint x="290" y="230" />
89+
<di:waypoint x="372" y="230" />
90+
</bpmndi:BPMNEdge>
91+
<bpmndi:BPMNEdge id="Flow_1rfbrlf_di" bpmnElement="Flow_1rfbrlf">
92+
<di:waypoint x="290" y="142" />
93+
<di:waypoint x="290" y="340" />
94+
<di:waypoint x="372" y="340" />
95+
</bpmndi:BPMNEdge>
96+
<bpmndi:BPMNEdge id="Flow_0mppjk9_di" bpmnElement="Flow_0mppjk9">
97+
<di:waypoint x="408" y="340" />
98+
<di:waypoint x="472" y="340" />
99+
</bpmndi:BPMNEdge>
100+
<bpmndi:BPMNEdge id="Flow_18v90rx_di" bpmnElement="Flow_18v90rx">
101+
<di:waypoint x="408" y="230" />
102+
<di:waypoint x="472" y="230" />
103+
</bpmndi:BPMNEdge>
104+
<bpmndi:BPMNEdge id="Flow_04a2n0x_di" bpmnElement="Flow_04a2n0x">
105+
<di:waypoint x="205" y="117" />
106+
<di:waypoint x="265" y="117" />
107+
</bpmndi:BPMNEdge>
108+
<bpmndi:BPMNEdge id="Flow_1xrb5hw_di" bpmnElement="Flow_1xrb5hw">
109+
<di:waypoint x="390" y="99" />
110+
<di:waypoint x="390" y="0" />
111+
<di:waypoint x="180" y="0" />
112+
<di:waypoint x="180" y="92" />
113+
</bpmndi:BPMNEdge>
105114
</bpmndi:BPMNPlane>
106115
</bpmndi:BPMNDiagram>
107116
</bpmn:definitions>

0 commit comments

Comments
 (0)