diff --git a/go-backend/internal/http/handler/flow_policy.go b/go-backend/internal/http/handler/flow_policy.go index c7a218d..460ee57 100644 --- a/go-backend/internal/http/handler/flow_policy.go +++ b/go-backend/internal/http/handler/flow_policy.go @@ -157,6 +157,12 @@ func (h *Handler) processPeerShareFlowFromForward(forwardID int64, serviceName s if h == nil || h.repo == nil || forwardID <= 0 { return } + + delta := item.D + item.U + if delta <= 0 { + return + } + forward, err := h.getForwardRecord(forwardID) if err != nil || forward == nil { // Forward not found in local database - might be a federation port-forward @@ -164,26 +170,22 @@ func (h *Handler) processPeerShareFlowFromForward(forwardID int64, serviceName s h.processPeerShareFlowByServiceName(serviceName, item) return } - tunnel, err := h.getTunnelRecord(forward.TunnelID) - if err != nil || tunnel == nil { - return - } tunnelName, err := h.repo.GetTunnelName(forward.TunnelID) if err != nil { + h.processPeerShareFlowByServiceName(serviceName, item) return } shareID, ok := parsePeerShareIDFromFederationTunnelName(tunnelName) if !ok { + h.processPeerShareFlowByServiceName(serviceName, item) return } - delta := item.D + item.U - if delta <= 0 { + if err := h.repo.AddPeerShareCurrentFlow(shareID, delta); err != nil { + h.processPeerShareFlowByServiceName(serviceName, item) return } - _ = h.repo.AddPeerShareCurrentFlow(shareID, delta) - share, err := h.repo.GetPeerShare(shareID) if err != nil || share == nil { return @@ -404,6 +406,11 @@ func (h *Handler) cleanOrphanedServices(nodeID int64, services []namedConfigItem if err != nil { return } + minUpdatedTime := time.Now().Add(-10 * time.Minute).UnixMilli() + hasUnboundForwardPeerRuntime, err := h.repo.HasRecentUnboundForwardPeerShareRuntimeOnNode(nodeID, minUpdatedTime) + if err != nil { + hasUnboundForwardPeerRuntime = false + } runtimeServiceSet := make(map[string]struct{}, len(runtimeServiceNames)) for _, serviceName := range runtimeServiceNames { serviceName = strings.TrimSpace(serviceName) @@ -432,6 +439,9 @@ func (h *Handler) cleanOrphanedServices(nodeID int64, services []namedConfigItem parts := strings.Split(name, "_") if len(parts) >= 3 { forwardID, err := strconv.ParseInt(parts[0], 10, 64) + if err == nil && forwardID > 0 && hasUnboundForwardPeerRuntime { + continue + } if err == nil && forwardID > 0 && !h.forwardExists(forwardID) { _, _ = h.sendNodeCommand(nodeID, "DeleteService", map[string]interface{}{"services": []string{name, parts[0] + "_" + parts[1] + "_" + parts[2], parts[0] + "_" + parts[1] + "_" + parts[2] + "_tcp", parts[0] + "_" + parts[1] + "_" + parts[2] + "_udp"}}, false, true) continue @@ -451,6 +461,9 @@ func (h *Handler) cleanOrphanedServices(nodeID int64, services []namedConfigItem continue } forwardID, err := strconv.ParseInt(parts[0], 10, 64) + if err == nil && forwardID > 0 && hasUnboundForwardPeerRuntime { + continue + } if err != nil || forwardID <= 0 || h.forwardExists(forwardID) { continue } diff --git a/go-backend/internal/http/handler/flow_policy_federation_test.go b/go-backend/internal/http/handler/flow_policy_federation_test.go index 0e38b4b..ce80980 100644 --- a/go-backend/internal/http/handler/flow_policy_federation_test.go +++ b/go-backend/internal/http/handler/flow_policy_federation_test.go @@ -177,6 +177,66 @@ func TestProcessFlowItemTracksPeerShareFlowByForwardServiceName(t *testing.T) { } } +func TestProcessFlowItemFallsBackToServiceNameWhenForwardIDCollidesAcrossPanels(t *testing.T) { + r, err := repo.Open(filepath.Join(t.TempDir(), "panel-forward-collision.db")) + if err != nil { + t.Fatalf("open repo: %v", err) + } + defer r.Close() + + now := time.Now().UnixMilli() + if err := r.CreatePeerShare(&repo.PeerShare{ + Name: "collision-share", + NodeID: 1, + Token: "collision-token", + MaxBandwidth: 0, + CurrentFlow: 0, + PortRangeStart: 31400, + PortRangeEnd: 31410, + IsActive: 1, + CreatedTime: now, + UpdatedTime: now, + }); err != nil { + t.Fatalf("create peer share: %v", err) + } + share, err := r.GetPeerShareByToken("collision-token") + if err != nil || share == nil { + t.Fatalf("load peer 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, "collision-r1", "collision-rk1", "", "forward", "", "20_2_10", "tcp", "fifo", 31401, "", 1, 1, now, now).Error; err != nil { + t.Fatalf("insert peer_share_runtime: %v", err) + } + + if err := r.DB().Exec(` + INSERT INTO tunnel(id, name, traffic_ratio, type, protocol, flow, created_time, updated_time, status, in_ip, inx) + VALUES(2, 'local-tunnel-with-colliding-forward-id', 1.0, 1, 'tls', 1, ?, ?, 1, NULL, 0) + `, now, now).Error; err != nil { + t.Fatalf("insert local tunnel: %v", err) + } + + if err := r.DB().Exec(` + INSERT INTO forward(id, user_id, user_name, name, tunnel_id, remote_addr, strategy, in_flow, out_flow, created_time, updated_time, status, inx) + VALUES(20, 1, 'local-user', 'local-f20', 2, '8.8.8.8:53', 'fifo', 0, 0, ?, ?, 1, 0) + `, now, now).Error; err != nil { + t.Fatalf("insert local forward: %v", err) + } + + h := &Handler{repo: r} + h.processFlowItem(flowItem{N: "20_2_10_tcp", U: 120, D: 80}) + + updatedShare, err := r.GetPeerShare(share.ID) + if err != nil || updatedShare == nil { + t.Fatalf("reload share: %v", err) + } + if updatedShare.CurrentFlow != 200 { + t.Fatalf("expected current_flow=200, got %d", updatedShare.CurrentFlow) + } +} + func TestProcessFlowItemSkipsPeerShareFlowWhenServiceNameIsAmbiguous(t *testing.T) { r, err := repo.Open(filepath.Join(t.TempDir(), "panel-forward-ambiguous.db")) if err != nil { @@ -299,3 +359,48 @@ func TestCleanOrphanedServicesSkipsFederationServicePrefix(t *testing.T) { h.cleanOrphanedServices(1, []namedConfigItem{{Name: "fed_svc_999_tcp"}}) } + +func TestCleanOrphanedServicesSkipsForwardPatternWhenNodeHasActivePeerShareForwardRuntime(t *testing.T) { + r, err := repo.Open(filepath.Join(t.TempDir(), "panel-cleanup-forward-runtime-empty-service.db")) + if err != nil { + t.Fatalf("open repo: %v", err) + } + defer r.Close() + + now := time.Now().UnixMilli() + if err := r.CreatePeerShare(&repo.PeerShare{ + Name: "cleanup-forward-runtime-empty-service", + NodeID: 1, + Token: "cleanup-forward-runtime-empty-service-token", + MaxBandwidth: 0, + CurrentFlow: 0, + PortRangeStart: 31420, + PortRangeEnd: 31430, + IsActive: 1, + CreatedTime: now, + UpdatedTime: now, + }); err != nil { + t.Fatalf("create peer share: %v", err) + } + share, err := r.GetPeerShareByToken("cleanup-forward-runtime-empty-service-token") + if err != nil || share == nil { + t.Fatalf("load peer 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, "cleanup-forward-empty-r1", "cleanup-forward-empty-rk1", "", "forward", "", "", "tcp", "fifo", 31421, "", 0, 1, now, now).Error; err != nil { + t.Fatalf("insert peer_share_runtime with empty service name: %v", err) + } + + h := &Handler{repo: r} + + defer func() { + if rec := recover(); rec != nil { + t.Fatalf("cleanOrphanedServices should skip forward-pattern services when active peer-share forward runtime exists; got panic: %v", rec) + } + }() + + h.cleanOrphanedServices(share.NodeID, []namedConfigItem{{Name: "20_2_10_tcp"}}) +} diff --git a/go-backend/internal/http/handler/jobs_test.go b/go-backend/internal/http/handler/jobs_test.go index aeef09b..98a2315 100644 --- a/go-backend/internal/http/handler/jobs_test.go +++ b/go-backend/internal/http/handler/jobs_test.go @@ -69,6 +69,13 @@ func TestRunResetAndExpiryJobResetsFlowAndDisablesExpiredRecords(t *testing.T) { t.Fatalf("insert expired user: %v", err) } + if err := r.DB().Exec(` + INSERT INTO user(id, user, pwd, role_id, exp_time, flow, in_flow, out_flow, flow_reset_time, num, created_time, updated_time, status) + VALUES(3, 'non_expiring_user', 'x', 1, 0, 100, 1000, 2000, 15, 1, ?, ?, 1) + `, nowMs, nowMs).Error; err != nil { + t.Fatalf("insert non-expiring user: %v", err) + } + if err := r.DB().Exec(` INSERT INTO tunnel(id, name, traffic_ratio, type, protocol, flow, created_time, updated_time, status, in_ip, inx) VALUES(1, 't1', 1.0, 1, 'tls', 1, ?, ?, 1, NULL, 0) @@ -83,6 +90,13 @@ func TestRunResetAndExpiryJobResetsFlowAndDisablesExpiredRecords(t *testing.T) { t.Fatalf("insert expired user_tunnel: %v", err) } + if err := r.DB().Exec(` + INSERT INTO user_tunnel(id, user_id, tunnel_id, speed_id, num, flow, in_flow, out_flow, flow_reset_time, exp_time, status) + VALUES(11, 3, 1, NULL, 1, 1, 300, 400, 15, 0, 1) + `).Error; err != nil { + t.Fatalf("insert non-expiring user_tunnel: %v", err) + } + if err := r.DB().Exec(` INSERT INTO forward(id, user_id, user_name, name, tunnel_id, remote_addr, strategy, in_flow, out_flow, created_time, updated_time, status, inx) VALUES(20, 2, 'expired_user', 'f1', 1, '1.1.1.1:443', 'fifo', 0, 0, ?, ?, 1, 0) @@ -90,6 +104,13 @@ func TestRunResetAndExpiryJobResetsFlowAndDisablesExpiredRecords(t *testing.T) { t.Fatalf("insert forward: %v", err) } + if err := r.DB().Exec(` + INSERT INTO forward(id, user_id, user_name, name, tunnel_id, remote_addr, strategy, in_flow, out_flow, created_time, updated_time, status, inx) + VALUES(21, 3, 'non_expiring_user', 'f2', 1, '1.1.1.1:443', 'fifo', 0, 0, ?, ?, 1, 1) + `, nowMs, nowMs).Error; err != nil { + t.Fatalf("insert non-expiring forward: %v", err) + } + h.runResetAndExpiryJob(now) userIn, userOut, userStatus := mustQueryInt64Int64Int(t, r, `SELECT in_flow, out_flow, status FROM user WHERE id = 2`) @@ -106,4 +127,19 @@ func TestRunResetAndExpiryJobResetsFlowAndDisablesExpiredRecords(t *testing.T) { if forwardStatus != 0 { t.Fatalf("expected forward status=0 after expiry handling, got %d", forwardStatus) } + + nonExpUserStatus := mustQueryInt(t, r, `SELECT status FROM user WHERE id = 3`) + if nonExpUserStatus != 1 { + t.Fatalf("expected non-expiring user to remain enabled, got status=%d", nonExpUserStatus) + } + + nonExpTunnelStatus := mustQueryInt(t, r, `SELECT status FROM user_tunnel WHERE id = 11`) + if nonExpTunnelStatus != 1 { + t.Fatalf("expected non-expiring user_tunnel to remain enabled, got status=%d", nonExpTunnelStatus) + } + + nonExpForwardStatus := mustQueryInt(t, r, `SELECT status FROM forward WHERE id = 21`) + if nonExpForwardStatus != 1 { + t.Fatalf("expected non-expiring forward to remain enabled, got status=%d", nonExpForwardStatus) + } } diff --git a/go-backend/internal/store/repo/repository.go b/go-backend/internal/store/repo/repository.go index e378bfb..0113be9 100644 --- a/go-backend/internal/store/repo/repository.go +++ b/go-backend/internal/store/repo/repository.go @@ -1304,6 +1304,20 @@ func (r *Repository) ListActiveForwardPeerShareRuntimeServiceNamesByNode(nodeID return names, nil } +func (r *Repository) HasRecentUnboundForwardPeerShareRuntimeOnNode(nodeID int64, minUpdatedTime int64) (bool, error) { + if r == nil || r.db == nil { + return false, errors.New("repository not initialized") + } + var count int64 + err := r.db.Model(&model.PeerShareRuntime{}). + Where("node_id = ? AND status = 1 AND role = ? AND applied = 0 AND updated_time >= ? AND (service_name = '' OR service_name IS NULL)", nodeID, "forward", minUpdatedTime). + Count(&count).Error + if err != nil { + return false, err + } + return count > 0, nil +} + func (r *Repository) GetActiveForwardPeerShareRuntimeByPort(shareID int64, port int) (*model.PeerShareRuntime, error) { if r == nil || r.db == nil { return nil, errors.New("repository not initialized") @@ -2327,7 +2341,7 @@ func (r *Repository) ListExpiredActiveUserIDs(nowMs int64) ([]int64, error) { } var ids []int64 err := r.db.Model(&model.User{}). - Where("role_id != 0 AND status = 1 AND exp_time IS NOT NULL AND exp_time < ?", nowMs). + Where("role_id != 0 AND status = 1 AND exp_time > 0 AND exp_time < ?", nowMs). Pluck("id", &ids).Error if err != nil { return nil, err @@ -2347,7 +2361,7 @@ func (r *Repository) ListExpiredActiveUserTunnels(nowMs int64) ([]model.ExpiredU return nil, errors.New("repository not initialized") } var uts []model.UserTunnel - err := r.db.Where("status = 1 AND exp_time IS NOT NULL AND exp_time < ?", nowMs).Find(&uts).Error + err := r.db.Where("status = 1 AND exp_time > 0 AND exp_time < ?", nowMs).Find(&uts).Error if err != nil { return nil, err } diff --git a/vite-frontend/src/components/search-bar.tsx b/vite-frontend/src/components/search-bar.tsx index 2a333c6..3a8643e 100644 --- a/vite-frontend/src/components/search-bar.tsx +++ b/vite-frontend/src/components/search-bar.tsx @@ -55,7 +55,6 @@ export function SearchBar({ transition={{ duration: 0.18, ease: [0.25, 0.46, 0.45, 0.94] }} > ) { +function AlertTitle({ + className, + children, + ...props +}: React.ComponentProps<"h5">) { return (
+ > + {children} + ); } diff --git a/vite-frontend/src/components/ui/card.tsx b/vite-frontend/src/components/ui/card.tsx index 525a9b3..de46041 100644 --- a/vite-frontend/src/components/ui/card.tsx +++ b/vite-frontend/src/components/ui/card.tsx @@ -25,7 +25,11 @@ function CardHeader({ className, ...props }: React.ComponentProps<"div">) { ); } -function CardTitle({ className, ...props }: React.ComponentProps<"h3">) { +function CardTitle({ + className, + children, + ...props +}: React.ComponentProps<"h3">) { return (+ v{version} + {updateAvailable && latestUpdateVersion && ( + + {latestUpdateVersion} + + )} +
++ Powered by{" "} + + FLVX + +
+- v{siteConfig.version} -
@@ -405,17 +403,12 @@ export default function AdminLayout({- Powered by{" "} - - FLVX - -
++ 更新通道 +
++ 稳定版仅匹配纯数字版本;开发版仅匹配包含 alpha / beta / rc + 的版本。 +
+按用户筛选
按隧道筛选
- Powered by{" "} - - FLVX - -
-- v{isWebViewFunc() ? siteConfig.app_version : siteConfig.version} -
-+ 版本提示会根据该通道检查最新版本。 +
+