diff --git a/lib/pubsubclient3/.github/doxygen-header.html b/lib/pubsubclient3/.github/doxygen-header.html new file mode 100644 index 000000000..fbe29e8a2 --- /dev/null +++ b/lib/pubsubclient3/.github/doxygen-header.html @@ -0,0 +1,77 @@ + + + + + + + + +$projectname: $title +$title + + + + + + + + + + + + +$treeview +$search +$mathjax +$darkmode + +$extrastylesheet + + + +
+ + +
+ + +
+ + + + + + + + + + + + + + + + + + + + + + + +
+
$projectname $projectnumber
+ +
$projectbrief
+ + +
$searchbox
$searchbox
+
+ + diff --git a/lib/pubsubclient3/.github/doxygen-update-version.sh b/lib/pubsubclient3/.github/doxygen-update-version.sh new file mode 100644 index 000000000..a2ec04181 --- /dev/null +++ b/lib/pubsubclient3/.github/doxygen-update-version.sh @@ -0,0 +1,9 @@ +#!/bin/bash +# Update the Doxyfile PROJECT_NUMBER with the version from library.properties +# Exit immediately if a command exits with a non-zero status. +set -e +# Make sure we are inside the github workspace +cd $GITHUB_WORKSPACE +version=`grep "version=" library.properties | cut -d "=" -f 2-` +echo "Update current version in Doxyfile: $version" +sed -i "s/PROJECT_NUMBER.*/PROJECT_NUMBER = \"v$version\"/" Doxyfile diff --git a/lib/pubsubclient3/.github/workflows/build.yml b/lib/pubsubclient3/.github/workflows/build.yml index f8c919060..110a8f0ec 100644 --- a/lib/pubsubclient3/.github/workflows/build.yml +++ b/lib/pubsubclient3/.github/workflows/build.yml @@ -3,6 +3,10 @@ name: build on: push: branches: [ "master" ] + paths: + - 'examples/**' + - 'src/**' + - 'tests/**' pull_request: branches: [ "master" ] diff --git a/lib/pubsubclient3/.github/workflows/gh-pages.yml b/lib/pubsubclient3/.github/workflows/gh-pages.yml new file mode 100644 index 000000000..6eed1ea16 --- /dev/null +++ b/lib/pubsubclient3/.github/workflows/gh-pages.yml @@ -0,0 +1,65 @@ +# This workflow generates Doxygen API documentation and deploys it to GitHub Pages. +# I didn't manage to build a Jekyll site and deploy the content to GitHub Pages. +# So now I use a static index.html file in the gh-pages branch. +# It runs on pushes to the master branch and can also be triggered manually. +name: Generate Doxygen API and deploy content to GitHub Pages + +on: + # Runs on pushes targeting the default branch + push: + branches: ["master"] + + # Allows you to run this workflow manually from the Actions tab + workflow_dispatch: + +# Sets permissions of the GITHUB_TOKEN to allow deployment to GitHub Pages +permissions: + contents: write + pages: write + id-token: write + +# Allow only one concurrent deployment, skipping runs queued between the run in-progress and latest queued. +# However, do NOT cancel in-progress runs as we want to allow these production deployments to complete. +concurrency: + group: "pages" + cancel-in-progress: false + +jobs: + # Single deploy job since we're just deploying + deploy: + if: github.repository == 'hmueller01/pubsubclient3' + environment: + name: github-pages + url: ${{ steps.deployment.outputs.page_url }} + runs-on: ubuntu-latest + steps: + - name: Checkout + uses: actions/checkout@v4 + - name: Update Doxyfile PROJECT_NUMBER + run: bash .github/doxygen-update-version.sh + - name: Doxygenize + #uses: langroodi/doxygenize@v1.7.1 # waiting for update ... + uses: langroodi/doxygenize@6e920681c5d838e9a1b5cd273b814d1e4023b63d + #uses: hmueller01/doxygenize@update-alpine-linux+patches + with: + # Doxygen configuration file path + doxygenconf: ./Doxyfile + # Generated HTML documents output folder + htmloutput: ./site/api + # GitHub Pages branch name + ghpagesbranch: gh-pages + # GitHub Pages directory path + ghpagesdir: ./site/api + # - name: List site contents + # run: | + # find ./site + - name: Setup Pages + uses: actions/configure-pages@v5 + - name: Upload artifact + uses: actions/upload-pages-artifact@v3 + with: + #path: '.' # Upload entire repository + path: './site' + - name: Deploy to GitHub Pages + id: deployment + uses: actions/deploy-pages@v4 diff --git a/lib/pubsubclient3/.gitignore b/lib/pubsubclient3/.gitignore index a42cc406e..9ffb670e2 100644 --- a/lib/pubsubclient3/.gitignore +++ b/lib/pubsubclient3/.gitignore @@ -3,3 +3,4 @@ tests/bin .piolibdeps .clang_complete .gcc-flags.json +site/api diff --git a/lib/pubsubclient3/Doxyfile b/lib/pubsubclient3/Doxyfile new file mode 100644 index 000000000..8e0e065b8 --- /dev/null +++ b/lib/pubsubclient3/Doxyfile @@ -0,0 +1,39 @@ +# Minimal Doxyfile for Doxygen 1.13.2 + +# Project related configuration options +PROJECT_NAME = "PubSubClient3" +PROJECT_NUMBER = "automatically updated by doxygen-update-version.sh" +PROJECT_BRIEF = "Located at GitHub" +OUTPUT_DIRECTORY = ./site +INPUT = ./src +RECURSIVE = YES + +# Build related configuration options +EXTRACT_ALL = YES +EXTRACT_PRIVATE = NO +EXTRACT_STATIC = NO +EXCLUDE_SYMBOLS = CHECK_STRING_LENGTH + +# HTML output +GENERATE_HTML = YES +HTML_OUTPUT = api +HTML_COLORSTYLE = AUTO_LIGHT +HTML_HEADER = .github/doxygen-header.html + +# LaTeX output (disable if not needed) +GENERATE_LATEX = NO + +# Source code browsing +SOURCE_BROWSER = YES + +# Warnings +WARNINGS = YES + +# Dot/Graphviz support (disable if not needed) +HAVE_DOT = NO + +# File patterns (optional, default covers most C/C++ files) +FILE_PATTERNS = *.h *.cpp *.c *.hpp + +# Language +OUTPUT_LANGUAGE = English diff --git a/lib/pubsubclient3/README.md b/lib/pubsubclient3/README.md index 8233d8e11..f285bee6e 100644 --- a/lib/pubsubclient3/README.md +++ b/lib/pubsubclient3/README.md @@ -24,12 +24,19 @@ I appreciate every contribution to this library. The library comes with a number of example sketches. See File > Examples > PubSubClient within the Arduino application. -Full API documentation is available here: https://pubsubclient.knolleary.net +Full API documentation is available here: https://hmueller01.github.io/pubsubclient3/api ## Limitations - The client is based on the [MQTT Version 3.1.1 specification](https://docs.oasis-open.org/mqtt/mqtt/v3.1.1/os/mqtt-v3.1.1-os.html) with some limitations. - - It can only publish QoS 0 messages. It can subscribe at QoS 0 or QoS 1. + - It can publish at QoS 0, 1 or 2. + + **WARNING:** No retransmission is supported to keep the library as much memory friendly as possible. + + Note: Without retransmission support, the publish QoS is only meaningful when the broker sends your + message to a subscriber, supposing that the subscriber subscribes with a QoS greater then or equal to + the publish QoS. Consider that MQTT runs over TCP, so retransmission isn't really required in most cases, especially when publishing to the broker. + - It can subscribe at QoS 0 or QoS 1. - The maximum message size, including header, is **256 bytes** by default. This is configurable via `MQTT_MAX_PACKET_SIZE` in `PubSubClient.h` or can be changed by calling `PubSubClient::setBufferSize(size)`. diff --git a/lib/pubsubclient3/examples/mqtt_auth/mqtt_auth.ino b/lib/pubsubclient3/examples/mqtt_auth/mqtt_auth.ino index cdeb7959d..5d5d95ce4 100644 --- a/lib/pubsubclient3/examples/mqtt_auth/mqtt_auth.ino +++ b/lib/pubsubclient3/examples/mqtt_auth/mqtt_auth.ino @@ -16,7 +16,7 @@ byte mac[] = {0xDE, 0xED, 0xBA, 0xFE, 0xFE, 0xED}; IPAddress ip(172, 16, 0, 100); IPAddress server(172, 16, 0, 2); -void callback(char* topic, uint8_t* payload, size_t length) { +void callback(char* topic, uint8_t* payload, size_t plength) { // handle message arrived } diff --git a/lib/pubsubclient3/examples/mqtt_basic/mqtt_basic.ino b/lib/pubsubclient3/examples/mqtt_basic/mqtt_basic.ino index f790aab6a..b38716899 100644 --- a/lib/pubsubclient3/examples/mqtt_basic/mqtt_basic.ino +++ b/lib/pubsubclient3/examples/mqtt_basic/mqtt_basic.ino @@ -21,11 +21,11 @@ byte mac[] = {0xDE, 0xED, 0xBA, 0xFE, 0xFE, 0xED}; IPAddress ip(172, 16, 0, 100); IPAddress server(172, 16, 0, 2); -void callback(char* topic, uint8_t* payload, size_t length) { +void callback(char* topic, uint8_t* payload, size_t plength) { Serial.print("Message arrived ["); Serial.print(topic); Serial.print("] "); - for (size_t i = 0; i < length; i++) { + for (size_t i = 0; i < plength; i++) { Serial.print((char)payload[i]); } Serial.println(); diff --git a/lib/pubsubclient3/examples/mqtt_esp8266/mqtt_esp8266.ino b/lib/pubsubclient3/examples/mqtt_esp8266/mqtt_esp8266.ino index b41192ab4..9ffe0e725 100644 --- a/lib/pubsubclient3/examples/mqtt_esp8266/mqtt_esp8266.ino +++ b/lib/pubsubclient3/examples/mqtt_esp8266/mqtt_esp8266.ino @@ -66,11 +66,11 @@ void setup_wifi() { Serial.println(WiFi.localIP()); } -void callback(char* topic, uint8_t* payload, size_t length) { +void callback(char* topic, uint8_t* payload, size_t plength) { Serial.print("Message arrived ["); Serial.print(topic); Serial.print("] "); - for (size_t i = 0; i < length; i++) { + for (size_t i = 0; i < plength; i++) { Serial.print((char)payload[i]); } Serial.println(); diff --git a/lib/pubsubclient3/examples/mqtt_large_message/mqtt_large_message.ino b/lib/pubsubclient3/examples/mqtt_large_message/mqtt_large_message.ino index 48049c3e9..3fc23c8c5 100644 --- a/lib/pubsubclient3/examples/mqtt_large_message/mqtt_large_message.ino +++ b/lib/pubsubclient3/examples/mqtt_large_message/mqtt_large_message.ino @@ -69,11 +69,11 @@ void setup_wifi() { Serial.println(WiFi.localIP()); } -void callback(char* topic, uint8_t* payload, size_t length) { +void callback(char* topic, uint8_t* payload, size_t plength) { Serial.print("Message arrived ["); Serial.print(topic); Serial.print("] "); - for (size_t i = 0; i < length; i++) { + for (size_t i = 0; i < plength; i++) { Serial.print((char)payload[i]); } Serial.println(); diff --git a/lib/pubsubclient3/examples/mqtt_publish_in_callback/mqtt_publish_in_callback.ino b/lib/pubsubclient3/examples/mqtt_publish_in_callback/mqtt_publish_in_callback.ino index 733bea3f5..f30159706 100644 --- a/lib/pubsubclient3/examples/mqtt_publish_in_callback/mqtt_publish_in_callback.ino +++ b/lib/pubsubclient3/examples/mqtt_publish_in_callback/mqtt_publish_in_callback.ino @@ -23,22 +23,22 @@ IPAddress ip(172, 16, 0, 100); IPAddress server(172, 16, 0, 2); // Callback function header -void callback(char* topic, uint8_t* payload, size_t length); +void callback(char* topic, uint8_t* payload, size_t plength); EthernetClient ethClient; PubSubClient client(server, 1883, callback, ethClient); // Callback function -void callback(char* topic, uint8_t* payload, size_t length) { +void callback(char* topic, uint8_t* payload, size_t plength) { // In order to republish this payload, a copy must be made // as the orignal payload buffer will be overwritten whilst // constructing the PUBLISH packet. // Allocate the correct amount of memory for the payload copy - byte* p = (byte*)malloc(length); + byte* p = (byte*)malloc(plength); // Copy the payload to the new buffer - memcpy(p, payload, length); - client.publish("outTopic", p, length); + memcpy(p, payload, plength); + client.publish("outTopic", p, plength); // Free the memory free(p); } diff --git a/lib/pubsubclient3/examples/mqtt_reconnect_nonblocking/mqtt_reconnect_nonblocking.ino b/lib/pubsubclient3/examples/mqtt_reconnect_nonblocking/mqtt_reconnect_nonblocking.ino index 9d5c60098..d0797419b 100644 --- a/lib/pubsubclient3/examples/mqtt_reconnect_nonblocking/mqtt_reconnect_nonblocking.ino +++ b/lib/pubsubclient3/examples/mqtt_reconnect_nonblocking/mqtt_reconnect_nonblocking.ino @@ -16,7 +16,7 @@ byte mac[] = {0xDE, 0xED, 0xBA, 0xFE, 0xFE, 0xED}; IPAddress ip(172, 16, 0, 100); IPAddress server(172, 16, 0, 2); -void callback(char* topic, uint8_t* payload, size_t length) { +void callback(char* topic, uint8_t* payload, size_t plength) { // handle message arrived } diff --git a/lib/pubsubclient3/examples/mqtt_stream/mqtt_stream.ino b/lib/pubsubclient3/examples/mqtt_stream/mqtt_stream.ino index 5e620a0e3..54ab20fa0 100644 --- a/lib/pubsubclient3/examples/mqtt_stream/mqtt_stream.ino +++ b/lib/pubsubclient3/examples/mqtt_stream/mqtt_stream.ino @@ -21,11 +21,11 @@ IPAddress server(172, 16, 0, 2); SRAM sram(4, SRAM_1024); -void callback(char* topic, uint8_t* payload, size_t length) { +void callback(char* topic, uint8_t* payload, size_t plength) { sram.seek(1); // do something with the message - for (size_t i = 0; i < length; i++) { + for (size_t i = 0; i < plength; i++) { Serial.write(sram.read()); } Serial.println(); diff --git a/lib/pubsubclient3/src/PubSubClient.cpp b/lib/pubsubclient3/src/PubSubClient.cpp index dd1e5bf94..96aee156e 100644 --- a/lib/pubsubclient3/src/PubSubClient.cpp +++ b/lib/pubsubclient3/src/PubSubClient.cpp @@ -1,12 +1,28 @@ -/* - PubSubClient.cpp - A simple client for MQTT. - Nick O'Leary, Holger Mueller - http://knolleary.net - https://github.com/hmueller01/pubsubclient3 -*/ +/** + * @file PubSubClient.cpp + * @brief A simple client for MQTT. + * @author Nicholas O'Leary - http://knolleary.net + * @author Holger Mueller - https://github.com/hmueller01/pubsubclient3 + * @copyright MIT License 2008-2025 + * + * This file is part of the PubSubClient library. + */ #include "PubSubClient.h" +/** + * @brief Macro to check if a string 's' can be safely added to the MQTT buffer. + * + * If either check fails, the client connection is stopped and the function returns false. + * @param l current length in the buffer + * @param s string to check + */ +#define CHECK_STRING_LENGTH(l, s) \ + if ((!s) || (l + 2 + strnlen(s, this->bufferSize) > this->bufferSize)) { \ + _client->stop(); \ + return false; \ + } + PubSubClient::PubSubClient() { setBufferSize(MQTT_MAX_PACKET_SIZE); setKeepAlive(MQTT_KEEPALIVE); @@ -95,20 +111,20 @@ PubSubClient::~PubSubClient() { } bool PubSubClient::connect(const char* id) { - return connect(id, nullptr, nullptr, 0, 0, 0, 0, 1); + return connect(id, nullptr, nullptr, nullptr, MQTT_QOS0, false, nullptr, true); } bool PubSubClient::connect(const char* id, const char* user, const char* pass) { - return connect(id, user, pass, 0, 0, 0, 0, 1); + return connect(id, user, pass, nullptr, MQTT_QOS0, false, nullptr, true); } bool PubSubClient::connect(const char* id, const char* willTopic, uint8_t willQos, bool willRetain, const char* willMessage) { - return connect(id, nullptr, nullptr, willTopic, willQos, willRetain, willMessage, 1); + return connect(id, nullptr, nullptr, willTopic, willQos, willRetain, willMessage, true); } bool PubSubClient::connect(const char* id, const char* user, const char* pass, const char* willTopic, uint8_t willQos, bool willRetain, const char* willMessage) { - return connect(id, user, pass, willTopic, willQos, willRetain, willMessage, 1); + return connect(id, user, pass, willTopic, willQos, willRetain, willMessage, true); } bool PubSubClient::connect(const char* id, const char* user, const char* pass, const char* willTopic, uint8_t willQos, bool willRetain, @@ -157,20 +173,20 @@ bool PubSubClient::connect(const char* id, const char* user, const char* pass, c this->buffer[length++] = keepAlive & 0xFF; CHECK_STRING_LENGTH(length, id) - length = writeString(id, this->buffer, length); + length = writeString(id, this->buffer, length, this->bufferSize); if (willTopic) { CHECK_STRING_LENGTH(length, willTopic) - length = writeString(willTopic, this->buffer, length); + length = writeString(willTopic, this->buffer, length, this->bufferSize); CHECK_STRING_LENGTH(length, willMessage) - length = writeString(willMessage, this->buffer, length); + length = writeString(willMessage, this->buffer, length, this->bufferSize); } if (user) { CHECK_STRING_LENGTH(length, user) - length = writeString(user, this->buffer, length); + length = writeString(user, this->buffer, length, this->bufferSize); if (pass) { CHECK_STRING_LENGTH(length, pass) - length = writeString(pass, this->buffer, length); + length = writeString(pass, this->buffer, length, this->bufferSize); } } @@ -314,8 +330,8 @@ size_t PubSubClient::readPacket(uint8_t* hdrLen) { if (!readByte(this->buffer, &len)) return 0; skip = (this->buffer[*hdrLen + 1] << 8) + this->buffer[*hdrLen + 2]; start = 2; - if (this->buffer[0] & MQTTQOS1) { - // skip message id + if (MQTT_HDR_GET_QOS(this->buffer[0]) > MQTT_QOS0) { + // skip msgId (packet identifier) for QoS 1 and 2 messages skip += 2; } } @@ -355,7 +371,7 @@ bool PubSubClient::handlePacket(uint8_t hdrLen, size_t length) { switch (type) { case MQTTPUBLISH: if (callback) { - // MQTT Publish packet: See section 3.3 MQTT v3.1 protocol specification: + // MQTT Publish packet: See section 3.3 MQTT v3.1.1 protocol specification: // - Header: 1 byte // - Remaining header length: hdrLen bytes, multibyte field (1 .. MQTT_MAX_HEADER_SIZE - 1) // - Topic length: 2 bytes (starts at buffer[hdrLen + 1]) @@ -377,7 +393,7 @@ bool PubSubClient::handlePacket(uint8_t hdrLen, size_t length) { memmove(topic, topic + 1, topicLen); // move topic inside buffer 1 byte to front topic[topicLen] = '\0'; // end the topic as a 'C' string with \x00 - if ((this->buffer[0] & 0x06) == MQTTQOS0) { + if (MQTT_HDR_GET_QOS(this->buffer[0]) == MQTT_QOS0) { // No msgId for QOS == 0 callback(topic, payload, payloadLen); } else { @@ -399,7 +415,37 @@ bool PubSubClient::handlePacket(uint8_t hdrLen, size_t length) { } } break; + case MQTTPUBACK: + // MQTT Publish Acknowledgment (QoS 1 publish received): See section 3.4 MQTT v3.1.1 protocol specification + if (length < 4) { + ERROR_PSC_PRINTF_P("handlePacket(): Received PUBACK packet with length %zu, expected at least 4 bytes\n", length); + return false; + } + // No futher action here, as resending is not supported. + break; + case MQTTPUBREC: + // MQTT Publish Received (QoS 2 publish received, part 1): See section 3.5 MQTT v3.1.1 protocol specification + if (length < 4) { + ERROR_PSC_PRINTF_P("handlePacket(): Received PUBREC packet with length %zu, expected at least 4 bytes\n", length); + return false; + } + // MQTT Publish Release (QoS 2 publish received, part 2): See section 3.6 MQTT v3.1.1 protocol specification + buffer[0] = MQTTPUBREL | 2; // PUBREL with bit 1 set + // bytes 1-3 of PUBREL are the same as of PUBREC + if (_client->write(buffer, 4) == 4) { + lastOutActivity = millis(); + } + break; + case MQTTPUBCOMP: + // MQTT Publish Complete (QoS 2 publish received, part 3): See section 3.7 MQTT v3.1.1 protocol specification + if (length < 4) { + ERROR_PSC_PRINTF_P("handlePacket(): Received PUBCOMP packet with length %zu, expected at least 4 bytes\n", length); + return false; + } + // No futher action here, as resending is not supported. + break; case MQTTPINGREQ: + // MQTT Ping Request: See section 3.12 MQTT v3.1.1 protocol specification this->buffer[0] = MQTTPINGRESP; this->buffer[1] = 0; if (_client->write(this->buffer, 2) == 2) { @@ -456,19 +502,27 @@ bool PubSubClient::loop() { } bool PubSubClient::publish(const char* topic, const char* payload) { - return publish(topic, (const uint8_t*)payload, payload ? strnlen(payload, MQTT_MAX_POSSIBLE_PACKET_SIZE) : 0, false); + return publish(topic, payload, MQTT_QOS0, false); } bool PubSubClient::publish(const char* topic, const char* payload, bool retained) { - return publish(topic, (const uint8_t*)payload, payload ? strnlen(payload, MQTT_MAX_POSSIBLE_PACKET_SIZE) : 0, retained); + return publish(topic, payload, MQTT_QOS0, retained); +} + +bool PubSubClient::publish(const char* topic, const char* payload, uint8_t qos, bool retained) { + return publish(topic, (const uint8_t*)payload, payload ? strnlen(payload, MQTT_MAX_POSSIBLE_PACKET_SIZE) : 0, qos, retained); } bool PubSubClient::publish(const char* topic, const uint8_t* payload, size_t plength) { - return publish(topic, payload, plength, false); + return publish(topic, payload, plength, MQTT_QOS0, false); } bool PubSubClient::publish(const char* topic, const uint8_t* payload, size_t plength, bool retained) { - if (beginPublish(topic, plength, retained)) { + return publish(topic, payload, plength, MQTT_QOS0, retained); +} + +bool PubSubClient::publish(const char* topic, const uint8_t* payload, size_t plength, uint8_t qos, bool retained) { + if (beginPublish(topic, plength, qos, retained)) { size_t rc = write(payload, plength); lastOutActivity = millis(); return endPublish() && (rc == plength); @@ -477,11 +531,19 @@ bool PubSubClient::publish(const char* topic, const uint8_t* payload, size_t ple } bool PubSubClient::publish_P(const char* topic, const char* payload, bool retained) { - return publish_P(topic, (const uint8_t*)payload, payload ? strnlen_P(payload, MQTT_MAX_POSSIBLE_PACKET_SIZE) : 0, retained); + return publish_P(topic, payload, MQTT_QOS0, retained); +} + +bool PubSubClient::publish_P(const char* topic, const char* payload, uint8_t qos, bool retained) { + return publish_P(topic, (const uint8_t*)payload, payload ? strnlen_P(payload, MQTT_MAX_POSSIBLE_PACKET_SIZE) : 0, qos, retained); } bool PubSubClient::publish_P(const char* topic, const uint8_t* payload, size_t plength, bool retained) { - if (beginPublish(topic, plength, retained)) { + return publish_P(topic, payload, plength, MQTT_QOS0, retained); +} + +bool PubSubClient::publish_P(const char* topic, const uint8_t* payload, size_t plength, uint8_t qos, bool retained) { + if (beginPublish(topic, plength, qos, retained)) { size_t rc = 0; for (size_t i = 0; i < plength; i++) { rc += _client->write((uint8_t)pgm_read_byte_near(payload + i)); @@ -493,14 +555,26 @@ bool PubSubClient::publish_P(const char* topic, const uint8_t* payload, size_t p } bool PubSubClient::beginPublish(const char* topic, size_t plength, bool retained) { + return beginPublish(topic, plength, MQTT_QOS0, retained); +} + +bool PubSubClient::beginPublish(const char* topic, size_t plength, uint8_t qos, bool retained) { if (!topic) return false; + if (strlen(topic) == 0) return false; // empty topic is not allowed + if (qos > MQTT_QOS2) { // only valid QoS supported + this->_qos = MQTT_QOS0; // reset QoS to 0, that endPublish() will not send a nextMsgId + ERROR_PSC_PRINTF_P("beginPublish() called with invalid QoS %u\n", qos); + return false; + } + this->_qos = qos; // save the QoS for later endPublish() operation // check if the header and the topic (including 2 length bytes) fit into the buffer if (connected() && MQTT_MAX_HEADER_SIZE + strlen(topic) + 2 <= this->bufferSize) { // first write the topic at the end of the maximal variable header (MQTT_MAX_HEADER_SIZE) to the buffer - size_t topicLen = writeString(topic, this->buffer, MQTT_MAX_HEADER_SIZE) - MQTT_MAX_HEADER_SIZE; + size_t topicLen = writeString(topic, this->buffer, MQTT_MAX_HEADER_SIZE, this->bufferSize) - MQTT_MAX_HEADER_SIZE; // we now know the length of the topic string (lenght + 2 bytes signalling the length) and can build the variable header information - const uint8_t header = MQTTPUBLISH | (retained ? MQTTRETAINED : 0); - uint8_t hdrLen = buildHeader(header, this->buffer, topicLen + plength); + const uint8_t header = MQTTPUBLISH | MQTT_QOS_GET_HDR(qos) | (retained ? MQTTRETAINED : 0); + const size_t nextMsgLen = (qos) ? 2 : 0; // add 2 bytes for the nextMsgId if QoS > 0 + uint8_t hdrLen = buildHeader(header, this->buffer, topicLen + plength + nextMsgLen); if (hdrLen == 0) return false; // exit here in case of header generation failure // as the header length is variable, it starts at MQTT_MAX_HEADER_SIZE - hdrLen (see buildHeader() documentation) size_t rc = _client->write(this->buffer + (MQTT_MAX_HEADER_SIZE - hdrLen), hdrLen + topicLen); @@ -511,8 +585,19 @@ bool PubSubClient::beginPublish(const char* topic, size_t plength, bool retained } bool PubSubClient::endPublish() { - flushBuffer(); - return connected(); + if (connected()) { + if (this->_qos > MQTT_QOS0) { + // QoS == 1 or 2, send the msgId + uint8_t buf[2]; + writeNextMsgId(buf, 0, 2); + size_t rc = _client->write(buf, 2); + lastOutActivity = millis(); + return (rc == 2); + } + // QoS == 0, no msgId to send + return true; + } + return false; } /** @@ -550,19 +635,13 @@ uint8_t PubSubClient::buildHeader(uint8_t header, uint8_t* buf, size_t length) { } size_t PubSubClient::write(uint8_t data) { - const size_t rc = appendBuffer(data); - if (rc != 0) { - lastOutActivity = millis(); - } - return rc; + lastOutActivity = millis(); + return _client->write(data); } size_t PubSubClient::write(const uint8_t* buffer, size_t size) { - const size_t rc = appendBuffer(buffer, size); - if (rc != 0) { - lastOutActivity = millis(); - } - return rc; + lastOutActivity = millis(); + return _client->write(buffer, size); } /** @@ -605,24 +684,10 @@ bool PubSubClient::write(uint8_t header, uint8_t* buf, size_t length) { } /** - * @brief Write an UTF-8 encoded string to the give buffer and position. The string can have a length of 0 to 65535 bytes. The buffer is prefixed with two - * bytes representing the length of the string. See section 1.5.3 of MQTT v3.1.1 protocol specification. - * @note If the string does not fit in the buffer (bufferSize) or is longer than 65535 bytes nothing is written to the buffer and the returned position - * is unchanged. - * - * @param string 'C' string of the data that shall be written in the buffer. - * @param buf Buffer to write the string into. - * @param pos Position in the buffer to write the string. - * @return New position in the buffer (pos + 2 + string length), or pos if a buffer overrun would occur. - */ -size_t PubSubClient::writeString(const char* string, uint8_t* buf, size_t pos) { - return writeString(string, buf, pos, this->bufferSize); -} - -/** - * @brief Write an UTF-8 encoded string to the give buffer and position. The string can have a length of 0 to 65535 bytes. The buffer is prefixed with two - * bytes representing the length of the string. See section 1.5.3 of MQTT v3.1.1 protocol specification. - * @note If the string does not fit in the buffer or is longer than 65535 bytes nothing is written to the buffer and the returned position is unchanged. + * @brief Write an UTF-8 encoded string to the give buffer and position. The string can have a length of 0 to 65535 bytes. The buffer is prefixed with + * two bytes representing the length of the string. See section 1.5.3 of MQTT v3.1.1 protocol specification. + * @note If the string does not fit in the buffer or is longer than 65535 bytes nothing is written to the buffer and the returned position is + * unchanged. * * @param string 'C' string of the data that shall be written in the buffer. * @param buf Buffer to write the string into. @@ -645,61 +710,46 @@ size_t PubSubClient::writeString(const char* string, uint8_t* buf, size_t pos, s return pos; } -size_t PubSubClient::appendBuffer(uint8_t data) { - buffer[_bufferWritePos] = data; - ++_bufferWritePos; - if (_bufferWritePos >= bufferSize) { - if (flushBuffer() == 0) return 0; +/** + * @brief Write nextMsgId to the give buffer and position. + * @note If the nextMsgId (2 bytes) does not fit in the buffer nothing is written to the buffer and the returned position is unchanged. + * + * @param buf Buffer to write the nextMsgId into. + * @param pos Position in the buffer to write the nextMsgId. + * @param size Maximal size of the buffer. + * @return New position in the buffer (pos + 2), or pos if a buffer overrun would occur. + */ +size_t PubSubClient::writeNextMsgId(uint8_t* buf, size_t pos, size_t size) { + if (pos + 2 <= size) { + nextMsgId = (++nextMsgId == 0) ? 1 : nextMsgId; // increment msgId (must not be 0, so start at 1) + buf[pos++] = (uint8_t)(nextMsgId >> 8); + buf[pos++] = (uint8_t)(nextMsgId & 0xFF); + } else { + ERROR_PSC_PRINTF_P("writeNextMsgId(): buffer (%zu) does not fit into buf (%zu)\n", pos + 2, size); } - return 1; -} - -size_t PubSubClient::appendBuffer(const uint8_t *data, size_t size) { - for (size_t i = 0; i < size; ++i) { - if (appendBuffer(data[i]) == 0) return i; - } - return size; -} - -size_t PubSubClient::flushBuffer() { - size_t rc = 0; - if (_bufferWritePos > 0) { - if (connected()) { - rc = _client->write(buffer, _bufferWritePos); - if (rc != 0) { - lastOutActivity = millis(); - } - } - _bufferWritePos = 0; - } - return rc; + return pos; } bool PubSubClient::subscribe(const char* topic) { - return subscribe(topic, 0); + return subscribe(topic, MQTT_QOS0); } bool PubSubClient::subscribe(const char* topic, uint8_t qos) { if (!topic) return false; - if (qos > 1) return false; // only QoS 0 and 1 supported + if (qos > MQTT_QOS1) return false; // only QoS 0 and 1 supported size_t topicLen = strnlen(topic, this->bufferSize); - if (this->bufferSize < MQTT_MAX_HEADER_SIZE + 2 + 2 + topicLen) { - // Too long: header + nextMsgId (2) + topic length bytes (2) + topicLen + if (this->bufferSize < MQTT_MAX_HEADER_SIZE + 2 + 2 + topicLen + 1) { + // Too long: header + nextMsgId (2) + topic length bytes (2) + topicLen + QoS (1) return false; } if (connected()) { // Leave room in the buffer for header and variable length field uint16_t length = MQTT_MAX_HEADER_SIZE; - nextMsgId++; - if (nextMsgId == 0) { - nextMsgId = 1; - } - this->buffer[length++] = (nextMsgId >> 8); - this->buffer[length++] = (nextMsgId & 0xFF); - length = writeString(topic, this->buffer, length); + length = writeNextMsgId(buffer, length, this->bufferSize); // buffer size is checked before + length = writeString(topic, this->buffer, length, this->bufferSize); this->buffer[length++] = qos; - return write(MQTTSUBSCRIBE | MQTTQOS1, this->buffer, length - MQTT_MAX_HEADER_SIZE); + return write(MQTTSUBSCRIBE | MQTT_QOS_GET_HDR(MQTT_QOS1), this->buffer, length - MQTT_MAX_HEADER_SIZE); } return false; } @@ -714,14 +764,9 @@ bool PubSubClient::unsubscribe(const char* topic) { } if (connected()) { uint16_t length = MQTT_MAX_HEADER_SIZE; - nextMsgId++; - if (nextMsgId == 0) { - nextMsgId = 1; - } - this->buffer[length++] = (nextMsgId >> 8); - this->buffer[length++] = (nextMsgId & 0xFF); - length = writeString(topic, this->buffer, length); - return write(MQTTUNSUBSCRIBE | MQTTQOS1, this->buffer, length - MQTT_MAX_HEADER_SIZE); + length = writeNextMsgId(buffer, length, this->bufferSize); // buffer size is checked before + length = writeString(topic, this->buffer, length, this->bufferSize); + return write(MQTTUNSUBSCRIBE | MQTT_QOS_GET_HDR(MQTT_QOS1), this->buffer, length - MQTT_MAX_HEADER_SIZE); } return false; } diff --git a/lib/pubsubclient3/src/PubSubClient.h b/lib/pubsubclient3/src/PubSubClient.h index db77df9e9..2e8842997 100644 --- a/lib/pubsubclient3/src/PubSubClient.h +++ b/lib/pubsubclient3/src/PubSubClient.h @@ -1,9 +1,12 @@ -/* - PubSubClient.h - A simple client for MQTT. - Nick O'Leary, Holger Mueller - http://knolleary.net - https://github.com/hmueller01/pubsubclient3 -*/ +/** + * @file PubSubClient.h + * @brief A simple client for MQTT. + * @author Nicholas O'Leary - http://knolleary.net + * @author Holger Mueller - https://github.com/hmueller01/pubsubclient3 + * @copyright MIT License 2008-2025 + * + * This file is part of the PubSubClient library. + */ #ifndef PubSubClient_h #define PubSubClient_h @@ -14,20 +17,22 @@ #include "IPAddress.h" #include "Stream.h" -#define MQTT_VERSION_3_1 3 -#define MQTT_VERSION_3_1_1 4 +#define MQTT_VERSION_3_1 3 ///< Defines MQTT 3.1 protocol version, see #MQTT_VERSION +#define MQTT_VERSION_3_1_1 4 ///< Defines MQTT 3.1.1 protocol version, see #MQTT_VERSION //< @note The following #define directives can be used to configure the library. /** - * @brief Sets the version of the MQTT protocol to use. - * @note Default value is MQTT_VERSION_3_1_1 for MQTT 3.1.1. + * @brief Sets the version of the MQTT protocol to use (3.1 or 3.1.1). [#MQTT_VERSION_3_1, #MQTT_VERSION_3_1_1]. + * @note Default value is #MQTT_VERSION_3_1_1 for MQTT 3.1.1. */ #ifndef MQTT_VERSION #define MQTT_VERSION MQTT_VERSION_3_1_1 #endif -// MQTT_MAX_POSSIBLE_PACKET_SIZE : Maximum packet size defined by MQTT protocol. +/** + * @brief Maximum packet size defined by MQTT protocol. + */ #ifndef MQTT_MAX_POSSIBLE_PACKET_SIZE #define MQTT_MAX_POSSIBLE_PACKET_SIZE 268435455 #endif @@ -65,23 +70,32 @@ /** * @brief Sets the maximum number of bytes passed to the network client in each write call. * Some hardware has a limit to how much data can be passed to them in one go, - * such as the Arduino Wifi Shield. + * such as the Arduino Wifi Shield e.g. use 80. * @note Defaults to undefined, which passes the entire packet in each write call. */ -// #define MQTT_MAX_TRANSFER_SIZE 80 +#ifndef MQTT_MAX_TRANSFER_SIZE // just a hack that it gets shown in Doxygen +#define MQTT_MAX_TRANSFER_SIZE 80 +#undef MQTT_MAX_TRANSFER_SIZE +#endif -// Possible values for client.state() -#define MQTT_CONNECTION_TIMEOUT -4 -#define MQTT_CONNECTION_LOST -3 -#define MQTT_CONNECT_FAILED -2 -#define MQTT_DISCONNECTED -1 -#define MQTT_CONNECTED 0 -#define MQTT_CONNECT_BAD_PROTOCOL 1 -#define MQTT_CONNECT_BAD_CLIENT_ID 2 -#define MQTT_CONNECT_UNAVAILABLE 3 -#define MQTT_CONNECT_BAD_CREDENTIALS 4 -#define MQTT_CONNECT_UNAUTHORIZED 5 +/** + * @defgroup group_state state() result + * @brief These values indicate the current PubSubClient::state() of the client. + * @{ + */ +#define MQTT_CONNECTION_TIMEOUT -4 ///< The network connection timed out or server didn't respond within the keepalive time. +#define MQTT_CONNECTION_LOST -3 ///< The network connection was lost/broken. +#define MQTT_CONNECT_FAILED -2 ///< The network connection failed. +#define MQTT_DISCONNECTED -1 ///< The client is disconnected cleanly. +#define MQTT_CONNECTED 0 ///< The client is connected. +#define MQTT_CONNECT_BAD_PROTOCOL 1 ///< The server does not support the requested MQTT version. +#define MQTT_CONNECT_BAD_CLIENT_ID 2 ///< The server rejected the client identifier. +#define MQTT_CONNECT_UNAVAILABLE 3 ///< The server was unable to accept the connection. +#define MQTT_CONNECT_BAD_CREDENTIALS 4 ///< The username or password is not valid. +#define MQTT_CONNECT_UNAUTHORIZED 5 ///< The client is not authorized to connect to the server. +/** @} */ +/// \cond #define MQTTRETAINED 1 // Retained flag in the header #define MQTTCONNECT 1 << 4 // Client request to connect to Server #define MQTTCONNACK 2 << 4 // Connect Acknowledgment @@ -98,31 +112,41 @@ #define MQTTPINGRESP 13 << 4 // PING Response #define MQTTDISCONNECT 14 << 4 // Client is Disconnecting #define MQTTRESERVED 15 << 4 // Reserved +/// \endcond -#define MQTTQOS0 (0 << 1) -#define MQTTQOS1 (1 << 1) -#define MQTTQOS2 (2 << 1) +/** + * @defgroup group_qos QoS levels + * @brief Quality of Service (QoS) levels for MQTT messages. + * @{ + */ +#define MQTT_QOS0 ((uint8_t)0) ///< Quality of Service 0: At most once +#define MQTT_QOS1 ((uint8_t)1) ///< Quality of Service 1: At least once +#define MQTT_QOS2 ((uint8_t)2) ///< Quality of Service 2: Exactly once +/// \cond +#define MQTT_QOS_GET_HDR(qos) (((qos) & 0x03) << 1) // Get QoS header bits from QoS value +#define MQTT_HDR_GET_QOS(header) (((header) & 0x06 ) >> 1) // Get QoS value from MQTT header +/// \endcond +/** @} */ -// Maximum size of fixed header and variable length size header +/// \cond Maximum size of fixed header and variable length size header #define MQTT_MAX_HEADER_SIZE 5 +/// \endcond -#if defined(__has_include) && __has_include() && !defined(NOFUNCTIONAL) -#include +/// \anchor callback /** * @brief Define the signature required by any callback function. - * @note The parameters are TOPIC, PAYLOAD, and LENGTH, respectively. + * @param topic The topic of the message. + * @param payload The payload of the message. + * @param plength The length of the payload. */ -#define MQTT_CALLBACK_SIGNATURE std::function callback +#if defined(__has_include) && __has_include() && !defined(NOFUNCTIONAL) +#include +#define MQTT_CALLBACK_SIGNATURE std::function callback #else -#define MQTT_CALLBACK_SIGNATURE void (*callback)(char*, uint8_t*, size_t) +#define MQTT_CALLBACK_SIGNATURE void (*callback)(char* topic, uint8_t* payload, size_t plength) #endif -#define CHECK_STRING_LENGTH(l, s) \ - if ((!s) || (l + 2 + strnlen(s, this->bufferSize) > this->bufferSize)) { \ - _client->stop(); \ - return false; \ - } - +/// \cond #ifdef DEBUG_ESP_PORT #ifdef DEBUG_PUBSUBCLIENT #define DEBUG_PSC_PRINTF(fmt, ...) DEBUG_ESP_PORT.printf(("PubSubClient: " fmt), ##__VA_ARGS__) @@ -143,6 +167,7 @@ #define ERROR_PSC_PRINTF_P(fmt, ...) #endif #endif +/// \endcond /** * @class PubSubClient @@ -165,7 +190,7 @@ class PubSubClient : public Print { uint16_t port{}; Stream* stream{}; int _state{MQTT_DISCONNECTED}; - int _bufferWritePos = 0; + uint8_t _qos{MQTT_QOS0}; size_t readPacket(uint8_t* hdrLen); bool handlePacket(uint8_t hdrLen, size_t len); @@ -173,37 +198,106 @@ class PubSubClient : public Print { bool readByte(uint8_t* result, size_t* pos); uint8_t buildHeader(uint8_t header, uint8_t* buf, size_t length); bool write(uint8_t header, uint8_t* buf, size_t length); - size_t writeString(const char* string, uint8_t* buf, size_t pos); size_t writeString(const char* string, uint8_t* buf, size_t pos, size_t size); - - // Add to buffer and flush if full (only to be used with beginPublish/endPublish)Add commentMore actions - size_t appendBuffer(uint8_t data); - size_t appendBuffer(const uint8_t *data, size_t size); - size_t flushBuffer(); + size_t writeNextMsgId(uint8_t* buf, size_t pos, size_t size); public: /** * @brief Creates an uninitialised client instance. - * @note Before it can be used, - * it must be configured using the property setters setClient and setServer. + * @note Before it can be used, it must be configured using the property setters setClient() and setServer(). */ PubSubClient(); /** * @brief Creates a partially initialised client instance. * @param client The network client to use, for example WiFiClient. - * @note Before it can be used, - * it must be configured with the property setter setServer. + * @note Before it can be used, it must be configured with the property setter setServer(). */ PubSubClient(Client& client); + /** + * @brief Creates a fully configured client instance. + * @param addr The address of the server. + * @param port The port to connect to. + * @param client The network client to use, for example WiFiClient. + */ + PubSubClient(IPAddress addr, uint16_t port, Client& client); + + /** + * @brief Creates a fully configured client instance. + * @param addr The address of the server. + * @param port The port to connect to. + * @param client The network client to use, for example WiFiClient. + * @param stream A stream to write received messages to. + */ + PubSubClient(IPAddress addr, uint16_t port, Client& client, Stream& stream); + + /** + * @brief Creates a fully configured client instance. + * @param addr The address of the server. + * @param port The port to connect to. + * @param callback Pointer to a message \ref callback function. + * Called when a message arrives for a subscription created by this client. + * @param client The network client to use, for example WiFiClient. + */ + PubSubClient(IPAddress addr, uint16_t port, MQTT_CALLBACK_SIGNATURE, Client& client); + + /** + * @brief Creates a fully configured client instance. + * @param addr The address of the server. + * @param port The port to connect to. + * @param callback Pointer to a message \ref callback function. + * Called when a message arrives for a subscription created by this client. + * @param client The network client to use, for example WiFiClient. + * @param stream A stream to write received messages to. + */ + PubSubClient(IPAddress addr, uint16_t port, MQTT_CALLBACK_SIGNATURE, Client& client, Stream& stream); + + /** + * @brief Creates a fully configured client instance. + * @param ip The address of the server. + * @param port The port to connect to. + * @param client The network client to use, for example WiFiClient. + */ + PubSubClient(uint8_t* ip, uint16_t port, Client& client); + + /** + * @brief Creates a fully configured client instance. + * @param ip The address of the server. + * @param port The port to connect to. + * @param client The network client to use, for example WiFiClient. + * @param stream A stream to write received messages to. + */ + PubSubClient(uint8_t* ip, uint16_t port, Client& client, Stream& stream); + + /** + * @brief Creates a fully configured client instance. + * @param ip The address of the server. + * @param port The port to connect to. + * @param callback Pointer to a message \ref callback function. + * Called when a message arrives for a subscription created by this client. + * @param client The network client to use, for example WiFiClient. + */ + PubSubClient(uint8_t* ip, uint16_t port, MQTT_CALLBACK_SIGNATURE, Client& client); + + /** + * @brief Creates a fully configured client instance. + * @param ip The address of the server. + * @param port The port to connect to. + * @param callback Pointer to a message \ref callback function. + * Called when a message arrives for a subscription created by this client. + * @param client The network client to use, for example WiFiClient. + * @param stream A stream to write received messages to. + */ + PubSubClient(uint8_t* ip, uint16_t port, MQTT_CALLBACK_SIGNATURE, Client& client, Stream& stream); + /** * @brief Creates a fully configured client instance. * @param domain The address of the server. * @param port The port to connect to. * @param client The network client to use, for example WiFiClient. */ - PubSubClient(IPAddress, uint16_t, Client& client); + PubSubClient(const char* domain, uint16_t port, Client& client); /** * @brief Creates a fully configured client instance. @@ -212,104 +306,28 @@ class PubSubClient : public Print { * @param client The network client to use, for example WiFiClient. * @param stream A stream to write received messages to. */ - PubSubClient(IPAddress, uint16_t, Client& client, Stream&); + PubSubClient(const char* domain, uint16_t port, Client& client, Stream& stream); /** * @brief Creates a fully configured client instance. * @param domain The address of the server. * @param port The port to connect to. - * @param callback Pointer to a message callback function. + * @param callback Pointer to a message \ref callback function. * Called when a message arrives for a subscription created by this client. * @param client The network client to use, for example WiFiClient. */ - PubSubClient(IPAddress, uint16_t, MQTT_CALLBACK_SIGNATURE, Client& client); + PubSubClient(const char* domain, uint16_t port, MQTT_CALLBACK_SIGNATURE, Client& client); /** * @brief Creates a fully configured client instance. * @param domain The address of the server. * @param port The port to connect to. - * @param callback Pointer to a message callback function. + * @param callback Pointer to a message \ref callback function. * Called when a message arrives for a subscription created by this client. * @param client The network client to use, for example WiFiClient. * @param stream A stream to write received messages to. */ - PubSubClient(IPAddress, uint16_t, MQTT_CALLBACK_SIGNATURE, Client& client, Stream&); - - /** - * @brief Creates a fully configured client instance. - * @param domain The address of the server. - * @param port The port to connect to. - * @param client The network client to use, for example WiFiClient. - */ - PubSubClient(uint8_t*, uint16_t, Client& client); - - /** - * @brief Creates a fully configured client instance. - * @param domain The address of the server. - * @param port The port to connect to. - * @param client The network client to use, for example WiFiClient. - * @param stream A stream to write received messages to. - */ - PubSubClient(uint8_t*, uint16_t, Client& client, Stream&); - - /** - * @brief Creates a fully configured client instance. - * @param domain The address of the server. - * @param port The port to connect to. - * @param callback Pointer to a message callback function. - * Called when a message arrives for a subscription created by this client. - * @param client The network client to use, for example WiFiClient. - */ - PubSubClient(uint8_t*, uint16_t, MQTT_CALLBACK_SIGNATURE, Client& client); - - /** - * @brief Creates a fully configured client instance. - * @param domain The address of the server. - * @param port The port to connect to. - * @param callback Pointer to a message callback function. - * Called when a message arrives for a subscription created by this client. - * @param client The network client to use, for example WiFiClient. - * @param stream A stream to write received messages to. - */ - PubSubClient(uint8_t*, uint16_t, MQTT_CALLBACK_SIGNATURE, Client& client, Stream&); - - /** - * @brief Creates a fully configured client instance. - * @param domain The address of the server. - * @param port The port to connect to. - * @param client The network client to use, for example WiFiClient. - */ - PubSubClient(const char*, uint16_t, Client& client); - - /** - * @brief Creates a fully configured client instance. - * @param domain The address of the server. - * @param port The port to connect to. - * @param client The network client to use, for example WiFiClient. - * @param stream A stream to write received messages to. - */ - PubSubClient(const char*, uint16_t, Client& client, Stream&); - - /** - * @brief Creates a fully configured client instance. - * @param domain The address of the server. - * @param port The port to connect to. - * @param callback Pointer to a message callback function. - * Called when a message arrives for a subscription created by this client. - * @param client The network client to use, for example WiFiClient. - */ - PubSubClient(const char*, uint16_t, MQTT_CALLBACK_SIGNATURE, Client& client); - - /** - * @brief Creates a fully configured client instance. - * @param domain The address of the server. - * @param port The port to connect to. - * @param callback Pointer to a message callback function. - * Called when a message arrives for a subscription created by this client. - * @param client The network client to use, for example WiFiClient. - * @param stream A stream to write received messages to. - */ - PubSubClient(const char*, uint16_t, MQTT_CALLBACK_SIGNATURE, Client& client, Stream&); + PubSubClient(const char* domain, uint16_t port, MQTT_CALLBACK_SIGNATURE, Client& client, Stream& stream); /** * @brief Destructor for the PubSubClient class. @@ -319,7 +337,7 @@ class PubSubClient : public Print { /** * @brief Sets the server details. * @param ip The address of the server. - * @param port The port to connect to. + * @param port The port to connect to. * @return The client instance, allowing the function to be chained. */ PubSubClient& setServer(IPAddress ip, uint16_t port); @@ -327,22 +345,22 @@ class PubSubClient : public Print { /** * @brief Sets the server details. * @param ip The address of the server. - * @param port The port to connect to. + * @param port The port to connect to. * @return The client instance, allowing the function to be chained. */ PubSubClient& setServer(uint8_t* ip, uint16_t port); /** * @brief Sets the server details. - * @param ip The address of the server. - * @param port The port to connect to. + * @param domain The address of the server. + * @param port The port to connect to. * @return The client instance, allowing the function to be chained. */ PubSubClient& setServer(const char* domain, uint16_t port); /** * @brief Sets the message callback function. - * @param callback Pointer to a message callback function. + * @param callback Pointer to a message \ref callback function. * Called when a message arrives for a subscription created by this client. * @return The client instance, allowing the function to be chained. */ @@ -398,7 +416,7 @@ class PubSubClient : public Print { size_t getBufferSize(); /** - * @brief Connects the client. + * @brief Connects the client using a clean session without username and password. * @param id The client ID to use when connecting to the server. * @return true If client succeeded in establishing a connection to the broker. * false If client failed to establish a connection to the broker. @@ -406,48 +424,60 @@ class PubSubClient : public Print { bool connect(const char* id); /** - * @brief Connects the client. + * @brief Connects the client using a clean session with username and password. * @param id The client ID to use when connecting to the server. - * @param user The username to use. If NULL, no username or password is used. - * @param pass The password to use. If NULL, no password is used. + * @param user The username to use. + * @param pass The password to use. + * @note If **user** is NULL, no username or password is used. + * @note If **pass** is NULL, no password is used. * @return true If client succeeded in establishing a connection to the broker. * false If client failed to establish a connection to the broker. */ bool connect(const char* id, const char* user, const char* pass); /** - * @brief Connects the client. + * @brief Connects the client using a clean session and will. * @param id The client ID to use when connecting to the server. * @param willTopic The topic to be used by the will message. * @param willQos The quality of service to be used by the will message. [0, 1, 2]. * @param willRetain Publish the will message with the retain flag. + * @param willMessage The message to be used by the will message. + * @note If **willTopic** is NULL, no will message is sent. * @return true If client succeeded in establishing a connection to the broker. * false If client failed to establish a connection to the broker. */ bool connect(const char* id, const char* willTopic, uint8_t willQos, bool willRetain, const char* willMessage); /** - * @brief Connects the client. + * @brief Connects the client using a clean session with username, password and will. * @param id The client ID to use when connecting to the server. - * @param user The username to use. If NULL, no username or password is used. - * @param pass The password to use. If NULL, no password is used. + * @param user The username to use. + * @param pass The password to use. * @param willTopic The topic to be used by the will message. * @param willQos The quality of service to be used by the will message. [0, 1, 2]. * @param willRetain Publish the will message with the retain flag. + * @param willMessage The message to be used by the will message. + * @note If **user** is NULL, no username or password is used. + * @note If **pass** is NULL, no password is used. + * @note If **willTopic** is NULL, no will message is sent. * @return true If client succeeded in establishing a connection to the broker. * false If client failed to establish a connection to the broker. */ bool connect(const char* id, const char* user, const char* pass, const char* willTopic, uint8_t willQos, bool willRetain, const char* willMessage); /** - * @brief Connects the client. + * @brief Connects the client with all possible parameters (user, password, will and session). * @param id The client ID to use when connecting to the server. - * @param user The username to use. If NULL, no username or password is used. - * @param pass The password to use. If NULL, no password is used. + * @param user The username to use. + * @param pass The password to use. * @param willTopic The topic to be used by the will message. * @param willQos The quality of service to be used by the will message. [0, 1, 2]. * @param willRetain Publish the will message with the retain flag. - * @param cleanSession Connect with a clean session. + * @param willMessage The message to be used by the will message. + * @param cleanSession True to connect with a clean session. + * @note If **user** is NULL, no username or password is used. + * @note If **pass** is NULL, no password is used. + * @note If **willTopic** is NULL, no will message is sent. * @return true If client succeeded in establishing a connection to the broker. * false If client failed to establish a connection to the broker. */ @@ -460,7 +490,7 @@ class PubSubClient : public Print { void disconnect(); /** - * @brief Publishes a non retained message to the specified topic. + * @brief Publishes a non retained message to the specified topic using QoS 0. * @param topic The topic to publish to. * @param payload The message to publish. * @return true If the publish succeeded. @@ -469,7 +499,7 @@ class PubSubClient : public Print { bool publish(const char* topic, const char* payload); /** - * @brief Publishes a message to the specified topic. + * @brief Publishes a message to the specified topic using QoS 0. * @param topic The topic to publish to. * @param payload The message to publish. * @param retained Publish the message with the retain flag. @@ -479,7 +509,18 @@ class PubSubClient : public Print { bool publish(const char* topic, const char* payload, bool retained); /** - * @brief Publishes a non retained message to the specified topic. + * @brief Publishes a message to the specified topic. + * @param topic The topic to publish to. + * @param payload The message to publish. + * @param qos The quality of service (\ref group_qos) to publish at. [0, 1, 2]. + * @param retained Publish the message with the retain flag. + * @return true If the publish succeeded. + * false If the publish failed, either connection lost or message too large. + */ + bool publish(const char* topic, const char* payload, uint8_t qos, bool retained); + + /** + * @brief Publishes a non retained message to the specified topic using QoS 0. * @param topic The topic to publish to. * @param payload The message to publish. * @param plength The length of the payload. @@ -489,7 +530,7 @@ class PubSubClient : public Print { bool publish(const char* topic, const uint8_t* payload, size_t plength); /** - * @brief Publishes a message to the specified topic. + * @brief Publishes a message to the specified topic using QoS 0. * @param topic The topic to publish to. * @param payload The message to publish. * @param plength The length of the payload. @@ -500,7 +541,19 @@ class PubSubClient : public Print { bool publish(const char* topic, const uint8_t* payload, size_t plength, bool retained); /** - * @brief Publishes a message stored in PROGMEM to the specified topic. + * @brief Publishes a message to the specified topic. + * @param topic The topic to publish to. + * @param payload The message to publish. + * @param plength The length of the payload. + * @param qos The quality of service (\ref group_qos) to publish at. [0, 1, 2]. + * @param retained Publish the message with the retain flag. + * @return true If the publish succeeded. + * false If the publish failed, either connection lost or message too large. + */ + bool publish(const char* topic, const uint8_t* payload, size_t plength, uint8_t qos, bool retained); + + /** + * @brief Publishes a message stored in PROGMEM to the specified topic using QoS 0. * @param topic The topic to publish to. * @param payload The message to publish. * @param retained Publish the message with the retain flag. @@ -513,6 +566,17 @@ class PubSubClient : public Print { * @brief Publishes a message stored in PROGMEM to the specified topic. * @param topic The topic to publish to. * @param payload The message to publish. + * @param qos The quality of service (\ref group_qos) to publish at. [0, 1, 2]. + * @param retained Publish the message with the retain flag. + * @return true If the publish succeeded. + * false If the publish failed, either connection lost or message too large. + */ + bool publish_P(const char* topic, const char* payload, uint8_t qos, bool retained); + + /** + * @brief Publishes a message stored in PROGMEM to the specified topic using QoS 0. + * @param topic The topic to publish to. + * @param payload The message to publish. * @param plength The length of the payload. * @param retained Publish the message with the retain flag. * @return true If the publish succeeded. @@ -521,7 +585,19 @@ class PubSubClient : public Print { bool publish_P(const char* topic, const uint8_t* payload, size_t plength, bool retained); /** - * @brief Start to publish a message. + * @brief Publishes a message stored in PROGMEM to the specified topic. + * @param topic The topic to publish to. + * @param payload The message to publish. + * @param plength The length of the payload. + * @param qos The quality of service (\ref group_qos) to publish at. [0, 1, 2]. + * @param retained Publish the message with the retain flag. + * @return true If the publish succeeded. + * false If the publish failed, either connection lost or message too large. + */ + bool publish_P(const char* topic, const uint8_t* payload, size_t plength, uint8_t qos, bool retained); + + /** + * @brief Start to publish a message using QoS 0. * This API: * beginPublish(...) * one or more calls to write(...) @@ -536,6 +612,23 @@ class PubSubClient : public Print { */ bool beginPublish(const char* topic, size_t plength, bool retained); + /** + * @brief Start to publish a message. + * This API: + * beginPublish(...) + * one or more calls to write(...) + * endPublish() + * Allows for arbitrarily large payloads to be sent without them having to be copied into + * a new buffer and held in memory at one time. + * @param topic The topic to publish to. + * @param plength The length of the payload. + * @param qos The quality of service (\ref group_qos) to publish at. [0, 1, 2]. + * @param retained Publish the message with the retain flag. + * @return true If the publish succeeded. + * false If the publish failed, either connection lost or message too large. + */ + bool beginPublish(const char* topic, size_t plength, uint8_t qos, bool retained); + /** * @brief Finish sending a message that was started with a call to beginPublish. * @return true If the publish succeeded. @@ -545,17 +638,13 @@ class PubSubClient : public Print { /** * @brief Writes a single byte as a component of a publish started with a call to beginPublish. - * For performance reasons, this will be appended to the internal buffer, - * which will be flushed when full or on a call to endPublish(). - * @param byte A byte to write to the publish payload. + * @param data A byte to write to the publish payload. * @return The number of bytes written. */ - virtual size_t write(uint8_t); + virtual size_t write(uint8_t data); /** * @brief Writes an array of bytes as a component of a publish started with a call to beginPublish. - * For performance reasons, this will be appended to the internal buffer, - * which will be flushed when full or on a call to endPublish(). * @param buffer The bytes to write. * @param size The length of the payload to be sent. * @return The number of bytes written. @@ -563,7 +652,7 @@ class PubSubClient : public Print { virtual size_t write(const uint8_t* buffer, size_t size); /** - * @brief Subscribes to messages published to the specified topic. + * @brief Subscribes to messages published to the specified topic using QoS 0. * @param topic The topic to subscribe to. * @return true If sending the subscribe succeeded. * false If sending the subscribe failed, either connection lost or message too large. @@ -605,16 +694,7 @@ class PubSubClient : public Print { * @brief Returns the current state of the client. * If a connection attempt fails, this can be used to get more information about the failure. * @note All of the values have corresponding constants defined in PubSubClient.h. - * @return -4 : MQTT_CONNECTION_TIMEOUT - The server didn't respond within the keepalive time. - * -3 : MQTT_CONNECTION_LOST - The network connection was broken. - * -2 : MQTT_CONNECT_FAILED - The network connection failed. - * -1 : MQTT_DISCONNECTED - The client is disconnected cleanly. - * 0 : MQTT_CONNECTED - The client is connected. - * 1 : MQTT_CONNECT_BAD_PROTOCOL - The server doesn't support the requested version of MQTT. - * 2 : MQTT_CONNECT_BAD_CLIENT_ID - The server rejected the client identifier. - * 3 : MQTT_CONNECT_UNAVAILABLE - The server was unable to accept the connection. - * 4 : MQTT_CONNECT_BAD_CREDENTIALS - The username/password were rejected. - * 5 : MQTT_CONNECT_UNAUTHORIZED - The client was not authorized to connect. + * @return See \ref group_state */ int state(); }; diff --git a/lib/pubsubclient3/tests/src/connect_spec.cpp b/lib/pubsubclient3/tests/src/connect_spec.cpp index 75dc819b5..b272e92f5 100644 --- a/lib/pubsubclient3/tests/src/connect_spec.cpp +++ b/lib/pubsubclient3/tests/src/connect_spec.cpp @@ -7,7 +7,7 @@ byte server[] = {172, 16, 0, 2}; // function declarations -void callback(char* topic, uint8_t* payload, size_t length); +void callback(char* topic, uint8_t* payload, size_t plength); int test_connect_fails_no_network(); int test_connect_fails_on_no_response(); int test_connect_properly_formatted(); @@ -23,7 +23,7 @@ int test_connect_with_will_username_password(); int test_connect_disconnect_connect(); int test_connect_custom_keepalive(); -void callback(_UNUSED_ char* topic, _UNUSED_ uint8_t* payload, _UNUSED_ size_t length) { +void callback(_UNUSED_ char* topic, _UNUSED_ uint8_t* payload, _UNUSED_ size_t plength) { // handle message arrived } diff --git a/lib/pubsubclient3/tests/src/keepalive_spec.cpp b/lib/pubsubclient3/tests/src/keepalive_spec.cpp index 58e3f206c..37b99be72 100644 --- a/lib/pubsubclient3/tests/src/keepalive_spec.cpp +++ b/lib/pubsubclient3/tests/src/keepalive_spec.cpp @@ -9,14 +9,14 @@ byte server[] = {172, 16, 0, 2}; // function declarations -void callback(char* topic, uint8_t* payload, size_t length); +void callback(char* topic, uint8_t* payload, size_t plength); int test_keepalive_pings_idle(); int test_keepalive_pings_with_outbound_qos0(); int test_keepalive_pings_with_inbound_qos0(); int test_keepalive_no_pings_inbound_qos1(); int test_keepalive_disconnects_hung(); -void callback(_UNUSED_ char* topic, _UNUSED_ uint8_t* payload, _UNUSED_ size_t length) { +void callback(_UNUSED_ char* topic, _UNUSED_ uint8_t* payload, _UNUSED_ size_t plength) { // handle message arrived } @@ -54,7 +54,7 @@ int test_keepalive_pings_idle() { } int test_keepalive_pings_with_outbound_qos0() { - IT("keeps a connection alive that only sends qos0 (takes 1 minute)"); + IT("keeps a connection alive that only sends QoS 0 (takes 1 minute)"); ShimClient shimClient; shimClient.setAllowConnect(true); @@ -90,7 +90,7 @@ int test_keepalive_pings_with_outbound_qos0() { } int test_keepalive_pings_with_inbound_qos0() { - IT("keeps a connection alive that only receives qos0 (takes 1 minute)"); + IT("keeps a connection alive that only receives QoS 0 (takes 1 minute)"); ShimClient shimClient; shimClient.setAllowConnect(true); @@ -123,7 +123,7 @@ int test_keepalive_pings_with_inbound_qos0() { } int test_keepalive_no_pings_inbound_qos1() { - IT("does not send pings for connections with inbound qos1 (takes 1 minute)"); + IT("does not send pings for connections with inbound QoS 1 (takes 1 minute)"); ShimClient shimClient; shimClient.setAllowConnect(true); diff --git a/lib/pubsubclient3/tests/src/publish_spec.cpp b/lib/pubsubclient3/tests/src/publish_spec.cpp index 4346a721a..241ee57e2 100644 --- a/lib/pubsubclient3/tests/src/publish_spec.cpp +++ b/lib/pubsubclient3/tests/src/publish_spec.cpp @@ -7,7 +7,7 @@ byte server[] = {172, 16, 0, 2}; // function declarations -void callback(char* topic, uint8_t* payload, size_t length); +void callback(char* topic, uint8_t* payload, size_t plength); int test_publish(); int test_publish_bytes(); int test_publish_retained(); @@ -16,8 +16,14 @@ int test_publish_not_connected(); int test_publish_too_long(); int test_publish_P(); int test_publish_P_too_long(); +int test_publish_empty_topic(); +int test_publish_null_payload(); +int test_publish_qos1(); +int test_publish_qos2(); +int test_publish_P_qos1(); +int test_publish_P_qos2(); -void callback(_UNUSED_ char* topic, _UNUSED_ uint8_t* payload, _UNUSED_ size_t length) { +void callback(_UNUSED_ char* topic, _UNUSED_ uint8_t* payload, _UNUSED_ size_t plength) { // handle message arrived } @@ -27,14 +33,14 @@ int test_publish() { shimClient.setAllowConnect(true); byte connack[] = {0x20, 0x02, 0x00, 0x00}; - shimClient.respond(connack, 4); + shimClient.respond(connack, sizeof(connack)); PubSubClient client(server, 1883, callback, shimClient); bool rc = client.connect("client_test1"); IS_TRUE(rc); - byte publish[] = {0x30, 0xe, 0x0, 0x5, 0x74, 0x6f, 0x70, 0x69, 0x63, 0x70, 0x61, 0x79, 0x6c, 0x6f, 0x61, 0x64}; - shimClient.expect(publish, 16); + byte publish[] = {0x30, 0x0e, 0x00, 0x05, 't', 'o', 'p', 'i', 'c', 'p', 'a', 'y', 'l', 'o', 'a', 'd'}; + shimClient.expect(publish, sizeof(publish)); rc = client.publish("topic", "payload"); IS_TRUE(rc); @@ -49,18 +55,18 @@ int test_publish_bytes() { ShimClient shimClient; shimClient.setAllowConnect(true); - byte payload[] = {0x01, 0x02, 0x03, 0x0, 0x05}; - size_t length = 5; + byte payload[] = {0x01, 0x02, 0x03, 0x00, 0x05}; + size_t length = sizeof(payload); byte connack[] = {0x20, 0x02, 0x00, 0x00}; - shimClient.respond(connack, 4); + shimClient.respond(connack, sizeof(connack)); PubSubClient client(server, 1883, callback, shimClient); bool rc = client.connect("client_test1"); IS_TRUE(rc); - byte publish[] = {0x30, 0xc, 0x0, 0x5, 0x74, 0x6f, 0x70, 0x69, 0x63, 0x1, 0x2, 0x3, 0x0, 0x5}; - shimClient.expect(publish, 14); + byte publish[] = {0x30, 0x0c, 0x00, 0x05, 't', 'o', 'p', 'i', 'c', 0x01, 0x02, 0x03, 0x00, 0x05}; + shimClient.expect(publish, sizeof(publish)); rc = client.publish("topic", payload, length); IS_TRUE(rc); @@ -75,18 +81,18 @@ int test_publish_retained() { ShimClient shimClient; shimClient.setAllowConnect(true); - byte payload[] = {0x01, 0x02, 0x03, 0x0, 0x05}; - size_t length = 5; + byte payload[] = {0x01, 0x02, 0x03, 0x00, 0x05}; + size_t length = sizeof(payload); byte connack[] = {0x20, 0x02, 0x00, 0x00}; - shimClient.respond(connack, 4); + shimClient.respond(connack, sizeof(connack)); PubSubClient client(server, 1883, callback, shimClient); bool rc = client.connect("client_test1"); IS_TRUE(rc); - byte publish[] = {0x31, 0xc, 0x0, 0x5, 0x74, 0x6f, 0x70, 0x69, 0x63, 0x1, 0x2, 0x3, 0x0, 0x5}; - shimClient.expect(publish, 14); + byte publish[] = {0x31, 0x0c, 0x00, 0x05, 't', 'o', 'p', 'i', 'c', 0x1, 0x2, 0x3, 0x0, 0x5}; + shimClient.expect(publish, sizeof(publish)); rc = client.publish("topic", payload, length, true); IS_TRUE(rc); @@ -102,14 +108,14 @@ int test_publish_retained_2() { shimClient.setAllowConnect(true); byte connack[] = {0x20, 0x02, 0x00, 0x00}; - shimClient.respond(connack, 4); + shimClient.respond(connack, sizeof(connack)); PubSubClient client(server, 1883, callback, shimClient); bool rc = client.connect("client_test1"); IS_TRUE(rc); - byte publish[] = {0x31, 0xc, 0x0, 0x5, 0x74, 0x6f, 0x70, 0x69, 0x63, 'A', 'B', 'C', 'D', 'E'}; - shimClient.expect(publish, 14); + byte publish[] = {0x31, 0x0c, 0x00, 0x05, 't', 'o', 'p', 'i', 'c', 'A', 'B', 'C', 'D', 'E'}; + shimClient.expect(publish, sizeof(publish)); rc = client.publish("topic", "ABCDE", true); IS_TRUE(rc); @@ -165,18 +171,18 @@ int test_publish_P() { ShimClient shimClient; shimClient.setAllowConnect(true); - byte payload[] = {0x01, 0x02, 0x03, 0x0, 0x05}; - size_t length = 5; + byte payload[] = {0x01, 0x02, 0x03, 0x00, 0x05}; + size_t length = sizeof(payload); byte connack[] = {0x20, 0x02, 0x00, 0x00}; - shimClient.respond(connack, 4); + shimClient.respond(connack, sizeof(connack)); PubSubClient client(server, 1883, callback, shimClient); bool rc = client.connect("client_test1"); IS_TRUE(rc); - byte publish[] = {0x31, 0xc, 0x0, 0x5, 0x74, 0x6f, 0x70, 0x69, 0x63, 0x1, 0x2, 0x3, 0x0, 0x5}; - shimClient.expect(publish, 14); + byte publish[] = {0x31, 0x0c, 0x00, 0x05, 't', 'o', 'p', 'i', 'c', 0x01, 0x02, 0x03, 0x00, 0x05}; + shimClient.expect(publish, sizeof(publish)); rc = client.publish_P("topic", payload, length, true); IS_TRUE(rc); @@ -187,7 +193,7 @@ int test_publish_P() { } int test_publish_P_too_long() { - IT("publish using PROGMEM fails when topic are too long"); + IT("publish using PROGMEM fails when topic is too long"); ShimClient shimClient; shimClient.setAllowConnect(true); @@ -198,7 +204,7 @@ int test_publish_P_too_long() { char payload[] = "123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890"; byte connack[] = {0x20, 0x02, 0x00, 0x00}; - shimClient.respond(connack, 4); + shimClient.respond(connack, sizeof(connack)); PubSubClient client(server, 1883, callback, shimClient); client.setBufferSize(64); @@ -213,15 +219,164 @@ int test_publish_P_too_long() { END_IT } +int test_publish_empty_topic() { + IT("publish fails when topic is empty"); + ShimClient shimClient; + shimClient.setAllowConnect(true); + + byte connack[] = {0x20, 0x02, 0x00, 0x00}; + shimClient.respond(connack, 4); + + PubSubClient client(server, 1883, callback, shimClient); + bool rc = client.connect("client_test1"); + IS_TRUE(rc); + + rc = client.publish("", "payload"); + IS_FALSE(rc); + + IS_FALSE(shimClient.error()); + + END_IT +} + +int test_publish_null_payload() { + IT("publishes with null payload"); + ShimClient shimClient; + shimClient.setAllowConnect(true); + + byte connack[] = {0x20, 0x02, 0x00, 0x00}; + shimClient.respond(connack, sizeof(connack)); + + PubSubClient client(server, 1883, callback, shimClient); + bool rc = client.connect("client_test1"); + IS_TRUE(rc); + + byte publish[] = {0x30, 0x07, 0x00, 0x05, 't', 'o', 'p', 'i', 'c'}; + shimClient.expect(publish, sizeof(publish)); + + rc = client.publish("topic", nullptr); + IS_TRUE(rc); + + IS_FALSE(shimClient.error()); + + END_IT +} + +int test_publish_qos1() { + IT("publishes with QoS 1"); + ShimClient shimClient; + shimClient.setAllowConnect(true); + + byte connack[] = {0x20, 0x02, 0x00, 0x00}; + shimClient.respond(connack, 4); + + PubSubClient client(server, 1883, callback, shimClient); + bool rc = client.connect("client_test1"); + IS_TRUE(rc); + + // Example publish packet for QoS 1 (0x32) + byte publish[] = {0x32, 0x10, 0x00, 0x05, 't', 'o', 'p', 'i', 'c', 'p', 'a', 'y', 'l', 'o', 'a', 'd', 0x00, 0x02}; + shimClient.expect(publish, sizeof(publish)); + + rc = client.publish("topic", "payload", MQTT_QOS1, false); + IS_TRUE(rc); + + IS_FALSE(shimClient.error()); + + END_IT +} + +int test_publish_qos2() { + IT("publishes with QoS 2"); + ShimClient shimClient; + shimClient.setAllowConnect(true); + + byte connack[] = {0x20, 0x02, 0x00, 0x00}; + shimClient.respond(connack, 4); + + PubSubClient client(server, 1883, callback, shimClient); + bool rc = client.connect("client_test1"); + IS_TRUE(rc); + + // Example publish packet for QoS 2 (0x34) + byte publish[] = {0x34, 0x10, 0x00, 0x05, 't', 'o', 'p', 'i', 'c', 'p', 'a', 'y', 'l', 'o', 'a', 'd', 0x00, 0x02}; + shimClient.expect(publish, sizeof(publish)); + + rc = client.publish("topic", "payload", MQTT_QOS2, false); + IS_TRUE(rc); + + IS_FALSE(shimClient.error()); + + END_IT +} + +int test_publish_P_qos1() { + IT("publishes using PROGMEM with QoS 1 retained"); + ShimClient shimClient; + shimClient.setAllowConnect(true); + + byte payload[] = {0x01, 0x02, 0x03, 0x04, 0x05}; + size_t length = sizeof(payload); + + byte connack[] = {0x20, 0x02, 0x00, 0x00}; + shimClient.respond(connack, sizeof(connack)); + + PubSubClient client(server, 1883, callback, shimClient); + bool rc = client.connect("client_test1"); + IS_TRUE(rc); + + byte publish[] = {0x33, 0x0e, 0x00, 0x05, 't', 'o', 'p', 'i', 'c', 0x01, 0x02, 0x03, 0x04, 0x05, 0x00, 0x02}; + shimClient.expect(publish, sizeof(publish)); + + rc = client.publish_P("topic", payload, length, MQTT_QOS1, true); + IS_TRUE(rc); + + IS_FALSE(shimClient.error()); + + END_IT +} + +int test_publish_P_qos2() { + IT("publishes using PROGMEM with QoS 2 retained"); + ShimClient shimClient; + shimClient.setAllowConnect(true); + + byte payload[] = {0x01, 0x02, 0x03, 0x04, 0x05}; + size_t length = sizeof(payload); + + byte connack[] = {0x20, 0x02, 0x00, 0x00}; + shimClient.respond(connack, sizeof(connack)); + + PubSubClient client(server, 1883, callback, shimClient); + bool rc = client.connect("client_test1"); + IS_TRUE(rc); + + byte publish[] = {0x35, 0x0e, 0x00, 0x05, 't', 'o', 'p', 'i', 'c', 0x01, 0x02, 0x03, 0x04, 0x05, 0x00, 0x02}; + shimClient.expect(publish, sizeof(publish)); + + rc = client.publish_P("topic", payload, length, MQTT_QOS2, true); + IS_TRUE(rc); + + IS_FALSE(shimClient.error()); + + END_IT +} + int main() { SUITE("Publish"); test_publish(); test_publish_bytes(); test_publish_retained(); test_publish_retained_2(); + test_publish_qos1(); + test_publish_qos2(); + test_publish_null_payload(); test_publish_not_connected(); + test_publish_empty_topic(); test_publish_too_long(); test_publish_P(); + test_publish_P_qos1(); + test_publish_P_qos2(); test_publish_P_too_long(); FINISH diff --git a/lib/pubsubclient3/tests/src/receive_spec.cpp b/lib/pubsubclient3/tests/src/receive_spec.cpp index 3635e9778..8cbe0efae 100644 --- a/lib/pubsubclient3/tests/src/receive_spec.cpp +++ b/lib/pubsubclient3/tests/src/receive_spec.cpp @@ -18,7 +18,7 @@ char lastPayload[1024]; size_t lastLength; // function declarations -void callback(char* topic, uint8_t* payload, size_t length); +void callback(char* topic, uint8_t* payload, size_t plength); void reset_callback(); int test_receive_callback(); int test_receive_stream(); @@ -36,12 +36,12 @@ void reset_callback() { lastLength = 0; } -void callback(char* topic, uint8_t* payload, size_t length) { - TRACE("Callback received topic=[" << topic << "] length=" << length << "\n") +void callback(char* topic, uint8_t* payload, size_t plength) { + TRACE("Callback received topic=[" << topic << "] plength=" << plength << "\n") callback_called = true; strcpy(lastTopic, topic); - memcpy(lastPayload, payload, length); - lastLength = length; + memcpy(lastPayload, payload, plength); + lastLength = plength; } int test_receive_callback() { diff --git a/lib/pubsubclient3/tests/src/subscribe_spec.cpp b/lib/pubsubclient3/tests/src/subscribe_spec.cpp index 8ae83e7ad..c00eddf75 100644 --- a/lib/pubsubclient3/tests/src/subscribe_spec.cpp +++ b/lib/pubsubclient3/tests/src/subscribe_spec.cpp @@ -7,7 +7,7 @@ byte server[] = {172, 16, 0, 2}; // function declarations -void callback(char* topic, uint8_t* payload, size_t length); +void callback(char* topic, uint8_t* payload, size_t plength); int test_subscribe_no_qos(); int test_subscribe_qos_1(); int test_subscribe_not_connected(); @@ -16,12 +16,12 @@ int test_subscribe_too_long(); int test_unsubscribe(); int test_unsubscribe_not_connected(); -void callback(_UNUSED_ char* topic, _UNUSED_ uint8_t* payload, _UNUSED_ size_t length) { +void callback(_UNUSED_ char* topic, _UNUSED_ uint8_t* payload, _UNUSED_ size_t plength) { // handle message arrived } int test_subscribe_no_qos() { - IT("subscribe without qos defaults to 0"); + IT("subscribe without QoS (defaults to 0)"); ShimClient shimClient; shimClient.setAllowConnect(true); @@ -46,7 +46,7 @@ int test_subscribe_no_qos() { } int test_subscribe_qos_1() { - IT("subscribes qos 1"); + IT("subscribe with QoS 1"); ShimClient shimClient; shimClient.setAllowConnect(true); @@ -62,7 +62,7 @@ int test_subscribe_qos_1() { byte suback[] = {0x90, 0x3, 0x0, 0x2, 0x1}; shimClient.respond(suback, 5); - rc = client.subscribe("topic", 1); + rc = client.subscribe("topic", MQTT_QOS1); IS_TRUE(rc); IS_FALSE(shimClient.error()); @@ -85,7 +85,7 @@ int test_subscribe_not_connected() { } int test_subscribe_invalid_qos() { - IT("subscribe fails with invalid qos values"); + IT("subscribe fails with invalid QoS values"); ShimClient shimClient; shimClient.setAllowConnect(true); @@ -96,7 +96,7 @@ int test_subscribe_invalid_qos() { bool rc = client.connect("client_test1"); IS_TRUE(rc); - rc = client.subscribe("topic", 2); + rc = client.subscribe("topic", MQTT_QOS2); IS_FALSE(rc); rc = client.subscribe("topic", 254); IS_FALSE(rc); @@ -119,13 +119,13 @@ int test_subscribe_too_long() { bool rc = client.connect("client_test1"); IS_TRUE(rc); - // max length should be allowed - // 0 1 2 3 4 5 6 7 8 9 0 1 2 - rc = client.subscribe("12345678901234567890123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456789"); + // max length should be allowed (buffer size - MQTT_MAX_HEADER_SIZE - msgId (2) - topic length bytes (2) - QoS (1) = 118) + // 0 1 2 3 4 5 6 7 8 9 0 1 2 + rc = client.subscribe("1234567890123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456789012345678"); IS_TRUE(rc); - // 0 1 2 3 4 5 6 7 8 9 0 1 2 - rc = client.subscribe("123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890"); + // 0 1 2 3 4 5 6 7 8 9 0 1 2 + rc = client.subscribe("12345678901234567890123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456789"); IS_FALSE(rc); IS_FALSE(shimClient.error());