fix(region): support kube cluster create

This commit is contained in:
ioito
2023-05-12 10:15:46 +08:00
parent 1ee6cb3f2e
commit 336f92f111
42 changed files with 2041 additions and 32 deletions
+1
View File
@@ -23,6 +23,7 @@ import (
func init() {
cmd := shell.NewResourceCmd(&modules.KubeClusters)
cmd.List(&compute.KubeClusterListOptions{})
cmd.Create(&compute.KubeClusterCreateOptions{})
cmd.Show(&compute.KubeClusterIdOption{})
cmd.Delete(&compute.KubeClusterIdOption{})
cmd.Perform("syncstatus", &compute.KubeClusterIdOption{})
@@ -23,6 +23,8 @@ import (
func init() {
cmd := shell.NewResourceCmd(&modules.KubeNodePools)
cmd.List(&compute.KubeNodePoolListOptions{})
cmd.Create(&compute.KubeNodePoolCreateOptions{})
cmd.Show(&compute.KubeNodePoolIdOption{})
cmd.Delete(&compute.KubeNodePoolIdOption{})
cmd.Perform("syncstatus", &compute.KubeNodePoolIdOption{})
}
+2 -2
View File
@@ -84,14 +84,14 @@ require (
k8s.io/client-go v0.19.3
k8s.io/cluster-bootstrap v0.19.3
moul.io/http2curl/v2 v2.3.0
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230515072146-3a7171baa793
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230519063348-6f0a3a041cdb
yunion.io/x/executor v0.0.0-20211018100936-39a2cd966656
yunion.io/x/jsonutils v1.0.1-0.20230428104347-7c2fdff8e8e7
yunion.io/x/log v1.0.1-0.20230411060016-feb3f46ab361
yunion.io/x/ovsdb v0.0.0-20230306173834-f164f413a900
yunion.io/x/pkg v1.0.1-0.20230504073602-0a74096f836a
yunion.io/x/s3cli v0.0.0-20190917004522-13ac36d8687e
yunion.io/x/sqlchemy v1.1.2-0.20230512065832-5323af46107e
yunion.io/x/sqlchemy v1.1.2-0.20230518004216-245cb3766e38
yunion.io/x/structarg v0.0.0-20220312084958-9c6c79c7d1c6
)
+4 -4
View File
@@ -1185,8 +1185,8 @@ sigs.k8s.io/structured-merge-diff/v4 v4.0.1/go.mod h1:bJZC9H9iH24zzfZ/41RGcq60oK
sigs.k8s.io/yaml v1.1.0/go.mod h1:UJmg0vDUVViEyp3mgSv9WPwZCDxu4rQW1olrI1uml+o=
sigs.k8s.io/yaml v1.2.0 h1:kr/MCeFWJWTwyaHoR9c8EjH9OumOmoF9YGiZd7lFm/Q=
sigs.k8s.io/yaml v1.2.0/go.mod h1:yfXDCHCao9+ENCvLSE62v9VSji2MKu5jeNfTrofGhJc=
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230515072146-3a7171baa793 h1:MlMEK/RorFch+yizvlV5QHVevoKLVSxuJc/XhoGOO2o=
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230515072146-3a7171baa793/go.mod h1:crMeQeaNaZefTXfXbQkoj5SStggqkSNVABHtYBFjM3Y=
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230519063348-6f0a3a041cdb h1:nICglDMht07iM2IFOGud6mcRk6QbVH7mi2r3gyuAylo=
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230519063348-6f0a3a041cdb/go.mod h1:crMeQeaNaZefTXfXbQkoj5SStggqkSNVABHtYBFjM3Y=
yunion.io/x/executor v0.0.0-20211018100936-39a2cd966656 h1:0zlZD5uhZoIHgLVAWCz2aHaYk2ZrNsACCYD7R6EIBII=
yunion.io/x/executor v0.0.0-20211018100936-39a2cd966656/go.mod h1:Uxuou9WQIeJXNpy7t2fPLL0BYLvLiMvGQwY7Qc6aSws=
yunion.io/x/jsonutils v0.0.0-20190625054549-a964e1e8a051/go.mod h1:4N0/RVzsYL3kH3WE/H1BjUQdFiWu50JGCFQuuy+Z634=
@@ -1209,7 +1209,7 @@ yunion.io/x/pkg v1.0.1-0.20230504073602-0a74096f836a/go.mod h1:ksCJVQ+DwKrJ5QBEo
yunion.io/x/s3cli v0.0.0-20190917004522-13ac36d8687e h1:v+EzIadodSwkdZ/7bremd7J8J50Cise/HCylsOJngmo=
yunion.io/x/s3cli v0.0.0-20190917004522-13ac36d8687e/go.mod h1:0iFKpOs1y4lbCxeOmq3Xx/0AcQoewVPwj62eRluioEo=
yunion.io/x/sqlchemy v1.0.1/go.mod h1:FTdwPdGhMgh4E+UFXc9klI1Ok34fMuybTT+jLhOaIjI=
yunion.io/x/sqlchemy v1.1.2-0.20230512065832-5323af46107e h1:BZhnpXqJJbOuEGoXBb3uvsRz3t6dbJyG1otGqxbFtc4=
yunion.io/x/sqlchemy v1.1.2-0.20230512065832-5323af46107e/go.mod h1:uuPVZEyEq3sWd5vf9VjGSy6lZzof22X87OEHw9sddJQ=
yunion.io/x/sqlchemy v1.1.2-0.20230518004216-245cb3766e38 h1:cRAYmZsSIeoh+WsihZH7Cee/qst/A8KL53JUtW7pxcA=
yunion.io/x/sqlchemy v1.1.2-0.20230518004216-245cb3766e38/go.mod h1:uuPVZEyEq3sWd5vf9VjGSy6lZzof22X87OEHw9sddJQ=
yunion.io/x/structarg v0.0.0-20220312084958-9c6c79c7d1c6 h1:WuWXhY3DvhdRTzWCJ/kwt3Ss6KIq7+KqJwb+esvNGwU=
yunion.io/x/structarg v0.0.0-20220312084958-9c6c79c7d1c6/go.mod h1:EP6NSv2C0zzqBDTKumv8hPWLb3XvgMZDHQRfyuOrQng=
+50
View File
@@ -15,6 +15,11 @@
package compute
import (
"reflect"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/gotypes"
"yunion.io/x/onecloud/pkg/apis"
)
@@ -38,9 +43,23 @@ type KubeClusterListInput struct {
RegionalFilterListInput
ManagedResourceListInput
VpcFilterListInput
}
type KubeClusterCreateInput struct {
apis.EnabledStatusInfrasResourceBaseCreateInput
Version string `json:"version"`
// required: true
NetworkIds SKubeNetworkIds `json:"network_ids"`
// swagger:ignore
ManagerId string `json:"manager_id"`
// swagger:ignore
CloudregionId string `json:"cloudregion_id"`
// required: true
VpcResourceInput
RoleName string `json:"role_name"`
}
type KubeClusterDetails struct {
@@ -49,6 +68,7 @@ type KubeClusterDetails struct {
SKubeCluster
ManagedResourceInfo
CloudregionResourceInfo
VpcResourceInfo
}
func (self KubeClusterDetails) GetMetricTags() map[string]string {
@@ -93,3 +113,33 @@ type KubeClusterDeleteInput struct {
// default: false
Retain bool `json:"retain"`
}
type SInstanceTypes []string
func (kn SInstanceTypes) String() string {
return jsonutils.Marshal(kn).String()
}
func (kn SInstanceTypes) IsZero() bool {
return len(kn) == 0
}
type SKubeNetworkIds []string
func (kn SKubeNetworkIds) String() string {
return jsonutils.Marshal(kn).String()
}
func (kn SKubeNetworkIds) IsZero() bool {
return len(kn) == 0
}
func init() {
gotypes.RegisterSerializable(reflect.TypeOf(&SKubeNetworkIds{}), func() gotypes.ISerializable {
return &SKubeNetworkIds{}
})
gotypes.RegisterSerializable(reflect.TypeOf(&SInstanceTypes{}), func() gotypes.ISerializable {
return &SInstanceTypes{}
})
}
+16
View File
@@ -25,6 +25,22 @@ type KubeNodePoolListInput struct {
}
type KubeNodePoolCreateInput struct {
apis.StatusStandaloneResourceCreateInput
NetworkIds SKubeNetworkIds `json:"network_ids"`
InstanceTypes SInstanceTypes `json:"instance_types"`
// default: 2
MinInstanceCount int `json:"min_instance_count"`
// default: 2
MaxInstanceCount int `json:"max_instance_count"`
// 预期节点数量, 不得小于min_instance_count
// default: 2
DesiredInstanceCount int `json:"desired_instance_count"`
RootDiskSizeGb int `json:"root_disk_size_gb"`
CloudKubeClusterId string `json:"cloud_kube_cluster_id"`
}
type KubeNodePoolDetails struct {
+1 -1
View File
@@ -829,7 +829,7 @@ func NewIPv4AddrValidator(key string) *ValidatorIPv4Addr {
var ValidateModel = func(userCred mcclient.TokenCredential, manager db.IStandaloneModelManager, id *string) (db.IModel, error) {
if len(*id) == 0 {
return nil, httperrors.NewMissingParameterError(manager.Keyword())
return nil, httperrors.NewMissingParameterError(manager.Keyword() + "_id")
}
model, err := manager.FetchByIdOrName(userCred, *id)
+168 -4
View File
@@ -33,6 +33,7 @@ import (
"yunion.io/x/onecloud/pkg/cloudcommon/db/lockman"
"yunion.io/x/onecloud/pkg/cloudcommon/db/quotas"
"yunion.io/x/onecloud/pkg/cloudcommon/db/taskman"
"yunion.io/x/onecloud/pkg/cloudcommon/validators"
"yunion.io/x/onecloud/pkg/compute/options"
"yunion.io/x/onecloud/pkg/httperrors"
"yunion.io/x/onecloud/pkg/mcclient"
@@ -45,6 +46,7 @@ type SKubeClusterManager struct {
db.SEnabledStatusInfrasResourceBaseManager
db.SExternalizedResourceBaseManager
SManagedResourceBaseManager
SVpcResourceBaseManager
SCloudregionResourceBaseManager
}
@@ -65,11 +67,20 @@ func init() {
type SKubeCluster struct {
db.SEnabledStatusInfrasResourceBase
db.SExternalizedResourceBase
SManagedResourceBase
SVpcResourceBase `wdith:"36" charset:"ascii" nullable:"false" list:"domain" create:"domain_required"`
SCloudregionResourceBase `width:"36" charset:"ascii" nullable:"false" list:"domain" create:"domain_required" default:"default"`
Version string `width:"12" charset:"utf8" nullable:"false" list:"admin" create:"domain_optional"`
// 本地KubeserverId
ExternalClusterId string `width:"36" charset:"ascii" nullable:"false" list:"admin"`
NetworkIds *api.SKubeNetworkIds `list:"user" update:"user" create:"required"`
}
func (self *SKubeCluster) GetCloudproviderId() string {
return self.ManagerId
}
func (manager *SKubeClusterManager) GetContextManagers() [][]db.IModelManager {
@@ -97,6 +108,15 @@ func (self *SKubeCluster) GetRegion() (*SCloudregion, error) {
return region.(*SCloudregion), nil
}
func (self *SKubeCluster) GetNetworks() ([]SNetwork, error) {
networks := []SNetwork{}
if self.NetworkIds == nil {
return networks, nil
}
q := NetworkManager.Query().In("id", *self.NetworkIds)
return networks, db.FetchModelObjects(NetworkManager, q, &networks)
}
func (manager *SKubeClusterManager) FetchCustomizeColumns(
ctx context.Context,
userCred mcclient.TokenCredential,
@@ -108,12 +128,14 @@ func (manager *SKubeClusterManager) FetchCustomizeColumns(
rows := make([]api.KubeClusterDetails, len(objs))
stdRows := manager.SEnabledStatusInfrasResourceBaseManager.FetchCustomizeColumns(ctx, userCred, query, objs, fields, isList)
managerRows := manager.SManagedResourceBaseManager.FetchCustomizeColumns(ctx, userCred, query, objs, fields, isList)
vpcRows := manager.SVpcResourceBaseManager.FetchCustomizeColumns(ctx, userCred, query, objs, fields, isList)
regionRows := manager.SCloudregionResourceBaseManager.FetchCustomizeColumns(ctx, userCred, query, objs, fields, isList)
for i := range rows {
rows[i] = api.KubeClusterDetails{
EnabledStatusInfrasResourceBaseDetails: stdRows[i],
ManagedResourceInfo: managerRows[i],
CloudregionResourceInfo: regionRows[i],
VpcResourceInfo: vpcRows[i],
}
}
return rows
@@ -298,6 +320,40 @@ func (self *SKubeCluster) SyncAllWithCloudKubeCluster(ctx context.Context, userC
func (self *SKubeCluster) SyncWithCloudKubeCluster(ctx context.Context, userCred mcclient.TokenCredential, ext cloudprovider.ICloudKubeCluster, provider *SCloudprovider) error {
diff, err := db.UpdateWithLock(ctx, self, func() error {
self.Status = ext.GetStatus()
if version := ext.GetVersion(); len(version) > 0 {
self.Version = ext.GetVersion()
}
if self.NetworkIds != nil && len(*self.NetworkIds) > 0 {
networkIds := ext.GetNetworkIds()
netIds := api.SKubeNetworkIds{}
for i := range networkIds {
netObj, err := db.FetchByExternalIdAndManagerId(NetworkManager, networkIds[i], func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
wires := WireManager.Query().SubQuery()
vpcs := VpcManager.Query().SubQuery()
return q.Join(wires, sqlchemy.Equals(wires.Field("id"), q.Field("wire_id"))).
Join(vpcs, sqlchemy.Equals(vpcs.Field("id"), wires.Field("vpc_id"))).
Filter(sqlchemy.Equals(vpcs.Field("manager_id"), self.ManagerId))
})
if err != nil {
break
}
netIds = append(netIds, netObj.GetId())
}
if len(networkIds) == len(netIds) && len(netIds) > 0 {
self.NetworkIds = &netIds
}
}
if vpcId := ext.GetVpcId(); len(vpcId) > 0 && len(self.VpcId) == 0 {
vpcObj, _ := db.FetchByExternalIdAndManagerId(VpcManager, vpcId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", provider.Id)
})
if vpcObj != nil {
self.ManagerId = vpcObj.GetId()
}
}
return nil
})
if err != nil {
@@ -323,8 +379,37 @@ func (self *SCloudregion) newFromCloudKubeCluster(ctx context.Context, userCred
cluster.ManagerId = provider.Id
cluster.ExternalId = ext.GetGlobalId()
cluster.Enabled = tristate.True
cluster.Version = ext.GetVersion()
cluster.Status = ext.GetStatus()
networkIds := ext.GetNetworkIds()
netIds := api.SKubeNetworkIds{}
for i := range networkIds {
netObj, err := db.FetchByExternalIdAndManagerId(NetworkManager, networkIds[i], func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
wires := WireManager.Query().SubQuery()
vpcs := VpcManager.Query().SubQuery()
return q.Join(wires, sqlchemy.Equals(wires.Field("id"), q.Field("wire_id"))).
Join(vpcs, sqlchemy.Equals(vpcs.Field("id"), wires.Field("vpc_id"))).
Filter(sqlchemy.Equals(vpcs.Field("manager_id"), cluster.ManagerId))
})
if err != nil {
break
}
netIds = append(netIds, netObj.GetId())
}
if len(networkIds) == len(netIds) && len(netIds) > 0 {
cluster.NetworkIds = &netIds
}
if vpcId := ext.GetVpcId(); len(vpcId) > 0 {
vpcObj, _ := db.FetchByExternalIdAndManagerId(VpcManager, vpcId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", provider.Id)
})
if vpcObj != nil {
cluster.ManagerId = vpcObj.GetId()
}
}
var err = func() error {
lockman.LockRawObject(ctx, KubeClusterManager.Keyword(), "name")
defer lockman.ReleaseRawObject(ctx, KubeClusterManager.Keyword(), "name")
@@ -359,8 +444,67 @@ func (manager *SKubeClusterManager) ValidateCreateData(
ownerId mcclient.IIdentityProvider,
query jsonutils.JSONObject,
input api.KubeClusterCreateInput,
) (api.KubeClusterCreateInput, error) {
return input, httperrors.NewNotImplementedError("Not Implemented")
) (*api.KubeClusterCreateInput, error) {
var err error
input.EnabledStatusInfrasResourceBaseCreateInput, err = manager.SEnabledStatusInfrasResourceBaseManager.ValidateCreateData(ctx, userCred, ownerId, query, input.EnabledStatusInfrasResourceBaseCreateInput)
if err != nil {
return nil, err
}
if input.Enabled == nil && input.Disabled == nil {
enabled := true
input.Enabled = &enabled
}
if len(input.VpcId) == 0 {
return nil, httperrors.NewMissingParameterError("vpc_id")
}
vpcObj, err := validators.ValidateModel(userCred, VpcManager, &input.VpcId)
if err != nil {
return nil, err
}
vpc := vpcObj.(*SVpc)
input.CloudregionId = vpc.CloudregionId
input.ManagerId = vpc.ManagerId
networks, err := vpc.GetNetworks()
if err != nil {
return nil, errors.Wrapf(err, "GetNetworks")
}
nets := map[string]bool{}
for _, net := range networks {
nets[net.Id] = true
}
if input.NetworkIds == nil {
return nil, httperrors.NewMissingParameterError("network_ids")
}
for i := range input.NetworkIds {
_, err = validators.ValidateModel(userCred, NetworkManager, &input.NetworkIds[i])
if err != nil {
return nil, err
}
_, ok := nets[input.NetworkIds[i]]
if !ok {
return nil, httperrors.NewInputParameterError("network %s not belong to vpc %s", input.NetworkIds[i], vpc.Name)
}
}
region, err := vpc.GetRegion()
if err != nil {
return nil, errors.Wrapf(err, "GetRegion")
}
return region.GetDriver().ValidateCreateKubeClusterData(ctx, userCred, ownerId, &input)
}
func (self *SKubeCluster) PostCreate(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, data jsonutils.JSONObject) {
self.SEnabledStatusInfrasResourceBase.PostCreate(ctx, userCred, ownerId, query, data)
self.StartKubeClusterCreateTask(ctx, userCred, data)
}
func (self *SKubeCluster) StartKubeClusterCreateTask(ctx context.Context, userCred mcclient.TokenCredential, data jsonutils.JSONObject) error {
params := data.(*jsonutils.JSONDict)
task, err := taskman.TaskManager.NewTask(ctx, "KubeClusterCreateTask", self, userCred, params, "", "", nil)
if err != nil {
return errors.Wrapf(err, "NewTask")
}
self.SetStatus(userCred, api.KUBE_CLUSTER_STATUS_CREATING, "")
return task.ScheduleRun(nil)
}
func (self *SKubeCluster) GetIRegion(ctx context.Context) (cloudprovider.ICloudRegion, error) {
@@ -489,6 +633,11 @@ func (manager *SKubeClusterManager) ListItemFilter(
return nil, errors.Wrap(err, "SCloudregionResourceBaseManager.ListItemFilter")
}
q, err = manager.SVpcResourceBaseManager.ListItemFilter(ctx, q, userCred, query.VpcFilterListInput)
if err != nil {
return nil, errors.Wrap(err, "SVpcResourceBaseManager.ListItemFilter")
}
return q, nil
}
@@ -514,6 +663,11 @@ func (manager *SKubeClusterManager) QueryDistinctExtraField(q *sqlchemy.SQuery,
return q, nil
}
q, err = manager.SVpcResourceBaseManager.QueryDistinctExtraField(q, field)
if err == nil {
return q, nil
}
}
return q, httperrors.ErrNotFound
}
@@ -536,6 +690,10 @@ func (manager *SKubeClusterManager) OrderByExtraFields(
if err != nil {
return nil, errors.Wrap(err, "SCloudregionResourceBaseManager.OrderByExtraFields")
}
q, err = manager.SVpcResourceBaseManager.OrderByExtraFields(ctx, q, userCred, query.VpcFilterListInput)
if err != nil {
return nil, errors.Wrap(err, "SCloudregionResourceBaseManager.OrderByExtraFields")
}
return q, nil
}
@@ -547,7 +705,7 @@ func (self *SKubeCluster) PerformSyncstatus(ctx context.Context, userCred mcclie
func (cluster *SKubeCluster) GetQuotaKeys() quotas.SDomainRegionalCloudResourceKeys {
region, _ := cluster.GetRegion()
manager := cluster.GetCloudprovider()
manager := cluster.SManagedResourceBase.GetCloudprovider()
ownerId := cluster.GetOwnerId()
regionKeys := fetchRegionalQuotaKeys(rbacscope.ScopeDomain, ownerId, region, manager)
keys := quotas.SDomainRegionalCloudResourceKeys{}
@@ -592,6 +750,12 @@ func (manager *SKubeClusterManager) ListItemExportKeys(ctx context.Context,
return nil, errors.Wrap(err, "SCloudregionResourceBaseManager.ListItemExportKeys")
}
}
if keys.ContainsAny(manager.SVpcResourceBaseManager.GetExportKeys()...) {
q, err = manager.SVpcResourceBaseManager.ListItemExportKeys(ctx, q, userCred, keys)
if err != nil {
return nil, errors.Wrap(err, "SVpcResourceBaseManager.ListItemExportKeys")
}
}
if keys.ContainsAny(manager.SManagedResourceBaseManager.GetExportKeys()...) {
q, err = manager.SManagedResourceBaseManager.ListItemExportKeys(ctx, q, userCred, keys)
if err != nil {
+173 -3
View File
@@ -30,6 +30,7 @@ import (
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/cloudcommon/db/lockman"
"yunion.io/x/onecloud/pkg/cloudcommon/db/taskman"
"yunion.io/x/onecloud/pkg/cloudcommon/validators"
"yunion.io/x/onecloud/pkg/httperrors"
"yunion.io/x/onecloud/pkg/mcclient"
"yunion.io/x/onecloud/pkg/util/stringutils2"
@@ -58,6 +59,15 @@ type SKubeNodePool struct {
db.SStatusStandaloneResourceBase
db.SExternalizedResourceBase
NetworkIds *api.SKubeNetworkIds `list:"user" update:"user" create:"required"`
InstanceTypes *api.SInstanceTypes `list:"user" update:"user" create:"required"`
MinInstanceCount int `nullable:"false" list:"user" create:"optional" default:"2"`
MaxInstanceCount int `nullable:"false" list:"user" create:"optional" default:"2"`
DesiredInstanceCount int `nullable:"false" list:"user" create:"optional" default:"2"`
RootDiskSizeGb int `nullable:"false" list:"user" create:"optional" default:"100"`
CloudKubeClusterId string `width:"36" charset:"ascii" name:"cloud_kube_cluster_id" nullable:"false" list:"user" create:"required" index:"true"`
}
@@ -87,6 +97,14 @@ func (self *SKubeNodePool) GetKubeCluster() (*SKubeCluster, error) {
return cluster.(*SKubeCluster), nil
}
func (self *SKubeNodePool) GetRegion() (*SCloudregion, error) {
cluster, err := self.GetKubeCluster()
if err != nil {
return nil, errors.Wrapf(err, "GetKubeCluster")
}
return cluster.GetRegion()
}
func (self *SKubeNodePool) GetOwnerId() mcclient.IIdentityProvider {
cluster, err := self.GetKubeCluster()
if err != nil {
@@ -96,6 +114,40 @@ func (self *SKubeNodePool) GetOwnerId() mcclient.IIdentityProvider {
return cluster.GetOwnerId()
}
func (self *SKubeNodePool) GetNetworks() ([]SNetwork, error) {
ret := []SNetwork{}
q := NetworkManager.Query().In("id", self.NetworkIds)
return ret, db.FetchModelObjects(NetworkManager, q, &ret)
}
func (self *SKubeNodePool) GetIKubeCluster(ctx context.Context) (cloudprovider.ICloudKubeCluster, error) {
cluster, err := self.GetKubeCluster()
if err != nil {
return nil, err
}
return cluster.GetIKubeCluster(ctx)
}
func (self *SKubeNodePool) GetIKubeNodePool(ctx context.Context) (cloudprovider.ICloudKubeNodePool, error) {
if len(self.ExternalId) == 0 {
return nil, errors.Wrapf(cloudprovider.ErrNotFound, "empty external id")
}
cluster, err := self.GetIKubeCluster(ctx)
if err != nil {
return nil, errors.Wrapf(err, "GetIKubeCluster")
}
pools, err := cluster.GetIKubeNodePools()
if err != nil {
return nil, err
}
for i := range pools {
if pools[i].GetGlobalId() == self.ExternalId {
return pools[i], nil
}
}
return nil, errors.Wrapf(cloudprovider.ErrNotFound, self.ExternalId)
}
func (manager *SKubeNodePoolManager) FetchOwnerId(ctx context.Context, data jsonutils.JSONObject) (mcclient.IIdentityProvider, error) {
info := struct{ KubeClusterId string }{}
data.Unmarshal(&info)
@@ -121,6 +173,11 @@ func (manager *SKubeNodePoolManager) FilterByOwner(q *sqlchemy.SQuery, man db.Fi
return q
}
// 同步Kube Node Pool 状态
func (self *SKubeNodePool) PerformSyncstatus(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input apis.SyncstatusInput) (jsonutils.JSONObject, error) {
return nil, StartResourceSyncStatusTask(ctx, userCred, self, "KubeNodePoolSyncstatusTask", "")
}
func (self *SKubeNodePool) ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.KubeNodePoolUpdateInput) (api.KubeNodePoolUpdateInput, error) {
var err error
input.StatusStandaloneResourceBaseUpdateInput, err = self.SStatusStandaloneResourceBase.ValidateUpdateData(ctx, userCred, query, input.StatusStandaloneResourceBaseUpdateInput)
@@ -234,8 +291,49 @@ func (manager *SKubeNodePoolManager) FilterByUniqValues(q *sqlchemy.SQuery, valu
return q
}
func (manager *SKubeNodePoolManager) ValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, input api.KubeNodePoolCreateInput) (api.KubeNodePoolCreateInput, error) {
return input, httperrors.NewNotImplementedError("Not Implemented")
func (manager *SKubeNodePoolManager) ValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, input api.KubeNodePoolCreateInput) (*api.KubeNodePoolCreateInput, error) {
var err error
input.StatusStandaloneResourceCreateInput, err = manager.SStatusStandaloneResourceBaseManager.ValidateCreateData(ctx, userCred, ownerId, query, input.StatusStandaloneResourceCreateInput)
if err != nil {
return nil, err
}
clusterObj, err := validators.ValidateModel(userCred, KubeClusterManager, &input.CloudKubeClusterId)
if err != nil {
return nil, err
}
cluster := clusterObj.(*SKubeCluster)
for i := range input.NetworkIds {
_, err = validators.ValidateModel(userCred, NetworkManager, &input.NetworkIds[i])
if err != nil {
return nil, err
}
}
if len(input.NetworkIds) == 0 {
return nil, httperrors.NewMissingParameterError("network_ids")
}
if len(input.InstanceTypes) == 0 {
return nil, httperrors.NewMissingParameterError("instance_types")
}
region, err := cluster.GetRegion()
if err != nil {
return nil, err
}
return region.GetDriver().ValidateCreateKubeNodePoolData(ctx, userCred, ownerId, &input)
}
func (self *SKubeNodePool) PostCreate(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, data jsonutils.JSONObject) {
self.SStatusStandaloneResourceBase.PostCreate(ctx, userCred, ownerId, query, data)
self.StartKubeNodePoolCreateTask(ctx, userCred, data)
}
func (self *SKubeNodePool) StartKubeNodePoolCreateTask(ctx context.Context, userCred mcclient.TokenCredential, data jsonutils.JSONObject) error {
params := data.(*jsonutils.JSONDict)
task, err := taskman.TaskManager.NewTask(ctx, "KubeNodePoolCreateTask", self, userCred, params, "", "", nil)
if err != nil {
return errors.Wrapf(err, "NewTask")
}
self.SetStatus(userCred, apis.STATUS_CREATING, "")
return task.ScheduleRun(nil)
}
func (self *SKubeCluster) SyncKubeNodePools(ctx context.Context, userCred mcclient.TokenCredential, exts []cloudprovider.ICloudKubeNodePool) compare.SyncResult {
@@ -290,8 +388,51 @@ func (self *SKubeCluster) SyncKubeNodePools(ctx context.Context, userCred mcclie
}
func (self *SKubeNodePool) SyncWithCloudKubeNodePool(ctx context.Context, userCred mcclient.TokenCredential, ext cloudprovider.ICloudKubeNodePool) error {
_, err := db.UpdateWithLock(ctx, self, func() error {
cluster, err := self.GetKubeCluster()
if err != nil {
return errors.Wrapf(err, "GetKubeCluster")
}
_, err = db.UpdateWithLock(ctx, self, func() error {
self.Status = ext.GetStatus()
networkIds := ext.GetNetworkIds()
netIds := api.SKubeNetworkIds{}
for i := range networkIds {
netObj, err := db.FetchByExternalIdAndManagerId(NetworkManager, networkIds[i], func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
wires := WireManager.Query().SubQuery()
vpcs := VpcManager.Query().SubQuery()
return q.Join(wires, sqlchemy.Equals(wires.Field("id"), q.Field("wire_id"))).
Join(vpcs, sqlchemy.Equals(vpcs.Field("id"), wires.Field("vpc_id"))).
Filter(sqlchemy.Equals(vpcs.Field("manager_id"), cluster.ManagerId))
})
if err != nil {
break
}
netIds = append(netIds, netObj.GetId())
}
if len(networkIds) == len(netIds) && len(netIds) > 0 {
self.NetworkIds = &netIds
}
instanceTypes := api.SInstanceTypes{}
for _, instanceType := range ext.GetInstanceTypes() {
instanceTypes = append(instanceTypes, instanceType)
}
if len(instanceTypes) > 0 {
self.InstanceTypes = &instanceTypes
}
if minSize := ext.GetMinInstanceCount(); minSize > 0 {
self.MinInstanceCount = minSize
}
if maxSize := ext.GetMaxInstanceCount(); maxSize > 0 {
self.MaxInstanceCount = maxSize
}
if desiredSize := ext.GetDesiredInstanceCount(); desiredSize > 0 {
self.DesiredInstanceCount = desiredSize
}
if rootSize := ext.GetRootDiskSizeGb(); rootSize > 0 {
self.RootDiskSizeGb = rootSize
}
return nil
})
if err != nil {
@@ -312,6 +453,35 @@ func (self *SKubeCluster) newFromCloudKubeNodePool(ctx context.Context, userCred
pool.CloudKubeClusterId = self.Id
pool.ExternalId = ext.GetGlobalId()
networkIds := ext.GetNetworkIds()
netIds := api.SKubeNetworkIds{}
for i := range networkIds {
netObj, err := db.FetchByExternalIdAndManagerId(NetworkManager, networkIds[i], func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
wires := WireManager.Query().SubQuery()
vpcs := VpcManager.Query().SubQuery()
return q.Join(wires, sqlchemy.Equals(wires.Field("id"), q.Field("wire_id"))).
Join(vpcs, sqlchemy.Equals(vpcs.Field("id"), wires.Field("vpc_id"))).
Filter(sqlchemy.Equals(vpcs.Field("manager_id"), self.ManagerId))
})
if err != nil {
break
}
netIds = append(netIds, netObj.GetId())
}
if len(netIds) > 0 {
pool.NetworkIds = &netIds
}
pool.MinInstanceCount = ext.GetMinInstanceCount()
pool.MaxInstanceCount = ext.GetMaxInstanceCount()
pool.DesiredInstanceCount = ext.GetDesiredInstanceCount()
pool.RootDiskSizeGb = ext.GetRootDiskSizeGb()
instanceTypes := api.SInstanceTypes{}
for _, instanceType := range ext.GetInstanceTypes() {
instanceTypes = append(instanceTypes, instanceType)
}
pool.InstanceTypes = &instanceTypes
err := KubeNodePoolManager.TableSpec().Insert(ctx, &pool)
if err != nil {
return nil, errors.Wrapf(err, "Insert")
+10 -2
View File
@@ -39,7 +39,8 @@ type IRegionDriver interface {
IElasticcacheBackup
IDBInstanceDriver
IElasticSearchDriver
IKafkaDrive
IKafkaDriver
IKubeClusterDriver
ValidateCreateLoadbalancerData(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, data *api.LoadbalancerCreateInput) (*api.LoadbalancerCreateInput, error)
RequestCreateLoadbalancerInstance(ctx context.Context, userCred mcclient.TokenCredential, lb *SLoadbalancer, input *api.LoadbalancerCreateInput, task taskman.ITask) error
@@ -265,10 +266,17 @@ type IElasticSearchDriver interface {
RequestRemoteUpdateElasticSearch(ctx context.Context, userCred mcclient.TokenCredential, elasticcache *SElasticSearch, replaceTags bool, task taskman.ITask) error
}
type IKafkaDrive interface {
type IKafkaDriver interface {
RequestRemoteUpdateKafka(ctx context.Context, userCred mcclient.TokenCredential, kafka *SKafka, replaceTags bool, task taskman.ITask) error
}
type IKubeClusterDriver interface {
ValidateCreateKubeClusterData(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, input *api.KubeClusterCreateInput) (*api.KubeClusterCreateInput, error)
ValidateCreateKubeNodePoolData(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, input *api.KubeNodePoolCreateInput) (*api.KubeNodePoolCreateInput, error)
RequestCreateKubeCluster(ctx context.Context, userCred mcclient.TokenCredential, cluster *SKubeCluster, task taskman.ITask) error
RequestCreateKubeNodePool(ctx context.Context, userCred mcclient.TokenCredential, pool *SKubeNodePool, task taskman.ITask) error
}
var regionDrivers map[string]IRegionDriver
func init() {
+17
View File
@@ -481,6 +481,7 @@ func (self *SBaseRegionDriver) RequestSyncBackupStorageStatus(ctx context.Contex
func (self *SBaseRegionDriver) RequestPackInstanceBackup(ctx context.Context, ib *models.SInstanceBackup, task taskman.ITask, packageName string) error {
return errors.Wrapf(cloudprovider.ErrNotImplemented, "RequestPackInstanceBackup")
}
func (self *SBaseRegionDriver) RequestUnpackInstanceBackup(ctx context.Context, ib *models.SInstanceBackup, task taskman.ITask, packageName string, metadataOnly bool) error {
return errors.Wrapf(cloudprovider.ErrNotImplemented, "RequestUnpackInstanceBackup")
}
@@ -492,3 +493,19 @@ func (self *SBaseRegionDriver) RequestRemoteUpdateElasticSearch(ctx context.Cont
func (self *SBaseRegionDriver) RequestRemoteUpdateKafka(ctx context.Context, userCred mcclient.TokenCredential, kafka *models.SKafka, replaceTags bool, task taskman.ITask) error {
return errors.Wrapf(cloudprovider.ErrNotImplemented, "RequestRemoteUpdateElasticSearch")
}
func (self *SBaseRegionDriver) ValidateCreateKubeClusterData(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, input *api.KubeClusterCreateInput) (*api.KubeClusterCreateInput, error) {
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "ValidateCreateKubeClusterData")
}
func (self *SBaseRegionDriver) ValidateCreateKubeNodePoolData(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, input *api.KubeNodePoolCreateInput) (*api.KubeNodePoolCreateInput, error) {
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "ValidateCreateKubeNodePoolData")
}
func (self *SBaseRegionDriver) RequestCreateKubeCluster(ctx context.Context, userCred mcclient.TokenCredential, cluster *models.SKubeCluster, task taskman.ITask) error {
return errors.Wrapf(cloudprovider.ErrNotImplemented, "RequestCreateKubeCluster")
}
func (self *SBaseRegionDriver) RequestCreateKubeNodePool(ctx context.Context, userCred mcclient.TokenCredential, pool *models.SKubeNodePool, task taskman.ITask) error {
return errors.Wrapf(cloudprovider.ErrNotImplemented, "RequestCreateKubeNodePool")
}
+108
View File
@@ -3139,3 +3139,111 @@ func (self *SManagedVirtualizationRegionDriver) RequestRemoteUpdateKafka(ctx con
})
return nil
}
func (self *SManagedVirtualizationRegionDriver) ValidateCreateKubeClusterData(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, input *api.KubeClusterCreateInput) (*api.KubeClusterCreateInput, error) {
return input, nil
}
func (self *SManagedVirtualizationRegionDriver) RequestCreateKubeCluster(ctx context.Context, userCred mcclient.TokenCredential, cluster *models.SKubeCluster, task taskman.ITask) error {
taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) {
opts := &cloudprovider.KubeClusterCreateOptions{
NAME: cluster.Name,
Desc: cluster.Description,
Version: cluster.Version,
NetworkIds: []string{},
}
vpc, err := cluster.GetVpc()
if err != nil {
return nil, errors.Wrapf(err, "cluster.GetVpc")
}
opts.VpcId = vpc.ExternalId
networks, err := cluster.GetNetworks()
if err != nil {
return nil, errors.Wrapf(err, "GetNetworks")
}
for _, net := range networks {
opts.NetworkIds = append(opts.NetworkIds, net.ExternalId)
}
opts.Tags, _ = cluster.GetAllUserMetadata()
params := task.GetParams()
opts.ServiceCIDR, _ = params.GetString("service_cidr")
opts.RoleName, _ = params.GetString("role_name")
opts.PrivateAccess, _ = params.Bool("private_access")
opts.PublicAccess, _ = params.Bool("public_access")
iregion, err := cluster.GetIRegion(ctx)
if err != nil {
return nil, errors.Wrapf(err, "GetIRegion")
}
icluster, err := iregion.CreateIKubeCluster(opts)
if err != nil {
return nil, errors.Wrapf(err, "CreateIKubeCluster")
}
err = db.SetExternalId(cluster, userCred, icluster.GetGlobalId())
if err != nil {
return nil, errors.Wrapf(err, "db.SetExternalId")
}
err = cloudprovider.WaitStatusWithSync(icluster, api.KUBE_CLUSTER_STATUS_RUNNING, func(status string) {
cluster.SetStatus(userCred, status, "")
}, time.Second*30, time.Hour*1)
if err != nil {
return nil, errors.Wrapf(err, "wait cluster status timeout, current status: %s", icluster.GetStatus())
}
return nil, nil
})
return nil
}
func (self *SManagedVirtualizationRegionDriver) ValidateCreateKubeNodePoolData(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, input *api.KubeNodePoolCreateInput) (*api.KubeNodePoolCreateInput, error) {
return input, nil
}
func (self *SManagedVirtualizationRegionDriver) RequestCreateKubeNodePool(ctx context.Context, userCred mcclient.TokenCredential, pool *models.SKubeNodePool, task taskman.ITask) error {
taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) {
opts := &cloudprovider.KubeNodePoolCreateOptions{
NAME: pool.Name,
Desc: pool.Description,
MinInstanceCount: pool.MinInstanceCount,
MaxInstanceCount: pool.MaxInstanceCount,
DesiredInstanceCount: pool.DesiredInstanceCount,
RootDiskSizeGb: pool.RootDiskSizeGb,
NetworkIds: []string{},
InstanceTypes: []string{},
}
if pool.InstanceTypes != nil {
for _, instanceType := range *pool.InstanceTypes {
opts.InstanceTypes = append(opts.InstanceTypes, instanceType)
}
}
networks, err := pool.GetNetworks()
if err != nil {
return nil, errors.Wrapf(err, "GetNetworks")
}
for _, net := range networks {
opts.NetworkIds = append(opts.NetworkIds, net.ExternalId)
}
opts.Tags, _ = pool.GetAllUserMetadata()
icluster, err := pool.GetIKubeCluster(ctx)
if err != nil {
return nil, errors.Wrapf(err, "GetIKubeCluster")
}
ipool, err := icluster.CreateIKubeNodePool(opts)
if err != nil {
return nil, errors.Wrapf(err, "CreateIKubeNodePool")
}
err = db.SetExternalId(pool, userCred, ipool.GetGlobalId())
if err != nil {
return nil, errors.Wrapf(err, "db.SetExternalId")
}
err = cloudprovider.WaitStatus(ipool, api.KUBE_CLUSTER_STATUS_RUNNING, time.Second*30, time.Hour*1)
if err != nil {
return nil, errors.Wrapf(err, "wait cluster status timeout, current status: %s", icluster.GetStatus())
}
return nil, pool.SetStatus(userCred, api.KUBE_CLUSTER_STATUS_RUNNING, "")
})
return nil
}
@@ -0,0 +1,68 @@
// Copyright 2019 Yunion
//
// 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 tasks
import (
"context"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
"yunion.io/x/onecloud/pkg/apis"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/cloudcommon/db/taskman"
"yunion.io/x/onecloud/pkg/compute/models"
"yunion.io/x/onecloud/pkg/util/logclient"
)
type KubeClusterCreateTask struct {
taskman.STask
}
func init() {
taskman.RegisterTask(KubeClusterCreateTask{})
}
func (self *KubeClusterCreateTask) taskFail(ctx context.Context, cluster *models.SKubeCluster, err error) {
cluster.SetStatus(self.UserCred, apis.STATUS_CREATE_FAILED, err.Error())
db.OpsLog.LogEvent(cluster, db.ACT_ALLOCATE, err, self.GetUserCred())
logclient.AddActionLogWithStartable(self, cluster, logclient.ACT_ALLOCATE, err, self.UserCred, false)
self.SetStageFailed(ctx, jsonutils.NewString(err.Error()))
}
func (self *KubeClusterCreateTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
cluster := obj.(*models.SKubeCluster)
region, err := cluster.GetRegion()
if err != nil {
self.taskFail(ctx, cluster, errors.Wrapf(err, "GetRegion"))
return
}
self.SetStage("OnKubeClusterCreateComplate", nil)
err = region.GetDriver().RequestCreateKubeCluster(ctx, self.UserCred, cluster, self)
if err != nil {
self.taskFail(ctx, cluster, errors.Wrapf(err, "RequestCreateKubeCluster"))
return
}
}
func (self *KubeClusterCreateTask) OnKubeClusterCreateComplate(ctx context.Context, cluster *models.SKubeCluster, data jsonutils.JSONObject) {
self.SetStageComplete(ctx, nil)
}
func (self *KubeClusterCreateTask) OnKubeClusterCreateComplateFailed(ctx context.Context, cluster *models.SKubeCluster, reason jsonutils.JSONObject) {
self.taskFail(ctx, cluster, errors.Errorf(reason.String()))
}
@@ -0,0 +1,68 @@
// Copyright 2019 Yunion
//
// 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 tasks
import (
"context"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
"yunion.io/x/onecloud/pkg/apis"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/cloudcommon/db/taskman"
"yunion.io/x/onecloud/pkg/compute/models"
"yunion.io/x/onecloud/pkg/util/logclient"
)
type KubeNodePoolCreateTask struct {
taskman.STask
}
func init() {
taskman.RegisterTask(KubeNodePoolCreateTask{})
}
func (self *KubeNodePoolCreateTask) taskFail(ctx context.Context, pool *models.SKubeNodePool, err error) {
pool.SetStatus(self.UserCred, apis.STATUS_CREATE_FAILED, err.Error())
db.OpsLog.LogEvent(pool, db.ACT_ALLOCATE, err, self.GetUserCred())
logclient.AddActionLogWithStartable(self, pool, logclient.ACT_ALLOCATE, err, self.UserCred, false)
self.SetStageFailed(ctx, jsonutils.NewString(err.Error()))
}
func (self *KubeNodePoolCreateTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
pool := obj.(*models.SKubeNodePool)
region, err := pool.GetRegion()
if err != nil {
self.taskFail(ctx, pool, errors.Wrapf(err, "GetRegion"))
return
}
self.SetStage("OnKubeNodePoolCreateComplate", nil)
err = region.GetDriver().RequestCreateKubeNodePool(ctx, self.UserCred, pool, self)
if err != nil {
self.taskFail(ctx, pool, errors.Wrapf(err, "RequestCreateKubeNodePool"))
return
}
}
func (self *KubeNodePoolCreateTask) OnKubeNodePoolCreateComplate(ctx context.Context, pool *models.SKubeNodePool, data jsonutils.JSONObject) {
self.SetStageComplete(ctx, nil)
}
func (self *KubeNodePoolCreateTask) OnKubeNodePoolCreateComplateFailed(ctx context.Context, pool *models.SKubeNodePool, reason jsonutils.JSONObject) {
self.taskFail(ctx, pool, errors.Errorf(reason.String()))
}
@@ -0,0 +1,78 @@
// Copyright 2019 Yunion
//
// 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 tasks
import (
"context"
"time"
"yunion.io/x/cloudmux/pkg/cloudprovider"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
"yunion.io/x/onecloud/pkg/apis"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/cloudcommon/db/taskman"
"yunion.io/x/onecloud/pkg/compute/models"
"yunion.io/x/onecloud/pkg/util/logclient"
)
type KubeNodePoolDeleteTask struct {
taskman.STask
}
func init() {
taskman.RegisterTask(KubeNodePoolDeleteTask{})
}
func (self *KubeNodePoolDeleteTask) taskFailed(ctx context.Context, pool *models.SKubeNodePool, err error) {
pool.SetStatus(self.UserCred, apis.STATUS_DELETE_FAILED, err.Error())
db.OpsLog.LogEvent(pool, db.ACT_DELOCATE_FAIL, err, self.UserCred)
logclient.AddActionLogWithStartable(self, pool, logclient.ACT_DELETE, err, self.UserCred, false)
self.SetStageFailed(ctx, jsonutils.NewString(err.Error()))
}
func (self *KubeNodePoolDeleteTask) OnInit(ctx context.Context, obj db.IStandaloneModel, body jsonutils.JSONObject) {
pool := obj.(*models.SKubeNodePool)
ipool, err := pool.GetIKubeNodePool(ctx)
if err != nil {
if errors.Cause(err) == cloudprovider.ErrNotFound {
self.taskComplete(ctx, pool)
return
}
self.taskFailed(ctx, pool, errors.Wrapf(err, "GetIKubeNodePool"))
return
}
err = ipool.Delete()
if err != nil {
self.taskFailed(ctx, pool, errors.Wrapf(err, "Delete"))
return
}
err = cloudprovider.WaitDeleted(ipool, time.Second*10, time.Minute*10)
if err != nil {
self.taskFailed(ctx, pool, errors.Wrapf(err, "WaitDeleted"))
return
}
self.taskComplete(ctx, pool)
}
func (self *KubeNodePoolDeleteTask) taskComplete(ctx context.Context, pool *models.SKubeNodePool) {
logclient.AddActionLogWithStartable(self, pool, logclient.ACT_DELETE, pool.GetShortDesc(ctx), self.UserCred, true)
pool.RealDelete(ctx, self.GetUserCred())
self.SetStageComplete(ctx, nil)
}
@@ -0,0 +1,60 @@
// Copyright 2019 Yunion
//
// 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 tasks
import (
"context"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
"yunion.io/x/onecloud/pkg/apis"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/cloudcommon/db/taskman"
"yunion.io/x/onecloud/pkg/compute/models"
"yunion.io/x/onecloud/pkg/util/logclient"
)
type KubeNodePoolSyncstatusTask struct {
taskman.STask
}
func init() {
taskman.RegisterTask(KubeNodePoolSyncstatusTask{})
}
func (self *KubeNodePoolSyncstatusTask) taskFailed(ctx context.Context, pool *models.SKubeNodePool, err error) {
pool.SetStatus(self.GetUserCred(), apis.STATUS_UNKNOWN, err.Error())
db.OpsLog.LogEvent(pool, db.ACT_SYNC_STATUS, pool.GetShortDesc(ctx), self.GetUserCred())
logclient.AddActionLogWithContext(ctx, pool, logclient.ACT_SYNC_STATUS, err, self.UserCred, false)
self.SetStageFailed(ctx, jsonutils.NewString(err.Error()))
}
func (self *KubeNodePoolSyncstatusTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
pool := obj.(*models.SKubeNodePool)
iNodePool, err := pool.GetIKubeNodePool(ctx)
if err != nil {
self.taskFailed(ctx, pool, errors.Wrapf(err, "GetIKubeNodePool"))
return
}
err = pool.SetStatus(self.UserCred, iNodePool.GetStatus(), "")
if err != nil {
self.taskFailed(ctx, pool, errors.Wrapf(err, "SetStatus"))
return
}
self.SetStageComplete(ctx, nil)
}
@@ -49,3 +49,15 @@ type KubeClusterConfigOptions struct {
func (opts *KubeClusterConfigOptions) Params() (jsonutils.JSONObject, error) {
return jsonutils.Marshal(opts), nil
}
type KubeClusterCreateOptions struct {
options.BaseCreateOptions
Version string
VpcId string
NetworkIds []string
RoleName string
}
func (opts *KubeClusterCreateOptions) Params() (jsonutils.JSONObject, error) {
return jsonutils.Marshal(opts), nil
}
@@ -39,3 +39,18 @@ func (opts *KubeNodePoolIdOption) GetId() string {
func (opts *KubeNodePoolIdOption) Params() (jsonutils.JSONObject, error) {
return nil, nil
}
type KubeNodePoolCreateOptions struct {
options.BaseCreateOptions
NetworkIds []string `metavar:"NETWORK"`
InstanceTypes []string `metavar:"INSTANCE_TYPE"`
MinInstanceCount int
MaxInstanceCount int
DesiredInstanceCount int
RootDiskSizeGb int
CloudKubeClusterId string `metavar:"CLUSTER"`
}
func (opts *KubeNodePoolCreateOptions) Params() (jsonutils.JSONObject, error) {
return jsonutils.Marshal(opts), nil
}
+2 -2
View File
@@ -1457,7 +1457,7 @@ sigs.k8s.io/structured-merge-diff/v4/value
# sigs.k8s.io/yaml v1.2.0
## explicit; go 1.12
sigs.k8s.io/yaml
# yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230515072146-3a7171baa793
# yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230519063348-6f0a3a041cdb
## explicit; go 1.18
yunion.io/x/cloudmux/pkg/apis
yunion.io/x/cloudmux/pkg/apis/billing
@@ -1603,7 +1603,7 @@ yunion.io/x/pkg/utils
# yunion.io/x/s3cli v0.0.0-20190917004522-13ac36d8687e
## explicit; go 1.12
yunion.io/x/s3cli
# yunion.io/x/sqlchemy v1.1.2-0.20230512065832-5323af46107e
# yunion.io/x/sqlchemy v1.1.2-0.20230518004216-245cb3766e38
## explicit; go 1.17
yunion.io/x/sqlchemy
yunion.io/x/sqlchemy/backends
+15
View File
@@ -13,3 +13,18 @@
// limitations under the License.
package compute
const (
KUBE_CLUSTER_STATUS_RUNNING = "running"
KUBE_CLUSTER_STATUS_CREATING = "creating"
KUBE_CLUSTER_STATUS_DELETING = "deleting"
KUBE_CLUSTER_STATUS_ABNORMAL = "abnormal"
// 升级中
KUBE_CLUSTER_STATUS_UPDATING = "updating"
// 升级失败
KUBE_CLUSTER_STATUS_UPDATING_FAILED = "updating_failed"
// 伸缩中
KUBE_CLUSTER_STATUS_SCALING = "scaling"
// 停止
KUBE_CLUSTER_STATUS_STOPPED = "stopped"
)
+29
View File
@@ -22,3 +22,32 @@ type SKubeconfig struct {
Config string `json:"config"`
Expiration time.Time `json:"expiration"`
}
type KubeClusterCreateOptions struct {
NAME string
Desc string
VpcId string
Version string
NetworkIds []string
Tags map[string]string
ServiceCIDR string
PrivateAccess bool
PublicAccess bool
RoleName string
}
type KubeNodePoolCreateOptions struct {
NAME string
Desc string
MinInstanceCount int
MaxInstanceCount int
DesiredInstanceCount int
RootDiskSizeGb int
InstanceTypes []string
NetworkIds []string
Tags map[string]string
}
+16
View File
@@ -189,6 +189,7 @@ type ICloudRegion interface {
GetICloudKubeClusters() ([]ICloudKubeCluster, error)
GetICloudKubeClusterById(id string) (ICloudKubeCluster, error)
CreateIKubeCluster(opts *KubeClusterCreateOptions) (ICloudKubeCluster, error)
GetICloudTablestores() ([]ICloudTablestore, error)
@@ -1610,7 +1611,12 @@ type ICloudKubeCluster interface {
GetKubeConfig(private bool, expireMinutes int) (*SKubeconfig, error)
GetVersion() string
GetVpcId() string
GetNetworkIds() []string
GetIKubeNodePools() ([]ICloudKubeNodePool, error)
CreateIKubeNodePool(opts *KubeNodePoolCreateOptions) (ICloudKubeNodePool, error)
GetIKubeNodes() ([]ICloudKubeNode, error)
Delete(isRetain bool) error
@@ -1624,6 +1630,16 @@ type ICloudKubeNode interface {
type ICloudKubeNodePool interface {
ICloudResource
GetMinInstanceCount() int
GetMaxInstanceCount() int
GetDesiredInstanceCount() int
GetRootDiskSizeGb() int
GetInstanceTypes() []string
GetNetworkIds() []string
Delete() error
}
type ICloudTablestore interface {
+19 -4
View File
@@ -21,7 +21,7 @@ import (
"yunion.io/x/pkg/errors"
)
func WaitStatus(res ICloudResource, expect string, interval time.Duration, timeout time.Duration) error {
func WaitStatusWithSync(res ICloudResource, expect string, sync func(status string), interval time.Duration, timeout time.Duration) error {
startTime := time.Now()
for time.Now().Sub(startTime) < timeout {
err := res.Refresh()
@@ -29,6 +29,9 @@ func WaitStatus(res ICloudResource, expect string, interval time.Duration, timeo
return err
}
log.Infof("%s status %s expect %s", res.GetName(), res.GetStatus(), expect)
if sync != nil {
sync(res.GetStatus())
}
if res.GetStatus() == expect {
return nil
}
@@ -37,16 +40,24 @@ func WaitStatus(res ICloudResource, expect string, interval time.Duration, timeo
return ErrTimeout
}
func WaitMultiStatus(res ICloudResource, expects []string, interval time.Duration, timeout time.Duration) error {
func WaitStatus(res ICloudResource, expect string, interval time.Duration, timeout time.Duration) error {
return WaitStatusWithSync(res, expect, nil, interval, timeout)
}
func WaitMultiStatusWithSync(res ICloudResource, expects []string, sync func(string), interval time.Duration, timeout time.Duration) error {
startTime := time.Now()
for time.Now().Sub(startTime) < timeout {
err := res.Refresh()
if err != nil {
return errors.Wrap(err, "resource.Refresh()")
}
log.Infof("%s status %s expect %s", res.GetName(), res.GetStatus(), expects)
status := res.GetStatus()
log.Infof("%s status %s expect %s", res.GetName(), status, expects)
if sync != nil {
sync(status)
}
for _, expect := range expects {
if res.GetStatus() == expect {
if status == expect {
return nil
}
}
@@ -55,6 +66,10 @@ func WaitMultiStatus(res ICloudResource, expects []string, interval time.Duratio
return errors.Wrap(errors.ErrTimeout, "WaitMultistatus")
}
func WaitMultiStatus(res ICloudResource, expects []string, interval time.Duration, timeout time.Duration) error {
return WaitMultiStatusWithSync(res, expects, nil, interval, timeout)
}
func WaitStatusWithDelay(res ICloudResource, expect string, delay time.Duration, interval time.Duration, timeout time.Duration) error {
time.Sleep(delay)
return WaitStatus(res, expect, interval, timeout)
+19
View File
@@ -131,6 +131,21 @@ func (self *SKubeCluster) GetKubeConfig(private bool, expireMinutes int) (*cloud
return self.region.GetKubeConfig(self.ClusterId, private, expireMinutes)
}
func (self *SKubeCluster) GetVersion() string {
return self.CurrentVersion
}
func (self *SKubeCluster) GetVpcId() string {
return self.VpcId
}
func (self *SKubeCluster) GetNetworkIds() []string {
if len(self.VswitchId) > 0 {
return []string{self.VswitchId}
}
return []string{}
}
func (self *SKubeCluster) Delete(isRetain bool) error {
return self.region.DeleteKubeCluster(self.ClusterId, isRetain)
}
@@ -256,3 +271,7 @@ func (self *SRegion) DeleteKubeCluster(id string, isRetain bool) error {
_, err := self.k8sRequest("DeleteCluster", params)
return errors.Wrapf(err, "DeleteCluster")
}
func (self *SKubeCluster) CreateIKubeNodePool(opts *cloudprovider.KubeNodePoolCreateOptions) (cloudprovider.ICloudKubeNodePool, error) {
return nil, cloudprovider.ErrNotImplemented
}
+29 -1
View File
@@ -59,7 +59,7 @@ type SKubeNodePool struct {
AutoRenewPeriod int `json:"auto_renew_period"`
ImageId string `json:"image_id"`
InstanceChargeType string `json:"instance_charge_type"`
Instance_types []string `json:"instance_types"`
InstanceTypes []string `json:"instance_types"`
MultiAzPolicy string `json:"multi_az_policy"`
OnDemandBaseCapacity int `json:"on_demand_base_capacity"`
OnDemandPercentageAboveBaseCapacity int `json:"on_demand_percentage_above_base_capacity"`
@@ -126,6 +126,34 @@ func (self *SKubeNodePool) GetStatus() string {
return self.Status.State
}
func (self *SKubeNodePool) GetMinInstanceCount() int {
return self.AutoScaling.MinInstances
}
func (self *SKubeNodePool) GetMaxInstanceCount() int {
return self.AutoScaling.MaxInstances
}
func (self *SKubeNodePool) GetDesiredInstanceCount() int {
return 0
}
func (self *SKubeNodePool) GetRootDiskSizeGb() int {
return self.ScalingGroup.SystemDiskSize
}
func (self *SKubeNodePool) Delete() error {
return cloudprovider.ErrNotImplemented
}
func (self *SKubeNodePool) GetInstanceTypes() []string {
return self.ScalingGroup.InstanceTypes
}
func (self *SKubeNodePool) GetNetworkIds() []string {
return self.ScalingGroup.VswitchIds
}
func (self *SKubeCluster) GetIKubeNodePools() ([]cloudprovider.ICloudKubeNodePool, error) {
pools, err := self.region.GetKubeNodePools(self.ClusterId)
if err != nil {
-5
View File
@@ -61,10 +61,6 @@ const (
DefaultAssumeRoleName = "OrganizationAccountAccessRole"
)
var (
DEBUG = false
)
type AwsClientConfig struct {
cpcfg cloudprovider.ProviderConfig
@@ -95,7 +91,6 @@ func (cfg *AwsClientConfig) CloudproviderConfig(cpcfg cloudprovider.ProviderConf
func (cfg *AwsClientConfig) Debug(debug bool) *AwsClientConfig {
cfg.debug = debug
DEBUG = debug
return cfg
}
+227
View File
@@ -0,0 +1,227 @@
// Copyright 2019 Yunion
//
// 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 aws
import (
"io/ioutil"
"strings"
"github.com/aws/aws-sdk-go/aws"
"github.com/aws/aws-sdk-go/aws/awserr"
"github.com/aws/aws-sdk-go/aws/client"
"github.com/aws/aws-sdk-go/aws/client/metadata"
"github.com/aws/aws-sdk-go/aws/corehandlers"
"github.com/aws/aws-sdk-go/aws/request"
v4 "github.com/aws/aws-sdk-go/aws/signer/v4"
"github.com/aws/aws-sdk-go/private/protocol/restjson"
"yunion.io/x/cloudmux/pkg/cloudprovider"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
)
// eks only
func (self *SAwsClient) invoke(regionId, serviceName, serviceId, apiVersion string, apiName, path string, params map[string]interface{}, retval interface{}, assumeRole bool) error {
if len(regionId) == 0 {
regionId = self.getDefaultRegionId()
}
session, err := self.getAwsSession(regionId, assumeRole)
if err != nil {
return err
}
c := session.ClientConfig(serviceName)
metadata := metadata.ClientInfo{
ServiceName: serviceName,
ServiceID: serviceId,
SigningName: c.SigningName,
SigningRegion: c.SigningRegion,
Endpoint: c.Endpoint,
APIVersion: apiVersion,
}
if self.debug {
logLevel := aws.LogLevelType(uint(aws.LogDebugWithRequestErrors) + uint(aws.LogDebugWithHTTPBody))
c.Config.LogLevel = &logLevel
}
client := client.New(*c.Config, metadata, c.Handlers)
client.Handlers.Sign.PushBackNamed(v4.SignRequestHandler)
client.Handlers.Build.PushBackNamed(JsonBuildHandler)
client.Handlers.Unmarshal.PushBackNamed(UnmarshalJsonHandler)
client.Handlers.UnmarshalMeta.PushBackNamed(restjson.UnmarshalMetaHandler)
client.Handlers.UnmarshalError.PushBackNamed(UnmarshalJsonErrorHandler)
client.Handlers.Validate.Remove(corehandlers.ValidateEndpointHandler)
return jsonInvoke(client, apiName, path, params, retval, true)
}
func jsonInvoke(cli *client.Client, apiName, path string, params map[string]interface{}, retval interface{}, debug bool) error {
method := "POST"
for _, key := range []string{"List", "Describe"} {
if strings.HasPrefix(apiName, key) {
method = "GET"
break
}
}
if strings.HasPrefix(apiName, "Delete") {
method = "DELETE"
}
if method == "GET" || method == "DELETE" {
for k, v := range params {
if strings.Contains(path, "{"+k+"}") {
path = strings.Replace(path, "{"+k+"}", v.(string), 1)
delete(params, k)
}
}
}
op := &request.Operation{
Name: apiName,
HTTPMethod: method,
HTTPPath: path,
Paginator: &request.Paginator{
InputTokens: []string{"nextToken"},
OutputTokens: []string{"nextToken"},
LimitToken: "maxResults",
TruncationToken: "",
},
}
req := cli.NewRequest(op, params, retval)
err := req.Send()
if err != nil {
if e, ok := err.(awserr.RequestFailure); ok && e.StatusCode() == 404 {
return cloudprovider.ErrNotFound
}
return err
}
return nil
}
var JsonBuildHandler = request.NamedHandler{
Name: "yunion.json.Build",
Fn: JsonBuild,
}
func JsonBuild(r *request.Request) {
if r.Config.LogLevel != nil && r.Config.LogLevel.AtLeast(aws.LogDebugWithHTTPBody) {
log.Debugf("body: %s", jsonutils.Marshal(r.Params).PrettyString())
}
if r.Params != nil && r.HTTPRequest.Method == "POST" {
r.SetBufferBody([]byte(jsonutils.Marshal(r.Params).String()))
}
if v := r.HTTPRequest.Header.Get("Content-Type"); len(v) == 0 {
r.HTTPRequest.Header.Set("Content-Type", "application/json")
}
if r.ClientInfo.TargetPrefix != "" {
target := r.ClientInfo.TargetPrefix + "." + r.Operation.Name
r.HTTPRequest.Header.Add("X-Amz-Target", target)
}
if ct, v := r.HTTPRequest.Header.Get("Content-Type"), r.ClientInfo.JSONVersion; len(ct) == 0 && len(v) != 0 {
jsonVersion := r.ClientInfo.JSONVersion
r.HTTPRequest.Header.Set("Content-Type", "application/x-amz-json-"+jsonVersion)
}
}
var UnmarshalJsonHandler = request.NamedHandler{Name: "yunion.query.UnmarshalJson", Fn: UnmarshalJson}
func UnmarshalJson(r *request.Request) {
defer r.HTTPResponse.Body.Close()
if r.DataFilled() {
body, err := ioutil.ReadAll(r.HTTPResponse.Body)
if err != nil {
r.Error = awserr.NewRequestFailure(
awserr.New("ioutil.ReadAll", "read response body", err),
r.HTTPResponse.StatusCode,
r.RequestID,
)
return
}
if r.Config.LogLevel != nil && r.Config.LogLevel.AtLeast(aws.LogDebugWithHTTPBody) {
log.Debugf("response: \n%s", string(body))
}
obj, err := jsonutils.Parse(body)
if err != nil {
r.Error = awserr.NewRequestFailure(
awserr.New("DecodeElement", "failed decoding Query response", err),
r.HTTPResponse.StatusCode,
r.RequestID,
)
return
}
err = obj.Unmarshal(r.Data)
if err != nil {
r.Error = awserr.NewRequestFailure(
awserr.New("DecodeElement", "failed decoding Query response", err),
r.HTTPResponse.StatusCode,
r.RequestID,
)
return
}
}
}
var UnmarshalJsonErrorHandler = request.NamedHandler{
Name: "awssdk.ec2query.UnmarshalJsonError",
Fn: UnmarshalJsonError,
}
type sAwsInvokeError struct {
Message string
}
func (self sAwsInvokeError) Error() string {
return jsonutils.Marshal(self).String()
}
func UnmarshalJsonError(r *request.Request) {
defer r.HTTPResponse.Body.Close()
result, err := ioutil.ReadAll(r.HTTPResponse.Body)
if err != nil {
r.Error = errors.Wrapf(err, "ioutil.ReadAll")
return
}
if r.HTTPResponse.StatusCode == 404 {
r.Error = errors.Wrapf(cloudprovider.ErrNotFound, string(result))
return
}
obj, err := jsonutils.Parse(result)
if err != nil {
r.Error = errors.Wrapf(err, "jsonutils.Parse")
return
}
if r.Config.LogLevel != nil && r.Config.LogLevel.AtLeast(aws.LogDebugWithHTTPBody) {
log.Debugf("response error: %s", obj.PrettyString())
}
respErr := &sAwsInvokeError{}
err = obj.Unmarshal(respErr)
if err != nil {
r.Error = errors.Wrapf(err, obj.String())
return
}
r.Error = respErr
return
}
+2 -2
View File
@@ -44,7 +44,7 @@ func Unmarshal(r *request.Request) {
defer r.HTTPResponse.Body.Close()
if r.DataFilled() {
var decoder *xml.Decoder
if DEBUG {
if r.Config.LogLevel != nil && r.Config.LogLevel.AtLeast(aws.LogDebugWithHTTPBody) {
body, err := ioutil.ReadAll(r.HTTPResponse.Body)
if err != nil {
r.Error = awserr.NewRequestFailure(
@@ -126,7 +126,7 @@ func Build(r *request.Request) {
}
}
if DEBUG {
if r.Config.LogLevel != nil && r.Config.LogLevel.AtLeast(aws.LogDebugWithHTTPBody) {
log.Debugf("params: %s", body.Encode())
}
+422
View File
@@ -0,0 +1,422 @@
// Copyright 2019 Yunion
//
// 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 aws
import (
"fmt"
"strings"
"time"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/utils"
api "yunion.io/x/cloudmux/pkg/apis/compute"
"yunion.io/x/cloudmux/pkg/cloudprovider"
"yunion.io/x/cloudmux/pkg/multicloud"
)
type SKubeCluster struct {
multicloud.SResourceBase
AwsTags
region *SRegion
Name string
Arn string `json:"arn"`
CreatedAt float64 `json:"createdAt"`
Version string `json:"version"`
Endpoint string `json:"endpoint"`
RoleArn string `json:"roleArn"`
ResourcesVpcConfig struct {
SubnetIds []string `json:"subnetIds"`
SecurityGroupIds []string `json:"securityGroupIds"`
ClusterSecurityGroupId string `json:"clusterSecurityGroupId"`
VpcId string `json:"vpcId"`
EndpointPublicAccess bool `json:"endpointPublicAccess"`
EndpointPrivateAccess bool `json:"endpointPrivateAccess"`
PublicAccessCidrs []string `json:"publicAccessCidrs"`
} `json:"resourcesVpcConfig"`
KubernetesNetworkConfig struct {
ServiceIpv4CIDR string `json:"serviceIpv4Cidr"`
ServiceIpv6CIDR string `json:"serviceIpv6Cidr"`
IPFamily string `json:"ipFamily"`
} `json:"kubernetesNetworkConfig"`
Logging struct {
ClusterLogging []ClusterLogging `json:"clusterLogging"`
} `json:"logging"`
Identity string `json:"identity"`
Status string `json:"status"`
CertificateAuthority struct {
Data string
} `json:"certificateAuthority"`
ClientRequestToken string `json:"clientRequestToken"`
PlatformVersion string `json:"platformVersion"`
EncryptionConfig string `json:"encryptionConfig"`
ConnectorConfig string `json:"connectorConfig"`
Id string `json:"id"`
Health string `json:"health"`
}
type ClusterLogging struct {
Types []string `json:"types"`
Enabled bool `json:"enabled"`
}
func (self *SKubeCluster) GetName() string {
return self.Name
}
func (self *SKubeCluster) GetId() string {
return self.Name
}
func (self *SKubeCluster) GetGlobalId() string {
return self.GetId()
}
func (self *SKubeCluster) GetEnabled() bool {
return true
}
func (self *SKubeCluster) GetStatus() string {
if len(self.Status) == 0 {
self.Refresh()
}
switch self.Status {
case "ACTIVE":
return api.KUBE_CLUSTER_STATUS_RUNNING
case "DELETING":
return api.KUBE_CLUSTER_STATUS_DELETING
default:
return strings.ToLower(self.Status)
}
}
func (self *SKubeCluster) GetKubeConfig(private bool, expireMinutes int) (*cloudprovider.SKubeconfig, error) {
if len(self.CertificateAuthority.Data) == 0 {
self.Refresh()
}
eksId := fmt.Sprintf("%s:%s:cluster/%s", self.region.RegionId, self.region.client.ownerId, self.Name)
config := fmt.Sprintf(`apiVersion: v1
clusters:
- cluster:
server: %s
certificate-authority-data: %s
name: arn:aws:eks:%s
contexts:
- context:
cluster: arn:aws:eks:%s
user: arn:aws:eks:%s
name: arn:aws:eks:%s
current-context: arn:aws:eks:%s
kind: Config
preferences: {}
users:
- name: arn:aws:eks:%s
user:
exec:
apiVersion: client.authentication.k8s.io/v1beta1
command: aws-iam-authenticator
args:
- "token"
- "-i"
- "%s"`, self.Endpoint, self.CertificateAuthority.Data, eksId, eksId, eksId, eksId, eksId, eksId, self.Name)
return &cloudprovider.SKubeconfig{
Config: config,
}, nil
}
func (self *SKubeCluster) GetIKubeNodePools() ([]cloudprovider.ICloudKubeNodePool, error) {
ret := []cloudprovider.ICloudKubeNodePool{}
nextToken := ""
for {
part, nextToken, err := self.region.GetNodegroups(self.Name, nextToken)
if err != nil {
return nil, errors.Wrapf(err, "GetNodegroups")
}
for i := range part {
ret = append(ret, &part[i])
}
if len(nextToken) == 0 {
break
}
}
return ret, nil
}
func (self *SKubeCluster) GetIKubeNodes() ([]cloudprovider.ICloudKubeNode, error) {
return nil, cloudprovider.ErrNotImplemented
}
func (self *SKubeCluster) Delete(isRetain bool) error {
return self.region.DeleteKubeCluster(self.Name)
}
func (self *SKubeCluster) GetCreatedAt() time.Time {
return time.Unix(int64(self.CreatedAt), 0)
}
func (self *SKubeCluster) GetVpcId() string {
if len(self.ResourcesVpcConfig.VpcId) == 0 {
self.Refresh()
}
return self.ResourcesVpcConfig.VpcId
}
func (self *SKubeCluster) GetVersion() string {
if len(self.Version) == 0 {
self.Refresh()
}
return self.Version
}
func (self *SKubeCluster) GetNetworkIds() []string {
if len(self.ResourcesVpcConfig.SubnetIds) == 0 {
self.Refresh()
}
return self.ResourcesVpcConfig.SubnetIds
}
func (self *SKubeCluster) Refresh() error {
cluster, err := self.region.GetKubeCluster(self.Name)
if err != nil {
return err
}
return jsonutils.Update(self, cluster)
}
func (self *SRegion) GetKubeClusters(nextToken string) ([]SKubeCluster, string, error) {
ret := struct {
Clusters []string
NextToken string
}{}
params := map[string]interface{}{
"include": "all",
}
if len(nextToken) > 0 {
params["nextToken"] = nextToken
}
result := []SKubeCluster{}
err := self.eksRequest("ListClusters", "/clusters", params, &ret)
if err != nil {
return nil, "", errors.Wrapf(err, "ListClusters")
}
for i := range ret.Clusters {
result = append(result, SKubeCluster{
region: self,
Name: ret.Clusters[i],
})
}
return result, ret.NextToken, nil
}
func (self *SRegion) GetKubeCluster(name string) (*SKubeCluster, error) {
params := map[string]interface{}{
"name": name,
}
ret := struct {
Cluster SKubeCluster
}{}
err := self.eksRequest("DescribeCluster", "/clusters/{name}", params, &ret)
if err != nil {
return nil, errors.Wrapf(err, "DescribeCluster")
}
ret.Cluster.region = self
return &ret.Cluster, nil
}
func (self *SRegion) DeleteKubeCluster(name string) error {
params := map[string]interface{}{
"name": name,
}
ret := struct {
}{}
return self.eksRequest("DeleteCluster", "/clusters/{name}", params, &ret)
}
func (self *SRegion) GetICloudKubeClusters() ([]cloudprovider.ICloudKubeCluster, error) {
ret := []cloudprovider.ICloudKubeCluster{}
nextToken := ""
for {
part, nextToken, err := self.GetKubeClusters(nextToken)
if err != nil {
return nil, errors.Wrapf(err, "GetKubeClusters")
}
for i := range part {
part[i].region = self
ret = append(ret, &part[i])
}
if len(nextToken) == 0 {
break
}
}
return ret, nil
}
func (self *SRegion) GetICloudKubeClusterById(id string) (cloudprovider.ICloudKubeCluster, error) {
cluster, err := self.GetKubeCluster(id)
if err != nil {
return nil, err
}
return cluster, nil
}
func (self *SRegion) CreateIKubeCluster(opts *cloudprovider.KubeClusterCreateOptions) (cloudprovider.ICloudKubeCluster, error) {
cluster, err := self.CreateKubeCluster(opts)
if err != nil {
return nil, err
}
return cluster, nil
}
func (self *SRegion) CreateKubeCluster(opts *cloudprovider.KubeClusterCreateOptions) (*SKubeCluster, error) {
if !opts.PrivateAccess && !opts.PublicAccess { // avoid occur 'Private and public endpoint access cannot be false' error
opts.PrivateAccess = true
}
params := map[string]interface{}{
"name": opts.NAME,
"clientRequestToken": utils.GenRequestId(20),
"resourcesVpcConfig": map[string]interface{}{
"endpointPrivateAccess": opts.PrivateAccess,
"endpointPublicAccess": opts.PublicAccess,
"subnetIds": opts.NetworkIds,
},
"tags": opts.Tags,
}
if len(opts.RoleName) == 0 {
opts.RoleName = "eksClusterRole"
}
role, err := func() (*SRole, error) {
role, err := self.client.GetRole(opts.RoleName)
if err != nil {
if errors.Cause(err) == cloudprovider.ErrNotFound {
params := map[string]string{
"RoleName": opts.RoleName,
"Description": opts.Desc,
"AssumeRolePolicyDocument": k8sRole,
}
role := struct {
Role SRole
}{}
err := self.client.iamRequest("CreateRole", params, &role)
if err != nil {
return nil, errors.Wrapf(err, "CreateRole")
}
role.Role.client = self.client
return &role.Role, nil
}
return nil, errors.Wrapf(err, "GetRole(%s)", opts.RoleName)
}
return role, nil
}()
if err != nil {
return nil, err
}
err = self.client.AttachRolePolicy(opts.RoleName, self.client.getIamArn("AmazonEKSClusterPolicy"))
if err != nil {
return nil, errors.Wrapf(err, "AttachRolePolicy")
}
params["roleArn"] = role.Arn
if len(opts.ServiceCIDR) > 0 {
params["kubernetesNetworkConfig"] = map[string]interface{}{
"ipFamily": "ipv4",
"serviceIpv4Cidr": opts.ServiceCIDR,
}
}
if len(opts.Version) > 0 {
params["version"] = opts.Version
}
ret := struct {
Cluster SKubeCluster
}{}
err = self.eksRequest("CreateCluster", "/clusters", params, &ret)
if err != nil {
return nil, err
}
ret.Cluster.region = self
return &ret.Cluster, nil
}
func (self *SRegion) CreateNodegroup(cluster string, opts *cloudprovider.KubeNodePoolCreateOptions) (*SNodeGroup, error) {
params := map[string]interface{}{
"nodegroupName": opts.NAME,
"clientRequestToken": utils.GenRequestId(20),
"diskSize": opts.RootDiskSizeGb,
"instanceTypes": opts.InstanceTypes,
"tags": opts.Tags,
"scalingConfig": map[string]interface{}{
"desiredSize": opts.DesiredInstanceCount,
"maxSize": opts.MaxInstanceCount,
"minSize": opts.MinInstanceCount,
},
"subnets": opts.NetworkIds,
}
roleName := "AmazonEKSNodeRole"
role, err := func() (*SRole, error) {
role, err := self.client.GetRole(roleName)
if err != nil {
if errors.Cause(err) == cloudprovider.ErrNotFound {
params := map[string]string{
"RoleName": roleName,
"Description": opts.Desc,
"AssumeRolePolicyDocument": nodeRole,
}
role := struct {
Role SRole
}{}
err := self.client.iamRequest("CreateRole", params, &role)
if err != nil {
return nil, errors.Wrapf(err, "CreateRole")
}
role.Role.client = self.client
return &role.Role, nil
}
return nil, errors.Wrapf(err, "GetRole(%s)", roleName)
}
return role, nil
}()
if err != nil {
return nil, errors.Wrapf(err, "Create role")
}
for _, policy := range []string{"AmazonEKSWorkerNodePolicy", "AmazonEC2ContainerRegistryReadOnly"} {
err = self.client.AttachRolePolicy(roleName, self.client.getIamArn(policy))
if err != nil {
return nil, errors.Wrapf(err, "AttachRolePolicy %s", policy)
}
}
params["nodeRole"] = role.Arn
ret := struct {
Nodegroup SNodeGroup
}{}
err = self.eksRequest("CreateNodegroup", fmt.Sprintf("/clusters/%s/node-groups", cluster), params, &ret)
if err != nil {
return nil, err
}
ret.Nodegroup.region = self
return &ret.Nodegroup, nil
}
func (self *SKubeCluster) CreateIKubeNodePool(opts *cloudprovider.KubeNodePoolCreateOptions) (cloudprovider.ICloudKubeNodePool, error) {
nodegroup, err := self.region.CreateNodegroup(self.Name, opts)
if err != nil {
return nil, errors.Wrapf(err, "CreateNodegroup")
}
return nodegroup, nil
}
+176
View File
@@ -0,0 +1,176 @@
// Copyright 2019 Yunion
//
// 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 aws
import (
"fmt"
"strings"
api "yunion.io/x/cloudmux/pkg/apis/compute"
"yunion.io/x/cloudmux/pkg/multicloud"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
)
type SNodeGroup struct {
multicloud.SResourceBase
region *SRegion
AwsTags
ClusterName string
NodegroupName string
DiskSize int
Status string
InstanceTypes []string
Subnets []string
ScalingConfig struct {
DesiredSize int
MaxSize int
MinSize int
}
}
func (self *SNodeGroup) GetId() string {
return self.NodegroupName
}
func (self *SNodeGroup) GetName() string {
return self.NodegroupName
}
func (self *SNodeGroup) GetGlobalId() string {
return self.NodegroupName
}
func (self *SNodeGroup) Refresh() error {
ng, err := self.region.GetNodegroup(self.ClusterName, self.NodegroupName)
if err != nil {
return err
}
self.InstanceTypes = nil
self.Subnets = nil
return jsonutils.Update(self, ng)
}
func (self *SNodeGroup) GetStatus() string {
if len(self.Status) == 0 {
self.Refresh()
}
switch strings.ToLower(self.Status) {
case "active":
return api.KUBE_CLUSTER_STATUS_RUNNING
case "creating":
return api.KUBE_CLUSTER_STATUS_CREATING
}
return strings.ToLower(self.Status)
}
func (self *SNodeGroup) GetMinInstanceCount() int {
if len(self.Subnets) == 0 {
self.Refresh()
}
return self.ScalingConfig.MinSize
}
func (self *SNodeGroup) GetMaxInstanceCount() int {
if len(self.Subnets) == 0 {
self.Refresh()
}
return self.ScalingConfig.MaxSize
}
func (self *SNodeGroup) GetDesiredInstanceCount() int {
if len(self.Subnets) == 0 {
self.Refresh()
}
return self.ScalingConfig.DesiredSize
}
func (self *SNodeGroup) GetRootDiskSizeGb() int {
if len(self.Subnets) == 0 {
self.Refresh()
}
return self.DiskSize
}
func (self *SNodeGroup) Delete() error {
return self.region.DeleteNodegroup(self.ClusterName, self.NodegroupName)
}
func (self *SNodeGroup) GetNetworkIds() []string {
if len(self.Subnets) == 0 {
self.Refresh()
}
return self.Subnets
}
func (self *SNodeGroup) GetInstanceTypes() []string {
if len(self.InstanceTypes) == 0 {
self.Refresh()
}
return self.InstanceTypes
}
func (self *SRegion) GetNodegroup(cluster, name string) (*SNodeGroup, error) {
params := map[string]interface{}{
"name": cluster,
"nodegroupName": name,
}
ret := struct {
Nodegroup SNodeGroup
}{}
err := self.eksRequest("DescribeNodegroup", "/clusters/{name}/node-groups/{nodegroupName}", params, &ret)
if err != nil {
return nil, errors.Wrapf(err, "DescribeNodegroup")
}
ret.Nodegroup.region = self
return &ret.Nodegroup, nil
}
func (self *SRegion) GetNodegroups(cluster, nextToken string) ([]SNodeGroup, string, error) {
params := map[string]interface{}{}
if len(nextToken) > 0 {
params["nextToken"] = nextToken
}
ret := struct {
Nodegroups []string
NextToken string
}{}
resource := fmt.Sprintf("/clusters/%s/node-groups", cluster)
err := self.eksRequest("ListNodegroups", resource, params, &ret)
if err != nil {
return nil, "", errors.Wrapf(err, "DescribeCluster")
}
result := []SNodeGroup{}
for i := range ret.Nodegroups {
result = append(result, SNodeGroup{
region: self,
ClusterName: cluster,
NodegroupName: ret.Nodegroups[i],
})
}
return result, ret.NextToken, nil
}
func (self *SRegion) DeleteNodegroup(cluster, name string) error {
params := map[string]interface{}{
"name": cluster,
"nodegroupName": name,
}
ret := struct {
Nodegroup SNodeGroup
}{}
return self.eksRequest("DeleteNodegroup", "/clusters/{name}/node-groups/{nodegroupName}", params, &ret)
}
+2
View File
@@ -28,6 +28,8 @@ import (
var (
samlRole = `{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Action":"sts:AssumeRoleWithSAML","Principal":{"Federated":"%s"},"Condition":{"StringEquals":{"SAML:aud":["%s"]}}}]}`
k8sRole = `{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":{"Service":"eks.amazonaws.com"},"Action":"sts:AssumeRole"}]}`
nodeRole = `{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":{"Service":"ec2.amazonaws.com.cn"},"Action":"sts:AssumeRole"}]}`
)
type SRole struct {
+7
View File
@@ -142,6 +142,9 @@ const (
ELB_SERVICE_NAME = "elasticloadbalancing"
ELB_SERVICE_ID = "Elastic Load Balancing v2"
EKS_SERVICE_NAME = "eks"
EKS_SERVICE_ID = "EKS"
)
type SRegion struct {
@@ -270,6 +273,10 @@ func (self *SAwsClient) monitorRequest(regionId, apiName string, params map[stri
return self.request(regionId, CLOUDWATCH_SERVICE_NAME, CLOUDWATCH_SERVICE_ID, "2010-08-01", apiName, params, retval, true)
}
func (self *SRegion) eksRequest(apiName, path string, params map[string]interface{}, retval interface{}) error {
return self.client.invoke(self.RegionId, EKS_SERVICE_NAME, EKS_SERVICE_ID, "2017-11-01", apiName, path, params, retval, true)
}
/////////////////////////////////////////////////////////////////////////////
func (self *SRegion) fetchZones() error {
ec2Client, err := self.getEc2Client()
+16
View File
@@ -161,6 +161,18 @@ func (self *SKubeCluster) GetKubeConfig(private bool, expireMinute int) (*cloudp
return self.region.GetKubeConfig(self.Id)
}
func (self *SKubeCluster) GetVersion() string {
return self.Properties.KubernetesVersion
}
func (self *SKubeCluster) GetVpcId() string {
return ""
}
func (self *SKubeCluster) GetNetworkIds() []string {
return []string{}
}
func (self *SRegion) GetICloudKubeClusters() ([]cloudprovider.ICloudKubeCluster, error) {
clusters, err := self.GetKubeClusters()
if err != nil {
@@ -218,3 +230,7 @@ func (self *SRegion) GetKubeConfig(id string) (*cloudprovider.SKubeconfig, error
result.Config = string(config)
return result, err
}
func (self *SKubeCluster) CreateIKubeNodePool(opts *cloudprovider.KubeNodePoolCreateOptions) (cloudprovider.ICloudKubeNodePool, error) {
return nil, cloudprovider.ErrNotImplemented
}
+29
View File
@@ -17,6 +17,7 @@ package azure
import (
"strings"
"yunion.io/x/cloudmux/pkg/cloudprovider"
"yunion.io/x/cloudmux/pkg/multicloud"
)
@@ -62,3 +63,31 @@ func (self *SKubeNodePool) GetGlobalId() string {
func (self *SKubeNodePool) GetStatus() string {
return strings.ToLower(self.PowerState.Code)
}
func (self *SKubeNodePool) GetMinInstanceCount() int {
return 0
}
func (self *SKubeNodePool) GetMaxInstanceCount() int {
return 0
}
func (self *SKubeNodePool) GetDesiredInstanceCount() int {
return 0
}
func (self *SKubeNodePool) GetRootDiskSizeGb() int {
return self.OsDiskSizeGB
}
func (self *SKubeNodePool) Delete() error {
return cloudprovider.ErrNotImplemented
}
func (self *SKubeNodePool) GetInstanceTypes() []string {
return []string{}
}
func (self *SKubeNodePool) GetNetworkIds() []string {
return []string{}
}
+1 -1
View File
@@ -36,7 +36,7 @@ func (h *SHost) GetId() string {
}
func (h *SHost) GetName() string {
return fmt.Sprintf("%s-%s", h.zone.region.cpcfg.Name, h.zone.GetName())
return fmt.Sprintf("%s-%s", h.zone.region.Name, h.zone.GetName())
}
func (h *SHost) GetGlobalId() string {
+1
View File
@@ -32,6 +32,7 @@ const (
JDCLOUD_DEFAULT_REGION = "cn-north-1"
)
// https://docs.jdcloud.com/cn/common-declaration/api/introduction
var regionList = map[string]string{
"cn-north-1": "华北-北京",
"cn-east-1": "华东-宿迁",
+2 -1
View File
@@ -17,6 +17,7 @@ package jdcloud
import (
"fmt"
api "yunion.io/x/cloudmux/pkg/apis/compute"
"yunion.io/x/cloudmux/pkg/cloudprovider"
"yunion.io/x/cloudmux/pkg/multicloud"
)
@@ -102,7 +103,7 @@ func (z *SZone) GetGlobalId() string {
}
func (z *SZone) GetStatus() string {
return "enable"
return api.ZONE_ENABLE
}
func (z *SZone) Refresh() error {
+18
View File
@@ -37,6 +37,8 @@ type SKubeCluster struct {
ClusterOs string
ClusterType string
ClusterNetworkSettings struct {
VpcId string
Subnets []string
}
ClusterNodeNum int
ProjectId string
@@ -81,6 +83,18 @@ func (self *SKubeCluster) Refresh() error {
return jsonutils.Update(self, cluster)
}
func (self *SKubeCluster) GetVersion() string {
return self.ClusterVersion
}
func (self *SKubeCluster) GetVpcId() string {
return self.ClusterNetworkSettings.VpcId
}
func (self *SKubeCluster) GetNetworkIds() []string {
return self.ClusterNetworkSettings.Subnets
}
func (self *SKubeCluster) GetKubeConfig(private bool, expireMinutes int) (*cloudprovider.SKubeconfig, error) {
return self.region.GetKubeConfig(self.ClusterId, private)
}
@@ -187,3 +201,7 @@ func (self *SRegion) DeleteKubeCluster(id string, isRetain bool) error {
_, err := self.tkeRequest("DeleteCluster", params)
return errors.Wrapf(err, "DeleteCluster")
}
func (self *SKubeCluster) CreateIKubeNodePool(opts *cloudprovider.KubeNodePoolCreateOptions) (cloudprovider.ICloudKubeNodePool, error) {
return nil, cloudprovider.ErrNotImplemented
}
+28
View File
@@ -83,6 +83,34 @@ func (self *SKubeNodePool) GetStatus() string {
return self.LifeState
}
func (self *SKubeNodePool) GetMinInstanceCount() int {
return self.MinNodesNum
}
func (self *SKubeNodePool) GetMaxInstanceCount() int {
return self.MaxNodesNum
}
func (self *SKubeNodePool) GetDesiredInstanceCount() int {
return self.DesiredNodesNum
}
func (self *SKubeNodePool) GetRootDiskSizeGb() int {
return 0
}
func (self *SKubeNodePool) Delete() error {
return cloudprovider.ErrNotImplemented
}
func (self *SKubeNodePool) GetInstanceTypes() []string {
return []string{}
}
func (self *SKubeNodePool) GetNetworkIds() []string {
return []string{}
}
func (self *SKubeCluster) GetIKubeNodePools() ([]cloudprovider.ICloudKubeNodePool, error) {
pools, err := self.region.GetKubeNodePools(self.ClusterId)
if err != nil {
+4
View File
@@ -247,6 +247,10 @@ func (self *SRegion) GetICloudKubeClusters() ([]cloudprovider.ICloudKubeCluster,
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "GetICloudKubeClusters")
}
func (self *SRegion) CreateIKubeCluster(opts *cloudprovider.KubeClusterCreateOptions) (cloudprovider.ICloudKubeCluster, error) {
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "CreateIKubeCluster")
}
func (self *SRegion) GetICloudKubeClusterById(id string) (cloudprovider.ICloudKubeCluster, error) {
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "GetICloudKubeClusterById")
}
+68
View File
@@ -0,0 +1,68 @@
// Copyright 2019 Yunion
//
// 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 sqlchemy
import (
"fmt"
"reflect"
"strings"
"yunion.io/x/log"
)
func getSQLFilters(filter map[string]interface{}) ([]string, []interface{}) {
conds := make([]string, 0, len(filter))
params := make([]interface{}, 0, len(filter))
for k, v := range filter {
if reflect.TypeOf(v).Kind() == reflect.Slice || reflect.TypeOf(v).Kind() == reflect.Array {
value := reflect.ValueOf(v)
if value.Len() == 0 {
continue
}
arr := make([]string, value.Len())
for i := 0; i < value.Len(); i++ {
arr[i] = "?"
params = append(params, value.Index(i).Interface())
}
conds = append(conds, fmt.Sprintf("`%s` in (%s)", k, strings.Join(arr, ", ")))
} else {
conds = append(conds, fmt.Sprintf("`%s` = ?", k))
params = append(params, v)
}
}
return conds, params
}
func (ts *STableSpec) DeleteFrom(filters map[string]interface{}) error {
buf := strings.Builder{}
buf.WriteString("DELETE FROM `")
buf.WriteString(ts.Name())
buf.WriteString("`")
conds, params := getSQLFilters(filters)
if len(conds) > 0 {
buf.WriteString(" WHERE ")
buf.WriteString(strings.Join(conds, " AND "))
}
if DEBUG_SQLCHEMY {
log.Infof("Update: %s %s", buf.String(), params)
}
_, err := ts.Database().TxExec(buf.String(), params...)
return err
}
+56
View File
@@ -0,0 +1,56 @@
// Copyright 2019 Yunion
//
// 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 sqlchemy
import (
"fmt"
"strings"
"yunion.io/x/log"
)
func (ts *STableSpec) UpdateBatch(data map[string]interface{}, filter map[string]interface{}) error {
if len(data) <= 0 {
return nil
}
params := make([]interface{}, 0, len(data))
setter := make([]string, 0, len(data))
for k, v := range data {
setter = append(setter, fmt.Sprintf("`%s` = ?", k))
params = append(params, v)
}
conds, condparams := getSQLFilters(filter)
params = append(params, condparams...)
buf := strings.Builder{}
buf.WriteString("UPDATE `")
buf.WriteString(ts.Name())
buf.WriteString("` SET ")
buf.WriteString(strings.Join(setter, ", "))
if len(conds) > 0 {
buf.WriteString(" WHERE ")
buf.WriteString(strings.Join(conds, " AND "))
}
if DEBUG_SQLCHEMY {
log.Infof("Update: %s %s", buf.String(), params)
}
_, err := ts.Database().Exec(buf.String(), params...)
return err
}