diff --git a/controllers/helpers_test.go b/controllers/helpers_test.go index 07e65f9dfa..c2578d7aec 100644 --- a/controllers/helpers_test.go +++ b/controllers/helpers_test.go @@ -36,12 +36,12 @@ import ( credentialtypes "github.com/nutanix-cloud-native/prism-go-client/environment/credentials" prismclientv3 "github.com/nutanix-cloud-native/prism-go-client/v3" clusterModels "github.com/nutanix/ntnx-api-golang-clients/clustermgmt-go-client/v4/models/clustermgmt/v4/config" + dataPoliciesModels "github.com/nutanix/ntnx-api-golang-clients/datapolicies-go-client/v4/models/datapolicies/v4/config" iamModels "github.com/nutanix/ntnx-api-golang-clients/iam-go-client/v4/models/iam/v4/authn" subnetModels "github.com/nutanix/ntnx-api-golang-clients/networking-go-client/v4/models/networking/v4/config" prismNetworkingModels "github.com/nutanix/ntnx-api-golang-clients/networking-go-client/v4/models/prism/v4/config" prismModels "github.com/nutanix/ntnx-api-golang-clients/prism-go-client/v4/models/prism/v4/config" prismErrors "github.com/nutanix/ntnx-api-golang-clients/prism-go-client/v4/models/prism/v4/error" - dataPoliciesModels "github.com/nutanix/ntnx-api-golang-clients/datapolicies-go-client/v4/models/datapolicies/v4/config" vmmModels "github.com/nutanix/ntnx-api-golang-clients/vmm-go-client/v4/models/vmm/v4/ahv/config" policyModels "github.com/nutanix/ntnx-api-golang-clients/vmm-go-client/v4/models/vmm/v4/ahv/policies" imageModels "github.com/nutanix/ntnx-api-golang-clients/vmm-go-client/v4/models/vmm/v4/content" @@ -3372,3 +3372,138 @@ func TestGetVMUUID(t *testing.T) { }) } } + +func TestEnrichTaskErrorWithFailedSubtasks(t *testing.T) { + parentUUID := "parent-create-vm" + childUUID := "child-vnic" + parentErr := errors.New("task parent-create-vm failed: Failed to perform the operation on the VM with UUID 'EMPTY' as '' is not defined") + dhcpMsg := "Cannot allocate address! No DHCP pool is defined or the DHCP pool is exhausted" + + t.Run("appends failed child subtask errors from list", func(t *testing.T) { + ctrl := gomock.NewController(t) + defer ctrl.Finish() + + mockClient := NewMockConvergedClient(ctrl) + failedStatus := prismModels.TASKSTATUS_FAILED + mockClient.MockTasks.EXPECT().List(gomock.Any(), gomock.Any()).Return([]prismModels.Task{ + { + ExtId: ptr.To(childUUID), + Status: &failedStatus, + OperationDescription: ptr.To("Create VM vNIC port"), + ErrorMessages: []prismErrors.AppMessage{ + {Message: ptr.To(dhcpMsg)}, + }, + }, + }, nil) + // Nested collect for the failed child lists its own children. + mockClient.MockTasks.EXPECT().List(gomock.Any(), gomock.Any()).Return([]prismModels.Task{}, nil) + + err := enrichTaskErrorWithFailedSubtasks(context.Background(), mockClient.Client, parentUUID, parentErr) + require.Error(t, err) + assert.ErrorContains(t, err, parentErr.Error()) + assert.ErrorContains(t, err, "failed_subtasks:") + assert.ErrorContains(t, err, dhcpMsg) + assert.ErrorContains(t, err, "Create VM vNIC port") + }) + + t.Run("does not fall back to Get when list returns empty", func(t *testing.T) { + ctrl := gomock.NewController(t) + defer ctrl.Finish() + + mockClient := NewMockConvergedClient(ctrl) + mockClient.MockTasks.EXPECT().List(gomock.Any(), gomock.Any()).Return([]prismModels.Task{}, nil) + + err := enrichTaskErrorWithFailedSubtasks(context.Background(), mockClient.Client, parentUUID, parentErr) + require.Error(t, err) + assert.Equal(t, parentErr, err) + }) + + t.Run("falls back to parent subtask references when list errors", func(t *testing.T) { + ctrl := gomock.NewController(t) + defer ctrl.Finish() + + mockClient := NewMockConvergedClient(ctrl) + failedStatus := prismModels.TASKSTATUS_FAILED + mockClient.MockTasks.EXPECT().List(gomock.Any(), gomock.Any()).Return(nil, errors.New("list failed")) + mockClient.MockTasks.EXPECT().Get(gomock.Any(), parentUUID).Return(&prismModels.Task{ + ExtId: ptr.To(parentUUID), + SubTasks: []prismModels.TaskReferenceInternal{ + {ExtId: ptr.To(childUUID)}, + }, + }, nil) + mockClient.MockTasks.EXPECT().Get(gomock.Any(), childUUID).Return(&prismModels.Task{ + ExtId: ptr.To(childUUID), + Status: &failedStatus, + OperationDescription: ptr.To("Create VM vNIC port"), + LegacyErrorMessage: ptr.To(dhcpMsg), + }, nil) + mockClient.MockTasks.EXPECT().List(gomock.Any(), gomock.Any()).Return([]prismModels.Task{}, nil) + + err := enrichTaskErrorWithFailedSubtasks(context.Background(), mockClient.Client, parentUUID, parentErr) + require.Error(t, err) + assert.ErrorContains(t, err, dhcpMsg) + }) + + t.Run("deduplicates identical ErrorMessages and LegacyErrorMessage", func(t *testing.T) { + ctrl := gomock.NewController(t) + defer ctrl.Finish() + + mockClient := NewMockConvergedClient(ctrl) + failedStatus := prismModels.TASKSTATUS_FAILED + mockClient.MockTasks.EXPECT().List(gomock.Any(), gomock.Any()).Return([]prismModels.Task{ + { + ExtId: ptr.To(childUUID), + Status: &failedStatus, + OperationDescription: ptr.To("Create VM vNIC port"), + ErrorMessages: []prismErrors.AppMessage{ + {Message: ptr.To(dhcpMsg)}, + }, + LegacyErrorMessage: ptr.To(dhcpMsg), + }, + }, nil) + mockClient.MockTasks.EXPECT().List(gomock.Any(), gomock.Any()).Return([]prismModels.Task{}, nil) + + err := enrichTaskErrorWithFailedSubtasks(context.Background(), mockClient.Client, parentUUID, parentErr) + require.Error(t, err) + assert.Equal(t, 1, strings.Count(err.Error(), dhcpMsg)) + }) + + t.Run("returns original error when parent error is nil", func(t *testing.T) { + err := enrichTaskErrorWithFailedSubtasks(context.Background(), nil, parentUUID, nil) + assert.NoError(t, err) + }) +} + +func TestWaitForConvergedOperation_EnrichesFailedWait(t *testing.T) { + ctrl := gomock.NewController(t) + defer ctrl.Finish() + + ctx := context.Background() + parentUUID := "parent-task" + childUUID := "child-task" + waitErr := errors.New("task parent-task failed: generic parent error") + childMsg := "Cannot allocate address! No DHCP pool is defined or the DHCP pool is exhausted" + + mockClient := NewMockConvergedClient(ctrl) + mockOp := mockconverged.NewMockOperation[vmmModels.Vm](ctrl) + mockOp.EXPECT().Wait(ctx).Return(nil, waitErr) + mockOp.EXPECT().UUID().Return(parentUUID) + + failedStatus := prismModels.TASKSTATUS_FAILED + mockClient.MockTasks.EXPECT().List(gomock.Any(), gomock.Any()).Return([]prismModels.Task{ + { + ExtId: ptr.To(childUUID), + Status: &failedStatus, + OperationDescription: ptr.To("Create VM vNIC port"), + ErrorMessages: []prismErrors.AppMessage{ + {Message: ptr.To(childMsg)}, + }, + }, + }, nil) + mockClient.MockTasks.EXPECT().List(gomock.Any(), gomock.Any()).Return([]prismModels.Task{}, nil) + + _, err := waitForConvergedOperation(ctx, mockClient.Client, mockOp) + require.Error(t, err) + assert.ErrorContains(t, err, "generic parent error") + assert.ErrorContains(t, err, childMsg) +} diff --git a/controllers/nutanixmachine_controller.go b/controllers/nutanixmachine_controller.go index 991e24509f..885eadc712 100644 --- a/controllers/nutanixmachine_controller.go +++ b/controllers/nutanixmachine_controller.go @@ -927,13 +927,7 @@ func (r *NutanixMachineReconciler) getOrCreateVM(rctx *nctx.MachineContext) (*vm // if VM exists if vmFound != nil { log.Info(fmt.Sprintf("vm %s found with UUID %s", *vmFound.Name, rctx.NutanixMachine.Status.VmUUID)) - - v1beta1conditions.MarkTrue(rctx.NutanixMachine, infrav1.VMProvisionedCondition) - v1beta2conditions.Set(rctx.NutanixMachine, metav1.Condition{ - Type: string(infrav1.VMProvisionedCondition), - Status: metav1.ConditionTrue, - Reason: capiv1beta1.ProvisionedV1Beta2Reason, - }) + markVMProvisioned(rctx) return vmFound, nil } @@ -1033,13 +1027,9 @@ func (r *NutanixMachineReconciler) getOrCreateVM(rctx *nctx.MachineContext) (*vm // Create the actual VM/Machine log.Info(fmt.Sprintf("Creating VM with name %s for cluster %s", vmName, rctx.NutanixCluster.Name)) - vm, err = convergedClient.VMs.Create(ctx, vm) + vm, err = createAndWaitForVM(ctx, rctx, vm, vmName) if err != nil { - errorMsg := fmt.Errorf("failed to create VM %s: %w", vmName, err) - if !isRetryableAPIError(err) { - rctx.SetFailureStatus(createErrorFailureReason, errorMsg) - } - return nil, errorMsg + return nil, err } vmUuid := *vm.ExtId @@ -1059,13 +1049,43 @@ func (r *NutanixMachineReconciler) getOrCreateVM(rctx *nctx.MachineContext) (*vm return nil, err } + markVMProvisioned(rctx) + return vm, nil +} + +func markVMProvisioned(rctx *nctx.MachineContext) { v1beta1conditions.MarkTrue(rctx.NutanixMachine, infrav1.VMProvisionedCondition) v1beta2conditions.Set(rctx.NutanixMachine, metav1.Condition{ Type: string(infrav1.VMProvisionedCondition), Status: metav1.ConditionTrue, Reason: capiv1beta1.ProvisionedV1Beta2Reason, }) - return vm, nil +} + +func createAndWaitForVM(ctx context.Context, rctx *nctx.MachineContext, vm *vmmconfig.Vm, vmName string) (*vmmconfig.Vm, error) { + convergedClient := rctx.ConvergedClient + createOp, err := convergedClient.VMs.CreateAsync(ctx, vm) + if err != nil { + return nil, vmCreateFailure(rctx, vmName, err) + } + createdVMs, err := waitForConvergedOperation(ctx, convergedClient, createOp) + if err != nil { + return nil, vmCreateFailure(rctx, vmName, err) + } + if len(createdVMs) != 1 || createdVMs[0] == nil { + errorMsg := fmt.Errorf("failed to create VM %s: operation completed but expected exactly 1 VM, got %d", vmName, len(createdVMs)) + rctx.SetFailureStatus(createErrorFailureReason, errorMsg) + return nil, errorMsg + } + return createdVMs[0], nil +} + +func vmCreateFailure(rctx *nctx.MachineContext, vmName string, err error) error { + errorMsg := fmt.Errorf("failed to create VM %s: %w", vmName, err) + if !isRetryableAPIError(err) { + rctx.SetFailureStatus(createErrorFailureReason, errorMsg) + } + return errorMsg } // addCustomAttributes sets custom attributes on the VM, including the provider ID. diff --git a/controllers/nutanixmachine_controller_test.go b/controllers/nutanixmachine_controller_test.go index 67611ac610..3520a61583 100644 --- a/controllers/nutanixmachine_controller_test.go +++ b/controllers/nutanixmachine_controller_test.go @@ -3033,11 +3033,13 @@ func TestNutanixMachineReconciler_getOrCreateVM(t *testing.T) { // 2. GetTaskUUIDFromVM after VM creation returns task with UUID mockConvergedClient.MockTasks.EXPECT().List(ctx, gomock.Any()).Return([]prismModels.Task{}, nil) - // Mock CreateVM + // Mock CreateVM (async + wait) createdVM := vmmModels.NewVm() createdVM.Name = ptr.To(vmName) createdVM.ExtId = ptr.To(vmUUID) - mockConvergedClient.MockVMs.EXPECT().Create(ctx, gomock.Any()).Return(createdVM, nil) + mockCreateOp := mockconverged.NewMockOperation[vmmModels.Vm](ctrl) + mockCreateOp.EXPECT().Wait(ctx).Return([]*vmmModels.Vm{createdVM}, nil) + mockConvergedClient.MockVMs.EXPECT().CreateAsync(ctx, gomock.Any()).Return(mockCreateOp, nil) // Create machine context rctx := &nctx.MachineContext{ @@ -3308,6 +3310,82 @@ func TestNutanixMachineReconciler_getOrCreateVM(t *testing.T) { }) } +func Test_createAndWaitForVM(t *testing.T) { + ctx := context.Background() + vmName := "test-vm" + createdVM := vmmModels.NewVm() + createdVM.Name = ptr.To(vmName) + createdVM.ExtId = ptr.To("vm-uuid") + + tests := []struct { + name string + waitResult []*vmmModels.Vm + wantErrSubstr string + wantFailure bool + wantReturnedVM bool + }{ + { + name: "returns the VM when wait yields exactly one", + waitResult: []*vmmModels.Vm{createdVM}, + wantReturnedVM: true, + }, + { + name: "fails when wait yields no VMs", + waitResult: []*vmmModels.Vm{}, + wantErrSubstr: "expected exactly 1 VM, got 0", + wantFailure: true, + }, + { + name: "fails when wait yields a nil VM", + waitResult: []*vmmModels.Vm{nil}, + wantErrSubstr: "expected exactly 1 VM, got 1", + wantFailure: true, + }, + { + name: "fails when wait yields more than one VM", + waitResult: []*vmmModels.Vm{createdVM, createdVM}, + wantErrSubstr: "expected exactly 1 VM, got 2", + wantFailure: true, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + ctrl := gomock.NewController(t) + defer ctrl.Finish() + + mockConvergedClient := NewMockConvergedClient(ctrl) + mockCreateOp := mockconverged.NewMockOperation[vmmModels.Vm](ctrl) + mockCreateOp.EXPECT().Wait(ctx).Return(tt.waitResult, nil) + mockConvergedClient.MockVMs.EXPECT().CreateAsync(ctx, gomock.Any()).Return(mockCreateOp, nil) + + ntnxMachine := &infrav1.NutanixMachine{} + rctx := &nctx.MachineContext{ + Context: ctx, + NutanixMachine: ntnxMachine, + ConvergedClient: mockConvergedClient.Client, + } + + vm, err := createAndWaitForVM(ctx, rctx, vmmModels.NewVm(), vmName) + if tt.wantReturnedVM { + require.NoError(t, err) + require.NotNil(t, vm) + assert.Equal(t, createdVM, vm) + assert.Nil(t, ntnxMachine.Status.FailureReason) + return + } + + require.Error(t, err) + assert.Nil(t, vm) + assert.ErrorContains(t, err, tt.wantErrSubstr) + if tt.wantFailure { + require.NotNil(t, ntnxMachine.Status.FailureReason) + assert.Equal(t, createErrorFailureReason, *ntnxMachine.Status.FailureReason) + } + }) + } +} + func TestNutanixMachineReconciler_addCustomAttributes(t *testing.T) { const ( vmName = "test-vm" diff --git a/controllers/task_errors.go b/controllers/task_errors.go new file mode 100644 index 0000000000..2e8b57d836 --- /dev/null +++ b/controllers/task_errors.go @@ -0,0 +1,180 @@ +/* +Copyright 2026 Nutanix + +Licensed 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 controllers + +import ( + "context" + "fmt" + "strings" + + prismModels "github.com/nutanix/ntnx-api-golang-clients/prism-go-client/v4/models/prism/v4/config" + "k8s.io/utils/ptr" + ctrl "sigs.k8s.io/controller-runtime" + + "github.com/nutanix-cloud-native/prism-go-client/converged" + v4Converged "github.com/nutanix-cloud-native/prism-go-client/converged/v4" +) + +const failedSubtasksFilter = "parentTask/extId eq '%s' and status eq Prism.Config.TaskStatus'FAILED'" + +// waitForConvergedOperation waits for a Prism task and, on failure, appends +// error messages from failed child subtasks so callers see the actionable +// infrastructure reason (for example DHCP pool exhaustion) rather than only +// the generic parent-task error. +func waitForConvergedOperation[T any](ctx context.Context, client *v4Converged.Client, op converged.Operation[T]) ([]*T, error) { + if op == nil { + return nil, fmt.Errorf("operation is nil") + } + result, err := op.Wait(ctx) + if err != nil { + return nil, enrichTaskErrorWithFailedSubtasks(ctx, client, op.UUID(), err) + } + return result, nil +} + +func enrichTaskErrorWithFailedSubtasks(ctx context.Context, client *v4Converged.Client, parentTaskUUID string, parentErr error) error { + if parentErr == nil || parentTaskUUID == "" || client == nil { + return parentErr + } + + subtaskErrs := collectFailedSubtaskErrors(ctx, client, parentTaskUUID, map[string]struct{}{}) + if len(subtaskErrs) == 0 { + return parentErr + } + + log := ctrl.LoggerFrom(ctx) + log.Error(parentErr, "Prism parent task failed; including failed subtask errors", + "taskUUID", parentTaskUUID, + "subtaskErrors", subtaskErrs, + ) + return fmt.Errorf("%w; failed_subtasks: %s", parentErr, strings.Join(subtaskErrs, "; ")) +} + +func collectFailedSubtaskErrors(ctx context.Context, client *v4Converged.Client, parentTaskUUID string, visited map[string]struct{}) []string { + if parentTaskUUID == "" { + return nil + } + if visited == nil { + visited = map[string]struct{}{} + } + if _, seen := visited[parentTaskUUID]; seen { + return nil + } + visited[parentTaskUUID] = struct{}{} + + log := ctrl.LoggerFrom(ctx) + children, err := listFailedChildTasks(ctx, client, parentTaskUUID) + if err != nil { + log.Error(err, "failed to list Prism subtasks while collecting failure details", "parentTaskUUID", parentTaskUUID) + children = failedSubtasksFromParentGet(ctx, client, parentTaskUUID) + } + + var msgs []string + for i := range children { + child := children[i] + if child.Status == nil || *child.Status != prismModels.TASKSTATUS_FAILED { + continue + } + detail := formatTaskError(&child) + childUUID := ptr.Deref(child.ExtId, "") + log.Info("failed Prism subtask", + "parentTaskUUID", parentTaskUUID, + "subtaskUUID", childUUID, + "operation", taskOperation(&child), + "error", detail, + ) + msgs = append(msgs, detail) + if childUUID != "" { + msgs = append(msgs, collectFailedSubtaskErrors(ctx, client, childUUID, visited)...) + } + } + return msgs +} + +func listFailedChildTasks(ctx context.Context, client *v4Converged.Client, parentTaskUUID string) ([]prismModels.Task, error) { + return client.Tasks.List(ctx, converged.WithFilter(fmt.Sprintf(failedSubtasksFilter, parentTaskUUID))) +} + +func failedSubtasksFromParentGet(ctx context.Context, client *v4Converged.Client, parentTaskUUID string) []prismModels.Task { + log := ctrl.LoggerFrom(ctx) + parent, err := client.Tasks.Get(ctx, parentTaskUUID) + if err != nil { + log.Error(err, "failed to get Prism parent task while collecting subtask failure details", "parentTaskUUID", parentTaskUUID) + return nil + } + if parent == nil || len(parent.SubTasks) == 0 { + return nil + } + + children := make([]prismModels.Task, 0, len(parent.SubTasks)) + for _, ref := range parent.SubTasks { + if ref.ExtId == nil || *ref.ExtId == "" { + continue + } + child, err := client.Tasks.Get(ctx, *ref.ExtId) + if err != nil { + log.Error(err, "failed to get Prism subtask while collecting failure details", "subtaskUUID", *ref.ExtId) + continue + } + if child != nil { + children = append(children, *child) + } + } + return children +} + +func taskOperation(task *prismModels.Task) string { + if task == nil { + return "unknown" + } + if desc := ptr.Deref(task.OperationDescription, ""); desc != "" { + return desc + } + if op := ptr.Deref(task.Operation, ""); op != "" { + return op + } + return "unknown" +} + +func formatTaskError(task *prismModels.Task) string { + op := taskOperation(task) + + var parts []string + if task != nil { + for _, msg := range task.ErrorMessages { + if msg.Message != nil && *msg.Message != "" { + parts = appendUnique(parts, *msg.Message) + } + } + if legacy := ptr.Deref(task.LegacyErrorMessage, ""); legacy != "" { + parts = appendUnique(parts, legacy) + } + } + if len(parts) == 0 { + parts = append(parts, "no error message provided") + } + return fmt.Sprintf("[%s] %s", op, strings.Join(parts, "; ")) +} + +func appendUnique(parts []string, msg string) []string { + for _, existing := range parts { + if existing == msg { + return parts + } + } + return append(parts, msg) +}