diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheAdapter.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheAdapter.java index 6b32f9408fb02..74285ab5cc957 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheAdapter.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheAdapter.java @@ -2991,7 +2991,9 @@ public CacheMetricsImpl metrics0() { e = ex; } - throw new NodeStoppingException(e); + throw e == null + ? new NodeStoppingException("Failed to acquire lock (node is stopping) [keys=" + keys + ']') + : new NodeStoppingException(e); } finally { if (isInterrupted) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheExplicitLockSpan.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheExplicitLockSpan.java index 53c8ad6f5774d..3d4bff80a533b 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheExplicitLockSpan.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheExplicitLockSpan.java @@ -236,17 +236,21 @@ public Collection candidates() { * Marks all candidates added for given key as owned. * * @param key Key. + * @return {@code True} if candidate owned, {@code false} otherwise. */ - public void markOwned(IgniteTxKey key) { + public boolean markOwned(IgniteTxKey key) { lock(); try { Deque deque = cands.get(key); - assert deque != null; + if (deque == null) + return false; for (GridCacheMvccCandidate cand : deque) cand.setOwner(); + + return true; } finally { unlock(); diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheMvccManager.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheMvccManager.java index 9fe3462af0546..90f19c05463b0 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheMvccManager.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheMvccManager.java @@ -1031,14 +1031,14 @@ public boolean isLockedByThread(IgniteTxKey key, long threadId) { * * @param key Key. * @param threadId Thread id. + * @return {@code False} if there is no explicit lock candidate for the given thread and key. */ - public void markExplicitOwner(IgniteTxKey key, long threadId) { + public boolean markExplicitOwner(IgniteTxKey key, long threadId) { assert threadId > 0; GridCacheExplicitLockSpan span = pendingExplicit.get(threadId); - if (span != null) - span.markOwned(key); + return span != null && span.markOwned(key); } /** diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/colocated/GridDhtColocatedLockFuture.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/colocated/GridDhtColocatedLockFuture.java index 3b7158b4e5dd6..9d4e07a3070fe 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/colocated/GridDhtColocatedLockFuture.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/colocated/GridDhtColocatedLockFuture.java @@ -1329,10 +1329,18 @@ private void markLocalDhtLocksAcquired(Collection keys) { } else { for (KeyCacheObject key : keys) - cctx.mvcc().markExplicitOwner(cctx.txKey(key), threadId); + markExplicitOwner(key); } } + /** */ + private void markExplicitOwner(KeyCacheObject key) { + boolean marked = cctx.mvcc().markExplicitOwner(cctx.txKey(key), threadId); + + assert marked || isDone() : "Explicit lock candidate not found for active lock future [key=" + key + + ", fut=" + this + ']'; + } + /** * Tries to map this future in assumption that local node is primary for all keys passed in. * If node is not primary for one of the keys, then mapping is reverted and full remote mapping is performed. @@ -1775,7 +1783,7 @@ void onResult(GridNearLockResponse res) { log.debug("Processed response for entry [res=" + res + ", entry=" + entry + ']'); } else - cctx.mvcc().markExplicitOwner(cctx.txKey(k), threadId); + markExplicitOwner(k); if (retval && cctx.events().isRecordable(EVT_CACHE_OBJECT_READ)) { cctx.events().addEvent(cctx.affinity().partition(k), diff --git a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/distributed/dht/ExplicitLockCancelOnNodeStopTest.java b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/distributed/dht/ExplicitLockCancelOnNodeStopTest.java new file mode 100644 index 0000000000000..7caf82227057d --- /dev/null +++ b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/distributed/dht/ExplicitLockCancelOnNodeStopTest.java @@ -0,0 +1,121 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.ignite.internal.processors.cache.distributed.dht; + +import java.util.List; +import java.util.concurrent.atomic.AtomicReference; +import org.apache.ignite.Ignite; +import org.apache.ignite.IgniteCache; +import org.apache.ignite.configuration.CacheConfiguration; +import org.apache.ignite.configuration.IgniteConfiguration; +import org.apache.ignite.failure.FailureHandler; +import org.apache.ignite.failure.TestFailureHandler; +import org.apache.ignite.internal.NodeStoppingException; +import org.apache.ignite.internal.TestRecordingCommunicationSpi; +import org.apache.ignite.internal.util.lang.RunnableX; +import org.apache.ignite.internal.util.typedef.X; +import org.apache.ignite.lifecycle.LifecycleEventType; +import org.apache.ignite.testframework.junits.common.GridCommonAbstractTest; +import org.junit.Test; + +import static org.apache.ignite.cache.CacheAtomicityMode.TRANSACTIONAL; +import static org.apache.ignite.internal.TestRecordingCommunicationSpi.spi; + +/** + * Tests that an explicit lock request cancelled on a stopping node does not cause a critical failure + * when the response for the in-flight local DHT lock arrives after the cancellation. + */ +public class ExplicitLockCancelOnNodeStopTest extends GridCommonAbstractTest { + /** */ + private volatile RunnableX beforeStop; + + /** */ + private final TestFailureHandler failureHnd = new TestFailureHandler(false); + + /** {@inheritDoc} */ + @Override protected IgniteConfiguration getConfiguration(String igniteInstanceName) throws Exception { + return super.getConfiguration(igniteInstanceName) + .setCommunicationSpi(new TestRecordingCommunicationSpi()) + .setCacheConfiguration(new CacheConfiguration<>(DEFAULT_CACHE_NAME) + .setAtomicityMode(TRANSACTIONAL) + .setBackups(1)) + .setLifecycleBeans(evt -> { + if (evt == LifecycleEventType.BEFORE_NODE_STOP && getTestIgniteInstanceName(0).equals(igniteInstanceName)) + beforeStop.run(); + }); + } + + /** {@inheritDoc} */ + @Override protected FailureHandler getFailureHandler(String igniteInstanceName) { + return failureHnd; + } + + /** {@inheritDoc} */ + @Override protected void afterTest() throws Exception { + stopAllGrids(); + + super.afterTest(); + } + + /** + * Scenario: + *
    + *
  1. A thread holds an explicit lock on one key so its explicit lock span stays non-empty.
  2. + *
  3. The node starts stopping and, from the {@code BEFORE_NODE_STOP} callback, the same thread tries to + * lock another key. The local DHT lock is acquired and a lock request is sent to the backup.
  4. + *
  5. Since the node is already stopping, the lock future is cancelled immediately which removes + * the explicit lock candidate from the span.
  6. + *
  7. The backup response arrives afterwards and the completed DHT lock future tries to mark + * the already removed candidate as owned.
  8. + *
+ */ + @Test + public void testLockOnStoppingNode() throws Exception { + Ignite srv = startGrid(0); + Ignite backup = startGrid(1); + + awaitPartitionMapExchange(); + + IgniteCache cache = srv.cache(DEFAULT_CACHE_NAME); + + List keys = primaryKeys(cache, 2); + + cache.lock(keys.get(0)).lock(); + + spi(backup).blockMessages(GridDhtLockResponse.class, srv.name()); + + AtomicReference lockErr = new AtomicReference<>(); + + beforeStop = () -> { + try { + cache.lock(keys.get(1)).tryLock(); + } + catch (Throwable e) { + lockErr.set(e); + } + + spi(backup).waitForBlocked(); + spi(backup).stopBlock(); + }; + + stopGrid(0); + + assertTrue("Lock on stopping node must fail", X.hasCause(lockErr.get(), NodeStoppingException.class)); + assertNull("No failures", failureHnd.failureContext()); + } +} diff --git a/modules/core/src/test/java/org/apache/ignite/testsuites/IgniteCacheTestSuite14.java b/modules/core/src/test/java/org/apache/ignite/testsuites/IgniteCacheTestSuite14.java index 2f8fc5a17dc71..21e88dbbe1130 100644 --- a/modules/core/src/test/java/org/apache/ignite/testsuites/IgniteCacheTestSuite14.java +++ b/modules/core/src/test/java/org/apache/ignite/testsuites/IgniteCacheTestSuite14.java @@ -48,6 +48,7 @@ import org.apache.ignite.internal.processors.cache.distributed.IgniteCacheClientNodePartitionsExchangeTest; import org.apache.ignite.internal.processors.cache.distributed.IgniteCacheServerNodeConcurrentStart; import org.apache.ignite.internal.processors.cache.distributed.dht.CachePartitionPartialCountersMapSelfTest; +import org.apache.ignite.internal.processors.cache.distributed.dht.ExplicitLockCancelOnNodeStopTest; import org.apache.ignite.internal.processors.cache.distributed.dht.GridCacheColocatedDebugTest; import org.apache.ignite.internal.processors.cache.distributed.dht.GridCacheColocatedPreloadRestartSelfTest; import org.apache.ignite.internal.processors.cache.distributed.dht.GridCacheColocatedPrimarySyncSelfTest; @@ -125,6 +126,7 @@ public static List> suite(Collection ignoredTests) { GridTestUtils.addTestIfNeeded(suite, GridCachePartitionedMultiNodeSelfTest.class, ignoredTests); GridTestUtils.addTestIfNeeded(suite, GridCachePartitionedExplicitLockNodeFailureSelfTest.class, ignoredTests); GridTestUtils.addTestIfNeeded(suite, CacheLockReleaseNodeLeaveTest.class, ignoredTests); + GridTestUtils.addTestIfNeeded(suite, ExplicitLockCancelOnNodeStopTest.class, ignoredTests); GridTestUtils.addTestIfNeeded(suite, GridCachePartitionedNestedTxTest.class, ignoredTests); GridTestUtils.addTestIfNeeded(suite, GridCachePartitionedTxConcurrentGetTest.class, ignoredTests); GridTestUtils.addTestIfNeeded(suite, GridCachePartitionedTxReadTest.class, ignoredTests);