fix(region): support ocean base (#23587)

This commit is contained in:
屈轩
2025-10-22 14:43:32 +08:00
committed by GitHub
parent b1ef423c2a
commit a9649db617
23 changed files with 2061 additions and 35 deletions
+3
View File
@@ -63,6 +63,7 @@ func init() {
cmd.CreateWithKeyword("create-oracle", &options.SOracleCloudAccountCreateOptions{})
cmd.CreateWithKeyword("create-cephfs", &options.SCephFSCloudAccountCreateOptions{})
cmd.CreateWithKeyword("create-cnware", &options.SCNwareCloudAccountCreateOptions{})
cmd.CreateWithKeyword("create-oceanbase", &options.SOceanbaseCloudAccountCreateOptions{})
cmd.UpdateWithKeyword("update-vmware", &options.SVMwareCloudAccountUpdateOptions{})
cmd.UpdateWithKeyword("update-aliyun", &options.SAliyunCloudAccountUpdateOptions{})
@@ -90,6 +91,7 @@ func init() {
cmd.UpdateWithKeyword("update-cucloud", &options.SCucloudCloudAccountUpdateOptions{})
cmd.UpdateWithKeyword("update-qingcloud", &options.SQingCloudCloudAccountUpdateOptions{})
cmd.UpdateWithKeyword("update-cnware", &options.SCNwareCloudAccountUpdateOptions{})
cmd.UpdateWithKeyword("update-oceanbase", &options.SOceanbaseCloudAccountUpdateOptions{})
cmd.Perform("update-credential", &options.CloudaccountUpdateCredentialOptions{})
@@ -136,6 +138,7 @@ func init() {
cmd.PerformWithKeyword("test-connectivity-ctyun", "test-connectivity", &options.SCtyunCloudAccountUpdateCredentialOptions{})
cmd.PerformWithKeyword("test-connectivity-jdcloud", "test-connectivity", &options.SJDcloudCloudAccountUpdateCredentialOptions{})
cmd.PerformWithKeyword("test-connectivity-cloudpods", "test-connectivity", &options.SCloudpodsCloudAccountUpdateCredentialOptions{})
cmd.PerformWithKeyword("test-connectivity-oceanbase", "test-connectivity", &options.SOceanbaseCloudAccountUpdateCredentialOptions{})
cmd.Perform("enable", &options.SCloudAccountIdOptions{})
cmd.Perform("disable", &options.SCloudAccountIdOptions{})
+2 -2
View File
@@ -96,7 +96,7 @@ require (
k8s.io/cri-api v0.22.17
k8s.io/klog/v2 v2.20.0
moul.io/http2curl/v2 v2.3.0
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20251022024201-332dfd944a29
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20251022062416-16b2a67e2e75
yunion.io/x/executor v0.0.0-20250518005516-5402e9e0bed0
yunion.io/x/jsonutils v1.0.1-0.20250507052344-1abcf4f443b1
yunion.io/x/log v1.0.1-0.20240305175729-7cf2d6cd5a91
@@ -237,6 +237,7 @@ require (
github.com/grpc-ecosystem/grpc-opentracing v0.0.0-20180507213350-8e809c8a8645 // indirect
github.com/huandu/xstrings v1.2.0 // indirect
github.com/huaweicloud/huaweicloud-sdk-go v1.0.26 // indirect
github.com/icholy/digest v1.1.0 // indirect
github.com/imdario/mergo v0.3.6 // indirect
github.com/jdcloud-api/jdcloud-sdk-go v1.55.0 // indirect
github.com/jmespath/go-jmespath v0.4.0 // indirect
@@ -344,7 +345,6 @@ require (
gopkg.in/inf.v0 v0.9.1 // indirect
gopkg.in/ini.v1 v1.62.0 // indirect
gopkg.in/yaml.v3 v3.0.1 // indirect
gotest.tools/v3 v3.5.1 // indirect
k8s.io/utils v0.0.0-20201110183641-67b214c5f920 // indirect
sigs.k8s.io/structured-merge-diff/v4 v4.0.1 // indirect
sigs.k8s.io/yaml v1.2.0 // indirect
+5 -2
View File
@@ -563,6 +563,8 @@ github.com/huandu/xstrings v1.2.0/go.mod h1:DvyZB1rfVYsBIigL8HwpZgxHwXozlTgGqn63
github.com/huaweicloud/huaweicloud-sdk-go v1.0.26 h1:aWTl4Ng9lZUW1DupYyqYiqkTwuTKMmqUZBKLmX0iWxo=
github.com/huaweicloud/huaweicloud-sdk-go v1.0.26/go.mod h1:YHXxw/bm7AohI4jTY8Z43JYw+DPCDLqjxoLTqmcyja4=
github.com/ianlancetaylor/demangle v0.0.0-20181102032728-5e5cf60278f6/go.mod h1:aSSvb/t6k1mPoxDqO4vJh6VOCGPwU4O0C2/Eqndh1Sc=
github.com/icholy/digest v1.1.0 h1:HfGg9Irj7i+IX1o1QAmPfIBNu/Q5A5Tu3n/MED9k9H4=
github.com/icholy/digest v1.1.0/go.mod h1:QNrsSGQ5v7v9cReDI0+eyjsXGUoRSUZQHeQ5C4XLa0Y=
github.com/imdario/mergo v0.3.5/go.mod h1:2EnlNZ0deacrJVfApfmtdGgDfMuh/nq6Ok1EcJh5FfA=
github.com/imdario/mergo v0.3.6 h1:xTNEAn+kxVO7dTZGu0CegyqKZmoWFI0rF8UxjlB2d28=
github.com/imdario/mergo v0.3.6/go.mod h1:2EnlNZ0deacrJVfApfmtdGgDfMuh/nq6Ok1EcJh5FfA=
@@ -1416,6 +1418,7 @@ gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C
gopkg.in/yaml.v3 v3.0.0-20210107192922-496545a6307b/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
gotest.tools v2.2.0+incompatible h1:VsBPFP1AI068pPrMxtb/S8Zkgf9xEmTLJjfM+P5UIEo=
gotest.tools v2.2.0+incompatible/go.mod h1:DsYFclhRJ6vuDpmuTbkuFWG+y2sxOXAzmJt81HFBacw=
gotest.tools/v3 v3.5.1 h1:EENdUnS3pdur5nybKYIh2Vfgc8IUNBjxDPSjtiJcOzU=
gotest.tools/v3 v3.5.1/go.mod h1:isy3WKz7GK6uNw/sbHzfKBLvlvXwUyV06n6brMxxopU=
@@ -1455,8 +1458,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.20251022024201-332dfd944a29 h1:AFF6zLcHleDwLkRJZkTKBr0e7CRpRTGCUgzRsKLIhf0=
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20251022024201-332dfd944a29/go.mod h1:qn8eC11Gn4/41YJRzQcwB/G1Z3rQRxRrDNC8LDvQR48=
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20251022062416-16b2a67e2e75 h1:gqdyuR8AMGahxxzrtmxdeiEIM2cZB9QhWjOP3416VUA=
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20251022062416-16b2a67e2e75/go.mod h1:nya/IL1IXSMmFgB48P27CppyGpkT5coG150PD+xVJrc=
yunion.io/x/executor v0.0.0-20250518005516-5402e9e0bed0 h1:msG4SiDSVU7CrXH06WuHlNEZXIooTcmNbfrIGHuIHBU=
yunion.io/x/executor v0.0.0-20250518005516-5402e9e0bed0/go.mod h1:Uxuou9WQIeJXNpy7t2fPLL0BYLvLiMvGQwY7Qc6aSws=
yunion.io/x/jsonutils v0.0.0-20190625054549-a964e1e8a051/go.mod h1:4N0/RVzsYL3kH3WE/H1BjUQdFiWu50JGCFQuuy+Z634=
+2
View File
@@ -77,6 +77,7 @@ const (
CLOUD_PROVIDER_CAS = compute.CLOUD_PROVIDER_CAS
CLOUD_PROVIDER_CLOUDFLARE = compute.CLOUD_PROVIDER_CLOUDFLARE
CLOUD_PROVIDER_CNWARE = compute.CLOUD_PROVIDER_CNWARE
CLOUD_PROVIDER_OCEANBASE = compute.CLOUD_PROVIDER_OCEANBASE
CLOUD_PROVIDER_GENERICS3 = compute.CLOUD_PROVIDER_GENERICS3
CLOUD_PROVIDER_CEPH = compute.CLOUD_PROVIDER_CEPH
@@ -186,6 +187,7 @@ var (
CLOUD_PROVIDER_CAS,
CLOUD_PROVIDER_CLOUDFLARE,
CLOUD_PROVIDER_CNWARE,
CLOUD_PROVIDER_OCEANBASE,
}
CLOUD_PROVIDER_HOST_TYPE_MAP = map[string][]string{
+8 -28
View File
@@ -1623,13 +1623,13 @@ func (self *SDBInstance) SetZoneIds(extInstance cloudprovider.ICloudDBInstance)
return errors.Wrapf(err, "GetZones")
}
var setZoneId = func(input string, output *string) {
*output = input
for _, zone := range zones {
if strings.HasSuffix(zone.ExternalId, input) {
*output = zone.Id
break
}
}
return
}
zone1 := extInstance.GetZone1Id()
if len(zone1) > 0 {
@@ -1705,26 +1705,13 @@ func (self *SDBInstance) SyncWithCloudDBInstance(ctx context.Context, userCred m
self.CreatedAt = createdAt
}
if len(self.VpcId) == 0 {
if vpcId := ext.GetIVpcId(); len(vpcId) > 0 {
vpc, err := db.FetchByExternalIdAndManagerId(VpcManager, vpcId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", provider.Id)
})
if err != nil {
log.Errorf("FetchVpcId(%s) error: %v", vpcId, err)
} else {
self.VpcId = vpc.GetId()
}
}
}
if len(self.VpcId) == 0 {
region, err := self.GetRegion()
if vpcId := ext.GetIVpcId(); len(vpcId) > 0 {
self.VpcId = vpcId
vpc, err := db.FetchByExternalIdAndManagerId(VpcManager, vpcId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", provider.Id)
})
if err != nil {
return err
}
vpc, err := VpcManager.GetOrCreateVpcForClassicNetwork(ctx, provider, region)
if err != nil {
log.Errorf("failed to create classic vpc for region %s error: %v", region.Name, err)
log.Errorf("FetchVpcId(%s) error: %v", vpcId, err)
} else {
self.VpcId = vpc.GetId()
}
@@ -1792,6 +1779,7 @@ func (manager *SDBInstanceManager) newFromCloudDBInstance(ctx context.Context, u
instance.SetZoneIds(extInstance)
if vpcId := extInstance.GetIVpcId(); len(vpcId) > 0 {
instance.VpcId = vpcId
vpc, err := db.FetchByExternalIdAndManagerId(VpcManager, vpcId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", provider.Id)
})
@@ -1801,14 +1789,6 @@ func (manager *SDBInstanceManager) newFromCloudDBInstance(ctx context.Context, u
instance.VpcId = vpc.GetId()
}
}
if len(instance.VpcId) == 0 {
vpc, err := VpcManager.GetOrCreateVpcForClassicNetwork(ctx, provider, region)
if err != nil {
log.Errorf("failed to create classic vpc for region %s error: %v", region.Name, err)
} else {
instance.VpcId = vpc.GetId()
}
}
if createdAt := extInstance.GetCreatedAt(); !createdAt.IsZero() {
instance.CreatedAt = createdAt
+33
View File
@@ -0,0 +1,33 @@
// 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 regiondrivers
import (
api "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/compute/models"
)
type SOceanbaseRegionDriver struct {
SManagedVirtualizationRegionDriver
}
func init() {
driver := SOceanbaseRegionDriver{}
models.RegisterRegionDriver(&driver)
}
func (self *SOceanbaseRegionDriver) GetProvider() string {
return api.CLOUD_PROVIDER_OCEANBASE
}
@@ -1599,3 +1599,33 @@ type SCNwareCloudAccountUpdateOptions struct {
func (opts *SCNwareCloudAccountUpdateOptions) Params() (jsonutils.JSONObject, error) {
return jsonutils.Marshal(opts), nil
}
type SOceanbaseCloudAccountCreateOptions struct {
SCloudAccountCreateBaseOptions
SAccessKeyCredential
}
func (opts *SOceanbaseCloudAccountCreateOptions) Params() (jsonutils.JSONObject, error) {
params := jsonutils.Marshal(opts)
params.(*jsonutils.JSONDict).Add(jsonutils.NewString("OceanBase"), "provider")
return params, nil
}
type SOceanbaseCloudAccountUpdateOptions struct {
SCloudAccountUpdateBaseOptions
}
func (opts *SOceanbaseCloudAccountUpdateOptions) Params() (jsonutils.JSONObject, error) {
params := jsonutils.Marshal(opts).(*jsonutils.JSONDict)
return params, nil
}
type SOceanbaseCloudAccountUpdateCredentialOptions struct {
SCloudAccountIdOptions
SAccessKeyCredential
}
func (opts *SOceanbaseCloudAccountUpdateCredentialOptions) Params() (jsonutils.JSONObject, error) {
return jsonutils.Marshal(opts), nil
}
+1
View File
@@ -0,0 +1 @@
/.idea
+21
View File
@@ -0,0 +1,21 @@
MIT License
Copyright (c) 2020 Ilia Choly
Permission is hereby granted, free of charge, to any person obtaining a copy
of this software and associated documentation files (the "Software"), to deal
in the Software without restriction, including without limitation the rights
to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
copies of the Software, and to permit persons to whom the Software is
furnished to do so, subject to the following conditions:
The above copyright notice and this permission notice shall be included in all
copies or substantial portions of the Software.
THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
SOFTWARE.
+164
View File
@@ -0,0 +1,164 @@
# HTTP Digest Access Authentication
[![go.dev reference](https://img.shields.io/badge/go.dev-reference-007d9c?logo=go&logoColor=white&style=flat-square)](https://pkg.go.dev/github.com/icholy/digest)
> This package provides a http.RoundTripper implementation which re-uses digest challenges
``` go
package main
import (
"net/http"
"github.com/icholy/digest"
)
func main() {
client := &http.Client{
Transport: &digest.Transport{
Username: "foo",
Password: "bar",
},
}
res, err := client.Get("http://localhost:8080/some_outdated_service")
if err != nil {
panic(err)
}
defer res.Body.Close()
}
```
## Using Cookies
If you're using an `http.CookieJar` the `digest.Transport` needs a reference to it.
``` go
package main
import (
"net/http"
"net/http/cookiejar"
"github.com/icholy/digest"
)
func main() {
jar, _ := cookiejar.New(nil)
client := &http.Client{
Transport: &digest.Transport{
Jar: jar,
Username: "foo",
Password: "bar",
},
}
res, err := client.Get("http://localhost:8080/digest_with_cookies")
if err != nil {
panic(err)
}
defer res.Body.Close()
}
```
## Custom Authenticate Header
``` go
package main
import (
"net/http"
"github.com/icholy/digest"
)
func main() {
client := &http.Client{
Transport: &digest.Transport{
Username: "foo",
Password: "bar",
FindChallenge: func(h http.Header) (*digest.Challenge, error) {
value := h.Get("Custom-Authenticate-Header")
if value == "" {
return nil, digest.ErrNoChallenge
}
return digest.ParseChallenge(value)
},
},
}
res, err := client.Get("http://localhost:8080/non_compliant")
if err != nil {
panic(err)
}
defer res.Body.Close()
}
```
## Override Digest Options
``` go
package main
import (
"fmt"
"net/http"
"github.com/icholy/digest"
)
func main() {
client := &http.Client{
Transport: &digest.Transport{
Digest: func(req *http.Request, chal *digest.Challenge, opt digest.Options) (*digest.Credentials, error) {
switch req.URL.Hostname() {
case "badauth.org":
opt.Username = "foo"
opt.Password = "bar"
case "poorsecurity.com":
opt.Username = "zoo"
opt.Password = "boo"
default:
return nil, fmt.Errorf("unsuported host: %q", req.URL)
}
return digest.Digest(chal, opt)
},
},
}
res, err := client.Get("http://poorsecurity.com/legacy.php")
if err != nil {
panic(err)
}
defer res.Body.Close()
}
```
## Low Level API
``` go
func main() {
// get the challenge from a 401 response
header := res.Header.Get("WWW-Authenticate")
chal, _ := digest.ParseChallenge(header)
// use it to create credentials for the next request
cred, _ := digest.Digest(chal, digest.Options{
Username: "foo",
Password: "bar",
Method: req.Method,
URI: req.URL.RequestURI(),
GetBody: req.GetBody,
Count: 1,
})
req.Header.Set("Authorization", cred.String())
// if you use the same challenge again, you must increment the Count
cred2, _ := digest.Digest(chal, digest.Options{
Username: "foo",
Password: "bar",
Method: req2.Method,
URI: req2.URL.RequestURI(),
GetBody: req2.GetBody,
Count: 2,
})
req2.Header.Set("Authorization", cred.String())
}
```
+155
View File
@@ -0,0 +1,155 @@
package digest
import (
"errors"
"fmt"
"net/http"
"strings"
"github.com/icholy/digest/internal/param"
)
// Challenge is a challenge sent in the WWW-Authenticate header
type Challenge struct {
Realm string
Domain []string
Nonce string
Opaque string
Stale bool
Algorithm string
QOP []string
Charset string
Userhash bool
}
// SupportsQOP returns true if the challenge advertises support
// for the provided qop value
func (c *Challenge) SupportsQOP(qop string) bool {
for _, v := range c.QOP {
if v == qop {
return true
}
}
return false
}
// ParseChallenge parses the WWW-Authenticate header challenge
func ParseChallenge(s string) (*Challenge, error) {
s, ok := strings.CutPrefix(s, Prefix)
if !ok {
return nil, errors.New("digest: invalid challenge prefix")
}
pp, err := param.Parse(s)
if err != nil {
return nil, fmt.Errorf("digest: invalid challenge: %w", err)
}
var c Challenge
for _, p := range pp {
switch p.Key {
case "realm":
c.Realm = p.Value
case "domain":
c.Domain = strings.Fields(p.Value)
case "nonce":
c.Nonce = p.Value
case "algorithm":
c.Algorithm = p.Value
case "stale":
c.Stale = strings.ToLower(p.Value) == "true"
case "opaque":
c.Opaque = p.Value
case "qop":
c.QOP = strings.Split(p.Value, ",")
case "charset":
c.Charset = p.Value
case "userhash":
c.Userhash = strings.ToLower(p.Value) == "true"
}
}
return &c, nil
}
// String returns the foramtted header value
func (c *Challenge) String() string {
var pp []param.Param
pp = append(pp, param.Param{
Key: "realm",
Value: c.Realm,
Quote: true,
})
if len(c.Domain) != 0 {
pp = append(pp, param.Param{
Key: "domain",
Value: strings.Join(c.Domain, " "),
Quote: true,
})
}
pp = append(pp, param.Param{
Key: "nonce",
Value: c.Nonce,
Quote: true,
})
if c.Opaque != "" {
pp = append(pp, param.Param{
Key: "opaque",
Value: c.Opaque,
Quote: true,
})
}
if c.Stale {
pp = append(pp, param.Param{
Key: "stale",
Value: "true",
})
}
if c.Algorithm != "" {
pp = append(pp, param.Param{
Key: "algorithm",
Value: c.Algorithm,
})
}
if len(c.QOP) != 0 {
pp = append(pp, param.Param{
Key: "qop",
Value: strings.Join(c.QOP, ","),
Quote: true,
})
}
if c.Charset != "" {
pp = append(pp, param.Param{
Key: "charset",
Value: c.Charset,
})
}
if c.Userhash {
pp = append(pp, param.Param{
Key: "userhash",
Value: "true",
})
}
return Prefix + param.Format(pp...)
}
// ErrNoChallenge indicates that no WWW-Authenticate headers were found.
var ErrNoChallenge = errors.New("digest: no challenge found")
// FindChallenge returns the first supported challenge in the headers
func FindChallenge(h http.Header) (*Challenge, error) {
var last error
for _, header := range h.Values("WWW-Authenticate") {
if !IsDigest(header) {
continue
}
chal, err := ParseChallenge(header)
if err == nil && CanDigest(chal) {
return chal, nil
}
if err != nil {
last = err
}
}
if last != nil {
return nil, last
}
return nil, ErrNoChallenge
}
+142
View File
@@ -0,0 +1,142 @@
package digest
import (
"errors"
"fmt"
"strconv"
"strings"
"github.com/icholy/digest/internal/param"
)
// Credentials is a parsed version of the Authorization header
type Credentials struct {
Username string
Realm string
Nonce string
URI string
Response string
Algorithm string
Cnonce string
Opaque string
QOP string
Nc int
Userhash bool
}
// ParseCredentials parses the Authorization header value into credentials
func ParseCredentials(s string) (*Credentials, error) {
s, ok := strings.CutPrefix(s, Prefix)
if !ok {
return nil, errors.New("digest: invalid credentials prefix")
}
pp, err := param.Parse(s)
if err != nil {
return nil, fmt.Errorf("digest: invalid credentials: %w", err)
}
var c Credentials
for _, p := range pp {
switch p.Key {
case "username":
c.Username = p.Value
case "realm":
c.Realm = p.Value
case "nonce":
c.Nonce = p.Value
case "uri":
c.URI = p.Value
case "response":
c.Response = p.Value
case "algorithm":
c.Algorithm = p.Value
case "cnonce":
c.Cnonce = p.Value
case "opaque":
c.Opaque = p.Value
case "qop":
c.QOP = p.Value
case "nc":
nc, err := strconv.ParseInt(p.Value, 16, 32)
if err != nil {
return nil, fmt.Errorf("digest: invalid nc: %w", err)
}
c.Nc = int(nc)
case "userhash":
c.Userhash = strings.ToLower(p.Value) == "true"
}
}
return &c, nil
}
// String formats the credentials into the header format
func (c *Credentials) String() string {
var pp []param.Param
pp = append(pp,
param.Param{
Key: "username",
Value: c.Username,
Quote: true,
},
param.Param{
Key: "realm",
Value: c.Realm,
Quote: true,
},
param.Param{
Key: "nonce",
Value: c.Nonce,
Quote: true,
},
param.Param{
Key: "uri",
Value: c.URI,
Quote: true,
},
)
if c.Algorithm != "" {
pp = append(pp, param.Param{
Key: "algorithm",
Value: c.Algorithm,
})
}
if c.QOP != "" {
pp = append(pp, param.Param{
Key: "cnonce",
Value: c.Cnonce,
Quote: true,
})
}
if c.Opaque != "" {
pp = append(pp, param.Param{
Key: "opaque",
Value: c.Opaque,
Quote: true,
})
}
if c.QOP != "" {
pp = append(pp,
param.Param{
Key: "qop",
Value: c.QOP,
},
param.Param{
Key: "nc",
Value: fmt.Sprintf("%08x", c.Nc),
},
)
}
if c.Userhash {
pp = append(pp, param.Param{
Key: "userhash",
Value: "true",
})
}
// The RFC does not specify an order, but some implementations expect the response to be at the end.
// See: https://github.com/icholy/digest/issues/8
pp = append(pp, param.Param{
Key: "response",
Value: c.Response,
Quote: true,
})
return Prefix + param.Format(pp...)
}
+164
View File
@@ -0,0 +1,164 @@
package digest
import (
"crypto/md5"
"crypto/rand"
"crypto/sha256"
"crypto/sha512"
"encoding/hex"
"fmt"
"hash"
"io"
"net/http"
"strings"
)
// Prefix for digest authentication headers
const Prefix = "Digest "
// IsDigest returns true if the header value is a digest auth header
func IsDigest(header string) bool {
return strings.HasPrefix(header, Prefix)
}
// Options for creating a credentials
type Options struct {
Method string
URI string
GetBody func() (io.ReadCloser, error)
Count int
Username string
Password string
// The following are provided for advanced use cases where the client needs
// to override the default digest calculation behavior. Most users should
// leave these fields unset.
A1 string
Cnonce string
}
// CanDigest checks if the algorithm and qop are supported
func CanDigest(c *Challenge) bool {
switch strings.ToUpper(c.Algorithm) {
case "", "MD5", "SHA-256", "SHA-512", "SHA-512-256":
default:
return false
}
return len(c.QOP) == 0 || c.SupportsQOP("auth") || c.SupportsQOP("auth-int")
}
// Digest creates credentials from a challenge and request options.
// Note: if you want to re-use a challenge, you must increment the Count.
func Digest(chal *Challenge, o Options) (*Credentials, error) {
cred := &Credentials{
Username: o.Username,
URI: o.URI,
Cnonce: o.Cnonce,
Nc: o.Count,
Realm: chal.Realm,
Nonce: chal.Nonce,
Algorithm: chal.Algorithm,
Opaque: chal.Opaque,
Userhash: chal.Userhash,
}
// we re-use the same hash.Hash
var h hash.Hash
switch strings.ToUpper(cred.Algorithm) {
case "", "MD5":
h = md5.New()
case "SHA-256":
h = sha256.New()
case "SHA-512":
h = sha512.New()
case "SHA-512-256":
h = sha512.New512_256()
default:
return nil, fmt.Errorf("digest: unsupported algorithm: %q", cred.Algorithm)
}
// hash the username if requested
if cred.Userhash {
cred.Username = hashf(h, "%s:%s", o.Username, cred.Realm)
}
// generate the a1 hash if one was not provided
a1 := o.A1
if a1 == "" {
a1 = hashf(h, "%s:%s:%s", o.Username, cred.Realm, o.Password)
}
// generate the response
switch {
case len(chal.QOP) == 0:
cred.Response = hashf(h, "%s:%s:%s",
a1,
cred.Nonce,
hashf(h, "%s:%s", o.Method, o.URI), // A2
)
case chal.SupportsQOP("auth"):
cred.QOP = "auth"
if cred.Cnonce == "" {
cred.Cnonce = cnonce()
}
if cred.Nc == 0 {
cred.Nc = 1
}
cred.Response = hashf(h, "%s:%s:%08x:%s:%s:%s",
a1,
cred.Nonce,
cred.Nc,
cred.Cnonce,
cred.QOP,
hashf(h, "%s:%s", o.Method, o.URI), // A2
)
case chal.SupportsQOP("auth-int"):
cred.QOP = "auth-int"
if cred.Cnonce == "" {
cred.Cnonce = cnonce()
}
if cred.Nc == 0 {
cred.Nc = 1
}
hbody, err := hashbody(h, o.GetBody)
if err != nil {
return nil, fmt.Errorf("digest: failed to read body for auth-int: %w", err)
}
cred.Response = hashf(h, "%s:%s:%08x:%s:%s:%s",
a1,
cred.Nonce,
cred.Nc,
cred.Cnonce,
cred.QOP,
hashf(h, "%s:%s:%s", o.Method, o.URI, hbody), // A2
)
default:
return nil, fmt.Errorf("digest: unsupported qop: %q", strings.Join(chal.QOP, ","))
}
return cred, nil
}
func hashf(h hash.Hash, format string, args ...interface{}) string {
h.Reset()
fmt.Fprintf(h, format, args...)
return hex.EncodeToString(h.Sum(nil))
}
func hashbody(h hash.Hash, getbody func() (io.ReadCloser, error)) (string, error) {
h.Reset()
if getbody != nil {
r, err := getbody()
if err != nil {
return "", err
}
defer r.Close()
if r != http.NoBody {
if _, err := io.Copy(h, r); err != nil {
return "", err
}
}
}
return hex.EncodeToString(h.Sum(nil)), nil
}
func cnonce() string {
b := make([]byte, 8)
io.ReadFull(rand.Reader, b)
return hex.EncodeToString(b)
}
+186
View File
@@ -0,0 +1,186 @@
package param
import (
"bufio"
"fmt"
"io"
"strings"
)
// Param is a key/value header parameter
type Param struct {
Key string
Value string
Quote bool
}
// String returns the formatted parameter
func (p Param) String() string {
if p.Quote {
return fmt.Sprintf("%s=%q", p.Key, p.Value)
}
return fmt.Sprintf("%s=%s", p.Key, p.Value)
}
// Format formats the parameters to be included in the header
func Format(pp ...Param) string {
var b strings.Builder
for i, p := range pp {
if i > 0 {
b.WriteString(", ")
}
b.WriteString(p.String())
}
return b.String()
}
// Parse parses the header parameters
func Parse(s string) ([]Param, error) {
var pp []Param
br := bufio.NewReader(strings.NewReader(s))
for i := 0; true; i++ {
// skip whitespace
if err := skipWhite(br); err != nil {
return nil, err
}
// see if there's more to read
if _, err := br.Peek(1); err == io.EOF {
break
}
// read key/value pair
p, err := parseParam(br, i == 0)
if err != nil {
return nil, fmt.Errorf("param: %w", err)
}
pp = append(pp, p)
}
return pp, nil
}
func parseIdent(br *bufio.Reader) (string, error) {
var ident []byte
for {
b, err := br.ReadByte()
if err == io.EOF {
break
}
if err != nil {
return "", err
}
if !(('a' <= b && b <= 'z') || ('A' <= b && b <= 'Z') || '0' <= b && b <= '9' || b == '-') {
if err := br.UnreadByte(); err != nil {
return "", err
}
break
}
ident = append(ident, b)
}
return string(ident), nil
}
func parseByte(br *bufio.Reader, expect byte) error {
b, err := br.ReadByte()
if err != nil {
if err == io.EOF {
return fmt.Errorf("expected '%c', got EOF", expect)
}
return err
}
if b != expect {
return fmt.Errorf("expected '%c', got '%c'", expect, b)
}
return nil
}
func parseString(br *bufio.Reader) (string, error) {
var s []rune
// read the open quote
if err := parseByte(br, '"'); err != nil {
return "", err
}
// read the string
var escaped bool
for {
r, _, err := br.ReadRune()
if err != nil {
return "", err
}
if escaped {
s = append(s, r)
escaped = false
continue
}
if r == '\\' {
escaped = true
continue
}
// closing quote
if r == '"' {
break
}
s = append(s, r)
}
return string(s), nil
}
func skipWhite(br *bufio.Reader) error {
for {
b, err := br.ReadByte()
if err != nil {
if err == io.EOF {
return nil
}
return err
}
if b != ' ' {
return br.UnreadByte()
}
}
}
func parseParam(br *bufio.Reader, first bool) (Param, error) {
// skip whitespace
if err := skipWhite(br); err != nil {
return Param{}, err
}
if !first {
// read the comma separator
if err := parseByte(br, ','); err != nil {
return Param{}, err
}
// skip whitespace
if err := skipWhite(br); err != nil {
return Param{}, err
}
}
// read the key
key, err := parseIdent(br)
if err != nil {
return Param{}, err
}
// skip whitespace
if err := skipWhite(br); err != nil {
return Param{}, err
}
// read the equals sign
if err := parseByte(br, '='); err != nil {
return Param{}, err
}
// skip whitespace
if err := skipWhite(br); err != nil {
return Param{}, err
}
// read the value
var value string
var quote bool
if b, _ := br.Peek(1); len(b) == 1 && b[0] == '"' {
quote = true
value, err = parseString(br)
} else {
value, err = parseIdent(br)
}
if err != nil {
return Param{}, err
}
return Param{Key: key, Value: value, Quote: quote}, nil
}
+237
View File
@@ -0,0 +1,237 @@
package digest
import (
"bytes"
"io"
"net/http"
"sync"
)
// cchal is a cached challenge and the number of times it's been used.
type cchal struct {
c *Challenge
n int
}
// Transport implements http.RoundTripper
type Transport struct {
Username string
Password string
// Digest computes the digest credentials.
// If nil, the Digest function is used.
Digest func(*http.Request, *Challenge, Options) (*Credentials, error)
// FindChallenge extracts the challenge from the request headers.
// If nil, the FindChallenge function is used.
FindChallenge func(http.Header) (*Challenge, error)
// Transport specifies the mechanism by which individual
// HTTP requests are made.
// If nil, DefaultTransport is used.
Transport http.RoundTripper
// Jar specifies the cookie jar.
//
// The Jar is used to insert relevant cookies into every
// outbound Request and is updated with the cookie values
// of every inbound Response. The Jar is consulted for every
// redirect that the Client follows.
//
// If Jar is nil, cookies are only sent if they are explicitly
// set on the Request.
Jar http.CookieJar
// NoReuse prevents the transport from reusing challenges.
NoReuse bool
// cache of challenges indexed by host
cache map[string]*cchal
cacheMu sync.Mutex
}
// save parses the digest challenge from the response
// and adds it to the cache
func (t *Transport) save(res *http.Response) error {
// save cookies
if t.Jar != nil {
t.Jar.SetCookies(res.Request.URL, res.Cookies())
}
// find and save digest challenge
find := t.FindChallenge
if find == nil {
find = FindChallenge
}
chal, err := find(res.Header)
t.cacheMu.Lock()
defer t.cacheMu.Unlock()
if t.cache == nil {
t.cache = map[string]*cchal{}
}
// TODO: if the challenge contains a domain, we should be using that
// to match against outgoing requests. We're currently ignoring
// it and just matching the hostname. That being said, none of
// the major browsers respect the domain either.
host := res.Request.URL.Hostname()
if err != nil {
// if save is being invoked, the existing cached challenge didn't work
delete(t.cache, host)
return err
}
t.cache[host] = &cchal{c: chal}
return nil
}
// digest creates credentials from the cached challenge
func (t *Transport) digest(req *http.Request, chal *Challenge, count int) (*Credentials, error) {
opt := Options{
Method: req.Method,
URI: req.URL.RequestURI(),
GetBody: req.GetBody,
Count: count,
Username: t.Username,
Password: t.Password,
}
if t.Digest != nil {
return t.Digest(req, chal, opt)
}
return Digest(chal, opt)
}
// challenge returns a cached challenge and count for the provided request
func (t *Transport) challenge(req *http.Request) (*Challenge, int, bool) {
t.cacheMu.Lock()
defer t.cacheMu.Unlock()
host := req.URL.Hostname()
cc, ok := t.cache[host]
if !ok {
return nil, 0, false
}
if t.NoReuse {
delete(t.cache, host)
}
cc.n++
return cc.c, cc.n, true
}
// prepare attempts to find a cached challenge that matches the
// requested domain, and use it to set the Authorization header
func (t *Transport) prepare(req *http.Request) error {
// add cookies
if t.Jar != nil {
for _, cookie := range t.Jar.Cookies(req.URL) {
req.AddCookie(cookie)
}
}
// add auth
chal, count, ok := t.challenge(req)
if !ok {
return nil
}
cred, err := t.digest(req, chal, count)
if err != nil {
return err
}
if cred != nil {
req.Header.Set("Authorization", cred.String())
}
return nil
}
// RoundTrip will try to authorize the request using a cached challenge.
// If that doesn't work and we receive a 401, we'll try again using that challenge.
func (t *Transport) RoundTrip(req *http.Request) (*http.Response, error) {
// use the configured transport if there is one
tr := t.Transport
if tr == nil {
tr = http.DefaultTransport
}
// don't modify the original request
clone, err := cloner(req)
if err != nil {
return nil, err
}
// make a copy of the request
first, err := clone()
if err != nil {
return nil, err
}
// prepare the first request using a cached challenge
if err := t.prepare(first); err != nil {
return nil, err
}
// the first request will either succeed or return a 401
res, err := tr.RoundTrip(first)
if err != nil || res.StatusCode != http.StatusUnauthorized {
return res, err
}
// drain and close the first message body
_, _ = io.Copy(io.Discard, res.Body)
_ = res.Body.Close()
// save the challenge for future use
if err := t.save(res); err != nil {
if err == ErrNoChallenge {
return res, nil
}
return nil, err
}
// make a second copy of the request
second, err := clone()
if err != nil {
return nil, err
}
// prepare the second request based on the new challenge
if err := t.prepare(second); err != nil {
return nil, err
}
return tr.RoundTrip(second)
}
// CloseIdleConnections delegates the call to the underlying transport.
func (t *Transport) CloseIdleConnections() {
tr := t.Transport
if tr == nil {
tr = http.DefaultTransport
}
type closeIdler interface {
CloseIdleConnections()
}
if tr, ok := tr.(closeIdler); ok {
tr.CloseIdleConnections()
}
}
// cloner returns a function which makes clones of the provided request
func cloner(req *http.Request) (func() (*http.Request, error), error) {
getbody := req.GetBody
// if there's no GetBody function set we have to copy the body
// into memory to use for future clones
if getbody == nil {
if req.Body == nil || req.Body == http.NoBody {
getbody = func() (io.ReadCloser, error) {
return http.NoBody, nil
}
} else {
body, err := io.ReadAll(req.Body)
if err != nil {
return nil, err
}
if err := req.Body.Close(); err != nil {
return nil, err
}
getbody = func() (io.ReadCloser, error) {
return io.NopCloser(bytes.NewReader(body)), nil
}
}
}
return func() (*http.Request, error) {
clone := req.Clone(req.Context())
body, err := getbody()
if err != nil {
return nil, err
}
clone.Body = body
clone.GetBody = getbody
return clone, nil
}, nil
}
+7 -3
View File
@@ -900,6 +900,10 @@ github.com/huandu/xstrings
# github.com/huaweicloud/huaweicloud-sdk-go v1.0.26
## explicit
github.com/huaweicloud/huaweicloud-sdk-go/auth/aksk
# github.com/icholy/digest v1.1.0
## explicit; go 1.20
github.com/icholy/digest
github.com/icholy/digest/internal/param
# github.com/imdario/mergo v0.3.6
## explicit
github.com/imdario/mergo
@@ -1832,8 +1836,6 @@ gopkg.in/yaml.v2
# gopkg.in/yaml.v3 v3.0.1
## explicit
gopkg.in/yaml.v3
# gotest.tools/v3 v3.5.1
## explicit; go 1.17
# k8s.io/api v0.19.3
## explicit; go 1.15
k8s.io/api/admissionregistration/v1
@@ -2012,7 +2014,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.20251022024201-332dfd944a29
# yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20251022062416-16b2a67e2e75
## explicit; go 1.24
yunion.io/x/cloudmux/pkg/apis
yunion.io/x/cloudmux/pkg/apis/billing
@@ -2073,6 +2075,8 @@ yunion.io/x/cloudmux/pkg/multicloud/objectstore/ceph/provider
yunion.io/x/cloudmux/pkg/multicloud/objectstore/provider
yunion.io/x/cloudmux/pkg/multicloud/objectstore/xsky
yunion.io/x/cloudmux/pkg/multicloud/objectstore/xsky/provider
yunion.io/x/cloudmux/pkg/multicloud/oceanbase
yunion.io/x/cloudmux/pkg/multicloud/oceanbase/provider
yunion.io/x/cloudmux/pkg/multicloud/openstack
yunion.io/x/cloudmux/pkg/multicloud/openstack/oscli
yunion.io/x/cloudmux/pkg/multicloud/openstack/provider
+1
View File
@@ -51,6 +51,7 @@ const (
CLOUD_PROVIDER_UIS = "UIS"
CLOUD_PROVIDER_CAS = "CAS"
CLOUD_PROVIDER_CNWARE = "CNware"
CLOUD_PROVIDER_OCEANBASE = "OceanBase"
CLOUD_PROVIDER_GENERICS3 = "S3"
CLOUD_PROVIDER_CEPH = "Ceph"
+1
View File
@@ -37,6 +37,7 @@ import (
_ "yunion.io/x/cloudmux/pkg/multicloud/objectstore/ceph/provider"
_ "yunion.io/x/cloudmux/pkg/multicloud/objectstore/provider"
_ "yunion.io/x/cloudmux/pkg/multicloud/objectstore/xsky/provider"
_ "yunion.io/x/cloudmux/pkg/multicloud/oceanbase/provider"
_ "yunion.io/x/cloudmux/pkg/multicloud/openstack/provider"
_ "yunion.io/x/cloudmux/pkg/multicloud/oracle/provider"
_ "yunion.io/x/cloudmux/pkg/multicloud/proxmox/provider" // private clouds
+370
View File
@@ -0,0 +1,370 @@
// 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 oceanbase
import (
"context"
"fmt"
"net/url"
"strings"
"time"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
billing_api "yunion.io/x/cloudmux/pkg/apis/billing"
api "yunion.io/x/cloudmux/pkg/apis/compute"
"yunion.io/x/cloudmux/pkg/cloudprovider"
"yunion.io/x/cloudmux/pkg/multicloud"
)
type SDBInstance struct {
multicloud.SDBInstanceBase
region *SRegion
AvailableZones []string
DeployMode string
DeployType string
NodeNum int
CreateTime time.Time
Status string
Vpc string
StopTime time.Time
PayType string
CloudProvider string
CPU int
Mem int
CPUArchitecture string
DiskType string
Resource struct {
Cpu struct {
TotalCpu float64
UsedCpu float64
}
Memory struct {
TotalMemory float64
UsedMemory float64
}
Disk struct {
TotalDataSize float64
TotalDiskSize float64
UsedDiskSize float64
}
}
InstanceName string
InstanceId string
InstanceType string
InstanceClass string
SaleChannel string
Series string
StandbyInstanceIds []string
State string
StorageArchitecture string
TagList []string
UsedDiskSize float64
Version string
VpcId string
DiskSize int
Region string
}
func (region *SRegion) GetDBInstances() ([]SDBInstance, error) {
params := url.Values{}
params.Set("pageSize", "10")
pageNum := 1
ret := []SDBInstance{}
for {
params.Set("pageNumber", fmt.Sprintf("%d", pageNum))
resp, err := region.list("/api/v2/instances", params)
if err != nil {
return nil, err
}
part := struct {
Data struct {
DataList []SDBInstance
Total int
}
}{}
err = resp.Unmarshal(&part)
if err != nil {
return nil, err
}
ret = append(ret, part.Data.DataList...)
if len(ret) >= part.Data.Total || len(part.Data.DataList) == 0 {
break
}
pageNum++
}
return ret, nil
}
func (region *SRegion) GetDBInstance(id string) (*SDBInstance, error) {
resp, err := region.list(fmt.Sprintf("/api/v2/instances/%s", id), url.Values{})
if err != nil {
return nil, err
}
dbinstance := &SDBInstance{region: region}
err = resp.Unmarshal(&dbinstance, "data")
if err != nil {
return nil, err
}
return dbinstance, nil
}
func (region *SRegion) DeleteDBInstance(id string) error {
_, err := region.delete("/api/v2/instances", map[string]interface{}{"instanceId": id})
return err
}
func (region *SRegion) StartDBInstance(id string) error {
_, err := region.put(fmt.Sprintf("/api/v2/instances/%s/startCluster", id), nil)
return err
}
func (region *SRegion) StopDBInstance(id string) error {
_, err := region.put(fmt.Sprintf("/api/v2/instances/%s/stopCluster", id), nil)
return err
}
func (rds *SDBInstance) GetGlobalId() string {
return rds.InstanceId
}
func (rds *SDBInstance) GetName() string {
return rds.InstanceName
}
func (rds *SDBInstance) GetStatus() string {
switch rds.Status {
case "ONLINE":
return api.DBINSTANCE_RUNNING
case "PENDING_STOP", "STOPPED", "PENDING_START":
return api.DBINSTANCE_REBOOTING
case "PENDING_DELETE":
return api.DBINSTANCE_DELETING
case "PENDING_CREATE":
return api.DBINSTANCE_DEPLOYING
default:
return api.DBINSTANCE_UNKNOWN
}
}
func (rds *SDBInstance) GetCreatedAt() time.Time {
return rds.CreateTime
}
func (rds *SDBInstance) GetExpiredAt() time.Time {
return time.Time{}
}
func (rds *SDBInstance) GetBillingType() string {
if rds.PayType == "POSTPAY" {
return billing_api.BILLING_TYPE_POSTPAID
}
return billing_api.BILLING_TYPE_PREPAID
}
func (rds *SDBInstance) GetProjectId() string {
return ""
}
// ICloudResource 接口方法
func (rds *SDBInstance) GetId() string {
return rds.InstanceId
}
func (rds *SDBInstance) GetDescription() string {
return ""
}
func (rds *SDBInstance) Refresh() error {
instance, err := rds.region.GetDBInstance(rds.InstanceId)
if err != nil {
return err
}
return jsonutils.Update(rds, instance)
}
func (rds *SDBInstance) GetSysTags() map[string]string {
return map[string]string{}
}
func (rds *SDBInstance) GetTags() (map[string]string, error) {
return map[string]string{}, nil
}
func (rds *SDBInstance) SetTags(tags map[string]string, replace bool) error {
return cloudprovider.ErrNotSupported
}
// 资源相关方法
func (rds *SDBInstance) GetVcpuCount() int {
return rds.CPU
}
func (rds *SDBInstance) GetVmemSizeMB() int {
return rds.Mem * 1024
}
func (rds *SDBInstance) GetDiskSizeGB() int {
return rds.DiskSize
}
func (rds *SDBInstance) GetDiskSizeUsedMB() int {
return int(rds.UsedDiskSize * 1024)
}
func (rds *SDBInstance) GetInstanceType() string {
return rds.InstanceClass
}
// 网络相关方法
func (rds *SDBInstance) GetPort() int {
return 2881 // OceanBase 默认端口
}
func (rds *SDBInstance) GetConnectionStr() string {
return ""
}
func (rds *SDBInstance) GetInternalConnectionStr() string {
return ""
}
func (rds *SDBInstance) GetIVpcId() string {
return rds.VpcId
}
// 可用区相关方法
func (rds *SDBInstance) GetZone1Id() string {
if len(rds.AvailableZones) > 0 {
return fmt.Sprintf("%s-%s", rds.Region, rds.AvailableZones[0])
}
return ""
}
func (rds *SDBInstance) GetZone2Id() string {
if len(rds.AvailableZones) > 1 {
return fmt.Sprintf("%s-%s", rds.Region, rds.AvailableZones[1])
}
return ""
}
func (rds *SDBInstance) GetZone3Id() string {
if len(rds.AvailableZones) > 2 {
return fmt.Sprintf("%s-%s", rds.Region, rds.AvailableZones[2])
}
return ""
}
// 引擎相关方法
func (rds *SDBInstance) GetEngine() string {
return "OceanBase"
}
func (rds *SDBInstance) GetEngineVersion() string {
return rds.Version
}
func (rds *SDBInstance) GetCategory() string {
return strings.ToLower(rds.InstanceType)
}
func (rds *SDBInstance) GetStorageType() string {
return rds.DiskType
}
func (rds *SDBInstance) GetMaintainTime() string {
return ""
}
func (rds *SDBInstance) GetIops() int {
return 0
}
// 操作相关方法
func (rds *SDBInstance) Reboot() error {
return errors.Wrapf(cloudprovider.ErrNotImplemented, "Reboot")
}
func (rds *SDBInstance) Delete() error {
return rds.region.DeleteDBInstance(rds.InstanceId)
}
func (rds *SDBInstance) ChangeConfig(ctx context.Context, config *cloudprovider.SManagedDBInstanceChangeConfig) error {
return errors.Wrapf(cloudprovider.ErrNotImplemented, "ChangeConfig")
}
func (rds *SDBInstance) Update(ctx context.Context, input cloudprovider.SDBInstanceUpdateOptions) error {
return errors.Wrapf(cloudprovider.ErrNotImplemented, "Update")
}
// 管理相关方法
func (rds *SDBInstance) GetMasterInstanceId() string {
return ""
}
func (rds *SDBInstance) GetSecurityGroupIds() ([]string, error) {
return []string{}, nil
}
func (rds *SDBInstance) SetSecurityGroups(ids []string) error {
return errors.Wrapf(cloudprovider.ErrNotImplemented, "SetSecurityGroups")
}
func (rds *SDBInstance) GetDBNetworks() ([]cloudprovider.SDBInstanceNetwork, error) {
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "GetDBNetworks")
}
func (rds *SDBInstance) GetIDBInstanceParameters() ([]cloudprovider.ICloudDBInstanceParameter, error) {
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "GetIDBInstanceParameters")
}
func (rds *SDBInstance) GetIDBInstanceDatabases() ([]cloudprovider.ICloudDBInstanceDatabase, error) {
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "GetIDBInstanceDatabases")
}
func (rds *SDBInstance) GetIDBInstanceAccounts() ([]cloudprovider.ICloudDBInstanceAccount, error) {
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "GetIDBInstanceAccounts")
}
func (rds *SDBInstance) GetIDBInstanceBackups() ([]cloudprovider.ICloudDBInstanceBackup, error) {
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "GetIDBInstanceBackups")
}
func (rds *SDBInstance) OpenPublicConnection() error {
return errors.Wrapf(cloudprovider.ErrNotImplemented, "OpenPublicConnection")
}
func (rds *SDBInstance) ClosePublicConnection() error {
return errors.Wrapf(cloudprovider.ErrNotImplemented, "ClosePublicConnection")
}
func (rds *SDBInstance) CreateDatabase(conf *cloudprovider.SDBInstanceDatabaseCreateConfig) error {
return errors.Wrapf(cloudprovider.ErrNotImplemented, "CreateDatabase")
}
func (rds *SDBInstance) CreateAccount(conf *cloudprovider.SDBInstanceAccountCreateConfig) error {
return errors.Wrapf(cloudprovider.ErrNotImplemented, "CreateAccount")
}
func (rds *SDBInstance) CreateIBackup(conf *cloudprovider.SDBInstanceBackupCreateConfig) (string, error) {
return "", errors.Wrapf(cloudprovider.ErrNotImplemented, "CreateIBackup")
}
func (rds *SDBInstance) RecoveryFromBackup(conf *cloudprovider.SDBInstanceRecoveryConfig) error {
return errors.Wrapf(cloudprovider.ErrNotImplemented, "RecoveryFromBackup")
}
+200
View File
@@ -0,0 +1,200 @@
// 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 oceanbase
import (
"context"
"crypto/tls"
"fmt"
"net/http"
"net/url"
"strings"
"sync"
"time"
"github.com/icholy/digest"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/gotypes"
"yunion.io/x/pkg/util/httputils"
api "yunion.io/x/cloudmux/pkg/apis/compute"
"yunion.io/x/cloudmux/pkg/cloudprovider"
)
const (
OB_DEFAULT_REGION_NAME = "OceanBase Cloud"
)
type OceanBaseClientConfig struct {
cpcfg cloudprovider.ProviderConfig
accessKeyId string
accessKeySecret string
debug bool
}
type SOceanBaseClient struct {
*OceanBaseClientConfig
client *http.Client
lock sync.Mutex
ctx context.Context
}
func NewOceanBaseClientConfig(accessKeyId, accessKeySecret string) *OceanBaseClientConfig {
cfg := &OceanBaseClientConfig{
accessKeyId: accessKeyId,
accessKeySecret: accessKeySecret,
}
return cfg
}
func (cfg *OceanBaseClientConfig) Debug(debug bool) *OceanBaseClientConfig {
cfg.debug = debug
return cfg
}
func (cfg *OceanBaseClientConfig) CloudproviderConfig(cpcfg cloudprovider.ProviderConfig) *OceanBaseClientConfig {
cfg.cpcfg = cpcfg
return cfg
}
func NewOceanBaseClient(cfg *OceanBaseClientConfig) (*SOceanBaseClient, error) {
client := &SOceanBaseClient{
OceanBaseClientConfig: cfg,
ctx: context.Background(),
}
client.ctx = context.WithValue(client.ctx, "time", time.Now())
_, err := client.list("/api/v2/instances", nil)
return client, err
}
func (cli *SOceanBaseClient) GetRegion() *SRegion {
return &SRegion{
client: cli,
}
}
func (cli *SOceanBaseClient) getDefaultClient() *http.Client {
cli.lock.Lock()
defer cli.lock.Unlock()
if !gotypes.IsNil(cli.client) {
return cli.client
}
cli.client = httputils.GetAdaptiveTimeoutClient()
httputils.SetClientProxyFunc(cli.client, cli.cpcfg.ProxyFunc)
ts, _ := cli.client.Transport.(*http.Transport)
ts.TLSClientConfig = &tls.Config{InsecureSkipVerify: true}
t := &digest.Transport{
Username: cli.accessKeyId,
Password: cli.accessKeySecret,
Transport: cloudprovider.GetCheckTransport(ts, func(req *http.Request) (func(resp *http.Response) error, error) {
if cli.cpcfg.ReadOnly {
if req.Method == "GET" {
return nil, nil
}
return nil, errors.Wrapf(cloudprovider.ErrAccountReadOnly, "%s %s", req.Method, req.URL.Path)
}
return nil, nil
}),
}
cli.client.Transport = t
return cli.client
}
type sObError struct {
StatusCode int `json:"statusCode"`
method httputils.THttpMethod
url string
body jsonutils.JSONObject
}
func (e *sObError) Error() string {
return jsonutils.Marshal(e).String()
}
func (e *sObError) ParseErrorFromJsonResponse(statusCode int, status string, body jsonutils.JSONObject) error {
if body != nil {
body.Unmarshal(e)
}
e.StatusCode = statusCode
log.Infof("%s %s body: %s error: %v", e.method, e.url, e.body, e.Error())
if e.StatusCode == 404 {
return errors.Wrapf(cloudprovider.ErrNotFound, "%s", e.Error())
}
return e
}
func (cli *SOceanBaseClient) Do(req *http.Request) (*http.Response, error) {
client := cli.getDefaultClient()
return client.Do(req)
}
func (cli *SOceanBaseClient) list(resource string, params url.Values) (jsonutils.JSONObject, error) {
return cli.request(httputils.GET, resource, params, nil)
}
func (cli *SOceanBaseClient) delete(resource string, body map[string]interface{}) (jsonutils.JSONObject, error) {
return cli.request(httputils.DELETE, resource, nil, body)
}
func (cli *SOceanBaseClient) put(resource string, body map[string]interface{}) (jsonutils.JSONObject, error) {
return cli.request(httputils.PUT, resource, nil, body)
}
func (cli *SOceanBaseClient) request(method httputils.THttpMethod, resource string, params url.Values, body map[string]interface{}) (jsonutils.JSONObject, error) {
if body == nil {
body = map[string]interface{}{}
}
uri := fmt.Sprintf("https://api-cloud-cn.oceanbase.com/%s", strings.TrimPrefix(resource, "/"))
if len(params) > 0 {
uri = fmt.Sprintf("%s?%s", uri, params.Encode())
}
req := httputils.NewJsonRequest(method, uri, body)
bErr := &sObError{method: method, url: uri, body: jsonutils.Marshal(body)}
client := httputils.NewJsonClient(cli)
_, resp, err := client.Send(cli.ctx, req, bErr, cli.debug)
if err != nil {
return nil, err
}
if !jsonutils.QueryBoolean(resp, "success", true) {
return nil, fmt.Errorf("request failed: %s", resp.String())
}
return resp, nil
}
func (cli *SOceanBaseClient) GetSubAccounts() ([]cloudprovider.SSubAccount, error) {
subAccount := cloudprovider.SSubAccount{}
subAccount.Id = cli.GetAccountId()
subAccount.Name = cli.cpcfg.Name
subAccount.Account = cli.accessKeyId
subAccount.HealthStatus = api.CLOUD_PROVIDER_HEALTH_NORMAL
return []cloudprovider.SSubAccount{subAccount}, nil
}
func (cli *SOceanBaseClient) GetAccountId() string {
return ""
}
func (cli *SOceanBaseClient) GetCapabilities() []string {
caps := []string{
cloudprovider.CLOUD_CAPABILITY_RDS,
}
return caps
}
@@ -0,0 +1,179 @@
// 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 provider
import (
"context"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
api "yunion.io/x/cloudmux/pkg/apis/compute"
"yunion.io/x/cloudmux/pkg/cloudprovider"
"yunion.io/x/cloudmux/pkg/multicloud/oceanbase"
)
type SOceanBaseProviderFactory struct {
cloudprovider.SPublicCloudBaseProviderFactory
}
func (self *SOceanBaseProviderFactory) GetId() string {
return api.CLOUD_PROVIDER_OCEANBASE
}
func (self *SOceanBaseProviderFactory) GetName() string {
return api.CLOUD_PROVIDER_OCEANBASE
}
func (self *SOceanBaseProviderFactory) ValidateCreateCloudaccountData(ctx context.Context, input cloudprovider.SCloudaccountCredential) (cloudprovider.SCloudaccount, error) {
output := cloudprovider.SCloudaccount{}
if len(input.AccessKeyId) == 0 {
return output, errors.Wrap(cloudprovider.ErrMissingParameter, "access_key_id")
}
if len(input.AccessKeySecret) == 0 {
return output, errors.Wrap(cloudprovider.ErrMissingParameter, "access_key_secret")
}
output.AccessUrl = input.Environment
output.Account = input.AccessKeyId
output.Secret = input.AccessKeySecret
return output, nil
}
func (self *SOceanBaseProviderFactory) ValidateUpdateCloudaccountCredential(ctx context.Context, input cloudprovider.SCloudaccountCredential, cloudaccount string) (cloudprovider.SCloudaccount, error) {
output := cloudprovider.SCloudaccount{}
if len(input.AccessKeyId) == 0 {
return output, errors.Wrap(cloudprovider.ErrMissingParameter, "access_key_id")
}
if len(input.AccessKeySecret) == 0 {
return output, errors.Wrap(cloudprovider.ErrMissingParameter, "access_key_secret")
}
output = cloudprovider.SCloudaccount{
Account: input.AccessKeyId,
Secret: input.AccessKeySecret,
}
return output, nil
}
func (self *SOceanBaseProviderFactory) GetProvider(cfg cloudprovider.ProviderConfig) (cloudprovider.ICloudProvider, error) {
client, err := oceanbase.NewOceanBaseClient(
oceanbase.NewOceanBaseClientConfig(
cfg.Account,
cfg.Secret,
).CloudproviderConfig(cfg),
)
if err != nil {
return nil, err
}
return &SOceanBaseProvider{
SBaseProvider: cloudprovider.NewBaseProvider(self),
client: client,
}, nil
}
func (self *SOceanBaseProviderFactory) GetClientRC(info cloudprovider.SProviderInfo) (map[string]string, error) {
return map[string]string{
"OCEANBASE_ACCESS_KEY_ID": info.Account,
"OCEANBASE_ACCESS_KEY_SECRET": info.Secret,
}, nil
}
func init() {
factory := SOceanBaseProviderFactory{}
cloudprovider.RegisterFactory(&factory)
}
type SOceanBaseProvider struct {
cloudprovider.SBaseProvider
client *oceanbase.SOceanBaseClient
}
func (self *SOceanBaseProvider) GetSysInfo() (jsonutils.JSONObject, error) {
return jsonutils.NewDict(), nil
}
func (self *SOceanBaseProvider) GetVersion() string {
return ""
}
func (self *SOceanBaseProvider) GetSubAccounts() ([]cloudprovider.SSubAccount, error) {
return self.client.GetSubAccounts()
}
func (self *SOceanBaseProvider) GetAccountId() string {
return self.client.GetAccountId()
}
func (self *SOceanBaseProvider) GetIRegions() ([]cloudprovider.ICloudRegion, error) {
return []cloudprovider.ICloudRegion{
self.client.GetRegion(),
}, nil
}
func (self *SOceanBaseProvider) GetIRegionById(extId string) (cloudprovider.ICloudRegion, error) {
regions, err := self.GetIRegions()
if err != nil {
return nil, err
}
for i := range regions {
if regions[i].GetGlobalId() == extId {
return regions[i], nil
}
}
return nil, cloudprovider.ErrNotFound
}
func (self *SOceanBaseProvider) GetBalance() (*cloudprovider.SBalanceInfo, error) {
return &cloudprovider.SBalanceInfo{
Currency: "CNY",
Status: api.CLOUD_PROVIDER_HEALTH_NORMAL,
}, cloudprovider.ErrNotSupported
}
func (self *SOceanBaseProvider) GetIProjects() ([]cloudprovider.ICloudProject, error) {
return []cloudprovider.ICloudProject{}, nil
}
func (self *SOceanBaseProvider) GetStorageClasses(regionId string) []string {
return []string{}
}
func (self *SOceanBaseProvider) GetBucketCannedAcls(regionId string) []string {
return []string{}
}
func (self *SOceanBaseProvider) GetObjectCannedAcls(regionId string) []string {
return []string{}
}
func (self *SOceanBaseProvider) CreateIProject(name string) (cloudprovider.ICloudProject, error) {
return nil, cloudprovider.ErrNotImplemented
}
func (self *SOceanBaseProvider) GetCapabilities() []string {
return self.client.GetCapabilities()
}
func (self *SOceanBaseProvider) GetIamLoginUrl() string {
return ""
}
func (self *SOceanBaseProvider) GetCloudRegionExternalIdPrefix() string {
return api.CLOUD_PROVIDER_OCEANBASE + "/"
}
func (self *SOceanBaseProvider) GetMetrics(opts *cloudprovider.MetricListOptions) ([]cloudprovider.MetricValues, error) {
return nil, cloudprovider.ErrNotImplemented
}
+112
View File
@@ -0,0 +1,112 @@
// 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 oceanbase
import (
"fmt"
"net/url"
api "yunion.io/x/cloudmux/pkg/apis/compute"
"yunion.io/x/cloudmux/pkg/cloudprovider"
"yunion.io/x/cloudmux/pkg/multicloud"
"yunion.io/x/jsonutils"
)
type SRegion struct {
multicloud.SRegion
multicloud.SNoObjectStorageRegion
multicloud.SNoLbRegion
multicloud.SNoEipRegion
multicloud.SNoSecurityGroupRegion
multicloud.SNoVpcRegion
multicloud.SRegionZoneBase
client *SOceanBaseClient
}
func (region *SRegion) GetId() string {
return fmt.Sprintf("%s/Default", api.CLOUD_PROVIDER_OCEANBASE)
}
func (region *SRegion) GetName() string {
return OB_DEFAULT_REGION_NAME
}
func (region *SRegion) GetGlobalId() string {
return region.GetId()
}
func (region *SRegion) GetProvider() string {
return api.CLOUD_PROVIDER_OCEANBASE
}
func (region *SRegion) GetClient() *SOceanBaseClient {
return region.client
}
func (region *SRegion) GetStatus() string {
return api.CLOUD_REGION_STATUS_INSERVER
}
func (region *SRegion) list(resource string, params url.Values) (jsonutils.JSONObject, error) {
return region.client.list(resource, params)
}
func (region *SRegion) delete(resource string, body map[string]interface{}) (jsonutils.JSONObject, error) {
return region.client.delete(resource, body)
}
func (region *SRegion) put(resource string, body map[string]interface{}) (jsonutils.JSONObject, error) {
return region.client.put(resource, body)
}
func (region *SRegion) GetCloudEnv() string {
return api.CLOUD_PROVIDER_OCEANBASE
}
func (region *SRegion) GetI18n() cloudprovider.SModelI18nTable {
table := cloudprovider.SModelI18nTable{}
return table
}
func (region *SRegion) GetGeographicInfo() cloudprovider.SGeographicInfo {
return cloudprovider.SGeographicInfo{}
}
func (region *SRegion) GetCapabilities() []string {
return region.client.GetCapabilities()
}
func (region *SRegion) GetIDBInstances() ([]cloudprovider.ICloudDBInstance, error) {
dbinstances, err := region.GetDBInstances()
if err != nil {
return nil, err
}
ret := []cloudprovider.ICloudDBInstance{}
for i := range dbinstances {
dbinstances[i].region = region
ret = append(ret, &dbinstances[i])
}
return ret, nil
}
func (region *SRegion) GetIDBInstanceById(id string) (cloudprovider.ICloudDBInstance, error) {
ret, err := region.GetDBInstance(id)
if err != nil {
return nil, err
}
return ret, nil
}
+38
View File
@@ -404,3 +404,41 @@ func (self *SRegion) GetILoadBalancerHealthChecks() ([]cloudprovider.ICloudLoadb
func (self *SRegion) CreateILoadBalancerHealthCheck(healthCheck *cloudprovider.SLoadbalancerHealthCheck) (cloudprovider.ICloudLoadbalancerHealthCheck, error) {
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "CreateILoadBalancerHealthCheck")
}
type SNoEipRegion struct{}
func (region *SNoEipRegion) GetIEips() ([]cloudprovider.ICloudEIP, error) {
return nil, cloudprovider.ErrNotSupported
}
func (region *SNoEipRegion) GetIEipById(id string) (cloudprovider.ICloudEIP, error) {
return nil, cloudprovider.ErrNotSupported
}
func (region *SNoEipRegion) CreateEIP(opts *cloudprovider.SEip) (cloudprovider.ICloudEIP, error) {
return nil, cloudprovider.ErrNotSupported
}
type SNoSecurityGroupRegion struct{}
func (region *SNoSecurityGroupRegion) CreateISecurityGroup(opts *cloudprovider.SecurityGroupCreateInput) (cloudprovider.ICloudSecurityGroup, error) {
return nil, cloudprovider.ErrNotSupported
}
func (region *SNoSecurityGroupRegion) GetISecurityGroupById(id string) (cloudprovider.ICloudSecurityGroup, error) {
return nil, cloudprovider.ErrNotSupported
}
type SNoVpcRegion struct{}
func (region *SNoVpcRegion) CreateIVpc(opts *cloudprovider.VpcCreateOptions) (cloudprovider.ICloudVpc, error) {
return nil, cloudprovider.ErrNotSupported
}
func (region *SNoVpcRegion) GetIVpcs() ([]cloudprovider.ICloudVpc, error) {
return nil, cloudprovider.ErrNotSupported
}
func (region *SNoVpcRegion) GetIVpcById(id string) (cloudprovider.ICloudVpc, error) {
return nil, cloudprovider.ErrNotSupported
}