From 9dc99b9ecc1056b4a93dbdf5d254f4b153af3f64 Mon Sep 17 00:00:00 2001 From: orignal Date: Mon, 3 Aug 2026 17:35:37 -0400 Subject: [PATCH] multiple requests at the time --- libi2pd_client/Torrents.cpp | 59 ++++++++++++++++++------------------- libi2pd_client/Torrents.h | 11 +++---- 2 files changed, 34 insertions(+), 36 deletions(-) diff --git a/libi2pd_client/Torrents.cpp b/libi2pd_client/Torrents.cpp index a6d349e6..849906c6 100644 --- a/libi2pd_client/Torrents.cpp +++ b/libi2pd_client/Torrents.cpp @@ -176,8 +176,11 @@ namespace torrents std::fill_n (m_Blocks->begin () + startBlock, numBlocks, BlockStatus::Available); if (std::find_if (m_Blocks->begin (), m_Blocks->end (), [](BlockStatus status) { return status != BlockStatus::Available; }) == m_Blocks->end ()) + { // all blocks are available + LogPrint (eLogDebug, "Torrents: piece complete"); m_Blocks = nullptr; + } } } @@ -191,7 +194,6 @@ namespace torrents f.write ((const char *)m_Data, m_Size); delete[] m_Data; m_Data = nullptr; m_Blocks = nullptr; - m_Connections.clear (); } } @@ -254,16 +256,6 @@ namespace torrents it = BlockStatus::Missing; } - void Piece::AddConnection (std::shared_ptr connection) - { - m_Connections.emplace_back (connection); - } - - void Piece::RemoveConnection (std::shared_ptr connection) - { - m_Connections.remove_if ([connection](std::weak_ptr c) { return c.lock () == connection; }); - } - Torrent::Torrent (std::string_view buf): m_Length (0), m_PieceLength (0), m_Interval (MIN_TRACKER_REQUESTS_INTERVAL), m_NextTrackerRequestTime (0) @@ -411,7 +403,7 @@ namespace torrents uint32_t ind = 0; for (auto& it: m_Pieces) { - if (!it.IsComplete ()) + if (!it.IsComplete () && conn->IsPieceAvailable (ind)) { auto [offset, len] = it.GetNextBlockToRequest (); if (len > 0) @@ -433,7 +425,7 @@ namespace torrents PeerConnection::PeerConnection (i2p::client::I2PService * owner, std::shared_ptr stream): i2p::client::I2PServiceHandler (owner), m_Stream (stream), m_ReceiveBufferOffset (0), m_IsHandshakeSent (false), m_IsEstablished (false), m_IsChoked (true), - m_LastReceiveTime (0), m_LastSendTime (0) + m_LastReceiveTime (0), m_LastSendTime (0), m_NumRequests (0) { } @@ -464,6 +456,12 @@ namespace torrents return static_cast(GetOwner ()); } + bool PeerConnection::IsPieceAvailable (size_t ind) const + { + if (ind >= m_RemoteBitfield.size ()) return false; + return m_RemoteBitfield.test (ind); + } + void PeerConnection::WriteToStream (const uint8_t * buf, size_t len) { LogPrint (eLogDebug, "Torrents: Sending ", len, " bytes"); @@ -603,7 +601,7 @@ namespace torrents break; case eMessageTypeUnchoke: m_IsChoked = false; - RequestNextBlock (); // TODO: remove later + RequestNextBlocks (); break; case eMessageTypeInterested: LogPrint (eLogInfo, "Torrents: Interested message is not implemented"); @@ -705,10 +703,7 @@ namespace torrents { if (idx >= numPieces) break; if (buf[i] & bit) - { m_RemoteBitfield.set (idx); - m_Torrent->GetPiece (idx).AddConnection (shared_from_this ()); - } bit >>= 1; idx++; } @@ -731,12 +726,6 @@ namespace torrents size_t numPieces = m_Torrent->GetNumPieces (); m_RemoteBitfield.resize (numPieces); m_RemoteBitfield.set (); - for (size_t i = 0; i < numPieces; i++) - { - Piece& piece = m_Torrent->GetPiece (i); - if (!piece.IsComplete ()) - piece.AddConnection (shared_from_this ()); - } } void PeerConnection::HandleHaveNoneMsg () @@ -745,12 +734,6 @@ namespace torrents size_t numPieces = m_Torrent->GetNumPieces (); m_RemoteBitfield.resize (numPieces); m_RemoteBitfield.reset (); - for (size_t i = 0; i < numPieces; i++) - { - Piece& piece = m_Torrent->GetPiece (i); - if (!piece.IsComplete ()) - piece.RemoveConnection (shared_from_this ()); - } } void PeerConnection::HandlePieceMsg (const uint8_t * buf, size_t len) @@ -774,6 +757,8 @@ namespace torrents }); } } + if (m_NumRequests > 0) m_NumRequests--; + RequestNextBlocks (); } void PeerConnection::SendPieceMsg (uint32_t index, uint32_t offset, const uint8_t * data, size_t len) @@ -851,12 +836,24 @@ namespace torrents WriteToStream (buf, INTERESTED_MSG_LENGTH); } - void PeerConnection::RequestNextBlock () + bool PeerConnection::RequestNextBlock () { - if (!m_Torrent) return; + if (!m_Torrent) return false; auto [index, offset, len] = m_Torrent->GetNextBlockToRequest (shared_from_this ()); if (len > 0) + { SendRequestMsg (index, offset, len); + m_NumRequests++; + } + return len > 0; + } + + void PeerConnection::RequestNextBlocks () + { + while (m_NumRequests < MAX_NUM_REQUESTS) + { + if (!RequestNextBlock ()) break; + } } TorrentsTunnel::TorrentsTunnel (std::shared_ptr localDestination, diff --git a/libi2pd_client/Torrents.h b/libi2pd_client/Torrents.h index 7f87837d..5dc45008 100644 --- a/libi2pd_client/Torrents.h +++ b/libi2pd_client/Torrents.h @@ -44,6 +44,7 @@ namespace torrents constexpr int PEER_CONNECTION_MAX_IDLE = 3600; // in seconds constexpr int PEER_KEEP_ALIVE_INTERVAL = 120; // in seconds constexpr int PEER_KEEP_ALIVE_CHECK_TIMEOUT = 15; // in seconds + constexpr size_t MAX_NUM_REQUESTS = 8; constexpr size_t HANDSHAKE_MSG_LENGTH = 68; constexpr size_t INTERESTED_MSG_LENGTH = 5; @@ -91,9 +92,6 @@ namespace torrents std::pair GetNextBlockToRequest (); // return (offset, len) of next buffer, len = 0 if no next buffer void ClearAllRequests (); - void AddConnection (std::shared_ptr connection); - void RemoveConnection (std::shared_ptr connection); - private: bool IsAvailable (int block) const; @@ -104,7 +102,6 @@ namespace torrents size_t m_Size; uint8_t * m_Data, m_Hash[SHA_DIGEST_LENGTH]; std::unique_ptr > m_Blocks; - std::list > m_Connections; // for incomplete pieces only }; class Torrent final @@ -164,6 +161,8 @@ namespace torrents void ReceiveHandshake (); void CheckKeepAlive (uint64_t ts); + bool IsPieceAvailable (size_t ind) const; + private: void Terminate (); @@ -189,7 +188,8 @@ namespace torrents void SendRequestMsg (uint32_t index, uint32_t offset, uint32_t len); void SendInterestedMsg (); - void RequestNextBlock (); + bool RequestNextBlock (); + void RequestNextBlocks (); private: @@ -201,6 +201,7 @@ namespace torrents boost::dynamic_bitset<> m_RemoteBitfield; bool m_IsHandshakeSent, m_IsEstablished, m_IsChoked; uint64_t m_LastReceiveTime, m_LastSendTime; // monotonic seconds + size_t m_NumRequests; }; class TorrentsTunnel final: public i2p::client::I2PService