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
@@ -0,0 +1,45 @@
// 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 com.cloud.hypervisor.kvm.storage;

import java.util.concurrent.ConcurrentHashMap;

/**
* One monitor per storage pool uuid, shared by {@link KVMStoragePoolManager} and
* {@link LibvirtStorageAdaptor} so that both take the same lock for a pool.
*
* The adaptor holds it across the whole of createStoragePool and deleteStoragePool. The refcount
* alone cannot keep a pool alive: deleteStoragePool decides to tear the pool down when the count
* reaches zero and then destroys and unmounts it, and a createStoragePool that finds the still
* active pool and takes a reference in between would have it torn down underneath it.
*
* Whoever also takes the manager wide lock must take this one first. A teardown can sit in an
* umount for as long as the storage takes to answer, and a create of that pool that waited for it
* while holding the manager wide lock would hold up the creates of every other pool on the host.
*
* Entries are never removed, as there is one per pool the host has ever used.
*/
final class KVMStoragePoolLocks {
private static final ConcurrentHashMap<String, Object> LOCKS = new ConcurrentHashMap<>();

private KVMStoragePoolLocks() {
}

static Object get(String uuid) {
return LOCKS.computeIfAbsent(uuid, k -> new Object());
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -387,8 +387,19 @@ public KVMStoragePool createStoragePool(String name, String host, int port, Stri
return createStoragePool(name, host, port, path, userInfo, type, details, true);
}

private KVMStoragePool createStoragePool(String name, String host, int port, String path, String userInfo, StoragePoolType type, Map<String, String> details, boolean primaryStorage) {
/*
* Take the pool's own lock before the manager wide one. A create of a pool that is being torn
* down waits here for the teardown, which can sit in an umount, and must not hold up the
* creates of every other pool on the host while it does. See KVMStoragePoolLocks.
*/
synchronized (KVMStoragePoolLocks.get(name)) {
return createStoragePoolSynchronized(name, host, port, path, userInfo, type, details, primaryStorage);
}
}

//Note: due to bug CLOUDSTACK-4459, createStoragepool can be called in parallel, so need to be synced.
private synchronized KVMStoragePool createStoragePool(String name, String host, int port, String path, String userInfo, StoragePoolType type, Map<String, String> details, boolean primaryStorage) {
protected synchronized KVMStoragePool createStoragePoolSynchronized(String name, String host, int port, String path, String userInfo, StoragePoolType type, Map<String, String> details, boolean primaryStorage) {
StorageAdaptor adaptor = getStorageAdaptor(type);
KVMStoragePool pool = adaptor.createStoragePool(name, host, port, path, userInfo, type, details, primaryStorage);
if (pool instanceof LibvirtStoragePool) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -700,40 +700,41 @@ public KVMPhysicalDisk getPhysicalDisk(String volumeUuid, KVMStoragePool pool) {
* adjust refcount
*/
private int adjustStoragePoolRefCount(String uuid, int adjustment) {
final String mutexKey = storagePoolRefCounts.keySet().stream()
.filter(k -> k.equals(uuid))
.findFirst()
.orElse(uuid);
synchronized (mutexKey) {
// some access on the storagePoolRefCounts.key(mutexKey) element
int refCount = storagePoolRefCounts.computeIfAbsent(mutexKey, k -> 0);
refCount += adjustment;
if (refCount < 1) {
storagePoolRefCounts.remove(mutexKey);
} else {
storagePoolRefCounts.put(mutexKey, refCount);
}
return refCount;
}
/*
* compute() is atomic for the key, so concurrent callers cannot lose an
* update. Returning null from the remapping function removes the entry,
* which keeps the map free of pools that are no longer in use.
*/
Integer refCount = storagePoolRefCounts.compute(uuid, (key, count) -> {
Comment thread
bhouse-nexthop marked this conversation as resolved.
int adjusted = (count == null ? 0 : count) + adjustment;
return adjusted < 1 ? null : adjusted;
});
return refCount == null ? 0 : refCount;
}
/**
* Thread-safe increment storage pool usage refcount
* @param uuid UUID of the storage pool to increment the count
*/
private void incStoragePoolRefCount(String uuid) {
protected void incStoragePoolRefCount(String uuid) {
adjustStoragePoolRefCount(uuid, 1);
}
/**
* Thread-safe decrement storage pool usage refcount for the given uuid and return if storage pool still in use.
* @param uuid UUID of the storage pool to decrement the count
* @return true if the storage pool is still used, else false.
*/
private boolean decStoragePoolRefCount(String uuid) {
protected boolean decStoragePoolRefCount(String uuid) {
return adjustStoragePoolRefCount(uuid, -1) > 0;
}

@Override
public KVMStoragePool createStoragePool(String name, String host, int port, String path, String userInfo, StoragePoolType type, Map<String, String> details, boolean isPrimaryStorage) {
synchronized (KVMStoragePoolLocks.get(name)) {
return createStoragePoolLocked(name, host, port, path, userInfo, type, details, isPrimaryStorage);
}
}

protected KVMStoragePool createStoragePoolLocked(String name, String host, int port, String path, String userInfo, StoragePoolType type, Map<String, String> details, boolean isPrimaryStorage) {
logger.info("Attempting to create storage pool {} ({}) in libvirt", name, type);
StoragePool sp;
Connect conn;
Expand Down Expand Up @@ -905,6 +906,12 @@ private boolean destroyStoragePoolHandleException(Connect conn, String uuid)

@Override
public boolean deleteStoragePool(String uuid) {
synchronized (KVMStoragePoolLocks.get(uuid)) {
return deleteStoragePoolLocked(uuid);
}
}

protected boolean deleteStoragePoolLocked(String uuid) {
logger.info("Attempting to remove storage pool " + uuid + " from libvirt");

// decrement and check if storage pool still in use
Expand Down Expand Up @@ -948,13 +955,29 @@ public boolean deleteStoragePool(String uuid) {
String targetPath = _mountPoint + File.separator + uuid;
logger.error("deleteStoragePool removed pool from libvirt, but libvirt had trouble unmounting the pool. Trying umount location " + targetPath +
" again in a few seconds");
String result = Script.runSimpleBashScript("sleep 5 && umount " + targetPath);
if (result == null) {
/*
* runSimpleBashScript() returns null both when the command fails,
* because runScript() discards the output on a non-zero exit, and
* when it succeeds without printing anything. Its result therefore
* cannot say whether the umount worked. It is still used to run the
* umount, because it logs the failure reason, which is the useful
* diagnostic, but the outcome is taken from whether the path is
* still a mount point. That also covers the pool having been
* unmounted by something else in the meantime.
*/
Script.runSimpleBashScript("sleep 5 && umount " + targetPath);
if (Script.runSimpleBashScriptForExitValue("mountpoint -q " + targetPath) != 0) {
logger.info("Succeeded in unmounting " + targetPath);
destroyStoragePoolHandleException(conn, uuid);
return true;
}
logger.error("Failed to unmount " + targetPath);
/*
* Do not throw here. deleteStoragePool() is called from finally
* blocks, where a throw would discard the result of an operation
* that has already succeeded.
*/
logger.error("Failed to unmount " + targetPath + ", it is still a mount point");
return false;
}
throw new CloudRuntimeException(e.toString(), e);
}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
// 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 com.cloud.hypervisor.kvm.storage;

import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyBoolean;
import static org.mockito.ArgumentMatchers.anyInt;
import static org.mockito.ArgumentMatchers.eq;

import java.util.UUID;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;

import org.junit.Assert;
import org.junit.Test;
import org.mockito.Mockito;

import com.cloud.storage.Storage.StoragePoolType;

public class KVMStoragePoolManagerTest {

@Test(timeout = 60000)
public void testCreateWaitingForPoolLockDoesNotHoldManagerLock() throws Exception {
/*
* A teardown of a pool holds that pool's lock, and can sit in an umount for as long as the
* storage takes to answer. A create of the same pool has to wait for it, but must do so
* without holding the manager wide lock, or the creates of every other pool on the host wait
* behind the umount too.
*/
final KVMStoragePoolManager manager = Mockito.mock(KVMStoragePoolManager.class, Mockito.CALLS_REAL_METHODS);
final KVMStoragePool pool = Mockito.mock(KVMStoragePool.class);
final String uuid = String.valueOf(UUID.randomUUID());
Mockito.doReturn(pool).when(manager).createStoragePoolSynchronized(eq(uuid), any(), anyInt(), any(), any(), any(), any(), anyBoolean());

final AtomicReference<KVMStoragePool> created = new AtomicReference<>();
final Thread create;
synchronized (KVMStoragePoolLocks.get(uuid)) {
// the teardown in progress
create = new Thread(() -> created.set(manager.createStoragePool(uuid, "127.0.0.1", 0, "/export/primary",
null, StoragePoolType.NetworkFilesystem)));
create.start();
while (create.getState() != Thread.State.BLOCKED) {
Assert.assertTrue("create finished while the pool's lock was held", create.isAlive());
Thread.sleep(10);
}

final CountDownLatch managerLockTaken = new CountDownLatch(1);
final Thread otherPool = new Thread(() -> {
synchronized (manager) {
managerLockTaken.countDown();
}
});
otherPool.start();
Assert.assertTrue("the manager wide lock was held by a create waiting for its pool's lock",
managerLockTaken.await(30, TimeUnit.SECONDS));
Assert.assertNull(created.get());
}

create.join(30000);
Assert.assertSame(pool, created.get());
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -17,16 +17,31 @@

package com.cloud.hypervisor.kvm.storage;

import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyBoolean;
import static org.mockito.ArgumentMatchers.anyInt;
import static org.mockito.ArgumentMatchers.anyList;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.Mockito.never;

import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.BrokenBarrierException;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.CyclicBarrier;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicBoolean;

import org.junit.After;
import org.junit.Assert;
import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
Expand Down Expand Up @@ -176,4 +191,116 @@ public void testUpdateLocalPoolIops_NullResultFromScript() {

Mockito.verify(mockPool, never()).setUsedIops(anyLong());
}

@Test(timeout = 120000)
public void testStoragePoolRefCountCountsEveryConcurrentIncrement() throws Exception {
final LibvirtStorageAdaptor adaptor = new LibvirtStorageAdaptor(null);
final int threads = 16;
final int rounds = 500;
final CyclicBarrier barrier = new CyclicBarrier(threads);
final ExecutorService executor = Executors.newFixedThreadPool(threads);

try {
for (int round = 0; round < rounds; round++) {
// A fresh uuid each round, so every round starts with no entry for the pool.
final String uuid = String.valueOf(UUID.randomUUID());
final List<Future<?>> futures = new ArrayList<>();

for (int i = 0; i < threads; i++) {
futures.add(executor.submit(() -> {
/*
* Every caller arrives with its own String instance, the way the
* agent does when the uuid is parsed out of a separate command
* payload for each request. The instances are equal but they are
* not the same object.
*/
final String ownInstance = new String(uuid);
try {
barrier.await();
} catch (InterruptedException | BrokenBarrierException e) {
Thread.currentThread().interrupt();
throw new IllegalStateException(e);
}
adaptor.incStoragePoolRefCount(ownInstance);
}));
}
for (Future<?> future : futures) {
future.get(60, TimeUnit.SECONDS);
}

// Every increment must be counted, so the pool stays in use until the last release.
for (int i = 1; i < threads; i++) {
Assert.assertTrue("Round " + round + ": pool should still be in use after " + i
+ " of " + threads + " releases", adaptor.decStoragePoolRefCount(uuid));
}
Assert.assertFalse("Round " + round + ": pool should no longer be in use after the last release",
adaptor.decStoragePoolRefCount(uuid));
}
} finally {
executor.shutdownNow();
}
}

/**
* Holds a delete of {@code deleteUuid} part way through its teardown, then starts a create of
* {@code createUuid} and reports whether that create completed while the delete was still held.
*/
private boolean createRanDuringDelete(String deleteUuid, String createUuid) throws Exception {
final LibvirtStorageAdaptor adaptor = Mockito.spy(new LibvirtStorageAdaptor(null));
final CountDownLatch deleteEntered = new CountDownLatch(1);
final CountDownLatch releaseDelete = new CountDownLatch(1);
final AtomicBoolean deleteFinished = new AtomicBoolean();
final AtomicBoolean createRanBeforeDeleteFinished = new AtomicBoolean();

Mockito.doAnswer(invocation -> {
deleteEntered.countDown();
releaseDelete.await();
deleteFinished.set(true);
return true;
}).when(adaptor).deleteStoragePoolLocked(deleteUuid);
Mockito.doAnswer(invocation -> {
createRanBeforeDeleteFinished.set(!deleteFinished.get());
return mockPool;
}).when(adaptor).createStoragePoolLocked(Mockito.eq(createUuid), any(), anyInt(), any(), any(), any(), any(), anyBoolean());

final ExecutorService executor = Executors.newFixedThreadPool(2);
try {
final Future<Boolean> delete = executor.submit(() -> adaptor.deleteStoragePool(deleteUuid));
Assert.assertTrue("delete never started", deleteEntered.await(30, TimeUnit.SECONDS));

final Future<KVMStoragePool> create = executor.submit(() -> adaptor.createStoragePool(createUuid, "127.0.0.1", 0,
"/export/secondary", null, Storage.StoragePoolType.NetworkFilesystem, null, false));
try {
create.get(2, TimeUnit.SECONDS);
} catch (TimeoutException e) {
// still waiting on the delete, which is what a create of the same pool must do
}

releaseDelete.countDown();
Assert.assertTrue(delete.get(30, TimeUnit.SECONDS));
Assert.assertSame(mockPool, create.get(30, TimeUnit.SECONDS));
return createRanBeforeDeleteFinished.get();
} finally {
releaseDelete.countDown();
executor.shutdownNow();
}
}

@Test(timeout = 120000)
public void testCreateStoragePoolWaitsForDeleteOfSamePool() throws Exception {
/*
* A delete that has dropped the last reference goes on to destroy and unmount the pool. A
* create that found the still active pool and took a reference in between would have it torn
* down underneath it, so the create has to wait for the teardown to finish.
*/
final String uuid = String.valueOf(UUID.randomUUID());
Assert.assertFalse("create of a pool ran while that pool was being deleted",
createRanDuringDelete(uuid, new String(uuid)));
}

@Test(timeout = 120000)
public void testCreateStoragePoolDoesNotWaitForDeleteOfAnotherPool() throws Exception {
Assert.assertTrue("create of one pool waited for the delete of another",
createRanDuringDelete(String.valueOf(UUID.randomUUID()), String.valueOf(UUID.randomUUID())));
}
}
Loading