Compare commits

...

9 Commits

Author SHA1 Message Date
sagit b32133f81a Merge pull request #92 from Sagit-chu/opencode/lucky-otter
fix(upgrade): stabilize batch node upgrades
2026-02-12 13:37:43 +08:00
sagit ff57bca505 Merge branch 'main' into opencode/lucky-otter 2026-02-12 13:11:45 +08:00
sagit cdb2914dbf fix(upgrade): stabilize batch node upgrades under long-running operations 2026-02-12 04:57:58 +00:00
sagit b56d0a28e7 Merge pull request #91 from Sagit-chu/opencode/lucky-otter
fix(ws): prevent monitor websocket reconnect loop
2026-02-12 12:28:53 +08:00
sagit bdfc704f95 Merge branch 'main' into opencode/lucky-otter 2026-02-12 12:27:23 +08:00
sagit 9223892ca5 fix(ws): prevent monitor websocket reconnect loop 2026-02-12 04:25:47 +00:00
sagit 6f205df37c Merge pull request #90 from Sagit-chu/opencode/lucky-otter
fix(ws): stabilize node connectivity with ping/pong keepalive
2026-02-12 10:29:14 +08:00
sagit 04266165df Merge branch 'main' into opencode/lucky-otter 2026-02-12 10:28:04 +08:00
sagit dbd5773717 fix(ws): stabilize node connectivity with ping/pong keepalive 2026-02-12 02:27:10 +00:00
6 changed files with 104 additions and 22 deletions
+26 -12
View File
@@ -6,6 +6,7 @@ import (
"io"
"net/http"
"strings"
"sync"
"time"
"go-backend/internal/http/response"
@@ -16,6 +17,8 @@ const (
githubProxy = "https://gcode.hostcentral.cc"
githubAPIBase = "https://api.github.com"
githubHTMLBase = "https://github.com"
upgradeTimeout = 5 * time.Minute
batchWorkers = 5
)
func (h *Handler) nodeUpgrade(w http.ResponseWriter, r *http.Request) {
@@ -59,7 +62,7 @@ func (h *Handler) nodeUpgrade(w http.ResponseWriter, r *http.Request) {
result, err := h.wsServer.SendCommand(req.ID, "UpgradeAgent", map[string]interface{}{
"downloadUrl": downloadURL,
"checksumUrl": checksumURL,
}, 120*time.Second)
}, upgradeTimeout)
if err != nil {
response.WriteJSON(w, response.Err(-2, fmt.Sprintf("升级失败: %v", err)))
return
@@ -173,18 +176,29 @@ func (h *Handler) nodeBatchUpgrade(w http.ResponseWriter, r *http.Request) {
Message string `json:"message"`
}
results := make([]upgradeResult, 0, len(req.IDs))
for _, id := range req.IDs {
result, err := h.wsServer.SendCommand(id, "UpgradeAgent", map[string]interface{}{
"downloadUrl": downloadURL,
"checksumUrl": checksumURL,
}, 120*time.Second)
if err != nil {
results = append(results, upgradeResult{ID: id, Success: false, Message: err.Error()})
} else {
results = append(results, upgradeResult{ID: id, Success: true, Message: result.Message})
}
results := make([]upgradeResult, len(req.IDs))
sem := make(chan struct{}, batchWorkers)
var wg sync.WaitGroup
for i, id := range req.IDs {
wg.Add(1)
go func(index int, nodeID int64) {
defer wg.Done()
sem <- struct{}{}
defer func() { <-sem }()
result, err := h.wsServer.SendCommand(nodeID, "UpgradeAgent", map[string]interface{}{
"downloadUrl": downloadURL,
"checksumUrl": checksumURL,
}, upgradeTimeout)
if err != nil {
results[index] = upgradeResult{ID: nodeID, Success: false, Message: err.Error()}
return
}
results[index] = upgradeResult{ID: nodeID, Success: true, Message: result.Message}
}(i, id)
}
wg.Wait()
response.WriteJSON(w, response.OK(map[string]interface{}{
"version": version,
+49
View File
@@ -54,6 +54,12 @@ type pendingRequest struct {
ch chan CommandResult
}
const (
wsPingPeriod = 15 * time.Second
wsPongWait = 45 * time.Second
wsWriteWait = 5 * time.Second
)
type CommandResult struct {
Type string `json:"type"`
Success bool `json:"success"`
@@ -120,12 +126,19 @@ func (s *Server) handleAdmin(w http.ResponseWriter, r *http.Request) {
return
}
cw := &connWrap{conn: conn}
_ = conn.SetReadDeadline(time.Now().Add(wsPongWait))
conn.SetPongHandler(func(string) error {
return conn.SetReadDeadline(time.Now().Add(wsPongWait))
})
done := make(chan struct{})
go startKeepalive(cw, done)
s.mu.Lock()
s.admins[cw] = struct{}{}
s.mu.Unlock()
defer func() {
close(done)
s.mu.Lock()
delete(s.admins, cw)
s.mu.Unlock()
@@ -145,6 +158,12 @@ func (s *Server) handleNode(w http.ResponseWriter, r *http.Request, nodeID int64
return
}
cw := &connWrap{conn: conn}
_ = conn.SetReadDeadline(time.Now().Add(wsPongWait))
conn.SetPongHandler(func(string) error {
return conn.SetReadDeadline(time.Now().Add(wsPongWait))
})
done := make(chan struct{})
go startKeepalive(cw, done)
version := r.URL.Query().Get("version")
httpVal := parseIntDefault(r.URL.Query().Get("http"), 0)
@@ -165,6 +184,7 @@ func (s *Server) handleNode(w http.ResponseWriter, r *http.Request, nodeID int64
s.broadcastStatus(nodeID, 1)
defer func() {
close(done)
needOfflineBroadcast := false
s.mu.Lock()
current, ok := s.nodes[nodeID]
@@ -272,7 +292,9 @@ func (s *Server) SendCommand(nodeID int64, cmdType string, data interface{}, tim
}
ns.conn.mu.Lock()
_ = ns.conn.conn.SetWriteDeadline(time.Now().Add(wsWriteWait))
err = ns.conn.conn.WriteMessage(websocket.TextMessage, messageData)
_ = ns.conn.conn.SetWriteDeadline(time.Time{})
ns.conn.mu.Unlock()
if err != nil {
cleanup()
@@ -409,7 +431,9 @@ func (s *Server) broadcastToAdmins(message string) {
for _, c := range admins {
c.mu.Lock()
_ = c.conn.SetWriteDeadline(time.Now().Add(wsWriteWait))
err := c.conn.WriteMessage(websocket.TextMessage, []byte(message))
_ = c.conn.SetWriteDeadline(time.Time{})
c.mu.Unlock()
if err != nil {
log.Printf("websocket broadcast failed: %v", err)
@@ -442,3 +466,28 @@ func parseIntDefault(v string, fallback int) int {
}
return x
}
func startKeepalive(cw *connWrap, done <-chan struct{}) {
if cw == nil || cw.conn == nil {
return
}
ticker := time.NewTicker(wsPingPeriod)
defer ticker.Stop()
for {
select {
case <-done:
return
case <-ticker.C:
cw.mu.Lock()
_ = cw.conn.SetWriteDeadline(time.Now().Add(wsWriteWait))
err := cw.conn.WriteMessage(websocket.PingMessage, nil)
_ = cw.conn.SetWriteDeadline(time.Time{})
cw.mu.Unlock()
if err != nil {
_ = cw.conn.Close()
return
}
}
}
}
+16 -5
View File
@@ -91,6 +91,11 @@ type TcpPingResponse struct {
RequestId string `json:"requestId,omitempty"`
}
const (
reporterReadWait = 60 * time.Second
reporterWriteWait = 5 * time.Second
)
type WebSocketReporter struct {
url string
addr string // 保存服务器地址
@@ -243,6 +248,14 @@ func (w *WebSocketReporter) connect() error {
w.conn = conn
w.connected = true
_ = conn.SetReadDeadline(time.Now().Add(reporterReadWait))
conn.SetPingHandler(func(appData string) error {
_ = conn.SetReadDeadline(time.Now().Add(reporterReadWait))
return conn.WriteControl(websocket.PongMessage, []byte(appData), time.Now().Add(reporterWriteWait))
})
conn.SetPongHandler(func(string) error {
return conn.SetReadDeadline(time.Now().Add(reporterReadWait))
})
// 设置关闭处理器来检测连接状态
w.conn.SetCloseHandler(func(code int, text string) error {
@@ -383,7 +396,7 @@ func (w *WebSocketReporter) receiveMessages() {
}
// 设置读取超时
conn.SetReadDeadline(time.Now().Add(30 * time.Second))
conn.SetReadDeadline(time.Now().Add(reporterReadWait))
messageType, message, err := conn.ReadMessage()
if err != nil {
@@ -472,9 +485,8 @@ func (w *WebSocketReporter) handleReceivedMessage(messageType int, message []byt
}
if cmdMsg.Type != "call" {
// TcpPing 诊断命令异步执行,避免阻塞其他命令
// 其他状态变更命令保持同步,确保顺序执行
if cmdMsg.Type == "TcpPing" {
if cmdMsg.Type == "TcpPing" || cmdMsg.Type == "UpgradeAgent" || cmdMsg.Type == "RollbackAgent" {
go w.routeCommand(cmdMsg)
} else {
w.routeCommand(cmdMsg)
@@ -489,9 +501,8 @@ func (w *WebSocketReporter) handleReceivedMessage(messageType int, message []byt
return
}
if cmdMsg.Type != "call" {
// TcpPing 诊断命令异步执行,避免阻塞其他命令
// 其他状态变更命令保持同步,确保顺序执行
if cmdMsg.Type == "TcpPing" {
if cmdMsg.Type == "TcpPing" || cmdMsg.Type == "UpgradeAgent" || cmdMsg.Type == "RollbackAgent" {
go w.routeCommand(cmdMsg)
} else {
w.routeCommand(cmdMsg)
+3 -1
View File
@@ -87,6 +87,8 @@ http {
proxy_http_version 1.1;
proxy_set_header Upgrade $http_upgrade;
proxy_set_header Connection "upgrade";
proxy_read_timeout 3600s;
proxy_send_timeout 3600s;
proxy_set_header Host $host;
proxy_set_header X-Real-IP $remote_addr;
@@ -94,4 +96,4 @@ http {
proxy_set_header X-Forwarded-Proto $scheme;
}
}
}
}
+2 -2
View File
@@ -42,9 +42,9 @@ export const checkNodeStatus = (nodeId?: number) => {
};
export const upgradeNode = (id: number, version?: string) =>
Network.post("/node/upgrade", { id, version: version || "" });
Network.post("/node/upgrade", { id, version: version || "" }, { timeout: 5 * 60 * 1000 });
export const batchUpgradeNodes = (ids: number[], version?: string) =>
Network.post("/node/batch-upgrade", { ids, version: version || "" });
Network.post("/node/batch-upgrade", { ids, version: version || "" }, { timeout: 15 * 60 * 1000 });
export const getNodeReleases = () => Network.post("/node/releases");
export const rollbackNode = (id: number) =>
Network.post("/node/rollback", { id });
+8 -2
View File
@@ -43,6 +43,10 @@ interface ApiResponse<T = any> {
data: T;
}
interface RequestOptions {
timeout?: number;
}
// 处理token失效的逻辑
function handleTokenExpired() {
// 清除localStorage中的token
@@ -71,6 +75,7 @@ const Network = {
get: function <T = any>(
path: string = "",
data: any = {},
options: RequestOptions = {},
): Promise<ApiResponse<T>> {
return new Promise(function (resolve) {
// 如果baseURL是默认值且是WebView环境,说明没有设置面板地址
@@ -83,7 +88,7 @@ const Network = {
axios
.get(path, {
params: data,
timeout: 30000,
timeout: options.timeout ?? 30000,
headers: {
Authorization: window.localStorage.getItem("token"),
},
@@ -117,6 +122,7 @@ const Network = {
post: function <T = any>(
path: string = "",
data: any = {},
options: RequestOptions = {},
): Promise<ApiResponse<T>> {
return new Promise(function (resolve) {
// 如果baseURL是默认值且是WebView环境,说明没有设置面板地址
@@ -128,7 +134,7 @@ const Network = {
axios
.post(path, data, {
timeout: 30000,
timeout: options.timeout ?? 30000,
headers: {
Authorization: window.localStorage.getItem("token"),
"Content-Type": "application/json",