mirror of
https://github.com/espressif/openthread.git
synced 2026-10-07 00:07:42 +00:00
[coap] ensure correct dequeue time for entries in response queue (#4686)
This commit changes the handling of timer for CoAP response queue to ensure that correct dequeue time is used. The existing code assumed that entries in the queue would be in order of their dequeue time (due to exchange lifetime being fixed). With the addition of configurable CoAP tx parameters feature, this is no longer valid. This commit changes the code to restart the timer with the earliest dequeue time. Also when evicting a message from the queue (due to the number of entries reaching the max limit of `kMaxCachedResponses`) we now search and find the message with earliest dequeue time to remove (instead of removing the entry at the head of queue).
This commit is contained in:
committed by
Jonathan Hui
parent
1d463868d9
commit
bf1cf97b0f
+53
-46
@@ -834,44 +834,58 @@ void ResponsesQueue::EnqueueResponse(Message & aMessage,
|
|||||||
const Ip6::MessageInfo &aMessageInfo,
|
const Ip6::MessageInfo &aMessageInfo,
|
||||||
const TxParameters & aTxParameters)
|
const TxParameters & aTxParameters)
|
||||||
{
|
{
|
||||||
otError error = OT_ERROR_NONE;
|
Message * responseCopy;
|
||||||
Message * responseCopy = NULL;
|
|
||||||
uint16_t messageCount;
|
|
||||||
uint16_t bufferCount;
|
|
||||||
uint32_t exchangeLifetime = aTxParameters.CalculateExchangeLifetime();
|
|
||||||
ResponseMetadata metadata;
|
ResponseMetadata metadata;
|
||||||
|
|
||||||
metadata.mDequeueTime = TimerMilli::GetNow() + exchangeLifetime;
|
metadata.mDequeueTime = TimerMilli::GetNow() + aTxParameters.CalculateExchangeLifetime();
|
||||||
metadata.mMessageInfo = aMessageInfo;
|
metadata.mMessageInfo = aMessageInfo;
|
||||||
|
|
||||||
// Return success if matched response already exists in the cache.
|
|
||||||
VerifyOrExit(FindMatchedResponse(aMessage, aMessageInfo) == NULL);
|
VerifyOrExit(FindMatchedResponse(aMessage, aMessageInfo) == NULL);
|
||||||
|
|
||||||
mQueue.GetInfo(messageCount, bufferCount);
|
UpdateQueue();
|
||||||
|
|
||||||
if (messageCount >= kMaxCachedResponses)
|
|
||||||
{
|
|
||||||
DequeueOldestResponse();
|
|
||||||
}
|
|
||||||
|
|
||||||
VerifyOrExit((responseCopy = aMessage.Clone()) != NULL);
|
VerifyOrExit((responseCopy = aMessage.Clone()) != NULL);
|
||||||
|
|
||||||
SuccessOrExit(error = metadata.AppendTo(*responseCopy));
|
VerifyOrExit(metadata.AppendTo(*responseCopy) == OT_ERROR_NONE, responseCopy->Free());
|
||||||
|
|
||||||
mQueue.Enqueue(*responseCopy);
|
mQueue.Enqueue(*responseCopy);
|
||||||
|
|
||||||
if (!mTimer.IsRunning())
|
mTimer.FireAtIfEarlier(metadata.mDequeueTime);
|
||||||
{
|
|
||||||
mTimer.Start(exchangeLifetime);
|
|
||||||
}
|
|
||||||
|
|
||||||
exit:
|
exit:
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
if (error != OT_ERROR_NONE && responseCopy != NULL)
|
void ResponsesQueue::UpdateQueue(void)
|
||||||
|
{
|
||||||
|
uint16_t msgCount = 0;
|
||||||
|
Message * earliestMsg = NULL;
|
||||||
|
TimeMilli earliestDequeueTime(0);
|
||||||
|
|
||||||
|
// Check the number of messages in the queue and if number is at
|
||||||
|
// `kMaxCachedResponses` remove the one with earliest dequeue
|
||||||
|
// time.
|
||||||
|
|
||||||
|
for (Message *message = static_cast<Message *>(mQueue.GetHead()); message != NULL;
|
||||||
|
message = static_cast<Message *>(message->GetNext()))
|
||||||
{
|
{
|
||||||
responseCopy->Free();
|
ResponseMetadata metadata;
|
||||||
|
|
||||||
|
metadata.ReadFrom(*message);
|
||||||
|
|
||||||
|
if ((earliestMsg == NULL) || (metadata.mDequeueTime < earliestDequeueTime))
|
||||||
|
{
|
||||||
|
earliestMsg = message;
|
||||||
|
earliestDequeueTime = metadata.mDequeueTime;
|
||||||
|
}
|
||||||
|
|
||||||
|
msgCount++;
|
||||||
}
|
}
|
||||||
|
|
||||||
return;
|
if (msgCount >= kMaxCachedResponses)
|
||||||
|
{
|
||||||
|
DequeueResponse(*earliestMsg);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
void ResponsesQueue::DequeueResponse(Message &aMessage)
|
void ResponsesQueue::DequeueResponse(Message &aMessage)
|
||||||
@@ -880,17 +894,6 @@ void ResponsesQueue::DequeueResponse(Message &aMessage)
|
|||||||
aMessage.Free();
|
aMessage.Free();
|
||||||
}
|
}
|
||||||
|
|
||||||
void ResponsesQueue::DequeueOldestResponse(void)
|
|
||||||
{
|
|
||||||
Message *message;
|
|
||||||
|
|
||||||
VerifyOrExit((message = static_cast<Message *>(mQueue.GetHead())) != NULL);
|
|
||||||
DequeueResponse(*message);
|
|
||||||
|
|
||||||
exit:
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
void ResponsesQueue::DequeueAllResponses(void)
|
void ResponsesQueue::DequeueAllResponses(void)
|
||||||
{
|
{
|
||||||
Message *message;
|
Message *message;
|
||||||
@@ -908,23 +911,34 @@ void ResponsesQueue::HandleTimer(Timer &aTimer)
|
|||||||
|
|
||||||
void ResponsesQueue::HandleTimer(void)
|
void ResponsesQueue::HandleTimer(void)
|
||||||
{
|
{
|
||||||
Message * message;
|
TimeMilli now = TimerMilli::GetNow();
|
||||||
ResponseMetadata metadata;
|
TimeMilli nextDequeueTime = now.GetDistantFuture();
|
||||||
|
Message * nextMessage;
|
||||||
|
|
||||||
while ((message = static_cast<Message *>(mQueue.GetHead())) != NULL)
|
for (Message *message = static_cast<Message *>(mQueue.GetHead()); message != NULL; message = nextMessage)
|
||||||
{
|
{
|
||||||
|
ResponseMetadata metadata;
|
||||||
|
|
||||||
|
nextMessage = static_cast<Message *>(message->GetNext());
|
||||||
|
|
||||||
metadata.ReadFrom(*message);
|
metadata.ReadFrom(*message);
|
||||||
|
|
||||||
if (TimerMilli::GetNow() >= metadata.mDequeueTime)
|
if (now >= metadata.mDequeueTime)
|
||||||
{
|
{
|
||||||
DequeueResponse(*message);
|
DequeueResponse(*message);
|
||||||
|
continue;
|
||||||
}
|
}
|
||||||
else
|
|
||||||
|
if (metadata.mDequeueTime < nextDequeueTime)
|
||||||
{
|
{
|
||||||
mTimer.Start(metadata.GetRemainingTime());
|
nextDequeueTime = metadata.mDequeueTime;
|
||||||
break;
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if (nextDequeueTime < now.GetDistantFuture())
|
||||||
|
{
|
||||||
|
mTimer.FireAt(nextDequeueTime);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
void ResponsesQueue::ResponseMetadata::ReadFrom(const Message &aMessage)
|
void ResponsesQueue::ResponseMetadata::ReadFrom(const Message &aMessage)
|
||||||
@@ -935,13 +949,6 @@ void ResponsesQueue::ResponseMetadata::ReadFrom(const Message &aMessage)
|
|||||||
aMessage.Read(length - sizeof(*this), sizeof(*this), this);
|
aMessage.Read(length - sizeof(*this), sizeof(*this), this);
|
||||||
}
|
}
|
||||||
|
|
||||||
uint32_t ResponsesQueue::ResponseMetadata::GetRemainingTime(void) const
|
|
||||||
{
|
|
||||||
TimeMilli now = TimerMilli::GetNow();
|
|
||||||
|
|
||||||
return (mDequeueTime > now) ? mDequeueTime - now : 0;
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Return product of @p aValueA and @p aValueB if no overflow otherwise 0.
|
/// Return product of @p aValueA and @p aValueB if no overflow otherwise 0.
|
||||||
static uint32_t Multiply(uint32_t aValueA, uint32_t aValueB)
|
static uint32_t Multiply(uint32_t aValueA, uint32_t aValueB)
|
||||||
{
|
{
|
||||||
|
|||||||
@@ -210,12 +210,6 @@ public:
|
|||||||
*/
|
*/
|
||||||
void EnqueueResponse(Message &aMessage, const Ip6::MessageInfo &aMessageInfo, const TxParameters &aTxParameters);
|
void EnqueueResponse(Message &aMessage, const Ip6::MessageInfo &aMessageInfo, const TxParameters &aTxParameters);
|
||||||
|
|
||||||
/**
|
|
||||||
* This method removes the oldest response from the cache.
|
|
||||||
*
|
|
||||||
*/
|
|
||||||
void DequeueOldestResponse(void);
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* This method removes all responses from the cache.
|
* This method removes all responses from the cache.
|
||||||
*
|
*
|
||||||
@@ -252,9 +246,8 @@ private:
|
|||||||
|
|
||||||
struct ResponseMetadata
|
struct ResponseMetadata
|
||||||
{
|
{
|
||||||
otError AppendTo(Message &aMessage) const { return aMessage.Append(this, sizeof(*this)); }
|
otError AppendTo(Message &aMessage) const { return aMessage.Append(this, sizeof(*this)); }
|
||||||
void ReadFrom(const Message &aMessage);
|
void ReadFrom(const Message &aMessage);
|
||||||
uint32_t GetRemainingTime(void) const;
|
|
||||||
|
|
||||||
TimeMilli mDequeueTime;
|
TimeMilli mDequeueTime;
|
||||||
Ip6::MessageInfo mMessageInfo;
|
Ip6::MessageInfo mMessageInfo;
|
||||||
@@ -262,6 +255,7 @@ private:
|
|||||||
|
|
||||||
const Message *FindMatchedResponse(const Message &aRequest, const Ip6::MessageInfo &aMessageInfo) const;
|
const Message *FindMatchedResponse(const Message &aRequest, const Ip6::MessageInfo &aMessageInfo) const;
|
||||||
void DequeueResponse(Message &aMessage);
|
void DequeueResponse(Message &aMessage);
|
||||||
|
void UpdateQueue(void);
|
||||||
|
|
||||||
static void HandleTimer(Timer &aTimer);
|
static void HandleTimer(Timer &aTimer);
|
||||||
void HandleTimer(void);
|
void HandleTimer(void);
|
||||||
|
|||||||
Reference in New Issue
Block a user