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.
This commit is contained in:
Fletcher Dunn
2020-04-06 08:48:59 -07:00
parent 97f5ad7d64
commit 1ce965d63e
7 changed files with 289 additions and 217 deletions
+1
View File
@@ -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"
+1
View File
@@ -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',
@@ -13,6 +13,7 @@
#include <tier1/utlhashmap.h>
#include <tier1/netadr.h>
#include "steamnetworkingsockets_lowlevel.h"
#include "../steamnetworkingsockets_thinker.h"
#include "keypair.h"
#include "crypto.h"
#include "crypto_25519.h"
@@ -22,6 +22,7 @@
#include "steamnetworkingsockets_lowlevel.h"
#include "../steamnetworkingsockets_internal.h"
#include "../steamnetworkingsockets_thinker.h"
#include <vstdlib/random.h>
#include <tier1/utlpriorityqueue.h>
#include <tier1/utllinkedlist.h>
@@ -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<IThinker*,ThinkerLess,ThinkerSetIndex> 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
@@ -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
@@ -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 <tier1/utlpriorityqueue.h>
#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<IThinker*,ThinkerLess,ThinkerSetIndex> 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
@@ -0,0 +1,80 @@
//====== Copyright Valve Corporation, All rights reserved. ====================
#ifndef STEAMNETWORKINGSOCKETS_THINKER_H
#define STEAMNETWORKINGSOCKETS_THINKER_H
#pragma once
#include <steam/steamnetworkingtypes.h>
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