mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-19 02:37:24 +08:00
fix(host): add vm clonse progress and mbps
This commit is contained in:
@@ -127,6 +127,9 @@ type SGuest struct {
|
||||
// 备份机所在宿主机Id
|
||||
BackupHostId string `width:"36" charset:"ascii" nullable:"true" list:"user" get:"user"`
|
||||
|
||||
// 迁移或克隆的速度
|
||||
ProgressMbps float64 `nullable:"false" default:"0" list:"user" create:"optional" update:"user"`
|
||||
|
||||
Vga string `width:"36" charset:"ascii" nullable:"true" list:"user" update:"user" create:"optional"`
|
||||
Vdi string `width:"36" charset:"ascii" nullable:"true" list:"user" update:"user" create:"optional"`
|
||||
Machine string `width:"36" charset:"ascii" nullable:"true" list:"user" update:"user" create:"optional"`
|
||||
@@ -993,6 +996,17 @@ func ValidateMemCpuData(vmemSize, vcpuCount int, hypervisor string) (int, int, e
|
||||
return vmemSize, vcpuCount, nil
|
||||
}
|
||||
|
||||
func (self *SGuest) PreUpdate(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) {
|
||||
// 减少更新日志
|
||||
if data.Contains("progress_mbps") {
|
||||
db.Update(self, func() error {
|
||||
self.ProgressMbps, _ = data.Float("progress_mbps")
|
||||
return nil
|
||||
})
|
||||
}
|
||||
self.SVirtualResourceBase.PreUpdate(ctx, userCred, query, data)
|
||||
}
|
||||
|
||||
func (self *SGuest) ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.ServerUpdateInput) (api.ServerUpdateInput, error) {
|
||||
if len(input.Name) > 0 && len(input.Name) < 2 {
|
||||
return input, httperrors.NewInputParameterError("name is too short")
|
||||
|
||||
@@ -420,6 +420,7 @@ func (s *SKVMGuestInstance) StartMonitorWithImportGuestSocketFile(ctx context.Co
|
||||
timeutils2.AddTimeout(100*time.Millisecond, func() {
|
||||
s.Monitor = monitor.NewQmpMonitor(
|
||||
s.GetName(),
|
||||
s.Id,
|
||||
s.onImportGuestMonitorDisConnect, // on monitor disconnect
|
||||
func(err error) { s.onImportGuestMonitorTimeout(ctx, err) }, // on monitor timeout
|
||||
func() { s.onImportGuestMonitorConnected(ctx) }, // on monitor connected
|
||||
@@ -436,6 +437,7 @@ func (s *SKVMGuestInstance) StartMonitor(ctx context.Context) {
|
||||
var mon monitor.Monitor
|
||||
mon = monitor.NewQmpMonitor(
|
||||
s.GetName(),
|
||||
s.Id,
|
||||
s.onMonitorDisConnect, // on monitor disconnect
|
||||
func(err error) { s.onMonitorTimeout(ctx, err) }, // on monitor timeout
|
||||
func() { s.onMonitorConnected(ctx) }, // on monitor connected
|
||||
@@ -446,6 +448,7 @@ func (s *SKVMGuestInstance) StartMonitor(ctx context.Context) {
|
||||
log.Errorf("Guest %s qmp monitor connect failed %s, try hmp", s.GetName(), err)
|
||||
mon = monitor.NewHmpMonitor(
|
||||
s.GetName(),
|
||||
s.Id,
|
||||
s.onMonitorDisConnect, // on monitor disconnect
|
||||
func(err error) { s.onMonitorTimeout(ctx, err) }, // on monitor timeout
|
||||
func() { s.onMonitorConnected(ctx) }, // on monitor connected
|
||||
|
||||
@@ -148,6 +148,14 @@ func UpdateServerStatus(ctx context.Context, sid, status, reason string) (jsonut
|
||||
return modules.Servers.PerformAction(GetComputeSession(ctx), sid, "status", stats)
|
||||
}
|
||||
|
||||
func UpdateServerProgress(ctx context.Context, sid string, progress, progressMbps float64) (jsonutils.JSONObject, error) {
|
||||
params := map[string]float64{
|
||||
"progress": progress,
|
||||
"progress_mbps": progressMbps,
|
||||
}
|
||||
return modules.Servers.Update(GetComputeSession(ctx), sid, jsonutils.Marshal(params))
|
||||
}
|
||||
|
||||
func ResponseOk(ctx context.Context, w http.ResponseWriter) {
|
||||
Response(ctx, w, map[string]string{"result": "ok"})
|
||||
}
|
||||
|
||||
@@ -36,9 +36,9 @@ type HmpMonitor struct {
|
||||
callbackQueue []StringCallback
|
||||
}
|
||||
|
||||
func NewHmpMonitor(server string, OnMonitorDisConnect, OnMonitorTimeout MonitorErrorFunc, OnMonitorConnected MonitorSuccFunc) *HmpMonitor {
|
||||
func NewHmpMonitor(server, sid string, OnMonitorDisConnect, OnMonitorTimeout MonitorErrorFunc, OnMonitorConnected MonitorSuccFunc) *HmpMonitor {
|
||||
return &HmpMonitor{
|
||||
SBaseMonitor: *NewBaseMonitor(server, OnMonitorConnected, OnMonitorDisConnect, OnMonitorTimeout),
|
||||
SBaseMonitor: *NewBaseMonitor(server, sid, OnMonitorConnected, OnMonitorDisConnect, OnMonitorTimeout),
|
||||
commandQueue: make([]string, 0),
|
||||
callbackQueue: make([]StringCallback, 0),
|
||||
}
|
||||
|
||||
@@ -23,7 +23,7 @@ func TestHmpMonitor_Connect(t *testing.T) {
|
||||
onConnected := func() { t.Logf("Monitor Connected") }
|
||||
onDisConnect := func(error) { t.Logf("Monitor DisConnect") }
|
||||
onTimeout := func(error) { t.Logf("Monitor Timeout") }
|
||||
m := NewHmpMonitor("fake_server", onDisConnect, onTimeout, onConnected)
|
||||
m := NewHmpMonitor("fake_server", "", onDisConnect, onTimeout, onConnected)
|
||||
var host = "127.0.0.1"
|
||||
var port = 55901
|
||||
m.Connect(host, port)
|
||||
|
||||
@@ -47,6 +47,7 @@ type BlockJob struct {
|
||||
start time.Time
|
||||
preOffset int64
|
||||
now time.Time
|
||||
speedMbps float64
|
||||
}
|
||||
|
||||
type blockSizeByte int64
|
||||
@@ -76,6 +77,7 @@ func (self *BlockJob) PreOffset(preOffset int64) {
|
||||
second := time.Now().Sub(self.now).Seconds()
|
||||
if second > 0 {
|
||||
speed := float64(self.Offset-preOffset) / second
|
||||
self.speedMbps = speed / 1024 / 1024
|
||||
avgSpeed := float64(self.Offset) / time.Now().Sub(self.start).Seconds()
|
||||
log.Infof(`[%s / %s] server %s block job for %s speed: %s/s(avg: %s/s)`, blockSizeByte(self.Offset).String(), blockSizeByte(self.Len).String(), self.server, self.Device, blockSizeByte(speed).String(), blockSizeByte(avgSpeed).String())
|
||||
}
|
||||
@@ -144,6 +146,7 @@ type SBaseMonitor struct {
|
||||
OnMonitorTimeout MonitorErrorFunc
|
||||
|
||||
server string
|
||||
sid string
|
||||
|
||||
QemuVersion string
|
||||
connected bool
|
||||
@@ -155,12 +158,13 @@ type SBaseMonitor struct {
|
||||
reading bool
|
||||
}
|
||||
|
||||
func NewBaseMonitor(server string, OnMonitorConnected MonitorSuccFunc, OnMonitorDisConnect, OnMonitorTimeout MonitorErrorFunc) *SBaseMonitor {
|
||||
func NewBaseMonitor(server, sid string, OnMonitorConnected MonitorSuccFunc, OnMonitorDisConnect, OnMonitorTimeout MonitorErrorFunc) *SBaseMonitor {
|
||||
return &SBaseMonitor{
|
||||
OnMonitorConnected: OnMonitorConnected,
|
||||
OnMonitorDisConnect: OnMonitorDisConnect,
|
||||
OnMonitorTimeout: OnMonitorTimeout,
|
||||
server: server,
|
||||
sid: sid,
|
||||
timeout: true,
|
||||
mutex: &sync.Mutex{},
|
||||
}
|
||||
|
||||
@@ -24,10 +24,14 @@ import (
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"golang.org/x/net/context"
|
||||
|
||||
"yunion.io/x/jsonutils"
|
||||
"yunion.io/x/log"
|
||||
"yunion.io/x/pkg/errors"
|
||||
"yunion.io/x/pkg/utils"
|
||||
|
||||
"yunion.io/x/onecloud/pkg/hostman/hostutils"
|
||||
)
|
||||
|
||||
// https://github.com/qemu/qemu/blob/master/docs/interop/qmp-spec.txt
|
||||
@@ -109,10 +113,10 @@ type QmpMonitor struct {
|
||||
jobs map[string]BlockJob
|
||||
}
|
||||
|
||||
func NewQmpMonitor(server string, OnMonitorDisConnect, OnMonitorTimeout MonitorErrorFunc,
|
||||
func NewQmpMonitor(server, sid string, OnMonitorDisConnect, OnMonitorTimeout MonitorErrorFunc,
|
||||
OnMonitorConnected MonitorSuccFunc, qmpEventFunc qmpEventCallback) *QmpMonitor {
|
||||
m := &QmpMonitor{
|
||||
SBaseMonitor: *NewBaseMonitor(server, OnMonitorConnected, OnMonitorDisConnect, OnMonitorTimeout),
|
||||
SBaseMonitor: *NewBaseMonitor(server, sid, OnMonitorConnected, OnMonitorDisConnect, OnMonitorTimeout),
|
||||
qmpEventFunc: qmpEventFunc,
|
||||
commandQueue: make([]*Command, 0),
|
||||
callbackQueue: make([]qmpMonitorCallBack, 0),
|
||||
@@ -671,6 +675,21 @@ func (m *QmpMonitor) blockJobs(res *Response) ([]BlockJob, error) {
|
||||
}
|
||||
jobs := []BlockJob{}
|
||||
ret.Unmarshal(&jobs)
|
||||
defer func() {
|
||||
mbps, progress := 0.0, 0.0
|
||||
totalSize, totalOffset := int64(1), int64(0)
|
||||
for _, job := range m.jobs {
|
||||
mbps += job.speedMbps
|
||||
totalSize += job.Len
|
||||
totalOffset += job.Offset
|
||||
}
|
||||
if len(m.jobs) == 0 && len(jobs) == 0 {
|
||||
progress = 100.0
|
||||
} else {
|
||||
progress = float64(totalOffset) / float64(totalSize) * 100
|
||||
}
|
||||
hostutils.UpdateServerProgress(context.Background(), m.sid, progress, mbps)
|
||||
}()
|
||||
for i := range jobs {
|
||||
job := jobs[i]
|
||||
job.server = m.server
|
||||
|
||||
@@ -25,7 +25,7 @@ func TestQmpMonitor_Connect(t *testing.T) {
|
||||
onConnected := func() { log.Infof("Monitor Connected") }
|
||||
onDisConnect := func(error) { log.Infof("Monitor DisConnect") }
|
||||
onTimeout := func(error) { log.Infof("Monitor Timeout") }
|
||||
m := NewQmpMonitor("", onDisConnect, onTimeout, onConnected, nil)
|
||||
m := NewQmpMonitor("", "", onDisConnect, onTimeout, onConnected, nil)
|
||||
var host = "127.0.0.1"
|
||||
var port = 56101
|
||||
m.Connect(host, port)
|
||||
|
||||
Reference in New Issue
Block a user