diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index b373eeea27..f9fd73a84a 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -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") diff --git a/pkg/hostman/guestman/qemu-kvm.go b/pkg/hostman/guestman/qemu-kvm.go index a1939296cc..4852bfe730 100644 --- a/pkg/hostman/guestman/qemu-kvm.go +++ b/pkg/hostman/guestman/qemu-kvm.go @@ -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 diff --git a/pkg/hostman/hostutils/hostutils.go b/pkg/hostman/hostutils/hostutils.go index d5214154ea..1c2de8b9b3 100644 --- a/pkg/hostman/hostutils/hostutils.go +++ b/pkg/hostman/hostutils/hostutils.go @@ -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"}) } diff --git a/pkg/hostman/monitor/hmp.go b/pkg/hostman/monitor/hmp.go index 473e6da396..bae240035d 100644 --- a/pkg/hostman/monitor/hmp.go +++ b/pkg/hostman/monitor/hmp.go @@ -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), } diff --git a/pkg/hostman/monitor/hmp_test.go b/pkg/hostman/monitor/hmp_test.go index a355859bc9..565ff4b5d9 100644 --- a/pkg/hostman/monitor/hmp_test.go +++ b/pkg/hostman/monitor/hmp_test.go @@ -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) diff --git a/pkg/hostman/monitor/monitor.go b/pkg/hostman/monitor/monitor.go index e3364f06f9..7dbd44393f 100644 --- a/pkg/hostman/monitor/monitor.go +++ b/pkg/hostman/monitor/monitor.go @@ -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{}, } diff --git a/pkg/hostman/monitor/qmp.go b/pkg/hostman/monitor/qmp.go index 863a476609..e920070872 100644 --- a/pkg/hostman/monitor/qmp.go +++ b/pkg/hostman/monitor/qmp.go @@ -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 diff --git a/pkg/hostman/monitor/qmp_test.go b/pkg/hostman/monitor/qmp_test.go index 66dfa5aa27..8dd9c5c3cd 100644 --- a/pkg/hostman/monitor/qmp_test.go +++ b/pkg/hostman/monitor/qmp_test.go @@ -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)