// Copyright (c) 2026 The Limenka developers // Distributed under the MIT software license, see the accompanying // file COPYING or http://www.opensource.org/licenses/mit-license.php. #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include namespace node { namespace { // --------------------------------------------------------------------------- // Foreign network table: message start + default port per chain. // --------------------------------------------------------------------------- const std::vector KNOWN_NETWORKS{ {"mainnet", {0xf9, 0xbe, 0xb4, 0xd9}, 8333, {}}, {"testnet3", {0x07, 0x09, 0x11, 0x0b}, 18333, {}}, {"testnet4", {0x1c, 0x16, 0x3f, 0x28}, 48333, {}}, {"signet", {0x0a, 0x03, 0xcf, 0x40}, 38333, {}}, {"regtest", {0xfa, 0xbf, 0xb5, 0xda}, 18444, {}}, // Limenka's own fork network: lets the bridge run against a second // limenka node (functional tests, cross-node mesh). {"fork", {0xfa, 0xbf, 0xd5, 0xc6}, 19444, {}}, }; constexpr std::chrono::seconds RECONNECT_INTERVAL{5}; constexpr size_t SEEN_MAX{100000}; constexpr size_t TX_CACHE_MAX{10000}; constexpr size_t MAX_INV_PER_MESSAGE{50000}; constexpr size_t INVENTORY_BROADCAST_MAX_LOCAL{1000}; // Flood guard per peer: at most this many transactions accepted per window. constexpr size_t MAX_TX_PER_WINDOW{1000}; constexpr std::chrono::seconds TX_WINDOW{10}; } // namespace // --------------------------------------------------------------------------- // Peer: one connection to one foreign node. Reads and writes are async on // the bridge io_context; all socket operations are serialized through it. // --------------------------------------------------------------------------- class MempoolBridge::Peer final : public std::enable_shared_from_this { public: Peer(MempoolBridge& bridge, const BridgeNetwork& net, const CService& remote) : m_bridge{bridge}, m_net{net}, m_remote{remote}, m_socket{*bridge.m_io} { } void Start() { auto self{shared_from_this()}; boost::system::error_code ec; const auto ep = ResolveEndpoint(ec); if (ec) { m_dead = true; return; } m_socket.async_connect(ep, [self](const boost::system::error_code& err) { if (err) { LogPrintf("mempoolbridge: connect %s failed: %s\n", self->m_remote.ToStringAddrPort(), err.message()); self->m_dead = true; return; } LogPrintf("mempoolbridge: tcp connected to %s\n", self->m_remote.ToStringAddrPort()); self->SendVersion(); self->DoReadHeader(); }); } bool Dead() const { return m_dead; } bool Ready() const { return m_ready; } const CService& Remote() const { return m_remote; } const std::string& NetworkName() const { return m_net.name; } void SendInv(const CTransactionRef& tx) { // Announce by wtxid (MSG_WTX, BIP339): peers serve witness // transactions only when requested by wtxid - a txid request gets a // witness-stripped tx, which would fail limenka validation. DataStream ss; ss << CompactSizeWriter(1) << CInv{MSG_WTX, tx->GetWitnessHash()}; SendMessage(NetMsgType::INV, ss); } void SendTx(const CTransactionRef& tx) { DataStream ss; ss << TX_WITH_WITNESS(*tx); SendMessage(NetMsgType::TX, ss); } void Disconnect() { boost::asio::post(*m_bridge.m_io, [self = shared_from_this()]() { boost::system::error_code ec; self->m_socket.close(ec); self->m_dead = true; }); } private: MempoolBridge& m_bridge; const BridgeNetwork& m_net; CService m_remote; boost::asio::ip::tcp::socket m_socket; bool m_dead{false}; bool m_ready{false}; bool m_version_sent{false}; std::array m_hdr{}; std::vector m_payload; std::deque> m_send_q; bool m_writing{false}; // Flood guard. int64_t m_window_start{0}; size_t m_window_tx{0}; boost::asio::ip::tcp::endpoint ResolveEndpoint(boost::system::error_code& ec) { if (m_remote.IsIPv4()) { struct in_addr a4; m_remote.GetInAddr(&a4); boost::asio::ip::address_v4::bytes_type b{}; std::memcpy(b.data(), &a4, 4); return {boost::asio::ip::address_v4{b}, m_remote.GetPort()}; } struct in6_addr a6; m_remote.GetIn6Addr(&a6); boost::asio::ip::address_v6::bytes_type b{}; std::memcpy(b.data(), &a6, 16); return {boost::asio::ip::address_v6{b}, m_remote.GetPort()}; } void SendVersion() { if (m_version_sent) return; m_version_sent = true; uint64_t nonce; GetRandBytes({reinterpret_cast(&nonce), sizeof(nonce)}); DataStream ss; ss << PROTOCOL_VERSION; ss << uint64_t{NODE_WITNESS}; ss << GetTime().count(); SerAddr(ss, m_remote); // their address SerAddr(ss, m_remote); // our address (same socket target) ss << nonce; ss << std::string{"/limenka-bridge:1.0.0/"}; ss << int32_t{0}; // start height: we do not sync blocks ss << uint8_t{1}; // relay: yes, we want their mempool SendMessage(NetMsgType::VERSION, ss); } static void SerAddr(DataStream& ss, const CService& svc) { ss << uint64_t{NODE_NONE}; std::array ip{}; if (svc.IsIPv4()) { ip[10] = 0xff; ip[11] = 0xff; struct in_addr a4; svc.GetInAddr(&a4); std::memcpy(ip.data() + 12, &a4, 4); } else { struct in6_addr a6; svc.GetIn6Addr(&a6); std::memcpy(ip.data(), &a6, 16); } ss << ip; ss << uint16_t{svc.GetPort()}; } void SendMessage(const char* type, DataStream& payload) { CMessageHeader hdr{m_net.magic, type, static_cast(payload.size())}; const uint256 checksum = Hash(MakeByteSpan(payload)); std::memcpy(hdr.pchChecksum, checksum.begin(), CMessageHeader::CHECKSUM_SIZE); std::vector out; out.reserve(CMessageHeader::HEADER_SIZE + payload.size()); VectorWriter vw{out, 0, hdr}; vw.write(MakeByteSpan(payload)); auto self{shared_from_this()}; boost::asio::post(*m_bridge.m_io, [self, bytes = std::move(out)]() mutable { if (self->m_dead) return; self->m_send_q.push_back(std::move(bytes)); self->MaybeWrite(); }); } void MaybeWrite() { if (m_writing || m_send_q.empty()) return; m_writing = true; auto self{shared_from_this()}; const auto& bytes = m_send_q.front(); boost::asio::async_write(m_socket, boost::asio::buffer(bytes), [self](const boost::system::error_code& err, size_t) { self->m_writing = false; self->m_send_q.pop_front(); if (err) { self->m_dead = true; return; } self->MaybeWrite(); }); } void DoReadHeader() { auto self{shared_from_this()}; boost::asio::async_read(m_socket, boost::asio::buffer(m_hdr), [self](const boost::system::error_code& err, size_t n) { if (err) { self->m_dead = true; return; } if (n != CMessageHeader::HEADER_SIZE) { self->m_dead = true; return; } if (std::memcmp(self->m_hdr.data(), self->m_net.magic.data(), 4) != 0) { // Wrong network magic on this connection. self->m_dead = true; return; } uint32_t len; std::memcpy(&len, self->m_hdr.data() + CMessageHeader::MESSAGE_SIZE_OFFSET, 4); if (len > MAX_PROTOCOL_MESSAGE_LENGTH) { self->m_dead = true; return; } self->m_payload.resize(len); if (len == 0) { self->DispatchMessage(); self->DoReadHeader(); return; } boost::asio::async_read(self->m_socket, boost::asio::buffer(self->m_payload), [self](const boost::system::error_code& err2, size_t n2) { if (err2 || n2 != self->m_payload.size()) { self->m_dead = true; return; } const uint256 checksum = Hash(MakeByteSpan(self->m_payload)); if (std::memcmp(self->m_hdr.data() + CMessageHeader::CHECKSUM_OFFSET, checksum.begin(), CMessageHeader::CHECKSUM_SIZE) != 0) { self->m_dead = true; return; } self->DispatchMessage(); self->DoReadHeader(); }); }); } void DispatchMessage() { std::string type{ reinterpret_cast(m_hdr.data() + 4), strnlen(reinterpret_cast(m_hdr.data() + 4), CMessageHeader::MESSAGE_TYPE_SIZE)}; try { HandleMessage(type); } catch (const std::exception& e) { LogDebug(BCLog::NET, "mempoolbridge: %s message from %s failed: %s\n", type, m_remote.ToStringAddrPort(), e.what()); } } void HandleMessage(const std::string& type) { if (type == NetMsgType::VERSION) { // BIP339: wtxidrelay must be sent between VERSION and VERACK. // We announce and request by wtxid so peers serve us // full-witness transactions. DataStream empty; SendMessage(NetMsgType::WTXIDRELAY, empty); SendMessage(NetMsgType::VERACK, empty); return; } if (type == NetMsgType::VERACK) { m_ready = true; LogPrintf("mempoolbridge: connected %s peer %s\n", m_net.name, m_remote.ToStringAddrPort()); return; } if (type == NetMsgType::PING) { DataStream out; if (m_payload.size() >= 8) { DataStream ss{m_payload}; uint64_t nonce{0}; ss >> nonce; out << nonce; } SendMessage(NetMsgType::PONG, out); return; } if (type == NetMsgType::PONG) return; if (type == NetMsgType::INV) { HandleInv(); return; } if (type == NetMsgType::TX) { HandleTx(); return; } if (type == NetMsgType::GETDATA) { HandleGetData(); return; } if (type == NetMsgType::GETHEADERS || type == NetMsgType::HEADERS) return; // Everything else (addr, getaddr, feefilter, sendheaders, sendcmpct, // wtxidrelay, notfound, filter*, block announcements) is ignored: // the bridge deals in transaction gossip only. } void HandleInv() { DataStream ss{m_payload}; const uint64_t count = ReadCompactSize(ss); const size_t n = std::min(count, MAX_INV_PER_MESSAGE); std::vector to_fetch; to_fetch.reserve(std::min(n, INVENTORY_BROADCAST_MAX_LOCAL)); for (size_t i = 0; i < n; ++i) { CInv inv; ss >> inv; if (to_fetch.size() >= INVENTORY_BROADCAST_MAX_LOCAL) continue; if (!(inv.type == MSG_TX || inv.type == MSG_WTX || inv.type == MSG_WITNESS_TX)) continue; if (m_bridge.HasTx(inv.hash)) continue; to_fetch.push_back(inv); } if (to_fetch.empty()) return; DataStream out; out << to_fetch; SendMessage(NetMsgType::GETDATA, out); } void HandleTx() { const int64_t now = GetTime().count(); if (now - m_window_start > TX_WINDOW.count()) { m_window_start = now; m_window_tx = 0; } if (++m_window_tx > MAX_TX_PER_WINDOW) return; // flood guard DataStream ss{m_payload}; CMutableTransaction mtx; ss >> TX_WITH_WITNESS(mtx); if (ss.size() != 0) return; // trailing bytes - malformed m_bridge.OnForeignTx(m_net.name, MakeTransactionRef(std::move(mtx))); } void HandleGetData() { DataStream ss{m_payload}; const uint64_t count = ReadCompactSize(ss); const size_t n = std::min(count, MAX_INV_PER_MESSAGE); for (size_t i = 0; i < n; ++i) { CInv inv; ss >> inv; if (inv.type != MSG_TX && inv.type != MSG_WTX && inv.type != MSG_WITNESS_TX) continue; if (auto tx = m_bridge.GetTx(inv.hash)) SendTx(tx); } } }; // --------------------------------------------------------------------------- // Bridge // --------------------------------------------------------------------------- MempoolBridge::MempoolBridge(node::NodeContext& node) : m_node{node} {} MempoolBridge::~MempoolBridge() { Stop(); } std::optional> MempoolBridge::ParseBridgePeer(const std::string& spec) { const size_t c1 = spec.find(':'); if (c1 == std::string::npos) return std::nullopt; const std::string net_name = spec.substr(0, c1); std::string hostport = spec.substr(c1 + 1); if (hostport.empty()) return std::nullopt; const auto it = std::find_if(KNOWN_NETWORKS.begin(), KNOWN_NETWORKS.end(), [&](const BridgeNetwork& n) { return n.name == net_name; }); if (it == KNOWN_NETWORKS.end()) return std::nullopt; // If no port and no brackets, append the network default. const size_t c2 = hostport.rfind(':'); if (c2 == std::string::npos && hostport.find(']') == std::string::npos) { hostport = strprintf("%s:%u", hostport, it->default_port); } const auto svc = Lookup(hostport, it->default_port, /*fAllowLookup=*/true); if (!svc) return std::nullopt; return std::make_pair(net_name, *svc); } bool MempoolBridge::Start(const std::vector& peer_specs) { if (m_running.exchange(true)) return true; { LOCK(m_net_mutex); m_networks = KNOWN_NETWORKS; for (const auto& spec : peer_specs) { const auto parsed = ParseBridgePeer(spec); if (!parsed) { LogPrintf("mempoolbridge: ignoring invalid -bridgepeer=%s\n", spec); continue; } const auto it = std::find_if(m_networks.begin(), m_networks.end(), [&](const BridgeNetwork& n) { return n.name == parsed->first; }); if (it != m_networks.end()) it->peers.push_back(parsed->second); } m_networks.erase(std::remove_if(m_networks.begin(), m_networks.end(), [](const BridgeNetwork& n) { return n.peers.empty(); }), m_networks.end()); } if (m_networks.empty()) { m_running = false; return false; } m_io = std::make_unique(); if (m_node.validation_signals) { m_node.validation_signals->RegisterSharedValidationInterface(shared_from_this()); } m_thread = std::thread([this]() { while (m_running) { MaintainConnections(); // Park on a timer so a failed connect batch cannot spin the // loop (run() returns only when the timer fires or Stop()). boost::asio::steady_timer timer{*m_io, RECONNECT_INTERVAL}; timer.async_wait([](const boost::system::error_code&) {}); try { m_io->restart(); m_io->run(); } catch (const std::exception& e) { LogPrintf("mempoolbridge: io loop error: %s\n", e.what()); break; } } }); LogPrintf("mempoolbridge: started with %d network(s)\n", int(m_networks.size())); return true; } void MempoolBridge::Stop() { if (!m_running.exchange(false)) return; if (m_io) m_io->stop(); if (m_thread.joinable()) m_thread.join(); { LOCK(m_net_mutex); for (auto& peer : m_peers) peer->Disconnect(); m_peers.clear(); } if (m_node.validation_signals) { m_node.validation_signals->UnregisterSharedValidationInterface(shared_from_this()); } LogPrintf("mempoolbridge: stopped\n"); } void MempoolBridge::TransactionAddedToMempool(const NewMempoolTransactionInfo& info, uint64_t mempool_sequence) { if (!m_io) return; const CTransactionRef tx = info.info.m_tx; const uint256 wtxid = tx->GetWitnessHash(); { // Serve this tx to foreign peers by txid or wtxid (BIP339). LOCK(m_cache_mutex); if (m_tx_cache.size() > TX_CACHE_MAX) m_tx_cache.clear(); m_tx_cache[tx->GetHash()] = tx; m_tx_cache[wtxid] = tx; } std::vector nets; { LOCK(m_net_mutex); nets = m_networks; } auto self{shared_from_this()}; for (const auto& net : nets) { if (!NoteSeen(net.name, wtxid)) continue; boost::asio::post(*m_io, [self, tx, wtxid, net]() { self->AnnounceToNetwork(net, tx, wtxid); }); } } void MempoolBridge::AnnounceToNetwork(const BridgeNetwork& net, const CTransactionRef& tx, const uint256& wtxid) { std::vector> peers; { LOCK(m_net_mutex); for (const auto& peer : m_peers) { if (peer->NetworkName() == net.name && peer->Ready() && !peer->Dead()) { peers.push_back(peer); } } } if (peers.empty()) return; LogDebug(BCLog::MEMPOOL, "mempoolbridge: announcing %s to %s (%d peer(s))\n", tx->GetHash().ToString(), net.name, int(peers.size())); for (const auto& peer : peers) peer->SendInv(tx); } bool MempoolBridge::NoteSeen(const std::string& network, const uint256& wtxid) { LOCK(m_seen_mutex); auto& set = m_seen[network]; if (set.size() > SEEN_MAX) set.clear(); return set.insert(wtxid).second; } void MempoolBridge::OnForeignTx(const std::string& from_network, CTransactionRef tx) { const uint256 txid = tx->GetHash(); const uint256 wtxid = tx->GetWitnessHash(); // Never echo a transaction back to the network it arrived from, and // remember it so the mempool announcement does not re-send it there. NoteSeen(from_network, wtxid); { LOCK(m_cache_mutex); if (m_tx_cache.size() > TX_CACHE_MAX) m_tx_cache.clear(); m_tx_cache[txid] = tx; m_tx_cache[wtxid] = tx; } ChainstateManager* chainman = m_node.chainman.get(); if (!chainman || !chainman->ActiveChainstate().GetMempool()) return; LogDebug(BCLog::MEMPOOL, "mempoolbridge: tx %s from %s, validating\n", txid.ToString(), from_network); { LOCK(::cs_main); const MempoolAcceptResult result = AcceptToMemoryPool(chainman->ActiveChainstate(), tx, GetTime().count(), /*bypass_limits=*/false, /*test_accept=*/false); if (result.m_result_type == MempoolAcceptResult::ResultType::VALID) { LogDebug(BCLog::MEMPOOL, "mempoolbridge: tx %s from %s accepted into limenka mempool\n", txid.ToString(), from_network); } else { LogDebug(BCLog::MEMPOOL, "mempoolbridge: tx %s from %s rejected: %s\n", txid.ToString(), from_network, result.m_state.ToString()); } } } bool MempoolBridge::HasTx(const uint256& hash) const { { LOCK(m_cache_mutex); if (m_tx_cache.count(hash)) return true; } ChainstateManager* chainman = m_node.chainman.get(); if (!chainman || !chainman->ActiveChainstate().GetMempool()) return false; CTxMemPool& pool = *chainman->ActiveChainstate().GetMempool(); LOCK(pool.cs); return pool.GetEntry(Txid::FromUint256(hash)) != nullptr; } CTransactionRef MempoolBridge::GetTx(const uint256& hash) const { { LOCK(m_cache_mutex); const auto it = m_tx_cache.find(hash); if (it != m_tx_cache.end()) return it->second; } ChainstateManager* chainman = m_node.chainman.get(); if (!chainman || !chainman->ActiveChainstate().GetMempool()) return nullptr; CTxMemPool& pool = *chainman->ActiveChainstate().GetMempool(); LOCK(pool.cs); const CTxMemPoolEntry* entry = pool.GetEntry(Txid::FromUint256(hash)); if (!entry) return nullptr; return entry->GetSharedTx(); } void MempoolBridge::MaintainConnections() { if (!m_io) return; LOCK(m_net_mutex); // Drop dead peers. for (auto it = m_peers.begin(); it != m_peers.end();) { if ((*it)->Dead()) { it = m_peers.erase(it); } else { ++it; } } // One live connection per configured peer address. for (const auto& net : m_networks) { for (const auto& addr : net.peers) { const bool exists = std::any_of(m_peers.begin(), m_peers.end(), [&](const std::shared_ptr& p) { return p->NetworkName() == net.name && p->Remote() == addr && !p->Dead(); }); if (exists) continue; LogPrintf("mempoolbridge: connecting to %s %s\n", net.name, addr.ToStringAddrPort()); auto peer = std::make_shared(*this, net, addr); m_peers.push_back(peer); peer->Start(); } } } } // namespace node