This is an automated email from the ASF dual-hosted git repository.

zhongxjian pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/dubbo-kubernetes.git


The following commit(s) were added to refs/heads/master by this push:
     new 3d48f435 [horus] Node downtime logic implementation (#358)
3d48f435 is described below

commit 3d48f435ab24b583e703d10e366fa44af0fa01fb
Author: mfordjody <[email protected]>
AuthorDate: Sat Sep 21 15:11:15 2024 +0800

    [horus] Node downtime logic implementation (#358)
---
 app/horus/basic/config/file.go     |  1 +
 app/horus/cmd/main.go              | 10 ++++++
 app/horus/core/horuser/downtime.go | 66 ++++++++++++++++++++++++++++++++++++++
 3 files changed, 77 insertions(+)

diff --git a/app/horus/basic/config/file.go b/app/horus/basic/config/file.go
index 5c319db1..d5c52fc6 100644
--- a/app/horus/basic/config/file.go
+++ b/app/horus/basic/config/file.go
@@ -71,5 +71,6 @@ type DowntimeConfiguration struct {
        PromQueryTimeSecond int64                  `yaml:"promQueryTimeSecond"`
        KubeMultiple        map[string]string      `yaml:"kubeMultiple"`
        CheckQL             []string               `yaml:"checkQL"`
+       NodeNameToIPs       string                 `yaml:"nodeNameToIPs"`
        DingTalk            *DingTalkConfiguration `yaml:"dingTalk"`
 }
diff --git a/app/horus/cmd/main.go b/app/horus/cmd/main.go
index cfdcac29..aed1bc02 100644
--- a/app/horus/cmd/main.go
+++ b/app/horus/cmd/main.go
@@ -117,6 +117,16 @@ func main() {
                }
                return nil
        })
+       group.Add(func() error {
+               if c.CustomModular.Enabled {
+                       klog.Info("horus down time manager start success.")
+                       err := horus.DownTimeManager(ctx)
+                       if err != nil {
+                               klog.Errorf("horus down time manager start 
failed error:%v", err)
+                       }
+               }
+               return nil
+       })
        group.Wait()
 }
 
diff --git a/app/horus/core/horuser/downtime.go 
b/app/horus/core/horuser/downtime.go
index 1df32c0a..00c67d39 100644
--- a/app/horus/core/horuser/downtime.go
+++ b/app/horus/core/horuser/downtime.go
@@ -17,12 +17,22 @@ package horuser
 
 import (
        "context"
+       "fmt"
+       "github.com/apache/dubbo-kubernetes/app/horus/basic/db"
+       "github.com/apache/dubbo-kubernetes/app/horus/core/alert"
        "k8s.io/apimachinery/pkg/util/wait"
        "k8s.io/klog"
        "sync"
        "time"
 )
 
+const (
+       NODE_DOWN        = "node_down"
+       NODE_DOWN_REASON = "node_down_aegis"
+       POWER_OFF        = "POWER_OFF"
+       POWER_ON         = "POWER_ON"
+)
+
 func (h *Horuser) DownTimeManager(ctx context.Context) error {
        go wait.UntilWithContext(ctx, h.DownTimeCheck, 
time.Duration(h.cc.NodeDownTime.CheckIntervalSecond)*time.Second)
        <-ctx.Done()
@@ -80,4 +90,60 @@ func (h *Horuser) DownTimeNodes(clusterName, addr string) {
                        continue
                }
        }
+
+       WithDownNodeIPs := map[string]string{}
+       for node, count := range resMap {
+               if count < checkQl {
+                       klog.Errorf("downtimeNodes node not reach threshold")
+                       klog.Infof("clusterName:%v nodeName:%v threshold:%v 
count:%v", clusterName, node, checkQl, count)
+                       continue
+               }
+               toNodeNameips := fmt.Sprintf(h.cc.NodeDownTime.NodeNameToIPs, 
node)
+               res, err := h.InstantQuery(toNodeNameips, addr, clusterName, 
h.cc.NodeDownTime.PromQueryTimeSecond)
+               if err != nil {
+                       klog.Errorf("downtimeNodes InstantQuery NodeName To IPs 
empty err:%v", err)
+                       klog.Infof("clusterName:%v toNodeNameips:%v err:%v", 
clusterName, toNodeNameips, err)
+                       continue
+               }
+               str := ""
+               for _, v := range res {
+                       v := v
+                       str = string(v.Metric["instance"])
+               }
+               WithDownNodeIPs[node] = str
+       }
+       WithDownNodeIPsMsg := fmt.Sprintf("【%s】\n【集群:%v】\n【宕机:%v】\n", 
h.cc.NodeDownTime.DingTalk.Title, clusterName, len(WithDownNodeIPs))
+       newfound := 0
+       for nodeName, nodeIP := range WithDownNodeIPs {
+               today := time.Now().Format("2006-01-02")
+               err := h.Cordon(nodeName, clusterName, NODE_DOWN)
+               if err != nil {
+                       klog.Errorf("Cordon node err:%v", err)
+                       klog.Infof("clusterName:%v nodeName:%v", clusterName, 
nodeName)
+               }
+               write := db.NodeDataInfo{
+                       NodeName:    nodeName,
+                       NodeIP:      nodeIP,
+                       ClusterName: clusterName,
+                       ModuleName:  NODE_DOWN,
+               }
+               exist, _ := write.Check()
+               if exist {
+                       continue
+               }
+               newfound++
+               WithDownNodeIPsMsg += fmt.Sprintf("node:%v ip:%v", nodeName, 
nodeIP)
+               write.Reason = NODE_DOWN_REASON
+               write.FirstDate = today
+               _, err = write.Add()
+               if err != nil {
+                       klog.Errorf("NodeDownTimeCheckOnCluster abnormal 
cordonNode AddOrGetOne err:%v", err)
+                       klog.Infof("cluster:%v node:%v", clusterName, nodeName)
+               }
+               klog.Infof("NodeDownTimeCheckOnCluster abnormal cordonNode 
AddOrGetOne cluster:%v node:%v", clusterName, nodeName)
+               if newfound > 0 {
+                       klog.Infof("NodeDownTimeCheckOnCluster get 
toNodeNameips result msg:%v clusterName:%v count:%v detail:%v", 
WithDownNodeIPsMsg, clusterName, len(nodeIP), nodeName)
+                       alert.DingTalkSend(h.cc.NodeDownTime.DingTalk, 
WithDownNodeIPsMsg)
+               }
+       }
 }

Reply via email to