reconcile_skip_test.go 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353
  1. package runtime
  2. import (
  3. "context"
  4. "net/http"
  5. "net/http/httptest"
  6. "strings"
  7. "sync/atomic"
  8. "testing"
  9. "github.com/mhsanaei/3x-ui/v3/internal/database/model"
  10. )
  11. // TestReconcileInbound_SkipsUnchanged proves the delta-skip: a second reconcile
  12. // of an unchanged inbound that the node still reports sends no push, while a
  13. // content change or an absent-on-node inbound forces a fresh push.
  14. func TestReconcileInbound_SkipsUnchanged(t *testing.T) {
  15. var pushes atomic.Int32
  16. srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
  17. if r.Method == http.MethodPost && (strings.Contains(r.URL.Path, "/panel/api/inbounds/update/") ||
  18. strings.Contains(r.URL.Path, "/panel/api/inbounds/add")) {
  19. pushes.Add(1)
  20. }
  21. w.Header().Set("Content-Type", "application/json")
  22. _, _ = w.Write([]byte(`{"success":true}`))
  23. }))
  24. defer srv.Close()
  25. r := NewRemote(nodeForPlainServer(t, srv, "verify", "tok"), nil)
  26. ib := &model.Inbound{Tag: "in-1", Protocol: model.VLESS, Port: 443, Settings: `{"clients":[]}`}
  27. // Pre-seed the tag→id cache so resolveRemoteID needs no network round-trip.
  28. r.cacheSet(ib.Tag, 7)
  29. // First reconcile: node doesn't report it yet → must push and record the fp.
  30. if pushed, err := r.ReconcileInbound(context.Background(), ib, false); err != nil || !pushed {
  31. t.Fatalf("first reconcile: pushed=%v err=%v, want push", pushed, err)
  32. }
  33. if got := pushes.Load(); got != 1 {
  34. t.Fatalf("after first reconcile pushes=%d, want 1", got)
  35. }
  36. // Second reconcile: unchanged and present on node → skip.
  37. if pushed, err := r.ReconcileInbound(context.Background(), ib, true); err != nil || pushed {
  38. t.Fatalf("second reconcile: pushed=%v err=%v, want skip", pushed, err)
  39. }
  40. if got := pushes.Load(); got != 1 {
  41. t.Fatalf("unchanged reconcile pushed again: pushes=%d, want 1", got)
  42. }
  43. // Content change → push again even though it's present on node.
  44. ib.Settings = `{"clients":[{"email":"a@x"}]}`
  45. if pushed, err := r.ReconcileInbound(context.Background(), ib, true); err != nil || !pushed {
  46. t.Fatalf("changed reconcile: pushed=%v err=%v, want push", pushed, err)
  47. }
  48. if got := pushes.Load(); got != 2 {
  49. t.Fatalf("changed reconcile pushes=%d, want 2", got)
  50. }
  51. // Absent on node (e.g. node restarted/lost it) → re-push even if fp matches.
  52. if pushed, err := r.ReconcileInbound(context.Background(), ib, false); err != nil || !pushed {
  53. t.Fatalf("absent-on-node reconcile: pushed=%v err=%v, want push", pushed, err)
  54. }
  55. if got := pushes.Load(); got != 3 {
  56. t.Fatalf("absent-on-node reconcile pushes=%d, want 3", got)
  57. }
  58. }
  59. type nodeCallCounts struct {
  60. adds atomic.Int32
  61. inboundUpdates atomic.Int32
  62. clientMutations atomic.Int32
  63. }
  64. func newCountingNodeServer(t *testing.T, clientsResp string) (*httptest.Server, *nodeCallCounts) {
  65. t.Helper()
  66. counts := &nodeCallCounts{}
  67. srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
  68. w.Header().Set("Content-Type", "application/json")
  69. switch {
  70. case strings.Contains(r.URL.Path, "/panel/api/inbounds/add"):
  71. counts.adds.Add(1)
  72. case strings.Contains(r.URL.Path, "/panel/api/inbounds/update/"):
  73. counts.inboundUpdates.Add(1)
  74. case strings.Contains(r.URL.Path, "/panel/api/clients/"):
  75. counts.clientMutations.Add(1)
  76. _, _ = w.Write([]byte(clientsResp))
  77. return
  78. }
  79. _, _ = w.Write([]byte(`{"success":true}`))
  80. }))
  81. t.Cleanup(srv.Close)
  82. return srv, counts
  83. }
  84. func perClientMutationCases(client model.Client) []struct {
  85. name string
  86. run func(*Remote, *model.Inbound) error
  87. } {
  88. return []struct {
  89. name string
  90. run func(*Remote, *model.Inbound) error
  91. }{
  92. {"add", func(r *Remote, ib *model.Inbound) error {
  93. return r.AddClient(context.Background(), ib, client)
  94. }},
  95. {"delete", func(r *Remote, ib *model.Inbound) error {
  96. return r.DeleteUser(context.Background(), ib, client.Email)
  97. }},
  98. {"update", func(r *Remote, ib *model.Inbound) error {
  99. return r.UpdateUser(context.Background(), ib, client.Email, client)
  100. }},
  101. }
  102. }
  103. // TestPerClientMutationsAloneDoNotSeedReconcileFingerprint: a per-client RPC
  104. // only proves one client's slice converged, so without an explicit advance the
  105. // dirty-reconcile backup must still send the full inbound.
  106. func TestPerClientMutationsAloneDoNotSeedReconcileFingerprint(t *testing.T) {
  107. srv, counts := newCountingNodeServer(t, `{"success":true}`)
  108. client := model.Client{ID: "11111111-1111-1111-1111-111111111111", Email: "a@x", SubID: "s", Enable: true}
  109. for _, tt := range perClientMutationCases(client) {
  110. t.Run(tt.name, func(t *testing.T) {
  111. counts.inboundUpdates.Store(0)
  112. counts.clientMutations.Store(0)
  113. r := NewRemote(nodeForPlainServer(t, srv, "verify", "tok"), nil)
  114. ib := &model.Inbound{Tag: "in-" + tt.name, Protocol: model.VLESS, Port: 443, Settings: `{"clients":[{"id":"11111111-1111-1111-1111-111111111111","email":"a@x","subId":"s","enable":true}]}`}
  115. r.cacheSet(ib.Tag, 7)
  116. if err := tt.run(r, ib); err != nil {
  117. t.Fatalf("%s client mutation: %v", tt.name, err)
  118. }
  119. if got := counts.clientMutations.Load(); got != 1 {
  120. t.Fatalf("%s client mutation requests=%d, want 1", tt.name, got)
  121. }
  122. if pushed, err := r.ReconcileInbound(context.Background(), ib, true); err != nil || !pushed {
  123. t.Fatalf("%s reconcile after unadvanced mutation: pushed=%v err=%v, want full push", tt.name, pushed, err)
  124. }
  125. if got := counts.inboundUpdates.Load(); got != 1 {
  126. t.Fatalf("%s reconcile sent %d full inbound updates, want 1", tt.name, got)
  127. }
  128. })
  129. }
  130. }
  131. // TestAdvancePushedInboundEnablesReconcileSkip: when the node provably held the
  132. // pre-edit payload and every per-client push succeeded, advancing the
  133. // fingerprint lets the next reconcile skip the redundant full push.
  134. func TestAdvancePushedInboundEnablesReconcileSkip(t *testing.T) {
  135. srv, counts := newCountingNodeServer(t, `{"success":true}`)
  136. client := model.Client{ID: "11111111-1111-1111-1111-111111111111", Email: "a@x", SubID: "s", Enable: true}
  137. prevByOp := map[string]string{
  138. "add": `{"clients":[]}`,
  139. "delete": `{"clients":[{"email":"a@x","enable":true}]}`,
  140. "update": `{"clients":[{"email":"a@x","enable":false}]}`,
  141. }
  142. newByOp := map[string]string{
  143. "add": `{"clients":[{"email":"a@x","enable":true}]}`,
  144. "delete": `{"clients":[]}`,
  145. "update": `{"clients":[{"email":"a@x","enable":true}]}`,
  146. }
  147. for _, tt := range perClientMutationCases(client) {
  148. t.Run(tt.name, func(t *testing.T) {
  149. counts.inboundUpdates.Store(0)
  150. r := NewRemote(nodeForPlainServer(t, srv, "verify", "tok"), nil)
  151. prevIb := &model.Inbound{Tag: "in-adv-" + tt.name, Protocol: model.VLESS, Port: 443, Settings: prevByOp[tt.name]}
  152. ib := &model.Inbound{Tag: prevIb.Tag, Protocol: model.VLESS, Port: 443, Settings: newByOp[tt.name]}
  153. r.cacheSet(ib.Tag, 7)
  154. r.recordPushedInbound(prevIb)
  155. if err := tt.run(r, ib); err != nil {
  156. t.Fatalf("%s client mutation: %v", tt.name, err)
  157. }
  158. r.AdvancePushedInbound(prevIb, ib)
  159. if pushed, err := r.ReconcileInbound(context.Background(), ib, true); err != nil || pushed {
  160. t.Fatalf("%s reconcile after advance: pushed=%v err=%v, want skip", tt.name, pushed, err)
  161. }
  162. if got := counts.inboundUpdates.Load(); got != 0 {
  163. t.Fatalf("%s reconcile sent %d full inbound updates, want 0", tt.name, got)
  164. }
  165. })
  166. }
  167. }
  168. // TestAdvancePushedInboundRequiresMatchingPreviousFingerprint: if changes were
  169. // folded to dirty (or an earlier push failed), the recorded fingerprint no
  170. // longer matches the pre-edit payload; a later successful client push must not
  171. // mask the pending reconcile.
  172. func TestAdvancePushedInboundRequiresMatchingPreviousFingerprint(t *testing.T) {
  173. srv, counts := newCountingNodeServer(t, `{"success":true}`)
  174. r := NewRemote(nodeForPlainServer(t, srv, "verify", "tok"), nil)
  175. staleIb := &model.Inbound{Tag: "in-stale", Protocol: model.VLESS, Port: 443, Settings: `{"clients":[{"email":"folded@x"}]}`}
  176. prevIb := &model.Inbound{Tag: "in-stale", Protocol: model.VLESS, Port: 443, Settings: `{"clients":[]}`}
  177. ib := &model.Inbound{Tag: "in-stale", Protocol: model.VLESS, Port: 443, Settings: `{"clients":[{"email":"a@x"}]}`}
  178. r.cacheSet(ib.Tag, 7)
  179. r.recordPushedInbound(staleIb)
  180. if err := r.UpdateUser(context.Background(), ib, "a@x", model.Client{Email: "a@x"}); err != nil {
  181. t.Fatalf("client mutation: %v", err)
  182. }
  183. r.AdvancePushedInbound(prevIb, ib)
  184. if pushed, err := r.ReconcileInbound(context.Background(), ib, true); err != nil || !pushed {
  185. t.Fatalf("reconcile with unproven pre-state: pushed=%v err=%v, want full push", pushed, err)
  186. }
  187. if got := counts.inboundUpdates.Load(); got != 1 {
  188. t.Fatalf("reconcile sent %d full inbound updates, want 1", got)
  189. }
  190. }
  191. // TestAdoptedSerializationChainKeepsReconcileSkip: after a push the node
  192. // re-serializes settings its own way and the master adopts that form back into
  193. // its DB; stamping the adopted payload keeps edit->advance->skip alive instead
  194. // of degrading every edit to a full reconcile push.
  195. func TestAdoptedSerializationChainKeepsReconcileSkip(t *testing.T) {
  196. srv, counts := newCountingNodeServer(t, `{"success":true}`)
  197. r := NewRemote(nodeForPlainServer(t, srv, "verify", "tok"), nil)
  198. pushForm := &model.Inbound{Tag: "in-adopt", Protocol: model.VLESS, Port: 443, Settings: `{"clients":[{"email":"a@x"}]}`}
  199. adoptedForm := &model.Inbound{Tag: "in-adopt", Protocol: model.VLESS, Port: 443, Settings: "{\n \"clients\": [{\"email\": \"a@x\"}]\n}"}
  200. edited := &model.Inbound{Tag: "in-adopt", Protocol: model.VLESS, Port: 443, Settings: "{\n \"clients\": [{\"comment\": \"c\", \"email\": \"a@x\"}]\n}"}
  201. r.cacheSet(pushForm.Tag, 7)
  202. r.recordPushedInbound(pushForm)
  203. r.RecordAdoptedInbound(adoptedForm)
  204. if err := r.UpdateUser(context.Background(), edited, "a@x", model.Client{Email: "a@x"}); err != nil {
  205. t.Fatalf("client mutation: %v", err)
  206. }
  207. r.AdvancePushedInbound(adoptedForm, edited)
  208. if pushed, err := r.ReconcileInbound(context.Background(), edited, true); err != nil || pushed {
  209. t.Fatalf("reconcile after adopted-form advance: pushed=%v err=%v, want skip", pushed, err)
  210. }
  211. if got := counts.inboundUpdates.Load(); got != 0 {
  212. t.Fatalf("full inbound updates=%d, want 0", got)
  213. }
  214. }
  215. func TestDeleteUserNotFoundHandling(t *testing.T) {
  216. t.Run("envelope not-found counts as already deleted", func(t *testing.T) {
  217. srv, counts := newCountingNodeServer(t, `{"success":false,"msg":"client not found"}`)
  218. r := NewRemote(nodeForPlainServer(t, srv, "verify", "tok"), nil)
  219. ib := &model.Inbound{Tag: "in-delete-missing", Protocol: model.VLESS, Port: 443, Settings: `{"clients":[]}`}
  220. r.cacheSet(ib.Tag, 7)
  221. if err := r.DeleteUser(context.Background(), ib, "missing@x"); err != nil {
  222. t.Fatalf("DeleteUser missing client: %v", err)
  223. }
  224. if pushed, err := r.ReconcileInbound(context.Background(), ib, true); err != nil || !pushed {
  225. t.Fatalf("reconcile after missing-client delete: pushed=%v err=%v, want full push", pushed, err)
  226. }
  227. if got := counts.inboundUpdates.Load(); got != 1 {
  228. t.Fatalf("reconcile sent %d full inbound updates, want 1", got)
  229. }
  230. })
  231. t.Run("http 404 from an old node stays an error", func(t *testing.T) {
  232. srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
  233. if strings.Contains(r.URL.Path, "/panel/api/clients/") {
  234. http.Error(w, "404 page not found", http.StatusNotFound)
  235. return
  236. }
  237. w.Header().Set("Content-Type", "application/json")
  238. _, _ = w.Write([]byte(`{"success":true}`))
  239. }))
  240. defer srv.Close()
  241. r := NewRemote(nodeForPlainServer(t, srv, "verify", "tok"), nil)
  242. ib := &model.Inbound{Tag: "in-delete-404", Protocol: model.VLESS, Port: 443, Settings: `{"clients":[]}`}
  243. r.cacheSet(ib.Tag, 7)
  244. err := r.DeleteUser(context.Background(), ib, "missing@x")
  245. if err == nil || !strings.Contains(err.Error(), "HTTP 404") {
  246. t.Fatalf("DeleteUser against old node = %v, want HTTP 404 error", err)
  247. }
  248. })
  249. }
  250. // TestDelInboundDropsReconcileFingerprint: deleting an inbound must forget its
  251. // fingerprint so a later same-tag inbound with identical content is re-pushed.
  252. func TestDelInboundDropsReconcileFingerprint(t *testing.T) {
  253. srv, counts := newCountingNodeServer(t, `{"success":true}`)
  254. r := NewRemote(nodeForPlainServer(t, srv, "verify", "tok"), nil)
  255. ib := &model.Inbound{Tag: "in-del", Protocol: model.VLESS, Port: 443, Settings: `{"clients":[]}`}
  256. r.cacheSet(ib.Tag, 7)
  257. if pushed, err := r.ReconcileInbound(context.Background(), ib, true); err != nil || !pushed {
  258. t.Fatalf("initial reconcile: pushed=%v err=%v, want push", pushed, err)
  259. }
  260. if err := r.DelInbound(context.Background(), ib); err != nil {
  261. t.Fatalf("DelInbound: %v", err)
  262. }
  263. r.cacheSet(ib.Tag, 7)
  264. if pushed, err := r.ReconcileInbound(context.Background(), ib, true); err != nil || !pushed {
  265. t.Fatalf("reconcile after DelInbound: pushed=%v err=%v, want full push", pushed, err)
  266. }
  267. if got := counts.inboundUpdates.Load(); got != 2 {
  268. t.Fatalf("full inbound updates=%d, want 2", got)
  269. }
  270. }
  271. func TestUpdateInboundFallbackAddSeedsReconcileFingerprint(t *testing.T) {
  272. srv, counts := newCountingNodeServer(t, `{"success":true}`)
  273. r := NewRemote(nodeForPlainServer(t, srv, "verify", "tok"), nil)
  274. ib := &model.Inbound{Tag: "in-add-fallback", Protocol: model.VLESS, Port: 443, Settings: `{"clients":[]}`}
  275. if err := r.UpdateInbound(context.Background(), ib, ib); err != nil {
  276. t.Fatalf("UpdateInbound fallback add: %v", err)
  277. }
  278. if got := counts.adds.Load(); got != 1 {
  279. t.Fatalf("fallback add requests=%d, want 1", got)
  280. }
  281. if pushed, err := r.ReconcileInbound(context.Background(), ib, true); err != nil || pushed {
  282. t.Fatalf("reconcile after fallback add: pushed=%v err=%v, want skip", pushed, err)
  283. }
  284. if got := counts.inboundUpdates.Load(); got != 0 {
  285. t.Fatalf("reconcile sent %d full inbound updates, want 0", got)
  286. }
  287. }
  288. // An inbound deleted on the node must be re-created by the next reconcile; a
  289. // cached tag→id from before the delete used to send update/<gone id> forever.
  290. func TestReconcileInbound_RecreatesInboundTheNodeLost(t *testing.T) {
  291. var adds, staleUpdates atomic.Int32
  292. srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
  293. w.Header().Set("Content-Type", "application/json")
  294. switch {
  295. case strings.Contains(r.URL.Path, "/panel/api/inbounds/list"):
  296. _, _ = w.Write([]byte(`{"success":true,"obj":[]}`))
  297. case strings.Contains(r.URL.Path, "/panel/api/inbounds/update/"):
  298. staleUpdates.Add(1)
  299. _, _ = w.Write([]byte(`{"success":false,"msg":"record not found"}`))
  300. case strings.Contains(r.URL.Path, "/panel/api/inbounds/add"):
  301. adds.Add(1)
  302. _, _ = w.Write([]byte(`{"success":true,"obj":{"id":9,"tag":"in-1"}}`))
  303. default:
  304. _, _ = w.Write([]byte(`{"success":true}`))
  305. }
  306. }))
  307. defer srv.Close()
  308. r := NewRemote(nodeForPlainServer(t, srv, "verify", "tok"), nil)
  309. ib := &model.Inbound{Tag: "n1-in-1", Protocol: model.VLESS, Port: 443, Settings: `{"clients":[]}`}
  310. r.cacheSet("in-1", 7)
  311. if pushed, err := r.ReconcileInbound(context.Background(), ib, false); err != nil || !pushed {
  312. t.Fatalf("reconcile of a lost inbound: pushed=%v err=%v, want a re-create", pushed, err)
  313. }
  314. if staleUpdates.Load() != 0 || adds.Load() != 1 {
  315. t.Fatalf("updates to the stale id=%d adds=%d, want 0 and 1", staleUpdates.Load(), adds.Load())
  316. }
  317. }