From fca9678452824ea91b13a3c31578f216865d7ba2 Mon Sep 17 00:00:00 2001 From: Kareem Date: Thu, 24 Sep 2026 16:08:33 -0700 Subject: [PATCH 1/8] wolfMQTT broker: stop the outbound drain when the Session hand-off takes its queue. Thanks to Gwanhyun Lee for the report. --- src/mqtt_broker.c | 22 +++++ tests/test_broker_connect.c | 157 ++++++++++++++++++++++++++++++++++++ 2 files changed, 179 insertions(+) diff --git a/src/mqtt_broker.c b/src/mqtt_broker.c index fcea7ee23..045c74fe0 100644 --- a/src/mqtt_broker.c +++ b/src/mqtt_broker.c @@ -1999,6 +1999,22 @@ static void BrokerClient_FreeOutQueue(BrokerClient* bc) bc->out_q_pending_len = 0; } +/* A write can service this same client's close callback before it returns - + * the WebSocket transport runs lws_service inline - and that hands the whole + * outbound queue to the Session carrier. Entries the drain is walking were + * reached from bc->out_q_head, so an empty head while one is still held means + * the carrier owns them now: they must not be unlinked, freed, or re-linked + * into bc. Returns non-zero when the drain has to stop for that reason. */ +static int BrokerClient_OutQueueMoved(BrokerClient* bc, int enc_len) +{ + if (bc->out_q_head != NULL) { + return 0; + } + bc->out_q_pending_len = 0; + BROKER_FORCE_ZERO(bc->tx_buf, enc_len); + return 1; +} + /* Send as many QUEUED entries from out_q as the inflight cap allows. * * Ordering: walks from out_q_head, never reorders. Already-sent entries @@ -2059,6 +2075,9 @@ static int BrokerClient_DrainOutQueue(BrokerClient* bc) return MQTT_CODE_ERROR_SYSTEM; } wr_rc = MqttPacket_Write(&bc->client, bc->tx_buf, rel_rc); + if (BrokerClient_OutQueueMoved(bc, rel_rc)) { + return sent; + } if (wr_rc == MQTT_CODE_CONTINUE) { bc->out_q_pending_len = rel_rc; return wr_rc; @@ -2152,6 +2171,9 @@ static int BrokerClient_DrainOutQueue(BrokerClient* bc) { int wr_rc; wr_rc = MqttPacket_Write(&bc->client, bc->tx_buf, enc_rc); + if (BrokerClient_OutQueueMoved(bc, enc_rc)) { + return sent; + } /* Scrub the forwarded PUBLISH (which may carry a replayed will or * an application payload) once the buffer is idle. Skip only the * MQTT_CODE_CONTINUE case, where a non-blocking or TLS-async send diff --git a/tests/test_broker_connect.c b/tests/test_broker_connect.c index 37557a96f..664f91f17 100644 --- a/tests/test_broker_connect.c +++ b/tests/test_broker_connect.c @@ -63,6 +63,8 @@ static void* g_capture_alloc_ptr; static byte g_capture_freed[64]; static size_t g_capture_freed_len; static int g_capture_free_seen; +static void* g_free_watch_ptr; +static int g_free_watch_count; void* wolfmqtt_test_broker_malloc(size_t size) { @@ -93,6 +95,9 @@ void wolfmqtt_test_broker_free(void* ptr) { if (ptr != NULL) { g_alloc_free_count++; + if (ptr == g_free_watch_ptr) { + g_free_watch_count++; + } if (ptr == g_capture_alloc_ptr) { g_capture_freed_len = g_capture_alloc_size; XMEMCPY(g_capture_freed, ptr, g_capture_freed_len); @@ -108,6 +113,21 @@ unsigned long wolfmqtt_test_broker_time_s(void) return g_broker_time_s; } +#ifndef WOLFMQTT_STATIC_MEMORY +/* Count frees of one specific allocation, so a test can assert on the fate of + * a single object rather than on a total that unrelated activity moves. */ +static void broker_test_watch_free(void* ptr) +{ + g_free_watch_ptr = ptr; + g_free_watch_count = 0; +} + +static int broker_test_watched_frees(void) +{ + return g_free_watch_count; +} +#endif + #ifndef WOLFMQTT_STATIC_MEMORY static void broker_test_fail_alloc_after(int successful_allocations) { @@ -1979,6 +1999,142 @@ TEST(fanout_survives_reentrant_sub_free) MqttBroker_Free(&broker); } +/* The outbound drain holds a queue entry across the write that sends it. That + * write can service this client's own close callback before it returns - the + * WebSocket transport runs lws_service inline - and the callback hands the + * whole queue to the carrier that keeps the Session alive while the client is + * gone. The entry the drain is holding belongs to that carrier from then on. + * + * BrokerSubs_OrphanClient performs the hand-off through file-local helpers, so + * the hook below stages the same ownership transfer directly: the queue moves + * to a carrier and the client's own queue fields are cleared, which is what + * stops its teardown from freeing the entries a second time. */ +static MqttBroker* g_handoff_broker; +static BrokerClient* g_handoff_client; +static void handoff_out_queue_mid_write(void) +{ + MqttBroker* broker = g_handoff_broker; + BrokerClient* bc = g_handoff_client; + BrokerOrphanSession* o; + size_t id_len; + + if (broker == NULL || bc == NULL || bc->client_id == NULL || + bc->out_q_head == NULL) { + return; + } + o = (BrokerOrphanSession*)WOLFMQTT_MALLOC(sizeof(*o)); + if (o == NULL) { + return; + } + XMEMSET(o, 0, sizeof(*o)); + id_len = XSTRLEN(bc->client_id); + o->client_id = (char*)WOLFMQTT_MALLOC(id_len + 1); + if (o->client_id == NULL) { + WOLFMQTT_FREE(o); + return; + } + XMEMCPY(o->client_id, bc->client_id, id_len + 1); + o->protocol_level = bc->protocol_level; + o->session_expiry_sec = bc->session_expiry_sec; + o->orphan_since = wolfmqtt_test_broker_time_s(); + + /* Watch the entry the drain is holding, so the test can tell whether the + * drain went on to free memory the carrier owns. */ + broker_test_watch_free(bc->out_q_head); + + o->out_q_head = bc->out_q_head; + o->out_q_tail = bc->out_q_tail; + o->out_q_count = bc->out_q_count; + o->out_q_inflight = bc->out_q_inflight; + bc->out_q_head = NULL; + bc->out_q_tail = NULL; + bc->out_q_count = 0; + bc->out_q_inflight = 0; + + o->next = broker->orphan_sessions; + broker->orphan_sessions = o; + broker->orphan_session_count++; +} + +TEST(drain_stops_when_session_handoff_takes_queue) +{ + MqttBroker broker; + MqttBrokerNet net; + BrokerClient* sub_bc; + BrokerOrphanSession* o; + int i; + /* Publisher "P". */ + static const byte connect_pub[] = { + 0x10, 0x0D, 0x00, 0x04, 'M', 'Q', 'T', 'T', 0x04, 0x02, 0x00, 0x3C, + 0x00, 0x01, 'P' + }; + /* Subscriber "S" with CleanSession=0, so the Session it owns is the kind + * that survives a close and has a carrier to move the queue into. */ + static const byte connect_sub[] = { + 0x10, 0x0D, 0x00, 0x04, 'M', 'Q', 'T', 'T', 0x04, 0x00, 0x00, 0x3C, + 0x00, 0x01, 'S' + }; + /* SUBSCRIBE packet_id=1, filter "x", granted QoS 0: the delivery the + * drain completes and then unlinks. */ + static const byte subscribe_x[] = { + 0x82, 0x06, 0x00, 0x01, 0x00, 0x01, 'x', 0x00 + }; + /* QoS 0 PUBLISH, topic "x", payload "ABC". */ + static const byte publish_x[] = { + 0x30, 0x06, 0x00, 0x01, 'x', 'A', 'B', 'C' + }; + + install_mock_net(&net); + XMEMSET(&broker, 0, sizeof(broker)); + ASSERT_EQ(MQTT_CODE_SUCCESS, MqttBroker_Init(&broker, &net)); + ASSERT_EQ(MQTT_CODE_SUCCESS, MqttBroker_Start(&broker)); + + reset_mock_clients(2); + mock_client_input_append(0, connect_pub, sizeof(connect_pub)); + mock_client_input_append(1, connect_sub, sizeof(connect_sub)); + mock_client_input_append(1, subscribe_x, sizeof(subscribe_x)); + for (i = 0; i < 16; i++) { + (void)MqttBroker_Step(&broker); + } + sub_bc = find_broker_client(&broker, "S"); + ASSERT_NOT_NULL(sub_bc); + ASSERT_NULL(sub_bc->out_q_head); + + broker_test_watch_free(NULL); + g_handoff_broker = &broker; + g_handoff_client = sub_bc; + g_write_publish_hook = handoff_out_queue_mid_write; + + mock_client_input_append(0, publish_x, sizeof(publish_x)); + (void)MqttBroker_Step(&broker); + + /* The hand-off really did run from inside the write. */ + ASSERT_NULL(g_write_publish_hook); + + /* The carrier holds the queue it took. */ + o = broker.orphan_sessions; + ASSERT_NOT_NULL(o); + ASSERT_EQ(1, broker.orphan_session_count); + ASSERT_EQ(1, o->out_q_count); + ASSERT_NOT_NULL(o->out_q_head); + + /* The drain must not free an entry it no longer owns. */ + ASSERT_EQ(0, broker_test_watched_frees()); + + /* And it must leave no trace of that entry on the client. */ + ASSERT_EQ(0, sub_bc->out_q_count); + ASSERT_EQ(0, sub_bc->out_q_inflight); + ASSERT_NULL(sub_bc->out_q_head); + ASSERT_NULL(sub_bc->out_q_tail); + ASSERT_EQ(0, sub_bc->out_q_pending_len); + + g_handoff_broker = NULL; + g_handoff_client = NULL; + broker_test_watch_free(NULL); + MqttBroker_Stop(&broker); + MqttBroker_Free(&broker); +} + #ifdef WOLFMQTT_BROKER_WILL /* Same hazard on the Will fan-out, which runs its own subscription walk. The * publisher's socket drops, the broker fans its Will out to three subscribers, @@ -10516,6 +10672,7 @@ int main(int argc, char** argv) RUN_TEST(online_qos1_flood_disconnects_slow_v311_subscriber); RUN_TEST(online_qos1_at_cap_keeps_subscriber); RUN_TEST(outbound_packet_ids_are_scoped_per_session); + RUN_TEST(drain_stops_when_session_handoff_takes_queue); #ifdef WOLFMQTT_NONBLOCK RUN_TEST(outbound_queue_short_write_resumes_on_next_step); #ifdef WOLFMQTT_BROKER_RETAINED From 83ebb13756cfa14f781cfc9f2b81d62654ba7754 Mon Sep 17 00:00:00 2001 From: Kareem Date: Thu, 24 Sep 2026 16:10:05 -0700 Subject: [PATCH 2/8] wolfMQTT broker: claim the Will before publishing it. Thanks to Gwanhyun Lee for the report. --- src/mqtt_broker.c | 7 +++ tests/test_broker_connect.c | 112 ++++++++++++++++++++++++++++++++++++ 2 files changed, 119 insertions(+) diff --git a/src/mqtt_broker.c b/src/mqtt_broker.c index 045c74fe0..7497025b8 100644 --- a/src/mqtt_broker.c +++ b/src/mqtt_broker.c @@ -5904,6 +5904,13 @@ static void BrokerClient_PublishWill(MqttBroker* broker, BrokerClient* bc) (int)bc->sock, BrokerLog_Sanitize(bc->will_topic), (unsigned)bc->will_payload_len); + /* Claim the Will before fan-out. A fan-out write can service this client's + * close callback inline, which re-enters here; the claim makes that nested + * call return at the has_will check above, so this frame stays the single + * owner of the topic and payload it lends to the encoder below and the + * Will is published once. */ + bc->has_will = 0; + BrokerClient_PublishWillImmediate(broker, bc->will_topic, bc->will_payload, bc->will_payload_len, bc->will_qos, bc->will_retain); diff --git a/tests/test_broker_connect.c b/tests/test_broker_connect.c index 664f91f17..6ec1f54dd 100644 --- a/tests/test_broker_connect.c +++ b/tests/test_broker_connect.c @@ -2209,7 +2209,118 @@ TEST(will_fanout_survives_reentrant_sub_free) MqttBroker_Stop(&broker); MqttBroker_Free(&broker); } + +/* A Will is published by fanning it out to matching subscribers, and only + * cleared once that returns. A write inside the fan-out can service the + * owner's own close callback - the WebSocket transport runs lws_service + * inline - and that callback publishes the Will of the client that just + * closed. It must find nothing left to publish: the fan-out below is still + * lending the topic and payload to the encoder, and a nested publish would + * both deliver the Will a second time and free those buffers underneath it. + * + * The hook samples the ownership flag the nested path tests, from inside the + * first Will delivery. */ +static BrokerClient* g_will_owner; +static int g_will_claimed_during_fanout; +static int g_will_topic_live_during_fanout; +static void sample_will_state_mid_write(void) +{ + if (g_will_owner == NULL) { + return; + } + g_will_claimed_during_fanout = !g_will_owner->has_will; + g_will_topic_live_during_fanout = (g_will_owner->will_topic != NULL); +} + +TEST(takeover_will_claimed_before_fanout) +{ + MqttBroker broker; + MqttBrokerNet net; + int i; + /* Subscribers "A" and "B", CleanSession=1. */ + static const byte connect_a[] = { + 0x10, 0x0D, 0x00, 0x04, 'M', 'Q', 'T', 'T', 0x04, 0x02, 0x00, 0x3C, + 0x00, 0x01, 'A' + }; + static const byte connect_b[] = { + 0x10, 0x0D, 0x00, 0x04, 'M', 'Q', 'T', 'T', 0x04, 0x02, 0x00, 0x3C, + 0x00, 0x01, 'B' + }; + /* SUBSCRIBE packet_id=1, filter "w", QoS 0. */ + static const byte subscribe_w[] = { + 0x82, 0x06, 0x00, 0x01, 0x00, 0x01, 'w', 0x00 + }; + /* v3.1.1 CONNECT for "W" carrying a Will: flags 0x06 = CleanSession + + * Will Flag, Will Topic "w", Will Payload "bye". Remaining Length 21. */ + static const byte connect_w_will[] = { + 0x10, 0x15, 0x00, 0x04, 'M', 'Q', 'T', 'T', 0x04, 0x06, 0x00, 0x3C, + 0x00, 0x01, 'W', + 0x00, 0x01, 'w', + 0x00, 0x03, 'b', 'y', 'e' + }; + /* Second CONNECT with the same ClientId takes the Session over, which + * publishes the old connection's Will. */ + static const byte connect_w_dup[] = { + 0x10, 0x0D, 0x00, 0x04, 'M', 'Q', 'T', 'T', 0x04, 0x02, 0x00, 0x3C, + 0x00, 0x01, 'W' + }; + + install_mock_net(&net); + XMEMSET(&broker, 0, sizeof(broker)); + ASSERT_EQ(MQTT_CODE_SUCCESS, MqttBroker_Init(&broker, &net)); + ASSERT_EQ(MQTT_CODE_SUCCESS, MqttBroker_Start(&broker)); + + reset_mock_clients(3); + mock_client_input_append(0, connect_a, sizeof(connect_a)); + mock_client_input_append(0, subscribe_w, sizeof(subscribe_w)); + mock_client_input_append(1, connect_b, sizeof(connect_b)); + mock_client_input_append(1, subscribe_w, sizeof(subscribe_w)); + mock_client_input_append(2, connect_w_will, sizeof(connect_w_will)); + for (i = 0; i < 24; i++) { + (void)MqttBroker_Step(&broker); + } + + g_will_owner = find_broker_client(&broker, "W"); + ASSERT_NOT_NULL(g_will_owner); + ASSERT_EQ(1, (int)g_will_owner->has_will); + ASSERT_NOT_NULL(g_will_owner->will_topic); + /* Nothing has been published yet, so the hook below cannot be consumed by + * an earlier delivery. */ + ASSERT_EQ(0, count_packets_of_type(g_clients[0].out_buf, + g_clients[0].out_len, MQTT_PACKET_TYPE_PUBLISH)); + + g_will_claimed_during_fanout = -1; + g_will_topic_live_during_fanout = -1; + g_write_publish_hook = sample_will_state_mid_write; + + mock_client_input_append(3, connect_w_dup, sizeof(connect_w_dup)); + g_clients_active = 4; + for (i = 0; i < 24; i++) { + (void)MqttBroker_Step(&broker); + } + + /* The Will really was fanned out, so the sample is meaningful. */ + ASSERT_NULL(g_write_publish_hook); + + /* Claimed before the first write: a close callback delivered during the + * fan-out returns without republishing or freeing anything. */ + ASSERT_EQ(1, g_will_claimed_during_fanout); + /* And the buffers the fan-out is lending out are still owned, so the + * encoder is reading live memory. */ + ASSERT_EQ(1, g_will_topic_live_during_fanout); + + /* Exactly one copy reaches each subscriber [MQTT-3.1.2-8]. */ + ASSERT_EQ(1, count_packets_of_type(g_clients[0].out_buf, + g_clients[0].out_len, MQTT_PACKET_TYPE_PUBLISH)); + ASSERT_EQ(1, count_packets_of_type(g_clients[1].out_buf, + g_clients[1].out_len, MQTT_PACKET_TYPE_PUBLISH)); + + g_will_owner = NULL; + MqttBroker_Stop(&broker); + MqttBroker_Free(&broker); +} #endif /* WOLFMQTT_BROKER_WILL */ + #endif /* !WOLFMQTT_STATIC_MEMORY */ TEST(qos2_duplicate_publish_dedup) @@ -10655,6 +10766,7 @@ int main(int argc, char** argv) RUN_TEST(fanout_survives_reentrant_sub_free); #ifdef WOLFMQTT_BROKER_WILL RUN_TEST(will_fanout_survives_reentrant_sub_free); + RUN_TEST(takeover_will_claimed_before_fanout); #endif #endif RUN_TEST(qos2_duplicate_publish_dedup); From 05abc8050ae2305353346db25fe612c11e588017 Mon Sep 17 00:00:00 2001 From: Kareem Date: Thu, 24 Sep 2026 16:13:00 -0700 Subject: [PATCH 3/8] wolfMQTT client: do not report a successful cancel while a reader owns the request. Thanks to Gwanhyun Lee for the report. --- src/mqtt_client.c | 91 ++++++++++++++++++++++------------------ tests/test_mqtt_client.c | 70 +++++++++++++++++++++++++++++++ wolfmqtt/mqtt_client.h | 6 ++- 3 files changed, 126 insertions(+), 41 deletions(-) diff --git a/src/mqtt_client.c b/src/mqtt_client.c index 1ce10bb4f..c20a630ef 100644 --- a/src/mqtt_client.c +++ b/src/mqtt_client.c @@ -5335,6 +5335,57 @@ int MqttClient_CancelMessage(MqttClient *client, MqttObject* msg) PRINTF("Cancel Msg: %p", msg); #endif +#ifdef WOLFMQTT_MULTITHREAD + /* Remove any pending responses expected. Runs before the resets below so + * that a refusal leaves the message exactly as it was found. + * + * A reading thread claims an entry with packetProcessing while it decodes + * the response into the packet_obj this message owns, having dropped + * lockClient first. Success is the caller's signal to release or reuse the + * object, so it is withheld while that claim stands; the entry stays + * listed and the caller retries until the reader marks it done. */ + rc = wm_SemLock(&client->lockClient); + if (rc != MQTT_CODE_SUCCESS) { + return rc; + } + + for (tmpResp = client->firstPendResp; + tmpResp != NULL; + tmpResp = tmpResp->next) + { + #ifdef WOLFMQTT_DEBUG_CLIENT + PRINTF("\tMsg: %p (obj %p), Type %s (%d), ID %d, InProc %d, Done %d", + tmpResp, tmpResp->packet_obj, + MqttPacket_TypeDesc(tmpResp->packet_type), + tmpResp->packet_type, tmpResp->packet_id, + tmpResp->packetProcessing, tmpResp->packetDone); + #endif + if ((size_t)tmpResp->packet_obj == (size_t)msg || + (size_t)tmpResp - OFFSETOF(MqttMessage, pendResp) == (size_t)msg) { + #ifdef WOLFMQTT_DEBUG_CLIENT + PRINTF("Found Cancel Msg: %p (obj %p), Type %s (%d), ID %d, " + "InProc %d, Done %d", + tmpResp, tmpResp->packet_obj, + MqttPacket_TypeDesc(tmpResp->packet_type), + tmpResp->packet_type, tmpResp->packet_id, + tmpResp->packetProcessing, tmpResp->packetDone); + #endif + if (tmpResp->packetProcessing && !tmpResp->packetDone) { + wm_SemUnlock(&client->lockClient); + return MQTT_CODE_CONTINUE; + } + /* Do not credit any reserved Receive Maximum unit here: the PUBLISH + * may already be on the wire, where the server keeps counting it + * [MQTT-4.9], so crediting it on a local cancel could exceed the + * negotiated quota. The unit is released on the acknowledgement, or + * recovered when the connection resets server_recv_max. */ + MqttClient_RespList_Remove(client, tmpResp); + break; + } + } + wm_SemUnlock(&client->lockClient); +#endif /* WOLFMQTT_MULTITHREAD */ + /* Whether this message's packet finished going out. MQTT_MSG_WAIT is only * reached once the whole Control Packet has been written. */ onWire = (mms_stat->write == MQTT_MSG_WAIT) ? 1 : 0; @@ -5377,46 +5428,6 @@ int MqttClient_CancelMessage(MqttClient *client, MqttObject* msg) mms_stat->recvQuotaHeld = 0; #endif -#ifdef WOLFMQTT_MULTITHREAD - /* Remove any pending responses expected */ - rc = wm_SemLock(&client->lockClient); - if (rc != MQTT_CODE_SUCCESS) { - return rc; - } - - for (tmpResp = client->firstPendResp; - tmpResp != NULL; - tmpResp = tmpResp->next) - { - #ifdef WOLFMQTT_DEBUG_CLIENT - PRINTF("\tMsg: %p (obj %p), Type %s (%d), ID %d, InProc %d, Done %d", - tmpResp, tmpResp->packet_obj, - MqttPacket_TypeDesc(tmpResp->packet_type), - tmpResp->packet_type, tmpResp->packet_id, - tmpResp->packetProcessing, tmpResp->packetDone); - #endif - if ((size_t)tmpResp->packet_obj == (size_t)msg || - (size_t)tmpResp - OFFSETOF(MqttMessage, pendResp) == (size_t)msg) { - #ifdef WOLFMQTT_DEBUG_CLIENT - PRINTF("Found Cancel Msg: %p (obj %p), Type %s (%d), ID %d, " - "InProc %d, Done %d", - tmpResp, tmpResp->packet_obj, - MqttPacket_TypeDesc(tmpResp->packet_type), - tmpResp->packet_type, tmpResp->packet_id, - tmpResp->packetProcessing, tmpResp->packetDone); - #endif - /* Do not credit any reserved Receive Maximum unit here: the PUBLISH - * may already be on the wire, where the server keeps counting it - * [MQTT-4.9], so crediting it on a local cancel could exceed the - * negotiated quota. The unit is released on the acknowledgement, or - * recovered when the connection resets server_recv_max. */ - MqttClient_RespList_Remove(client, tmpResp); - break; - } - } - wm_SemUnlock(&client->lockClient); -#endif /* WOLFMQTT_MULTITHREAD */ - /* cancel any active flags / locks */ if (mms_stat->isReadActive) { #ifdef WOLFMQTT_DEBUG_CLIENT diff --git a/tests/test_mqtt_client.c b/tests/test_mqtt_client.c index 1cea25033..2f4f4c200 100644 --- a/tests/test_mqtt_client.c +++ b/tests/test_mqtt_client.c @@ -1783,6 +1783,72 @@ TEST(cancel_message_retain_is_idempotent) ASSERT_EQ(MQTT_CODE_SUCCESS, rc); ASSERT_EQ(4, test_client.server_recv_max); } + +#if defined(WOLFMQTT_MULTITHREAD) && defined(WOLFMQTT_NONBLOCK) && \ + (WOLFMQTT_MAX_QOS >= 1) +/* A write-only publish leaves its response for another thread to process. That + * reader claims the pending entry with packetProcessing, then drops the client + * lock and decodes the peer's answer into the packet_obj this message owns. + * + * Reporting success to a cancel during that window would be wrong: the caller + * reads success as permission to release or reuse the request, and the decode + * is still writing through it. The claim is honoured instead - the entry stays + * listed, the message is left exactly as it was found, and the caller retries + * until the reader marks the response done. */ +TEST(cancel_refuses_message_while_response_decodes) +{ + int rc; + int i; + int write_state; + /* static so the registered pendResp does not point into freed stack after + * the test returns. */ + static MqttPublish publish; + static byte payload[] = "hello"; + + rc = test_init_client(); + ASSERT_EQ(MQTT_CODE_SUCCESS, rc); + test_client_connect_sent(); + + test_net.write = mock_net_write_accept; + test_net.read = mock_net_read; + + XMEMSET(&publish, 0, sizeof(publish)); + publish.qos = MQTT_QOS_1; + publish.packet_id = 31; + publish.topic_name = "test/topic"; + publish.buffer = payload; + publish.total_len = (word32)(sizeof(payload) - 1); + publish.buffer_len = publish.total_len; + + rc = MQTT_CODE_CONTINUE; + for (i = 0; i < 20 && rc == MQTT_CODE_CONTINUE; i++) { + rc = MqttClient_Publish_WriteOnly(&test_client, &publish, NULL); + } + /* The PUBLISH is out and its acknowledgement is registered for whichever + * thread reads it. */ + ASSERT_TRUE(test_client.firstPendResp == &publish.pendResp); + ASSERT_EQ(0, (int)publish.pendResp.packetProcessing); + write_state = (int)publish.stat.write; + + /* A reader claims the entry, exactly as MqttClient_WaitType does before it + * drops the client lock to decode. */ + publish.pendResp.packetProcessing = 1; + + rc = MqttClient_CancelMessage(&test_client, (MqttObject*)&publish); + ASSERT_EQ(MQTT_CODE_CONTINUE, rc); + /* Still listed, so the reader's pointers stay good, and untouched, so the + * retry starts from the same place. */ + ASSERT_TRUE(test_client.firstPendResp == &publish.pendResp); + ASSERT_EQ(write_state, (int)publish.stat.write); + + /* Once the reader is done with the object the cancel completes. */ + publish.pendResp.packetDone = 1; + rc = MqttClient_CancelMessage(&test_client, (MqttObject*)&publish); + ASSERT_EQ(MQTT_CODE_SUCCESS, rc); + ASSERT_NULL(test_client.firstPendResp); + ASSERT_EQ(MQTT_MSG_BEGIN, (int)publish.stat.write); +} +#endif /* WOLFMQTT_MULTITHREAD && WOLFMQTT_NONBLOCK && MAX_QOS >= 1 */ #endif /* WOLFMQTT_MULTITHREAD || WOLFMQTT_NONBLOCK */ /* A QoS>0 v5 publish that fails on the wire (unsent) must give its reserved @@ -7741,6 +7807,10 @@ void run_mqtt_client_tests(void) RUN_TEST(cancel_message_retains_recv_quota_on_wire); RUN_TEST(cancel_message_retain_is_idempotent); RUN_TEST(cancel_message_reuse_does_not_bypass_recv_quota); +#if defined(WOLFMQTT_MULTITHREAD) && defined(WOLFMQTT_NONBLOCK) && \ + (WOLFMQTT_MAX_QOS >= 1) + RUN_TEST(cancel_refuses_message_while_response_decodes); +#endif #endif RUN_TEST(publish_qos1_v5_write_failure_restores_recv_quota); #if defined(WOLFMQTT_MULTITHREAD) && defined(WOLFMQTT_NONBLOCK) && \ diff --git a/wolfmqtt/mqtt_client.h b/wolfmqtt/mqtt_client.h index 562019a0c..1ee28a59a 100644 --- a/wolfmqtt/mqtt_client.h +++ b/wolfmqtt/mqtt_client.h @@ -888,9 +888,13 @@ WOLFMQTT_API int MqttClient_WaitMessage_ex( /*! \brief In a multi-threaded and non-blocking mode this allows you to cancel an MQTT object that was previously submitted. * \note This is a blocking function that will wait for MqttNet.read + * \note Only MQTT_CODE_SUCCESS releases the object back to the caller. + MQTT_CODE_CONTINUE means another thread is still processing + this message's response: the object must not be freed or + reused, and the call should be retried. * \param client Pointer to MqttClient structure * \param msg Pointer to MqttObject structure - * \return MQTT_CODE_SUCCESS or MQTT_CODE_ERROR_* + * \return MQTT_CODE_SUCCESS, MQTT_CODE_CONTINUE or MQTT_CODE_ERROR_* (see enum MqttPacketResponseCodes) */ WOLFMQTT_API int MqttClient_CancelMessage( From 333e5cc47c02da86841b0dae3bac8f6d28d93391 Mon Sep 17 00:00:00 2001 From: Kareem Date: Thu, 24 Sep 2026 16:26:18 -0700 Subject: [PATCH 4/8] wolfMQTT client: take the client lock over the Session replay pool in MqttClient_Connect. Thanks to Gwanhyun Lee for the report. --- src/mqtt_client.c | 75 +++++++++++++++++++++++++----- tests/include.am | 4 +- tests/test_mqtt_client.c | 98 ++++++++++++++++++++++++++++++++++++++++ 3 files changed, 164 insertions(+), 13 deletions(-) diff --git a/src/mqtt_client.c b/src/mqtt_client.c index c20a630ef..5472fb14e 100644 --- a/src/mqtt_client.c +++ b/src/mqtt_client.c @@ -797,6 +797,24 @@ static void MqttClient_Replay_AddSafe(MqttClient* client, MqttPublish* publish) #endif } +/* Unlike the other wrappers this one reports a lock failure. Skipping the + * reset would leave the previous Session's messages retained, and a later + * reconnect would replay them into a Session they do not belong to. */ +static int MqttClient_Replay_ResetSafe(MqttClient* client) +{ +#ifdef WOLFMQTT_MULTITHREAD + int rc = wm_SemLock(&client->lockClient); + if (rc != MQTT_CODE_SUCCESS) { + return rc; + } +#endif + MqttClient_Replay_Reset(client); +#ifdef WOLFMQTT_MULTITHREAD + wm_SemUnlock(&client->lockClient); +#endif + return MQTT_CODE_SUCCESS; +} + static void MqttClient_Replay_RemoveSafe(MqttClient* client, word16 packet_id) { #ifdef WOLFMQTT_MULTITHREAD @@ -861,8 +879,12 @@ static int MqttClient_SendIds_Find(const MqttClient* client, word16 packet_id) * MQTT_CODE_SUCCESS when it was free (or when isRetransmit says this is a * re-send of the same Control Packet, which [MQTT-2.3.1-3] requires to keep * its original identifier), and MQTT_CODE_ERROR_PACKET_ID when it is still - * awaiting its acknowledgement. */ -static int MqttClient_SendIdReserve(MqttClient* client, word16 packet_id, + * awaiting its acknowledgement. + * + * The _Locked form is for a caller that already holds client->lockClient - the + * Session resume walks the replay pool under it - since the semaphore is not + * recursive. */ +static int MqttClient_SendIdReserve_Locked(MqttClient* client, word16 packet_id, void* owner, int isRetransmit, MqttPacketType ack_type) { int rc = MQTT_CODE_SUCCESS; @@ -871,12 +893,6 @@ static int MqttClient_SendIdReserve(MqttClient* client, word16 packet_id, if (packet_id == 0) { return MQTT_CODE_SUCCESS; /* nothing to track */ } -#ifdef WOLFMQTT_MULTITHREAD - rc = wm_SemLock(&client->lockClient); - if (rc != MQTT_CODE_SUCCESS) { - return rc; - } -#endif i = MqttClient_SendIds_Find(client, packet_id); if (i >= 0) { /* Only a re-send of the same Control Packet may keep an identifier @@ -906,6 +922,25 @@ static int MqttClient_SendIdReserve(MqttClient* client, word16 packet_id, * packets in flight, so the check is best effort past that point - * raise MQTT_MAX_SEND_INFLIGHT to widen the window. */ } + return rc; +} + +static int MqttClient_SendIdReserve(MqttClient* client, word16 packet_id, + void* owner, int isRetransmit, MqttPacketType ack_type) +{ + int rc; + + if (packet_id == 0) { + return MQTT_CODE_SUCCESS; /* nothing to track */ + } +#ifdef WOLFMQTT_MULTITHREAD + rc = wm_SemLock(&client->lockClient); + if (rc != MQTT_CODE_SUCCESS) { + return rc; + } +#endif + rc = MqttClient_SendIdReserve_Locked(client, packet_id, owner, + isRetransmit, ack_type); #ifdef WOLFMQTT_MULTITHREAD wm_SemUnlock(&client->lockClient); #endif @@ -3594,8 +3629,14 @@ int MqttClient_Connect(MqttClient *client, MqttConnect *mc_connect) * into it would inject one Session's messages into another. */ if (!(mc_connect->ack.flags & MQTT_CONNECT_ACK_FLAG_SESSION_PRESENT) || !session_id_matched) { - /* A fresh Session starts with no outbound state to re-send. */ - MqttClient_Replay_Reset(client); + /* A fresh Session starts with no outbound state to re-send. + * Under the client lock: a send thread may be replacing a slot's + * topic and payload through MqttClient_Replay_AddSafe at the same + * moment, and both sides free what they find there. */ + rc = MqttClient_Replay_ResetSafe(client); + if (rc != MQTT_CODE_SUCCESS) { + return rc; + } } else { int i; @@ -3603,7 +3644,14 @@ int MqttClient_Connect(MqttClient *client, MqttConnect *mc_connect) /* The server resumed the Session, so these messages are still in * flight as far as it is concerned: keep their Packet Identifiers * reserved (MqttClient_Connect cleared the table above) and - * re-send them [MQTT-4.4.0-1]. */ + * re-send them [MQTT-4.4.0-1]. Walked under the client lock for + * the same reason as the reset above. */ +#ifdef WOLFMQTT_MULTITHREAD + rc = wm_SemLock(&client->lockClient); + if (rc != MQTT_CODE_SUCCESS) { + return rc; + } +#endif for (i = 0; i < MQTT_MAX_REPLAY_MSGS; i++) { if (client->replay[i].packet_id == 0) { continue; @@ -3617,7 +3665,7 @@ int MqttClient_Connect(MqttClient *client, MqttConnect *mc_connect) MqttClient_Replay_FreeSlot(&client->replay[i]); continue; } - (void)MqttClient_SendIdReserve(client, + (void)MqttClient_SendIdReserve_Locked(client, client->replay[i].packet_id, &client->replay[i], 1, client->replay[i].pubrelSent ? MQTT_PACKET_TYPE_PUBLISH_COMP : @@ -3625,6 +3673,9 @@ int MqttClient_Connect(MqttClient *client, MqttConnect *mc_connect) MQTT_PACKET_TYPE_PUBLISH_COMP : MQTT_PACKET_TYPE_PUBLISH_ACK)); } +#ifdef WOLFMQTT_MULTITHREAD + wm_SemUnlock(&client->lockClient); +#endif client->replayIdx = 0; mc_connect->stat.write = MQTT_MSG_PAYLOAD; rc = MqttClient_ReplaySession(client, mc_connect); diff --git a/tests/include.am b/tests/include.am index 9d79bef16..fd4143112 100644 --- a/tests/include.am +++ b/tests/include.am @@ -14,7 +14,8 @@ tests_unit_tests_SOURCES = \ src/mqtt_socket.c \ src/mqtt_sn_client.c \ src/mqtt_sn_packet.c -tests_unit_tests_CPPFLAGS = -I$(top_srcdir) $(AM_CPPFLAGS) +tests_unit_tests_CPPFLAGS = -I$(top_srcdir) \ + -include $(top_srcdir)/tests/test_client_alloc.h $(AM_CPPFLAGS) tests_unit_tests_LDADD = $(PTHREAD_LIBS) # Exercise peer-name verification with a trusted local certificate and an @@ -54,6 +55,7 @@ endif # Test framework headers noinst_HEADERS += tests/unit_test.h tests/test_broker_alloc.h \ + tests/test_client_alloc.h \ tests/test_mqtt_props_config.h tests/test_mqtt_curl_config.h # MQTT-SN packet decoder tests. The SN_Decode_* routines are WOLFMQTT_LOCAL diff --git a/tests/test_mqtt_client.c b/tests/test_mqtt_client.c index 2f4f4c200..3c70ec1eb 100644 --- a/tests/test_mqtt_client.c +++ b/tests/test_mqtt_client.c @@ -48,6 +48,47 @@ static MqttNet test_net; static byte test_tx_buf[TEST_TX_BUF_SIZE]; static byte test_rx_buf[TEST_RX_BUF_SIZE]; +/* The suite routes the library's allocator through here (see + * tests/test_client_alloc.h) so a test can watch one specific allocation and + * record the state the library was in when it released it. Everything else is + * a straight pass-through. */ +static void* g_free_watch_ptr; +static int g_free_watch_count; +static int g_free_watch_lock_count; + +void* wolfmqtt_test_client_malloc(size_t size) +{ + return malloc(size); +} + +/* How many times the client lock is held right now, or -1 where the build + * keeps no count to read. */ +static int test_client_lock_depth(void) +{ +#if defined(WOLFMQTT_MULTITHREAD) && defined(WOLFMQTT_POSIX_SEMAPHORES) && \ + !defined(WOLFMQTT_NO_COND_SIGNAL) + return test_client.lockClient.lockCount; +#else + return -1; +#endif +} + +void wolfmqtt_test_client_free(void* ptr) +{ + if (ptr != NULL && ptr == g_free_watch_ptr) { + g_free_watch_count++; + g_free_watch_lock_count = test_client_lock_depth(); + } + free(ptr); +} + +UT_MAYBE_UNUSED static void test_client_watch_free(void* ptr) +{ + g_free_watch_ptr = ptr; + g_free_watch_count = 0; + g_free_watch_lock_count = -1; +} + /* Mock network callbacks - just return errors since we're not actually * connecting to anything */ static int mock_net_connect(void *context, const char* host, word16 port, @@ -3372,6 +3413,59 @@ TEST(reconnect_without_session_present_replays_nothing) ASSERT_EQ(1, g_frames_written); /* the CONNECT only */ } +#if defined(WOLFMQTT_MULTITHREAD) && defined(WOLFMQTT_POSIX_SEMAPHORES) && \ + !defined(WOLFMQTT_NO_COND_SIGNAL) +/* Discarding the retained copies releases the topic and payload of every + * slot. A send thread reaches those same slots through the replay store's + * lock-taking wrapper while the connection is being re-established - the + * client admits send calls once CONNECT has reached the transport - and both + * sides free what they find there. So the discard has to hold the client lock + * as well, or the two can free the same buffer or leave one pointing at + * memory the other released. + * + * The POSIX lock keeps its own depth, and the allocator hook samples it at the + * moment the slot's topic is released, which is the access that has to be + * covered. */ +TEST(fresh_session_reset_frees_replay_under_client_lock) +{ + int rc; + MqttConnect connect; + MqttPublish publish; + static byte payload[] = "hello"; + + rc = test_init_client(); + ASSERT_EQ(MQTT_CODE_SUCCESS, rc); +#ifdef WOLFMQTT_V5 + test_client.protocol_level = MQTT_CONNECT_PROTOCOL_LEVEL_4; +#endif + ASSERT_EQ(MQTT_CODE_SUCCESS, run_initial_connect(&connect)); + + init_qos_publish(&publish, MQTT_QOS_1, 0x1234, "sensor/temp", + payload, (word32)(sizeof(payload) - 1)); + rc = run_publish_unacked(&publish); + ASSERT_EQ(MQTT_CODE_ERROR_NETWORK, rc); + ASSERT_EQ(MQTT_CODE_SUCCESS, MqttClient_NetDisconnect(&test_client)); + + /* The unacknowledged PUBLISH is retained, so the discard has a real + * allocation to release. */ + ASSERT_EQ(0x1234, (int)test_client.replay[0].packet_id); + ASSERT_NOT_NULL(test_client.replay[0].topic); + test_client_watch_free(test_client.replay[0].topic); + + /* Session Present = 0 is a different Session, so nothing is carried. */ + rc = run_reconnect(&connect, 0); + ASSERT_EQ(MQTT_CODE_SUCCESS, rc); + ASSERT_EQ(0, (int)test_client.replay[0].packet_id); + + /* It really was released on this path, exactly once. */ + ASSERT_EQ(1, g_free_watch_count); + /* And under the client lock. */ + ASSERT_EQ(1, g_free_watch_lock_count); + + test_client_watch_free(NULL); +} +#endif + /* A completed exchange leaves Session state, so an acknowledged PUBLISH is * not replayed [MQTT-4.4.0-1] covers unacknowledged messages only. */ TEST(reconnect_does_not_replay_acked_publish) @@ -7706,6 +7800,10 @@ void run_mqtt_client_tests(void) #if !defined(WOLFMQTT_NO_SESSION_REPLAY) && WOLFMQTT_MAX_QOS >= 1 RUN_TEST(reconnect_replays_unacked_qos1_publish); RUN_TEST(reconnect_without_session_present_replays_nothing); +#if defined(WOLFMQTT_MULTITHREAD) && defined(WOLFMQTT_POSIX_SEMAPHORES) && \ + !defined(WOLFMQTT_NO_COND_SIGNAL) + RUN_TEST(fresh_session_reset_frees_replay_under_client_lock); +#endif RUN_TEST(reconnect_does_not_replay_acked_publish); #if WOLFMQTT_MAX_QOS >= 2 RUN_TEST(reconnect_replays_unacked_pubrel); From ef23354977bafa22951926d47d2f9d33207fabf6 Mon Sep 17 00:00:00 2001 From: Kareem Date: Fri, 2 Oct 2026 12:56:54 -0700 Subject: [PATCH 5/8] Address code review feedback and minimize comments. --- src/mqtt_broker.c | 86 +++++++++++++++++---- src/mqtt_client.c | 50 +++++++++---- tests/include.am | 3 +- tests/test_broker_connect.c | 144 ++++++++++++++++++++++++++---------- tests/test_mqtt_client.c | 142 ++++++++++++++++++++++++++--------- wolfmqtt/mqtt_client.h | 7 +- 6 files changed, 322 insertions(+), 110 deletions(-) diff --git a/src/mqtt_broker.c b/src/mqtt_broker.c index 7497025b8..65a5cabd3 100644 --- a/src/mqtt_broker.c +++ b/src/mqtt_broker.c @@ -1999,17 +1999,30 @@ static void BrokerClient_FreeOutQueue(BrokerClient* bc) bc->out_q_pending_len = 0; } -/* A write can service this same client's close callback before it returns - - * the WebSocket transport runs lws_service inline - and that hands the whole - * outbound queue to the Session carrier. Entries the drain is walking were - * reached from bc->out_q_head, so an empty head while one is still held means - * the carrier owns them now: they must not be unlinked, freed, or re-linked - * into bc. Returns non-zero when the drain has to stop for that reason. */ -static int BrokerClient_OutQueueMoved(BrokerClient* bc, int enc_len) +#ifndef WOLFMQTT_STATIC_MEMORY +static void BrokerOrphan_ReconcileSent(MqttBroker* broker, + const char* client_id, BrokerOutPub* sent); +#endif + +/* A write can service this client's close callback before returning - the + * WebSocket transport runs lws_service inline - which hands the queue to the + * Session carrier. The walked entries came from bc->out_q_head, so an empty + * head while one is still held means the carrier owns them now. Returns + * non-zero when the drain must stop and leave them to it. */ +static int BrokerClient_OutQueueMoved(BrokerClient* bc, BrokerOutPub* sent, + int completed, int enc_len) { if (bc->out_q_head != NULL) { return 0; } +#ifndef WOLFMQTT_STATIC_MEMORY + if (completed) { + BrokerOrphan_ReconcileSent(bc->broker, bc->client_id, sent); + } +#else + (void)sent; + (void)completed; +#endif bc->out_q_pending_len = 0; BROKER_FORCE_ZERO(bc->tx_buf, enc_len); return 1; @@ -2075,7 +2088,8 @@ static int BrokerClient_DrainOutQueue(BrokerClient* bc) return MQTT_CODE_ERROR_SYSTEM; } wr_rc = MqttPacket_Write(&bc->client, bc->tx_buf, rel_rc); - if (BrokerClient_OutQueueMoved(bc, rel_rc)) { + if (BrokerClient_OutQueueMoved(bc, cur, + wr_rc == rel_rc, rel_rc)) { return sent; } if (wr_rc == MQTT_CODE_CONTINUE) { @@ -2171,7 +2185,8 @@ static int BrokerClient_DrainOutQueue(BrokerClient* bc) { int wr_rc; wr_rc = MqttPacket_Write(&bc->client, bc->tx_buf, enc_rc); - if (BrokerClient_OutQueueMoved(bc, enc_rc)) { + if (BrokerClient_OutQueueMoved(bc, cur, + wr_rc == enc_rc, enc_rc)) { return sent; } /* Scrub the forwarded PUBLISH (which may carry a replayed will or @@ -3484,6 +3499,52 @@ static BrokerOrphanSession* BrokerOrphan_Find(MqttBroker* broker, return NULL; } +/* Account for a write that completed after the hand-off moved the queue here. + * A delivered QoS 0 entry is spent - MQTT 3.1.1 section 4.3.1 allows no retry - + * so retire it rather than let the resume send it again. A QoS > 0 entry is + * still unacknowledged and stays, but its resume copy is a re-delivery and + * needs DUP [MQTT-3.3.1-1]. Found by walking the carrier, so an entry the + * hand-off dropped on its way in is simply absent. */ +static void BrokerOrphan_ReconcileSent(MqttBroker* broker, + const char* client_id, BrokerOutPub* sent) +{ + BrokerOrphanSession* o; + BrokerOutPub* prev = NULL; + BrokerOutPub* cur; + + if (sent == NULL || !BROKER_STR_VALID(client_id)) { + return; + } + o = BrokerOrphan_Find(broker, client_id); + if (o == NULL) { + return; + } + for (cur = o->out_q_head; cur != NULL; prev = cur, cur = cur->next) { + if (cur != sent) { + continue; + } + if (cur->qos != MQTT_QOS_0) { + cur->retransmit_dup = 1; + return; + } + if (prev == NULL) { + o->out_q_head = cur->next; + } + else { + prev->next = cur->next; + } + if (o->out_q_tail == cur) { + o->out_q_tail = prev; + } + o->out_q_count--; + WBLOG_DBG(broker, + "broker: retiring delivered qos0 topic=%s client_id=%s", + BrokerLog_Sanitize(cur->topic), BrokerLog_Sanitize(client_id)); + BrokerOutPub_Free(cur); + return; + } +} + /* Free everything an orphan owns (queue entries + client_id) but do * NOT unlink from broker->orphan_sessions; the caller does that. */ static void BrokerOrphan_FreeContents(BrokerOrphanSession* o) @@ -5904,11 +5965,8 @@ static void BrokerClient_PublishWill(MqttBroker* broker, BrokerClient* bc) (int)bc->sock, BrokerLog_Sanitize(bc->will_topic), (unsigned)bc->will_payload_len); - /* Claim the Will before fan-out. A fan-out write can service this client's - * close callback inline, which re-enters here; the claim makes that nested - * call return at the has_will check above, so this frame stays the single - * owner of the topic and payload it lends to the encoder below and the - * Will is published once. */ + /* Claimed before fan-out, whose writes can service this client's close + * callback and re-enter here; that nested call then returns above. */ bc->has_will = 0; BrokerClient_PublishWillImmediate(broker, bc->will_topic, diff --git a/src/mqtt_client.c b/src/mqtt_client.c index 5472fb14e..2c732f4c7 100644 --- a/src/mqtt_client.c +++ b/src/mqtt_client.c @@ -88,6 +88,9 @@ static int MqttClient_AuthEx(MqttClient *client, MqttAuth* auth, #if !defined(WOLFMQTT_MULTITHREAD) && !defined(WOLFMQTT_NONBLOCK) static int MqttClient_CancelMessage(MqttClient *client, MqttObject* msg); #endif +#ifndef WOLFMQTT_NO_SESSION_REPLAY +static int MqttClient_SendIds_Find(const MqttClient* client, word16 packet_id); +#endif #ifdef WOLFMQTT_MULTITHREAD #ifdef WOLFMQTT_USER_THREADING @@ -797,18 +800,31 @@ static void MqttClient_Replay_AddSafe(MqttClient* client, MqttPublish* publish) #endif } -/* Unlike the other wrappers this one reports a lock failure. Skipping the - * reset would leave the previous Session's messages retained, and a later - * reconnect would replay them into a Session they do not belong to. */ +/* Discard the retained copies of a Session the server did not resume. A send + * is allowed once CONNECT reaches the transport, so an entry may instead belong + * to the Session being established: MqttClient_Connect empties the Packet + * Identifier table before sending CONNECT, so a still-reserved identifier marks + * one of those, and it is kept for a later resume [MQTT-4.4.0-1]. Reports a + * lock failure, unlike the wrappers above: a skipped discard would replay the + * previous Session's messages into this one. */ static int MqttClient_Replay_ResetSafe(MqttClient* client) { + int i; #ifdef WOLFMQTT_MULTITHREAD int rc = wm_SemLock(&client->lockClient); if (rc != MQTT_CODE_SUCCESS) { return rc; } #endif - MqttClient_Replay_Reset(client); + for (i = 0; i < MQTT_MAX_REPLAY_MSGS; i++) { + if (client->replay[i].packet_id != 0 && + MqttClient_SendIds_Find(client, + client->replay[i].packet_id) >= 0) { + continue; /* published on this connection */ + } + MqttClient_Replay_FreeSlot(&client->replay[i]); + } + client->replayIdx = MQTT_MAX_REPLAY_MSGS; #ifdef WOLFMQTT_MULTITHREAD wm_SemUnlock(&client->lockClient); #endif @@ -879,11 +895,8 @@ static int MqttClient_SendIds_Find(const MqttClient* client, word16 packet_id) * MQTT_CODE_SUCCESS when it was free (or when isRetransmit says this is a * re-send of the same Control Packet, which [MQTT-2.3.1-3] requires to keep * its original identifier), and MQTT_CODE_ERROR_PACKET_ID when it is still - * awaiting its acknowledgement. - * - * The _Locked form is for a caller that already holds client->lockClient - the - * Session resume walks the replay pool under it - since the semaphore is not - * recursive. */ + * awaiting its acknowledgement. The _Locked form is for a caller already + * holding client->lockClient, which is not recursive. */ static int MqttClient_SendIdReserve_Locked(MqttClient* client, word16 packet_id, void* owner, int isRetransmit, MqttPacketType ack_type) { @@ -2400,6 +2413,16 @@ static int MqttClient_WaitType(MqttClient *client, void *packet_obj, rc = MQTT_CODE_SUCCESS; } else { + #ifdef WOLFMQTT_MULTITHREAD + /* Terminal, so the claim must not outlive this reader. The + * CONTINUE path above keeps it: that reader resumes. */ + if (pendResp != NULL) { + if (wm_SemLock(&client->lockClient) == 0) { + pendResp->packetProcessing = 0; + wm_SemUnlock(&client->lockClient); + } + } + #endif /* error, break */ break; } @@ -3629,10 +3652,7 @@ int MqttClient_Connect(MqttClient *client, MqttConnect *mc_connect) * into it would inject one Session's messages into another. */ if (!(mc_connect->ack.flags & MQTT_CONNECT_ACK_FLAG_SESSION_PRESENT) || !session_id_matched) { - /* A fresh Session starts with no outbound state to re-send. - * Under the client lock: a send thread may be replacing a slot's - * topic and payload through MqttClient_Replay_AddSafe at the same - * moment, and both sides free what they find there. */ + /* A fresh Session starts with no outbound state to re-send. */ rc = MqttClient_Replay_ResetSafe(client); if (rc != MQTT_CODE_SUCCESS) { return rc; @@ -3644,8 +3664,8 @@ int MqttClient_Connect(MqttClient *client, MqttConnect *mc_connect) /* The server resumed the Session, so these messages are still in * flight as far as it is concerned: keep their Packet Identifiers * reserved (MqttClient_Connect cleared the table above) and - * re-send them [MQTT-4.4.0-1]. Walked under the client lock for - * the same reason as the reset above. */ + * re-send them [MQTT-4.4.0-1]. A send thread reaches these same + * slots, so the walk holds the client lock. */ #ifdef WOLFMQTT_MULTITHREAD rc = wm_SemLock(&client->lockClient); if (rc != MQTT_CODE_SUCCESS) { diff --git a/tests/include.am b/tests/include.am index fd4143112..69577694c 100644 --- a/tests/include.am +++ b/tests/include.am @@ -15,7 +15,7 @@ tests_unit_tests_SOURCES = \ src/mqtt_sn_client.c \ src/mqtt_sn_packet.c tests_unit_tests_CPPFLAGS = -I$(top_srcdir) \ - -include $(top_srcdir)/tests/test_client_alloc.h $(AM_CPPFLAGS) + -include $(top_srcdir)/tests/test_broker_alloc.h $(AM_CPPFLAGS) tests_unit_tests_LDADD = $(PTHREAD_LIBS) # Exercise peer-name verification with a trusted local certificate and an @@ -55,7 +55,6 @@ endif # Test framework headers noinst_HEADERS += tests/unit_test.h tests/test_broker_alloc.h \ - tests/test_client_alloc.h \ tests/test_mqtt_props_config.h tests/test_mqtt_curl_config.h # MQTT-SN packet decoder tests. The SN_Decode_* routines are WOLFMQTT_LOCAL diff --git a/tests/test_broker_connect.c b/tests/test_broker_connect.c index 6ec1f54dd..7008d97b8 100644 --- a/tests/test_broker_connect.c +++ b/tests/test_broker_connect.c @@ -114,8 +114,8 @@ unsigned long wolfmqtt_test_broker_time_s(void) } #ifndef WOLFMQTT_STATIC_MEMORY -/* Count frees of one specific allocation, so a test can assert on the fate of - * a single object rather than on a total that unrelated activity moves. */ +/* Count frees of one allocation, so a test can assert on a single object + * rather than a total that unrelated activity moves. */ static void broker_test_watch_free(void* ptr) { g_free_watch_ptr = ptr; @@ -1999,16 +1999,10 @@ TEST(fanout_survives_reentrant_sub_free) MqttBroker_Free(&broker); } -/* The outbound drain holds a queue entry across the write that sends it. That - * write can service this client's own close callback before it returns - the - * WebSocket transport runs lws_service inline - and the callback hands the - * whole queue to the carrier that keeps the Session alive while the client is - * gone. The entry the drain is holding belongs to that carrier from then on. - * - * BrokerSubs_OrphanClient performs the hand-off through file-local helpers, so - * the hook below stages the same ownership transfer directly: the queue moves - * to a carrier and the client's own queue fields are cleared, which is what - * stops its teardown from freeing the entries a second time. */ +/* The drain holds a queue entry across the write that sends it, and that write + * can service this client's close callback, which hands the whole queue to the + * carrier holding the Session. BrokerSubs_OrphanClient does that through + * file-local helpers, so the hook stages the same transfer directly. */ static MqttBroker* g_handoff_broker; static BrokerClient* g_handoff_client; static void handoff_out_queue_mid_write(void) @@ -2038,8 +2032,7 @@ static void handoff_out_queue_mid_write(void) o->session_expiry_sec = bc->session_expiry_sec; o->orphan_since = wolfmqtt_test_broker_time_s(); - /* Watch the entry the drain is holding, so the test can tell whether the - * drain went on to free memory the carrier owns. */ + /* Watch the entry the drain holds, to see if it frees the carrier's. */ broker_test_watch_free(bc->out_q_head); o->out_q_head = bc->out_q_head; @@ -2068,14 +2061,12 @@ TEST(drain_stops_when_session_handoff_takes_queue) 0x10, 0x0D, 0x00, 0x04, 'M', 'Q', 'T', 'T', 0x04, 0x02, 0x00, 0x3C, 0x00, 0x01, 'P' }; - /* Subscriber "S" with CleanSession=0, so the Session it owns is the kind - * that survives a close and has a carrier to move the queue into. */ + /* CleanSession=0, so the Session survives and has a carrier. */ static const byte connect_sub[] = { 0x10, 0x0D, 0x00, 0x04, 'M', 'Q', 'T', 'T', 0x04, 0x00, 0x00, 0x3C, 0x00, 0x01, 'S' }; - /* SUBSCRIBE packet_id=1, filter "x", granted QoS 0: the delivery the - * drain completes and then unlinks. */ + /* SUBSCRIBE packet_id=1, filter "x", granted QoS 0. */ static const byte subscribe_x[] = { 0x82, 0x06, 0x00, 0x01, 0x00, 0x01, 'x', 0x00 }; @@ -2108,20 +2099,26 @@ TEST(drain_stops_when_session_handoff_takes_queue) mock_client_input_append(0, publish_x, sizeof(publish_x)); (void)MqttBroker_Step(&broker); - /* The hand-off really did run from inside the write. */ + /* The hand-off ran from inside the write. */ ASSERT_NULL(g_write_publish_hook); + /* The delivery completed, which is what makes the entry spent. */ + ASSERT_EQ(1, count_packets_of_type(g_clients[1].out_buf, + g_clients[1].out_len, MQTT_PACKET_TYPE_PUBLISH)); - /* The carrier holds the queue it took. */ + /* The carrier took the queue, and the delivered QoS 0 entry is retired + * from it: MQTT 3.1.1 section 4.3.1 allows no retry. */ o = broker.orphan_sessions; ASSERT_NOT_NULL(o); ASSERT_EQ(1, broker.orphan_session_count); - ASSERT_EQ(1, o->out_q_count); - ASSERT_NOT_NULL(o->out_q_head); + ASSERT_EQ(0, o->out_q_count); + ASSERT_NULL(o->out_q_head); + ASSERT_NULL(o->out_q_tail); - /* The drain must not free an entry it no longer owns. */ - ASSERT_EQ(0, broker_test_watched_frees()); + /* Released once by the owner that retired it: more would mean the drain + * freed a carrier entry, none would leak it. */ + ASSERT_EQ(1, broker_test_watched_frees()); - /* And it must leave no trace of that entry on the client. */ + /* And no trace of it left on the client. */ ASSERT_EQ(0, sub_bc->out_q_count); ASSERT_EQ(0, sub_bc->out_q_inflight); ASSERT_NULL(sub_bc->out_q_head); @@ -2135,6 +2132,82 @@ TEST(drain_stops_when_session_handoff_takes_queue) MqttBroker_Free(&broker); } +/* The QoS > 0 half. A completed at-least-once delivery is still unacknowledged + * so it stays with the carrier, but the subscriber has seen it, so the resume + * copy must be marked a duplicate [MQTT-3.3.1-1]. */ +TEST(handoff_marks_sent_qos1_for_redelivery) +{ + MqttBroker broker; + MqttBrokerNet net; + BrokerClient* sub_bc; + BrokerOrphanSession* o; + int i; + static const byte connect_pub[] = { + 0x10, 0x0D, 0x00, 0x04, 'M', 'Q', 'T', 'T', 0x04, 0x02, 0x00, 0x3C, + 0x00, 0x01, 'P' + }; + /* CleanSession=0, so the Session survives and has a carrier. */ + static const byte connect_sub[] = { + 0x10, 0x0D, 0x00, 0x04, 'M', 'Q', 'T', 'T', 0x04, 0x00, 0x00, 0x3C, + 0x00, 0x01, 'S' + }; + /* SUBSCRIBE packet_id=1, filter "x", granted QoS 1. */ + static const byte subscribe_x[] = { + 0x82, 0x06, 0x00, 0x01, 0x00, 0x01, 'x', 0x01 + }; + /* QoS 1 PUBLISH, Packet Identifier 7, topic "x", payload "ABC". */ + static const byte publish_x[] = { + 0x32, 0x08, 0x00, 0x01, 'x', 0x00, 0x07, 'A', 'B', 'C' + }; + + install_mock_net(&net); + XMEMSET(&broker, 0, sizeof(broker)); + ASSERT_EQ(MQTT_CODE_SUCCESS, MqttBroker_Init(&broker, &net)); + ASSERT_EQ(MQTT_CODE_SUCCESS, MqttBroker_Start(&broker)); + + reset_mock_clients(2); + mock_client_input_append(0, connect_pub, sizeof(connect_pub)); + mock_client_input_append(1, connect_sub, sizeof(connect_sub)); + mock_client_input_append(1, subscribe_x, sizeof(subscribe_x)); + for (i = 0; i < 16; i++) { + (void)MqttBroker_Step(&broker); + } + sub_bc = find_broker_client(&broker, "S"); + ASSERT_NOT_NULL(sub_bc); + ASSERT_NULL(sub_bc->out_q_head); + + broker_test_watch_free(NULL); + g_handoff_broker = &broker; + g_handoff_client = sub_bc; + g_write_publish_hook = handoff_out_queue_mid_write; + + mock_client_input_append(0, publish_x, sizeof(publish_x)); + (void)MqttBroker_Step(&broker); + + ASSERT_NULL(g_write_publish_hook); + ASSERT_EQ(1, count_packets_of_type(g_clients[1].out_buf, + g_clients[1].out_len, MQTT_PACKET_TYPE_PUBLISH)); + + /* Still the carrier's, awaiting its PUBACK, and nothing freed. */ + o = broker.orphan_sessions; + ASSERT_NOT_NULL(o); + ASSERT_EQ(1, o->out_q_count); + ASSERT_NOT_NULL(o->out_q_head); + ASSERT_EQ(0, broker_test_watched_frees()); + /* Flagged, so the resume sends it with DUP set. */ + ASSERT_EQ(1, (int)o->out_q_head->retransmit_dup); + ASSERT_EQ(MQTT_QOS_1, (int)o->out_q_head->qos); + + ASSERT_EQ(0, sub_bc->out_q_count); + ASSERT_NULL(sub_bc->out_q_head); + + g_handoff_broker = NULL; + g_handoff_client = NULL; + broker_test_watch_free(NULL); + MqttBroker_Stop(&broker); + MqttBroker_Free(&broker); +} + #ifdef WOLFMQTT_BROKER_WILL /* Same hazard on the Will fan-out, which runs its own subscription walk. The * publisher's socket drops, the broker fans its Will out to three subscribers, @@ -2210,16 +2283,10 @@ TEST(will_fanout_survives_reentrant_sub_free) MqttBroker_Free(&broker); } -/* A Will is published by fanning it out to matching subscribers, and only - * cleared once that returns. A write inside the fan-out can service the - * owner's own close callback - the WebSocket transport runs lws_service - * inline - and that callback publishes the Will of the client that just - * closed. It must find nothing left to publish: the fan-out below is still - * lending the topic and payload to the encoder, and a nested publish would - * both deliver the Will a second time and free those buffers underneath it. - * - * The hook samples the ownership flag the nested path tests, from inside the - * first Will delivery. */ +/* The Will is fanned out before it is cleared, and a write inside that fan-out + * can service the owner's own close callback, which publishes the same Will. + * It must find nothing left to publish [MQTT-3.1.2-8]. The hook samples the + * flag that nested path tests, from inside the first delivery. */ static BrokerClient* g_will_owner; static int g_will_claimed_during_fanout; static int g_will_topic_live_during_fanout; @@ -2284,8 +2351,7 @@ TEST(takeover_will_claimed_before_fanout) ASSERT_NOT_NULL(g_will_owner); ASSERT_EQ(1, (int)g_will_owner->has_will); ASSERT_NOT_NULL(g_will_owner->will_topic); - /* Nothing has been published yet, so the hook below cannot be consumed by - * an earlier delivery. */ + /* Nothing published yet, so the hook cannot be consumed early. */ ASSERT_EQ(0, count_packets_of_type(g_clients[0].out_buf, g_clients[0].out_len, MQTT_PACKET_TYPE_PUBLISH)); @@ -2305,8 +2371,7 @@ TEST(takeover_will_claimed_before_fanout) /* Claimed before the first write: a close callback delivered during the * fan-out returns without republishing or freeing anything. */ ASSERT_EQ(1, g_will_claimed_during_fanout); - /* And the buffers the fan-out is lending out are still owned, so the - * encoder is reading live memory. */ + /* And the buffers it lends the encoder are still owned. */ ASSERT_EQ(1, g_will_topic_live_during_fanout); /* Exactly one copy reaches each subscriber [MQTT-3.1.2-8]. */ @@ -10785,6 +10850,7 @@ int main(int argc, char** argv) RUN_TEST(online_qos1_at_cap_keeps_subscriber); RUN_TEST(outbound_packet_ids_are_scoped_per_session); RUN_TEST(drain_stops_when_session_handoff_takes_queue); + RUN_TEST(handoff_marks_sent_qos1_for_redelivery); #ifdef WOLFMQTT_NONBLOCK RUN_TEST(outbound_queue_short_write_resumes_on_next_step); #ifdef WOLFMQTT_BROKER_RETAINED diff --git a/tests/test_mqtt_client.c b/tests/test_mqtt_client.c index 3c70ec1eb..7121c8148 100644 --- a/tests/test_mqtt_client.c +++ b/tests/test_mqtt_client.c @@ -48,21 +48,18 @@ static MqttNet test_net; static byte test_tx_buf[TEST_TX_BUF_SIZE]; static byte test_rx_buf[TEST_RX_BUF_SIZE]; -/* The suite routes the library's allocator through here (see - * tests/test_client_alloc.h) so a test can watch one specific allocation and - * record the state the library was in when it released it. Everything else is - * a straight pass-through. */ +/* The library allocator is routed through here (tests/test_broker_alloc.h) so + * a test can watch one allocation and record the state it was released in. */ static void* g_free_watch_ptr; static int g_free_watch_count; static int g_free_watch_lock_count; -void* wolfmqtt_test_client_malloc(size_t size) +void* wolfmqtt_test_broker_malloc(size_t size) { return malloc(size); } -/* How many times the client lock is held right now, or -1 where the build - * keeps no count to read. */ +/* Lock depth, or -1 where the build keeps no count to read. */ static int test_client_lock_depth(void) { #if defined(WOLFMQTT_MULTITHREAD) && defined(WOLFMQTT_POSIX_SEMAPHORES) && \ @@ -73,7 +70,7 @@ static int test_client_lock_depth(void) #endif } -void wolfmqtt_test_client_free(void* ptr) +void wolfmqtt_test_broker_free(void* ptr) { if (ptr != NULL && ptr == g_free_watch_ptr) { g_free_watch_count++; @@ -1827,15 +1824,10 @@ TEST(cancel_message_retain_is_idempotent) #if defined(WOLFMQTT_MULTITHREAD) && defined(WOLFMQTT_NONBLOCK) && \ (WOLFMQTT_MAX_QOS >= 1) -/* A write-only publish leaves its response for another thread to process. That - * reader claims the pending entry with packetProcessing, then drops the client - * lock and decodes the peer's answer into the packet_obj this message owns. - * - * Reporting success to a cancel during that window would be wrong: the caller - * reads success as permission to release or reuse the request, and the decode - * is still writing through it. The claim is honoured instead - the entry stays - * listed, the message is left exactly as it was found, and the caller retries - * until the reader marks the response done. */ +/* A write-only publish leaves its response to another thread, which claims the + * pending entry, drops the client lock and decodes into the packet_obj this + * message owns. A successful cancel is the caller's cue to release or reuse the + * object, so it is withheld until that claim is gone. */ TEST(cancel_refuses_message_while_response_decodes) { int rc; @@ -1865,14 +1857,12 @@ TEST(cancel_refuses_message_while_response_decodes) for (i = 0; i < 20 && rc == MQTT_CODE_CONTINUE; i++) { rc = MqttClient_Publish_WriteOnly(&test_client, &publish, NULL); } - /* The PUBLISH is out and its acknowledgement is registered for whichever - * thread reads it. */ + /* Out, with its acknowledgement registered for whichever thread reads. */ ASSERT_TRUE(test_client.firstPendResp == &publish.pendResp); ASSERT_EQ(0, (int)publish.pendResp.packetProcessing); write_state = (int)publish.stat.write; - /* A reader claims the entry, exactly as MqttClient_WaitType does before it - * drops the client lock to decode. */ + /* Claim the entry as MqttClient_WaitType does before it decodes. */ publish.pendResp.packetProcessing = 1; rc = MqttClient_CancelMessage(&test_client, (MqttObject*)&publish); @@ -1889,6 +1879,7 @@ TEST(cancel_refuses_message_while_response_decodes) ASSERT_NULL(test_client.firstPendResp); ASSERT_EQ(MQTT_MSG_BEGIN, (int)publish.stat.write); } + #endif /* WOLFMQTT_MULTITHREAD && WOLFMQTT_NONBLOCK && MAX_QOS >= 1 */ #endif /* WOLFMQTT_MULTITHREAD || WOLFMQTT_NONBLOCK */ @@ -3413,19 +3404,98 @@ TEST(reconnect_without_session_present_replays_nothing) ASSERT_EQ(1, g_frames_written); /* the CONNECT only */ } +/* A send is allowed once CONNECT reaches the transport, so a publish can be + * retained while this handshake still waits for CONNACK. It belongs to the + * Session being established, so a Session Present = 0 answer must keep it for + * a later resume [MQTT-4.4.0-1]. The write callback seeds what such a publish + * leaves behind, at the only point it is reachable: the Packet Identifier + * table is empty and CONNECT is on the wire. */ +static int g_seed_concurrent_publish; +static byte g_seed_payload[] = "seeded"; + +static int mock_net_write_seed_publish(void *context, const byte* buf, + int buf_len, int timeout_ms) +{ + int rc = mock_net_write_accept(context, buf, buf_len, timeout_ms); + + if (g_seed_concurrent_publish && buf_len > 0 && + (buf[0] >> 4) == MQTT_PACKET_TYPE_CONNECT) { +#ifndef WOLFMQTT_STATIC_MEMORY + char* topic; + byte* payload; +#endif + + g_seed_concurrent_publish = 0; + test_client.replay[1].packet_id = 0x5678; + test_client.replay[1].qos = MQTT_QOS_1; + test_client.replay[1].haveCopy = 1; + test_client.replay[1].payload_len = (word32)(sizeof(g_seed_payload) - 1); +#ifdef WOLFMQTT_STATIC_MEMORY + XMEMCPY(test_client.replay[1].topic, "new/1", 6); + XMEMCPY(test_client.replay[1].payload, g_seed_payload, + sizeof(g_seed_payload) - 1); +#else + topic = (char*)WOLFMQTT_MALLOC(6); + payload = (byte*)WOLFMQTT_MALLOC(sizeof(g_seed_payload) - 1); + if (topic == NULL || payload == NULL) { + WOLFMQTT_FREE(topic); + WOLFMQTT_FREE(payload); + test_client.replay[1].packet_id = 0; + return rc; + } + XMEMCPY(topic, "new/1", 6); + XMEMCPY(payload, g_seed_payload, sizeof(g_seed_payload) - 1); + test_client.replay[1].topic = topic; + test_client.replay[1].payload = payload; +#endif + /* The reservation that marks the entry as this connection's. */ + test_client.send_inflight[0].packet_id = 0x5678; + test_client.send_inflight[0].ack_type = MQTT_PACKET_TYPE_PUBLISH_ACK; + } + return rc; +} + +TEST(fresh_session_reset_keeps_publish_from_this_connection) +{ + int rc; + MqttConnect connect; + MqttPublish publish; + static byte payload[] = "hello"; + + rc = test_init_client(); + ASSERT_EQ(MQTT_CODE_SUCCESS, rc); +#ifdef WOLFMQTT_V5 + test_client.protocol_level = MQTT_CONNECT_PROTOCOL_LEVEL_4; +#endif + ASSERT_EQ(MQTT_CODE_SUCCESS, run_initial_connect(&connect)); + + /* An unacknowledged publish from the Session that is about to end. */ + init_qos_publish(&publish, MQTT_QOS_1, 0x1234, "sensor/temp", + payload, (word32)(sizeof(payload) - 1)); + rc = run_publish_unacked(&publish); + ASSERT_EQ(MQTT_CODE_ERROR_NETWORK, rc); + ASSERT_EQ(MQTT_CODE_SUCCESS, MqttClient_NetDisconnect(&test_client)); + ASSERT_EQ(0x1234, (int)test_client.replay[0].packet_id); + + g_seed_concurrent_publish = 1; + rc = run_reconnect_with(&connect, 0, mock_net_write_seed_publish); + ASSERT_EQ(MQTT_CODE_SUCCESS, rc); + /* The seed really was planted, so the assertions below mean something. */ + ASSERT_EQ(0, g_seed_concurrent_publish); + + /* The ended Session's message is gone. */ + ASSERT_EQ(0, (int)test_client.replay[0].packet_id); + /* The one published on this connection is kept, with its copy intact. */ + ASSERT_EQ(0x5678, (int)test_client.replay[1].packet_id); + ASSERT_STR_EQ("new/1", test_client.replay[1].topic); +} + #if defined(WOLFMQTT_MULTITHREAD) && defined(WOLFMQTT_POSIX_SEMAPHORES) && \ - !defined(WOLFMQTT_NO_COND_SIGNAL) -/* Discarding the retained copies releases the topic and payload of every - * slot. A send thread reaches those same slots through the replay store's - * lock-taking wrapper while the connection is being re-established - the - * client admits send calls once CONNECT has reached the transport - and both - * sides free what they find there. So the discard has to hold the client lock - * as well, or the two can free the same buffer or leave one pointing at - * memory the other released. - * - * The POSIX lock keeps its own depth, and the allocator hook samples it at the - * moment the slot's topic is released, which is the access that has to be - * covered. */ + !defined(WOLFMQTT_NO_COND_SIGNAL) && !defined(WOLFMQTT_STATIC_MEMORY) +/* A send thread reaches the same replay slots under the client lock, and both + * sides release what they find there, so the discard has to hold it too. The + * POSIX lock keeps its own depth; the allocator hook samples it as the slot's + * topic is released, which is the access that has to be covered. */ TEST(fresh_session_reset_frees_replay_under_client_lock) { int rc; @@ -3446,8 +3516,7 @@ TEST(fresh_session_reset_frees_replay_under_client_lock) ASSERT_EQ(MQTT_CODE_ERROR_NETWORK, rc); ASSERT_EQ(MQTT_CODE_SUCCESS, MqttClient_NetDisconnect(&test_client)); - /* The unacknowledged PUBLISH is retained, so the discard has a real - * allocation to release. */ + /* Retained, so the discard has a real allocation to release. */ ASSERT_EQ(0x1234, (int)test_client.replay[0].packet_id); ASSERT_NOT_NULL(test_client.replay[0].topic); test_client_watch_free(test_client.replay[0].topic); @@ -7800,8 +7869,9 @@ void run_mqtt_client_tests(void) #if !defined(WOLFMQTT_NO_SESSION_REPLAY) && WOLFMQTT_MAX_QOS >= 1 RUN_TEST(reconnect_replays_unacked_qos1_publish); RUN_TEST(reconnect_without_session_present_replays_nothing); + RUN_TEST(fresh_session_reset_keeps_publish_from_this_connection); #if defined(WOLFMQTT_MULTITHREAD) && defined(WOLFMQTT_POSIX_SEMAPHORES) && \ - !defined(WOLFMQTT_NO_COND_SIGNAL) + !defined(WOLFMQTT_NO_COND_SIGNAL) && !defined(WOLFMQTT_STATIC_MEMORY) RUN_TEST(fresh_session_reset_frees_replay_under_client_lock); #endif RUN_TEST(reconnect_does_not_replay_acked_publish); diff --git a/wolfmqtt/mqtt_client.h b/wolfmqtt/mqtt_client.h index 1ee28a59a..465380235 100644 --- a/wolfmqtt/mqtt_client.h +++ b/wolfmqtt/mqtt_client.h @@ -888,10 +888,9 @@ WOLFMQTT_API int MqttClient_WaitMessage_ex( /*! \brief In a multi-threaded and non-blocking mode this allows you to cancel an MQTT object that was previously submitted. * \note This is a blocking function that will wait for MqttNet.read - * \note Only MQTT_CODE_SUCCESS releases the object back to the caller. - MQTT_CODE_CONTINUE means another thread is still processing - this message's response: the object must not be freed or - reused, and the call should be retried. + * \note MQTT_CODE_CONTINUE means another thread is still processing this + message's response: retry, and do not free or reuse the object + until MQTT_CODE_SUCCESS. * \param client Pointer to MqttClient structure * \param msg Pointer to MqttObject structure * \return MQTT_CODE_SUCCESS, MQTT_CODE_CONTINUE or MQTT_CODE_ERROR_* From e5d1ac517673f6fa6878c27bd431a2fe1ffec222 Mon Sep 17 00:00:00 2001 From: Kareem Date: Tue, 6 Oct 2026 10:10:35 -0700 Subject: [PATCH 6/8] Code review feedback --- src/mqtt_broker.c | 5 + src/mqtt_client.c | 95 +++++++++---- tests/test_broker_connect.c | 145 +++++++++++++++++++ tests/test_mqtt_client.c | 271 +++++++++++++++++++++++++++++++----- wolfmqtt/mqtt_client.h | 6 + 5 files changed, 460 insertions(+), 62 deletions(-) diff --git a/src/mqtt_broker.c b/src/mqtt_broker.c index 65a5cabd3..bc3b830bc 100644 --- a/src/mqtt_broker.c +++ b/src/mqtt_broker.c @@ -3525,6 +3525,11 @@ static void BrokerOrphan_ReconcileSent(MqttBroker* broker, } if (cur->qos != MQTT_QOS_0) { cur->retransmit_dup = 1; + #ifdef WOLFMQTT_BROKER_PERSIST + /* BrokerOrphan_Take shadow-wrote this entry before the write + * returned, so the stored copy still says it was never sent. */ + (void)BrokerPersist_PutOutPub(broker, client_id, cur); + #endif return; } if (prev == NULL) { diff --git a/src/mqtt_client.c b/src/mqtt_client.c index 2c732f4c7..07de35393 100644 --- a/src/mqtt_client.c +++ b/src/mqtt_client.c @@ -88,6 +88,8 @@ static int MqttClient_AuthEx(MqttClient *client, MqttAuth* auth, #if !defined(WOLFMQTT_MULTITHREAD) && !defined(WOLFMQTT_NONBLOCK) static int MqttClient_CancelMessage(MqttClient *client, MqttObject* msg); #endif +static int MqttClient_CancelMessageEx(MqttClient *client, MqttObject* msg, + int force); #ifndef WOLFMQTT_NO_SESSION_REPLAY static int MqttClient_SendIds_Find(const MqttClient* client, word16 packet_id); #endif @@ -736,6 +738,7 @@ static void MqttClient_Replay_Add(MqttClient* client, MqttPublish* publish) #endif slot->payload_len = publish->total_len; slot->haveCopy = 1; + slot->onThisConn = 1; } /* QoS 2 advanced past PUBREC, so [MQTT-4.4.0-1] replays the PUBREL rather @@ -757,6 +760,7 @@ static void MqttClient_Replay_PubRelSent(MqttClient* client, word16 packet_id) MqttClient_Replay_FreeSlot(slot); slot->packet_id = packet_id; slot->qos = MQTT_QOS_2; + slot->onThisConn = 1; } slot->pubrelSent = 1; } @@ -800,13 +804,28 @@ static void MqttClient_Replay_AddSafe(MqttClient* client, MqttPublish* publish) #endif } -/* Discard the retained copies of a Session the server did not resume. A send - * is allowed once CONNECT reaches the transport, so an entry may instead belong - * to the Session being established: MqttClient_Connect empties the Packet - * Identifier table before sending CONNECT, so a still-reserved identifier marks - * one of those, and it is kept for a later resume [MQTT-4.4.0-1]. Reports a - * lock failure, unlike the wrappers above: a skipped discard would replay the - * previous Session's messages into this one. */ +/* Everything retained so far belongs to the Session that is ending; a send + * admitted once CONNECT is on the wire tags itself. */ +static void MqttClient_Replay_NewConnSafe(MqttClient* client) +{ + int i; +#ifdef WOLFMQTT_MULTITHREAD + if (wm_SemLock(&client->lockClient) != MQTT_CODE_SUCCESS) { + return; + } +#endif + for (i = 0; i < MQTT_MAX_REPLAY_MSGS; i++) { + client->replay[i].onThisConn = 0; + } +#ifdef WOLFMQTT_MULTITHREAD + wm_SemUnlock(&client->lockClient); +#endif +} + +/* Discard the retained copies of a Session the server did not resume, keeping + * anything published on this connection: a later resume still owes that + * [MQTT-4.4.0-1]. Reports a lock failure, unlike the wrappers above: a skipped + * discard would replay the old Session's messages into this one. */ static int MqttClient_Replay_ResetSafe(MqttClient* client) { int i; @@ -817,10 +836,8 @@ static int MqttClient_Replay_ResetSafe(MqttClient* client) } #endif for (i = 0; i < MQTT_MAX_REPLAY_MSGS; i++) { - if (client->replay[i].packet_id != 0 && - MqttClient_SendIds_Find(client, - client->replay[i].packet_id) >= 0) { - continue; /* published on this connection */ + if (client->replay[i].onThisConn) { + continue; } MqttClient_Replay_FreeSlot(&client->replay[i]); } @@ -2971,7 +2988,8 @@ static int MqttClient_Replay_EncodeNext(MqttClient* client, word16* dropped_id) * still acknowledge the message it stands for, and that acknowledgement * is what releases both the slot and its Packet Identifier * [MQTT-2.3.1-3]. */ - if (slot->packet_id == 0 || (!slot->pubrelSent && !slot->haveCopy)) { + if (slot->packet_id == 0 || slot->onThisConn || + (!slot->pubrelSent && !slot->haveCopy)) { rc = 0; keep = 1; } @@ -3262,6 +3280,9 @@ int MqttClient_Connect(MqttClient *client, MqttConnect *mc_connect) * and outside the WOLFMQTT_V5 block below, since the table exists in * every build. */ MqttClient_SendIdsReset(client); +#ifndef WOLFMQTT_NO_SESSION_REPLAY + MqttClient_Replay_NewConnSafe(client); +#endif #ifdef WOLFMQTT_V5 #ifdef WOLFMQTT_MULTITHREAD @@ -3454,7 +3475,8 @@ int MqttClient_Connect(MqttClient *client, MqttConnect *mc_connect) /* The handshake state set above stands or not on whether bytes * reached the peer, so a retry is refused only when some of this * CONNECT is already out there. */ - MqttClient_CancelMessage(client, (MqttObject*)mc_connect); + (void)MqttClient_CancelMessageEx(client, + (MqttObject*)mc_connect, 1); return rc; } @@ -3673,7 +3695,9 @@ int MqttClient_Connect(MqttClient *client, MqttConnect *mc_connect) } #endif for (i = 0; i < MQTT_MAX_REPLAY_MSGS; i++) { - if (client->replay[i].packet_id == 0) { + if (client->replay[i].packet_id == 0 || + client->replay[i].onThisConn) { + /* Not the resumed Session's: already in flight here. */ continue; } if (!client->replay[i].pubrelSent && @@ -3685,6 +3709,14 @@ int MqttClient_Connect(MqttClient *client, MqttConnect *mc_connect) MqttClient_Replay_FreeSlot(&client->replay[i]); continue; } + if (MqttClient_SendIds_Find(client, + client->replay[i].packet_id) >= 0) { + /* A live request holds the identifier and + * [MQTT-2.3.1-2] allows one use at a time, so drop this + * rather than take its reservation over. */ + MqttClient_Replay_FreeSlot(&client->replay[i]); + continue; + } (void)MqttClient_SendIdReserve_Locked(client, client->replay[i].packet_id, &client->replay[i], 1, client->replay[i].pubrelSent ? @@ -4244,7 +4276,8 @@ static int MqttPublishMsg(MqttClient *client, MqttPublish *publish, * message's ownership of it either way. */ MqttClient_RestoreRecvQuota(client, publish); #endif - MqttClient_CancelMessage(client, (MqttObject*)publish); + (void)MqttClient_CancelMessageEx(client, + (MqttObject*)publish, 1); if (wrote > 0) { /* Part of the PUBLISH reached the server, so it may have * seen the Packet Identifier. Reclaim the reservation the @@ -4291,7 +4324,8 @@ static int MqttPublishMsg(MqttClient *client, MqttPublish *publish, * above. */ MqttClient_RestoreRecvQuota(client, publish); #endif - MqttClient_CancelMessage(client, (MqttObject*)publish); + (void)MqttClient_CancelMessageEx(client, + (MqttObject*)publish, 1); /* Reaching the payload means the fixed header - and with it * the Packet Identifier - is already on the wire, so keep the * reservation the cancel dropped [MQTT-2.3.1-3]. */ @@ -4540,7 +4574,7 @@ int MqttClient_Subscribe(MqttClient *client, MqttSubscribe *subscribe) wrote = (rc == xfer) ? xfer : client->write.pos; MqttWriteStop(client, &subscribe->stat); if (rc != xfer) { - MqttClient_CancelMessage(client, (MqttObject*)subscribe); + (void)MqttClient_CancelMessageEx(client, (MqttObject*)subscribe, 1); if (wrote > 0) { /* Part of the SUBSCRIBE reached the server, so it may have * seen the Packet Identifier. Reclaim the reservation the @@ -4695,7 +4729,8 @@ int MqttClient_Unsubscribe(MqttClient *client, MqttUnsubscribe *unsubscribe) wrote = (rc == xfer) ? xfer : client->write.pos; MqttWriteStop(client, &unsubscribe->stat); if (rc != xfer) { - MqttClient_CancelMessage(client, (MqttObject*)unsubscribe); + (void)MqttClient_CancelMessageEx(client, + (MqttObject*)unsubscribe, 1); if (wrote > 0) { /* Part of the UNSUBSCRIBE reached the server, so it may have * seen the Packet Identifier. Reclaim the reservation the @@ -4859,7 +4894,7 @@ int MqttClient_Ping_ex(MqttClient *client, MqttPing* ping) #endif MqttWriteStop(client, &ping->stat); if (rc != xfer) { - MqttClient_CancelMessage(client, (MqttObject*)ping); + (void)MqttClient_CancelMessageEx(client, (MqttObject*)ping, 1); return rc; } @@ -5383,16 +5418,19 @@ int MqttClient_WaitMessage(MqttClient *client, int timeout_ms) return MqttClient_WaitMessage_ex(client, &client->msg, timeout_ms); } -#if !defined(WOLFMQTT_MULTITHREAD) && !defined(WOLFMQTT_NONBLOCK) -static -#endif -int MqttClient_CancelMessage(MqttClient *client, MqttObject* msg) +/* force: detach even while a reader holds its claim. Failure cleanup is + * handing the object back with an error, so it can leave nothing linked; only + * the public entry point can ask the caller to wait. */ +static int MqttClient_CancelMessageEx(MqttClient *client, MqttObject* msg, + int force) { int rc = MQTT_CODE_SUCCESS; MqttMsgStat* mms_stat; int onWire; #ifdef WOLFMQTT_MULTITHREAD MqttPendResp* tmpResp; +#else + (void)force; /* nothing can hold a claim without a reader thread */ #endif if (client == NULL || msg == NULL) { @@ -5441,7 +5479,8 @@ int MqttClient_CancelMessage(MqttClient *client, MqttObject* msg) tmpResp->packet_type, tmpResp->packet_id, tmpResp->packetProcessing, tmpResp->packetDone); #endif - if (tmpResp->packetProcessing && !tmpResp->packetDone) { + if (!force && tmpResp->packetProcessing && + !tmpResp->packetDone) { wm_SemUnlock(&client->lockClient); return MQTT_CODE_CONTINUE; } @@ -5520,6 +5559,14 @@ int MqttClient_CancelMessage(MqttClient *client, MqttObject* msg) return rc; } +#if !defined(WOLFMQTT_MULTITHREAD) && !defined(WOLFMQTT_NONBLOCK) +static +#endif +int MqttClient_CancelMessage(MqttClient *client, MqttObject* msg) +{ + return MqttClient_CancelMessageEx(client, msg, 0); +} + #ifdef WOLFMQTT_NONBLOCK static inline int IsMessageActive(MqttObject *msg) { diff --git a/tests/test_broker_connect.c b/tests/test_broker_connect.c index 7008d97b8..9b69cc6d3 100644 --- a/tests/test_broker_connect.c +++ b/tests/test_broker_connect.c @@ -2047,6 +2047,24 @@ static void handoff_out_queue_mid_write(void) o->next = broker->orphan_sessions; broker->orphan_sessions = o; broker->orphan_session_count++; + +#ifdef WOLFMQTT_BROKER_PERSIST + /* BrokerOrphan_Take shadow-writes here, before the write it interrupted + * has returned, so the stored copy records an unsent entry. */ + (void)BrokerPersist_PutOrphanSession(broker, o->client_id, + o->protocol_level, o->session_expiry_sec, o->orphan_since); + { + BrokerOutPub* e; + word64 seq = 1; + + for (e = o->out_q_head; e != NULL; e = e->next) { + e->enqueue_seq = seq++; + if (e->qos > MQTT_QOS_0) { + (void)BrokerPersist_PutOutPub(broker, o->client_id, e); + } + } + } +#endif } TEST(drain_stops_when_session_handoff_takes_queue) @@ -9213,6 +9231,31 @@ static int persist_order_put(void* ctx, byte ns, const byte* key, return MQTT_CODE_SUCCESS; } +/* Keyed variant: a record written again under the same key replaces it, as a + * real backend does. persist_order_put appends, which would turn a re-write + * into a second queue entry. */ +static int persist_keyed_put(void* ctx, byte ns, const byte* key, + word16 key_len, const byte* blob, word32 blob_len) +{ + PersistOrderStore* store = (PersistOrderStore*)ctx; + int i; + + if (ns == BROKER_PERSIST_NS_OUTQ) { + for (i = 0; i < store->outq_count; i++) { + if (store->outq[i].key_len == key_len && + XMEMCMP(store->outq[i].key, key, key_len) == 0) { + if (blob_len > sizeof(store->outq[i].blob)) { + return MQTT_CODE_ERROR_OUT_OF_BUFFER; + } + XMEMCPY(store->outq[i].blob, blob, blob_len); + store->outq[i].blob_len = blob_len; + return MQTT_CODE_SUCCESS; + } + } + } + return persist_order_put(ctx, ns, key, key_len, blob, blob_len); +} + static int persist_order_get(void* ctx, byte ns, const byte* key, word16 key_len, byte* out, word32* inout_len) { @@ -9393,6 +9436,107 @@ TEST(persist_partial_publish_restart_keeps_dup) MqttBroker_Stop(&restored); MqttBroker_Free(&restored); } + +/* The same requirement when the Session hand-off takes the queue mid-write. + * BrokerOrphan_Take shadow-writes each entry before that write returns, so the + * stored copy still says it was never sent; the completed delivery has to be + * written through, or a restart before reconnect re-sends it without DUP + * [MQTT-3.3.1-1]. */ +TEST(persist_handoff_sent_qos1_restart_keeps_dup) +{ + MqttBroker source; + MqttBroker restored; + MqttBrokerNet net; + MqttBrokerPersistHooks hooks; + PersistOrderStore store; + BrokerClient* sub_bc; + PublishInfo info; + word16 packet_id; + int i; + static const byte connect_pub[] = { + 0x10, 0x0D, 0x00, 0x04, 'M', 'Q', 'T', 'T', 0x04, 0x02, 0x00, 0x3C, + 0x00, 0x01, 'P' + }; + /* CleanSession=0, so the Session survives and has a carrier. */ + static const byte connect_sub[] = { + 0x10, 0x0D, 0x00, 0x04, 'M', 'Q', 'T', 'T', 0x04, 0x00, 0x00, 0x3C, + 0x00, 0x01, 'S' + }; + static const byte subscribe_x[] = { + 0x82, 0x06, 0x00, 0x01, 0x00, 0x01, 'x', 0x01 + }; + static const byte publish_x[] = { + 0x32, 0x08, 0x00, 0x01, 'x', 0x00, 0x07, 'A', 'B', 'C' + }; + + install_mock_net(&net); + XMEMSET(&source, 0, sizeof(source)); + XMEMSET(&restored, 0, sizeof(restored)); + XMEMSET(&hooks, 0, sizeof(hooks)); + XMEMSET(&store, 0, sizeof(store)); + hooks.kv_put = persist_keyed_put; + hooks.kv_get = persist_order_get; + hooks.kv_iter = persist_order_iter; + hooks.ctx = &store; +#ifdef WOLFMQTT_BROKER_PERSIST_ENCRYPT + hooks.derive_key = persist_order_derive_key; +#endif + + ASSERT_EQ(MQTT_CODE_SUCCESS, MqttBroker_Init(&source, &net)); + ASSERT_EQ(MQTT_CODE_SUCCESS, MqttBroker_SetPersistHooks(&source, &hooks)); + ASSERT_EQ(MQTT_CODE_SUCCESS, BrokerPersist_Restore(&source)); + ASSERT_EQ(MQTT_CODE_SUCCESS, MqttBroker_Start(&source)); + reset_mock_clients(2); + mock_client_input_append(0, connect_pub, sizeof(connect_pub)); + mock_client_input_append(1, connect_sub, sizeof(connect_sub)); + mock_client_input_append(1, subscribe_x, sizeof(subscribe_x)); + for (i = 0; i < 16; i++) { + (void)MqttBroker_Step(&source); + } + sub_bc = find_broker_client(&source, "S"); + ASSERT_NOT_NULL(sub_bc); + + /* The hand-off runs from inside the write, which then completes. */ + broker_test_watch_free(NULL); + g_handoff_broker = &source; + g_handoff_client = sub_bc; + g_write_publish_hook = handoff_out_queue_mid_write; + mock_client_input_append(0, publish_x, sizeof(publish_x)); + (void)MqttBroker_Step(&source); + ASSERT_NULL(g_write_publish_hook); + g_handoff_broker = NULL; + g_handoff_client = NULL; + + ASSERT_NOT_NULL(source.orphan_sessions); + ASSERT_NOT_NULL(source.orphan_sessions->out_q_head); + ASSERT_EQ(1, (int)source.orphan_sessions->out_q_head->retransmit_dup); + packet_id = source.orphan_sessions->out_q_head->packet_id; + ASSERT_EQ(1, store.outq_count); + + MqttBroker_Stop(&source); + MqttBroker_Free(&source); + + /* Restart before the subscriber reconnects. */ + ASSERT_EQ(MQTT_CODE_SUCCESS, MqttBroker_Init(&restored, &net)); + ASSERT_EQ(MQTT_CODE_SUCCESS, + MqttBroker_SetPersistHooks(&restored, &hooks)); + ASSERT_EQ(MQTT_CODE_SUCCESS, BrokerPersist_Restore(&restored)); + ASSERT_EQ(MQTT_CODE_SUCCESS, MqttBroker_Start(&restored)); + reset_mock_clients(1); + mock_client_input_append(0, connect_sub, sizeof(connect_sub)); + for (i = 0; i < 16; i++) { + (void)MqttBroker_Step(&restored); + } + + info = first_publish_info(g_clients[0].out_buf, g_clients[0].out_len); + ASSERT_TRUE(info.found); + /* 0x3A = PUBLISH | DUP | QoS 1. */ + ASSERT_EQ(0x3A, info.first_byte); + ASSERT_EQ(packet_id, info.packet_id); + + MqttBroker_Stop(&restored); + MqttBroker_Free(&restored); +} #endif /* Non-persisted QoS 0 nodes still occupy FIFO positions. Assigning sequence @@ -11081,6 +11225,7 @@ int main(int argc, char** argv) RUN_TEST(persist_restore_packet_id_wrap_preserves_fifo); #ifdef WOLFMQTT_NONBLOCK RUN_TEST(persist_partial_publish_restart_keeps_dup); + RUN_TEST(persist_handoff_sent_qos1_restart_keeps_dup); #endif RUN_TEST(persist_mixed_qos_queue_preserves_fifo); RUN_TEST(orphan_reclaim_keeps_persisted_outq); diff --git a/tests/test_mqtt_client.c b/tests/test_mqtt_client.c index 7121c8148..55439e6db 100644 --- a/tests/test_mqtt_client.c +++ b/tests/test_mqtt_client.c @@ -1881,6 +1881,60 @@ TEST(cancel_refuses_message_while_response_decodes) } #endif /* WOLFMQTT_MULTITHREAD && WOLFMQTT_NONBLOCK && MAX_QOS >= 1 */ +#if defined(WOLFMQTT_MULTITHREAD) && (WOLFMQTT_MAX_QOS >= 1) +/* The library's own failure cleanup is a different case from a caller asking + * to reclaim an object: it is handing the object back with an error, so the + * pending response must not stay linked whatever a reader is doing. A caller + * that frees on that error would otherwise leave a dangling list node. + * + * The write callback claims the entry the way a reader does, from inside the + * write that is about to fail. */ +static MqttPublish* g_claim_on_write; +static int mock_net_write_claim_then_fail(void *context, const byte* buf, + int buf_len, int timeout_ms) +{ + (void)context; (void)timeout_ms; + if (g_claim_on_write != NULL && buf_len > 0 && + (buf[0] >> 4) == MQTT_PACKET_TYPE_PUBLISH) { + g_claim_on_write->pendResp.packetProcessing = 1; + g_claim_on_write = NULL; + } + return MQTT_CODE_ERROR_NETWORK; +} + +TEST(write_failure_cleanup_unlinks_claimed_pending_response) +{ + int rc; + static MqttPublish publish; + static byte payload[] = "hello"; + + rc = test_init_client(); + ASSERT_EQ(MQTT_CODE_SUCCESS, rc); + test_client_connect_sent(); + + XMEMSET(&publish, 0, sizeof(publish)); + publish.qos = MQTT_QOS_1; + publish.packet_id = 41; + publish.topic_name = "test/topic"; + publish.buffer = payload; + publish.total_len = (word32)(sizeof(payload) - 1); + publish.buffer_len = publish.total_len; + + g_claim_on_write = &publish; + test_net.write = mock_net_write_claim_then_fail; + test_net.read = mock_net_read; + + rc = MqttClient_Publish(&test_client, &publish); + ASSERT_EQ(MQTT_CODE_ERROR_NETWORK, rc); + /* The claim really was taken, so the assertion below means something. */ + ASSERT_NULL(g_claim_on_write); + ASSERT_EQ(1, (int)publish.pendResp.packetProcessing); + + /* Nothing left linked for the caller to dangle. */ + ASSERT_NULL(test_client.firstPendResp); +} +#endif /* WOLFMQTT_MULTITHREAD && MAX_QOS >= 1 */ + #endif /* WOLFMQTT_MULTITHREAD || WOLFMQTT_NONBLOCK */ /* A QoS>0 v5 publish that fails on the wire (unsent) must give its reserved @@ -3406,51 +3460,74 @@ TEST(reconnect_without_session_present_replays_nothing) /* A send is allowed once CONNECT reaches the transport, so a publish can be * retained while this handshake still waits for CONNACK. It belongs to the - * Session being established, so a Session Present = 0 answer must keep it for - * a later resume [MQTT-4.4.0-1]. The write callback seeds what such a publish - * leaves behind, at the only point it is reachable: the Packet Identifier - * table is empty and CONNECT is on the wire. */ -static int g_seed_concurrent_publish; + * Session being established, which the entry's onThisConn tag records. The + * write callback seeds what such a send leaves behind, at the only point it is + * reachable: the tags have just been cleared and CONNECT is on the wire. */ +static int g_seed_slot = -1; /* slot to seed, -1 seeds nothing */ +static word16 g_seed_packet_id; +static byte g_seed_on_this_conn; +static byte g_seed_reserve_ack_type; /* 0 reserves no identifier */ static byte g_seed_payload[] = "seeded"; +static void seed_replay_slot(int idx, word16 packet_id, byte on_this_conn) +{ + MqttReplayMsg* slot = &test_client.replay[idx]; + word32 len = (word32)(sizeof(g_seed_payload) - 1); + + slot->packet_id = packet_id; + slot->qos = MQTT_QOS_1; + slot->haveCopy = 1; + slot->onThisConn = on_this_conn; + slot->payload_len = len; +#ifdef WOLFMQTT_STATIC_MEMORY + XMEMCPY(slot->topic, "new/1", 6); + XMEMCPY(slot->payload, g_seed_payload, len); +#else + slot->topic = (char*)WOLFMQTT_MALLOC(6); + slot->payload = (byte*)WOLFMQTT_MALLOC(len); + if (slot->topic == NULL || slot->payload == NULL) { + WOLFMQTT_FREE(slot->topic); + WOLFMQTT_FREE(slot->payload); + slot->topic = NULL; + slot->payload = NULL; + slot->packet_id = 0; + return; + } + XMEMCPY(slot->topic, "new/1", 6); + XMEMCPY(slot->payload, g_seed_payload, len); +#endif +} + +/* A SUBSCRIBE admitted once CONNECT is on the wire, reduced to the state it + * leaves behind: its Packet Identifier reservation. */ +static word16 g_reserve_sub_id = 0x1234; +static int mock_net_write_reserve_subscribe_id(void *context, const byte* buf, + int buf_len, int timeout_ms) +{ + int rc = mock_net_write_accept(context, buf, buf_len, timeout_ms); + + if (buf_len > 0 && (buf[0] >> 4) == MQTT_PACKET_TYPE_CONNECT) { + test_client.send_inflight[0].packet_id = g_reserve_sub_id; + test_client.send_inflight[0].ack_type = MQTT_PACKET_TYPE_SUBSCRIBE_ACK; + test_client.send_inflight[0].owner = &test_client.msg; + } + return rc; +} + static int mock_net_write_seed_publish(void *context, const byte* buf, int buf_len, int timeout_ms) { int rc = mock_net_write_accept(context, buf, buf_len, timeout_ms); - if (g_seed_concurrent_publish && buf_len > 0 && + if (g_seed_slot >= 0 && buf_len > 0 && (buf[0] >> 4) == MQTT_PACKET_TYPE_CONNECT) { -#ifndef WOLFMQTT_STATIC_MEMORY - char* topic; - byte* payload; -#endif - - g_seed_concurrent_publish = 0; - test_client.replay[1].packet_id = 0x5678; - test_client.replay[1].qos = MQTT_QOS_1; - test_client.replay[1].haveCopy = 1; - test_client.replay[1].payload_len = (word32)(sizeof(g_seed_payload) - 1); -#ifdef WOLFMQTT_STATIC_MEMORY - XMEMCPY(test_client.replay[1].topic, "new/1", 6); - XMEMCPY(test_client.replay[1].payload, g_seed_payload, - sizeof(g_seed_payload) - 1); -#else - topic = (char*)WOLFMQTT_MALLOC(6); - payload = (byte*)WOLFMQTT_MALLOC(sizeof(g_seed_payload) - 1); - if (topic == NULL || payload == NULL) { - WOLFMQTT_FREE(topic); - WOLFMQTT_FREE(payload); - test_client.replay[1].packet_id = 0; - return rc; + seed_replay_slot(g_seed_slot, g_seed_packet_id, g_seed_on_this_conn); + if (g_seed_reserve_ack_type != 0) { + test_client.send_inflight[0].packet_id = g_seed_packet_id; + test_client.send_inflight[0].ack_type = g_seed_reserve_ack_type; + test_client.send_inflight[0].owner = &test_client.msg; } - XMEMCPY(topic, "new/1", 6); - XMEMCPY(payload, g_seed_payload, sizeof(g_seed_payload) - 1); - test_client.replay[1].topic = topic; - test_client.replay[1].payload = payload; -#endif - /* The reservation that marks the entry as this connection's. */ - test_client.send_inflight[0].packet_id = 0x5678; - test_client.send_inflight[0].ack_type = MQTT_PACKET_TYPE_PUBLISH_ACK; + g_seed_slot = -1; } return rc; } @@ -3477,11 +3554,14 @@ TEST(fresh_session_reset_keeps_publish_from_this_connection) ASSERT_EQ(MQTT_CODE_SUCCESS, MqttClient_NetDisconnect(&test_client)); ASSERT_EQ(0x1234, (int)test_client.replay[0].packet_id); - g_seed_concurrent_publish = 1; + g_seed_slot = 1; + g_seed_packet_id = 0x5678; + g_seed_on_this_conn = 1; + g_seed_reserve_ack_type = MQTT_PACKET_TYPE_PUBLISH_ACK; rc = run_reconnect_with(&connect, 0, mock_net_write_seed_publish); ASSERT_EQ(MQTT_CODE_SUCCESS, rc); /* The seed really was planted, so the assertions below mean something. */ - ASSERT_EQ(0, g_seed_concurrent_publish); + ASSERT_EQ(-1, g_seed_slot); /* The ended Session's message is gone. */ ASSERT_EQ(0, (int)test_client.replay[0].packet_id); @@ -3490,6 +3570,115 @@ TEST(fresh_session_reset_keeps_publish_from_this_connection) ASSERT_STR_EQ("new/1", test_client.replay[1].topic); } +/* A reserved Packet Identifier does not prove where a replay entry came from: + * the table also holds SUBSCRIBE and UNSUBSCRIBE reservations. A request sent + * after CONNECT that happens to reuse an ended Session's identifier must not + * make a fresh CONNACK keep that Session's PUBLISH, which could otherwise be + * replayed into a later Session. */ +TEST(fresh_session_reset_drops_old_publish_behind_subscribe_id) +{ + int rc; + MqttConnect connect; + MqttPublish publish; + static byte payload[] = "hello"; + + rc = test_init_client(); + ASSERT_EQ(MQTT_CODE_SUCCESS, rc); +#ifdef WOLFMQTT_V5 + test_client.protocol_level = MQTT_CONNECT_PROTOCOL_LEVEL_4; +#endif + ASSERT_EQ(MQTT_CODE_SUCCESS, run_initial_connect(&connect)); + + init_qos_publish(&publish, MQTT_QOS_1, 0x1234, "sensor/temp", + payload, (word32)(sizeof(payload) - 1)); + rc = run_publish_unacked(&publish); + ASSERT_EQ(MQTT_CODE_ERROR_NETWORK, rc); + ASSERT_EQ(MQTT_CODE_SUCCESS, MqttClient_NetDisconnect(&test_client)); + ASSERT_EQ(0x1234, (int)test_client.replay[0].packet_id); + + /* A SUBSCRIBE on the new connection takes the same identifier. Nothing is + * seeded into the pool: slot 0 is already the ended Session's entry. */ + g_seed_slot = -1; + test_net.write = mock_net_write_accept; + rc = run_reconnect_with(&connect, 0, mock_net_write_reserve_subscribe_id); + ASSERT_EQ(MQTT_CODE_SUCCESS, rc); + ASSERT_EQ(0x1234, (int)test_client.send_inflight[0].packet_id); + + /* Discarded: it was never this connection's, whoever holds the id now. */ + ASSERT_EQ(0, (int)test_client.replay[0].packet_id); +} + +/* The resumed-Session walk must re-send only what the previous connection left + * unacknowledged. An entry published after CONNECT is already in flight here, + * so replaying it would deliver it twice. */ +TEST(resume_does_not_replay_publish_from_this_connection) +{ + int rc; + MqttConnect connect; + + rc = test_init_client(); + ASSERT_EQ(MQTT_CODE_SUCCESS, rc); +#ifdef WOLFMQTT_V5 + test_client.protocol_level = MQTT_CONNECT_PROTOCOL_LEVEL_4; +#endif + ASSERT_EQ(MQTT_CODE_SUCCESS, run_initial_connect(&connect)); + ASSERT_EQ(MQTT_CODE_SUCCESS, MqttClient_NetDisconnect(&test_client)); + + g_seed_slot = 0; + g_seed_packet_id = 0x5678; + g_seed_on_this_conn = 1; + g_seed_reserve_ack_type = MQTT_PACKET_TYPE_PUBLISH_ACK; + g_frames_written = 0; + rc = run_reconnect_with(&connect, 1, mock_net_write_seed_publish); + ASSERT_EQ(MQTT_CODE_SUCCESS, rc); + ASSERT_EQ(-1, g_seed_slot); + + /* The CONNECT only: nothing from this connection is re-sent. */ + ASSERT_EQ(1, g_frames_written); + /* And the entry stays retained, awaiting its own acknowledgement. */ + ASSERT_EQ(0x5678, (int)test_client.replay[0].packet_id); +} + +/* [MQTT-2.3.1-2] allows one use of an identifier at a time. If a request on + * the new connection already holds the one an ended Session's entry needs, the + * replay cannot go out, and it must not take that live request's reservation + * over either. */ +TEST(resume_drops_replay_whose_id_is_taken_on_this_connection) +{ + int rc; + MqttConnect connect; + MqttPublish publish; + static byte payload[] = "hello"; + + rc = test_init_client(); + ASSERT_EQ(MQTT_CODE_SUCCESS, rc); +#ifdef WOLFMQTT_V5 + test_client.protocol_level = MQTT_CONNECT_PROTOCOL_LEVEL_4; +#endif + ASSERT_EQ(MQTT_CODE_SUCCESS, run_initial_connect(&connect)); + + init_qos_publish(&publish, MQTT_QOS_1, 0x1234, "sensor/temp", + payload, (word32)(sizeof(payload) - 1)); + rc = run_publish_unacked(&publish); + ASSERT_EQ(MQTT_CODE_ERROR_NETWORK, rc); + ASSERT_EQ(MQTT_CODE_SUCCESS, MqttClient_NetDisconnect(&test_client)); + ASSERT_EQ(0x1234, (int)test_client.replay[0].packet_id); + + g_seed_slot = -1; + g_frames_written = 0; + rc = run_reconnect_with(&connect, 1, mock_net_write_reserve_subscribe_id); + ASSERT_EQ(MQTT_CODE_SUCCESS, rc); + + /* The CONNECT only: the entry could not be re-sent. */ + ASSERT_EQ(1, g_frames_written); + /* Dropped rather than replayed behind the live request. */ + ASSERT_EQ(0, (int)test_client.replay[0].packet_id); + /* And the SUBSCRIBE still owns its reservation, unchanged. */ + ASSERT_EQ(0x1234, (int)test_client.send_inflight[0].packet_id); + ASSERT_EQ(MQTT_PACKET_TYPE_SUBSCRIBE_ACK, + (int)test_client.send_inflight[0].ack_type); +} + #if defined(WOLFMQTT_MULTITHREAD) && defined(WOLFMQTT_POSIX_SEMAPHORES) && \ !defined(WOLFMQTT_NO_COND_SIGNAL) && !defined(WOLFMQTT_STATIC_MEMORY) /* A send thread reaches the same replay slots under the client lock, and both @@ -7870,6 +8059,9 @@ void run_mqtt_client_tests(void) RUN_TEST(reconnect_replays_unacked_qos1_publish); RUN_TEST(reconnect_without_session_present_replays_nothing); RUN_TEST(fresh_session_reset_keeps_publish_from_this_connection); + RUN_TEST(fresh_session_reset_drops_old_publish_behind_subscribe_id); + RUN_TEST(resume_does_not_replay_publish_from_this_connection); + RUN_TEST(resume_drops_replay_whose_id_is_taken_on_this_connection); #if defined(WOLFMQTT_MULTITHREAD) && defined(WOLFMQTT_POSIX_SEMAPHORES) && \ !defined(WOLFMQTT_NO_COND_SIGNAL) && !defined(WOLFMQTT_STATIC_MEMORY) RUN_TEST(fresh_session_reset_frees_replay_under_client_lock); @@ -7979,6 +8171,9 @@ void run_mqtt_client_tests(void) (WOLFMQTT_MAX_QOS >= 1) RUN_TEST(cancel_refuses_message_while_response_decodes); #endif +#if defined(WOLFMQTT_MULTITHREAD) && (WOLFMQTT_MAX_QOS >= 1) + RUN_TEST(write_failure_cleanup_unlinks_claimed_pending_response); +#endif #endif RUN_TEST(publish_qos1_v5_write_failure_restores_recv_quota); #if defined(WOLFMQTT_MULTITHREAD) && defined(WOLFMQTT_NONBLOCK) && \ diff --git a/wolfmqtt/mqtt_client.h b/wolfmqtt/mqtt_client.h index 465380235..1ec1b582b 100644 --- a/wolfmqtt/mqtt_client.h +++ b/wolfmqtt/mqtt_client.h @@ -335,6 +335,12 @@ typedef struct _MqttReplayMsg { * for the pool or came from a payload callback, which has nothing to * copy; such an entry can still replay a PUBREL but not a PUBLISH. */ byte haveCopy; + /* Published on the Network Connection being established rather than + * carried from the Session that ended. Sends are admitted once CONNECT is + * on the wire, so both kinds can be in the pool when CONNACK arrives, and + * only the latter are the old Session's to discard or re-send. Cleared on + * every entry when a handshake starts. */ + byte onThisConn; } MqttReplayMsg; #endif /* !WOLFMQTT_NO_SESSION_REPLAY */ From 908f9512e4cdbee3f0b7c855a575e89694bec033 Mon Sep 17 00:00:00 2001 From: Kareem Date: Tue, 6 Oct 2026 11:59:50 -0700 Subject: [PATCH 7/8] Fix compilation error --- tests/test_broker_connect.c | 55 +++++++++++++++++++------------------ 1 file changed, 28 insertions(+), 27 deletions(-) diff --git a/tests/test_broker_connect.c b/tests/test_broker_connect.c index 9b69cc6d3..2a7a9bde4 100644 --- a/tests/test_broker_connect.c +++ b/tests/test_broker_connect.c @@ -9231,31 +9231,6 @@ static int persist_order_put(void* ctx, byte ns, const byte* key, return MQTT_CODE_SUCCESS; } -/* Keyed variant: a record written again under the same key replaces it, as a - * real backend does. persist_order_put appends, which would turn a re-write - * into a second queue entry. */ -static int persist_keyed_put(void* ctx, byte ns, const byte* key, - word16 key_len, const byte* blob, word32 blob_len) -{ - PersistOrderStore* store = (PersistOrderStore*)ctx; - int i; - - if (ns == BROKER_PERSIST_NS_OUTQ) { - for (i = 0; i < store->outq_count; i++) { - if (store->outq[i].key_len == key_len && - XMEMCMP(store->outq[i].key, key, key_len) == 0) { - if (blob_len > sizeof(store->outq[i].blob)) { - return MQTT_CODE_ERROR_OUT_OF_BUFFER; - } - XMEMCPY(store->outq[i].blob, blob, blob_len); - store->outq[i].blob_len = blob_len; - return MQTT_CODE_SUCCESS; - } - } - } - return persist_order_put(ctx, ns, key, key_len, blob, blob_len); -} - static int persist_order_get(void* ctx, byte ns, const byte* key, word16 key_len, byte* out, word32* inout_len) { @@ -9437,6 +9412,33 @@ TEST(persist_partial_publish_restart_keeps_dup) MqttBroker_Free(&restored); } +#endif + +/* Keyed variant: a record written again under the same key replaces it, as a + * real backend does. persist_order_put appends, which would turn a re-write + * into a second queue entry. */ +static int persist_keyed_put(void* ctx, byte ns, const byte* key, + word16 key_len, const byte* blob, word32 blob_len) +{ + PersistOrderStore* store = (PersistOrderStore*)ctx; + int i; + + if (ns == BROKER_PERSIST_NS_OUTQ) { + for (i = 0; i < store->outq_count; i++) { + if (store->outq[i].key_len == key_len && + XMEMCMP(store->outq[i].key, key, key_len) == 0) { + if (blob_len > sizeof(store->outq[i].blob)) { + return MQTT_CODE_ERROR_OUT_OF_BUFFER; + } + XMEMCPY(store->outq[i].blob, blob, blob_len); + store->outq[i].blob_len = blob_len; + return MQTT_CODE_SUCCESS; + } + } + } + return persist_order_put(ctx, ns, key, key_len, blob, blob_len); +} + /* The same requirement when the Session hand-off takes the queue mid-write. * BrokerOrphan_Take shadow-writes each entry before that write returns, so the * stored copy still says it was never sent; the completed delivery has to be @@ -9537,7 +9539,6 @@ TEST(persist_handoff_sent_qos1_restart_keeps_dup) MqttBroker_Stop(&restored); MqttBroker_Free(&restored); } -#endif /* Non-persisted QoS 0 nodes still occupy FIFO positions. Assigning sequence * numbers only to durable nodes can make a later offline QoS 1 enqueue reuse @@ -11225,8 +11226,8 @@ int main(int argc, char** argv) RUN_TEST(persist_restore_packet_id_wrap_preserves_fifo); #ifdef WOLFMQTT_NONBLOCK RUN_TEST(persist_partial_publish_restart_keeps_dup); - RUN_TEST(persist_handoff_sent_qos1_restart_keeps_dup); #endif + RUN_TEST(persist_handoff_sent_qos1_restart_keeps_dup); RUN_TEST(persist_mixed_qos_queue_preserves_fifo); RUN_TEST(orphan_reclaim_keeps_persisted_outq); #endif From c805c8d7f65fe45d06de188de30afb4d292545f7 Mon Sep 17 00:00:00 2001 From: Kareem Date: Tue, 6 Oct 2026 14:56:56 -0700 Subject: [PATCH 8/8] Code review feedback --- src/mqtt_client.c | 20 ++++++++++---- tests/test_broker_connect.c | 8 ++++++ tests/test_mqtt_client.c | 55 +++++++++++++++++++++++++++++++++++++ 3 files changed, 77 insertions(+), 6 deletions(-) diff --git a/src/mqtt_client.c b/src/mqtt_client.c index 07de35393..21c77ade7 100644 --- a/src/mqtt_client.c +++ b/src/mqtt_client.c @@ -685,6 +685,9 @@ static void MqttClient_Replay_Add(MqttClient* client, MqttPublish* publish) slot->packet_id = publish->packet_id; slot->qos = (byte)publish->qos; slot->retain = publish->retain; + /* Before the returns below: they leave the slot in use, and a PUBREL can + * still be recorded against it. */ + slot->onThisConn = 1; /* A streamed publish delivers its payload through a callback, so there is * nothing here to copy; the entry still tracks the QoS 2 PUBREL stage. @@ -738,7 +741,6 @@ static void MqttClient_Replay_Add(MqttClient* client, MqttPublish* publish) #endif slot->payload_len = publish->total_len; slot->haveCopy = 1; - slot->onThisConn = 1; } /* QoS 2 advanced past PUBREC, so [MQTT-4.4.0-1] replays the PUBREL rather @@ -805,13 +807,15 @@ static void MqttClient_Replay_AddSafe(MqttClient* client, MqttPublish* publish) } /* Everything retained so far belongs to the Session that is ending; a send - * admitted once CONNECT is on the wire tags itself. */ -static void MqttClient_Replay_NewConnSafe(MqttClient* client) + * admitted once CONNECT is on the wire tags itself. Reports a lock failure: + * stale tags would have the handshake skip or keep the wrong entries. */ +static int MqttClient_Replay_NewConnSafe(MqttClient* client) { int i; #ifdef WOLFMQTT_MULTITHREAD - if (wm_SemLock(&client->lockClient) != MQTT_CODE_SUCCESS) { - return; + int rc = wm_SemLock(&client->lockClient); + if (rc != MQTT_CODE_SUCCESS) { + return rc; } #endif for (i = 0; i < MQTT_MAX_REPLAY_MSGS; i++) { @@ -820,6 +824,7 @@ static void MqttClient_Replay_NewConnSafe(MqttClient* client) #ifdef WOLFMQTT_MULTITHREAD wm_SemUnlock(&client->lockClient); #endif + return MQTT_CODE_SUCCESS; } /* Discard the retained copies of a Session the server did not resume, keeping @@ -3281,7 +3286,10 @@ int MqttClient_Connect(MqttClient *client, MqttConnect *mc_connect) * every build. */ MqttClient_SendIdsReset(client); #ifndef WOLFMQTT_NO_SESSION_REPLAY - MqttClient_Replay_NewConnSafe(client); + rc = MqttClient_Replay_NewConnSafe(client); + if (rc != MQTT_CODE_SUCCESS) { + return rc; + } #endif #ifdef WOLFMQTT_V5 diff --git a/tests/test_broker_connect.c b/tests/test_broker_connect.c index 2a7a9bde4..7bdb8a01b 100644 --- a/tests/test_broker_connect.c +++ b/tests/test_broker_connect.c @@ -2150,6 +2150,7 @@ TEST(drain_stops_when_session_handoff_takes_queue) MqttBroker_Free(&broker); } +#if WOLFMQTT_MAX_QOS >= 1 /* The QoS > 0 half. A completed at-least-once delivery is still unacknowledged * so it stays with the carrier, but the subscriber has seen it, so the resume * copy must be marked a duplicate [MQTT-3.3.1-1]. */ @@ -2225,6 +2226,7 @@ TEST(handoff_marks_sent_qos1_for_redelivery) MqttBroker_Stop(&broker); MqttBroker_Free(&broker); } +#endif /* WOLFMQTT_MAX_QOS >= 1 */ #ifdef WOLFMQTT_BROKER_WILL /* Same hazard on the Will fan-out, which runs its own subscription walk. The @@ -9439,6 +9441,7 @@ static int persist_keyed_put(void* ctx, byte ns, const byte* key, return persist_order_put(ctx, ns, key, key_len, blob, blob_len); } +#if WOLFMQTT_MAX_QOS >= 1 /* The same requirement when the Session hand-off takes the queue mid-write. * BrokerOrphan_Take shadow-writes each entry before that write returns, so the * stored copy still says it was never sent; the completed delivery has to be @@ -9539,6 +9542,7 @@ TEST(persist_handoff_sent_qos1_restart_keeps_dup) MqttBroker_Stop(&restored); MqttBroker_Free(&restored); } +#endif /* WOLFMQTT_MAX_QOS >= 1 */ /* Non-persisted QoS 0 nodes still occupy FIFO positions. Assigning sequence * numbers only to durable nodes can make a later offline QoS 1 enqueue reuse @@ -10995,7 +10999,9 @@ int main(int argc, char** argv) RUN_TEST(online_qos1_at_cap_keeps_subscriber); RUN_TEST(outbound_packet_ids_are_scoped_per_session); RUN_TEST(drain_stops_when_session_handoff_takes_queue); +#if WOLFMQTT_MAX_QOS >= 1 RUN_TEST(handoff_marks_sent_qos1_for_redelivery); +#endif #ifdef WOLFMQTT_NONBLOCK RUN_TEST(outbound_queue_short_write_resumes_on_next_step); #ifdef WOLFMQTT_BROKER_RETAINED @@ -11227,7 +11233,9 @@ int main(int argc, char** argv) #ifdef WOLFMQTT_NONBLOCK RUN_TEST(persist_partial_publish_restart_keeps_dup); #endif +#if WOLFMQTT_MAX_QOS >= 1 RUN_TEST(persist_handoff_sent_qos1_restart_keeps_dup); +#endif RUN_TEST(persist_mixed_qos_queue_preserves_fifo); RUN_TEST(orphan_reclaim_keeps_persisted_outq); #endif diff --git a/tests/test_mqtt_client.c b/tests/test_mqtt_client.c index 55439e6db..934bf7f7e 100644 --- a/tests/test_mqtt_client.c +++ b/tests/test_mqtt_client.c @@ -4039,6 +4039,60 @@ TEST(reconnect_does_not_replay_streamed_publish) ASSERT_EQ(MQTT_PACKET_TYPE_CONNECT, g_replay_types[0]); } +/* MqttClient_Replay_Add gives up on the copy for a streamed publish, a v5 + * publish carrying properties, and an oversized one, but the slot it already + * claimed stays in use and a PUBREL can still be recorded against it. The tag + * that says which connection the entry belongs to therefore has to be set + * before those returns, or a send admitted during a handshake is mistaken for + * the previous Session's and re-sent. */ +TEST(replay_uncopied_slot_still_tagged_to_this_connection) +{ + int rc; + int i; + int slot = -1; + MqttConnect connect; + MqttPublish publish; + + rc = test_init_client(); + ASSERT_EQ(MQTT_CODE_SUCCESS, rc); +#ifdef WOLFMQTT_V5 + test_client.protocol_level = MQTT_CONNECT_PROTOCOL_LEVEL_4; +#endif + ASSERT_EQ(MQTT_CODE_SUCCESS, run_initial_connect(&connect)); + + XMEMSET(g_replay_stream_buf, 0, sizeof(g_replay_stream_buf)); + test_net.write = mock_net_write_accept; + test_net.read = mock_net_read; /* network error: no PUBACK */ + + /* Streamed: buffer_len is one chunk, total_len the whole message, so the + * store has nothing to copy. */ + XMEMSET(&publish, 0, sizeof(publish)); + publish.qos = MQTT_QOS_1; + publish.packet_id = 0x4321; + publish.topic_name = "sensor/temp"; + publish.buffer = g_replay_stream_buf; + publish.buffer_len = (word32)REPLAY_STREAM_CHUNK; + publish.total_len = (word32)REPLAY_BIG_PAYLOAD_LEN; + + rc = MQTT_CODE_CONTINUE; + for (i = 0; i < 20 && rc == MQTT_CODE_CONTINUE; i++) { + rc = MqttClient_Publish_ex(&test_client, &publish, replay_stream_cb); + } + ASSERT_EQ(MQTT_CODE_ERROR_NETWORK, rc); + + for (i = 0; i < MQTT_MAX_REPLAY_MSGS; i++) { + if (test_client.replay[i].packet_id == 0x4321) { + slot = i; + break; + } + } + /* The slot is in use even though the copy was skipped. */ + ASSERT_TRUE(slot >= 0); + ASSERT_EQ(0, (int)test_client.replay[slot].haveCopy); + /* And it is tagged to this connection. */ + ASSERT_EQ(1, (int)test_client.replay[slot].onThisConn); +} + /* An entry the replay cannot rebuild is dropped from the Session, so nothing * will ever acknowledge it and the Packet Identifier the reconnect reserved for * it must be given back [MQTT-2.3.1-3]. Left reserved, the next legitimate @@ -8073,6 +8127,7 @@ void run_mqtt_client_tests(void) RUN_TEST(reconnect_drops_replay_larger_than_tx_buf); RUN_TEST(reconnect_replays_at_most_pool_size); RUN_TEST(reconnect_does_not_replay_streamed_publish); + RUN_TEST(replay_uncopied_slot_still_tagged_to_this_connection); RUN_TEST(reconnect_dropped_replay_frees_packet_id); RUN_TEST(reconnect_with_new_client_id_does_not_replay); RUN_TEST(reconnect_with_hash_colliding_client_id_does_not_replay);