Compare commits

...

1 Commits

Author SHA1 Message Date
sagit 6189fe23f1 feat: include forward ports in federation remote usage and display share flow (#209)
* feat(backend): include forward ports in federation remote usage list

Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-opencode)

Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>

* feat(frontend): display federation share flow in forward list

Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-opencode)

Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>

---------

Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
2026-02-25 16:59:39 +08:00
5 changed files with 319 additions and 16 deletions
+29 -1
View File
@@ -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 {
@@ -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 {
@@ -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 {
+154 -15
View File
@@ -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<Forward[]> => {
const shareIds = new Set<number>();
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<number, number>();
usageRes.data.forEach((item: Record<string, unknown>) => {
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<number, number>();
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() {
</TableCell>
<TableCell className="whitespace-nowrap">
<span className="text-sm font-medium text-default-600 font-mono">
{formatFlow((forward.inFlow || 0) + (forward.outFlow || 0))}
{formatFlow(getForwardDisplayFlow(forward))}
</span>
</TableCell>
<TableCell>
@@ -1606,24 +1723,46 @@ export default function ForwardPage() {
>
{strategyDisplay.text}
</Chip>
<div className="flex items-center gap-1">
{(forward.inFlow || 0) + (forward.outFlow || 0) > 0 ? (
<>
<div className="flex items-center gap-1">
<Chip
className="text-xs whitespace-nowrap"
color="primary"
size="sm"
variant="flat"
>
↑{formatFlow(forward.inFlow || 0)}
</Chip>
</div>
<Chip
className="text-xs whitespace-nowrap"
color="success"
size="sm"
variant="flat"
>
↓{formatFlow(forward.outFlow || 0)}
</Chip>
</>
) : (forward.federationShareFlow || 0) > 0 ? (
<Chip
className="text-xs whitespace-nowrap"
color="primary"
color="secondary"
size="sm"
variant="flat"
>
↑{formatFlow(forward.inFlow || 0)}
共享 {formatFlow(forward.federationShareFlow || 0)}
</Chip>
</div>
<Chip
className="text-xs whitespace-nowrap"
color="success"
size="sm"
variant="flat"
>
↓{formatFlow(forward.outFlow || 0)}
</Chip>
) : (
<Chip
className="text-xs whitespace-nowrap"
color="default"
size="sm"
variant="flat"
>
总流量 {formatFlow(0)}
</Chip>
)}
</div>
</div>
@@ -375,6 +375,9 @@ export default function PanelSharingPage() {
};
const formatChainType = (chainType: number, hopInx: number) => {
if (chainType === 1) {
return "入口节点";
}
if (chainType === 2) {
return `中继跳点 #${hopInx}`;
}