This is an automated email from the ASF dual-hosted git repository.
shishkovilja pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ignite.git
The following commit(s) were added to refs/heads/master by this push:
new f83f8a55305 IGNITE-28790 Use MessageSerializer to transfer continuous
query DTO's in DataBag. (#13426)
f83f8a55305 is described below
commit f83f8a55305d36c7a25dd5cdeb88621763611b7b
Author: Vladimir Steshin <[email protected]>
AuthorDate: Wed Aug 12 17:02:08 2026 +0300
IGNITE-28790 Use MessageSerializer to transfer continuous query DTO's in
DataBag. (#13426)
---
.../ignite/internal/CoreMessagesProvider.java | 18 +
.../ignite/internal/GridEventConsumeHandler.java | 76 ++--
.../ignite/internal/GridMessageListenHandler.java | 113 +++---
.../deployment/GridDeploymentInfoMessage.java | 2 +-
.../processors/affinity/GridAffinityUtils.java | 2 +-
.../CacheContinuousQueryDeployableObject.java | 37 +-
.../continuous/CacheContinuousQueryHandler.java | 333 ++++++++----------
.../continuous/ContinousRoutineDiscoveryData.java | 72 ++++
.../ContinousRoutineDiscoveryDataItem.java | 96 +++++
.../continuous/ContinousRoutineLocalInfo.java | 133 +++++++
.../continuous/ContinuousRoutineInfo.java | 44 ++-
.../ContinuousRoutinesCommonDiscoveryData.java | 14 +-
...ContinuousRoutinesJoiningNodeDiscoveryData.java | 14 +-
.../continuous/GridContinuousHandler.java | 4 +-
.../continuous/GridContinuousProcessor.java | 387 +++------------------
.../processors/continuous/StartRequestData.java | 5 +-
.../org/apache/ignite/spi/IgniteSpiAdapter.java | 4 +-
.../ignite/spi/discovery/DiscoveryDataBag.java | 8 -
.../ignite/spi/discovery/tcp/ClientImpl.java | 2 +-
.../ignite/spi/discovery/tcp/ServerImpl.java | 4 +-
.../ignite/spi/discovery/tcp/TcpDiscoverySpi.java | 13 +
.../CacheContinuousQueryEntriesExpireTest.java | 4 +-
...ueryRemoteFilterMissingInClassPathSelfTest.java | 53 +--
...coveryDataDeserializationFailureHanderTest.java | 1 +
...GridCacheContinuousQueryNodesFilteringTest.java | 15 +-
.../continuous/GridEventConsumeSelfTest.java | 9 +-
.../zk/internal/ZookeeperDiscoveryImpl.java | 2 +-
27 files changed, 729 insertions(+), 736 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 000cbd47f02..d95ce2dab4e 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
@@ -178,7 +178,9 @@ import
org.apache.ignite.internal.processors.cache.query.GridCacheQueryRequest;
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.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;
@@ -199,7 +201,13 @@ import
org.apache.ignite.internal.processors.cluster.ClusterMetricsUpdateMessage
import
org.apache.ignite.internal.processors.cluster.ClusterUpdateNotifierDataBagItem;
import org.apache.ignite.internal.processors.cluster.DiscoveryDataClusterState;
import org.apache.ignite.internal.processors.cluster.NodeFullMetricsMessage;
+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.ContinuousRoutinesCommonDiscoveryData;
+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;
@@ -627,6 +635,16 @@ public class CoreMessagesProvider extends
AbstractMessageFactoryProvider {
register(QueryProposalsDataBagItem.class);
register(QueryEntityMessage.class);
register(QueryEntityExMessage.class);
+ register(ContinuousRoutineInfo.class);
+ register(ContinuousRoutinesJoiningNodeDiscoveryData.class);
+ register(CacheContinuousQueryDeployableObject.class);
+ register(CacheContinuousQueryHandler.class);
+ register(GridEventConsumeHandler.class);
+ register(GridMessageListenHandler.class);
+ register(ContinousRoutineLocalInfo.class);
+ register(ContinousRoutineDiscoveryDataItem.class);
+ register(ContinousRoutineDiscoveryData.class);
+ register(ContinuousRoutinesCommonDiscoveryData.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 8ebea159666..af4d85c204f 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,12 +61,9 @@ import static org.apache.ignite.events.EventType.EVTS_ALL;
/**
* 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<UUID, Event> DFLT_CALLBACK = new
P2<UUID, Event>() {
+ private static final IgniteBiPredicate<UUID, Event> DFLT_CALLBACK = new
P2<>() {
@Override public boolean apply(UUID uuid, Event e) {
return true;
}
@@ -76,28 +73,32 @@ class GridEventConsumeHandler implements
GridContinuousHandler {
private IgniteBiPredicate<UUID, Event> cb;
/** Filter. */
- private IgnitePredicate<Event> filter;
+ @Nullable volatile IgnitePredicate<Event> filter;
- /** Serialized filter. */
- private byte[] filterBytes;
+ /** Marshaled {@link #filter}. */
+ @Order(0)
+ @Nullable volatile byte[] filterBytes;
- /** Deployment class name. */
- private String clsName;
+ /** Deployment class name. Is {@code null} if P2P deployment is disabled.
*/
+ @Order(1)
+ @Nullable volatile String clsName;
- /** Deployment info. */
- private GridDeploymentInfo depInfo;
+ /** Deployment info. Is {@code null} if P2P deployment is disabled. */
+ @Order(2)
+ @Nullable volatile GridDeploymentInfoMessage depInfo;
/** Types. */
- private int[] types;
+ @Order(3)
+ int[] types;
/** Listener. */
private GridLocalEventListener lsnr;
/** P2P unmarshalling future. */
- private IgniteInternalFuture<Void> p2pUnmarshalFut = new
GridFinishedFuture<>();
+ private volatile IgniteInternalFuture<Void> p2pUnmarshalFut = new
GridFinishedFuture<>();
/**
- * Required by {@link Externalizable}.
+ * Empty constructor for serialization purposes.
*/
public GridEventConsumeHandler() {
// No-op.
@@ -225,8 +226,6 @@ class GridEventConsumeHandler implements
GridContinuousHandler {
EventWrapper wrapper = new
EventWrapper(evt);
if (evt instanceof CacheEvent)
{
- String cacheName =
((CacheEvent)evt).cacheName();
-
ClusterNode node =
ctx.discovery().node(t3.get1());
if (node == null)
@@ -393,6 +392,10 @@ class GridEventConsumeHandler implements
GridContinuousHandler {
assert ctx != null;
assert ctx.config().isPeerClassLoadingEnabled();
+ // TODO : Remove this check after
https://issues.apache.org/jira/browse/IGNITE-28945
+ if (filterBytes != null)
+ return;
+
if (filter != null) {
Class cls = U.detectClass(filter);
@@ -415,6 +418,10 @@ class GridEventConsumeHandler implements
GridContinuousHandler {
assert ctx != null;
assert ctx.config().isPeerClassLoadingEnabled();
+ // TODO : Remove this check after
https://issues.apache.org/jira/browse/IGNITE-28945
+ if (filter != null)
+ return;
+
if (filterBytes != null) {
try {
GridDeployment dep = ctx.deploy().globalDeployment(depInfo,
clsName, nodeId);
@@ -467,36 +474,27 @@ class GridEventConsumeHandler implements
GridContinuousHandler {
}
/** {@inheritDoc} */
- @Override public void writeExternal(ObjectOutput out) throws IOException {
- boolean b = filterBytes != null;
+ @Override public void marshal(Marshaller marsh) throws
IgniteCheckedException {
+ assert (clsName == null) == (depInfo == null);
- out.writeBoolean(b);
-
- if (b) {
- U.writeByteArray(out, filterBytes);
- U.writeString(out, clsName);
- out.writeObject(depInfo);
- }
- else
- out.writeObject(filter);
-
- out.writeObject(types);
+ /** Are marshaled in {@link #p2pUnmarshal(UUID, GridKernalContext)}. */
+ if (filter != null && depInfo == null)
+ filterBytes = marsh.marshal(filter);
}
/** {@inheritDoc} */
- @Override public void readExternal(ObjectInput in) throws IOException,
ClassNotFoundException {
- boolean b = in.readBoolean();
+ @Override public void unmarshal(Marshaller marsh, ClassLoader clsLdr)
throws IgniteCheckedException {
+ assert (clsName == null) == (depInfo == null);
- if (b) {
+ /** Are unmarshaled in {@link #p2pUnmarshal(UUID, GridKernalContext)}.
*/
+ if (depInfo != null) {
p2pUnmarshalFut = new GridFutureAdapter<>();
- filterBytes = U.readByteArray(in);
- clsName = U.readString(in);
- depInfo = (GridDeploymentInfo)in.readObject();
+
+ return;
}
- else
- filter = (IgnitePredicate<Event>)in.readObject();
- 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 f3397e5d214..c1bd4a683c1 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,40 @@ import
org.apache.ignite.internal.util.lang.GridPeerDeployAware;
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 volatile @Nullable Object topic;
- /** */
- private Object topic;
+ /** Marshalled {@link #topic}. */
+ @Order(0)
+ @Nullable volatile byte[] topicBytes;
/** */
- private IgniteBiPredicate<UUID, Object> pred;
+ private volatile IgniteBiPredicate<UUID, Object> pred;
- /** */
- private byte[] topicBytes;
+ /** Marshalled {@link #pred}. */
+ @Order(1)
+ volatile byte[] predBytes;
- /** */
- private byte[] predBytes;
+ /** Class name of {@link #pred}. Is {@code null} if the P2P deployment is
disabled. */
+ @Order(2)
+ @Nullable volatile String clsName;
- /** */
- private String clsName;
-
- /** */
- private GridDeploymentInfoMessage depInfo;
-
- /** */
- private boolean depEnabled;
+ /** P2P deploy info of {@link #pred}. Is {@code null} if the P2P
deployment is disabled. */
+ @Order(3)
+ @Nullable volatile GridDeploymentInfoMessage predDepInfo;
/** P2P unmarshalling future. */
- private IgniteInternalFuture<Void> p2pUnmarshalFut = new
GridFinishedFuture<>();
+ private volatile IgniteInternalFuture<Void> p2pUnmarshalFut = new
GridFinishedFuture<>();
/**
- * Required by {@link Externalizable}.
+ * Empty constructor for serialization purposes
*/
public GridMessageListenHandler() {
// No-op.
@@ -151,6 +146,10 @@ public class GridMessageListenHandler implements
GridContinuousHandler {
assert ctx != null;
assert ctx.config().isPeerClassLoadingEnabled();
+ // TODO : Remove this check after
https://issues.apache.org/jira/browse/IGNITE-28945
+ if (predDepInfo != null)
+ return;
+
if (topic != null)
topicBytes = U.marshal(ctx.marshaller(), topic);
@@ -166,9 +165,19 @@ public class GridMessageListenHandler implements
GridContinuousHandler {
if (dep == null)
throw new IgniteDeploymentCheckedException("Failed to deploy
message listener.");
- depInfo = new GridDeploymentInfoMessage(dep);
+ predDepInfo = new GridDeploymentInfoMessage(dep);
+ }
+
+ /** {@inheritDoc} */
+ @Override public void marshal(Marshaller marsh) throws
IgniteCheckedException {
+ /** Are marshaled in {@link #p2pMarshal(GridKernalContext)}. */
+ if (predDepInfo != null)
+ return;
+
+ if (topic != null)
+ topicBytes = marsh.marshal(topic);
- depEnabled = true;
+ predBytes = marsh.marshal(pred);
}
/** {@inheritDoc} */
@@ -177,8 +186,12 @@ public class GridMessageListenHandler implements
GridContinuousHandler {
assert ctx != null;
assert ctx.config().isPeerClassLoadingEnabled();
+ // TODO : Remove this check after
https://issues.apache.org/jira/browse/IGNITE-28945
+ if (pred != null)
+ return;
+
try {
- ClassLoader ldr = ctx.deploy().globalDeployment(depInfo, clsName,
nodeId).classLoader();
+ ClassLoader ldr = ctx.deploy().globalDeployment(predDepInfo,
clsName, nodeId).classLoader();
if (topicBytes != null)
topic = U.unmarshal(ctx, topicBytes, U.resolveClassLoader(ldr,
ctx.config()));
@@ -199,6 +212,21 @@ public class GridMessageListenHandler implements
GridContinuousHandler {
((GridFutureAdapter)p2pUnmarshalFut).onDone();
}
+ /** {@inheritDoc} */
+ @Override public void unmarshal(Marshaller marsh, ClassLoader clsLdr)
throws IgniteCheckedException {
+ /** Are unmarshaled in {@link #p2pUnmarshal(UUID, GridKernalContext)}.
*/
+ if (predDepInfo != null) {
+ p2pUnmarshalFut = new GridFutureAdapter<>();
+
+ return;
+ }
+
+ if (topicBytes != null)
+ topic = marsh.unmarshal(topicBytes, clsLdr);
+
+ pred = marsh.unmarshal(predBytes, clsLdr);
+ }
+
/** {@inheritDoc} */
@Override public GridContinuousBatch createBatch() {
return new GridContinuousBatchAdapter();
@@ -229,39 +257,6 @@ public class GridMessageListenHandler implements
GridContinuousHandler {
}
}
- /** {@inheritDoc} */
- @Override public void writeExternal(ObjectOutput out) throws IOException {
- out.writeBoolean(depEnabled);
-
- if (depEnabled) {
- U.writeByteArray(out, topicBytes);
- U.writeByteArray(out, predBytes);
- U.writeString(out, clsName);
- out.writeObject(depInfo);
- }
- else {
- out.writeObject(topic);
- out.writeObject(pred);
- }
- }
-
- /** {@inheritDoc} */
- @Override public void readExternal(ObjectInput in) throws IOException,
ClassNotFoundException {
- depEnabled = in.readBoolean();
-
- if (depEnabled) {
- p2pUnmarshalFut = new GridFutureAdapter<>();
- topicBytes = U.readByteArray(in);
- predBytes = U.readByteArray(in);
- clsName = U.readString(in);
- depInfo = (GridDeploymentInfoMessage)in.readObject();
- }
- else {
- topic = in.readObject();
- pred = (IgniteBiPredicate<UUID, Object>)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/GridDeploymentInfoMessage.java
b/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentInfoMessage.java
index 4a4f4e4a6a9..3c00a13799a 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentInfoMessage.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentInfoMessage.java
@@ -30,7 +30,7 @@ import
org.apache.ignite.plugin.extensions.communication.Message;
/**
* Deployment of classes, as it travels inside the messages carrying them.
*/
-public class GridDeploymentInfoMessage implements Message, GridDeploymentInfo,
Serializable {
+public final class GridDeploymentInfoMessage implements Message,
GridDeploymentInfo, Serializable {
/** */
private static final long serialVersionUID = 0L;
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 2eda9688193..7e0246b7897 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 @@ class GridAffinityUtils {
}
/**
- * 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 5f31963ba21..57843e86d1b 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.GridDeploymentInfoMessage;
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 final 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)
+ GridDeploymentInfoMessage depInfo;
/**
- * Required by {@link Externalizable}.
+ * Empty constructor for serialization purposes.
*/
public CacheContinuousQueryDeployableObject() {
// No-op.
@@ -93,20 +90,6 @@ class CacheContinuousQueryDeployableObject implements
Externalizable {
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 805ebda1e00..ac541a10a41 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,10 +17,6 @@
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;
import java.util.Collections;
@@ -51,6 +47,9 @@ import org.apache.ignite.events.CacheQueryExecutedEvent;
import org.apache.ignite.events.CacheQueryReadEvent;
import org.apache.ignite.internal.GridKernalContext;
import org.apache.ignite.internal.IgniteInternalFuture;
+import org.apache.ignite.internal.MarshallableMessage;
+import org.apache.ignite.internal.Marshalled;
+import org.apache.ignite.internal.Order;
import org.apache.ignite.internal.cluster.ClusterTopologyCheckedException;
import org.apache.ignite.internal.managers.communication.GridIoPolicy;
import org.apache.ignite.internal.managers.communication.MessageMarshalling;
@@ -83,6 +82,7 @@ import org.apache.ignite.internal.util.typedef.internal.U;
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 +96,7 @@ import static
org.apache.ignite.internal.processors.cache.query.continuous.Cache
/**
* Continuous query handler.
*/
-public class CacheContinuousQueryHandler<K, V> implements
GridContinuousHandler {
- /** */
- private static final long serialVersionUID = 0L;
-
+public final class CacheContinuousQueryHandler<K, V> implements
GridContinuousHandler, MarshallableMessage {
/** @see #IGNITE_CONTINUOUS_QUERY_BACKUP_ACK_THRESHOLD */
public static final int DFLT_CONTINUOUS_QUERY_BACKUP_ACK_THRESHOLD = 100;
@@ -131,8 +128,8 @@ public class CacheContinuousQueryHandler<K, V> implements
GridContinuousHandler
* Transformer implementation for processing received remote events.
* They are already transformed so we simply return transformed value for
event.
*/
- private transient IgniteClosure<CacheEntryEvent<? extends K, ? extends V>,
?> returnValTrans =
- new IgniteClosure<CacheEntryEvent<? extends K, ? extends V>, Object>()
{
+ private IgniteClosure<CacheEntryEvent<? extends K, ? extends V>, ?>
returnValTrans =
+ new IgniteClosure<>() {
@Override public Object apply(CacheEntryEvent<? extends K, ?
extends V> evt) {
assert evt.getKey() == null;
@@ -141,124 +138,153 @@ public class CacheContinuousQueryHandler<K, V>
implements GridContinuousHandler
};
/** Cache name. */
- private String cacheName;
+ @Order(0)
+ String cacheName;
/** Topic for ordered messages. */
- private Object topic;
+ @Marshalled("topicBytes")
+ Object topic;
+
+ /** Marshalled {@link #topic}. */
+ @Order(1)
+ byte[] topicBytes;
/** P2P unmarshalling future. */
- protected transient IgniteInternalFuture<Void> p2pUnmarshalFut = new
GridFinishedFuture<>();
+ protected volatile IgniteInternalFuture<Void> p2pUnmarshalFut = new
GridFinishedFuture<>();
/** Initialization future. */
- protected transient IgniteInternalFuture<Void> initFut;
+ protected IgniteInternalFuture<Void> initFut;
/** Local listener. */
- private transient CacheEntryUpdatedListener<K, V> locLsnr;
+ private CacheEntryUpdatedListener<K, V> locLsnr;
/** Remote filter. */
- private CacheEntryEventSerializableFilter<K, V> rmtFilter;
+ volatile CacheEntryEventSerializableFilter<K, V> rmtFilter;
- /** Deployable object for filter. */
- private CacheContinuousQueryDeployableObject rmtFilterDep;
+ /** Deployable object for {@link #rmtFilter}. Is {@code null} if no
external marsshalling used. */
+ @Order(2)
+ @Nullable volatile CacheContinuousQueryDeployableObject rmtFilterDep;
+
+ /** Marshalled {@link #rmtFilter} if {@link #rmtFilterDep} is {@code
null}. */
+ @Order(3)
+ @Nullable volatile byte[] rmtFilterBytes;
/** Remote filter factory. */
- private Factory<? extends CacheEntryEventFilter> rmtFilterFactory;
+ @Nullable volatile Factory<? extends CacheEntryEventFilter>
rmtFilterFactory;
+
+ /** Deployable object for {@link #rmtFilterFactory}. Is {@code null} if no
external marsshalling used. */
+ @Order(4)
+ volatile CacheContinuousQueryDeployableObject rmtFilterFactoryDep;
- /** Deployable object for filter factory. */
- private CacheContinuousQueryDeployableObject rmtFilterFactoryDep;
+ /** Marshalled {@link #rmtFilterFactory} if {@link #rmtFilterFactoryDep}
is {@code null}. */
+ @Order(5)
+ @Nullable volatile byte[] rmtFilterFactoryBytes;
/** Remote filter created by {@link #rmtFilterFactory}. */
- private transient CacheEntryEventFilter rmtFilterFromFactory;
+ private CacheEntryEventFilter rmtFilterFromFactory;
/** Event types for JCache API. */
- private byte types;
+ @Order(6)
+ byte types;
/** Remote transformer factory. */
- private Factory<? extends IgniteClosure<CacheEntryEvent<? extends K, ?
extends V>, ?>> rmtTransFactory;
+ volatile Factory<? extends IgniteClosure<CacheEntryEvent<? extends K, ?
extends V>, ?>> rmtTransFactory;
+
+ /** Deployable object for {@link #rmtTransFactory}. Is {@code null} if no
external marsshalling used. */
+ @Order(7)
+ volatile CacheContinuousQueryDeployableObject rmtTransFactoryDep;
- /** Deployable object for transformer factory. */
- private CacheContinuousQueryDeployableObject rmtTransFactoryDep;
+ /** Marshalled {@link #rmtTransFactory} if {@link #rmtTransFactoryDep} is
{@code null}. */
+ @Order(8)
+ @Nullable volatile byte[] rmtTransFactoryBytes;
/** Remote transformer created by {@link #rmtTransFactory}. */
- private transient IgniteClosure<CacheEntryEvent<? extends K, ? extends V>,
?> rmtTrans;
+ private IgniteClosure<CacheEntryEvent<? extends K, ? extends V>, ?>
rmtTrans;
/** Local listener for transformed events. */
- private transient EventListener<?> locTransLsnr;
+ private EventListener<?> locTransLsnr;
/** Internal flag. */
- private boolean internal;
+ @Order(9)
+ boolean internal;
/** Notify existing flag. */
- private boolean notifyExisting;
+ @Order(10)
+ boolean notifyExisting;
/** Old value required flag. */
- private boolean oldValRequired;
+ @Order(11)
+ boolean oldValRequired;
/** Synchronous flag. */
- private boolean sync;
+ @Order(12)
+ boolean sync;
/** Ignore expired events flag. */
- private boolean ignoreExpired;
+ @Order(13)
+ boolean ignoreExpired;
/** Task name hash code. */
- private int taskHash;
+ @Order(14)
+ int taskHash;
/** Whether to skip primary check for REPLICATED cache. */
- private transient boolean skipPrimaryCheck;
+ boolean skipPrimaryCheck;
/** */
- private transient boolean locOnly;
+ private boolean locOnly;
/** */
- private boolean keepBinary;
+ @Order(15)
+ boolean keepBinary;
/** */
- private transient ConcurrentMap<Integer,
CacheContinuousQueryPartitionRecovery> rcvs;
+ private ConcurrentMap<Integer, CacheContinuousQueryPartitionRecovery> rcvs;
/** */
- private transient ConcurrentMap<Integer, CacheContinuousQueryEventBuffer>
entryBufs;
+ private ConcurrentMap<Integer, CacheContinuousQueryEventBuffer> entryBufs;
/** */
- private transient CacheContinuousQueryAcknowledgeBuffer ackBuf;
+ private CacheContinuousQueryAcknowledgeBuffer ackBuf;
/** */
- private transient int cacheId;
+ private int cacheId;
/** */
- private transient volatile Map<Integer, Long> initUpdCntrs;
+ private volatile Map<Integer, Long> initUpdCntrs;
/** */
- private transient volatile Map<UUID, Map<Integer, Long>>
initUpdCntrsPerNode;
+ private volatile Map<UUID, Map<Integer, Long>> 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<Integer, Long> locInitUpdCntrs;
+ private volatile Map<Integer, Long> 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.
@@ -411,11 +437,6 @@ public class CacheContinuousQueryHandler<K, V> implements
GridContinuousHandler
this.locOnly = locOnly;
}
- /** @return {@code True} if handler are local only, {@code false}
otherwise. */
- public boolean localOnly() {
- return locOnly;
- }
-
/**
* @param taskHash Task hash.
*/
@@ -942,8 +963,6 @@ public class CacheContinuousQueryHandler<K, V> implements
GridContinuousHandler
* @throws IgniteCheckedException In case of error.
*/
void waitTopologyFuture(GridKernalContext ctx) throws
IgniteCheckedException {
- GridCacheContext<K, V> cctx = cacheContext(ctx);
-
AffinityTopologyVersion topVer = initTopVer;
cacheContext(ctx).shared().exchange().affinityReadyFuture(topVer).get();
@@ -1392,73 +1411,89 @@ public class CacheContinuousQueryHandler<K, V>
implements GridContinuousHandler
assert ctx != null;
assert ctx.config().isPeerClassLoadingEnabled();
- if (requiresDeployment(rmtFilter))
+ // TODO : Remove this check after
https://issues.apache.org/jira/browse/IGNITE-28945
+ if (rmtFilterDep != null || rmtFilterFactoryDep != null ||
rmtTransFactoryDep != null)
+ return;
+
+ /**
+ * 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);
}
+ /** {@inheritDoc} */
+ @Override public void marshal(Marshaller marsh) throws
IgniteCheckedException {
+ if (rmtFilter != null && rmtFilterDep == null)
+ rmtFilterBytes = marsh.marshal(rmtFilter);
+
+ if (rmtFilterFactory != null && rmtFilterFactoryDep == null)
+ rmtFilterFactoryBytes = marsh.marshal(rmtFilterFactory);
+
+ if (rmtTransFactory != null && rmtTransFactoryDep == null)
+ 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();
- if (rmtFilterDep != null)
- rmtFilter = p2pUnmarshal(rmtFilterDep, nodeId, ctx);
+ // TODO : Remove this check after
https://issues.apache.org/jira/browse/IGNITE-28945
+ if (rmtFilter != null || rmtFilterFactory != null || rmtTransFactory
!= null)
+ return;
- if (rmtFilterFactoryDep != null)
- rmtFilterFactory = p2pUnmarshal(rmtFilterFactoryDep, nodeId, ctx);
+ try {
+ if (rmtFilterDep != null)
+ rmtFilter = rmtFilterDep.unmarshal(nodeId, ctx);
- if (rmtTransFactoryDep != null)
- rmtTransFactory = p2pUnmarshal(rmtTransFactoryDep, nodeId, ctx);
+ if (rmtFilterFactoryDep != null)
+ rmtFilterFactory = rmtFilterFactoryDep.unmarshal(nodeId, ctx);
- if (!p2pUnmarshalFut.isDone())
- ((GridFutureAdapter)p2pUnmarshalFut).onDone();
- }
+ if (rmtTransFactoryDep != null)
+ rmtTransFactory = rmtTransFactoryDep.unmarshal(nodeId, ctx);
- /**
- * @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);
- }
-
- /**
- * @param depObj Deployable object to unmarshal.
- * @param nodeId Sender node Id.
- * @param ctx Kernal context.
- * @param <T> Result type.
- * @return Unmarshalled object.
- * @throws IgniteCheckedException In case of unmarshalling failures.
- */
- protected <T> 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);
+ if (!p2pUnmarshalFut.isDone())
+ ((GridFutureAdapter)p2pUnmarshalFut).onDone();
+ }
+ catch (IgniteCheckedException e) {
+ ((GridFutureAdapter<?>)p2pUnmarshalFut).onDone(e);
- throw e;
- }
- catch (ExceptionInInitializerError e) {
- IgniteCheckedException err = new
IgniteCheckedException("Failed to unmarshal deployable object.", e);
+ throw e;
+ }
+ catch (ExceptionInInitializerError e) {
+ IgniteCheckedException err = new IgniteCheckedException("Failed to
unmarshal deployable object.", e);
- ((GridFutureAdapter)p2pUnmarshalFut).onDone(err);
+ ((GridFutureAdapter<?>)p2pUnmarshalFut).onDone(err);
- throw err;
- }
+ throw err;
}
- else
- return null;
}
/** {@inheritDoc} */
@@ -1554,87 +1589,6 @@ public class CacheContinuousQueryHandler<K, V>
implements GridContinuousHandler
return S.toString(CacheContinuousQueryHandler.class, this);
}
- /** {@inheritDoc} */
- @Override public void writeExternal(ObjectOutput out) throws IOException {
- U.writeString(out, cacheName);
- out.writeObject(topic);
-
- writeDeployable(out, rmtFilter, rmtFilterDep);
-
- out.writeBoolean(internal);
- out.writeBoolean(notifyExisting);
- out.writeBoolean(oldValRequired);
- out.writeBoolean(sync);
- out.writeBoolean(ignoreExpired);
- out.writeInt(taskHash);
- out.writeBoolean(keepBinary);
-
- writeDeployable(out, rmtFilterFactory, rmtFilterFactoryDep);
-
- out.writeByte(types);
-
- writeDeployable(out, rmtTransFactory, rmtTransFactoryDep);
- }
-
- /** {@inheritDoc} */
- @Override public void readExternal(ObjectInput in) throws IOException,
ClassNotFoundException {
- cacheName = U.readString(in);
- topic = in.readObject();
-
- boolean b = in.readBoolean();
-
- if (b) {
- rmtFilterDep =
(CacheContinuousQueryDeployableObject)in.readObject();
-
- p2pUnmarshalFut = new GridFutureAdapter<>();
- }
- else
- rmtFilter = (CacheEntryEventSerializableFilter<K,
V>)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<? extends
IgniteClosure<CacheEntryEvent<? extends K, ? extends V>, ?>>)in.readObject();
-
- cacheId = CU.cacheId(cacheName);
- }
-
- /** */
- private static void writeDeployable(ObjectOutput out, Object obj,
CacheContinuousQueryDeployableObject dep) throws IOException {
- boolean b = dep != null;
-
- out.writeBoolean(b);
-
- out.writeObject(b ? dep : obj);
- }
-
/**
* @param ctx Kernal context.
* @return Cache context.
@@ -1812,9 +1766,4 @@ public class CacheContinuousQueryHandler<K, V> implements
GridContinuousHandler
Map<Integer, CacheContinuousQueryEventBuffer>
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/ContinousRoutineDiscoveryData.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryData.java
new file mode 100644
index 00000000000..44f0a2f231a
--- /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;
+
+/** Continous routine Discovery data. */
+public final class ContinousRoutineDiscoveryData implements Message {
+ /** Node ID. */
+ @Order(0)
+ UUID nodeId;
+
+ /** Items. */
+ @GridToStringInclude
+ @Order(1)
+ Collection<ContinousRoutineDiscoveryDataItem> items;
+
+ /** */
+ @Order(2)
+ Map<UUID, Map<UUID, ContinousRoutineLocalInfo>> clientInfos;
+
+ /** Empty constructor for serialization purposes. */
+ public ContinousRoutineDiscoveryData() {
+ // No-op.
+ }
+
+ /**
+ * @param nodeId Node ID.
+ * @param clientInfos Client information.
+ */
+ ContinousRoutineDiscoveryData(UUID nodeId, Map<UUID, Map<UUID,
ContinousRoutineLocalInfo>> 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 00000000000..650229ae0ff
--- /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;
+
+/** Continous routine Discovery data item. */
+public final class ContinousRoutineDiscoveryDataItem implements Message {
+ /** */
+ @Order(0)
+ UUID routineId;
+
+ /** */
+ @Marshalled("prjPredBytes")
+ IgnitePredicate<ClusterNode> 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<ClusterNode> 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 00000000000..8bcbc4f8bd9
--- /dev/null
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineLocalInfo.java
@@ -0,0 +1,133 @@
+/*
+ * 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.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;
+
+/** Continous routine local info Discovery data. */
+@UseBinaryMarshaller
+public final class ContinousRoutineLocalInfo implements Message,
GridContinuousProcessor.RoutineInfo {
+ /** Source node id. */
+ @Order(0)
+ UUID nodeId;
+
+ /** Projection predicate. */
+ @Marshalled("prjPredBytes")
+ IgnitePredicate<ClusterNode> 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<ClusterNode> 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 207a8f4fda9..8c5a920ffef 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,45 +17,53 @@
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 final class ContinuousRoutineInfo implements Message {
/** */
+ @Order(0)
UUID srcNodeId;
/** */
- final UUID routineId;
+ @Order(1)
+ UUID routineId;
/** */
- final byte[] hnd;
+ @Order(2)
+ GridContinuousHandler 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.
* @param routineId Routine ID.
- * @param hnd Marshalled handler.
+ * @param hnd Handler.
* @param nodeFilter Marshalled node filter.
* @param bufSize Handler buffer size.
* @param interval Time interval.
@@ -64,7 +72,7 @@ class ContinuousRoutineInfo implements Serializable {
ContinuousRoutineInfo(
UUID srcNodeId,
UUID routineId,
- byte[] hnd,
+ GridContinuousHandler hnd,
byte[] nodeFilter,
int bufSize,
long interval,
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesCommonDiscoveryData.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesCommonDiscoveryData.java
index d29de89b9b1..7b6ad6a3396 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesCommonDiscoveryData.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesCommonDiscoveryData.java
@@ -17,19 +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 ContinuousRoutinesCommonDiscoveryData implements Serializable {
+public final class ContinuousRoutinesCommonDiscoveryData implements Message {
/** */
- private static final long serialVersionUID = 0L;
+ @Order(0)
+ List<ContinuousRoutineInfo> startedRoutines;
- /** */
- final List<ContinuousRoutineInfo> startedRoutines;
+ /** Empty constructor for serialization purposes. */
+ public ContinuousRoutinesCommonDiscoveryData() {
+ // No-op.
+ }
/**
* @param startedRoutines Routines started in cluster.
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 9be6ef8e07e..ec40f56ca1f 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,19 +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 final class ContinuousRoutinesJoiningNodeDiscoveryData implements
Message {
/** */
- private static final long serialVersionUID = 0L;
+ @Order(0)
+ List<ContinuousRoutineInfo> startedRoutines;
- /** */
- final List<ContinuousRoutineInfo> startedRoutines;
+ /** Empty constructor for serialization purposes. */
+ public ContinuousRoutinesJoiningNodeDiscoveryData() {
+ // No-op.
+ }
/**
* @param startedRoutines Routines registered on nodes, to be started in
cluster.
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 b1a3812f612..03ca0defa6d 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 111d3b74049..7570b94f37e 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,12 +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.io.Serializable;
-import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.HashMap;
@@ -80,7 +74,6 @@ import
org.apache.ignite.internal.util.future.GridFinishedFuture;
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 +118,10 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
public static final String CQ_SYS_VIEW_DESC = "Continuous queries";
/** Local infos. */
- private final ConcurrentMap<UUID, LocalRoutineInfo> locInfos = new
ConcurrentHashMap<>();
+ private final ConcurrentMap<UUID, ContinousRoutineLocalInfo> locInfos =
new ConcurrentHashMap<>();
/** Local infos. */
- private final ConcurrentMap<UUID, Map<UUID, LocalRoutineInfo>> clientInfos
= new ConcurrentHashMap<>();
+ private final ConcurrentMap<UUID, Map<UUID, ContinousRoutineLocalInfo>>
clientInfos = new ConcurrentHashMap<>();
/** Remote infos. */
private final ConcurrentMap<UUID, RemoteRoutineInfo> rmtInfos = new
ConcurrentHashMap<>();
@@ -336,12 +329,12 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
}
/** */
- Map<UUID, LocalRoutineInfo> localRoutineInfos() {
+ Map<UUID, ContinousRoutineLocalInfo> localRoutineInfos() {
return Collections.unmodifiableMap(locInfos);
}
/** */
- Map<UUID, Map<UUID, LocalRoutineInfo>> clientRoutineInfos() {
+ Map<UUID, Map<UUID, ContinousRoutineLocalInfo>> clientRoutineInfos() {
return Collections.unmodifiableMap(clientInfos);
}
@@ -408,7 +401,7 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
return;
}
- Serializable data = getDiscoveryData(dataBag.joiningNodeId());
+ ContinousRoutineDiscoveryData data =
getDiscoveryData(dataBag.joiningNodeId());
if (data != null)
dataBag.addJoiningNodeData(CONTINUOUS_PROC.ordinal(), data);
@@ -422,7 +415,7 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
return;
}
- Serializable data = getDiscoveryData(dataBag.joiningNodeId());
+ ContinousRoutineDiscoveryData data =
getDiscoveryData(dataBag.joiningNodeId());
if (data != null)
dataBag.addNodeSpecificData(CONTINUOUS_PROC.ordinal(), data);
@@ -431,7 +424,7 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
/**
* @param joiningNodeId Joining node id.
*/
- private Serializable getDiscoveryData(UUID joiningNodeId) {
+ private @Nullable ContinousRoutineDiscoveryData getDiscoveryData(UUID
joiningNodeId) {
if (log.isDebugEnabled()) {
log.debug("collectDiscoveryData [node=" + joiningNodeId +
", loc=" + ctx.localNodeId() +
@@ -441,26 +434,22 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
}
if (!joiningNodeId.equals(ctx.localNodeId()) || !locInfos.isEmpty()) {
- Map<UUID, Map<UUID, LocalRoutineInfo>> clientInfos0 =
copyClientInfos(clientInfos);
+ Map<UUID, Map<UUID, ContinousRoutineLocalInfo>> clientInfos0 =
copyClientInfos(clientInfos);
if (joiningNodeId.equals(ctx.localNodeId()) &&
ctx.discovery().localNode().isClient()) {
- Map<UUID, LocalRoutineInfo> infos = copyLocalInfos(locInfos);
+ Map<UUID, ContinousRoutineLocalInfo> 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<UUID, LocalRoutineInfo> e : locInfos.entrySet()) {
+ for (Map.Entry<UUID, ContinousRoutineLocalInfo> 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();
-
- data.addItem(new DiscoveryDataItem(routineId,
+ data.addItem(new ContinousRoutineDiscoveryDataItem(routineId,
info.prjPred,
info.hnd,
info.bufSize,
@@ -477,13 +466,13 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
/**
* @param clientInfos Client infos.
*/
- private Map<UUID, Map<UUID, LocalRoutineInfo>> copyClientInfos(Map<UUID,
Map<UUID, LocalRoutineInfo>> clientInfos) {
- Map<UUID, Map<UUID, LocalRoutineInfo>> res =
U.newHashMap(clientInfos.size());
+ private Map<UUID, Map<UUID, ContinousRoutineLocalInfo>>
copyClientInfos(Map<UUID, Map<UUID, ContinousRoutineLocalInfo>> clientInfos) {
+ Map<UUID, Map<UUID, ContinousRoutineLocalInfo>> res =
U.newHashMap(clientInfos.size());
- for (Map.Entry<UUID, Map<UUID, LocalRoutineInfo>> e :
clientInfos.entrySet()) {
- Map<UUID, LocalRoutineInfo> cp = U.newHashMap(e.getValue().size());
+ for (Map.Entry<UUID, Map<UUID, ContinousRoutineLocalInfo>> e :
clientInfos.entrySet()) {
+ Map<UUID, ContinousRoutineLocalInfo> cp =
U.newHashMap(e.getValue().size());
- for (Map.Entry<UUID, LocalRoutineInfo> e0 :
e.getValue().entrySet())
+ for (Map.Entry<UUID, ContinousRoutineLocalInfo> e0 :
e.getValue().entrySet())
cp.put(e0.getKey(), e0.getValue());
res.put(e.getKey(), cp);
@@ -495,10 +484,10 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
/**
* @param locInfos Locale infos.
*/
- private Map<UUID, LocalRoutineInfo> copyLocalInfos(Map<UUID,
LocalRoutineInfo> locInfos) {
- Map<UUID, LocalRoutineInfo> res = U.newHashMap(locInfos.size());
+ private Map<UUID, ContinousRoutineLocalInfo> copyLocalInfos(Map<UUID,
ContinousRoutineLocalInfo> locInfos) {
+ Map<UUID, ContinousRoutineLocalInfo> res =
U.newHashMap(locInfos.size());
- for (Map.Entry<UUID, LocalRoutineInfo> e : locInfos.entrySet())
+ for (Map.Entry<UUID, ContinousRoutineLocalInfo> e :
locInfos.entrySet())
res.put(e.getKey(), e.getValue());
return res;
@@ -547,10 +536,10 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
}
}
else {
- Map<UUID, DiscoveryData> nodeSpecData = data.nodeSpecificData();
+ Map<UUID, ContinousRoutineDiscoveryData> nodeSpecData =
data.nodeSpecificData();
if (nodeSpecData != null) {
- for (DiscoveryData val : nodeSpecData.values())
+ for (ContinousRoutineDiscoveryData val : nodeSpecData.values())
onDiscoveryDataReceivedMutable(val);
}
}
@@ -562,37 +551,37 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
*
* @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<UUID, Map<UUID, LocalRoutineInfo>> entry :
data.clientInfos.entrySet()) {
+ for (Map.Entry<UUID, Map<UUID, ContinousRoutineLocalInfo>> entry :
data.clientInfos.entrySet()) {
UUID clientNodeId = entry.getKey();
if (!ctx.localNodeId().equals(clientNodeId)) {
- Map<UUID, LocalRoutineInfo> clientRoutineMap =
entry.getValue();
+ Map<UUID, ContinousRoutineLocalInfo> clientRoutineMap =
entry.getValue();
- for (Map.Entry<UUID, LocalRoutineInfo> e :
clientRoutineMap.entrySet()) {
+ for (Map.Entry<UUID, ContinousRoutineLocalInfo> 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<UUID, LocalRoutineInfo> map =
+ Map<UUID, ContinousRoutineLocalInfo> map =
clientInfos.computeIfAbsent(clientNodeId, k -> new
HashMap<>());
map.putAll(entry.getValue());
@@ -628,23 +617,8 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
return;
}
- GridContinuousHandler hnd;
-
- try {
- hnd = U.unmarshal(marsh, routineInfo.hnd,
U.resolveClassLoader(ctx.config()));
- }
- catch (IgniteCheckedException e) {
- U.error(log, "Failed to unmarshal continuous routine handler [" +
- "routineId=" + routineInfo.routineId +
- ", srcNodeId=" + routineInfo.srcNodeId + ']', e);
-
- ctx.failure().process(new
FailureContext(FailureType.CRITICAL_ERROR, e));
-
- return;
- }
-
registerHandlerOnJoin(routineInfo.srcNodeId, routineInfo.routineId,
nodeFilter,
- hnd, routineInfo.bufSize, routineInfo.interval,
routineInfo.autoUnsubscribe);
+ routineInfo.hnd, routineInfo.bufSize, routineInfo.interval,
routineInfo.autoUnsubscribe);
}
/**
@@ -782,7 +756,7 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
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(
@@ -822,14 +796,13 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
long interval,
boolean autoUnsubscribe)
throws IgniteCheckedException {
- byte[] hndBytes = marsh.marshal(hnd);
byte[] filterBytes = nodeFilter != null ? marsh.marshal(nodeFilter) :
null;
return new ContinuousRoutineInfo(
srcNodeId,
routineId,
- hndBytes,
+ hnd,
filterBytes,
bufSize,
interval,
@@ -858,15 +831,12 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
// 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 LocalRoutineInfo(ctx.localNodeId(), prjPred, hnd, bufSize,
interval, autoUnsubscribe));
+ new ContinousRoutineLocalInfo(ctx.localNodeId(), prjPred, hnd,
bufSize, interval, autoUnsubscribe));
if (locOnly) {
try {
@@ -986,8 +956,6 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
hnd.p2pMarshal(ctx);
}
- reqData.hndBytes = U.marshal(marsh, hnd);
-
if (nodeFilter != null)
reqData.nodeFilterBytes = U.marshal(marsh, nodeFilter);
@@ -1066,7 +1034,7 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
boolean stop = false;
// Unregister routine locally.
- LocalRoutineInfo routine = locInfos.remove(routineId);
+ ContinousRoutineLocalInfo routine = locInfos.remove(routineId);
if (routine != null) {
stop = true;
@@ -1131,7 +1099,7 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
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);
@@ -1256,7 +1224,7 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
unregisterRemote(e.getKey());
}
- for (LocalRoutineInfo routine : locInfos.values())
+ for (ContinousRoutineLocalInfo routine : locInfos.values())
routine.hnd.onClientDisconnected();
rmtInfos.clear();
@@ -1323,7 +1291,7 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
unregisterRemote(routineId);
}
- for (Map<UUID, LocalRoutineInfo> clientInfo : clientInfos.values()) {
+ for (Map<UUID, ContinousRoutineLocalInfo> clientInfo :
clientInfos.values()) {
if (clientInfo.remove(msg.routineId()) != null)
break;
}
@@ -1359,17 +1327,13 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
data.nodeFilter = U.unmarshal(marsh, data.nodeFilterBytes,
ctx.deploy().classLoader(data.depInfo, data.clsName, sndId));
- if (data.hndBytes != null) {
- data.hnd = U.unmarshal(marsh, data.hndBytes,
U.resolveClassLoader(ctx.config()));
+ if (ctx.config().isPeerClassLoadingEnabled())
+ data.hnd.p2pUnmarshal(sndId, ctx);
- if (ctx.config().isPeerClassLoadingEnabled())
- data.hnd.p2pUnmarshal(sndId, ctx);
+ if (data.keepBinary) {
+ assert data.hnd instanceof CacheContinuousQueryHandler : data.hnd;
- if (data.keepBinary) {
- assert data.hnd instanceof CacheContinuousQueryHandler :
data.hnd;
-
- ((CacheContinuousQueryHandler<?, ?>)data.hnd).keepBinary(true);
- }
+ ((CacheContinuousQueryHandler<?, ?>)data.hnd).keepBinary(true);
}
}
@@ -1409,17 +1373,17 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
GridContinuousHandler hnd = data.handler();
if (node.isClient()) {
- Map<UUID, LocalRoutineInfo> clientRoutineMap =
clientInfos.get(node.id());
+ Map<UUID, ContinousRoutineLocalInfo> clientRoutineMap =
clientInfos.get(node.id());
if (clientRoutineMap == null) {
clientRoutineMap = new HashMap<>();
- Map<UUID, LocalRoutineInfo> old = clientInfos.put(node.id(),
clientRoutineMap);
+ Map<UUID, ContinousRoutineLocalInfo> 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(),
@@ -1457,7 +1421,7 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
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) {
@@ -1503,7 +1467,7 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
ContinuousRoutineInfo routineInfo = new ContinuousRoutineInfo(snd.id(),
msg.routineId(),
- reqData.hndBytes,
+ reqData.hnd,
reqData.nodeFilterBytes,
reqData.bufferSize(),
reqData.interval(),
@@ -1642,7 +1606,7 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
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);
@@ -1804,7 +1768,7 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
@SuppressWarnings("TooBroadScope")
private void unregisterRemote(UUID routineId) {
RemoteRoutineInfo remote;
- LocalRoutineInfo loc;
+ ContinousRoutineLocalInfo loc;
stopLock.lock();
@@ -2006,100 +1970,6 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
boolean delayedRegister();
}
- /**
- * Local routine info.
- */
- public static class LocalRoutineInfo implements Serializable, RoutineInfo {
- /** */
- private static final long serialVersionUID = 0L;
-
- /** Source node id. */
- private final UUID nodeId;
-
- /** Projection predicate. */
- private final IgnitePredicate<ClusterNode> prjPred;
-
- /** Continuous routine handler. */
- private final GridContinuousHandler hnd;
-
- /** Buffer size. */
- private final int bufSize;
-
- /** Time interval. */
- private final long interval;
-
- /** Automatic unsubscribe flag. */
- private 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<ClusterNode> 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.
*/
@@ -2321,157 +2191,6 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
}
}
- /**
- * Discovery data.
- */
- private static class DiscoveryData implements Externalizable {
- /** */
- private static final long serialVersionUID = 0L;
-
- /** Node ID. */
- private UUID nodeId;
-
- /** Items. */
- @GridToStringInclude
- private Collection<DiscoveryDataItem> items;
-
- /** */
- private Map<UUID, Map<UUID, LocalRoutineInfo>> clientInfos;
-
- /**
- * Required by {@link Externalizable}.
- */
- public DiscoveryData() {
- // No-op.
- }
-
- /**
- * @param nodeId Node ID.
- * @param clientInfos Client information.
- */
- DiscoveryData(UUID nodeId, Map<UUID, Map<UUID, LocalRoutineInfo>>
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 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);
- }
- }
-
- /**
- * Discovery data item.
- */
- private static class DiscoveryDataItem implements Externalizable {
- /** */
- private static final long serialVersionUID = 0L;
-
- /** Consume ID. */
- private UUID routineId;
-
- /** Projection predicate. */
- private IgnitePredicate<ClusterNode> 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<ClusterNode> 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<ClusterNode>)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.
*/
@@ -2556,7 +2275,7 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
Map<Integer, Long> 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/main/java/org/apache/ignite/internal/processors/continuous/StartRequestData.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/StartRequestData.java
index 595aaa56c58..2d52f2dfe0b 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/StartRequestData.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/StartRequestData.java
@@ -44,11 +44,8 @@ public class StartRequestData implements Message {
GridDeploymentInfoMessage depInfo;
/** Handler, restored by the processor reading this request. */
- GridContinuousHandler hnd;
-
- /** Serialized handler. */
@Order(3)
- byte[] hndBytes;
+ GridContinuousHandler hnd;
/** Buffer size. */
@Order(4)
diff --git
a/modules/core/src/main/java/org/apache/ignite/spi/IgniteSpiAdapter.java
b/modules/core/src/main/java/org/apache/ignite/spi/IgniteSpiAdapter.java
index 850b196a747..79abc47ac7f 100644
--- a/modules/core/src/main/java/org/apache/ignite/spi/IgniteSpiAdapter.java
+++ b/modules/core/src/main/java/org/apache/ignite/spi/IgniteSpiAdapter.java
@@ -81,7 +81,7 @@ public abstract class IgniteSpiAdapter implements IgniteSpi {
protected IgniteLogger log;
/** Ignite instance. */
- protected Ignite ignite;
+ protected IgniteEx ignite;
/** Ignite instance name. */
protected String igniteInstanceName;
@@ -268,7 +268,7 @@ public abstract class IgniteSpiAdapter implements IgniteSpi
{
*/
@IgniteInstanceResource
protected void injectResources(Ignite ignite) {
- this.ignite = ignite;
+ this.ignite = (IgniteEx)ignite;
if (ignite != null && igniteInstanceName == null)
igniteInstanceName = ignite.name();
diff --git
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/DiscoveryDataBag.java
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/DiscoveryDataBag.java
index a2f3af751af..df44689dff8 100644
---
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/DiscoveryDataBag.java
+++
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/DiscoveryDataBag.java
@@ -280,14 +280,6 @@ public class DiscoveryDataBag {
commonData.put(cmpId, data);
}
- /**
- * @param cmpId Component ID.
- * @param data Serializable data.
- */
- public void addNodeSpecificData(Integer cmpId, Serializable data) {
- addNodeSpecificData(cmpId, new SerializableDataBagItemWrapper(data));
- }
-
/**
* @param cmpId Component ID.
* @param data Message data.
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 b894f0b8772..e6d27b3bb48 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
@@ -864,7 +864,7 @@ class ClientImpl extends TcpDiscoveryImpl {
}
/**
- * 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 e84c716c8f2..89cc89e3359 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
@@ -1672,7 +1672,7 @@ class ServerImpl extends TcpDiscoveryImpl {
}
/**
- * 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.
@@ -1693,7 +1693,7 @@ class ServerImpl extends TcpDiscoveryImpl {
}
/**
- * 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/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoverySpi.java
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoverySpi.java
index 9c52f9504be..b92ca657ad7 100644
---
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoverySpi.java
+++
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoverySpi.java
@@ -54,6 +54,7 @@ import org.apache.ignite.cluster.ClusterNode;
import org.apache.ignite.configuration.AddressResolver;
import org.apache.ignite.configuration.IgniteConfiguration;
import org.apache.ignite.failure.FailureContext;
+import org.apache.ignite.failure.FailureType;
import org.apache.ignite.internal.IgniteEx;
import org.apache.ignite.internal.IgniteInterruptedCheckedException;
import
org.apache.ignite.internal.managers.communication.UnknownMessageException;
@@ -1888,6 +1889,18 @@ public class TcpDiscoverySpi extends IgniteSpiAdapter
implements IgniteDiscovery
if (msg != null && sslMsgPattern.matcher(msg).matches())
streamCorruptedCause.initCause(new SSLException("Detected
SSL alert in StreamCorruptedException"));
}
+
+ if (X.hasCause(e, ClassNotFoundException.class)) {
+ LT.error(log, e, "Failed to read message due to an unknown
class to unmarshal received. Unable to " +
+ "process the Discovery protocol. Stopping the Discovery
SPI and invoking the failure handler. " +
+ "RmtAddr=" + sock.getRemoteSocketAddress() + ", rmtPort="
+ sock.getPort() + ']');
+
+ ignite.context().failure().process(new
FailureContext(FailureType.CRITICAL_ERROR, e));
+
+ // Prevents following cycling attempts to reconnect and logs
flooding.
+ spiStop();
+ }
+
throw e;
}
finally {
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 736b1742bef..317076bc92a 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.cache.query.ContinuousQuery;
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 class CacheContinuousQueryEntriesExpireTest extends
GridCommonAbstractTes
for (int i = 0; i < 1_000; i++)
cache.put(i, i);
- ConcurrentMap<UUID, GridContinuousProcessor.LocalRoutineInfo> locInfos
=
+ ConcurrentMap<UUID, ContinousRoutineLocalInfo> locInfos =
getFieldValue(srv1.context().continuous(), "locInfos");
assertEquals(1, locInfos.size());
diff --git
a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/ContinuousQueryRemoteFilterMissingInClassPathSelfTest.java
b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/ContinuousQueryRemoteFilterMissingInClassPathSelfTest.java
index 62873ce51fd..b2b9a27d605 100644
---
a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/ContinuousQueryRemoteFilterMissingInClassPathSelfTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/ContinuousQueryRemoteFilterMissingInClassPathSelfTest.java
@@ -28,7 +28,6 @@ import javax.cache.event.CacheEntryListenerException;
import javax.cache.event.CacheEntryUpdatedListener;
import org.apache.ignite.Ignite;
import org.apache.ignite.IgniteCache;
-import org.apache.ignite.IgniteCheckedException;
import org.apache.ignite.IgniteLogger;
import org.apache.ignite.cache.CacheEntryEventSerializableFilter;
import org.apache.ignite.cache.CacheMode;
@@ -37,13 +36,14 @@ import org.apache.ignite.configuration.CacheConfiguration;
import org.apache.ignite.configuration.IgniteConfiguration;
import org.apache.ignite.internal.util.typedef.X;
import org.apache.ignite.testframework.GridStringLogger;
-import org.apache.ignite.testframework.GridTestUtils;
import org.apache.ignite.testframework.ListeningTestLogger;
import org.apache.ignite.testframework.LogListener;
import org.apache.ignite.testframework.config.GridTestProperties;
import org.apache.ignite.testframework.junits.common.GridCommonAbstractTest;
import org.junit.Test;
+import static
org.apache.ignite.testframework.GridTestUtils.assertThrowsWithCause;
+
/**
*
*/
@@ -110,21 +110,8 @@ public class
ContinuousQueryRemoteFilterMissingInClassPathSelfTest extends GridC
* @throws Exception If fail.
*/
@Test
- public void testClientJoinsMissingClassWarning() throws Exception {
- setExternalLoader = true;
- Ignite ignite0 = startGrid(1);
-
- executeContinuousQuery(ignite0.cache(DEFAULT_CACHE_NAME));
-
- log = new GridStringLogger();
- setExternalLoader = false;
-
- startClientGrid(2);
-
- String logStr = log.toString();
-
- assertTrue(logStr.contains("Failed to unmarshal continuous query
remote filter on client node. " +
- "Can be ignored.") || logStr.contains("Failed to unmarshal
continuous routine handler"));
+ public void testClientJoinsMissingClass() throws Exception {
+ doTestNodeJoinsWithNoClassLoaderForContinuousQuery(false);
}
/**
@@ -151,6 +138,11 @@ public class
ContinuousQueryRemoteFilterMissingInClassPathSelfTest extends GridC
*/
@Test
public void testServerJoinsMissingClassException() throws Exception {
+ doTestNodeJoinsWithNoClassLoaderForContinuousQuery(true);
+ }
+
+ /** */
+ private void doTestNodeJoinsWithNoClassLoaderForContinuousQuery(boolean
server) throws Exception {
setExternalLoader = true;
Ignite ignite0 = startGrid(1);
@@ -160,19 +152,30 @@ public class
ContinuousQueryRemoteFilterMissingInClassPathSelfTest extends GridC
log = listeningLog;
- LogListener lsnr = LogListener.matches(logStr ->
- logStr.contains("class org.apache.ignite.IgniteCheckedException: "
+
- "Failed to find class with given class loader for
unmarshalling")
- || logStr.contains("Failed to unmarshal continuous routine
handler"
- )).build();
+ LogListener lsnr1 = LogListener.matches("Failed to initialize a
continuous query").build();
+ LogListener lsnr2 = LogListener.matches("ClassNotFoundException: " +
EXT_FILTER_CLASS).build();
- listeningLog.registerListener(lsnr);
+ listeningLog.registerListener(lsnr1);
+ listeningLog.registerListener(lsnr2);
setExternalLoader = false;
- GridTestUtils.assertThrows(log, () -> startGrid(2),
IgniteCheckedException.class, "Failed to start");
+ if (server) {
+ assertThrowsWithCause(() -> startGrid(2),
ClassNotFoundException.class);
- assertTrue(lsnr.check());
+ assertTrue(lsnr1.check());
+ assertTrue(lsnr2.check());
+ }
+ else {
+ startClientGrid(2);
+
+ /**
+ * Client successfuly starts because continous query isn't
actually deployed on a client. It is filtered out
+ * by {@link AttributeNodeFilter} with {@link
IgniteNodeAttributes#ATTR_CLIENT_MODE}.
+ **/
+ assertFalse(lsnr1.check());
+ assertFalse(lsnr2.check());
+ }
}
/**
diff --git
a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/DiscoveryDataDeserializationFailureHanderTest.java
b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/DiscoveryDataDeserializationFailureHanderTest.java
index e52d06fa0d3..f0a0c281782 100644
---
a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/DiscoveryDataDeserializationFailureHanderTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/DiscoveryDataDeserializationFailureHanderTest.java
@@ -57,6 +57,7 @@ public class DiscoveryDataDeserializationFailureHanderTest
extends GridCommonAbs
@Test
public void testFailureHander() throws Exception {
CountDownLatch latch = new CountDownLatch(1);
+
failureHnd = new TestFailureHandler(latch);
Ignite node1 = startGrid(1);
diff --git
a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/GridCacheContinuousQueryNodesFilteringTest.java
b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/GridCacheContinuousQueryNodesFilteringTest.java
index 7688a7b12ab..7efdad76be4 100644
---
a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/GridCacheContinuousQueryNodesFilteringTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/GridCacheContinuousQueryNodesFilteringTest.java
@@ -34,6 +34,7 @@ import org.apache.ignite.configuration.CacheConfiguration;
import org.apache.ignite.configuration.IgniteConfiguration;
import org.apache.ignite.failure.FailureHandler;
import org.apache.ignite.failure.TestFailureHandler;
+import org.apache.ignite.internal.util.typedef.X;
import org.apache.ignite.lang.IgnitePredicate;
import org.apache.ignite.testframework.GridStringLogger;
import org.apache.ignite.testframework.junits.common.GridCommonAbstractTest;
@@ -74,9 +75,17 @@ public class GridCacheContinuousQueryNodesFilteringTest
extends GridCommonAbstra
IgniteConfiguration node2Cfg = getConfiguration("node2", true,
null)
.setFailureHandler(failHnd);
- try (Ignite node2 = startGrid(node2Cfg)) {
- assertTrue("Failure handler hasn't been invoked on the joined
node.",
- latch.await(5, TimeUnit.SECONDS));
+ try {
+ startGrid(node2Cfg);
+ }
+ catch (Throwable t) {
+ // Node start failure with unknow class of continous query
filter is ok.
+ if (!X.hasCause(t, "Class not found for continuous query
remote filter [name=" + ENTRY_FILTER_CLS_NAME,
+ IgniteException.class))
+ throw t;
+ }
+ finally {
+ assertTrue("Failure handler hasn't been invoked on the joined
node.", latch.await(5, TimeUnit.SECONDS));
}
}
}
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 47b107edeee..c16bf257775 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_FINISHED;
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,10 +156,10 @@ public class GridEventConsumeSelfTest extends
GridCommonAbstractTest {
* @param proc Continuous processor.
* @return Local event routines.
*/
- private Collection<LocalRoutineInfo> localRoutines(GridContinuousProcessor
proc) {
- return F.view(U.<Map<UUID, LocalRoutineInfo>>field(proc,
"locInfos").values(),
- new IgnitePredicate<LocalRoutineInfo>() {
- @Override public boolean apply(LocalRoutineInfo info) {
+ private Collection<ContinousRoutineLocalInfo>
localRoutines(GridContinuousProcessor proc) {
+ return F.view(U.<Map<UUID, ContinousRoutineLocalInfo>>field(proc,
"locInfos").values(),
+ new IgnitePredicate<>() {
+ @Override public boolean apply(ContinousRoutineLocalInfo info)
{
return info.handler().isEvents();
}
});
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 eabb741a605..87799c52e3c 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 @@ public class ZookeeperDiscoveryImpl {
}
/**
- * 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.