Drop the HTTP fallback: directives arrive on <root>/<id>/alexaDirective and reports leave over MQTT

The library fetched every directive over HTTP (GET /Alex2ESP/<token>, the token taken from <root>/<id>/alexaDirective_e) and posted every report back over HTTP. The detour was built in 2024 because AsyncMqttClient hands a message larger than one TCP segment to onMessage in fragments. Alex2MQTT publishes the whole directive on <root>/<id>/alexaDirective, so the fragments are reassembled instead: by index/total, into one heap block of total + 1 bytes that is freed as soon as loop() has parsed it.

Removed: ESP8266HTTPClient with processHttpGet/processHttpPost and their two HTTPClient and two WiFiClient members; the 5-slot packet, topic and receive queues, combinedData, receivePayload and AlexaStatusMessage::outputString (17,264 bytes of .bss); the {MESSAGE_PAYLOAD_SPLIT} and {REPLACE_WITH_DATETIME} conventions; the subscription to the _e token topic; AlexaUtils.cpp (AlexaUtils::printMemoryInfo stays, in the header).

Receive: the bridge subscribes <root>/discover and <root>/+/alexaDirective. The MQTT callback only collects; loop() parses and dispatches, and answers discovery. A directive over ALEX2ESP_MAX_DIRECTIVE (2047) is refused. One directive is in flight: a second one that arrives before loop() ran is dropped. Both are logged. A byte-identical repeat of the waiting directive and a repeat of one of the last four messageIds (kept as 64-bit hashes, 32 bytes) are ignored, because the broker mirror delivers every message twice. The wildcard subscription also delivers the directives of other boards of the account; they are recognised by their topic before anything is allocated.

Publish: measureJson first. A message is refused when it is over ALEX2ESP_MAX_MESSAGE (2047), when its document overflowed, when there is no session, when the packet does not fit the largest free block, or when AsyncMqttClient returns 0. Every refusal is logged and returned as an AlexaSendResult; AlexaStatusMessage::send() keeps its bool. The limit now also applies to the discovery object of a device.

Time: timeOfSample is an ISO 8601 UTC instant from the clock of the board (gmtime_r + snprintf, no strftime). begin() calls configTime(0, 0, "pool.ntp.org", "time.nist.gov") and loop() holds the first connect until the clock is set or 5 s have passed. setTimeSource(false) leaves the clock to the sketch. AddContextProp() fills timeOfSample in when a hand-built property has none or carries the old placeholder.

Log: AlexaLog, with a level at run time (setLogLevel) and a ceiling at build time (ALEX2ESP_LOG_MAX), replaces the Serial prints and the Alex2ESP_DEBUG define; the new receive and publish paths need a line for every refusal. The log and the time stamp call vsnprintf/snprintf with a PROGMEM format and not the _P variants: newlib keeps those in one object with printf_P and sprintf_P, which links the FILE-based printf. Measured on basicLight: 5,776 bytes of flash.

For a sketch: the 1.x API is unchanged, the five examples compile unmodified and without warnings. loop() no longer blocks for two HTTP round trips per directive. send() publishes at once and returns false without a session, where 1.1.0 queued the report. The MQTT session opens from loop(), up to 5 s after begin(). getState() is CONNECTED once both subscriptions are acknowledged. The 1.2.0 section of readme.md lists every change; its transport section is rewritten.

Measured for d1_mini with empty credentials (PlatformIO 6.2.0, espressif8266 4.2.1, Arduino core 3.1.2), static RAM / flash in bytes, 1.1.0 -> this commit:
  basicLight           52,768 / 350,885 -> 34,116 / 336,757
  lightWithBrightness  52,880 / 354,729 -> 34,232 / 340,517
  lightWithColorTemp   53,028 / 355,389 -> 34,380 / 341,177
  tempSensor           52,676 / 349,441 -> 34,024 / 335,457
  blindControl         52,900 / 353,069 -> 34,256 / 338,925
basicLight with -DALEX2ESP_LOG_MAX=0: 34,108 / 333,289. SNTP and the time stamp are 1,848 bytes of the flash figure.

Tests: platformio.ini with [env:native] and test/test_bridge_logic, 32 host tests of the logic in src/AlexaBridgeLogic.cpp (reassembly at every fragment size, both limits, repeats, publish checks, topics, time stamps). They also pass under -fsanitize=address,undefined. .gitignore no longer hides /test.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
David 2026-09-28 14:31:57 +00:00
parent 21cc693e74
commit 8ea777806a
20 changed files with 1840 additions and 792 deletions

2
.gitignore vendored
View file

@ -1,2 +1,2 @@
/.vscode /.vscode
/test /.pio

View file

@ -23,6 +23,10 @@ EndpointHealth KEYWORD1
PowerController KEYWORD1 PowerController KEYWORD1
TemperatureSensorScale KEYWORD1 TemperatureSensorScale KEYWORD1
AlexaUtils KEYWORD1 AlexaUtils KEYWORD1
AlexaLog KEYWORD1
AlexaLogLevel KEYWORD1
AlexaSendResult KEYWORD1
AlexaTransport KEYWORD1
####################################### #######################################
# Methods and Functions (KEYWORD2) # Methods and Functions (KEYWORD2)
@ -33,9 +37,14 @@ loop KEYWORD2
getState KEYWORD2 getState KEYWORD2
getDisconnectReason KEYWORD2 getDisconnectReason KEYWORD2
getDevice KEYWORD2 getDevice KEYWORD2
setLogLevel KEYWORD2
setTimeSource KEYWORD2
publish KEYWORD2
timestamp KEYWORD2
setName KEYWORD2 setName KEYWORD2
getName KEYWORD2 getName KEYWORD2
getEndpointId KEYWORD2 getEndpointId KEYWORD2
hasEndpointId KEYWORD2
setDisplayCategory KEYWORD2 setDisplayCategory KEYWORD2
getDisplayCategory KEYWORD2 getDisplayCategory KEYWORD2
setDescription KEYWORD2 setDescription KEYWORD2
@ -72,10 +81,14 @@ AddColorTemperatureControllerProp KEYWORD2
AddToggleControllerProp KEYWORD2 AddToggleControllerProp KEYWORD2
AddContextProp KEYWORD2 AddContextProp KEYWORD2
send KEYWORD2 send KEYWORD2
printMemoryInfo KEYWORD2
####################################### #######################################
# Constants (LITERAL1) # 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 MAX_EVENTS LITERAL1

17
platformio.ini Normal file
View file

@ -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 = -<*> +<AlexaBridgeLogic.cpp>
lib_deps =
bblanchon/ArduinoJson@^7
build_flags = -std=gnu++17 -Wall -Wextra

View file

@ -135,9 +135,16 @@ An `ActionMapping` takes an optional third argument, the directive payload as JS
## How it talks to Alex2MQTT ## How it talks to Alex2MQTT
- **Discovery.** On `<root>/discover` the library answers with one discovery object per device, published straight to `<root>/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 <endpointId>` 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. Everything goes over MQTT (port 1883 of `alex2mqtt.stormysdream.club`); the library makes no HTTP requests.
- **Directives.** A directive arrives as a short id on `<root>/<endpointId>/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 `<root>/<endpointId>/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. - **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 `<root>/discover` and `<root>/+/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 `<root>/discover` the library answers with one discovery object per device on `<root>/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 <endpointId>` 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 `<root>/<endpointId>/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 `<root>/<endpointId>/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=<bytes>` and `-DALEX2ESP_MAX_MESSAGE=<bytes>` 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 `<root>/<endpointId>/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["namespace"] = "Alexa.PowerController";
doc["name"] = "powerState"; doc["name"] = "powerState";
doc["value"] = "ON"; doc["value"] = "ON";
doc["timeOfSample"] = "{REPLACE_WITH_DATETIME}"; // Alex2MQTT server wil do the replace
doc["uncertaintyInMilliseconds"] = 0; doc["uncertaintyInMilliseconds"] = 0;
AddContextProp(doc.as<JsonObject>()) AddContextProp(doc.as<JsonObject>())
``` ```
`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 ## Changelog
### 1.2.0
Behaviour changes:
- Directives arrive over MQTT. The library subscribes to `<root>/+/alexaDirective`, where Alex2MQTT has always published the whole directive, instead of fetching it over HTTP with the token from `<root>/<endpointId>/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 `<root>/<endpointId>/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 ### 1.1.0
Behaviour changes: Behaviour changes:
@ -253,7 +283,7 @@ Packaging: `library.json` and `library.properties` restored with the Forgejo URL
--- ---
## Contributing ## 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.
--- ---

View file

@ -1,39 +1,77 @@
/* /*
* @title Alex2ESP Library * @title Alex2ESP Library
* @version 1.1.0
* @author David * @author David
* @license MIT * @license MIT
* @contributors chaos511 * @contributors chaos511
* *
* @description The Alex2ESP library is a companion to the Alex2MQTT Alexa Skill, * @description The bridge: MQTT session, discovery, directive receive and dispatch, publishing, time.
* providing seamless integration between ESP-based devices and the Alex2MQTT server. * The topics are listed in Alex2ESP.h.
* 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.
*/ */
#include "Alex2ESP.h" #include "Alex2ESP.h"
#include <time.h>
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() 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) 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->rootTopic = rootTopic;
this->mqttUsername = username; discoverTopic = this->rootTopic + "/discover";
this->mqttPassword = password; discoverTopicSend = this->rootTopic + "/discover_r";
directiveFilter = this->rootTopic + "/+/alexaDirective";
discoverTopic = (String(this->rootTopic) + String("/discover"));
discoverTopicSend = (String(this->rootTopic) + String("/discover_r"));
TopicESP = (String(this->rootTopic) + String("/+/alexaDirective_e"));
// Log the root topic (never the credentials) // Log the root topic (never the credentials)
AlexaUtils::log("Setting root topic: "); ALEX2ESP_LOGI("root topic %s", this->rootTopic.c_str());
AlexaUtils::logln(rootTopic);
AlexaUtils::printMemoryInfo(); // 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 // Set up MQTT client callbacks
mqttClient.onConnect([this](bool sessionPresent) 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); }); { this->onMessage(topic, payload, properties, len, index, total); });
// Configure MQTT client // Configure MQTT client
mqttClient.setServer(mqttServer, mqttPort); mqttClient.setServer(MQTT_SERVER, MQTT_PORT);
mqttClient.setCredentials(username, password); mqttClient.setCredentials(username, password);
// TODO: Throw error if dns fails or has no internet access? // loop() connects: see connectWhenClockIsSet()
_state = Alex2ESPState::CONNECTING; beginTime = millis();
mqttClient.connect(); state = Alex2ESPState::INITIALIZED;
} }
void Alex2ESP::onMqttConnect(bool sessionPresent) void Alex2ESP::setLogLevel(AlexaLogLevel level)
{ {
_state = Alex2ESPState::SUBSCRIBING; AlexaLog::setLevel(level);
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);
} }
void Alex2ESP::onSubscribe(uint16_t packetId, uint8_t qos) void Alex2ESP::setTimeSource(bool useSntp)
{ {
_state = Alex2ESPState::CONNECTED; this->useSntp = useSntp;
AlexaUtils::log("Subscribed with packetId: ");
AlexaUtils::log(packetId);
AlexaUtils::log(" and QoS: ");
AlexaUtils::logln(qos);
} }
// Internal: Handle disconnection Alex2ESPState Alex2ESP::getState() const
// TODO: auto reconnect?
void Alex2ESP::onMqttDisconnect(AsyncMqttClientDisconnectReason reason)
{ {
_state = Alex2ESPState::DISCONNECTED; return state;
_disconnectReason = reason;
} }
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"); return disconnectReason;
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;
} }
// Answer a Discover: one discovery object per device, published straight to <root>/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) AlexaDevice *Alex2ESP::getDevice(const String &name, const String &endpointId)
{ {
// Check if a device with the given endpointId already exists // 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 existing;
{ }
return &device; // Return the existing device
} 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 // If the device doesn't exist, create a new one
devices.emplace_back(name, rootTopic, endpointId); devices.emplace_back(name, rootTopic, endpointId, this);
ALEX2ESP_LOGI("device %s (%s) created", endpointId.c_str(), name.c_str());
Serial.print("Created new device: ");
Serial.print(name);
Serial.print(" with endpointId: ");
Serial.println(endpointId);
// Return a pointer to the newly created device // Return a pointer to the newly created device
return &devices.back(); return &devices.back();
} }
// void Alex2ESP::clearCache(boolean skipDocClear) AlexaDevice *Alex2ESP::findDevice(const char *endpointId, size_t length)
// {
// 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
{ {
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() 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(); lastReconnectTime = millis();
reconnectAttempt++;
mqttClient.connect(); 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"); ALEX2ESP_LOGD("subscribed to %s", topic.c_str());
httpPOST.setAuthorization(mqttUsername, mqttPassword); if (subscriptionsPending > 0 && --subscriptionsPending == 0)
httpPOST.addHeader("Content-Type", "text/plain"); {
state = Alex2ESPState::CONNECTED;
AlexaUtils::log("Sending Data: "); ALEX2ESP_LOGI("ready: %u device(s)", (unsigned)devices.size());
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();
} }
} }
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"); // The payload is unused: answer once per message even when TCP split it
if (index == 0)
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)
{ {
Serial.print("Error code: "); discoveryRequested = true;
Serial.println(httpResponseCode); discoveryStarted = millis();
if (retryCountGET >= MAX_RETRY_COUNT) }
{ return;
AlexaUtils::dequeueReceive(uuid, true); }
}
retryCountGET++; size_t endpointLength = 0;
const char *endpointId = AlexaBridgeLogic::directiveEndpoint(topic, rootTopic.c_str(), &endpointLength);
if (endpointId == nullptr)
{
return;
}
if (findDevice(endpointId, endpointLength) == nullptr)
{
// <root>/+/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<PGM_P>(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<JsonObjectConst>())
{
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 <root>/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 else
{ {
retryCountGET = 0; // This object can never be sent: say so and announce the others
AlexaUtils::dequeueReceive(uuid, true); logRefusal(discoverTopicSend.c_str(), result, length);
String payloadStr = httpGET.getString(); ALEX2ESP_LOGE("device %s is not announced", device.getEndpointId().c_str());
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"]));
}
}
}
}
}
} }
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<char *>(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() void Alex2ESP::loop()
{ {
connectWhenClockIsSet();
handleMqttReconnection(); handleMqttReconnection();
finishDiscovery(); continueDiscovery();
if(clearToSend){ processDirective();
processHttpPost();
processHttpGet();
}
return;
} }

View file

@ -1,15 +1,15 @@
/* /*
* @title Alex2ESP Library * @title Alex2ESP Library
* @version 1.1.0
* @author David * @author David
* @license MIT * @license MIT
* @contributors chaos511 * @contributors chaos511
* *
* @description The Alex2ESP library is a companion to the Alex2MQTT Alexa Skill, * @description Companion library of the Alex2MQTT Alexa skill: the devices a sketch declares become Alexa
* providing seamless integration between ESP-based devices and the Alex2MQTT server. * endpoints through the MQTT broker at alex2mqtt.stormysdream.club.
* This library connects to alex2mqtt.stormysdream.club, where the skill is hosted, *
* allowing your devices to communicate effortlessly with the Alexa Voice Service * <root>/discover in answered with one discovery object per device on <root>/discover_r
* using MQTT as the backbone. * <root>/<endpointId>/alexaDirective in the directive, handed to the device's ReportState / Event handler
* <root>/<endpointId>/alexaResponce out the report the handler built with buildStatusMessage()
*/ */
#ifndef ALEX2ESP_H #ifndef ALEX2ESP_H
@ -18,98 +18,113 @@
#include <Arduino.h> #include <Arduino.h>
#include <AsyncMqttClient.h> #include <AsyncMqttClient.h>
#include <ArduinoJson.h> #include <ArduinoJson.h>
#include "AlexaBridgeLogic.h"
#include "AlexaDevice.h" #include "AlexaDevice.h"
#include "AlexaInterface.h" #include "AlexaInterface.h"
#include "AlexaLog.h"
#include "AlexaTransport.h"
#include "AlexaUtils.h" #include "AlexaUtils.h"
#include <deque> #include <deque>
#ifdef ESP32 // Largest directive the bridge accepts, in bytes. Override with -DALEX2ESP_MAX_DIRECTIVE=<bytes> in build_flags.
#include <HTTPClient.h> // ESP32: untested #ifndef ALEX2ESP_MAX_DIRECTIVE
#else #define ALEX2ESP_MAX_DIRECTIVE 2047
#include <ESP8266HTTPClient.h>
#endif #endif
enum class Alex2ESPState enum class Alex2ESPState
{ {
UNINITIALIZED, UNINITIALIZED, // begin() has not been called
INITIALIZED, INITIALIZED, // begin() has been called; the first connect waits for the clock
CONNECTING, CONNECTING,
SUBSCRIBING, SUBSCRIBING, // session open, the subscriptions are not acknowledged yet
CONNECTED, CONNECTED, // subscribed: discovery requests and directives arrive
DISCONNECTED DISCONNECTED
}; };
class Alex2ESP class Alex2ESP : public AlexaTransport
{ {
public: public:
// Constructor // Constructor
Alex2ESP(); 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); void begin(const char *username, const char *password, const char *rootTopic);
Alex2ESPState getState() const; 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(); void loop();
AsyncMqttClientDisconnectReason getDisconnectReason() const; AsyncMqttClientDisconnectReason getDisconnectReason() const;
// Returns the device with this endpointId, creating it on first use. The pointer stays valid for the lifetime // 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. // 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); 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: 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 static const unsigned long DISCOVERY_WINDOW_MS = 5000; // How long the backend keeps collecting a discovery answer
int retryCountPOST=0; static const unsigned long DISCOVERY_RETRY_MS = 20; // Pause before a refused discovery publish is tried again
int retryCountGET=0;
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 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 discoverTopic; // The topic we listen on for discovery messages
String discoverTopicSend; // The topic we send discovery messages String discoverTopicSend; // The topic we send discovery messages
String directiveFilter; // The subscription that delivers the directives of every endpoint
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;
std::deque<AlexaDevice> devices; // Collection of devices (deque: pointers handed out by getDevice stay valid) std::deque<AlexaDevice> devices; // Collection of devices (deque: pointers handed out by getDevice stay valid)
Alex2ESPState _state; Alex2ESPState state;
AsyncMqttClientDisconnectReason _disconnectReason; 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 bool discoveryRequested; // Set by the MQTT callback, taken by loop()
const char *mqttServer = "alex2mqtt.stormysdream.club"; bool discoveryActive; // An answer is going out
uint16_t mqttPort = 1883; 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 onMqttConnect(bool sessionPresent);
void onMqttDisconnect(AsyncMqttClientDisconnectReason reason); void onMqttDisconnect(AsyncMqttClientDisconnectReason reason);
void onSubscribe(uint16_t packetId, uint8_t qos); 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); void onMessage(char *topic, char *payload, AsyncMqttClientMessageProperties properties, size_t length, size_t index, size_t total);
//loop processing function //loop processing function
void connectWhenClockIsSet();
void handleMqttReconnection(); void handleMqttReconnection();
void continueDiscovery();
void publishDiscovery(); void publishDiscovery();
void finishDiscovery(); void processDirective();
void processHttpPost();
void processHttpGet();
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 #endif // ALEX2ESP_H

300
src/AlexaBridgeLogic.cpp Normal file
View file

@ -0,0 +1,300 @@
#include "AlexaBridgeLogic.h"
#include "AlexaCompat.h"
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
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<char *>(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<uint8_t>(*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);
// <root> '/' <at least one character> <suffix>
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;
}

123
src/AlexaBridgeLogic.h Normal file
View file

@ -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 <stddef.h>
#include <stdint.h>
#include <time.h>
#include <ArduinoJson.h>
#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 "<root>/<endpointId>/alexaDirective": a pointer into the topic and the id's length.
// nullptr for any other topic, among them <root>/discover and the <root>/<endpointId>/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

14
src/AlexaCompat.h Normal file
View file

@ -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 <Arduino.h>
#else
#ifndef PSTR
#define PSTR(text) (text)
#endif
#endif
#endif // ALEXA_COMPAT_H

View file

@ -1,7 +1,8 @@
#include "AlexaDevice.h" #include "AlexaDevice.h"
#include "AlexaLog.h"
AlexaDevice::AlexaDevice(const String& name, const String& rootTopic, const String& endpointId) AlexaDevice::AlexaDevice(const String& name, const String& rootTopic, const String& endpointId, AlexaTransport* transport)
: name(name), endpointId(endpointId), rootTopic(rootTopic) { : name(name), endpointId(endpointId), rootTopic(rootTopic), transport(transport) {
// Initialize event arrays to nullptr // Initialize event arrays to nullptr
for (int i = 0; i < MAX_EVENTS; ++i) { for (int i = 0; i < MAX_EVENTS; ++i) {
eventNames[i] = nullptr; // Set all event names to nullptr eventNames[i] = nullptr; // Set all event names to nullptr
@ -22,6 +23,10 @@ String AlexaDevice::getEndpointId() const {
return endpointId; 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 { DisplayCategory AlexaDevice::getDisplayCategory() const {
return displayCategory; return displayCategory;
} }
@ -78,8 +83,7 @@ AlexaInterface* AlexaDevice::addCapability(AlexaInterfaceType type) {
// If the interface doesn't exist, create a new one // If the interface doesn't exist, create a new one
capabilities.emplace_back(type); capabilities.emplace_back(type);
Serial.print("Created new interface with type: "); ALEX2ESP_LOGI("%s: capability %s added", endpointId.c_str(), capabilities.back().getTypeString().c_str());
Serial.println(capabilities.back().getTypeString());
// Return a pointer to the newly created device // Return a pointer to the newly created device
return &capabilities.back(); return &capabilities.back();
@ -158,9 +162,8 @@ void AlexaDevice::triggerEvent(const char* eventName, const JsonDocument& direct
return; 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) { if (warnIfMissing) {
Serial.print("No event registered for: "); ALEX2ESP_LOGE("%s: no handler registered for %s, directive not handled", endpointId.c_str(), eventName);
Serial.println(eventName);
} }
} }

View file

@ -10,6 +10,7 @@
#include <string> #include <string>
#include <deque> #include <deque>
#include "AlexaStatusMessage.h" #include "AlexaStatusMessage.h"
#include "AlexaTransport.h"
#define MAX_EVENTS 10 #define MAX_EVENTS 10
@ -147,13 +148,18 @@ public:
class AlexaDevice { class AlexaDevice {
public: 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); void setName(const String& name);
String getName() const; String getName() const;
String getEndpointId() 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; DisplayCategory getDisplayCategory() const;
void setDisplayCategory(DisplayCategory category); void setDisplayCategory(DisplayCategory category);
@ -183,13 +189,14 @@ public:
void triggerEvent(const char* eventName, const JsonDocument& directive, const AlexaInterfaceType& type, bool warnIfMissing = true) const; void triggerEvent(const char* eventName, const JsonDocument& directive, const AlexaInterfaceType& type, bool warnIfMissing = true) const;
AlexaStatusMessage buildStatusMessage(const String& correlationToken,const bool isResponse=false) { AlexaStatusMessage buildStatusMessage(const String& correlationToken,const bool isResponse=false) {
return AlexaStatusMessage(correlationToken,rootTopic,endpointId,isResponse); return AlexaStatusMessage(correlationToken,rootTopic,endpointId,isResponse,transport);
} }
private: private:
String name; String name;
String endpointId; String endpointId;
String rootTopic; String rootTopic;
AlexaTransport* transport;
const char* eventNames[MAX_EVENTS]; const char* eventNames[MAX_EVENTS];
void (*eventCallbacks[MAX_EVENTS])(const JsonDocument&, const AlexaInterfaceType&); void (*eventCallbacks[MAX_EVENTS])(const JsonDocument&, const AlexaInterfaceType&);

View file

@ -6,6 +6,7 @@
#include <unordered_map> #include <unordered_map>
#include <Arduino.h> #include <Arduino.h>
#include <ArduinoJson.h> #include <ArduinoJson.h>
#include "AlexaLog.h"
// Enum for Alexa Interface Types // Enum for Alexa Interface Types
enum class AlexaInterfaceType enum class AlexaInterfaceType
@ -617,8 +618,7 @@ public:
if (deserializeJson(payloadDoc, directive.payload) == DeserializationError::Ok && payloadDoc.is<JsonObject>()) { if (deserializeJson(payloadDoc, directive.payload) == DeserializationError::Ok && payloadDoc.is<JsonObject>()) {
directiveObj["payload"] = payloadDoc.as<JsonObject>(); directiveObj["payload"] = payloadDoc.as<JsonObject>();
} else { } else {
Serial.print("[Alex2ESP] action mapping payload is not a JSON object, omitted: "); ALEX2ESP_LOGE("action mapping payload is not a JSON object, omitted: %s", directive.payload.c_str());
Serial.println(directive.payload);
} }
} }
return doc; return doc;

49
src/AlexaLog.cpp Normal file
View file

@ -0,0 +1,49 @@
#include "AlexaLog.h"
#include <stdarg.h>
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);
}

56
src/AlexaLog.h Normal file
View file

@ -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 <Arduino.h>
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=<n> 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

View file

@ -1,3 +1,90 @@
#include "AlexaStatusMessage.h" #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>();
JsonObject event_header = event["header"].to<JsonObject>();
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<JsonObject>();
event_endpoint["endpointId"] = endpointId;
event["payload"].to<JsonObject>(); // required by Alexa.Response / StateReport, empty when there is nothing to add
contextProperties = doc["context"]["properties"].to<JsonArray>();
}
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<unsigned int>(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;
}

View file

@ -4,10 +4,9 @@
#include <Arduino.h> #include <Arduino.h>
#include <ArduinoJson.h> #include <ArduinoJson.h>
#include "AlexaInterface.h" #include "AlexaInterface.h"
#include "AlexaTransport.h"
#include "AlexaUtils.h" #include "AlexaUtils.h"
#define MAX_STATUS_REPORT_SIZE 2048
enum class EndpointHealth enum class EndpointHealth
{ {
OK, OK,
@ -28,32 +27,9 @@ enum class TemperatureSensorScale
class AlexaStatusMessage class AlexaStatusMessage
{ {
public: public:
AlexaStatusMessage(const String &correlationToken, const String &rootTopic, const String &endpointId, const bool isResponse) // 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.
this->endpointId = endpointId; AlexaStatusMessage(const String &correlationToken, const String &rootTopic, const String &endpointId, const bool isResponse, AlexaTransport *transport = nullptr);
this->rootTopic = rootTopic;
JsonObject event = doc["event"].to<JsonObject>();
JsonObject event_header = event["header"].to<JsonObject>();
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<JsonObject>();
event_endpoint["endpointId"] = endpointId;
event["payload"].to<JsonObject>(); // required by Alexa.Response / StateReport, empty when there is nothing to add
contextProperties = doc["context"]["properties"].to<JsonArray>();
}
AlexaStatusMessage &AddHealthProp(EndpointHealth endpointHealth, unsigned int uncertaintyInMs = 0) AlexaStatusMessage &AddHealthProp(EndpointHealth endpointHealth, unsigned int uncertaintyInMs = 0)
{ {
@ -93,66 +69,29 @@ public:
return AddProperty(AlexaInterfaceType::TOGGLE_CONTROLLER, "toggleState", value, uncertaintyInMs,instanceName); return AddProperty(AlexaInterfaceType::TOGGLE_CONTROLLER, "toggleState", value, uncertaintyInMs,instanceName);
} }
AlexaStatusMessage &AddContextProp(const JsonObject &property) // 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}".
contextProperties.add(property); AlexaStatusMessage &AddContextProp(const JsonObject &property);
return *this; // Return a reference to the current object
}
// Queue the report for delivery. Returns false (and says so on Serial) when the report does not fit the // Publishes the report on <root>/<endpointId>/alexaResponce now. Returns false when nothing was sent: no
// MAX_STATUS_REPORT_SIZE buffer or the send queue is full: a report is never sent truncated. // session with the broker, the MQTT client or the heap cannot take the report, or it is over
bool send() // ALEX2ESP_MAX_MESSAGE bytes. The reason is on Serial; 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;
}
private: private:
String rootTopic; String rootTopic;
String endpointId; String endpointId;
AlexaTransport *transport;
JsonDocument doc; JsonDocument doc;
JsonArray contextProperties; JsonArray contextProperties;
static char outputString[MAX_STATUS_REPORT_SIZE];
// TODO: make real uuid4 gen function String generateMessageId();
String generateMessageId() void setTimeOfSample(JsonObject property);
{
char buffer[38];
const char charset[] = "0123456789abcdefABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz"; // Allowed characters
// Seed the random number generator (optional)
srand(static_cast<unsigned int>(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);
}
template <typename T> template <typename T>
AlexaStatusMessage &AddProperty(AlexaInterfaceType type, const String &propertyName, const T &value, unsigned int uncertaintyInMs = 0,const String &instanceName="") AlexaStatusMessage &AddProperty(AlexaInterfaceType type, const String &propertyName, const T &value, unsigned int uncertaintyInMs = 0,const String &instanceName="")
{ {
JsonDocument prop; JsonObject prop = contextProperties.add<JsonObject>();
prop["namespace"] = AlexaInterfaceUtils::toString(type); prop["namespace"] = AlexaInterfaceUtils::toString(type);
prop["name"] = propertyName; prop["name"] = propertyName;
@ -161,10 +100,8 @@ private:
} }
prop["value"] = value; prop["value"] = value;
prop["timeOfSample"] = "{REPLACE_WITH_DATETIME}"; setTimeOfSample(prop);
prop["uncertaintyInMilliseconds"] = uncertaintyInMs; prop["uncertaintyInMilliseconds"] = uncertaintyInMs;
contextProperties.add(prop.as<JsonObject>());
return *this; return *this;
} }
}; };

43
src/AlexaTransport.h Normal file
View file

@ -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 <stddef.h>
#include <stdint.h>
#include <ArduinoJson.h>
// Largest message the bridge publishes, in bytes of JSON: a report or the discovery object of one device.
// Override with -DALEX2ESP_MAX_MESSAGE=<bytes> 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

View file

@ -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
}

View file

@ -1,89 +1,26 @@
#ifndef ALEXA_UTILS_H #ifndef ALEXA_UTILS_H
#define ALEXA_UTILS_H #define ALEXA_UTILS_H
#include <vector>
#include <string>
#include <Arduino.h> #include <Arduino.h>
#include <AsyncMqttClient.h>
#include <cstring> // 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 class AlexaUtils
{ {
public: 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() static void printMemoryInfo()
{ {
size_t freeHeap = ESP.getFreeHeap();
#ifdef ESP8266 #ifdef ESP8266
size_t freeStack = ESP.getFreeContStack(); // Requires ESP8266 core 3.0.0+ size_t freeStack = ESP.getFreeContStack(); // Requires ESP8266 core 3.0.0+
#else #else
size_t freeStack = 0; // ESP32: untested, no getFreeContStack() size_t freeStack = 0; // ESP32: untested, no getFreeContStack()
#endif #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", Serial.printf("freeHeap: %u, freeStack: %u, freeROM: %u, usedROM: %u\n",
freeHeap, (unsigned)ESP.getFreeHeap(),
freeStack, (unsigned)freeStack,
freeSketchSpace, (unsigned)ESP.getFreeSketchSpace(),
sketchSize); (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 #endif // ALEXA_UTILS_H

View file

@ -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 <unity.h>
#include <ArduinoJson.h>
#include <stdlib.h>
#include <string.h>
#include <string>
#include "AlexaBridgeLogic.h"
typedef AlexaDirectiveBuffer::Result Result;
#define ASSERT_RESULT(expected, actual) TEST_ASSERT_EQUAL_INT(static_cast<int>(expected), static_cast<int>(actual))
// A directive as Alex2MQTT publishes it on <root>/<endpointId>/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();
}