diff --git a/.gitignore b/.gitignore index 627feb9..28843fb 100644 --- a/.gitignore +++ b/.gitignore @@ -1,2 +1,2 @@ /.vscode -/test \ No newline at end of file +/.pio diff --git a/keywords.txt b/keywords.txt index ca507d6..f86f130 100644 --- a/keywords.txt +++ b/keywords.txt @@ -23,6 +23,10 @@ EndpointHealth KEYWORD1 PowerController KEYWORD1 TemperatureSensorScale KEYWORD1 AlexaUtils KEYWORD1 +AlexaLog KEYWORD1 +AlexaLogLevel KEYWORD1 +AlexaSendResult KEYWORD1 +AlexaTransport KEYWORD1 ####################################### # Methods and Functions (KEYWORD2) @@ -33,9 +37,14 @@ loop KEYWORD2 getState KEYWORD2 getDisconnectReason KEYWORD2 getDevice KEYWORD2 +setLogLevel KEYWORD2 +setTimeSource KEYWORD2 +publish KEYWORD2 +timestamp KEYWORD2 setName KEYWORD2 getName KEYWORD2 getEndpointId KEYWORD2 +hasEndpointId KEYWORD2 setDisplayCategory KEYWORD2 getDisplayCategory KEYWORD2 setDescription KEYWORD2 @@ -72,10 +81,14 @@ AddColorTemperatureControllerProp KEYWORD2 AddToggleControllerProp KEYWORD2 AddContextProp KEYWORD2 send KEYWORD2 +printMemoryInfo KEYWORD2 ####################################### # Constants (LITERAL1) ####################################### -MAX_STATUS_REPORT_SIZE LITERAL1 +ALEX2ESP_MAX_DIRECTIVE LITERAL1 +ALEX2ESP_MAX_MESSAGE LITERAL1 +ALEX2ESP_LOG_MAX LITERAL1 +ALEXA_TIMESTAMP_SIZE LITERAL1 MAX_EVENTS LITERAL1 diff --git a/platformio.ini b/platformio.ini new file mode 100644 index 0000000..9b73647 --- /dev/null +++ b/platformio.ini @@ -0,0 +1,17 @@ +; Development project of the library itself. A sketch does not need this file: it gets the library through +; lib_deps (readme.md, Quick Start). +; +; pio test -e native host tests of the bridge logic in src/AlexaBridgeLogic.cpp; no board, no broker + +[platformio] +default_envs = native + +[env:native] +platform = native +test_framework = unity +; Only the sources without a hardware dependency are compiled with the tests +test_build_src = yes +build_src_filter = -<*> + +lib_deps = + bblanchon/ArduinoJson@^7 +build_flags = -std=gnu++17 -Wall -Wextra diff --git a/readme.md b/readme.md index 4ea9603..1e4782a 100644 --- a/readme.md +++ b/readme.md @@ -135,9 +135,16 @@ An `ActionMapping` takes an optional third argument, the directive payload as JS ## How it talks to Alex2MQTT -- **Discovery.** On `/discover` the library answers with one discovery object per device, published straight to `/discover_r` over MQTT (no HTTP round trip, no queue slot). The backend accepts one endpoint object per message and collects everything that arrives within 1 s for Alexa's discovery answer (up to 5 s for its proactive AddOrUpdate push), so all devices are published back to back the moment the request arrives. Each object sits on the heap (about 1 KB) until the broker acknowledges it; when the MQTT client refuses another one (free heap under 4 KB), the library prints `[Alex2ESP] discovery publish deferred at ` and sends the rest from `loop()` as the queue drains, for up to 5 s after the request. `[Alex2ESP] discovery gave up: N device(s) not announced` means those devices missed this answer - on the backend's proactive discovery that can remove them from Alexa until the next one. -- **Directives.** A directive arrives as a short id on `//alexaDirective_e`; the library fetches the full directive over HTTP, fires `ReportState` or `Event` (and `DirectiveReceived`, if registered), and your status report is queued for an HTTP POST that the backend republishes on `//alexaResponce`, replacing `{REPLACE_WITH_DATETIME}` with the current time. -- **Limits.** The send queue holds 5 reports of up to 2047 bytes each. `send()` returns `false` and prints a `[Alex2ESP]` line on Serial when a report does not fit or the queue is full; nothing is ever sent truncated. Debug logging of the library's internals is compiled in by defining `Alex2ESP_DEBUG` in `AlexaUtils.cpp`; credentials are never printed. +Everything goes over MQTT (port 1883 of `alex2mqtt.stormysdream.club`); the library makes no HTTP requests. + +- **Session.** `begin()` starts SNTP (`pool.ntp.org`, `time.nist.gov`) and returns; `loop()` opens the MQTT session once the clock is set, or after 5 s without an answer, and subscribes to `/discover` and `/+/alexaDirective`. `getState()` is `CONNECTED` when the broker has acknowledged both subscriptions. A sketch that sets the clock itself (its own `configTime()` with a time zone, an RTC) calls `alexClient.setTimeSource(false)` before `begin()`. +- **Discovery.** On `/discover` the library answers with one discovery object per device on `/discover_r`. The backend accepts one endpoint object per message and collects everything that arrives within 1 s for Alexa's discovery answer (up to 5 s for its proactive AddOrUpdate push), so all devices are published back to back from the next `loop()`. Each object sits on the heap (about 1 KB) until the MQTT client has sent it; when the client cannot take another one (free heap under 4 KB), the library prints `[Alex2ESP] discovery deferred at ` and sends the rest from `loop()` as the queue drains, for up to 5 s after the request. `[Alex2ESP] error: discovery gave up: N device(s) not announced` means those devices missed this answer - on the backend's proactive discovery that can remove them from Alexa until the next one. +- **Directives.** The directive arrives as JSON on `//alexaDirective`. A directive larger than one TCP segment arrives in fragments, which are put together in one heap block that exists only until `loop()` has parsed it. `loop()` then fires `ReportState` or `Event` (and `DirectiveReceived`, if registered) with the directive: `directive["header"]`, `directive["endpoint"]`, `directive["payload"]`. One directive is handled at a time; a directive that arrives twice (a broker that mirrors its topics delivers every message twice) is handled once. +- **Reports.** `send()` publishes the report on `//alexaResponce` at once. The backend waits 7 s for it, so answer from the event handler. Every property carries the board's UTC time as `timeOfSample`. +- **Limits.** A directive and a report (or the discovery object of one device) may be 2047 bytes each; `-DALEX2ESP_MAX_DIRECTIVE=` and `-DALEX2ESP_MAX_MESSAGE=` in `build_flags` change that. `send()` returns `false` when the report was not sent: no session with the broker, the MQTT client or the heap cannot take it, or it is too large. Nothing is ever sent truncated, and nothing is dropped without a line on Serial. +- **Serial output.** Every line of the library starts with `[Alex2ESP]`, a problem with `[Alex2ESP] error:`. `alexClient.setLogLevel(AlexaLogLevel::ERROR)` leaves only the problems, `AlexaLogLevel::NONE` nothing; the default, `AlexaLogLevel::INFO`, adds the session, discovery and one line per directive (`[Alex2ESP] ESP-01 <- Alexa.PowerController.TurnOn`). `AlexaLogLevel::DEBUG` (sizes and free heap per message) has to be compiled in with `-DALEX2ESP_LOG_MAX=3`; `-DALEX2ESP_LOG_MAX=0` compiles every line out. Credentials and correlation tokens are never printed. + +Boards that run 1.1.0 or older keep working: the backend still publishes the token on `//alexaDirective_e` and serves the HTTP routes they use. --- @@ -222,15 +229,38 @@ For instance, to manually report the state of a PowerController, you can use the doc["namespace"] = "Alexa.PowerController"; doc["name"] = "powerState"; doc["value"] = "ON"; - doc["timeOfSample"] = "{REPLACE_WITH_DATETIME}"; // Alex2MQTT server wil do the replace doc["uncertaintyInMilliseconds"] = 0; AddContextProp(doc.as()) ``` +`AddContextProp` sets `timeOfSample` to the current time when the object has none. (Until 1.1.0 a sketch wrote `"{REPLACE_WITH_DATETIME}"` there for the server to replace; the library now replaces that itself.) + --- ## Changelog +### 1.2.0 + +Behaviour changes: +- Directives arrive over MQTT. The library subscribes to `/+/alexaDirective`, where Alex2MQTT has always published the whole directive, instead of fetching it over HTTP with the token from `//alexaDirective_e`. The HTTP detour dates from 2024, when a directive larger than one TCP segment reached the MQTT callback in pieces; the pieces are now put together by their offset and the total length, in one heap block that lives until `loop()` has parsed the directive. `loop()` no longer stalls for two HTTP round trips per directive. +- Reports leave over MQTT. `send()` publishes on `//alexaResponce` at once, where 1.1.0 queued the report for an HTTP POST from a later `loop()`. It returns `false` when there is no session with the broker, when the MQTT client or the heap cannot take the report, or when the report is over 2047 bytes (`ALEX2ESP_MAX_MESSAGE`); each case prints its reason. The 5-slot send queue is gone. +- `timeOfSample` is the board's own time in UTC, for example `2026-09-28T13:05:09Z`. 1.1.0 sent the placeholder `{REPLACE_WITH_DATETIME}`, which only the backend's HTTP route replaced. `begin()` starts SNTP (`pool.ntp.org`, `time.nist.gov`) and no longer connects itself: `loop()` opens the MQTT session once the clock is set, or after 5 s without an answer, so the session comes up a few seconds later than before. `alexClient.setTimeSource(false)` before `begin()` leaves the clock to the sketch. `AddContextProp()` fills `timeOfSample` in when the property has none or carries the old placeholder. +- A directive is parsed and handed to the sketch from `loop()`, one at a time. One that arrives while the previous one still waits for `loop()` is dropped, as is one over 2047 bytes (`ALEX2ESP_MAX_DIRECTIVE`); both print an error. +- A directive that arrives twice is handled once: the Alex2MQTT broker currently delivers every message twice through a mirror. The `messageId`s of the last four directives are remembered; the repeat prints `repeated directive ... ignored`. +- The subscription delivers the directives of every endpoint of the account. Those for endpoints of another board are recognised by their topic and neither buffered nor parsed. +- Discovery is answered from `loop()`, not inside the MQTT callback. A discovery object over 2047 bytes is refused with an error; the other devices are still announced. +- `getState()` stays `INITIALIZED` until the first connect, and becomes `CONNECTED` when the broker has acknowledged both subscriptions (1.1.0: the first of them). A refused subscription prints an error. +- Serial output goes through one log with levels. `alexClient.setLogLevel()` takes `AlexaLogLevel::NONE`, `ERROR`, `INFO` (the default) or `DEBUG`; `-DALEX2ESP_LOG_MAX=<0..3>` in `build_flags` sets the highest level that is compiled in (default 2, `INFO`). This replaces the `Alex2ESP_DEBUG` define inside `AlexaUtils.cpp`. 1.1.0 printed the memory figures for every MQTT message and the whole directive, correlation token included, for every directive; both are gone. +- `getDevice()` before `begin()` prints an error: the device would have no root topic. A second `begin()` is ignored with an error. +- New: `Alex2ESP::setLogLevel()`, `Alex2ESP::setTimeSource()`, `AlexaDevice::hasEndpointId()`, `AlexaLog`, `AlexaSendResult`. +- Removed: the queues and buffers of `AlexaUtils` (`enqueue`, `dequeue`, `dequeueVals`, `enqueueReceive`, `dequeueReceive`, `isQueueEmpty`, `isQueueFull`, `isReceiveQueueEmpty`, `isReceiveQueueFull`, `receivePayload`, `nextMessageId`) and its `log`/`logln`, which printed nothing unless the library was edited; `AlexaUtils::printMemoryInfo()` stays. `MAX_STATUS_REPORT_SIZE` (the limit is `ALEX2ESP_MAX_MESSAGE`). The library no longer includes `ESP8266HTTPClient`. + +Memory: `examples/basicLight.cpp` for a D1 mini takes 34,116 bytes of static RAM (1.1.0: 52,768) and 336,757 bytes of flash (1.1.0: 350,885), as PlatformIO reports them (espressif8266 4.2.1, Arduino core 3.1.2). The static RAM was the five 2 KB queue slots, three more 2 KB buffers and the two HTTP clients. SNTP and the time stamp are 1.8 KB of the flash figure. + +Tests: `pio test -e native` in the repository runs 32 host tests of the receive and publish logic (reassembly of fragments, the two size limits, repeated directives, topics, time stamps). No board is needed. + +Boards that run 1.1.0 are not affected: the backend keeps the token topic and the HTTP routes. + ### 1.1.0 Behaviour changes: @@ -253,7 +283,7 @@ Packaging: `library.json` and `library.properties` restored with the Forgejo URL --- ## Contributing -Feel free to submit pull requests or issues for feature requests and bug fixes. +Feel free to submit pull requests or issues for feature requests and bug fixes. `pio test -e native` in the repository root runs the host tests (no board needed); they have to pass. --- diff --git a/src/Alex2ESP.cpp b/src/Alex2ESP.cpp index 219d950..fb164cb 100644 --- a/src/Alex2ESP.cpp +++ b/src/Alex2ESP.cpp @@ -1,39 +1,77 @@ /* * @title Alex2ESP Library - * @version 1.1.0 * @author David * @license MIT * @contributors chaos511 * - * @description The Alex2ESP library is a companion to the Alex2MQTT Alexa Skill, - * providing seamless integration between ESP-based devices and the Alex2MQTT server. - * This library connects to alex2mqtt.stormysdream.club, where the skill is hosted, - * allowing your devices to communicate effortlessly with the Alexa Voice Service - * using MQTT as the backbone. + * @description The bridge: MQTT session, discovery, directive receive and dispatch, publishing, time. + * The topics are listed in Alex2ESP.h. */ #include "Alex2ESP.h" +#include + +namespace +{ + const char MQTT_SERVER[] = "alex2mqtt.stormysdream.club"; + const uint16_t MQTT_PORT = 1883; + const char NTP_SERVER_1[] = "pool.ntp.org"; + const char NTP_SERVER_2[] = "time.nist.gov"; + + const uint8_t SUBSCRIPTION_REFUSED = 0x80; // Return code of a SUBACK for a subscription the broker denies + + // Room the MQTT client needs on top of topic and payload for the packet it builds (headers, its own object) + const size_t PACKET_OVERHEAD = 64; + + size_t largestFreeBlock() + { +#ifdef ESP8266 + return ESP.getMaxFreeBlockSize(); +#else + return ESP.getMaxAllocHeap(); // ESP32: untested +#endif + } +} Alex2ESP::Alex2ESP() - : rootTopic(), mqttUsername(nullptr), mqttPassword(nullptr), lastReconnectTime(0), reconnectAttempt(0), _state(Alex2ESPState::UNINITIALIZED), _disconnectReason(AsyncMqttClientDisconnectReason::TCP_DISCONNECTED) {} + : rootTopic(), + state(Alex2ESPState::UNINITIALIZED), + disconnectReason(AsyncMqttClientDisconnectReason::TCP_DISCONNECTED), + useSntp(true), + beginTime(0), + lastReconnectTime(0), + discoverSubscription(0), + directiveSubscription(0), + subscriptionsPending(0), + discoveryRequested(false), + discoveryActive(false), + discoveryDeferred(false), + discoveryNext(0), + discoveryAnnounced(0), + discoveryStarted(0), + discoveryLastAttempt(0), + directive(ALEX2ESP_MAX_DIRECTIVE) {} void Alex2ESP::begin(const char *username, const char *password, const char *rootTopic) { + if (state != Alex2ESPState::UNINITIALIZED) + { + ALEX2ESP_LOGE("begin() called again, ignored"); + return; + } - _state = Alex2ESPState::INITIALIZED; - _disconnectReason = AsyncMqttClientDisconnectReason::TCP_DISCONNECTED; this->rootTopic = rootTopic; - this->mqttUsername = username; - this->mqttPassword = password; - - discoverTopic = (String(this->rootTopic) + String("/discover")); - discoverTopicSend = (String(this->rootTopic) + String("/discover_r")); - - TopicESP = (String(this->rootTopic) + String("/+/alexaDirective_e")); + discoverTopic = this->rootTopic + "/discover"; + discoverTopicSend = this->rootTopic + "/discover_r"; + directiveFilter = this->rootTopic + "/+/alexaDirective"; // Log the root topic (never the credentials) - AlexaUtils::log("Setting root topic: "); - AlexaUtils::logln(rootTopic); - AlexaUtils::printMemoryInfo(); + ALEX2ESP_LOGI("root topic %s", this->rootTopic.c_str()); + + // Reports carry the time of their samples, so the board needs the time of day: UTC from SNTP + if (useSntp) + { + configTime(0, 0, NTP_SERVER_1, NTP_SERVER_2); + } // Set up MQTT client callbacks mqttClient.onConnect([this](bool sessionPresent) @@ -46,396 +84,427 @@ void Alex2ESP::begin(const char *username, const char *password, const char *roo { this->onMessage(topic, payload, properties, len, index, total); }); // Configure MQTT client - mqttClient.setServer(mqttServer, mqttPort); + mqttClient.setServer(MQTT_SERVER, MQTT_PORT); mqttClient.setCredentials(username, password); - // TODO: Throw error if dns fails or has no internet access? - _state = Alex2ESPState::CONNECTING; - mqttClient.connect(); + // loop() connects: see connectWhenClockIsSet() + beginTime = millis(); + state = Alex2ESPState::INITIALIZED; } -void Alex2ESP::onMqttConnect(bool sessionPresent) +void Alex2ESP::setLogLevel(AlexaLogLevel level) { - _state = Alex2ESPState::SUBSCRIBING; - - AlexaUtils::log("Connected to MQTT server, Now subscribing to: "); - AlexaUtils::log(discoverTopic.c_str()); - AlexaUtils::log(" "); - AlexaUtils::logln(TopicESP.c_str()); - - mqttClient.subscribe(discoverTopic.c_str(), 1); - mqttClient.subscribe(TopicESP.c_str(), 1); + AlexaLog::setLevel(level); } -void Alex2ESP::onSubscribe(uint16_t packetId, uint8_t qos) +void Alex2ESP::setTimeSource(bool useSntp) { - _state = Alex2ESPState::CONNECTED; - - AlexaUtils::log("Subscribed with packetId: "); - AlexaUtils::log(packetId); - AlexaUtils::log(" and QoS: "); - AlexaUtils::logln(qos); + this->useSntp = useSntp; } -// Internal: Handle disconnection -// TODO: auto reconnect? -void Alex2ESP::onMqttDisconnect(AsyncMqttClientDisconnectReason reason) +Alex2ESPState Alex2ESP::getState() const { - _state = Alex2ESPState::DISCONNECTED; - _disconnectReason = reason; + return state; } -void Alex2ESP::onMessage(char *topic, char *payload, AsyncMqttClientMessageProperties properties, size_t length, size_t index, size_t total) +AsyncMqttClientDisconnectReason Alex2ESP::getDisconnectReason() const { - AlexaUtils::logln("OnMessage Start"); - AlexaUtils::printMemoryInfo(); - - clearToSend = false; - - if (strcmp(topic, discoverTopic.c_str()) == 0) - { - // The payload is unused: answer once per message even when TCP split it, starting over if an answer is still going out - if (index == 0) - { - discoveryNext = 0; - discoveryPending = false; - discoveryStarted = millis(); - publishDiscovery(); - } - } - else if (index != 0 || total != length) - { - // AsyncMqttClient hands a large publish over in fragments; a directive id never is, so a fragment can only be part of one - AlexaUtils::logln("Ignoring fragmented directive id"); - } - else - { - // The payload is not NUL-terminated and belongs to the MQTT client: copy exactly `length` bytes out of it - String uuid; - uuid.concat(payload, length); - - for (auto &device : devices) - { - String directiveTopic = String(rootTopic) + "/" + device.getEndpointId() + "/alexaDirective_e"; - - if (strcmp(topic, directiveTopic.c_str()) == 0) - { - AlexaUtils::logln(topic); - AlexaUtils::logln(uuid); - if (!AlexaUtils::enqueueReceive(uuid.c_str())) - { - Serial.println("[Alex2ESP] receive queue full, directive dropped"); - } - } - } - } - AlexaUtils::logln("OnMessage End"); - AlexaUtils::printMemoryInfo(); - clearToSend = true; + return disconnectReason; } -// Answer a Discover: one discovery object per device, published straight to /discover_r over MQTT (async, no -// HTTP round trip, no queue slot). The backend accepts one endpoint object per message and collects everything that -// arrives within 1 s for Alexa's answer (5 s for its proactive push). publish() copies the object into the client's -// out-queue, where it stays on the heap until the broker has acknowledged it, and returns 0 once the free heap drops -// under 4 KB (or the client is not connected) - so rather than skip that device we stop, remember where we got to and -// let loop() carry on once the queue has drained, inside the backend's 5 s window. -void Alex2ESP::publishDiscovery() -{ - while (discoveryNext < devices.size()) - { - const AlexaDevice &device = devices[discoveryNext]; - AlexaUtils::logln("Getting device json"); - JsonDocument jsonData = device.getDeviceJSON(); - String jsonString; - serializeJson(jsonData, jsonString); - jsonData.clear(); - if (mqttClient.publish(discoverTopicSend.c_str(), 0, false, jsonString.c_str(), jsonString.length()) == 0) - { - if (!discoveryPending) - { - Serial.print("[Alex2ESP] discovery publish deferred at "); - Serial.println(device.getEndpointId()); - } - discoveryPending = true; - return; - } - discoveryNext++; - } - discoveryNext = 0; - discoveryPending = false; -} - -// Finish a discovery answer that publishDiscovery() had to cut short, or give it up once the backend has stopped listening -void Alex2ESP::finishDiscovery() -{ - if (!discoveryPending) - { - return; - } - if (millis() - discoveryStarted > DISCOVERY_WINDOW_MS) - { - Serial.printf("[Alex2ESP] discovery gave up: %u device(s) not announced\n", (unsigned)(devices.size() - discoveryNext)); - discoveryNext = 0; - discoveryPending = false; - } - else if (mqttClient.connected()) - { - publishDiscovery(); - } -} - -// boolean Alex2ESP::messagePackExtractor(char *topic, char *payload, size_t length, size_t index, size_t total) -// { -// payload[length] = '\0'; // Truncate the payload - -// uint8_t messageId = payload[0]; // Message ID (byte 0) -// uint8_t fragmentId = payload[1]; // Fragment ID (byte 1) -// uint8_t totalFragments = payload[2]; // Total Fragments (byte 2) - -// if (totalFragments > MAX_FRAGMENT_COUNT) -// { -// return false; //Message too long, skip. -// } - -// if (fragmentId == 0) -// { -// clearCache(); -// cache.totalParts = totalFragments; -// cache.messageId = String(messageId); -// } - -// if (cache.messageId != String(messageId)) -// { -// clearCache(); -// cache.messageId = String(messageId); -// cache.totalParts = totalFragments; -// } - -// if (fragmentId < cache.totalParts) -// { -// strncpy(cache.message + cache.fragmentOffset, payload + 3, length - 3); // Copy the fragment data -// cache.fragmentOffset=cache.fragmentOffset+(length-3); // -// cache.receivedParts++; -// Serial.printf("Received fragment %d of %d, length %d receivedParts: %d \n", fragmentId + 1, totalFragments, length, cache.receivedParts); - -// if (isMessageComplete()) -// { -// AlexaUtils::logln("Full message received."); -// AlexaUtils::logln(cache.message); - -// DeserializationError error = deserializeJson(inputDoc, cache.message); -// if (error) -// { -// AlexaUtils::log("Failed to parse message: "); -// AlexaUtils::logln(String(error.f_str())); -// clearCache(); -// return false; -// } -// else -// { -// AlexaUtils::logln("JSON message parsed successfully!"); -// clearCache(true); -// return true; -// } -// } -// } -// return false; -// } - AlexaDevice *Alex2ESP::getDevice(const String &name, const String &endpointId) { // Check if a device with the given endpointId already exists - for (auto &device : devices) + AlexaDevice *existing = findDevice(endpointId.c_str(), endpointId.length()); + if (existing != nullptr) { - if (device.getEndpointId() == endpointId) - { - return &device; // Return the existing device - } + return existing; + } + + if (state == Alex2ESPState::UNINITIALIZED) + { + ALEX2ESP_LOGE("device %s created before begin(): it has no root topic, its reports will not arrive", endpointId.c_str()); } // If the device doesn't exist, create a new one - devices.emplace_back(name, rootTopic, endpointId); - - Serial.print("Created new device: "); - Serial.print(name); - Serial.print(" with endpointId: "); - Serial.println(endpointId); + devices.emplace_back(name, rootTopic, endpointId, this); + ALEX2ESP_LOGI("device %s (%s) created", endpointId.c_str(), name.c_str()); // Return a pointer to the newly created device return &devices.back(); } -// void Alex2ESP::clearCache(boolean skipDocClear) -// { -// if (!skipDocClear) -// { -// inputDoc.clear(); -// } -// cache.receivedParts = 0; // Reset the count of received parts -// cache.messageId = ""; // Clear the message ID -// cache.totalParts = 0; // Reset the total parts count -// cache.fragmentOffset = 0; -// memset(cache.message, 0, sizeof(cache.message)); // Clear the message buffer -// } - -// bool Alex2ESP::isMessageComplete() -// { -// return cache.receivedParts == cache.totalParts; // Check if we have received all parts -// } - -Alex2ESPState Alex2ESP::getState() const +AlexaDevice *Alex2ESP::findDevice(const char *endpointId, size_t length) { - return _state; + for (auto &device : devices) + { + if (device.hasEndpointId(endpointId, length)) + { + return &device; + } + } + return nullptr; } -AsyncMqttClientDisconnectReason Alex2ESP::getDisconnectReason() const +// The first connect waits until the clock is set, for CLOCK_WAIT_MS at most: a report sent before SNTP has answered +// would carry a time of sample in 1970. +void Alex2ESP::connectWhenClockIsSet() { - return _disconnectReason; + if (state != Alex2ESPState::INITIALIZED) + { + return; + } + if (!AlexaBridgeLogic::clockIsSet(time(nullptr))) + { + if (millis() - beginTime < CLOCK_WAIT_MS) + { + return; + } + ALEX2ESP_LOGE("clock not set after %lu ms: connecting, reports carry a wrong time until it is set", CLOCK_WAIT_MS); + } + + state = Alex2ESPState::CONNECTING; + lastReconnectTime = millis(); + mqttClient.connect(); } void Alex2ESP::handleMqttReconnection() { - if (!mqttClient.connected() && _disconnectReason == AsyncMqttClientDisconnectReason::TCP_DISCONNECTED) + if (state == Alex2ESPState::UNINITIALIZED || state == Alex2ESPState::INITIALIZED) { - if (millis() - lastReconnectTime > 5000) + return; + } + if (!mqttClient.connected() && disconnectReason == AsyncMqttClientDisconnectReason::TCP_DISCONNECTED) + { + if (millis() - lastReconnectTime > RECONNECT_INTERVAL_MS) { lastReconnectTime = millis(); - reconnectAttempt++; mqttClient.connect(); } } } -void Alex2ESP::processHttpPost() + +void Alex2ESP::onMqttConnect(bool sessionPresent) { - if (!AlexaUtils::isQueueEmpty()) + state = Alex2ESPState::SUBSCRIBING; + + subscriptionsPending = 2; + discoverSubscription = mqttClient.subscribe(discoverTopic.c_str(), 1); + directiveSubscription = mqttClient.subscribe(directiveFilter.c_str(), 1); + if (discoverSubscription == 0 || directiveSubscription == 0) { - AlexaUtils::logln("processHttpPost"); + ALEX2ESP_LOGE("the MQTT client refused to subscribe (%u bytes of heap free): no directives until it reconnects", (unsigned)ESP.getFreeHeap()); + return; + } + ALEX2ESP_LOGI("connected to %s, subscribing", MQTT_SERVER); +} - const char *payload = AlexaUtils::dequeueVals(false); +void Alex2ESP::onSubscribe(uint16_t packetId, uint8_t qos) +{ + if (packetId != discoverSubscription && packetId != directiveSubscription) + { + return; + } + const String &topic = (packetId == discoverSubscription) ? discoverTopic : directiveFilter; + if (qos == SUBSCRIPTION_REFUSED) + { + ALEX2ESP_LOGE("the broker refused the subscription to %s", topic.c_str()); + return; + } - httpPOST.begin(wifiPOST, "http://alex2mqtt.stormysdream.club/Alex2ESP"); - httpPOST.setAuthorization(mqttUsername, mqttPassword); - httpPOST.addHeader("Content-Type", "text/plain"); - - AlexaUtils::log("Sending Data: "); - AlexaUtils::logln(payload); - - int httpResponseCode = httpPOST.POST(payload); - - if (httpResponseCode != 200) - { - - if (retryCountPOST >= MAX_RETRY_COUNT) - { - AlexaUtils::dequeueVals(true); - } - - AlexaUtils::log("POST Error code: "); - AlexaUtils::logln(httpResponseCode); - retryCountPOST++; - if (httpResponseCode > 0) - { - String response = httpPOST.getString(); - Serial.println("Response:"); - Serial.println(response); - } - } - else - { - retryCountPOST = 0; - AlexaUtils::dequeueVals(true); - } - httpPOST.end(); + ALEX2ESP_LOGD("subscribed to %s", topic.c_str()); + if (subscriptionsPending > 0 && --subscriptionsPending == 0) + { + state = Alex2ESPState::CONNECTED; + ALEX2ESP_LOGI("ready: %u device(s)", (unsigned)devices.size()); } } -void Alex2ESP::processHttpGet() + +void Alex2ESP::onMqttDisconnect(AsyncMqttClientDisconnectReason reason) { - if (!AlexaUtils::isReceiveQueueEmpty()) + state = Alex2ESPState::DISCONNECTED; + disconnectReason = reason; + directive.cancelArrival(); + ALEX2ESP_LOGI("disconnected (reason %u)", (unsigned)reason); +} + +// Runs in the network context, once per fragment of a message. It only takes notes: loop() answers a Discover and +// parses and dispatches a directive, on the sketch's own stack. +void Alex2ESP::onMessage(char *topic, char *payload, AsyncMqttClientMessageProperties properties, size_t length, size_t index, size_t total) +{ + if (strcmp(topic, discoverTopic.c_str()) == 0) { - AlexaUtils::logln("processHttpGet"); - - String uuid; - AlexaUtils::dequeueReceive(uuid, false); - AlexaUtils::log("Sending get request for uuid: "); - // AlexaUtils::logln(uuid); - - httpGET.begin(wifiGET, "http://alex2mqtt.stormysdream.club/Alex2ESP/" + uuid); - httpGET.setAuthorization(mqttUsername, mqttPassword); - int httpResponseCode = httpGET.GET(); - - Serial.print("HTTP Response code: "); - Serial.println(httpResponseCode); - if (httpResponseCode != 200) + // The payload is unused: answer once per message even when TCP split it + if (index == 0) { - Serial.print("Error code: "); - Serial.println(httpResponseCode); - if (retryCountGET >= MAX_RETRY_COUNT) - { - AlexaUtils::dequeueReceive(uuid, true); - } + discoveryRequested = true; + discoveryStarted = millis(); + } + return; + } - retryCountGET++; + size_t endpointLength = 0; + const char *endpointId = AlexaBridgeLogic::directiveEndpoint(topic, rootTopic.c_str(), &endpointLength); + if (endpointId == nullptr) + { + return; + } + if (findDevice(endpointId, endpointLength) == nullptr) + { + // /+/alexaDirective delivers the directives of every endpoint of the account, also those of other + // boards: not an error, and not worth a buffer + if (index == 0) + { + ALEX2ESP_LOGD("directive on %s is for another board", topic); + } + return; + } + + switch (directive.append(payload, length, index, total)) + { + case AlexaDirectiveBuffer::Result::TOO_LARGE: + ALEX2ESP_LOGE("directive of %u bytes on %s dropped: the limit is %u (ALEX2ESP_MAX_DIRECTIVE)", (unsigned)total, topic, (unsigned)ALEX2ESP_MAX_DIRECTIVE); + break; + case AlexaDirectiveBuffer::Result::BUSY: + ALEX2ESP_LOGE("directive on %s dropped: the one before it still waits for loop()", topic); + break; + case AlexaDirectiveBuffer::Result::NO_MEMORY: + ALEX2ESP_LOGE("directive of %u bytes on %s dropped: no memory (%u bytes of heap free)", (unsigned)total, topic, (unsigned)ESP.getFreeHeap()); + break; + case AlexaDirectiveBuffer::Result::OUT_OF_ORDER: + ALEX2ESP_LOGE("directive on %s dropped: the fragment at byte %u of %u does not continue it", topic, (unsigned)index, (unsigned)total); + break; + case AlexaDirectiveBuffer::Result::EMPTY: + ALEX2ESP_LOGE("directive on %s dropped: it is empty", topic); + break; + case AlexaDirectiveBuffer::Result::DUPLICATE: + ALEX2ESP_LOGI("repeated directive on %s ignored", topic); + break; + case AlexaDirectiveBuffer::Result::INCOMPLETE: + case AlexaDirectiveBuffer::Result::COMPLETE: + case AlexaDirectiveBuffer::Result::IGNORED: + break; + } +} + +void Alex2ESP::processDirective() +{ + if (!directive.ready()) + { + return; + } + + JsonDocument message; + size_t length = directive.length(); + DeserializationError error = deserializeJson(message, directive.data(), length); + // The text is not needed any more. Releasing it before the handler runs frees the heap for the report and + // lets the next directive arrive while the sketch deals with this one. + directive.release(); + + if (error) + { + char reason[24]; + strncpy_P(reason, reinterpret_cast(error.f_str()), sizeof(reason) - 1); + reason[sizeof(reason) - 1] = '\0'; + ALEX2ESP_LOGE("directive of %u bytes dropped: not JSON (%s)", (unsigned)length, reason); + return; + } + + // Alex2MQTT publishes the directive itself: {"header": ..., "endpoint": ..., "payload": ...} + const JsonDocument &received = message; + if (!received["header"].is()) + { + ALEX2ESP_LOGE("directive of %u bytes dropped: it has no header", (unsigned)length); + return; + } + const char *endpointId = received["endpoint"]["endpointId"] | ""; + const char *interfaceName = received["header"]["namespace"] | ""; + const char *name = received["header"]["name"] | ""; + const char *messageId = received["header"]["messageId"] | ""; + + AlexaDevice *device = findDevice(endpointId, strlen(endpointId)); + if (device == nullptr) + { + ALEX2ESP_LOGE("directive %s.%s dropped: it names the endpoint '%s', which this board does not have", interfaceName, name, endpointId); + return; + } + if (recentIds.seenBefore(messageId)) + { + ALEX2ESP_LOGI("%s: repeated directive %s.%s ignored", endpointId, interfaceName, name); + return; + } + + ALEX2ESP_LOGI("%s <- %s.%s", endpointId, interfaceName, name); + ALEX2ESP_LOGD("%u bytes, %u bytes of heap free", (unsigned)length, (unsigned)ESP.getFreeHeap()); + + device->triggerEvent("DirectiveReceived", message, AlexaInterfaceType::UNKNOWN, false); + if (strcmp(interfaceName, "Alexa") == 0 && strcmp(name, "ReportState") == 0) + { + device->triggerEvent("ReportState", message, AlexaInterfaceType::UNKNOWN); + } + else + { + device->triggerEvent("Event", message, AlexaInterfaceUtils::fromString(interfaceName)); + } +} + +// Answer a Discover: one discovery object per device, published straight to /discover_r. The backend accepts +// one endpoint object per message and collects everything that arrives within 1 s for Alexa's answer (5 s for its +// proactive push), and it takes the answers to one request as the complete list of the board. So every Discover +// starts the list from the top, also when the answer to the one before is still going out. +void Alex2ESP::continueDiscovery() +{ + if (discoveryRequested) + { + discoveryRequested = false; + discoveryActive = true; + discoveryDeferred = false; + discoveryNext = 0; + discoveryAnnounced = 0; + } + if (!discoveryActive) + { + return; + } + + if (millis() - discoveryStarted > DISCOVERY_WINDOW_MS) + { + // The backend has stopped listening + ALEX2ESP_LOGE("discovery gave up: %u device(s) not announced", (unsigned)(devices.size() - discoveryAnnounced)); + discoveryActive = false; + return; + } + if (discoveryDeferred && millis() - discoveryLastAttempt < DISCOVERY_RETRY_MS) + { + return; + } + publishDiscovery(); +} + +// publish() copies the object into the client's out-queue, where it stays on the heap until it has been written to +// the socket, and the client refuses another one once the free heap drops under 4 KB - so rather than skip that +// device we stop, remember where we got to and let loop() carry on once the queue has drained. +void Alex2ESP::publishDiscovery() +{ + while (discoveryNext < devices.size()) + { + const AlexaDevice &device = devices[discoveryNext]; + JsonDocument json = device.getDeviceJSON(); + size_t length = 0; + AlexaSendResult result = trySend(discoverTopicSend.c_str(), json, &length); + + if (result == AlexaSendResult::REFUSED || result == AlexaSendResult::NOT_CONNECTED) + { + if (!discoveryDeferred) + { + ALEX2ESP_LOGI("discovery deferred at %s (%u bytes of heap free)", device.getEndpointId().c_str(), (unsigned)ESP.getFreeHeap()); + } + discoveryDeferred = true; + discoveryLastAttempt = millis(); + return; + } + + if (result == AlexaSendResult::OK) + { + discoveryAnnounced++; } else { - retryCountGET = 0; - AlexaUtils::dequeueReceive(uuid, true); - String payloadStr = httpGET.getString(); - size_t length = payloadStr.length(); - - if (length < (size_t)(AlexaUtils::MAX_PAYLOAD_LENGTH - 1)) - { - payloadStr.toCharArray(AlexaUtils::receivePayload, length + 1); - payloadStr = ""; - AlexaUtils::receivePayload[length + 1] = '\0'; - // AlexaUtils::logln(AlexaUtils::receivePayload); - - DeserializationError error = deserializeJson(inputDoc, AlexaUtils::receivePayload); - if (error) - { - AlexaUtils::log("Failed to parse message: "); - AlexaUtils::logln(String(error.f_str())); - inputDoc.clear(); - } - else - { - AlexaUtils::logln("JSON message parsed successfully!"); - - for (auto &device : devices) - { - if (strcmp(device.getEndpointId().c_str(), inputDoc["directive"]["endpoint"]["endpointId"] | "") == 0) - { - serializeJson(inputDoc, Serial); - Serial.println(); - device.triggerEvent("DirectiveReceived", inputDoc["directive"], AlexaInterfaceType::UNKNOWN, false); - - if (inputDoc["directive"]["header"]["name"] == "ReportState") - { - device.triggerEvent("ReportState", inputDoc["directive"], AlexaInterfaceType::UNKNOWN); - } - else - { - device.triggerEvent("Event", inputDoc["directive"], AlexaInterfaceUtils::fromString(inputDoc["directive"]["header"]["namespace"])); - } - } - } - } - } + // This object can never be sent: say so and announce the others + logRefusal(discoverTopicSend.c_str(), result, length); + ALEX2ESP_LOGE("device %s is not announced", device.getEndpointId().c_str()); } - httpGET.end(); + discoveryNext++; + } + + discoveryActive = false; + ALEX2ESP_LOGI("discovery: %u of %u device(s) announced", (unsigned)discoveryAnnounced, (unsigned)devices.size()); +} + +AlexaSendResult Alex2ESP::publish(const char *topic, JsonDocument &doc) +{ + size_t length = 0; + AlexaSendResult result = trySend(topic, doc, &length); + if (result == AlexaSendResult::OK) + { + ALEX2ESP_LOGD("%u bytes -> %s, %u bytes of heap free", (unsigned)length, topic, (unsigned)ESP.getFreeHeap()); + } + else + { + logRefusal(topic, result, length); + } + return result; +} + +void Alex2ESP::timestamp(char *buffer, size_t size) +{ + AlexaBridgeLogic::formatTimestamp(time(nullptr), buffer, size); +} + +// Sends the document whole or not at all, and says which. The caller reports a refusal. +AlexaSendResult Alex2ESP::trySend(const char *topic, JsonDocument &doc, size_t *length) +{ + AlexaSendResult result = AlexaBridgeLogic::checkMessage(doc, ALEX2ESP_MAX_MESSAGE, length); + if (result != AlexaSendResult::OK) + { + return result; + } + if (!mqttClient.connected()) + { + return AlexaSendResult::NOT_CONNECTED; + } + + char *text = static_cast(malloc(*length + 1)); + if (text == nullptr) + { + return AlexaSendResult::REFUSED; + } + serializeJson(doc, text, *length + 1); + + // The client copies topic and text into a packet of its own and returns 0 when it is not connected or the free + // heap is under 4 KB. It does not check that the packet fits the largest free block, and it cannot report an + // allocation that fails (no exceptions on this target: the board resets). So that is checked here. + result = AlexaSendResult::REFUSED; + if (largestFreeBlock() >= *length + strlen(topic) + PACKET_OVERHEAD && mqttClient.publish(topic, 0, false, text, *length) != 0) + { + result = AlexaSendResult::OK; + } + free(text); + return result; +} + +void Alex2ESP::logRefusal(const char *topic, AlexaSendResult result, size_t length) +{ + switch (result) + { + case AlexaSendResult::NOT_CONNECTED: + ALEX2ESP_LOGE("%s: %u bytes not sent, no session with the broker", topic, (unsigned)length); + break; + case AlexaSendResult::REFUSED: + ALEX2ESP_LOGE("%s: %u bytes not sent, the MQTT client or the heap is full (%u bytes free)", topic, (unsigned)length, (unsigned)ESP.getFreeHeap()); + break; + case AlexaSendResult::TOO_LARGE: + if (length == 0) + { + ALEX2ESP_LOGE("%s: not sent, the message ran out of memory while it was built", topic); + } + else + { + ALEX2ESP_LOGE("%s: %u bytes not sent, the limit is %u (ALEX2ESP_MAX_MESSAGE)", topic, (unsigned)length, (unsigned)ALEX2ESP_MAX_MESSAGE); + } + break; + case AlexaSendResult::EMPTY: + ALEX2ESP_LOGE("%s: not sent, the message is empty", topic); + break; + case AlexaSendResult::OK: + break; } } void Alex2ESP::loop() { - + connectWhenClockIsSet(); handleMqttReconnection(); - finishDiscovery(); - if(clearToSend){ - processHttpPost(); - processHttpGet(); - } - - return; + continueDiscovery(); + processDirective(); } diff --git a/src/Alex2ESP.h b/src/Alex2ESP.h index 6681eaa..ca44e50 100644 --- a/src/Alex2ESP.h +++ b/src/Alex2ESP.h @@ -1,15 +1,15 @@ /* * @title Alex2ESP Library - * @version 1.1.0 * @author David * @license MIT * @contributors chaos511 * - * @description The Alex2ESP library is a companion to the Alex2MQTT Alexa Skill, - * providing seamless integration between ESP-based devices and the Alex2MQTT server. - * This library connects to alex2mqtt.stormysdream.club, where the skill is hosted, - * allowing your devices to communicate effortlessly with the Alexa Voice Service - * using MQTT as the backbone. + * @description Companion library of the Alex2MQTT Alexa skill: the devices a sketch declares become Alexa + * endpoints through the MQTT broker at alex2mqtt.stormysdream.club. + * + * /discover in answered with one discovery object per device on /discover_r + * //alexaDirective in the directive, handed to the device's ReportState / Event handler + * //alexaResponce out the report the handler built with buildStatusMessage() */ #ifndef ALEX2ESP_H @@ -18,98 +18,113 @@ #include #include #include +#include "AlexaBridgeLogic.h" #include "AlexaDevice.h" #include "AlexaInterface.h" +#include "AlexaLog.h" +#include "AlexaTransport.h" #include "AlexaUtils.h" #include -#ifdef ESP32 -#include // ESP32: untested -#else -#include +// Largest directive the bridge accepts, in bytes. Override with -DALEX2ESP_MAX_DIRECTIVE= in build_flags. +#ifndef ALEX2ESP_MAX_DIRECTIVE +#define ALEX2ESP_MAX_DIRECTIVE 2047 #endif - enum class Alex2ESPState { - UNINITIALIZED, - INITIALIZED, + UNINITIALIZED, // begin() has not been called + INITIALIZED, // begin() has been called; the first connect waits for the clock CONNECTING, - SUBSCRIBING, - CONNECTED, + SUBSCRIBING, // session open, the subscriptions are not acknowledged yet + CONNECTED, // subscribed: discovery requests and directives arrive DISCONNECTED }; -class Alex2ESP +class Alex2ESP : public AlexaTransport { public: // Constructor Alex2ESP(); - // Begin function for initialization: MQTT username, MQTT password, root topic (the same order as alex2node) + // Begin function for initialization: MQTT username, MQTT password, root topic (the same order as alex2node). + // The username and the password are not copied: they have to stay valid for as long as the client is used. + // Starts SNTP; loop() opens the MQTT session once the clock is set, or after 5 s without an answer. void begin(const char *username, const char *password, const char *rootTopic); Alex2ESPState getState() const; + + // Call from the sketch's loop(): connects, answers discovery requests and hands directives to the devices. + // It does not block; the handlers of the sketch run inside it. void loop(); AsyncMqttClientDisconnectReason getDisconnectReason() const; // Returns the device with this endpointId, creating it on first use. The pointer stays valid for the lifetime // of the client: devices live in a std::deque, which never relocates its elements when another one is added. + // Call it after begin(): a device takes the root topic of its reports when it is created. AlexaDevice *getDevice(const String &name, const String &endpointId); + // What the library prints on Serial: AlexaLogLevel::NONE, ERROR, INFO (the default) or DEBUG + void setLogLevel(AlexaLogLevel level); + + // false before begin(): the sketch sets the clock itself (its own configTime() with a time zone, an RTC). + // begin() then leaves SNTP alone; the library only reads time(). + void setTimeSource(bool useSntp); + + // AlexaTransport: what the devices' reports are sent and stamped with + AlexaSendResult publish(const char *topic, JsonDocument &doc) override; + void timestamp(char *buffer, size_t size) override; + private: - static const int MAX_RETRY_COUNT = 2; // Define maximum retry count + static const unsigned long CLOCK_WAIT_MS = 5000; // How long the first connect waits for SNTP + static const unsigned long RECONNECT_INTERVAL_MS = 5000; static const unsigned long DISCOVERY_WINDOW_MS = 5000; // How long the backend keeps collecting a discovery answer - int retryCountPOST=0; - int retryCountGET=0; + static const unsigned long DISCOVERY_RETRY_MS = 20; // Pause before a refused discovery publish is tried again - size_t discoveryNext = 0; // Next device to announce while a discovery answer is still going out - bool discoveryPending = false; // publishDiscovery() stopped early (client out-queue full); loop() finishes it - unsigned long discoveryStarted = 0; // millis() when the Discover arrived - - boolean clearToSend=false; AsyncMqttClient mqttClient; // MQTT client instance - String rootTopic; // Root topic for communication + String rootTopic; // Root topic for communication String discoverTopic; // The topic we listen on for discovery messages String discoverTopicSend; // The topic we send discovery messages - - String TopicESP; // The topic we listen on for esp messages - - const char *mqttUsername; // Username for authentication - const char *mqttPassword; // Password for authentication - - unsigned long lastReconnectTime; - int reconnectAttempt; - - HTTPClient httpGET; - HTTPClient httpPOST; - - WiFiClient wifiGET; - WiFiClient wifiPOST; - JsonDocument inputDoc; + String directiveFilter; // The subscription that delivers the directives of every endpoint std::deque devices; // Collection of devices (deque: pointers handed out by getDevice stay valid) - Alex2ESPState _state; - AsyncMqttClientDisconnectReason _disconnectReason; + Alex2ESPState state; + AsyncMqttClientDisconnectReason disconnectReason; + bool useSntp; + unsigned long beginTime; // millis() when begin() ran + unsigned long lastReconnectTime; + uint16_t discoverSubscription; // Packet ids of the two SUBSCRIBEs, to match their acknowledgements + uint16_t directiveSubscription; + uint8_t subscriptionsPending; - // MQTT connection details - const char *mqttServer = "alex2mqtt.stormysdream.club"; - uint16_t mqttPort = 1883; + bool discoveryRequested; // Set by the MQTT callback, taken by loop() + bool discoveryActive; // An answer is going out + bool discoveryDeferred; // The answer had to pause: the MQTT client or the heap was full + size_t discoveryNext; // Next device to announce + size_t discoveryAnnounced; // Devices announced in this answer + unsigned long discoveryStarted; // millis() when the Discover arrived + unsigned long discoveryLastAttempt; // millis() of the publish that was refused - // Internal event handlers + AlexaDirectiveBuffer directive; // The directive that waits for loop(), or is still arriving + AlexaRecentIds recentIds; // messageIds of the last directives handled + + // Internal event handlers (called by the MQTT client from the network context: they only take notes) void onMqttConnect(bool sessionPresent); void onMqttDisconnect(AsyncMqttClientDisconnectReason reason); void onSubscribe(uint16_t packetId, uint8_t qos); void onMessage(char *topic, char *payload, AsyncMqttClientMessageProperties properties, size_t length, size_t index, size_t total); //loop processing function + void connectWhenClockIsSet(); void handleMqttReconnection(); + void continueDiscovery(); void publishDiscovery(); - void finishDiscovery(); - void processHttpPost(); - void processHttpGet(); + void processDirective(); + AlexaDevice *findDevice(const char *endpointId, size_t length); + AlexaSendResult trySend(const char *topic, JsonDocument &doc, size_t *length); + void logRefusal(const char *topic, AlexaSendResult result, size_t length); }; #endif // ALEX2ESP_H diff --git a/src/AlexaBridgeLogic.cpp b/src/AlexaBridgeLogic.cpp new file mode 100644 index 0000000..5aebf1f --- /dev/null +++ b/src/AlexaBridgeLogic.cpp @@ -0,0 +1,300 @@ +#include "AlexaBridgeLogic.h" +#include "AlexaCompat.h" +#include +#include +#include + +AlexaDirectiveBuffer::AlexaDirectiveBuffer(size_t capacityBytes) + : capacity(capacityBytes) {} + +AlexaDirectiveBuffer::~AlexaDirectiveBuffer() +{ + discard(); +} + +AlexaDirectiveBuffer::Result AlexaDirectiveBuffer::append(const char *data, size_t length, size_t index, size_t total) +{ + if (index == 0) + { + return start(data, length, total); + } + return proceed(data, length, index, total); +} + +// First fragment: decides what happens to the whole message +AlexaDirectiveBuffer::Result AlexaDirectiveBuffer::start(const char *data, size_t length, size_t total) +{ + // A message that was still arriving ends here: the connection dropped in the middle of it + if (!waiting) + { + discard(); + } + arrival = Arrival::NONE; + + if (total == 0) + { + return Result::EMPTY; + } + if (total > capacity) + { + return refuse(Result::TOO_LARGE); + } + if (data == nullptr || length > total) + { + return refuse(Result::OUT_OF_ORDER); + } + + if (waiting) + { + // loop() has not read the previous directive yet. The repeat of it is dropped without a loss; anything + // else has no place to go. + if (total != size || memcmp(text, data, length) != 0) + { + return refuse(Result::BUSY); + } + if (length == total) + { + return Result::DUPLICATE; + } + offset = length; + arrival = Arrival::COMPARING; + return Result::INCOMPLETE; + } + + text = static_cast(malloc(total + 1)); + if (text == nullptr) + { + return refuse(Result::NO_MEMORY); + } + size = total; + offset = 0; + return store(data, length); +} + +// Any later fragment +AlexaDirectiveBuffer::Result AlexaDirectiveBuffer::proceed(const char *data, size_t length, size_t index, size_t total) +{ + if (arrival == Arrival::SKIPPING) + { + return Result::IGNORED; + } + if (arrival == Arrival::NONE) + { + return refuse(Result::OUT_OF_ORDER); + } + + // The fragment has to continue exactly where the last one ended and stay inside the message + if (data == nullptr || total != size || index != offset || length > size - offset) + { + if (!waiting) + { + discard(); + } + return refuse(Result::OUT_OF_ORDER); + } + + if (arrival == Arrival::COMPARING) + { + if (memcmp(text + offset, data, length) == 0) + { + offset += length; + if (offset < size) + { + return Result::INCOMPLETE; + } + arrival = Arrival::NONE; + if (!waiting) + { + discard(); + } + return Result::DUPLICATE; + } + if (waiting) + { + return refuse(Result::BUSY); + } + // Not a repeat after all, and the directive it was compared with has been read in the meantime: its block + // has the right size and already holds the bytes that matched, so it takes the rest of this message. + arrival = Arrival::STORING; + } + return store(data, length); +} + +AlexaDirectiveBuffer::Result AlexaDirectiveBuffer::store(const char *data, size_t length) +{ + memcpy(text + offset, data, length); + offset += length; + if (offset < size) + { + arrival = Arrival::STORING; + return Result::INCOMPLETE; + } + text[size] = '\0'; + waiting = true; + arrival = Arrival::NONE; + return Result::COMPLETE; +} + +AlexaDirectiveBuffer::Result AlexaDirectiveBuffer::refuse(Result reason) +{ + arrival = Arrival::SKIPPING; + return reason; +} + +void AlexaDirectiveBuffer::release() +{ + if (!waiting) + { + return; + } + waiting = false; + // A repeat that is still being compared needs the text until its last fragment + if (arrival != Arrival::COMPARING) + { + discard(); + } +} + +void AlexaDirectiveBuffer::cancelArrival() +{ + arrival = Arrival::NONE; + if (!waiting) + { + discard(); + } +} + +void AlexaDirectiveBuffer::discard() +{ + free(text); + text = nullptr; + size = 0; + offset = 0; +} + +bool AlexaRecentIds::seenBefore(const char *id) +{ + if (id == nullptr || id[0] == '\0') + { + return false; + } + + // FNV-1a, 64 bit + uint64_t hash = 0xcbf29ce484222325ULL; + for (const char *c = id; *c != '\0'; c++) + { + hash ^= static_cast(*c); + hash *= 0x100000001b3ULL; + } + + for (uint8_t i = 0; i < count; i++) + { + if (hashes[i] == hash) + { + return true; + } + } + hashes[next] = hash; + next = (next + 1) % CAPACITY; + if (count < CAPACITY) + { + count++; + } + return false; +} + +bool AlexaBridgeLogic::clockIsSet(time_t now) +{ + return now >= CLOCK_SET_AFTER; +} + +bool AlexaBridgeLogic::formatTimestamp(time_t instant, char *buffer, size_t size) +{ + if (buffer == nullptr || size == 0) + { + return false; + } + buffer[0] = '\0'; + if (size < ALEXA_TIMESTAMP_SIZE) + { + return false; + } + + // gmtime_r and snprintf instead of strftime, which links the time zone and locale tables: 5.5 KB of flash. + // snprintf and not snprintf_P: on the ESP8266 snprintf reads a PROGMEM format itself, and snprintf_P comes + // in one object with printf_P and sprintf_P, which link the FILE-based printf: 5.8 KB. + struct tm utc; + if (gmtime_r(&instant, &utc) == nullptr) + { + return false; + } + int written = snprintf(buffer, size, PSTR("%04d-%02d-%02dT%02d:%02d:%02dZ"), + utc.tm_year + 1900, utc.tm_mon + 1, utc.tm_mday, utc.tm_hour, utc.tm_min, utc.tm_sec); + if (written != ALEXA_TIMESTAMP_SIZE - 1) + { + buffer[0] = '\0'; + return false; + } + return true; +} + +const char *AlexaBridgeLogic::directiveEndpoint(const char *topic, const char *rootTopic, size_t *length) +{ + static const char suffix[] = "/alexaDirective"; + const size_t suffixLength = sizeof(suffix) - 1; + + if (topic == nullptr || rootTopic == nullptr || length == nullptr) + { + return nullptr; + } + size_t rootLength = strlen(rootTopic); + size_t topicLength = strlen(topic); + // '/' + if (topicLength < rootLength + 1 + 1 + suffixLength) + { + return nullptr; + } + if (memcmp(topic, rootTopic, rootLength) != 0 || topic[rootLength] != '/') + { + return nullptr; + } + if (memcmp(topic + topicLength - suffixLength, suffix, suffixLength) != 0) + { + return nullptr; + } + + const char *endpointId = topic + rootLength + 1; + size_t endpointLength = topicLength - suffixLength - rootLength - 1; + if (memchr(endpointId, '/', endpointLength) != nullptr) + { + return nullptr; + } + *length = endpointLength; + return endpointId; +} + +AlexaSendResult AlexaBridgeLogic::checkMessage(const JsonDocument &doc, size_t maxLength, size_t *length) +{ + if (length != nullptr) + { + *length = 0; + } + if (doc.overflowed()) + { + return AlexaSendResult::TOO_LARGE; + } + if (doc.isNull()) + { + return AlexaSendResult::EMPTY; + } + size_t needed = measureJson(doc); + if (length != nullptr) + { + *length = needed; + } + if (needed > maxLength) + { + return AlexaSendResult::TOO_LARGE; + } + return AlexaSendResult::OK; +} diff --git a/src/AlexaBridgeLogic.h b/src/AlexaBridgeLogic.h new file mode 100644 index 0000000..de310a0 --- /dev/null +++ b/src/AlexaBridgeLogic.h @@ -0,0 +1,123 @@ +// The parts of the bridge that are plain logic: reassembling a directive from the fragments the MQTT client hands +// over, telling a repeated directive from a new one, the limits on what is received and sent, topics and time +// stamps. Nothing here touches the MQTT client, Wi-Fi or Serial, so the same code runs in the host tests +// (test/test_bridge_logic, pio test -e native) with the byte sequences a broker would deliver. +#ifndef ALEXA_BRIDGE_LOGIC_H +#define ALEXA_BRIDGE_LOGIC_H + +#include +#include +#include +#include +#include "AlexaTransport.h" + +// One directive on its way from the MQTT callback to loop(). +// +// AsyncMqttClient delivers a message larger than one TCP segment in fragments: the callback runs once per fragment +// with the fragment's offset (index) and the length of the whole message (total). The fragments are collected into +// one heap block of total + 1 bytes, allocated when the first fragment arrives and freed by release(). A message is +// refused as a whole when it is larger than the capacity or when another directive is still waiting; nothing is +// ever written outside the block. +class AlexaDirectiveBuffer +{ +public: + enum class Result : uint8_t + { + INCOMPLETE, // fragment taken, the message continues + COMPLETE, // last fragment taken: the directive waits for loop() + DUPLICATE, // byte for byte the directive that waits or was just handled; the repeat is discarded + IGNORED, // fragment of a message that was refused at its first fragment + EMPTY, // message without a payload + TOO_LARGE, // total is over the capacity + BUSY, // a different directive is still waiting for loop() + NO_MEMORY, // no heap for total + 1 bytes + OUT_OF_ORDER // fragment that does not continue the message being received + }; + + // capacityBytes: the largest message that is accepted (ALEX2ESP_MAX_DIRECTIVE in the bridge) + explicit AlexaDirectiveBuffer(size_t capacityBytes); + ~AlexaDirectiveBuffer(); + AlexaDirectiveBuffer(const AlexaDirectiveBuffer &) = delete; + AlexaDirectiveBuffer &operator=(const AlexaDirectiveBuffer &) = delete; + + // One call per fragment, in the order of arrival. Only TOO_LARGE, BUSY, NO_MEMORY, OUT_OF_ORDER and EMPTY + // report a message that is lost; each is returned once per message, its remaining fragments are IGNORED. + Result append(const char *data, size_t length, size_t index, size_t total); + + // A complete directive waits: data() is its text (NUL-terminated), length() its size without the NUL + bool ready() const { return waiting; } + const char *data() const { return waiting ? text : nullptr; } + size_t length() const { return waiting ? size : 0; } + + // The waiting directive has been read: the buffer takes the next one + void release(); + + // The connection is gone: a message that was still arriving will not be completed. A directive that waits stays. + void cancelArrival(); + +private: + // What happens to the fragments of the message that is arriving + enum class Arrival : uint8_t + { + NONE, // no message is arriving + STORING, // they are copied into text + COMPARING, // they are compared with the complete directive in text, which they have matched so far + SKIPPING // they are dropped, the message was refused + }; + + Result start(const char *data, size_t length, size_t total); + Result proceed(const char *data, size_t length, size_t index, size_t total); + Result store(const char *data, size_t length); + Result refuse(Result reason); + void discard(); + + const size_t capacity; + char *text = nullptr; + size_t size = 0; // length of the message text holds or is receiving + size_t offset = 0; // bytes of the arriving message stored or compared so far + bool waiting = false; + Arrival arrival = Arrival::NONE; +}; + +// The messageIds of the last directives that were dispatched. A broker that mirrors its topics to a second broker +// and back delivers every message twice (the Alex2MQTT broker does), and a sketch must not act on a directive +// twice. Four ids are kept as 64-bit FNV-1a hashes: 32 bytes instead of the 148 that four UUID strings take. +class AlexaRecentIds +{ +public: + static const uint8_t CAPACITY = 4; + + // True when the id is one of the last CAPACITY ids given to this function. Otherwise it is remembered in place + // of the oldest one. An empty id is never remembered and never a repeat. + bool seenBefore(const char *id); + +private: + uint64_t hashes[CAPACITY] = {0, 0, 0, 0}; + uint8_t count = 0; + uint8_t next = 0; +}; + +namespace AlexaBridgeLogic +{ + // First second of 2024. time() counts from 1970 at boot until SNTP has answered, so anything earlier than the + // library itself means that the clock has not been set. + const time_t CLOCK_SET_AFTER = 1704067200; + + bool clockIsSet(time_t now); + + // Writes the instant as "YYYY-MM-DDThh:mm:ssZ" (UTC). Returns false and leaves an empty string when the buffer + // is smaller than ALEXA_TIMESTAMP_SIZE or the year has more than four digits. + bool formatTimestamp(time_t instant, char *buffer, size_t size); + + // The endpoint id in "//alexaDirective": a pointer into the topic and the id's length. + // nullptr for any other topic, among them /discover and the //alexaDirective_e token + // topic that 1.x boards use. + const char *directiveEndpoint(const char *topic, const char *rootTopic, size_t *length); + + // Decides whether a document may be published and measures it: TOO_LARGE when it is over maxLength bytes or + // ran out of memory while it was built (it would be sent truncated), EMPTY when it holds nothing, otherwise OK + // with the serialised size in length. + AlexaSendResult checkMessage(const JsonDocument &doc, size_t maxLength, size_t *length); +} + +#endif // ALEXA_BRIDGE_LOGIC_H diff --git a/src/AlexaCompat.h b/src/AlexaCompat.h new file mode 100644 index 0000000..86ed3b6 --- /dev/null +++ b/src/AlexaCompat.h @@ -0,0 +1,14 @@ +// Lets the sources that have no hardware dependency compile on a host (pio test -e native), where the +// program-memory helpers of the Arduino cores do not exist. On a board this is Arduino.h and nothing else. +#ifndef ALEXA_COMPAT_H +#define ALEXA_COMPAT_H + +#if defined(ARDUINO) +#include +#else +#ifndef PSTR +#define PSTR(text) (text) +#endif +#endif + +#endif // ALEXA_COMPAT_H diff --git a/src/AlexaDevice.cpp b/src/AlexaDevice.cpp index 9f1c278..95d6ca7 100644 --- a/src/AlexaDevice.cpp +++ b/src/AlexaDevice.cpp @@ -1,7 +1,8 @@ #include "AlexaDevice.h" +#include "AlexaLog.h" -AlexaDevice::AlexaDevice(const String& name, const String& rootTopic, const String& endpointId) - : name(name), endpointId(endpointId), rootTopic(rootTopic) { +AlexaDevice::AlexaDevice(const String& name, const String& rootTopic, const String& endpointId, AlexaTransport* transport) + : name(name), endpointId(endpointId), rootTopic(rootTopic), transport(transport) { // Initialize event arrays to nullptr for (int i = 0; i < MAX_EVENTS; ++i) { eventNames[i] = nullptr; // Set all event names to nullptr @@ -22,6 +23,10 @@ String AlexaDevice::getEndpointId() const { return endpointId; } +bool AlexaDevice::hasEndpointId(const char* id, size_t length) const { + return id != nullptr && endpointId.length() == length && memcmp(endpointId.c_str(), id, length) == 0; +} + DisplayCategory AlexaDevice::getDisplayCategory() const { return displayCategory; } @@ -78,8 +83,7 @@ AlexaInterface* AlexaDevice::addCapability(AlexaInterfaceType type) { // If the interface doesn't exist, create a new one capabilities.emplace_back(type); - Serial.print("Created new interface with type: "); - Serial.println(capabilities.back().getTypeString()); + ALEX2ESP_LOGI("%s: capability %s added", endpointId.c_str(), capabilities.back().getTypeString().c_str()); // Return a pointer to the newly created device return &capabilities.back(); @@ -158,9 +162,8 @@ void AlexaDevice::triggerEvent(const char* eventName, const JsonDocument& direct return; } } - // If no event was found, print an error message + // No handler: the directive goes unanswered, which Alexa reports as a device that does not respond if (warnIfMissing) { - Serial.print("No event registered for: "); - Serial.println(eventName); + ALEX2ESP_LOGE("%s: no handler registered for %s, directive not handled", endpointId.c_str(), eventName); } } diff --git a/src/AlexaDevice.h b/src/AlexaDevice.h index 65157aa..4f87793 100644 --- a/src/AlexaDevice.h +++ b/src/AlexaDevice.h @@ -10,6 +10,7 @@ #include #include #include "AlexaStatusMessage.h" +#include "AlexaTransport.h" #define MAX_EVENTS 10 @@ -147,13 +148,18 @@ public: class AlexaDevice { public: - AlexaDevice(const String& name, const String& rootTopic, const String& endpointId); + // Devices are created by Alex2ESP::getDevice(), which passes the bridge as the transport that publishes the + // device's reports. A device built without one can describe itself, but its reports cannot be sent. + AlexaDevice(const String& name, const String& rootTopic, const String& endpointId, AlexaTransport* transport = nullptr); void setName(const String& name); String getName() const; String getEndpointId() const; + // Compares the endpoint id with `length` characters at `id` (not NUL-terminated: a part of an MQTT topic) + bool hasEndpointId(const char* id, size_t length) const; + DisplayCategory getDisplayCategory() const; void setDisplayCategory(DisplayCategory category); @@ -183,13 +189,14 @@ public: void triggerEvent(const char* eventName, const JsonDocument& directive, const AlexaInterfaceType& type, bool warnIfMissing = true) const; AlexaStatusMessage buildStatusMessage(const String& correlationToken,const bool isResponse=false) { - return AlexaStatusMessage(correlationToken,rootTopic,endpointId,isResponse); + return AlexaStatusMessage(correlationToken,rootTopic,endpointId,isResponse,transport); } private: String name; String endpointId; String rootTopic; + AlexaTransport* transport; const char* eventNames[MAX_EVENTS]; void (*eventCallbacks[MAX_EVENTS])(const JsonDocument&, const AlexaInterfaceType&); diff --git a/src/AlexaInterface.h b/src/AlexaInterface.h index 88cc4d3..46c03f7 100644 --- a/src/AlexaInterface.h +++ b/src/AlexaInterface.h @@ -6,6 +6,7 @@ #include #include #include +#include "AlexaLog.h" // Enum for Alexa Interface Types enum class AlexaInterfaceType @@ -617,8 +618,7 @@ public: if (deserializeJson(payloadDoc, directive.payload) == DeserializationError::Ok && payloadDoc.is()) { directiveObj["payload"] = payloadDoc.as(); } else { - Serial.print("[Alex2ESP] action mapping payload is not a JSON object, omitted: "); - Serial.println(directive.payload); + ALEX2ESP_LOGE("action mapping payload is not a JSON object, omitted: %s", directive.payload.c_str()); } } return doc; diff --git a/src/AlexaLog.cpp b/src/AlexaLog.cpp new file mode 100644 index 0000000..460f6d5 --- /dev/null +++ b/src/AlexaLog.cpp @@ -0,0 +1,49 @@ +#include "AlexaLog.h" +#include + +AlexaLogLevel AlexaLog::level = AlexaLogLevel::INFO; +Print *AlexaLog::output = &Serial; + +void AlexaLog::setLevel(AlexaLogLevel newLevel) +{ + level = newLevel; +} + +AlexaLogLevel AlexaLog::getLevel() +{ + return level; +} + +void AlexaLog::setOutput(Print *newOutput) +{ + output = newOutput; +} + +bool AlexaLog::enabled(AlexaLogLevel lineLevel) +{ + return output != nullptr && lineLevel != AlexaLogLevel::NONE && lineLevel <= level; +} + +void AlexaLog::write(AlexaLogLevel lineLevel, PGM_P format, ...) +{ + if (!enabled(lineLevel)) + { + return; + } + + // vsnprintf and not vsnprintf_P: on the ESP8266 vsnprintf reads a PROGMEM format itself (vsnprintf_P is a jump + // to it), and vsnprintf_P comes in one object with printf_P and sprintf_P, which link the FILE-based printf: + // 5.8 KB of flash. + char line[160]; + va_list arguments; + va_start(arguments, format); + int written = vsnprintf(line, sizeof(line), format, arguments); + va_end(arguments); + if (written < 0) + { + return; + } + + output->print(lineLevel == AlexaLogLevel::ERROR ? F("[Alex2ESP] error: ") : F("[Alex2ESP] ")); + output->println(line); +} diff --git a/src/AlexaLog.h b/src/AlexaLog.h new file mode 100644 index 0000000..0e044f8 --- /dev/null +++ b/src/AlexaLog.h @@ -0,0 +1,56 @@ +// Serial diagnostics of the library. A line is printed when its level is within the level chosen at run time +// (Alex2ESP::setLogLevel, INFO unless changed) and within the ceiling compiled in (ALEX2ESP_LOG_MAX); a line above +// the ceiling costs no flash. Whatever the library refuses or drops is reported at ERROR. Credentials and +// correlation tokens are never printed. +#ifndef ALEXA_LOG_H +#define ALEXA_LOG_H + +#include + +enum class AlexaLogLevel : uint8_t +{ + NONE = 0, // nothing + ERROR = 1, // something was refused, dropped or not sent + INFO = 2, // session, discovery and one line per directive + DEBUG = 3 // sizes and free heap per message; needs ALEX2ESP_LOG_MAX=3 +}; + +// Highest level that is compiled in: 0 none, 1 ERROR, 2 INFO, 3 DEBUG. Override with -DALEX2ESP_LOG_MAX= in +// build_flags. The Arduino IDE has no per-sketch flags and builds with this default. +#ifndef ALEX2ESP_LOG_MAX +#define ALEX2ESP_LOG_MAX 2 +#endif + +class AlexaLog +{ +public: + static void setLevel(AlexaLogLevel level); + static AlexaLogLevel getLevel(); + + // Where the lines go: Serial unless changed, nullptr for nowhere + static void setOutput(Print *output); + + static bool enabled(AlexaLogLevel level); + + // printf-style, the format is a PROGMEM string. A line longer than 159 characters is cut there. + static void write(AlexaLogLevel level, PGM_P format, ...); + +private: + static AlexaLogLevel level; + static Print *output; +}; + +#define ALEX2ESP_LOG(levelNumber, levelName, format, ...) \ + do \ + { \ + if (ALEX2ESP_LOG_MAX >= (levelNumber) && AlexaLog::enabled(levelName)) \ + { \ + AlexaLog::write(levelName, PSTR(format), ##__VA_ARGS__); \ + } \ + } while (0) + +#define ALEX2ESP_LOGE(format, ...) ALEX2ESP_LOG(1, AlexaLogLevel::ERROR, format, ##__VA_ARGS__) +#define ALEX2ESP_LOGI(format, ...) ALEX2ESP_LOG(2, AlexaLogLevel::INFO, format, ##__VA_ARGS__) +#define ALEX2ESP_LOGD(format, ...) ALEX2ESP_LOG(3, AlexaLogLevel::DEBUG, format, ##__VA_ARGS__) + +#endif // ALEXA_LOG_H diff --git a/src/AlexaStatusMessage.cpp b/src/AlexaStatusMessage.cpp index 3ab44b7..05232be 100644 --- a/src/AlexaStatusMessage.cpp +++ b/src/AlexaStatusMessage.cpp @@ -1,3 +1,90 @@ #include "AlexaStatusMessage.h" +#include "AlexaLog.h" -char AlexaStatusMessage::outputString[MAX_STATUS_REPORT_SIZE]; +// What a 1.x sketch wrote where the time belongs when it built a property by hand: the backend's HTTP route +// replaced it. Reports leave over MQTT now, where nothing rewrites them, so the library fills the time in. +static const char TIME_PLACEHOLDER[] PROGMEM = "{REPLACE_WITH_DATETIME}"; + +AlexaStatusMessage::AlexaStatusMessage(const String &correlationToken, const String &rootTopic, const String &endpointId, const bool isResponse, AlexaTransport *transport) + : rootTopic(rootTopic), endpointId(endpointId), transport(transport) +{ + JsonObject event = doc["event"].to(); + + JsonObject event_header = event["header"].to(); + event_header["namespace"] = "Alexa"; + if (isResponse) + { + event_header["name"] = "Response"; + } + else + { + event_header["name"] = "StateReport"; + } + event_header["payloadVersion"] = "3"; + event_header["messageId"] = generateMessageId(); + event_header["correlationToken"] = correlationToken; + + JsonObject event_endpoint = event["endpoint"].to(); + event_endpoint["endpointId"] = endpointId; + event["payload"].to(); // required by Alexa.Response / StateReport, empty when there is nothing to add + contextProperties = doc["context"]["properties"].to(); +} + +AlexaStatusMessage &AlexaStatusMessage::AddContextProp(const JsonObject &property) +{ + if (contextProperties.add(property)) + { + JsonObject added = contextProperties[contextProperties.size() - 1]; + const char *timeOfSample = added["timeOfSample"] | ""; + if (timeOfSample[0] == '\0' || strcmp_P(timeOfSample, TIME_PLACEHOLDER) == 0) + { + setTimeOfSample(added); + } + } + return *this; // a property that did not fit leaves the document marked as overflowed: send() refuses it +} + +bool AlexaStatusMessage::send() +{ + if (transport == nullptr) + { + ALEX2ESP_LOGE("report for %s not sent: its device was not created by getDevice()", endpointId.c_str()); + doc.clear(); + return false; + } + + String topic = rootTopic + "/" + endpointId + "/alexaResponce"; + AlexaSendResult result = transport->publish(topic.c_str(), doc); + doc.clear(); + return result == AlexaSendResult::OK; +} + +// TODO: make real uuid4 gen function +String AlexaStatusMessage::generateMessageId() +{ + char buffer[38]; + const char charset[] = "0123456789abcdefABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz"; // Allowed characters + + // Seed the random number generator (optional) + srand(static_cast(time(nullptr))); + + // Generate 37 random characters + for (int i = 0; i < 37; ++i) + { + buffer[i] = charset[rand() % (sizeof(charset) - 1)]; // Pick a random character + } + + buffer[37] = '\0'; // Null-terminate the string + return String(buffer); +} + +void AlexaStatusMessage::setTimeOfSample(JsonObject property) +{ + // A char array is copied into the document; the buffer does not have to outlive this call + char now[ALEXA_TIMESTAMP_SIZE] = ""; + if (transport != nullptr) + { + transport->timestamp(now, sizeof(now)); + } + property["timeOfSample"] = now; +} diff --git a/src/AlexaStatusMessage.h b/src/AlexaStatusMessage.h index c8c215c..37b6a42 100644 --- a/src/AlexaStatusMessage.h +++ b/src/AlexaStatusMessage.h @@ -4,10 +4,9 @@ #include #include #include "AlexaInterface.h" +#include "AlexaTransport.h" #include "AlexaUtils.h" -#define MAX_STATUS_REPORT_SIZE 2048 - enum class EndpointHealth { OK, @@ -28,32 +27,9 @@ enum class TemperatureSensorScale class AlexaStatusMessage { public: - AlexaStatusMessage(const String &correlationToken, const String &rootTopic, const String &endpointId, const bool isResponse) - { - this->endpointId = endpointId; - this->rootTopic = rootTopic; - - JsonObject event = doc["event"].to(); - - JsonObject event_header = event["header"].to(); - event_header["namespace"] = "Alexa"; - if (isResponse) - { - event_header["name"] = "Response"; - } - else - { - event_header["name"] = "StateReport"; - } - event_header["payloadVersion"] = "3"; - event_header["messageId"] = generateMessageId(); - event_header["correlationToken"] = correlationToken; - - JsonObject event_endpoint = event["endpoint"].to(); - event_endpoint["endpointId"] = endpointId; - event["payload"].to(); // required by Alexa.Response / StateReport, empty when there is nothing to add - contextProperties = doc["context"]["properties"].to(); - } + // Built by AlexaDevice::buildStatusMessage(), which passes the bridge as the transport: it publishes the report + // and supplies the time of every property. A message without a transport cannot be sent. + AlexaStatusMessage(const String &correlationToken, const String &rootTopic, const String &endpointId, const bool isResponse, AlexaTransport *transport = nullptr); AlexaStatusMessage &AddHealthProp(EndpointHealth endpointHealth, unsigned int uncertaintyInMs = 0) { @@ -93,66 +69,29 @@ public: return AddProperty(AlexaInterfaceType::TOGGLE_CONTROLLER, "toggleState", value, uncertaintyInMs,instanceName); } - AlexaStatusMessage &AddContextProp(const JsonObject &property) - { - contextProperties.add(property); - return *this; // Return a reference to the current object - } + // Adds a property the sketch built itself. Its timeOfSample is set to the current time when the object has + // none or still carries the 1.x placeholder "{REPLACE_WITH_DATETIME}". + AlexaStatusMessage &AddContextProp(const JsonObject &property); - // Queue the report for delivery. Returns false (and says so on Serial) when the report does not fit the - // MAX_STATUS_REPORT_SIZE buffer or the send queue is full: a report is never sent truncated. - bool send() - { - doc.shrinkToFit(); - String topic = rootTopic + "/" + endpointId + "/alexaResponce"; - size_t needed = measureJson(doc); - if (needed > sizeof(outputString) - 1) - { - Serial.printf("[Alex2ESP] status report for %s is %u bytes, limit is %u - not sent\n", - endpointId.c_str(), (unsigned)needed, (unsigned)(sizeof(outputString) - 1)); - doc.clear(); - return false; - } - serializeJson(doc, outputString); - doc.clear(); - if (!AlexaUtils::enqueue(outputString, topic.c_str())) - { - Serial.println("[Alex2ESP] send queue full, status report dropped"); - return false; - } - return true; - } + // Publishes the report on //alexaResponce now. Returns false when nothing was sent: no + // session with the broker, the MQTT client or the heap cannot take the report, or it is over + // ALEX2ESP_MAX_MESSAGE bytes. The reason is on Serial; a report is never sent truncated. + bool send(); private: String rootTopic; String endpointId; + AlexaTransport *transport; JsonDocument doc; JsonArray contextProperties; - static char outputString[MAX_STATUS_REPORT_SIZE]; - // TODO: make real uuid4 gen function - String generateMessageId() - { - char buffer[38]; - const char charset[] = "0123456789abcdefABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz"; // Allowed characters - - // Seed the random number generator (optional) - srand(static_cast(time(nullptr))); - - // Generate 37 random characters - for (int i = 0; i < 37; ++i) - { - buffer[i] = charset[rand() % (sizeof(charset) - 1)]; // Pick a random character - } - - buffer[37] = '\0'; // Null-terminate the string - return String(buffer); - } + String generateMessageId(); + void setTimeOfSample(JsonObject property); template AlexaStatusMessage &AddProperty(AlexaInterfaceType type, const String &propertyName, const T &value, unsigned int uncertaintyInMs = 0,const String &instanceName="") { - JsonDocument prop; + JsonObject prop = contextProperties.add(); prop["namespace"] = AlexaInterfaceUtils::toString(type); prop["name"] = propertyName; @@ -161,10 +100,8 @@ private: } prop["value"] = value; - prop["timeOfSample"] = "{REPLACE_WITH_DATETIME}"; + setTimeOfSample(prop); prop["uncertaintyInMilliseconds"] = uncertaintyInMs; - - contextProperties.add(prop.as()); return *this; } }; diff --git a/src/AlexaTransport.h b/src/AlexaTransport.h new file mode 100644 index 0000000..47ad4e1 --- /dev/null +++ b/src/AlexaTransport.h @@ -0,0 +1,43 @@ +// What a device and its reports need from the bridge: a publisher and a clock. Alex2ESP implements both over MQTT +// and SNTP; AlexaDevice and AlexaStatusMessage only ever see this interface. +#ifndef ALEXA_TRANSPORT_H +#define ALEXA_TRANSPORT_H + +#include +#include +#include + +// Largest message the bridge publishes, in bytes of JSON: a report or the discovery object of one device. +// Override with -DALEX2ESP_MAX_MESSAGE= in build_flags. +#ifndef ALEX2ESP_MAX_MESSAGE +#define ALEX2ESP_MAX_MESSAGE 2047 +#endif + +// "YYYY-MM-DDThh:mm:ssZ" and the terminating NUL +#define ALEXA_TIMESTAMP_SIZE 21 + +// Outcome of a publish. Anything but OK means that nothing was sent; the reason is on Serial at the ERROR level. +enum class AlexaSendResult : uint8_t +{ + OK, // handed to the MQTT client + NOT_CONNECTED, // no session with the broker + REFUSED, // the MQTT client or the heap cannot take the message now; a later attempt may succeed + TOO_LARGE, // over ALEX2ESP_MAX_MESSAGE, or the document ran out of memory while it was built + EMPTY // the document holds nothing +}; + +class AlexaTransport +{ +public: + // Serialises the document and publishes it on the topic. A document is sent whole or not at all. + virtual AlexaSendResult publish(const char *topic, JsonDocument &doc) = 0; + + // Writes the current time as ISO 8601 in UTC, e.g. "2026-09-28T13:05:09Z". The buffer holds + // ALEXA_TIMESTAMP_SIZE bytes or more; a smaller one is left as an empty string. + virtual void timestamp(char *buffer, size_t size) = 0; + +protected: + ~AlexaTransport() {} +}; + +#endif // ALEXA_TRANSPORT_H diff --git a/src/AlexaUtils.cpp b/src/AlexaUtils.cpp deleted file mode 100644 index cbfff31..0000000 --- a/src/AlexaUtils.cpp +++ /dev/null @@ -1,227 +0,0 @@ -#include "AlexaUtils.h" - -// #define Alex2ESP_DEBUG - -uint8_t AlexaUtils::nextMessageId = 0; - -char AlexaUtils::receiveQueue[MAX_QUEUE_LENGTH][16]={}; -char AlexaUtils::receivePayload[MAX_PAYLOAD_LENGTH]; -char AlexaUtils::topicQueue[MAX_QUEUE_LENGTH][MAX_TOPIC_LENGTH] = {}; -char AlexaUtils::packetQueue[MAX_QUEUE_LENGTH][MAX_PACKET_LENGTH] = {}; -char AlexaUtils::combinedData[MAX_TOPIC_LENGTH + MAX_PACKET_LENGTH + 32]; - -int AlexaUtils::queueStart = 0; -int AlexaUtils::queueEnd = 0; -int AlexaUtils::queueCount = 0; - - -int AlexaUtils::queueStartReceive = 0; -int AlexaUtils::queueEndReceive = 0; -int AlexaUtils::queueCountReceive = 0; - -// Enqueue a packet into the receive queue -bool AlexaUtils::enqueueReceive(const char* packet) { - if (isReceiveQueueFull()) { - return false; // Queue is full - } - AlexaUtils::log("Enqueing: "); - AlexaUtils::logln(packet); - - // Copy the packet into the receive queue - strncpy(receiveQueue[queueEndReceive], packet, sizeof(receiveQueue[0]) - 1); - receiveQueue[queueEndReceive][sizeof(receiveQueue[0]) - 1] = '\0'; // Ensure null-termination - - // Update the queueEnd and queueCount - queueEndReceive = (queueEndReceive + 1) % MAX_QUEUE_LENGTH; - queueCountReceive++; - - return true; -} - -// Dequeue a packet from the receive queue -bool AlexaUtils::dequeueReceive(String& packet,bool remove) { - if (isReceiveQueueEmpty()) { - return false; // Queue is empty - } - - // Read the front of the receive queue - packet = String(receiveQueue[queueStartReceive]); - - if(remove){ - // Clear the dequeued slot (optional but good for debugging) - memset(receiveQueue[queueStartReceive], 0, sizeof(receiveQueue[0])); - - // Update the queueStart and queueCount - queueStartReceive = (queueStartReceive + 1) % MAX_QUEUE_LENGTH; - queueCountReceive--; - } - return true; -} - -// Check if the receive queue is empty -bool AlexaUtils::isReceiveQueueEmpty() { - return queueCountReceive == 0; -} - -// Check if the receive queue is full -bool AlexaUtils::isReceiveQueueFull() { - return queueCountReceive == MAX_QUEUE_LENGTH; -} - - -// Enqueue a topic and packet into the queue -bool AlexaUtils::enqueue(const char* packet,const char* topic) { - if (isQueueFull()) { - return false; // Queue is full - } - - // Refuse anything that would be cut short by the slot size - a truncated JSON packet is worthless - if (strlen(topic) > MAX_TOPIC_LENGTH - 1 || strlen(packet) > MAX_PACKET_LENGTH - 1) { - Serial.printf("[Alex2ESP] packet of %u bytes does not fit the %d byte send slot - dropped\n", - (unsigned)strlen(packet), MAX_PACKET_LENGTH - 1); - return false; - } - - // Copy the topic and packet into the respective arrays - strncpy(topicQueue[queueEnd], topic, MAX_TOPIC_LENGTH - 1); - topicQueue[queueEnd][MAX_TOPIC_LENGTH - 1] = '\0'; // Ensure null-termination - - strncpy(packetQueue[queueEnd], packet, MAX_PACKET_LENGTH - 1); - packetQueue[queueEnd][MAX_PACKET_LENGTH - 1] = '\0'; // Ensure null-termination - - // Update the queueEnd and queueCount - queueEnd = (queueEnd + 1) % MAX_QUEUE_LENGTH; - queueCount++; - - return true; -} - -// Dequeue a topic and packet from the queue -bool AlexaUtils::dequeue(String& topic, String& packet) { - if (isQueueEmpty()) { - return false; // Queue is empty - } - - // Read the front of the queue - topic = String(topicQueue[queueStart]); - packet = String(packetQueue[queueStart]); - - // Clear the dequeued slot (optional but good for debugging) - memset(topicQueue[queueStart], 0, MAX_TOPIC_LENGTH); - memset(packetQueue[queueStart], 0, MAX_PACKET_LENGTH); - - // Update the queueStart and queueCount - queueStart = (queueStart + 1) % MAX_QUEUE_LENGTH; - queueCount--; - - return true; -} - - -const char* AlexaUtils::dequeueVals(bool remove) { - if (isQueueEmpty()) { - return nullptr; // Queue is empty - } - - // Reset combinedData before use to ensure no leftover data from previous calls - memset(combinedData, 0, sizeof(combinedData)); - - // Get lengths of the topic and packet - int topicLen = strlen(topicQueue[queueStart]); - int packetLen = strlen(packetQueue[queueStart]); - - // Calculate available space in the combinedData buffer - size_t availableSpace = sizeof(combinedData) - 1; // Reserve space for the null terminator - - // Check if the combined data can fit in the buffer - if ((topicLen + packetLen + strlen(MESSAGE_PAYLOAD_SPLIT)) > availableSpace) { - // Return null if the combined data is too large for the buffer - AlexaUtils::logln("Error: Combined data exceeds buffer size."); - return nullptr; - } - - // Concatenate the topic, separator, and packet into combinedData - snprintf(combinedData, sizeof(combinedData), "%s%s%s", topicQueue[queueStart], MESSAGE_PAYLOAD_SPLIT, packetQueue[queueStart]); - - if (remove) { - // Clear the dequeued slot - memset(topicQueue[queueStart], 0, MAX_TOPIC_LENGTH); - memset(packetQueue[queueStart], 0, MAX_PACKET_LENGTH); - - // Update the queueStart and queueCount - queueStart = (queueStart + 1) % MAX_QUEUE_LENGTH; - queueCount--; - } - - // Return the pointer to the combined data - return combinedData; -} - - - -// Check if the queue is empty -bool AlexaUtils::isQueueEmpty() { - return queueCount == 0; -} - -// Check if the queue is full -bool AlexaUtils::isQueueFull() { - return queueCount == MAX_QUEUE_LENGTH; -} - - - -void AlexaUtils::log(const char *string) { -#ifdef Alex2ESP_DEBUG - Serial.print(string); -#endif -} - -void AlexaUtils::logln(const char *string) { -#ifdef Alex2ESP_DEBUG - Serial.println(string); -#endif -} - -void AlexaUtils::log(int num) { -#ifdef Alex2ESP_DEBUG - Serial.print(num); -#endif -} - -void AlexaUtils::logln(int num) { -#ifdef Alex2ESP_DEBUG - Serial.println(num); -#endif -} - -void AlexaUtils::log(uint16_t num) { -#ifdef Alex2ESP_DEBUG - Serial.print(num); -#endif -} - -void AlexaUtils::logln(uint16_t num) { -#ifdef Alex2ESP_DEBUG - Serial.println(num); -#endif -} - -void AlexaUtils::log(uint8_t num) { -#ifdef Alex2ESP_DEBUG - Serial.print(num); -#endif -} - -void AlexaUtils::logln(uint8_t num) { -#ifdef Alex2ESP_DEBUG - Serial.println(num); -#endif -} - - -void AlexaUtils::logln(String string) { -#ifdef Alex2ESP_DEBUG - Serial.println(string); -#endif -} diff --git a/src/AlexaUtils.h b/src/AlexaUtils.h index 74dfd19..fffdaa8 100644 --- a/src/AlexaUtils.h +++ b/src/AlexaUtils.h @@ -1,89 +1,26 @@ #ifndef ALEXA_UTILS_H #define ALEXA_UTILS_H -#include -#include #include -#include -#include // Include for strcpy and memset +// A 1.x name, kept for sketches that print the memory figures. The send and receive queues that lived here +// belonged to the HTTP fallback and went with it in 1.2.0; the library's own diagnostics are in AlexaLog.h. class AlexaUtils { public: - static const int MAX_PAYLOAD_LENGTH = 2048; // Maximum length for receive packets - static char receivePayload[MAX_PAYLOAD_LENGTH]; - - static bool enqueueReceive(const char* packet); - static bool dequeueReceive(String& packet,bool remove); - static bool isReceiveQueueEmpty(); - static bool isReceiveQueueFull(); - - // Queue a packet for the HTTP sender. Returns false (nothing queued) when the queue is full or the - // packet/topic would not fit its slot - a packet is never truncated. - static bool enqueue(const char* packet, const char* topic); - static bool dequeue(String& topic, String& packet); - static bool isQueueEmpty(); - static bool isQueueFull(); - static const char* dequeueVals(bool remove); - - static void log(const char *string); - static void logln(const char *string); - - static void log(int num); - static void logln(int num); - - static void log(uint16_t num); - static void logln(uint16_t num); - - static void log(uint8_t num); - static void logln(uint8_t num); - - static void logln(String string); - - // Incrementing message ID - static uint8_t nextMessageId; - - // Static function to set the MQTT client reference static void printMemoryInfo() { - size_t freeHeap = ESP.getFreeHeap(); #ifdef ESP8266 size_t freeStack = ESP.getFreeContStack(); // Requires ESP8266 core 3.0.0+ #else size_t freeStack = 0; // ESP32: untested, no getFreeContStack() #endif - size_t sketchSize = ESP.getSketchSize(); - size_t freeSketchSpace = ESP.getFreeSketchSpace(); - - // Print in custom format Serial.printf("freeHeap: %u, freeStack: %u, freeROM: %u, usedROM: %u\n", - freeHeap, - freeStack, - freeSketchSpace, - sketchSize); + (unsigned)ESP.getFreeHeap(), + (unsigned)freeStack, + (unsigned)ESP.getFreeSketchSpace(), + (unsigned)ESP.getSketchSize()); } - - -private: - static const int MAX_QUEUE_LENGTH = 5; // Define maximum queue length - static const int MAX_TOPIC_LENGTH = 128; // Maximum length for topics - static const int MAX_PACKET_LENGTH = 2048; // Maximum length for sending packets - static constexpr char MESSAGE_PAYLOAD_SPLIT[32] = "{MESSAGE_PAYLOAD_SPLIT}"; - - - // Static arrays for topics and packets - static char topicQueue[MAX_QUEUE_LENGTH][MAX_TOPIC_LENGTH]; - static char packetQueue[MAX_QUEUE_LENGTH][MAX_PACKET_LENGTH]; - static char combinedData[MAX_TOPIC_LENGTH + MAX_PACKET_LENGTH + 32]; - - static char receiveQueue[MAX_QUEUE_LENGTH][16]; - static int queueStart; // Index of the front of the queue - static int queueEnd; // Index of the back of the queue - static int queueCount; // Number of items in the queue - - static int queueStartReceive; - static int queueEndReceive; - static int queueCountReceive; }; #endif // ALEXA_UTILS_H diff --git a/test/test_bridge_logic/test_main.cpp b/test/test_bridge_logic/test_main.cpp new file mode 100644 index 0000000..a4bfc7a --- /dev/null +++ b/test/test_bridge_logic/test_main.cpp @@ -0,0 +1,575 @@ +// Host tests of the bridge logic (src/AlexaBridgeLogic.cpp): what the MQTT callback decides about the fragments a +// broker delivers, which directives are repeats, what may be published, topics and time stamps. +// pio test -e native +#include +#include +#include +#include +#include +#include "AlexaBridgeLogic.h" + +typedef AlexaDirectiveBuffer::Result Result; + +#define ASSERT_RESULT(expected, actual) TEST_ASSERT_EQUAL_INT(static_cast(expected), static_cast(actual)) + +// A directive as Alex2MQTT publishes it on //alexaDirective +static std::string directive(const char *name, const char *messageId) +{ + return std::string("{\"header\":{\"namespace\":\"Alexa.PowerController\",\"name\":\"") + name + + "\",\"payloadVersion\":\"3\",\"messageId\":\"" + messageId + + "\",\"correlationToken\":\"AAAAAAAAAQBnlqNbYnB0dHNmYW5zbGF0ZQ==\"}," + "\"endpoint\":{\"endpointId\":\"ESP-01\",\"cookie\":{}},\"payload\":{}}"; +} + +static const char ID_1[] = "1bd5d003-31b9-476f-ad03-71d471922820"; +static const char ID_2[] = "7c0e1f6a-52d4-4b8e-9a3c-0d9f4e2b6a11"; + +// Hands the message over as the MQTT client does: pieces of fragmentSize bytes with their offset and the total. +// Returns the result of the last piece; every piece before it has to be INCOMPLETE. +static Result deliver(AlexaDirectiveBuffer &buffer, const std::string &message, size_t fragmentSize) +{ + Result result = Result::INCOMPLETE; + for (size_t index = 0; index < message.size(); index += fragmentSize) + { + ASSERT_RESULT(Result::INCOMPLETE, result); + size_t length = message.size() - index < fragmentSize ? message.size() - index : fragmentSize; + result = buffer.append(message.data() + index, length, index, message.size()); + } + return result; +} + +static void assertWaiting(const AlexaDirectiveBuffer &buffer, const std::string &message) +{ + TEST_ASSERT_TRUE(buffer.ready()); + TEST_ASSERT_EQUAL_UINT(message.size(), buffer.length()); + TEST_ASSERT_EQUAL_STRING(message.c_str(), buffer.data()); // also proves the terminating NUL +} + +void setUp() {} +void tearDown() {} + +// --- reassembly --- + +void test_directive_in_one_piece_is_complete() +{ + AlexaDirectiveBuffer buffer(2047); + std::string turnOn = directive("TurnOn", ID_1); + + ASSERT_RESULT(Result::COMPLETE, buffer.append(turnOn.data(), turnOn.size(), 0, turnOn.size())); + assertWaiting(buffer, turnOn); + + buffer.release(); + TEST_ASSERT_FALSE(buffer.ready()); + TEST_ASSERT_NULL(buffer.data()); + TEST_ASSERT_EQUAL_UINT(0, buffer.length()); +} + +void test_directive_in_fragments_is_reassembled_and_parses() +{ + std::string turnOn = directive("TurnOn", ID_1); + + // Every fragment size from one byte per fragment to the message in two halves + for (size_t fragmentSize = 1; fragmentSize < turnOn.size(); fragmentSize++) + { + AlexaDirectiveBuffer buffer(2047); + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, turnOn, fragmentSize)); + assertWaiting(buffer, turnOn); + } + + AlexaDirectiveBuffer buffer(2047); + deliver(buffer, turnOn, 100); + JsonDocument parsed; + TEST_ASSERT_TRUE(deserializeJson(parsed, buffer.data(), buffer.length()) == DeserializationError::Ok); + TEST_ASSERT_EQUAL_STRING("TurnOn", parsed["header"]["name"]); + TEST_ASSERT_EQUAL_STRING("ESP-01", parsed["endpoint"]["endpointId"]); +} + +void test_nothing_waits_before_the_last_fragment() +{ + AlexaDirectiveBuffer buffer(2047); + std::string turnOn = directive("TurnOn", ID_1); + + ASSERT_RESULT(Result::INCOMPLETE, buffer.append(turnOn.data(), 50, 0, turnOn.size())); + TEST_ASSERT_FALSE(buffer.ready()); + TEST_ASSERT_NULL(buffer.data()); +} + +// --- limits --- + +void test_directive_at_the_limit_is_accepted() +{ + AlexaDirectiveBuffer buffer(2047); + std::string atLimit(2047, 'x'); + + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, atLimit, 536)); + assertWaiting(buffer, atLimit); +} + +void test_directive_over_the_limit_is_refused_once() +{ + AlexaDirectiveBuffer buffer(2047); + std::string tooLarge(2048, 'x'); + + ASSERT_RESULT(Result::TOO_LARGE, buffer.append(tooLarge.data(), 536, 0, tooLarge.size())); + // Its other fragments still arrive: dropped without another report + ASSERT_RESULT(Result::IGNORED, buffer.append(tooLarge.data() + 536, 536, 536, tooLarge.size())); + ASSERT_RESULT(Result::IGNORED, buffer.append(tooLarge.data() + 1072, 976, 1072, tooLarge.size())); + TEST_ASSERT_FALSE(buffer.ready()); +} + +void test_refused_directive_does_not_block_the_next_one() +{ + AlexaDirectiveBuffer buffer(2047); + std::string tooLarge(5000, 'x'); + std::string turnOn = directive("TurnOn", ID_1); + + ASSERT_RESULT(Result::TOO_LARGE, buffer.append(tooLarge.data(), 1460, 0, tooLarge.size())); + ASSERT_RESULT(Result::IGNORED, buffer.append(tooLarge.data() + 1460, 1460, 1460, tooLarge.size())); + // The connection drops before the rest of it arrives; the next message is a normal directive + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, turnOn, 100)); + assertWaiting(buffer, turnOn); +} + +void test_empty_message_is_reported() +{ + AlexaDirectiveBuffer buffer(2047); + + // AsyncMqttClient hands a message without a payload over as (nullptr, 0, 0, 0) + ASSERT_RESULT(Result::EMPTY, buffer.append(nullptr, 0, 0, 0)); + TEST_ASSERT_FALSE(buffer.ready()); +} + +// --- one directive at a time --- + +void test_second_directive_before_loop_is_dropped_and_the_first_kept() +{ + AlexaDirectiveBuffer buffer(2047); + std::string turnOn = directive("TurnOn", ID_1); + std::string turnOff = directive("TurnOff", ID_2); + + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, turnOn, 100)); + ASSERT_RESULT(Result::BUSY, buffer.append(turnOff.data(), 100, 0, turnOff.size())); + ASSERT_RESULT(Result::IGNORED, buffer.append(turnOff.data() + 100, turnOff.size() - 100, 100, turnOff.size())); + assertWaiting(buffer, turnOn); + + // Once loop() has read the first, the buffer takes directives again + buffer.release(); + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, turnOff, 100)); + assertWaiting(buffer, turnOff); +} + +void test_second_directive_of_the_same_size_is_dropped() +{ + AlexaDirectiveBuffer buffer(2047); + std::string first = directive("TurnOn", ID_1); + std::string second = directive("TurnOn", ID_2); // same length, another messageId + TEST_ASSERT_EQUAL_UINT(first.size(), second.size()); + + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, first, 100)); + ASSERT_RESULT(Result::BUSY, buffer.append(second.data(), second.size(), 0, second.size())); + assertWaiting(buffer, first); +} + +// --- the repeat a mirroring broker delivers --- + +void test_repeat_of_the_waiting_directive_is_a_duplicate() +{ + AlexaDirectiveBuffer buffer(2047); + std::string turnOn = directive("TurnOn", ID_1); + + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, turnOn, 100)); + ASSERT_RESULT(Result::DUPLICATE, buffer.append(turnOn.data(), turnOn.size(), 0, turnOn.size())); + ASSERT_RESULT(Result::DUPLICATE, deliver(buffer, turnOn, 64)); + assertWaiting(buffer, turnOn); +} + +void test_repeat_that_differs_in_a_later_fragment_is_dropped() +{ + AlexaDirectiveBuffer buffer(2047); + std::string first = directive("TurnOn", ID_1); + std::string second = first; + second[second.size() - 10] = '#'; // the same up to its last fragment + + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, first, 100)); + ASSERT_RESULT(Result::INCOMPLETE, buffer.append(second.data(), 100, 0, second.size())); + ASSERT_RESULT(Result::INCOMPLETE, buffer.append(second.data() + 100, 100, 100, second.size())); + ASSERT_RESULT(Result::BUSY, buffer.append(second.data() + 200, second.size() - 200, 200, second.size())); + assertWaiting(buffer, first); +} + +void test_repeat_still_arriving_when_the_first_is_read() +{ + AlexaDirectiveBuffer buffer(2047); + std::string turnOn = directive("TurnOn", ID_1); + + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, turnOn, 100)); + ASSERT_RESULT(Result::INCOMPLETE, buffer.append(turnOn.data(), 100, 0, turnOn.size())); + buffer.release(); // loop() ran between two fragments of the repeat + TEST_ASSERT_FALSE(buffer.ready()); + ASSERT_RESULT(Result::INCOMPLETE, buffer.append(turnOn.data() + 100, 100, 100, turnOn.size())); + ASSERT_RESULT(Result::DUPLICATE, buffer.append(turnOn.data() + 200, turnOn.size() - 200, 200, turnOn.size())); + TEST_ASSERT_FALSE(buffer.ready()); + + // and the buffer is free for the next directive + std::string turnOff = directive("TurnOff", ID_2); + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, turnOff, 100)); + assertWaiting(buffer, turnOff); +} + +void test_directive_that_starts_like_the_one_just_read_is_not_lost() +{ + AlexaDirectiveBuffer buffer(2047); + std::string first = directive("TurnOn", ID_1); + std::string second = first; + second[second.size() - 10] = '#'; + + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, first, 100)); + ASSERT_RESULT(Result::INCOMPLETE, buffer.append(second.data(), 100, 0, second.size())); + buffer.release(); // the first has been read: nothing waits, so the second has a place to go + ASSERT_RESULT(Result::INCOMPLETE, buffer.append(second.data() + 100, 100, 100, second.size())); + ASSERT_RESULT(Result::COMPLETE, buffer.append(second.data() + 200, second.size() - 200, 200, second.size())); + assertWaiting(buffer, second); +} + +// --- fragments that do not fit --- + +void test_fragment_with_a_gap_drops_the_message() +{ + AlexaDirectiveBuffer buffer(2047); + std::string turnOn = directive("TurnOn", ID_1); + + ASSERT_RESULT(Result::INCOMPLETE, buffer.append(turnOn.data(), 100, 0, turnOn.size())); + ASSERT_RESULT(Result::OUT_OF_ORDER, buffer.append(turnOn.data() + 150, 50, 150, turnOn.size())); + ASSERT_RESULT(Result::IGNORED, buffer.append(turnOn.data() + 200, turnOn.size() - 200, 200, turnOn.size())); + TEST_ASSERT_FALSE(buffer.ready()); + + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, turnOn, 100)); + assertWaiting(buffer, turnOn); +} + +void test_fragment_beyond_the_total_is_refused() +{ + AlexaDirectiveBuffer buffer(2047); + std::string message(300, 'x'); + + // The block holds 200 + 1 bytes; a fragment that would end at 250 must not be copied + ASSERT_RESULT(Result::INCOMPLETE, buffer.append(message.data(), 100, 0, 200)); + ASSERT_RESULT(Result::OUT_OF_ORDER, buffer.append(message.data() + 100, 150, 100, 200)); + TEST_ASSERT_FALSE(buffer.ready()); + + // A first fragment longer than its own total + ASSERT_RESULT(Result::OUT_OF_ORDER, buffer.append(message.data(), 300, 0, 200)); + TEST_ASSERT_FALSE(buffer.ready()); +} + +void test_fragment_with_another_total_is_refused() +{ + AlexaDirectiveBuffer buffer(2047); + std::string message(300, 'x'); + + ASSERT_RESULT(Result::INCOMPLETE, buffer.append(message.data(), 100, 0, 200)); + ASSERT_RESULT(Result::OUT_OF_ORDER, buffer.append(message.data() + 100, 100, 100, 300)); + TEST_ASSERT_FALSE(buffer.ready()); +} + +void test_fragment_without_a_start_is_refused_once() +{ + AlexaDirectiveBuffer buffer(2047); + std::string turnOn = directive("TurnOn", ID_1); + + ASSERT_RESULT(Result::OUT_OF_ORDER, buffer.append(turnOn.data() + 100, 100, 100, turnOn.size())); + ASSERT_RESULT(Result::IGNORED, buffer.append(turnOn.data() + 200, turnOn.size() - 200, 200, turnOn.size())); + TEST_ASSERT_FALSE(buffer.ready()); +} + +void test_new_message_replaces_one_that_never_completed() +{ + AlexaDirectiveBuffer buffer(2047); + std::string turnOn = directive("TurnOn", ID_1); + std::string turnOff = directive("TurnOff", ID_2); + + ASSERT_RESULT(Result::INCOMPLETE, buffer.append(turnOn.data(), 100, 0, turnOn.size())); + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, turnOff, 100)); + assertWaiting(buffer, turnOff); +} + +void test_lost_connection_drops_the_arriving_message_only() +{ + AlexaDirectiveBuffer buffer(2047); + std::string turnOn = directive("TurnOn", ID_1); + std::string turnOff = directive("TurnOff", ID_2); + + // Nothing waits, a message is arriving + ASSERT_RESULT(Result::INCOMPLETE, buffer.append(turnOn.data(), 100, 0, turnOn.size())); + buffer.cancelArrival(); + ASSERT_RESULT(Result::OUT_OF_ORDER, buffer.append(turnOn.data() + 100, 100, 100, turnOn.size())); + TEST_ASSERT_FALSE(buffer.ready()); + + // A directive waits, its repeat is arriving + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, turnOff, 100)); + ASSERT_RESULT(Result::INCOMPLETE, buffer.append(turnOff.data(), 100, 0, turnOff.size())); + buffer.cancelArrival(); + assertWaiting(buffer, turnOff); +} + +// --- repeated messageIds --- + +void test_repeated_message_id_is_recognised() +{ + AlexaRecentIds recent; + + TEST_ASSERT_FALSE(recent.seenBefore(ID_1)); + TEST_ASSERT_TRUE(recent.seenBefore(ID_1)); + TEST_ASSERT_FALSE(recent.seenBefore(ID_2)); + TEST_ASSERT_TRUE(recent.seenBefore(ID_1)); + TEST_ASSERT_TRUE(recent.seenBefore(ID_2)); +} + +void test_only_the_last_four_ids_are_remembered() +{ + AlexaRecentIds recent; + const char *ids[] = {"id-1", "id-2", "id-3", "id-4", "id-5"}; + + for (const char *id : ids) + { + TEST_ASSERT_FALSE(recent.seenBefore(id)); + } + // id-5 took the place of id-1 + TEST_ASSERT_TRUE(recent.seenBefore("id-2")); + TEST_ASSERT_TRUE(recent.seenBefore("id-3")); + TEST_ASSERT_TRUE(recent.seenBefore("id-4")); + TEST_ASSERT_TRUE(recent.seenBefore("id-5")); + TEST_ASSERT_FALSE(recent.seenBefore("id-1")); +} + +void test_recognising_a_repeat_does_not_use_a_place() +{ + AlexaRecentIds recent; + + TEST_ASSERT_FALSE(recent.seenBefore("id-1")); + for (int i = 0; i < 10; i++) + { + TEST_ASSERT_TRUE(recent.seenBefore("id-1")); + } + TEST_ASSERT_FALSE(recent.seenBefore("id-2")); + TEST_ASSERT_FALSE(recent.seenBefore("id-3")); + TEST_ASSERT_FALSE(recent.seenBefore("id-4")); + TEST_ASSERT_TRUE(recent.seenBefore("id-1")); +} + +void test_directive_without_a_message_id_is_never_a_repeat() +{ + AlexaRecentIds recent; + + TEST_ASSERT_FALSE(recent.seenBefore("")); + TEST_ASSERT_FALSE(recent.seenBefore("")); + TEST_ASSERT_FALSE(recent.seenBefore(nullptr)); + TEST_ASSERT_FALSE(recent.seenBefore(nullptr)); +} + +// --- what may be published --- + +void test_message_within_the_limit_is_measured() +{ + JsonDocument doc; + doc["event"]["header"]["name"] = "Response"; + size_t length = 0; + + ASSERT_RESULT(AlexaSendResult::OK, AlexaBridgeLogic::checkMessage(doc, 2047, &length)); + TEST_ASSERT_EQUAL_UINT(strlen("{\"event\":{\"header\":{\"name\":\"Response\"}}}"), length); +} + +void test_message_at_the_limit_is_accepted_and_one_byte_more_is_refused() +{ + // {"v":"xxx...x"} is 8 bytes and the text + JsonDocument doc; + size_t length = 0; + + doc["v"] = std::string(2047 - 8, 'x'); + ASSERT_RESULT(AlexaSendResult::OK, AlexaBridgeLogic::checkMessage(doc, 2047, &length)); + TEST_ASSERT_EQUAL_UINT(2047, length); + + doc["v"] = std::string(2048 - 8, 'x'); + ASSERT_RESULT(AlexaSendResult::TOO_LARGE, AlexaBridgeLogic::checkMessage(doc, 2047, &length)); + TEST_ASSERT_EQUAL_UINT(2048, length); +} + +void test_empty_document_is_not_published() +{ + JsonDocument doc; + size_t length = 99; + + ASSERT_RESULT(AlexaSendResult::EMPTY, AlexaBridgeLogic::checkMessage(doc, 2047, &length)); + TEST_ASSERT_EQUAL_UINT(0, length); +} + +// An allocator that runs dry, as the heap of the board does +struct ScarceAllocator : ArduinoJson::Allocator +{ + size_t left; + explicit ScarceAllocator(size_t bytes) : left(bytes) {} + void *allocate(size_t size) override + { + if (size > left) + { + return nullptr; + } + left -= size; + return malloc(size); + } + void deallocate(void *pointer) override { free(pointer); } + void *reallocate(void *pointer, size_t size) override + { + if (size > left) + { + return nullptr; + } + left -= size; + return realloc(pointer, size); + } +}; + +void test_document_that_ran_out_of_memory_is_not_published() +{ + ScarceAllocator allocator(8192); + JsonDocument doc(&allocator); + for (int i = 0; i < 200 && !doc.overflowed(); i++) + { + doc["context"]["properties"][i]["value"] = std::string(100, 'x'); + } + TEST_ASSERT_TRUE(doc.overflowed()); + + size_t length = 99; + ASSERT_RESULT(AlexaSendResult::TOO_LARGE, AlexaBridgeLogic::checkMessage(doc, 1000000, &length)); + TEST_ASSERT_EQUAL_UINT(0, length); // 0 tells the two reasons for TOO_LARGE apart +} + +// --- topics --- + +void test_endpoint_is_taken_from_the_directive_topic() +{ + size_t length = 0; + const char *topic = "AEXAMPLEROOT/ESP-01/alexaDirective"; + + const char *endpointId = AlexaBridgeLogic::directiveEndpoint(topic, "AEXAMPLEROOT", &length); + TEST_ASSERT_NOT_NULL(endpointId); + TEST_ASSERT_EQUAL_UINT(6, length); + TEST_ASSERT_EQUAL_STRING_LEN("ESP-01", endpointId, length); + TEST_ASSERT_EQUAL_PTR(topic + strlen("AEXAMPLEROOT/"), endpointId); + + // An empty root topic (begin() with "") still has the separator + TEST_ASSERT_NOT_NULL(AlexaBridgeLogic::directiveEndpoint("/ESP-01/alexaDirective", "", &length)); + TEST_ASSERT_EQUAL_UINT(6, length); +} + +void test_other_topics_are_not_directive_topics() +{ + size_t length = 0; + + // the token topic of the 1.x HTTP fallback, published next to every directive + TEST_ASSERT_NULL(AlexaBridgeLogic::directiveEndpoint("root/ESP-01/alexaDirective_e", "root", &length)); + TEST_ASSERT_NULL(AlexaBridgeLogic::directiveEndpoint("root/discover", "root", &length)); + TEST_ASSERT_NULL(AlexaBridgeLogic::directiveEndpoint("root/ESP-01/alexaResponce", "root", &length)); + // another root, also one that only starts like ours + TEST_ASSERT_NULL(AlexaBridgeLogic::directiveEndpoint("other/ESP-01/alexaDirective", "root", &length)); + TEST_ASSERT_NULL(AlexaBridgeLogic::directiveEndpoint("rootless/ESP-01/alexaDirective", "root", &length)); + // no endpoint id, more than one level + TEST_ASSERT_NULL(AlexaBridgeLogic::directiveEndpoint("root//alexaDirective", "root", &length)); + TEST_ASSERT_NULL(AlexaBridgeLogic::directiveEndpoint("root/alexaDirective", "root", &length)); + TEST_ASSERT_NULL(AlexaBridgeLogic::directiveEndpoint("root/a/b/alexaDirective", "root", &length)); + TEST_ASSERT_NULL(AlexaBridgeLogic::directiveEndpoint("", "root", &length)); + TEST_ASSERT_NULL(AlexaBridgeLogic::directiveEndpoint(nullptr, "root", &length)); + TEST_ASSERT_NULL(AlexaBridgeLogic::directiveEndpoint("root/ESP-01/alexaDirective", nullptr, &length)); + TEST_ASSERT_NULL(AlexaBridgeLogic::directiveEndpoint("root/ESP-01/alexaDirective", "root", nullptr)); +} + +// --- time --- + +void test_timestamp_is_iso_8601_in_utc() +{ + char stamp[ALEXA_TIMESTAMP_SIZE]; + + TEST_ASSERT_TRUE(AlexaBridgeLogic::formatTimestamp(0, stamp, sizeof(stamp))); + TEST_ASSERT_EQUAL_STRING("1970-01-01T00:00:00Z", stamp); + + TEST_ASSERT_TRUE(AlexaBridgeLogic::formatTimestamp(1709210096, stamp, sizeof(stamp))); + TEST_ASSERT_EQUAL_STRING("2024-02-29T12:34:56Z", stamp); // a leap day + + TEST_ASSERT_TRUE(AlexaBridgeLogic::formatTimestamp(1790600709, stamp, sizeof(stamp))); + TEST_ASSERT_EQUAL_STRING("2026-09-28T13:05:09Z", stamp); + + TEST_ASSERT_TRUE(AlexaBridgeLogic::formatTimestamp(1798761599, stamp, sizeof(stamp))); + TEST_ASSERT_EQUAL_STRING("2026-12-31T23:59:59Z", stamp); +} + +void test_timestamp_needs_a_buffer_of_its_size() +{ + char small[ALEXA_TIMESTAMP_SIZE - 1]; + memset(small, 'x', sizeof(small)); + + TEST_ASSERT_FALSE(AlexaBridgeLogic::formatTimestamp(1790600709, small, sizeof(small))); + TEST_ASSERT_EQUAL_STRING("", small); + TEST_ASSERT_FALSE(AlexaBridgeLogic::formatTimestamp(1790600709, nullptr, ALEXA_TIMESTAMP_SIZE)); + + // The year 10000 does not fit the format + char stamp[ALEXA_TIMESTAMP_SIZE + 4]; + TEST_ASSERT_FALSE(AlexaBridgeLogic::formatTimestamp(253402300800LL, stamp, sizeof(stamp))); + TEST_ASSERT_EQUAL_STRING("", stamp); +} + +void test_clock_counts_as_set_from_2024() +{ + TEST_ASSERT_FALSE(AlexaBridgeLogic::clockIsSet(0)); + TEST_ASSERT_FALSE(AlexaBridgeLogic::clockIsSet(5)); // seconds since boot, before SNTP has answered + TEST_ASSERT_FALSE(AlexaBridgeLogic::clockIsSet(1704067199)); // 2023-12-31T23:59:59Z + TEST_ASSERT_TRUE(AlexaBridgeLogic::clockIsSet(1704067200)); // 2024-01-01T00:00:00Z + TEST_ASSERT_TRUE(AlexaBridgeLogic::clockIsSet(1790600709)); +} + +int main(int, char **) +{ + UNITY_BEGIN(); + + RUN_TEST(test_directive_in_one_piece_is_complete); + RUN_TEST(test_directive_in_fragments_is_reassembled_and_parses); + RUN_TEST(test_nothing_waits_before_the_last_fragment); + + RUN_TEST(test_directive_at_the_limit_is_accepted); + RUN_TEST(test_directive_over_the_limit_is_refused_once); + RUN_TEST(test_refused_directive_does_not_block_the_next_one); + RUN_TEST(test_empty_message_is_reported); + + RUN_TEST(test_second_directive_before_loop_is_dropped_and_the_first_kept); + RUN_TEST(test_second_directive_of_the_same_size_is_dropped); + + RUN_TEST(test_repeat_of_the_waiting_directive_is_a_duplicate); + RUN_TEST(test_repeat_that_differs_in_a_later_fragment_is_dropped); + RUN_TEST(test_repeat_still_arriving_when_the_first_is_read); + RUN_TEST(test_directive_that_starts_like_the_one_just_read_is_not_lost); + + RUN_TEST(test_fragment_with_a_gap_drops_the_message); + RUN_TEST(test_fragment_beyond_the_total_is_refused); + RUN_TEST(test_fragment_with_another_total_is_refused); + RUN_TEST(test_fragment_without_a_start_is_refused_once); + RUN_TEST(test_new_message_replaces_one_that_never_completed); + RUN_TEST(test_lost_connection_drops_the_arriving_message_only); + + RUN_TEST(test_repeated_message_id_is_recognised); + RUN_TEST(test_only_the_last_four_ids_are_remembered); + RUN_TEST(test_recognising_a_repeat_does_not_use_a_place); + RUN_TEST(test_directive_without_a_message_id_is_never_a_repeat); + + RUN_TEST(test_message_within_the_limit_is_measured); + RUN_TEST(test_message_at_the_limit_is_accepted_and_one_byte_more_is_refused); + RUN_TEST(test_empty_document_is_not_published); + RUN_TEST(test_document_that_ran_out_of_memory_is_not_published); + + RUN_TEST(test_endpoint_is_taken_from_the_directive_topic); + RUN_TEST(test_other_topics_are_not_directive_topics); + + RUN_TEST(test_timestamp_is_iso_8601_in_utc); + RUN_TEST(test_timestamp_needs_a_buffer_of_its_size); + RUN_TEST(test_clock_counts_as_set_from_2024); + + return UNITY_END(); +}