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:
parent
21cc693e74
commit
8ea777806a
20 changed files with 1840 additions and 792 deletions
2
.gitignore
vendored
2
.gitignore
vendored
|
|
@ -1,2 +1,2 @@
|
|||
/.vscode
|
||||
/test
|
||||
/.pio
|
||||
|
|
|
|||
15
keywords.txt
15
keywords.txt
|
|
@ -23,6 +23,10 @@ EndpointHealth KEYWORD1
|
|||
PowerController KEYWORD1
|
||||
TemperatureSensorScale KEYWORD1
|
||||
AlexaUtils KEYWORD1
|
||||
AlexaLog KEYWORD1
|
||||
AlexaLogLevel KEYWORD1
|
||||
AlexaSendResult KEYWORD1
|
||||
AlexaTransport KEYWORD1
|
||||
|
||||
#######################################
|
||||
# Methods and Functions (KEYWORD2)
|
||||
|
|
@ -33,9 +37,14 @@ loop KEYWORD2
|
|||
getState KEYWORD2
|
||||
getDisconnectReason KEYWORD2
|
||||
getDevice KEYWORD2
|
||||
setLogLevel KEYWORD2
|
||||
setTimeSource KEYWORD2
|
||||
publish KEYWORD2
|
||||
timestamp KEYWORD2
|
||||
setName KEYWORD2
|
||||
getName KEYWORD2
|
||||
getEndpointId KEYWORD2
|
||||
hasEndpointId KEYWORD2
|
||||
setDisplayCategory KEYWORD2
|
||||
getDisplayCategory KEYWORD2
|
||||
setDescription KEYWORD2
|
||||
|
|
@ -72,10 +81,14 @@ AddColorTemperatureControllerProp KEYWORD2
|
|||
AddToggleControllerProp KEYWORD2
|
||||
AddContextProp KEYWORD2
|
||||
send KEYWORD2
|
||||
printMemoryInfo KEYWORD2
|
||||
|
||||
#######################################
|
||||
# Constants (LITERAL1)
|
||||
#######################################
|
||||
|
||||
MAX_STATUS_REPORT_SIZE LITERAL1
|
||||
ALEX2ESP_MAX_DIRECTIVE LITERAL1
|
||||
ALEX2ESP_MAX_MESSAGE LITERAL1
|
||||
ALEX2ESP_LOG_MAX LITERAL1
|
||||
ALEXA_TIMESTAMP_SIZE LITERAL1
|
||||
MAX_EVENTS LITERAL1
|
||||
|
|
|
|||
17
platformio.ini
Normal file
17
platformio.ini
Normal 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
|
||||
40
readme.md
40
readme.md
|
|
@ -135,9 +135,16 @@ An `ActionMapping` takes an optional third argument, the directive payload as JS
|
|||
|
||||
## How it talks to Alex2MQTT
|
||||
|
||||
- **Discovery.** On `<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.
|
||||
- **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.
|
||||
Everything goes over MQTT (port 1883 of `alex2mqtt.stormysdream.club`); the library makes no HTTP requests.
|
||||
|
||||
- **Session.** `begin()` starts SNTP (`pool.ntp.org`, `time.nist.gov`) and returns; `loop()` opens the MQTT session once the clock is set, or after 5 s without an answer, and subscribes to `<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["name"] = "powerState";
|
||||
doc["value"] = "ON";
|
||||
doc["timeOfSample"] = "{REPLACE_WITH_DATETIME}"; // Alex2MQTT server wil do the replace
|
||||
doc["uncertaintyInMilliseconds"] = 0;
|
||||
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
|
||||
|
||||
### 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
|
||||
|
||||
Behaviour changes:
|
||||
|
|
@ -253,7 +283,7 @@ Packaging: `library.json` and `library.properties` restored with the Forgejo URL
|
|||
---
|
||||
|
||||
## Contributing
|
||||
Feel free to submit pull requests or issues for feature requests and bug fixes.
|
||||
Feel free to submit pull requests or issues for feature requests and bug fixes. `pio test -e native` in the repository root runs the host tests (no board needed); they have to pass.
|
||||
|
||||
---
|
||||
|
||||
|
|
|
|||
727
src/Alex2ESP.cpp
727
src/Alex2ESP.cpp
|
|
@ -1,39 +1,77 @@
|
|||
/*
|
||||
* @title Alex2ESP Library
|
||||
* @version 1.1.0
|
||||
* @author David
|
||||
* @license MIT
|
||||
* @contributors chaos511
|
||||
*
|
||||
* @description The Alex2ESP library is a companion to the Alex2MQTT Alexa Skill,
|
||||
* providing seamless integration between ESP-based devices and the Alex2MQTT server.
|
||||
* This library connects to alex2mqtt.stormysdream.club, where the skill is hosted,
|
||||
* allowing your devices to communicate effortlessly with the Alexa Voice Service
|
||||
* using MQTT as the backbone.
|
||||
* @description The bridge: MQTT session, discovery, directive receive and dispatch, publishing, time.
|
||||
* The topics are listed in Alex2ESP.h.
|
||||
*/
|
||||
#include "Alex2ESP.h"
|
||||
#include <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()
|
||||
: rootTopic(), mqttUsername(nullptr), mqttPassword(nullptr), lastReconnectTime(0), reconnectAttempt(0), _state(Alex2ESPState::UNINITIALIZED), _disconnectReason(AsyncMqttClientDisconnectReason::TCP_DISCONNECTED) {}
|
||||
: rootTopic(),
|
||||
state(Alex2ESPState::UNINITIALIZED),
|
||||
disconnectReason(AsyncMqttClientDisconnectReason::TCP_DISCONNECTED),
|
||||
useSntp(true),
|
||||
beginTime(0),
|
||||
lastReconnectTime(0),
|
||||
discoverSubscription(0),
|
||||
directiveSubscription(0),
|
||||
subscriptionsPending(0),
|
||||
discoveryRequested(false),
|
||||
discoveryActive(false),
|
||||
discoveryDeferred(false),
|
||||
discoveryNext(0),
|
||||
discoveryAnnounced(0),
|
||||
discoveryStarted(0),
|
||||
discoveryLastAttempt(0),
|
||||
directive(ALEX2ESP_MAX_DIRECTIVE) {}
|
||||
|
||||
void Alex2ESP::begin(const char *username, const char *password, const char *rootTopic)
|
||||
{
|
||||
if (state != Alex2ESPState::UNINITIALIZED)
|
||||
{
|
||||
ALEX2ESP_LOGE("begin() called again, ignored");
|
||||
return;
|
||||
}
|
||||
|
||||
_state = Alex2ESPState::INITIALIZED;
|
||||
_disconnectReason = AsyncMqttClientDisconnectReason::TCP_DISCONNECTED;
|
||||
this->rootTopic = rootTopic;
|
||||
this->mqttUsername = username;
|
||||
this->mqttPassword = password;
|
||||
|
||||
discoverTopic = (String(this->rootTopic) + String("/discover"));
|
||||
discoverTopicSend = (String(this->rootTopic) + String("/discover_r"));
|
||||
|
||||
TopicESP = (String(this->rootTopic) + String("/+/alexaDirective_e"));
|
||||
discoverTopic = this->rootTopic + "/discover";
|
||||
discoverTopicSend = this->rootTopic + "/discover_r";
|
||||
directiveFilter = this->rootTopic + "/+/alexaDirective";
|
||||
|
||||
// Log the root topic (never the credentials)
|
||||
AlexaUtils::log("Setting root topic: ");
|
||||
AlexaUtils::logln(rootTopic);
|
||||
AlexaUtils::printMemoryInfo();
|
||||
ALEX2ESP_LOGI("root topic %s", this->rootTopic.c_str());
|
||||
|
||||
// Reports carry the time of their samples, so the board needs the time of day: UTC from SNTP
|
||||
if (useSntp)
|
||||
{
|
||||
configTime(0, 0, NTP_SERVER_1, NTP_SERVER_2);
|
||||
}
|
||||
|
||||
// Set up MQTT client callbacks
|
||||
mqttClient.onConnect([this](bool sessionPresent)
|
||||
|
|
@ -46,396 +84,427 @@ void Alex2ESP::begin(const char *username, const char *password, const char *roo
|
|||
{ this->onMessage(topic, payload, properties, len, index, total); });
|
||||
|
||||
// Configure MQTT client
|
||||
mqttClient.setServer(mqttServer, mqttPort);
|
||||
mqttClient.setServer(MQTT_SERVER, MQTT_PORT);
|
||||
mqttClient.setCredentials(username, password);
|
||||
|
||||
// TODO: Throw error if dns fails or has no internet access?
|
||||
_state = Alex2ESPState::CONNECTING;
|
||||
mqttClient.connect();
|
||||
// loop() connects: see connectWhenClockIsSet()
|
||||
beginTime = millis();
|
||||
state = Alex2ESPState::INITIALIZED;
|
||||
}
|
||||
|
||||
void Alex2ESP::onMqttConnect(bool sessionPresent)
|
||||
void Alex2ESP::setLogLevel(AlexaLogLevel level)
|
||||
{
|
||||
_state = Alex2ESPState::SUBSCRIBING;
|
||||
|
||||
AlexaUtils::log("Connected to MQTT server, Now subscribing to: ");
|
||||
AlexaUtils::log(discoverTopic.c_str());
|
||||
AlexaUtils::log(" ");
|
||||
AlexaUtils::logln(TopicESP.c_str());
|
||||
|
||||
mqttClient.subscribe(discoverTopic.c_str(), 1);
|
||||
mqttClient.subscribe(TopicESP.c_str(), 1);
|
||||
AlexaLog::setLevel(level);
|
||||
}
|
||||
|
||||
void Alex2ESP::onSubscribe(uint16_t packetId, uint8_t qos)
|
||||
void Alex2ESP::setTimeSource(bool useSntp)
|
||||
{
|
||||
_state = Alex2ESPState::CONNECTED;
|
||||
|
||||
AlexaUtils::log("Subscribed with packetId: ");
|
||||
AlexaUtils::log(packetId);
|
||||
AlexaUtils::log(" and QoS: ");
|
||||
AlexaUtils::logln(qos);
|
||||
this->useSntp = useSntp;
|
||||
}
|
||||
|
||||
// Internal: Handle disconnection
|
||||
// TODO: auto reconnect?
|
||||
void Alex2ESP::onMqttDisconnect(AsyncMqttClientDisconnectReason reason)
|
||||
Alex2ESPState Alex2ESP::getState() const
|
||||
{
|
||||
_state = Alex2ESPState::DISCONNECTED;
|
||||
_disconnectReason = reason;
|
||||
return state;
|
||||
}
|
||||
|
||||
void Alex2ESP::onMessage(char *topic, char *payload, AsyncMqttClientMessageProperties properties, size_t length, size_t index, size_t total)
|
||||
AsyncMqttClientDisconnectReason Alex2ESP::getDisconnectReason() const
|
||||
{
|
||||
AlexaUtils::logln("OnMessage Start");
|
||||
AlexaUtils::printMemoryInfo();
|
||||
|
||||
clearToSend = false;
|
||||
|
||||
if (strcmp(topic, discoverTopic.c_str()) == 0)
|
||||
{
|
||||
// The payload is unused: answer once per message even when TCP split it, starting over if an answer is still going out
|
||||
if (index == 0)
|
||||
{
|
||||
discoveryNext = 0;
|
||||
discoveryPending = false;
|
||||
discoveryStarted = millis();
|
||||
publishDiscovery();
|
||||
return disconnectReason;
|
||||
}
|
||||
}
|
||||
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)
|
||||
{
|
||||
// Check if a device with the given endpointId already exists
|
||||
for (auto &device : devices)
|
||||
AlexaDevice *existing = findDevice(endpointId.c_str(), endpointId.length());
|
||||
if (existing != nullptr)
|
||||
{
|
||||
if (device.getEndpointId() == endpointId)
|
||||
{
|
||||
return &device; // Return the existing device
|
||||
return existing;
|
||||
}
|
||||
|
||||
if (state == Alex2ESPState::UNINITIALIZED)
|
||||
{
|
||||
ALEX2ESP_LOGE("device %s created before begin(): it has no root topic, its reports will not arrive", endpointId.c_str());
|
||||
}
|
||||
|
||||
// If the device doesn't exist, create a new one
|
||||
devices.emplace_back(name, rootTopic, endpointId);
|
||||
|
||||
Serial.print("Created new device: ");
|
||||
Serial.print(name);
|
||||
Serial.print(" with endpointId: ");
|
||||
Serial.println(endpointId);
|
||||
devices.emplace_back(name, rootTopic, endpointId, this);
|
||||
ALEX2ESP_LOGI("device %s (%s) created", endpointId.c_str(), name.c_str());
|
||||
|
||||
// Return a pointer to the newly created device
|
||||
return &devices.back();
|
||||
}
|
||||
|
||||
// void Alex2ESP::clearCache(boolean skipDocClear)
|
||||
// {
|
||||
// if (!skipDocClear)
|
||||
// {
|
||||
// inputDoc.clear();
|
||||
// }
|
||||
// cache.receivedParts = 0; // Reset the count of received parts
|
||||
// cache.messageId = ""; // Clear the message ID
|
||||
// cache.totalParts = 0; // Reset the total parts count
|
||||
// cache.fragmentOffset = 0;
|
||||
// memset(cache.message, 0, sizeof(cache.message)); // Clear the message buffer
|
||||
// }
|
||||
|
||||
// bool Alex2ESP::isMessageComplete()
|
||||
// {
|
||||
// return cache.receivedParts == cache.totalParts; // Check if we have received all parts
|
||||
// }
|
||||
|
||||
Alex2ESPState Alex2ESP::getState() const
|
||||
AlexaDevice *Alex2ESP::findDevice(const char *endpointId, size_t length)
|
||||
{
|
||||
return _state;
|
||||
for (auto &device : devices)
|
||||
{
|
||||
if (device.hasEndpointId(endpointId, length))
|
||||
{
|
||||
return &device;
|
||||
}
|
||||
}
|
||||
return nullptr;
|
||||
}
|
||||
|
||||
AsyncMqttClientDisconnectReason Alex2ESP::getDisconnectReason() const
|
||||
// The first connect waits until the clock is set, for CLOCK_WAIT_MS at most: a report sent before SNTP has answered
|
||||
// would carry a time of sample in 1970.
|
||||
void Alex2ESP::connectWhenClockIsSet()
|
||||
{
|
||||
return _disconnectReason;
|
||||
if (state != Alex2ESPState::INITIALIZED)
|
||||
{
|
||||
return;
|
||||
}
|
||||
if (!AlexaBridgeLogic::clockIsSet(time(nullptr)))
|
||||
{
|
||||
if (millis() - beginTime < CLOCK_WAIT_MS)
|
||||
{
|
||||
return;
|
||||
}
|
||||
ALEX2ESP_LOGE("clock not set after %lu ms: connecting, reports carry a wrong time until it is set", CLOCK_WAIT_MS);
|
||||
}
|
||||
|
||||
state = Alex2ESPState::CONNECTING;
|
||||
lastReconnectTime = millis();
|
||||
mqttClient.connect();
|
||||
}
|
||||
|
||||
void Alex2ESP::handleMqttReconnection()
|
||||
{
|
||||
if (!mqttClient.connected() && _disconnectReason == AsyncMqttClientDisconnectReason::TCP_DISCONNECTED)
|
||||
if (state == Alex2ESPState::UNINITIALIZED || state == Alex2ESPState::INITIALIZED)
|
||||
{
|
||||
if (millis() - lastReconnectTime > 5000)
|
||||
return;
|
||||
}
|
||||
if (!mqttClient.connected() && disconnectReason == AsyncMqttClientDisconnectReason::TCP_DISCONNECTED)
|
||||
{
|
||||
if (millis() - lastReconnectTime > RECONNECT_INTERVAL_MS)
|
||||
{
|
||||
lastReconnectTime = millis();
|
||||
reconnectAttempt++;
|
||||
mqttClient.connect();
|
||||
}
|
||||
}
|
||||
}
|
||||
void Alex2ESP::processHttpPost()
|
||||
|
||||
void Alex2ESP::onMqttConnect(bool sessionPresent)
|
||||
{
|
||||
if (!AlexaUtils::isQueueEmpty())
|
||||
state = Alex2ESPState::SUBSCRIBING;
|
||||
|
||||
subscriptionsPending = 2;
|
||||
discoverSubscription = mqttClient.subscribe(discoverTopic.c_str(), 1);
|
||||
directiveSubscription = mqttClient.subscribe(directiveFilter.c_str(), 1);
|
||||
if (discoverSubscription == 0 || directiveSubscription == 0)
|
||||
{
|
||||
AlexaUtils::logln("processHttpPost");
|
||||
|
||||
const char *payload = AlexaUtils::dequeueVals(false);
|
||||
|
||||
httpPOST.begin(wifiPOST, "http://alex2mqtt.stormysdream.club/Alex2ESP");
|
||||
httpPOST.setAuthorization(mqttUsername, mqttPassword);
|
||||
httpPOST.addHeader("Content-Type", "text/plain");
|
||||
|
||||
AlexaUtils::log("Sending Data: ");
|
||||
AlexaUtils::logln(payload);
|
||||
|
||||
int httpResponseCode = httpPOST.POST(payload);
|
||||
|
||||
if (httpResponseCode != 200)
|
||||
{
|
||||
|
||||
if (retryCountPOST >= MAX_RETRY_COUNT)
|
||||
{
|
||||
AlexaUtils::dequeueVals(true);
|
||||
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);
|
||||
}
|
||||
|
||||
AlexaUtils::log("POST Error code: ");
|
||||
AlexaUtils::logln(httpResponseCode);
|
||||
retryCountPOST++;
|
||||
if (httpResponseCode > 0)
|
||||
void Alex2ESP::onSubscribe(uint16_t packetId, uint8_t qos)
|
||||
{
|
||||
String response = httpPOST.getString();
|
||||
Serial.println("Response:");
|
||||
Serial.println(response);
|
||||
if (packetId != discoverSubscription && packetId != directiveSubscription)
|
||||
{
|
||||
return;
|
||||
}
|
||||
}
|
||||
else
|
||||
const String &topic = (packetId == discoverSubscription) ? discoverTopic : directiveFilter;
|
||||
if (qos == SUBSCRIPTION_REFUSED)
|
||||
{
|
||||
retryCountPOST = 0;
|
||||
AlexaUtils::dequeueVals(true);
|
||||
}
|
||||
httpPOST.end();
|
||||
}
|
||||
}
|
||||
void Alex2ESP::processHttpGet()
|
||||
{
|
||||
if (!AlexaUtils::isReceiveQueueEmpty())
|
||||
{
|
||||
AlexaUtils::logln("processHttpGet");
|
||||
|
||||
String uuid;
|
||||
AlexaUtils::dequeueReceive(uuid, false);
|
||||
AlexaUtils::log("Sending get request for uuid: ");
|
||||
// AlexaUtils::logln(uuid);
|
||||
|
||||
httpGET.begin(wifiGET, "http://alex2mqtt.stormysdream.club/Alex2ESP/" + uuid);
|
||||
httpGET.setAuthorization(mqttUsername, mqttPassword);
|
||||
int httpResponseCode = httpGET.GET();
|
||||
|
||||
Serial.print("HTTP Response code: ");
|
||||
Serial.println(httpResponseCode);
|
||||
if (httpResponseCode != 200)
|
||||
{
|
||||
Serial.print("Error code: ");
|
||||
Serial.println(httpResponseCode);
|
||||
if (retryCountGET >= MAX_RETRY_COUNT)
|
||||
{
|
||||
AlexaUtils::dequeueReceive(uuid, true);
|
||||
ALEX2ESP_LOGE("the broker refused the subscription to %s", topic.c_str());
|
||||
return;
|
||||
}
|
||||
|
||||
retryCountGET++;
|
||||
ALEX2ESP_LOGD("subscribed to %s", topic.c_str());
|
||||
if (subscriptionsPending > 0 && --subscriptionsPending == 0)
|
||||
{
|
||||
state = Alex2ESPState::CONNECTED;
|
||||
ALEX2ESP_LOGI("ready: %u device(s)", (unsigned)devices.size());
|
||||
}
|
||||
}
|
||||
else
|
||||
{
|
||||
retryCountGET = 0;
|
||||
AlexaUtils::dequeueReceive(uuid, true);
|
||||
String payloadStr = httpGET.getString();
|
||||
size_t length = payloadStr.length();
|
||||
|
||||
if (length < (size_t)(AlexaUtils::MAX_PAYLOAD_LENGTH - 1))
|
||||
void Alex2ESP::onMqttDisconnect(AsyncMqttClientDisconnectReason reason)
|
||||
{
|
||||
payloadStr.toCharArray(AlexaUtils::receivePayload, length + 1);
|
||||
payloadStr = "";
|
||||
AlexaUtils::receivePayload[length + 1] = '\0';
|
||||
// AlexaUtils::logln(AlexaUtils::receivePayload);
|
||||
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)
|
||||
{
|
||||
// The payload is unused: answer once per message even when TCP split it
|
||||
if (index == 0)
|
||||
{
|
||||
discoveryRequested = true;
|
||||
discoveryStarted = millis();
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
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();
|
||||
|
||||
DeserializationError error = deserializeJson(inputDoc, AlexaUtils::receivePayload);
|
||||
if (error)
|
||||
{
|
||||
AlexaUtils::log("Failed to parse message: ");
|
||||
AlexaUtils::logln(String(error.f_str()));
|
||||
inputDoc.clear();
|
||||
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
|
||||
{
|
||||
AlexaUtils::logln("JSON message parsed successfully!");
|
||||
device->triggerEvent("Event", message, AlexaInterfaceUtils::fromString(interfaceName));
|
||||
}
|
||||
}
|
||||
|
||||
for (auto &device : devices)
|
||||
// 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 (strcmp(device.getEndpointId().c_str(), inputDoc["directive"]["endpoint"]["endpointId"] | "") == 0)
|
||||
if (discoveryRequested)
|
||||
{
|
||||
serializeJson(inputDoc, Serial);
|
||||
Serial.println();
|
||||
device.triggerEvent("DirectiveReceived", inputDoc["directive"], AlexaInterfaceType::UNKNOWN, false);
|
||||
discoveryRequested = false;
|
||||
discoveryActive = true;
|
||||
discoveryDeferred = false;
|
||||
discoveryNext = 0;
|
||||
discoveryAnnounced = 0;
|
||||
}
|
||||
if (!discoveryActive)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
if (inputDoc["directive"]["header"]["name"] == "ReportState")
|
||||
if (millis() - discoveryStarted > DISCOVERY_WINDOW_MS)
|
||||
{
|
||||
device.triggerEvent("ReportState", inputDoc["directive"], AlexaInterfaceType::UNKNOWN);
|
||||
// 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
|
||||
{
|
||||
device.triggerEvent("Event", inputDoc["directive"], AlexaInterfaceUtils::fromString(inputDoc["directive"]["header"]["namespace"]));
|
||||
// This object can never be sent: say so and announce the others
|
||||
logRefusal(discoverTopicSend.c_str(), result, length);
|
||||
ALEX2ESP_LOGE("device %s is not announced", device.getEndpointId().c_str());
|
||||
}
|
||||
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;
|
||||
}
|
||||
httpGET.end();
|
||||
|
||||
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()
|
||||
{
|
||||
|
||||
connectWhenClockIsSet();
|
||||
handleMqttReconnection();
|
||||
finishDiscovery();
|
||||
if(clearToSend){
|
||||
processHttpPost();
|
||||
processHttpGet();
|
||||
}
|
||||
|
||||
return;
|
||||
continueDiscovery();
|
||||
processDirective();
|
||||
}
|
||||
|
|
|
|||
113
src/Alex2ESP.h
113
src/Alex2ESP.h
|
|
@ -1,15 +1,15 @@
|
|||
/*
|
||||
* @title Alex2ESP Library
|
||||
* @version 1.1.0
|
||||
* @author David
|
||||
* @license MIT
|
||||
* @contributors chaos511
|
||||
*
|
||||
* @description The Alex2ESP library is a companion to the Alex2MQTT Alexa Skill,
|
||||
* providing seamless integration between ESP-based devices and the Alex2MQTT server.
|
||||
* This library connects to alex2mqtt.stormysdream.club, where the skill is hosted,
|
||||
* allowing your devices to communicate effortlessly with the Alexa Voice Service
|
||||
* using MQTT as the backbone.
|
||||
* @description Companion library of the Alex2MQTT Alexa skill: the devices a sketch declares become Alexa
|
||||
* endpoints through the MQTT broker at alex2mqtt.stormysdream.club.
|
||||
*
|
||||
* <root>/discover in answered with one discovery object per device on <root>/discover_r
|
||||
* <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
|
||||
|
|
@ -18,98 +18,113 @@
|
|||
#include <Arduino.h>
|
||||
#include <AsyncMqttClient.h>
|
||||
#include <ArduinoJson.h>
|
||||
#include "AlexaBridgeLogic.h"
|
||||
#include "AlexaDevice.h"
|
||||
#include "AlexaInterface.h"
|
||||
#include "AlexaLog.h"
|
||||
#include "AlexaTransport.h"
|
||||
#include "AlexaUtils.h"
|
||||
#include <deque>
|
||||
|
||||
#ifdef ESP32
|
||||
#include <HTTPClient.h> // ESP32: untested
|
||||
#else
|
||||
#include <ESP8266HTTPClient.h>
|
||||
// Largest directive the bridge accepts, in bytes. Override with -DALEX2ESP_MAX_DIRECTIVE=<bytes> in build_flags.
|
||||
#ifndef ALEX2ESP_MAX_DIRECTIVE
|
||||
#define ALEX2ESP_MAX_DIRECTIVE 2047
|
||||
#endif
|
||||
|
||||
|
||||
enum class Alex2ESPState
|
||||
{
|
||||
UNINITIALIZED,
|
||||
INITIALIZED,
|
||||
UNINITIALIZED, // begin() has not been called
|
||||
INITIALIZED, // begin() has been called; the first connect waits for the clock
|
||||
CONNECTING,
|
||||
SUBSCRIBING,
|
||||
CONNECTED,
|
||||
SUBSCRIBING, // session open, the subscriptions are not acknowledged yet
|
||||
CONNECTED, // subscribed: discovery requests and directives arrive
|
||||
DISCONNECTED
|
||||
};
|
||||
|
||||
class Alex2ESP
|
||||
class Alex2ESP : public AlexaTransport
|
||||
{
|
||||
public:
|
||||
// Constructor
|
||||
Alex2ESP();
|
||||
|
||||
// Begin function for initialization: MQTT username, MQTT password, root topic (the same order as alex2node)
|
||||
// Begin function for initialization: MQTT username, MQTT password, root topic (the same order as alex2node).
|
||||
// The username and the password are not copied: they have to stay valid for as long as the client is used.
|
||||
// Starts SNTP; loop() opens the MQTT session once the clock is set, or after 5 s without an answer.
|
||||
void begin(const char *username, const char *password, const char *rootTopic);
|
||||
Alex2ESPState getState() const;
|
||||
|
||||
// Call from the sketch's loop(): connects, answers discovery requests and hands directives to the devices.
|
||||
// It does not block; the handlers of the sketch run inside it.
|
||||
void loop();
|
||||
|
||||
AsyncMqttClientDisconnectReason getDisconnectReason() const;
|
||||
|
||||
// Returns the device with this endpointId, creating it on first use. The pointer stays valid for the lifetime
|
||||
// of the client: devices live in a std::deque, which never relocates its elements when another one is added.
|
||||
// Call it after begin(): a device takes the root topic of its reports when it is created.
|
||||
AlexaDevice *getDevice(const String &name, const String &endpointId);
|
||||
|
||||
// What the library prints on Serial: AlexaLogLevel::NONE, ERROR, INFO (the default) or DEBUG
|
||||
void setLogLevel(AlexaLogLevel level);
|
||||
|
||||
// false before begin(): the sketch sets the clock itself (its own configTime() with a time zone, an RTC).
|
||||
// begin() then leaves SNTP alone; the library only reads time().
|
||||
void setTimeSource(bool useSntp);
|
||||
|
||||
// AlexaTransport: what the devices' reports are sent and stamped with
|
||||
AlexaSendResult publish(const char *topic, JsonDocument &doc) override;
|
||||
void timestamp(char *buffer, size_t size) override;
|
||||
|
||||
private:
|
||||
static const int MAX_RETRY_COUNT = 2; // Define maximum retry count
|
||||
static const unsigned long CLOCK_WAIT_MS = 5000; // How long the first connect waits for SNTP
|
||||
static const unsigned long RECONNECT_INTERVAL_MS = 5000;
|
||||
static const unsigned long DISCOVERY_WINDOW_MS = 5000; // How long the backend keeps collecting a discovery answer
|
||||
int retryCountPOST=0;
|
||||
int retryCountGET=0;
|
||||
static const unsigned long DISCOVERY_RETRY_MS = 20; // Pause before a refused discovery publish is tried again
|
||||
|
||||
size_t discoveryNext = 0; // Next device to announce while a discovery answer is still going out
|
||||
bool discoveryPending = false; // publishDiscovery() stopped early (client out-queue full); loop() finishes it
|
||||
unsigned long discoveryStarted = 0; // millis() when the Discover arrived
|
||||
|
||||
boolean clearToSend=false;
|
||||
AsyncMqttClient mqttClient; // MQTT client instance
|
||||
String rootTopic; // Root topic for communication
|
||||
String discoverTopic; // The topic we listen on for discovery messages
|
||||
String discoverTopicSend; // The topic we send discovery messages
|
||||
|
||||
String TopicESP; // The topic we listen on for esp messages
|
||||
|
||||
const char *mqttUsername; // Username for authentication
|
||||
const char *mqttPassword; // Password for authentication
|
||||
|
||||
unsigned long lastReconnectTime;
|
||||
int reconnectAttempt;
|
||||
|
||||
HTTPClient httpGET;
|
||||
HTTPClient httpPOST;
|
||||
|
||||
WiFiClient wifiGET;
|
||||
WiFiClient wifiPOST;
|
||||
JsonDocument inputDoc;
|
||||
String directiveFilter; // The subscription that delivers the directives of every endpoint
|
||||
|
||||
std::deque<AlexaDevice> devices; // Collection of devices (deque: pointers handed out by getDevice stay valid)
|
||||
|
||||
Alex2ESPState _state;
|
||||
AsyncMqttClientDisconnectReason _disconnectReason;
|
||||
Alex2ESPState state;
|
||||
AsyncMqttClientDisconnectReason disconnectReason;
|
||||
bool useSntp;
|
||||
unsigned long beginTime; // millis() when begin() ran
|
||||
unsigned long lastReconnectTime;
|
||||
uint16_t discoverSubscription; // Packet ids of the two SUBSCRIBEs, to match their acknowledgements
|
||||
uint16_t directiveSubscription;
|
||||
uint8_t subscriptionsPending;
|
||||
|
||||
// MQTT connection details
|
||||
const char *mqttServer = "alex2mqtt.stormysdream.club";
|
||||
uint16_t mqttPort = 1883;
|
||||
bool discoveryRequested; // Set by the MQTT callback, taken by loop()
|
||||
bool discoveryActive; // An answer is going out
|
||||
bool discoveryDeferred; // The answer had to pause: the MQTT client or the heap was full
|
||||
size_t discoveryNext; // Next device to announce
|
||||
size_t discoveryAnnounced; // Devices announced in this answer
|
||||
unsigned long discoveryStarted; // millis() when the Discover arrived
|
||||
unsigned long discoveryLastAttempt; // millis() of the publish that was refused
|
||||
|
||||
// Internal event handlers
|
||||
AlexaDirectiveBuffer directive; // The directive that waits for loop(), or is still arriving
|
||||
AlexaRecentIds recentIds; // messageIds of the last directives handled
|
||||
|
||||
// Internal event handlers (called by the MQTT client from the network context: they only take notes)
|
||||
void onMqttConnect(bool sessionPresent);
|
||||
void onMqttDisconnect(AsyncMqttClientDisconnectReason reason);
|
||||
void onSubscribe(uint16_t packetId, uint8_t qos);
|
||||
void onMessage(char *topic, char *payload, AsyncMqttClientMessageProperties properties, size_t length, size_t index, size_t total);
|
||||
|
||||
//loop processing function
|
||||
void connectWhenClockIsSet();
|
||||
void handleMqttReconnection();
|
||||
void continueDiscovery();
|
||||
void publishDiscovery();
|
||||
void finishDiscovery();
|
||||
void processHttpPost();
|
||||
void processHttpGet();
|
||||
void processDirective();
|
||||
|
||||
AlexaDevice *findDevice(const char *endpointId, size_t length);
|
||||
AlexaSendResult trySend(const char *topic, JsonDocument &doc, size_t *length);
|
||||
void logRefusal(const char *topic, AlexaSendResult result, size_t length);
|
||||
};
|
||||
|
||||
#endif // ALEX2ESP_H
|
||||
|
|
|
|||
300
src/AlexaBridgeLogic.cpp
Normal file
300
src/AlexaBridgeLogic.cpp
Normal 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
123
src/AlexaBridgeLogic.h
Normal 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
14
src/AlexaCompat.h
Normal 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
|
||||
|
|
@ -1,7 +1,8 @@
|
|||
#include "AlexaDevice.h"
|
||||
#include "AlexaLog.h"
|
||||
|
||||
AlexaDevice::AlexaDevice(const String& name, const String& rootTopic, const String& endpointId)
|
||||
: name(name), endpointId(endpointId), rootTopic(rootTopic) {
|
||||
AlexaDevice::AlexaDevice(const String& name, const String& rootTopic, const String& endpointId, AlexaTransport* transport)
|
||||
: name(name), endpointId(endpointId), rootTopic(rootTopic), transport(transport) {
|
||||
// Initialize event arrays to nullptr
|
||||
for (int i = 0; i < MAX_EVENTS; ++i) {
|
||||
eventNames[i] = nullptr; // Set all event names to nullptr
|
||||
|
|
@ -22,6 +23,10 @@ String AlexaDevice::getEndpointId() const {
|
|||
return endpointId;
|
||||
}
|
||||
|
||||
bool AlexaDevice::hasEndpointId(const char* id, size_t length) const {
|
||||
return id != nullptr && endpointId.length() == length && memcmp(endpointId.c_str(), id, length) == 0;
|
||||
}
|
||||
|
||||
DisplayCategory AlexaDevice::getDisplayCategory() const {
|
||||
return displayCategory;
|
||||
}
|
||||
|
|
@ -78,8 +83,7 @@ AlexaInterface* AlexaDevice::addCapability(AlexaInterfaceType type) {
|
|||
// If the interface doesn't exist, create a new one
|
||||
capabilities.emplace_back(type);
|
||||
|
||||
Serial.print("Created new interface with type: ");
|
||||
Serial.println(capabilities.back().getTypeString());
|
||||
ALEX2ESP_LOGI("%s: capability %s added", endpointId.c_str(), capabilities.back().getTypeString().c_str());
|
||||
|
||||
// Return a pointer to the newly created device
|
||||
return &capabilities.back();
|
||||
|
|
@ -158,9 +162,8 @@ void AlexaDevice::triggerEvent(const char* eventName, const JsonDocument& direct
|
|||
return;
|
||||
}
|
||||
}
|
||||
// If no event was found, print an error message
|
||||
// No handler: the directive goes unanswered, which Alexa reports as a device that does not respond
|
||||
if (warnIfMissing) {
|
||||
Serial.print("No event registered for: ");
|
||||
Serial.println(eventName);
|
||||
ALEX2ESP_LOGE("%s: no handler registered for %s, directive not handled", endpointId.c_str(), eventName);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -10,6 +10,7 @@
|
|||
#include <string>
|
||||
#include <deque>
|
||||
#include "AlexaStatusMessage.h"
|
||||
#include "AlexaTransport.h"
|
||||
|
||||
|
||||
#define MAX_EVENTS 10
|
||||
|
|
@ -147,13 +148,18 @@ public:
|
|||
|
||||
class AlexaDevice {
|
||||
public:
|
||||
AlexaDevice(const String& name, const String& rootTopic, const String& endpointId);
|
||||
// Devices are created by Alex2ESP::getDevice(), which passes the bridge as the transport that publishes the
|
||||
// device's reports. A device built without one can describe itself, but its reports cannot be sent.
|
||||
AlexaDevice(const String& name, const String& rootTopic, const String& endpointId, AlexaTransport* transport = nullptr);
|
||||
|
||||
void setName(const String& name);
|
||||
String getName() const;
|
||||
|
||||
String getEndpointId() const;
|
||||
|
||||
// Compares the endpoint id with `length` characters at `id` (not NUL-terminated: a part of an MQTT topic)
|
||||
bool hasEndpointId(const char* id, size_t length) const;
|
||||
|
||||
DisplayCategory getDisplayCategory() const;
|
||||
void setDisplayCategory(DisplayCategory category);
|
||||
|
||||
|
|
@ -183,13 +189,14 @@ public:
|
|||
void triggerEvent(const char* eventName, const JsonDocument& directive, const AlexaInterfaceType& type, bool warnIfMissing = true) const;
|
||||
|
||||
AlexaStatusMessage buildStatusMessage(const String& correlationToken,const bool isResponse=false) {
|
||||
return AlexaStatusMessage(correlationToken,rootTopic,endpointId,isResponse);
|
||||
return AlexaStatusMessage(correlationToken,rootTopic,endpointId,isResponse,transport);
|
||||
}
|
||||
|
||||
private:
|
||||
String name;
|
||||
String endpointId;
|
||||
String rootTopic;
|
||||
AlexaTransport* transport;
|
||||
const char* eventNames[MAX_EVENTS];
|
||||
void (*eventCallbacks[MAX_EVENTS])(const JsonDocument&, const AlexaInterfaceType&);
|
||||
|
||||
|
|
|
|||
|
|
@ -6,6 +6,7 @@
|
|||
#include <unordered_map>
|
||||
#include <Arduino.h>
|
||||
#include <ArduinoJson.h>
|
||||
#include "AlexaLog.h"
|
||||
|
||||
// Enum for Alexa Interface Types
|
||||
enum class AlexaInterfaceType
|
||||
|
|
@ -617,8 +618,7 @@ public:
|
|||
if (deserializeJson(payloadDoc, directive.payload) == DeserializationError::Ok && payloadDoc.is<JsonObject>()) {
|
||||
directiveObj["payload"] = payloadDoc.as<JsonObject>();
|
||||
} else {
|
||||
Serial.print("[Alex2ESP] action mapping payload is not a JSON object, omitted: ");
|
||||
Serial.println(directive.payload);
|
||||
ALEX2ESP_LOGE("action mapping payload is not a JSON object, omitted: %s", directive.payload.c_str());
|
||||
}
|
||||
}
|
||||
return doc;
|
||||
|
|
|
|||
49
src/AlexaLog.cpp
Normal file
49
src/AlexaLog.cpp
Normal 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
56
src/AlexaLog.h
Normal 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
|
||||
|
|
@ -1,3 +1,90 @@
|
|||
#include "AlexaStatusMessage.h"
|
||||
#include "AlexaLog.h"
|
||||
|
||||
char AlexaStatusMessage::outputString[MAX_STATUS_REPORT_SIZE];
|
||||
// What a 1.x sketch wrote where the time belongs when it built a property by hand: the backend's HTTP route
|
||||
// replaced it. Reports leave over MQTT now, where nothing rewrites them, so the library fills the time in.
|
||||
static const char TIME_PLACEHOLDER[] PROGMEM = "{REPLACE_WITH_DATETIME}";
|
||||
|
||||
AlexaStatusMessage::AlexaStatusMessage(const String &correlationToken, const String &rootTopic, const String &endpointId, const bool isResponse, AlexaTransport *transport)
|
||||
: rootTopic(rootTopic), endpointId(endpointId), transport(transport)
|
||||
{
|
||||
JsonObject event = doc["event"].to<JsonObject>();
|
||||
|
||||
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;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -4,10 +4,9 @@
|
|||
#include <Arduino.h>
|
||||
#include <ArduinoJson.h>
|
||||
#include "AlexaInterface.h"
|
||||
#include "AlexaTransport.h"
|
||||
#include "AlexaUtils.h"
|
||||
|
||||
#define MAX_STATUS_REPORT_SIZE 2048
|
||||
|
||||
enum class EndpointHealth
|
||||
{
|
||||
OK,
|
||||
|
|
@ -28,32 +27,9 @@ enum class TemperatureSensorScale
|
|||
class AlexaStatusMessage
|
||||
{
|
||||
public:
|
||||
AlexaStatusMessage(const String &correlationToken, const String &rootTopic, const String &endpointId, const bool isResponse)
|
||||
{
|
||||
this->endpointId = endpointId;
|
||||
this->rootTopic = rootTopic;
|
||||
|
||||
JsonObject event = doc["event"].to<JsonObject>();
|
||||
|
||||
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>();
|
||||
}
|
||||
// Built by AlexaDevice::buildStatusMessage(), which passes the bridge as the transport: it publishes the report
|
||||
// and supplies the time of every property. A message without a transport cannot be sent.
|
||||
AlexaStatusMessage(const String &correlationToken, const String &rootTopic, const String &endpointId, const bool isResponse, AlexaTransport *transport = nullptr);
|
||||
|
||||
AlexaStatusMessage &AddHealthProp(EndpointHealth endpointHealth, unsigned int uncertaintyInMs = 0)
|
||||
{
|
||||
|
|
@ -93,66 +69,29 @@ public:
|
|||
return AddProperty(AlexaInterfaceType::TOGGLE_CONTROLLER, "toggleState", value, uncertaintyInMs,instanceName);
|
||||
}
|
||||
|
||||
AlexaStatusMessage &AddContextProp(const JsonObject &property)
|
||||
{
|
||||
contextProperties.add(property);
|
||||
return *this; // Return a reference to the current object
|
||||
}
|
||||
// Adds a property the sketch built itself. Its timeOfSample is set to the current time when the object has
|
||||
// none or still carries the 1.x placeholder "{REPLACE_WITH_DATETIME}".
|
||||
AlexaStatusMessage &AddContextProp(const JsonObject &property);
|
||||
|
||||
// Queue the report for delivery. Returns false (and says so on Serial) when the report does not fit the
|
||||
// MAX_STATUS_REPORT_SIZE buffer or the send queue is full: a report is never sent truncated.
|
||||
bool send()
|
||||
{
|
||||
doc.shrinkToFit();
|
||||
String topic = rootTopic + "/" + endpointId + "/alexaResponce";
|
||||
size_t needed = measureJson(doc);
|
||||
if (needed > sizeof(outputString) - 1)
|
||||
{
|
||||
Serial.printf("[Alex2ESP] status report for %s is %u bytes, limit is %u - not sent\n",
|
||||
endpointId.c_str(), (unsigned)needed, (unsigned)(sizeof(outputString) - 1));
|
||||
doc.clear();
|
||||
return false;
|
||||
}
|
||||
serializeJson(doc, outputString);
|
||||
doc.clear();
|
||||
if (!AlexaUtils::enqueue(outputString, topic.c_str()))
|
||||
{
|
||||
Serial.println("[Alex2ESP] send queue full, status report dropped");
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
// Publishes the report on <root>/<endpointId>/alexaResponce now. Returns false when nothing was sent: no
|
||||
// session with the broker, the MQTT client or the heap cannot take the report, or it is over
|
||||
// ALEX2ESP_MAX_MESSAGE bytes. The reason is on Serial; a report is never sent truncated.
|
||||
bool send();
|
||||
|
||||
private:
|
||||
String rootTopic;
|
||||
String endpointId;
|
||||
AlexaTransport *transport;
|
||||
JsonDocument doc;
|
||||
JsonArray contextProperties;
|
||||
static char outputString[MAX_STATUS_REPORT_SIZE];
|
||||
|
||||
// TODO: make real uuid4 gen function
|
||||
String generateMessageId()
|
||||
{
|
||||
char buffer[38];
|
||||
const char charset[] = "0123456789abcdefABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz"; // Allowed characters
|
||||
|
||||
// Seed the random number generator (optional)
|
||||
srand(static_cast<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);
|
||||
}
|
||||
String generateMessageId();
|
||||
void setTimeOfSample(JsonObject property);
|
||||
|
||||
template <typename T>
|
||||
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["name"] = propertyName;
|
||||
|
||||
|
|
@ -161,10 +100,8 @@ private:
|
|||
}
|
||||
prop["value"] = value;
|
||||
|
||||
prop["timeOfSample"] = "{REPLACE_WITH_DATETIME}";
|
||||
setTimeOfSample(prop);
|
||||
prop["uncertaintyInMilliseconds"] = uncertaintyInMs;
|
||||
|
||||
contextProperties.add(prop.as<JsonObject>());
|
||||
return *this;
|
||||
}
|
||||
};
|
||||
|
|
|
|||
43
src/AlexaTransport.h
Normal file
43
src/AlexaTransport.h
Normal 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
|
||||
|
|
@ -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
|
||||
}
|
||||
|
|
@ -1,89 +1,26 @@
|
|||
#ifndef ALEXA_UTILS_H
|
||||
#define ALEXA_UTILS_H
|
||||
|
||||
#include <vector>
|
||||
#include <string>
|
||||
#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
|
||||
{
|
||||
public:
|
||||
static const int MAX_PAYLOAD_LENGTH = 2048; // Maximum length for receive packets
|
||||
static char receivePayload[MAX_PAYLOAD_LENGTH];
|
||||
|
||||
static bool enqueueReceive(const char* packet);
|
||||
static bool dequeueReceive(String& packet,bool remove);
|
||||
static bool isReceiveQueueEmpty();
|
||||
static bool isReceiveQueueFull();
|
||||
|
||||
// Queue a packet for the HTTP sender. Returns false (nothing queued) when the queue is full or the
|
||||
// packet/topic would not fit its slot - a packet is never truncated.
|
||||
static bool enqueue(const char* packet, const char* topic);
|
||||
static bool dequeue(String& topic, String& packet);
|
||||
static bool isQueueEmpty();
|
||||
static bool isQueueFull();
|
||||
static const char* dequeueVals(bool remove);
|
||||
|
||||
static void log(const char *string);
|
||||
static void logln(const char *string);
|
||||
|
||||
static void log(int num);
|
||||
static void logln(int num);
|
||||
|
||||
static void log(uint16_t num);
|
||||
static void logln(uint16_t num);
|
||||
|
||||
static void log(uint8_t num);
|
||||
static void logln(uint8_t num);
|
||||
|
||||
static void logln(String string);
|
||||
|
||||
// Incrementing message ID
|
||||
static uint8_t nextMessageId;
|
||||
|
||||
// Static function to set the MQTT client reference
|
||||
static void printMemoryInfo()
|
||||
{
|
||||
size_t freeHeap = ESP.getFreeHeap();
|
||||
#ifdef ESP8266
|
||||
size_t freeStack = ESP.getFreeContStack(); // Requires ESP8266 core 3.0.0+
|
||||
#else
|
||||
size_t freeStack = 0; // ESP32: untested, no getFreeContStack()
|
||||
#endif
|
||||
size_t sketchSize = ESP.getSketchSize();
|
||||
size_t freeSketchSpace = ESP.getFreeSketchSpace();
|
||||
|
||||
// Print in custom format
|
||||
Serial.printf("freeHeap: %u, freeStack: %u, freeROM: %u, usedROM: %u\n",
|
||||
freeHeap,
|
||||
freeStack,
|
||||
freeSketchSpace,
|
||||
sketchSize);
|
||||
(unsigned)ESP.getFreeHeap(),
|
||||
(unsigned)freeStack,
|
||||
(unsigned)ESP.getFreeSketchSpace(),
|
||||
(unsigned)ESP.getSketchSize());
|
||||
}
|
||||
|
||||
|
||||
private:
|
||||
static const int MAX_QUEUE_LENGTH = 5; // Define maximum queue length
|
||||
static const int MAX_TOPIC_LENGTH = 128; // Maximum length for topics
|
||||
static const int MAX_PACKET_LENGTH = 2048; // Maximum length for sending packets
|
||||
static constexpr char MESSAGE_PAYLOAD_SPLIT[32] = "{MESSAGE_PAYLOAD_SPLIT}";
|
||||
|
||||
|
||||
// Static arrays for topics and packets
|
||||
static char topicQueue[MAX_QUEUE_LENGTH][MAX_TOPIC_LENGTH];
|
||||
static char packetQueue[MAX_QUEUE_LENGTH][MAX_PACKET_LENGTH];
|
||||
static char combinedData[MAX_TOPIC_LENGTH + MAX_PACKET_LENGTH + 32];
|
||||
|
||||
static char receiveQueue[MAX_QUEUE_LENGTH][16];
|
||||
static int queueStart; // Index of the front of the queue
|
||||
static int queueEnd; // Index of the back of the queue
|
||||
static int queueCount; // Number of items in the queue
|
||||
|
||||
static int queueStartReceive;
|
||||
static int queueEndReceive;
|
||||
static int queueCountReceive;
|
||||
};
|
||||
|
||||
#endif // ALEXA_UTILS_H
|
||||
|
|
|
|||
575
test/test_bridge_logic/test_main.cpp
Normal file
575
test/test_bridge_logic/test_main.cpp
Normal 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();
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue