|
| 1 | +use discv5::{ |
| 2 | + enr::NodeId, |
| 3 | + kbucket::{ |
| 4 | + Entry, FailureReason, Filter, InsertResult, KBucketsTable, Key, NodeStatus, |
| 5 | + MAX_NODES_PER_BUCKET, |
| 6 | + }, |
| 7 | + ConnectionDirection, ConnectionState, TalkRequest, |
| 8 | +}; |
| 9 | +use futures::channel::oneshot; |
| 10 | +use parking_lot::RwLock; |
| 11 | +use ssz::Encode; |
1 | 12 | use std::{ |
2 | 13 | collections::{BTreeMap, HashSet}, |
3 | 14 | fmt::{Debug, Display}, |
4 | 15 | marker::{PhantomData, Sync}, |
5 | 16 | sync::Arc, |
6 | 17 | time::Duration, |
7 | 18 | }; |
8 | | - |
9 | | -use discv5::{ |
10 | | - enr::NodeId, |
11 | | - kbucket::{Filter, KBucketsTable, NodeStatus, MAX_NODES_PER_BUCKET}, |
12 | | - TalkRequest, |
13 | | -}; |
14 | | -use futures::channel::oneshot; |
15 | | -use parking_lot::RwLock; |
16 | | -use ssz::Encode; |
17 | 19 | use tokio::sync::mpsc::UnboundedSender; |
18 | 20 | use tracing::{debug, error, info, warn}; |
19 | 21 | use utp_rs::socket::UtpSocket; |
@@ -276,6 +278,97 @@ where |
276 | 278 | .collect() |
277 | 279 | } |
278 | 280 |
|
| 281 | + /// `AddEnr` adds requested `enr` to our kbucket. |
| 282 | + pub fn add_enr(&self, enr: Enr) -> Result<(), OverlayRequestError> { |
| 283 | + let key = Key::from(enr.node_id()); |
| 284 | + match self.kbuckets.write().insert_or_update( |
| 285 | + &key, |
| 286 | + Node { |
| 287 | + enr, |
| 288 | + data_radius: Distance::MAX, |
| 289 | + }, |
| 290 | + NodeStatus { |
| 291 | + state: ConnectionState::Disconnected, |
| 292 | + direction: ConnectionDirection::Incoming, |
| 293 | + }, |
| 294 | + ) { |
| 295 | + InsertResult::Inserted |
| 296 | + | InsertResult::Pending { .. } |
| 297 | + | InsertResult::StatusUpdated { .. } |
| 298 | + | InsertResult::ValueUpdated |
| 299 | + | InsertResult::Updated { .. } |
| 300 | + | InsertResult::UpdatedPending => Ok(()), |
| 301 | + InsertResult::Failed(FailureReason::BucketFull) => { |
| 302 | + Err(OverlayRequestError::Failure("The bucket was full.".into())) |
| 303 | + } |
| 304 | + InsertResult::Failed(FailureReason::BucketFilter) => Err(OverlayRequestError::Failure( |
| 305 | + "The node didn't pass the bucket filter.".into(), |
| 306 | + )), |
| 307 | + InsertResult::Failed(FailureReason::TableFilter) => Err(OverlayRequestError::Failure( |
| 308 | + "The node didn't pass the table filter.".into(), |
| 309 | + )), |
| 310 | + InsertResult::Failed(FailureReason::InvalidSelfUpdate) => { |
| 311 | + Err(OverlayRequestError::Failure("Cannot update self.".into())) |
| 312 | + } |
| 313 | + InsertResult::Failed(_) => { |
| 314 | + Err(OverlayRequestError::Failure("Failed to insert ENR".into())) |
| 315 | + } |
| 316 | + } |
| 317 | + } |
| 318 | + |
| 319 | + /// `GetEnr` gets requested `enr` from our kbucket. |
| 320 | + pub fn get_enr(&self, node_id: NodeId) -> Result<Enr, OverlayRequestError> { |
| 321 | + if node_id == self.local_enr().node_id() { |
| 322 | + return Ok(self.local_enr()); |
| 323 | + } |
| 324 | + let key = Key::from(node_id); |
| 325 | + if let Entry::Present(entry, _) = self.kbuckets.write().entry(&key) { |
| 326 | + return Ok(entry.value().enr()); |
| 327 | + } |
| 328 | + Err(OverlayRequestError::Failure("Couldn't get ENR".into())) |
| 329 | + } |
| 330 | + |
| 331 | + /// `DeleteEnr` deletes requested `enr` from our kbucket. |
| 332 | + pub fn delete_enr(&self, node_id: NodeId) -> bool { |
| 333 | + let key = &Key::from(node_id); |
| 334 | + self.kbuckets.write().remove(key) |
| 335 | + } |
| 336 | + |
| 337 | + /// `LookupEnr` finds requested `enr` from our kbucket, FindNode, and RecursiveFindNode. |
| 338 | + pub async fn lookup_enr(&self, node_id: NodeId) -> Result<Enr, OverlayRequestError> { |
| 339 | + if node_id == self.local_enr().node_id() { |
| 340 | + return Ok(self.local_enr()); |
| 341 | + } |
| 342 | + |
| 343 | + let enr = self.get_enr(node_id); |
| 344 | + |
| 345 | + // try to find more up to date enr |
| 346 | + if let Ok(enr) = enr.clone() { |
| 347 | + let nodes = self.send_find_nodes(enr, vec![0]).await?; |
| 348 | + let enr_highest_seq = nodes.enrs.into_iter().max_by(|a, b| a.seq().cmp(&b.seq())); |
| 349 | + |
| 350 | + if let Some(enr_highest_seq) = enr_highest_seq { |
| 351 | + return Ok(enr_highest_seq.into()); |
| 352 | + } |
| 353 | + } |
| 354 | + |
| 355 | + let lookup_node_enr = self.lookup_node(node_id).await; |
| 356 | + let lookup_node_enr = lookup_node_enr |
| 357 | + .into_iter() |
| 358 | + .max_by(|a, b| a.seq().cmp(&b.seq())); |
| 359 | + if let Some(lookup_node_enr) = lookup_node_enr { |
| 360 | + let mut enr_seq = 0; |
| 361 | + if let Ok(enr) = enr.clone() { |
| 362 | + enr_seq = enr.seq(); |
| 363 | + } |
| 364 | + if lookup_node_enr.seq() > enr_seq { |
| 365 | + return Ok(lookup_node_enr); |
| 366 | + } |
| 367 | + } |
| 368 | + |
| 369 | + enr |
| 370 | + } |
| 371 | + |
279 | 372 | /// Sends a `Ping` request to `enr`. |
280 | 373 | pub async fn send_ping(&self, enr: Enr) -> Result<Pong, OverlayRequestError> { |
281 | 374 | // Construct the request. |
|
0 commit comments