This is an automated email from the ASF dual-hosted git repository.
anton-vinogradov 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 32d3cad3f56 IGNITE-28528 Revise marshalling of the node filter in
StartRequestData (#13440)
32d3cad3f56 is described below
commit 32d3cad3f5630ac72598434ed71bc299434aec58
Author: Anton Vinogradov <[email protected]>
AuthorDate: Mon Aug 10 14:36:50 2026 +0300
IGNITE-28528 Revise marshalling of the node filter in StartRequestData
(#13440)
---
.../ignite/internal/CoreMessagesProvider.java | 4 +-
.../ignite/internal/GridEventConsumeHandler.java | 10 +-
.../ignite/internal/GridJobExecuteRequest.java | 102 ++++++---------------
.../ignite/internal/GridMessageListenHandler.java | 16 +---
.../managers/communication/GridIoManager.java | 19 +---
.../managers/communication/GridIoUserMessage.java | 79 +++-------------
...nfoBean.java => GridDeploymentInfoMessage.java} | 16 ++--
.../managers/deployment/GridDeploymentManager.java | 73 ++++++++++++++-
.../deployment/GridDeploymentMetadata.java | 20 ----
.../eventstorage/GridEventStorageManager.java | 27 +-----
.../eventstorage/GridEventStorageRequest.java | 64 ++-----------
.../processors/affinity/GridAffinityUtils.java | 3 +-
.../cache/GridCacheDeploymentManager.java | 21 ++---
.../processors/cache/GridCacheMessage.java | 10 +-
.../CacheContinuousQueryDeployableObject.java | 10 +-
.../continuous/GridContinuousProcessor.java | 47 ++++++++--
.../processors/continuous/StartRequestData.java | 75 ++-------------
.../datastreamer/DataStreamProcessor.java | 25 ++---
.../processors/datastreamer/DataStreamerImpl.java | 5 +-
.../datastreamer/DataStreamerRequest.java | 63 +++----------
.../internal/processors/job/GridJobProcessor.java | 18 +---
.../internal/processors/task/GridTaskWorker.java | 5 +-
.../main/resources/META-INF/classnames.properties | 2 +-
.../datastreamer/DataStreamerImplSelfTest.java | 5 +-
.../p2p/ClassLoadingProblemExceptionTest.java | 4 +-
25 files changed, 248 insertions(+), 475 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 a460ed789d2..3b066879eea 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
@@ -30,7 +30,7 @@ import
org.apache.ignite.internal.managers.communication.GridIoUserMessage;
import org.apache.ignite.internal.managers.communication.IgniteIoTestMessage;
import org.apache.ignite.internal.managers.communication.IgniteMessageFactory;
import org.apache.ignite.internal.managers.communication.SessionChannelMessage;
-import org.apache.ignite.internal.managers.deployment.GridDeploymentInfoBean;
+import
org.apache.ignite.internal.managers.deployment.GridDeploymentInfoMessage;
import org.apache.ignite.internal.managers.deployment.GridDeploymentRequest;
import org.apache.ignite.internal.managers.deployment.GridDeploymentResponse;
import
org.apache.ignite.internal.managers.encryption.ChangeCacheEncryptionRequest;
@@ -686,7 +686,7 @@ public class CoreMessagesProvider extends
AbstractMessageFactoryProvider {
// [12200 - 12300]: Binary, classloading and marshalling messages.
msgIdx = 12200;
- register(GridDeploymentInfoBean.class);
+ register(GridDeploymentInfoMessage.class);
register(GridDeploymentRequest.class);
register(GridDeploymentResponse.class);
register(MissingMappingRequestMessage.class);
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 cad29b7b129..8ebea159666 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
@@ -34,7 +34,7 @@ import org.apache.ignite.events.Event;
import org.apache.ignite.internal.cluster.ClusterTopologyCheckedException;
import org.apache.ignite.internal.managers.deployment.GridDeployment;
import org.apache.ignite.internal.managers.deployment.GridDeploymentInfo;
-import org.apache.ignite.internal.managers.deployment.GridDeploymentInfoBean;
+import
org.apache.ignite.internal.managers.deployment.GridDeploymentInfoMessage;
import org.apache.ignite.internal.managers.deployment.P2PClassLoadingIssues;
import org.apache.ignite.internal.managers.eventstorage.GridLocalEventListener;
import org.apache.ignite.internal.processors.affinity.AffinityTopologyVersion;
@@ -403,7 +403,7 @@ class GridEventConsumeHandler implements
GridContinuousHandler {
if (dep == null)
throw new IgniteDeploymentCheckedException("Failed to deploy
event filter: " + filter);
- depInfo = new GridDeploymentInfoBean(dep);
+ depInfo = new GridDeploymentInfoMessage(dep);
filterBytes = U.marshal(ctx.marshaller(), filter);
}
@@ -417,11 +417,7 @@ class GridEventConsumeHandler implements
GridContinuousHandler {
if (filterBytes != null) {
try {
- GridDeployment dep =
ctx.deploy().getGlobalDeployment(depInfo.deployMode(), clsName, clsName,
- depInfo.userVersion(), nodeId, depInfo.classLoaderId(),
depInfo.participants(), null);
-
- if (dep == null)
- throw new IgniteDeploymentCheckedException("Failed to
obtain deployment for class: " + clsName);
+ GridDeployment dep = ctx.deploy().globalDeployment(depInfo,
clsName, nodeId);
filter = U.unmarshal(ctx, filterBytes,
U.resolveClassLoader(dep.classLoader(), ctx.config()));
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/GridJobExecuteRequest.java
b/modules/core/src/main/java/org/apache/ignite/internal/GridJobExecuteRequest.java
index 098e8514f68..c431610e1ea 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/GridJobExecuteRequest.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/GridJobExecuteRequest.java
@@ -24,10 +24,10 @@ import java.util.UUID;
import org.apache.ignite.cluster.ClusterNode;
import org.apache.ignite.compute.ComputeJob;
import org.apache.ignite.compute.ComputeJobSibling;
-import org.apache.ignite.configuration.DeploymentMode;
+import org.apache.ignite.internal.managers.deployment.GridDeploymentInfo;
+import
org.apache.ignite.internal.managers.deployment.GridDeploymentInfoMessage;
import org.apache.ignite.internal.processors.affinity.AffinityTopologyVersion;
import org.apache.ignite.internal.util.tostring.GridToStringExclude;
-import org.apache.ignite.internal.util.tostring.GridToStringInclude;
import org.apache.ignite.internal.util.typedef.internal.S;
import org.apache.ignite.internal.util.typedef.internal.U;
import org.apache.ignite.lang.IgnitePredicate;
@@ -70,22 +70,17 @@ public class GridJobExecuteRequest implements
ExecutorAwareMessage, DeferredUnma
@Order(5)
String taskName;
- /** */
+ /** Deployment of the task classes. */
@Order(6)
- String userVer;
+ GridDeploymentInfoMessage depInfo;
/** */
@Order(7)
String taskClsName;
- /** Node class loader participants. */
- @GridToStringInclude
- @Order(8)
- Map<UUID, IgniteUuid> ldrParticipants;
-
/** */
@GridToStringExclude
- @Order(9)
+ @Order(8)
byte[] sesAttrsBytes;
/** */
@@ -95,7 +90,7 @@ public class GridJobExecuteRequest implements
ExecutorAwareMessage, DeferredUnma
/** */
@GridToStringExclude
- @Order(10)
+ @Order(9)
byte[] jobAttrsBytes;
/** */
@@ -104,7 +99,7 @@ public class GridJobExecuteRequest implements
ExecutorAwareMessage, DeferredUnma
Map<? extends Serializable, ? extends Serializable> jobAttrs;
/** Checkpoint SPI name. */
- @Order(11)
+ @Order(10)
String cpSpi;
/** Left unset for a continuous task: such a job requests its siblings
from the task node instead. */
@@ -112,38 +107,30 @@ public class GridJobExecuteRequest implements
ExecutorAwareMessage, DeferredUnma
Collection<ComputeJobSibling> siblings;
/** */
- @Order(12)
+ @Order(11)
byte[] siblingsBytes;
/** Transient since needs to hold local creation time. */
private final long createTime = U.currentTimeMillis();
/** */
- @Order(13)
- IgniteUuid clsLdrId;
-
- /** */
- @Order(14)
- DeploymentMode depMode;
-
- /** */
- @Order(15)
+ @Order(12)
boolean dynamicSiblings;
/** */
- @Order(16)
+ @Order(13)
boolean forceLocDep;
/** */
- @Order(17)
+ @Order(14)
boolean sesFullSup;
/** */
- @Order(18)
+ @Order(15)
boolean internal;
/** */
- @Order(19)
+ @Order(16)
Collection<UUID> top;
/** */
@@ -151,23 +138,23 @@ public class GridJobExecuteRequest implements
ExecutorAwareMessage, DeferredUnma
IgnitePredicate<ClusterNode> topPred;
/** */
- @Order(20)
+ @Order(17)
byte[] topPredBytes;
/** */
- @Order(21)
+ @Order(18)
int[] cacheIds;
/** */
- @Order(22)
+ @Order(19)
int part;
/** */
- @Order(23)
+ @Order(20)
AffinityTopologyVersion topVer;
/** */
- @Order(24)
+ @Order(21)
String execName;
/**
@@ -181,7 +168,7 @@ public class GridJobExecuteRequest implements
ExecutorAwareMessage, DeferredUnma
* @param sesId Task session ID.
* @param jobId Job ID.
* @param taskName Task name.
- * @param userVer Code version.
+ * @param depInfo Deployment of the task classes.
* @param taskClsName Fully qualified task name.
* @param job Job.
* @param startTaskTime Task execution start time.
@@ -192,10 +179,7 @@ public class GridJobExecuteRequest implements
ExecutorAwareMessage, DeferredUnma
* @param sesAttrs Session attributes.
* @param jobAttrs Job attributes.
* @param cpSpi Collision SPI.
- * @param clsLdrId Task local class loader id.
- * @param depMode Task deployment mode.
* @param dynamicSiblings {@code True} if siblings are dynamic.
- * @param ldrParticipants Other node class loader IDs that can also load
classes.
* @param forceLocDep {@code True} If remote node should ignore deployment
settings.
* @param sesFullSup {@code True} if session attributes are disabled.
* @param internal {@code True} if internal job.
@@ -208,7 +192,7 @@ public class GridJobExecuteRequest implements
ExecutorAwareMessage, DeferredUnma
IgniteUuid sesId,
IgniteUuid jobId,
String taskName,
- String userVer,
+ GridDeploymentInfo depInfo,
String taskClsName,
ComputeJob job,
long startTaskTime,
@@ -219,10 +203,7 @@ public class GridJobExecuteRequest implements
ExecutorAwareMessage, DeferredUnma
Map<Object, Object> sesAttrs,
Map<? extends Serializable, ? extends Serializable> jobAttrs,
String cpSpi,
- IgniteUuid clsLdrId,
- DeploymentMode depMode,
boolean dynamicSiblings,
- Map<UUID, IgniteUuid> ldrParticipants,
boolean forceLocDep,
boolean sesFullSup,
boolean internal,
@@ -238,14 +219,12 @@ public class GridJobExecuteRequest implements
ExecutorAwareMessage, DeferredUnma
assert sesAttrs != null || !sesFullSup;
assert jobAttrs != null;
assert top != null || topPred != null;
- assert clsLdrId != null;
- assert userVer != null;
- assert depMode != null;
+ assert depInfo != null;
this.sesId = sesId;
this.jobId = jobId;
this.taskName = taskName;
- this.userVer = userVer;
+ this.depInfo = new GridDeploymentInfoMessage(depInfo);
this.taskClsName = taskClsName;
this.job = job;
this.startTaskTime = startTaskTime;
@@ -256,10 +235,7 @@ public class GridJobExecuteRequest implements
ExecutorAwareMessage, DeferredUnma
this.siblings = dynamicSiblings ? null : siblings;
this.sesAttrs = sesAttrs;
this.jobAttrs = jobAttrs;
- this.clsLdrId = clsLdrId;
- this.depMode = depMode;
this.dynamicSiblings = dynamicSiblings;
- this.ldrParticipants = ldrParticipants;
this.forceLocDep = forceLocDep;
this.sesFullSup = sesFullSup;
this.internal = internal;
@@ -285,6 +261,11 @@ public class GridJobExecuteRequest implements
ExecutorAwareMessage, DeferredUnma
return jobId;
}
+ /** @return Deployment of the task classes. */
+ public GridDeploymentInfo deploymentInfo() {
+ return depInfo;
+ }
+
/**
* @return Task class name.
*/
@@ -299,13 +280,6 @@ public class GridJobExecuteRequest implements
ExecutorAwareMessage, DeferredUnma
return taskName;
}
- /**
- * @return Task version.
- */
- public String userVersion() {
- return userVer;
- }
-
/**
* @return Grid job.
*/
@@ -364,27 +338,6 @@ public class GridJobExecuteRequest implements
ExecutorAwareMessage, DeferredUnma
return cpSpi;
}
- /**
- * @return Task local class loader id.
- */
- public IgniteUuid classLoaderId() {
- return clsLdrId;
- }
-
- /**
- * @return Deployment mode.
- */
- public DeploymentMode deploymentMode() {
- return depMode;
- }
-
- /**
- * @return Node class loader participant map.
- */
- public Map<UUID, IgniteUuid> loaderParticipants() {
- return ldrParticipants;
- }
-
/**
* @return Returns {@code true} if deployment should always be used.
*/
@@ -446,7 +399,6 @@ public class GridJobExecuteRequest implements
ExecutorAwareMessage, DeferredUnma
return topVer;
}
-
/** {@inheritDoc} */
@Override public String toString() {
return S.toString(GridJobExecuteRequest.class, this);
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 e59873f4b44..f3397e5d214 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
@@ -28,7 +28,7 @@ import java.util.UUID;
import org.apache.ignite.IgniteCheckedException;
import org.apache.ignite.IgniteException;
import org.apache.ignite.internal.managers.deployment.GridDeployment;
-import org.apache.ignite.internal.managers.deployment.GridDeploymentInfoBean;
+import
org.apache.ignite.internal.managers.deployment.GridDeploymentInfoMessage;
import org.apache.ignite.internal.processors.affinity.AffinityTopologyVersion;
import org.apache.ignite.internal.processors.continuous.GridContinuousBatch;
import
org.apache.ignite.internal.processors.continuous.GridContinuousBatchAdapter;
@@ -64,7 +64,7 @@ public class GridMessageListenHandler implements
GridContinuousHandler {
private String clsName;
/** */
- private GridDeploymentInfoBean depInfo;
+ private GridDeploymentInfoMessage depInfo;
/** */
private boolean depEnabled;
@@ -166,7 +166,7 @@ public class GridMessageListenHandler implements
GridContinuousHandler {
if (dep == null)
throw new IgniteDeploymentCheckedException("Failed to deploy
message listener.");
- depInfo = new GridDeploymentInfoBean(dep);
+ depInfo = new GridDeploymentInfoMessage(dep);
depEnabled = true;
}
@@ -178,13 +178,7 @@ public class GridMessageListenHandler implements
GridContinuousHandler {
assert ctx.config().isPeerClassLoadingEnabled();
try {
- GridDeployment dep =
ctx.deploy().getGlobalDeployment(depInfo.deployMode(), clsName, clsName,
- depInfo.userVersion(), nodeId, depInfo.classLoaderId(),
depInfo.participants(), null);
-
- if (dep == null)
- throw new IgniteDeploymentCheckedException("Failed to obtain
deployment for class: " + clsName);
-
- ClassLoader ldr = dep.classLoader();
+ ClassLoader ldr = ctx.deploy().globalDeployment(depInfo, clsName,
nodeId).classLoader();
if (topicBytes != null)
topic = U.unmarshal(ctx, topicBytes, U.resolveClassLoader(ldr,
ctx.config()));
@@ -260,7 +254,7 @@ public class GridMessageListenHandler implements
GridContinuousHandler {
topicBytes = U.readByteArray(in);
predBytes = U.readByteArray(in);
clsName = U.readString(in);
- depInfo = (GridDeploymentInfoBean)in.readObject();
+ depInfo = (GridDeploymentInfoMessage)in.readObject();
}
else {
topic = in.readObject();
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoManager.java
b/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoManager.java
index 320f387a1c6..5e226a9ccaa 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoManager.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoManager.java
@@ -2422,10 +2422,7 @@ public class GridIoManager extends
GridManagerAdapter<CommunicationSpi<Object>>
depClsName,
topic,
serTopic,
- dep != null ? dep.classLoaderId() : null,
- dep != null ? dep.deployMode() : null,
- dep != null ? dep.userVersion() : null,
- dep != null ? dep.participants() : null);
+ dep);
if (ordered)
sendOrderedMessageToGridTopic(nodes, TOPIC_COMM_USER, ioMsg,
PUBLIC_POOL, timeout, true);
@@ -3635,21 +3632,15 @@ public class GridIoManager extends
GridManagerAdapter<CommunicationSpi<Object>>
if (dep == null &&
ctx.config().isPeerClassLoadingEnabled() &&
ioMsg.deploymentClassName() != null) {
- dep = ctx.deploy().getGlobalDeployment(
- ioMsg.deploymentMode(),
- ioMsg.deploymentClassName(),
- ioMsg.deploymentClassName(),
- ioMsg.userVersion(),
- nodeId,
- ioMsg.classLoaderId(),
- ioMsg.loaderParticipants(),
- null);
+ dep =
ctx.deploy().globalDeployment(ioMsg.deploymentInfo(),
ioMsg.deploymentClassName(),
+ ioMsg.deploymentClassName(), nodeId);
- if (dep == null)
+ if (dep == null) {
throw new IgniteDeploymentCheckedException(
"Failed to obtain deployment information for
user message. " +
"If you are using custom message or topic
class, try implementing " +
"GridPeerDeployAware interface. [msg=" +
ioMsg + ']');
+ }
ioMsg.deployment(dep); // Cache deployment.
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoUserMessage.java
b/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoUserMessage.java
index afe4db51c2b..ef4f0b61a8b 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoUserMessage.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoUserMessage.java
@@ -17,16 +17,12 @@
package org.apache.ignite.internal.managers.communication;
-import java.util.Collections;
-import java.util.Map;
-import java.util.UUID;
-import org.apache.ignite.configuration.DeploymentMode;
import org.apache.ignite.internal.Order;
import org.apache.ignite.internal.UseBinaryMarshaller;
import org.apache.ignite.internal.managers.deployment.GridDeployment;
-import org.apache.ignite.internal.util.tostring.GridToStringInclude;
+import org.apache.ignite.internal.managers.deployment.GridDeploymentInfo;
+import
org.apache.ignite.internal.managers.deployment.GridDeploymentInfoMessage;
import org.apache.ignite.internal.util.typedef.internal.S;
-import org.apache.ignite.lang.IgniteUuid;
import org.apache.ignite.plugin.extensions.communication.Message;
import org.jetbrains.annotations.Nullable;
@@ -42,34 +38,21 @@ public class GridIoUserMessage implements Message {
@Order(0)
byte[] bodyBytes;
- /** Class loader ID. */
- @Order(1)
- IgniteUuid clsLdrId;
-
/** Message topic. */
private Object topic;
/** Serialized message topic. */
- @Order(2)
+ @Order(1)
byte[] topicBytes;
- /** Deployment mode. */
- @Order(3)
- DeploymentMode depMode;
+ /** Deployment of the message classes. */
+ @Order(2)
+ GridDeploymentInfoMessage depInfo;
/** Deployment class name. */
- @Order(4)
+ @Order(3)
String depClsName;
- /** User version. */
- @Order(5)
- String userVer;
-
- /** Node class loader participants. */
- @Order(6)
- @GridToStringInclude
- Map<UUID, IgniteUuid> ldrParties;
-
/** Message deployment. */
private GridDeployment dep;
@@ -79,10 +62,7 @@ public class GridIoUserMessage implements Message {
* @param depClsName Message body class name.
* @param topic Message topic.
* @param topicBytes Serialized message topic bytes.
- * @param clsLdrId Class loader ID.
- * @param depMode Deployment mode.
- * @param userVer User version.
- * @param ldrParties Node loader participant map.
+ * @param depInfo Deployment of the message classes.
*/
GridIoUserMessage(
Object body,
@@ -90,19 +70,13 @@ public class GridIoUserMessage implements Message {
@Nullable String depClsName,
@Nullable Object topic,
@Nullable byte[] topicBytes,
- @Nullable IgniteUuid clsLdrId,
- @Nullable DeploymentMode depMode,
- @Nullable String userVer,
- @Nullable Map<UUID, IgniteUuid> ldrParties) {
+ @Nullable GridDeploymentInfo depInfo) {
this.body = body;
this.bodyBytes = bodyBytes;
this.depClsName = depClsName;
this.topic = topic;
this.topicBytes = topicBytes;
- this.depMode = depMode;
- this.clsLdrId = clsLdrId;
- this.userVer = userVer;
- this.ldrParties = ldrParties;
+ this.depInfo = depInfo != null ? new
GridDeploymentInfoMessage(depInfo) : null;
}
/**
@@ -119,20 +93,6 @@ public class GridIoUserMessage implements Message {
return bodyBytes;
}
- /**
- * @return the Class loader ID.
- */
- @Nullable public IgniteUuid classLoaderId() {
- return clsLdrId;
- }
-
- /**
- * @return Deployment mode.
- */
- @Nullable public DeploymentMode deploymentMode() {
- return depMode;
- }
-
/**
* @return Message body class name.
*/
@@ -140,20 +100,6 @@ public class GridIoUserMessage implements Message {
return depClsName;
}
- /**
- * @return User version.
- */
- @Nullable public String userVersion() {
- return userVer;
- }
-
- /**
- * @return Node class loader participant map.
- */
- @Nullable public Map<UUID, IgniteUuid> loaderParticipants() {
- return ldrParties != null ? Collections.unmodifiableMap(ldrParties) :
null;
- }
-
/**
* @return Serialized message topic.
*/
@@ -182,6 +128,11 @@ public class GridIoUserMessage implements Message {
this.body = body;
}
+ /** @return Deployment of the message classes, or {@code null} when peer
class loading is off. */
+ @Nullable public GridDeploymentInfo deploymentInfo() {
+ return depInfo;
+ }
+
/**
* @return Message body.
*/
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentInfoBean.java
b/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentInfoMessage.java
similarity index 86%
rename from
modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentInfoBean.java
rename to
modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentInfoMessage.java
index ba9f48c3ca2..4a4f4e4a6a9 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentInfoBean.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentInfoMessage.java
@@ -28,9 +28,9 @@ import org.apache.ignite.lang.IgniteUuid;
import org.apache.ignite.plugin.extensions.communication.Message;
/**
- * Deployment info bean.
+ * Deployment of classes, as it travels inside the messages carrying them.
*/
-public class GridDeploymentInfoBean implements Message, GridDeploymentInfo,
Serializable {
+public class GridDeploymentInfoMessage implements Message, GridDeploymentInfo,
Serializable {
/** */
private static final long serialVersionUID = 0L;
@@ -54,7 +54,7 @@ public class GridDeploymentInfoBean implements Message,
GridDeploymentInfo, Seri
/**
* Empty constructor for a message factory.
*/
- public GridDeploymentInfoBean() {
+ public GridDeploymentInfoMessage() {
/* No-op. */
}
@@ -64,7 +64,7 @@ public class GridDeploymentInfoBean implements Message,
GridDeploymentInfo, Seri
* @param depMode Deployment mode.
* @param participants Participants.
*/
- public GridDeploymentInfoBean(
+ public GridDeploymentInfoMessage(
IgniteUuid clsLdrId,
String userVer,
DeploymentMode depMode,
@@ -79,7 +79,7 @@ public class GridDeploymentInfoBean implements Message,
GridDeploymentInfo, Seri
/**
* @param dep Grid deployment.
*/
- public GridDeploymentInfoBean(GridDeploymentInfo dep) {
+ public GridDeploymentInfoMessage(GridDeploymentInfo dep) {
clsLdrId = dep.classLoaderId();
depMode = dep.deployMode();
userVer = dep.userVersion();
@@ -118,12 +118,12 @@ public class GridDeploymentInfoBean implements Message,
GridDeploymentInfo, Seri
/** {@inheritDoc} */
@Override public boolean equals(Object o) {
- return o == this || o instanceof GridDeploymentInfoBean &&
- clsLdrId.equals(((GridDeploymentInfoBean)o).clsLdrId);
+ return o == this || o instanceof GridDeploymentInfoMessage &&
+ clsLdrId.equals(((GridDeploymentInfoMessage)o).clsLdrId);
}
/** {@inheritDoc} */
@Override public String toString() {
- return S.toString(GridDeploymentInfoBean.class, this);
+ return S.toString(GridDeploymentInfoMessage.class, this);
}
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentManager.java
b/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentManager.java
index 3069f201be9..5be5bd627c6 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentManager.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentManager.java
@@ -27,6 +27,7 @@ import org.apache.ignite.compute.ComputeTask;
import org.apache.ignite.compute.ComputeTaskName;
import org.apache.ignite.configuration.DeploymentMode;
import org.apache.ignite.internal.GridKernalContext;
+import org.apache.ignite.internal.IgniteDeploymentCheckedException;
import org.apache.ignite.internal.IgniteInternalFuture;
import org.apache.ignite.internal.managers.GridManagerAdapter;
import
org.apache.ignite.internal.managers.deployment.protocol.gg.GridProtocolHandler;
@@ -402,6 +403,73 @@ public class GridDeploymentManager extends
GridManagerAdapter<DeploymentSpi> {
return locStore.getDeployment(meta);
}
+ /**
+ * Resolves the class loader classes described by {@code depInfo} must be
read with. Blocks when the deployment
+ * has to be requested, so it must not be called from a socket-reading
thread.
+ *
+ * @param depInfo Deployment of the classes, or {@code null} when they
carry none.
+ * @param clsName Name of a class the deployment must be able to load.
+ * @param sndNodeId Node the classes came from.
+ * @return Class loader of the deployment, or the local one when there is
no deployment.
+ * @throws IgniteDeploymentCheckedException If the deployment cannot be
obtained.
+ */
+ public ClassLoader classLoader(@Nullable GridDeploymentInfo depInfo,
String clsName, UUID sndNodeId)
+ throws IgniteDeploymentCheckedException {
+ if (depInfo == null)
+ return U.resolveClassLoader(ctx.config());
+
+ return U.resolveClassLoader(globalDeployment(depInfo, clsName,
sndNodeId).classLoader(), ctx.config());
+ }
+
+ /**
+ * Resolves the deployment {@code depInfo} describes, for the classes of
{@code clsName}.
+ *
+ * @param depInfo Deployment of the classes, or {@code null} when the
message carries none.
+ * @param clsName Name of a class the deployment must be able to load.
+ * @param sndNodeId Node the classes came from. It is not always the node
that created the class loader: a node
+ * that got the classes by peer loading passes them on as a
participant of the same deployment.
+ * @return The deployment the classes are loaded with.
+ * @throws IgniteDeploymentCheckedException If there is no deployment to
resolve, it is gone, or peer class
+ * loading is off.
+ */
+ public GridDeployment globalDeployment(@Nullable GridDeploymentInfo
depInfo, String clsName, UUID sndNodeId)
+ throws IgniteDeploymentCheckedException {
+ GridDeployment dep = globalDeployment(depInfo, clsName, clsName,
sndNodeId);
+
+ if (dep == null) {
+ throw new IgniteDeploymentCheckedException("Failed to obtain
deployment for class (is peer class " +
+ "loading turned on?): " + clsName);
+ }
+
+ return dep;
+ }
+
+ /**
+ * Resolves the deployment {@code depInfo} describes, as
+ * {@link #globalDeployment(GridDeploymentInfo, String, UUID)} does, but
under {@code rsrcName} (a task may be
+ * deployed under a name of its own) and returns {@code null} instead of
throwing, for callers that have somewhere
+ * else to look.
+ *
+ * @param depInfo Deployment of the classes, or {@code null} when the
message carries none.
+ * @param rsrcName Name the classes are deployed under.
+ * @param clsName Name of a class the deployment must be able to load.
+ * @param sndNodeId Node the classes came from.
+ * @return The deployment, or {@code null} when there is none to resolve
or none is found.
+ */
+ @Nullable public GridDeployment globalDeployment(@Nullable
GridDeploymentInfo depInfo, String rsrcName,
+ String clsName, UUID sndNodeId) {
+ if (depInfo == null)
+ return null;
+
+ return getGlobalDeployment(depInfo.deployMode(),
+ rsrcName,
+ clsName,
+ depInfo.userVersion(),
+ sndNodeId,
+ depInfo.classLoaderId(),
+ depInfo.participants());
+ }
+
/**
* @param depMode Deployment mode.
* @param rsrcName Resource name (could be task name).
@@ -410,7 +478,6 @@ public class GridDeploymentManager extends
GridManagerAdapter<DeploymentSpi> {
* @param sndNodeId Sender node ID.
* @param clsLdrId Class loader ID.
* @param participants Node class loader participant map.
- * @param nodeFilter Node filter for class loader.
* @return Deployment class if found.
*/
@Nullable public GridDeployment getGlobalDeployment(
@@ -420,8 +487,7 @@ public class GridDeploymentManager extends
GridManagerAdapter<DeploymentSpi> {
String userVer,
UUID sndNodeId,
IgniteUuid clsLdrId,
- Map<UUID, IgniteUuid> participants,
- @Nullable IgnitePredicate<ClusterNode> nodeFilter) {
+ Map<UUID, IgniteUuid> participants) {
if (locDep != null)
return locDep;
@@ -439,7 +505,6 @@ public class GridDeploymentManager extends
GridManagerAdapter<DeploymentSpi> {
meta.senderNodeId(sndNodeId);
meta.classLoaderId(clsLdrId);
meta.participants(participants);
- meta.nodeFilter(nodeFilter);
if (!ctx.config().isPeerClassLoadingEnabled()) {
meta.record(true);
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentMetadata.java
b/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentMetadata.java
index 370825c5fee..c5f6857be9c 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentMetadata.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentMetadata.java
@@ -19,11 +19,9 @@ package org.apache.ignite.internal.managers.deployment;
import java.util.Map;
import java.util.UUID;
-import org.apache.ignite.cluster.ClusterNode;
import org.apache.ignite.configuration.DeploymentMode;
import org.apache.ignite.internal.util.tostring.GridToStringInclude;
import org.apache.ignite.internal.util.typedef.internal.S;
-import org.apache.ignite.lang.IgnitePredicate;
import org.apache.ignite.lang.IgniteUuid;
/**
@@ -61,9 +59,6 @@ public class GridDeploymentMetadata {
/** */
private boolean record;
- /** */
- private IgnitePredicate<ClusterNode> nodeFilter;
-
/**
*
*/
@@ -87,7 +82,6 @@ public class GridDeploymentMetadata {
participants = meta.participants();
parentLdr = meta.parentLoader();
record = meta.record();
- nodeFilter = meta.nodeFilter();
}
/**
@@ -271,20 +265,6 @@ public class GridDeploymentMetadata {
this.clsLdr = clsLdr;
}
- /**
- * @param nodeFilter Node filter.
- */
- public void nodeFilter(IgnitePredicate<ClusterNode> nodeFilter) {
- this.nodeFilter = nodeFilter;
- }
-
- /**
- * @return Node filter.
- */
- public IgnitePredicate<ClusterNode> nodeFilter() {
- return nodeFilter;
- }
-
/** {@inheritDoc} */
@Override public String toString() {
return S.toString(GridDeploymentMetadata.class, this, "seqNum",
clsLdrId != null ? clsLdrId.localId() : "n/a");
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/GridEventStorageManager.java
b/modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/GridEventStorageManager.java
index 1c883d1474a..62ee043907a 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/GridEventStorageManager.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/GridEventStorageManager.java
@@ -1070,13 +1070,7 @@ public class GridEventStorageManager extends
GridManagerAdapter<EventStorageSpi>
if (dep == null)
throw new IgniteDeploymentCheckedException("Failed to deploy
event filter: " + p);
- GridEventStorageRequest msg = new GridEventStorageRequest(
- resTopicId,
- p,
- dep.classLoaderId(),
- dep.deployMode(),
- dep.userVersion(),
- dep.participants());
+ GridEventStorageRequest msg = new
GridEventStorageRequest(resTopicId, p, dep);
sendMessage(nodes, TOPIC_EVENT, msg, PUBLIC_POOL);
@@ -1248,24 +1242,13 @@ public class GridEventStorageManager extends
GridManagerAdapter<EventStorageSpi>
Collection<Event> evts;
try {
- GridDeployment dep = ctx.deploy().getGlobalDeployment(
- req.deploymentMode(),
- req.filterClassName(),
- req.filterClassName(),
- req.userVersion(),
- nodeId,
- req.classLoaderId(),
- req.loaderParticipants(),
- null);
-
- if (dep == null)
- throw new IgniteDeploymentCheckedException("Failed to
obtain deployment for event filter " +
- "(is peer class loading turned on?): " + req);
-
- MessageMarshalling.unmarshal(req, ctx, null,
U.resolveClassLoader(dep.classLoader(), ctx.config()));
+ MessageMarshalling.unmarshal(req, ctx, null,
+ ctx.deploy().classLoader(req.deploymentInfo(),
req.filterClassName(), nodeId));
filter = (IgnitePredicate<Event>)req.filter();
+ GridDeployment dep =
ctx.deploy().globalDeployment(req.deploymentInfo(), req.filterClassName(),
nodeId);
+
// Resource injection.
ctx.resource().inject(dep,
dep.deployedClass(req.filterClassName()).get1(), filter);
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/GridEventStorageRequest.java
b/modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/GridEventStorageRequest.java
index a2fd6a0c26c..322d8bc2bc1 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/GridEventStorageRequest.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/GridEventStorageRequest.java
@@ -17,19 +17,15 @@
package org.apache.ignite.internal.managers.eventstorage;
-import java.util.Collections;
-import java.util.Map;
-import java.util.UUID;
-import org.apache.ignite.configuration.DeploymentMode;
import org.apache.ignite.internal.DeferredUnmarshalMessage;
import org.apache.ignite.internal.Marshalled;
import org.apache.ignite.internal.Order;
import org.apache.ignite.internal.UseBinaryMarshaller;
-import org.apache.ignite.internal.util.tostring.GridToStringInclude;
+import org.apache.ignite.internal.managers.deployment.GridDeploymentInfo;
+import
org.apache.ignite.internal.managers.deployment.GridDeploymentInfoMessage;
import org.apache.ignite.internal.util.typedef.internal.S;
import org.apache.ignite.lang.IgnitePredicate;
import org.apache.ignite.lang.IgniteUuid;
-import org.jetbrains.annotations.Nullable;
import static org.apache.ignite.internal.GridTopic.TOPIC_EVENT;
@@ -48,27 +44,14 @@ public class GridEventStorageRequest implements
DeferredUnmarshalMessage {
@Order(1)
byte[] filterBytes;
- /** */
+ /** Deployment of the filter classes. */
@Order(2)
- IgniteUuid clsLdrId;
+ GridDeploymentInfoMessage depInfo;
/** */
@Order(3)
- DeploymentMode depMode;
-
- /** */
- @Order(4)
String filterClsName;
- /** */
- @Order(5)
- String userVer;
-
- /** Node class loader participants. */
- @GridToStringInclude
- @Order(6)
- Map<UUID, IgniteUuid> ldrParties;
-
/** */
public GridEventStorageRequest() {
// No-op.
@@ -77,24 +60,12 @@ public class GridEventStorageRequest implements
DeferredUnmarshalMessage {
/**
* @param resTopicId Id of the node waiting for the response.
* @param filter Query filter.
- * @param clsLdrId Class loader ID.
- * @param depMode Deployment mode.
- * @param userVer User version.
- * @param ldrParties Node loader participant map.
+ * @param depInfo Deployment of the filter classes.
*/
- GridEventStorageRequest(
- IgniteUuid resTopicId,
- IgnitePredicate<?> filter,
- IgniteUuid clsLdrId,
- DeploymentMode depMode,
- String userVer,
- Map<UUID, IgniteUuid> ldrParties) {
+ GridEventStorageRequest(IgniteUuid resTopicId, IgnitePredicate<?> filter,
GridDeploymentInfo depInfo) {
this.resTopicId = resTopicId;
this.filter = filter;
- this.clsLdrId = clsLdrId;
- this.depMode = depMode;
- this.userVer = userVer;
- this.ldrParties = ldrParties;
+ this.depInfo = new GridDeploymentInfoMessage(depInfo);
filterClsName = filter.getClass().getName();
}
@@ -109,14 +80,9 @@ public class GridEventStorageRequest implements
DeferredUnmarshalMessage {
return filter;
}
- /** @return Class loader ID. */
- IgniteUuid classLoaderId() {
- return clsLdrId;
- }
-
- /** @return Deployment mode. */
- DeploymentMode deploymentMode() {
- return depMode;
+ /** @return Deployment of the filter classes. */
+ GridDeploymentInfo deploymentInfo() {
+ return depInfo;
}
/** @return Filter class name. */
@@ -124,16 +90,6 @@ public class GridEventStorageRequest implements
DeferredUnmarshalMessage {
return filterClsName;
}
- /** @return User version. */
- String userVersion() {
- return userVer;
- }
-
- /** @return Node class loader participant map. */
- @Nullable Map<UUID, IgniteUuid> loaderParticipants() {
- return ldrParties != null ? Collections.unmodifiableMap(ldrParties) :
null;
- }
-
/** {@inheritDoc} */
@Override public String toString() {
return S.toString(GridEventStorageRequest.class, this);
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 3937be9b792..2eda9688193 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
@@ -103,8 +103,7 @@ class GridAffinityUtils {
msg.userVersion(),
sndNodeId,
msg.classLoaderId(),
- msg.loaderParticipants(),
- null);
+ msg.loaderParticipants());
if (dep == null)
throw new IgniteDeploymentCheckedException("Failed to obtain
affinity object (is peer class loading turned on?): " +
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheDeploymentManager.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheDeploymentManager.java
index 125fe4ce947..46a7c046ea7 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheDeploymentManager.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheDeploymentManager.java
@@ -30,7 +30,7 @@ import org.apache.ignite.configuration.DeploymentMode;
import org.apache.ignite.events.DiscoveryEvent;
import org.apache.ignite.events.Event;
import org.apache.ignite.internal.managers.deployment.GridDeployment;
-import org.apache.ignite.internal.managers.deployment.GridDeploymentInfoBean;
+import
org.apache.ignite.internal.managers.deployment.GridDeploymentInfoMessage;
import org.apache.ignite.internal.managers.eventstorage.GridLocalEventListener;
import org.apache.ignite.internal.util.lang.GridPeerDeployAware;
import org.apache.ignite.internal.util.tostring.GridToStringInclude;
@@ -395,14 +395,14 @@ public class GridCacheDeploymentManager<K, V> extends
GridCacheSharedManagerAdap
// Only set deployment info if it was not set automatically.
if (deployable.deployInfo() == null) {
- GridDeploymentInfoBean dep = globalDeploymentInfo();
+ GridDeploymentInfoMessage dep = globalDeploymentInfo();
if (dep == null) {
GridDeployment locDep0 = locDep.get();
if (locDep0 != null) {
// Will copy sequence number to bean.
- dep = new GridDeploymentInfoBean(locDep0);
+ dep = new GridDeploymentInfoMessage(locDep0);
checkDeploymentIsCorrect(dep, deployable, false);
}
@@ -426,7 +426,7 @@ public class GridCacheDeploymentManager<K, V> extends
GridCacheSharedManagerAdap
* @param failIfNotCorrect Flag determining whether to throw exception or
just warn.
* @throws IgnitePeerToPeerClassLoadingException If deployment is
incorrect.
*/
- private void checkDeploymentIsCorrect(GridDeploymentInfoBean deployment,
GridCacheDeployable deployable,
+ private void checkDeploymentIsCorrect(GridDeploymentInfoMessage
deployment, GridCacheDeployable deployable,
boolean failIfNotCorrect)
throws IgnitePeerToPeerClassLoadingException {
if (deployment.participants() == null
@@ -445,7 +445,7 @@ public class GridCacheDeploymentManager<K, V> extends
GridCacheSharedManagerAdap
/**
* @return First global deployment.
*/
- @Nullable public GridDeploymentInfoBean globalDeploymentInfo() {
+ @Nullable public GridDeploymentInfoMessage globalDeploymentInfo() {
assert depEnabled;
// Do not return info if mode is CONTINUOUS.
@@ -456,14 +456,14 @@ public class GridCacheDeploymentManager<K, V> extends
GridCacheSharedManagerAdap
IgniteUuid locLdrId0 = localLdrId.get();
if (locLdrId0 != null) {
- GridDeploymentInfoBean deploymentInfoBean =
getDepBean(deps.get(localLdrId.get()));
+ GridDeploymentInfoMessage deploymentInfoBean =
getDepBean(deps.get(localLdrId.get()));
if (deploymentInfoBean != null)
return deploymentInfoBean;
}
for (CachedDeploymentInfo<K, V> d : deps.values()) {
- GridDeploymentInfoBean deploymentInfoBean = getDepBean(d);
+ GridDeploymentInfoMessage deploymentInfoBean = getDepBean(d);
if (deploymentInfoBean != null)
return deploymentInfoBean;
}
@@ -472,7 +472,7 @@ public class GridCacheDeploymentManager<K, V> extends
GridCacheSharedManagerAdap
}
/** */
- @Nullable private GridDeploymentInfoBean
getDepBean(CachedDeploymentInfo<K, V> d) {
+ @Nullable private GridDeploymentInfoMessage
getDepBean(CachedDeploymentInfo<K, V> d) {
if (d == null || cctx.discovery().node(d.senderId()) == null)
// Sender has left.
return null;
@@ -484,7 +484,7 @@ public class GridCacheDeploymentManager<K, V> extends
GridCacheSharedManagerAdap
for (UUID id : participants.keySet()) {
if (cctx.discovery().node(id) != null) {
// At least 1 participant is still in the grid.
- return new GridDeploymentInfoBean(
+ return new GridDeploymentInfoMessage(
d.loaderId(),
d.userVersion(),
d.mode(),
@@ -644,8 +644,7 @@ public class GridCacheDeploymentManager<K, V> extends
GridCacheSharedManagerAdap
userVer,
sndId,
ldrId,
- participants,
- F.<ClusterNode>alwaysTrue());
+ participants);
return d != null ? d.deployedClass(name) : null;
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheMessage.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheMessage.java
index 2f09c39eb9f..8cdceed58e3 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheMessage.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheMessage.java
@@ -29,7 +29,7 @@ import org.apache.ignite.internal.Order;
import org.apache.ignite.internal.StripedMessage;
import org.apache.ignite.internal.managers.deployment.GridDeployment;
import org.apache.ignite.internal.managers.deployment.GridDeploymentInfo;
-import org.apache.ignite.internal.managers.deployment.GridDeploymentInfoBean;
+import
org.apache.ignite.internal.managers.deployment.GridDeploymentInfoMessage;
import org.apache.ignite.internal.processors.affinity.AffinityTopologyVersion;
import org.apache.ignite.internal.processors.cache.transactions.IgniteTxEntry;
import org.apache.ignite.internal.util.tostring.GridToStringInclude;
@@ -65,7 +65,7 @@ public abstract class GridCacheMessage implements
DeferredUnmarshalMessage, Stri
/** */
@GridToStringInclude
@Order(1)
- public GridDeploymentInfoBean depInfo;
+ public GridDeploymentInfoMessage depInfo;
/** */
@GridToStringInclude
@@ -257,8 +257,8 @@ public abstract class GridCacheMessage implements
DeferredUnmarshalMessage, Stri
if (((GridDeployment)depInfo).local())
return;
- this.depInfo = depInfo instanceof GridDeploymentInfoBean ?
- (GridDeploymentInfoBean)depInfo : new
GridDeploymentInfoBean(depInfo);
+ this.depInfo = depInfo instanceof GridDeploymentInfoMessage ?
+ (GridDeploymentInfoMessage)depInfo : new
GridDeploymentInfoMessage(depInfo);
}
}
@@ -266,7 +266,7 @@ public abstract class GridCacheMessage implements
DeferredUnmarshalMessage, Stri
* @return Preset deployment info.
* @see GridCacheDeployable#deployInfo()
*/
- public GridDeploymentInfoBean deployInfo() {
+ public GridDeploymentInfoMessage deployInfo() {
return depInfo;
}
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 c4e4005095f..5f31963ba21 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
@@ -27,7 +27,7 @@ import org.apache.ignite.internal.GridKernalContext;
import org.apache.ignite.internal.IgniteDeploymentCheckedException;
import org.apache.ignite.internal.managers.deployment.GridDeployment;
import org.apache.ignite.internal.managers.deployment.GridDeploymentInfo;
-import org.apache.ignite.internal.managers.deployment.GridDeploymentInfoBean;
+import
org.apache.ignite.internal.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;
@@ -74,7 +74,7 @@ class CacheContinuousQueryDeployableObject implements
Externalizable {
if (dep == null)
throw new IgniteDeploymentCheckedException("Failed to deploy
object: " + obj);
- depInfo = new GridDeploymentInfoBean(dep);
+ depInfo = new GridDeploymentInfoMessage(dep);
bytes = U.marshal(ctx, obj);
}
@@ -88,11 +88,7 @@ class CacheContinuousQueryDeployableObject implements
Externalizable {
<T> T unmarshal(UUID nodeId, GridKernalContext ctx) throws
IgniteCheckedException {
assert ctx != null;
- GridDeployment dep =
ctx.deploy().getGlobalDeployment(depInfo.deployMode(), clsName, clsName,
- depInfo.userVersion(), nodeId, depInfo.classLoaderId(),
depInfo.participants(), null);
-
- if (dep == null)
- throw new IgniteDeploymentCheckedException("Failed to obtain
deployment for class: " + clsName);
+ GridDeployment dep = ctx.deploy().globalDeployment(depInfo, clsName,
nodeId);
return U.unmarshal(ctx, bytes, U.resolveClassLoader(dep.classLoader(),
ctx.config()));
}
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 9edaab63c7a..111d3b74049 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
@@ -58,7 +58,7 @@ import org.apache.ignite.internal.NodeStoppingException;
import org.apache.ignite.internal.cluster.ClusterTopologyCheckedException;
import org.apache.ignite.internal.managers.communication.GridMessageListener;
import org.apache.ignite.internal.managers.deployment.GridDeployment;
-import org.apache.ignite.internal.managers.deployment.GridDeploymentInfoBean;
+import
org.apache.ignite.internal.managers.deployment.GridDeploymentInfoMessage;
import org.apache.ignite.internal.managers.discovery.CustomEventListener;
import org.apache.ignite.internal.managers.discovery.DiscoCache;
import
org.apache.ignite.internal.managers.discovery.DiscoveryMessageResultsCollector;
@@ -951,7 +951,7 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
hnd = hnd.clone();
String clsName = null;
- GridDeploymentInfoBean dep = null;
+ GridDeploymentInfoMessage dep = null;
if (ctx.config().isPeerClassLoadingEnabled()) {
// Handle peer deployment for projection predicate.
@@ -965,7 +965,7 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
if (dep0 == null)
throw new IgniteDeploymentCheckedException("Failed to
deploy projection predicate: " + nodeFilter);
- dep = new GridDeploymentInfoBean(dep0);
+ dep = new GridDeploymentInfoMessage(dep0);
}
}
@@ -981,7 +981,15 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
reqData.deploymentInfo(dep);
}
- reqData.marshal(ctx);
+ if (ctx.config().isPeerClassLoadingEnabled()) {
+ // Handle peer deployment for other handler-specific objects.
+ hnd.p2pMarshal(ctx);
+ }
+
+ reqData.hndBytes = U.marshal(marsh, hnd);
+
+ if (nodeFilter != null)
+ reqData.nodeFilterBytes = U.marshal(marsh, nodeFilter);
if (!immutableDiscoCustomMsg) {
StartRoutineDiscoveryMessage msg = new
StartRoutineDiscoveryMessage(routineId, reqData, Mode.MUTABLE);
@@ -1338,6 +1346,33 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
}
}
+ /**
+ * Restores the objects a start request carries. The discovery layer reads
the message on the thread that reads
+ * the ring, where obtaining a deployment must not happen, so the request
keeps them serialized until here.
+ *
+ * @param msg Message carrying the request.
+ * @param sndId Node that started the routine.
+ */
+ private void unmarshalStartRequest(StartRoutineDiscoveryMessage msg, UUID
sndId) throws IgniteCheckedException {
+ StartRequestData data = msg.startRequestData();
+
+ 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 (data.keepBinary) {
+ assert data.hnd instanceof CacheContinuousQueryHandler :
data.hnd;
+
+ ((CacheContinuousQueryHandler<?, ?>)data.hnd).keepBinary(true);
+ }
+ }
+ }
+
/**
* @param node Sender.
* @param req Start request.
@@ -1353,7 +1388,7 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
IgniteCheckedException err = null;
try {
- data.unmarshal(ctx, node.id());
+ unmarshalStartRequest(req, node.id());
}
catch (IgniteCheckedException e) {
U.error(log, "Failed to unmarshal start request data [nodeId=" +
node.id() +
@@ -1495,7 +1530,7 @@ public class GridContinuousProcessor extends
GridProcessorAdapter {
Exception err = null;
try {
- reqData.unmarshal(ctx, snd.id());
+ unmarshalStartRequest(msg, snd.id());
}
catch (IgniteCheckedException e) {
err = e;
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 9571c231bbf..595aaa56c58 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
@@ -17,17 +17,10 @@
package org.apache.ignite.internal.processors.continuous;
-import java.util.UUID;
-import org.apache.ignite.IgniteCheckedException;
import org.apache.ignite.cluster.ClusterNode;
-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.GridDeploymentInfoBean;
-import
org.apache.ignite.internal.processors.cache.query.continuous.CacheContinuousQueryHandler;
+import
org.apache.ignite.internal.managers.deployment.GridDeploymentInfoMessage;
import org.apache.ignite.internal.util.typedef.internal.S;
-import org.apache.ignite.internal.util.typedef.internal.U;
import org.apache.ignite.lang.IgnitePredicate;
import org.apache.ignite.plugin.extensions.communication.Message;
@@ -35,8 +28,8 @@ import
org.apache.ignite.plugin.extensions.communication.Message;
* Start request data.
*/
public class StartRequestData implements Message {
- /** Node filter. */
- private IgnitePredicate<ClusterNode> nodeFilter;
+ /** Node filter, restored by the processor reading this request. */
+ IgnitePredicate<ClusterNode> nodeFilter;
/** Serialized node filter. */
@Order(0)
@@ -48,10 +41,10 @@ public class StartRequestData implements Message {
/** Deployment info. */
@Order(2)
- GridDeploymentInfoBean depInfo;
+ GridDeploymentInfoMessage depInfo;
- /** Handler. */
- private GridContinuousHandler hnd;
+ /** Handler, restored by the processor reading this request. */
+ GridContinuousHandler hnd;
/** Serialized handler. */
@Order(3)
@@ -119,7 +112,7 @@ public class StartRequestData implements Message {
/**
* @param depInfo New deployment info.
*/
- public void deploymentInfo(GridDeploymentInfoBean depInfo) {
+ public void deploymentInfo(GridDeploymentInfoMessage depInfo) {
this.depInfo = depInfo;
}
@@ -169,58 +162,4 @@ public class StartRequestData implements Message {
@Override public String toString() {
return S.toString(StartRequestData.class, this);
}
-
- /** */
- public void marshal(GridKernalContext ctx) throws IgniteCheckedException {
- if (hnd != null) {
- if (ctx.config().isPeerClassLoadingEnabled()) {
- // Handle peer deployment for other handler-specific objects.
- hnd.p2pMarshal(ctx);
- }
-
- hndBytes = U.marshal(ctx.marshaller(), hnd);
- }
-
- if (nodeFilter != null)
- nodeFilterBytes = U.marshal(ctx.marshaller(), nodeFilter);
- }
-
- /** */
- public void unmarshal(GridKernalContext ctx, UUID sndId) throws
IgniteCheckedException {
- if (ctx.config().isPeerClassLoadingEnabled() && clsName != null) {
- GridDeployment dep =
ctx.deploy().getGlobalDeployment(depInfo.deployMode(),
- clsName,
- clsName,
- depInfo.userVersion(),
- sndId,
- depInfo.classLoaderId(),
- depInfo.participants(),
- null);
-
- if (dep == null)
- throw new IgniteDeploymentCheckedException("Failed to obtain
deployment for class: " + clsName);
-
- nodeFilter = U.unmarshal(ctx.marshaller(),
- nodeFilterBytes,
- U.resolveClassLoader(dep.classLoader(), ctx.config()));
- }
- else {
- nodeFilter = U.unmarshal(ctx.marshaller(),
- nodeFilterBytes,
- U.resolveClassLoader(ctx.config()));
- }
-
- if (hndBytes != null) {
- hnd = U.unmarshal(ctx.marshaller(), hndBytes,
U.resolveClassLoader(ctx.config()));
-
- if (ctx.config().isPeerClassLoadingEnabled())
- hnd.p2pUnmarshal(sndId, ctx);
-
- if (keepBinary) {
- assert hnd instanceof CacheContinuousQueryHandler : hnd;
-
- ((CacheContinuousQueryHandler<?, ?>)hnd).keepBinary(true);
- }
- }
- }
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamProcessor.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamProcessor.java
index 464a74d82ee..92ec36cc210 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamProcessor.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamProcessor.java
@@ -22,6 +22,7 @@ import java.util.UUID;
import org.apache.ignite.IgniteCheckedException;
import org.apache.ignite.cluster.ClusterNode;
import org.apache.ignite.internal.GridKernalContext;
+import org.apache.ignite.internal.IgniteDeploymentCheckedException;
import org.apache.ignite.internal.IgniteInternalFuture;
import org.apache.ignite.internal.IgniteInterruptedCheckedException;
import org.apache.ignite.internal.cluster.ClusterTopologyCheckedException;
@@ -214,22 +215,16 @@ public class DataStreamProcessor extends
GridProcessorAdapter {
if (req.forceLocalDeployment())
clsLdr = U.gridClassLoader();
else {
- GridDeployment dep = ctx.deploy().getGlobalDeployment(
- req.deploymentMode(),
- req.sampleClassName(),
- req.sampleClassName(),
- req.userVersion(),
- nodeId,
- req.classLoaderId(),
- req.participants(),
- null);
-
- if (dep == null) {
- sendResponse(nodeId,
- topic,
- req.requestId(),
+ GridDeployment dep;
+
+ try {
+ dep = ctx.deploy().globalDeployment(req.deploymentInfo(),
req.sampleClassName(), nodeId);
+ }
+ catch (IgniteDeploymentCheckedException e) {
+ // The sender waits for an answer, so a missing deployment
is reported back, not thrown.
+ sendResponse(nodeId, topic, req.requestId(),
new IgniteCheckedException("Failed to get deployment
for request [sndId=" + nodeId +
- ", req=" + req + ']'));
+ ", req=" + req + ']', e));
return;
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImpl.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImpl.java
index 1040f06ddf8..ec42ac560ae 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImpl.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImpl.java
@@ -2000,11 +2000,8 @@ public class DataStreamerImpl<K, V> implements
IgniteDataStreamer<K, V>, Delayed
true,
skipStore,
keepBinary,
- dep != null ? dep.deployMode() : null,
+ dep,
dep != null ? jobPda0.deployClass().getName() : null,
- dep != null ? dep.userVersion() : null,
- dep != null ? dep.participants() : null,
- dep != null ? dep.classLoaderId() : null,
dep == null,
topVer,
(rcvr == ISOLATED_UPDATER) ? partId : NO_STRIPE);
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerRequest.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerRequest.java
index 88af958a503..6d941ce969a 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerRequest.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerRequest.java
@@ -18,15 +18,13 @@
package org.apache.ignite.internal.processors.datastreamer;
import java.util.Collection;
-import java.util.Map;
-import java.util.UUID;
-import org.apache.ignite.configuration.DeploymentMode;
import org.apache.ignite.internal.DeferredUnmarshalMessage;
import org.apache.ignite.internal.Order;
import org.apache.ignite.internal.StripedMessage;
+import org.apache.ignite.internal.managers.deployment.GridDeploymentInfo;
+import
org.apache.ignite.internal.managers.deployment.GridDeploymentInfoMessage;
import org.apache.ignite.internal.processors.affinity.AffinityTopologyVersion;
import org.apache.ignite.internal.processors.cache.GridCacheUtils;
-import org.apache.ignite.internal.util.tostring.GridToStringInclude;
import org.apache.ignite.internal.util.typedef.internal.S;
import org.apache.ignite.lang.IgniteUuid;
import org.apache.ignite.plugin.extensions.communication.CacheIdAware;
@@ -70,9 +68,9 @@ public class DataStreamerRequest implements
DeferredUnmarshalMessage, CacheIdAwa
@Order(7)
boolean keepBinary;
- /** */
+ /** Deployment of the streamed classes. */
@Order(8)
- DeploymentMode depMode;
+ GridDeploymentInfoMessage depInfo;
/** */
@Order(9)
@@ -80,27 +78,14 @@ public class DataStreamerRequest implements
DeferredUnmarshalMessage, CacheIdAwa
/** */
@Order(10)
- String userVer;
-
- /** Node class loader participants. */
- @GridToStringInclude
- @Order(11)
- Map<UUID, IgniteUuid> ldrParticipants;
-
- /** */
- @Order(12)
- IgniteUuid clsLdrId;
-
- /** */
- @Order(13)
boolean forceLocDep;
/** Topology version. */
- @Order(14)
+ @Order(11)
AffinityTopologyVersion topVer;
/** */
- @Order(15)
+ @Order(12)
int partId;
/** Empty constructor. */
@@ -117,11 +102,8 @@ public class DataStreamerRequest implements
DeferredUnmarshalMessage, CacheIdAwa
* @param ignoreDepOwnership Ignore ownership.
* @param skipStore Skip store flag.
* @param keepBinary Keep binary flag.
- * @param depMode Deployment mode.
+ * @param depInfo Deployment of the streamed classes.
* @param sampleClsName Sample class name.
- * @param userVer User version.
- * @param ldrParticipants Loader participants.
- * @param clsLdrId Class loader ID.
* @param forceLocDep Force local deployment.
* @param topVer Topology version.
* @param partId Partition ID.
@@ -135,11 +117,8 @@ public class DataStreamerRequest implements
DeferredUnmarshalMessage, CacheIdAwa
boolean ignoreDepOwnership,
boolean skipStore,
boolean keepBinary,
- DeploymentMode depMode,
+ GridDeploymentInfo depInfo,
String sampleClsName,
- String userVer,
- Map<UUID, IgniteUuid> ldrParticipants,
- IgniteUuid clsLdrId,
boolean forceLocDep,
@NotNull AffinityTopologyVersion topVer,
int partId
@@ -154,11 +133,8 @@ public class DataStreamerRequest implements
DeferredUnmarshalMessage, CacheIdAwa
this.ignoreDepOwnership = ignoreDepOwnership;
this.skipStore = skipStore;
this.keepBinary = keepBinary;
- this.depMode = depMode;
+ this.depInfo = depInfo != null ? new
GridDeploymentInfoMessage(depInfo) : null;
this.sampleClsName = sampleClsName;
- this.userVer = userVer;
- this.ldrParticipants = ldrParticipants;
- this.clsLdrId = clsLdrId;
this.forceLocDep = forceLocDep;
this.topVer = topVer;
this.partId = partId;
@@ -204,9 +180,9 @@ public class DataStreamerRequest implements
DeferredUnmarshalMessage, CacheIdAwa
return keepBinary;
}
- /** @return Deployment mode. */
- DeploymentMode deploymentMode() {
- return depMode;
+ /** @return Deployment of the streamed classes. */
+ GridDeploymentInfo deploymentInfo() {
+ return depInfo;
}
/** @return Sample class name. */
@@ -214,21 +190,6 @@ public class DataStreamerRequest implements
DeferredUnmarshalMessage, CacheIdAwa
return sampleClsName;
}
- /** @return User version. */
- String userVersion() {
- return userVer;
- }
-
- /** @return Participants. */
- Map<UUID, IgniteUuid> participants() {
- return ldrParticipants;
- }
-
- /** @return Class loader ID. */
- IgniteUuid classLoaderId() {
- return clsLdrId;
- }
-
/** @return {@code True} to force local deployment. */
boolean forceLocalDeployment() {
return forceLocDep;
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/job/GridJobProcessor.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/job/GridJobProcessor.java
index b81a8b43bf2..096206004f8 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/job/GridJobProcessor.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/job/GridJobProcessor.java
@@ -1209,15 +1209,7 @@ public class GridJobProcessor extends
GridProcessorAdapter {
GridDeployment tmpDep = req.forceLocalDeployment() ?
ctx.deploy().getLocalDeployment(req.taskClassName()) :
- ctx.deploy().getGlobalDeployment(
- req.deploymentMode(),
- req.taskName(),
- req.taskClassName(),
- req.userVersion(),
- node.id(),
- req.classLoaderId(),
- req.loaderParticipants(),
- null);
+ ctx.deploy().globalDeployment(req.deploymentInfo(),
req.taskName(), req.taskClassName(), node.id());
if (tmpDep == null) {
if (log.isDebugEnabled())
@@ -1225,7 +1217,7 @@ public class GridJobProcessor extends
GridProcessorAdapter {
// Check local tasks.
for (Map.Entry<String, GridDeployment> d :
ctx.task().getUsedDeploymentMap().entrySet()) {
- if
(d.getValue().classLoaderId().equals(req.classLoaderId())) {
+ if
(d.getValue().classLoaderId().equals(req.deploymentInfo().classLoaderId())) {
assert d.getValue().local();
tmpDep = d.getValue();
@@ -1284,7 +1276,7 @@ public class GridJobProcessor extends
GridProcessorAdapter {
catch (IgniteCheckedException e) {
IgniteException ex = new IgniteException("Failed to
deserialize task attributes " +
"[taskName=" + req.taskName() + ", taskClsName=" +
req.taskClassName() +
- ", codeVer=" + req.userVersion() + ", taskClsLdr="
+ dep.classLoader() + ']', e);
+ ", codeVer=" + req.deploymentInfo().userVersion()
+ ", taskClsLdr=" + dep.classLoader() + ']', e);
U.error(log, ex.getMessage(), e);
@@ -1376,9 +1368,7 @@ public class GridJobProcessor extends
GridProcessorAdapter {
// Deployment is null.
IgniteException ex = new IgniteDeploymentException("Task
was not deployed or was redeployed since " +
"task execution [taskName=" + req.taskName() + ",
taskClsName=" + req.taskClassName() +
- ", codeVer=" + req.userVersion() + ", clsLdrId=" +
req.classLoaderId() +
- ", seqNum=" + req.classLoaderId().localId() + ",
depMode=" + req.deploymentMode() +
- ", dep=" + dep + ']');
+ ", dep=" + req.deploymentInfo() + ", resolved=" + dep
+ ']');
U.error(log, ex.getMessage(), ex);
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/task/GridTaskWorker.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/task/GridTaskWorker.java
index 0f262d928f5..5f62ee2b9f0 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/task/GridTaskWorker.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/task/GridTaskWorker.java
@@ -1385,7 +1385,7 @@ public class GridTaskWorker<T, R> extends GridWorker
implements GridTimeoutObjec
ses.getId(),
res.getJobContext().getJobId(),
ses.getTaskName(),
- ses.getUserVersion(),
+ dep,
ses.getTaskClassName(),
res.getJob(),
ses.getStartTime(),
@@ -1396,10 +1396,7 @@ public class GridTaskWorker<T, R> extends GridWorker
implements GridTimeoutObjec
sesAttrs,
jobAttrs,
ses.getCheckpointSpi(),
- dep.classLoaderId(),
- dep.deployMode(),
continuous,
- dep.participants(),
forceLocDep,
ses.isFullSupport(),
internal,
diff --git a/modules/core/src/main/resources/META-INF/classnames.properties
b/modules/core/src/main/resources/META-INF/classnames.properties
index fc1966fef4b..b9e5b18ebca 100644
--- a/modules/core/src/main/resources/META-INF/classnames.properties
+++ b/modules/core/src/main/resources/META-INF/classnames.properties
@@ -711,7 +711,7 @@
org.apache.ignite.internal.managers.communication.SessionChannelMessage
org.apache.ignite.internal.managers.communication.TransmissionCancelledException
org.apache.ignite.internal.managers.communication.TransmissionMeta
org.apache.ignite.internal.managers.communication.TransmissionPolicy
-org.apache.ignite.internal.managers.deployment.GridDeploymentInfoBean
+org.apache.ignite.internal.managers.deployment.GridDeploymentInfoMessage
org.apache.ignite.internal.managers.deployment.GridDeploymentPerVersionStore$2
org.apache.ignite.internal.managers.deployment.GridDeploymentRequest
org.apache.ignite.internal.managers.deployment.GridDeploymentResponse
diff --git
a/modules/core/src/test/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImplSelfTest.java
b/modules/core/src/test/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImplSelfTest.java
index 2bc777180e0..b6e8e1efdd1 100644
---
a/modules/core/src/test/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImplSelfTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImplSelfTest.java
@@ -697,11 +697,8 @@ public class DataStreamerImplSelfTest extends
GridCommonAbstractTest {
req.ignoreDeploymentOwnership(),
req.skipStore(),
req.keepBinary(),
- req.deploymentMode(),
+ req.deploymentInfo(),
req.sampleClassName(),
- req.userVersion(),
- req.participants(),
- req.classLoaderId(),
req.forceLocalDeployment(),
staleTop,
-1);
diff --git
a/modules/core/src/test/java/org/apache/ignite/p2p/ClassLoadingProblemExceptionTest.java
b/modules/core/src/test/java/org/apache/ignite/p2p/ClassLoadingProblemExceptionTest.java
index 1af74a63c33..cfab17bfa56 100644
---
a/modules/core/src/test/java/org/apache/ignite/p2p/ClassLoadingProblemExceptionTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/p2p/ClassLoadingProblemExceptionTest.java
@@ -38,7 +38,7 @@ import org.apache.ignite.configuration.IgniteConfiguration;
import org.apache.ignite.internal.IgniteEx;
import org.apache.ignite.internal.managers.communication.GridIoMessage;
import org.apache.ignite.internal.managers.deployment.GridDeployment;
-import org.apache.ignite.internal.managers.deployment.GridDeploymentInfoBean;
+import
org.apache.ignite.internal.managers.deployment.GridDeploymentInfoMessage;
import org.apache.ignite.internal.managers.deployment.GridDeploymentManager;
import org.apache.ignite.internal.managers.deployment.GridDeploymentMetadata;
import org.apache.ignite.internal.managers.deployment.GridDeploymentStore;
@@ -197,7 +197,7 @@ public class ClassLoadingProblemExceptionTest extends
GridCommonAbstractTest imp
GridCacheQueryRequest qryReq = (GridCacheQueryRequest)m;
if (qryReq.deployInfo() != null) {
- qryReq.deploy(new GridDeploymentInfoBean(
+ qryReq.deploy(new GridDeploymentInfoMessage(
IgniteUuid.fromUuid(UUID.randomUUID()),
qryReq.deployInfo().userVersion(),
qryReq.deployInfo().deployMode(),