From dcad3e8a9c871d8d2252f45e98844a87ccbed032 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Thu, 30 Jul 2026 15:59:02 +0300 Subject: [PATCH 01/24] raw --- .../ignite/internal/CoreMessagesProvider.java | 7 ++ .../continuous/ContinuousRoutineInfo.java | 48 +++++----- ...nuousRoutinesJoiningNodeDiscoveryData.java | 22 ++--- .../continuous/GridContinuousProcessor.java | 88 ++++++++----------- 4 files changed, 81 insertions(+), 84 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java b/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java index 9cd8a97159dae..15d799833b6e0 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java @@ -192,8 +192,11 @@ import org.apache.ignite.internal.processors.cluster.ClusterUpdateNotifierDataBagItem; import org.apache.ignite.internal.processors.cluster.NodeFullMetricsMessage; import org.apache.ignite.internal.processors.cluster.NodeMetricsMessage; +import org.apache.ignite.internal.processors.continuous.ContinuousRoutineInfo; import org.apache.ignite.internal.processors.continuous.ContinuousRoutineStartResultMessage; +import org.apache.ignite.internal.processors.continuous.ContinuousRoutinesJoiningNodeDiscoveryData; import org.apache.ignite.internal.processors.continuous.GridContinuousMessage; +import org.apache.ignite.internal.processors.continuous.GridContinuousProcessor; import org.apache.ignite.internal.processors.continuous.StartRequestData; import org.apache.ignite.internal.processors.continuous.StartRoutineAckDiscoveryMessage; import org.apache.ignite.internal.processors.continuous.StartRoutineDiscoveryMessage; @@ -610,6 +613,10 @@ public CoreMessagesProvider(Marshaller dfltMarsh, Marshaller schemaAwareMarsh) { withNoSchema(QueryProposalsDataBagItem.class); withNoSchema(QueryEntityMessage.class); withNoSchema(QueryEntityExMessage.class); + withNoSchema(ContinuousRoutineInfo.class); + withNoSchema(ContinuousRoutinesJoiningNodeDiscoveryData.class); + withNoSchema(GridContinuousProcessor.LocalRoutineInfo.class); + withNoSchema(GridContinuousProcessor.DiscoveryData.class); // [11200 - 11300]: Compute, distributed process messages. msgIdx = 11200; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutineInfo.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutineInfo.java index 207a8f4fda984..d339acecdebef 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutineInfo.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutineInfo.java @@ -17,40 +17,48 @@ package org.apache.ignite.internal.processors.continuous; -import java.io.Serializable; import java.util.UUID; +import org.apache.ignite.internal.Order; import org.apache.ignite.internal.util.typedef.internal.S; +import org.apache.ignite.plugin.extensions.communication.Message; -/** - * - */ -class ContinuousRoutineInfo implements Serializable { - /** */ - private static final long serialVersionUID = 0L; - +/** */ +public class ContinuousRoutineInfo implements Message { /** */ + @Order(0) UUID srcNodeId; /** */ - final UUID routineId; + @Order(1) + UUID routineId; /** */ - final byte[] hnd; + @Order(2) + byte[] hnd; /** */ - final byte[] nodeFilter; + @Order(3) + byte[] nodeFilter; /** */ - final int bufSize; + @Order(4) + int bufSize; /** */ - final long interval; + @Order(5) + long interval; /** */ - final boolean autoUnsubscribe; + @Order(6) + boolean autoUnsubscribe; - /** */ - transient boolean disconnected; + /** Transient. */ + boolean disconnected; + + /** Empty constructor for serialization purposes. */ + public ContinuousRoutineInfo() { + // No-op. + } /** * @param srcNodeId Source node ID. @@ -79,16 +87,12 @@ class ContinuousRoutineInfo implements Serializable { this.autoUnsubscribe = autoUnsubscribe; } - /** - * @param srcNodeId Source node ID. - */ + /** @param srcNodeId Source node ID. */ void sourceNodeId(UUID srcNodeId) { this.srcNodeId = srcNodeId; } - /** - * - */ + /** */ void onDisconnected() { disconnected = true; } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesJoiningNodeDiscoveryData.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesJoiningNodeDiscoveryData.java index 9be6ef8e07eab..8c40160ed119b 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesJoiningNodeDiscoveryData.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesJoiningNodeDiscoveryData.java @@ -17,23 +17,23 @@ package org.apache.ignite.internal.processors.continuous; -import java.io.Serializable; import java.util.List; +import org.apache.ignite.internal.Order; import org.apache.ignite.internal.util.typedef.internal.S; +import org.apache.ignite.plugin.extensions.communication.Message; -/** - * - */ -public class ContinuousRoutinesJoiningNodeDiscoveryData implements Serializable { +/** */ +public class ContinuousRoutinesJoiningNodeDiscoveryData implements Message { /** */ - private static final long serialVersionUID = 0L; + @Order(0) + List startedRoutines; - /** */ - final List startedRoutines; + /** Empty constructor for serialization purposes. */ + public ContinuousRoutinesJoiningNodeDiscoveryData() { + // No-op. + } - /** - * @param startedRoutines Routines registered on nodes, to be started in cluster. - */ + /** @param startedRoutines Routines registered on nodes, to be started in cluster. */ ContinuousRoutinesJoiningNodeDiscoveryData(List startedRoutines) { this.startedRoutines = startedRoutines; } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java index fecd065ff7551..16e7dba962d2d 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java @@ -21,7 +21,6 @@ import java.io.IOException; import java.io.ObjectInput; import java.io.ObjectOutput; -import java.io.Serializable; import java.util.ArrayList; import java.util.Collection; import java.util.Collections; @@ -55,6 +54,7 @@ import org.apache.ignite.internal.IgniteInternalFuture; import org.apache.ignite.internal.IgniteInterruptedCheckedException; import org.apache.ignite.internal.NodeStoppingException; +import org.apache.ignite.internal.Order; import org.apache.ignite.internal.cluster.ClusterTopologyCheckedException; import org.apache.ignite.internal.managers.communication.ErrorMessage; import org.apache.ignite.internal.managers.communication.GridMessageListener; @@ -408,7 +408,7 @@ public void unlockStopping() { return; } - Serializable data = getDiscoveryData(dataBag.joiningNodeId()); + var data = getDiscoveryData(dataBag.joiningNodeId()); if (data != null) dataBag.addJoiningNodeData(CONTINUOUS_PROC.ordinal(), data); @@ -422,7 +422,7 @@ public void unlockStopping() { return; } - Serializable data = getDiscoveryData(dataBag.joiningNodeId()); + Message data = getDiscoveryData(dataBag.joiningNodeId()); if (data != null) dataBag.addNodeSpecificData(CONTINUOUS_PROC.ordinal(), data); @@ -431,13 +431,14 @@ public void unlockStopping() { /** * @param joiningNodeId Joining node id. */ - private Serializable getDiscoveryData(UUID joiningNodeId) { + private @Nullable DiscoveryData getDiscoveryData(UUID joiningNodeId) { if (log.isDebugEnabled()) { log.debug("collectDiscoveryData [node=" + joiningNodeId + - ", loc=" + ctx.localNodeId() + - ", locInfos=" + locInfos + - ", clientInfos=" + clientInfos + - ']'); + ", loc=" + ctx.localNodeId() + + ", locInfos=" + locInfos + + ", clientInfos=" + clientInfos + + ']' + ); } if (!joiningNodeId.equals(ctx.localNodeId()) || !locInfos.isEmpty()) { @@ -458,14 +459,16 @@ private Serializable getDiscoveryData(UUID joiningNodeId) { assert !ctx.config().isPeerClassLoadingEnabled() || !(info.hnd instanceof CacheContinuousQueryHandler) || - ((CacheContinuousQueryHandler)info.hnd).isMarshalled(); + ((CacheContinuousQueryHandler)info.hnd).isMarshalled(); - data.addItem(new DiscoveryDataItem(routineId, + data.addItem(new DiscoveryDataItem( + routineId, info.prjPred, info.hnd, info.bufSize, info.interval, - info.autoUnsubscribe)); + info.autoUnsubscribe + )); } return data; @@ -1971,30 +1974,31 @@ public static interface RoutineInfo { boolean delayedRegister(); } - /** - * Local routine info. - */ - public static class LocalRoutineInfo implements Serializable, RoutineInfo { - /** */ - private static final long serialVersionUID = 0L; - + /** Local routine info. */ + public static class LocalRoutineInfo implements Message, RoutineInfo { /** Source node id. */ - private final UUID nodeId; + @Order(0) + UUID nodeId; /** Projection predicate. */ - private final IgnitePredicate prjPred; + @Order(1) + IgnitePredicate prjPred; /** Continuous routine handler. */ - private final GridContinuousHandler hnd; + @Order(2) + GridContinuousHandler hnd; /** Buffer size. */ - private final int bufSize; + @Order(3) + int bufSize; /** Time interval. */ - private final long interval; + @Order(4) + long interval; /** Automatic unsubscribe flag. */ - private boolean autoUnsubscribe; + @Order(5) + boolean autoUnsubscribe; /** * @param nodeId Node id. @@ -2286,26 +2290,22 @@ IgniteBiTuple checkInterval() { } } - /** - * Discovery data. - */ - private static class DiscoveryData implements Externalizable { - /** */ - private static final long serialVersionUID = 0L; - + /** Discovery data. */ + public static class DiscoveryData implements Message { /** Node ID. */ - private UUID nodeId; + @Order(0) + UUID nodeId; /** Items. */ @GridToStringInclude - private Collection items; + @Order(1) + Collection items; /** */ - private Map> clientInfos; + @Order(2) + Map> clientInfos; - /** - * Required by {@link Externalizable}. - */ + /** Empty constructor for serialization purposes. */ public DiscoveryData() { // No-op. } @@ -2331,20 +2331,6 @@ public void addItem(DiscoveryDataItem item) { items.add(item); } - /** {@inheritDoc} */ - @Override public void writeExternal(ObjectOutput out) throws IOException { - U.writeUuid(out, nodeId); - U.writeCollection(out, items); - U.writeMap(out, clientInfos); - } - - /** {@inheritDoc} */ - @Override public void readExternal(ObjectInput in) throws IOException, ClassNotFoundException { - nodeId = U.readUuid(in); - items = U.readCollection(in); - clientInfos = U.readMap(in); - } - /** {@inheritDoc} */ @Override public String toString() { return S.toString(DiscoveryData.class, this); From 26f0524d1cb4b44625143571b5bc663c5eb36dbb Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Thu, 30 Jul 2026 18:23:04 +0300 Subject: [PATCH 02/24] raw --- .../ignite/internal/CoreMessagesProvider.java | 15 +- .../internal/GridEventConsumeHandler.java | 77 ++-- .../internal/GridMessageListenHandler.java | 109 +++--- .../deployment/GridDeploymentInfoBean.java | 2 +- .../CacheContinuousQueryDeployableObject.java | 37 +- .../CacheContinuousQueryHandler.java | 283 ++++----------- .../CacheContinuousQueryHandlerMessage.java | 209 +++++++++++ .../ContinousRoutineDiscoveryData.java | 72 ++++ .../ContinousRoutineDiscoveryDataItem.java | 96 +++++ .../continuous/ContinousRoutineLocalInfo.java | 131 +++++++ .../continuous/ContinuousRoutineInfo.java | 12 +- ...nuousRoutinesJoiningNodeDiscoveryData.java | 8 +- .../continuous/GridContinuousHandler.java | 4 +- .../continuous/GridContinuousProcessor.java | 341 +++--------------- ...CacheContinuousQueryEntriesExpireTest.java | 4 +- .../continuous/GridEventConsumeSelfTest.java | 13 +- 16 files changed, 774 insertions(+), 639 deletions(-) create mode 100644 modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandlerMessage.java create mode 100644 modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryData.java create mode 100644 modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryDataItem.java create mode 100644 modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineLocalInfo.java diff --git a/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java b/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java index 15d799833b6e0..0501ccd36f44f 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java @@ -174,7 +174,9 @@ import org.apache.ignite.internal.processors.cache.query.GridCacheQueryResponse; import org.apache.ignite.internal.processors.cache.query.GridCacheSqlQuery; import org.apache.ignite.internal.processors.cache.query.continuous.CacheContinuousQueryBatchAck; +import org.apache.ignite.internal.processors.cache.query.continuous.CacheContinuousQueryDeployableObject; import org.apache.ignite.internal.processors.cache.query.continuous.CacheContinuousQueryEntry; +import org.apache.ignite.internal.processors.cache.query.continuous.CacheContinuousQueryHandlerMessage; import org.apache.ignite.internal.processors.cache.transactions.IgniteTxEntry; import org.apache.ignite.internal.processors.cache.transactions.IgniteTxKey; import org.apache.ignite.internal.processors.cache.transactions.TxEntryValueHolder; @@ -192,11 +194,13 @@ import org.apache.ignite.internal.processors.cluster.ClusterUpdateNotifierDataBagItem; import org.apache.ignite.internal.processors.cluster.NodeFullMetricsMessage; import org.apache.ignite.internal.processors.cluster.NodeMetricsMessage; +import org.apache.ignite.internal.processors.continuous.ContinousRoutineDiscoveryData; +import org.apache.ignite.internal.processors.continuous.ContinousRoutineDiscoveryDataItem; +import org.apache.ignite.internal.processors.continuous.ContinousRoutineLocalInfo; import org.apache.ignite.internal.processors.continuous.ContinuousRoutineInfo; import org.apache.ignite.internal.processors.continuous.ContinuousRoutineStartResultMessage; import org.apache.ignite.internal.processors.continuous.ContinuousRoutinesJoiningNodeDiscoveryData; import org.apache.ignite.internal.processors.continuous.GridContinuousMessage; -import org.apache.ignite.internal.processors.continuous.GridContinuousProcessor; import org.apache.ignite.internal.processors.continuous.StartRequestData; import org.apache.ignite.internal.processors.continuous.StartRoutineAckDiscoveryMessage; import org.apache.ignite.internal.processors.continuous.StartRoutineDiscoveryMessage; @@ -615,8 +619,13 @@ public CoreMessagesProvider(Marshaller dfltMarsh, Marshaller schemaAwareMarsh) { withNoSchema(QueryEntityExMessage.class); withNoSchema(ContinuousRoutineInfo.class); withNoSchema(ContinuousRoutinesJoiningNodeDiscoveryData.class); - withNoSchema(GridContinuousProcessor.LocalRoutineInfo.class); - withNoSchema(GridContinuousProcessor.DiscoveryData.class); + withNoSchema(CacheContinuousQueryDeployableObject.class); + withSchema(CacheContinuousQueryHandlerMessage.class); + withSchema(GridEventConsumeHandler.class); + withSchema(GridMessageListenHandler.class); + withSchema(ContinousRoutineLocalInfo.class); + withSchema(ContinousRoutineDiscoveryDataItem.class); + withNoSchema(ContinousRoutineDiscoveryData.class); // [11200 - 11300]: Compute, distributed process messages. msgIdx = 11200; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java index cad29b7b12909..485b22d50e9d3 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java @@ -61,10 +61,7 @@ /** * Continuous routine handler for remote event listening. */ -class GridEventConsumeHandler implements GridContinuousHandler { - /** */ - private static final long serialVersionUID = 0L; - +public final class GridEventConsumeHandler implements GridContinuousHandler, MarshallableMessage { /** Default callback. */ private static final IgniteBiPredicate DFLT_CALLBACK = new P2() { @Override public boolean apply(UUID uuid, Event e) { @@ -76,19 +73,32 @@ class GridEventConsumeHandler implements GridContinuousHandler { private IgniteBiPredicate cb; /** Filter. */ - private IgnitePredicate filter; + @Nullable IgnitePredicate filter; - /** Serialized filter. */ - private byte[] filterBytes; + /** Serialized {@link #filter}. */ + @Order(0) + @Nullable byte[] filterBytes; - /** Deployment class name. */ - private String clsName; + /** Deployment class name. Is {@code null} if P2P deployment is disabled. */ + @Order(1) + @Nullable String clsName; - /** Deployment info. */ - private GridDeploymentInfo depInfo; + /** Deployment info. Is {@code null} if P2P deployment is disabled. */ + @Order(2) + @Nullable GridDeploymentInfoBean depInfo; /** Types. */ - private int[] types; + @Order(3) + int[] types; + + /** + * Lever of own marshaling. + * + * @see #p2pMarshal(GridKernalContext) + * @see #marshal(Marshaller) + */ + @Order(4) + boolean externalMarshaling; /** Listener. */ private GridLocalEventListener lsnr; @@ -225,8 +235,6 @@ private void initFilter(IgnitePredicate filter, GridKernalContext ctx) th EventWrapper wrapper = new EventWrapper(evt); if (evt instanceof CacheEvent) { - String cacheName = ((CacheEvent)evt).cacheName(); - ClusterNode node = ctx.discovery().node(t3.get1()); if (node == null) @@ -390,11 +398,10 @@ private boolean filterDropsEvent(Event evt) { /** {@inheritDoc} */ @Override public void p2pMarshal(GridKernalContext ctx) throws IgniteCheckedException { - assert ctx != null; assert ctx.config().isPeerClassLoadingEnabled(); if (filter != null) { - Class cls = U.detectClass(filter); + Class cls = U.detectClass(filter); clsName = cls.getName(); @@ -406,13 +413,14 @@ private boolean filterDropsEvent(Event evt) { depInfo = new GridDeploymentInfoBean(dep); filterBytes = U.marshal(ctx.marshaller(), filter); + + externalMarshaling = true; } } /** {@inheritDoc} */ @Override public void p2pUnmarshal(UUID nodeId, GridKernalContext ctx) throws IgniteCheckedException { assert nodeId != null; - assert ctx != null; assert ctx.config().isPeerClassLoadingEnabled(); if (filterBytes != null) { @@ -471,36 +479,21 @@ private boolean filterDropsEvent(Event evt) { } /** {@inheritDoc} */ - @Override public void writeExternal(ObjectOutput out) throws IOException { - boolean b = filterBytes != null; - - out.writeBoolean(b); + @Override public void marshal(Marshaller marsh) throws IgniteCheckedException { + assert clsName == null ^ depInfo == null; + assert depInfo == null ^ !externalMarshaling; - if (b) { - U.writeByteArray(out, filterBytes); - U.writeString(out, clsName); - out.writeObject(depInfo); - } - else - out.writeObject(filter); - - out.writeObject(types); + if (filter != null && !externalMarshaling) + filterBytes = marsh.marshal(filter); } /** {@inheritDoc} */ - @Override public void readExternal(ObjectInput in) throws IOException, ClassNotFoundException { - boolean b = in.readBoolean(); - - if (b) { - p2pUnmarshalFut = new GridFutureAdapter<>(); - filterBytes = U.readByteArray(in); - clsName = U.readString(in); - depInfo = (GridDeploymentInfo)in.readObject(); - } - else - filter = (IgnitePredicate)in.readObject(); + @Override public void unmarshal(Marshaller marsh, ClassLoader clsLdr) throws IgniteCheckedException { + if (externalMarshaling) + return; - types = (int[])in.readObject(); + if (filterBytes != null) + filter = marsh.unmarshal(filterBytes, clsLdr); } /** diff --git a/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java index e59873f4b449f..155e0eacef842 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java @@ -17,10 +17,6 @@ package org.apache.ignite.internal; -import java.io.Externalizable; -import java.io.IOException; -import java.io.ObjectInput; -import java.io.ObjectOutput; import java.util.Collection; import java.util.Collections; import java.util.Map; @@ -39,41 +35,48 @@ import org.apache.ignite.internal.util.typedef.internal.S; import org.apache.ignite.internal.util.typedef.internal.U; import org.apache.ignite.lang.IgniteBiPredicate; +import org.apache.ignite.marshaller.Marshaller; import org.jetbrains.annotations.Nullable; /** * Continuous handler for message subscription. */ -public class GridMessageListenHandler implements GridContinuousHandler { +public final class GridMessageListenHandler implements GridContinuousHandler, MarshallableMessage { /** */ - private static final long serialVersionUID = 0L; + private @Nullable Object topic; - /** */ - private Object topic; + /** Marshalled {@link #topic}. */ + @Order(0) + @Nullable byte[] topicBytes; /** */ private IgniteBiPredicate pred; - /** */ - private byte[] topicBytes; - - /** */ - private byte[] predBytes; + /** Marshalled {@link #pred}. */ + @Order(1) + byte[] predBytes; - /** */ - private String clsName; + /** Is {@code null} if the P2P deployment is disabled. */ + @Order(2) + @Nullable String clsName; - /** */ - private GridDeploymentInfoBean depInfo; + /** Is {@code null} if the P2P deployment is disabled. */ + @Order(3) + @Nullable GridDeploymentInfoBean depInfo; - /** */ - private boolean depEnabled; + /** + * Lever of the own marshaling. + * + * @see #p2pMarshal(GridKernalContext) + */ + @Order(4) + boolean externalMarshal; /** P2P unmarshalling future. */ - private IgniteInternalFuture p2pUnmarshalFut = new GridFinishedFuture<>(); + private final IgniteInternalFuture p2pUnmarshalFut = new GridFinishedFuture<>(); /** - * Required by {@link Externalizable}. + * Empty constructor for serialization purposes */ public GridMessageListenHandler() { // No-op. @@ -168,7 +171,7 @@ public GridMessageListenHandler(@Nullable Object topic, IgniteBiPredicate(); - topicBytes = U.readByteArray(in); - predBytes = U.readByteArray(in); - clsName = U.readString(in); - depInfo = (GridDeploymentInfoBean)in.readObject(); - } - else { - topic = in.readObject(); - pred = (IgniteBiPredicate)in.readObject(); - } - } - /** {@inheritDoc} */ @Override public String toString() { return S.toString(GridMessageListenHandler.class, this); diff --git a/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentInfoBean.java b/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentInfoBean.java index ba9f48c3ca24a..b0aab5d4548da 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentInfoBean.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentInfoBean.java @@ -30,7 +30,7 @@ /** * Deployment info bean. */ -public class GridDeploymentInfoBean implements Message, GridDeploymentInfo, Serializable { +public final class GridDeploymentInfoBean implements Message, GridDeploymentInfo, Serializable { /** */ private static final long serialVersionUID = 0L; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryDeployableObject.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryDeployableObject.java index c4e4005095ff4..647e4feeae175 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryDeployableObject.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryDeployableObject.java @@ -17,40 +17,37 @@ package org.apache.ignite.internal.processors.cache.query.continuous; -import java.io.Externalizable; -import java.io.IOException; -import java.io.ObjectInput; -import java.io.ObjectOutput; import java.util.UUID; import org.apache.ignite.IgniteCheckedException; import org.apache.ignite.internal.GridKernalContext; import org.apache.ignite.internal.IgniteDeploymentCheckedException; +import org.apache.ignite.internal.Order; import org.apache.ignite.internal.managers.deployment.GridDeployment; -import org.apache.ignite.internal.managers.deployment.GridDeploymentInfo; import org.apache.ignite.internal.managers.deployment.GridDeploymentInfoBean; import org.apache.ignite.internal.util.tostring.GridToStringExclude; import org.apache.ignite.internal.util.typedef.internal.S; import org.apache.ignite.internal.util.typedef.internal.U; +import org.apache.ignite.plugin.extensions.communication.Message; /** * Deployable object. */ -class CacheContinuousQueryDeployableObject implements Externalizable { - /** */ - private static final long serialVersionUID = 0L; - +public class CacheContinuousQueryDeployableObject implements Message { /** Serialized object. */ @GridToStringExclude - private byte[] bytes; + @Order(0) + byte[] bytes; /** Deployment class name. */ - private String clsName; + @Order(1) + String clsName; /** Deployment info. */ - private GridDeploymentInfo depInfo; + @Order(2) + GridDeploymentInfoBean depInfo; /** - * Required by {@link Externalizable}. + * Empty constructor for serialization purposes. */ public CacheContinuousQueryDeployableObject() { // No-op. @@ -97,20 +94,6 @@ T unmarshal(UUID nodeId, GridKernalContext ctx) throws IgniteCheckedExceptio return U.unmarshal(ctx, bytes, U.resolveClassLoader(dep.classLoader(), ctx.config())); } - /** {@inheritDoc} */ - @Override public void writeExternal(ObjectOutput out) throws IOException { - U.writeByteArray(out, bytes); - U.writeString(out, clsName); - out.writeObject(depInfo); - } - - /** {@inheritDoc} */ - @Override public void readExternal(ObjectInput in) throws IOException, ClassNotFoundException { - bytes = U.readByteArray(in); - clsName = U.readString(in); - depInfo = (GridDeploymentInfo)in.readObject(); - } - /** {@inheritDoc} */ @Override public String toString() { return S.toString(CacheContinuousQueryDeployableObject.class, this); diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java index 30668521d71dc..1620a145f1bd0 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java @@ -17,9 +17,7 @@ package org.apache.ignite.internal.processors.cache.query.continuous; -import java.io.Externalizable; import java.io.IOException; -import java.io.ObjectInput; import java.io.ObjectOutput; import java.util.ArrayList; import java.util.Collection; @@ -83,6 +81,7 @@ import org.apache.ignite.lang.IgniteAsyncCallback; import org.apache.ignite.lang.IgniteBiTuple; import org.apache.ignite.lang.IgniteClosure; +import org.apache.ignite.marshaller.Marshaller; import org.jetbrains.annotations.NotNull; import org.jetbrains.annotations.Nullable; @@ -96,10 +95,7 @@ /** * Continuous query handler. */ -public class CacheContinuousQueryHandler implements GridContinuousHandler { - /** */ - private static final long serialVersionUID = 0L; - +public final class CacheContinuousQueryHandler extends CacheContinuousQueryHandlerMessage implements GridContinuousHandler { /** @see #IGNITE_CONTINUOUS_QUERY_BACKUP_ACK_THRESHOLD */ public static final int DFLT_CONTINUOUS_QUERY_BACKUP_ACK_THRESHOLD = 100; @@ -131,8 +127,8 @@ public class CacheContinuousQueryHandler implements GridContinuousHandler * Transformer implementation for processing received remote events. * They are already transformed so we simply return transformed value for event. */ - private transient IgniteClosure, ?> returnValTrans = - new IgniteClosure, Object>() { + private IgniteClosure, ?> returnValTrans = + new IgniteClosure<>() { @Override public Object apply(CacheEntryEvent evt) { assert evt.getKey() == null; @@ -140,125 +136,77 @@ public class CacheContinuousQueryHandler implements GridContinuousHandler } }; - /** Cache name. */ - private String cacheName; - - /** Topic for ordered messages. */ - private Object topic; - /** P2P unmarshalling future. */ - protected transient IgniteInternalFuture p2pUnmarshalFut = new GridFinishedFuture<>(); + protected IgniteInternalFuture p2pUnmarshalFut = new GridFinishedFuture<>(); /** Initialization future. */ - protected transient IgniteInternalFuture initFut; + protected IgniteInternalFuture initFut; /** Local listener. */ - private transient CacheEntryUpdatedListener locLsnr; - - /** Remote filter. */ - private CacheEntryEventSerializableFilter rmtFilter; - - /** Deployable object for filter. */ - private CacheContinuousQueryDeployableObject rmtFilterDep; - - /** Remote filter factory. */ - private Factory rmtFilterFactory; - - /** Deployable object for filter factory. */ - private CacheContinuousQueryDeployableObject rmtFilterFactoryDep; + private CacheEntryUpdatedListener locLsnr; /** Remote filter created by {@link #rmtFilterFactory}. */ - private transient CacheEntryEventFilter rmtFilterFromFactory; - - /** Event types for JCache API. */ - private byte types; - - /** Remote transformer factory. */ - private Factory, ?>> rmtTransFactory; - - /** Deployable object for transformer factory. */ - private CacheContinuousQueryDeployableObject rmtTransFactoryDep; + private CacheEntryEventFilter rmtFilterFromFactory; /** Remote transformer created by {@link #rmtTransFactory}. */ - private transient IgniteClosure, ?> rmtTrans; + private IgniteClosure, ?> rmtTrans; /** Local listener for transformed events. */ - private transient EventListener locTransLsnr; - - /** Internal flag. */ - private boolean internal; - - /** Notify existing flag. */ - private boolean notifyExisting; - - /** Old value required flag. */ - private boolean oldValRequired; - - /** Synchronous flag. */ - private boolean sync; - - /** Ignore expired events flag. */ - private boolean ignoreExpired; - - /** Task name hash code. */ - private int taskHash; + private EventListener locTransLsnr; /** Whether to skip primary check for REPLICATED cache. */ - private transient boolean skipPrimaryCheck; - - /** */ - private transient boolean locOnly; - + boolean skipPrimaryCheck; + /** */ - private boolean keepBinary; + private boolean locOnly; /** */ - private transient ConcurrentMap rcvs; + private ConcurrentMap rcvs; /** */ - private transient ConcurrentMap entryBufs; + private ConcurrentMap entryBufs; /** */ - private transient CacheContinuousQueryAcknowledgeBuffer ackBuf; + private CacheContinuousQueryAcknowledgeBuffer ackBuf; /** */ - private transient int cacheId; + private int cacheId; /** */ - private transient volatile Map initUpdCntrs; + private volatile Map initUpdCntrs; /** */ - private transient volatile Map> initUpdCntrsPerNode; + private volatile Map> initUpdCntrsPerNode; /** */ - private transient volatile AffinityTopologyVersion initTopVer; + private volatile AffinityTopologyVersion initTopVer; /** */ - private transient volatile boolean nodeLeft; + private volatile boolean nodeLeft; /** */ - private transient boolean ignoreClsNotFound; + private boolean ignoreClsNotFound; /** */ - transient boolean asyncCb; + boolean asyncCb; /** */ - private transient UUID nodeId; + private UUID nodeId; /** */ - private transient UUID routineId; + private UUID routineId; /** Local update counters values on listener start. Used for skipping events fired before the listener start. */ - private transient volatile Map locInitUpdCntrs; + private volatile Map locInitUpdCntrs; /** */ - private transient GridKernalContext ctx; + private GridKernalContext ctx; /** */ - private transient IgniteLogger log; + private IgniteLogger log; /** - * Required by {@link Externalizable}. + * Empty constructor for serialization purposes. */ public CacheContinuousQueryHandler() { // No-op. @@ -1389,76 +1337,72 @@ CacheContinuousQueryEventBuffer partitionBuffer(GridCacheContext cctx, int /** {@inheritDoc} */ @Override public void p2pMarshal(GridKernalContext ctx) throws IgniteCheckedException { - assert ctx != null; assert ctx.config().isPeerClassLoadingEnabled(); - if (requiresDeployment(rmtFilter)) - rmtFilterDep = new CacheContinuousQueryDeployableObject(rmtFilter, ctx); + externalMarshal(ctx); + } + + /** Processes {@link #p2pUnmarshalFut} over the super's method. */ + @Override protected CacheContinuousQueryDeployableObject marshalDeployable( + Object deployable, + GridKernalContext ctx + ) throws IgniteCheckedException { + CacheContinuousQueryDeployableObject res = super.marshalDeployable(deployable, ctx); - if (requiresDeployment(rmtFilterFactory)) - rmtFilterFactoryDep = new CacheContinuousQueryDeployableObject(rmtFilterFactory, ctx); + if (p2pUnmarshalFut == null) + p2pUnmarshalFut = new GridFutureAdapter<>(); + else if (p2pUnmarshalFut.isDone()) + p2pUnmarshalFut = new GridFutureAdapter<>(); - if (requiresDeployment(rmtTransFactory)) - rmtTransFactoryDep = new CacheContinuousQueryDeployableObject(rmtTransFactory, ctx); + return res; } /** {@inheritDoc} */ @Override public void p2pUnmarshal(UUID nodeId, GridKernalContext ctx) throws IgniteCheckedException { - assert nodeId != null; - assert ctx != null; assert ctx.config().isPeerClassLoadingEnabled(); - if (rmtFilterDep != null) - rmtFilter = p2pUnmarshal(rmtFilterDep, nodeId, ctx); - - if (rmtFilterFactoryDep != null) - rmtFilterFactory = p2pUnmarshal(rmtFilterFactoryDep, nodeId, ctx); - - if (rmtTransFactoryDep != null) - rmtTransFactory = p2pUnmarshal(rmtTransFactoryDep, nodeId, ctx); + externalUnmarshal(nodeId, ctx); if (!p2pUnmarshalFut.isDone()) ((GridFutureAdapter)p2pUnmarshalFut).onDone(); } - /** - * @return Whether the handler is marshalled for peer class loading. - */ - public boolean isMarshalled() { - return (!requiresDeployment(rmtFilter) || rmtFilterDep != null) - && (!requiresDeployment(rmtFilterFactory) || rmtFilterFactoryDep != null) - && (!requiresDeployment(rmtTransFactory) || rmtTransFactoryDep != null); + /**{@inheritDoc} */ + @Override protected T externalUnmarshal( + CacheContinuousQueryDeployableObject depObj, + UUID nodeId, + GridKernalContext ctx + ) throws IgniteCheckedException { + try { + return super.externalUnmarshal(depObj, nodeId, ctx); + } + catch (IgniteCheckedException e) { + ((GridFutureAdapter)p2pUnmarshalFut).onDone(e); + + throw e; + } + catch (ExceptionInInitializerError e) { + IgniteCheckedException err = new IgniteCheckedException("Failed to unmarshal deployable object.", e); + + ((GridFutureAdapter)p2pUnmarshalFut).onDone(err); + + throw err; + } } - /** - * @param depObj Deployable object to unmarshal. - * @param nodeId Sender node Id. - * @param ctx Kernal context. - * @param Result type. - * @return Unmarshalled object. - * @throws IgniteCheckedException In case of unmarshalling failures. - */ - protected T p2pUnmarshal(CacheContinuousQueryDeployableObject depObj, - UUID nodeId, GridKernalContext ctx) throws IgniteCheckedException { - if (depObj != null) { - try { - return depObj.unmarshal(nodeId, ctx); - } - catch (IgniteCheckedException e) { - ((GridFutureAdapter)p2pUnmarshalFut).onDone(e); + /** {@inheritDoc} */ + @Override public void marshal(Marshaller marsh) throws IgniteCheckedException { + /** @see #marshalDeployable(Object, GridKernalContext) */ + p2pUnmarshalFut = null; - throw e; - } - catch (ExceptionInInitializerError e) { - IgniteCheckedException err = new IgniteCheckedException("Failed to unmarshal deployable object.", e); + super.marshal(marsh); + } - ((GridFutureAdapter)p2pUnmarshalFut).onDone(err); + /** {@inheritDoc} */ + @Override public void unmarshal(Marshaller marsh, ClassLoader clsLdr) throws IgniteCheckedException { + super.unmarshal(marsh, clsLdr); - throw err; - } - } - else - return null; + cacheId = CU.cacheId(cacheName); } /** {@inheritDoc} */ @@ -1554,78 +1498,6 @@ private void sendBackupAcknowledge(final IgniteBiTuple, Set(); - } - else - rmtFilter = (CacheEntryEventSerializableFilter)in.readObject(); - - internal = in.readBoolean(); - notifyExisting = in.readBoolean(); - oldValRequired = in.readBoolean(); - sync = in.readBoolean(); - ignoreExpired = in.readBoolean(); - taskHash = in.readInt(); - keepBinary = in.readBoolean(); - - b = in.readBoolean(); - - if (b) { - rmtFilterFactoryDep = (CacheContinuousQueryDeployableObject)in.readObject(); - - if (p2pUnmarshalFut.isDone()) - p2pUnmarshalFut = new GridFutureAdapter<>(); - } - else - rmtFilterFactory = (Factory)in.readObject(); - - types = in.readByte(); - - b = in.readBoolean(); - - if (b) { - rmtTransFactoryDep = (CacheContinuousQueryDeployableObject)in.readObject(); - - if (p2pUnmarshalFut.isDone()) - p2pUnmarshalFut = new GridFutureAdapter<>(); - } - else - rmtTransFactory = (Factory, ?>>)in.readObject(); - - cacheId = CU.cacheId(cacheName); - } - /** */ private static void writeDeployable(ObjectOutput out, Object obj, CacheContinuousQueryDeployableObject dep) throws IOException { boolean b = dep != null; @@ -1812,9 +1684,4 @@ private Object transform(IgniteClosure Map partitionContinuesQueryEntryBuffers() { return Collections.unmodifiableMap(entryBufs); } - - /** */ - private static boolean requiresDeployment(@Nullable Object obj) { - return obj != null && !U.isGrid(obj.getClass()); - } } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandlerMessage.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandlerMessage.java new file mode 100644 index 0000000000000..340d8bbef1ba2 --- /dev/null +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandlerMessage.java @@ -0,0 +1,209 @@ +/* + * 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.apache.ignite.internal.processors.cache.query.continuous; + +import java.util.UUID; +import javax.cache.configuration.Factory; +import javax.cache.event.CacheEntryEvent; +import javax.cache.event.CacheEntryEventFilter; +import org.apache.ignite.IgniteCheckedException; +import org.apache.ignite.cache.CacheEntryEventSerializableFilter; +import org.apache.ignite.internal.GridKernalContext; +import org.apache.ignite.internal.MarshallableMessage; +import org.apache.ignite.internal.Marshalled; +import org.apache.ignite.internal.Order; +import org.apache.ignite.internal.util.typedef.internal.U; +import org.apache.ignite.lang.IgniteClosure; +import org.apache.ignite.marshaller.Marshaller; +import org.jetbrains.annotations.Nullable; + +/** */ +public class CacheContinuousQueryHandlerMessage implements MarshallableMessage { + /** Remote filter. */ + CacheEntryEventSerializableFilter rmtFilter; + + /** Lever of own marshaling. */ + @Order(0) + public boolean[] externalMarshaling; + + /** Deployable object for {@link #rmtFilter}. Is {@code null} if no external marsshalling used. */ + @Order(1) + @Nullable CacheContinuousQueryDeployableObject rmtFilterDep; + + /** Marshalled {@link #rmtFilter} if {@link #rmtFilterDep} is {@code null}. */ + @Order(2) + @Nullable byte[] rmtFilterBytes; + + /** Remote filter factory. */ + @Nullable Factory rmtFilterFactory; + + /** Deployable object for {@link #rmtFilterFactory}. Is {@code null} if no external marsshalling used. */ + @Order(3) + CacheContinuousQueryDeployableObject rmtFilterFactoryDep; + + /** Marshalled {@link #rmtFilterFactory} if {@link #rmtFilterFactoryDep} is {@code null}. */ + @Order(4) + @Nullable byte[] rmtFilterFactoryBytes; + + /** Remote transformer factory. */ + Factory, ?>> rmtTransFactory; + + /** Deployable object for {@link #rmtTransFactory}. Is {@code null} if no external marsshalling used. */ + @Order(5) + CacheContinuousQueryDeployableObject rmtTransFactoryDep; + + /** Marshalled {@link #rmtTransFactory} if {@link #rmtTransFactoryDep} is {@code null}. */ + @Order(6) + @Nullable byte[] rmtTransFactoryBytes; + + /** Cache name. */ + @Order(7) + String cacheName; + + /** Topic for ordered messages. */ + @Marshalled("topicBytes") + Object topic; + + /** Marshalled {@link #topic}. */ + @Order(8) + byte[] topicBytes; + + /** Internal flag. */ + @Order(9) + boolean internal; + + /** Notify existing flag. */ + @Order(10) + boolean notifyExisting; + + /** Old value required flag. */ + @Order(11) + boolean oldValRequired; + + /** Synchronous flag. */ + @Order(12) + boolean sync; + + /** Ignore expired events flag. */ + @Order(13) + boolean ignoreExpired; + + /** Task name hash code. */ + @Order(14) + int taskHash; + + /** */ + @Order(15) + boolean keepBinary; + + /** Event types for JCache API. */ + @Order(16) + byte types; + + /** External marshaling. */ + protected void externalMarshal(GridKernalContext ctx) throws IgniteCheckedException { + if (requiresDeployment(rmtFilter)) + rmtFilterDep = marshalDeployable(rmtFilter, ctx); + + if (requiresDeployment(rmtFilterFactory)) + rmtFilterFactoryDep = marshalDeployable(rmtFilterFactory, ctx); + + if (requiresDeployment(rmtTransFactory)) + rmtTransFactoryDep = marshalDeployable(rmtTransFactory, ctx); + } + + /** Allows to od additional work along marshaling {@code deployable}. */ + protected CacheContinuousQueryDeployableObject marshalDeployable( + Object deployable, + GridKernalContext ctx + ) throws IgniteCheckedException { + return new CacheContinuousQueryDeployableObject(deployable, ctx); + } + + /** {@inheritDoc} */ + @Override public void marshal(Marshaller marsh) throws IgniteCheckedException { + externalMarshaling = new boolean[3]; + + if (rmtFilterDep == null && rmtFilter != null) { + rmtFilterBytes = marsh.marshal(rmtFilter); + + externalMarshaling[0] = true; + } + + if (rmtFilterFactoryDep == null && rmtFilterFactory != null) { + rmtFilterFactoryBytes = marsh.marshal(rmtFilterFactory); + + externalMarshaling[1] = true; + } + + if (rmtTransFactoryDep == null && rmtTransFactory != null) { + rmtTransFactoryBytes = marsh.marshal(rmtTransFactory); + + externalMarshaling[2] = true; + } + } + + /** External unmarshaling. */ + protected void externalUnmarshal(UUID nodeId, GridKernalContext ctx) throws IgniteCheckedException { + if (rmtFilterDep != null) + rmtFilter = externalUnmarshal(rmtFilterDep, nodeId, ctx); + + if (rmtFilterFactoryDep != null) + rmtFilterFactory = externalUnmarshal(rmtFilterFactoryDep, nodeId, ctx); + + if (rmtTransFactoryDep != null) + rmtTransFactory = externalUnmarshal(rmtTransFactoryDep, nodeId, ctx); + } + + /** {@inheritDoc} */ + @Override public void unmarshal(Marshaller marsh, ClassLoader clsLdr) throws IgniteCheckedException { + assert externalMarshaling != null && externalMarshaling.length == 3; + + if (externalMarshaling[0]) { + assert rmtFilterBytes != null; + + rmtFilter = marsh.unmarshal(rmtFilterBytes, clsLdr); + } + + if (externalMarshaling[1]) { + assert rmtFilterFactoryBytes != null; + + rmtFilterFactory = marsh.unmarshal(rmtFilterFactoryBytes, clsLdr); + } + + if (externalMarshaling[2]) { + assert rmtTransFactoryBytes != null; + + rmtTransFactory = marsh.unmarshal(rmtTransFactoryBytes, clsLdr); + } + } + + /** */ + protected T externalUnmarshal( + CacheContinuousQueryDeployableObject depObj, + UUID nodeId, + GridKernalContext ctx + ) throws IgniteCheckedException { + return depObj.unmarshal(nodeId, ctx); + } + + /** */ + private static boolean requiresDeployment(@Nullable Object obj) { + return obj != null && !U.isGrid(obj.getClass()); + } +} diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryData.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryData.java new file mode 100644 index 0000000000000..d6d1340649f4b --- /dev/null +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryData.java @@ -0,0 +1,72 @@ +/* + * 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.apache.ignite.internal.processors.continuous; + +import java.util.ArrayList; +import java.util.Collection; +import java.util.Map; +import java.util.UUID; +import org.apache.ignite.internal.Order; +import org.apache.ignite.internal.util.tostring.GridToStringInclude; +import org.apache.ignite.internal.util.typedef.internal.S; +import org.apache.ignite.plugin.extensions.communication.Message; + +/** Discovery data. */ +public class ContinousRoutineDiscoveryData implements Message { + /** Node ID. */ + @Order(0) + UUID nodeId; + + /** Items. */ + @GridToStringInclude + @Order(1) + Collection items; + + /** */ + @Order(2) + Map> clientInfos; + + /** Empty constructor for serialization purposes. */ + public ContinousRoutineDiscoveryData() { + // No-op. + } + + /** + * @param nodeId Node ID. + * @param clientInfos Client information. + */ + ContinousRoutineDiscoveryData(UUID nodeId, Map> clientInfos) { + assert nodeId != null; + + this.nodeId = nodeId; + + this.clientInfos = clientInfos; + + items = new ArrayList<>(); + } + + /** @param item Item. */ + public void addItem(ContinousRoutineDiscoveryDataItem item) { + items.add(item); + } + + /** {@inheritDoc} */ + @Override public String toString() { + return S.toString(ContinousRoutineDiscoveryData.class, this); + } +} diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryDataItem.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryDataItem.java new file mode 100644 index 0000000000000..ce4ead23f101f --- /dev/null +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryDataItem.java @@ -0,0 +1,96 @@ +/* + * 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.apache.ignite.internal.processors.continuous; + +import java.util.UUID; +import org.apache.ignite.cluster.ClusterNode; +import org.apache.ignite.internal.Marshalled; +import org.apache.ignite.internal.Order; +import org.apache.ignite.internal.util.typedef.internal.S; +import org.apache.ignite.lang.IgnitePredicate; +import org.apache.ignite.plugin.extensions.communication.Message; +import org.jetbrains.annotations.Nullable; + +/** Discovery data item. */ +public class ContinousRoutineDiscoveryDataItem implements Message { + /** */ + @Order(0) + UUID routineId; + + /** */ + @Marshalled("prjPredBytes") + IgnitePredicate prjPred; + + /** Marshalled {@link #prjPred}. */ + @Order(1) + byte[] prjPredBytes; + + /** Handler. */ + @Order(2) + GridContinuousHandler hnd; + + /** Buffer size. */ + @Order(3) + int bufSize; + + /** Time interval. */ + @Order(4) + long interval; + + /** Automatic unsubscribe flag. */ + @Order(5) + boolean autoUnsubscribe; + + /** Empty constructor for serialization porposes. */ + public ContinousRoutineDiscoveryDataItem() { + // No-op. + } + + /** + * @param routineId Consume ID. + * @param prjPred Projection predicate. + * @param hnd Handler. + * @param bufSize Buffer size. + * @param interval Time interval. + * @param autoUnsubscribe Automatic unsubscribe flag. + */ + ContinousRoutineDiscoveryDataItem(UUID routineId, + @Nullable IgnitePredicate prjPred, + GridContinuousHandler hnd, + int bufSize, + long interval, + boolean autoUnsubscribe + ) { + assert routineId != null; + assert hnd != null; + assert bufSize > 0; + assert interval >= 0; + + this.routineId = routineId; + this.prjPred = prjPred; + this.hnd = hnd; + this.bufSize = bufSize; + this.interval = interval; + this.autoUnsubscribe = autoUnsubscribe; + } + + /** {@inheritDoc} */ + @Override public String toString() { + return S.toString(ContinousRoutineDiscoveryDataItem.class, this); + } +} diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineLocalInfo.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineLocalInfo.java new file mode 100644 index 0000000000000..0ce5cb077da06 --- /dev/null +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineLocalInfo.java @@ -0,0 +1,131 @@ +/* + * 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.apache.ignite.internal.processors.continuous; + +import java.util.UUID; +import org.apache.ignite.cluster.ClusterNode; +import org.apache.ignite.internal.Marshalled; +import org.apache.ignite.internal.Order; +import org.apache.ignite.internal.util.typedef.internal.S; +import org.apache.ignite.lang.IgnitePredicate; +import org.apache.ignite.plugin.extensions.communication.Message; +import org.jetbrains.annotations.Nullable; + +/** Local routine info. */ +public class ContinousRoutineLocalInfo implements Message, GridContinuousProcessor.RoutineInfo { + /** Source node id. */ + @Order(0) + UUID nodeId; + + /** Projection predicate. */ + @Marshalled("prjPredBytes") + IgnitePredicate prjPred; + + /** Marshalled {@link #prjPred}. */ + @Order(1) + byte[] prjPredBytes; + + /** Continuous routine handler. */ + @Order(2) + GridContinuousHandler hnd; + + /** Buffer size. */ + @Order(3) + int bufSize; + + /** Time interval. */ + @Order(4) + long interval; + + /** Automatic unsubscribe flag. */ + @Order(5) + boolean autoUnsubscribe; + + /** Empty constructor for serialization purposes. */ + public ContinousRoutineLocalInfo() { + // No-op. + } + + /** + * @param nodeId Node id. + * @param prjPred Projection predicate. + * @param hnd Continuous routine handler. + * @param bufSize Buffer size. + * @param interval Interval. + * @param autoUnsubscribe Automatic unsubscribe flag. + */ + ContinousRoutineLocalInfo( + UUID nodeId, + @Nullable IgnitePredicate prjPred, + GridContinuousHandler hnd, + int bufSize, + long interval, + boolean autoUnsubscribe + ) { + assert hnd != null; + assert bufSize > 0; + assert interval >= 0; + + this.nodeId = nodeId; + this.prjPred = prjPred; + this.hnd = hnd; + this.bufSize = bufSize; + this.interval = interval; + this.autoUnsubscribe = autoUnsubscribe; + } + + /** {@inheritDoc} */ + @Override public GridContinuousHandler handler() { + return hnd; + } + + /** {@inheritDoc} */ + @Override public int bufferSize() { + return bufSize; + } + + /** {@inheritDoc} */ + @Override public long interval() { + return interval; + } + + /** {@inheritDoc} */ + @Override public boolean autoUnsubscribe() { + return autoUnsubscribe; + } + + /** {@inheritDoc} */ + @Override public long lastSendTime() { + return -1; + } + + /** {@inheritDoc} */ + @Override public boolean delayedRegister() { + return false; + } + + /** {@inheritDoc} */ + @Override public UUID nodeId() { + return nodeId; + } + + /** {@inheritDoc} */ + @Override public String toString() { + return S.toString(ContinousRoutineLocalInfo.class, this); + } +} diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutineInfo.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutineInfo.java index d339acecdebef..9253f6fdb0500 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutineInfo.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutineInfo.java @@ -22,7 +22,9 @@ import org.apache.ignite.internal.util.typedef.internal.S; import org.apache.ignite.plugin.extensions.communication.Message; -/** */ +/** + * + */ public class ContinuousRoutineInfo implements Message { /** */ @Order(0) @@ -87,12 +89,16 @@ public ContinuousRoutineInfo() { this.autoUnsubscribe = autoUnsubscribe; } - /** @param srcNodeId Source node ID. */ + /** + * @param srcNodeId Source node ID. + */ void sourceNodeId(UUID srcNodeId) { this.srcNodeId = srcNodeId; } - /** */ + /** + * + */ void onDisconnected() { disconnected = true; } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesJoiningNodeDiscoveryData.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesJoiningNodeDiscoveryData.java index 8c40160ed119b..1849e472672de 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesJoiningNodeDiscoveryData.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesJoiningNodeDiscoveryData.java @@ -22,7 +22,9 @@ import org.apache.ignite.internal.util.typedef.internal.S; import org.apache.ignite.plugin.extensions.communication.Message; -/** */ +/** + * + */ public class ContinuousRoutinesJoiningNodeDiscoveryData implements Message { /** */ @Order(0) @@ -33,7 +35,9 @@ public ContinuousRoutinesJoiningNodeDiscoveryData() { // No-op. } - /** @param startedRoutines Routines registered on nodes, to be started in cluster. */ + /** + * @param startedRoutines Routines registered on nodes, to be started in cluster. + */ ContinuousRoutinesJoiningNodeDiscoveryData(List startedRoutines) { this.startedRoutines = startedRoutines; } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousHandler.java index b1a3812f61217..03ca0defa6d37 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousHandler.java @@ -17,20 +17,20 @@ package org.apache.ignite.internal.processors.continuous; -import java.io.Externalizable; import java.util.Collection; import java.util.Map; import java.util.UUID; import org.apache.ignite.IgniteCheckedException; import org.apache.ignite.internal.GridKernalContext; import org.apache.ignite.internal.processors.affinity.AffinityTopologyVersion; +import org.apache.ignite.plugin.extensions.communication.Message; import org.jetbrains.annotations.Nullable; /** * Continuous routine handler. */ @SuppressWarnings("PublicInnerClass") -public interface GridContinuousHandler extends Externalizable, Cloneable { +public interface GridContinuousHandler extends Cloneable, Message { /** * Listener registration status. */ diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java index 16e7dba962d2d..773df58673eca 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java @@ -17,11 +17,6 @@ package org.apache.ignite.internal.processors.continuous; -import java.io.Externalizable; -import java.io.IOException; -import java.io.ObjectInput; -import java.io.ObjectOutput; -import java.util.ArrayList; import java.util.Collection; import java.util.Collections; import java.util.HashMap; @@ -38,6 +33,7 @@ import java.util.concurrent.locks.ReadWriteLock; import java.util.concurrent.locks.ReentrantLock; import java.util.concurrent.locks.ReentrantReadWriteLock; +import java.util.stream.Stream; import javax.cache.event.CacheEntryUpdatedListener; import org.apache.ignite.IgniteCheckedException; import org.apache.ignite.IgniteException; @@ -54,7 +50,6 @@ import org.apache.ignite.internal.IgniteInternalFuture; import org.apache.ignite.internal.IgniteInterruptedCheckedException; import org.apache.ignite.internal.NodeStoppingException; -import org.apache.ignite.internal.Order; import org.apache.ignite.internal.cluster.ClusterTopologyCheckedException; import org.apache.ignite.internal.managers.communication.ErrorMessage; import org.apache.ignite.internal.managers.communication.GridMessageListener; @@ -80,7 +75,6 @@ import org.apache.ignite.internal.util.future.GridFutureAdapter; import org.apache.ignite.internal.util.lang.GridPlainRunnable; import org.apache.ignite.internal.util.lang.gridfunc.ReadOnlyCollectionView2X; -import org.apache.ignite.internal.util.tostring.GridToStringInclude; import org.apache.ignite.internal.util.typedef.CI1; import org.apache.ignite.internal.util.typedef.F; import org.apache.ignite.internal.util.typedef.X; @@ -125,10 +119,10 @@ public class GridContinuousProcessor extends GridProcessorAdapter { public static final String CQ_SYS_VIEW_DESC = "Continuous queries"; /** Local infos. */ - private final ConcurrentMap locInfos = new ConcurrentHashMap<>(); + private final ConcurrentMap locInfos = new ConcurrentHashMap<>(); /** Local infos. */ - private final ConcurrentMap> clientInfos = new ConcurrentHashMap<>(); + private final ConcurrentMap> clientInfos = new ConcurrentHashMap<>(); /** Remote infos. */ private final ConcurrentMap rmtInfos = new ConcurrentHashMap<>(); @@ -336,12 +330,12 @@ Map remoteRoutineInfos() { } /** */ - Map localRoutineInfos() { + Map localRoutineInfos() { return Collections.unmodifiableMap(locInfos); } /** */ - Map> clientRoutineInfos() { + Map> clientRoutineInfos() { return Collections.unmodifiableMap(clientInfos); } @@ -408,7 +402,7 @@ public void unlockStopping() { return; } - var data = getDiscoveryData(dataBag.joiningNodeId()); + Message data = getDiscoveryData(dataBag.joiningNodeId()); if (data != null) dataBag.addJoiningNodeData(CONTINUOUS_PROC.ordinal(), data); @@ -431,44 +425,42 @@ public void unlockStopping() { /** * @param joiningNodeId Joining node id. */ - private @Nullable DiscoveryData getDiscoveryData(UUID joiningNodeId) { + private @Nullable ContinousRoutineDiscoveryData getDiscoveryData(UUID joiningNodeId) { if (log.isDebugEnabled()) { log.debug("collectDiscoveryData [node=" + joiningNodeId + - ", loc=" + ctx.localNodeId() + - ", locInfos=" + locInfos + - ", clientInfos=" + clientInfos + - ']' - ); + ", loc=" + ctx.localNodeId() + + ", locInfos=" + locInfos + + ", clientInfos=" + clientInfos + + ']'); } if (!joiningNodeId.equals(ctx.localNodeId()) || !locInfos.isEmpty()) { - Map> clientInfos0 = copyClientInfos(clientInfos); + Map> clientInfos0 = copyClientInfos(clientInfos); if (joiningNodeId.equals(ctx.localNodeId()) && ctx.discovery().localNode().isClient()) { - Map infos = copyLocalInfos(locInfos); + Map infos = copyLocalInfos(locInfos); clientInfos0.put(ctx.localNodeId(), infos); } - DiscoveryData data = new DiscoveryData(ctx.localNodeId(), clientInfos0); + ContinousRoutineDiscoveryData data = new ContinousRoutineDiscoveryData(ctx.localNodeId(), clientInfos0); // Collect listeners information (will be sent to joining node during discovery process). - for (Map.Entry e : locInfos.entrySet()) { + for (Map.Entry e : locInfos.entrySet()) { UUID routineId = e.getKey(); - LocalRoutineInfo info = e.getValue(); + ContinousRoutineLocalInfo info = e.getValue(); assert !ctx.config().isPeerClassLoadingEnabled() || !(info.hnd instanceof CacheContinuousQueryHandler) || - ((CacheContinuousQueryHandler)info.hnd).isMarshalled(); + (((CacheContinuousQueryHandler)info.hnd).externalMarshaling != null + && Stream.of(((CacheContinuousQueryHandler)info.hnd).externalMarshaling).findAny().isPresent()); - data.addItem(new DiscoveryDataItem( - routineId, + data.addItem(new ContinousRoutineDiscoveryDataItem(routineId, info.prjPred, info.hnd, info.bufSize, info.interval, - info.autoUnsubscribe - )); + info.autoUnsubscribe)); } return data; @@ -480,13 +472,13 @@ public void unlockStopping() { /** * @param clientInfos Client infos. */ - private Map> copyClientInfos(Map> clientInfos) { - Map> res = U.newHashMap(clientInfos.size()); + private Map> copyClientInfos(Map> clientInfos) { + Map> res = U.newHashMap(clientInfos.size()); - for (Map.Entry> e : clientInfos.entrySet()) { - Map cp = U.newHashMap(e.getValue().size()); + for (Map.Entry> e : clientInfos.entrySet()) { + Map cp = U.newHashMap(e.getValue().size()); - for (Map.Entry e0 : e.getValue().entrySet()) + for (Map.Entry e0 : e.getValue().entrySet()) cp.put(e0.getKey(), e0.getValue()); res.put(e.getKey(), cp); @@ -498,10 +490,10 @@ private Map> copyClientInfos(Map copyLocalInfos(Map locInfos) { - Map res = U.newHashMap(locInfos.size()); + private Map copyLocalInfos(Map locInfos) { + Map res = U.newHashMap(locInfos.size()); - for (Map.Entry e : locInfos.entrySet()) + for (Map.Entry e : locInfos.entrySet()) res.put(e.getKey(), e.getValue()); return res; @@ -550,10 +542,10 @@ private Map copyLocalInfos(Map l } } else { - Map nodeSpecData = data.nodeSpecificData(); + Map nodeSpecData = data.nodeSpecificData(); if (nodeSpecData != null) { - for (DiscoveryData val : nodeSpecData.values()) + for (ContinousRoutineDiscoveryData val : nodeSpecData.values()) onDiscoveryDataReceivedMutable(val); } } @@ -565,37 +557,37 @@ private Map copyLocalInfos(Map l * * @param data received discovery data. */ - private void onDiscoveryDataReceivedMutable(DiscoveryData data) { + private void onDiscoveryDataReceivedMutable(ContinousRoutineDiscoveryData data) { if (data != null) { - for (DiscoveryDataItem item : data.items) { + for (ContinousRoutineDiscoveryDataItem item : data.items) { if (!locInfos.containsKey(item.routineId)) { registerHandlerOnJoin(data.nodeId, item.routineId, item.prjPred, item.hnd, item.bufSize, item.interval, item.autoUnsubscribe); } if (!item.autoUnsubscribe) { - locInfos.putIfAbsent(item.routineId, new LocalRoutineInfo(data.nodeId, + locInfos.putIfAbsent(item.routineId, new ContinousRoutineLocalInfo(data.nodeId, item.prjPred, item.hnd, item.bufSize, item.interval, item.autoUnsubscribe)); } } // Process CQs started on clients. - for (Map.Entry> entry : data.clientInfos.entrySet()) { + for (Map.Entry> entry : data.clientInfos.entrySet()) { UUID clientNodeId = entry.getKey(); if (!ctx.localNodeId().equals(clientNodeId)) { - Map clientRoutineMap = entry.getValue(); + Map clientRoutineMap = entry.getValue(); - for (Map.Entry e : clientRoutineMap.entrySet()) { + for (Map.Entry e : clientRoutineMap.entrySet()) { UUID routineId = e.getKey(); - LocalRoutineInfo info = e.getValue(); + ContinousRoutineLocalInfo info = e.getValue(); registerHandlerOnJoin(clientNodeId, routineId, info.prjPred, info.hnd, info.bufSize, info.interval, info.autoUnsubscribe); } } - Map map = + Map map = clientInfos.computeIfAbsent(clientNodeId, k -> new HashMap<>()); map.putAll(entry.getValue()); @@ -785,7 +777,7 @@ public UUID registerStaticRoutine( final UUID routineId = UUID.randomUUID(); - LocalRoutineInfo routineInfo = new LocalRoutineInfo(ctx.localNodeId(), prjPred, hnd, 1, 0, true); + ContinousRoutineLocalInfo routineInfo = new ContinousRoutineLocalInfo(ctx.localNodeId(), prjPred, hnd, 1, 0, true); if (immutableDiscoCustomMsg) { routinesInfo.addRoutineInfo(createRoutineInfo( @@ -864,12 +856,13 @@ public IgniteInternalFuture startRoutine(GridContinuousHandler hnd, if (ctx.config().isPeerClassLoadingEnabled()) { hnd.p2pMarshal(ctx); - assert !(hnd instanceof CacheContinuousQueryHandler) || ((CacheContinuousQueryHandler)hnd).isMarshalled(); + assert !(hnd instanceof CacheContinuousQueryHandler) || (((CacheContinuousQueryHandler)hnd).externalMarshaling != null + && Stream.of(((CacheContinuousQueryHandler)hnd).externalMarshaling).findAny().isPresent()); } // Register routine locally. locInfos.put(routineId, - new LocalRoutineInfo(ctx.localNodeId(), prjPred, hnd, bufSize, interval, autoUnsubscribe)); + new ContinousRoutineLocalInfo(ctx.localNodeId(), prjPred, hnd, bufSize, interval, autoUnsubscribe)); if (locOnly) { try { @@ -1061,7 +1054,7 @@ public IgniteInternalFuture stopRoutine(UUID routineId) { boolean stop = false; // Unregister routine locally. - LocalRoutineInfo routine = locInfos.remove(routineId); + ContinousRoutineLocalInfo routine = locInfos.remove(routineId); if (routine != null) { stop = true; @@ -1126,7 +1119,7 @@ public void addBackupNotification(UUID nodeId, sendNotification(nodeId, routineId, null, toSnd, orderedTopic, true, null); } else { - LocalRoutineInfo locRoutineInfo = locInfos.get(routineId); + ContinousRoutineLocalInfo locRoutineInfo = locInfos.get(routineId); if (locRoutineInfo != null) locRoutineInfo.handler().notifyCallback(nodeId, routineId, objs, ctx); @@ -1251,7 +1244,7 @@ public void addNotification(UUID nodeId, unregisterRemote(e.getKey()); } - for (LocalRoutineInfo routine : locInfos.values()) + for (ContinousRoutineLocalInfo routine : locInfos.values()) routine.hnd.onClientDisconnected(); rmtInfos.clear(); @@ -1318,7 +1311,7 @@ private void processStopRequest(ClusterNode snd, StopRoutineDiscoveryMessage msg unregisterRemote(routineId); } - for (Map clientInfo : clientInfos.values()) { + for (Map clientInfo : clientInfos.values()) { if (clientInfo.remove(msg.routineId()) != null) break; } @@ -1377,17 +1370,17 @@ private void processStartRequestMutable(ClusterNode node, StartRoutineDiscoveryM GridContinuousHandler hnd = data.handler(); if (node.isClient()) { - Map clientRoutineMap = clientInfos.get(node.id()); + Map clientRoutineMap = clientInfos.get(node.id()); if (clientRoutineMap == null) { clientRoutineMap = new HashMap<>(); - Map old = clientInfos.put(node.id(), clientRoutineMap); + Map old = clientInfos.put(node.id(), clientRoutineMap); assert old == null; } - clientRoutineMap.put(routineId, new LocalRoutineInfo(node.id(), + clientRoutineMap.put(routineId, new ContinousRoutineLocalInfo(node.id(), data.nodeFilter(), hnd, data.bufferSize(), @@ -1425,7 +1418,7 @@ private void processStartRequestMutable(ClusterNode node, StartRoutineDiscoveryM if (!data.autoUnsubscribe()) // Register routine locally. - locInfos.putIfAbsent(routineId, new LocalRoutineInfo( + locInfos.putIfAbsent(routineId, new ContinousRoutineLocalInfo( node.id(), prjPred, hnd, data.bufferSize(), data.interval(), data.autoUnsubscribe())); } catch (IgniteCheckedException e) { @@ -1610,7 +1603,7 @@ private void processNotification(UUID nodeId, GridContinuousMessage msg) { UUID routineId = msg.routineId(); try { - LocalRoutineInfo routine = locInfos.get(routineId); + ContinousRoutineLocalInfo routine = locInfos.get(routineId); if (routine != null) routine.hnd.notifyCallback(nodeId, routineId, (Collection)msg.data(), ctx); @@ -1772,7 +1765,7 @@ private void unregisterHandler(UUID routineId, GridContinuousHandler hnd, boolea @SuppressWarnings("TooBroadScope") private void unregisterRemote(UUID routineId) { RemoteRoutineInfo remote; - LocalRoutineInfo loc; + ContinousRoutineLocalInfo loc; stopLock.lock(); @@ -1974,101 +1967,6 @@ public static interface RoutineInfo { boolean delayedRegister(); } - /** Local routine info. */ - public static class LocalRoutineInfo implements Message, RoutineInfo { - /** Source node id. */ - @Order(0) - UUID nodeId; - - /** Projection predicate. */ - @Order(1) - IgnitePredicate prjPred; - - /** Continuous routine handler. */ - @Order(2) - GridContinuousHandler hnd; - - /** Buffer size. */ - @Order(3) - int bufSize; - - /** Time interval. */ - @Order(4) - long interval; - - /** Automatic unsubscribe flag. */ - @Order(5) - boolean autoUnsubscribe; - - /** - * @param nodeId Node id. - * @param prjPred Projection predicate. - * @param hnd Continuous routine handler. - * @param bufSize Buffer size. - * @param interval Interval. - * @param autoUnsubscribe Automatic unsubscribe flag. - */ - LocalRoutineInfo( - UUID nodeId, - @Nullable IgnitePredicate prjPred, - GridContinuousHandler hnd, - int bufSize, - long interval, - boolean autoUnsubscribe - ) { - assert hnd != null; - assert bufSize > 0; - assert interval >= 0; - - this.nodeId = nodeId; - this.prjPred = prjPred; - this.hnd = hnd; - this.bufSize = bufSize; - this.interval = interval; - this.autoUnsubscribe = autoUnsubscribe; - } - - /** {@inheritDoc} */ - @Override public GridContinuousHandler handler() { - return hnd; - } - - /** {@inheritDoc} */ - @Override public int bufferSize() { - return bufSize; - } - - /** {@inheritDoc} */ - @Override public long interval() { - return interval; - } - - /** {@inheritDoc} */ - @Override public boolean autoUnsubscribe() { - return autoUnsubscribe; - } - - /** {@inheritDoc} */ - @Override public long lastSendTime() { - return -1; - } - - /** {@inheritDoc} */ - @Override public boolean delayedRegister() { - return false; - } - - /** {@inheritDoc} */ - @Override public UUID nodeId() { - return nodeId; - } - - /** {@inheritDoc} */ - @Override public String toString() { - return S.toString(LocalRoutineInfo.class, this); - } - } - /** * Remote routine info. */ @@ -2290,139 +2188,6 @@ IgniteBiTuple checkInterval() { } } - /** Discovery data. */ - public static class DiscoveryData implements Message { - /** Node ID. */ - @Order(0) - UUID nodeId; - - /** Items. */ - @GridToStringInclude - @Order(1) - Collection items; - - /** */ - @Order(2) - Map> clientInfos; - - /** Empty constructor for serialization purposes. */ - public DiscoveryData() { - // No-op. - } - - /** - * @param nodeId Node ID. - * @param clientInfos Client information. - */ - DiscoveryData(UUID nodeId, Map> clientInfos) { - assert nodeId != null; - - this.nodeId = nodeId; - - this.clientInfos = clientInfos; - - items = new ArrayList<>(); - } - - /** - * @param item Item. - */ - public void addItem(DiscoveryDataItem item) { - items.add(item); - } - - /** {@inheritDoc} */ - @Override public String toString() { - return S.toString(DiscoveryData.class, this); - } - } - - /** - * Discovery data item. - */ - private static class DiscoveryDataItem implements Externalizable { - /** */ - private static final long serialVersionUID = 0L; - - /** Consume ID. */ - private UUID routineId; - - /** Projection predicate. */ - private IgnitePredicate prjPred; - - /** Handler. */ - private GridContinuousHandler hnd; - - /** Buffer size. */ - private int bufSize; - - /** Time interval. */ - private long interval; - - /** Automatic unsubscribe flag. */ - private boolean autoUnsubscribe; - - /** - * Required by {@link Externalizable}. - */ - public DiscoveryDataItem() { - // No-op. - } - - /** - * @param routineId Consume ID. - * @param prjPred Projection predicate. - * @param hnd Handler. - * @param bufSize Buffer size. - * @param interval Time interval. - * @param autoUnsubscribe Automatic unsubscribe flag. - */ - DiscoveryDataItem(UUID routineId, - @Nullable IgnitePredicate prjPred, - GridContinuousHandler hnd, - int bufSize, - long interval, - boolean autoUnsubscribe - ) { - assert routineId != null; - assert hnd != null; - assert bufSize > 0; - assert interval >= 0; - - this.routineId = routineId; - this.prjPred = prjPred; - this.hnd = hnd; - this.bufSize = bufSize; - this.interval = interval; - this.autoUnsubscribe = autoUnsubscribe; - } - - /** {@inheritDoc} */ - @Override public void writeExternal(ObjectOutput out) throws IOException { - U.writeUuid(out, routineId); - out.writeObject(prjPred); - out.writeObject(hnd); - out.writeInt(bufSize); - out.writeLong(interval); - out.writeBoolean(autoUnsubscribe); - } - - /** {@inheritDoc} */ - @Override public void readExternal(ObjectInput in) throws IOException, ClassNotFoundException { - routineId = U.readUuid(in); - prjPred = (IgnitePredicate)in.readObject(); - hnd = (GridContinuousHandler)in.readObject(); - bufSize = in.readInt(); - interval = in.readLong(); - autoUnsubscribe = in.readBoolean(); - } - - /** {@inheritDoc} */ - @Override public String toString() { - return S.toString(DiscoveryDataItem.class, this); - } - } - /** * Future for start routine. */ @@ -2507,7 +2272,7 @@ private void onAllRemoteRegistered( Map cntrs) { try { if (errs == null || errs.isEmpty()) { - LocalRoutineInfo routine = locInfos.get(routineId); + ContinousRoutineLocalInfo routine = locInfos.get(routineId); // Update partition counters. if (routine != null && routine.handler().isQuery()) { diff --git a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryEntriesExpireTest.java b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryEntriesExpireTest.java index 736b1742befa5..317076bc92a47 100644 --- a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryEntriesExpireTest.java +++ b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryEntriesExpireTest.java @@ -28,7 +28,7 @@ import org.apache.ignite.configuration.CacheConfiguration; import org.apache.ignite.internal.IgniteEx; import org.apache.ignite.internal.processors.cache.GridCacheContext; -import org.apache.ignite.internal.processors.continuous.GridContinuousProcessor; +import org.apache.ignite.internal.processors.continuous.ContinousRoutineLocalInfo; import org.apache.ignite.internal.util.typedef.F; import org.apache.ignite.internal.util.typedef.internal.CU; import org.apache.ignite.testframework.GridTestUtils; @@ -57,7 +57,7 @@ public void testBackupQOnEntriesExpire() throws Exception { for (int i = 0; i < 1_000; i++) cache.put(i, i); - ConcurrentMap locInfos = + ConcurrentMap locInfos = getFieldValue(srv1.context().continuous(), "locInfos"); assertEquals(1, locInfos.size()); diff --git a/modules/core/src/test/java/org/apache/ignite/internal/processors/continuous/GridEventConsumeSelfTest.java b/modules/core/src/test/java/org/apache/ignite/internal/processors/continuous/GridEventConsumeSelfTest.java index 47b107edeee24..96e6e1d50cb12 100644 --- a/modules/core/src/test/java/org/apache/ignite/internal/processors/continuous/GridEventConsumeSelfTest.java +++ b/modules/core/src/test/java/org/apache/ignite/internal/processors/continuous/GridEventConsumeSelfTest.java @@ -63,7 +63,6 @@ import static org.apache.ignite.events.EventType.EVT_JOB_STARTED; import static org.apache.ignite.events.EventType.EVT_NODE_FAILED; import static org.apache.ignite.events.EventType.EVT_NODE_LEFT; -import static org.apache.ignite.internal.processors.continuous.GridContinuousProcessor.LocalRoutineInfo; import static org.apache.ignite.testframework.GridTestUtils.noop; /** @@ -157,13 +156,15 @@ public class GridEventConsumeSelfTest extends GridCommonAbstractTest { * @param proc Continuous processor. * @return Local event routines. */ - private Collection localRoutines(GridContinuousProcessor proc) { - return F.view(U.>field(proc, "locInfos").values(), - new IgnitePredicate() { - @Override public boolean apply(LocalRoutineInfo info) { + private Collection localRoutines(GridContinuousProcessor proc) { + return F.view( + U.>field(proc, "locInfos").values(), + new IgnitePredicate<>() { + @Override public boolean apply(ContinousRoutineLocalInfo info) { return info.handler().isEvents(); } - }); + } + ); } /** From aec0349d7764890100f3b7184077578a73899a52 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Fri, 31 Jul 2026 14:29:16 +0300 Subject: [PATCH 03/24] fix --- .../internal/GridEventConsumeHandler.java | 8 +++--- .../CacheContinuousQueryHandler.java | 26 ++++++++++++++----- .../CacheContinuousQueryHandlerMessage.java | 14 +++++----- .../continuous/GridContinuousProcessor.java | 7 ++--- 4 files changed, 32 insertions(+), 23 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java index 485b22d50e9d3..aaab1dc2db066 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java @@ -63,7 +63,7 @@ */ public final class GridEventConsumeHandler implements GridContinuousHandler, MarshallableMessage { /** Default callback. */ - private static final IgniteBiPredicate DFLT_CALLBACK = new P2() { + private static final IgniteBiPredicate DFLT_CALLBACK = new P2<>() { @Override public boolean apply(UUID uuid, Event e) { return true; } @@ -75,7 +75,7 @@ public final class GridEventConsumeHandler implements GridContinuousHandler, Mar /** Filter. */ @Nullable IgnitePredicate filter; - /** Serialized {@link #filter}. */ + /** Marshaled {@link #filter}. */ @Order(0) @Nullable byte[] filterBytes; @@ -104,10 +104,10 @@ public final class GridEventConsumeHandler implements GridContinuousHandler, Mar private GridLocalEventListener lsnr; /** P2P unmarshalling future. */ - private IgniteInternalFuture p2pUnmarshalFut = new GridFinishedFuture<>(); + private final IgniteInternalFuture p2pUnmarshalFut = new GridFinishedFuture<>(); /** - * Required by {@link Externalizable}. + * Empty constructor for serialization purposes. */ public GridEventConsumeHandler() { // No-op. diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java index 1620a145f1bd0..6437f66e5851a 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java @@ -890,8 +890,6 @@ private void prepareEntry(GridCacheContext cctx, UUID nodeId, CacheContinuousQue * @throws IgniteCheckedException In case of error. */ void waitTopologyFuture(GridKernalContext ctx) throws IgniteCheckedException { - GridCacheContext cctx = cacheContext(ctx); - AffinityTopologyVersion topVer = initTopVer; cacheContext(ctx).shared().exchange().affinityReadyFuture(topVer).get(); @@ -1339,10 +1337,15 @@ CacheContinuousQueryEventBuffer partitionBuffer(GridCacheContext cctx, int @Override public void p2pMarshal(GridKernalContext ctx) throws IgniteCheckedException { assert ctx.config().isPeerClassLoadingEnabled(); - externalMarshal(ctx); + marshalExternally(ctx); } - /** Processes {@link #p2pUnmarshalFut} over the super's method. */ + /** + * Processes {@link #p2pUnmarshalFut} over the super's method. + * + * @param deployable Deployable object. + * @param ctx Kernal context. + */ @Override protected CacheContinuousQueryDeployableObject marshalDeployable( Object deployable, GridKernalContext ctx @@ -1361,20 +1364,20 @@ else if (p2pUnmarshalFut.isDone()) @Override public void p2pUnmarshal(UUID nodeId, GridKernalContext ctx) throws IgniteCheckedException { assert ctx.config().isPeerClassLoadingEnabled(); - externalUnmarshal(nodeId, ctx); + unmarshalExternally(nodeId, ctx); if (!p2pUnmarshalFut.isDone()) ((GridFutureAdapter)p2pUnmarshalFut).onDone(); } /**{@inheritDoc} */ - @Override protected T externalUnmarshal( + @Override protected T unmarshalExternally( CacheContinuousQueryDeployableObject depObj, UUID nodeId, GridKernalContext ctx ) throws IgniteCheckedException { try { - return super.externalUnmarshal(depObj, nodeId, ctx); + return super.unmarshalExternally(depObj, nodeId, ctx); } catch (IgniteCheckedException e) { ((GridFutureAdapter)p2pUnmarshalFut).onDone(e); @@ -1405,6 +1408,15 @@ else if (p2pUnmarshalFut.isDone()) cacheId = CU.cacheId(cacheName); } + /** + * @return Whether the handler is marshalled for peer class loading. + */ + public boolean isMarshalled() { + return (!requiresDeployment(rmtFilter) || rmtFilterDep != null) + && (!requiresDeployment(rmtFilterFactory) || rmtFilterFactoryDep != null) + && (!requiresDeployment(rmtTransFactory) || rmtTransFactoryDep != null); + } + /** {@inheritDoc} */ @Override public GridContinuousBatch createBatch() { return new GridContinuousQueryBatch(); diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandlerMessage.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandlerMessage.java index 340d8bbef1ba2..8aedf15c10ec0 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandlerMessage.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandlerMessage.java @@ -116,7 +116,7 @@ public class CacheContinuousQueryHandlerMessage implements MarshallableMes byte types; /** External marshaling. */ - protected void externalMarshal(GridKernalContext ctx) throws IgniteCheckedException { + protected void marshalExternally(GridKernalContext ctx) throws IgniteCheckedException { if (requiresDeployment(rmtFilter)) rmtFilterDep = marshalDeployable(rmtFilter, ctx); @@ -159,15 +159,15 @@ protected CacheContinuousQueryDeployableObject marshalDeployable( } /** External unmarshaling. */ - protected void externalUnmarshal(UUID nodeId, GridKernalContext ctx) throws IgniteCheckedException { + protected void unmarshalExternally(UUID nodeId, GridKernalContext ctx) throws IgniteCheckedException { if (rmtFilterDep != null) - rmtFilter = externalUnmarshal(rmtFilterDep, nodeId, ctx); + rmtFilter = unmarshalExternally(rmtFilterDep, nodeId, ctx); if (rmtFilterFactoryDep != null) - rmtFilterFactory = externalUnmarshal(rmtFilterFactoryDep, nodeId, ctx); + rmtFilterFactory = unmarshalExternally(rmtFilterFactoryDep, nodeId, ctx); if (rmtTransFactoryDep != null) - rmtTransFactory = externalUnmarshal(rmtTransFactoryDep, nodeId, ctx); + rmtTransFactory = unmarshalExternally(rmtTransFactoryDep, nodeId, ctx); } /** {@inheritDoc} */ @@ -194,7 +194,7 @@ protected void externalUnmarshal(UUID nodeId, GridKernalContext ctx) throws Igni } /** */ - protected T externalUnmarshal( + protected T unmarshalExternally( CacheContinuousQueryDeployableObject depObj, UUID nodeId, GridKernalContext ctx @@ -203,7 +203,7 @@ protected T externalUnmarshal( } /** */ - private static boolean requiresDeployment(@Nullable Object obj) { + protected static boolean requiresDeployment(@Nullable Object obj) { return obj != null && !U.isGrid(obj.getClass()); } } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java index 773df58673eca..cdedd8cdf1c40 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java @@ -33,7 +33,6 @@ import java.util.concurrent.locks.ReadWriteLock; import java.util.concurrent.locks.ReentrantLock; import java.util.concurrent.locks.ReentrantReadWriteLock; -import java.util.stream.Stream; import javax.cache.event.CacheEntryUpdatedListener; import org.apache.ignite.IgniteCheckedException; import org.apache.ignite.IgniteException; @@ -452,8 +451,7 @@ public void unlockStopping() { assert !ctx.config().isPeerClassLoadingEnabled() || !(info.hnd instanceof CacheContinuousQueryHandler) || - (((CacheContinuousQueryHandler)info.hnd).externalMarshaling != null - && Stream.of(((CacheContinuousQueryHandler)info.hnd).externalMarshaling).findAny().isPresent()); + ((CacheContinuousQueryHandler)info.hnd).isMarshalled(); data.addItem(new ContinousRoutineDiscoveryDataItem(routineId, info.prjPred, @@ -856,8 +854,7 @@ public IgniteInternalFuture startRoutine(GridContinuousHandler hnd, if (ctx.config().isPeerClassLoadingEnabled()) { hnd.p2pMarshal(ctx); - assert !(hnd instanceof CacheContinuousQueryHandler) || (((CacheContinuousQueryHandler)hnd).externalMarshaling != null - && Stream.of(((CacheContinuousQueryHandler)hnd).externalMarshaling).findAny().isPresent()); + assert !(hnd instanceof CacheContinuousQueryHandler) || ((CacheContinuousQueryHandler)hnd).isMarshalled(); } // Register routine locally. From f1151845a0dabc59d359300a9fbb1d7d0b922856 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Fri, 31 Jul 2026 15:13:23 +0300 Subject: [PATCH 04/24] + master --- .../org/apache/ignite/internal/GridEventConsumeHandler.java | 1 + .../org/apache/ignite/internal/GridMessageListenHandler.java | 1 + .../query/continuous/CacheContinuousQueryHandlerMessage.java | 2 ++ .../continuous/ContinousRoutineDiscoveryDataItem.java | 2 ++ .../processors/continuous/ContinousRoutineLocalInfo.java | 2 ++ 5 files changed, 8 insertions(+) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java index aaab1dc2db066..31d37a7ac6934 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java @@ -61,6 +61,7 @@ /** * Continuous routine handler for remote event listening. */ +@UseBinaryMarshaller public final class GridEventConsumeHandler implements GridContinuousHandler, MarshallableMessage { /** Default callback. */ private static final IgniteBiPredicate DFLT_CALLBACK = new P2<>() { diff --git a/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java index 155e0eacef842..af57760df2891 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java @@ -41,6 +41,7 @@ /** * Continuous handler for message subscription. */ +@UseBinaryMarshaller public final class GridMessageListenHandler implements GridContinuousHandler, MarshallableMessage { /** */ private @Nullable Object topic; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandlerMessage.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandlerMessage.java index 8aedf15c10ec0..d14294105f753 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandlerMessage.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandlerMessage.java @@ -27,12 +27,14 @@ import org.apache.ignite.internal.MarshallableMessage; import org.apache.ignite.internal.Marshalled; import org.apache.ignite.internal.Order; +import org.apache.ignite.internal.UseBinaryMarshaller; import org.apache.ignite.internal.util.typedef.internal.U; import org.apache.ignite.lang.IgniteClosure; import org.apache.ignite.marshaller.Marshaller; import org.jetbrains.annotations.Nullable; /** */ +@UseBinaryMarshaller public class CacheContinuousQueryHandlerMessage implements MarshallableMessage { /** Remote filter. */ CacheEntryEventSerializableFilter rmtFilter; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryDataItem.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryDataItem.java index ce4ead23f101f..247c5be996cb2 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryDataItem.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryDataItem.java @@ -21,12 +21,14 @@ import org.apache.ignite.cluster.ClusterNode; import org.apache.ignite.internal.Marshalled; import org.apache.ignite.internal.Order; +import org.apache.ignite.internal.UseBinaryMarshaller; import org.apache.ignite.internal.util.typedef.internal.S; import org.apache.ignite.lang.IgnitePredicate; import org.apache.ignite.plugin.extensions.communication.Message; import org.jetbrains.annotations.Nullable; /** Discovery data item. */ +@UseBinaryMarshaller public class ContinousRoutineDiscoveryDataItem implements Message { /** */ @Order(0) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineLocalInfo.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineLocalInfo.java index 0ce5cb077da06..900600c54954b 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineLocalInfo.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineLocalInfo.java @@ -21,12 +21,14 @@ import org.apache.ignite.cluster.ClusterNode; import org.apache.ignite.internal.Marshalled; import org.apache.ignite.internal.Order; +import org.apache.ignite.internal.UseBinaryMarshaller; import org.apache.ignite.internal.util.typedef.internal.S; import org.apache.ignite.lang.IgnitePredicate; import org.apache.ignite.plugin.extensions.communication.Message; import org.jetbrains.annotations.Nullable; /** Local routine info. */ +@UseBinaryMarshaller public class ContinousRoutineLocalInfo implements Message, GridContinuousProcessor.RoutineInfo { /** Source node id. */ @Order(0) From 0587132d024e76799a6e42df77b0de44caf12d87 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Fri, 31 Jul 2026 15:27:36 +0300 Subject: [PATCH 05/24] fix --- .../internal/GridMessageListenHandler.java | 17 ++++++++++------- 1 file changed, 10 insertions(+), 7 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java index af57760df2891..51dd2d007d38b 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java @@ -74,7 +74,7 @@ public final class GridMessageListenHandler implements GridContinuousHandler, Ma boolean externalMarshal; /** P2P unmarshalling future. */ - private final IgniteInternalFuture p2pUnmarshalFut = new GridFinishedFuture<>(); + private IgniteInternalFuture p2pUnmarshalFut = new GridFinishedFuture<>(); /** * Empty constructor for serialization purposes @@ -206,13 +206,13 @@ public GridMessageListenHandler(@Nullable Object topic, IgniteBiPredicate(); + return; + } topic = marsh.unmarshal(topicBytes, clsLdr); pred = marsh.unmarshal(predBytes, clsLdr); From a50ab96f309132daaffd049fc217b5e72ff5618f Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Fri, 31 Jul 2026 18:19:36 +0300 Subject: [PATCH 06/24] in-progress --- .../internal/GridMessageListenHandler.java | 72 +++++++++++-------- .../continuous/ContinuousRoutineInfo.java | 6 +- .../continuous/GridContinuousHandler.java | 5 ++ .../continuous/GridContinuousProcessor.java | 32 ++------- .../continuous/StartRequestData.java | 21 ++---- 5 files changed, 60 insertions(+), 76 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java index 51dd2d007d38b..ae8e9254b7660 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java @@ -44,26 +44,26 @@ @UseBinaryMarshaller public final class GridMessageListenHandler implements GridContinuousHandler, MarshallableMessage { /** */ - private @Nullable Object topic; + private volatile @Nullable Object topic; /** Marshalled {@link #topic}. */ @Order(0) - @Nullable byte[] topicBytes; + volatile @Nullable byte[] topicBytes; /** */ - private IgniteBiPredicate pred; + private volatile IgniteBiPredicate pred; /** Marshalled {@link #pred}. */ @Order(1) - byte[] predBytes; + volatile byte[] predBytes; - /** Is {@code null} if the P2P deployment is disabled. */ + /** Class name of {@link #pred}. Is {@code null} if the P2P deployment is disabled. */ @Order(2) - @Nullable String clsName; + volatile @Nullable String clsName; - /** Is {@code null} if the P2P deployment is disabled. */ + /** P2P deploy info of {@link #pred}. Is {@code null} if the P2P deployment is disabled. */ @Order(3) - @Nullable GridDeploymentInfoBean depInfo; + volatile @Nullable GridDeploymentInfoBean predDepInfo; /** * Lever of the own marshaling. @@ -71,7 +71,7 @@ public final class GridMessageListenHandler implements GridContinuousHandler, Ma * @see #p2pMarshal(GridKernalContext) */ @Order(4) - boolean externalMarshal; + volatile boolean externalMarshal; /** P2P unmarshalling future. */ private IgniteInternalFuture p2pUnmarshalFut = new GridFinishedFuture<>(); @@ -158,6 +158,8 @@ public GridMessageListenHandler(@Nullable Object topic, IgniteBiPredicate Date: Mon, 3 Aug 2026 16:45:18 +0300 Subject: [PATCH 08/24] fix --- .../org/apache/ignite/internal/GridEventConsumeHandler.java | 6 ------ 1 file changed, 6 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java index 1e8f9fc0c8c64..d0c79bc093b72 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java @@ -402,8 +402,6 @@ private boolean filterDropsEvent(Event evt) { assert ctx.config().isPeerClassLoadingEnabled(); if (filter != null) { - assert filterBytes == null : "Duplicated p2p-marshalling, " + getClass().getSimpleName(); - Class cls = U.detectClass(filter); clsName = cls.getName(); @@ -426,12 +424,8 @@ private boolean filterDropsEvent(Event evt) { assert nodeId != null; assert ctx.config().isPeerClassLoadingEnabled(); assert externalMarshal : "Is not p2p-marshaled " + getClass().getSimpleName(); - assert !p2pUnmarshalFut.isDone() && p2pUnmarshalFut instanceof GridFutureAdapter : - "Can't p2p-unmarshal, the p2p-umarshalling future seems to be already done, " + getClass().getSimpleName(); if (filterBytes != null) { - assert filter == null : "Duplicated p2p-unmarshalling, " + getClass().getSimpleName(); - try { GridDeployment dep = ctx.deploy().getGlobalDeployment(depInfo.deployMode(), clsName, clsName, depInfo.userVersion(), nodeId, depInfo.classLoaderId(), depInfo.participants(), null); From fa3f2c4b17f399221147446ee14cd83e08b3aa5f Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Mon, 3 Aug 2026 17:59:58 +0300 Subject: [PATCH 09/24] fix --- .../CacheContinuousQueryHandler.java | 99 ++++++------------- 1 file changed, 28 insertions(+), 71 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java index 901944b2f577c..27e000707d205 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java @@ -144,42 +144,38 @@ public final class CacheContinuousQueryHandler implements GridContinuousHa /** Remote filter. */ volatile CacheEntryEventSerializableFilter rmtFilter; - /** Lever of own marshaling. */ - @Order(0) - protected volatile boolean[] externalMarshal; - /** Deployable object for {@link #rmtFilter}. Is {@code null} if no external marsshalling used. */ - @Order(1) + @Order(0) @Nullable volatile CacheContinuousQueryDeployableObject rmtFilterDep; /** Marshalled {@link #rmtFilter} if {@link #rmtFilterDep} is {@code null}. */ - @Order(2) + @Order(1) @Nullable volatile byte[] rmtFilterBytes; /** Remote filter factory. */ @Nullable volatile Factory rmtFilterFactory; /** Deployable object for {@link #rmtFilterFactory}. Is {@code null} if no external marsshalling used. */ - @Order(3) + @Order(2) volatile CacheContinuousQueryDeployableObject rmtFilterFactoryDep; /** Marshalled {@link #rmtFilterFactory} if {@link #rmtFilterFactoryDep} is {@code null}. */ - @Order(4) + @Order(3) @Nullable volatile byte[] rmtFilterFactoryBytes; /** Remote transformer factory. */ volatile Factory, ?>> rmtTransFactory; /** Deployable object for {@link #rmtTransFactory}. Is {@code null} if no external marsshalling used. */ - @Order(5) + @Order(4) volatile CacheContinuousQueryDeployableObject rmtTransFactoryDep; /** Marshalled {@link #rmtTransFactory} if {@link #rmtTransFactoryDep} is {@code null}. */ - @Order(6) + @Order(5) @Nullable volatile byte[] rmtTransFactoryBytes; /** Cache name. */ - @Order(7) + @Order(6) String cacheName; /** Topic for ordered messages. */ @@ -187,39 +183,39 @@ public final class CacheContinuousQueryHandler implements GridContinuousHa Object topic; /** Marshalled {@link #topic}. */ - @Order(8) + @Order(7) byte[] topicBytes; /** Internal flag. */ - @Order(9) + @Order(8) boolean internal; /** Notify existing flag. */ - @Order(10) + @Order(9) boolean notifyExisting; /** Old value required flag. */ - @Order(11) + @Order(10) boolean oldValRequired; /** Synchronous flag. */ - @Order(12) + @Order(11) boolean sync; /** Ignore expired events flag. */ - @Order(13) + @Order(12) boolean ignoreExpired; /** Task name hash code. */ - @Order(14) + @Order(13) int taskHash; /** */ - @Order(15) + @Order(14) boolean keepBinary; /** Event types for JCache API. */ - @Order(16) + @Order(15) byte types; /** P2P unmarshalling future. */ @@ -1424,33 +1420,13 @@ CacheContinuousQueryEventBuffer partitionBuffer(GridCacheContext cctx, int assert ctx.config().isPeerClassLoadingEnabled(); if (requiresDeployment(rmtFilter)) - rmtFilterDep = marshalDeployable(rmtFilter, ctx); + rmtFilterDep = new CacheContinuousQueryDeployableObject(rmtFilter, ctx); if (requiresDeployment(rmtFilterFactory)) - rmtFilterFactoryDep = marshalDeployable(rmtFilterFactory, ctx); + rmtFilterFactoryDep = new CacheContinuousQueryDeployableObject(rmtFilterFactory, ctx); if (requiresDeployment(rmtTransFactory)) - rmtTransFactoryDep = marshalDeployable(rmtTransFactory, ctx); - } - - /** - * Processes {@link #p2pUnmarshalFut} over the super's method. - * - * @param deployable Deployable object. - * @param ctx Kernal context. - */ - private CacheContinuousQueryDeployableObject marshalDeployable( - Object deployable, - GridKernalContext ctx - ) throws IgniteCheckedException { - CacheContinuousQueryDeployableObject res = new CacheContinuousQueryDeployableObject(deployable, ctx); - - if (p2pUnmarshalFut == null) - p2pUnmarshalFut = new GridFutureAdapter<>(); - else if (p2pUnmarshalFut.isDone()) - p2pUnmarshalFut = new GridFutureAdapter<>(); - - return res; + rmtTransFactoryDep = new CacheContinuousQueryDeployableObject(rmtTransFactory, ctx); } /** */ @@ -1463,25 +1439,14 @@ private static boolean requiresDeployment(@Nullable Object obj) { /** @see #marshalDeployable(Object, GridKernalContext) */ p2pUnmarshalFut = null; - externalMarshal = new boolean[3]; - - if (rmtFilterDep == null && rmtFilter != null) { + if (rmtFilter != null && rmtFilterDep == null) rmtFilterBytes = marsh.marshal(rmtFilter); - externalMarshal[0] = true; - } - - if (rmtFilterFactoryDep == null && rmtFilterFactory != null) { + if (rmtFilterFactory != null && rmtFilterFactoryDep == null) rmtFilterFactoryBytes = marsh.marshal(rmtFilterFactory); - externalMarshal[1] = true; - } - - if (rmtTransFactoryDep == null && rmtTransFactory != null) { + if (rmtTransFactory != null && rmtTransFactoryDep == null) rmtTransFactoryBytes = marsh.marshal(rmtTransFactory); - - externalMarshal[2] = true; - } } /** {@inheritDoc} */ @@ -1517,25 +1482,17 @@ private static boolean requiresDeployment(@Nullable Object obj) { /** {@inheritDoc} */ @Override public void unmarshal(Marshaller marsh, ClassLoader clsLdr) throws IgniteCheckedException { - assert externalMarshal != null && externalMarshal.length == 3; - - if (externalMarshal[0]) { - assert rmtFilterBytes != null; - + if (rmtFilterBytes != null && rmtFilterDep == null) rmtFilter = marsh.unmarshal(rmtFilterBytes, clsLdr); - } - - if (externalMarshal[1]) { - assert rmtFilterFactoryBytes != null; + if (rmtFilterFactoryBytes != null && rmtFilterFactoryDep == null) rmtFilterFactory = marsh.unmarshal(rmtFilterFactoryBytes, clsLdr); - } - - if (externalMarshal[2]) { - assert rmtTransFactoryBytes != null; + if (rmtTransFactoryBytes != null && rmtTransFactoryDep == null) rmtTransFactory = marsh.unmarshal(rmtTransFactoryBytes, clsLdr); - } + + if (rmtFilterDep != null || rmtFilterFactoryDep != null || rmtTransFactoryDep != null) + p2pUnmarshalFut = new GridFutureAdapter<>(); cacheId = CU.cacheId(cacheName); } From 522850604ca1cccb4ca3049b7891f86e0ab940d1 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Mon, 3 Aug 2026 18:26:15 +0300 Subject: [PATCH 10/24] revert field reordering --- .../CacheContinuousQueryHandler.java | 99 ++++++++++--------- 1 file changed, 50 insertions(+), 49 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java index 27e000707d205..23e19163773ae 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java @@ -141,107 +141,108 @@ public final class CacheContinuousQueryHandler implements GridContinuousHa } }; + /** Cache name. */ + @Order(0) + String cacheName; + + /** Topic for ordered messages. */ + @Marshalled("topicBytes") + Object topic; + + /** Marshalled {@link #topic}. */ + @Order(1) + byte[] topicBytes; + + /** P2P unmarshalling future. */ + protected volatile IgniteInternalFuture p2pUnmarshalFut = new GridFinishedFuture<>(); + + /** Initialization future. */ + protected IgniteInternalFuture initFut; + + /** Local listener. */ + private CacheEntryUpdatedListener locLsnr; + /** Remote filter. */ volatile CacheEntryEventSerializableFilter rmtFilter; /** Deployable object for {@link #rmtFilter}. Is {@code null} if no external marsshalling used. */ - @Order(0) + @Order(2) @Nullable volatile CacheContinuousQueryDeployableObject rmtFilterDep; /** Marshalled {@link #rmtFilter} if {@link #rmtFilterDep} is {@code null}. */ - @Order(1) + @Order(3) @Nullable volatile byte[] rmtFilterBytes; /** Remote filter factory. */ @Nullable volatile Factory rmtFilterFactory; /** Deployable object for {@link #rmtFilterFactory}. Is {@code null} if no external marsshalling used. */ - @Order(2) + @Order(4) volatile CacheContinuousQueryDeployableObject rmtFilterFactoryDep; /** Marshalled {@link #rmtFilterFactory} if {@link #rmtFilterFactoryDep} is {@code null}. */ - @Order(3) + @Order(5) @Nullable volatile byte[] rmtFilterFactoryBytes; + /** Remote filter created by {@link #rmtFilterFactory}. */ + private CacheEntryEventFilter rmtFilterFromFactory; + + /** Event types for JCache API. */ + @Order(6) + byte types; + /** Remote transformer factory. */ volatile Factory, ?>> rmtTransFactory; + /** Deployable object for {@link #rmtTransFactory}. Is {@code null} if no external marsshalling used. */ - @Order(4) + @Order(7) volatile CacheContinuousQueryDeployableObject rmtTransFactoryDep; /** Marshalled {@link #rmtTransFactory} if {@link #rmtTransFactoryDep} is {@code null}. */ - @Order(5) + @Order(8) @Nullable volatile byte[] rmtTransFactoryBytes; - /** Cache name. */ - @Order(6) - String cacheName; - - /** Topic for ordered messages. */ - @Marshalled("topicBytes") - Object topic; + /** Remote transformer created by {@link #rmtTransFactory}. */ + private IgniteClosure, ?> rmtTrans; - /** Marshalled {@link #topic}. */ - @Order(7) - byte[] topicBytes; + /** Local listener for transformed events. */ + private EventListener locTransLsnr; /** Internal flag. */ - @Order(8) + @Order(9) boolean internal; /** Notify existing flag. */ - @Order(9) + @Order(10) boolean notifyExisting; /** Old value required flag. */ - @Order(10) + @Order(11) boolean oldValRequired; /** Synchronous flag. */ - @Order(11) + @Order(12) boolean sync; /** Ignore expired events flag. */ - @Order(12) + @Order(13) boolean ignoreExpired; /** Task name hash code. */ - @Order(13) - int taskHash; - - /** */ @Order(14) - boolean keepBinary; - - /** Event types for JCache API. */ - @Order(15) - byte types; - - /** P2P unmarshalling future. */ - protected volatile IgniteInternalFuture p2pUnmarshalFut = new GridFinishedFuture<>(); - - /** Initialization future. */ - protected IgniteInternalFuture initFut; - - /** Local listener. */ - private CacheEntryUpdatedListener locLsnr; - - /** Remote filter created by {@link #rmtFilterFactory}. */ - private CacheEntryEventFilter rmtFilterFromFactory; - - /** Remote transformer created by {@link #rmtTransFactory}. */ - private IgniteClosure, ?> rmtTrans; - - /** Local listener for transformed events. */ - private EventListener locTransLsnr; + int taskHash; /** Whether to skip primary check for REPLICATED cache. */ boolean skipPrimaryCheck; - + /** */ private boolean locOnly; + /** */ + @Order(15) + boolean keepBinary; + /** */ private ConcurrentMap rcvs; From f0710e67e02d2f676777a9ce19889d4aad9f247d Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Mon, 3 Aug 2026 18:47:35 +0300 Subject: [PATCH 11/24] fixes --- .../internal/GridEventConsumeHandler.java | 21 ++------- .../internal/GridMessageListenHandler.java | 22 ++------- .../CacheContinuousQueryHandler.java | 47 ++++++++++--------- .../continuous/GridContinuousProcessor.java | 4 +- 4 files changed, 36 insertions(+), 58 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java index d0c79bc093b72..e0dc2a86834bf 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java @@ -92,15 +92,6 @@ public final class GridEventConsumeHandler implements GridContinuousHandler, Mar @Order(3) int[] types; - /** - * Lever of own marshaling. - * - * @see #p2pMarshal(GridKernalContext) - * @see #marshal(Marshaller) - */ - @Order(4) - volatile boolean externalMarshal; - /** Listener. */ private GridLocalEventListener lsnr; @@ -414,8 +405,6 @@ private boolean filterDropsEvent(Event evt) { depInfo = new GridDeploymentInfoBean(dep); filterBytes = U.marshal(ctx.marshaller(), filter); - - externalMarshal = true; } } @@ -423,9 +412,10 @@ private boolean filterDropsEvent(Event evt) { @Override public void p2pUnmarshal(UUID nodeId, GridKernalContext ctx) throws IgniteCheckedException { assert nodeId != null; assert ctx.config().isPeerClassLoadingEnabled(); - assert externalMarshal : "Is not p2p-marshaled " + getClass().getSimpleName(); if (filterBytes != null) { + assert filter == null : "Already P2P-unmarshaled " + getClass().getSimpleName(); + try { GridDeployment dep = ctx.deploy().getGlobalDeployment(depInfo.deployMode(), clsName, clsName, depInfo.userVersion(), nodeId, depInfo.classLoaderId(), depInfo.participants(), null); @@ -483,19 +473,18 @@ private boolean filterDropsEvent(Event evt) { /** {@inheritDoc} */ @Override public void marshal(Marshaller marsh) throws IgniteCheckedException { assert (clsName == null) == (depInfo == null); - assert (depInfo == null) == !externalMarshal; - if (filter != null && !externalMarshal) + /** Are marshaled in {@link #p2pUnmarshal(UUID, GridKernalContext)}. */ + if (filter != null && depInfo == null) filterBytes = marsh.marshal(filter); } /** {@inheritDoc} */ @Override public void unmarshal(Marshaller marsh, ClassLoader clsLdr) throws IgniteCheckedException { assert (clsName == null) == (depInfo == null); - assert (depInfo == null) == !externalMarshal; /** Are unmarshaled in {@link #p2pUnmarshal(UUID, GridKernalContext)}. */ - if (externalMarshal) { + if (depInfo != null) { p2pUnmarshalFut = new GridFutureAdapter<>(); return; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java index 0529bc1a5e004..17c82c2e092ed 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java @@ -65,14 +65,6 @@ public final class GridMessageListenHandler implements GridContinuousHandler, Ma @Order(3) @Nullable volatile GridDeploymentInfoBean predDepInfo; - /** - * Lever of the own marshaling. - * - * @see #p2pMarshal(GridKernalContext) - */ - @Order(4) - volatile boolean externalMarshal; - /** P2P unmarshalling future. */ private volatile IgniteInternalFuture p2pUnmarshalFut = new GridFinishedFuture<>(); @@ -154,6 +146,7 @@ public GridMessageListenHandler(@Nullable Object topic, IgniteBiPredicate(); return; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java index 23e19163773ae..f75099e73692e 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java @@ -1418,6 +1418,7 @@ CacheContinuousQueryEventBuffer partitionBuffer(GridCacheContext cctx, int /** {@inheritDoc} */ @Override public void p2pMarshal(GridKernalContext ctx) throws IgniteCheckedException { + assert ctx != null; assert ctx.config().isPeerClassLoadingEnabled(); if (requiresDeployment(rmtFilter)) @@ -1430,11 +1431,6 @@ CacheContinuousQueryEventBuffer partitionBuffer(GridCacheContext cctx, int rmtTransFactoryDep = new CacheContinuousQueryDeployableObject(rmtTransFactory, ctx); } - /** */ - private static boolean requiresDeployment(@Nullable Object obj) { - return obj != null && !U.isGrid(obj.getClass()); - } - /** {@inheritDoc} */ @Override public void marshal(Marshaller marsh) throws IgniteCheckedException { /** @see #marshalDeployable(Object, GridKernalContext) */ @@ -1450,8 +1446,27 @@ private static boolean requiresDeployment(@Nullable Object obj) { rmtTransFactoryBytes = marsh.marshal(rmtTransFactory); } + /** {@inheritDoc} */ + @Override public void unmarshal(Marshaller marsh, ClassLoader clsLdr) throws IgniteCheckedException { + if (rmtFilterBytes != null && rmtFilterDep == null) + rmtFilter = marsh.unmarshal(rmtFilterBytes, clsLdr); + + if (rmtFilterFactoryBytes != null && rmtFilterFactoryDep == null) + rmtFilterFactory = marsh.unmarshal(rmtFilterFactoryBytes, clsLdr); + + if (rmtTransFactoryBytes != null && rmtTransFactoryDep == null) + rmtTransFactory = marsh.unmarshal(rmtTransFactoryBytes, clsLdr); + + if (rmtFilterDep != null || rmtFilterFactoryDep != null || rmtTransFactoryDep != null) + p2pUnmarshalFut = new GridFutureAdapter<>(); + + cacheId = CU.cacheId(cacheName); + } + /** {@inheritDoc} */ @Override public void p2pUnmarshal(UUID nodeId, GridKernalContext ctx) throws IgniteCheckedException { + assert nodeId != null; + assert ctx != null; assert ctx.config().isPeerClassLoadingEnabled(); try { @@ -1481,23 +1496,6 @@ private static boolean requiresDeployment(@Nullable Object obj) { } } - /** {@inheritDoc} */ - @Override public void unmarshal(Marshaller marsh, ClassLoader clsLdr) throws IgniteCheckedException { - if (rmtFilterBytes != null && rmtFilterDep == null) - rmtFilter = marsh.unmarshal(rmtFilterBytes, clsLdr); - - if (rmtFilterFactoryBytes != null && rmtFilterFactoryDep == null) - rmtFilterFactory = marsh.unmarshal(rmtFilterFactoryBytes, clsLdr); - - if (rmtTransFactoryBytes != null && rmtTransFactoryDep == null) - rmtTransFactory = marsh.unmarshal(rmtTransFactoryBytes, clsLdr); - - if (rmtFilterDep != null || rmtFilterFactoryDep != null || rmtTransFactoryDep != null) - p2pUnmarshalFut = new GridFutureAdapter<>(); - - cacheId = CU.cacheId(cacheName); - } - /** * @return Whether the handler is marshalled for peer class loading. */ @@ -1786,4 +1784,9 @@ private Object transform(IgniteClosure Map partitionContinuesQueryEntryBuffers() { return Collections.unmodifiableMap(entryBufs); } + + /** */ + private static boolean requiresDeployment(@Nullable Object obj) { + return obj != null && !U.isGrid(obj.getClass()); + } } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java index 7b4ae2a56ebf5..e2d3350a7a201 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java @@ -401,7 +401,7 @@ public void unlockStopping() { return; } - Message data = getDiscoveryData(dataBag.joiningNodeId()); + ContinousRoutineDiscoveryData data = getDiscoveryData(dataBag.joiningNodeId()); if (data != null) dataBag.addJoiningNodeData(CONTINUOUS_PROC.ordinal(), data); @@ -415,7 +415,7 @@ public void unlockStopping() { return; } - Message data = getDiscoveryData(dataBag.joiningNodeId()); + ContinousRoutineDiscoveryData data = getDiscoveryData(dataBag.joiningNodeId()); if (data != null) dataBag.addNodeSpecificData(CONTINUOUS_PROC.ordinal(), data); From 64808c1d9334e4e62454345f180a1c47c5233807 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Mon, 3 Aug 2026 18:52:17 +0300 Subject: [PATCH 12/24] fixes --- .../org/apache/ignite/internal/GridEventConsumeHandler.java | 2 -- .../org/apache/ignite/internal/GridMessageListenHandler.java | 2 -- 2 files changed, 4 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java index e0dc2a86834bf..031a1a94a9675 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java @@ -414,8 +414,6 @@ private boolean filterDropsEvent(Event evt) { assert ctx.config().isPeerClassLoadingEnabled(); if (filterBytes != null) { - assert filter == null : "Already P2P-unmarshaled " + getClass().getSimpleName(); - try { GridDeployment dep = ctx.deploy().getGlobalDeployment(depInfo.deployMode(), clsName, clsName, depInfo.userVersion(), nodeId, depInfo.classLoaderId(), depInfo.participants(), null); diff --git a/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java index 17c82c2e092ed..2efee6e39e937 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java @@ -146,7 +146,6 @@ public GridMessageListenHandler(@Nullable Object topic, IgniteBiPredicate Date: Mon, 3 Aug 2026 19:21:38 +0300 Subject: [PATCH 13/24] fix --- .../cache/query/continuous/CacheContinuousQueryHandler.java | 2 -- 1 file changed, 2 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java index f75099e73692e..fdb4728673f70 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java @@ -1434,8 +1434,6 @@ CacheContinuousQueryEventBuffer partitionBuffer(GridCacheContext cctx, int /** {@inheritDoc} */ @Override public void marshal(Marshaller marsh) throws IgniteCheckedException { /** @see #marshalDeployable(Object, GridKernalContext) */ - p2pUnmarshalFut = null; - if (rmtFilter != null && rmtFilterDep == null) rmtFilterBytes = marsh.marshal(rmtFilter); From dc22e2012dcdb24deae8ad4e5e1540481788404b Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Mon, 3 Aug 2026 19:22:54 +0300 Subject: [PATCH 14/24] fix --- .../cache/query/continuous/CacheContinuousQueryHandler.java | 1 - 1 file changed, 1 deletion(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java index fdb4728673f70..67c63ee574481 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java @@ -1433,7 +1433,6 @@ CacheContinuousQueryEventBuffer partitionBuffer(GridCacheContext cctx, int /** {@inheritDoc} */ @Override public void marshal(Marshaller marsh) throws IgniteCheckedException { - /** @see #marshalDeployable(Object, GridKernalContext) */ if (rmtFilter != null && rmtFilterDep == null) rmtFilterBytes = marsh.marshal(rmtFilter); From a65c3355ad50347d8990ac9807a4f50f0bf72e17 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Tue, 4 Aug 2026 10:41:47 +0300 Subject: [PATCH 15/24] fix --- .../CacheContinuousQueryHandler.java | 22 +++++++++---------- .../continuous/GridContinuousProcessor.java | 9 +------- 2 files changed, 12 insertions(+), 19 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java index 67c63ee574481..90c58917f6d4a 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java @@ -88,6 +88,7 @@ import org.apache.ignite.marshaller.Marshaller; import org.jetbrains.annotations.NotNull; import org.jetbrains.annotations.Nullable; +import org.jetbrains.annotations.TestOnly; import static javax.cache.event.EventType.EXPIRED; import static javax.cache.event.EventType.REMOVED; @@ -1421,13 +1422,18 @@ CacheContinuousQueryEventBuffer partitionBuffer(GridCacheContext cctx, int assert ctx != null; assert ctx.config().isPeerClassLoadingEnabled(); - if (requiresDeployment(rmtFilter)) + /** + * Some filters, factories might be an Ignite-internals and do not require external marshaling. But there is no + * quarantine that a user-defuned class is not included in a wrap like {@link SecurityAwareFilter}. Hence, we always + * externally-marshall here. + */ + if (rmtFilter != null) rmtFilterDep = new CacheContinuousQueryDeployableObject(rmtFilter, ctx); - if (requiresDeployment(rmtFilterFactory)) + if (rmtFilterFactory != null) rmtFilterFactoryDep = new CacheContinuousQueryDeployableObject(rmtFilterFactory, ctx); - if (requiresDeployment(rmtTransFactory)) + if (rmtTransFactory != null) rmtTransFactoryDep = new CacheContinuousQueryDeployableObject(rmtTransFactory, ctx); } @@ -1496,10 +1502,9 @@ CacheContinuousQueryEventBuffer partitionBuffer(GridCacheContext cctx, int /** * @return Whether the handler is marshalled for peer class loading. */ + @TestOnly public boolean isMarshalled() { - return (!requiresDeployment(rmtFilter) || rmtFilterDep != null) - && (!requiresDeployment(rmtFilterFactory) || rmtFilterFactoryDep != null) - && (!requiresDeployment(rmtTransFactory) || rmtTransFactoryDep != null); + return rmtFilterDep != null || rmtFilterFactoryDep != null || rmtTransFactoryDep != null; } /** {@inheritDoc} */ @@ -1781,9 +1786,4 @@ private Object transform(IgniteClosure Map partitionContinuesQueryEntryBuffers() { return Collections.unmodifiableMap(entryBufs); } - - /** */ - private static boolean requiresDeployment(@Nullable Object obj) { - return obj != null && !U.isGrid(obj.getClass()); - } } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java index e2d3350a7a201..5a7127bba5399 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java @@ -449,10 +449,6 @@ public void unlockStopping() { UUID routineId = e.getKey(); ContinousRoutineLocalInfo info = e.getValue(); - assert !ctx.config().isPeerClassLoadingEnabled() || - !(info.hnd instanceof CacheContinuousQueryHandler) || - ((CacheContinuousQueryHandler)info.hnd).isMarshalled(); - data.addItem(new ContinousRoutineDiscoveryDataItem(routineId, info.prjPred, info.hnd, @@ -835,12 +831,9 @@ public IgniteInternalFuture startRoutine(GridContinuousHandler hnd, // Generate ID. final UUID routineId = UUID.randomUUID(); - if (ctx.config().isPeerClassLoadingEnabled()) { + if (ctx.config().isPeerClassLoadingEnabled()) hnd.p2pMarshal(ctx); - assert !(hnd instanceof CacheContinuousQueryHandler) || ((CacheContinuousQueryHandler)hnd).isMarshalled(); - } - // Register routine locally. locInfos.put(routineId, new ContinousRoutineLocalInfo(ctx.localNodeId(), prjPred, hnd, bufSize, interval, autoUnsubscribe)); From 70c8d1b3a966bf690bfb4490fdbc0455644a65cb Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Tue, 4 Aug 2026 11:20:10 +0300 Subject: [PATCH 16/24] fix --- .../query/continuous/CacheContinuousQueryHandler.java | 9 --------- 1 file changed, 9 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java index 90c58917f6d4a..8e57506a75a91 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java @@ -88,7 +88,6 @@ import org.apache.ignite.marshaller.Marshaller; import org.jetbrains.annotations.NotNull; import org.jetbrains.annotations.Nullable; -import org.jetbrains.annotations.TestOnly; import static javax.cache.event.EventType.EXPIRED; import static javax.cache.event.EventType.REMOVED; @@ -1499,14 +1498,6 @@ CacheContinuousQueryEventBuffer partitionBuffer(GridCacheContext cctx, int } } - /** - * @return Whether the handler is marshalled for peer class loading. - */ - @TestOnly - public boolean isMarshalled() { - return rmtFilterDep != null || rmtFilterFactoryDep != null || rmtTransFactoryDep != null; - } - /** {@inheritDoc} */ @Override public GridContinuousBatch createBatch() { return new GridContinuousQueryBatch(); From 18d02680fa791ba6645841e476575ae353240330 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Tue, 4 Aug 2026 11:30:57 +0300 Subject: [PATCH 17/24] fix --- .../internal/GridMessageListenHandler.java | 16 +++++++++++----- 1 file changed, 11 insertions(+), 5 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java index 2efee6e39e937..487a0f2d83896 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java @@ -168,12 +168,14 @@ public GridMessageListenHandler(@Nullable Object topic, IgniteBiPredicate(); return; } - topic = marsh.unmarshal(topicBytes, clsLdr); + assert predBytes != null; + + if (topicBytes != null) + topic = marsh.unmarshal(topicBytes, clsLdr); + pred = marsh.unmarshal(predBytes, clsLdr); } From 952ce55afd7fa2a163926b88a795513b02885931 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Tue, 4 Aug 2026 13:55:05 +0300 Subject: [PATCH 18/24] fix --- .../org/apache/ignite/internal/GridMessageListenHandler.java | 2 -- 1 file changed, 2 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java index 487a0f2d83896..bf220092bb0bb 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java @@ -217,8 +217,6 @@ public GridMessageListenHandler(@Nullable Object topic, IgniteBiPredicate(); return; From 1292013cf452546bbd2f0feaabba8bc46eaa220c Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Tue, 4 Aug 2026 14:15:38 +0300 Subject: [PATCH 19/24] fix --- .../java/org/apache/ignite/internal/GridEventConsumeHandler.java | 1 - .../org/apache/ignite/internal/GridMessageListenHandler.java | 1 - .../cache/query/continuous/CacheContinuousQueryHandler.java | 1 - 3 files changed, 3 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java index 031a1a94a9675..6c8a1dd7d7086 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java @@ -61,7 +61,6 @@ /** * Continuous routine handler for remote event listening. */ -@UseBinaryMarshaller public final class GridEventConsumeHandler implements GridContinuousHandler, MarshallableMessage { /** Default callback. */ private static final IgniteBiPredicate DFLT_CALLBACK = new P2<>() { diff --git a/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java index bf220092bb0bb..4cfa498b05fc8 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java @@ -41,7 +41,6 @@ /** * Continuous handler for message subscription. */ -@UseBinaryMarshaller public final class GridMessageListenHandler implements GridContinuousHandler, MarshallableMessage { /** */ private volatile @Nullable Object topic; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java index 8e57506a75a91..6c503832e6478 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java @@ -99,7 +99,6 @@ /** * Continuous query handler. */ -@UseBinaryMarshaller public final class CacheContinuousQueryHandler implements GridContinuousHandler, MarshallableMessage { /** @see #IGNITE_CONTINUOUS_QUERY_BACKUP_ACK_THRESHOLD */ public static final int DFLT_CONTINUOUS_QUERY_BACKUP_ACK_THRESHOLD = 100; From 84da647b35a876a5187142f4fac4a9661002a4b7 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Tue, 4 Aug 2026 15:03:32 +0300 Subject: [PATCH 20/24] fix --- .../org/apache/ignite/internal/GridMessageListenHandler.java | 4 ---- 1 file changed, 4 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java index 4cfa498b05fc8..a68abb6192871 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java @@ -170,8 +170,6 @@ public GridMessageListenHandler(@Nullable Object topic, IgniteBiPredicate Date: Tue, 4 Aug 2026 15:21:54 +0300 Subject: [PATCH 21/24] + master --- .../org/apache/ignite/internal/CoreMessagesProvider.java | 7 ------- 1 file changed, 7 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java b/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java index 55dfef184fe1a..a6e4c5e79a255 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java @@ -179,10 +179,8 @@ import org.apache.ignite.internal.processors.cache.query.GridCacheSqlQuery; import org.apache.ignite.internal.processors.cache.query.continuous.CacheContinuousQueryBatchAck; import org.apache.ignite.internal.processors.cache.query.continuous.CacheContinuousQueryDeployableObject; -import org.apache.ignite.internal.processors.cache.query.continuous.CacheContinuousQueryDeployableObject; import org.apache.ignite.internal.processors.cache.query.continuous.CacheContinuousQueryEntry; import org.apache.ignite.internal.processors.cache.query.continuous.CacheContinuousQueryHandler; -import org.apache.ignite.internal.processors.cache.query.continuous.CacheContinuousQueryHandler; import org.apache.ignite.internal.processors.cache.transactions.IgniteTxEntry; import org.apache.ignite.internal.processors.cache.transactions.IgniteTxKey; import org.apache.ignite.internal.processors.cache.transactions.TxEntryValueHolder; @@ -204,13 +202,8 @@ import org.apache.ignite.internal.processors.continuous.ContinousRoutineDiscoveryDataItem; import org.apache.ignite.internal.processors.continuous.ContinousRoutineLocalInfo; import org.apache.ignite.internal.processors.continuous.ContinuousRoutineInfo; -import org.apache.ignite.internal.processors.continuous.ContinousRoutineDiscoveryData; -import org.apache.ignite.internal.processors.continuous.ContinousRoutineDiscoveryDataItem; -import org.apache.ignite.internal.processors.continuous.ContinousRoutineLocalInfo; -import org.apache.ignite.internal.processors.continuous.ContinuousRoutineInfo; import org.apache.ignite.internal.processors.continuous.ContinuousRoutineStartResultMessage; import org.apache.ignite.internal.processors.continuous.ContinuousRoutinesJoiningNodeDiscoveryData; -import org.apache.ignite.internal.processors.continuous.ContinuousRoutinesJoiningNodeDiscoveryData; import org.apache.ignite.internal.processors.continuous.GridContinuousMessage; import org.apache.ignite.internal.processors.continuous.StartRequestData; import org.apache.ignite.internal.processors.continuous.StartRoutineAckDiscoveryMessage; From bccd8ebc6db30f8d33809a76fdb2a6ffac0735ac Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Tue, 4 Aug 2026 16:38:03 +0300 Subject: [PATCH 22/24] minority --- .../org/apache/ignite/internal/GridEventConsumeHandler.java | 4 +++- .../processors/continuous/GridEventConsumeSelfTest.java | 6 ++---- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java index 6c8a1dd7d7086..9077a987802e8 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java @@ -389,10 +389,11 @@ private boolean filterDropsEvent(Event evt) { /** {@inheritDoc} */ @Override public void p2pMarshal(GridKernalContext ctx) throws IgniteCheckedException { + assert ctx != null; assert ctx.config().isPeerClassLoadingEnabled(); if (filter != null) { - Class cls = U.detectClass(filter); + Class cls = U.detectClass(filter); clsName = cls.getName(); @@ -410,6 +411,7 @@ private boolean filterDropsEvent(Event evt) { /** {@inheritDoc} */ @Override public void p2pUnmarshal(UUID nodeId, GridKernalContext ctx) throws IgniteCheckedException { assert nodeId != null; + assert ctx != null; assert ctx.config().isPeerClassLoadingEnabled(); if (filterBytes != null) { diff --git a/modules/core/src/test/java/org/apache/ignite/internal/processors/continuous/GridEventConsumeSelfTest.java b/modules/core/src/test/java/org/apache/ignite/internal/processors/continuous/GridEventConsumeSelfTest.java index 96e6e1d50cb12..c16bf2577751d 100644 --- a/modules/core/src/test/java/org/apache/ignite/internal/processors/continuous/GridEventConsumeSelfTest.java +++ b/modules/core/src/test/java/org/apache/ignite/internal/processors/continuous/GridEventConsumeSelfTest.java @@ -157,14 +157,12 @@ public class GridEventConsumeSelfTest extends GridCommonAbstractTest { * @return Local event routines. */ private Collection localRoutines(GridContinuousProcessor proc) { - return F.view( - U.>field(proc, "locInfos").values(), + return F.view(U.>field(proc, "locInfos").values(), new IgnitePredicate<>() { @Override public boolean apply(ContinousRoutineLocalInfo info) { return info.handler().isEvents(); } - } - ); + }); } /** From be1c9a3e9b5a09cbb8c66f6b6a13df410e69c997 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Tue, 4 Aug 2026 16:38:03 +0300 Subject: [PATCH 23/24] fixes --- .../org/apache/ignite/internal/GridEventConsumeHandler.java | 4 +++- .../cache/query/continuous/CacheContinuousQueryHandler.java | 1 - .../continuous/ContinousRoutineDiscoveryDataItem.java | 2 -- .../processors/continuous/GridEventConsumeSelfTest.java | 6 ++---- 4 files changed, 5 insertions(+), 8 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java index 6c8a1dd7d7086..9077a987802e8 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java @@ -389,10 +389,11 @@ private boolean filterDropsEvent(Event evt) { /** {@inheritDoc} */ @Override public void p2pMarshal(GridKernalContext ctx) throws IgniteCheckedException { + assert ctx != null; assert ctx.config().isPeerClassLoadingEnabled(); if (filter != null) { - Class cls = U.detectClass(filter); + Class cls = U.detectClass(filter); clsName = cls.getName(); @@ -410,6 +411,7 @@ private boolean filterDropsEvent(Event evt) { /** {@inheritDoc} */ @Override public void p2pUnmarshal(UUID nodeId, GridKernalContext ctx) throws IgniteCheckedException { assert nodeId != null; + assert ctx != null; assert ctx.config().isPeerClassLoadingEnabled(); if (filterBytes != null) { diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java index 6c503832e6478..dd44adf227f4b 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java @@ -52,7 +52,6 @@ import org.apache.ignite.internal.MarshallableMessage; import org.apache.ignite.internal.Marshalled; import org.apache.ignite.internal.Order; -import org.apache.ignite.internal.UseBinaryMarshaller; import org.apache.ignite.internal.cluster.ClusterTopologyCheckedException; import org.apache.ignite.internal.managers.communication.GridIoPolicy; import org.apache.ignite.internal.managers.communication.MessageMarshalling; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryDataItem.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryDataItem.java index 247c5be996cb2..ce4ead23f101f 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryDataItem.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryDataItem.java @@ -21,14 +21,12 @@ import org.apache.ignite.cluster.ClusterNode; import org.apache.ignite.internal.Marshalled; import org.apache.ignite.internal.Order; -import org.apache.ignite.internal.UseBinaryMarshaller; import org.apache.ignite.internal.util.typedef.internal.S; import org.apache.ignite.lang.IgnitePredicate; import org.apache.ignite.plugin.extensions.communication.Message; import org.jetbrains.annotations.Nullable; /** Discovery data item. */ -@UseBinaryMarshaller public class ContinousRoutineDiscoveryDataItem implements Message { /** */ @Order(0) diff --git a/modules/core/src/test/java/org/apache/ignite/internal/processors/continuous/GridEventConsumeSelfTest.java b/modules/core/src/test/java/org/apache/ignite/internal/processors/continuous/GridEventConsumeSelfTest.java index 96e6e1d50cb12..c16bf2577751d 100644 --- a/modules/core/src/test/java/org/apache/ignite/internal/processors/continuous/GridEventConsumeSelfTest.java +++ b/modules/core/src/test/java/org/apache/ignite/internal/processors/continuous/GridEventConsumeSelfTest.java @@ -157,14 +157,12 @@ public class GridEventConsumeSelfTest extends GridCommonAbstractTest { * @return Local event routines. */ private Collection localRoutines(GridContinuousProcessor proc) { - return F.view( - U.>field(proc, "locInfos").values(), + return F.view(U.>field(proc, "locInfos").values(), new IgnitePredicate<>() { @Override public boolean apply(ContinousRoutineLocalInfo info) { return info.handler().isEvents(); } - } - ); + }); } /** From 436ae3abda05fb2f8447f11db2f238945a460ce0 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Tue, 4 Aug 2026 19:38:19 +0300 Subject: [PATCH 24/24] minority --- .../internal/processors/affinity/GridAffinityUtils.java | 2 +- .../continuous/CacheContinuousQueryDeployableObject.java | 2 +- .../cache/query/continuous/CacheContinuousQueryHandler.java | 1 - .../processors/continuous/ContinousRoutineDiscoveryData.java | 2 +- .../continuous/ContinousRoutineDiscoveryDataItem.java | 2 +- .../processors/continuous/ContinousRoutineLocalInfo.java | 2 +- .../internal/processors/continuous/ContinuousRoutineInfo.java | 2 +- .../ContinuousRoutinesJoiningNodeDiscoveryData.java | 2 +- .../java/org/apache/ignite/spi/discovery/tcp/ClientImpl.java | 2 +- .../java/org/apache/ignite/spi/discovery/tcp/ServerImpl.java | 4 ++-- .../spi/discovery/zk/internal/ZookeeperDiscoveryImpl.java | 2 +- 11 files changed, 11 insertions(+), 12 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/affinity/GridAffinityUtils.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/affinity/GridAffinityUtils.java index 3937be9b792a4..8b6b95b6f0e2d 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/affinity/GridAffinityUtils.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/affinity/GridAffinityUtils.java @@ -86,7 +86,7 @@ private static GridAffinityMessage affinityMessage(GridKernalContext ctx, Object } /** - * Unmarshalls transfer object from remote node within a given context. + * Unmarshals transfer object from remote node within a given context. * * @param ctx Grid kernal context that provides deployment and marshalling services. * @param sndNodeId {@link UUID} of the sender node. diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryDeployableObject.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryDeployableObject.java index 647e4feeae175..032eb62a2736f 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryDeployableObject.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryDeployableObject.java @@ -32,7 +32,7 @@ /** * Deployable object. */ -public class CacheContinuousQueryDeployableObject implements Message { +public final class CacheContinuousQueryDeployableObject implements Message { /** Serialized object. */ @GridToStringExclude @Order(0) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java index dd44adf227f4b..a7914be622fea 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java @@ -192,7 +192,6 @@ public final class CacheContinuousQueryHandler implements GridContinuousHa /** Remote transformer factory. */ volatile Factory, ?>> rmtTransFactory; - /** Deployable object for {@link #rmtTransFactory}. Is {@code null} if no external marsshalling used. */ @Order(7) volatile CacheContinuousQueryDeployableObject rmtTransFactoryDep; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryData.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryData.java index d6d1340649f4b..3ce690a283f1e 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryData.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryData.java @@ -27,7 +27,7 @@ import org.apache.ignite.plugin.extensions.communication.Message; /** Discovery data. */ -public class ContinousRoutineDiscoveryData implements Message { +public final class ContinousRoutineDiscoveryData implements Message { /** Node ID. */ @Order(0) UUID nodeId; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryDataItem.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryDataItem.java index ce4ead23f101f..9da8946bf6617 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryDataItem.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryDataItem.java @@ -27,7 +27,7 @@ import org.jetbrains.annotations.Nullable; /** Discovery data item. */ -public class ContinousRoutineDiscoveryDataItem implements Message { +public final class ContinousRoutineDiscoveryDataItem implements Message { /** */ @Order(0) UUID routineId; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineLocalInfo.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineLocalInfo.java index 900600c54954b..18d67fad06f59 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineLocalInfo.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineLocalInfo.java @@ -29,7 +29,7 @@ /** Local routine info. */ @UseBinaryMarshaller -public class ContinousRoutineLocalInfo implements Message, GridContinuousProcessor.RoutineInfo { +public final class ContinousRoutineLocalInfo implements Message, GridContinuousProcessor.RoutineInfo { /** Source node id. */ @Order(0) UUID nodeId; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutineInfo.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutineInfo.java index 86f7b0db9803b..77147f4c47eb5 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutineInfo.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutineInfo.java @@ -25,7 +25,7 @@ /** * */ -public class ContinuousRoutineInfo implements Message { +public final class ContinuousRoutineInfo implements Message { /** */ @Order(0) UUID srcNodeId; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesJoiningNodeDiscoveryData.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesJoiningNodeDiscoveryData.java index 1849e472672de..ec40f56ca1f3f 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesJoiningNodeDiscoveryData.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesJoiningNodeDiscoveryData.java @@ -25,7 +25,7 @@ /** * */ -public class ContinuousRoutinesJoiningNodeDiscoveryData implements Message { +public final class ContinuousRoutinesJoiningNodeDiscoveryData implements Message { /** */ @Order(0) List startedRoutines; diff --git a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ClientImpl.java b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ClientImpl.java index 43d4f1b4867a3..c31cf0a8a6cb3 100644 --- a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ClientImpl.java +++ b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ClientImpl.java @@ -847,7 +847,7 @@ private static void sleepEx(long millis, Runnable before, Runnable after) throws } /** - * Marshalls credentials with discovery SPI marshaller (will replace attribute value). + * Marshals credentials with discovery SPI marshaller (will replace attribute value). * * @param node Node to marshall credentials for. * @throws IgniteSpiException If marshalling failed. diff --git a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ServerImpl.java b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ServerImpl.java index 1d70baed0aa48..46ef634bb0c3e 100644 --- a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ServerImpl.java +++ b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ServerImpl.java @@ -1665,7 +1665,7 @@ else if (U.millisSinceNanos(joinStartNanos) > spi.joinTimeout) } /** - * Marshalls credentials with discovery SPI marshaller (will replace attribute value). + * Marshals credentials with discovery SPI marshaller (will replace attribute value). * * @param node Node to marshall credentials for. * @param cred Credentials for marshall. @@ -1686,7 +1686,7 @@ private void marshalCredentials(TcpDiscoveryNode node, SecurityCredentials cred) } /** - * Unmarshalls credentials with discovery SPI marshaller (will not replace attribute value). + * Unmarshals credentials with discovery SPI marshaller (will not replace attribute value). * * @param node Node to unmarshall credentials for. * @return Security credentials. diff --git a/modules/zookeeper/src/main/java/org/apache/ignite/spi/discovery/zk/internal/ZookeeperDiscoveryImpl.java b/modules/zookeeper/src/main/java/org/apache/ignite/spi/discovery/zk/internal/ZookeeperDiscoveryImpl.java index eabb741a60564..87799c52e3cd4 100644 --- a/modules/zookeeper/src/main/java/org/apache/ignite/spi/discovery/zk/internal/ZookeeperDiscoveryImpl.java +++ b/modules/zookeeper/src/main/java/org/apache/ignite/spi/discovery/zk/internal/ZookeeperDiscoveryImpl.java @@ -1156,7 +1156,7 @@ private SecurityCredentials unmarshalCredentials(ZookeeperClusterNode node) thro } /** - * Marshalls credentials with discovery SPI marshaller (will replace attribute value). + * Marshals credentials with discovery SPI marshaller (will replace attribute value). * * @param node Node to marshall credentials for. * @throws IgniteSpiException If marshalling failed.