2020-05-09 15:31:34 +08:00
|
|
|
package msgclient
|
|
|
|
|
|
|
|
import (
|
2020-05-18 02:10:45 +08:00
|
|
|
"context"
|
2020-05-09 15:43:53 +08:00
|
|
|
"crypto/md5"
|
|
|
|
"hash"
|
2020-05-09 15:31:34 +08:00
|
|
|
"io"
|
|
|
|
"log"
|
|
|
|
"time"
|
|
|
|
|
|
|
|
"github.com/chrislusf/seaweedfs/weed/messaging/broker"
|
|
|
|
"github.com/chrislusf/seaweedfs/weed/pb/messaging_pb"
|
|
|
|
)
|
|
|
|
|
|
|
|
type SubChannel struct {
|
2020-05-09 15:43:53 +08:00
|
|
|
ch chan []byte
|
|
|
|
stream messaging_pb.SeaweedMessaging_SubscribeClient
|
|
|
|
md5hash hash.Hash
|
2020-05-18 02:10:45 +08:00
|
|
|
cancel context.CancelFunc
|
2020-05-09 15:31:34 +08:00
|
|
|
}
|
|
|
|
|
2020-05-13 12:26:02 +08:00
|
|
|
func (mc *MessagingClient) NewSubChannel(subscriberId, chanName string) (*SubChannel, error) {
|
2020-05-09 15:31:34 +08:00
|
|
|
tp := broker.TopicPartition{
|
|
|
|
Namespace: "chan",
|
|
|
|
Topic: chanName,
|
|
|
|
Partition: 0,
|
|
|
|
}
|
|
|
|
grpcConnection, err := mc.findBroker(tp)
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
2020-05-18 02:10:45 +08:00
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
|
|
sc, err := setupSubscriberClient(ctx, grpcConnection, tp, subscriberId, time.Unix(0, 0))
|
2020-05-09 15:31:34 +08:00
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
t := &SubChannel{
|
2020-05-09 15:43:53 +08:00
|
|
|
ch: make(chan []byte),
|
|
|
|
stream: sc,
|
|
|
|
md5hash: md5.New(),
|
2020-05-18 02:10:45 +08:00
|
|
|
cancel: cancel,
|
2020-05-09 15:31:34 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
go func() {
|
|
|
|
for {
|
|
|
|
resp, subErr := t.stream.Recv()
|
|
|
|
if subErr == io.EOF {
|
|
|
|
return
|
|
|
|
}
|
|
|
|
if subErr != nil {
|
|
|
|
log.Printf("fail to receive from netchan %s: %v", chanName, subErr)
|
|
|
|
return
|
|
|
|
}
|
2020-05-16 23:57:29 +08:00
|
|
|
if resp.Data == nil {
|
|
|
|
// this could be heartbeat from broker
|
|
|
|
continue
|
|
|
|
}
|
2020-05-09 15:31:34 +08:00
|
|
|
if resp.Data.IsClose {
|
|
|
|
t.stream.Send(&messaging_pb.SubscriberMessage{
|
|
|
|
IsClose: true,
|
|
|
|
})
|
|
|
|
close(t.ch)
|
2020-05-18 02:10:45 +08:00
|
|
|
cancel()
|
2020-05-09 15:31:34 +08:00
|
|
|
return
|
|
|
|
}
|
|
|
|
t.ch <- resp.Data.Value
|
2020-05-10 18:48:35 +08:00
|
|
|
t.md5hash.Write(resp.Data.Value)
|
2020-05-09 15:31:34 +08:00
|
|
|
}
|
|
|
|
}()
|
|
|
|
|
|
|
|
return t, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func (sc *SubChannel) Channel() chan []byte {
|
|
|
|
return sc.ch
|
|
|
|
}
|
2020-05-09 15:43:53 +08:00
|
|
|
|
|
|
|
func (sc *SubChannel) Md5() []byte {
|
|
|
|
return sc.md5hash.Sum(nil)
|
|
|
|
}
|
2020-05-18 02:10:45 +08:00
|
|
|
|
|
|
|
func (sc *SubChannel) Cancel() {
|
|
|
|
sc.cancel()
|
|
|
|
}
|