From 565d7329672d68818afeb8e1d44c4fb120ba2c63 Mon Sep 17 00:00:00 2001 From: sagit Date: Mon, 9 Feb 2026 08:12:10 +0000 Subject: [PATCH 1/4] refactor(backend): reimplement speed limit logic 1. Refactor speed limit CRUD to sync with agents immediately via WebSocket (AddLimiters/DeleteLimiters). 2. Update unit conversion to match GOST v3 requirements (Mbps -> MB/s). 3. Update service config generation to reference Limiter IDs instead of hardcoded values. --- .../internal/http/handler/control_plane.go | 61 ++++++++++++++----- go-backend/internal/http/handler/mutations.go | 17 +++++- 2 files changed, 61 insertions(+), 17 deletions(-) diff --git a/go-backend/internal/http/handler/control_plane.go b/go-backend/internal/http/handler/control_plane.go index 174405a..0ece85d 100644 --- a/go-backend/internal/http/handler/control_plane.go +++ b/go-backend/internal/http/handler/control_plane.go @@ -234,9 +234,9 @@ func (h *Handler) getNodeRecord(nodeID int64) (*nodeRecord, error) { return &n, nil } -func (h *Handler) resolveUserTunnelAndLimiter(userID, tunnelID int64) (int64, *int, error) { +func (h *Handler) resolveUserTunnelAndLimiter(userID, tunnelID int64) (int64, *int64, error) { row := h.repo.DB().QueryRow(` - SELECT ut.id, sl.speed + SELECT ut.id, sl.id FROM user_tunnel ut LEFT JOIN speed_limit sl ON sl.id = ut.speed_id WHERE ut.user_id = ? AND ut.tunnel_id = ? @@ -244,18 +244,18 @@ func (h *Handler) resolveUserTunnelAndLimiter(userID, tunnelID int64) (int64, *i LIMIT 1 `, userID, tunnelID) var userTunnelID int64 - var speed sql.NullInt64 - err := row.Scan(&userTunnelID, &speed) + var limiterID sql.NullInt64 + err := row.Scan(&userTunnelID, &limiterID) if err != nil { if errors.Is(err, sql.ErrNoRows) { return 0, nil, nil } return 0, nil, err } - if !speed.Valid || speed.Int64 <= 0 { + if !limiterID.Valid || limiterID.Int64 <= 0 { return userTunnelID, nil, nil } - v := int(speed.Int64) + v := limiterID.Int64 return userTunnelID, &v, nil } @@ -328,7 +328,7 @@ func (h *Handler) syncForwardServices(forward *forwardRecord, method string, all return errors.New("转发入口端口不存在") } - userTunnelID, limiter, err := h.resolveUserTunnelAndLimiter(forward.UserID, forward.TunnelID) + userTunnelID, limiterID, err := h.resolveUserTunnelAndLimiter(forward.UserID, forward.TunnelID) if err != nil { return err } @@ -339,7 +339,7 @@ func (h *Handler) syncForwardServices(forward *forwardRecord, method string, all if err != nil { return err } - services := buildForwardServiceConfigs(serviceBase, forward, tunnel, node, fp.Port, limiter) + services := buildForwardServiceConfigs(serviceBase, forward, tunnel, node, fp.Port, limiterID) _, err = h.sendNodeCommand(node.ID, method, services, true, false) if err != nil && allowFallbackAdd && method == "UpdateService" { _, err = h.sendNodeCommand(node.ID, "AddService", services, true, false) @@ -1003,7 +1003,7 @@ func isNotFoundError(err error) bool { return strings.Contains(msg, "not found") || strings.Contains(msg, "不存在") } -func buildForwardServiceConfigs(baseName string, forward *forwardRecord, tunnel *tunnelRecord, node *nodeRecord, port int, limiter *int) []map[string]interface{} { +func buildForwardServiceConfigs(baseName string, forward *forwardRecord, tunnel *tunnelRecord, node *nodeRecord, port int, limiterID *int64) []map[string]interface{} { protocols := []string{"tcp", "udp"} services := make([]map[string]interface{}, 0, 2) targets := splitRemoteTargets(forward.RemoteAddr) @@ -1044,11 +1044,8 @@ func buildForwardServiceConfigs(baseName string, forward *forwardRecord, tunnel if tunnel != nil && tunnel.Type == 1 && strings.TrimSpace(node.InterfaceName) != "" { service["metadata"] = map[string]interface{}{"interface": node.InterfaceName} } - if limiter != nil && *limiter > 0 { - // Convert Mbps to Bytes/s - // 1 Mbps = 1,000,000 bits/s = 125,000 Bytes/s - // We use decimal Mbps standard as is common in networking - service["limiter"] = strconv.Itoa(*limiter * 125000) + if limiterID != nil && *limiterID > 0 { + service["limiter"] = strconv.FormatInt(*limiterID, 10) } services = append(services, service) } @@ -1111,3 +1108,39 @@ func asBool(v interface{}, def bool) bool { return def } } + +func (h *Handler) sendLimiterConfig(limiterID int64, speedMbps int, tunnelID int64) error { + rate := float64(speedMbps) / 8.0 + limitStr := fmt.Sprintf("$ %.1fMB %.1fMB", rate, rate) + + payload := map[string]interface{}{ + "name": strconv.FormatInt(limiterID, 10), + "limits": []string{limitStr}, + } + + nodes, err := h.tunnelEntryNodeIDs(tunnelID) + if err != nil { + return err + } + + for _, nodeID := range nodes { + _, _ = h.sendNodeCommand(nodeID, "AddLimiters", payload, false, false) + } + return nil +} + +func (h *Handler) sendDeleteLimiterConfig(limiterID int64, tunnelID int64) error { + payload := map[string]interface{}{ + "limiter": strconv.FormatInt(limiterID, 10), + } + + nodes, err := h.tunnelEntryNodeIDs(tunnelID) + if err != nil { + return err + } + + for _, nodeID := range nodes { + _, _ = h.sendNodeCommand(nodeID, "DeleteLimiters", payload, false, true) + } + return nil +} diff --git a/go-backend/internal/http/handler/mutations.go b/go-backend/internal/http/handler/mutations.go index 31da320..61849a3 100644 --- a/go-backend/internal/http/handler/mutations.go +++ b/go-backend/internal/http/handler/mutations.go @@ -1469,12 +1469,15 @@ func (h *Handler) speedLimitCreate(w http.ResponseWriter, r *http.Request) { return } now := time.Now().UnixMilli() - _, err := h.repo.DB().Exec(`INSERT INTO speed_limit(name, speed, tunnel_id, tunnel_name, created_time, updated_time, status) VALUES(?, ?, ?, ?, ?, ?, ?)`, - name, asInt(req["speed"], 100), tunnelID, tunnelName, now, now, asInt(req["status"], 1)) + speed := asInt(req["speed"], 100) + res, err := h.repo.DB().Exec(`INSERT INTO speed_limit(name, speed, tunnel_id, tunnel_name, created_time, updated_time, status) VALUES(?, ?, ?, ?, ?, ?, ?)`, + name, speed, tunnelID, tunnelName, now, now, asInt(req["status"], 1)) if err != nil { response.WriteJSON(w, response.Err(-2, err.Error())) return } + id, _ := res.LastInsertId() + _ = h.sendLimiterConfig(id, speed, tunnelID) response.WriteJSON(w, response.OKEmpty()) } @@ -1496,12 +1499,14 @@ func (h *Handler) speedLimitUpdate(w http.ResponseWriter, r *http.Request) { response.WriteJSON(w, response.ErrDefault("隧道不存在")) return } + speed := asInt(req["speed"], 100) _, err := h.repo.DB().Exec(`UPDATE speed_limit SET name=?, speed=?, tunnel_id=?, tunnel_name=?, status=?, updated_time=? WHERE id=?`, - asString(req["name"]), asInt(req["speed"], 100), tunnelID, tunnelName, asInt(req["status"], 1), time.Now().UnixMilli(), id) + asString(req["name"]), speed, tunnelID, tunnelName, asInt(req["status"], 1), time.Now().UnixMilli(), id) if err != nil { response.WriteJSON(w, response.Err(-2, err.Error())) return } + _ = h.sendLimiterConfig(id, speed, tunnelID) response.WriteJSON(w, response.OKEmpty()) } @@ -1510,11 +1515,17 @@ func (h *Handler) speedLimitDelete(w http.ResponseWriter, r *http.Request) { if id <= 0 { return } + var tunnelID int64 + _ = h.repo.DB().QueryRow(`SELECT tunnel_id FROM speed_limit WHERE id = ?`, id).Scan(&tunnelID) + _, err := h.repo.DB().Exec(`DELETE FROM speed_limit WHERE id = ?`, id) if err != nil { response.WriteJSON(w, response.Err(-2, err.Error())) return } + if tunnelID > 0 { + _ = h.sendDeleteLimiterConfig(id, tunnelID) + } response.WriteJSON(w, response.OKEmpty()) } From 3d7a0b697d2d42835094fa2bd5df5a9b209528e9 Mon Sep 17 00:00:00 2001 From: sagit Date: Mon, 9 Feb 2026 08:47:13 +0000 Subject: [PATCH 2/4] feat(backend): sync limiters on agent connect Implemented full sync of speed limit configurations when an Agent connects via WebSocket. This ensures that even fresh or restarted agents receive the necessary limiter configurations. --- .../internal/http/handler/control_plane.go | 27 +++++++++++++++++++ go-backend/internal/http/handler/handler.go | 4 ++- go-backend/internal/ws/server.go | 6 +++++ 3 files changed, 36 insertions(+), 1 deletion(-) diff --git a/go-backend/internal/http/handler/control_plane.go b/go-backend/internal/http/handler/control_plane.go index 0ece85d..c4b3f76 100644 --- a/go-backend/internal/http/handler/control_plane.go +++ b/go-backend/internal/http/handler/control_plane.go @@ -1144,3 +1144,30 @@ func (h *Handler) sendDeleteLimiterConfig(limiterID int64, tunnelID int64) error } return nil } + +func (h *Handler) onNodeConnected(nodeID int64) { + // Sync limiters + rows, err := h.repo.DB().Query(` + SELECT DISTINCT sl.id, sl.speed + FROM speed_limit sl + JOIN chain_tunnel ct ON ct.tunnel_id = sl.tunnel_id + WHERE ct.node_id = ? AND ct.chain_type = 1 AND sl.status = 1 + `, nodeID) + + if err == nil { + defer rows.Close() + for rows.Next() { + var id int64 + var speed int + if err := rows.Scan(&id, &speed); err == nil { + rate := float64(speed) / 8.0 + limitStr := fmt.Sprintf("$ %.1fMB %.1fMB", rate, rate) + payload := map[string]interface{}{ + "name": strconv.FormatInt(id, 10), + "limits": []string{limitStr}, + } + _, _ = h.sendNodeCommand(nodeID, "AddLimiters", payload, false, false) + } + } + } +} diff --git a/go-backend/internal/http/handler/handler.go b/go-backend/internal/http/handler/handler.go index 8236cfe..ac48973 100644 --- a/go-backend/internal/http/handler/handler.go +++ b/go-backend/internal/http/handler/handler.go @@ -62,11 +62,13 @@ type flowItem struct { } func New(repo *sqlite.Repository, jwtSecret string) *Handler { - return &Handler{ + h := &Handler{ repo: repo, jwtSecret: jwtSecret, wsServer: ws.NewServer(repo, jwtSecret), } + h.wsServer.OnNodeConnected = h.onNodeConnected + return h } func (h *Handler) WebSocketHandler() http.Handler { diff --git a/go-backend/internal/ws/server.go b/go-backend/internal/ws/server.go index f36ec97..1ff4fd6 100644 --- a/go-backend/internal/ws/server.go +++ b/go-backend/internal/ws/server.go @@ -71,6 +71,8 @@ type Server struct { nodes map[int64]*nodeSession byConn map[*websocket.Conn]*nodeSession pending map[string]pendingRequest + + OnNodeConnected func(nodeID int64) } func NewServer(repo *sqlite.Repository, jwtSecret string) *Server { @@ -164,6 +166,10 @@ func (s *Server) handleNode(w http.ResponseWriter, r *http.Request, nodeID int64 _ = s.repo.UpdateNodeOnline(nodeID, 1, version, httpVal, tlsVal, socksVal) s.broadcastStatus(nodeID, 1) + if s.OnNodeConnected != nil { + go s.OnNodeConnected(nodeID) + } + defer func() { needOfflineBroadcast := false s.mu.Lock() From 3420dc5460308f8250658c77ab17e68b47861403 Mon Sep 17 00:00:00 2001 From: sagit Date: Mon, 9 Feb 2026 09:06:59 +0000 Subject: [PATCH 3/4] fix(backend): sync limiter on association instead of connection Reverted the full sync on connection hook. Instead, ensureLimiterOnNode is called within syncForwardServices to push limiter configuration immediately before pushing the service configuration that references it. --- .../internal/http/handler/control_plane.go | 55 ++++++++----------- go-backend/internal/http/handler/handler.go | 4 +- go-backend/internal/ws/server.go | 6 -- 3 files changed, 23 insertions(+), 42 deletions(-) diff --git a/go-backend/internal/http/handler/control_plane.go b/go-backend/internal/http/handler/control_plane.go index c4b3f76..75452eb 100644 --- a/go-backend/internal/http/handler/control_plane.go +++ b/go-backend/internal/http/handler/control_plane.go @@ -234,9 +234,9 @@ func (h *Handler) getNodeRecord(nodeID int64) (*nodeRecord, error) { return &n, nil } -func (h *Handler) resolveUserTunnelAndLimiter(userID, tunnelID int64) (int64, *int64, error) { +func (h *Handler) resolveUserTunnelAndLimiter(userID, tunnelID int64) (int64, *int64, *int, error) { row := h.repo.DB().QueryRow(` - SELECT ut.id, sl.id + SELECT ut.id, sl.id, sl.speed FROM user_tunnel ut LEFT JOIN speed_limit sl ON sl.id = ut.speed_id WHERE ut.user_id = ? AND ut.tunnel_id = ? @@ -245,18 +245,20 @@ func (h *Handler) resolveUserTunnelAndLimiter(userID, tunnelID int64) (int64, *i `, userID, tunnelID) var userTunnelID int64 var limiterID sql.NullInt64 - err := row.Scan(&userTunnelID, &limiterID) + var speed sql.NullInt64 + err := row.Scan(&userTunnelID, &limiterID, &speed) if err != nil { if errors.Is(err, sql.ErrNoRows) { - return 0, nil, nil + return 0, nil, nil, nil } - return 0, nil, err + return 0, nil, nil, err } if !limiterID.Valid || limiterID.Int64 <= 0 { - return userTunnelID, nil, nil + return userTunnelID, nil, nil, nil } v := limiterID.Int64 - return userTunnelID, &v, nil + s := int(speed.Int64) + return userTunnelID, &v, &s, nil } func (h *Handler) listUserTunnelIDs(userID, tunnelID int64) ([]int64, error) { @@ -328,13 +330,17 @@ func (h *Handler) syncForwardServices(forward *forwardRecord, method string, all return errors.New("转发入口端口不存在") } - userTunnelID, limiterID, err := h.resolveUserTunnelAndLimiter(forward.UserID, forward.TunnelID) + userTunnelID, limiterID, speed, err := h.resolveUserTunnelAndLimiter(forward.UserID, forward.TunnelID) if err != nil { return err } serviceBase := buildForwardServiceBase(forward.ID, forward.UserID, userTunnelID) for _, fp := range ports { + if limiterID != nil && speed != nil { + h.ensureLimiterOnNode(fp.NodeID, *limiterID, *speed) + } + node, err := h.getNodeRecord(fp.NodeID) if err != nil { return err @@ -362,7 +368,7 @@ func (h *Handler) controlForwardServices(forward *forwardRecord, commandType str if len(ports) == 0 { return nil } - userTunnelID, _, err := h.resolveUserTunnelAndLimiter(forward.UserID, forward.TunnelID) + userTunnelID, _, _, err := h.resolveUserTunnelAndLimiter(forward.UserID, forward.TunnelID) if err != nil { return err } @@ -1145,29 +1151,12 @@ func (h *Handler) sendDeleteLimiterConfig(limiterID int64, tunnelID int64) error return nil } -func (h *Handler) onNodeConnected(nodeID int64) { - // Sync limiters - rows, err := h.repo.DB().Query(` - SELECT DISTINCT sl.id, sl.speed - FROM speed_limit sl - JOIN chain_tunnel ct ON ct.tunnel_id = sl.tunnel_id - WHERE ct.node_id = ? AND ct.chain_type = 1 AND sl.status = 1 - `, nodeID) - - if err == nil { - defer rows.Close() - for rows.Next() { - var id int64 - var speed int - if err := rows.Scan(&id, &speed); err == nil { - rate := float64(speed) / 8.0 - limitStr := fmt.Sprintf("$ %.1fMB %.1fMB", rate, rate) - payload := map[string]interface{}{ - "name": strconv.FormatInt(id, 10), - "limits": []string{limitStr}, - } - _, _ = h.sendNodeCommand(nodeID, "AddLimiters", payload, false, false) - } - } +func (h *Handler) ensureLimiterOnNode(nodeID int64, limiterID int64, speed int) { + rate := float64(speed) / 8.0 + limitStr := fmt.Sprintf("$ %.1fMB %.1fMB", rate, rate) + payload := map[string]interface{}{ + "name": strconv.FormatInt(limiterID, 10), + "limits": []string{limitStr}, } + _, _ = h.sendNodeCommand(nodeID, "AddLimiters", payload, false, false) } diff --git a/go-backend/internal/http/handler/handler.go b/go-backend/internal/http/handler/handler.go index ac48973..8236cfe 100644 --- a/go-backend/internal/http/handler/handler.go +++ b/go-backend/internal/http/handler/handler.go @@ -62,13 +62,11 @@ type flowItem struct { } func New(repo *sqlite.Repository, jwtSecret string) *Handler { - h := &Handler{ + return &Handler{ repo: repo, jwtSecret: jwtSecret, wsServer: ws.NewServer(repo, jwtSecret), } - h.wsServer.OnNodeConnected = h.onNodeConnected - return h } func (h *Handler) WebSocketHandler() http.Handler { diff --git a/go-backend/internal/ws/server.go b/go-backend/internal/ws/server.go index 1ff4fd6..f36ec97 100644 --- a/go-backend/internal/ws/server.go +++ b/go-backend/internal/ws/server.go @@ -71,8 +71,6 @@ type Server struct { nodes map[int64]*nodeSession byConn map[*websocket.Conn]*nodeSession pending map[string]pendingRequest - - OnNodeConnected func(nodeID int64) } func NewServer(repo *sqlite.Repository, jwtSecret string) *Server { @@ -166,10 +164,6 @@ func (s *Server) handleNode(w http.ResponseWriter, r *http.Request, nodeID int64 _ = s.repo.UpdateNodeOnline(nodeID, 1, version, httpVal, tlsVal, socksVal) s.broadcastStatus(nodeID, 1) - if s.OnNodeConnected != nil { - go s.OnNodeConnected(nodeID) - } - defer func() { needOfflineBroadcast := false s.mu.Lock() From 47b1663938a7180897469c52c074c5c7eb5df351 Mon Sep 17 00:00:00 2001 From: sagit Date: Mon, 9 Feb 2026 10:46:58 +0000 Subject: [PATCH 4/4] feat: add Shanghai timezone to docker-compose and update github mirror --- docker-compose-v4.yml | 1 + docker-compose-v6.yml | 1 + install.sh | 2 +- panel_install.sh | 2 +- 4 files changed, 4 insertions(+), 2 deletions(-) diff --git a/docker-compose-v4.yml b/docker-compose-v4.yml index f13eb1b..1169509 100644 --- a/docker-compose-v4.yml +++ b/docker-compose-v4.yml @@ -12,6 +12,7 @@ services: JWT_SECRET: ${JWT_SECRET} LOG_DIR: /app/logs SERVER_ADDR: :6365 + TZ: Asia/Shanghai ports: - "${BACKEND_PORT}:6365" volumes: diff --git a/docker-compose-v6.yml b/docker-compose-v6.yml index 15ba827..e0832eb 100644 --- a/docker-compose-v6.yml +++ b/docker-compose-v6.yml @@ -12,6 +12,7 @@ services: JWT_SECRET: ${JWT_SECRET} LOG_DIR: /app/logs SERVER_ADDR: :6365 + TZ: Asia/Shanghai ports: - "${BACKEND_PORT}:6365" volumes: diff --git a/install.sh b/install.sh index 8f0cca9..e09394a 100644 --- a/install.sh +++ b/install.sh @@ -28,7 +28,7 @@ COUNTRY=$(curl -s https://ipinfo.io/country) maybe_proxy_url() { local url="$1" if [ "$COUNTRY" = "CN" ]; then - echo "https://ghfast.top/${url}" + echo "https://gcode.hostcentral.cc/${url}" else echo "$url" fi diff --git a/panel_install.sh b/panel_install.sh index e977ba9..bb77c25 100755 --- a/panel_install.sh +++ b/panel_install.sh @@ -15,7 +15,7 @@ COUNTRY=$(curl -s https://ipinfo.io/country) maybe_proxy_url() { local url="$1" if [ "$COUNTRY" = "CN" ]; then - echo "https://ghfast.top/${url}" + echo "https://gcode.hostcentral.cc/${url}" else echo "$url" fi