From 908b9c3ce90c55bb19100ab2f46c11ce7deb88c5 Mon Sep 17 00:00:00 2001 From: wanyaoqi Date: Mon, 11 Nov 2019 11:55:37 +0800 Subject: [PATCH] fix instance-snapshot create --- pkg/baremetal/tasks/bm_register.go | 4 +- pkg/compute/guestdrivers/kvm.go | 8 +++ pkg/compute/models/guest_actions.go | 56 ++++++++++++++----- pkg/compute/models/guests.go | 18 ++++++ pkg/compute/models/instance_snapshots.go | 39 +++++++++++-- pkg/compute/models/snapshots.go | 10 ++++ .../tasks/instance_snapshot_and_clone_task.go | 16 +++--- pkg/hostman/guestman/guesttasks.go | 10 ++-- pkg/hostman/monitor/hmp.go | 2 +- pkg/hostman/monitor/qmp.go | 2 +- pkg/mcclient/modules/mod_servers.go | 2 +- 11 files changed, 132 insertions(+), 35 deletions(-) diff --git a/pkg/baremetal/tasks/bm_register.go b/pkg/baremetal/tasks/bm_register.go index 033b5e34ae..ef9fba2dc8 100644 --- a/pkg/baremetal/tasks/bm_register.go +++ b/pkg/baremetal/tasks/bm_register.go @@ -161,8 +161,8 @@ func (s *sBaremetalRegisterTask) updateIpmiInfo(cli *ssh.Client) { var conf *types.SIPMILanConfig for _, lanChannel := range ipmitool.GetLanChannels(sysInfo) { - conf, err = ipmitool.GetLanConfig(ipmiTool, lanChannel) - if err == nil || len(conf.Mac) == 0 { + conf, _ = ipmitool.GetLanConfig(ipmiTool, lanChannel) + if conf == nil || len(conf.Mac) == 0 { continue } } diff --git a/pkg/compute/guestdrivers/kvm.go b/pkg/compute/guestdrivers/kvm.go index be7bbb6157..d6fab80b53 100644 --- a/pkg/compute/guestdrivers/kvm.go +++ b/pkg/compute/guestdrivers/kvm.go @@ -287,6 +287,14 @@ func (self *SKVMGuestDriver) RequestSyncstatusOnHost(ctx context.Context, guest } func (self *SKVMGuestDriver) OnDeleteGuestFinalCleanup(ctx context.Context, guest *models.SGuest, userCred mcclient.TokenCredential) error { + if ispId := guest.GetMetadata("__base_instance_snapshot_id", userCred); len(ispId) > 0 { + ispM, err := models.InstanceSnapshotManager.FetchById(ispId) + if err == nil { + isp := ispM.(*models.SInstanceSnapshot) + isp.DecRefCount(ctx, userCred) + } + guest.SetMetadata(ctx, "__base_instance_snapshot_id", "", userCred) + } return nil } diff --git a/pkg/compute/models/guest_actions.go b/pkg/compute/models/guest_actions.go index 47752e9102..0e4abb5ae7 100644 --- a/pkg/compute/models/guest_actions.go +++ b/pkg/compute/models/guest_actions.go @@ -54,6 +54,7 @@ import ( "yunion.io/x/onecloud/pkg/util/billing" "yunion.io/x/onecloud/pkg/util/httputils" "yunion.io/x/onecloud/pkg/util/logclient" + "yunion.io/x/onecloud/pkg/util/rand" "yunion.io/x/onecloud/pkg/util/rbacutils" "yunion.io/x/onecloud/pkg/util/seclib2" ) @@ -2495,6 +2496,13 @@ func (self *SGuest) PerformStatus(ctx context.Context, userCred mcclient.TokenCr if len(self.GetMetadata("__mirror_job_status", userCred)) == 0 { self.SetMetadata(ctx, "__mirror_job_status", "ready", userCred) } + } else if ispId := self.GetMetadata("__base_instance_snapshot_id", userCred); len(ispId) > 0 { + ispM, err := InstanceSnapshotManager.FetchById(ispId) + if err == nil { + isp := ispM.(*SInstanceSnapshot) + isp.DecRefCount(ctx, userCred) + } + self.SetMetadata(ctx, "__base_instance_snapshot_id", "", userCred) } if preStatus != self.Status && !self.isNotRunningStatus(preStatus) && self.isNotRunningStatus(self.Status) { @@ -4007,7 +4015,7 @@ func (self *SGuest) PerformInstanceSnapshot( return nil, err } name, _ := data.GetString("name") - instanceSnapshot, err := InstanceSnapshotManager.CreateInstanceSnapshot(ctx, ownerId, self, name) + instanceSnapshot, err := InstanceSnapshotManager.CreateInstanceSnapshot(ctx, ownerId, self, name, false) if err != nil { QuotaManager.CancelPendingUsage( ctx, userCred, rbacutils.ScopeProject, ownerId, self.GetQuotaPlatformID(), pendingUsage, pendingUsage) @@ -4093,23 +4101,47 @@ func (self *SGuest) PerformSnapshotAndClone( if err != nil { return nil, httperrors.NewMissingParameterError("name") } - - pendingUsage, err := self.validateCreateInstanceSnapshot(ctx, userCred, query, data) + count, err := data.Int("count") if err != nil { - return nil, err + count = 1 + } else if count <= 0 { + return nil, httperrors.NewInputParameterError("count must > 0") } lockman.LockClass(ctx, InstanceSnapshotManager, self.ProjectId) defer lockman.ReleaseClass(ctx, InstanceSnapshotManager, self.ProjectId) - instanceSnapshotName, err := db.GenerateName(InstanceSnapshotManager, self.GetOwnerId(), "Snapshot-For-"+newlyGuestName) + // validate create instance snapshot and set snapshot pending usage + snapshotUsage, err := self.validateCreateInstanceSnapshot(ctx, userCred, query, data) if err != nil { - QuotaManager.CancelPendingUsage( - ctx, userCred, rbacutils.ScopeProject, self.GetOwnerId(), + return nil, err + } + // set guest pending usage + pendingUsage, err := self.getGuestUsage(int(count)) + if err != nil { + QuotaManager.CancelPendingUsage(ctx, userCred, rbacutils.ScopeProject, self.GetOwnerId(), + self.GetQuotaPlatformID(), snapshotUsage, snapshotUsage) + return nil, err + } + err = QuotaManager.CheckSetPendingQuota(ctx, userCred, + rbacutils.ScopeProject, self.GetOwnerId(), self.GetQuotaPlatformID(), pendingUsage) + if err != nil { + QuotaManager.CancelPendingUsage(ctx, userCred, rbacutils.ScopeProject, self.GetOwnerId(), + self.GetQuotaPlatformID(), snapshotUsage, snapshotUsage) + return nil, httperrors.NewOutOfQuotaError("Check set pending quota error %s", err) + } + pendingUsage.Snapshot = snapshotUsage.Snapshot + + instanceSnapshotName, err := db.GenerateName(InstanceSnapshotManager, self.GetOwnerId(), + fmt.Sprintf("%s-%s", newlyGuestName, rand.String(8))) + if err != nil { + QuotaManager.CancelPendingUsage(ctx, userCred, rbacutils.ScopeProject, self.GetOwnerId(), self.GetQuotaPlatformID(), pendingUsage, pendingUsage) return nil, httperrors.NewInternalServerError("Generate snapshot name failed %s", err) } - instanceSnapshot, err := InstanceSnapshotManager.CreateInstanceSnapshot(ctx, self.GetOwnerId(), self, instanceSnapshotName) + instanceSnapshot, err := InstanceSnapshotManager.CreateInstanceSnapshot( + ctx, self.GetOwnerId(), self, instanceSnapshotName, + jsonutils.QueryBoolean(data, "auto_delete_instance_snapshot", false)) if err != nil { QuotaManager.CancelPendingUsage( ctx, userCred, rbacutils.ScopeProject, self.GetOwnerId(), @@ -4155,11 +4187,8 @@ func (manager *SGuestManager) CreateGuestFromInstanceSnapshot( if err != nil { return nil, nil, fmt.Errorf("No new guest name provider") } - if err := db.NewNameValidator(manager, isp.GetOwnerId(), guestName, ""); err != nil { - guestName, err = db.GenerateName2(manager, isp.GetOwnerId(), guestName, nil, index) - if err != nil { - return nil, nil, err - } + if guestName, err = db.GenerateName2(manager, isp.GetOwnerId(), guestName, nil, index); err != nil { + return nil, nil, err } guestParams.Set("name", jsonutils.NewString(guestName)) @@ -4172,6 +4201,7 @@ func (manager *SGuestManager) CreateGuestFromInstanceSnapshot( if isp.ServerMetadata != nil { metadata := make(map[string]interface{}, 0) isp.ServerMetadata.Unmarshal(metadata) + metadata["__base_instance_snapshot_id"] = isp.Id guest.SetAllMetadata(ctx, metadata, userCred) } return guest, guestParams, nil diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index eb91141c0e..e9f0911a9b 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -4762,6 +4762,24 @@ func (self *SGuest) GetDiskSnapshotsNotInInstanceSnapshots() ([]SSnapshot, error return snapshots, nil } +func (self *SGuest) getGuestUsage(guestCount int) (*SQuota, error) { + usage := new(SQuota) + usage.Cpu = int(self.VcpuCount) * guestCount + usage.Memory = int(self.VmemSize * guestCount) + diskSize := self.getDiskSize() + if diskSize < 0 { + return nil, httperrors.NewInternalServerError("fetch disk size failed") + } + usage.Storage = self.getDiskSize() * guestCount + netCount, err := self.NetworkCount() + if err != nil && err != sql.ErrNoRows { + return nil, err + } + usage.Port = netCount + usage.Bw = self.getBandwidth(false) + return usage, err +} + func (self *SGuestManager) checkGuestImage(ctx context.Context, input *api.ServerCreateInput) error { // There is no need to check the availability of guest imag if input.Disks is empty if len(input.Disks) == 0 { diff --git a/pkg/compute/models/instance_snapshots.go b/pkg/compute/models/instance_snapshots.go index 2b6c398209..e079a10dd2 100644 --- a/pkg/compute/models/instance_snapshots.go +++ b/pkg/compute/models/instance_snapshots.go @@ -27,6 +27,7 @@ import ( "yunion.io/x/onecloud/pkg/apis/compute" schedapi "yunion.io/x/onecloud/pkg/apis/scheduler" "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "yunion.io/x/onecloud/pkg/cloudcommon/db/quotas" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" "yunion.io/x/onecloud/pkg/httperrors" @@ -51,6 +52,8 @@ type SInstanceSnapshot struct { GuestId string `width:"36" charset:"ascii" nullable:"false" list:"user" create:"required" index:"true"` ServerConfig jsonutils.JSONObject `nullable:"true" list:"user"` ServerMetadata jsonutils.JSONObject `nullable:"true" list:"user"` + AutoDelete bool `default:"false" update:"user" list:"user"` + RefCount int `default:"0" list:"user"` } type SInstanceSnapshotManager struct { @@ -141,7 +144,7 @@ func (self *SInstanceSnapshot) StartCreateInstanceSnapshotTask( } func (manager *SInstanceSnapshotManager) CreateInstanceSnapshot( - ctx context.Context, ownerId mcclient.IIdentityProvider, guest *SGuest, name string, + ctx context.Context, ownerId mcclient.IIdentityProvider, guest *SGuest, name string, autoDelete bool, ) (*SInstanceSnapshot, error) { instanceSnapshot := &SInstanceSnapshot{} instanceSnapshot.SetModelManager(manager, instanceSnapshot) @@ -149,6 +152,7 @@ func (manager *SInstanceSnapshotManager) CreateInstanceSnapshot( instanceSnapshot.ProjectId = ownerId.GetProjectId() instanceSnapshot.DomainId = ownerId.GetProjectDomainId() instanceSnapshot.GuestId = guest.Id + instanceSnapshot.AutoDelete = autoDelete guestSchedInput := guest.ToSchedDesc() for i := 0; i < len(guestSchedInput.Disks); i++ { @@ -206,9 +210,13 @@ func (self *SInstanceSnapshot) ToInstanceCreateInput( serverConfig.Disks[i].SnapshotId = isjs[serverConfig.Disks[i].Index].SnapshotId } sourceInput.Disks = serverConfig.Disks - sourceInput.VmemSize = serverConfig.Memory - sourceInput.VcpuCount = serverConfig.Ncpu - sourceInput.Networks = serverConfig.Networks + if sourceInput.VmemSize == 0 { + sourceInput.VmemSize = serverConfig.Memory + } + if sourceInput.VcpuCount == 0 { + sourceInput.VcpuCount = serverConfig.Ncpu + } + // sourceInput.Networks = serverConfig.Networks return sourceInput, nil } @@ -270,3 +278,26 @@ func (self *SInstanceSnapshot) RealDelete(ctx context.Context, userCred mcclient func (self *SInstanceSnapshot) Delete(ctx context.Context, userCred mcclient.TokenCredential) error { return nil } + +func (self *SInstanceSnapshot) AddRefCount(ctx context.Context) error { + lockman.LockObject(ctx, self) + defer lockman.ReleaseObject(ctx, self) + _, err := db.Update(self, func() error { + self.RefCount += 1 + return nil + }) + return err +} + +func (self *SInstanceSnapshot) DecRefCount(ctx context.Context, userCred mcclient.TokenCredential) error { + lockman.LockObject(ctx, self) + defer lockman.ReleaseObject(ctx, self) + _, err := db.Update(self, func() error { + self.RefCount -= 1 + return nil + }) + if err == nil && self.RefCount == 0 && self.AutoDelete { + self.StartInstanceSnapshotDeleteTask(ctx, userCred, "") + } + return err +} diff --git a/pkg/compute/models/snapshots.go b/pkg/compute/models/snapshots.go index b873b974b7..d6aacd023f 100644 --- a/pkg/compute/models/snapshots.go +++ b/pkg/compute/models/snapshots.go @@ -127,6 +127,16 @@ func (manager *SSnapshotManager) ListItemFilter(ctx context.Context, q *sqlchemy q = q.In("disk_id", sq) } + if isInstanceSnapshot, err := query.Bool("is_instance_snapshot"); err == nil { + insjsq := InstanceSnapshotJointManager.Query().SubQuery() + if !isInstanceSnapshot { + q = q.LeftJoin(insjsq, sqlchemy.Equals(q.Field("id"), insjsq.Field("snapshot_id"))). + Filter(sqlchemy.IsNull(insjsq.Field("snapshot_id"))) + } else { + q = q.Join(insjsq, sqlchemy.Equals(q.Field("id"), insjsq.Field("snapshot_id"))) + } + } + /*if provider, err := query.GetString("provider"); err == nil { cloudproviderTbl := CloudproviderManager.Query().SubQuery() sq := cloudproviderTbl.Query(cloudproviderTbl.Field("id")).Equals("provider", provider) diff --git a/pkg/compute/tasks/instance_snapshot_and_clone_task.go b/pkg/compute/tasks/instance_snapshot_and_clone_task.go index b80c99ef1e..5f9580fda6 100644 --- a/pkg/compute/tasks/instance_snapshot_and_clone_task.go +++ b/pkg/compute/tasks/instance_snapshot_and_clone_task.go @@ -114,16 +114,20 @@ func (self *InstanceSnapshotAndCloneTask) OnCreateInstanceSnapshot( func (self *InstanceSnapshotAndCloneTask) doGuestCreate( ctx context.Context, isp *models.SInstanceSnapshot, params jsonutils.JSONObject, count int) error { - dictParmas := params.(*jsonutils.JSONDict) - var errStr string + + var ( + dictParmas = params.(*jsonutils.JSONDict) + errStr string + ) for i := 0; i < count; i++ { newGuest, input, err := models.GuestManager.CreateGuestFromInstanceSnapshot( - ctx, self.UserCred, dictParmas.DeepCopy().(*jsonutils.JSONDict), isp, i) + ctx, self.UserCred, dictParmas.DeepCopy().(*jsonutils.JSONDict), isp, i+1) if err != nil { log.Errorln(err) errStr += err.Error() + "\n" continue } + isp.AddRefCount(ctx) models.GuestManager.OnCreateComplete(ctx, []db.IModel{newGuest}, self.UserCred, nil, input) } if len(errStr) > 0 { @@ -132,12 +136,6 @@ func (self *InstanceSnapshotAndCloneTask) doGuestCreate( return nil } -func (self *InstanceSnapshotAndCloneTask) OnGuestCreated( - ctx context.Context, isp *models.SInstanceSnapshot, data jsonutils.JSONObject) { - - self.taskComplete(ctx, isp, data) -} - func (self *InstanceSnapshotAndCloneTask) OnCreateInstanceSnapshotFailed( ctx context.Context, isp *models.SInstanceSnapshot, data jsonutils.JSONObject) { self.taskFailed(ctx, isp, data.String()) diff --git a/pkg/hostman/guestman/guesttasks.go b/pkg/hostman/guestman/guesttasks.go index efa87055d1..7e8ed7bcc4 100644 --- a/pkg/hostman/guestman/guesttasks.go +++ b/pkg/hostman/guestman/guesttasks.go @@ -481,7 +481,6 @@ func (s *SGuestResumeTask) onResumeSucc(res string) { } func (s *SGuestResumeTask) onStartRunning() { - // s.removeStatefile() XXX 可能不用了,先注释了 if s.ctx != nil && len(appctx.AppContextTaskId(s.ctx)) > 0 { hostutils.TaskComplete(s.ctx, nil) } @@ -489,17 +488,19 @@ func (s *SGuestResumeTask) onStartRunning() { s.SetVncPassword() } s.OnResumeSyncMetadataInfo() - s.SyncStatus() s.optimizeOom() s.doBlockIoThrottle() timeutils2.AddTimeout(time.Second*5, s.SetCgroup) disksIdx := s.GetNeedMergeBackingFileDiskIndexs() if len(disksIdx) > 0 { - timeutils2.AddTimeout(time.Second*5, func() { s.startStreamDisks(disksIdx) }) + s.startStreamDisks(disksIdx) } else if options.HostOptions.AutoMergeBackingTemplate { + s.SyncStatus() timeutils2.AddTimeout( time.Second*time.Duration(options.HostOptions.AutoMergeDelaySeconds), func() { s.startStreamDisks(nil) }) + } else { + s.SyncStatus() } } @@ -524,6 +525,7 @@ func (s *SGuestResumeTask) startStreamDisks(disksIdx []int) { func (s *SGuestResumeTask) onStreamComplete(disksIdx []int) { if len(disksIdx) == 0 { + // if disks idx length == 0 indicate merge backing template s.SyncStatus() } else { s.streamDisksComplete(s.ctx) @@ -637,7 +639,7 @@ func (s *SGuestStreamDisksTask) startWaitBlockStream(res string) { } func (s *SGuestStreamDisksTask) checkStreamJobs(jobs int) { - if jobs == 0 { + if jobs == 0 && s.c != nil { close(s.c) s.startDoBlockStream() } diff --git a/pkg/hostman/monitor/hmp.go b/pkg/hostman/monitor/hmp.go index 6ad9639b77..c698710b2d 100644 --- a/pkg/hostman/monitor/hmp.go +++ b/pkg/hostman/monitor/hmp.go @@ -395,7 +395,7 @@ func (m *HmpMonitor) DriveMirror(callback StringCallback, drive, target, syncMod func (m *HmpMonitor) BlockStream(drive string, callback StringCallback) { var ( - speed = 30 // MB/s speed limit 31457280 bytes/s + speed = 100 // limit 100 MB/s cmd = fmt.Sprintf("block_stream %s %d", drive, speed) ) m.Query(cmd, callback) diff --git a/pkg/hostman/monitor/qmp.go b/pkg/hostman/monitor/qmp.go index 43da9385bb..363a615dfe 100644 --- a/pkg/hostman/monitor/qmp.go +++ b/pkg/hostman/monitor/qmp.go @@ -713,7 +713,7 @@ func (m *QmpMonitor) DriveMirror(callback StringCallback, drive, target, syncMod func (m *QmpMonitor) BlockStream(drive string, callback StringCallback) { var ( - speed = 30 * 1024 * 1024 // qmp speed default unit is byte + speed = 100 * 1024 * 1024 // limit 100 MB/s cb = func(res *Response) { callback(m.actionResult(res)) } diff --git a/pkg/mcclient/modules/mod_servers.go b/pkg/mcclient/modules/mod_servers.go index 8486e5e6f3..3aff2c1889 100644 --- a/pkg/mcclient/modules/mod_servers.go +++ b/pkg/mcclient/modules/mod_servers.go @@ -48,7 +48,7 @@ func (this *ServerManager) GetLoginInfo(s *mcclient.ClientSession, id string, pa ret.Add(v, "updated") } - loginKey, _ := data.GetString("metadata", "login_key") + loginKey, e := data.GetString("metadata", "login_key") if e != nil { return nil, fmt.Errorf("No login key: %s", e) }