diff --git a/internal/sp/service/resource_manager/service_type_instance.go b/internal/sp/service/resource_manager/service_type_instance.go index 27b007c..bbc875e 100644 --- a/internal/sp/service/resource_manager/service_type_instance.go +++ b/internal/sp/service/resource_manager/service_type_instance.go @@ -82,32 +82,43 @@ func (s *InstanceService) CreateInstance(ctx context.Context, request *resource_ return nil, service.NewValidationError("spec.service_type must not be empty") } - // Send request to provider endpoint with the resolved ID - providerResponse, err := s.createInstanceWithProvider(ctx, provider.Endpoint, request, instanceID) - if err != nil { - log.Error("Provider provisioning failed", "instance_id", *instanceID, "provider_name", providerName, "error", err) - return nil, service.NewProviderError(fmt.Sprintf("Error from Provider (%s): %v", providerName, err)) - } - - // Create instance in database + // Persist a placeholder row *before* dispatching to the provider. This closes + // the race where a status event for this instance (e.g. from StatusConsumer) + // arrives while the create request is still in flight: the row already + // exists, so the update applies instead of being silently and permanently + // discarded as ErrInstanceNotFound (see enhancements/sp-resource-manager and + // sp-resource-status-reader design docs). instance := model.ServiceTypeInstance{ ID: *instanceID, ProviderName: providerName, ServiceType: serviceType, - Status: providerResponse.Status, + Status: rmstore.StatusPending, Spec: request.Spec, } created, err := s.store.ServiceTypeInstance().Create(ctx, instance) if err != nil { log.Error("Failed to create instance in store", "instance_id", *instanceID, "error", err) - return nil, service.NewInternalError(fmt.Sprintf("failed to create database record for instance %s: %v", providerResponse.ID, err)) + return nil, service.NewInternalError(fmt.Sprintf("failed to create database record for instance %s: %v", *instanceID, err)) + } + + // Send request to provider endpoint with the resolved ID + if _, err := s.createInstanceWithProvider(ctx, provider.Endpoint, request, instanceID); err != nil { + log.Error("Provider provisioning failed", "instance_id", *instanceID, "provider_name", providerName, "error", err) + if delErr := s.store.ServiceTypeInstance().HardDelete(ctx, *instanceID); delErr != nil { + log.Error("Failed to roll back placeholder instance after provider failure", "instance_id", *instanceID, "error", delErr) + } + return nil, service.NewProviderError(fmt.Sprintf("Error from Provider (%s): %v", providerName, err)) } + // The row stays PENDING here - StatusConsumer owns every status transition + // from this point on. The provider's synchronous response isn't persisted: + // it's immediately superseded by StatusConsumer's async update in practice, + // and nothing downstream reads it off the create response. log.Info("Instance created successfully", "instance_id", created.ID, "provider_name", providerName, - "status", providerResponse.Status, + "status", created.Status, ) return ModelToAPI(created), nil } diff --git a/internal/sp/service/resource_manager/service_type_instance_test.go b/internal/sp/service/resource_manager/service_type_instance_test.go index a2f4340..59ecc52 100644 --- a/internal/sp/service/resource_manager/service_type_instance_test.go +++ b/internal/sp/service/resource_manager/service_type_instance_test.go @@ -13,6 +13,7 @@ import ( rmsvc "github.com/dcm-project/control-plane/internal/sp/service/resource_manager" "github.com/dcm-project/control-plane/internal/sp/store" "github.com/dcm-project/control-plane/internal/sp/store/model" + "github.com/dcm-project/control-plane/internal/sp/testutil" "github.com/go-resty/resty/v2" "github.com/google/uuid" . "github.com/onsi/ginkgo/v2" @@ -69,7 +70,7 @@ var _ = Describe("InstanceService", func() { } Expect(db.Create(&provider).Error).NotTo(HaveOccurred()) - dataStore = store.NewStore(db) + dataStore = store.NewStore(db, store.WithServiceTypeInstanceRetry(testutil.FastServiceTypeInstanceRetry()...)) instanceService = rmsvc.NewInstanceService(dataStore, resty.New(). SetTimeout(5*time.Second). SetRetryCount(0)) @@ -342,18 +343,17 @@ var _ = Describe("InstanceService", func() { Expect(providerCalled).To(BeFalse()) }) - It("returns internal error with instance ID when DB insert fails", func() { - var instanceID string - var providerCallCount int - mockProviderWithID := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { - providerCallCount++ - instanceID = uuid.New().String() - - if providerCallCount == 1 { - sqlDB, _ := db.DB() - _ = sqlDB.Close() - } - + It("applies a status update that arrives while the create request is still in flight to the provider (BAC-1)", func() { + // The mock provider's handler is the exact dispatch point: while the + // synchronous create call to the provider is still outstanding, simulate + // a status event for this same instance arriving via the real production + // code path (StatusConsumer.handleMessage calls this same store method). + // The handler runs on httptest's own goroutine, so the result must cross + // back to the test over a channel rather than a shared variable. + statusUpdateErrCh := make(chan error, 1) + mockProviderMidDispatch := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + instanceID := r.URL.Query().Get("id") + statusUpdateErrCh <- dataStore.ServiceTypeInstance().UpdateStatus(ctx, instanceID, "RUNNING", "provisioning started") w.Header().Set("Content-Type", "application/json") w.WriteHeader(http.StatusOK) _ = json.NewEncoder(w).Encode(map[string]string{ @@ -361,19 +361,70 @@ var _ = Describe("InstanceService", func() { "status": "PROVISIONING", }) })) - defer mockProviderWithID.Close() + defer mockProviderMidDispatch.Close() + + providerMidDispatch := model.Provider{ + ID: uuid.New().String(), + Name: "provider-mid-dispatch", + ServiceType: "vm", + Endpoint: mockProviderMidDispatch.URL, + HealthStatus: model.HealthStatusReady, + } + Expect(db.Create(&providerMidDispatch).Error).NotTo(HaveOccurred()) + + req := &resource_manager.ServiceTypeInstance{ + ProviderName: "provider-mid-dispatch", + Spec: map[string]interface{}{"cpu": 1, "service_type": "vm"}, + } + + _, err := instanceService.CreateInstance(ctx, req, nil) + Expect(err).NotTo(HaveOccurred()) + + statusUpdateErr := <-statusUpdateErrCh + Expect(statusUpdateErr).NotTo(HaveOccurred()) + }) + + It("does not leave an orphaned instance visible after a provider failure (BAC-2)", func() { + mockProviderFailing := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + _, _ = w.Write([]byte(`{"error": "boom"}`)) + })) + defer mockProviderFailing.Close() - providerWithID := model.Provider{ + providerFailing := model.Provider{ ID: uuid.New().String(), - Name: "provider-db-fail", + Name: "provider-failing", ServiceType: "vm", - Endpoint: mockProviderWithID.URL, + Endpoint: mockProviderFailing.URL, HealthStatus: model.HealthStatusReady, } - Expect(db.Create(&providerWithID).Error).NotTo(HaveOccurred()) + Expect(db.Create(&providerFailing).Error).NotTo(HaveOccurred()) + specifiedID := uuid.New().String() req := &resource_manager.ServiceTypeInstance{ - ProviderName: "provider-db-fail", + ProviderName: "provider-failing", + Spec: map[string]interface{}{"cpu": 1, "service_type": "vm"}, + } + + _, err := instanceService.CreateInstance(ctx, req, &specifiedID) + Expect(err).To(HaveOccurred()) + + _, err = instanceService.GetInstance(ctx, specifiedID, false) + Expect(err).To(HaveOccurred()) + var svcErr *service.ServiceError + Expect(err).To(BeAssignableToTypeOf(svcErr)) + errors.As(err, &svcErr) + Expect(svcErr.Code).To(Equal(service.ErrCodeNotFound)) + }) + + It("returns internal error naming the instance id and never contacts the provider when the DB persist fails (BAC-4)", func() { + // Drop only the service_type_instances table so the provider lookup + // (a separate table) still succeeds, and the failure is isolated to the + // initial persist step - before the provider is ever dispatched to. + Expect(db.Migrator().DropTable(&model.ServiceTypeInstance{})).To(Succeed()) + + req := &resource_manager.ServiceTypeInstance{ + ProviderName: "test-provider", Spec: map[string]interface{}{"cpu": 2, "service_type": "vm"}, } @@ -385,7 +436,8 @@ var _ = Describe("InstanceService", func() { errors.As(err, &svcErr) Expect(svcErr.Code).To(Equal(service.ErrCodeInternal)) Expect(svcErr.Message).To(ContainSubstring("failed to create database record")) - Expect(svcErr.Message).To(ContainSubstring(instanceID)) + Expect(svcErr.Message).To(MatchRegexp(`[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}`)) + Expect(providerCalled).To(BeFalse()) }) }) diff --git a/internal/sp/store/resource_manager/service_instance.go b/internal/sp/store/resource_manager/service_instance.go index d54b4e7..029d181 100644 --- a/internal/sp/store/resource_manager/service_instance.go +++ b/internal/sp/store/resource_manager/service_instance.go @@ -192,6 +192,14 @@ const ( DeletionStatusPendingProvider = "PENDING_PROVIDER" ) +// StatusPending is the placeholder status an instance record is created with +// before its create request is dispatched to the provider. Persisting this +// row first (rather than after a successful dispatch) closes the window +// where a status event for the instance could arrive before any record of +// it exists - see enhancements/sp-resource-manager and +// sp-resource-status-reader design docs. +const StatusPending = "PENDING" + func (s *ServiceTypeInstanceStore) MarkForDeletion(ctx context.Context, id string) error { now := time.Now() result := s.db.WithContext(ctx).