micro-client-grpc/codec.go

195 lines
4.1 KiB
Go
Raw Normal View History

2019-06-03 20:44:43 +03:00
package grpc
import (
"fmt"
2019-06-17 22:05:58 +03:00
"strings"
2019-06-03 20:44:43 +03:00
b "bytes"
"github.com/golang/protobuf/jsonpb"
2019-06-03 20:44:43 +03:00
"github.com/golang/protobuf/proto"
jsoniter "github.com/json-iterator/go"
2019-06-03 20:44:43 +03:00
"github.com/micro/go-micro/codec"
2019-06-17 22:05:58 +03:00
"github.com/micro/go-micro/codec/bytes"
2019-06-03 20:44:43 +03:00
"github.com/micro/go-micro/codec/jsonrpc"
"github.com/micro/go-micro/codec/protorpc"
2019-06-17 22:05:58 +03:00
"google.golang.org/grpc"
2019-06-18 20:51:52 +03:00
"google.golang.org/grpc/encoding"
2019-06-03 20:44:43 +03:00
)
type jsonCodec struct{}
type protoCodec struct{}
type bytesCodec struct{}
type wrapCodec struct{ encoding.Codec }
var jsonpbMarshaler = &jsonpb.Marshaler{}
2019-06-03 20:44:43 +03:00
var (
defaultGRPCCodecs = map[string]encoding.Codec{
"application/json": jsonCodec{},
"application/proto": protoCodec{},
"application/protobuf": protoCodec{},
"application/octet-stream": protoCodec{},
2019-08-27 10:13:58 +03:00
"application/grpc": protoCodec{},
2019-06-03 20:44:43 +03:00
"application/grpc+json": jsonCodec{},
"application/grpc+proto": protoCodec{},
"application/grpc+bytes": bytesCodec{},
}
defaultRPCCodecs = map[string]codec.NewCodec{
"application/json": jsonrpc.NewCodec,
"application/json-rpc": jsonrpc.NewCodec,
"application/protobuf": protorpc.NewCodec,
"application/proto-rpc": protorpc.NewCodec,
"application/octet-stream": protorpc.NewCodec,
}
json = jsoniter.ConfigCompatibleWithStandardLibrary
)
// UseNumber fix unmarshal Number(8234567890123456789) to interface(8.234567890123457e+18)
func UseNumber() {
json = jsoniter.Config{
UseNumber: true,
EscapeHTML: true,
SortMapKeys: true,
ValidateJsonRawMessage: true,
}.Froze()
}
func (w wrapCodec) String() string {
return w.Codec.Name()
}
2019-06-17 22:05:58 +03:00
func (w wrapCodec) Marshal(v interface{}) ([]byte, error) {
b, ok := v.(*bytes.Frame)
if ok {
return b.Data, nil
}
return w.Codec.Marshal(v)
}
func (w wrapCodec) Unmarshal(data []byte, v interface{}) error {
b, ok := v.(*bytes.Frame)
if ok {
b.Data = data
return nil
}
return w.Codec.Unmarshal(data, v)
}
2019-06-03 20:44:43 +03:00
func (protoCodec) Marshal(v interface{}) ([]byte, error) {
2019-06-17 22:05:58 +03:00
b, ok := v.(*bytes.Frame)
if ok {
return b.Data, nil
}
2019-06-03 20:44:43 +03:00
return proto.Marshal(v.(proto.Message))
}
func (protoCodec) Unmarshal(data []byte, v interface{}) error {
return proto.Unmarshal(data, v.(proto.Message))
}
func (protoCodec) Name() string {
return "proto"
}
func (bytesCodec) Marshal(v interface{}) ([]byte, error) {
b, ok := v.(*[]byte)
if !ok {
return nil, fmt.Errorf("failed to marshal: %v is not type of *[]byte", v)
}
return *b, nil
}
func (bytesCodec) Unmarshal(data []byte, v interface{}) error {
b, ok := v.(*[]byte)
if !ok {
return fmt.Errorf("failed to unmarshal: %v is not type of *[]byte", v)
}
*b = data
return nil
}
func (bytesCodec) Name() string {
return "bytes"
}
func (jsonCodec) Marshal(v interface{}) ([]byte, error) {
if pb, ok := v.(proto.Message); ok {
s, err := jsonpbMarshaler.MarshalToString(pb)
return []byte(s), err
}
2019-06-03 20:44:43 +03:00
return json.Marshal(v)
}
func (jsonCodec) Unmarshal(data []byte, v interface{}) error {
if pb, ok := v.(proto.Message); ok {
return jsonpb.Unmarshal(b.NewReader(data), pb)
}
2019-06-03 20:44:43 +03:00
return json.Unmarshal(data, v)
}
func (jsonCodec) Name() string {
return "json"
}
2019-06-17 22:05:58 +03:00
type grpcCodec struct {
// headers
2019-06-18 20:51:52 +03:00
id string
target string
method string
2019-06-17 22:05:58 +03:00
endpoint string
s grpc.ClientStream
c encoding.Codec
}
func (g *grpcCodec) ReadHeader(m *codec.Message, mt codec.MessageType) error {
md, err := g.s.Header()
if err != nil {
return err
}
if m == nil {
m = new(codec.Message)
}
if m.Header == nil {
m.Header = make(map[string]string)
}
for k, v := range md {
m.Header[k] = strings.Join(v, ",")
}
m.Id = g.id
m.Target = g.target
m.Method = g.method
m.Endpoint = g.endpoint
return nil
}
func (g *grpcCodec) ReadBody(v interface{}) error {
2019-06-18 20:51:52 +03:00
if f, ok := v.(*bytes.Frame); ok {
return g.s.RecvMsg(f)
2019-06-17 22:05:58 +03:00
}
2019-06-18 20:51:52 +03:00
return g.s.RecvMsg(v)
2019-06-17 22:05:58 +03:00
}
func (g *grpcCodec) Write(m *codec.Message, v interface{}) error {
// if we don't have a body
2019-06-18 20:51:52 +03:00
if v != nil {
return g.s.SendMsg(v)
2019-06-17 22:05:58 +03:00
}
// write the body using the framing codec
2019-06-18 20:51:52 +03:00
return g.s.SendMsg(&bytes.Frame{m.Body})
2019-06-17 22:05:58 +03:00
}
func (g *grpcCodec) Close() error {
return g.s.CloseSend()
}
func (g *grpcCodec) String() string {
return g.c.Name()
}