node_sync_test.go 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344
  1. package nodee2e
  2. import (
  3. "encoding/json"
  4. "fmt"
  5. "net/http"
  6. "slices"
  7. "strconv"
  8. "testing"
  9. "time"
  10. )
  11. // TestNodeSync walks one master/node pair per enrollment scope through every
  12. // operation that must converge onto the node. Each subtest names its invariant.
  13. func TestNodeSync(t *testing.T) {
  14. bin := panelBinary(t)
  15. for _, scope := range []string{"admin", "node-sync"} {
  16. t.Run("enrolled with "+scope+" token", func(t *testing.T) {
  17. t.Parallel()
  18. runNodeSyncScenarios(t, bin, scope)
  19. })
  20. }
  21. }
  22. type pair struct {
  23. t *testing.T
  24. master *panel
  25. node *panel
  26. nodeID int
  27. }
  28. func (pr *pair) nodeInbound(port int) (inboundView, bool) { return pr.node.inboundOnPort(port) }
  29. func (pr *pair) masterInbound(port int) (inboundView, bool) {
  30. for _, ib := range pr.master.inbounds() {
  31. if ib.Port == port && ib.NodeID != nil && *ib.NodeID == pr.nodeID {
  32. return ib, true
  33. }
  34. }
  35. return inboundView{}, false
  36. }
  37. func (pr *pair) waitNode(what string, port int, ok func(inboundView) bool) {
  38. pr.t.Helper()
  39. eventually(pr.t, settleTimeout, what, func() (bool, string) {
  40. ib, found := pr.nodeInbound(port)
  41. if !found {
  42. return ok(inboundView{}) && false, fmt.Sprintf("node has no inbound on %d", port)
  43. }
  44. return ok(ib), fmt.Sprintf("node inbound %d: enable=%v remark=%q emails=%v", port, ib.Enable, ib.Remark, ib.emails())
  45. })
  46. }
  47. func (pr *pair) waitNodeAbsent(what string, port int) {
  48. pr.t.Helper()
  49. eventually(pr.t, settleTimeout, what, func() (bool, string) {
  50. ib, found := pr.nodeInbound(port)
  51. return !found, fmt.Sprintf("node still has inbound %d with %v", port, ib.emails())
  52. })
  53. }
  54. func (pr *pair) bulkAttach(emails []string, masterInboundID int) {
  55. pr.t.Helper()
  56. var res struct {
  57. Attached []string `json:"attached"`
  58. Errors []string `json:"errors"`
  59. }
  60. obj := pr.master.call(http.MethodPost, "/panel/api/clients/bulkAttach", map[string]any{"emails": emails, "inboundIds": []int{masterInboundID}})
  61. if err := json.Unmarshal(obj, &res); err != nil {
  62. pr.t.Fatalf("decode bulkAttach: %v", err)
  63. }
  64. if len(res.Errors) != 0 || len(res.Attached) != len(emails) {
  65. pr.t.Fatalf("bulkAttach attached=%d/%d errors=%v", len(res.Attached), len(emails), res.Errors)
  66. }
  67. }
  68. func emailRange(prefix string, from, to int) []string {
  69. out := make([]string, 0, to-from+1)
  70. for i := from; i <= to; i++ {
  71. out = append(out, prefix+strconv.Itoa(i))
  72. }
  73. return out
  74. }
  75. func hasAll(have []string, want ...string) bool {
  76. for _, w := range want {
  77. if !slices.Contains(have, w) {
  78. return false
  79. }
  80. }
  81. return true
  82. }
  83. func runNodeSyncScenarios(t *testing.T, bin, scope string) {
  84. master := newPanel(t, bin, "master")
  85. node := newPanel(t, bin, "node")
  86. linkToken := node.mintToken("master-link", scope)
  87. master.start()
  88. node.start()
  89. pr := &pair{t: t, master: master, node: node}
  90. var (
  91. adoptedPort = freePort(t)
  92. madePort = freePort(t)
  93. lostPort = freePort(t)
  94. droppedPort = freePort(t)
  95. offlinePort = freePort(t)
  96. unmanagedPort = freePort(t)
  97. localPort = freePort(t)
  98. )
  99. node.call(http.MethodPost, "/panel/api/inbounds/add", vlessInbound("pre-existing", adoptedPort, nil))
  100. local := master.call(http.MethodPost, "/panel/api/inbounds/add", vlessInbound("local-pool", localPort, nil))
  101. var localIb inboundView
  102. _ = json.Unmarshal(local, &localIb)
  103. for _, email := range emailRange("p", 1, 45) {
  104. master.call(http.MethodPost, "/panel/api/clients/add", map[string]any{
  105. "client": map[string]any{"email": email, "enable": true}, "inboundIds": []int{localIb.Id},
  106. })
  107. }
  108. var nodeView struct {
  109. Id int `json:"id"`
  110. }
  111. obj := master.call(http.MethodPost, "/panel/api/nodes/add", map[string]any{
  112. "name": "n1", "scheme": "http", "address": "127.0.0.1", "port": node.port, "basePath": "/",
  113. "apiToken": linkToken, "enable": true, "allowPrivateAddress": true,
  114. })
  115. if err := json.Unmarshal(obj, &nodeView); err != nil || nodeView.Id == 0 {
  116. t.Fatalf("decode node add: %v (%s)", err, obj)
  117. }
  118. pr.nodeID = nodeView.Id
  119. var adoptedID, madeID int
  120. t.Run("an inbound already on the node is adopted by the master", func(t *testing.T) {
  121. pr.t = t
  122. eventually(t, settleTimeout, "master adopts the node inbound", func() (bool, string) {
  123. ib, ok := pr.masterInbound(adoptedPort)
  124. adoptedID = ib.Id
  125. return ok, "not adopted yet"
  126. })
  127. })
  128. if adoptedID == 0 {
  129. t.Fatal("no adopted inbound; later scenarios depend on it")
  130. }
  131. t.Run("an inbound created on the master for the node lands there with its clients", func(t *testing.T) {
  132. pr.t = t
  133. obj := master.call(http.MethodPost, "/panel/api/inbounds/add",
  134. vlessInbound("made-on-master", madePort, &pr.nodeID, vlessClient("m1"), vlessClient("m2")))
  135. var ib inboundView
  136. _ = json.Unmarshal(obj, &ib)
  137. madeID = ib.Id
  138. pr.waitNode("node holds the master-made inbound", madePort, func(ib inboundView) bool {
  139. return hasAll(ib.emails(), "m1", "m2")
  140. })
  141. })
  142. t.Run("editing the inbound on the master updates the node and keeps its clients", func(t *testing.T) {
  143. pr.t = t
  144. body := vlessInbound("renamed-on-master", madePort, &pr.nodeID)
  145. body["settings"] = `{"decryption":"none"}`
  146. master.call(http.MethodPost, "/panel/api/inbounds/update/"+strconv.Itoa(madeID), body)
  147. pr.waitNode("node shows the new remark with both clients", madePort, func(ib inboundView) bool {
  148. return ib.Remark == "renamed-on-master" && hasAll(ib.emails(), "m1", "m2")
  149. })
  150. })
  151. t.Run("attaching a few existing clients reaches the node", func(t *testing.T) {
  152. pr.t = t
  153. pr.bulkAttach([]string{"p1", "p2", "p3"}, madeID)
  154. pr.waitNode("node holds p1..p3", madePort, func(ib inboundView) bool {
  155. return hasAll(ib.emails(), "m1", "m2", "p1", "p2", "p3")
  156. })
  157. })
  158. t.Run("attaching more clients than the per-client push limit reaches the node", func(t *testing.T) {
  159. pr.t = t
  160. emails := emailRange("p", 4, 43)
  161. pr.bulkAttach(emails, adoptedID)
  162. pr.waitNode("node holds all 40", adoptedPort, func(ib inboundView) bool {
  163. return hasAll(ib.emails(), emails...)
  164. })
  165. })
  166. t.Run("disabling a client on the master disables it on the node", func(t *testing.T) {
  167. pr.t = t
  168. mib, _ := pr.masterInbound(madePort)
  169. entry := mib.client("p1")
  170. if entry == nil {
  171. t.Fatalf("master inbound has no p1: %v", mib.emails())
  172. }
  173. entry["enable"] = false
  174. master.call(http.MethodPost, "/panel/api/clients/update/p1", entry)
  175. pr.waitNode("node p1 disabled", madePort, func(ib inboundView) bool {
  176. c := ib.client("p1")
  177. return c != nil && c["enable"] == false
  178. })
  179. })
  180. t.Run("detaching a client from the node inbound removes it there", func(t *testing.T) {
  181. pr.t = t
  182. master.call(http.MethodPost, "/panel/api/clients/p2/detach", map[string]any{"inboundIds": []int{madeID}})
  183. pr.waitNode("node drops p2", madePort, func(ib inboundView) bool {
  184. return ib.client("p2") == nil && ib.client("p3") != nil
  185. })
  186. })
  187. t.Run("deleting a client on the master removes it from the node", func(t *testing.T) {
  188. pr.t = t
  189. master.call(http.MethodPost, "/panel/api/clients/del/p4", nil)
  190. pr.waitNode("node drops p4", adoptedPort, func(ib inboundView) bool {
  191. return ib.client("p4") == nil && ib.client("p5") != nil
  192. })
  193. })
  194. t.Run("switching the inbound off on the master switches it off on the node", func(t *testing.T) {
  195. pr.t = t
  196. master.call(http.MethodPost, "/panel/api/inbounds/setEnable/"+strconv.Itoa(madeID), map[string]any{"enable": false})
  197. pr.waitNode("node inbound disabled", madePort, func(ib inboundView) bool { return !ib.Enable })
  198. })
  199. t.Run("node traffic reaches the master and a master reset clears the node", func(t *testing.T) {
  200. pr.t = t
  201. node.call(http.MethodPost, "/panel/api/clients/updateTraffic/p5", map[string]any{"upload": 1000, "download": 2000})
  202. usage := func(p *panel) int64 {
  203. var tr struct{ Up, Down int64 }
  204. _ = json.Unmarshal(p.call(http.MethodGet, "/panel/api/clients/traffic/p5", nil), &tr)
  205. return tr.Up + tr.Down
  206. }
  207. eventually(t, settleTimeout, "master sees p5's node traffic", func() (bool, string) {
  208. u := usage(master)
  209. return u == 3000, fmt.Sprintf("master p5 usage %d", u)
  210. })
  211. master.call(http.MethodPost, "/panel/api/clients/resetTraffic/p5", nil)
  212. eventually(t, settleTimeout, "node p5 usage reset", func() (bool, string) {
  213. u := usage(node)
  214. return u == 0, fmt.Sprintf("node p5 usage %d", u)
  215. })
  216. time.Sleep(12 * time.Second)
  217. if u := usage(master); u != 0 {
  218. t.Fatalf("master p5 usage %d after reset settled, want 0", u)
  219. }
  220. })
  221. t.Run("an inbound deleted on the node is removed from the master too", func(t *testing.T) {
  222. pr.t = t
  223. master.call(http.MethodPost, "/panel/api/inbounds/add", vlessInbound("deleted-on-node", lostPort, &pr.nodeID, vlessClient("l1")))
  224. pr.waitNode("node holds the inbound", lostPort, func(ib inboundView) bool { return ib.client("l1") != nil })
  225. nib, _ := pr.nodeInbound(lostPort)
  226. node.call(http.MethodPost, "/panel/api/inbounds/del/"+strconv.Itoa(nib.Id), nil)
  227. eventually(t, settleTimeout, "master mirrors the node-side delete (#6219)", func() (bool, string) {
  228. _, still := pr.masterInbound(lostPort)
  229. return !still, "master still has the inbound"
  230. })
  231. })
  232. t.Run("deleting the inbound on the master removes it from the node", func(t *testing.T) {
  233. pr.t = t
  234. obj := master.call(http.MethodPost, "/panel/api/inbounds/add", vlessInbound("deleted-on-master", droppedPort, &pr.nodeID, vlessClient("d1")))
  235. var ib inboundView
  236. _ = json.Unmarshal(obj, &ib)
  237. pr.waitNode("node holds the inbound", droppedPort, func(ib inboundView) bool { return ib.client("d1") != nil })
  238. master.call(http.MethodPost, "/panel/api/inbounds/del/"+strconv.Itoa(ib.Id), nil)
  239. pr.waitNodeAbsent("node drops the inbound", droppedPort)
  240. })
  241. t.Run("a change made while the node is down reaches it once it is back", func(t *testing.T) {
  242. pr.t = t
  243. node.stop()
  244. pr.bulkAttach([]string{"p44"}, adoptedID)
  245. node.start()
  246. pr.waitNode("node holds p44 after restart", adoptedPort, func(ib inboundView) bool {
  247. return ib.client("p44") != nil
  248. })
  249. })
  250. t.Run("an inbound the node lost while down is re-created with the master's pending change", func(t *testing.T) {
  251. pr.t = t
  252. nib, ok := pr.nodeInbound(adoptedPort)
  253. if !ok {
  254. t.Fatal("node has no adopted inbound to lose")
  255. }
  256. node.stop()
  257. node.deleteInboundRow(nib.Id)
  258. pr.bulkAttach([]string{"p45"}, adoptedID)
  259. node.start()
  260. pr.waitNode("node re-creates the inbound with p44 and p45", adoptedPort, func(ib inboundView) bool {
  261. return ib.client("p44") != nil && ib.client("p45") != nil
  262. })
  263. })
  264. t.Run("an inbound created and edited while the node is down lands once it is back", func(t *testing.T) {
  265. pr.t = t
  266. node.stop()
  267. obj := master.call(http.MethodPost, "/panel/api/inbounds/add", vlessInbound("made-while-down", offlinePort, &pr.nodeID, vlessClient("o1")))
  268. var ib inboundView
  269. _ = json.Unmarshal(obj, &ib)
  270. body := vlessInbound("edited-while-down", offlinePort, &pr.nodeID)
  271. body["settings"] = `{"decryption":"none"}`
  272. master.call(http.MethodPost, "/panel/api/inbounds/update/"+strconv.Itoa(ib.Id), body)
  273. node.start()
  274. pr.waitNode("node holds the edited inbound with o1", offlinePort, func(ib inboundView) bool {
  275. return ib.Remark == "edited-while-down" && ib.client("o1") != nil
  276. })
  277. })
  278. t.Run("a change made while the node is disabled on the master lands once it is re-enabled", func(t *testing.T) {
  279. pr.t = t
  280. nodePath := "/panel/api/nodes/setEnable/" + strconv.Itoa(pr.nodeID)
  281. master.call(http.MethodPost, nodePath, map[string]any{"enable": false})
  282. pr.bulkAttach([]string{"p6"}, madeID)
  283. time.Sleep(6 * time.Second)
  284. if ib, _ := pr.nodeInbound(madePort); ib.client("p6") != nil {
  285. t.Fatal("a disabled node received a push")
  286. }
  287. master.call(http.MethodPost, nodePath, map[string]any{"enable": true})
  288. pr.waitNode("node holds p6 after re-enable", madePort, func(ib inboundView) bool { return ib.client("p6") != nil })
  289. })
  290. t.Run("selected sync mode leaves the node's unselected inbounds alone", func(t *testing.T) {
  291. pr.t = t
  292. var selected []string
  293. for _, ib := range master.inbounds() {
  294. if ib.NodeID != nil && *ib.NodeID == pr.nodeID {
  295. selected = append(selected, ib.Tag)
  296. }
  297. }
  298. master.call(http.MethodPost, "/panel/api/nodes/update/"+strconv.Itoa(pr.nodeID), map[string]any{
  299. "name": "n1", "scheme": "http", "address": "127.0.0.1", "port": node.port, "basePath": "/",
  300. "enable": true, "allowPrivateAddress": true, "inboundSyncMode": "selected", "inboundTags": selected,
  301. })
  302. node.call(http.MethodPost, "/panel/api/inbounds/add", vlessInbound("node-only", unmanagedPort, nil, vlessClient("u1")))
  303. pr.bulkAttach([]string{"p7"}, madeID)
  304. pr.waitNode("selected inbound still converges", madePort, func(ib inboundView) bool { return ib.client("p7") != nil })
  305. time.Sleep(12 * time.Second)
  306. if _, adopted := pr.masterInbound(unmanagedPort); adopted {
  307. t.Fatal("master adopted an unselected node inbound")
  308. }
  309. if ib, ok := pr.nodeInbound(unmanagedPort); !ok || ib.client("u1") == nil {
  310. t.Fatal("reconcile swept or rewrote an unselected node inbound")
  311. }
  312. })
  313. }