Compare commits

...

2 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
sagit 7ba90e8696 fix(backend): resolve federation forward traffic stats and listener disappearance (#208)
* fix(backend): add repository methods for federation forward runtime management

- GetActiveForwardPeerShareRuntimeByServiceName: lookup runtime by share_id and service_name
- MarkForwardPeerShareRuntimeReleasedByServiceName: release runtime by service_name
- ListActiveForwardPeerShareRuntimesByNodeAndServiceName: node-scoped query for flow processing

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

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

* fix(backend): bind and release federation forward runtimes on service commands

- bindPeerShareForwardRuntimeServices: create runtime if missing, update ServiceName/Port/Applied/Status
- releasePeerShareForwardRuntimeServices: handle deleteservice command to mark runtime released
- parseFederationForwardServiceNamesForRelease: extract service names from delete payload
- Tests: bind creates runtime when missing, release marks runtime as released

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

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

* fix(backend): scope federation flow lookup by node to avoid cross-share collisions

- flowUpload: use GetNodeBySecret to extract nodeID for flow processing
- processFlowItem: accept nodeID parameter and pass to flow handlers
- processPeerShareFlowByServiceName: try node-scoped query first, fallback to global
- Add warning log when multiple runtimes match (ambiguous)
- Tests: update all processFlowItem calls with nodeID parameter

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

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

* test(contract): adjust federation dual panel contract expectations

Update assertion for entry share runtime binding behavior after fix

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 14:10:24 +08:00
10 changed files with 647 additions and 41 deletions
+126 -4
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 {
@@ -1268,6 +1296,8 @@ func (h *Handler) federationRuntimeCommand(w http.ResponseWriter, r *http.Reques
}
if strings.EqualFold(cmd, "addservice") || strings.EqualFold(cmd, "updateservice") {
h.bindPeerShareForwardRuntimeServices(share, req.Data)
} else if strings.EqualFold(cmd, "deleteservice") {
h.releasePeerShareForwardRuntimeServices(share, req.Data)
}
response.WriteJSON(w, response.OK(res))
}
@@ -1326,6 +1356,45 @@ func parseFederationForwardServiceBindings(data interface{}) []federationForward
return bindings
}
func parseFederationForwardServiceNamesForRelease(data interface{}) []string {
names := make(map[string]struct{})
appendName := func(raw string) {
name := normalizeForwardRuntimeServiceName(raw)
if name == "" {
return
}
if _, _, _, ok := parseFlowServiceIDs(name); !ok {
return
}
names[name] = struct{}{}
}
for _, svcMap := range extractFederationServiceEntries(data) {
appendName(asString(svcMap["name"]))
}
if dataMap, ok := data.(map[string]interface{}); ok {
for _, item := range asAnySlice(dataMap["services"]) {
appendName(asString(item))
}
}
for _, item := range asAnySlice(data) {
appendName(asString(item))
}
if len(names) == 0 {
return nil
}
out := make([]string, 0, len(names))
for name := range names {
out = append(out, name)
}
sort.Strings(out)
return out
}
func (h *Handler) bindPeerShareForwardRuntimeServices(share *repo.PeerShare, data interface{}) {
if h == nil || h.repo == nil || share == nil {
return
@@ -1338,13 +1407,66 @@ func (h *Handler) bindPeerShareForwardRuntimeServices(share *repo.PeerShare, dat
now := time.Now().UnixMilli()
for _, binding := range bindings {
runtime, err := h.repo.GetActiveForwardPeerShareRuntimeByPort(share.ID, binding.Port)
if err != nil || runtime == nil || runtime.Status != 1 {
if err != nil {
continue
}
if runtime.ServiceName == binding.Name && runtime.Applied == 1 {
if runtime == nil {
runtime, err = h.repo.GetActiveForwardPeerShareRuntimeByServiceName(share.ID, binding.Name)
if err != nil {
continue
}
}
if runtime == nil {
_ = h.repo.CreatePeerShareRuntime(&repo.PeerShareRuntime{
ShareID: share.ID,
NodeID: share.NodeID,
ReservationID: randomToken(24),
ResourceKey: fmt.Sprintf("forward-runtime:%d:%s:%d:%s", share.ID, binding.Name, binding.Port, randomToken(8)),
BindingID: "",
Role: "forward",
ChainName: "",
ServiceName: binding.Name,
Protocol: "tcp",
Strategy: "fifo",
Port: binding.Port,
Target: "",
Applied: 1,
Status: 1,
CreatedTime: now,
UpdatedTime: now,
})
continue
}
_ = h.repo.UpdatePeerShareRuntimeServiceName(runtime.ID, binding.Name, now)
if runtime.ServiceName == binding.Name && runtime.Applied == 1 && runtime.Port == binding.Port && runtime.Status == 1 {
continue
}
runtime.ServiceName = binding.Name
runtime.Port = binding.Port
runtime.Applied = 1
runtime.Status = 1
runtime.UpdatedTime = now
if strings.TrimSpace(runtime.Protocol) == "" {
runtime.Protocol = "tcp"
}
if strings.TrimSpace(runtime.Strategy) == "" {
runtime.Strategy = "fifo"
}
_ = h.repo.UpdatePeerShareRuntime(runtime)
}
}
func (h *Handler) releasePeerShareForwardRuntimeServices(share *repo.PeerShare, data interface{}) {
if h == nil || h.repo == nil || share == nil {
return
}
names := parseFederationForwardServiceNamesForRelease(data)
if len(names) == 0 {
return
}
now := time.Now().UnixMilli()
for _, name := range names {
_ = h.repo.MarkForwardPeerShareRuntimeReleasedByServiceName(share.ID, name, now)
}
}
@@ -722,6 +722,123 @@ func TestBindPeerShareForwardRuntimeServicesAcceptsTopLevelServiceArray(t *testi
}
}
func TestBindPeerShareForwardRuntimeServicesCreatesRuntimeWhenMissing(t *testing.T) {
r, err := repo.Open(filepath.Join(t.TempDir(), "panel-bind-create-runtime.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.CreatePeerShare(&repo.PeerShare{
Name: "bind-create-runtime-share",
NodeID: 1,
Token: "bind-create-runtime-token",
MaxBandwidth: 0,
CurrentFlow: 0,
PortRangeStart: 26300,
PortRangeEnd: 26320,
IsActive: 1,
CreatedTime: now,
UpdatedTime: now,
}); err != nil {
t.Fatalf("create share: %v", err)
}
share, err := r.GetPeerShareByToken("bind-create-runtime-token")
if err != nil || share == nil {
t.Fatalf("load share: %v", err)
}
h.bindPeerShareForwardRuntimeServices(share, map[string]interface{}{
"services": []interface{}{
map[string]interface{}{"name": "55_2_10_tcp", "addr": "[::]:26301"},
},
})
var count int64
if err := r.DB().Raw(`SELECT COUNT(1) FROM peer_share_runtime WHERE share_id = ? AND role = ? AND status = 1`, share.ID, "forward").Scan(&count).Error; err != nil {
t.Fatalf("query runtime count: %v", err)
}
if count != 1 {
t.Fatalf("expected 1 active forward runtime row, got %d", count)
}
var serviceName string
var port int
var applied int
if err := r.DB().Raw(`SELECT service_name, port, applied FROM peer_share_runtime WHERE share_id = ? AND role = ? ORDER BY id DESC LIMIT 1`, share.ID, "forward").Row().Scan(&serviceName, &port, &applied); err != nil {
t.Fatalf("query created runtime: %v", err)
}
if serviceName != "55_2_10" {
t.Fatalf("expected service_name=55_2_10, got %q", serviceName)
}
if port != 26301 {
t.Fatalf("expected port=26301, got %d", port)
}
if applied != 1 {
t.Fatalf("expected applied=1, got %d", applied)
}
}
func TestReleasePeerShareForwardRuntimeServicesMarksRuntimeReleased(t *testing.T) {
r, err := repo.Open(filepath.Join(t.TempDir(), "panel-release-runtime.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.CreatePeerShare(&repo.PeerShare{
Name: "release-runtime-share",
NodeID: 1,
Token: "release-runtime-token",
MaxBandwidth: 0,
CurrentFlow: 0,
PortRangeStart: 26400,
PortRangeEnd: 26420,
IsActive: 1,
CreatedTime: now,
UpdatedTime: now,
}); err != nil {
t.Fatalf("create share: %v", err)
}
share, err := r.GetPeerShareByToken("release-runtime-token")
if err != nil || share == nil {
t.Fatalf("load share: %v", err)
}
if err := r.DB().Exec(`
INSERT INTO peer_share_runtime(share_id, node_id, reservation_id, resource_key, binding_id, role, chain_name, service_name, protocol, strategy, port, target, applied, status, created_time, updated_time)
VALUES(?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
`, share.ID, share.NodeID, "release-r1", "release-rk1", "", "forward", "", "77_2_10", "tcp", "fifo", 26401, "", 1, 1, now, now).Error; err != nil {
t.Fatalf("insert runtime: %v", err)
}
h.releasePeerShareForwardRuntimeServices(share, map[string]interface{}{
"services": []interface{}{"77_2_10_tcp"},
})
var status int
var applied int
var serviceName string
if err := r.DB().Raw(`SELECT status, applied, service_name FROM peer_share_runtime WHERE share_id = ? AND role = ? ORDER BY id DESC LIMIT 1`, share.ID, "forward").Row().Scan(&status, &applied, &serviceName); err != nil {
t.Fatalf("query released runtime: %v", err)
}
if status != 0 {
t.Fatalf("expected status=0 after release, got %d", status)
}
if applied != 0 {
t.Fatalf("expected applied=0 after release, got %d", applied)
}
if serviceName != "" {
t.Fatalf("expected service_name cleared after release, got %q", serviceName)
}
}
func TestValidateFederationCommandPortsAcceptsTopLevelServiceArray(t *testing.T) {
share := &repo.PeerShare{
PortRangeStart: 26200,
@@ -824,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 {
+41 -14
View File
@@ -2,9 +2,12 @@ package handler
import (
"encoding/json"
"log"
"strconv"
"strings"
"time"
"go-backend/internal/store/model"
)
const bytesPerGB int64 = 1024 * 1024 * 1024
@@ -30,7 +33,7 @@ type namedConfigItem struct {
Name string `json:"name"`
}
func (h *Handler) processFlowItem(item flowItem) {
func (h *Handler) processFlowItem(nodeID int64, item flowItem) {
serviceName := strings.TrimSpace(item.N)
if serviceName == "" || serviceName == "web_api" {
return
@@ -40,7 +43,7 @@ func (h *Handler) processFlowItem(item flowItem) {
if ok {
inFlow, outFlow := h.scaleFlowByTunnel(forwardID, item.D, item.U)
_ = h.repo.AddFlow(forwardID, userID, userTunnelID, inFlow, outFlow)
h.processPeerShareFlowFromForward(forwardID, serviceName, item)
h.processPeerShareFlowFromForward(forwardID, nodeID, serviceName, item)
if userTunnelID > 0 {
h.enforceFlowPolicies(userID, userTunnelID)
@@ -153,7 +156,7 @@ func (h *Handler) processPeerShareFlow(runtimeID int64, item flowItem) {
h.enforcePeerShareFlowLimit(share.ID)
}
func (h *Handler) processPeerShareFlowFromForward(forwardID int64, serviceName string, item flowItem) {
func (h *Handler) processPeerShareFlowFromForward(forwardID int64, nodeID int64, serviceName string, item flowItem) {
if h == nil || h.repo == nil || forwardID <= 0 {
return
}
@@ -167,22 +170,22 @@ func (h *Handler) processPeerShareFlowFromForward(forwardID int64, serviceName s
if err != nil || forward == nil {
// Forward not found in local database - might be a federation port-forward
// Try to find by service name in peer_share_runtime
h.processPeerShareFlowByServiceName(serviceName, item)
h.processPeerShareFlowByServiceName(nodeID, serviceName, item)
return
}
tunnelName, err := h.repo.GetTunnelName(forward.TunnelID)
if err != nil {
h.processPeerShareFlowByServiceName(serviceName, item)
h.processPeerShareFlowByServiceName(nodeID, serviceName, item)
return
}
shareID, ok := parsePeerShareIDFromFederationTunnelName(tunnelName)
if !ok {
h.processPeerShareFlowByServiceName(serviceName, item)
h.processPeerShareFlowByServiceName(nodeID, serviceName, item)
return
}
if err := h.repo.AddPeerShareCurrentFlow(shareID, delta); err != nil {
h.processPeerShareFlowByServiceName(serviceName, item)
h.processPeerShareFlowByServiceName(nodeID, serviceName, item)
return
}
@@ -207,7 +210,7 @@ func normalizeForwardRuntimeServiceName(serviceName string) string {
return name
}
func (h *Handler) processPeerShareFlowByServiceName(serviceName string, item flowItem) {
func (h *Handler) processPeerShareFlowByServiceName(nodeID int64, serviceName string, item flowItem) {
if h == nil || h.repo == nil || strings.TrimSpace(serviceName) == "" {
return
}
@@ -218,17 +221,41 @@ func (h *Handler) processPeerShareFlowByServiceName(serviceName string, item flo
}
normalized := normalizeForwardRuntimeServiceName(serviceName)
runtimes, err := h.repo.ListActiveForwardPeerShareRuntimesByServiceName(normalized)
if err != nil {
return
}
if len(runtimes) == 0 && normalized != serviceName {
runtimes, err = h.repo.ListActiveForwardPeerShareRuntimesByServiceName(serviceName)
var runtimes []model.PeerShareRuntime
var err error
// Try node-scoped query first if nodeID is valid
if nodeID > 0 {
runtimes, err = h.repo.ListActiveForwardPeerShareRuntimesByNodeAndServiceName(nodeID, normalized)
if err != nil {
return
}
if len(runtimes) == 0 && normalized != serviceName {
runtimes, err = h.repo.ListActiveForwardPeerShareRuntimesByNodeAndServiceName(nodeID, serviceName)
if err != nil {
return
}
}
}
// Fallback to global query if node-scoped query returned nothing or nodeID is invalid
if len(runtimes) == 0 {
runtimes, err = h.repo.ListActiveForwardPeerShareRuntimesByServiceName(normalized)
if err != nil {
return
}
if len(runtimes) == 0 && normalized != serviceName {
runtimes, err = h.repo.ListActiveForwardPeerShareRuntimesByServiceName(serviceName)
if err != nil {
return
}
}
}
if len(runtimes) != 1 {
if len(runtimes) > 1 {
log.Printf("WARN: ambiguous peer share runtime match for service=%s nodeID=%d count=%d", serviceName, nodeID, len(runtimes))
}
return
}
runtime := runtimes[0]
@@ -44,7 +44,7 @@ func TestProcessFlowItemTracksPeerShareFlowAndEnforcesLimit(t *testing.T) {
}
h := &Handler{repo: r}
h.processFlowItem(flowItem{N: "fed_svc_17", U: 1200, D: 900})
h.processFlowItem(1, flowItem{N: "fed_svc_17", U: 1200, D: 900})
updatedShare, err := r.GetPeerShare(share.ID)
if err != nil || updatedShare == nil {
@@ -120,7 +120,7 @@ func TestProcessFlowItemTracksPeerShareFlowForFederationPortForward(t *testing.T
}
h := &Handler{repo: r}
h.processFlowItem(flowItem{N: "20_2_10", U: 120, D: 80})
h.processFlowItem(1, flowItem{N: "20_2_10", U: 120, D: 80})
updatedShare, err := r.GetPeerShare(share.ID)
if err != nil || updatedShare == nil {
@@ -166,7 +166,7 @@ func TestProcessFlowItemTracksPeerShareFlowByForwardServiceName(t *testing.T) {
}
h := &Handler{repo: r}
h.processFlowItem(flowItem{N: "20_2_10_tcp", U: 120, D: 80})
h.processFlowItem(1, flowItem{N: "20_2_10_tcp", U: 120, D: 80})
updatedShare, err := r.GetPeerShare(share.ID)
if err != nil || updatedShare == nil {
@@ -226,7 +226,7 @@ func TestProcessFlowItemFallsBackToServiceNameWhenForwardIDCollidesAcrossPanels(
}
h := &Handler{repo: r}
h.processFlowItem(flowItem{N: "20_2_10_tcp", U: 120, D: 80})
h.processFlowItem(1, flowItem{N: "20_2_10_tcp", U: 120, D: 80})
updatedShare, err := r.GetPeerShare(share.ID)
if err != nil || updatedShare == nil {
@@ -288,7 +288,7 @@ func TestProcessFlowItemSkipsPeerShareFlowWhenServiceNameIsAmbiguous(t *testing.
}
h := &Handler{repo: r}
h.processFlowItem(flowItem{N: "99_2_10_tcp", U: 120, D: 80})
h.processFlowItem(1, flowItem{N: "99_2_10_tcp", U: 120, D: 80})
updatedA, _ := r.GetPeerShare(shareA.ID)
updatedB, _ := r.GetPeerShare(shareB.ID)
+3 -2
View File
@@ -704,7 +704,8 @@ func (h *Handler) flowConfig(w http.ResponseWriter, r *http.Request) {
func (h *Handler) flowUpload(w http.ResponseWriter, r *http.Request) {
secret := r.URL.Query().Get("secret")
if ok, _ := h.repo.NodeExistsBySecret(secret); !ok {
node, _ := h.repo.GetNodeBySecret(secret)
if node == nil {
w.Header().Set("Content-Type", "text/plain; charset=utf-8")
_, _ = w.Write([]byte("ok"))
return
@@ -715,7 +716,7 @@ func (h *Handler) flowUpload(w http.ResponseWriter, r *http.Request) {
var items []flowItem
if json.Unmarshal([]byte(raw), &items) == nil {
for _, item := range items {
h.processFlowItem(item)
h.processFlowItem(node.ID, item)
}
}
}
@@ -1287,6 +1287,28 @@ func (r *Repository) ListActiveForwardPeerShareRuntimesByServiceName(serviceName
return items, nil
}
func (r *Repository) ListActiveForwardPeerShareRuntimesByNodeAndServiceName(nodeID int64, serviceName string) ([]model.PeerShareRuntime, error) {
if r == nil || r.db == nil {
return nil, errors.New("repository not initialized")
}
serviceName = strings.TrimSpace(serviceName)
if serviceName == "" {
return []model.PeerShareRuntime{}, nil
}
var items []model.PeerShareRuntime
err := r.db.Where("node_id = ? AND service_name = ? AND status = 1 AND role = ?", nodeID, serviceName, "forward").
Order("id ASC").
Find(&items).Error
if err != nil {
return nil, err
}
if items == nil {
items = make([]model.PeerShareRuntime, 0)
}
return items, nil
}
func (r *Repository) ListActiveForwardPeerShareRuntimeServiceNamesByNode(nodeID int64) ([]string, error) {
if r == nil || r.db == nil {
return nil, errors.New("repository not initialized")
@@ -1333,6 +1355,27 @@ func (r *Repository) GetActiveForwardPeerShareRuntimeByPort(shareID int64, port
return &item, nil
}
func (r *Repository) GetActiveForwardPeerShareRuntimeByServiceName(shareID int64, serviceName string) (*model.PeerShareRuntime, error) {
if r == nil || r.db == nil {
return nil, errors.New("repository not initialized")
}
serviceName = strings.TrimSpace(serviceName)
if shareID <= 0 || serviceName == "" {
return nil, nil
}
var item model.PeerShareRuntime
err := r.db.Where("share_id = ? AND service_name = ? AND status = 1 AND role = ?", shareID, serviceName, "forward").
Order("id ASC").
First(&item).Error
if errors.Is(err, gorm.ErrRecordNotFound) {
return nil, nil
}
if err != nil {
return nil, err
}
return &item, nil
}
func (r *Repository) ExistsActivePeerShareRuntimeOnNodePort(nodeID int64, port int) (bool, error) {
if r == nil || r.db == nil {
return false, errors.New("repository not initialized")
@@ -1376,6 +1419,27 @@ func (r *Repository) MarkPeerShareRuntimeReleasedByPort(shareID int64, port int,
}).Error
}
func (r *Repository) MarkForwardPeerShareRuntimeReleasedByServiceName(shareID int64, serviceName string, updatedTime int64) error {
if r == nil || r.db == nil {
return errors.New("repository not initialized")
}
serviceName = strings.TrimSpace(serviceName)
if shareID <= 0 || serviceName == "" {
return nil
}
if updatedTime <= 0 {
updatedTime = unixMilliNow()
}
return r.db.Model(&model.PeerShareRuntime{}).
Where("share_id = ? AND status = 1 AND role = ? AND service_name = ?", shareID, "forward", serviceName).
Updates(map[string]interface{}{
"status": 0,
"applied": 0,
"service_name": "",
"updated_time": updatedTime,
}).Error
}
// ─── FederationTunnelBinding ─────────────────────────────────────────
func (r *Repository) UpsertFederationTunnelBinding(item *model.FederationTunnelBinding) error {
@@ -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 {
@@ -166,7 +166,7 @@ func TestFederationDualPanelMiddleExitAutoPortContract(t *testing.T) {
assertCount(t, providerRepo, `SELECT COUNT(1) FROM peer_share_runtime WHERE share_id = ? AND status = 1 AND applied = 1`, middleShareID, 1)
assertCount(t, providerRepo, `SELECT COUNT(1) FROM peer_share_runtime WHERE share_id = ? AND status = 1 AND applied = 1`, exitShareID, 1)
assertCount(t, providerRepo, `SELECT COUNT(1) FROM peer_share_runtime WHERE share_id = ?`, entryShareID, 0)
assertCount(t, providerRepo, `SELECT COUNT(1) FROM peer_share_runtime WHERE share_id = ? AND status = 1 AND applied = 1`, entryShareID, 1)
}
func TestFederationDualPanelRemoteDiagnosisContract(t *testing.T) {
+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}`;
}