From 21dec133cb48f572ba0332d869985c45204e1c34 Mon Sep 17 00:00:00 2001 From: Joseph Callen Date: Fri, 11 Sep 2026 16:33:37 -0400 Subject: [PATCH 01/11] vsphere: raise controller cache resync from 10m to 30m Periodic re-verification of healthy machines is the dominant vCenter API load source. Transient states already requeue at 30s and watches are event-driven, so a 30m resync cuts sustained load ~3x without changing provisioning or failure-reaction latency. --- cmd/vsphere/main.go | 6 +++--- cmd/vsphere/main_test.go | 15 +++++++++++++++ 2 files changed, 18 insertions(+), 3 deletions(-) create mode 100644 cmd/vsphere/main_test.go diff --git a/cmd/vsphere/main.go b/cmd/vsphere/main.go index 27d602810..3209beebc 100644 --- a/cmd/vsphere/main.go +++ b/cmd/vsphere/main.go @@ -34,7 +34,7 @@ import ( "github.com/openshift/machine-api-operator/pkg/version" ) -const timeout = 10 * time.Minute +const syncPeriod = 30 * time.Minute func main() { var printVersion bool @@ -123,7 +123,7 @@ func main() { } cfg := config.GetConfigOrDie() - syncPeriod := timeout + syncPeriodRef := syncPeriod le := util.GetLeaderElectionConfig(cfg, configv1.LeaderElection{ Disable: !*leaderElect, @@ -136,7 +136,7 @@ func main() { }, HealthProbeBindAddress: *healthAddr, Cache: cache.Options{ - SyncPeriod: &syncPeriod, + SyncPeriod: &syncPeriodRef, }, LeaderElection: *leaderElect, LeaderElectionNamespace: *leaderElectResourceNamespace, diff --git a/cmd/vsphere/main_test.go b/cmd/vsphere/main_test.go new file mode 100644 index 000000000..fc04173d1 --- /dev/null +++ b/cmd/vsphere/main_test.go @@ -0,0 +1,15 @@ +package main + +import ( + "testing" + "time" +) + +// Resync must stay well above the node-drain window so periodic +// reconciliation does not saturate the shared vCenter. Transient +// states are covered by the 30s requeue-until-running, not by resync. +func TestSyncPeriodFloor(t *testing.T) { + if syncPeriod < 10*time.Minute { + t.Fatalf("syncPeriod %s is below the 10m floor; do not lower it without re-analyzing vCenter API load", syncPeriod) + } +} From 263252aee3f089b693753acea1e8292e74ff30f0 Mon Sep 17 00:00:00 2001 From: Joseph Callen Date: Fri, 11 Sep 2026 16:36:21 -0400 Subject: [PATCH 02/11] vsphere: add --max-concurrent-reconciles (default 10) Parallelize machine reconciles so a saturated node drain (N sequential 30-min drains) finishes in ~3 minutes instead of ~30. Combined with the reduced per-reconcile vCenter calls and 30m resync, the burst is a small fraction of the previous sustained load. --- cmd/vsphere/main.go | 13 ++++++++++++- cmd/vsphere/main_test.go | 22 ++++++++++++++++++++++ 2 files changed, 34 insertions(+), 1 deletion(-) diff --git a/cmd/vsphere/main.go b/cmd/vsphere/main.go index 3209beebc..32fea83ba 100644 --- a/cmd/vsphere/main.go +++ b/cmd/vsphere/main.go @@ -36,6 +36,14 @@ import ( const syncPeriod = 30 * time.Minute +// registerControllerFlags registers machine controller tuning flags on fs. +func registerControllerFlags(fs *flag.FlagSet) *int { + return fs.Int("max-concurrent-reconciles", 10, + "Maximum number of parallel Machine reconciles. Higher values drain a "+ + "cluster faster but issue the same vCenter calls faster; keep 10 for "+ + "shared vCenter environments.") +} + func main() { var printVersion bool flag.BoolVar(&printVersion, "version", false, "print version and exit") @@ -99,6 +107,8 @@ func main() { "The address for health checking.", ) + maxConcurrentReconciles := registerControllerFlags(flag.CommandLine) + majorVersion := version.Version.Major if majorVersion == 0 { @@ -203,7 +213,8 @@ func main() { klog.Fatalf("unable to add ipamv1beta1 to scheme: %v", err) } - if err := capimachine.AddWithActuator(mgr, machineActuator, defaultMutableGate); err != nil { + if err := capimachine.AddWithActuatorOpts(mgr, machineActuator, + controller.Options{MaxConcurrentReconciles: *maxConcurrentReconciles}, defaultMutableGate); err != nil { klog.Fatal(err) } diff --git a/cmd/vsphere/main_test.go b/cmd/vsphere/main_test.go index fc04173d1..31f1dd313 100644 --- a/cmd/vsphere/main_test.go +++ b/cmd/vsphere/main_test.go @@ -1,6 +1,7 @@ package main import ( + "flag" "testing" "time" ) @@ -13,3 +14,24 @@ func TestSyncPeriodFloor(t *testing.T) { t.Fatalf("syncPeriod %s is below the 10m floor; do not lower it without re-analyzing vCenter API load", syncPeriod) } } + +func TestMaxConcurrentReconcilesDefault(t *testing.T) { + // The flag is registered in main(); register it in a test flagset + // by calling the helper that wires flags (extracted below). + fs := flag.NewFlagSet("test", flag.ContinueOnError) + maxConcurrent := registerControllerFlags(fs) // see Step 3 + if *maxConcurrent != 10 { + t.Errorf("default max-concurrent-reconciles = %d, want 10", *maxConcurrent) + } +} + +func TestMaxConcurrentReconcilesCustom(t *testing.T) { + fs := flag.NewFlagSet("test", flag.ContinueOnError) + maxConcurrent := registerControllerFlags(fs) + if err := fs.Parse([]string{"--max-concurrent-reconciles=5"}); err != nil { + t.Fatalf("unexpected error parsing flags: %v", err) + } + if *maxConcurrent != 5 { + t.Errorf("expected max-concurrent-reconciles = 5, got %d", *maxConcurrent) + } +} From 2fead2708ab182d0dbe8f3a9ab5adde625ec33d4 Mon Sep 17 00:00:00 2001 From: Joseph Callen Date: Fri, 11 Sep 2026 16:44:31 -0400 Subject: [PATCH 03/11] vsphere: assert 30m sync period in main test --- cmd/vsphere/main_test.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/cmd/vsphere/main_test.go b/cmd/vsphere/main_test.go index 31f1dd313..218327f55 100644 --- a/cmd/vsphere/main_test.go +++ b/cmd/vsphere/main_test.go @@ -10,8 +10,8 @@ import ( // reconciliation does not saturate the shared vCenter. Transient // states are covered by the 30s requeue-until-running, not by resync. func TestSyncPeriodFloor(t *testing.T) { - if syncPeriod < 10*time.Minute { - t.Fatalf("syncPeriod %s is below the 10m floor; do not lower it without re-analyzing vCenter API load", syncPeriod) + if syncPeriod != 30*time.Minute { + t.Fatalf("syncPeriod %s is not 30m; expected 30m to reduce vCenter API load", syncPeriod) } } From 21001ab1db4829b9a2148c642d682acf78a8bd87 Mon Sep 17 00:00:00 2001 From: Joseph Callen Date: Fri, 11 Sep 2026 16:48:17 -0400 Subject: [PATCH 04/11] vsphere: clear finished task ref on update to skip GetTask --- pkg/controller/vsphere/reconciler.go | 17 +++- pkg/controller/vsphere/reconciler_test.go | 112 ++++++++++++++++++++++ 2 files changed, 128 insertions(+), 1 deletion(-) diff --git a/pkg/controller/vsphere/reconciler.go b/pkg/controller/vsphere/reconciler.go index 7af805f43..737e73b78 100644 --- a/pkg/controller/vsphere/reconciler.go +++ b/pkg/controller/vsphere/reconciler.go @@ -272,6 +272,11 @@ func (r *Reconciler) update() error { }) return err } + // Task history eviction or a session restart can make a task + // ref permanently unavailable. Clear it so future resyncs do + // not keep issuing the same GetTask request. + klog.Infof("%v: task %s no longer found, clearing TaskRef", r.machine.GetName(), r.providerStatus.TaskRef) + r.providerStatus.TaskRef = "" } if moTask != nil { if taskIsFinished, err := taskIsFinished(moTask); err != nil { @@ -283,6 +288,11 @@ func (r *Reconciler) update() error { return fmt.Errorf("%v task %v finished with error: %w", moTask.Info.DescriptionId, moTask.Reference().Value, err) } else if !taskIsFinished { return fmt.Errorf("%v task %v has not finished", moTask.Info.DescriptionId, moTask.Reference().Value) + } else { + // A completed task can never transition again. Clear its ref + // so steady-state resyncs skip GetTask entirely. + klog.Infof("%v: task %v has completed, clearing TaskRef", r.machine.GetName(), moTask.Reference().Value) + r.providerStatus.TaskRef = "" } } } @@ -856,7 +866,12 @@ func constructKargsFromNetworkConfig(s *machineScope) (string, error) { } func isRetrieveMONotFound(taskRef string, err error) bool { - return err.Error() == fmt.Sprintf("ServerFaultCode: The object 'vim.Task:%v' has already been deleted or has not been completely created", taskRef) + if err == nil { + return false + } + errMessage := err.Error() + return errMessage == fmt.Sprintf("ServerFaultCode: The object 'vim.Task:%v' has already been deleted or has not been completely created", taskRef) || + errMessage == "ServerFaultCode: The object has already been deleted or has not been completely created" } func getHwVersion(ctx context.Context, vm *object.VirtualMachine) (int, error) { diff --git a/pkg/controller/vsphere/reconciler_test.go b/pkg/controller/vsphere/reconciler_test.go index 3fc67d993..9eeb5eb38 100644 --- a/pkg/controller/vsphere/reconciler_test.go +++ b/pkg/controller/vsphere/reconciler_test.go @@ -3585,3 +3585,115 @@ func TestReconcilePowerStateAnnontation(t *testing.T) { } // See https://github.com/vmware/govmomi/blob/master/simulator/example_extend_test.go#L33:6 for extending behaviour example + +func TestUpdateClearsFinishedTaskRef(t *testing.T) { + model, sess, server := initSimulator(t) + defer model.Remove() + defer server.Close() + + host, port, err := net.SplitHostPort(server.URL.Host) + if err != nil { + t.Fatal(err) + } + password, _ := server.URL.User.Password() + namespace := "test" + credentialsSecretName := "test" + credentialsSecret := &corev1.Secret{ + ObjectMeta: metav1.ObjectMeta{ + Name: credentialsSecretName, + Namespace: namespace, + }, + Data: map[string][]byte{ + fmt.Sprintf("%s.username", host): []byte(server.URL.User.Username()), + fmt.Sprintf("%s.password", host): []byte(password), + }, + } + configMap := &corev1.ConfigMap{ + ObjectMeta: metav1.ObjectMeta{ + Name: OpenshiftConfigManagedConfigMap, + Namespace: openshiftConfigNamespaceForTest, + }, + Data: map[string]string{ + OpenshiftConfigManagedCloudConfigKey: fmt.Sprintf(testConfigFmt, port, credentialsSecretName, namespace), + }, + } + if _, err := createTagAndCategory(sess, tagToCategoryName("CLUSTERID"), "CLUSTERID"); err != nil { + t.Fatalf("cannot create tag and category: %v", err) + } + + vm := model.Map().Any("VirtualMachine").(*simulator.VirtualMachine) + vm.Config.InstanceUuid = "a5764857-ae35-34dc-8f25-a9c9e73aa898" + vmObj := object.NewVirtualMachine(sess.Client.Client, vm.Reference()) + powerOffTask, err := vmObj.PowerOff(context.Background()) + if err != nil { + t.Fatal(err) + } + if err := object.NewTask(sess.Client.Client, powerOffTask.Reference()).Wait(context.Background()); err != nil { + t.Fatal(err) + } + task, err := vmObj.PowerOn(context.Background()) + if err != nil { + t.Fatal(err) + } + if err := object.NewTask(sess.Client.Client, task.Reference()).Wait(context.Background()); err != nil { + t.Fatal(err) + } + + rawProviderSpec, err := RawExtensionFromProviderSpec(&machinev1.VSphereMachineProviderSpec{ + Workspace: &machinev1.Workspace{Server: host}, + CredentialsSecret: &corev1.LocalObjectReference{ + Name: credentialsSecretName, + }, + Template: vm.Name, + Network: machinev1.NetworkSpec{ + Devices: []machinev1.NetworkDeviceSpec{{NetworkName: "test"}}, + }, + }) + if err != nil { + t.Fatal(err) + } + + for _, tc := range []struct { + name string + taskRef string + }{ + {name: "finished task", taskRef: task.Reference().Value}, + {name: "stale missing task", taskRef: "task-99999"}, + } { + t.Run(tc.name, func(t *testing.T) { + machineObj := &machinev1.Machine{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test-" + strings.ReplaceAll(tc.name, " ", "-"), + Namespace: namespace, + Labels: map[string]string{ + machinev1.MachineClusterIDLabel: "CLUSTERID", + }, + UID: apimachinerytypes.UID(vm.Config.InstanceUuid), + }, + Spec: machinev1.MachineSpec{ + ProviderSpec: machinev1.ProviderSpec{Value: rawProviderSpec}, + }, + } + client := fake.NewClientBuilder().WithScheme(scheme.Scheme).WithRuntimeObjects( + credentialsSecret, configMap).Build() + scope, err := newMachineScope(machineScopeParams{ + client: client, + Context: context.Background(), + machine: machineObj, + apiReader: client, + openshiftConfigNameSpace: openshiftConfigNamespaceForTest, + }) + if err != nil { + t.Fatal(err) + } + scope.providerStatus.TaskRef = tc.taskRef + + if err := newReconciler(scope).update(); err != nil { + t.Fatalf("update() error: %v", err) + } + if scope.providerStatus.TaskRef != "" { + t.Errorf("TaskRef not cleared after finished/stale task, got %q", scope.providerStatus.TaskRef) + } + }) + } +} From 8daa9d2bbc4ae03f7dd768024775868f6f44482c Mon Sep 17 00:00:00 2001 From: Joseph Callen Date: Fri, 11 Sep 2026 16:54:04 -0400 Subject: [PATCH 05/11] vsphere: skip region/zone tag walk when labels already set Once both labels are on the Machine the ancestry tag walk (SOAP HostSystem + property-collector Ancestors + REST ListAttachedTags per ancestor) yields the same immutable result; skip it on resync. --- pkg/controller/vsphere/reconciler.go | 8 ++++++ pkg/controller/vsphere/reconciler_test.go | 30 +++++++++++++++++++++++ 2 files changed, 38 insertions(+) diff --git a/pkg/controller/vsphere/reconciler.go b/pkg/controller/vsphere/reconciler.go index 737e73b78..e11722b67 100644 --- a/pkg/controller/vsphere/reconciler.go +++ b/pkg/controller/vsphere/reconciler.go @@ -601,6 +601,14 @@ func (r *Reconciler) reconcileRegionAndZoneLabels(vm *virtualMachine) error { return nil } + // Region/zone come from tags on the VM's ancestry and are immutable + // after provisioning. If the labels are already set, skip the tag + // traversal (HostSystem + Ancestors + N REST tag calls per resync). + if r.machine.Labels[machinecontroller.MachineRegionLabelName] != "" && + r.machine.Labels[machinecontroller.MachineAZLabelName] != "" { + return nil + } + regionLabel := r.vSphereConfig.Labels.Region zoneLabel := r.vSphereConfig.Labels.Zone diff --git a/pkg/controller/vsphere/reconciler_test.go b/pkg/controller/vsphere/reconciler_test.go index 9eeb5eb38..9d7b2feaa 100644 --- a/pkg/controller/vsphere/reconciler_test.go +++ b/pkg/controller/vsphere/reconciler_test.go @@ -3697,3 +3697,33 @@ func TestUpdateClearsFinishedTaskRef(t *testing.T) { }) } } + +func TestReconcileRegionAndZoneLabelsSkipsWhenSet(t *testing.T) { + machine := &machinev1.Machine{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test", + Labels: map[string]string{ + machinecontroller.MachineRegionLabelName: "east", + machinecontroller.MachineAZLabelName: "a", + }, + }, + } + r := &Reconciler{ + machineScope: &machineScope{ + machine: machine, + providerStatus: &machinev1.VSphereMachineProviderStatus{}, + vSphereConfig: &vsphere.Config{ + Labels: vsphere.Labels{Region: "region", Zone: "zone"}, + }, + }, + } + // No session: if the function touches the session it panics; the + // guard must return before any vCenter call. + if err := r.reconcileRegionAndZoneLabels(nil); err != nil { + t.Fatalf("expected nil, got %v", err) + } + if machine.Labels[machinecontroller.MachineRegionLabelName] != "east" || + machine.Labels[machinecontroller.MachineAZLabelName] != "a" { + t.Errorf("labels were modified: %v", machine.Labels) + } +} From 931729d04962724b458692669aa1ee58f786ea1d Mon Sep 17 00:00:00 2001 From: Joseph Callen Date: Fri, 11 Sep 2026 17:00:56 -0400 Subject: [PATCH 06/11] vsphere: dedup power-state/name/UUID fetches on steady-state resync --- pkg/controller/vsphere/reconciler.go | 50 +++++++++++++---------- pkg/controller/vsphere/reconciler_test.go | 41 ++++++++++++++++++- 2 files changed, 68 insertions(+), 23 deletions(-) diff --git a/pkg/controller/vsphere/reconciler.go b/pkg/controller/vsphere/reconciler.go index e11722b67..21dc61716 100644 --- a/pkg/controller/vsphere/reconciler.go +++ b/pkg/controller/vsphere/reconciler.go @@ -631,6 +631,10 @@ func (r *Reconciler) reconcileRegionAndZoneLabels(vm *virtualMachine) error { } func (r *Reconciler) reconcileProviderID(vm *virtualMachine) error { + if r.machine.Spec.ProviderID != nil && *r.machine.Spec.ProviderID != "" { + return nil + } + providerID, err := convertUUIDToProviderID(vm.Obj.UUID(vm.Context)) if err != nil { return err @@ -649,7 +653,7 @@ func convertUUIDToProviderID(UUID string) (string, error) { } func (r *Reconciler) reconcileNetwork(vm *virtualMachine) error { - currentNetworkStatusList, err := vm.getNetworkStatusList(r.session.Client.Client) + currentNetworkStatusList, vmName, err := vm.getNetworkStatusList(r.session.Client.Client) if err != nil { return fmt.Errorf("error getting network status: %v", err) } @@ -671,15 +675,6 @@ func (r *Reconciler) reconcileNetwork(vm *virtualMachine) error { } } - // Using Name() if InventoryPath is empty will return empty name - // see: https://github.com/vmware/govmomi/blob/master/object/common.go#L66-L75 - // Using ObjectName() as it will query from VirtualMachine properties - - vmName, err := vm.Obj.ObjectName(vm.Context) - if err != nil { - return fmt.Errorf("error getting virtual machine name: %v", err) - } - ipAddrs = append(ipAddrs, corev1.NodeAddress{ Type: corev1.NodeInternalDNS, Address: vmName, @@ -1461,8 +1456,10 @@ func setProviderStatus(taskRef string, condition metav1.Condition, scope *machin klog.Infof("%s: Updating provider status", scope.machine.Name) if vm != nil { - id := vm.Obj.UUID(scope.Context) - scope.providerStatus.InstanceID = &id + if scope.providerStatus.InstanceID == nil || *scope.providerStatus.InstanceID == "" { + id := vm.Obj.UUID(scope.Context) + scope.providerStatus.InstanceID = &id + } // This can return an error if machine is being deleted powerState, err := vm.getPowerState() @@ -1501,6 +1498,12 @@ type virtualMachine struct { context.Context Ref types.ManagedObjectReference Obj *object.VirtualMachine + + // powerState cache: one PowerState call per reconcile pass. + // The struct is built fresh per reconcile (see update()/exists()), + // so the cache never leaks across passes. + ps types.VirtualMachinePowerState + psKnown bool } // getHostSystemAncestors looks up and returns vm's host system ancestors, such as "Cluster" and "Datacenter". @@ -1591,18 +1594,20 @@ func (vm *virtualMachine) powerOffVM() (string, error) { } func (vm *virtualMachine) getPowerState() (types.VirtualMachinePowerState, error) { + if vm.psKnown { + return vm.ps, nil + } + powerState, err := vm.Obj.PowerState(vm.Context) if err != nil { return "", err } switch powerState { - case types.VirtualMachinePowerStatePoweredOn: - return types.VirtualMachinePowerStatePoweredOn, nil - case types.VirtualMachinePowerStatePoweredOff: - return types.VirtualMachinePowerStatePoweredOff, nil - case types.VirtualMachinePowerStateSuspended: - return types.VirtualMachinePowerStateSuspended, nil + case types.VirtualMachinePowerStatePoweredOn, types.VirtualMachinePowerStatePoweredOff, types.VirtualMachinePowerStateSuspended: + vm.ps = powerState + vm.psKnown = true + return powerState, nil default: return "", fmt.Errorf("unexpected power state %q for vm %v", powerState, vm) } @@ -1729,20 +1734,21 @@ type NetworkStatus struct { NetworkName string } -func (vm *virtualMachine) getNetworkStatusList(client *vim25.Client) ([]NetworkStatus, error) { +func (vm *virtualMachine) getNetworkStatusList(client *vim25.Client) ([]NetworkStatus, string, error) { var obj mo.VirtualMachine var pc = property.DefaultCollector(client) var props = []string{ "config.hardware.device", "guest.net", + "name", } if err := pc.RetrieveOne(vm.Context, vm.Ref, props, &obj); err != nil { - return nil, fmt.Errorf("unable to fetch props %v for vm %v: %w", props, vm.Ref, err) + return nil, "", fmt.Errorf("unable to fetch props %v for vm %v: %w", props, vm.Ref, err) } klog.V(3).Infof("Getting network status: object reference: %v", obj.Reference().Value) if obj.Config == nil { - return nil, errors.New("config.hardware.device is nil") + return nil, "", errors.New("config.hardware.device is nil") } var networkStatusList []NetworkStatus @@ -1769,7 +1775,7 @@ func (vm *virtualMachine) getNetworkStatusList(client *vim25.Client) ([]NetworkS } } - return networkStatusList, nil + return networkStatusList, obj.Name, nil } type attachedDisk struct { diff --git a/pkg/controller/vsphere/reconciler_test.go b/pkg/controller/vsphere/reconciler_test.go index 9d7b2feaa..be5a1665b 100644 --- a/pkg/controller/vsphere/reconciler_test.go +++ b/pkg/controller/vsphere/reconciler_test.go @@ -1492,7 +1492,7 @@ func TestGetNetworkStatusList(t *testing.T) { } // validations - networkStatusList, err := vm.getNetworkStatusList(session.Client.Client) + networkStatusList, _, err := vm.getNetworkStatusList(session.Client.Client) if err != nil { t.Fatal(err) } @@ -3727,3 +3727,42 @@ func TestReconcileRegionAndZoneLabelsSkipsWhenSet(t *testing.T) { t.Errorf("labels were modified: %v", machine.Labels) } } + +func TestReconcileProviderIDSkipsWhenSet(t *testing.T) { + pid := "vsphere://564d...c7f6" + machine := &machinev1.Machine{} + machine.Spec.ProviderID = &pid + r := &Reconciler{ + machineScope: &machineScope{machine: machine, providerStatus: &machinev1.VSphereMachineProviderStatus{}}, + } + // vm == nil: if the function calls into the VM client it panics; + // the guard must return first. + if err := r.reconcileProviderID(nil); err != nil { + t.Fatalf("expected nil, got %v", err) + } +} + +func TestGetPowerStateCachedWithinPass(t *testing.T) { + _, sess, server := initSimulator(t) + defer server.Close() + ctx := context.Background() + + vmObj, err := sess.Finder.VirtualMachine(ctx, "DC0/host/DC0_H0/VM0") + if err != nil { + // adjust inventory path to the sim topology used by this suite + t.Skipf("no default VM in sim: %v", err) + } + vm := &virtualMachine{Context: ctx, Obj: vmObj, Ref: vmObj.Reference()} + + first, err := vm.getPowerState() + if err != nil { + t.Fatal(err) + } + second, err := vm.getPowerState() + if err != nil { + t.Fatal(err) + } + if first != second { + t.Errorf("cached and fresh power states differ: %s vs %s", first, second) + } +} From 5ae318a8670d1f832be424461ae9cf1561d630d1 Mon Sep 17 00:00:00 2001 From: Joseph Callen Date: Fri, 11 Sep 2026 17:07:31 -0400 Subject: [PATCH 07/11] vsphere: reset machine identity when re-cloning VM --- pkg/controller/vsphere/reconciler.go | 4 +++ pkg/controller/vsphere/reconciler_test.go | 34 +++++++++++++++++++++++ 2 files changed, 38 insertions(+) diff --git a/pkg/controller/vsphere/reconciler.go b/pkg/controller/vsphere/reconciler.go index 21dc61716..af2b61f0f 100644 --- a/pkg/controller/vsphere/reconciler.go +++ b/pkg/controller/vsphere/reconciler.go @@ -157,6 +157,10 @@ func (r *Reconciler) create() error { } klog.Infof("%v: cloning", r.machine.GetName()) + // A new clone has a different identity. Clear values from a VM that + // may have disappeared so the next update records the new VM identity. + r.machine.Spec.ProviderID = nil + r.providerStatus.InstanceID = nil task, err := clone(r.machineScope) if err != nil { metrics.RegisterFailedInstanceCreate(&metrics.MachineLabels{ diff --git a/pkg/controller/vsphere/reconciler_test.go b/pkg/controller/vsphere/reconciler_test.go index be5a1665b..3d5f7b4df 100644 --- a/pkg/controller/vsphere/reconciler_test.go +++ b/pkg/controller/vsphere/reconciler_test.go @@ -1466,6 +1466,40 @@ func createDataDiskDefinitions(numOfDataDisks int) []machinev1.VSphereDisk { return disks } +func TestSetProviderStatusPreservesExistingInstanceID(t *testing.T) { + model, sess, server := initSimulator(t) + defer model.Remove() + defer server.Close() + + managedObj := model.Map().Any("VirtualMachine").(*simulator.VirtualMachine) + vmRef := managedObj.Reference() + vm := &virtualMachine{ + Context: context.Background(), + Obj: object.NewVirtualMachine(sess.Client.Client, vmRef), + Ref: vmRef, + } + + const existingInstanceID = "existing-instance-id" + scope := &machineScope{ + Context: context.Background(), + machine: &machinev1.Machine{ObjectMeta: metav1.ObjectMeta{Name: "test-machine"}}, + providerStatus: &machinev1.VSphereMachineProviderStatus{ + InstanceID: func() *string { v := existingInstanceID; return &v }(), + }, + } + + if err := setProviderStatus("", conditionSuccess(), scope, vm); err != nil { + t.Fatal(err) + } + got := "" + if scope.providerStatus.InstanceID != nil { + got = *scope.providerStatus.InstanceID + } + if got != existingInstanceID { + t.Errorf("InstanceID changed from %q to %q", existingInstanceID, got) + } +} + func TestGetNetworkStatusList(t *testing.T) { model, session, server := initSimulator(t) defer model.Remove() From 5e4bc6ca34155ef9b9aba826990cfb1491b6fb2c Mon Sep 17 00:00:00 2001 From: Joseph Callen Date: Fri, 11 Sep 2026 17:22:29 -0400 Subject: [PATCH 08/11] vsphere: guard TaskIDCache with a mutex --- pkg/controller/vsphere/actuator.go | 32 +++++++++++++++++++++---- pkg/controller/vsphere/actuator_test.go | 21 ++++++++++++++++ 2 files changed, 49 insertions(+), 4 deletions(-) diff --git a/pkg/controller/vsphere/actuator.go b/pkg/controller/vsphere/actuator.go index 6b9aad746..fd25d220a 100644 --- a/pkg/controller/vsphere/actuator.go +++ b/pkg/controller/vsphere/actuator.go @@ -5,6 +5,7 @@ package vsphere import ( "context" "fmt" + "sync" "time" "k8s.io/component-base/featuregate" @@ -33,6 +34,7 @@ type Actuator struct { apiReader runtimeclient.Reader eventRecorder events.EventRecorder TaskIDCache map[string]string + taskIDCacheMu sync.Mutex FeatureGates featuregate.MutableFeatureGate openshiftConfigNamespace string } @@ -59,6 +61,28 @@ func NewActuator(params ActuatorParams) *Actuator { } } +func (a *Actuator) getTaskID(machineName string) (string, bool) { + a.taskIDCacheMu.Lock() + defer a.taskIDCacheMu.Unlock() + value, ok := a.TaskIDCache[machineName] + return value, ok +} + +func (a *Actuator) setTaskID(machineName, taskID string) { + a.taskIDCacheMu.Lock() + defer a.taskIDCacheMu.Unlock() + if a.TaskIDCache == nil { + a.TaskIDCache = make(map[string]string) + } + a.TaskIDCache[machineName] = taskID +} + +func (a *Actuator) clearTaskID(machineName string) { + a.taskIDCacheMu.Lock() + defer a.taskIDCacheMu.Unlock() + delete(a.TaskIDCache, machineName) +} + // Set corresponding event based on error. It also returns the original error // for convenience, so callers can do "return handleMachineError(...)". func (a *Actuator) handleMachineError(machine *machinev1.Machine, err error, eventAction string) error { @@ -88,7 +112,7 @@ func (a *Actuator) Create(ctx context.Context, machine *machinev1.Machine) error // Ensure we're not reconciling a stale machine by checking our task-id. // This is a workaround for a cache race condition. - if val, ok := a.TaskIDCache[machine.Name]; ok { + if val, ok := a.getTaskID(machine.Name); ok { if val != scope.providerStatus.TaskRef { klog.Errorf("%s: machine object missing expected provider task ID, requeue", machine.GetName()) return &machinecontroller.RequeueAfterError{RequeueAfter: requeueAfterSeconds * time.Second} @@ -99,7 +123,7 @@ func (a *Actuator) Create(ctx context.Context, machine *machinev1.Machine) error err = newReconciler(scope).create() // save the taskRef in our cache in case of any error with patch. if scope.providerStatus.TaskRef != "" { - a.TaskIDCache[machine.Name] = scope.providerStatus.TaskRef + a.setTaskID(machine.Name, scope.providerStatus.TaskRef) } if err != nil { fmtErr := fmt.Errorf(reconcilerFailFmt, machine.GetName(), createEventAction, err) @@ -134,7 +158,7 @@ func (a *Actuator) Exists(ctx context.Context, machine *machinev1.Machine) (bool func (a *Actuator) Update(ctx context.Context, machine *machinev1.Machine) error { klog.Infof("%s: actuator updating machine", machine.GetName()) // Cleanup TaskIDCache so we don't continually grow - delete(a.TaskIDCache, machine.Name) + a.clearTaskID(machine.Name) scope, err := newMachineScope(machineScopeParams{ Context: ctx, @@ -176,7 +200,7 @@ func (a *Actuator) Delete(ctx context.Context, machine *machinev1.Machine) error klog.Infof("%s: actuator deleting machine", machine.GetName()) // Cleanup TaskIDCache so we don't continually grow // Cleanup here as well in case Update() was never successfully called. - delete(a.TaskIDCache, machine.Name) + a.clearTaskID(machine.Name) scope, err := newMachineScope(machineScopeParams{ Context: ctx, diff --git a/pkg/controller/vsphere/actuator_test.go b/pkg/controller/vsphere/actuator_test.go index 38118eed5..8f805f770 100644 --- a/pkg/controller/vsphere/actuator_test.go +++ b/pkg/controller/vsphere/actuator_test.go @@ -5,6 +5,7 @@ import ( "fmt" "net" "path/filepath" + "sync" "testing" "time" @@ -422,3 +423,23 @@ func TestMachineEvents(t *testing.T) { }) } } + +func TestTaskIDCacheConcurrentAccess(t *testing.T) { + actuator := &Actuator{TaskIDCache: make(map[string]string)} + + const workers = 100 + var wg sync.WaitGroup + wg.Add(workers) + for i := 0; i < workers; i++ { + go func(i int) { + defer wg.Done() + machineName := fmt.Sprintf("machine-%d", i) + actuator.setTaskID(machineName, "task") + if taskID, ok := actuator.getTaskID(machineName); !ok || taskID != "task" { + t.Errorf("getTaskID(%q) = %q, %t; want task, true", machineName, taskID, ok) + } + actuator.clearTaskID(machineName) + }(i) + } + wg.Wait() +} From bb4e6417df754a22c7b4ac729ae98d56ddbda0bc Mon Sep 17 00:00:00 2001 From: Joseph Callen Date: Fri, 11 Sep 2026 17:22:35 -0400 Subject: [PATCH 09/11] vsphere: clear terminal failed task ref in update --- pkg/controller/vsphere/reconciler.go | 3 +++ pkg/controller/vsphere/reconciler_test.go | 19 ++++++++++++++++--- 2 files changed, 19 insertions(+), 3 deletions(-) diff --git a/pkg/controller/vsphere/reconciler.go b/pkg/controller/vsphere/reconciler.go index af2b61f0f..57c252892 100644 --- a/pkg/controller/vsphere/reconciler.go +++ b/pkg/controller/vsphere/reconciler.go @@ -289,6 +289,9 @@ func (r *Reconciler) update() error { Namespace: r.machine.Namespace, Reason: "Task finished with error", }) + // A terminally failed task cannot transition to success. Clear + // its ref so retries reconcile the VM instead of polling it forever. + r.providerStatus.TaskRef = "" return fmt.Errorf("%v task %v finished with error: %w", moTask.Info.DescriptionId, moTask.Reference().Value, err) } else if !taskIsFinished { return fmt.Errorf("%v task %v has not finished", moTask.Info.DescriptionId, moTask.Reference().Value) diff --git a/pkg/controller/vsphere/reconciler_test.go b/pkg/controller/vsphere/reconciler_test.go index 3d5f7b4df..13ef52d11 100644 --- a/pkg/controller/vsphere/reconciler_test.go +++ b/pkg/controller/vsphere/reconciler_test.go @@ -3673,6 +3673,12 @@ func TestUpdateClearsFinishedTaskRef(t *testing.T) { t.Fatal(err) } + failedTask := simulator.CreateTask(vm, "failedTask", func(*simulator.Task) (types.AnyType, types.BaseMethodFault) { + return nil, &types.InvalidArgument{} + }) + failedTaskRef := failedTask.Run(model.Service.Context) + failedTask.Wait() + rawProviderSpec, err := RawExtensionFromProviderSpec(&machinev1.VSphereMachineProviderSpec{ Workspace: &machinev1.Workspace{Server: host}, CredentialsSecret: &corev1.LocalObjectReference{ @@ -3688,11 +3694,13 @@ func TestUpdateClearsFinishedTaskRef(t *testing.T) { } for _, tc := range []struct { - name string - taskRef string + name string + taskRef string + expectError bool }{ {name: "finished task", taskRef: task.Reference().Value}, {name: "stale missing task", taskRef: "task-99999"}, + {name: "failed task", taskRef: failedTaskRef.Value, expectError: true}, } { t.Run(tc.name, func(t *testing.T) { machineObj := &machinev1.Machine{ @@ -3722,7 +3730,12 @@ func TestUpdateClearsFinishedTaskRef(t *testing.T) { } scope.providerStatus.TaskRef = tc.taskRef - if err := newReconciler(scope).update(); err != nil { + err = newReconciler(scope).update() + if tc.expectError { + if err == nil { + t.Fatal("update() succeeded for failed task") + } + } else if err != nil { t.Fatalf("update() error: %v", err) } if scope.providerStatus.TaskRef != "" { From 2142c1a393d4139c67e4356af44792c10c94dc7c Mon Sep 17 00:00:00 2001 From: Joseph Callen Date: Fri, 11 Sep 2026 17:27:36 -0400 Subject: [PATCH 10/11] vsphere: clear TaskRef only for terminal failed tasks --- pkg/controller/vsphere/reconciler.go | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/pkg/controller/vsphere/reconciler.go b/pkg/controller/vsphere/reconciler.go index 57c252892..4450372f8 100644 --- a/pkg/controller/vsphere/reconciler.go +++ b/pkg/controller/vsphere/reconciler.go @@ -289,9 +289,11 @@ func (r *Reconciler) update() error { Namespace: r.machine.Namespace, Reason: "Task finished with error", }) - // A terminally failed task cannot transition to success. Clear - // its ref so retries reconcile the VM instead of polling it forever. - r.providerStatus.TaskRef = "" + if taskIsFinished { + // A terminally failed task cannot transition to success. Clear + // its ref so retries reconcile the VM instead of polling it forever. + r.providerStatus.TaskRef = "" + } return fmt.Errorf("%v task %v finished with error: %w", moTask.Info.DescriptionId, moTask.Reference().Value, err) } else if !taskIsFinished { return fmt.Errorf("%v task %v has not finished", moTask.Info.DescriptionId, moTask.Reference().Value) From 8c2c11d97eb6381ff314eec8192f5128313955b3 Mon Sep 17 00:00:00 2001 From: Joseph Callen Date: Mon, 14 Sep 2026 08:50:19 -0400 Subject: [PATCH 11/11] vsphere: address review findings from bot review - powerOnVM/powerOffVM: invalidate cached power state so a subsequent getPowerState in the same pass sees the new state - create(): treat RetrieveMO NotFound from GetTask like update()/ delete(); probe for the VM and reconcile it into steady state - getPowerState: replace single-case switch with a plain if - --max-concurrent-reconciles: validate [1, 100] in main; add table test - tests: table test for isRetrieveMONotFound; strengthen TestGetPowerStateCachedWithinPass to mutate simulator state Skipped from the bot commit: namespace-keyed TaskIDCache (machines live in one namespace; name-keying is sufficient) and TestTaskRefGoneDuringCreate (its model.SetTask/GetVMByName helpers do not exist in the vendored govmomi simulator). --- cmd/vsphere/main.go | 11 ++++ cmd/vsphere/main_test.go | 22 +++++++ pkg/controller/vsphere/reconciler.go | 54 ++++++++++++---- pkg/controller/vsphere/reconciler_test.go | 75 ++++++++++++++++++++--- 4 files changed, 141 insertions(+), 21 deletions(-) diff --git a/cmd/vsphere/main.go b/cmd/vsphere/main.go index 32fea83ba..276e95d41 100644 --- a/cmd/vsphere/main.go +++ b/cmd/vsphere/main.go @@ -44,6 +44,13 @@ func registerControllerFlags(fs *flag.FlagSet) *int { "shared vCenter environments.") } +func validateMaxConcurrentReconciles(n int) error { + if n < 1 || n > 100 { + return fmt.Errorf("--max-concurrent-reconciles must be in [1, 100]; got %d", n) + } + return nil +} + func main() { var printVersion bool flag.BoolVar(&printVersion, "version", false, "print version and exit") @@ -127,6 +134,10 @@ func main() { flag.Parse() + if err := validateMaxConcurrentReconciles(*maxConcurrentReconciles); err != nil { + klog.Fatalf("%v", err) + } + if printVersion { fmt.Println(version.String) os.Exit(0) diff --git a/cmd/vsphere/main_test.go b/cmd/vsphere/main_test.go index 218327f55..81edb7f22 100644 --- a/cmd/vsphere/main_test.go +++ b/cmd/vsphere/main_test.go @@ -35,3 +35,25 @@ func TestMaxConcurrentReconcilesCustom(t *testing.T) { t.Errorf("expected max-concurrent-reconciles = 5, got %d", *maxConcurrent) } } + +func TestValidateMaxConcurrentReconciles(t *testing.T) { + for _, tc := range []struct { + name string + val int + wantErr bool + }{ + {name: "default", val: 10}, + {name: "min", val: 1}, + {name: "max", val: 100}, + {name: "zero", val: 0, wantErr: true}, + {name: "negative", val: -1, wantErr: true}, + {name: "over limit", val: 101, wantErr: true}, + } { + t.Run(tc.name, func(t *testing.T) { + err := validateMaxConcurrentReconciles(tc.val) + if (err != nil) != tc.wantErr { + t.Errorf("validateMaxConcurrentReconciles(%d) err = %v, wantErr %v", tc.val, err, tc.wantErr) + } + }) + } +} diff --git a/pkg/controller/vsphere/reconciler.go b/pkg/controller/vsphere/reconciler.go index 4450372f8..88b4a4d5f 100644 --- a/pkg/controller/vsphere/reconciler.go +++ b/pkg/controller/vsphere/reconciler.go @@ -181,15 +181,36 @@ func (r *Reconciler) create() error { moTask, err := r.session.GetTask(r.Context, r.providerStatus.TaskRef) if err != nil { - metrics.RegisterFailedInstanceCreate(&metrics.MachineLabels{ - Name: r.machine.Name, - Namespace: r.machine.Namespace, - Reason: "GetTask finished with error", - }) - return err + if !isRetrieveMONotFound(r.providerStatus.TaskRef, err) { + metrics.RegisterFailedInstanceCreate(&metrics.MachineLabels{ + Name: r.machine.Name, + Namespace: r.machine.Namespace, + Reason: "GetTask finished with error", + }) + return err + } + // Task history eviction or a session restart can make the clone + // task ref permanently unavailable. Clear it; the moTask == nil + // block below probes for the VM instead of failing forever. + klog.Infof("%v: task %s no longer found, clearing TaskRef", r.machine.GetName(), r.providerStatus.TaskRef) + r.providerStatus.TaskRef = "" } if moTask == nil { + // The clone task is gone from vCenter. If the VM exists the clone + // succeeded; reconcile it into steady state instead of failing on + // the missing task. + if vmRef, err := findVM(r.machineScope); err != nil { + return err + } else if vmRef != (types.ManagedObjectReference{}) { + klog.Infof("%v: clone task gone but VM found, reconciling VM state", r.machine.GetName()) + vm := &virtualMachine{ + Context: r.machineScope.Context, + Obj: object.NewVirtualMachine(r.machineScope.session.Client.Client, vmRef), + Ref: vmRef, + } + return r.reconcileMachineWithCloudState(vm, r.providerStatus.TaskRef) + } // Possible eventual consistency problem from vsphere // TODO: change error message here to indicate this might be expected. return fmt.Errorf("unexpected moTask nil") @@ -1461,6 +1482,9 @@ func taskIsFinished(task *mo.Task) (bool, error) { } } +// setProviderStatus updates the Machine's ProviderStatus with a task reference, +// a condition, and (optionally) instance metadata. It is called from create(), +// update(), and delete() flows; pass "" for taskRef when there is no active task. func setProviderStatus(taskRef string, condition metav1.Condition, scope *machineScope, vm *virtualMachine) error { klog.Infof("%s: Updating provider status", scope.machine.Name) @@ -1587,6 +1611,9 @@ func (vm *virtualMachine) getRegionAndZone(tagsMgr *session.CachingTagsManager, } func (vm *virtualMachine) powerOnVM() (string, error) { + // Invalidate the power-state cache so a subsequent getPowerState in + // the same reconcile pass observes the new state, not the stale one. + vm.psKnown = false task, err := vm.Obj.PowerOn(vm.Context) if err != nil { return "", err @@ -1595,6 +1622,9 @@ func (vm *virtualMachine) powerOnVM() (string, error) { } func (vm *virtualMachine) powerOffVM() (string, error) { + // Invalidate the power-state cache so a subsequent getPowerState in + // the same reconcile pass observes the new state, not the stale one. + vm.psKnown = false task, err := vm.Obj.PowerOff(vm.Context) if err != nil { return "", err @@ -1612,14 +1642,14 @@ func (vm *virtualMachine) getPowerState() (types.VirtualMachinePowerState, error return "", err } - switch powerState { - case types.VirtualMachinePowerStatePoweredOn, types.VirtualMachinePowerStatePoweredOff, types.VirtualMachinePowerStateSuspended: - vm.ps = powerState - vm.psKnown = true - return powerState, nil - default: + if powerState != types.VirtualMachinePowerStatePoweredOn && + powerState != types.VirtualMachinePowerStatePoweredOff && + powerState != types.VirtualMachinePowerStateSuspended { return "", fmt.Errorf("unexpected power state %q for vm %v", powerState, vm) } + vm.ps = powerState + vm.psKnown = true + return powerState, nil } // reconcileTags ensures that the required tags are present on the virtual machine, eg the Cluster ID diff --git a/pkg/controller/vsphere/reconciler_test.go b/pkg/controller/vsphere/reconciler_test.go index 13ef52d11..94b88c60f 100644 --- a/pkg/controller/vsphere/reconciler_test.go +++ b/pkg/controller/vsphere/reconciler_test.go @@ -3790,26 +3790,83 @@ func TestReconcileProviderIDSkipsWhenSet(t *testing.T) { } func TestGetPowerStateCachedWithinPass(t *testing.T) { - _, sess, server := initSimulator(t) + model, sess, server := initSimulator(t) + defer model.Remove() defer server.Close() ctx := context.Background() - vmObj, err := sess.Finder.VirtualMachine(ctx, "DC0/host/DC0_H0/VM0") - if err != nil { - // adjust inventory path to the sim topology used by this suite - t.Skipf("no default VM in sim: %v", err) - } - vm := &virtualMachine{Context: ctx, Obj: vmObj, Ref: vmObj.Reference()} + simVM := model.Map().Any("VirtualMachine").(*simulator.VirtualMachine) + vmObj := object.NewVirtualMachine(sess.Client.Client, simVM.Reference()) + vm := &virtualMachine{Context: ctx, Obj: vmObj, Ref: simVM.Reference()} first, err := vm.getPowerState() if err != nil { t.Fatal(err) } + + // Mutate the simulator's power state in-process so a non-caching + // implementation would observe a different value on the next call. + simVM.Runtime.PowerState = types.VirtualMachinePowerStateSuspended + + // The second call must still return the cached value, proving the + // cache is used instead of re-querying vCenter within the pass. second, err := vm.getPowerState() if err != nil { t.Fatal(err) } - if first != second { - t.Errorf("cached and fresh power states differ: %s vs %s", first, second) + if second != first { + t.Errorf("cached power state = %s, want %s (simulator state changed to suspended)", second, first) + } +} + +func TestIsRetrieveMONotFound(t *testing.T) { + taskRef := "task-12345" + expectedErr := fmt.Sprintf("ServerFaultCode: The object 'vim.Task:%v' has already been deleted or has not been completely created", taskRef) + + tests := []struct { + name string + taskRef string + err error + want bool + }{ + { + name: "nil error returns false", + taskRef: taskRef, + err: nil, + want: false, + }, + { + name: "RetrieveMO NotFound with full message returns true", + taskRef: taskRef, + err: errors.New(expectedErr), + want: true, + }, + { + name: "RetrieveMO NotFound with generic message returns true", + taskRef: taskRef, + err: errors.New("ServerFaultCode: The object has already been deleted or has not been completely created"), + want: true, + }, + { + name: "other error returns false", + taskRef: taskRef, + err: errors.New("some other error"), + want: false, + }, + { + name: "different task ref in message returns false", + taskRef: taskRef, + err: errors.New("ServerFaultCode: The object 'vim.Task:task-99999' has already been deleted or has not been completely created"), + want: false, + }, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + got := isRetrieveMONotFound(tc.taskRef, tc.err) + if got != tc.want { + t.Errorf("isRetrieveMONotFound(%q, %v) = %v, want %v", tc.taskRef, tc.err, got, tc.want) + } + }) } }