/* * Copyright (c) 2026, The PurpleI2P Project * * This file is part of Purple i2pd project and licensed under BSD3 * * See full license text in LICENSE file at top of project tree */ #ifndef NO_TORRENTS #include #include #include #include #include #include "Log.h" #include "Timestamp.h" #include "I2PEndian.h" #include "AddressBook.h" #include "ClientContext.h" #include "TorrentsTunnel.h" namespace i2p { namespace torrents { TorrentsTunnel::TorrentsTunnel (std::string_view name, std::shared_ptr localDestination, std::string_view torrentsDir, std::string_view trackers, bool dht): i2p::client::I2PService (localDestination), m_Name (name), m_PeerID ("-I2PD-"), m_TorrentsDir (torrentsDir), m_TrackerRequestsCheckTimer (GetService ()), m_KeepAliveCheckTimer (GetService ()), m_ReconnectCheckTimer (GetService ()), m_TorrentsStatusUpdateTimer (GetService ()) { if (localDestination) m_PeerID += localDestination->GetIdentHash ().ToBase64 (); m_PeerID.resize (20, '0'); if (!trackers.empty ()) { // parse param std::vector trackersList; boost::split(trackersList, trackers, boost::is_any_of(","), boost::token_compress_on); // exclude duplicates std::unordered_map hosts; for (const auto& it: trackersList) { i2p::http::URL url (it); auto address = i2p::client::context.GetAddressBook ().GetAddress (url.host); if (address && address->IsIdentHash ()) { auto [it1, inserted] = hosts.emplace (address->identHash, url); if (!inserted && url.schema == "udp" && it1->second.schema == "http") it1->second = url; // replace http address by udp address } else LogPrint (eLogInfo, "TorrentsTunnel: Unexepcted tracker host ", url.host); } // create common trackers for (const auto& it: hosts) m_Trackers.emplace_back (TrackerInfo{ it.second, true, 0, 0, 0 }); } if (dht) m_DHT = std::make_unique(*this, TORRENT_PORT + (localDestination ? localDestination->GetRng()() % (65535 - TORRENT_PORT - 1) : 1)); } void TorrentsTunnel::Start () { i2p::client::I2PService::Start (); m_DiskIOService.Start (); auto dgramDest = GetLocalDestination ()->CreateDatagramDestination (true, i2p::datagram::eDatagramV1); // V1 for DHT if (dgramDest) dgramDest->SetRawReceiver (std::bind (&TorrentsTunnel::HandleRecvFromI2PRaw, std::static_pointer_cast(shared_from_this ()), std::placeholders::_1, std::placeholders::_2, std::placeholders::_3, std::placeholders::_4)); if (m_DHT) m_DHT->Start (); Accept (); if (!m_TorrentsDir.empty() && std::filesystem::exists (m_TorrentsDir) && std::filesystem::is_directory (m_TorrentsDir)) { for (const auto& it: std::filesystem::directory_iterator (m_TorrentsDir)) if (std::filesystem::is_regular_file (it.status()) && it.path ().extension () == ".torrent") ReadTorrentFile (it.path ()); } ScheduleTrackerRequestsCheck (); ScheduleReconnectCheck (); ScheduleKeepAliveCheck (); ScheduleStatusUpdate (); } void TorrentsTunnel::Stop () { if (m_DHT) m_DHT->Stop (); auto localDestination = GetLocalDestination (); if (localDestination) { localDestination->StopAcceptingStreams (); auto dgramDest = localDestination->GetDatagramDestination (); if (dgramDest) dgramDest->ResetRawReceiver (); } m_TrackerRequestsCheckTimer.cancel (); m_KeepAliveCheckTimer.cancel (); m_ReconnectCheckTimer.cancel (); m_TorrentsStatusUpdateTimer.cancel (); for (auto it: m_Torrents) StopTorrent (it.second); m_Torrents.clear (); m_DiskIOService.Stop (); i2p::client::I2PService::ClearHandlers (); // close connections i2p::client::I2PService::Stop (); } void TorrentsTunnel::ReadTorrentFile (const std::filesystem::path& torrentFilePath) { std::shared_ptr torrent; std::ifstream s(torrentFilePath, std::ifstream::binary); if (s) { s.seekg (0,std::ios::end); size_t len = s.tellg (); if (len > 0) { s.seekg(0, std::ios::beg); char * buf = new char[len]; s.read(buf, len); torrent = std::make_shared(std::string_view{buf, len}); delete[] buf; } else LogPrint (eLogError, "TorrentsTunnel: Empty file ", torrentFilePath); } else LogPrint (eLogError, "TorrentsTunnel: Can't open file ", torrentFilePath); if (torrent && !torrent->IsValid ()) { LogPrint (eLogError, "TorrentsTunnel: Invalid torrent file ", torrentFilePath, ". Skipped"); torrent = nullptr; } if (torrent) { torrent->SetFullPath (m_TorrentsDir/std::filesystem::path (torrent->GetName ())); InsertTorrent (torrent); InitTorrentFiles (torrent); } } void TorrentsTunnel::SaveTorrentFile (std::shared_ptr torrent) { if (!torrent) return; auto torrentFilePath = torrent->GetFullPath (); torrentFilePath += ".torrent"; std::ofstream f(torrentFilePath, std::ofstream::binary); if (f) { auto content = torrent->CreateTorrentFileContent (); f.write (content.data (), content.size ()); } else LogPrint (eLogError, "TorrentsTunnel: Can't open ", torrentFilePath); } void TorrentsTunnel::InitTorrentFiles (std::shared_ptr torrent) { if (!torrent || torrent->GetError () == eTorrentErrorMalformedMetaInfo) return; bool completed = true; for (auto it: torrent->GetFiles ()) { if (torrent->IsSingleFile ()) it->SetFullPath (torrent->GetFullPath ()); else it->UpdateFullPath (torrent->GetFullPath ()); if (!std::filesystem::exists (it->GetFullFilePath ())) { auto partFilePath = it->GetFullFilePath (); partFilePath += ".part"; if (!std::filesystem::exists (partFilePath)) CreateAndReserveFile (partFilePath, it->GetFileLength ()); completed = false; } } if (completed) torrent->SetComplete (); auto resumeFilePath = torrent->GetFullPath (); resumeFilePath += ".resume"; if (std::filesystem::exists (resumeFilePath)) { if (!torrent->IsComplete ()) { std::ifstream rs(resumeFilePath, std::ifstream::binary); if (rs) { rs.seekg (0,std::ios::end); size_t l = rs.tellg (); if (l > 0) { rs.seekg(0, std::ios::beg); std::vector bitfield(l); rs.read((char *)bitfield.data (), l); if (torrent->ApplyBitfield (bitfield)) CompleteTorrent (torrent); } } } else std::filesystem::remove (resumeFilePath); } } bool TorrentsTunnel::CreateAndReserveFile (const std::filesystem::path& filePath, size_t reserve) { if (std::filesystem::exists (filePath)) return false; auto subdirs = filePath.parent_path (); if (!subdirs.empty ()) { // try to create all subdirs try { std::filesystem::create_directories (subdirs); } catch (std::exception& ex) { LogPrint (eLogError, "TorrentsTunnel: Can't create subdirs ", subdirs, " : ", ex.what()); return false; } } // create file std::ofstream f(filePath, std::ios::binary); if (!f) return false; f.close (); if (reserve > 0) { // resize try { std::filesystem::resize_file (filePath, reserve); } catch (std::exception& ex) { LogPrint (eLogError, "TorrentsTunnel: Can't resize file ", filePath, " to ", reserve, " : ", ex.what()); return false; } } return true; } void TorrentsTunnel::CompleteTorrent (std::shared_ptr torrent) { boost::asio::post (GetDiskIOService (), [this, torrent]() { bool completed = true; for (auto it: torrent->GetFiles ()) { auto partFilePath = it->GetFullFilePath (); partFilePath += ".part"; if (std::filesystem::exists (partFilePath)) { it->Close (); // make sure file is closed before renaming std::error_code ec; std::filesystem::rename (partFilePath, it->GetFullFilePath (), ec); if (ec) { completed = false; LogPrint (eLogError, "TorrentsTunnel: Can't rename ", partFilePath); } } if (completed) it->Complete (); } if (completed) { torrent->SetComplete (); LogPrint (eLogInfo, "TorrentsTunnel: Download complete ", torrent->GetFullPath ()); boost::asio::post (GetService (), [this, torrent]() { // inform trackers that we are done RequestTorrentTrackers (torrent, eTrackerAnnounceEventCompleted); // close connections with seeds and reset stats for remaining auto conns = torrent->GetConnections (); for (auto it: conns) { if (it->GetRemoteBitfield ().all ()) // seed it->Close (); else it->ResetStats (); } }); } }); } std::shared_ptr TorrentsTunnel::FindTorrent (const Torrent::InfoHash& infoHash) const { std::lock_guard l(m_TorrentsMutex); auto it = m_Torrents.find (infoHash); if (it != m_Torrents.end ()) return it->second; return nullptr; } std::shared_ptr TorrentsTunnel::FindTorrentByID (int id) const { std::lock_guard l(m_TorrentsMutex); auto it = m_TorrentsByID.find (id); if (it != m_TorrentsByID.end ()) return it->second.lock (); return nullptr; } std::vector TorrentsTunnel::GetTorrentIDs () const { std::vector ids; std::lock_guard l(m_TorrentsMutex); for (const auto& it: m_TorrentsByID) if (!it.second.expired ()) ids.push_back (it.first); return ids; } std::list > TorrentsTunnel:: GetTorrents () const { std::list > torrents; std::lock_guard l(m_TorrentsMutex); for (const auto& it: m_Torrents) torrents.push_back (it.second); return torrents; } std::pair, int> TorrentsTunnel::AddTorrent (std::string_view torrentFileContent) { auto torrent = std::make_shared (torrentFileContent); auto it = m_Torrents.find (torrent->GetInfoHash ()); if (it == m_Torrents.end ()) { torrent->SetFullPath (m_TorrentsDir/std::filesystem::path (torrent->GetName ())); auto [id, inserted] = InsertTorrent (torrent); if (inserted) boost::asio::post (GetDiskIOService (), [this, torrent]() { SaveTorrentFile (torrent); InitTorrentFiles (torrent); }); if (!torrent->GetError ()) boost::asio::post (GetService (), [this, torrent] { RequestTorrentTrackers (torrent, eTrackerAnnounceEventNone); }); return { torrent, id }; } else return { it->second, FindTorrentID (it->second) }; return { torrent, 0 }; } std::pair, int> TorrentsTunnel::AddMagnet (std::string_view magnet) { // magnet:?xt=urn:btih:&dn=&tr= static constexpr std::string_view magnetPrefix { "magnet:?" }; if (magnet.starts_with (magnetPrefix)) { Torrent::InfoHash infoHash; bool isInfoHashFound = false; std::string announce, name, hexStr; magnet = magnet.substr (magnetPrefix.size ()); while (!magnet.empty()) { auto pos = magnet.find ('&'); std::string_view param; if (pos != std::string::npos) { param = magnet.substr (0, pos); magnet = magnet.substr (pos + 1); } else { param = magnet; magnet = ""; } static constexpr std::string_view hashPrefix { "xt=urn:btih:" }; static constexpr std::string_view trackerPrefix { "tr=" }; static constexpr std::string_view namePrefix { "dn=" }; if (param.starts_with (hashPrefix)) { hexStr = param.substr (hashPrefix.size (), infoHash.size ()*2); try { boost::algorithm::unhex (hexStr.begin(), hexStr.end(), infoHash.begin()); isInfoHashFound = true; } catch (std::exception& ex) { LogPrint (eLogInfo, "TorentsTunnel: Can't unhex magnet hash ", hexStr); } } else if (param.starts_with (trackerPrefix)) announce = i2p::http::UrlDecode (param.substr (trackerPrefix.size ())); else if (param.starts_with (namePrefix)) name = i2p::http::UrlDecode (param.substr (namePrefix.size ())); } if (isInfoHashFound) { auto it = m_Torrents.find (infoHash); if (it == m_Torrents.end ()) { auto torrent = std::make_shared (infoHash); if (!announce.empty ()) torrent->SetAnnounce (announce); torrent->SetName (name.empty () ? hexStr : name); auto [id, inserted] = InsertTorrent (torrent); if (inserted) boost::asio::post (GetService (), [this, torrent] { RequestTorrentTrackers (torrent, eTrackerAnnounceEventNone); }); return { torrent, id }; } else return { it->second, FindTorrentID (it->second) }; } } return { nullptr, 0 }; } std::pair TorrentsTunnel::InsertTorrent (std::shared_ptr torrent) { if (!torrent) return { 0, false }; std::lock_guard l(m_TorrentsMutex); if (m_Torrents.emplace (torrent->GetInfoHash (), torrent).second) { int id = 1; if (!m_TorrentsByID.empty ()) id = m_TorrentsByID.rbegin ()->first + 1; m_TorrentsByID.emplace (id, torrent); if (!torrent->GetAnnounce ().empty ()) torrent->SetAnnounceTrackerID (AddTracker (torrent->GetAnnounce (), false)); // add announce to trackers return { id, true }; } else // already exists, find id return { FindTorrentID (torrent), false }; return { 0, false }; } bool TorrentsTunnel::RemoveTorrent (int id, bool deleteFiles) { std::shared_ptr torrent; { std::lock_guard l(m_TorrentsMutex); auto it = m_TorrentsByID.find (id); if (it == m_TorrentsByID.end ()) return false; torrent = it->second.lock (); m_TorrentsByID.erase (it); if (!torrent) return false; m_Torrents.erase (torrent->GetInfoHash ()); } boost::asio::post (GetService (), [this, torrent, deleteFiles]() { RemoveTorrent (torrent, deleteFiles); }); return true; } void TorrentsTunnel::RemoveTorrent (std::shared_ptr torrent, bool deleteFiles) { if (!torrent) return; StopTorrent (torrent); boost::asio::post (GetDiskIOService (), [torrent, deleteFiles]() { auto fullPath = torrent->GetFullPath (); auto torrentFilePath = fullPath; torrentFilePath += ".torrent"; std::error_code ec; std::filesystem::remove (torrentFilePath, ec); if (ec) LogPrint (eLogError, "TorrentsTunnel: Can't delete ", torrentFilePath); auto resumeFilePath = fullPath; resumeFilePath += ".resume"; if (std::filesystem::exists (resumeFilePath)) { std::filesystem::remove (resumeFilePath, ec); if (ec) LogPrint (eLogError, "TorrentsTunnel: Can't delete ", resumeFilePath); } if (deleteFiles) { if (std::filesystem::exists (fullPath)) { std::filesystem::remove_all (fullPath, ec); if (ec) LogPrint (eLogError, "TorrentsTunnel: Can't delete ", fullPath); } auto partFilePath = fullPath; partFilePath += ".part"; if (std::filesystem::exists (partFilePath)) { std::filesystem::remove (partFilePath, ec); if (ec) LogPrint (eLogError, "TorrentsTunnel: Can't delete ", partFilePath); } } }); } int TorrentsTunnel::FindTorrentID (std::shared_ptr torrent) const { if (torrent) { auto it = std::find_if (m_TorrentsByID.begin (), m_TorrentsByID.end (), [torrent](const auto& torrentByID) { return torrentByID.second.lock () == torrent; }); if (it != m_TorrentsByID.end ()) return it->first; } return 0; } bool TorrentsTunnel::StopTorrent (int id) { std::shared_ptr torrent; { std::lock_guard l(m_TorrentsMutex); auto it = m_TorrentsByID.find (id); if (it == m_TorrentsByID.end ()) return false; torrent = it->second.lock (); if (!torrent) return false; } boost::asio::post (GetService (), [this, torrent]() { StopTorrent (torrent); }); return true; } void TorrentsTunnel::StopTorrent (std::shared_ptr torrent) { if (!torrent) return; torrent->SetStopped (true); torrent->UpdateStatus (i2p::util::GetMonotonicSeconds ()); // inform trackers that we stopped RequestTorrentTrackers (torrent, eTrackerAnnounceEventStopped); // close connections auto connections = torrent->GetConnections (); for (auto it: connections) it->Close (); } bool TorrentsTunnel::StartTorrent (int id) { std::shared_ptr torrent; { std::lock_guard l(m_TorrentsMutex); auto it = m_TorrentsByID.find (id); if (it == m_TorrentsByID.end ()) return false; torrent = it->second.lock (); if (!torrent) return false; } boost::asio::post (GetService (), [this, torrent]() { torrent->SetStopped (false); // inform trackers that we started RequestTorrentTrackers (torrent, eTrackerAnnounceEventStarted); }); return true; } void TorrentsTunnel::UpdateTorrentInfo (std::shared_ptr torrent, std::string_view info) { if (!torrent || torrent->IsStopped ()) return; torrent->ParseInfo (info); torrent->SetFullPath (m_TorrentsDir/std::filesystem::path (torrent->GetName ())); boost::asio::post (GetDiskIOService (), [this, torrent]() { SaveTorrentFile (torrent); InitTorrentFiles (torrent); }); } void TorrentsTunnel::Accept () { auto localDestination = GetLocalDestination (); if (localDestination) { if (!localDestination->IsAcceptingStreams ()) // set it as default if not set yet localDestination->AcceptStreams ([this](std::shared_ptr stream) { if (stream) { auto conn = std::make_shared (shared_from_this (), stream); AddHandler (conn); conn->ReceiveHandshake (); } }); } else LogPrint (eLogError, "TorrentsTunnel: Local destination not set"); } void TorrentsTunnel::RequestTorrentTrackers (std::shared_ptr torrent, TrackerAnnounceEvent event) { for (size_t i = 0; i < m_Trackers.size (); i++) RequestTracker (i, torrent, event); } bool TorrentsTunnel::RequestTracker (size_t trackerID, std::shared_ptr torrent, TrackerAnnounceEvent event) { if (!torrent || trackerID >= m_Trackers.size() || (!std::get<1>(m_Trackers[trackerID]) && std::get<0>(m_Trackers[trackerID]).host != i2p::http::URL (torrent->GetAnnounce ()).host)) // not common tracker or with torrent's announce return false; i2p::http::URL reqURL = std::get<0>(m_Trackers[trackerID]); if (!reqURL.host.ends_with (".i2p")) { LogPrint (eLogWarning, "TorrentsTunnel: Non-I2P address ", reqURL.host, " for torrent ", torrent->GetName ()); return false; } if (reqURL.schema == "udp") { auto ts = i2p::util::GetMonotonicMilliseconds (); auto& [announce, common, connectionID, expiration, fromPort] = m_Trackers[trackerID]; if (connectionID && ts <= expiration) SendAnnounceToDatagramTracker (trackerID, connectionID, torrent, reqURL.host, reqURL.port, fromPort, event); else { if (!expiration || ts >= expiration) // don't connect if pending connection { connectionID = 0; expiration = ts + DATAGRAM_TRACKER_TRANSACTION_TIMEOUT; ConnectToDatagramTracker (trackerID, reqURL.host, reqURL.port); } } return true; } std::map params; params.emplace ("info_hash", torrent->GetHexStringInfoHash ()); params.emplace ("peer_id", m_PeerID); params.emplace ("ip", GetLocalDestination ()->GetIdentity ()->ToBase64 () + ".i2p"); params.emplace ("port", std::to_string (TORRENT_PORT)); // 6881 params.emplace ("compact", "1"); params.emplace ("uploaded", std::to_string (torrent->GetUploaded ())); params.emplace ("downloaded", std::to_string (torrent->GetLength () - torrent->GetLeft ())); params.emplace ("left", std::to_string (torrent->GetLeft ())); int numWant = 0; if (!torrent->IsComplete () && (event == eTrackerAnnounceEventNone || event == eTrackerAnnounceEventStarted)) numWant = TRACKER_MAX_NUM_WANT; params.emplace ("numwant", std::to_string (numWant)); if (event != eTrackerAnnounceEventNone) params.emplace ("event", TrackerAnnounceEventStr[event]); reqURL.create_query (params); auto req = std::make_shared >(boost::beast::http::verb::get, reqURL.to_string (true), 11); // HTTP 1.1 req->set (boost::beast::http::field::host, reqURL.host); req->set (boost::beast::http::field::user_agent, "I2PSocketEepGet"); req->keep_alive (false); // Connection: close CreateStream ([this, req, torrent, trackerID](std::shared_ptr stream) { if (stream) { auto httpStream = std::make_shared(stream); boost::beast::http::async_write (*httpStream, *req, std::bind (&TorrentsTunnel::TrackerRequestSent, this, std::placeholders::_1, std::placeholders::_2, httpStream, torrent, req, trackerID)); } }, reqURL.host, reqURL.port); return true; } void TorrentsTunnel::TrackerRequestSent (const boost::beast::error_code& ecode, size_t bytes_transferred, std::shared_ptr httpStream, std::shared_ptr torrent, std::shared_ptr > req, size_t trackerID) { if (!ecode) { // receive auto buf = std::make_shared (); auto res = std::make_shared >(); boost::beast::http::async_read (*httpStream, *buf, *res, [this, httpStream, torrent, buf, res, trackerID](const boost::beast::error_code& ecode, size_t bytes_transferred) { httpStream->GetStream ()->AsyncClose (); if (!ecode) { if (res->result () == boost::beast::http::status::ok) { torrent->ParseTrackerResponse (trackerID, res->body ()); ConnectToPeers (torrent, trackerID); } else LogPrint (eLogWarning, "TorrentsTunnel: Tracker ", trackerID, " response code ", res->result_int()); } }); } } void TorrentsTunnel::ConnectToPeer (std::shared_ptr torrent, const i2p::data::IdentHash& peer) { if (!torrent || torrent->IsConnectedToPeer (peer)) return; LogPrint (eLogDebug, "TorrentsTunnel: Connecting to peer ", peer.ToBase32 () + ".b32.i2p"); if (peer == GetLocalDestination ()->GetIdentHash ()) { LogPrint (eLogInfo, "TorrentsTunnel: Can't connect to self"); return; } CreateStream ([this, torrent, peer](std::shared_ptr stream) { if (stream) { LogPrint (eLogDebug, "TorrentsTunnel: Connected to peer ", peer.ToBase32 () + ".b32.i2p"); auto connection = std::make_shared(shared_from_this (), stream, torrent); AddHandler (connection); connection->Connect (); } else LogPrint (eLogInfo, "TorrentsTunnel: Can't connect to peer ", peer.ToBase32 () + ".b32.i2p"); }, std::make_shared(peer), TORRENT_PORT); } size_t TorrentsTunnel::ConnectToPeers (std::shared_ptr torrent) { if (!torrent) return 0; auto peersToConnect = torrent->GetNonConnectedPeers (); if (!peersToConnect.empty ()) { for (const auto& it: peersToConnect) ConnectToPeer (torrent, it); } return peersToConnect.size (); } size_t TorrentsTunnel::ConnectToPeers (std::shared_ptr torrent, size_t trackerID) { if (!torrent) return 0; auto peersToConnect = torrent->GetNonConnectedPeers (trackerID); if (!peersToConnect.empty ()) { for (const auto& it: peersToConnect) ConnectToPeer (torrent, it); } return peersToConnect.size (); } void TorrentsTunnel::ConnectToNewPeers (std::shared_ptr torrent, std::unordered_set& newPeers) { if (!torrent) return; if (!newPeers.empty ()) { for (const auto& it: newPeers) if (!torrent->IsConnectedToPeer (it)) ConnectToPeer (torrent, it); } } void TorrentsTunnel::ScheduleTrackerRequestsCheck () { m_TrackerRequestsCheckTimer.expires_after (std::chrono::milliseconds(TRACKER_REQUESTS_CHECK_TIMEOUT)); m_TrackerRequestsCheckTimer.async_wait (std::bind (&TorrentsTunnel::HandleTrackerRequestsCheckTimer, this, std::placeholders::_1)); } void TorrentsTunnel::HandleTrackerRequestsCheckTimer (const boost::system::error_code& ecode) { if (ecode != boost::asio::error::operation_aborted) { auto ts = i2p::util::GetMonotonicMilliseconds (); if (!m_DatragramTrackerTransactions.empty ()) { // cleanup expired transactions auto it = m_DatragramTrackerTransactions.begin (); while (it != m_DatragramTrackerTransactions.end ()) { if (ts > std::get<3>(it->second) + DATAGRAM_TRACKER_TRANSACTION_TIMEOUT) it = m_DatragramTrackerTransactions.erase (it); else it++; } } for (auto it: m_Torrents) { if (!it.second->IsStopped ()) { for (size_t i = 0; i < m_Trackers.size (); i++) { if (!it.second->GetNextTrackerRequestTime (i) && // first time (std::get<1>(m_Trackers[i]) || it.second->GetAnnounceTrackerID () == (int)i)) { auto initialInterval = GetLocalDestination ()->GetRng()() % TRACKER_INITIAL_REQUEST_INTERVAL_VARIANCE; if (initialInterval <= TRACKER_REQUESTS_CHECK_TIMEOUT) initialInterval = 0; // request immeditely it.second->SetNextTrackerRequestTime (i, ts + initialInterval); } if (ts >= it.second->GetNextTrackerRequestTime (i)) { if (i2p::util::GetSecondsSinceEpoch () > it.second->GetLastTrackerUpdateTime (i) + 2*it.second->GetInterval (i)/1000 && it.second->GetTrackerError (i).empty ()) { it.second->SetTrackerError (i, "No response"); it.second->SetInterval (i, MIN_TRACKER_REQUESTS_INTERVAL); // try again shortly } if (RequestTracker (i, it.second, eTrackerAnnounceEventNone)) { auto nextInterval = it.second->GetInterval (i) + GetLocalDestination ()->GetRng()() % TRACKER_REQUESTS_INTERVAL_VARIANCE; 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); } } } } ScheduleTrackerRequestsCheck (); } } void TorrentsTunnel::ScheduleKeepAliveCheck () { m_KeepAliveCheckTimer.expires_after (std::chrono::seconds(PEER_KEEP_ALIVE_CHECK_INTERVAL)); m_KeepAliveCheckTimer.async_wait (std::bind (&TorrentsTunnel::HandleKeepAliveCheckTimer, this, std::placeholders::_1)); } void TorrentsTunnel::HandleKeepAliveCheckTimer (const boost::system::error_code& ecode) { if (ecode != boost::asio::error::operation_aborted) { auto ts = i2p::util::GetMonotonicSeconds (); IterateHandlers ([ts](std::shared_ptr handler) { if (handler) std::static_pointer_cast(handler)->CheckKeepAlive (ts); }); ScheduleKeepAliveCheck (); } } void TorrentsTunnel::ScheduleReconnectCheck () { m_ReconnectCheckTimer.cancel (); m_ReconnectCheckTimer.expires_after (std::chrono::seconds(RECONNECT_CHECK_INTERVAL)); m_ReconnectCheckTimer.async_wait (std::bind (&TorrentsTunnel::HandleReconnectCheckTimer, this, std::placeholders::_1)); } void TorrentsTunnel::HandleReconnectCheckTimer (const boost::system::error_code& ecode) { if (ecode != boost::asio::error::operation_aborted) { auto ts = i2p::util::GetMonotonicSeconds (); for (auto it: m_Torrents) if (ts > it.second->GetNextReconnectTime ()) { if (!it.second->IsComplete () && !it.second->IsStopped ()) { auto numPeers = ConnectToPeers (it.second); if (numPeers) LogPrint (eLogDebug, "TorrentsTunnel: Reconnecting to ", numPeers, " peers"); } it.second->SetNextReconnectTime (ts + RECONNECT_INTERVAL + GetLocalDestination ()->GetRng ()() % RECONNECT_INTERVAL_VARIANCE); } ScheduleReconnectCheck (); } } void TorrentsTunnel::ScheduleStatusUpdate () { m_TorrentsStatusUpdateTimer.cancel (); m_TorrentsStatusUpdateTimer.expires_after (std::chrono::seconds(TORRENTS_STATUS_UPDATE_CHECK_INTERVAL)); m_TorrentsStatusUpdateTimer.async_wait (std::bind (&TorrentsTunnel::HandleTorrentsStatusUpdateTimer, this, std::placeholders::_1)); } void TorrentsTunnel::HandleTorrentsStatusUpdateTimer (const boost::system::error_code& ecode) { if (ecode != boost::asio::error::operation_aborted) { auto ts = i2p::util::GetMonotonicSeconds (); for (auto it: m_Torrents) if (ts > it.second->GetNextUpdateStatusTime ()) { if (!it.second->IsStopped () && (it.second->IsActive () || !it.second->IsComplete ())) { if (it.second->UpdateStatus (ts)) { if (!it.second->IsComplete ()) CompleteTorrent (it.second); } else UpdatePeersPerPiece (it.second); } boost::asio::post (GetDiskIOService (), [torrent = it.second, ts]() { for (auto it: torrent->GetFiles ()) if (ts > it->GetLastAccessTime () + TORRENT_FILE_INACTIVITY_TIMEOUT) it->Close (); }); it.second->SetNextUpdateStatusTime (ts + TORRENTS_STATUS_UPDATE_INTERVAL + GetLocalDestination ()->GetRng ()() % TORRENTS_STATUS_UPDATE_INTERVAL_VARIANCE); } ScheduleStatusUpdate (); } } void TorrentsTunnel::UpdatePeersPerPiece (std::shared_ptr torrent) { if (!torrent) return; torrent->StartCountingPeers (); IterateHandlers ([torrent](std::shared_ptr handler) { if (handler) { auto conn = std::static_pointer_cast(handler); if (conn->GetTorrent () == torrent) torrent->ApplyPeerRemoteBitfield (conn->GetRemoteBitfield ()); } }); } void TorrentsTunnel::HandleRecvFromI2PRaw (uint16_t fromPort, uint16_t toPort, const uint8_t * buf, size_t len) { if (m_DHT && toPort == m_DHT->GetRPort ()) { m_DHT->HandleRawDatagram (buf, len); return; } // response from tracker if (len < 8) return; uint32_t action = bufbe32toh (buf); switch (action) { case eDatagramTrackerActionConnect: HandleConnectResponse (buf + 4, len - 4); break; case eDatagramTrackerActionAnnounce: HandleAnnounceResponse (buf + 4, len - 4); break; case eDatagramTrackerActionError: HandleErrorResponse (buf + 4, len - 4); break; default: LogPrint (eLogInfo, "TorrentsTunnel: Unexpected action ", action, " from tracker"); } } void TorrentsTunnel::HandleErrorResponse (const uint8_t * buf, size_t len) { uint32_t transactionID = bufbe32toh (buf); LogPrint (eLogDebug, "TorrentsTunnel: Datagram tracker action error response ", transactionID); auto it = m_DatragramTrackerTransactions.find (transactionID); if (it == m_DatragramTrackerTransactions.end ()) { LogPrint (eLogInfo, "TorrentsTunnel: Datagram tracker transaction ", transactionID, " not found"); return; } std::string_view error ((const char *)(buf + 4), len - 4); LogPrint (eLogInfo, "TorrentsTunnel: Datagram tracker error response: ", error); auto [trackerID, fromPort, torrent, ts] = it->second; if (!torrent.expired ()) // response to announce torrent.lock ()->SetTrackerError (trackerID, error); m_DatragramTrackerTransactions.erase (it); } void TorrentsTunnel::ConnectToDatagramTracker (size_t trackerID, std::string_view dest, uint16_t port) { LogPrint (eLogDebug, "TorrentsTunnel: Connecting to datagram tracker ", dest, ":", port); auto address = i2p::client::context.GetAddressBook ().GetAddress (dest); if (address && address->IsIdentHash ()) { auto localDestination = GetLocalDestination (); auto dgramDest = localDestination->GetDatagramDestination (); if (dgramDest) { uint8_t connectRequest[16]; htobe64buf (connectRequest, 0x41727101980); // protocol_id htobe32buf (connectRequest + 8, eDatagramTrackerActionConnect); // action uint32_t transactionID = localDestination->GetRng()(); htobe32buf (connectRequest + 12, transactionID); // transactionID uint16_t fromPort = localDestination->GetRng()() % 1000 + 6000; if (m_DHT && fromPort == m_DHT->GetRPort ()) fromPort++; auto session = dgramDest->GetSession (address->identHash); if (session) { m_DatragramTrackerTransactions.emplace (transactionID, std::make_tuple ( trackerID, fromPort, std::shared_ptr(nullptr), i2p::util::GetMonotonicMilliseconds ())); session->SetVersion (i2p::datagram::eDatagramV2); // send datagram2 dgramDest->SendDatagram (session, connectRequest, 16, fromPort, port); dgramDest->FlushSendQueue (session); } else LogPrint (eLogInfo, "TorrentsTunnel: Can't obtain datagram session to ", dest); } else LogPrint (eLogError, "TorrentsTunnel: Datagram destination is not avaliable"); } else LogPrint (eLogInfo, "TorrentsTunnel: Tracker not found: ", dest); } void TorrentsTunnel::HandleConnectResponse (const uint8_t * buf, size_t len) { if (len < 12) { LogPrint (eLogInfo, "TorrentsTunnel: Unexpected connect response length ", len + 4); return; } uint32_t transactionID = bufbe32toh (buf); LogPrint (eLogDebug, "TorrentsTunnel: Datagram tracker action connect response ", transactionID); auto it = m_DatragramTrackerTransactions.find (transactionID); if (it == m_DatragramTrackerTransactions.end ()) { LogPrint (eLogInfo, "TorrentsTunnel: Datagram tracker transaction ", transactionID, " not found"); return; } uint64_t connectionID = bufbe64toh (buf + 4); int lifetime = DATAGRAM_TRACKER_CONNECTION_EXPIRATION; if (len >= 14) { lifetime = std::max (bufbe16toh (buf + 12)*1000, DATAGRAM_TRACKER_CONNECTION_EXPIRATION); LogPrint (eLogDebug, "TorrentsTunnel: Datagram tracker connection lifetime set to ", lifetime/1000, " seconds"); } auto [trackerID, fromPort, torrent, ts] = it->second; if (trackerID < m_Trackers.size ()) { auto& [announce, common, connID, connExpiration, connFromPort] = m_Trackers[trackerID]; connID = connectionID; connExpiration = i2p::util::GetMonotonicMilliseconds () + lifetime; connFromPort = fromPort; } m_DatragramTrackerTransactions.erase (it); } void TorrentsTunnel::SendAnnounceToDatagramTracker (size_t trackerID, uint64_t connectionID, std::shared_ptr torrent, std::string_view dest, uint16_t port, uint16_t fromPort, TrackerAnnounceEvent event) { LogPrint (eLogDebug, "TorrentsTunnel: Sending announce to datagram tracker ", dest, ":", port); auto address = i2p::client::context.GetAddressBook ().GetAddress (dest); if (address && address->IsIdentHash ()) { auto localDestination = GetLocalDestination (); auto dgramDest = localDestination->GetDatagramDestination (); if (dgramDest) { uint8_t announce[98]; htobe64buf (announce, connectionID); // connection_id htobe32buf (announce + 8, eDatagramTrackerActionAnnounce); // action uint32_t transactionID = localDestination->GetRng()(); htobe32buf (announce + 12, transactionID); // transaction_id memcpy (announce + 16, torrent->GetInfoHash ().data (), 20); // info_hash memcpy (announce + 36, m_PeerID.data (), 20); // peer_id auto left = torrent->GetLeft (); htobe64buf (announce + 56, torrent->GetLength () - left); // downloaded htobe64buf (announce + 64, left); // left htobe64buf (announce + 72, torrent->GetUploaded ()); // uploaded htobe32buf (announce + 80, event); // event htobe32buf (announce + 84, 0); // IP address 0:not used htobe32buf (announce + 88, 0); // key, ignored int numWant = 0; if (!torrent->IsComplete () && (event == eTrackerAnnounceEventNone || event == eTrackerAnnounceEventStarted)) numWant = TRACKER_MAX_NUM_WANT; htobe32buf (announce + 92, numWant); // num_want htobe16buf (announce + 96, fromPort); // from port auto session = dgramDest->GetSession (address->identHash); if (session) { m_DatragramTrackerTransactions.emplace (transactionID, std::make_tuple ( trackerID, fromPort, torrent, i2p::util::GetMonotonicMilliseconds ())); session->SetVersion (i2p::datagram::eDatagramV3); // send datagram3 dgramDest->SendDatagram (session, announce, 98, fromPort, port); dgramDest->FlushSendQueue (session); } else LogPrint (eLogInfo, "TorrentsTunnel: Can't obtain datagram session"); } } else LogPrint (eLogInfo, "TorrentsTunnel: Tracker not found: ", dest); } void TorrentsTunnel::HandleAnnounceResponse (const uint8_t * buf, size_t len) { if (len < 16) { LogPrint (eLogInfo, "TorrentsTunnel: Unexpected announce response length ", len + 4); return; } uint32_t transactionID = bufbe32toh (buf); LogPrint (eLogDebug, "TorrentsTunnel: Datagram tracker action announce response ", transactionID); auto it = m_DatragramTrackerTransactions.find (transactionID); if (it == m_DatragramTrackerTransactions.end ()) { LogPrint (eLogInfo, "TorrentsTunnel: Datagram tracker transaction ", transactionID, " not found"); return; } auto [trackerID, fromPort, torrent, ts] = it->second; if (!torrent.expired ()) { uint32_t interval = bufbe32toh (buf + 4); uint32_t leechers = bufbe32toh (buf + 8); uint32_t seeders = bufbe32toh (buf + 12); torrent.lock ()->HandleDatagramTrackerResponse (trackerID, interval, buf + 16, len - 16, seeders, leechers); } m_DatragramTrackerTransactions.erase (it); } int TorrentsTunnel::AddTracker (std::string_view announce, bool isCommon) { if (announce.empty ()) return -1; i2p::http::URL reqURL (announce); if (!reqURL.host.ends_with (".i2p")) { LogPrint (eLogInfo, "TorrentsTunnel: Non-I2P address ", reqURL.host, " in tracker announce"); return -1; } auto reqAddr = i2p::client::context.GetAddressBook ().GetAddress (reqURL.host); if (reqAddr) { auto it = std::find_if (m_Trackers.begin (), m_Trackers.end (), [host = std::string_view (reqURL.host), reqAddr](const TrackerInfo& tracker) { const auto& url = std::get<0>(tracker); if (url.host == host) return true; auto addr = i2p::client::context.GetAddressBook ().GetAddress (url.host); if (!addr) return false; return *addr == *reqAddr; }); if (it != m_Trackers.end ()) // existing return it - m_Trackers.begin(); } m_Trackers.emplace_back (TrackerInfo{ announce, isCommon, 0, 0, 0 }); return m_Trackers.size () - 1; } void TorrentsTunnel::SendDHTPingQuery (const i2p::data::IdentHash& toIdent, uint16_t toPort) { boost::asio::post (GetService (), [this, toIdent, toPort]() { if (m_DHT) m_DHT->SendPingQuery (toIdent, toPort); }); } } } #endif // NO_TORRENTS