From 1ce965d63e31d22d6c221fc7e4e48c3dd45ea2db Mon Sep 17 00:00:00 2001 From: Fletcher Dunn Date: Mon, 6 Apr 2020 08:43:30 -0700 Subject: [PATCH] Break out IThinker system to separate files. I am using this on the relay server to make periodic thinking of clients more efficient, instead of sweeping and polling. --- src/CMakeLists.txt | 1 + src/meson.build | 1 + .../steamnetworkingsockets_connections.h | 1 + .../steamnetworkingsockets_lowlevel.cpp | 160 +------------- .../steamnetworkingsockets_lowlevel.h | 68 +----- .../steamnetworkingsockets_thinker.cpp | 195 ++++++++++++++++++ .../steamnetworkingsockets_thinker.h | 80 +++++++ 7 files changed, 289 insertions(+), 217 deletions(-) create mode 100644 src/steamnetworkingsockets/steamnetworkingsockets_thinker.cpp create mode 100644 src/steamnetworkingsockets/steamnetworkingsockets_thinker.h diff --git a/src/CMakeLists.txt b/src/CMakeLists.txt index f62e2d4..dbd29c4 100644 --- a/src/CMakeLists.txt +++ b/src/CMakeLists.txt @@ -26,6 +26,7 @@ set(GNS_SRCS "steamnetworkingsockets/steamnetworkingsockets_certstore.cpp" "steamnetworkingsockets/steamnetworkingsockets_shared.cpp" "steamnetworkingsockets/steamnetworkingsockets_stats.cpp" + "steamnetworkingsockets/steamnetworkingsockets_thinker.cpp" "tier0/dbg.cpp" "tier0/platformtime.cpp" "tier1/bitstring.cpp" diff --git a/src/meson.build b/src/meson.build index 283f696..719c96e 100644 --- a/src/meson.build +++ b/src/meson.build @@ -82,6 +82,7 @@ sources = [ 'steamnetworkingsockets/steamnetworkingsockets_certstore.cpp', 'steamnetworkingsockets/steamnetworkingsockets_shared.cpp', 'steamnetworkingsockets/steamnetworkingsockets_stats.cpp', + 'steamnetworkingsockets/steamnetworkingsockets_thinker.cpp', 'tier0/dbg.cpp', 'tier0/platformtime.cpp', 'tier1/bitstring.cpp', diff --git a/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_connections.h b/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_connections.h index 1c5a44c..f20aec9 100644 --- a/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_connections.h +++ b/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_connections.h @@ -13,6 +13,7 @@ #include #include #include "steamnetworkingsockets_lowlevel.h" +#include "../steamnetworkingsockets_thinker.h" #include "keypair.h" #include "crypto.h" #include "crypto_25519.h" diff --git a/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_lowlevel.cpp b/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_lowlevel.cpp index da10af0..9b945bc 100644 --- a/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_lowlevel.cpp +++ b/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_lowlevel.cpp @@ -22,6 +22,7 @@ #include "steamnetworkingsockets_lowlevel.h" #include "../steamnetworkingsockets_internal.h" +#include "../steamnetworkingsockets_thinker.h" #include #include #include @@ -565,7 +566,7 @@ static CPacketLagger s_packetLagQueue; static std::thread *s_pThreadSteamDatagram = nullptr; -static void WakeSteamDatagramThread() +void WakeSteamDatagramThread() { #if defined( _WIN32 ) if ( s_hEventWakeThread != INVALID_HANDLE_VALUE ) @@ -1202,156 +1203,6 @@ void ProcessPendingDestroyClosedRawUDPSockets() s_vecRawSocketsPendingDeletion.RemoveAll(); } -///////////////////////////////////////////////////////////////////////////// -// -// Periodic processing -// -///////////////////////////////////////////////////////////////////////////// - -struct ThinkerLess -{ - bool operator()( const IThinker *a, const IThinker *b ) const - { - return a->GetNextThinkTime() > b->GetNextThinkTime(); - } -}; -class ThinkerSetIndex -{ -public: - static void SetIndex( IThinker *p, int idx ) { p->m_queueIndex = idx; } -}; - -static CUtlPriorityQueue s_queueThinkers; - -IThinker::IThinker() -: m_usecNextThinkTime( k_nThinkTime_Never ) -, m_queueIndex( -1 ) -{ -} - -IThinker::~IThinker() -{ - ClearNextThinkTime(); -} - -#ifdef __GNUC__ - // older steamrt:scout gcc requires this also, probably getting confused by unbalanced push/pop - #pragma GCC diagnostic ignored "-Wstrict-overflow" -#endif - -void IThinker::SetNextThinkTime( SteamNetworkingMicroseconds usecTargetThinkTime ) -{ - // Protect against us blowing up because of an invalid think time - if ( usecTargetThinkTime <= 0 ) - { - AssertMsg1( false, "Attempt to set target think time to %lld", (long long)usecTargetThinkTime ); - usecTargetThinkTime = Plat_USTime() + 5000; - } - - // Clearing it? - if ( usecTargetThinkTime == k_nThinkTime_Never ) - { - if ( m_queueIndex >= 0 ) - { - Assert( s_queueThinkers.Element( m_queueIndex ) == this ); - s_queueThinkers.RemoveAt( m_queueIndex ); - Assert( m_queueIndex == -1 ); - } - - m_usecNextThinkTime = k_nThinkTime_Never; - return; - } - - // Save current time when the next thinker wants service - SteamNetworkingMicroseconds usecNextWake = ( s_queueThinkers.Count() > 0 ) ? s_queueThinkers.ElementAtHead()->GetNextThinkTime() : k_nThinkTime_Never; - - // Not currently scheduled? - if ( m_queueIndex < 0 ) - { - Assert( m_usecNextThinkTime == k_nThinkTime_Never ); - m_usecNextThinkTime = usecTargetThinkTime; - s_queueThinkers.Insert( this ); - } - else - { - - // We're already scheduled. - Assert( s_queueThinkers.Element( m_queueIndex ) == this ); - Assert( m_usecNextThinkTime != k_nThinkTime_Never ); - - // Set the new schedule time - m_usecNextThinkTime = usecTargetThinkTime; - - // And update our position in the queue - s_queueThinkers.RevaluateElement( m_queueIndex ); - } - - // Check that we know our place - Assert( m_queueIndex >= 0 ); - Assert( s_queueThinkers.Element( m_queueIndex ) == this ); - - // Do we need service before we were previously schedule to wake up? - // If so, wake the thread now so that it can redo its schedule work - // NOTE: On Windows we could use a waitable timer. This would avoid - // waking up the service thread just to re-schedule when it should - // wake up for real. - if ( m_usecNextThinkTime < usecNextWake ) - WakeSteamDatagramThread(); -} - -void ProcessThinkers() -{ - - // Until the queue is empty - int nIterations = 0; - while ( s_queueThinkers.Count() > 0 ) - { - - // Grab the head element - IThinker *pNextThinker = s_queueThinkers.ElementAtHead(); - - // Refetch timestamp each time. The reason is that certain thinkers - // may pass through to other systems (e.g. fake lag) that fetch the time. - // If we don't update the time here, that code may have used the newer - // timestamp (e.g. to mark when a packet was received) and then - // in our next iteration, we will use an older timestamp to process - // a thinker. - SteamNetworkingMicroseconds usecNow = SteamNetworkingSockets_GetLocalTimestamp(); - - // Scheduled too far in the future? - if ( pNextThinker->GetNextThinkTime() >= usecNow ) - { - // Keep waiting - break; - } - - ++nIterations; - if ( nIterations > 10000 ) - { - AssertMsg1( false, "Processed thinkers %d times -- probably one thinker keeps requesting an immediate wakeup call.", nIterations ); - break; - } - - // Go ahead and clear his think time now and remove him - // from the heap. He needs to schedule a new think time - // if heeds service again. For thinkers that need frequent - // service, removing them and then re-inserting them when - // they reschedule is a bit of extra work that could be - // optimized by trying to not remove them now, but adjusting - // them once we know when they want to think. But this - // is probably just a bit too complicated for the expected - // benefit. If the number of total Thinkers is relatively - // small (which it probably will be), the heap operations - // are probably negligible. - pNextThinker->ClearNextThinkTime(); - - // Execute callback. (Note: this could result - // in self-destruction or essentially any change - // to the rest of the queue.) - pNextThinker->Think( usecNow ); - } -} - ///////////////////////////////////////////////////////////////////////////// // // Service thread @@ -1370,9 +1221,9 @@ static bool SteamNetworkingSockets_InternalPoll( int msWait, bool bManualPoll ) Assert( SteamDatagramTransportLock::s_nLocked == 1 ); // exactly once // Figure out how long to sleep - if ( s_queueThinkers.Count() > 0 ) + IThinker *pNextThinker = Thinker_GetNextScheduled(); + if ( pNextThinker ) { - IThinker *pNextThinker = s_queueThinkers.ElementAtHead(); // Calc wait time to wake up as late as possible, // rounded up to the nearest millisecond. @@ -1432,7 +1283,7 @@ static bool SteamNetworkingSockets_InternalPoll( int msWait, bool bManualPoll ) } // Check for periodic processing - ProcessThinkers(); + Thinker_ProcessThinkers(); // Close any sockets pending delete, if we discarded a server // We can close the sockets safely now, because we know we're @@ -2019,7 +1870,6 @@ void SteamNetworkingSocketsLowLevelDecRef() void SteamNetworkingSocketsLowLevelValidate( CValidator &validator ) { ValidateRecursive( s_vecRawSockets ); - ValidateObj( s_queueThinkers ); } #endif diff --git a/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_lowlevel.h b/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_lowlevel.h index 1d54095..185f7a4 100644 --- a/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_lowlevel.h +++ b/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_lowlevel.h @@ -224,66 +224,6 @@ private: static void CallbackRecvPacket( const void *pPkt, int cbPkt, const netadr_t &adrFrom, CSharedSocket *pSock ); }; -///////////////////////////////////////////////////////////////////////////// -// -// Periodic processing -// -///////////////////////////////////////////////////////////////////////////// - -const SteamNetworkingMicroseconds k_nThinkTime_Never = INT64_MAX; -class ThinkerSetIndex; - -class IThinker -{ -public: - virtual ~IThinker(); - - /// Callback to do whatever periodic processing you need. If you don't - /// explicitly call SetNextThinkTime inside this function, then thinking - /// will be disabled. - /// - /// Think callbacks will always happen from the service thread, - /// with the lock held. - /// - /// Note that we assume a limited precision of the thread scheduler, - /// and you won't get your callback exactly when you request. - virtual void Think( SteamNetworkingMicroseconds usecNow ) = 0; - - /// Called to set when you next want to get your Think() callback. - /// You should assume that, due to scheduler inaccuracy, you could - /// get your callback 1 or 2 ms late. - void SetNextThinkTime( SteamNetworkingMicroseconds usecTargetThinkTime ); - - /// Adjust schedule time to the earlier of the current schedule time, - /// or the given time. - inline void EnsureMinThinkTime( SteamNetworkingMicroseconds usecTargetThinkTime ) - { - if ( usecTargetThinkTime < m_usecNextThinkTime ) - SetNextThinkTime( usecTargetThinkTime ); - } - - /// Clear the next think time. You won't get a callback. - void ClearNextThinkTime() { SetNextThinkTime( k_nThinkTime_Never ); } - - /// Request an immediate wakeup. - void SetNextThinkTimeASAP() { EnsureMinThinkTime( 1 ); } - - /// Fetch time when the next Think() call is currently scheduled to - /// happen. - inline SteamNetworkingMicroseconds GetNextThinkTime() const { return m_usecNextThinkTime; } - - /// Return true if we are scheduled to get our callback - inline bool IsScheduled() const { return m_usecNextThinkTime != k_nThinkTime_Never; } - -protected: - IThinker(); - -private: - SteamNetworkingMicroseconds m_usecNextThinkTime; - int m_queueIndex; - friend class ThinkerSetIndex; -}; - ///////////////////////////////////////////////////////////////////////////// // // Misc low level service thread stuff @@ -353,10 +293,14 @@ extern void SteamNetworkingSocketsLowLevelValidate( CValidator &validator ); #endif /// Fetch current time -SteamNetworkingMicroseconds SteamNetworkingSockets_GetLocalTimestamp(); +extern SteamNetworkingMicroseconds SteamNetworkingSockets_GetLocalTimestamp(); /// Set debug output hook -void SteamNetworkingSockets_SetDebugOutputFunction( ESteamNetworkingSocketsDebugOutputType eDetailLevel, FSteamNetworkingSocketsDebugOutput pfnFunc ); +extern void SteamNetworkingSockets_SetDebugOutputFunction( ESteamNetworkingSocketsDebugOutputType eDetailLevel, FSteamNetworkingSocketsDebugOutput pfnFunc ); + +/// Wake up the service thread ASAP. Intended to be called from other threads, +/// but is safe to call from the service thread as well. +extern void WakeSteamDatagramThread(); } // namespace SteamNetworkingSocketsLib diff --git a/src/steamnetworkingsockets/steamnetworkingsockets_thinker.cpp b/src/steamnetworkingsockets/steamnetworkingsockets_thinker.cpp new file mode 100644 index 0000000..8a03d1c --- /dev/null +++ b/src/steamnetworkingsockets/steamnetworkingsockets_thinker.cpp @@ -0,0 +1,195 @@ +//====== Copyright Valve Corporation, All rights reserved. ==================== + +#ifdef __GNUC__ + // src/public/tier0/basetypes.h:104:30: error: assuming signed overflow does not occur when assuming that (X + c) < X is always false [-Werror=strict-overflow] + // current steamrt:scout gcc "g++ (SteamRT 4.8.4-1ubuntu15~12.04+steamrt1.2+srt1) 4.8.4" requires this at the top due to optimizations + #pragma GCC diagnostic ignored "-Wstrict-overflow" +#endif + +#include + +#include "steamnetworkingsockets_thinker.h" + +#ifdef IS_STEAMDATAGRAMROUTER + #include "router/sdr.h" +#else + #include "clientlib/steamnetworkingsockets_lowlevel.h" +#endif + +// memdbgon must be the last include file in a .cpp file!!! +#include "tier0/memdbgon.h" + +namespace SteamNetworkingSocketsLib { + +///////////////////////////////////////////////////////////////////////////// +// +// Periodic processing +// +///////////////////////////////////////////////////////////////////////////// + +struct ThinkerLess +{ + bool operator()( const IThinker *a, const IThinker *b ) const + { + return a->GetNextThinkTime() > b->GetNextThinkTime(); + } +}; +class ThinkerSetIndex +{ +public: + static void SetIndex( IThinker *p, int idx ) { p->m_queueIndex = idx; } +}; + +static CUtlPriorityQueue s_queueThinkers; + +IThinker::IThinker() +: m_usecNextThinkTime( k_nThinkTime_Never ) +, m_queueIndex( -1 ) +{ +} + +IThinker::~IThinker() +{ + ClearNextThinkTime(); +} + +#ifdef __GNUC__ + // older steamrt:scout gcc requires this also, probably getting confused by unbalanced push/pop + #pragma GCC diagnostic ignored "-Wstrict-overflow" +#endif + +void IThinker::SetNextThinkTime( SteamNetworkingMicroseconds usecTargetThinkTime ) +{ + // Protect against us blowing up because of an invalid think time + if ( usecTargetThinkTime <= 0 ) + { +#ifndef _DEBUG // Running into this a lot while streaming and breaking in the debugger... + AssertMsg1( false, "Attempt to set target think time to %lld", (long long)usecTargetThinkTime ); +#endif + usecTargetThinkTime = Plat_USTime() + 5000; + } + + // Clearing it? + if ( usecTargetThinkTime == k_nThinkTime_Never ) + { + if ( m_queueIndex >= 0 ) + { + Assert( s_queueThinkers.Element( m_queueIndex ) == this ); + s_queueThinkers.RemoveAt( m_queueIndex ); + Assert( m_queueIndex == -1 ); + } + + m_usecNextThinkTime = k_nThinkTime_Never; + return; + } + + // Save current time when the next thinker wants service + #ifndef IS_STEAMDATAGRAMROUTER + SteamNetworkingMicroseconds usecNextWake = ( s_queueThinkers.Count() > 0 ) ? s_queueThinkers.ElementAtHead()->GetNextThinkTime() : k_nThinkTime_Never; + #endif + + // Not currently scheduled? + if ( m_queueIndex < 0 ) + { + Assert( m_usecNextThinkTime == k_nThinkTime_Never ); + m_usecNextThinkTime = usecTargetThinkTime; + s_queueThinkers.Insert( this ); + } + else + { + + // We're already scheduled. + Assert( s_queueThinkers.Element( m_queueIndex ) == this ); + Assert( m_usecNextThinkTime != k_nThinkTime_Never ); + + // Set the new schedule time + m_usecNextThinkTime = usecTargetThinkTime; + + // And update our position in the queue + s_queueThinkers.RevaluateElement( m_queueIndex ); + } + + // Check that we know our place + Assert( m_queueIndex >= 0 ); + Assert( s_queueThinkers.Element( m_queueIndex ) == this ); + + #ifndef IS_STEAMDATAGRAMROUTER + // Do we need service before we were previously schedule to wake up? + // If so, wake the thread now so that it can redo its schedule work + // NOTE: On Windows we could use a waitable timer. This would avoid + // waking up the service thread just to re-schedule when it should + // wake up for real. + if ( m_usecNextThinkTime < usecNextWake ) + WakeSteamDatagramThread(); + #endif +} + +IThinker *Thinker_GetNextScheduled() +{ + if ( s_queueThinkers.Count() == 0 ) + return nullptr; + return s_queueThinkers.ElementAtHead(); +} + +void Thinker_ProcessThinkers() +{ + + // Until the queue is empty + int nIterations = 0; + while ( s_queueThinkers.Count() > 0 ) + { + + // Grab the head element + IThinker *pNextThinker = s_queueThinkers.ElementAtHead(); + + // Refetch timestamp each time. The reason is that certain thinkers + // may pass through to other systems (e.g. fake lag) that fetch the time. + // If we don't update the time here, that code may have used the newer + // timestamp (e.g. to mark when a packet was received) and then + // in our next iteration, we will use an older timestamp to process + // a thinker. + SteamNetworkingMicroseconds usecNow = SteamNetworkingSockets_GetLocalTimestamp(); + + // Scheduled too far in the future? + if ( pNextThinker->GetNextThinkTime() >= usecNow ) + { + // Keep waiting + break; + } + + ++nIterations; + if ( nIterations > 10000 ) + { + AssertMsg1( false, "Processed thinkers %d times -- probably one thinker keeps requesting an immediate wakeup call.", nIterations ); + break; + } + + // Go ahead and clear his think time now and remove him + // from the heap. He needs to schedule a new think time + // if heeds service again. For thinkers that need frequent + // service, removing them and then re-inserting them when + // they reschedule is a bit of extra work that could be + // optimized by trying to not remove them now, but adjusting + // them once we know when they want to think. But this + // is probably just a bit too complicated for the expected + // benefit. If the number of total Thinkers is relatively + // small (which it probably will be), the heap operations + // are probably negligible. + pNextThinker->ClearNextThinkTime(); + + // Execute callback. (Note: this could result + // in self-destruction or essentially any change + // to the rest of the queue.) + pNextThinker->Think( usecNow ); + } +} + +#ifdef DBGFLAG_VALIDATE +void Thinker_ValidateStatics( CValidator &validator ) +{ + ValidateObj( s_queueThinkers ); +} +#endif + +} // namespace SteamNetworkingSocketsLib + diff --git a/src/steamnetworkingsockets/steamnetworkingsockets_thinker.h b/src/steamnetworkingsockets/steamnetworkingsockets_thinker.h new file mode 100644 index 0000000..4ff96aa --- /dev/null +++ b/src/steamnetworkingsockets/steamnetworkingsockets_thinker.h @@ -0,0 +1,80 @@ +//====== Copyright Valve Corporation, All rights reserved. ==================== + +#ifndef STEAMNETWORKINGSOCKETS_THINKER_H +#define STEAMNETWORKINGSOCKETS_THINKER_H +#pragma once + +#include + +namespace SteamNetworkingSocketsLib { + +///////////////////////////////////////////////////////////////////////////// +// +// Periodic processing +// +///////////////////////////////////////////////////////////////////////////// + +const SteamNetworkingMicroseconds k_nThinkTime_Never = INT64_MAX; +class ThinkerSetIndex; + +class IThinker +{ +public: + virtual ~IThinker(); + + /// Callback to do whatever periodic processing you need. If you don't + /// explicitly call SetNextThinkTime inside this function, then thinking + /// will be disabled. + /// + /// Think callbacks will always happen from the service thread, + /// with the lock held. + /// + /// Note that we assume a limited precision of the thread scheduler, + /// and you won't get your callback exactly when you request. + virtual void Think( SteamNetworkingMicroseconds usecNow ) = 0; + + /// Called to set when you next want to get your Think() callback. + /// You should assume that, due to scheduler inaccuracy, you could + /// get your callback 1 or 2 ms late. + void SetNextThinkTime( SteamNetworkingMicroseconds usecTargetThinkTime ); + + /// Adjust schedule time to the earlier of the current schedule time, + /// or the given time. + inline void EnsureMinThinkTime( SteamNetworkingMicroseconds usecTargetThinkTime ) + { + if ( usecTargetThinkTime < m_usecNextThinkTime ) + SetNextThinkTime( usecTargetThinkTime ); + } + + /// Clear the next think time. You won't get a callback. + void ClearNextThinkTime() { SetNextThinkTime( k_nThinkTime_Never ); } + + /// Request an immediate wakeup. + void SetNextThinkTimeASAP() { EnsureMinThinkTime( 1 ); } + + /// Fetch time when the next Think() call is currently scheduled to + /// happen. + inline SteamNetworkingMicroseconds GetNextThinkTime() const { return m_usecNextThinkTime; } + + /// Return true if we are scheduled to get our callback + inline bool IsScheduled() const { return m_usecNextThinkTime != k_nThinkTime_Never; } + +protected: + IThinker(); + +private: + SteamNetworkingMicroseconds m_usecNextThinkTime; + int m_queueIndex; + friend class ThinkerSetIndex; +}; + +extern IThinker *Thinker_GetNextScheduled(); +extern void Thinker_ProcessThinkers(); + +#ifdef DBGFLAG_VALIDATE +extern void Thinker_ValidateStatics( CValidator &validator ); +#endif + +} // namespace SteamNetworkingSocketsLib + +#endif // STEAMNETWORKINGSOCKETS_THINKER_H