This is an automated email from the ASF dual-hosted git repository.
Rikkola 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 426327f6519 refactor: extract PropagationEntry from drools-core to
drools-base (#6890)
426327f6519 is described below
commit 426327f651993429499bcbada5b3776653ddbf4f
Author: Toni Rikkola <[email protected]>
AuthorDate: Fri Aug 14 11:03:18 2026 +0300
refactor: extract PropagationEntry from drools-core to drools-base (#6890)
Move PropagationEntry interface and AbstractPropagationEntry to
org.drools.base.phreak so that drools-base sequencing infrastructure
can depend on them without a circular dependency on drools-core.
Changes:
- PropagationEntry<T extends ValueResolver> replaces the hardcoded
ReteEvaluator parameter; all existing implementations bind T=ReteEvaluator
- Inner classes extracted to top-level files under
org.drools.core.phreak.actions:
Insert, Delete, Update, PartitionedUpdate, PartitionedDelete,
ExecuteQuery,
AbstractPartitionedPropagationEntry, PropagationEntryWithResult
- AbstractPropagationEntry extracted to org.drools.base.phreak.actions
- All call sites updated across drools-core, drools-kiesession, drools-tms,
drools-serialization-protobuf, drools-reliability, drools-ruleunits,
kogito-jbpm
- execute() default is now a pure dispatch (onWorkingMemoryAction hook
moved to propagation list implementations that hold a ReteEvaluator)
- ExecuteQuery constructor drops the InternalRuleBase parameter; fetches
rule base via reteEvaluator.getKnowledgeBase() at execution time
No logic changes - pure extraction and relocation.
https://github.com/apache/incubator-kie-issues/issues/2389
Assisted-by: Bob (IBM AI Assistant)
Co-authored-by: Toni Rikkola <[email protected]>
---
.../org/drools/base/phreak/PropagationEntry.java | 33 +-
.../phreak/actions/AbstractPropagationEntry.java | 59 +++
.../org/drools/core/common/ActivationsManager.java | 4 +-
.../drools/core/common/AgendaGroupQueueImpl.java | 9 +-
.../org/drools/core/common/InternalAgenda.java | 4 +-
.../drools/core/common/InternalWorkingMemory.java | 4 +-
.../java/org/drools/core/common/ReteEvaluator.java | 12 +-
.../drools/core/common/WorkingMemoryAction.java | 4 +-
.../drools/core/impl/ActivationsManagerImpl.java | 10 +-
.../core/impl/WorkingMemoryReteExpireAction.java | 8 +-
.../org/drools/core/phreak/PhreakTimerNode.java | 4 +-
.../org/drools/core/phreak/PropagationEntry.java | 499 ---------------------
.../org/drools/core/phreak/PropagationList.java | 11 +-
.../org/drools/core/phreak/ReactiveObjectUtil.java | 3 +-
.../phreak/SynchronizedBypassPropagationList.java | 6 +-
.../core/phreak/SynchronizedPropagationList.java | 34 +-
.../core/phreak/ThreadUnsafePropagationList.java | 10 +-
.../AbstractPartitionedPropagationEntry.java} | 24 +-
.../org/drools/core/phreak/actions/Delete.java | 64 +++
.../drools/core/phreak/actions/ExecuteQuery.java | 82 ++++
.../org/drools/core/phreak/actions/Insert.java | 150 +++++++
.../core/phreak/actions/PartitionedDelete.java | 57 +++
.../core/phreak/actions/PartitionedUpdate.java | 63 +++
.../actions/PropagationEntryWithResult.java} | 38 +-
.../org/drools/core/phreak/actions/Update.java | 103 +++++
.../org/drools/core/reteoo/AsyncReceiveNode.java | 5 +-
.../CompositePartitionAwareObjectSinkAdapter.java | 7 +-
.../org/drools/core/reteoo/EntryPointNode.java | 17 +-
.../org/drools/core/rule/SlidingTimeWindow.java | 5 +-
.../core/time/EnqueuedSelfRemovalJobContext.java | 4 +-
.../drools/core/time/impl/JDKTimerServiceTest.java | 2 +-
.../kiesession/agenda/CompositeDefaultAgenda.java | 6 +-
.../drools/kiesession/agenda/DefaultAgenda.java | 27 +-
.../StatefulKnowledgeSessionForRHS.java | 6 +-
.../session/StatefulKnowledgeSessionImpl.java | 21 +-
.../drools/kiesession/ReteooWorkingMemoryTest.java | 5 +-
.../reliability/core/ReliablePropagationList.java | 11 +-
.../core/ReliableSessionInitializer.java | 11 +-
.../impl/sessions/RuleUnitExecutorImpl.java | 7 +-
.../protobuf/ProtobufOutputMarshaller.java | 6 +-
.../protobuf/WorkingMemoryReteAssertAction.java | 5 +-
.../mvel/compiler/command/PropagationListTest.java | 5 +-
.../simple/BeliefSystemLogicalCallback.java | 5 +-
.../kogito/codegen/tests/CallActivityTaskIT.java | 2 +-
.../jbpm/process/instance/LightProcessRuntime.java | 4 +-
.../jbpm/process/instance/ProcessRuntimeImpl.java | 4 +-
.../instance/event/DefaultSignalManager.java | 6 +-
47 files changed, 794 insertions(+), 672 deletions(-)
diff --git
a/drools-core/src/main/java/org/drools/core/phreak/PropagationList.java
b/drools-base/src/main/java/org/drools/base/phreak/PropagationEntry.java
similarity index 60%
copy from drools-core/src/main/java/org/drools/core/phreak/PropagationList.java
copy to drools-base/src/main/java/org/drools/base/phreak/PropagationEntry.java
index e09fcaff788..040131a9858 100644
--- a/drools-core/src/main/java/org/drools/core/phreak/PropagationList.java
+++ b/drools-base/src/main/java/org/drools/base/phreak/PropagationEntry.java
@@ -1,4 +1,4 @@
-/*
+/**
* 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
@@ -16,33 +16,30 @@
* specific language governing permissions and limitations
* under the License.
*/
-package org.drools.core.phreak;
+package org.drools.base.phreak;
-import java.util.Iterator;
+import org.drools.base.base.ValueResolver;
-public interface PropagationList {
- void addEntry(PropagationEntry propagationEntry);
+public interface PropagationEntry<T extends ValueResolver> {
- PropagationEntry takeAll();
+ default void execute(T t) {
+ internalExecute(t);
+ }
- void flush();
- void flush( PropagationEntry currentHead );
+ void internalExecute(T t);
- void reset();
+ PropagationEntry<T> getNext();
- boolean isEmpty();
+ void setNext(PropagationEntry<T> next);
- boolean hasEntriesDeferringExpiration();
+ boolean requiresImmediateFlushing();
- Iterator<PropagationEntry> iterator();
+ boolean isCalledFromRHS();
- void waitOnRest();
+ boolean isPartitionSplittable();
- void notifyWaitOnRest();
+ PropagationEntry<T> getSplitForPartition(int partitionNr);
- void onEngineInactive();
+ boolean defersExpiration();
- void dispose();
-
- void setFiringUntilHalt( boolean firingUntilHalt );
}
diff --git
a/drools-base/src/main/java/org/drools/base/phreak/actions/AbstractPropagationEntry.java
b/drools-base/src/main/java/org/drools/base/phreak/actions/AbstractPropagationEntry.java
new file mode 100644
index 00000000000..eadfde1a629
--- /dev/null
+++
b/drools-base/src/main/java/org/drools/base/phreak/actions/AbstractPropagationEntry.java
@@ -0,0 +1,59 @@
+/*
+ * 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.drools.base.phreak.actions;
+
+import org.drools.base.base.ValueResolver;
+import org.drools.base.phreak.PropagationEntry;
+
+public abstract class AbstractPropagationEntry<T extends ValueResolver>
implements PropagationEntry<T> {
+ protected PropagationEntry<T> next;
+
+ public void setNext(PropagationEntry<T> next) {
+ this.next = next;
+ }
+
+ public PropagationEntry<T> getNext() {
+ return next;
+ }
+
+ @Override
+ public boolean requiresImmediateFlushing() {
+ return false;
+ }
+
+ @Override
+ public boolean isCalledFromRHS() {
+ return false;
+ }
+
+ @Override
+ public boolean isPartitionSplittable() {
+ return false;
+ }
+
+ @Override
+ public boolean defersExpiration() {
+ return false;
+ }
+
+ @Override
+ public PropagationEntry<T> getSplitForPartition(int partitionNr) {
+ throw new UnsupportedOperationException();
+ }
+}
diff --git
a/drools-core/src/main/java/org/drools/core/common/ActivationsManager.java
b/drools-core/src/main/java/org/drools/core/common/ActivationsManager.java
index 1e810d1c7ac..714818f0488 100644
--- a/drools-core/src/main/java/org/drools/core/common/ActivationsManager.java
+++ b/drools-core/src/main/java/org/drools/core/common/ActivationsManager.java
@@ -21,7 +21,7 @@ package org.drools.core.common;
import org.drools.base.common.NetworkNode;
import org.drools.core.event.AgendaEventSupport;
import org.drools.core.phreak.ExecutableEntry;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
import org.drools.core.phreak.RuleAgendaItem;
import org.drools.core.reteoo.PathMemory;
import org.drools.core.reteoo.RuleTerminalNodeLeftTuple;
@@ -87,7 +87,7 @@ public interface ActivationsManager {
int fireAllRules(AgendaFilter agendaFilter, int fireLimit);
- void addPropagation(PropagationEntry propagationEntry);
+ void addPropagation(PropagationEntry<ReteEvaluator> propagationEntry);
default void stageLeftTuple(RuleAgendaItem ruleAgendaItem, InternalMatch
justified) {
if (!ruleAgendaItem.isQueued()) {
diff --git
a/drools-core/src/main/java/org/drools/core/common/AgendaGroupQueueImpl.java
b/drools-core/src/main/java/org/drools/core/common/AgendaGroupQueueImpl.java
index 43fa4de13d3..7f57cc14663 100644
--- a/drools-core/src/main/java/org/drools/core/common/AgendaGroupQueueImpl.java
+++ b/drools-core/src/main/java/org/drools/core/common/AgendaGroupQueueImpl.java
@@ -27,7 +27,8 @@ import java.util.concurrent.ConcurrentHashMap;
import org.drools.core.conflict.RuleAgendaConflictResolver;
import org.drools.core.impl.InternalRuleBase;
import org.drools.core.marshalling.MarshallerReaderContext;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
+import org.drools.base.phreak.actions.AbstractPropagationEntry;
import org.drools.core.phreak.RuleAgendaItem;
import org.drools.core.util.ArrayQueue;
import org.drools.core.util.Queue;
@@ -112,7 +113,7 @@ public class AgendaGroupQueueImpl
reteEvaluator.addPropagation( new ClearAction( this.name ) );
}
- public class ClearAction extends PropagationEntry.AbstractPropagationEntry
{
+ public class ClearAction extends AbstractPropagationEntry<ReteEvaluator> {
private final String name;
@@ -130,7 +131,7 @@ public class AgendaGroupQueueImpl
reteEvaluator.addPropagation( new SetFocusAction( this.name ) );
}
- public class SetFocusAction extends
PropagationEntry.AbstractPropagationEntry {
+ public class SetFocusAction extends
AbstractPropagationEntry<ReteEvaluator> {
private final String name;
@@ -267,7 +268,7 @@ public class AgendaGroupQueueImpl
}
public static class DeactivateCallback
- extends PropagationEntry.AbstractPropagationEntry
+ extends AbstractPropagationEntry<ReteEvaluator>
implements WorkingMemoryAction {
private static final long serialVersionUID = 510l;
diff --git
a/drools-core/src/main/java/org/drools/core/common/InternalAgenda.java
b/drools-core/src/main/java/org/drools/core/common/InternalAgenda.java
index f49662ccf9c..a532189c512 100644
--- a/drools-core/src/main/java/org/drools/core/common/InternalAgenda.java
+++ b/drools-core/src/main/java/org/drools/core/common/InternalAgenda.java
@@ -21,7 +21,7 @@ package org.drools.core.common;
import java.util.Iterator;
import java.util.Map;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
import org.kie.api.runtime.rule.Agenda;
import org.kie.api.runtime.rule.AgendaFilter;
@@ -147,7 +147,7 @@ public interface InternalAgenda extends Agenda,
ActivationsManager {
boolean isRuleActiveInRuleFlowGroup(String ruleflowGroupName, String
ruleName, String processInstanceId);
void notifyWaitOnRest();
- Iterator<PropagationEntry> getActionsIterator();
+ Iterator<PropagationEntry<ReteEvaluator>> getActionsIterator();
boolean hasPendingPropagations();
boolean isParallelAgenda();
diff --git
a/drools-core/src/main/java/org/drools/core/common/InternalWorkingMemory.java
b/drools-core/src/main/java/org/drools/core/common/InternalWorkingMemory.java
index 77f70726ff6..7b0f6c41747 100644
---
a/drools-core/src/main/java/org/drools/core/common/InternalWorkingMemory.java
+++
b/drools-core/src/main/java/org/drools/core/common/InternalWorkingMemory.java
@@ -27,7 +27,7 @@ import org.drools.core.WorkingMemory;
import org.drools.core.WorkingMemoryEntryPoint;
import org.drools.core.event.AgendaEventSupport;
import org.drools.core.event.RuleRuntimeEventSupport;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
import org.drools.core.rule.consequence.InternalMatch;
import org.drools.core.runtime.process.InternalProcessRuntime;
import org.kie.api.runtime.Channel;
@@ -110,7 +110,7 @@ public interface InternalWorkingMemory
void deactivate();
boolean tryDeactivate();
- Iterator<? extends PropagationEntry> getActionsIterator();
+ Iterator<? extends PropagationEntry<ReteEvaluator>> getActionsIterator();
void removeGlobal(String identifier);
diff --git
a/drools-core/src/main/java/org/drools/core/common/ReteEvaluator.java
b/drools-core/src/main/java/org/drools/core/common/ReteEvaluator.java
index 30a5803ae97..9ed93e4284a 100644
--- a/drools-core/src/main/java/org/drools/core/common/ReteEvaluator.java
+++ b/drools-core/src/main/java/org/drools/core/common/ReteEvaluator.java
@@ -31,7 +31,7 @@ import org.drools.core.event.AgendaEventSupport;
import org.drools.core.event.RuleEventListenerSupport;
import org.drools.core.event.RuleRuntimeEventSupport;
import org.drools.core.impl.InternalRuleBase;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
import org.drools.core.phreak.RuleNetworkEvaluator;
import org.drools.core.reteoo.ObjectTypeConf;
import org.drools.core.reteoo.RuntimeComponentFactory;
@@ -97,7 +97,7 @@ public interface ReteEvaluator extends ValueResolver {
return timerService != null ? timerService.getTimerJobInstances(id) :
Collections.emptyList();
}
- void addPropagation(PropagationEntry propagationEntry);
+ void addPropagation(PropagationEntry<ReteEvaluator> propagationEntry);
long getNextPropagationIdCounter();
@@ -142,16 +142,16 @@ public interface ReteEvaluator extends ValueResolver {
int fireAllRules(AgendaFilter agendaFilter);
int fireAllRules(AgendaFilter agendaFilter, int max);
- default void setWorkingMemoryActionListener(Consumer<PropagationEntry>
listener) {
+ default void
setWorkingMemoryActionListener(Consumer<PropagationEntry<ReteEvaluator>>
listener) {
throw new UnsupportedOperationException();
}
- default Consumer<PropagationEntry> getWorkingMemoryActionListener() {
+ default Consumer<PropagationEntry<ReteEvaluator>>
getWorkingMemoryActionListener() {
return null;
}
- default void onWorkingMemoryAction(PropagationEntry entry) {
- Consumer<PropagationEntry> listener = getWorkingMemoryActionListener();
+ default void onWorkingMemoryAction(PropagationEntry<ReteEvaluator> entry) {
+ Consumer<PropagationEntry<ReteEvaluator>> listener =
getWorkingMemoryActionListener();
if (listener != null) {
listener.accept(entry);
}
diff --git
a/drools-core/src/main/java/org/drools/core/common/WorkingMemoryAction.java
b/drools-core/src/main/java/org/drools/core/common/WorkingMemoryAction.java
index 63497ff8e2e..341b34da384 100644
--- a/drools-core/src/main/java/org/drools/core/common/WorkingMemoryAction.java
+++ b/drools-core/src/main/java/org/drools/core/common/WorkingMemoryAction.java
@@ -18,9 +18,9 @@
*/
package org.drools.core.common;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
-public interface WorkingMemoryAction extends PropagationEntry {
+public interface WorkingMemoryAction extends PropagationEntry<ReteEvaluator> {
short WorkingMemoryReteAssertAction = 1;
short DeactivateCallback = 2;
short PropagateAction = 3;
diff --git
a/drools-core/src/main/java/org/drools/core/impl/ActivationsManagerImpl.java
b/drools-core/src/main/java/org/drools/core/impl/ActivationsManagerImpl.java
index 05250612393..97fe44cf0c0 100644
--- a/drools-core/src/main/java/org/drools/core/impl/ActivationsManagerImpl.java
+++ b/drools-core/src/main/java/org/drools/core/impl/ActivationsManagerImpl.java
@@ -42,7 +42,7 @@ import org.drools.core.concurrent.GroupEvaluator;
import org.drools.core.concurrent.SequentialGroupEvaluator;
import org.drools.core.event.AgendaEventSupport;
import org.drools.core.phreak.ExecutableEntry;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
import org.drools.core.phreak.PropagationList;
import org.drools.core.phreak.RuleAgendaItem;
import org.drools.core.phreak.RuleExecutor;
@@ -274,7 +274,7 @@ public class ActivationsManagerImpl implements
ActivationsManager {
}
@Override
- public void addPropagation(PropagationEntry propagationEntry) {
+ public void addPropagation(PropagationEntry<ReteEvaluator>
propagationEntry) {
propagationList.addEntry( propagationEntry );
}
@@ -286,7 +286,7 @@ public class ActivationsManagerImpl implements
ActivationsManager {
private int fireLoop(AgendaFilter agendaFilter, int fireLimit) {
firing = true;
int fireCount = 0;
- PropagationEntry head = propagationList.takeAll();
+ PropagationEntry<ReteEvaluator> head = propagationList.takeAll();
int returnedFireCount;
boolean limitReached = fireLimit == 0; // -1 or > 0 will return false.
No reason for user to give 0, just handled for completeness.
@@ -345,8 +345,8 @@ public class ActivationsManagerImpl implements
ActivationsManager {
}
}
- private PropagationEntry handleRest() {
- PropagationEntry head = propagationList.takeAll();
+ private PropagationEntry<ReteEvaluator> handleRest() {
+ PropagationEntry<ReteEvaluator> head = propagationList.takeAll();
if (head == null) {
firing = false;
}
diff --git
a/drools-core/src/main/java/org/drools/core/impl/WorkingMemoryReteExpireAction.java
b/drools-core/src/main/java/org/drools/core/impl/WorkingMemoryReteExpireAction.java
index 4021af67721..5841d09dfdc 100644
---
a/drools-core/src/main/java/org/drools/core/impl/WorkingMemoryReteExpireAction.java
+++
b/drools-core/src/main/java/org/drools/core/impl/WorkingMemoryReteExpireAction.java
@@ -29,14 +29,16 @@ import org.drools.core.common.PropagationContext;
import org.drools.core.common.ReteEvaluator;
import org.drools.core.common.WorkingMemoryAction;
import org.drools.core.marshalling.MarshallerReaderContext;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
+import org.drools.base.phreak.actions.AbstractPropagationEntry;
+import org.drools.core.phreak.actions.AbstractPartitionedPropagationEntry;
import org.drools.core.reteoo.ObjectTypeNode;
import org.drools.core.reteoo.RightTuple;
import static
org.drools.core.common.PhreakPropagationContextFactory.createPropagationContextForFact;
public class WorkingMemoryReteExpireAction
- extends PropagationEntry.AbstractPropagationEntry
+ extends AbstractPropagationEntry<ReteEvaluator>
implements WorkingMemoryAction, Externalizable {
protected DefaultEventHandle factHandle;
@@ -127,7 +129,7 @@ public class WorkingMemoryReteExpireAction
this.factHandle = (DefaultEventHandle) in.readObject();
}
- public static class PartitionAwareWorkingMemoryReteExpireAction extends
PropagationEntry.AbstractPartitionedPropagationEntry {
+ public static class PartitionAwareWorkingMemoryReteExpireAction extends
AbstractPartitionedPropagationEntry<ReteEvaluator> {
private final DefaultEventHandle factHandle;
private final ObjectTypeNode node;
diff --git
a/drools-core/src/main/java/org/drools/core/phreak/PhreakTimerNode.java
b/drools-core/src/main/java/org/drools/core/phreak/PhreakTimerNode.java
index 5b7c9b916ca..481c136a1d5 100644
--- a/drools-core/src/main/java/org/drools/core/phreak/PhreakTimerNode.java
+++ b/drools-core/src/main/java/org/drools/core/phreak/PhreakTimerNode.java
@@ -41,6 +41,8 @@ import org.drools.core.reteoo.TupleFactory;
import org.drools.core.reteoo.TupleImpl;
import org.drools.core.time.Job;
import org.drools.core.time.JobContext;
+import org.drools.base.phreak.PropagationEntry;
+import org.drools.base.phreak.actions.AbstractPropagationEntry;
import org.drools.core.time.TimerService;
import org.drools.core.time.impl.DefaultJobHandle;
import org.drools.core.util.index.TupleList;
@@ -367,7 +369,7 @@ public class PhreakTimerNode {
}
public static class TimerAction
- extends PropagationEntry.AbstractPropagationEntry
+ extends AbstractPropagationEntry<ReteEvaluator>
implements WorkingMemoryAction {
private final TimerNodeJobContext timerJobCtx;
diff --git
a/drools-core/src/main/java/org/drools/core/phreak/PropagationEntry.java
b/drools-core/src/main/java/org/drools/core/phreak/PropagationEntry.java
deleted file mode 100644
index da5dba2050f..00000000000
--- a/drools-core/src/main/java/org/drools/core/phreak/PropagationEntry.java
+++ /dev/null
@@ -1,499 +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.drools.core.phreak;
-
-import java.io.Externalizable;
-import java.io.IOException;
-import java.io.ObjectInput;
-import java.io.ObjectOutput;
-import java.util.concurrent.CountDownLatch;
-
-import org.drools.base.reteoo.NodeTypeEnums;
-import org.drools.core.base.DroolsQueryImpl;
-import org.drools.core.common.DefaultEventHandle;
-import org.drools.core.common.InternalFactHandle;
-import org.drools.core.common.PropagationContext;
-import org.drools.core.common.ReteEvaluator;
-import org.drools.core.impl.InternalRuleBase;
-import org.drools.core.impl.WorkingMemoryReteExpireAction;
-import org.drools.core.reteoo.ClassObjectTypeConf;
-import org.drools.core.reteoo.CompositePartitionAwareObjectSinkAdapter;
-import org.drools.core.reteoo.EntryPointNode;
-import org.drools.core.reteoo.LeftInputAdapterNode;
-import org.drools.core.reteoo.LeftTupleSource;
-import org.drools.core.reteoo.ModifyPreviousTuples;
-import org.drools.core.reteoo.ObjectTypeConf;
-import org.drools.core.reteoo.ObjectTypeNode;
-import org.drools.core.reteoo.PathMemory;
-import org.drools.core.reteoo.QueryTerminalNode;
-import org.drools.core.reteoo.TerminalNode;
-import org.drools.core.time.JobContext;
-import org.drools.core.time.impl.DefaultJobHandle;
-import org.drools.core.time.impl.PointInTimeTrigger;
-import org.kie.api.prototype.PrototypeEventInstance;
-
-import static org.drools.base.rule.TypeDeclaration.NEVER_EXPIRES;
-import static
org.drools.core.reteoo.EntryPointNode.removeRightTuplesMatchingOTN;
-
-public interface PropagationEntry {
-
- default void execute(ReteEvaluator reteEvaluator) {
- internalExecute(reteEvaluator);
- reteEvaluator.onWorkingMemoryAction(this);
- }
-
- void internalExecute(ReteEvaluator reteEvaluator);
-
- PropagationEntry getNext();
- void setNext(PropagationEntry next);
-
- boolean requiresImmediateFlushing();
-
- boolean isCalledFromRHS();
-
- boolean isPartitionSplittable();
- PropagationEntry getSplitForPartition(int partitionNr);
-
- boolean defersExpiration();
-
- abstract class AbstractPropagationEntry implements PropagationEntry {
- protected PropagationEntry next;
-
- public void setNext(PropagationEntry next) {
- this.next = next;
- }
-
- public PropagationEntry getNext() {
- return next;
- }
-
- @Override
- public boolean requiresImmediateFlushing() {
- return false;
- }
-
- @Override
- public boolean isCalledFromRHS() {
- return false;
- }
-
- @Override
- public boolean isPartitionSplittable() {
- return false;
- }
-
- @Override
- public boolean defersExpiration() {
- return false;
- }
-
- @Override
- public PropagationEntry getSplitForPartition(int partitionNr) {
- throw new UnsupportedOperationException();
- }
-
- public InternalFactHandle getHandle() {
- throw new UnsupportedOperationException( "This method is not
supported for this type of PropagationEntry" );
- }
- }
-
- abstract class AbstractPartitionedPropagationEntry extends
AbstractPropagationEntry {
- protected final int partition;
-
- protected AbstractPartitionedPropagationEntry( int partition ) {
- this.partition = partition;
- }
-
- protected boolean isMainPartition() {
- return partition == 0;
- }
- }
-
- abstract class PropagationEntryWithResult<T> extends
PropagationEntry.AbstractPropagationEntry {
- private final CountDownLatch done = new CountDownLatch( 1 );
-
- private T result;
-
- public final T getResult() {
- try {
- done.await();
- } catch (InterruptedException e) {
- throw new RuntimeException( e );
- }
- return result;
- }
-
- protected void done(T result) {
- this.result = result;
- done.countDown();
- }
-
- @Override
- public boolean requiresImmediateFlushing() {
- return true;
- }
- }
-
- class ExecuteQuery extends
PropagationEntry.PropagationEntryWithResult<QueryTerminalNode[]> {
-
- private final String queryName;
- private final DroolsQueryImpl queryObject;
- private final InternalFactHandle handle;
- private final PropagationContext pCtx;
- private final boolean calledFromRHS;
- private InternalRuleBase ruleBase;
-
- public ExecuteQuery(InternalRuleBase ruleBase, String queryName,
DroolsQueryImpl queryObject, InternalFactHandle handle, PropagationContext
pCtx, boolean calledFromRHS) {
- this.ruleBase = ruleBase;
- this.queryName = queryName;
- this.queryObject = queryObject;
- this.handle = handle;
- this.pCtx = pCtx;
- this.calledFromRHS = calledFromRHS;
- }
-
- @Override
- public void internalExecute(ReteEvaluator reteEvaluator ) {
- QueryTerminalNode[] tnodes =
ruleBase.getReteooBuilder().getTerminalNodesForQuery( queryName );
- if ( tnodes == null ) {
- throw new RuntimeException( "Query '" + queryName + "' does
not exist" );
- }
-
- QueryTerminalNode tnode = tnodes[0];
-
- if (queryObject.getElements().length !=
tnode.getQuery().getParameters().length) {
- throw new RuntimeException( "Query '" + queryName + "' has
been invoked with a wrong number of arguments. Expected " +
- tnode.getQuery().getParameters().length + ", actual "
+ queryObject.getElements().length );
- }
-
- LeftTupleSource lts = tnode.getLeftTupleSource();
- while ( !NodeTypeEnums.isLeftInputAdapterNode(lts)) {
- lts = lts.getLeftTupleSource();
- }
- LeftInputAdapterNode lian = (LeftInputAdapterNode) lts;
- LeftInputAdapterNode.LiaNodeMemory lmem =
reteEvaluator.getNodeMemory( lian );
- if ( lmem.getSegmentMemory() == null ) {
-
reteEvaluator.getSegmentMemorySupport().getOrCreateSegmentMemory(lts, lmem);
- }
-
- LeftInputAdapterNode.doInsertObject( handle, pCtx, lian,
reteEvaluator, lmem, false, queryObject.isOpen() );
-
- for ( PathMemory rm : lmem.getSegmentMemory().getPathMemories() ) {
- RuleAgendaItem evaluator =
reteEvaluator.getActivationsManager().createRuleAgendaItem( Integer.MAX_VALUE,
rm, (TerminalNode) rm.getPathEndNode() );
- evaluator.getRuleExecutor().setDirty( true );
- evaluator.getRuleExecutor().evaluateNetworkAndFire(
reteEvaluator, null, 0, -1 );
- }
-
- done(tnodes);
- }
-
- @Override
- public boolean isCalledFromRHS() {
- return calledFromRHS;
- }
- }
-
- class Insert extends AbstractPropagationEntry implements Externalizable {
- private static final ObjectTypeNode.ExpireJob job = new
ObjectTypeNode.ExpireJob();
-
- private InternalFactHandle handle;
- private PropagationContext context;
- private ObjectTypeConf objectTypeConf;
-
- public Insert() { }
-
- public Insert( InternalFactHandle handle, PropagationContext context,
ReteEvaluator reteEvaluator, ObjectTypeConf objectTypeConf) {
- this.handle = handle;
- this.context = context;
- this.objectTypeConf = objectTypeConf;
-
- if ( handle.isEvent() ) {
- scheduleExpiration(reteEvaluator, handle, context,
objectTypeConf, reteEvaluator.getTimerService().getCurrentTime());
- }
- }
-
- public static void execute( InternalFactHandle handle,
PropagationContext context, ReteEvaluator reteEvaluator, ObjectTypeConf
objectTypeConf) {
- if ( handle.isEvent() ) {
- scheduleExpiration(reteEvaluator, handle, context,
objectTypeConf, reteEvaluator.getTimerService().getCurrentTime());
- }
- propagate( handle, context, reteEvaluator, objectTypeConf );
- }
-
- private static void propagate( InternalFactHandle handle,
PropagationContext context, ReteEvaluator reteEvaluator, ObjectTypeConf
objectTypeConf ) {
- if (objectTypeConf == null) {
- // it can be null after deserialization
- objectTypeConf =
handle.getEntryPoint(reteEvaluator).getObjectTypeConfigurationRegistry().getOrCreateObjectTypeConf(handle.getEntryPointId(),
handle.getObject());
- }
- for ( ObjectTypeNode otn : objectTypeConf.getObjectTypeNodes() ) {
- otn.propagateAssert( handle, context, reteEvaluator );
- }
- if ( isOrphanHandle(handle, reteEvaluator) ) {
- handle.setDisconnected(true);
-
handle.getEntryPoint(reteEvaluator).getObjectStore().removeHandle( handle );
- if (handle instanceof DefaultEventHandle eventHandle) {
- eventHandle.unscheduleAllJobs(reteEvaluator);
- }
- }
- }
-
- private static boolean isOrphanHandle(InternalFactHandle handle,
ReteEvaluator reteEvaluator) {
- return !handle.hasMatches() &&
!reteEvaluator.getKnowledgeBase().getKieBaseConfiguration().isMutabilityEnabled();
- }
-
- public void internalExecute(ReteEvaluator reteEvaluator ) {
- propagate( handle, context, reteEvaluator, objectTypeConf );
- }
-
- private static void scheduleExpiration(ReteEvaluator reteEvaluator,
InternalFactHandle handle, PropagationContext context, ObjectTypeConf
objectTypeConf, long insertionTime) {
- for ( ObjectTypeNode otn : objectTypeConf.getObjectTypeNodes() ) {
- long expirationOffset = objectTypeConf.isPrototype() ?
((PrototypeEventInstance) handle.getObject()).getExpiration() :
otn.getExpirationOffset();
- scheduleExpiration( reteEvaluator, handle, context, otn,
insertionTime, expirationOffset );
- }
- if ( objectTypeConf.getConcreteObjectTypeNode() == null ) {
- long expirationOffset = objectTypeConf.isPrototype() ?
((PrototypeEventInstance) handle.getObject()).getExpiration() :
((ClassObjectTypeConf) objectTypeConf).getExpirationOffset();
- scheduleExpiration( reteEvaluator, handle, context, null,
insertionTime, expirationOffset);
- }
- }
-
- private static void scheduleExpiration( ReteEvaluator reteEvaluator,
InternalFactHandle handle, PropagationContext context, ObjectTypeNode otn, long
insertionTime, long expirationOffset ) {
- if ( expirationOffset == NEVER_EXPIRES || expirationOffset ==
Long.MAX_VALUE || context.getReaderContext() != null ) {
- return;
- }
-
- // DROOLS-455 the calculation of the effectiveEnd may overflow and
become negative
- DefaultEventHandle eventFactHandle = (DefaultEventHandle) handle;
- long nextTimestamp = getNextTimestamp( insertionTime,
expirationOffset, eventFactHandle );
-
- WorkingMemoryReteExpireAction action = new
WorkingMemoryReteExpireAction((DefaultEventHandle) handle, otn );
- if (nextTimestamp <=
reteEvaluator.getTimerService().getCurrentTime()) {
- reteEvaluator.addPropagation( action );
- } else {
- JobContext jobctx = new ObjectTypeNode.ExpireJobContext(
action, reteEvaluator );
- DefaultJobHandle jobHandle = (DefaultJobHandle)
reteEvaluator.getTimerService()
-
.scheduleJob( job, jobctx, PointInTimeTrigger.createPointInTimeTrigger(
nextTimestamp, null ) );
- jobctx.setJobHandle( jobHandle );
- eventFactHandle.addJob( jobHandle );
- }
- }
-
- private static long getNextTimestamp( long insertionTime, long
expirationOffset, DefaultEventHandle eventFactHandle) {
- long effectiveEnd = eventFactHandle.getEndTimestamp() +
expirationOffset;
- return Math.max( insertionTime, effectiveEnd >= 0 ? effectiveEnd :
Long.MAX_VALUE );
- }
-
- @Override
- public String toString() {
- return "Insert of " + handle.getObject();
- }
-
- @Override
- public InternalFactHandle getHandle() {
- return handle;
- }
-
- @Override
- public void writeExternal(ObjectOutput out) throws IOException {
- out.writeObject(next);
- out.writeObject(handle);
- out.writeObject(context);
- }
-
- @Override
- public void readExternal(ObjectInput in) throws IOException,
ClassNotFoundException {
- this.next = (PropagationEntry) in.readObject();
- this.handle = (InternalFactHandle) in.readObject();
- this.context = (PropagationContext) in.readObject();
- }
- }
-
- class Update extends AbstractPropagationEntry implements Externalizable {
- private InternalFactHandle handle;
- private PropagationContext context;
- private ObjectTypeConf objectTypeConf;
-
- public Update(){}
-
- public Update(InternalFactHandle handle, PropagationContext context,
ObjectTypeConf objectTypeConf) {
- this.handle = handle;
- this.context = context;
- this.objectTypeConf = objectTypeConf;
- }
-
- public void internalExecute(ReteEvaluator reteEvaluator) {
- execute(handle, context, objectTypeConf, reteEvaluator);
- }
-
- public static void execute(InternalFactHandle handle,
PropagationContext pctx, ObjectTypeConf objectTypeConf, ReteEvaluator
reteEvaluator) {
- if (objectTypeConf == null) {
- // it can be null after deserialization
- objectTypeConf =
handle.getEntryPoint(reteEvaluator).getObjectTypeConfigurationRegistry().getOrCreateObjectTypeConf(handle.getEntryPointId(),
handle.getObject());
- }
- // make a reference to the previous tuples, then null then on the
handle
- ModifyPreviousTuples modifyPreviousTuples = new
ModifyPreviousTuples( handle.detachLinkedTuples() );
- ObjectTypeNode[] cachedNodes = objectTypeConf.getObjectTypeNodes();
- for ( int i = 0, length = cachedNodes.length; i < length; i++ ) {
- cachedNodes[i].modifyObject( handle, modifyPreviousTuples,
pctx, reteEvaluator );
- if (i < cachedNodes.length - 1) {
- removeRightTuplesMatchingOTN( pctx, reteEvaluator,
modifyPreviousTuples, cachedNodes[i], 0 );
- }
- }
- modifyPreviousTuples.retractTuples(pctx, reteEvaluator);
- }
-
- @Override
- public boolean isPartitionSplittable() {
- return true;
- }
-
- @Override
- public PropagationEntry getSplitForPartition( int partitionNr ) {
- return new PartitionedUpdate( handle, context, objectTypeConf,
partitionNr );
- }
-
- @Override
- public InternalFactHandle getHandle() {
- return handle;
- }
-
- @Override
- public String toString() {
- return "Update of " + handle.getObject();
- }
-
- @Override
- public void writeExternal(ObjectOutput out) throws IOException {
- out.writeObject(next);
- out.writeObject(handle);
- out.writeObject(context);
- }
-
- @Override
- public void readExternal(ObjectInput in) throws IOException,
ClassNotFoundException {
- this.next = (PropagationEntry) in.readObject();
- this.handle = (InternalFactHandle) in.readObject();
- this.context = (PropagationContext) in.readObject();
- }
- }
-
- class PartitionedUpdate extends AbstractPartitionedPropagationEntry {
- private final InternalFactHandle handle;
- private final PropagationContext context;
- private final ObjectTypeConf objectTypeConf;
-
- PartitionedUpdate(InternalFactHandle handle, PropagationContext
context, ObjectTypeConf objectTypeConf, int partition) {
- super( partition );
- this.handle = handle;
- this.context = context;
- this.objectTypeConf = objectTypeConf;
- }
-
- public void internalExecute(ReteEvaluator reteEvaluator) {
- ModifyPreviousTuples modifyPreviousTuples = new
ModifyPreviousTuples( handle.detachLinkedTuplesForPartition(partition) );
- ObjectTypeNode[] cachedNodes = objectTypeConf.getObjectTypeNodes();
- for ( int i = 0, length = cachedNodes.length; i < length; i++ ) {
- ObjectTypeNode otn = cachedNodes[i];
- ( (CompositePartitionAwareObjectSinkAdapter)
otn.getObjectSinkPropagator() )
- .propagateModifyObjectForPartition( handle,
modifyPreviousTuples,
-
context.adaptModificationMaskForObjectType(otn.getObjectType(), reteEvaluator),
- reteEvaluator,
partition );
- if (i < cachedNodes.length - 1) {
- removeRightTuplesMatchingOTN( context, reteEvaluator,
modifyPreviousTuples, otn, partition );
- }
- }
- modifyPreviousTuples.retractTuples(context, reteEvaluator);
- }
-
- @Override
- public String toString() {
- return "Update of " + handle.getObject() + " for partition " +
partition;
- }
- }
-
- class Delete extends AbstractPropagationEntry {
- private final EntryPointNode epn;
- private final InternalFactHandle handle;
- private final PropagationContext context;
- private final ObjectTypeConf objectTypeConf;
-
- public Delete(EntryPointNode epn, InternalFactHandle handle,
PropagationContext context, ObjectTypeConf objectTypeConf) {
- this.epn = epn;
- this.handle = handle;
- this.context = context;
- this.objectTypeConf = objectTypeConf;
- }
-
- public void internalExecute(ReteEvaluator reteEvaluator) {
- execute(reteEvaluator, epn, handle, context, objectTypeConf);
- }
-
- public static void execute(ReteEvaluator reteEvaluator, EntryPointNode
epn, InternalFactHandle handle, PropagationContext context, ObjectTypeConf
objectTypeConf) {
- epn.propagateRetract(handle, context, objectTypeConf,
reteEvaluator);
- }
-
- @Override
- public boolean isPartitionSplittable() {
- return true;
- }
-
- @Override
- public PropagationEntry getSplitForPartition( int partitionNr ) {
- return new PartitionedDelete( handle, context, objectTypeConf,
partitionNr );
- }
-
- @Override
- public String toString() {
- return "Delete of " + handle.getObject();
- }
- }
-
- class PartitionedDelete extends AbstractPartitionedPropagationEntry {
- private final InternalFactHandle handle;
- private final PropagationContext context;
- private final ObjectTypeConf objectTypeConf;
-
- PartitionedDelete(InternalFactHandle handle, PropagationContext
context, ObjectTypeConf objectTypeConf, int partition) {
- super( partition );
- this.handle = handle;
- this.context = context;
- this.objectTypeConf = objectTypeConf;
- }
-
- public void internalExecute(ReteEvaluator reteEvaluator) {
- ObjectTypeNode[] cachedNodes = objectTypeConf.getObjectTypeNodes();
-
- if ( cachedNodes == null ) {
- // it is possible that there are no ObjectTypeNodes for an
object being retracted
- return;
- }
-
- for ( ObjectTypeNode cachedNode : cachedNodes ) {
- cachedNode.retractObject( handle, context, reteEvaluator,
partition );
- }
-
- if (handle.isEvent() && isMainPartition()) {
- ((DefaultEventHandle) handle).unscheduleAllJobs(reteEvaluator);
- }
- }
-
- @Override
- public String toString() {
- return "Delete of " + handle.getObject() + " for partition " +
partition;
- }
- }
-}
diff --git
a/drools-core/src/main/java/org/drools/core/phreak/PropagationList.java
b/drools-core/src/main/java/org/drools/core/phreak/PropagationList.java
index e09fcaff788..b93c7c91335 100644
--- a/drools-core/src/main/java/org/drools/core/phreak/PropagationList.java
+++ b/drools-core/src/main/java/org/drools/core/phreak/PropagationList.java
@@ -20,13 +20,16 @@ package org.drools.core.phreak;
import java.util.Iterator;
+import org.drools.base.phreak.PropagationEntry;
+import org.drools.core.common.ReteEvaluator;
+
public interface PropagationList {
- void addEntry(PropagationEntry propagationEntry);
+ void addEntry(PropagationEntry<ReteEvaluator> propagationEntry);
- PropagationEntry takeAll();
+ PropagationEntry<ReteEvaluator> takeAll();
void flush();
- void flush( PropagationEntry currentHead );
+ void flush( PropagationEntry<ReteEvaluator> currentHead );
void reset();
@@ -34,7 +37,7 @@ public interface PropagationList {
boolean hasEntriesDeferringExpiration();
- Iterator<PropagationEntry> iterator();
+ Iterator<PropagationEntry<ReteEvaluator>> iterator();
void waitOnRest();
diff --git
a/drools-core/src/main/java/org/drools/core/phreak/ReactiveObjectUtil.java
b/drools-core/src/main/java/org/drools/core/phreak/ReactiveObjectUtil.java
index 38ab7938185..a8ab02946d6 100644
--- a/drools-core/src/main/java/org/drools/core/phreak/ReactiveObjectUtil.java
+++ b/drools-core/src/main/java/org/drools/core/phreak/ReactiveObjectUtil.java
@@ -21,6 +21,7 @@ package org.drools.core.phreak;
import java.util.Collection;
import org.drools.base.phreak.ReactiveObject;
+import org.drools.base.phreak.actions.AbstractPropagationEntry;
import org.drools.base.reteoo.BaseTuple;
import org.drools.core.common.BetaConstraints;
import org.drools.core.common.InternalFactHandle;
@@ -64,7 +65,7 @@ public class ReactiveObjectUtil {
}
}
- static class ReactivePropagation extends
PropagationEntry.AbstractPropagationEntry {
+ static class ReactivePropagation extends
AbstractPropagationEntry<ReteEvaluator> {
private final Object object;
private final ReactiveFromNodeLeftTuple leftTuple;
diff --git
a/drools-core/src/main/java/org/drools/core/phreak/SynchronizedBypassPropagationList.java
b/drools-core/src/main/java/org/drools/core/phreak/SynchronizedBypassPropagationList.java
index 7ee672ac556..50fe6551749 100644
---
a/drools-core/src/main/java/org/drools/core/phreak/SynchronizedBypassPropagationList.java
+++
b/drools-core/src/main/java/org/drools/core/phreak/SynchronizedBypassPropagationList.java
@@ -20,6 +20,7 @@ package org.drools.core.phreak;
import java.util.concurrent.atomic.AtomicBoolean;
+import org.drools.base.phreak.PropagationEntry;
import org.drools.core.common.ReteEvaluator;
public class SynchronizedBypassPropagationList extends
SynchronizedPropagationList {
@@ -31,13 +32,14 @@ public class SynchronizedBypassPropagationList extends
SynchronizedPropagationLi
}
@Override
- public void addEntry(final PropagationEntry propagationEntry) {
+ public void addEntry(final PropagationEntry<ReteEvaluator>
propagationEntry) {
reteEvaluator.getActivationsManager().executeTask( new
ExecutableEntry() {
@Override
public void execute() {
if (executing.compareAndSet( false, true )) {
try {
propagationEntry.execute( reteEvaluator );
+ reteEvaluator.onWorkingMemoryAction( propagationEntry );
} finally {
executing.set( false );
flush();
@@ -62,7 +64,7 @@ public class SynchronizedBypassPropagationList extends
SynchronizedPropagationLi
@Override
public void flush() {
if (!executing.get()) {
- PropagationEntry head = takeAll();
+ PropagationEntry<ReteEvaluator> head = takeAll();
while (head != null) {
flush( head );
head = takeAll();
diff --git
a/drools-core/src/main/java/org/drools/core/phreak/SynchronizedPropagationList.java
b/drools-core/src/main/java/org/drools/core/phreak/SynchronizedPropagationList.java
index f2440aaa028..73a56937b3d 100644
---
a/drools-core/src/main/java/org/drools/core/phreak/SynchronizedPropagationList.java
+++
b/drools-core/src/main/java/org/drools/core/phreak/SynchronizedPropagationList.java
@@ -20,6 +20,7 @@ package org.drools.core.phreak;
import java.util.Iterator;
+import org.drools.base.phreak.PropagationEntry;
import org.drools.core.common.ReteEvaluator;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -30,8 +31,8 @@ public class SynchronizedPropagationList implements
PropagationList {
protected ReteEvaluator reteEvaluator;
- protected volatile PropagationEntry head;
- protected volatile PropagationEntry tail;
+ protected volatile PropagationEntry<ReteEvaluator> head;
+ protected volatile PropagationEntry<ReteEvaluator> tail;
protected volatile boolean disposed = false;
@@ -46,10 +47,11 @@ public class SynchronizedPropagationList implements
PropagationList {
public SynchronizedPropagationList(){}
@Override
- public void addEntry(final PropagationEntry entry) {
+ public void addEntry(final PropagationEntry<ReteEvaluator> entry) {
if (entry.requiresImmediateFlushing()) {
if (entry.isCalledFromRHS()) {
entry.execute(reteEvaluator);
+ reteEvaluator.onWorkingMemoryAction(entry);
} else {
reteEvaluator.getActivationsManager().executeTask( new
ExecutableEntry() {
@Override
@@ -59,6 +61,7 @@ public class SynchronizedPropagationList implements
PropagationList {
} else {
entry.execute( reteEvaluator );
}
+ reteEvaluator.onWorkingMemoryAction(entry);
}
@Override
@@ -72,7 +75,7 @@ public class SynchronizedPropagationList implements
PropagationList {
}
}
- synchronized void internalAddEntry( PropagationEntry entry ) {
+ synchronized void internalAddEntry( PropagationEntry<ReteEvaluator> entry
) {
if ( head == null ) {
head = entry;
if (firingUntilHalt) {
@@ -96,13 +99,14 @@ public class SynchronizedPropagationList implements
PropagationList {
}
@Override
- public void flush(PropagationEntry currentHead) {
+ public void flush(PropagationEntry<ReteEvaluator> currentHead) {
flush( reteEvaluator, currentHead );
}
- private void flush( ReteEvaluator reteEvaluator, PropagationEntry
currentHead ) {
- for (PropagationEntry entry = currentHead; !disposed && entry != null;
entry = entry.getNext()) {
+ private void flush( ReteEvaluator reteEvaluator,
PropagationEntry<ReteEvaluator> currentHead ) {
+ for (PropagationEntry<ReteEvaluator> entry = currentHead; !disposed &&
entry != null; entry = entry.getNext()) {
entry.execute(reteEvaluator);
+ reteEvaluator.onWorkingMemoryAction(entry);
}
}
@@ -111,8 +115,8 @@ public class SynchronizedPropagationList implements
PropagationList {
}
@Override
- public synchronized PropagationEntry takeAll() {
- PropagationEntry currentHead = head;
+ public synchronized PropagationEntry<ReteEvaluator> takeAll() {
+ PropagationEntry<ReteEvaluator> currentHead = head;
head = null;
tail = null;
hasEntriesDeferringExpiration = false;
@@ -146,15 +150,15 @@ public class SynchronizedPropagationList implements
PropagationList {
}
@Override
- public synchronized Iterator<PropagationEntry> iterator() {
+ public synchronized Iterator<PropagationEntry<ReteEvaluator>> iterator() {
return new PropagationEntryIterator(head);
}
- public static class PropagationEntryIterator implements
Iterator<PropagationEntry> {
+ public static class PropagationEntryIterator implements
Iterator<PropagationEntry<ReteEvaluator>> {
- private PropagationEntry next;
+ private PropagationEntry<ReteEvaluator> next;
- public PropagationEntryIterator(PropagationEntry head) {
+ public PropagationEntryIterator(PropagationEntry<ReteEvaluator> head) {
this.next = head;
}
@@ -164,8 +168,8 @@ public class SynchronizedPropagationList implements
PropagationList {
}
@Override
- public PropagationEntry next() {
- PropagationEntry current = next;
+ public PropagationEntry<ReteEvaluator> next() {
+ PropagationEntry<ReteEvaluator> current = next;
next = current.getNext();
return current;
}
diff --git
a/drools-core/src/main/java/org/drools/core/phreak/ThreadUnsafePropagationList.java
b/drools-core/src/main/java/org/drools/core/phreak/ThreadUnsafePropagationList.java
index d4f066d3338..1ebc867a894 100644
---
a/drools-core/src/main/java/org/drools/core/phreak/ThreadUnsafePropagationList.java
+++
b/drools-core/src/main/java/org/drools/core/phreak/ThreadUnsafePropagationList.java
@@ -21,6 +21,7 @@ package org.drools.core.phreak;
import java.util.Collections;
import java.util.Iterator;
+import org.drools.base.phreak.PropagationEntry;
import org.drools.core.common.ReteEvaluator;
public class ThreadUnsafePropagationList implements PropagationList {
@@ -32,12 +33,13 @@ public class ThreadUnsafePropagationList implements
PropagationList {
}
@Override
- public void addEntry( PropagationEntry propagationEntry ) {
+ public void addEntry( PropagationEntry<ReteEvaluator> propagationEntry ) {
propagationEntry.execute( reteEvaluator );
+ reteEvaluator.onWorkingMemoryAction( propagationEntry );
}
@Override
- public PropagationEntry takeAll() {
+ public PropagationEntry<ReteEvaluator> takeAll() {
return null;
}
@@ -46,7 +48,7 @@ public class ThreadUnsafePropagationList implements
PropagationList {
}
@Override
- public void flush( PropagationEntry currentHead ) {
+ public void flush( PropagationEntry<ReteEvaluator> currentHead ) {
}
@Override
@@ -64,7 +66,7 @@ public class ThreadUnsafePropagationList implements
PropagationList {
}
@Override
- public Iterator<PropagationEntry> iterator() {
+ public Iterator<PropagationEntry<ReteEvaluator>> iterator() {
return Collections.emptyIterator();
}
diff --git
a/drools-core/src/main/java/org/drools/core/common/WorkingMemoryAction.java
b/drools-core/src/main/java/org/drools/core/phreak/actions/AbstractPartitionedPropagationEntry.java
similarity index 61%
copy from
drools-core/src/main/java/org/drools/core/common/WorkingMemoryAction.java
copy to
drools-core/src/main/java/org/drools/core/phreak/actions/AbstractPartitionedPropagationEntry.java
index 63497ff8e2e..ba961cbea21 100644
--- a/drools-core/src/main/java/org/drools/core/common/WorkingMemoryAction.java
+++
b/drools-core/src/main/java/org/drools/core/phreak/actions/AbstractPartitionedPropagationEntry.java
@@ -16,17 +16,19 @@
* specific language governing permissions and limitations
* under the License.
*/
-package org.drools.core.common;
+package org.drools.core.phreak.actions;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.base.ValueResolver;
+import org.drools.base.phreak.actions.AbstractPropagationEntry;
-public interface WorkingMemoryAction extends PropagationEntry {
- short WorkingMemoryReteAssertAction = 1;
- short DeactivateCallback = 2;
- short PropagateAction = 3;
- short LogicalRetractCallback = 4;
- short WorkingMemoryReteExpireAction = 5;
- short SignalProcessInstanceAction = 6;
- short SignalAction = 7;
- short WorkingMemoryBehahviourRetract = 8;
+public abstract class AbstractPartitionedPropagationEntry<T extends
ValueResolver> extends AbstractPropagationEntry<T> {
+ protected final int partition;
+
+ protected AbstractPartitionedPropagationEntry(int partition) {
+ this.partition = partition;
+ }
+
+ protected boolean isMainPartition() {
+ return partition == 0;
+ }
}
diff --git
a/drools-core/src/main/java/org/drools/core/phreak/actions/Delete.java
b/drools-core/src/main/java/org/drools/core/phreak/actions/Delete.java
new file mode 100644
index 00000000000..20f4aa7ea12
--- /dev/null
+++ b/drools-core/src/main/java/org/drools/core/phreak/actions/Delete.java
@@ -0,0 +1,64 @@
+/*
+ * 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.drools.core.phreak.actions;
+
+import org.drools.base.phreak.PropagationEntry;
+import org.drools.base.phreak.actions.AbstractPropagationEntry;
+import org.drools.core.common.InternalFactHandle;
+import org.drools.core.common.PropagationContext;
+import org.drools.core.common.ReteEvaluator;
+import org.drools.core.reteoo.EntryPointNode;
+import org.drools.core.reteoo.ObjectTypeConf;
+
+public class Delete extends AbstractPropagationEntry<ReteEvaluator> {
+ private final EntryPointNode epn;
+ private final InternalFactHandle handle;
+ private final PropagationContext context;
+ private final ObjectTypeConf objectTypeConf;
+
+ public Delete(EntryPointNode epn, InternalFactHandle handle,
PropagationContext context, ObjectTypeConf objectTypeConf) {
+ this.epn = epn;
+ this.handle = handle;
+ this.context = context;
+ this.objectTypeConf = objectTypeConf;
+ }
+
+ public void internalExecute(ReteEvaluator reteEvaluator) {
+ execute(reteEvaluator, epn, handle, context, objectTypeConf);
+ }
+
+ public static void execute(ReteEvaluator reteEvaluator, EntryPointNode
epn, InternalFactHandle handle, PropagationContext context, ObjectTypeConf
objectTypeConf) {
+ epn.propagateRetract(handle, context, objectTypeConf, reteEvaluator);
+ }
+
+ @Override
+ public boolean isPartitionSplittable() {
+ return true;
+ }
+
+ @Override
+ public PropagationEntry<ReteEvaluator> getSplitForPartition(int
partitionNr) {
+ return new PartitionedDelete(handle, context, objectTypeConf,
partitionNr);
+ }
+
+ @Override
+ public String toString() {
+ return "Delete of " + handle.getObject();
+ }
+}
diff --git
a/drools-core/src/main/java/org/drools/core/phreak/actions/ExecuteQuery.java
b/drools-core/src/main/java/org/drools/core/phreak/actions/ExecuteQuery.java
new file mode 100644
index 00000000000..1b6dac1a0af
--- /dev/null
+++ b/drools-core/src/main/java/org/drools/core/phreak/actions/ExecuteQuery.java
@@ -0,0 +1,82 @@
+/*
+ * 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.drools.core.phreak.actions;
+
+import org.drools.base.reteoo.NodeTypeEnums;
+import org.drools.core.base.DroolsQueryImpl;
+import org.drools.core.common.InternalFactHandle;
+import org.drools.core.common.PropagationContext;
+import org.drools.core.common.ReteEvaluator;
+import org.drools.core.phreak.RuleAgendaItem;
+import org.drools.core.reteoo.LeftInputAdapterNode;
+import org.drools.core.reteoo.LeftTupleSource;
+import org.drools.core.reteoo.PathMemory;
+import org.drools.core.reteoo.QueryTerminalNode;
+import org.drools.core.reteoo.TerminalNode;
+
+public class ExecuteQuery extends PropagationEntryWithResult<ReteEvaluator,
QueryTerminalNode[]> {
+
+ private final String queryName;
+ private final DroolsQueryImpl queryObject;
+ private final InternalFactHandle handle;
+ private final PropagationContext pCtx;
+ private final boolean calledFromRHS;
+
+ public ExecuteQuery(String queryName, DroolsQueryImpl queryObject,
InternalFactHandle handle, PropagationContext pCtx, boolean calledFromRHS) {
+ this.queryName = queryName;
+ this.queryObject = queryObject;
+ this.handle = handle;
+ this.pCtx = pCtx;
+ this.calledFromRHS = calledFromRHS;
+ }
+
+ @Override
+ public void internalExecute(ReteEvaluator reteEvaluator) {
+ QueryTerminalNode[] tnodes =
reteEvaluator.getKnowledgeBase().getReteooBuilder().getTerminalNodesForQuery(queryName);
+ if (tnodes == null) {
+ throw new RuntimeException("Query '" + queryName + "' does not
exist");
+ }
+ QueryTerminalNode tnode = tnodes[0];
+ if (queryObject.getElements().length !=
tnode.getQuery().getParameters().length) {
+ throw new RuntimeException("Query '" + queryName + "' has been
invoked with a wrong number of arguments. Expected " +
+ tnode.getQuery().getParameters().length
+ ", actual " + queryObject.getElements().length);
+ }
+ LeftTupleSource lts = tnode.getLeftTupleSource();
+ while (!NodeTypeEnums.isLeftInputAdapterNode(lts)) {
+ lts = lts.getLeftTupleSource();
+ }
+ LeftInputAdapterNode lian = (LeftInputAdapterNode) lts;
+ LeftInputAdapterNode.LiaNodeMemory lmem =
reteEvaluator.getNodeMemory(lian);
+ if (lmem.getSegmentMemory() == null) {
+
reteEvaluator.getSegmentMemorySupport().getOrCreateSegmentMemory(lts, lmem);
+ }
+ LeftInputAdapterNode.doInsertObject(handle, pCtx, lian, reteEvaluator,
lmem, false, queryObject.isOpen());
+ for (PathMemory rm : lmem.getSegmentMemory().getPathMemories()) {
+ RuleAgendaItem evaluator =
reteEvaluator.getActivationsManager().createRuleAgendaItem(Integer.MAX_VALUE,
rm, (TerminalNode) rm.getPathEndNode());
+ evaluator.getRuleExecutor().setDirty(true);
+ evaluator.getRuleExecutor().evaluateNetworkAndFire(reteEvaluator,
null, 0, -1);
+ }
+ done(tnodes);
+ }
+
+ @Override
+ public boolean isCalledFromRHS() {
+ return calledFromRHS;
+ }
+}
diff --git
a/drools-core/src/main/java/org/drools/core/phreak/actions/Insert.java
b/drools-core/src/main/java/org/drools/core/phreak/actions/Insert.java
new file mode 100644
index 00000000000..c52034cb762
--- /dev/null
+++ b/drools-core/src/main/java/org/drools/core/phreak/actions/Insert.java
@@ -0,0 +1,150 @@
+/*
+ * 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.drools.core.phreak.actions;
+
+import org.drools.base.phreak.actions.AbstractPropagationEntry;
+import org.drools.core.common.DefaultEventHandle;
+import org.drools.core.common.InternalFactHandle;
+import org.drools.core.common.PropagationContext;
+import org.drools.core.common.ReteEvaluator;
+import org.drools.core.impl.WorkingMemoryReteExpireAction;
+import org.drools.base.phreak.PropagationEntry;
+import org.drools.core.reteoo.ClassObjectTypeConf;
+import org.drools.core.reteoo.ObjectTypeConf;
+import org.drools.core.reteoo.ObjectTypeNode;
+import org.drools.core.time.JobContext;
+import org.drools.core.time.impl.DefaultJobHandle;
+import org.drools.core.time.impl.PointInTimeTrigger;
+import org.kie.api.prototype.PrototypeEventInstance;
+
+import java.io.Externalizable;
+import java.io.IOException;
+import java.io.ObjectInput;
+import java.io.ObjectOutput;
+
+import static org.drools.base.rule.TypeDeclaration.NEVER_EXPIRES;
+
+public class Insert extends AbstractPropagationEntry<ReteEvaluator> implements
Externalizable {
+ private static final ObjectTypeNode.ExpireJob job = new
ObjectTypeNode.ExpireJob();
+
+ private InternalFactHandle handle;
+ private PropagationContext context;
+ private ObjectTypeConf objectTypeConf;
+
+ public Insert() {}
+
+ public Insert(InternalFactHandle handle, PropagationContext context,
ReteEvaluator reteEvaluator, ObjectTypeConf objectTypeConf) {
+ this.handle = handle;
+ this.context = context;
+ this.objectTypeConf = objectTypeConf;
+ if (handle.isEvent()) {
+ scheduleExpiration(reteEvaluator, handle, context, objectTypeConf,
reteEvaluator.getTimerService().getCurrentTime());
+ }
+ }
+
+ public static void execute(InternalFactHandle handle, PropagationContext
context, ReteEvaluator reteEvaluator, ObjectTypeConf objectTypeConf) {
+ if (handle.isEvent()) {
+ scheduleExpiration(reteEvaluator, handle, context, objectTypeConf,
reteEvaluator.getTimerService().getCurrentTime());
+ }
+ propagate(handle, context, reteEvaluator, objectTypeConf);
+ }
+
+ private static void propagate(InternalFactHandle handle,
PropagationContext context, ReteEvaluator reteEvaluator, ObjectTypeConf
objectTypeConf) {
+ if (objectTypeConf == null) {
+ objectTypeConf =
handle.getEntryPoint(reteEvaluator).getObjectTypeConfigurationRegistry().getOrCreateObjectTypeConf(handle.getEntryPointId(),
handle.getObject());
+ }
+ for (ObjectTypeNode otn : objectTypeConf.getObjectTypeNodes()) {
+ otn.propagateAssert(handle, context, reteEvaluator);
+ }
+ if (isOrphanHandle(handle, reteEvaluator)) {
+ handle.setDisconnected(true);
+
handle.getEntryPoint(reteEvaluator).getObjectStore().removeHandle(handle);
+ if (handle instanceof DefaultEventHandle eventHandle) {
+ eventHandle.unscheduleAllJobs(reteEvaluator);
+ }
+ }
+ }
+
+ private static boolean isOrphanHandle(InternalFactHandle handle,
ReteEvaluator reteEvaluator) {
+ return !handle.hasMatches() &&
!reteEvaluator.getKnowledgeBase().getKieBaseConfiguration().isMutabilityEnabled();
+ }
+
+ public void internalExecute(ReteEvaluator reteEvaluator) {
+ propagate(handle, context, reteEvaluator, objectTypeConf);
+ }
+
+ private static void scheduleExpiration(ReteEvaluator reteEvaluator,
InternalFactHandle handle, PropagationContext context, ObjectTypeConf
objectTypeConf, long insertionTime) {
+ for (ObjectTypeNode otn : objectTypeConf.getObjectTypeNodes()) {
+ long expirationOffset = objectTypeConf.isPrototype() ?
((PrototypeEventInstance) handle.getObject()).getExpiration() :
otn.getExpirationOffset();
+ scheduleExpiration(reteEvaluator, handle, context, otn,
insertionTime, expirationOffset);
+ }
+ if (objectTypeConf.getConcreteObjectTypeNode() == null) {
+ long expirationOffset = objectTypeConf.isPrototype() ?
((PrototypeEventInstance) handle.getObject()).getExpiration() :
((ClassObjectTypeConf) objectTypeConf).getExpirationOffset();
+ scheduleExpiration(reteEvaluator, handle, context, null,
insertionTime, expirationOffset);
+ }
+ }
+
+ private static void scheduleExpiration(ReteEvaluator reteEvaluator,
InternalFactHandle handle, PropagationContext context, ObjectTypeNode otn, long
insertionTime, long expirationOffset) {
+ if (expirationOffset == NEVER_EXPIRES || expirationOffset ==
Long.MAX_VALUE || context.getReaderContext() != null) {
+ return;
+ }
+ DefaultEventHandle eventFactHandle = (DefaultEventHandle) handle;
+ long nextTimestamp = getNextTimestamp(insertionTime,
expirationOffset, eventFactHandle);
+ WorkingMemoryReteExpireAction action = new
WorkingMemoryReteExpireAction((DefaultEventHandle) handle, otn);
+ if (nextTimestamp <= reteEvaluator.getTimerService().getCurrentTime())
{
+ reteEvaluator.addPropagation(action);
+ } else {
+ JobContext jobctx = new ObjectTypeNode.ExpireJobContext(action,
reteEvaluator);
+ DefaultJobHandle jobHandle = (DefaultJobHandle)
reteEvaluator.getTimerService()
+ .scheduleJob(job, jobctx,
PointInTimeTrigger.createPointInTimeTrigger(nextTimestamp, null));
+ jobctx.setJobHandle(jobHandle);
+ eventFactHandle.addJob(jobHandle);
+ }
+ }
+
+ private static long getNextTimestamp(long insertionTime, long
expirationOffset, DefaultEventHandle eventFactHandle) {
+ long effectiveEnd = eventFactHandle.getEndTimestamp() +
expirationOffset;
+ return Math.max(insertionTime, effectiveEnd >= 0 ? effectiveEnd :
Long.MAX_VALUE);
+ }
+
+ @Override
+ public String toString() {
+ return "Insert of " + handle.getObject();
+ }
+
+ public InternalFactHandle getHandle() {
+ return handle;
+ }
+
+ @Override
+ public void writeExternal(ObjectOutput out) throws IOException {
+ out.writeObject(next);
+ out.writeObject(handle);
+ out.writeObject(context);
+ }
+
+ @Override
+ public void readExternal(ObjectInput in) throws IOException,
ClassNotFoundException {
+ @SuppressWarnings("unchecked")
+ PropagationEntry<ReteEvaluator> next =
(PropagationEntry<ReteEvaluator>) in.readObject();
+ this.next = next;
+ this.handle = (InternalFactHandle) in.readObject();
+ this.context = (PropagationContext) in.readObject();
+ }
+}
diff --git
a/drools-core/src/main/java/org/drools/core/phreak/actions/PartitionedDelete.java
b/drools-core/src/main/java/org/drools/core/phreak/actions/PartitionedDelete.java
new file mode 100644
index 00000000000..e257d31a23a
--- /dev/null
+++
b/drools-core/src/main/java/org/drools/core/phreak/actions/PartitionedDelete.java
@@ -0,0 +1,57 @@
+/*
+ * 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.drools.core.phreak.actions;
+
+import org.drools.core.common.DefaultEventHandle;
+import org.drools.core.common.InternalFactHandle;
+import org.drools.core.common.PropagationContext;
+import org.drools.core.common.ReteEvaluator;
+import org.drools.core.reteoo.ObjectTypeConf;
+import org.drools.core.reteoo.ObjectTypeNode;
+
+public class PartitionedDelete extends
AbstractPartitionedPropagationEntry<ReteEvaluator> {
+ private final InternalFactHandle handle;
+ private final PropagationContext context;
+ private final ObjectTypeConf objectTypeConf;
+
+ PartitionedDelete(InternalFactHandle handle, PropagationContext context,
ObjectTypeConf objectTypeConf, int partition) {
+ super(partition);
+ this.handle = handle;
+ this.context = context;
+ this.objectTypeConf = objectTypeConf;
+ }
+
+ public void internalExecute(ReteEvaluator reteEvaluator) {
+ ObjectTypeNode[] cachedNodes = objectTypeConf.getObjectTypeNodes();
+ if (cachedNodes == null) {
+ return;
+ }
+ for (ObjectTypeNode cachedNode : cachedNodes) {
+ cachedNode.retractObject(handle, context, reteEvaluator,
partition);
+ }
+ if (handle.isEvent() && isMainPartition()) {
+ ((DefaultEventHandle) handle).unscheduleAllJobs(reteEvaluator);
+ }
+ }
+
+ @Override
+ public String toString() {
+ return "Delete of " + handle.getObject() + " for partition " +
partition;
+ }
+}
diff --git
a/drools-core/src/main/java/org/drools/core/phreak/actions/PartitionedUpdate.java
b/drools-core/src/main/java/org/drools/core/phreak/actions/PartitionedUpdate.java
new file mode 100644
index 00000000000..038bec1175c
--- /dev/null
+++
b/drools-core/src/main/java/org/drools/core/phreak/actions/PartitionedUpdate.java
@@ -0,0 +1,63 @@
+/*
+ * 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.drools.core.phreak.actions;
+
+import org.drools.core.common.InternalFactHandle;
+import org.drools.core.common.PropagationContext;
+import org.drools.core.common.ReteEvaluator;
+import org.drools.core.reteoo.CompositePartitionAwareObjectSinkAdapter;
+import org.drools.core.reteoo.ModifyPreviousTuples;
+import org.drools.core.reteoo.ObjectTypeConf;
+import org.drools.core.reteoo.ObjectTypeNode;
+
+import static
org.drools.core.reteoo.EntryPointNode.removeRightTuplesMatchingOTN;
+
+public class PartitionedUpdate extends
AbstractPartitionedPropagationEntry<ReteEvaluator> {
+ private final InternalFactHandle handle;
+ private final PropagationContext context;
+ private final ObjectTypeConf objectTypeConf;
+
+ PartitionedUpdate(InternalFactHandle handle, PropagationContext context,
ObjectTypeConf objectTypeConf, int partition) {
+ super(partition);
+ this.handle = handle;
+ this.context = context;
+ this.objectTypeConf = objectTypeConf;
+ }
+
+ public void internalExecute(ReteEvaluator reteEvaluator) {
+ ModifyPreviousTuples modifyPreviousTuples = new
ModifyPreviousTuples(handle.detachLinkedTuplesForPartition(partition));
+ ObjectTypeNode[] cachedNodes =
objectTypeConf.getObjectTypeNodes();
+ for (int i = 0, length = cachedNodes.length; i < length; i++) {
+ ObjectTypeNode otn = cachedNodes[i];
+ ((CompositePartitionAwareObjectSinkAdapter)
otn.getObjectSinkPropagator())
+ .propagateModifyObjectForPartition(handle,
modifyPreviousTuples,
+
context.adaptModificationMaskForObjectType(otn.getObjectType(), reteEvaluator),
+ reteEvaluator,
partition);
+ if (i < cachedNodes.length - 1) {
+ removeRightTuplesMatchingOTN(context, reteEvaluator,
modifyPreviousTuples, otn, partition);
+ }
+ }
+ modifyPreviousTuples.retractTuples(context, reteEvaluator);
+ }
+
+ @Override
+ public String toString() {
+ return "Update of " + handle.getObject() + " for partition " +
partition;
+ }
+}
diff --git
a/drools-core/src/main/java/org/drools/core/time/EnqueuedSelfRemovalJobContext.java
b/drools-core/src/main/java/org/drools/core/phreak/actions/PropagationEntryWithResult.java
similarity index 52%
copy from
drools-core/src/main/java/org/drools/core/time/EnqueuedSelfRemovalJobContext.java
copy to
drools-core/src/main/java/org/drools/core/phreak/actions/PropagationEntryWithResult.java
index 08b2018431f..db46750a351 100644
---
a/drools-core/src/main/java/org/drools/core/time/EnqueuedSelfRemovalJobContext.java
+++
b/drools-core/src/main/java/org/drools/core/phreak/actions/PropagationEntryWithResult.java
@@ -16,26 +16,34 @@
* specific language governing permissions and limitations
* under the License.
*/
-package org.drools.core.time;
+package org.drools.core.phreak.actions;
-import java.util.Map;
+import org.drools.base.base.ValueResolver;
+import org.drools.base.phreak.actions.AbstractPropagationEntry;
-import org.drools.core.common.ReteEvaluator;
-import org.drools.core.phreak.PropagationEntry;
-import org.drools.core.time.impl.TimerJobInstance;
+import java.util.concurrent.CountDownLatch;
-public class EnqueuedSelfRemovalJobContext extends SelfRemovalJobContext {
- public EnqueuedSelfRemovalJobContext( JobContext jobContext, Map<Long,
TimerJobInstance> timerInstances ) {
- super( jobContext, timerInstances );
+public abstract class PropagationEntryWithResult<T extends ValueResolver, R>
extends AbstractPropagationEntry<T> {
+ private final CountDownLatch done = new CountDownLatch(1);
+
+ private R result;
+
+ public final R getResult() {
+ try {
+ done.await();
+ } catch (InterruptedException e) {
+ throw new RuntimeException(e);
+ }
+ return result;
+ }
+
+ protected void done(R result) {
+ this.result = result;
+ done.countDown();
}
@Override
- public void remove() {
- getReteEvaluator().addPropagation( new
PropagationEntry.AbstractPropagationEntry() {
- @Override
- public void internalExecute(ReteEvaluator reteEvaluator) {
- timerInstances.remove( jobContext.getJobHandle().getId() );
- }
- } );
+ public boolean requiresImmediateFlushing() {
+ return true;
}
}
diff --git
a/drools-core/src/main/java/org/drools/core/phreak/actions/Update.java
b/drools-core/src/main/java/org/drools/core/phreak/actions/Update.java
new file mode 100644
index 00000000000..cda29bcae8b
--- /dev/null
+++ b/drools-core/src/main/java/org/drools/core/phreak/actions/Update.java
@@ -0,0 +1,103 @@
+/*
+ * 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.drools.core.phreak.actions;
+
+import org.drools.base.phreak.PropagationEntry;
+import org.drools.base.phreak.actions.AbstractPropagationEntry;
+import org.drools.core.common.InternalFactHandle;
+import org.drools.core.common.PropagationContext;
+import org.drools.core.common.ReteEvaluator;
+import org.drools.core.reteoo.ModifyPreviousTuples;
+import org.drools.core.reteoo.ObjectTypeConf;
+import org.drools.core.reteoo.ObjectTypeNode;
+
+import java.io.Externalizable;
+import java.io.IOException;
+import java.io.ObjectInput;
+import java.io.ObjectOutput;
+
+import static
org.drools.core.reteoo.EntryPointNode.removeRightTuplesMatchingOTN;
+
+public class Update extends AbstractPropagationEntry<ReteEvaluator> implements
Externalizable {
+ private InternalFactHandle handle;
+ private PropagationContext context;
+ private ObjectTypeConf objectTypeConf;
+
+ public Update() {}
+
+ public Update(InternalFactHandle handle, PropagationContext context,
ObjectTypeConf objectTypeConf) {
+ this.handle = handle;
+ this.context = context;
+ this.objectTypeConf = objectTypeConf;
+ }
+
+ public void internalExecute(ReteEvaluator reteEvaluator) {
+ execute(handle, context, objectTypeConf, reteEvaluator);
+ }
+
+ public InternalFactHandle getHandle() {
+ return handle;
+ }
+
+ public static void execute(InternalFactHandle handle, PropagationContext
pctx, ObjectTypeConf objectTypeConf, ReteEvaluator reteEvaluator) {
+ if (objectTypeConf == null) {
+ objectTypeConf =
handle.getEntryPoint(reteEvaluator).getObjectTypeConfigurationRegistry().getOrCreateObjectTypeConf(handle.getEntryPointId(),
handle.getObject());
+ }
+ ModifyPreviousTuples modifyPreviousTuples = new
ModifyPreviousTuples(handle.detachLinkedTuples());
+ ObjectTypeNode[] cachedNodes =
objectTypeConf.getObjectTypeNodes();
+ for (int i = 0, length = cachedNodes.length; i < length; i++) {
+ cachedNodes[i].modifyObject(handle, modifyPreviousTuples, pctx,
reteEvaluator);
+ if (i < cachedNodes.length - 1) {
+ removeRightTuplesMatchingOTN(pctx, reteEvaluator,
modifyPreviousTuples, cachedNodes[i], 0);
+ }
+ }
+ modifyPreviousTuples.retractTuples(pctx, reteEvaluator);
+ }
+
+ @Override
+ public boolean isPartitionSplittable() {
+ return true;
+ }
+
+ @Override
+ public PropagationEntry<ReteEvaluator> getSplitForPartition(int
partitionNr) {
+ return new PartitionedUpdate(handle, context, objectTypeConf,
partitionNr);
+ }
+
+ @Override
+ public String toString() {
+ return "Update of " + handle.getObject();
+ }
+
+ @Override
+ public void writeExternal(ObjectOutput out) throws IOException {
+ out.writeObject(next);
+ out.writeObject(handle);
+ out.writeObject(context);
+ }
+
+ @Override
+ public void readExternal(ObjectInput in) throws IOException,
ClassNotFoundException {
+ @SuppressWarnings("unchecked")
+ PropagationEntry<ReteEvaluator> next =
(PropagationEntry<ReteEvaluator>) in.readObject();
+ this.next = next;
+ this.handle = (InternalFactHandle) in.readObject();
+ this.context = (PropagationContext) in.readObject();
+ }
+}
diff --git
a/drools-core/src/main/java/org/drools/core/reteoo/AsyncReceiveNode.java
b/drools-core/src/main/java/org/drools/core/reteoo/AsyncReceiveNode.java
index 823eb499171..7fe6cad7128 100644
--- a/drools-core/src/main/java/org/drools/core/reteoo/AsyncReceiveNode.java
+++ b/drools-core/src/main/java/org/drools/core/reteoo/AsyncReceiveNode.java
@@ -34,7 +34,8 @@ import org.drools.core.common.Memory;
import org.drools.core.common.MemoryFactory;
import org.drools.core.common.ReteEvaluator;
import org.drools.core.common.UpdateContext;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
+import org.drools.base.phreak.actions.AbstractPropagationEntry;
import org.drools.core.reteoo.builder.BuildContext;
import org.drools.core.util.AbstractLinkedListNode;
import org.drools.core.util.index.TupleList;
@@ -110,7 +111,7 @@ public class AsyncReceiveNode extends LeftTupleSource
return objectTypeConf;
}
- public static class AsyncReceiveAction extends
PropagationEntry.AbstractPropagationEntry {
+ public static class AsyncReceiveAction extends
AbstractPropagationEntry<ReteEvaluator> {
private final AsyncReceiveNode asyncReceiveNode;
private final Object object;
diff --git
a/drools-core/src/main/java/org/drools/core/reteoo/CompositePartitionAwareObjectSinkAdapter.java
b/drools-core/src/main/java/org/drools/core/reteoo/CompositePartitionAwareObjectSinkAdapter.java
index f92107ab61b..f63cee1ac1e 100644
---
a/drools-core/src/main/java/org/drools/core/reteoo/CompositePartitionAwareObjectSinkAdapter.java
+++
b/drools-core/src/main/java/org/drools/core/reteoo/CompositePartitionAwareObjectSinkAdapter.java
@@ -35,7 +35,8 @@ import org.drools.core.common.BaseNode;
import org.drools.core.common.InternalFactHandle;
import org.drools.core.common.PropagationContext;
import org.drools.core.common.ReteEvaluator;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
+import org.drools.base.phreak.actions.AbstractPropagationEntry;
import org.drools.core.reteoo.CompositeObjectSinkAdapter.FieldIndex;
public class CompositePartitionAwareObjectSinkAdapter implements
ObjectSinkPropagator {
@@ -128,7 +129,7 @@ public class CompositePartitionAwareObjectSinkAdapter
implements ObjectSinkPropa
}
}
- public static class Insert extends
PropagationEntry.AbstractPropagationEntry {
+ public static class Insert extends AbstractPropagationEntry<ReteEvaluator>
{
private final ObjectSinkPropagator propagator;
private final InternalFactHandle factHandle;
@@ -151,7 +152,7 @@ public class CompositePartitionAwareObjectSinkAdapter
implements ObjectSinkPropa
}
}
- public static class HashedInsert extends
PropagationEntry.AbstractPropagationEntry {
+ public static class HashedInsert extends
AbstractPropagationEntry<ReteEvaluator> {
private final AlphaNode sink;
private final InternalFactHandle factHandle;
diff --git
a/drools-core/src/main/java/org/drools/core/reteoo/EntryPointNode.java
b/drools-core/src/main/java/org/drools/core/reteoo/EntryPointNode.java
index 24c65b13ca5..750911ce61c 100644
--- a/drools-core/src/main/java/org/drools/core/reteoo/EntryPointNode.java
+++ b/drools-core/src/main/java/org/drools/core/reteoo/EntryPointNode.java
@@ -39,7 +39,10 @@ import
org.drools.core.common.ObjectTypeConfigurationRegistry;
import org.drools.core.common.PropagationContext;
import org.drools.core.common.ReteEvaluator;
import org.drools.core.impl.InternalRuleBase;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
+import org.drools.core.phreak.actions.Delete;
+import org.drools.core.phreak.actions.Insert;
+import org.drools.core.phreak.actions.Update;
import org.drools.core.reteoo.builder.BuildContext;
import org.drools.util.bitmask.BitMask;
import org.slf4j.Logger;
@@ -198,9 +201,9 @@ public class EntryPointNode extends ObjectSource implements
ObjectSink {
// In case of parallel execution the
CompositePartitionAwareObjectSinkAdapter
// used by the OTNs will take care of enqueueing this insertion on
the propagation queues
// of the different agendas
- PropagationEntry.Insert.execute( handle, context, reteEvaluator,
objectTypeConf );
+ Insert.execute( handle, context, reteEvaluator, objectTypeConf );
} else {
- reteEvaluator.addPropagation( new PropagationEntry.Insert( handle,
context, reteEvaluator, objectTypeConf ) );
+ reteEvaluator.addPropagation( new Insert( handle, context,
reteEvaluator, objectTypeConf ) );
}
}
@@ -214,9 +217,9 @@ public class EntryPointNode extends ObjectSource implements
ObjectSink {
}
if (reteEvaluator.isThreadSafe()) {
- reteEvaluator.addPropagation( new PropagationEntry.Update( handle,
pctx, objectTypeConf ) );
+ reteEvaluator.addPropagation( new Update( handle, pctx,
objectTypeConf ) );
} else {
- PropagationEntry.Update.execute( handle, pctx, objectTypeConf,
reteEvaluator );
+ Update.execute( handle, pctx, objectTypeConf, reteEvaluator );
}
}
@@ -298,7 +301,7 @@ public class EntryPointNode extends ObjectSource implements
ObjectSink {
log.trace( "Delete {}", handle.toString() );
}
- reteEvaluator.addPropagation(new PropagationEntry.Delete(this, handle,
context, objectTypeConf));
+ reteEvaluator.addPropagation(new Delete(this, handle, context,
objectTypeConf));
}
public void immediateDeleteObject(InternalFactHandle handle,
PropagationContext context,
@@ -307,7 +310,7 @@ public class EntryPointNode extends ObjectSource implements
ObjectSink {
log.trace( "Delete {}", handle.toString() );
}
- PropagationEntry.Delete.execute(reteEvaluator, this, handle, context,
objectTypeConf);
+ Delete.execute(reteEvaluator, this, handle, context, objectTypeConf);
}
public void propagateRetract(InternalFactHandle handle, PropagationContext
context, ObjectTypeConf objectTypeConf, ReteEvaluator reteEvaluator) {
diff --git
a/drools-core/src/main/java/org/drools/core/rule/SlidingTimeWindow.java
b/drools-core/src/main/java/org/drools/core/rule/SlidingTimeWindow.java
index e69e49222f5..b8f06f00c51 100644
--- a/drools-core/src/main/java/org/drools/core/rule/SlidingTimeWindow.java
+++ b/drools-core/src/main/java/org/drools/core/rule/SlidingTimeWindow.java
@@ -33,7 +33,8 @@ import org.drools.core.common.PropagationContext;
import org.drools.core.common.ReteEvaluator;
import org.drools.core.common.WorkingMemoryAction;
import org.drools.core.marshalling.MarshallerReaderContext;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
+import org.drools.base.phreak.actions.AbstractPropagationEntry;
import org.drools.core.reteoo.ObjectTypeNode;
import org.drools.core.reteoo.WindowNode;
import org.drools.core.reteoo.WindowNode.WindowMemory;
@@ -357,7 +358,7 @@ public class SlidingTimeWindow
}
public static class BehaviorExpireWMAction
- extends PropagationEntry.AbstractPropagationEntry
+ extends AbstractPropagationEntry<ReteEvaluator>
implements WorkingMemoryAction {
protected BehaviorRuntime behavior;
protected BehaviorContext context;
diff --git
a/drools-core/src/main/java/org/drools/core/time/EnqueuedSelfRemovalJobContext.java
b/drools-core/src/main/java/org/drools/core/time/EnqueuedSelfRemovalJobContext.java
index 08b2018431f..6c9ecc5298b 100644
---
a/drools-core/src/main/java/org/drools/core/time/EnqueuedSelfRemovalJobContext.java
+++
b/drools-core/src/main/java/org/drools/core/time/EnqueuedSelfRemovalJobContext.java
@@ -21,7 +21,7 @@ package org.drools.core.time;
import java.util.Map;
import org.drools.core.common.ReteEvaluator;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.actions.AbstractPropagationEntry;
import org.drools.core.time.impl.TimerJobInstance;
public class EnqueuedSelfRemovalJobContext extends SelfRemovalJobContext {
@@ -31,7 +31,7 @@ public class EnqueuedSelfRemovalJobContext extends
SelfRemovalJobContext {
@Override
public void remove() {
- getReteEvaluator().addPropagation( new
PropagationEntry.AbstractPropagationEntry() {
+ getReteEvaluator().addPropagation( new
AbstractPropagationEntry<ReteEvaluator>() {
@Override
public void internalExecute(ReteEvaluator reteEvaluator) {
timerInstances.remove( jobContext.getJobHandle().getId() );
diff --git
a/drools-core/src/test/java/org/drools/core/time/impl/JDKTimerServiceTest.java
b/drools-core/src/test/java/org/drools/core/time/impl/JDKTimerServiceTest.java
index 059ed19d2e7..09d8bf175dd 100644
---
a/drools-core/src/test/java/org/drools/core/time/impl/JDKTimerServiceTest.java
+++
b/drools-core/src/test/java/org/drools/core/time/impl/JDKTimerServiceTest.java
@@ -36,7 +36,7 @@ import org.drools.core.SessionConfiguration;
import org.drools.core.common.InternalWorkingMemory;
import org.drools.core.common.ReteEvaluator;
import org.drools.core.impl.RuleBaseFactory;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
import org.drools.core.time.Job;
import org.drools.core.time.JobContext;
import org.drools.core.time.TimerService;
diff --git
a/drools-kiesession/src/main/java/org/drools/kiesession/agenda/CompositeDefaultAgenda.java
b/drools-kiesession/src/main/java/org/drools/kiesession/agenda/CompositeDefaultAgenda.java
index 12d6d25977e..903df2ed91e 100644
---
a/drools-kiesession/src/main/java/org/drools/kiesession/agenda/CompositeDefaultAgenda.java
+++
b/drools-kiesession/src/main/java/org/drools/kiesession/agenda/CompositeDefaultAgenda.java
@@ -32,7 +32,7 @@ import org.drools.core.common.RuleFlowGroup;
import org.drools.core.event.AgendaEventSupport;
import org.drools.core.impl.InternalRuleBase;
import org.drools.core.phreak.ExecutableEntry;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
import org.drools.core.phreak.PropagationList;
import org.drools.core.phreak.RuleAgendaItem;
import org.drools.core.reteoo.PathMemory;
@@ -242,7 +242,7 @@ public class CompositeDefaultAgenda implements
Externalizable, InternalAgenda {
}
@Override
- public void addPropagation( PropagationEntry propagationEntry ) {
+ public void addPropagation( PropagationEntry<ReteEvaluator>
propagationEntry ) {
if (propagationEntry.isPartitionSplittable()) {
for ( int i = 0; i < agendas.length; i++ ) {
agendas[i].addPropagation(
propagationEntry.getSplitForPartition( i ) );
@@ -274,7 +274,7 @@ public class CompositeDefaultAgenda implements
Externalizable, InternalAgenda {
}
@Override
- public Iterator<PropagationEntry> getActionsIterator() {
+ public Iterator<PropagationEntry<ReteEvaluator>> getActionsIterator() {
return new CompositeIterator<>( Stream.of( agendas ).map(
DefaultAgenda::getActionsIterator ).toArray(Iterator[]::new) );
}
diff --git
a/drools-kiesession/src/main/java/org/drools/kiesession/agenda/DefaultAgenda.java
b/drools-kiesession/src/main/java/org/drools/kiesession/agenda/DefaultAgenda.java
index f65b418eb42..c3192b696ec 100644
---
a/drools-kiesession/src/main/java/org/drools/kiesession/agenda/DefaultAgenda.java
+++
b/drools-kiesession/src/main/java/org/drools/kiesession/agenda/DefaultAgenda.java
@@ -49,7 +49,8 @@ import org.drools.core.concurrent.SequentialGroupEvaluator;
import org.drools.core.event.AgendaEventSupport;
import org.drools.core.impl.InternalRuleBase;
import org.drools.core.phreak.ExecutableEntry;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
+import org.drools.base.phreak.actions.AbstractPropagationEntry;
import org.drools.core.phreak.PropagationList;
import org.drools.core.phreak.RuleAgendaItem;
import org.drools.core.phreak.RuleExecutor;
@@ -572,7 +573,7 @@ public class DefaultAgenda implements InternalAgenda {
private int fireLoop(AgendaFilter agendaFilter, int fireLimit, RestHandler
restHandler, boolean isInternalFire) {
int fireCount = 0;
try {
- PropagationEntry head = takePropagationHead();
+ PropagationEntry<ReteEvaluator> head = takePropagationHead();
int returnedFireCount;
boolean limitReached = fireLimit == 0; // -1 or > 0 will return
false. No reason for user to give 0, just handled for completeness.
@@ -643,7 +644,7 @@ public class DefaultAgenda implements InternalAgenda {
return fireCount;
}
- private PropagationEntry takePropagationHead() {
+ private PropagationEntry<ReteEvaluator> takePropagationHead() {
if (executionStateMachine.getCurrentState().isHalting()) {
return null;
}
@@ -654,13 +655,13 @@ public class DefaultAgenda implements InternalAgenda {
RestHandler FIRE_ALL_RULES = new FireAllRulesRestHandler();
RestHandler FIRE_UNTIL_HALT = new FireUntilHaltRestHandler();
- PropagationEntry handleRest(DefaultAgenda agenda, boolean
isInternalFire);
+ PropagationEntry<ReteEvaluator> handleRest(DefaultAgenda agenda,
boolean isInternalFire);
class FireAllRulesRestHandler implements RestHandler {
@Override
- public PropagationEntry handleRest(DefaultAgenda agenda, boolean
isInternalFire) {
+ public PropagationEntry<ReteEvaluator> handleRest(DefaultAgenda
agenda, boolean isInternalFire) {
synchronized
(agenda.executionStateMachine.getStateMachineLock()) {
- PropagationEntry head = agenda.propagationList.takeAll();
+ PropagationEntry<ReteEvaluator> head =
agenda.propagationList.takeAll();
if (isInternalFire && head == null) {
agenda.internalHalt();
}
@@ -671,14 +672,14 @@ public class DefaultAgenda implements InternalAgenda {
class FireUntilHaltRestHandler implements RestHandler {
@Override
- public PropagationEntry handleRest(DefaultAgenda agenda, boolean
isInternalFire) {
+ public PropagationEntry<ReteEvaluator> handleRest(DefaultAgenda
agenda, boolean isInternalFire) {
boolean deactivated = false;
if (isInternalFire &&
agenda.executionStateMachine.getCurrentState() ==
ExecutionStateMachine.ExecutionState.FIRING_UNTIL_HALT) {
agenda.executionStateMachine.inactiveOnFireUntilHalt();
deactivated = true;
}
- PropagationEntry head;
+ PropagationEntry<ReteEvaluator> head;
// this must use the same sync target as takeAllPropagations,
to ensure this entire block is atomic, up to the point of wait
synchronized (agenda.propagationList) {
head = agenda.takePropagationHead();
@@ -748,7 +749,7 @@ public class DefaultAgenda implements InternalAgenda {
return executionStateMachine.tryDeactivate();
}
- static class Halt extends PropagationEntry.AbstractPropagationEntry {
+ static class Halt extends AbstractPropagationEntry<ReteEvaluator> {
private final ExecutionStateMachine executionStateMachine;
@@ -768,7 +769,7 @@ public class DefaultAgenda implements InternalAgenda {
}
}
- static class ImmediateHalt extends
PropagationEntry.AbstractPropagationEntry {
+ static class ImmediateHalt extends AbstractPropagationEntry<ReteEvaluator>
{
private final ExecutionStateMachine executionStateMachine;
private final PropagationList propagationList;
@@ -796,7 +797,7 @@ public class DefaultAgenda implements InternalAgenda {
// This will place a halt command on the propagation queue
// that will allow the engine to halt safely
if ( isFiring() ) {
- PropagationEntry halt = executionStateMachine.getCurrentState() ==
ExecutionStateMachine.ExecutionState.FIRING_ALL_RULES ?
+ PropagationEntry<ReteEvaluator> halt =
executionStateMachine.getCurrentState() ==
ExecutionStateMachine.ExecutionState.FIRING_ALL_RULES ?
new ImmediateHalt(executionStateMachine, propagationList) :
new Halt(executionStateMachine);
propagationList.addEntry(halt);
@@ -855,7 +856,7 @@ public class DefaultAgenda implements InternalAgenda {
}
@Override
- public void addPropagation(PropagationEntry propagationEntry) {
+ public void addPropagation(PropagationEntry<ReteEvaluator>
propagationEntry) {
propagationList.addEntry( propagationEntry );
}
@@ -870,7 +871,7 @@ public class DefaultAgenda implements InternalAgenda {
}
@Override
- public Iterator<PropagationEntry> getActionsIterator() {
+ public Iterator<PropagationEntry<ReteEvaluator>> getActionsIterator() {
return propagationList.iterator();
}
diff --git
a/drools-kiesession/src/main/java/org/drools/kiesession/consequence/StatefulKnowledgeSessionForRHS.java
b/drools-kiesession/src/main/java/org/drools/kiesession/consequence/StatefulKnowledgeSessionForRHS.java
index b077395b8b6..6c9dbd56f2e 100644
---
a/drools-kiesession/src/main/java/org/drools/kiesession/consequence/StatefulKnowledgeSessionForRHS.java
+++
b/drools-kiesession/src/main/java/org/drools/kiesession/consequence/StatefulKnowledgeSessionForRHS.java
@@ -45,7 +45,7 @@ import org.drools.core.common.SegmentMemorySupport;
import org.drools.core.event.AgendaEventSupport;
import org.drools.core.event.RuleEventListenerSupport;
import org.drools.core.event.RuleRuntimeEventSupport;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
import org.drools.core.phreak.RuleNetworkEvaluator;
import org.drools.core.reteoo.EntryPointNode;
import org.drools.core.reteoo.TerminalNode;
@@ -679,7 +679,7 @@ public class StatefulKnowledgeSessionForRHS
delegate.closeLiveQuery(factHandle);
}
- public void addPropagation(PropagationEntry propagationEntry) {
+ public void addPropagation(PropagationEntry<ReteEvaluator>
propagationEntry) {
delegate.addPropagation(propagationEntry);
}
@@ -699,7 +699,7 @@ public class StatefulKnowledgeSessionForRHS
return delegate.tryDeactivate();
}
- public Iterator<? extends PropagationEntry> getActionsIterator() {
+ public Iterator<? extends PropagationEntry<ReteEvaluator>>
getActionsIterator() {
return delegate.getActionsIterator();
}
diff --git
a/drools-kiesession/src/main/java/org/drools/kiesession/session/StatefulKnowledgeSessionImpl.java
b/drools-kiesession/src/main/java/org/drools/kiesession/session/StatefulKnowledgeSessionImpl.java
index 14479ab79cd..22dfd2687c7 100644
---
a/drools-kiesession/src/main/java/org/drools/kiesession/session/StatefulKnowledgeSessionImpl.java
+++
b/drools-kiesession/src/main/java/org/drools/kiesession/session/StatefulKnowledgeSessionImpl.java
@@ -64,7 +64,10 @@ import org.drools.core.impl.AbstractRuntime;
import org.drools.core.impl.EnvironmentFactory;
import org.drools.core.management.DroolsManagementAgent;
import org.drools.core.marshalling.MarshallerReaderContext;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
+import org.drools.base.phreak.actions.AbstractPropagationEntry;
+import org.drools.core.phreak.actions.ExecuteQuery;
+import org.drools.core.phreak.actions.PropagationEntryWithResult;
import org.drools.core.phreak.RuleAgendaItem;
import org.drools.core.phreak.RuleNetworkEvaluator;
import org.drools.core.phreak.RuleNetworkEvaluatorImpl;
@@ -254,7 +257,7 @@ public class StatefulKnowledgeSessionImpl extends
AbstractRuntime
private NamedEntryPointsManager entryPointsManager;
- private Consumer<PropagationEntry> workingMemoryActionListener;
+ private Consumer<PropagationEntry<ReteEvaluator>>
workingMemoryActionListener;
private boolean tmsEnabled;
@@ -465,12 +468,12 @@ public class StatefulKnowledgeSessionImpl extends
AbstractRuntime
}
@Override
- public Consumer<PropagationEntry> getWorkingMemoryActionListener() {
+ public Consumer<PropagationEntry<ReteEvaluator>>
getWorkingMemoryActionListener() {
return workingMemoryActionListener;
}
@Override
- public void setWorkingMemoryActionListener(Consumer<PropagationEntry>
workingMemoryActionListener) {
+ public void
setWorkingMemoryActionListener(Consumer<PropagationEntry<ReteEvaluator>>
workingMemoryActionListener) {
this.workingMemoryActionListener = workingMemoryActionListener;
}
@@ -730,7 +733,7 @@ public class StatefulKnowledgeSessionImpl extends
AbstractRuntime
}
private QueryTerminalNode[] evalQuery(final String queryName, final
DroolsQueryImpl queryObject, final InternalFactHandle handle, final
PropagationContext pCtx, final boolean isCalledFromRHS) {
- PropagationEntry.ExecuteQuery executeQuery = new
PropagationEntry.ExecuteQuery( kBase, queryName, queryObject, handle, pCtx,
isCalledFromRHS);
+ ExecuteQuery executeQuery = new ExecuteQuery( queryName, queryObject,
handle, pCtx, isCalledFromRHS);
addPropagation( executeQuery );
return executeQuery.getResult();
}
@@ -747,7 +750,7 @@ public class StatefulKnowledgeSessionImpl extends
AbstractRuntime
}
}
- private class ExecuteCloseLiveQuery extends
PropagationEntry.PropagationEntryWithResult<Void> {
+ private class ExecuteCloseLiveQuery extends
PropagationEntryWithResult<ReteEvaluator, Void> {
private final InternalFactHandle factHandle;
@@ -1251,7 +1254,7 @@ public class StatefulKnowledgeSessionImpl extends
AbstractRuntime
}
public void submit(AtomicAction action) {
- agenda.addPropagation( new PropagationEntry.AbstractPropagationEntry()
{
+ agenda.addPropagation( new AbstractPropagationEntry<ReteEvaluator>() {
@Override
public void internalExecute(ReteEvaluator reteEvaluator ) {
action.execute( (KieSession)reteEvaluator );
@@ -1613,7 +1616,7 @@ public class StatefulKnowledgeSessionImpl extends
AbstractRuntime
}
@Override
- public void addPropagation(PropagationEntry propagationEntry) {
+ public void addPropagation(PropagationEntry<ReteEvaluator>
propagationEntry) {
agenda.addPropagation( propagationEntry );
}
@@ -1627,7 +1630,7 @@ public class StatefulKnowledgeSessionImpl extends
AbstractRuntime
}
@Override
- public Iterator<? extends PropagationEntry> getActionsIterator() {
+ public Iterator<? extends PropagationEntry<ReteEvaluator>>
getActionsIterator() {
return agenda.getActionsIterator();
}
diff --git
a/drools-kiesession/src/test/java/org/drools/kiesession/ReteooWorkingMemoryTest.java
b/drools-kiesession/src/test/java/org/drools/kiesession/ReteooWorkingMemoryTest.java
index 40cc28adb01..b7903954acc 100644
---
a/drools-kiesession/src/test/java/org/drools/kiesession/ReteooWorkingMemoryTest.java
+++
b/drools-kiesession/src/test/java/org/drools/kiesession/ReteooWorkingMemoryTest.java
@@ -30,7 +30,8 @@ import org.drools.base.common.RuleBasePartitionId;
import org.drools.core.common.TruthMaintenanceSystem;
import org.drools.core.common.TruthMaintenanceSystemFactory;
import org.drools.core.common.WorkingMemoryAction;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
+import org.drools.base.phreak.actions.AbstractPropagationEntry;
import org.drools.core.reteoo.EntryPointNode;
import org.drools.core.reteoo.Rete;
import org.drools.core.reteoo.builder.NodeFactory;
@@ -201,7 +202,7 @@ public class ReteooWorkingMemoryTest {
}
private static class ReentrantAction
- extends PropagationEntry.AbstractPropagationEntry
+ extends AbstractPropagationEntry<ReteEvaluator>
implements WorkingMemoryAction {
// I am using AtomicInteger just as an int wrapper... nothing to do
with concurrency here
public AtomicInteger counter = new AtomicInteger(0);
diff --git
a/drools-reliability/drools-reliability-core/src/main/java/org/drools/reliability/core/ReliablePropagationList.java
b/drools-reliability/drools-reliability-core/src/main/java/org/drools/reliability/core/ReliablePropagationList.java
index cd456546ed4..50b46194eb4 100644
---
a/drools-reliability/drools-reliability-core/src/main/java/org/drools/reliability/core/ReliablePropagationList.java
+++
b/drools-reliability/drools-reliability-core/src/main/java/org/drools/reliability/core/ReliablePropagationList.java
@@ -20,7 +20,7 @@ package org.drools.reliability.core;
import org.drools.core.common.ReteEvaluator;
import org.drools.core.common.Storage;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
import org.drools.core.phreak.SynchronizedPropagationList;
import java.io.Externalizable;
@@ -47,8 +47,8 @@ public class ReliablePropagationList extends
SynchronizedPropagationList impleme
}
@Override
- public synchronized PropagationEntry takeAll() {
- PropagationEntry p = super.takeAll();
+ public synchronized PropagationEntry<ReteEvaluator> takeAll() {
+ PropagationEntry<ReteEvaluator> p = super.takeAll();
Storage<String, Object> componentsStorage =
StorageManagerFactory.get().getStorageManager().getOrCreateStorageForSession(this.reteEvaluator,
"components");
componentsStorage.put(PROPAGATION_LIST, this);
return p;
@@ -63,9 +63,10 @@ public class ReliablePropagationList extends
SynchronizedPropagationList impleme
out.writeBoolean(firingUntilHalt);
}
@Override
+ @SuppressWarnings("unchecked")
public void readExternal(ObjectInput in) throws IOException,
ClassNotFoundException {
- this.head = (PropagationEntry) in.readObject();
- this.tail = (PropagationEntry) in.readObject();
+ this.head = (PropagationEntry<ReteEvaluator>) in.readObject();
+ this.tail = (PropagationEntry<ReteEvaluator>) in.readObject();
this.disposed = in.readBoolean();
this.hasEntriesDeferringExpiration = in.readBoolean();
this.firingUntilHalt = in.readBoolean();
diff --git
a/drools-reliability/drools-reliability-core/src/main/java/org/drools/reliability/core/ReliableSessionInitializer.java
b/drools-reliability/drools-reliability-core/src/main/java/org/drools/reliability/core/ReliableSessionInitializer.java
index dac9faf940f..412cbdacac2 100644
---
a/drools-reliability/drools-reliability-core/src/main/java/org/drools/reliability/core/ReliableSessionInitializer.java
+++
b/drools-reliability/drools-reliability-core/src/main/java/org/drools/reliability/core/ReliableSessionInitializer.java
@@ -27,8 +27,11 @@ import org.drools.core.WorkingMemoryEntryPoint;
import org.drools.core.common.InternalFactHandle;
import org.drools.core.common.InternalWorkingMemory;
import org.drools.core.common.InternalWorkingMemoryEntryPoint;
+import org.drools.core.common.ReteEvaluator;
import org.drools.core.common.Storage;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
+import org.drools.core.phreak.actions.Insert;
+import org.drools.core.phreak.actions.Update;
import org.drools.reliability.core.util.ReliabilityUtils;
import org.kie.api.event.rule.AfterMatchFiredEvent;
import org.kie.api.event.rule.DefaultAgendaEventListener;
@@ -84,9 +87,9 @@ public class ReliableSessionInitializer {
return session;
}
- private void onWorkingMemoryAction(InternalWorkingMemory session,
PropagationEntry entry) {
- if (entry instanceof PropagationEntry.Insert || entry instanceof
PropagationEntry.Update) {
- InternalFactHandle fh =
((PropagationEntry.AbstractPropagationEntry) entry).getHandle();
+ private void onWorkingMemoryAction(InternalWorkingMemory session,
PropagationEntry<ReteEvaluator> entry) {
+ if (entry instanceof Insert || entry instanceof Update) {
+ InternalFactHandle fh = entry instanceof Insert ? ((Insert)
entry).getHandle() : ((Update) entry).getHandle();
if (fh.isValid()) {
WorkingMemoryEntryPoint ep = fh.getEntryPoint(session);
((SimpleReliableObjectStore)
ep.getObjectStore()).putIntoPersistedStorage(fh, true);
diff --git
a/drools-ruleunits/drools-ruleunits-impl/src/main/java/org/drools/ruleunits/impl/sessions/RuleUnitExecutorImpl.java
b/drools-ruleunits/drools-ruleunits-impl/src/main/java/org/drools/ruleunits/impl/sessions/RuleUnitExecutorImpl.java
index 968acd5e18c..72308a32138 100644
---
a/drools-ruleunits/drools-ruleunits-impl/src/main/java/org/drools/ruleunits/impl/sessions/RuleUnitExecutorImpl.java
+++
b/drools-ruleunits/drools-ruleunits-impl/src/main/java/org/drools/ruleunits/impl/sessions/RuleUnitExecutorImpl.java
@@ -53,7 +53,8 @@ import org.drools.core.event.RuleEventListenerSupport;
import org.drools.core.event.RuleRuntimeEventSupport;
import org.drools.core.impl.ActivationsManagerImpl;
import org.drools.core.impl.InternalRuleBase;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
+import org.drools.core.phreak.actions.ExecuteQuery;
import org.drools.core.phreak.RuleNetworkEvaluator;
import org.drools.core.phreak.RuleNetworkEvaluatorImpl;
import org.drools.core.phreak.SegmentMemorySupportImpl;
@@ -231,7 +232,7 @@ public class RuleUnitExecutorImpl implements ReteEvaluator {
}
@Override
- public void addPropagation(PropagationEntry propagationEntry) {
+ public void addPropagation(PropagationEntry<ReteEvaluator>
propagationEntry) {
activationsManager.addPropagation( propagationEntry );
}
@@ -337,7 +338,7 @@ public class RuleUnitExecutorImpl implements ReteEvaluator {
final PropagationContext pCtx = new
PhreakPropagationContext(getNextPropagationIdCounter(),
PropagationContext.Type.INSERTION, null, null, handle,
getDefaultEntryPointId());
- PropagationEntry.ExecuteQuery executeQuery = new
PropagationEntry.ExecuteQuery( ruleBase, queryName, queryObject, handle, pCtx,
false);
+ ExecuteQuery executeQuery = new ExecuteQuery( queryName, queryObject,
handle, pCtx, false);
addPropagation( executeQuery );
TerminalNode[] terminalNodes = executeQuery.getResult();
diff --git
a/drools-serialization-protobuf/src/main/java/org/drools/serialization/protobuf/ProtobufOutputMarshaller.java
b/drools-serialization-protobuf/src/main/java/org/drools/serialization/protobuf/ProtobufOutputMarshaller.java
index 3b57ec5a258..64bae8ee43d 100644
---
a/drools-serialization-protobuf/src/main/java/org/drools/serialization/protobuf/ProtobufOutputMarshaller.java
+++
b/drools-serialization-protobuf/src/main/java/org/drools/serialization/protobuf/ProtobufOutputMarshaller.java
@@ -42,7 +42,7 @@ import org.drools.core.common.RuleFlowGroup;
import org.drools.core.common.TruthMaintenanceSystem;
import org.drools.core.common.TruthMaintenanceSystemFactory;
import org.drools.core.marshalling.MarshallerWriteContext;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
import org.drools.core.phreak.RuleAgendaItem;
import org.drools.core.process.WorkItem;
import org.drools.core.reteoo.TupleImpl;
@@ -419,14 +419,14 @@ public class ProtobufOutputMarshaller {
public static void writeActionQueue( MarshallerWriteContext context,
ProtobufMessages.RuleData.Builder
_session) throws IOException {
- Iterator<? extends PropagationEntry> i =
context.getWorkingMemory().getActionsIterator();
+ Iterator<? extends PropagationEntry<?>> i =
context.getWorkingMemory().getActionsIterator();
if ( !i.hasNext() ) {
return;
}
ProtobufMessages.ActionQueue.Builder _queue =
ProtobufMessages.ActionQueue.newBuilder();
while ( i.hasNext() ) {
- PropagationEntry entry = i.next();
+ PropagationEntry<?> entry = i.next();
if (entry instanceof ProtobufWorkingMemoryAction) {
_queue.addAction(((ProtobufWorkingMemoryAction)
entry).serialize(context));
}
diff --git
a/drools-serialization-protobuf/src/main/java/org/drools/serialization/protobuf/WorkingMemoryReteAssertAction.java
b/drools-serialization-protobuf/src/main/java/org/drools/serialization/protobuf/WorkingMemoryReteAssertAction.java
index 724898cd791..42e513b9572 100644
---
a/drools-serialization-protobuf/src/main/java/org/drools/serialization/protobuf/WorkingMemoryReteAssertAction.java
+++
b/drools-serialization-protobuf/src/main/java/org/drools/serialization/protobuf/WorkingMemoryReteAssertAction.java
@@ -27,14 +27,15 @@ import org.drools.core.common.WorkingMemoryAction;
import org.drools.base.definitions.InternalKnowledgePackage;
import org.drools.base.definitions.rule.impl.RuleImpl;
import org.drools.core.marshalling.MarshallerReaderContext;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
+import org.drools.base.phreak.actions.AbstractPropagationEntry;
import org.drools.core.reteoo.RuntimeComponentFactory;
import org.drools.core.common.PropagationContext;
import org.drools.core.reteoo.TerminalNode;
import org.drools.core.reteoo.Tuple;
public class WorkingMemoryReteAssertAction
- extends PropagationEntry.AbstractPropagationEntry
+ extends AbstractPropagationEntry<ReteEvaluator>
implements WorkingMemoryAction {
protected InternalFactHandle factHandle;
diff --git
a/drools-test-coverage/test-compiler-integration/src/test/java/org/drools/mvel/compiler/command/PropagationListTest.java
b/drools-test-coverage/test-compiler-integration/src/test/java/org/drools/mvel/compiler/command/PropagationListTest.java
index 02d0b2e55a9..e4a0437f286 100644
---
a/drools-test-coverage/test-compiler-integration/src/test/java/org/drools/mvel/compiler/command/PropagationListTest.java
+++
b/drools-test-coverage/test-compiler-integration/src/test/java/org/drools/mvel/compiler/command/PropagationListTest.java
@@ -25,7 +25,8 @@ import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import org.drools.core.common.ReteEvaluator;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
+import org.drools.base.phreak.actions.AbstractPropagationEntry;
import org.drools.core.phreak.PropagationList;
import org.drools.core.phreak.SynchronizedPropagationList;
import org.junit.jupiter.api.Disabled;
@@ -126,7 +127,7 @@ public class PropagationListTest {
};
}
- public static class TestEntry extends
PropagationEntry.AbstractPropagationEntry {
+ public static class TestEntry extends
AbstractPropagationEntry<ReteEvaluator> {
final Checker checker;
final int i;
diff --git
a/drools-tms/src/main/java/org/drools/tms/beliefsystem/simple/BeliefSystemLogicalCallback.java
b/drools-tms/src/main/java/org/drools/tms/beliefsystem/simple/BeliefSystemLogicalCallback.java
index c8673c729a8..a00c4b119eb 100644
---
a/drools-tms/src/main/java/org/drools/tms/beliefsystem/simple/BeliefSystemLogicalCallback.java
+++
b/drools-tms/src/main/java/org/drools/tms/beliefsystem/simple/BeliefSystemLogicalCallback.java
@@ -24,7 +24,8 @@ import org.drools.core.common.ReteEvaluator;
import org.drools.core.common.SuperCacheFixer;
import org.drools.core.common.WorkingMemoryAction;
import org.drools.core.marshalling.MarshallerReaderContext;
-import org.drools.core.phreak.PropagationEntry;
+import org.drools.base.phreak.PropagationEntry;
+import org.drools.base.phreak.actions.AbstractPropagationEntry;
import org.drools.core.reteoo.ObjectTypeConf;
import org.drools.core.rule.consequence.InternalMatch;
import org.drools.kiesession.entrypoints.NamedEntryPoint;
@@ -35,7 +36,7 @@ import java.io.IOException;
import static
org.drools.base.reteoo.PropertySpecificUtil.allSetButTraitBitMask;
-public class BeliefSystemLogicalCallback extends
PropagationEntry.AbstractPropagationEntry implements WorkingMemoryAction {
+public class BeliefSystemLogicalCallback extends
AbstractPropagationEntry<ReteEvaluator> implements WorkingMemoryAction {
protected InternalFactHandle handle;
protected PropagationContext context;
diff --git
a/kogito-codegen-modules/kogito-codegen-processes-integration-tests/src/test/java/org/kie/kogito/codegen/tests/CallActivityTaskIT.java
b/kogito-codegen-modules/kogito-codegen-processes-integration-tests/src/test/java/org/kie/kogito/codegen/tests/CallActivityTaskIT.java
index e06b5c9799f..ecf97ad13c8 100644
---
a/kogito-codegen-modules/kogito-codegen-processes-integration-tests/src/test/java/org/kie/kogito/codegen/tests/CallActivityTaskIT.java
+++
b/kogito-codegen-modules/kogito-codegen-processes-integration-tests/src/test/java/org/kie/kogito/codegen/tests/CallActivityTaskIT.java
@@ -217,7 +217,7 @@ public class CallActivityTaskIT extends AbstractCodegenIT {
}
/**
- * Verifies that a subprocess whose process id contains hyphens
+ * Verifies that a subprocess whose process id contains hyphens
*/
@Test
public void testCallActivityWithHyphenatedSubProcessId() throws Exception {
diff --git
a/kogito-jbpm/jbpm-flow/src/main/java/org/jbpm/process/instance/LightProcessRuntime.java
b/kogito-jbpm/jbpm-flow/src/main/java/org/jbpm/process/instance/LightProcessRuntime.java
index 0148e40713c..bd236f1cc02 100755
---
a/kogito-jbpm/jbpm-flow/src/main/java/org/jbpm/process/instance/LightProcessRuntime.java
+++
b/kogito-jbpm/jbpm-flow/src/main/java/org/jbpm/process/instance/LightProcessRuntime.java
@@ -24,10 +24,10 @@ import java.util.List;
import java.util.Map;
import org.drools.base.definitions.rule.impl.RuleImpl;
+import org.drools.base.phreak.actions.AbstractPropagationEntry;
import org.drools.core.common.InternalKnowledgeRuntime;
import org.drools.core.common.ReteEvaluator;
import org.drools.core.common.WorkingMemoryAction;
-import org.drools.core.phreak.PropagationEntry;
import org.jbpm.process.core.event.EventFilter;
import org.jbpm.process.core.event.EventTypeFilter;
import org.jbpm.ruleflow.core.RuleFlowProcess;
@@ -393,7 +393,7 @@ public class LightProcessRuntime extends
AbstractProcessRuntime {
this.processInstanceManager.clearProcessInstancesState();
}
- public class SignalManagerSignalAction extends
PropagationEntry.AbstractPropagationEntry implements WorkingMemoryAction {
+ public class SignalManagerSignalAction extends
AbstractPropagationEntry<ReteEvaluator> implements WorkingMemoryAction {
private String type;
private Object event;
diff --git
a/kogito-jbpm/jbpm-flow/src/main/java/org/jbpm/process/instance/ProcessRuntimeImpl.java
b/kogito-jbpm/jbpm-flow/src/main/java/org/jbpm/process/instance/ProcessRuntimeImpl.java
index 0cd50cde8d8..4b6cc0e8e60 100755
---
a/kogito-jbpm/jbpm-flow/src/main/java/org/jbpm/process/instance/ProcessRuntimeImpl.java
+++
b/kogito-jbpm/jbpm-flow/src/main/java/org/jbpm/process/instance/ProcessRuntimeImpl.java
@@ -24,11 +24,11 @@ import java.util.List;
import java.util.Map;
import org.drools.base.definitions.rule.impl.RuleImpl;
+import org.drools.base.phreak.actions.AbstractPropagationEntry;
import org.drools.core.common.InternalKnowledgeRuntime;
import org.drools.core.common.InternalWorkingMemory;
import org.drools.core.common.ReteEvaluator;
import org.drools.core.common.WorkingMemoryAction;
-import org.drools.core.phreak.PropagationEntry;
import org.drools.core.time.TimerService;
import org.drools.core.time.impl.CommandServiceTimerJobFactoryManager;
import org.drools.core.time.impl.ThreadSafeTrackableTimeJobFactoryManager;
@@ -465,7 +465,7 @@ public class ProcessRuntimeImpl extends
AbstractProcessRuntime {
}
}
- public class SignalManagerSignalAction extends
PropagationEntry.AbstractPropagationEntry implements WorkingMemoryAction {
+ public class SignalManagerSignalAction extends
AbstractPropagationEntry<ReteEvaluator> implements WorkingMemoryAction {
private String type;
private Object event;
diff --git
a/kogito-jbpm/jbpm-flow/src/main/java/org/jbpm/process/instance/event/DefaultSignalManager.java
b/kogito-jbpm/jbpm-flow/src/main/java/org/jbpm/process/instance/event/DefaultSignalManager.java
index bfcd99c803c..337ab320daf 100755
---
a/kogito-jbpm/jbpm-flow/src/main/java/org/jbpm/process/instance/event/DefaultSignalManager.java
+++
b/kogito-jbpm/jbpm-flow/src/main/java/org/jbpm/process/instance/event/DefaultSignalManager.java
@@ -26,13 +26,13 @@ import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.CopyOnWriteArrayList;
+import org.drools.base.phreak.actions.AbstractPropagationEntry;
import org.drools.core.common.InternalKnowledgeRuntime;
import org.drools.core.common.InternalWorkingMemory;
import org.drools.core.common.ReteEvaluator;
import org.drools.core.common.WorkingMemoryAction;
import org.drools.core.marshalling.MarshallerReaderContext;
import org.drools.core.marshalling.MarshallerWriteContext;
-import org.drools.core.phreak.PropagationEntry;
import org.jbpm.process.instance.InternalProcessRuntime;
import org.kie.api.runtime.process.EventListener;
import org.kie.api.runtime.process.ProcessInstance;
@@ -96,7 +96,7 @@ public class DefaultSignalManager implements SignalManager {
}
}
- public static class SignalProcessInstanceAction extends
PropagationEntry.AbstractPropagationEntry implements WorkingMemoryAction {
+ public static class SignalProcessInstanceAction extends
AbstractPropagationEntry<ReteEvaluator> implements WorkingMemoryAction {
private String processInstanceId;
private String type;
@@ -160,7 +160,7 @@ public class DefaultSignalManager implements SignalManager {
}
- public static class SignalAction extends
PropagationEntry.AbstractPropagationEntry implements WorkingMemoryAction {
+ public static class SignalAction extends
AbstractPropagationEntry<ReteEvaluator> implements WorkingMemoryAction {
private String type;
private Object event;
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]