24#if defined(ARDUINO_ARCH_ESP32) || defined(ESP32)
25#define PUBSUB_WAIT_YIELD() vTaskDelay(1)
27#define PUBSUB_WAIT_YIELD() yield()
37#define CHECK_STRING_LENGTH(l, s) \
38 if ((!s) || (l + 2 + strnlen(s, _bufferSize) > _bufferSize)) { \
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;
133 if (!_buffer)
return false;
137 if (_client->connected()) {
139 }
else if (_port != 0) {
141 result = _client->connect(_domain, _port);
143 result = _client->connect(_ip, _port);
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};
156 memcpy(_buffer + MQTT_MAX_HEADER_SIZE, protocol,
sizeof(protocol));
158 size_t length = MQTT_MAX_HEADER_SIZE +
sizeof(protocol);
159 uint8_t flags = 0x00;
161 flags = (0x01 << 2) | (willQos << 3) | (willRetain << 5);
164 flags = flags | (0x01 << 1);
167 flags = flags | (0x01 << 7);
169 flags = flags | (0x01 << 6);
172 const uint16_t keepAlive = _keepAliveMillis / 1000;
173 _buffer[length++] = flags;
174 _buffer[length++] = keepAlive >> 8;
175 _buffer[length++] = keepAlive & 0xFF;
177 CHECK_STRING_LENGTH(length,
id)
178 length = writeString(
id, length);
180 CHECK_STRING_LENGTH(length, willTopic)
181 length = writeString(willTopic, length);
182 CHECK_STRING_LENGTH(length, willMessage)
183 length = writeString(willMessage, length);
187 CHECK_STRING_LENGTH(length, user)
188 length = writeString(user, length);
190 CHECK_STRING_LENGTH(length, pass)
191 length = writeString(pass, length);
195 if (!writeControlPacket(MQTTCONNECT, length - MQTT_MAX_HEADER_SIZE)) {
200 _lastInActivity = _lastOutActivity = millis();
201 _pingOutstanding =
false;
203 while (!_client->available()) {
205 unsigned long t = millis();
206 if (t - _lastInActivity >= _socketTimeoutMillis) {
207 DEBUG_PSC_PRINTF(
"connect aborting due to timeout\n");
214 size_t len = readPacket(&hdrLen);
217 if (_buffer[3] == 0) {
218 _lastInActivity = millis();
225 DEBUG_PSC_PRINTF(
"connect aborting due to protocol error\n");
236 if (!_client)
return false;
237 if (!_buffer)
return false;
239 if (_client->connected()) {
242 DEBUG_PSC_PRINTF(
"lost connection (client may have more details)\n");
245 _pingOutstanding =
false;
251 DEBUG_PSC_PRINTF(
"disconnect called\n");
253 if (_client && _buffer) {
254 _buffer[0] = MQTTDISCONNECT;
256 _client->write(_buffer, 2);
259 _lastInActivity = _lastOutActivity = millis();
261 _pingOutstanding =
false;
270bool PubSubClient::readByte(uint8_t* result) {
271 if (!_client)
return false;
273 unsigned long previousMillis = millis();
274 while (!_client->available()) {
276 unsigned long currentMillis = millis();
277 if (currentMillis - previousMillis >= _socketTimeoutMillis) {
281 int rc = _client->read();
285 *result = (uint8_t)rc;
295bool PubSubClient::readByte(uint8_t* result,
size_t* pos) {
296 uint8_t* write_address = &(result[*pos]);
297 if (readByte(write_address)) {
310size_t PubSubClient::readPacket(uint8_t* hdrLen) {
312 if (!readByte(_buffer, &len))
return 0;
313 bool isPublish = (_buffer[0] & 0xF0) == MQTTPUBLISH;
314 uint32_t multiplier = 1;
321 if (len == MQTT_MAX_HEADER_SIZE) {
323 DEBUG_PSC_PRINTF(
"readPacket detected packet of invalid length\n");
328 if (!readByte(&digit))
return 0;
329 _buffer[len++] = digit;
330 length += (digit & 0x7F) * multiplier;
332 }
while ((digit & 0x80) != 0);
333 *hdrLen = (uint8_t)(len - 1);
335 DEBUG_PSC_PRINTF(
"readPacket received packet of length %zu (isPublish = %u)\n", length, isPublish);
339 if (!readByte(_buffer, &len))
return 0;
340 if (!readByte(_buffer, &len))
return 0;
341 skip = (_buffer[*hdrLen + 1] << 8) + _buffer[*hdrLen + 2];
343 if (MQTT_HDR_GET_QOS(_buffer[0]) >
MQTT_QOS0) {
350 for (
size_t i = start; i < length; i++) {
351 if (!readByte(&digit))
return 0;
353 if (isPublish && (idx - *hdrLen - 2 > skip)) {
354 _stream->write(digit);
358 if (len < _bufferSize) {
359 _buffer[len++] = digit;
364 if (!_stream && (idx > _bufferSize)) {
365 DEBUG_PSC_PRINTF(
"readPacket ignoring packet of size %zu exceeding buffer of size %zu\n", length, _bufferSize);
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) {
383 DEBUG_PSC_PRINTF(
"handlePacket(): packet length %zu exceeds buffer size %zu\n", length, _bufferSize);
399 if (length < hdrLen + 3ul) {
400 DEBUG_PSC_PRINTF(
"handlePacket(): Packet too short to contain topic length field\n");
403 const uint16_t topicLen = (_buffer[hdrLen + 1] << 8) + _buffer[hdrLen + 2];
404 char* topic = (
char*)(_buffer + hdrLen + 3 - 1);
409 const size_t payloadOffset = (size_t)hdrLen + 3 + topicLen;
410 const size_t payloadLen = length - payloadOffset;
411 uint8_t*
const payload = _buffer + payloadOffset;
414 if (length < payloadOffset) {
415 ERROR_PSC_PRINTF_P(
"handlePacket(): Suspicious topicLen (%u) points outside of received buffer length (%zu)\n", topicLen, length);
418 memmove(topic, topic + 1, topicLen);
419 topic[topicLen] =
'\0';
421 if (MQTT_HDR_GET_QOS(_buffer[0]) ==
MQTT_QOS0) {
423 callback(topic, payload, payloadLen);
426 if (payloadLen < 2) {
427 DEBUG_PSC_PRINTF(
"handlePacket(): Missing msgId in QoS 1/2 message\n");
430 const uint8_t publishQos = MQTT_HDR_GET_QOS(_buffer[0]);
431 const uint16_t msgId = (_buffer[payloadOffset] << 8) + _buffer[payloadOffset + 1];
432 callback(topic, payload + 2, payloadLen - 2);
436 _buffer[0] = (publishQos ==
MQTT_QOS1) ? MQTTPUBACK : MQTTPUBREC;
438 _buffer[2] = (uint8_t)(msgId >> 8);
439 _buffer[3] = (uint8_t)(msgId & 0xFF);
440 if (_client->write(_buffer, 4) == 4) {
441 _lastOutActivity = millis();
449 ERROR_PSC_PRINTF_P(
"handlePacket(): Received PUBACK packet with length %zu, expected at least 4 bytes\n", length);
458 ERROR_PSC_PRINTF_P(
"handlePacket(): Received PUBREC packet with length %zu, expected at least 4 bytes\n", length);
462 _buffer[0] = MQTTPUBREL | 2;
464 if (_client->write(_buffer, 4) == 4) {
465 _lastOutActivity = millis();
472 ERROR_PSC_PRINTF_P(
"handlePacket(): Received PUBREL packet with length %zu, expected at least 4 bytes\n", length);
476 _buffer[0] = MQTTPUBCOMP;
478 if (_client->write(_buffer, 4) == 4) {
479 _lastOutActivity = millis();
486 ERROR_PSC_PRINTF_P(
"handlePacket(): Received PUBCOMP packet with length %zu, expected at least 4 bytes\n", length);
493 _buffer[0] = MQTTPINGRESP;
495 if (_client->write(_buffer, 2) == 2) {
496 _lastOutActivity = millis();
500 _pingOutstanding =
false;
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");
518 _pingOutstanding =
false;
520 }
else if (_bufferWritePos > 0) {
522 if (flushBuffer() == 0) {
525 _pingOutstanding =
false;
529 _buffer[0] = MQTTPINGREQ;
531 if (_client->write(_buffer, 2) == 2) {
532 _lastInActivity = _lastOutActivity = t;
533 _pingOutstanding =
true;
537 if (_client->available()) {
539 size_t len = readPacket(&hdrLen);
542 ret = handlePacket(hdrLen, len);
557 size_t rc =
write(payload, plength);
563bool PubSubClient::publish(
const __FlashStringHelper* topic,
const uint8_t* payload,
size_t plength, uint8_t qos,
bool retained) {
565 size_t rc =
write(payload, plength);
573 size_t rc =
write_P(payload, plength);
579bool PubSubClient::publish_P(
const __FlashStringHelper* topic,
const uint8_t* payload,
size_t plength, uint8_t qos,
bool retained) {
581 size_t rc =
write_P(payload, plength);
598bool PubSubClient::beginPublishImpl(
bool progmem,
const char* topic,
size_t plength, uint8_t qos,
bool retained) {
599 if (!topic)
return false;
602 size_t topicLen = progmem ? strlen_P(topic) : strlen(topic);
603 if (topicLen == 0)
return false;
606 ERROR_PSC_PRINTF_P(
"beginPublish() called with invalid QoS %u\n", qos);
610 const size_t nextMsgLen = (qos >
MQTT_QOS0) ? 2 : 0;
612 if (
connected() && (MQTT_MAX_HEADER_SIZE + topicLen + 2 + nextMsgLen <= _bufferSize)) {
614 topicLen = writeStringImpl(progmem, topic, MQTT_MAX_HEADER_SIZE) - MQTT_MAX_HEADER_SIZE;
617 writeNextMsgId(MQTT_MAX_HEADER_SIZE + topicLen);
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;
624 size_t rc = _client->write(_buffer + (MQTT_MAX_HEADER_SIZE - hdrLen), hdrLen + topicLen + nextMsgLen);
625 _lastOutActivity = millis();
626 return (rc == (hdrLen + topicLen + nextMsgLen));
633 if (_bufferWritePos > 0) {
635 if (flushBuffer() == 0)
return false;
651uint8_t PubSubClient::buildHeader(uint8_t header,
size_t length) {
652 uint8_t hdrBuf[MQTT_MAX_HEADER_SIZE - 1];
662 hdrBuf[hdrLen++] = digit;
663 }
while ((len > 0) && (hdrLen < MQTT_MAX_HEADER_SIZE - 1));
666 ERROR_PSC_PRINTF_P(
"buildHeader: header=0x%02X, length too big %zu, left %zu\n", header, length, len);
670 _buffer[MQTT_MAX_HEADER_SIZE - 1 - hdrLen] = header;
671 memcpy(_buffer + MQTT_MAX_HEADER_SIZE - hdrLen, hdrBuf, hdrLen);
676 return appendBuffer(data);
680 for (
size_t i = 0; i < size; i++) {
681 if (appendBuffer(buf[i]) == 0)
return i;
687 for (
size_t i = 0; i < size; i++) {
688 if (appendBuffer((uint8_t)pgm_read_byte_near(buf + i)) == 0)
return i;
701bool PubSubClient::writeControlPacket(uint8_t header,
size_t length) {
702 uint8_t hdrLen = buildHeader(header, length);
703 if (hdrLen == 0)
return false;
705 return writeBuffer(MQTT_MAX_HEADER_SIZE - hdrLen, hdrLen + length);
715size_t PubSubClient::writeBuffer(
size_t pos,
size_t size) {
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;
722 while ((bytesRemaining > 0) && result) {
724 size_t bytesWritten = _client->write(writeBuf, bytesToWrite);
725 result = (bytesWritten == bytesToWrite);
726 bytesRemaining -= bytesWritten;
727 writeBuf += bytesWritten;
729 _lastOutActivity = millis();
733 rc = result ? size : 0;
735 rc = _client->write(_buffer + pos, size);
737 _lastOutActivity = millis();
758size_t PubSubClient::writeStringImpl(
bool progmem,
const char*
string,
size_t pos) {
759 if (!
string)
return pos;
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);
766 memcpy_P(_buffer + pos,
string, sLen);
768 memcpy(_buffer + pos,
string, sLen);
772 ERROR_PSC_PRINTF_P(
"writeStringImpl(): string (%zu) does not fit into buf (%zu)\n", pos + 2 + sLen, _bufferSize);
787inline size_t PubSubClient::writeString(
const char*
string,
size_t pos) {
788 return writeStringImpl(
false,
string, pos);
798size_t PubSubClient::writeNextMsgId(
size_t pos) {
799 if ((pos + 2) <= _bufferSize) {
800 _nextMsgId = (++_nextMsgId == 0) ? 1 : _nextMsgId;
801 _buffer[pos++] = (uint8_t)(_nextMsgId >> 8);
802 _buffer[pos++] = (uint8_t)(_nextMsgId & 0xFF);
804 ERROR_PSC_PRINTF_P(
"writeNextMsgId(): buffer overrun (%zu) \n", pos + 2);
815size_t PubSubClient::appendBuffer(uint8_t data) {
816 _buffer[_bufferWritePos++] = data;
817 if (_bufferWritePos >= _bufferSize) {
818 if (flushBuffer() == 0)
return 0;
829size_t PubSubClient::flushBuffer() {
832 rc = writeBuffer(0, _bufferWritePos);
846bool PubSubClient::subscribeImpl(
bool progmem,
const char* topic, uint8_t qos) {
847 if (!topic)
return false;
851 size_t topicLen = progmem ? strnlen_P(topic, _bufferSize) : strnlen(topic, _bufferSize);
852 if (_bufferSize < MQTT_MAX_HEADER_SIZE + 2 + 2 + topicLen + 1) {
858 size_t length = MQTT_MAX_HEADER_SIZE;
859 length = writeNextMsgId(length);
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);
874bool PubSubClient::unsubscribeImpl(
bool progmem,
const char* topic) {
875 if (!topic)
return false;
878 size_t topicLen = progmem ? strnlen_P(topic, _bufferSize) : strnlen(topic, _bufferSize);
879 if (_bufferSize < MQTT_MAX_HEADER_SIZE + 2 + 2 + topicLen) {
884 size_t length = MQTT_MAX_HEADER_SIZE;
885 length = writeNextMsgId(length);
886 length = writeStringImpl(progmem, topic, length);
887 return writeControlPacket(MQTTUNSUBSCRIBE | MQTT_QOS_GET_HDR(
MQTT_QOS1), length - MQTT_MAX_HEADER_SIZE);
893 IPAddress addr(ip[0], ip[1], ip[2], ip[3]);
906 char* newDomain =
nullptr;
908 newDomain = (
char*)realloc(_domain, strlen(domain) + 1);
911 strcpy(newDomain, domain);
923 this->callback = callback;
938 if (_bufferSize == size)
return true;
948 }
else if (size >= MQTT_MIN_BUFFER_SIZE) {
950 uint8_t* newBuffer = (uint8_t*)realloc(_buffer, size);
956 return (_bufferSize == size);
964 _keepAliveMillis = keepAlive * 1000UL;
969 _socketTimeoutMillis = timeout * 1000UL;
#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_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.