wtf... its getting worse. Lets see if we can simplify it a lot

This commit is contained in:
Igor Katson 2023-11-29 14:48:22 +00:00
parent 6518dc6eff
commit 826d1b8f1d
No known key found for this signature in database
GPG key ID: B4EC22B66D61A3F5

View file

@ -80,6 +80,13 @@ fn make_rate_limiter() -> RateLimiter {
.build() .build()
} }
struct InfoHashMeta {
seen_peers: IndexSet<SocketAddr>,
subscriber: tokio::sync::broadcast::Sender<SocketAddr>,
closest_responding_nodes: Vec<MaybeUsefulNode>,
join_handle: tokio::task::JoinHandle<()>,
}
pub struct DhtState { pub struct DhtState {
id: Id20, id: Id20,
next_transaction_id: AtomicU16, next_transaction_id: AtomicU16,
@ -98,10 +105,8 @@ pub struct DhtState {
rate_limiter: RateLimiter, rate_limiter: RateLimiter,
sender: UnboundedSender<WorkerSendRequest>, sender: UnboundedSender<WorkerSendRequest>,
seen_peers: DashMap<Id20, IndexSet<SocketAddr>>, // Per-torrent stats.
info_hash_meta: DashMap<Id20, InfoHashMeta>,
closest_responding_nodes_for_info_hash: DashMap<Id20, Vec<MaybeUsefulNode>>,
get_peers_subscribers: DashMap<Id20, tokio::sync::broadcast::Sender<SocketAddr>>,
} }
impl DhtState { impl DhtState {
@ -119,10 +124,8 @@ impl DhtState {
routing_table: RwLock::new(routing_table), routing_table: RwLock::new(routing_table),
sender, sender,
listen_addr, listen_addr,
seen_peers: Default::default(),
get_peers_subscribers: Default::default(),
rate_limiter: make_rate_limiter(), rate_limiter: make_rate_limiter(),
closest_responding_nodes_for_info_hash: Default::default(), info_hash_meta: Default::default(),
recent_requests: Default::default(), recent_requests: Default::default(),
} }
} }
@ -318,8 +321,8 @@ impl DhtState {
Ok(()) Ok(())
} }
MessageKind::GetPeersRequest(req) => { MessageKind::GetPeersRequest(req) => {
let peers = self.seen_peers.get(&req.info_hash).map(|peers| { let peers = self.info_hash_meta.get(&req.info_hash).map(|meta| {
peers meta.seen_peers
.iter() .iter()
.copied() .copied()
.filter_map(|a| match a { .filter_map(|a| match a {
@ -384,7 +387,11 @@ impl DhtState {
DhtStats { DhtStats {
id: self.id, id: self.id,
outstanding_requests: self.inflight_by_transaction_id.len(), outstanding_requests: self.inflight_by_transaction_id.len(),
seen_peers: self.seen_peers.iter().map(|e| e.value().len()).sum(), seen_peers: self
.info_hash_meta
.iter()
.map(|e| e.value().seen_peers.len())
.sum(),
recent_requests: self.recent_requests.len(), recent_requests: self.recent_requests.len(),
routing_table_size: self.routing_table.read().len(), routing_table_size: self.routing_table.read().len(),
} }
@ -399,36 +406,35 @@ impl DhtState {
tokio::sync::broadcast::Receiver<SocketAddr>, tokio::sync::broadcast::Receiver<SocketAddr>,
)> { )> {
use dashmap::mapref::entry::Entry; use dashmap::mapref::entry::Entry;
match self.get_peers_subscribers.entry(info_hash) { match self.info_hash_meta.entry(info_hash) {
Entry::Occupied(o) => { Entry::Occupied(o) => {
let pos = self.seen_peers.get(&info_hash).and_then(|p| { let seen_peers = &o.get().seen_peers;
if p.is_empty() { let pos = if seen_peers.is_empty() {
None None
} else { } else {
Some((0, p.len())) Some((0, seen_peers.len()))
} };
}); let rx = o.get().subscriber.subscribe();
let rx = o.get().subscribe();
Ok((pos, rx)) Ok((pos, rx))
} }
Entry::Vacant(v) => { Entry::Vacant(v) => {
// DHT sends peers REALLY fast, so ideally the consumer of this broadcast should not lag behind. // DHT sends peers REALLY fast, so ideally the consumer of this broadcast should not lag behind.
// In case it does though we have PeerStream to replay. // In case it does though we have PeerStream to replay.
let (tx, rx) = tokio::sync::broadcast::channel(100);
v.insert(tx);
let this = self.clone(); let this = self.clone();
spawn( let join_handle = spawn(
error_span!("peers_requester", info_hash = format!("{:?}", info_hash)), error_span!("peers_requester", info_hash = format!("{:?}", info_hash)),
async move { async move {
let mut iteration = 0usize; let mut iteration = 0usize;
loop { loop {
if !this.get_peers_subscribers.contains_key(&info_hash) { let meta = match this.info_hash_meta.get(&info_hash) {
debug!("no more subscribers, closing peers_requester"); Some(meta) => meta,
return Ok(()); None => {
} debug!("no more subscribers, closing peers_requester");
return Ok(());
}
};
trace!("iteration {iteration}"); trace!("iteration {iteration}");
// We don't need to allocate/collect here, but the borrow checker is not happy otherwise.
let nodes_to_query = this let nodes_to_query = this
.routing_table .routing_table
.read() .read()
@ -440,19 +446,26 @@ impl DhtState {
for (id, addr) in nodes_to_query { for (id, addr) in nodes_to_query {
this.send_find_peers_if_not_yet(info_hash, id, addr)?; this.send_find_peers_if_not_yet(info_hash, id, addr)?;
} }
if let Some(e) = for MaybeUsefulNode { id, addr, .. } in
this.closest_responding_nodes_for_info_hash.get(&info_hash) meta.closest_responding_nodes.iter()
{ {
for MaybeUsefulNode { id, addr, .. } in e.value().iter() { this.send_find_peers_if_not_yet(info_hash, *id, *addr)?;
this.send_find_peers_if_not_yet(info_hash, *id, *addr)?;
}
} }
drop(meta);
tokio::time::sleep(REQUERY_INTERVAL).await; tokio::time::sleep(REQUERY_INTERVAL).await;
iteration += 1; iteration += 1;
} }
}, },
); );
let (tx, rx) = tokio::sync::broadcast::channel(100);
v.insert(InfoHashMeta {
seen_peers: Default::default(),
subscriber: tx,
closest_responding_nodes: Default::default(),
join_handle,
});
Ok((None, rx)) Ok((None, rx))
} }
} }
@ -559,10 +572,8 @@ impl DhtState {
target: Id20, target: Id20,
nodes: CompactNodeInfo, nodes: CompactNodeInfo,
) -> anyhow::Result<()> { ) -> anyhow::Result<()> {
// We don't need to allocate/collect here, but the borrow checker is not happy
// otherwise when we iterate self.searching_for_peers and mutating self in the loop.
let searching_for_peers = self let searching_for_peers = self
.get_peers_subscribers .info_hash_meta
.iter() .iter()
.map(|e| *e.key()) .map(|e| *e.key())
.collect::<Vec<_>>(); .collect::<Vec<_>>();
@ -596,41 +607,29 @@ impl DhtState {
info_hash: Id20, info_hash: Id20,
node_id: Id20, node_id: Id20,
addr: SocketAddr, addr: SocketAddr,
closest_nodes: &mut Vec<MaybeUsefulNode>,
) -> bool { ) -> bool {
use dashmap::mapref::entry::Entry; closest_nodes.push(MaybeUsefulNode {
let n = MaybeUsefulNode {
id: node_id, id: node_id,
addr, addr,
last_response: None, last_response: None,
returned_peers: false, returned_peers: false,
}; });
match self.closest_responding_nodes_for_info_hash.entry(info_hash) {
Entry::Occupied(mut occ) => {
// How many nodes to query per torrent.
const LIMIT: usize = 256;
let v = occ.get_mut();
v.push(n);
v.sort_by_key(|n| {
let has_returned_peers_desc = Reverse(n.returned_peers);
let has_responded_desc = Reverse(n.last_response.is_some() as u8);
let distance = n.id.distance(&info_hash);
(has_returned_peers_desc, has_responded_desc, distance)
});
if v.len() > LIMIT {
let popped = v.pop().unwrap();
if popped.id == node_id {
return false;
}
}
true const LIMIT: usize = 256;
} closest_nodes.sort_by_key(|n| {
Entry::Vacant(v) => { let has_returned_peers_desc = Reverse(n.returned_peers);
v.insert(vec![n]); let has_responded_desc = Reverse(n.last_response.is_some() as u8);
true let distance = n.id.distance(&info_hash);
(has_returned_peers_desc, has_responded_desc, distance)
});
if closest_nodes.len() > LIMIT {
let popped = closest_nodes.pop().unwrap();
if popped.id == node_id {
return false;
} }
} }
true
} }
fn on_found_peers_or_nodes( fn on_found_peers_or_nodes(
@ -643,72 +642,54 @@ impl DhtState {
self.routing_table_add_node(source, source_addr); self.routing_table_add_node(source, source_addr);
use dashmap::mapref::entry::Entry; use dashmap::mapref::entry::Entry;
let bsender = match self.get_peers_subscribers.entry(info_hash) { let mut meta = match self.info_hash_meta.entry(info_hash) {
Entry::Occupied(o) => o, Entry::Occupied(o) => o,
Entry::Vacant(_) => { Entry::Vacant(_) => {
warn!( warn!(
"ignoring get_peers response, no subscribers for {:?}", "ignoring found_peers response, no subscribers for {:?}",
info_hash info_hash
); );
return Ok(()); return Ok(());
} }
}; };
let meta_mut = meta.get_mut();
{ {
let n = MaybeUsefulNode { let now = Some(Instant::now());
id: source, let returned_peers = data.values.as_ref().map(|p| !p.is_empty()).unwrap_or(false);
addr: source_addr,
last_response: Some(Instant::now()), if let Some(existing_useful_node) = meta_mut
returned_peers: data.values.as_ref().map(|p| !p.is_empty()).unwrap_or(false), .closest_responding_nodes
}; .iter_mut()
match self.closest_responding_nodes_for_info_hash.entry(info_hash) { .find(|n| n.id == source && n.addr == source_addr)
Entry::Occupied(mut useful_nodes) => { {
if let Some(useful_node) = useful_nodes existing_useful_node.last_response = now;
.get_mut() existing_useful_node.returned_peers |= returned_peers;
.iter_mut() } else {
.find(|n| n.id == source && n.addr == source_addr) meta_mut.closest_responding_nodes.push(MaybeUsefulNode {
{ id: source,
useful_node.last_response = Some(Instant::now()); addr: source_addr,
} else { last_response: now,
useful_nodes.get_mut().push(n); returned_peers,
} });
} }
Entry::Vacant(v) => {
v.insert(vec![n]);
}
};
} }
if let Some(peers) = data.values { if let Some(peers) = data.values {
let mut seen = self.seen_peers.entry(info_hash).or_default();
for peer in peers.iter() { for peer in peers.iter() {
if peer.addr.port() < 1024 { if peer.addr.port() < 1024 {
debug!("bad peer port, ignoring: {}", peer.addr); debug!("bad peer port, ignoring: {}", peer.addr);
continue; continue;
} }
let addr = SocketAddr::V4(peer.addr); let addr = SocketAddr::V4(peer.addr);
if seen.insert(addr) { if meta_mut.seen_peers.insert(addr) {
match bsender.get().send(addr) { match meta_mut.subscriber.send(addr) {
Ok(_) => {} Ok(_) => {}
Err(_) => { Err(_) => {
debug!("no more subscribers for {:?}, cleaning up", info_hash); debug!("no more subscribers for {:?}, cleaning up", info_hash);
// bsender.remove(); meta_mut.join_handle.abort();
meta.remove();
// let this = self.clone();
// spawn(
// error_span!("cleanup", info_hash = format!("{info_hash:?}")),
// async move {
// tokio::time::sleep(Duration::from_secs(10)).await;
// if !this.get_peers_subscribers.contains_key(&info_hash) {
// debug!("no more subscribers for {:?}, removed it from seen peers", info_hash);
// this.seen_peers.remove(&info_hash);
// this.closest_responding_nodes_for_info_hash
// .remove(&info_hash);
// }
// Ok(())
// },
// );
return Ok(()); return Ok(());
} }
} }
@ -721,6 +702,7 @@ impl DhtState {
info_hash, info_hash,
node.id, node.id,
node.addr.into(), node.addr.into(),
&mut meta_mut.closest_responding_nodes,
) { ) {
self.routing_table_add_node(node.id, node.addr.into()); self.routing_table_add_node(node.id, node.addr.into());
self.send_find_peers_if_not_yet(info_hash, node.id, node.addr.into())?; self.send_find_peers_if_not_yet(info_hash, node.id, node.addr.into())?;
@ -984,13 +966,15 @@ impl Stream for PeerStream {
) -> Poll<Option<Self::Item>> { ) -> Poll<Option<Self::Item>> {
loop { loop {
if let Some((pos, end)) = self.initial_peers_pos.take() { if let Some((pos, end)) = self.initial_peers_pos.take() {
let addr = *self let addr = match self
.state .state
.seen_peers .info_hash_meta
.get(&self.info_hash) .get(&self.info_hash)
.unwrap() .and_then(|meta| meta.seen_peers.get_index(pos).copied())
.get_index(pos) {
.unwrap(); Some(addr) => addr,
None => return Poll::Ready(None),
};
if pos + 1 < end { if pos + 1 < end {
self.initial_peers_pos = Some((pos + 1, end)); self.initial_peers_pos = Some((pos + 1, end));
} }