Files
EdgeNode/internal/nodes/task_sync_api_nodes.go

123 lines
2.5 KiB
Go
Raw Normal View History

2021-02-24 11:01:06 +08:00
package nodes
import (
"github.com/TeaOSLab/EdgeCommon/pkg/rpc/pb"
"github.com/TeaOSLab/EdgeNode/internal/configs"
"github.com/TeaOSLab/EdgeNode/internal/events"
"github.com/TeaOSLab/EdgeNode/internal/goman"
2021-02-24 11:01:06 +08:00
"github.com/TeaOSLab/EdgeNode/internal/rpc"
"github.com/TeaOSLab/EdgeNode/internal/trackers"
"github.com/TeaOSLab/EdgeNode/internal/utils"
2021-02-24 11:01:06 +08:00
"github.com/iwind/TeaGo/Tea"
"github.com/iwind/TeaGo/logs"
"time"
)
2022-01-12 20:31:04 +08:00
var sharedSyncAPINodesTask = NewSyncAPINodesTask()
2021-02-24 11:01:06 +08:00
func init() {
events.On(events.EventStart, func() {
goman.New(func() {
2022-01-12 20:31:04 +08:00
sharedSyncAPINodesTask.Start()
})
2021-02-24 11:01:06 +08:00
})
2022-01-12 20:31:04 +08:00
events.On(events.EventQuit, func() {
sharedSyncAPINodesTask.Stop()
})
2021-02-24 11:01:06 +08:00
}
// SyncAPINodesTask API节点同步任务
2021-02-24 11:01:06 +08:00
type SyncAPINodesTask struct {
2022-01-12 20:31:04 +08:00
ticker *time.Ticker
2021-02-24 11:01:06 +08:00
}
func NewSyncAPINodesTask() *SyncAPINodesTask {
return &SyncAPINodesTask{}
}
func (this *SyncAPINodesTask) Start() {
2022-01-12 20:31:04 +08:00
this.ticker = time.NewTicker(5 * time.Minute)
2021-02-24 11:01:06 +08:00
if Tea.IsTesting() {
// 快速测试
2022-01-12 20:31:04 +08:00
this.ticker = time.NewTicker(1 * time.Minute)
2021-02-24 11:01:06 +08:00
}
2022-01-12 20:31:04 +08:00
for range this.ticker.C {
2021-02-24 11:01:06 +08:00
err := this.Loop()
if err != nil {
logs.Println("[TASK][SYNC_API_NODES_TASK]" + err.Error())
}
}
}
2022-01-12 20:31:04 +08:00
func (this *SyncAPINodesTask) Stop() {
if this.ticker != nil {
this.ticker.Stop()
}
}
2021-02-24 11:01:06 +08:00
func (this *SyncAPINodesTask) Loop() error {
// 如果有节点定制的API节点地址
var hasCustomizedAPINodeAddrs = sharedNodeConfig != nil && len(sharedNodeConfig.APINodeAddrs) > 0
config, err := configs.LoadAPIConfig()
if err != nil {
return err
}
// 是否禁止自动升级
if config.RPC.DisableUpdate {
return nil
}
var tr = trackers.Begin("SYNC_API_NODES")
defer tr.End()
2021-02-24 11:01:06 +08:00
// 获取所有可用的节点
rpcClient, err := rpc.SharedRPC()
if err != nil {
return err
}
2022-08-24 20:04:46 +08:00
resp, err := rpcClient.APINodeRPC.FindAllEnabledAPINodes(rpcClient.Context(), &pb.FindAllEnabledAPINodesRequest{})
2021-02-24 11:01:06 +08:00
if err != nil {
return err
}
var newEndpoints = []string{}
2021-11-05 15:37:07 +08:00
for _, node := range resp.ApiNodes {
2021-02-24 11:01:06 +08:00
if !node.IsOn {
continue
}
newEndpoints = append(newEndpoints, node.AccessAddrs...)
}
// 和现有的对比
if utils.EqualStrings(newEndpoints, config.RPC.Endpoints) {
2021-02-24 11:01:06 +08:00
return nil
}
// 测试是否有API节点可用
var hasOk = rpcClient.TestEndpoints(newEndpoints)
if !hasOk {
return nil
}
2021-02-24 11:01:06 +08:00
// 修改RPC对象配置
config.RPC.Endpoints = newEndpoints
// 更新当前RPC
if !hasCustomizedAPINodeAddrs {
err = rpcClient.UpdateConfig(config)
if err != nil {
return err
}
2021-02-24 11:01:06 +08:00
}
// 保存到文件
err = config.WriteFile(Tea.ConfigFile("api.yaml"))
if err != nil {
return err
}
return nil
}