mirror of
https://github.com/Sagit-chu/flvx.git
synced 2026-09-28 23:56:36 +08:00
Compare commits
11 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| c8c1841058 | |||
| 1c596fae4b | |||
| 2ff52e3275 | |||
| 7efb49bdab | |||
| a00b20abf3 | |||
| 1450b25475 | |||
| b815be54b8 | |||
| 75edeb9afa | |||
| 7c54192055 | |||
| 7ba68778c1 | |||
| ef613c1518 |
@@ -223,21 +223,27 @@ func (h *Handler) listUserTunnelIDsByUser(userID int64) ([]int64, error) {
|
||||
}
|
||||
|
||||
func (h *Handler) syncForwardServices(forward *forwardRecord, method string, allowFallbackAdd bool) error {
|
||||
_, err := h.syncForwardServicesWithWarnings(forward, method, allowFallbackAdd)
|
||||
return err
|
||||
}
|
||||
|
||||
func (h *Handler) syncForwardServicesWithWarnings(forward *forwardRecord, method string, allowFallbackAdd bool) ([]string, error) {
|
||||
if h == nil || forward == nil {
|
||||
return errors.New("invalid forward sync context")
|
||||
return nil, errors.New("invalid forward sync context")
|
||||
}
|
||||
|
||||
tunnel, err := h.getTunnelRecord(forward.TunnelID)
|
||||
if err != nil {
|
||||
return err
|
||||
return nil, err
|
||||
}
|
||||
ports, err := h.listForwardPorts(forward.ID)
|
||||
if err != nil {
|
||||
return err
|
||||
return nil, err
|
||||
}
|
||||
if len(ports) == 0 {
|
||||
return errors.New("转发入口端口不存在")
|
||||
return nil, errors.New("转发入口端口不存在")
|
||||
}
|
||||
warnings := make([]string, 0)
|
||||
|
||||
// Determine limiter from forward's SpeedID first, fallback to UserTunnel's limiter
|
||||
var limiterID *int64
|
||||
@@ -258,7 +264,7 @@ func (h *Handler) syncForwardServices(forward *forwardRecord, method string, all
|
||||
var utSpeed *int
|
||||
_, utLimiterID, utSpeed, err = h.resolveUserTunnelAndLimiter(forward.UserID, forward.TunnelID)
|
||||
if err != nil {
|
||||
return err
|
||||
return nil, err
|
||||
}
|
||||
limiterID = utLimiterID
|
||||
speed = utSpeed
|
||||
@@ -267,28 +273,138 @@ func (h *Handler) syncForwardServices(forward *forwardRecord, method string, all
|
||||
serviceBase := buildForwardServiceBase(forward.ID, forward.UserID, 0)
|
||||
tunnelTLSProtocol, err := h.isTunnelSelectedTLSProtocol(forward.TunnelID)
|
||||
if err != nil {
|
||||
return err
|
||||
return nil, err
|
||||
}
|
||||
|
||||
for _, fp := range ports {
|
||||
if limiterID != nil && speed != nil {
|
||||
if err := h.ensureLimiterOnNode(fp.NodeID, *limiterID, *speed); err != nil {
|
||||
return err
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
|
||||
node, err := h.getNodeRecord(fp.NodeID)
|
||||
if err != nil {
|
||||
return err
|
||||
return nil, err
|
||||
}
|
||||
services := buildForwardServiceConfigs(serviceBase, forward, tunnel, node, fp.Port, strings.TrimSpace(fp.InIP), limiterID, tunnelTLSProtocol)
|
||||
_, err = h.sendNodeCommand(node.ID, method, services, true, false)
|
||||
if err != nil && allowFallbackAdd && method == "UpdateService" {
|
||||
_, err = h.sendNodeCommand(node.ID, "AddService", services, true, false)
|
||||
}
|
||||
if err != nil {
|
||||
return fmt.Errorf("节点 %s 下发失败: %w", node.Name, err)
|
||||
if err != nil && strings.EqualFold(strings.TrimSpace(method), "UpdateService") && isAddressAlreadyInUseError(err) {
|
||||
err = h.rebindForwardServiceOnSelfOccupiedPort(forward, node, fp.Port, services)
|
||||
}
|
||||
if err != nil && strings.EqualFold(strings.TrimSpace(method), "UpdateService") && isCannotAssignRequestedAddressError(err) {
|
||||
var warning string
|
||||
warning, err = h.fallbackForwardPortToDefaultBind(forward, tunnel, node, fp, serviceBase, limiterID, tunnelTLSProtocol)
|
||||
if err == nil && warning != "" {
|
||||
warnings = append(warnings, warning)
|
||||
}
|
||||
}
|
||||
if err != nil {
|
||||
return warnings, fmt.Errorf("节点 %s 下发失败: %w", node.Name, err)
|
||||
}
|
||||
}
|
||||
return warnings, nil
|
||||
}
|
||||
|
||||
func (h *Handler) fallbackForwardPortToDefaultBind(forward *forwardRecord, tunnel *tunnelRecord, node *nodeRecord, fp forwardPortRecord, serviceBase string, limiterID *int64, tunnelTLSProtocol bool) (string, error) {
|
||||
if h == nil || forward == nil || tunnel == nil || node == nil {
|
||||
return "", errors.New("invalid bind fallback context")
|
||||
}
|
||||
if fp.Port <= 0 {
|
||||
return "", errors.New("invalid forward port")
|
||||
}
|
||||
explicitBindIP := strings.TrimSpace(fp.InIP)
|
||||
if explicitBindIP == "" {
|
||||
return "", errors.New("default bind address cannot be assigned")
|
||||
}
|
||||
|
||||
if err := h.deleteForwardServicesOnNode(forward, node.ID); err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
time.Sleep(150 * time.Millisecond)
|
||||
defaultServices := buildForwardServiceConfigs(serviceBase, forward, tunnel, node, fp.Port, "", limiterID, tunnelTLSProtocol)
|
||||
if _, err := h.sendNodeCommand(node.ID, "AddService", defaultServices, true, false); err != nil {
|
||||
return "", err
|
||||
}
|
||||
if err := h.repo.UpdateForwardPortBindIP(forward.ID, node.ID, fp.Port, ""); err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
warning := fmt.Sprintf("节点 %s 监听IP %s 不在主机网卡地址中,已自动回退为默认监听IP", strings.TrimSpace(node.Name), explicitBindIP)
|
||||
return warning, nil
|
||||
}
|
||||
|
||||
func (h *Handler) rebindForwardServiceOnSelfOccupiedPort(forward *forwardRecord, node *nodeRecord, port int, services []map[string]interface{}) error {
|
||||
if h == nil || forward == nil || node == nil {
|
||||
return errors.New("invalid self-occupy rebind context")
|
||||
}
|
||||
if port <= 0 {
|
||||
return errors.New("invalid forward port")
|
||||
}
|
||||
|
||||
hasOtherForward, err := h.repo.HasOtherForwardOnNodePort(node.ID, port, forward.ID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if hasOtherForward {
|
||||
return fmt.Errorf("端口 %d 已被其他转发占用", port)
|
||||
}
|
||||
|
||||
if err := h.deleteForwardServicesOnNode(forward, node.ID); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
time.Sleep(150 * time.Millisecond)
|
||||
|
||||
_, err = h.sendNodeCommand(node.ID, "AddService", services, true, false)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (h *Handler) deleteForwardServicesOnNode(forward *forwardRecord, nodeID int64) error {
|
||||
if h == nil || forward == nil {
|
||||
return errors.New("invalid forward delete context")
|
||||
}
|
||||
|
||||
userTunnelID, _, _, err := h.resolveUserTunnelAndLimiter(forward.UserID, forward.TunnelID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
userTunnelIDs, err := h.listUserTunnelIDs(forward.UserID, forward.TunnelID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
allUserTunnelIDs, err := h.listUserTunnelIDsByUser(forward.UserID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
candidateTunnelIDs := make([]int64, 0, len(userTunnelIDs)+len(allUserTunnelIDs))
|
||||
candidateTunnelIDs = append(candidateTunnelIDs, userTunnelIDs...)
|
||||
candidateTunnelIDs = append(candidateTunnelIDs, allUserTunnelIDs...)
|
||||
bases := buildForwardServiceBaseCandidates(forward.ID, forward.UserID, userTunnelID, candidateTunnelIDs)
|
||||
|
||||
var lastErr error
|
||||
for _, base := range bases {
|
||||
names := buildForwardControlServiceNames(base, "DeleteService")
|
||||
payload := map[string]interface{}{
|
||||
"services": names,
|
||||
}
|
||||
_, cmdErr := h.sendNodeCommand(nodeID, "DeleteService", payload, false, true)
|
||||
if cmdErr == nil {
|
||||
return nil
|
||||
}
|
||||
lastErr = cmdErr
|
||||
}
|
||||
|
||||
if lastErr != nil {
|
||||
return lastErr
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -1303,6 +1419,42 @@ func isAlreadyExistsMessage(message string) bool {
|
||||
return strings.Contains(msg, "already exists") || strings.Contains(msg, "已存在")
|
||||
}
|
||||
|
||||
func isBindAddressInUseError(err error) bool {
|
||||
if err == nil {
|
||||
return false
|
||||
}
|
||||
msg := strings.ToLower(strings.TrimSpace(err.Error()))
|
||||
if msg == "" {
|
||||
return false
|
||||
}
|
||||
return isAddressAlreadyInUseMessage(msg) || strings.Contains(msg, "cannot assign requested address")
|
||||
}
|
||||
|
||||
func isAddressAlreadyInUseError(err error) bool {
|
||||
if err == nil {
|
||||
return false
|
||||
}
|
||||
return isAddressAlreadyInUseMessage(strings.ToLower(strings.TrimSpace(err.Error())))
|
||||
}
|
||||
|
||||
func isAddressAlreadyInUseMessage(msg string) bool {
|
||||
if msg == "" {
|
||||
return false
|
||||
}
|
||||
return strings.Contains(msg, "address already in use")
|
||||
}
|
||||
|
||||
func isCannotAssignRequestedAddressError(err error) bool {
|
||||
if err == nil {
|
||||
return false
|
||||
}
|
||||
msg := strings.ToLower(strings.TrimSpace(err.Error()))
|
||||
if msg == "" {
|
||||
return false
|
||||
}
|
||||
return strings.Contains(msg, "cannot assign requested address")
|
||||
}
|
||||
|
||||
func buildForwardServiceConfigs(baseName string, forward *forwardRecord, tunnel *tunnelRecord, node *nodeRecord, port int, bindIP string, limiterID *int64, tunnelTLSProtocol bool) []map[string]interface{} {
|
||||
protocols := []string{"tcp", "udp"}
|
||||
services := make([]map[string]interface{}, 0, 2)
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package handler
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"reflect"
|
||||
"testing"
|
||||
)
|
||||
@@ -66,6 +67,39 @@ func TestIsAlreadyExistsMessage(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestIsBindAddressInUseError(t *testing.T) {
|
||||
if !isBindAddressInUseError(errors.New("listen tcp [::]:10001: bind: address already in use")) {
|
||||
t.Fatalf("address already in use should be detected")
|
||||
}
|
||||
if !isBindAddressInUseError(errors.New("listen tcp4 13.228.170.187:16765: bind: cannot assign requested address")) {
|
||||
t.Fatalf("cannot assign requested address should be detected")
|
||||
}
|
||||
if isBindAddressInUseError(errors.New("service demo already exists")) {
|
||||
t.Fatalf("already exists should not be treated as bind conflict")
|
||||
}
|
||||
if isBindAddressInUseError(nil) {
|
||||
t.Fatalf("nil error should not be treated as bind conflict")
|
||||
}
|
||||
}
|
||||
|
||||
func TestIsAddressAlreadyInUseError(t *testing.T) {
|
||||
if !isAddressAlreadyInUseError(errors.New("listen tcp [::]:10001: bind: address already in use")) {
|
||||
t.Fatalf("address already in use should be detected")
|
||||
}
|
||||
if isAddressAlreadyInUseError(errors.New("listen tcp4 13.228.170.187:16765: bind: cannot assign requested address")) {
|
||||
t.Fatalf("cannot assign requested address should not be treated as address-in-use")
|
||||
}
|
||||
}
|
||||
|
||||
func TestIsCannotAssignRequestedAddressError(t *testing.T) {
|
||||
if !isCannotAssignRequestedAddressError(errors.New("listen tcp4 13.228.170.187:16765: bind: cannot assign requested address")) {
|
||||
t.Fatalf("cannot assign requested address should be detected")
|
||||
}
|
||||
if isCannotAssignRequestedAddressError(errors.New("listen tcp [::]:10001: bind: address already in use")) {
|
||||
t.Fatalf("address already in use should not be treated as cannot-assign")
|
||||
}
|
||||
}
|
||||
|
||||
func TestBuildForwardServiceConfigs_UsesBindIPForListen(t *testing.T) {
|
||||
forward := &forwardRecord{RemoteAddr: "1.2.3.4:80", Strategy: "fifo", TunnelID: 7}
|
||||
node := &nodeRecord{TCPListenAddr: "[::]", UDPListenAddr: "[::]"}
|
||||
|
||||
@@ -1285,7 +1285,12 @@ func (h *Handler) forwardUpdate(w http.ResponseWriter, r *http.Request) {
|
||||
port = h.pickTunnelPort(tunnelID)
|
||||
}
|
||||
}
|
||||
inIp := asString(req["inIp"])
|
||||
hasInIP := false
|
||||
inIp := ""
|
||||
if rawInIP, ok := req["inIp"]; ok {
|
||||
hasInIP = true
|
||||
inIp = asString(rawInIP)
|
||||
}
|
||||
fwdEntryNodes, _ := h.tunnelEntryNodeIDs(tunnelID)
|
||||
for _, nodeID := range fwdEntryNodes {
|
||||
node, nodeErr := h.getNodeRecord(nodeID)
|
||||
@@ -1302,7 +1307,14 @@ func (h *Handler) forwardUpdate(w http.ResponseWriter, r *http.Request) {
|
||||
response.WriteJSON(w, response.Err(-2, err.Error()))
|
||||
return
|
||||
}
|
||||
if err := h.replaceForwardPorts(id, tunnelID, port, inIp); err != nil {
|
||||
if hasInIP {
|
||||
err = h.replaceForwardPorts(id, tunnelID, port, inIp)
|
||||
} else if tunnelID != forward.TunnelID {
|
||||
err = h.replaceForwardPorts(id, tunnelID, port, "")
|
||||
} else {
|
||||
err = h.replaceForwardPortsPreservingInIP(id, tunnelID, port, oldPorts)
|
||||
}
|
||||
if err != nil {
|
||||
h.rollbackForwardMutation(forward, oldPorts)
|
||||
response.WriteJSON(w, response.Err(-2, err.Error()))
|
||||
return
|
||||
@@ -1313,11 +1325,16 @@ func (h *Handler) forwardUpdate(w http.ResponseWriter, r *http.Request) {
|
||||
response.WriteJSON(w, response.Err(-2, err.Error()))
|
||||
return
|
||||
}
|
||||
if err := h.syncForwardServices(updatedForward, "UpdateService", true); err != nil {
|
||||
warnings, err := h.syncForwardServicesWithWarnings(updatedForward, "UpdateService", true)
|
||||
if err != nil {
|
||||
h.rollbackForwardMutation(forward, oldPorts)
|
||||
response.WriteJSON(w, response.ErrDefault(err.Error()))
|
||||
return
|
||||
}
|
||||
if len(warnings) > 0 {
|
||||
response.WriteJSON(w, response.OK(map[string]interface{}{"warnings": warnings}))
|
||||
return
|
||||
}
|
||||
response.WriteJSON(w, response.OKEmpty())
|
||||
}
|
||||
|
||||
@@ -3016,38 +3033,62 @@ func parsePorts(portRange string) ([]int, error) {
|
||||
return ports, nil
|
||||
}
|
||||
|
||||
type forwardPortReplaceEntry = struct {
|
||||
NodeID int64
|
||||
Port int
|
||||
InIP string
|
||||
}
|
||||
|
||||
func buildForwardPortEntriesWithPreservedInIP(entryNodeIDs []int64, oldPorts []forwardPortRecord, port int) []forwardPortReplaceEntry {
|
||||
preservedByNode := make(map[int64]string)
|
||||
for _, fp := range oldPorts {
|
||||
current, exists := preservedByNode[fp.NodeID]
|
||||
if !exists {
|
||||
preservedByNode[fp.NodeID] = fp.InIP
|
||||
continue
|
||||
}
|
||||
if strings.TrimSpace(current) == "" && strings.TrimSpace(fp.InIP) != "" {
|
||||
preservedByNode[fp.NodeID] = fp.InIP
|
||||
}
|
||||
}
|
||||
|
||||
entries := make([]forwardPortReplaceEntry, 0, len(entryNodeIDs))
|
||||
for _, nid := range entryNodeIDs {
|
||||
entries = append(entries, forwardPortReplaceEntry{
|
||||
NodeID: nid,
|
||||
Port: port,
|
||||
InIP: preservedByNode[nid],
|
||||
})
|
||||
}
|
||||
|
||||
return entries
|
||||
}
|
||||
|
||||
func (h *Handler) replaceForwardPorts(forwardID, tunnelID int64, port int, inIp string) error {
|
||||
entryNodes, err := h.tunnelEntryNodeIDs(tunnelID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
entries := make([]struct {
|
||||
NodeID int64
|
||||
Port int
|
||||
InIP string
|
||||
}, len(entryNodes))
|
||||
entries := make([]forwardPortReplaceEntry, len(entryNodes))
|
||||
for i, nid := range entryNodes {
|
||||
entries[i] = struct {
|
||||
NodeID int64
|
||||
Port int
|
||||
InIP string
|
||||
}{NodeID: nid, Port: port, InIP: inIp}
|
||||
entries[i] = forwardPortReplaceEntry{NodeID: nid, Port: port, InIP: inIp}
|
||||
}
|
||||
return h.repo.ReplaceForwardPorts(forwardID, entries)
|
||||
}
|
||||
|
||||
func (h *Handler) replaceForwardPortsPreservingInIP(forwardID, tunnelID int64, port int, oldPorts []forwardPortRecord) error {
|
||||
entryNodes, err := h.tunnelEntryNodeIDs(tunnelID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
entries := buildForwardPortEntriesWithPreservedInIP(entryNodes, oldPorts, port)
|
||||
return h.repo.ReplaceForwardPorts(forwardID, entries)
|
||||
}
|
||||
|
||||
func (h *Handler) replaceForwardPortsWithRecords(forwardID int64, ports []forwardPortRecord) error {
|
||||
entries := make([]struct {
|
||||
NodeID int64
|
||||
Port int
|
||||
InIP string
|
||||
}, len(ports))
|
||||
entries := make([]forwardPortReplaceEntry, len(ports))
|
||||
for i, fp := range ports {
|
||||
entries[i] = struct {
|
||||
NodeID int64
|
||||
Port int
|
||||
InIP string
|
||||
}{NodeID: fp.NodeID, Port: fp.Port, InIP: fp.InIP}
|
||||
entries[i] = forwardPortReplaceEntry{NodeID: fp.NodeID, Port: fp.Port, InIP: fp.InIP}
|
||||
}
|
||||
return h.repo.ReplaceForwardPorts(forwardID, entries)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,39 @@
|
||||
package handler
|
||||
|
||||
import "testing"
|
||||
|
||||
func TestBuildForwardPortEntriesWithPreservedInIP(t *testing.T) {
|
||||
entryNodeIDs := []int64{10, 20, 30}
|
||||
oldPorts := []forwardPortRecord{
|
||||
{NodeID: 10, Port: 10001, InIP: ""},
|
||||
{NodeID: 10, Port: 10002, InIP: "10.0.0.10"},
|
||||
{NodeID: 20, Port: 10003, InIP: "10.0.0.20"},
|
||||
}
|
||||
|
||||
entries := buildForwardPortEntriesWithPreservedInIP(entryNodeIDs, oldPorts, 18080)
|
||||
if len(entries) != 3 {
|
||||
t.Fatalf("expected 3 entries, got %d", len(entries))
|
||||
}
|
||||
|
||||
if entries[0].NodeID != 10 || entries[0].Port != 18080 || entries[0].InIP != "10.0.0.10" {
|
||||
t.Fatalf("unexpected first entry: %+v", entries[0])
|
||||
}
|
||||
if entries[1].NodeID != 20 || entries[1].Port != 18080 || entries[1].InIP != "10.0.0.20" {
|
||||
t.Fatalf("unexpected second entry: %+v", entries[1])
|
||||
}
|
||||
if entries[2].NodeID != 30 || entries[2].Port != 18080 || entries[2].InIP != "" {
|
||||
t.Fatalf("unexpected third entry: %+v", entries[2])
|
||||
}
|
||||
}
|
||||
|
||||
func TestBuildForwardPortEntriesWithPreservedInIP_EmptyOldPorts(t *testing.T) {
|
||||
entryNodeIDs := []int64{99}
|
||||
entries := buildForwardPortEntriesWithPreservedInIP(entryNodeIDs, nil, 17000)
|
||||
|
||||
if len(entries) != 1 {
|
||||
t.Fatalf("expected 1 entry, got %d", len(entries))
|
||||
}
|
||||
if entries[0].NodeID != 99 || entries[0].Port != 17000 || entries[0].InIP != "" {
|
||||
t.Fatalf("unexpected entry: %+v", entries[0])
|
||||
}
|
||||
}
|
||||
@@ -111,6 +111,25 @@ func (r *Repository) ListForwardPorts(forwardID int64) ([]model.ForwardPortRecor
|
||||
return rows, nil
|
||||
}
|
||||
|
||||
func (r *Repository) HasOtherForwardOnNodePort(nodeID int64, port int, currentForwardID int64) (bool, error) {
|
||||
if r == nil || r.db == nil {
|
||||
return false, errors.New("repository not initialized")
|
||||
}
|
||||
if nodeID <= 0 || port <= 0 {
|
||||
return false, nil
|
||||
}
|
||||
|
||||
var count int64
|
||||
err := r.db.Model(&model.ForwardPort{}).
|
||||
Where("node_id = ? AND port = ? AND forward_id <> ?", nodeID, port, currentForwardID).
|
||||
Count(&count).Error
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
|
||||
return count > 0, nil
|
||||
}
|
||||
|
||||
func (r *Repository) GetTunnelOutProtocol(tunnelID int64) (string, error) {
|
||||
if r == nil || r.db == nil {
|
||||
return "", errors.New("repository not initialized")
|
||||
|
||||
@@ -720,6 +720,18 @@ func (r *Repository) ReplaceForwardPorts(forwardID int64, entries []struct {
|
||||
})
|
||||
}
|
||||
|
||||
func (r *Repository) UpdateForwardPortBindIP(forwardID, nodeID int64, port int, inIP string) error {
|
||||
if r == nil || r.db == nil {
|
||||
return errors.New("repository not initialized")
|
||||
}
|
||||
if forwardID <= 0 || nodeID <= 0 || port <= 0 {
|
||||
return nil
|
||||
}
|
||||
return r.db.Model(&model.ForwardPort{}).
|
||||
Where("forward_id = ? AND node_id = ? AND port = ?", forwardID, nodeID, port).
|
||||
Update("in_ip", sql.NullString{String: inIP, Valid: strings.TrimSpace(inIP) != ""}).Error
|
||||
}
|
||||
|
||||
func (r *Repository) RollbackForwardFields(id, userID int64, userName, name string, tunnelID int64, remoteAddr, strategy string, status int, speedID interface{}, now int64) {
|
||||
if r == nil || r.db == nil {
|
||||
return
|
||||
|
||||
@@ -0,0 +1,7 @@
|
||||
- [x] Review current forward import flow and confirm ny import uses tunnel selection
|
||||
- [x] Define ny compatibility update with tunnel-first behavior and auto port assignment fallback
|
||||
- [x] Update ny parser to accept alias fields and optional `listen_port`
|
||||
- [x] Keep import execution bound to selected tunnel and remove entry-selection dependency from ux copy
|
||||
- [x] Update ny import help text to document optional port auto assignment
|
||||
- [x] Add parser tests for alias-field compatibility and missing-port auto assignment
|
||||
- [x] Validate updated import parser tests locally
|
||||
@@ -0,0 +1,11 @@
|
||||
# 003 Forward Edit Bind IP Preserve
|
||||
|
||||
## Checklist
|
||||
|
||||
- [x] Confirm forward edit flow and identify why untouched listen IP gets overwritten.
|
||||
- [x] Update frontend forward edit submit logic to only send `inIp` when user explicitly changes listen IP.
|
||||
- [x] On tunnel switch in edit form, reset listen IP to default unless user reselects.
|
||||
- [x] Update backend forward update logic to preserve existing `forward_port.in_ip` when request omits `inIp` and tunnel is unchanged.
|
||||
- [x] Keep backend behavior explicit: if `inIp` is sent (including empty), apply requested value; if tunnel changed with no `inIp`, use default bind.
|
||||
- [x] Add regression tests for preserved bind-IP reconstruction helper behavior.
|
||||
- [x] Run focused frontend/backend checks for touched files.
|
||||
@@ -0,0 +1,11 @@
|
||||
# 004 Forward Explicit Bind Self-Occupy Release
|
||||
|
||||
## Checklist
|
||||
|
||||
- [x] Confirm current forward edit/save failure path and lock strategy: explicit bind always stays explicit.
|
||||
- [x] Add repository query to detect whether a node+port is occupied by other forwards (excluding current forward).
|
||||
- [x] Enhance forward service sync to treat address-in-use as a recoverable case when only self occupies the port.
|
||||
- [x] On self-occupy conflict, proactively delete current forward services on target node and retry AddService.
|
||||
- [x] Keep hard failure when the same node+port is occupied by other forwards.
|
||||
- [x] Add focused unit tests for new error classification helpers.
|
||||
- [x] Run focused backend tests for touched handler/repo packages.
|
||||
@@ -0,0 +1,11 @@
|
||||
# 005 Forward Invalid BindIP Fallback Default
|
||||
|
||||
## Checklist
|
||||
|
||||
- [x] Split forward service bind failures into address-in-use and cannot-assign classes.
|
||||
- [x] Keep self-occupy release/rebind only for address-in-use conflicts.
|
||||
- [x] Add fallback path for cannot-assign: switch to default listener bind and retry service creation.
|
||||
- [x] Persist fallback result to DB by clearing `forward_port.in_ip` for affected node+port.
|
||||
- [x] Return non-blocking warning in forward update response when fallback occurs.
|
||||
- [x] Show warning toast in forward edit UI while still treating operation as success.
|
||||
- [x] Run focused backend tests for touched handler/repo packages.
|
||||
@@ -563,11 +563,6 @@ export default function ForwardPage() {
|
||||
const [importData, setImportData] = useState("");
|
||||
const [importLoading, setImportLoading] = useState(false);
|
||||
const [importFormat, setImportFormat] = useState<ImportFormat>("flvx");
|
||||
const [selectedEntryNode, setSelectedEntryNode] = useState<number | null>(
|
||||
null,
|
||||
);
|
||||
const [matchedTunnels, setMatchedTunnels] = useState<Tunnel[]>([]);
|
||||
const [tunnelSelectModalOpen, setTunnelSelectModalOpen] = useState(false);
|
||||
const [selectedTunnelForImport, setSelectedTunnelForImport] = useState<
|
||||
number | null
|
||||
>(null);
|
||||
@@ -591,6 +586,7 @@ export default function ForwardPage() {
|
||||
strategy: "fifo",
|
||||
speedId: null,
|
||||
});
|
||||
const [inIpTouched, setInIpTouched] = useState(false);
|
||||
|
||||
// 表单验证错误
|
||||
const [errors, setErrors] = useState<{ [key: string]: string }>({});
|
||||
@@ -1281,6 +1277,7 @@ export default function ForwardPage() {
|
||||
// 新增转发
|
||||
const handleAdd = () => {
|
||||
setIsEdit(false);
|
||||
setInIpTouched(false);
|
||||
setForm({
|
||||
name: "",
|
||||
tunnelId: null,
|
||||
@@ -1298,6 +1295,7 @@ export default function ForwardPage() {
|
||||
// 编辑转发
|
||||
const handleEdit = (forward: Forward) => {
|
||||
setIsEdit(true);
|
||||
setInIpTouched(false);
|
||||
setForm({
|
||||
id: forward.id,
|
||||
userId: forward.userId,
|
||||
@@ -1362,11 +1360,17 @@ export default function ForwardPage() {
|
||||
const nextTunnelId = parseInt(tunnelId);
|
||||
const options = tunnelInIpOptionMap.get(nextTunnelId) || [];
|
||||
|
||||
setForm((prev) => ({
|
||||
...prev,
|
||||
tunnelId: nextTunnelId,
|
||||
inIp: options.includes(prev.inIp) ? prev.inIp : "",
|
||||
}));
|
||||
setInIpTouched(false);
|
||||
|
||||
setForm((prev) => {
|
||||
const tunnelChanged = prev.tunnelId !== nextTunnelId;
|
||||
|
||||
return {
|
||||
...prev,
|
||||
tunnelId: nextTunnelId,
|
||||
inIp: tunnelChanged ? "" : options.includes(prev.inIp) ? prev.inIp : "",
|
||||
};
|
||||
});
|
||||
};
|
||||
|
||||
// 提交表单
|
||||
@@ -1393,7 +1397,7 @@ export default function ForwardPage() {
|
||||
name: form.name,
|
||||
tunnelId: form.tunnelId,
|
||||
inPort: form.inPort,
|
||||
inIp: form.inIp || undefined,
|
||||
...(inIpTouched ? { inIp: form.inIp || "" } : {}),
|
||||
remoteAddr: processedRemoteAddr,
|
||||
strategy: addressCount > 1 ? form.strategy : "fifo",
|
||||
speedId: normalizeSpeedId(form.speedId),
|
||||
@@ -1415,6 +1419,18 @@ export default function ForwardPage() {
|
||||
}
|
||||
|
||||
if (res.code === 0) {
|
||||
const warningItems = Array.isArray((res as any).data?.warnings)
|
||||
? (res as any).data.warnings
|
||||
.map((item: unknown) => (typeof item === "string" ? item.trim() : ""))
|
||||
.filter((item: string) => item)
|
||||
: [];
|
||||
|
||||
warningItems.forEach((warning: string) => {
|
||||
toast(warning, {
|
||||
icon: "⚠️",
|
||||
duration: 5000,
|
||||
});
|
||||
});
|
||||
toast.success(isEdit ? "修改成功" : "创建成功");
|
||||
setModalOpen(false);
|
||||
loadData();
|
||||
@@ -4293,6 +4309,8 @@ export default function ForwardPage() {
|
||||
onSelectionChange={(keys) => {
|
||||
const selectedKey = Array.from(keys)[0] as string;
|
||||
|
||||
setInIpTouched(true);
|
||||
|
||||
setForm((prev) => ({
|
||||
...prev,
|
||||
inIp: selectedKey === "__default__" ? "" : selectedKey,
|
||||
@@ -4625,10 +4643,10 @@ export default function ForwardPage() {
|
||||
) : (
|
||||
<>
|
||||
<p className="text-small text-default-500">
|
||||
ny格式:JSON对象,支持多个目标地址(负载均衡)
|
||||
ny格式:JSON对象,支持多个目标地址(负载均衡),按所选隧道导入
|
||||
</p>
|
||||
<p className="text-small text-default-400">
|
||||
格式:{"dest":["地址:端口"],"listen_port":端口,"name":"名称"}
|
||||
格式:{"dest":["地址:端口"],"listen_port":端口,"name":"名称"}(listen_port可省略,自动分配端口)
|
||||
</p>
|
||||
</>
|
||||
)}
|
||||
@@ -4648,8 +4666,6 @@ export default function ForwardPage() {
|
||||
if (selectedKey) {
|
||||
setImportFormat(selectedKey);
|
||||
setSelectedTunnelForImport(null);
|
||||
setSelectedEntryNode(null);
|
||||
setMatchedTunnels([]);
|
||||
setImportData("");
|
||||
setImportResults([]);
|
||||
}
|
||||
@@ -4663,114 +4679,34 @@ export default function ForwardPage() {
|
||||
</SelectItem>
|
||||
</Select>
|
||||
|
||||
{/* flvx格式:隧道选择 */}
|
||||
{importFormat === "flvx" && (
|
||||
<Select
|
||||
isRequired
|
||||
label="选择导入隧道"
|
||||
placeholder="请选择要导入的隧道"
|
||||
selectedKeys={
|
||||
selectedTunnelForImport
|
||||
? [selectedTunnelForImport.toString()]
|
||||
: []
|
||||
}
|
||||
variant="bordered"
|
||||
onSelectionChange={(keys) => {
|
||||
const selectedKey = Array.from(keys)[0] as string;
|
||||
{/* 隧道选择 - 两种格式都需要 */}
|
||||
<Select
|
||||
isRequired
|
||||
label="选择导入隧道"
|
||||
placeholder="请选择要导入的隧道"
|
||||
selectedKeys={
|
||||
selectedTunnelForImport
|
||||
? [selectedTunnelForImport.toString()]
|
||||
: []
|
||||
}
|
||||
variant="bordered"
|
||||
onSelectionChange={(keys) => {
|
||||
const selectedKey = Array.from(keys)[0] as string;
|
||||
|
||||
setSelectedTunnelForImport(
|
||||
selectedKey ? parseInt(selectedKey) : null,
|
||||
);
|
||||
}}
|
||||
>
|
||||
{tunnels.map((tunnel) => (
|
||||
<SelectItem
|
||||
key={tunnel.id.toString()}
|
||||
textValue={tunnel.name}
|
||||
>
|
||||
{tunnel.name}
|
||||
</SelectItem>
|
||||
))}
|
||||
</Select>
|
||||
)}
|
||||
|
||||
{/* ny格式:入口节点选择 */}
|
||||
{importFormat === "ny" && (
|
||||
<Select
|
||||
isRequired
|
||||
label="选择入口节点"
|
||||
placeholder="请选择入口节点"
|
||||
selectedKeys={
|
||||
selectedEntryNode ? [selectedEntryNode.toString()] : []
|
||||
}
|
||||
variant="bordered"
|
||||
onSelectionChange={(keys) => {
|
||||
const selectedKey = Array.from(keys)[0] as string;
|
||||
const nodeId = selectedKey ? parseInt(selectedKey) : null;
|
||||
|
||||
setSelectedEntryNode(nodeId);
|
||||
setSelectedTunnelForImport(null);
|
||||
|
||||
if (nodeId) {
|
||||
const matched = allTunnels.filter(
|
||||
(t) =>
|
||||
t.type === 1 &&
|
||||
t.inNodeId?.some((n) => n.nodeId === nodeId),
|
||||
);
|
||||
|
||||
setMatchedTunnels(matched);
|
||||
|
||||
if (matched.length === 0) {
|
||||
toast.error(
|
||||
"该入口节点没有匹配的隧道,请先创建端口转发类型的隧道",
|
||||
);
|
||||
} else if (matched.length === 1) {
|
||||
setSelectedTunnelForImport(matched[0].id);
|
||||
} else {
|
||||
setTunnelSelectModalOpen(true);
|
||||
}
|
||||
} else {
|
||||
setMatchedTunnels([]);
|
||||
}
|
||||
}}
|
||||
>
|
||||
{nodes.map((node) => (
|
||||
<SelectItem key={node.id.toString()} textValue={node.name}>
|
||||
{node.name}
|
||||
</SelectItem>
|
||||
))}
|
||||
</Select>
|
||||
)}
|
||||
|
||||
{/* ny格式:显示匹配的隧道 */}
|
||||
{importFormat === "ny" && matchedTunnels.length > 0 && (
|
||||
<div className="text-xs text-default-500">
|
||||
{matchedTunnels.length === 1 ? (
|
||||
<span>
|
||||
已匹配隧道:<strong>{matchedTunnels[0].name}</strong>
|
||||
</span>
|
||||
) : (
|
||||
<span>
|
||||
找到 {matchedTunnels.length}{" "}
|
||||
个匹配隧道,请点击下方按钮选择
|
||||
</span>
|
||||
)}
|
||||
</div>
|
||||
)}
|
||||
|
||||
{/* ny格式:多隧道选择按钮 */}
|
||||
{importFormat === "ny" &&
|
||||
matchedTunnels.length > 1 &&
|
||||
!selectedTunnelForImport && (
|
||||
<Button
|
||||
color="primary"
|
||||
size="sm"
|
||||
variant="flat"
|
||||
onPress={() => setTunnelSelectModalOpen(true)}
|
||||
setSelectedTunnelForImport(
|
||||
selectedKey ? parseInt(selectedKey) : null,
|
||||
);
|
||||
}}
|
||||
>
|
||||
{tunnels.map((tunnel) => (
|
||||
<SelectItem
|
||||
key={tunnel.id.toString()}
|
||||
textValue={tunnel.name}
|
||||
>
|
||||
选择隧道({matchedTunnels.length}个可选)
|
||||
</Button>
|
||||
)}
|
||||
{tunnel.name}
|
||||
</SelectItem>
|
||||
))}
|
||||
</Select>
|
||||
|
||||
{/* 输入区域 */}
|
||||
<Textarea
|
||||
@@ -4783,7 +4719,7 @@ export default function ForwardPage() {
|
||||
placeholder={
|
||||
importFormat === "flvx"
|
||||
? "请输入要导入的转发数据,格式:目标地址|转发名称|入口端口"
|
||||
: '请输入ny格式数据,每行一个JSON对象,如:{"dest":["1.2.3.4:80"],"listen_port":8080,"name":"转发1"}'
|
||||
: '请输入ny格式数据,每行一个JSON对象,如:{"dest":["1.2.3.4:80"],"listen_port":8080,"name":"转发1"};listen_port可省略自动分配'
|
||||
}
|
||||
value={importData}
|
||||
variant="flat"
|
||||
@@ -4887,11 +4823,7 @@ export default function ForwardPage() {
|
||||
</Button>
|
||||
<Button
|
||||
color="warning"
|
||||
isDisabled={
|
||||
!importData.trim() ||
|
||||
!selectedTunnelForImport ||
|
||||
(importFormat === "ny" && !selectedEntryNode)
|
||||
}
|
||||
isDisabled={!importData.trim() || !selectedTunnelForImport}
|
||||
isLoading={importLoading}
|
||||
onPress={executeImport}
|
||||
>
|
||||
@@ -4901,54 +4833,6 @@ export default function ForwardPage() {
|
||||
</ModalContent>
|
||||
</Modal>
|
||||
|
||||
{/* 隧道选择模态框(ny格式多隧道匹配时使用) */}
|
||||
<Modal
|
||||
backdrop="blur"
|
||||
isOpen={tunnelSelectModalOpen}
|
||||
placement="center"
|
||||
size="md"
|
||||
onClose={() => setTunnelSelectModalOpen(false)}
|
||||
>
|
||||
<ModalContent>
|
||||
<ModalHeader>选择隧道</ModalHeader>
|
||||
<ModalBody>
|
||||
<p className="text-sm text-default-500 mb-3">
|
||||
找到多个使用该入口节点的隧道,请选择一个:
|
||||
</p>
|
||||
<div className="space-y-2">
|
||||
{matchedTunnels.map((tunnel) => (
|
||||
<Button
|
||||
key={tunnel.id}
|
||||
className="w-full justify-start"
|
||||
color={
|
||||
selectedTunnelForImport === tunnel.id
|
||||
? "primary"
|
||||
: "default"
|
||||
}
|
||||
variant={
|
||||
selectedTunnelForImport === tunnel.id ? "solid" : "bordered"
|
||||
}
|
||||
onPress={() => {
|
||||
setSelectedTunnelForImport(tunnel.id);
|
||||
setTunnelSelectModalOpen(false);
|
||||
}}
|
||||
>
|
||||
{tunnel.name}
|
||||
</Button>
|
||||
))}
|
||||
</div>
|
||||
</ModalBody>
|
||||
<ModalFooter>
|
||||
<Button
|
||||
variant="light"
|
||||
onPress={() => setTunnelSelectModalOpen(false)}
|
||||
>
|
||||
取消
|
||||
</Button>
|
||||
</ModalFooter>
|
||||
</ModalContent>
|
||||
</Modal>
|
||||
|
||||
{/* 诊断结果模态框 */}
|
||||
<Modal
|
||||
backdrop="blur"
|
||||
|
||||
@@ -46,6 +46,33 @@ test("parseNyFormatData returns validation errors for invalid fields", () => {
|
||||
assert.match(result[1].error || "", /目标地址格式错误/);
|
||||
});
|
||||
|
||||
test("parseNyFormatData allows missing listen_port for auto assignment", () => {
|
||||
const input = '{"dest":["1.1.1.1:1000"],"name":"No Port"}';
|
||||
|
||||
const result = parseNyFormatData(input);
|
||||
|
||||
assert.equal(result.length, 1);
|
||||
assert.equal(result[0].error, undefined);
|
||||
assert.equal(result[0].parsed?.listen_port, null);
|
||||
});
|
||||
|
||||
test("parseNyFormatData supports ny alias fields", () => {
|
||||
const input =
|
||||
'{"dst":["2.2.2.2:2000"],"listenPort":"3000","forward_name":"Alias A"}\n{"target":"3.3.3.3:4000,4.4.4.4:5000","port":6000,"forwardName":"Alias B"}';
|
||||
|
||||
const result = parseNyFormatData(input);
|
||||
|
||||
assert.equal(result.length, 2);
|
||||
assert.equal(result[0].error, undefined);
|
||||
assert.equal(result[1].error, undefined);
|
||||
assert.deepEqual(result[0].parsed?.dest, ["2.2.2.2:2000"]);
|
||||
assert.equal(result[0].parsed?.listen_port, 3000);
|
||||
assert.equal(result[0].parsed?.name, "Alias A");
|
||||
assert.deepEqual(result[1].parsed?.dest, ["3.3.3.3:4000", "4.4.4.4:5000"]);
|
||||
assert.equal(result[1].parsed?.listen_port, 6000);
|
||||
assert.equal(result[1].parsed?.name, "Alias B");
|
||||
});
|
||||
|
||||
test("convertNyItemToForwardInput maps ny fields correctly", () => {
|
||||
const mapped = convertNyItemToForwardInput({
|
||||
dest: ["1.1.1.1:1111", "2.2.2.2:2222"],
|
||||
@@ -60,3 +87,18 @@ test("convertNyItemToForwardInput maps ny fields correctly", () => {
|
||||
strategy: "fifo",
|
||||
});
|
||||
});
|
||||
|
||||
test("convertNyItemToForwardInput keeps null inPort for auto assignment", () => {
|
||||
const mapped = convertNyItemToForwardInput({
|
||||
dest: ["1.1.1.1:1111"],
|
||||
listen_port: null,
|
||||
name: "No Port",
|
||||
});
|
||||
|
||||
assert.deepEqual(mapped, {
|
||||
name: "No Port",
|
||||
inPort: null,
|
||||
remoteAddr: "1.1.1.1:1111",
|
||||
strategy: "fifo",
|
||||
});
|
||||
});
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
export interface NyImportItem {
|
||||
dest: string[];
|
||||
listen_port: number;
|
||||
listen_port: number | null;
|
||||
name: string;
|
||||
}
|
||||
|
||||
@@ -12,6 +12,70 @@ export interface ParsedNyImportLine {
|
||||
|
||||
const ADDRESS_PATTERN = /^[^:]+:\d+$/;
|
||||
|
||||
const getAliasField = (
|
||||
item: Record<string, unknown>,
|
||||
aliases: string[],
|
||||
): unknown => {
|
||||
for (const alias of aliases) {
|
||||
if (Object.prototype.hasOwnProperty.call(item, alias)) {
|
||||
return item[alias];
|
||||
}
|
||||
}
|
||||
|
||||
return undefined;
|
||||
};
|
||||
|
||||
const normalizeDestList = (value: unknown): string[] | null => {
|
||||
if (Array.isArray(value)) {
|
||||
const normalized = value.map((itemValue) =>
|
||||
typeof itemValue === "string" ? itemValue.trim() : "",
|
||||
);
|
||||
|
||||
if (normalized.some((itemValue) => itemValue === "")) {
|
||||
return null;
|
||||
}
|
||||
|
||||
return normalized;
|
||||
}
|
||||
|
||||
if (typeof value === "string") {
|
||||
const normalized = value
|
||||
.split(",")
|
||||
.map((itemValue) => itemValue.trim())
|
||||
.filter((itemValue) => itemValue !== "");
|
||||
|
||||
return normalized.length > 0 ? normalized : null;
|
||||
}
|
||||
|
||||
return null;
|
||||
};
|
||||
|
||||
const normalizeListenPort = (value: unknown): number | null | undefined => {
|
||||
if (value === undefined || value === null || value === "") {
|
||||
return null;
|
||||
}
|
||||
|
||||
if (typeof value === "number") {
|
||||
return Number.isInteger(value) ? value : undefined;
|
||||
}
|
||||
|
||||
if (typeof value === "string") {
|
||||
const trimmed = value.trim();
|
||||
|
||||
if (!trimmed) {
|
||||
return null;
|
||||
}
|
||||
|
||||
if (!/^\d+$/.test(trimmed)) {
|
||||
return undefined;
|
||||
}
|
||||
|
||||
return Number.parseInt(trimmed, 10);
|
||||
}
|
||||
|
||||
return undefined;
|
||||
};
|
||||
|
||||
const isValidListenPort = (value: unknown): value is number => {
|
||||
return (
|
||||
typeof value === "number" &&
|
||||
@@ -27,11 +91,19 @@ const validateNyItem = (line: string, value: unknown): ParsedNyImportLine => {
|
||||
}
|
||||
|
||||
const item = value as Record<string, unknown>;
|
||||
const dest = item.dest;
|
||||
const listenPort = item.listen_port;
|
||||
const name = item.name;
|
||||
const dest = getAliasField(item, ["dest", "dst", "target", "targets"]);
|
||||
const listenPortRaw = getAliasField(item, [
|
||||
"listen_port",
|
||||
"listenPort",
|
||||
"port",
|
||||
"in_port",
|
||||
"inPort",
|
||||
]);
|
||||
const name = getAliasField(item, ["name", "forward_name", "forwardName"]);
|
||||
const normalizedDest = normalizeDestList(dest);
|
||||
const normalizedListenPort = normalizeListenPort(listenPortRaw);
|
||||
|
||||
if (!Array.isArray(dest) || dest.length === 0) {
|
||||
if (!normalizedDest || normalizedDest.length === 0) {
|
||||
return { line, error: "dest数组为空或格式错误" };
|
||||
}
|
||||
|
||||
@@ -39,16 +111,15 @@ const validateNyItem = (line: string, value: unknown): ParsedNyImportLine => {
|
||||
return { line, error: "name不能为空" };
|
||||
}
|
||||
|
||||
if (!isValidListenPort(listenPort)) {
|
||||
return { line, error: "listen_port必须为1-65535之间的数字" };
|
||||
if (normalizedListenPort === undefined) {
|
||||
return { line, error: "listen_port格式错误,应为1-65535之间的数字" };
|
||||
}
|
||||
|
||||
const normalizedDest = dest.map((itemValue) =>
|
||||
typeof itemValue === "string" ? itemValue.trim() : "",
|
||||
);
|
||||
|
||||
if (normalizedDest.some((itemValue) => itemValue === "")) {
|
||||
return { line, error: "dest中包含空地址" };
|
||||
if (
|
||||
normalizedListenPort !== null &&
|
||||
!isValidListenPort(normalizedListenPort)
|
||||
) {
|
||||
return { line, error: "listen_port必须为1-65535之间的数字" };
|
||||
}
|
||||
|
||||
const invalid = normalizedDest.find(
|
||||
@@ -63,7 +134,7 @@ const validateNyItem = (line: string, value: unknown): ParsedNyImportLine => {
|
||||
line,
|
||||
parsed: {
|
||||
dest: normalizedDest,
|
||||
listen_port: listenPort,
|
||||
listen_port: normalizedListenPort,
|
||||
name: name.trim(),
|
||||
},
|
||||
};
|
||||
|
||||
Reference in New Issue
Block a user