11# # Native Nim MQTT client library and binaries
22# #
3- # # zevv (https://github.com/zevv) & ThomasTJdev (https://github.com/ThomasTJdev)
3+ # # zevv (https://github.com/zevv) & ThomasTJdev (https://github.com/ThomasTJdev) & python36 (https://github.com/python36)
44
55import
66 strutils,
4141 inWork: bool
4242 hasNewWorks: bool
4343 keepAlive: uint16
44+ maxInflightMessages: int
4445 willFlag: bool
4546 willQoS: uint8
4647 willRetain: bool
@@ -337,10 +338,17 @@ proc nextMsgId(ctx: MqttCtx): MsgId =
337338 inc ctx.msgIdSeq
338339 return ctx.msgIdSeq
339340
341+ proc hasInflightSlots (ctx: MqttCtx ): bool =
342+ var cnt = 0
343+ for w in ctx.workQueue.values ():
344+ if w.qos in {1 , 2 } and (w.typ != Publish or w.state == WorkSent ):
345+ inc cnt
346+ if cnt == ctx.maxInflightMessages:
347+ return false
348+ true
340349
341350proc sendDisconnect (ctx: MqttCtx ): Future [bool ] {.async .}
342351
343-
344352proc close (ctx: MqttCtx , reason: string ) {.async .} =
345353 if ctx.state in {Connecting , Connected }:
346354 ctx.state = Disconnecting
@@ -350,7 +358,6 @@ proc close(ctx: MqttCtx, reason: string) {.async.} =
350358 ctx.s.close ()
351359 ctx.state = Disconnected
352360
353-
354361proc send (ctx: MqttCtx , pkt: Pkt ): Future [bool ] {.async .} =
355362 # # Send the packet
356363 if ctx.state notin {Connecting , Connected , Disconnecting }:
@@ -378,7 +385,6 @@ proc send(ctx: MqttCtx, pkt: Pkt): Future[bool] {.async.} =
378385
379386 return true
380387
381-
382388proc recv (ctx: MqttCtx ): Future [Pkt ] {.async .} =
383389 # # Receive and parse the packet
384390 if ctx.state notin {Connecting ,Connected }:
@@ -442,7 +448,6 @@ proc recv(ctx: MqttCtx): Future[Pkt] {.async.} =
442448 ctx.dmp " rx> " & $ pkt
443449 return pkt
444450
445-
446451proc sendConnect (ctx: MqttCtx ): Future [bool ] =
447452 var flags: uint8
448453 flags = flags or CleanSession .uint8
@@ -479,7 +484,6 @@ proc sendConnect(ctx: MqttCtx): Future[bool] =
479484 ctx.state = Connecting
480485 result = ctx.send (pkt)
481486
482-
483487proc sendDisconnect (ctx: MqttCtx ): Future [bool ] =
484488 let pkt = newPkt (Disconnect , 0 )
485489 result = ctx.send (pkt)
@@ -616,17 +620,26 @@ proc work(ctx: MqttCtx) {.async.} =
616620 continue
617621
618622 if work.wk == PubWork and work.state == WorkNew :
619- if work.typ == Publish and work.qos == 0 :
620- if await ctx.sendWork (work): ctx.workQueue.del msgId
623+ if work.typ == Publish :
624+ if work.qos == 0 :
625+ if await ctx.sendWork (work):
626+ ctx.workQueue.del msgId
627+
628+ elif hasInflightSlots (ctx):
629+ if await ctx.sendWork (work):
630+ work.state = WorkSent
621631
622632 elif work.typ == PubAck and work.qos == 1 :
623- if await ctx.sendWork (work): ctx.workQueue.del msgId
633+ if await ctx.sendWork (work):
634+ ctx.workQueue.del msgId
624635
625636 elif work.typ == PubComp and work.qos == 2 :
626- if await ctx.sendWork (work): ctx.workQueue.del msgId
637+ if await ctx.sendWork (work):
638+ ctx.workQueue.del msgId
627639
628640 else :
629- if await ctx.sendWork (work): work.state = WorkSent
641+ if await ctx.sendWork (work):
642+ work.state = WorkSent
630643
631644 # when not defined(broker):
632645 elif work.wk == SubWork and work.state == WorkNew :
@@ -884,6 +897,7 @@ proc onPubAck(ctx: MqttCtx, pkt: Pkt) {.async.} =
884897 assert ctx.workQueue[msgId].state == WorkSent
885898 assert ctx.workQueue[msgId].qos == 1
886899 ctx.workQueue.del msgId
900+ await ctx.work ()
887901
888902proc onPubRec (ctx: MqttCtx , pkt: Pkt ) {.async .} =
889903 let (msgId, _) = pkt.getu16 (0 )
@@ -910,6 +924,7 @@ proc onPubComp(ctx: MqttCtx, pkt: Pkt) {.async.} =
910924 assert ctx.workQueue[msgId].state == WorkSent
911925 assert ctx.workQueue[msgId].qos == 2
912926 ctx.workQueue.del msgId
927+ await ctx.work ()
913928
914929# when defined(broker):
915930proc onSubscribe (ctx: MqttCtx , pkt: Pkt ) {.async .} =
@@ -1069,7 +1084,7 @@ proc runPing(ctx: MqttCtx) {.async.} =
10691084 await ctx.work ()
10701085
10711086proc connectBroker (ctx: MqttCtx ) {.async .} =
1072- # # Connect to the broker
1087+ # # Connect to the broker.
10731088 if ctx.keepAlive == 0 :
10741089 ctx.keepAlive = 60
10751090
@@ -1094,7 +1109,7 @@ proc connectBroker(ctx: MqttCtx) {.async.} =
10941109
10951110
10961111proc runConnect (ctx: MqttCtx ) {.async .} =
1097- # # Auto-connect and reconnect to broker
1112+ # # Auto-connect and reconnect to broker.
10981113
10991114 while true :
11001115 if ctx.state == Disabled :
@@ -1127,40 +1142,44 @@ proc runConnect(ctx: MqttCtx) {.async.} =
11271142#
11281143
11291144proc newMqttCtx * (clientId: string ): MqttCtx =
1130- # # Initiate a new MQTT client
1131- MqttCtx (clientId: clientId, state: Disconnected )
1145+ # # Initiate a new MQTT client.
1146+ MqttCtx (clientId: clientId, state: Disconnected , maxInflightMessages: 20 )
11321147
1133- proc set_ping_interval * (ctx: MqttCtx , txInterval: int = 60 ) =
1148+ proc setPingInterval * (ctx: MqttCtx , txInterval: int = 60 ) =
11341149 # # Set the clients ping interval in seconds. Default is 60 seconds.
11351150 if txInterval > 0 and txInterval < 65535 :
11361151 ctx.keepAlive = txInterval.uint16
11371152
1138- proc set_host * (ctx: MqttCtx , host: string , port: int = 1883 , sslOn= false ) =
1139- # # Set the MQTT host
1153+ proc setHost * (ctx: MqttCtx , host: string , port: int = 1883 , sslOn= false ) =
1154+ # # Set the MQTT host.
11401155 ctx.host = host
11411156 ctx.port = Port (port)
11421157 ctx.sslOn = sslOn
11431158
1144- proc set_ssl_certificates * (ctx: MqttCtx , sslCert: string , sslKey: string ) =
1159+ proc setSSLCertificates * (ctx: MqttCtx , sslCert: string , sslKey: string ) =
11451160 # Sets the SSL Certificate and Key to use when connecting to the remote broker
11461161 # for mutal TLS authentication
11471162 ctx.sslCert = sslCert
11481163 ctx.sslKey = sslKey
11491164
1150- proc set_auth * (ctx: MqttCtx , username: string , password: string ) =
1165+ proc setAuth * (ctx: MqttCtx , username: string , password: string ) =
11511166 # # Set the authentication for the host.
11521167 ctx.username = username
11531168 ctx.password = password
11541169
1155- proc set_will * (ctx: MqttCtx , topic, msg: string , qos= 0 , retain= false ) =
1170+ proc setWill * (ctx: MqttCtx , topic, msg: string , qos= 0 , retain= false ) =
11561171 # # Set the clients will.
11571172 ctx.willFlag = true
11581173 ctx.willTopic = topic
11591174 ctx.willMsg = msg
11601175 ctx.willQoS = qos.uint8
11611176 ctx.willRetain = retain
11621177
1163- proc set_verbosity * (ctx: MqttCtx , verbosity: int ) =
1178+ proc setMaxInflightMessages * (ctx: MqttCtx , maxInflightMessages: int ) =
1179+ # # Sets the maximum number of unacknowledged MQTT messages (QoS 1 and QoS 2). Default = 20.
1180+ ctx.maxInflightMessages = maxInflightMessages
1181+
1182+ proc setVerbosity * (ctx: MqttCtx , verbosity: int ) =
11641183 # # Set the verbosity.
11651184 ctx.verbosity = verbosity
11661185
0 commit comments