From 11dc21e46fcf475d5b631b07a097edcb92c0dbaf Mon Sep 17 00:00:00 2001 From: sagit Date: Sun, 8 Feb 2026 09:03:56 +0000 Subject: [PATCH] feat: implement consumer side panel peering logic --- .sisyphus/ralph-loop.local.md | 9 + go-backend/internal/http/client/federation.go | 115 +++++++ .../internal/http/handler/federation.go | 249 ++++++++++++++ go-backend/internal/http/handler/handler.go | 3 + go-backend/internal/http/handler/mutations.go | 48 ++- .../internal/store/sqlite/repository.go | 180 +++++++++- vite-frontend/src/App.tsx | 11 +- vite-frontend/src/api/index.ts | 17 + vite-frontend/src/layouts/admin.tsx | 10 + vite-frontend/src/pages/node.tsx | 16 +- vite-frontend/src/pages/panel-sharing.tsx | 313 ++++++++++++++++++ 11 files changed, 954 insertions(+), 17 deletions(-) create mode 100644 .sisyphus/ralph-loop.local.md create mode 100644 go-backend/internal/http/client/federation.go create mode 100644 go-backend/internal/http/handler/federation.go create mode 100644 vite-frontend/src/pages/panel-sharing.tsx diff --git a/.sisyphus/ralph-loop.local.md b/.sisyphus/ralph-loop.local.md new file mode 100644 index 0000000..c919270 --- /dev/null +++ b/.sisyphus/ralph-loop.local.md @@ -0,0 +1,9 @@ +--- +active: true +iteration: 1 +max_iterations: 100 +completion_promise: "DONE" +started_at: "2026-02-08T08:59:10.522Z" +session_id: "ses_3c3ac2d4bffe6K5IO2BeAEWe7g" +--- +完善消费者 (Consumer) 侧的对接逻辑,完成后提交全部变更并push diff --git a/go-backend/internal/http/client/federation.go b/go-backend/internal/http/client/federation.go new file mode 100644 index 0000000..c57a0d2 --- /dev/null +++ b/go-backend/internal/http/client/federation.go @@ -0,0 +1,115 @@ +package client + +import ( + "encoding/json" + "fmt" + "io" + "net/http" + "strings" + "time" +) + +type FederationClient struct { + client *http.Client +} + +type RemoteNodeInfo struct { + ShareID int64 `json:"shareId"` + ShareName string `json:"shareName"` + NodeID int64 `json:"nodeId"` + NodeName string `json:"nodeName"` + ServerIP string `json:"serverIp"` + Status int `json:"status"` + MaxBandwidth int64 `json:"maxBandwidth"` + ExpiryTime int64 `json:"expiryTime"` + PortRangeStart int `json:"portRangeStart"` + PortRangeEnd int `json:"portRangeEnd"` +} + +type RemoteTunnelResponse struct { + TunnelID int64 `json:"tunnelId"` +} + +func NewFederationClient() *FederationClient { + return &FederationClient{ + client: &http.Client{ + Timeout: 10 * time.Second, + }, + } +} + +func (c *FederationClient) Connect(url, token string) (*RemoteNodeInfo, error) { + url = strings.TrimSuffix(url, "/") + req, err := http.NewRequest("POST", url+"/api/v1/federation/connect", nil) + if err != nil { + return nil, err + } + req.Header.Set("Authorization", "Bearer "+token) + req.Header.Set("Content-Type", "application/json") + + resp, err := c.client.Do(req) + if err != nil { + return nil, err + } + defer resp.Body.Close() + + if resp.StatusCode != 200 { + body, _ := io.ReadAll(resp.Body) + return nil, fmt.Errorf("remote error %d: %s", resp.StatusCode, string(body)) + } + + var res struct { + Code int `json:"code"` + Msg string `json:"msg"` + Data RemoteNodeInfo `json:"data"` + } + if err := json.NewDecoder(resp.Body).Decode(&res); err != nil { + return nil, err + } + if res.Code != 0 { + return nil, fmt.Errorf("remote api error: %s", res.Msg) + } + + return &res.Data, nil +} + +func (c *FederationClient) CreateTunnel(url, token string, protocol string, remotePort int, target string) (*RemoteTunnelResponse, error) { + url = strings.TrimSuffix(url, "/") + payload := map[string]interface{}{ + "protocol": protocol, + "remotePort": remotePort, + "target": target, + } + bodyBytes, _ := json.Marshal(payload) + req, err := http.NewRequest("POST", url+"/api/v1/federation/tunnel/create", strings.NewReader(string(bodyBytes))) + if err != nil { + return nil, err + } + req.Header.Set("Authorization", "Bearer "+token) + req.Header.Set("Content-Type", "application/json") + + resp, err := c.client.Do(req) + if err != nil { + return nil, err + } + defer resp.Body.Close() + + if resp.StatusCode != 200 { + body, _ := io.ReadAll(resp.Body) + return nil, fmt.Errorf("remote error %d: %s", resp.StatusCode, string(body)) + } + + var res struct { + Code int `json:"code"` + Msg string `json:"msg"` + Data RemoteTunnelResponse `json:"data"` + } + if err := json.NewDecoder(resp.Body).Decode(&res); err != nil { + return nil, err + } + if res.Code != 0 { + return nil, fmt.Errorf("remote api error: %s", res.Msg) + } + + return &res.Data, nil +} diff --git a/go-backend/internal/http/handler/federation.go b/go-backend/internal/http/handler/federation.go new file mode 100644 index 0000000..f8eb3c7 --- /dev/null +++ b/go-backend/internal/http/handler/federation.go @@ -0,0 +1,249 @@ +package handler + +import ( + "encoding/json" + "fmt" + "net/http" + "strings" + "time" + + "go-backend/internal/http/client" + "go-backend/internal/http/response" +) + +type federationTunnelRequest struct { + Protocol string `json:"protocol"` + RemotePort int `json:"remotePort"` + Target string `json:"target"` +} + +type nodeImportRequest struct { + RemoteURL string `json:"remoteUrl"` + Token string `json:"token"` +} + +func (h *Handler) nodeImport(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + response.WriteJSON(w, response.ErrDefault("Invalid method")) + return + } + + var req nodeImportRequest + if err := decodeJSON(r.Body, &req); err != nil { + response.WriteJSON(w, response.ErrDefault("Invalid JSON")) + return + } + + if req.RemoteURL == "" || req.Token == "" { + response.WriteJSON(w, response.ErrDefault("Remote URL and Token are required")) + return + } + + fc := client.NewFederationClient() + info, err := fc.Connect(req.RemoteURL, req.Token) + if err != nil { + response.WriteJSON(w, response.Err(-2, "Failed to connect: "+err.Error())) + return + } + + // Prepare config json for local storage (metadata about limits) + configData := map[string]interface{}{ + "shareId": info.ShareID, + "maxBandwidth": info.MaxBandwidth, + "expiryTime": info.ExpiryTime, + "portRangeStart": info.PortRangeStart, + "portRangeEnd": info.PortRangeEnd, + } + configBytes, _ := json.Marshal(configData) + + db := h.repo.DB() + inx := nextIndex(db, "node") + now := time.Now().UnixMilli() + + _, err = db.Exec(` + INSERT INTO node(name, secret, server_ip, server_ip_v4, server_ip_v6, port, interface_name, version, http, tls, socks, created_time, updated_time, status, tcp_listen_addr, udp_listen_addr, inx, is_remote, remote_url, remote_token, remote_config) + VALUES(?, ?, ?, ?, ?, ?, ?, ?, 0, 0, 0, ?, ?, ?, ?, ?, ?, 1, ?, ?, ?) + `, + fmt.Sprintf("%s (Remote)", info.NodeName), + randomToken(16), // Dummy secret + info.ServerIP, + "", "", // v4/v6 unknown, use server_ip + "0", // port range not applicable for remote + "", + "", + now, now, + info.Status, + "[::]", "[::]", + inx, + req.RemoteURL, + req.Token, + string(configBytes), + ) + + if err != nil { + response.WriteJSON(w, response.Err(-2, "Database error: "+err.Error())) + return + } + + response.WriteJSON(w, response.OKEmpty()) +} + +func (h *Handler) authPeer(next http.HandlerFunc) http.HandlerFunc { + return func(w http.ResponseWriter, r *http.Request) { + authHeader := r.Header.Get("Authorization") + if authHeader == "" { + response.WriteJSON(w, response.Err(401, "Missing Authorization header")) + return + } + + parts := strings.Split(authHeader, " ") + if len(parts) != 2 || parts[0] != "Bearer" { + response.WriteJSON(w, response.Err(401, "Invalid Authorization format")) + return + } + + token := parts[1] + share, err := h.repo.GetPeerShareByToken(token) + if err != nil { + response.WriteJSON(w, response.Err(-2, err.Error())) + return + } + if share == nil { + response.WriteJSON(w, response.Err(401, "Invalid token")) + return + } + + if share.IsActive == 0 { + response.WriteJSON(w, response.Err(403, "Share is disabled")) + return + } + + if share.ExpiryTime > 0 && share.ExpiryTime < time.Now().UnixMilli() { + response.WriteJSON(w, response.Err(403, "Share expired")) + return + } + + next(w, r) + } +} + +func (h *Handler) federationConnect(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + response.WriteJSON(w, response.ErrDefault("Invalid method")) + return + } + + token := extractBearerToken(r) + share, err := h.repo.GetPeerShareByToken(token) + if err != nil || share == nil { + response.WriteJSON(w, response.Err(401, "Unauthorized")) + return + } + + var nodeName string + var serverIP string + var status int + + err = h.repo.DB().QueryRow("SELECT name, server_ip, status FROM node WHERE id = ?", share.NodeID).Scan(&nodeName, &serverIP, &status) + if err != nil { + response.WriteJSON(w, response.Err(-2, "Node not found")) + return + } + + response.WriteJSON(w, response.OK(map[string]interface{}{ + "shareId": share.ID, + "shareName": share.Name, + "nodeId": share.NodeID, + "nodeName": nodeName, + "serverIp": serverIP, + "status": status, + "maxBandwidth": share.MaxBandwidth, + "expiryTime": share.ExpiryTime, + "portRangeStart": share.PortRangeStart, + "portRangeEnd": share.PortRangeEnd, + })) +} + +func (h *Handler) federationTunnelCreate(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + response.WriteJSON(w, response.ErrDefault("Invalid method")) + return + } + + token := extractBearerToken(r) + share, err := h.repo.GetPeerShareByToken(token) + if err != nil || share == nil { + response.WriteJSON(w, response.Err(401, "Unauthorized")) + return + } + + var req federationTunnelRequest + if err := decodeJSON(r.Body, &req); err != nil { + response.WriteJSON(w, response.ErrDefault("Invalid JSON")) + return + } + + if req.RemotePort < share.PortRangeStart || req.RemotePort > share.PortRangeEnd { + response.WriteJSON(w, response.Err(403, "Port out of range")) + return + } + + tunnelType := 1 + if strings.ToLower(req.Protocol) == "udp" { + tunnelType = 2 + } + + tx, err := h.repo.DB().Begin() + if err != nil { + response.WriteJSON(w, response.Err(-2, err.Error())) + return + } + defer tx.Rollback() + + now := time.Now().UnixMilli() + res, err := tx.Exec(`INSERT INTO tunnel (name, type, protocol, flow, created_time, updated_time, status, in_ip) VALUES (?, ?, ?, 0, ?, ?, 1, ?)`, + fmt.Sprintf("Share-%d-Port-%d", share.ID, req.RemotePort), + tunnelType, + req.Protocol, + now, + now, + "", + ) + if err != nil { + response.WriteJSON(w, response.Err(-2, err.Error())) + return + } + + tunnelID, _ := res.LastInsertId() + + _, err = tx.Exec(`INSERT INTO chain_tunnel (tunnel_id, chain_type, node_id, port, strategy, inx, protocol) VALUES (?, 1, ?, ?, 'fifo', 0, ?)`, + tunnelID, + share.NodeID, + req.RemotePort, + req.Protocol, + ) + if err != nil { + response.WriteJSON(w, response.Err(-2, err.Error())) + return + } + + if err := tx.Commit(); err != nil { + response.WriteJSON(w, response.Err(-2, err.Error())) + return + } + + h.wsServer.SendCommand(share.NodeID, "reload", nil, time.Second*5) + + response.WriteJSON(w, response.OK(map[string]interface{}{ + "tunnelId": tunnelID, + })) +} + +func extractBearerToken(r *http.Request) string { + authHeader := r.Header.Get("Authorization") + parts := strings.Split(authHeader, " ") + if len(parts) == 2 && parts[0] == "Bearer" { + return parts[1] + } + return "" +} diff --git a/go-backend/internal/http/handler/handler.go b/go-backend/internal/http/handler/handler.go index 8236cfe..f05b7ce 100644 --- a/go-backend/internal/http/handler/handler.go +++ b/go-backend/internal/http/handler/handler.go @@ -143,6 +143,9 @@ func (h *Handler) Register(mux *http.ServeMux) { mux.HandleFunc("/api/v1/group/permission/assign", h.groupPermissionAssign) mux.HandleFunc("/api/v1/group/permission/remove", h.groupPermissionRemove) mux.HandleFunc("/api/v1/open_api/sub_store", h.openAPISubStore) + mux.HandleFunc("/api/v1/federation/connect", h.authPeer(h.federationConnect)) + mux.HandleFunc("/api/v1/federation/tunnel/create", h.authPeer(h.federationTunnelCreate)) + mux.HandleFunc("/api/v1/federation/node/import", h.nodeImport) mux.HandleFunc("/flow/test", h.flowTest) mux.HandleFunc("/flow/config", h.flowConfig) diff --git a/go-backend/internal/http/handler/mutations.go b/go-backend/internal/http/handler/mutations.go index ab4d701..eba38a2 100644 --- a/go-backend/internal/http/handler/mutations.go +++ b/go-backend/internal/http/handler/mutations.go @@ -15,6 +15,7 @@ import ( "strings" "time" + "go-backend/internal/http/client" "go-backend/internal/http/response" "go-backend/internal/security" ) @@ -273,8 +274,8 @@ func (h *Handler) nodeCreate(w http.ResponseWriter, r *http.Request) { now := time.Now().UnixMilli() inx := nextIndex(db, "node") _, err := db.Exec(` - INSERT INTO node(name, secret, server_ip, server_ip_v4, server_ip_v6, port, interface_name, version, http, tls, socks, created_time, updated_time, status, tcp_listen_addr, udp_listen_addr, inx) - VALUES(?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + INSERT INTO node(name, secret, server_ip, server_ip_v4, server_ip_v6, port, interface_name, version, http, tls, socks, created_time, updated_time, status, tcp_listen_addr, udp_listen_addr, inx, is_remote, remote_url, remote_token, remote_config) + VALUES(?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) `, name, randomToken(16), @@ -293,6 +294,10 @@ func (h *Handler) nodeCreate(w http.ResponseWriter, r *http.Request) { defaultString(asString(req["tcpListenAddr"]), "[::]"), defaultString(asString(req["udpListenAddr"]), "[::]"), inx, + asInt(req["isRemote"], 0), + nullableText(asString(req["remoteUrl"])), + nullableText(asString(req["remoteToken"])), + nullableText(asString(req["remoteConfig"])), ) if err != nil { response.WriteJSON(w, response.Err(-2, err.Error())) @@ -509,6 +514,45 @@ func (h *Handler) tunnelCreate(w http.ResponseWriter, r *http.Request) { inIP = buildTunnelInIP(runtimeState.InNodes, runtimeState.Nodes) } + if len(runtimeState.InNodes) > 0 { + firstNodeID := runtimeState.InNodes[0].NodeID + var isRemote int + var rUrl, rToken sql.NullString + if err := h.repo.DB().QueryRow("SELECT is_remote, remote_url, remote_token FROM node WHERE id = ?", firstNodeID).Scan(&isRemote, &rUrl, &rToken); err == nil && isRemote == 1 { + fc := client.NewFederationClient() + + targetProto := "tcp" + targetPort := 0 + targetAddr := "" + + if typeVal == 1 { + if len(runtimeState.OutNodes) > 0 { + outNode := runtimeState.OutNodes[0] + outNodeRec := runtimeState.Nodes[outNode.NodeID] + targetAddr = processServerAddress(outNodeRec.ServerIP) + if outNode.Port > 0 { + targetAddr = fmt.Sprintf("%s:%d", targetAddr, outNode.Port) + } + } + if len(runtimeState.InNodes) > 0 { + inNodesRaw := asMapSlice(req["inNodeId"]) + if len(inNodesRaw) > 0 { + targetPort = asInt(inNodesRaw[0]["port"], 0) + targetProto = defaultString(asString(inNodesRaw[0]["protocol"]), "tcp") + } + } + + if targetPort > 0 && targetAddr != "" { + _, err := fc.CreateTunnel(rUrl.String, rToken.String, targetProto, targetPort, targetAddr) + if err != nil { + response.WriteJSON(w, response.ErrDefault("Remote tunnel creation failed: "+err.Error())) + return + } + } + } + } + } + res, err := tx.Exec(`INSERT INTO tunnel(name, traffic_ratio, type, protocol, flow, created_time, updated_time, status, in_ip, inx) VALUES(?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, name, trafficRatio, typeVal, "tls", flow, now, now, status, nullableText(inIP), inx) if err != nil { diff --git a/go-backend/internal/store/sqlite/repository.go b/go-backend/internal/store/sqlite/repository.go index 2c81562..187ff97 100644 --- a/go-backend/internal/store/sqlite/repository.go +++ b/go-backend/internal/store/sqlite/repository.go @@ -95,13 +95,32 @@ type StatisticsFlow struct { } type Node struct { - ID int64 - Secret string - Version sql.NullString - HTTP int - TLS int - Socks int - Status int + ID int64 + Secret string + Version sql.NullString + HTTP int + TLS int + Socks int + Status int + IsRemote int + RemoteURL sql.NullString + RemoteToken sql.NullString + RemoteConfig sql.NullString +} + +type PeerShare struct { + ID int64 `json:"id"` + Name string `json:"name"` + NodeID int64 `json:"nodeId"` + Token string `json:"token"` + MaxBandwidth int64 `json:"maxBandwidth"` + ExpiryTime int64 `json:"expiryTime"` + PortRangeStart int `json:"portRangeStart"` + PortRangeEnd int `json:"portRangeEnd"` + CurrentFlow int64 `json:"currentFlow"` + IsActive int `json:"isActive"` + CreatedTime int64 `json:"createdTime"` + UpdatedTime int64 `json:"updatedTime"` } func Open(path string) (*Repository, error) { @@ -124,6 +143,11 @@ func Open(path string) (*Repository, error) { return nil, err } + if err := ensurePeerSchema(db); err != nil { + _ = db.Close() + return nil, err + } + return &Repository{db: db}, nil } @@ -389,9 +413,9 @@ func (r *Repository) GetNodeBySecret(secret string) (*Node, error) { return nil, errors.New("repository not initialized") } - row := r.db.QueryRow(`SELECT id, secret, version, http, tls, socks, status FROM node WHERE secret = ? LIMIT 1`, secret) + row := r.db.QueryRow(`SELECT id, secret, version, http, tls, socks, status, is_remote, remote_url, remote_token, remote_config FROM node WHERE secret = ? LIMIT 1`, secret) var n Node - if err := row.Scan(&n.ID, &n.Secret, &n.Version, &n.HTTP, &n.TLS, &n.Socks, &n.Status); err != nil { + if err := row.Scan(&n.ID, &n.Secret, &n.Version, &n.HTTP, &n.TLS, &n.Socks, &n.Status, &n.IsRemote, &n.RemoteURL, &n.RemoteToken, &n.RemoteConfig); err != nil { if errors.Is(err, sql.ErrNoRows) { return nil, nil } @@ -454,7 +478,7 @@ func (r *Repository) ListNodes() ([]map[string]interface{}, error) { } rows, err := r.db.Query(` - SELECT id, inx, name, server_ip, server_ip_v4, server_ip_v6, port, tcp_listen_addr, udp_listen_addr, version, http, tls, socks, status + SELECT id, inx, name, server_ip, server_ip_v4, server_ip_v6, port, tcp_listen_addr, udp_listen_addr, version, http, tls, socks, status, is_remote, remote_url, remote_token, remote_config FROM node ORDER BY inx ASC, id ASC `) @@ -467,10 +491,10 @@ func (r *Repository) ListNodes() ([]map[string]interface{}, error) { for rows.Next() { var id, inx int64 var name, serverIP, port string - var serverIPV4, serverIPV6, tcpListen, udpListen, version sql.NullString - var httpVal, tlsVal, socksVal, status int + var serverIPV4, serverIPV6, tcpListen, udpListen, version, remoteURL, remoteToken, remoteConfig sql.NullString + var httpVal, tlsVal, socksVal, status, isRemote int - if err := rows.Scan(&id, &inx, &name, &serverIP, &serverIPV4, &serverIPV6, &port, &tcpListen, &udpListen, &version, &httpVal, &tlsVal, &socksVal, &status); err != nil { + if err := rows.Scan(&id, &inx, &name, &serverIP, &serverIPV4, &serverIPV6, &port, &tcpListen, &udpListen, &version, &httpVal, &tlsVal, &socksVal, &status, &isRemote, &remoteURL, &remoteToken, &remoteConfig); err != nil { return nil, err } @@ -490,6 +514,10 @@ func (r *Repository) ListNodes() ([]map[string]interface{}, error) { "tls": tlsVal, "socks": socksVal, "status": status, + "isRemote": isRemote, + "remoteUrl": nullableString(remoteURL), + "remoteToken": nullableString(remoteToken), + "remoteConfig": nullableString(remoteConfig), }) } @@ -1186,6 +1214,132 @@ func bootstrapSchema(db *sql.DB) error { return nil } +func ensurePeerSchema(db *sql.DB) error { + if db == nil { + return errors.New("nil db") + } + + _, err := db.Exec(`CREATE TABLE IF NOT EXISTS peer_share ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + name TEXT NOT NULL, + node_id INTEGER NOT NULL, + token TEXT NOT NULL UNIQUE, + max_bandwidth INTEGER DEFAULT 0, + expiry_time INTEGER DEFAULT 0, + port_range_start INTEGER DEFAULT 0, + port_range_end INTEGER DEFAULT 0, + current_flow INTEGER DEFAULT 0, + is_active INTEGER DEFAULT 1, + created_time INTEGER NOT NULL, + updated_time INTEGER NOT NULL + )`) + if err != nil { + return fmt.Errorf("create peer_share: %w", err) + } + + columns := map[string]string{ + "is_remote": "INTEGER DEFAULT 0", + "remote_url": "TEXT", + "remote_token": "TEXT", + "remote_config": "TEXT", + } + + for col, typ := range columns { + var dummy interface{} + err := db.QueryRow(fmt.Sprintf("SELECT %s FROM node LIMIT 1", col)).Scan(&dummy) + if err != nil { + if strings.Contains(err.Error(), "no such column") { + _, err = db.Exec(fmt.Sprintf("ALTER TABLE node ADD COLUMN %s %s", col, typ)) + if err != nil { + log.Printf("failed to add column %s: %v", col, err) + } + } + } + } + return nil +} + +func (r *Repository) CreatePeerShare(share *PeerShare) error { + if r == nil || r.db == nil { + return errors.New("repository not initialized") + } + _, err := r.db.Exec(` + INSERT INTO peer_share(name, node_id, token, max_bandwidth, expiry_time, port_range_start, port_range_end, current_flow, is_active, created_time, updated_time) + VALUES(?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + `, share.Name, share.NodeID, share.Token, share.MaxBandwidth, share.ExpiryTime, share.PortRangeStart, share.PortRangeEnd, share.CurrentFlow, share.IsActive, share.CreatedTime, share.UpdatedTime) + return err +} + +func (r *Repository) UpdatePeerShare(share *PeerShare) error { + if r == nil || r.db == nil { + return errors.New("repository not initialized") + } + _, err := r.db.Exec(` + UPDATE peer_share SET name=?, max_bandwidth=?, expiry_time=?, port_range_start=?, port_range_end=?, is_active=?, updated_time=? + WHERE id=? + `, share.Name, share.MaxBandwidth, share.ExpiryTime, share.PortRangeStart, share.PortRangeEnd, share.IsActive, share.UpdatedTime, share.ID) + return err +} + +func (r *Repository) DeletePeerShare(id int64) error { + if r == nil || r.db == nil { + return errors.New("repository not initialized") + } + _, err := r.db.Exec(`DELETE FROM peer_share WHERE id=?`, id) + return err +} + +func (r *Repository) GetPeerShare(id int64) (*PeerShare, error) { + if r == nil || r.db == nil { + return nil, errors.New("repository not initialized") + } + row := r.db.QueryRow(`SELECT id, name, node_id, token, max_bandwidth, expiry_time, port_range_start, port_range_end, current_flow, is_active, created_time, updated_time FROM peer_share WHERE id = ?`, id) + var s PeerShare + if err := row.Scan(&s.ID, &s.Name, &s.NodeID, &s.Token, &s.MaxBandwidth, &s.ExpiryTime, &s.PortRangeStart, &s.PortRangeEnd, &s.CurrentFlow, &s.IsActive, &s.CreatedTime, &s.UpdatedTime); err != nil { + if errors.Is(err, sql.ErrNoRows) { + return nil, nil + } + return nil, err + } + return &s, nil +} + +func (r *Repository) GetPeerShareByToken(token string) (*PeerShare, error) { + if r == nil || r.db == nil { + return nil, errors.New("repository not initialized") + } + row := r.db.QueryRow(`SELECT id, name, node_id, token, max_bandwidth, expiry_time, port_range_start, port_range_end, current_flow, is_active, created_time, updated_time FROM peer_share WHERE token = ?`, token) + var s PeerShare + if err := row.Scan(&s.ID, &s.Name, &s.NodeID, &s.Token, &s.MaxBandwidth, &s.ExpiryTime, &s.PortRangeStart, &s.PortRangeEnd, &s.CurrentFlow, &s.IsActive, &s.CreatedTime, &s.UpdatedTime); err != nil { + if errors.Is(err, sql.ErrNoRows) { + return nil, nil + } + return nil, err + } + return &s, nil +} + +func (r *Repository) ListPeerShares() ([]PeerShare, error) { + if r == nil || r.db == nil { + return nil, errors.New("repository not initialized") + } + rows, err := r.db.Query(`SELECT id, name, node_id, token, max_bandwidth, expiry_time, port_range_start, port_range_end, current_flow, is_active, created_time, updated_time FROM peer_share ORDER BY id DESC`) + if err != nil { + return nil, err + } + defer rows.Close() + + var shares []PeerShare + for rows.Next() { + var s PeerShare + if err := rows.Scan(&s.ID, &s.Name, &s.NodeID, &s.Token, &s.MaxBandwidth, &s.ExpiryTime, &s.PortRangeStart, &s.PortRangeEnd, &s.CurrentFlow, &s.IsActive, &s.CreatedTime, &s.UpdatedTime); err != nil { + return nil, err + } + shares = append(shares, s) + } + return shares, nil +} + var osMkdirAll = func(path string) error { return os.MkdirAll(path, 0o755) } diff --git a/vite-frontend/src/App.tsx b/vite-frontend/src/App.tsx index ab3c8a8..d7c1528 100644 --- a/vite-frontend/src/App.tsx +++ b/vite-frontend/src/App.tsx @@ -12,6 +12,7 @@ import GroupPage from "@/pages/group"; import ProfilePage from "@/pages/profile"; import LimitPage from "@/pages/limit"; import ConfigPage from "@/pages/config"; +import PanelSharingPage from "@/pages/panel-sharing"; import { SettingsPage } from "@/pages/settings"; import AdminLayout from "@/layouts/admin"; import H5Layout from "@/layouts/h5"; @@ -235,12 +236,20 @@ function App() { /> + } path="/config" /> + + + + } + path="/panel-sharing" + /> } path="/settings" /> ); diff --git a/vite-frontend/src/api/index.ts b/vite-frontend/src/api/index.ts index 4d972de..59b0ffc 100644 --- a/vite-frontend/src/api/index.ts +++ b/vite-frontend/src/api/index.ts @@ -187,3 +187,20 @@ export const assignGroupPermission = (data: { }) => Network.post("/group/permission/assign", data); export const removeGroupPermission = (id: number) => Network.post("/group/permission/remove", { id }); + +// 面板共享 (Federation) 接口 +export const getPeerShareList = () => Network.post("/federation/share/list"); +export const createPeerShare = (data: { + name: string; + nodeId: number; + maxBandwidth?: number; + expiryTime?: number; + portRangeStart?: number; + portRangeEnd?: number; +}) => Network.post("/federation/share/create", data); +export const deletePeerShare = (id: number) => + Network.post("/federation/share/delete", { id }); +export const importRemoteNode = (data: { + remoteUrl: string; + token: string; +}) => Network.post("/federation/node/import", data); diff --git a/vite-frontend/src/layouts/admin.tsx b/vite-frontend/src/layouts/admin.tsx index 8880ad1..6837f9c 100644 --- a/vite-frontend/src/layouts/admin.tsx +++ b/vite-frontend/src/layouts/admin.tsx @@ -144,6 +144,16 @@ export default function AdminLayout({ ), adminOnly: true, }, + { + path: "/panel-sharing", + label: "共享", + icon: ( + + + + ), + adminOnly: true, + }, { path: "/config", label: "设置", diff --git a/vite-frontend/src/pages/node.tsx b/vite-frontend/src/pages/node.tsx index 3bde440..9ffe687 100644 --- a/vite-frontend/src/pages/node.tsx +++ b/vite-frontend/src/pages/node.tsx @@ -62,7 +62,9 @@ interface Node { http?: number; // 0 关 1 开 tls?: number; // 0 关 1 开 socks?: number; // 0 关 1 开 - status: number; // 1: 在线, 0: 离线 + status: number; + isRemote?: number; + remoteUrl?: string; connectionStatus: "online" | "offline"; systemInfo?: { cpuUsage: number; @@ -1133,6 +1135,16 @@ export default function NodePage() {
+ {node.isRemote === 1 && ( + + 远程 + + )}
handleEdit(node)} diff --git a/vite-frontend/src/pages/panel-sharing.tsx b/vite-frontend/src/pages/panel-sharing.tsx new file mode 100644 index 0000000..4f5f5ef --- /dev/null +++ b/vite-frontend/src/pages/panel-sharing.tsx @@ -0,0 +1,313 @@ +import React, { useState, useEffect } from "react"; +import { Button } from "@heroui/button"; +import { Card, CardBody, CardHeader } from "@heroui/card"; +import { Tabs, Tab } from "@heroui/tabs"; +import { Input } from "@heroui/input"; +import { + Modal, + ModalContent, + ModalHeader, + ModalBody, + ModalFooter, +} from "@heroui/modal"; +import { Select, SelectItem } from "@heroui/select"; +import { toast } from "react-hot-toast"; +import { + getNodeList, + createPeerShare, + getPeerShareList, + deletePeerShare, + importRemoteNode, +} from "@/api"; + +interface Node { + id: number; + name: string; +} + +interface PeerShare { + id: number; + name: string; + token: string; + maxBandwidth: number; + expiryTime: number; + portRangeStart: number; + portRangeEnd: number; + isActive: number; +} + +export default function PanelSharingPage() { + const [selectedTab, setSelectedTab] = useState("my-shares"); + const [shares, setShares] = useState([]); + const [nodes, setNodes] = useState([]); + const [loading, setLoading] = useState(false); + + // Modals + const [createShareOpen, setCreateShareOpen] = useState(false); + const [importNodeOpen, setImportNodeOpen] = useState(false); + + // Forms + const [shareForm, setShareForm] = useState({ + name: "", + nodeId: "", + maxBandwidth: 0, + expiryDays: 30, + portRangeStart: 10000, + portRangeEnd: 20000, + }); + + const [importForm, setImportForm] = useState({ + remoteUrl: "", + token: "", + }); + + useEffect(() => { + if (selectedTab === "my-shares") { + loadShares(); + loadNodes(); + } + }, [selectedTab]); + + const loadShares = async () => { + setLoading(true); + try { + const res = await getPeerShareList(); + if (res.code === 0) { + setShares(res.data || []); + } else { + toast.error(res.msg || "加载分享列表失败"); + } + } finally { + setLoading(false); + } + }; + + const loadNodes = async () => { + try { + const res = await getNodeList(); + if (res.code === 0) { + setNodes(res.data || []); + } + } catch { + // ignore + } + }; + + const handleCreateShare = async () => { + if (!shareForm.name || !shareForm.nodeId) { + toast.error("请填写必要信息"); + return; + } + try { + const expiryTime = + Date.now() + shareForm.expiryDays * 24 * 60 * 60 * 1000; + const res = await createPeerShare({ + name: shareForm.name, + nodeId: parseInt(shareForm.nodeId), + maxBandwidth: shareForm.maxBandwidth * 1024 * 1024 * 1024, + expiryTime: shareForm.expiryDays === 0 ? 0 : expiryTime, + portRangeStart: shareForm.portRangeStart, + portRangeEnd: shareForm.portRangeEnd, + }); + if (res.code === 0) { + toast.success("创建成功"); + setCreateShareOpen(false); + loadShares(); + } else { + toast.error(res.msg || "创建失败"); + } + } catch { + toast.error("网络错误"); + } + }; + + const handleDeleteShare = async (id: number) => { + try { + const res = await deletePeerShare(id); + if (res.code === 0) { + toast.success("删除成功"); + loadShares(); + } else { + toast.error(res.msg || "删除失败"); + } + } catch { + toast.error("网络错误"); + } + }; + + const handleImportNode = async () => { + if (!importForm.remoteUrl || !importForm.token) { + toast.error("请填写完整信息"); + return; + } + try { + // Automatically add http/https if missing + let url = importForm.remoteUrl.trim(); + if (!url.startsWith("http")) { + url = "http://" + url; + } + + const res = await importRemoteNode({ + remoteUrl: url, + token: importForm.token.trim(), + }); + if (res.code === 0) { + toast.success("导入成功,请前往节点列表查看"); + setImportNodeOpen(false); + setImportForm({ remoteUrl: "", token: "" }); + } else { + toast.error(res.msg || "导入失败"); + } + } catch { + toast.error("网络错误"); + } + }; + + const copyToken = (token: string) => { + navigator.clipboard.writeText(token); + toast.success("Token已复制"); + }; + + return ( +
+
+

面板共享 (Panel Peering)

+
+ + setSelectedTab(k as string)} + > + + + +
+ +
+ + {loading ? ( +
加载中...
+ ) : shares.length === 0 ? ( +
暂无分享
+ ) : ( +
+ {shares.map((share) => ( + + +

{share.name}

+ +
+ +

端口范围: {share.portRangeStart} - {share.portRangeEnd}

+

过期时间: {share.expiryTime === 0 ? "永久" : new Date(share.expiryTime).toLocaleDateString()}

+
+ + +
+
+
+ ))} +
+ )} +
+
+
+ + + +
+ +
+
+

已导入的节点将显示在“节点管理”页面,带有“远程”标记。

+

请使用其创建隧道。

+
+
+
+
+
+ + {/* Create Share Modal */} + setCreateShareOpen(false)}> + + 创建分享 + + setShareForm({ ...shareForm, name: e.target.value })} + /> + +
+ setShareForm({ ...shareForm, portRangeStart: parseInt(e.target.value) })} + /> + setShareForm({ ...shareForm, portRangeEnd: parseInt(e.target.value) })} + /> +
+ setShareForm({ ...shareForm, expiryDays: parseInt(e.target.value) })} + /> +
+ + + + +
+
+ + {/* Import Node Modal */} + setImportNodeOpen(false)}> + + 导入远程节点 + + setImportForm({ ...importForm, remoteUrl: e.target.value })} + /> + setImportForm({ ...importForm, token: e.target.value })} + /> + + + + + + + +
+ ); +} \ No newline at end of file