mirror of
https://github.com/Sagit-chu/flvx.git
synced 2026-09-28 07:36:38 +08:00
Compare commits
2 Commits
2.1.5-rc11
...
2.1.5-rc13
| Author | SHA1 | Date | |
|---|---|---|---|
| 6189fe23f1 | |||
| 7ba90e8696 |
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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}`;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user