Bladeren bron

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.
MHSanaei 6 uur geleden
bovenliggende
commit
2ffc694a6d

+ 2 - 0
internal/util/wirecodec/wirecodec.go

@@ -16,6 +16,8 @@ const (
 	HashHeader = "X-Config-Sha256"
 	// CapsHeader is set by a node on its API responses to advertise support.
 	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 = "zstd"
 	// CapZstd is the capability token advertised in CapsHeader.

+ 4 - 1
internal/web/controller/inbound.go

@@ -7,6 +7,7 @@ import (
 	"strings"
 
 	"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/service"
 	"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 {
 	svc := a.inboundService
 	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
 }
 

+ 79 - 0
internal/web/controller/inbound_master_push_test.go

@@ -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
internal/web/runtime/reconcile_skip_test.go

@@ -17,7 +17,8 @@ import (
 func TestReconcileInbound_SkipsUnchanged(t *testing.T) {
 	var pushes atomic.Int32
 	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)
 		}
 		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":[]}`}
 	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)
 	}
 	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)
 	}
 }
+
+// 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 - 0
internal/web/runtime/remote.go

@@ -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("Accept", "application/json")
+	req.Header.Set(wirecodec.MasterPushHeader, "1")
 	if contentType != "" {
 		req.Header.Set("Content-Type", contentType)
 	}
@@ -359,6 +360,15 @@ func (r *Remote) cacheDel(tag string) {
 	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) {
 	if err := r.refreshRemoteIDs(ctx); err != nil {
 		return nil, err
@@ -494,6 +504,8 @@ func (r *Remote) ReconcileInbound(ctx context.Context, ib *model.Inbound, exists
 		if ok && prev == fp {
 			return false, nil
 		}
+	} else {
+		r.forgetTag(ib.Tag)
 	}
 	if err := r.UpdateInbound(ctx, ib, ib); err != nil {
 		return false, err