Skip to content
Closed
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
72 changes: 47 additions & 25 deletions source/core_mqtt_agent.c
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
/*
* coreMQTT Agent <DEVELOPMENT BRANCH>
* 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
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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 ) );

Expand Down Expand Up @@ -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;
Expand All @@ -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
Expand Down Expand Up @@ -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;
}

/*-----------------------------------------------------------*/
Expand Down Expand Up @@ -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 )
{
Expand Down Expand Up @@ -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 );
Expand Down Expand Up @@ -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 )
Expand Down
16 changes: 10 additions & 6 deletions source/core_mqtt_agent_command_functions.c
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
/*
* coreMQTT Agent <DEVELOPMENT BRANCH>
* 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
Expand Down Expand Up @@ -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 );
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand All @@ -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 )
Expand Down Expand Up @@ -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;
Expand Down
Loading