Refactor control messages for the 9999th time.

Maybe one day I will be happy with it.  But this definitely cleaned up a
bunch of redundant code and made things simpler.

- Renamed NeedToSend to GetSendReasonOrUpdateNextThinkTime.  It still returns
  a reason code and expected reply urgency, as before.  Now it also updates
  the next time we need to check in, if nothing needs to be sent now.
- Use GetSendReasonOrUpdateNextThinkTime in the base connection class thinking
  code.  There was custom code that was doing exactly what this function did.
- Remove code where different transports were checking the end-to-end stats to
  see when they might need to wake up, to send end-to-end stats.  That's not
  necessary, the base connection class will handle that.  Whenever we actually
  send a packet, we check for opportunitistically sending stats at all levels.
  That is important and hasn't changed.  But when it comes to figuring out
  when to think next and send a packet, we don't need the transport to worry
  about end-to-end stats.
- Removed NeedToSendEndToEndStatsOrAcks.  This was called by tansports, to
  figure out when they needed to wake up to send SNP acks.  But that's
  already handledby the base connection class.
- In a few places we were handling replying to control messages immediately in
  the code that receives those stats.  That's not necessary.  We can just use
  the normal wakeup mechanism.  This is less code and works just as good.
- Simplify and optimize how we return which stats need to be sent.  Previously
  it was two function calls, and all call sites were using them in sequence.
  Now it is a single funciton call that returns a bitmask.  Also, by saving
  this bitmask and passing it into the method that actually populates the
  message, we can avoid redoing this work.  We can also be smarter when
  trimming opportunistic stats that don't fit.
This commit is contained in:
Fletcher Dunn
2020-05-17 12:01:06 -07:00
parent 7ba2324f52
commit 705bc0e4bb
6 changed files with 231 additions and 263 deletions
@@ -2713,67 +2713,34 @@ void CSteamNetworkConnectionBase::CheckConnectionStateAndSetNextThinkTime( Steam
UpdateMinThinkTime( usecNow + 50*1000 );
}
// Check for keepalives of varying urgency.
// Ping aggressively because connection appears to be timing out?
if ( m_statsEndToEnd.m_nReplyTimeoutsSinceLastRecv > 0 || m_statsEndToEnd.m_usecWhenTimeoutStarted > 0 )
// Check for sending keepalives and stats
if ( bCanSendEndToEnd )
{
SteamNetworkingMicroseconds usecSendAggressivePing = Max( m_statsEndToEnd.m_usecTimeLastRecv, m_statsEndToEnd.m_usecLastSendPacketExpectingImmediateReply ) + k_usecAggressivePingInterval;
if ( usecNow >= usecSendAggressivePing )
{
if ( bCanSendEndToEnd )
{
if ( m_statsEndToEnd.m_nReplyTimeoutsSinceLastRecv == 1 )
SpewVerbose( "[%s] Reply timeout, last recv %.1fms ago. Sending keepalive.\n", GetDescription(), ( usecNow - m_statsEndToEnd.m_usecTimeLastRecv ) * 1e-3 );
else
SpewMsg( "[%s] %d reply timeouts, last recv %.1fms ago. Sending keepalive.\n", GetDescription(), m_statsEndToEnd.m_nReplyTimeoutsSinceLastRecv, ( usecNow - m_statsEndToEnd.m_usecTimeLastRecv ) * 1e-3 );
Assert( m_statsEndToEnd.BNeedToSendPingImmediate( usecNow ) ); // Make sure logic matches
m_pTransport->SendEndToEndStatsMsg( k_EStatsReplyRequest_Immediate, usecNow, "E2ETimingOutKeepalive" );
AssertMsg( !m_statsEndToEnd.BNeedToSendPingImmediate( usecNow ), "SendEndToEndStatsMsg didn't do its job!" );
Assert( m_statsEndToEnd.m_usecInFlightReplyTimeout > 0 );
}
else
{
// Nothing we can do right now. Just check back in a little bit.
UpdateMinThinkTime( usecNow+20*1000 );
}
}
else
{
UpdateMinThinkTime( usecSendAggressivePing );
}
}
// Ordinary keepalive?
if ( m_statsEndToEnd.m_usecInFlightReplyTimeout == 0 )
{
// FIXME We really should be a lot better here with an adaptive keepalive time. If they have been
// sending us a steady stream of packets, we could expect it to continue at a high rate, so that we
// can begin to detect a dropped connection much more quickly. But if the connection is mostly idle, we want
// to make sure we use a relatively long keepalive.
SteamNetworkingMicroseconds usecSendKeepalive = m_statsEndToEnd.m_usecTimeLastRecv+k_usecKeepAliveInterval;
if ( usecNow >= usecSendKeepalive )
// Urgent keepalive because we are timing out?
SteamNetworkingMicroseconds usecStatsNextThinkTime = k_nThinkTime_Never;
EStatsReplyRequest eReplyRequested;
const char *pszStatsReason = m_statsEndToEnd.GetSendReasonOrUpdateNextThinkTime( usecNow, eReplyRequested, usecStatsNextThinkTime );
if ( pszStatsReason )
{
if ( bCanSendEndToEnd )
{
Assert( m_statsEndToEnd.BNeedToSendKeepalive( usecNow ) ); // Make sure logic matches
m_pTransport->SendEndToEndStatsMsg( k_EStatsReplyRequest_DelayedOK, usecNow, "E2EKeepalive" );
AssertMsg( !m_statsEndToEnd.BNeedToSendKeepalive( usecNow ), "SendEndToEndStatsMsg didn't do its job!" );
}
else
{
// Nothing we can do right now. Just check back in a little bit.
UpdateMinThinkTime( usecNow+20*1000 );
}
// Spew if we're dropping replies
if ( m_statsEndToEnd.m_nReplyTimeoutsSinceLastRecv == 1 )
SpewVerbose( "[%s] Reply timeout, last recv %.1fms ago. Sending keepalive.\n", GetDescription(), ( usecNow - m_statsEndToEnd.m_usecTimeLastRecv ) * 1e-3 );
else if ( m_statsEndToEnd.m_nReplyTimeoutsSinceLastRecv > 0 )
SpewMsg( "[%s] %d reply timeouts, last recv %.1fms ago. Sending keepalive.\n", GetDescription(), m_statsEndToEnd.m_nReplyTimeoutsSinceLastRecv, ( usecNow - m_statsEndToEnd.m_usecTimeLastRecv ) * 1e-3 );
// Send it
m_pTransport->SendEndToEndStatsMsg( eReplyRequested, usecNow, pszStatsReason );
// Re-calculate next think time
usecStatsNextThinkTime = k_nThinkTime_Never;
const char *pszStatsReason2 = m_statsEndToEnd.GetSendReasonOrUpdateNextThinkTime( usecNow, eReplyRequested, usecStatsNextThinkTime );
AssertMsg1( pszStatsReason2 == nullptr && usecStatsNextThinkTime > usecNow, "Stats sending didn't clear stats need to send reason %s!", pszStatsReason2 ? pszStatsReason2 : "??" );
}
else
{
// Not right now, but schedule a wakeup call to do it
UpdateMinThinkTime( usecSendKeepalive );
}
}
else
{
UpdateMinThinkTime( m_statsEndToEnd.m_usecInFlightReplyTimeout );
// Make sure we are scheduled to wake up the next time we need to take action
UpdateMinThinkTime( usecStatsNextThinkTime );
}
}
@@ -290,6 +290,16 @@ public:
/// (The wire protocol doesn't care about local states such as linger)
bool BStateIsConnectedForWirePurposes() const { return m_eConnectionState == k_ESteamNetworkingConnectionState_Connected || m_eConnectionState == k_ESteamNetworkingConnectionState_Linger; }
/// Return true if the connection is still "active" in some way.
bool BStateIsActive() const
{
return
m_eConnectionState == k_ESteamNetworkingConnectionState_Connecting
|| m_eConnectionState == k_ESteamNetworkingConnectionState_FindingRoute
|| m_eConnectionState == k_ESteamNetworkingConnectionState_Connected
|| m_eConnectionState == k_ESteamNetworkingConnectionState_Linger;
}
/// Reason connection ended
ESteamNetConnectionEnd GetConnectionEndReason() const { return m_eEndReason; }
const char *GetConnectionEndDebugString() const { return m_szEndDebug; }
@@ -465,16 +475,6 @@ public:
}
}
/// Check if we need to send stats or acks. If so, return a reason string
// FIXME - This needs to be refactored. There is some redundancy in the different
// transport code that uses it
const char *NeedToSendEndToEndStatsOrAcks( SteamNetworkingMicroseconds usecNow )
{
if ( m_receiverState.TimeWhenFlushAcks() <= usecNow )
return "SNPFlushAcks";
return m_statsEndToEnd.NeedToSend( usecNow );
}
inline const CMsgSteamDatagramSessionCryptInfoSigned &GetSignedCryptLocal() { return m_msgSignedCryptLocal; }
inline const CMsgSteamDatagramCertificateSigned &GetSignedCertLocal() { return m_msgSignedCertLocal; }
inline bool BCertHasIdentity() const { return m_bCertHasIdentity; }
@@ -601,23 +601,15 @@ void CConnectionTransportUDPBase::RecvStats( const CMsgSteamSockets_UDP_Stats &m
);
// Check if we need to reply, either now or later
if ( m_connection.BStateIsConnectedForWirePurposes() )
if ( m_connection.BStateIsActive() )
{
// Check for queuing outgoing acks
bool bImmediate = ( msgStatsIn.flags() & msgStatsIn.ACK_REQUEST_IMMEDIATE ) != 0;
if ( ( msgStatsIn.flags() & msgStatsIn.ACK_REQUEST_E2E ) || msgStatsIn.has_stats() )
{
bool bImmediate = ( msgStatsIn.flags() & msgStatsIn.ACK_REQUEST_IMMEDIATE ) != 0;
m_connection.QueueEndToEndAck( bImmediate, usecNow );
}
// Do we need to send an immediate reply?
const char *pszReason = m_connection.NeedToSendEndToEndStatsOrAcks( usecNow );
if ( pszReason )
{
// Send a stats message
SendStatsMsg( k_EStatsReplyRequest_NothingToSend, usecNow, pszReason );
}
}
}
@@ -674,7 +666,6 @@ void CConnectionTransportUDPBase::Received_Data( const uint8 *pPkt, int cbPkt, S
{
case k_ESteamNetworkingConnectionState_Dead:
case k_ESteamNetworkingConnectionState_None:
case k_ESteamNetworkingConnectionState_FindingRoute: // not used for raw UDP
default:
Assert( false );
return;
@@ -694,6 +685,7 @@ void CConnectionTransportUDPBase::Received_Data( const uint8 *pPkt, int cbPkt, S
case k_ESteamNetworkingConnectionState_Linger:
case k_ESteamNetworkingConnectionState_Connected:
case k_ESteamNetworkingConnectionState_FindingRoute: // not used for raw UDP, but might be used for derived class
// We'll process the chunk
break;
@@ -972,37 +964,26 @@ void UDPSendPacketContext_t::Populate( size_t cbHdrtReserve, EStatsReplyRequest
m_nFlags = nFlags;
// Need to send any connection stats stats?
if ( statsEndToEnd.BNeedToSendStats( m_usecNow ) )
m_nStatsNeed = statsEndToEnd.GetStatsSendNeed( m_usecNow );
if ( m_nStatsNeed & k_nSendStats_Due )
{
m_nStatsNeed = 2;
statsEndToEnd.PopulateMessage( *msg.mutable_stats(), m_usecNow );
statsEndToEnd.PopulateMessage( m_nStatsNeed, *msg.mutable_stats(), m_usecNow );
if ( nReadyToSendTracer > 0 )
nFlags |= msg.ACK_REQUEST_E2E;
SlamFlagsAndCalcSize();
CalcMaxEncryptedPayloadSize( cbHdrtReserve, &connection );
m_nFlags |= msg.ACK_REQUEST_E2E;
}
else
{
// Populate flags now, based on what is implied from what we HAVE to send
SlamFlagsAndCalcSize();
CalcMaxEncryptedPayloadSize( cbHdrtReserve, &connection );
// Would we like to try to send some additional stats, if there is room?
if ( statsEndToEnd.BReadyToSendStats( m_usecNow ) )
{
if ( nReadyToSendTracer > 0 )
nFlags |= msg.ACK_REQUEST_E2E;
statsEndToEnd.PopulateMessage( *msg.mutable_stats(), m_usecNow );
SlamFlagsAndCalcSize();
m_nStatsNeed = 1;
}
else
{
// No need to send any stats right now
m_nStatsNeed = 0;
}
// Populate flags now, based on what is implied from what we HAVE to send
SlamFlagsAndCalcSize();
CalcMaxEncryptedPayloadSize( cbHdrtReserve, &connection );
// Would we like to try to send some additional stats, if there is room?
if ( m_nStatsNeed & k_nSendStats_Ready )
{
if ( nReadyToSendTracer > 0 )
m_nFlags |= msg.ACK_REQUEST_E2E;
statsEndToEnd.PopulateMessage( m_nStatsNeed & k_nSendStats_Ready, *msg.mutable_stats(), m_usecNow );
SlamFlagsAndCalcSize();
}
}
@@ -1010,29 +991,42 @@ void UDPSendPacketContext_t::Trim( int cbHdrOutSpaceRemaining )
{
while ( m_cbTotalSize > cbHdrOutSpaceRemaining )
{
if ( msg.has_stats() )
if ( !msg.has_stats() )
{
AssertMsg( m_nStatsNeed == 1, "We didn't reserve enough space for stats!" );
if ( msg.stats().has_instantaneous() && msg.stats().has_lifetime() )
// Nothing left to clear!? We shouldn't get here!
AssertMsg( false, "Serialized stats message still won't fit, ever after clearing everything?" );
m_cbTotalSize = 0;
break;
}
if ( m_nStatsNeed & k_nSendStats_Instantanous_Ready )
{
msg.mutable_stats()->clear_instantaneous();
m_nStatsNeed &= ~k_nSendStats_Instantanous_Ready;
}
else if ( m_nStatsNeed & k_nSendStats_Lifetime_Ready )
{
msg.mutable_stats()->clear_lifetime();
m_nStatsNeed &= ~k_nSendStats_Lifetime_Ready;
}
else
{
AssertMsg( false, "We didn't reserve enough space for stats!" );
if ( m_nStatsNeed & k_nSendStats_Instantanous_Due )
{
// Trying to send both - clear instantaneous
msg.mutable_stats()->clear_instantaneous();
m_nStatsNeed &= ~k_nSendStats_Instantanous_Due;
}
else
{
// Trying to send just one or the other. Clear the whole container.
msg.clear_stats();
m_nStatsNeed = 0;
}
SlamFlagsAndCalcSize();
continue;
}
// Nothing left to clear!? We shouldn't get here!
AssertMsg( false, "Serialized stats message still won't fit, ever after clearing everything?" );
m_cbTotalSize = 0;
break;
if ( m_nStatsNeed == 0 )
msg.clear_stats();
SlamFlagsAndCalcSize();
}
}
@@ -1176,40 +1170,6 @@ void CConnectionTransportUDP::SendEndToEndConnectRequest( SteamNetworkingMicrose
m_connection.m_statsEndToEnd.TrackSentPingRequest( usecNow, false );
}
void CSteamNetworkConnectionUDP::ThinkConnection( SteamNetworkingMicroseconds usecNow )
{
// FIXME - We should refactor this, maybe promote this to the base class.
// There's really nothing specific to plain UDP transport here.
// Check if we have stats we need to flush out
if ( !m_statsEndToEnd.IsPassive() && m_pTransport )
{
// Do we need to send something immediately, for any reason?
const char *pszReason = NeedToSendEndToEndStatsOrAcks( usecNow );
if ( pszReason )
{
m_pTransport->SendEndToEndStatsMsg( k_EStatsReplyRequest_NothingToSend, usecNow, pszReason );
// Make sure that took care of what we needed!
Assert( !NeedToSendEndToEndStatsOrAcks( usecNow ) );
}
// Make sure we are scheduled to think the next time we need to
SteamNetworkingMicroseconds usecNextStatsThink = m_statsEndToEnd.GetNextThinkTime( usecNow );
if ( usecNextStatsThink <= usecNow )
{
AssertMsg( false, "We didn't send all the stats we needed to!" );
}
else
{
EnsureMinThinkTime( usecNextStatsThink );
}
}
}
bool CSteamNetworkConnectionUDP::BBeginAccept(
CSteamNetworkListenSocketDirectUDP *pParent,
const netadr_t &adrFrom,
@@ -197,7 +197,6 @@ public:
/// Implements CSteamNetworkConnectionBase
virtual EResult AcceptConnection( SteamNetworkingMicroseconds usecNow ) override;
virtual void ThinkConnection( SteamNetworkingMicroseconds usecNow ) override;
virtual void GetConnectionTypeDescription( ConnectionTypeDescription_t &szDescription ) const override;
virtual EUnsignedCert AllowRemoteUnsignedCert() override;
virtual EUnsignedCert AllowLocalUnsignedCert() override;
@@ -13,6 +13,7 @@
#include "percentile_generator.h"
#include "steamnetworking_stats.h"
#include "steamnetworkingsockets_internal.h"
#include "steamnetworkingsockets_thinker.h"
//#include <google/protobuf/repeated_field.h> // FIXME - should only need this!
#include <tier0/memdbgoff.h>
@@ -40,12 +41,10 @@ const SteamNetworkingMicroseconds k_usecLinkStatsMaxPingRequestInterval = 7 * k_
/// Client should send instantaneous connection quality stats
/// at approximately this interval
const SteamNetworkingMicroseconds k_usecLinkStatsInstantaneousReportMinInterval = 17 * k_nMillion;
const SteamNetworkingMicroseconds k_usecLinkStatsInstantaneousReportInterval = 20 * k_nMillion;
const SteamNetworkingMicroseconds k_usecLinkStatsInstantaneousReportMaxInterval = 30 * k_nMillion;
/// Client will report lifetime connection stats at approximately this interval
const SteamNetworkingMicroseconds k_usecLinkStatsLifetimeReportMinInterval = 102 * k_nMillion;
const SteamNetworkingMicroseconds k_usecLinkStatsLifetimeReportInterval = 120 * k_nMillion;
const SteamNetworkingMicroseconds k_usecLinkStatsLifetimeReportMaxInterval = 140 * k_nMillion;
@@ -384,6 +383,16 @@ private:
float m_flTokenDeficitFromFull;
};
// Bitmask returned by GetStatsSendNeed
constexpr int k_nSendStats_Instantanous_Due = 1;
constexpr int k_nSendStats_Instantanous_Ready = 2;
constexpr int k_nSendStats_Lifetime_Due = 4;
constexpr int k_nSendStats_Lifetime_Ready = 8;
constexpr int k_nSendStats_Instantanous = k_nSendStats_Instantanous_Due|k_nSendStats_Instantanous_Ready;
constexpr int k_nSendStats_Lifetime = k_nSendStats_Lifetime_Due|k_nSendStats_Lifetime_Ready;
constexpr int k_nSendStats_Due = k_nSendStats_Instantanous_Due|k_nSendStats_Lifetime_Due;
constexpr int k_nSendStats_Ready = k_nSendStats_Instantanous_Ready|k_nSendStats_Lifetime_Ready;
/// Base class used to handle link quality calculations.
///
@@ -570,10 +579,9 @@ struct LinkStatsTrackerBase
/// connectivity as well.
inline bool BNeedToSendPingImmediate( SteamNetworkingMicroseconds usecNow ) const
{
return
!m_bPassive
&& m_nReplyTimeoutsSinceLastRecv > 0 // We're timing out
&& m_usecLastSendPacketExpectingImmediateReply+k_usecAggressivePingInterval < usecNow; // we haven't just recently sent an aggressive ping.
if ( m_bPassive || m_nReplyTimeoutsSinceLastRecv == 0 )
return false;
return usecNow < m_usecLastSendPacketExpectingImmediateReply+k_usecAggressivePingInterval;
}
/// Check if we should send a keepalive ping. In this case we haven't heard from the peer in a while,
@@ -586,26 +594,10 @@ struct LinkStatsTrackerBase
&& m_usecTimeLastRecv + k_usecKeepAliveInterval < usecNow; // haven't heard from the peer recently
}
/// Check if we have data worth sending, if we have a good
/// opportunity (inline in a data packet) to do it.
inline bool BReadyToSendStats( SteamNetworkingMicroseconds usecNow )
{
bool bResult = false;
if ( m_pktNumInFlight == 0 && !m_bPassive )
{
if ( m_usecPeerAckedInstaneous + k_usecLinkStatsInstantaneousReportInterval < usecNow && BCheckHaveDataToSendInstantaneous( usecNow ) )
bResult = true ;
if ( m_usecPeerAckedLifetime + k_usecLinkStatsLifetimeReportInterval < usecNow && BCheckHaveDataToSendLifetime( usecNow ) )
bResult = true;
}
return bResult;
}
/// Fill out message with everything we'd like to send. We don't assume that we will
/// actually send it. (We might be looking for a good opportunity, and the data we want
/// to send doesn't fit.)
void PopulateMessage( CMsgSteamDatagramConnectionQuality &msg, SteamNetworkingMicroseconds usecNow );
void PopulateMessage( int nNeedFlags, CMsgSteamDatagramConnectionQuality &msg, SteamNetworkingMicroseconds usecNow );
void PopulateLifetimeMessage( CMsgSteamDatagramLinkLifetimeStats &msg );
/// Called when we send any message for which we expect some sort of reply. (But maybe not an ack.)
void TrackSentMessageExpectingReply( SteamNetworkingMicroseconds usecNow, bool bAllowDelayedReply );
@@ -693,8 +685,8 @@ struct LinkStatsTrackerBase
m_bInFlightInstantaneous = m_bInFlightLifetime = false;
}
/// Check if we really need to flush out stats now.
bool BNeedToSendStats( SteamNetworkingMicroseconds usecNow );
/// Get urgency level to send instantaneous/lifetime stats.
int GetStatsSendNeed( SteamNetworkingMicroseconds usecNow );
/// Describe this stats tracker, for debugging, asserts, etc
virtual std::string Describe() const = 0;
@@ -755,19 +747,16 @@ protected:
/// Are we in "passive" state? When we are "active", we expect that our peer is awake
/// and will reply to our messages, and that we should be actively sending our peer
/// connection quality statistics and keepalives. When we are passive, we still measure
/// statistics and can receive messages from the peer, and send acknowledgements as necessary.
/// statistics and can receive messages from the peer, and send acknowledgments as necessary.
/// but we will indicate that keepalives or stats need to be sent to the peer.
bool m_bPassive;
/// Called to switch the pasive state. (Should only be called on an actual state change.)
/// Called to switch the passive state. (Should only be called on an actual state change.)
void SetPassiveInternal( bool bFlag, SteamNetworkingMicroseconds usecNow );
/// Check if we really need to flush out stats now. Derived class should provide the reason strings.
/// (See the code.)
const char *NeedToSendStats( SteamNetworkingMicroseconds usecNow, const char *const arpszReasonStrings[4] );
/// Get time when we need to take action or think
SteamNetworkingMicroseconds GetNextThinkTimeInternal( SteamNetworkingMicroseconds usecNow ) const;
const char *InternalGetSendStatsReasonOrUpdateNextThinkTime( SteamNetworkingMicroseconds usecNow, const char *const arpszReasonStrings[4], SteamNetworkingMicroseconds &inOutNextThinkTime );
/// Called when we send a packet for which we expect a reply and
/// for which we expect to get latency info.
@@ -786,6 +775,36 @@ protected:
pThis->m_ping.ReceivedPing( nPingMS, usecNow );
}
inline bool BInternalNeedToSendPingImmediate( SteamNetworkingMicroseconds usecNow, SteamNetworkingMicroseconds &inOutNextThinkTime )
{
if ( m_nReplyTimeoutsSinceLastRecv == 0 )
return false;
SteamNetworkingMicroseconds usecUrgentPing = m_usecLastSendPacketExpectingImmediateReply+k_usecAggressivePingInterval;
if ( usecUrgentPing <= usecNow )
return true;
if ( usecUrgentPing < inOutNextThinkTime )
inOutNextThinkTime = usecUrgentPing;
return false;
}
inline bool BInternalNeedToSendKeepAlive( SteamNetworkingMicroseconds usecNow, SteamNetworkingMicroseconds &inOutNextThinkTime )
{
if ( m_usecInFlightReplyTimeout == 0 )
{
SteamNetworkingMicroseconds usecKeepAlive = m_usecTimeLastRecv + k_usecKeepAliveInterval;
if ( usecKeepAlive <= usecNow )
return true;
if ( usecKeepAlive < inOutNextThinkTime )
inOutNextThinkTime = usecKeepAlive;
}
else
{
if ( m_usecInFlightReplyTimeout < inOutNextThinkTime )
inOutNextThinkTime = m_usecInFlightReplyTimeout;
}
return false;
}
private:
// Number of lifetime sequenced packets received, and overall packets sent,
@@ -862,20 +881,32 @@ struct LinkStatsTrackerEndToEnd : public LinkStatsTrackerBase
/// Called when we get a speed sample
void UpdateSpeeds( int nTXSpeed, int nRXSpeed );
/// Do we need to send anything? Return the reason code, or NULL if
/// we don't need to send anything right now
inline const char *NeedToSend( SteamNetworkingMicroseconds usecNow )
/// Do we need to send any stats?
inline const char *GetSendReasonOrUpdateNextThinkTime( SteamNetworkingMicroseconds usecNow, EStatsReplyRequest &eReplyRequested, SteamNetworkingMicroseconds &inOutNextThinkTime )
{
if ( m_bPassive )
{
if ( m_usecInFlightReplyTimeout > 0 && m_usecInFlightReplyTimeout < inOutNextThinkTime )
inOutNextThinkTime = m_usecInFlightReplyTimeout;
eReplyRequested = k_EStatsReplyRequest_NothingToSend;
return nullptr;
}
// Connectivity check because we appear to be timing out?
if ( BNeedToSendPingImmediate( usecNow ) )
// Urgent ping?
if ( BInternalNeedToSendPingImmediate( usecNow, inOutNextThinkTime ) )
{
eReplyRequested = k_EStatsReplyRequest_Immediate;
return "E2EUrgentPing";
}
// Ordinary keepalive?
if ( BNeedToSendKeepalive( usecNow ) )
return "E2EKeepalive";
// Keepalive?
if ( BInternalNeedToSendKeepAlive( usecNow, inOutNextThinkTime ) )
{
eReplyRequested = k_EStatsReplyRequest_DelayedOK;
return "E2EKeepAlive";
}
// Stats?
// Connection stats?
static const char *arpszReasons[4] =
{
nullptr,
@@ -883,7 +914,15 @@ struct LinkStatsTrackerEndToEnd : public LinkStatsTrackerBase
"E2ELifetimeStats",
"E2EAllStats"
};
return LinkStatsTrackerBase::NeedToSendStats( usecNow, arpszReasons );
const char *pszReason = LinkStatsTrackerBase::InternalGetSendStatsReasonOrUpdateNextThinkTime( usecNow, arpszReasons, inOutNextThinkTime );
if ( pszReason )
{
eReplyRequested = k_EStatsReplyRequest_DelayedOK;
return pszReason;
}
eReplyRequested = k_EStatsReplyRequest_NothingToSend;
return nullptr;
}
/// Describe this stats tracker, for debugging, asserts, etc
@@ -903,21 +942,6 @@ protected:
}
}
inline SteamNetworkingMicroseconds GetNextThinkTimeInternal( SteamNetworkingMicroseconds usecNow ) const
{
SteamNetworkingMicroseconds usecResult = LinkStatsTrackerBase::GetNextThinkTimeInternal( usecNow );
if ( !m_bPassive )
{
if ( !m_usecInFlightReplyTimeout )
{
// Time when BNeedToSendKeepalive will return true
usecResult = std::min( usecResult, m_usecTimeLastRecv + k_usecKeepAliveInterval );
}
}
return usecResult;
}
private:
void UpdateSpeedInterval( SteamNetworkingMicroseconds usecNow );
@@ -942,7 +966,6 @@ struct LinkStatsTracker final : public TLinkStatsTracker
inline bool IsPassive() const { return TLinkStatsTracker::m_bPassive; }
inline void TrackSentMessageExpectingSeqNumAck( SteamNetworkingMicroseconds usecNow, bool bAllowDelayedReply ) { TLinkStatsTracker::TrackSentMessageExpectingSeqNumAckInternal( this, usecNow, bAllowDelayedReply ); }
inline void TrackSentPingRequest( SteamNetworkingMicroseconds usecNow, bool bAllowDelayedReply ) { TLinkStatsTracker::TrackSentPingRequestInternal( this, usecNow, bAllowDelayedReply ); }
inline SteamNetworkingMicroseconds GetNextThinkTime( SteamNetworkingMicroseconds usecNow ) const { return TLinkStatsTracker::GetNextThinkTimeInternal( usecNow ); }
inline void ReceivedPing( int nPingMS, SteamNetworkingMicroseconds usecNow ) { TLinkStatsTracker::ReceivedPingInternal( this, nPingMS, usecNow ); }
inline void InFlightReplyTimeout( SteamNetworkingMicroseconds usecNow ) { TLinkStatsTracker::InFlightReplyTimeoutInternal( this, usecNow ); }
@@ -987,6 +1010,15 @@ struct LinkStatsTracker final : public TLinkStatsTracker
return bResult;
}
// Shortcut when we know that we aren't going to send now, but we want to know when to wakeup and do so
inline SteamNetworkingMicroseconds GetNextThinkTime( SteamNetworkingMicroseconds usecNow )
{
SteamNetworkingMicroseconds usecNextThink = k_nThinkTime_Never;
EStatsReplyRequest eReplyRequested;
if ( TLinkStatsTracker::GetSendReasonOrUpdateNextThinkTime( usecNow, eReplyRequested, usecNextThink ) )
return k_nThinkTime_ASAP;
return usecNextThink;
}
};
@@ -467,7 +467,7 @@ bool LinkStatsTrackerBase::BCheckHaveDataToSendInstantaneous( SteamNetworkingMic
// Calculate threshold based on how much time has elapsed and a very low packet rate
int64 usecElapsed = usecNow - m_usecPeerAckedInstaneous;
Assert( usecElapsed >= k_usecLinkStatsInstantaneousReportMinInterval ); // don't call this unless you know it's been long enough!
Assert( usecElapsed >= k_usecLinkStatsInstantaneousReportInterval ); // don't call this unless you know it's been long enough!
int nThreshold = usecElapsed / k_usecActiveConnectionSendInterval;
// Has there been any traffic worth reporting on in this interval?
@@ -500,75 +500,85 @@ bool LinkStatsTrackerBase::BCheckHaveDataToSendLifetime( SteamNetworkingMicrosec
return false;
}
bool LinkStatsTrackerBase::BNeedToSendStats( SteamNetworkingMicroseconds usecNow )
int LinkStatsTrackerBase::GetStatsSendNeed( SteamNetworkingMicroseconds usecNow )
{
// Message already in flight?
if ( m_pktNumInFlight != 0 || m_bPassive )
return false;
bool bNeedToSendInstantaneous = ( m_usecPeerAckedInstaneous + k_usecLinkStatsInstantaneousReportMaxInterval < usecNow ) && BCheckHaveDataToSendInstantaneous( usecNow );
bool bNeedToSendLifetime = ( m_usecPeerAckedLifetime + k_usecLinkStatsLifetimeReportMaxInterval < usecNow ) && BCheckHaveDataToSendLifetime( usecNow );
return bNeedToSendInstantaneous || bNeedToSendLifetime;
}
int nResult = 0;
const char *LinkStatsTrackerBase::NeedToSendStats( SteamNetworkingMicroseconds usecNow, const char *const arpszReasonStrings[4] )
{
// Message already in flight?
if ( m_pktNumInFlight != 0 || m_bPassive )
return nullptr;
int n = 0;
if ( m_usecPeerAckedInstaneous + k_usecLinkStatsInstantaneousReportMaxInterval < usecNow && BCheckHaveDataToSendInstantaneous( usecNow ) )
n |= 1;
if ( m_usecPeerAckedLifetime + k_usecLinkStatsLifetimeReportMaxInterval < usecNow && BCheckHaveDataToSendLifetime( usecNow ) )
n |= 2;
return arpszReasonStrings[n];
}
SteamNetworkingMicroseconds LinkStatsTrackerBase::GetNextThinkTimeInternal( SteamNetworkingMicroseconds usecNow ) const
{
SteamNetworkingMicroseconds usecResult = INT64_MAX;
if ( !m_bPassive )
if ( m_pktNumInFlight == 0 && !m_bPassive )
{
// Expecting a reply?
if ( m_usecInFlightReplyTimeout )
if ( m_usecPeerAckedInstaneous + k_usecLinkStatsInstantaneousReportInterval < usecNow && BCheckHaveDataToSendInstantaneous( usecNow ) )
{
usecResult = std::min( usecResult, m_usecInFlightReplyTimeout );
if ( m_usecPeerAckedInstaneous + k_usecLinkStatsInstantaneousReportMaxInterval < usecNow )
nResult |= k_nSendStats_Instantanous_Due;
else
nResult |= k_nSendStats_Instantanous_Ready;
}
// Time when BNeedToSendPingImmediate will return true
if ( m_nReplyTimeoutsSinceLastRecv > 0 )
usecResult = std::min( usecResult, m_usecLastSendPacketExpectingImmediateReply+k_usecAggressivePingInterval );
// Time when we need to flush stats
if ( m_pktNumInFlight == 0 )
if ( m_usecPeerAckedLifetime + k_usecLinkStatsLifetimeReportInterval < usecNow && BCheckHaveDataToSendLifetime( usecNow ) )
{
usecResult = std::min( usecResult, m_usecPeerAckedInstaneous + k_usecLinkStatsInstantaneousReportMaxInterval );
usecResult = std::min( usecResult, m_usecPeerAckedLifetime + k_usecLinkStatsLifetimeReportMaxInterval );
if ( m_usecPeerAckedInstaneous + k_usecLinkStatsLifetimeReportMaxInterval < usecNow )
nResult |= k_nSendStats_Lifetime_Due;
else
nResult |= k_nSendStats_Lifetime_Ready;
}
}
return usecResult;
return nResult;
}
void LinkStatsTrackerBase::PopulateMessage( CMsgSteamDatagramConnectionQuality &msg, SteamNetworkingMicroseconds usecNow )
const char *LinkStatsTrackerBase::InternalGetSendStatsReasonOrUpdateNextThinkTime( SteamNetworkingMicroseconds usecNow, const char *const arpszReasonStrings[4], SteamNetworkingMicroseconds &inOutNextThinkTime )
{
if ( m_pktNumInFlight == 0 && !m_bPassive )
if ( m_bPassive )
return nullptr;
if ( m_usecInFlightReplyTimeout > 0 && m_usecInFlightReplyTimeout < inOutNextThinkTime )
inOutNextThinkTime = m_usecInFlightReplyTimeout;
// Message already in flight?
if ( m_pktNumInFlight )
return nullptr;
int n = 0;
if ( m_usecPeerAckedInstaneous + k_usecLinkStatsInstantaneousReportMaxInterval < usecNow && BCheckHaveDataToSendInstantaneous( usecNow ) )
{
n |= 1;
}
else
{
SteamNetworkingMicroseconds usecNextCheck = m_usecPeerAckedInstaneous + k_usecLinkStatsInstantaneousReportMaxInterval;
if ( usecNextCheck < inOutNextThinkTime )
inOutNextThinkTime = usecNextCheck;
}
if ( m_usecPeerAckedLifetime + k_usecLinkStatsLifetimeReportMaxInterval < usecNow && BCheckHaveDataToSendLifetime( usecNow ) )
{
n |= 2;
}
else
{
SteamNetworkingMicroseconds usecNextCheck = m_usecPeerAckedLifetime + k_usecLinkStatsLifetimeReportMaxInterval;
if ( usecNextCheck < inOutNextThinkTime )
inOutNextThinkTime = usecNextCheck;
}
return arpszReasonStrings[n];
}
// Ready to send instantaneous stats?
if ( m_usecPeerAckedInstaneous + k_usecLinkStatsInstantaneousReportMinInterval < usecNow && BCheckHaveDataToSendInstantaneous( usecNow ) )
{
// !KLUDGE! Go through public struct as intermediary to keep code simple.
SteamDatagramLinkInstantaneousStats sInstant;
GetInstantaneousStats( sInstant );
LinkStatsInstantaneousStructToMsg( sInstant, *msg.mutable_instantaneous() );
}
void LinkStatsTrackerBase::PopulateMessage( int nNeedFlags, CMsgSteamDatagramConnectionQuality &msg, SteamNetworkingMicroseconds usecNow )
{
Assert( m_pktNumInFlight == 0 && !m_bPassive );
// Ready to send lifetime stats?
if ( m_usecPeerAckedLifetime + k_usecLinkStatsLifetimeReportMinInterval < usecNow && BCheckHaveDataToSendLifetime( usecNow ) )
{
PopulateLifetimeMessage( *msg.mutable_lifetime() );
}
// Ready to send instantaneous stats?
if ( nNeedFlags & k_nSendStats_Instantanous )
{
// !KLUDGE! Go through public struct as intermediary to keep code simple.
SteamDatagramLinkInstantaneousStats sInstant;
GetInstantaneousStats( sInstant );
LinkStatsInstantaneousStructToMsg( sInstant, *msg.mutable_instantaneous() );
}
// Ready to send lifetime stats?
if ( nNeedFlags & k_nSendStats_Lifetime )
{
PopulateLifetimeMessage( *msg.mutable_lifetime() );
}
}