Skip to content
Open
Show file tree
Hide file tree
Changes from 2 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 @@ -86,6 +86,14 @@ public class LibvirtStorageAdaptor implements StorageAdaptor {
private String _mountPoint = "/mnt";
private String _manageSnapshotPath;
private static final ConcurrentHashMap<String, Integer> storagePoolRefCounts = new ConcurrentHashMap<>();
/*
* One monitor per pool, held 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.
* Entries are never removed, as there is one per pool the host has ever used.
*/
private static final ConcurrentHashMap<String, Object> storagePoolLocks = new ConcurrentHashMap<>();

private String rbdTemplateSnapName = "cloudstack-base-snap";
private static final int RBD_FEATURE_LAYERING = 1;
Expand Down Expand Up @@ -700,40 +708,45 @@ 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;
}

private static Object getStoragePoolLock(String uuid) {
return storagePoolLocks.computeIfAbsent(uuid, k -> new Object());
}

@Override
public KVMStoragePool createStoragePool(String name, String host, int port, String path, String userInfo, StoragePoolType type, Map<String, String> details, boolean isPrimaryStorage) {
synchronized (getStoragePoolLock(name)) {
Comment thread
bhouse-nexthop marked this conversation as resolved.
Outdated
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 +918,12 @@ private boolean destroyStoragePoolHandleException(Connect conn, String uuid)

@Override
public boolean deleteStoragePool(String uuid) {
synchronized (getStoragePoolLock(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 +967,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
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