mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
fix(region): cloudpods network scheduler
This commit is contained in:
@@ -306,6 +306,8 @@ type ICloudHost interface {
|
||||
GetIHostNics() ([]ICloudHostNetInterface, error)
|
||||
|
||||
GetSchedtags() ([]string, error)
|
||||
|
||||
GetOvnVersion() string // just for cloudpods host
|
||||
}
|
||||
|
||||
type ICloudVM interface {
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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()))
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -34,3 +34,7 @@ func (self *SHostBase) GetReservedMemoryMb() int {
|
||||
func (self *SHostBase) GetSchedtags() ([]string, error) {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
func (self *SHostBase) GetOvnVersion() string {
|
||||
return ""
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user