From 67013852d081ae98fd624b97aaadbf703fda0cb2 Mon Sep 17 00:00:00 2001 From: orignal Date: Tue, 3 Mar 2026 18:59:44 -0500 Subject: [PATCH] limit receive queue size --- libi2pd/Streaming.cpp | 37 +++++++++++++++++++++++++++++++------ libi2pd/Streaming.h | 3 ++- 2 files changed, 33 insertions(+), 7 deletions(-) diff --git a/libi2pd/Streaming.cpp b/libi2pd/Streaming.cpp index 068283b5..ac18cabf 100644 --- a/libi2pd/Streaming.cpp +++ b/libi2pd/Streaming.cpp @@ -89,7 +89,7 @@ namespace stream m_LastConfirmedReceivedSequenceNumber (0), // for limit inbound speed m_Status (eStreamStatusNew), m_IsIncoming (false), m_IsAckSendScheduled (false), m_IsNAcked (false), m_IsFirstACK (false), m_IsResendNeeded (false), m_IsFirstRttSample (false), m_IsSendTime (true), - m_IsWinDropped (true), m_IsChoking2 (false), m_IsClientChoked (false), m_IsClientChoked2 (false), + m_IsWinDropped (true), m_IsChoking2 (false), m_IsChoking3 (false), m_IsClientChoked (false), m_IsClientChoked2 (false), m_IsTimeOutResend (false), m_IsImmediateAckRequested (false), m_IsRemoteLeaseChangeInProgress (false), m_IsBufferEmpty (false), m_IsJavaClient (false), m_DontSign (local.GetOwner ()->IsStreamingDontSign ()), m_LocalDestination (local), m_RemoteLeaseSet (remote), m_ReceiveTimer (m_Service), @@ -123,7 +123,7 @@ namespace stream m_LastConfirmedReceivedSequenceNumber (0), // for limit inbound speed m_Status (eStreamStatusNew), m_IsIncoming (true), m_IsAckSendScheduled (false), m_IsNAcked (false), m_IsFirstACK (false), m_IsResendNeeded (false), m_IsFirstRttSample (false), m_IsSendTime (true), - m_IsWinDropped (true), m_IsChoking2 (false), m_IsClientChoked (false), m_IsClientChoked2 (false), + m_IsWinDropped (true), m_IsChoking2 (false), m_IsChoking3 (false), m_IsClientChoked (false), m_IsClientChoked2 (false), m_IsTimeOutResend (false), m_IsImmediateAckRequested (false), m_IsRemoteLeaseChangeInProgress (false), m_IsBufferEmpty (false), m_IsJavaClient (false), m_DontSign (local.GetOwner ()->IsStreamingDontSign ()), m_LocalDestination (local),m_ReceiveTimer (m_Service), m_SendTimer (m_Service), @@ -249,6 +249,21 @@ namespace stream } LogPrint (eLogDebug, "Streaming: Received seqn=", receivedSeqn, " on sSID=", m_SendStreamID); + if (m_ReceiveQueue.size () > m_MaxWindowSize*3) + { + LogPrint (eLogDebug, "Streaming: ReceiveQueue is full, delete packet"); + m_LocalDestination.DeletePacket (packet); + m_IsChoking3 = true; + if (!m_IsAckSendScheduled) + { + SendQuickAck (); + auto ackTimeout = m_RTT/10; + if (ackTimeout > m_AckDelay) ackTimeout = m_AckDelay; + ScheduleAck (ackTimeout); + } + return; + } + if (receivedSeqn == m_LastReceivedSequenceNumber + 1) { // we have received next in sequence message @@ -434,7 +449,14 @@ namespace stream } if (delayRequested >= DELAY_CHOKING) { - if (delayRequested == 65535) + if (delayRequested == DELAY_CHOKING_3) + { + m_NumResendAttempts = 0; + m_WindowDropTargetSize = MIN_WINDOW_SIZE; + m_IsClientChoked = true; + ResetWindowSize (); + } + else if (delayRequested == DELAY_CHOKING_2) { m_IsClientChoked2 = true; m_DropWindowDelaySequenceNumber = m_SequenceNumber-1; @@ -1188,7 +1210,7 @@ namespace stream htobe32buf (packet + size, lastReceivedSeqn); size += 4; // ack Through uint8_t numNacks = 0; - bool choking = m_IsChoking2; + bool choking = (m_IsChoking2 || m_IsChoking3); if (lastReceivedSeqn > m_LastReceivedSequenceNumber) { // fill NACKs @@ -1238,7 +1260,9 @@ namespace stream if (choking || requestImmediateAck) { htobe16buf (packet + size, 2); // 2 bytes delay interval - if (m_IsChoking2) + if (m_IsChoking3) + htobe16buf (packet + size + 2, DELAY_CHOKING_3); // set choking3 + else if (m_IsChoking2) htobe16buf (packet + size + 2, DELAY_CHOKING_2); // set choking2 else htobe16buf (packet + size + 2, choking ? DELAY_CHOKING : 0); // set choking or immediate ack interval @@ -1258,6 +1282,7 @@ namespace stream m_LastACKSendTime = ts; // for limit inbound speed m_LastConfirmedReceivedSequenceNumber = lastReceivedSeqn; // for limit inbound speed m_IsChoking2 = false; + m_IsChoking3 = false; LogPrint (eLogDebug, "Streaming: Quick Ack sent. ", (int)numNacks, " NACKs"); } @@ -1771,7 +1796,7 @@ namespace stream UpdatePacingTime (); } } - else if (m_IsTimeOutResend) + else if (m_IsTimeOutResend && m_NumResendAttempts > 1) { m_RTO = INITIAL_RTO; // drop RTO to initial upon tunnels pair change m_WindowDropTargetSize = INITIAL_WINDOW_SIZE; diff --git a/libi2pd/Streaming.h b/libi2pd/Streaming.h index 49bfd28d..176e9ecb 100644 --- a/libi2pd/Streaming.h +++ b/libi2pd/Streaming.h @@ -78,6 +78,7 @@ namespace stream const uint16_t DELAY_CHOKING = 60000; // in milliseconds const uint16_t DELAY_CHOKING_JAVA = 61000; // in milliseconds const uint16_t DELAY_CHOKING_2 = 65535; // in milliseconds + const uint16_t DELAY_CHOKING_3 = 65534; // in milliseconds const uint64_t SEND_INTERVAL = 10000; // in microseconds const uint64_t SEND_INTERVAL_VARIANCE = 2000; // in microseconds const uint64_t REQUEST_IMMEDIATE_ACK_INTERVAL = 7500; // in milliseconds @@ -274,7 +275,7 @@ namespace stream int32_t m_LastConfirmedReceivedSequenceNumber; // for limit inbound speed StreamStatus m_Status; bool m_IsIncoming, m_IsAckSendScheduled, m_IsNAcked, m_IsFirstACK, m_IsResendNeeded, - m_IsFirstRttSample, m_IsSendTime, m_IsWinDropped, m_IsChoking2, m_IsClientChoked, + m_IsFirstRttSample, m_IsSendTime, m_IsWinDropped, m_IsChoking2, m_IsChoking3, m_IsClientChoked, m_IsClientChoked2, m_IsTimeOutResend, m_IsImmediateAckRequested, m_IsRemoteLeaseChangeInProgress, m_IsBufferEmpty, m_IsJavaClient, m_DontSign; StreamingDestination& m_LocalDestination;