node_heartbeat_job.go 4.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172
  1. package job
  2. import (
  3. "context"
  4. "fmt"
  5. "sort"
  6. "strconv"
  7. "strings"
  8. "sync"
  9. "time"
  10. "github.com/mhsanaei/3x-ui/v3/internal/database/model"
  11. "github.com/mhsanaei/3x-ui/v3/internal/eventbus"
  12. "github.com/mhsanaei/3x-ui/v3/internal/logger"
  13. "github.com/mhsanaei/3x-ui/v3/internal/util/common"
  14. "github.com/mhsanaei/3x-ui/v3/internal/web/service"
  15. "github.com/mhsanaei/3x-ui/v3/internal/web/websocket"
  16. )
  17. const (
  18. nodeHeartbeatConcurrency = 32
  19. nodeHeartbeatRequestTimeout = 4 * time.Second
  20. // Past this many same-direction transitions in one tick, one summary event goes out:
  21. // per-node events overflow the notifier queues and every chat's rate limit.
  22. nodeTransitionBurst = 5
  23. nodeTransitionBurstNames = 10
  24. )
  25. type NodeHeartbeatJob struct {
  26. nodeService service.NodeService
  27. running sync.Mutex
  28. }
  29. func NewNodeHeartbeatJob() *NodeHeartbeatJob {
  30. return &NodeHeartbeatJob{}
  31. }
  32. func (j *NodeHeartbeatJob) Run() {
  33. if !j.running.TryLock() {
  34. return
  35. }
  36. defer j.running.Unlock()
  37. nodes, err := j.nodeService.GetAll()
  38. if err != nil {
  39. logger.Warning("node heartbeat: load nodes failed:", err)
  40. return
  41. }
  42. j.nodeService.RetainEnabledNodeDescendants(nodes)
  43. if len(nodes) == 0 {
  44. return
  45. }
  46. sem := make(chan struct{}, nodeHeartbeatConcurrency)
  47. var wg sync.WaitGroup
  48. var transitionsMu sync.Mutex
  49. var transitions []eventbus.Event
  50. for _, n := range nodes {
  51. if !n.Enable {
  52. continue
  53. }
  54. wg.Add(1)
  55. sem <- struct{}{}
  56. n := n
  57. common.GoRecover("node-heartbeat:"+n.Name, func() {
  58. defer wg.Done()
  59. defer func() { <-sem }()
  60. if event := j.probeOne(n); event != nil {
  61. transitionsMu.Lock()
  62. transitions = append(transitions, *event)
  63. transitionsMu.Unlock()
  64. }
  65. })
  66. }
  67. wg.Wait()
  68. publishNodeTransitions(transitions)
  69. if !websocket.HasClients() {
  70. return
  71. }
  72. updated, err := j.nodeService.GetNodeTreeView()
  73. if err != nil {
  74. logger.Warning("node heartbeat: load nodes for broadcast failed:", err)
  75. return
  76. }
  77. websocket.BroadcastNodes(updated)
  78. }
  79. func (j *NodeHeartbeatJob) probeOne(n *model.Node) *eventbus.Event {
  80. ctx, cancel := context.WithTimeout(context.Background(), nodeHeartbeatRequestTimeout)
  81. defer cancel()
  82. prevStatus := n.Status
  83. patch, err := j.nodeService.Probe(ctx, n)
  84. if err != nil {
  85. patch.Status = "offline"
  86. } else {
  87. patch.Status = "online"
  88. }
  89. if updErr := j.nodeService.UpdateHeartbeat(n.Id, patch); updErr != nil {
  90. logger.Warning("node heartbeat: update node", n.Id, "failed:", updErr)
  91. }
  92. // Learn the nodes this node manages so the panel can surface them as
  93. // transitive sub-nodes (#4983). Fresh context — the probe budget above may
  94. // be spent. Drop them when the node is unreachable.
  95. if patch.Status == "online" {
  96. dctx, dcancel := context.WithTimeout(context.Background(), nodeHeartbeatRequestTimeout)
  97. j.nodeService.RefreshDescendants(dctx, n)
  98. dcancel()
  99. } else {
  100. j.nodeService.ClearDescendants(n.Id)
  101. }
  102. return nodeTransitionEvent(n, prevStatus, patch)
  103. }
  104. // nodeTransitionEvent is node.down / node.up on a genuine state change only; an unknown
  105. // previous status (fresh start) counts as not-online, so it never yields node.down.
  106. func nodeTransitionEvent(n *model.Node, prevStatus string, patch service.HeartbeatPatch) *eventbus.Event {
  107. var eventType eventbus.EventType
  108. switch {
  109. case prevStatus == "online" && patch.Status == "offline":
  110. eventType = eventbus.EventNodeDown
  111. case prevStatus != "online" && patch.Status == "online":
  112. eventType = eventbus.EventNodeUp
  113. default:
  114. return nil
  115. }
  116. source := n.Name
  117. if source == "" {
  118. source = "node-" + strconv.Itoa(n.Id)
  119. }
  120. return &eventbus.Event{
  121. Type: eventType,
  122. Source: source,
  123. Data: &eventbus.NodeHealthData{
  124. NodeId: n.Id,
  125. LatencyMs: patch.LatencyMs,
  126. CpuPct: patch.CpuPct,
  127. MemPct: patch.MemPct,
  128. XrayState: patch.XrayState,
  129. XrayError: patch.XrayError,
  130. },
  131. }
  132. }
  133. // publishNodeTransitions sends one tick's transitions, folding a same-direction burst
  134. // (a master-side blip flips every node at once) into one event naming the nodes.
  135. func publishNodeTransitions(events []eventbus.Event) {
  136. if EventBus == nil {
  137. return
  138. }
  139. namesByType := make(map[eventbus.EventType][]string)
  140. for _, e := range events {
  141. namesByType[e.Type] = append(namesByType[e.Type], e.Source)
  142. }
  143. for _, e := range events {
  144. if len(namesByType[e.Type]) <= nodeTransitionBurst {
  145. EventBus.Publish(e)
  146. }
  147. }
  148. for _, eventType := range []eventbus.EventType{eventbus.EventNodeDown, eventbus.EventNodeUp} {
  149. names := namesByType[eventType]
  150. if len(names) <= nodeTransitionBurst {
  151. continue
  152. }
  153. sort.Strings(names)
  154. source := strings.Join(names[:min(len(names), nodeTransitionBurstNames)], ", ")
  155. if extra := len(names) - nodeTransitionBurstNames; extra > 0 {
  156. source += fmt.Sprintf(" (+%d)", extra)
  157. }
  158. EventBus.Publish(eventbus.Event{Type: eventType, Source: source})
  159. }
  160. }