From 705bc0e4bbbb891cf797dfb8e5595efaefeb1e31 Mon Sep 17 00:00:00 2001 From: Fletcher Dunn Date: Sun, 17 May 2020 12:01:06 -0700 Subject: [PATCH] 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. --- .../steamnetworkingsockets_connections.cpp | 81 +++------- .../steamnetworkingsockets_connections.h | 20 +-- .../clientlib/steamnetworkingsockets_udp.cpp | 130 ++++++---------- .../clientlib/steamnetworkingsockets_udp.h | 1 - .../steamnetworking_statsutils.h | 146 +++++++++++------- .../steamnetworkingsockets_stats.cpp | 116 +++++++------- 6 files changed, 231 insertions(+), 263 deletions(-) diff --git a/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_connections.cpp b/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_connections.cpp index 9d1cc26..0d37c4b 100644 --- a/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_connections.cpp +++ b/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_connections.cpp @@ -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 ); } } diff --git a/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_connections.h b/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_connections.h index cfe8298..58cd044 100644 --- a/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_connections.h +++ b/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_connections.h @@ -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; } diff --git a/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_udp.cpp b/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_udp.cpp index 93cd3a3..0832818 100644 --- a/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_udp.cpp +++ b/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_udp.cpp @@ -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, diff --git a/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_udp.h b/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_udp.h index 71612a1..c36716f 100644 --- a/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_udp.h +++ b/src/steamnetworkingsockets/clientlib/steamnetworkingsockets_udp.h @@ -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; diff --git a/src/steamnetworkingsockets/steamnetworking_statsutils.h b/src/steamnetworkingsockets/steamnetworking_statsutils.h index 22dd775..e5af032 100644 --- a/src/steamnetworkingsockets/steamnetworking_statsutils.h +++ b/src/steamnetworkingsockets/steamnetworking_statsutils.h @@ -13,6 +13,7 @@ #include "percentile_generator.h" #include "steamnetworking_stats.h" #include "steamnetworkingsockets_internal.h" +#include "steamnetworkingsockets_thinker.h" //#include // FIXME - should only need this! #include @@ -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; + } }; diff --git a/src/steamnetworkingsockets/steamnetworkingsockets_stats.cpp b/src/steamnetworkingsockets/steamnetworkingsockets_stats.cpp index b9cd9cd..8b4f5bd 100644 --- a/src/steamnetworkingsockets/steamnetworkingsockets_stats.cpp +++ b/src/steamnetworkingsockets/steamnetworkingsockets_stats.cpp @@ -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() ); } }