diff --git a/internal/web/service/inbound_node.go b/internal/web/service/inbound_node.go index 6c56a8cd8..327659ef1 100644 --- a/internal/web/service/inbound_node.go +++ b/internal/web/service/inbound_node.go @@ -698,7 +698,12 @@ func (s *InboundService) setRemoteTrafficLocked(nodeID int, snap *runtime.Traffi var compatible []*model.Inbound for i := range central { candidate := ¢ral[i] - if candidate.OriginNodeGuid == origin && + // An empty origin is a master-created inbound not yet tag-matched: hosted here. + candidateOrigin := candidate.OriginNodeGuid + if candidateOrigin == "" { + candidateOrigin = selfKey + } + if candidateOrigin == origin && candidate.Port == snapIb.Port && candidate.Protocol == snapIb.Protocol && strings.TrimSpace(candidate.Listen) == strings.TrimSpace(snapIb.Listen) { @@ -862,6 +867,24 @@ func (s *InboundService) setRemoteTrafficLocked(nodeID int, snap *runtime.Traffi } } + // The inbound each reported email now lives under, so a swept inbound's + // accumulator rows follow the email instead of being dropped and re-seeded at 0. + snapEmailHome := make(map[string]int, len(snapEmailsAll)) + for _, snapIb := range snap.Inbounds { + if snapIb == nil { + continue + } + home, ok := tagToCentral[snapIb.Tag] + if !ok { + continue + } + for i := range snapIb.ClientStats { + if _, taken := snapEmailHome[snapIb.ClientStats[i].Email]; !taken { + snapEmailHome[snapIb.ClientStats[i].Email] = home.Id + } + } + } + for _, c := range central { if dirty { continue @@ -886,7 +909,7 @@ func (s *InboundService) setRemoteTrafficLocked(nodeID int, snap *runtime.Traffi if unmanagedTag(c.Tag) { continue } - // This drops the central inbound and its clients' traffic history, so say + // This drops the central inbound and its unreported clients' history, so say // so: silent removal is indistinguishable from an inbound never arriving. logger.Warningf("setRemoteTraffic: node %d no longer reports inbound %q (id %d, port %d) — removing it centrally", nodeID, c.Tag, c.Id, c.Port) var goneEmails []string @@ -922,11 +945,30 @@ func (s *InboundService) setRemoteTrafficLocked(nodeID int, snap *runtime.Traffi return false, sErr } delEmails := make([]string, 0, len(goneEmails)) + rehome := make(map[int][]string) for _, e := range goneEmails { + if home, still := snapEmailHome[e]; still { + rehome[home] = append(rehome[home], e) + if row, ok := centralCS[csKey{c.Id, e}]; ok { + delete(centralCS, csKey{c.Id, e}) + row.InboundId = home + centralCS[csKey{home, e}] = row + } + continue + } if !sharedEmails[strings.ToLower(strings.TrimSpace(e))] { delEmails = append(delEmails, e) } } + for home, emails := range rehome { + for _, batch := range chunkStrings(emails, sqliteMaxVars) { + if err := tx.Model(xray.ClientTraffic{}). + Where("inbound_id = ? AND email IN ?", c.Id, batch). + Update("inbound_id", home).Error; err != nil { + return false, err + } + } + } for _, batch := range chunkStrings(delEmails, sqliteMaxVars) { if err := tx.Where("inbound_id = ? AND email IN ?", c.Id, batch). Delete(&xray.ClientTraffic{}).Error; err != nil { diff --git a/internal/web/service/node_inbound_replace_test.go b/internal/web/service/node_inbound_replace_test.go new file mode 100644 index 000000000..9d580500d --- /dev/null +++ b/internal/web/service/node_inbound_replace_test.go @@ -0,0 +1,90 @@ +package service + +import ( + "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" +) + +// A snapshot tag that matches neither the central tag nor an alias replaces the +// inbound in one tick; the client's accumulated usage must survive the swap. +func TestSetRemoteTraffic_InboundReplacedKeepsClientHistory(t *testing.T) { + db := initTrafficTestDB(t) + const nodeID = 2 + if err := db.Create(&model.Node{Id: nodeID, Name: "node", Address: "10.0.0.2", Port: 2053, ApiToken: "t", Guid: "node-guid"}).Error; err != nil { + t.Fatalf("create node: %v", err) + } + const email = "baba" + settings := `{"clients":[{"email":"baba","enable":true}]}` + createNodeInboundWithClient(t, db, nodeID, "n2-in-2053-tcp", 2053, email) + svc := &InboundService{} + sync := func(up, down int64) { + t.Helper() + snap := &runtime.TrafficSnapshot{Inbounds: []*model.Inbound{{ + Tag: "in-2053-tcp", OriginNodeGuid: "node-guid", Enable: true, Port: 2053, Protocol: model.VLESS, + Settings: settings, ClientStats: []xray.ClientTraffic{{Email: email, Up: up, Down: down, Enable: true}}, + }}} + if _, err := svc.setRemoteTrafficLocked(nodeID, snap, false, false); err != nil { + t.Fatalf("setRemoteTrafficLocked: %v", err) + } + } + + sync(100, 200) + sync(600, 1200) + assertUpDown(t, readTraffic(t, db, email), 500, 1000, "before the replace") + + // Desync the central row so the next snapshot neither tag-matches nor aliases it. + if err := db.Model(&model.Inbound{}).Where("node_id = ?", nodeID). + Updates(map[string]any{"tag": "n2-legacy", "origin_node_guid": "stale-guid"}).Error; err != nil { + t.Fatalf("desync central inbound: %v", err) + } + + sync(650, 1300) + assertUpDown(t, readTraffic(t, db, email), 550, 1100, "replace tick") + sync(700, 1400) + assertUpDown(t, readTraffic(t, db, email), 600, 1200, "tick after the replace") + + var ib model.Inbound + if err := db.Where("node_id = ?", nodeID).First(&ib).Error; err != nil { + t.Fatalf("read node inbound: %v", err) + } + if ib.Tag != "in-2053-tcp" { + t.Fatalf("fixture did not replace the inbound: surviving tag %q", ib.Tag) + } + if ct := readTraffic(t, db, email); ct.InboundId != ib.Id { + t.Errorf("client_traffics.inbound_id = %d, want the surviving inbound %d", ct.InboundId, ib.Id) + } +} + +// A node inbound created on the master has no origin until its first tag match; +// that empty origin is still this node, so a renamed tag aliases instead of replacing. +func TestSetRemoteTraffic_AliasesInboundWithEmptyOrigin(t *testing.T) { + db := initTrafficTestDB(t) + const nodeID = 2 + if err := db.Create(&model.Node{Id: nodeID, Name: "node", Address: "10.0.0.2", Port: 2053, ApiToken: "t", Guid: "node-guid"}).Error; err != nil { + t.Fatalf("create node: %v", err) + } + createNodeInbound(t, db, nodeID, "n2-in-2053-tcp", 2053) + var before model.Inbound + if err := db.Where("node_id = ?", nodeID).First(&before).Error; err != nil { + t.Fatalf("read node inbound: %v", err) + } + + snap := &runtime.TrafficSnapshot{Inbounds: []*model.Inbound{{ + Tag: "in-2053-tcp-2", OriginNodeGuid: "node-guid", Enable: true, Port: 2053, Protocol: model.VLESS, + Settings: `{"clients":[]}`, + }}} + if _, err := (&InboundService{}).setRemoteTrafficLocked(nodeID, snap, false, false); err != nil { + t.Fatalf("setRemoteTrafficLocked: %v", err) + } + + var rows []model.Inbound + if err := db.Where("node_id = ?", nodeID).Find(&rows).Error; err != nil { + t.Fatalf("list node inbounds: %v", err) + } + if len(rows) != 1 || rows[0].Id != before.Id { + t.Fatalf("node inbounds = %#v, want only the original id %d", rows, before.Id) + } +}