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 83cfc507f6000..69e92e29ce0fe 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.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 08befbe25008a..34b4f45488a6b 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 void resetMetrics() { 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 GridIoMessage createGridIoMessage( 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 8296b60f81e92..c83f98feb9f40 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 @@ 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 ddbf0d3d96f7a..a997e9e6a048d 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.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 @@ private SecurityContext securityContext(UUID subjId) { @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 81777e5629fc2..79e869756c4ef 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 @@ 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 @@ /** * {@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 7c6899fec0eb7..870f8e1a56630 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 11f56e032bf6d..e38b9fd18ddbf 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 @@ * * @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> attrs = new ConcurrentSkipListMap<>(); + private volatile OperationContextAttribute[] 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 { * *

Registered attribute value is automatically captured and propagated between cluster nodes * during the messages transmission.

+ * + * @see DistributedAttributeRegistry */ - public void registerDistributedAttribute(int id, OperationContextAttribute attr) { + public synchronized void registerDistributedAttribute(int id, OperationContextAttribute 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[] locRegisteredAttrs = registeredAttrs; + + OperationContextAttribute[] 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 void registerDistributedAttribute(int id, OperationCo * * @see OperationContext#get(OperationContextAttribute) */ - public @Nullable OperationContextMessage collectDistributedAttributes() { - OperationContextMessage res = null; + public @Nullable OperationContextMessage collectDistributedAttributeValues() { + OperationContextAttribute[] locRegisteredAttrs = registeredAttrs; + + if (locRegisteredAttrs.length == 0) + return null; + + byte bitmap = 0; List vals = null; - for (Map.Entry> e : attrs.entrySet()) { - OperationContextAttribute attr = e.getValue(); + for (int id = 0; id < locRegisteredAttrs.length; id++) { + OperationContextAttribute 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[] 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 attr = (OperationContextAttribute)attrs.get(maskIdx++); + OperationContextAttribute attr = (OperationContextAttribute)locRegisteredAttrs[attrId++]; assert attr != null; @@ -139,7 +155,7 @@ public Scope restoreDistributedAttributes(@Nullable OperationContextMessage msg) } /** 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 9dc9fdb82ccea..01dea73efd58b 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 @@ 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 86a4f190026d2..def066f89ad93 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 @@ private class SocketWriter extends IgniteSpiThread { * @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 @@ private MessageWorker(IgniteLogger log) { ? (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 c95bba54b6c14..49b2f6e47e2c0 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 @@ void addMessage(TcpDiscoveryAbstractMessage msg, boolean ignoreHighPriority, boo } if (!fromSocket) - msg.opCtxMsg = operationCtxDispatcher.collectDistributedAttributes(); + msg.opCtxMsg = operationCtxDispatcher.collectDistributedAttributeValues(); if (msg instanceof TraceableMessage tMsg) { @@ -3293,7 +3293,7 @@ else if (msg instanceof TcpDiscoveryAuthFailedMessage) if (msg == WAKEUP) return; - try (Scope ignored = operationCtxDispatcher.restoreDistributedAttributes(msg.opCtxMsg)) { + try (Scope ignored = operationCtxDispatcher.restoreRemoteAttributeValues(msg.opCtxMsg)) { processMessage0(msg); } } @@ -6186,6 +6186,7 @@ private void processCustomMessage(TcpDiscoveryCustomEventMessage msg, boolean wa 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 5f09498060e2d..25e67f806100a 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.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 0000000000000..17bcf34b58a2d --- /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 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 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 PTR_ATTR = newInstance(DFLT_PTR); + + /** */ + public static final OperationContextAttribute 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 68b638bf610f0..0000000000000 --- 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 dAttr1 = - OperationContextAttribute.newInstance(new InetSocketAddressMessage(InetAddress.getLoopbackAddress(), 80)); - - OperationContextAttribute dAttr2 = OperationContextAttribute.newInstance(new GridIntList(1)); - - OperationContextAttribute 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 dAttr1, - InetSocketAddressMessage valToSend1, - OperationContextAttribute dAttr2, - GridIntList valToSend2 - ) throws Exception { - Set 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 dAttr1, - InetSocketAddressMessage valToSend1, - OperationContextAttribute 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 attr1, - @Nullable OperationContextAttribute 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 865d9ba164444..a96b9ab06d8fa 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.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 MessagesPluginProvider(Class... msgs) { /** {@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 f7c0c688926bb..e93cd40187c0f 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.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 @@ 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 4658c9aa75b85..36940bfaa0a96 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 a4a657d0877cf..6e0b545adae27 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.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 boolean knownNode(UUID nodeId) { /** */ 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 @@ private void notifyCustomEvent(final ZkDiscoveryCustomEventData evtData, Discove 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 d468897a0e3f3..b6b6226a79b13 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.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 @@ 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 c3d33bf581535..0000000000000 --- 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); - } -}