We now support configurable codecs for broker

This commit is contained in:
Asim Aslam 2016-12-06 19:55:09 +00:00 committed by Vasiliy Tolstov
parent 90f40cbad4
commit f8714d4cd9

12
nats.go
View File

@ -1,10 +1,10 @@
package nats package nats
import ( import (
"encoding/json"
"strings" "strings"
"github.com/micro/go-micro/broker" "github.com/micro/go-micro/broker"
"github.com/micro/go-micro/broker/codec/json"
"github.com/micro/go-micro/cmd" "github.com/micro/go-micro/cmd"
"github.com/nats-io/nats" "github.com/nats-io/nats"
) )
@ -100,7 +100,7 @@ func (n *nbroker) Options() broker.Options {
} }
func (n *nbroker) Publish(topic string, msg *broker.Message, opts ...broker.PublishOption) error { func (n *nbroker) Publish(topic string, msg *broker.Message, opts ...broker.PublishOption) error {
b, err := json.Marshal(msg) b, err := n.opts.Codec.Marshal(msg)
if err != nil { if err != nil {
return err return err
} }
@ -118,7 +118,7 @@ func (n *nbroker) Subscribe(topic string, handler broker.Handler, opts ...broker
fn := func(msg *nats.Msg) { fn := func(msg *nats.Msg) {
var m *broker.Message var m *broker.Message
if err := json.Unmarshal(msg.Data, &m); err != nil { if err := n.opts.Codec.Unmarshal(msg.Data, &m); err != nil {
return return
} }
handler(&publication{m: m, t: topic}) handler(&publication{m: m, t: topic})
@ -143,7 +143,11 @@ func (n *nbroker) String() string {
} }
func NewBroker(opts ...broker.Option) broker.Broker { func NewBroker(opts ...broker.Option) broker.Broker {
var options broker.Options options := broker.Options{
// Default codec
Codec: json.NewCodec(),
}
for _, o := range opts { for _, o := range opts {
o(&options) o(&options)
} }