From 8d18ede79d37b973a6e89c0dcf6ffa2e7cf52f0d Mon Sep 17 00:00:00 2001 From: wanyaoqi Date: Thu, 9 Mar 2023 12:10:45 +0800 Subject: [PATCH] fix(region,host): host health misc fix Signed-off-by: wanyaoqi --- build/docker/Dockerfile.host-health | 3 + cmd/host-health/main.go | 128 ++++++++++++++ pkg/cloudmon/misc/pinger.go | 6 +- pkg/cloudmon/misc/pingutils.go | 10 +- pkg/compute/models/host_health.go | 32 ++-- pkg/compute/models/hosts.go | 38 +++- pkg/hostman/guestman/guestman.go | 14 +- pkg/hostman/host_health/health_manager.go | 164 ++++++++---------- pkg/hostman/host_services.go | 64 ------- pkg/hostman/hosthandler/handler.go | 30 ++-- pkg/hostman/hostinfo/hostconsts/hostconsts.go | 3 +- pkg/hostman/hostinfo/hostinfo.go | 15 +- pkg/hostman/hostutils/hostutils.go | 18 ++ 13 files changed, 310 insertions(+), 215 deletions(-) create mode 100644 build/docker/Dockerfile.host-health create mode 100644 cmd/host-health/main.go diff --git a/build/docker/Dockerfile.host-health b/build/docker/Dockerfile.host-health new file mode 100644 index 0000000000..ad056c181c --- /dev/null +++ b/build/docker/Dockerfile.host-health @@ -0,0 +1,3 @@ +FROM registry.cn-beijing.aliyuncs.com/yunionio/onecloud-base:v0.3.5-1 + +ADD ./_output/alpine-build/bin/host-health /opt/yunion/bin/host-health diff --git a/cmd/host-health/main.go b/cmd/host-health/main.go new file mode 100644 index 0000000000..62552a75d0 --- /dev/null +++ b/cmd/host-health/main.go @@ -0,0 +1,128 @@ +package main + +import ( + "context" + "io/ioutil" + "os" + "path/filepath" + + execlient "yunion.io/x/executor/client" + "yunion.io/x/log" + "yunion.io/x/pkg/errors" + + app_common "yunion.io/x/onecloud/pkg/cloudcommon/app" + "yunion.io/x/onecloud/pkg/cloudcommon/service" + "yunion.io/x/onecloud/pkg/hostman/host_health" + "yunion.io/x/onecloud/pkg/hostman/options" + "yunion.io/x/onecloud/pkg/httperrors" + "yunion.io/x/onecloud/pkg/util/atexit" + "yunion.io/x/onecloud/pkg/util/procutils" + "yunion.io/x/onecloud/pkg/util/sysutils" +) + +type SHostHealthService struct { + *service.SServiceBase +} + +func (host *SHostHealthService) InitService() { + options.Init() + isRoot := sysutils.IsRootPermission() + if !isRoot { + log.Fatalf("host service must running with root permissions") + } + + if len(options.HostOptions.DeployServerSocketPath) == 0 { + log.Fatalf("missing deploy server socket path") + } + + // options.HostOptions.EnableRbac = false // disable rbac + // init base option for pid file + host.SServiceBase.O = &options.HostOptions.BaseOptions + + log.Infof("exec socket path: %s", options.HostOptions.ExecutorSocketPath) + if options.HostOptions.EnableRemoteExecutor { + execlient.Init(options.HostOptions.ExecutorSocketPath) + procutils.SetRemoteExecutor() + } +} + +func (host *SHostHealthService) RunService() { + hn, err := os.Hostname() + if err != nil { + log.Fatalf("fail to get hostname %s", err) + } + app_common.InitAuth(&options.HostOptions.CommonOptions, func() { + log.Infof("Auth complete!!") + + if err := host.initEtcdConfig(); err != nil { + log.Fatalln("Init etcd config:", err) + } + + if len(options.HostOptions.EtcdEndpoints) > 0 { + _, err := host_health.InitHostHealthManager(hn) + if err != nil { + log.Fatalf("Init host health manager failed %s", err) + } + } + }) + select {} +} + +func writeFile(dir, file string, data []byte) (string, error) { + p := filepath.Join(dir, file) + return p, ioutil.WriteFile(p, data, 0600) +} + +func (host *SHostHealthService) initEtcdConfig() error { + etcdEndpoint, err := app_common.FetchEtcdServiceInfo() + if err != nil { + if errors.Cause(err) == httperrors.ErrNotFound { + return nil + } + return errors.Wrap(err, "fetch etcd service info") + } + if etcdEndpoint == nil { + return nil + } + if len(options.HostOptions.EtcdEndpoints) == 0 { + options.HostOptions.EtcdEndpoints = []string{etcdEndpoint.Url} + if len(etcdEndpoint.CertId) > 0 { + dir, err := ioutil.TempDir("", "etcd-cluster-tls") + if err != nil { + return errors.Wrap(err, "create dir etcd cluster tls") + } + options.HostOptions.EtcdCert, err = writeFile(dir, "etcd.crt", []byte(etcdEndpoint.Certificate)) + if err != nil { + return errors.Wrap(err, "write file certificate") + } + options.HostOptions.EtcdKey, err = writeFile(dir, "etcd.key", []byte(etcdEndpoint.PrivateKey)) + if err != nil { + return errors.Wrap(err, "write file private key") + } + options.HostOptions.EtcdCacert, err = writeFile(dir, "etcd-ca.crt", []byte(etcdEndpoint.CaCertificate)) + if err != nil { + return errors.Wrap(err, "write file cacert") + } + options.HostOptions.EtcdUseTLS = true + } + } + return nil +} + +func (host *SHostHealthService) OnExitService() {} + +func StartService() { + var srv = &SHostHealthService{} + srv.SServiceBase = &service.SServiceBase{ + Service: srv, + } + srv.StartService() +} + +func main() { + defer atexit.Handle() + + go procutils.WaitZombieLoop(context.TODO()) + + StartService() +} diff --git a/pkg/cloudmon/misc/pinger.go b/pkg/cloudmon/misc/pinger.go index b1f1f219e5..529efe060c 100644 --- a/pkg/cloudmon/misc/pinger.go +++ b/pkg/cloudmon/misc/pinger.go @@ -132,7 +132,11 @@ func pingProbeNetwork(s *mcclient.ClientSession, net sNetwork) ([]influxdb.SMetr addrStr := addr.String() pingAddrs = append(pingAddrs, addrStr) } - pingResults, err := Ping(pingAddrs) + pingResults, err := Ping(pingAddrs, + options.Options.PingProbeOptions.ProbeCount, + options.Options.PingProbeOptions.TimeoutSecond, + options.Options.PingProbeOptions.Debug, + ) if err != nil { return nil, errors.Wrap(err, "Ping") } diff --git a/pkg/cloudmon/misc/pingutils.go b/pkg/cloudmon/misc/pingutils.go index e7e2027c31..df79628385 100644 --- a/pkg/cloudmon/misc/pingutils.go +++ b/pkg/cloudmon/misc/pingutils.go @@ -22,8 +22,6 @@ import ( "github.com/tatsushid/go-fastping" "yunion.io/x/pkg/errors" - - "yunion.io/x/onecloud/pkg/cloudmon/options" ) type SPingResult struct { @@ -72,13 +70,13 @@ func (pr SPingResult) String() string { return fmt.Sprintf("%d packets transmitted, %d received, %d%% packet loss, rtt min/avg/max = %d/%d/%d ms", pr.count, len(pr.rtt), pr.Loss(), min/time.Millisecond, avg/time.Millisecond, max/time.Millisecond) } -func Ping(addrList []string) (map[string]*SPingResult, error) { +func Ping(addrList []string, probeCount, timeoutSecond int, debug bool) (map[string]*SPingResult, error) { p := fastping.NewPinger() - count := options.Options.PingProbeOptions.ProbeCount - timeout := time.Second * time.Duration(options.Options.PingProbeOptions.TimeoutSecond) + count := probeCount + timeout := time.Second * time.Duration(timeoutSecond) p.MaxRTT = timeout p.Size = 64 - p.Debug = options.Options.PingProbeOptions.Debug + p.Debug = debug result := make(map[string]*SPingResult) for _, addr := range addrList { result[addr] = NewPingResult(addr, count) diff --git a/pkg/compute/models/host_health.go b/pkg/compute/models/host_health.go index 52999801eb..0a4796b73f 100644 --- a/pkg/compute/models/host_health.go +++ b/pkg/compute/models/host_health.go @@ -22,11 +22,11 @@ import ( "yunion.io/x/log" "yunion.io/x/pkg/errors" - "yunion.io/x/pkg/utils" api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "yunion.io/x/onecloud/pkg/cloudcommon/etcd" + "yunion.io/x/onecloud/pkg/cloudmon/misc" "yunion.io/x/onecloud/pkg/mcclient/auth" ) @@ -99,8 +99,8 @@ func (h *SHostHealthChecker) startWatcher(ctx context.Context, hostname string) } if err := h.cli.Watch( ctx, key, - h.onHostOnline(ctx, hostname), - h.onHostOffline(ctx, hostname), + h.onHostOnlineCreated(ctx, hostname), + h.onHostOnlineModified(ctx, hostname), h.onHostOfflineDeleted(ctx, hostname), ); err != nil { return err @@ -131,15 +131,22 @@ func (h *SHostHealthChecker) onHostUnhealthy(ctx context.Context, hostname strin lockman.LockRawObject(ctx, api.HOST_HEALTH_LOCK_PREFIX, hostname) defer lockman.ReleaseRawObject(ctx, api.HOST_HEALTH_LOCK_PREFIX, hostname) host := HostManager.FetchHostByHostname(hostname) - if host != nil && !utils.IsInStringArray(host.RemoteHealthStatus(ctx), - []string{api.HOST_HEALTH_STATUS_RECONNECTING, api.HOST_HEALTH_STATUS_RUNNING}, - ) { - // in case hostagent health manager in status reconnecting - host.OnHostDown(ctx, auth.AdminCredential()) + if host != nil { + pingRes, err := misc.Ping([]string{host.AccessIp}, 3, 10, true) + if err != nil { + log.Errorf("failed ping dest host %s", hostname) + return + } + if ps := pingRes[host.AccessIp]; ps.Loss() < 100 { + log.Infof("ping host %s access ip %s succeed %s, skip host down", hostname, host.AccessIp, ps) + } else { + log.Errorf("ping host %s access ip %s failed %s, host down", hostname, host.AccessIp, ps) + host.OnHostDown(ctx, auth.AdminCredential()) + } } } -func (h *SHostHealthChecker) onHostOnline(ctx context.Context, hostname string) etcd.TEtcdCreateEventFunc { +func (h *SHostHealthChecker) onHostOnlineCreated(ctx context.Context, hostname string) etcd.TEtcdCreateEventFunc { return func(ctx context.Context, key, value []byte) { log.Infof("Got host online %s", hostname) if v, ok := h.hc.Load(hostname); ok { @@ -163,10 +170,10 @@ func (h *SHostHealthChecker) processHostOffline(ctx context.Context, hostname st }() } -func (h *SHostHealthChecker) onHostOffline(ctx context.Context, hostname string) etcd.TEtcdModifyEventFunc { +func (h *SHostHealthChecker) onHostOnlineModified(ctx context.Context, hostname string) etcd.TEtcdModifyEventFunc { return func(ctx context.Context, key, oldvalue, value []byte) { - log.Errorf("watch host key modified %s %s %s", key, oldvalue, value) - h.processHostOffline(ctx, hostname) + log.Infof("watch host key modified %s %s %s", key, oldvalue, value) + h.onHostOnlineCreated(ctx, hostname) } } @@ -178,6 +185,7 @@ func (h *SHostHealthChecker) onHostOfflineDeleted(ctx context.Context, hostname } func (h *SHostHealthChecker) WatchHost(ctx context.Context, hostname string) error { + h.onHostOnlineCreated(ctx, hostname) h.cli.Unwatch(hostKey(hostname)) return h.startWatcher(ctx, hostname) } diff --git a/pkg/compute/models/hosts.go b/pkg/compute/models/hosts.go index c11ca97e50..a88209e5c8 100644 --- a/pkg/compute/models/hosts.go +++ b/pkg/compute/models/hosts.go @@ -1121,6 +1121,31 @@ func (self *SHostManager) IsNewNameUnique(name string, userCred mcclient.TokenCr return cnt == 0, nil } +func (self *SHostManager) GetPropertyK8sMasterNodeIps(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) (jsonutils.JSONObject, error) { + cli, err := tokens.GetCoreClient() + if err != nil { + return nil, errors.Wrap(err, "get k8s client") + } + nodes, err := cli.Nodes().List(context.Background(), metav1.ListOptions{ + LabelSelector: "node-role.kubernetes.io/master", + }) + if err != nil { + return nil, errors.Wrap(err, "list master nodes") + } + ips := make([]string, 0) + for i := range nodes.Items { + for j := range nodes.Items[i].Status.Addresses { + if nodes.Items[i].Status.Addresses[j].Type == v1.NodeInternalIP { + ips = append(ips, nodes.Items[i].Status.Addresses[j].Address) + } + } + } + log.Infof("k8s master nodes ips %v", ips) + res := jsonutils.NewDict() + res.Set("ips", jsonutils.Marshal(ips)) + return res, nil +} + func (self *SHostManager) GetPropertyBmStartRegisterScript(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) (jsonutils.JSONObject, error) { regionUri, err := auth.GetPublicServiceURL("compute_v2", options.Options.Region, "") if err != nil { @@ -4068,15 +4093,18 @@ func (self *SHost) PerformAutoMigrateOnHostDown( meta[api.HOSTMETA_AUTO_MIGRATE_ON_HOST_SHUTDOWN] = "disable" } + data := jsonutils.NewDict() if input.AutoMigrateOnHostDown == "enable" { + data.Set("shutdown_servers", jsonutils.JSONTrue) meta[api.HOSTMETA_AUTO_MIGRATE_ON_HOST_DOWN] = "enable" - _, err := self.Request(ctx, userCred, "POST", "/hosts/shutdown-servers-on-host-down", - mcclient.GetTokenHeaders(userCred), nil) - if err != nil { - return nil, err - } } else if input.AutoMigrateOnHostDown == "disable" { meta[api.HOSTMETA_AUTO_MIGRATE_ON_HOST_DOWN] = "disable" + data.Set("shutdown_servers", jsonutils.JSONFalse) + } + _, err := self.Request(ctx, userCred, "POST", "/hosts/shutdown-servers-on-host-down", + mcclient.GetTokenHeaders(userCred), data) + if err != nil { + return nil, err } return nil, self.SetAllMetadata(ctx, meta, userCred) diff --git a/pkg/hostman/guestman/guestman.go b/pkg/hostman/guestman/guestman.go index d769ba19c0..4a9d0b54f9 100644 --- a/pkg/hostman/guestman/guestman.go +++ b/pkg/hostman/guestman/guestman.go @@ -29,7 +29,6 @@ import ( "yunion.io/x/jsonutils" "yunion.io/x/log" "yunion.io/x/pkg/errors" - "yunion.io/x/pkg/util/regutils" "yunion.io/x/pkg/util/seclib" "yunion.io/x/pkg/utils" @@ -53,7 +52,6 @@ import ( modules "yunion.io/x/onecloud/pkg/mcclient/modules/compute" "yunion.io/x/onecloud/pkg/util/cgrouputils" "yunion.io/x/onecloud/pkg/util/cgrouputils/cpuset" - "yunion.io/x/onecloud/pkg/util/fileutils2" "yunion.io/x/onecloud/pkg/util/netutils2" "yunion.io/x/onecloud/pkg/util/procutils" "yunion.io/x/onecloud/pkg/util/timeutils2" @@ -342,17 +340,7 @@ func (m *SGuestManager) CPUSetRemove(ctx context.Context, sid string) error { } func (m *SGuestManager) IsGuestDir(f os.FileInfo) bool { - if !regutils.MatchUUID(f.Name()) { - return false - } - if !f.Mode().IsDir() && f.Mode()&os.ModeSymlink == 0 { - return false - } - descFile := path.Join(m.ServersPath, f.Name(), "desc") - if !fileutils2.Exists(descFile) { - return false - } - return true + return hostutils.IsGuestDir(f, m.ServersPath) } func (m *SGuestManager) IsGuestExist(sid string) bool { diff --git a/pkg/hostman/host_health/health_manager.go b/pkg/hostman/host_health/health_manager.go index 1d3904c80f..4f0684ff8b 100644 --- a/pkg/hostman/host_health/health_manager.go +++ b/pkg/hostman/host_health/health_manager.go @@ -17,22 +17,23 @@ package host_health import ( "context" "fmt" - "os" - "strconv" - "strings" + "io/ioutil" + "path" "time" "yunion.io/x/log" "yunion.io/x/pkg/errors" - "yunion.io/x/pkg/utils" api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/etcd" common_options "yunion.io/x/onecloud/pkg/cloudcommon/options" - "yunion.io/x/onecloud/pkg/hostman/guestman/types" + "yunion.io/x/onecloud/pkg/cloudmon/misc" "yunion.io/x/onecloud/pkg/hostman/hostinfo/hostconsts" + "yunion.io/x/onecloud/pkg/hostman/hostutils" "yunion.io/x/onecloud/pkg/hostman/options" + modules "yunion.io/x/onecloud/pkg/mcclient/modules/compute" "yunion.io/x/onecloud/pkg/util/fileutils2" + "yunion.io/x/onecloud/pkg/util/procutils" ) type SHostHealthManager struct { @@ -41,24 +42,31 @@ type SHostHealthManager struct { timeout int requestExpend int - hostId string - status string - onHostDown string + hostId string + status string + + masterNodesIps []string } var ( - manager *SHostHealthManager - HostDownActions = []string{hostconsts.SHUTDOWN_SERVERS} + manager *SHostHealthManager ) -func InitHostHealthManager(hostId, onHostDown string) (*SHostHealthManager, error) { +func InitHostHealthManager(hostId string) (*SHostHealthManager, error) { if manager != nil { return manager, nil } var m = SHostHealthManager{} - var dialTimeout, requestTimeout = 3, 2 + masterNodesIps, err := m.masterNodesInternalIps() + if err != nil { + return nil, err + } else if len(masterNodesIps) == 0 { + return nil, errors.Errorf("failed get k8s master nodes") + } + m.masterNodesIps = masterNodesIps + var dialTimeout, requestTimeout = 3, 2 cfg, err := NewEtcdOptions( &options.HostOptions.EtcdOptions, options.HostOptions.HostLeaseTimeout, @@ -73,7 +81,6 @@ func InitHostHealthManager(hostId, onHostDown string) (*SHostHealthManager, erro } m.cli = etcd.Default() - m.onHostDown = onHostDown m.hostId = hostId m.requestExpend = requestTimeout m.timeout = options.HostOptions.HostHealthTimeout - options.HostOptions.HostLeaseTimeout @@ -81,6 +88,7 @@ func InitHostHealthManager(hostId, onHostDown string) (*SHostHealthManager, erro if err := m.StartHealthCheck(); err != nil { return nil, err } + log.Infof("put key %s success", m.GetKey()) m.status = api.HOST_HEALTH_STATUS_RUNNING manager = &m return manager, nil @@ -105,32 +113,33 @@ func NewEtcdOptions( func (m *SHostHealthManager) StartHealthCheck() error { return m.cli.PutSession(context.Background(), - fmt.Sprintf("%s/%s", api.HOST_HEALTH_PREFIX, m.hostId), - api.HOST_HEALTH_STATUS_RUNNING, + m.GetKey(), api.HOST_HEALTH_STATUS_RUNNING, ) } +func (m *SHostHealthManager) GetKey() string { + return fmt.Sprintf("%s/%s", api.HOST_HEALTH_PREFIX, m.hostId) +} + func (m *SHostHealthManager) OnKeepaliveFailure() { m.status = api.HOST_HEALTH_STATUS_RECONNECTING - nicRecord := m.recordNic() ctx, cancel := context.WithTimeout(context.Background(), time.Second*time.Duration(m.timeout)) defer cancel() err := m.cli.RestartSessionWithContext(ctx) if err == nil { if err := m.cli.PutSession(context.Background(), - fmt.Sprintf("%s/%s", api.HOST_HEALTH_PREFIX, m.hostId), - api.HOST_HEALTH_STATUS_RUNNING, + m.GetKey(), api.HOST_HEALTH_STATUS_RUNNING, ); err != nil { log.Errorf("put host key failed %s", err) } else { m.status = api.HOST_HEALTH_STATUS_RUNNING - log.Infof("etcd client restart session success") + log.Infof("etcd client restart session put %s success", m.GetKey()) return } } log.Errorf("keep etcd lease failed: %s", err) - if m.networkAvailable(nicRecord) { + if m.networkAvailable() { log.Infof("network is available, try reconnect") // may be etcd not work m.Reconnect() @@ -141,70 +150,45 @@ func (m *SHostHealthManager) OnKeepaliveFailure() { } } -func (m *SHostHealthManager) recordNic() map[string]int { - nicRecord := make(map[string]int) - for _, n := range options.HostOptions.Networks { - data := strings.Split(n, "/") - interf := data[0] - rx, err := fileutils2.FileGetContents( - fmt.Sprintf("/sys/class/net/%s/statistics/rx_bytes", interf), - ) - if err != nil { - log.Errorf("failed get nic rx %s statistics %s", interf, err) - continue - } - tx, err := fileutils2.FileGetContents( - fmt.Sprintf("/sys/class/net/%s/statistics/tx_bytes", interf), - ) - if err != nil { - log.Errorf("failed get nic tx %s statistics %s", interf, err) - continue - } - irx, err := strconv.Atoi(strings.TrimSpace(rx)) - if err != nil { - log.Errorf("failed convert rx %s %s", rx, err) - } - itx, err := strconv.Atoi(strings.TrimSpace(tx)) - if err != nil { - log.Errorf("failed convert tx %s %s", tx, err) - } - nicRecord[interf] = irx + itx +func (m *SHostHealthManager) networkAvailable() bool { + res, err := misc.Ping(m.masterNodesIps, 3, 10, true) + if err != nil { + log.Errorf("failed ping master nodes %s", res) + return true } - return nicRecord -} - -func (m *SHostHealthManager) networkAvailable(oldRecord map[string]int) bool { - newRecord := m.recordNic() - for _, n := range options.HostOptions.Networks { - data := strings.Split(n, "/") - interf := data[0] - - oldR, ok := oldRecord[interf] - if !ok { - continue - } - newR, ok := newRecord[interf] - if !ok { - log.Errorf("nic %s record not found", n) - continue - } - - if newR != oldR { + for _, v := range res { + if v.Loss() < 100 { return true } } return false } +func (m *SHostHealthManager) masterNodesInternalIps() ([]string, error) { + result, err := modules.Hosts.Get(hostutils.GetComputeSession(context.Background()), "k8s-master-node-ips", nil) + if err != nil { + return nil, err + } + ips := make([]string, 0) + err = result.Unmarshal(&ips, "ips") + if err != nil { + return nil, errors.Wrap(err, "unmarshal master node ips") + } + return ips, nil +} + func (m *SHostHealthManager) OnUnhealth() { - if m.onHostDown == hostconsts.SHUTDOWN_SERVERS { - log.Errorf("Host unhealthy, going to shotdown servers") - m.shutdownServers() + p := path.Join(options.HostOptions.ServersPath, hostconsts.HOST_HEALTH_FILENAME) + if fileutils2.Exists(p) { + if act, err := fileutils2.FileGetContents(p); err != nil { + log.Errorf(" failed read file %s: %s", p, err) + } else if act == hostconsts.SHUTDOWN_SERVERS { + log.Errorf("Host unhealthy, going to shutdown servers") + m.shutdownServers() + } } // reconnect wait for network available m.Reconnect() - utils.DumpAllGoroutineStack(log.Logger().Out) - os.Exit(1) } func (m *SHostHealthManager) Reconnect() { @@ -223,32 +207,34 @@ func (m *SHostHealthManager) Reconnect() { log.Infof("restart ression success") if err := m.cli.PutSession( - context.Background(), fmt.Sprintf("%s/%s", api.HOST_HEALTH_PREFIX, m.hostId), - api.HOST_HEALTH_STATUS_RUNNING, + context.Background(), m.GetKey(), api.HOST_HEALTH_STATUS_RUNNING, ); err != nil { log.Errorf("put host key failed %s", err) go m.Reconnect() return } - log.Infof("put key %s/%s success", api.HOST_HEALTH_PREFIX, m.hostId) + log.Infof("put key %s success", m.GetKey()) m.status = api.HOST_HEALTH_STATUS_RUNNING } -func (m *SHostHealthManager) SetOnHostDown(onHostDown string) { - m.onHostDown = onHostDown -} - -// shutdown servers used shared storage func (m *SHostHealthManager) shutdownServers() { - types.HealthCheckReactor.ShutdownServers() -} - -func SetOnHostDown(onHostDown string) error { - if manager != nil { - manager.SetOnHostDown(onHostDown) - return nil + files, err := ioutil.ReadDir(options.HostOptions.ServersPath) + if err != nil { + log.Errorf("failed walk dir %s: %s", options.HostOptions.ServersPath, err) + return + } + for i := range files { + if hostutils.IsGuestDir(files[i], options.HostOptions.ServersPath) { + stopvm := path.Join(options.HostOptions.ServersPath, files[i].Name(), "stopvm") + if fileutils2.Exists(stopvm) { + log.Infof("start exec stopvm script for guest %s", files[i].Name()) + out, err := procutils.NewRemoteCommandAsFarAsPossible("bash", stopvm, "--force").Output() + if err != nil { + log.Errorf("failed exec stopvm script for guest %s: %s %s", files[i].Name(), out, err) + } + } + } } - return fmt.Errorf("host health manager not init") } func GetHealthStatus() string { diff --git a/pkg/hostman/host_services.go b/pkg/hostman/host_services.go index 2eddf35a33..3d49697309 100644 --- a/pkg/hostman/host_services.go +++ b/pkg/hostman/host_services.go @@ -15,13 +15,8 @@ package hostman import ( - "io/ioutil" - "os" - "path/filepath" - execlient "yunion.io/x/executor/client" "yunion.io/x/log" - "yunion.io/x/pkg/errors" "yunion.io/x/onecloud/pkg/appsrv" app_common "yunion.io/x/onecloud/pkg/cloudcommon/app" @@ -31,7 +26,6 @@ import ( "yunion.io/x/onecloud/pkg/hostman/guestman" "yunion.io/x/onecloud/pkg/hostman/guestman/desc" "yunion.io/x/onecloud/pkg/hostman/guestman/guesthandlers" - "yunion.io/x/onecloud/pkg/hostman/host_health" "yunion.io/x/onecloud/pkg/hostman/hostdeployer/deployclient" "yunion.io/x/onecloud/pkg/hostman/hosthandler" "yunion.io/x/onecloud/pkg/hostman/hostinfo" @@ -43,7 +37,6 @@ import ( "yunion.io/x/onecloud/pkg/hostman/storageman" "yunion.io/x/onecloud/pkg/hostman/storageman/diskhandlers" "yunion.io/x/onecloud/pkg/hostman/storageman/storagehandler" - "yunion.io/x/onecloud/pkg/httperrors" "yunion.io/x/onecloud/pkg/util/procutils" "yunion.io/x/onecloud/pkg/util/sysutils" ) @@ -77,11 +70,6 @@ func (host *SHostService) InitService() { func (host *SHostService) OnExitService() {} func (host *SHostService) RunService() { - hn, err := os.Hostname() - if err != nil { - log.Fatalf("fail to get hostname %s", err) - } - app := app_common.InitApp(&options.HostOptions.BaseOptions, false) cronManager := cronman.InitCronJobManager(false, options.HostOptions.CronJobWorkerCount) hostutils.Init() @@ -93,17 +81,6 @@ func (host *SHostService) RunService() { app_common.InitAuth(&options.HostOptions.CommonOptions, func() { log.Infof("Auth complete!!") - - if err := host.initEtcdConfig(); err != nil { - log.Fatalln("Init etcd config:", err) - } - - if len(options.HostOptions.EtcdEndpoints) > 0 { - _, err := host_health.InitHostHealthManager(hn, "") - if err != nil { - log.Fatalf("Init host health manager failed %s", err) - } - } }) deployclient.Init(options.HostOptions.DeployServerSocketPath) @@ -166,47 +143,6 @@ func (host *SHostService) initHandlers(app *appsrv.Application) { app_common.ExportOptionsHandler(app, &options.HostOptions) } -func (host *SHostService) initEtcdConfig() error { - etcdEndpoint, err := app_common.FetchEtcdServiceInfo() - if err != nil { - if errors.Cause(err) == httperrors.ErrNotFound { - return nil - } - return errors.Wrap(err, "fetch etcd service info") - } - if etcdEndpoint == nil { - return nil - } - if len(options.HostOptions.EtcdEndpoints) == 0 { - options.HostOptions.EtcdEndpoints = []string{etcdEndpoint.Url} - if len(etcdEndpoint.CertId) > 0 { - dir, err := ioutil.TempDir("", "etcd-cluster-tls") - if err != nil { - return errors.Wrap(err, "create dir etcd cluster tls") - } - options.HostOptions.EtcdCert, err = writeFile(dir, "etcd.crt", []byte(etcdEndpoint.Certificate)) - if err != nil { - return errors.Wrap(err, "write file certificate") - } - options.HostOptions.EtcdKey, err = writeFile(dir, "etcd.key", []byte(etcdEndpoint.PrivateKey)) - if err != nil { - return errors.Wrap(err, "write file private key") - } - options.HostOptions.EtcdCacert, err = writeFile(dir, "etcd-ca.crt", []byte(etcdEndpoint.CaCertificate)) - if err != nil { - return errors.Wrap(err, "write file cacert") - } - options.HostOptions.EtcdUseTLS = true - } - } - return nil -} - -func writeFile(dir, file string, data []byte) (string, error) { - p := filepath.Join(dir, file) - return p, ioutil.WriteFile(p, data, 0600) -} - func StartService() { var srv = &SHostService{} srv.SServiceBase = &service.SServiceBase{ diff --git a/pkg/hostman/hosthandler/handler.go b/pkg/hostman/hosthandler/handler.go index 180d4136df..ed4c1ecef4 100644 --- a/pkg/hostman/hosthandler/handler.go +++ b/pkg/hostman/hosthandler/handler.go @@ -22,10 +22,10 @@ import ( "yunion.io/x/jsonutils" "yunion.io/x/onecloud/pkg/appsrv" - "yunion.io/x/onecloud/pkg/hostman/host_health" "yunion.io/x/onecloud/pkg/hostman/hostinfo" "yunion.io/x/onecloud/pkg/hostman/hostinfo/hostconsts" "yunion.io/x/onecloud/pkg/hostman/hostutils" + "yunion.io/x/onecloud/pkg/httperrors" "yunion.io/x/onecloud/pkg/mcclient/auth" ) @@ -37,14 +37,10 @@ var ( func AddHostHandler(prefix string, app *appsrv.Application) { for _, keyword := range keyWords { - app.AddHandler("POST", fmt.Sprintf("%s/%s/shutdown-servers-on-host-down", prefix, keyword), - auth.Authenticate(setOnHostDown)) - app.AddHandler("GET", fmt.Sprintf("%s/%s/health-status", prefix, keyword), - auth.Authenticate(getHealthManagerStatus)) - for action, f := range map[string]actionFunc{ - "sync": hostSync, - "probe-isolated-devices": hostProbeIsolatedDevices, + "sync": hostSync, + "probe-isolated-devices": hostProbeIsolatedDevices, + "shutdown-servers-on-host-down": setOnHostDown, } { app.AddHandler("POST", fmt.Sprintf("%s/%s//%s", prefix, keyword, action), @@ -54,17 +50,17 @@ func AddHostHandler(prefix string, app *appsrv.Application) { } } -func setOnHostDown(ctx context.Context, w http.ResponseWriter, r *http.Request) { - if err := host_health.SetOnHostDown(hostconsts.SHUTDOWN_SERVERS); err != nil { - hostutils.Response(ctx, w, err) - return +func setOnHostDown(ctx context.Context, hostId string, body jsonutils.JSONObject) (interface{}, error) { + if !body.Contains("shutdown_servers") { + return nil, httperrors.NewMissingParameterError("shutdown_servers") } - hostutils.ResponseOk(ctx, w) -} -func getHealthManagerStatus(ctx context.Context, w http.ResponseWriter, r *http.Request) { - status := host_health.GetHealthStatus() - hostutils.Response(ctx, w, map[string]string{"status": status}) + if jsonutils.QueryBoolean(body, "shutdown_servers", false) { + hostinfo.Instance().SetOnHostDown(hostconsts.SHUTDOWN_SERVERS) + } else { + hostinfo.Instance().SetOnHostDown("") + } + return nil, nil } func hostActions(f actionFunc) appsrv.FilterHandler { diff --git a/pkg/hostman/hostinfo/hostconsts/hostconsts.go b/pkg/hostman/hostinfo/hostconsts/hostconsts.go index 98bc1a3d14..de48cc515a 100644 --- a/pkg/hostman/hostinfo/hostconsts/hostconsts.go +++ b/pkg/hostman/hostinfo/hostconsts/hostconsts.go @@ -28,7 +28,8 @@ const ( TELEGRAF_TAG_ONECLOUD_HOST_TYPE_CONTROLLER = "controller" TELEGRAF_TAG_ONECLOUD_HOST_TYPE_LBAGENT = "lbagent" - SHUTDOWN_SERVERS = "shutdown-servers" + SHUTDOWN_SERVERS = "shutdown-servers" + HOST_HEALTH_FILENAME = ".host-health" HOST_CGROUP = "cloudpods.hostagent" HOST_RESERVED_CPUSET = "cloudpods.hostagent.reserved" diff --git a/pkg/hostman/hostinfo/hostinfo.go b/pkg/hostman/hostinfo/hostinfo.go index 499b6bd103..2bb6834eda 100644 --- a/pkg/hostman/hostinfo/hostinfo.go +++ b/pkg/hostman/hostinfo/hostinfo.go @@ -46,7 +46,6 @@ import ( "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/cloudcommon/types" "yunion.io/x/onecloud/pkg/hostman/guestfs/fsdriver" - "yunion.io/x/onecloud/pkg/hostman/host_health" "yunion.io/x/onecloud/pkg/hostman/hostinfo/hostbridge" "yunion.io/x/onecloud/pkg/hostman/hostinfo/hostconsts" "yunion.io/x/onecloud/pkg/hostman/hostutils" @@ -1395,6 +1394,11 @@ func (h *SHostInfo) updateHostMetadata(hostname string) error { // return err // } +func (h *SHostInfo) SetOnHostDown(action string) error { + h.onHostDown = action + return fileutils2.FilePutContents(path.Join(options.HostOptions.ServersPath, hostconsts.HOST_HEALTH_FILENAME), h.onHostDown, false) +} + func (h *SHostInfo) onUpdateHostInfoSucc(hostbody jsonutils.JSONObject) { h.HostId, _ = hostbody.GetString("id") hostname, _ := hostbody.GetString("name") @@ -1403,10 +1407,9 @@ func (h *SHostInfo) onUpdateHostInfoSucc(hostbody jsonutils.JSONObject) { return } if jsonutils.QueryBoolean(hostbody, "auto_migrate_on_host_down", false) { - h.onHostDown = hostconsts.SHUTDOWN_SERVERS + h.SetOnHostDown(hostconsts.SHUTDOWN_SERVERS) } - log.Infof("on host down %s", h.onHostDown) - host_health.SetOnHostDown(h.onHostDown) + log.Infof("host health manager on host down %s", h.onHostDown) // fetch host reserved cpus info reservedCpusStr, _ := hostbody.GetString("metadata", api.HOSTMETA_RESERVED_CPUS_INFO) @@ -2046,10 +2049,8 @@ func (h *SHostInfo) stop() { func (h *SHostInfo) unregister() { isLog := false for { - updateHealthStatus := true input := api.HostOfflineInput{ - UpdateHealthStatus: &updateHealthStatus, - Reason: "host stop", + Reason: "host stop", } _, err := modules.Hosts.PerformAction(h.GetSession(), h.HostId, "offline", jsonutils.Marshal(input)) if err != nil { diff --git a/pkg/hostman/hostutils/hostutils.go b/pkg/hostman/hostutils/hostutils.go index 87d09b0bfa..b64ad30c85 100644 --- a/pkg/hostman/hostutils/hostutils.go +++ b/pkg/hostman/hostutils/hostutils.go @@ -18,10 +18,13 @@ import ( "context" "fmt" "net/http" + "os" + "path" "yunion.io/x/jsonutils" "yunion.io/x/log" "yunion.io/x/pkg/appctx" + "yunion.io/x/pkg/util/regutils" "yunion.io/x/onecloud/pkg/apis" hostapi "yunion.io/x/onecloud/pkg/apis/host" @@ -38,6 +41,7 @@ import ( modules "yunion.io/x/onecloud/pkg/mcclient/modules/compute" "yunion.io/x/onecloud/pkg/mcclient/modules/k8s" "yunion.io/x/onecloud/pkg/util/cgrouputils/cpuset" + "yunion.io/x/onecloud/pkg/util/fileutils2" ) type IHost interface { @@ -159,6 +163,20 @@ func UpdateServerProgress(ctx context.Context, sid string, progress, progressMbp return modules.Servers.Update(GetComputeSession(ctx), sid, jsonutils.Marshal(params)) } +func IsGuestDir(f os.FileInfo, serversPath string) bool { + if !regutils.MatchUUID(f.Name()) { + return false + } + if !f.Mode().IsDir() && f.Mode()&os.ModeSymlink == 0 { + return false + } + descFile := path.Join(serversPath, f.Name(), "desc") + if !fileutils2.Exists(descFile) { + return false + } + return true +} + func ResponseOk(ctx context.Context, w http.ResponseWriter) { Response(ctx, w, map[string]string{"result": "ok"}) }