Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -236,17 +236,21 @@ public Collection<GridCacheMvccCandidate> 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<GridCacheMvccCandidate> deque = cands.get(key);

assert deque != null;
if (deque == null)
return false;

for (GridCacheMvccCandidate cand : deque)
cand.setOwner();

return true;
}
finally {
unlock();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1329,10 +1329,18 @@ private void markLocalDhtLocksAcquired(Collection<KeyCacheObject> 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.
Expand Down Expand Up @@ -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),
Expand Down
Original file line number Diff line number Diff line change
@@ -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:
* <ol>
* <li>A thread holds an explicit lock on one key so its explicit lock span stays non-empty.</li>
* <li>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.</li>
* <li>Since the node is already stopping, the lock future is cancelled immediately which removes
* the explicit lock candidate from the span.</li>
* <li>The backup response arrives afterwards and the completed DHT lock future tries to mark
* the already removed candidate as owned.</li>
* </ol>
*/
@Test
public void testLockOnStoppingNode() throws Exception {
Ignite srv = startGrid(0);
Ignite backup = startGrid(1);

awaitPartitionMapExchange();

IgniteCache<Integer, Integer> cache = srv.cache(DEFAULT_CACHE_NAME);

List<Integer> keys = primaryKeys(cache, 2);

cache.lock(keys.get(0)).lock();

spi(backup).blockMessages(GridDhtLockResponse.class, srv.name());

AtomicReference<Throwable> 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());
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -125,6 +126,7 @@ public static List<Class<?>> suite(Collection<Class> 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);
Expand Down
Loading