1
0

server.go 4.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174
  1. package gateway
  2. import (
  3. "context"
  4. "fmt"
  5. "net/http"
  6. "strings"
  7. "time"
  8. "github.com/fullstorydev/grpcurl"
  9. "github.com/golang/protobuf/jsonpb"
  10. "github.com/jhump/protoreflect/grpcreflect"
  11. "github.com/zeromicro/go-zero/core/logx"
  12. "github.com/zeromicro/go-zero/core/mr"
  13. "github.com/zeromicro/go-zero/gateway/internal"
  14. "github.com/zeromicro/go-zero/rest"
  15. "github.com/zeromicro/go-zero/rest/httpx"
  16. "github.com/zeromicro/go-zero/zrpc"
  17. "google.golang.org/grpc/reflection/grpc_reflection_v1alpha"
  18. )
  19. type (
  20. // Server is a gateway server.
  21. Server struct {
  22. *rest.Server
  23. upstreams []Upstream
  24. timeout time.Duration
  25. processHeader func(http.Header) []string
  26. }
  27. // Option defines the method to customize Server.
  28. Option func(svr *Server)
  29. )
  30. // MustNewServer creates a new gateway server.
  31. func MustNewServer(c GatewayConf, opts ...Option) *Server {
  32. svr := &Server{
  33. Server: rest.MustNewServer(c.RestConf),
  34. upstreams: c.Upstreams,
  35. timeout: c.Timeout,
  36. }
  37. for _, opt := range opts {
  38. opt(svr)
  39. }
  40. return svr
  41. }
  42. // Start starts the gateway server.
  43. func (s *Server) Start() {
  44. logx.Must(s.build())
  45. s.Server.Start()
  46. }
  47. // Stop stops the gateway server.
  48. func (s *Server) Stop() {
  49. s.Server.Stop()
  50. }
  51. func (s *Server) build() error {
  52. return mr.MapReduceVoid(func(source chan<- interface{}) {
  53. for _, up := range s.upstreams {
  54. source <- up
  55. }
  56. }, func(item interface{}, writer mr.Writer, cancel func(error)) {
  57. up := item.(Upstream)
  58. cli := zrpc.MustNewClient(up.Grpc)
  59. source, err := s.createDescriptorSource(cli, up)
  60. if err != nil {
  61. cancel(err)
  62. return
  63. }
  64. methods, err := internal.GetMethods(source)
  65. if err != nil {
  66. cancel(err)
  67. return
  68. }
  69. resolver := grpcurl.AnyResolverFromDescriptorSource(source)
  70. for _, m := range methods {
  71. if len(m.HttpMethod) > 0 && len(m.HttpPath) > 0 {
  72. writer.Write(rest.Route{
  73. Method: m.HttpMethod,
  74. Path: m.HttpPath,
  75. Handler: s.buildHandler(source, resolver, cli, m.RpcPath),
  76. })
  77. }
  78. }
  79. methodSet := make(map[string]struct{})
  80. for _, m := range methods {
  81. methodSet[m.RpcPath] = struct{}{}
  82. }
  83. for _, m := range up.Mapping {
  84. if _, ok := methodSet[m.RpcPath]; !ok {
  85. cancel(fmt.Errorf("rpc method %s not found", m.RpcPath))
  86. return
  87. }
  88. writer.Write(rest.Route{
  89. Method: strings.ToUpper(m.Method),
  90. Path: m.Path,
  91. Handler: s.buildHandler(source, resolver, cli, m.RpcPath),
  92. })
  93. }
  94. }, func(pipe <-chan interface{}, cancel func(error)) {
  95. for item := range pipe {
  96. route := item.(rest.Route)
  97. s.Server.AddRoute(route)
  98. }
  99. })
  100. }
  101. func (s *Server) buildHandler(source grpcurl.DescriptorSource, resolver jsonpb.AnyResolver,
  102. cli zrpc.Client, rpcPath string) func(http.ResponseWriter, *http.Request) {
  103. return func(w http.ResponseWriter, r *http.Request) {
  104. handler := &grpcurl.DefaultEventHandler{
  105. Out: w,
  106. Formatter: grpcurl.NewJSONFormatter(true,
  107. grpcurl.AnyResolverFromDescriptorSource(source)),
  108. }
  109. parser, err := internal.NewRequestParser(r, resolver)
  110. if err != nil {
  111. httpx.Error(w, err)
  112. return
  113. }
  114. timeout := internal.GetTimeout(r.Header, s.timeout)
  115. ctx, can := context.WithTimeout(r.Context(), timeout)
  116. defer can()
  117. w.Header().Set(httpx.ContentType, httpx.JsonContentType)
  118. if err := grpcurl.InvokeRPC(ctx, source, cli.Conn(), rpcPath, s.prepareMetadata(r.Header),
  119. handler, parser.Next); err != nil {
  120. httpx.Error(w, err)
  121. }
  122. }
  123. }
  124. func (s *Server) createDescriptorSource(cli zrpc.Client, up Upstream) (grpcurl.DescriptorSource, error) {
  125. var source grpcurl.DescriptorSource
  126. var err error
  127. if len(up.ProtoSet) > 0 {
  128. source, err = grpcurl.DescriptorSourceFromProtoSets(up.ProtoSet)
  129. if err != nil {
  130. return nil, err
  131. }
  132. } else {
  133. refCli := grpc_reflection_v1alpha.NewServerReflectionClient(cli.Conn())
  134. client := grpcreflect.NewClient(context.Background(), refCli)
  135. source = grpcurl.DescriptorSourceFromServer(context.Background(), client)
  136. }
  137. return source, nil
  138. }
  139. func (s *Server) prepareMetadata(header http.Header) []string {
  140. vals := internal.ProcessHeaders(header)
  141. if s.processHeader != nil {
  142. vals = append(vals, s.processHeader(header)...)
  143. }
  144. return vals
  145. }
  146. // WithHeaderProcessor sets a processor to process request headers.
  147. // The returned headers are used as metadata to invoke the RPC.
  148. func WithHeaderProcessor(processHeader func(http.Header) []string) func(*Server) {
  149. return func(s *Server) {
  150. s.processHeader = processHeader
  151. }
  152. }