fix(node): keep client traffic history when a node inbound is replaced

When a node snapshot reports an inbound under a tag that matches neither
the central tag nor an origin/port alias, SetRemoteTraffic adopts it as a
new central inbound and, in the same tick, sweeps the old one. The sweep
deleted the old inbound's client_traffics rows even though the snapshot
still reported those emails under the new inbound. existingEmails was
computed before the sweep, so the rows were not recreated that tick; on
the next tick the inbound was no longer new, the rows came back seeded at
0/0 and the node baseline jumped to the node's totals. Quota accounting
lost every byte the clients had used.

The sweep now moves rows for emails the snapshot still reports to the
inbound now carrying them, instead of deleting them. Only emails the node
no longer reports at all lose their row, mirroring the existing baseline
rule.

The alias check also treated a master-created node inbound that had not
been tag-matched yet (empty origin_node_guid) as foreign, so a renamed tag
on its first sync replaced it instead of aliasing it. An empty origin now
counts as this node.

PostgreSQL lane not run locally; the change is a plain UPDATE through GORM.

Closes #6749
This commit is contained in:
MHSanaei
2026-10-06 12:35:59 +02:00
parent d57dcf824b
commit 6be3c438e1
2 changed files with 134 additions and 2 deletions
+44 -2
View File
@@ -698,7 +698,12 @@ func (s *InboundService) setRemoteTrafficLocked(nodeID int, snap *runtime.Traffi
var compatible []*model.Inbound
for i := range central {
candidate := &central[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 {
@@ -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)
}
}