Merge pull request #16217 from wanyaoqi/automated-cherry-pick-of-#16189-upstream-release-3.10

Automated cherry pick of #16189: fix(region,host): host health misc fix
This commit is contained in:
Zexi Li
2023-03-14 12:33:45 +08:00
committed by GitHub
13 changed files with 310 additions and 215 deletions
+3
View File
@@ -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
+128
View File
@@ -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()
}
+5 -1
View File
@@ -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")
}
+4 -6
View File
@@ -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)
+20 -12
View File
@@ -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)
}
+33 -5
View File
@@ -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)
+1 -13
View File
@@ -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 {
+75 -89
View File
@@ -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 {
-64
View File
@@ -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{
+13 -17
View File
@@ -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/<sid>/%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 {
@@ -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"
+8 -7
View File
@@ -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 {
+18
View File
@@ -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"})
}