mirror of
https://github.com/fhmq/hmq.git
synced 2026-08-29 22:13:11 +00:00
muit channel
This commit is contained in:
+42
-19
@@ -27,6 +27,11 @@ var (
|
|||||||
brokerLog = logger.Get().Named("Broker")
|
brokerLog = logger.Get().Named("Broker")
|
||||||
)
|
)
|
||||||
|
|
||||||
|
const (
|
||||||
|
MessagePoolNum = 1024
|
||||||
|
MessageNum = 1024
|
||||||
|
)
|
||||||
|
|
||||||
type Message struct {
|
type Message struct {
|
||||||
client *client
|
client *client
|
||||||
packet packets.ControlPacket
|
packet packets.ControlPacket
|
||||||
@@ -45,12 +50,21 @@ type Broker struct {
|
|||||||
remotes sync.Map
|
remotes sync.Map
|
||||||
nodes map[string]interface{}
|
nodes map[string]interface{}
|
||||||
clusterChannel chan *Message
|
clusterChannel chan *Message
|
||||||
clientChannel chan *Message
|
messagePool []chan *Message
|
||||||
sl *Sublist
|
sl *Sublist
|
||||||
rl *RetainList
|
rl *RetainList
|
||||||
queues map[string]int
|
queues map[string]int
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func newMessagePool() []chan *Message {
|
||||||
|
mp := make([]chan *Message, 0)
|
||||||
|
for i := 0; i < MessagePoolNum; i++ {
|
||||||
|
tempCh := make(chan *Message, MessageNum)
|
||||||
|
mp = append(mp, tempCh)
|
||||||
|
}
|
||||||
|
return mp
|
||||||
|
}
|
||||||
|
|
||||||
func NewBroker(config *Config) (*Broker, error) {
|
func NewBroker(config *Config) (*Broker, error) {
|
||||||
b := &Broker{
|
b := &Broker{
|
||||||
id: GenUniqueId(),
|
id: GenUniqueId(),
|
||||||
@@ -61,7 +75,7 @@ func NewBroker(config *Config) (*Broker, error) {
|
|||||||
nodes: make(map[string]interface{}),
|
nodes: make(map[string]interface{}),
|
||||||
queues: make(map[string]int),
|
queues: make(map[string]int),
|
||||||
clusterChannel: make(chan *Message),
|
clusterChannel: make(chan *Message),
|
||||||
clientChannel: make(chan *Message),
|
messagePool: newMessagePool(),
|
||||||
}
|
}
|
||||||
if b.config.TlsPort != "" {
|
if b.config.TlsPort != "" {
|
||||||
tlsconfig, err := NewTLSConfig(b.config.TlsInfo)
|
tlsconfig, err := NewTLSConfig(b.config.TlsInfo)
|
||||||
@@ -84,15 +98,19 @@ func NewBroker(config *Config) (*Broker, error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (b *Broker) StartDispatcher() {
|
func (b *Broker) StartDispatcher() {
|
||||||
for {
|
for i := 0; i < MessagePoolNum; i++ {
|
||||||
msg, ok := <-b.clientChannel
|
go func(idx int) {
|
||||||
if !ok {
|
for {
|
||||||
brokerLog.Error("read message from client channel error")
|
msg, ok := <-b.messagePool[idx]
|
||||||
return
|
if !ok {
|
||||||
}
|
brokerLog.Error("read message from client channel error")
|
||||||
b.wpool.Submit(func() {
|
return
|
||||||
ProcessMessage(msg)
|
}
|
||||||
})
|
b.wpool.Submit(func() {
|
||||||
|
ProcessMessage(msg)
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}(i)
|
||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
@@ -167,9 +185,10 @@ func (b *Broker) StartWebsocketListening() {
|
|||||||
|
|
||||||
func (b *Broker) wsHandler(ws *websocket.Conn) {
|
func (b *Broker) wsHandler(ws *websocket.Conn) {
|
||||||
// io.Copy(ws, ws)
|
// io.Copy(ws, ws)
|
||||||
atomic.AddUint64(&b.cid, 1)
|
|
||||||
ws.PayloadType = websocket.BinaryFrame
|
ws.PayloadType = websocket.BinaryFrame
|
||||||
b.handleConnection(CLIENT, ws, b.cid)
|
|
||||||
|
idx := atomic.AddUint64(&b.cid, 1)
|
||||||
|
b.handleConnection(CLIENT, ws, idx)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (b *Broker) StartClientListening(Tls bool) {
|
func (b *Broker) StartClientListening(Tls bool) {
|
||||||
@@ -207,8 +226,8 @@ func (b *Broker) StartClientListening(Tls bool) {
|
|||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
tmpDelay = ACCEPT_MIN_SLEEP
|
tmpDelay = ACCEPT_MIN_SLEEP
|
||||||
atomic.AddUint64(&b.cid, 1)
|
idx := atomic.AddUint64(&b.cid, 1)
|
||||||
go b.handleConnection(CLIENT, conn, b.cid)
|
go b.handleConnection(CLIENT, conn, idx)
|
||||||
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -252,7 +271,6 @@ func (b *Broker) StartClusterListening() {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
var idx uint64 = 0
|
|
||||||
tmpDelay := 10 * ACCEPT_MIN_SLEEP
|
tmpDelay := 10 * ACCEPT_MIN_SLEEP
|
||||||
for {
|
for {
|
||||||
conn, err := l.Accept()
|
conn, err := l.Accept()
|
||||||
@@ -271,7 +289,7 @@ func (b *Broker) StartClusterListening() {
|
|||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
tmpDelay = ACCEPT_MIN_SLEEP
|
tmpDelay = ACCEPT_MIN_SLEEP
|
||||||
|
idx := atomic.AddUint64(&b.cid, 1)
|
||||||
go b.handleConnection(ROUTER, conn, idx)
|
go b.handleConnection(ROUTER, conn, idx)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -356,7 +374,9 @@ func (b *Broker) handleConnection(typ int, conn net.Conn, idx uint64) {
|
|||||||
b.routes.Store(cid, c)
|
b.routes.Store(cid, c)
|
||||||
}
|
}
|
||||||
|
|
||||||
c.readLoop(b.clientChannel)
|
mpool := b.messagePool[idx%MessagePoolNum]
|
||||||
|
|
||||||
|
c.readLoop(mpool)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (b *Broker) ConnectToDiscovery() {
|
func (b *Broker) ConnectToDiscovery() {
|
||||||
@@ -479,7 +499,10 @@ func (b *Broker) connectRouter(id, addr string) {
|
|||||||
|
|
||||||
c.SendConnect()
|
c.SendConnect()
|
||||||
|
|
||||||
go c.readLoop(b.clientChannel)
|
idx := atomic.AddUint64(&b.cid, 1)
|
||||||
|
mpool := b.messagePool[idx%MessagePoolNum]
|
||||||
|
|
||||||
|
go c.readLoop(mpool)
|
||||||
go c.StartPing()
|
go c.StartPing()
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user