diff --git a/src/mqtt_broker.c b/src/mqtt_broker.c index d3187c780..9cbecc258 100644 --- a/src/mqtt_broker.c +++ b/src/mqtt_broker.c @@ -7180,11 +7180,53 @@ int MqttBroker_Stop(MqttBroker* broker) return MQTT_CODE_SUCCESS; } +static void BrokerSubs_FreeAll(MqttBroker* broker) +{ +#ifdef WOLFMQTT_STATIC_MEMORY + BROKER_FORCE_ZERO(broker->subs, sizeof(broker->subs)); +#else + while (broker->subs != NULL) { + BrokerSub* next = broker->subs->next; + if (broker->subs->filter != NULL) { + BROKER_FORCE_ZERO(broker->subs->filter, + XSTRLEN(broker->subs->filter) + 1); + WOLFMQTT_FREE(broker->subs->filter); + } + if (broker->subs->client_id != NULL) { + BROKER_FORCE_ZERO(broker->subs->client_id, + XSTRLEN(broker->subs->client_id) + 1); + WOLFMQTT_FREE(broker->subs->client_id); + } + WOLFMQTT_FREE(broker->subs); + broker->subs = next; + } +#endif +} + +#ifdef WOLFMQTT_BROKER_PERSIST +/* BrokerPersist_Restore is called on an empty broker before its first start. */ +WOLFMQTT_LOCAL void BrokerPersist_RestoreRollback(MqttBroker* broker) +{ + if (broker == NULL) { + return; + } + BrokerSubs_FreeAll(broker); +#ifndef WOLFMQTT_STATIC_MEMORY + BrokerOrphan_FreeAll(broker); +#endif +#ifdef WOLFMQTT_BROKER_RETAINED + BrokerRetained_FreeAll(broker); +#endif + broker->next_packet_id = 1; +} +#endif + int MqttBroker_Free(MqttBroker* broker) { if (broker == NULL) { return MQTT_CODE_ERROR_BAD_ARG; } + broker->running = 0; /* Disconnect and free all clients and subscriptions */ #ifdef WOLFMQTT_STATIC_MEMORY @@ -7202,21 +7244,9 @@ int MqttBroker_Free(MqttBroker* broker) BrokerSubs_RemoveClient(broker, broker->clients); BrokerClient_Remove(broker, broker->clients); } - /* Free any orphaned subs (e.g. from clean_session=0 clients) */ - while (broker->subs) { - BrokerSub* next = broker->subs->next; - if (broker->subs->filter) { - BROKER_FORCE_ZERO(broker->subs->filter, - XSTRLEN(broker->subs->filter) + 1); - WOLFMQTT_FREE(broker->subs->filter); - } - if (broker->subs->client_id) { - WOLFMQTT_FREE(broker->subs->client_id); - } - WOLFMQTT_FREE(broker->subs); - broker->subs = next; - } #endif + /* Free any orphaned subs (e.g. from clean_session=0 clients). */ + BrokerSubs_FreeAll(broker); /* Clean up pending wills and retained messages */ BrokerPendingWill_FreeAll(broker); @@ -7224,6 +7254,10 @@ int MqttBroker_Free(MqttBroker* broker) #ifndef WOLFMQTT_STATIC_MEMORY BrokerOrphan_FreeAll(broker); #endif +#ifdef WOLFMQTT_BROKER_PERSIST + broker->persist_restored = 0; + broker->persist = NULL; +#endif #ifdef ENABLE_MQTT_TLS if (broker->tls_ctx != NULL) { @@ -7601,7 +7635,14 @@ int wolfmqtt_broker(int argc, char** argv) * "DEV-KEY" log line below makes the choice obvious. */ persist_hooks.derive_key = wolfmqtt_broker_dev_derive_key; #endif - (void)MqttBroker_SetPersistHooks(&broker, &persist_hooks); + rc = MqttBroker_SetPersistHooks(&broker, &persist_hooks); + if (rc != MQTT_CODE_SUCCESS) { + PRINTF("broker: persist hook install failed rc=%d", rc); + MqttBroker_Free(&broker); + MqttBrokerNet_PersistPosix_Free(&persist_hooks); + BROKER_WIPE_AUTH_PASS(); + return rc; + } PRINTF("broker: persist enabled dir=%s%s", persist_dir, #if defined(WOLFMQTT_BROKER_PERSIST_ENCRYPT) && \ defined(WOLFMQTT_BROKER_PERSIST_ENCRYPT_DEV_KEY) diff --git a/src/mqtt_broker_persist.c b/src/mqtt_broker_persist.c index de22f1f0e..01bbc79dc 100644 --- a/src/mqtt_broker_persist.c +++ b/src/mqtt_broker_persist.c @@ -564,6 +564,9 @@ int MqttBroker_SetPersistHooks(MqttBroker* broker, if (broker == NULL) { return MQTT_CODE_ERROR_BAD_ARG; } + if (hooks != NULL && (broker->running || broker->persist_restored)) { + return MQTT_CODE_ERROR_BAD_ARG; + } broker->persist = hooks; return MQTT_CODE_SUCCESS; } @@ -1954,11 +1957,7 @@ static int wmqb_wipe_all(MqttBroker* broker) #endif } -/* Restore is intended to be called exactly once per process (from - * MqttBroker_Start). Calling it more than once will re-insert - * already-restored subs / retained / OUTQ entries because the splice - * paths below do not check for duplicates against current in-memory - * state. */ +/* Restore runs once for each initialized broker. */ int BrokerPersist_Restore(MqttBroker* broker) { const MqttBrokerPersistHooks* h; @@ -1969,6 +1968,9 @@ int BrokerPersist_Restore(MqttBroker* broker) if (broker == NULL || broker->persist == NULL) { return 0; } + if (broker->persist_restored) { + return 0; + } h = broker->persist; rc = wmqb_meta_check(broker, &meta_present); @@ -1983,7 +1985,11 @@ int BrokerPersist_Restore(MqttBroker* broker) WMQB_LOG_ERR(broker, "broker: persist schema mismatch - wiping all records"); (void)wmqb_wipe_all(broker); - return wmqb_meta_write(broker); + rc = wmqb_meta_write(broker); + if (rc == 0) { + broker->persist_restored = 1; + } + return rc; } if (rc != 0) { /* Real backend error (I/O failure, permission denied, etc.). @@ -1996,7 +2002,11 @@ int BrokerPersist_Restore(MqttBroker* broker) } if (!meta_present) { /* First run - no state to restore. Just stamp META. */ - return wmqb_meta_write(broker); + rc = wmqb_meta_write(broker); + if (rc == 0) { + broker->persist_restored = 1; + } + return rc; } XMEMSET(&ctx, 0, sizeof(ctx)); @@ -2004,8 +2014,14 @@ int BrokerPersist_Restore(MqttBroker* broker) #ifndef WOLFMQTT_STATIC_MEMORY /* Sessions first so subs and OUTQ entries can find their owner. */ if (h->kv_iter != NULL) { - (void)wmqb_kv_iter(broker, BROKER_PERSIST_NS_SESSION, + rc = wmqb_kv_iter(broker, BROKER_PERSIST_NS_SESSION, wmqb_iter_session_cb, &ctx); + if (rc != 0) { + WMQB_LOG_ERR(broker, + "broker: persist restore sessions failed rc=%d", rc); + BrokerPersist_RestoreRollback(broker); + return rc; + } WMQB_LOG_INFO(broker, "broker: persist restore sessions loaded=%d skipped=%d", ctx.loaded, ctx.skipped); @@ -2015,8 +2031,14 @@ int BrokerPersist_Restore(MqttBroker* broker) #endif #ifdef WOLFMQTT_BROKER_RETAINED if (h->kv_iter != NULL) { - (void)wmqb_kv_iter(broker, BROKER_PERSIST_NS_RETAINED, + rc = wmqb_kv_iter(broker, BROKER_PERSIST_NS_RETAINED, wmqb_iter_retained_cb, &ctx); + if (rc != 0) { + WMQB_LOG_ERR(broker, + "broker: persist restore retained failed rc=%d", rc); + BrokerPersist_RestoreRollback(broker); + return rc; + } WMQB_LOG_INFO(broker, "broker: persist restore retained loaded=%d skipped=%d", ctx.loaded, ctx.skipped); @@ -2025,8 +2047,14 @@ int BrokerPersist_Restore(MqttBroker* broker) } #endif if (h->kv_iter != NULL) { - (void)wmqb_kv_iter(broker, BROKER_PERSIST_NS_SUBS, + rc = wmqb_kv_iter(broker, BROKER_PERSIST_NS_SUBS, wmqb_iter_subs_cb, &ctx); + if (rc != 0) { + WMQB_LOG_ERR(broker, + "broker: persist restore subs failed rc=%d", rc); + BrokerPersist_RestoreRollback(broker); + return rc; + } WMQB_LOG_INFO(broker, "broker: persist restore subs loaded=%d skipped=%d", ctx.loaded, ctx.skipped); @@ -2035,8 +2063,14 @@ int BrokerPersist_Restore(MqttBroker* broker) } #ifndef WOLFMQTT_STATIC_MEMORY if (h->kv_iter != NULL) { - (void)wmqb_kv_iter(broker, BROKER_PERSIST_NS_OUTQ, + rc = wmqb_kv_iter(broker, BROKER_PERSIST_NS_OUTQ, wmqb_iter_outq_cb, &ctx); + if (rc != 0) { + WMQB_LOG_ERR(broker, + "broker: persist restore outq failed rc=%d", rc); + BrokerPersist_RestoreRollback(broker); + return rc; + } WMQB_LOG_INFO(broker, "broker: persist restore outq loaded=%d skipped=%d", ctx.loaded, ctx.skipped); @@ -2046,6 +2080,7 @@ int BrokerPersist_Restore(MqttBroker* broker) * persisted OUTQ records via the existing helpers. */ wmqb_restore_expiry_sweep(broker); #endif + broker->persist_restored = 1; return 0; } diff --git a/wolfmqtt/mqtt_broker.h b/wolfmqtt/mqtt_broker.h index e638ea0fd..9412ae4ec 100644 --- a/wolfmqtt/mqtt_broker.h +++ b/wolfmqtt/mqtt_broker.h @@ -320,8 +320,8 @@ typedef struct MqttBrokerNet { /* The persistence layer is intentionally hook-based so the broker can run * on top of POSIX files, embedded flash, an external KV store, or an * in-RAM stub used by tests. Each hook returns 0 on success or a negative - * error code (broker logs and skips persist for that record - the - * in-memory state is still authoritative). + * error code. Shadow-write failures are logged while in-memory state remains + * authoritative. Iterator failures during restore abort MqttBroker_Start. * * Both a key/value API and a streaming API are provided. The broker will * use whichever family the registered hook implements; any individual @@ -755,6 +755,9 @@ typedef struct MqttBroker { * walk the orphan list on every single call. */ WOLFMQTT_BROKER_TIME_T orphan_last_expire_check; #endif +#ifdef WOLFMQTT_BROKER_PERSIST + byte persist_restored; +#endif } MqttBroker; /* -------------------------------------------------------------------------- */ @@ -803,10 +806,10 @@ WOLFMQTT_API int MqttBrokerNet_Init(MqttBrokerNet* net); #endif #ifdef WOLFMQTT_BROKER_PERSIST -/* Install persistence hooks on the broker. Must be called before - * MqttBroker_Start. Passing NULL clears any previously installed hooks - * (reverts to in-memory-only behavior). The MqttBrokerPersistHooks - * struct must outlive the broker. */ +/* Install non-NULL persistence hooks after MqttBroker_Init and before the + * first MqttBroker_Start, or after MqttBroker_Free. Passing NULL may detach + * the current hooks at any time. The MqttBrokerPersistHooks struct must + * outlive the broker while installed. */ WOLFMQTT_API int MqttBroker_SetPersistHooks(MqttBroker* broker, const MqttBrokerPersistHooks* hooks); @@ -863,11 +866,13 @@ WOLFMQTT_LOCAL int BrokerPersist_DelOutPub(MqttBroker* broker, WOLFMQTT_LOCAL int BrokerPersist_DelOutQueue(MqttBroker* broker, const char* client_id); -/* Startup-time restore: iterate persisted records and rebuild the - * in-memory tables. Called from MqttBroker_Init when hooks are - * installed. Wipes everything and re-stamps the META namespace if - * the persisted schema version doesn't match. */ +/* Startup-time restore on a freshly initialized, empty broker. Rebuilds the + * in-memory tables before the first MqttBroker_Start accepts clients. Wipes + * and re-stamps META when the persisted schema version does not match. */ WOLFMQTT_LOCAL int BrokerPersist_Restore(MqttBroker* broker); +/* Discard state loaded by a failed startup restore. Defined in + * mqtt_broker.c because it uses the broker's shared teardown helpers. */ +WOLFMQTT_LOCAL void BrokerPersist_RestoreRollback(MqttBroker* broker); #endif /* WOLFMQTT_BROKER_PERSIST */ #ifndef WOLFMQTT_STATIC_MEMORY