You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
135 lines
3.3 KiB
135 lines
3.3 KiB
package deliver |
|
|
|
import ( |
|
"context" |
|
"fmt" |
|
"github.com/golang/protobuf/proto" |
|
"google.golang.org/grpc" |
|
"google.golang.org/grpc/credentials/insecure" |
|
"google.golang.org/grpc/resolver" |
|
"reflect" |
|
"sonet/api/gen/postal" |
|
"sonet/pkg/grpc/balancer" |
|
"sonet/pkg/grpc/discovery" |
|
"sonet/pkg/utils/logger" |
|
"time" |
|
) |
|
|
|
type Status int16 |
|
|
|
const ( |
|
StatusSuccess Status = 1 |
|
StatusError Status = 2 |
|
StatusReceiverOffline Status = 10 |
|
) |
|
|
|
// Deliver n包,通知消息投递 |
|
type Deliver struct { |
|
svcName string |
|
postal postal.PostalClient |
|
} |
|
|
|
func NewDeliver(msgInServiceName string) *Deliver { |
|
return &Deliver{ |
|
svcName: msgInServiceName, |
|
} |
|
} |
|
|
|
func (d *Deliver) InitWithResolver(ctx context.Context, builder resolver.Builder) (err error) { |
|
balancer.InitConsistentHashBuilder() |
|
|
|
postalUrl := discovery.EtcdDialUrl(postal.Postal_ServiceDesc.ServiceName) |
|
conn, err := grpc.DialContext(ctx, postalUrl, |
|
grpc.WithTransportCredentials(insecure.NewCredentials()), |
|
// consistent hash lb |
|
grpc.WithDefaultServiceConfig(fmt.Sprintf(`{"loadBalancingPolicy":"%s"}`, balancer.ConsistentHash)), |
|
grpc.WithResolvers(builder), |
|
) |
|
if err != nil { |
|
return |
|
} |
|
d.postal = postal.NewPostalClient(conn) |
|
return |
|
} |
|
|
|
func (d *Deliver) InitWithAddr(postalAddr string) (err error) { |
|
// Conn *grpc.ClientConn |
|
conn, err := grpc.Dial(postalAddr, grpc.WithTransportCredentials(insecure.NewCredentials())) |
|
if err != nil { |
|
return |
|
} |
|
d.postal = postal.NewPostalClient(conn) |
|
return |
|
} |
|
|
|
func (d *Deliver) Deliver(ctx context.Context, msg proto.Message, receiver string, options ...Option) (Status, error) { |
|
return d.deliver0(ctx, msg, []string{receiver}, options...) |
|
} |
|
|
|
func (d *Deliver) DeliverBatch(ctx context.Context, msg proto.Message, receivers []string, options ...Option) (Status, error) { |
|
return d.deliver0(ctx, msg, receivers, options...) |
|
} |
|
|
|
func (d *Deliver) deliver0(ctx context.Context, msg proto.Message, receivers []string, options ...Option) (status Status, err error) { |
|
opts := defaultOptions |
|
if options != nil { |
|
for _, opt := range options { |
|
opt.f(&opts) |
|
} |
|
} |
|
if receivers == nil || len(receivers) == 0 { |
|
// return StatusError, errors.New("receivers is empty") |
|
return StatusSuccess, nil |
|
} |
|
|
|
// encode msg |
|
body, err := proto.Marshal(msg) |
|
if err != nil { |
|
return StatusError, err |
|
} |
|
msgName := reflect.TypeOf(msg).Elem().Name() |
|
message := &postal.Message{ |
|
Time: time.Now().UnixMilli(), |
|
Svc: d.svcName, |
|
Msg: msgName, |
|
Body: body, |
|
} |
|
// deliver to gateway |
|
if len(receivers) == 1 { |
|
// deliver one receiver |
|
reqDeliver := &postal.ReqDeliver{ |
|
Receiver: receivers[0], |
|
Msg: message, |
|
} |
|
|
|
ctx = context.WithValue(ctx, balancer.ConsistentHashKey, reqDeliver.Receiver) |
|
res, err := d.postal.Deliver(ctx, reqDeliver) |
|
if err != nil { |
|
return StatusError, err |
|
} |
|
// TODO res code |
|
if res.Ok { |
|
status = StatusSuccess |
|
} else { |
|
status = StatusError |
|
} |
|
logger.Info("deliver result: ", err, res) |
|
} else { |
|
|
|
// deliver batch receiver |
|
ctx = context.WithValue(ctx, balancer.ConsistentHashKey, receivers[0]) |
|
req := &postal.ReqDeliverBatch{Receivers: receivers, Msg: message} |
|
res, err := d.postal.DeliverBatch(ctx, req) |
|
if err != nil { |
|
logger.Error("deliver error: ", err) |
|
return StatusError, err |
|
} |
|
if res.Ok { |
|
status = StatusSuccess |
|
} else { |
|
status = StatusError |
|
} |
|
logger.Info("deliver batch result: ", res) |
|
} |
|
return |
|
}
|
|
|