diff --git a/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/storage/KVMStoragePoolLocks.java b/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/storage/KVMStoragePoolLocks.java new file mode 100644 index 000000000000..dd204e0212f7 --- /dev/null +++ b/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/storage/KVMStoragePoolLocks.java @@ -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 LOCKS = new ConcurrentHashMap<>(); + + private KVMStoragePoolLocks() { + } + + static Object get(String uuid) { + return LOCKS.computeIfAbsent(uuid, k -> new Object()); + } +} diff --git a/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/storage/KVMStoragePoolManager.java b/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/storage/KVMStoragePoolManager.java index 35cc864268c3..a660ca5baee7 100644 --- a/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/storage/KVMStoragePoolManager.java +++ b/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/storage/KVMStoragePoolManager.java @@ -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 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 details, boolean primaryStorage) { + protected synchronized KVMStoragePool createStoragePoolSynchronized(String name, String host, int port, String path, String userInfo, StoragePoolType type, Map details, boolean primaryStorage) { StorageAdaptor adaptor = getStorageAdaptor(type); KVMStoragePool pool = adaptor.createStoragePool(name, host, port, path, userInfo, type, details, primaryStorage); if (pool instanceof LibvirtStoragePool) { diff --git a/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/storage/LibvirtStorageAdaptor.java b/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/storage/LibvirtStorageAdaptor.java index 059f4f8b67af..83545eef9b72 100644 --- a/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/storage/LibvirtStorageAdaptor.java +++ b/plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/storage/LibvirtStorageAdaptor.java @@ -700,27 +700,22 @@ 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) -> { + 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); } /** @@ -728,12 +723,18 @@ private void incStoragePoolRefCount(String uuid) { * @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 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 details, boolean isPrimaryStorage) { logger.info("Attempting to create storage pool {} ({}) in libvirt", name, type); StoragePool sp; Connect conn; @@ -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 @@ -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); } diff --git a/plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/storage/KVMStoragePoolManagerTest.java b/plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/storage/KVMStoragePoolManagerTest.java new file mode 100644 index 000000000000..bbb1b4d722af --- /dev/null +++ b/plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/storage/KVMStoragePoolManagerTest.java @@ -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 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()); + } +} diff --git a/plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/storage/LibvirtStorageAdaptorTest.java b/plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/storage/LibvirtStorageAdaptorTest.java index 88346abd0176..abc49eb402b8 100644 --- a/plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/storage/LibvirtStorageAdaptorTest.java +++ b/plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/storage/LibvirtStorageAdaptorTest.java @@ -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; @@ -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> 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 delete = executor.submit(() -> adaptor.deleteStoragePool(deleteUuid)); + Assert.assertTrue("delete never started", deleteEntered.await(30, TimeUnit.SECONDS)); + + final Future 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()))); + } }