Skip to content

Commit 1a893ea

Browse files
João JandreJoaoJandre
authored andcommitted
Fix concurrency issue
1 parent 0f97674 commit 1a893ea

3 files changed

Lines changed: 35 additions & 12 deletions

File tree

plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/BlockCommitListener.java

Lines changed: 19 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -27,24 +27,39 @@
2727
import org.libvirt.event.BlockJobStatus;
2828
import org.libvirt.event.BlockJobType;
2929

30+
import java.util.concurrent.Semaphore;
31+
import java.util.concurrent.TimeUnit;
32+
3033
public class BlockCommitListener implements BlockJobListener {
3134
private String result;
3235
private String vmName;
3336

3437
private Logger logger;
3538
private String logid;
39+
private Semaphore semaphore;
3640

3741
protected BlockCommitListener(String vmName, String logid) {
3842
this.vmName = vmName;
3943
this.logid = logid;
4044
this.logger = LogManager.getLogger(getClass());
45+
this.semaphore = new Semaphore(0);
4146
this.result = String.format("Failed to block commit disk of VM [%s]. Libvirt did not launch an event for it.", vmName);
4247
}
4348

44-
protected String getResult() {
49+
protected String getResult(int timeout) {
50+
this.waitBlockCommit(timeout);
4551
return result;
4652
}
4753

54+
protected void waitBlockCommit(int timeout) {
55+
try {
56+
logger.debug("Trying to acquire result semaphore. If the correct event was not launched, will wait for [{}] seconds before giving up.", timeout);
57+
this.semaphore.tryAcquire(timeout, TimeUnit.SECONDS);
58+
} catch (InterruptedException ex) {
59+
logger.error("Thread that was tracking the progress for the block commit job of vm {} was interrupted.", vmName, ex);
60+
}
61+
}
62+
4863
@Override
4964
public void onEvent(Domain domain, String diskPath, BlockJobType type, BlockJobStatus status) {
5065
if (!BlockJobType.COMMIT.equals(type) && !BlockJobType.ACTIVE_COMMIT.equals(type)) {
@@ -56,17 +71,20 @@ public void onEvent(Domain domain, String diskPath, BlockJobType type, BlockJobS
5671
switch (status) {
5772
case COMPLETED:
5873
result = null;
74+
semaphore.release();
5975
return;
6076
case READY:
6177
try {
6278
logger.debug("Pivoting disk [{}] of VM [{}].", diskPath, vmName);
6379
domain.blockJobAbort(diskPath, Domain.BlockJobAbortFlags.PIVOT);
6480
} catch (LibvirtException ex) {
6581
result = String.format("Failed to pivot disk due to [%s].", ex.getMessage());
82+
semaphore.release();
6683
}
6784
return;
6885
default:
6986
result = String.format("Failed to block commit disk with status [%s].", status);
87+
semaphore.release();
7088
}
7189
}
7290
}

plugins/hypervisors/kvm/src/main/java/com/cloud/hypervisor/kvm/resource/LibvirtComputingResource.java

Lines changed: 10 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -6652,6 +6652,8 @@ protected void mergeSnapshotIntoBaseFileWithEventsAndConfigurableTimeout(Domain
66526652
}
66536653

66546654
BlockCommitListener blockCommitListener = getBlockCommitListener(vmName);
6655+
int remainingTimeout = 0;
6656+
String mergeResult = "Got an exception during block commit wait. Check earlier logs for more info.";
66556657
try {
66566658
vm.addBlockJobListener(blockCommitListener);
66576659

@@ -6660,12 +6662,12 @@ protected void mergeSnapshotIntoBaseFileWithEventsAndConfigurableTimeout(Domain
66606662

66616663
vm.blockCommit(diskLabel, baseFilePath, topFilePath, 0, commitFlags);
66626664

6663-
checkBlockCommitProgress(vm, diskLabel, vmName, snapshotName, topFilePath, baseFilePath);
6665+
remainingTimeout = checkBlockCommitProgress(vm, diskLabel, vmName, snapshotName, topFilePath, baseFilePath);
6666+
mergeResult = blockCommitListener.getResult(remainingTimeout);
66646667
} finally {
66656668
vm.removeBlockJobListener(blockCommitListener);
66666669
}
66676670

6668-
String mergeResult = blockCommitListener.getResult();
66696671
if (mergeResult != null) {
66706672
String commitError = String.format("Failed the block commit of top file [%s] into base file [%s] for snapshot [%s] of VM [%s]. The job will be left running to avoid" +
66716673
" data corruption, but ACS will return an error and volume [%s] will need to be normalized manually. If the commit involved the active image, the pivot will" +
@@ -6730,7 +6732,7 @@ protected BlockCommitListener getBlockCommitListener(String vmName) {
67306732
return new BlockCommitListener(vmName, ThreadContext.get("logcontextid"));
67316733
}
67326734

6733-
protected void checkBlockCommitProgress(Domain vm, String diskLabel, String vmName, String snapshotName, String topFilePath, String baseFilePath) {
6735+
protected int checkBlockCommitProgress(Domain vm, String diskLabel, String vmName, String snapshotName, String topFilePath, String baseFilePath) {
67346736
int timeout = qcow2DeltaMergeTimeout;
67356737
DomainBlockJobInfo result;
67366738
long lastCommittedBytes = 0;
@@ -6751,12 +6753,12 @@ protected void checkBlockCommitProgress(Domain vm, String diskLabel, String vmNa
67516753
result = vm.getBlockJobInfo(diskLabel, 0);
67526754
} catch (LibvirtException ex) {
67536755
logger.warn("Exception while getting block job info {}: [{}].", partialLog, ex.getMessage(), ex);
6754-
return;
6756+
return timeout;
67556757
}
67566758

67576759
if (result == null || result.type == 0 && result.end == 0 && result.cur == 0) {
67586760
logger.debug("Block commit job {} has already finished.", partialLog);
6759-
return;
6761+
return timeout;
67606762
}
67616763

67626764
long currentCommittedBytes = result.cur;
@@ -6766,7 +6768,9 @@ protected void checkBlockCommitProgress(Domain vm, String diskLabel, String vmNa
67666768
lastCommittedBytes = currentCommittedBytes;
67676769
endBytes = result.end;
67686770
}
6769-
logger.warn("Block commit {} has timed out after waiting at least {} seconds. The progress of the operation was [{}] of [{}].", partialLog, qcow2DeltaMergeTimeout, lastCommittedBytes, endBytes);
6771+
logger.warn("Block commit {} has timed out after waiting at least {} seconds. The progress of the operation was [{}] of [{}].", partialLog,
6772+
qcow2DeltaMergeTimeout, lastCommittedBytes, endBytes);
6773+
return 0;
67706774
}
67716775

67726776
/**

plugins/hypervisors/kvm/src/test/java/com/cloud/hypervisor/kvm/resource/LibvirtComputingResourceTest.java

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@
2727
import static org.junit.Assert.assertTrue;
2828
import static org.junit.Assert.fail;
2929
import static org.mockito.ArgumentMatchers.any;
30+
import static org.mockito.ArgumentMatchers.anyInt;
3031
import static org.mockito.ArgumentMatchers.nullable;
3132
import static org.mockito.Mockito.doNothing;
3233
import static org.mockito.Mockito.doReturn;
@@ -6697,7 +6698,7 @@ public void mergeSnapshotIntoBaseFileTestActiveAndDeleteFlags() throws Exception
66976698
threadContextMockedStatic.when(() ->
66986699
ThreadContext.get(Mockito.anyString())).thenReturn("logid");
66996700
Mockito.doReturn(blockCommitListenerMock).when(libvirtComputingResourceSpy).getBlockCommitListener(Mockito.any());
6700-
Mockito.doReturn(null).when(blockCommitListenerMock).getResult();
6701+
Mockito.doReturn(null).when(blockCommitListenerMock).getResult(anyInt());
67016702
Mockito.doNothing().when(domainMock).addBlockJobListener(blockCommitListenerMock);
67026703
Mockito.doReturn(null).when(domainMock).getBlockJobInfo(Mockito.anyString(), Mockito.anyInt());
67036704
Mockito.doNothing().when(domainMock).removeBlockJobListener(blockCommitListenerMock);
@@ -6726,7 +6727,7 @@ public void mergeSnapshotIntoBaseFileTestActiveFlag() throws Exception {
67266727
threadContextMockedStatic.when(() ->
67276728
ThreadContext.get(Mockito.anyString())).thenReturn("logid");
67286729
Mockito.doReturn(blockCommitListenerMock).when(libvirtComputingResourceSpy).getBlockCommitListener(Mockito.any());
6729-
Mockito.doReturn(null).when(blockCommitListenerMock).getResult();
6730+
Mockito.doReturn(null).when(blockCommitListenerMock).getResult(anyInt());
67306731
Mockito.doNothing().when(domainMock).addBlockJobListener(blockCommitListenerMock);
67316732
Mockito.doNothing().when(domainMock).removeBlockJobListener(blockCommitListenerMock);
67326733
Mockito.doNothing().when(libvirtComputingResourceSpy).manuallyDeleteUnusedSnapshotFile(Mockito.anyBoolean(), Mockito.anyString());
@@ -6753,7 +6754,7 @@ public void mergeSnapshotIntoBaseFileTestDeleteFlag() throws Exception {
67536754
libvirtUtilitiesHelperMockedStatic.when(() -> LibvirtUtilitiesHelper.isLibvirtSupportingFlagDeleteOnCommandVirshBlockcommit(Mockito.any())).thenReturn(true);
67546755
threadContextMockedStatic.when(() -> ThreadContext.get(Mockito.anyString())).thenReturn("logid");
67556756
Mockito.doReturn(blockCommitListenerMock).when(libvirtComputingResourceSpy).getBlockCommitListener(Mockito.any());
6756-
Mockito.doReturn(null).when(blockCommitListenerMock).getResult();
6757+
Mockito.doReturn(null).when(blockCommitListenerMock).getResult(anyInt());
67576758
Mockito.doNothing().when(domainMock).addBlockJobListener(blockCommitListenerMock);
67586759
Mockito.doReturn(null).when(domainMock).getBlockJobInfo(Mockito.anyString(), Mockito.anyInt());
67596760
Mockito.doNothing().when(domainMock).removeBlockJobListener(blockCommitListenerMock);
@@ -6781,7 +6782,7 @@ public void mergeSnapshotIntoBaseFileTestNoFlags() throws Exception {
67816782
libvirtUtilitiesHelperMockedStatic.when(() -> LibvirtUtilitiesHelper.isLibvirtSupportingFlagDeleteOnCommandVirshBlockcommit(Mockito.any())).thenReturn(false);
67826783
threadContextMockedStatic.when(() -> ThreadContext.get(Mockito.anyString())).thenReturn("logid");
67836784
Mockito.doReturn(blockCommitListenerMock).when(libvirtComputingResourceSpy).getBlockCommitListener(Mockito.any());
6784-
Mockito.doReturn(null).when(blockCommitListenerMock).getResult();
6785+
Mockito.doReturn(null).when(blockCommitListenerMock).getResult(anyInt());
67856786
Mockito.doNothing().when(domainMock).addBlockJobListener(blockCommitListenerMock);
67866787
Mockito.doReturn(null).when(domainMock).getBlockJobInfo(Mockito.anyString(), Mockito.anyInt());
67876788
Mockito.doNothing().when(domainMock).removeBlockJobListener(blockCommitListenerMock);
@@ -6811,7 +6812,7 @@ public void mergeSnapshotIntoBaseFileTestMergeFailsThrowException() throws Excep
68116812
Mockito.doNothing().when(domainMock).removeBlockJobListener(Mockito.any());
68126813

68136814
Mockito.doReturn(blockCommitListenerMock).when(libvirtComputingResourceSpy).getBlockCommitListener(Mockito.any());
6814-
Mockito.doReturn("Failed").when(blockCommitListenerMock).getResult();
6815+
Mockito.doReturn("Failed").when(blockCommitListenerMock).getResult(anyInt());
68156816

68166817
String diskLabel = "vda";
68176818
String baseFilePath = "/file";

0 commit comments

Comments
 (0)