get peers throught DHT and announce own torrents

This commit is contained in:
orignal
2026-09-29 18:59:46 -04:00
parent 14eb418c5e
commit 5bf161e29c
5 changed files with 115 additions and 55 deletions
+14 -12
View File
@@ -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") }
+6 -1
View File
@@ -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<PeerConnection> conn);
void RemoveConnection (std::shared_ptr<PeerConnection> conn);
std::list<std::shared_ptr<PeerConnection> > GetConnections ();
@@ -360,7 +365,7 @@ namespace torrents
bool m_IsComplete, m_IsStopped, m_IsSingleFile;
std::list<std::shared_ptr<TorrentFile> > 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;
};
+84 -39
View File
@@ -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<std::pair<NodeID, Distance> > RoutingTable::FindClosestNodes (const Torrent::InfoHash& infoHash, size_t num) const
{
std::list<std::pair<NodeID, Distance> > 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<NodeID> 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<Node>(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<Node>(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<i2p::data::IdentHash> 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<i2p::data::IdentHash> 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> 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<Torrent>{}));
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> 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<DHTTorrent> 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<const Node> 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> 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);
}
}
}
}
+6 -3
View File
@@ -118,10 +118,11 @@ namespace torrents
size_t GetNumBuckets () const;
bool AddNode (const NodeID& id);
void RemoveNode (const NodeID& id);
std::list<std::pair<NodeID, Distance> > FindClosestNodes (const Torrent::InfoHash& infoHash, size_t num = 1) const;
std::optional<NodeID> FindClosestNode (const Torrent::InfoHash& infoHash) const;
std::list<std::pair<NodeID, NodeID> > GetExploratoryTargets (std::mt19937& rng) const; // (target, node to send find_node to)
NodeID FindClosestNodeInBucket (const NodeID& target) const;
std::optional<NodeID> FindClosestNodeInBucket (const NodeID& target) const;
std::list<NodeID> 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> 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> torrent = nullptr);
void SendFindNodeQuery (const NodeID& target, const i2p::data::IdentHash& toIdent, uint16_t toPort);
void SendGetPeersQuery (std::shared_ptr<Torrent> 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<DHTTorrent> torrent,
uint64_t token, const i2p::data::IdentHash& toIdent, uint16_t toPort);
void SendGetPeersResponse (std::string_view transactionID, std::shared_ptr<const Node> node,
uint64_t token, const i2p::data::IdentHash& toIdent, uint16_t toPort);
void SendFindNodeResponse (std::string_view transactionID, std::shared_ptr<const Node> 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);
+5
View File
@@ -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);
}
}
}
}