redirect_test.go 4.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203
  1. package cluster
  2. import (
  3. "context"
  4. "fmt"
  5. "io"
  6. "net/http"
  7. "sync"
  8. "testing"
  9. "time"
  10. "github.com/gin-gonic/gin"
  11. "github.com/langgenius/dify-plugin-daemon/internal/utils/network"
  12. "github.com/langgenius/dify-plugin-daemon/pkg/entities/endpoint_entities"
  13. )
  14. type SimulationCheckServer struct {
  15. http.Server
  16. port uint16
  17. }
  18. func createSimulationSevers(nums int, register_callback func(i int, c *gin.Engine)) ([]*SimulationCheckServer, error) {
  19. gin.SetMode(gin.ReleaseMode)
  20. engines := make([]*gin.Engine, nums)
  21. servers := make([]*SimulationCheckServer, nums)
  22. for i := 0; i < nums; i++ {
  23. engines[i] = gin.Default()
  24. register_callback(i, engines[i])
  25. }
  26. // get random port
  27. ports := make([]uint16, nums)
  28. for i := 0; i < nums; i++ {
  29. port, err := network.GetRandomPort()
  30. if err != nil {
  31. return nil, err
  32. }
  33. ports[i] = port
  34. }
  35. for i := 0; i < nums; i++ {
  36. srv := &SimulationCheckServer{
  37. Server: http.Server{
  38. Addr: fmt.Sprintf(":%d", ports[i]),
  39. Handler: engines[i],
  40. },
  41. port: ports[i],
  42. }
  43. servers[i] = srv
  44. go func(i int) {
  45. srv.ListenAndServe()
  46. }(i)
  47. }
  48. return servers, nil
  49. }
  50. func closeSimulationHealthCheckSevers(servers []*SimulationCheckServer) {
  51. for _, server := range servers {
  52. server.Shutdown(context.Background())
  53. }
  54. }
  55. func TestRedirectTraffic(t *testing.T) {
  56. clearClusterState()
  57. // create 2 nodes cluster
  58. cluster, err := createSimulationCluster(2)
  59. if err != nil {
  60. t.Fatal(err)
  61. }
  62. // wait for voting
  63. wg := sync.WaitGroup{}
  64. wg.Add(len(cluster))
  65. // wait for all voting processes complete
  66. for _, node := range cluster {
  67. node := node
  68. go func() {
  69. defer wg.Done()
  70. <-node.NotifyVotingCompleted()
  71. }()
  72. }
  73. node1RecvReqs := make(chan struct{})
  74. node1RecvCorrectReqs := make(chan struct{})
  75. defer close(node1RecvReqs)
  76. defer close(node1RecvCorrectReqs)
  77. // create 2 simulated servers
  78. servers, err := createSimulationSevers(2, func(i int, c *gin.Engine) {
  79. c.GET("/plugin/invoke/tool", func(c *gin.Context) {
  80. if i == 0 {
  81. // redirect to node 1
  82. statusCode, headers, reader, err := cluster[i].RedirectRequest(cluster[1].id, c.Request)
  83. if err != nil {
  84. c.String(http.StatusInternalServerError, err.Error())
  85. return
  86. }
  87. c.Status(statusCode)
  88. for k, v := range headers {
  89. for _, vv := range v {
  90. c.Header(k, vv)
  91. }
  92. }
  93. io.Copy(c.Writer, reader)
  94. } else {
  95. c.String(http.StatusOK, "ok")
  96. node1RecvReqs <- struct{}{}
  97. }
  98. })
  99. c.GET("/health/check", func(c *gin.Context) {
  100. c.JSON(http.StatusOK, gin.H{
  101. "status": "ok",
  102. })
  103. })
  104. })
  105. if err != nil {
  106. t.Fatal(err)
  107. }
  108. defer closeSimulationHealthCheckSevers(servers)
  109. // change port
  110. for i, node := range cluster {
  111. node.port = servers[i].port
  112. }
  113. // launch cluster
  114. launchSimulationCluster(cluster)
  115. defer closeSimulationCluster(cluster, t)
  116. // wait for all nodes to be ready
  117. wg.Wait()
  118. // wait for node status to by synchronized
  119. wg = sync.WaitGroup{}
  120. wg.Add(len(cluster))
  121. // wait for all voting processes complete
  122. for _, node := range cluster {
  123. node := node
  124. go func() {
  125. defer wg.Done()
  126. <-node.NotifyNodeUpdateCompleted()
  127. }()
  128. }
  129. wg.Wait()
  130. // request to node 0
  131. go func() {
  132. for i := 0; i < 10; i++ {
  133. resp, err := http.Get(fmt.Sprintf("http://localhost:%d/plugin/invoke/tool", servers[0].port))
  134. if err != nil {
  135. t.Error(err)
  136. }
  137. content, err := io.ReadAll(resp.Body)
  138. if err != nil {
  139. t.Error(err)
  140. }
  141. if string(content) == "ok" {
  142. node1RecvCorrectReqs <- struct{}{}
  143. }
  144. }
  145. }()
  146. // check if node 1 received the request
  147. recvCount := 0
  148. correctCount := 0
  149. for {
  150. select {
  151. case <-node1RecvReqs:
  152. recvCount++
  153. case <-node1RecvCorrectReqs:
  154. correctCount++
  155. if correctCount == 10 {
  156. return
  157. }
  158. case <-time.After(5 * time.Second):
  159. t.Fatal("node 1 did not receive correct requests")
  160. }
  161. }
  162. }
  163. func TestRedirectTrafficWithQueryParams(t *testing.T) {
  164. request, err := http.NewRequest("GET", "http://localhost:8080/plugin/invoke/tool?a=1&b=2", nil)
  165. if err != nil {
  166. t.Fatal(err)
  167. }
  168. request.Header.Set(endpoint_entities.HeaderXOriginalHost, "localhost:8080")
  169. ip := address{
  170. Ip: "127.0.0.1",
  171. Port: 8080,
  172. }
  173. redirectedRequest := constructRedirectUrl(ip, request)
  174. if redirectedRequest != "http://127.0.0.1:8080/plugin/invoke/tool?a=1&b=2" {
  175. t.Fatal("redirected request is not correct")
  176. }
  177. }