diff --git a/engine/schema/src/main/java/com/cloud/usage/dao/UsageBackupDao.java b/engine/schema/src/main/java/com/cloud/usage/dao/UsageBackupDao.java index 2a24016653db..5afdf737edae 100644 --- a/engine/schema/src/main/java/com/cloud/usage/dao/UsageBackupDao.java +++ b/engine/schema/src/main/java/com/cloud/usage/dao/UsageBackupDao.java @@ -24,7 +24,8 @@ import com.cloud.utils.db.GenericDao; public interface UsageBackupDao extends GenericDao { - void updateMetrics(Long vmId, Long backupOfferingId, Long size, Long virtualSize); + List listActiveUsage(Long vmId, Long backupOfferingId); + void updateMetrics(Long vmId, Long backupOfferingId, Long size, Long virtualSize, Date eventDate); void removeUsage(Long accountId, Long vmId, Long backupOfferingId, Date eventDate); List getUsageRecords(Long accountId, Date startDate, Date endDate); } diff --git a/engine/schema/src/main/java/com/cloud/usage/dao/UsageBackupDaoImpl.java b/engine/schema/src/main/java/com/cloud/usage/dao/UsageBackupDaoImpl.java index 3f852b0cfb5a..81d268a3cecf 100644 --- a/engine/schema/src/main/java/com/cloud/usage/dao/UsageBackupDaoImpl.java +++ b/engine/schema/src/main/java/com/cloud/usage/dao/UsageBackupDaoImpl.java @@ -30,6 +30,7 @@ import com.cloud.exception.CloudException; import com.cloud.usage.UsageBackupVO; import com.cloud.utils.DateUtil; +import com.cloud.utils.db.Filter; import com.cloud.utils.db.GenericDaoBase; import com.cloud.utils.db.SearchCriteria; import com.cloud.utils.db.TransactionLegacy; @@ -42,19 +43,50 @@ public class UsageBackupDaoImpl extends GenericDaoBase impl " OR ((created <= ?) AND (removed >= ?)))"; @Override - public void updateMetrics(final Long vmId, Long backupOfferingId, final Long size, final Long virtualSize) { - try (TransactionLegacy txn = TransactionLegacy.open(TransactionLegacy.USAGE_DB)) { - SearchCriteria sc = this.createSearchCriteria(); - sc.addAnd("vmId", SearchCriteria.Op.EQ, vmId); - sc.addAnd("backupOfferingId", SearchCriteria.Op.EQ, backupOfferingId); - UsageBackupVO vo = findOneBy(sc); - if (vo != null) { - vo.setSize(size); - vo.setProtectedSize(virtualSize); - update(vo.getId(), vo); + public List listActiveUsage(Long vmId, Long backupOfferingId) { + SearchCriteria sc = this.createSearchCriteria(); + sc.addAnd("vmId", SearchCriteria.Op.EQ, vmId); + sc.addAnd("backupOfferingId", SearchCriteria.Op.EQ, backupOfferingId); + sc.addAnd("removed", SearchCriteria.Op.NULL); + return listBy(sc, new Filter(UsageBackupVO.class, "created", false)); + } + + @Override + public void updateMetrics(final Long vmId, final Long backupOfferingId, final Long size, final Long virtualSize, final Date eventDate) { + final long newSize = size != null ? size : 0L; + final long newProtectedSize = virtualSize != null ? virtualSize : 0L; + TransactionLegacy txn = TransactionLegacy.open(TransactionLegacy.USAGE_DB); + try { + txn.start(); + List activeUsage = listActiveUsage(vmId, backupOfferingId); + if (activeUsage.isEmpty()) { + logger.warn("No active backup usage for VM [{}] and backup offering [{}], ignoring backup metrics of size [{}] and protected size [{}].", + vmId, backupOfferingId, newSize, newProtectedSize); + txn.commit(); + return; + } + + UsageBackupVO latest = activeUsage.get(0); + if (activeUsage.size() == 1 && latest.getSize() == newSize && latest.getProtectedSize() == newProtectedSize) { + txn.commit(); + return; + } + + // Close the active rows and open one with the new size; this also merges duplicates. + for (UsageBackupVO usage : activeUsage) { + usage.setRemoved(eventDate); + update(usage.getId(), usage); } + UsageBackupVO newUsage = new UsageBackupVO(latest.getZoneId(), latest.getAccountId(), latest.getDomainId(), vmId, backupOfferingId, eventDate); + newUsage.setSize(newSize); + newUsage.setProtectedSize(newProtectedSize); + persist(newUsage); + txn.commit(); } catch (final Exception e) { + txn.rollback(); logger.error("Error updating backup metrics: " + e.getMessage(), e); + } finally { + txn.close(); } } diff --git a/engine/schema/src/main/java/org/apache/cloudstack/backup/BackupUsageMetricVO.java b/engine/schema/src/main/java/org/apache/cloudstack/backup/BackupUsageMetricVO.java new file mode 100644 index 000000000000..694e041d828a --- /dev/null +++ b/engine/schema/src/main/java/org/apache/cloudstack/backup/BackupUsageMetricVO.java @@ -0,0 +1,107 @@ +// 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.cloudstack.backup; + +import java.util.Date; + +import javax.persistence.Column; +import javax.persistence.Entity; +import javax.persistence.GeneratedValue; +import javax.persistence.GenerationType; +import javax.persistence.Id; +import javax.persistence.Table; +import javax.persistence.Temporal; +import javax.persistence.TemporalType; + +import org.apache.cloudstack.api.InternalIdentity; + +/** + * The backup usage metric last published for a VM and backup offering. + */ +@Entity +@Table(name = "backup_usage_metric") +public class BackupUsageMetricVO implements InternalIdentity { + + @Id + @GeneratedValue(strategy = GenerationType.IDENTITY) + @Column(name = "id") + private long id; + + @Column(name = "vm_id") + private long vmId; + + @Column(name = "backup_offering_id") + private long backupOfferingId; + + @Column(name = "size") + private long size; + + @Column(name = "protected_size") + private long protectedSize; + + @Column(name = "updated") + @Temporal(value = TemporalType.TIMESTAMP) + private Date updated; + + protected BackupUsageMetricVO() { + } + + public BackupUsageMetricVO(long vmId, long backupOfferingId, long size, long protectedSize, Date updated) { + this.vmId = vmId; + this.backupOfferingId = backupOfferingId; + this.size = size; + this.protectedSize = protectedSize; + this.updated = updated; + } + + @Override + public long getId() { + return id; + } + + public long getVmId() { + return vmId; + } + + public long getBackupOfferingId() { + return backupOfferingId; + } + + public long getSize() { + return size; + } + + public void setSize(long size) { + this.size = size; + } + + public long getProtectedSize() { + return protectedSize; + } + + public void setProtectedSize(long protectedSize) { + this.protectedSize = protectedSize; + } + + public Date getUpdated() { + return updated; + } + + public void setUpdated(Date updated) { + this.updated = updated; + } +} diff --git a/engine/schema/src/main/java/org/apache/cloudstack/backup/dao/BackupUsageMetricDao.java b/engine/schema/src/main/java/org/apache/cloudstack/backup/dao/BackupUsageMetricDao.java new file mode 100644 index 000000000000..c068f4efca5a --- /dev/null +++ b/engine/schema/src/main/java/org/apache/cloudstack/backup/dao/BackupUsageMetricDao.java @@ -0,0 +1,27 @@ +// 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.cloudstack.backup.dao; + +import org.apache.cloudstack.backup.BackupUsageMetricVO; + +import com.cloud.utils.db.GenericDao; + +public interface BackupUsageMetricDao extends GenericDao { + BackupUsageMetricVO findByVmAndOffering(long vmId, long backupOfferingId); + + int removeByVmAndOffering(long vmId, long backupOfferingId); +} diff --git a/engine/schema/src/main/java/org/apache/cloudstack/backup/dao/BackupUsageMetricDaoImpl.java b/engine/schema/src/main/java/org/apache/cloudstack/backup/dao/BackupUsageMetricDaoImpl.java new file mode 100644 index 000000000000..9f40d80a4fbb --- /dev/null +++ b/engine/schema/src/main/java/org/apache/cloudstack/backup/dao/BackupUsageMetricDaoImpl.java @@ -0,0 +1,54 @@ +// 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.cloudstack.backup.dao; + +import javax.annotation.PostConstruct; + +import org.apache.cloudstack.backup.BackupUsageMetricVO; + +import com.cloud.utils.db.GenericDaoBase; +import com.cloud.utils.db.SearchBuilder; +import com.cloud.utils.db.SearchCriteria; + +public class BackupUsageMetricDaoImpl extends GenericDaoBase implements BackupUsageMetricDao { + private SearchBuilder vmAndOfferingSearch; + + @PostConstruct + protected void init() { + vmAndOfferingSearch = createSearchBuilder(); + vmAndOfferingSearch.and("vmId", vmAndOfferingSearch.entity().getVmId(), SearchCriteria.Op.EQ); + vmAndOfferingSearch.and("backupOfferingId", vmAndOfferingSearch.entity().getBackupOfferingId(), SearchCriteria.Op.EQ); + vmAndOfferingSearch.done(); + } + + private SearchCriteria createVmAndOfferingCriteria(long vmId, long backupOfferingId) { + SearchCriteria sc = vmAndOfferingSearch.create(); + sc.setParameters("vmId", vmId); + sc.setParameters("backupOfferingId", backupOfferingId); + return sc; + } + + @Override + public BackupUsageMetricVO findByVmAndOffering(long vmId, long backupOfferingId) { + return findOneBy(createVmAndOfferingCriteria(vmId, backupOfferingId)); + } + + @Override + public int removeByVmAndOffering(long vmId, long backupOfferingId) { + return remove(createVmAndOfferingCriteria(vmId, backupOfferingId)); + } +} diff --git a/engine/schema/src/main/resources/META-INF/cloudstack/core/spring-engine-schema-core-daos-context.xml b/engine/schema/src/main/resources/META-INF/cloudstack/core/spring-engine-schema-core-daos-context.xml index 0656d5e3c440..fe9cb76da02c 100644 --- a/engine/schema/src/main/resources/META-INF/cloudstack/core/spring-engine-schema-core-daos-context.xml +++ b/engine/schema/src/main/resources/META-INF/cloudstack/core/spring-engine-schema-core-daos-context.xml @@ -273,6 +273,7 @@ + diff --git a/engine/schema/src/main/resources/META-INF/db/schema-42210to42220.sql b/engine/schema/src/main/resources/META-INF/db/schema-42210to42220.sql index f6d73def74fd..4ed2008867e7 100644 --- a/engine/schema/src/main/resources/META-INF/db/schema-42210to42220.sql +++ b/engine/schema/src/main/resources/META-INF/db/schema-42210to42220.sql @@ -18,3 +18,15 @@ --; -- Schema upgrade from 4.22.1.0 to 4.22.2.0 --; + +-- Last backup usage metric published per VM and backup offering +CREATE TABLE IF NOT EXISTS `cloud`.`backup_usage_metric` ( + `id` bigint unsigned NOT NULL auto_increment COMMENT 'id', + `vm_id` bigint unsigned NOT NULL COMMENT 'VM ID', + `backup_offering_id` bigint unsigned NOT NULL COMMENT 'Backup offering ID', + `size` bigint unsigned NOT NULL COMMENT 'Backup size last published', + `protected_size` bigint unsigned NOT NULL COMMENT 'Protected size last published', + `updated` datetime NOT NULL COMMENT 'Date the metric was last published', + PRIMARY KEY (`id`), + UNIQUE KEY `uk_backup_usage_metric__vm_id__backup_offering_id` (`vm_id`, `backup_offering_id`) +) ENGINE=InnoDB DEFAULT CHARSET=utf8; diff --git a/engine/schema/src/test/java/com/cloud/usage/dao/UsageBackupDaoImplTest.java b/engine/schema/src/test/java/com/cloud/usage/dao/UsageBackupDaoImplTest.java new file mode 100644 index 000000000000..1844ad13cd5a --- /dev/null +++ b/engine/schema/src/test/java/com/cloud/usage/dao/UsageBackupDaoImplTest.java @@ -0,0 +1,175 @@ +// 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.usage.dao; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; + +import java.util.ArrayList; +import java.util.Date; +import java.util.List; + +import org.junit.After; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Mock; +import org.mockito.MockedStatic; +import org.mockito.Mockito; +import org.mockito.Spy; +import org.mockito.junit.MockitoJUnitRunner; + +import com.cloud.usage.UsageBackupVO; +import com.cloud.utils.db.TransactionLegacy; + +@RunWith(MockitoJUnitRunner.class) +public class UsageBackupDaoImplTest { + + private static final long ZONE_ID = 1L; + private static final long ACCOUNT_ID = 2L; + private static final long DOMAIN_ID = 3L; + private static final long VM_ID = 4L; + private static final long OFFERING_ID = 5L; + + @Mock + private TransactionLegacy transactionMock; + + @Spy + private UsageBackupDaoImpl usageBackupDao; + + private MockedStatic transactionLegacyMock; + + private final Date assigned = new Date(1_000_000L); + private final Date eventDate = new Date(2_000_000L); + + @Before + public void setUp() { + transactionLegacyMock = Mockito.mockStatic(TransactionLegacy.class); + transactionLegacyMock.when(() -> TransactionLegacy.open(TransactionLegacy.USAGE_DB)).thenReturn(transactionMock); + } + + @After + public void tearDown() { + transactionLegacyMock.close(); + } + + private UsageBackupVO usage(long id, long size, long protectedSize, Date created) { + return new UsageBackupVO(id, ZONE_ID, ACCOUNT_ID, DOMAIN_ID, VM_ID, OFFERING_ID, size, protectedSize, created, null); + } + + private void mockActiveUsage(UsageBackupVO... rows) { + List list = new ArrayList<>(List.of(rows)); + doReturn(list).when(usageBackupDao).listActiveUsage(VM_ID, OFFERING_ID); + } + + @Test + public void updateMetricsIgnoresUnchangedSize() { + mockActiveUsage(usage(10L, 100L, 1000L, assigned)); + + usageBackupDao.updateMetrics(VM_ID, OFFERING_ID, 100L, 1000L, eventDate); + + verify(usageBackupDao, never()).update(anyLong(), any(UsageBackupVO.class)); + verify(usageBackupDao, never()).persist(any(UsageBackupVO.class)); + } + + @Test + public void updateMetricsClosesActiveRowAndOpensNewOneOnSizeChange() { + UsageBackupVO active = usage(10L, 100L, 1000L, assigned); + mockActiveUsage(active); + doReturn(true).when(usageBackupDao).update(anyLong(), any(UsageBackupVO.class)); + doReturn(null).when(usageBackupDao).persist(any(UsageBackupVO.class)); + + usageBackupDao.updateMetrics(VM_ID, OFFERING_ID, 9900L, 99000L, eventDate); + + Assert.assertEquals(eventDate, active.getRemoved()); + Assert.assertEquals(100L, active.getSize()); + verify(usageBackupDao).update(10L, active); + + ArgumentCaptor captor = ArgumentCaptor.forClass(UsageBackupVO.class); + verify(usageBackupDao).persist(captor.capture()); + UsageBackupVO created = captor.getValue(); + Assert.assertEquals(9900L, created.getSize()); + Assert.assertEquals(99000L, created.getProtectedSize()); + Assert.assertEquals(eventDate, created.getCreated()); + Assert.assertNull(created.getRemoved()); + Assert.assertEquals(ZONE_ID, created.getZoneId()); + Assert.assertEquals(ACCOUNT_ID, created.getAccountId()); + Assert.assertEquals(DOMAIN_ID, created.getDomainId()); + Assert.assertEquals(VM_ID, created.getVmId()); + Assert.assertEquals(OFFERING_ID, created.getBackupOfferingId()); + } + + @Test + public void updateMetricsClosesRowStartingAtEventDate() { + UsageBackupVO active = usage(10L, 0L, 0L, eventDate); + mockActiveUsage(active); + doReturn(true).when(usageBackupDao).update(anyLong(), any(UsageBackupVO.class)); + doReturn(null).when(usageBackupDao).persist(any(UsageBackupVO.class)); + + usageBackupDao.updateMetrics(VM_ID, OFFERING_ID, 100L, 1000L, eventDate); + + Assert.assertEquals(eventDate, active.getRemoved()); + verify(usageBackupDao).update(10L, active); + ArgumentCaptor captor = ArgumentCaptor.forClass(UsageBackupVO.class); + verify(usageBackupDao).persist(captor.capture()); + Assert.assertEquals(100L, captor.getValue().getSize()); + Assert.assertEquals(1000L, captor.getValue().getProtectedSize()); + Assert.assertEquals(eventDate, captor.getValue().getCreated()); + } + + @Test + public void updateMetricsIgnoresVmWithoutActiveUsage() { + mockActiveUsage(); + + usageBackupDao.updateMetrics(VM_ID, OFFERING_ID, 100L, 1000L, eventDate); + + verify(usageBackupDao, never()).update(anyLong(), any(UsageBackupVO.class)); + verify(usageBackupDao, never()).persist(any(UsageBackupVO.class)); + } + + @Test + public void updateMetricsMergesDuplicateActiveRows() { + UsageBackupVO newer = usage(11L, 100L, 1000L, new Date(1_500_000L)); + UsageBackupVO older = usage(10L, 100L, 1000L, assigned); + mockActiveUsage(newer, older); + doReturn(true).when(usageBackupDao).update(anyLong(), any(UsageBackupVO.class)); + doReturn(null).when(usageBackupDao).persist(any(UsageBackupVO.class)); + + usageBackupDao.updateMetrics(VM_ID, OFFERING_ID, 100L, 1000L, eventDate); + + Assert.assertEquals(eventDate, newer.getRemoved()); + Assert.assertEquals(eventDate, older.getRemoved()); + verify(usageBackupDao, times(2)).update(anyLong(), any(UsageBackupVO.class)); + verify(usageBackupDao, times(1)).persist(any(UsageBackupVO.class)); + } + + @Test + public void updateMetricsTreatsNullSizesAsZero() { + mockActiveUsage(usage(10L, 0L, 0L, assigned)); + + usageBackupDao.updateMetrics(VM_ID, OFFERING_ID, null, null, eventDate); + + verify(usageBackupDao, never()).update(anyLong(), any(UsageBackupVO.class)); + verify(usageBackupDao, never()).persist(any(UsageBackupVO.class)); + } +} diff --git a/server/src/main/java/org/apache/cloudstack/backup/BackupManagerImpl.java b/server/src/main/java/org/apache/cloudstack/backup/BackupManagerImpl.java index 867f4111b66a..beac12575cb3 100644 --- a/server/src/main/java/org/apache/cloudstack/backup/BackupManagerImpl.java +++ b/server/src/main/java/org/apache/cloudstack/backup/BackupManagerImpl.java @@ -69,6 +69,7 @@ import org.apache.cloudstack.backup.dao.BackupDetailsDao; import org.apache.cloudstack.backup.dao.BackupOfferingDao; import org.apache.cloudstack.backup.dao.BackupScheduleDao; +import org.apache.cloudstack.backup.dao.BackupUsageMetricDao; import org.apache.cloudstack.context.CallContext; import org.apache.cloudstack.framework.config.ConfigKey; import org.apache.cloudstack.framework.jobs.AsyncJobDispatcher; @@ -182,6 +183,8 @@ public class BackupManagerImpl extends ManagerBase implements BackupManager { @Inject private BackupDetailsDao backupDetailsDao; @Inject + private BackupUsageMetricDao backupUsageMetricDao; + @Inject private BackupScheduleDao backupScheduleDao; @Inject private BackupOfferingDao backupOfferingDao; @@ -543,9 +546,7 @@ public boolean removeVMFromBackupOffering(final Long vmId, final boolean forced) if ((result || forced) && vmInstanceDao.update(vm.getId(), vm)) { final List backups = backupDao.listByVmId(null, vm.getId()); if (backups.size() == 0) { - UsageEventUtils.publishUsageEvent(EventTypes.EVENT_VM_BACKUP_OFFERING_REMOVED_AND_BACKUPS_DELETED, vm.getAccountId(), vm.getDataCenterId(), vm.getId(), - "Backup-" + vm.getHostName() + "-" + vm.getUuid(), backupOfferingId, null, null, - Backup.class.getSimpleName(), vm.getUuid()); + publishBackupOfferingUsageRemoved(vm, backupOfferingId); } final List backupSchedules = backupScheduleDao.listByVM(vm.getId()); for(BackupSchedule backupSchedule: backupSchedules) { @@ -1653,11 +1654,21 @@ private void checkAndGenerateUsageForLastBackupDeletedAfterOfferingRemove(Virtua (vm.getBackupOfferingId() == null || vm.getBackupOfferingId() != backup.getBackupOfferingId())) { List backups = backupDao.listByVmIdAndOffering(vm.getDataCenterId(), vm.getId(), backup.getBackupOfferingId()); if (backups.size() == 0) { + publishBackupOfferingUsageRemoved(vm, backup.getBackupOfferingId()); + } + } + } + + private void publishBackupOfferingUsageRemoved(final VirtualMachine vm, final long backupOfferingId) { + Transaction.execute(new TransactionCallbackNoReturn() { + @Override + public void doInTransactionWithoutResult(TransactionStatus status) { UsageEventUtils.publishUsageEvent(EventTypes.EVENT_VM_BACKUP_OFFERING_REMOVED_AND_BACKUPS_DELETED, vm.getAccountId(), vm.getDataCenterId(), vm.getId(), "Backup-" + vm.getHostName() + "-" + vm.getUuid(), - backup.getBackupOfferingId(), null, null, Backup.class.getSimpleName(), vm.getUuid()); + backupOfferingId, null, null, Backup.class.getSimpleName(), vm.getUuid()); + backupUsageMetricDao.removeByVmAndOffering(vm.getId(), backupOfferingId); } - } + }); } @Override @@ -1985,11 +1996,7 @@ protected void runInContext() { continue; } - backupProvider.syncBackupStorageStats(dataCenter.getId()); - - syncOutOfBandBackups(backupProvider, dataCenter); - - updateBackupUsageRecords(backupProvider, dataCenter); + syncZone(backupProvider, dataCenter); } } catch (final Throwable t) { logger.error(String.format("Error trying to run backup-sync background task due to: [%s].", t.getMessage()), t); @@ -2014,7 +2021,30 @@ private void syncOutOfBandBackups(final BackupProvider backupProvider, DataCente } } - private void updateBackupUsageRecords(final BackupProvider backupProvider, DataCenter dataCenter) { + private void syncZone(final BackupProvider backupProvider, final DataCenter dataCenter) { + // Every management server runs this task. The lock lets one of them at a time sync a zone, so out-of-band + // backups are not added or removed twice and the usage metrics are compared against the last published ones. + GlobalLock lock = GlobalLock.getInternLock("backup.sync." + dataCenter.getId()); + try { + if (!lock.lock(5)) { + logger.debug("Backups of zone {} are being synced by another management server, skipping.", dataCenter); + return; + } + try { + backupProvider.syncBackupStorageStats(dataCenter.getId()); + + syncOutOfBandBackups(backupProvider, dataCenter); + + updateBackupUsageRecords(backupProvider, dataCenter); + } finally { + lock.unlock(); + } + } finally { + lock.releaseRef(); + } + } + + protected void updateBackupUsageRecords(final BackupProvider backupProvider, DataCenter dataCenter) { List vmIdsWithBackups = backupDao.listVmIdsWithBackupsInZone(dataCenter.getId()); List vmsWithBackups; if (vmIdsWithBackups.size() == 0) { @@ -2029,7 +2059,7 @@ private void updateBackupUsageRecords(final BackupProvider backupProvider, DataC Map> backupOfferingToSizeMap = new HashMap<>(); List backups = backupDao.listByVmId(null, vm.getId()); - if (backups.isEmpty() && vm.getBackupOfferingId() != null) { + if (vm.getBackupOfferingId() != null) { backupOfferingToSizeMap.put(vm.getBackupOfferingId(), new Pair<>(0L, 0L)); } for (final Backup backup: backups) { @@ -2053,14 +2083,34 @@ private void updateBackupUsageRecords(final BackupProvider backupProvider, DataC for (final Map.Entry> entry : backupOfferingToSizeMap.entrySet()) { Long offeringId = entry.getKey(); Pair sizes = entry.getValue(); - Long backupSize = sizes.first(); - Long protectedSize = sizes.second(); + publishBackupUsageMetricIfChanged(vm, offeringId, sizes.first(), sizes.second()); + } + } + } + + protected void publishBackupUsageMetricIfChanged(final VirtualMachine vm, final long offeringId, final long size, final long protectedSize) { + final BackupUsageMetricVO lastPublished = backupUsageMetricDao.findByVmAndOffering(vm.getId(), offeringId); + if (lastPublished != null && lastPublished.getSize() == size && lastPublished.getProtectedSize() == protectedSize) { + return; + } + // The event and the value the next sync compares against are committed together. + Transaction.execute(new TransactionCallbackNoReturn() { + @Override + public void doInTransactionWithoutResult(TransactionStatus status) { UsageEventUtils.publishUsageEvent(EventTypes.EVENT_VM_BACKUP_USAGE_METRIC, vm.getAccountId(), vm.getDataCenterId(), vm.getId(), "Backup-" + vm.getHostName() + "-" + vm.getUuid(), - offeringId, null, backupSize, protectedSize, + offeringId, null, size, protectedSize, Backup.class.getSimpleName(), vm.getUuid()); + if (lastPublished == null) { + backupUsageMetricDao.persist(new BackupUsageMetricVO(vm.getId(), offeringId, size, protectedSize, new Date())); + } else { + lastPublished.setSize(size); + lastPublished.setProtectedSize(protectedSize); + lastPublished.setUpdated(new Date()); + backupUsageMetricDao.update(lastPublished.getId(), lastPublished); + } } - } + }); } private Backup checkAndUpdateIfBackupEntryExistsForRestorePoint(Backup.RestorePoint restorePoint, List backupsInDb, VirtualMachine vm) { diff --git a/server/src/test/java/org/apache/cloudstack/backup/BackupManagerTest.java b/server/src/test/java/org/apache/cloudstack/backup/BackupManagerTest.java index db75602b600d..e00a67ccb56c 100644 --- a/server/src/test/java/org/apache/cloudstack/backup/BackupManagerTest.java +++ b/server/src/test/java/org/apache/cloudstack/backup/BackupManagerTest.java @@ -35,6 +35,7 @@ import java.util.ArrayList; import java.util.Collections; +import java.util.Date; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -47,6 +48,7 @@ import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; +import org.mockito.ArgumentCaptor; import org.mockito.InjectMocks; import org.mockito.Mock; import org.mockito.MockedStatic; @@ -68,6 +70,7 @@ import org.apache.cloudstack.api.response.BackupResponse; import org.apache.cloudstack.backup.dao.BackupDao; import org.apache.cloudstack.backup.dao.BackupDetailsDao; +import org.apache.cloudstack.backup.dao.BackupUsageMetricDao; import org.apache.cloudstack.backup.dao.BackupOfferingDao; import org.apache.cloudstack.backup.dao.BackupScheduleDao; import org.apache.cloudstack.context.CallContext; @@ -152,6 +155,9 @@ public class BackupManagerTest { @Mock BackupDetailsDao backupDetailsDao; + @Mock + BackupUsageMetricDao backupUsageMetricDao; + @Mock BackupProvider backupProvider; @@ -898,6 +904,133 @@ public void testBackupSyncTask() { } } + private void verifyBackupUsageMetricPublished(MockedStatic usageEventUtilsMocked, int times, Long size, Long protectedSize) { + usageEventUtilsMocked.verify(() -> UsageEventUtils.publishUsageEvent(Mockito.eq(EventTypes.EVENT_VM_BACKUP_USAGE_METRIC), Mockito.anyLong(), + Mockito.anyLong(), Mockito.anyLong(), Mockito.any(), Mockito.anyLong(), Mockito.any(), Mockito.eq(size), Mockito.eq(protectedSize), + Mockito.any(), Mockito.any()), times(times)); + } + + private DataCenterVO mockZoneWithBackedUpVm(Long dataCenterId, Long vmId, Long offeringId, Long size, Long protectedSize) { + DataCenterVO dataCenter = mock(DataCenterVO.class); + when(dataCenter.getId()).thenReturn(dataCenterId); + VMInstanceVO vm = mock(VMInstanceVO.class); + when(vm.getId()).thenReturn(vmId); + when(vm.getBackupOfferingId()).thenReturn(offeringId); + when(backupDao.listVmIdsWithBackupsInZone(dataCenterId)).thenReturn(List.of(vmId)); + when(vmInstanceDao.listByIdsIncludingRemoved(List.of(vmId))).thenReturn(List.of(vm)); + when(vmInstanceDao.listByZoneAndBackupOffering(dataCenterId, null)).thenReturn(List.of(vm)); + + BackupVO backup = new BackupVO(); + backup.setBackupOfferingId(offeringId); + backup.setSize(size); + backup.setProtectedSize(protectedSize); + when(backupDao.listByVmId(null, vmId)).thenReturn(List.of(backup)); + return dataCenter; + } + + @Test + public void updateBackupUsageRecordsPublishesAndStoresMetricWhenNoneWasPublished() { + DataCenterVO dataCenter = mockZoneWithBackedUpVm(1L, 2L, 3L, 100L, 1000L); + when(backupUsageMetricDao.findByVmAndOffering(2L, 3L)).thenReturn(null); + + try (MockedStatic usageEventUtilsMocked = Mockito.mockStatic(UsageEventUtils.class)) { + backupManager.new BackupSyncTask(backupManager).updateBackupUsageRecords(mock(BackupProvider.class), dataCenter); + + verifyBackupUsageMetricPublished(usageEventUtilsMocked, 1, 100L, 1000L); + ArgumentCaptor captor = ArgumentCaptor.forClass(BackupUsageMetricVO.class); + verify(backupUsageMetricDao).persist(captor.capture()); + Assert.assertEquals(2L, captor.getValue().getVmId()); + Assert.assertEquals(3L, captor.getValue().getBackupOfferingId()); + Assert.assertEquals(100L, captor.getValue().getSize()); + Assert.assertEquals(1000L, captor.getValue().getProtectedSize()); + } + } + + @Test + public void updateBackupUsageRecordsSkipsUnchangedMetric() { + DataCenterVO dataCenter = mockZoneWithBackedUpVm(1L, 2L, 3L, 100L, 1000L); + when(backupUsageMetricDao.findByVmAndOffering(2L, 3L)).thenReturn(new BackupUsageMetricVO(2L, 3L, 100L, 1000L, new Date())); + + try (MockedStatic usageEventUtilsMocked = Mockito.mockStatic(UsageEventUtils.class)) { + backupManager.new BackupSyncTask(backupManager).updateBackupUsageRecords(mock(BackupProvider.class), dataCenter); + + verifyBackupUsageMetricPublished(usageEventUtilsMocked, 0, 100L, 1000L); + verify(backupUsageMetricDao, never()).persist(any(BackupUsageMetricVO.class)); + verify(backupUsageMetricDao, never()).update(Mockito.anyLong(), any(BackupUsageMetricVO.class)); + } + } + + @Test + public void updateBackupUsageRecordsPublishesWhenAnotherServerPublishedADifferentValue() { + // Another management server published 200 and went down; the size is back to 100. + DataCenterVO dataCenter = mockZoneWithBackedUpVm(1L, 2L, 3L, 100L, 1000L); + BackupUsageMetricVO lastPublished = new BackupUsageMetricVO(2L, 3L, 200L, 1000L, new Date()); + when(backupUsageMetricDao.findByVmAndOffering(2L, 3L)).thenReturn(lastPublished); + + try (MockedStatic usageEventUtilsMocked = Mockito.mockStatic(UsageEventUtils.class)) { + backupManager.new BackupSyncTask(backupManager).updateBackupUsageRecords(mock(BackupProvider.class), dataCenter); + + verifyBackupUsageMetricPublished(usageEventUtilsMocked, 1, 100L, 1000L); + verify(backupUsageMetricDao).update(lastPublished.getId(), lastPublished); + Assert.assertEquals(100L, lastPublished.getSize()); + } + } + + @Test + public void backupSyncTaskSkipsZoneSyncedByAnotherServer() { + Long dataCenterId = 1L; + overrideBackupFrameworkConfigValue(); + DataCenterVO dataCenter = mock(DataCenterVO.class); + when(dataCenter.getId()).thenReturn(dataCenterId); + when(dataCenterDao.listAllZones()).thenReturn(List.of(dataCenter)); + Mockito.doReturn(backupProvider).when(backupManager).getBackupProvider(dataCenterId); + mockedGlobalLocks.add("backup.sync." + dataCenterId); + + try (MockedStatic ignored = Mockito.mockStatic(UsageEventUtils.class)) { + backupManager.new BackupSyncTask(backupManager).runInContext(); + + verify(backupProvider, never()).syncBackupStorageStats(dataCenterId); + verify(vmInstanceDao, never()).listByZoneAndBackupOffering(dataCenterId, null); + verify(backupDao, never()).listVmIdsWithBackupsInZone(dataCenterId); + } + } + + @Test + public void updateBackupUsageRecordsReportsCurrentOfferingWithoutBackups() { + Long dataCenterId = 1L; + Long vmId = 2L; + Long oldOfferingId = 3L; + Long currentOfferingId = 4L; + + DataCenterVO dataCenter = mock(DataCenterVO.class); + when(dataCenter.getId()).thenReturn(dataCenterId); + VMInstanceVO vm = mock(VMInstanceVO.class); + when(vm.getId()).thenReturn(vmId); + when(vm.getBackupOfferingId()).thenReturn(currentOfferingId); + when(backupDao.listVmIdsWithBackupsInZone(dataCenterId)).thenReturn(List.of(vmId)); + when(vmInstanceDao.listByIdsIncludingRemoved(List.of(vmId))).thenReturn(List.of(vm)); + when(vmInstanceDao.listByZoneAndBackupOffering(dataCenterId, null)).thenReturn(List.of(vm)); + + // Backups kept from an offering the VM had before; none for its current offering. + BackupVO oldBackup = new BackupVO(); + oldBackup.setBackupOfferingId(oldOfferingId); + oldBackup.setSize(100L); + oldBackup.setProtectedSize(1000L); + when(backupDao.listByVmId(null, vmId)).thenReturn(List.of(oldBackup)); + + BackupManagerImpl.BackupSyncTask backupSyncTask = backupManager.new BackupSyncTask(backupManager); + try (MockedStatic usageEventUtilsMocked = Mockito.mockStatic(UsageEventUtils.class)) { + backupSyncTask.updateBackupUsageRecords(mock(BackupProvider.class), dataCenter); + + usageEventUtilsMocked.verify(() -> UsageEventUtils.publishUsageEvent(Mockito.eq(EventTypes.EVENT_VM_BACKUP_USAGE_METRIC), Mockito.anyLong(), + Mockito.anyLong(), Mockito.eq(vmId), Mockito.any(), Mockito.eq(currentOfferingId), Mockito.any(), Mockito.eq(0L), Mockito.eq(0L), + Mockito.any(), Mockito.any())); + usageEventUtilsMocked.verify(() -> UsageEventUtils.publishUsageEvent(Mockito.eq(EventTypes.EVENT_VM_BACKUP_USAGE_METRIC), Mockito.anyLong(), + Mockito.anyLong(), Mockito.eq(vmId), Mockito.any(), Mockito.eq(oldOfferingId), Mockito.any(), Mockito.eq(100L), Mockito.eq(1000L), + Mockito.any(), Mockito.any())); + } + } + @Test public void checkCallerAccessToBackupScheduleVmTestExecuteAccessCheckMethods() { long vmId = 1L; @@ -1277,6 +1410,7 @@ public void testRemoveVMFromBackupOffering() { verify(backupScheduleDao, times(1)).remove(backupScheduleId); usageEventUtilsMocked.verify(() -> UsageEventUtils.publishUsageEvent(EventTypes.EVENT_VM_BACKUP_OFFERING_REMOVED_AND_BACKUPS_DELETED, accountId, zoneId, vmId, resourceName, offeringId, null, null, Backup.class.getSimpleName(), vmUuid)); + verify(backupUsageMetricDao).removeByVmAndOffering(vmId, offeringId); } } @@ -1581,6 +1715,7 @@ public void testDeleteBackupVmNotFound() throws ResourceAllocationException { verify(backupDao).remove(backupId); usageEventUtilsMocked.verify(() -> UsageEventUtils.publishUsageEvent(EventTypes.EVENT_VM_BACKUP_OFFERING_REMOVED_AND_BACKUPS_DELETED, accountId, zoneId, vmId, resourceName, backupOfferingId, null, null, Backup.class.getSimpleName(), vmUuid)); + verify(backupUsageMetricDao).removeByVmAndOffering(vmId, backupOfferingId); } } diff --git a/usage/src/main/java/com/cloud/usage/UsageManagerImpl.java b/usage/src/main/java/com/cloud/usage/UsageManagerImpl.java index eab371ab353f..1f96ff7a8d1b 100644 --- a/usage/src/main/java/com/cloud/usage/UsageManagerImpl.java +++ b/usage/src/main/java/com/cloud/usage/UsageManagerImpl.java @@ -2038,7 +2038,7 @@ private void createVmSnapshotOnPrimaryEvent(UsageEventVO event) { } } - private void createBackupEvent(final UsageEventVO event) { + protected void createBackupEvent(final UsageEventVO event) { Long vmId = event.getResourceId(); Long zoneId = event.getZoneId(); Long accountId = event.getAccountId(); @@ -2048,12 +2048,18 @@ private void createBackupEvent(final UsageEventVO event) { Date created = event.getCreateDate(); if (EventTypes.EVENT_VM_BACKUP_OFFERING_ASSIGN.equals(event.getType())) { + // Removing the offering while keeping the backups leaves the usage active, so re-assigning the + // same offering must not add a second row that bills the same backups again. + if (!usageBackupDao.listActiveUsage(vmId, backupOfferingId).isEmpty()) { + logger.debug("VM [{}] already has active usage for backup offering [{}], not creating another usage entry.", vmId, backupOfferingId); + return; + } final UsageBackupVO backupVO = new UsageBackupVO(zoneId, accountId, domainId, vmId, backupOfferingId, created); usageBackupDao.persist(backupVO); } else if (EventTypes.EVENT_VM_BACKUP_OFFERING_REMOVED_AND_BACKUPS_DELETED.equals(event.getType())) { usageBackupDao.removeUsage(accountId, vmId, backupOfferingId, event.getCreateDate()); } else if (EventTypes.EVENT_VM_BACKUP_USAGE_METRIC.equals(event.getType())) { - usageBackupDao.updateMetrics(vmId, backupOfferingId, event.getSize(), event.getVirtualSize()); + usageBackupDao.updateMetrics(vmId, backupOfferingId, event.getSize(), event.getVirtualSize(), created); } } diff --git a/usage/src/test/java/com/cloud/usage/UsageManagerImplTest.java b/usage/src/test/java/com/cloud/usage/UsageManagerImplTest.java index 03939e430d25..81439ae0b995 100644 --- a/usage/src/test/java/com/cloud/usage/UsageManagerImplTest.java +++ b/usage/src/test/java/com/cloud/usage/UsageManagerImplTest.java @@ -17,9 +17,11 @@ package com.cloud.usage; import java.util.ArrayList; +import java.util.Date; import java.util.List; import com.cloud.event.dao.UsageEventDetailsDao; +import com.cloud.usage.dao.UsageBackupDao; import com.cloud.usage.dao.UsageVMSnapshotDao; import org.junit.Before; import org.junit.Test; @@ -58,6 +60,9 @@ public class UsageManagerImplTest { @Mock private AccountDao accountDaoMock; + @Mock + private UsageBackupDao usageBackupDaoMock; + @Mock private UsageVPNUserVO vpnUserMock; @@ -243,4 +248,46 @@ public void handleVMSnapshotEventTestEventIsNeitherAddNorRemove() { Mockito.verify(usageManagerImpl, Mockito.never()).createUsageVpnUser(usageEventVOMock,accountMock); Mockito.verify(usageManagerImpl, Mockito.never()).deleteUsageVpnUser(usageEventVOMock, accountMock); } + + private void mockBackupEvent(String type) { + Mockito.when(usageEventVOMock.getType()).thenReturn(type); + Mockito.when(usageEventVOMock.getResourceId()).thenReturn(10L); + Mockito.when(usageEventVOMock.getZoneId()).thenReturn(1L); + Mockito.when(usageEventVOMock.getAccountId()).thenReturn(accountMockId); + Mockito.when(usageEventVOMock.getOfferingId()).thenReturn(20L); + } + + @Test + public void createBackupEventTestAssignCreatesUsage() { + mockBackupEvent(EventTypes.EVENT_VM_BACKUP_OFFERING_ASSIGN); + Mockito.when(usageEventVOMock.getCreateDate()).thenReturn(new Date()); + Mockito.when(usageBackupDaoMock.listActiveUsage(10L, 20L)).thenReturn(new ArrayList<>()); + + usageManagerImpl.createBackupEvent(usageEventVOMock); + + Mockito.verify(usageBackupDaoMock).persist(Mockito.any(UsageBackupVO.class)); + } + + @Test + public void createBackupEventTestAssignSkipsWhenUsageIsActive() { + mockBackupEvent(EventTypes.EVENT_VM_BACKUP_OFFERING_ASSIGN); + Mockito.when(usageBackupDaoMock.listActiveUsage(10L, 20L)).thenReturn(List.of(Mockito.mock(UsageBackupVO.class))); + + usageManagerImpl.createBackupEvent(usageEventVOMock); + + Mockito.verify(usageBackupDaoMock, Mockito.never()).persist(Mockito.any(UsageBackupVO.class)); + } + + @Test + public void createBackupEventTestMetricPassesEventDate() { + Date eventDate = new Date(); + mockBackupEvent(EventTypes.EVENT_VM_BACKUP_USAGE_METRIC); + Mockito.when(usageEventVOMock.getCreateDate()).thenReturn(eventDate); + Mockito.when(usageEventVOMock.getSize()).thenReturn(100L); + Mockito.when(usageEventVOMock.getVirtualSize()).thenReturn(1000L); + + usageManagerImpl.createBackupEvent(usageEventVOMock); + + Mockito.verify(usageBackupDaoMock).updateMetrics(10L, 20L, 100L, 1000L, eventDate); + } }