diff --git a/pkg/cloudcommon/dhcp/server.go b/pkg/cloudcommon/dhcp/server.go index 1b1e3bd83e..3abdcc4500 100644 --- a/pkg/cloudcommon/dhcp/server.go +++ b/pkg/cloudcommon/dhcp/server.go @@ -53,7 +53,6 @@ func (s *DHCPServer) ListenAndServe(handler DHCPHandler) error { func (s *DHCPServer) serveDHCP(handler DHCPHandler) error { for { pkt, addr, intf, err := s.conn.RecvDHCP() - log.Errorf("DHCP Receive from %s", addr.String()) if err != nil { return fmt.Errorf("Receiving DHCP packet: %s", err) } diff --git a/pkg/cloudcommon/service/services.go b/pkg/cloudcommon/service/services.go index ac09102f1c..b6d7cc66a1 100644 --- a/pkg/cloudcommon/service/services.go +++ b/pkg/cloudcommon/service/services.go @@ -11,14 +11,16 @@ import ( type SServiceBase struct{} -func (s *SServiceBase) RegisterSignals(quitHandler signalutils.Trap) { +func (s *SServiceBase) RegisterQuitSignals(quitHandler signalutils.Trap) { quitSignals := []os.Signal{syscall.SIGHUP, syscall.SIGINT, syscall.SIGQUIT, syscall.SIGTERM} signalutils.RegisterSignal(quitHandler, quitSignals...) + signalutils.StartTrap() +} + +func (s *SServiceBase) RegisterSIGUSR1() { // dump goroutine stack signalutils.RegisterSignal(func() { utils.DumpAllGoroutineStack(log.Logger().Out) }, syscall.SIGUSR1) - - signalutils.StartTrap() } diff --git a/pkg/hostman/downloader/downloader.go b/pkg/hostman/downloader/downloader.go index ab718c7862..4b4b25f590 100644 --- a/pkg/hostman/downloader/downloader.go +++ b/pkg/hostman/downloader/downloader.go @@ -46,6 +46,7 @@ func (d *SDownloadProvider) Start( d.w.Header().Add(k, headers.Get(k)) } + log.Infof("Downloader Start Transfer %s, compress %t", downloadFilePath, d.compress) fi, err := os.Open(downloadFilePath) if err != nil { log.Errorln(err) @@ -54,11 +55,12 @@ func (d *SDownloadProvider) Start( defer fi.Close() var ( - end = false - chunk = make([]byte, CHUNK_SIZE) - writer io.Writer = d.w - startTime = time.Now() - sendBytes = 0 + end = false + chunk = make([]byte, CHUNK_SIZE) + writer io.Writer = d.w + startTime = time.Now() + sendBytes = 0 + writeChunk []byte ) if d.compress { @@ -68,19 +70,23 @@ func (d *SDownloadProvider) Start( return err } writer = zw - defer zw.Flush() // it's cool defer zw.Close() + defer zw.Flush() // it's cool } for !end { - if _, err := fi.Read(chunk); err == io.EOF { - end = true - } else if err != nil && err != io.EOF { - log.Errorln(err) - return err + size, err := fi.Read(chunk) + if err != nil { + if err != io.EOF { + log.Errorln(err) + return err + } else { + end = true + } } - if size, err := writer.Write(chunk); err != nil { + writeChunk = chunk[:size] + if size, err = writer.Write(writeChunk); err != nil { log.Errorln(err) return err } else { @@ -100,7 +106,7 @@ func (d *SDownloadProvider) Start( sendMb := float64(sendBytes) / 1000.0 / 1000.0 timeDur := time.Now().Sub(startTime) - log.Infof("Send data: %fMB rate: %fMB/sec", sendMb/timeDur.Seconds()) + log.Infof("Send data: %fMB rate: %fMB/sec", sendMb, sendMb/timeDur.Seconds()) if onDownloadComplete != nil { onDownloadComplete() diff --git a/pkg/hostman/downloader/downloadhandler.go b/pkg/hostman/downloader/downloadhandler.go index 38261153a1..dfa964aa69 100644 --- a/pkg/hostman/downloader/downloadhandler.go +++ b/pkg/hostman/downloader/downloadhandler.go @@ -74,7 +74,9 @@ func download(ctx context.Context, w http.ResponseWriter, r *http.Request) { if !fileutils2.Exists(hand.downloadFilePath()) { httperrors.NotFoundError(w, "Image cache %s not found", id) } else { - hand.Start() + if err := hand.Start(); err != nil { + hostutils.Response(ctx, w, err) + } } case "servers": hand := NewGuestDownloadProvider(w, compress, rateLimit, id) @@ -144,7 +146,7 @@ func snapshotPrecheck( params, _, _ = appsrv.FetchEnv(ctx, w, r) storageId = params[""] diskId = params[""] - snapshotId = params["snapshotId"] + snapshotId = params[""] ) storage := storageman.GetManager().GetStorage(storageId) diff --git a/pkg/hostman/downloader/image_downloader.go b/pkg/hostman/downloader/image_downloader.go index 4f2b069072..bdeced2153 100644 --- a/pkg/hostman/downloader/image_downloader.go +++ b/pkg/hostman/downloader/image_downloader.go @@ -1,7 +1,6 @@ package downloader import ( - "fmt" "net/http" "os" @@ -40,7 +39,10 @@ func (i *SImageDownloadProvider) downloadFilePath() string { } func (i *SImageDownloadProvider) prepareDownload() error { - log.Infof(fmt.Sprintf("Compress %s to %s", i.fullPath(), i.downloadFilePath())) + if i.fullPath() != i.downloadFilePath() { + log.Infof("Compress %s %s to %s", i.compressFormat, i.fullPath(), i.downloadFilePath()) + } + switch i.compressFormat { case "qcow2": img, err := qemuimg.NewQemuImage(i.fullPath()) diff --git a/pkg/hostman/guestman/guesthandler.go b/pkg/hostman/guestman/guesthandler.go index d3d383727e..69bed3f526 100644 --- a/pkg/hostman/guestman/guesthandler.go +++ b/pkg/hostman/guestman/guesthandler.go @@ -15,11 +15,35 @@ import ( "yunion.io/x/onecloud/pkg/mcclient/auth" ) -var keyWords = []string{"servers"} - type strDict map[string]string type actionFunc func(context.Context, string, jsonutils.JSONObject) (interface{}, error) +var ( + keyWords = []string{"servers"} + actionFuncs = map[string]actionFunc{ + "create": guestCreate, + "deploy": guestDeploy, + "start": guestStart, + "stop": guestStop, + "monitor": guestMonitor, + "sync": guestSync, + "suspend": guestSuspend, + + "snapshot": guestSnapshot, + "delete-snapshot": guestDeleteSnapshot, + "reload-disk-snapshot": guestReloadDiskSnapshot, + // "remove-statefile": guestRemoveStatefile, + // "io-throttle": guestIoThrottle, + + "src-prepare-migrate": guestSrcPrepareMigrate, + "dest-prepare-migrate": guestDestPrepareMigrate, + "live-migrate": guestLiveMigrate, + "resume": guestResume, + // "start-nbd-server": guestStartNbdServer, + "drive-mirror": guestDriveMirror, + } +) + func AddGuestTaskHandler(prefix string, app *appsrv.Application) { for _, keyWord := range keyWords { app.AddHandler("GET", @@ -396,26 +420,3 @@ func guestDeleteSnapshot(ctx context.Context, sid string, body jsonutils.JSONObj hostutils.DelayTask(ctx, guestManger.DeleteSnapshot, params) return nil, nil } - -var actionFuncs = map[string]actionFunc{ - "create": guestCreate, - "deploy": guestDeploy, - "start": guestStart, - "stop": guestStop, - "monitor": guestMonitor, - "sync": guestSync, - "suspend": guestSuspend, - - "snapshot": guestSnapshot, - "delete-snapshot": guestDeleteSnapshot, - "reload-disk-snapshot": guestReloadDiskSnapshot, - // "remove-statefile": guestRemoveStatefile, - // "io-throttle": guestIoThrottle, - - "src-prepare-migrate": guestSrcPrepareMigrate, - "dest-prepare-migrate": guestDestPrepareMigrate, - "live-migrate": guestLiveMigrate, - "resume": guestResume, - // "start-nbd-server": guestStartNbdServer, - "drive-mirror": guestDriveMirror, -} diff --git a/pkg/hostman/guestman/guestman.go b/pkg/hostman/guestman/guestman.go index 2a4e5d09fc..4513c7121e 100644 --- a/pkg/hostman/guestman/guestman.go +++ b/pkg/hostman/guestman/guestman.go @@ -158,9 +158,11 @@ func (m *SGuestManager) StartCpusetBalancer() { return } go func() { - time.Sleep(time.Second * 120) - if options.HostOptions.EnableCpuBinding { - m.cpusetBalance() + for { + if options.HostOptions.EnableCpuBinding { + m.cpusetBalance() + } + time.Sleep(time.Second * 120) } }() } diff --git a/pkg/hostman/guestman/qemu-kvm.go b/pkg/hostman/guestman/qemu-kvm.go index edc32572e9..1093289b15 100644 --- a/pkg/hostman/guestman/qemu-kvm.go +++ b/pkg/hostman/guestman/qemu-kvm.go @@ -336,14 +336,36 @@ func (s *SKVMGuestInstance) StartMonitor(ctx context.Context) { func (s *SKVMGuestInstance) delayStartMonitor(ctx context.Context) { if options.HostOptions.EnableQmpMonitor && s.GetQmpMonitorPort(-1) > 0 { s.Monitor = monitor.NewQmpMonitor( - s.onMonitorDisConnect, - func(err error) { s.onMonitorTimeout(ctx, err) }, - func() { s.onMonitorConnected(ctx) }, + s.onMonitorDisConnect, // on monitor disconnect + func(err error) { s.onMonitorTimeout(ctx, err) }, // on monitor timeout + func() { s.onMonitorConnected(ctx) }, // on monitor connected + s.onReceiveQMPEvent, // on reveive qmp event ) s.Monitor.Connect("127.0.0.1", s.GetQmpMonitorPort(-1)) } } +func (s *SKVMGuestInstance) onReceiveQMPEvent(event *monitor.Event) { + if event.Event == "BLOCK_JOB_READY" && s.IsMaster() { + if itype, ok := event.Data["type"]; ok { + stype, _ := itype.(string) + if stype == "mirror" { + if s.mirrorJobSuccCount != nil { + *s.mirrorJobSuccCount += 1 + } else { + s.mirrorJobSuccCount = new(int) + *s.mirrorJobSuccCount = 1 + } + if *s.mirrorJobSuccCount == s.DiskCount() { + hostutils.UpdateServerStatus(context.Background(), s.GetId(), "running") + } + } + } + } else if event.Event == "BLOCK_JOB_ERROR" && s.IsMaster() { + hostutils.UpdateServerStatus(context.Background(), s.GetId(), "mirror_failed") + } +} + func (s *SKVMGuestInstance) onMonitorConnected(ctx context.Context) { log.Infof("Monitor connected ...") s.Monitor.GetVersion(func(v string) { diff --git a/pkg/hostman/guestman/qemu-kvmhelper.go b/pkg/hostman/guestman/qemu-kvmhelper.go index c2120faa9b..a1a8479c30 100644 --- a/pkg/hostman/guestman/qemu-kvmhelper.go +++ b/pkg/hostman/guestman/qemu-kvmhelper.go @@ -3,6 +3,7 @@ package guestman import ( "fmt" "path" + "time" "yunion.io/x/jsonutils" "yunion.io/x/log" @@ -623,6 +624,18 @@ func (s *SKVMGuestInstance) generateStopScript(data *jsonutils.JSONDict) string return cmd } +func (s *SKVMGuestInstance) presendArpForNic(nic jsonutils.JSONObject) { + +} + func (s *SKVMGuestInstance) StartPresendArp() { - // TODO go func + go func() { + for i := 0; i < 5; i++ { + nics, _ := s.Desc.GetArray("nics") + for _, nic := range nics { + s.presendArpForNic(nic) + } + time.Sleep(1 * time.Second) + } + }() } diff --git a/pkg/hostman/host_services.go b/pkg/hostman/host_services.go index db967a21a1..9e715f3268 100644 --- a/pkg/hostman/host_services.go +++ b/pkg/hostman/host_services.go @@ -25,19 +25,16 @@ type SHostService struct { func (host *SHostService) StartService() { cloudcommon.ParseOptions(&options.HostOptions, os.Args, "host.conf", "host") - - // disable rbac - options.HostOptions.EnableRbac = false + options.HostOptions.EnableRbac = false // disable rbac app := cloudcommon.InitApp(&options.HostOptions.CommonOptions, false) - hostInstance := hostinfo.Instance() if err := hostInstance.Init(); err != nil { log.Fatalf(err.Error()) } - // register quit handler - host.RegisterSignals(func() { + host.RegisterSIGUSR1() + host.RegisterQuitSignals(func() { // register quit handler if host.isExiting { return } else { @@ -64,15 +61,14 @@ func (host *SHostService) StartService() { } guestman.Init(hostInstance, options.HostOptions.ServersPath) - cloudcommon.InitAuth(&options.HostOptions.CommonOptions, func() { log.Infof("Auth complete!!") - hostInstance.StartRegister(5, guestman.GetGuestManager().Bootstrap) + // ??? Why wait 5 seconds + hostInstance.StartRegister(2, guestman.GetGuestManager().Bootstrap) }) - host.initHandlers(app) - <-hostinfo.Instance().IsRegistered // wait host and guest init + <-hostinfo.Instance().IsRegistered // wait host and guest init cloudcommon.ServeForever(app, &options.HostOptions.CommonOptions) } diff --git a/pkg/hostman/monitor/qmp.go b/pkg/hostman/monitor/qmp.go index 980ba9c720..058b083a68 100644 --- a/pkg/hostman/monitor/qmp.go +++ b/pkg/hostman/monitor/qmp.go @@ -30,6 +30,7 @@ Not support oob yet */ type qmpMonitorCallBack func(*Response) +type qmpEventCallback func(*Event) type Response struct { Return []byte @@ -83,14 +84,16 @@ func (e *Error) Error() string { type QmpMonitor struct { SBaseMonitor + qmpEventFunc qmpEventCallback commandQueue []*Command callbackQueue []qmpMonitorCallBack } func NewQmpMonitor(OnMonitorDisConnect, OnMonitorTimeout MonitorErrorFunc, - OnMonitorConnected MonitorSuccFunc) *QmpMonitor { + OnMonitorConnected MonitorSuccFunc, qmpEventFunc qmpEventCallback) *QmpMonitor { m := &QmpMonitor{ SBaseMonitor: *NewBaseMonitor(OnMonitorConnected, OnMonitorDisConnect, OnMonitorTimeout), + qmpEventFunc: qmpEventFunc, commandQueue: make([]*Command, 0), callbackQueue: make([]qmpMonitorCallBack, 0), } @@ -197,8 +200,11 @@ func (m *QmpMonitor) read(r io.Reader) { m.reading = false } -func (m QmpMonitor) watchEvent(event *Event) { +func (m *QmpMonitor) watchEvent(event *Event) { log.Infof(event.String()) + if m.qmpEventFunc != nil { + go m.qmpEventFunc(event) + } } func (m *QmpMonitor) write(cmd []byte) error { diff --git a/pkg/hostman/monitor/qmp_test.go b/pkg/hostman/monitor/qmp_test.go index a546ef5ba3..7d1c08cc17 100644 --- a/pkg/hostman/monitor/qmp_test.go +++ b/pkg/hostman/monitor/qmp_test.go @@ -11,7 +11,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) + m := NewQmpMonitor(onDisConnect, onTimeout, onConnected, nil) var host = "127.0.0.1" var port = 56101 m.Connect(host, port) diff --git a/pkg/hostman/storageman/storagelocal.go b/pkg/hostman/storageman/storagelocal.go index 574da52b2a..fcb66b3175 100644 --- a/pkg/hostman/storageman/storagelocal.go +++ b/pkg/hostman/storageman/storagelocal.go @@ -55,8 +55,7 @@ func (s *SLocalStorage) GetSnapshotDir() string { } func (s *SLocalStorage) GetSnapshotPathByIds(diskId, snapshotId string) string { - return path.Join(s.GetSnapshotDir(), - diskId+options.HostOptions.SnapshotDirSuffix, snapshotId) + return path.Join(s.GetSnapshotDir(), diskId+options.HostOptions.SnapshotDirSuffix, snapshotId) } func (s *SLocalStorage) SyncStorageInfo() (jsonutils.JSONObject, error) { diff --git a/pkg/util/cgrouputils/cgrouputils.go b/pkg/util/cgrouputils/cgrouputils.go index 193b24f82f..e774e86961 100644 --- a/pkg/util/cgrouputils/cgrouputils.go +++ b/pkg/util/cgrouputils/cgrouputils.go @@ -228,7 +228,7 @@ func (c *CGroupTask) MoveTasksToRoot() { func (c *CGroupTask) RemoveTask() bool { if c.taskIsExist() { c.MoveTasksToRoot() - if err := os.RemoveAll(c.TaskPath()); err != nil { + if err := os.Remove(c.TaskPath()); err != nil { log.Errorf("Remove task %s", err) return false } @@ -529,8 +529,11 @@ func Init() bool { } func CgroupSet(pid string, coreNum int) bool { - tasks := []ICGroupTask{&CGroupCPUTask{&CGroupTask{}}, &CGroupIOTask{&CGroupTask{}}, - &CGroupMemoryTask{&CGroupTask{}}} + tasks := []ICGroupTask{ + &CGroupCPUTask{&CGroupTask{}}, + &CGroupIOTask{&CGroupTask{}}, + &CGroupMemoryTask{&CGroupTask{}}, + } for _, hand := range tasks { hand.SetHand(hand) hand.SetPid(pid) @@ -551,9 +554,13 @@ func CgroupIoHardlimitSet( } func CgroupDestroy(pid string) bool { - tasks := []ICGroupTask{&CGroupCPUTask{&CGroupTask{}}, &CGroupIOTask{&CGroupTask{}}, - &CGroupMemoryTask{&CGroupTask{}}, &CGroupCPUSetTask{&CGroupTask{}, ""}, - &CGroupIOHardlimitTask{CGroupIOTask: &CGroupIOTask{&CGroupTask{}}}} + tasks := []ICGroupTask{ + &CGroupCPUTask{&CGroupTask{}}, + &CGroupIOTask{&CGroupTask{}}, + &CGroupMemoryTask{&CGroupTask{}}, + &CGroupCPUSetTask{&CGroupTask{}, ""}, + &CGroupIOHardlimitTask{CGroupIOTask: &CGroupIOTask{&CGroupTask{}}}, + } for _, hand := range tasks { hand.SetHand(hand) hand.SetPid(pid) @@ -565,9 +572,13 @@ func CgroupDestroy(pid string) bool { } func CgroupCleanAll() { - tasks := []ICGroupTask{&CGroupCPUTask{&CGroupTask{}}, &CGroupIOTask{&CGroupTask{}}, - &CGroupMemoryTask{&CGroupTask{}}, &CGroupCPUSetTask{CGroupTask: &CGroupTask{}}, - &CGroupIOHardlimitTask{CGroupIOTask: &CGroupIOTask{&CGroupTask{}}}} + tasks := []ICGroupTask{ + &CGroupCPUTask{&CGroupTask{}}, + &CGroupIOTask{&CGroupTask{}}, + &CGroupMemoryTask{&CGroupTask{}}, + &CGroupCPUSetTask{CGroupTask: &CGroupTask{}}, + &CGroupIOHardlimitTask{CGroupIOTask: &CGroupIOTask{&CGroupTask{}}}, + } for _, hand := range tasks { hand.SetHand(hand) CleanupNonexistPids(hand.Module()) diff --git a/pkg/util/cgrouputils/cpusetutils.go b/pkg/util/cgrouputils/cpusetutils.go index 5a6e2c4b20..99317dc4f6 100644 --- a/pkg/util/cgrouputils/cpusetutils.go +++ b/pkg/util/cgrouputils/cpusetutils.go @@ -210,7 +210,7 @@ func FetchHistoryUtil() map[string][]float64 { return utilHistory } var objmap map[string]*json.RawMessage - if err := json.Unmarshal([]byte(contents), objmap); err != nil { + if err := json.Unmarshal([]byte(contents), &objmap); err != nil { log.Errorf("FetchHistoryUtil error: %s", err) return utilHistory } diff --git a/pkg/util/fileutils2/fileutils.go b/pkg/util/fileutils2/fileutils.go index b5e291537e..d4d88f7f37 100644 --- a/pkg/util/fileutils2/fileutils.go +++ b/pkg/util/fileutils2/fileutils.go @@ -399,7 +399,6 @@ func GetDevSector512Count(dev string) int { return size } -// TODO test func ResizeDiskFs(diskPath string, sizeMb int) error { var cmds = []string{"parted", "-a", "none", "-s", diskPath, "--", "unit", "s", "print"} lines, err := procutils.NewCommand(cmds[0], cmds[1:]...).Run() @@ -434,7 +433,7 @@ func ResizeDiskFs(diskPath string, sizeMb int) error { return err } for _, s := range []string{"r", "e", "Y", "w", "Y", "Y"} { - io.WriteString(stdin, s) + io.WriteString(stdin, fmt.Sprintf("%s\n", s)) } stdoutPut, err := ioutil.ReadAll(outb) if err != nil { @@ -445,12 +444,15 @@ func ResizeDiskFs(diskPath string, sizeMb int) error { return err } log.Infof("gdisk: %s %s", stdoutPut, stderrOutPut) - proc.Wait() + if err = proc.Wait(); err != nil { + log.Errorln(err) + return err + } } if len(parts) > 0 && (label == "gpt" || (label == "msdos" && parts[len(parts)-1][5] == "primary")) { var ( - part = parts[len(parts)] + part = parts[len(parts)-1] end int ) if sizeMb > 0 { @@ -472,10 +474,10 @@ func ResizeDiskFs(diskPath string, sizeMb int) error { if len(part[1]) > 0 { cmds = append(cmds, "set", part[0], "boot", "on") } - _, err := procutils.NewCommand(cmds[0], cmds[1:]...).Run() + output, err := procutils.NewCommand(cmds[0], cmds[1:]...).Run() if err != nil { - log.Errorln(err) - return err + log.Errorln(output) + return fmt.Errorf("%s", output) } if len(part[6]) > 0 { err := ResizePartitionFs(part[7], part[6])