From dbee3ec258b1c726fa24987c3afe9407a04da084 Mon Sep 17 00:00:00 2001 From: Qu Xuan Date: Thu, 16 Sep 2021 10:36:09 +0800 Subject: [PATCH] fix(hostman): rbd command --- build/docker/Dockerfile.host | 1 - build/docker/multi-arch/Dockerfile.host | 1 - pkg/hostman/guestman/qemu-kvm.go | 2 +- pkg/hostman/storageman/disk_rbd.go | 20 +- pkg/hostman/storageman/imagecache_rbd.go | 2 +- pkg/hostman/storageman/storage_rbd.go | 693 +++++++---------------- pkg/util/cephutils/ceph.go | 479 ++++++++++++++++ pkg/util/cephutils/doc.go | 1 + 8 files changed, 704 insertions(+), 495 deletions(-) create mode 100644 pkg/util/cephutils/ceph.go create mode 100644 pkg/util/cephutils/doc.go diff --git a/build/docker/Dockerfile.host b/build/docker/Dockerfile.host index 11a4dadbf9..26c003bf62 100644 --- a/build/docker/Dockerfile.host +++ b/build/docker/Dockerfile.host @@ -4,6 +4,5 @@ MAINTAINER "Yaoqi Wan wanyaoqi@yunionyun.com" ENV TZ Asia/Shanghai -RUN apk add librados librbd RUN mkdir -p /opt/yunion/bin ADD ./_output/alpine-build/bin/host /opt/yunion/bin/host diff --git a/build/docker/multi-arch/Dockerfile.host b/build/docker/multi-arch/Dockerfile.host index b1d6abe04e..b8a22726a2 100644 --- a/build/docker/multi-arch/Dockerfile.host +++ b/build/docker/multi-arch/Dockerfile.host @@ -13,6 +13,5 @@ MAINTAINER "Yaoqi Wan wanyaoqi@yunionyun.com" ENV TZ Asia/Shanghai -RUN apk add librados librbd RUN mkdir -p /opt/yunion/bin COPY --from=build /root/go/src/yunion.io/x/onecloud/_output/bin/host /opt/yunion/bin/host diff --git a/pkg/hostman/guestman/qemu-kvm.go b/pkg/hostman/guestman/qemu-kvm.go index 9045452234..078ac15ba7 100644 --- a/pkg/hostman/guestman/qemu-kvm.go +++ b/pkg/hostman/guestman/qemu-kvm.go @@ -846,7 +846,7 @@ func (s *SKVMGuestInstance) DeployFs(deployInfo *deployapi.DeployInfo) (jsonutil diskPath, _ := disks[0].GetString("path") disk := storageman.GetManager().GetDiskByPath(diskPath) if disk == nil { - return nil, fmt.Errorf("Cannot find disk index 0") + return nil, fmt.Errorf("Cannot find disk %s index 0", diskPath) } return disk.DeployGuestFs(disk.GetPath(), s.Desc, deployInfo) } else { diff --git a/pkg/hostman/storageman/disk_rbd.go b/pkg/hostman/storageman/disk_rbd.go index 552a5ba899..6678a5ea8b 100644 --- a/pkg/hostman/storageman/disk_rbd.go +++ b/pkg/hostman/storageman/disk_rbd.go @@ -20,8 +20,6 @@ import ( "context" "fmt" - "github.com/ceph/go-ceph/rbd" - "yunion.io/x/jsonutils" "yunion.io/x/log" "yunion.io/x/pkg/errors" @@ -29,6 +27,7 @@ import ( api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/appctx" + "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/hostman/hostutils" ) @@ -48,12 +47,14 @@ func (d *SRBDDisk) GetType() string { func (d *SRBDDisk) Probe() error { storage := d.Storage.(*SRbdStorage) - storageConf := d.Storage.GetStorageConf() - pool, _ := storageConf.GetString("pool") - _, err := storage.withImage(pool, d.Id, func(image *rbd.Image) (interface{}, error) { - return image.GetSize() - }) - return err + exist, err := storage.IsImageExist(d.Id) + if err != nil { + return errors.Wrapf(err, "IsImageExist") + } + if !exist { + return cloudprovider.ErrNotFound + } + return nil } func (d *SRBDDisk) getPath() string { @@ -75,11 +76,12 @@ func (d *SRBDDisk) GetDiskDesc() jsonutils.JSONObject { storage := d.Storage.(*SRbdStorage) storageConf := d.Storage.GetStorageConf() pool, _ := storageConf.GetString("pool") + sizeMb, _ := storage.getImageSizeMb(pool, d.Id) desc := map[string]interface{}{ "disk_id": d.Id, "disk_format": "raw", "disk_path": d.GetPath(), - "disk_size": storage.getImageSizeMb(pool, d.Id), + "disk_size": sizeMb, } return jsonutils.Marshal(desc) } diff --git a/pkg/hostman/storageman/imagecache_rbd.go b/pkg/hostman/storageman/imagecache_rbd.go index 39a5657256..7cd22c9106 100644 --- a/pkg/hostman/storageman/imagecache_rbd.go +++ b/pkg/hostman/storageman/imagecache_rbd.go @@ -113,7 +113,7 @@ func (r *SRbdImageCache) GetDesc() *remotefile.SImageDesc { imageCacheManger := r.Manager.(*SRbdImageCacheManager) storage := imageCacheManger.storage.(*SRbdStorage) - size := storage.getImageSizeMb(imageCacheManger.Pool, r.GetName()) + size, _ := storage.getImageSizeMb(imageCacheManger.Pool, r.GetName()) return &remotefile.SImageDesc{ Size: int64(size), Name: r.imageName, diff --git a/pkg/hostman/storageman/storage_rbd.go b/pkg/hostman/storageman/storage_rbd.go index 48bcc07948..44095119f8 100644 --- a/pkg/hostman/storageman/storage_rbd.go +++ b/pkg/hostman/storageman/storage_rbd.go @@ -21,24 +21,21 @@ import ( "fmt" "os" "strings" - "time" - - "github.com/ceph/go-ceph/rados" - "github.com/ceph/go-ceph/rbd" - "github.com/pkg/errors" "yunion.io/x/jsonutils" "yunion.io/x/log" + "yunion.io/x/pkg/errors" + "yunion.io/x/pkg/gotypes" "yunion.io/x/pkg/utils" api "yunion.io/x/onecloud/pkg/apis/compute" - "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "yunion.io/x/onecloud/pkg/cloudprovider" deployapi "yunion.io/x/onecloud/pkg/hostman/hostdeployer/apis" "yunion.io/x/onecloud/pkg/hostman/hostdeployer/deployclient" "yunion.io/x/onecloud/pkg/hostman/hostutils" "yunion.io/x/onecloud/pkg/hostman/options" "yunion.io/x/onecloud/pkg/mcclient/modules" + "yunion.io/x/onecloud/pkg/util/cephutils" "yunion.io/x/onecloud/pkg/util/procutils" "yunion.io/x/onecloud/pkg/util/qemutils" ) @@ -48,22 +45,24 @@ const ( RBD_ORDER = 22 //为rbd对应到rados中每个对象的大小,默认为4MB ) -var ( - ErrNoSuchImage = errors.New("no such image") - ErrNoSuchSnapshot = errors.New("no such snapshot") -) +type sStorageConf struct { + MonHost string + Key string + Pool string + RadosMonOpTimeout int64 + RadosOsdOpTimeout int64 + ClientMountTimeout int64 +} type SRbdStorage struct { SBaseStorage + sStorageConf } func NewRBDStorage(manager *SStorageManager, path string) *SRbdStorage { var ret = new(SRbdStorage) ret.SBaseStorage = *NewBaseStorage(manager, path) - err := procutils.NewRemoteCommandAsFarAsPossible("mkdir", "-p", "/etc/ceph").Run() - if err != nil { - log.Errorf("Failed to mkdir /etc/ceph: %s", err) - } + ret.sStorageConf = sStorageConf{} return ret } @@ -90,28 +89,21 @@ func (s *SRbdStorage) GetSnapshotPathByIds(diskId, snapshotId string) string { return "" } +func (s *SRbdStorage) GetClient() (*cephutils.CephClient, error) { + return cephutils.NewClient(s.MonHost, s.Key, s.Pool) +} + func (s *SRbdStorage) IsSnapshotExist(diskId, snapshotId string) (bool, error) { - var exist bool - pool, _ := s.StorageConf.GetString("pool") - _, err := s.withImage(pool, diskId, - func(src *rbd.Image) (interface{}, error) { - sps, err := src.GetSnapshotNames() - if err != nil { - return nil, errors.Wrap(err, "get snapshot names") - } - for i := 0; i < len(sps); i++ { - if sps[i].Name == snapshotId { - exist = true - break - } - } - return nil, nil - }, - ) + client, err := s.GetClient() if err != nil { - return false, err + return false, errors.Wrapf(err, "GetClient") } - return exist, nil + defer client.Close() + img, err := client.GetImage(diskId) + if err != nil { + return false, errors.Wrapf(err, "GetImage") + } + return img.IsSnapshotExist(snapshotId) } func (s *SRbdStorage) GetSnapshotDir() string { @@ -132,488 +124,217 @@ func (s *SRbdStorage) GetImgsaveBackupPath() string { //Tip Configuration values containing :, @, or = can be escaped with a leading \ character. func (s *SRbdStorage) getStorageConfString() string { - conf := "" - for _, key := range []string{"mon_host", "key"} { - if value, _ := s.StorageConf.GetString(key); len(value) > 0 { - if key == "mon_host" { - value = strings.Replace(value, ",", `\;`, -1) - } - for _, keyworkd := range []string{":", "@", "="} { - if strings.Index(value, keyworkd) != -1 { - value = strings.Replace(value, keyworkd, fmt.Sprintf(`\%s`, keyworkd), -1) - } - } - conf += fmt.Sprintf(":%s=%s", key, value) + conf := []string{} + conf = append(conf, "mon_host="+strings.ReplaceAll(s.MonHost, ",", `\;`)) + key := s.Key + if len(key) > 0 { + for _, k := range []string{":", "@", "="} { + key = strings.ReplaceAll(key, k, fmt.Sprintf(`\%s`, k)) } + conf = append(conf, "key="+key) } - 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, + for k, timeout := range map[string]int64{ + "rados_mon_op_timeout": s.RadosMonOpTimeout, + "rados_osd_op_timeout": s.RadosOsdOpTimeout, + "client_mount_timeout": s.ClientMountTimeout, } { - if timeout, _ := s.StorageConf.Int(key); timeout > 0 { - conf += fmt.Sprintf(":%s=%d", key, timeout) - } else { - conf += fmt.Sprintf(":%s=%d", key, _timeout) - } + conf = append(conf, fmt.Sprintf("%s=%d", k, timeout)) } - return conf + return ":" + strings.Join(conf, ":") } -func (s *SRbdStorage) getImageSizeMb(pool string, name string) uint64 { - size, err := s.withImage(pool, name, func(image *rbd.Image) (interface{}, error) { - size, err := image.GetSize() - if err != nil { - return nil, err - } - return size / 1024 / 1024, nil - }) +func (s *SRbdStorage) listImages(pool string) ([]string, error) { + client, err := s.GetClient() if err != nil { - log.Errorf("get image error: %v", err) - return 0 + return nil, errors.Wrapf(err, "GetClient") } - return size.(uint64) + client.SetPool(pool) + defer client.Close() + return client.ListImages() +} + +func (s *SRbdStorage) IsImageExist(name string) (bool, error) { + images, err := s.listImages(s.Pool) + if err != nil { + return false, errors.Wrapf(err, "listImages") + } + if utils.IsInStringArray(name, images) { + return true, nil + } + return false, nil +} + +func (s *SRbdStorage) getImageSizeMb(pool string, name string) (uint64, error) { + client, err := s.GetClient() + if err != nil { + return 0, errors.Wrapf(err, "GetClient") + } + defer client.Close() + client.SetPool(pool) + img, err := client.GetImage(name) + if err != nil { + return 0, errors.Wrapf(err, "GetImage") + } + info, err := img.GetInfo() + if err != nil { + return 0, errors.Wrapf(err, "GetInfo") + } + return uint64(info.SizeByte) / 1024 / 1024, nil } func (s *SRbdStorage) resizeImage(pool string, name string, sizeMb uint64) error { - _, err := s.withImage(pool, name, func(image *rbd.Image) (interface{}, error) { - return nil, image.Resize(sizeMb * 1024 * 1024) - }) - return err + client, err := s.GetClient() + if err != nil { + return errors.Wrapf(err, "GetClient") + } + defer client.Close() + client.SetPool(pool) + img, err := client.GetImage(name) + if err != nil { + return errors.Wrapf(err, "GetImage") + } + return img.Resize(int64(sizeMb)) } func (s *SRbdStorage) deleteImage(pool string, name string) error { - _, err := s.withIOContext(pool, func(ioctx *rados.IOContext) (interface{}, error) { - names, err := rbd.GetImageNames(ioctx) - if err != nil { - return nil, err + client, err := s.GetClient() + if err != nil { + return errors.Wrapf(err, "GetClient") + } + defer client.Close() + client.SetPool(pool) + img, err := client.GetImage(name) + if err != nil { + if errors.Cause(err) == cloudprovider.ErrNotFound { + return nil } - if !utils.IsInStringArray(name, names) { - return nil, nil - } - - image := rbd.GetImage(ioctx, name) - 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 - }) - return err -} - -// 比较费时 -func (s *SRbdStorage) copyImage(srcPool string, srcImage string, destPool string, destImage string) error { - _, err := s.withImage(srcPool, srcImage, func(src *rbd.Image) (interface{}, error) { - imageSize, err := src.GetSize() - if err != nil { - return nil, err - } - if err := s.createImage(destPool, destImage, imageSize/1024/1024); err != nil { - log.Errorf("create image dest pool: %s dest image: %s image size: %dMb error: %v", destPool, destImage, imageSize/1024/1024, err) - return nil, err - } - _, err = s.withImage(destPool, destImage, func(dest *rbd.Image) (interface{}, error) { - return nil, src.Copy(*dest) - }) - return nil, err - }) - return err + return errors.Wrapf(err, "GetImage") + } + return img.Delete() } // 速度快 func (s *SRbdStorage) cloneImage(ctx context.Context, srcPool string, srcImage string, destPool string, destImage string) error { - _, err := s.withImage(srcPool, srcImage, func(src *rbd.Image) (interface{}, error) { - err := func() error { - lockman.LockRawObject(ctx, "rbd_image_cache", srcImage) - defer lockman.ReleaseRawObject(ctx, "rbd_image_cache", srcImage) - snapInfos, err := src.GetSnapshotNames() - if err != nil { - return errors.Wrap(err, "image.GetSnapshotNames()") - } - var snapshot *rbd.Snapshot = nil - for _, snap := range snapInfos { - if snap.Name == destImage { - snapshot = src.GetSnapshot(destImage) - break - } - } - if snapshot == nil { - snapshot, err = src.CreateSnapshot(destImage) - if err != nil { - return errors.Wrap(err, "src.CreateSnapshot") - } - } - - isProtect, err := snapshot.IsProtected() - if err != nil { - return errors.Wrap(err, "snapshot.IsProtected") - } - if !isProtect { - if err := snapshot.Protect(); err != nil { - return errors.Wrap(err, "snapshot.Protect") - } - } - return nil - }() - if err != nil { - return nil, errors.Wrap(err, "get snapshot") - } - - return s.withIOContext(destPool, func(ioctx *rados.IOContext) (interface{}, error) { - _, err := src.Clone(destImage, ioctx, destImage, RBD_FEATURE, RBD_ORDER) - if err != nil { - return nil, errors.Wrapf(err, "src.Clone") - } - return nil, nil - }) - }) - return err + client, err := s.GetClient() + if err != nil { + return errors.Wrapf(err, "GetClient") + } + defer client.Close() + client.SetPool(srcPool) + img, err := client.GetImage(srcImage) + if err != nil { + return errors.Wrapf(err, "GetImage") + } + return img.Clone(ctx, destPool, destImage) } func (s *SRbdStorage) cloneFromSnapshot(srcImage, srcPool, srcSnapshot, newImage, pool string) error { - _, err := s.withImage(srcPool, srcImage, func(src *rbd.Image) (interface{}, error) { - snapshot := src.GetSnapshot(srcSnapshot) - isProtect, err := snapshot.IsProtected() - if err != nil { - return nil, errors.Wrap(err, "snapshot is protected") - } - if !isProtect { - if err := snapshot.Protect(); err != nil { - return nil, errors.Wrap(err, "snapshot protect") - } - defer snapshot.Unprotect() - } - return s.withIOContext(pool, func(ioctx *rados.IOContext) (interface{}, error) { - _, err := src.Clone(srcSnapshot, ioctx, newImage, RBD_FEATURE, RBD_ORDER) - if err != nil { - return nil, errors.Wrap(err, "clone from snapshot") - } - return nil, nil - }) - }) - + client, err := s.GetClient() if err != nil { - return errors.Wrap(err, "clone from snapshot") + return errors.Wrapf(err, "GetClient") } - return nil -} - -func (s *SRbdStorage) withImage(pool string, name string, doFunc func(*rbd.Image) (interface{}, error)) (interface{}, error) { - return s.withIOContext(pool, func(ioctx *rados.IOContext) (interface{}, error) { - image := rbd.GetImage(ioctx, name) - if err := image.Open(); err != nil { - log.Errorf("open image %s name error: %v", name, err) - return nil, err - } - defer image.Close() - return doFunc(image) - }) -} - -func (s *SRbdStorage) withIOContext(pool string, doFunc func(*rados.IOContext) (interface{}, error)) (interface{}, error) { - return s.withCluster(func(conn *rados.Conn) (interface{}, error) { - ioctx, err := conn.OpenIOContext(pool) - if err != nil { - return nil, errors.Wrapf(err, "conn.OpenIOContext(%s)", pool) - } - return doFunc(ioctx) - }) -} - -func (s *SRbdStorage) listImages(pool string) ([]string, error) { - images, err := s.withIOContext(pool, func(ioctx *rados.IOContext) (interface{}, error) { - return rbd.GetImageNames(ioctx) - }) + defer client.Close() + client.SetPool(srcPool) + img, err := client.GetImage(srcImage) if err != nil { - return nil, err + return errors.Wrapf(err, "GetImage") } - return images.([]string), nil -} - -func (s *SRbdStorage) withCluster(doFunc func(*rados.Conn) (interface{}, error)) (interface{}, error) { - conn, _ := rados.NewConn() - for _, key := range []string{"mon_host", "key"} { - if value, _ := s.StorageConf.GetString(key); len(value) > 0 { - if err := conn.SetConfigOption(key, value); err != nil { - return nil, err - } - } + snap, err := img.GetSnapshot(srcSnapshot) + if err != nil { + return errors.Wrapf(err, "GetSnapshot") } - 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 - } - if err := conn.SetConfigOption(key, fmt.Sprintf("%d", timeout)); err != nil { - return nil, err - } - } - if err := conn.Connect(); err != nil { - return nil, errors.Wrapf(err, "conn.Connect() %s", s.StorageName) - } - defer conn.Shutdown() - return doFunc(conn) + return snap.Clone(pool, newImage) } func (s *SRbdStorage) createImage(pool string, name string, sizeMb uint64) error { - _, err := s.withIOContext(pool, func(ioctx *rados.IOContext) (interface{}, error) { - image, err := rbd.Create(ioctx, name, sizeMb*1024*1024, RBD_ORDER, RBD_FEATURE) - if err != nil { - return nil, err - } - defer image.Close() - return nil, nil - }) + client, err := s.GetClient() + if err != nil { + return errors.Wrapf(err, "GetClient") + } + defer client.Close() + client.SetPool(pool) + _, err = client.CreateImage(name, int64(sizeMb)) return err } func (s *SRbdStorage) renameImage(pool string, src string, dest string) error { - _, err := s.withImage(pool, src, func(image *rbd.Image) (interface{}, error) { - return nil, image.Rename(dest) - }) - return err + client, err := s.GetClient() + if err != nil { + return errors.Wrapf(err, "GetClient") + } + defer client.Close() + client.SetPool(pool) + img, err := client.GetImage(src) + if err != nil { + return errors.Wrapf(err, "GetImage") + } + return img.Rename(dest) } func (s *SRbdStorage) resetDisk(pool string, diskId string, snapshotId string) error { - _, err := s.withImage(pool, diskId, func(image *rbd.Image) (interface{}, error) { - snap := image.GetSnapshot(snapshotId) - return nil, snap.Rollback() - }) - return err + client, err := s.GetClient() + if err != nil { + return errors.Wrapf(err, "GetClient") + } + client.SetPool(pool) + defer client.Close() + img, err := client.GetImage(diskId) + if err != nil { + return errors.Wrapf(err, "GetImage") + } + snap, err := img.GetSnapshot(snapshotId) + if err != nil { + return errors.Wrapf(err, "GetSnapshot") + } + return snap.Rollback() } func (s *SRbdStorage) createSnapshot(pool string, diskId string, snapshotId string) error { - _, err := s.withImage(pool, diskId, func(image *rbd.Image) (interface{}, error) { - return image.CreateSnapshot(snapshotId) - }) + client, err := s.GetClient() + if err != nil { + return errors.Wrapf(err, "GetClient") + } + client.SetPool(pool) + defer client.Close() + img, err := client.GetImage(diskId) + if err != nil { + return errors.Wrapf(err, "GetImage") + } + _, err = img.CreateSnapshot(snapshotId) return err } func (s *SRbdStorage) deleteSnapshot(pool string, diskId string, snapshotId string) error { - _, err := s.withImage(pool, diskId, func(image *rbd.Image) (interface{}, error) { - snapshots, err := image.GetSnapshotNames() - if err != nil { - return nil, errors.Wrapf(err, "image.GetSnapshotNames()") - } - var snapshot *rbd.Snapshot - for _, snapInfo := range snapshots { - if snapInfo.Name == snapshotId { - snapshot = image.GetSnapshot(snapshotId) - break - } - } - if snapshot == nil { //not found snapshot - return nil, nil - } - isProtect, err := snapshot.IsProtected() - if err != nil { - return nil, errors.Wrap(err, "snapshot is protected") - } - if isProtect { - image.Close() - if err := image.Open(snapshotId); err != nil { - return nil, errors.Wrap(err, "image open snapshot") - } - pools, childImgs, err := image.ListChildren() - if err != nil { - return nil, errors.Wrap(err, "image list children") - } - for i, iPool := range pools { - _, err = s.withIOContext(iPool, func(ioctx *rados.IOContext) (interface{}, error) { - _image := rbd.GetImage(ioctx, childImgs[i]) - err = _image.Open() - if err != nil { - return nil, errors.Wrap(err, "open child image") - } - defer _image.Close() - - log.Infof("start flatten %s/%s@%s => %s/%s", pool, diskId, snapshotId, iPool, childImgs[i]) - err := _image.Flatten() - if err != nil { - return nil, errors.Wrap(err, "child image flatten") - } - return nil, nil - }) - if err != nil { - return nil, errors.Wrapf(err, "flatten child %s/%s", iPool, childImgs[i]) - } - } - 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", snapshotId) - } - } - err = snapshot.Remove() - if err != nil { - return nil, errors.Wrap(err, "snapshot remove") - } - return nil, nil - }) - return err -} - -type SMonCommand struct { - Prefix string - Pool string - Format string -} - -type sCephCapacity struct { - CapacitySizeByte int64 - UsedCapacitySizeByte int64 -} - -type sCephDfResult struct { - Stats struct { - TotalBytes int64 `json:"total_bytes"` - TotalAvailBytes int64 `json:"total_avail_bytes"` - TotalUsedBytes int64 `json:"total_used_bytes"` - TotalUsedRawBytes int64 `json:"total_used_raw_bytes"` - TotalUsedRawRatio float64 `json:"total_used_raw_ratio"` - NumOsds int `json:"num_osds"` - NumPerPoolOsds int `json:"num_per_pool_osds"` - } - Pools []struct { - Name string `json:"name"` - ID int `json:"id"` - Stats struct { - Stored int `json:"stored"` - Objects int `json:"objects"` - KbUsed int `json:"kb_used"` - BytesUsed int64 `json:"bytes_used"` - PercentUsed int `json:"percent_used"` - MaxAvail int64 `json:"max_avail"` - } `json:"stats"` - } -} - -func (s *SRbdStorage) getCapacity() (*sCephCapacity, error) { - ret := &sCephCapacity{} - _, err := s.withCluster(func(conn *rados.Conn) (interface{}, error) { - poolName, _ := s.StorageConf.GetString("pool") - cmd := SMonCommand{Prefix: "df", Format: "json"} - buffer, _, err := conn.MonCommand([]byte(jsonutils.Marshal(cmd).String())) - if err != nil { - return nil, errors.Wrapf(err, "MonCommand") - } - result, err := jsonutils.Parse(buffer) - if err != nil { - return nil, errors.Wrapf(err, string(buffer)) - } - df := &sCephDfResult{} - result.Unmarshal(df) - ret.CapacitySizeByte = df.Stats.TotalAvailBytes - ret.UsedCapacitySizeByte = df.Stats.TotalUsedBytes - for _, pool := range df.Pools { - if pool.Name == poolName { - ret.UsedCapacitySizeByte = pool.Stats.BytesUsed - } - } - return nil, nil - }) + client, err := s.GetClient() if err != nil { - return ret, errors.Wrap(err, "getCapacity") + return errors.Wrapf(err, "GetClient") } - return ret, nil + client.SetPool(pool) + defer client.Close() + img, err := client.GetImage(diskId) + if err != nil { + return errors.Wrapf(err, "GetImage") + } + snap, err := img.GetSnapshot(snapshotId) + if err != nil { + return errors.Wrapf(err, "GetSnapshot") + } + return snap.Delete() } func (s *SRbdStorage) SyncStorageSize() error { content := jsonutils.NewDict() - capacity, err := s.getCapacity() + client, err := s.GetClient() if err != nil { - return errors.Wrapf(err, "getCapacity") + return errors.Wrapf(err, "GetClient") } - content.Set("capacity", jsonutils.NewInt(int64(capacity.CapacitySizeByte/1024/1024))) - content.Set("actual_capacity_used", jsonutils.NewInt(int64(capacity.UsedCapacitySizeByte/1024/1024))) + defer client.Close() + capacity, err := client.GetCapacity() + if err != nil { + return errors.Wrapf(err, "GetCapacity") + } + content.Set("capacity", jsonutils.NewInt(int64(capacity.CapacitySizeKb/1024))) + content.Set("actual_capacity_used", jsonutils.NewInt(int64(capacity.UsedCapacitySizeKb/1024))) _, err = modules.Storages.Put( hostutils.GetComputeSession(context.Background()), s.StorageId, content) @@ -623,14 +344,20 @@ func (s *SRbdStorage) SyncStorageSize() error { func (s *SRbdStorage) SyncStorageInfo() (jsonutils.JSONObject, error) { content := map[string]interface{}{} if len(s.StorageId) > 0 { - capacity, err := s.getCapacity() + client, err := s.GetClient() if err != nil { return modules.Storages.PerformAction(hostutils.GetComputeSession(context.Background()), s.StorageId, "offline", nil) } + defer client.Close() + capacity, err := client.GetCapacity() + if err != nil { + return modules.Storages.PerformAction(hostutils.GetComputeSession(context.Background()), s.StorageId, "offline", nil) + } + content = map[string]interface{}{ "name": s.StorageName, - "capacity": capacity.CapacitySizeByte / 1024 / 1024, - "actual_capacity_used": capacity.UsedCapacitySizeByte / 1024 / 1024, + "capacity": capacity.CapacitySizeKb / 1024, + "actual_capacity_used": capacity.UsedCapacitySizeKb / 1024, "status": api.STORAGE_ONLINE, "zone": s.GetZoneName(), } @@ -646,9 +373,6 @@ func (s *SRbdStorage) GetDiskById(diskId string) (IDisk, error) { if s.Disks[i].GetId() == diskId { err := s.Disks[i].Probe() if err != nil { - if errors.Cause(err) == rbd.RbdErrorNotFound { - return nil, cloudprovider.ErrNotFound - } return nil, errors.Wrapf(err, "disk.Prob") } return s.Disks[i], nil @@ -671,20 +395,12 @@ func (s *SRbdStorage) CreateDisk(diskId string) IDisk { } func (s *SRbdStorage) Accessible() error { - var c = make(chan error) - go func() { - _, err := s.withCluster(func(conn *rados.Conn) (interface{}, error) { - return conn.ListPools() - }) - c <- err - }() - var err error - select { - case err = <-c: - break - case <-time.After(time.Second * 30): - err = ErrStorageTimeout + client, err := s.GetClient() + if err != nil { + return errors.Wrapf(err, "GetClient") } + defer client.Close() + _, err = client.GetCapacity() return err } @@ -814,8 +530,21 @@ func (s *SRbdStorage) CreateDiskFromSnapshot( func (s *SRbdStorage) SetStorageInfo(storageId, storageName string, conf jsonutils.JSONObject) error { s.StorageId = storageId s.StorageName = storageName + if gotypes.IsNil(conf) { + return fmt.Errorf("empty storage conf for storage %s(%s)", storageName, storageId) + } if dconf, ok := conf.(*jsonutils.JSONDict); ok { s.StorageConf = dconf } + conf.Unmarshal(&s.sStorageConf) + if s.RadosMonOpTimeout == 0 { + s.RadosMonOpTimeout = api.RBD_DEFAULT_MON_TIMEOUT + } + if s.RadosOsdOpTimeout == 0 { + s.RadosOsdOpTimeout = api.RBD_DEFAULT_OSD_TIMEOUT + } + if s.ClientMountTimeout == 0 { + s.ClientMountTimeout = api.RBD_DEFAULT_MOUNT_TIMEOUT + } return nil } diff --git a/pkg/util/cephutils/ceph.go b/pkg/util/cephutils/ceph.go new file mode 100644 index 0000000000..99f78c8c9c --- /dev/null +++ b/pkg/util/cephutils/ceph.go @@ -0,0 +1,479 @@ +// 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 cephutils + +import ( + "context" + "fmt" + "io/ioutil" + "os" + "strings" + + "github.com/pkg/errors" + + "yunion.io/x/jsonutils" + "yunion.io/x/pkg/utils" + + "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" + "yunion.io/x/onecloud/pkg/cloudprovider" + "yunion.io/x/onecloud/pkg/util/procutils" +) + +type CephClient struct { + monHost string + key string + pool string + cephConf string + keyConf string +} + +func (self *CephClient) Close() error { + if len(self.keyConf) > 0 { + os.Remove(self.keyConf) + } + return os.Remove(self.cephConf) +} + +func (self *CephClient) SetPool(pool string) { + self.pool = pool +} + +type cephStats struct { + Stats struct { + TotalBytes int64 `json:"total_bytes"` + TotalAvailBytes int64 `json:"total_avail_bytes"` + TotalUsedBytes int64 `json:"total_used_bytes"` + TotalUsedRawBytes int64 `json:"total_used_raw_bytes"` + TotalUsedRawRatio float64 `json:"total_used_raw_ratio"` + NumOsds int `json:"num_osds"` + NumPerPoolOsds int `json:"num_per_pool_osds"` + } `json:"stats"` + StatsByClass struct { + Hdd struct { + TotalBytes int64 `json:"total_bytes"` + TotalAvailBytes int64 `json:"total_avail_bytes"` + TotalUsedBytes int64 `json:"total_used_bytes"` + TotalUsedRawBytes int64 `json:"total_used_raw_bytes"` + TotalUsedRawRatio float64 `json:"total_used_raw_ratio"` + } `json:"hdd"` + } `json:"stats_by_class"` + Pools []struct { + Name string `json:"name"` + ID int `json:"id"` + Stats struct { + Stored int `json:"stored"` + Objects int `json:"objects"` + KbUsed int `json:"kb_used"` + BytesUsed int `json:"bytes_used"` + PercentUsed int `json:"percent_used"` + MaxAvail int64 `json:"max_avail"` + } `json:"stats"` + } `json:"pools"` +} + +type SCapacity struct { + CapacitySizeKb int64 + UsedCapacitySizeKb int64 +} + +func (self *CephClient) output(name string, opts []string) (jsonutils.JSONObject, error) { + opts = append([]string{"--format", "json"}, opts...) + resp, err := procutils.NewRemoteCommandAsFarAsPossible(name, opts...).Output() + if err != nil { + return nil, errors.Wrapf(err, "CreateImage %s", string(resp)) + } + return jsonutils.Parse(resp) +} + +func (self *CephClient) run(name string, opts []string) error { + output, err := procutils.NewRemoteCommandAsFarAsPossible(name, opts...).Output() + if err != nil { + return errors.Wrapf(err, "CreateImage %s", string(output)) + } + return nil +} + +func (self *CephClient) options() []string { + opts := []string{"--conf", self.cephConf} + if len(self.keyConf) > 0 { + opts = append(opts, []string{"--keyring", self.keyConf}...) + } + return opts +} + +func (self *CephClient) CreateImage(name string, sizeMb int64) (*SImage, error) { + opts := self.options() + image := &SImage{name: name, client: self} + opts = append(opts, []string{"create", image.GetName(), "--size", fmt.Sprintf("%dM", sizeMb)}...) + return image, self.run("rbd", opts) +} + +func (self *CephClient) GetCapacity() (*SCapacity, error) { + result := &SCapacity{} + opts := self.options() + opts = append(opts, "df") + resp, err := self.output("ceph", opts) + if err != nil { + return nil, errors.Wrapf(err, "output") + } + stats := cephStats{} + err = resp.Unmarshal(&stats) + if err != nil { + return nil, errors.Wrapf(err, "ret.Unmarshal") + } + result.CapacitySizeKb = stats.Stats.TotalAvailBytes / 1024 + result.UsedCapacitySizeKb = stats.Stats.TotalUsedBytes / 1024 + for _, pool := range stats.Pools { + if pool.Name == self.pool { + result.UsedCapacitySizeKb = int64(pool.Stats.KbUsed) + } + } + return result, nil +} + +func writeFile(pattern string, content string) (string, error) { + file, err := ioutil.TempFile("", pattern) + if err != nil { + return "", errors.Wrapf(err, "TempFile") + } + defer file.Close() + name := file.Name() + _, err = file.Write([]byte(content)) + if err != nil { + return name, errors.Wrapf(err, "write") + } + return name, nil +} + +func NewClient(monHost, key, pool string) (*CephClient, error) { + client := &CephClient{ + monHost: monHost, + key: key, + pool: pool, + } + var err error + monHosts := []string{} + for _, monHost := range strings.Split(client.monHost, ",") { + monHosts = append(monHosts, fmt.Sprintf(`[%s]`, monHost)) + } + conf := fmt.Sprintf(`[global] +mon host = %s +rados mon op timeout = 5 +rados osd_op timeout = 1200 +client mount timeout = 120 +`, strings.Join(monHosts, ",")) + client.cephConf, err = writeFile("ceph.*.conf", conf) + if err != nil { + return nil, errors.Wrapf(err, "write file") + } + if len(client.key) > 0 { + keyring := fmt.Sprintf(`[client.admin] + key = %s +`, client.key) + client.keyConf, err = writeFile("ceph.*.keyring", keyring) + if err != nil { + return nil, errors.Wrapf(err, "write keyring") + } + } + return client, nil +} + +type SImage struct { + name string + client *CephClient +} + +func (self *SImage) GetName() string { + return fmt.Sprintf("%s/%s", self.client.pool, self.name) +} + +func (self *CephClient) ListImages() ([]string, error) { + result := []string{} + opts := self.options() + opts = append(opts, []string{"ls", self.pool}...) + resp, err := self.output("rbd", opts) + if err != nil { + return nil, err + } + err = resp.Unmarshal(&result) + if err != nil { + return nil, errors.Wrapf(err, "ret.Unmarshal") + } + return result, nil +} + +func (self *CephClient) GetImage(name string) (*SImage, error) { + images, err := self.ListImages() + if err != nil { + return nil, errors.Wrapf(err, "ListImages") + } + if !utils.IsInStringArray(name, images) { + return nil, cloudprovider.ErrNotFound + } + return &SImage{name: name, client: self}, nil +} + +type SImageInfo struct { + Name string `json:"name"` + ID string `json:"id"` + SizeByte int64 `json:"size"` + Objects int `json:"objects"` + Order int `json:"order"` + ObjectSize int `json:"object_size"` + SnapshotCount int `json:"snapshot_count"` + BlockNamePrefix string `json:"block_name_prefix"` + Format int `json:"format"` + Features []string `json:"features"` + OpFeatures []interface{} `json:"op_features"` + Flags []interface{} `json:"flags"` + CreateTimestamp string `json:"create_timestamp"` + AccessTimestamp string `json:"access_timestamp"` + ModifyTimestamp string `json:"modify_timestamp"` +} + +func (self *SImage) options() []string { + return self.client.options() +} + +func (self *SImage) GetInfo() (*SImageInfo, error) { + opts := self.options() + opts = append(opts, []string{"info", self.GetName()}...) + resp, err := self.client.output("rbd", opts) + if err != nil { + return nil, err + } + info := &SImageInfo{} + return info, resp.Unmarshal(info) +} + +func (self *SImage) ListSnapshots() ([]SSnapshot, error) { + opts := self.options() + opts = append(opts, []string{"snap", "ls", self.GetName()}...) + resp, err := self.client.output("rbd", opts) + if err != nil { + return nil, err + } + result := []SSnapshot{} + err = resp.Unmarshal(&result) + if err != nil { + return nil, errors.Wrapf(err, "ret.Unmarshal") + } + for i := range result { + result[i].image = self + } + return result, nil +} + +func (self *SImage) GetSnapshot(name string) (*SSnapshot, error) { + snaps, err := self.ListSnapshots() + if err != nil { + return nil, errors.Wrapf(err, "ListSnapshots") + } + for i := range snaps { + if snaps[i].Name == name { + return &snaps[i], nil + } + } + return nil, cloudprovider.ErrNotFound +} + +func (self *SImage) IsSnapshotExist(name string) (bool, error) { + _, err := self.GetSnapshot(name) + if err != nil { + if errors.Cause(err) == cloudprovider.ErrNotFound { + return false, nil + } + return false, errors.Wrapf(err, "GetSnapshot") + } + return true, nil +} + +type SSnapshot struct { + Name string + Id string + Size int64 + Protected bool + Timestamp string + + image *SImage +} + +func (self *SSnapshot) Rollback() error { + if !self.Protected { + return nil + } + opts := self.options() + opts = append(opts, []string{"snap", "rollback", self.GetName()}...) + return self.image.client.run("rbd", opts) +} + +func (self *SSnapshot) options() []string { + return self.image.options() +} + +func (self *SSnapshot) GetName() string { + return fmt.Sprintf("%s@%s", self.image.GetName(), self.Name) +} + +func (self *SSnapshot) Unprotect() error { + if !self.Protected { + return nil + } + opts := self.options() + opts = append(opts, []string{"snap", "unprotect", self.GetName()}...) + return self.image.client.run("rbd", opts) +} + +func (self *SSnapshot) Protect() error { + if self.Protected { + return nil + } + opts := self.options() + opts = append(opts, []string{"snap", "protect", self.GetName()}...) + return self.image.client.run("rbd", opts) +} + +func (self *SSnapshot) Remove() error { + opts := self.options() + opts = append(opts, []string{"snap", "rm", self.GetName()}...) + return self.image.client.run("rbd", opts) +} + +func (self *SSnapshot) Delete() error { + pool := self.image.client.pool + defer self.image.client.SetPool(pool) + if self.Protected { + err := self.Unprotect() + if err != nil { + return errors.Wrapf(err, "Unprotect") + } + chidren, err := self.ListChildren() + if err != nil { + return errors.Wrapf(err, "") + } + + for i := range chidren { + self.image.client.SetPool(chidren[i].Pool) + image, err := self.image.client.GetImage(chidren[i].Image) + if err != nil { + return errors.Wrapf(err, "GetImage(%s/%s)", chidren[i].Pool, chidren[i].Image) + } + err = image.Flatten() + if err != nil { + return errors.Wrapf(err, "Flatten") + } + } + } + return self.Remove() +} + +type SChildren struct { + Pool string + PoolNamespace string + Image string +} + +func (self *SSnapshot) ListChildren() ([]SChildren, error) { + opts := self.options() + opts = append(opts, []string{"children", self.GetName()}...) + resp, err := self.image.client.output("rbd", opts) + if err != nil { + return nil, errors.Wrapf(err, "ListChildren") + } + chidren := []SChildren{} + return chidren, resp.Unmarshal(&chidren) +} + +func (self *SImage) Resize(sizeMb int64) error { + opts := self.options() + opts = append(opts, []string{"resize", self.GetName(), "--size", fmt.Sprintf("%dM", sizeMb)}...) + return self.client.run("rbd", opts) +} + +func (self *SImage) Remove() error { + opts := self.options() + opts = append(opts, []string{"rm", self.GetName()}...) + return self.client.run("rbd", opts) +} + +func (self *SImage) Flatten() error { + opts := self.options() + opts = append(opts, []string{"flatten", self.GetName()}...) + return self.client.run("rbd", opts) +} + +func (self *SImage) Delete() error { + snapshots, err := self.ListSnapshots() + if err != nil { + return errors.Wrapf(err, "ListSnapshots") + } + for i := range snapshots { + err := snapshots[i].Delete() + if err != nil { + return errors.Wrapf(err, "delete snapshot %s", snapshots[i].GetName()) + } + } + return self.Remove() +} + +func (self *SImage) Rename(name string) error { + opts := self.options() + opts = append(opts, []string{"rename", self.GetName(), fmt.Sprintf("%s/%s", self.client.pool, name)}...) + return self.client.run("rbd", opts) +} + +func (self *SImage) CreateSnapshot(name string) (*SSnapshot, error) { + snap := &SSnapshot{Name: name, image: self} + opts := self.options() + opts = append(opts, []string{"snap", "create", snap.GetName()}...) + return snap, self.client.run("rbd", opts) +} + +func (self *SImage) Clone(ctx context.Context, pool, name string) error { + lockman.LockRawObject(ctx, "rbd_image_cache", self.GetName()) + defer lockman.ReleaseRawObject(ctx, "rbd_image_cache", self.GetName()) + + var findOrCreateSnap = func() (*SSnapshot, error) { + snaps, err := self.ListSnapshots() + if err != nil { + return nil, errors.Wrapf(err, "ListSnapshots") + } + for i := range snaps { + if snaps[i].Name == name { + return &snaps[i], nil + } + } + snap, err := self.CreateSnapshot(name) + if err != nil { + return nil, errors.Wrapf(err, "CreateSnapshot") + } + return snap, nil + } + snap, err := findOrCreateSnap() + if err != nil { + return errors.Wrapf(err, "findOrCreateSnap") + } + return snap.Clone(pool, name) +} + +func (self *SSnapshot) Clone(pool, name string) error { + err := self.Protect() + if err != nil { + return errors.Wrapf(err, "Protect") + } + opts := self.options() + opts = append(opts, []string{"clone", self.GetName(), fmt.Sprintf("%s/%s", pool, name)}...) + return self.image.client.run("rbd", opts) +} diff --git a/pkg/util/cephutils/doc.go b/pkg/util/cephutils/doc.go new file mode 100644 index 0000000000..c88b9adddd --- /dev/null +++ b/pkg/util/cephutils/doc.go @@ -0,0 +1 @@ +package cephutils // import "yunion.io/x/onecloud/pkg/util/cephutils"