node_record.go 2.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150
  1. package analytic
  2. import (
  3. "context"
  4. "net/http"
  5. "sync"
  6. "time"
  7. "github.com/0xJacky/Nginx-UI/model"
  8. "github.com/0xJacky/Nginx-UI/query"
  9. "github.com/gorilla/websocket"
  10. "github.com/uozi-tech/cosy/logger"
  11. )
  12. var (
  13. ctx, cancel = context.WithCancel(context.Background())
  14. wg sync.WaitGroup
  15. restartMu sync.Mutex // Add mutex to prevent concurrent restarts
  16. )
  17. func RestartRetrieveNodesStatus() {
  18. restartMu.Lock() // Acquire lock before modifying shared resources
  19. defer restartMu.Unlock()
  20. // Cancel previous context to stop all operations
  21. cancel()
  22. // Wait for previous goroutines to finish
  23. wg.Wait()
  24. // Create new context for this run
  25. ctx, cancel = context.WithCancel(context.Background())
  26. wg.Add(1)
  27. go func() {
  28. defer wg.Done()
  29. RetrieveNodesStatus()
  30. }()
  31. }
  32. func RetrieveNodesStatus() {
  33. logger.Info("RetrieveNodesStatus start")
  34. defer logger.Info("RetrieveNodesStatus exited")
  35. mutex.Lock()
  36. if NodeMap == nil {
  37. NodeMap = make(TNodeMap)
  38. }
  39. mutex.Unlock()
  40. env := query.Environment
  41. envs, err := env.Where(env.Enabled.Is(true)).Find()
  42. if err != nil {
  43. logger.Error(err)
  44. return
  45. }
  46. var wg sync.WaitGroup
  47. defer wg.Wait()
  48. for _, env := range envs {
  49. wg.Add(1)
  50. go func(e *model.Environment) {
  51. defer wg.Done()
  52. retryTicker := time.NewTicker(5 * time.Second)
  53. defer retryTicker.Stop()
  54. for {
  55. select {
  56. case <-ctx.Done():
  57. return
  58. default:
  59. if err := nodeAnalyticRecord(e, ctx); err != nil {
  60. logger.Error(err)
  61. if NodeMap[env.ID] != nil {
  62. mutex.Lock()
  63. NodeMap[env.ID].Status = false
  64. mutex.Unlock()
  65. }
  66. select {
  67. case <-retryTicker.C:
  68. case <-ctx.Done():
  69. return
  70. }
  71. }
  72. }
  73. }
  74. }(env)
  75. }
  76. <-ctx.Done()
  77. }
  78. func nodeAnalyticRecord(env *model.Environment, ctx context.Context) error {
  79. scopeCtx, cancel := context.WithCancel(ctx)
  80. defer cancel()
  81. node, err := InitNode(env)
  82. mutex.Lock()
  83. NodeMap[env.ID] = node
  84. mutex.Unlock()
  85. if err != nil {
  86. return err
  87. }
  88. u, err := env.GetWebSocketURL("/api/analytic/intro")
  89. if err != nil {
  90. return err
  91. }
  92. header := http.Header{}
  93. header.Set("X-Node-Secret", env.Token)
  94. dial := &websocket.Dialer{
  95. Proxy: http.ProxyFromEnvironment,
  96. HandshakeTimeout: 5 * time.Second,
  97. }
  98. c, _, err := dial.Dial(u, header)
  99. if err != nil {
  100. return err
  101. }
  102. defer c.Close()
  103. go func() {
  104. <-scopeCtx.Done()
  105. _ = c.Close()
  106. }()
  107. var nodeStat NodeStat
  108. for {
  109. err = c.ReadJSON(&nodeStat)
  110. if err != nil {
  111. return err
  112. }
  113. // set online
  114. nodeStat.Status = true
  115. nodeStat.ResponseAt = time.Now()
  116. mutex.Lock()
  117. NodeMap[env.ID].NodeStat = nodeStat
  118. mutex.Unlock()
  119. }
  120. }