This is an automated email from the ASF dual-hosted git repository.
petrov-mg 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 47fb44979ef IGNITE-28902 Fixed Operation Context propagation for
discovery ack messages via TCP Discovery SPI (#13381)
47fb44979ef is described below
commit 47fb44979efe5be386859f77c00f66aff61a5a40
Author: Mikhail Petrov <[email protected]>
AuthorDate: Thu Jul 23 00:41:08 2026 +0300
IGNITE-28902 Fixed Operation Context propagation for discovery ack messages
via TCP Discovery SPI (#13381)
---
.../ignite/internal/CoreMessagesProvider.java | 1 +
.../managers/communication/GridIoManager.java | 4 +-
.../managers/communication/GridIoMessage.java | 2 +-
.../security/IgniteSecurityProcessor.java | 4 +-
.../security/SecurityContextWrapper.java | 6 +-
...bute.java => DistributedAttributeRegistry.java} | 21 +-
.../thread/context/OperationContextDispatcher.java | 88 ++++---
.../context}/OperationContextMessage.java | 15 +-
.../ignite/spi/discovery/tcp/ClientImpl.java | 4 +-
.../ignite/spi/discovery/tcp/ServerImpl.java | 5 +-
.../tcp/messages/TcpDiscoveryAbstractMessage.java | 2 +-
.../OperationContextAttributePropagationTest.java | 266 ++++++++++++++++++++
.../OperationContextSendAttributesTest.java | 278 ---------------------
.../apache/ignite/spi/MessagesPluginProvider.java | 9 +-
.../ignite/testsuites/SecurityTestSuite.java | 4 +-
.../ZkOperationContextAwareCustomMessage.java | 2 +-
.../zk/internal/ZookeeperDiscoveryImpl.java | 6 +-
.../zk/ZookeeperDiscoverySpiTestSuite4.java | 4 +-
.../ZkOperationContextSendAttributesTest.java | 32 ---
19 files changed, 364 insertions(+), 389 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 83cfc507f60..69e92e29ce0 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
@@ -259,6 +259,7 @@ import
org.apache.ignite.internal.processors.service.ServiceSingleNodeDeployment
import
org.apache.ignite.internal.processors.service.ServiceSingleNodeDeploymentResultBatch;
import org.apache.ignite.internal.processors.service.ServiceTopology;
import
org.apache.ignite.internal.processors.service.ServiceUndeploymentRequest;
+import org.apache.ignite.internal.thread.context.OperationContextMessage;
import org.apache.ignite.internal.util.GridByteArrayList;
import org.apache.ignite.internal.util.GridIntList;
import org.apache.ignite.internal.util.GridPartitionStateMap;
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 08befbe2500..34b4f45488a 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
@@ -459,7 +459,7 @@ public class GridIoManager extends
GridManagerAdapter<CommunicationSpi<Object>>
try {
GridIoMessage msg0 = (GridIoMessage)msg;
- try (Scope ignored =
ctx.operationContextDispatcher().restoreDistributedAttributes(msg0.opCtxMsg)) {
+ try (Scope ignored =
ctx.operationContextDispatcher().restoreRemoteAttributeValues(msg0.opCtxMsg)) {
onMessage0(nodeId, msg0, msgC);
}
}
@@ -2051,7 +2051,7 @@ public class GridIoManager extends
GridManagerAdapter<CommunicationSpi<Object>>
res = new GridIoMessage(plc, topic, msg, ordered, timeout,
skipOnTimeout);
- res.opCtxMsg =
ctx.operationContextDispatcher().collectDistributedAttributes();
+ res.opCtxMsg =
ctx.operationContextDispatcher().collectDistributedAttributeValues();
return res;
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoMessage.java
b/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoMessage.java
index 8296b60f81e..c83f98feb9f 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoMessage.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoMessage.java
@@ -19,11 +19,11 @@ package org.apache.ignite.internal.managers.communication;
import org.apache.ignite.internal.ExecutorAwareMessage;
import org.apache.ignite.internal.GridTopicMessage;
-import org.apache.ignite.internal.OperationContextMessage;
import org.apache.ignite.internal.Order;
import org.apache.ignite.internal.processors.cache.GridCacheMessage;
import org.apache.ignite.internal.processors.datastreamer.DataStreamerRequest;
import org.apache.ignite.internal.processors.tracing.messages.SpanTransport;
+import org.apache.ignite.internal.thread.context.OperationContextMessage;
import org.apache.ignite.internal.util.nio.GridNioServer.MessageWrapper;
import org.apache.ignite.internal.util.tostring.GridToStringInclude;
import org.apache.ignite.internal.util.typedef.internal.S;
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/security/IgniteSecurityProcessor.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/security/IgniteSecurityProcessor.java
index ddbf0d3d96f..a997e9e6a04 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/security/IgniteSecurityProcessor.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/security/IgniteSecurityProcessor.java
@@ -56,7 +56,7 @@ import static
org.apache.ignite.internal.processors.security.SecurityUtils.IGNIT
import static
org.apache.ignite.internal.processors.security.SecurityUtils.MSG_SEC_PROC_CLS_IS_INVALID;
import static
org.apache.ignite.internal.processors.security.SecurityUtils.hasSecurityManager;
import static
org.apache.ignite.internal.processors.security.SecurityUtils.nodeSecurityContext;
-import static
org.apache.ignite.internal.thread.context.DistributedOperationContextAttribute.SECURITY;
+import static
org.apache.ignite.internal.thread.context.DistributedAttributeRegistry.SECURITY;
import static
org.apache.ignite.plugin.security.SecurityPermission.ADMIN_USER_ACCESS;
import static
org.apache.ignite.plugin.security.SecurityPermission.JOIN_AS_SERVER;
@@ -253,7 +253,7 @@ public class IgniteSecurityProcessor extends
IgniteSecurityAdapter {
@Override public void start() throws IgniteCheckedException {
super.start();
-
ctx.operationContextDispatcher().registerDistributedAttribute(SECURITY.id(),
SEC_CTX_ATTR);
+
ctx.operationContextDispatcher().registerDistributedAttribute(SECURITY,
SEC_CTX_ATTR);
ctx.addNodeAttribute(ATTR_GRID_SEC_PROC_CLASS,
secPrc.getClass().getName());
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/security/SecurityContextWrapper.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/security/SecurityContextWrapper.java
index 81777e5629f..79e869756c4 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/security/SecurityContextWrapper.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/security/SecurityContextWrapper.java
@@ -19,7 +19,7 @@ package org.apache.ignite.internal.processors.security;
import java.util.UUID;
import org.apache.ignite.internal.Order;
-import
org.apache.ignite.internal.thread.context.DistributedOperationContextAttribute;
+import org.apache.ignite.internal.thread.context.DistributedAttributeRegistry;
import org.apache.ignite.internal.thread.context.OperationContextDispatcher;
import org.apache.ignite.plugin.extensions.communication.Message;
import org.apache.ignite.plugin.security.SecuritySubject;
@@ -27,8 +27,8 @@ import org.apache.ignite.plugin.security.SecuritySubject;
/**
* {@link SecurityContext} attribute value holder and message for {@link
SecuritySubject}'s id.
*
- * @see OperationContextDispatcher#collectDistributedAttributes()
- * @see DistributedOperationContextAttribute#SECURITY
+ * @see OperationContextDispatcher#collectDistributedAttributeValues()
+ * @see DistributedAttributeRegistry#SECURITY
*/
public class SecurityContextWrapper implements Message {
/** A value of {@link SecuritySubject#id()} */
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/thread/context/DistributedOperationContextAttribute.java
b/modules/core/src/main/java/org/apache/ignite/internal/thread/context/DistributedAttributeRegistry.java
similarity index 68%
rename from
modules/core/src/main/java/org/apache/ignite/internal/thread/context/DistributedOperationContextAttribute.java
rename to
modules/core/src/main/java/org/apache/ignite/internal/thread/context/DistributedAttributeRegistry.java
index 7c6899fec0e..870f8e1a566 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/thread/context/DistributedOperationContextAttribute.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/thread/context/DistributedAttributeRegistry.java
@@ -18,19 +18,12 @@
package org.apache.ignite.internal.thread.context;
import org.apache.ignite.internal.processors.security.SecurityContext;
-import org.apache.ignite.internal.processors.security.SecurityContextWrapper;
-/** Ids of Ignite's known distributed operation context attributes. */
-public enum DistributedOperationContextAttribute {
- /**
- * Distributed {@link SecurityContext}.
- *
- * @see SecurityContextWrapper
- */
- SECURITY;
-
- /** Cluster-wide id of distributed attribute. */
- public byte id() {
- return (byte)ordinal();
- }
+/**
+ * Declares reserved distributed IDs used to consistently identify {@link
OperationContext} attributes across
+ * all nodes in the cluster.
+ */
+public class DistributedAttributeRegistry {
+ /** Reserved for {@link SecurityContext} propagation. */
+ public static final byte SECURITY = 0;
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/thread/context/OperationContextDispatcher.java
b/modules/core/src/main/java/org/apache/ignite/internal/thread/context/OperationContextDispatcher.java
index 11f56e032bf..e38b9fd18dd 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/thread/context/OperationContextDispatcher.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/thread/context/OperationContextDispatcher.java
@@ -17,11 +17,9 @@
package org.apache.ignite.internal.thread.context;
import java.util.ArrayList;
+import java.util.Arrays;
import java.util.List;
-import java.util.Map;
-import java.util.concurrent.ConcurrentSkipListMap;
import org.apache.ignite.IgniteException;
-import org.apache.ignite.internal.OperationContextMessage;
import org.apache.ignite.internal.util.typedef.F;
import org.apache.ignite.plugin.extensions.communication.Message;
import org.jetbrains.annotations.Nullable;
@@ -44,17 +42,16 @@ import org.jetbrains.annotations.Nullable;
*
* @see OperationContext
* @see OperationContextMessage
- * @see DistributedOperationContextAttribute
*/
public class OperationContextDispatcher {
/** Maximal number of supported distributed attributes. */
static final byte MAX_ATTRS_CNT = Byte.SIZE;
/** Registered distributed attributes by their cluster-wide id. */
- private final Map<Byte, OperationContextAttribute<? extends Message>>
attrs = new ConcurrentSkipListMap<>();
+ private volatile OperationContextAttribute<? extends Message>[]
registeredAttrs = new OperationContextAttribute[0];
/** Whether the registration of new distributed attributes is allowed. */
- private volatile boolean regFinished;
+ private boolean regFinished;
/**
* Registers an attribute of {@link OperationContext} with the specified
distributed ID.
@@ -64,15 +61,27 @@ public class OperationContextDispatcher {
*
* <p>Registered attribute value is automatically captured and propagated
between cluster nodes
* during the messages transmission.</p>
+ *
+ * @see DistributedAttributeRegistry
*/
- public <T extends Message> void registerDistributedAttribute(int id,
OperationContextAttribute<T> attr) {
+ public synchronized <T extends Message> void
registerDistributedAttribute(int id, OperationContextAttribute<T> attr) {
if (regFinished)
throw new IgniteException("Initialization of distributed operation
context attributes has already finished.");
- assert id >= 0 && id < MAX_ATTRS_CNT : "Invalid distributed attributed
id [id=" + id + ']';
+ assert 0 <= id && id < MAX_ATTRS_CNT : "Invalid distributed attributed
id [id=" + id + ']';
+
+ OperationContextAttribute<? extends Message>[] locRegisteredAttrs =
registeredAttrs;
+
+ OperationContextAttribute<? extends Message>[] copy = Arrays.copyOf(
+ locRegisteredAttrs,
+ Math.max(locRegisteredAttrs.length, id + 1));
- if (attrs.putIfAbsent((byte)id, attr) != null)
+ if (copy[id] != null)
throw new IgniteException("Duplicated distributed attribute id
[id=" + id + ']');
+
+ copy[id] = attr;
+
+ registeredAttrs = copy;
}
/**
@@ -80,55 +89,62 @@ public class OperationContextDispatcher {
*
* @see OperationContext#get(OperationContextAttribute)
*/
- public @Nullable OperationContextMessage collectDistributedAttributes() {
- OperationContextMessage res = null;
+ public @Nullable OperationContextMessage
collectDistributedAttributeValues() {
+ OperationContextAttribute<? extends Message>[] locRegisteredAttrs =
registeredAttrs;
+
+ if (locRegisteredAttrs.length == 0)
+ return null;
+
+ byte bitmap = 0;
List<Message> vals = null;
- for (Map.Entry<Byte, OperationContextAttribute<? extends Message>> e :
attrs.entrySet()) {
- OperationContextAttribute<? extends Message> attr = e.getValue();
+ for (int id = 0; id < locRegisteredAttrs.length; id++) {
+ OperationContextAttribute<? extends Message> attr =
locRegisteredAttrs[id];
+
+ if (attr == null)
+ continue;
Message curVal = OperationContext.get(attr);
- if (curVal != attr.initialValue()) {
- if (res == null) {
- res = new OperationContextMessage();
+ if (curVal == attr.initialValue())
+ continue;
- vals = new ArrayList<>(MAX_ATTRS_CNT / 2);
- }
+ if (vals == null)
+ vals = new ArrayList<>(MAX_ATTRS_CNT / 2);
- byte mask = (byte)(1 << e.getKey());
+ byte mask = (byte)(1 << id);
- assert (res.idBitmap & mask) == 0;
+ assert (bitmap & mask) == 0;
- vals.add(curVal);
- res.idBitmap |= mask;
- }
+ vals.add(curVal);
+ bitmap |= mask;
}
- if (res != null)
- res.vals = vals.toArray(new Message[vals.size()]);
-
- return res;
+ return bitmap == 0 ? null : new OperationContextMessage(bitmap,
vals.toArray(Message[]::new));
}
/** Restores distributed {@link OperationContextAttribute} values received
from a remote node. */
- public Scope restoreDistributedAttributes(@Nullable
OperationContextMessage msg) {
+ public Scope restoreRemoteAttributeValues(@Nullable
OperationContextMessage msg) {
if (msg == null)
return Scope.NOOP_SCOPE;
+ OperationContextAttribute<? extends Message>[] locRegisteredAttrs =
registeredAttrs;
+
assert msg.idBitmap != 0;
- assert !F.isEmpty(msg.vals);
- assert msg.vals.length <= MAX_ATTRS_CNT;
+ assert !F.isEmpty(msg.attrs);
+ assert msg.attrs.length <= MAX_ATTRS_CNT;
OperationContext.ContextUpdater updater =
OperationContext.ContextUpdater.create();
- for (byte valIdx = 0, maskIdx = 0; valIdx < msg.vals.length; ++valIdx)
{
- Message curVal = msg.vals[valIdx];
+ for (byte valIdx = 0, attrId = 0; valIdx < msg.attrs.length; ++valIdx)
{
+ Message curVal = msg.attrs[valIdx];
+
+ while ((msg.idBitmap & (1 << attrId)) == 0)
+ ++attrId;
- while ((msg.idBitmap & (1 << maskIdx)) == 0)
- ++maskIdx;
+ assert attrId < locRegisteredAttrs.length;
- OperationContextAttribute<Message> attr =
(OperationContextAttribute<Message>)attrs.get(maskIdx++);
+ OperationContextAttribute<Message> attr =
(OperationContextAttribute<Message>)locRegisteredAttrs[attrId++];
assert attr != null;
@@ -139,7 +155,7 @@ public class OperationContextDispatcher {
}
/** Restricts further registration of distributed attributes. */
- public void finishRegistration() {
+ public synchronized void finishRegistration() {
regFinished = true;
}
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/OperationContextMessage.java
b/modules/core/src/main/java/org/apache/ignite/internal/thread/context/OperationContextMessage.java
similarity index 82%
rename from
modules/core/src/main/java/org/apache/ignite/internal/OperationContextMessage.java
rename to
modules/core/src/main/java/org/apache/ignite/internal/thread/context/OperationContextMessage.java
index 9dc9fdb82cc..01dea73efd5 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/OperationContextMessage.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/thread/context/OperationContextMessage.java
@@ -15,10 +15,9 @@
* limitations under the License.
*/
-package org.apache.ignite.internal;
+package org.apache.ignite.internal.thread.context;
-import org.apache.ignite.internal.thread.context.OperationContext;
-import org.apache.ignite.internal.thread.context.OperationContextDispatcher;
+import org.apache.ignite.internal.Order;
import org.apache.ignite.plugin.extensions.communication.Message;
/**
@@ -29,14 +28,20 @@ import
org.apache.ignite.plugin.extensions.communication.Message;
public class OperationContextMessage implements Message {
/** Values of operation context attributes. */
@Order(0)
- public Message[] vals;
+ Message[] attrs;
/** Bitmap of effective attributes ids. */
@Order(1)
- public byte idBitmap;
+ byte idBitmap;
/** Empty constructor for serialization purposes. */
public OperationContextMessage() {
// No-op.
}
+
+ /** */
+ public OperationContextMessage(byte idBitmap, Message[] attrs) {
+ this.attrs = attrs;
+ this.idBitmap = idBitmap;
+ }
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ClientImpl.java
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ClientImpl.java
index 86a4f190026..def066f89ad 100644
---
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ClientImpl.java
+++
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ClientImpl.java
@@ -1311,7 +1311,7 @@ class ClientImpl extends TcpDiscoveryImpl {
* @param msg Message.
*/
private void sendMessage(TcpDiscoveryAbstractMessage msg) {
- msg.opCtxMsg =
operationCtxDispatcher.collectDistributedAttributes();
+ msg.opCtxMsg =
operationCtxDispatcher.collectDistributedAttributeValues();
synchronized (mux) {
queue.add(msg);
@@ -1764,7 +1764,7 @@ class ClientImpl extends TcpDiscoveryImpl {
? (TcpDiscoveryAbstractMessage)msg
: null;
- try (Scope ignored =
operationCtxDispatcher.restoreDistributedAttributes(dm == null ? null :
dm.opCtxMsg)) {
+ try (Scope ignored =
operationCtxDispatcher.restoreRemoteAttributeValues(dm == null ? null :
dm.opCtxMsg)) {
if (msg instanceof JoinTimeout) {
int joinCnt0 = ((JoinTimeout)msg).joinCnt;
diff --git
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ServerImpl.java
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ServerImpl.java
index c95bba54b6c..49b2f6e47e2 100644
---
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ServerImpl.java
+++
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ServerImpl.java
@@ -3021,7 +3021,7 @@ class ServerImpl extends TcpDiscoveryImpl {
}
if (!fromSocket)
- msg.opCtxMsg =
operationCtxDispatcher.collectDistributedAttributes();
+ msg.opCtxMsg =
operationCtxDispatcher.collectDistributedAttributeValues();
if (msg instanceof TraceableMessage tMsg) {
@@ -3293,7 +3293,7 @@ class ServerImpl extends TcpDiscoveryImpl {
if (msg == WAKEUP)
return;
- try (Scope ignored =
operationCtxDispatcher.restoreDistributedAttributes(msg.opCtxMsg)) {
+ try (Scope ignored =
operationCtxDispatcher.restoreRemoteAttributeValues(msg.opCtxMsg)) {
processMessage0(msg);
}
}
@@ -6186,6 +6186,7 @@ class ServerImpl extends TcpDiscoveryImpl {
getLocalNodeId(), nextMsg);
ackMsg.topologyVersion(msg.topologyVersion());
+ ackMsg.opCtxMsg =
operationCtxDispatcher.collectDistributedAttributeValues();
processCustomMessage(ackMsg, waitForNotification);
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/messages/TcpDiscoveryAbstractMessage.java
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/messages/TcpDiscoveryAbstractMessage.java
index 5f09498060e..25e67f80610 100644
---
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/messages/TcpDiscoveryAbstractMessage.java
+++
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/messages/TcpDiscoveryAbstractMessage.java
@@ -21,8 +21,8 @@ import java.io.Externalizable;
import java.util.HashSet;
import java.util.Set;
import java.util.UUID;
-import org.apache.ignite.internal.OperationContextMessage;
import org.apache.ignite.internal.Order;
+import org.apache.ignite.internal.thread.context.OperationContextMessage;
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;
diff --git
a/modules/core/src/test/java/org/apache/ignite/internal/thread/context/OperationContextAttributePropagationTest.java
b/modules/core/src/test/java/org/apache/ignite/internal/thread/context/OperationContextAttributePropagationTest.java
new file mode 100644
index 00000000000..17bcf34b58a
--- /dev/null
+++
b/modules/core/src/test/java/org/apache/ignite/internal/thread/context/OperationContextAttributePropagationTest.java
@@ -0,0 +1,266 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.ignite.internal.thread.context;
+
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.function.Consumer;
+import org.apache.ignite.Ignite;
+import org.apache.ignite.IgniteException;
+import org.apache.ignite.cluster.ClusterNode;
+import org.apache.ignite.configuration.IgniteConfiguration;
+import org.apache.ignite.internal.GridKernalContext;
+import org.apache.ignite.internal.IgniteEx;
+import org.apache.ignite.internal.managers.communication.GridMessageListener;
+import org.apache.ignite.internal.managers.communication.IgniteIoTestMessage;
+import org.apache.ignite.internal.processors.authentication.User;
+import org.apache.ignite.internal.processors.cache.persistence.wal.WALPointer;
+import
org.apache.ignite.internal.processors.security.TestDiscoveryAcknowledgeMessage;
+import org.apache.ignite.internal.processors.security.TestDiscoveryMessage;
+import org.apache.ignite.internal.util.typedef.G;
+import org.apache.ignite.plugin.AbstractTestPluginProvider;
+import org.apache.ignite.plugin.PluginContext;
+import org.apache.ignite.spi.MessagesPluginProvider;
+import org.apache.ignite.testframework.junits.common.GridCommonAbstractTest;
+import org.junit.Test;
+
+import static java.util.concurrent.TimeUnit.MILLISECONDS;
+import static org.apache.ignite.internal.GridTopic.TOPIC_IO_TEST;
+import static
org.apache.ignite.internal.thread.context.OperationContextAttribute.newInstance;
+import static
org.apache.ignite.internal.thread.context.OperationContextAttributePropagationTest.TestIgniteComponent.DFLT_PTR;
+import static
org.apache.ignite.internal.thread.context.OperationContextAttributePropagationTest.TestIgniteComponent.DFLT_USR;
+import static
org.apache.ignite.internal.thread.context.OperationContextAttributePropagationTest.TestIgniteComponent.PTR_ATTR;
+import static
org.apache.ignite.internal.thread.context.OperationContextAttributePropagationTest.TestIgniteComponent.USR_ATTR;
+import static
org.apache.ignite.internal.thread.context.OperationContextDispatcher.MAX_ATTRS_CNT;
+import static org.apache.ignite.testframework.GridTestUtils.assertThrows;
+import static
org.apache.ignite.testframework.GridTestUtils.assertThrowsAnyCause;
+import static org.apache.ignite.testframework.GridTestUtils.waitForCondition;
+
+/** */
+public class OperationContextAttributePropagationTest extends
GridCommonAbstractTest {
+ /** */
+ private volatile Consumer<Integer> discoMsgLsnr;
+
+ /** {@inheritDoc} */
+ @Override protected void afterTest() throws Exception {
+ super.afterTest();
+
+ stopAllGrids();
+ }
+
+ /** {@inheritDoc} */
+ @Override protected IgniteConfiguration getConfiguration(String
igniteInstanceName) throws Exception {
+ IgniteConfiguration cfg = super.getConfiguration(igniteInstanceName);
+
+ cfg.setPluginProviders(
+ new TestIgniteComponent(),
+ new MessagesPluginProvider(TestDiscoveryMessage.class,
TestDiscoveryAcknowledgeMessage.class));
+
+ return cfg;
+ }
+
+ /** */
+ @Test
+ public void testSendAttributesByDiscovery() throws Exception {
+ prepareCluster();
+
+ doTestOperationContextAttributesPropagationThroughDiscovery(new
WALPointer(1, 1, 1), User.create("1", "1"));
+ }
+
+ /** */
+ @Test
+ public void testSendAttributesByCommunication() throws Exception {
+ prepareCluster();
+
+ doTestOperationContextAttributesPropagationThroughCommunication(new
WALPointer(1, 1, 1), User.create("1", "1"));
+ }
+
+ /** */
+ private void prepareCluster() throws Exception {
+ startGrids(2);
+ startClientGrid(2);
+
+ assertThrows(
+ null,
+ () ->
grid(0).context().operationContextDispatcher().registerDistributedAttribute(1,
null),
+ IgniteException.class,
+ "Initialization of distributed operation context attributes has
already finished"
+ );
+
+ for (int nodeIdx = 0; nodeIdx < 3; nodeIdx++) {
+ int finalNodeIdx = nodeIdx;
+
+
grid(nodeIdx).context().discovery().setCustomEventListener(TestDiscoveryMessage.class,
(topVer, snd, msg) -> {
+ if (discoMsgLsnr != null)
+ discoMsgLsnr.accept(finalNodeIdx);
+ });
+
+
grid(nodeIdx).context().discovery().setCustomEventListener(TestDiscoveryAcknowledgeMessage.class,
(topVer, snd, msg) -> {
+ if (discoMsgLsnr != null)
+ discoMsgLsnr.accept(finalNodeIdx);
+ });
+ }
+ }
+
+ /** */
+ private void
doTestOperationContextAttributesPropagationThroughDiscovery(WALPointer ptrVal,
User usrVal) throws Exception {
+ for (int nodeIdx = 0; nodeIdx < G.allGrids().size(); ++nodeIdx) {
+ try (Scope ignored = OperationContext.set(PTR_ATTR, ptrVal)) {
+ checkOperationContextDiscoveryTransmission(nodeIdx, ptrVal,
DFLT_USR);
+ }
+
+ try (Scope ignored = OperationContext.set(USR_ATTR, usrVal)) {
+ checkOperationContextDiscoveryTransmission(nodeIdx, DFLT_PTR,
usrVal);
+ }
+
+ try (Scope ignored = OperationContext.set(PTR_ATTR, ptrVal,
USR_ATTR, usrVal)) {
+ checkOperationContextDiscoveryTransmission(nodeIdx, ptrVal,
usrVal);
+ }
+
+ checkOperationContextDiscoveryTransmission(nodeIdx, DFLT_PTR,
DFLT_USR);
+ }
+ }
+
+ /** */
+ private void
doTestOperationContextAttributesPropagationThroughCommunication(WALPointer
ptrVal, User usrVal) throws Exception {
+ for (int fromIdx = 0; fromIdx < 3; ++fromIdx) {
+ for (int toIdx = 0; toIdx < 3; ++toIdx) {
+ if (fromIdx == toIdx)
+ continue;
+
+ try (Scope ignored = OperationContext.set(PTR_ATTR, ptrVal)) {
+ checkOperationContextCommunicationTransmission(fromIdx,
toIdx, ptrVal, DFLT_USR);
+ }
+
+ try (Scope ignored = OperationContext.set(USR_ATTR, usrVal)) {
+ checkOperationContextCommunicationTransmission(fromIdx,
toIdx, DFLT_PTR, usrVal);
+ }
+
+ try (Scope ignored = OperationContext.set(PTR_ATTR, ptrVal,
USR_ATTR, usrVal)) {
+ checkOperationContextCommunicationTransmission(fromIdx,
toIdx, ptrVal, usrVal);
+ }
+
+ checkOperationContextCommunicationTransmission(fromIdx, toIdx,
DFLT_PTR, DFLT_USR);
+ }
+ }
+ }
+
+ /** */
+ private void checkOperationContextDiscoveryTransmission(int sndIdx,
WALPointer expPtrVal, User expUsrVal) throws Exception {
+ Map<Integer, AtomicInteger> checkedNodes = new ConcurrentHashMap<>();
+
+ discoMsgLsnr = nodeIdx -> {
+ assertEquals(expUsrVal, OperationContext.get(USR_ATTR));
+ assertEquals(expPtrVal, OperationContext.get(PTR_ATTR));
+
+ checkedNodes.computeIfAbsent(nodeIdx, k -> new
AtomicInteger()).incrementAndGet();
+ };
+
+ try {
+ grid(sndIdx).context().discovery().sendCustomEvent(new
TestDiscoveryMessage());
+
+ assertTrue(waitForCondition(() ->
+ checkedNodes.size() == 3 &&
+
checkedNodes.values().stream().mapToInt(AtomicInteger::get).allMatch(v -> v ==
2),
+ getTestTimeout(),
+ 50));
+ }
+ finally {
+ discoMsgLsnr = null;
+ }
+ }
+
+ /** */
+ private void checkOperationContextCommunicationTransmission(
+ int fromIdx,
+ int toIdx,
+ WALPointer expPtrVal,
+ User expUsrVal
+ ) throws Exception {
+ IgniteEx from = grid(fromIdx);
+ IgniteEx to = grid(toIdx);
+
+ CountDownLatch rcvLatch = new CountDownLatch(2);
+
+ GridMessageListener lsnr = (nodeId, msg, plc) -> {
+ if (msg instanceof IgniteIoTestMessage &&
((IgniteIoTestMessage)msg).request()) {
+ assertEquals(expUsrVal, OperationContext.get(USR_ATTR));
+ assertEquals(expPtrVal, OperationContext.get(PTR_ATTR));
+
+ rcvLatch.countDown();
+ }
+ };
+
+ to.context().io().addMessageListener(TOPIC_IO_TEST, lsnr);
+
+ try {
+ from.context().io().sendIoTest(node(from, to), null, false);
+ from.context().io().sendIoTest(node(from, to), null, true);
+
+ assertTrue(rcvLatch.await(getTestTimeout(), MILLISECONDS));
+ }
+ finally {
+ assertTrue(to.context().io().removeMessageListener(TOPIC_IO_TEST,
lsnr));
+ }
+ }
+
+ /** Prevents {@link ClusterNode#isLocal()} to be negative. */
+ private ClusterNode node(Ignite from, Ignite to) {
+ return from.cluster().node(((IgniteEx)to).localNode().id());
+ }
+
+ /** */
+ static class TestIgniteComponent extends AbstractTestPluginProvider {
+ /** */
+ public static final WALPointer DFLT_PTR = new WALPointer(0, 0, 0);
+
+ /** */
+ public static final User DFLT_USR = User.create("0", "0");
+
+ /** */
+ public static final OperationContextAttribute<WALPointer> PTR_ATTR =
newInstance(DFLT_PTR);
+
+ /** */
+ public static final OperationContextAttribute<User> USR_ATTR =
newInstance(DFLT_USR);
+
+ /** {@inheritDoc} */
+ @Override public String name() {
+ return "TestDistributedOperationContextAttributesRegistrator";
+ }
+
+ /** {@inheritDoc} */
+ @Override public void start(PluginContext ctx) {
+ GridKernalContext kctx = ((IgniteEx)ctx.grid()).context();
+
+
kctx.operationContextDispatcher().registerDistributedAttribute(MAX_ATTRS_CNT -
1, USR_ATTR);
+ kctx.operationContextDispatcher().registerDistributedAttribute(0,
PTR_ATTR);
+
+ assertThrowsAnyCause(
+ log,
+ () -> {
+
kctx.operationContextDispatcher().registerDistributedAttribute(MAX_ATTRS_CNT -
1, PTR_ATTR);
+ return null;
+
+ }, IgniteException.class,
+ "Duplicated distributed attribute id"
+ );
+ }
+ }
+}
diff --git
a/modules/core/src/test/java/org/apache/ignite/internal/thread/context/OperationContextSendAttributesTest.java
b/modules/core/src/test/java/org/apache/ignite/internal/thread/context/OperationContextSendAttributesTest.java
deleted file mode 100644
index 68b638bf610..00000000000
---
a/modules/core/src/test/java/org/apache/ignite/internal/thread/context/OperationContextSendAttributesTest.java
+++ /dev/null
@@ -1,278 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements. See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.apache.ignite.internal.thread.context;
-
-import java.net.InetAddress;
-import java.util.Set;
-import java.util.UUID;
-import java.util.concurrent.ConcurrentHashMap;
-import java.util.concurrent.CountDownLatch;
-import org.apache.ignite.Ignite;
-import org.apache.ignite.IgniteException;
-import org.apache.ignite.cluster.ClusterNode;
-import org.apache.ignite.configuration.IgniteConfiguration;
-import org.apache.ignite.internal.GridKernalContext;
-import org.apache.ignite.internal.GridTopic;
-import org.apache.ignite.internal.IgniteEx;
-import org.apache.ignite.internal.managers.communication.GridMessageListener;
-import org.apache.ignite.internal.managers.communication.IgniteIoTestMessage;
-import org.apache.ignite.internal.managers.discovery.CustomEventListener;
-import org.apache.ignite.internal.processors.affinity.AffinityTopologyVersion;
-import org.apache.ignite.internal.processors.cache.DynamicCacheChangeBatch;
-import org.apache.ignite.internal.util.GridByteArrayList;
-import org.apache.ignite.internal.util.GridIntList;
-import org.apache.ignite.internal.util.typedef.G;
-import org.apache.ignite.plugin.AbstractTestPluginProvider;
-import org.apache.ignite.plugin.PluginContext;
-import org.apache.ignite.plugin.PluginProvider;
-import org.apache.ignite.spi.discovery.tcp.messages.InetSocketAddressMessage;
-import org.apache.ignite.testframework.junits.common.GridCommonAbstractTest;
-import org.junit.Test;
-import org.springframework.lang.Nullable;
-
-import static java.util.concurrent.TimeUnit.MILLISECONDS;
-import static org.apache.ignite.testframework.GridTestUtils.assertThrows;
-import static
org.apache.ignite.testframework.GridTestUtils.assertThrowsAnyCause;
-import static org.apache.ignite.testframework.GridTestUtils.waitForCondition;
-
-/** */
-public class OperationContextSendAttributesTest extends GridCommonAbstractTest
{
- /** */
- private PluginProvider pluginProvider;
-
- /** {@inheritDoc} */
- @Override protected void afterTest() throws Exception {
- super.afterTest();
-
- stopAllGrids();
- }
-
- /** {@inheritDoc} */
- @Override protected IgniteConfiguration getConfiguration(String
igniteInstanceName) throws Exception {
- IgniteConfiguration cfg = super.getConfiguration(igniteInstanceName);
-
- assert pluginProvider != null;
-
- cfg.setPluginProviders(pluginProvider);
-
- return cfg;
- }
-
- /** */
- @Test
- public void testSendAttributesByDiscovery() throws Exception {
- doTestOperationContextAttributesPropagation(true);
- }
-
- /** */
- @Test
- public void testSendAttributesByCommunication() throws Exception {
- doTestOperationContextAttributesPropagation(false);
- }
-
- /** */
- protected void doTestOperationContextAttributesPropagation(boolean
discovery) throws Exception {
- OperationContextAttribute<InetSocketAddressMessage> dAttr1 =
- OperationContextAttribute.newInstance(new
InetSocketAddressMessage(InetAddress.getLoopbackAddress(), 80));
-
- OperationContextAttribute<GridIntList> dAttr2 =
OperationContextAttribute.newInstance(new GridIntList(1));
-
- OperationContextAttribute<GridByteArrayList> otherTestAttr =
OperationContextAttribute.newInstance(new GridByteArrayList());
-
- pluginProvider = new AbstractTestPluginProvider() {
- @Override public String name() {
- return "TestDistributedOperationContextAttributesRegistrator";
- }
-
- @Override public void start(PluginContext ctx) {
- GridKernalContext kctx = ((IgniteEx)ctx.grid()).context();
-
- int dAttr1Id = OperationContextDispatcher.MAX_ATTRS_CNT - 2;
- int dAttr2Id = OperationContextDispatcher.MAX_ATTRS_CNT - 1;
-
-
kctx.operationContextDispatcher().registerDistributedAttribute(dAttr1Id,
dAttr1);
-
kctx.operationContextDispatcher().registerDistributedAttribute(dAttr2Id,
dAttr2);
-
- assertThrowsAnyCause(
- log,
- () -> {
-
kctx.operationContextDispatcher().registerDistributedAttribute(dAttr2Id,
otherTestAttr);
- return null;
-
- }, IgniteException.class,
- "Duplicated distributed attribute id"
- );
- }
- };
-
- // Local attribute 1.
- OperationContextAttribute.newInstance(1000);
-
- startGrids(2);
- startClientGrid(2);
-
- assertThrows(
- null,
- () ->
grid(0).context().operationContextDispatcher().registerDistributedAttribute(1,
null),
- IgniteException.class,
- "Initialization of distributed operation context attributes has
already finished"
- );
-
- // Local attribute 2.
- OperationContextAttribute.newInstance("locaAttr2");
-
- InetSocketAddressMessage valToSend1 = new
InetSocketAddressMessage(dAttr1.initialValue().address(), 443);
- GridIntList valToSend2 = new GridIntList(2);
-
- if (discovery)
-
doTestOperationContextAttributesPropagationThroughDiscovery(dAttr1, valToSend1,
dAttr2, valToSend2);
- else
-
doTestOperationContextAttributesPropagationThroughCommunication(dAttr1,
valToSend1, dAttr2, valToSend2);
- }
-
- /** */
- private void doTestOperationContextAttributesPropagationThroughDiscovery(
- OperationContextAttribute<InetSocketAddressMessage> dAttr1,
- InetSocketAddressMessage valToSend1,
- OperationContextAttribute<GridIntList> dAttr2,
- GridIntList valToSend2
- ) throws Exception {
- Set<Integer> checkedNodes = ConcurrentHashMap.newKeySet();
-
- for (int i = 0; i < G.allGrids().size(); ++i) {
- int i0 = i;
-
- grid(i).context().discovery().setCustomEventListener(
- DynamicCacheChangeBatch.class, new CustomEventListener<>() {
- @Override public void
onCustomEvent(AffinityTopologyVersion topVer, ClusterNode snd,
- DynamicCacheChangeBatch msg) {
-
- InetSocketAddressMessage receivedVal1 =
OperationContext.get(dAttr1);
- GridIntList receivedVal2 =
OperationContext.get(dAttr2);
-
- assertTrue(receivedVal1 != null && valToSend1.port()
== receivedVal1.port());
- assertTrue(receivedVal1 != null &&
valToSend1.address().equals(receivedVal1.address()));
-
- assertEquals(valToSend2, receivedVal2);
-
- checkedNodes.add(i0);
- }
- });
- }
-
- // Send from the coordinator.
- try (Scope ignored = OperationContext.set(dAttr1, valToSend1, dAttr2,
valToSend2)) {
- grid(0).createCache(defaultCacheConfiguration());
- }
-
- assertTrue(waitForCondition(() -> checkedNodes.size() == 3,
getTestTimeout(), 50));
- checkedNodes.clear();
-
- // Send from a server.
- try (Scope ignored = OperationContext.set(dAttr1, valToSend1, dAttr2,
valToSend2)) {
- grid(1).destroyCache(DEFAULT_CACHE_NAME);
- }
-
- assertTrue(waitForCondition(() -> checkedNodes.size() == 3,
getTestTimeout(), 50));
- checkedNodes.clear();
-
- // Send from a client.
- try (Scope ignored = OperationContext.set(dAttr1, valToSend1, dAttr2,
valToSend2)) {
- grid(2).createCache(defaultCacheConfiguration());
- }
-
- assertTrue(waitForCondition(() -> checkedNodes.size() == 3,
getTestTimeout(), 50));
- checkedNodes.clear();
- }
-
- /** */
- private void
doTestOperationContextAttributesPropagationThroughCommunication(
- OperationContextAttribute<InetSocketAddressMessage> dAttr1,
- InetSocketAddressMessage valToSend1,
- OperationContextAttribute<GridIntList> dAttr2,
- GridIntList valToSend2
- ) throws Exception {
- // Coordinator -> Server, Coordinator -> Client, Server -> Client,
Client -> Server, etc.
- for (int fromIdx = 0; fromIdx < 3; ++fromIdx) {
- for (int toIdx = 0; toIdx < 3; ++toIdx) {
- if (fromIdx == toIdx)
- continue;
-
- // One value.
- try (Scope ignored = OperationContext.set(dAttr1, valToSend1))
{
- checkOperationContextCommunicationTransmission(fromIdx,
toIdx, dAttr1, null);
- }
-
- // A couple of values.
- try (Scope ignored = OperationContext.set(dAttr1, valToSend1,
dAttr2, valToSend2)) {
- checkOperationContextCommunicationTransmission(fromIdx,
toIdx, dAttr1, dAttr2);
- }
- }
- }
- }
-
- /** */
- private void checkOperationContextCommunicationTransmission(
- int gridFromIdx,
- int gridToIdx,
- OperationContextAttribute<InetSocketAddressMessage> attr1,
- @Nullable OperationContextAttribute<GridIntList> attr2
- ) throws Exception {
- IgniteEx from = grid(gridFromIdx);
- IgniteEx to = grid(gridToIdx);
-
- CountDownLatch rcvLatch = new CountDownLatch(2);
-
- InetSocketAddressMessage expVal1 = OperationContext.get(attr1);
- GridIntList expVal2 = attr2 == null ? null :
OperationContext.get(attr2);
-
- GridMessageListener lsnr = new GridMessageListener() {
- @Override public void onMessage(UUID nodeId, Object msg, byte plc)
{
- if (msg instanceof IgniteIoTestMessage &&
((IgniteIoTestMessage)msg).request()) {
- InetSocketAddressMessage receivedVal1 =
OperationContext.get(attr1);
- GridIntList receivedVal2 = attr2 == null ? null :
OperationContext.get(attr2);
-
- assertTrue(receivedVal1 != null && expVal1.port() ==
receivedVal1.port());
- assertTrue(receivedVal1 != null &&
expVal1.address().equals(receivedVal1.address()));
-
- if (attr2 != null)
- assertEquals(expVal2, receivedVal2);
-
- rcvLatch.countDown();
- }
- }
- };
-
- to.context().io().addMessageListener(GridTopic.TOPIC_IO_TEST, lsnr);
-
- try {
- from.context().io().sendIoTest(node(from, to), null, false);
- from.context().io().sendIoTest(node(from, to), null, true);
-
- assertTrue(rcvLatch.await(getTestTimeout(), MILLISECONDS));
- }
- finally {
-
assertTrue(to.context().io().removeMessageListener(GridTopic.TOPIC_IO_TEST,
lsnr));
- }
- }
-
- /** Prevents {@link ClusterNode#isLocal()} to be negative. */
- private ClusterNode node(Ignite from, Ignite to) {
- return from.cluster().node(((IgniteEx)to).localNode().id());
- }
-}
diff --git
a/modules/core/src/test/java/org/apache/ignite/spi/MessagesPluginProvider.java
b/modules/core/src/test/java/org/apache/ignite/spi/MessagesPluginProvider.java
index 865d9ba1644..a96b9ab06d8 100644
---
a/modules/core/src/test/java/org/apache/ignite/spi/MessagesPluginProvider.java
+++
b/modules/core/src/test/java/org/apache/ignite/spi/MessagesPluginProvider.java
@@ -26,6 +26,7 @@ import org.apache.ignite.plugin.ExtensionRegistry;
import org.apache.ignite.plugin.PluginContext;
import org.apache.ignite.plugin.extensions.communication.Message;
import
org.apache.ignite.plugin.extensions.communication.MessageFactoryProvider;
+import org.apache.ignite.spi.discovery.DiscoverySpi;
import org.apache.ignite.spi.discovery.tcp.TestTcpDiscoverySpi;
import static org.apache.ignite.testframework.GridTestUtils.loadSerializer;
@@ -73,9 +74,11 @@ public class MessagesPluginProvider extends
AbstractTestPluginProvider {
/** {@inheritDoc} */
@Override public void start(PluginContext ctx) throws
IgniteCheckedException {
- // Register messages into the discovery protocol.
- TestTcpDiscoverySpi discoSpi =
(TestTcpDiscoverySpi)ctx.igniteConfiguration().getDiscoverySpi();
+ DiscoverySpi discoSpi = ctx.igniteConfiguration().getDiscoverySpi();
- discoSpi.messageFactory(msgFactoryProvider, ctx.igniteConfiguration());
+ if (discoSpi instanceof TestTcpDiscoverySpi testDiscoSpi) {
+ // Register messages into the discovery protocol.
+ testDiscoSpi.messageFactory(msgFactoryProvider,
ctx.igniteConfiguration());
+ }
}
}
diff --git
a/modules/core/src/test/java/org/apache/ignite/testsuites/SecurityTestSuite.java
b/modules/core/src/test/java/org/apache/ignite/testsuites/SecurityTestSuite.java
index f7c0c688926..e93cd40187c 100644
---
a/modules/core/src/test/java/org/apache/ignite/testsuites/SecurityTestSuite.java
+++
b/modules/core/src/test/java/org/apache/ignite/testsuites/SecurityTestSuite.java
@@ -73,8 +73,8 @@ import
org.apache.ignite.internal.processors.security.scheduler.SchedulerRemoteS
import
org.apache.ignite.internal.processors.security.service.ServiceAuthorizationTest;
import
org.apache.ignite.internal.processors.security.service.ServiceStaticConfigTest;
import
org.apache.ignite.internal.processors.security.snapshot.SnapshotPermissionCheckTest;
+import
org.apache.ignite.internal.thread.context.OperationContextAttributePropagationTest;
import
org.apache.ignite.internal.thread.context.OperationContextAttributesTest;
-import
org.apache.ignite.internal.thread.context.OperationContextSendAttributesTest;
import org.apache.ignite.ssl.MultipleSSLContextsTest;
import org.apache.ignite.tools.junit.JUnitTeamcityReporter;
import org.junit.BeforeClass;
@@ -148,7 +148,7 @@ import org.junit.runners.Suite;
SecurityContextInternalFuturePropagationTest.class,
NodeConnectionCertificateCapturingTest.class,
OperationContextAttributesTest.class,
- OperationContextSendAttributesTest.class,
+ OperationContextAttributePropagationTest.class,
})
public class SecurityTestSuite {
/** */
diff --git
a/modules/zookeeper/src/main/java/org/apache/ignite/spi/discovery/zk/internal/ZkOperationContextAwareCustomMessage.java
b/modules/zookeeper/src/main/java/org/apache/ignite/spi/discovery/zk/internal/ZkOperationContextAwareCustomMessage.java
index 4658c9aa75b..36940bfaa0a 100644
---
a/modules/zookeeper/src/main/java/org/apache/ignite/spi/discovery/zk/internal/ZkOperationContextAwareCustomMessage.java
+++
b/modules/zookeeper/src/main/java/org/apache/ignite/spi/discovery/zk/internal/ZkOperationContextAwareCustomMessage.java
@@ -17,10 +17,10 @@
package org.apache.ignite.spi.discovery.zk.internal;
-import org.apache.ignite.internal.OperationContextMessage;
import org.apache.ignite.internal.Order;
import org.apache.ignite.internal.thread.context.OperationContext;
import org.apache.ignite.internal.thread.context.OperationContextDispatcher;
+import org.apache.ignite.internal.thread.context.OperationContextMessage;
import org.apache.ignite.plugin.extensions.communication.MessageFactory;
import org.apache.ignite.spi.discovery.DiscoverySpiCustomMessage;
import org.apache.ignite.spi.discovery.tcp.TcpDiscoverySpi;
diff --git
a/modules/zookeeper/src/main/java/org/apache/ignite/spi/discovery/zk/internal/ZookeeperDiscoveryImpl.java
b/modules/zookeeper/src/main/java/org/apache/ignite/spi/discovery/zk/internal/ZookeeperDiscoveryImpl.java
index a4a657d0877..6e0b545adae 100644
---
a/modules/zookeeper/src/main/java/org/apache/ignite/spi/discovery/zk/internal/ZookeeperDiscoveryImpl.java
+++
b/modules/zookeeper/src/main/java/org/apache/ignite/spi/discovery/zk/internal/ZookeeperDiscoveryImpl.java
@@ -63,11 +63,11 @@ import
org.apache.ignite.internal.IgniteFutureTimeoutCheckedException;
import org.apache.ignite.internal.IgniteInternalFuture;
import org.apache.ignite.internal.IgniteKernal;
import org.apache.ignite.internal.IgnitionEx;
-import org.apache.ignite.internal.OperationContextMessage;
import org.apache.ignite.internal.cluster.ClusterTopologyCheckedException;
import org.apache.ignite.internal.events.DiscoveryCustomEvent;
import org.apache.ignite.internal.processors.security.SecurityContext;
import org.apache.ignite.internal.thread.context.OperationContextDispatcher;
+import org.apache.ignite.internal.thread.context.OperationContextMessage;
import org.apache.ignite.internal.thread.context.Scope;
import org.apache.ignite.internal.thread.pool.IgniteThreadPoolExecutor;
import org.apache.ignite.internal.util.GridLongList;
@@ -669,7 +669,7 @@ public class ZookeeperDiscoveryImpl {
/** */
public void sendCustomEvent(DiscoverySpiCustomMessage msg) {
- OperationContextMessage opCtx =
opCtxDispatcher.collectDistributedAttributes();
+ OperationContextMessage opCtx =
opCtxDispatcher.collectDistributedAttributeValues();
if (opCtx != null)
sendCustomMessage(new ZkOperationContextAwareCustomMessage(msg,
opCtx));
@@ -3548,7 +3548,7 @@ public class ZookeeperDiscoveryImpl {
IgniteFuture<?> fut;
- try (Scope ignored =
opCtxDispatcher.restoreDistributedAttributes(opCtxMsg)) {
+ try (Scope ignored =
opCtxDispatcher.restoreRemoteAttributeValues(opCtxMsg)) {
fut = lsnr.onDiscovery(
new DiscoveryNotification(
DiscoveryCustomEvent.EVT_DISCOVERY_CUSTOM_EVT,
diff --git
a/modules/zookeeper/src/test/java/org/apache/ignite/spi/discovery/zk/ZookeeperDiscoverySpiTestSuite4.java
b/modules/zookeeper/src/test/java/org/apache/ignite/spi/discovery/zk/ZookeeperDiscoverySpiTestSuite4.java
index d468897a0e3..b6b6226a79b 100644
---
a/modules/zookeeper/src/test/java/org/apache/ignite/spi/discovery/zk/ZookeeperDiscoverySpiTestSuite4.java
+++
b/modules/zookeeper/src/test/java/org/apache/ignite/spi/discovery/zk/ZookeeperDiscoverySpiTestSuite4.java
@@ -29,8 +29,8 @@ import
org.apache.ignite.internal.processors.metastorage.DistributedMetaStorageP
import
org.apache.ignite.internal.processors.metastorage.DistributedMetaStorageTest;
import
org.apache.ignite.internal.processors.security.cluster.ActivationOnJoinWithoutPermissionsWithPersistenceTest;
import
org.apache.ignite.internal.processors.security.cluster.NodeJoinPermissionsTest;
+import
org.apache.ignite.internal.thread.context.OperationContextAttributePropagationTest;
import org.apache.ignite.spi.discovery.DiscoverySpiDataExchangeTest;
-import
org.apache.ignite.spi.discovery.zk.internal.ZkOperationContextSendAttributesTest;
import org.junit.BeforeClass;
import org.junit.runner.RunWith;
import org.junit.runners.Suite;
@@ -51,7 +51,7 @@ import org.junit.runners.Suite;
DistributedMetaStoragePersistentTest.class,
IgniteNodeValidationFailedEventTest.class,
DiscoverySpiDataExchangeTest.class,
- ZkOperationContextSendAttributesTest.class,
+ OperationContextAttributePropagationTest.class,
CacheCreateDestroyEventSecurityContextTest.class,
NodeJoinPermissionsTest.class,
ActivationOnJoinWithoutPermissionsWithPersistenceTest.class,
diff --git
a/modules/zookeeper/src/test/java/org/apache/ignite/spi/discovery/zk/internal/ZkOperationContextSendAttributesTest.java
b/modules/zookeeper/src/test/java/org/apache/ignite/spi/discovery/zk/internal/ZkOperationContextSendAttributesTest.java
deleted file mode 100644
index c3d33bf5815..00000000000
---
a/modules/zookeeper/src/test/java/org/apache/ignite/spi/discovery/zk/internal/ZkOperationContextSendAttributesTest.java
+++ /dev/null
@@ -1,32 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements. See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.apache.ignite.spi.discovery.zk.internal;
-
-import
org.apache.ignite.internal.thread.context.OperationContextSendAttributesTest;
-
-import static org.junit.Assume.assumeTrue;
-
-/** */
-public class ZkOperationContextSendAttributesTest extends
OperationContextSendAttributesTest {
- /** {@inheritDoc} */
- @Override protected void
doTestOperationContextAttributesPropagation(boolean discovery) throws Exception
{
- assumeTrue(discovery);
-
- super.doTestOperationContextAttributesPropagation(true);
- }
-}