Branch data Line data Source code
1 : : // Copyright (c) 2023-present The Bitcoin Core developers
2 : : // Distributed under the MIT software license, see the accompanying
3 : : // file COPYING or https://opensource.org/license/mit/.
4 : :
5 : : #include <private_broadcast.h>
6 : :
7 : : #include <util/check.h>
8 : :
9 : : #include <algorithm>
10 : : #include <ranges>
11 : :
12 : :
13 : 10021 : PrivateBroadcast::AddResult PrivateBroadcast::Add(const CTransactionRef& tx)
14 : : EXCLUSIVE_LOCKS_REQUIRED(!m_mutex)
15 : : {
16 : 10021 : LOCK(m_mutex);
17 : : // Cleanup finished transactions
18 [ + - ]: 10021 : std::erase_if(m_transactions, [this](const auto& entry) {
19 : 50075007 : const auto& state{entry.second};
20 [ + + + + : 50075008 : return state.resolved && !IsPending(state) &&
+ - ]
21 [ - + ]: 2 : std::ranges::all_of(state.send_statuses, [](const auto& status) { return status.disconnected; });
22 : : });
23 : :
24 [ + + ]: 10021 : if (const auto it{m_transactions.find(tx)}; it != m_transactions.end()) {
25 [ - + + + ]: 7 : if (it->second.send_statuses.size() < m_max_send_attempts) return AddResult::AlreadyPresent;
26 : :
27 : : // A transaction that has reached m_max_send_attempts can be explicitly retried by adding it again.
28 : 2 : it->second.time_added = NodeClock::now();
29 : 2 : it->second.send_statuses.clear();
30 : 2 : it->second.planned_sends = INITIAL_CONNECTION_COUNT;
31 : 2 : it->second.resolved = false;
32 : 2 : return AddResult::Added;
33 : : }
34 : :
35 [ + + ]: 10014 : if (m_transactions.size() >= m_max_transactions) return AddResult::QueueFull;
36 : :
37 [ + - ]: 10009 : m_transactions.try_emplace(tx);
38 : : return AddResult::Added;
39 : 10021 : }
40 : :
41 : 7 : std::optional<size_t> PrivateBroadcast::Remove(const CTransactionRef& tx)
42 : : EXCLUSIVE_LOCKS_REQUIRED(!m_mutex)
43 : : {
44 : 7 : LOCK(m_mutex);
45 : 7 : const auto handle{m_transactions.extract(tx)};
46 [ + + ]: 7 : if (handle) {
47 [ + + ]: 5 : const auto& state{handle.mapped()};
48 [ + + ]: 5 : const size_t planned{std::min(state.planned_sends, m_max_send_attempts)};
49 [ - + ]: 5 : const size_t actual{state.send_statuses.size()};
50 [ + - ]: 5 : return planned > actual ? planned - actual : 0;
51 : : }
52 : 2 : return std::nullopt;
53 [ + - ]: 7 : }
54 : :
55 : 8 : bool PrivateBroadcast::MarkResolved(const CTransactionRef& tx)
56 : : {
57 : 8 : LOCK(m_mutex);
58 : 8 : const auto it{m_transactions.find(tx)};
59 [ + + + - ]: 9 : if (it == m_transactions.end()) return false;
60 : 1 : it->second.resolved = true;
61 : 1 : return true;
62 : 8 : }
63 : :
64 : 4 : void PrivateBroadcast::NodeDisconnected(NodeId nodeid)
65 : : {
66 : 4 : LOCK(m_mutex);
67 [ + - + + ]: 4 : if (const auto entry{GetSendStatusByNode(nodeid)}) {
68 : 1 : entry->send_status.disconnected = true;
69 : : }
70 : 4 : }
71 : :
72 : 3 : bool PrivateBroadcast::TryGrantRetry(const CTransactionRef& tx)
73 : : {
74 : 3 : LOCK(m_mutex);
75 : 3 : const auto it{m_transactions.find(tx)};
76 [ + - + - ]: 3 : if (it == m_transactions.end() ||
77 [ + - + - ]: 6 : it->second.resolved ||
78 [ + - ]: 3 : IsPending(it->second) ||
79 [ + + ]: 3 : it->second.planned_sends >= m_max_send_attempts) return false;
80 : 2 : ++it->second.planned_sends;
81 : 2 : return true;
82 : 3 : }
83 : :
84 : 17 : std::optional<CTransactionRef> PrivateBroadcast::PickTxForSend(const NodeId& will_send_to_nodeid, const CService& will_send_to_address)
85 : : EXCLUSIVE_LOCKS_REQUIRED(!m_mutex)
86 : : {
87 : 17 : LOCK(m_mutex);
88 : :
89 [ + - - + ]: 17 : if (GetSendStatusByNode(will_send_to_nodeid).has_value()) { // nodeid reuse, shouldn't send >1 tx to a given node
90 : 0 : Assume(false);
91 : 0 : return std::nullopt;
92 : : }
93 : :
94 [ + - ]: 36 : auto pending_transactions{m_transactions | std::views::filter([this](const auto& entry) { return IsPending(entry.second); })};
95 [ + - ]: 17 : const auto it{std::ranges::max_element(
96 : : pending_transactions,
97 : 3 : [](const auto& a, const auto& b) { return a < b; },
98 : 3 : [](const auto& el) { return DerivePriority(el.second.send_statuses); })};
99 : :
100 [ + + ]: 17 : if (it != pending_transactions.end()) {
101 : 13 : auto& [tx, state]{*it};
102 [ + - ]: 13 : state.send_statuses.emplace_back(will_send_to_nodeid, will_send_to_address, NodeClock::now());
103 : 13 : return tx;
104 : : }
105 : :
106 : 4 : return std::nullopt;
107 : 17 : }
108 : :
109 : 5 : std::optional<CTransactionRef> PrivateBroadcast::GetTxForNode(const NodeId& nodeid)
110 : : EXCLUSIVE_LOCKS_REQUIRED(!m_mutex)
111 : : {
112 : 5 : LOCK(m_mutex);
113 [ + - ]: 5 : const auto tx_and_status{GetSendStatusByNode(nodeid)};
114 [ + + ]: 5 : if (tx_and_status.has_value()) {
115 : 3 : return tx_and_status.value().tx;
116 : : }
117 : 2 : return std::nullopt;
118 : 5 : }
119 : :
120 : 3 : void PrivateBroadcast::NodeConfirmedReception(const NodeId& nodeid)
121 : : EXCLUSIVE_LOCKS_REQUIRED(!m_mutex)
122 : : {
123 : 3 : LOCK(m_mutex);
124 [ + - ]: 3 : const auto tx_and_status{GetSendStatusByNode(nodeid)};
125 [ + + ]: 3 : if (tx_and_status.has_value()) {
126 [ - + + - ]: 5 : tx_and_status.value().send_status.confirmed = NodeClock::now();
127 : : }
128 : 3 : }
129 : :
130 : 14 : bool PrivateBroadcast::HavePendingTransactions()
131 : : EXCLUSIVE_LOCKS_REQUIRED(!m_mutex)
132 : : {
133 : 14 : LOCK(m_mutex);
134 [ + - + - ]: 24 : return std::ranges::any_of(m_transactions, [this](const auto& entry) { return IsPending(entry.second); });
135 : 14 : }
136 : :
137 : 10 : std::vector<CTransactionRef> PrivateBroadcast::GetStale() const
138 : : EXCLUSIVE_LOCKS_REQUIRED(!m_mutex)
139 : : {
140 : 10 : LOCK(m_mutex);
141 : 10 : const auto now{NodeClock::now()};
142 : 10 : std::vector<CTransactionRef> stale;
143 [ + + + - ]: 24 : for (const auto& [tx, state] : m_transactions) {
144 [ + + + - ]: 14 : if (state.resolved || state.planned_sends >= m_max_send_attempts) continue;
145 [ + - ]: 13 : const Priority p{DerivePriority(state.send_statuses)};
146 [ + + ]: 13 : if (p.num_confirmed == 0) {
147 [ + + + - ]: 11 : if (state.time_added < now - INITIAL_STALE_DURATION) stale.push_back(tx);
148 : : } else {
149 [ + + + - ]: 2 : if (p.last_confirmed < now - STALE_DURATION) stale.push_back(tx);
150 : : }
151 : : }
152 [ + - ]: 10 : return stale;
153 : 10 : }
154 : :
155 : 18 : std::vector<PrivateBroadcast::TxBroadcastInfo> PrivateBroadcast::GetBroadcastInfo() const
156 : : EXCLUSIVE_LOCKS_REQUIRED(!m_mutex)
157 : : {
158 : 18 : LOCK(m_mutex);
159 : 18 : std::vector<TxBroadcastInfo> entries;
160 [ + - ]: 18 : entries.reserve(m_transactions.size());
161 : :
162 [ + + + + ]: 70031 : for (const auto& [tx, state] : m_transactions) {
163 [ + + ]: 70013 : if (state.resolved) continue;
164 : 70012 : std::vector<PeerSendInfo> peers;
165 [ - + + - ]: 70012 : peers.reserve(state.send_statuses.size());
166 [ + + ]: 70025 : for (const auto& status : state.send_statuses) {
167 [ + - ]: 26 : peers.emplace_back(PeerSendInfo{.address = status.address, .sent = status.picked, .received = status.confirmed});
168 : : }
169 [ - + + - ]: 70012 : const size_t attempts_remaining{m_max_send_attempts - std::min(state.send_statuses.size(), m_max_send_attempts)};
170 [ + - + - ]: 140024 : entries.emplace_back(TxBroadcastInfo{.tx = tx, .time_added = state.time_added, .attempts_remaining = attempts_remaining, .peers = std::move(peers)});
171 : 70012 : }
172 : :
173 [ + - ]: 18 : return entries;
174 : 18 : }
175 : :
176 : 34 : bool PrivateBroadcast::IsPending(const TxSendStatus& status) const
177 : : {
178 : : // Deliberately ignore resolved so all initially scheduled connections can complete.
179 [ + + ]: 34 : const size_t limit{std::min(status.planned_sends, m_max_send_attempts)};
180 [ - + ]: 34 : return status.send_statuses.size() < limit;
181 : : }
182 : :
183 : 19 : PrivateBroadcast::Priority PrivateBroadcast::DerivePriority(const std::vector<SendStatus>& sent_to)
184 : : {
185 : 19 : Priority p;
186 [ - + ]: 19 : p.num_picked = sent_to.size();
187 [ + + ]: 30 : for (const auto& send_status : sent_to) {
188 : 11 : p.last_picked = std::max(p.last_picked, send_status.picked);
189 [ + + ]: 11 : if (send_status.confirmed.has_value()) {
190 : 2 : ++p.num_confirmed;
191 : 2 : p.last_confirmed = std::max(p.last_confirmed, send_status.confirmed.value());
192 : : }
193 : : }
194 : 19 : return p;
195 : : }
196 : :
197 : 29 : std::optional<PrivateBroadcast::TxAndSendStatusForNode> PrivateBroadcast::GetSendStatusByNode(const NodeId& nodeid)
198 : : EXCLUSIVE_LOCKS_REQUIRED(m_mutex)
199 : : {
200 : 29 : AssertLockHeld(m_mutex);
201 [ + + ]: 54 : for (auto& [tx, state] : m_transactions) {
202 [ + + ]: 64 : for (auto& send_status : state.send_statuses) {
203 [ + + ]: 39 : if (send_status.nodeid == nodeid) {
204 : 6 : return TxAndSendStatusForNode{.tx = tx, .send_status = send_status};
205 : : }
206 : : }
207 : : }
208 : 23 : return std::nullopt;
209 : : }
|