diff --git a/pkg/cloudprovider/resources.go b/pkg/cloudprovider/resources.go index 5fbbddabd4..6133b991e7 100644 --- a/pkg/cloudprovider/resources.go +++ b/pkg/cloudprovider/resources.go @@ -306,6 +306,8 @@ type ICloudHost interface { GetIHostNics() ([]ICloudHostNetInterface, error) GetSchedtags() ([]string, error) + + GetOvnVersion() string // just for cloudpods host } type ICloudVM interface { diff --git a/pkg/compute/models/hosts.go b/pkg/compute/models/hosts.go index 61a35f785c..fcc56b5929 100644 --- a/pkg/compute/models/hosts.go +++ b/pkg/compute/models/hosts.go @@ -1772,6 +1772,7 @@ func (self *SHost) syncWithCloudHost(ctx context.Context, userCred mcclient.Toke self.StorageSize = extHost.GetStorageSizeMB() self.StorageType = extHost.GetStorageType() self.HostType = extHost.GetHostType() + self.OvnVersion = extHost.GetOvnVersion() if cpuCmt := extHost.GetCpuCmtbound(); cpuCmt > 0 { self.CpuCmtbound = cpuCmt @@ -1991,6 +1992,7 @@ func (manager *SHostManager) NewFromCloudHost(ctx context.Context, userCred mccl host.ZoneId = izone.Id host.HostType = extHost.GetHostType() + host.OvnVersion = extHost.GetOvnVersion() host.Status = extHost.GetStatus() host.HostStatus = extHost.GetHostStatus() diff --git a/pkg/compute/models/regiondrivers.go b/pkg/compute/models/regiondrivers.go index 1b1b199b61..a5697bfe84 100644 --- a/pkg/compute/models/regiondrivers.go +++ b/pkg/compute/models/regiondrivers.go @@ -146,6 +146,8 @@ type IRegionDriver interface { RequestSyncNatGatewayStatus(ctx context.Context, userCred mcclient.TokenCredential, natgateway *SNatGateway, task taskman.ITask) error RequestSyncBucketStatus(ctx context.Context, userCred mcclient.TokenCredential, bucket *SBucket, task taskman.ITask) error RequestSyncDBInstanceBackupStatus(ctx context.Context, userCred mcclient.TokenCredential, backup *SDBInstanceBackup, task taskman.ITask) error + + RequestCreateNetwork(ctx context.Context, userCred mcclient.TokenCredential, network *SNetwork) error } type IDBInstanceDriver interface { diff --git a/pkg/compute/regiondrivers/base.go b/pkg/compute/regiondrivers/base.go index 99a0c8898c..29426099f8 100644 --- a/pkg/compute/regiondrivers/base.go +++ b/pkg/compute/regiondrivers/base.go @@ -459,3 +459,7 @@ func (self *SBaseRegionDriver) ValidateCreateWafInstanceData(ctx context.Context func (self *SBaseRegionDriver) ValidateCreateWafRuleData(ctx context.Context, userCred mcclient.TokenCredential, waf *models.SWafInstance, input api.WafRuleCreateInput) (api.WafRuleCreateInput, error) { return input, errors.Wrapf(cloudprovider.ErrNotImplemented, "ValidateCreateWafRuleData") } + +func (self *SBaseRegionDriver) RequestCreateNetwork(ctx context.Context, userCred mcclient.TokenCredential, net *models.SNetwork) error { + return errors.Wrapf(cloudprovider.ErrNotImplemented, "RequestCreateNetwork") +} diff --git a/pkg/compute/regiondrivers/managedvirtual.go b/pkg/compute/regiondrivers/managedvirtual.go index 762f8c2c5d..a5550441d7 100644 --- a/pkg/compute/regiondrivers/managedvirtual.go +++ b/pkg/compute/regiondrivers/managedvirtual.go @@ -3121,3 +3121,46 @@ func (self *SManagedVirtualizationRegionDriver) RequestAssociatEip(ctx context.C }) return nil } + +func (self *SManagedVirtualizationRegionDriver) RequestCreateNetwork(ctx context.Context, userCred mcclient.TokenCredential, net *models.SNetwork) error { + wire, err := net.GetWire() + if err != nil { + return errors.Wrapf(err, "GetWire") + } + + iwire, err := wire.GetIWire() + if err != nil { + return errors.Wrapf(err, "GetIWire") + } + + prefix, err := net.GetPrefix() + if err != nil { + return errors.Wrapf(err, "GetPrefix") + } + + opts := cloudprovider.SNetworkCreateOptions{ + Name: net.Name, + Cidr: prefix.String(), + Desc: net.Description, + } + + provider := wire.GetCloudprovider() + opts.ProjectId, err = provider.SyncProject(ctx, userCred, net.ProjectId) + if err != nil { + log.Errorf("failed to sync project %s for create %s network %s error: %v", net.ProjectId, provider.Provider, net.Name, err) + } + + inet, err := iwire.CreateINetwork(&opts) + if err != nil { + return errors.Wrapf(err, "CreateINetwork") + } + + db.SetExternalId(net, userCred, inet.GetGlobalId()) + + err = cloudprovider.WaitStatus(inet, api.NETWORK_STATUS_AVAILABLE, 10*time.Second, 5*time.Minute) + if err != nil { + return errors.Wrapf(err, "Wait network available after 6 minutes status: %s", inet.GetStatus()) + } + + return net.SyncWithCloudNetwork(ctx, userCred, inet, nil, nil) +} diff --git a/pkg/compute/tasks/network_create_task.go b/pkg/compute/tasks/network_create_task.go index 9efc6a628d..67b71f1b11 100644 --- a/pkg/compute/tasks/network_create_task.go +++ b/pkg/compute/tasks/network_create_task.go @@ -16,17 +16,14 @@ package tasks import ( "context" - "fmt" - "time" "yunion.io/x/jsonutils" - "yunion.io/x/log" + "yunion.io/x/pkg/errors" api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" - "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/util/logclient" ) @@ -39,76 +36,40 @@ func init() { taskman.RegisterTask(NetworkCreateTask{}) } -func (self *NetworkCreateTask) taskFailed(ctx context.Context, network *models.SNetwork, event string, err error) { - log.Errorf("network create task fail on %s: %s", event, err.Error()) - reason := jsonutils.NewDict() - reason.Set("event", jsonutils.NewString(event)) - reason.Set("reason", jsonutils.NewString(err.Error())) - network.SetStatus(self.UserCred, api.NETWORK_STATUS_FAILED, reason.String()) - db.OpsLog.LogEvent(network, db.ACT_ALLOCATE_FAIL, reason, self.UserCred) - logclient.AddActionLogWithStartable(self, network, logclient.ACT_CREATE, reason, self.UserCred, false) - self.SetStageFailed(ctx, reason) +func (self *NetworkCreateTask) taskFailed(ctx context.Context, network *models.SNetwork, err error) { + network.SetStatus(self.UserCred, api.NETWORK_STATUS_FAILED, err.Error()) + db.OpsLog.LogEvent(network, db.ACT_ALLOCATE_FAIL, err, self.UserCred) + logclient.AddActionLogWithStartable(self, network, logclient.ACT_CREATE, err, self.UserCred, false) + self.SetStageFailed(ctx, jsonutils.NewString(err.Error())) } func (self *NetworkCreateTask) OnInit(ctx context.Context, obj db.IStandaloneModel, body jsonutils.JSONObject) { - network := obj.(*models.SNetwork) + net := obj.(*models.SNetwork) - network.SetStatus(self.UserCred, api.NETWORK_STATUS_PENDING, "") + net.SetStatus(self.UserCred, api.NETWORK_STATUS_PENDING, "") - wire, _ := network.GetWire() - if wire == nil { - self.taskFailed(ctx, network, "getwire", fmt.Errorf("no vpc")) + region, err := net.GetRegion() + if err != nil { + self.taskFailed(ctx, net, errors.Wrapf(err, "GetRegion")) return } - iwire, err := wire.GetIWire() + driver, err := region.GetRegionDriver() if err != nil { - self.taskFailed(ctx, network, "getiwire", err) + self.taskFailed(ctx, net, errors.Wrapf(err, "GetRegionDriver")) return } - prefix, err := network.GetPrefix() + err = driver.RequestCreateNetwork(ctx, self.GetUserCred(), net) if err != nil { - self.taskFailed(ctx, network, "getprefix", err) + self.taskFailed(ctx, net, errors.Wrapf(err, "RequestCreateNetwork")) return } - opts := cloudprovider.SNetworkCreateOptions{ - Name: network.Name, - Cidr: prefix.String(), - Desc: network.Description, - } - - provider := wire.GetCloudprovider() - opts.ProjectId, err = provider.SyncProject(ctx, self.GetUserCred(), network.ProjectId) - if err != nil { - log.Errorf("failed to sync project %s for create %s network %s error: %v", network.ProjectId, provider.Provider, network.Name, err) - } - - inet, err := iwire.CreateINetwork(&opts) - if err != nil { - self.taskFailed(ctx, network, "createinetwork", err) - return - } - db.SetExternalId(network, self.UserCred, inet.GetGlobalId()) - - err = cloudprovider.WaitStatus(inet, api.NETWORK_STATUS_AVAILABLE, 10*time.Second, 300*time.Second) - if err != nil { - self.taskFailed(ctx, network, "waitstatu", err) - return - } - - err = network.SyncWithCloudNetwork(ctx, self.UserCred, inet, nil, nil) - - if err != nil { - self.taskFailed(ctx, network, "SyncWithCloudNetwork", err) - return - } - - network.ClearSchedDescCache() - logclient.AddActionLogWithStartable(self, network, logclient.ACT_CREATE, "", self.UserCred, true) + net.ClearSchedDescCache() + logclient.AddActionLogWithStartable(self, net, logclient.ACT_CREATE, "", self.UserCred, true) notifyclient.EventNotify(ctx, self.UserCred, notifyclient.SEventNotifyParam{ - Obj: network, + Obj: net, Action: notifyclient.ActionCreate, }) self.SetStageComplete(ctx, nil) diff --git a/pkg/compute/tasks/vpc_create_task.go b/pkg/compute/tasks/vpc_create_task.go index da4282c983..10f4d043ff 100644 --- a/pkg/compute/tasks/vpc_create_task.go +++ b/pkg/compute/tasks/vpc_create_task.go @@ -18,7 +18,6 @@ import ( "context" "yunion.io/x/jsonutils" - "yunion.io/x/log" "yunion.io/x/pkg/errors" api "yunion.io/x/onecloud/pkg/apis/compute" @@ -36,12 +35,11 @@ func init() { taskman.RegisterTask(VpcCreateTask{}) } -func (self *VpcCreateTask) TaskFailed(ctx context.Context, vpc *models.SVpc, err jsonutils.JSONObject) { - log.Errorf("vpc create task fail: %s", err) - vpc.SetStatus(self.UserCred, api.VPC_STATUS_FAILED, err.String()) +func (self *VpcCreateTask) taskFailed(ctx context.Context, vpc *models.SVpc, err error) { + vpc.SetStatus(self.UserCred, api.VPC_STATUS_FAILED, err.Error()) db.OpsLog.LogEvent(vpc, db.ACT_ALLOCATE_FAIL, err, self.UserCred) logclient.AddActionLogWithStartable(self, vpc, logclient.ACT_ALLOCATE, err, self.UserCred, false) - self.SetStageFailed(ctx, err) + self.SetStageFailed(ctx, jsonutils.NewString(err.Error())) } func (self *VpcCreateTask) OnInit(ctx context.Context, obj db.IStandaloneModel, body jsonutils.JSONObject) { @@ -50,13 +48,13 @@ func (self *VpcCreateTask) OnInit(ctx context.Context, obj db.IStandaloneModel, region, err := vpc.GetRegion() if err != nil { - self.TaskFailed(ctx, vpc, jsonutils.NewString(errors.Wrap(err, "vpc.GetRegion").Error())) + self.taskFailed(ctx, vpc, errors.Wrapf(err, "GetRegion")) return } self.SetStage("OnCreateVpcComplete", nil) err = region.GetDriver().RequestCreateVpc(ctx, self.UserCred, region, vpc, self) if err != nil { - self.TaskFailed(ctx, vpc, jsonutils.NewString(err.Error())) + self.taskFailed(ctx, vpc, errors.Wrapf(err, "RequestCreateVpc")) return } } @@ -67,5 +65,5 @@ func (self *VpcCreateTask) OnCreateVpcComplete(ctx context.Context, vpc *models. } func (self *VpcCreateTask) OnCreateVpcCompleteFailed(ctx context.Context, vpc *models.SVpc, data jsonutils.JSONObject) { - self.TaskFailed(ctx, vpc, data) + self.taskFailed(ctx, vpc, errors.Errorf(data.String())) } diff --git a/pkg/multicloud/cloudpods/host.go b/pkg/multicloud/cloudpods/host.go index b954433fbc..69b2981f75 100644 --- a/pkg/multicloud/cloudpods/host.go +++ b/pkg/multicloud/cloudpods/host.go @@ -52,6 +52,10 @@ func (self *SHost) GetAccessIp() string { return self.AccessIp } +func (self *SHost) GetOvnVersion() string { + return self.OvnVersion +} + func (self *SHost) Refresh() error { host, err := self.zone.region.GetHost(self.Id) if err != nil { diff --git a/pkg/multicloud/host_base.go b/pkg/multicloud/host_base.go index c5d72f56cb..b0594f2a61 100644 --- a/pkg/multicloud/host_base.go +++ b/pkg/multicloud/host_base.go @@ -34,3 +34,7 @@ func (self *SHostBase) GetReservedMemoryMb() int { func (self *SHostBase) GetSchedtags() ([]string, error) { return nil, nil } + +func (self *SHostBase) GetOvnVersion() string { + return "" +} diff --git a/pkg/scheduler/algorithm/predicates/network_predicate.go b/pkg/scheduler/algorithm/predicates/network_predicate.go index d5ae981141..711da3cf0a 100644 --- a/pkg/scheduler/algorithm/predicates/network_predicate.go +++ b/pkg/scheduler/algorithm/predicates/network_predicate.go @@ -111,7 +111,7 @@ func IsNetworksAvailable(c core.Candidater, data *api.SchedInfo, req *computeapi ovnNetworks := []*api.CandidateNetwork{} for i := len(networks) - 1; i >= 0; i -= 1 { net := networks[i] - if net.Provider == computeapi.CLOUD_PROVIDER_ONECLOUD { + if net.Provider == computeapi.CLOUD_PROVIDER_ONECLOUD || net.Provider == computeapi.CLOUD_PROVIDER_CLOUDPODS { networks = append(networks[:i], networks[i+1:]...) ovnNetworks = append(ovnNetworks, net) }