mirror of
https://github.com/tinygo-org/drivers.git
synced 2026-08-20 06:28:59 +00:00
mqtt: use buffered channels for incoming messages to handle bursts
Signed-off-by: Ron Evans <ron@hybridgroup.com>
This commit is contained in:
committed by
Daniel Esteban
parent
6380ad5ed5
commit
12ac4c2c06
+3
-3
@@ -77,9 +77,9 @@ func (c *mqttclient) Connect() Token {
|
|||||||
}
|
}
|
||||||
|
|
||||||
c.mid = 1
|
c.mid = 1
|
||||||
c.inbound = make(chan packets.ControlPacket)
|
c.inbound = make(chan packets.ControlPacket, 10)
|
||||||
c.stop = make(chan struct{})
|
c.stop = make(chan struct{})
|
||||||
c.incomingPubChan = make(chan *packets.PublishPacket)
|
c.incomingPubChan = make(chan *packets.PublishPacket, 10)
|
||||||
c.msgRouter.matchAndDispatch(c.incomingPubChan, c.opts.Order, c)
|
c.msgRouter.matchAndDispatch(c.incomingPubChan, c.opts.Order, c)
|
||||||
|
|
||||||
// send the MQTT connect message
|
// send the MQTT connect message
|
||||||
@@ -98,7 +98,7 @@ func (c *mqttclient) Connect() Token {
|
|||||||
connectPkt.ClientIdentifier = c.opts.ClientID
|
connectPkt.ClientIdentifier = c.opts.ClientID
|
||||||
connectPkt.ProtocolVersion = byte(c.opts.ProtocolVersion)
|
connectPkt.ProtocolVersion = byte(c.opts.ProtocolVersion)
|
||||||
connectPkt.ProtocolName = "MQTT"
|
connectPkt.ProtocolName = "MQTT"
|
||||||
connectPkt.Keepalive = 30
|
connectPkt.Keepalive = 60
|
||||||
|
|
||||||
err = connectPkt.Write(c.conn)
|
err = connectPkt.Write(c.conn)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
Reference in New Issue
Block a user