nezha/cmd/dashboard/rpc/rpc.go

137 lines
4.3 KiB
Go
Raw Normal View History

2019-12-08 03:59:58 -05:00
package rpc
import (
2024-10-23 08:37:29 -04:00
"fmt"
"net/http"
"time"
2019-12-08 03:59:58 -05:00
"google.golang.org/grpc"
2024-10-23 08:37:29 -04:00
"github.com/hashicorp/go-uuid"
"github.com/naiba/nezha/model"
2024-10-23 08:37:29 -04:00
"github.com/naiba/nezha/pkg/utils"
"github.com/naiba/nezha/proto"
2020-11-10 21:07:45 -05:00
rpcService "github.com/naiba/nezha/service/rpc"
2022-01-08 22:54:14 -05:00
"github.com/naiba/nezha/service/singleton"
2019-12-08 03:59:58 -05:00
)
2024-10-23 00:55:10 -04:00
func ServeRPC() *grpc.Server {
2019-12-08 03:59:58 -05:00
server := grpc.NewServer()
rpcService.NezhaHandlerSingleton = rpcService.NewNezhaHandler()
2024-10-23 08:37:29 -04:00
proto.RegisterNezhaServiceServer(server, rpcService.NezhaHandlerSingleton)
2024-10-22 11:44:50 -04:00
return server
2019-12-08 03:59:58 -05:00
}
2024-10-24 12:13:45 -04:00
func DispatchTask(serviceSentinelDispatchBus <-chan model.Service) {
workedServerIndex := 0
for task := range serviceSentinelDispatchBus {
round := 0
endIndex := workedServerIndex
2022-01-08 22:54:14 -05:00
singleton.SortedServerLock.RLock()
// 如果已经轮了一整圈又轮到自己,没有合适机器去请求,跳出循环
for round < 1 || workedServerIndex < endIndex {
// 如果到了圈尾,再回到圈头,圈数加一,游标重置
2022-01-08 22:54:14 -05:00
if workedServerIndex >= len(singleton.SortedServerList) {
workedServerIndex = 0
round++
continue
}
// 如果服务器不在线,跳过这个服务器
2022-01-08 22:54:14 -05:00
if singleton.SortedServerList[workedServerIndex].TaskStream == nil {
workedServerIndex++
continue
}
// 如果此任务不可使用此服务器请求,跳过这个服务器(有些 IPv6 only 开了 NAT64 的机器请求 IPv4 总会出问题)
2024-10-24 12:13:45 -04:00
if (task.Cover == model.ServiceCoverAll && task.SkipServers[singleton.SortedServerList[workedServerIndex].ID]) ||
(task.Cover == model.ServiceCoverIgnoreAll && !task.SkipServers[singleton.SortedServerList[workedServerIndex].ID]) {
workedServerIndex++
continue
}
2024-10-24 12:13:45 -04:00
if task.Cover == model.ServiceCoverIgnoreAll && task.SkipServers[singleton.SortedServerList[workedServerIndex].ID] {
singleton.SortedServerList[workedServerIndex].TaskStream.Send(task.PB())
workedServerIndex++
continue
}
2024-10-24 12:13:45 -04:00
if task.Cover == model.ServiceCoverAll && !task.SkipServers[singleton.SortedServerList[workedServerIndex].ID] {
singleton.SortedServerList[workedServerIndex].TaskStream.Send(task.PB())
workedServerIndex++
continue
}
// 找到合适机器执行任务,跳出循环
// singleton.SortedServerList[workedServerIndex].TaskStream.Send(task.PB())
// workedServerIndex++
// break
}
2022-01-08 22:54:14 -05:00
singleton.SortedServerLock.RUnlock()
}
}
func DispatchKeepalive() {
2022-01-08 22:54:14 -05:00
singleton.Cron.AddFunc("@every 60s", func() {
singleton.SortedServerLock.RLock()
defer singleton.SortedServerLock.RUnlock()
for i := 0; i < len(singleton.SortedServerList); i++ {
if singleton.SortedServerList[i] == nil || singleton.SortedServerList[i].TaskStream == nil {
continue
}
2024-10-23 08:37:29 -04:00
singleton.SortedServerList[i].TaskStream.Send(&proto.Task{Type: model.TaskTypeKeepalive})
}
})
}
2024-10-23 08:37:29 -04:00
func ServeNAT(w http.ResponseWriter, r *http.Request, natConfig *model.NAT) {
singleton.ServerLock.RLock()
server := singleton.ServerList[natConfig.ServerID]
singleton.ServerLock.RUnlock()
if server == nil || server.TaskStream == nil {
w.WriteHeader(http.StatusServiceUnavailable)
w.Write([]byte("server not found or not connected"))
return
}
streamId, err := uuid.GenerateUUID()
if err != nil {
w.WriteHeader(http.StatusServiceUnavailable)
w.Write([]byte(fmt.Sprintf("stream id error: %v", err)))
return
}
rpcService.NezhaHandlerSingleton.CreateStream(streamId)
defer rpcService.NezhaHandlerSingleton.CloseStream(streamId)
taskData, err := utils.Json.Marshal(model.TaskNAT{
StreamID: streamId,
Host: natConfig.Host,
})
if err != nil {
w.WriteHeader(http.StatusServiceUnavailable)
w.Write([]byte(fmt.Sprintf("task data error: %v", err)))
return
}
if err := server.TaskStream.Send(&proto.Task{
Type: model.TaskTypeNAT,
Data: string(taskData),
}); err != nil {
w.WriteHeader(http.StatusServiceUnavailable)
w.Write([]byte(fmt.Sprintf("send task error: %v", err)))
return
}
wWrapped, err := utils.NewRequestWrapper(r, w)
if err != nil {
w.WriteHeader(http.StatusServiceUnavailable)
w.Write([]byte(fmt.Sprintf("request wrapper error: %v", err)))
return
}
if err := rpcService.NezhaHandlerSingleton.UserConnected(streamId, wWrapped); err != nil {
w.WriteHeader(http.StatusServiceUnavailable)
w.Write([]byte(fmt.Sprintf("user connected error: %v", err)))
return
}
rpcService.NezhaHandlerSingleton.StartStream(streamId, time.Second*10)
}