diff --git a/go-backend/internal/http/handler/federation.go b/go-backend/internal/http/handler/federation.go index 97c7228..fe2cf57 100644 --- a/go-backend/internal/http/handler/federation.go +++ b/go-backend/internal/http/handler/federation.go @@ -477,9 +477,14 @@ func (h *Handler) federationRemoteUsageList(w http.ResponseWriter, r *http.Reque response.WriteJSON(w, response.Err(-2, err.Error())) return } + forwardPortRows, err := h.repo.ListActiveForwardPortsForNode(nodeID) + if err != nil { + response.WriteJSON(w, response.Err(-2, err.Error())) + return + } usedSet := make(map[int]struct{}) - bindings := make([]remoteUsageBindingItem, 0, len(bindingRows)) + bindings := make([]remoteUsageBindingItem, 0, len(bindingRows)+len(forwardPortRows)) for _, b := range bindingRows { bindings = append(bindings, remoteUsageBindingItem{ BindingID: b.ID, @@ -496,6 +501,29 @@ func (h *Handler) federationRemoteUsageList(w http.ResponseWriter, r *http.Reque usedSet[b.AllocatedPort] = struct{}{} } } + for _, fp := range forwardPortRows { + bindings = append(bindings, remoteUsageBindingItem{ + BindingID: -fp.ForwardID, + TunnelID: fp.TunnelID, + TunnelName: fp.TunnelName, + ChainType: 1, + HopInx: 0, + AllocatedPort: fp.Port, + ResourceKey: fmt.Sprintf("forward:%d", fp.ForwardID), + RemoteBindingID: "", + UpdatedTime: fp.UpdatedTime, + }) + if fp.Port > 0 { + usedSet[fp.Port] = struct{}{} + } + } + + sort.Slice(bindings, func(i, j int) bool { + if bindings[i].AllocatedPort == bindings[j].AllocatedPort { + return bindings[i].BindingID < bindings[j].BindingID + } + return bindings[i].AllocatedPort < bindings[j].AllocatedPort + }) usedPorts := make([]int, 0, len(usedSet)) for port := range usedSet { diff --git a/go-backend/internal/http/handler/federation_share_test.go b/go-backend/internal/http/handler/federation_share_test.go index dfbccc0..c0a054a 100644 --- a/go-backend/internal/http/handler/federation_share_test.go +++ b/go-backend/internal/http/handler/federation_share_test.go @@ -941,6 +941,110 @@ func TestFederationRemoteUsageList(t *testing.T) { } } +func TestFederationRemoteUsageListIncludesForwardPorts(t *testing.T) { + r, err := repo.Open(filepath.Join(t.TempDir(), "panel-forward-usage.db")) + if err != nil { + t.Fatalf("open sqlite: %v", err) + } + t.Cleanup(func() { _ = r.Close() }) + + h := New(r, "test-jwt-secret") + now := time.Now().UnixMilli() + + if err := r.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(?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + `, "forward-usage-remote-node", "forward-usage-secret", "10.60.70.80", "10.60.70.80", "", "33000-33010", "", "v1", 1, 1, 1, now, now, 1, "[::]", "[::]", 0, 1, "", "", `{"shareId":99,"maxBandwidth":0,"currentFlow":0,"portRangeStart":33000,"portRangeEnd":33010}`).Error; err != nil { + t.Fatalf("insert remote node: %v", err) + } + + var nodeID int64 + if err := r.DB().Raw(`SELECT id FROM node WHERE name = ? ORDER BY id DESC LIMIT 1`, "forward-usage-remote-node").Row().Scan(&nodeID); err != nil { + t.Fatalf("query node id: %v", err) + } + + if err := r.DB().Exec(` + INSERT INTO tunnel(name, type, protocol, flow, created_time, updated_time, status, in_ip, inx) + VALUES(?, ?, ?, ?, ?, ?, ?, ?, ?) + `, "forward-usage-tunnel", 1, "tls", 1, now, now, 1, "", 0).Error; err != nil { + t.Fatalf("insert tunnel: %v", err) + } + + var tunnelID int64 + if err := r.DB().Raw(`SELECT id FROM tunnel WHERE name = ? ORDER BY id DESC LIMIT 1`, "forward-usage-tunnel").Row().Scan(&tunnelID); err != nil { + t.Fatalf("query tunnel id: %v", err) + } + + if err := r.DB().Exec(` + INSERT INTO forward(user_id, user_name, name, tunnel_id, remote_addr, strategy, in_flow, out_flow, created_time, updated_time, status, inx) + VALUES(?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + `, 1, "tester", "forward-usage-item", tunnelID, "1.1.1.1:443", "fifo", 0, 0, now, now, 1, 0).Error; err != nil { + t.Fatalf("insert forward: %v", err) + } + + var forwardID int64 + if err := r.DB().Raw(`SELECT id FROM forward WHERE name = ? ORDER BY id DESC LIMIT 1`, "forward-usage-item").Row().Scan(&forwardID); err != nil { + t.Fatalf("query forward id: %v", err) + } + + if err := r.DB().Exec(`INSERT INTO forward_port(forward_id, node_id, port) VALUES(?, ?, ?)`, forwardID, nodeID, 33001).Error; err != nil { + t.Fatalf("insert forward_port: %v", err) + } + + req := httptest.NewRequest(http.MethodPost, "/api/v1/federation/share/remote-usage/list", nil) + res := httptest.NewRecorder() + h.federationRemoteUsageList(res, req) + + if res.Code != http.StatusOK { + t.Fatalf("expected status %d, got %d", http.StatusOK, res.Code) + } + + var payload response.R + if err := json.NewDecoder(res.Body).Decode(&payload); err != nil { + t.Fatalf("decode response: %v", err) + } + if payload.Code != 0 { + t.Fatalf("expected response code 0, got %d (%s)", payload.Code, payload.Msg) + } + + rows, ok := payload.Data.([]interface{}) + if !ok || len(rows) == 0 { + t.Fatalf("expected non-empty usage list, got %T", payload.Data) + } + + first, ok := rows[0].(map[string]interface{}) + if !ok { + t.Fatalf("expected usage row map, got %T", rows[0]) + } + + usedPortsRaw, ok := first["usedPorts"].([]interface{}) + if !ok { + t.Fatalf("expected usedPorts array, got %T", first["usedPorts"]) + } + if len(usedPortsRaw) != 1 || int(usedPortsRaw[0].(float64)) != 33001 { + t.Fatalf("expected usedPorts [33001], got %v", usedPortsRaw) + } + + bindingsRaw, ok := first["bindings"].([]interface{}) + if !ok { + t.Fatalf("expected bindings array, got %T", first["bindings"]) + } + if len(bindingsRaw) != 1 { + t.Fatalf("expected 1 binding row from forward usage, got %d", len(bindingsRaw)) + } + + binding, ok := bindingsRaw[0].(map[string]interface{}) + if !ok { + t.Fatalf("expected binding row object, got %T", bindingsRaw[0]) + } + if int(binding["allocatedPort"].(float64)) != 33001 { + t.Fatalf("expected allocatedPort=33001, got %v", binding["allocatedPort"]) + } + if int(binding["chainType"].(float64)) != 1 { + t.Fatalf("expected chainType=1 for forward usage row, got %v", binding["chainType"]) + } +} + func TestAuthPeerAllowedIPs(t *testing.T) { r, err := repo.Open(filepath.Join(t.TempDir(), "panel.db")) if err != nil { diff --git a/go-backend/internal/store/repo/repository_federation.go b/go-backend/internal/store/repo/repository_federation.go index c3f34ab..438e7c0 100644 --- a/go-backend/internal/store/repo/repository_federation.go +++ b/go-backend/internal/store/repo/repository_federation.go @@ -38,6 +38,14 @@ type FederationBindingRow struct { UpdatedTime int64 } +type ActiveForwardPortRow struct { + ForwardID int64 + TunnelID int64 + TunnelName string + Port int + UpdatedTime int64 +} + // ListRemoteNodes returns all nodes with is_remote=1, ordered by id desc. func (r *Repository) ListRemoteNodes() ([]RemoteNodeRow, error) { if r == nil || r.db == nil { @@ -87,6 +95,27 @@ func (r *Repository) ListActiveBindingsForNode(nodeID int64) ([]FederationBindin return result, nil } +func (r *Repository) ListActiveForwardPortsForNode(nodeID int64) ([]ActiveForwardPortRow, error) { + if r == nil || r.db == nil { + return nil, errors.New("repository not initialized") + } + var result []ActiveForwardPortRow + err := r.db.Model(&model.ForwardPort{}). + Select("forward_port.forward_id, forward.tunnel_id, COALESCE(tunnel.name, '') AS tunnel_name, forward_port.port, forward.updated_time"). + Joins("JOIN forward ON forward.id = forward_port.forward_id"). + Joins("LEFT JOIN tunnel ON tunnel.id = forward.tunnel_id"). + Where("forward_port.node_id = ? AND forward_port.port > 0", nodeID). + Order("forward_port.port ASC, forward_port.id ASC"). + Find(&result).Error + if err != nil { + return nil, err + } + if result == nil { + result = make([]ActiveForwardPortRow, 0) + } + return result, nil +} + // GetNodeBasicInfo returns the name, server_ip, and status for a given node. func (r *Repository) GetNodeBasicInfo(nodeID int64) (*NodeBasicInfo, error) { if r == nil || r.db == nil { diff --git a/vite-frontend/src/pages/forward.tsx b/vite-frontend/src/pages/forward.tsx index e5aa696..5df33b7 100644 --- a/vite-frontend/src/pages/forward.tsx +++ b/vite-frontend/src/pages/forward.tsx @@ -50,6 +50,7 @@ import { Checkbox } from "@/shadcn-bridge/heroui/checkbox"; import { createForward, getForwardList, + getPeerRemoteUsageList, updateForward, deleteForward, forceDeleteForward, @@ -98,6 +99,7 @@ interface Forward { inFlow: number; outFlow: number; serviceRunning: boolean; + federationShareFlow?: number; createdTime: string; userName?: string; userId?: number; @@ -219,6 +221,119 @@ export default function ForwardPage() { ); const [batchLoading, setBatchLoading] = useState(false); + const parseShareIdFromTunnelName = (tunnelName: string): number | null => { + const normalized = (tunnelName || "").trim(); + + if (!normalized.startsWith("Share-")) { + return null; + } + + const raw = normalized.slice("Share-".length); + const idx = raw.indexOf("-Port-"); + + if (idx <= 0) { + return null; + } + + const shareId = Number(raw.slice(0, idx).trim()); + + return Number.isFinite(shareId) && shareId > 0 ? shareId : null; + }; + + const mergeFederationShareFlow = async ( + forwardsData: Forward[], + ): Promise => { + const shareIds = new Set(); + + forwardsData.forEach((forward) => { + const shareId = parseShareIdFromTunnelName(forward.tunnelName || ""); + + if (shareId) { + shareIds.add(shareId); + } + }); + + if (shareIds.size === 0) { + return forwardsData; + } + + try { + const usageRes = await getPeerRemoteUsageList(); + + if (usageRes.code !== 0 || !Array.isArray(usageRes.data)) { + return forwardsData; + } + + const flowByShare = new Map(); + + usageRes.data.forEach((item: Record) => { + const shareId = Number(item.shareId || 0); + const currentFlow = Number(item.currentFlow || 0); + + if ( + Number.isFinite(shareId) && + shareId > 0 && + Number.isFinite(currentFlow) && + currentFlow > 0 + ) { + flowByShare.set(shareId, currentFlow); + } + }); + + const forwardCountByShare = new Map(); + + forwardsData.forEach((forward) => { + const shareId = parseShareIdFromTunnelName(forward.tunnelName || ""); + + if (!shareId || !flowByShare.has(shareId)) { + return; + } + + forwardCountByShare.set( + shareId, + (forwardCountByShare.get(shareId) || 0) + 1, + ); + }); + + return forwardsData.map((forward) => { + const shareId = parseShareIdFromTunnelName(forward.tunnelName || ""); + + if (!shareId) { + return { ...forward, federationShareFlow: undefined }; + } + + const shareFlow = flowByShare.get(shareId) || 0; + + if (shareFlow <= 0) { + return { ...forward, federationShareFlow: undefined }; + } + + const directFlow = (forward.inFlow || 0) + (forward.outFlow || 0); + + if (directFlow > 0) { + return { ...forward, federationShareFlow: undefined }; + } + + const count = forwardCountByShare.get(shareId) || 1; + const estimated = Math.max(1, Math.floor(shareFlow / count)); + + return { ...forward, federationShareFlow: estimated }; + }); + } catch { + return forwardsData; + } + }; + + const getForwardDisplayFlow = (forward: Forward): number => { + const directFlow = (forward.inFlow || 0) + (forward.outFlow || 0); + + if (directFlow > 0) { + return directFlow; + } + + return forward.federationShareFlow || 0; + }; + useEffect(() => { loadData(); }, []); @@ -249,12 +364,14 @@ export default function ForwardPage() { serviceRunning: forward.status === 1, })) || []; - setForwards(forwardsData); + const mergedForwards = await mergeFederationShareFlow(forwardsData); + + setForwards(mergedForwards); // 初始化拖拽排序顺序 const currentUserId = JwtUtil.getUserIdFromToken(); const { order, fromDatabase } = buildForwardOrder( - forwardsData, + mergedForwards, currentUserId, ); @@ -1359,7 +1476,7 @@ export default function ForwardPage() { - {formatFlow((forward.inFlow || 0) + (forward.outFlow || 0))} + {formatFlow(getForwardDisplayFlow(forward))} @@ -1606,24 +1723,46 @@ export default function ForwardPage() { > {strategyDisplay.text} -
+ {(forward.inFlow || 0) + (forward.outFlow || 0) > 0 ? ( + <> +
+ + ↑{formatFlow(forward.inFlow || 0)} + +
+ + ↓{formatFlow(forward.outFlow || 0)} + + + ) : (forward.federationShareFlow || 0) > 0 ? ( - ↑{formatFlow(forward.inFlow || 0)} + 共享 {formatFlow(forward.federationShareFlow || 0)} -
- - ↓{formatFlow(forward.outFlow || 0)} - + ) : ( + + 总流量 {formatFlow(0)} + + )} diff --git a/vite-frontend/src/pages/panel-sharing.tsx b/vite-frontend/src/pages/panel-sharing.tsx index dd253bd..b191ee5 100644 --- a/vite-frontend/src/pages/panel-sharing.tsx +++ b/vite-frontend/src/pages/panel-sharing.tsx @@ -375,6 +375,9 @@ export default function PanelSharingPage() { }; const formatChainType = (chainType: number, hopInx: number) => { + if (chainType === 1) { + return "入口节点"; + } if (chainType === 2) { return `中继跳点 #${hopInx}`; }