This commit is contained in:
joy.zhou
2019-07-31 12:10:26 +08:00
parent 3a9bbcb58a
commit 913bff4ed5
345 changed files with 114771 additions and 3099 deletions
+39 -305
View File
@@ -42,7 +42,6 @@ type Broker struct {
tlsConfig *tls.Config
wpool *pool.WorkerPool
clients sync.Map
routes sync.Map
remotes sync.Map
nodes map[string]interface{}
clusterPool chan *Message
@@ -115,7 +114,7 @@ func (b *Broker) SubmitWork(clientId string, msg *Message) {
b.wpool = pool.New(b.config.Worker)
}
if msg.client.typ == CLUSTER {
if msg.client.typ == ROUTER {
b.clusterPool <- msg
} else {
b.wpool.Submit(clientId, func() {
@@ -131,16 +130,14 @@ func (b *Broker) Start() {
return
}
go initRPCService()
go InitHTTPMoniter(b)
//listen clinet over tcp
if b.config.Port != "" {
go b.StartClientListening(false)
}
//listen for cluster
if b.config.Cluster.Port != "" {
go b.StartClusterListening()
//connet to router
if b.config.Router != "" {
go b.ConnectToDiscovery()
go b.processClusterInfo()
}
//listen for websocket
@@ -153,29 +150,19 @@ func (b *Broker) Start() {
go b.StartClientListening(true)
}
//connect on other node in cluster
if b.config.Router != "" {
go b.processClusterInfo()
b.ConnectToDiscovery()
//listen clinet over tcp
if b.config.Port != "" {
go b.StartClientListening(false)
}
//system monitor
go StateMonitor()
// if b.config.Debug {
// startPProf()
// }
}
func startPProf() {
go func() {
http.ListenAndServe(":10060", nil)
}()
}
func StateMonitor() {
v, _ := mem.VirtualMemory()
timeSticker := time.NewTicker(time.Second * 30)
timeSticker := time.NewTicker(time.Second * 60)
for {
select {
case <-timeSticker.C:
@@ -278,39 +265,6 @@ func TlsTimeout(conn *tls.Conn) {
}
}
func (b *Broker) StartClusterListening() {
var hp string = b.config.Cluster.Host + ":" + b.config.Cluster.Port
log.Info("Start Listening cluster on ", zap.String("hp", hp))
l, e := net.Listen("tcp", hp)
if e != nil {
log.Error("Error listening on ", zap.Error(e))
return
}
tmpDelay := 10 * ACCEPT_MIN_SLEEP
for {
conn, err := l.Accept()
if err != nil {
if ne, ok := err.(net.Error); ok && ne.Temporary() {
log.Error("Temporary Client Accept Error(%v), sleeping %dms",
zap.Error(ne), zap.Duration("sleeping", tmpDelay/time.Millisecond))
time.Sleep(tmpDelay)
tmpDelay *= 2
if tmpDelay > ACCEPT_MAX_SLEEP {
tmpDelay = ACCEPT_MAX_SLEEP
}
} else {
log.Error("Accept error: %v", zap.Error(err))
}
continue
}
tmpDelay = ACCEPT_MIN_SLEEP
go b.handleConnection(ROUTER, conn)
}
}
func (b *Broker) handleConnection(typ int, conn net.Conn) {
//process connect packet
packet, err := packets.ReadPacket(conn)
@@ -417,281 +371,61 @@ func (b *Broker) handleConnection(typ int, conn net.Conn) {
}
}
b.clients.Store(cid, c)
b.OnlineOfflineNotification(cid, true)
case ROUTER:
old, exist = b.routes.Load(cid)
if exist {
log.Warn("router exist, close old...")
ol, ok := old.(*client)
if ok {
ol.Close()
}
}
b.routes.Store(cid, c)
//TODO notify othernode to close connect
}
c.readLoop()
}
func (b *Broker) ConnectToDiscovery() {
var conn net.Conn
var err error
var tempDelay time.Duration = 0
for {
conn, err = net.Dial("tcp", b.config.Router)
if err != nil {
log.Error("Error trying to connect to route: ", zap.Error(err))
log.Debug("Connect to route timeout ,retry...")
if 0 == tempDelay {
tempDelay = 1 * time.Second
} else {
tempDelay *= 2
}
if max := 20 * time.Second; tempDelay > max {
tempDelay = max
}
time.Sleep(tempDelay)
continue
}
break
}
log.Debug("connect to router success :", zap.String("Router", b.config.Router))
cid := b.id
info := info{
clientID: cid,
keepalive: 60,
}
c := &client{
typ: CLUSTER,
broker: b,
conn: conn,
info: info,
}
c.init()
c.SendConnect()
c.SendInfo()
go c.readLoop()
go c.StartPing()
}
func (b *Broker) processClusterInfo() {
for {
msg, ok := <-b.clusterPool
if !ok {
log.Error("read message from cluster channel error")
return
}
ProcessMessage(msg)
}
}
func (b *Broker) connectRouter(id, addr string) {
var conn net.Conn
var err error
var timeDelay time.Duration = 0
retryTimes := 0
max := 32 * time.Second
for {
if !b.checkNodeExist(id, addr) {
return
}
conn, err = net.Dial("tcp", addr)
if err != nil {
log.Error("Error trying to connect to route: ", zap.Error(err))
if retryTimes > 50 {
return
}
log.Debug("Connect to route timeout ,retry...")
if 0 == timeDelay {
timeDelay = 1 * time.Second
} else {
timeDelay *= 2
}
if timeDelay > max {
timeDelay = max
}
time.Sleep(timeDelay)
retryTimes++
continue
}
break
}
route := route{
remoteID: id,
remoteUrl: addr,
}
cid := GenUniqueId()
info := info{
clientID: cid,
keepalive: 60,
}
c := &client{
broker: b,
typ: REMOTE,
conn: conn,
route: route,
info: info,
}
c.init()
b.remotes.Store(cid, c)
c.SendConnect()
// mpool := b.messagePool[fnv1a.HashString64(cid)%MessagePoolNum]
go c.readLoop()
go c.StartPing()
}
func (b *Broker) checkNodeExist(id, url string) bool {
if id == b.id {
return false
}
for k, v := range b.nodes {
if k == id {
return true
}
//skip
l, ok := v.(string)
if ok {
if url == l {
return true
}
}
}
return false
}
func (b *Broker) CheckRemoteExist(remoteID, url string) bool {
exist := false
b.remotes.Range(func(key, value interface{}) bool {
v, ok := value.(*client)
if ok {
if v.route.remoteUrl == url {
v.route.remoteID = remoteID
exist = true
return false
}
}
return true
})
return exist
}
func (b *Broker) SendLocalSubsToRouter(c *client) {
subInfo := packets.NewControlPacket(packets.Subscribe).(*packets.SubscribePacket)
b.clients.Range(func(key, value interface{}) bool {
client, ok := value.(*client)
if ok {
subs := client.subMap
for _, sub := range subs {
subInfo.Topics = append(subInfo.Topics, sub.topic)
subInfo.Qoss = append(subInfo.Qoss, sub.qos)
}
}
return true
})
if len(subInfo.Topics) > 0 {
err := c.WriterPacket(subInfo)
if err != nil {
log.Error("Send localsubs To Router error :", zap.Error(err))
}
}
}
func (b *Broker) BroadcastInfoMessage(remoteID string, msg *packets.PublishPacket) {
b.routes.Range(func(key, value interface{}) bool {
r, ok := value.(*client)
if ok {
if r.route.remoteID == remoteID {
return true
}
r.WriterPacket(msg)
}
return true
})
// log.Info("BroadcastInfoMessage success ")
}
func (b *Broker) BroadcastSubOrUnsubMessage(packet packets.ControlPacket) {
b.routes.Range(func(key, value interface{}) bool {
r, ok := value.(*client)
if ok {
r.WriterPacket(packet)
}
return true
})
// log.Info("BroadcastSubscribeMessage remotes: ", s.remotes)
}
func (b *Broker) removeClient(c *client) {
clientId := string(c.info.clientID)
typ := c.typ
switch typ {
case CLIENT:
b.clients.Delete(clientId)
case ROUTER:
b.routes.Delete(clientId)
case REMOTE:
b.remotes.Delete(clientId)
}
// log.Info("delete client ,", clientId)
b.clients.Delete(clientId)
}
func (b *Broker) PublishMessage(packet *packets.PublishPacket) {
{
//do retain
if packet.Retain {
if err := b.topicsMgr.Retain(packet); err != nil {
log.Error("Error retaining message: ", zap.Error(err))
}
}
}
var subs []interface{}
var qoss []byte
b.mu.Lock()
err := b.topicsMgr.Subscribers([]byte(packet.TopicName), packet.Qos, &subs, &qoss)
b.mu.Unlock()
if err != nil {
log.Error("search sub client error, ", zap.Error(err))
return
}
for _, sub := range subs {
if len(subs) == 0 {
return
}
var qsub []int
for i, sub := range subs {
s, ok := sub.(*subscription)
if ok {
err := s.client.WriterPacket(packet)
if err != nil {
log.Error("write message error, ", zap.Error(err))
if s.share {
qsub = append(qsub, i)
} else {
publish(s, packet)
}
}
}
}
func (b *Broker) BroadcastUnSubscribe(subs map[string]*subscription) {
unsub := packets.NewControlPacket(packets.Unsubscribe).(*packets.UnsubscribePacket)
for topic, _ := range subs {
unsub.Topics = append(unsub.Topics, topic)
}
if len(unsub.Topics) > 0 {
b.BroadcastSubOrUnsubMessage(unsub)
if len(qsub) > 0 {
idx := r.Intn(len(qsub))
sub := subs[qsub[idx]].(*subscription)
publish(sub, packet)
}
}
func (b *Broker) OnlineOfflineNotification(clientID string, online bool) {
+4 -197
View File
@@ -26,11 +26,8 @@ const (
BrokerInfoTopic = "broker000100101info"
// CLIENT is an end user.
CLIENT = 0
// ROUTER is another router in the cluster.
// ROUTER is client in the router.
ROUTER = 1
//REMOTE is the router connect to other cluster
REMOTE = 2
CLUSTER = 3
)
const (
@@ -193,8 +190,6 @@ func (c *client) ProcessPublish(packet *packets.PublishPacket) {
case CLIENT:
c.processClientPublish(packet)
case ROUTER:
c.processRouterPublish(packet)
case CLUSTER:
c.processRemotePublish(packet)
}
@@ -213,31 +208,6 @@ func (c *client) processRemotePublish(packet *packets.PublishPacket) {
}
func (c *client) processRouterPublish(packet *packets.PublishPacket) {
if c.status == Disconnected {
return
}
switch packet.Qos {
case QosAtMostOnce:
c.ProcessPublishMessage(packet)
case QosAtLeastOnce:
puback := packets.NewControlPacket(packets.Puback).(*packets.PubackPacket)
puback.MessageID = packet.MessageID
if err := c.WriterPacket(puback); err != nil {
log.Error("send puback error, ", zap.Error(err), zap.String("ClientID", c.info.clientID))
return
}
c.ProcessPublishMessage(packet)
case QosExactlyOnce:
return
default:
log.Error("publish with unknown qos", zap.String("ClientID", c.info.clientID))
return
}
}
func (c *client) processClientPublish(packet *packets.PublishPacket) {
if c.status == Disconnected {
return
@@ -262,7 +232,7 @@ func (c *client) processClientPublish(packet *packets.PublishPacket) {
switch packet.Qos {
case QosAtMostOnce:
c.ProcessPublishMessage(packet)
c.broker.PublishMessage(packet)
case QosAtLeastOnce:
puback := packets.NewControlPacket(packets.Puback).(*packets.PubackPacket)
puback.MessageID = packet.MessageID
@@ -270,7 +240,7 @@ func (c *client) processClientPublish(packet *packets.PublishPacket) {
log.Error("send puback error, ", zap.Error(err), zap.String("ClientID", c.info.clientID))
return
}
c.ProcessPublishMessage(packet)
c.broker.PublishMessage(packet)
case QosExactlyOnce:
return
default:
@@ -280,64 +250,10 @@ func (c *client) processClientPublish(packet *packets.PublishPacket) {
}
func (c *client) ProcessPublishMessage(packet *packets.PublishPacket) {
b := c.broker
if b == nil {
return
}
typ := c.typ
if packet.Retain {
if err := c.topicsMgr.Retain(packet); err != nil {
log.Error("Error retaining message: ", zap.Error(err), zap.String("ClientID", c.info.clientID))
}
}
err := c.topicsMgr.Subscribers([]byte(packet.TopicName), packet.Qos, &c.subs, &c.qoss)
if err != nil {
log.Error("Error retrieving subscribers list: ", zap.String("ClientID", c.info.clientID))
return
}
// fmt.Println("psubs num: ", len(c.subs))
if len(c.subs) == 0 {
return
}
var qsub []int
for i, sub := range c.subs {
s, ok := sub.(*subscription)
if ok {
if s.client.typ == ROUTER {
if typ != CLIENT {
continue
}
}
if s.share {
qsub = append(qsub, i)
} else {
publish(s, packet)
}
}
}
if len(qsub) > 0 {
idx := r.Intn(len(qsub))
sub := c.subs[qsub[idx]].(*subscription)
publish(sub, packet)
}
}
func (c *client) ProcessSubscribe(packet *packets.SubscribePacket) {
switch c.typ {
case CLIENT:
c.processClientSubscribe(packet)
case ROUTER:
c.processRouterSubscribe(packet)
}
}
@@ -417,8 +333,6 @@ func (c *client) processClientSubscribe(packet *packets.SubscribePacket) {
log.Error("send suback error, ", zap.Error(err), zap.String("ClientID", c.info.clientID))
return
}
//broadcast subscribe message
go b.BroadcastSubOrUnsubMessage(packet)
//process retain message
for _, rm := range c.rmsgs {
@@ -430,106 +344,10 @@ func (c *client) processClientSubscribe(packet *packets.SubscribePacket) {
}
}
func (c *client) processRouterSubscribe(packet *packets.SubscribePacket) {
if c.status == Disconnected {
return
}
b := c.broker
if b == nil {
return
}
topics := packet.Topics
qoss := packet.Qoss
suback := packets.NewControlPacket(packets.Suback).(*packets.SubackPacket)
suback.MessageID = packet.MessageID
var retcodes []byte
for i, topic := range topics {
t := topic
groupName := ""
share := false
if strings.HasPrefix(topic, "$share/") {
substr := groupCompile.FindStringSubmatch(topic)
if len(substr) != 3 {
retcodes = append(retcodes, QosFailure)
continue
}
share = true
groupName = substr[1]
topic = substr[2]
}
sub := &subscription{
topic: topic,
qos: qoss[i],
client: c,
share: share,
groupName: groupName,
}
rqos, err := c.topicsMgr.Subscribe([]byte(topic), qoss[i], sub)
if err != nil {
log.Error("subscribe error, ", zap.Error(err), zap.String("ClientID", c.info.clientID))
retcodes = append(retcodes, QosFailure)
continue
}
c.subMap[t] = sub
addSubMap(c.routeSubMap, topic)
retcodes = append(retcodes, rqos)
}
suback.ReturnCodes = retcodes
err := c.WriterPacket(suback)
if err != nil {
log.Error("send suback error, ", zap.Error(err), zap.String("ClientID", c.info.clientID))
return
}
}
func (c *client) ProcessUnSubscribe(packet *packets.UnsubscribePacket) {
switch c.typ {
case CLIENT:
c.processClientUnSubscribe(packet)
case ROUTER:
c.processRouterUnSubscribe(packet)
}
}
func (c *client) processRouterUnSubscribe(packet *packets.UnsubscribePacket) {
if c.status == Disconnected {
return
}
b := c.broker
if b == nil {
return
}
topics := packet.Topics
for _, topic := range topics {
sub, exist := c.subMap[topic]
if exist {
retainNum := delSubMap(c.routeSubMap, topic)
if retainNum > 0 {
continue
}
c.topicsMgr.Unsubscribe([]byte(sub.topic), sub)
delete(c.subMap, topic)
}
}
unsuback := packets.NewControlPacket(packets.Unsuback).(*packets.UnsubackPacket)
unsuback.MessageID = packet.MessageID
err := c.WriterPacket(unsuback)
if err != nil {
log.Error("send unsuback error, ", zap.Error(err), zap.String("ClientID", c.info.clientID))
return
}
}
@@ -574,8 +392,6 @@ func (c *client) processClientUnSubscribe(packet *packets.UnsubscribePacket) {
log.Error("send unsuback error, ", zap.Error(err), zap.String("ClientID", c.info.clientID))
return
}
// //process ubsubscribe message
b.BroadcastSubOrUnsubMessage(packet)
}
func (c *client) ProcessPing() {
@@ -598,9 +414,6 @@ func (c *client) Close() {
c.cancelFunc()
c.status = Disconnected
//wait for message complete
// time.Sleep(1 * time.Second)
// c.status = Disconnected
b := c.broker
b.Publish(&plugins.Elements{
@@ -627,7 +440,6 @@ func (c *client) Close() {
}
if c.typ == CLIENT {
b.BroadcastUnSubscribe(subs)
//offline notification
b.OnlineOfflineNotification(c.info.clientID, false)
}
@@ -636,14 +448,9 @@ func (c *client) Close() {
b.PublishMessage(c.info.willMsg)
}
if c.typ == CLUSTER {
if c.typ == ROUTER {
b.ConnectToDiscovery()
}
//do reconnect
if c.typ == REMOTE {
go b.connectRouter(c.route.remoteID, c.route.remoteUrl)
}
}
}
+1 -1
View File
@@ -104,7 +104,7 @@ func (c *client) ProcessInfo(packet *packets.PublishPacket) {
if ok {
exist := b.CheckRemoteExist(rid, url)
if !exist {
b.connectRouter(rid, url)
//todo new rpc client
}
}
+106
View File
@@ -0,0 +1,106 @@
package broker
import (
"net"
"time"
"go.uber.org/zap"
)
func (b *Broker) processClusterInfo() {
for {
msg, ok := <-b.clusterPool
if !ok {
log.Error("read message from cluster channel error")
return
}
ProcessMessage(msg)
}
}
func (b *Broker) ConnectToDiscovery() {
var conn net.Conn
var err error
var tempDelay time.Duration = 0
for {
conn, err = net.Dial("tcp", b.config.Router)
if err != nil {
log.Error("Error trying to connect to route: ", zap.Error(err))
log.Debug("Connect to route timeout ,retry...")
if 0 == tempDelay {
tempDelay = 1 * time.Second
} else {
tempDelay *= 2
}
if max := 20 * time.Second; tempDelay > max {
tempDelay = max
}
time.Sleep(tempDelay)
continue
}
break
}
log.Debug("connect to router success :", zap.String("Router", b.config.Router))
cid := b.id
info := info{
clientID: cid,
keepalive: 60,
}
c := &client{
typ: ROUTER,
broker: b,
conn: conn,
info: info,
}
c.init()
c.SendConnect()
c.SendInfo()
go c.readLoop()
go c.StartPing()
}
func (b *Broker) checkNodeExist(id, url string) bool {
if id == b.id {
return false
}
for k, v := range b.nodes {
if k == id {
return true
}
//skip
l, ok := v.(string)
if ok {
if url == l {
return true
}
}
}
return false
}
func (b *Broker) CheckRemoteExist(remoteID, url string) bool {
exist := false
b.remotes.Range(func(key, value interface{}) bool {
v, ok := value.(*client)
if ok {
if v.route.remoteUrl == url {
v.route.remoteID = remoteID
exist = true
return false
}
}
return true
})
return exist
}
+68
View File
@@ -0,0 +1,68 @@
package broker
import (
"context"
"net"
"time"
"github.com/eclipse/paho.mqtt.golang/packets"
pb "github.com/fhmq/hmq/grpc"
"go.uber.org/zap"
"google.golang.org/grpc"
"google.golang.org/grpc/keepalive"
"google.golang.org/grpc/reflection"
)
var (
rpcClient = make(map[string]pb.HMQServiceClient)
)
func initRPCService() {
lis, err := net.Listen("tcp", ":10011")
if err != nil {
log.Error("failed to listen: ", zap.Error(err))
return
}
s := grpc.NewServer(grpc.KeepaliveParams(keepalive.ServerParameters{
Time: 30 * time.Minute,
}))
pb.RegisterHMQServiceServer(s, &HMQ{})
reflection.Register(s)
if err := s.Serve(lis); err != nil {
log.Error("failed to serve: ", zap.Error(err))
}
}
func initRPCClient(url string) {
conn, err := grpc.Dial(url,
grpc.WithInsecure(),
grpc.WithKeepaliveParams(keepalive.ClientParameters{ // avoid 'code = Unavailable desc = transport is closing' error
Time: 30 * time.Minute,
}))
if err != nil {
log.Error("create connect rpc service failed", zap.String("url", url), zap.Error(err))
}
cli := pb.NewHMQServiceClient(conn)
rpcClient[url] = cli
}
type HMQ struct {
}
func (h *HMQ) QuerySubscribe(ctx context.Context, in *pb.QuerySubscribeRequest) (*pb.Response, error) {
return nil, nil
}
func (h *HMQ) QueryConnect(ctx context.Context, in *pb.QueryConnectRequest) (*pb.Response, error) {
return nil, nil
}
func (h *HMQ) DeliverMessage(ctx context.Context, in *pb.DeliverMessageRequest) (*pb.Response, error) {
return nil, nil
}
func (b *Broker) DeliverMessage(packet *packets.PublishPacket) {
//TODO Query and Deliver Message
}