Skip to content

Commit 6ef60bb

Browse files
authored
[Prism] Fix potential side-effect in TriggerAfterEach.onFire. (#36166)
* Fix potential side-effect in TriggerAfterEach.onFire. * Add a unit test.
1 parent d0e48e2 commit 6ef60bb

2 files changed

Lines changed: 28 additions & 1 deletion

File tree

sdks/go/pkg/beam/runners/prism/internal/engine/strategy.go

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -302,15 +302,23 @@ func (t *TriggerAfterEach) onFire(state *StateData) {
302302
if !t.shouldFire(state) {
303303
return
304304
}
305-
for _, sub := range t.SubTriggers {
305+
for i, sub := range t.SubTriggers {
306306
if state.getTriggerState(sub).finished {
307307
continue
308308
}
309309
sub.onFire(state)
310+
// If the sub-trigger didn't finish, we return, waiting for it to finish on a subsequent call.
310311
if !state.getTriggerState(sub).finished {
311312
return
312313
}
314+
315+
// If the sub-trigger finished, we check if it's the last one.
316+
// If it's not the last one, we return, waiting for the next onFire call to advance to the next sub-trigger.
317+
if i < len(t.SubTriggers)-1 {
318+
return
319+
}
313320
}
321+
// clear and reset when all sub-triggers have fired.
314322
triggerClearAndFinish(t, state)
315323
}
316324

sdks/go/pkg/beam/runners/prism/internal/engine/strategy_test.go

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -122,6 +122,25 @@ func TestTriggers_isReady(t *testing.T) {
122122
{triggerInput{newElementCount: 1}, false},
123123
{triggerInput{newElementCount: 1}, false},
124124
},
125+
}, {
126+
name: "afterEach_2_Always_1",
127+
trig: &TriggerAfterEach{
128+
SubTriggers: []Trigger{
129+
&TriggerElementCount{2},
130+
&TriggerAfterAny{SubTriggers: []Trigger{&TriggerAlways{}}},
131+
&TriggerElementCount{1},
132+
},
133+
},
134+
inputs: []io{
135+
{triggerInput{newElementCount: 1}, false},
136+
{triggerInput{newElementCount: 1}, true}, // first is ready
137+
{triggerInput{newElementCount: 1}, true}, // second is ready
138+
{triggerInput{newElementCount: 1}, true}, // third is ready
139+
{triggerInput{newElementCount: 1}, false}, // never resets after this.
140+
{triggerInput{newElementCount: 1}, false},
141+
{triggerInput{newElementCount: 1}, false},
142+
{triggerInput{newElementCount: 1}, false},
143+
},
125144
}, {
126145
name: "afterAny_2_3_4",
127146
trig: &TriggerAfterAny{

0 commit comments

Comments
 (0)