DHT announce compiles
This commit is contained in:
parent
80df2c1001
commit
162afe3056
1 changed files with 84 additions and 12 deletions
|
|
@ -1,6 +1,7 @@
|
||||||
use std::{
|
use std::{
|
||||||
cmp::Reverse,
|
cmp::Reverse,
|
||||||
net::SocketAddr,
|
net::{SocketAddr, SocketAddrV4},
|
||||||
|
str::FromStr,
|
||||||
sync::{
|
sync::{
|
||||||
atomic::{AtomicU16, Ordering},
|
atomic::{AtomicU16, Ordering},
|
||||||
Arc,
|
Arc,
|
||||||
|
|
@ -11,8 +12,8 @@ use std::{
|
||||||
|
|
||||||
use crate::{
|
use crate::{
|
||||||
bprotocol::{
|
bprotocol::{
|
||||||
self, CompactNodeInfo, ErrorDescription, FindNodeRequest, GetPeersRequest, Message,
|
self, AnnouncePeer, CompactNodeInfo, ErrorDescription, FindNodeRequest, GetPeersRequest,
|
||||||
MessageKind, Node, PingRequest, Response,
|
Message, MessageKind, Node, PingRequest, Response,
|
||||||
},
|
},
|
||||||
peer_store::PeerStore,
|
peer_store::PeerStore,
|
||||||
routing_table::{InsertResult, NodeStatus, RoutingTable},
|
routing_table::{InsertResult, NodeStatus, RoutingTable},
|
||||||
|
|
@ -93,17 +94,52 @@ trait RecursiveRequestCallbacks: Sized + Send + Sync + 'static {
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
struct RecursiveRequestCallbacksGetPeers {}
|
struct RecursiveRequestCallbacksGetPeers {
|
||||||
|
// Id20::from_str("00000fffffffffffffffffffffffffffffffffff").unwrap()
|
||||||
|
min_distance_to_announce: Id20,
|
||||||
|
}
|
||||||
|
|
||||||
impl RecursiveRequestCallbacks for RecursiveRequestCallbacksGetPeers {
|
impl RecursiveRequestCallbacks for RecursiveRequestCallbacksGetPeers {
|
||||||
fn on_request_start(&self, _: &RecursiveRequest<Self>, _: Id20, _: SocketAddr) {}
|
fn on_request_start(&self, _: &RecursiveRequest<Self>, _: Id20, _: SocketAddr) {}
|
||||||
|
|
||||||
fn on_request_end(
|
fn on_request_end(
|
||||||
&self,
|
&self,
|
||||||
_: &RecursiveRequest<Self>,
|
req: &RecursiveRequest<Self>,
|
||||||
_: Id20,
|
target_node: Id20,
|
||||||
_: SocketAddr,
|
addr: SocketAddr,
|
||||||
_: &anyhow::Result<ResponseOrError>,
|
resp: &anyhow::Result<ResponseOrError>,
|
||||||
) {
|
) {
|
||||||
|
let announce_addr = match req.dht.announce_addr {
|
||||||
|
Some(a) => a,
|
||||||
|
None => return,
|
||||||
|
};
|
||||||
|
let resp = match resp {
|
||||||
|
Ok(ResponseOrError::Response(resp)) => resp,
|
||||||
|
_ => return,
|
||||||
|
};
|
||||||
|
let token = match &resp.token {
|
||||||
|
Some(token) => token,
|
||||||
|
None => return,
|
||||||
|
};
|
||||||
|
if req.info_hash.distance(&target_node) > self.min_distance_to_announce {
|
||||||
|
trace!(
|
||||||
|
"not announcing, {:?} is too far from {:?}",
|
||||||
|
target_node,
|
||||||
|
req.info_hash
|
||||||
|
);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
let (tid, message) = req.dht.create_request(Request::Announce {
|
||||||
|
info_hash: req.info_hash,
|
||||||
|
token: token.clone(),
|
||||||
|
addr: announce_addr,
|
||||||
|
});
|
||||||
|
|
||||||
|
let _ = req.dht.worker_sender.send(WorkerSendRequest {
|
||||||
|
our_tid: Some(tid),
|
||||||
|
message,
|
||||||
|
addr,
|
||||||
|
});
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -165,7 +201,12 @@ impl RequestPeersStream {
|
||||||
useful_nodes: RwLock::new(Vec::new()),
|
useful_nodes: RwLock::new(Vec::new()),
|
||||||
peer_tx,
|
peer_tx,
|
||||||
node_tx,
|
node_tx,
|
||||||
callbacks: RecursiveRequestCallbacksGetPeers {},
|
callbacks: RecursiveRequestCallbacksGetPeers {
|
||||||
|
min_distance_to_announce: Id20::from_str(
|
||||||
|
"000000ffffffffffffffffffffffffffffffffff",
|
||||||
|
)
|
||||||
|
.unwrap(),
|
||||||
|
},
|
||||||
});
|
});
|
||||||
let join_handle = rp.request_peers_forever(node_rx);
|
let join_handle = rp.request_peers_forever(node_rx);
|
||||||
Self {
|
Self {
|
||||||
|
|
@ -351,7 +392,7 @@ impl<C: RecursiveRequestCallbacks> RecursiveRequest<C> {
|
||||||
self.callbacks.on_request_start(self, id, addr);
|
self.callbacks.on_request_start(self, id, addr);
|
||||||
}
|
}
|
||||||
|
|
||||||
let response = self.dht.request(self.request, addr).await.map(|r| {
|
let response = self.dht.request(self.request.clone(), addr).await.map(|r| {
|
||||||
self.mark_node_responded(addr, &r);
|
self.mark_node_responded(addr, &r);
|
||||||
r
|
r
|
||||||
});
|
});
|
||||||
|
|
@ -359,7 +400,7 @@ impl<C: RecursiveRequestCallbacks> RecursiveRequest<C> {
|
||||||
self.callbacks.on_request_end(self, id, addr, &response);
|
self.callbacks.on_request_end(self, id, addr, &response);
|
||||||
}
|
}
|
||||||
|
|
||||||
let response = match self.dht.request(self.request, addr).await {
|
let response = match self.dht.request(self.request.clone(), addr).await {
|
||||||
Ok(ResponseOrError::Response(r)) => r,
|
Ok(ResponseOrError::Response(r)) => r,
|
||||||
Ok(ResponseOrError::Error(e)) => bail!("error response: {:?}", e),
|
Ok(ResponseOrError::Error(e)) => bail!("error response: {:?}", e),
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
|
|
@ -493,6 +534,7 @@ pub struct DhtState {
|
||||||
worker_sender: UnboundedSender<WorkerSendRequest>,
|
worker_sender: UnboundedSender<WorkerSendRequest>,
|
||||||
|
|
||||||
pub(crate) peer_store: PeerStore,
|
pub(crate) peer_store: PeerStore,
|
||||||
|
announce_addr: Option<SocketAddrV4>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl DhtState {
|
impl DhtState {
|
||||||
|
|
@ -502,6 +544,7 @@ impl DhtState {
|
||||||
routing_table: Option<RoutingTable>,
|
routing_table: Option<RoutingTable>,
|
||||||
listen_addr: SocketAddr,
|
listen_addr: SocketAddr,
|
||||||
peer_store: PeerStore,
|
peer_store: PeerStore,
|
||||||
|
announce_addr: Option<SocketAddrV4>,
|
||||||
) -> Self {
|
) -> Self {
|
||||||
let routing_table = routing_table.unwrap_or_else(|| RoutingTable::new(id, None));
|
let routing_table = routing_table.unwrap_or_else(|| RoutingTable::new(id, None));
|
||||||
Self {
|
Self {
|
||||||
|
|
@ -513,6 +556,7 @@ impl DhtState {
|
||||||
listen_addr,
|
listen_addr,
|
||||||
rate_limiter: make_rate_limiter(),
|
rate_limiter: make_rate_limiter(),
|
||||||
peer_store,
|
peer_store,
|
||||||
|
announce_addr,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -581,6 +625,22 @@ impl DhtState {
|
||||||
ip: None,
|
ip: None,
|
||||||
kind: MessageKind::PingRequest(PingRequest { id: self.id }),
|
kind: MessageKind::PingRequest(PingRequest { id: self.id }),
|
||||||
},
|
},
|
||||||
|
Request::Announce {
|
||||||
|
info_hash,
|
||||||
|
token,
|
||||||
|
addr,
|
||||||
|
} => Message {
|
||||||
|
kind: MessageKind::AnnouncePeer(AnnouncePeer {
|
||||||
|
id: self.id,
|
||||||
|
implied_port: 0,
|
||||||
|
info_hash,
|
||||||
|
port: addr.port(),
|
||||||
|
token,
|
||||||
|
}),
|
||||||
|
transaction_id: ByteString::from(transaction_id_buf.as_ref()),
|
||||||
|
version: None,
|
||||||
|
ip: None,
|
||||||
|
},
|
||||||
};
|
};
|
||||||
(transaction_id, message)
|
(transaction_id, message)
|
||||||
}
|
}
|
||||||
|
|
@ -744,10 +804,15 @@ impl DhtState {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
|
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
|
||||||
enum Request {
|
enum Request {
|
||||||
GetPeers(Id20),
|
GetPeers(Id20),
|
||||||
FindNode(Id20),
|
FindNode(Id20),
|
||||||
|
Announce {
|
||||||
|
info_hash: Id20,
|
||||||
|
token: ByteString,
|
||||||
|
addr: SocketAddrV4,
|
||||||
|
},
|
||||||
Ping,
|
Ping,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -1095,6 +1160,13 @@ impl DhtState {
|
||||||
config.routing_table,
|
config.routing_table,
|
||||||
listen_addr,
|
listen_addr,
|
||||||
config.peer_store.unwrap_or_else(|| PeerStore::new(peer_id)),
|
config.peer_store.unwrap_or_else(|| PeerStore::new(peer_id)),
|
||||||
|
config.announce_addr.and_then(|a| match a {
|
||||||
|
SocketAddr::V4(v4) => Some(v4),
|
||||||
|
SocketAddr::V6(_) => {
|
||||||
|
warn!("libqrqbit-dht doesn't support announcing IPv6 addresses");
|
||||||
|
None
|
||||||
|
}
|
||||||
|
}),
|
||||||
));
|
));
|
||||||
|
|
||||||
spawn(error_span!("dht"), {
|
spawn(error_span!("dht"), {
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue