Queue up to eight directives for loop(); drop the repeat of a directive while it arrives
8ea7778kept one directive for loop() and recognised the repeat that the broker mirror delivers in two ways: bytes compared with the directive that still waited, and the messageId once loop() had parsed it. A repeat that arrived after its directive had been read therefore took the one place until the next loop(), and a different directive behind it was dropped. Alexa sends a group command as one directive per endpoint, so a board with several endpoints lost directives whenever they arrived faster than loop() ran; 1.1.0 queued five tokens. AlexaDirectiveBuffer is now a queue. The arriving message is collected in a heap block of its own and hashed (FNV-1a, 64 bit) as its fragments come in. At the last fragment it is one of three things: a repeat, when the hash is among the last 16 that were queued, and then it is freed and never takes a place; a directive, which is queued; or lost, when eight directives (ALEX2ESP_MAX_QUEUED_DIRECTIVES) or 8188 bytes (ALEX2ESP_MAX_QUEUED_BYTES, four times the largest directive) already wait. Whether there is a place is decided at the last fragment, because loop() may have read a directive by then. A message for which there is no memory is still hashed, so its loss is reported only when it is not a repeat. loop() parses one directive per call, in the order of arrival, and frees its block before the handler runs. The messageId check stays for a directive that comes again with other bytes; it remembers 16 ids instead of 4 and shares the ring with the buffer (AlexaRecentHashes). The two limits are in src/AlexaLimits.h, with #error for values that cannot work. The constructor of the buffer takes the allocator, malloc by default, so the tests can let it fail. For a sketch: nothing to change. The error line of a directive that finds no place reads "directive of N bytes on <topic> dropped: 8 directives, M bytes, already wait for loop()". The heap holds up to 8188 bytes of waiting directives and one arriving directive of up to 2048, where it held one directive. Measured with the bridge built for the host against a fake MQTT client (not in the repository), directives of 793 bytes, every message delivered twice,8ea7778-> this commit: repeat of D1 and a new D2 after D1 was handled D2 dropped -> D2 handled group of 5 in one burst 1 of 5 handled, 8 error lines -> 5 of 5, none group of 8 in one burst 8 of 8 handled, no error line group of 10 in one burst, not mirrored 8 of 10 handled, 2 error lines group of 10, a loop() after every fourth message 10 of 10 handled, no error line Tests: 41 host tests (33 before). New: the order of the queue and the reuse of its places, the ninth directive, a place that becomes free while a directive arrives, the limit in bytes, the largest directive in an empty queue, the repeat of a waiting directive and of one that was read, a group of five with repeats, a repeat when no place is free, a directive that differs in one byte, how long a repeat is remembered, a directive and a repeat without memory, the FNV-1a test vectors. They also pass under -fsanitize=address,undefined. Seven faults planted in a copy of AlexaBridgeLogic.cpp (no repeat check, no limit in bytes, no limit in places, no bounds check, a lost directive remembered, release() that keeps the bytes, last in first out) were each noticed: six by failing tests, the missing bounds check by AddressSanitizer as a heap-buffer-overflow. Built for d1_mini with empty credentials (PlatformIO 6.2.0, espressif8266 4.2.1), static RAM / flash in bytes,e48f858-> this commit, no warnings: basicLight 34,116 / 336,757 -> 34,444 / 337,145 lightWithBrightness 34,232 / 340,517 -> 34,560 / 340,921 lightWithColorTemp 34,380 / 341,177 -> 34,708 / 341,581 tempSensor 34,024 / 335,457 -> 34,352 / 335,845 blindControl 34,256 / 338,925 -> 34,584 / 339,313 The 328 bytes of RAM are the two rings of 16 hashes (256) and the eight places of the queue. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
parent
e48f8580f9
commit
6fe8311a6f
8 changed files with 574 additions and 275 deletions
|
|
@ -49,7 +49,7 @@ Alex2ESP::Alex2ESP()
|
|||
discoveryAnnounced(0),
|
||||
discoveryStarted(0),
|
||||
discoveryLastAttempt(0),
|
||||
directive(ALEX2ESP_MAX_DIRECTIVE) {}
|
||||
directives(ALEX2ESP_MAX_DIRECTIVE, ALEX2ESP_MAX_QUEUED_BYTES) {}
|
||||
|
||||
void Alex2ESP::begin(const char *username, const char *password, const char *rootTopic)
|
||||
{
|
||||
|
|
@ -224,12 +224,12 @@ void Alex2ESP::onMqttDisconnect(AsyncMqttClientDisconnectReason reason)
|
|||
{
|
||||
state = Alex2ESPState::DISCONNECTED;
|
||||
disconnectReason = reason;
|
||||
directive.cancelArrival();
|
||||
directives.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.
|
||||
// parses and dispatches the directives, 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)
|
||||
|
|
@ -260,13 +260,13 @@ void Alex2ESP::onMessage(char *topic, char *payload, AsyncMqttClientMessagePrope
|
|||
return;
|
||||
}
|
||||
|
||||
switch (directive.append(payload, length, index, total))
|
||||
switch (directives.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);
|
||||
case AlexaDirectiveBuffer::Result::QUEUE_FULL:
|
||||
ALEX2ESP_LOGE("directive of %u bytes on %s dropped: %u directives, %u bytes, already wait for loop()", (unsigned)total, topic, (unsigned)directives.waitingDirectives(), (unsigned)directives.waitingBytes());
|
||||
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());
|
||||
|
|
@ -289,17 +289,17 @@ void Alex2ESP::onMessage(char *topic, char *payload, AsyncMqttClientMessagePrope
|
|||
|
||||
void Alex2ESP::processDirective()
|
||||
{
|
||||
if (!directive.ready())
|
||||
if (!directives.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();
|
||||
size_t length = directives.length();
|
||||
DeserializationError error = deserializeJson(message, directives.data(), length);
|
||||
// The document has its own copy of every string, so the text is not needed any more: releasing it before the
|
||||
// handler runs frees the heap for the report and the place in the queue for the next directive.
|
||||
directives.release();
|
||||
|
||||
if (error)
|
||||
{
|
||||
|
|
|
|||
|
|
@ -49,8 +49,8 @@ public:
|
|||
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.
|
||||
// Call from the sketch's loop(): connects, answers discovery requests and hands one directive per call to its
|
||||
// device. It does not block; the handlers of the sketch run inside it.
|
||||
void loop();
|
||||
|
||||
AsyncMqttClientDisconnectReason getDisconnectReason() const;
|
||||
|
|
@ -102,8 +102,8 @@ private:
|
|||
unsigned long discoveryStarted; // millis() when the Discover arrived
|
||||
unsigned long discoveryLastAttempt; // millis() of the publish that was refused
|
||||
|
||||
AlexaDirectiveBuffer directive; // The directive that waits for loop(), or is still arriving
|
||||
AlexaRecentIds recentIds; // messageIds of the last directives handled
|
||||
AlexaDirectiveBuffer directives; // The directives that wait for loop(), and the one that is 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);
|
||||
|
|
|
|||
|
|
@ -4,12 +4,40 @@
|
|||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
|
||||
AlexaDirectiveBuffer::AlexaDirectiveBuffer(size_t capacityBytes)
|
||||
: capacity(capacityBytes) {}
|
||||
bool AlexaRecentHashes::contains(uint64_t hash) const
|
||||
{
|
||||
for (uint8_t i = 0; i < count; i++)
|
||||
{
|
||||
if (hashes[i] == hash)
|
||||
{
|
||||
return true;
|
||||
}
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
void AlexaRecentHashes::remember(uint64_t hash)
|
||||
{
|
||||
hashes[next] = hash;
|
||||
next = (next + 1) % CAPACITY;
|
||||
if (count < CAPACITY)
|
||||
{
|
||||
count++;
|
||||
}
|
||||
}
|
||||
|
||||
AlexaDirectiveBuffer::AlexaDirectiveBuffer(size_t largestDirective, size_t mostBytesWaiting, Allocator allocator)
|
||||
: largestDirective(largestDirective),
|
||||
mostBytesWaiting(mostBytesWaiting < largestDirective ? largestDirective : mostBytesWaiting),
|
||||
allocator(allocator) {}
|
||||
|
||||
AlexaDirectiveBuffer::~AlexaDirectiveBuffer()
|
||||
{
|
||||
discard();
|
||||
discardArriving();
|
||||
while (ready())
|
||||
{
|
||||
release();
|
||||
}
|
||||
}
|
||||
|
||||
AlexaDirectiveBuffer::Result AlexaDirectiveBuffer::append(const char *data, size_t length, size_t index, size_t total)
|
||||
|
|
@ -21,117 +49,102 @@ AlexaDirectiveBuffer::Result AlexaDirectiveBuffer::append(const char *data, size
|
|||
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)
|
||||
// First fragment
|
||||
AlexaDirectiveBuffer::Result AlexaDirectiveBuffer::start(const char *data, size_t length, size_t messageLength)
|
||||
{
|
||||
// A message that was still arriving ends here: the connection dropped in the middle of it
|
||||
if (!waiting)
|
||||
{
|
||||
discard();
|
||||
}
|
||||
arrival = Arrival::NONE;
|
||||
discardArriving();
|
||||
arrival = Arrival::IDLE;
|
||||
|
||||
if (total == 0)
|
||||
if (messageLength == 0)
|
||||
{
|
||||
return Result::EMPTY;
|
||||
}
|
||||
if (total > capacity)
|
||||
if (messageLength > largestDirective)
|
||||
{
|
||||
return refuse(Result::TOO_LARGE);
|
||||
}
|
||||
if (data == nullptr || length > total)
|
||||
if (data == nullptr || length > messageLength)
|
||||
{
|
||||
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;
|
||||
total = messageLength;
|
||||
offset = 0;
|
||||
return store(data, length);
|
||||
hash = AlexaBridgeLogic::HASH_OF_NOTHING;
|
||||
// Whether the queue has a place is decided at the last fragment: loop() may have read a directive by then
|
||||
arriving = static_cast<char *>(allocator(messageLength + 1));
|
||||
arrival = (arriving != nullptr) ? Arrival::STORING : Arrival::HASHING;
|
||||
return take(data, length);
|
||||
}
|
||||
|
||||
// Any later fragment
|
||||
AlexaDirectiveBuffer::Result AlexaDirectiveBuffer::proceed(const char *data, size_t length, size_t index, size_t total)
|
||||
AlexaDirectiveBuffer::Result AlexaDirectiveBuffer::proceed(const char *data, size_t length, size_t index, size_t messageLength)
|
||||
{
|
||||
if (arrival == Arrival::SKIPPING)
|
||||
{
|
||||
return Result::IGNORED;
|
||||
}
|
||||
if (arrival == Arrival::NONE)
|
||||
if (arrival == Arrival::IDLE)
|
||||
{
|
||||
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 (data == nullptr || messageLength != total || index != offset || length > total - offset)
|
||||
{
|
||||
if (!waiting)
|
||||
{
|
||||
discard();
|
||||
}
|
||||
discardArriving();
|
||||
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);
|
||||
return take(data, length);
|
||||
}
|
||||
|
||||
AlexaDirectiveBuffer::Result AlexaDirectiveBuffer::store(const char *data, size_t length)
|
||||
AlexaDirectiveBuffer::Result AlexaDirectiveBuffer::take(const char *data, size_t length)
|
||||
{
|
||||
memcpy(text + offset, data, length);
|
||||
offset += length;
|
||||
if (offset < size)
|
||||
if (arrival == Arrival::STORING)
|
||||
{
|
||||
memcpy(arriving + offset, data, length);
|
||||
}
|
||||
hash = AlexaBridgeLogic::hashBytes(data, length, hash);
|
||||
offset += length;
|
||||
if (offset < total)
|
||||
{
|
||||
arrival = Arrival::STORING;
|
||||
return Result::INCOMPLETE;
|
||||
}
|
||||
text[size] = '\0';
|
||||
waiting = true;
|
||||
arrival = Arrival::NONE;
|
||||
return finish();
|
||||
}
|
||||
|
||||
// Last fragment: the message is a repeat, a directive for the queue, or lost
|
||||
AlexaDirectiveBuffer::Result AlexaDirectiveBuffer::finish()
|
||||
{
|
||||
const bool stored = (arrival == Arrival::STORING);
|
||||
arrival = Arrival::IDLE;
|
||||
|
||||
if (recent.contains(hash))
|
||||
{
|
||||
discardArriving();
|
||||
return Result::DUPLICATE;
|
||||
}
|
||||
|
||||
// A directive that is lost is not remembered: when its repeat finds memory and a place, the repeat is handled
|
||||
if (!stored)
|
||||
{
|
||||
return Result::NO_MEMORY;
|
||||
}
|
||||
if (count == CAPACITY || bytes + total > mostBytesWaiting)
|
||||
{
|
||||
discardArriving();
|
||||
return Result::QUEUE_FULL;
|
||||
}
|
||||
|
||||
arriving[total] = '\0';
|
||||
Entry &place = waiting[(first + count) % CAPACITY];
|
||||
place.text = arriving;
|
||||
place.length = total;
|
||||
arriving = nullptr;
|
||||
count++;
|
||||
bytes += total;
|
||||
recent.remember(hash);
|
||||
return Result::COMPLETE;
|
||||
}
|
||||
|
||||
|
|
@ -143,33 +156,26 @@ AlexaDirectiveBuffer::Result AlexaDirectiveBuffer::refuse(Result reason)
|
|||
|
||||
void AlexaDirectiveBuffer::release()
|
||||
{
|
||||
if (!waiting)
|
||||
if (count == 0)
|
||||
{
|
||||
return;
|
||||
}
|
||||
waiting = false;
|
||||
// A repeat that is still being compared needs the text until its last fragment
|
||||
if (arrival != Arrival::COMPARING)
|
||||
{
|
||||
discard();
|
||||
}
|
||||
free(waiting[first].text);
|
||||
bytes -= waiting[first].length;
|
||||
first = (first + 1) % CAPACITY;
|
||||
count--;
|
||||
}
|
||||
|
||||
void AlexaDirectiveBuffer::cancelArrival()
|
||||
{
|
||||
arrival = Arrival::NONE;
|
||||
if (!waiting)
|
||||
{
|
||||
discard();
|
||||
}
|
||||
discardArriving();
|
||||
arrival = Arrival::IDLE;
|
||||
}
|
||||
|
||||
void AlexaDirectiveBuffer::discard()
|
||||
void AlexaDirectiveBuffer::discardArriving()
|
||||
{
|
||||
free(text);
|
||||
text = nullptr;
|
||||
size = 0;
|
||||
offset = 0;
|
||||
free(arriving);
|
||||
arriving = nullptr;
|
||||
}
|
||||
|
||||
bool AlexaRecentIds::seenBefore(const char *id)
|
||||
|
|
@ -179,28 +185,23 @@ bool AlexaRecentIds::seenBefore(const char *id)
|
|||
return false;
|
||||
}
|
||||
|
||||
// FNV-1a, 64 bit
|
||||
uint64_t hash = 0xcbf29ce484222325ULL;
|
||||
for (const char *c = id; *c != '\0'; c++)
|
||||
uint64_t hash = AlexaBridgeLogic::hashBytes(id, strlen(id));
|
||||
if (recent.contains(hash))
|
||||
{
|
||||
hash ^= static_cast<uint8_t>(*c);
|
||||
return true;
|
||||
}
|
||||
recent.remember(hash);
|
||||
return false;
|
||||
}
|
||||
|
||||
uint64_t AlexaBridgeLogic::hashBytes(const char *data, size_t length, uint64_t hash)
|
||||
{
|
||||
for (size_t i = 0; i < length; i++)
|
||||
{
|
||||
hash ^= static_cast<uint8_t>(data[i]);
|
||||
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;
|
||||
return hash;
|
||||
}
|
||||
|
||||
bool AlexaBridgeLogic::clockIsSet(time_t now)
|
||||
|
|
|
|||
|
|
@ -1,23 +1,51 @@
|
|||
// 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.
|
||||
// over, queueing directives for loop(), 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 <stdlib.h>
|
||||
#include <time.h>
|
||||
#include <ArduinoJson.h>
|
||||
#include "AlexaLimits.h"
|
||||
#include "AlexaTransport.h"
|
||||
|
||||
// One directive on its way from the MQTT callback to loop().
|
||||
// The last values it was given, as many as CAPACITY: what a repeat is recognised by.
|
||||
class AlexaRecentHashes
|
||||
{
|
||||
public:
|
||||
// Twice the directives that may wait for loop(): the repeat of a directive is still recognised when a full
|
||||
// queue of other directives has been read since.
|
||||
static const uint8_t CAPACITY = 2 * ALEX2ESP_MAX_QUEUED_DIRECTIVES;
|
||||
|
||||
bool contains(uint64_t hash) const;
|
||||
|
||||
// Once CAPACITY values are held the new one takes the place of the oldest
|
||||
void remember(uint64_t hash);
|
||||
|
||||
private:
|
||||
uint64_t hashes[CAPACITY] = {};
|
||||
uint8_t count = 0;
|
||||
uint8_t next = 0;
|
||||
};
|
||||
|
||||
// The directives on their 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.
|
||||
// one heap block of total + 1 bytes, allocated when the first fragment arrives. The complete directive waits in a
|
||||
// queue until loop() has read it and calls release(), which frees the block. A queue, because several directives
|
||||
// arrive before loop() runs again when Alexa switches a group or the sketch is busy.
|
||||
//
|
||||
// A broker that mirrors its topics to a second broker and back delivers every message twice (the Alex2MQTT broker
|
||||
// does). The bytes of a message are hashed while they arrive, and a message with the hash of one that was queued
|
||||
// lately is a repeat: it is discarded at its last fragment and never takes a place in the queue.
|
||||
//
|
||||
// The heap this takes is bounded: the waiting directives by their number and by their bytes, the arriving one by
|
||||
// the limit for a directive. Nothing is ever written outside a block.
|
||||
class AlexaDirectiveBuffer
|
||||
{
|
||||
public:
|
||||
|
|
@ -25,80 +53,118 @@ public:
|
|||
{
|
||||
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
|
||||
DUPLICATE, // last fragment of a repeat: byte for byte a directive that was queued lately
|
||||
IGNORED, // fragment of a message that was refused at an earlier 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
|
||||
TOO_LARGE, // total is over the limit for a directive
|
||||
QUEUE_FULL, // last fragment of a directive that found no place: as many directives or bytes wait as may
|
||||
NO_MEMORY, // last fragment of a directive that found no heap for its 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);
|
||||
// Directives that may wait for loop()
|
||||
static const uint8_t CAPACITY = ALEX2ESP_MAX_QUEUED_DIRECTIVES;
|
||||
|
||||
// Returns a block that free() takes, or nullptr
|
||||
typedef void *(*Allocator)(size_t size);
|
||||
|
||||
// largestDirective: the largest message that is accepted (ALEX2ESP_MAX_DIRECTIVE in the bridge).
|
||||
// mostBytesWaiting: what the waiting directives may take together (ALEX2ESP_MAX_QUEUED_BYTES); at least
|
||||
// largestDirective, so that the empty queue takes any directive that is accepted.
|
||||
// allocator: malloc, unless a test wants it to fail.
|
||||
AlexaDirectiveBuffer(size_t largestDirective, size_t mostBytesWaiting, Allocator allocator = malloc);
|
||||
~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.
|
||||
// One call per fragment, in the order of arrival. EMPTY, TOO_LARGE, OUT_OF_ORDER, QUEUE_FULL and NO_MEMORY
|
||||
// report a message that is lost, each once per message. The first three are returned by the fragment that
|
||||
// shows the fault, and the fragments after it are IGNORED. QUEUE_FULL and NO_MEMORY are returned by the last
|
||||
// fragment: a place may become free while the message arrives, and a message that has no place or no memory
|
||||
// may still turn out to be a repeat, which is no loss.
|
||||
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 directive that has waited longest: data() is its text (NUL-terminated), length() its size without the NUL
|
||||
bool ready() const { return count > 0; }
|
||||
const char *data() const { return count > 0 ? waiting[first].text : nullptr; }
|
||||
size_t length() const { return count > 0 ? waiting[first].length : 0; }
|
||||
|
||||
// The waiting directive has been read: the buffer takes the next one
|
||||
// That directive has been read: its block is freed, the next one is up
|
||||
void release();
|
||||
|
||||
// The connection is gone: a message that was still arriving will not be completed. A directive that waits stays.
|
||||
// The connection is gone: a message that was still arriving will not be completed. The directives that wait stay.
|
||||
void cancelArrival();
|
||||
|
||||
// For the log: how many directives wait, and their bytes
|
||||
size_t waitingDirectives() const { return count; }
|
||||
size_t waitingBytes() const { return bytes; }
|
||||
|
||||
private:
|
||||
struct Entry
|
||||
{
|
||||
char *text;
|
||||
size_t length;
|
||||
};
|
||||
|
||||
// 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
|
||||
IDLE, // no message is arriving
|
||||
STORING, // they are copied into the block and hashed
|
||||
HASHING, // there was no memory for the block: they are hashed, which tells a repeat from a loss
|
||||
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 start(const char *data, size_t length, size_t messageLength);
|
||||
Result proceed(const char *data, size_t length, size_t index, size_t messageLength);
|
||||
Result take(const char *data, size_t length);
|
||||
Result finish();
|
||||
Result refuse(Result reason);
|
||||
void discard();
|
||||
void discardArriving();
|
||||
|
||||
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;
|
||||
const size_t largestDirective;
|
||||
const size_t mostBytesWaiting;
|
||||
const Allocator allocator;
|
||||
|
||||
Entry waiting[CAPACITY]; // a ring: the oldest directive at first, count of them
|
||||
uint8_t first = 0;
|
||||
uint8_t count = 0;
|
||||
size_t bytes = 0; // of the waiting directives, without their NULs
|
||||
|
||||
char *arriving = nullptr; // block of the message that is arriving, while it is STORING
|
||||
size_t total = 0; // length of that message
|
||||
size_t offset = 0; // bytes of it taken so far
|
||||
uint64_t hash = 0; // of those bytes
|
||||
Arrival arrival = Arrival::IDLE;
|
||||
|
||||
AlexaRecentHashes recent; // the directives queued lately, by the hash of their text
|
||||
};
|
||||
|
||||
// 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.
|
||||
// The messageIds of the last directives that were dispatched. The repeat a mirroring broker delivers is discarded
|
||||
// by AlexaDirectiveBuffer when it is the same bytes; this recognises the directive that comes again with other
|
||||
// bytes, so that a sketch never acts on a messageId twice. The ids are kept as hashes: 8 bytes each instead of the
|
||||
// 37 a UUID string takes.
|
||||
class AlexaRecentIds
|
||||
{
|
||||
public:
|
||||
static const uint8_t CAPACITY = 4;
|
||||
static const uint8_t CAPACITY = AlexaRecentHashes::CAPACITY;
|
||||
|
||||
// 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;
|
||||
AlexaRecentHashes recent;
|
||||
};
|
||||
|
||||
namespace AlexaBridgeLogic
|
||||
{
|
||||
// FNV-1a, 64 bit. Two different texts have the same hash with a probability of 2^-64.
|
||||
const uint64_t HASH_OF_NOTHING = 0xcbf29ce484222325ULL;
|
||||
|
||||
// The hash of `length` bytes at `data`, or, given the hash of what came before them, of both together
|
||||
uint64_t hashBytes(const char *data, size_t length, uint64_t hash = HASH_OF_NOTHING);
|
||||
|
||||
// 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;
|
||||
|
|
|
|||
|
|
@ -19,4 +19,25 @@
|
|||
#define ALEX2ESP_MAX_MESSAGE (ALEX2ESP_MAX_DIRECTIVE + 1024)
|
||||
#endif
|
||||
|
||||
// How many directives may wait for loop(), and how many bytes they may take on the heap together. Alexa sends a
|
||||
// group command ("turn off the kitchen") as one directive per endpoint, all within a second, so a board with
|
||||
// several endpoints receives several directives while its sketch is still busy with the first. The defaults hold
|
||||
// eight directives of up to 1023 bytes each. One more directive, of up to ALEX2ESP_MAX_DIRECTIVE bytes, is on the
|
||||
// heap while it arrives.
|
||||
#ifndef ALEX2ESP_MAX_QUEUED_DIRECTIVES
|
||||
#define ALEX2ESP_MAX_QUEUED_DIRECTIVES 8
|
||||
#endif
|
||||
|
||||
#ifndef ALEX2ESP_MAX_QUEUED_BYTES
|
||||
#define ALEX2ESP_MAX_QUEUED_BYTES (4 * ALEX2ESP_MAX_DIRECTIVE)
|
||||
#endif
|
||||
|
||||
#if ALEX2ESP_MAX_QUEUED_DIRECTIVES < 1 || ALEX2ESP_MAX_QUEUED_DIRECTIVES > 64
|
||||
#error "ALEX2ESP_MAX_QUEUED_DIRECTIVES has to be between 1 and 64"
|
||||
#endif
|
||||
|
||||
#if ALEX2ESP_MAX_QUEUED_BYTES < ALEX2ESP_MAX_DIRECTIVE
|
||||
#error "ALEX2ESP_MAX_QUEUED_BYTES has to be ALEX2ESP_MAX_DIRECTIVE or more: the largest directive has to fit the empty queue"
|
||||
#endif
|
||||
|
||||
#endif // ALEXA_LIMITS_H
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue