mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
Merge pull request #1965 from ioito/hotfix/qx-storage-optimized
支持设置ceph精准超时时间
This commit is contained in:
+33
-47
@@ -53,34 +53,23 @@ func init() {
|
||||
})
|
||||
|
||||
type StorageUpdateOptions struct {
|
||||
ID string `help:"ID or Name of storage to update"`
|
||||
Name string `help:"New Name of storage"`
|
||||
Desc string `help:"Description" metavar:"<DESCRIPTION>"`
|
||||
CommitBound float64 `help:"Upper bound of storage overcommit rate"`
|
||||
StorageType string `help:"Storage type" choices:"local|nas|vsan|rbd|baremetal"`
|
||||
MediumType string `help:"Medium type, either ssd or rotate" choices:"ssd|rotate"`
|
||||
Reserved string `help:"Reserved storage space"`
|
||||
ID string `help:"ID or Name of storage to update"`
|
||||
Name string `help:"New Name of storage"`
|
||||
Desc string `help:"Description"`
|
||||
CommitBound float64 `help:"Upper bound of storage overcommit rate"`
|
||||
MediumType string `help:"Medium type, either ssd or rotate" choices:"ssd|rotate"`
|
||||
RbdRadosMonOpTimeout int64 `help:"ceph rados_mon_op_timeout"`
|
||||
RbdRadosOsdOpTimeout int64 `help:"ceph rados_osd_op_timeout"`
|
||||
RbdClientMountTimeout int64 `help:"ceph client_mount_timeout"`
|
||||
RbdKey string `help:"ceph rbd key"`
|
||||
Reserved string `help:"Reserved storage space"`
|
||||
}
|
||||
R(&StorageUpdateOptions{}, "storage-update", "Update a storage", func(s *mcclient.ClientSession, args *StorageUpdateOptions) error {
|
||||
params := jsonutils.NewDict()
|
||||
if len(args.Name) > 0 {
|
||||
params.Add(jsonutils.NewString(args.Name), "name")
|
||||
}
|
||||
if len(args.Desc) > 0 {
|
||||
params.Add(jsonutils.NewString(args.Desc), "description")
|
||||
}
|
||||
if args.CommitBound > 0 {
|
||||
params.Add(jsonutils.NewFloat(args.CommitBound), "cmtbound")
|
||||
}
|
||||
if len(args.StorageType) > 0 {
|
||||
params.Add(jsonutils.NewString(args.StorageType), "storage_type")
|
||||
}
|
||||
if len(args.MediumType) > 0 {
|
||||
params.Add(jsonutils.NewString(args.MediumType), "medium_type")
|
||||
}
|
||||
if len(args.Reserved) > 0 {
|
||||
params.Add(jsonutils.NewString(args.Reserved), "reserved")
|
||||
params, err := options.StructToParams(args)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
result, err := modules.Storages.Update(s, args.ID, params)
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -90,37 +79,34 @@ func init() {
|
||||
})
|
||||
|
||||
type StorageCreateOptions struct {
|
||||
NAME string `help:"Name of the Storage"`
|
||||
ZONE string `help:"Zone id of storage"`
|
||||
Capacity int64 `help:"Capacity of the Storage"`
|
||||
MediumType string `help:"Medium type, either ssd or rotate" choices:"ssd|rotate"`
|
||||
StorageType string `help:"Storage type" choices:"local|nas|vsan|rbd|nfs|gpfs|baremetal"`
|
||||
MonHost string `help:"Ceph mon_host config"`
|
||||
Key string `help:"Ceph key config"`
|
||||
Pool string `help:"Ceph Pool Name"`
|
||||
NfsHost string `help:"NFS host"`
|
||||
NfsSharedDir string `help:"NFS shared dir"`
|
||||
NAME string `help:"Name of the Storage"`
|
||||
ZONE string `help:"Zone id of storage"`
|
||||
Capacity int64 `help:"Capacity of the Storage"`
|
||||
MediumType string `help:"Medium type, either ssd or rotate" choices:"ssd|rotate"`
|
||||
StorageType string `help:"Storage type" choices:"local|nas|vsan|rbd|nfs|gpfs|baremetal"`
|
||||
RbdMonHost string `help:"Ceph mon_host config"`
|
||||
RbdRadosMonOpTimeout int64 `help:"ceph rados_mon_op_timeout"`
|
||||
RbdRadosOsdOpTimeout int64 `help:"ceph rados_osd_op_timeout"`
|
||||
RbdClientMountTimeout int64 `help:"ceph client_mount_timeout"`
|
||||
RbdKey string `help:"Ceph key config"`
|
||||
RbdPool string `help:"Ceph Pool Name"`
|
||||
NfsHost string `help:"NFS host"`
|
||||
NfsSharedDir string `help:"NFS shared dir"`
|
||||
}
|
||||
R(&StorageCreateOptions{}, "storage-create", "Create a Storage", func(s *mcclient.ClientSession, args *StorageCreateOptions) error {
|
||||
params := jsonutils.NewDict()
|
||||
params.Add(jsonutils.NewString(args.NAME), "name")
|
||||
params.Add(jsonutils.NewString(args.ZONE), "zone")
|
||||
params.Add(jsonutils.NewInt(args.Capacity), "capacity")
|
||||
params.Add(jsonutils.NewString(args.StorageType), "storage_type")
|
||||
params.Add(jsonutils.NewString(args.MediumType), "medium_type")
|
||||
params, err := options.StructToParams(args)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if args.StorageType == "rbd" {
|
||||
if args.MonHost == "" || args.Key == "" || args.Pool == "" {
|
||||
if args.RbdMonHost == "" || args.RbdKey == "" || args.RbdPool == "" {
|
||||
return fmt.Errorf("Not enough arguments, missing mon_host、key or pool")
|
||||
}
|
||||
params.Add(jsonutils.NewString(args.MonHost), "rbd_mon_host")
|
||||
params.Add(jsonutils.NewString(args.Key), "rbd_key")
|
||||
params.Add(jsonutils.NewString(args.Pool), "rbd_pool")
|
||||
} else if args.StorageType == "nfs" {
|
||||
if len(args.NfsHost) == 0 || len(args.NfsSharedDir) == 0 {
|
||||
return fmt.Errorf("Storage type nfs missing conf host or shared dir")
|
||||
}
|
||||
params.Add(jsonutils.NewString(args.NfsHost), "nfs_host")
|
||||
params.Add(jsonutils.NewString(args.NfsSharedDir), "nfs_shared_dir")
|
||||
}
|
||||
storage, err := modules.Storages.Create(s, params)
|
||||
if err != nil {
|
||||
|
||||
@@ -81,6 +81,12 @@ const (
|
||||
DISK_TYPE_HYBRID = "hybrid"
|
||||
)
|
||||
|
||||
const (
|
||||
RBD_DEFAULT_MON_TIMEOUT = 5 //5 seconds 连接超时时间
|
||||
RBD_DEFAULT_OSD_TIMEOUT = 20 * 60 //20 minute 操作超时时间
|
||||
RBD_DEFAULT_MOUNT_TIMEOUT = 2 * 60 //CephFS挂载超时时间, 目前未使用
|
||||
)
|
||||
|
||||
var (
|
||||
DISK_TYPES = []string{DISK_TYPE_ROTATE, DISK_TYPE_SSD, DISK_TYPE_HYBRID}
|
||||
STORAGE_LOCAL_TYPES = []string{STORAGE_LOCAL, STORAGE_BAREMETAL, STORAGE_UCLOUD_LOCAL_NORMAL, STORAGE_UCLOUD_LOCAL_SSD, STORAGE_UCLOUD_EXCLUSIVE_LOCAL_DISK}
|
||||
|
||||
@@ -20,6 +20,7 @@ import (
|
||||
"yunion.io/x/jsonutils"
|
||||
"yunion.io/x/log"
|
||||
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/db/taskman"
|
||||
"yunion.io/x/onecloud/pkg/mcclient"
|
||||
)
|
||||
|
||||
@@ -27,6 +28,9 @@ type IStorageDriver interface {
|
||||
GetStorageType() string
|
||||
|
||||
ValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error)
|
||||
ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, data *jsonutils.JSONDict, storage *SStorage) (*jsonutils.JSONDict, error)
|
||||
|
||||
DoStorageUpdateTask(ctx context.Context, userCred mcclient.TokenCredential, storage *SStorage, task taskman.ITask) error
|
||||
|
||||
PostCreate(ctx context.Context, userCred mcclient.TokenCredential, storage *SStorage, data jsonutils.JSONObject)
|
||||
}
|
||||
|
||||
@@ -30,6 +30,7 @@ import (
|
||||
api "yunion.io/x/onecloud/pkg/apis/compute"
|
||||
"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/cloudprovider"
|
||||
"yunion.io/x/onecloud/pkg/compute/options"
|
||||
"yunion.io/x/onecloud/pkg/httperrors"
|
||||
@@ -63,7 +64,7 @@ type SStorage struct {
|
||||
|
||||
Capacity int64 `nullable:"false" list:"admin" update:"admin" create:"admin_required"` // Column(Integer, nullable=False) # capacity of disk in MB
|
||||
Reserved int64 `nullable:"true" default:"0" list:"admin" update:"admin"` // Column(Integer, nullable=True, default=0)
|
||||
StorageType string `width:"32" charset:"ascii" nullable:"false" list:"user" update:"admin" create:"admin_required"` // Column(VARCHAR(32, charset='ascii'), nullable=False)
|
||||
StorageType string `width:"32" charset:"ascii" nullable:"false" list:"user" create:"admin_required"` // Column(VARCHAR(32, charset='ascii'), nullable=False)
|
||||
MediumType string `width:"32" charset:"ascii" nullable:"false" list:"user" update:"admin" create:"admin_required"` // Column(VARCHAR(32, charset='ascii'), nullable=False)
|
||||
Cmtbound float32 `nullable:"true" default:"1" list:"admin" update:"admin"` // Column(Float, nullable=True)
|
||||
StorageConf jsonutils.JSONObject `nullable:"true" get:"admin" update:"admin"` // = Column(JSONEncodedDict, nullable=True)
|
||||
@@ -102,6 +103,14 @@ func (self *SStorage) AllowUpdateItem(ctx context.Context, userCred mcclient.Tok
|
||||
return db.IsAdminAllowUpdate(userCred, self)
|
||||
}
|
||||
|
||||
func (self *SStorage) ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) {
|
||||
driver := GetStorageDriver(self.StorageType)
|
||||
if driver != nil {
|
||||
return driver.ValidateUpdateData(ctx, userCred, data, self)
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
func (self *SStorage) PostUpdate(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) {
|
||||
self.SStandaloneResourceBase.PostUpdate(ctx, userCred, query, data)
|
||||
|
||||
@@ -113,6 +122,19 @@ func (self *SStorage) PostUpdate(ctx context.Context, userCred mcclient.TokenCre
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if update, _ := data.Bool("update_storage_conf"); update {
|
||||
self.StartStorageUpdateTask(ctx, userCred)
|
||||
}
|
||||
}
|
||||
|
||||
func (self *SStorage) StartStorageUpdateTask(ctx context.Context, userCred mcclient.TokenCredential) error {
|
||||
task, err := taskman.TaskManager.NewTask(ctx, "StorageUpdateTask", self, userCred, nil, "", "", nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
task.ScheduleRun(nil)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (self *SStorage) AllowDeleteItem(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool {
|
||||
|
||||
@@ -20,6 +20,7 @@ import (
|
||||
|
||||
"yunion.io/x/jsonutils"
|
||||
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/db/taskman"
|
||||
"yunion.io/x/onecloud/pkg/compute/models"
|
||||
"yunion.io/x/onecloud/pkg/mcclient"
|
||||
)
|
||||
@@ -34,3 +35,12 @@ func (self *SBaseStorageDriver) ValidateCreateData(ctx context.Context, userCred
|
||||
func (self *SBaseStorageDriver) PostCreate(ctx context.Context, userCred mcclient.TokenCredential, storage *models.SStorage, data jsonutils.JSONObject) {
|
||||
|
||||
}
|
||||
|
||||
func (self *SBaseStorageDriver) ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, data *jsonutils.JSONDict, storage *models.SStorage) (*jsonutils.JSONDict, error) {
|
||||
return data, nil
|
||||
}
|
||||
|
||||
func (self *SBaseStorageDriver) DoStorageUpdateTask(ctx context.Context, userCred mcclient.TokenCredential, storage *models.SStorage, task taskman.ITask) error {
|
||||
task.ScheduleRun(nil)
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -24,6 +24,7 @@ import (
|
||||
|
||||
api "yunion.io/x/onecloud/pkg/apis/compute"
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/db"
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/db/taskman"
|
||||
"yunion.io/x/onecloud/pkg/compute/models"
|
||||
"yunion.io/x/onecloud/pkg/httperrors"
|
||||
"yunion.io/x/onecloud/pkg/mcclient"
|
||||
@@ -55,10 +56,16 @@ func (self *SRbdStorageDriver) ValidateCreateData(ctx context.Context, userCred
|
||||
conf.Add(jsonutils.NewString(key), "key")
|
||||
}
|
||||
|
||||
if timeout, _ := data.Int("rbd_timeout"); timeout > 0 {
|
||||
conf.Add(jsonutils.NewInt(timeout), "rados_osd_op_timeout")
|
||||
conf.Add(jsonutils.NewInt(timeout), "rados_mon_op_timeout")
|
||||
conf.Add(jsonutils.NewInt(timeout), "client_mount_timeout")
|
||||
for k, v := range map[string]int64{
|
||||
"rbd_rados_mon_op_timeout": api.RBD_DEFAULT_MON_TIMEOUT,
|
||||
"rbd_rados_osd_op_timeout": api.RBD_DEFAULT_OSD_TIMEOUT,
|
||||
"rbd_client_mount_timeout": api.RBD_DEFAULT_MOUNT_TIMEOUT,
|
||||
} {
|
||||
if timeout, _ := data.Int(k); timeout > 0 {
|
||||
conf.Add(jsonutils.NewInt(timeout), strings.TrimPrefix(k, "rbd_"))
|
||||
} else {
|
||||
conf.Add(jsonutils.NewInt(v), strings.TrimPrefix(k, "rbd_"))
|
||||
}
|
||||
}
|
||||
|
||||
storages := []models.SStorage{}
|
||||
@@ -82,6 +89,37 @@ func (self *SRbdStorageDriver) ValidateCreateData(ctx context.Context, userCred
|
||||
return data, nil
|
||||
}
|
||||
|
||||
func (self *SRbdStorageDriver) ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, data *jsonutils.JSONDict, storage *models.SStorage) (*jsonutils.JSONDict, error) {
|
||||
conf, ok := storage.StorageConf.(*jsonutils.JSONDict)
|
||||
if !ok {
|
||||
conf = jsonutils.NewDict()
|
||||
}
|
||||
data.Set("update_storage_conf", jsonutils.JSONFalse)
|
||||
for _, k := range []string{"rbd_rados_mon_op_timeout", "rbd_rados_osd_op_timeout", "rbd_client_mount_timeout"} {
|
||||
if timeout, _ := data.Int(k); timeout > 0 {
|
||||
conf.Set(strings.TrimPrefix(k, "rbd_"), jsonutils.NewInt(timeout))
|
||||
data.Set("update_storage_conf", jsonutils.JSONTrue)
|
||||
}
|
||||
}
|
||||
|
||||
if key, _ := data.GetString("rbd_key"); len(key) > 0 {
|
||||
conf.Set("key", jsonutils.NewString(key))
|
||||
data.Set("update_storage_conf", jsonutils.JSONTrue)
|
||||
}
|
||||
|
||||
if update, _ := data.Bool("update_storage_conf"); update {
|
||||
_, err := storage.GetModelManager().TableSpec().Update(storage, func() error {
|
||||
storage.StorageConf = conf
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
return nil, httperrors.NewGeneralError(err)
|
||||
}
|
||||
}
|
||||
|
||||
return data, nil
|
||||
}
|
||||
|
||||
func (self *SRbdStorageDriver) PostCreate(ctx context.Context, userCred mcclient.TokenCredential, storage *models.SStorage, data jsonutils.JSONObject) {
|
||||
storages := []models.SStorage{}
|
||||
q := models.StorageManager.Query().Equals("storage_type", api.STORAGE_RBD)
|
||||
@@ -124,3 +162,12 @@ func (self *SRbdStorageDriver) PostCreate(ctx context.Context, userCred mcclient
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (self *SRbdStorageDriver) DoStorageUpdateTask(ctx context.Context, userCred mcclient.TokenCredential, storage *models.SStorage, task taskman.ITask) error {
|
||||
subtask, err := taskman.TaskManager.NewTask(ctx, "RbdStorageUpdateTask", storage, task.GetUserCred(), task.GetParams(), task.GetTaskId(), "", nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
subtask.ScheduleRun(nil)
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -0,0 +1,82 @@
|
||||
// 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"
|
||||
"fmt"
|
||||
|
||||
"yunion.io/x/jsonutils"
|
||||
"yunion.io/x/log"
|
||||
|
||||
"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/mcclient"
|
||||
"yunion.io/x/onecloud/pkg/util/httputils"
|
||||
)
|
||||
|
||||
func init() {
|
||||
taskman.RegisterTask(StorageUpdateTask{})
|
||||
taskman.RegisterTask(RbdStorageUpdateTask{})
|
||||
}
|
||||
|
||||
type StorageUpdateTask struct {
|
||||
taskman.STask
|
||||
}
|
||||
|
||||
func (self *StorageUpdateTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
|
||||
self.SetStage("OnStorageUpdate", nil)
|
||||
storage := obj.(*models.SStorage)
|
||||
dirver := models.GetStorageDriver(storage.StorageType)
|
||||
if dirver != nil {
|
||||
err := dirver.DoStorageUpdateTask(ctx, self.UserCred, storage, self)
|
||||
if err != nil {
|
||||
self.SetStageFailed(ctx, err.Error())
|
||||
}
|
||||
}
|
||||
self.SetStageComplete(ctx, nil)
|
||||
}
|
||||
|
||||
func (self *StorageUpdateTask) OnStorageUpdate(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
|
||||
self.SetStageComplete(ctx, nil)
|
||||
}
|
||||
|
||||
func (self *StorageUpdateTask) OnStorageUpdateFailed(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
|
||||
self.SetStageFailed(ctx, data.String())
|
||||
}
|
||||
|
||||
type RbdStorageUpdateTask struct {
|
||||
taskman.STask
|
||||
}
|
||||
|
||||
func (self *RbdStorageUpdateTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
|
||||
storage := obj.(*models.SStorage)
|
||||
hosts := storage.GetAllAttachingHosts()
|
||||
|
||||
for _, host := range hosts {
|
||||
log.Infof("Updata rbd Storage [%s] on host %s ...", storage.Name, host.Name)
|
||||
url := fmt.Sprintf("%s/storages/update", host.ManagerUri)
|
||||
headers := mcclient.GetTokenHeaders(self.GetUserCred())
|
||||
body := jsonutils.Marshal(map[string]interface{}{
|
||||
"storage_id": storage.Id,
|
||||
"storage_conf": storage.StorageConf,
|
||||
})
|
||||
_, _, err := httputils.JSONRequest(httputils.GetDefaultClient(), ctx, "POST", url, headers, body, false)
|
||||
//这里尽可能的更新所有在线的hoststorage信息,仅打印warning信息
|
||||
log.Warningf("update rbd storage info for host %s(%s) error: %v", host.Name, host.Id, err)
|
||||
}
|
||||
self.SetStageComplete(ctx, nil)
|
||||
}
|
||||
@@ -21,6 +21,7 @@ import (
|
||||
"fmt"
|
||||
"os"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/ceph/go-ceph/rados"
|
||||
"github.com/ceph/go-ceph/rbd"
|
||||
@@ -41,9 +42,8 @@ import (
|
||||
)
|
||||
|
||||
const (
|
||||
RBD_FEATURE = 3
|
||||
RBD_ORDER = 22 //为rbd对应到rados中每个对象的大小,默认为4MB
|
||||
DEFAULT_TIMEOUT = 240 //4 minutes
|
||||
RBD_FEATURE = 3
|
||||
RBD_ORDER = 22 //为rbd对应到rados中每个对象的大小,默认为4MB
|
||||
)
|
||||
|
||||
var (
|
||||
@@ -116,12 +116,16 @@ func (s *SRbdStorage) getStorageConfString() string {
|
||||
conf += fmt.Sprintf(":%s=%s", key, value)
|
||||
}
|
||||
}
|
||||
for _, key := range []string{"rados_osd_op_timeout", "rados_mon_op_timeout", "client_mount_timeout"} {
|
||||
var timeout int64
|
||||
if timeout, _ = s.StorageConf.Int(key); timeout == 0 {
|
||||
timeout = DEFAULT_TIMEOUT
|
||||
for key, _timeout := range map[string]int64{
|
||||
"rados_mon_op_timeout": api.RBD_DEFAULT_MON_TIMEOUT,
|
||||
"rados_osd_op_timeout": api.RBD_DEFAULT_OSD_TIMEOUT,
|
||||
"client_mount_timeout": api.RBD_DEFAULT_MOUNT_TIMEOUT,
|
||||
} {
|
||||
if timeout, _ := s.StorageConf.Int(key); timeout > 0 {
|
||||
conf += fmt.Sprintf(":%s=%d", key, timeout)
|
||||
} else {
|
||||
conf += fmt.Sprintf(":%s=%d", key, _timeout)
|
||||
}
|
||||
conf += fmt.Sprintf(":%s=%d", key, timeout)
|
||||
}
|
||||
return conf
|
||||
}
|
||||
@@ -159,9 +163,83 @@ func (s *SRbdStorage) deleteImage(pool string, name string) error {
|
||||
}
|
||||
|
||||
image := rbd.GetImage(ioctx, name)
|
||||
if err := image.Remove(); err != nil {
|
||||
log.Errorf("remove image %s from pool %s error: %v", name, pool, err)
|
||||
return nil, err
|
||||
err = image.Open()
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "image.Open()")
|
||||
}
|
||||
|
||||
//需要先删除image底下的snap
|
||||
snapInfos, err := image.GetSnapshotNames()
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "image.GetSnapshotNames()")
|
||||
}
|
||||
for _, snapInfo := range snapInfos {
|
||||
image.Close()
|
||||
|
||||
err = image.Open(snapInfo.Name)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "image.Open(%s)", snapInfo.Name)
|
||||
}
|
||||
|
||||
pools, images, err := image.ListChildren()
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "image.ListChildren")
|
||||
}
|
||||
|
||||
for i, _pool := range pools {
|
||||
//需要解除snap底下的image关系
|
||||
_, err = s.withIOContext(_pool, func(ioctx *rados.IOContext) (interface{}, error) {
|
||||
_image := rbd.GetImage(ioctx, images[i])
|
||||
err = _image.Open()
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "_image.Open()")
|
||||
}
|
||||
defer _image.Close()
|
||||
log.Debugf("start flatten %s/%s@%s => %s/%s", pool, name, snapInfo.Name, _pool, images[i])
|
||||
err := _image.Flatten()
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "_image.Flatten")
|
||||
}
|
||||
return nil, nil
|
||||
})
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "flatten child %s/%s", _pool, images[i])
|
||||
}
|
||||
}
|
||||
|
||||
snapshot := image.GetSnapshot(snapInfo.Name)
|
||||
protect, err := snapshot.IsProtected()
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "snapshot.IsProtected() %s", snapInfo.Name)
|
||||
}
|
||||
|
||||
if protect {
|
||||
for i := 0; i < 3; i++ {
|
||||
err = snapshot.Unprotect()
|
||||
if err == nil {
|
||||
break
|
||||
}
|
||||
//Resource busy
|
||||
if strings.Contains(err.Error(), "16") {
|
||||
log.Warningf("snapshot is busy, try unprotect after %d seconds", (i+1)*5)
|
||||
time.Sleep(time.Second * time.Duration(i+1) * 5)
|
||||
continue
|
||||
}
|
||||
return nil, errors.Wrapf(err, "snapshot.Unprotect() %s", snapInfo.Name)
|
||||
}
|
||||
}
|
||||
|
||||
err = snapshot.Remove()
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "snapshot.Remove() %s", snapInfo.Name)
|
||||
}
|
||||
}
|
||||
|
||||
image.Close()
|
||||
|
||||
err = image.Remove()
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "image.Remove() %s/%s", pool, name)
|
||||
}
|
||||
return nil, nil
|
||||
})
|
||||
@@ -192,37 +270,24 @@ func (s *SRbdStorage) cloneImage(srcPool string, srcImage string, destPool strin
|
||||
_, err := s.withImage(srcPool, srcImage, func(src *rbd.Image) (interface{}, error) {
|
||||
snapshot, err := src.CreateSnapshot(destImage)
|
||||
if err != nil {
|
||||
log.Errorf("create snapshot error: %v", err)
|
||||
return nil, err
|
||||
return nil, errors.Wrap(err, "src.CreateSnapshot")
|
||||
}
|
||||
defer snapshot.Remove()
|
||||
|
||||
isProtect, err := snapshot.IsProtected()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if !isProtect {
|
||||
if err := snapshot.Protect(); err != nil {
|
||||
log.Errorf("snapshot protect error: %v", err)
|
||||
return nil, err
|
||||
return nil, errors.Wrap(err, "snapshot.Protect")
|
||||
}
|
||||
}
|
||||
defer snapshot.Unprotect()
|
||||
|
||||
return s.withIOContext(destPool, func(ioctx *rados.IOContext) (interface{}, error) {
|
||||
dest, err := src.Clone(destImage, ioctx, destImage, RBD_FEATURE, RBD_ORDER)
|
||||
_, err := src.Clone(destImage, ioctx, destImage, RBD_FEATURE, RBD_ORDER)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
err = dest.Open()
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "cloneImage.Open")
|
||||
}
|
||||
defer dest.Close()
|
||||
|
||||
err = dest.Flatten()
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "cloneImage.Flatten")
|
||||
return nil, errors.Wrapf(err, "src.Clone")
|
||||
}
|
||||
return nil, nil
|
||||
})
|
||||
@@ -254,8 +319,7 @@ func (s *SRbdStorage) withIOContext(pool string, doFunc func(*rados.IOContext) (
|
||||
return s.withCluster(func(conn *rados.Conn) (interface{}, error) {
|
||||
ioctx, err := conn.OpenIOContext(pool)
|
||||
if err != nil {
|
||||
log.Errorf("get ioctx for pool %s error: %v", pool, err)
|
||||
return nil, err
|
||||
return nil, errors.Wrapf(err, "conn.OpenIOContext(%s)", pool)
|
||||
}
|
||||
return doFunc(ioctx)
|
||||
})
|
||||
@@ -280,7 +344,12 @@ func (s *SRbdStorage) withCluster(doFunc func(*rados.Conn) (interface{}, error))
|
||||
}
|
||||
}
|
||||
}
|
||||
for key, timeout := range map[string]int64{"rados_osd_op_timeout": 3, "rados_mon_op_timeout": 3, "client_mount_timeout": 3} {
|
||||
for key, timeout := range map[string]int64{
|
||||
"rados_osd_op_timeout": api.RBD_DEFAULT_OSD_TIMEOUT,
|
||||
"rados_mon_op_timeout": api.RBD_DEFAULT_MON_TIMEOUT,
|
||||
"client_mount_timeout": api.RBD_DEFAULT_MOUNT_TIMEOUT,
|
||||
} {
|
||||
|
||||
_timeout, _ := s.StorageConf.Int(key)
|
||||
if _timeout > 0 {
|
||||
timeout = _timeout
|
||||
@@ -290,8 +359,7 @@ func (s *SRbdStorage) withCluster(doFunc func(*rados.Conn) (interface{}, error))
|
||||
}
|
||||
}
|
||||
if err := conn.Connect(); err != nil {
|
||||
log.Errorf("connect rbd cluster %s error: %v", s.StorageName, err)
|
||||
return nil, err
|
||||
return nil, errors.Wrapf(err, "conn.Connect() %s", s.StorageName)
|
||||
}
|
||||
defer conn.Shutdown()
|
||||
return doFunc(conn)
|
||||
@@ -372,8 +440,7 @@ func (s *SRbdStorage) getCapacity() (uint64, error) {
|
||||
return uint64(maxBytes) / 1024, nil
|
||||
})
|
||||
if err != nil {
|
||||
log.Errorf("get capacity error: %v", err)
|
||||
return 0, err
|
||||
return 0, errors.Wrap(err, "getCapacity")
|
||||
}
|
||||
sizeKb := _sizeKb.(uint64)
|
||||
return sizeKb / 1024, nil
|
||||
|
||||
Reference in New Issue
Block a user