F-10352 - Report unknown PUBREC packet identifiers

pull/577/head
Aidan Garske 2026-08-21 10:42:56 -07:00
parent d8f0d57f25
commit c85b4878a4
2 changed files with 179 additions and 0 deletions

View File

@ -1678,6 +1678,18 @@ wait_again:
{
MqttPublishResp resp;
MqttPacketType use_packet_type;
#ifdef WOLFMQTT_V5
int pubrec_tracked = 0;
/* A synchronous QoS 2 publish waits for PUBCOMP while processing
* the intermediate PUBREC for the same Packet Identifier. */
if (packet_type == MQTT_PACKET_TYPE_PUBLISH_REC &&
client->protocol_level >= MQTT_CONNECT_PROTOCOL_LEVEL_5 &&
wait_type == MQTT_PACKET_TYPE_PUBLISH_COMP &&
wait_packet_id != 0 && wait_packet_id == packet_id) {
pubrec_tracked = 1;
}
#endif
/* Determine if we received data for this request */
if ((wait_type == MQTT_PACKET_TYPE_ANY ||
@ -1738,6 +1750,30 @@ wait_again:
waitMatchFound = 0;
}
}
#ifdef WOLFMQTT_V5
/* A reader thread receives the intermediate PUBREC while the
* publisher's pending entry is keyed by the final PUBCOMP. */
if (!pubrec_tracked &&
packet_type == MQTT_PACKET_TYPE_PUBLISH_REC &&
client->protocol_level >=
MQTT_CONNECT_PROTOCOL_LEVEL_5) {
MqttPendResp* qos2Resp;
for (qos2Resp = client->firstPendResp;
qos2Resp != NULL; qos2Resp = qos2Resp->next) {
if (qos2Resp->packet_type ==
MQTT_PACKET_TYPE_PUBLISH_COMP &&
qos2Resp->packet_id == packet_id &&
!qos2Resp->packetDone) {
pubrec_tracked = 1;
/* This belongs to the publisher's flow, not the
* generic WaitMessage caller. */
waitMatchFound = 0;
break;
}
}
}
#endif
wm_SemUnlock(&client->lockClient);
}
else {
@ -1766,6 +1802,18 @@ wait_again:
rc = MqttClient_HandlePacket(client, use_packet_type,
use_packet_obj, &resp, timeout_ms);
#ifdef WOLFMQTT_V5
/* [MQTT-3.6.2.1] An unsolicited PUBREC is answered with PUBREL
* reason 0x92 instead of falsely advancing an unknown flow. */
if (rc >= 0 &&
packet_type == MQTT_PACKET_TYPE_PUBLISH_REC &&
client->protocol_level >= MQTT_CONNECT_PROTOCOL_LEVEL_5 &&
resp.packet_type == MQTT_PACKET_TYPE_PUBLISH_REL &&
!pubrec_tracked) {
resp.reason_code = MQTT_REASON_PACKET_ID_NOT_FOUND;
}
#endif
/* if using the shared packet object, make sure the original
* state is correct for publish payload 2 (continued) */
if (use_packet_obj != NULL && use_packet_obj != mms_stat &&

View File

@ -2667,6 +2667,7 @@ TEST(publish_qos2_v5_broker_rejection_returns_publish_rejected)
TEST(publish_qos2_v5_success_returns_success)
{
int rc;
int i;
MqttPublish publish;
static byte payload[] = "hello";
/* PUBREC success (no reason byte): type=0x50, remain=2, packet_id=14.
@ -2675,6 +2676,11 @@ TEST(publish_qos2_v5_success_returns_success)
0x50, 0x02, 0x00, 0x0E,
0x70, 0x02, 0x00, 0x0E
};
static const byte success_pubrel[] = { 0x62, 0x02, 0x00, 0x0E };
static const byte stray_pubrec[] = { 0x50, 0x02, 0x00, 0x0E };
static const byte unknown_pubrel[] = {
0x62, 0x03, 0x00, 0x0E, MQTT_REASON_PACKET_ID_NOT_FOUND
};
XMEMSET(&publish, 0, sizeof(publish));
publish.qos = MQTT_QOS_2;
@ -2690,6 +2696,25 @@ TEST(publish_qos2_v5_success_returns_success)
ASSERT_EQ(MQTT_CODE_SUCCESS, rc);
ASSERT_EQ(MQTT_REASON_SUCCESS, publish.resp.reason_code);
ASSERT_TRUE(g_pubrel_written);
ASSERT_EQ((int)sizeof(success_pubrel), connect_mock_xfer);
ASSERT_MEM_EQ(success_pubrel, connect_mock_sent, sizeof(success_pubrel));
/* The matching PUBCOMP completed the tracked flow. Repeating its PUBREC
* through the generic receive path must now report the ID as unknown. */
XMEMCPY(g_canned_buf, stray_pubrec, sizeof(stray_pubrec));
g_canned_len = (int)sizeof(stray_pubrec);
g_canned_pos = 0;
connect_mock_xfer = 0;
XMEMSET(connect_mock_sent, 0, sizeof(connect_mock_sent));
rc = MQTT_CODE_CONTINUE;
for (i = 0; i < 20 && rc == MQTT_CODE_CONTINUE; i++) {
rc = MqttClient_WaitMessage(&test_client, TEST_CMD_TIMEOUT_MS);
}
ASSERT_EQ(MQTT_CODE_SUCCESS, rc);
ASSERT_EQ((int)sizeof(unknown_pubrel), connect_mock_xfer);
ASSERT_MEM_EQ(unknown_pubrel, connect_mock_sent, sizeof(unknown_pubrel));
}
/* The primary QoS 2 rejection point is the PUBREC: a v5 broker reports
@ -3116,6 +3141,72 @@ TEST(publish_qos2_v5_pubrec_rejection_multithread_reader)
/* The publisher's own struct is NOT updated on this path. */
ASSERT_EQ(MQTT_REASON_SUCCESS, publish.resp.reason_code);
}
/* A reader thread must recognize the PUBCOMP-keyed pending entry as the
* outbound state for an intermediate PUBREC, then stop recognizing it after
* the matching PUBCOMP completes and the publisher consumes that response. */
TEST(publish_qos2_v5_pubrec_state_multithread_reader)
{
int rc;
int i;
static MqttPublish publish;
static byte payload[] = "hello";
static const byte pubrec[] = { 0x50, 0x02, 0x00, 0x0F };
static const byte pubcomp[] = { 0x70, 0x02, 0x00, 0x0F };
static const byte success_pubrel[] = { 0x62, 0x02, 0x00, 0x0F };
static const byte unknown_pubrel[] = {
0x62, 0x03, 0x00, 0x0F, MQTT_REASON_PACKET_ID_NOT_FOUND
};
rc = test_init_client();
ASSERT_EQ(MQTT_CODE_SUCCESS, rc);
test_client.protocol_level = MQTT_CONNECT_PROTOCOL_LEVEL_5;
test_net.write = mock_net_write_accept;
test_net.read = mock_net_read_canned;
XMEMSET(&publish, 0, sizeof(publish));
publish.qos = MQTT_QOS_2;
publish.packet_id = 15;
publish.topic_name = "test/topic";
publish.buffer = payload;
publish.total_len = (word32)(sizeof(payload) - 1);
publish.buffer_len = publish.total_len;
rc = MqttClient_Publish_WriteOnly(&test_client, &publish, NULL);
ASSERT_EQ(MQTT_CODE_CONTINUE, rc);
XMEMCPY(g_canned_buf, pubrec, sizeof(pubrec));
g_canned_len = (int)sizeof(pubrec);
g_canned_pos = 0;
connect_mock_xfer = 0;
rc = MqttClient_WaitMessage(&test_client, TEST_CMD_TIMEOUT_MS);
ASSERT_EQ(MQTT_CODE_CONTINUE, rc);
ASSERT_EQ((int)sizeof(success_pubrel), connect_mock_xfer);
ASSERT_MEM_EQ(success_pubrel, connect_mock_sent, sizeof(success_pubrel));
XMEMCPY(g_canned_buf, pubcomp, sizeof(pubcomp));
g_canned_len = (int)sizeof(pubcomp);
g_canned_pos = 0;
rc = MqttClient_WaitMessage(&test_client, TEST_CMD_TIMEOUT_MS);
ASSERT_EQ(MQTT_CODE_CONTINUE, rc);
ASSERT_EQ((int)sizeof(pubcomp), g_canned_pos);
for (i = 0; i < 20 && rc == MQTT_CODE_CONTINUE; i++) {
rc = MqttClient_Publish_WriteOnly(&test_client, &publish, NULL);
}
ASSERT_EQ(MQTT_CODE_SUCCESS, rc);
ASSERT_NULL(test_client.firstPendResp);
XMEMCPY(g_canned_buf, pubrec, sizeof(pubrec));
g_canned_len = (int)sizeof(pubrec);
g_canned_pos = 0;
connect_mock_xfer = 0;
XMEMSET(connect_mock_sent, 0, sizeof(connect_mock_sent));
rc = MqttClient_WaitMessage(&test_client, TEST_CMD_TIMEOUT_MS);
ASSERT_EQ(MQTT_CODE_SUCCESS, rc);
ASSERT_EQ((int)sizeof(unknown_pubrel), connect_mock_xfer);
ASSERT_MEM_EQ(unknown_pubrel, connect_mock_sent, sizeof(unknown_pubrel));
}
#endif /* WOLFMQTT_MULTITHREAD && WOLFMQTT_NONBLOCK && WOLFMQTT_MAX_QOS >= 2 */
#ifdef WOLFMQTT_MULTITHREAD
@ -4160,6 +4251,42 @@ TEST(wait_message_pubrec_emits_pubrel)
ASSERT_EQ(MQTT_PACKET_TYPE_PUBLISH_REL, g_last_ack_written);
}
#ifdef WOLFMQTT_V5
/* [MQTT-3.6.2.1] permits 0x92 in PUBREL when the receiver has no matching
* outbound QoS 2 flow. This fixed wire fixture is an unsolicited successful
* PUBREC; the exact response must report that its Packet Identifier is unknown. */
TEST(wait_message_v5_unmatched_pubrec_reports_unknown_id)
{
int rc;
int i;
static const byte pubrec[] = { 0x50, 0x02, 0x12, 0x34 };
static const byte expected_pubrel[] = {
0x62, 0x03, 0x12, 0x34, MQTT_REASON_PACKET_ID_NOT_FOUND
};
rc = test_init_client();
ASSERT_EQ(MQTT_CODE_SUCCESS, rc);
test_client.protocol_level = MQTT_CONNECT_PROTOCOL_LEVEL_5;
test_net.write = mock_net_write_accept;
test_net.read = mock_net_read_canned;
XMEMCPY(g_canned_buf, pubrec, sizeof(pubrec));
g_canned_len = (int)sizeof(pubrec);
g_canned_pos = 0;
connect_mock_xfer = 0;
XMEMSET(connect_mock_sent, 0, sizeof(connect_mock_sent));
rc = MQTT_CODE_CONTINUE;
for (i = 0; i < 20 && rc == MQTT_CODE_CONTINUE; i++) {
rc = MqttClient_WaitMessage(&test_client, TEST_CMD_TIMEOUT_MS);
}
ASSERT_EQ(MQTT_CODE_SUCCESS, rc);
ASSERT_EQ((int)sizeof(expected_pubrel), connect_mock_xfer);
ASSERT_MEM_EQ(expected_pubrel, connect_mock_sent,
sizeof(expected_pubrel));
}
#endif /* WOLFMQTT_V5 */
TEST(wait_message_pubrel_emits_pubcomp)
{
int rc;
@ -4917,6 +5044,7 @@ void run_mqtt_client_tests(void)
#if defined(WOLFMQTT_MULTITHREAD) && defined(WOLFMQTT_NONBLOCK) && \
WOLFMQTT_MAX_QOS >= 2
RUN_TEST(publish_qos2_v5_pubrec_rejection_multithread_reader);
RUN_TEST(publish_qos2_v5_pubrec_state_multithread_reader);
#endif
#ifdef WOLFMQTT_MULTITHREAD
#ifdef WOLFMQTT_NONBLOCK
@ -4960,6 +5088,9 @@ void run_mqtt_client_tests(void)
#endif
#endif
RUN_TEST(wait_message_pubrec_emits_pubrel);
#ifdef WOLFMQTT_V5
RUN_TEST(wait_message_v5_unmatched_pubrec_reports_unknown_id);
#endif
RUN_TEST(wait_message_pubrel_emits_pubcomp);
RUN_TEST(wait_message_puback_emits_no_ack);
RUN_TEST(wait_message_pubcomp_emits_no_ack);