diff --git a/internal/web/node_contract_test.go b/internal/web/node_contract_test.go index d76d7a939..072337f1a 100644 --- a/internal/web/node_contract_test.go +++ b/internal/web/node_contract_test.go @@ -362,6 +362,7 @@ var remoteMethodsOutsideContract = map[string]string{ "AdoptInboundAlias": "local alias bookkeeping", "AdoptedInboundAliases": "local alias bookkeeping", "AdvancePushedInbound": "local fingerprint bookkeeping", + "ForgetPushedInbound": "local fingerprint bookkeeping", "UpdatePanel": "replaces the node binary; node-sync is denied it on purpose (#6201)", } diff --git a/internal/web/runtime/remote.go b/internal/web/runtime/remote.go index 9044c73e4..479669d7b 100644 --- a/internal/web/runtime/remote.go +++ b/internal/web/runtime/remote.go @@ -528,6 +528,17 @@ func (r *Remote) RecordAdoptedInbound(ib *model.Inbound) { r.recordPushedInbound(ib) } +// ForgetPushedInbound drops the reconcile-skip fingerprint once the node is seen +// without the payload it stamped, so the next reconcile re-sends the inbound. +func (r *Remote) ForgetPushedInbound(tag string) { + prefix := nodeInboundTagPrefix(r.node.Id) + bare := strings.TrimPrefix(tag, prefix) + r.mu.Lock() + delete(r.pushedFP, bare) + delete(r.pushedFP, prefix+bare) + r.mu.Unlock() +} + // AdoptInboundAlias records a deployed alias without mutating either panel. // The runtime association is rediscovered after a master restart. func (r *Remote) AdoptInboundAlias(ib *model.Inbound, remote RemoteInboundOption) { diff --git a/internal/web/service/client_sync_orphan_test.go b/internal/web/service/client_sync_orphan_test.go index bb057e974..1a0a1b8ff 100644 --- a/internal/web/service/client_sync_orphan_test.go +++ b/internal/web/service/client_sync_orphan_test.go @@ -30,8 +30,8 @@ func backdateOrphanMark(t *testing.T, db *gorm.DB, email string) { } } -// The merge must soft-orphan, not delete: everything stays recoverable until -// the grace period has elapsed and the reaper confirms nothing reclaimed it. +// A partial snapshot (node alive, still serving another client) authoritatively drops one; +// the merge soft-orphans, recoverable until the grace elapses and the reaper confirms it. func TestSyncOrphanSurvivesMergeUntilGraceElapses(t *testing.T) { db := initTrafficTestDB(t) svc := &InboundService{} @@ -40,16 +40,20 @@ func TestSyncOrphanSurvivesMergeUntilGraceElapses(t *testing.T) { seedNodeRow(t, db, &model.Node{Id: 1, Name: "n1", Address: "127.0.0.1", Port: 2096, ApiToken: "tok", Enable: true}) const email = "gone@x" - createNodeInboundWithClient(t, db, 1, "n1-in", 41001, email) - settings := fmt.Sprintf(`{"clients":[{"email":%q,"enable":true}]}`, email) - syncNodeWithSettings(t, svc, 1, "n1-in", settings, + const keep = "keep@x" + createNodeInboundWithClient(t, db, 1, "n1-in", 41001, keep) + bothSettings := fmt.Sprintf(`{"clients":[{"email":%q,"enable":true},{"email":%q,"enable":true}]}`, keep, email) + syncNodeWithSettings(t, svc, 1, "n1-in", bothSettings, + xray.ClientTraffic{Email: keep, Enable: true}, xray.ClientTraffic{Email: email, Up: 5, Down: 5, Enable: true}) if rec, traf := countClientRows(t, db, email); rec != 1 || traf != 1 { t.Fatalf("setup: clients=%d client_traffics=%d, want 1/1", rec, traf) } - if _, err := svc.setRemoteTrafficLocked(1, snapshotWithoutClients(t, "n1-in"), false, false); err != nil { + keepOnly := fmt.Sprintf(`{"clients":[{"email":%q,"enable":true}]}`, keep) + if _, err := svc.setRemoteTrafficLocked(1, snapshotWithClients(t, "n1-in", keepOnly, + xray.ClientTraffic{Email: keep, Enable: true}), false, false); err != nil { t.Fatalf("orphaning merge: %v", err) } if rec, traf := countClientRows(t, db, email); rec != 1 || traf != 1 { @@ -84,7 +88,7 @@ func TestSyncOrphanSurvivesMergeUntilGraceElapses(t *testing.T) { } } -// A client the node reports again was never gone: clearing the mark is what +// A client the node reports again (partial snapshot) was never gone: clearing the mark // turns a bad merge into a recoverable blip instead of a delayed deletion. func TestSyncOrphanMarkClearedOnReattach(t *testing.T) { db := initTrafficTestDB(t) @@ -94,19 +98,24 @@ func TestSyncOrphanMarkClearedOnReattach(t *testing.T) { seedNodeRow(t, db, &model.Node{Id: 1, Name: "n1", Address: "127.0.0.1", Port: 2096, ApiToken: "tok", Enable: true}) const email = "flaky@x" - createNodeInboundWithClient(t, db, 1, "n1-in", 41001, email) - settings := fmt.Sprintf(`{"clients":[{"email":%q,"enable":true}]}`, email) - syncNodeWithSettings(t, svc, 1, "n1-in", settings, + const keep = "keep@x" + createNodeInboundWithClient(t, db, 1, "n1-in", 41001, keep) + bothSettings := fmt.Sprintf(`{"clients":[{"email":%q,"enable":true},{"email":%q,"enable":true}]}`, keep, email) + syncNodeWithSettings(t, svc, 1, "n1-in", bothSettings, + xray.ClientTraffic{Email: keep, Enable: true}, xray.ClientTraffic{Email: email, Up: 5, Down: 5, Enable: true}) - if _, err := svc.setRemoteTrafficLocked(1, snapshotWithoutClients(t, "n1-in"), false, false); err != nil { + keepOnly := fmt.Sprintf(`{"clients":[{"email":%q,"enable":true}]}`, keep) + if _, err := svc.setRemoteTrafficLocked(1, snapshotWithClients(t, "n1-in", keepOnly, + xray.ClientTraffic{Email: keep, Enable: true}), false, false); err != nil { t.Fatalf("orphaning merge: %v", err) } if readOrphanMark(t, db, email) <= 0 { t.Fatal("setup: expected the merge to mark the client") } - syncNodeWithSettings(t, svc, 1, "n1-in", settings, + syncNodeWithSettings(t, svc, 1, "n1-in", bothSettings, + xray.ClientTraffic{Email: keep, Enable: true}, xray.ClientTraffic{Email: email, Up: 6, Down: 6, Enable: true}) if orphanedAt := readOrphanMark(t, db, email); orphanedAt != 0 { diff --git a/internal/web/service/inbound_amneziawg_test.go b/internal/web/service/inbound_amneziawg_test.go index bdf1d2106..cb0c65da7 100644 --- a/internal/web/service/inbound_amneziawg_test.go +++ b/internal/web/service/inbound_amneziawg_test.go @@ -287,6 +287,9 @@ func TestNormalizeAmneziaWGSettings_CanonicalizesClientAllowedIPs(t *testing.T) } func TestGetAmneziaWGLogs_ClampsCountAndFiltersEvents(t *testing.T) { + // GetAmneziaWGLogs appends peer handshake activity, which reads the DB; + // own a throwaway one so -shuffle can't leave us the global nil DB. + setupConflictDB(t) logger.InitLogger(logging.DEBUG) logger.Info("amneziawg: started interface awg1 for inbound 1") logger.Info("xray: unrelated line that must never show up here") diff --git a/internal/web/service/inbound_node.go b/internal/web/service/inbound_node.go index 1feec0f46..8d4a6d209 100644 --- a/internal/web/service/inbound_node.go +++ b/internal/web/service/inbound_node.go @@ -427,6 +427,20 @@ func adoptedWireInbound(c, snapIb *model.Inbound, adoptedSettings string) *model return &a } +// snapshotDropsEveryHubClient reports a node that lists no clients where the hub +// still links some: a reset or half-started node, never an authoritative removal. +func snapshotDropsEveryHubClient(tx *gorm.DB, inboundID int, wireSettings string) bool { + clients, err := ParseInboundSettingsClients(wireSettings) + if err != nil || len(clients) > 0 { + return false + } + var links int64 + if err := tx.Table("client_inbounds").Where("inbound_id = ?", inboundID).Count(&links).Error; err != nil { + return false + } + return links > 0 +} + // clientEmailsOwnedElsewhere returns the emails attached only to inbounds of // other nodes: email is unique, so adopting one would overwrite a client this // node does not serve. Attached nowhere means soft-orphaned, hence adoptable. @@ -644,6 +658,7 @@ func (s *InboundService) setRemoteTrafficLocked(nodeID int, snap *runtime.Traffi wireSettings string } var pendingAdopts []pendingAdopt + degradedInbounds := map[int]string{} newInboundIDs := make(map[int]struct{}) @@ -782,7 +797,9 @@ func (s *InboundService) setRemoteTrafficLocked(nodeID int, snap *runtime.Traffi adoptedSettings = deduped } updates := map[string]any{} - if !dirty { + if !dirty && snapshotDropsEveryHubClient(tx, c.Id, adoptedSettings) { + degradedInbounds[c.Id] = c.Tag + } else if !dirty { // Defer lifecycle lift until after client_traffics absorbs this tick's // deltas so quota stale-disable matches SQL (#6228). pendingAdopts = append(pendingAdopts, pendingAdopt{ @@ -1121,6 +1138,9 @@ func (s *InboundService) setRemoteTrafficLocked(nodeID int, snap *runtime.Traffi if k.inboundID != c.Id { continue } + if _, degraded := degradedInbounds[c.Id]; degraded { + continue + } if _, kept := snapEmails[k.email]; kept { continue } @@ -1228,6 +1248,13 @@ func (s *InboundService) setRemoteTrafficLocked(nodeID int, snap *runtime.Traffi applyMasterClientLifecycle(&clients[i], existing, csPtr) filtered = append(filtered, clients[i]) } + // A degraded node (reset/restart/removal) reports zero clients for an inbound the + // hub populates; adopting it empties links and ReapSyncOrphans deletes shared clients (#6734). + if _, degraded := degradedInbounds[c.Id]; degraded { + logger.Warningf("setRemoteTraffic: node %d reported zero clients for tag %q while the hub has %d attached — keeping them and re-pushing", nodeID, snapIb.Tag, len(oldEmailsRows)) + syncFailedInbounds[c.Id] = struct{}{} + continue + } localEmails := make([]string, 0, len(filtered)) for i := range filtered { if filtered[i].Email != "" { @@ -1337,6 +1364,21 @@ func (s *InboundService) setRemoteTrafficLocked(nodeID int, snap *runtime.Traffi } committed = true + if len(degradedInbounds) > 0 { + if mgr := runtime.GetManager(); mgr != nil { + if rt, rtErr := mgr.RuntimeFor(&nodeID); rtErr == nil { + if rem, ok := rt.(*runtime.Remote); ok { + for _, tag := range degradedInbounds { + rem.ForgetPushedInbound(tag) + } + } + } + } + if err := (&NodeService{}).MarkNodeDirty(nodeID); err != nil { + logger.Warningf("setRemoteTraffic: mark node %d dirty after an empty snapshot failed: %v", nodeID, err) + } + } + if lifecycleLifted && !dirty { var already model.Node if err := database.GetDB().Select("config_dirty").Where("id = ?", nodeID).First(&already).Error; err == nil && already.ConfigDirty { diff --git a/internal/web/service/node_degraded_snapshot_test.go b/internal/web/service/node_degraded_snapshot_test.go new file mode 100644 index 000000000..a9cf4c9f4 --- /dev/null +++ b/internal/web/service/node_degraded_snapshot_test.go @@ -0,0 +1,205 @@ +package service + +import ( + "context" + "encoding/json" + "net/http" + "net/http/httptest" + "strings" + "sync" + "testing" + + "github.com/mhsanaei/3x-ui/v3/internal/database/model" + "github.com/mhsanaei/3x-ui/v3/internal/web/runtime" + "github.com/mhsanaei/3x-ui/v3/internal/xray" + + "gorm.io/gorm" +) + +// linkCount returns how many client_inbounds links a client currently has, +// across every inbound — the value ReapSyncOrphans checks before deleting. +func linkCount(t *testing.T, db *gorm.DB, email string) int64 { + t.Helper() + var n int64 + if err := db.Table("client_inbounds"). + Joins("JOIN clients ON clients.id = client_inbounds.client_id"). + Where("clients.email = ?", email). + Count(&n).Error; err != nil { + t.Fatalf("count links for %q: %v", email, err) + } + return n +} + +// A degraded node reporting zero clients for an inbound the hub populates must +// keep its links and never orphan-mark, or SyncInbound/ReapSyncOrphans delete the row. +func TestSetRemoteTraffic_EmptySnapshotKeepsClients(t *testing.T) { + db := initTrafficTestDB(t) + svc := &InboundService{} + + seedNodeRow(t, db, &model.Node{Id: 1, Name: "n1", Address: "127.0.0.1", Port: 2096, ApiToken: "tok", Enable: true}) + createNodeInboundWithClient(t, db, 1, "n1-in", 41001, "svc@x") + + settings := `{"clients":[{"email":"svc@x","enable":true}]}` + if _, err := svc.setRemoteTrafficLocked(1, snapshotWithClients(t, "n1-in", settings, + xray.ClientTraffic{Email: "svc@x", Enable: true}), false, false); err != nil { + t.Fatalf("seed sync: %v", err) + } + if n := linkCount(t, db, "svc@x"); n != 1 { + t.Fatalf("setup: svc@x links=%d, want 1", n) + } + + // The node returns an empty snapshot — the trigger that deleted real clients. + if _, err := svc.setRemoteTrafficLocked(1, snapshotWithoutClients(t, "n1-in"), false, false); err != nil { + t.Fatalf("empty-snapshot sync: %v", err) + } + + if rec, _ := countClientRows(t, db, "svc@x"); rec != 1 { + t.Fatalf("empty snapshot deleted the client row: clients=%d, want 1", rec) + } + if n := linkCount(t, db, "svc@x"); n != 1 { + t.Fatalf("empty snapshot stripped the client link: links=%d, want 1", n) + } + if at := readOrphanMark(t, db, "svc@x"); at != 0 { + t.Fatalf("empty snapshot orphan-marked a live client: sync_orphaned_at=%d, want 0", at) + } + // The hub must keep the client in the inbound's settings, or reconcile re-pushes + // an empty blob to the node and the clients never come back (#6734). + var ib model.Inbound + if err := db.Where("tag = ?", "n1-in").First(&ib).Error; err != nil { + t.Fatalf("read central inbound: %v", err) + } + if !strings.Contains(ib.Settings, "svc@x") { + t.Fatalf("empty snapshot blanked the inbound settings: %q", ib.Settings) + } +} + +// The guard is narrow: a snapshot still carrying a client is authoritative, so a +// client the node really dropped is unlinked and orphan-marked; only all-empty is degraded. +func TestSetRemoteTraffic_PartialSnapshotStillPrunes(t *testing.T) { + db := initTrafficTestDB(t) + svc := &InboundService{} + + seedNodeRow(t, db, &model.Node{Id: 1, Name: "n1", Address: "127.0.0.1", Port: 2096, ApiToken: "tok", Enable: true}) + createNodeInboundWithClient(t, db, 1, "n1-in", 41001, "keep@x") + + bothSettings := `{"clients":[{"email":"keep@x","enable":true},{"email":"drop@x","enable":true}]}` + if _, err := svc.setRemoteTrafficLocked(1, snapshotWithClients(t, "n1-in", bothSettings, + xray.ClientTraffic{Email: "keep@x", Enable: true}, + xray.ClientTraffic{Email: "drop@x", Enable: true}), false, false); err != nil { + t.Fatalf("seed sync: %v", err) + } + if n := linkCount(t, db, "drop@x"); n != 1 { + t.Fatalf("setup: drop@x links=%d, want 1", n) + } + + // Node now reports only keep@x — drop@x was genuinely removed there. + keepOnlySettings := `{"clients":[{"email":"keep@x","enable":true}]}` + if _, err := svc.setRemoteTrafficLocked(1, snapshotWithClients(t, "n1-in", keepOnlySettings, + xray.ClientTraffic{Email: "keep@x", Enable: true}), false, false); err != nil { + t.Fatalf("partial-snapshot sync: %v", err) + } + + if n := linkCount(t, db, "keep@x"); n != 1 { + t.Fatalf("partial snapshot dropped a reported client: keep@x links=%d, want 1", n) + } + if n := linkCount(t, db, "drop@x"); n != 0 { + t.Fatalf("partial snapshot kept an unreported client linked: drop@x links=%d, want 0", n) + } + if at := readOrphanMark(t, db, "drop@x"); at <= 0 { + t.Fatalf("partial snapshot did not orphan-mark the removed client: sync_orphaned_at=%d, want >0", at) + } +} + +// Keeping the hub's settings is not recovery: the node is only healed once the +// hub actually re-pushes them, which needs a dirty node and a stale fingerprint. +func TestSetRemoteTraffic_EmptySnapshotRepushesHubClients(t *testing.T) { + db := initTrafficTestDB(t) + svc := &InboundService{} + + var mu sync.Mutex + var pushed []string + writeOK := func(w http.ResponseWriter, obj any) { + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(map[string]any{"success": true, "msg": "", "obj": obj}) + } + mux := http.NewServeMux() + mux.HandleFunc("/panel/api/inbounds/list", func(w http.ResponseWriter, _ *http.Request) { + writeOK(w, []map[string]any{{"id": 7, "tag": "deg-in", "port": 41001, "protocol": "vless"}}) + }) + mux.HandleFunc("/panel/api/inbounds/update/", func(w http.ResponseWriter, r *http.Request) { + if err := r.ParseForm(); err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + mu.Lock() + pushed = append(pushed, r.PostForm.Get("settings")) + mu.Unlock() + writeOK(w, nil) + }) + ts := httptest.NewServer(mux) + t.Cleanup(ts.Close) + + node := reconcileTestNode(t, ts, "deg-node", "all", nil) + settings := `{"clients":[{"email":"svc@x","enable":true,"id":"11111111-1111-1111-1111-111111111111"}]}` + nid := node.Id + if err := db.Create(&model.Inbound{UserId: 1, Tag: "deg-in", Enable: true, Port: 41001, Protocol: model.VLESS, NodeID: &nid, Settings: settings}).Error; err != nil { + t.Fatalf("create inbound: %v", err) + } + rt := runtime.NewRemote(node, nil) + mgr := runtime.NewManager(runtime.LocalDeps{}) + mgr.SetRuntimeOverride(node.Id, rt) + runtime.SetManager(mgr) + t.Cleanup(func() { runtime.SetManager(nil) }) + + if _, err := svc.setRemoteTrafficLocked(node.Id, snapshotWithClients(t, "deg-in", settings, + xray.ClientTraffic{Email: "svc@x", Enable: true}), false, false); err != nil { + t.Fatalf("seed sync: %v", err) + } + if err := svc.ReconcileNode(context.Background(), rt, node); err != nil { + t.Fatalf("first reconcile: %v", err) + } + mu.Lock() + pushed = nil + mu.Unlock() + + if _, err := svc.setRemoteTrafficLocked(node.Id, snapshotWithoutClients(t, "deg-in"), false, false); err != nil { + t.Fatalf("empty-snapshot sync: %v", err) + } + var after model.Node + if err := db.Where("id = ?", node.Id).First(&after).Error; err != nil { + t.Fatalf("reload node: %v", err) + } + if !after.ConfigDirty { + t.Fatal("empty snapshot left the node clean: the job never reconciles it, so the node stays without its clients") + } + if err := svc.ReconcileNode(context.Background(), rt, &after); err != nil { + t.Fatalf("reconcile after empty snapshot: %v", err) + } + mu.Lock() + defer mu.Unlock() + if len(pushed) != 1 || !strings.Contains(pushed[0], "svc@x") { + t.Fatalf("reconcile after empty snapshot pushed %d settings payload(s) %q, want one carrying svc@x", len(pushed), pushed) + } +} + +// The traffic a client used while its node reported nothing must still count +// once the node reports it again. +func TestSetRemoteTraffic_EmptySnapshotKeepsTrafficBaseline(t *testing.T) { + db := initTrafficTestDB(t) + svc := &InboundService{} + + seedNodeRow(t, db, &model.Node{Id: 1, Name: "n1", Address: "127.0.0.1", Port: 2096, ApiToken: "tok", Enable: true}) + createNodeInboundWithClient(t, db, 1, "n1-in", 41001, "svc@x") + settings := `{"clients":[{"email":"svc@x","enable":true}]}` + for _, used := range []int64{100, 200} { + syncNodeWithSettings(t, svc, 1, "n1-in", settings, xray.ClientTraffic{Email: "svc@x", Up: used, Down: used, Enable: true}) + } + before := readTraffic(t, db, "svc@x") + + if _, err := svc.setRemoteTrafficLocked(1, snapshotWithoutClients(t, "n1-in"), false, false); err != nil { + t.Fatalf("empty-snapshot sync: %v", err) + } + syncNodeWithSettings(t, svc, 1, "n1-in", settings, xray.ClientTraffic{Email: "svc@x", Up: 250, Down: 250, Enable: true}) + + assertUpDown(t, readTraffic(t, db, "svc@x"), before.Up+50, before.Down+50, "after the node recovered") +}