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 a460ed789d24b..3b066879eea24 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.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 CoreMessagesProvider(Marshaller dfltMarsh, Marshaller schemaAwareMarsh) { // [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 cad29b7b12909..8ebea159666f4 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.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 @@ private boolean filterDropsEvent(Event evt) { 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 @@ private boolean filterDropsEvent(Event evt) { 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 098e8514f681c..c431610e1ea44 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 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 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 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 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 top; /** */ @@ -151,23 +138,23 @@ public class GridJobExecuteRequest implements ExecutorAwareMessage, DeferredUnma IgnitePredicate 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 GridJobExecuteRequest() { * @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 GridJobExecuteRequest() { * @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 GridJobExecuteRequest( IgniteUuid sesId, IgniteUuid jobId, String taskName, - String userVer, + GridDeploymentInfo depInfo, String taskClsName, ComputeJob job, long startTaskTime, @@ -219,10 +203,7 @@ public GridJobExecuteRequest( Map sesAttrs, Map jobAttrs, String cpSpi, - IgniteUuid clsLdrId, - DeploymentMode depMode, boolean dynamicSiblings, - Map ldrParticipants, boolean forceLocDep, boolean sesFullSup, boolean internal, @@ -238,14 +219,12 @@ public GridJobExecuteRequest( 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 GridJobExecuteRequest( 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 IgniteUuid jobId() { return jobId; } + /** @return Deployment of the task classes. */ + public GridDeploymentInfo deploymentInfo() { + return depInfo; + } + /** * @return Task class name. */ @@ -299,13 +280,6 @@ public String taskName() { return taskName; } - /** - * @return Task version. - */ - public String userVersion() { - return userVer; - } - /** * @return Grid job. */ @@ -364,27 +338,6 @@ public String checkpointSpi() { 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 loaderParticipants() { - return ldrParticipants; - } - /** * @return Returns {@code true} if deployment should always be used. */ @@ -446,7 +399,6 @@ public AffinityTopologyVersion topologyVersion() { 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 e59873f4b449f..f3397e5d21485 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 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 GridMessageListenHandler(@Nullable Object topic, IgniteBiPredicate nodes, 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 @@ private class GridUserMessageListener implements GridMessageListener { 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 afe4db51c2b5b..ef4f0b61a8bca 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 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 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 GridIoUserMessage() { 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 GridIoUserMessage() { return depClsName; } - /** - * @return User version. - */ - @Nullable public String userVersion() { - return userVer; - } - - /** - * @return Node class loader participant map. - */ - @Nullable public Map loaderParticipants() { - return ldrParties != null ? Collections.unmodifiableMap(ldrParties) : null; - } - /** * @return Serialized message topic. */ @@ -182,6 +128,11 @@ public void body(Object body) { 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 ba9f48c3ca24a..4a4f4e4a6a92c 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.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 GridDeploymentInfoBean() { * @param depMode Deployment mode. * @param participants Participants. */ - public GridDeploymentInfoBean( + public GridDeploymentInfoMessage( IgniteUuid clsLdrId, String userVer, DeploymentMode depMode, @@ -79,7 +79,7 @@ public GridDeploymentInfoBean( /** * @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 GridDeploymentInfoBean(GridDeploymentInfo dep) { /** {@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 3069f201be9e8..5be5bd627c6a2 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.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 @@ private GridDeployment checkDeployment(GridDeployment deployment, String store) 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 @@ private GridDeployment checkDeployment(GridDeployment deployment, String store) * @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 @@ private GridDeployment checkDeployment(GridDeployment deployment, String store) String userVer, UUID sndNodeId, IgniteUuid clsLdrId, - Map participants, - @Nullable IgnitePredicate nodeFilter) { + Map participants) { if (locDep != null) return locDep; @@ -439,7 +505,6 @@ private GridDeployment checkDeployment(GridDeployment deployment, String store) 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 370825c5fee71..c5f6857be9cb9 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 @@ 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 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 void classLoader(ClassLoader clsLdr) { this.clsLdr = clsLdr; } - /** - * @param nodeFilter Node filter. - */ - public void nodeFilter(IgnitePredicate nodeFilter) { - this.nodeFilter = nodeFilter; - } - - /** - * @return Node filter. - */ - public IgnitePredicate 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 1c883d1474aa3..62ee043907a7d 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 @@ private List query(IgnitePredicate p, Collection 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)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 a2fd6a0c26c3c..322d8bc2bc14d 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 ldrParties; - /** */ public GridEventStorageRequest() { // No-op. @@ -77,24 +60,12 @@ public GridEventStorageRequest() { /** * @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 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 @@ IgnitePredicate filter() { 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 @@ String filterClassName() { return filterClsName; } - /** @return User version. */ - String userVersion() { - return userVer; - } - - /** @return Node class loader participant map. */ - @Nullable Map 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 3937be9b792a4..2eda968819377 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 @@ static Object unmarshall(GridKernalContext ctx, UUID sndNodeId, GridAffinityMess 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 125fe4ce94743..46a7c046ea721 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.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 void prepare(GridCacheDeployable deployable) throws IgnitePeerToPeerClass // 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 void prepare(GridCacheDeployable deployable) throws IgnitePeerToPeerClass * @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 @@ private void checkDeploymentIsCorrect(GridDeploymentInfoBean deployment, GridCac /** * @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 @@ private void checkDeploymentIsCorrect(GridDeploymentInfoBean deployment, GridCac 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 d : deps.values()) { - GridDeploymentInfoBean deploymentInfoBean = getDepBean(d); + GridDeploymentInfoMessage deploymentInfoBean = getDepBean(d); if (deploymentInfoBean != null) return deploymentInfoBean; } @@ -472,7 +472,7 @@ private void checkDeploymentIsCorrect(GridDeploymentInfoBean deployment, GridCac } /** */ - @Nullable private GridDeploymentInfoBean getDepBean(CachedDeploymentInfo d) { + @Nullable private GridDeploymentInfoMessage getDepBean(CachedDeploymentInfo d) { if (d == null || cctx.discovery().node(d.senderId()) == null) // Sender has left. return null; @@ -484,7 +484,7 @@ private void checkDeploymentIsCorrect(GridDeploymentInfoBean deployment, GridCac 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 @@ else if (err == null) userVer, sndId, ldrId, - participants, - F.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 2f09c39eb9fc9..8cdceed58e393 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.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 final void deploy(GridDeploymentInfo depInfo) { 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 final void deploy(GridDeploymentInfo depInfo) { * @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 c4e4005095ff4..5f31963ba210a 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.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 @@ protected CacheContinuousQueryDeployableObject(Object obj, GridKernalContext ctx 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 @@ protected CacheContinuousQueryDeployableObject(Object obj, GridKernalContext ctx 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 9edaab63c7a3b..111d3b740493f 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.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 @@ private AbstractContinuousMessage createStartMessage(UUID routineId, 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 @@ private AbstractContinuousMessage createStartMessage(UUID routineId, if (dep0 == null) throw new IgniteDeploymentCheckedException("Failed to deploy projection predicate: " + nodeFilter); - dep = new GridDeploymentInfoBean(dep0); + dep = new GridDeploymentInfoMessage(dep0); } } @@ -981,7 +981,15 @@ private AbstractContinuousMessage createStartMessage(UUID routineId, 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 @@ private void processStartAckRequest(AffinityTopologyVersion topVer, } } + /** + * 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 @@ private void processStartRequestMutable(ClusterNode node, StartRoutineDiscoveryM 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 @@ private void processStartRequestImmutable(final AffinityTopologyVersion topVer, 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 9571c231bbf9a..595aaa56c5865 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 @@ * Start request data. */ public class StartRequestData implements Message { - /** Node filter. */ - private IgnitePredicate nodeFilter; + /** Node filter, restored by the processor reading this request. */ + IgnitePredicate 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 void className(String clsName) { /** * @param depInfo New deployment info. */ - public void deploymentInfo(GridDeploymentInfoBean depInfo) { + public void deploymentInfo(GridDeploymentInfoMessage depInfo) { this.depInfo = depInfo; } @@ -169,58 +162,4 @@ public void autoUnsubscribe(boolean autoUnsubscribe) { @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 464a74d82ee25..92ec36cc210a6 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 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 @@ private void processRequest(final UUID nodeId, final DataStreamerRequest req) { 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 1040f06ddf83b..ec42ac560ae84 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 @@ private void submit( 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 88af958a50373..6d941ce969a2f 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 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 DataStreamerRequest() { * @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 DataStreamerRequest( boolean ignoreDepOwnership, boolean skipStore, boolean keepBinary, - DeploymentMode depMode, + GridDeploymentInfo depInfo, String sampleClsName, - String userVer, - Map ldrParticipants, - IgniteUuid clsLdrId, boolean forceLocDep, @NotNull AffinityTopologyVersion topVer, int partId @@ -154,11 +133,8 @@ public DataStreamerRequest( 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 @@ boolean keepBinary() { 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 @@ String sampleClassName() { return sampleClsName; } - /** @return User version. */ - String userVersion() { - return userVer; - } - - /** @return Participants. */ - Map 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 b81a8b43bf21e..096206004f8d3 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 void processJobExecuteRequest(ClusterNode node, final GridJobExecuteReque 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 void processJobExecuteRequest(ClusterNode node, final GridJobExecuteReque // Check local tasks. for (Map.Entry 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 void processJobExecuteRequest(ClusterNode node, final GridJobExecuteReque 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 @@ else if (jobAlwaysActivate) { // 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 0f262d928f5a9..5f62ee2b9f04e 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 @@ private void sendRequest(ComputeJobResult res) { ses.getId(), res.getJobContext().getJobId(), ses.getTaskName(), - ses.getUserVersion(), + dep, ses.getTaskClassName(), res.getJob(), ses.getStartTime(), @@ -1396,10 +1396,7 @@ private void sendRequest(ComputeJobResult res) { 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 fc1966fef4b4e..b9e5b18ebcaef 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 2bc777180e0ff..b6e8e1efdd13c 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 @@ private static class StaleTopologyCommunicationSpi extends TcpCommunicationSpi { 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 1af74a63c3340..cfab17bfa5668 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.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 @@ private class TestCommunicationSpi extends TcpCommunicationSpi { 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(),