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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
87 changes: 87 additions & 0 deletions src/mqtt_broker.c
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Comment thread
kareem-wolfssl marked this conversation as resolved.
}
/* 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
Expand Down Expand Up @@ -3462,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;
}
Comment on lines +3526 to +3529
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)
Expand Down Expand Up @@ -5882,6 +5965,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);
Expand Down
184 changes: 133 additions & 51 deletions src/mqtt_client.c
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -797,6 +800,37 @@ 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. */
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].packet_id != 0 &&
MqttClient_SendIds_Find(client,
client->replay[i].packet_id) >= 0) {
continue; /* published on this connection */
Comment on lines +820 to +823
}
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
Expand Down Expand Up @@ -861,8 +895,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;
Expand All @@ -871,12 +906,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
Expand Down Expand Up @@ -906,6 +935,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
Expand Down Expand Up @@ -2365,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;
}
Expand Down Expand Up @@ -3595,15 +3653,25 @@ 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;

/* 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) {
continue;
Expand All @@ -3617,14 +3685,17 @@ 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 :
((client->replay[i].qos == MQTT_QOS_2) ?
MQTT_PACKET_TYPE_PUBLISH_COMP :
MQTT_PACKET_TYPE_PUBLISH_ACK));
Comment on lines +3688 to 3694
}
#ifdef WOLFMQTT_MULTITHREAD
wm_SemUnlock(&client->lockClient);
#endif
client->replayIdx = 0;
mc_connect->stat.write = MQTT_MSG_PAYLOAD;
rc = MqttClient_ReplaySession(client, mc_connect);
Expand Down Expand Up @@ -5335,6 +5406,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) {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

New CONTINUE return from CancelMessage is ignored by internal error paths, leaving a stale pendResp linked · API contract violations

MqttClient_CancelMessage can now return MQTT_CODE_CONTINUE without unlinking the pendResp or resetting stat. The failure paths in Publish (4247, 4294), Subscribe (4543), Unsubscribe (4698), Ping (4862) and Connect (3457) ignore that return and report the error anyway. The caller can then free or reuse an object that is still on firstPendResp, which leaves a dangling list node.

Suggested fix: Update every internal caller of MqttClient_CancelMessage to handle MQTT_CODE_CONTINUE. Either wait until the reader sets packetDone and cancel again, or keep the object owned instead of returning an error.

Related known findings (similar but distinct; listed for context, not part of this finding)

  • F-13261 (open): File/function: same function MqttClient_CancelMessage as F-13261. Operation: F-13261 is the active-read release path omitting an rx_buf scrub; candidate is the packetProcessing-and-not-done early return skipping RespList_Remove when callers ignore MQTT_CODE_CONTINUE. Root cause: F-13261 is missing buffer zeroization (info disclosure); candidate is a missing caller-side contract check leaving a stale linked-list node (dangling pointer / potential use-after-free). Patch: one requires clearing rx_buf

wm_SemUnlock(&client->lockClient);
return MQTT_CODE_CONTINUE;
Comment thread
kareem-wolfssl marked this conversation as resolved.
}
/* 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;
Expand Down Expand Up @@ -5377,46 +5499,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
Expand Down
3 changes: 2 additions & 1 deletion tests/include.am
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading