Force grpc client/server to use grpc codec for broker

This commit is contained in:
Asim Aslam 2019-08-18 11:28:21 +01:00 committed by Vasiliy Tolstov
parent f9dcbb71be
commit 4a4159a9f6

View File

@ -1,18 +1,15 @@
package grpc package grpc
import ( import (
"bytes"
"context" "context"
"fmt" "fmt"
"reflect" "reflect"
"strings" "strings"
"github.com/micro/go-micro/broker" "github.com/micro/go-micro/broker"
"github.com/micro/go-micro/codec"
"github.com/micro/go-micro/metadata" "github.com/micro/go-micro/metadata"
"github.com/micro/go-micro/registry" "github.com/micro/go-micro/registry"
"github.com/micro/go-micro/server" "github.com/micro/go-micro/server"
"github.com/micro/go-micro/util/buf"
) )
const ( const (
@ -175,7 +172,7 @@ func (g *grpcServer) createSubHandler(sb *subscriber, opts server.Options) broke
msg.Header["Content-Type"] = defaultContentType msg.Header["Content-Type"] = defaultContentType
ct = defaultContentType ct = defaultContentType
} }
cf, err := g.newCodec(ct) cf, err := g.newGRPCCodec(ct)
if err != nil { if err != nil {
return err return err
} }
@ -205,15 +202,7 @@ func (g *grpcServer) createSubHandler(sb *subscriber, opts server.Options) broke
req = req.Elem() req = req.Elem()
} }
b := buf.New(bytes.NewBuffer(msg.Body)) if err := cf.Unmarshal(msg.Body, req.Interface()); err != nil {
co := cf(b)
defer co.Close()
if err := co.ReadHeader(&codec.Message{}, codec.Event); err != nil {
return err
}
if err := co.ReadBody(req.Interface()); err != nil {
return err return err
} }