From 3d7a0b697d2d42835094fa2bd5df5a9b209528e9 Mon Sep 17 00:00:00 2001 From: sagit Date: Mon, 9 Feb 2026 08:47:13 +0000 Subject: [PATCH] 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()