diff --git a/src/mqtt_broker.c b/src/mqtt_broker.c index fcea7ee23..bc3b830bc 100644 --- a/src/mqtt_broker.c +++ b/src/mqtt_broker.c @@ -1999,6 +1999,35 @@ static void BrokerClient_FreeOutQueue(BrokerClient* bc) bc->out_q_pending_len = 0; } +#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; +} + /* 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 +2088,10 @@ 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, cur, + wr_rc == rel_rc, rel_rc)) { + return sent; + } if (wr_rc == MQTT_CODE_CONTINUE) { bc->out_q_pending_len = rel_rc; return wr_rc; @@ -2152,6 +2185,10 @@ static int BrokerClient_DrainOutQueue(BrokerClient* bc) { int wr_rc; wr_rc = MqttPacket_Write(&bc->client, bc->tx_buf, 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 * an application payload) once the buffer is idle. Skip only the * MQTT_CODE_CONTINUE case, where a non-blocking or TLS-async send @@ -3462,6 +3499,57 @@ 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; + #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) { + 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) @@ -5882,6 +5970,10 @@ static void BrokerClient_PublishWill(MqttBroker* broker, BrokerClient* bc) (int)bc->sock, BrokerLog_Sanitize(bc->will_topic), (unsigned)bc->will_payload_len); + /* 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, bc->will_payload, bc->will_payload_len, bc->will_qos, bc->will_retain); diff --git a/src/mqtt_client.c b/src/mqtt_client.c index 1ce10bb4f..21c77ade7 100644 --- a/src/mqtt_client.c +++ b/src/mqtt_client.c @@ -88,6 +88,11 @@ 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 #ifdef WOLFMQTT_MULTITHREAD #ifdef WOLFMQTT_USER_THREADING @@ -680,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. @@ -754,6 +762,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; } @@ -797,6 +806,53 @@ static void MqttClient_Replay_AddSafe(MqttClient* client, MqttPublish* publish) #endif } +/* Everything retained so far belongs to the Session that is ending; a send + * 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 + int rc = wm_SemLock(&client->lockClient); + if (rc != MQTT_CODE_SUCCESS) { + return rc; + } +#endif + for (i = 0; i < MQTT_MAX_REPLAY_MSGS; i++) { + client->replay[i].onThisConn = 0; + } +#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 + * 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; +#ifdef WOLFMQTT_MULTITHREAD + int 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].onThisConn) { + continue; + } + MqttClient_Replay_FreeSlot(&client->replay[i]); + } + client->replayIdx = MQTT_MAX_REPLAY_MSGS; +#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 +917,9 @@ 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 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) { int rc = MQTT_CODE_SUCCESS; @@ -871,12 +928,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 +957,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 @@ -2365,6 +2435,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; } @@ -2913,7 +2993,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; } @@ -3204,6 +3285,12 @@ 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 + rc = MqttClient_Replay_NewConnSafe(client); + if (rc != MQTT_CODE_SUCCESS) { + return rc; + } +#endif #ifdef WOLFMQTT_V5 #ifdef WOLFMQTT_MULTITHREAD @@ -3396,7 +3483,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; } @@ -3595,7 +3683,10 @@ int MqttClient_Connect(MqttClient *client, MqttConnect *mc_connect) 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); + rc = MqttClient_Replay_ResetSafe(client); + if (rc != MQTT_CODE_SUCCESS) { + return rc; + } } else { int i; @@ -3603,9 +3694,18 @@ 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]. 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) { + return rc; + } +#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 && @@ -3617,7 +3717,15 @@ int MqttClient_Connect(MqttClient *client, MqttConnect *mc_connect) MqttClient_Replay_FreeSlot(&client->replay[i]); continue; } - (void)MqttClient_SendIdReserve(client, + 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 ? MQTT_PACKET_TYPE_PUBLISH_COMP : @@ -3625,6 +3733,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); @@ -4173,7 +4284,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 @@ -4220,7 +4332,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]. */ @@ -4469,7 +4582,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 @@ -4624,7 +4737,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 @@ -4788,7 +4902,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; } @@ -5312,16 +5426,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) { @@ -5335,6 +5452,58 @@ 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 (!force && 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 +5546,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 @@ -5438,6 +5567,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/include.am b/tests/include.am index 9d79bef16..69577694c 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_broker_alloc.h $(AM_CPPFLAGS) tests_unit_tests_LDADD = $(PTHREAD_LIBS) # Exercise peer-name verification with a trusted local certificate and an diff --git a/tests/test_broker_connect.c b/tests/test_broker_connect.c index 37557a96f..7bdb8a01b 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 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; + 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,235 @@ TEST(fanout_survives_reentrant_sub_free) MqttBroker_Free(&broker); } +/* 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) +{ + 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 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; + 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++; + +#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) +{ + 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' + }; + /* 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. */ + 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 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 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(0, o->out_q_count); + ASSERT_NULL(o->out_q_head); + ASSERT_NULL(o->out_q_tail); + + /* 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 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); + 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); +} + +#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]. */ +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); +} +#endif /* WOLFMQTT_MAX_QOS >= 1 */ + #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, @@ -2053,7 +2302,110 @@ TEST(will_fanout_survives_reentrant_sub_free) MqttBroker_Stop(&broker); MqttBroker_Free(&broker); } + +/* 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; +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 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)); + + 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 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]. */ + 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) @@ -9061,8 +9413,137 @@ TEST(persist_partial_publish_restart_keeps_dup) MqttBroker_Stop(&restored); 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); +} + +#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 + * 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 /* 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 * an earlier sequence and restore ahead of older messages. */ @@ -10499,6 +10980,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); @@ -10516,6 +10998,10 @@ 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); +#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 @@ -10747,6 +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 1cea25033..934bf7f7e 100644 --- a/tests/test_mqtt_client.c +++ b/tests/test_mqtt_client.c @@ -48,6 +48,44 @@ static MqttNet test_net; static byte test_tx_buf[TEST_TX_BUF_SIZE]; static byte test_rx_buf[TEST_RX_BUF_SIZE]; +/* 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_broker_malloc(size_t size) +{ + return malloc(size); +} + +/* 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) && \ + !defined(WOLFMQTT_NO_COND_SIGNAL) + return test_client.lockClient.lockCount; +#else + return -1; +#endif +} + +void wolfmqtt_test_broker_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, @@ -1783,6 +1821,120 @@ 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 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; + 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); + } + /* 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; + + /* Claim the entry as MqttClient_WaitType does before it decodes. */ + 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 */ +#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 @@ -3306,6 +3458,272 @@ 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, 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_slot >= 0 && buf_len > 0 && + (buf[0] >> 4) == MQTT_PACKET_TYPE_CONNECT) { + 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; + } + g_seed_slot = -1; + } + 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_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(-1, g_seed_slot); + + /* 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); +} + +/* 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 + * 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; + 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)); + + /* 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) @@ -3621,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 @@ -7640,6 +8112,14 @@ 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); + 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); +#endif RUN_TEST(reconnect_does_not_replay_acked_publish); #if WOLFMQTT_MAX_QOS >= 2 RUN_TEST(reconnect_replays_unacked_pubrel); @@ -7647,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); @@ -7741,6 +8222,13 @@ 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 +#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 562019a0c..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 */ @@ -888,9 +894,12 @@ 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 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 or MQTT_CODE_ERROR_* + * \return MQTT_CODE_SUCCESS, MQTT_CODE_CONTINUE or MQTT_CODE_ERROR_* (see enum MqttPacketResponseCodes) */ WOLFMQTT_API int MqttClient_CancelMessage(