2019-06-03 18:44:43 +01:00
|
|
|
// Package grpc transparently forwards the grpc protocol using a go-micro client.
|
|
|
|
package grpc
|
|
|
|
|
|
|
|
import (
|
|
|
|
"context"
|
|
|
|
"io"
|
|
|
|
"strings"
|
|
|
|
|
2020-07-27 13:22:00 +01:00
|
|
|
"github.com/micro/go-micro/v3/client"
|
|
|
|
"github.com/micro/go-micro/v3/client/grpc"
|
|
|
|
"github.com/micro/go-micro/v3/codec"
|
|
|
|
"github.com/micro/go-micro/v3/proxy"
|
|
|
|
"github.com/micro/go-micro/v3/server"
|
2019-06-03 18:44:43 +01:00
|
|
|
)
|
|
|
|
|
2019-06-06 17:55:32 +01:00
|
|
|
// Proxy will transparently proxy requests to the backend.
|
2019-06-03 18:44:43 +01:00
|
|
|
// If no backend is specified it will call a service using the client.
|
2019-06-06 17:58:21 +01:00
|
|
|
// If the service matches the Name it will use the server.DefaultRouter.
|
2019-06-06 17:55:32 +01:00
|
|
|
type Proxy struct {
|
2019-06-07 13:42:39 +01:00
|
|
|
// The proxy options
|
2019-12-16 14:55:47 +00:00
|
|
|
options proxy.Options
|
2019-06-03 18:44:43 +01:00
|
|
|
|
|
|
|
// Endpoint specified the fixed endpoint to call.
|
|
|
|
Endpoint string
|
|
|
|
|
|
|
|
// The client to use for outbound requests
|
|
|
|
Client client.Client
|
|
|
|
}
|
|
|
|
|
|
|
|
// read client request and write to server
|
|
|
|
func readLoop(r server.Request, s client.Stream) error {
|
|
|
|
// request to backend server
|
|
|
|
req := s.Request()
|
|
|
|
|
|
|
|
for {
|
|
|
|
// get data from client
|
|
|
|
// no need to decode it
|
|
|
|
body, err := r.Read()
|
|
|
|
if err == io.EOF {
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
// get the header from client
|
|
|
|
hdr := r.Header()
|
|
|
|
msg := &codec.Message{
|
|
|
|
Type: codec.Request,
|
|
|
|
Header: hdr,
|
|
|
|
Body: body,
|
|
|
|
}
|
|
|
|
// write the raw request
|
|
|
|
err = req.Codec().Write(msg, nil)
|
|
|
|
if err == io.EOF {
|
|
|
|
return nil
|
|
|
|
} else if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2019-11-25 16:31:43 +00:00
|
|
|
// ProcessMessage acts as a message exchange and forwards messages to ongoing topics
|
|
|
|
// TODO: should we look at p.Endpoint and only send to the local endpoint? probably
|
|
|
|
func (p *Proxy) ProcessMessage(ctx context.Context, msg server.Message) error {
|
|
|
|
// TODO: check that we're not broadcast storming by sending to the same topic
|
|
|
|
// that we're actually subscribed to
|
|
|
|
|
|
|
|
// directly publish to the local client
|
|
|
|
return p.Client.Publish(ctx, msg)
|
2019-08-23 14:05:11 +01:00
|
|
|
}
|
|
|
|
|
2019-06-06 17:55:32 +01:00
|
|
|
// ServeRequest honours the server.Proxy interface
|
|
|
|
func (p *Proxy) ServeRequest(ctx context.Context, req server.Request, rsp server.Response) error {
|
2019-06-03 18:44:43 +01:00
|
|
|
// set default client
|
|
|
|
if p.Client == nil {
|
2019-06-07 13:42:39 +01:00
|
|
|
p.Client = grpc.NewClient()
|
2019-06-03 18:44:43 +01:00
|
|
|
}
|
|
|
|
|
|
|
|
opts := []client.CallOption{}
|
|
|
|
|
|
|
|
// service name
|
|
|
|
service := req.Service()
|
|
|
|
endpoint := req.Endpoint()
|
|
|
|
|
|
|
|
// call a specific backend
|
2019-06-07 13:42:39 +01:00
|
|
|
if len(p.Endpoint) > 0 {
|
2019-06-03 18:44:43 +01:00
|
|
|
// address:port
|
2019-06-07 13:42:39 +01:00
|
|
|
if parts := strings.Split(p.Endpoint, ":"); len(parts) > 1 {
|
|
|
|
opts = append(opts, client.WithAddress(p.Endpoint))
|
2019-06-03 18:44:43 +01:00
|
|
|
// use as service name
|
|
|
|
} else {
|
2019-06-07 13:42:39 +01:00
|
|
|
service = p.Endpoint
|
2019-06-03 18:44:43 +01:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
// create new request with raw bytes body
|
2019-06-18 18:51:52 +01:00
|
|
|
creq := p.Client.NewRequest(service, endpoint, nil, client.WithContentType(req.ContentType()))
|
2019-06-03 18:44:43 +01:00
|
|
|
|
|
|
|
// create new stream
|
|
|
|
stream, err := p.Client.Stream(ctx, creq, opts...)
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
defer stream.Close()
|
|
|
|
|
|
|
|
// create client request read loop
|
|
|
|
go readLoop(req, stream)
|
|
|
|
|
|
|
|
// get raw response
|
|
|
|
resp := stream.Response()
|
|
|
|
|
|
|
|
// create server response write loop
|
|
|
|
for {
|
|
|
|
// read backend response body
|
|
|
|
body, err := resp.Read()
|
|
|
|
if err == io.EOF {
|
|
|
|
return nil
|
|
|
|
} else if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
// read backend response header
|
|
|
|
hdr := resp.Header()
|
|
|
|
|
|
|
|
// write raw response header to client
|
|
|
|
rsp.WriteHeader(hdr)
|
|
|
|
|
|
|
|
// write raw response body to client
|
|
|
|
err = rsp.Write(body)
|
|
|
|
if err == io.EOF {
|
|
|
|
return nil
|
|
|
|
} else if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2019-12-16 15:18:20 +00:00
|
|
|
func (p *Proxy) String() string {
|
2019-12-16 17:36:47 +00:00
|
|
|
return "grpc"
|
2019-12-16 15:18:20 +00:00
|
|
|
}
|
|
|
|
|
2019-06-06 17:55:32 +01:00
|
|
|
// NewProxy returns a new grpc proxy server
|
2019-12-16 14:55:47 +00:00
|
|
|
func NewProxy(opts ...proxy.Option) proxy.Proxy {
|
|
|
|
var options proxy.Options
|
|
|
|
for _, o := range opts {
|
|
|
|
o(&options)
|
2019-06-06 17:55:32 +01:00
|
|
|
}
|
2019-06-07 13:42:39 +01:00
|
|
|
|
2019-12-16 14:55:47 +00:00
|
|
|
p := new(Proxy)
|
|
|
|
p.Endpoint = options.Endpoint
|
|
|
|
p.Client = options.Client
|
2019-06-07 13:42:39 +01:00
|
|
|
|
|
|
|
return p
|
2019-06-06 17:55:32 +01:00
|
|
|
}
|
|
|
|
|
|
|
|
// NewSingleHostProxy returns a router which sends requests to a single backend
|
|
|
|
func NewSingleHostProxy(url string) *Proxy {
|
|
|
|
return &Proxy{
|
2019-06-07 13:42:39 +01:00
|
|
|
Endpoint: url,
|
2019-06-03 18:44:43 +01:00
|
|
|
}
|
|
|
|
}
|