PubSubClient3 v3.3.2
Located at <a href='https://github.com/hmueller01/pubsubclient3'>GitHub</a>
 
Loading...
Searching...
No Matches
PubSubClient.cpp
Go to the documentation of this file.
1
10
11#include "PubSubClient.h"
12
24#if defined(ARDUINO_ARCH_ESP32) || defined(ESP32)
25#define PUBSUB_WAIT_YIELD() vTaskDelay(1)
26#else
27#define PUBSUB_WAIT_YIELD() yield()
28#endif
29
37#define CHECK_STRING_LENGTH(l, s) \
38 if ((!s) || (l + 2 + strnlen(s, _bufferSize) > _bufferSize)) { \
39 _client->stop(); \
40 return false; \
41 }
42
48
50 setClient(client);
51}
52
53PubSubClient::PubSubClient(IPAddress addr, uint16_t port, Client& client) : PubSubClient() {
54 setServer(addr, port);
55 setClient(client);
56}
57
58PubSubClient::PubSubClient(IPAddress addr, uint16_t port, Client& client, Stream& stream) : PubSubClient() {
59 setServer(addr, port);
60 setClient(client);
61 setStream(stream);
62}
63
64PubSubClient::PubSubClient(IPAddress addr, uint16_t port, MQTT_CALLBACK_SIGNATURE, Client& client) : PubSubClient() {
65 setServer(addr, port);
66 setCallback(callback);
67 setClient(client);
68}
69
70PubSubClient::PubSubClient(IPAddress addr, uint16_t port, MQTT_CALLBACK_SIGNATURE, Client& client, Stream& stream) : PubSubClient() {
71 setServer(addr, port);
72 setCallback(callback);
73 setClient(client);
74 setStream(stream);
75}
76
77PubSubClient::PubSubClient(uint8_t* ip, uint16_t port, Client& client) : PubSubClient() {
78 setServer(ip, port);
79 setClient(client);
80}
81
82PubSubClient::PubSubClient(uint8_t* ip, uint16_t port, Client& client, Stream& stream) : PubSubClient() {
83 setServer(ip, port);
84 setClient(client);
85 setStream(stream);
86}
87
88PubSubClient::PubSubClient(uint8_t* ip, uint16_t port, MQTT_CALLBACK_SIGNATURE, Client& client) : PubSubClient() {
89 setServer(ip, port);
90 setCallback(callback);
91 setClient(client);
92}
93
94PubSubClient::PubSubClient(uint8_t* ip, uint16_t port, MQTT_CALLBACK_SIGNATURE, Client& client, Stream& stream) : PubSubClient() {
95 setServer(ip, port);
96 setCallback(callback);
97 setClient(client);
98 setStream(stream);
99}
100
101PubSubClient::PubSubClient(const char* domain, uint16_t port, Client& client) : PubSubClient() {
102 setServer(domain, port);
103 setClient(client);
104}
105
106PubSubClient::PubSubClient(const char* domain, uint16_t port, Client& client, Stream& stream) : PubSubClient() {
107 setServer(domain, port);
108 setClient(client);
109 setStream(stream);
110}
111
112PubSubClient::PubSubClient(const char* domain, uint16_t port, MQTT_CALLBACK_SIGNATURE, Client& client) : PubSubClient() {
113 setServer(domain, port);
114 setCallback(callback);
115 setClient(client);
116}
117
118PubSubClient::PubSubClient(const char* domain, uint16_t port, MQTT_CALLBACK_SIGNATURE, Client& client, Stream& stream) : PubSubClient() {
119 setServer(domain, port);
120 setCallback(callback);
121 setClient(client);
122 setStream(stream);
123}
124
126 free(_domain);
127 free(_buffer);
128}
129
130bool PubSubClient::connect(const char* id, const char* user, const char* pass, const char* willTopic, uint8_t willQos, bool willRetain,
131 const char* willMessage, bool cleanSession) {
132 if (!_client) return false; // do not crash if client not set
133 if (!_buffer) return false; // do not crash if buffer allocation failed at construction
134 if (!connected()) {
135 int result = 0;
136
137 if (_client->connected()) {
138 result = 1;
139 } else if (_port != 0) {
140 if (_domain) {
141 result = _client->connect(_domain, _port);
142 } else {
143 result = _client->connect(_ip, _port);
144 }
145 }
146
147 if (result == 1) {
148 _nextMsgId = 1; // init msgId (packet identifier)
149
150#if MQTT_VERSION == MQTT_VERSION_3_1
151 const uint8_t protocol[9] = {0x00, 0x06, 'M', 'Q', 'I', 's', 'd', 'p', MQTT_VERSION};
152#elif MQTT_VERSION == MQTT_VERSION_3_1_1
153 const uint8_t protocol[7] = {0x00, 0x04, 'M', 'Q', 'T', 'T', MQTT_VERSION};
154#endif
155 // Leave room in the _buffer for header and variable length field
156 memcpy(_buffer + MQTT_MAX_HEADER_SIZE, protocol, sizeof(protocol));
157
158 size_t length = MQTT_MAX_HEADER_SIZE + sizeof(protocol);
159 uint8_t flags = 0x00;
160 if (willTopic) {
161 flags = (0x01 << 2) | (willQos << 3) | (willRetain << 5); // set will flag bit 2, will QoS and will retain bit 5
162 }
163 if (cleanSession) {
164 flags = flags | (0x01 << 1); // set clean session bit 1
165 }
166 if (user) {
167 flags = flags | (0x01 << 7); // set user name flag bit 7
168 if (pass) {
169 flags = flags | (0x01 << 6); // set password flag bit 6
170 }
171 }
172 const uint16_t keepAlive = _keepAliveMillis / 1000;
173 _buffer[length++] = flags;
174 _buffer[length++] = keepAlive >> 8;
175 _buffer[length++] = keepAlive & 0xFF;
176
177 CHECK_STRING_LENGTH(length, id)
178 length = writeString(id, length);
179 if (willTopic) {
180 CHECK_STRING_LENGTH(length, willTopic)
181 length = writeString(willTopic, length);
182 CHECK_STRING_LENGTH(length, willMessage)
183 length = writeString(willMessage, length);
184 }
185
186 if (user) {
187 CHECK_STRING_LENGTH(length, user)
188 length = writeString(user, length);
189 if (pass) {
190 CHECK_STRING_LENGTH(length, pass)
191 length = writeString(pass, length);
192 }
193 }
194
195 if (!writeControlPacket(MQTTCONNECT, length - MQTT_MAX_HEADER_SIZE)) {
196 _state = MQTT_CONNECT_FAILED;
197 _client->stop();
198 return false;
199 }
200 _lastInActivity = _lastOutActivity = millis();
201 _pingOutstanding = false;
202
203 while (!_client->available()) {
205 unsigned long t = millis();
206 if (t - _lastInActivity >= _socketTimeoutMillis) {
207 DEBUG_PSC_PRINTF("connect aborting due to timeout\n");
209 _client->stop();
210 return false;
211 }
212 }
213 uint8_t hdrLen;
214 size_t len = readPacket(&hdrLen);
215
216 if (len == 4) {
217 if (_buffer[3] == 0) {
218 _lastInActivity = millis();
219 _state = MQTT_CONNECTED;
220 return true;
221 } else {
222 _state = _buffer[3];
223 }
224 }
225 DEBUG_PSC_PRINTF("connect aborting due to protocol error\n");
226 _client->stop();
227 } else {
228 _state = MQTT_CONNECT_FAILED;
229 }
230 return false;
231 }
232 return true;
233}
234
236 if (!_client) return false;
237 if (!_buffer) return false; // we can't be connected if we don't have a buffer to read into
238
239 if (_client->connected()) {
240 return (_state == MQTT_CONNECTED);
241 } else if (_state == MQTT_CONNECTED) {
242 DEBUG_PSC_PRINTF("lost connection (client may have more details)\n");
243 _state = MQTT_CONNECTION_LOST;
244 _client->stop();
245 _pingOutstanding = false;
246 }
247 return false;
248}
249
251 DEBUG_PSC_PRINTF("disconnect called\n");
252 _state = MQTT_DISCONNECTED;
253 if (_client && _buffer) { // guard against null buffer if allocation failed at construction
254 _buffer[0] = MQTTDISCONNECT;
255 _buffer[1] = 0;
256 _client->write(_buffer, 2);
257 _client->flush();
258 _client->stop();
259 _lastInActivity = _lastOutActivity = millis();
260 }
261 _pingOutstanding = false;
262}
263
270bool PubSubClient::readByte(uint8_t* result) {
271 if (!_client) return false; // do not crash if client not set
272
273 unsigned long previousMillis = millis();
274 while (!_client->available()) {
276 unsigned long currentMillis = millis();
277 if (currentMillis - previousMillis >= _socketTimeoutMillis) {
278 return false;
279 }
280 }
281 int rc = _client->read();
282 if (rc < 0) {
283 return false;
284 }
285 *result = (uint8_t)rc;
286 return true;
287}
288
295bool PubSubClient::readByte(uint8_t* result, size_t* pos) {
296 uint8_t* write_address = &(result[*pos]);
297 if (readByte(write_address)) {
298 (*pos)++;
299 return true;
300 }
301 return false;
302}
303
310size_t PubSubClient::readPacket(uint8_t* hdrLen) {
311 size_t len = 0;
312 if (!readByte(_buffer, &len)) return 0;
313 bool isPublish = (_buffer[0] & 0xF0) == MQTTPUBLISH;
314 uint32_t multiplier = 1;
315 size_t length = 0;
316 uint8_t digit = 0;
317 uint16_t skip = 0;
318 uint8_t start = 0;
319
320 do {
321 if (len == MQTT_MAX_HEADER_SIZE) {
322 // Invalid remaining length encoding - kill the connection
323 DEBUG_PSC_PRINTF("readPacket detected packet of invalid length\n");
324 _state = MQTT_DISCONNECTED;
325 _client->stop();
326 return 0;
327 }
328 if (!readByte(&digit)) return 0;
329 _buffer[len++] = digit;
330 length += (digit & 0x7F) * multiplier; // length is coded in the lower 7 bits
331 multiplier <<= 7; // multiplier *= 128
332 } while ((digit & 0x80) != 0); // do while 8th continuation bit is set
333 *hdrLen = (uint8_t)(len - 1);
334
335 DEBUG_PSC_PRINTF("readPacket received packet of length %zu (isPublish = %u)\n", length, isPublish);
336
337 if (isPublish) {
338 // Read in topic length to calculate bytes to skip over for Stream writing
339 if (!readByte(_buffer, &len)) return 0;
340 if (!readByte(_buffer, &len)) return 0;
341 skip = (_buffer[*hdrLen + 1] << 8) + _buffer[*hdrLen + 2];
342 start = 2;
343 if (MQTT_HDR_GET_QOS(_buffer[0]) > MQTT_QOS0) {
344 // skip msgId (packet identifier) for QoS 1 and 2 messages
345 skip += 2;
346 }
347 }
348 size_t idx = len;
349
350 for (size_t i = start; i < length; i++) {
351 if (!readByte(&digit)) return 0;
352 if (_stream) {
353 if (isPublish && (idx - *hdrLen - 2 > skip)) {
354 _stream->write(digit);
355 }
356 }
357
358 if (len < _bufferSize) {
359 _buffer[len++] = digit;
360 }
361 idx++;
362 }
363
364 if (!_stream && (idx > _bufferSize)) {
365 DEBUG_PSC_PRINTF("readPacket ignoring packet of size %zu exceeding buffer of size %zu\n", length, _bufferSize);
366 len = 0; // This will cause the packet to be ignored.
367 }
368 return len;
369}
370
378bool PubSubClient::handlePacket(uint8_t hdrLen, size_t length) {
379 uint8_t type = _buffer[0] & 0xF0;
380 DEBUG_PSC_PRINTF("handlePacket(): received message of type %u\n", type);
381 if (length > _bufferSize) {
382 // This should never happen as readPacket() prevents buffer overflow, but we check again here to be sure and prevent any buffer overflows.
383 DEBUG_PSC_PRINTF("handlePacket(): packet length %zu exceeds buffer size %zu\n", length, _bufferSize);
384 return false;
385 }
386 switch (type) {
387 case MQTTPUBLISH:
388 if (callback) {
389 // MQTT Publish packet: See section 3.3 MQTT v3.1.1 protocol specification:
390 // - Header: 1 byte
391 // - Remaining header length: hdrLen bytes, multibyte field (1 .. MQTT_MAX_HEADER_SIZE - 1)
392 // - Topic length: 2 bytes (starts at _buffer[hdrLen + 1])
393 // - Topic: topicLen bytes (starts at _buffer[hdrLen + 3])
394 // - Packet Identifier (msgId): 0 bytes for QoS 0, 2 bytes for QoS 1 and 2 (starts at _buffer[hdrLen + 3 + topicLen])
395 // - Payload (for QoS = 0): length - (hdrLen + 3 + topicLen) bytes (starts at _buffer[hdrLen + 3 + topicLen])
396 // - Payload (for QoS > 0): length - (hdrLen + 5 + topicLen) bytes (starts at _buffer[hdrLen + 5 + topicLen])
397 // To get a null terminated 'C' topic string we move the topic 1 byte to the front (overwriting the LSB of the topic lenght)
398 // Guard 1: ensure topic length bytes are readable
399 if (length < hdrLen + 3ul) {
400 DEBUG_PSC_PRINTF("handlePacket(): Packet too short to contain topic length field\n");
401 return false;
402 }
403 const uint16_t topicLen = (_buffer[hdrLen + 1] << 8) + _buffer[hdrLen + 2]; // topic length in bytes
404 char* topic = (char*)(_buffer + hdrLen + 3 - 1); // set the topic in the LSB of the topic lenght, as we move it there later
405 // Use size_t here: hdrLen + 3 + topicLen can reach 65539, which
406 // wraps a uint16_t (e.g. topicLen = 0xFFFF wraps to 3), silently
407 // defeating Guard 2 below and letting an attacker-chosen topicLen
408 // reach the memmove/write further down.
409 const size_t payloadOffset = (size_t)hdrLen + 3 + topicLen; // payload starts after header and topic (if there is no packet identifier)
410 const size_t payloadLen = length - payloadOffset;
411 uint8_t* const payload = _buffer + payloadOffset;
412
413 // Guard 2: ensure topic fits in buffer
414 if (length < payloadOffset) {
415 ERROR_PSC_PRINTF_P("handlePacket(): Suspicious topicLen (%u) points outside of received buffer length (%zu)\n", topicLen, length);
416 return false;
417 }
418 memmove(topic, topic + 1, topicLen); // move topic inside buffer 1 byte to front to get space for null termination
419 topic[topicLen] = '\0'; // end the topic as a 'C' string with null termination
420
421 if (MQTT_HDR_GET_QOS(_buffer[0]) == MQTT_QOS0) {
422 // No msgId for QOS == 0
423 callback(topic, payload, payloadLen);
424 } else {
425 // For QOS 1 and 2 we have a msgId (packet identifier) after the topic at the current payloadOffset
426 if (payloadLen < 2) { // payload must be >= 2, as we have the msgId before the actual payload
427 DEBUG_PSC_PRINTF("handlePacket(): Missing msgId in QoS 1/2 message\n");
428 return false;
429 }
430 const uint8_t publishQos = MQTT_HDR_GET_QOS(_buffer[0]); // save QoS before _buffer[0] is overwritten
431 const uint16_t msgId = (_buffer[payloadOffset] << 8) + _buffer[payloadOffset + 1];
432 callback(topic, payload + 2, payloadLen - 2); // remove the msgId from the callback payload
433
434 // QoS 1: respond with PUBACK
435 // QoS 2: respond with PUBREC (first step of the QoS 2 subscriber handshake)
436 _buffer[0] = (publishQos == MQTT_QOS1) ? MQTTPUBACK : MQTTPUBREC;
437 _buffer[1] = 2;
438 _buffer[2] = (uint8_t)(msgId >> 8);
439 _buffer[3] = (uint8_t)(msgId & 0xFF);
440 if (_client->write(_buffer, 4) == 4) {
441 _lastOutActivity = millis();
442 }
443 }
444 }
445 break;
446 case MQTTPUBACK:
447 // MQTT Publish Acknowledgment (QoS 1 publish received): See section 3.4 MQTT v3.1.1 protocol specification
448 if (length < 4) {
449 ERROR_PSC_PRINTF_P("handlePacket(): Received PUBACK packet with length %zu, expected at least 4 bytes\n", length);
450 return false;
451 }
452 // No futher action here, as resending is not supported.
453 break;
454 case MQTTPUBREC:
455 // MQTT Publish Received (QoS 2 publisher handshake, part 1): broker acknowledges our QoS 2 PUBLISH.
456 // See section 3.5 MQTT v3.1.1 protocol specification.
457 if (length < 4) {
458 ERROR_PSC_PRINTF_P("handlePacket(): Received PUBREC packet with length %zu, expected at least 4 bytes\n", length);
459 return false;
460 }
461 // MQTT Publish Release (QoS 2 publisher handshake, part 2): See section 3.6 MQTT v3.1.1 protocol specification
462 _buffer[0] = MQTTPUBREL | 2; // PUBREL fixed header: bit 1 must be set per spec
463 // bytes 1-3 of PUBREL are the same as of PUBREC (remaining length + msgId)
464 if (_client->write(_buffer, 4) == 4) {
465 _lastOutActivity = millis();
466 }
467 break;
468 case MQTTPUBREL:
469 // MQTT Publish Release (QoS 2 subscriber handshake, part 2): broker releases the message to us.
470 // See section 3.6 MQTT v3.1.1 protocol specification.
471 if (length < 4) {
472 ERROR_PSC_PRINTF_P("handlePacket(): Received PUBREL packet with length %zu, expected at least 4 bytes\n", length);
473 return false;
474 }
475 // MQTT Publish Complete (QoS 2 subscriber handshake, part 3): See section 3.7 MQTT v3.1.1 protocol specification
476 _buffer[0] = MQTTPUBCOMP;
477 // bytes 1-3 of PUBCOMP are the same as of PUBREL (remaining length + msgId)
478 if (_client->write(_buffer, 4) == 4) {
479 _lastOutActivity = millis();
480 }
481 break;
482 case MQTTPUBCOMP:
483 // MQTT Publish Complete (QoS 2 publisher handshake, part 3): broker confirms delivery of our QoS 2 PUBLISH.
484 // See section 3.7 MQTT v3.1.1 protocol specification.
485 if (length < 4) {
486 ERROR_PSC_PRINTF_P("handlePacket(): Received PUBCOMP packet with length %zu, expected at least 4 bytes\n", length);
487 return false;
488 }
489 // No futher action here, as resending is not supported.
490 break;
491 case MQTTPINGREQ:
492 // MQTT Ping Request: See section 3.12 MQTT v3.1.1 protocol specification
493 _buffer[0] = MQTTPINGRESP;
494 _buffer[1] = 0;
495 if (_client->write(_buffer, 2) == 2) {
496 _lastOutActivity = millis();
497 }
498 break;
499 case MQTTPINGRESP:
500 _pingOutstanding = false;
501 break;
502 default:
503 break;
504 }
505 return true;
506}
507
509 if (!connected()) return false;
510
511 bool ret = true;
512 const unsigned long t = millis();
513 if (_keepAliveMillis && ((t - _lastInActivity > _keepAliveMillis) || (t - _lastOutActivity > _keepAliveMillis))) {
514 if (_pingOutstanding) {
515 DEBUG_PSC_PRINTF("loop aborting due to timeout\n");
517 _client->stop();
518 _pingOutstanding = false;
519 return false;
520 } else if (_bufferWritePos > 0) {
521 // There is still data in the _buffer to be sent, so send it now instead of a ping
522 if (flushBuffer() == 0) {
524 _client->stop();
525 _pingOutstanding = false;
526 return false;
527 }
528 } else {
529 _buffer[0] = MQTTPINGREQ;
530 _buffer[1] = 0;
531 if (_client->write(_buffer, 2) == 2) {
532 _lastInActivity = _lastOutActivity = t;
533 _pingOutstanding = true;
534 }
535 }
536 }
537 if (_client->available()) {
538 uint8_t hdrLen;
539 size_t len = readPacket(&hdrLen);
540 if (len > 0) {
541 _lastInActivity = t;
542 ret = handlePacket(hdrLen, len);
543 if (!ret) {
544 _state = MQTT_DISCONNECTED;
545 _client->stop();
546 }
547 } else if (!connected()) {
548 // readPacket has closed the connection
549 return false;
550 }
551 }
552 return ret;
553}
554
555bool PubSubClient::publish(const char* topic, const uint8_t* payload, size_t plength, uint8_t qos, bool retained) {
556 if (beginPublish(topic, plength, qos, retained)) {
557 size_t rc = write(payload, plength);
558 return endPublish() && (rc == plength);
559 }
560 return false;
561}
562
563bool PubSubClient::publish(const __FlashStringHelper* topic, const uint8_t* payload, size_t plength, uint8_t qos, bool retained) {
564 if (beginPublish(topic, plength, qos, retained)) {
565 size_t rc = write(payload, plength);
566 return endPublish() && (rc == plength);
567 }
568 return false;
569}
570
571bool PubSubClient::publish_P(const char* topic, const uint8_t* payload, size_t plength, uint8_t qos, bool retained) {
572 if (beginPublish(topic, plength, qos, retained)) {
573 size_t rc = write_P(payload, plength);
574 return endPublish() && (rc == plength);
575 }
576 return false;
577}
578
579bool PubSubClient::publish_P(const __FlashStringHelper* topic, const uint8_t* payload, size_t plength, uint8_t qos, bool retained) {
580 if (beginPublish(topic, plength, qos, retained)) {
581 size_t rc = write_P(payload, plength);
582 return endPublish() && (rc == plength);
583 }
584 return false;
585}
586
598bool PubSubClient::beginPublishImpl(bool progmem, const char* topic, size_t plength, uint8_t qos, bool retained) {
599 if (!topic) return false;
600
601 // get topic length depending on storage (RAM vs PROGMEM)
602 size_t topicLen = progmem ? strlen_P(topic) : strlen(topic);
603 if (topicLen == 0) return false; // empty topic is not allowed
604
605 if (qos > MQTT_QOS2) { // only valid QoS supported
606 ERROR_PSC_PRINTF_P("beginPublish() called with invalid QoS %u\n", qos);
607 return false;
608 }
609
610 const size_t nextMsgLen = (qos > MQTT_QOS0) ? 2 : 0; // add 2 bytes for nextMsgId if QoS > 0
611 // check if the header, the topic (including 2 length bytes) and nextMsgId fit into the _buffer
612 if (connected() && (MQTT_MAX_HEADER_SIZE + topicLen + 2 + nextMsgLen <= _bufferSize)) {
613 // first write the topic at the end of the maximal variable header (MQTT_MAX_HEADER_SIZE) to the _buffer
614 topicLen = writeStringImpl(progmem, topic, MQTT_MAX_HEADER_SIZE) - MQTT_MAX_HEADER_SIZE;
615 if (qos > MQTT_QOS0) {
616 // if QoS 1 or 2, we need to send the nextMsgId (packet identifier) after topic
617 writeNextMsgId(MQTT_MAX_HEADER_SIZE + topicLen);
618 }
619 // we now know the length of the topic string (length + 2 bytes signalling the length) and can build the variable header information
620 const uint8_t header = MQTTPUBLISH | MQTT_QOS_GET_HDR(qos) | (retained ? MQTTRETAINED : 0);
621 uint8_t hdrLen = buildHeader(header, topicLen + nextMsgLen + plength);
622 if (hdrLen == 0) return false; // exit here in case of header generation failure
623 // as the header length is variable, it starts at MQTT_MAX_HEADER_SIZE - hdrLen (see buildHeader() documentation)
624 size_t rc = _client->write(_buffer + (MQTT_MAX_HEADER_SIZE - hdrLen), hdrLen + topicLen + nextMsgLen);
625 _lastOutActivity = millis();
626 return (rc == (hdrLen + topicLen + nextMsgLen));
627 }
628 return false;
629}
630
632 if (connected()) {
633 if (_bufferWritePos > 0) {
634 // still data in the _buffer to be sent
635 if (flushBuffer() == 0) return false;
636 }
637 return true;
638 }
639 return false;
640}
641
651uint8_t PubSubClient::buildHeader(uint8_t header, size_t length) {
652 uint8_t hdrBuf[MQTT_MAX_HEADER_SIZE - 1];
653 uint8_t hdrLen = 0;
654 uint8_t digit;
655 size_t len = length;
656 do {
657 digit = len & 0x7F; // digit = len % 128
658 len >>= 7; // len = len / 128
659 if (len > 0) {
660 digit |= 0x80;
661 }
662 hdrBuf[hdrLen++] = digit;
663 } while ((len > 0) && (hdrLen < MQTT_MAX_HEADER_SIZE - 1));
664
665 if (len > 0) {
666 ERROR_PSC_PRINTF_P("buildHeader: header=0x%02X, length too big %zu, left %zu\n", header, length, len);
667 return 0;
668 }
669
670 _buffer[MQTT_MAX_HEADER_SIZE - 1 - hdrLen] = header;
671 memcpy(_buffer + MQTT_MAX_HEADER_SIZE - hdrLen, hdrBuf, hdrLen);
672 return hdrLen + 1; // Full header size is variable length bit plus the 1-byte fixed header
673}
674
675size_t PubSubClient::write(uint8_t data) {
676 return appendBuffer(data);
677}
678
679size_t PubSubClient::write(const uint8_t* buf, size_t size) {
680 for (size_t i = 0; i < size; i++) {
681 if (appendBuffer(buf[i]) == 0) return i;
682 }
683 return size;
684}
685
686size_t PubSubClient::write_P(const uint8_t* buf, size_t size) {
687 for (size_t i = 0; i < size; i++) {
688 if (appendBuffer((uint8_t)pgm_read_byte_near(buf + i)) == 0) return i;
689 }
690 return size;
691}
692
701bool PubSubClient::writeControlPacket(uint8_t header, size_t length) {
702 uint8_t hdrLen = buildHeader(header, length);
703 if (hdrLen == 0) return false; // exit here in case of header generation failure
704
705 return writeBuffer(MQTT_MAX_HEADER_SIZE - hdrLen, hdrLen + length);
706}
707
715size_t PubSubClient::writeBuffer(size_t pos, size_t size) {
716 size_t rc = 0;
717 if (_client && (size > 0) && (pos + size <= _bufferSize)) {
718#ifdef MQTT_MAX_TRANSFER_SIZE
719 uint8_t* writeBuf = _buffer + pos;
720 size_t bytesRemaining = size;
721 bool result = true;
722 while ((bytesRemaining > 0) && result) {
723 size_t bytesToWrite = (bytesRemaining > MQTT_MAX_TRANSFER_SIZE) ? MQTT_MAX_TRANSFER_SIZE : bytesRemaining;
724 size_t bytesWritten = _client->write(writeBuf, bytesToWrite);
725 result = (bytesWritten == bytesToWrite);
726 bytesRemaining -= bytesWritten;
727 writeBuf += bytesWritten;
728 if (result) {
729 _lastOutActivity = millis();
730 }
731 yield();
732 }
733 rc = result ? size : 0; // if result is false indicate a write error
734#else
735 rc = _client->write(_buffer + pos, size);
736 if (rc == size) {
737 _lastOutActivity = millis();
738 } else {
739 rc = 0; // indicate a write error
740 }
741#endif
742 }
743 return rc;
744}
745
758size_t PubSubClient::writeStringImpl(bool progmem, const char* string, size_t pos) {
759 if (!string) return pos;
760
761 size_t sLen = progmem ? strlen_P(string) : strlen(string);
762 if ((pos + 2 + sLen <= _bufferSize) && (sLen <= 0xFFFF)) {
763 _buffer[pos++] = (uint8_t)(sLen >> 8);
764 _buffer[pos++] = (uint8_t)(sLen & 0xFF);
765 if (progmem) {
766 memcpy_P(_buffer + pos, string, sLen);
767 } else {
768 memcpy(_buffer + pos, string, sLen);
769 }
770 pos += sLen;
771 } else {
772 ERROR_PSC_PRINTF_P("writeStringImpl(): string (%zu) does not fit into buf (%zu)\n", pos + 2 + sLen, _bufferSize);
773 }
774 return pos;
775}
776
787inline size_t PubSubClient::writeString(const char* string, size_t pos) {
788 return writeStringImpl(false, string, pos);
789}
790
798size_t PubSubClient::writeNextMsgId(size_t pos) {
799 if ((pos + 2) <= _bufferSize) {
800 _nextMsgId = (++_nextMsgId == 0) ? 1 : _nextMsgId; // increment msgId (must not be 0, so start at 1)
801 _buffer[pos++] = (uint8_t)(_nextMsgId >> 8);
802 _buffer[pos++] = (uint8_t)(_nextMsgId & 0xFF);
803 } else {
804 ERROR_PSC_PRINTF_P("writeNextMsgId(): buffer overrun (%zu) \n", pos + 2);
805 }
806 return pos;
807}
808
815size_t PubSubClient::appendBuffer(uint8_t data) {
816 _buffer[_bufferWritePos++] = data;
817 if (_bufferWritePos >= _bufferSize) {
818 if (flushBuffer() == 0) return 0;
819 }
820 return 1;
821}
822
829size_t PubSubClient::flushBuffer() {
830 size_t rc = 0;
831 if (connected()) {
832 rc = writeBuffer(0, _bufferWritePos);
833 }
834 _bufferWritePos = 0;
835 return rc;
836}
837
846bool PubSubClient::subscribeImpl(bool progmem, const char* topic, uint8_t qos) {
847 if (!topic) return false;
848 if (qos > MQTT_QOS1) return false; // only QoS 0 and 1 supported
849
850 // get topic length depending on storage (RAM vs PROGMEM)
851 size_t topicLen = progmem ? strnlen_P(topic, _bufferSize) : strnlen(topic, _bufferSize);
852 if (_bufferSize < MQTT_MAX_HEADER_SIZE + 2 + 2 + topicLen + 1) {
853 // Too long: header + nextMsgId (2) + topic length bytes (2) + topicLen + QoS (1)
854 return false;
855 }
856 if (connected()) {
857 // Leave room in the _buffer for header and variable length field
858 size_t length = MQTT_MAX_HEADER_SIZE;
859 length = writeNextMsgId(length); // _buffer size is checked before
860 length = writeStringImpl(progmem, topic, length);
861 _buffer[length++] = qos;
862 return writeControlPacket(MQTTSUBSCRIBE | MQTT_QOS_GET_HDR(MQTT_QOS1), length - MQTT_MAX_HEADER_SIZE);
863 }
864 return false;
865}
866
874bool PubSubClient::unsubscribeImpl(bool progmem, const char* topic) {
875 if (!topic) return false;
876
877 // get topic length depending on storage (RAM vs PROGMEM)
878 size_t topicLen = progmem ? strnlen_P(topic, _bufferSize) : strnlen(topic, _bufferSize);
879 if (_bufferSize < MQTT_MAX_HEADER_SIZE + 2 + 2 + topicLen) {
880 // Too long: header + nextMsgId (2) + topic length bytes (2) + topicLen
881 return false;
882 }
883 if (connected()) {
884 size_t length = MQTT_MAX_HEADER_SIZE;
885 length = writeNextMsgId(length); // _buffer size is checked before
886 length = writeStringImpl(progmem, topic, length);
887 return writeControlPacket(MQTTUNSUBSCRIBE | MQTT_QOS_GET_HDR(MQTT_QOS1), length - MQTT_MAX_HEADER_SIZE);
888 }
889 return false;
890}
891
892PubSubClient& PubSubClient::setServer(uint8_t* ip, uint16_t port) {
893 IPAddress addr(ip[0], ip[1], ip[2], ip[3]);
894 return setServer(addr, port);
895}
896
897PubSubClient& PubSubClient::setServer(IPAddress ip, uint16_t port) {
898 _ip = ip;
899 _port = port;
900 free(_domain);
901 _domain = nullptr;
902 return *this;
903}
904
905PubSubClient& PubSubClient::setServer(const char* domain, uint16_t port) {
906 char* newDomain = nullptr;
907 if (domain) {
908 newDomain = (char*)realloc(_domain, strlen(domain) + 1);
909 }
910 if (newDomain) {
911 strcpy(newDomain, domain);
912 _domain = newDomain;
913 _port = port;
914 } else {
915 free(_domain);
916 _domain = nullptr;
917 _port = 0;
918 }
919 return *this;
920}
921
922PubSubClient& PubSubClient::setCallback(MQTT_CALLBACK_SIGNATURE) {
923 this->callback = callback;
924 return *this;
925}
926
928 _client = &client;
929 return *this;
930}
931
933 _stream = &stream;
934 return *this;
935}
936
938 if (_bufferSize == size) return true; // if size is unchanged, do nothing and return true
939 // TODO: Adding a check for MQTT_CONNECTED would break the current setBufferSize() behavior. Implement with PubSubClient3 V4.
940 // if (_state == MQTT_CONNECTED) return false; // only allow to change the buffer size if the client is not connected
941 if (size == 0) {
942 // allow to free the buffer if the size is set to 0
943 if (_state != MQTT_CONNECTED) { // only if disconnected, otherwise we would free the buffer while it is still in use
944 free(_buffer);
945 _buffer = nullptr;
946 _bufferSize = 0;
947 }
948 } else if (size >= MQTT_MIN_BUFFER_SIZE) {
949 // buffer must be large enough to hold at least a minimal MQTT packet
950 uint8_t* newBuffer = (uint8_t*)realloc(_buffer, size); // realloc() is nullptr safe, so it will allocate a new buffer in this case‚
951 if (newBuffer) {
952 _buffer = newBuffer;
953 _bufferSize = size;
954 }
955 }
956 return (_bufferSize == size);
957}
958
960 return _bufferSize;
961}
962
964 _keepAliveMillis = keepAlive * 1000UL;
965 return *this;
966}
967
969 _socketTimeoutMillis = timeout * 1000UL;
970 return *this;
971}
972
974 return _state;
975}
#define PUBSUB_WAIT_YIELD()
Give other tasks a chance to run while waiting for data.
A simple client for MQTT.
#define MQTT_SOCKET_TIMEOUT
Sets the timeout, in seconds, when reading from the network. This also applies as the timeout for cal...
#define MQTT_MAX_TRANSFER_SIZE
Sets the maximum number of bytes passed to the network client in each write call. Some hardware has a...
#define MQTT_VERSION
Sets the version of the MQTT protocol to use (3.1 or 3.1.1). [MQTT_VERSION_3_1, MQTT_VERSION_3_1_1].
#define MQTT_MAX_PACKET_SIZE
Sets the largest packet size, in bytes, the client will handle. Any packet received that exceeds this...
#define MQTT_KEEPALIVE
Sets the keepalive interval, in seconds, the client will use. This is used to maintain the connection...
bool loop()
This should be called regularly to allow the client to process incoming messages and maintain its con...
PubSubClient & setCallback(MQTT_CALLBACK_SIGNATURE)
Sets the message callback function.
PubSubClient & setServer(IPAddress ip, uint16_t port)
Sets the server details.
bool publish_P(const char *topic, PGM_P payload, bool retained)
Publishes a message stored in PROGMEM to the specified topic using QoS 0.
bool beginPublish(const char *topic, size_t plength, bool retained)
Start to publish a message using QoS 0. This API: beginPublish(...) one or more calls to write(....
virtual size_t write(uint8_t data)
Writes a single byte as a component of a publish started with a call to beginPublish....
bool publish(const char *topic, const char *payload)
Publishes a non retained message to the specified topic using QoS 0.
~PubSubClient()
Destructor for the PubSubClient class.
PubSubClient & setSocketTimeout(uint16_t timeout)
Sets the socket timeout used by the client. This determines how long the client will wait for incomin...
PubSubClient()
Creates an uninitialised client instance.
PubSubClient & setKeepAlive(uint16_t keepAlive)
Sets the keep alive interval used by the client. This value should only be changed when the client is...
bool connected()
Checks whether the client is connected to the server.
int state()
Returns the current state of the client. If a connection attempt fails, this can be used to get more ...
bool endPublish()
Finish sending a message that was started with a call to beginPublish.
size_t write_P(PGM_P string)
Writes a string in PROGMEM as a component of a publish started with a call to beginPublish....
PubSubClient & setClient(Client &client)
Sets the network client instance to use.
void disconnect()
Disconnects the client.
size_t getBufferSize()
Gets the current size of the internal buffer.
bool setBufferSize(size_t size)
Sets the size, in bytes, of the internal send and receive buffer. This must be large enough to contai...
bool connect(const char *id)
Connects the client using a clean session without username and password.
PubSubClient & setStream(Stream &stream)
Sets the stream to write received messages to.
#define MQTT_QOS0
Quality of Service 0: At most once.
#define MQTT_QOS2
#define MQTT_QOS1
Quality of Service 1: At least once.
#define MQTT_CONNECTION_TIMEOUT
The network connection timed out or server didn't respond within the keepalive time.
#define MQTT_CONNECTED
The client is connected.
#define MQTT_CONNECT_FAILED
The network connection failed.
#define MQTT_CONNECTION_LOST
The network connection was lost/broken.
#define MQTT_DISCONNECTED
The client is disconnected cleanly.