-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathhandler.go
More file actions
130 lines (117 loc) · 2.84 KB
/
Copy pathhandler.go
File metadata and controls
130 lines (117 loc) · 2.84 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
package serverFinder
import (
"encoding/json"
"net/http"
"sync"
"time"
"github.com/cnlesscode/serverFinder/client"
"github.com/google/uuid"
"github.com/gorilla/websocket"
)
var upgrader = websocket.Upgrader{
ReadBufferSize: 1024,
WriteBufferSize: 1024,
CheckOrigin: func(r *http.Request) bool {
return true // 允许所有来源的连接
},
}
// 监听客户端连接池
var ConnsMu sync.RWMutex
var ListenClients = make(map[string]map[string]*websocket.Conn)
func Handler(w http.ResponseWriter, r *http.Request) {
// 初始化 url 参数
addr := r.URL.Query().Get("addr")
mainKey := r.URL.Query().Get("mainKey")
action := r.URL.Query().Get("action")
listen := r.URL.Query().Get("listen")
if mainKey == "" || action == "" || addr == "" {
return
}
switch action {
// 服务注册
case "register":
connUUID := uuid.New().String()
// 升级协议
conn, err := upgrader.Upgrade(w, r, nil)
if err != nil {
return
}
// 保存注册节点数据
SetItem(mainKey, addr, time.Now().Unix())
// 注册连接是否同时用于监听
if listen == "true" {
AddListener(mainKey, connUUID, conn)
}
// 连接被关闭
defer func() {
RemoveItem(mainKey, addr)
if listen == "true" {
RemoveListener(mainKey, connUUID)
}
}()
webSocketReadLoopHandle(conn)
// 获取数据
case "get":
data, ok := Get(mainKey)
if ok {
messageByte, err := json.Marshal(data)
if err != nil {
return
}
w.Write(messageByte)
}
// 监听数据变化
case "listen":
connUUID := uuid.New().String()
// 升级协议
conn, err := upgrader.Upgrade(w, r, nil)
if err != nil {
return
}
// 记录监听连接
AddListener(mainKey, connUUID, conn)
// 连接被关闭
defer func() {
RemoveListener(mainKey, connUUID)
}()
webSocketReadLoopHandle(conn)
}
}
// 保存连接到监听连接池
func AddListener(mainKey, id string, conn *websocket.Conn) {
// 记录监听连接
ConnsMu.Lock()
defer ConnsMu.Unlock()
if _, ok := ListenClients[mainKey]; ok {
ListenClients[mainKey][id] = conn
} else {
ListenClients[mainKey] = map[string]*websocket.Conn{}
ListenClients[mainKey][id] = conn
}
}
// 删除监听连接
func RemoveListener(mainKey, id string) {
ConnsMu.Lock()
defer ConnsMu.Unlock()
if conns, exists := ListenClients[mainKey]; exists {
delete(conns, id)
// 清理空的 mainKey 条目,防止内存泄漏
if len(conns) == 0 {
delete(ListenClients, mainKey)
}
}
}
func webSocketReadLoopHandle(conn *websocket.Conn) {
defer conn.Close()
conn.SetReadDeadline(time.Now().Add(client.ReadDeadlineTimer * time.Second))
conn.SetPingHandler(func(appData string) error {
conn.SetReadDeadline(time.Now().Add(client.ReadDeadlineTimer * time.Second))
return conn.WriteMessage(websocket.PongMessage, []byte(appData))
})
for {
_, _, err := conn.ReadMessage()
if err != nil {
break
}
}
}