This is an automated email from the ASF dual-hosted git repository.
pefernan pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/incubator-kie.git
The following commit(s) were added to refs/heads/main by this push:
new a92f9bcce9e [incubator-kie#7094] Avoid redundant process instance
unmarshalling when signaling via `ProcessService` (#7092)
a92f9bcce9e is described below
commit a92f9bcce9ed00f58e3576c7ff91e769a8035322
Author: Pere Fernández <[email protected]>
AuthorDate: Tue Sep 15 11:41:17 2026 +0200
[incubator-kie#7094] Avoid redundant process instance unmarshalling when
signaling via `ProcessService` (#7092)
* [incubator-kie#7094] Avoid redundant process instance unmarshalling when
signaling via ProcessService
* fix: resolve variable expressions in EventNode event descriptions for
correct signal routing; remove acceptingEventType from ProcessInstances and
drop eventTypes from ProcessInstance API
* Add fail fast condition to avoid silently handling null signalName in
ProcessServiceImpl
* Removed unnecessary 'Message' check in ProcessServiceImpl, added null
check
---
.../org/kie/kogito/process/ProcessInstances.java | 16 --
.../instance/impl/WorkflowProcessInstanceImpl.java | 40 ++--
.../workflow/instance/node/EventNodeInstance.java | 23 +-
.../kogito/process/impl/ProcessServiceImpl.java | 15 +-
.../ProcessInstancesAcceptingEventTypeTest.java | 257 ---------------------
.../process/impl/ProcessServiceImplSignalTest.java | 217 ++++++++---------
.../java/org/jbpm/bpmn2/IntermediateEventTest.java | 10 +
.../java/org/jbpm/bpmn2/SLAComplianceTest.java | 2 -
8 files changed, 152 insertions(+), 428 deletions(-)
diff --git
a/kogito-api/kogito-api/src/main/java/org/kie/kogito/process/ProcessInstances.java
b/kogito-api/kogito-api/src/main/java/org/kie/kogito/process/ProcessInstances.java
index a6bf3c6dc34..9570d167f06 100644
---
a/kogito-api/kogito-api/src/main/java/org/kie/kogito/process/ProcessInstances.java
+++
b/kogito-api/kogito-api/src/main/java/org/kie/kogito/process/ProcessInstances.java
@@ -56,20 +56,4 @@ public interface ProcessInstances<T> {
}
Stream<ProcessInstance<T>> waitingForEventType(String eventType,
ProcessInstanceReadMode mode);
-
- default Stream<ProcessInstance<T>> acceptingEventType(String signalName,
String id) {
- return findById(id, ProcessInstanceReadMode.MUTABLE)
- .filter(pi -> {
- // Check if waiting for event (traditional signal event)
- boolean isWaitingForSignal = Stream.concat(
- waitingForEventType(signalName,
ProcessInstanceReadMode.READ_ONLY),
- waitingForEventType("Message-" + signalName,
ProcessInstanceReadMode.READ_ONLY)).anyMatch(p -> p.id().equals(id));
-
- boolean isAdHocNode = pi.adHocFragments().stream()
- .anyMatch(fragment ->
fragment.getName().equals(signalName));
-
- return isWaitingForSignal || isAdHocNode;
- })
- .stream();
- }
}
diff --git
a/kogito-jbpm/jbpm-flow/src/main/java/org/jbpm/workflow/instance/impl/WorkflowProcessInstanceImpl.java
b/kogito-jbpm/jbpm-flow/src/main/java/org/jbpm/workflow/instance/impl/WorkflowProcessInstanceImpl.java
index 8e296db9c7e..9a30f446515 100755
---
a/kogito-jbpm/jbpm-flow/src/main/java/org/jbpm/workflow/instance/impl/WorkflowProcessInstanceImpl.java
+++
b/kogito-jbpm/jbpm-flow/src/main/java/org/jbpm/workflow/instance/impl/WorkflowProcessInstanceImpl.java
@@ -1016,7 +1016,7 @@ public abstract class WorkflowProcessInstanceImpl extends
ProcessInstanceImpl im
return Collections.emptySet();
}
VariableScope variableScope = (VariableScope) ((ContextContainer)
getProcess()).getDefaultContext(VariableScope.VARIABLE_SCOPE);
- Set<EventDescription<?>> eventDesciptions = new LinkedHashSet<>();
+ Set<EventDescription<?>> eventDescriptions = new LinkedHashSet<>();
List<KogitoEventListener> activeListeners =
eventListeners.values().stream()
.flatMap(List::stream)
@@ -1024,9 +1024,9 @@ public abstract class WorkflowProcessInstanceImpl extends
ProcessInstanceImpl im
activeListeners.addAll(externalEventListeners.values().stream()
.flatMap(List::stream)
- .collect(Collectors.toList()));
+ .toList());
- activeListeners.forEach(el ->
eventDesciptions.addAll(el.getEventDescriptions()));
+ activeListeners.forEach(el ->
eventDescriptions.addAll(el.getEventDescriptions()));
((org.jbpm.workflow.core.WorkflowProcess)
getProcess()).getNodesRecursively().stream().filter(n -> n instanceof
EventNodeInterface).forEach(n -> {
@@ -1037,8 +1037,7 @@ public abstract class WorkflowProcessInstanceImpl extends
ProcessInstanceImpl im
dataType = new NamedDataType(eventVar.getName(),
eventVar.getType());
}
}
- if (n instanceof BoundaryEventNode) {
- BoundaryEventNode boundaryEventNode = (BoundaryEventNode) n;
+ if (n instanceof BoundaryEventNode boundaryEventNode) {
StateBasedNodeInstance attachedToNodeInstance =
(StateBasedNodeInstance) getNodeInstances(true).stream()
.filter(ni ->
ni.getNode().getUniqueId().equals(boundaryEventNode.getAttachedToNodeId())).findFirst().orElse(null);
if (attachedToNodeInstance != null) {
@@ -1054,12 +1053,11 @@ public abstract class WorkflowProcessInstanceImpl
extends ProcessInstanceImpl im
eventName = TIMER_TRIGGERED_EVENT;
}
- eventDesciptions.add(new BaseEventDescription(eventName,
n.getUniqueId(), n.getName(), eventType, null, getStringId(), dataType,
properties));
+ eventDescriptions.add(new BaseEventDescription(eventName,
n.getUniqueId(), n.getName(), eventType, null, getStringId(), dataType,
properties));
}
- } else if (n instanceof EventSubProcessNode) {
- EventSubProcessNode eventSubProcessNode =
(EventSubProcessNode) n;
+ } else if (n instanceof EventSubProcessNode eventSubProcessNode) {
org.kie.api.definition.process.Node startNode =
eventSubProcessNode.findStartNode();
Map<Timer, DroolsAction> timers =
eventSubProcessNode.getTimers();
if (timers != null && !timers.isEmpty()) {
@@ -1067,33 +1065,23 @@ public abstract class WorkflowProcessInstanceImpl
extends ProcessInstanceImpl im
Map<String, String> timerProperties =
((StateBasedNodeInstance) ni).extractTimerEventInformation();
if (timerProperties != null) {
-
- eventDesciptions.add(new
BaseEventDescription(TIMER_TRIGGERED_EVENT, (String) startNode.getUniqueId(),
startNode.getName(), "timer", ni.getStringId(),
+ eventDescriptions.add(new
BaseEventDescription(TIMER_TRIGGERED_EVENT, startNode.getUniqueId(),
startNode.getName(), "timer", ni.getStringId(),
getStringId(), null, timerProperties));
}
});
} else {
-
for (String eventName : eventSubProcessNode.getEvents()) {
-
- eventDesciptions.add(new
BaseEventDescription(eventName, (String) startNode.getUniqueId(),
startNode.getName(), "signal", null, getStringId(), dataType));
+ eventDescriptions.add(new
BaseEventDescription(eventName, startNode.getUniqueId(), startNode.getName(),
"signal", null, getStringId(), dataType));
}
-
}
} else if (n instanceof EventNode) {
- NamedDataType finalDataType = dataType;
- getNodeInstances(n.getId()).forEach(ni -> eventDesciptions.add(
- new BaseEventDescription(
- ((EventNode) n).getType(),
- n.getUniqueId(),
- n.getName(),
- (String)
n.getMetaData().getOrDefault(EVENT_TYPE, EVENT_TYPE_SIGNAL),
- ni.getStringId(),
- getStringId(),
- finalDataType)));
+ getNodeInstances(n.getId()).stream()
+ .filter(KogitoEventListener.class::isInstance)
+ .map(KogitoEventListener.class::cast)
+ .forEach(ni ->
eventDescriptions.addAll(ni.getEventDescriptions()));
} else if (n instanceof StateNode) {
- getNodeInstances(n.getId()).forEach(ni -> eventDesciptions.add(
+ getNodeInstances(n.getId()).forEach(ni ->
eventDescriptions.add(
new BaseEventDescription(
(String) n.getMetaData().get(CONDITION),
n.getUniqueId(),
@@ -1106,7 +1094,7 @@ public abstract class WorkflowProcessInstanceImpl extends
ProcessInstanceImpl im
});
- return eventDesciptions;
+ return eventDescriptions;
}
@Override
diff --git
a/kogito-jbpm/jbpm-flow/src/main/java/org/jbpm/workflow/instance/node/EventNodeInstance.java
b/kogito-jbpm/jbpm-flow/src/main/java/org/jbpm/workflow/instance/node/EventNodeInstance.java
index 10cc2587e0d..596aa9c77eb 100755
---
a/kogito-jbpm/jbpm-flow/src/main/java/org/jbpm/workflow/instance/node/EventNodeInstance.java
+++
b/kogito-jbpm/jbpm-flow/src/main/java/org/jbpm/workflow/instance/node/EventNodeInstance.java
@@ -27,6 +27,7 @@ import java.util.regex.Matcher;
import org.jbpm.process.core.context.variable.Variable;
import org.jbpm.process.core.context.variable.VariableScope;
import org.jbpm.process.instance.InternalProcessRuntime;
+import org.jbpm.process.instance.context.variable.VariableScopeInstance;
import org.jbpm.util.PatternConstants;
import org.jbpm.workflow.core.Node;
import org.jbpm.workflow.core.impl.NodeIoHelper;
@@ -44,6 +45,8 @@ import org.kie.kogito.process.NamedDataType;
import org.kie.kogito.timer.TimerInstance;
import static java.util.Objects.isNull;
+import static org.jbpm.ruleflow.core.Metadata.EVENT_TYPE;
+import static org.jbpm.ruleflow.core.Metadata.EVENT_TYPE_SIGNAL;
import static
org.jbpm.workflow.instance.impl.DummyEventListener.EMPTY_EVENT_LISTENER;
import static
org.jbpm.workflow.instance.node.TimerNodeInstance.TIMER_TRIGGERED_EVENT;
import static org.kie.kogito.internal.utils.ConversionUtils.isEmpty;
@@ -272,12 +275,22 @@ public class EventNodeInstance extends
ExtendedNodeInstanceImpl implements Kogit
@Override
public Set<EventDescription<?>> getEventDescriptions() {
NamedDataType dataType = null;
- if (getEventNode().getVariableName() != null) {
- VariableScope variableScope = (VariableScope)
getEventNode().getContext(VariableScope.VARIABLE_SCOPE);
- Variable variable =
variableScope.findVariable(getEventNode().getVariableName());
- dataType = new NamedDataType(variable.getName(),
variable.getType());
+ String variableName = getEventNode().getVariableName();
+ if (variableName != null) {
+ VariableScopeInstance variableScopeInstance =
(VariableScopeInstance) resolveContextInstance(VariableScope.VARIABLE_SCOPE,
variableName);
+ if (variableScopeInstance == null) {
+ variableScopeInstance = (VariableScopeInstance)
getProcessInstance().getContextInstance(VariableScope.VARIABLE_SCOPE);
+ }
+ if (variableScopeInstance != null) {
+ Variable variable =
variableScopeInstance.getVariableScope().findVariable(variableName);
+ if (variable != null) {
+ dataType = new NamedDataType(variable.getName(),
variable.getType());
+ }
+ }
}
- return Collections.singleton(new BaseEventDescription(getEventType(),
getNodeDefinitionId(), getNodeName(), "signal", getStringId(),
getProcessInstance().getStringId(), dataType));
+ EventNode eventNode = getEventNode();
+ return Collections.singleton(new BaseEventDescription(getEventType(),
eventNode.getUniqueId(), eventNode.getName(),
+ (String) eventNode.getMetaData().getOrDefault(EVENT_TYPE,
EVENT_TYPE_SIGNAL), getStringId(), getProcessInstance().getStringId(),
dataType));
}
@Override
diff --git
a/kogito-jbpm/jbpm-flow/src/main/java/org/kie/kogito/process/impl/ProcessServiceImpl.java
b/kogito-jbpm/jbpm-flow/src/main/java/org/kie/kogito/process/impl/ProcessServiceImpl.java
index 8198e7e0b22..abf3f87bb36 100644
---
a/kogito-jbpm/jbpm-flow/src/main/java/org/kie/kogito/process/impl/ProcessServiceImpl.java
+++
b/kogito-jbpm/jbpm-flow/src/main/java/org/kie/kogito/process/impl/ProcessServiceImpl.java
@@ -20,6 +20,7 @@ package org.kie.kogito.process.impl;
import java.util.List;
import java.util.Map;
+import java.util.Objects;
import java.util.Optional;
import java.util.function.Function;
import java.util.function.Predicate;
@@ -160,11 +161,19 @@ public class ProcessServiceImpl implements ProcessService
{
@Override
public <T extends MappableToModel<R>, R> Optional<R>
signalProcessInstance(Process<T> process, String id, Object data, String
signalName) {
+ Objects.requireNonNull(signalName, "signalName must not be null");
return UnitOfWorkExecutor.executeInUnitOfWork(
application.unitOfWorkManager(),
- () -> process
- .instances().acceptingEventType(signalName, id)
- .findFirst()
+ () -> process.instances()
+ .findById(id)
+ .filter(pi -> {
+ if (pi.events().stream().anyMatch(e ->
signalName.equals(e.getEvent()))) {
+ return true;
+ }
+ return pi.adHocFragments()
+ .stream()
+ .anyMatch(f ->
f.getName().equals(signalName));
+ })
.map(pi -> {
pi.send(SignalFactory.of(signalName, data));
return pi.checkError().variables().toModel();
diff --git
a/kogito-jbpm/jbpm-flow/src/test/java/org/kie/kogito/process/impl/ProcessInstancesAcceptingEventTypeTest.java
b/kogito-jbpm/jbpm-flow/src/test/java/org/kie/kogito/process/impl/ProcessInstancesAcceptingEventTypeTest.java
deleted file mode 100644
index ad76444fb26..00000000000
---
a/kogito-jbpm/jbpm-flow/src/test/java/org/kie/kogito/process/impl/ProcessInstancesAcceptingEventTypeTest.java
+++ /dev/null
@@ -1,257 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements. See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership. The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * "License"); you may not use this file except in compliance
- * with the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing,
- * software distributed under the License is distributed on an
- * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
- * KIND, either express or implied. See the License for the
- * specific language governing permissions and limitations
- * under the License.
- */
-package org.kie.kogito.process.impl;
-
-import java.util.Collections;
-import java.util.stream.Stream;
-
-import org.jbpm.workflow.instance.WorkflowProcessInstance;
-import org.junit.jupiter.api.BeforeEach;
-import org.junit.jupiter.api.Test;
-import org.kie.kogito.Model;
-import org.kie.kogito.process.ProcessInstance;
-import org.kie.kogito.process.flexible.AdHocFragment;
-
-import static org.assertj.core.api.Assertions.assertThat;
-import static org.mockito.ArgumentMatchers.any;
-import static org.mockito.Mockito.mock;
-import static org.mockito.Mockito.when;
-
-/**
- * Unit tests for acceptingEventType() method in ProcessInstances
implementations.
- * Tests both traditional signal events and ad hoc node triggering.
- */
-class ProcessInstancesAcceptingEventTypeTest {
-
- private MapProcessInstances<Model> processInstances;
- private AbstractProcess<Model> process;
- private WorkflowProcessInstance workflowProcessInstance;
- private AbstractProcessInstance<Model> processInstance;
-
- @BeforeEach
- void setup() {
- process = mock(AbstractProcess.class);
- processInstances = new MapProcessInstances<>(process);
- workflowProcessInstance = mock(WorkflowProcessInstance.class);
- processInstance = mock(AbstractProcessInstance.class);
- }
-
- @Test
- void testAcceptingEventType_WaitingForSignalEvent() {
- // Given: Process instance waiting for a signal event
- String processInstanceId = "test-accepting-event-instance-1";
- String signalName = "HelloMartin";
-
-
when(workflowProcessInstance.getStringId()).thenReturn(processInstanceId);
- when(workflowProcessInstance.getEventTypes()).thenReturn(new String[]
{ signalName, "OtherSignal" });
- when(processInstance.id()).thenReturn(processInstanceId);
-
when(processInstance.internalGetProcessInstance()).thenReturn(workflowProcessInstance);
-
when(processInstance.adHocFragments()).thenReturn(Collections.emptyList());
-
when(process.createInstance(any(WorkflowProcessInstance.class))).thenAnswer(invocation
-> processInstance);
-
when(process.createReadOnlyInstance(any(WorkflowProcessInstance.class))).thenAnswer(invocation
-> processInstance);
-
- // Store the instance
- processInstances.create(processInstanceId, processInstance);
-
- // When: Check if accepting the signal
- Stream<ProcessInstance<Model>> result =
processInstances.acceptingEventType(signalName, processInstanceId);
-
- // Then: Should return the process instance
- assertThat(result).isNotNull();
- assertThat(result.count()).isEqualTo(1);
- }
-
- @Test
- void testAcceptingEventType_AdHocNode() {
- // Given: Process instance with ad hoc node
- String processInstanceId = "test-accepting-event-instance-2";
- String adHocNodeName = "AdHocTask";
-
- AdHocFragment adHocFragment = new AdHocFragment.Builder(adHocNodeName)
- .withName(adHocNodeName)
- .withAutoStart(false)
- .build();
-
-
when(workflowProcessInstance.getStringId()).thenReturn(processInstanceId);
- when(workflowProcessInstance.getEventTypes()).thenReturn(new String[]
{});
- when(processInstance.id()).thenReturn(processInstanceId);
-
when(processInstance.internalGetProcessInstance()).thenReturn(workflowProcessInstance);
-
when(processInstance.adHocFragments()).thenReturn(Collections.singletonList(adHocFragment));
-
when(process.createInstance(any(WorkflowProcessInstance.class))).thenReturn(processInstance);
-
- // Store the instance
- processInstances.create(processInstanceId, processInstance);
-
- // When: Check if accepting the ad hoc node signal
- Stream<ProcessInstance<Model>> result =
processInstances.acceptingEventType(adHocNodeName, processInstanceId);
-
- // Then: Should return the process instance
- assertThat(result).isNotNull();
- assertThat(result.count()).isEqualTo(1);
- }
-
- @Test
- void testAcceptingEventType_BothSignalAndAdHoc() {
- // Given: Process instance with both signal event and ad hoc node
- String processInstanceId = "test-accepting-event-instance-3";
- String signalName = "HelloMartin";
- String adHocNodeName = "AdHocTask";
-
- AdHocFragment adHocFragment = new AdHocFragment.Builder(adHocNodeName)
- .withName(adHocNodeName)
- .withAutoStart(false)
- .build();
-
-
when(workflowProcessInstance.getStringId()).thenReturn(processInstanceId);
- when(workflowProcessInstance.getEventTypes()).thenReturn(new String[]
{ signalName });
- when(processInstance.id()).thenReturn(processInstanceId);
-
when(processInstance.internalGetProcessInstance()).thenReturn(workflowProcessInstance);
-
when(processInstance.adHocFragments()).thenReturn(Collections.singletonList(adHocFragment));
- // Ensure the mock always returns a valid instance, not just on first
call
-
when(process.createInstance(any(WorkflowProcessInstance.class))).thenAnswer(invocation
-> processInstance);
-
when(process.createReadOnlyInstance(any(WorkflowProcessInstance.class))).thenAnswer(invocation
-> processInstance);
-
- // Store the instance
- processInstances.create(processInstanceId, processInstance);
-
- // When: Check if accepting the signal
- Stream<ProcessInstance<Model>> resultSignal =
processInstances.acceptingEventType(signalName, processInstanceId);
-
- // Then: Should return the process instance for signal
- assertThat(resultSignal).isNotNull();
- assertThat(resultSignal.count()).isEqualTo(1);
-
- // When: Check if accepting the ad hoc node
- Stream<ProcessInstance<Model>> resultAdHoc =
processInstances.acceptingEventType(adHocNodeName, processInstanceId);
-
- // Then: Should return the process instance for ad hoc node
- assertThat(resultAdHoc).isNotNull();
- assertThat(resultAdHoc.count()).isEqualTo(1);
- }
-
- @Test
- void testAcceptingEventType_NotAccepting() {
- // Given: Process instance not waiting for signal and no ad hoc node
- String processInstanceId = "test-accepting-event-instance-4";
- String waitingFor = "HelloMartin";
- String notWaitingFor = "ByeMartin";
-
-
when(workflowProcessInstance.getStringId()).thenReturn(processInstanceId);
- when(workflowProcessInstance.getEventTypes()).thenReturn(new String[]
{ waitingFor });
- when(processInstance.id()).thenReturn(processInstanceId);
-
when(processInstance.internalGetProcessInstance()).thenReturn(workflowProcessInstance);
-
when(processInstance.adHocFragments()).thenReturn(Collections.emptyList());
-
when(process.createInstance(any(WorkflowProcessInstance.class))).thenReturn(processInstance);
-
- // Store the instance
- processInstances.create(processInstanceId, processInstance);
-
- // When: Check if accepting a signal it's not waiting for
- Stream<ProcessInstance<Model>> result =
processInstances.acceptingEventType(notWaitingFor, processInstanceId);
-
- // Then: Should return empty stream
- assertThat(result).isNotNull();
- assertThat(result.count()).isZero();
- }
-
- @Test
- void testAcceptingEventType_NonExistentInstance() {
- // Given: Non-existent process instance
- String processInstanceId = "non-existent";
- String signalName = "HelloMartin";
-
- // When: Check if accepting signal for non-existent instance
- Stream<ProcessInstance<Model>> result =
processInstances.acceptingEventType(signalName, processInstanceId);
-
- // Then: Should return empty stream
- assertThat(result).isNotNull();
- assertThat(result.count()).isZero();
- }
-
- @Test
- void testAcceptingEventType_MultipleAdHocNodes() {
- // Given: Process instance with multiple ad hoc nodes
- String processInstanceId = "test-accepting-event-instance-5";
- String adHocNode1 = "AdHocTask1";
- String adHocNode2 = "AdHocTask2";
- String adHocNode3 = "AdHocTask3";
-
- AdHocFragment fragment1 = new AdHocFragment.Builder(adHocNode1)
- .withName(adHocNode1)
- .withAutoStart(false)
- .build();
- AdHocFragment fragment2 = new AdHocFragment.Builder(adHocNode2)
- .withName(adHocNode2)
- .withAutoStart(true)
- .build();
- AdHocFragment fragment3 = new AdHocFragment.Builder(adHocNode3)
- .withName(adHocNode3)
- .withAutoStart(false)
- .build();
-
-
when(workflowProcessInstance.getStringId()).thenReturn(processInstanceId);
- when(workflowProcessInstance.getEventTypes()).thenReturn(new String[]
{});
- when(processInstance.id()).thenReturn(processInstanceId);
-
when(processInstance.internalGetProcessInstance()).thenReturn(workflowProcessInstance);
-
when(processInstance.adHocFragments()).thenReturn(java.util.Arrays.asList(fragment1,
fragment2, fragment3));
-
when(process.createInstance(any(WorkflowProcessInstance.class))).thenReturn(processInstance);
-
- // Store the instance
- processInstances.create(processInstanceId, processInstance);
-
- // When/Then: Check each ad hoc node
- Stream<ProcessInstance<Model>> result1 =
processInstances.acceptingEventType(adHocNode1, processInstanceId);
- assertThat(result1.count()).isEqualTo(1);
-
- Stream<ProcessInstance<Model>> result2 =
processInstances.acceptingEventType(adHocNode2, processInstanceId);
- assertThat(result2.count()).isEqualTo(1);
-
- Stream<ProcessInstance<Model>> result3 =
processInstances.acceptingEventType(adHocNode3, processInstanceId);
- assertThat(result3.count()).isEqualTo(1);
-
- // When: Check non-existent ad hoc node
- Stream<ProcessInstance<Model>> resultNone =
processInstances.acceptingEventType("NonExistent", processInstanceId);
- assertThat(resultNone.count()).isZero();
- }
-
- @Test
- void testAcceptingEventType_NullEventTypes() {
- // Given: Process instance with null event types
- String processInstanceId = "test-accepting-event-instance-6";
- String signalName = "HelloMartin";
-
-
when(workflowProcessInstance.getStringId()).thenReturn(processInstanceId);
- when(workflowProcessInstance.getEventTypes()).thenReturn(new String[]
{ null, signalName, null });
- when(processInstance.id()).thenReturn(processInstanceId);
-
when(processInstance.internalGetProcessInstance()).thenReturn(workflowProcessInstance);
-
when(processInstance.adHocFragments()).thenReturn(Collections.emptyList());
-
when(process.createInstance(any(WorkflowProcessInstance.class))).thenAnswer(invocation
-> processInstance);
-
when(process.createReadOnlyInstance(any(WorkflowProcessInstance.class))).thenAnswer(invocation
-> processInstance);
-
- // Store the instance
- processInstances.create(processInstanceId, processInstance);
-
- // When: Check if accepting the signal
- Stream<ProcessInstance<Model>> result =
processInstances.acceptingEventType(signalName, processInstanceId);
-
- // Then: Should handle null values and return the instance
- assertThat(result).isNotNull();
- assertThat(result.count()).isEqualTo(1);
- }
-}
diff --git
a/kogito-jbpm/jbpm-flow/src/test/java/org/kie/kogito/process/impl/ProcessServiceImplSignalTest.java
b/kogito-jbpm/jbpm-flow/src/test/java/org/kie/kogito/process/impl/ProcessServiceImplSignalTest.java
index 5533ee20d7f..3b7c6e735d2 100644
---
a/kogito-jbpm/jbpm-flow/src/test/java/org/kie/kogito/process/impl/ProcessServiceImplSignalTest.java
+++
b/kogito-jbpm/jbpm-flow/src/test/java/org/kie/kogito/process/impl/ProcessServiceImplSignalTest.java
@@ -19,8 +19,9 @@
package org.kie.kogito.process.impl;
import java.util.Collections;
+import java.util.List;
import java.util.Optional;
-import java.util.stream.Stream;
+import java.util.Set;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
@@ -28,43 +29,43 @@ import org.kie.kogito.Application;
import org.kie.kogito.MappableToModel;
import org.kie.kogito.Model;
import org.kie.kogito.config.ConfigBean;
+import org.kie.kogito.process.BaseEventDescription;
+import org.kie.kogito.process.EventDescription;
import org.kie.kogito.process.Process;
import org.kie.kogito.process.ProcessInstance;
-import org.kie.kogito.process.ProcessInstanceReadMode;
import org.kie.kogito.process.ProcessInstances;
+import org.kie.kogito.process.flexible.AdHocFragment;
import org.kie.kogito.uow.UnitOfWorkManager;
-import org.kie.kogito.uow.WorkUnit;
import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.ArgumentMatchers.any;
-import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.doAnswer;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
/**
- * Integration tests for ProcessServiceImpl signal handling.
- * Tests validation of signals for both traditional signal events and ad hoc
nodes.
+ * Unit tests for ProcessServiceImpl.signalProcessInstance.
+ * Verifies signal routing for traditional signal events, message events, ad
hoc nodes,
+ * non-existent instances, and non-matching signals.
*/
class ProcessServiceImplSignalTest {
private ProcessServiceImpl processService;
- private Application application;
private Process<TestModel> process;
private ProcessInstances<TestModel> processInstances;
private ProcessInstance<TestModel> processInstance;
- private UnitOfWorkManager unitOfWorkManager;
- private ConfigBean configBean;
@BeforeEach
+ @SuppressWarnings("unchecked")
void setup() {
- application = mock(Application.class);
+ Application application = mock(Application.class);
process = mock(Process.class);
processInstances = mock(ProcessInstances.class);
processInstance = mock(ProcessInstance.class);
- unitOfWorkManager = mock(UnitOfWorkManager.class);
- configBean = mock(ConfigBean.class);
+ UnitOfWorkManager unitOfWorkManager = mock(UnitOfWorkManager.class);
+ ConfigBean configBean = mock(ConfigBean.class);
when(application.unitOfWorkManager()).thenReturn(unitOfWorkManager);
when(application.config()).thenReturn(mock(org.kie.kogito.Config.class));
@@ -72,166 +73,144 @@ class ProcessServiceImplSignalTest {
when(configBean.processInstanceLimit()).thenReturn((short) 100);
when(process.instances()).thenReturn(processInstances);
- // Setup UnitOfWorkManager to execute code immediately
+ // Make the UoW execute the supplied callable immediately
org.kie.kogito.uow.UnitOfWork unitOfWork =
mock(org.kie.kogito.uow.UnitOfWork.class);
when(unitOfWorkManager.newUnitOfWork()).thenReturn(unitOfWork);
when(unitOfWorkManager.currentUnitOfWork()).thenReturn(unitOfWork);
doAnswer(invocation -> {
- org.kie.kogito.uow.WorkUnit<?> workUnit =
invocation.getArgument(0);
- workUnit.perform();
+ invocation.<org.kie.kogito.uow.WorkUnit<?>>
getArgument(0).perform();
return null;
}).when(unitOfWork).intercept(any());
processService = new ProcessServiceImpl(application);
}
+ // --- helpers ---
+
+ private static EventDescription<?> eventDesc(String eventName) {
+ return new BaseEventDescription(eventName, "node-1", "Node", "signal",
"ni-1", "pi-1", null);
+ }
+
+ private void givenInstanceWithEvents(String id, Set<EventDescription<?>>
events, List<AdHocFragment> adHocFragments) {
+
when(processInstances.findById(id)).thenReturn(Optional.of(processInstance));
+ when(processInstance.events()).thenReturn(events);
+ when(processInstance.adHocFragments()).thenReturn(adHocFragments);
+ when(processInstance.checkError()).thenReturn(processInstance);
+ when(processInstance.variables()).thenReturn(new TestModel());
+ }
+
+ // --- signal event tests ---
+
@Test
- void testSignalProcessInstance_WithTraditionalSignalEvent() {
- // Given: Process instance waiting for a signal event
- String processInstanceId = "test-signal-instance1";
+ void signalAccepted_whenInstanceWaitingForSignalEvent() {
+ String id = "pi-1";
String signalName = "HelloMartin";
- Object signalData = "test-data";
- TestModel model = new TestModel();
- when(processInstances.findById(eq(processInstanceId),
any(ProcessInstanceReadMode.class)))
- .thenReturn(Optional.of(processInstance));
- when(processInstances.acceptingEventType(signalName,
processInstanceId))
- .thenReturn(Stream.of(processInstance));
- when(processInstance.checkError()).thenReturn(processInstance);
- when(processInstance.variables()).thenReturn(model);
+ givenInstanceWithEvents(id, Set.of(eventDesc(signalName)),
Collections.emptyList());
- // When: Signal is sent
- Optional<TestModel> result =
processService.signalProcessInstance(process, processInstanceId, signalData,
signalName);
+ Optional<TestModel> result =
processService.signalProcessInstance(process, id, "data", signalName);
- // Then: Signal should be accepted
assertThat(result).isPresent();
verify(processInstance).send(any());
}
@Test
- void testSignalProcessInstance_WithAdHocNode() {
- // Given: Process instance with ad hoc node
- String processInstanceId = "test-signal-instance2";
- String adHocNodeName = "AdHocTask";
- Object signalData = Collections.emptyMap();
- TestModel model = new TestModel();
-
- when(processInstances.findById(eq(processInstanceId),
any(ProcessInstanceReadMode.class)))
- .thenReturn(Optional.of(processInstance));
- when(processInstances.acceptingEventType(adHocNodeName,
processInstanceId))
- .thenReturn(Stream.of(processInstance));
- when(processInstance.checkError()).thenReturn(processInstance);
- when(processInstance.variables()).thenReturn(model);
+ void signalAccepted_whenInstanceHasMatchingAdHocFragment() {
+ String id = "pi-3";
+ String adHocName = "AdHocTask";
+ AdHocFragment fragment = new
AdHocFragment.Builder(adHocName).withName(adHocName).withAutoStart(false).build();
+
+ givenInstanceWithEvents(id, Collections.emptySet(), List.of(fragment));
- // When: Signal is sent to trigger ad hoc node
- Optional<TestModel> result =
processService.signalProcessInstance(process, processInstanceId, signalData,
adHocNodeName);
+ Optional<TestModel> result =
processService.signalProcessInstance(process, id, null, adHocName);
- // Then: Signal should be accepted
assertThat(result).isPresent();
verify(processInstance).send(any());
}
@Test
- void testSignalProcessInstance_InvalidSignal_ReturnsEmpty() {
- // Given: Process instance not accepting the signal
- String processInstanceId = "test-signal-instance3";
- String invalidSignal = "InvalidSignal";
- Object signalData = "test-data";
+ void signalAccepted_whenInstanceHasBothSignalEventAndAdHocFragment() {
+ String id = "pi-4";
+ String signalName = "HelloMartin";
+ String adHocName = "AdHocTask";
+ AdHocFragment fragment = new
AdHocFragment.Builder(adHocName).withName(adHocName).withAutoStart(false).build();
+
+ givenInstanceWithEvents(id, Set.of(eventDesc(signalName)),
List.of(fragment));
- when(processInstances.findById(eq(processInstanceId),
any(ProcessInstanceReadMode.class)))
- .thenReturn(Optional.of(processInstance));
- when(processInstances.acceptingEventType(invalidSignal,
processInstanceId))
- .thenReturn(Stream.empty());
+ assertThat(processService.signalProcessInstance(process, id, "data",
signalName)).isPresent();
+ assertThat(processService.signalProcessInstance(process, id, null,
adHocName)).isPresent();
+ }
+
+ @Test
+ void signalAccepted_whenInstanceHasMultipleAdHocFragments() {
+ String id = "pi-5";
+ String node1 = "AdHocTask1";
+ String node2 = "AdHocTask2";
+ AdHocFragment f1 = new
AdHocFragment.Builder(node1).withName(node1).withAutoStart(false).build();
+ AdHocFragment f2 = new
AdHocFragment.Builder(node2).withName(node2).withAutoStart(true).build();
- // When: Invalid signal is sent
- Optional<TestModel> result =
processService.signalProcessInstance(process, processInstanceId, signalData,
invalidSignal);
+ givenInstanceWithEvents(id, Collections.emptySet(), List.of(f1, f2));
- // Then: Should return empty Optional
- assertThat(result).isEmpty();
+ assertThat(processService.signalProcessInstance(process, id, null,
node1)).isPresent();
+ assertThat(processService.signalProcessInstance(process, id, null,
node2)).isPresent();
}
@Test
- void testSignalProcessInstance_NonExistentInstance_ReturnsEmpty() {
- // Given: Non-existent process instance
- String processInstanceId = "non-existent";
- String signalName = "HelloMartin";
- Object signalData = "test-data";
+ void
signalRejected_whenEventDescriptionContainsRawExpressionInsteadOfResolvedName()
{
+ // Regression: before the fix, getEventDescriptions() could return the
raw #{signalName}
+ // expression instead of the resolved value. The filter must not match
it.
+ String id = "pi-expr";
+ String signalName = "myVarSignal";
- when(processInstances.findById(eq(processInstanceId),
any(ProcessInstanceReadMode.class)))
- .thenReturn(Optional.empty());
- when(processInstances.acceptingEventType(signalName,
processInstanceId))
- .thenReturn(Stream.empty());
+ givenInstanceWithEvents(id, Set.of(eventDesc("#{signalName}")),
Collections.emptyList());
- // When: Signal is sent to non-existent instance
- Optional<TestModel> result =
processService.signalProcessInstance(process, processInstanceId, signalData,
signalName);
+ Optional<TestModel> result =
processService.signalProcessInstance(process, id, "data", signalName);
- // Then: Should return empty Optional
assertThat(result).isEmpty();
}
+ // --- rejection / empty result tests ---
+
@Test
- void testSignalProcessInstance_MultipleAdHocNodes_AcceptsCorrectOne() {
- // Given: Process instance with multiple ad hoc nodes
- String processInstanceId = "test-signal-instance4";
- String adHocNode1 = "AdHocTask1";
- String adHocNode2 = "AdHocTask2";
- TestModel model = new TestModel();
-
- when(processInstances.findById(eq(processInstanceId),
any(ProcessInstanceReadMode.class)))
- .thenReturn(Optional.of(processInstance));
- when(processInstances.acceptingEventType(adHocNode1,
processInstanceId))
- .thenReturn(Stream.of(processInstance));
- when(processInstances.acceptingEventType(adHocNode2,
processInstanceId))
- .thenReturn(Stream.of(processInstance));
- when(processInstance.checkError()).thenReturn(processInstance);
- when(processInstance.variables()).thenReturn(model);
+ void signalRejected_whenInstanceNotWaitingForSignal() {
+ String id = "pi-6";
- // When: Signal first ad hoc node
- Optional<TestModel> result1 =
processService.signalProcessInstance(process, processInstanceId, null,
adHocNode1);
+ givenInstanceWithEvents(id, Set.of(eventDesc("OtherSignal")),
Collections.emptyList());
- // Then: Should accept
- assertThat(result1).isPresent();
+ Optional<TestModel> result =
processService.signalProcessInstance(process, id, "data", "HelloMartin");
- // When: Signal second ad hoc node
- Optional<TestModel> result2 =
processService.signalProcessInstance(process, processInstanceId, null,
adHocNode2);
+ assertThat(result).isEmpty();
+ }
- // Then: Should accept
- assertThat(result2).isPresent();
+ @Test
+ void signalThrows_whenSignalNameIsNull() {
+ assertThatThrownBy(() -> processService.signalProcessInstance(process,
"pi-1", "data", null))
+ .isInstanceOf(NullPointerException.class)
+ .hasMessageContaining("signalName");
}
@Test
- void testSignalProcessInstance_BothSignalAndAdHoc_AcceptsBoth() {
- // Given: Process instance with both signal event and ad hoc node
- String processInstanceId = "test-signal-instance5";
- String signalName = "HelloMartin";
- String adHocNodeName = "AdHocTask";
- TestModel model = new TestModel();
-
- when(processInstances.findById(eq(processInstanceId),
any(ProcessInstanceReadMode.class)))
- .thenReturn(Optional.of(processInstance));
- when(processInstances.acceptingEventType(signalName,
processInstanceId))
- .thenReturn(Stream.of(processInstance));
- when(processInstances.acceptingEventType(adHocNodeName,
processInstanceId))
- .thenReturn(Stream.of(processInstance));
- when(processInstance.checkError()).thenReturn(processInstance);
- when(processInstance.variables()).thenReturn(model);
+ void signalRejected_whenInstanceDoesNotExist() {
+
when(processInstances.findById("non-existent")).thenReturn(Optional.empty());
- // When: Signal traditional event
- Optional<TestModel> resultSignal =
processService.signalProcessInstance(process, processInstanceId, "data",
signalName);
+ Optional<TestModel> result =
processService.signalProcessInstance(process, "non-existent", "data",
"HelloMartin");
- // Then: Should accept
- assertThat(resultSignal).isPresent();
+ assertThat(result).isEmpty();
+ }
+
+ @Test
+ void signalRejected_whenInstanceHasNoEventsAndNoAdHocFragments() {
+ String id = "pi-7";
- // When: Signal ad hoc node
- Optional<TestModel> resultAdHoc =
processService.signalProcessInstance(process, processInstanceId, null,
adHocNodeName);
+ givenInstanceWithEvents(id, Collections.emptySet(),
Collections.emptyList());
- // Then: Should accept
- assertThat(resultAdHoc).isPresent();
+ Optional<TestModel> result =
processService.signalProcessInstance(process, id, "data", "HelloMartin");
+
+ assertThat(result).isEmpty();
}
- /**
- * Test model for testing
- */
+ // --- test model ---
+
static class TestModel implements MappableToModel<TestModel>, Model {
@Override
public TestModel toModel() {
diff --git
a/kogito-jbpm/jbpm-tests/src/test/java/org/jbpm/bpmn2/IntermediateEventTest.java
b/kogito-jbpm/jbpm-tests/src/test/java/org/jbpm/bpmn2/IntermediateEventTest.java
index 6c87bdfdff8..d5a26665044 100755
---
a/kogito-jbpm/jbpm-tests/src/test/java/org/jbpm/bpmn2/IntermediateEventTest.java
+++
b/kogito-jbpm/jbpm-tests/src/test/java/org/jbpm/bpmn2/IntermediateEventTest.java
@@ -2546,6 +2546,16 @@ public class IntermediateEventTest extends
JbpmBpmn2TestCase {
assertThat(processInstance.status()).isEqualTo(ProcessInstance.STATE_ACTIVE);
+ // events() must expose the resolved value, not the raw #{signalName}
expression
+ Set<EventDescription<?>> eventDescriptions = processInstance.events();
+ assertThat(eventDescriptions)
+ .hasSize(1)
+ .extracting(EventDescription::getEvent).contains(signalVar);
+
+ // a signal with a different name must not advance the instance
+ processInstance.send(SignalFactory.of("wrongSignalName",
"WrongValue"));
+
assertThat(processInstance.status()).isEqualTo(ProcessInstance.STATE_ACTIVE);
+
processInstance.send(SignalFactory.of(signalVar, "SomeValue"));
assertThat(processInstance.status()).isEqualTo(ProcessInstance.STATE_COMPLETED);
diff --git
a/kogito-jbpm/jbpm-tests/src/test/java/org/jbpm/bpmn2/SLAComplianceTest.java
b/kogito-jbpm/jbpm-tests/src/test/java/org/jbpm/bpmn2/SLAComplianceTest.java
index 6d0897533c1..9da13c02b1c 100755
--- a/kogito-jbpm/jbpm-tests/src/test/java/org/jbpm/bpmn2/SLAComplianceTest.java
+++ b/kogito-jbpm/jbpm-tests/src/test/java/org/jbpm/bpmn2/SLAComplianceTest.java
@@ -248,7 +248,6 @@ public class SLAComplianceTest extends JbpmBpmn2TestCase {
assertThat(nodeSlaCompliance.get()).isEqualTo(org.kie.api.runtime.process.ProcessInstance.SLA_MET);
}
-
@Test
public void testSLAonProcessViolatedWithExpression() throws Exception {
AtomicInteger processSlaCompliance = new AtomicInteger(-1);
@@ -294,7 +293,6 @@ public class SLAComplianceTest extends JbpmBpmn2TestCase {
assertThat(processSlaCompliance.get()).isEqualTo(org.kie.api.runtime.process.ProcessInstance.SLA_VIOLATED);
}
-
@Test
public void testSLAonCatchEventViolated() throws Exception {
CountDownLatch latch = new CountDownLatch(1);
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]