From cf242802e65e439f7e8d63615f4412143af50719 Mon Sep 17 00:00:00 2001 From: Janne Valkealahti Date: Sun, 25 Oct 2020 08:58:24 +0000 Subject: [PATCH] Use reactive lifecycle for triggers - Change to use reactive method for lifecycle handling for triggers in an executor. --- .../support/ReactiveStateMachineExecutor.java | 36 ++++++++++--------- 1 file changed, 20 insertions(+), 16 deletions(-) diff --git a/spring-statemachine-core/src/main/java/org/springframework/statemachine/support/ReactiveStateMachineExecutor.java b/spring-statemachine-core/src/main/java/org/springframework/statemachine/support/ReactiveStateMachineExecutor.java index e73db270..1e63f010 100644 --- a/spring-statemachine-core/src/main/java/org/springframework/statemachine/support/ReactiveStateMachineExecutor.java +++ b/spring-statemachine-core/src/main/java/org/springframework/statemachine/support/ReactiveStateMachineExecutor.java @@ -28,6 +28,7 @@ import java.util.Queue; import java.util.Set; import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.stream.Collectors; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -116,8 +117,7 @@ public class ReactiveStateMachineExecutor extends LifecycleObjectSupport i @Override protected Mono doPreStartReactively() { return Mono.defer(() -> { - Mono mono = Mono.empty(); - startTriggers(); + Mono mono = startTriggers(); if (triggerDisposable == null) { triggerDisposable = triggerFlux.subscribe(); @@ -140,14 +140,14 @@ public class ReactiveStateMachineExecutor extends LifecycleObjectSupport i @Override protected Mono doPreStopReactively() { - return Mono.fromRunnable(() -> { - stopTriggers(); + Mono mono = Mono.fromRunnable(() -> { if (triggerDisposable != null) { triggerDisposable.dispose(); triggerDisposable = null; } initialHandled.set(false); }); + return stopTriggers().and(mono); } @Override @@ -451,20 +451,24 @@ public class ReactiveStateMachineExecutor extends LifecycleObjectSupport i } } - private void startTriggers() { - for (final Trigger trigger : triggerToTransitionMap.keySet()) { - if (trigger instanceof Lifecycle) { - ((Lifecycle) trigger).start(); - } - } + private Mono startTriggers() { + List smrl = triggerToTransitionMap.keySet().stream() + .filter(StateMachineReactiveLifecycle.class::isInstance) + .map(StateMachineReactiveLifecycle.class::cast) + .collect(Collectors.toList()); + return Flux.fromIterable(smrl) + .flatMap(StateMachineReactiveLifecycle::startReactively) + .then(); } - private void stopTriggers() { - for (final Trigger trigger : triggerToTransitionMap.keySet()) { - if (trigger instanceof Lifecycle) { - ((Lifecycle) trigger).stop(); - } - } + private Mono stopTriggers() { + List smrl = triggerToTransitionMap.keySet().stream() + .filter(StateMachineReactiveLifecycle.class::isInstance) + .map(StateMachineReactiveLifecycle.class::cast) + .collect(Collectors.toList()); + return Flux.fromIterable(smrl) + .flatMap(StateMachineReactiveLifecycle::stopReactively) + .then(); } private class TriggerQueueItem {