From 969e9ac67ee8f36f9bceab6b2f41eac0528f0525 Mon Sep 17 00:00:00 2001 From: Kareem Date: Mon, 24 Aug 2026 15:37:29 -0700 Subject: [PATCH] Ensure expired messages have their records removed. Thanks to Maksim Hayder for the report. --- src/mqtt_broker.c | 26 +++++++++++++ src/mqtt_broker_persist.c | 81 +++++++++++++++++++++++++++++++++++++-- 2 files changed, 103 insertions(+), 4 deletions(-) diff --git a/src/mqtt_broker.c b/src/mqtt_broker.c index 1936dba..0833feb 100644 --- a/src/mqtt_broker.c +++ b/src/mqtt_broker.c @@ -4574,6 +4574,23 @@ static int BrokerSubs_ReassociateClient(MqttBroker* broker, /* Retained message management */ /* -------------------------------------------------------------------------- */ #ifdef WOLFMQTT_BROKER_RETAINED +#ifdef WOLFMQTT_BROKER_PERSIST +/* Retire the stored record behind a retained message that expiry is about to + * drop from RAM. A backend refusal is reported rather than swallowed; the + * restore sweep retries it on the next start. */ +static void BrokerRetained_PersistExpired(MqttBroker* broker, + const char* topic) +{ + if (BrokerPersist_DelRetained(broker, topic) != 0) { + WBLOG_ERR(broker, "broker: retained expiry persist delete failed " + "topic=%s", BrokerLog_Sanitize(topic)); + } +} +#else + #define BrokerRetained_PersistExpired(b, t) \ + do { (void)(b); (void)(t); } while (0) +#endif + /* Free retained messages whose v5 Message Expiry Interval has elapsed so they * stop occupying a slot / the retained_count cap. Delivery also reaps expired * entries, but a publisher can hit the cap with only-expired entries before any @@ -4592,6 +4609,7 @@ static void BrokerRetained_ReapExpired(MqttBroker* broker, (now - rm->store_time) >= rm->expiry_sec) { WBLOG_DBG(broker, "broker: retained expired topic=%s", BrokerLog_Sanitize(rm->topic)); + BrokerRetained_PersistExpired(broker, rm->topic); BROKER_FORCE_ZERO(rm, sizeof(BrokerRetainedMsg)); } } @@ -4611,6 +4629,7 @@ static void BrokerRetained_ReapExpired(MqttBroker* broker, (now - cur->store_time) >= cur->expiry_sec) { WBLOG_DBG(broker, "broker: retained expired topic=%s", BrokerLog_Sanitize(cur->topic)); + BrokerRetained_PersistExpired(broker, cur->topic); if (prev != NULL) { prev->next = next; } @@ -5321,6 +5340,7 @@ static int BrokerRetained_DeliverToClient(MqttBroker* broker, (now - rm->store_time) >= rm->expiry_sec) { WBLOG_DBG(broker, "broker: retained expired topic=%s", BrokerLog_Sanitize(rm->topic)); + BrokerRetained_PersistExpired(broker, rm->topic); BROKER_FORCE_ZERO(rm, sizeof(BrokerRetainedMsg)); continue; } @@ -5431,6 +5451,12 @@ static int BrokerRetained_DeliverToClient(MqttBroker* broker, (rm->expiry_sec > 0 && now >= rm->store_time && (now - rm->store_time) >= rm->expiry_sec)) { + /* A pending_delete node already had its record retired by + * BrokerRetained_Delete; expiry has to retire its own. Done ahead + * of the deferral below so the flag keeps that one meaning. */ + if (!rm->pending_delete) { + BrokerRetained_PersistExpired(broker, rm->topic); + } if (broker->retained_delivering > 1) { rm->pending_delete = 1; rm_prev = rm; diff --git a/src/mqtt_broker_persist.c b/src/mqtt_broker_persist.c index 6db9c8f..6827aeb 100644 --- a/src/mqtt_broker_persist.c +++ b/src/mqtt_broker_persist.c @@ -1285,6 +1285,10 @@ struct wmqb_ended_session { }; #endif +/* wmqb_decode_and_insert_retained: record elapsed, distinct from the cap skip + * (1) so only genuinely dead keys are dropped from the store. */ +#define WMQB_RETAINED_EXPIRED 2 + /* Restore iterator context. Used for retained-msg, subs, session, * and OUTQ callbacks. */ struct wmqb_restore_ctx { @@ -1296,6 +1300,13 @@ struct wmqb_restore_ctx { /* Session keys whose zero Session Expiry ended them before this restart */ struct wmqb_ended_session* ended; int ended_count; +#endif + /* Retained keys whose Message Expiry Interval elapsed before this restart */ +#ifdef WOLFMQTT_STATIC_MEMORY + byte expired_key[BROKER_MAX_TOPIC_LEN]; + word16 expired_key_len; +#else + struct wmqb_wipe_key* expired; #endif }; @@ -1563,7 +1574,7 @@ static int wmqb_decode_and_insert_retained(MqttBroker* broker, if (expiry > 0 && now >= (WOLFMQTT_BROKER_TIME_T)store_time && (now - (WOLFMQTT_BROKER_TIME_T)store_time) >= expiry) { - return 1; + return WMQB_RETAINED_EXPIRED; } #ifdef WOLFMQTT_STATIC_MEMORY @@ -1660,17 +1671,77 @@ static int wmqb_iter_retained_cb(const byte* key, word16 key_len, { struct wmqb_restore_ctx* c = (struct wmqb_restore_ctx*)cb_ctx; int rc; - (void)key; (void)key_len; +#ifndef WOLFMQTT_STATIC_MEMORY + struct wmqb_wipe_key* node; +#endif + rc = wmqb_decode_and_insert_retained(c->broker, blob, blob_len); if (rc == 0) { c->loaded++; + return 0; } - else { - c->skipped++; + c->skipped++; + /* Elapsed while the broker was down: restore drops it from memory, so its + * record has to go too. Stash the key - the backend iterator must not be + * mutated while it runs. */ + if (rc != WMQB_RETAINED_EXPIRED || key == NULL || key_len == 0) { + return 0; } +#ifdef WOLFMQTT_STATIC_MEMORY + /* No allocator: carry one key per restore, the rest drain on later ones. */ + if (c->expired_key_len == 0 && key_len <= sizeof(c->expired_key)) { + XMEMCPY(c->expired_key, key, key_len); + c->expired_key_len = key_len; + } +#else + node = (struct wmqb_wipe_key*)WOLFMQTT_MALLOC(sizeof(*node)); + if (node != NULL) { + node->key = (byte*)WOLFMQTT_MALLOC(key_len); + if (node->key == NULL) { + WOLFMQTT_FREE(node); + } + else { + XMEMCPY(node->key, key, key_len); + node->key_len = key_len; + node->next = c->expired; + c->expired = node; + } + } +#endif return 0; /* always continue */ } +/* Retire the keys wmqb_iter_retained_cb stashed. del==0 just releases them, + * for the abort path where no on-disk state may change. */ +static void wmqb_restore_drop_expired_retained(MqttBroker* broker, + struct wmqb_restore_ctx* c, byte del) +{ +#ifdef WOLFMQTT_STATIC_MEMORY + if (del && c->expired_key_len > 0 && + wmqb_kv_del_commit(broker, BROKER_PERSIST_NS_RETAINED, + c->expired_key, c->expired_key_len) != 0) { + WMQB_LOG_ERR(broker, "broker: persist expired retained delete " + "failed - deferred to next restore"); + } + c->expired_key_len = 0; +#else + struct wmqb_wipe_key* cur = c->expired; + + while (cur != NULL) { + struct wmqb_wipe_key* next = cur->next; + if (del && wmqb_kv_del_commit(broker, BROKER_PERSIST_NS_RETAINED, + cur->key, cur->key_len) != 0) { + WMQB_LOG_ERR(broker, "broker: persist expired retained delete " + "failed - deferred to next restore"); + } + WOLFMQTT_FREE(cur->key); + WOLFMQTT_FREE(cur); + cur = next; + } + c->expired = NULL; +#endif +} + /* Allocate orphan subs from a decoded NS_SUBS blob. The blob key carries * the client_id - subs created here have client=NULL, client_id set; * the existing BrokerSubs_ReassociateClient path on reconnect rebinds @@ -2561,9 +2632,11 @@ int BrokerPersist_Restore(MqttBroker* broker) if (rc != 0) { WMQB_LOG_ERR(broker, "broker: persist restore retained failed rc=%d", rc); + wmqb_restore_drop_expired_retained(broker, &ctx, 0); BrokerPersist_RestoreRollback(broker); return rc; } + wmqb_restore_drop_expired_retained(broker, &ctx, 1); WMQB_LOG_INFO(broker, "broker: persist restore retained loaded=%d skipped=%d", ctx.loaded, ctx.skipped);