From 4ed2c8c222a8d1ebcf1909b9f24435ef54860e94 Mon Sep 17 00:00:00 2001 From: Zhanglong Xia Date: Thu, 20 Dec 2018 03:28:34 +0800 Subject: [PATCH] [qos] forward messages based on the priority of the message (#3317) --- src/core/common/message.hpp | 6 + src/core/openthread-core-default-config.h | 15 +- src/core/thread/mesh_forwarder.cpp | 85 +++++++++--- src/core/thread/mesh_forwarder.hpp | 160 +++++++++++++++++++--- src/core/thread/mesh_forwarder_ftd.cpp | 154 ++++++++++++++++++++- 5 files changed, 378 insertions(+), 42 deletions(-) diff --git a/src/core/common/message.hpp b/src/core/common/message.hpp index 816c09a4e..1c2e1b501 100644 --- a/src/core/common/message.hpp +++ b/src/core/common/message.hpp @@ -576,6 +576,12 @@ public: */ void SetTimeout(uint8_t aTimeout) { mBuffer.mHead.mInfo.mTimeout = aTimeout; } + /** + * This method decrements the timeout. + * + */ + void DecrementTimeout(void) { mBuffer.mHead.mInfo.mTimeout--; } + /** * This method returns the interface ID. * diff --git a/src/core/openthread-core-default-config.h b/src/core/openthread-core-default-config.h index 87b983ff0..95de5f217 100644 --- a/src/core/openthread-core-default-config.h +++ b/src/core/openthread-core-default-config.h @@ -392,11 +392,11 @@ /** * @def OPENTHREAD_CONFIG_6LOWPAN_REASSEMBLY_TIMEOUT * - * The 6LoWPAN fragment reassembly timeout in seconds. + * The reassembly timeout between 6LoWPAN fragments in seconds. * */ #ifndef OPENTHREAD_CONFIG_6LOWPAN_REASSEMBLY_TIMEOUT -#define OPENTHREAD_CONFIG_6LOWPAN_REASSEMBLY_TIMEOUT 5 +#define OPENTHREAD_CONFIG_6LOWPAN_REASSEMBLY_TIMEOUT 2 #endif /** @@ -1919,4 +1919,15 @@ #ifndef OPENTHREAD_CONFIG_TIME_SYNC_JUMP_NOTIF_MIN_US #define OPENTHREAD_CONFIG_TIME_SYNC_JUMP_NOTIF_MIN_US 10000 #endif + +/** + * @def OPENTHREAD_CONFIG_NUM_FRAGMENT_PRIORITY_ENTRIES + * + * The number of fragment priority entries. + * + */ +#ifndef OPENTHREAD_CONFIG_NUM_FRAGMENT_PRIORITY_ENTRIES +#define OPENTHREAD_CONFIG_NUM_FRAGMENT_PRIORITY_ENTRIES 8 +#endif + #endif // OPENTHREAD_CORE_DEFAULT_CONFIG_H_ diff --git a/src/core/thread/mesh_forwarder.cpp b/src/core/thread/mesh_forwarder.cpp index 957a436ea..654897d8f 100644 --- a/src/core/thread/mesh_forwarder.cpp +++ b/src/core/thread/mesh_forwarder.cpp @@ -59,7 +59,7 @@ namespace ot { MeshForwarder::MeshForwarder(Instance &aInstance) : InstanceLocator(aInstance) , mDiscoverTimer(aInstance, &MeshForwarder::HandleDiscoverTimer, this) - , mReassemblyTimer(aInstance, &MeshForwarder::HandleReassemblyTimer, this) + , mUpdateTimer(aInstance, &MeshForwarder::HandleUpdateTimer, this) , mMessageNextOffset(0) , mSendMessage(NULL) , mSendMessageIsARetransmission(false) @@ -91,6 +91,10 @@ MeshForwarder::MeshForwarder(Instance &aInstance) mIpCounters.mRxSuccess = 0; mIpCounters.mTxFailure = 0; mIpCounters.mRxFailure = 0; + +#if OPENTHREAD_FTD + memset(mFragmentEntries, 0, sizeof(mFragmentEntries)); +#endif } otError MeshForwarder::Start(void) @@ -115,7 +119,7 @@ otError MeshForwarder::Stop(void) VerifyOrExit(mEnabled == true); mDataPollManager.StopPolling(); - mReassemblyTimer.Stop(); + mUpdateTimer.Stop(); if (mScanning) { @@ -134,6 +138,10 @@ otError MeshForwarder::Stop(void) message->Free(); } +#if OPENTHREAD_FTD + memset(mFragmentEntries, 0, sizeof(mFragmentEntries)); +#endif + mEnabled = false; mSendMessage = NULL; netif.GetMac().SetRxOnWhenIdle(false); @@ -447,6 +455,18 @@ otError MeshForwarder::GetMacDestinationAddress(const Ip6::Address &aIp6Addr, Ma return OT_ERROR_NONE; } +otError MeshForwarder::GetMeshHeader(const uint8_t *&aFrame, uint8_t &aFrameLength, Lowpan::MeshHeader &aMeshHeader) +{ + otError error; + + VerifyOrExit(aFrameLength >= 1 && reinterpret_cast(aFrame)->IsMeshHeader(), + error = OT_ERROR_NOT_FOUND); + SuccessOrExit(error = aMeshHeader.Init(aFrame, aFrameLength)); + +exit: + return error; +} + otError MeshForwarder::SkipMeshHeader(const uint8_t *&aFrame, uint8_t &aFrameLength) { otError error = OT_ERROR_NONE; @@ -462,6 +482,21 @@ exit: return error; } +otError MeshForwarder::GetFragmentHeader(const uint8_t * aFrame, + uint8_t aFrameLength, + Lowpan::FragmentHeader &aFragmentHeader) +{ + otError error = OT_ERROR_NONE; + + VerifyOrExit(aFrameLength >= 1 && reinterpret_cast(aFrame)->IsFragmentHeader(), + error = OT_ERROR_NOT_FOUND); + + SuccessOrExit(error = aFragmentHeader.Init(aFrame, aFrameLength)); + +exit: + return error; +} + otError MeshForwarder::DecompressIp6Header(const uint8_t * aFrame, uint8_t aFrameLength, const Mac::Address &aMacSource, @@ -478,13 +513,10 @@ otError MeshForwarder::DecompressIp6Header(const uint8_t * aFrame, SuccessOrExit(error = SkipMeshHeader(aFrame, aFrameLength)); - if (aFrameLength >= 1 && reinterpret_cast(aFrame)->IsFragmentHeader()) + if (GetFragmentHeader(aFrame, aFrameLength, fragmentHeader) == OT_ERROR_NONE) { - SuccessOrExit(error = fragmentHeader.Init(aFrame, aFrameLength)); - // only the first fragment header is followed by a LOWPAN_IPHC header VerifyOrExit(fragmentHeader.GetDatagramOffset() == 0, error = OT_ERROR_NOT_FOUND); - aFrame += fragmentHeader.GetHeaderLength(); aFrameLength -= fragmentHeader.GetHeaderLength(); } @@ -1351,9 +1383,9 @@ void MeshForwarder::HandleFragment(uint8_t * aFrame, mReassemblyList.Enqueue(*message); - if (!mReassemblyTimer.IsRunning()) + if (!mUpdateTimer.IsRunning()) { - mReassemblyTimer.Start(kStateUpdatePeriod); + mUpdateTimer.Start(kStateUpdatePeriod); } } else @@ -1391,6 +1423,7 @@ void MeshForwarder::HandleFragment(uint8_t * aFrame, message->Write(message->GetOffset(), aFrameLength, aFrame); message->MoveOffset(aFrameLength); message->AddRss(aLinkInfo.mRss); + message->SetTimeout(kReassemblyTimeout); } exit: @@ -1435,31 +1468,42 @@ void MeshForwarder::ClearReassemblyList(void) } } -void MeshForwarder::HandleReassemblyTimer(Timer &aTimer) +void MeshForwarder::HandleUpdateTimer(Timer &aTimer) { - aTimer.GetOwner().HandleReassemblyTimer(); + aTimer.GetOwner().HandleUpdateTimer(); } -void MeshForwarder::HandleReassemblyTimer(void) +void MeshForwarder::HandleUpdateTimer(void) +{ + bool shouldRun = false; + +#if OPENTHREAD_FTD + shouldRun = UpdateFragmentLifetime(); +#endif + + if (UpdateReassemblyList() || shouldRun) + { + mUpdateTimer.Start(kStateUpdatePeriod); + } +} + +bool MeshForwarder::UpdateReassemblyList(void) { Message *next = NULL; - uint8_t timeout; for (Message *message = mReassemblyList.GetHead(); message; message = next) { - next = message->GetNext(); - timeout = message->GetTimeout(); + next = message->GetNext(); - if (timeout > 0) + if (message->GetTimeout() > 0) { - message->SetTimeout(timeout - 1); + message->DecrementTimeout(); } else { mReassemblyList.Dequeue(*message); LogMessage(kMessageReassemblyDrop, *message, NULL, OT_ERROR_REASSEMBLY_TIMEOUT); - if (message->GetType() == Message::kTypeIp6) { mIpCounters.mRxFailure++; @@ -1469,10 +1513,7 @@ void MeshForwarder::HandleReassemblyTimer(void) } } - if (mReassemblyList.GetHead() != NULL) - { - mReassemblyTimer.Start(kStateUpdatePeriod); - } + return mReassemblyList.GetHead() != NULL; } void MeshForwarder::HandleLowpanHC(uint8_t * aFrame, @@ -1547,7 +1588,7 @@ otError MeshForwarder::HandleDatagram(Message & aMessage, return netif.GetIp6().HandleDatagram(aMessage, &netif, netif.GetInterfaceId(), &aLinkInfo, false); } -otError MeshForwarder::GetFramePriority(uint8_t * aFrame, +otError MeshForwarder::GetFramePriority(const uint8_t * aFrame, uint8_t aFrameLength, const Mac::Address &aMacSource, const Mac::Address &aMacDest, diff --git a/src/core/thread/mesh_forwarder.hpp b/src/core/thread/mesh_forwarder.hpp index 511c19c57..091da3438 100644 --- a/src/core/thread/mesh_forwarder.hpp +++ b/src/core/thread/mesh_forwarder.hpp @@ -64,6 +64,103 @@ enum * @{ */ +/** + * This class reprents an IPv6 fragment priority entry + * + */ +class FragmentPriorityEntry +{ +public: + /** + * This method returns the fragment datagram tag value. + * + * @returns The fragment datagram tag value. + * + */ + uint16_t GetDatagramTag(void) const { return mDatagramTag; } + + /** + * This method sets the fragment datagram tag value. + * + * @param[in] aDatagramTag The fragment datagram tag value. + * + */ + void SetDatagramTag(uint16_t aDatagramTag) { mDatagramTag = aDatagramTag; } + + /** + * This method returns the source Rloc16 of the fragment. + * + * @returns The source Rloc16 value. + * + */ + uint16_t GetSrcRloc16(void) const { return mSrcRloc16; } + + /** + * This method sets the source Rloc16 value of the fragment. + * + * @param[in] aSrcRloc16 The source Rloc16 value. + * + */ + void SetSrcRloc16(uint16_t aSrcRloc16) { mSrcRloc16 = aSrcRloc16; } + + /** + * This method returns the fragment priority value. + * + * @returns The fragment priority value. + * + */ + uint8_t GetPriority(void) const { return mPriority; } + + /** + * This method sets the fragment priotity value. + * + * @param[in] aPriority The fragment priority value. + * + */ + void SetPriority(uint8_t aPriority) { mPriority = aPriority; } + + /** + * This method returns the fragment priority entry's remaining lifetime. + * + * @returns The fragment priority entry's remaining lifetime. + * + */ + uint8_t GetLifetime(void) const { return mLifetime; } + + /** + * This method sets the remaining lifetime of the fragment priority entry. + * + * @param[in] aLifetime The remaining lifetime of the fragment priority entry (in seconds). + * + */ + void SetLifetime(uint8_t aLifetime) + { + if (aLifetime > kMaxLifeTime) + { + aLifetime = kMaxLifeTime; + } + + mLifetime = aLifetime; + } + + /** + * This method decrements the entry lifetime. + * + */ + void DecrementLifetime(void) { mLifetime--; } + +private: + enum + { + kMaxLifeTime = 5, ///< The maximum lifetime of the fragment entry (in seconds). + }; + + uint16_t mSrcRloc16; ///< The source Rloc16 of the datagram. + uint16_t mDatagramTag; ///< The datagram tag of the fragment header. + uint8_t mPriority : 3; ///< The priority level of the first fragment. + uint8_t mLifetime : 3; ///< The lifetime of the entry (in seconds). 0 means the entry is invalid. +}; + /** * This class implements mesh forwarding within Thread. * @@ -238,11 +335,18 @@ public: private: enum { - kStateUpdatePeriod = 1000, ///< State update period in milliseconds. + kStateUpdatePeriod = 1000, ///< State update period in milliseconds. + kDefaultMsgPriority = Message::kPriorityNormal, ///< Default message priority. }; enum { + /** + * The number of fragment priority entries. + * + */ + kNumFragmentPriorityEntries = OPENTHREAD_CONFIG_NUM_FRAGMENT_PRIORITY_ENTRIES, + /** * Maximum number of tx attempts by `MeshForwarder` for an outbound indirect frame (for a sleepy child). The * `MeshForwader` attempts occur following the reception of a new data request command (a new data poll) from @@ -277,6 +381,7 @@ private: const Mac::Address &aMeshSource, const Mac::Address &aMeshDest); + otError GetMeshHeader(const uint8_t *&aFrame, uint8_t &aFrameLength, Lowpan::MeshHeader &aMeshHeader); otError SkipMeshHeader(const uint8_t *&aFrame, uint8_t &aFrameLength); otError DecompressIp6Header(const uint8_t * aFrame, uint8_t aFrameLength, @@ -313,18 +418,29 @@ private: const Mac::Address & aMacDest, const otThreadLinkInfo &aLinkInfo); void HandleDataRequest(const Mac::Address &aMacSource, const otThreadLinkInfo &aLinkInfo); - otError SendPoll(Message &aMessage, Mac::Frame &aFrame); - otError SendMesh(Message &aMessage, Mac::Frame &aFrame); - otError SendFragment(Message &aMessage, Mac::Frame &aFrame); - otError SendEmptyFrame(Mac::Frame &aFrame, bool aAckRequest); - otError UpdateIp6Route(Message &aMessage); - otError UpdateIp6RouteFtd(Ip6::Header &ip6Header); - otError UpdateMeshRoute(Message &aMessage); - otError HandleDatagram(Message &aMessage, const otThreadLinkInfo &aLinkInfo, const Mac::Address &aMacSource); - void ClearReassemblyList(void); - otError RemoveMessageFromSleepyChild(Message &aMessage, Child &aChild); - void RemoveMessage(Message &aMessage); - void HandleDiscoverComplete(void); + + static otError GetFragmentHeader(const uint8_t * aFrame, + uint8_t aFrameLength, + Lowpan::FragmentHeader &aFragmentHeader); + + otError SendPoll(Message &aMessage, Mac::Frame &aFrame); + otError SendMesh(Message &aMessage, Mac::Frame &aFrame); + otError SendFragment(Message &aMessage, Mac::Frame &aFrame); + otError SendEmptyFrame(Mac::Frame &aFrame, bool aAckRequest); + otError UpdateIp6Route(Message &aMessage); + otError UpdateIp6RouteFtd(Ip6::Header &ip6Header); + otError UpdateMeshRoute(Message &aMessage); + bool UpdateReassemblyList(void); + bool UpdateFragmentLifetime(void); + void UpdateFragmentPriority(Lowpan::FragmentHeader &aFragmentHeader, + uint8_t aFragmentLength, + uint16_t aSrcRloc16, + uint8_t aPriority); + otError HandleDatagram(Message &aMessage, const otThreadLinkInfo &aLinkInfo, const Mac::Address &aMacSource); + void ClearReassemblyList(void); + otError RemoveMessageFromSleepyChild(Message &aMessage, Child &aChild); + void RemoveMessage(Message &aMessage); + void HandleDiscoverComplete(void); void HandleReceivedFrame(Mac::Frame &aFrame); otError HandleFrameRequest(Mac::Frame &aFrame); @@ -333,16 +449,25 @@ private: static void HandleDiscoverTimer(Timer &aTimer); void HandleDiscoverTimer(void); - static void HandleReassemblyTimer(Timer &aTimer); - void HandleReassemblyTimer(void); + static void HandleUpdateTimer(Timer &aTimer); + void HandleUpdateTimer(void); static void ScheduleTransmissionTask(Tasklet &aTasklet); void ScheduleTransmissionTask(void); - otError GetFramePriority(uint8_t * aFrame, + otError GetFramePriority(const uint8_t * aFrame, uint8_t aFrameLength, const Mac::Address &aMacSource, const Mac::Address &aMacDest, uint8_t & aPriority); + otError GetFragmentPriority(Lowpan::FragmentHeader &aFragmentHeader, uint16_t aSrcRloc16, uint8_t &aPriority); + otError GetForwardFramePriority(const uint8_t * aFrame, + uint8_t aFrameLength, + const Mac::Address &aMacDest, + const Mac::Address &aMacSource, + uint8_t & aPriority); + + FragmentPriorityEntry *FindFragmentPriorityEntry(uint16_t aTag, uint16_t aSrcRloc16); + FragmentPriorityEntry *GetUnusedFragementPriorityEntry(void); otError GetDestinationRlocByServiceAloc(uint16_t aServiceAloc, uint16_t &aMeshDest); @@ -409,7 +534,7 @@ private: #endif // #if (OPENTHREAD_CONFIG_LOG_LEVEL >= OT_LOG_LEVEL_NOTE) && (OPENTHREAD_CONFIG_LOG_MAC == 1) TimerMilli mDiscoverTimer; - TimerMilli mReassemblyTimer; + TimerMilli mUpdateTimer; PriorityQueue mSendQueue; MessageQueue mReassemblyList; @@ -441,6 +566,7 @@ private: otIpCounters mIpCounters; #if OPENTHREAD_FTD + FragmentPriorityEntry mFragmentEntries[kNumFragmentPriorityEntries]; MessageQueue mResolvingQueue; SourceMatchController mSourceMatchController; uint32_t mSendMessageFrameCounter; diff --git a/src/core/thread/mesh_forwarder_ftd.cpp b/src/core/thread/mesh_forwarder_ftd.cpp index 3e6fd2e43..918da51bb 100644 --- a/src/core/thread/mesh_forwarder_ftd.cpp +++ b/src/core/thread/mesh_forwarder_ftd.cpp @@ -963,6 +963,8 @@ void MeshForwarder::HandleMesh(uint8_t * aFrame, } else if (meshHeader.GetHopsLeft() > 0) { + uint8_t priority = kDefaultMsgPriority; + netif.GetMle().ResolveRoutingLoops(aMacSource.GetShort(), meshDest.GetShort()); SuccessOrExit(error = CheckReachability(aFrame, aFrameLength, meshSource, meshDest)); @@ -970,7 +972,8 @@ void MeshForwarder::HandleMesh(uint8_t * aFrame, meshHeader.SetHopsLeft(meshHeader.GetHopsLeft() - 1); meshHeader.AppendTo(aFrame); - VerifyOrExit((message = GetInstance().GetMessagePool().New(Message::kType6lowpan, 0)) != NULL, + GetForwardFramePriority(aFrame, aFrameLength, meshDest, meshSource, priority); + VerifyOrExit((message = GetInstance().GetMessagePool().New(Message::kType6lowpan, priority)) != NULL, error = OT_ERROR_NO_BUFS); SuccessOrExit(error = message->SetLength(aFrameLength)); message->Write(0, aFrameLength, aFrame); @@ -1023,6 +1026,155 @@ exit: return; } +bool MeshForwarder::UpdateFragmentLifetime(void) +{ + bool shouldRun = false; + + for (size_t i = 0; i < OT_ARRAY_LENGTH(mFragmentEntries); i++) + { + if (mFragmentEntries[i].GetLifetime() != 0) + { + mFragmentEntries[i].DecrementLifetime(); + + if (mFragmentEntries[i].GetLifetime() != 0) + { + shouldRun = true; + } + } + } + + return shouldRun; +} + +void MeshForwarder::UpdateFragmentPriority(Lowpan::FragmentHeader &aFragmentHeader, + uint8_t aFragmentLength, + uint16_t aSrcRloc16, + uint8_t aPriority) +{ + FragmentPriorityEntry *entry; + + if (aFragmentHeader.GetDatagramOffset() == 0) + { + VerifyOrExit((entry = GetUnusedFragementPriorityEntry()) != NULL); + + entry->SetDatagramTag(aFragmentHeader.GetDatagramTag()); + entry->SetSrcRloc16(aSrcRloc16); + entry->SetPriority(aPriority); + entry->SetLifetime(kReassemblyTimeout); + + if (!mUpdateTimer.IsRunning()) + { + mUpdateTimer.Start(kStateUpdatePeriod); + } + } + else + { + VerifyOrExit((entry = FindFragmentPriorityEntry(aFragmentHeader.GetDatagramTag(), aSrcRloc16)) != NULL); + + entry->SetLifetime(kReassemblyTimeout); + + if (aFragmentHeader.GetDatagramOffset() + aFragmentLength >= aFragmentHeader.GetDatagramSize()) + { + entry->SetLifetime(0); + } + } + +exit: + return; +} + +FragmentPriorityEntry *MeshForwarder::FindFragmentPriorityEntry(uint16_t aTag, uint16_t aSrcRloc16) +{ + size_t i; + + for (i = 0; i < OT_ARRAY_LENGTH(mFragmentEntries); i++) + { + if ((mFragmentEntries[i].GetLifetime() != 0) && (mFragmentEntries[i].GetDatagramTag() == aTag) && + (mFragmentEntries[i].GetSrcRloc16() == aSrcRloc16)) + { + break; + } + } + + return (i >= OT_ARRAY_LENGTH(mFragmentEntries)) ? NULL : &mFragmentEntries[i]; +} + +FragmentPriorityEntry *MeshForwarder::GetUnusedFragementPriorityEntry(void) +{ + size_t i; + + for (i = 0; i < OT_ARRAY_LENGTH(mFragmentEntries); i++) + { + if (mFragmentEntries[i].GetLifetime() == 0) + { + break; + } + } + + return (i >= OT_ARRAY_LENGTH(mFragmentEntries)) ? NULL : &mFragmentEntries[i]; +} + +otError MeshForwarder::GetFragmentPriority(Lowpan::FragmentHeader &aFragmentHeader, + uint16_t aSrcRloc16, + uint8_t & aPriority) +{ + otError error = OT_ERROR_NONE; + FragmentPriorityEntry *entry; + + VerifyOrExit((entry = FindFragmentPriorityEntry(aFragmentHeader.GetDatagramTag(), aSrcRloc16)) != NULL, + error = OT_ERROR_NOT_FOUND); + aPriority = entry->GetPriority(); + +exit: + return error; +} + +otError MeshForwarder::GetForwardFramePriority(const uint8_t * aFrame, + uint8_t aFrameLength, + const Mac::Address &aMacDest, + const Mac::Address &aMacSource, + uint8_t & aPriority) +{ + otError error = OT_ERROR_NONE; + bool isFragment = false; + Lowpan::MeshHeader meshHeader; + Lowpan::FragmentHeader fragmentHeader; + + SuccessOrExit(error = GetMeshHeader(aFrame, aFrameLength, meshHeader)); + aFrame += meshHeader.GetHeaderLength(); + aFrameLength -= meshHeader.GetHeaderLength(); + + if (GetFragmentHeader(aFrame, aFrameLength, fragmentHeader) == OT_ERROR_NONE) + { + isFragment = true; + aFrame += fragmentHeader.GetHeaderLength(); + aFrameLength -= fragmentHeader.GetHeaderLength(); + + if (fragmentHeader.GetDatagramOffset() > 0) + { + // Get priority from the pre-buffered info + ExitNow(error = GetFragmentPriority(fragmentHeader, meshHeader.GetSource(), aPriority)); + } + } + + // Get priority from Ipv6 header or UDP destination port directly + error = GetFramePriority(aFrame, aFrameLength, aMacSource, aMacDest, aPriority); + +exit: + if (error != OT_ERROR_NONE) + { + otLogNoteMac("Failed to get forwarded frame priority, error:%s, len:%d, dst:%s, src:%s", + otThreadErrorToString(error), aFrameLength, aMacDest.ToString().AsCString(), + aMacSource.ToString().AsCString()); + } + else if (isFragment) + { + UpdateFragmentPriority(fragmentHeader, aFrameLength, meshHeader.GetSource(), aPriority); + } + + return error; +} + #if OPENTHREAD_ENABLE_SERVICE otError MeshForwarder::GetDestinationRlocByServiceAloc(uint16_t aServiceAloc, uint16_t &aMeshDest) {