diff --git a/keywords.txt b/keywords.txt index f86f130..d483099 100644 --- a/keywords.txt +++ b/keywords.txt @@ -89,6 +89,8 @@ printMemoryInfo KEYWORD2 ALEX2ESP_MAX_DIRECTIVE LITERAL1 ALEX2ESP_MAX_MESSAGE LITERAL1 +ALEX2ESP_MAX_QUEUED_DIRECTIVES LITERAL1 +ALEX2ESP_MAX_QUEUED_BYTES LITERAL1 ALEX2ESP_LOG_MAX LITERAL1 ALEXA_TIMESTAMP_SIZE LITERAL1 MAX_EVENTS LITERAL1 diff --git a/readme.md b/readme.md index 6d09567..beb82b9 100644 --- a/readme.md +++ b/readme.md @@ -139,9 +139,9 @@ Everything goes over MQTT (port 1883 of `alex2mqtt.stormysdream.club`); the libr - **Session.** `begin()` starts SNTP (`pool.ntp.org`, `time.nist.gov`) and returns; `loop()` opens the MQTT session once the clock is set, or after 5 s without an answer, and subscribes to `/discover` and `/+/alexaDirective`. `getState()` is `CONNECTED` when the broker has acknowledged both subscriptions. A sketch that sets the clock itself (its own `configTime()` with a time zone, an RTC) calls `alexClient.setTimeSource(false)` before `begin()`. - **Discovery.** On `/discover` the library answers with one discovery object per device on `/discover_r`. The backend accepts one endpoint object per message and collects everything that arrives within 1 s for Alexa's discovery answer (up to 5 s for its proactive AddOrUpdate push), so all devices are published back to back from the next `loop()`. Each object sits on the heap (about 1 KB) until the MQTT client has sent it; when the client cannot take another one (free heap under 4 KB), the library prints `[Alex2ESP] discovery deferred at ` and sends the rest from `loop()` as the queue drains, for up to 5 s after the request. `[Alex2ESP] error: discovery gave up: N device(s) not announced` means those devices missed this answer - on the backend's proactive discovery that can remove them from Alexa until the next one. -- **Directives.** The directive arrives as JSON on `//alexaDirective`. A directive larger than one TCP segment arrives in fragments, which are put together in one heap block that exists only until `loop()` has parsed it. `loop()` then fires `ReportState` or `Event` (and `DirectiveReceived`, if registered) with the directive: `directive["header"]`, `directive["endpoint"]`, `directive["payload"]`. One directive is handled at a time; a directive that arrives twice (a broker that mirrors its topics delivers every message twice) is handled once. +- **Directives.** The directive arrives as JSON on `//alexaDirective`. A directive larger than one TCP segment arrives in fragments, which are put together in one heap block that exists only until `loop()` has parsed it. `loop()` then fires `ReportState` or `Event` (and `DirectiveReceived`, if registered) with the directive: `directive["header"]`, `directive["endpoint"]`, `directive["payload"]`. Every call of `loop()` handles one directive, in the order of arrival. Up to eight directives wait for it: Alexa sends a group command ("turn off the kitchen") as one directive per endpoint, and they arrive faster than a busy sketch calls `loop()`. A directive that arrives twice (a broker that mirrors its topics delivers every message twice) is handled once; the repeat is recognised while it arrives and takes no place among the waiting ones. - **Reports.** `send()` publishes the report on `//alexaResponce` at once. The backend waits 7 s for it, so answer from the event handler. Every property carries the board's UTC time as `timeOfSample`. -- **Limits.** A directive may be 2047 bytes (`ALEX2ESP_MAX_DIRECTIVE`), a report or the discovery object of one device 3071 (`ALEX2ESP_MAX_MESSAGE`). The second limit is the first plus 1024 unless it is set: an answer repeats the `correlationToken` of its directive, which is most of a large directive, and adds 140 to 170 bytes per property, so the largest directive can be answered with six properties. `-DALEX2ESP_MAX_DIRECTIVE=` and `-DALEX2ESP_MAX_MESSAGE=` in `build_flags` change the limits. `send()` returns `false` when the report was not sent: no session with the broker, the MQTT client or the heap cannot take it, or it is too large. Nothing is ever sent truncated, and nothing is dropped without a line on Serial. +- **Limits.** A directive may be 2047 bytes (`ALEX2ESP_MAX_DIRECTIVE`), a report or the discovery object of one device 3071 (`ALEX2ESP_MAX_MESSAGE`). The second limit is the first plus 1024 unless it is set: an answer repeats the `correlationToken` of its directive, which is most of a large directive, and adds 140 to 170 bytes per property, so the largest directive can be answered with six properties. Eight directives (`ALEX2ESP_MAX_QUEUED_DIRECTIVES`) of 8188 bytes together (`ALEX2ESP_MAX_QUEUED_BYTES`, four times the largest directive) may wait for `loop()`; one that finds no place is dropped with `[Alex2ESP] error: directive of N bytes on dropped: ...`. `-D=` in `build_flags` changes a limit. `send()` returns `false` when the report was not sent: no session with the broker, the MQTT client or the heap cannot take it, or it is too large. Nothing is ever sent truncated, and nothing is dropped without a line on Serial. - **Serial output.** Every line of the library starts with `[Alex2ESP]`, a problem with `[Alex2ESP] error:`. `alexClient.setLogLevel(AlexaLogLevel::ERROR)` leaves only the problems, `AlexaLogLevel::NONE` nothing; the default, `AlexaLogLevel::INFO`, adds the session, discovery and one line per directive (`[Alex2ESP] ESP-01 <- Alexa.PowerController.TurnOn`). `AlexaLogLevel::DEBUG` (sizes and free heap per message) has to be compiled in with `-DALEX2ESP_LOG_MAX=3`; `-DALEX2ESP_LOG_MAX=0` compiles every line out. Credentials and correlation tokens are never printed. A sketch that defines a macro named `DEBUG`, `ERROR` or `INFO` cannot write the level of that name; it passes the number instead, for example `alexClient.setLogLevel(static_cast(3))` for `DEBUG`. Boards that run 1.1.0 or older keep working: the backend still publishes the token on `//alexaDirective_e` and serves the HTTP routes they use. @@ -245,8 +245,8 @@ Behaviour changes: - Directives arrive over MQTT. The library subscribes to `/+/alexaDirective`, where Alex2MQTT has always published the whole directive, instead of fetching it over HTTP with the token from `//alexaDirective_e`. The HTTP detour dates from 2024, when a directive larger than one TCP segment reached the MQTT callback in pieces; the pieces are now put together by their offset and the total length, in one heap block that lives until `loop()` has parsed the directive. `loop()` no longer stalls for two HTTP round trips per directive. - Reports leave over MQTT. `send()` publishes on `//alexaResponce` at once, where 1.1.0 queued the report for an HTTP POST from a later `loop()`. It returns `false` when there is no session with the broker, when the MQTT client or the heap cannot take the report, or when the report is over 3071 bytes (`ALEX2ESP_MAX_MESSAGE`, by default 1024 more than the largest directive, so that every directive that is accepted can be answered); each case prints its reason. The 5-slot send queue is gone. - `timeOfSample` is the board's own time in UTC, for example `2026-09-28T13:05:09Z`. 1.1.0 sent the placeholder `{REPLACE_WITH_DATETIME}`, which only the backend's HTTP route replaced. `begin()` starts SNTP (`pool.ntp.org`, `time.nist.gov`) and no longer connects itself: `loop()` opens the MQTT session once the clock is set, or after 5 s without an answer, so the session comes up a few seconds later than before. `alexClient.setTimeSource(false)` before `begin()` leaves the clock to the sketch. `AddContextProp()` fills `timeOfSample` in when the property has none or carries the old placeholder. -- A directive is parsed and handed to the sketch from `loop()`, one at a time. One that arrives while the previous one still waits for `loop()` is dropped, as is one over 2047 bytes (`ALEX2ESP_MAX_DIRECTIVE`); both print an error. -- A directive that arrives twice is handled once: the Alex2MQTT broker currently delivers every message twice through a mirror. The `messageId`s of the last four directives are remembered; the repeat prints `repeated directive ... ignored`. +- A directive is parsed and handed to the sketch from `loop()`, one per call, in the order of arrival. Up to eight directives of 8188 bytes together wait for it on the heap (`ALEX2ESP_MAX_QUEUED_DIRECTIVES`, `ALEX2ESP_MAX_QUEUED_BYTES`; 1.1.0 queued five tokens), because Alexa sends a group command as one directive per endpoint. A directive that finds no place is dropped, as is one over 2047 bytes (`ALEX2ESP_MAX_DIRECTIVE`); both print an error. +- A directive that arrives twice is handled once: the Alex2MQTT broker currently delivers every message twice through a mirror. The repeat is recognised while it arrives, by a hash of its bytes, and takes no place among the waiting directives; a directive that comes again with other bytes is recognised by its `messageId`. The last 16 directives are remembered; the repeat prints `repeated directive ... ignored`. - The subscription delivers the directives of every endpoint of the account. Those for endpoints of another board are recognised by their topic and neither buffered nor parsed. - Discovery is answered from `loop()`, not inside the MQTT callback. A discovery object over 3071 bytes (`ALEX2ESP_MAX_MESSAGE`) is refused with an error; the other devices are still announced. - `getState()` stays `INITIALIZED` until the first connect, and becomes `CONNECTED` when the broker has acknowledged both subscriptions (1.1.0: the first of them). A refused subscription prints an error. @@ -255,9 +255,9 @@ Behaviour changes: - New: `Alex2ESP::setLogLevel()`, `Alex2ESP::setTimeSource()`, `AlexaDevice::hasEndpointId()`, `AlexaLog`, `AlexaSendResult`. - Removed: the queues and buffers of `AlexaUtils` (`enqueue`, `dequeue`, `dequeueVals`, `enqueueReceive`, `dequeueReceive`, `isQueueEmpty`, `isQueueFull`, `isReceiveQueueEmpty`, `isReceiveQueueFull`, `receivePayload`, `nextMessageId`) and its `log`/`logln`, which printed nothing unless the library was edited; `AlexaUtils::printMemoryInfo()` stays. `MAX_STATUS_REPORT_SIZE` (the limit is `ALEX2ESP_MAX_MESSAGE`). The library no longer includes `ESP8266HTTPClient`. -Memory: `examples/basicLight.cpp` for a D1 mini takes 34,116 bytes of static RAM (1.1.0: 52,768) and 336,757 bytes of flash (1.1.0: 350,885), as PlatformIO reports them (espressif8266 4.2.1, Arduino core 3.1.2). The static RAM was the five 2 KB queue slots, three more 2 KB buffers and the two HTTP clients. SNTP and the time stamp are 1.8 KB of the flash figure. +Memory: `examples/basicLight.cpp` for a D1 mini takes 34,444 bytes of static RAM (1.1.0: 52,768) and 337,145 bytes of flash (1.1.0: 350,885), as PlatformIO reports them (espressif8266 4.2.1, Arduino core 3.1.2). The static RAM was the five 2 KB queue slots, three more 2 KB buffers and the two HTTP clients. SNTP and the time stamp are 1.8 KB of the flash figure. -Tests: `pio test -e native` in the repository runs 33 host tests of the receive and publish logic (reassembly of fragments, the two size limits, repeated directives, topics, time stamps). No board is needed. +Tests: `pio test -e native` in the repository runs 41 host tests of the receive and publish logic (reassembly of fragments, the directives that wait for `loop()`, repeated directives, the size limits, a heap without room, topics, time stamps). No board is needed. Boards that run 1.1.0 are not affected: the backend keeps the token topic and the HTTP routes. diff --git a/src/Alex2ESP.cpp b/src/Alex2ESP.cpp index fb164cb..5041aa5 100644 --- a/src/Alex2ESP.cpp +++ b/src/Alex2ESP.cpp @@ -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) { diff --git a/src/Alex2ESP.h b/src/Alex2ESP.h index 84cb631..c101921 100644 --- a/src/Alex2ESP.h +++ b/src/Alex2ESP.h @@ -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); diff --git a/src/AlexaBridgeLogic.cpp b/src/AlexaBridgeLogic.cpp index 5aebf1f..86188ed 100644 --- a/src/AlexaBridgeLogic.cpp +++ b/src/AlexaBridgeLogic.cpp @@ -4,12 +4,40 @@ #include #include -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(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(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(*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(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) diff --git a/src/AlexaBridgeLogic.h b/src/AlexaBridgeLogic.h index de310a0..0981f95 100644 --- a/src/AlexaBridgeLogic.h +++ b/src/AlexaBridgeLogic.h @@ -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 #include +#include #include #include +#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; diff --git a/src/AlexaLimits.h b/src/AlexaLimits.h index af9c9b7..c6fe5ae 100644 --- a/src/AlexaLimits.h +++ b/src/AlexaLimits.h @@ -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 diff --git a/test/test_bridge_logic/test_main.cpp b/test/test_bridge_logic/test_main.cpp index 00af990..6bbdb07 100644 --- a/test/test_bridge_logic/test_main.cpp +++ b/test/test_bridge_logic/test_main.cpp @@ -1,8 +1,10 @@ // Host tests of the bridge logic (src/AlexaBridgeLogic.cpp): what the MQTT callback decides about the fragments a -// broker delivers, which directives are repeats, what may be published, topics and time stamps. +// broker delivers, how directives wait for loop(), which of them are repeats, what may be published, topics and +// time stamps. // pio test -e native #include #include +#include #include #include #include @@ -26,6 +28,19 @@ static std::string directive(const char *name, const char *messageId, static const char ID_1[] = "1bd5d003-31b9-476f-ad03-71d471922820"; static const char ID_2[] = "7c0e1f6a-52d4-4b8e-9a3c-0d9f4e2b6a11"; +// The limits the bridge passes to its buffer by default +static const size_t LARGEST = 2047; +static const size_t MOST_BYTES = 4 * 2047; +static const size_t PLACES = AlexaDirectiveBuffer::CAPACITY; + +// The n-th directive of a series: each has a messageId of its own, as each directive of Alexa has +static std::string numbered(unsigned n) +{ + char messageId[37]; + snprintf(messageId, sizeof(messageId), "00000000-0000-4000-8000-%012u", n); + return directive(n % 2 == 0 ? "TurnOn" : "TurnOff", messageId); +} + // Hands the message over as the MQTT client does: pieces of fragmentSize bytes with their offset and the total. // Returns the result of the last piece; every piece before it has to be INCOMPLETE. static Result deliver(AlexaDirectiveBuffer &buffer, const std::string &message, size_t fragmentSize) @@ -54,7 +69,7 @@ void tearDown() {} void test_directive_in_one_piece_is_complete() { - AlexaDirectiveBuffer buffer(2047); + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES); std::string turnOn = directive("TurnOn", ID_1); ASSERT_RESULT(Result::COMPLETE, buffer.append(turnOn.data(), turnOn.size(), 0, turnOn.size())); @@ -73,12 +88,12 @@ void test_directive_in_fragments_is_reassembled_and_parses() // Every fragment size from one byte per fragment to the message in two halves for (size_t fragmentSize = 1; fragmentSize < turnOn.size(); fragmentSize++) { - AlexaDirectiveBuffer buffer(2047); + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES); ASSERT_RESULT(Result::COMPLETE, deliver(buffer, turnOn, fragmentSize)); assertWaiting(buffer, turnOn); } - AlexaDirectiveBuffer buffer(2047); + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES); deliver(buffer, turnOn, 100); JsonDocument parsed; TEST_ASSERT_TRUE(deserializeJson(parsed, buffer.data(), buffer.length()) == DeserializationError::Ok); @@ -88,7 +103,7 @@ void test_directive_in_fragments_is_reassembled_and_parses() void test_nothing_waits_before_the_last_fragment() { - AlexaDirectiveBuffer buffer(2047); + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES); std::string turnOn = directive("TurnOn", ID_1); ASSERT_RESULT(Result::INCOMPLETE, buffer.append(turnOn.data(), 50, 0, turnOn.size())); @@ -100,7 +115,7 @@ void test_nothing_waits_before_the_last_fragment() void test_directive_at_the_limit_is_accepted() { - AlexaDirectiveBuffer buffer(2047); + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES); std::string atLimit(2047, 'x'); ASSERT_RESULT(Result::COMPLETE, deliver(buffer, atLimit, 536)); @@ -109,7 +124,7 @@ void test_directive_at_the_limit_is_accepted() void test_directive_over_the_limit_is_refused_once() { - AlexaDirectiveBuffer buffer(2047); + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES); std::string tooLarge(2048, 'x'); ASSERT_RESULT(Result::TOO_LARGE, buffer.append(tooLarge.data(), 536, 0, tooLarge.size())); @@ -121,7 +136,7 @@ void test_directive_over_the_limit_is_refused_once() void test_refused_directive_does_not_block_the_next_one() { - AlexaDirectiveBuffer buffer(2047); + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES); std::string tooLarge(5000, 'x'); std::string turnOn = directive("TurnOn", ID_1); @@ -134,110 +149,271 @@ void test_refused_directive_does_not_block_the_next_one() void test_empty_message_is_reported() { - AlexaDirectiveBuffer buffer(2047); + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES); // AsyncMqttClient hands a message without a payload over as (nullptr, 0, 0, 0) ASSERT_RESULT(Result::EMPTY, buffer.append(nullptr, 0, 0, 0)); TEST_ASSERT_FALSE(buffer.ready()); } -// --- one directive at a time --- +// --- directives wait for loop() --- -void test_second_directive_before_loop_is_dropped_and_the_first_kept() +void test_directives_wait_in_the_order_they_arrived() { - AlexaDirectiveBuffer buffer(2047); - std::string turnOn = directive("TurnOn", ID_1); - std::string turnOff = directive("TurnOff", ID_2); + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES); + size_t bytes = 0; - ASSERT_RESULT(Result::COMPLETE, deliver(buffer, turnOn, 100)); - ASSERT_RESULT(Result::BUSY, buffer.append(turnOff.data(), 100, 0, turnOff.size())); - ASSERT_RESULT(Result::IGNORED, buffer.append(turnOff.data() + 100, turnOff.size() - 100, 100, turnOff.size())); - assertWaiting(buffer, turnOn); + for (unsigned n = 0; n < PLACES; n++) + { + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, numbered(n), 100)); + bytes += numbered(n).size(); + TEST_ASSERT_EQUAL_UINT(n + 1, buffer.waitingDirectives()); + TEST_ASSERT_EQUAL_UINT(bytes, buffer.waitingBytes()); + } + for (unsigned n = 0; n < PLACES; n++) + { + assertWaiting(buffer, numbered(n)); + buffer.release(); + } + TEST_ASSERT_FALSE(buffer.ready()); + TEST_ASSERT_EQUAL_UINT(0, buffer.waitingDirectives()); + TEST_ASSERT_EQUAL_UINT(0, buffer.waitingBytes()); - // Once loop() has read the first, the buffer takes directives again - buffer.release(); - ASSERT_RESULT(Result::COMPLETE, deliver(buffer, turnOff, 100)); - assertWaiting(buffer, turnOff); + // The places are a ring: they are used again, in the same order + for (unsigned round = 1; round <= 3; round++) + { + for (unsigned n = 0; n < PLACES - 1; n++) + { + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, numbered(100 * round + n), 100)); + } + for (unsigned n = 0; n < PLACES - 1; n++) + { + assertWaiting(buffer, numbered(100 * round + n)); + buffer.release(); + } + } + TEST_ASSERT_FALSE(buffer.ready()); } -void test_second_directive_of_the_same_size_is_dropped() +void test_directive_beyond_the_places_is_dropped_and_the_others_kept() { - AlexaDirectiveBuffer buffer(2047); - std::string first = directive("TurnOn", ID_1); - std::string second = directive("TurnOn", ID_2); // same length, another messageId - TEST_ASSERT_EQUAL_UINT(first.size(), second.size()); + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES); + for (unsigned n = 0; n < PLACES; n++) + { + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, numbered(n), 100)); + } - ASSERT_RESULT(Result::COMPLETE, deliver(buffer, first, 100)); - ASSERT_RESULT(Result::BUSY, buffer.append(second.data(), second.size(), 0, second.size())); - assertWaiting(buffer, first); + // The loss is reported once, by the last fragment + ASSERT_RESULT(Result::QUEUE_FULL, deliver(buffer, numbered(PLACES), 100)); + TEST_ASSERT_EQUAL_UINT(PLACES, buffer.waitingDirectives()); + assertWaiting(buffer, numbered(0)); + + // A directive that was dropped is not taken for a repeat when it comes again and a place is free + buffer.release(); + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, numbered(PLACES), 100)); + for (unsigned n = 1; n <= PLACES; n++) + { + assertWaiting(buffer, numbered(n)); + buffer.release(); + } + TEST_ASSERT_FALSE(buffer.ready()); +} + +void test_place_that_becomes_free_while_a_directive_arrives_is_used() +{ + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES); + for (unsigned n = 0; n < PLACES; n++) + { + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, numbered(n), 100)); + } + std::string late = numbered(PLACES); + + ASSERT_RESULT(Result::INCOMPLETE, buffer.append(late.data(), 100, 0, late.size())); + buffer.release(); // loop() ran between two fragments + ASSERT_RESULT(Result::INCOMPLETE, buffer.append(late.data() + 100, 100, 100, late.size())); + ASSERT_RESULT(Result::COMPLETE, buffer.append(late.data() + 200, late.size() - 200, 200, late.size())); + TEST_ASSERT_EQUAL_UINT(PLACES, buffer.waitingDirectives()); +} + +void test_waiting_directives_are_limited_in_bytes() +{ + AlexaDirectiveBuffer buffer(LARGEST, 3000); + std::string first(1500, 'a'); + std::string second(1500, 'b'); + std::string third(1, 'c'); + + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, first, 536)); + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, second, 536)); // 3000 bytes wait: at the limit + ASSERT_RESULT(Result::QUEUE_FULL, deliver(buffer, third, 536)); + TEST_ASSERT_EQUAL_UINT(2, buffer.waitingDirectives()); + TEST_ASSERT_EQUAL_UINT(3000, buffer.waitingBytes()); + + buffer.release(); + TEST_ASSERT_EQUAL_UINT(1500, buffer.waitingBytes()); + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, third, 536)); + assertWaiting(buffer, second); +} + +void test_empty_queue_takes_the_largest_directive() +{ + // A limit for the waiting bytes under the limit for one directive is raised to it + AlexaDirectiveBuffer buffer(LARGEST, 100); + std::string largest(LARGEST, 'x'); + + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, largest, 536)); + assertWaiting(buffer, largest); + ASSERT_RESULT(Result::QUEUE_FULL, deliver(buffer, std::string("y"), 536)); } // --- the repeat a mirroring broker delivers --- -void test_repeat_of_the_waiting_directive_is_a_duplicate() +void test_repeat_of_a_waiting_directive_is_discarded() { - AlexaDirectiveBuffer buffer(2047); + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES); std::string turnOn = directive("TurnOn", ID_1); ASSERT_RESULT(Result::COMPLETE, deliver(buffer, turnOn, 100)); ASSERT_RESULT(Result::DUPLICATE, buffer.append(turnOn.data(), turnOn.size(), 0, turnOn.size())); - ASSERT_RESULT(Result::DUPLICATE, deliver(buffer, turnOn, 64)); + ASSERT_RESULT(Result::DUPLICATE, deliver(buffer, turnOn, 64)); // in other fragments than the first time + TEST_ASSERT_EQUAL_UINT(1, buffer.waitingDirectives()); assertWaiting(buffer, turnOn); } -void test_repeat_that_differs_in_a_later_fragment_is_dropped() +void test_repeat_of_a_directive_that_was_read_does_not_hold_up_the_next() { - AlexaDirectiveBuffer buffer(2047); - std::string first = directive("TurnOn", ID_1); - std::string second = first; - second[second.size() - 10] = '#'; // the same up to its last fragment - - ASSERT_RESULT(Result::COMPLETE, deliver(buffer, first, 100)); - ASSERT_RESULT(Result::INCOMPLETE, buffer.append(second.data(), 100, 0, second.size())); - ASSERT_RESULT(Result::INCOMPLETE, buffer.append(second.data() + 100, 100, 100, second.size())); - ASSERT_RESULT(Result::BUSY, buffer.append(second.data() + 200, second.size() - 200, 200, second.size())); - assertWaiting(buffer, first); -} - -void test_repeat_still_arriving_when_the_first_is_read() -{ - AlexaDirectiveBuffer buffer(2047); + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES); std::string turnOn = directive("TurnOn", ID_1); + std::string turnOff = directive("TurnOff", ID_2); ASSERT_RESULT(Result::COMPLETE, deliver(buffer, turnOn, 100)); - ASSERT_RESULT(Result::INCOMPLETE, buffer.append(turnOn.data(), 100, 0, turnOn.size())); - buffer.release(); // loop() ran between two fragments of the repeat + buffer.release(); // loop() has handled it + // Its repeat and the next directive arrive together + ASSERT_RESULT(Result::DUPLICATE, deliver(buffer, turnOn, 100)); TEST_ASSERT_FALSE(buffer.ready()); - ASSERT_RESULT(Result::INCOMPLETE, buffer.append(turnOn.data() + 100, 100, 100, turnOn.size())); - ASSERT_RESULT(Result::DUPLICATE, buffer.append(turnOn.data() + 200, turnOn.size() - 200, 200, turnOn.size())); - TEST_ASSERT_FALSE(buffer.ready()); - - // and the buffer is free for the next directive - std::string turnOff = directive("TurnOff", ID_2); ASSERT_RESULT(Result::COMPLETE, deliver(buffer, turnOff, 100)); assertWaiting(buffer, turnOff); } -void test_directive_that_starts_like_the_one_just_read_is_not_lost() +void test_group_of_five_with_repeats_in_one_burst() { - AlexaDirectiveBuffer buffer(2047); + TEST_ASSERT_TRUE(PLACES >= 5); + + // Every directive followed by its repeat + AlexaDirectiveBuffer interleaved(LARGEST, MOST_BYTES); + for (unsigned n = 0; n < 5; n++) + { + ASSERT_RESULT(Result::COMPLETE, deliver(interleaved, numbered(n), 536)); + ASSERT_RESULT(Result::DUPLICATE, deliver(interleaved, numbered(n), 536)); + } + TEST_ASSERT_EQUAL_UINT(5, interleaved.waitingDirectives()); + for (unsigned n = 0; n < 5; n++) + { + assertWaiting(interleaved, numbered(n)); + interleaved.release(); + } + + // The five directives, then the five repeats + AlexaDirectiveBuffer trailing(LARGEST, MOST_BYTES); + for (unsigned n = 0; n < 5; n++) + { + ASSERT_RESULT(Result::COMPLETE, deliver(trailing, numbered(n), 536)); + } + for (unsigned n = 0; n < 5; n++) + { + ASSERT_RESULT(Result::DUPLICATE, deliver(trailing, numbered(n), 536)); + } + TEST_ASSERT_EQUAL_UINT(5, trailing.waitingDirectives()); +} + +void test_repeat_is_recognised_when_no_place_is_free() +{ + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES); + for (unsigned n = 0; n < PLACES; n++) + { + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, numbered(n), 100)); + } + + // Not a loss, so not QUEUE_FULL + ASSERT_RESULT(Result::DUPLICATE, deliver(buffer, numbered(0), 100)); + ASSERT_RESULT(Result::DUPLICATE, deliver(buffer, numbered(PLACES - 1), 100)); + TEST_ASSERT_EQUAL_UINT(PLACES, buffer.waitingDirectives()); +} + +void test_directive_that_differs_in_one_byte_is_not_a_repeat() +{ + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES); std::string first = directive("TurnOn", ID_1); std::string second = first; - second[second.size() - 10] = '#'; + second[second.size() - 10] = '#'; // the same length, the same up to its last fragment ASSERT_RESULT(Result::COMPLETE, deliver(buffer, first, 100)); - ASSERT_RESULT(Result::INCOMPLETE, buffer.append(second.data(), 100, 0, second.size())); - buffer.release(); // the first has been read: nothing waits, so the second has a place to go - ASSERT_RESULT(Result::INCOMPLETE, buffer.append(second.data() + 100, 100, 100, second.size())); - ASSERT_RESULT(Result::COMPLETE, buffer.append(second.data() + 200, second.size() - 200, 200, second.size())); + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, second, 100)); + assertWaiting(buffer, first); + buffer.release(); assertWaiting(buffer, second); } +void test_repeat_is_forgotten_after_as_many_directives_as_are_remembered() +{ + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES); + const unsigned remembered = AlexaRecentHashes::CAPACITY; + + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, numbered(0), 100)); + buffer.release(); + for (unsigned n = 1; n < remembered; n++) + { + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, numbered(n), 100)); + buffer.release(); + } + ASSERT_RESULT(Result::DUPLICATE, deliver(buffer, numbered(0), 100)); // the oldest that is remembered + + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, numbered(remembered), 100)); // takes its place + buffer.release(); + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, numbered(0), 100)); +} + +// --- no memory --- + +static bool memoryLeft = true; + +static void *scarceMemory(size_t size) +{ + return memoryLeft ? malloc(size) : nullptr; +} + +void test_directive_without_memory_is_reported_by_its_last_fragment() +{ + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES, scarceMemory); + std::string turnOn = directive("TurnOn", ID_1); + + memoryLeft = false; + ASSERT_RESULT(Result::NO_MEMORY, deliver(buffer, turnOn, 100)); + TEST_ASSERT_FALSE(buffer.ready()); + + // It was not remembered: when it comes again and there is memory, it is handled + memoryLeft = true; + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, turnOn, 100)); + assertWaiting(buffer, turnOn); +} + +void test_repeat_without_memory_is_still_a_repeat() +{ + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES, scarceMemory); + std::string turnOn = directive("TurnOn", ID_1); + + memoryLeft = true; + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, turnOn, 100)); + memoryLeft = false; + ASSERT_RESULT(Result::DUPLICATE, deliver(buffer, turnOn, 100)); + memoryLeft = true; + assertWaiting(buffer, turnOn); +} + // --- fragments that do not fit --- void test_fragment_with_a_gap_drops_the_message() { - AlexaDirectiveBuffer buffer(2047); + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES); std::string turnOn = directive("TurnOn", ID_1); ASSERT_RESULT(Result::INCOMPLETE, buffer.append(turnOn.data(), 100, 0, turnOn.size())); @@ -251,7 +427,7 @@ void test_fragment_with_a_gap_drops_the_message() void test_fragment_beyond_the_total_is_refused() { - AlexaDirectiveBuffer buffer(2047); + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES); std::string message(300, 'x'); // The block holds 200 + 1 bytes; a fragment that would end at 250 must not be copied @@ -266,7 +442,7 @@ void test_fragment_beyond_the_total_is_refused() void test_fragment_with_another_total_is_refused() { - AlexaDirectiveBuffer buffer(2047); + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES); std::string message(300, 'x'); ASSERT_RESULT(Result::INCOMPLETE, buffer.append(message.data(), 100, 0, 200)); @@ -276,7 +452,7 @@ void test_fragment_with_another_total_is_refused() void test_fragment_without_a_start_is_refused_once() { - AlexaDirectiveBuffer buffer(2047); + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES); std::string turnOn = directive("TurnOn", ID_1); ASSERT_RESULT(Result::OUT_OF_ORDER, buffer.append(turnOn.data() + 100, 100, 100, turnOn.size())); @@ -286,18 +462,19 @@ void test_fragment_without_a_start_is_refused_once() void test_new_message_replaces_one_that_never_completed() { - AlexaDirectiveBuffer buffer(2047); + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES); std::string turnOn = directive("TurnOn", ID_1); std::string turnOff = directive("TurnOff", ID_2); ASSERT_RESULT(Result::INCOMPLETE, buffer.append(turnOn.data(), 100, 0, turnOn.size())); ASSERT_RESULT(Result::COMPLETE, deliver(buffer, turnOff, 100)); + TEST_ASSERT_EQUAL_UINT(1, buffer.waitingDirectives()); assertWaiting(buffer, turnOff); } void test_lost_connection_drops_the_arriving_message_only() { - AlexaDirectiveBuffer buffer(2047); + AlexaDirectiveBuffer buffer(LARGEST, MOST_BYTES); std::string turnOn = directive("TurnOn", ID_1); std::string turnOff = directive("TurnOff", ID_2); @@ -307,11 +484,15 @@ void test_lost_connection_drops_the_arriving_message_only() ASSERT_RESULT(Result::OUT_OF_ORDER, buffer.append(turnOn.data() + 100, 100, 100, turnOn.size())); TEST_ASSERT_FALSE(buffer.ready()); - // A directive waits, its repeat is arriving + // A directive waits, another one is arriving ASSERT_RESULT(Result::COMPLETE, deliver(buffer, turnOff, 100)); - ASSERT_RESULT(Result::INCOMPLETE, buffer.append(turnOff.data(), 100, 0, turnOff.size())); + ASSERT_RESULT(Result::INCOMPLETE, buffer.append(turnOn.data(), 100, 0, turnOn.size())); buffer.cancelArrival(); + TEST_ASSERT_EQUAL_UINT(1, buffer.waitingDirectives()); assertWaiting(buffer, turnOff); + + // What arrived of the message that was cut off does not make the whole message a repeat + ASSERT_RESULT(Result::COMPLETE, deliver(buffer, turnOn, 100)); } // --- repeated messageIds --- @@ -327,36 +508,55 @@ void test_repeated_message_id_is_recognised() TEST_ASSERT_TRUE(recent.seenBefore(ID_2)); } -void test_only_the_last_four_ids_are_remembered() +void test_only_the_last_ids_are_remembered() { AlexaRecentIds recent; - const char *ids[] = {"id-1", "id-2", "id-3", "id-4", "id-5"}; + const unsigned remembered = AlexaRecentIds::CAPACITY; + char id[16]; - for (const char *id : ids) + for (unsigned n = 0; n <= remembered; n++) { + snprintf(id, sizeof(id), "id-%u", n); TEST_ASSERT_FALSE(recent.seenBefore(id)); } - // id-5 took the place of id-1 - TEST_ASSERT_TRUE(recent.seenBefore("id-2")); - TEST_ASSERT_TRUE(recent.seenBefore("id-3")); - TEST_ASSERT_TRUE(recent.seenBefore("id-4")); - TEST_ASSERT_TRUE(recent.seenBefore("id-5")); - TEST_ASSERT_FALSE(recent.seenBefore("id-1")); + // The last one took the place of id-0 + for (unsigned n = 1; n <= remembered; n++) + { + snprintf(id, sizeof(id), "id-%u", n); + TEST_ASSERT_TRUE(recent.seenBefore(id)); + } + TEST_ASSERT_FALSE(recent.seenBefore("id-0")); } void test_recognising_a_repeat_does_not_use_a_place() { AlexaRecentIds recent; + const unsigned remembered = AlexaRecentIds::CAPACITY; + char id[16]; - TEST_ASSERT_FALSE(recent.seenBefore("id-1")); - for (int i = 0; i < 10; i++) + TEST_ASSERT_FALSE(recent.seenBefore("id-0")); + for (int i = 0; i < 100; i++) { - TEST_ASSERT_TRUE(recent.seenBefore("id-1")); + TEST_ASSERT_TRUE(recent.seenBefore("id-0")); } - TEST_ASSERT_FALSE(recent.seenBefore("id-2")); - TEST_ASSERT_FALSE(recent.seenBefore("id-3")); - TEST_ASSERT_FALSE(recent.seenBefore("id-4")); - TEST_ASSERT_TRUE(recent.seenBefore("id-1")); + for (unsigned n = 1; n < remembered; n++) + { + snprintf(id, sizeof(id), "id-%u", n); + TEST_ASSERT_FALSE(recent.seenBefore(id)); + } + TEST_ASSERT_TRUE(recent.seenBefore("id-0")); +} + +void test_hash_is_fnv_1a_64() +{ + // Test vectors of the FNV reference code + TEST_ASSERT_EQUAL_HEX64(0xcbf29ce484222325ULL, AlexaBridgeLogic::hashBytes("", 0)); + TEST_ASSERT_EQUAL_HEX64(0xaf63dc4c8601ec8cULL, AlexaBridgeLogic::hashBytes("a", 1)); + TEST_ASSERT_EQUAL_HEX64(0x85944171f73967e8ULL, AlexaBridgeLogic::hashBytes("foobar", 6)); + + // In pieces, as the fragments of a message are hashed + uint64_t hash = AlexaBridgeLogic::hashBytes("foo", 3); + TEST_ASSERT_EQUAL_HEX64(0x85944171f73967e8ULL, AlexaBridgeLogic::hashBytes("bar", 3, hash)); } void test_directive_without_a_message_id_is_never_a_repeat() @@ -489,7 +689,7 @@ void test_largest_directive_can_be_answered() std::string largest = directive("TurnOn", ID_1, std::string(ALEX2ESP_MAX_DIRECTIVE - withoutToken.size(), 'T')); TEST_ASSERT_EQUAL_UINT(ALEX2ESP_MAX_DIRECTIVE, largest.size()); - AlexaDirectiveBuffer buffer(ALEX2ESP_MAX_DIRECTIVE); + AlexaDirectiveBuffer buffer(ALEX2ESP_MAX_DIRECTIVE, ALEX2ESP_MAX_QUEUED_BYTES); ASSERT_RESULT(Result::COMPLETE, deliver(buffer, largest, 536)); JsonDocument received; TEST_ASSERT_TRUE(deserializeJson(received, buffer.data(), buffer.length()) == DeserializationError::Ok); @@ -609,13 +809,21 @@ int main(int, char **) RUN_TEST(test_refused_directive_does_not_block_the_next_one); RUN_TEST(test_empty_message_is_reported); - RUN_TEST(test_second_directive_before_loop_is_dropped_and_the_first_kept); - RUN_TEST(test_second_directive_of_the_same_size_is_dropped); + RUN_TEST(test_directives_wait_in_the_order_they_arrived); + RUN_TEST(test_directive_beyond_the_places_is_dropped_and_the_others_kept); + RUN_TEST(test_place_that_becomes_free_while_a_directive_arrives_is_used); + RUN_TEST(test_waiting_directives_are_limited_in_bytes); + RUN_TEST(test_empty_queue_takes_the_largest_directive); - RUN_TEST(test_repeat_of_the_waiting_directive_is_a_duplicate); - RUN_TEST(test_repeat_that_differs_in_a_later_fragment_is_dropped); - RUN_TEST(test_repeat_still_arriving_when_the_first_is_read); - RUN_TEST(test_directive_that_starts_like_the_one_just_read_is_not_lost); + RUN_TEST(test_repeat_of_a_waiting_directive_is_discarded); + RUN_TEST(test_repeat_of_a_directive_that_was_read_does_not_hold_up_the_next); + RUN_TEST(test_group_of_five_with_repeats_in_one_burst); + RUN_TEST(test_repeat_is_recognised_when_no_place_is_free); + RUN_TEST(test_directive_that_differs_in_one_byte_is_not_a_repeat); + RUN_TEST(test_repeat_is_forgotten_after_as_many_directives_as_are_remembered); + + RUN_TEST(test_directive_without_memory_is_reported_by_its_last_fragment); + RUN_TEST(test_repeat_without_memory_is_still_a_repeat); RUN_TEST(test_fragment_with_a_gap_drops_the_message); RUN_TEST(test_fragment_beyond_the_total_is_refused); @@ -625,8 +833,9 @@ int main(int, char **) RUN_TEST(test_lost_connection_drops_the_arriving_message_only); RUN_TEST(test_repeated_message_id_is_recognised); - RUN_TEST(test_only_the_last_four_ids_are_remembered); + RUN_TEST(test_only_the_last_ids_are_remembered); RUN_TEST(test_recognising_a_repeat_does_not_use_a_place); + RUN_TEST(test_hash_is_fnv_1a_64); RUN_TEST(test_directive_without_a_message_id_is_never_a_repeat); RUN_TEST(test_message_within_the_limit_is_measured);