mirror of
https://github.com/fhmq/hmq.git
synced 2026-08-31 23:04:52 +00:00
update queue shub
This commit is contained in:
+92
-25
@@ -3,6 +3,7 @@
|
||||
package broker
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/tls"
|
||||
"fmt"
|
||||
"net"
|
||||
@@ -387,7 +388,16 @@ func (b *Broker) removeClient(c *client) {
|
||||
b.clients.Delete(clientId)
|
||||
}
|
||||
|
||||
func (b *Broker) PublishMessage(packet *packets.PublishPacket, deliver bool) {
|
||||
func (b *Broker) OnlineOfflineNotification(clientID string, online bool) {
|
||||
packet := packets.NewControlPacket(packets.Publish).(*packets.PublishPacket)
|
||||
packet.TopicName = "$SYS/broker/connection/clients/" + clientID
|
||||
packet.Qos = 0
|
||||
packet.Payload = []byte(fmt.Sprintf(`{"clientID":"%s","online":%v,"timestamp":"%s"}`, clientID, online, time.Now().UTC().Format(time.RFC3339)))
|
||||
|
||||
b.PublishMessage(packet)
|
||||
}
|
||||
|
||||
func (b *Broker) PublishDeliverdMessage(packet *packets.PublishPacket, share bool) {
|
||||
{
|
||||
//do retain
|
||||
if packet.Retain {
|
||||
@@ -397,13 +407,6 @@ func (b *Broker) PublishMessage(packet *packets.PublishPacket, deliver bool) {
|
||||
}
|
||||
}
|
||||
|
||||
{
|
||||
//deliver message to other node
|
||||
if deliver {
|
||||
go b.DeliverMessage(packet)
|
||||
}
|
||||
}
|
||||
|
||||
var subs []interface{}
|
||||
var qoss []byte
|
||||
err := b.topicsMgr.Subscribers([]byte(packet.TopicName), packet.Qos, &subs, &qoss)
|
||||
@@ -416,33 +419,97 @@ func (b *Broker) PublishMessage(packet *packets.PublishPacket, deliver bool) {
|
||||
return
|
||||
}
|
||||
|
||||
var qsub []int
|
||||
for i, sub := range subs {
|
||||
var qsub []*subscription
|
||||
for _, sub := range subs {
|
||||
s, ok := sub.(*subscription)
|
||||
if ok {
|
||||
if s.share {
|
||||
qsub = append(qsub, i)
|
||||
qsub = append(qsub, s)
|
||||
} else {
|
||||
publish(s, packet)
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
if len(qsub) > 0 {
|
||||
idx := r.Intn(len(qsub))
|
||||
sub := subs[qsub[idx]].(*subscription)
|
||||
if share {
|
||||
target := r.Intn(len(qsub))
|
||||
sub := qsub[target]
|
||||
publish(sub, packet)
|
||||
}
|
||||
}
|
||||
|
||||
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
|
||||
err := b.topicsMgr.Subscribers([]byte(packet.TopicName), packet.Qos, &subs, &qoss)
|
||||
if err != nil {
|
||||
log.Error("search sub client error, ", zap.Error(err))
|
||||
return
|
||||
}
|
||||
|
||||
if len(subs) == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
var qsub []*subscription
|
||||
for _, sub := range subs {
|
||||
s, ok := sub.(*subscription)
|
||||
if ok {
|
||||
if s.share {
|
||||
qsub = append(qsub, s)
|
||||
} else {
|
||||
publish(s, packet)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
b.ProcessRemote(packet, qsub)
|
||||
|
||||
}
|
||||
|
||||
func (b *Broker) ProcessRemote(packet *packets.PublishPacket, loaclShareSub []*subscription) {
|
||||
shareRemoteID := ""
|
||||
remoteSubInfo := b.QuerySubscribe(packet.TopicName, packet.Qos)
|
||||
totalShare := len(loaclShareSub)
|
||||
for _, v := range remoteSubInfo {
|
||||
totalShare = totalShare + v.shareSubCount
|
||||
}
|
||||
target := r.Intn(totalShare)
|
||||
|
||||
if target < len(loaclShareSub) {
|
||||
shareRemoteID = b.id
|
||||
} else {
|
||||
target = target - len(loaclShareSub)
|
||||
for k, v := range remoteSubInfo {
|
||||
if target < v.shareSubCount {
|
||||
shareRemoteID = k
|
||||
return
|
||||
}
|
||||
target = target - v.shareSubCount
|
||||
}
|
||||
}
|
||||
|
||||
//send remote
|
||||
for id, sub := range remoteSubInfo {
|
||||
rpcCli := b.rpcClient[id]
|
||||
if sub.subCount > 0 {
|
||||
rpcCli.DeliverMessage(context.Background(), &pb.DeliverMessageRequest{Topic: packet.TopicName, Payload: packet.Payload, Share: shareRemoteID == id})
|
||||
}
|
||||
}
|
||||
|
||||
//send local
|
||||
if shareRemoteID == b.id {
|
||||
sub := loaclShareSub[target]
|
||||
publish(sub, packet)
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
func (b *Broker) OnlineOfflineNotification(clientID string, online bool) {
|
||||
packet := packets.NewControlPacket(packets.Publish).(*packets.PublishPacket)
|
||||
packet.TopicName = "$SYS/broker/connection/clients/" + clientID
|
||||
packet.Qos = 0
|
||||
packet.Payload = []byte(fmt.Sprintf(`{"clientID":"%s","online":%v,"timestamp":"%s"}`, clientID, online, time.Now().UTC().Format(time.RFC3339)))
|
||||
|
||||
b.PublishMessage(packet, true)
|
||||
}
|
||||
|
||||
+3
-3
@@ -228,7 +228,7 @@ func (c *client) processClientPublish(packet *packets.PublishPacket) {
|
||||
|
||||
switch packet.Qos {
|
||||
case QosAtMostOnce:
|
||||
c.broker.PublishMessage(packet, true)
|
||||
c.broker.PublishMessage(packet)
|
||||
case QosAtLeastOnce:
|
||||
puback := packets.NewControlPacket(packets.Puback).(*packets.PubackPacket)
|
||||
puback.MessageID = packet.MessageID
|
||||
@@ -236,7 +236,7 @@ func (c *client) processClientPublish(packet *packets.PublishPacket) {
|
||||
log.Error("send puback error, ", zap.Error(err), zap.String("ClientID", c.info.clientID))
|
||||
return
|
||||
}
|
||||
c.broker.PublishMessage(packet, true)
|
||||
c.broker.PublishMessage(packet)
|
||||
case QosExactlyOnce:
|
||||
return
|
||||
default:
|
||||
@@ -439,7 +439,7 @@ func (c *client) Close() {
|
||||
//offline notification
|
||||
b.OnlineOfflineNotification(c.info.clientID, false)
|
||||
if c.info.willMsg != nil {
|
||||
b.PublishMessage(c.info.willMsg, true)
|
||||
b.PublishMessage(c.info.willMsg)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+10
-1
@@ -11,6 +11,7 @@ import (
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
"os"
|
||||
"strconv"
|
||||
|
||||
"github.com/fhmq/hmq/logger"
|
||||
"go.uber.org/zap"
|
||||
@@ -73,7 +74,6 @@ func ConfigureConfig(args []string) (*Config, error) {
|
||||
fs.StringVar(&config.Port, "port", "1883", "Port to listen on.")
|
||||
fs.StringVar(&config.Port, "p", "1883", "Port to listen on.")
|
||||
fs.StringVar(&config.Host, "host", "0.0.0.0", "Network host to listen on")
|
||||
fs.StringVar(&config.RpcPort, "rpc", "10011", "Port to listen on.")
|
||||
fs.StringVar(&config.Router, "r", "", "Router who maintenance cluster info")
|
||||
fs.StringVar(&config.Router, "router", "", "Router who maintenance cluster info")
|
||||
fs.StringVar(&config.WsPort, "ws", "", "port for ws to listen on")
|
||||
@@ -164,6 +164,11 @@ func (config *Config) check() error {
|
||||
config.TlsHost = "0.0.0.0"
|
||||
}
|
||||
}
|
||||
|
||||
if config.Router != "" {
|
||||
config.RpcPort = strconv.Itoa(randInt())
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -205,3 +210,7 @@ func NewTLSConfig(tlsInfo TLSInfo) (*tls.Config, error) {
|
||||
|
||||
return &config, nil
|
||||
}
|
||||
|
||||
func randInt() int {
|
||||
return r.Intn(1000) + 10000
|
||||
}
|
||||
|
||||
+25
-60
@@ -48,8 +48,8 @@ type HMQ struct {
|
||||
b *Broker
|
||||
}
|
||||
|
||||
func (h *HMQ) QuerySubscribe(ctx context.Context, in *pb.QuerySubscribeRequest) (*pb.Response, error) {
|
||||
resp := &pb.Response{
|
||||
func (h *HMQ) QuerySubscribe(ctx context.Context, in *pb.QuerySubscribeRequest) (*pb.SubscribeResponse, error) {
|
||||
resp := &pb.SubscribeResponse{
|
||||
RetCode: 0,
|
||||
}
|
||||
topic := in.Topic
|
||||
@@ -71,6 +71,18 @@ func (h *HMQ) QuerySubscribe(ctx context.Context, in *pb.QuerySubscribeRequest)
|
||||
if len(subs) == 0 {
|
||||
resp.RetCode = 404
|
||||
}
|
||||
resp.SubCount = int32(len(subs))
|
||||
|
||||
var qsub int32
|
||||
for _, sub := range subs {
|
||||
s, ok := sub.(*subscription)
|
||||
if ok {
|
||||
if s.share {
|
||||
qsub++
|
||||
}
|
||||
}
|
||||
}
|
||||
resp.ShareSubCount = qsub
|
||||
|
||||
return resp, nil
|
||||
}
|
||||
@@ -96,7 +108,7 @@ func (h *HMQ) DeliverMessage(ctx context.Context, in *pb.DeliverMessageRequest)
|
||||
p.TopicName = in.Topic
|
||||
p.Payload = in.Payload
|
||||
p.Retain = false
|
||||
b.PublishMessage(p, false)
|
||||
b.PublishDeliverdMessage(p, in.Share)
|
||||
|
||||
resp := &pb.Response{
|
||||
RetCode: 0,
|
||||
@@ -104,57 +116,8 @@ func (h *HMQ) DeliverMessage(ctx context.Context, in *pb.DeliverMessageRequest)
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (h *HMQ) QueryShareSubscribe(ctx context.Context, in *pb.QueryShareSubscribeRequest) (*pb.ShareSubscribeResponse, error) {
|
||||
resp := &pb.ShareSubscribeResponse{
|
||||
RetCode: 0,
|
||||
}
|
||||
topic := in.Topic
|
||||
qos := in.Qos
|
||||
if qos > 1 {
|
||||
resp.RetCode = 404
|
||||
return resp, nil
|
||||
}
|
||||
func (b *Broker) DeliverMessage(packet *packets.PublishPacket, shareRemoteID string) {
|
||||
|
||||
b := h.b
|
||||
var qoss []byte
|
||||
var subs []interface{}
|
||||
err := b.topicsMgr.Subscribers([]byte(topic), byte(qos), &subs, &qoss)
|
||||
if err != nil {
|
||||
log.Error("search sub client error, ", zap.Error(err))
|
||||
resp.RetCode = 404
|
||||
}
|
||||
|
||||
if len(subs) == 0 {
|
||||
resp.RetCode = 404
|
||||
}
|
||||
|
||||
var qsub int32
|
||||
for _, sub := range subs {
|
||||
s, ok := sub.(*subscription)
|
||||
if ok {
|
||||
if s.share {
|
||||
qsub++
|
||||
}
|
||||
}
|
||||
}
|
||||
resp.ShareSubCount = qsub
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (b *Broker) DeliverMessage(packet *packets.PublishPacket) {
|
||||
for _, client := range b.rpcClient {
|
||||
|
||||
resp, err := client.QuerySubscribe(context.Background(), &pb.QuerySubscribeRequest{Topic: packet.TopicName, Qos: int32(packet.Qos)})
|
||||
if err != nil {
|
||||
log.Error("rpc request error:", zap.Error(err))
|
||||
continue
|
||||
}
|
||||
|
||||
if resp.RetCode == 0 {
|
||||
client.DeliverMessage(context.Background(), &pb.DeliverMessageRequest{Topic: packet.TopicName, Payload: packet.Payload})
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
func (b *Broker) QueryConnect(clientID string) {
|
||||
@@ -163,18 +126,20 @@ func (b *Broker) QueryConnect(clientID string) {
|
||||
}
|
||||
}
|
||||
|
||||
func (b *Broker) QueryShareSubscribe(topic string, qos byte) map[string]int32 {
|
||||
result := make(map[string]int32)
|
||||
type remoteSubInfo struct {
|
||||
subCount int
|
||||
shareSubCount int
|
||||
}
|
||||
|
||||
func (b *Broker) QuerySubscribe(topic string, qos byte) map[string]remoteSubInfo {
|
||||
result := make(map[string]remoteSubInfo)
|
||||
for id, client := range b.rpcClient {
|
||||
resp, err := client.QueryShareSubscribe(context.Background(), &pb.QueryShareSubscribeRequest{Topic: topic, Qos: int32(qos)})
|
||||
resp, err := client.QuerySubscribe(context.Background(), &pb.QuerySubscribeRequest{Topic: topic, Qos: int32(qos)})
|
||||
if err != nil {
|
||||
log.Error("rpc request error:", zap.Error(err))
|
||||
continue
|
||||
}
|
||||
if resp.ShareSubCount > 0 {
|
||||
result[id] = resp.ShareSubCount
|
||||
}
|
||||
|
||||
result[id] = remoteSubInfo{subCount: int(resp.SubCount), shareSubCount: int(resp.ShareSubCount)}
|
||||
}
|
||||
return result
|
||||
}
|
||||
|
||||
@@ -17,7 +17,6 @@ Logging Options:
|
||||
|
||||
Cluster Options:
|
||||
-r, --router <rurl> Router who maintenance cluster info
|
||||
-rpc, <rpc-port> Cluster listen port for others
|
||||
|
||||
Common Options:
|
||||
-h, --help Show this message
|
||||
|
||||
Reference in New Issue
Block a user