1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
//! A walk's state and the frontier read off it (bl-d00f, bl-d9c1): every node
//! heard of, which were asked, which replied — and, against the flight, the
//! next node to ask and whether anything on the frontier is still in the air.
//! `lookup` is the loop that drives it.
use super::bencode::Dict;
use super::flight::{Flight, Query};
use super::krpc::{self, Message, Node, NodeId};
use std::collections::{BTreeMap, BTreeSet};
use std::net::SocketAddr;
/// What a walk found: every reply, closest first, and every error a node
/// answered instead of a reply.
#[derive(Debug, Default)]
pub(crate) struct Outcome {
pub(crate) replies: Vec<(Node, Dict)>,
pub(crate) errors: Vec<String>,
}
/// The frontier as it stands between two events of the walk.
pub(crate) struct Frontier {
/// The closest frontier node not yet asked.
pub(crate) next: Option<SocketAddr>,
/// A frontier node is in the air: its answer may still move the frontier.
pub(crate) waiting: bool,
/// Walk queries in the air — the door's are not, and hold no α slot.
pub(crate) walking: usize,
/// A door query is in the air: its answer may still seed the pool.
pub(crate) door: bool,
}
pub(crate) struct Walk {
target: NodeId,
k: usize,
pub(crate) asked: BTreeSet<SocketAddr>,
replied: BTreeSet<SocketAddr>,
/// Door addresses that have answered this walk — the only ones a
/// re-ask spends a query on (bl-f519).
opened: BTreeSet<SocketAddr>,
/// Keyed `(distance, address)`: a repeated entry is learned once, one id
/// at two addresses is two nodes, and iteration is closest first.
pool: BTreeMap<([u8; 20], SocketAddr), Node>,
pub(crate) out: Outcome,
pub(crate) claims: Vec<SocketAddr>,
}
impl Walk {
pub(crate) fn new(target: NodeId, k: usize) -> Walk {
Walk {
target,
k,
asked: BTreeSet::new(),
replied: BTreeSet::new(),
opened: BTreeSet::new(),
pool: BTreeMap::new(),
out: Outcome::default(),
claims: Vec::new(),
}
}
/// The frontier: the K closest nodes that replied, are not yet asked, or
/// are still in the air. A node asked and silent past its deadline, one
/// that answered only an error, and one the socket refused have left it,
/// so none holds a slot a live node past it would take.
pub(crate) fn frontier(&self, flight: &Flight) -> Frontier {
let airborne: BTreeSet<SocketAddr> = flight
.values()
.filter(|q| !q.door)
.map(|q| q.addr)
.collect();
let front: Vec<SocketAddr> = self
.pool
.values()
.map(|n| n.addr)
.filter(|a| self.replied.contains(a) || !self.asked.contains(a) || airborne.contains(a))
.take(self.k)
.collect();
Frontier {
next: front.iter().copied().find(|a| !self.asked.contains(a)),
waiting: front.iter().any(|a| airborne.contains(a)),
walking: flight.values().filter(|q| !q.door).count(),
door: flight.values().any(|q| q.door),
}
}
/// Whether a door address is worth a query: never asked, or it answered.
/// A router silent past its deadline has left the door for this walk,
/// as a silent node leaves the frontier (bl-f519).
pub(crate) fn knocks(&self, addr: SocketAddr) -> bool {
!self.asked.contains(&addr) || self.opened.contains(&addr)
}
/// Whether the door should be asked again: it has named someone, and
/// fewer than K nodes past it have replied.
pub(crate) fn dry(&self) -> bool {
!self.pool.is_empty() && self.replied.len() < self.k
}
/// Take one answer in. The door's reply seeds the pool and its claim
/// votes, but the door is never a result; a walk node's reply joins the
/// pool, the replied set and the outcome, and its error the outcome's
/// errors. A reply that names no id is heard and is nothing.
pub(crate) fn heard(&mut self, query: Query, message: Message) {
match message {
Message::Reply { r, ip, .. } => {
let Some(id) = r
.get(b"id".as_slice())
.and_then(|v| v.as_bytes())
.and_then(NodeId::parse)
else {
return;
};
let node = Node {
id,
addr: query.addr,
};
self.claims.extend(ip);
let this = (!query.door).then_some(node);
for near in krpc::nodes_of(&r).into_iter().chain(this) {
self.pool
.insert((near.id.distance(&self.target), near.addr), near);
}
if query.door {
self.opened.insert(query.addr);
} else {
self.replied.insert(query.addr);
self.out.replies.push((node, r));
}
}
Message::Error { code, message, .. } if !query.door => {
self.out
.errors
.push(format!("{}: {code} {message}", query.addr));
}
Message::Error { .. } => {}
}
}
}