diff --git a/mqtt_client.c b/mqtt_client.c index bf8bbde..fcca7da 100644 --- a/mqtt_client.c +++ b/mqtt_client.c @@ -77,6 +77,16 @@ static int esp_mqtt_handle_transport_read_error(int err, esp_mqtt_client_handle_ return -2; } +// Reset per-message pending state before composing a new outbound control packet +static inline void mqtt_reset_pending_message(esp_mqtt_client_handle_t client) +{ + client->mqtt_state.pending_msg_id = 0; + client->mqtt_state.pending_msg_type = 0; + client->mqtt_state.pending_publish_qos = 0; + client->mqtt_state.connection.outbound_message.fragmented_msg_total_length = 0; + client->mqtt_state.connection.outbound_message.fragmented_msg_data_offset = 0; +} + #if MQTT_ENABLE_SSL enum esp_mqtt_ssl_cert_key_api { MQTT_SSL_DATA_API_CA_CERT, @@ -837,6 +847,7 @@ static esp_err_t esp_mqtt_connect(esp_mqtt_client_handle_t client, int timeout_m { int read_len, connect_rsp_code = 0; client->wait_for_ping_resp = false; + mqtt_reset_pending_message(client); if (client->mqtt_state.connection.information.protocol_ver == MQTT_PROTOCOL_V_5) { #ifdef MQTT_PROTOCOL_5 @@ -1974,7 +1985,8 @@ static void esp_mqtt_task(void *pv) } } - if (client->mqtt_state.pending_publish_qos > 0) { + if (client->mqtt_state.pending_publish_qos > 0 && + mqtt_get_type(client->mqtt_state.connection.outbound_message.data) == MQTT_MSG_TYPE_PUBLISH) { outbox_set_pending(client->outbox, client->mqtt_state.pending_msg_id, TRANSMITTED); #ifdef MQTT_PROTOCOL_5 @@ -1995,7 +2007,9 @@ static void esp_mqtt_task(void *pv) if (mqtt_resend_queued(client, item) == ESP_OK) { #ifdef MQTT_PROTOCOL_5 - if (client->mqtt_state.connection.information.protocol_ver == MQTT_PROTOCOL_V_5) { + if (client->mqtt_state.connection.information.protocol_ver == MQTT_PROTOCOL_V_5 && + client->mqtt_state.pending_publish_qos > 0 && + mqtt_get_type(client->mqtt_state.connection.outbound_message.data) == MQTT_MSG_TYPE_PUBLISH) { esp_mqtt5_increment_packet_counter(client); } @@ -2009,7 +2023,12 @@ static void esp_mqtt_task(void *pv) if (mqtt_resend_pubrel(client, item) == ESP_OK) { #ifdef MQTT_PROTOCOL_5 - if (client->mqtt_state.connection.information.protocol_ver == MQTT_PROTOCOL_V_5) { + // Do not count PUBREL as a new inflight PUBLISH + // Only PUBLISH QoS>0 contributes to inflight limitation + // (outbound_message here is PUBREL, so this condition will be false) + if (client->mqtt_state.connection.information.protocol_ver == MQTT_PROTOCOL_V_5 && + client->mqtt_state.pending_publish_qos > 0 && + mqtt_get_type(client->mqtt_state.connection.outbound_message.data) == MQTT_MSG_TYPE_PUBLISH) { esp_mqtt5_increment_packet_counter(client); } @@ -2324,6 +2343,8 @@ int esp_mqtt_client_unsubscribe(esp_mqtt_client_handle_t client, const char *top } MQTT_API_LOCK(client); + // Reset pending state to avoid inheriting previous PUBLISH QoS or type + mqtt_reset_pending_message(client); if (client->mqtt_state.connection.information.protocol_ver == MQTT_PROTOCOL_V_5) { #ifdef MQTT_PROTOCOL_5