Copilot commented on code in PR #13907:
URL: https://github.com/apache/cloudstack/pull/13907#discussion_r4157640677
##########
server/src/main/java/com/cloud/vm/UserVmManagerImpl.java:
##########
@@ -3892,6 +3898,8 @@ public boolean deleteVmGroup(long groupId) {
public boolean addInstanceToGroup(final long userVmId, String groupName) {
UserVmVO vm = _vmDao.findById(userVmId);
+
instanceBootGroupMembershipGuard.validateVmEligibleForGroupMembership(userVmId);
Review Comment:
This new guard protects VM destruction and VM/group membership additions,
but deleting an `InstanceGroup` still removes its mappings and the group
without removing its `InstanceBootGroupMember` row. That leaves a boot-group
member pointing at a deleted group; later start/stop/list operations can show
the stale member and a start tier with no VMs. Block group deletion while it is
a boot-group member or remove that member atomically.
##########
server/src/main/java/org/apache/cloudstack/vm/bootgroup/InstanceBootGroupApiServiceImpl.java:
##########
@@ -0,0 +1,912 @@
+// 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.vm.bootgroup;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.Comparator;
+import java.util.List;
+import java.util.Objects;
+import java.util.stream.Collectors;
+
+import javax.inject.Inject;
+
+import
org.apache.cloudstack.api.command.user.bootgroup.AddMemberToInstanceBootGroupCmd;
+import
org.apache.cloudstack.api.command.user.bootgroup.CreateInstanceBootGroupCmd;
+import
org.apache.cloudstack.api.command.user.bootgroup.CreateInstanceBootGroupReadinessRuleCmd;
+import
org.apache.cloudstack.api.command.user.bootgroup.DeleteInstanceBootGroupCmd;
+import
org.apache.cloudstack.api.command.user.bootgroup.DeleteInstanceBootGroupReadinessRuleCmd;
+import
org.apache.cloudstack.api.command.user.bootgroup.ListInstanceBootGroupMembersCmd;
+import
org.apache.cloudstack.api.command.user.bootgroup.ListInstanceBootGroupReadinessRulesCmd;
+import
org.apache.cloudstack.api.command.user.bootgroup.ListInstanceBootGroupsCmd;
+import
org.apache.cloudstack.api.command.user.bootgroup.RebootInstanceBootGroupCmd;
+import
org.apache.cloudstack.api.command.user.bootgroup.RemoveInstanceBootGroupMemberCmd;
+import
org.apache.cloudstack.api.command.user.bootgroup.StartInstanceBootGroupCmd;
+import
org.apache.cloudstack.api.command.user.bootgroup.StopInstanceBootGroupCmd;
+import
org.apache.cloudstack.api.command.user.bootgroup.UpdateInstanceBootGroupCmd;
+import
org.apache.cloudstack.api.command.user.bootgroup.UpdateInstanceBootGroupMemberCmd;
+import
org.apache.cloudstack.api.command.user.bootgroup.UpdateInstanceBootGroupReadinessRuleCmd;
+import org.apache.cloudstack.api.query.dao.InstanceBootGroupJoinDao;
+import org.apache.cloudstack.api.query.vo.InstanceBootGroupJoinVO;
+import org.apache.cloudstack.api.response.InstanceBootGroupMemberChildResponse;
+import org.apache.cloudstack.api.response.InstanceBootGroupMemberResponse;
+import
org.apache.cloudstack.api.response.InstanceBootGroupReadinessRuleResponse;
+import org.apache.cloudstack.api.response.InstanceBootGroupResponse;
+import org.apache.cloudstack.api.response.ListResponse;
+import org.apache.cloudstack.context.CallContext;
+import
org.apache.cloudstack.vm.bootgroup.readiness.InstanceBootGroupReadinessRule;
+import
org.apache.cloudstack.vm.bootgroup.readiness.InstanceBootGroupReadinessRuleService;
+import org.apache.cloudstack.vm.bootgroup.readiness.ReadinessChecker;
+import org.apache.commons.lang3.EnumUtils;
+import org.apache.commons.lang3.StringUtils;
+import org.jetbrains.annotations.NotNull;
+import org.springframework.stereotype.Component;
+
+import com.cloud.api.ApiResponseHelper;
+import com.cloud.event.ActionEvent;
+import com.cloud.event.EventTypes;
+import com.cloud.exception.InvalidParameterValueException;
+import com.cloud.exception.PermissionDeniedException;
+import com.cloud.projects.Project;
+import com.cloud.user.Account;
+import com.cloud.user.AccountManager;
+import com.cloud.uservm.UserVm;
+import com.cloud.utils.Pair;
+import com.cloud.utils.Ternary;
+import com.cloud.utils.component.PluggableService;
+import com.cloud.utils.db.Filter;
+import com.cloud.utils.db.SearchBuilder;
+import com.cloud.utils.db.SearchCriteria;
+import com.cloud.utils.db.Transaction;
+import com.cloud.utils.db.TransactionCallback;
+import com.cloud.vm.InstanceGroupVMMapVO;
+import com.cloud.vm.InstanceGroupVO;
+import com.cloud.vm.UserVmVO;
+import com.cloud.vm.dao.InstanceBootGroupDao;
+import com.cloud.vm.dao.InstanceBootGroupDetailsDao;
+import com.cloud.vm.dao.InstanceBootGroupMemberDao;
+import com.cloud.vm.dao.InstanceBootGroupReadinessCheckResultDao;
+import com.cloud.vm.dao.InstanceBootGroupReadinessRuleDao;
+import com.cloud.vm.dao.InstanceBootGroupReadinessRuleDetailsDao;
+import com.cloud.vm.dao.InstanceGroupDao;
+import com.cloud.vm.dao.InstanceGroupVMMapDao;
+import com.cloud.vm.dao.UserVmDao;
+
+/**
+ * API-facing half of the Instance Boot Group feature: ACL, param validation,
response building,
+ * command registration. Delegates orchestration/hypervisor work to {@link
InstanceBootGroupManager}
+ * and membership eligibility checks to {@link
InstanceBootGroupMembershipGuard}.
+ */
+@Component
+public class InstanceBootGroupApiServiceImpl implements
InstanceBootGroupService, PluggableService {
+
+ @Inject
+ private InstanceBootGroupDao instanceBootGroupDao;
+
+ @Inject
+ private InstanceBootGroupJoinDao instanceBootGroupJoinDao;
+
+ @Inject
+ private InstanceBootGroupMemberDao instanceBootGroupMemberDao;
+
+ @Inject
+ private AccountManager accountManager;
+
+ @Inject
+ private UserVmDao userVmDao;
+
+ @Inject
+ private InstanceGroupDao instanceGroupDao;
+
+ @Inject
+ private InstanceBootGroupManager instanceBootGroupManager;
+
+ @Inject
+ private InstanceBootGroupMembershipGuard instanceBootGroupMembershipGuard;
+
+ @Inject
+ private InstanceBootGroupReadinessRuleService
instanceBootGroupReadinessRuleService;
+
+ @Inject
+ private InstanceBootGroupReadinessRuleDao
instanceBootGroupReadinessRuleDao;
+
+ @Inject
+ private InstanceBootGroupReadinessRuleDetailsDao
instanceBootGroupReadinessRuleDetailsDao;
+
+ @Inject
+ private InstanceBootGroupReadinessCheckResultDao
instanceBootGroupReadinessCheckResultDao;
+
+ @Inject
+ private InstanceBootGroupDetailsDao instanceBootGroupDetailsDao;
+
+ @Inject
+ private InstanceGroupVMMapDao instanceGroupVMMapDao;
+
+ @NotNull
+ protected InstanceBootGroupVO getGroupAndCheckAccess(long id) {
+ InstanceBootGroupVO group = instanceBootGroupDao.findById(id);
+ if (group == null) {
+ throw new InvalidParameterValueException("Unable to find instance
boot group with ID: " + id);
+ }
+ Account caller = CallContext.current().getCallingAccount();
+ accountManager.checkAccess(caller, null, true, group);
+ return group;
+ }
+
+ protected InstanceBootGroupResponse
createInstanceBootGroupResponse(InstanceBootGroupJoinVO bootGroup) {
+ InstanceBootGroupResponse response = new InstanceBootGroupResponse();
+ response.setId(bootGroup.getUuid());
+ response.setName(bootGroup.getName());
+ response.setDescription(bootGroup.getDescription());
+ response.setCreated(bootGroup.getCreated());
+ ApiResponseHelper.populateOwner(response, bootGroup);
+
+ String timeoutOverride =
instanceBootGroupDetailsDao.getDetail(bootGroup.getId(),
InstanceBootGroupManagerImpl.ReadinessAttemptTimeoutSeconds.key());
+ response.setReadinessAttemptTimeoutSeconds(timeoutOverride != null ?
Long.parseLong(timeoutOverride) :
InstanceBootGroupManagerImpl.ReadinessAttemptTimeoutSeconds.value());
+ String maxRetryOverride =
instanceBootGroupDetailsDao.getDetail(bootGroup.getId(),
InstanceBootGroupManagerImpl.ReadinessMaxRetryAttempts.key());
+ response.setReadinessMaxRetryAttempts(maxRetryOverride != null ?
Long.parseLong(maxRetryOverride) :
InstanceBootGroupManagerImpl.ReadinessMaxRetryAttempts.value());
+ String rebootOnRetryOverride =
instanceBootGroupDetailsDao.getDetail(bootGroup.getId(),
InstanceBootGroupManagerImpl.ReadinessRebootOnRetry.key());
+ response.setReadinessRebootOnRetry(rebootOnRetryOverride != null ?
Boolean.parseBoolean(rebootOnRetryOverride) :
InstanceBootGroupManagerImpl.ReadinessRebootOnRetry.value());
+ String initialDelayOverride =
instanceBootGroupDetailsDao.getDetail(bootGroup.getId(),
InstanceBootGroupManagerImpl.ReadinessInitialDelaySeconds.key());
+ response.setReadinessInitialDelaySeconds(initialDelayOverride != null
? Long.parseLong(initialDelayOverride) :
InstanceBootGroupManagerImpl.ReadinessInitialDelaySeconds.value());
+
+ response.setObjectName("instancebootgroup");
+ return response;
+ }
+
+ @Override
+ @ActionEvent(eventType = EventTypes.EVENT_INSTANCE_BOOT_GROUP_CREATE,
eventDescription = "creating Instance Boot Group")
+ public InstanceBootGroup
createInstanceBootGroup(CreateInstanceBootGroupCmd cmd) {
+ Account caller = CallContext.current().getCallingAccount();
+ Account owner = accountManager.finalizeOwner(caller,
cmd.getAccountName(), cmd.getDomainId(), cmd.getProjectId());
+
+ if (instanceBootGroupDao.isNameInUse(owner.getId(), cmd.getName())) {
+ throw new InvalidParameterValueException("An instance boot group
with name '" + cmd.getName() + "' already exists in this account");
+ }
+
+ return Transaction.execute((TransactionCallback<InstanceBootGroupVO>)
status -> {
+ InstanceBootGroupVO group = new InstanceBootGroupVO(cmd.getName(),
cmd.getDescription(), owner.getId(), owner.getDomainId());
+ group = instanceBootGroupDao.persist(group);
+ CallContext.current().setEventResourceId(group.getId());
+
+ if (cmd.getReadinessAttemptTimeoutSeconds() != null) {
+ setOrClearOverride(group.getId(),
InstanceBootGroupManagerImpl.ReadinessAttemptTimeoutSeconds.key(),
cmd.getReadinessAttemptTimeoutSeconds());
+ }
+ if (cmd.getReadinessMaxRetryAttempts() != null) {
+ setOrClearOverride(group.getId(),
InstanceBootGroupManagerImpl.ReadinessMaxRetryAttempts.key(),
cmd.getReadinessMaxRetryAttempts());
+ }
+ if (cmd.getReadinessRebootOnRetry() != null) {
+ instanceBootGroupDetailsDao.setDetail(group.getId(),
InstanceBootGroupManagerImpl.ReadinessRebootOnRetry.key(),
String.valueOf(cmd.getReadinessRebootOnRetry()));
+ }
+ if (cmd.getReadinessInitialDelaySeconds() != null) {
+ setOrClearOverride(group.getId(),
InstanceBootGroupManagerImpl.ReadinessInitialDelaySeconds.key(),
cmd.getReadinessInitialDelaySeconds());
+ }
+
+ return group;
+ });
+ }
+
+ @Override
+ @ActionEvent(eventType = EventTypes.EVENT_INSTANCE_BOOT_GROUP_DELETE,
eventDescription = "deleting Instance Boot Group")
+ public boolean deleteInstanceBootGroup(DeleteInstanceBootGroupCmd cmd) {
+ InstanceBootGroupVO group = getGroupAndCheckAccess(cmd.getId());
+ return Transaction.execute((TransactionCallback<Boolean>) status -> {
+ instanceBootGroupMemberDao.deleteByBootGroupId(group.getId());
+ instanceBootGroupDao.remove(group.getId());
+ return true;
+ });
+ }
+
+ @Override
+ @ActionEvent(eventType = EventTypes.EVENT_INSTANCE_BOOT_GROUP_UPDATE,
eventDescription = "updating Instance Boot Group")
+ public InstanceBootGroup
updateInstanceBootGroup(UpdateInstanceBootGroupCmd cmd) {
+ InstanceBootGroupVO group = getGroupAndCheckAccess(cmd.getId());
+
+ if (cmd.getName() != null && !Objects.equals(cmd.getName(),
group.getName())) {
+ Account owner = accountManager.getAccount(group.getAccountId());
+ if (instanceBootGroupDao.isNameInUse(owner.getId(),
cmd.getName())) {
+ throw new InvalidParameterValueException("An instance boot
group with name '" + cmd.getName() + "' already exists in this account");
+ }
+ group.setName(cmd.getName());
+ }
+ if (cmd.getDescription() != null) {
+ group.setDescription(cmd.getDescription());
+ }
+
+ return Transaction.execute((TransactionCallback<InstanceBootGroupVO>)
status -> {
+ if (cmd.getReadinessAttemptTimeoutSeconds() != null) {
+ setOrClearOverride(group.getId(),
InstanceBootGroupManagerImpl.ReadinessAttemptTimeoutSeconds.key(),
cmd.getReadinessAttemptTimeoutSeconds());
+ }
+ if (cmd.getReadinessMaxRetryAttempts() != null) {
+ setOrClearOverride(group.getId(),
InstanceBootGroupManagerImpl.ReadinessMaxRetryAttempts.key(),
cmd.getReadinessMaxRetryAttempts());
+ }
+ if (cmd.getReadinessRebootOnRetry() != null) {
+ instanceBootGroupDetailsDao.setDetail(group.getId(),
InstanceBootGroupManagerImpl.ReadinessRebootOnRetry.key(),
String.valueOf(cmd.getReadinessRebootOnRetry()));
+ }
+ if (cmd.getReadinessInitialDelaySeconds() != null) {
+ setOrClearOverride(group.getId(),
InstanceBootGroupManagerImpl.ReadinessInitialDelaySeconds.key(),
cmd.getReadinessInitialDelaySeconds());
+ }
+
+ instanceBootGroupDao.update(group.getId(), group);
+ return instanceBootGroupDao.findById(group.getId());
+ });
+ }
+
+ private void setOrClearOverride(long bootGroupId, String key, long value) {
+ if (value < 0) {
+ instanceBootGroupDetailsDao.setDetail(bootGroupId, key, null);
+ } else {
+ instanceBootGroupDetailsDao.setDetail(bootGroupId, key,
String.valueOf(value));
+ }
+ }
+
+ @Override
+ public ListResponse<InstanceBootGroupResponse>
listInstanceBootGroups(ListInstanceBootGroupsCmd cmd) {
+ final CallContext ctx = CallContext.current();
+ final Account caller = ctx.getCallingAccount();
+ final Long id = cmd.getId();
+ final String keyword = cmd.getKeyword();
+
Review Comment:
`ListInstanceBootGroupsCmd` exposes `virtualmachineid` and
`instancegroupid`, but this service never reads either parameter or adds a
member constraint to the search. Requests using those filters therefore return
every accessible boot group instead of only groups containing the requested VM
or instance group. Apply the member filter in the query/DAO and adjust the
count before building the response.
##########
server/src/main/java/org/apache/cloudstack/vm/bootgroup/InstanceBootGroupManagerImpl.java:
##########
@@ -0,0 +1,604 @@
+// 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.vm.bootgroup;
+
+import java.util.ArrayList;
+import java.util.Collection;
+import java.util.Comparator;
+import java.util.HashMap;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.TreeMap;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.stream.Collectors;
+
+import javax.inject.Inject;
+import javax.naming.ConfigurationException;
+
+import org.apache.cloudstack.api.ApiCommandResourceType;
+import org.apache.cloudstack.context.CallContext;
+import org.apache.cloudstack.framework.config.ConfigKey;
+import org.apache.cloudstack.framework.config.Configurable;
+import org.apache.cloudstack.managed.context.ManagedContextRunnable;
+import
org.apache.cloudstack.vm.bootgroup.readiness.InstanceBootGroupReadinessRule;
+import
org.apache.cloudstack.vm.bootgroup.readiness.InstanceBootGroupReadinessRuleService;
+import org.springframework.stereotype.Component;
+
+import com.cloud.utils.component.ManagerBase;
+import com.cloud.utils.concurrency.NamedThreadFactory;
+import com.cloud.utils.exception.CloudRuntimeException;
+import com.cloud.vm.InstanceGroup;
+import com.cloud.vm.UserVmService;
+import com.cloud.vm.UserVmVO;
+import com.cloud.vm.VirtualMachine;
+import com.cloud.vm.VirtualMachineManager;
+import com.cloud.vm.dao.InstanceBootGroupDetailsDao;
+import com.cloud.vm.dao.InstanceBootGroupMemberDao;
+import com.cloud.vm.dao.InstanceGroupDao;
+import com.cloud.vm.dao.InstanceGroupVMMapDao;
+import com.cloud.vm.dao.UserVmDao;
+
+/**
+ * Backend/orchestration half of the Instance Boot Group feature — tier
concurrency, hypervisor
+ * start/stop/reboot calls, and readiness-gated tier progression. API-cmd
handling (ACL, validation,
+ * response building, command registration) lives in {@code
InstanceBootGroupApiServiceImpl}, which
+ * delegates here with resolved domain objects.
+ *
+ * <p>Per-VM timeout/reboot-attempt bookkeeping during a start is kept purely
in-memory, scoped to the
+ * async job thread executing the start — it is not persisted. Surviving a
management-server restart
+ * mid-run is explicitly not a goal here; if the process restarts, the job
(and this bookkeeping) is
+ * simply lost, same as any other in-flight async job. Current readiness is
queryable at any time via
+ * {@code listInstanceBootGroupMembers?details=readiness}, not via a separate
run-history API.</p>
+ */
+@Component
+public class InstanceBootGroupManagerImpl extends ManagerBase implements
InstanceBootGroupManager, Configurable {
+
+ public static final ConfigKey<Long> ReadinessAttemptTimeoutSeconds = new
ConfigKey<>("Advanced", Long.class,
+ "instance.boot.group.readiness.timeout.seconds", "300",
+ "How long to wait (in seconds) for an instance to become ready
during boot group orchestration before starting a new readiness retry attempt.
Overridable per boot group.", true);
+
+ public static final ConfigKey<Long> ReadinessMaxRetryAttempts = new
ConfigKey<>("Advanced", Long.class,
+ "instance.boot.group.readiness.max.retry.attempts", "5",
+ "Maximum number of readiness retry attempts for an instance that
fails to become ready during boot group orchestration before the boot group
start is halted. Overridable per boot group.", true);
+
+ public static final ConfigKey<Long> ReadinessPollIntervalSeconds = new
ConfigKey<>("Advanced", Long.class,
+ "instance.boot.group.readiness.poll.interval.seconds", "10",
+ "How often (in seconds) to re-check instance/instance-group
readiness during boot group orchestration, including the minimum pause after a
readiness retry attempt that did not reboot the instance before repeating the
check that just failed. A very low value can cause rapid repeated
(\"hammering\") readiness retries against an instance/VR/host. Global only, not
overridable per boot group.", true);
+
+ public static final ConfigKey<Long> ReadinessInitialDelaySeconds = new
ConfigKey<>("Advanced", Long.class,
+ "instance.boot.group.readiness.initial.delay.seconds", "30",
+ "How long to wait (in seconds) after starting or rebooting an
instance before its first readiness check of that attempt, giving the guest
OS/agent/network time to come up. Overridable per boot group.", true);
+
+ public static final ConfigKey<Boolean> ReadinessRebootOnRetry = new
ConfigKey<>("Advanced", Boolean.class,
+ "instance.boot.group.readiness.reboot.on.retry", "false",
+ "Whether to reboot an instance between readiness retry attempts
during boot group orchestration, instead of just waiting longer. Overridable
per boot group.", true);
+
+ public static final ConfigKey<Long> ReadinessCheckConcurrency = new
ConfigKey<>("Advanced", Long.class,
+ "instance.boot.group.readiness.check.concurrency", "10",
+ "Maximum number of instances within a single boot-order tier whose
readiness is checked concurrently during boot group orchestration, so one slow
check cannot delay every other instance's check in the same poll. Global only,
not overridable per boot group.", true);
+
+ public static final ConfigKey<Long> MaxMembersPerBootGroup = new
ConfigKey<>("Advanced", Long.class,
+ "instance.boot.group.max.members", "10",
+ "Maximum number of members that can be added to a single instance
boot group.", true, ConfigKey.Scope.Domain);
+
+ @Inject
+ private InstanceBootGroupMemberDao instanceBootGroupMemberDao;
+
+ @Inject
+ private UserVmService userVmService;
+
+ @Inject
+ private UserVmDao userVmDao;
+
+ @Inject
+ private InstanceGroupDao instanceGroupDao;
+
+ @Inject
+ private InstanceGroupVMMapDao instanceGroupVMMapDao;
+
+ @Inject
+ private VirtualMachineManager virtualMachineManager;
+
+ @Inject
+ private InstanceBootGroupReadinessRuleService
instanceBootGroupReadinessRuleService;
+
+ @Inject
+ private InstanceBootGroupDetailsDao instanceBootGroupDetailsDao;
+
+ @Override
+ public boolean configure(String name, Map<String, Object> params) throws
ConfigurationException {
+ VirtualMachine.State.getStateMachine().registerListener(new
InstanceBootGroupVmStateListener(instanceBootGroupReadinessRuleService));
+ return true;
+ }
+
+ @Override
+ public String getConfigComponentName() {
+ return InstanceBootGroupManagerImpl.class.getSimpleName();
+ }
+
+ @Override
+ public ConfigKey<?>[] getConfigKeys() {
+ return new ConfigKey<?>[]{ReadinessAttemptTimeoutSeconds,
ReadinessMaxRetryAttempts, ReadinessPollIntervalSeconds,
ReadinessInitialDelaySeconds, ReadinessRebootOnRetry, ReadinessCheckConcurrency,
+ MaxMembersPerBootGroup};
+ }
+
+ /** In-memory-only per-VM progress for a single start attempt — never
persisted. Package-visible
+ * (rather than private) purely so tests can construct one and inspect
retry/give-up state
+ * directly instead of via reflection. */
+ protected static final class VmProgress {
+ private final long vmId;
+ private final Long bootGroupMemberId;
+ private boolean ready;
+ /** Set when this VM has exhausted its retry attempts but belongs to
an InstanceGroup
+ * member, so the group's own readiness rule (e.g. a quorum rule)
gets the final say
+ * instead of this one VM halting the whole boot group. */
+ protected boolean gaveUp;
+ protected int retryAttempts;
+ /** Anchor for the per-attempt timeout window; reset on every retry,
rebooted or not. */
+ private long enteredWaitAtMs;
+ /** Anchor for the initial-delay grace period; only reset on the
initial start and on an
+ * actual reboot — a no-op retry (reboot-on-retry disabled) leaves
this alone, since there's
+ * no fresh boot to wait out. */
+ private long lastBootedAtMs;
+
+ protected VmProgress(long vmId, Long bootGroupMemberId) {
+ this.vmId = vmId;
+ this.bootGroupMemberId = bootGroupMemberId;
+ }
+
+ private String getAttemptsLog(long maxAttempts) {
+ return String.format("%d/%d", retryAttempts + 1, maxAttempts);
+ }
+ }
+
+ @Override
+ public void startInstanceBootGroup(InstanceBootGroupVO group) {
+ List<InstanceBootGroupMemberVO> members =
instanceBootGroupMemberDao.listByBootGroupId(group.getId());
+ Map<Integer, List<InstanceBootGroupMemberVO>> tiers =
groupByOrder(members);
+ logger.info("Starting {}: {} tier(s), {} member(s) total", group,
tiers.size(), members.size());
+ long groupStartedAtMs = System.currentTimeMillis();
+
+ for (Map.Entry<Integer, List<InstanceBootGroupMemberVO>> tierEntry :
tiers.entrySet()) {
+ int tierOrder = tierEntry.getKey();
+ List<InstanceBootGroupMemberVO> tierMembers = tierEntry.getValue();
+
+ Map<Long, VmProgress> progressByVmId = new LinkedHashMap<>();
+ for (InstanceBootGroupMemberVO member : tierMembers) {
+ for (Long vmId : resolveVmIds(List.of(member))) {
+ progressByVmId.put(vmId, new VmProgress(vmId,
member.getId()));
+ }
+ }
+ List<Long> tierVmIds = new ArrayList<>(progressByVmId.keySet());
+ logger.info("Starting tier {} of {}: {} member(s), {} VM(s)",
tierOrder, group, tierMembers.size(), tierVmIds.size());
+ long tierStartedAtMs = System.currentTimeMillis();
+
+ try {
+ runTierConcurrently(tierVmIds, group, "start", vmId -> {
+ UserVmVO vm = userVmDao.findById(vmId);
+ boolean alreadyRunning = vm != null &&
VirtualMachine.State.Running.equals(vm.getState());
+ if (vm != null && !alreadyRunning) {
+ userVmService.startVirtualMachine(vm, null);
+ }
+ anchorInitialDelay(group, progressByVmId.get(vmId), vm,
alreadyRunning);
+ });
+ } catch (CloudRuntimeException e) {
+ halt(group, "Failed to start a VM in tier " + tierOrder + ": "
+ e.getMessage());
+ throw e;
+ }
+
+ waitForTierReady(group, tierOrder, tierMembers, progressByVmId);
+ logger.info("Tier {} of {} is ready ({}ms)", tierOrder, group,
System.currentTimeMillis() - tierStartedAtMs);
+ }
+
+ logger.info("{} start completed ({}ms)", group,
System.currentTimeMillis() - groupStartedAtMs);
+ }
+
+ /**
+ * If the VM was already running, anchors the initial-delay grace period
to when CloudStack last
+ * confirmed its power state rather than to "now" — so it waits out only
what's left of the
+ * delay (or none) instead of a full fresh wait it doesn't need.
+ */
+ private void anchorInitialDelay(InstanceBootGroupVO group, VmProgress
progress, UserVmVO vm, boolean alreadyRunning) {
+ long now = System.currentTimeMillis();
+ progress.enteredWaitAtMs = now;
+ if (alreadyRunning && vm.getPowerStateUpdateTime() != null) {
+ progress.lastBootedAtMs = vm.getPowerStateUpdateTime().getTime();
+ logger.debug("{} was already running (power state last confirmed
{}); readiness checks begin after any remaining portion of the {}s initial
delay",
+ vm, vm.getPowerStateUpdateTime(),
effectiveInitialDelaySeconds(group));
+ } else {
+ progress.lastBootedAtMs = now;
+ logger.debug("{} start action completed; readiness checks begin
after the {}s initial delay", vm, effectiveInitialDelaySeconds(group));
+ }
+ }
+
+ /**
+ * Polls every not-yet-settled VM in the tier concurrently (bounded by
+ * {@code ReadinessCheckConcurrency}) until the whole tier — VMs and any
InstanceGroup
+ * members — reports ready, or a halt is triggered.
+ */
+ private void waitForTierReady(InstanceBootGroupVO group, int tierOrder,
List<InstanceBootGroupMemberVO> tierMembers, Map<Long, VmProgress>
progressByVmId) {
+ Map<Long, InstanceBootGroupMemberVO> memberById = new HashMap<>();
+ Map<Long, Boolean> membersReadyStatus = new ConcurrentHashMap<>();
+ for (InstanceBootGroupMemberVO member : tierMembers) {
+ memberById.put(member.getId(), member);
+ membersReadyStatus.put(member.getId(), false);
+ }
+ final long effectiveMaxRetryAttempts =
effectiveMaxRetryAttempts(group);
+ final long effectiveTimeoutSeconds = effectiveTimeoutSeconds(group);
+ final long effectivePollIntervalSeconds =
effectivePollIntervalSeconds();
+ final long pollIntervalMs = effectivePollIntervalSeconds * 1000L;
+ final boolean effectiveRebootOnRetry = effectiveRebootOnRetry(group);
+ int concurrency = (int) Math.max(1, Math.min(progressByVmId.size(),
effectiveReadinessCheckConcurrency()));
+ // Bound the polling loop: initial delay + maxAttempts full timeout
windows + inter-poll sleeps between them.
+ // At least one timeout window is always budgeted so a single check is
never starved out, even if
+ // maxAttempts is configured to 0.
+ long maxWaitMs = (effectiveInitialDelaySeconds(group) + Math.max(1,
effectiveMaxRetryAttempts) * effectiveTimeoutSeconds
+ + Math.max(0, effectiveMaxRetryAttempts - 1) *
effectivePollIntervalSeconds) * 1000L;
+ long deadline = System.currentTimeMillis() + maxWaitMs;
+ logger.debug("Waiting for tier {} of {} to become ready: {} VM(s)
tracked, timeout={}s, pollInterval={}s, checkConcurrency={}, maxWait={}ms",
+ tierOrder, group, progressByVmId.size(),
effectiveTimeoutSeconds, effectivePollIntervalSeconds, concurrency, maxWaitMs);
+
+ CallContext callerContext = CallContext.current();
+ ExecutorService readinessExecutor =
Executors.newFixedThreadPool(concurrency, new
NamedThreadFactory("InstanceBootGroup-readiness-" + tierOrder));
+ try {
+ while (System.currentTimeMillis() < deadline) {
+ List<Future<Void>> futures = new ArrayList<>();
+ for (VmProgress progress : progressByVmId.values()) {
+ if (progress.ready || progress.gaveUp) {
+ continue;
+ }
+ futures.add(readinessExecutor.submit(() -> {
+ CallContext.register(callerContext,
ApiCommandResourceType.VirtualMachine);
+ try {
+ checkVmReadiness(group, progress, memberById,
membersReadyStatus,
+ effectiveMaxRetryAttempts,
effectiveTimeoutSeconds, effectiveRebootOnRetry);
+ } finally {
+ CallContext.unregister();
+ }
+ return null;
+ }));
+ }
+ for (Future<Void> future : futures) {
+ try {
+ future.get();
+ } catch (ExecutionException e) {
+ Throwable cause = e.getCause() != null ? e.getCause()
: e;
+ if (cause instanceof CloudRuntimeException) {
+ throw (CloudRuntimeException) cause;
+ }
+ throw new CloudRuntimeException("Failed to evaluate
readiness for a VM in tier " + tierOrder + " of " + group.getName() + ": " +
cause.getMessage(), cause);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new CloudRuntimeException("Interrupted while
evaluating readiness for tier " + tierOrder + " of " + group.getName(), e);
+ }
+ }
+
+ checkInstanceGroupMembersReady(group, tierMembers,
progressByVmId, membersReadyStatus);
+
+ if
(membersReadyStatus.values().stream().allMatch(Boolean::booleanValue)) {
+ return;
+ }
+
+ sleep(pollIntervalMs);
+ }
+ String reason = String.format("Tier %d of boot group '%s' did not
become ready within the maximum wait of %dms", tierOrder, group.getName(),
maxWaitMs);
+ logger.error(reason);
+ halt(group, reason);
+ throw new CloudRuntimeException(reason);
+ } finally {
+ readinessExecutor.shutdown();
+ }
+ }
+
+ /**
+ * Runs on one of {@code waitForTierReady}'s pooled threads for a single
VM: gates on the
+ * initial-delay window, dispatches this poll's check with the remaining
time budget, and treats
+ * Error the same as NotReady — both get a retry before anything halts.
+ *
+ * <p>Package-visible (rather than private) so tests can call it directly
instead of via
+ * reflection.</p>
+ */
+ protected void checkVmReadiness(InstanceBootGroupVO group, VmProgress
progress, Map<Long, InstanceBootGroupMemberVO> memberById,
+ Map<Long, Boolean> membersReadyStatus, long
effectiveMaxRetryAttempts, long effectiveTimeoutSeconds,
+ boolean effectiveRebootOnRetry) {
+ UserVmVO vm = userVmDao.findById(progress.vmId);
+ long elapsedMs = System.currentTimeMillis() - progress.enteredWaitAtMs;
+ long elapsedSinceBootMs = System.currentTimeMillis() -
progress.lastBootedAtMs;
+ long initialDelayMs = effectiveInitialDelaySeconds(group) * 1000L;
+ if (elapsedSinceBootMs < initialDelayMs) {
+ logger.debug("{} still within the initial delay window ({}ms
elapsed of {}ms since last boot) — skipping readiness check this poll. Attempt:
{}",
+ vm, elapsedSinceBootMs, initialDelayMs,
progress.getAttemptsLog(effectiveMaxRetryAttempts));
+ return;
+ }
+
+ long remainingMs = Math.max(0, effectiveTimeoutSeconds * 1000L -
elapsedMs);
+ String attemptLabel =
progress.getAttemptsLog(effectiveMaxRetryAttempts);
+ logger.debug("Evaluating readiness of {} for {} ({}ms since this
attempt started, {}ms remaining budget). Attempt: {}",
+ vm, group, elapsedMs, remainingMs, attemptLabel);
+ InstanceBootGroupReadinessRule.Status readiness =
instanceBootGroupReadinessRuleService.evaluateVmReadiness(group.getId(),
progress.vmId, remainingMs, attemptLabel);
+ if (readiness == InstanceBootGroupReadinessRule.Status.Ready) {
+ progress.ready = true;
+ logger.debug("{} is ready for {}", vm, group);
+ InstanceBootGroupMemberVO member =
memberById.get(progress.bootGroupMemberId);
+ if (member != null &&
InstanceBootGroupMember.MemberType.VirtualMachine.equals(member.getMemberType()))
{
+ membersReadyStatus.put(progress.bootGroupMemberId, true);
+ }
+ return;
+ }
+
+ if (progress.retryAttempts < effectiveMaxRetryAttempts - 1) {
+ long now = System.currentTimeMillis();
+ if (effectiveRebootOnRetry) {
+ rebootVm(progress.vmId);
+ progress.lastBootedAtMs = now;
+ }
+ logger.debug("{} readiness retry attempt {} of {} with a
reboot={}",
+ vm, progress.getAttemptsLog(effectiveMaxRetryAttempts),
group, effectiveRebootOnRetry);
+ progress.retryAttempts++;
+ progress.enteredWaitAtMs = now;
+ } else {
+ progress.gaveUp = true;
+ InstanceBootGroupMemberVO member =
memberById.get(progress.bootGroupMemberId);
+ if (member != null &&
InstanceBootGroupMember.MemberType.InstanceGroup.equals(member.getMemberType()))
{
+ logger.warn("{} failed readiness after {} retry attempts;
giving up on it and deferring to {}'s own readiness rule",
+ vm,
progress.getAttemptsLog(effectiveMaxRetryAttempts), member);
+ } else {
+ String reason = String.format("Instance '%s' failed readiness
after %s retry attempts",
+ vm.getName(),
progress.getAttemptsLog(effectiveMaxRetryAttempts));
+ logger.warn("{} failed readiness after {} retry attempts;
halting {}",
+ vm,
progress.getAttemptsLog(effectiveMaxRetryAttempts), group);
+ halt(group, reason);
+ throw new CloudRuntimeException(reason);
+ }
+ }
+ }
+
+ /**
+ * Sequential pass over the tier's InstanceGroup members, run once all of
this poll's per-VM
+ * tasks finish. An empty member list is treated as settled — {@code
allMatch()} on an empty
+ * stream is vacuously true either way, so there's nothing left to wait
for.
+ */
+ private void checkInstanceGroupMembersReady(InstanceBootGroupVO group,
List<InstanceBootGroupMemberVO> tierMembers,
+ Map<Long, VmProgress> progressByVmId, Map<Long, Boolean>
membersReadyStatus) {
+ for (InstanceBootGroupMemberVO member : tierMembers) {
+ if
(!InstanceBootGroupMember.MemberType.InstanceGroup.equals(member.getMemberType())
|| membersReadyStatus.getOrDefault(member.getId(), false)) {
+ continue;
+ }
+ InstanceGroup instanceGroup =
instanceGroupDao.findById(member.getMemberId());
+ Collection<VmProgress> memberProgresses =
progressByVmId.values().stream()
+ .filter(p -> member.getId() == (p.bootGroupMemberId ==
null ? -1 : p.bootGroupMemberId))
+ .collect(Collectors.toList());
+ if (!memberProgresses.isEmpty() &&
memberProgresses.stream().allMatch(p -> !p.ready && !p.gaveUp)) {
+ logger.debug("{} part of {} has no VMs that are ready or have
exhausted their retries yet", instanceGroup, group);
+ continue;
+ }
+
+ Set<Long> permanentlyFailedVmIds = memberProgresses.stream()
+ .filter(p -> p.gaveUp)
+ .map(p -> p.vmId)
+ .collect(Collectors.toSet());
+ InstanceBootGroupReadinessRule.Status groupStatus =
+
instanceBootGroupReadinessRuleService.evaluateInstanceGroupReadiness(group.getId(),
member.getMemberId(), permanentlyFailedVmIds);
+ if (groupStatus == InstanceBootGroupReadinessRule.Status.Ready) {
+ membersReadyStatus.put(member.getId(), true);
+ logger.info("{} part of {} reached readiness state Ready",
instanceGroup, group);
+ continue;
+ }
+ if
(InstanceBootGroupReadinessRule.Status.Error.equals(groupStatus) ||
+
(InstanceBootGroupReadinessRule.Status.NotReady.equals(groupStatus) &&
+ memberProgresses.stream().allMatch(p -> p.ready ||
p.gaveUp))) {
+ String reason = String.format("Instance group '%s' failed its
own readiness rules", instanceGroup.getName());
+ logger.error("{} failed its own readiness rules; halting {}",
instanceGroup, group);
+ halt(group, reason);
+ throw new CloudRuntimeException(reason);
+ }
+ }
+ }
+
+ /**
+ * Only stops the orchestration loop — never stops a VM, since every VM
touched by this point may
+ * already be running and tearing it down would be destructive, not
recoverable.
+ */
+ private void halt(InstanceBootGroupVO group, String reason) {
+ logger.warn("Halting {} start: {}", group, reason);
+ }
+
+ private void rebootVm(long vmId) {
+ UserVmVO vm = userVmDao.findById(vmId);
+ if (vm == null) {
+ logger.warn("Cannot reboot instance id {} for a boot group
readiness retry: VM not found", vmId);
+ return;
+ }
+ logger.debug("Rebooting {} for a boot group readiness retry attempt",
vm);
+ try {
+ virtualMachineManager.reboot(vm.getUuid(), null);
+ } catch (Exception e) {
+ throw new CloudRuntimeException("Failed to reboot VM " + vm + "
during boot group readiness retry: " + e.getMessage(), e);
+ }
+ }
+
+ private void sleep(long millis) {
+ try {
+ Thread.sleep(millis);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new CloudRuntimeException("Interrupted while waiting for
boot group tier readiness", e);
+ }
+ }
+
+ private long effectiveTimeoutSeconds(InstanceBootGroupVO group) {
+ String override = instanceBootGroupDetailsDao.getDetail(group.getId(),
ReadinessAttemptTimeoutSeconds.key());
+ return override != null ? Long.parseLong(override) :
ReadinessAttemptTimeoutSeconds.value();
+ }
+
+ private long effectiveMaxRetryAttempts(InstanceBootGroupVO group) {
+ String override = instanceBootGroupDetailsDao.getDetail(group.getId(),
ReadinessMaxRetryAttempts.key());
+ return override != null ? Long.parseLong(override) :
ReadinessMaxRetryAttempts.value();
+ }
+
+ private long effectivePollIntervalSeconds() {
+ return ReadinessPollIntervalSeconds.value();
+ }
+
+ private long effectiveReadinessCheckConcurrency() {
+ return ReadinessCheckConcurrency.value();
+ }
+
+ private long effectiveInitialDelaySeconds(InstanceBootGroupVO group) {
+ String override = instanceBootGroupDetailsDao.getDetail(group.getId(),
ReadinessInitialDelaySeconds.key());
+ return override != null ? Long.parseLong(override) :
ReadinessInitialDelaySeconds.value();
+ }
+
+ private boolean effectiveRebootOnRetry(InstanceBootGroupVO group) {
+ String override = instanceBootGroupDetailsDao.getDetail(group.getId(),
ReadinessRebootOnRetry.key());
+ return override != null ? Boolean.parseBoolean(override) :
ReadinessRebootOnRetry.value();
+ }
+
+ @Override
+ public void stopInstanceBootGroup(InstanceBootGroupVO group, boolean
forced) {
+ List<InstanceBootGroupMemberVO> members =
instanceBootGroupMemberDao.listByBootGroupId(group.getId());
+ Map<Integer, List<InstanceBootGroupMemberVO>> tiers =
groupByOrderDescending(members);
+ logger.info("Stopping {}: {} tier(s), {} member(s) total, forced={}",
group, tiers.size(), members.size(), forced);
+ long groupStoppedAtMs = System.currentTimeMillis();
+
+ for (Map.Entry<Integer, List<InstanceBootGroupMemberVO>> tier :
tiers.entrySet()) {
+ List<Long> vmIds = resolveVmIds(tier.getValue());
+ runTierConcurrently(vmIds, group, "stop", vmId -> {
+ UserVmVO vm = userVmDao.findById(vmId);
+ if (vm != null && vm.getState() !=
com.cloud.vm.VirtualMachine.State.Stopped) {
+ userVmService.stopVirtualMachine(vmId, forced);
+ }
+ });
Review Comment:
The stop command's description promises that stopping continues through all
tiers when an instance fails, but `runTierConcurrently` propagates the first VM
exception and exits this loop. A failure in a higher tier therefore prevents
all lower tiers from being stopped. Collect/log per-VM failures for a tier and
continue through the remaining tiers, then report the aggregate result.
##########
server/src/main/java/org/apache/cloudstack/vm/bootgroup/InstanceBootGroupManagerImpl.java:
##########
@@ -0,0 +1,604 @@
+// 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.vm.bootgroup;
+
+import java.util.ArrayList;
+import java.util.Collection;
+import java.util.Comparator;
+import java.util.HashMap;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.TreeMap;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.stream.Collectors;
+
+import javax.inject.Inject;
+import javax.naming.ConfigurationException;
+
+import org.apache.cloudstack.api.ApiCommandResourceType;
+import org.apache.cloudstack.context.CallContext;
+import org.apache.cloudstack.framework.config.ConfigKey;
+import org.apache.cloudstack.framework.config.Configurable;
+import org.apache.cloudstack.managed.context.ManagedContextRunnable;
+import
org.apache.cloudstack.vm.bootgroup.readiness.InstanceBootGroupReadinessRule;
+import
org.apache.cloudstack.vm.bootgroup.readiness.InstanceBootGroupReadinessRuleService;
+import org.springframework.stereotype.Component;
+
+import com.cloud.utils.component.ManagerBase;
+import com.cloud.utils.concurrency.NamedThreadFactory;
+import com.cloud.utils.exception.CloudRuntimeException;
+import com.cloud.vm.InstanceGroup;
+import com.cloud.vm.UserVmService;
+import com.cloud.vm.UserVmVO;
+import com.cloud.vm.VirtualMachine;
+import com.cloud.vm.VirtualMachineManager;
+import com.cloud.vm.dao.InstanceBootGroupDetailsDao;
+import com.cloud.vm.dao.InstanceBootGroupMemberDao;
+import com.cloud.vm.dao.InstanceGroupDao;
+import com.cloud.vm.dao.InstanceGroupVMMapDao;
+import com.cloud.vm.dao.UserVmDao;
+
+/**
+ * Backend/orchestration half of the Instance Boot Group feature — tier
concurrency, hypervisor
+ * start/stop/reboot calls, and readiness-gated tier progression. API-cmd
handling (ACL, validation,
+ * response building, command registration) lives in {@code
InstanceBootGroupApiServiceImpl}, which
+ * delegates here with resolved domain objects.
+ *
+ * <p>Per-VM timeout/reboot-attempt bookkeeping during a start is kept purely
in-memory, scoped to the
+ * async job thread executing the start — it is not persisted. Surviving a
management-server restart
+ * mid-run is explicitly not a goal here; if the process restarts, the job
(and this bookkeeping) is
+ * simply lost, same as any other in-flight async job. Current readiness is
queryable at any time via
+ * {@code listInstanceBootGroupMembers?details=readiness}, not via a separate
run-history API.</p>
+ */
+@Component
+public class InstanceBootGroupManagerImpl extends ManagerBase implements
InstanceBootGroupManager, Configurable {
+
+ public static final ConfigKey<Long> ReadinessAttemptTimeoutSeconds = new
ConfigKey<>("Advanced", Long.class,
+ "instance.boot.group.readiness.timeout.seconds", "300",
+ "How long to wait (in seconds) for an instance to become ready
during boot group orchestration before starting a new readiness retry attempt.
Overridable per boot group.", true);
+
+ public static final ConfigKey<Long> ReadinessMaxRetryAttempts = new
ConfigKey<>("Advanced", Long.class,
+ "instance.boot.group.readiness.max.retry.attempts", "5",
+ "Maximum number of readiness retry attempts for an instance that
fails to become ready during boot group orchestration before the boot group
start is halted. Overridable per boot group.", true);
+
+ public static final ConfigKey<Long> ReadinessPollIntervalSeconds = new
ConfigKey<>("Advanced", Long.class,
+ "instance.boot.group.readiness.poll.interval.seconds", "10",
+ "How often (in seconds) to re-check instance/instance-group
readiness during boot group orchestration, including the minimum pause after a
readiness retry attempt that did not reboot the instance before repeating the
check that just failed. A very low value can cause rapid repeated
(\"hammering\") readiness retries against an instance/VR/host. Global only, not
overridable per boot group.", true);
+
+ public static final ConfigKey<Long> ReadinessInitialDelaySeconds = new
ConfigKey<>("Advanced", Long.class,
+ "instance.boot.group.readiness.initial.delay.seconds", "30",
+ "How long to wait (in seconds) after starting or rebooting an
instance before its first readiness check of that attempt, giving the guest
OS/agent/network time to come up. Overridable per boot group.", true);
+
+ public static final ConfigKey<Boolean> ReadinessRebootOnRetry = new
ConfigKey<>("Advanced", Boolean.class,
+ "instance.boot.group.readiness.reboot.on.retry", "false",
+ "Whether to reboot an instance between readiness retry attempts
during boot group orchestration, instead of just waiting longer. Overridable
per boot group.", true);
+
+ public static final ConfigKey<Long> ReadinessCheckConcurrency = new
ConfigKey<>("Advanced", Long.class,
+ "instance.boot.group.readiness.check.concurrency", "10",
+ "Maximum number of instances within a single boot-order tier whose
readiness is checked concurrently during boot group orchestration, so one slow
check cannot delay every other instance's check in the same poll. Global only,
not overridable per boot group.", true);
+
+ public static final ConfigKey<Long> MaxMembersPerBootGroup = new
ConfigKey<>("Advanced", Long.class,
+ "instance.boot.group.max.members", "10",
+ "Maximum number of members that can be added to a single instance
boot group.", true, ConfigKey.Scope.Domain);
+
+ @Inject
+ private InstanceBootGroupMemberDao instanceBootGroupMemberDao;
+
+ @Inject
+ private UserVmService userVmService;
+
+ @Inject
+ private UserVmDao userVmDao;
+
+ @Inject
+ private InstanceGroupDao instanceGroupDao;
+
+ @Inject
+ private InstanceGroupVMMapDao instanceGroupVMMapDao;
+
+ @Inject
+ private VirtualMachineManager virtualMachineManager;
+
+ @Inject
+ private InstanceBootGroupReadinessRuleService
instanceBootGroupReadinessRuleService;
+
+ @Inject
+ private InstanceBootGroupDetailsDao instanceBootGroupDetailsDao;
+
+ @Override
+ public boolean configure(String name, Map<String, Object> params) throws
ConfigurationException {
+ VirtualMachine.State.getStateMachine().registerListener(new
InstanceBootGroupVmStateListener(instanceBootGroupReadinessRuleService));
+ return true;
+ }
+
+ @Override
+ public String getConfigComponentName() {
+ return InstanceBootGroupManagerImpl.class.getSimpleName();
+ }
+
+ @Override
+ public ConfigKey<?>[] getConfigKeys() {
+ return new ConfigKey<?>[]{ReadinessAttemptTimeoutSeconds,
ReadinessMaxRetryAttempts, ReadinessPollIntervalSeconds,
ReadinessInitialDelaySeconds, ReadinessRebootOnRetry, ReadinessCheckConcurrency,
+ MaxMembersPerBootGroup};
+ }
+
+ /** In-memory-only per-VM progress for a single start attempt — never
persisted. Package-visible
+ * (rather than private) purely so tests can construct one and inspect
retry/give-up state
+ * directly instead of via reflection. */
+ protected static final class VmProgress {
+ private final long vmId;
+ private final Long bootGroupMemberId;
+ private boolean ready;
+ /** Set when this VM has exhausted its retry attempts but belongs to
an InstanceGroup
+ * member, so the group's own readiness rule (e.g. a quorum rule)
gets the final say
+ * instead of this one VM halting the whole boot group. */
+ protected boolean gaveUp;
+ protected int retryAttempts;
+ /** Anchor for the per-attempt timeout window; reset on every retry,
rebooted or not. */
+ private long enteredWaitAtMs;
+ /** Anchor for the initial-delay grace period; only reset on the
initial start and on an
+ * actual reboot — a no-op retry (reboot-on-retry disabled) leaves
this alone, since there's
+ * no fresh boot to wait out. */
+ private long lastBootedAtMs;
+
+ protected VmProgress(long vmId, Long bootGroupMemberId) {
+ this.vmId = vmId;
+ this.bootGroupMemberId = bootGroupMemberId;
+ }
+
+ private String getAttemptsLog(long maxAttempts) {
+ return String.format("%d/%d", retryAttempts + 1, maxAttempts);
+ }
+ }
+
+ @Override
+ public void startInstanceBootGroup(InstanceBootGroupVO group) {
+ List<InstanceBootGroupMemberVO> members =
instanceBootGroupMemberDao.listByBootGroupId(group.getId());
+ Map<Integer, List<InstanceBootGroupMemberVO>> tiers =
groupByOrder(members);
+ logger.info("Starting {}: {} tier(s), {} member(s) total", group,
tiers.size(), members.size());
+ long groupStartedAtMs = System.currentTimeMillis();
+
+ for (Map.Entry<Integer, List<InstanceBootGroupMemberVO>> tierEntry :
tiers.entrySet()) {
+ int tierOrder = tierEntry.getKey();
+ List<InstanceBootGroupMemberVO> tierMembers = tierEntry.getValue();
+
+ Map<Long, VmProgress> progressByVmId = new LinkedHashMap<>();
+ for (InstanceBootGroupMemberVO member : tierMembers) {
+ for (Long vmId : resolveVmIds(List.of(member))) {
+ progressByVmId.put(vmId, new VmProgress(vmId,
member.getId()));
+ }
+ }
+ List<Long> tierVmIds = new ArrayList<>(progressByVmId.keySet());
+ logger.info("Starting tier {} of {}: {} member(s), {} VM(s)",
tierOrder, group, tierMembers.size(), tierVmIds.size());
+ long tierStartedAtMs = System.currentTimeMillis();
+
+ try {
+ runTierConcurrently(tierVmIds, group, "start", vmId -> {
+ UserVmVO vm = userVmDao.findById(vmId);
+ boolean alreadyRunning = vm != null &&
VirtualMachine.State.Running.equals(vm.getState());
+ if (vm != null && !alreadyRunning) {
+ userVmService.startVirtualMachine(vm, null);
+ }
+ anchorInitialDelay(group, progressByVmId.get(vmId), vm,
alreadyRunning);
+ });
+ } catch (CloudRuntimeException e) {
+ halt(group, "Failed to start a VM in tier " + tierOrder + ": "
+ e.getMessage());
+ throw e;
+ }
+
+ waitForTierReady(group, tierOrder, tierMembers, progressByVmId);
+ logger.info("Tier {} of {} is ready ({}ms)", tierOrder, group,
System.currentTimeMillis() - tierStartedAtMs);
+ }
+
+ logger.info("{} start completed ({}ms)", group,
System.currentTimeMillis() - groupStartedAtMs);
+ }
+
+ /**
+ * If the VM was already running, anchors the initial-delay grace period
to when CloudStack last
+ * confirmed its power state rather than to "now" — so it waits out only
what's left of the
+ * delay (or none) instead of a full fresh wait it doesn't need.
+ */
+ private void anchorInitialDelay(InstanceBootGroupVO group, VmProgress
progress, UserVmVO vm, boolean alreadyRunning) {
+ long now = System.currentTimeMillis();
+ progress.enteredWaitAtMs = now;
+ if (alreadyRunning && vm.getPowerStateUpdateTime() != null) {
+ progress.lastBootedAtMs = vm.getPowerStateUpdateTime().getTime();
+ logger.debug("{} was already running (power state last confirmed
{}); readiness checks begin after any remaining portion of the {}s initial
delay",
+ vm, vm.getPowerStateUpdateTime(),
effectiveInitialDelaySeconds(group));
+ } else {
+ progress.lastBootedAtMs = now;
+ logger.debug("{} start action completed; readiness checks begin
after the {}s initial delay", vm, effectiveInitialDelaySeconds(group));
+ }
+ }
+
+ /**
+ * Polls every not-yet-settled VM in the tier concurrently (bounded by
+ * {@code ReadinessCheckConcurrency}) until the whole tier — VMs and any
InstanceGroup
+ * members — reports ready, or a halt is triggered.
+ */
+ private void waitForTierReady(InstanceBootGroupVO group, int tierOrder,
List<InstanceBootGroupMemberVO> tierMembers, Map<Long, VmProgress>
progressByVmId) {
+ Map<Long, InstanceBootGroupMemberVO> memberById = new HashMap<>();
+ Map<Long, Boolean> membersReadyStatus = new ConcurrentHashMap<>();
+ for (InstanceBootGroupMemberVO member : tierMembers) {
+ memberById.put(member.getId(), member);
+ membersReadyStatus.put(member.getId(), false);
+ }
+ final long effectiveMaxRetryAttempts =
effectiveMaxRetryAttempts(group);
+ final long effectiveTimeoutSeconds = effectiveTimeoutSeconds(group);
+ final long effectivePollIntervalSeconds =
effectivePollIntervalSeconds();
+ final long pollIntervalMs = effectivePollIntervalSeconds * 1000L;
+ final boolean effectiveRebootOnRetry = effectiveRebootOnRetry(group);
+ int concurrency = (int) Math.max(1, Math.min(progressByVmId.size(),
effectiveReadinessCheckConcurrency()));
+ // Bound the polling loop: initial delay + maxAttempts full timeout
windows + inter-poll sleeps between them.
+ // At least one timeout window is always budgeted so a single check is
never starved out, even if
+ // maxAttempts is configured to 0.
+ long maxWaitMs = (effectiveInitialDelaySeconds(group) + Math.max(1,
effectiveMaxRetryAttempts) * effectiveTimeoutSeconds
+ + Math.max(0, effectiveMaxRetryAttempts - 1) *
effectivePollIntervalSeconds) * 1000L;
+ long deadline = System.currentTimeMillis() + maxWaitMs;
+ logger.debug("Waiting for tier {} of {} to become ready: {} VM(s)
tracked, timeout={}s, pollInterval={}s, checkConcurrency={}, maxWait={}ms",
+ tierOrder, group, progressByVmId.size(),
effectiveTimeoutSeconds, effectivePollIntervalSeconds, concurrency, maxWaitMs);
+
+ CallContext callerContext = CallContext.current();
+ ExecutorService readinessExecutor =
Executors.newFixedThreadPool(concurrency, new
NamedThreadFactory("InstanceBootGroup-readiness-" + tierOrder));
+ try {
+ while (System.currentTimeMillis() < deadline) {
+ List<Future<Void>> futures = new ArrayList<>();
+ for (VmProgress progress : progressByVmId.values()) {
+ if (progress.ready || progress.gaveUp) {
+ continue;
+ }
+ futures.add(readinessExecutor.submit(() -> {
+ CallContext.register(callerContext,
ApiCommandResourceType.VirtualMachine);
+ try {
+ checkVmReadiness(group, progress, memberById,
membersReadyStatus,
+ effectiveMaxRetryAttempts,
effectiveTimeoutSeconds, effectiveRebootOnRetry);
+ } finally {
+ CallContext.unregister();
+ }
+ return null;
+ }));
+ }
+ for (Future<Void> future : futures) {
+ try {
+ future.get();
+ } catch (ExecutionException e) {
+ Throwable cause = e.getCause() != null ? e.getCause()
: e;
+ if (cause instanceof CloudRuntimeException) {
+ throw (CloudRuntimeException) cause;
+ }
+ throw new CloudRuntimeException("Failed to evaluate
readiness for a VM in tier " + tierOrder + " of " + group.getName() + ": " +
cause.getMessage(), cause);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new CloudRuntimeException("Interrupted while
evaluating readiness for tier " + tierOrder + " of " + group.getName(), e);
+ }
+ }
+
+ checkInstanceGroupMembersReady(group, tierMembers,
progressByVmId, membersReadyStatus);
+
+ if
(membersReadyStatus.values().stream().allMatch(Boolean::booleanValue)) {
+ return;
+ }
+
+ sleep(pollIntervalMs);
+ }
+ String reason = String.format("Tier %d of boot group '%s' did not
become ready within the maximum wait of %dms", tierOrder, group.getName(),
maxWaitMs);
+ logger.error(reason);
+ halt(group, reason);
+ throw new CloudRuntimeException(reason);
+ } finally {
+ readinessExecutor.shutdown();
+ }
+ }
+
+ /**
+ * Runs on one of {@code waitForTierReady}'s pooled threads for a single
VM: gates on the
+ * initial-delay window, dispatches this poll's check with the remaining
time budget, and treats
+ * Error the same as NotReady — both get a retry before anything halts.
+ *
+ * <p>Package-visible (rather than private) so tests can call it directly
instead of via
+ * reflection.</p>
+ */
+ protected void checkVmReadiness(InstanceBootGroupVO group, VmProgress
progress, Map<Long, InstanceBootGroupMemberVO> memberById,
+ Map<Long, Boolean> membersReadyStatus, long
effectiveMaxRetryAttempts, long effectiveTimeoutSeconds,
+ boolean effectiveRebootOnRetry) {
+ UserVmVO vm = userVmDao.findById(progress.vmId);
+ long elapsedMs = System.currentTimeMillis() - progress.enteredWaitAtMs;
+ long elapsedSinceBootMs = System.currentTimeMillis() -
progress.lastBootedAtMs;
+ long initialDelayMs = effectiveInitialDelaySeconds(group) * 1000L;
+ if (elapsedSinceBootMs < initialDelayMs) {
+ logger.debug("{} still within the initial delay window ({}ms
elapsed of {}ms since last boot) — skipping readiness check this poll. Attempt:
{}",
+ vm, elapsedSinceBootMs, initialDelayMs,
progress.getAttemptsLog(effectiveMaxRetryAttempts));
+ return;
+ }
+
+ long remainingMs = Math.max(0, effectiveTimeoutSeconds * 1000L -
elapsedMs);
+ String attemptLabel =
progress.getAttemptsLog(effectiveMaxRetryAttempts);
+ logger.debug("Evaluating readiness of {} for {} ({}ms since this
attempt started, {}ms remaining budget). Attempt: {}",
+ vm, group, elapsedMs, remainingMs, attemptLabel);
+ InstanceBootGroupReadinessRule.Status readiness =
instanceBootGroupReadinessRuleService.evaluateVmReadiness(group.getId(),
progress.vmId, remainingMs, attemptLabel);
+ if (readiness == InstanceBootGroupReadinessRule.Status.Ready) {
+ progress.ready = true;
+ logger.debug("{} is ready for {}", vm, group);
+ InstanceBootGroupMemberVO member =
memberById.get(progress.bootGroupMemberId);
+ if (member != null &&
InstanceBootGroupMember.MemberType.VirtualMachine.equals(member.getMemberType()))
{
+ membersReadyStatus.put(progress.bootGroupMemberId, true);
+ }
+ return;
+ }
+
+ if (progress.retryAttempts < effectiveMaxRetryAttempts - 1) {
+ long now = System.currentTimeMillis();
+ if (effectiveRebootOnRetry) {
+ rebootVm(progress.vmId);
+ progress.lastBootedAtMs = now;
+ }
+ logger.debug("{} readiness retry attempt {} of {} with a
reboot={}",
+ vm, progress.getAttemptsLog(effectiveMaxRetryAttempts),
group, effectiveRebootOnRetry);
+ progress.retryAttempts++;
+ progress.enteredWaitAtMs = now;
+ } else {
+ progress.gaveUp = true;
+ InstanceBootGroupMemberVO member =
memberById.get(progress.bootGroupMemberId);
+ if (member != null &&
InstanceBootGroupMember.MemberType.InstanceGroup.equals(member.getMemberType()))
{
+ logger.warn("{} failed readiness after {} retry attempts;
giving up on it and deferring to {}'s own readiness rule",
+ vm,
progress.getAttemptsLog(effectiveMaxRetryAttempts), member);
+ } else {
+ String reason = String.format("Instance '%s' failed readiness
after %s retry attempts",
+ vm.getName(),
progress.getAttemptsLog(effectiveMaxRetryAttempts));
+ logger.warn("{} failed readiness after {} retry attempts;
halting {}",
+ vm,
progress.getAttemptsLog(effectiveMaxRetryAttempts), group);
+ halt(group, reason);
+ throw new CloudRuntimeException(reason);
+ }
+ }
+ }
+
+ /**
+ * Sequential pass over the tier's InstanceGroup members, run once all of
this poll's per-VM
+ * tasks finish. An empty member list is treated as settled — {@code
allMatch()} on an empty
+ * stream is vacuously true either way, so there's nothing left to wait
for.
+ */
+ private void checkInstanceGroupMembersReady(InstanceBootGroupVO group,
List<InstanceBootGroupMemberVO> tierMembers,
+ Map<Long, VmProgress> progressByVmId, Map<Long, Boolean>
membersReadyStatus) {
+ for (InstanceBootGroupMemberVO member : tierMembers) {
+ if
(!InstanceBootGroupMember.MemberType.InstanceGroup.equals(member.getMemberType())
|| membersReadyStatus.getOrDefault(member.getId(), false)) {
+ continue;
+ }
+ InstanceGroup instanceGroup =
instanceGroupDao.findById(member.getMemberId());
+ Collection<VmProgress> memberProgresses =
progressByVmId.values().stream()
+ .filter(p -> member.getId() == (p.bootGroupMemberId ==
null ? -1 : p.bootGroupMemberId))
+ .collect(Collectors.toList());
+ if (!memberProgresses.isEmpty() &&
memberProgresses.stream().allMatch(p -> !p.ready && !p.gaveUp)) {
+ logger.debug("{} part of {} has no VMs that are ready or have
exhausted their retries yet", instanceGroup, group);
+ continue;
+ }
+
+ Set<Long> permanentlyFailedVmIds = memberProgresses.stream()
+ .filter(p -> p.gaveUp)
+ .map(p -> p.vmId)
+ .collect(Collectors.toSet());
+ InstanceBootGroupReadinessRule.Status groupStatus =
+
instanceBootGroupReadinessRuleService.evaluateInstanceGroupReadiness(group.getId(),
member.getMemberId(), permanentlyFailedVmIds);
+ if (groupStatus == InstanceBootGroupReadinessRule.Status.Ready) {
+ membersReadyStatus.put(member.getId(), true);
+ logger.info("{} part of {} reached readiness state Ready",
instanceGroup, group);
+ continue;
+ }
+ if
(InstanceBootGroupReadinessRule.Status.Error.equals(groupStatus) ||
+
(InstanceBootGroupReadinessRule.Status.NotReady.equals(groupStatus) &&
+ memberProgresses.stream().allMatch(p -> p.ready ||
p.gaveUp))) {
+ String reason = String.format("Instance group '%s' failed its
own readiness rules", instanceGroup.getName());
+ logger.error("{} failed its own readiness rules; halting {}",
instanceGroup, group);
+ halt(group, reason);
+ throw new CloudRuntimeException(reason);
+ }
+ }
+ }
+
+ /**
+ * Only stops the orchestration loop — never stops a VM, since every VM
touched by this point may
+ * already be running and tearing it down would be destructive, not
recoverable.
+ */
+ private void halt(InstanceBootGroupVO group, String reason) {
+ logger.warn("Halting {} start: {}", group, reason);
+ }
+
+ private void rebootVm(long vmId) {
+ UserVmVO vm = userVmDao.findById(vmId);
+ if (vm == null) {
+ logger.warn("Cannot reboot instance id {} for a boot group
readiness retry: VM not found", vmId);
+ return;
+ }
+ logger.debug("Rebooting {} for a boot group readiness retry attempt",
vm);
+ try {
+ virtualMachineManager.reboot(vm.getUuid(), null);
+ } catch (Exception e) {
+ throw new CloudRuntimeException("Failed to reboot VM " + vm + "
during boot group readiness retry: " + e.getMessage(), e);
+ }
+ }
+
+ private void sleep(long millis) {
+ try {
+ Thread.sleep(millis);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new CloudRuntimeException("Interrupted while waiting for
boot group tier readiness", e);
+ }
+ }
+
+ private long effectiveTimeoutSeconds(InstanceBootGroupVO group) {
+ String override = instanceBootGroupDetailsDao.getDetail(group.getId(),
ReadinessAttemptTimeoutSeconds.key());
+ return override != null ? Long.parseLong(override) :
ReadinessAttemptTimeoutSeconds.value();
+ }
+
+ private long effectiveMaxRetryAttempts(InstanceBootGroupVO group) {
+ String override = instanceBootGroupDetailsDao.getDetail(group.getId(),
ReadinessMaxRetryAttempts.key());
+ return override != null ? Long.parseLong(override) :
ReadinessMaxRetryAttempts.value();
+ }
+
+ private long effectivePollIntervalSeconds() {
+ return ReadinessPollIntervalSeconds.value();
+ }
+
+ private long effectiveReadinessCheckConcurrency() {
+ return ReadinessCheckConcurrency.value();
+ }
+
+ private long effectiveInitialDelaySeconds(InstanceBootGroupVO group) {
+ String override = instanceBootGroupDetailsDao.getDetail(group.getId(),
ReadinessInitialDelaySeconds.key());
+ return override != null ? Long.parseLong(override) :
ReadinessInitialDelaySeconds.value();
+ }
+
+ private boolean effectiveRebootOnRetry(InstanceBootGroupVO group) {
+ String override = instanceBootGroupDetailsDao.getDetail(group.getId(),
ReadinessRebootOnRetry.key());
+ return override != null ? Boolean.parseBoolean(override) :
ReadinessRebootOnRetry.value();
+ }
+
+ @Override
+ public void stopInstanceBootGroup(InstanceBootGroupVO group, boolean
forced) {
+ List<InstanceBootGroupMemberVO> members =
instanceBootGroupMemberDao.listByBootGroupId(group.getId());
+ Map<Integer, List<InstanceBootGroupMemberVO>> tiers =
groupByOrderDescending(members);
+ logger.info("Stopping {}: {} tier(s), {} member(s) total, forced={}",
group, tiers.size(), members.size(), forced);
+ long groupStoppedAtMs = System.currentTimeMillis();
+
+ for (Map.Entry<Integer, List<InstanceBootGroupMemberVO>> tier :
tiers.entrySet()) {
+ List<Long> vmIds = resolveVmIds(tier.getValue());
+ runTierConcurrently(vmIds, group, "stop", vmId -> {
+ UserVmVO vm = userVmDao.findById(vmId);
+ if (vm != null && vm.getState() !=
com.cloud.vm.VirtualMachine.State.Stopped) {
+ userVmService.stopVirtualMachine(vmId, forced);
+ }
+ });
+ }
+
+ logger.info("{} stop completed ({}ms)", group,
System.currentTimeMillis() - groupStoppedAtMs);
+ }
+
+ @Override
+ public void rebootInstanceBootGroup(InstanceBootGroupVO group, boolean
forced) {
+ logger.info("Rebooting {}: stopping, then starting", group);
+ stopInstanceBootGroup(group, forced);
+ startInstanceBootGroup(group);
+ }
+
+ /**
+ * Runs {@code action} for every VM in a tier concurrently and aborts on
the first failure. Each
+ * thread gets a copied {@link CallContext} — without one, a VM lifecycle
action routed through
+ * the job-queue path fails to submit its sub-job ("no lock found").
+ */
+ private void runTierConcurrently(List<Long> vmIds, InstanceBootGroupVO
group,
+ String actionName, VmAction action) {
+ if (vmIds.isEmpty()) {
+ return;
+ }
+
+ logger.debug("Running '{}' action for a tier of {}: {} VM id(s) {}",
+ actionName, group, vmIds.size(), vmIds);
+ long actionStartedAtMs = System.currentTimeMillis();
+ CallContext callerContext = CallContext.current();
+ int threadCount = Math.min(vmIds.size(),
ReadinessCheckConcurrency.value().intValue());
+ ExecutorService executor = Executors.newFixedThreadPool(
+ threadCount, new NamedThreadFactory("InstanceBootGroup-" +
actionName));
Review Comment:
When `ReadinessCheckConcurrency` is configured as 0 or a negative value,
this computes a non-positive pool size and `Executors.newFixedThreadPool`
throws before any start/stop action runs. The ConfigKey has no lower-bound
validation, so clamp the pool size to at least one (as the readiness polling
path already does).
##########
server/src/main/java/org/apache/cloudstack/vm/bootgroup/readiness/InstanceBootGroupReadinessRuleManagerImpl.java:
##########
@@ -0,0 +1,618 @@
+// 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.vm.bootgroup.readiness;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.Date;
+import java.util.EnumSet;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.stream.Collectors;
+
+import javax.inject.Inject;
+
+import org.apache.cloudstack.vm.bootgroup.InstanceBootGroupMember;
+import org.apache.cloudstack.vm.bootgroup.InstanceBootGroupMemberVO;
+import
org.apache.cloudstack.vm.bootgroup.InstanceBootGroupReadinessCheckResultVO;
+import org.apache.cloudstack.vm.bootgroup.InstanceBootGroupReadinessRuleVO;
+import org.apache.commons.collections.MapUtils;
+import org.apache.commons.lang3.StringUtils;
+import org.springframework.stereotype.Component;
+
+import com.cloud.exception.InvalidParameterValueException;
+import com.cloud.hypervisor.Hypervisor.HypervisorType;
+import com.cloud.utils.component.ManagerBase;
+import com.cloud.utils.db.Transaction;
+import com.cloud.utils.db.TransactionCallback;
+import com.cloud.vm.InstanceGroupVMMapVO;
+import com.cloud.vm.UserVmVO;
+import com.cloud.vm.VirtualMachine;
+import com.cloud.vm.dao.InstanceBootGroupMemberDao;
+import com.cloud.vm.dao.InstanceBootGroupReadinessCheckResultDao;
+import com.cloud.vm.dao.InstanceBootGroupReadinessRuleDao;
+import com.cloud.vm.dao.InstanceBootGroupReadinessRuleDetailsDao;
+import com.cloud.vm.dao.InstanceGroupDao;
+import com.cloud.vm.dao.InstanceGroupVMMapDao;
+import com.cloud.vm.dao.UserVmDao;
+
+@Component
+public class InstanceBootGroupReadinessRuleManagerImpl extends ManagerBase
implements InstanceBootGroupReadinessRuleService {
+
+ private static final String THRESHOLD_TYPE_KEY = "threshold_type";
+ private static final String THRESHOLD_VALUE_KEY = "threshold_value";
+ private static final String PORT_KEY = "port";
+ private static final String PROTOCOL_KEY = "protocol";
+
+ private static final Map<InstanceBootGroupMember.MemberType,
Set<InstanceBootGroupReadinessRule.RuleType>> VALID_RULE_TYPES_BY_ITEM_TYPE =
Map.of(
+ InstanceBootGroupMember.MemberType.VirtualMachine, EnumSet.of(
+ InstanceBootGroupReadinessRule.RuleType.GuestAgentLiveness,
+ InstanceBootGroupReadinessRule.RuleType.Ping,
+ InstanceBootGroupReadinessRule.RuleType.PortCheck),
+ InstanceBootGroupMember.MemberType.InstanceGroup, EnumSet.of(
+ InstanceBootGroupReadinessRule.RuleType.GuestAgentLiveness,
+ InstanceBootGroupReadinessRule.RuleType.Ping,
+ InstanceBootGroupReadinessRule.RuleType.PortCheck,
+ InstanceBootGroupReadinessRule.RuleType.MemberQuorum));
+
+ /**
+ * Rule types an item may have at most one of — Ping/GuestAgentLiveness
each check a single fixed
+ * target on the VM, and MemberQuorum aggregates the whole InstanceGroup,
so a second one would
+ * just be redundant. PortCheck (different ports) is not a singleton.
+ */
+ private static final Set<InstanceBootGroupReadinessRule.RuleType>
SINGLETON_RULE_TYPES = EnumSet.of(
+ InstanceBootGroupReadinessRule.RuleType.Ping,
+ InstanceBootGroupReadinessRule.RuleType.GuestAgentLiveness,
+ InstanceBootGroupReadinessRule.RuleType.MemberQuorum);
+
+ @Inject
+ private InstanceBootGroupReadinessRuleDao
instanceBootGroupReadinessRuleDao;
+
+ @Inject
+ private InstanceBootGroupReadinessRuleDetailsDao
instanceBootGroupReadinessRuleDetailsDao;
+
+ @Inject
+ private InstanceBootGroupReadinessCheckResultDao
instanceBootGroupReadinessCheckResultDao;
+
+ @Inject
+ private InstanceBootGroupMemberDao instanceBootGroupMemberDao;
+
+ @Inject
+ private InstanceGroupVMMapDao instanceGroupVMMapDao;
+
+ @Inject
+ private InstanceGroupDao instanceGroupDao;
+
+ @Inject
+ private UserVmDao userVmDao;
+
+ private List<ReadinessChecker> readinessCheckers;
+
+ private Map<InstanceBootGroupReadinessRule.RuleType, ReadinessChecker>
checkersByRuleType;
+
+ protected void updateCheckersByRuleType(boolean forced) {
+ if (MapUtils.isNotEmpty(checkersByRuleType) && !forced) {
+ return;
+ }
+ checkersByRuleType = new HashMap<>();
+ for (ReadinessChecker checker : readinessCheckers) {
+ checkersByRuleType.put(checker.getRuleType(), checker);
+ }
+ }
+
+ protected ReadinessChecker
getCheckerByRuleType(InstanceBootGroupReadinessRule.RuleType ruleType) {
+ updateCheckersByRuleType(false);
+ return checkersByRuleType.get(ruleType);
+ }
+
+ public List<ReadinessChecker> getReadinessCheckers() {
+ return readinessCheckers;
+ }
+
+ public void setReadinessCheckers(List<ReadinessChecker> readinessCheckers)
{
+ this.readinessCheckers = readinessCheckers;
+ updateCheckersByRuleType(true);
+ }
+
+ @Override
+ public InstanceBootGroupReadinessRule createReadinessRule(long
bootGroupId, InstanceBootGroupMember.MemberType itemType, long itemId,
+ InstanceBootGroupReadinessRule.RuleType ruleType, String name,
boolean enabled, Map<String, String> details) {
+ validateRuleTypeForItemType(itemType, ruleType);
+ validateItemBelongsToBootGroup(bootGroupId, itemType, itemId);
+ validateSingletonRuleType(bootGroupId, itemType, itemId, ruleType);
+ validateGuestAgentLivenessSupported(itemType, itemId, ruleType);
+ validateRuleTypeSpecificDetails(ruleType, details);
+
+ String effectiveName = StringUtils.isNotBlank(name) ? name :
String.format("%s-%s-%d", ruleType.name(), itemType.name(), itemId);
+ InstanceBootGroupReadinessRuleVO rule = new
InstanceBootGroupReadinessRuleVO(effectiveName, bootGroupId, itemType, itemId,
ruleType, enabled);
+ rule = instanceBootGroupReadinessRuleDao.persist(rule);
+
+ if (details != null) {
+ for (Map.Entry<String, String> entry : details.entrySet()) {
+
instanceBootGroupReadinessRuleDetailsDao.addDetail(rule.getId(),
entry.getKey(), entry.getValue(), true);
+ }
+ }
+ return rule;
+ }
+
+ @Override
+ public InstanceBootGroupReadinessRule updateReadinessRule(long ruleId,
String name, Boolean enabled, Map<String, String> details) {
+ InstanceBootGroupReadinessRuleVO rule =
instanceBootGroupReadinessRuleDao.findById(ruleId);
+ if (rule == null) {
+ throw new InvalidParameterValueException("Unable to find a
readiness rule with ID: " + ruleId);
+ }
+
+ if (StringUtils.isNotBlank(name)) {
+ rule.setName(name);
+ }
+ if (enabled != null) {
+ rule.setEnabled(enabled);
+ }
+ if (details != null) {
+ Map<String, String> mergedDetails = new
HashMap<>(instanceBootGroupReadinessRuleDetailsDao.getDetails(rule.getId()));
+ mergedDetails.putAll(details);
+ validateRuleTypeSpecificDetails(rule.getRuleType(), mergedDetails);
+ }
+ instanceBootGroupReadinessRuleDao.update(rule.getId(), rule);
+
+ if (details != null) {
+ for (Map.Entry<String, String> entry : details.entrySet()) {
+
instanceBootGroupReadinessRuleDetailsDao.addDetail(rule.getId(),
entry.getKey(), entry.getValue(), true);
+ }
+ }
+ return instanceBootGroupReadinessRuleDao.findById(rule.getId());
+ }
+
+ @Override
+ public boolean deleteReadinessRule(long ruleId) {
+ InstanceBootGroupReadinessRuleVO rule =
instanceBootGroupReadinessRuleDao.findById(ruleId);
+ if (rule == null) {
+ throw new InvalidParameterValueException("Unable to find a
readiness rule with ID: " + ruleId);
+ }
+ return Transaction.execute((TransactionCallback<Boolean>) status -> {
+
instanceBootGroupReadinessRuleDetailsDao.removeDetails(rule.getId());
+
instanceBootGroupReadinessCheckResultDao.deleteByRuleId(rule.getId());
+ instanceBootGroupReadinessRuleDao.remove(rule.getId());
+ return true;
+ });
+ }
+
+ @Override
+ public InstanceBootGroupReadinessRule findById(long ruleId) {
+ return instanceBootGroupReadinessRuleDao.findById(ruleId);
+ }
+
+ @Override
+ public Map<String, String> getRuleDetails(long ruleId) {
+ return instanceBootGroupReadinessRuleDetailsDao.getDetails(ruleId);
+ }
+
+ @Override
+ public InstanceBootGroupReadinessRule.Status evaluateVmReadiness(long
bootGroupId, long vmId, long remainingMs, String attemptLabel) {
+ return resolveVmReadiness(bootGroupId, vmId, true, remainingMs,
attemptLabel);
+ }
+
+ /**
+ * Same aggregation as {@link #evaluateVmReadiness}, but reads each rule's
last-cached result
+ * instead of dispatching a fresh check — for callers that run after the
per-VM loop has already
+ * dispatched this poll, so a live re-check would just repeat the same
remote command.
+ */
+ private InstanceBootGroupReadinessRule.Status getCachedVmReadiness(long
bootGroupId, long vmId) {
+ return resolveVmReadiness(bootGroupId, vmId, false, Long.MAX_VALUE,
null);
+ }
+
+ /**
+ * Shared by the dispatching and cache-reading paths. A non-Running VM is
always NotReady regardless
+ * of any cached rule result. Direct rules cache at vmId 0, inherited
group rules per member vmId;
+ * when dispatching, the remaining budget shrinks across a VM's rules so N
slow rules can't each burn
+ * the full per-attempt timeout.
+ */
+ private InstanceBootGroupReadinessRule.Status resolveVmReadiness(long
bootGroupId, long vmId, boolean dispatch, long remainingMs, String
attemptLabel) {
+ UserVmVO vm = userVmDao.findById(vmId);
+ if (vm == null) {
+ logger.debug("VM id {} not found while evaluating readiness for
boot group id {}; treating as NotReady", vmId, bootGroupId);
+ return InstanceBootGroupReadinessRule.Status.NotReady;
+ }
+ if (!VirtualMachine.State.Running.equals(vm.getState())) {
+ logger.debug("{} is not Running (state={}) for boot group id {};
treating as NotReady without evaluating or dispatching its readiness rules",
vm, vm.getState(), bootGroupId);
+ return InstanceBootGroupReadinessRule.Status.NotReady;
+ }
+ List<InstanceBootGroupReadinessRuleVO> directRules =
instanceBootGroupReadinessRuleDao.listEnabledByItem(bootGroupId,
InstanceBootGroupMember.MemberType.VirtualMachine, vmId);
+ List<InstanceBootGroupReadinessRuleVO> inheritedRules =
findInheritedGroupRuleVOs(bootGroupId, vmId);
+
+ if (directRules.isEmpty() && inheritedRules.isEmpty()) {
+ logger.debug("{} has no readiness rules for boot group id {};
derived readiness Ready from its current state ({})", vm, bootGroupId,
vm.getState());
+ return InstanceBootGroupReadinessRule.Status.Ready;
+ }
+
+ if (dispatch) {
+ logger.debug("Evaluating readiness of {} for boot group id {}: {}
direct rule(s), {} inherited rule(s), {}ms remaining budget",
+ vm, bootGroupId, directRules.size(),
inheritedRules.size(), remainingMs);
+ }
+
+ boolean anyError = false;
+ boolean anyNotReady = false;
+ long remainingBudgetMs = remainingMs;
+ for (InstanceBootGroupReadinessRuleVO rule : directRules) {
+ long startedAtMs = System.currentTimeMillis();
+ InstanceBootGroupReadinessRule.Status status = dispatch ?
evaluateAndCacheRule(rule, vmId, 0, remainingBudgetMs, attemptLabel) :
readCachedRuleStatus(rule.getId(), 0);
+ if (dispatch) {
+ remainingBudgetMs = Math.max(0, remainingBudgetMs -
(System.currentTimeMillis() - startedAtMs));
+ }
+ if (status == InstanceBootGroupReadinessRule.Status.Error) {
+ anyError = true;
+ } else if (status != InstanceBootGroupReadinessRule.Status.Ready) {
+ anyNotReady = true;
+ }
+ }
+ for (InstanceBootGroupReadinessRuleVO rule : inheritedRules) {
+ long startedAtMs = System.currentTimeMillis();
+ InstanceBootGroupReadinessRule.Status status = dispatch ?
evaluateAndCacheRule(rule, vmId, vmId, remainingBudgetMs, attemptLabel) :
readCachedRuleStatus(rule.getId(), vmId);
+ if (dispatch) {
+ remainingBudgetMs = Math.max(0, remainingBudgetMs -
(System.currentTimeMillis() - startedAtMs));
+ }
+ if (status == InstanceBootGroupReadinessRule.Status.Error) {
+ anyError = true;
+ } else if (status != InstanceBootGroupReadinessRule.Status.Ready) {
+ anyNotReady = true;
+ }
+ }
+
+ InstanceBootGroupReadinessRule.Status overallStatus = anyError ?
InstanceBootGroupReadinessRule.Status.Error
+ : (anyNotReady ?
InstanceBootGroupReadinessRule.Status.NotReady :
InstanceBootGroupReadinessRule.Status.Ready);
+ if (dispatch) {
+ logger.debug("{} readiness for boot group id {} evaluated as {}",
vm, bootGroupId, overallStatus);
+ }
+ return overallStatus;
+ }
+
+ private InstanceBootGroupReadinessRule.Status readCachedRuleStatus(long
ruleId, long cacheVmId) {
+ InstanceBootGroupReadinessCheckResultVO cached =
instanceBootGroupReadinessCheckResultDao.findByRuleAndVm(ruleId, cacheVmId);
+ return (cached != null && cached.getStatus() != null) ?
cached.getStatus() : InstanceBootGroupReadinessRule.Status.Unknown;
+ }
+
+ /**
+ * Dispatches (or, for a missing checker, synthesizes) one rule's check
and persists the result —
+ * the single place a check actually happens, so every caller shares this
one log/cache path.
+ * @param attemptLabel e.g. {@code "2/5"}, appended to the persisted
message; pass {@code null} to skip.
+ */
+ private InstanceBootGroupReadinessRule.Status
evaluateAndCacheRule(InstanceBootGroupReadinessRuleVO rule, long vmId, long
cacheVmId, long remainingMs, String attemptLabel) {
+ logger.debug("Evaluating rule {} against VM id {} with {}ms remaining
budget", () -> rule, () -> userVmDao.findById(vmId), () -> remainingMs);
+ ReadinessChecker checker = getCheckerByRuleType(rule.getRuleType());
+ InstanceBootGroupReadinessRule.Status status;
+ String message;
+ if (checker == null) {
+ status = InstanceBootGroupReadinessRule.Status.Error;
+ message = "No checker implemented yet for rule type " +
rule.getRuleType();
+ } else {
+ Map<String, String> details =
instanceBootGroupReadinessRuleDetailsDao.getDetails(rule.getId());
+ ReadinessChecker.Result result = checker.check(rule, details,
vmId, remainingMs);
+ status = result.getStatus();
+ message = result.getMessage();
+ }
+ String finalMessage = StringUtils.isNotBlank(attemptLabel) ? message +
" (attempt " + attemptLabel + ")" : message;
+ logger.debug("Rule {} evaluated against VM id {}: status={},
message={}", () -> rule, () -> userVmDao.findById(vmId), () -> status, () ->
finalMessage);
+ instanceBootGroupReadinessCheckResultDao.upsert(rule.getId(),
cacheVmId, status, finalMessage, new Date());
+ return status;
+ }
+
+ @Override
+ public List<InstanceBootGroupReadinessRule> findInheritedGroupRules(long
bootGroupId, long vmId) {
+ return new ArrayList<>(findInheritedGroupRuleVOs(bootGroupId, vmId));
+ }
+
+ /**
+ * Resets a VM's own cached rule results (direct and inherited) to Unknown
when it starts, so a
+ * VM restarted outside boot group orchestration can't keep reporting a
stale Ready from before
+ * it stopped. A no-op for a VM with no boot-group involvement at all.
+ */
+ @Override
+ public void invalidateCachedReadinessOnRestart(long vmId) {
+ InstanceBootGroupMemberVO directMember =
instanceBootGroupMemberDao.findByMember(InstanceBootGroupMember.MemberType.VirtualMachine,
vmId);
+ if (directMember != null) {
+
invalidateRuleResults(instanceBootGroupReadinessRuleDao.listEnabledByItem(directMember.getBootGroupId(),
InstanceBootGroupMember.MemberType.VirtualMachine, vmId), 0L);
+ }
+ for (InstanceGroupVMMapVO mapping :
instanceGroupVMMapDao.listByInstanceId(vmId)) {
+ InstanceBootGroupMemberVO groupMember =
instanceBootGroupMemberDao.findByMember(InstanceBootGroupMember.MemberType.InstanceGroup,
mapping.getGroupId());
+ if (groupMember == null) {
+ continue;
+ }
+ List<InstanceBootGroupReadinessRuleVO> memberTargetedRules =
instanceBootGroupReadinessRuleDao.listEnabledByItem(groupMember.getBootGroupId(),
+ InstanceBootGroupMember.MemberType.InstanceGroup,
mapping.getGroupId()).stream()
+ .filter(rule -> rule.getRuleType().isMemberTargeted())
+ .collect(Collectors.toList());
+ invalidateRuleResults(memberTargetedRules, vmId);
+ }
+ }
+
+ private void invalidateRuleResults(List<InstanceBootGroupReadinessRuleVO>
rules, long cacheVmId) {
+ for (InstanceBootGroupReadinessRuleVO rule : rules) {
+ instanceBootGroupReadinessCheckResultDao.upsert(rule.getId(),
cacheVmId, InstanceBootGroupReadinessRule.Status.Unknown,
+ "Instance (re)started; not yet re-verified this session",
new Date());
+ }
+ }
+
+ /**
+ * A VM inherits its owning InstanceGroup's
Ping/PortCheck/GuestAgentLiveness rules (not
+ * MemberQuorum/CustomScript, which only ever make sense at group scope) —
resolved by finding
+ * the InstanceGroup, among any this VM belongs to, that is itself a
member of this boot group.
+ */
+ private List<InstanceBootGroupReadinessRuleVO>
findInheritedGroupRuleVOs(long bootGroupId, long vmId) {
+ for (InstanceGroupVMMapVO mapping :
instanceGroupVMMapDao.listByInstanceId(vmId)) {
+ InstanceBootGroupMemberVO groupMember =
instanceBootGroupMemberDao.findByMember(InstanceBootGroupMember.MemberType.InstanceGroup,
mapping.getGroupId());
+ if (groupMember != null && groupMember.getBootGroupId() ==
bootGroupId) {
+ return
instanceBootGroupReadinessRuleDao.listEnabledByItem(bootGroupId,
InstanceBootGroupMember.MemberType.InstanceGroup, mapping.getGroupId()).stream()
+ .filter(rule -> rule.getRuleType().isMemberTargeted())
+ .collect(Collectors.toList());
+ }
+ }
+ return Collections.emptyList();
+ }
+
+ /**
+ * AND of the group's own rules and every member's cached readiness
(read-only — the per-VM loop
+ * already dispatched this poll). A MemberQuorum rule's own
tolerance-aware verdict decides
+ * Ready/Error on its own; without one, a member still mid-retry only
counts as NotReady, never Error.
+ * Once a MemberQuorum rule governs the group, every other member-targeted
rule's own all-members
+ * aggregate becomes informational only — it's still evaluated and shown,
but no longer gates the
+ * overall verdict, since that's exactly what attaching a quorum rule is
meant to relax.
+ */
+ @Override
+ public InstanceBootGroupReadinessRule.Status
evaluateInstanceGroupReadiness(long bootGroupId, long instanceGroupId,
Set<Long> permanentlyFailedVmIds) {
+ List<InstanceBootGroupReadinessRuleVO> groupRules =
instanceBootGroupReadinessRuleDao.listEnabledByItem(bootGroupId,
InstanceBootGroupMember.MemberType.InstanceGroup, instanceGroupId);
+ List<InstanceGroupVMMapVO> members =
instanceGroupVMMapDao.listByGroupId(instanceGroupId);
+ boolean hasMemberQuorumRule = groupRules.stream().anyMatch(rule ->
rule.getRuleType() == InstanceBootGroupReadinessRule.RuleType.MemberQuorum);
+
+ logger.debug("Evaluating readiness of instance group id {} for boot
group id {}: {} member VM(s), {} own rule(s), quorum-governed={}",
+ () -> instanceGroupDao.findById(instanceGroupId), () ->
bootGroupId, () -> members.size(), () -> groupRules.size(), () ->
hasMemberQuorumRule);
+
+ boolean anyError = false;
+ boolean anyNotReady = false;
+
+ for (InstanceGroupVMMapVO member : members) {
+ InstanceBootGroupReadinessRule.Status vmStatus =
getCachedVmReadiness(bootGroupId, member.getInstanceId());
+ if (hasMemberQuorumRule || vmStatus ==
InstanceBootGroupReadinessRule.Status.Ready) {
+ continue;
+ }
+ if (permanentlyFailedVmIds.contains(member.getInstanceId())) {
+ anyError = true;
+ } else {
+ anyNotReady = true;
+ }
+ }
+
+ for (InstanceBootGroupReadinessRuleVO rule : groupRules) {
+ ReadinessChecker.Result result;
+ boolean memberTargeted = rule.getRuleType().isMemberTargeted();
+ if (rule.getRuleType() ==
InstanceBootGroupReadinessRule.RuleType.MemberQuorum) {
+ logger.debug("Evaluating group-scoped rule {} for instance
group id {} via member quorum", rule, instanceGroupId);
+ Map<String, String> details =
instanceBootGroupReadinessRuleDetailsDao.getDetails(rule.getId());
+ result = evaluateInstanceQuorum(bootGroupId, instanceGroupId,
details, permanentlyFailedVmIds);
+ } else if (memberTargeted) {
+ logger.debug("Evaluating group-scoped rule {} for instance
group id {} by aggregating its {} member(s)' own cached results", rule,
instanceGroupId, members.size());
+ result = aggregateMemberTargetedGroupRule(rule, members);
+ } else {
+ result = new
ReadinessChecker.Result(InstanceBootGroupReadinessRule.Status.Error,
+ "No evaluator implemented yet for rule type " +
rule.getRuleType());
+ }
+ logger.debug("Group-scoped rule {} evaluated for instance group id
{}: status={}, message={}", rule, instanceGroupId, result.getStatus(),
result.getMessage());
+ instanceBootGroupReadinessCheckResultDao.upsert(rule.getId(), 0,
result.getStatus(), result.getMessage(), new Date());
+
+ if (memberTargeted && hasMemberQuorumRule) {
+ continue;
+ }
+ if (result.getStatus() ==
InstanceBootGroupReadinessRule.Status.Error) {
+ anyError = true;
+ } else if (result.getStatus() !=
InstanceBootGroupReadinessRule.Status.Ready) {
+ anyNotReady = true;
+ }
+ }
+
+ InstanceBootGroupReadinessRule.Status overallStatus = anyError ?
InstanceBootGroupReadinessRule.Status.Error
+ : (anyNotReady ?
InstanceBootGroupReadinessRule.Status.NotReady :
InstanceBootGroupReadinessRule.Status.Ready);
+ logger.debug("Instance group id {} readiness for boot group id {}
evaluated as {}", instanceGroupId, bootGroupId, overallStatus);
+ return overallStatus;
+ }
+
+ /**
+ * Reads the per-member cached results {@link #evaluateVmReadiness}
already wrote for this rule,
+ * rather than dispatching it again.
+ */
+ private ReadinessChecker.Result
aggregateMemberTargetedGroupRule(InstanceBootGroupReadinessRuleVO rule,
List<InstanceGroupVMMapVO> members) {
+ if (members.isEmpty()) {
+ return new
ReadinessChecker.Result(InstanceBootGroupReadinessRule.Status.NotReady,
"Instance group has no members");
+ }
+ boolean anyError = false;
+ int readyCount = 0;
+ for (InstanceGroupVMMapVO member : members) {
+ InstanceBootGroupReadinessCheckResultVO cached =
instanceBootGroupReadinessCheckResultDao.findByRuleAndVm(rule.getId(),
member.getInstanceId());
+ InstanceBootGroupReadinessRule.Status status = (cached != null &&
cached.getStatus() != null) ? cached.getStatus() :
InstanceBootGroupReadinessRule.Status.Unknown;
+ if (status == InstanceBootGroupReadinessRule.Status.Ready) {
+ readyCount++;
+ } else if (status == InstanceBootGroupReadinessRule.Status.Error) {
+ anyError = true;
+ }
+ }
+ int total = members.size();
+ String message = String.format("%d of %d member(s) ready via %s",
readyCount, total, rule.getRuleType().name());
+ if (anyError) {
+ return new
ReadinessChecker.Result(InstanceBootGroupReadinessRule.Status.Error, message);
+ }
+ return new ReadinessChecker.Result(readyCount == total ?
InstanceBootGroupReadinessRule.Status.Ready :
InstanceBootGroupReadinessRule.Status.NotReady, message);
+ }
+
+ /**
+ * Pure computation, no dispatch: counts members currently READY against
the configured threshold.
+ * @param permanentlyFailedVmIds excluded from "achievable" so a hopeless
quorum reports Error
+ * instead of NotReady once it can never be met, even with every
remaining member succeeding.
+ */
+ private ReadinessChecker.Result evaluateInstanceQuorum(long bootGroupId,
long instanceGroupId, Map<String, String> details, Set<Long>
permanentlyFailedVmIds) {
+ List<InstanceGroupVMMapVO> members =
instanceGroupVMMapDao.listByGroupId(instanceGroupId);
+ int total = members.size();
+ if (total == 0) {
+ return new
ReadinessChecker.Result(InstanceBootGroupReadinessRule.Status.NotReady,
"Instance group has no members");
+ }
+
+ long readyCount = members.stream()
+ .filter(member -> getCachedVmReadiness(bootGroupId,
member.getInstanceId()) == InstanceBootGroupReadinessRule.Status.Ready)
+ .count();
+ long permanentlyFailedCount = members.stream()
+ .filter(member ->
permanentlyFailedVmIds.contains(member.getInstanceId()))
+ .count();
+ long achievableCount = total - permanentlyFailedCount;
+
+ String thresholdType = details == null ? null :
details.get(THRESHOLD_TYPE_KEY);
+ String thresholdValue = details == null ? null :
details.get(THRESHOLD_VALUE_KEY);
+
+ boolean met;
+ boolean achievable;
+ try {
+ if ("PERCENTAGE".equalsIgnoreCase(thresholdType)) {
+ double thresholdPercentage =
Double.parseDouble(thresholdValue);
+ met = (readyCount * 100.0 / total) >= thresholdPercentage;
+ achievable = (achievableCount * 100.0 / total) >=
thresholdPercentage;
+ } else {
+ long thresholdCount = Long.parseLong(thresholdValue);
+ met = readyCount >= thresholdCount;
+ achievable = achievableCount >= thresholdCount;
+ }
+ } catch (NumberFormatException e) {
+ return new
ReadinessChecker.Result(InstanceBootGroupReadinessRule.Status.Error, "Invalid
threshold configuration: " + thresholdType + "=" + thresholdValue);
+ }
+
+ String message = String.format("%d/%d members ready (%s threshold
%s)", readyCount, total, thresholdType, thresholdValue);
+ if (!met && !achievable) {
+ String reason = String.format("%s; unreachable — %d/%d member(s)
have permanently failed readiness", message, permanentlyFailedCount, total);
+ logger.debug("Instance group id {} quorum check: {}",
instanceGroupId, reason);
+ return new
ReadinessChecker.Result(InstanceBootGroupReadinessRule.Status.Error, reason);
+ }
+ logger.debug("Instance group id {} quorum check: {}, met={}",
instanceGroupId, message, met);
+ return new ReadinessChecker.Result(met ?
InstanceBootGroupReadinessRule.Status.Ready :
InstanceBootGroupReadinessRule.Status.NotReady, message);
+ }
+
+ private void
validateRuleTypeSpecificDetails(InstanceBootGroupReadinessRule.RuleType
ruleType, Map<String, String> details) {
+ if (ruleType == InstanceBootGroupReadinessRule.RuleType.MemberQuorum) {
+ validateInstanceQuorumDetails(details);
+ } else if (ruleType ==
InstanceBootGroupReadinessRule.RuleType.PortCheck) {
+ validatePortCheckDetails(details);
+ }
+ }
+
+ private void validateInstanceQuorumDetails(Map<String, String> details) {
+ String thresholdType = details == null ? null :
details.get(THRESHOLD_TYPE_KEY);
+ String thresholdValue = details == null ? null :
details.get(THRESHOLD_VALUE_KEY);
+ if (StringUtils.isBlank(thresholdType) ||
StringUtils.isBlank(thresholdValue)) {
+ throw new InvalidParameterValueException(String.format("%s rules
require '%s' (COUNT or PERCENTAGE) and '%s' details",
InstanceBootGroupReadinessRule.RuleType.MemberQuorum.name(),
THRESHOLD_TYPE_KEY, THRESHOLD_VALUE_KEY));
+ }
+ if (!"COUNT".equalsIgnoreCase(thresholdType) &&
!"PERCENTAGE".equalsIgnoreCase(thresholdType)) {
+ throw new InvalidParameterValueException(THRESHOLD_TYPE_KEY + "
must be COUNT or PERCENTAGE");
+ }
+ try {
+ if ("PERCENTAGE".equalsIgnoreCase(thresholdType)) {
+ Double.parseDouble(thresholdValue);
+ } else {
+ Long.parseLong(thresholdValue);
+ }
Review Comment:
The threshold parser accepts negative values. A `COUNT=-1` or
`PERCENTAGE=-10` quorum is then immediately met even when no member is ready,
because the evaluation compares the ready count/percentage directly to that
negative threshold. Reject negative threshold values during validation.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]