From 2e3df4e99a21ae7ed737da043b6fb64124bffff9 Mon Sep 17 00:00:00 2001 From: ioito Date: Wed, 24 Jul 2019 20:29:15 +0800 Subject: [PATCH] =?UTF-8?q?=E6=94=AF=E6=8C=81=E8=AE=BE=E7=BD=AEceph?= =?UTF-8?q?=E7=B2=BE=E5=87=86=E8=B6=85=E6=97=B6=E6=97=B6=E9=97=B4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- cmd/climc/shell/storages.go | 80 ++++++------- pkg/apis/compute/storage_const.go | 6 + pkg/compute/models/storagedrivers.go | 4 + pkg/compute/models/storages.go | 24 +++- pkg/compute/storagedrivers/base.go | 10 ++ pkg/compute/storagedrivers/rbd.go | 55 ++++++++- pkg/compute/tasks/storage_update_task.go | 82 +++++++++++++ pkg/hostman/storageman/storage_rbd.go | 139 +++++++++++++++++------ 8 files changed, 312 insertions(+), 88 deletions(-) create mode 100644 pkg/compute/tasks/storage_update_task.go diff --git a/cmd/climc/shell/storages.go b/cmd/climc/shell/storages.go index 2b59c92f5a..eb14008e5d 100644 --- a/cmd/climc/shell/storages.go +++ b/cmd/climc/shell/storages.go @@ -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:""` - 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 { diff --git a/pkg/apis/compute/storage_const.go b/pkg/apis/compute/storage_const.go index 113e89259b..52e9b85466 100644 --- a/pkg/apis/compute/storage_const.go +++ b/pkg/apis/compute/storage_const.go @@ -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} diff --git a/pkg/compute/models/storagedrivers.go b/pkg/compute/models/storagedrivers.go index cb0d986db7..e8ed9da34d 100644 --- a/pkg/compute/models/storagedrivers.go +++ b/pkg/compute/models/storagedrivers.go @@ -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) } diff --git a/pkg/compute/models/storages.go b/pkg/compute/models/storages.go index ffc68eac25..6ba0d0c787 100644 --- a/pkg/compute/models/storages.go +++ b/pkg/compute/models/storages.go @@ -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 { diff --git a/pkg/compute/storagedrivers/base.go b/pkg/compute/storagedrivers/base.go index 36c007e0f6..121cc0d678 100644 --- a/pkg/compute/storagedrivers/base.go +++ b/pkg/compute/storagedrivers/base.go @@ -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 +} diff --git a/pkg/compute/storagedrivers/rbd.go b/pkg/compute/storagedrivers/rbd.go index 50a60054b2..90639a319f 100644 --- a/pkg/compute/storagedrivers/rbd.go +++ b/pkg/compute/storagedrivers/rbd.go @@ -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 +} diff --git a/pkg/compute/tasks/storage_update_task.go b/pkg/compute/tasks/storage_update_task.go new file mode 100644 index 0000000000..d1fe2c714e --- /dev/null +++ b/pkg/compute/tasks/storage_update_task.go @@ -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) +} diff --git a/pkg/hostman/storageman/storage_rbd.go b/pkg/hostman/storageman/storage_rbd.go index e1827b0e66..e6a37622de 100644 --- a/pkg/hostman/storageman/storage_rbd.go +++ b/pkg/hostman/storageman/storage_rbd.go @@ -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