From c252ba6ed2ca5942fb52fe84035ee2376e762a67 Mon Sep 17 00:00:00 2001 From: Abtin Keshavarzian Date: Tue, 10 Aug 2021 18:22:46 -0700 Subject: [PATCH] [message] add `DequeueAndFree()` & `DequeueAndFreeAll()` to queues (#6902) This commit adds two new helper methods in `MessageQueue` and `PriorityQueue` which provide functionality that is commonly used in other modules: `DequeueAndFree()` which dequeues and frees a message from the queue), and `DequeueAndFreeAll()` which removes and frees all messages in the queue). --- src/core/coap/coap.cpp | 10 ++----- src/core/coap/coap_message.hpp | 8 +++++ src/core/coap/coap_secure.cpp | 9 +----- src/core/common/message.cpp | 32 ++++++++++++++++++++ src/core/common/message.hpp | 28 ++++++++++++++++++ src/core/meshcop/joiner_router.cpp | 3 +- src/core/net/dns_client.cpp | 3 +- src/core/net/ip6.cpp | 15 +++------- src/core/net/sntp_client.cpp | 5 +--- src/core/thread/mesh_forwarder.cpp | 41 +++++++------------------- src/core/thread/mesh_forwarder_ftd.cpp | 3 +- src/core/thread/mle.cpp | 3 +- 12 files changed, 90 insertions(+), 70 deletions(-) diff --git a/src/core/coap/coap.cpp b/src/core/coap/coap.cpp index e19261b59..111383ef5 100644 --- a/src/core/coap/coap.cpp +++ b/src/core/coap/coap.cpp @@ -1528,18 +1528,12 @@ void ResponsesQueue::UpdateQueue(void) void ResponsesQueue::DequeueResponse(Message &aMessage) { - mQueue.Dequeue(aMessage); - aMessage.Free(); + mQueue.DequeueAndFree(aMessage); } void ResponsesQueue::DequeueAllResponses(void) { - Message *message; - - while ((message = mQueue.GetHead()) != nullptr) - { - DequeueResponse(*message); - } + mQueue.DequeueAndFreeAll(); } void ResponsesQueue::HandleTimer(Timer &aTimer) diff --git a/src/core/coap/coap_message.hpp b/src/core/coap/coap_message.hpp index 760971111..9924f4478 100644 --- a/src/core/coap/coap_message.hpp +++ b/src/core/coap/coap_message.hpp @@ -1024,6 +1024,14 @@ public: * */ void Dequeue(Message &aMessage) { ot::MessageQueue::Dequeue(aMessage); } + + /** + * This method removes a message from the queue and frees it. + * + * @param[in] aMessage The message to remove and free. + * + */ + void DequeueAndFree(Message &aMessage) { ot::MessageQueue::DequeueAndFree(aMessage); } }; /** diff --git a/src/core/coap/coap_secure.cpp b/src/core/coap/coap_secure.cpp index c0a23642a..c9f45f9d2 100644 --- a/src/core/coap/coap_secure.cpp +++ b/src/core/coap/coap_secure.cpp @@ -84,16 +84,9 @@ exit: void CoapSecure::Stop(void) { - ot::Message *message; - mDtls.Close(); - while ((message = mTransmitQueue.GetHead()) != nullptr) - { - mTransmitQueue.Dequeue(*message); - message->Free(); - } - + mTransmitQueue.DequeueAndFreeAll(); ClearRequestsAndResponses(); } diff --git a/src/core/common/message.cpp b/src/core/common/message.cpp index 60eab7ac8..40ded9a26 100644 --- a/src/core/common/message.cpp +++ b/src/core/common/message.cpp @@ -801,6 +801,22 @@ void MessageQueue::Dequeue(Message &aMessage) aMessage.SetMessageQueue(nullptr); } +void MessageQueue::DequeueAndFree(Message &aMessage) +{ + Dequeue(aMessage); + aMessage.Free(); +} + +void MessageQueue::DequeueAndFreeAll(void) +{ + Message *message; + + while ((message = GetHead()) != nullptr) + { + DequeueAndFree(*message); + } +} + void MessageQueue::GetInfo(uint16_t &aMessageCount, uint16_t &aBufferCount) const { aMessageCount = 0; @@ -940,6 +956,22 @@ void PriorityQueue::Dequeue(Message &aMessage) aMessage.SetMessageQueue(nullptr); } +void PriorityQueue::DequeueAndFree(Message &aMessage) +{ + Dequeue(aMessage); + aMessage.Free(); +} + +void PriorityQueue::DequeueAndFreeAll(void) +{ + Message *message; + + while ((message = GetHead()) != nullptr) + { + DequeueAndFree(*message); + } +} + void PriorityQueue::GetInfo(uint16_t &aMessageCount, uint16_t &aBufferCount) const { aMessageCount = 0; diff --git a/src/core/common/message.hpp b/src/core/common/message.hpp index c33d03a2d..1f1d9ab4d 100644 --- a/src/core/common/message.hpp +++ b/src/core/common/message.hpp @@ -1436,6 +1436,20 @@ public: */ void Dequeue(Message &aMessage); + /** + * This method removes a message from the queue and frees it. + * + * @param[in] aMessage The message to remove and free. + * + */ + void DequeueAndFree(Message &aMessage); + + /** + * This method removes and frees all messages from the queue. + * + */ + void DequeueAndFreeAll(void); + /** * This method returns the number of messages and buffers enqueued. * @@ -1515,6 +1529,20 @@ public: */ void Dequeue(Message &aMessage); + /** + * This method removes a message from the queue and frees it. + * + * @param[in] aMessage The message to remove and free. + * + */ + void DequeueAndFree(Message &aMessage); + + /** + * This method removes and frees all messages from the queue. + * + */ + void DequeueAndFreeAll(void); + /** * This method returns the number of messages and buffers enqueued. * diff --git a/src/core/meshcop/joiner_router.cpp b/src/core/meshcop/joiner_router.cpp index 13dff8ca2..5d683cdbf 100644 --- a/src/core/meshcop/joiner_router.cpp +++ b/src/core/meshcop/joiner_router.cpp @@ -273,8 +273,7 @@ void JoinerRouter::SendDelayedJoinerEntrust(void) } else { - mDelayedJoinEnts.Dequeue(*message); - message->Free(); + mDelayedJoinEnts.DequeueAndFree(*message); Get().SetKek(metadata.mKek); diff --git a/src/core/net/dns_client.cpp b/src/core/net/dns_client.cpp index bcb8c40ed..383a8d873 100644 --- a/src/core/net/dns_client.cpp +++ b/src/core/net/dns_client.cpp @@ -745,8 +745,7 @@ exit: void Client::FreeQuery(Query &aQuery) { - mQueries.Dequeue(aQuery); - aQuery.Free(); + mQueries.DequeueAndFree(aQuery); } void Client::SendQuery(Query &aQuery, QueryInfo &aInfo, bool aUpdateTimer) diff --git a/src/core/net/ip6.cpp b/src/core/net/ip6.cpp index 628392afd..b728926cf 100644 --- a/src/core/net/ip6.cpp +++ b/src/core/net/ip6.cpp @@ -781,9 +781,9 @@ exit: { if (message != nullptr) { - mReassemblyList.Dequeue(*message); - message->Free(); + mReassemblyList.DequeueAndFree(*message); } + otLogWarnIp6("Reassembly failed: %s", ErrorToString(error)); } @@ -798,13 +798,7 @@ exit: void Ip6::CleanupFragmentationBuffer(void) { - Message *message; - - while ((message = mReassemblyList.GetHead()) != nullptr) - { - mReassemblyList.Dequeue(*message); - message->Free(); - } + mReassemblyList.DequeueAndFreeAll(); } void Ip6::HandleTimeTick(void) @@ -834,8 +828,7 @@ void Ip6::UpdateReassemblyList(void) otLogNoteIp6("Reassembly timeout."); SendIcmpError(*message, Icmp::Header::kTypeTimeExceeded, Icmp::Header::kCodeFragmReasTimeEx); - mReassemblyList.Dequeue(*message); - message->Free(); + mReassemblyList.DequeueAndFree(*message); } } } diff --git a/src/core/net/sntp_client.cpp b/src/core/net/sntp_client.cpp index 1269cbe82..83b83e98a 100644 --- a/src/core/net/sntp_client.cpp +++ b/src/core/net/sntp_client.cpp @@ -206,16 +206,13 @@ exit: void Client::DequeueMessage(Message &aMessage) { - mPendingQueries.Dequeue(aMessage); - if (mRetransmissionTimer.IsRunning() && (mPendingQueries.GetHead() == nullptr)) { // No more requests pending, stop the timer. mRetransmissionTimer.Stop(); } - // Free the message memory. - aMessage.Free(); + mPendingQueries.DequeueAndFree(aMessage); } Error Client::SendMessage(Message &aMessage, const Ip6::MessageInfo &aMessageInfo) diff --git a/src/core/thread/mesh_forwarder.cpp b/src/core/thread/mesh_forwarder.cpp index 838ec8bcf..dcca5f25c 100644 --- a/src/core/thread/mesh_forwarder.cpp +++ b/src/core/thread/mesh_forwarder.cpp @@ -119,25 +119,14 @@ void MeshForwarder::Start(void) void MeshForwarder::Stop(void) { - Message *message; - VerifyOrExit(mEnabled); mDataPollSender.StopPolling(); Get().UnregisterReceiver(TimeTicker::kMeshForwarder); Get().Stop(); - while ((message = mSendQueue.GetHead()) != nullptr) - { - mSendQueue.Dequeue(*message); - message->Free(); - } - - while ((message = mReassemblyList.GetHead()) != nullptr) - { - mReassemblyList.Dequeue(*message); - message->Free(); - } + mSendQueue.DequeueAndFreeAll(); + mReassemblyList.DequeueAndFreeAll(); #if OPENTHREAD_FTD mIndirectSender.Stop(); @@ -223,9 +212,8 @@ void MeshForwarder::RemoveMessage(Message &aMessage) } } - queue->Dequeue(aMessage); LogMessage(kMessageEvict, aMessage, nullptr, kErrorNoBufs); - aMessage.Free(); + queue->DequeueAndFree(aMessage); } void MeshForwarder::ResumeMessageTransmissions(void) @@ -320,9 +308,8 @@ Message *MeshForwarder::GetDirectTransmission(void) #endif default: - mSendQueue.Dequeue(*curMessage); LogMessage(kMessageDrop, *curMessage, nullptr, error); - curMessage->Free(); + mSendQueue.DequeueAndFree(*curMessage); continue; } } @@ -1112,8 +1099,7 @@ void MeshForwarder::RemoveMessageIfNoPendingTx(Message &aMessage) mMessageNextOffset = 0; } - mSendQueue.Dequeue(aMessage); - aMessage.Free(); + mSendQueue.DequeueAndFree(aMessage); exit: return; @@ -1326,23 +1312,17 @@ exit: void MeshForwarder::ClearReassemblyList(void) { - Message *message; - Message *next; - - for (message = mReassemblyList.GetHead(); message; message = next) + for (const Message *message = mReassemblyList.GetHead(); message != nullptr; message = message->GetNext()) { - next = message->GetNext(); - mReassemblyList.Dequeue(*message); - LogMessage(kMessageReassemblyDrop, *message, nullptr, kErrorNoFrameReceived); if (message->GetType() == Message::kTypeIp6) { mIpCounters.mRxFailure++; } - - message->Free(); } + + mReassemblyList.DequeueAndFreeAll(); } void MeshForwarder::HandleTimeTick(void) @@ -1375,15 +1355,14 @@ bool MeshForwarder::UpdateReassemblyList(void) } else { - mReassemblyList.Dequeue(*message); - LogMessage(kMessageReassemblyDrop, *message, nullptr, kErrorReassemblyTimeout); + if (message->GetType() == Message::kTypeIp6) { mIpCounters.mRxFailure++; } - message->Free(); + mReassemblyList.DequeueAndFree(*message); } } diff --git a/src/core/thread/mesh_forwarder_ftd.cpp b/src/core/thread/mesh_forwarder_ftd.cpp index 2c817e05a..f6d36876e 100644 --- a/src/core/thread/mesh_forwarder_ftd.cpp +++ b/src/core/thread/mesh_forwarder_ftd.cpp @@ -329,9 +329,8 @@ void MeshForwarder::RemoveDataResponseMessages(void) mSendMessage = nullptr; } - mSendQueue.Dequeue(*message); LogMessage(kMessageDrop, *message, nullptr, kErrorNone); - message->Free(); + mSendQueue.DequeueAndFree(*message); } } diff --git a/src/core/thread/mle.cpp b/src/core/thread/mle.cpp index 0e5d59103..8a98c8acf 100644 --- a/src/core/thread/mle.cpp +++ b/src/core/thread/mle.cpp @@ -1989,8 +1989,7 @@ void Mle::RemoveDelayedDataResponseMessage(void) if (message->GetSubType() == Message::kSubTypeMleDataResponse) { - mDelayedResponses.Dequeue(*message); - message->Free(); + mDelayedResponses.DequeueAndFree(*message); Log(kMessageRemoveDelayed, kTypeDataResponse, metadata.mDestination); // no more than one multicast MLE Data Response in Delayed Message Queue.