diff --git a/libi2pd_client/Torrents.cpp b/libi2pd_client/Torrents.cpp index 1f669d1b..6115e532 100644 --- a/libi2pd_client/Torrents.cpp +++ b/libi2pd_client/Torrents.cpp @@ -446,7 +446,9 @@ namespace torrents Torrent::Torrent (): m_AnnounceTrackerID (-1), m_Length (0), m_PieceLength (0), m_IsComplete (false), m_IsStopped (false), m_IsSingleFile (true), m_Uploaded (0), m_Downloaded (0), - m_NextUpdateStatusTime (0), m_NextReconnectTime (0), m_Error (eTorrentErrorNoError) + m_NextUpdateStatusTime (0), m_NextReconnectTime (0), + m_NextDHTUpdateTime (i2p::util::GetMonotonicSeconds () + DHT_TORRENT_INITIAL_UPDATE_INTERVAL), + m_Error (eTorrentErrorNoError) { } @@ -2180,21 +2182,21 @@ namespace torrents void PeerConnection::SendExtendedMsg (uint8_t extendedMsgID, std::string_view payload, std::string_view data) { - const static std::string messages = GetTorrentsTunnel ()->SupportsDHT () ? - CreateDictionary ({ - { EXTENSION_NAME_I2P_DHT, CreateInteger (EXTENSION_MSGID_I2P_DHT) }, - { EXTENSION_NAME_I2P_PEX, CreateInteger (EXTENSION_MSGID_I2P_PEX) }, - { EXTENSION_NAME_UT_METADATA, CreateInteger (EXTENSION_MSGID_UT_METADATA) } - }): - CreateDictionary ({ - { EXTENSION_NAME_I2P_PEX, CreateInteger (EXTENSION_MSGID_I2P_PEX) }, - { EXTENSION_NAME_UT_METADATA, CreateInteger (EXTENSION_MSGID_UT_METADATA) } - }); std::string str; if (!extendedMsgID) // handshake { str = CreateDictionary ({ - { "m", messages }, + { "m", GetTorrentsTunnel ()->SupportsDHT () ? + CreateDictionary ({ + { EXTENSION_NAME_I2P_DHT, CreateInteger (EXTENSION_MSGID_I2P_DHT) }, + { EXTENSION_NAME_I2P_PEX, CreateInteger (EXTENSION_MSGID_I2P_PEX) }, + { EXTENSION_NAME_UT_METADATA, CreateInteger (EXTENSION_MSGID_UT_METADATA) } + }): + CreateDictionary ({ + { EXTENSION_NAME_I2P_PEX, CreateInteger (EXTENSION_MSGID_I2P_PEX) }, + { EXTENSION_NAME_UT_METADATA, CreateInteger (EXTENSION_MSGID_UT_METADATA) } + }) + }, { "metadata_size", CreateInteger (m_Torrent->GetInfo ().size ()) }, { "reqq", CreateInteger (MAX_INCOMING_REQUESTS_QUEUE_SIZE) }, { "v", CreateByteString ("i2pd") } diff --git a/libi2pd_client/Torrents.h b/libi2pd_client/Torrents.h index df2b9c21..648c9d4a 100644 --- a/libi2pd_client/Torrents.h +++ b/libi2pd_client/Torrents.h @@ -68,6 +68,9 @@ namespace torrents constexpr uint16_t TORRENT_PORT = 6881; // not used by required by protocol constexpr int MIN_TRACKER_REQUESTS_INTERVAL = 15*1000; // in milliseconds constexpr int MAX_TRACKER_REQUESTS_INTERVAL = 24*3600*1000; // in milliseconds + constexpr int DHT_TORRENT_UPDATE_INTERVAL = 150; // in seconds + constexpr int DHT_TORRENT_UPDATE_INTERVAL_VARIANCE = 30; // in seconds + constexpr int DHT_TORRENT_INITIAL_UPDATE_INTERVAL = 100; // in seconds constexpr size_t PEER_CONNECTION_RECEIVE_BUFFER_SIZE = 65535; constexpr int PEER_CONNECTION_MAX_IDLE = 3600; // in seconds constexpr int PEER_KEEP_ALIVE_TIMEOUT = 120; // in seconds @@ -307,6 +310,8 @@ namespace torrents void SetNextUpdateStatusTime (uint64_t nextUpdateStatusTime) { m_NextUpdateStatusTime = nextUpdateStatusTime; } uint64_t GetNextReconnectTime () { return m_NextReconnectTime; } void SetNextReconnectTime (uint64_t nextReconnectTime) { m_NextReconnectTime = nextReconnectTime; } + uint64_t GetNextDHTUpdateTime () { return m_NextDHTUpdateTime; } + void SetNextDHTUpdateTime (uint64_t nextDHTUpdateTime) { m_NextDHTUpdateTime = nextDHTUpdateTime; } bool AddConnection (std::shared_ptr conn); void RemoveConnection (std::shared_ptr conn); std::list > GetConnections (); @@ -360,7 +365,7 @@ namespace torrents bool m_IsComplete, m_IsStopped, m_IsSingleFile; std::list > m_Files; size_t m_Uploaded, m_Downloaded; - uint64_t m_NextUpdateStatusTime, m_NextReconnectTime; // in monotonic seconds + uint64_t m_NextUpdateStatusTime, m_NextReconnectTime, m_NextDHTUpdateTime; // in monotonic seconds TorrentError m_Error; }; diff --git a/libi2pd_client/TorrentsDHT.cpp b/libi2pd_client/TorrentsDHT.cpp index fc1f7643..e15f36db 100644 --- a/libi2pd_client/TorrentsDHT.cpp +++ b/libi2pd_client/TorrentsDHT.cpp @@ -230,6 +230,15 @@ namespace torrents return true; } + void RoutingTable::RemoveNode (const NodeID& id) + { + auto bucket = FindBucket (id); + if (!bucket) return; + bucket->nodes.erase (id); + if (bucket->nodes.empty () && bucket != m_Buckets) + RemoveEmptyBuckets (); + } + std::list > RoutingTable::FindClosestNodes (const Torrent::InfoHash& infoHash, size_t num) const { std::list > ret; @@ -270,23 +279,25 @@ namespace torrents if (!bucket->IsFull () && !bucket->nodes.empty ()) { auto randomID = bucket->GetRandomID (rng); - auto it = bucket->nodes.begin (); - while (it != bucket->nodes.end () && randomID > it->first) it++; - if (it != bucket->nodes.end ()) - ret.emplace_back (std::make_pair (randomID, it->first)); - else - ret.emplace_back (std::make_pair (randomID, bucket->nodes.rbegin ()->first)); + auto closestNodeID = FindClosestNodeInBucket (randomID); + if (closestNodeID) + ret.emplace_back (std::make_pair (randomID, *closestNodeID)); } bucket = bucket->next; } return ret; } - NodeID RoutingTable::FindClosestNodeInBucket (const NodeID& target) const + std::optional RoutingTable::FindClosestNodeInBucket (const NodeID& target) const { auto bucket = FindBucket (target); - if (!bucket || bucket->nodes.empty ()) return m_OurNode; - return bucket->nodes.begin ()->first; + if (!bucket || bucket->nodes.empty ()) return {}; + auto it = bucket->nodes.begin (); + while (it != bucket->nodes.end () && target > it->first) it++; + if (it != bucket->nodes.end ()) + return it->first; + else + return bucket->nodes.rbegin ()->first; } DHTTorrent::DHTTorrent (): @@ -500,7 +511,7 @@ namespace torrents } m_RoutingTable->RemoveEmptyBuckets (); if (numLoaded > 0) - LogPrint (eLogInfo, "TorrentsDHT: ", numLoaded, " DHT nodes loaded"); + LogPrint (eLogInfo, "TorrentsDHT: ", numLoaded, " DHT nodes loaded to ", m_RoutingTable->GetNumBuckets (), " buckets"); } } @@ -728,9 +739,17 @@ namespace torrents LogPrint (eLogDebug, "TorrentsDHT: Find node query received"); if (m_RoutingTable) { - auto it = m_Nodes.find (m_RoutingTable->FindClosestNodeInBucket (target)); - if (it != m_Nodes.end ()) - SendFindNodeResponse (transactionID, it->second, fromIdent, fromPort + 1); // to rport + auto closestNodeID = m_RoutingTable->FindClosestNodeInBucket (target); + if (closestNodeID) + { + auto it = m_Nodes.find (*closestNodeID); + if (it != m_Nodes.end ()) + SendFindNodeResponse (transactionID, it->second->GetNodeInfo (), fromIdent, fromPort + 1); // to rport + } + else if (m_Tunnel.GetLocalDestination ()) + SendFindNodeResponse (transactionID, + Node (m_NodeID, m_Tunnel.GetLocalDestination ()->GetIdentHash (), m_Port).GetNodeInfo (), // send ours + fromIdent, fromPort + 1); // to rport } } @@ -774,31 +793,30 @@ namespace torrents } else if (query == "get_peers") { - LogPrint (eLogDebug, "TorrentsDHT: get_peers response received"); - if (!values.empty ()) + LogPrint (eLogDebug, "TorrentsDHT: get_peers response received from ", ident.ToBase64 ()); + if (!values.empty () && values[0].empty ()) // nodes { - if (values[0].empty ()) // nodes + auto node = std::make_shared(nodeInfo); + m_Nodes.emplace (node->id, node); + if (m_RoutingTable) m_RoutingTable->AddNode (node->id); + LogPrint (eLogDebug, "TorrentsDHT: Node ", node->peer.ToBase64 (), ":", node->port, " added"); + } + else //values + { + auto torrent = torrentw.lock (); + if (torrent) { - auto node = std::make_shared(nodeInfo); - m_Nodes.emplace (node->id, node); - if (m_RoutingTable) m_RoutingTable->AddNode (node->id); - } - else //values - { - auto torrent = torrentw.lock (); - if (torrent) + std::unordered_set newPeers; + for (auto it1: values) + if (it1.size () == i2p::data::IdentHash::len) + newPeers.emplace ((const uint8_t *)it1.data ()); + if (!newPeers.empty ()) { - std::unordered_set newPeers; - for (auto it1: values) - if (it1.size () == i2p::data::IdentHash::len) - newPeers.emplace ((const uint8_t *)it1.data ()); - if (!newPeers.empty ()) - { - LogPrint (eLogDebug, "TorrentsDHT: ", newPeers.size (), " new peers received"); - m_Tunnel.ConnectToNewPeers (torrent, newPeers); - } - SendAnnouncePeerQuery (torrent->GetInfoHash (), token, ident, port + 1); // to rport + LogPrint (eLogDebug, "TorrentsDHT: ", newPeers.size (), " new peers received"); + m_Tunnel.ConnectToNewPeers (torrent, newPeers); } + LogPrint (eLogDebug, "TorrentsDHT: Send announce to ", ident.ToBase64 ()); + SendAnnouncePeerQuery (torrent->GetInfoHash (), token, ident, port + 1); // to rport } } } @@ -809,6 +827,8 @@ namespace torrents m_Nodes.emplace (node->id, node); if (m_RoutingTable) m_RoutingTable->AddNode (node->id); } + else if (query == "announce_peer") + LogPrint (eLogDebug, "TorrentsDHT: Announce peer response received"); else LogPrint (eLogInfo, "TorrentsDHT: Response to unknown query ", query); } @@ -841,7 +861,7 @@ namespace torrents } void TorrentsDHT::SendQueryMsg (std::string_view query, std::string_view arguments, - const i2p::data::IdentHash& toIdent, uint16_t toPort, bool isRaw) + const i2p::data::IdentHash& toIdent, uint16_t toPort, bool isRaw, std::shared_ptr torrent) { uint16_t transactionID = m_Tunnel.GetLocalDestination () ? m_Tunnel.GetLocalDestination ()->GetRng ()() : 1; auto msg = CreateDictionary ({ @@ -850,7 +870,7 @@ namespace torrents { "t", CreateByteString (std::string_view ((const char *)&transactionID, 2)) }, { "y", CreateByteString ("q") } }); - m_Queries.insert_or_assign (transactionID, std::make_tuple (toIdent, toPort, query, std::shared_ptr{})); + m_Queries.insert_or_assign (transactionID, std::make_tuple (toIdent, toPort, query, torrent)); if (isRaw) SendRawDatagram (msg, toIdent, toPort); else @@ -884,6 +904,18 @@ namespace torrents transactionID, toIdent, toPort); } + void TorrentsDHT::SendGetPeersQuery (std::shared_ptr torrent, const i2p::data::IdentHash& toIdent, uint16_t toPort) + { + if (!torrent) return; + const auto& infoHash = torrent->GetInfoHash (); + SendQueryMsg ("get_peers", CreateDictionary ({ + { "id", CreateByteString (std::string_view ((const char *)m_NodeID.data (), m_NodeID.size ())) }, + { "info_hash", CreateByteString (std::string_view ((const char *)infoHash.data (), infoHash.size ())) } + }), + toIdent, toPort, false, torrent); + } + + void TorrentsDHT::SendGetPeersResponse (std::string_view transactionID, std::shared_ptr torrent, uint64_t token, const i2p::data::IdentHash& toIdent, uint16_t toPort) { @@ -918,11 +950,9 @@ namespace torrents toIdent, toPort); } - void TorrentsDHT::SendFindNodeResponse (std::string_view transactionID, std::shared_ptr node, + void TorrentsDHT::SendFindNodeResponse (std::string_view transactionID, const NodeInfo& nodeInfo, const i2p::data::IdentHash& toIdent, uint16_t toPort) { - if (!node) return; - NodeInfo nodeInfo = node->GetNodeInfo (); SendResponseMsg (CreateDictionary ({ { "id", CreateByteString (std::string_view ((const char *)m_NodeID.data (), m_NodeID.size ())) }, { "nodes", CreateByteString (std::string_view ((const char *)nodeInfo.data (), nodeInfo.size ())) } @@ -1006,6 +1036,21 @@ namespace torrents ScheduleDHTExpirationCheck (); } } + + void TorrentsDHT::GetPeersAndAnnounce (std::shared_ptr torrent) + { + if (!torrent || !m_RoutingTable) return; + auto queriedNodeID = m_RoutingTable->FindClosestNode (torrent->GetInfoHash ()); + if (!queriedNodeID) return; // DHT is empty + auto it = m_Nodes.find (*queriedNodeID); + if (it != m_Nodes.end ()) + SendGetPeersQuery (torrent, it->second->peer, it->second->port); + else + { + LogPrint (eLogError, "TorrentsDHT: No nodeInfo for node from routing table"); + m_RoutingTable->RemoveNode (*queriedNodeID); + } + } } } diff --git a/libi2pd_client/TorrentsDHT.h b/libi2pd_client/TorrentsDHT.h index 46dd20b5..eb96cf2d 100644 --- a/libi2pd_client/TorrentsDHT.h +++ b/libi2pd_client/TorrentsDHT.h @@ -118,10 +118,11 @@ namespace torrents size_t GetNumBuckets () const; bool AddNode (const NodeID& id); + void RemoveNode (const NodeID& id); std::list > FindClosestNodes (const Torrent::InfoHash& infoHash, size_t num = 1) const; std::optional FindClosestNode (const Torrent::InfoHash& infoHash) const; std::list > GetExploratoryTargets (std::mt19937& rng) const; // (target, node to send find_node to) - NodeID FindClosestNodeInBucket (const NodeID& target) const; + std::optional FindClosestNodeInBucket (const NodeID& target) const; std::list DeleteExpiredNodes (uint64_t ts); void RemoveEmptyBuckets (); @@ -172,6 +173,7 @@ namespace torrents void HandleRawDatagram (const uint8_t * buf, size_t len); void SendPingQuery (const i2p::data::IdentHash& toIdent, uint16_t toPort); + void GetPeersAndAnnounce (std::shared_ptr torrent); private: @@ -190,15 +192,16 @@ namespace torrents void SendDatagram (std::string_view msg, const i2p::data::IdentHash& toIdent, uint16_t toPort); void SendRawDatagram (std::string_view msg, const i2p::data::IdentHash& toIdent, uint16_t toPort); void SendQueryMsg (std::string_view query, std::string_view arguments, - const i2p::data::IdentHash& toIdent, uint16_t toPort, bool isRaw = false); + const i2p::data::IdentHash& toIdent, uint16_t toPort, bool isRaw = false, std::shared_ptr torrent = nullptr); void SendFindNodeQuery (const NodeID& target, const i2p::data::IdentHash& toIdent, uint16_t toPort); + void SendGetPeersQuery (std::shared_ptr torrent, const i2p::data::IdentHash& toIdent, uint16_t toPort); void SendResponseMsg (std::string_view response, std::string_view transactionID, const i2p::data::IdentHash& toIdent, uint16_t toPort); void SendPingResponse (std::string_view transactionID, const i2p::data::IdentHash& toIdent, uint16_t toPort); void SendGetPeersResponse (std::string_view transactionID, std::shared_ptr torrent, uint64_t token, const i2p::data::IdentHash& toIdent, uint16_t toPort); void SendGetPeersResponse (std::string_view transactionID, std::shared_ptr node, uint64_t token, const i2p::data::IdentHash& toIdent, uint16_t toPort); - void SendFindNodeResponse (std::string_view transactionID, std::shared_ptr node, + void SendFindNodeResponse (std::string_view transactionID, const NodeInfo& nodeInfo, const i2p::data::IdentHash& toIdent, uint16_t toPort); void SendAnnouncePeerQuery (const Torrent::InfoHash& infoHash, uint64_t token, const i2p::data::IdentHash& toIdent, uint16_t toPort); diff --git a/libi2pd_client/TorrentsTunnel.cpp b/libi2pd_client/TorrentsTunnel.cpp index ff92ff83..5168dcf0 100644 --- a/libi2pd_client/TorrentsTunnel.cpp +++ b/libi2pd_client/TorrentsTunnel.cpp @@ -788,6 +788,11 @@ namespace torrents it.second->SetNextTrackerRequestTime (i, ts + nextInterval); } } + if (SupportsDHT () && ts >= it.second->GetNextDHTUpdateTime ()*1000) + { + if (m_DHT) m_DHT->GetPeersAndAnnounce (it.second); + it.second->SetNextDHTUpdateTime (ts/1000 + DHT_TORRENT_UPDATE_INTERVAL + GetLocalDestination ()->GetRng()() % DHT_TORRENT_UPDATE_INTERVAL_VARIANCE); + } } } }