Skip to content
Merged
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
11 changes: 11 additions & 0 deletions BROKER.md
Original file line number Diff line number Diff line change
Expand Up @@ -111,6 +111,17 @@ When built with `WOLFMQTT_STATIC_MEMORY`, the broker uses fixed-size arrays inst
| `BROKER_TIMEOUT_MS` | 1000 | `select()` timeout |
| `BROKER_LISTEN_BACKLOG` | 128 | Listen queue depth |

Retaining a message is best-effort. A `RETAIN=1` PUBLISH whose retained copy
cannot be stored (the retained table is full at `BROKER_MAX_RETAINED`, the
payload exceeds `BROKER_MAX_PAYLOAD_LEN`, or an allocation fails) is still
delivered to every current subscriber and acknowledged with success at any QoS
or protocol level. Only the retained copy is skipped, so a later subscriber will
not receive that message until the topic is published again with room to store
it. The skipped retained copy is logged broker-side. This matches how servers
such as Mosquitto treat a full retained store: live delivery is never sacrificed
to retention. Raise `BROKER_MAX_RETAINED` (or `BROKER_MAX_PAYLOAD_LEN`) if
retained-topic capacity matters for your deployment.

The static offline queue is broker-owned fixed storage. Its dominant RAM cost
is approximately `sessions * messages * (topic length + data length)` bytes,
plus queue metadata. A new persistent CONNECT is refused when all session slots
Expand Down
14 changes: 14 additions & 0 deletions ChangeLog.md
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,20 @@
previous manual keep-alive loop under `WOLFMQTT_NO_TIME` (#501)

* API / Behavior Changes
- The broker now treats retaining a message as best-effort. When a
`RETAIN=1` PUBLISH cannot be stored (retained table full, oversized
payload, or allocation failure), the message is still delivered to all
current subscribers and acknowledged with success at every QoS and
protocol level; only the retained copy is skipped, and the skip is logged.
This matches Mosquitto and replaces the previous inconsistent handling
that could drop delivery on a retained-store failure. See `BROKER.md`.
- The client rejects an inbound v5 PUBLISH that carries a Topic Alias and no
longer advertises a nonzero Topic Alias Maximum in `CONNECT`. Inbound
alias resolution is not implemented, so the client advertises Topic Alias
Maximum 0 and a server that sends an alias anyway is treated as a protocol
error. Outbound Topic Alias (client to server) is unchanged. An
application that supplies a nonzero `MQTT_PROP_TOPIC_ALIAS_MAX` now gets
`MQTT_CODE_ERROR_PROPERTY` from `MqttClient_Connect`.
- A v5 `CONNECT` now advertises `Receive Maximum` set to
`MQTT_MAX_RECV_QOS2` (16 by default) unless the application supplied its
own `MQTT_PROP_RECEIVE_MAX`. This bounds the QoS 1 and QoS 2 PUBLISH
Expand Down
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -276,7 +276,7 @@ The following v5.0 specification features are supported by the wolfMQTT client:
* Maximum packet size
* Server assigned client identifier
* Subscription ID
* Topic Alias
* Topic Alias (outbound PUBLISH only; inbound aliases are rejected because the client advertises a Topic Alias Maximum of 0)

The v5 enabled wolfMQTT client was tested with the following MQTT v5 brokers:
* Mosquitto
Expand Down
5 changes: 3 additions & 2 deletions examples/mqttclient/mqttclient.c
Original file line number Diff line number Diff line change
Expand Up @@ -361,10 +361,11 @@ int mqttclient_test(MQTTCtx *mqttCtx)
prop->data_int = (word32)mqttCtx->max_packet_size;
}
{
/* Topic Alias Maximum */
/* Topic Alias Maximum. Advertise 0: the client does not resolve inbound
* Topic Aliases, so a conforming server must not send any. */
MqttProp* prop = MqttClient_PropsAdd(&mqttCtx->connect.props);
prop->type = MQTT_PROP_TOPIC_ALIAS_MAX;
prop->data_short = mqttCtx->topic_alias_max;
prop->data_short = 0;
}
if (mqttCtx->clean_session == 0) {
/* Session expiry interval */
Expand Down
5 changes: 3 additions & 2 deletions examples/pub-sub/mqtt-pub.c
Original file line number Diff line number Diff line change
Expand Up @@ -297,10 +297,11 @@ int pub_client(MQTTCtx *mqttCtx)
prop->data_int = (word32)mqttCtx->max_packet_size;
}
{
/* Topic Alias Maximum */
/* Topic Alias Maximum. Advertise 0: the client does not resolve inbound
* Topic Aliases, so a conforming server must not send any. */
MqttProp* prop = MqttClient_PropsAdd(&mqttCtx->connect.props);
prop->type = MQTT_PROP_TOPIC_ALIAS_MAX;
prop->data_short = mqttCtx->topic_alias_max;
prop->data_short = 0;
}
if (mqttCtx->clean_session == 0) {
/* Session expiry interval */
Expand Down
2 changes: 1 addition & 1 deletion examples/pub-sub/mqtt-sub.c
Original file line number Diff line number Diff line change
Expand Up @@ -364,7 +364,7 @@ int sub_client(MQTTCtx *mqttCtx)
/* Topic Alias Maximum */
MqttProp* prop = MqttClient_PropsAdd(&mqttCtx->connect.props);
prop->type = MQTT_PROP_TOPIC_ALIAS_MAX;
prop->data_short = mqttCtx->topic_alias_max;
prop->data_short = 0;
}
if (mqttCtx->clean_session == 0) {
/* Session expiry interval */
Expand Down
45 changes: 30 additions & 15 deletions src/mqtt_broker.c
Original file line number Diff line number Diff line change
Expand Up @@ -7114,6 +7114,19 @@ static int BrokerHandle_Connect(BrokerClient* bc, int rx_len,
ack.return_code = MQTT_CONNECT_ACK_CODE_ACCEPTED;
#ifdef WOLFMQTT_V5
ack.props = NULL;

/* Release the decoded CONNECT and Will properties before building the
* CONNACK: they share the fixed property pool, and a CONNECT carrying many
* User Properties would otherwise exhaust it and drop the mandatory
* Assigned Client Identifier added below. */
if (mc.props != NULL) {
(void)MqttProps_Free(mc.props);
mc.props = NULL;
}
if (lwt.props != NULL) {
(void)MqttProps_Free(lwt.props);
lwt.props = NULL;
}
#endif

#ifdef WOLFMQTT_V5
Expand All @@ -7132,6 +7145,17 @@ static int BrokerHandle_Connect(BrokerClient* bc, int rx_len,
prop->data_str.str = bc->client_id;
prop->data_str.len = (word16)XSTRLEN(bc->client_id);
}
else {
Comment thread
embhorn marked this conversation as resolved.
/* [MQTT-3.1.3-6] An empty-ClientId client must be told its
* assigned id; refuse rather than accept an unusable connection
* when the property pool cannot hold it. Clear the effective
* session expiry first so the refused connection tears its
* tentative session down instead of orphaning an unreachable
* one that could evict a valid session at capacity. */
bc->session_expiry_sec = 0;
ack.return_code = MQTT_REASON_SERVER_UNAVAILABLE;
goto send_connack;
}
}

/* Advertise feature availability */
Expand Down Expand Up @@ -7576,9 +7600,6 @@ static int BrokerHandle_Publish(BrokerClient* bc, int rx_len,
byte* payload = NULL;
char* topic = NULL;
MqttQoS eff_qos;
#if defined(WOLFMQTT_V5) && defined(WOLFMQTT_BROKER_RETAINED)
int retain_rc = MQTT_CODE_SUCCESS;
#endif
#if WOLFMQTT_MAX_QOS >= 2
int qos2_duplicate = 0;
#endif
Expand Down Expand Up @@ -7777,19 +7798,20 @@ static int BrokerHandle_Publish(BrokerClient* bc, int rx_len,
int ret_rc = BrokerRetained_Store(broker, topic, payload,
pub.total_len, pub.qos, expiry);
if (ret_rc != MQTT_CODE_SUCCESS) {
/* Retaining is best-effort: a store failure (table full,
* oversized payload, transient alloc) does not reject the
* PUBLISH. The message is still delivered live and
* acknowledged; only the retained copy is not kept. */
WBLOG_ERR(broker, "Retained store failed: %s",
MqttClient_ReturnCodeToString(ret_rc));
}
#ifdef WOLFMQTT_V5
retain_rc = ret_rc;
#endif
}
}
}
#endif /* WOLFMQTT_BROKER_RETAINED */

/* Fan-out is skipped for QoS 2 duplicates: subscribers already received
* the application message from the original PUBLISH ([MQTT-4.3.3]). */
/* Fan-out is skipped for QoS 2 duplicates: subscribers already received the
* application message from the original PUBLISH ([MQTT-4.3.3]). */
if (
#if WOLFMQTT_MAX_QOS >= 2
!qos2_duplicate &&
Expand Down Expand Up @@ -8105,13 +8127,6 @@ static int BrokerHandle_Publish(BrokerClient* bc, int rx_len,
#ifdef WOLFMQTT_V5
resp.protocol_level = bc->protocol_level;
resp.reason_code = MQTT_REASON_SUCCESS;
/* A retained-store failure must not be ACKed as success: tell the
* publisher the quota was exceeded [MQTT-3.4.2]. */
#ifdef WOLFMQTT_BROKER_RETAINED
if (retain_rc != MQTT_CODE_SUCCESS) {
resp.reason_code = MQTT_REASON_QUOTA_EXCEEDED;
}
#endif
resp.props = NULL;
#endif
rc = MqttEncode_PublishResp(bc->tx_buf, BROKER_CLIENT_TX_SZ(bc),
Expand Down
44 changes: 43 additions & 1 deletion src/mqtt_client.c
Original file line number Diff line number Diff line change
Expand Up @@ -1404,7 +1404,24 @@ static int MqttClient_DecodePacket(MqttClient* client, byte* rx_buf,
if (rc >= 0) {
packet_id = p_publish->packet_id;
#ifdef WOLFMQTT_V5
if (doProps) {
if (client->protocol_level >= MQTT_CONNECT_PROTOCOL_LEVEL_5) {
Comment thread
embhorn marked this conversation as resolved.
MqttProp* prop;
for (prop = p_publish->props; prop != NULL;
prop = prop->next) {
/* The client advertises Topic Alias Maximum 0 and keeps
* no inbound alias table, so it can neither resolve nor
* record an alias; reject rather than deliver an
* unresolved topic. */
if (prop->type == MQTT_PROP_TOPIC_ALIAS) {
MqttProps_Free(p_publish->props);
p_publish->props = NULL;
rc = MQTT_TRACE_ERROR(
MQTT_CODE_ERROR_MALFORMED_DATA);
break;
}
}
}
if (rc >= 0 && doProps) {
/* Retain returned properties until the message callback. */
int tmp = Handle_Props(client, p_publish->props,
(packet_obj != NULL),
Expand Down Expand Up @@ -3099,12 +3116,37 @@ int MqttClient_Connect(MqttClient *client, MqttConnect *mc_connect)
MqttProp* app_props = NULL;
int recv_max_added = 0;
#endif
#ifdef WOLFMQTT_V5
MqttProp* ta_prop;
int ta_count;
#endif

/* Validate required arguments */
if (client == NULL || mc_connect == NULL) {
return MQTT_TRACE_ERROR(MQTT_CODE_ERROR_BAD_ARG);
}

#ifdef WOLFMQTT_V5
/* The client does not resolve inbound Topic Aliases, so it must not
* advertise the capability. Reject a nonzero caller-supplied Topic Alias
* Maximum before sending CONNECT rather than advertising support it cannot
* honor and then rejecting the server's aliased PUBLISH. The scan is
* bounded like MqttEncode_Props so a cyclic list cannot spin. */
if (client->protocol_level >= MQTT_CONNECT_PROTOCOL_LEVEL_5) {
ta_count = 0;
for (ta_prop = mc_connect->props; ta_prop != NULL;
ta_prop = ta_prop->next) {
if (++ta_count > MQTT_MAX_PROPS) {
break;
}
if (ta_prop->type == MQTT_PROP_TOPIC_ALIAS_MAX &&
ta_prop->data_short != 0) {
return MQTT_TRACE_ERROR(MQTT_CODE_ERROR_PROPERTY);
}
}
}
#endif

#ifndef WOLFMQTT_NO_SESSION_REPLAY
if (mc_connect->stat.write == MQTT_MSG_PAYLOAD) {
/* MQTT_MSG_PAYLOAD is not part of the CONNECT write sequence
Expand Down
20 changes: 11 additions & 9 deletions src/mqtt_packet.c
Original file line number Diff line number Diff line change
Expand Up @@ -127,10 +127,6 @@ static const struct MqttPropMatrix gPropMatrix[] = {
{ MQTT_PROP_TYPE_MAX, MQTT_DATA_TYPE_NONE, 0 }
};

/* Maximum number of active properties - overridable */
#ifndef MQTT_MAX_PROPS
#define MQTT_MAX_PROPS 30
#endif

/* WOLFMQTT_DYN_PROP allows property allocation using malloc */
#ifndef WOLFMQTT_DYN_PROP
Expand Down Expand Up @@ -916,6 +912,12 @@ int MqttEncode_Props(MqttPacketType packet, MqttProp* props, byte* buf)
{
case MQTT_DATA_TYPE_BYTE:
Comment thread
embhorn marked this conversation as resolved.
{
/* Every MQTT 5 Byte property is Boolean-valued (0 or 1);
* Maximum QoS shares the same {0,1} domain. Reject any other
* value so the encoder never emits a Protocol Error. */
if (cur_prop->data_byte > 1) {
return MQTT_TRACE_ERROR(MQTT_CODE_ERROR_PROPERTY);
}
if (buf != NULL) {
*(buf++) = cur_prop->data_byte;
}
Expand Down Expand Up @@ -1165,11 +1167,11 @@ int MqttDecode_Props(MqttPacketType packet, MqttProp** props, byte* pbuf,
tmp++;
total++;
prop_len--;
/* [MQTT-3.1.2-28/29] Request Response/Problem Information
* MUST be 0 or 1; any other value is a Protocol Error. */
if ((cur_prop->type == MQTT_PROP_REQ_RESP_INFO ||
cur_prop->type == MQTT_PROP_REQ_PROB_INFO) &&
cur_prop->data_byte > 1) {
/* Every MQTT 5 Byte property is Boolean-valued (0 or 1), and
* Maximum QoS shares the same {0,1} domain; any other value is
* a Protocol Error. Mirrors MqttEncode_Props so a decoded
* property always re-encodes. */
if (cur_prop->data_byte > 1) {
rc = MQTT_TRACE_ERROR(MQTT_CODE_ERROR_PROPERTY);
}
break;
Expand Down
Loading
Loading