Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
71 changes: 56 additions & 15 deletions src/mqtt_broker.c
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Comment thread
aidangarske marked this conversation as resolved.
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);
Comment thread
aidangarske marked this conversation as resolved.
}
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
Expand All @@ -7202,28 +7244,20 @@ 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);
BrokerRetained_FreeAll(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) {
Expand Down Expand Up @@ -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)
Expand Down
57 changes: 46 additions & 11 deletions src/mqtt_broker_persist.c
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Expand Down Expand Up @@ -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;
Expand All @@ -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);
Expand All @@ -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.).
Expand All @@ -1996,16 +2002,26 @@ 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));
ctx.broker = 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);
Expand All @@ -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);
Expand All @@ -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);
Expand All @@ -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);
Expand All @@ -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;
}

Expand Down
25 changes: 15 additions & 10 deletions wolfmqtt/mqtt_broker.h
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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;

/* -------------------------------------------------------------------------- */
Expand Down Expand Up @@ -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);

Expand Down Expand Up @@ -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);
Comment thread
aidangarske marked this conversation as resolved.
#endif /* WOLFMQTT_BROKER_PERSIST */

#ifndef WOLFMQTT_STATIC_MEMORY
Expand Down
Loading