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
765
src/Alex2ESP.cpp
765
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();
|
||||
}
|
||||
}
|
||||
else if (index != 0 || total != length)
|
||||
{
|
||||
// AsyncMqttClient hands a large publish over in fragments; a directive id never is, so a fragment can only be part of one
|
||||
AlexaUtils::logln("Ignoring fragmented directive id");
|
||||
}
|
||||
else
|
||||
{
|
||||
// The payload is not NUL-terminated and belongs to the MQTT client: copy exactly `length` bytes out of it
|
||||
String uuid;
|
||||
uuid.concat(payload, length);
|
||||
|
||||
for (auto &device : devices)
|
||||
{
|
||||
String directiveTopic = String(rootTopic) + "/" + device.getEndpointId() + "/alexaDirective_e";
|
||||
|
||||
if (strcmp(topic, directiveTopic.c_str()) == 0)
|
||||
{
|
||||
AlexaUtils::logln(topic);
|
||||
AlexaUtils::logln(uuid);
|
||||
if (!AlexaUtils::enqueueReceive(uuid.c_str()))
|
||||
{
|
||||
Serial.println("[Alex2ESP] receive queue full, directive dropped");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
AlexaUtils::logln("OnMessage End");
|
||||
AlexaUtils::printMemoryInfo();
|
||||
clearToSend = true;
|
||||
return disconnectReason;
|
||||
}
|
||||
|
||||
// Answer a Discover: one discovery object per device, published straight to <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");
|
||||
ALEX2ESP_LOGE("the MQTT client refused to subscribe (%u bytes of heap free): no directives until it reconnects", (unsigned)ESP.getFreeHeap());
|
||||
return;
|
||||
}
|
||||
ALEX2ESP_LOGI("connected to %s, subscribing", MQTT_SERVER);
|
||||
}
|
||||
|
||||
const char *payload = AlexaUtils::dequeueVals(false);
|
||||
void Alex2ESP::onSubscribe(uint16_t packetId, uint8_t qos)
|
||||
{
|
||||
if (packetId != discoverSubscription && packetId != directiveSubscription)
|
||||
{
|
||||
return;
|
||||
}
|
||||
const String &topic = (packetId == discoverSubscription) ? discoverTopic : directiveFilter;
|
||||
if (qos == SUBSCRIPTION_REFUSED)
|
||||
{
|
||||
ALEX2ESP_LOGE("the broker refused the subscription to %s", topic.c_str());
|
||||
return;
|
||||
}
|
||||
|
||||
httpPOST.begin(wifiPOST, "http://alex2mqtt.stormysdream.club/Alex2ESP");
|
||||
httpPOST.setAuthorization(mqttUsername, mqttPassword);
|
||||
httpPOST.addHeader("Content-Type", "text/plain");
|
||||
|
||||
AlexaUtils::log("Sending Data: ");
|
||||
AlexaUtils::logln(payload);
|
||||
|
||||
int httpResponseCode = httpPOST.POST(payload);
|
||||
|
||||
if (httpResponseCode != 200)
|
||||
{
|
||||
|
||||
if (retryCountPOST >= MAX_RETRY_COUNT)
|
||||
{
|
||||
AlexaUtils::dequeueVals(true);
|
||||
}
|
||||
|
||||
AlexaUtils::log("POST Error code: ");
|
||||
AlexaUtils::logln(httpResponseCode);
|
||||
retryCountPOST++;
|
||||
if (httpResponseCode > 0)
|
||||
{
|
||||
String response = httpPOST.getString();
|
||||
Serial.println("Response:");
|
||||
Serial.println(response);
|
||||
}
|
||||
}
|
||||
else
|
||||
{
|
||||
retryCountPOST = 0;
|
||||
AlexaUtils::dequeueVals(true);
|
||||
}
|
||||
httpPOST.end();
|
||||
ALEX2ESP_LOGD("subscribed to %s", topic.c_str());
|
||||
if (subscriptionsPending > 0 && --subscriptionsPending == 0)
|
||||
{
|
||||
state = Alex2ESPState::CONNECTED;
|
||||
ALEX2ESP_LOGI("ready: %u device(s)", (unsigned)devices.size());
|
||||
}
|
||||
}
|
||||
void Alex2ESP::processHttpGet()
|
||||
|
||||
void Alex2ESP::onMqttDisconnect(AsyncMqttClientDisconnectReason reason)
|
||||
{
|
||||
if (!AlexaUtils::isReceiveQueueEmpty())
|
||||
state = Alex2ESPState::DISCONNECTED;
|
||||
disconnectReason = reason;
|
||||
directive.cancelArrival();
|
||||
ALEX2ESP_LOGI("disconnected (reason %u)", (unsigned)reason);
|
||||
}
|
||||
|
||||
// Runs in the network context, once per fragment of a message. It only takes notes: loop() answers a Discover and
|
||||
// parses and dispatches a directive, on the sketch's own stack.
|
||||
void Alex2ESP::onMessage(char *topic, char *payload, AsyncMqttClientMessageProperties properties, size_t length, size_t index, size_t total)
|
||||
{
|
||||
if (strcmp(topic, discoverTopic.c_str()) == 0)
|
||||
{
|
||||
AlexaUtils::logln("processHttpGet");
|
||||
|
||||
String uuid;
|
||||
AlexaUtils::dequeueReceive(uuid, false);
|
||||
AlexaUtils::log("Sending get request for uuid: ");
|
||||
// AlexaUtils::logln(uuid);
|
||||
|
||||
httpGET.begin(wifiGET, "http://alex2mqtt.stormysdream.club/Alex2ESP/" + uuid);
|
||||
httpGET.setAuthorization(mqttUsername, mqttPassword);
|
||||
int httpResponseCode = httpGET.GET();
|
||||
|
||||
Serial.print("HTTP Response code: ");
|
||||
Serial.println(httpResponseCode);
|
||||
if (httpResponseCode != 200)
|
||||
// The payload is unused: answer once per message even when TCP split it
|
||||
if (index == 0)
|
||||
{
|
||||
Serial.print("Error code: ");
|
||||
Serial.println(httpResponseCode);
|
||||
if (retryCountGET >= MAX_RETRY_COUNT)
|
||||
{
|
||||
AlexaUtils::dequeueReceive(uuid, true);
|
||||
}
|
||||
discoveryRequested = true;
|
||||
discoveryStarted = millis();
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
retryCountGET++;
|
||||
size_t endpointLength = 0;
|
||||
const char *endpointId = AlexaBridgeLogic::directiveEndpoint(topic, rootTopic.c_str(), &endpointLength);
|
||||
if (endpointId == nullptr)
|
||||
{
|
||||
return;
|
||||
}
|
||||
if (findDevice(endpointId, endpointLength) == nullptr)
|
||||
{
|
||||
// <root>/+/alexaDirective delivers the directives of every endpoint of the account, also those of other
|
||||
// boards: not an error, and not worth a buffer
|
||||
if (index == 0)
|
||||
{
|
||||
ALEX2ESP_LOGD("directive on %s is for another board", topic);
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
switch (directive.append(payload, length, index, total))
|
||||
{
|
||||
case AlexaDirectiveBuffer::Result::TOO_LARGE:
|
||||
ALEX2ESP_LOGE("directive of %u bytes on %s dropped: the limit is %u (ALEX2ESP_MAX_DIRECTIVE)", (unsigned)total, topic, (unsigned)ALEX2ESP_MAX_DIRECTIVE);
|
||||
break;
|
||||
case AlexaDirectiveBuffer::Result::BUSY:
|
||||
ALEX2ESP_LOGE("directive on %s dropped: the one before it still waits for loop()", topic);
|
||||
break;
|
||||
case AlexaDirectiveBuffer::Result::NO_MEMORY:
|
||||
ALEX2ESP_LOGE("directive of %u bytes on %s dropped: no memory (%u bytes of heap free)", (unsigned)total, topic, (unsigned)ESP.getFreeHeap());
|
||||
break;
|
||||
case AlexaDirectiveBuffer::Result::OUT_OF_ORDER:
|
||||
ALEX2ESP_LOGE("directive on %s dropped: the fragment at byte %u of %u does not continue it", topic, (unsigned)index, (unsigned)total);
|
||||
break;
|
||||
case AlexaDirectiveBuffer::Result::EMPTY:
|
||||
ALEX2ESP_LOGE("directive on %s dropped: it is empty", topic);
|
||||
break;
|
||||
case AlexaDirectiveBuffer::Result::DUPLICATE:
|
||||
ALEX2ESP_LOGI("repeated directive on %s ignored", topic);
|
||||
break;
|
||||
case AlexaDirectiveBuffer::Result::INCOMPLETE:
|
||||
case AlexaDirectiveBuffer::Result::COMPLETE:
|
||||
case AlexaDirectiveBuffer::Result::IGNORED:
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
void Alex2ESP::processDirective()
|
||||
{
|
||||
if (!directive.ready())
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
JsonDocument message;
|
||||
size_t length = directive.length();
|
||||
DeserializationError error = deserializeJson(message, directive.data(), length);
|
||||
// The text is not needed any more. Releasing it before the handler runs frees the heap for the report and
|
||||
// lets the next directive arrive while the sketch deals with this one.
|
||||
directive.release();
|
||||
|
||||
if (error)
|
||||
{
|
||||
char reason[24];
|
||||
strncpy_P(reason, reinterpret_cast<PGM_P>(error.f_str()), sizeof(reason) - 1);
|
||||
reason[sizeof(reason) - 1] = '\0';
|
||||
ALEX2ESP_LOGE("directive of %u bytes dropped: not JSON (%s)", (unsigned)length, reason);
|
||||
return;
|
||||
}
|
||||
|
||||
// Alex2MQTT publishes the directive itself: {"header": ..., "endpoint": ..., "payload": ...}
|
||||
const JsonDocument &received = message;
|
||||
if (!received["header"].is<JsonObjectConst>())
|
||||
{
|
||||
ALEX2ESP_LOGE("directive of %u bytes dropped: it has no header", (unsigned)length);
|
||||
return;
|
||||
}
|
||||
const char *endpointId = received["endpoint"]["endpointId"] | "";
|
||||
const char *interfaceName = received["header"]["namespace"] | "";
|
||||
const char *name = received["header"]["name"] | "";
|
||||
const char *messageId = received["header"]["messageId"] | "";
|
||||
|
||||
AlexaDevice *device = findDevice(endpointId, strlen(endpointId));
|
||||
if (device == nullptr)
|
||||
{
|
||||
ALEX2ESP_LOGE("directive %s.%s dropped: it names the endpoint '%s', which this board does not have", interfaceName, name, endpointId);
|
||||
return;
|
||||
}
|
||||
if (recentIds.seenBefore(messageId))
|
||||
{
|
||||
ALEX2ESP_LOGI("%s: repeated directive %s.%s ignored", endpointId, interfaceName, name);
|
||||
return;
|
||||
}
|
||||
|
||||
ALEX2ESP_LOGI("%s <- %s.%s", endpointId, interfaceName, name);
|
||||
ALEX2ESP_LOGD("%u bytes, %u bytes of heap free", (unsigned)length, (unsigned)ESP.getFreeHeap());
|
||||
|
||||
device->triggerEvent("DirectiveReceived", message, AlexaInterfaceType::UNKNOWN, false);
|
||||
if (strcmp(interfaceName, "Alexa") == 0 && strcmp(name, "ReportState") == 0)
|
||||
{
|
||||
device->triggerEvent("ReportState", message, AlexaInterfaceType::UNKNOWN);
|
||||
}
|
||||
else
|
||||
{
|
||||
device->triggerEvent("Event", message, AlexaInterfaceUtils::fromString(interfaceName));
|
||||
}
|
||||
}
|
||||
|
||||
// Answer a Discover: one discovery object per device, published straight to <root>/discover_r. The backend accepts
|
||||
// one endpoint object per message and collects everything that arrives within 1 s for Alexa's answer (5 s for its
|
||||
// proactive push), and it takes the answers to one request as the complete list of the board. So every Discover
|
||||
// starts the list from the top, also when the answer to the one before is still going out.
|
||||
void Alex2ESP::continueDiscovery()
|
||||
{
|
||||
if (discoveryRequested)
|
||||
{
|
||||
discoveryRequested = false;
|
||||
discoveryActive = true;
|
||||
discoveryDeferred = false;
|
||||
discoveryNext = 0;
|
||||
discoveryAnnounced = 0;
|
||||
}
|
||||
if (!discoveryActive)
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
if (millis() - discoveryStarted > DISCOVERY_WINDOW_MS)
|
||||
{
|
||||
// The backend has stopped listening
|
||||
ALEX2ESP_LOGE("discovery gave up: %u device(s) not announced", (unsigned)(devices.size() - discoveryAnnounced));
|
||||
discoveryActive = false;
|
||||
return;
|
||||
}
|
||||
if (discoveryDeferred && millis() - discoveryLastAttempt < DISCOVERY_RETRY_MS)
|
||||
{
|
||||
return;
|
||||
}
|
||||
publishDiscovery();
|
||||
}
|
||||
|
||||
// publish() copies the object into the client's out-queue, where it stays on the heap until it has been written to
|
||||
// the socket, and the client refuses another one once the free heap drops under 4 KB - so rather than skip that
|
||||
// device we stop, remember where we got to and let loop() carry on once the queue has drained.
|
||||
void Alex2ESP::publishDiscovery()
|
||||
{
|
||||
while (discoveryNext < devices.size())
|
||||
{
|
||||
const AlexaDevice &device = devices[discoveryNext];
|
||||
JsonDocument json = device.getDeviceJSON();
|
||||
size_t length = 0;
|
||||
AlexaSendResult result = trySend(discoverTopicSend.c_str(), json, &length);
|
||||
|
||||
if (result == AlexaSendResult::REFUSED || result == AlexaSendResult::NOT_CONNECTED)
|
||||
{
|
||||
if (!discoveryDeferred)
|
||||
{
|
||||
ALEX2ESP_LOGI("discovery deferred at %s (%u bytes of heap free)", device.getEndpointId().c_str(), (unsigned)ESP.getFreeHeap());
|
||||
}
|
||||
discoveryDeferred = true;
|
||||
discoveryLastAttempt = millis();
|
||||
return;
|
||||
}
|
||||
|
||||
if (result == AlexaSendResult::OK)
|
||||
{
|
||||
discoveryAnnounced++;
|
||||
}
|
||||
else
|
||||
{
|
||||
retryCountGET = 0;
|
||||
AlexaUtils::dequeueReceive(uuid, true);
|
||||
String payloadStr = httpGET.getString();
|
||||
size_t length = payloadStr.length();
|
||||
|
||||
if (length < (size_t)(AlexaUtils::MAX_PAYLOAD_LENGTH - 1))
|
||||
{
|
||||
payloadStr.toCharArray(AlexaUtils::receivePayload, length + 1);
|
||||
payloadStr = "";
|
||||
AlexaUtils::receivePayload[length + 1] = '\0';
|
||||
// AlexaUtils::logln(AlexaUtils::receivePayload);
|
||||
|
||||
DeserializationError error = deserializeJson(inputDoc, AlexaUtils::receivePayload);
|
||||
if (error)
|
||||
{
|
||||
AlexaUtils::log("Failed to parse message: ");
|
||||
AlexaUtils::logln(String(error.f_str()));
|
||||
inputDoc.clear();
|
||||
}
|
||||
else
|
||||
{
|
||||
AlexaUtils::logln("JSON message parsed successfully!");
|
||||
|
||||
for (auto &device : devices)
|
||||
{
|
||||
if (strcmp(device.getEndpointId().c_str(), inputDoc["directive"]["endpoint"]["endpointId"] | "") == 0)
|
||||
{
|
||||
serializeJson(inputDoc, Serial);
|
||||
Serial.println();
|
||||
device.triggerEvent("DirectiveReceived", inputDoc["directive"], AlexaInterfaceType::UNKNOWN, false);
|
||||
|
||||
if (inputDoc["directive"]["header"]["name"] == "ReportState")
|
||||
{
|
||||
device.triggerEvent("ReportState", inputDoc["directive"], AlexaInterfaceType::UNKNOWN);
|
||||
}
|
||||
else
|
||||
{
|
||||
device.triggerEvent("Event", inputDoc["directive"], AlexaInterfaceUtils::fromString(inputDoc["directive"]["header"]["namespace"]));
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
// This object can never be sent: say so and announce the others
|
||||
logRefusal(discoverTopicSend.c_str(), result, length);
|
||||
ALEX2ESP_LOGE("device %s is not announced", device.getEndpointId().c_str());
|
||||
}
|
||||
httpGET.end();
|
||||
discoveryNext++;
|
||||
}
|
||||
|
||||
discoveryActive = false;
|
||||
ALEX2ESP_LOGI("discovery: %u of %u device(s) announced", (unsigned)discoveryAnnounced, (unsigned)devices.size());
|
||||
}
|
||||
|
||||
AlexaSendResult Alex2ESP::publish(const char *topic, JsonDocument &doc)
|
||||
{
|
||||
size_t length = 0;
|
||||
AlexaSendResult result = trySend(topic, doc, &length);
|
||||
if (result == AlexaSendResult::OK)
|
||||
{
|
||||
ALEX2ESP_LOGD("%u bytes -> %s, %u bytes of heap free", (unsigned)length, topic, (unsigned)ESP.getFreeHeap());
|
||||
}
|
||||
else
|
||||
{
|
||||
logRefusal(topic, result, length);
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
void Alex2ESP::timestamp(char *buffer, size_t size)
|
||||
{
|
||||
AlexaBridgeLogic::formatTimestamp(time(nullptr), buffer, size);
|
||||
}
|
||||
|
||||
// Sends the document whole or not at all, and says which. The caller reports a refusal.
|
||||
AlexaSendResult Alex2ESP::trySend(const char *topic, JsonDocument &doc, size_t *length)
|
||||
{
|
||||
AlexaSendResult result = AlexaBridgeLogic::checkMessage(doc, ALEX2ESP_MAX_MESSAGE, length);
|
||||
if (result != AlexaSendResult::OK)
|
||||
{
|
||||
return result;
|
||||
}
|
||||
if (!mqttClient.connected())
|
||||
{
|
||||
return AlexaSendResult::NOT_CONNECTED;
|
||||
}
|
||||
|
||||
char *text = static_cast<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();
|
||||
}
|
||||
|
|
|
|||
115
src/Alex2ESP.h
115
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 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
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue