server.go 5.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227
  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/codes"
  18. "google.golang.org/grpc/reflection/grpc_reflection_v1alpha"
  19. )
  20. type (
  21. // Server is a gateway server.
  22. Server struct {
  23. c GatewayConf
  24. *rest.Server
  25. upstreams []*upstream
  26. timeout time.Duration
  27. processHeader func(http.Header) []string
  28. }
  29. // Option defines the method to customize Server.
  30. Option func(svr *Server)
  31. upstream struct {
  32. Upstream
  33. client zrpc.Client
  34. }
  35. )
  36. // MustNewServer creates a new gateway server.
  37. func MustNewServer(c GatewayConf, opts ...Option) *Server {
  38. svr := &Server{
  39. c: c,
  40. Server: rest.MustNewServer(c.RestConf),
  41. timeout: time.Duration(c.Timeout) * time.Millisecond,
  42. }
  43. for _, opt := range opts {
  44. opt(svr)
  45. }
  46. return svr
  47. }
  48. // Start starts the gateway server.
  49. func (s *Server) Start() {
  50. logx.Must(s.build())
  51. s.Server.Start()
  52. }
  53. // Stop stops the gateway server.
  54. func (s *Server) Stop() {
  55. s.Server.Stop()
  56. }
  57. func (s *Server) build() error {
  58. if err := s.buildClient(); err != nil {
  59. return err
  60. }
  61. return s.buildUpstream()
  62. }
  63. func (s *Server) buildClient() error {
  64. if err := s.ensureUpstreamNames(); err != nil {
  65. return err
  66. }
  67. return mr.MapReduceVoid(func(source chan<- Upstream) {
  68. for _, up := range s.c.Upstreams {
  69. source <- up
  70. }
  71. }, func(up Upstream, writer mr.Writer[*upstream], cancel func(error)) {
  72. target, err := up.Grpc.BuildTarget()
  73. if err != nil {
  74. cancel(err)
  75. }
  76. up.Name = target
  77. cli := zrpc.MustNewClient(up.Grpc)
  78. writer.Write(&upstream{
  79. Upstream: up,
  80. client: cli,
  81. })
  82. }, func(pipe <-chan *upstream, cancel func(error)) {
  83. for up := range pipe {
  84. s.upstreams = append(s.upstreams, up)
  85. }
  86. })
  87. }
  88. func (s *Server) buildUpstream() error {
  89. return mr.MapReduceVoid(func(source chan<- *upstream) {
  90. for _, up := range s.upstreams {
  91. source <- up
  92. }
  93. }, func(up *upstream, writer mr.Writer[rest.Route], cancel func(error)) {
  94. cli := up.client
  95. source, err := s.createDescriptorSource(cli, up.Upstream)
  96. if err != nil {
  97. cancel(fmt.Errorf("%s: %w", up.Name, err))
  98. return
  99. }
  100. methods, err := internal.GetMethods(source)
  101. if err != nil {
  102. cancel(fmt.Errorf("%s: %w", up.Name, err))
  103. return
  104. }
  105. resolver := grpcurl.AnyResolverFromDescriptorSource(source)
  106. for _, m := range methods {
  107. if len(m.HttpMethod) > 0 && len(m.HttpPath) > 0 {
  108. writer.Write(rest.Route{
  109. Method: m.HttpMethod,
  110. Path: m.HttpPath,
  111. Handler: s.buildHandler(source, resolver, cli, m.RpcPath),
  112. })
  113. }
  114. }
  115. methodSet := make(map[string]struct{})
  116. for _, m := range methods {
  117. methodSet[m.RpcPath] = struct{}{}
  118. }
  119. for _, m := range up.Mappings {
  120. if _, ok := methodSet[m.RpcPath]; !ok {
  121. cancel(fmt.Errorf("%s: rpc method %s not found", up.Name, m.RpcPath))
  122. return
  123. }
  124. writer.Write(rest.Route{
  125. Method: strings.ToUpper(m.Method),
  126. Path: m.Path,
  127. Handler: s.buildHandler(source, resolver, cli, m.RpcPath),
  128. })
  129. }
  130. }, func(pipe <-chan rest.Route, cancel func(error)) {
  131. for route := range pipe {
  132. s.Server.AddRoute(route)
  133. }
  134. })
  135. }
  136. func (s *Server) buildHandler(source grpcurl.DescriptorSource, resolver jsonpb.AnyResolver,
  137. cli zrpc.Client, rpcPath string) func(http.ResponseWriter, *http.Request) {
  138. return func(w http.ResponseWriter, r *http.Request) {
  139. parser, err := internal.NewRequestParser(r, resolver)
  140. if err != nil {
  141. httpx.ErrorCtx(r.Context(), w, err)
  142. return
  143. }
  144. timeout := internal.GetTimeout(r.Header, s.timeout)
  145. ctx, can := context.WithTimeout(r.Context(), timeout)
  146. defer can()
  147. w.Header().Set(httpx.ContentType, httpx.JsonContentType)
  148. handler := internal.NewEventHandler(w, resolver)
  149. if err := grpcurl.InvokeRPC(ctx, source, cli.Conn(), rpcPath, s.prepareMetadata(r.Header),
  150. handler, parser.Next); err != nil {
  151. httpx.ErrorCtx(r.Context(), w, err)
  152. }
  153. st := handler.Status
  154. if st.Code() != codes.OK {
  155. httpx.ErrorCtx(r.Context(), w, st.Err())
  156. }
  157. }
  158. }
  159. func (s *Server) createDescriptorSource(cli zrpc.Client, up Upstream) (grpcurl.DescriptorSource, error) {
  160. var source grpcurl.DescriptorSource
  161. var err error
  162. if len(up.ProtoSets) > 0 {
  163. source, err = grpcurl.DescriptorSourceFromProtoSets(up.ProtoSets...)
  164. if err != nil {
  165. return nil, err
  166. }
  167. } else {
  168. refCli := grpc_reflection_v1alpha.NewServerReflectionClient(cli.Conn())
  169. client := grpcreflect.NewClient(context.Background(), refCli)
  170. source = grpcurl.DescriptorSourceFromServer(context.Background(), client)
  171. }
  172. return source, nil
  173. }
  174. func (s *Server) ensureUpstreamNames() error {
  175. for _, up := range s.c.Upstreams {
  176. target, err := up.Grpc.BuildTarget()
  177. if err != nil {
  178. return err
  179. }
  180. up.Name = target
  181. }
  182. return nil
  183. }
  184. func (s *Server) prepareMetadata(header http.Header) []string {
  185. vals := internal.ProcessHeaders(header)
  186. if s.processHeader != nil {
  187. vals = append(vals, s.processHeader(header)...)
  188. }
  189. return vals
  190. }
  191. // WithHeaderProcessor sets a processor to process request headers.
  192. // The returned headers are used as metadata to invoke the RPC.
  193. func WithHeaderProcessor(processHeader func(http.Header) []string) func(*Server) {
  194. return func(s *Server) {
  195. s.processHeader = processHeader
  196. }
  197. }