Compare commits

..

3 Commits

Author SHA1 Message Date
MHSanaei 2ffc694a6d fix(nodes): let a master's push and reconcile converge node inbounds
Invariant: every inbound the master holds for a node, with its clients and
enable flag, converges onto that node at the next push or reconcile.

823db059 made an inbound save keep the stored clients and enable unless the
request is a master push, recognised only by a node-sync token scope. Nodes
are usually enrolled with an admin token (the -getApiToken and install
default), so their nodes treated every master push as a stale form save:
clients attached on the master (Attach existing clients, bulk attach above
the per-client threshold, any dirty-node reconcile) never reached the node,
yet the push was recorded as successful so reconcile stopped retrying.
Remote now marks every request with X-3x-Master-Push, and the node honours
it like the node-sync scope. The browser form never sends it, so the
stale-form protection for panel saves and plain API scripts is unchanged.

Separately, an inbound deleted on the node kept its tag->id in the
master's cache, so reconcile sent update/<deleted id> forever instead of
re-creating it. When the node reports no form of the tag, ReconcileInbound
now drops the cached id so the resolve re-reads the node and falls back to
add. Two runtime tests pinned the old update-to-cached-id path for an
absent inbound; they now count the add, or report the inbound present.
2026-10-05 12:17:59 +02:00
MHSanaei 815c9c5772 fix(tuic): accept the server's STOP_SENDING when tests close uni streams
The race job failed in TestAudit3ManagerEnsureActualSendersWithPersistentTraffic
with "close called for canceled stream 14". The server parses one command
per uni stream and then calls CancelRead, as the quinn reference server does
on drop, so its STOP_SENDING can reach the client before the client's own
Close and quic-go reports that Close as an error. The data was already read.

Every test that wrote a command on a uni stream and required Close to
succeed shared this race. closeUniStream accepts only a remote StreamError
on the stream's context, so any other Close failure still fails the test.
2026-10-03 15:48:05 +02:00
MHSanaei 05eb06f333 fix(tuic): wait for both traffic counters in the relay E2E tests
The race job failed on TestServerUDPDatagramE2E with Up:0 Down:1300.
BytesUp is added on the sending goroutine after the relay Send returns,
while BytesDown is added on the response goroutine, so the mock echo can
be counted and delivered before the upload is. The test drained the
counters once right after the reply and assumed both were present.

Production is unaffected: deltas left for the next collection window are
still summed. The TCP E2E test made the same assumption, so both now
accumulate drained deltas until up and down reach the payload size.
2026-10-03 13:43:06 +02:00
9 changed files with 186 additions and 33 deletions
+2 -6
View File
@@ -96,9 +96,7 @@ func TestAudit3ManagerEnsureActualSendersWithPersistentTraffic(t *testing.T) {
if _, err := stream.Write(frame.Bytes()); err != nil { if _, err := stream.Write(frame.Bytes()); err != nil {
t.Fatal(err) t.Fatal(err)
} }
if err := stream.Close(); err != nil { closeUniStream(t, stream)
t.Fatal(err)
}
response, err := p.client.AcceptUniStream(ctx) response, err := p.client.AcceptUniStream(ctx)
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
@@ -164,9 +162,7 @@ func TestAudit3ManagerEnsureActualSendersWithPersistentTraffic(t *testing.T) {
if _, err := auth.Write(authBytes); err != nil { if _, err := auth.Write(authBytes); err != nil {
t.Fatal(err) t.Fatal(err)
} }
if err := auth.Close(); err != nil { closeUniStream(t, auth)
t.Fatal(err)
}
waitForClientCongestionSender(t, server, client, served) waitForClientCongestionSender(t, server, client, served)
var serverConn *quic.Conn var serverConn *quic.Conn
+1 -3
View File
@@ -57,9 +57,7 @@ func audit3LogsStart(t *testing.T, level, marker, relayAddr string) (*Server, *c
if _, err := auth.Write(frame); err != nil { if _, err := auth.Write(frame); err != nil {
t.Fatal(err) t.Fatal(err)
} }
if err := auth.Close(); err != nil { closeUniStream(t, auth)
t.Fatal(err)
}
_, _ = authenticatedServerConnection(t, s, id) _, _ = authenticatedServerConnection(t, s, id)
return s, c, id, password, token return s, c, id, password, token
} }
+49 -18
View File
@@ -10,12 +10,14 @@ import (
"crypto/x509" "crypto/x509"
"crypto/x509/pkix" "crypto/x509/pkix"
"encoding/pem" "encoding/pem"
"errors"
"io" "io"
"math/big" "math/big"
"net" "net"
"testing" "testing"
"time" "time"
serverquic "github.com/apernet/quic-go"
"github.com/google/uuid" "github.com/google/uuid"
"github.com/quic-go/quic-go" "github.com/quic-go/quic-go"
) )
@@ -196,12 +198,51 @@ func testServerTCPConnectE2E(t *testing.T, controller string) {
t.Fatalf("expected active email alice@example.com, got %v", activeEmails) t.Fatalf("expected active email alice@example.com, got %v", activeEmails)
} }
deltas := server.CollectClientTraffic() waitForClientTraffic(t, server, "alice@example.com", int64(len(testMsg)))
if len(deltas) == 0 {
t.Fatalf("expected traffic deltas, got none")
} }
if deltas[0].Email != "alice@example.com" || deltas[0].Up < int64(len(testMsg)) || deltas[0].Down < int64(len(testMsg)) {
t.Fatalf("unexpected traffic deltas: %+v", deltas[0]) // closeUniStream tolerates only the server's STOP_SENDING: it cancels the read side of a
// uni stream once the command is parsed, which can land before the client's FIN.
func closeUniStream(t *testing.T, stream interface {
Close() error
Context() context.Context
},
) {
t.Helper()
err := stream.Close()
if err == nil {
return
}
cause := context.Cause(stream.Context())
var clientErr *quic.StreamError
var serverErr *serverquic.StreamError
if (errors.As(cause, &clientErr) && clientErr.Remote) || (errors.As(cause, &serverErr) && serverErr.Remote) {
return
}
t.Fatalf("close uni stream: %v (cause %v)", err, cause)
}
// waitForClientTraffic accumulates drained deltas because the up and down counters are
// bumped on different relay goroutines, so the echo can arrive before the upload is counted.
func waitForClientTraffic(t *testing.T, server *Server, email string, minBytes int64) {
t.Helper()
var up, down int64
deadline := time.Now().Add(4 * time.Second)
for {
for _, delta := range server.CollectClientTraffic() {
if delta.Email != email {
t.Fatalf("unexpected traffic delta for %q: %+v", delta.Email, delta)
}
up += delta.Up
down += delta.Down
}
if up >= minBytes && down >= minBytes {
return
}
if time.Now().After(deadline) {
t.Fatalf("traffic for %s = up %d, down %d; want both >= %d", email, up, down, minBytes)
}
time.Sleep(5 * time.Millisecond)
} }
} }
@@ -356,13 +397,7 @@ func testServerUDPDatagramE2E(t *testing.T, controller string) {
} }
// 4. Verify traffic // 4. Verify traffic
deltas := server.CollectClientTraffic() waitForClientTraffic(t, server, "bob@example.com", int64(len(udpMsg)))
if len(deltas) == 0 {
t.Fatalf("expected traffic deltas, got none")
}
if deltas[0].Email != "bob@example.com" || deltas[0].Up < int64(len(udpMsg)) || deltas[0].Down < int64(len(udpMsg)) {
t.Fatalf("unexpected traffic deltas: %+v", deltas[0])
}
} }
func TestServerUDPStreamE2E(t *testing.T) { func TestServerUDPStreamE2E(t *testing.T) {
@@ -433,9 +468,7 @@ func testServerUDPStreamE2E(t *testing.T, controller string) {
if _, err := authStream.Write(authPayload); err != nil { if _, err := authStream.Write(authPayload); err != nil {
t.Fatalf("write authentication payload failed: %v", err) t.Fatalf("write authentication payload failed: %v", err)
} }
if err := authStream.Close(); err != nil { closeUniStream(t, authStream)
t.Fatalf("close authentication stream failed: %v", err)
}
target := &Address{Type: AddrTypeIPv4, IP: net.ParseIP("8.8.8.8"), Port: 53} target := &Address{Type: AddrTypeIPv4, IP: net.ParseIP("8.8.8.8"), Port: 53}
udpMsg := bytes.Repeat([]byte("s"), 8500) udpMsg := bytes.Repeat([]byte("s"), 8500)
@@ -458,9 +491,7 @@ func testServerUDPStreamE2E(t *testing.T, controller string) {
if _, err := packetStream.Write(frame.Bytes()); err != nil { if _, err := packetStream.Write(frame.Bytes()); err != nil {
t.Fatalf("write packet frame failed: %v", err) t.Fatalf("write packet frame failed: %v", err)
} }
if err := packetStream.Close(); err != nil { closeUniStream(t, packetStream)
t.Fatalf("close packet stream failed: %v", err)
}
} }
replyReassembler := newPacketReassembler(maxUdpRelayPacketSize) replyReassembler := newPacketReassembler(maxUdpRelayPacketSize)
+1 -3
View File
@@ -57,9 +57,7 @@ func startLifecycleTestServer(t *testing.T, relayAddr, email string) (*Server, *
if _, err := stream.Write(auth); err != nil { if _, err := stream.Write(auth); err != nil {
t.Fatalf("write authentication: %v", err) t.Fatalf("write authentication: %v", err)
} }
if err := stream.Close(); err != nil { closeUniStream(t, stream)
t.Fatalf("close authentication stream: %v", err)
}
return server, client, clientID, password return server, client, clientID, password
} }
+2
View File
@@ -16,6 +16,8 @@ const (
HashHeader = "X-Config-Sha256" HashHeader = "X-Config-Sha256"
// CapsHeader is set by a node on its API responses to advertise support. // CapsHeader is set by a node on its API responses to advertise support.
CapsHeader = "X-3x-Node-Caps" CapsHeader = "X-3x-Node-Caps"
// MasterPushHeader marks a request as a master's push, whatever its token scope.
MasterPushHeader = "X-3x-Master-Push"
// EncodingZstd is the Content-Encoding value for a zstd-compressed body. // EncodingZstd is the Content-Encoding value for a zstd-compressed body.
EncodingZstd = "zstd" EncodingZstd = "zstd"
// CapZstd is the capability token advertised in CapsHeader. // CapZstd is the capability token advertised in CapsHeader.
+4 -1
View File
@@ -7,6 +7,7 @@ import (
"strings" "strings"
"github.com/mhsanaei/3x-ui/v3/internal/database/model" "github.com/mhsanaei/3x-ui/v3/internal/database/model"
"github.com/mhsanaei/3x-ui/v3/internal/util/wirecodec"
"github.com/mhsanaei/3x-ui/v3/internal/web/middleware" "github.com/mhsanaei/3x-ui/v3/internal/web/middleware"
"github.com/mhsanaei/3x-ui/v3/internal/web/service" "github.com/mhsanaei/3x-ui/v3/internal/web/service"
"github.com/mhsanaei/3x-ui/v3/internal/web/session" "github.com/mhsanaei/3x-ui/v3/internal/web/session"
@@ -64,7 +65,9 @@ func (a *InboundController) broadcastInboundsUpdate(userId int) {
func (a *InboundController) inboundServiceFor(c *gin.Context) *service.InboundService { func (a *InboundController) inboundServiceFor(c *gin.Context) *service.InboundService {
svc := a.inboundService svc := a.inboundService
scope, _ := c.Get("api_token_scope") scope, _ := c.Get("api_token_scope")
svc.FromNodeSync = scope == model.ApiScopeNodeSync // A master enrolled with an admin token (the -getApiToken default) has no
// node-sync scope, so it marks every request it sends instead.
svc.FromNodeSync = scope == model.ApiScopeNodeSync || c.GetHeader(wirecodec.MasterPushHeader) != ""
return &svc return &svc
} }
@@ -0,0 +1,79 @@
package controller
import (
"context"
"net/http/httptest"
"net/url"
"path/filepath"
"strconv"
"strings"
"testing"
"github.com/gin-gonic/gin"
"github.com/mhsanaei/3x-ui/v3/internal/database"
"github.com/mhsanaei/3x-ui/v3/internal/database/dbtest"
"github.com/mhsanaei/3x-ui/v3/internal/database/model"
"github.com/mhsanaei/3x-ui/v3/internal/util/crypto"
"github.com/mhsanaei/3x-ui/v3/internal/web/runtime"
)
// A node enrolled with an admin-scope token (the -getApiToken default) must
// still store the clients its master pushes; it used to keep its own list.
func TestMasterPushWithAdminTokenAppliesClients(t *testing.T) {
gin.SetMode(gin.TestMode)
dbDir := t.TempDir()
t.Setenv("XUI_DB_FOLDER", dbDir)
dbtest.InitDB(t, filepath.Join(dbDir, "x-ui.db"))
prev := runtime.GetManager()
runtime.SetManager(runtime.NewManager(runtime.LocalDeps{APIPort: func() int { return 0 }, SetNeedRestart: func() {}}))
t.Cleanup(func() { runtime.SetManager(prev) })
const token = "admin-node-token"
if err := database.GetDB().Create(&model.ApiToken{
Name: "node", Token: crypto.HashTokenSHA256(token), Enabled: true, Scope: model.ApiScopeAdmin,
}).Error; err != nil {
t.Fatalf("seed token: %v", err)
}
var owner model.User
if err := database.GetDB().First(&owner).Error; err != nil {
t.Fatalf("load panel user: %v", err)
}
const stream = `{"network":"tcp","security":"none","tcpSettings":{"header":{"type":"none"}}}`
stored := &model.Inbound{
UserId: owner.Id, Tag: "in-46001", Protocol: model.VLESS, Port: 46001, Enable: true,
Settings: `{"clients":[],"decryption":"none"}`, StreamSettings: stream, Sniffing: `{}`,
}
if err := database.GetDB().Create(stored).Error; err != nil {
t.Fatalf("seed node inbound: %v", err)
}
engine := gin.New()
a := &APIController{}
api := engine.Group("/panel/api")
api.Use(a.checkAPIAuth, a.enforceTokenScope)
NewInboundController(api.Group("/inbounds"))
srv := httptest.NewServer(engine)
defer srv.Close()
u, _ := url.Parse(srv.URL)
port, _ := strconv.Atoi(u.Port())
master := runtime.NewRemote(&model.Node{
Id: 1, Name: "n1", Scheme: "http", Address: u.Hostname(), Port: port,
BasePath: "/", ApiToken: token, Enable: true, AllowPrivateAddress: true,
}, nil)
pushed := *stored
pushed.Settings = `{"clients":[{"id":"7fa0b7d1-9b5f-47ad-bef2-6cb0c4a624be","email":"alice","enable":true,"subId":"s-alice"}],"decryption":"none"}`
if err := master.UpdateInbound(context.Background(), &pushed, &pushed); err != nil {
t.Fatalf("master push: %v", err)
}
var got model.Inbound
if err := database.GetDB().First(&got, stored.Id).Error; err != nil {
t.Fatalf("reload node inbound: %v", err)
}
if !strings.Contains(got.Settings, `"alice"`) {
t.Fatalf("node kept its own client list after a master push: %s", got.Settings)
}
}
+36 -2
View File
@@ -17,7 +17,8 @@ import (
func TestReconcileInbound_SkipsUnchanged(t *testing.T) { func TestReconcileInbound_SkipsUnchanged(t *testing.T) {
var pushes atomic.Int32 var pushes atomic.Int32
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.Method == http.MethodPost && strings.Contains(r.URL.Path, "/panel/api/inbounds/update/") { if r.Method == http.MethodPost && (strings.Contains(r.URL.Path, "/panel/api/inbounds/update/") ||
strings.Contains(r.URL.Path, "/panel/api/inbounds/add")) {
pushes.Add(1) pushes.Add(1)
} }
w.Header().Set("Content-Type", "application/json") w.Header().Set("Content-Type", "application/json")
@@ -283,7 +284,7 @@ func TestDelInboundDropsReconcileFingerprint(t *testing.T) {
ib := &model.Inbound{Tag: "in-del", Protocol: model.VLESS, Port: 443, Settings: `{"clients":[]}`} ib := &model.Inbound{Tag: "in-del", Protocol: model.VLESS, Port: 443, Settings: `{"clients":[]}`}
r.cacheSet(ib.Tag, 7) r.cacheSet(ib.Tag, 7)
if pushed, err := r.ReconcileInbound(context.Background(), ib, false); err != nil || !pushed { if pushed, err := r.ReconcileInbound(context.Background(), ib, true); err != nil || !pushed {
t.Fatalf("initial reconcile: pushed=%v err=%v, want push", pushed, err) t.Fatalf("initial reconcile: pushed=%v err=%v, want push", pushed, err)
} }
if err := r.DelInbound(context.Background(), ib); err != nil { if err := r.DelInbound(context.Background(), ib); err != nil {
@@ -317,3 +318,36 @@ func TestUpdateInboundFallbackAddSeedsReconcileFingerprint(t *testing.T) {
t.Fatalf("reconcile sent %d full inbound updates, want 0", got) t.Fatalf("reconcile sent %d full inbound updates, want 0", got)
} }
} }
// An inbound deleted on the node must be re-created by the next reconcile; a
// cached tag→id from before the delete used to send update/<gone id> forever.
func TestReconcileInbound_RecreatesInboundTheNodeLost(t *testing.T) {
var adds, staleUpdates atomic.Int32
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
switch {
case strings.Contains(r.URL.Path, "/panel/api/inbounds/list"):
_, _ = w.Write([]byte(`{"success":true,"obj":[]}`))
case strings.Contains(r.URL.Path, "/panel/api/inbounds/update/"):
staleUpdates.Add(1)
_, _ = w.Write([]byte(`{"success":false,"msg":"record not found"}`))
case strings.Contains(r.URL.Path, "/panel/api/inbounds/add"):
adds.Add(1)
_, _ = w.Write([]byte(`{"success":true,"obj":{"id":9,"tag":"in-1"}}`))
default:
_, _ = w.Write([]byte(`{"success":true}`))
}
}))
defer srv.Close()
r := NewRemote(nodeForPlainServer(t, srv, "verify", "tok"), nil)
ib := &model.Inbound{Tag: "n1-in-1", Protocol: model.VLESS, Port: 443, Settings: `{"clients":[]}`}
r.cacheSet("in-1", 7)
if pushed, err := r.ReconcileInbound(context.Background(), ib, false); err != nil || !pushed {
t.Fatalf("reconcile of a lost inbound: pushed=%v err=%v, want a re-create", pushed, err)
}
if staleUpdates.Load() != 0 || adds.Load() != 1 {
t.Fatalf("updates to the stale id=%d adds=%d, want 0 and 1", staleUpdates.Load(), adds.Load())
}
}
+12
View File
@@ -237,6 +237,7 @@ func (r *Remote) do(ctx context.Context, method, path string, body any) (*envelo
req.Header.Set("Authorization", "Bearer "+token) req.Header.Set("Authorization", "Bearer "+token)
} }
req.Header.Set("Accept", "application/json") req.Header.Set("Accept", "application/json")
req.Header.Set(wirecodec.MasterPushHeader, "1")
if contentType != "" { if contentType != "" {
req.Header.Set("Content-Type", contentType) req.Header.Set("Content-Type", contentType)
} }
@@ -359,6 +360,15 @@ func (r *Remote) cacheDel(tag string) {
delete(r.pushedFP, tag) delete(r.pushedFP, tag)
} }
// forgetTag drops every tag form cacheGetTag would match, once the node reports
// none of them, so the next resolve re-reads the node instead of a deleted id.
func (r *Remote) forgetTag(tag string) {
prefix := nodeInboundTagPrefix(r.node.Id)
bare := strings.TrimPrefix(tag, prefix)
r.cacheDel(bare)
r.cacheDel(prefix + bare)
}
func (r *Remote) ListRemoteTags(ctx context.Context) ([]string, error) { func (r *Remote) ListRemoteTags(ctx context.Context) ([]string, error) {
if err := r.refreshRemoteIDs(ctx); err != nil { if err := r.refreshRemoteIDs(ctx); err != nil {
return nil, err return nil, err
@@ -494,6 +504,8 @@ func (r *Remote) ReconcileInbound(ctx context.Context, ib *model.Inbound, exists
if ok && prev == fp { if ok && prev == fp {
return false, nil return false, nil
} }
} else {
r.forgetTag(ib.Tag)
} }
if err := r.UpdateInbound(ctx, ib, ib); err != nil { if err := r.UpdateInbound(ctx, ib, ib); err != nil {
return false, err return false, err