From 876cf9158bb5468ba4b66e9a268e650e1f0b267f Mon Sep 17 00:00:00 2001 From: Kody Stribrny Date: Fri, 6 Mar 2026 17:02:58 -0800 Subject: [PATCH] Rough attempt at coreMQTT v5 updates This supports the latest coreMQTT v5 version which is now MQTT v5 conformant. --- source/core_mqtt_agent.c | 72 ++++++++++++++-------- source/core_mqtt_agent_command_functions.c | 16 +++-- 2 files changed, 57 insertions(+), 31 deletions(-) diff --git a/source/core_mqtt_agent.c b/source/core_mqtt_agent.c index 5c0f840b..3f22c4f5 100644 --- a/source/core_mqtt_agent.c +++ b/source/core_mqtt_agent.c @@ -1,5 +1,5 @@ /* - * coreMQTT Agent + * coreMQTT Agent v1.2.0 * Copyright (C) 2021 Amazon.com, Inc. or its affiliates. All Rights Reserved. * * Permission is hereby granted, free of charge, to any person obtaining a copy of @@ -51,6 +51,21 @@ /*-----------------------------------------------------------*/ +#if ( MQTT_AGENT_USE_QOS_1_2_PUBLISH != 0 ) + +/** + * @brief Array used to maintain the outgoing publish records and their + * state by the coreMQTT library. + */ + static MQTTPubAckInfo_t pOutgoingPublishRecords[ MQTT_AGENT_MAX_OUTSTANDING_ACKS ]; + +/** + * @brief Array used to maintain the incoming publish records and their + * state by the coreMQTT library. + */ + static MQTTPubAckInfo_t pIncomingPublishRecords[ MQTT_AGENT_MAX_OUTSTANDING_ACKS ]; +#endif + /** * @brief Track an operation by adding it to a list, indicating it is anticipating * an acknowledgment. @@ -145,9 +160,12 @@ static MQTTStatus_t processCommand( MQTTAgentContext_t * pMqttAgentContext, * @param[in] pDeserializedInfo Pointer to deserialized information from * the incoming packet. */ -static void mqttEventCallback( MQTTContext_t * pMqttContext, +static bool mqttEventCallback( MQTTContext_t * pMqttContext, MQTTPacketInfo_t * pPacketInfo, - MQTTDeserializedInfo_t * pDeserializedInfo ); + MQTTDeserializedInfo_t * pDeserializedInfo, + MQTTSuccessFailReasonCode_t * pReasonCode, + MQTTPropBuilder_t * pSendPropsBuffer, + MQTTPropBuilder_t * pGetPropsBuffer ); /** * @brief Mark a command as complete after receiving an acknowledgment packet. @@ -547,9 +565,9 @@ static MQTTStatus_t processCommand( MQTTAgentContext_t * pMqttAgentContext, if( pCommand != NULL ) { - assert( ( uint32_t ) pCommand->commandType < ( uint32_t ) NUM_COMMANDS ); + assert( pCommand->commandType < NUM_COMMANDS ); - if( ( uint32_t ) pCommand->commandType < ( uint32_t ) NUM_COMMANDS ) + if( ( pCommand->commandType >= NONE ) && ( pCommand->commandType < NUM_COMMANDS ) ) { commandFunction = pCommandFunctionTable[ pCommand->commandType ]; pCommandArgs = pCommand->pArgs; @@ -598,12 +616,6 @@ static MQTTStatus_t processCommand( MQTTAgentContext_t * pMqttAgentContext, } while( pMqttAgentContext->packetReceivedInLoop ); } - if( operationStatus == MQTTNeedMoreBytes ) - { - /* Reset the operation status as MQTTNeedMoreBytes is not an error condition. */ - operationStatus = MQTTSuccess; - } - /* Set the flag to break from the command loop. */ *pEndLoop = ( commandOutParams.endLoop || ( operationStatus != MQTTSuccess ) ); @@ -642,17 +654,17 @@ static MQTTAgentContext_t * getAgentFromMQTTContext( MQTTContext_t * pMQTTContex MQTTAgentContext_t ctx = { 0 }; ptrdiff_t offset = ( ( uint8_t * ) &( ctx.mqttContext ) ) - ( ( uint8_t * ) &ctx ); - /* MISRA Ref 11.3.1 [Misaligned access] */ - /* More details at: https://github.com/FreeRTOS/coreMQTT-Agent/blob/main/MISRA.md#rule-113 */ - /* coverity[misra_c_2012_rule_11_3_violation] */ return ( MQTTAgentContext_t * ) &( ( ( uint8_t * ) pMQTTContext )[ 0 - offset ] ); } /*-----------------------------------------------------------*/ -static void mqttEventCallback( MQTTContext_t * pMqttContext, +static bool mqttEventCallback( MQTTContext_t * pMqttContext, MQTTPacketInfo_t * pPacketInfo, - MQTTDeserializedInfo_t * pDeserializedInfo ) + MQTTDeserializedInfo_t * pDeserializedInfo, + MQTTSuccessFailReasonCode_t * pReasonCode, + MQTTPropBuilder_t * pSendPropsBuffer, + MQTTPropBuilder_t * pGetPropsBuffer ) { MQTTAgentAckInfo_t * pAckInfo; uint16_t packetIdentifier = pDeserializedInfo->packetIdentifier; @@ -662,6 +674,10 @@ static void mqttEventCallback( MQTTContext_t * pMqttContext, assert( pMqttContext != NULL ); assert( pPacketInfo != NULL ); + ( void ) pReasonCode; + ( void ) pSendPropsBuffer; + ( void ) pGetPropsBuffer; + pAgentContext = getAgentFromMQTTContext( pMqttContext ); /* This callback executes from within MQTT_ProcessLoop(). Setting this flag @@ -710,12 +726,15 @@ static void mqttEventCallback( MQTTContext_t * pMqttContext, break; /* Any other packet type is invalid. */ + case MQTT_PACKET_TYPE_PINGRESP: default: LogError( ( "Unknown packet type received:(%02x).\n", pPacketInfo->type ) ); break; } } + + return true; } /*-----------------------------------------------------------*/ @@ -837,7 +856,7 @@ static MQTTStatus_t resendPublishes( MQTTAgentContext_t * pMqttAgentContext ) /* Set the DUP flag. */ pOriginalPublish = ( MQTTPublishInfo_t * ) ( pFoundAck->pOriginalCommand->pArgs ); pOriginalPublish->dup = true; - statusResult = MQTT_Publish( pMqttContext, pOriginalPublish, packetId ); + statusResult = MQTT_Publish( pMqttContext, pOriginalPublish, packetId, NULL ); if( statusResult != MQTTSuccess ) { @@ -952,6 +971,7 @@ static bool validateParams( MQTTAgentCommandType_t commandType, ( pSubscribeArgs->numSubscriptions != 0U ) ); break; + case PUBLISH: default: /* Publish, does not need to be cast since we do not check it. */ ret = ( pParams != NULL ); @@ -1001,16 +1021,18 @@ MQTTStatus_t MQTTAgent_Init( MQTTAgentContext_t * pMqttAgentContext, pNetworkBuffer ); #if ( MQTT_AGENT_USE_QOS_1_2_PUBLISH != 0 ) - { - if( returnStatus == MQTTSuccess ) { - returnStatus = MQTT_InitStatefulQoS( &( pMqttAgentContext->mqttContext ), - pMqttAgentContext->pOutgoingPublishRecords, - MQTT_AGENT_MAX_OUTSTANDING_ACKS, - pMqttAgentContext->pIncomingPublishRecords, - MQTT_AGENT_MAX_OUTSTANDING_ACKS ); + if( returnStatus == MQTTSuccess ) + { + returnStatus = MQTT_InitStatefulQoS( &( pMqttAgentContext->mqttContext ), + pOutgoingPublishRecords, + MQTT_AGENT_MAX_OUTSTANDING_ACKS, + pIncomingPublishRecords, + MQTT_AGENT_MAX_OUTSTANDING_ACKS, + NULL, + 0 ); + } } - } #endif /* if ( MQTT_AGENT_USE_QOS_1_2_PUBLISH != 0 ) */ if( returnStatus == MQTTSuccess ) diff --git a/source/core_mqtt_agent_command_functions.c b/source/core_mqtt_agent_command_functions.c index 9a4e4653..00b39ac7 100644 --- a/source/core_mqtt_agent_command_functions.c +++ b/source/core_mqtt_agent_command_functions.c @@ -1,5 +1,5 @@ /* - * coreMQTT Agent + * coreMQTT Agent v1.2.0 * Copyright (C) 2021 Amazon.com, Inc. or its affiliates. All Rights Reserved. * * Permission is hereby granted, free of charge, to any person obtaining a copy of @@ -77,7 +77,7 @@ MQTTStatus_t MQTTAgentCommand_Publish( MQTTAgentContext_t * pMqttAgentContext, } LogInfo( ( "Publishing message to %.*s.\n", ( int ) pPublishInfo->topicNameLength, pPublishInfo->pTopicName ) ); - ret = MQTT_Publish( &( pMqttAgentContext->mqttContext ), pPublishInfo, pReturnFlags->packetId ); + ret = MQTT_Publish( &( pMqttAgentContext->mqttContext ), pPublishInfo, pReturnFlags->packetId, NULL ); /* Add to pending ack list, or call callback if QoS 0. */ pReturnFlags->addAcknowledgment = ( pPublishInfo->qos != MQTTQoS0 ) && ( ret == MQTTSuccess ); @@ -106,7 +106,8 @@ MQTTStatus_t MQTTAgentCommand_Subscribe( MQTTAgentContext_t * pMqttAgentContext, ret = MQTT_Subscribe( &( pMqttAgentContext->mqttContext ), pSubscribeArgs->pSubscribeInfo, pSubscribeArgs->numSubscriptions, - pReturnFlags->packetId ); + pReturnFlags->packetId, + NULL ); pReturnFlags->addAcknowledgment = ( ret == MQTTSuccess ); pReturnFlags->runProcessLoop = true; @@ -134,7 +135,8 @@ MQTTStatus_t MQTTAgentCommand_Unsubscribe( MQTTAgentContext_t * pMqttAgentContex ret = MQTT_Unsubscribe( &( pMqttAgentContext->mqttContext ), pSubscribeArgs->pSubscribeInfo, pSubscribeArgs->numSubscriptions, - pReturnFlags->packetId ); + pReturnFlags->packetId, + NULL ); pReturnFlags->addAcknowledgment = ( ret == MQTTSuccess ); pReturnFlags->runProcessLoop = true; @@ -161,7 +163,9 @@ MQTTStatus_t MQTTAgentCommand_Connect( MQTTAgentContext_t * pMqttAgentContext, pConnectInfo->pConnectInfo, pConnectInfo->pWillInfo, pConnectInfo->timeoutMs, - &( pConnectInfo->sessionPresent ) ); + &( pConnectInfo->sessionPresent ), + NULL, + NULL ); /* Resume a session if one existed, else clear the list of acknowledgments. */ if( ret == MQTTSuccess ) @@ -189,7 +193,7 @@ MQTTStatus_t MQTTAgentCommand_Disconnect( MQTTAgentContext_t * pMqttAgentContext assert( pMqttAgentContext != NULL ); assert( pReturnFlags != NULL ); - ret = MQTT_Disconnect( &( pMqttAgentContext->mqttContext ) ); + ret = MQTT_Disconnect( &( pMqttAgentContext->mqttContext ), NULL, NULL ); ( void ) memset( pReturnFlags, 0x00, sizeof( MQTTAgentCommandFuncReturns_t ) ); pReturnFlags->endLoop = true;