Merge pull request #12910 from ioito/automated-cherry-pick-of-#12909-upstream-release-3.8

Automated cherry pick of #12909: fix(host): add speed info for block-job
This commit is contained in:
Zexi Li
2021-12-20 15:59:24 +08:00
committed by GitHub
5 changed files with 131 additions and 69 deletions
+9 -8
View File
@@ -966,7 +966,6 @@ func (s *SGuestReloadDiskTask) onResumeSucc(results string) {
}
func (s *SGuestReloadDiskTask) taskFailed(reason string) {
log.Errorf("SGuestReloadDiskTask error: %s", reason)
hostutils.TaskFailed(s.ctx, reason)
}
@@ -1000,20 +999,21 @@ func (s *SGuestDiskSnapshotTask) startSnapshot(device string) {
func (s *SGuestDiskSnapshotTask) onReloadBlkdevSucc(res string) {
var cb = s.onResumeSucc
if len(res) > 0 {
log.Errorf("Monitor reload blkdev error: %s", res)
cb = s.onSnapshotBlkdevFail
cb = func(string) {
s.onSnapshotBlkdevFail(fmt.Sprintf("onReloadBlkdevFail: %s", res))
}
}
s.Monitor.SimpleCommand("cont", cb)
}
func (s *SGuestDiskSnapshotTask) onSnapshotBlkdevFail(string) {
func (s *SGuestDiskSnapshotTask) onSnapshotBlkdevFail(reason string) {
snapshotDir := s.disk.GetSnapshotDir()
snapshotPath := path.Join(snapshotDir, s.snapshotId)
output, err := procutils.NewCommand("mv", "-f", snapshotPath, s.disk.GetPath()).Output()
if err != nil {
log.Errorf("mv %s to %s failed: %s, %s", snapshotPath, s.disk.GetPath(), err, output)
}
hostutils.TaskFailed(s.ctx, "Reload blkdev error")
hostutils.TaskFailed(s.ctx, fmt.Sprintf("Reload blkdev error: %s", reason))
}
func (s *SGuestDiskSnapshotTask) onResumeSucc(res string) {
@@ -1098,8 +1098,9 @@ func (s *SGuestSnapshotDeleteTask) doReloadDisk(device string) {
func (s *SGuestSnapshotDeleteTask) onReloadBlkdevSucc(err string) {
var callback = s.onResumeSucc
if len(err) > 0 {
log.Errorf("Reload blkdev failed: %s", err)
callback = s.onSnapshotBlkdevFail
callback = func(string) {
s.onSnapshotBlkdevFail(fmt.Sprintf("onReloadBlkdevFail %s", err))
}
}
s.Monitor.SimpleCommand("cont", callback)
}
@@ -1109,7 +1110,7 @@ func (s *SGuestSnapshotDeleteTask) onSnapshotBlkdevFail(res string) {
if output, err := procutils.NewCommand("mv", "-f", s.tmpPath, snapshotPath).Output(); err != nil {
log.Errorf("mv %s to %s failed: %s, %s", s.tmpPath, snapshotPath, err, output)
}
s.taskFailed("Reload blkdev failed")
s.taskFailed(fmt.Sprintf("Reload blkdev failed %s", res))
}
func (s *SGuestSnapshotDeleteTask) onResumeSucc(res string) {
+10 -18
View File
@@ -648,49 +648,41 @@ func (ms MirrorJob) InProcess() bool {
}
func (s *SKVMGuestInstance) MirrorJobStatus() MirrorJob {
res := make(chan *jsonutils.JSONArray)
s.Monitor.GetBlockJobs(func(jobs *jsonutils.JSONArray) {
res := make(chan []monitor.BlockJob)
s.Monitor.GetBlockJobs(func(jobs []monitor.BlockJob) {
res <- jobs
})
select {
case <-time.After(time.Second * 3):
return 0
case v := <-res:
if v != nil && v.Length() >= s.DiskCount() {
if len(v) >= s.DiskCount() {
mirrorSuccCount := 0
vs, _ := v.GetArray()
for _, val := range vs {
jobType, _ := val.GetString("type")
jobStatus, _ := val.GetString("status")
if jobType == "mirror" && jobStatus == "ready" {
for _, job := range v {
if job.Type == "mirror" && job.Status == "ready" {
mirrorSuccCount += 1
}
}
if mirrorSuccCount == s.DiskCount() {
return 1
} else {
return 0
}
} else {
return -1
return 0
}
return -1
}
}
func (s *SKVMGuestInstance) BlockJobsCount() int {
res := make(chan *jsonutils.JSONArray)
s.Monitor.GetBlockJobs(func(jobs *jsonutils.JSONArray) {
res := make(chan []monitor.BlockJob)
s.Monitor.GetBlockJobs(func(jobs []monitor.BlockJob) {
res <- jobs
})
select {
case <-time.After(time.Second * 3):
return -1
case v := <-res:
if v != nil && v.Length() > 0 {
return v.Length()
}
return len(v)
}
return 0
}
func (s *SKVMGuestInstance) CleanStartupTask() {
+14 -17
View File
@@ -364,29 +364,26 @@ func (m *HmpMonitor) GetBlockJobCounts(callback func(jobs int)) {
m.Query("info block-jobs", cb)
}
func (m *HmpMonitor) GetBlockJobs(callback func(*jsonutils.JSONArray)) {
func (m *HmpMonitor) GetBlockJobs(callback func([]BlockJob)) {
cb := func(output string) {
lines := strings.Split(strings.TrimSuffix(output, "\r\n"), "\r\n")
if lines[0] == "No active jobs" {
callback(nil)
} else {
res := jsonutils.NewArray()
re := regexp.MustCompile(`Type (?P<type>\w+), device (?P<device>\w+)`)
for i := 0; i < len(lines); i++ {
m := regutils2.GetParams(re, lines[i])
if len(m) > 0 {
jobType, _ := m["type"]
device, _ := m["device"]
jobInfo := jsonutils.NewDict()
jobInfo.Set("type", jsonutils.NewString(jobType))
jobInfo.Set("device", jsonutils.NewString(device))
res.Add(jobInfo)
}
}
callback(res)
return
}
jobs := []BlockJob{}
re := regexp.MustCompile(`Type (?P<type>\w+), device (?P<device>\w+)`)
for i := 0; i < len(lines); i++ {
m := regutils2.GetParams(re, lines[i])
if len(m) > 0 {
job := BlockJob{}
job.Type, _ = m["type"]
job.Device, _ = m["device"]
jobs = append(jobs, job)
}
}
callback(jobs)
}
m.Query("info block-jobs", cb)
}
+58 -1
View File
@@ -27,6 +27,63 @@ import (
type StringCallback func(string)
type BlockJob struct {
server string
Busy bool
// commit|
Type string
Len int64
Paused bool
Ready bool
// running|ready
Status string
// ok|
IoStatus string `json:"io-status"`
Offset int64
Device string
Speed int64
start time.Time
preOffset int64
now time.Time
}
type blockSizeByte int64
func (self blockSizeByte) String() string {
size := map[string]float64{
"Kb": 1024,
"Mb": 1024 * 1024,
"Gb": 1024 * 1024 * 1024,
"TB": 1024 * 1024 * 1024 * 1024,
}
for _, unit := range []string{"TB", "Gb", "Mb", "Kb"} {
if int64(self)/int64(size[unit]) > 0 {
return fmt.Sprintf("%.2f%s", float64(self)/size[unit], unit)
}
}
return fmt.Sprintf("%d", int64(self))
}
func (self *BlockJob) PreOffset(preOffset int64) {
if self.start.IsZero() {
self.start = time.Now()
self.now = time.Now()
self.preOffset = preOffset
return
}
second := time.Now().Sub(self.now).Seconds()
if second > 0 {
speed := float64(self.Offset-preOffset) / second
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())
}
self.preOffset = preOffset
self.now = time.Now()
return
}
type Monitor interface {
Connect(host string, port int) error
ConnectWithSocket(address string) error
@@ -40,7 +97,7 @@ type Monitor interface {
QueryStatus(StringCallback)
GetVersion(StringCallback)
GetBlockJobCounts(func(jobs int))
GetBlockJobs(func(*jsonutils.JSONArray))
GetBlockJobs(func([]BlockJob))
GetCpuCount(func(count int))
AddCpu(cpuIndex int, callback StringCallback)
+40 -25
View File
@@ -26,6 +26,7 @@ import (
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/utils"
)
@@ -105,6 +106,7 @@ type QmpMonitor struct {
qmpEventFunc qmpEventCallback
commandQueue []*Command
callbackQueue []qmpMonitorCallBack
jobs map[string]BlockJob
}
func NewQmpMonitor(server string, OnMonitorDisConnect, OnMonitorTimeout MonitorErrorFunc,
@@ -114,6 +116,7 @@ func NewQmpMonitor(server string, OnMonitorDisConnect, OnMonitorTimeout MonitorE
qmpEventFunc: qmpEventFunc,
commandQueue: make([]*Command, 0),
callbackQueue: make([]qmpMonitorCallBack, 0),
jobs: map[string]BlockJob{},
}
// On qmp init must set capabilities
@@ -658,40 +661,52 @@ func (m *QmpMonitor) MigrateStartPostcopy(callback StringCallback) {
m.Query(cmd, cb)
}
func (m *QmpMonitor) blockJobs(res *Response) ([]BlockJob, error) {
if res.ErrorVal != nil {
return nil, errors.Errorf("GetBlockJobs for %s %s", m.server, jsonutils.Marshal(res.ErrorVal).String())
}
ret, err := jsonutils.Parse(res.Return)
if err != nil {
return nil, errors.Wrapf(err, "GetBlockJobs for %s parse %s", m.server, res.Return)
}
jobs := []BlockJob{}
ret.Unmarshal(&jobs)
for i := range jobs {
job := jobs[i]
job.server = m.server
_job, ok := m.jobs[job.Device]
if !ok {
job.PreOffset(0)
m.jobs[job.Device] = job
continue
}
if _job.Status == "ready" {
delete(m.jobs, _job.Device)
continue
}
job.start, job.now = _job.start, _job.now
job.PreOffset(_job.Offset)
m.jobs[job.Device] = job
}
return jobs, nil
}
func (m *QmpMonitor) GetBlockJobCounts(callback func(jobs int)) {
var cb = func(res *Response) {
if res.ErrorVal != nil {
log.Errorf("GetBlockJobCounts error %s: %s", m.server, res.ErrorVal.Error())
jobs, err := m.blockJobs(res)
if err != nil {
callback(-1)
} else {
ret, err := jsonutils.Parse(res.Return)
if err != nil {
log.Errorf("Parse GetBlockJobCounts qmp res error %s: %s", m.server, err)
callback(-1)
} else {
jobs, _ := ret.GetArray()
callback(len(jobs))
}
return
}
callback(len(jobs))
}
m.Query(&Command{Execute: "query-block-jobs"}, cb)
}
func (m *QmpMonitor) GetBlockJobs(callback func(*jsonutils.JSONArray)) {
func (m *QmpMonitor) GetBlockJobs(callback func([]BlockJob)) {
var cb = func(res *Response) {
if res.ErrorVal != nil {
log.Errorf("GetBlockJobs error %s: %s", m.server, res.ErrorVal.Error())
callback(nil)
} else {
ret, err := jsonutils.Parse(res.Return)
if err != nil {
log.Errorf("Parse GetBlockJobs qmp res error %s: %s", m.server, err)
callback(nil)
} else {
jobs, _ := ret.(*jsonutils.JSONArray)
callback(jobs)
}
}
jobs, _ := m.blockJobs(res)
callback(jobs)
}
m.Query(&Command{Execute: "query-block-jobs"}, cb)
}