This commit is contained in:
wanyaoqi
2019-01-24 11:48:12 +08:00
parent cdf2d485d5
commit 8ac044c8c5
16 changed files with 148 additions and 85 deletions
-1
View File
@@ -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)
}
+5 -3
View File
@@ -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()
}
+19 -13
View File
@@ -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()
+4 -2
View File
@@ -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["<storageId>"]
diskId = params["<diskId>"]
snapshotId = params["snapshotId"]
snapshotId = params["<snapshotId>"]
)
storage := storageman.GetManager().GetStorage(storageId)
+4 -2
View File
@@ -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())
+26 -25
View File
@@ -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,
}
+5 -3
View File
@@ -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)
}
}()
}
+25 -3
View File
@@ -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) {
+14 -1
View File
@@ -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)
}
}()
}
+6 -10
View File
@@ -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)
}
+8 -2
View File
@@ -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 {
+1 -1
View File
@@ -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)
+1 -2
View File
@@ -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) {
+20 -9
View File
@@ -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())
+1 -1
View File
@@ -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
}
+9 -7
View File
@@ -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])